Skip to content
markpaper

src/transport/client.ts

v0.2.0 · 11.4 KB

Download file
// Gateway client: throttled, retrying `POST /query` and `POST /execute` with envelope parsing.
//
// Transport facts (live, 2026-07-24):
//   - the gateway answers 403 to clients that do not negotiate compression. Node's fetch (undici) sends
//     `Accept-Encoding: gzip, deflate, br` and decompresses transparently - the client never sets that
//     header by hand (a manual value can disable automatic decompression);
//   - a syntactically broken query gets PLAIN TEXT, not a JSON envelope; a proxy may truncate a 200.
//     Both are `NadoResponseParseError` (transient), never "the venue said no";
//   - every request runs under `AbortSignal.timeout` (10 s by default).

import { unwrapEnvelope } from './envelope.js';
import {
  isRateLimitError,
  isTransientError,
  NadoHttpError,
  NadoNetworkError,
  NadoResponseParseError,
  NadoTimeoutError,
} from './errors.js';
import { DEFAULT_RETRY, type RetryOptions, withRetry } from './retry.js';
import { type Priority, type WeightThrottle } from './throttle.js';
import {
  type ExecuteCallOptions,
  type ExecuteRequest,
  GATEWAY_URL,
  type GatewayCallOptions,
  type Network,
  type QueryRequest,
} from './types.js';
import { executeAction, executeWeight, queryWeight } from './weights.js';

/** Default per-attempt timeout: 10 s. */
export const DEFAULT_GATEWAY_TIMEOUT_MS = 10_000;

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

/** One finished attempt, for telemetry. */
export interface GatewayAttemptEvent {
  path: 'query' | 'execute';
  /** Query `type` or execute action. */
  requestType: string;
  label: string | undefined;
  weight: number;
  attempt: number;
  ok: boolean;
  status: number | undefined;
  durationMs: number;
  error: unknown;
}

/** Options of {@link createGatewayClient}. */
export interface GatewayClientOptions {
  /** Network for the default URL. Default `'mainnet'`. */
  network?: Network;
  /**
   * Gateway base URL (`.../v1`). When both `network` and `baseUrl` are given and `baseUrl` is the
   * OFFICIAL host of the other network, the constructor throws: testnet must never reach a mainnet
   * signer (or vice versa). Prefer verifying the chain id with `createNetworkVerifier` anyway.
   */
  baseUrl?: string;
  /** Custom fetch (proxy, egress binding). Never add `Accept-Encoding` to it. */
  fetch?: FetchLike;
  /** Per-attempt timeout, ms. Default 10 000. */
  timeoutMs?: number;
  /**
   * Explicit throttle to pay weights into. Pass one shared instance when several clients hit the same
   * gateway from one IP / wallet; pass `null` only when the caller throttles every call itself.
   */
  throttle: WeightThrottle | null;
  /** Retry policy (queries: transient errors; executes: 429 always, transient only when idempotent). `false` disables. */
  retry?: RetryOptions | false;
  /** JSON parser for response text. Default `JSON.parse` (nonces are strings on the wire, so it is lossless). */
  parseJson?: (text: string) => unknown;
  /** Called after every attempt; exceptions are swallowed. */
  onAttempt?: (event: GatewayAttemptEvent) => void;
}

/** A gateway client. */
export interface GatewayClient {
  readonly network: Network;
  /** Base URL (`.../v1`). */
  readonly url: string;
  readonly throttle: WeightThrottle | null;
  /** `POST /query`; resolves with the envelope's `data`. Throws `NadoRejection` on a failure envelope. */
  query<T = unknown>(body: QueryRequest, opts?: GatewayCallOptions & { priority?: Priority }): Promise<T>;
  /** `POST /execute`; resolves with the envelope's `data`. Throws `NadoRejection` on a failure envelope (nothing applied). */
  execute<T = unknown>(body: ExecuteRequest, opts?: ExecuteCallOptions & { priority?: Priority }): Promise<T>;
}

const MAX_TIMEOUT_MS = 4_294_967_295;

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 gateway base URL from `network` / `baseUrl` and refuses a mainnet/testnet mix.
 *
 * @throws Error on an invalid URL or a network mismatch with an official host.
 */
export function resolveGatewayUrl(opts: { network?: Network; baseUrl?: string } = {}): string {
  const network = opts.network ?? 'mainnet';
  if (opts.baseUrl === undefined) return GATEWAY_URL[network];
  const trimmed = opts.baseUrl.trim().replace(/\/+$/, '');
  const host = hostOf(trimmed);
  if (!host || !/^https?:$/i.test(new URL(trimmed).protocol)) throw new Error(`invalid Nado baseUrl: ${opts.baseUrl}`);
  const other: Network = network === 'mainnet' ? 'testnet' : 'mainnet';
  if (host === hostOf(GATEWAY_URL[other])) {
    throw new Error(`Nado endpoints mix mainnet and testnet: network=${network}, baseUrl=${opts.baseUrl}`);
  }
  return trimmed;
}

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 gateway client. Each attempt pays its weight into the throttle, takes a concurrency slot,
 * runs fetch under a timeout, parses the body and unwraps the envelope. Retries go back through the
 * throttle. The client resolves ONLY with the envelope's `data` and otherwise throws - a failed read is
 * never an empty account.
 */
export function createGatewayClient(options: GatewayClientOptions): GatewayClient {
  if (!options || !Object.prototype.hasOwnProperty.call(options, 'throttle')) {
    throw new TypeError('createGatewayClient requires an explicit throttle (or null for external throttling)');
  }
  const network = options.network ?? 'mainnet';
  const url = resolveGatewayUrl({ network, baseUrl: options.baseUrl });
  const fetchImpl: FetchLike = options.fetch ?? ((input, init) => fetch(input, init));
  const throttle = options.throttle;
  const defaultTimeout = checkTimeout(options.timeoutMs ?? DEFAULT_GATEWAY_TIMEOUT_MS);
  const retry = options.retry === false ? false : { ...DEFAULT_RETRY, ...(options.retry ?? {}) };
  const parse = options.parseJson ?? JSON.parse;
  const report = (event: GatewayAttemptEvent): void => {
    try {
      options.onAttempt?.(event);
    } catch {
      // Telemetry must never turn a successful call into a failure.
    }
  };

  async function attemptOnce(
    path: 'query' | 'execute',
    requestType: string,
    payload: string,
    weight: number,
    attempt: number,
    opts: GatewayCallOptions,
  ): Promise<unknown> {
    const target = `${url}/${path}`;
    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;

    const fail = (err: unknown): never => {
      report({
        path,
        requestType,
        label: opts.label,
        weight,
        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 NadoTimeoutError({ timeoutMs, url: target, requestType, cause: err });
      return new NadoNetworkError({ url: target, requestType, cause: err });
    };

    try {
      let res: Response;
      try {
        // No Accept-Encoding here: undici negotiates and decompresses on its own.
        res = await fetchImpl(target, {
          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 {
          bodyText = await res.text(); // reading the body releases the connection
        } catch {
          await res.body?.cancel().catch(() => {});
        }
        return fail(new NadoHttpError({ status: res.status, bodyText, url: target, requestType }));
      }
      let text: string;
      try {
        text = await res.text();
      } catch (err) {
        return fail(mapAbortOrNetwork(err));
      }
      let body: unknown;
      try {
        body = parse(text);
      } catch (err) {
        return fail(
          new NadoResponseParseError({ status: res.status, bodyText: text, url: target, requestType, cause: err }),
        );
      }
      let data: unknown;
      try {
        data = unwrapEnvelope(body, path, requestType);
      } catch (err) {
        return fail(err);
      }
      report({
        path,
        requestType,
        label: opts.label,
        weight,
        attempt,
        ok: true,
        status,
        durationMs: performance.now() - started,
        error: undefined,
      });
      return data;
    } finally {
      combined.dispose();
    }
  }

  async function send(
    path: 'query' | 'execute',
    requestType: string,
    body: unknown,
    weight: number,
    opts: GatewayCallOptions & { priority?: Priority },
    shouldRetry: (err: unknown) => boolean,
  ): Promise<unknown> {
    if (opts.timeoutMs !== undefined) checkTimeout(opts.timeoutMs); // fail before paying weight
    const payload = JSON.stringify(body);
    const run = (attempt: number): Promise<unknown> => {
      const exec = () => attemptOnce(path, requestType, payload, weight, attempt, opts);
      if (!throttle) return exec();
      const callOpts = { label: opts.label ?? requestType, signal: opts.signal, priority: opts.priority };
      return path === 'query' ? throttle.query(weight, exec, callOpts) : throttle.execute(weight, exec, callOpts);
    };
    return retry ? withRetry(run, { ...retry, signal: opts.signal, shouldRetry }) : run(0);
  }

  async function query<T = unknown>(
    body: QueryRequest,
    opts: GatewayCallOptions & { priority?: Priority } = {},
  ): Promise<T> {
    if (!body || typeof body.type !== 'string' || body.type.length === 0)
      throw new TypeError('query body must have a non-empty string "type"');
    const weight = opts.weight ?? queryWeight(body);
    return (await send('query', body.type, body, weight, opts, isTransientError)) as T;
  }

  async function execute<T = unknown>(
    body: ExecuteRequest,
    opts: ExecuteCallOptions & { priority?: Priority } = {},
  ): Promise<T> {
    const action = executeAction(body);
    if (!action) throw new TypeError('execute body must have exactly one action key');
    const weight = opts.weight ?? executeWeight(body);
    const shouldRetry = opts.idempotent === true ? isTransientError : isRateLimitError;
    return (await send('execute', action, body, weight, opts, shouldRetry)) as T;
  }

  return { network, url, throttle, query, execute };
}
All files