From 69271d33cac7cd2cf3592735792419401d5e5537 Mon Sep 17 00:00:00 2001 From: discountry Date: Wed, 21 Jan 2026 16:03:34 +0800 Subject: [PATCH] Refactor MakerPointsEngine and BinanceDepthTracker for improved connection management - Renamed connection state variable in MakerPointsEngine for clarity. - Added connection state change listeners in BinanceDepthTracker to handle connection status updates. - Implemented heartbeat monitoring and connection duration checks in BinanceDepthTracker to enhance WebSocket reliability. - Introduced data staleness checks and improved error handling for WebSocket connections. - Enhanced logging for connection events to provide better insights into connection status changes. --- src/strategy/common/binance-depth.ts | 236 +++++++++++++++++++++++++-- src/strategy/maker-points-engine.ts | 30 +++- 2 files changed, 244 insertions(+), 22 deletions(-) diff --git a/src/strategy/common/binance-depth.ts b/src/strategy/common/binance-depth.ts index bafbb89..b436b85 100644 --- a/src/strategy/common/binance-depth.ts +++ b/src/strategy/common/binance-depth.ts @@ -8,6 +8,23 @@ const WebSocketCtor: typeof globalThis.WebSocket = const DEFAULT_BASE_URL = "wss://fstream.binance.com/ws"; +// ========== Binance WebSocket 连接管理常量 ========== +// Binance 服务器每 3 分钟发送 ping,10 分钟无 pong 会断连 +// 我们设置 5 分钟作为心跳超时阈值(保守值) +const HEARTBEAT_TIMEOUT_MS = 5 * 60 * 1000; +// 心跳检查间隔(每 30 秒检查一次) +const HEARTBEAT_CHECK_INTERVAL_MS = 30_000; +// Binance 连接最长有效期 24 小时,我们设置 23 小时主动重连 +const MAX_CONNECTION_DURATION_MS = 23 * 60 * 60 * 1000; +// 数据过时阈值(毫秒)- 超过此时间未收到数据,标记为不可用 +const DATA_STALE_THRESHOLD_MS = 5_000; +// 基础重连延迟 +const RECONNECT_DELAY_BASE_MS = 3000; +// 最大重连延迟 +const RECONNECT_DELAY_MAX_MS = 60_000; + +export type BinanceConnectionState = "connected" | "disconnected" | "stale"; + export interface BinanceDepthSnapshot { symbol: string; buySum: number; @@ -18,13 +35,29 @@ export interface BinanceDepthSnapshot { updatedAt: number; } +export type BinanceConnectionListener = (state: BinanceConnectionState) => void; + export class BinanceDepthTracker { private ws: WebSocket | null = null; private reconnectTimer: ReturnType | null = null; - private reconnectDelayMs = 3000; + private reconnectDelayMs = RECONNECT_DELAY_BASE_MS; private stopped = false; private snapshot: BinanceDepthSnapshot | null = null; private listeners = new Set<(snapshot: BinanceDepthSnapshot) => void>(); + private connectionListeners = new Set(); + + // ========== 心跳与连接管理 ========== + // 上次收到消息的时间戳 + private lastMessageTime = 0; + // 心跳检查定时器 + private heartbeatTimer: ReturnType | null = null; + // 连接建立时间(用于日志记录) + // eslint-disable-next-line @typescript-eslint/no-unused-vars + private connectionStartTime = 0; + // 24 小时重连定时器 + private maxDurationTimer: ReturnType | null = null; + // 当前连接状态 + private connectionState: BinanceConnectionState = "disconnected"; constructor( private readonly symbol: string, @@ -43,18 +76,7 @@ export class BinanceDepthTracker { stop(): void { this.stopped = true; - if (this.reconnectTimer) { - clearTimeout(this.reconnectTimer); - this.reconnectTimer = null; - } - if (this.ws) { - try { - this.ws.close(); - } catch { - // Ignore close errors - } - this.ws = null; - } + this.cleanup(); } onUpdate(handler: (snapshot: BinanceDepthSnapshot) => void): void { @@ -65,37 +87,122 @@ export class BinanceDepthTracker { this.listeners.delete(handler); } + /** + * 监听连接状态变化 + */ + onConnectionChange(handler: BinanceConnectionListener): void { + this.connectionListeners.add(handler); + } + + offConnectionChange(handler: BinanceConnectionListener): void { + this.connectionListeners.delete(handler); + } + getSnapshot(): BinanceDepthSnapshot | null { return this.snapshot ? { ...this.snapshot } : null; } + /** + * 获取当前连接状态 + */ + getConnectionState(): BinanceConnectionState { + return this.connectionState; + } + + /** + * 检查数据是否过时 + */ + isDataStale(): boolean { + if (!this.snapshot) return true; + return Date.now() - this.snapshot.updatedAt > DATA_STALE_THRESHOLD_MS; + } + + private cleanup(): void { + // 停止心跳监控 + if (this.heartbeatTimer) { + clearInterval(this.heartbeatTimer); + this.heartbeatTimer = null; + } + // 停止 24 小时重连定时器 + if (this.maxDurationTimer) { + clearTimeout(this.maxDurationTimer); + this.maxDurationTimer = null; + } + // 停止重连定时器 + if (this.reconnectTimer) { + clearTimeout(this.reconnectTimer); + this.reconnectTimer = null; + } + // 关闭 WebSocket + if (this.ws) { + try { + this.ws.close(); + } catch { + // Ignore close errors + } + this.ws = null; + } + } + private connect(): void { if (this.ws || this.stopped) return; const url = this.buildUrl(); this.ws = new WebSocketCtor(url); const handleOpen = () => { - this.reconnectDelayMs = 3000; + this.reconnectDelayMs = RECONNECT_DELAY_BASE_MS; + this.connectionStartTime = Date.now(); + this.lastMessageTime = Date.now(); + this.updateConnectionState("connected"); + + // 启动心跳监控 + this.startHeartbeatMonitor(); + // 启动 24 小时自动重连定时器 + this.startMaxDurationTimer(); + + this.options?.logger?.("binanceDepth", "WebSocket connected"); }; const handleClose = () => { this.ws = null; + this.stopHeartbeatMonitor(); + this.stopMaxDurationTimer(); + this.updateConnectionState("disconnected"); + if (!this.stopped) { + this.options?.logger?.("binanceDepth", "WebSocket closed, scheduling reconnect"); this.scheduleReconnect(); } }; const handleError = (error: unknown) => { this.options?.logger?.("binanceDepth", error); + // 如果连接从未成功建立,需要清理并重连 + if (this.ws && this.connectionState === "disconnected") { + this.ws = null; + this.scheduleReconnect(); + } }; const handleMessage = (event: { data: unknown }) => { + this.lastMessageTime = Date.now(); + // 如果之前是 stale 状态,恢复为 connected + if (this.connectionState === "stale") { + this.updateConnectionState("connected"); + } this.handlePayload(event.data); }; + // 处理 Binance 服务器的 ping 帧 + // 根据文档:必须尽快回复 pong,payload 为 ping 的 payload 副本 const handlePing = (data: unknown) => { + this.lastMessageTime = Date.now(); if (this.ws && "pong" in this.ws && typeof this.ws.pong === "function") { - this.ws.pong(data as any); + try { + this.ws.pong(data as any); + } catch (error) { + this.options?.logger?.("binanceDepth pong", error); + } } }; @@ -130,11 +237,101 @@ export class BinanceDepthTracker { if (this.reconnectTimer || this.stopped) return; this.reconnectTimer = setTimeout(() => { this.reconnectTimer = null; - this.reconnectDelayMs = Math.min(this.reconnectDelayMs * 2, 60_000); + this.reconnectDelayMs = Math.min(this.reconnectDelayMs * 2, RECONNECT_DELAY_MAX_MS); this.connect(); }, this.reconnectDelayMs); } + /** + * 启动心跳监控 + * 根据 Binance 文档:10 分钟无 pong 会断连 + * 我们设置 5 分钟作为心跳超时阈值 + */ + private startHeartbeatMonitor(): void { + this.stopHeartbeatMonitor(); + this.heartbeatTimer = setInterval(() => { + const now = Date.now(); + const elapsed = now - this.lastMessageTime; + + // 检查数据是否过时(5 秒无数据) + if (elapsed > DATA_STALE_THRESHOLD_MS && this.connectionState === "connected") { + this.updateConnectionState("stale"); + this.options?.logger?.("binanceDepth", `Data stale: ${elapsed}ms since last message`); + } + + // 检查心跳超时(5 分钟无消息) + if (elapsed > HEARTBEAT_TIMEOUT_MS) { + this.options?.logger?.("binanceDepth", `Heartbeat timeout: ${elapsed}ms, forcing reconnect`); + this.forceReconnect("heartbeat_timeout"); + } + }, HEARTBEAT_CHECK_INTERVAL_MS); + } + + private stopHeartbeatMonitor(): void { + if (this.heartbeatTimer) { + clearInterval(this.heartbeatTimer); + this.heartbeatTimer = null; + } + } + + /** + * 启动 24 小时自动重连定时器 + * 根据 Binance 文档:连接最长有效期 24 小时 + * 我们设置 23 小时主动重连,避免被服务器断开 + */ + private startMaxDurationTimer(): void { + this.stopMaxDurationTimer(); + this.maxDurationTimer = setTimeout(() => { + this.options?.logger?.("binanceDepth", "Max connection duration reached (23h), reconnecting"); + this.forceReconnect("max_duration"); + }, MAX_CONNECTION_DURATION_MS); + } + + private stopMaxDurationTimer(): void { + if (this.maxDurationTimer) { + clearTimeout(this.maxDurationTimer); + this.maxDurationTimer = null; + } + } + + /** + * 强制重连 + */ + private forceReconnect(reason: string): void { + this.options?.logger?.("binanceDepth", `Force reconnect: ${reason}`); + this.stopHeartbeatMonitor(); + this.stopMaxDurationTimer(); + + if (this.ws) { + try { + this.ws.close(); + } catch { + // ignore + } + this.ws = null; + } + + this.updateConnectionState("disconnected"); + // 立即重连(不使用指数退避) + this.reconnectDelayMs = RECONNECT_DELAY_BASE_MS; + this.scheduleReconnect(); + } + + /** + * 更新连接状态并通知监听器 + */ + private updateConnectionState(state: BinanceConnectionState): void { + if (this.connectionState === state) return; + this.connectionState = state; + for (const listener of this.connectionListeners) { + try { + listener(state); + } catch (error) { + this.options?.logger?.("binanceDepth connectionListener", error); + } + } + } + private handlePayload(data: unknown): void { const payload = this.parsePayload(data); if (!payload) return; @@ -158,7 +355,11 @@ export class BinanceDepthTracker { updatedAt: Date.now(), }; for (const listener of this.listeners) { - listener({ ...this.snapshot }); + try { + listener({ ...this.snapshot }); + } catch (error) { + this.options?.logger?.("binanceDepth listener", error); + } } } @@ -174,4 +375,3 @@ export class BinanceDepthTracker { } } } - diff --git a/src/strategy/maker-points-engine.ts b/src/strategy/maker-points-engine.ts index ba73b5d..1befc1b 100644 --- a/src/strategy/maker-points-engine.ts +++ b/src/strategy/maker-points-engine.ts @@ -1,5 +1,5 @@ import type { MakerPointsConfig } from "../config"; -import type { ExchangeAdapter, ConnectionEventType } from "../exchanges/adapter"; +import type { ExchangeAdapter } from "../exchanges/adapter"; import type { AsterAccountSnapshot, AsterDepth, @@ -154,7 +154,7 @@ export class MakerPointsEngine { private lastPositionSide: "LONG" | "SHORT" | "FLAT" = "FLAT"; // 连接保护相关状态 - private connectionState: "connected" | "disconnected" = "connected"; + private standxConnectionState: "connected" | "disconnected" = "connected"; private reconnectResetPending = false; private lastRepriceQueryTime = 0; private readonly repriceQueryIntervalMs = 3000; // 最小查询间隔 @@ -179,6 +179,19 @@ export class MakerPointsEngine { this.feedStatus.binance = true; this.emitUpdate(); }); + // 监听 Binance 连接状态变化 + this.binanceDepth.onConnectionChange((state) => { + if (state === "disconnected") { + this.feedStatus.binance = false; + this.tradeLog.push("warn", "Binance 深度连接断开"); + } else if (state === "stale") { + this.tradeLog.push("warn", "Binance 深度数据过时"); + } else if (state === "connected") { + this.feedStatus.binance = true; + this.tradeLog.push("info", "Binance 深度连接恢复"); + } + this.emitUpdate(); + }); this.syncPrecision(); this.bootstrap(); } @@ -325,7 +338,7 @@ export class MakerPointsEngine { * 处理断连事件 */ private handleDisconnect(symbol: string): void { - this.connectionState = "disconnected"; + this.standxConnectionState = "disconnected"; this.tradeLog.push("warn", `WebSocket 断连 (${symbol}),启动断连保护`); this.notify({ type: "token_expired", @@ -342,7 +355,7 @@ export class MakerPointsEngine { * 重连后需要重新查询挂单并取消所有挂单 */ private async handleReconnect(symbol: string): Promise { - this.connectionState = "connected"; + this.standxConnectionState = "connected"; this.reconnectResetPending = true; this.tradeLog.push("info", `WebSocket 重连成功 (${symbol}),开始重连保护流程`); @@ -420,6 +433,10 @@ export class MakerPointsEngine { private async tick(): Promise { if (this.processing) return; + // 重连处理期间不执行主循环,避免状态竞争 + if (this.reconnectResetPending) return; + // 止损执行期间不执行主循环,避免订单冲突 + if (this.stopLossProcessing) return; this.processing = true; let hadRateLimit = false; try { @@ -747,6 +764,11 @@ export class MakerPointsEngine { } private async syncOrders(targets: DesiredOrder[], _closeOnly: boolean): Promise { + // 止损执行期间不进行挂单操作,避免订单冲突 + if (this.stopLossProcessing) return; + // 重连处理期间不进行挂单操作 + if (this.reconnectResetPending) return; + // 价格变化保护:如果需要 reprice 且距上次查询已过足够时间,先查询真实挂单 const shouldVerifyOrders = await this.verifyOrdersIfNeeded(); if (shouldVerifyOrders) {