src/ws/subscriptionKey.ts
v0.3.0 · 9.8 KB
// 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;
}