Skip to content
markpaper

src/transport/limiter.ts

v0.3.0 · 13.6 KB

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