mirror of
https://github.com/discountry/ritmex-bot.git
synced 2026-09-09 08:18:07 +00:00
Enhance WebSocket connection management and data handling
- Introduced constants for WebSocket reconnection delays, heartbeat timeout, and data staleness thresholds. - Implemented heartbeat monitoring to ensure timely reconnections on inactivity. - Added data staleness checks to trigger REST API calls when market or account data is outdated. - Enhanced the StandxGateway class with methods for managing heartbeat and data checks, improving overall connection reliability and data integrity.
This commit is contained in:
@@ -48,7 +48,20 @@ const DEFAULT_WS_URL = "wss://perps.standx.com/ws-stream/v1";
|
||||
const DEFAULT_KLINE_LIMIT = 200;
|
||||
const KLINE_REFRESH_MS = 30_000;
|
||||
const FUNDING_REFRESH_MS = 60_000;
|
||||
const WS_RECONNECT_DELAY = 2000;
|
||||
|
||||
// ========== WebSocket 连接管理常量 ==========
|
||||
// 基础重连延迟(毫秒)
|
||||
const WS_RECONNECT_DELAY_BASE = 2000;
|
||||
// 最大重连延迟(毫秒)- 指数退避上限
|
||||
const WS_RECONNECT_DELAY_MAX = 30_000;
|
||||
// 心跳超时(毫秒)- StandX 服务器 5 分钟无 pong 会断连,我们设置 2 分钟作为安全阈值
|
||||
const WS_HEARTBEAT_TIMEOUT = 120_000;
|
||||
// 心跳检查间隔(毫秒)- 每 30 秒检查一次是否收到消息
|
||||
const WS_HEARTBEAT_CHECK_INTERVAL = 30_000;
|
||||
// 数据过时阈值(毫秒)- 超过此时间未收到行情/仓位数据,启动 REST 主动拉取
|
||||
const WS_DATA_STALE_THRESHOLD = 3000;
|
||||
// REST 轮询间隔(毫秒)- WS 断连或数据过时时的 REST 拉取间隔
|
||||
const REST_POLL_INTERVAL = 2000;
|
||||
|
||||
const SUPPORTED_QUOTES = ["USD", "USDT", "USDC", "DUSD"];
|
||||
|
||||
@@ -100,7 +113,7 @@ class StandxRequestSigner {
|
||||
const requestId = crypto.randomUUID();
|
||||
const timestamp = Date.now();
|
||||
const signMessage = `${version},${requestId},${timestamp},${payload}`;
|
||||
const signatureBytes = await sign(Buffer.from(signMessage, "utf-8"), this.privateKey);
|
||||
const signatureBytes = sign(Buffer.from(signMessage, "utf-8"), this.privateKey);
|
||||
return {
|
||||
"x-request-sign-version": version,
|
||||
"x-request-id": requestId,
|
||||
@@ -416,6 +429,26 @@ export class StandxGateway {
|
||||
private marketReconnectTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
private readonly subscriptions = new Set<string>();
|
||||
|
||||
// ========== 心跳与连接管理 ==========
|
||||
// 上次收到消息的时间戳
|
||||
private lastMessageTime = 0;
|
||||
// 心跳检查定时器
|
||||
private heartbeatTimer: ReturnType<typeof setInterval> | null = null;
|
||||
// 重连次数(用于指数退避)
|
||||
private reconnectAttempts = 0;
|
||||
|
||||
// ========== 数据过时检测与 REST 备用 ==========
|
||||
// 上次收到行情数据(price/depth)的时间戳
|
||||
private lastMarketDataTime = 0;
|
||||
// 上次收到账户数据(position/balance)的时间戳
|
||||
private lastAccountDataTime = 0;
|
||||
// 数据过时检查定时器
|
||||
private dataStaleCheckTimer: ReturnType<typeof setInterval> | null = null;
|
||||
// REST 轮询定时器(WS 断连时启用)
|
||||
private restPollTimer: ReturnType<typeof setInterval> | null = null;
|
||||
// REST 轮询是否激活
|
||||
private restPollActive = false;
|
||||
|
||||
private readonly klineTimers = new Map<string, PollTimer>();
|
||||
private readonly fundingTimers = new Map<string, PollTimer>();
|
||||
|
||||
@@ -852,7 +885,18 @@ export class StandxGateway {
|
||||
const handleOpen = () => {
|
||||
this.marketWsReady = true;
|
||||
this.marketWsAuthed = false;
|
||||
// 重置重连计数和时间戳
|
||||
this.reconnectAttempts = 0;
|
||||
this.lastMessageTime = Date.now();
|
||||
this.lastMarketDataTime = Date.now();
|
||||
this.lastAccountDataTime = Date.now();
|
||||
this.logDebug("ws open");
|
||||
// 启动心跳监控
|
||||
this.startHeartbeatMonitor();
|
||||
// 启动数据过时检测
|
||||
this.startDataStaleCheck();
|
||||
// 停止 REST 轮询(WS 恢复后不再需要)
|
||||
this.stopRestPoll();
|
||||
this.sendAuthIfNeeded();
|
||||
};
|
||||
const handleClose = () => {
|
||||
@@ -861,6 +905,9 @@ export class StandxGateway {
|
||||
this.marketWsAuthed = false;
|
||||
this.marketWsAuthRequested = false;
|
||||
this.marketWs = null;
|
||||
// 停止心跳监控和数据过时检测
|
||||
this.stopHeartbeatMonitor();
|
||||
this.stopDataStaleCheck();
|
||||
this.logDebug("ws close");
|
||||
// 触发断连事件,启动断连保护
|
||||
if (wasReady) {
|
||||
@@ -899,15 +946,23 @@ export class StandxGateway {
|
||||
|
||||
private scheduleReconnect(): void {
|
||||
if (this.marketReconnectTimer) return;
|
||||
this.logDebug("scheduling reconnect in 2s");
|
||||
// 指数退避:delay = min(base * 2^attempts, max)
|
||||
const delay = Math.min(
|
||||
WS_RECONNECT_DELAY_BASE * Math.pow(2, this.reconnectAttempts),
|
||||
WS_RECONNECT_DELAY_MAX
|
||||
);
|
||||
this.reconnectAttempts += 1;
|
||||
this.logDebug(`scheduling reconnect in ${delay}ms (attempt ${this.reconnectAttempts})`);
|
||||
this.marketReconnectTimer = setTimeout(() => {
|
||||
this.marketReconnectTimer = null;
|
||||
this.logDebug("attempting reconnect");
|
||||
this.connectMarketWs();
|
||||
}, WS_RECONNECT_DELAY);
|
||||
}, delay);
|
||||
}
|
||||
|
||||
private handleMarketMessage(event: { data: any }): void {
|
||||
// 更新最后收到消息的时间(心跳监控)
|
||||
this.lastMessageTime = Date.now();
|
||||
this.logRawPayload(event.data);
|
||||
const payloads = parseJsonPayloads(event.data);
|
||||
if (payloads.length === 0) return;
|
||||
@@ -938,9 +993,11 @@ export class StandxGateway {
|
||||
return;
|
||||
}
|
||||
if (channel === "depth_book") {
|
||||
// 更新行情数据时间戳
|
||||
this.lastMarketDataTime = Date.now();
|
||||
const data = message.data as StandxDepthBook | undefined;
|
||||
const rawSymbol = data?.symbol ?? message?.symbol;
|
||||
if (!rawSymbol) return;
|
||||
if (!rawSymbol || !data) return;
|
||||
const bids = normalizeDepthLevels((data.bids ?? []).map(([price, qty]) => [String(price), String(qty)]), "bid");
|
||||
const asks = normalizeDepthLevels((data.asks ?? []).map(([price, qty]) => [String(price), String(qty)]), "ask");
|
||||
const decrossed = decrossDepthBook(bids, asks);
|
||||
@@ -969,13 +1026,15 @@ export class StandxGateway {
|
||||
lastUpdateId: Number(message.seq ?? Date.now()),
|
||||
bids: finalBids,
|
||||
asks: finalAsks,
|
||||
eventTime: toTimestamp(data.time),
|
||||
eventTime: Date.now(),
|
||||
symbol: rawSymbol,
|
||||
};
|
||||
this.emitDepth(rawSymbol, depth);
|
||||
return;
|
||||
}
|
||||
if (channel === "price") {
|
||||
// 更新行情数据时间戳
|
||||
this.lastMarketDataTime = Date.now();
|
||||
const data = message.data as StandxPrice | undefined;
|
||||
if (!data?.symbol) return;
|
||||
const ticker = this.mapTicker(data);
|
||||
@@ -983,6 +1042,8 @@ export class StandxGateway {
|
||||
return;
|
||||
}
|
||||
if (channel === "order") {
|
||||
// 更新账户数据时间戳
|
||||
this.lastAccountDataTime = Date.now();
|
||||
const payload = message.data as StandxOrder | StandxOrder[] | undefined;
|
||||
if (!payload) return;
|
||||
const items = Array.isArray(payload) ? payload : [payload];
|
||||
@@ -994,6 +1055,8 @@ export class StandxGateway {
|
||||
return;
|
||||
}
|
||||
if (channel === "position") {
|
||||
// 更新账户数据时间戳
|
||||
this.lastAccountDataTime = Date.now();
|
||||
const payload = message.data as StandxPosition | StandxPosition[] | undefined;
|
||||
if (!payload) return;
|
||||
const items = Array.isArray(payload) ? payload : [payload];
|
||||
@@ -1006,6 +1069,8 @@ export class StandxGateway {
|
||||
return;
|
||||
}
|
||||
if (channel === "balance") {
|
||||
// 更新账户数据时间戳
|
||||
this.lastAccountDataTime = Date.now();
|
||||
const payload = message.data as StandxBalance | StandxBalance[] | undefined;
|
||||
if (!payload) return;
|
||||
const items = Array.isArray(payload) ? payload : [payload];
|
||||
@@ -1048,6 +1113,7 @@ export class StandxGateway {
|
||||
if (!this.marketWsAuthed) return;
|
||||
for (const entry of this.subscriptions) {
|
||||
const [channel, symbol] = entry.split(":");
|
||||
if (!channel) continue;
|
||||
this.sendSubscribe({ channel, ...(symbol ? { symbol } : {}) });
|
||||
}
|
||||
}
|
||||
@@ -1057,6 +1123,194 @@ export class StandxGateway {
|
||||
this.marketWs?.send(JSON.stringify({ subscribe: stream }));
|
||||
}
|
||||
|
||||
// ========== 心跳监控 ==========
|
||||
/**
|
||||
* 启动心跳监控
|
||||
* 定期检查是否收到消息,如果超时则主动触发重连
|
||||
*/
|
||||
private startHeartbeatMonitor(): void {
|
||||
this.stopHeartbeatMonitor();
|
||||
this.heartbeatTimer = setInterval(() => {
|
||||
const now = Date.now();
|
||||
const elapsed = now - this.lastMessageTime;
|
||||
if (elapsed > WS_HEARTBEAT_TIMEOUT) {
|
||||
this.logDebug(`heartbeat timeout (${elapsed}ms since last message), forcing reconnect`);
|
||||
this.forceReconnect("heartbeat_timeout");
|
||||
}
|
||||
}, WS_HEARTBEAT_CHECK_INTERVAL);
|
||||
}
|
||||
|
||||
/**
|
||||
* 停止心跳监控
|
||||
*/
|
||||
private stopHeartbeatMonitor(): void {
|
||||
if (this.heartbeatTimer) {
|
||||
clearInterval(this.heartbeatTimer);
|
||||
this.heartbeatTimer = null;
|
||||
}
|
||||
}
|
||||
|
||||
// ========== 数据过时检测与 REST 备用拉取 ==========
|
||||
/**
|
||||
* 启动数据过时检测
|
||||
* 即使 WS 连接正常,如果超过 3 秒未收到行情/账户数据,也主动通过 REST 拉取
|
||||
*/
|
||||
private startDataStaleCheck(): void {
|
||||
this.stopDataStaleCheck();
|
||||
this.dataStaleCheckTimer = setInterval(() => {
|
||||
const now = Date.now();
|
||||
const marketStale = now - this.lastMarketDataTime > WS_DATA_STALE_THRESHOLD;
|
||||
const accountStale = now - this.lastAccountDataTime > WS_DATA_STALE_THRESHOLD;
|
||||
|
||||
if (marketStale || accountStale) {
|
||||
this.logDebug("data stale detected", {
|
||||
marketStaleMs: now - this.lastMarketDataTime,
|
||||
accountStaleMs: now - this.lastAccountDataTime,
|
||||
marketStale,
|
||||
accountStale,
|
||||
});
|
||||
// 主动通过 REST 拉取数据
|
||||
this.fetchStaleData(marketStale, accountStale);
|
||||
}
|
||||
}, 1000); // 每秒检查一次
|
||||
}
|
||||
|
||||
/**
|
||||
* 停止数据过时检测
|
||||
*/
|
||||
private stopDataStaleCheck(): void {
|
||||
if (this.dataStaleCheckTimer) {
|
||||
clearInterval(this.dataStaleCheckTimer);
|
||||
this.dataStaleCheckTimer = null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 主动拉取过时的数据
|
||||
*/
|
||||
private fetchStaleData(marketStale: boolean, accountStale: boolean): void {
|
||||
// 获取当前订阅的 symbols
|
||||
const symbols = new Set<string>();
|
||||
for (const key of this.subscriptions) {
|
||||
const [, symbol] = key.split(":");
|
||||
if (symbol) symbols.add(symbol);
|
||||
}
|
||||
|
||||
if (marketStale) {
|
||||
for (const symbol of symbols) {
|
||||
void this.fetchTickerSnapshot(symbol).catch((e) => this.logger("staleTickerFetch", e));
|
||||
void this.fetchDepthSnapshot(symbol).catch((e) => this.logger("staleDepthFetch", e));
|
||||
}
|
||||
// 更新时间戳避免重复拉取
|
||||
this.lastMarketDataTime = Date.now();
|
||||
}
|
||||
|
||||
if (accountStale) {
|
||||
void this.refreshAccountSnapshot().catch((e) => this.logger("staleAccountFetch", e));
|
||||
for (const symbol of symbols) {
|
||||
void this.refreshOpenOrders(symbol).catch((e) => this.logger("staleOrdersFetch", e));
|
||||
}
|
||||
// 更新时间戳避免重复拉取
|
||||
this.lastAccountDataTime = Date.now();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 启动 REST 轮询(WS 断连时使用)
|
||||
* 持续通过 REST API 拉取行情和账户数据,确保止损等逻辑能正常工作
|
||||
*/
|
||||
private startRestPoll(): void {
|
||||
if (this.restPollActive) return;
|
||||
this.restPollActive = true;
|
||||
this.logDebug("REST poll started (WS disconnected)");
|
||||
|
||||
const poll = async () => {
|
||||
if (!this.restPollActive) return;
|
||||
|
||||
// 获取当前订阅的 symbols
|
||||
const symbols = new Set<string>();
|
||||
for (const key of this.subscriptions) {
|
||||
const [, symbol] = key.split(":");
|
||||
if (symbol) symbols.add(symbol);
|
||||
}
|
||||
|
||||
// 拉取行情数据
|
||||
for (const symbol of symbols) {
|
||||
try {
|
||||
await this.fetchTickerSnapshot(symbol);
|
||||
this.lastMarketDataTime = Date.now();
|
||||
} catch (e) {
|
||||
this.logger("restPollTicker", e);
|
||||
}
|
||||
try {
|
||||
await this.fetchDepthSnapshot(symbol);
|
||||
} catch (e) {
|
||||
this.logger("restPollDepth", e);
|
||||
}
|
||||
}
|
||||
|
||||
// 拉取账户数据
|
||||
try {
|
||||
await this.refreshAccountSnapshot();
|
||||
this.lastAccountDataTime = Date.now();
|
||||
} catch (e) {
|
||||
this.logger("restPollAccount", e);
|
||||
}
|
||||
|
||||
// 继续下一次轮询
|
||||
if (this.restPollActive) {
|
||||
this.restPollTimer = setTimeout(() => void poll(), REST_POLL_INTERVAL);
|
||||
}
|
||||
};
|
||||
|
||||
void poll();
|
||||
}
|
||||
|
||||
/**
|
||||
* 停止 REST 轮询
|
||||
*/
|
||||
private stopRestPoll(): void {
|
||||
if (!this.restPollActive) return;
|
||||
this.restPollActive = false;
|
||||
if (this.restPollTimer) {
|
||||
clearTimeout(this.restPollTimer);
|
||||
this.restPollTimer = null;
|
||||
}
|
||||
this.logDebug("REST poll stopped (WS restored)");
|
||||
}
|
||||
|
||||
/**
|
||||
* 强制重连(用于心跳超时)
|
||||
* 与普通断连不同,这里不增加重连计数(因为是主动行为)
|
||||
*/
|
||||
private forceReconnect(reason: string): void {
|
||||
this.logDebug(`force reconnect: ${reason}`);
|
||||
// 停止监控
|
||||
this.stopHeartbeatMonitor();
|
||||
this.stopDataStaleCheck();
|
||||
// 关闭现有连接
|
||||
if (this.marketWs) {
|
||||
try {
|
||||
this.marketWs.close();
|
||||
} catch {
|
||||
// ignore close errors
|
||||
}
|
||||
this.marketWs = null;
|
||||
}
|
||||
// 重置状态
|
||||
const wasReady = this.marketWsReady;
|
||||
this.marketWsReady = false;
|
||||
this.marketWsAuthed = false;
|
||||
this.marketWsAuthRequested = false;
|
||||
// 触发断连事件
|
||||
if (wasReady) {
|
||||
this.onDisconnect();
|
||||
}
|
||||
// 立即重连(不使用指数退避,因为是主动行为)
|
||||
this.reconnectAttempts = 0;
|
||||
this.scheduleReconnect();
|
||||
}
|
||||
|
||||
private logDebug(context: string, detail?: unknown): void {
|
||||
if (!this.debugWs) return;
|
||||
if (detail === undefined) {
|
||||
@@ -1509,6 +1763,9 @@ export class StandxGateway {
|
||||
if (this.lastKnownOpenOrders.length > 0 && this.disconnectedSymbol) {
|
||||
this.startDisconnectCancelRetry(this.disconnectedSymbol);
|
||||
}
|
||||
|
||||
// 启动 REST 轮询,确保断连期间仍能获取行情和账户数据(用于止损等逻辑)
|
||||
this.startRestPoll();
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user