fix(runtime): defer heartbeat probes until websocket auth

This commit is contained in:
Jinwoo-H
2026-09-01 13:23:11 -04:00
parent 057d40aeb5
commit 577ce9ea94
4 changed files with 45 additions and 22 deletions
+18 -11
View File
@@ -2493,30 +2493,31 @@
"https://github.com/stablyai/orca/issues/11298",
"https://github.com/stablyai/orca/pull/11300"
],
"invariant": "Every accepted runtime socket installs message, pong, close, and error ownership before any heartbeat probe. With uninterrupted timer delivery, the first socket that arms an idle heartbeat is probed immediately and an unresponsive socket is reaped within one interval. Later sockets join the existing shared cadence without another timer or immediate sweep and are reaped within two intervals. Responsive sockets survive, pause recovery grants a fresh probe, and close or error-to-close releases connection listeners and timers.",
"oracle": "With one fake clock and exact socket identities, accept the first socket at 0 ms and require an immediate owned probe plus reaping at 100 ms when unresponsive. Keep a responsive first socket, accept an unresponsive later socket at 50 ms, require the same shared timer, its first probe at 100 ms, no early reap, and termination at 200 ms. Inject synchronous message, pong, close, and error events, then require exact heartbeat membership and zero retained timers/listeners after final close. Production transport tests independently cover real socket round trips, pre-auth and capacity bounds, revocation, shutdown, and half-open cleanup.",
"invariant": "Every accepted runtime socket installs message, pong, close, and error ownership before any heartbeat probe. Unauthenticated sockets receive no heartbeat control frames during E2EE and are bounded by the pre-auth timeout; authenticated sockets share one periodic cadence, tolerate missed probes and event-loop pause, and release listeners and timers on close or error.",
"oracle": "With one fake clock and exact socket identities, accept an authenticated first socket at 0 ms and require its first probe on the 100 ms tick. Accept an unauthenticated later socket at 50 ms and require no probe at the 100 ms tick; authenticate it, then require its first probe on the next shared tick and termination only after the configured consecutive-miss budget. Inject synchronous message, pong, close, and error events, then require exact heartbeat membership and zero retained timers/listeners after final close. Production transport tests independently cover real socket round trips, pre-auth and capacity bounds, revocation, shutdown, and half-open cleanup.",
"commands": [
"pnpm exec vitest run --config config/vitest.config.ts src/main/runtime/rpc/ws-transport-accept-order.test.ts src/main/runtime/rpc/remote-runtime-server-heartbeat.test.ts src/main/runtime/rpc/ws-transport.test.ts"
"pnpm exec vitest run --config config/vitest.config.ts src/main/runtime/rpc/ws-transport-accept-order.test.ts src/main/runtime/rpc/remote-runtime-server-heartbeat.test.ts src/main/runtime/rpc/ws-transport.test.ts src/main/runtime/rpc/ws-transport-transient-packet-loss.test.ts"
],
"testFiles": [
"src/main/runtime/rpc/ws-transport-accept-order.test.ts",
"src/main/runtime/rpc/remote-runtime-server-heartbeat.test.ts",
"src/main/runtime/rpc/ws-transport.test.ts"
"src/main/runtime/rpc/ws-transport.test.ts",
"src/main/runtime/rpc/ws-transport-transient-packet-loss.test.ts"
],
"assertionRefs": [
{
"file": "src/main/runtime/rpc/ws-transport-accept-order.test.ts",
"assertions": [
"the first synchronous probe observes message, pong, close, and error ownership",
"the first unresponsive socket is reaped at one interval",
"a later socket keeps the original shared timer, is first probed on the shared tick, and is reaped within two intervals",
"the first periodic probe observes message, pong, close, and error ownership",
"an authenticated unresponsive socket is reaped on the configured consecutive-miss budget",
"an unauthenticated later socket is not probed before authentication, then joins the original shared timer",
"final close releases heartbeat membership, listeners, and timers"
]
},
{
"file": "src/main/runtime/rpc/remote-runtime-server-heartbeat.test.ts",
"assertions": [
"one missed probe reaps only the unresponsive client",
"a single missed probe does not reap the unresponsive client",
"event-loop resume grants clients a fresh probe"
]
},
@@ -2527,6 +2528,12 @@
"pre-auth, raw TCP, and accepted WebSocket resource bounds remain enforced",
"error and close races finalize membership once"
]
},
{
"file": "src/main/runtime/rpc/ws-transport-transient-packet-loss.test.ts",
"assertions": [
"an authenticated real WebSocket survives one swallowed pong and responds to later probes"
]
}
],
"evidenceRuns": [
@@ -2534,10 +2541,10 @@
"date": "2026-07-29",
"runner": "local",
"platform": "macos",
"command": "pnpm exec vitest run --config config/vitest.config.ts src/main/runtime/rpc/ws-transport-accept-order.test.ts src/main/runtime/rpc/remote-runtime-server-heartbeat.test.ts src/main/runtime/rpc/ws-transport.test.ts",
"command": "pnpm exec vitest run --config config/vitest.config.ts src/main/runtime/rpc/ws-transport-accept-order.test.ts src/main/runtime/rpc/remote-runtime-server-heartbeat.test.ts src/main/runtime/rpc/ws-transport.test.ts src/main/runtime/rpc/ws-transport-transient-packet-loss.test.ts",
"result": "passed",
"durationSeconds": 0.91,
"summary": "Three runtime transport files and 35 tests passed with deterministic first- and later-socket cadence, exact listener/timer ownership, pause recovery, security bounds, and cleanup."
"summary": "Four runtime transport files and 42 tests passed with deterministic authenticated heartbeat cadence, handshake grace for later sockets, exact listener/timer ownership, pause recovery, transient packet-loss tolerance, security bounds, and cleanup."
}
],
"runtimeBudget": {
@@ -2550,7 +2557,7 @@
},
"redGreenEvidence": {
"status": "complete",
"evidence": "On latest main and with the listener-order fix disabled, the first synchronous probe observed no message, pong, close, or error owner. The structural listener-order fix passed those assertions. The published delayed-first-sweep alternative missed the first-socket one-interval cleanup bound. A later-socket oracle now separately pins the intended shared cadence at a 100 ms first probe and 200 ms reap after acceptance at 50 ms."
"evidence": "On the candidate before this review fix, an authenticated first socket kept the shared timer alive while an unauthenticated socket accepted at 50 ms was pinged at the 100 ms tick; the new oracle failed. After filtering heartbeat clients to the authenticated transport map, the same byte-identical oracle passes: no pre-auth ping, first post-auth probe on the next shared tick, and bounded cleanup."
},
"performanceBudget": {
"required": true,
@@ -63,13 +63,14 @@ describe('WebSocketTransport accepted socket ordering', () => {
socket.once('open', () => events.push('open'))
socket.emit('open')
lifecycle.handleConnection(socket)
transport.setClientId(socket, 'client')
await vi.advanceTimersByTimeAsync(100)
expect(firstProbeListeners).toEqual({ pong: 1, message: 1, close: 1, error: 1 })
expect(events.slice(0, 4)).toEqual(['open', 'ping', 'ready', 'pong'])
expect(socket.terminate).not.toHaveBeenCalled()
expect(lifecycle.heartbeatConnections.size).toBe(1)
expect(vi.getTimerCount()).toBe(2)
expect(vi.getTimerCount()).toBe(1)
events.push('close')
socket.emit('close')
@@ -94,11 +95,12 @@ describe('WebSocketTransport accepted socket ordering', () => {
socket.emit('close')
}
})
const { lifecycle } = createHarness(socket)
const { lifecycle, transport } = createHarness(socket)
socket.once('open', () => events.push('open'))
socket.emit('open')
lifecycle.handleConnection(socket)
transport.setClientId(socket, 'client')
expect(socket.terminate).not.toHaveBeenCalled()
// A single silent interval is not evidence: the socket is re-probed, not reaped.
await vi.advanceTimersByTimeAsync(100)
@@ -126,9 +128,10 @@ describe('WebSocketTransport accepted socket ordering', () => {
firstSocket.emit('pong')
}
})
const { lifecycle } = createHarness(firstSocket)
const { lifecycle, transport } = createHarness(firstSocket)
lifecycle.handleConnection(firstSocket)
transport.setClientId(firstSocket, 'first-client')
const sharedTimer = lifecycle.heartbeat.timer
expect(firstPingTimes).toEqual([])
@@ -148,21 +151,25 @@ describe('WebSocketTransport accepted socket ordering', () => {
expect(lifecycle.heartbeat.timer).toBe(sharedTimer)
expect(laterPingTimes).toEqual([])
expect(vi.getTimerCount()).toBe(3)
expect(vi.getTimerCount()).toBe(2)
await vi.advanceTimersByTimeAsync(50)
expect(laterPingTimes).toEqual([100])
expect(laterPingTimes).toEqual([])
expect(laterSocket.terminate).not.toHaveBeenCalled()
transport.setClientId(laterSocket, 'later-client')
await vi.advanceTimersByTimeAsync(100)
expect(laterPingTimes).toEqual([200])
// Re-probed on every sweep while its misses bank, so a recovered path can answer immediately.
await vi.advanceTimersByTimeAsync(299)
expect(laterPingTimes).toEqual([100, 200, 300])
expect(laterPingTimes).toEqual([200, 300, 400])
expect(laterSocket.terminate).not.toHaveBeenCalled()
await vi.advanceTimersByTimeAsync(1)
expect(laterSocket.terminate).toHaveBeenCalledTimes(1)
expect(laterReapTimes).toEqual([400])
expect(firstPingTimes).toEqual([100, 200, 300, 400])
expect(laterReapTimes).toEqual([500])
expect(firstPingTimes).toEqual([100, 200, 300, 400, 500])
expect(
['pong', 'message', 'close', 'error'].map((event) => laterSocket.listenerCount(event))
).toEqual([0, 0, 0, 0])
@@ -183,11 +190,12 @@ describe('WebSocketTransport accepted socket ordering', () => {
socket.emit('error', new Error('probe failed'))
}
})
const { lifecycle } = createHarness(socket)
const { lifecycle, transport } = createHarness(socket)
socket.once('open', () => events.push('open'))
socket.emit('open')
lifecycle.handleConnection(socket)
transport.setClientId(socket, 'client')
await vi.advanceTimersByTimeAsync(100)
events.push('close')
socket.emit('close')
@@ -68,6 +68,12 @@ describe('WebSocketTransport under transient packet loss', () => {
client.once('open', resolve)
client.once('error', reject)
})
// Heartbeats intentionally begin after authentication; stamp this synthetic peer as authenticated
// so the liveness oracle exercises the production heartbeat path rather than pre-auth expiry.
const wss = (transport as unknown as { wss: { clients: Set<WebSocket> } }).wss
for (const serverSocket of wss.clients) {
transport.setClientId(serverSocket, 'test-client')
}
// Bounded by counted probe events, never by elapsed time.
await vi.waitFor(() => expect(closed || probesReceived >= swallowedProbe + 2).toBe(true), {
+4 -2
View File
@@ -319,11 +319,13 @@ export class WebSocketTransport implements RpcTransport {
ws.on('close', finalizeConnection)
ws.on('error', onError)
// Why: every lifecycle event must have an owner before the first synchronous probe.
// Why: install lifecycle ownership before periodic heartbeat ticks can observe this socket.
this.heartbeatConnections.add(ws)
this.heartbeat.noteAlive(ws)
if (this.heartbeatConnections.size === 1) {
this.heartbeat.start(() => this.wss?.clients ?? [])
// Unauthenticated sockets are protected by the pre-auth timeout; heartbeat probes begin only
// after E2EE binds a client id, avoiding control frames during the handshake.
this.heartbeat.start(() => this.wsClientIds.keys())
}
}