src/rest/client.ts
v0.1.0 · 30.4 KB
import {
assertQfexWriteAllowed,
buildRestAuthHeaders,
classifyQfexAuthFailure,
createNonceSource,
maskSecrets,
type QfexAuthFailure,
type QfexKeyPair,
qfexRestIsWrite,
} from '../auth/auth.js';
import { type Dec, decCmp, decSign, parseDec, parseLossless } from '../numbers/index.js';
import { QFEX_HOSTS } from './hosts.js';
export interface QfexFetchResponse {
status: number;
headers: {
get(name: string): string | null;
};
text(): Promise<string>;
}
export type QfexFetch = (
url: string,
init: {
method: string;
headers: Record<string, string>;
body?: string;
signal: AbortSignal;
redirect: 'manual';
},
) => Promise<QfexFetchResponse>;
export interface QfexRequestOptions {
tries?: number;
timeoutMs?: number;
maxWaitMs?: number;
}
export type QfexHttpErrorKind = 'auth' | 'client' | 'server' | 'transport' | 'rate_limited' | 'stale_cache' | 'decode';
export interface QfexProblem {
type: string | null;
title: string | null;
status: number | null;
detail: string | null;
errors: string[];
}
export class QfexHttpError extends Error {
readonly kind: QfexHttpErrorKind;
readonly status: number;
readonly body: string;
readonly method: string;
readonly path: string;
readonly problem: QfexProblem | null;
readonly code: string | null;
readonly auth: QfexAuthFailure | null;
readonly retryAfterMs: number | null;
readonly attempts: number;
readonly outcomeUnknown: boolean;
readonly skewMs: number | null;
constructor(
message: string,
f: {
kind: QfexHttpErrorKind;
status: number;
method: string;
path: string;
body?: string;
problem?: QfexProblem | null;
code?: string | null;
auth?: QfexAuthFailure | null;
retryAfterMs?: number | null;
attempts: number;
outcomeUnknown?: boolean;
skewMs?: number | null;
},
) {
super(maskSecrets(message));
this.name = 'QfexHttpError';
this.kind = f.kind;
this.status = f.status;
this.method = f.method;
this.path = f.path;
this.body = maskSecrets(f.body ?? '');
this.problem = f.problem ?? null;
this.code = f.code ?? null;
this.auth = f.auth ?? null;
this.retryAfterMs = f.retryAfterMs ?? null;
this.attempts = f.attempts;
this.outcomeUnknown = f.outcomeUnknown ?? false;
this.skewMs = f.skewMs ?? null;
}
}
export interface QfexResponseMeta {
status: number;
sentAt: number;
receivedAt: number;
rttMs: number;
serverDateMs: number | null;
xCache: string | null;
ageSec: number | null;
cached: boolean;
attempts: number;
}
export interface QfexClockSkew {
skewMs: number;
loMs: number;
hiMs: number;
at: number;
rttMs: number;
path: string;
}
export interface RateBackoffBounds {
minMs: number;
maxMs: number;
}
export function skewAbsLowerBoundMs(s: Pick<QfexClockSkew, 'loMs' | 'hiMs'>): number {
if (s.loMs <= 0 && s.hiMs >= 0) return 0;
return Math.min(Math.abs(s.loMs), Math.abs(s.hiMs));
}
export function computeClockSkew(
dateHeader: string | null,
sentAt: number,
receivedAt: number,
path = '',
): QfexClockSkew | null {
if (!dateHeader) return null;
const d = Date.parse(dateHeader);
if (!Number.isFinite(d)) return null;
const loMs = d - receivedAt;
const hiMs = d + 1000 - sentAt;
return {
skewMs: Math.round((loMs + hiMs) / 2),
loMs,
hiMs,
at: receivedAt,
rttMs: Math.max(0, receivedAt - sentAt),
path,
};
}
export function isCdnCached(xCache: string | null, age: string | null): boolean {
if (xCache && /hit/i.test(xCache)) return true;
if (age !== null && age.trim() !== '') {
const t = age.trim();
if (!/^\d+$/.test(t)) return true;
return Number(t) > 0;
}
return false;
}
export function isCdnEdgeRejection(contentType: string | null, body: string | null | undefined): boolean {
if (contentType && /text\/html/i.test(contentType)) return true;
const t = String(body ?? '')
.trimStart()
.slice(0, 4000)
.toLowerCase();
return t.startsWith('<') || t.includes('generated by cloudfront') || t.includes('request could not be satisfied');
}
export function parseRetryAfterMs(
header: string | null,
body: string | null,
nowMs: number,
defaultMs: number,
bounds: RateBackoffBounds,
): number {
if (
!Number.isFinite(defaultMs) ||
defaultMs <= 0 ||
!Number.isFinite(bounds?.minMs) ||
bounds.minMs <= 0 ||
!Number.isFinite(bounds.maxMs) ||
bounds.maxMs < bounds.minMs
)
throw new TypeError('Required retry-after policy bounds');
const clamp = (ms: number): number => Math.min(bounds.maxMs, Math.max(bounds.minMs, Math.ceil(ms)));
const h = (header ?? '').trim();
if (h) {
if (/^\d+(?:\.\d+)?$/.test(h)) return clamp(Number(h) * 1000);
const d = Date.parse(h);
if (Number.isFinite(d)) return clamp(d - nowMs);
}
const m =
/retry[\s_-]*after\D{0,3}?(\d+(?:\.\d+)?)\s*(milliseconds?|millis|msec|ms|minutes?|mins?|seconds?|secs?|s)?\b/i.exec(
body ?? '',
);
if (m) {
const v = Number(m[1]);
const unit = (m[2] ?? 's').toLowerCase();
const ms = unit.startsWith('mi') && !unit.startsWith('mil') ? v * 60000 : unit.startsWith('m') ? v : v * 1000;
if (Number.isFinite(ms)) return clamp(ms);
}
return clamp(defaultMs);
}
export function parseQfexErrorBody(text: string): {
problem: QfexProblem | null;
code: string | null;
message: string | null;
} {
let obj: unknown;
try {
obj = JSON.parse(text);
} catch {
return { problem: null, code: null, message: null };
}
if (!obj || typeof obj !== 'object' || Array.isArray(obj)) return { problem: null, code: null, message: null };
const o = obj as Record<string, unknown>;
const str = (v: unknown): string | null => (typeof v === 'string' ? v : null);
const nested = o.err && typeof o.err === 'object' ? (o.err as Record<string, unknown>) : null;
const code = str(o.error_code) ?? str(nested?.error_code) ?? null;
let problem: QfexProblem | null = null;
if ('title' in o || 'detail' in o || ('status' in o && typeof o.status === 'number')) {
const errors: string[] = [];
if (Array.isArray(o.errors)) {
for (const it of o.errors) {
if (it && typeof it === 'object') {
const r = it as Record<string, unknown>;
const msg = str(r.message);
if (msg) errors.push(str(r.location) ? `${r.location}: ${msg}` : msg);
} else if (typeof it === 'string') errors.push(it);
}
}
problem = {
type: str(o.type),
title: str(o.title),
status: typeof o.status === 'number' && Number.isInteger(o.status) ? o.status : null,
detail: str(o.detail),
errors,
};
}
const message = str(o.message) ?? str(nested?.message) ?? problem?.detail ?? problem?.title ?? null;
return { problem, code, message };
}
function stripQuery(p: string): string {
const i = p.indexOf('?');
return i < 0 ? p : p.slice(0, i);
}
export function qfexPath(path: string, query?: Record<string, string | number | boolean | null | undefined>): string {
if (!path.startsWith('/')) throw new Error(`REST path must start with a slash: ${path}`);
if (!query) return path;
const parts: string[] = [];
for (const [k, v] of Object.entries(query)) {
if (v === undefined || v === null) continue;
parts.push(`${encodeURIComponent(k)}=${encodeURIComponent(String(v))}`);
}
return parts.length ? `${path}${path.includes('?') ? '&' : '?'}${parts.join('&')}` : path;
}
export interface QfexContract {
tickerId: string;
lastPrice: string | null;
indexPrice: string | null;
raw: Record<string, unknown>;
}
function positivePx(v: unknown): string | null {
if (typeof v !== 'string') return null;
const d = parseDec(v);
return d && decSign(d) > 0 ? v : null;
}
export function decodeQfexContracts(payload: unknown): QfexContract[] {
if (!payload || typeof payload !== 'object' || Array.isArray(payload))
throw new Error('Contracts response must be an object');
const rows = (
payload as {
data?: unknown;
}
).data;
if (!Array.isArray(rows)) throw new Error('Contracts data must be an array');
const out: QfexContract[] = [];
for (const r of rows) {
if (!r || typeof r !== 'object' || Array.isArray(r)) throw new Error('Contract row must be an object');
const row = r as Record<string, unknown>;
if (typeof row.ticker_id !== 'string' || row.ticker_id === '')
throw new Error('Contract ticker_id must be a nonempty string');
out.push({
tickerId: row.ticker_id,
lastPrice: positivePx(row.last_price),
indexPrice: positivePx(row.index_price),
raw: row,
});
}
return out;
}
export interface QfexBookLevel {
px: string;
sz: string;
pxDec: Dec;
szDec: Dec;
}
export interface QfexOrderbook {
tickerId: string;
timestampMs: number | null;
bids: QfexBookLevel[];
asks: QfexBookLevel[];
zeroLevelsDropped: number;
crossed: boolean;
}
function decodeSide(
rows: unknown,
side: string,
): {
levels: QfexBookLevel[];
zeros: number;
} {
if (!Array.isArray(rows)) throw new Error(`${side} must be an array`);
const levels: QfexBookLevel[] = [];
let zeros = 0;
for (const r of rows) {
if (!Array.isArray(r) || r.length < 2 || typeof r[0] !== 'string' || typeof r[1] !== 'string')
throw new Error(`${side} level must contain price and size strings`);
const pxDec = parseDec(r[0]);
const szDec = parseDec(r[1]);
if (!pxDec || !szDec)
throw new Error(
`${side} has an invalid decimal level: ${String(r[0]).slice(0, 20)}, ${String(r[1]).slice(0, 20)}`,
);
if (decSign(szDec) < 0) throw new Error(`${side} contains a negative size: ${r[1]}`);
if (decSign(szDec) === 0) {
zeros++;
continue;
}
if (decSign(pxDec) <= 0) throw new Error(`${side} contains a nonpositive price: ${r[0]}`);
levels.push({ px: r[0], sz: r[1], pxDec, szDec });
}
return { levels, zeros };
}
export function decodeQfexOrderbook(payload: unknown, expectTicker?: string): QfexOrderbook {
if (!payload || typeof payload !== 'object' || Array.isArray(payload))
throw new Error('Orderbook response must be an object');
const o = payload as Record<string, unknown>;
if (typeof o.ticker_id !== 'string' || o.ticker_id === '')
throw new Error('Orderbook ticker_id must be a nonempty string');
if (expectTicker !== undefined && o.ticker_id !== expectTicker)
throw new Error(`Orderbook ticker_id=${o.ticker_id}; expected ${expectTicker}`);
const b = decodeSide(o.bids, 'bids');
const a = decodeSide(o.asks, 'asks');
b.levels.sort((x, y) => decCmp(y.pxDec, x.pxDec));
a.levels.sort((x, y) => decCmp(x.pxDec, y.pxDec));
const ts = typeof o.timestamp === 'string' && /^\d+$/.test(o.timestamp) ? Number(o.timestamp) : null;
const crossed = b.levels.length > 0 && a.levels.length > 0 && decCmp(b.levels[0]!.pxDec, a.levels[0]!.pxDec) >= 0;
return {
tickerId: o.ticker_id,
timestampMs: ts !== null && Number.isSafeInteger(ts) ? ts : null,
bids: b.levels,
asks: a.levels,
zeroLevelsDropped: b.zeros + a.zeros,
crossed,
};
}
export interface RestClientOptions {
baseUrl: 'mainnet' | 'uat' | string;
fetch?: QfexFetch;
keys?: QfexKeyPair | (() => QfexKeyPair);
now?: () => number;
sleep?: (ms: number) => Promise<void>;
restMaxRps: number;
rateLimitDefaultBackoffMs: number;
rateBackoff: RateBackoffBounds;
retryDelayMs: (attempt: number) => number;
readTries: number;
timeoutMs: number;
maxWaitMs: number;
skewHistoryLimit: number;
readOnly: boolean;
strictAccountScope?: boolean;
onRateLimited?: (ms: number, message: string | null) => void;
onLog?: (message: string) => void;
}
export function createRestClient(options: RestClientOptions) {
for (const key of ['restMaxRps', 'rateLimitDefaultBackoffMs', 'readTries', 'timeoutMs', 'skewHistoryLimit'] as const)
if (!Number.isFinite(options?.[key]) || options[key] <= 0)
throw new TypeError('Required positive REST policy: ' + key);
if (!Number.isFinite(options.maxWaitMs) || options.maxWaitMs < 0 || typeof options.readOnly !== 'boolean')
throw new TypeError('Required REST readOnly and maxWaitMs');
if (typeof options.retryDelayMs !== 'function') throw new TypeError('Required REST retry delay policy');
parseRetryAfterMs(null, null, 0, options.rateLimitDefaultBackoffMs, options.rateBackoff);
if (!Number.isSafeInteger(options.readTries) || !Number.isSafeInteger(options.skewHistoryLimit))
throw new TypeError('REST retry and skew counts must be integers');
const base =
options.baseUrl === 'mainnet' || options.baseUrl === 'uat' ? QFEX_HOSTS[options.baseUrl].rest : options.baseUrl;
const parsed = new URL(base);
if (
!['http:', 'https:'].includes(parsed.protocol) ||
parsed.username ||
parsed.password ||
parsed.search ||
parsed.hash
)
throw new TypeError('Invalid REST base URL');
const log = {
warn: (text: string) => {
try {
options.onLog?.(text);
} catch {}
},
};
const nonce = createNonceSource();
const requireKeys = () => {
const keys = typeof options.keys === 'function' ? options.keys() : options.keys;
if (!keys) throw new TypeError('Signed request requires explicit keys');
return keys;
};
const settings = () => ({
apiUrl: base.replace(/\/+$/, ''),
fetch: options.fetch ?? ((url, init) => fetch(url, init)),
now: options.now ?? Date.now,
sleep: options.sleep ?? ((ms: number) => new Promise<void>((resolve) => setTimeout(resolve, ms))),
restMaxRps: options.restMaxRps,
rateLimitDefaultBackoffMs: options.rateLimitDefaultBackoffMs,
dryRun: () => options.readOnly,
restSlot: null as null | (() => Promise<void>),
onRateLimited: options.onRateLimited,
});
type Resolved = ReturnType<typeof settings>;
const bodyForLog = (text: string) => maskSecrets(text, { secret: requireKeysOptional()?.secret }).slice(0, 500);
const requireKeysOptional = () => (typeof options.keys === 'function' ? options.keys() : options.keys);
let nextRequestAt = 0;
let rateBackoffUntil = 0;
let lastCacheWarnAt = 0;
const skewHistory: QfexClockSkew[] = [];
const streak401 = { calls: 0, responses: 0, firstAt: null as number | null, lastAt: null as number | null };
const stats = {
requests: 0,
retries: 0,
rateLimited: 0,
cacheRejects: 0,
lastOkAt: null as number | null,
lastFailAt: null as number | null,
lastFail: null as string | null,
};
function qfexClientStats(): {
requests: number;
retries: number;
rateLimited: number;
cacheRejects: number;
rateBackoffUntil: number;
lastOkAt: number | null;
lastFailAt: number | null;
lastFail: string | null;
} {
return { ...stats, rateBackoffUntil };
}
function lastServerClockSkewMs(): QfexClockSkew | null {
return skewHistory.length ? { ...skewHistory[skewHistory.length - 1]! } : null;
}
function qfexClockSkewSamples(): QfexClockSkew[] {
return skewHistory.map((s) => ({ ...s }));
}
function qfexClockSkewMedianMs(): number | null {
if (!skewHistory.length) return null;
const v = skewHistory.map((s) => s.skewMs).sort((a, b) => a - b);
const m = v.length >> 1;
return v.length % 2 ? v[m]! : Math.round((v[m - 1]! + v[m]!) / 2);
}
function noteSkew(s: QfexClockSkew): void {
skewHistory.push(s);
if (skewHistory.length > options.skewHistoryLimit)
skewHistory.splice(0, skewHistory.length - options.skewHistoryLimit);
}
function qfexAuth401Streak(): {
calls: number;
responses: number;
firstAt: number | null;
lastAt: number | null;
} {
return { ...streak401 };
}
function note401Response(at: number): void {
streak401.responses++;
if (streak401.firstAt === null) streak401.firstAt = at;
streak401.lastAt = at;
}
function note401Call(): void {
streak401.calls++;
}
function clear401(): void {
streak401.calls = 0;
streak401.responses = 0;
streak401.firstAt = null;
streak401.lastAt = null;
}
interface Spec {
method: 'GET' | 'POST';
path: string;
signed: boolean;
accountId: string | null;
body?: string;
tries: number;
timeoutMs: number;
maxWaitMs: number;
}
const isWrite = (spec: Spec): boolean => qfexRestIsWrite(spec.method);
async function takeSlot(s: Resolved, spec: Spec, attempts: number): Promise<void> {
const gap = s.restMaxRps > 0 ? 1000 / s.restMaxRps : 0;
for (let round = 0; ; round++) {
const now = s.now();
const wait = rateBackoffUntil - now;
if (wait > spec.maxWaitMs || round >= 8) {
fail(
new QfexHttpError(
`qfex ${spec.method} ${stripQuery(spec.path)}: rate hold exceeds request wait budget (${Math.ceil(Math.max(0, wait) / 1000)} seconds)`,
{
kind: 'rate_limited',
status: 0,
method: spec.method,
path: stripQuery(spec.path),
attempts,
retryAfterMs: Math.max(0, wait),
},
),
now,
);
}
const at = Math.max(now, nextRequestAt, rateBackoffUntil);
nextRequestAt = at + gap;
if (at > now) await s.sleep(at - now);
if (s.now() >= rateBackoffUntil) break;
}
if (s.restSlot) await s.restSlot();
}
async function sendOnce(
s: Resolved,
spec: Spec,
headers: Record<string, string>,
): Promise<{
res: QfexFetchResponse;
text: string;
}> {
const ctrl = new AbortController();
const t = setTimeout(() => ctrl.abort(), spec.timeoutMs);
try {
const res = await s.fetch(`${s.apiUrl}${spec.path}`, {
method: spec.method,
headers,
body: spec.body,
signal: ctrl.signal,
redirect: 'manual',
});
const text = await res.text();
return { res, text };
} finally {
clearTimeout(t);
}
}
const transientDelayMs = (n: number): number => {
const delay = options.retryDelayMs(n);
if (!Number.isFinite(delay) || delay < 0) throw new TypeError('Retry delay must be finite and nonnegative');
return delay;
};
function fail(e: QfexHttpError, now: number): never {
stats.lastFailAt = now;
stats.lastFail = `${e.kind}: ${e.message}`.slice(0, 300);
throw e;
}
async function request(spec: Spec): Promise<{
data: unknown;
meta: QfexResponseMeta;
}> {
const s = settings();
const path = stripQuery(spec.path);
if (isWrite(spec)) assertQfexWriteAllowed(`${spec.method} ${path}`, s.dryRun());
let attempts = 0;
let transient = 0;
let authRetried = false;
for (;;) {
await takeSlot(s, spec, attempts);
const headers: Record<string, string> = { accept: 'application/json' };
if (spec.body !== undefined) headers['content-type'] = 'application/json';
if (spec.signed) {
Object.assign(headers, buildRestAuthHeaders(requireKeys(), spec.accountId, s.now(), nonce.next(s.now())));
headers['cache-control'] = 'no-cache';
}
const sentAt = s.now();
attempts++;
stats.requests++;
if (attempts > 1) stats.retries++;
let res: QfexFetchResponse;
let text: string;
try {
({ res, text } = await sendOnce(s, spec, headers));
} catch (err) {
const aborted =
(
err as {
name?: string;
}
)?.name === 'AbortError';
const why = aborted
? `request timed out after ${spec.timeoutMs} ms`
: bodyForLog(String((err as Error)?.message ?? err));
const e = new QfexHttpError(`qfex ${spec.method} ${path}: transport failure: ${why}`, {
kind: 'transport',
status: 0,
method: spec.method,
path,
attempts,
outcomeUnknown: isWrite(spec),
});
transient++;
if (isWrite(spec) || transient >= spec.tries) fail(e, s.now());
await s.sleep(transientDelayMs(transient));
continue;
}
const receivedAt = s.now();
const status = res.status;
const xCache = res.headers.get('x-cache');
const ageRaw = res.headers.get('age');
const cached = isCdnCached(xCache, ageRaw);
const dateHeader = res.headers.get('date');
const skew = cached ? null : computeClockSkew(dateHeader, sentAt, receivedAt, path);
if (skew) noteSkew(skew);
const ageNum = ageRaw !== null && /^\d+$/.test(ageRaw.trim()) ? Number(ageRaw.trim()) : null;
const serverDate = dateHeader ? Date.parse(dateHeader) : NaN;
const meta: QfexResponseMeta = {
status,
sentAt,
receivedAt,
rttMs: Math.max(0, receivedAt - sentAt),
serverDateMs: Number.isFinite(serverDate) ? serverDate : null,
xCache,
ageSec: ageNum,
cached,
attempts,
};
const base = { method: spec.method, path, attempts, body: bodyForLog(text), skewMs: skew?.skewMs ?? null };
if (spec.signed && cached) {
stats.cacheRejects++;
if (receivedAt - lastCacheWarnAt > 60000) {
lastCacheWarnAt = receivedAt;
log.warn(
`qfex ${spec.method} ${path}: cached signed response rejected (X-Cache=${xCache ?? '-'}, Age=${ageRaw ?? '-'})`,
);
}
const e = new QfexHttpError(
`qfex ${spec.method} ${path}: cached signed response rejected (X-Cache=${xCache ?? '-'}, Age=${ageRaw ?? '-'})`,
{
...base,
kind: 'stale_cache',
status,
outcomeUnknown: isWrite(spec),
},
);
transient++;
if (isWrite(spec) || transient >= spec.tries) fail(e, receivedAt);
await s.sleep(transientDelayMs(transient));
continue;
}
const parsedErr =
status >= 400
? parseQfexErrorBody(maskSecrets(text, { secret: requireKeysOptional()?.secret }))
: { problem: null, code: null, message: null };
if (status === 429 || (status >= 400 && parsedErr.code === 'RateLimited')) {
const waitMs = parseRetryAfterMs(
res.headers.get('retry-after'),
parsedErr.message ?? text,
receivedAt,
s.rateLimitDefaultBackoffMs,
options.rateBackoff,
);
rateBackoffUntil = Math.max(rateBackoffUntil, receivedAt + waitMs);
stats.rateLimited++;
try {
s.onRateLimited?.(waitMs, parsedErr.message);
} catch (hookErr) {
log.warn(
`Rate-limit callback failed: ${maskSecrets(String((hookErr as Error)?.message ?? hookErr), { secret: requireKeysOptional()?.secret })}`,
);
}
const e = new QfexHttpError(
`qfex ${spec.method} ${path}: HTTP ${status}; rate hold ${Math.ceil(waitMs / 1000)} seconds`,
{
...base,
kind: 'rate_limited',
status,
problem: parsedErr.problem,
code: parsedErr.code ?? 'RateLimited',
retryAfterMs: waitMs,
},
);
transient++;
if (transient >= spec.tries) fail(e, receivedAt);
continue;
}
if ((status === 401 || status === 403) && isCdnEdgeRejection(res.headers.get('content-type'), text)) {
const e = new QfexHttpError(`qfex ${spec.method} ${path}: HTTP ${status}; CDN edge rejection`, {
...base,
kind: 'client',
status,
});
fail(e, receivedAt);
}
const af = status === 401 || status === 403 ? classifyQfexAuthFailure(status, text) : null;
if (af) {
if (spec.signed && status === 401) note401Response(receivedAt);
if (spec.signed && af.retryFreshNonce && !authRetried) {
authRetried = true;
continue;
}
if (spec.signed && status === 401) note401Call();
const e = new QfexHttpError(`qfex ${spec.method} ${path}: HTTP ${status} ${bodyForLog(text).slice(0, 120)}`, {
...base,
kind: 'auth',
status,
auth: af,
problem: parsedErr.problem,
code: parsedErr.code,
});
fail(e, receivedAt);
}
if (status >= 500) {
const e = new QfexHttpError(
`qfex ${spec.method} ${path}: HTTP ${status}${parsedErr.message ? ` ${parsedErr.message}` : ''}`,
{
...base,
kind: 'server',
status,
problem: parsedErr.problem,
code: parsedErr.code,
outcomeUnknown: isWrite(spec),
},
);
transient++;
if (isWrite(spec) || transient >= spec.tries) fail(e, receivedAt);
await s.sleep(transientDelayMs(transient));
continue;
}
if (status >= 300 || status < 200) {
const detail = parsedErr.problem?.detail ?? parsedErr.message ?? bodyForLog(text).slice(0, 120);
const e = new QfexHttpError(`qfex ${spec.method} ${path}: HTTP ${status}${detail ? ` ${detail}` : ''}`, {
...base,
kind: 'client',
status,
problem: parsedErr.problem,
code: parsedErr.code,
});
fail(e, receivedAt);
}
let data: unknown = null;
if (text.trim() !== '') {
try {
data = parseLossless(text);
} catch (pe) {
const e = new QfexHttpError(
`qfex ${spec.method} ${path}: HTTP ${status}; invalid JSON (${String((pe as Error)?.message ?? pe).slice(0, 80)})`,
{
...base,
kind: 'decode',
status,
outcomeUnknown: isWrite(spec),
},
);
transient++;
if (isWrite(spec) || transient >= spec.tries) fail(e, receivedAt);
await s.sleep(transientDelayMs(transient));
continue;
}
}
if (spec.signed) clear401();
stats.lastOkAt = receivedAt;
return { data, meta };
}
}
function opts(
o: QfexRequestOptions | undefined,
defTries: number,
): {
tries: number;
timeoutMs: number;
maxWaitMs: number;
} {
const tries = Math.max(1, Math.floor(o?.tries ?? defTries));
return {
tries,
timeoutMs: Math.max(1, o?.timeoutMs ?? options.timeoutMs),
maxWaitMs: Math.max(0, o?.maxWaitMs ?? options.maxWaitMs),
};
}
function checkPath(path: string): void {
if (typeof path !== 'string' || !path.startsWith('/') || path.startsWith('//'))
throw new Error(`Invalid REST path: ${path}`);
}
function checkAccountScope(method: string, path: string, accountId: string | null): void {
if (accountId !== null) return;
if (method === 'GET' && stripQuery(path).startsWith('/user/subaccounts/')) return;
throw new Error(`qfex ${method} ${stripQuery(path)} requires an explicit account ID`);
}
async function qfexPublicGetWithMeta<T>(
path: string,
o?: QfexRequestOptions,
): Promise<{
data: T;
meta: QfexResponseMeta;
}> {
checkPath(path);
const r = await request({ method: 'GET', path, signed: false, accountId: null, ...opts(o, options.readTries) });
return { data: r.data as T, meta: r.meta };
}
async function qfexPublicGet<T>(path: string, o?: QfexRequestOptions): Promise<T> {
return (await qfexPublicGetWithMeta<T>(path, o)).data;
}
async function qfexUserGetWithMeta<T>(
path: string,
accountId: string | null,
o?: QfexRequestOptions,
): Promise<{
data: T;
meta: QfexResponseMeta;
}> {
checkPath(path);
if (options.strictAccountScope !== false) checkAccountScope('GET', path, accountId);
const r = await request({ method: 'GET', path, signed: true, accountId, ...opts(o, options.readTries) });
return { data: r.data as T, meta: r.meta };
}
async function qfexUserGet<T>(path: string, accountId: string | null, o?: QfexRequestOptions): Promise<T> {
return (await qfexUserGetWithMeta<T>(path, accountId, o)).data;
}
async function qfexUserPost<T>(path: string, accountId: string, body: unknown, o?: QfexRequestOptions): Promise<T> {
checkPath(path);
if (options.strictAccountScope !== false) checkAccountScope('POST', path, accountId ?? null);
const r = await request({
method: 'POST',
path,
signed: true,
accountId,
body: JSON.stringify(body ?? {}),
...opts(o, options.readTries),
});
return r.data as T;
}
function decodeErr(path: string, why: string): QfexHttpError {
return new QfexHttpError(`qfex GET ${path}: ${why}`, {
kind: 'decode',
status: 200,
method: 'GET',
path,
attempts: 1,
});
}
async function fetchQfexRefdata(
ticker?: string,
o?: QfexRequestOptions,
): Promise<{
payload: {
type?: string;
data: unknown[];
};
meta: QfexResponseMeta;
}> {
const path = qfexPath('/refdata', { ticker });
const { data, meta } = await qfexPublicGetWithMeta<unknown>(path, o);
if (!data || typeof data !== 'object' || Array.isArray(data))
throw decodeErr('/refdata', 'Refdata response must be an object');
const obj = data as {
type?: unknown;
data?: unknown;
};
if (obj.type !== undefined && obj.type !== 'refdata')
throw decodeErr('/refdata', `Unexpected refdata type: ${String(obj.type).slice(0, 40)}`);
if (!Array.isArray(obj.data)) throw decodeErr('/refdata', 'Refdata data must be an array');
return { payload: { type: typeof obj.type === 'string' ? obj.type : undefined, data: obj.data }, meta };
}
async function fetchQfexContracts(symbol?: string, o?: QfexRequestOptions): Promise<QfexContract[]> {
const data = await qfexPublicGet<unknown>(qfexPath('/md/contracts', { symbol }), o);
try {
return decodeQfexContracts(data);
} catch (e) {
throw decodeErr('/md/contracts', String((e as Error)?.message ?? e));
}
}
async function fetchQfexOrderbook(tickerId: string, o?: QfexRequestOptions): Promise<QfexOrderbook> {
if (typeof tickerId !== 'string' || !/^[A-Za-z0-9._]+-[A-Za-z0-9]+$/.test(tickerId))
throw new Error(`Invalid orderbook ticker_id: ${tickerId}`);
const path = `/md/orderbook/${encodeURIComponent(tickerId)}`;
const data = await qfexPublicGet<unknown>(path, o);
try {
return decodeQfexOrderbook(data, tickerId);
} catch (e) {
throw decodeErr(path, String((e as Error)?.message ?? e));
}
}
return {
publicGet: qfexPublicGet,
publicGetWithMeta: qfexPublicGetWithMeta,
userGet: qfexUserGet,
userGetWithMeta: qfexUserGetWithMeta,
userPost: qfexUserPost,
fetchRefdata: fetchQfexRefdata,
fetchContracts: fetchQfexContracts,
fetchOrderbook: fetchQfexOrderbook,
clockSkew: lastServerClockSkewMs,
clockSkewSamples: qfexClockSkewSamples,
clockSkewMedian: qfexClockSkewMedianMs,
auth401Streak: qfexAuth401Streak,
stats: qfexClientStats,
};
}
export type RestClient = ReturnType<typeof createRestClient>;