diff --git a/apps/desktop/src/lib/gateway-rpc.test.ts b/apps/desktop/src/lib/gateway-rpc.test.ts index 6c84c12eca..aa8a1ef9e2 100644 --- a/apps/desktop/src/lib/gateway-rpc.test.ts +++ b/apps/desktop/src/lib/gateway-rpc.test.ts @@ -1,9 +1,16 @@ +import { JsonRpcGatewayError } from '@hermes/shared' import { describe, expect, it } from 'vitest' import { isMissingPendingPromptRequest, isMissingRpcMethod } from './gateway-rpc' describe('isMissingRpcMethod', () => { - it('detects JSON-RPC method-not-found errors', () => { + it('trusts the JSON-RPC code over the message when the frame survived', () => { + expect(isMissingRpcMethod(new JsonRpcGatewayError('unknown method: projects.create', { code: -32601 }))).toBe(true) + // A tool result that merely mentions the phrase must not read as a capability verdict. + expect(isMissingRpcMethod(new JsonRpcGatewayError('unknown method in user script', { code: -32000 }))).toBe(false) + }) + + it('falls back to the message for codeless (IPC-flattened) errors', () => { expect(isMissingRpcMethod(new Error('unknown method: projects.create'))).toBe(true) expect(isMissingRpcMethod(new Error('Method not found'))).toBe(true) expect(isMissingRpcMethod(new Error('RPC failed: -32601'))).toBe(true) diff --git a/apps/desktop/src/lib/gateway-rpc.ts b/apps/desktop/src/lib/gateway-rpc.ts index 1b52ef1f33..14c2d8a2dd 100644 --- a/apps/desktop/src/lib/gateway-rpc.ts +++ b/apps/desktop/src/lib/gateway-rpc.ts @@ -1,5 +1,16 @@ -/** True when a JSON-RPC call failed because the backend predates the method. */ +import { JSON_RPC_METHOD_NOT_FOUND } from '@hermes/shared' + +/** True when a JSON-RPC call failed because the backend predates the method. + * The gateway answers -32601 (`tui_gateway/server.py::dispatch`) and the + * shared client keeps that code on the error; the message match is only for + * errors that lost their frame across the IPC bridge or a wrapped rethrow. */ export function isMissingRpcMethod(error: unknown): boolean { + const code = error && typeof error === 'object' ? (error as { code?: unknown }).code : undefined + + if (typeof code === 'number') { + return code === JSON_RPC_METHOD_NOT_FOUND + } + const message = error instanceof Error ? error.message : String(error) return /method not found|-32601|unknown method|no such method/i.test(message) diff --git a/apps/shared/package.json b/apps/shared/package.json index 9cc9a92d4e..c9b8318ac6 100644 --- a/apps/shared/package.json +++ b/apps/shared/package.json @@ -9,7 +9,9 @@ "./billing-policy": "./src/billing-policy.ts", "./charge-settlement": "./src/charge-settlement.ts", "./gateway-events": "./src/gateway-events.ts", - "./skin": "./src/skin.ts" + "./skin": "./src/skin.ts", + "./json-rpc-channel": "./src/json-rpc-channel.ts", + "./reconnect-backoff": "./src/reconnect-backoff.ts" }, "types": "./src/index.ts", "scripts": { diff --git a/apps/shared/src/index.ts b/apps/shared/src/index.ts index 3d82d72dad..539a3ebbd9 100644 --- a/apps/shared/src/index.ts +++ b/apps/shared/src/index.ts @@ -89,14 +89,25 @@ export { type WakeDetectedPayload } from './gateway-events' export { - type ConnectionState, - type GatewayClientOptions, + DEFAULT_HEARTBEAT_DEADLINE_MS, + DEFAULT_HEARTBEAT_INTERVAL_MS, type GatewayRequestId, - isGatewayWebSocketUrl, + JSON_RPC_METHOD_NOT_FOUND, + jsonRpcErrorFromFrame, type JsonRpcErrorPayload, type JsonRpcFrame, - JsonRpcGatewayClient, JsonRpcGatewayError, + JsonRpcRequestChannel, + type JsonRpcRequestChannelOptions, + type JsonRpcTransport, + wireFrameText +} from './json-rpc-channel' +export { + type ConnectionState, + type GatewayClientOptions, + GatewayEventHub, + isGatewayWebSocketUrl, + JsonRpcGatewayClient, type WebSocketLike } from './json-rpc-gateway' export { skillInvocationText } from './skill-scaffold' diff --git a/apps/shared/src/json-rpc-channel.test.ts b/apps/shared/src/json-rpc-channel.test.ts new file mode 100644 index 0000000000..2e8e7594f7 --- /dev/null +++ b/apps/shared/src/json-rpc-channel.test.ts @@ -0,0 +1,106 @@ +import { describe, expect, it, vi } from 'vitest' + +import { JsonRpcGatewayError, JsonRpcRequestChannel, type JsonRpcTransport } from './json-rpc-channel.js' + +const spyTransport = () => { + const sent: string[] = [] + const transport: JsonRpcTransport = { send: text => void sent.push(text) } + + return { sent, transport, last: () => JSON.parse(sent.at(-1) ?? '{}') as { id: string; method: string } } +} + +describe('JsonRpcRequestChannel', () => { + it('routes responses to the pending call and keeps JSON-RPC code/data on errors', async () => { + const events: string[] = [] + const channel = new JsonRpcRequestChannel({ onEvent: ev => void events.push(ev.type), requestIdPrefix: 'x' }) + const { transport, last } = spyTransport() + + channel.attach(transport) + + const ok = channel.request<{ ok: boolean }>('session.create', { cols: 80 }) + const okId = last().id + const failing = channel.request('projects.create') + const failId = last().id + + expect(okId).toBe('x1') + channel.handleFrame(JSON.stringify({ id: okId, jsonrpc: '2.0', result: { ok: true } })) + channel.handleFrame( + JSON.stringify({ + error: { code: -32601, data: { method: 'projects.create' }, message: 'unknown method: projects.create' }, + id: failId, + jsonrpc: '2.0' + }) + ) + channel.handleFrame(JSON.stringify({ jsonrpc: '2.0', method: 'event', params: { type: 'session.info', payload: {} } })) + // Non-JSON and unknown ids are ignored, never thrown. + expect(channel.handleFrame('not json')).toBeNull() + channel.handleFrame(JSON.stringify({ id: 'never-sent', jsonrpc: '2.0', result: 1 })) + + await expect(ok).resolves.toEqual({ ok: true }) + const error = (await failing.catch((e: unknown) => e)) as JsonRpcGatewayError + + expect(error).toBeInstanceOf(JsonRpcGatewayError) + expect(error.code).toBe(-32601) + expect(error.data).toEqual({ method: 'projects.create' }) + expect(events).toEqual(['session.info']) + }) + + it('detach fails every in-flight call and a per-call timeout names the method', async () => { + vi.useFakeTimers() + + try { + const channel = new JsonRpcRequestChannel({ requestTimeoutMs: 60_000 }) + const { transport } = spyTransport() + + channel.attach(transport) + + const slow = expect(channel.request('a.slow', {}, 1_000)).rejects.toThrow('request timed out after 1s: a.slow') + const untilDetach = channel.request('b.wait') + + await vi.advanceTimersByTimeAsync(1_000) + await slow + + channel.detach(new Error('gateway exited (1)')) + await expect(untilDetach).rejects.toThrow('gateway exited (1)') + await expect(channel.request('c.after')).rejects.toThrow('gateway not connected') + } finally { + vi.useRealTimers() + } + }) + + it('heartbeat pings while frames keep arriving and reports a silent transport', async () => { + vi.useFakeTimers() + + try { + const failures: string[] = [] + + const channel = new JsonRpcRequestChannel({ + heartbeatDeadlineMs: 300, + heartbeatIntervalMs: 100, + onHeartbeatFailure: e => void failures.push(e.message) + }) + + const { sent, transport } = spyTransport() + + channel.attach(transport) + channel.startHeartbeat() + + for (let i = 0; i < 6; i++) { + await vi.advanceTimersByTimeAsync(100) + channel.handleFrame(JSON.stringify({ jsonrpc: '2.0', method: 'event', params: { type: 'status.update' } })) + } + + expect(sent.filter(f => f.includes('gateway.ping')).length).toBeGreaterThanOrEqual(5) + expect(failures).toEqual([]) + + await vi.advanceTimersByTimeAsync(400) + expect(failures).toEqual(['WebSocket heartbeat acknowledgement timed out']) + // Failure stops the timer: no further pings after the report. + const pings = sent.length + await vi.advanceTimersByTimeAsync(500) + expect(sent.length).toBe(pings) + } finally { + vi.useRealTimers() + } + }) +}) diff --git a/apps/shared/src/json-rpc-channel.ts b/apps/shared/src/json-rpc-channel.ts new file mode 100644 index 0000000000..0dab532fc2 --- /dev/null +++ b/apps/shared/src/json-rpc-channel.ts @@ -0,0 +1,356 @@ +import type { GatewayEvent } from './gateway-events.js' + +export type GatewayRequestId = number | string + +export interface JsonRpcErrorPayload { + code?: number + data?: unknown + message?: string +} + +export interface JsonRpcFrame { + error?: JsonRpcErrorPayload + id?: GatewayRequestId | null + method?: string + params?: GatewayEvent + result?: unknown +} + +/** JSON-RPC error with optional structured `data` from the gateway. */ +export class JsonRpcGatewayError extends Error { + readonly code?: number + readonly data?: unknown + + constructor(message: string, options?: { code?: number; data?: unknown }) { + super(message) + this.name = 'JsonRpcGatewayError' + this.code = options?.code + this.data = options?.data + } +} + +/** JSON-RPC "method not found" (tui_gateway/server.py::dispatch `_err(rid, -32601, …)`). */ +export const JSON_RPC_METHOD_NOT_FOUND = -32601 + +/** Map a raw `error` member of a response frame to the typed error every surface inspects. */ +export function jsonRpcErrorFromFrame(raw: unknown, fallbackMessage = 'Hermes RPC failed'): JsonRpcGatewayError { + const err = (raw && typeof raw === 'object' ? raw : {}) as JsonRpcErrorPayload + + return new JsonRpcGatewayError(typeof err.message === 'string' && err.message ? err.message : fallbackMessage, { + code: typeof err.code === 'number' ? err.code : undefined, + data: err.data + }) +} + +/** + * Anything that can carry one serialized JSON-RPC frame to the gateway. The + * channel never learns whether that is a WebSocket, a child's stdin, or a + * test spy; the owner feeds inbound text back through `handleFrame`. + */ +export interface JsonRpcTransport { + send(text: string): void +} + +export interface JsonRpcRequestChannelOptions { + createRequestId?: (nextId: number) => GatewayRequestId + heartbeatDeadlineMs?: number + heartbeatIntervalMs?: number + /** Called when the heartbeat deadline passes or a heartbeat send throws; the owner drops the transport. */ + onHeartbeatFailure?: (error: Error) => void + /** Decoded `event` notification. */ + onEvent?: (event: GatewayEvent) => void + requestIdPrefix?: string + requestTimeoutMs?: number + /** `setTimeout`/`setInterval` handles are `unref`'d when the runtime supports it (Node) so a pending call cannot pin the process. */ + unrefTimers?: boolean +} + +interface PendingCall { + reject: (error: Error) => void + resolve: (value: unknown) => void + timer?: ReturnType +} + +const DEFAULT_REQUEST_TIMEOUT_MS = 120_000 +// 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. +export const DEFAULT_HEARTBEAT_INTERVAL_MS = 15_000 +export const DEFAULT_HEARTBEAT_DEADLINE_MS = 45_000 + +// Hoisted decoder: attach mode can drive high-frequency binary frames (tool +// deltas, reasoning streams) and a fresh TextDecoder per message is avoidable +// GC pressure; UTF-8 is stateless and frames arrive whole. +const wireDecoder = new TextDecoder() + +/** Decode a socket `message.data` (string / ArrayBuffer / view) to text; `null` for anything else. */ +export function wireFrameText(raw: unknown): string | null { + if (typeof raw === 'string') { + return raw + } + + if (raw instanceof ArrayBuffer || ArrayBuffer.isView(raw)) { + return wireDecoder.decode(raw as ArrayBuffer) + } + + return null +} + +const unrefTimer = (timer: unknown) => { + ;(timer as { unref?: () => void } | undefined)?.unref?.() +} + +/** + * The transport-agnostic half of a JSON-RPC gateway connection: request ids, + * the pending map with per-call timeouts and AbortSignal, response → typed + * error mapping, event-notification decoding, and the `gateway.ping` + * heartbeat. Owners (`JsonRpcGatewayClient` over WebSocket, the Ink TUI over + * stdio / an attached socket) supply a `JsonRpcTransport` per connection + * generation and call `handleFrame` for every inbound text frame. + */ +export class JsonRpcRequestChannel { + private nextId = 0 + private readonly pending = new Map() + private transport: JsonRpcTransport | null = null + private heartbeatTimer: ReturnType | null = null + private heartbeatSequence = 0 + private lastInboundAt = 0 + private readonly options: Required> & + Pick + + constructor(options: JsonRpcRequestChannelOptions = {}) { + this.options = { + createRequestId: options.createRequestId ?? ((nextId: number) => `${options.requestIdPrefix ?? 'r'}${nextId}`), + heartbeatDeadlineMs: options.heartbeatDeadlineMs ?? DEFAULT_HEARTBEAT_DEADLINE_MS, + heartbeatIntervalMs: options.heartbeatIntervalMs ?? DEFAULT_HEARTBEAT_INTERVAL_MS, + onEvent: options.onEvent, + onHeartbeatFailure: options.onHeartbeatFailure, + requestIdPrefix: options.requestIdPrefix ?? 'r', + requestTimeoutMs: options.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS, + unrefTimers: options.unrefTimers ?? false + } + } + + get defaultRequestTimeoutMs(): number { + return this.options.requestTimeoutMs + } + + get connected(): boolean { + return this.transport !== null + } + + /** Bind a new connection generation. Any previous generation's heartbeat stops; its pending calls are the owner's to reject. */ + attach(transport: JsonRpcTransport): void { + this.stopHeartbeat() + this.transport = transport + this.lastInboundAt = Date.now() + } + + /** Drop the transport and fail every in-flight call with `error`. */ + detach(error: Error): void { + this.stopHeartbeat() + this.transport = null + this.rejectAllPending(error) + } + + /** True while `transport` is the bound generation (owners gate stale socket callbacks on this). */ + owns(transport: JsonRpcTransport): boolean { + return this.transport === transport + } + + request( + method: string, + params: Record = {}, + timeoutMs = this.options.requestTimeoutMs, + signal?: AbortSignal, + notConnectedError: () => Error = () => new Error('gateway not connected') + ): Promise { + const transport = this.transport + + if (!transport) { + return Promise.reject(notConnectedError()) + } + + if (signal?.aborted) { + return Promise.reject(new DOMException('Aborted', 'AbortError')) + } + + const id = this.options.createRequestId(++this.nextId) + + return new Promise((resolve, reject) => { + let onAbort: (() => void) | undefined + + const detachAbort = () => { + if (onAbort && signal) { + signal.removeEventListener('abort', onAbort) + } + } + + const pending: PendingCall = { + resolve: value => { + detachAbort() + resolve(value as T) + }, + reject: error => { + detachAbort() + reject(error) + } + } + + if (timeoutMs > 0) { + pending.timer = setTimeout(() => { + if (this.pending.delete(id)) { + detachAbort() + // Include the configured timeout so a caller (or a user looking + // at an error toast) can tell whether the default window fired + // or a per-call override — e.g. /compress opts into 120s. + const seconds = Math.round(timeoutMs / 1000) + reject(new Error(`request timed out after ${seconds}s: ${method}`)) + } + }, timeoutMs) + + if (this.options.unrefTimers) { + unrefTimer(pending.timer) + } + } + + // Abort drops the pending call immediately (no dangling resolver/timer); + // server-side cancellation is a separate cooperative RPC where it matters. + if (signal) { + onAbort = () => { + this.clearPending(id) + detachAbort() + reject(new DOMException('Aborted', 'AbortError')) + } + + signal.addEventListener('abort', onAbort, { once: true }) + } + + this.pending.set(id, pending) + + try { + transport.send(JSON.stringify({ jsonrpc: '2.0', id, method, params })) + } catch (error) { + this.clearPending(id) + detachAbort() + reject(error instanceof Error ? error : new Error(String(error))) + } + }) + } + + /** + * Route one inbound frame: a response settles its pending call, an + * `event` notification reaches `onEvent`. Returns the decoded frame so the + * owner can act on it too (mirror it, record seq, …) or `null` when the + * text was not JSON. + */ + handleFrame(text: string): JsonRpcFrame | null { + this.lastInboundAt = Date.now() + + let frame: JsonRpcFrame + + try { + frame = JSON.parse(text) as JsonRpcFrame + } catch { + return null + } + + if (frame.id !== undefined && frame.id !== null) { + const call = this.pending.get(frame.id) + + if (call) { + this.clearPending(frame.id) + + if (frame.error) { + call.reject(jsonRpcErrorFromFrame(frame.error)) + } else { + call.resolve(frame.result) + } + } + + return frame + } + + if (frame.method === 'event' && frame.params && typeof frame.params.type === 'string') { + this.options.onEvent?.(frame.params) + } + + return frame + } + + /** + * Begin the `gateway.ping` keepalive on the bound transport. Only call when + * `gateway.ready.heartbeat` advertised support — an older backend would + * answer with -32601 and never count as alive. Any inbound frame counts as + * liveness; only a full deadline of silence drops the transport. + */ + startHeartbeat(): void { + this.stopHeartbeat() + this.lastInboundAt = Date.now() + + const transport = this.transport + + if (!transport || this.options.heartbeatIntervalMs <= 0 || this.options.heartbeatDeadlineMs <= 0) { + return + } + + this.heartbeatTimer = setInterval(() => { + if (this.transport !== transport) { + return + } + + if (Date.now() - this.lastInboundAt >= this.options.heartbeatDeadlineMs) { + this.failHeartbeat(new Error('WebSocket heartbeat acknowledgement timed out')) + + return + } + + try { + transport.send( + JSON.stringify({ jsonrpc: '2.0', id: `heartbeat-${++this.heartbeatSequence}`, method: 'gateway.ping', params: {} }) + ) + } catch (error) { + this.failHeartbeat(error instanceof Error ? error : new Error(String(error))) + } + }, this.options.heartbeatIntervalMs) + + if (this.options.unrefTimers) { + unrefTimer(this.heartbeatTimer) + } + } + + stopHeartbeat(): void { + if (this.heartbeatTimer !== null) { + clearInterval(this.heartbeatTimer) + this.heartbeatTimer = null + } + } + + private failHeartbeat(error: Error): void { + this.stopHeartbeat() + this.options.onHeartbeatFailure?.(error) + } + + private clearPending(id: GatewayRequestId): void { + const call = this.pending.get(id) + + if (call?.timer) { + clearTimeout(call.timer) + } + + this.pending.delete(id) + } + + private rejectAllPending(error: Error): void { + for (const [id, call] of this.pending) { + if (call.timer) { + clearTimeout(call.timer) + } + + this.pending.delete(id) + call.reject(error) + } + } +} diff --git a/apps/shared/src/json-rpc-gateway.ts b/apps/shared/src/json-rpc-gateway.ts index cff8f15691..65696c142f 100644 --- a/apps/shared/src/json-rpc-gateway.ts +++ b/apps/shared/src/json-rpc-gateway.ts @@ -1,45 +1,18 @@ import type { GatewayEvent, GatewayEventName } from './gateway-events.js' +import { + DEFAULT_HEARTBEAT_DEADLINE_MS, + DEFAULT_HEARTBEAT_INTERVAL_MS, + type GatewayRequestId, + JsonRpcRequestChannel, + type JsonRpcTransport, + wireFrameText +} from './json-rpc-channel.js' export type { GatewayEvent, GatewayEventName } from './gateway-events.js' - export type ConnectionState = 'idle' | 'connecting' | 'open' | 'closed' | 'error' -export type GatewayRequestId = number | string - -export interface JsonRpcErrorPayload { - code?: number - data?: unknown - message?: string -} - -export interface JsonRpcFrame { - error?: JsonRpcErrorPayload - id?: GatewayRequestId | null - method?: string - params?: GatewayEvent - result?: unknown -} - -/** JSON-RPC error with optional structured `data` from the gateway. */ -export class JsonRpcGatewayError extends Error { - readonly code?: number - readonly data?: unknown - - constructor(message: string, options?: { code?: number; data?: unknown }) { - super(message) - this.name = 'JsonRpcGatewayError' - this.code = options?.code - this.data = options?.data - } -} export type WebSocketLike = WebSocket -type PendingCall = { - reject: (error: Error) => void - resolve: (value: unknown) => void - timer?: ReturnType -} - export interface GatewayClientOptions { closedErrorMessage?: string connectErrorMessage?: string @@ -48,7 +21,9 @@ export interface GatewayClientOptions { heartbeatDeadlineMs?: number heartbeatIntervalMs?: number /** Return true to intercept the default closed-state transition. */ - onSocketClose?: (event: CloseEvent) => boolean | void + onSocketClose?: (event: { code: number }) => boolean | void + /** Fetch `session.events.since` after a reconnect (default). Off for notification-only feeds whose peer never answers RPCs. */ + replay?: boolean requestIdPrefix?: string requestTimeoutMs?: number socketFactory?: (url: string) => WebSocketLike @@ -56,14 +31,12 @@ export interface GatewayClientOptions { } const ANY = '*' +const DEFAULT_REQUEST_TIMEOUT_MS = 120_000 const isGatewayReady = (event: GatewayEvent): event is GatewayEvent<'gateway.ready'> => event.type === 'gateway.ready' -const DEFAULT_REQUEST_TIMEOUT_MS = 120_000 // Replay fetch after reconnect: bounded so a wedged backend can't hold the // guard open; generous enough for a 512-frame ring to drain. const REPLAY_REQUEST_TIMEOUT_MS = 10_000 -const DEFAULT_HEARTBEAT_INTERVAL_MS = 15_000 -const DEFAULT_HEARTBEAT_DEADLINE_MS = 45_000 // A reconnect after sleep/wake must not hang forever in 'connecting' (which // keeps the composer disabled and stuck on "Starting Hermes..."). If the open // handshake doesn't land in this window, fail to 'error' so callers can retry. @@ -82,14 +55,55 @@ export function isGatewayWebSocketUrl(value: unknown): value is string { } } +/** + * Typed fan-out of gateway `event` notifications: per-type handlers plus a + * `*` wildcard. Shared by the WebSocket client below and the Ink TUI's stdio + * client so both dispatch the same way. + */ +export class GatewayEventHub { + private readonly handlers = new Map void>>() + + on(type: K, handler: (event: GatewayEvent) => void): () => void { + let set = this.handlers.get(type) + + if (!set) { + set = new Set() + this.handlers.set(type, set) + } + + set.add(handler as (event: GatewayEvent) => void) + + return () => set?.delete(handler as (event: GatewayEvent) => void) + } + + onAny(handler: (event: GatewayEvent) => void): () => void { + // ANY is a client-side wildcard, not a wire name; it never reaches the typed map. + return this.on(ANY as GatewayEventName, handler as (event: GatewayEvent) => void) + } + + dispatch(event: GatewayEvent): void { + for (const handler of this.handlers.get(event.type) ?? []) { + handler(event) + } + + for (const handler of this.handlers.get(ANY) ?? []) { + handler(event) + } + } +} + +/** + * Bring a `JsonRpcRequestChannel` to a raw text sink — a WebSocket here, a + * child's stdin in the TUI. Kept separate from the socket so the channel never + * holds a reference to a specific socket generation. + */ +const socketTransport = (socket: WebSocketLike): JsonRpcTransport => ({ send: text => socket.send(text) }) + export class JsonRpcGatewayClient { - private nextId = 0 - private pending = new Map() private socket: WebSocketLike | null = null private state: ConnectionState = 'idle' - private heartbeatTimer: ReturnType | null = null - private heartbeatSequence = 0 - private lastInboundAt = 0 + private readonly channel: JsonRpcRequestChannel + private readonly events = new GatewayEventHub() /** Last observed event seq per session_id — drives lossless reconnect replay. */ private lastSeenSeq = new Map() /** Set while a post-reconnect replay fetch is in flight (dedup guard). */ @@ -110,7 +124,6 @@ export class JsonRpcGatewayClient { * silently believe nothing was missed. */ private replayEpoch: string | null = null - private readonly eventHandlers = new Map void>>() private readonly stateHandlers = new Set<(state: ConnectionState) => void>() private readonly options: Required> & Pick @@ -125,10 +138,19 @@ export class JsonRpcGatewayClient { heartbeatIntervalMs: options.heartbeatIntervalMs ?? DEFAULT_HEARTBEAT_INTERVAL_MS, notConnectedErrorMessage: options.notConnectedErrorMessage ?? 'gateway not connected', onSocketClose: options.onSocketClose ?? (() => false), + replay: options.replay ?? true, requestIdPrefix: options.requestIdPrefix ?? 'r', requestTimeoutMs: options.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS, socketFactory: options.socketFactory } + this.channel = new JsonRpcRequestChannel({ + createRequestId: this.options.createRequestId, + heartbeatDeadlineMs: this.options.heartbeatDeadlineMs, + heartbeatIntervalMs: this.options.heartbeatIntervalMs, + onEvent: event => this.handleEvent(event), + onHeartbeatFailure: error => this.invalidate(error.message), + requestTimeoutMs: this.options.requestTimeoutMs + }) } get connectionState(): ConnectionState { @@ -148,23 +170,27 @@ export class JsonRpcGatewayClient { throw invalidUrl() } - if (this.socket?.readyState === WebSocket.OPEN || this.state === 'connecting') { + if ((this.socket && this.socket.readyState === WebSocket.OPEN) || this.state === 'connecting') { return } this.setState('connecting') const socket = this.options.socketFactory?.(wsUrl) ?? new WebSocket(wsUrl) + const transport = socketTransport(socket) this.socket = socket - this.stopHeartbeat() + this.channel.stopHeartbeat() socket.addEventListener('message', message => { if (this.socket !== socket) { return } - this.lastInboundAt = Date.now() - this.handleMessage(message.data) + const text = wireFrameText(message.data) + + if (text !== null) { + this.channel.handleFrame(text) + } }) socket.addEventListener('close', event => { @@ -176,10 +202,7 @@ export class JsonRpcGatewayClient { return } - this.socket = null - this.stopHeartbeat() - this.setState('closed') - this.rejectAllPending(new Error(this.options.closedErrorMessage)) + this.dropSocket(new Error(this.options.closedErrorMessage)) }) await new Promise((resolve, reject) => { @@ -193,6 +216,7 @@ export class JsonRpcGatewayClient { socket.removeEventListener('open', onOpen) socket.removeEventListener('error', onError) + socket.removeEventListener('close', onClose) } const onOpen = () => { @@ -202,6 +226,7 @@ export class JsonRpcGatewayClient { settled = true cleanup() + this.channel.attach(transport) this.setState('open') resolve() // Lossless resume: drain events emitted while we were disconnected. @@ -221,8 +246,30 @@ export class JsonRpcGatewayClient { reject(new Error(this.options.connectErrorMessage)) } + // A server that closes during the handshake (auth gate, 4401/4403) + // may never fire `error`; without this the caller waits out the + // connect timeout for a verdict the socket already delivered. The + // permanent close listener above has already moved the generation to + // 'closed' (unless onSocketClose intercepted), so only settle here. + const onClose = () => { + if (settled) { + return + } + + settled = true + cleanup() + + if (this.socket === socket) { + this.socket = null + this.setState('error') + } + + reject(new Error(this.options.connectErrorMessage)) + } + socket.addEventListener('open', onOpen, { once: true }) socket.addEventListener('error', onError, { once: true }) + socket.addEventListener('close', onClose, { once: true }) if (this.options.connectTimeoutMs > 0) { timer = setTimeout(() => { @@ -253,20 +300,7 @@ export class JsonRpcGatewayClient { } close(): void { - const socket = this.socket - - if (!socket) { - return - } - - try { - socket.close() - } finally { - this.socket = null - this.stopHeartbeat() - this.setState('closed') - this.rejectAllPending(new Error(this.options.closedErrorMessage)) - } + this.invalidate() } /** @@ -280,25 +314,24 @@ export class JsonRpcGatewayClient { return } - this.invalidateSocket(socket, new Error(message)) + // Drop the generation BEFORE closing: a synchronous `close` event from + // the socket must hit the identity guard and not run the default + // closed-path a second time on top of whatever the owner redialed. + this.dropSocket(new Error(message)) + + try { + socket.close() + } catch { + // The generation was already invalidated; the reconnect owner can redial. + } } on(type: K, handler: (event: GatewayEvent) => void): () => void { - let handlers = this.eventHandlers.get(type) - - if (!handlers) { - handlers = new Set() - this.eventHandlers.set(type, handlers) - } - - handlers.add(handler as (event: GatewayEvent) => void) - - return () => handlers?.delete(handler as (event: GatewayEvent) => void) + return this.events.on(type, handler) } onAny(handler: (event: GatewayEvent) => void): () => void { - // ANY is a client-side wildcard, not a wire name; it never reaches the typed map. - return this.on(ANY as GatewayEventName, handler as (event: GatewayEvent) => void) + return this.events.onAny(handler) } onEvent(handler: (event: GatewayEvent) => void): () => void { @@ -318,152 +351,39 @@ export class JsonRpcGatewayClient { timeoutMs = this.options.requestTimeoutMs, signal?: AbortSignal ): Promise { - const socket = this.socket - - if (!socket || socket.readyState !== WebSocket.OPEN) { + if (!this.socket || this.socket.readyState !== WebSocket.OPEN) { return Promise.reject(new Error(this.options.notConnectedErrorMessage)) } - if (signal?.aborted) { - return Promise.reject(new DOMException('Aborted', 'AbortError')) - } - - const id = this.options.createRequestId(++this.nextId) - - return new Promise((resolve, reject) => { - let onAbort: (() => void) | undefined - - const detach = () => { - if (onAbort && signal) { - signal.removeEventListener('abort', onAbort) - } - } - - const pending: PendingCall = { - resolve: value => { - detach() - resolve(value as T) - }, - reject: error => { - detach() - reject(error) - } - } - - if (timeoutMs > 0) { - pending.timer = setTimeout(() => { - if (this.pending.delete(id)) { - detach() - // Include the configured timeout so a caller (or a user looking - // at an error toast) can tell whether the default 30s window - // fired or a per-call override — e.g. /compress opts into 120s. - const seconds = Math.round(timeoutMs / 1000) - reject(new Error(`request timed out after ${seconds}s: ${method}`)) - } - }, timeoutMs) - } - - // Abort drops the pending call immediately (no dangling resolver/timer); - // server-side cancellation is a separate cooperative RPC where it matters. - if (signal) { - onAbort = () => { - const call = this.pending.get(id) - - if (call?.timer) { - clearTimeout(call.timer) - } - - this.pending.delete(id) - detach() - reject(new DOMException('Aborted', 'AbortError')) - } - - signal.addEventListener('abort', onAbort, { once: true }) - } - - this.pending.set(id, pending) - - try { - socket.send( - JSON.stringify({ - jsonrpc: '2.0', - id, - method, - params - }) - ) - } catch (error) { - this.clearPending(id) - detach() - reject(error instanceof Error ? error : new Error(String(error))) - } - }) + return this.channel.request(method, params, timeoutMs, signal, () => new Error(this.options.notConnectedErrorMessage)) } - private handleMessage(raw: unknown): void { - const text = typeof raw === 'string' ? raw : String(raw) - let frame: JsonRpcFrame + private handleEvent(event: GatewayEvent): void { + if (isGatewayReady(event)) { + if (event.payload?.heartbeat === true) { + this.channel.startHeartbeat() + } - try { - frame = JSON.parse(text) as JsonRpcFrame - } catch { - return + const epoch = event.payload?.replay_epoch + + if (typeof epoch === 'string' && epoch) { + this.adoptReplayEpoch(epoch) + } } - if (frame.id !== undefined && frame.id !== null) { - const call = this.pending.get(frame.id) + const sid = event.session_id + const seqValue = event.seq - if (!call) { - return - } - - this.clearPending(frame.id) - - if (frame.error) { - call.reject( - new JsonRpcGatewayError(frame.error.message || 'Hermes RPC failed', { - code: typeof frame.error.code === 'number' ? frame.error.code : undefined, - data: frame.error.data - }) - ) - } else { - call.resolve(frame.result) - } + if (this.replayHold && sid && typeof seqValue === 'number' && this.replayHold.has(sid)) { + // Replay in flight for this session: park the frame; flushReplayHold + // dispatches it after the replayed gap, gated on seq. + this.replayHold.get(sid)?.push(event) return } - if (frame.method === 'event' && frame.params?.type) { - if (isGatewayReady(frame.params)) { - if (frame.params.payload?.heartbeat === true) { - const socket = this.socket - - if (socket) { - this.startHeartbeat(socket) - } - } - - const epoch = frame.params.payload?.replay_epoch - - if (typeof epoch === 'string' && epoch) { - this.adoptReplayEpoch(epoch) - } - } - - const sid = frame.params.session_id - const seqValue = frame.params.seq - - if (this.replayHold && sid && typeof seqValue === 'number' && this.replayHold.has(sid)) { - // Replay in flight for this session: park the frame; flushReplayHold - // dispatches it after the replayed gap, gated on seq. - this.replayHold.get(sid)?.push(frame.params) - - return - } - - this.recordSeq(frame.params) - this.dispatchEvent(frame.params) - } + this.recordSeq(event) + this.dispatchEvent(event) } /** @@ -498,7 +418,7 @@ export class JsonRpcGatewayClient { * Best-effort: failures are swallowed (the next reconnect retries). */ private async fetchReplay(): Promise { - if (this.replayInFlight || this.lastSeenSeq.size === 0) { + if (!this.options.replay || this.replayInFlight || this.lastSeenSeq.size === 0) { return } @@ -618,94 +538,15 @@ export class JsonRpcGatewayClient { } } - private startHeartbeat(socket: WebSocketLike): void { - this.stopHeartbeat() - this.lastInboundAt = Date.now() - - if (this.options.heartbeatIntervalMs <= 0 || this.options.heartbeatDeadlineMs <= 0) { - return - } - - this.heartbeatTimer = setInterval(() => { - if (this.socket !== socket || socket.readyState !== WebSocket.OPEN) { - return - } - - if (Date.now() - this.lastInboundAt >= this.options.heartbeatDeadlineMs) { - this.invalidateSocket(socket, new Error('WebSocket heartbeat acknowledgement timed out')) - - return - } - - try { - socket.send( - JSON.stringify({ - jsonrpc: '2.0', - id: `heartbeat-${++this.heartbeatSequence}`, - method: 'gateway.ping', - params: {} - }) - ) - } catch (error) { - this.invalidateSocket(socket, error instanceof Error ? error : new Error(String(error))) - } - }, this.options.heartbeatIntervalMs) - } - - private stopHeartbeat(): void { - if (this.heartbeatTimer !== null) { - clearInterval(this.heartbeatTimer) - this.heartbeatTimer = null - } - } - - private invalidateSocket(socket: WebSocketLike, error: Error): void { - if (this.socket !== socket) { - return - } - + /** Forget the current socket generation, fail its calls, and go 'closed'. */ + private dropSocket(error: Error): void { this.socket = null - this.stopHeartbeat() - - try { - socket.close() - } catch { - // The generation was already invalidated; the reconnect owner can redial. - } - + this.channel.detach(error) this.setState('closed') - this.rejectAllPending(error) - } - - private clearPending(id: GatewayRequestId): void { - const call = this.pending.get(id) - - if (call?.timer) { - clearTimeout(call.timer) - } - - this.pending.delete(id) } private dispatchEvent(event: GatewayEvent): void { - for (const handler of this.eventHandlers.get(event.type) ?? []) { - handler(event) - } - - for (const handler of this.eventHandlers.get(ANY) ?? []) { - handler(event) - } - } - - private rejectAllPending(error: Error): void { - for (const [id, call] of this.pending) { - if (call.timer) { - clearTimeout(call.timer) - } - - call.reject(error) - this.pending.delete(id) - } + this.events.dispatch(event) } private setState(state: ConnectionState): void { diff --git a/ui-tui/src/__tests__/gatewayClient.test.ts b/ui-tui/src/__tests__/gatewayClient.test.ts index bce7602cc4..f164a561ad 100644 --- a/ui-tui/src/__tests__/gatewayClient.test.ts +++ b/ui-tui/src/__tests__/gatewayClient.test.ts @@ -402,6 +402,34 @@ describe('GatewayClient websocket attach mode', () => { gw.kill() }) + it('surfaces JSON-RPC error code and data to callers (shared error mapping)', async () => { + process.env.HERMES_TUI_GATEWAY_URL = 'ws://gateway.test/api/ws?token=abc' + const gw = new GatewayClient() + + gw.start() + const socket = FakeWebSocket.instances[0]! + + socket.open() + const req = gw.request('projects.create', {}) + await vi.waitFor(() => expect(socket.sent.length).toBeGreaterThan(0)) + + const frame = JSON.parse(socket.sent[0] ?? '{}') as { id: string } + socket.message( + JSON.stringify({ + error: { code: -32601, data: { method: 'projects.create' }, message: 'unknown method: projects.create' }, + id: frame.id, + jsonrpc: '2.0' + }) + ) + + await expect(req).rejects.toMatchObject({ + code: -32601, + data: { method: 'projects.create' }, + message: 'unknown method: projects.create' + }) + gw.kill() + }) + it('uses the undici WebSocket fallback when global WebSocket is unavailable', () => { process.env.HERMES_TUI_GATEWAY_URL = 'ws://gateway.test/api/ws?token=hunter2&channel=secret' delete (globalThis as { WebSocket?: unknown }).WebSocket @@ -517,14 +545,31 @@ describe('GatewayClient websocket attach mode', () => { params: { type: 'gateway.ready', payload: { heartbeat: true } } }) ) + // A live gateway answers every ping (tui_gateway/ws.py replies inline); + // the shared channel counts any inbound frame as liveness, so a socket + // whose pings keep getting acked must never trip the deadline. + const acked: string[] = [] + + const ackPings = () => { + for (const raw of socket.sent) { + const frame = JSON.parse(raw) as { id: string; method: string } + + if (frame.method === 'gateway.ping' && !acked.includes(frame.id)) { + acked.push(frame.id) + socket.message(JSON.stringify({ id: frame.id, jsonrpc: '2.0', result: { ok: true } })) + } + } + } + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_INTERVAL_MS) + expect(JSON.parse(socket.sent.at(-1) ?? '{}')).toMatchObject({ method: 'gateway.ping' }) - const heartbeat = JSON.parse(socket.sent.at(-1) ?? '{}') as { id: string; method: string } + for (let elapsed = 0; elapsed < WS_HEARTBEAT_DEAD_MS * 2; elapsed += WS_HEARTBEAT_INTERVAL_MS) { + ackPings() + await vi.advanceTimersByTimeAsync(WS_HEARTBEAT_INTERVAL_MS) + } - 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(acked.length).toBeGreaterThan(2) expect(socket.readyState).toBe(FakeWebSocket.OPEN) expect(FakeWebSocket.instances).toHaveLength(1) } finally { diff --git a/ui-tui/src/gatewayClient.ts b/ui-tui/src/gatewayClient.ts index fa4e5bb832..14d9350df8 100644 --- a/ui-tui/src/gatewayClient.ts +++ b/ui-tui/src/gatewayClient.ts @@ -4,6 +4,14 @@ import { existsSync } from 'node:fs' import { delimiter, resolve } from 'node:path' import { createInterface } from 'node:readline' +import type { GatewayEvent } from '@hermes/shared/gateway-events' +import { + DEFAULT_HEARTBEAT_DEADLINE_MS, + DEFAULT_HEARTBEAT_INTERVAL_MS, + JsonRpcRequestChannel, + wireFrameText +} from '@hermes/shared/json-rpc-channel' +import { reconnectBackoffDelayMs } from '@hermes/shared/reconnect-backoff' import { WebSocket as UndiciWebSocket } from 'undici' import type { AnyGatewayEvent } from './gatewayTypes.js' @@ -21,15 +29,14 @@ 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. +// Keepalive + dead-connection detection (issue #32997) lives in +// @hermes/shared's JsonRpcRequestChannel; these re-exports keep the TUI's +// timing constants readable at their call sites and in tests. +export const WS_HEARTBEAT_INTERVAL_MS = DEFAULT_HEARTBEAT_INTERVAL_MS +export const WS_HEARTBEAT_DEAD_MS = DEFAULT_HEARTBEAT_DEADLINE_MS +// Exponential backoff for reconnect attempts after a transport drop. No +// jitter: a single TUI process has nobody to desynchronize from, and the +// deterministic ladder is what the activity feed reports. export const RECONNECT_BASE_MS = 1_000 export const RECONNECT_MAX_MS = 30_000 @@ -80,29 +87,6 @@ const resolvePython = (root: string) => { return hit || (process.platform === 'win32' ? 'python' : 'python3') } -const asGatewayEvent = (value: unknown): AnyGatewayEvent | null => - value && typeof value === 'object' && !Array.isArray(value) && typeof (value as { type?: unknown }).type === 'string' - ? (value as AnyGatewayEvent) - : null - -// Hoisted decoder: attach mode can drive high-frequency binary frames -// (tool deltas, reasoning streams) and constructing a fresh TextDecoder -// per message creates avoidable GC pressure. One module-level instance -// is fine because UTF-8 is stateless and we always pass entire frames. -const _wireDecoder = new TextDecoder() - -const asWireText = (raw: unknown): string | null => { - if (typeof raw === 'string') { - return raw - } - - if (raw instanceof ArrayBuffer || ArrayBuffer.isView(raw)) { - return _wireDecoder.decode(raw as any as ArrayBuffer) - } - - return null -} - // Matches `://user:pass@host…` style user-info segments in // otherwise-malformed URLs that the WHATWG `URL` parser can't accept. // Used by the `redactUrl` fallback so embedded credentials are @@ -135,14 +119,6 @@ const redactUrl = (raw: string): string => { } } -interface Pending { - id: string - method: string - reject: (e: Error) => void - resolve: (v: unknown) => void - timeout: ReturnType -} - export class GatewayClient extends EventEmitter { private proc: ChildProcess | null = null private ws: WebSocket | null = null @@ -150,9 +126,17 @@ export class GatewayClient extends EventEmitter { private sidecarWs: WebSocket | null = null private attachUrl: null | string = null private sidecarUrl: null | string = null - private reqId = 0 private logs = new CircularBuffer(MAX_GATEWAY_LOG_LINES) - private pending = new Map() + // Request ids, pending map, timeouts, error mapping and the gateway.ping + // heartbeat are shared with the desktop/web WebSocket client; this class + // only owns the two transports (child stdio, attached socket) and the + // buffered-event replay that Ink's mount order needs. + private readonly channel = new JsonRpcRequestChannel({ + onEvent: ev => this.publish(ev as AnyGatewayEvent), + onHeartbeatFailure: () => this.onHeartbeatFailure(), + requestTimeoutMs: REQUEST_TIMEOUT_MS, + unrefTimers: true + }) private bufferedEvents = new CircularBuffer(MAX_BUFFERED_EVENTS) private pendingExit: number | null | undefined private ready = false @@ -161,13 +145,8 @@ 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 @@ -187,8 +166,8 @@ export class GatewayClient extends EventEmitter { this.readyTimer = null } - if (ev.payload?.heartbeat && this.ws?.readyState === WS_OPEN) { - this.startHeartbeat(this.ws) + if ((ev as GatewayEvent<'gateway.ready'>).payload?.heartbeat === true && this.ws?.readyState === WS_OPEN) { + this.channel.startHeartbeat() } } @@ -235,71 +214,22 @@ 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 - } + // The shared heartbeat found no inbound frame for a full deadline: force the + // socket closed so the ordinary close path reconnects (issue #32997). + private onHeartbeatFailure() { + const ws = this.ws - 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 + if (!ws) { + return } - this.heartbeatPendingId = null - this.heartbeatSentAt = 0 + this.lifecycle('[lifecycle] websocket silent drop detected (heartbeat ack timeout); forcing reconnect') + + try { + ws.close() + } catch { + // ignore + } } private scheduleReconnect() { @@ -307,7 +237,12 @@ export class GatewayClient extends EventEmitter { return } - const delay = Math.min(RECONNECT_BASE_MS * 2 ** this.reconnectAttempts, RECONNECT_MAX_MS) + const delay = reconnectBackoffDelayMs(this.reconnectAttempts, { + baseDelayMs: RECONNECT_BASE_MS, + capMs: RECONNECT_MAX_MS, + jitter: false + }) + 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 } }) @@ -338,7 +273,7 @@ export class GatewayClient extends EventEmitter { // handlers (now identity-gated to ignore unrelated transports) // never fire `rejectPending`, leaving callers hanging on promises // attached to a discarded child / socket. - this.rejectPending(new Error('gateway restarting')) + this.channel.detach(new Error('gateway restarting')) this.ready = false this.subscribed = false // Invalidate any pending deferred drain() flush from a prior transport so @@ -378,7 +313,7 @@ export class GatewayClient extends EventEmitter { this.clearReadyTimer() this.closeSidecarSocket() this.lifecycle(`[lifecycle] transport exit code=${code ?? 'null'} reason=${reason ?? 'none'}`) - this.rejectPending(new Error(reason || `gateway exited${code === null ? '' : ` (${code})`}`)) + this.channel.detach(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 @@ -450,27 +385,30 @@ export class GatewayClient extends EventEmitter { } private handleWebSocketFrame(raw: unknown) { - this.lastActivityAt = Date.now() - const text = asWireText(raw) + const text = wireFrameText(raw) if (!text) { return } - try { - const frame = JSON.parse(text) as Record + const frame = this.channel.handleFrame(text) - if (frame.method === 'event') { - this.mirrorEventToSidecar(text) - } + if (!frame) { + this.protocolError('malformed websocket frame', text, '(empty frame)') - this.dispatch(frame) - } catch { - const preview = text.trim().slice(0, MAX_LOG_PREVIEW) || '(empty frame)' - - this.pushLog(`[protocol] malformed websocket frame: ${preview}`) - this.publish({ type: 'gateway.protocol_error', payload: { preview } }) + return } + + if (frame.method === 'event') { + this.mirrorEventToSidecar(text) + } + } + + private protocolError(what: string, text: string, emptyLabel: string) { + const preview = text.trim().slice(0, MAX_LOG_PREVIEW) || emptyLabel + + this.pushLog(`[protocol] ${what}: ${preview}`) + this.publish({ type: 'gateway.protocol_error', payload: { preview } }) } private startSpawnedGateway(root: string) { @@ -487,15 +425,13 @@ export class GatewayClient extends EventEmitter { this.proc = spawn(python, ['-m', 'tui_gateway.entry'], { cwd, env, stdio: ['pipe', 'pipe', 'pipe'] }) this.lifecycle(`[lifecycle] spawned gateway child ${describeChild(this.proc)} python=${python} cwd=${cwd}`) + const stdin = this.proc.stdin! + this.channel.attach({ send: text => void stdin.write(text + '\n') }) + this.stdoutRl = createInterface({ input: this.proc.stdout! }) this.stdoutRl.on('line', raw => { - try { - this.dispatch(JSON.parse(raw)) - } catch { - const preview = raw.trim().slice(0, MAX_LOG_PREVIEW) || '(empty line)' - - this.pushLog(`[protocol] malformed stdout: ${preview}`) - this.publish({ type: 'gateway.protocol_error', payload: { preview } }) + if (!this.channel.handleFrame(raw)) { + this.protocolError('malformed stdout', raw, '(empty line)') } }) @@ -574,6 +510,11 @@ export class GatewayClient extends EventEmitter { let settled = false this.ws = ws + // Bind the channel to the socket as soon as it exists (not on open): + // RPCs issued while CONNECTING await wsConnectPromise and then must + // reach *this* generation; a stale generation's late frames are + // already filtered by the `this.ws !== ws` guards below. + this.channel.attach({ send: text => ws.send(text) }) const connectPromise = new Promise((resolve, reject) => { ws.addEventListener( @@ -584,7 +525,6 @@ export class GatewayClient extends EventEmitter { resolve() } - this.lastActivityAt = Date.now() this.clearReconnect() this.connectSidecarMirror() }, @@ -637,7 +577,6 @@ 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})` : ''}`) @@ -666,7 +605,6 @@ export class GatewayClient extends EventEmitter { 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)}`) @@ -686,50 +624,6 @@ export class GatewayClient extends EventEmitter { this.startSpawnedGateway(root) } - 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) { - this.settle(p, msg.error ? this.toError(msg.error) : null, msg.result) - - return - } - - if (msg.method === 'event') { - const ev = asGatewayEvent(msg.params) - - if (ev) { - this.publish(ev) - } - } - } - - private toError(raw: unknown): Error { - const err = raw as { message?: unknown } | null | undefined - - return new Error(typeof err?.message === 'string' ? err.message : 'request failed') - } - - private settle(p: Pending, err: Error | null, result: unknown) { - clearTimeout(p.timeout) - this.pending.delete(p.id) - - if (err) { - p.reject(err) - } else { - p.resolve(result) - } - } - private pushLog(line: string) { this.logs.push(truncateLine(line)) } @@ -742,26 +636,6 @@ export class GatewayClient extends EventEmitter { recordParentLifecycle(line) } - private rejectPending(err: Error) { - for (const p of this.pending.values()) { - clearTimeout(p.timeout) - p.reject(err) - } - - this.pending.clear() - } - - // Arrow class-field — stable identity, so `setTimeout(this.onTimeout, …, id)` - // doesn't allocate a bound function per request. - private onTimeout = (id: string) => { - const p = this.pending.get(id) - - if (p) { - this.pending.delete(id) - p.reject(new Error(`timeout: ${p.method}`)) - } - } - drain() { // Defer the buffered-event replay to the next microtask, and DO NOT flip // `subscribed` until that microtask runs. @@ -836,39 +710,9 @@ export class GatewayClient extends EventEmitter { return this.ws } - private requestOverWebSocket(method: string, params: Record = {}): Promise { - return this.ensureAttachedWebSocket(method).then( - ws => - new Promise((resolve, reject) => { - const id = `r${++this.reqId}` - const timeout = setTimeout(this.onTimeout, REQUEST_TIMEOUT_MS, id) + private notConnected = (method: string) => new Error(`gateway not connected: ${method}`) - timeout.unref?.() - this.pending.set(id, { - id, - method, - reject, - resolve: v => resolve(v as T), - timeout - }) - - try { - ws.send(JSON.stringify({ id, jsonrpc: '2.0', method, params })) - } catch (e) { - const pending = this.pending.get(id) - - if (pending) { - clearTimeout(pending.timeout) - this.pending.delete(id) - } - - reject(e instanceof Error ? e : new Error(String(e))) - } - }) - ) - } - - request(method: string, params: Record = {}): Promise { + request(method: string, params: Record = {}, timeoutMs?: number): Promise { const attachUrl = resolveGatewayAttachUrl() if (attachUrl) { @@ -877,11 +721,13 @@ export class GatewayClient extends EventEmitter { // switching from spawned-gateway mode to attach mode also // tears down the old Python child. Merely closing `this.ws` // would leave a previously spawned gateway process alive. - this.rejectPending(new Error('gateway attach url changed')) + this.channel.detach(new Error('gateway attach url changed')) this.start() } - return this.requestOverWebSocket(method, params) + return this.ensureAttachedWebSocket(method).then(() => + this.channel.request(method, params, timeoutMs, undefined, () => this.notConnected(method)) + ) } if (!this.proc?.stdin || this.proc.killed || this.proc.exitCode !== null) { @@ -892,40 +738,12 @@ export class GatewayClient extends EventEmitter { return Promise.reject(new Error('gateway not running')) } - const id = `r${++this.reqId}` - - return new Promise((resolve, reject) => { - const timeout = setTimeout(this.onTimeout, REQUEST_TIMEOUT_MS, id) - - timeout.unref?.() - - this.pending.set(id, { - id, - method, - reject, - resolve: v => resolve(v as T), - timeout - }) - - try { - this.proc!.stdin!.write(JSON.stringify({ id, jsonrpc: '2.0', method, params }) + '\n') - } catch (e) { - const pending = this.pending.get(id) - - if (pending) { - clearTimeout(pending.timeout) - this.pending.delete(id) - } - - reject(e instanceof Error ? e : new Error(String(e))) - } - }) + return this.channel.request(method, params, timeoutMs, undefined, () => this.notConnected(method)) } kill(reason = 'requested') { this.disposed = true this.clearReconnect() - this.stopHeartbeat() const proc = this.proc const killed = proc?.kill() @@ -939,6 +757,6 @@ export class GatewayClient extends EventEmitter { // and we just nulled `this.ws`, so it will short-circuit and // skip handleTransportExit. Reject pending RPCs explicitly so // attach-mode promises do not hang after an intentional kill. - this.rejectPending(new Error('gateway closed')) + this.channel.detach(new Error('gateway closed')) } }