Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
68 changes: 57 additions & 11 deletions src/handlers/subscribe/handler.ts
Original file line number Diff line number Diff line change
@@ -1,15 +1,15 @@
import * as grpc from "@grpc/grpc-js";
import type { Exchange } from "@usherlabs/ccxt";
import { authenticateRequest } from "../../helpers/auth";
import {
normalizeBinanceExecutionReport,
normalizeBinanceSpotBalanceEvent,
} from "../../helpers/binance-user-data-normalization";
import {
BinanceSpotUserDataStream,
isBinanceBalanceUserDataEvent,
isBinanceOrderUserDataEvent,
} from "../../helpers/binance-user-data-stream";
import {
normalizeBinanceExecutionReport,
normalizeBinanceSpotBalanceEvent,
} from "../../helpers/binance-user-data-normalization";
import {
type BrokerPoolEntry,
createBroker,
Expand Down Expand Up @@ -46,6 +46,10 @@ import {
} from "../../helpers/order-book";
import type { OtelMetrics } from "../../helpers/otel";
import { getErrorMessage } from "../../helpers/shared/errors";
import type {
UserDataStreamSupervisor,
UserDataSubscription,
} from "../../helpers/user-data-stream-supervisor";
import type { SubscribeRequest, SubscribeResponse } from "../types";
import { SubscribeBrokerLifecycle } from "./broker-lifecycle";

Expand All @@ -55,6 +59,7 @@ export type SubscribeDeps = {
otelMetrics?: OtelMetrics;
brokerArchiver?: BrokerExecutionArchiver;
brokerLifecycle?: SubscribeBrokerLifecycle;
userDataStreamSupervisor?: UserDataStreamSupervisor;
};

type SubscribeCall = grpc.ServerWritableStream<
Expand Down Expand Up @@ -179,16 +184,22 @@ async function streamBinanceUserData(
deploymentId: string;
assetType: BrokerMarketType;
},
userDataSource?: UserDataSubscription,
knownMarketId?: string | null,
): Promise<void> {
const userDataStream = new BinanceSpotUserDataStream(broker);
call.once("close", () => userDataStream.close());
call.once("cancelled", () => userDataStream.close());
call.once("error", () => userDataStream.close());
const userDataStream =
userDataSource ?? new BinanceSpotUserDataStream(broker);
if (!userDataSource) {
call.once("close", () => userDataStream.close());
call.once("cancelled", () => userDataStream.close());
call.once("error", () => userDataStream.close());
}

const marketId =
subscriptionType === SubscriptionType.ORDERS
knownMarketId ??
(subscriptionType === SubscriptionType.ORDERS
? await getBinanceMarketId(broker, symbol)
: null;
: null);

try {
for await (const message of userDataStream) {
Expand Down Expand Up @@ -319,7 +330,13 @@ async function runCcxtSubscribeLoop(
}

export function createSubscribeHandler(deps: SubscribeDeps) {
const { brokers, whitelistIps, otelMetrics, brokerArchiver } = deps;
const {
brokers,
whitelistIps,
otelMetrics,
brokerArchiver,
userDataStreamSupervisor,
} = deps;
const brokerLifecycle =
deps.brokerLifecycle ?? new SubscribeBrokerLifecycle();
return async (call: SubscribeCall) => {
Expand Down Expand Up @@ -485,13 +502,42 @@ export function createSubscribeHandler(deps: SubscribeDeps) {
});
return;
}
const marketId =
subscriptionType === SubscriptionType.ORDERS
? await getBinanceMarketId(accountBroker, resolvedSymbol)
: undefined;
let userDataSource: UserDataSubscription | undefined;
if (selectedBrokerAccount) {
if (!userDataStreamSupervisor) {
await writeSubscribeError(call, isStreamClosed, {
data: JSON.stringify({
error: "Configured account user-data supervisor is unavailable",
}),
timestamp: Date.now(),
symbol: resolvedSymbol,
type: subscriptionType,
});
return;
}
userDataSource = userDataStreamSupervisor.subscribe({
exchange: normalizedCex,
accountSelector: selectedBrokerAccount.label,
kind:
subscriptionType === SubscriptionType.BALANCE
? "balance"
: "orders",
marketId,
});
}
await streamBinanceUserData(
call,
accountBroker,
resolvedSymbol,
subscriptionType,
isStreamClosed,
streamArchiveContext,
userDataSource,
marketId,
);
return;
}
Expand Down
40 changes: 37 additions & 3 deletions src/helpers/binance-user-data-stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,23 @@ export type BinanceUserDataEvent = {
event: Record<string, unknown>;
};

export type BinanceUserDataStreamFailureKind =
| "auth_failed"
| "transport_error"
| "remote_closed"
| "protocol_error"
| "backpressure";

export type BinanceUserDataStreamObserver = {
onConnected?: () => void;
onAuthenticated?: () => void;
onEvent?: (event: BinanceUserDataEvent) => void;
onFailure?: (failure: {
kind: BinanceUserDataStreamFailureKind;
reason: string;
}) => void;
};

type BinanceUserDataMessage =
| {
id?: string | null;
Expand Down Expand Up @@ -39,6 +56,7 @@ type WebSocketFactory = (url: string) => WebSocketLike;

type BinanceSpotUserDataStreamOptions = {
maxBufferedEvents?: number;
observer?: BinanceUserDataStreamObserver;
};

let createWebSocket: WebSocketFactory = (url) =>
Expand Down Expand Up @@ -240,6 +258,7 @@ export class BinanceSpotUserDataStream
private readonly requestId =
`user-data-${Date.now()}-${userDataRequestCounter++}`;
private readonly maxBufferedEvents: number;
private readonly observer?: BinanceUserDataStreamObserver;
private readonly queue: BinanceUserDataEvent[] = [];
private readonly waiters: Array<{
resolve: (event: BinanceUserDataEvent | null) => void;
Expand All @@ -256,15 +275,22 @@ export class BinanceSpotUserDataStream
this.maxBufferedEvents =
options.maxBufferedEvents ??
DEFAULT_BINANCE_USER_DATA_MAX_BUFFERED_EVENTS;
this.observer = options.observer;
this.secretValues = [
getOptionalExchangeString(exchange, "apiKey"),
getOptionalExchangeString(exchange, "secret"),
].filter((value): value is string => value !== null);
this.ws = createWebSocket(getBinanceSpotWsApiUrl(exchange));
this.ws.on("open", () => this.subscribe());
this.ws.on("open", () => {
this.observer?.onConnected?.();
this.subscribe();
});
this.ws.on("message", (data) => this.handleMessage(data));
this.ws.on("error", (error) =>
this.fail(formatBinanceUserDataWebSocketError(error, this.secretValues)),
this.fail(
formatBinanceUserDataWebSocketError(error, this.secretValues),
"transport_error",
),
);
this.ws.on("close", (code, reason) => this.handleClose(code, reason));
}
Expand Down Expand Up @@ -299,6 +325,7 @@ export class BinanceSpotUserDataStream
}
this.fail(
formatBinanceUserDataWebSocketClose(code, reason, this.secretValues),
"remote_closed",
);
}

Expand Down Expand Up @@ -334,6 +361,7 @@ export class BinanceSpotUserDataStream
error instanceof Error
? error
: new Error("Invalid Binance user-data message"),
"protocol_error",
);
return;
}
Expand All @@ -346,10 +374,12 @@ export class BinanceSpotUserDataStream
message.error?.message ??
`Binance user-data subscription failed with status ${message.status}`,
),
"auth_failed",
);
return;
}
this.subscriptionId = message.result?.subscriptionId ?? null;
this.observer?.onAuthenticated?.();
return;
}

Expand All @@ -369,6 +399,7 @@ export class BinanceSpotUserDataStream
? `${errorMessage} (code ${errorCode})`
: errorMessage,
),
"protocol_error",
);
return;
}
Expand All @@ -387,6 +418,7 @@ export class BinanceSpotUserDataStream
if (this.closed) {
return;
}
this.observer?.onEvent?.(event);

const waiter = this.waiters.shift();
if (waiter) {
Expand All @@ -398,6 +430,7 @@ export class BinanceSpotUserDataStream
new Error(
`Binance user-data stream buffered event limit exceeded (${this.maxBufferedEvents}); downstream consumer is not keeping up`,
),
"backpressure",
);
return;
}
Expand All @@ -420,11 +453,12 @@ export class BinanceSpotUserDataStream
});
}

private fail(error: Error): void {
private fail(error: Error, kind: BinanceUserDataStreamFailureKind): void {
if (this.closeError) {
return;
}
this.closeError = error;
this.observer?.onFailure?.({ kind, reason: error.message });
this.closed = true;
this.queue.length = 0;
this.flushWaiters();
Expand Down
Loading
Loading