Skip to content
markpaper

src/transport/throttle.ts

v0.2.0 · 10.4 KB

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