Skip to content
markpaper

src/testing/fake.ts

v0.1.0 · 5.3 KB

Download file
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 };
All files