export type GatewayEventName = | 'gateway.ready' | 'session.info' | 'session.usage' | 'message.start' | 'message.delta' | 'message.interim' | 'message.complete' | 'thinking.delta' | 'reasoning.delta' | 'reasoning.available' | 'status.update' | 'tool.start' | 'tool.progress' | 'tool.complete' | 'tool.generating' | 'todo.updated' | 'clarify.request' | 'approval.request' | 'sudo.request' | 'secret.request' | 'background.complete' | 'error' | 'skin.changed' | (string & {}) export interface GatewayEvent

{ payload?: P /** Renderer-side source tag added by the Desktop gateway registry. */ profile?: string /** Registry connection whose socket delivered the event (renderer-side tag; * absent for the local/legacy primary path). */ connectionId?: string session_id?: string type: GatewayEventName } 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 connectTimeoutMs?: number createRequestId?: (nextId: number) => GatewayRequestId heartbeatDeadlineMs?: number heartbeatIntervalMs?: number /** Return true to intercept the default closed-state transition. */ onSocketClose?: (event: CloseEvent) => boolean | void requestIdPrefix?: string requestTimeoutMs?: number socketFactory?: (url: string) => WebSocketLike notConnectedErrorMessage?: string } const ANY = '*' 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. const DEFAULT_CONNECT_TIMEOUT_MS = 15_000 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 /** 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). */ private replayInFlight = false /** * While a replay fetch is in flight, live seq'd frames for the sessions * being replayed are parked here instead of dispatching immediately. * Without this hold, a live frame racing the replay response is dispatched * twice (once live, once when the replay returns the same seq) or, worse, * advances the watermark so the gap events the replay carries get skipped. */ private replayHold: Map | null = null /** * Server process identity for the replay contract (from gateway.ready / * session.events.since). Seq counters are in-process on the backend, so a * restart resets them while we still hold high watermarks — without this * check events_since(sid, 97) returns [] + truncated=false forever and we * 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 constructor(options: GatewayClientOptions = {}) { this.options = { closedErrorMessage: options.closedErrorMessage ?? 'WebSocket closed', connectErrorMessage: options.connectErrorMessage ?? 'WebSocket connection failed', connectTimeoutMs: options.connectTimeoutMs ?? DEFAULT_CONNECT_TIMEOUT_MS, createRequestId: options.createRequestId ?? ((nextId: number) => `${options.requestIdPrefix ?? 'r'}${nextId}`), heartbeatDeadlineMs: options.heartbeatDeadlineMs ?? DEFAULT_HEARTBEAT_DEADLINE_MS, heartbeatIntervalMs: options.heartbeatIntervalMs ?? DEFAULT_HEARTBEAT_INTERVAL_MS, notConnectedErrorMessage: options.notConnectedErrorMessage ?? 'gateway not connected', onSocketClose: options.onSocketClose ?? (() => false), requestIdPrefix: options.requestIdPrefix ?? 'r', requestTimeoutMs: options.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS, socketFactory: options.socketFactory } } get connectionState(): ConnectionState { return this.state } async connect(wsUrl: string): Promise { // Refuse garbage; WebSocket coerces non-strings into // `ws:///[object%20Object]` (#68250 stale-emit boot loop). const invalidUrl = () => { const got = typeof wsUrl === 'string' ? JSON.stringify(wsUrl) : `type "${typeof wsUrl}"` return new Error(`gateway connect() requires a ws:// or wss:// URL string, got ${got}`) } if (typeof wsUrl !== 'string') { throw invalidUrl() } let url: URL try { url = new URL(wsUrl) } catch { throw invalidUrl() } if (url.protocol !== 'ws:' && url.protocol !== 'wss:') { throw invalidUrl() } if (this.socket?.readyState === WebSocket.OPEN || this.state === 'connecting') { return } this.setState('connecting') const socket = this.options.socketFactory?.(wsUrl) ?? new WebSocket(wsUrl) this.socket = socket this.stopHeartbeat() socket.addEventListener('message', message => { if (this.socket !== socket) { return } this.lastInboundAt = Date.now() this.handleMessage(message.data) }) socket.addEventListener('close', event => { if (this.socket !== socket) { return } if (this.options.onSocketClose(event)) { return } this.socket = null this.stopHeartbeat() this.setState('closed') this.rejectAllPending(new Error(this.options.closedErrorMessage)) }) await new Promise((resolve, reject) => { let settled = false let timer: ReturnType | undefined const cleanup = () => { if (timer !== undefined) { clearTimeout(timer) } socket.removeEventListener('open', onOpen) socket.removeEventListener('error', onError) } const onOpen = () => { if (settled || this.socket !== socket) { return } settled = true cleanup() this.setState('open') resolve() // Lossless resume: drain events emitted while we were disconnected. // Fire-and-forget so connect() latency is unaffected; only runs when // we actually observed seq'd events before the drop. void this.fetchReplay() } const onError = () => { if (settled || this.socket !== socket) { return } settled = true cleanup() this.setState('error') reject(new Error(this.options.connectErrorMessage)) } socket.addEventListener('open', onOpen, { once: true }) socket.addEventListener('error', onError, { once: true }) if (this.options.connectTimeoutMs > 0) { timer = setTimeout(() => { if (settled) { return } settled = true cleanup() // Drop the half-open socket so the next connect() starts clean // instead of short-circuiting on a zombie 'connecting' state. if (this.socket === socket) { try { socket.close() } catch { // ignore } this.socket = null this.setState('error') } reject(new Error(this.options.connectErrorMessage)) }, this.options.connectTimeoutMs) } }) } 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)) } } /** * Invalidate the current socket generation after an ambiguous transport * outcome. The outer connection owner decides whether/when to reconnect. */ invalidate(message = this.options.closedErrorMessage): void { const socket = this.socket if (!socket) { return } this.invalidateSocket(socket, new Error(message)) } on

(type: GatewayEventName, 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) } onAny(handler: (event: GatewayEvent) => void): () => void { return this.on(ANY as GatewayEventName, handler) } onEvent(handler: (event: GatewayEvent) => void): () => void { return this.onAny(handler) } onState(handler: (state: ConnectionState) => void): () => void { this.stateHandlers.add(handler) handler(this.state) return () => this.stateHandlers.delete(handler) } request( method: string, params: Record = {}, timeoutMs = this.options.requestTimeoutMs, signal?: AbortSignal ): Promise { const socket = this.socket if (!socket || 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))) } }) } private handleMessage(raw: unknown): void { const text = typeof raw === 'string' ? raw : String(raw) let frame: JsonRpcFrame try { frame = JSON.parse(text) as JsonRpcFrame } catch { return } if (frame.id !== undefined && frame.id !== null) { const call = this.pending.get(frame.id) 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) } return } if (frame.method === 'event' && frame.params?.type) { if (frame.params.type === 'gateway.ready') { if (this.gatewayReadyAdvertisesHeartbeat(frame.params.payload)) { const socket = this.socket if (socket) { this.startHeartbeat(socket) } } const epoch = (frame.params.payload as { replay_epoch?: unknown } | undefined)?.replay_epoch if (typeof epoch === 'string' && epoch) { this.adoptReplayEpoch(epoch) } } const sid = frame.params.session_id const seqValue = (frame.params as { seq?: unknown }).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) } } /** * Track each session's last observed event seq. Events without a seq * (legacy backend, session-less globals) leave the map untouched. */ private recordSeq(event: GatewayEvent): void { const sid = event.session_id const seq = (event as { seq?: unknown }).seq if (!sid || typeof seq !== 'number' || !Number.isFinite(seq)) { return } const prev = this.lastSeenSeq.get(sid) ?? 0 if (seq > prev) { this.lastSeenSeq.set(sid, seq) } } /** Test/telemetry hook: current last-seen seq map snapshot. */ getSeqWatermarks(): Record { return Object.fromEntries(this.lastSeenSeq) } /** * After a reconnect, ask the gateway to replay every event newer than our * per-session watermarks. Replayed frames go through the SAME dispatchEvent * path as live frames — dedupe happens naturally because recordSeq ignores * non-increasing seqs and downstream stores key on event identity. * Best-effort: failures are swallowed (the next reconnect retries). */ private async fetchReplay(): Promise { if (this.replayInFlight || this.lastSeenSeq.size === 0) { return } this.replayInFlight = true // Park live frames for the sessions we're about to replay so a frame // racing the replay response can't dispatch ahead of (or duplicate) the // gap events. Sessions without watermarks are unaffected. const hold = new Map() for (const sid of this.lastSeenSeq.keys()) { hold.set(sid, []) } this.replayHold = hold try { const entries = Object.entries(this.getSeqWatermarks()) // One RPC per known session keeps params flat; sessions are few (<20). const results = await Promise.allSettled( entries.map(([sid, lastSeen]) => this.request<{ events?: Array<{ type: string; session_id?: string; seq?: number; payload?: unknown }> }>( 'session.events.since', { session_id: sid, last_seen: lastSeen }, REPLAY_REQUEST_TIMEOUT_MS ) ) ) for (const result of results) { if (result.status !== 'fulfilled' || !Array.isArray(result.value?.events)) { continue } const epoch = (result.value as { epoch?: unknown }).epoch if (typeof epoch === 'string' && epoch && this.replayEpoch && epoch !== this.replayEpoch) { // Backend restarted: its seq numbering reset, so our watermarks — // and this replay window — are meaningless. Drop them and start // fresh under the new epoch. this.adoptReplayEpoch(epoch) continue } if (typeof epoch === 'string' && epoch && !this.replayEpoch) { this.replayEpoch = epoch } for (const event of result.value.events) { if (!event?.type) { continue } this.dispatchIfNewer(event as GatewayEvent) } } } catch { // Replay is an optimization over lossy-reconnect; never surface errors. } finally { this.flushReplayHold() this.replayInFlight = false } } /** * Dispatch an event only when its seq advances the session watermark. * Seq-less events always dispatch (no ordering contract to violate). */ private dispatchIfNewer(event: GatewayEvent): void { const sid = event.session_id const seq = (event as { seq?: unknown }).seq if (sid && typeof seq === 'number' && Number.isFinite(seq)) { const prev = this.lastSeenSeq.get(sid) ?? 0 if (seq <= prev) { return } this.lastSeenSeq.set(sid, seq) } this.dispatchEvent(event) } /** * Record the server's replay epoch; on change (backend restart) the old * seq watermarks describe a numbering that no longer exists — clear them * so the next reconnect doesn't silently believe it missed nothing. */ private adoptReplayEpoch(epoch: string): void { if (this.replayEpoch === epoch) { return } if (this.replayEpoch !== null) { this.lastSeenSeq.clear() } this.replayEpoch = epoch } /** Release frames parked during a replay fetch, seq-gated against dupes. */ private flushReplayHold(): void { const hold = this.replayHold this.replayHold = null if (!hold) { return } for (const parked of hold.values()) { for (const event of parked) { this.dispatchIfNewer(event) } } } private gatewayReadyAdvertisesHeartbeat(payload: unknown): boolean { return Boolean(payload && typeof payload === 'object' && (payload as { heartbeat?: unknown }).heartbeat === true) } 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 } this.socket = null this.stopHeartbeat() try { socket.close() } catch { // The generation was already invalidated; the reconnect owner can redial. } 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) } } private setState(state: ConnectionState): void { if (this.state === state) { return } this.state = state for (const handler of this.stateHandlers) { handler(state) } } }