refactor(ts): ui-tui rides apps/shared's JSON-RPC request channel; one pending map, one heartbeat, typed RPC errors
Two independent JSON-RPC client cores existed for one backend: apps/shared's
JsonRpcGatewayClient (desktop, web) and ui-tui/src/gatewayClient.ts, which
re-implemented request ids, the pending map with timeouts, response->error
mapping, event decoding and the gateway.ping heartbeat (~200 LOC, drifted).
Split the transport-agnostic half out of the shared client into
JsonRpcRequestChannel (apps/shared/src/json-rpc-channel.ts): the owner binds a
JsonRpcTransport { send(text) } per connection generation and feeds inbound
text through handleFrame(). JsonRpcGatewayClient keeps only the WebSocket
lifecycle, seq replay and the typed event hub on top of it; the Ink TUI keeps
only its two transports (spawned child stdio, attached socket) and its
mount-order event buffering, and delegates everything else.
Behavior change:
- TUI RPC errors now carry the JSON-RPC `code` / `data` (JsonRpcGatewayError)
instead of a bare Error(message); the TUI's timeout text is now the shared
"request timed out after Ns: <method>" (was "timeout: <method>", matched by
no caller) and callers may pass a per-call timeout.
- TUI heartbeat liveness counts any inbound frame (shared semantics) rather
than tracking one in-flight ping id; the interval/deadline are unchanged
and pings no longer carry the unread `last_activity_ms` param.
- Desktop isMissingRpcMethod reads the -32601 code first and only regexes the
message for code-less (IPC-flattened) errors, so a tool result that merely
mentions "unknown method" no longer reads as a capability verdict.
- Shared connect() now settles on a `close` during the handshake (auth-gate
4401/4403) instead of waiting out the 15s connect timeout, and
invalidate()/close() drop the socket generation before calling close() so a
synchronous close event cannot run the closed-path twice.
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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'
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -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<typeof setTimeout>
|
||||
}
|
||||
|
||||
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<GatewayRequestId, PendingCall>()
|
||||
private transport: JsonRpcTransport | null = null
|
||||
private heartbeatTimer: ReturnType<typeof setInterval> | null = null
|
||||
private heartbeatSequence = 0
|
||||
private lastInboundAt = 0
|
||||
private readonly options: Required<Omit<JsonRpcRequestChannelOptions, 'onEvent' | 'onHeartbeatFailure'>> &
|
||||
Pick<JsonRpcRequestChannelOptions, 'onEvent' | 'onHeartbeatFailure'>
|
||||
|
||||
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<T>(
|
||||
method: string,
|
||||
params: Record<string, unknown> = {},
|
||||
timeoutMs = this.options.requestTimeoutMs,
|
||||
signal?: AbortSignal,
|
||||
notConnectedError: () => Error = () => new Error('gateway not connected')
|
||||
): Promise<T> {
|
||||
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<T>((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)
|
||||
}
|
||||
}
|
||||
}
|
||||
+138
-297
@@ -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<typeof setTimeout>
|
||||
}
|
||||
|
||||
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<string, Set<(event: GatewayEvent) => void>>()
|
||||
|
||||
on<K extends GatewayEventName>(type: K, handler: (event: GatewayEvent<K>) => 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<GatewayEventName>) => 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<GatewayRequestId, PendingCall>()
|
||||
private socket: WebSocketLike | null = null
|
||||
private state: ConnectionState = 'idle'
|
||||
private heartbeatTimer: ReturnType<typeof setInterval> | 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<string, number>()
|
||||
/** 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<string, Set<(event: GatewayEvent) => void>>()
|
||||
private readonly stateHandlers = new Set<(state: ConnectionState) => void>()
|
||||
private readonly options: Required<Omit<GatewayClientOptions, 'socketFactory'>> &
|
||||
Pick<GatewayClientOptions, 'socketFactory'>
|
||||
@@ -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<void>((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<K extends GatewayEventName>(type: K, handler: (event: GatewayEvent<K>) => 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<GatewayEventName>) => 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<T> {
|
||||
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<T>((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<T>(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<void> {
|
||||
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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
+83
-265
@@ -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 `<scheme>://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<typeof setTimeout>
|
||||
}
|
||||
|
||||
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<string>(MAX_GATEWAY_LOG_LINES)
|
||||
private pending = new Map<string, Pending>()
|
||||
// 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<AnyGatewayEvent>(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<typeof createInterface> | null = null
|
||||
private stderrRl: ReturnType<typeof createInterface> | null = null
|
||||
private heartbeatTimer: ReturnType<typeof setInterval> | null = null
|
||||
private reconnectTimer: ReturnType<typeof setTimeout> | 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<string, unknown>
|
||||
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<void>((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<string, unknown>) {
|
||||
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<T = unknown>(method: string, params: Record<string, unknown> = {}): Promise<T> {
|
||||
return this.ensureAttachedWebSocket(method).then(
|
||||
ws =>
|
||||
new Promise<T>((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<T = unknown>(method: string, params: Record<string, unknown> = {}): Promise<T> {
|
||||
request<T = unknown>(method: string, params: Record<string, unknown> = {}, timeoutMs?: number): Promise<T> {
|
||||
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<T>(method, params)
|
||||
return this.ensureAttachedWebSocket(method).then(() =>
|
||||
this.channel.request<T>(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<T>((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<T>(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'))
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user