Bourse.RateLimiter (bourse v0.9.0)

Copy Markdown View Source

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_id is the exchange string ID ("binance", "bybit", etc.)
  • credential_key is either:
    • The API key string (for authenticated requests) -- isolates per-user limits
    • :public atom (for public requests) -- shared pool for unauthenticated requests
  • bucket_axis is 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

bucket_check()

@type bucket_check() :: {key(), rate_limit() | nil, number()}

A single bucket capacity check: {key, rate_limit, cost}.

key()

@type key() ::
  {String.t(), String.t() | :public}
  | {String.t(), String.t() | :public, String.t()}

Rate limiter key: {exchange_id, api_key | :public, bucket_axis}.

rate_limit()

@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

check_rate(key, rate_limit, cost \\ 1, name \\ __MODULE__)

@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.

check_rates(bucket_checks)

@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.

check_rates(bucket_checks, opts)

@spec check_rates([bucket_check()], GenServer.name() | keyword()) ::
  :ok | {:delay, pos_integer()} | {:error, Bourse.Error.t()}

check_rates(bucket_checks, name, opts)

@spec check_rates([bucket_check()], GenServer.name(), keyword()) ::
  :ok | {:delay, pos_integer()} | {:error, Bourse.Error.t()}

child_spec(init_arg)

@spec child_spec(keyword()) :: Supervisor.child_spec()

Returns a child specification for starting the rate limiter under a supervisor.

get_cost(key, period, name \\ __MODULE__)

@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.

record_request(key, cost \\ 1, name \\ __MODULE__)

@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.

reset(key, name \\ __MODULE__)

@spec reset(key(), GenServer.name()) :: :ok

Resets rate limit tracking for a key.

reset_all(name \\ __MODULE__)

@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.

reset_exchange(exchange_id, name \\ __MODULE__)

@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.

start_link(opts \\ [])

@spec start_link(keyword()) :: GenServer.on_start()

Starts the rate limiter.

wait_for_capacity(key, rate_limit, cost \\ 1, name \\ __MODULE__, opts \\ [])

@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.