Skip to content
markpaper

src/ws/subscriptionKey.ts

v0.3.0 · 9.8 KB

Download file
// Canonical subscription identity, address validation, ack matching and frame routing.

import type { WsSubscription } from './types.js';

const ADDRESS_RE = /^0x[0-9a-f]{40}$/;
const HEX_RE = /^0x[0-9a-f]+$/i;

/** Thrown synchronously by `subscribe()` for a malformed subscription (nothing is sent). */
export class WsSubscriptionError extends Error {
  override readonly name = 'WsSubscriptionError';
}

/**
 * Normalizes a user address to lowercase and validates `^0x[0-9a-f]{40}$`.
 *
 * Why: the server may echo addresses in a different case (ack keys in the original case leave a
 * client "unhealthy" forever), and a placeholder like `0x...` copied from an example env gets a
 * `channel:error` and an endless reconnect loop. Validate ALWAYS, even in dry-run.
 *
 * @throws {WsSubscriptionError} when the value is not a 20-byte hex address.
 */
export function normalizeAddress(user: string): string {
  const normalized = String(user).trim().toLowerCase();
  if (!ADDRESS_RE.test(normalized)) {
    throw new WsSubscriptionError(`invalid user address in subscription: expected 0x + 40 hex chars`);
  }
  return normalized;
}

/** Returns true when `value` is a syntactically valid 20-byte hex address (any case). */
export function isValidAddress(value: unknown): boolean {
  return typeof value === 'string' && ADDRESS_RE.test(value.trim().toLowerCase());
}

type Json = null | boolean | number | string | Json[] | { [key: string]: Json };

function normalizeValue(value: unknown): Json | undefined {
  if (value === undefined || value === null) return undefined;
  if (typeof value === 'string') return HEX_RE.test(value) ? value.toLowerCase() : value;
  if (typeof value === 'number' || typeof value === 'boolean') return value;
  if (Array.isArray(value)) return value.map((v) => normalizeValue(v) ?? null);
  if (typeof value === 'object') {
    const out: { [key: string]: Json } = {};
    for (const key of Object.keys(value as object).sort()) {
      const v = normalizeValue((value as Record<string, unknown>)[key]);
      if (v !== undefined) out[key] = v;
    }
    return out;
  }
  return undefined;
}

/**
 * `dex: ""` means the main dex: the server treats `{type:"allMids",dex:""}` as `{type:"allMids"}`
 * (live, 2026-09-15: the second one got `Already subscribed: {"type":"allMids"}`) and main-dex frames
 * carry no `dex` field. Without dropping it, two registry entries share one server subscription and
 * the second never gets its own ack.
 */
function dropMainDex(value: Json | undefined): Json | undefined {
  if (value && typeof value === 'object' && !Array.isArray(value) && value.dex === '') {
    const { dex: _main, ...rest } = value;
    return rest;
  }
  return value;
}

/**
 * Returns the canonical form of a subscription: keys sorted, `null`/`undefined` fields dropped,
 * `dex: ""` (main dex) dropped, `user` validated and lowercased. This is also the object that goes on
 * the wire.
 *
 * Coins are NOT validated here (that needs `meta`), and they are case-sensitive: live on 2026-09-15 a
 * `l2Book`/`trades`/`bbo` subscription to an unknown coin or a wrong-case one (`btc`, `xyz:tsla`) made
 * the server drop the whole socket at once (close 1006, no error frame, later subscriptions in the
 * same burst lost). The client quarantines such a subscription after repeated drops, but pass exact
 * names: `BTC`, `xyz:TSLA` (HIP-3), `@107` / `PURR/USDC` (spot).
 *
 * @throws {WsSubscriptionError} on a missing `type` or an invalid `user`.
 */
export function canonicalizeSubscription<S extends WsSubscription>(subscription: S): S {
  if (!subscription || typeof subscription !== 'object' || typeof subscription.type !== 'string' || !subscription.type) {
    throw new WsSubscriptionError('subscription must be an object with a non-empty string `type`');
  }
  const raw = subscription as Record<string, unknown>;
  if ('user' in raw && raw.user !== undefined) {
    if (typeof raw.user !== 'string') throw new WsSubscriptionError('subscription `user` must be a string');
    normalizeAddress(raw.user);
  }
  return dropMainDex(normalizeValue(subscription)) as unknown as S;
}

/**
 * Stable identity of a subscription (canonical JSON). Two modules subscribing to the same payload
 * share one server subscription; the client refcounts handlers so one module's unsubscribe does not
 * remove the other's data (dedupe-by-JSON without a refcount has exactly that bug).
 */
export function subscriptionKey(subscription: WsSubscription): string {
  return JSON.stringify(canonicalizeSubscription(subscription));
}

/** Lowercased `user` of a user-specific subscription, `undefined` for market subscriptions. */
export function trackedUserOf(subscription: WsSubscription): string | undefined {
  const user = (subscription as Record<string, unknown>).user;
  return typeof user === 'string' ? user.trim().toLowerCase() : undefined;
}

/** Channels on which frames of this subscription arrive. */
export function channelsOf(subscription: WsSubscription): string[] {
  switch (subscription.type) {
    case 'userEvents':
      return ['user'];
    case 'activeAssetCtx':
      return ['activeAssetCtx', 'activeSpotAssetCtx'];
    default:
      return [subscription.type];
  }
}

function str(value: unknown): string | undefined {
  return typeof value === 'string' ? value : undefined;
}

function sameUser(a: unknown, b: unknown): boolean {
  return typeof a === 'string' && typeof b === 'string' && a.toLowerCase() === b.toLowerCase();
}

/**
 * Decides whether a data frame belongs to a subscription (several subscriptions of the same
 * channel can share one socket). Frames without an identifying field (`orderUpdates`, `userEvents`,
 * `allDexsAssetCtxs`) match every subscription of their channel.
 */
export function frameMatchesSubscription(subscription: WsSubscription, channel: string, data: unknown): boolean {
  if (!channelsOf(subscription).includes(channel)) return false;
  const sub = subscription as Record<string, unknown>;
  const d = (data && typeof data === 'object' ? data : {}) as Record<string, unknown>;
  switch (subscription.type) {
    case 'trades': {
      if (!Array.isArray(data)) return false;
      const first = data[0] as { coin?: unknown } | undefined;
      return first !== undefined && first.coin === sub.coin;
    }
    case 'l2Book':
    case 'bbo':
    case 'activeAssetCtx':
      return d.coin === sub.coin;
    case 'candle':
      return d.s === sub.coin && d.i === sub.interval;
    case 'allMids':
      return (str(d.dex) ?? '') === (str(sub.dex) ?? '');
    case 'orderUpdates':
    case 'userEvents':
    case 'allDexsAssetCtxs':
      return true;
    case 'webData3': {
      const state = d.userState as Record<string, unknown> | undefined;
      return state === undefined || sameUser(state.user, sub.user);
    }
    default: {
      if (typeof sub.user === 'string' && typeof d.user === 'string' && !sameUser(d.user, sub.user)) return false;
      if (typeof sub.coin === 'string' && typeof d.coin === 'string' && d.coin !== sub.coin) return false;
      return true;
    }
  }
}

function isSubset(subset: unknown, superset: unknown): boolean {
  if (typeof subset === 'string' && typeof superset === 'string') {
    return HEX_RE.test(subset) && HEX_RE.test(superset) ? subset.toLowerCase() === superset.toLowerCase() : subset === superset;
  }
  if (typeof subset !== 'object' || typeof superset !== 'object' || subset === null || superset === null) {
    return subset === superset;
  }
  if (Array.isArray(subset)) {
    return Array.isArray(superset) && subset.length === superset.length && subset.every((v, i) => isSubset(v, superset[i]));
  }
  const sup = superset as Record<string, unknown>;
  return Object.entries(subset as Record<string, unknown>).every(([k, v]) => k in sup && isSubset(v, sup[k]));
}

function specificity(value: unknown): number {
  if (typeof value !== 'object' || value === null) return 1;
  return Object.values(value).reduce<number>((sum, v) => sum + specificity(v), 0);
}

const LOOSE_FIELDS = ['type', 'coin', 'user', 'interval', 'dex'] as const;

/**
 * Loose ack identity `{type, coin, user(lowercase), interval, dex}`. The server normalizes its echo
 * (fields may be added, unknown ones dropped, address case changed), so exact JSON equality is not
 * enough on its own.
 */
export function looseAckKey(subscription: unknown): string {
  const s = (subscription && typeof subscription === 'object' ? subscription : {}) as Record<string, unknown>;
  const out: Record<string, unknown> = {};
  for (const f of LOOSE_FIELDS) {
    const v = s[f];
    if (v === undefined || v === null || v === '') continue;
    out[f] = f === 'user' && typeof v === 'string' ? v.toLowerCase() : v;
  }
  return JSON.stringify(out);
}

/**
 * Finds which pending subscription an echoed subscription (from `subscriptionResponse` or an error
 * text) refers to: exact canonical key, then the most specific pending payload contained in the
 * echo, then the loose `{type, coin, user, interval, dex}` identity.
 *
 * Why exact matching matters: in a multiplexed socket a frame for one user must never confirm the
 * subscription of another, and an unrelated frame must never make a partial set look healthy.
 */
export function matchEcho(
  pending: ReadonlyArray<{ key: string; subscription: WsSubscription }>,
  echo: unknown,
): string | undefined {
  if (!echo || typeof echo !== 'object') return undefined;
  const normalized = dropMainDex(normalizeValue(echo));
  const exact = JSON.stringify(normalized);
  const hit = pending.find((p) => p.key === exact);
  if (hit) return hit.key;

  let best: string | undefined;
  let bestScore = -1;
  for (const p of pending) {
    if (!isSubset(p.subscription, normalized)) continue;
    const score = specificity(p.subscription);
    if (score > bestScore) { best = p.key; bestScore = score; }
  }
  if (best !== undefined) return best;

  const loose = looseAckKey(echo);
  return pending.find((p) => looseAckKey(p.subscription) === loose)?.key;
}
All files