162 lines
5.2 KiB
Elixir
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
|