src/testing/fake.ts
v0.1.0 · 5.3 KB
import { EventEmitter } from 'node:events';
import WebSocket from 'ws';
import type { QfexOrderEvent } from '../frames/index.js';
import { parseLossless } from '../numbers/index.js';
import { FIXTURE_ACCOUNT, fixtureOrder, wireOrder } from './fixtures.js';
/** Worst-case controls are simulation hypotheses, not claims about venue behaviour. */
export interface FakeOptions {
ackDelayMs: number;
readDelayMs: number;
dropAddReplies: boolean;
rejectAfterAck?: string;
canonicalAccountId: string;
}
export function createFakeExchange(options: FakeOptions) {
const sockets = new Set<FakeSocket>();
const orders = new Map<string, QfexOrderEvent>();
const sent: Array<Record<string, unknown>> = [];
const timers = new Set<ReturnType<typeof setTimeout>>();
let sequence = 0;
let closing = false;
function later(ms: number, fn: () => void) {
const timer = setTimeout(() => {
timers.delete(timer);
if (!closing) fn();
}, ms);
timers.add(timer);
}
class FakeSocket extends EventEmitter {
readyState: number = WebSocket.CONNECTING;
bufferedAmount = 0;
constructor() {
super();
sockets.add(this);
queueMicrotask(() => {
if (closing) return;
this.readyState = WebSocket.OPEN;
this.emit('open');
});
}
frame(value: unknown) {
queueMicrotask(() => {
if (this.readyState === WebSocket.OPEN) this.emit('message', Buffer.from(JSON.stringify(value)), false);
});
}
send(text: string, callback?: (error?: Error) => void) {
if (this.readyState !== WebSocket.OPEN) throw new Error('Fake socket closed');
if (text === 'ping') {
this.frame('pong');
callback?.();
return;
}
const request = parseLossless(text) as { type: string; params?: Record<string, unknown> };
sent.push(request);
const params = request.params ?? {};
switch (request.type) {
case 'auth':
this.frame({ authenticated: params.account_id === options.canonicalAccountId });
break;
case 'subscribe':
for (const channel of (params.channels as string[] | undefined) ?? []) this.frame({ subscribed: channel });
break;
case 'get_user_orders':
later(options.readDelayMs, () => {
const list = [...orders.values()];
const offset = Number(params.offset),
limit = Number(params.limit);
this.frame({
all_orders_response: {
orders: list.slice(offset, offset + limit).map(wireOrder),
twaps: [],
stop_orders: [],
},
});
});
break;
case 'get_user_leverage':
this.frame({ user_leverage_response: [{ symbol: 'BTC-USD', leverage: '8' }] });
break;
case 'get_available_leverage_levels':
this.frame({ available_leverage_levels_response: [{ symbol: 'BTC-USD', leverage: '8' }] });
break;
case 'add_order': {
const row = fixtureOrder({
orderId: '20000000-0000-4000-8000-' + String(++sequence).padStart(12, '0'),
symbol: String(params.symbol),
side: params.side as 'BUY' | 'SELL',
qty: String(params.quantity),
remaining: String(params.quantity),
price: String(params.price),
cloid: String(params.client_order_id),
tif: String(params.order_time_in_force),
reduceOnly: params.reduce_only === true,
});
orders.set(row.orderId, row);
if (!options.dropAddReplies)
later(options.ackDelayMs, () => {
this.frame({ order_response: wireOrder(row) });
if (options.rejectAfterAck) {
orders.delete(row.orderId);
later(options.ackDelayMs, () =>
this.frame({ order_response: wireOrder({ ...row, status: options.rejectAfterAck! }) }),
);
}
});
break;
}
case 'cancel_order': {
const id = String(params.order_id);
const row = orders.get(id);
if (row) {
orders.delete(id);
this.frame({ order_response: wireOrder({ ...row, status: 'CANCELLED' }) });
} else this.frame({ err: { error_code: 'InvalidOrderId', incoming_message: request } });
break;
}
default:
this.frame({ err: { error_code: 'InvalidParameter', incoming_message: request } });
}
callback?.();
}
ping() {
queueMicrotask(() => this.emit('pong'));
}
close() {
if (this.readyState === WebSocket.CLOSED) return;
this.readyState = WebSocket.CLOSED;
sockets.delete(this);
queueMicrotask(() => this.emit('close', 1000, Buffer.alloc(0)));
}
terminate() {
this.close();
}
}
return {
webSocketFactory: () => new FakeSocket() as unknown as WebSocket,
orders,
sent,
emit: (value: unknown) => {
for (const socket of sockets) socket.frame(value);
},
drop: () => {
for (const socket of [...sockets]) socket.close();
},
close() {
closing = true;
for (const timer of timers) clearTimeout(timer);
timers.clear();
for (const socket of [...sockets]) socket.close();
},
};
}
export { FIXTURE_ACCOUNT };