Files
web-terminal/relay-run/src/wiring/socket-pipe.ts
Yaojia Wang c1c837c54f fix(relay-run): buffer early ws frames in wsToWebSocketLike
The message listener was attached only when the mux consumer called onMessage — which
happens inside attach(), AFTER the async mTLS-registry lookup. On loopback the agent's FIRST
heartbeat ping arrived in that gap and ws dropped it (no listener), so the agent's heartbeat
died on the first miss at 15s and the tunnel flapped forever. Attach the listener at
construction and buffer early frames (and an early close) until the consumer wires up.
2026-07-07 04:59:53 +02:00

129 lines
4.2 KiB
TypeScript

/**
* `WebSocketLike` adapters. Two of them:
* - `makeSocketPair()` — an in-process duplex pair used to drive the REAL relay-node/mux in tests
* (no sockets, deterministic). Delivery is synchronous but non-reentrant (a per-endpoint
* trampoline) so a send emitted while draining the peer's inbox can't recurse unboundedly.
* - `wsToWebSocketLike()` — wraps a real `ws` WebSocket to the same interface for `main.ts`.
* Bytes are COPIED on send (the mux slices/retains buffers) so no caller can alias another's memory.
*/
import type { WebSocket as WsWebSocket } from 'ws'
import type { WebSocketLike } from 'term-relay/data-plane/ws-like.js'
interface Endpoint extends WebSocketLike {
/** Internal: deliver bytes that the peer sent to THIS endpoint. */
_deliver(data: Uint8Array): void
_closed(): void
}
function makeEndpoint(): Endpoint {
let onMessage: ((data: Uint8Array) => void) | null = null
let onClose: (() => void) | null = null
let peer: Endpoint | null = null
const inbox: Uint8Array[] = []
let draining = false
let closed = false
const ep: Endpoint & { _link(p: Endpoint): void } = {
_link(p: Endpoint) {
peer = p
},
send(data: Uint8Array) {
if (closed || peer === null) return
peer._deliver(new Uint8Array(data)) // copy: never alias the sender's buffer
},
close(_code?: number) {
if (closed) return
closed = true
onClose?.()
peer?._closed()
},
onMessage(handler) {
onMessage = handler
},
onClose(handler) {
onClose = handler
},
_deliver(data: Uint8Array) {
inbox.push(data)
if (draining) return
draining = true
try {
while (inbox.length > 0) {
const next = inbox.shift()!
onMessage?.(next)
}
} finally {
draining = false
}
},
_closed() {
if (closed) return
closed = true
onClose?.()
},
}
return ep as Endpoint & { _link(p: Endpoint): void }
}
/** An in-process duplex `WebSocketLike` pair: `a.send` → `b`'s message handler and vice versa. */
export function makeSocketPair(): { a: WebSocketLike; b: WebSocketLike } {
const a = makeEndpoint() as Endpoint & { _link(p: Endpoint): void }
const b = makeEndpoint() as Endpoint & { _link(p: Endpoint): void }
a._link(b)
b._link(a)
return { a, b }
}
function toUint8(data: unknown): Uint8Array | null {
if (data instanceof Uint8Array) return data
if (data instanceof ArrayBuffer) return new Uint8Array(data)
if (Array.isArray(data)) return Buffer.concat(data.map((d) => (d instanceof Uint8Array ? d : Buffer.from(d))))
if (Buffer.isBuffer(data)) return new Uint8Array(data)
return null
}
/**
* Adapt a live `ws` WebSocket to `WebSocketLike` (binary frames only).
*
* The `ws.on('message'|'close')` listeners are attached NOW (at construction), not when the consumer
* later calls `onMessage`/`onClose`. The composition root builds this synchronously on the socket's
* 'connection' event, but the mux consumer only wires its handler AFTER an async mTLS-registry lookup
* (the bridge). Any frame that arrives in that gap — notably the agent's FIRST heartbeat ping, sent
* immediately on open — would otherwise be dropped by `ws` (no listener yet), starving the heartbeat
* and flapping the tunnel every 15 s. We buffer early frames (and an early close) and flush on wire-up.
*/
export function wsToWebSocketLike(ws: WsWebSocket): WebSocketLike {
let onMessage: ((data: Uint8Array) => void) | null = null
let onClose: (() => void) | null = null
const backlog: Uint8Array[] = []
let closedEarly = false
ws.on('message', (data: unknown) => {
const bytes = toUint8(data)
if (bytes === null) return
if (onMessage === null) backlog.push(bytes)
else onMessage(bytes)
})
ws.on('close', () => {
if (onClose === null) closedEarly = true
else onClose()
})
return {
send(data: Uint8Array) {
ws.send(data, { binary: true })
},
close(code?: number) {
ws.close(code)
},
onMessage(handler) {
onMessage = handler
while (backlog.length > 0) handler(backlog.shift()!)
},
onClose(handler) {
onClose = handler
if (closedEarly) handler()
},
}
}