Skip to content
markpaper

src/throttle/index.ts

v0.1.0 · 4.9 KB

Download file
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>;
All files