Author SHA1 Message Date
discountry 2a1cf24510 Enhance Binance depth health monitoring and defense mode logic
- Updated translations for Binance depth status messages to include depth window information.
- Modified `MakerPointsEngine` to incorporate health checks for the Binance depth tracker, including handling of unhealthy states.
- Improved defense mode activation logic to respond to Binance depth health status, ensuring appropriate logging and notifications.
- Added integration tests for defense mode behavior based on Binance depth health, validating transitions into and out of defense mode.
- Refactored `BinanceDepthTracker` to support health checks and improved connection management.
2026-02-07 21:54:22 +08:00
5 changed files with 623 additions and 108 deletions
+2 -2
View File
@@ -248,8 +248,8 @@ const translations: Record<string, TranslationEntry> = {
en: "Quote mode: {mode} | BUY {buy} | SELL {sell}", en: "Quote mode: {mode} | BUY {buy} | SELL {sell}",
}, },
"makerPoints.binanceLine": { "makerPoints.binanceLine": {
zh: "Binance 深度: 买10 {buy} 10 {sell} 状态: {status}", zh: "Binance 深度(±9bps): 买 {buy} 卖 {sell} 状态: {status}",
en: "Binance depth: bid10 {buy} | ask10 {sell} | Status: {status}", en: "Binance depth (±9bps): bid {buy} | ask {sell} | Status: {status}",
}, },
"makerPoints.bandDepthLine": { "makerPoints.bandDepthLine": {
zh: "StandX 档位 {band}bps 深度: 买 {buy} 卖 {sell}", zh: "StandX 档位 {band}bps 深度: 买 {buy} 卖 {sell}",
+428 -104
View File
@@ -1,28 +1,29 @@
import NodeWebSocket from "ws"; import NodeWebSocket from "ws";
import { computeDepthStats, type DepthImbalance } from "../../utils/depth"; import type { AsterDepthLevel } from "../../exchanges/types";
import type { DepthImbalance } from "../../utils/depth";
const WebSocketCtor: typeof globalThis.WebSocket = const WebSocketCtor: typeof globalThis.WebSocket =
typeof globalThis.WebSocket !== "undefined" typeof globalThis.WebSocket !== "undefined"
? globalThis.WebSocket ? globalThis.WebSocket
: ((NodeWebSocket as unknown) as typeof globalThis.WebSocket); : ((NodeWebSocket as unknown) as typeof globalThis.WebSocket);
const DEFAULT_BASE_URL = "wss://stream.binance.com:9443/ws"; const DEFAULT_WS_BASE_URL = "wss://stream.binance.com:9443/ws";
const DEFAULT_REST_BASE_URL = "https://api.binance.com";
// ========== Binance WebSocket 连接管理常量 ==========
// Binance 会发送 ping,若长时间无消息则认为连接异常
// 我们设置 5 分钟作为心跳超时阈值(保守值)
const HEARTBEAT_TIMEOUT_MS = 5 * 60 * 1000; const HEARTBEAT_TIMEOUT_MS = 5 * 60 * 1000;
// 心跳检查间隔(每 30 秒检查一次)
const HEARTBEAT_CHECK_INTERVAL_MS = 30_000; const HEARTBEAT_CHECK_INTERVAL_MS = 30_000;
// Binance 连接最长有效期 24 小时,我们设置 23 小时主动重连
const MAX_CONNECTION_DURATION_MS = 23 * 60 * 60 * 1000; const MAX_CONNECTION_DURATION_MS = 23 * 60 * 60 * 1000;
// 数据过时阈值(毫秒)- 超过此时间未收到数据,标记为不可用
const DATA_STALE_THRESHOLD_MS = 5_000; const DATA_STALE_THRESHOLD_MS = 5_000;
// 基础重连延迟
const RECONNECT_DELAY_BASE_MS = 3000; const RECONNECT_DELAY_BASE_MS = 3000;
// 最大重连延迟
const RECONNECT_DELAY_MAX_MS = 60_000; const RECONNECT_DELAY_MAX_MS = 60_000;
const DEFAULT_REFRESH_SYNC_INTERVAL_MS = 30_000;
const DEFAULT_DEPTH_WINDOW_BPS = 9;
const DEFAULT_IMBALANCE_RATIO = 9;
const MAX_BUFFER_SIZE = 5000;
const SYNC_SNAPSHOT_MAX_RETRIES = 5;
const REST_FAILURE_DEFENSE_THRESHOLD = 1;
export type BinanceConnectionState = "connected" | "disconnected" | "stale"; export type BinanceConnectionState = "connected" | "disconnected" | "stale";
export interface BinanceDepthSnapshot { export interface BinanceDepthSnapshot {
@@ -33,49 +34,91 @@ export interface BinanceDepthSnapshot {
skipSellSide: boolean; skipSellSide: boolean;
imbalance: DepthImbalance; imbalance: DepthImbalance;
updatedAt: number; updatedAt: number;
windowBps: number;
localLastUpdateId: number;
}
export interface BinanceDepthHealth {
started: boolean;
connected: boolean;
orderBookReady: boolean;
restHealthy: boolean;
healthy: boolean;
reason: string | null;
lastEventAt: number;
lastSnapshotAt: number;
lastRestSyncAt: number;
localLastUpdateId: number;
} }
export type BinanceConnectionListener = (state: BinanceConnectionState) => void; export type BinanceConnectionListener = (state: BinanceConnectionState) => void;
interface DepthUpdateEvent {
U: number;
u: number;
bids: AsterDepthLevel[];
asks: AsterDepthLevel[];
}
interface DepthSnapshotResponse {
lastUpdateId: number;
bids: AsterDepthLevel[];
asks: AsterDepthLevel[];
}
export class BinanceDepthTracker { export class BinanceDepthTracker {
private ws: WebSocket | null = null; private ws: WebSocket | null = null;
private reconnectTimer: ReturnType<typeof setTimeout> | null = null; private reconnectTimer: ReturnType<typeof setTimeout> | null = null;
private reconnectDelayMs = RECONNECT_DELAY_BASE_MS; private reconnectDelayMs = RECONNECT_DELAY_BASE_MS;
private stopped = false; private stopped = false;
private started = false;
private snapshot: BinanceDepthSnapshot | null = null; private snapshot: BinanceDepthSnapshot | null = null;
private listeners = new Set<(snapshot: BinanceDepthSnapshot) => void>(); private listeners = new Set<(snapshot: BinanceDepthSnapshot) => void>();
private connectionListeners = new Set<BinanceConnectionListener>(); private connectionListeners = new Set<BinanceConnectionListener>();
// ========== 心跳与连接管理 ==========
// 上次收到消息的时间戳
private lastMessageTime = 0; private lastMessageTime = 0;
// 心跳检查定时器
private heartbeatTimer: ReturnType<typeof setInterval> | null = null; private heartbeatTimer: ReturnType<typeof setInterval> | null = null;
// 连接建立时间(用于日志记录)
// eslint-disable-next-line @typescript-eslint/no-unused-vars
private connectionStartTime = 0;
// 24 小时重连定时器
private maxDurationTimer: ReturnType<typeof setTimeout> | null = null; private maxDurationTimer: ReturnType<typeof setTimeout> | null = null;
// 当前连接状态 private refreshSyncTimer: ReturnType<typeof setInterval> | null = null;
private connectionState: BinanceConnectionState = "disconnected"; private connectionState: BinanceConnectionState = "disconnected";
private bidBook = new Map<string, number>();
private askBook = new Map<string, number>();
private localLastUpdateId = 0;
private orderBookReady = false;
private eventBuffer: DepthUpdateEvent[] = [];
private syncInFlight: Promise<void> | null = null;
private lastEventAt = 0;
private lastSnapshotAt = 0;
private lastRestSyncAt = 0;
private restConsecutiveFailures = 0;
private restLastError: string | null = null;
constructor( constructor(
private readonly symbol: string, private readonly symbol: string,
private readonly options?: { private readonly options?: {
baseUrl?: string; baseUrl?: string;
restBaseUrl?: string;
levels?: number; levels?: number;
ratio?: number; ratio?: number;
speedMs?: number; speedMs?: number;
depthWindowBps?: number;
refreshSyncMs?: number;
logger?: (context: string, error: unknown) => void; logger?: (context: string, error: unknown) => void;
} }
) {} ) {}
start(): void { start(): void {
this.started = true;
this.stopped = false; this.stopped = false;
this.connect(); this.connect();
this.startRefreshSyncTimer();
} }
stop(): void { stop(): void {
this.started = false;
this.stopped = true; this.stopped = true;
this.cleanup(); this.cleanup();
} }
@@ -88,9 +131,6 @@ export class BinanceDepthTracker {
this.listeners.delete(handler); this.listeners.delete(handler);
} }
/**
* 监听连接状态变化
*/
onConnectionChange(handler: BinanceConnectionListener): void { onConnectionChange(handler: BinanceConnectionListener): void {
this.connectionListeners.add(handler); this.connectionListeners.add(handler);
} }
@@ -103,69 +143,111 @@ export class BinanceDepthTracker {
return this.snapshot ? { ...this.snapshot } : null; return this.snapshot ? { ...this.snapshot } : null;
} }
/**
* 获取当前连接状态
*/
getConnectionState(): BinanceConnectionState { getConnectionState(): BinanceConnectionState {
return this.connectionState; return this.connectionState;
} }
/**
* 检查数据是否过时
*/
isDataStale(): boolean { isDataStale(): boolean {
if (!this.snapshot) return true; if (!this.snapshot) return true;
return Date.now() - this.snapshot.updatedAt > DATA_STALE_THRESHOLD_MS; return Date.now() - this.snapshot.updatedAt > DATA_STALE_THRESHOLD_MS;
} }
isHealthy(): boolean {
return this.getHealth().healthy;
}
getHealth(): BinanceDepthHealth {
if (!this.started) {
return {
started: false,
connected: false,
orderBookReady: false,
restHealthy: true,
healthy: true,
reason: null,
lastEventAt: this.lastEventAt,
lastSnapshotAt: this.lastSnapshotAt,
lastRestSyncAt: this.lastRestSyncAt,
localLastUpdateId: this.localLastUpdateId,
};
}
const restHealthy = this.restConsecutiveFailures < REST_FAILURE_DEFENSE_THRESHOLD;
let reason: string | null = null;
if (this.connectionState !== "connected") {
reason = `ws_${this.connectionState}`;
} else if (!this.orderBookReady) {
reason = "orderbook_not_ready";
} else if (this.isDataStale()) {
reason = "orderbook_stale";
} else if (!restHealthy) {
reason = this.restLastError ? `rest_sync_failed:${this.restLastError}` : "rest_sync_failed";
}
return {
started: true,
connected: this.connectionState === "connected",
orderBookReady: this.orderBookReady,
restHealthy,
healthy: reason == null,
reason,
lastEventAt: this.lastEventAt,
lastSnapshotAt: this.lastSnapshotAt,
lastRestSyncAt: this.lastRestSyncAt,
localLastUpdateId: this.localLastUpdateId,
};
}
private cleanup(): void { private cleanup(): void {
// 停止心跳监控
if (this.heartbeatTimer) { if (this.heartbeatTimer) {
clearInterval(this.heartbeatTimer); clearInterval(this.heartbeatTimer);
this.heartbeatTimer = null; this.heartbeatTimer = null;
} }
// 停止 24 小时重连定时器
if (this.maxDurationTimer) { if (this.maxDurationTimer) {
clearTimeout(this.maxDurationTimer); clearTimeout(this.maxDurationTimer);
this.maxDurationTimer = null; this.maxDurationTimer = null;
} }
// 停止重连定时器 if (this.refreshSyncTimer) {
clearInterval(this.refreshSyncTimer);
this.refreshSyncTimer = null;
}
if (this.reconnectTimer) { if (this.reconnectTimer) {
clearTimeout(this.reconnectTimer); clearTimeout(this.reconnectTimer);
this.reconnectTimer = null; this.reconnectTimer = null;
} }
// 关闭 WebSocket
if (this.ws) { if (this.ws) {
try { try {
this.ws.close(); this.ws.close();
} catch { } catch {
// Ignore close errors // ignore close errors
} }
this.ws = null; this.ws = null;
} }
this.updateConnectionState("disconnected");
} }
private connect(): void { private connect(): void {
if (this.ws || this.stopped) return; if (this.ws || this.stopped) return;
const url = this.buildUrl(); const url = this.buildWsUrl();
this.ws = new WebSocketCtor(url); this.ws = new WebSocketCtor(url);
const handleOpen = () => { const handleOpen = () => {
this.reconnectDelayMs = RECONNECT_DELAY_BASE_MS; this.reconnectDelayMs = RECONNECT_DELAY_BASE_MS;
this.connectionStartTime = Date.now();
this.lastMessageTime = Date.now(); this.lastMessageTime = Date.now();
this.lastEventAt = Date.now();
this.orderBookReady = false;
this.eventBuffer = [];
this.updateConnectionState("connected"); this.updateConnectionState("connected");
// 启动心跳监控
this.startHeartbeatMonitor(); this.startHeartbeatMonitor();
// 启动 24 小时自动重连定时器
this.startMaxDurationTimer(); this.startMaxDurationTimer();
this.options?.logger?.("binanceDepth", "WebSocket connected"); this.options?.logger?.("binanceDepth", "WebSocket connected");
}; };
const handleClose = () => { const handleClose = () => {
this.ws = null; this.ws = null;
this.orderBookReady = false;
this.eventBuffer = [];
this.stopHeartbeatMonitor(); this.stopHeartbeatMonitor();
this.stopMaxDurationTimer(); this.stopMaxDurationTimer();
this.updateConnectionState("disconnected"); this.updateConnectionState("disconnected");
@@ -178,7 +260,6 @@ export class BinanceDepthTracker {
const handleError = (error: unknown) => { const handleError = (error: unknown) => {
this.options?.logger?.("binanceDepth", error); this.options?.logger?.("binanceDepth", error);
// 如果连接从未成功建立,需要清理并重连
if (this.ws && this.connectionState === "disconnected") { if (this.ws && this.connectionState === "disconnected") {
this.ws = null; this.ws = null;
this.scheduleReconnect(); this.scheduleReconnect();
@@ -187,20 +268,18 @@ export class BinanceDepthTracker {
const handleMessage = (event: { data: unknown }) => { const handleMessage = (event: { data: unknown }) => {
this.lastMessageTime = Date.now(); this.lastMessageTime = Date.now();
// 如果之前是 stale 状态,恢复为 connected this.lastEventAt = Date.now();
if (this.connectionState === "stale") { if (this.connectionState === "stale") {
this.updateConnectionState("connected"); this.updateConnectionState("connected");
} }
this.handlePayload(event.data); this.handlePayload(event.data);
}; };
// 处理 Binance 服务器的 ping 帧
// 根据文档:必须尽快回复 pongpayload 为 ping 的 payload 副本
const handlePing = (data: unknown) => { const handlePing = (data: unknown) => {
this.lastMessageTime = Date.now(); this.lastMessageTime = Date.now();
if (this.ws && "pong" in this.ws && typeof this.ws.pong === "function") { if (this.ws && "pong" in this.ws && typeof this.ws.pong === "function") {
try { try {
this.ws.pong(data as any); this.ws.pong(data as never);
} catch (error) { } catch (error) {
this.options?.logger?.("binanceDepth pong", error); this.options?.logger?.("binanceDepth pong", error);
} }
@@ -209,33 +288,40 @@ export class BinanceDepthTracker {
if ("addEventListener" in this.ws && typeof this.ws.addEventListener === "function") { if ("addEventListener" in this.ws && typeof this.ws.addEventListener === "function") {
this.ws.addEventListener("open", handleOpen); this.ws.addEventListener("open", handleOpen);
this.ws.addEventListener("message", handleMessage as any); this.ws.addEventListener("message", handleMessage as never);
this.ws.addEventListener("close", handleClose); this.ws.addEventListener("close", handleClose);
this.ws.addEventListener("error", handleError as any); this.ws.addEventListener("error", handleError as never);
this.ws.addEventListener("ping", handlePing as any); this.ws.addEventListener("ping", handlePing as never);
} else if ("on" in this.ws && typeof (this.ws as any).on === "function") { } else if ("on" in this.ws && typeof (this.ws as { on?: unknown }).on === "function") {
const nodeSocket = this.ws as any; const nodeSocket = this.ws as { on: (event: string, listener: (...args: unknown[]) => void) => void };
nodeSocket.on("open", handleOpen); nodeSocket.on("open", handleOpen);
nodeSocket.on("message", (data: unknown) => handleMessage({ data })); nodeSocket.on("message", (data: unknown) => handleMessage({ data }));
nodeSocket.on("close", handleClose); nodeSocket.on("close", handleClose);
nodeSocket.on("error", handleError); nodeSocket.on("error", handleError);
nodeSocket.on("ping", handlePing); nodeSocket.on("ping", handlePing);
} else { } else {
(this.ws as any).onopen = handleOpen; const genericSocket = this.ws as any;
(this.ws as any).onmessage = handleMessage; genericSocket.onopen = handleOpen;
(this.ws as any).onclose = handleClose; genericSocket.onmessage = handleMessage;
(this.ws as any).onerror = handleError; genericSocket.onclose = handleClose;
genericSocket.onerror = handleError;
} }
} }
private buildUrl(): string { private buildWsUrl(): string {
const base = this.options?.baseUrl ?? DEFAULT_BASE_URL; const baseRaw = (this.options?.baseUrl ?? DEFAULT_WS_BASE_URL).replace(/\/+$/, "");
const levels = this.options?.levels ?? 10; const base = baseRaw.endsWith("/ws") || baseRaw.includes("/stream") ? baseRaw : `${baseRaw}/ws`;
const speed = this.options?.speedMs ?? 100; const speed = this.options?.speedMs ?? 100;
const stream = `${this.symbol.toLowerCase()}@depth${levels}@${speed}ms`; const stream = `${this.symbol.toLowerCase()}@depth@${speed}ms`;
return `${base}/${stream}`; return `${base}/${stream}`;
} }
private buildRestDepthUrl(): string {
const base = (this.options?.restBaseUrl ?? process.env.BINANCE_REST_URL ?? DEFAULT_REST_BASE_URL).replace(/\/+$/, "");
const symbol = this.symbol.toUpperCase();
return `${base}/api/v3/depth?symbol=${encodeURIComponent(symbol)}&limit=5000`;
}
private scheduleReconnect(): void { private scheduleReconnect(): void {
if (this.reconnectTimer || this.stopped) return; if (this.reconnectTimer || this.stopped) return;
this.reconnectTimer = setTimeout(() => { this.reconnectTimer = setTimeout(() => {
@@ -245,24 +331,17 @@ export class BinanceDepthTracker {
}, this.reconnectDelayMs); }, this.reconnectDelayMs);
} }
/**
* 启动心跳监控
* 根据 Binance 文档:长时间无 pong 会断连
* 我们设置 5 分钟作为心跳超时阈值
*/
private startHeartbeatMonitor(): void { private startHeartbeatMonitor(): void {
this.stopHeartbeatMonitor(); this.stopHeartbeatMonitor();
this.heartbeatTimer = setInterval(() => { this.heartbeatTimer = setInterval(() => {
const now = Date.now(); const now = Date.now();
const elapsed = now - this.lastMessageTime; const elapsed = now - this.lastMessageTime;
// 检查数据是否过时(5 秒无数据)
if (elapsed > DATA_STALE_THRESHOLD_MS && this.connectionState === "connected") { if (elapsed > DATA_STALE_THRESHOLD_MS && this.connectionState === "connected") {
this.updateConnectionState("stale"); this.updateConnectionState("stale");
this.options?.logger?.("binanceDepth", `Data stale: ${elapsed}ms since last message`); this.options?.logger?.("binanceDepth", `Data stale: ${elapsed}ms since last message`);
} }
// 检查心跳超时(5 分钟无消息)
if (elapsed > HEARTBEAT_TIMEOUT_MS) { if (elapsed > HEARTBEAT_TIMEOUT_MS) {
this.options?.logger?.("binanceDepth", `Heartbeat timeout: ${elapsed}ms, forcing reconnect`); this.options?.logger?.("binanceDepth", `Heartbeat timeout: ${elapsed}ms, forcing reconnect`);
this.forceReconnect("heartbeat_timeout"); this.forceReconnect("heartbeat_timeout");
@@ -277,11 +356,6 @@ export class BinanceDepthTracker {
} }
} }
/**
* 启动 24 小时自动重连定时器
* 根据 Binance 文档:连接最长有效期 24 小时
* 我们设置 23 小时主动重连,避免被服务器断开
*/
private startMaxDurationTimer(): void { private startMaxDurationTimer(): void {
this.stopMaxDurationTimer(); this.stopMaxDurationTimer();
this.maxDurationTimer = setTimeout(() => { this.maxDurationTimer = setTimeout(() => {
@@ -297,9 +371,15 @@ export class BinanceDepthTracker {
} }
} }
/** private startRefreshSyncTimer(): void {
* 强制重连 if (this.refreshSyncTimer) return;
*/ const refreshSyncMs = Math.max(5000, this.options?.refreshSyncMs ?? DEFAULT_REFRESH_SYNC_INTERVAL_MS);
this.refreshSyncTimer = setInterval(() => {
if (!this.started || this.stopped || !this.orderBookReady) return;
this.ensureSynced("periodic_refresh");
}, refreshSyncMs);
}
private forceReconnect(reason: string): void { private forceReconnect(reason: string): void {
this.options?.logger?.("binanceDepth", `Force reconnect: ${reason}`); this.options?.logger?.("binanceDepth", `Force reconnect: ${reason}`);
this.stopHeartbeatMonitor(); this.stopHeartbeatMonitor();
@@ -314,15 +394,13 @@ export class BinanceDepthTracker {
this.ws = null; this.ws = null;
} }
this.orderBookReady = false;
this.eventBuffer = [];
this.updateConnectionState("disconnected"); this.updateConnectionState("disconnected");
// 立即重连(不使用指数退避)
this.reconnectDelayMs = RECONNECT_DELAY_BASE_MS; this.reconnectDelayMs = RECONNECT_DELAY_BASE_MS;
this.scheduleReconnect(); this.scheduleReconnect();
} }
/**
* 更新连接状态并通知监听器
*/
private updateConnectionState(state: BinanceConnectionState): void { private updateConnectionState(state: BinanceConnectionState): void {
if (this.connectionState === state) return; if (this.connectionState === state) return;
this.connectionState = state; this.connectionState = state;
@@ -336,27 +414,197 @@ export class BinanceDepthTracker {
} }
private handlePayload(data: unknown): void { private handlePayload(data: unknown): void {
const payload = this.parsePayload(data); const event = this.parseDepthEvent(data);
if (!payload) return; if (!event) return;
const bids = Array.isArray(payload.b) ? payload.b : Array.isArray(payload.bids) ? payload.bids : [];
const asks = Array.isArray(payload.a) ? payload.a : Array.isArray(payload.asks) ? payload.asks : []; if (!this.orderBookReady) {
const depth = { this.eventBuffer.push(event);
lastUpdateId: Number(payload.lastUpdateId ?? payload.u ?? Date.now()), if (this.eventBuffer.length > MAX_BUFFER_SIZE) {
bids, this.eventBuffer.splice(0, this.eventBuffer.length - MAX_BUFFER_SIZE);
asks, }
}; this.ensureSynced("bootstrap");
const levels = this.options?.levels ?? 10; return;
const ratio = this.options?.ratio ?? 3; }
const stats = computeDepthStats(depth, levels, ratio);
const applied = this.applyDepthEvent(event);
if (!applied) {
this.options?.logger?.(
"binanceDepth",
`Detected update gap: local=${this.localLastUpdateId}, event=[${event.U},${event.u}], resyncing`
);
this.orderBookReady = false;
this.eventBuffer = [event];
this.ensureSynced("sequence_gap");
return;
}
this.emitDepthSnapshot();
}
private ensureSynced(reason: string): void {
if (this.syncInFlight || this.stopped || !this.started) return;
this.syncInFlight = (async () => {
try {
if (!this.orderBookReady) {
await this.bootstrapOrderBookFromSnapshot(reason);
return;
}
await this.refreshOrderBookFromSnapshot(reason);
} finally {
this.syncInFlight = null;
}
})();
}
private async bootstrapOrderBookFromSnapshot(reason: string): Promise<void> {
if (this.eventBuffer.length === 0) {
return;
}
for (let attempt = 0; attempt < SYNC_SNAPSHOT_MAX_RETRIES; attempt += 1) {
const firstBuffered = this.eventBuffer[0];
if (!firstBuffered) return;
const snapshot = await this.fetchDepthSnapshot(reason);
if (!snapshot) return;
if (snapshot.lastUpdateId < firstBuffered.U) {
continue;
}
this.resetOrderBook(snapshot);
const buffered = this.eventBuffer.filter((event) => event.u > snapshot.lastUpdateId);
if (buffered.length > 0) {
const nextEvent = buffered[0];
if (!nextEvent) return;
const nextUpdateId = snapshot.lastUpdateId + 1;
if (nextEvent.U > nextUpdateId || nextEvent.u < nextUpdateId) {
continue;
}
let failed = false;
for (const event of buffered) {
if (!this.applyDepthEvent(event)) {
failed = true;
break;
}
}
if (failed) {
continue;
}
}
this.orderBookReady = true;
this.eventBuffer = [];
this.emitDepthSnapshot();
return;
}
this.options?.logger?.("binanceDepth", "Bootstrap orderbook failed after retries");
}
private async refreshOrderBookFromSnapshot(reason: string): Promise<void> {
const snapshot = await this.fetchDepthSnapshot(reason);
if (!snapshot) return;
if (snapshot.lastUpdateId < this.localLastUpdateId) {
return;
}
this.resetOrderBook(snapshot);
this.orderBookReady = true;
this.emitDepthSnapshot();
}
private resetOrderBook(snapshot: DepthSnapshotResponse): void {
this.bidBook.clear();
this.askBook.clear();
this.applyLevels(this.bidBook, snapshot.bids);
this.applyLevels(this.askBook, snapshot.asks);
this.localLastUpdateId = snapshot.lastUpdateId;
this.lastSnapshotAt = Date.now();
}
private applyDepthEvent(event: DepthUpdateEvent): boolean {
if (!this.localLastUpdateId) return false;
if (event.u < this.localLastUpdateId) {
return true;
}
if (event.U > this.localLastUpdateId + 1) {
return false;
}
this.applyLevels(this.bidBook, event.bids);
this.applyLevels(this.askBook, event.asks);
this.localLastUpdateId = event.u;
return true;
}
private applyLevels(book: Map<string, number>, levels: AsterDepthLevel[]): void {
for (const level of levels) {
const priceRaw = level?.[0];
const qtyRaw = level?.[1];
const price = Number(priceRaw);
const qty = Number(qtyRaw);
if (!priceRaw || !Number.isFinite(price) || price <= 0) continue;
if (!Number.isFinite(qty) || qty < 0) continue;
if (qty === 0) {
book.delete(priceRaw);
} else {
book.set(priceRaw, qty);
}
}
}
private emitDepthSnapshot(): void {
const bestBid = this.findBestPrice(this.bidBook, "bid");
const bestAsk = this.findBestPrice(this.askBook, "ask");
if (bestBid == null || bestAsk == null || bestBid <= 0 || bestAsk <= 0 || bestAsk < bestBid) {
return;
}
const windowBps = Math.max(1, this.options?.depthWindowBps ?? DEFAULT_DEPTH_WINDOW_BPS);
const ratio = Math.max(1.01, this.options?.ratio ?? DEFAULT_IMBALANCE_RATIO);
const bidWindowMin = bestBid * (1 - windowBps / 10_000);
const askWindowMax = bestAsk * (1 + windowBps / 10_000);
let buySum = 0;
let sellSum = 0;
for (const [priceRaw, qty] of this.bidBook.entries()) {
const price = Number(priceRaw);
if (!Number.isFinite(price) || price < bidWindowMin) continue;
buySum += qty;
}
for (const [priceRaw, qty] of this.askBook.entries()) {
const price = Number(priceRaw);
if (!Number.isFinite(price) || price > askWindowMax) continue;
sellSum += qty;
}
const skipSellSide = sellSum === 0 || buySum > sellSum * ratio;
const skipBuySide = buySum === 0 || sellSum > buySum * ratio;
let imbalance: DepthImbalance = "balanced";
if (buySum > sellSum * ratio) {
imbalance = "buy_dominant";
} else if (sellSum > buySum * ratio) {
imbalance = "sell_dominant";
}
this.snapshot = { this.snapshot = {
symbol: this.symbol, symbol: this.symbol,
buySum: stats.buySum, buySum,
sellSum: stats.sellSum, sellSum,
skipBuySide: stats.skipBuySide, skipBuySide,
skipSellSide: stats.skipSellSide, skipSellSide,
imbalance: stats.imbalance, imbalance,
updatedAt: Date.now(), updatedAt: Date.now(),
windowBps,
localLastUpdateId: this.localLastUpdateId,
}; };
for (const listener of this.listeners) { for (const listener of this.listeners) {
try { try {
listener({ ...this.snapshot }); listener({ ...this.snapshot });
@@ -366,24 +614,100 @@ export class BinanceDepthTracker {
} }
} }
private parsePayload( private findBestPrice(book: Map<string, number>, side: "bid" | "ask"): number | null {
data: unknown let best: number | null = null;
): { b?: [string, string][]; a?: [string, string][]; bids?: [string, string][]; asks?: [string, string][]; u?: number; lastUpdateId?: number } | null {
for (const [priceRaw, qty] of book.entries()) {
if (!Number.isFinite(qty) || qty <= 0) continue;
const price = Number(priceRaw);
if (!Number.isFinite(price) || price <= 0) continue;
if (best == null) {
best = price;
continue;
}
if (side === "bid") {
if (price > best) best = price;
} else if (price < best) {
best = price;
}
}
return best;
}
private async fetchDepthSnapshot(reason: string): Promise<DepthSnapshotResponse | null> {
try {
const response = await fetch(this.buildRestDepthUrl(), {
method: "GET",
headers: { "content-type": "application/json" },
});
if (!response.ok) {
throw new Error(`HTTP ${response.status}`);
}
const json = (await response.json()) as {
lastUpdateId?: number;
bids?: Array<[string, string]>;
asks?: Array<[string, string]>;
};
const lastUpdateId = Number(json.lastUpdateId);
if (!Number.isFinite(lastUpdateId) || lastUpdateId <= 0) {
throw new Error("invalid lastUpdateId");
}
const bids = Array.isArray(json.bids) ? (json.bids as AsterDepthLevel[]) : [];
const asks = Array.isArray(json.asks) ? (json.asks as AsterDepthLevel[]) : [];
this.lastRestSyncAt = Date.now();
this.restConsecutiveFailures = 0;
this.restLastError = null;
return { lastUpdateId, bids, asks };
} catch (error) {
this.restConsecutiveFailures += 1;
this.restLastError = this.extractMessage(error);
this.options?.logger?.("binanceDepth", `REST sync failed (${reason}): ${this.restLastError}`);
return null;
}
}
private parseDepthEvent(data: unknown): DepthUpdateEvent | null {
try { try {
const text = typeof data === "string" ? data : Buffer.isBuffer(data) ? data.toString("utf-8") : null; const text = typeof data === "string" ? data : Buffer.isBuffer(data) ? data.toString("utf-8") : null;
if (!text) return null; if (!text) return null;
const parsed = JSON.parse(text); const parsed = JSON.parse(text) as unknown;
if (!parsed || typeof parsed !== "object") return null; if (!parsed || typeof parsed !== "object") return null;
return parsed as {
b?: [string, string][]; const maybeCombined = parsed as { data?: unknown };
a?: [string, string][]; const payload =
bids?: [string, string][]; maybeCombined.data && typeof maybeCombined.data === "object"
asks?: [string, string][]; ? (maybeCombined.data as Record<string, unknown>)
u?: number; : (parsed as Record<string, unknown>);
lastUpdateId?: number;
}; const eventType = typeof payload.e === "string" ? payload.e : "";
if (eventType && eventType !== "depthUpdate") {
return null;
}
const U = Number(payload.U);
const u = Number(payload.u);
if (!Number.isFinite(U) || !Number.isFinite(u)) {
return null;
}
const bidsRaw = Array.isArray(payload.b) ? payload.b : [];
const asksRaw = Array.isArray(payload.a) ? payload.a : [];
const bids = bidsRaw.filter((level): level is AsterDepthLevel => Array.isArray(level)) as AsterDepthLevel[];
const asks = asksRaw.filter((level): level is AsterDepthLevel => Array.isArray(level)) as AsterDepthLevel[];
return { U, u, bids, asks };
} catch { } catch {
return null; return null;
} }
} }
private extractMessage(error: unknown): string {
if (error instanceof Error) return error.message;
return String(error);
}
} }
+17 -2
View File
@@ -203,8 +203,10 @@ export class MakerPointsEngine {
this.qtyStep = Math.max(1e-9, this.config.qtyStep); this.qtyStep = Math.max(1e-9, this.config.qtyStep);
this.binanceDepth = new BinanceDepthTracker(resolveBinanceSymbol(this.config.symbol), { this.binanceDepth = new BinanceDepthTracker(resolveBinanceSymbol(this.config.symbol), {
baseUrl: process.env.BINANCE_SPOT_WS_URL ?? process.env.BINANCE_WS_URL, baseUrl: process.env.BINANCE_SPOT_WS_URL ?? process.env.BINANCE_WS_URL,
restBaseUrl: process.env.BINANCE_REST_URL,
levels: 20, levels: 20,
ratio: 9, ratio: 9,
depthWindowBps: 9,
speedMs: 100, speedMs: 100,
logger: (context, error) => { logger: (context, error) => {
this.tradeLog.push("warn", `Binance ${context} 异常: ${extractMessage(error)}`); this.tradeLog.push("warn", `Binance ${context} 异常: ${extractMessage(error)}`);
@@ -221,6 +223,7 @@ export class MakerPointsEngine {
this.feedStatus.binance = false; this.feedStatus.binance = false;
this.tradeLog.push("warn", "Binance 深度连接断开"); this.tradeLog.push("warn", "Binance 深度连接断开");
} else if (state === "stale") { } else if (state === "stale") {
this.feedStatus.binance = false;
this.tradeLog.push("warn", "Binance 深度数据过时"); this.tradeLog.push("warn", "Binance 深度数据过时");
} else if (state === "connected") { } else if (state === "connected") {
this.feedStatus.binance = true; this.feedStatus.binance = true;
@@ -1654,6 +1657,8 @@ export class MakerPointsEngine {
const now = Date.now(); const now = Date.now();
const standxDepthStale = this.lastStandxDepthTime > 0 && (now - this.lastStandxDepthTime) > DATA_STALE_THRESHOLD_MS; const standxDepthStale = this.lastStandxDepthTime > 0 && (now - this.lastStandxDepthTime) > DATA_STALE_THRESHOLD_MS;
const binanceStale = this.lastBinanceDepthTime > 0 && (now - this.lastBinanceDepthTime) > DATA_STALE_THRESHOLD_MS; const binanceStale = this.lastBinanceDepthTime > 0 && (now - this.lastBinanceDepthTime) > DATA_STALE_THRESHOLD_MS;
const binanceHealth = this.binanceDepth.getHealth();
const binanceUnhealthy = !binanceHealth.healthy;
const standxAccountAge = this.lastStandxAccountTime > 0 ? now - this.lastStandxAccountTime : 0; const standxAccountAge = this.lastStandxAccountTime > 0 ? now - this.lastStandxAccountTime : 0;
const standxAccountStaleByAge = this.lastStandxAccountTime > 0 && standxAccountAge > ACCOUNT_DATA_STALE_THRESHOLD_MS; const standxAccountStaleByAge = this.lastStandxAccountTime > 0 && standxAccountAge > ACCOUNT_DATA_STALE_THRESHOLD_MS;
@@ -1675,6 +1680,7 @@ export class MakerPointsEngine {
const shouldDefend = const shouldDefend =
standxDepthStale || standxDepthStale ||
binanceStale || binanceStale ||
binanceUnhealthy ||
standxAccountStale || standxAccountStale ||
accountInvalid || accountInvalid ||
standxRestUnhealthy || standxRestUnhealthy ||
@@ -1692,6 +1698,8 @@ export class MakerPointsEngine {
standxRestLastError: this.standxRestLastError, standxRestLastError: this.standxRestLastError,
marginModeNotIsolated, marginModeNotIsolated,
marginMode, marginMode,
binanceUnhealthy,
binanceHealthReason: binanceHealth.reason,
standxDepthAge: this.lastStandxDepthTime > 0 ? now - this.lastStandxDepthTime : 0, standxDepthAge: this.lastStandxDepthTime > 0 ? now - this.lastStandxDepthTime : 0,
binanceAge: this.lastBinanceDepthTime > 0 ? now - this.lastBinanceDepthTime : 0, binanceAge: this.lastBinanceDepthTime > 0 ? now - this.lastBinanceDepthTime : 0,
standxAccountAge, standxAccountAge,
@@ -1710,6 +1718,8 @@ export class MakerPointsEngine {
private enterDefenseMode(staleInfo: { private enterDefenseMode(staleInfo: {
standxDepthStale: boolean; standxDepthStale: boolean;
binanceStale: boolean; binanceStale: boolean;
binanceUnhealthy?: boolean;
binanceHealthReason?: string | null;
standxAccountStale: boolean; standxAccountStale: boolean;
accountInvalid: boolean; accountInvalid: boolean;
standxRestUnhealthy: boolean; standxRestUnhealthy: boolean;
@@ -1744,8 +1754,13 @@ export class MakerPointsEngine {
if (staleInfo.binanceStale) { if (staleInfo.binanceStale) {
staleItems.push(`Binance深度(${Math.round(staleInfo.binanceAge / 1000)}s)`); staleItems.push(`Binance深度(${Math.round(staleInfo.binanceAge / 1000)}s)`);
} }
if (staleInfo.binanceUnhealthy && staleInfo.binanceHealthReason) {
staleItems.push(`Binance簿记异常(${staleInfo.binanceHealthReason})`);
}
this.tradeLog.push("warn", `数据过时检测: ${staleItems.join(", ")},进入防御模式`); const staleSummary = staleItems.length > 0 ? staleItems.join(", ") : "unknown";
this.tradeLog.push("warn", `数据过时检测: ${staleSummary},进入防御模式`);
// 发送通知 // 发送通知
if (!this.defenseModeNotified) { if (!this.defenseModeNotified) {
@@ -1754,7 +1769,7 @@ export class MakerPointsEngine {
level: "warn", level: "warn",
symbol: this.config.symbol, symbol: this.config.symbol,
title: "防御模式", title: "防御模式",
message: `数据推送中断: ${staleItems.join(", ")},已取消所有挂单`, message: `数据推送中断: ${staleSummary},已取消所有挂单`,
details: staleInfo, details: staleInfo,
}); });
this.defenseModeNotified = true; this.defenseModeNotified = true;
@@ -0,0 +1,172 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import type { ExchangeAdapter } from "../src/exchanges/adapter";
import type { AsterAccountSnapshot, AsterDepth, AsterKline, AsterOrder, AsterTicker } from "../src/exchanges/types";
import { MakerPointsEngine } from "../src/strategy/maker-points-engine";
const ORIGINAL_FETCH = globalThis.fetch;
class StubAdapter implements ExchangeAdapter {
id = "standx";
supportsTrailingStops(): boolean {
return false;
}
watchAccount(_cb: (snapshot: AsterAccountSnapshot) => void): void {}
watchOrders(_cb: (orders: AsterOrder[]) => void): void {}
watchDepth(_symbol: string, _cb: (depth: AsterDepth) => void): void {}
watchTicker(_symbol: string, _cb: (ticker: AsterTicker) => void): void {}
watchKlines(_symbol: string, _interval: string, _cb: (klines: AsterKline[]) => void): void {}
async createOrder(): Promise<AsterOrder> {
throw new Error("not implemented");
}
async cancelOrder(): Promise<void> {}
async cancelOrders(): Promise<void> {}
async cancelAllOrders(): Promise<void> {}
async queryAccountSnapshot(): Promise<AsterAccountSnapshot | null> {
return {
canTrade: true,
canDeposit: true,
canWithdraw: true,
updateTime: Date.now(),
totalWalletBalance: "0",
totalUnrealizedProfit: "0",
positions: [],
assets: [],
marketType: "perp",
};
}
}
afterEach(() => {
globalThis.fetch = ORIGINAL_FETCH;
vi.restoreAllMocks();
vi.useRealTimers();
});
describe("MakerPointsEngine Binance depth health defense", () => {
it("enters defense mode when Binance depth tracker is unhealthy", () => {
vi.useFakeTimers();
globalThis.fetch = vi.fn(async () => {
throw new Error("network blocked in test");
}) as any;
const engine = new MakerPointsEngine(
{
symbol: "BTC-USD",
perOrderAmount: 0.01,
closeThreshold: 0,
stopLossUsd: 1,
refreshIntervalMs: 500,
maxLogEntries: 20,
maxCloseSlippagePct: 0.05,
priceTick: 0.1,
qtyStep: 0.001,
enableBand0To10: true,
enableBand10To30: false,
enableBand30To100: false,
band0To10Amount: 0.01,
band10To30Amount: 0.01,
band30To100Amount: 0.01,
minRepriceBps: 3,
enableBinanceDepthCancel: true,
filterMinDepth: 0,
},
new StubAdapter()
);
const now = Date.now();
(engine as any).lastStandxDepthTime = now;
(engine as any).lastStandxAccountTime = now;
(engine as any).lastBinanceDepthTime = now;
(engine as any).binanceDepth = {
getHealth: () => ({
started: true,
connected: true,
orderBookReady: false,
restHealthy: false,
healthy: false,
reason: "orderbook_not_ready",
lastEventAt: now,
lastSnapshotAt: 0,
lastRestSyncAt: 0,
localLastUpdateId: 0,
}),
stop: () => {},
};
(engine as any).checkDataStaleAndDefense();
expect((engine as any).defenseMode).toBe(true);
const logs = ((engine as any).tradeLog.all() as Array<{ detail: string }>).map((entry) => entry.detail);
expect(logs.some((detail) => detail.includes("Binance簿记异常(orderbook_not_ready)"))).toBe(true);
engine.stop();
});
it("exits defense mode after Binance depth health recovers", () => {
vi.useFakeTimers();
globalThis.fetch = vi.fn(async () => {
throw new Error("network blocked in test");
}) as any;
const engine = new MakerPointsEngine(
{
symbol: "BTC-USD",
perOrderAmount: 0.01,
closeThreshold: 0,
stopLossUsd: 1,
refreshIntervalMs: 500,
maxLogEntries: 20,
maxCloseSlippagePct: 0.05,
priceTick: 0.1,
qtyStep: 0.001,
enableBand0To10: true,
enableBand10To30: false,
enableBand30To100: false,
band0To10Amount: 0.01,
band10To30Amount: 0.01,
band30To100Amount: 0.01,
minRepriceBps: 3,
enableBinanceDepthCancel: true,
filterMinDepth: 0,
},
new StubAdapter()
);
const now = Date.now();
(engine as any).lastStandxDepthTime = now;
(engine as any).lastStandxAccountTime = now;
(engine as any).lastBinanceDepthTime = now;
let unhealthy = true;
(engine as any).binanceDepth = {
getHealth: () => ({
started: true,
connected: true,
orderBookReady: !unhealthy,
restHealthy: !unhealthy,
healthy: !unhealthy,
reason: unhealthy ? "orderbook_not_ready" : null,
lastEventAt: now,
lastSnapshotAt: now,
lastRestSyncAt: now,
localLastUpdateId: unhealthy ? 0 : 100,
}),
stop: () => {},
};
(engine as any).checkDataStaleAndDefense();
expect((engine as any).defenseMode).toBe(true);
unhealthy = false;
(engine as any).checkDataStaleAndDefense();
expect((engine as any).defenseMode).toBe(false);
engine.stop();
});
});
+4
View File
@@ -94,6 +94,8 @@ describe("MakerPointsEngine defense-mode REST polling", () => {
(engine as any).enterDefenseMode({ (engine as any).enterDefenseMode({
standxDepthStale: true, standxDepthStale: true,
binanceStale: false, binanceStale: false,
binanceUnhealthy: false,
binanceHealthReason: null,
standxAccountStale: false, standxAccountStale: false,
accountInvalid: false, accountInvalid: false,
standxRestUnhealthy: false, standxRestUnhealthy: false,
@@ -146,6 +148,8 @@ describe("MakerPointsEngine defense-mode REST polling", () => {
(engine as any).enterDefenseMode({ (engine as any).enterDefenseMode({
standxDepthStale: true, standxDepthStale: true,
binanceStale: false, binanceStale: false,
binanceUnhealthy: false,
binanceHealthReason: null,
standxAccountStale: false, standxAccountStale: false,
accountInvalid: false, accountInvalid: false,
standxRestUnhealthy: false, standxRestUnhealthy: false,