src/throttle/index.ts
v0.1.0 · 4.9 KB
import { parseRetryAfterMs } from '../rest/client.js';
export type QfexBucket = 'general' | 'cancel';
// Published per-user buckets, weights and budget interval: docs.qfex.com/websocket/rate.
export const QFEX_WEIGHTS = Object.freeze({
addOrder: 1,
cancelOrder: 1,
getUserOrders: 5,
getOrder: 2,
leverage: 0.1,
});
export const QFEX_GENERAL_BUDGET = 12_000;
export const QFEX_BUDGET_INTERVAL_MS = 60_000;
export function createRateBudget(options: {
budgetGeneral: number;
budgetCancel: number;
fuseFraction: number;
defaultBackoffMs: number;
rateBackoff: { minMs: number; maxMs: number };
emaTimeConstantMs: number;
now?: () => number;
}) {
for (const name of [
'budgetGeneral',
'budgetCancel',
'fuseFraction',
'defaultBackoffMs',
'emaTimeConstantMs',
] as const)
if (!(options?.[name] > 0) || !Number.isFinite(options[name])) throw new TypeError('Required rate policy: ' + name);
if (options.fuseFraction > 1) throw new TypeError('Fuse fraction must not exceed one');
parseRetryAfterMs(null, null, 0, options.defaultBackoffMs, options.rateBackoff);
const now = options.now ?? Date.now;
const state = { general: { value: 0, at: now(), until: 0 }, cancel: { value: 0, at: now(), until: 0 } };
const value = (bucket: QfexBucket, at: number) =>
state[bucket].value * Math.exp(-Math.max(0, at - state[bucket].at) / options.emaTimeConstantMs);
function usage(at = now()) {
return {
general: value('general', at) / options.budgetGeneral,
cancel: value('cancel', at) / options.budgetCancel,
};
}
return {
usage,
noteWeight(bucket: QfexBucket, weight: number, at = now()) {
if (!(weight > 0) || !Number.isFinite(weight)) return;
state[bucket].value = value(bucket, at) + weight;
state[bucket].at = Math.max(at, state[bucket].at);
},
noteRateLimited(bucket: QfexBucket | 'all', message: string | null, at = now()) {
const ms = parseRetryAfterMs(null, message, at, options.defaultBackoffMs, options.rateBackoff);
for (const key of bucket === 'all' ? (['general', 'cancel'] as const) : [bucket])
state[key].until = Math.max(state[key].until, at + ms);
return ms;
},
hold(bucket: QfexBucket, at = now()): string | null {
if (at < state[bucket].until) return 'Rate limit hold for ' + (state[bucket].until - at) + ' ms';
if (bucket === 'general' && usage(at).general > options.fuseFraction) return 'Local EMA budget fuse';
return null;
},
};
}
export type RateBudget = ReturnType<typeof createRateBudget>;
export interface QfexAddSlot {
cloid: string;
since: number;
release(): void;
}
export function createAddSlot(options: {
staleMs: number;
now?: () => number;
setTimer?: (ms: number, callback: () => void) => { cancel(): void };
}) {
if (!(options?.staleMs > 0) || !Number.isFinite(options.staleMs)) throw new TypeError('Required stale slot policy');
const now = options.now ?? Date.now;
const setTimer =
options.setTimer ??
((ms: number, fn: () => void) => {
const timer = setTimeout(fn, ms);
return { cancel: () => clearTimeout(timer) };
});
let holder: { cloid: string; since: number; token: number } | null = null;
let token = 0;
const queue: Array<{ cloid: string; resolve: (slot: QfexAddSlot | null) => void; timer: { cancel(): void } | null }> =
[];
function grant(cloid: string): QfexAddSlot {
const selected = { cloid, since: now(), token: ++token };
holder = selected;
return {
cloid,
since: selected.since,
release() {
if (holder?.token !== selected.token) return;
holder = null;
handover();
},
};
}
function handover() {
const next = queue.shift();
if (next) {
next.timer?.cancel();
next.resolve(grant(next.cloid));
}
}
function sweep() {
if (holder && now() - holder.since >= options.staleMs) {
holder = null;
handover();
}
}
function tryAcquire(cloid: string): QfexAddSlot | null {
sweep();
return holder || queue.length ? null : grant(cloid);
}
function acquire(cloid: string, maxWaitMs: number): Promise<QfexAddSlot | null> {
const slot = tryAcquire(cloid);
if (slot || maxWaitMs <= 0) return Promise.resolve(slot);
return new Promise((resolve) => {
const item = { cloid, resolve, timer: null as { cancel(): void } | null };
item.timer = setTimer(maxWaitMs, () => {
const index = queue.indexOf(item);
if (index < 0) return;
queue.splice(index, 1);
sweep();
resolve(holder || queue.length ? null : grant(cloid));
});
queue.push(item);
});
}
return {
acquire,
tryAcquire,
inFlight: () => (holder ? { cloid: holder.cloid, sinceMs: Math.max(0, now() - holder.since) } : null),
close() {
for (const item of queue.splice(0)) {
item.timer?.cancel();
item.resolve(null);
}
holder = null;
},
};
}
export type AddSlot = ReturnType<typeof createAddSlot>;