mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(relay): reap a client that has stopped answering instead of holding its leases forever
This commit is contained in:
@@ -1,4 +1,4 @@
|
||||
import { FrameDecoder, KEEPALIVE_SEND_MS, encodeKeepAliveFrame } from './protocol'
|
||||
import { FrameDecoder, KEEPALIVE_SEND_MS, TIMEOUT_MS, encodeKeepAliveFrame } from './protocol'
|
||||
import type { PtyConsumerCloseCause } from '../shared/pty-consumer-session-contract'
|
||||
import type {
|
||||
DispatcherClientWriter,
|
||||
@@ -122,6 +122,8 @@ export abstract class RelayDispatcherClientLifecycle extends RelayDispatcherClie
|
||||
bulkChain: Promise.resolve(),
|
||||
nextOutgoingSeq: 1,
|
||||
highestReceivedSeq: 0,
|
||||
lastReceivedAt: null,
|
||||
keepaliveObserved: false,
|
||||
generation: 0,
|
||||
closed: false,
|
||||
droppedNotificationLog: null,
|
||||
@@ -144,6 +146,8 @@ export abstract class RelayDispatcherClientLifecycle extends RelayDispatcherClie
|
||||
protected resetClient(client: RelayClient): void {
|
||||
client.nextOutgoingSeq = 1
|
||||
client.highestReceivedSeq = 0
|
||||
client.lastReceivedAt = null
|
||||
client.keepaliveObserved = false
|
||||
client.decoder.reset()
|
||||
client.generation++
|
||||
client.closed = false
|
||||
@@ -163,14 +167,27 @@ export abstract class RelayDispatcherClientLifecycle extends RelayDispatcherClie
|
||||
}
|
||||
|
||||
protected startKeepalive(): void {
|
||||
let lastTickAt = Date.now()
|
||||
this.keepaliveTimer = setInterval(() => {
|
||||
if (this.disposed) {
|
||||
return
|
||||
}
|
||||
const now = Date.now()
|
||||
// Why this threshold and not TIMEOUT_MS: a healthy client answers the PREVIOUS tick, so its
|
||||
// lastReceivedAt is already up to KEEPALIVE_SEND_MS + RTT old. A tick gap beyond
|
||||
// TIMEOUT_MS - KEEPALIVE_SEND_MS therefore pushes staleness past the window on its own, and a
|
||||
// rebase armed at TIMEOUT_MS would not have fired -- reaping every client after a host
|
||||
// suspend, a VM migration, or the relay's own event loop stalling. Mirrors the client's
|
||||
// WAKE_GAP_MS guard (ssh-channel-multiplexer.ts).
|
||||
const resumedAfterPause = now - lastTickAt >= TIMEOUT_MS - KEEPALIVE_SEND_MS
|
||||
lastTickAt = now
|
||||
for (const client of this.clients.values()) {
|
||||
if (client.closed) {
|
||||
continue
|
||||
}
|
||||
if (resumedAfterPause && client.lastReceivedAt !== null) {
|
||||
client.lastReceivedAt = now
|
||||
}
|
||||
client.writer.enqueue(
|
||||
'liveness',
|
||||
() => {
|
||||
@@ -180,11 +197,49 @@ export abstract class RelayDispatcherClientLifecycle extends RelayDispatcherClie
|
||||
13
|
||||
)
|
||||
}
|
||||
this.reapSilentClients(now)
|
||||
}, KEEPALIVE_SEND_MS)
|
||||
// Why: unref so the keepalive interval doesn't pin the event loop and block process exit.
|
||||
this.keepaliveTimer.unref()
|
||||
}
|
||||
|
||||
/**
|
||||
* Drop the transport of a client that has gone silent. The relay's writer parks forever on a
|
||||
* half-open link and nothing else ever notices, so an abandoned viewer kept its owner lease and
|
||||
* left every PTY it held paused — the shape behind the "SSH degrades until I cannot connect at
|
||||
* all" reports. Reaping is a statement about the TRANSPORT only: the cause stays the cautious
|
||||
* 'local' default because silence is not evidence the peer died, and the PTYs stay live for the
|
||||
* replacement client to reclaim (docs/reference/ssh-execution-boundary.md).
|
||||
*/
|
||||
private reapSilentClients(now: number): void {
|
||||
for (const client of Array.from(this.clients.values())) {
|
||||
// Why the primary is exempt: closing it tears down the relay's own stdin/stdout, and nothing
|
||||
// in production revives it -- setWrite() has no non-test caller. The leak this exists for is
|
||||
// a socket client holding an owner lease, and the launch channel's own liveness is already
|
||||
// owned by the client-side dead-link check.
|
||||
if (client === this.primaryClient) {
|
||||
continue
|
||||
}
|
||||
// Why the null check is not just defensive: a relay is launched before its client finishes
|
||||
// handshaking, and on a slow link that can exceed the window. Reaping a client that has never
|
||||
// spoken would break the connect it is still completing, so silence only counts against a
|
||||
// client that has already proven it can talk.
|
||||
// Why keepaliveObserved gates this: not every client speaks the keepalive protocol. The
|
||||
// remote `orca` CLI sends one `orca.cli` request and waits for a result budgeted in minutes
|
||||
// (src/relay/remote-cli-timeout.ts), so judging it on inbound silence would kill
|
||||
// `terminal wait`, `--wait` and `orchestration ask` after 20s.
|
||||
if (
|
||||
client.closed ||
|
||||
!client.keepaliveObserved ||
|
||||
client.lastReceivedAt === null ||
|
||||
now - client.lastReceivedAt <= TIMEOUT_MS
|
||||
) {
|
||||
continue
|
||||
}
|
||||
this.closeClient(client, new Error('Relay client stopped answering'), true)
|
||||
}
|
||||
}
|
||||
|
||||
protected closeClient(
|
||||
client: RelayClient,
|
||||
error: Error,
|
||||
|
||||
@@ -51,6 +51,15 @@ export type RelayClient = {
|
||||
bulkChain: Promise<void>
|
||||
nextOutgoingSeq: number
|
||||
highestReceivedSeq: number
|
||||
// Why: the relay had no inbound-liveness signal at all, so a half-open client was never reaped
|
||||
// and kept its owner lease and paused PTYs indefinitely.
|
||||
lastReceivedAt: number | null
|
||||
// Why silence is only held against a client that sends keepalives: not every client speaks that
|
||||
// protocol. The remote `orca` CLI opens the socket, sends one `orca.cli` request and then waits
|
||||
// for a result that is deliberately budgeted in minutes (src/relay/remote-cli-timeout.ts), so
|
||||
// judging it on inbound silence would kill `terminal wait`, `--wait` and `orchestration ask`
|
||||
// after 20s. Only a client that has proven it participates is eligible.
|
||||
keepaliveObserved: boolean
|
||||
generation: number
|
||||
closed: boolean
|
||||
droppedNotificationLog: DroppedProducerNotificationLog | null
|
||||
|
||||
@@ -15,11 +15,14 @@ import { RelayDispatcherCapacitySignals } from './dispatcher-capacity-signals'
|
||||
|
||||
export abstract class RelayDispatcherFrameCodec extends RelayDispatcherCapacitySignals {
|
||||
protected handleFrame(client: RelayClient, frame: DecodedFrame): void {
|
||||
// Before the KeepAlive early return: a keepalive is the only proof a quiet client is still there.
|
||||
client.lastReceivedAt = Date.now()
|
||||
if (frame.id > client.highestReceivedSeq) {
|
||||
client.highestReceivedSeq = frame.id
|
||||
}
|
||||
|
||||
if (frame.type === MessageType.KeepAlive) {
|
||||
client.keepaliveObserved = true
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,131 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { RelayDispatcher } from './dispatcher'
|
||||
import { encodeJsonRpcFrame, encodeKeepAliveFrame, KEEPALIVE_SEND_MS, TIMEOUT_MS } from './protocol'
|
||||
|
||||
// The relay had no inbound-liveness signal at all: its writer parks forever on a half-open link, so
|
||||
// an abandoned viewer kept its owner lease and left the PTYs it held paused until the process died.
|
||||
describe('RelayDispatcher silent-client reaper', () => {
|
||||
let dispatcher: RelayDispatcher
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers()
|
||||
vi.setSystemTime(0)
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
dispatcher.dispose()
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('detaches a client that spoke once and then stopped answering', () => {
|
||||
const detachListener = vi.fn()
|
||||
dispatcher = new RelayDispatcher(() => true)
|
||||
dispatcher.onClientDetached(detachListener)
|
||||
const clientId = dispatcher.attachClient(() => true)
|
||||
dispatcher.feedClient(clientId, encodeKeepAliveFrame(1, 0))
|
||||
|
||||
vi.advanceTimersByTime(TIMEOUT_MS + KEEPALIVE_SEND_MS * 2)
|
||||
|
||||
// 'local', not a peer close: silence is not evidence the peer died, and a consumer that read it
|
||||
// as one would shorten the owner grace on a session that is still there.
|
||||
expect(detachListener).toHaveBeenCalledWith(clientId, 'local')
|
||||
})
|
||||
|
||||
it('keeps a quiet but answering client attached', () => {
|
||||
const detachListener = vi.fn()
|
||||
dispatcher = new RelayDispatcher(() => true)
|
||||
dispatcher.onClientDetached(detachListener)
|
||||
const clientId = dispatcher.attachClient(() => true)
|
||||
|
||||
// A client with nothing to say still answers the keepalive; that is the only proof required.
|
||||
for (let tick = 0; tick < 10; tick += 1) {
|
||||
vi.advanceTimersByTime(KEEPALIVE_SEND_MS)
|
||||
dispatcher.feedClient(clientId, encodeKeepAliveFrame(tick + 1, 0))
|
||||
}
|
||||
|
||||
// Asserted against this client specifically: the unattached primary sink has no peer answering
|
||||
// it in this harness, so it is expected to be reaped and says nothing about the case under test.
|
||||
expect(detachListener).not.toHaveBeenCalledWith(clientId, expect.anything())
|
||||
})
|
||||
|
||||
it('does not reap every client on the first tick after the host slept', () => {
|
||||
const detachListener = vi.fn()
|
||||
dispatcher = new RelayDispatcher(() => true)
|
||||
dispatcher.onClientDetached(detachListener)
|
||||
dispatcher.attachClient(() => true)
|
||||
|
||||
// One tick fires far late because the process was paused, not because the peers went away.
|
||||
vi.setSystemTime(10 * 60_000)
|
||||
vi.advanceTimersByTime(KEEPALIVE_SEND_MS)
|
||||
|
||||
expect(detachListener).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('never reaps a client that has not spoken yet', () => {
|
||||
// A relay is launched before its client finishes handshaking, and on a slow link that can
|
||||
// outlast the window. Reaping there would break the connect the client is still completing.
|
||||
const detachListener = vi.fn()
|
||||
dispatcher = new RelayDispatcher(() => true)
|
||||
dispatcher.onClientDetached(detachListener)
|
||||
const clientId = dispatcher.attachClient(() => true)
|
||||
|
||||
vi.advanceTimersByTime(TIMEOUT_MS * 5)
|
||||
|
||||
expect(detachListener).not.toHaveBeenCalledWith(clientId, expect.anything())
|
||||
})
|
||||
|
||||
it('never reaps a client that does not send keepalives at all', () => {
|
||||
// The remote `orca` CLI opens the socket, sends one `orca.cli` request and then waits for a
|
||||
// result budgeted in minutes (remote-cli-timeout.ts: 5min default, 10min for wait, 11min for
|
||||
// orchestration ask). It has no keepalive timer, so judging it on inbound silence would abort
|
||||
// `terminal wait`, `--wait` and `orchestration ask` after 20s.
|
||||
const detachListener = vi.fn()
|
||||
dispatcher = new RelayDispatcher(() => true)
|
||||
dispatcher.onClientDetached(detachListener)
|
||||
const clientId = dispatcher.attachClient(() => true)
|
||||
dispatcher.feedClient(
|
||||
clientId,
|
||||
encodeJsonRpcFrame({ jsonrpc: '2.0', id: 1, method: 'orca.cli', params: {} }, 1, 0)
|
||||
)
|
||||
|
||||
vi.advanceTimersByTime(TIMEOUT_MS * 20)
|
||||
|
||||
expect(detachListener).not.toHaveBeenCalledWith(clientId, expect.anything())
|
||||
})
|
||||
|
||||
it('does not reap a healthy client when the relay itself stalls for most of the window', () => {
|
||||
// The dead band this exists for: a healthy client answers the PREVIOUS tick, so its
|
||||
// lastReceivedAt is already ~KEEPALIVE_SEND_MS old. A tick gap short of TIMEOUT_MS still pushes
|
||||
// staleness past the window, so a rebase armed at TIMEOUT_MS would never fire and every client
|
||||
// would be reaped after a host suspend, VM migration, or an event-loop stall.
|
||||
const detachListener = vi.fn()
|
||||
dispatcher = new RelayDispatcher(() => true)
|
||||
dispatcher.onClientDetached(detachListener)
|
||||
const clientId = dispatcher.attachClient(() => true)
|
||||
|
||||
// The client answers at t=5s, then the next tick at t=10s finds it already ~5s stale — normal.
|
||||
vi.advanceTimersByTime(KEEPALIVE_SEND_MS)
|
||||
dispatcher.feedClient(clientId, encodeKeepAliveFrame(1, 0))
|
||||
vi.advanceTimersByTime(KEEPALIVE_SEND_MS)
|
||||
|
||||
// Now the relay stalls: the clock jumps but no tick runs, so the following tick lands 17s after
|
||||
// the last one. Staleness is 22s (past the window) while the tick gap is under TIMEOUT_MS, so a
|
||||
// rebase armed at TIMEOUT_MS would not fire and this healthy client would be reaped.
|
||||
vi.setSystemTime(Date.now() + 12_000)
|
||||
vi.advanceTimersByTime(KEEPALIVE_SEND_MS)
|
||||
|
||||
expect(detachListener).not.toHaveBeenCalledWith(clientId, expect.anything())
|
||||
})
|
||||
|
||||
it('never reaps the primary client, whose sink cannot be revived', () => {
|
||||
const detachListener = vi.fn()
|
||||
dispatcher = new RelayDispatcher(() => true)
|
||||
dispatcher.onClientDetached(detachListener)
|
||||
|
||||
vi.advanceTimersByTime(TIMEOUT_MS * 20)
|
||||
|
||||
// Client id 1 is the primary sink; closing it would tear down the relay's own stdin/stdout and
|
||||
// nothing in production calls setWrite() to bring it back.
|
||||
expect(detachListener).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user