Skip to content
markpaper

src/ws/limits.ts

v0.3.0 · 13.1 KB

Download file
// WebSocket limits of Hyperliquid and a per-egress-IP budget shared by clients.

/**
 * Numbers from practice and docs. All IP-scoped limits are counted per EGRESS IP across every
 * connection and process on it; opening more connections from the same IP adds no capacity, only
 * another egress IP does.
 */
export const WS_LIMITS = {
  /** The server closes a socket after 60 s without messages. */
  serverIdleTimeoutMs: 60_000,
  /** Client ping period: well below the idle timeout. */
  pingIntervalMs: 30_000,
  /** No inbound frame (pong included) for ping x 2 + 5 s => the socket is dead. */
  silenceTimeoutMs: 65_000,
  /** Hung handshake guard (the `handshakeTimeout` used with the `ws` package). */
  connectTimeoutMs: 15_000,
  /** Subscriptions per IP (docs, medium confidence; not verified live). */
  maxSubscriptionsPerIp: 1000,
  /**
   * Connections per IP. Later records say 10 (close 1008 above it), an early one says 100; plan for 10
   * until checked against the official rate-limits page.
   */
  maxConnectionsPerIp: 10,
  /** Unique users across user-specific subscriptions, DOCUMENTED value: the only guaranteed capacity. */
  uniqueUsersDocumented: 10,
  /**
   * Unique users MEASURED on 2026-07-17: ~20 per IP in total across all connections (fresh IP: 30 requested
   * -> 21 tracked, 9 refused; loaded IP: 27 -> 18 alive). Best-effort bonus HL may remove at any time.
   */
  uniqueUsersMeasured: 20,
  /** The number in the error text `Cannot track more than 15 total users`; it is NOT a per-connection cap. */
  uniqueUsersErrorText: 15,
  /** A slot stays held ~60 s after a reconnect, or after unsubscribing on one connection and subscribing on another. */
  ghostSlotMs: 60_000,
  /**
   * New connections per minute per IP.
   * @experimental Not recorded in the knowledge base; value from the public HL rate-limits page. A
   * client-side guard only delays connects, so a wrong number cannot cause a refusal by the server.
   */
  maxNewConnectionsPerMinute: 30,
  /**
   * Outgoing WS messages per minute per IP. The knowledge base confirms the limit exists but not the number.
   * @experimental Value from the public HL rate-limits page. Exceeding frames are delayed, not dropped.
   */
  maxMessagesPerMinute: 2000,
} as const;

export interface WsIpBudgetLimits {
  /** Max sockets open at once on this IP. Default 10. */
  maxConnections?: number;
  /** Max subscriptions on this IP. Default 1000. */
  maxSubscriptions?: number;
  /**
   * Hard cap of tracked user slots (live users + ghost slots). A new user above it is refused BEFORE
   * sending. Default 10 (documented). Raise up to ~20 (measured) knowingly: the surplus is not guaranteed.
   */
  maxUniqueUsers?: number;
  /** New users above this emit `userLimitWarning`. Default 10 (documented guarantee). */
  warnUniqueUsers?: number;
  /** How long a ghost slot is held. Default 60 000 ms. */
  ghostSlotMs?: number;
  /** @experimental Default 30. See {@link WS_LIMITS.maxNewConnectionsPerMinute}. */
  maxNewConnectionsPerMinute?: number;
  /** @experimental Default 2000. See {@link WS_LIMITS.maxMessagesPerMinute}. */
  maxMessagesPerMinute?: number;
}

export type WsBudgetWarning =
  | { kind: 'userLimitWarning'; used: number; limit: number }
  | { kind: 'migrationWarning'; user: string };

export type WsSubscribeCheck =
  | { ok: true; warnings: WsBudgetWarning[] }
  | { ok: false; reason: 'maxSubscriptions' | 'maxUniqueUsers'; used: number; limit: number };

export type WsConnectionCheck =
  | { ok: true }
  | { ok: false; reason: 'maxConnections' | 'newConnectionsPerMinute'; retryInMs: number };

export interface WsIpBudgetStats {
  connections: number;
  subscriptions: number;
  /** Distinct users with a live subscription. */
  uniqueUsers: number;
  /** Ghost slots currently counted against the user limit. */
  ghostSlots: number;
  /** `uniqueUsers + ghostSlots`: what the user limit is compared with. */
  usedUserSlots: number;
  limits: Required<WsIpBudgetLimits>;
}

/**
 * Accounting of one egress IP. Share ONE budget between every client that leaves through the same
 * IP (the server counts per IP, not per connection); give clients on different `localAddress`
 * egress IPs different budgets.
 */
export interface WsIpBudget {
  readonly limits: Required<WsIpBudgetLimits>;
  /** Asks to open a socket for `clientId`; records the attempt when allowed. */
  tryOpenConnection(clientId: string, now: number): WsConnectionCheck;
  /**
   * Marks the client's socket closed; its tracked users become ghost slots for `ghostSlotMs`.
   * Pass `heldSlots: false` when the socket never opened (handshake failure, connect timeout): no
   * subscription reached the server, so there is nothing to hold. Without it every failed connect
   * during an outage stacks another set of ghosts and the budget refuses new users for no reason.
   */
  connectionClosed(clientId: string, now: number, heldSlots?: boolean): void;
  /** Checks a NEW subscription before it is sent. */
  checkSubscribe(clientId: string, user: string | undefined, now: number): WsSubscribeCheck;
  addSubscription(clientId: string, key: string, user: string | undefined, now: number): void;
  removeSubscription(clientId: string, key: string, user: string | undefined, now: number): void;
  /** Drops every subscription of a client (on `close()`); users become ghost slots if it was connected. */
  releaseClient(clientId: string, now: number): void;
  /** Reserves one outgoing message; returns 0 when sent now is fine, otherwise ms to wait. */
  reserveMessage(now: number): number;
  stats(now: number): WsIpBudgetStats;
}

interface Ghost { user: string; clientId: string; until: number; kind: 'closed' | 'released' }

interface ClientState { connected: boolean; subs: Map<string, string | undefined> }

class SlidingWindow {
  private readonly times: number[] = [];
  constructor(private readonly limit: number, private readonly windowMs: number) {}
  private prune(now: number): void {
    while (this.times.length > 0 && (this.times[0] as number) <= now - this.windowMs) this.times.shift();
  }
  /** 0 if an event may happen now (and records it), otherwise ms until it may. */
  take(now: number): number {
    this.prune(now);
    if (this.times.length < this.limit) { this.times.push(now); return 0; }
    return Math.max(1, (this.times[0] as number) + this.windowMs - now);
  }
}

/**
 * Creates a per-IP budget.
 *
 * User slots are modeled on the observed HL mechanics:
 * - a user subscribed anywhere on the IP takes one slot, however many connections subscribe to it;
 * - `unsubscribe` frees the slot immediately on the SAME connection;
 * - after a reconnect the old connection keeps its users ~60 s (ghost slots), so right after a
 *   reconnect only part of N subscriptions comes back;
 * - unsubscribing on one connection and subscribing on another holds the old slot ~60 s: migrating
 *   subscriptions between shards can start a self-sustaining refusal storm even with demand under
 *   the cap.
 */
export function createWsIpBudget(limits: WsIpBudgetLimits = {}): WsIpBudget {
  const resolved: Required<WsIpBudgetLimits> = {
    maxConnections: limits.maxConnections ?? WS_LIMITS.maxConnectionsPerIp,
    maxSubscriptions: limits.maxSubscriptions ?? WS_LIMITS.maxSubscriptionsPerIp,
    maxUniqueUsers: limits.maxUniqueUsers ?? WS_LIMITS.uniqueUsersDocumented,
    warnUniqueUsers: limits.warnUniqueUsers ?? WS_LIMITS.uniqueUsersDocumented,
    ghostSlotMs: limits.ghostSlotMs ?? WS_LIMITS.ghostSlotMs,
    maxNewConnectionsPerMinute: limits.maxNewConnectionsPerMinute ?? WS_LIMITS.maxNewConnectionsPerMinute,
    maxMessagesPerMinute: limits.maxMessagesPerMinute ?? WS_LIMITS.maxMessagesPerMinute,
  };
  for (const [name, value] of Object.entries(resolved)) {
    if (!(value >= 0)) throw new RangeError(`ws budget: ${name} must be a non-negative number`);
  }

  const clients = new Map<string, ClientState>();
  let ghosts: Ghost[] = [];
  const newConnections = new SlidingWindow(resolved.maxNewConnectionsPerMinute, 60_000);
  const messages = new SlidingWindow(resolved.maxMessagesPerMinute, 60_000);

  const client = (id: string): ClientState => {
    let state = clients.get(id);
    if (!state) { state = { connected: false, subs: new Map() }; clients.set(id, state); }
    return state;
  };

  const pruneGhosts = (now: number) => { ghosts = ghosts.filter((g) => g.until > now); };

  const clientHasUser = (state: ClientState, user: string) => {
    for (const u of state.subs.values()) if (u === user) return true;
    return false;
  };

  const liveOnOtherClient = (user: string, clientId: string, extra?: { clientId: string; user: string }) => {
    if (extra && extra.user === user && extra.clientId !== clientId) return true;
    for (const [id, state] of clients) if (id !== clientId && clientHasUser(state, user)) return true;
    return false;
  };

  const liveUsers = (extra?: { clientId: string; user: string }) => {
    const set = new Set<string>();
    for (const state of clients.values()) for (const u of state.subs.values()) if (u !== undefined) set.add(u);
    if (extra) set.add(extra.user);
    return set;
  };

  const countGhosts = (list: Ghost[], extra?: { clientId: string; user: string }) =>
    list.filter((g) => g.kind === 'closed' || liveOnOtherClient(g.user, g.clientId, extra)).length;

  const totalSubscriptions = () => {
    let n = 0;
    for (const state of clients.values()) n += state.subs.size;
    return n;
  };

  function connectionClosed(clientId: string, now: number, heldSlots = true): void {
    const state = clients.get(clientId);
    if (!state || !state.connected) return;
    state.connected = false;
    if (!heldSlots) return;
    pruneGhosts(now);
    const users = new Set<string>();
    for (const u of state.subs.values()) if (u !== undefined) users.add(u);
    for (const user of users) ghosts.push({ user, clientId, until: now + resolved.ghostSlotMs, kind: 'closed' });
  }

  return {
    limits: resolved,

    tryOpenConnection(clientId, now) {
      const state = client(clientId);
      let open = 0;
      for (const [id, s] of clients) if (id !== clientId && s.connected) open++;
      if (open >= resolved.maxConnections) {
        return { ok: false, reason: 'maxConnections', retryInMs: resolved.ghostSlotMs };
      }
      const wait = newConnections.take(now);
      if (wait > 0) return { ok: false, reason: 'newConnectionsPerMinute', retryInMs: wait };
      state.connected = true;
      return { ok: true };
    },

    connectionClosed,

    checkSubscribe(clientId, user, now) {
      pruneGhosts(now);
      const subscriptions = totalSubscriptions();
      if (subscriptions + 1 > resolved.maxSubscriptions) {
        return { ok: false, reason: 'maxSubscriptions', used: subscriptions, limit: resolved.maxSubscriptions };
      }
      const warnings: WsBudgetWarning[] = [];
      if (user === undefined) return { ok: true, warnings };
      const state = client(clientId);
      if (clientHasUser(state, user) || liveOnOtherClient(user, clientId)) return { ok: true, warnings };

      const extra = { clientId, user };
      // Re-adding on the same connection frees its own released slot immediately.
      const projectedGhosts = ghosts.filter(
        (g) => !(g.kind === 'released' && g.clientId === clientId && g.user === user),
      );
      const used = liveUsers(extra).size + countGhosts(projectedGhosts, extra);

      if (ghosts.some((g) => g.kind === 'released' && g.user === user && g.clientId !== clientId)) {
        warnings.push({ kind: 'migrationWarning', user });
      }
      if (used > resolved.maxUniqueUsers) {
        return { ok: false, reason: 'maxUniqueUsers', used: used - 1, limit: resolved.maxUniqueUsers };
      }
      if (used > resolved.warnUniqueUsers) {
        warnings.push({ kind: 'userLimitWarning', used, limit: resolved.warnUniqueUsers });
      }
      return { ok: true, warnings };
    },

    addSubscription(clientId, key, user, now) {
      pruneGhosts(now);
      const state = client(clientId);
      state.subs.set(key, user);
      if (user !== undefined) {
        ghosts = ghosts.filter((g) => !(g.kind === 'released' && g.clientId === clientId && g.user === user));
      }
    },

    removeSubscription(clientId, key, user, now) {
      pruneGhosts(now);
      const state = clients.get(clientId);
      if (!state || !state.subs.delete(key)) return;
      if (user !== undefined && state.connected && !clientHasUser(state, user)) {
        ghosts.push({ user, clientId, until: now + resolved.ghostSlotMs, kind: 'released' });
      }
    },

    releaseClient(clientId, now) {
      const state = clients.get(clientId);
      if (!state) return;
      connectionClosed(clientId, now);
      state.subs.clear();
      clients.delete(clientId);
      // Ghosts of the gone client stay: the server still holds those slots for ~60 s.
    },

    reserveMessage(now) {
      return messages.take(now);
    },

    stats(now) {
      pruneGhosts(now);
      let connections = 0;
      for (const s of clients.values()) if (s.connected) connections++;
      const uniqueUsers = liveUsers().size;
      const ghostSlots = countGhosts(ghosts);
      return {
        connections,
        subscriptions: totalSubscriptions(),
        uniqueUsers,
        ghostSlots,
        usedUserSlots: uniqueUsers + ghostSlots,
        limits: resolved,
      };
    },
  };
}
All files