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