mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(ssh): declare a wedged relay link lost instead of suppressing the dead-link check
This commit is contained in:
@@ -0,0 +1,82 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { encodeKeepAliveFrame, KEEPALIVE_SEND_MS, TIMEOUT_MS } from './relay-protocol'
|
||||
import { SshChannelMultiplexer, type MultiplexerTransport } from './ssh-channel-multiplexer'
|
||||
|
||||
type WedgedTransport = MultiplexerTransport & { writes: Buffer[]; feed: (chunk: Buffer) => void }
|
||||
|
||||
/**
|
||||
* A transport that accepts the first write, then reports backpressure forever: no drain, and no
|
||||
* write settlement. This is a half-open TCP link — the socket buffer filled and the peer's FIN
|
||||
* never arrived — which is what sleep/resume and a dropped NAT mapping produce in the field.
|
||||
*/
|
||||
function createWedgedTransport(): WedgedTransport {
|
||||
const writes: Buffer[] = []
|
||||
let onData: (chunk: Buffer) => void = () => {}
|
||||
return {
|
||||
write: (data) => {
|
||||
writes.push(data)
|
||||
return false
|
||||
},
|
||||
onData: (callback) => {
|
||||
onData = callback
|
||||
},
|
||||
onClose: () => {},
|
||||
onDrain: () => () => {},
|
||||
supportsWriteSettlement: true,
|
||||
writes,
|
||||
feed: (chunk) => onData(chunk)
|
||||
}
|
||||
}
|
||||
|
||||
describe('SshChannelMultiplexer on a transport that saturates and never drains', () => {
|
||||
let transport: WedgedTransport
|
||||
let mux: SshChannelMultiplexer
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers()
|
||||
vi.setSystemTime(0)
|
||||
transport = createWedgedTransport()
|
||||
mux = new SshChannelMultiplexer(transport)
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
mux.dispose()
|
||||
vi.restoreAllMocks()
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('declares the link lost instead of suppressing the dead-link check forever', async () => {
|
||||
// Drive well past every health window: keepalive interval, dead-link timeout, and the
|
||||
// wake-gap grace that resets staleness after a suspend.
|
||||
await vi.advanceTimersByTimeAsync(TIMEOUT_MS * 10 + KEEPALIVE_SEND_MS)
|
||||
|
||||
expect(mux.isDisposed()).toBe(true)
|
||||
})
|
||||
|
||||
it('fails a request parked behind saturation rather than leaving it pending forever', async () => {
|
||||
const settled = vi.fn()
|
||||
mux.request('pty.spawn', {}).then(
|
||||
() => settled('resolved'),
|
||||
() => settled('rejected')
|
||||
)
|
||||
|
||||
await vi.advanceTimersByTimeAsync(TIMEOUT_MS * 10 + KEEPALIVE_SEND_MS)
|
||||
|
||||
expect(settled).toHaveBeenCalledWith('rejected')
|
||||
})
|
||||
|
||||
it('keeps a slow-but-alive peer connected while its own keepalives arrive', async () => {
|
||||
// The regression guard for the fix above: backpressure on our uplink is not evidence of
|
||||
// death, and the relay's own keepalive is what proves it.
|
||||
let seq = 1
|
||||
const inbound = setInterval(() => {
|
||||
transport.feed(encodeKeepAliveFrame(seq++, 0))
|
||||
}, KEEPALIVE_SEND_MS)
|
||||
try {
|
||||
await vi.advanceTimersByTimeAsync(TIMEOUT_MS * 10 + KEEPALIVE_SEND_MS)
|
||||
expect(mux.isDisposed()).toBe(false)
|
||||
} finally {
|
||||
clearInterval(inbound)
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -71,7 +71,6 @@ type MuxInternals = {
|
||||
disposeHandlers: unknown[]
|
||||
lastReceivedAt: number
|
||||
unackedTimestamps: Map<number, number>
|
||||
writerSaturated: boolean
|
||||
}
|
||||
|
||||
function getMuxInternals(instance: SshChannelMultiplexer): MuxInternals {
|
||||
@@ -354,9 +353,10 @@ describe('SshChannelMultiplexer', () => {
|
||||
expect(mux.isDisposed()).toBe(true)
|
||||
})
|
||||
|
||||
it('suppresses false death while locally saturated and rebases both clocks on drain', () => {
|
||||
it('survives local saturation while the peer keeps talking, and rebases both clocks on drain', () => {
|
||||
mux.dispose()
|
||||
let drain = (): void => {}
|
||||
let feed: (chunk: Buffer) => void = () => {}
|
||||
const written: Buffer[] = []
|
||||
const saturatedTransport: MultiplexerTransport = {
|
||||
write: (data) => {
|
||||
@@ -367,21 +367,29 @@ describe('SshChannelMultiplexer', () => {
|
||||
onDrain: (callback) => {
|
||||
drain = callback
|
||||
},
|
||||
onData: vi.fn(),
|
||||
onData: (callback) => {
|
||||
feed = callback
|
||||
},
|
||||
onClose: vi.fn()
|
||||
}
|
||||
mux = new SshChannelMultiplexer(saturatedTransport)
|
||||
|
||||
vi.advanceTimersByTime(5_000)
|
||||
expect(getMuxInternals(mux).writerSaturated).toBe(true)
|
||||
vi.advanceTimersByTime(25_000)
|
||||
// The writer parked after its first frame: that is the saturation this test is about.
|
||||
expect(written).toHaveLength(1)
|
||||
// Why: backpressure on our uplink is not evidence of death. The relay's own keepalive is,
|
||||
// and only that inbound traffic may keep the link alive — suppressing the check on
|
||||
// saturation alone wedged a half-open link forever (see the saturation-wedge suite).
|
||||
for (let tick = 0; tick < 5; tick++) {
|
||||
feed(encodeKeepAliveFrame(0, 0))
|
||||
vi.advanceTimersByTime(5_000)
|
||||
}
|
||||
expect(mux.isDisposed()).toBe(false)
|
||||
expect(written).toHaveLength(1)
|
||||
|
||||
drain()
|
||||
const resumedAt = Date.now()
|
||||
const internals = getMuxInternals(mux)
|
||||
expect(internals.writerSaturated).toBe(false)
|
||||
expect(internals.lastReceivedAt).toBe(resumedAt)
|
||||
expect(new Set(internals.unackedTimestamps.values())).toEqual(new Set([resumedAt]))
|
||||
|
||||
|
||||
@@ -102,7 +102,6 @@ export class SshChannelMultiplexer {
|
||||
private disposed = false
|
||||
private disposeReason: 'shutdown' | 'connection_lost' | null = null
|
||||
private decoderReadPaused = false
|
||||
private writerSaturated = false
|
||||
|
||||
// Track the oldest unacked outgoing message timestamp
|
||||
private unackedTimestamps = new Map<number, number>()
|
||||
@@ -587,7 +586,12 @@ export class SshChannelMultiplexer {
|
||||
|
||||
this.sendKeepAlive()
|
||||
|
||||
if (this.disposed || resumedAfterWake || this.decoderReadPaused || this.writerSaturated) {
|
||||
// Why: a saturated writer used to suppress this check outright, which wedged a half-open
|
||||
// link forever — no drain, so no frame ever left, and the writer's single-outstanding
|
||||
// liveness guard silenced the one probe that could have noticed. The relay sends its own
|
||||
// keepalive every KEEPALIVE_SEND_MS, so a slow-but-alive peer still refreshes
|
||||
// lastReceivedAt; only a link that delivers nothing inbound is declared lost.
|
||||
if (this.disposed || resumedAfterWake || this.decoderReadPaused) {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -650,7 +654,6 @@ export class SshChannelMultiplexer {
|
||||
}
|
||||
|
||||
private handleWriterSaturationChange(saturated: boolean): void {
|
||||
this.writerSaturated = saturated
|
||||
if (!saturated && !this.disposed) {
|
||||
this.rebaseHealthClocks(Date.now())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user