src/transport/limiter.ts
v0.3.0 · 13.6 KB
// Weight-based token bucket + priority concurrency semaphore for one egress IP.
//
// Hyperliquid documents ~1200 weight/min per IP (/info and /exchange combined), but it also
// enforces a shorter window. Sustained rate, burst capacity and concurrency are application policy;
// callers must configure them explicitly instead of inheriting one workload's tuned defaults.
import { MAX_INFO_WEIGHT } from './weights.js';
/**
* Queue priority.
*
* - `urgent`: reduce-only closes and protective orders. Without it a CLOSE can wait behind a batch
* of INCREASE orders and the exit comes late.
* - `high`: other orders and reads on the critical trading path.
* - `normal`: background reads, analytics.
*/
export type Priority = 'urgent' | 'high' | 'normal';
const TIERS: readonly Priority[] = ['urgent', 'high', 'normal'];
/** Options for {@link createWeightLimiter}. */
export interface WeightLimiterOptions {
/** Caller-selected sustained weight per minute (documented ceiling: about 1200 per IP). */
weightPerMinute: number;
/** Caller-selected bucket capacity. It must admit every request the caller sends. */
burstCapacity: number;
/** Caller-selected concurrent in-flight HTTP requests. */
maxConcurrent: number;
/** Monotonic clock in ms. Default `performance.now()`: a wall-clock NTP jump must not mint or freeze tokens. */
now?: () => number;
}
/** Options for a single acquisition. */
export interface AcquireOptions {
priority?: Priority;
/** Removes the waiter from the queue and rejects with the signal's reason. */
signal?: AbortSignal;
}
/** Snapshot of limiter state for health logs. */
export interface WeightLimiterStats {
/** Current tokens (may be negative after post-response charges). */
tokens: number;
capacity: number;
weightPerMinute: number;
urgentQueued: number;
highQueued: number;
normalQueued: number;
/** Requests currently holding a concurrency slot. */
inFlight: number;
/** Requests waiting for a concurrency slot. */
slotQueued: number;
/** How many acquisitions had to wait for tokens since creation. */
totalQueued: number;
/** Total weight consumed (acquired + charged) since creation. */
totalWeight: number;
}
/**
* Shared per-IP limiter. Create ONE per egress IP and route every REST call from that IP through
* it (all subsystems: info reads, exchange actions, scans). Two independent limiters on one IP
* together exceed the limit, and a close can then get stuck behind the 429s.
*/
export interface WeightLimiter {
readonly weightPerMinute: number;
readonly burstCapacity: number;
readonly maxConcurrent: number;
/**
* Waits until `weight` tokens are available (strict priority across tiers, FIFO within a tier)
* and consumes them. Rejects immediately when `weight > burstCapacity`: such a request could
* never be admitted and would hang forever without an error.
*/
acquire(weight: number, opts?: AcquireOptions): Promise<void>;
/** Runs `fn` inside a concurrency slot; slots are handed over by priority, atomically. */
withSlot<T>(fn: () => Promise<T>, opts?: AcquireOptions): Promise<T>;
/** `acquire(weight)` then `withSlot(fn)`: the usual way to send one request. */
schedule<T>(weight: number, fn: () => Promise<T>, opts?: AcquireOptions): Promise<T>;
/**
* Consumes `weight` tokens without waiting (the balance may go negative). Use it to pay a
* surcharge learned only after the response, e.g. per returned item.
*/
charge(weight: number): void;
stats(): WeightLimiterStats;
}
interface TokenWaiter {
weight: number;
resolve: () => void;
reject: (err: unknown) => void;
cleanup: () => void;
}
interface SlotWaiter {
resolve: () => void;
reject: (err: unknown) => void;
cleanup: () => void;
}
function abortReason(signal: AbortSignal): unknown {
return signal.reason ?? new DOMException('This operation was aborted', 'AbortError');
}
function assertPositive(name: string, value: number, integer: boolean): void {
if (!Number.isFinite(value) || value <= 0 || (integer && !Number.isInteger(value))) {
throw new RangeError(`${name} must be a positive ${integer ? 'integer' : 'number'}, got ${value}`);
}
}
/**
* Creates a token bucket by request weight plus a priority semaphore for concurrency.
*
* Three mechanisms work together: the bucket smooths the average rate (callers wait in a queue
* instead of getting 429), the semaphore bounds sockets and bursts of heavy reads, and retries
* with jittered backoff (see `withRetry`) cover what the short exchange window still rejects.
*/
export function createWeightLimiter(options: WeightLimiterOptions): WeightLimiter {
if (
!options ||
options.weightPerMinute === undefined ||
options.burstCapacity === undefined ||
options.maxConcurrent === undefined
) {
throw new TypeError('createWeightLimiter requires explicit weightPerMinute, burstCapacity and maxConcurrent');
}
const weightPerMinute = options.weightPerMinute;
const capacity = options.burstCapacity;
const maxConcurrent = options.maxConcurrent;
assertPositive('weightPerMinute', weightPerMinute, false);
assertPositive('burstCapacity', capacity, false);
assertPositive('maxConcurrent', maxConcurrent, true);
const now = options.now ?? (() => performance.now());
const ratePerMs = weightPerMinute / 60_000;
let tokens = capacity;
let last = now();
let timer: ReturnType<typeof setTimeout> | null = null;
let totalQueued = 0;
let totalWeight = 0;
const tokenQueues: Record<Priority, TokenWaiter[]> = { urgent: [], high: [], normal: [] };
let active = 0;
const slotQueues: Record<Priority, SlotWaiter[]> = { urgent: [], high: [], normal: [] };
function refill(): void {
const t = now();
const elapsed = t - last;
last = t;
if (elapsed > 0) tokens = Math.min(capacity, tokens + elapsed * ratePerMs);
}
function headQueue(): TokenWaiter[] | null {
for (const tier of TIERS) if (tokenQueues[tier].length > 0) return tokenQueues[tier];
return null;
}
function drain(): void {
refill();
let q = headQueue();
// Strictly by tier, FIFO inside a tier: a light normal read must not jump ahead of a queued close.
while (q) {
const head = q[0];
if (!head || tokens < head.weight) break;
q.shift();
tokens -= head.weight;
totalWeight += head.weight;
head.cleanup();
head.resolve();
q = headQueue();
}
const head = q?.[0];
if (head && timer === null) {
const need = head.weight - tokens;
const waitMs = Math.max(5, Math.ceil(need / ratePerMs));
timer = setTimeout(() => {
timer = null;
drain();
}, waitMs);
} else if (!head && timer !== null) {
clearTimeout(timer);
timer = null;
}
}
function acquire(weight: number, opts: AcquireOptions = {}): Promise<void> {
if (!Number.isFinite(weight) || weight < 0) {
return Promise.reject(new RangeError(`request weight must be a non-negative number, got ${weight}`));
}
if (weight > capacity) {
// A waiter heavier than the bucket can never be admitted (tokens are clamped to capacity).
return Promise.reject(
new RangeError(`request weight ${weight} exceeds the bucket capacity ${capacity} (chunk the request)`),
);
}
const signal = opts.signal;
if (signal?.aborted) return Promise.reject(abortReason(signal));
const priority = opts.priority ?? 'normal';
refill();
const nobodyWaiting = headQueue() === null;
if (nobodyWaiting && tokens >= weight) {
tokens -= weight;
totalWeight += weight;
return Promise.resolve();
}
totalQueued++;
return new Promise<void>((resolve, reject) => {
const waiter: TokenWaiter = { weight, resolve, reject, cleanup: () => {} };
if (signal) {
const onAbort = () => {
const q = tokenQueues[priority];
const i = q.indexOf(waiter);
if (i >= 0) q.splice(i, 1);
reject(abortReason(signal));
drain(); // the removed waiter may have been blocking the head of the line
};
signal.addEventListener('abort', onAbort, { once: true });
waiter.cleanup = () => signal.removeEventListener('abort', onAbort);
}
tokenQueues[priority].push(waiter);
drain();
});
}
function releaseSlot(): void {
// Atomic, priority-ordered hand-over: the slot passes directly to the next waiter without
// `active--`, so a fresh caller cannot slip in between and push `active` above the max.
// A plain FIFO here would break the bucket's priorities (a CLOSE would wait behind admitted reads).
for (const tier of TIERS) {
const next = slotQueues[tier].shift();
if (next) {
next.cleanup();
next.resolve();
return;
}
}
active--;
}
async function withSlot<T>(fn: () => Promise<T>, opts: AcquireOptions = {}): Promise<T> {
const signal = opts.signal;
if (signal?.aborted) throw abortReason(signal);
if (active < maxConcurrent) {
active++;
} else {
const priority = opts.priority ?? 'normal';
await new Promise<void>((resolve, reject) => {
const waiter: SlotWaiter = { resolve, reject, cleanup: () => {} };
if (signal) {
const onAbort = () => {
const q = slotQueues[priority];
const i = q.indexOf(waiter);
if (i >= 0) q.splice(i, 1);
reject(abortReason(signal));
};
signal.addEventListener('abort', onAbort, { once: true });
waiter.cleanup = () => signal.removeEventListener('abort', onAbort);
}
slotQueues[priority].push(waiter);
});
// Woken by a hand-over: the slot is already ours, `active` is unchanged.
}
try {
return await fn();
} finally {
releaseSlot();
}
}
async function schedule<T>(weight: number, fn: () => Promise<T>, opts: AcquireOptions = {}): Promise<T> {
await acquire(weight, opts);
return withSlot(fn, opts);
}
function charge(weight: number): void {
if (!Number.isFinite(weight) || weight <= 0) return;
refill();
tokens -= weight;
totalWeight += weight;
if (timer !== null) {
// The pending wake-up was computed for the old balance; recompute it.
clearTimeout(timer);
timer = null;
drain();
}
}
function stats(): WeightLimiterStats {
refill();
return {
tokens: Math.floor(tokens),
capacity,
weightPerMinute,
urgentQueued: tokenQueues.urgent.length,
highQueued: tokenQueues.high.length,
normalQueued: tokenQueues.normal.length,
inFlight: active,
slotQueued: slotQueues.urgent.length + slotQueues.high.length + slotQueues.normal.length,
totalQueued,
totalWeight,
};
}
return { weightPerMinute, burstCapacity: capacity, maxConcurrent, acquire, withSlot, schedule, charge, stats };
}
const REGISTRY_KEY = Symbol.for('@markpaper/hl-kit/weight-limiters');
function registry(): Map<string, { limiter: WeightLimiter; options: WeightLimiterOptions }> {
const g = globalThis as unknown as Record<symbol, Map<string, { limiter: WeightLimiter; options: WeightLimiterOptions }>>;
let map = g[REGISTRY_KEY];
if (!map) {
map = new Map();
g[REGISTRY_KEY] = map;
}
return map;
}
/** Egress key used by default: the process's default outbound IP. */
export const DEFAULT_EGRESS_KEY = 'default';
/**
* Process-wide limiter for one egress IP. The first call for a key must provide explicit options;
* later calls may omit them to retrieve the existing limiter. The registry lives on `globalThis`,
* so even duplicated copies of this package in `node_modules` share one bucket per key.
*
* Use a distinct `egressKey` ONLY for traffic that really leaves from a different IP (verified at
* startup with an echo request through that lane's `fetch`). Two buckets over one physical IP give
* double the rate against a single limit - a guaranteed 429 storm. Node/undici silently ignores
* `new Agent({ connect: { localAddress } })`, so a misconfigured lane does not raise an error.
*
* Other processes on the same machine are invisible to this registry: give them their own IP or a
* smaller `weightPerMinute` each.
*
* @throws TypeError when creating a key without explicit options.
* @throws RangeError when the limiter already exists with different explicit options.
*/
export function getSharedWeightLimiter(egressKey: string = DEFAULT_EGRESS_KEY, options?: WeightLimiterOptions): WeightLimiter {
const map = registry();
const existing = map.get(egressKey);
if (existing) {
if (options) {
const a = existing.limiter;
const differs =
(options.weightPerMinute !== undefined && options.weightPerMinute !== a.weightPerMinute) ||
(options.burstCapacity !== undefined && options.burstCapacity !== a.burstCapacity) ||
(options.maxConcurrent !== undefined && options.maxConcurrent !== a.maxConcurrent);
if (differs) {
throw new RangeError(`shared limiter "${egressKey}" already exists with different options`);
}
}
return existing.limiter;
}
if (!options) {
throw new TypeError(
`getSharedWeightLimiter requires explicit options when creating egress key "${egressKey}"`,
);
}
const limiter = createWeightLimiter(options);
map.set(egressKey, { limiter, options });
return limiter;
}
/** Removes shared limiters (all, or one key). Intended for tests and controlled restarts. */
export function resetSharedWeightLimiters(egressKey?: string): void {
const map = registry();
if (egressKey === undefined) map.clear();
else map.delete(egressKey);
}
/** Smallest bucket capacity that admits every known /info request (userRole weighs 60). */
export const MIN_SAFE_BURST_CAPACITY = MAX_INFO_WEIGHT;