Packages
setu_client
1.0.0
Production-grade Elixir client for the Setu API platform (UPI, BBPS, WhatsApp, Account Aggregator, KYC, eSign)
Current section
Files
Jump to
Current section
Files
lib/setu_client/rate_limiter.ex
defmodule SetuClient.RateLimiter do
@moduledoc """
Client-side token-bucket rate limiter implemented as an OTP GenServer.
Each unique `{client_id, environment}` pair gets an independent bucket.
Limits come from `SetuClient.Config` fields `:rate_limit_rps` and `:rate_limit_burst`.
"""
use GenServer
alias SetuClient.Config
alias SetuClient.Telemetry
@check_interval_ms 50
@type bucket_key :: {String.t(), atom()}
@type bucket :: %{tokens: float(), last_refill: integer()}
@type state :: %{buckets: %{optional(bucket_key()) => bucket()}}
# ── Public API ─────────────────────────────────────────────────────────────
@doc false
@spec start_link(keyword()) :: GenServer.on_start()
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@doc """
Acquires one token for `cfg`, blocking if the bucket is empty.
Returns `:ok` when a token is granted, `{:error, :timeout}` after `cfg.timeout` ms.
"""
@spec acquire(Config.t()) :: :ok | {:error, :timeout}
def acquire(%Config{} = cfg) do
deadline = System.monotonic_time(:millisecond) + cfg.timeout
do_acquire(bucket_key(cfg), cfg, deadline)
end
# ── GenServer callbacks ────────────────────────────────────────────────────
@impl GenServer
@spec init(keyword()) :: {:ok, state()}
def init(_opts), do: {:ok, %{buckets: %{}}}
@impl GenServer
def handle_call({:try_acquire, key, rps, burst}, _from, state) do
now_ms = System.monotonic_time(:millisecond)
default_bucket = %{tokens: burst * 1.0, last_refill: now_ms}
bucket = Map.get(state.buckets, key, default_bucket)
elapsed_s = (now_ms - bucket.last_refill) / 1_000
refilled = min(burst * 1.0, bucket.tokens + rps * elapsed_s)
{reply, new_tokens} =
if refilled >= 1.0 do
{:ok, refilled - 1.0}
else
{:empty, refilled}
end
new_state = put_in(state.buckets[key], %{tokens: new_tokens, last_refill: now_ms})
{:reply, reply, new_state}
end
# ── Private ────────────────────────────────────────────────────────────────
@spec do_acquire(bucket_key(), Config.t(), integer()) :: :ok | {:error, :timeout}
defp do_acquire(key, cfg, deadline) do
now = System.monotonic_time(:millisecond)
if now > deadline do
{:error, :timeout}
else
case GenServer.call(
__MODULE__,
{:try_acquire, key, cfg.rate_limit_rps, cfg.rate_limit_burst}
) do
:ok ->
:ok
:empty ->
Telemetry.emit_rate_limit_wait(@check_interval_ms)
Process.sleep(@check_interval_ms)
do_acquire(key, cfg, deadline)
end
end
end
@spec bucket_key(Config.t()) :: bucket_key()
defp bucket_key(%Config{client_id: id, environment: env}), do: {id || "", env}
end