From e3f695e5e00ef8718d8829fbe44fd3d2e36ed236 Mon Sep 17 00:00:00 2001 From: Eva <239388517+100yenadmin@users.noreply.github.com> Date: Wed, 19 Aug 2026 19:02:12 +0700 Subject: [PATCH] fix(tui): heartbeat and bounded reconnect for silent WebSocket drops the client half of the gateway.ping heartbeat contract (#89958); detects a silently-dropped socket via missed ping-acks and reconnects with bounded backoff; part of the #83166 recovery series. --- ui-tui/src/__tests__/gatewayClient.test.ts | 126 +++++++++++++++++- ui-tui/src/gatewayClient.ts | 143 +++++++++++++++++++++ ui-tui/src/gatewayTypes.ts | 3 +- 3 files changed, 270 insertions(+), 2 deletions(-) diff --git a/ui-tui/src/__tests__/gatewayClient.test.ts b/ui-tui/src/__tests__/gatewayClient.test.ts index 2a2384b38b..41c88d56f2 100644 --- a/ui-tui/src/__tests__/gatewayClient.test.ts +++ b/ui-tui/src/__tests__/gatewayClient.test.ts @@ -97,7 +97,7 @@ const { FakeWebSocket } = vi.hoisted(() => { vi.mock('undici', () => ({ WebSocket: FakeWebSocket })) -import { GatewayClient } from '../gatewayClient.js' +import { GatewayClient, RECONNECT_BASE_MS, RECONNECT_MAX_MS, WS_HEARTBEAT_DEAD_MS, WS_HEARTBEAT_INTERVAL_MS } from '../gatewayClient.js' describe('GatewayClient websocket attach mode', () => { const originalWebSocket = globalThis.WebSocket @@ -493,4 +493,128 @@ describe('GatewayClient websocket attach mode', () => { gw.kill() }) + + it('keeps a healthy idle websocket open when heartbeat acknowledgements arrive (issue #32997)', async () => { + vi.useFakeTimers() + process.env.HERMES_TUI_GATEWAY_URL = 'ws://gateway.test/api/ws?token=abc' + const gw = new GatewayClient() + + try { + gw.start() + const socket = FakeWebSocket.instances[0]! + + socket.open() + socket.message( + JSON.stringify({ + jsonrpc: '2.0', + method: 'event', + params: { type: 'gateway.ready', payload: { heartbeat: true } } + }) + ) + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_INTERVAL_MS) + + const heartbeat = JSON.parse(socket.sent.at(-1) ?? '{}') as { id: string; method: string } + + expect(heartbeat.method).toBe('gateway.ping') + socket.message(JSON.stringify({ id: heartbeat.id, jsonrpc: '2.0', result: { ok: true } })) + + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_DEAD_MS + WS_HEARTBEAT_INTERVAL_MS) + expect(socket.readyState).toBe(FakeWebSocket.OPEN) + expect(FakeWebSocket.instances).toHaveLength(1) + } finally { + gw.kill() + vi.useRealTimers() + } + }) + + it('auto-reconnects after a missing heartbeat acknowledgement (issue #32997)', async () => { + vi.useFakeTimers() + process.env.HERMES_TUI_GATEWAY_URL = 'ws://gateway.test/api/ws?token=abc' + const gw = new GatewayClient() + + try { + gw.start() + const first = FakeWebSocket.instances[0]! + + first.open() + first.message( + JSON.stringify({ + jsonrpc: '2.0', + method: 'event', + params: { type: 'gateway.ready', payload: { heartbeat: true } } + }) + ) + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_INTERVAL_MS) + expect(JSON.parse(first.sent.at(-1) ?? '{}')).toMatchObject({ method: 'gateway.ping' }) + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_DEAD_MS + WS_HEARTBEAT_INTERVAL_MS) + await vi.advanceTimersByTimeAsync(RECONNECT_BASE_MS) + expect(FakeWebSocket.instances.length).toBeGreaterThanOrEqual(2) + } finally { + gw.kill() + vi.useRealTimers() + } + }) + + it('does not heartbeat an older backend that omits the capability', async () => { + vi.useFakeTimers() + process.env.HERMES_TUI_GATEWAY_URL = 'ws://gateway.test/api/ws?token=abc' + const gw = new GatewayClient() + + try { + gw.start() + const socket = FakeWebSocket.instances[0]! + + socket.open() + socket.message( + JSON.stringify({ + jsonrpc: '2.0', + method: 'event', + params: { type: 'gateway.ready', payload: {} } + }) + ) + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_DEAD_MS + WS_HEARTBEAT_INTERVAL_MS) + expect(socket.readyState).toBe(FakeWebSocket.OPEN) + expect(socket.sent).toEqual([]) + expect(FakeWebSocket.instances).toHaveLength(1) + } finally { + gw.kill() + vi.useRealTimers() + } + }) + + it('does not double-reconnect when the exit subscriber restarts immediately', async () => { + vi.useFakeTimers() + process.env.HERMES_TUI_GATEWAY_URL = 'ws://gateway.test/api/ws?token=abc' + const gw = new GatewayClient() + + try { + gw.on('exit', () => gw.start()) + gw.start() + const first = FakeWebSocket.instances[0]! + + first.open() + gw.drain() + await Promise.resolve() + first.close(1011) + + expect(FakeWebSocket.instances).toHaveLength(2) + await vi.advanceTimersByTimeAsync(RECONNECT_BASE_MS) + expect(FakeWebSocket.instances).toHaveLength(2) + } finally { + gw.kill() + vi.useRealTimers() + } + }) + + it('does not auto-reconnect after an intentional kill() (issue #32997)', async () => { + vi.useFakeTimers() + process.env.HERMES_TUI_GATEWAY_URL = 'ws://gateway.test/api/ws?token=abc' + const gw = new GatewayClient() + gw.start() + FakeWebSocket.instances[0]!.open() + gw.kill() // sets disposed + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_DEAD_MS + RECONNECT_MAX_MS + 1000) + expect(FakeWebSocket.instances.length).toBe(1) // no reconnect attempted + vi.useRealTimers() + }) }) diff --git a/ui-tui/src/gatewayClient.ts b/ui-tui/src/gatewayClient.ts index 122b5089de..2730df10ee 100644 --- a/ui-tui/src/gatewayClient.ts +++ b/ui-tui/src/gatewayClient.ts @@ -21,6 +21,18 @@ const WS_OPEN = 1 const WS_CLOSING = 2 const WS_CLOSED = 3 +// Keepalive + dead-connection detection. A silent drop (macOS sleep, proxy +// idle timeout, VPN reconnect) kills the TCP socket without a `close` event, +// so the client hangs forever (issue #32997). Browser/undici WebSocket does +// not expose an acknowledged ping/pong API, so this uses a small JSON-RPC +// heartbeat that the TUI gateway explicitly answers. Healthy idle sockets stay +// open; only a missing heartbeat ack forces close -> reconnect. +export const WS_HEARTBEAT_INTERVAL_MS = 15_000 +export const WS_HEARTBEAT_DEAD_MS = 45_000 +// Exponential backoff for reconnect attempts after a transport drop. +export const RECONNECT_BASE_MS = 1_000 +export const RECONNECT_MAX_MS = 30_000 + const getWebSocketCtor = (): typeof WebSocket => typeof WebSocket === 'undefined' ? (UndiciWebSocket as unknown as typeof WebSocket) : WebSocket @@ -149,6 +161,15 @@ export class GatewayClient extends EventEmitter { private drainGeneration = 0 private stdoutRl: ReturnType | null = null private stderrRl: ReturnType | null = null + private heartbeatTimer: ReturnType | null = null + private reconnectTimer: ReturnType | null = null + private reconnectAttempts = 0 + private lastActivityAt = 0 + private heartbeatSeq = 0 + private heartbeatPendingId: string | null = null + private heartbeatSentAt = 0 + // Set on kill() so we never auto-reconnect after an intentional shutdown. + private disposed = false constructor() { super() @@ -165,6 +186,10 @@ export class GatewayClient extends EventEmitter { clearTimeout(this.readyTimer) this.readyTimer = null } + + if (ev.payload?.heartbeat && this.ws?.readyState === WS_OPEN) { + this.startHeartbeat(this.ws) + } } if (this.subscribed) { @@ -210,6 +235,96 @@ export class GatewayClient extends EventEmitter { } } + private startHeartbeat(ws: WebSocket) { + this.stopHeartbeat() + this.lastActivityAt = Date.now() + this.heartbeatPendingId = null + this.heartbeatSentAt = 0 + this.heartbeatTimer = setInterval(() => { + if (this.ws !== ws || ws.readyState !== WS_OPEN) { + return + } + + const now = Date.now() + + if (this.heartbeatPendingId && now - this.heartbeatSentAt > WS_HEARTBEAT_DEAD_MS) { + this.lifecycle('[lifecycle] websocket silent drop detected (heartbeat ack timeout); forcing reconnect') + this.stopHeartbeat() + + try { + ws.close() + } catch { + // ignore + } + + return + } + + if (this.heartbeatPendingId) { + return + } + + const id = `h${++this.heartbeatSeq}` + + this.heartbeatPendingId = id + this.heartbeatSentAt = now + + try { + ws.send(JSON.stringify({ id, jsonrpc: '2.0', method: 'gateway.ping', params: { last_activity_ms: this.lastActivityAt } })) + } catch { + this.lifecycle('[lifecycle] websocket heartbeat send failed; forcing reconnect') + this.stopHeartbeat() + + try { + ws.close() + } catch { + // ignore + } + } + }, WS_HEARTBEAT_INTERVAL_MS) + this.heartbeatTimer.unref?.() + } + + private stopHeartbeat() { + if (this.heartbeatTimer !== null) { + clearInterval(this.heartbeatTimer) + this.heartbeatTimer = null + } + + this.heartbeatPendingId = null + this.heartbeatSentAt = 0 + } + + private scheduleReconnect() { + if (this.disposed || this.reconnectTimer !== null) { + return + } + + const delay = Math.min(RECONNECT_BASE_MS * 2 ** this.reconnectAttempts, RECONNECT_MAX_MS) + this.reconnectAttempts += 1 + this.lifecycle(`[lifecycle] scheduling gateway reconnect in ${delay}ms (attempt ${this.reconnectAttempts})`) + this.publish({ type: 'gateway.reconnecting', payload: { attempt: this.reconnectAttempts, delay_ms: delay } }) + this.reconnectTimer = setTimeout(() => { + this.reconnectTimer = null + + if (this.disposed) { + return + } + + this.start() + }, delay) + this.reconnectTimer.unref?.() + } + + private clearReconnect() { + if (this.reconnectTimer !== null) { + clearTimeout(this.reconnectTimer) + this.reconnectTimer = null + } + + this.reconnectAttempts = 0 + } + private resetStartupState() { // Reject any in-flight RPCs left over from the previous transport // before we swap. Otherwise the old transport's stale exit/close @@ -258,6 +373,14 @@ export class GatewayClient extends EventEmitter { this.lifecycle(`[lifecycle] transport exit code=${code ?? 'null'} reason=${reason ?? 'none'}`) this.rejectPending(new Error(reason || `gateway exited${code === null ? '' : ` (${code})`}`)) + // Self-heal: a dropped transport (real close OR silent drop caught by the + // heartbeat) should reconnect instead of stranding the UI on a dead socket + // (issue #32997). Intentional shutdown sets `disposed` and skips this. + // Schedule before the synchronous 'exit' emission: useMainApp's existing + // recovery subscriber may call start() immediately, and start() cancels this + // timer so there is only one recovery owner. + this.scheduleReconnect() + if (this.subscribed) { this.emit('exit', code) } else { @@ -320,6 +443,7 @@ export class GatewayClient extends EventEmitter { } private handleWebSocketFrame(raw: unknown) { + this.lastActivityAt = Date.now() const text = asWireText(raw) if (!text) { @@ -453,6 +577,8 @@ export class GatewayClient extends EventEmitter { resolve() } + this.lastActivityAt = Date.now() + this.clearReconnect() this.connectSidecarMirror() }, { once: true } @@ -504,6 +630,7 @@ export class GatewayClient extends EventEmitter { } this.pushLog(`[lifecycle] websocket close code=${ev.code}`) + this.stopHeartbeat() this.ws = null this.wsConnectPromise = null this.handleTransportExit(ev.code, `gateway websocket closed${ev.code ? ` (${ev.code})` : ''}`) @@ -521,6 +648,9 @@ export class GatewayClient extends EventEmitter { } start() { + this.disposed = false + this.clearReconnect() + const root = process.env.HERMES_PYTHON_SRC_ROOT ?? resolve(import.meta.dirname, '../../') const attachUrl = resolveGatewayAttachUrl() const sidecarUrl = resolveSidecarUrl() @@ -528,6 +658,8 @@ export class GatewayClient extends EventEmitter { this.attachUrl = attachUrl this.sidecarUrl = sidecarUrl this.resetStartupState() + this.clearReconnect() + this.stopHeartbeat() if (this.proc && !this.proc.killed && this.proc.exitCode === null) { this.lifecycle(`[lifecycle] replacing live gateway child ${describeChild(this.proc)}`) @@ -549,6 +681,14 @@ export class GatewayClient extends EventEmitter { private dispatch(msg: Record) { const id = msg.id as string | undefined + + if (id && id === this.heartbeatPendingId) { + this.heartbeatPendingId = null + this.heartbeatSentAt = 0 + + return + } + const p = id ? this.pending.get(id) : undefined if (p) { @@ -776,6 +916,9 @@ export class GatewayClient extends EventEmitter { } kill(reason = 'requested') { + this.disposed = true + this.clearReconnect() + this.stopHeartbeat() const proc = this.proc const killed = proc?.kill() diff --git a/ui-tui/src/gatewayTypes.ts b/ui-tui/src/gatewayTypes.ts index 427af8bade..8d13e0f3f2 100644 --- a/ui-tui/src/gatewayTypes.ts +++ b/ui-tui/src/gatewayTypes.ts @@ -611,7 +611,7 @@ export interface SpawnTreeLoadResponse { } export type GatewayEvent = - | { payload?: { skin?: GatewaySkin }; session_id?: string; type: 'gateway.ready' } + | { payload?: { heartbeat?: boolean; skin?: GatewaySkin }; session_id?: string; type: 'gateway.ready' } | { payload?: GatewaySkin; session_id?: string; type: 'skin.changed' } | { payload: SessionInfo; session_id?: string; type: 'session.info' } | { payload?: { text?: string }; session_id?: string; type: 'thinking.delta' } @@ -649,6 +649,7 @@ export type GatewayEvent = } | { payload?: { reason?: string }; session_id?: string; type: 'dashboard.new_session_requested' } | { payload: { line: string }; session_id?: string; type: 'gateway.stderr' } + | { payload?: { attempt?: number; delay_ms?: number }; session_id?: string; type: 'gateway.reconnecting' } | { payload?: { level?: 'info' | 'warn' | 'error'; message?: string } session_id?: string