who_need_help/lib/mix/tasks/wnh.rate_limit_benchmark.ex

162 lines
5.2 KiB
Elixir

defmodule Mix.Tasks.Wnh.RateLimitBenchmark do
use Mix.Task
import Ecto.Query
alias WhoNeedHelp.Repo
alias WhoNeedHelp.Trust.{RateLimitBucket, RateLimiter}
@shortdoc "Measures the PostgreSQL rate limiter in an isolated database"
@moduledoc """
Measures the shared PostgreSQL limiter without choosing or enforcing a
product threshold.
The caller must provide the workload explicitly. Every run uses a unique
action name, restores the previous application policy, and deletes only the
buckets created by that action before it exits.
mix wnh.rate_limit_benchmark \
--profile hot \
--attempts 5000 \
--concurrency 20 \
--scopes 1 \
--output /benchmark-output/hot.json
"""
@impl Mix.Task
def run(arguments) do
Mix.Task.run("app.start")
{options, rest, invalid} =
OptionParser.parse(arguments,
strict: [
profile: :string,
attempts: :integer,
concurrency: :integer,
scopes: :integer,
output: :string
]
)
if rest != [] or invalid != [],
do: Mix.raise("invalid arguments: #{inspect(rest ++ invalid)}")
profile = Keyword.get(options, :profile) || Mix.raise("--profile is required")
attempts = positive!(options, :attempts)
concurrency = positive!(options, :concurrency)
scopes = positive!(options, :scopes)
output = Keyword.get(options, :output) || Mix.raise("--output PATH is required")
unless profile in ["hot", "distributed"],
do: Mix.raise("--profile must be hot or distributed")
if profile == "hot" and scopes != 1,
do: Mix.raise("the hot profile requires --scopes 1")
if scopes > attempts, do: Mix.raise("--scopes cannot exceed --attempts")
if Mix.env() == :prod,
do: Mix.raise("this benchmark refuses to run with MIX_ENV=prod")
if Mix.env() == :test do
Ecto.Adapters.SQL.Sandbox.mode(Repo, :auto)
end
action = "rate_limit_benchmark_" <> String.replace(Ecto.UUID.generate(), "-", "")
previous_policies = Application.get_env(:who_need_help, :rate_limit_policies, %{})
File.mkdir_p!(Path.dirname(output))
try do
Application.put_env(:who_need_help, :rate_limit_policies, %{
action => %{limit: attempts + 1, window_seconds: 3_600}
})
{elapsed_native, results} =
:timer.tc(fn ->
1..attempts
|> Task.async_stream(
fn index ->
scope = "benchmark-scope-#{rem(index - 1, scopes)}"
RateLimiter.check(action, scope)
end,
max_concurrency: concurrency,
ordered: false,
timeout: :infinity
)
|> Enum.to_list()
end)
counts = summarize_results(results)
bucket_count =
Repo.aggregate(from(bucket in RateLimitBucket, where: bucket.action == ^action), :count)
total_count =
Repo.one!(
from bucket in RateLimitBucket,
where: bucket.action == ^action,
select: coalesce(sum(bucket.count), 0)
)
elapsed_seconds = elapsed_native / 1_000_000
summary = %{
profile: profile,
attempts: attempts,
concurrency: concurrency,
scopes: scopes,
elapsed_ms: Float.round(elapsed_seconds * 1_000, 3),
operations_per_second: Float.round(attempts / elapsed_seconds, 2),
successful_calls: counts.ok,
rate_limited_calls: counts.rate_limited,
failed_tasks: counts.failed,
buckets_created: bucket_count,
persisted_attempt_count: total_count,
scheduler_count: System.schedulers_online(),
postgres_version: Repo.query!("SHOW server_version").rows |> hd() |> hd(),
captured_at: DateTime.utc_now() |> DateTime.truncate(:second) |> DateTime.to_iso8601()
}
{deleted, _} = delete_benchmark_rows(action)
remaining =
Repo.aggregate(from(bucket in RateLimitBucket, where: bucket.action == ^action), :count)
final_summary =
Map.merge(summary, %{cleanup_deleted: deleted, cleanup_remaining: remaining})
File.write!(output, Jason.encode_to_iodata!(final_summary, pretty: true))
Mix.shell().info(Jason.encode!(final_summary))
if counts.rate_limited != 0 or counts.failed != 0 or total_count != attempts or
remaining != 0 do
Mix.raise("rate limiter benchmark correctness checks failed")
end
after
delete_benchmark_rows(action)
Application.put_env(:who_need_help, :rate_limit_policies, previous_policies)
end
end
defp positive!(options, key) do
case Keyword.get(options, key) do
value when is_integer(value) and value > 0 -> value
_ -> Mix.raise("--#{key} must be a positive integer")
end
end
defp summarize_results(results) do
Enum.reduce(results, %{ok: 0, rate_limited: 0, failed: 0}, fn
{:ok, {:ok, _}}, counts -> Map.update!(counts, :ok, &(&1 + 1))
{:ok, {:error, :rate_limited}}, counts -> Map.update!(counts, :rate_limited, &(&1 + 1))
_other, counts -> Map.update!(counts, :failed, &(&1 + 1))
end)
end
defp delete_benchmark_rows(action) do
Repo.delete_all(from bucket in RateLimitBucket, where: bucket.action == ^action)
end
end