diff --git a/src/exchanges/lighter/gateway.ts b/src/exchanges/lighter/gateway.ts index 3447a2a..9ef8a63 100644 --- a/src/exchanges/lighter/gateway.ts +++ b/src/exchanges/lighter/gateway.ts @@ -119,6 +119,7 @@ const KLINE_DEFAULT_COUNT = 120; const DEFAULT_TICKER_POLL_MS = 3000; const DEFAULT_KLINE_POLL_MS = 15000; const WS_HEARTBEAT_INTERVAL_MS = 5_000; +const CLIENT_PING_INTERVAL_MS = 2_000; const WS_STALE_TIMEOUT_MS = 20_000; const FEED_STALE_TIMEOUT_MS = 8_000; const STALE_CHECK_INTERVAL_MS = 2_000; @@ -198,6 +199,7 @@ export class LighterGateway { private readonly wsUrl: string; private connectPromise: Promise | null = null; private heartbeatTimer: ReturnType | null = null; + private pingTimer: ReturnType | null = null; private lastMessageAt = 0; private accountDetails: LighterAccountDetails | null = null; @@ -528,6 +530,7 @@ export class LighterGateway { const cleanup = () => { ws.removeAllListeners(); this.stopHeartbeat(); + this.stopClientPing(); if (this.ws === ws) { this.ws = null; } @@ -541,6 +544,7 @@ export class LighterGateway { try { this.lastMessageAt = Date.now(); this.startHeartbeat(); + this.startClientPing(); await this.subscribeChannels(); this.startStaleMonitor(); settled = true; @@ -640,6 +644,7 @@ export class LighterGateway { this.logger("ws:terminate", error); } this.stopHeartbeat(); + this.stopClientPing(); this.scheduleReconnect(); } @@ -656,6 +661,7 @@ export class LighterGateway { this.logger("ws:terminate", error); } finally { this.stopHeartbeat(); + this.stopClientPing(); this.scheduleReconnect(); } 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 { try { const text = typeof data === "string" ? data : data.toString("utf8"); @@ -683,6 +709,9 @@ export class LighterGateway { switch (type) { case "connected": break; + case "ping": + this.handlePing(message); + break; case "subscribed/order_book": this.handleOrderBookSnapshot(message); break; @@ -709,6 +738,28 @@ export class LighterGateway { } } + private handlePing(message: Record | null | undefined): void { + const extraPayload: Record = {}; + 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 = {}): 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 { if (!message?.order_book) return; const incomingOffset = Number(message.offset ?? message.order_book?.offset ?? 0);