Skip to content
markpaper

src/transport/info-client.ts

v0.3.0 · 13.9 KB

Download file
// Throttled, retrying POST /info client implementing the shared InfoRequester contract.

import { API_URL, type InfoCallOptions, type InfoRequest, type InfoRequester, type Network } from './types.js';
import { HlHttpError, HlNetworkError, HlResponseParseError, HlTimeoutError, isTransientError } from './errors.js';
import { type Priority, type WeightLimiter } from './limiter.js';
import { DEFAULT_RETRY, withRetry, type RetryOptions } from './retry.js';
import { responseWeightSurcharge, weightOf } from './weights.js';

/** Default per-attempt timeout for /info, ms (knowledge base: 8-10 s). */
export const DEFAULT_INFO_TIMEOUT_MS = 10_000;

/** Minimal fetch signature the client needs; the global `fetch` satisfies it. */
export type FetchLike = (input: string, init: RequestInit) => Promise<Response>;

/** One finished attempt, for telemetry (per-label counters, health logs). */
export interface InfoAttemptEvent {
  type: string;
  label: string | undefined;
  /** Base weight paid before the attempt. */
  weight: number;
  /** Surcharge charged after the response (0 when disabled or not applicable). */
  extraWeight: number;
  /** 0-based attempt number. */
  attempt: number;
  ok: boolean;
  /** HTTP status when a response was received. */
  status: number | undefined;
  /** Time from fetch start to completion, ms (queue wait excluded). */
  durationMs: number;
  error: unknown;
}

/** Options for {@link createInfoClient}. */
export interface InfoClientOptions {
  /** Network for the default URL. Default `'mainnet'`. */
  network?: Network;
  /**
   * API base URL (e.g. a read-only proxy); `/info` is appended unless already present. When both
   * `network` and `baseUrl` are given and `baseUrl` is the OFFICIAL host of the other network, the
   * constructor throws: testnet data must never reach a mainnet signer (or vice versa).
   */
  baseUrl?: string;
  /**
   * Custom fetch. This is how to bind an egress IP: wrap the global fetch and pass an undici
   * dispatcher created with the TOP-LEVEL option, e.g.
   *
   * ```ts
   * import { Agent } from 'undici'; // pin a 6.x version; 7.x changed the handler API
   * const agent = new Agent({ localAddress: process.env.EGRESS_IP });
   * // NOT new Agent({ connect: { localAddress } }) - undici silently ignores it and the request
   * // leaves from the default IP (every lane then shares one 429 budget).
   * const laneFetch: FetchLike = (url, init) =>
   *   fetch(url, { ...init, dispatcher: agent } as RequestInit & { dispatcher?: unknown });
   * const limiter = getSharedWeightLimiter('lane-b', {
   *   weightPerMinute: config.weightPerMinute,
   *   burstCapacity: config.burstCapacity,
   *   maxConcurrent: config.maxConcurrent,
   * });
   * const info = createInfoClient({ fetch: laneFetch, limiter });
   * ```
   *
   * Verify the outbound IP of every lane at startup with one echo request before giving it its own
   * limiter key. The kit does not depend on undici.
   */
  fetch?: FetchLike;
  /**
   * Explicit limiter to pay request weight into. Pass `null` only when the caller already throttles
   * every call itself - a raw client outside the shared bucket silently spends the IP budget.
   */
  limiter: WeightLimiter | null;
  /** Per-attempt timeout, ms. Default 10 000. The timer starts when fetch starts, not while queued. */
  timeoutMs?: number;
  /**
   * Retry policy for 429 / 408 / 5xx / network / timeout / unparsable 200 (/info is idempotent).
   * Default: 4 retries, 250 ms x 2.5^n, jitter +-40%. `false` disables retries.
   */
  retry?: RetryOptions | false;
  /** Default queue priority. Default `'normal'`. */
  priority?: Priority;
  /**
   * Lower-case top-level string fields that look like addresses (`0x` + 40 hex). Default true:
   * all proven clients send and compare addresses in lower case.
   */
  lowercaseAddresses?: boolean;
  /**
   * Charge the documented per-item surcharge to the limiter after responses of history-like types.
   * Default true (over-charging only slows the bucket down; under-charging risks 429s).
   *
   * @experimental The surcharge is documented by HL but was never measured live; see `responseWeightSurcharge`.
   */
  chargeResponseWeight?: boolean;
  /**
   * JSON parser for the response text. Default `JSON.parse`. Hyperliquid order ids are 64-bit on the
   * protocol; `JSON.parse` silently rounds integers above 2^53, and a cancel by a rounded id "succeeds"
   * while the order stays on the book. Pass a lossless parser if you handle such ids.
   */
  parseJson?: (text: string) => unknown;
  /** Called after every attempt. Exceptions thrown by the callback are swallowed. */
  onAttempt?: (event: InfoAttemptEvent) => void;
}

/** Per-call options: the shared contract plus priority and a telemetry label. */
export interface InfoClientCallOptions extends InfoCallOptions {
  priority?: Priority;
  /** Free-form label for telemetry, e.g. `mon:clearinghouseState` or `trade:meta`. */
  label?: string;
}

/** An {@link InfoRequester} with its configuration exposed. */
export interface InfoClient {
  <T = unknown>(body: InfoRequest, opts?: InfoClientCallOptions): Promise<T>;
  /** Resolved /info URL. */
  readonly url: string;
  /** Limiter used by the client, or null. */
  readonly limiter: WeightLimiter | null;
}

const ADDRESS_RE = /^0x[0-9a-fA-F]{40}$/;

/** Largest delay `AbortSignal.timeout` accepts. */
const MAX_TIMEOUT_MS = 4_294_967_295;

/**
 * Validates a timeout and rounds it up to whole milliseconds. `AbortSignal.timeout` throws a
 * RangeError for fractional, negative or non-finite delays - inside a request that would surface as
 * a non-retryable failure of every call, so it is rejected up front instead.
 */
function checkTimeout(ms: number): number {
  if (typeof ms !== 'number' || !Number.isFinite(ms) || ms <= 0) {
    throw new RangeError(`timeoutMs must be a positive finite number, got ${String(ms)}`);
  }
  return Math.min(MAX_TIMEOUT_MS, Math.ceil(ms));
}

function hostOf(url: string): string | undefined {
  try {
    return new URL(url).host.toLowerCase();
  } catch {
    return undefined;
  }
}

/**
 * Resolves the /info URL from `network` / `baseUrl` and refuses a mainnet/testnet mix.
 *
 * @throws Error when `network` and an official `baseUrl` point to different networks, or the URL is invalid.
 */
export function resolveInfoUrl(opts: { network?: Network; baseUrl?: string } = {}): string {
  const network = opts.network;
  if (opts.baseUrl === undefined) return `${API_URL[network ?? 'mainnet']}/info`;
  const trimmed = opts.baseUrl.trim().replace(/\/+$/, '');
  const host = hostOf(trimmed);
  if (!host || !/^https?:$/i.test(new URL(trimmed).protocol)) {
    throw new Error(`invalid Hyperliquid baseUrl: ${opts.baseUrl}`);
  }
  if (network) {
    const other: Network = network === 'mainnet' ? 'testnet' : 'mainnet';
    if (host === hostOf(API_URL[other])) {
      throw new Error(`Hyperliquid endpoints mix mainnet and testnet: network=${network}, baseUrl=${opts.baseUrl}`);
    }
  }
  return /\/info$/i.test(trimmed) ? trimmed : `${trimmed}/info`;
}

/**
 * Normalizes an /info body without mutating it:
 * - removes `dex` when it is `''`, `null` or `undefined` - the main perp dex is addressed by
 *   omitting the key entirely, never by an empty string;
 * - optionally lower-cases top-level address strings.
 */
export function normalizeInfoBody(body: InfoRequest, lowercaseAddresses = true): InfoRequest {
  const out: InfoRequest = { type: body.type };
  for (const [key, value] of Object.entries(body)) {
    if (key === 'type') continue;
    if (key === 'dex' && (value === '' || value === null || value === undefined)) continue;
    out[key] = lowercaseAddresses && typeof value === 'string' && ADDRESS_RE.test(value) ? value.toLowerCase() : value;
  }
  return out;
}

function anySignal(signals: AbortSignal[]): { signal: AbortSignal; dispose: () => void } {
  if (signals.length === 1) return { signal: signals[0] as AbortSignal, dispose: () => {} };
  const ctl = new AbortController();
  const listeners: Array<[AbortSignal, () => void]> = [];
  for (const s of signals) {
    if (s.aborted) {
      ctl.abort(s.reason);
      break;
    }
    const on = () => ctl.abort(s.reason);
    s.addEventListener('abort', on, { once: true });
    listeners.push([s, on]);
  }
  return {
    signal: ctl.signal,
    dispose: () => {
      for (const [s, on] of listeners) s.removeEventListener('abort', on);
    },
  };
}

/**
 * Creates a throttled POST /info client that satisfies {@link InfoRequester}.
 *
 * Each attempt: pays the base weight into the limiter (priority queue), takes a concurrency slot,
 * runs fetch under `AbortSignal.timeout`, then charges any response surcharge. Retries (idempotent
 * /info) go back through the limiter, so backoff respects the budget.
 *
 * Failure semantics: the client resolves only with a parsed body and otherwise throws
 * ({@link HlHttpError}, {@link HlTimeoutError}, {@link HlNetworkError}, {@link HlResponseParseError},
 * or the caller's abort reason). It never returns `null`/`[]` on failure: a failed read is not an
 * empty account - treating a 429 as "no positions" silently skips trades.
 */
export function createInfoClient(options: InfoClientOptions): InfoClient {
  if (!options || !Object.prototype.hasOwnProperty.call(options, 'limiter')) {
    throw new TypeError('createInfoClient requires an explicit limiter (or null for external throttling)');
  }
  const url = resolveInfoUrl(options);
  const fetchImpl: FetchLike = options.fetch ?? ((input, init) => fetch(input, init));
  const limiter = options.limiter;
  const defaultTimeout = checkTimeout(options.timeoutMs ?? DEFAULT_INFO_TIMEOUT_MS);
  const retry = options.retry === false ? false : { ...DEFAULT_RETRY, ...(options.retry ?? {}) };
  const lowercase = options.lowercaseAddresses ?? true;
  const chargeResponse = options.chargeResponseWeight ?? true;
  const parse = options.parseJson ?? JSON.parse;
  const report = (event: InfoAttemptEvent): void => {
    try {
      options.onAttempt?.(event);
    } catch {
      // Telemetry must never turn a successful read into a failure (or into a retry that pays weight again).
    }
  };

  async function attemptOnce(
    type: string,
    payload: string,
    weight: number,
    attempt: number,
    opts: InfoClientCallOptions,
  ): Promise<unknown> {
    const timeoutMs = opts.timeoutMs === undefined ? defaultTimeout : checkTimeout(opts.timeoutMs);
    const timeoutSignal = AbortSignal.timeout(timeoutMs);
    const combined = anySignal(opts.signal ? [opts.signal, timeoutSignal] : [timeoutSignal]);
    const started = performance.now();
    let status: number | undefined;
    let extraWeight = 0;

    const fail = (err: unknown): never => {
      report({
        type, label: opts.label, weight, extraWeight, attempt, ok: false, status,
        durationMs: performance.now() - started, error: err,
      });
      throw err;
    };
    const mapAbortOrNetwork = (err: unknown): unknown => {
      if (opts.signal?.aborted) return opts.signal.reason ?? err;
      if (timeoutSignal.aborted) return new HlTimeoutError({ timeoutMs, url, requestType: type, cause: err });
      return new HlNetworkError({ url, requestType: type, cause: err });
    };

    try {
      let res: Response;
      try {
        res = await fetchImpl(url, {
          method: 'POST',
          headers: { 'content-type': 'application/json' },
          body: payload,
          signal: combined.signal,
        });
      } catch (err) {
        return fail(mapAbortOrNetwork(err));
      }
      status = res.status;

      if (!res.ok) {
        let bodyText = '';
        try {
          // Reading the body also releases the undici connection; an unread body keeps the socket busy.
          bodyText = await res.text();
        } catch {
          await res.body?.cancel().catch(() => {});
        }
        return fail(new HlHttpError({ status: res.status, bodyText, url, requestType: type }));
      }

      let text: string;
      try {
        text = await res.text();
      } catch (err) {
        return fail(mapAbortOrNetwork(err));
      }
      let data: unknown;
      try {
        data = parse(text);
      } catch (err) {
        return fail(new HlResponseParseError({ status: res.status, bodyText: text, url, requestType: type, cause: err }));
      }

      if (chargeResponse && limiter) {
        extraWeight = responseWeightSurcharge(type, data);
        if (extraWeight > 0) limiter.charge(extraWeight);
      }
      report({
        type, label: opts.label, weight, extraWeight, attempt, ok: true, status,
        durationMs: performance.now() - started, error: undefined,
      });
      return data;
    } finally {
      combined.dispose();
    }
  }

  const client = async <T = unknown>(body: InfoRequest, opts: InfoClientCallOptions = {}): Promise<T> => {
    if (!body || typeof body.type !== 'string' || body.type.length === 0) {
      throw new TypeError('info request body must have a non-empty string "type"');
    }
    if (opts.timeoutMs !== undefined) checkTimeout(opts.timeoutMs); // fail before paying any weight
    const normalized = normalizeInfoBody(body, lowercase);
    const type = normalized.type;
    const weight = opts.weight ?? weightOf(normalized);
    const payload = JSON.stringify(normalized);
    const priority = opts.priority ?? options.priority ?? 'normal';

    const run = (attempt: number): Promise<unknown> => {
      const exec = () => attemptOnce(type, payload, weight, attempt, opts);
      return limiter ? limiter.schedule(weight, exec, { priority, signal: opts.signal }) : exec();
    };

    const data = retry
      ? await withRetry(run, {
          ...retry,
          signal: opts.signal,
          // Only transport-level failures are retried; a RangeError from the limiter
          // (weight > capacity) or a 4xx such as 422 "unknown dex" is final.
          shouldRetry: (err) => isTransientError(err),
        })
      : await run(0);
    return data as T;
  };

  return Object.defineProperties(client, {
    url: { value: url, enumerable: true },
    limiter: { value: limiter, enumerable: true },
  }) as InfoClient;
}
All files