src/transport/client.ts
v0.2.0 · 11.4 KB
// 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 };
}