feat: implement client ping/pong mechanism in LighterGateway for improved WebSocket connection health

This commit is contained in:
discountry
2025-11-13 02:53:23 +08:00
parent a9fa7f2e19
commit 59ccd1fea8
+51
View File
@@ -119,6 +119,7 @@ const KLINE_DEFAULT_COUNT = 120;
const DEFAULT_TICKER_POLL_MS = 3000; const DEFAULT_TICKER_POLL_MS = 3000;
const DEFAULT_KLINE_POLL_MS = 15000; const DEFAULT_KLINE_POLL_MS = 15000;
const WS_HEARTBEAT_INTERVAL_MS = 5_000; const WS_HEARTBEAT_INTERVAL_MS = 5_000;
const CLIENT_PING_INTERVAL_MS = 2_000;
const WS_STALE_TIMEOUT_MS = 20_000; const WS_STALE_TIMEOUT_MS = 20_000;
const FEED_STALE_TIMEOUT_MS = 8_000; const FEED_STALE_TIMEOUT_MS = 8_000;
const STALE_CHECK_INTERVAL_MS = 2_000; const STALE_CHECK_INTERVAL_MS = 2_000;
@@ -198,6 +199,7 @@ export class LighterGateway {
private readonly wsUrl: string; private readonly wsUrl: string;
private connectPromise: Promise<void> | null = null; private connectPromise: Promise<void> | null = null;
private heartbeatTimer: ReturnType<typeof setInterval> | null = null; private heartbeatTimer: ReturnType<typeof setInterval> | null = null;
private pingTimer: ReturnType<typeof setInterval> | null = null;
private lastMessageAt = 0; private lastMessageAt = 0;
private accountDetails: LighterAccountDetails | null = null; private accountDetails: LighterAccountDetails | null = null;
@@ -528,6 +530,7 @@ export class LighterGateway {
const cleanup = () => { const cleanup = () => {
ws.removeAllListeners(); ws.removeAllListeners();
this.stopHeartbeat(); this.stopHeartbeat();
this.stopClientPing();
if (this.ws === ws) { if (this.ws === ws) {
this.ws = null; this.ws = null;
} }
@@ -541,6 +544,7 @@ export class LighterGateway {
try { try {
this.lastMessageAt = Date.now(); this.lastMessageAt = Date.now();
this.startHeartbeat(); this.startHeartbeat();
this.startClientPing();
await this.subscribeChannels(); await this.subscribeChannels();
this.startStaleMonitor(); this.startStaleMonitor();
settled = true; settled = true;
@@ -640,6 +644,7 @@ export class LighterGateway {
this.logger("ws:terminate", error); this.logger("ws:terminate", error);
} }
this.stopHeartbeat(); this.stopHeartbeat();
this.stopClientPing();
this.scheduleReconnect(); this.scheduleReconnect();
} }
@@ -656,6 +661,7 @@ export class LighterGateway {
this.logger("ws:terminate", error); this.logger("ws:terminate", error);
} finally { } finally {
this.stopHeartbeat(); this.stopHeartbeat();
this.stopClientPing();
this.scheduleReconnect(); this.scheduleReconnect();
} }
return; return;
@@ -675,6 +681,26 @@ export class LighterGateway {
} }
} }
private startClientPing(): void {
if (this.pingTimer) return;
this.pingTimer = setInterval(() => {
const ws = this.ws;
if (!ws || ws.readyState !== WebSocket.OPEN) return;
try {
ws.send(JSON.stringify({ type: "ping" }));
} catch (error) {
this.logger("ws:clientPing", error);
}
}, CLIENT_PING_INTERVAL_MS);
}
private stopClientPing(): void {
if (this.pingTimer) {
clearInterval(this.pingTimer);
this.pingTimer = null;
}
}
private handleMessage(data: WebSocket.RawData): void { private handleMessage(data: WebSocket.RawData): void {
try { try {
const text = typeof data === "string" ? data : data.toString("utf8"); const text = typeof data === "string" ? data : data.toString("utf8");
@@ -683,6 +709,9 @@ export class LighterGateway {
switch (type) { switch (type) {
case "connected": case "connected":
break; break;
case "ping":
this.handlePing(message);
break;
case "subscribed/order_book": case "subscribed/order_book":
this.handleOrderBookSnapshot(message); this.handleOrderBookSnapshot(message);
break; break;
@@ -709,6 +738,28 @@ export class LighterGateway {
} }
} }
private handlePing(message: Record<string, unknown> | null | undefined): void {
const extraPayload: Record<string, unknown> = {};
if (message && typeof message === "object") {
for (const [key, value] of Object.entries(message)) {
if (key === "type") continue;
extraPayload[key] = value;
}
}
this.sendPong(extraPayload);
}
private sendPong(extra: Record<string, unknown> = {}): void {
const ws = this.ws;
if (!ws || ws.readyState !== WebSocket.OPEN) return;
const payload = Object.keys(extra).length ? { ...extra, type: "pong" } : { type: "pong" };
try {
ws.send(JSON.stringify(payload));
} catch (error) {
this.logger("ws:pong", error);
}
}
private handleOrderBookSnapshot(message: any): void { private handleOrderBookSnapshot(message: any): void {
if (!message?.order_book) return; if (!message?.order_book) return;
const incomingOffset = Number(message.offset ?? message.order_book?.offset ?? 0); const incomingOffset = Number(message.offset ?? message.order_book?.offset ?? 0);