Per-credential token-bucket rate limiter for exchange API requests.
Each {exchange_id, credential_key, bucket_axis} key holds tokens that refill
at the authored refill_per_sec, capped at capacity (max_size). A request
is admitted when the bucket holds its cost. An over-capacity cost is not
exempt: the bucket accrues until it can pay, then goes to zero.
A waiter that will actually sleep (an explicit or default wait budget that
covers the delay) reserves the unpaid cost so interleaved cheaper traffic
cannot clamp accrual or spend the reserved tokens. Cheap burst stays capped
at authored capacity. Repeated checks of one key in a single check_rates/2
call are charged once at their combined cost; conflicting bucket definitions
for the same key fail before any spend.
Usage
key = {"okx", api_key, "request"}
case Bourse.RateLimiter.check_rate(key, %{capacity: 1, refill_per_sec: 9.09}, 1) do
:ok -> make_request()
{:delay, ms} -> Process.sleep(ms); make_request()
end%{requests: max, period: period_ms} is accepted as capacity max refilling
at max / (period_ms / 1000) tokens per second.
Credential Keys
The key is a tuple {exchange_id, credential_key, bucket_axis} where:
exchange_idis the exchange string ID ("binance","bybit", etc.)credential_keyis either:- The API key string (for authenticated requests) -- isolates per-user limits
:publicatom (for public requests) -- shared pool for unauthenticated requests
bucket_axisis the spec/header bucket axis ("request","ip","uid","order_weight", etc.)
Summary
Types
A single bucket capacity check: {key, rate_limit, cost}.
Rate limiter key: {exchange_id, api_key | :public, bucket_axis}.
Token-bucket configuration: authored capacity and refill rate.
Functions
Checks if a request can be made within rate limits.
Checks multiple bucket capacities atomically.
Returns a child specification for starting the rate limiter under a supervisor.
Gets tokens currently borrowed from the bucket (capacity - tokens after refill).
Records a request for a key with specified cost.
Resets rate limit tracking for a key.
Clears all rate-limit tracking state.
Clears rate-limit tracking for every bucket belonging to one exchange.
Starts the rate limiter.
Blocks until rate limit capacity is available, then records the request.
Types
@type bucket_check() :: {key(), rate_limit() | nil, number()}
A single bucket capacity check: {key, rate_limit, cost}.
Rate limiter key: {exchange_id, api_key | :public, bucket_axis}.
@type rate_limit() :: %{capacity: number(), refill_per_sec: number()} | %{requests: number(), period: pos_integer()} | %{requests: number()}
Token-bucket configuration: authored capacity and refill rate.
Functions
@spec check_rate(key(), rate_limit() | nil, number(), GenServer.name()) :: :ok | {:delay, pos_integer()}
Checks if a request can be made within rate limits.
Returns :ok if the bucket holds cost (and records it), or
{:delay, milliseconds} if the caller should wait for tokens to accrue.
A cost larger than capacity is limited: the caller waits until the bucket has accrued that cost. There is no skip-record exemption.
@spec check_rates([bucket_check()]) :: :ok | {:delay, pos_integer()} | {:error, Bourse.Error.t()}
Checks multiple bucket capacities atomically.
Returns :ok only when every bucket has capacity, recording every cost in the
same GenServer transition. Returns {:delay, milliseconds} without recording
any spend when at least one bucket cannot yet pay.
Repeated checks of the same key are coalesced into one combined cost. Two
checks that name the same key with different capacity or refill_per_sec
return {:error, %Bourse.Error{type: :invalid_parameters}} and spend nothing.
Pass max_wait_ms: when the caller will sleep the returned delay — the unpaid
cost is reserved so cheaper traffic cannot erase accrual. Probe-style calls
(no budget) do not reserve.
@spec check_rates([bucket_check()], GenServer.name() | keyword()) :: :ok | {:delay, pos_integer()} | {:error, Bourse.Error.t()}
@spec check_rates([bucket_check()], GenServer.name(), keyword()) :: :ok | {:delay, pos_integer()} | {:error, Bourse.Error.t()}
@spec child_spec(keyword()) :: Supervisor.child_spec()
Returns a child specification for starting the rate limiter under a supervisor.
@spec get_cost(key(), pos_integer(), GenServer.name()) :: number()
Gets tokens currently borrowed from the bucket (capacity - tokens after refill).
The period argument is unused; it remains so callers that passed a window
length keep compiling. Useful for debugging and monitoring.
@spec record_request(key(), number(), GenServer.name()) :: :ok
Records a request for a key with specified cost.
Called automatically by check_rate/4 when it returns :ok.
Exposed for manual tracking if needed.
@spec reset(key(), GenServer.name()) :: :ok
Resets rate limit tracking for a key.
@spec reset_all(GenServer.name()) :: :ok
Clears all rate-limit tracking state.
Whole-map reset; prefer reset_exchange/2 when only one venue's buckets
should be cleared.
@spec reset_exchange(String.t(), GenServer.name()) :: :ok
Clears rate-limit tracking for every bucket belonging to one exchange.
Bourse.Test.LiveGateIsolation uses this so a venue probe cannot enter its
bucket with capacity another probe already spent (Task 179) without
discarding the sibling venues' pacing at the same time — a global wipe lets a
heavy endpoint (okx system/status, authored cost 50 against a 9.09/s drain)
go out back to back and earn the venue's own 50011.
@spec start_link(keyword()) :: GenServer.on_start()
Starts the rate limiter.
@spec wait_for_capacity( key(), rate_limit() | nil, number(), GenServer.name(), keyword() ) :: :ok | {:error, Bourse.Error.t()}
Blocks until rate limit capacity is available, then records the request.
Returns {:error, %Bourse.Error{}} when the wait would exceed the per-call
budget (:max_wait_ms, default Bourse.Defaults.rate_limit_max_wait_ms/0).
The refusal names the required wait and spends nothing.