Skip to content
markpaper

src/rest/client.ts

v0.2.1 · 9.4 KB

Download file
// REST read client (knowledge base: instances-and-api.md §2, §2.5; rate-limits.md §5).
//
// Default transport policy: 15 s timeout, up to 3 attempts, pause 400 ms x attempt.
// 4xx is the exchange's answer and is NOT retried - 429 included: a read burst after a restart is
// normal for about a minute and is survived by the meta cache and the next tick, not by hammering.
// 401/403 on an authorized read refreshes the token and retries. 5xx, timeouts and network errors are
// retried. Reads never enter the 40/60 s write window.

import { type ParseOrderBookDetailsOptions, parseOrderBookDetails, parseOrderBooks } from '../markets/parse.js';
import type { LighterMarket, LighterOrderBookRow } from '../markets/types.js';
import { type LighterAccount, parseAccount } from './account.js';
import { type LighterActiveOrder, parseActiveOrders } from './activeOrders.js';
import { type AuthTokenCache, authHeader } from './authToken.js';
import { LighterHttpError, LighterTransportError } from './errors.js';
import { type FetchLike, INSTANCE_URLS, type LighterInstance, REST_PATHS } from './types.js';

export const DEFAULT_READ_TIMEOUT_MS = 15_000;
export const DEFAULT_READ_ATTEMPTS = 3;
export const DEFAULT_READ_RETRY_DELAY_MS = 400;

export interface ReadAttemptEvent {
  path: string;
  /** 1-based. */
  attempt: number;
  status: number | undefined;
  ok: boolean;
  durationMs: number;
  error: string | undefined;
}

export interface RestClientOptions {
  /** Instance name (`'robinhoodchain'` | `'mainnet'`) or a base URL. */
  baseUrl: string;
  /** Custom fetch (proxy or custom agent, tests). Default: global `fetch`. */
  fetch?: FetchLike;
  /** Per-attempt timeout. Default 15 000. */
  timeoutMs?: number;
  /** Total attempts for retryable failures. Default 3. `1` disables retries. */
  attempts?: number;
  /** Pause before attempt n+1 is `retryDelayMs * n`. Default 400. */
  retryDelayMs?: number;
  /** Token cache for private reads (`accountActiveOrders`). Without it those reads throw. */
  authToken?: AuthTokenCache;
  /** Sleep, for tests. */
  sleep?: (ms: number) => Promise<void>;
  /** Called after every attempt. Exceptions thrown by the callback are swallowed. */
  onAttempt?: (event: ReadAttemptEvent) => void;
}

export interface ReadOptions {
  /** Send the auth token (private read). */
  auth?: boolean;
  /** Caller's abort signal, combined with the per-attempt timeout. */
  signal?: AbortSignal;
}

export interface ReadResponse {
  status: number;
  /** Raw body: keep it for `parseJsonExact` when ids matter. */
  text: string;
}

export interface RestClient {
  readonly baseUrl: string;
  /** One GET with the retry policy; resolves only on 2xx. */
  get(path: string, opts?: ReadOptions): Promise<ReadResponse>;
  /** `get` + `JSON.parse`. Do NOT use for responses carrying order ids; use `accountActiveOrders`. */
  getJson<T = unknown>(path: string, opts?: ReadOptions): Promise<T>;
  /** `GET /api/v1/orderBooks`, decoded. */
  orderBooks(): Promise<LighterOrderBookRow[]>;
  /** `GET /api/v1/orderBookDetails`, decoded into a symbol map (perps only by default). */
  orderBookDetails(parse?: ParseOrderBookDetailsOptions): Promise<Map<string, LighterMarket>>;
  /** Raw `orderBookDetails` JSON, the shape `markets.createMarketCache` wants as `fetchDetails`. */
  orderBookDetailsRaw(): Promise<unknown>;
  /** `GET /api/v1/account?by=index&value=<accountIndex>`, decoded fail-closed. Public, no token. */
  account(accountIndex: number): Promise<LighterAccount>;
  /**
   * `GET /api/v1/accountActiveOrders?account_index=<accountIndex>` with the auth token; all markets
   * unless `marketIndex` is given. Ids are decoded from the raw text and stay exact.
   */
  accountActiveOrders(
    accountIndex: number,
    opts?: { marketIndex?: number; symbolByMarketId?: ReadonlyMap<number, string> },
  ): Promise<LighterActiveOrder[]>;
}

/** Instance name or URL -> base URL without a trailing slash. */
export function resolveBaseUrl(input: string): string {
  const trimmed = input.trim();
  const byName = (INSTANCE_URLS as Record<string, string>)[trimmed as LighterInstance];
  const url = (byName ?? trimmed).replace(/\/+$/, '');
  let parsed: URL;
  try {
    parsed = new URL(url);
  } catch {
    throw new Error(`invalid Lighter baseUrl: ${input}`);
  }
  if (!/^https?:$/.test(parsed.protocol)) throw new Error(`invalid Lighter baseUrl protocol: ${input}`);
  return url;
}

/** Which known instance a base URL belongs to, or `undefined` for a custom host. */
export function instanceOf(baseUrl: string): LighterInstance | undefined {
  const host = new URL(resolveBaseUrl(baseUrl)).host.toLowerCase();
  for (const [name, url] of Object.entries(INSTANCE_URLS)) {
    if (new URL(url).host.toLowerCase() === host) return name as LighterInstance;
  }
  return undefined;
}

const defaultSleep = (ms: number) => new Promise<void>((r) => setTimeout(r, ms));

function checkPositive(value: number, label: string): number {
  if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0) {
    throw new RangeError(`${label} must be a positive finite number, got ${String(value)}`);
  }
  return value;
}

function combineSignals(timeoutMs: number, outer?: AbortSignal): { signal: AbortSignal; dispose: () => void } {
  const ctl = new AbortController();
  const timer = setTimeout(() => ctl.abort(new Error(`timeout after ${timeoutMs} ms`)), timeoutMs);
  const onOuter = () => ctl.abort(outer?.reason);
  if (outer) {
    if (outer.aborted) ctl.abort(outer.reason);
    else outer.addEventListener('abort', onOuter, { once: true });
  }
  return {
    signal: ctl.signal,
    dispose: () => {
      clearTimeout(timer);
      outer?.removeEventListener('abort', onOuter);
    },
  };
}

export function createRestClient(opts: RestClientOptions): RestClient {
  const baseUrl = resolveBaseUrl(opts.baseUrl);
  const fetchImpl: FetchLike = opts.fetch ?? ((input, init) => fetch(input, init));
  const timeoutMs = checkPositive(opts.timeoutMs ?? DEFAULT_READ_TIMEOUT_MS, 'timeoutMs');
  const attempts = Math.max(1, Math.floor(opts.attempts ?? DEFAULT_READ_ATTEMPTS));
  const retryDelayMs = opts.retryDelayMs ?? DEFAULT_READ_RETRY_DELAY_MS;
  const sleep = opts.sleep ?? defaultSleep;

  const emit = (e: ReadAttemptEvent) => {
    try {
      opts.onAttempt?.(e);
    } catch {
      // telemetry must not break reads
    }
  };

  async function get(path: string, ro: ReadOptions = {}): Promise<ReadResponse> {
    if (ro.auth && !opts.authToken) throw new Error(`private read ${path.split('?')[0]} needs an authToken cache`);
    let lastError: unknown;
    for (let attempt = 1; attempt <= attempts; attempt++) {
      const started = Date.now();
      let status: number | undefined;
      try {
        const headers: Record<string, string> = { accept: 'application/json' };
        if (ro.auth && opts.authToken) Object.assign(headers, authHeader(await opts.authToken.get()));
        const { signal, dispose } = combineSignals(timeoutMs, ro.signal);
        let res: Response;
        let text: string;
        try {
          res = await fetchImpl(`${baseUrl}${path}`, { method: 'GET', headers, signal });
          text = await res.text();
        } finally {
          dispose();
        }
        status = res.status;
        if (res.ok) {
          emit({ path, attempt, status, ok: true, durationMs: Date.now() - started, error: undefined });
          return { status, text };
        }
        const httpError = new LighterHttpError(status, path, text);
        emit({ path, attempt, status, ok: false, durationMs: Date.now() - started, error: httpError.message });
        if ((status === 401 || status === 403) && ro.auth) {
          opts.authToken?.invalidate(); // token expired early: take a new one and retry
          lastError = httpError;
        } else if (status >= 400 && status < 500) {
          throw httpError; // the exchange's answer: never retried (429 included)
        } else {
          lastError = httpError;
        }
      } catch (err) {
        if (err instanceof LighterHttpError) throw err;
        if (ro.signal?.aborted) throw err; // the caller cancelled: not a transport failure
        lastError = err;
        emit({
          path,
          attempt,
          status,
          ok: false,
          durationMs: Date.now() - started,
          error: err instanceof Error ? err.message : String(err),
        });
      }
      if (attempt < attempts) await sleep(retryDelayMs * attempt);
    }
    if (lastError instanceof LighterHttpError) throw lastError;
    const message = lastError instanceof Error ? lastError.message : String(lastError);
    throw new LighterTransportError(path, attempts, message, { cause: lastError });
  }

  async function getJson<T>(path: string, ro?: ReadOptions): Promise<T> {
    const { text } = await get(path, ro);
    return JSON.parse(text) as T;
  }

  return {
    baseUrl,
    get,
    getJson,
    async orderBooks() {
      return parseOrderBooks(await getJson(REST_PATHS.orderBooks));
    },
    async orderBookDetails(parse) {
      return parseOrderBookDetails(await getJson(REST_PATHS.orderBookDetails), parse);
    },
    orderBookDetailsRaw() {
      return getJson(REST_PATHS.orderBookDetails);
    },
    async account(accountIndex) {
      return parseAccount(await getJson(REST_PATHS.account(accountIndex)));
    },
    async accountActiveOrders(accountIndex, o = {}) {
      const { text } = await get(REST_PATHS.accountActiveOrders(accountIndex, o.marketIndex), { auth: true });
      return parseActiveOrders(text, { symbolByMarketId: o.symbolByMarketId });
    },
  };
}
All files