src/transport/throttle.ts
v0.2.0 · 10.4 KB
// Own throttle for the gateway: one weight bucket for queries, one for executes, one concurrency
// semaphore where executes go first.
//
// Documented limits: queries 2400 weight/min per IP, executes 600/min per wallet (never measured live).
// Sustained rates, burst capacities and concurrency are workload policy. They are required inputs so
// the package does not turn one author's headroom choices into hidden defaults.
//
// A request heavier than the bucket capacity can never be admitted; it is rejected immediately with a
// RangeError instead of hanging the caller forever.
export type Priority = 'high' | 'normal';
/** Explicit options of {@link createWeightThrottle}. */
export interface WeightThrottleOptions {
/** Caller-selected sustained query weight per minute (documented ceiling: 2400 per IP). */
queriesPerMinute: number;
/** Caller-selected query bucket capacity (largest burst). */
queryBurst: number;
/** Caller-selected sustained execute weight per minute (documented ceiling: 600 per wallet). */
executesPerMinute: number;
/** Caller-selected execute bucket capacity. */
executeBurst: number;
/** Caller-selected concurrent in-flight HTTP requests across both families. */
maxConcurrent: number;
/** Monotonic clock in ms. Default `performance.now()`. */
now?: () => number;
}
/** Per-call options. */
export interface ThrottleCallOptions {
/** Counter label (`orders`, `subaccount_info`, `order:close`...). */
label?: string;
/** Removes the waiter from the queues and rejects with the signal's reason. */
signal?: AbortSignal;
/** Queue priority. Default: `'normal'` for queries, `'high'` for executes. */
priority?: Priority;
}
/** Snapshot of one bucket. */
export interface BucketStats {
tokens: number;
capacity: number;
perMinute: number;
queued: number;
totalQueued: number;
totalWeight: number;
}
/** Snapshot of the throttle. */
export interface WeightThrottleStats {
queries: BucketStats;
executes: BucketStats;
inFlight: number;
slotQueued: number;
}
export interface WeightThrottle {
/** Pays `weight` into the query bucket, then runs `fn` inside a concurrency slot (normal priority). */
query<T>(weight: number, fn: () => Promise<T>, opts?: ThrottleCallOptions): Promise<T>;
/** Pays `weight` into the execute bucket, then runs `fn` inside a concurrency slot (high priority). */
execute<T>(weight: number, fn: () => Promise<T>, opts?: ThrottleCallOptions): Promise<T>;
stats(): WeightThrottleStats;
/** Returns and resets the per-label call counters (log them periodically to see who spends the budget). */
drainCounts(): Array<{ label: string; n: number }>;
}
const TIERS: readonly Priority[] = ['high', 'normal'];
interface Waiter {
weight: number;
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 ${String(value)}`);
}
}
class TokenBucket {
private tokens: number;
private last: number;
private timer: ReturnType<typeof setTimeout> | null = null;
private totalQueued = 0;
private totalWeight = 0;
private readonly queues: Record<Priority, Waiter[]> = { high: [], normal: [] };
private readonly ratePerMs: number;
constructor(
readonly perMinute: number,
readonly capacity: number,
private readonly now: () => number,
private readonly family: string,
) {
this.tokens = capacity;
this.last = now();
this.ratePerMs = perMinute / 60_000;
}
private refill(): void {
const t = this.now();
const elapsed = t - this.last;
this.last = t;
if (elapsed > 0) this.tokens = Math.min(this.capacity, this.tokens + elapsed * this.ratePerMs);
}
private head(): Waiter[] | null {
for (const tier of TIERS) if (this.queues[tier].length > 0) return this.queues[tier];
return null;
}
private drain(): void {
this.refill();
let q = this.head();
while (q) {
const w = q[0];
if (!w || this.tokens < w.weight) break;
q.shift();
this.tokens -= w.weight;
this.totalWeight += w.weight;
w.cleanup();
w.resolve();
q = this.head();
}
const next = q?.[0];
if (next && this.timer === null) {
const waitMs = Math.max(5, Math.ceil((next.weight - this.tokens) / this.ratePerMs));
this.timer = setTimeout(() => {
this.timer = null;
this.drain();
}, waitMs);
} else if (!next && this.timer !== null) {
clearTimeout(this.timer);
this.timer = null;
}
}
acquire(weight: number, priority: Priority, signal?: AbortSignal): Promise<void> {
if (!Number.isFinite(weight) || weight < 0)
return Promise.reject(
new RangeError(`${this.family} weight must be a non-negative number, got ${String(weight)}`),
);
if (weight > this.capacity) {
return Promise.reject(
new RangeError(
`${this.family} weight ${weight} exceeds the bucket capacity ${this.capacity} (chunk the request)`,
),
);
}
if (signal?.aborted) return Promise.reject(abortReason(signal));
this.refill();
if (this.head() === null && this.tokens >= weight) {
this.tokens -= weight;
this.totalWeight += weight;
return Promise.resolve();
}
this.totalQueued++;
return new Promise<void>((resolve, reject) => {
const waiter: Waiter = { weight, resolve, reject, cleanup: () => {} };
if (signal) {
const onAbort = () => {
const q = this.queues[priority];
const i = q.indexOf(waiter);
if (i >= 0) q.splice(i, 1);
reject(abortReason(signal));
this.drain();
};
signal.addEventListener('abort', onAbort, { once: true });
waiter.cleanup = () => signal.removeEventListener('abort', onAbort);
}
this.queues[priority].push(waiter);
this.drain();
});
}
stats(): BucketStats {
this.refill();
return {
tokens: Math.floor(this.tokens),
capacity: this.capacity,
perMinute: this.perMinute,
queued: this.queues.high.length + this.queues.normal.length,
totalQueued: this.totalQueued,
totalWeight: this.totalWeight,
};
}
}
interface SlotWaiter {
resolve: () => void;
reject: (err: unknown) => void;
cleanup: () => void;
}
class Semaphore {
private active = 0;
private readonly queues: Record<Priority, SlotWaiter[]> = { high: [], normal: [] };
constructor(private readonly max: number) {}
get inFlight(): number {
return this.active;
}
get queued(): number {
return this.queues.high.length + this.queues.normal.length;
}
private release(): void {
// Hand the slot straight to the next waiter by priority (no `active--` in between, so a fresh caller
// cannot slip past a queued execute).
for (const tier of TIERS) {
const next = this.queues[tier].shift();
if (next) {
next.cleanup();
next.resolve();
return;
}
}
this.active--;
}
async run<T>(fn: () => Promise<T>, priority: Priority, signal?: AbortSignal): Promise<T> {
if (signal?.aborted) throw abortReason(signal);
if (this.active < this.max) {
this.active++;
} else {
await new Promise<void>((resolve, reject) => {
const waiter: SlotWaiter = { resolve, reject, cleanup: () => {} };
if (signal) {
const onAbort = () => {
const q = this.queues[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);
}
this.queues[priority].push(waiter);
});
}
try {
return await fn();
} finally {
this.release();
}
}
}
/** Creates the two-bucket throttle described in the module header. */
export function createWeightThrottle(options: WeightThrottleOptions): WeightThrottle {
if (
!options ||
options.queriesPerMinute === undefined ||
options.queryBurst === undefined ||
options.executesPerMinute === undefined ||
options.executeBurst === undefined ||
options.maxConcurrent === undefined
) {
throw new TypeError(
'createWeightThrottle requires explicit queriesPerMinute, queryBurst, executesPerMinute, executeBurst and maxConcurrent',
);
}
const qpm = options.queriesPerMinute;
const qBurst = options.queryBurst;
const epm = options.executesPerMinute;
const eBurst = options.executeBurst;
const maxConcurrent = options.maxConcurrent;
assertPositive('queriesPerMinute', qpm, false);
assertPositive('queryBurst', qBurst, false);
assertPositive('executesPerMinute', epm, false);
assertPositive('executeBurst', eBurst, false);
assertPositive('maxConcurrent', maxConcurrent, true);
const now = options.now ?? (() => performance.now());
const queries = new TokenBucket(qpm, qBurst, now, 'query');
const executes = new TokenBucket(epm, eBurst, now, 'execute');
const semaphore = new Semaphore(maxConcurrent);
const counts = new Map<string, number>();
const bump = (label: string): void => {
counts.set(label, (counts.get(label) ?? 0) + 1);
};
async function run<T>(
bucket: TokenBucket,
weight: number,
fn: () => Promise<T>,
opts: ThrottleCallOptions,
defaultPriority: Priority,
defaultLabel: string,
): Promise<T> {
const priority = opts.priority ?? defaultPriority;
await bucket.acquire(weight, priority, opts.signal);
bump(opts.label ?? defaultLabel);
return semaphore.run(fn, priority, opts.signal);
}
return {
query: (weight, fn, opts = {}) => run(queries, weight, fn, opts, 'normal', 'query'),
execute: (weight, fn, opts = {}) => run(executes, weight, fn, opts, 'high', 'execute'),
stats: () => ({
queries: queries.stats(),
executes: executes.stats(),
inFlight: semaphore.inFlight,
slotQueued: semaphore.queued,
}),
drainCounts: () => {
const out = [...counts.entries()].map(([label, n]) => ({ label, n })).sort((a, b) => b.n - a.n);
counts.clear();
return out;
},
};
}