Skip to content

Commit e761af6

Browse files
committed
feat: supervise configured user data streams
1 parent 9747687 commit e761af6

8 files changed

Lines changed: 1276 additions & 76 deletions

src/handlers/subscribe/handler.ts

Lines changed: 57 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,15 @@
11
import * as grpc from "@grpc/grpc-js";
22
import type { Exchange } from "@usherlabs/ccxt";
33
import { authenticateRequest } from "../../helpers/auth";
4+
import {
5+
normalizeBinanceExecutionReport,
6+
normalizeBinanceSpotBalanceEvent,
7+
} from "../../helpers/binance-user-data-normalization";
48
import {
59
BinanceSpotUserDataStream,
610
isBinanceBalanceUserDataEvent,
711
isBinanceOrderUserDataEvent,
812
} from "../../helpers/binance-user-data-stream";
9-
import {
10-
normalizeBinanceExecutionReport,
11-
normalizeBinanceSpotBalanceEvent,
12-
} from "../../helpers/binance-user-data-normalization";
1313
import {
1414
type BrokerPoolEntry,
1515
createBroker,
@@ -46,6 +46,10 @@ import {
4646
} from "../../helpers/order-book";
4747
import type { OtelMetrics } from "../../helpers/otel";
4848
import { getErrorMessage } from "../../helpers/shared/errors";
49+
import type {
50+
UserDataStreamSupervisor,
51+
UserDataSubscription,
52+
} from "../../helpers/user-data-stream-supervisor";
4953
import type { SubscribeRequest, SubscribeResponse } from "../types";
5054
import { SubscribeBrokerLifecycle } from "./broker-lifecycle";
5155

@@ -55,6 +59,7 @@ export type SubscribeDeps = {
5559
otelMetrics?: OtelMetrics;
5660
brokerArchiver?: BrokerExecutionArchiver;
5761
brokerLifecycle?: SubscribeBrokerLifecycle;
62+
userDataStreamSupervisor?: UserDataStreamSupervisor;
5863
};
5964

6065
type SubscribeCall = grpc.ServerWritableStream<
@@ -179,16 +184,22 @@ async function streamBinanceUserData(
179184
deploymentId: string;
180185
assetType: BrokerMarketType;
181186
},
187+
userDataSource?: UserDataSubscription,
188+
knownMarketId?: string | null,
182189
): Promise<void> {
183-
const userDataStream = new BinanceSpotUserDataStream(broker);
184-
call.once("close", () => userDataStream.close());
185-
call.once("cancelled", () => userDataStream.close());
186-
call.once("error", () => userDataStream.close());
190+
const userDataStream =
191+
userDataSource ?? new BinanceSpotUserDataStream(broker);
192+
if (!userDataSource) {
193+
call.once("close", () => userDataStream.close());
194+
call.once("cancelled", () => userDataStream.close());
195+
call.once("error", () => userDataStream.close());
196+
}
187197

188198
const marketId =
189-
subscriptionType === SubscriptionType.ORDERS
199+
knownMarketId ??
200+
(subscriptionType === SubscriptionType.ORDERS
190201
? await getBinanceMarketId(broker, symbol)
191-
: null;
202+
: null);
192203

193204
try {
194205
for await (const message of userDataStream) {
@@ -319,7 +330,13 @@ async function runCcxtSubscribeLoop(
319330
}
320331

321332
export function createSubscribeHandler(deps: SubscribeDeps) {
322-
const { brokers, whitelistIps, otelMetrics, brokerArchiver } = deps;
333+
const {
334+
brokers,
335+
whitelistIps,
336+
otelMetrics,
337+
brokerArchiver,
338+
userDataStreamSupervisor,
339+
} = deps;
323340
const brokerLifecycle =
324341
deps.brokerLifecycle ?? new SubscribeBrokerLifecycle();
325342
return async (call: SubscribeCall) => {
@@ -485,13 +502,42 @@ export function createSubscribeHandler(deps: SubscribeDeps) {
485502
});
486503
return;
487504
}
505+
const marketId =
506+
subscriptionType === SubscriptionType.ORDERS
507+
? await getBinanceMarketId(accountBroker, resolvedSymbol)
508+
: undefined;
509+
let userDataSource: UserDataSubscription | undefined;
510+
if (selectedBrokerAccount) {
511+
if (!userDataStreamSupervisor) {
512+
await writeSubscribeError(call, isStreamClosed, {
513+
data: JSON.stringify({
514+
error: "Configured account user-data supervisor is unavailable",
515+
}),
516+
timestamp: Date.now(),
517+
symbol: resolvedSymbol,
518+
type: subscriptionType,
519+
});
520+
return;
521+
}
522+
userDataSource = userDataStreamSupervisor.subscribe({
523+
exchange: normalizedCex,
524+
accountSelector: selectedBrokerAccount.label,
525+
kind:
526+
subscriptionType === SubscriptionType.BALANCE
527+
? "balance"
528+
: "orders",
529+
marketId,
530+
});
531+
}
488532
await streamBinanceUserData(
489533
call,
490534
accountBroker,
491535
resolvedSymbol,
492536
subscriptionType,
493537
isStreamClosed,
494538
streamArchiveContext,
539+
userDataSource,
540+
marketId,
495541
);
496542
return;
497543
}

src/helpers/binance-user-data-stream.ts

Lines changed: 37 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,23 @@ export type BinanceUserDataEvent = {
1111
event: Record<string, unknown>;
1212
};
1313

14+
export type BinanceUserDataStreamFailureKind =
15+
| "auth_failed"
16+
| "transport_error"
17+
| "remote_closed"
18+
| "protocol_error"
19+
| "backpressure";
20+
21+
export type BinanceUserDataStreamObserver = {
22+
onConnected?: () => void;
23+
onAuthenticated?: () => void;
24+
onEvent?: (event: BinanceUserDataEvent) => void;
25+
onFailure?: (failure: {
26+
kind: BinanceUserDataStreamFailureKind;
27+
reason: string;
28+
}) => void;
29+
};
30+
1431
type BinanceUserDataMessage =
1532
| {
1633
id?: string | null;
@@ -39,6 +56,7 @@ type WebSocketFactory = (url: string) => WebSocketLike;
3956

4057
type BinanceSpotUserDataStreamOptions = {
4158
maxBufferedEvents?: number;
59+
observer?: BinanceUserDataStreamObserver;
4260
};
4361

4462
let createWebSocket: WebSocketFactory = (url) =>
@@ -240,6 +258,7 @@ export class BinanceSpotUserDataStream
240258
private readonly requestId =
241259
`user-data-${Date.now()}-${userDataRequestCounter++}`;
242260
private readonly maxBufferedEvents: number;
261+
private readonly observer?: BinanceUserDataStreamObserver;
243262
private readonly queue: BinanceUserDataEvent[] = [];
244263
private readonly waiters: Array<{
245264
resolve: (event: BinanceUserDataEvent | null) => void;
@@ -256,15 +275,22 @@ export class BinanceSpotUserDataStream
256275
this.maxBufferedEvents =
257276
options.maxBufferedEvents ??
258277
DEFAULT_BINANCE_USER_DATA_MAX_BUFFERED_EVENTS;
278+
this.observer = options.observer;
259279
this.secretValues = [
260280
getOptionalExchangeString(exchange, "apiKey"),
261281
getOptionalExchangeString(exchange, "secret"),
262282
].filter((value): value is string => value !== null);
263283
this.ws = createWebSocket(getBinanceSpotWsApiUrl(exchange));
264-
this.ws.on("open", () => this.subscribe());
284+
this.ws.on("open", () => {
285+
this.observer?.onConnected?.();
286+
this.subscribe();
287+
});
265288
this.ws.on("message", (data) => this.handleMessage(data));
266289
this.ws.on("error", (error) =>
267-
this.fail(formatBinanceUserDataWebSocketError(error, this.secretValues)),
290+
this.fail(
291+
formatBinanceUserDataWebSocketError(error, this.secretValues),
292+
"transport_error",
293+
),
268294
);
269295
this.ws.on("close", (code, reason) => this.handleClose(code, reason));
270296
}
@@ -299,6 +325,7 @@ export class BinanceSpotUserDataStream
299325
}
300326
this.fail(
301327
formatBinanceUserDataWebSocketClose(code, reason, this.secretValues),
328+
"remote_closed",
302329
);
303330
}
304331

@@ -334,6 +361,7 @@ export class BinanceSpotUserDataStream
334361
error instanceof Error
335362
? error
336363
: new Error("Invalid Binance user-data message"),
364+
"protocol_error",
337365
);
338366
return;
339367
}
@@ -346,10 +374,12 @@ export class BinanceSpotUserDataStream
346374
message.error?.message ??
347375
`Binance user-data subscription failed with status ${message.status}`,
348376
),
377+
"auth_failed",
349378
);
350379
return;
351380
}
352381
this.subscriptionId = message.result?.subscriptionId ?? null;
382+
this.observer?.onAuthenticated?.();
353383
return;
354384
}
355385

@@ -369,6 +399,7 @@ export class BinanceSpotUserDataStream
369399
? `${errorMessage} (code ${errorCode})`
370400
: errorMessage,
371401
),
402+
"protocol_error",
372403
);
373404
return;
374405
}
@@ -387,6 +418,7 @@ export class BinanceSpotUserDataStream
387418
if (this.closed) {
388419
return;
389420
}
421+
this.observer?.onEvent?.(event);
390422

391423
const waiter = this.waiters.shift();
392424
if (waiter) {
@@ -398,6 +430,7 @@ export class BinanceSpotUserDataStream
398430
new Error(
399431
`Binance user-data stream buffered event limit exceeded (${this.maxBufferedEvents}); downstream consumer is not keeping up`,
400432
),
433+
"backpressure",
401434
);
402435
return;
403436
}
@@ -420,11 +453,12 @@ export class BinanceSpotUserDataStream
420453
});
421454
}
422455

423-
private fail(error: Error): void {
456+
private fail(error: Error, kind: BinanceUserDataStreamFailureKind): void {
424457
if (this.closeError) {
425458
return;
426459
}
427460
this.closeError = error;
461+
this.observer?.onFailure?.({ kind, reason: error.message });
428462
this.closed = true;
429463
this.queue.length = 0;
430464
this.flushWaiters();

0 commit comments

Comments
 (0)