mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 08:03:12 +00:00
fix(remote-terminal): keep the stream stall deadline armed on unacknowledged credit
A paired-runtime terminal could stall silently with a live socket, a live PTY and no transport error (#11265). Two compounding defects on the read side: - The stream watchdog re-armed its 30s delivery deadline from zero on every settled delivery, so sibling traffic postponed the verdict indefinitely, and it cleared the timer entirely once renderer parse credit hit zero. Re-arming required inbound output -- the exact thing an exhausted host ACK window stops -- so once an ACK went missing nothing could ever detect the stall. The deadline is now anchored to the oldest unsettled delivery and stays armed while delivered bytes remain unacknowledged to the host. - flushOutputAcknowledgement zeroed pendingAckBytes before knowing the ACK frame was accepted, permanently shrinking the host's send window. Unsent bytes are re-charged and the flush timer re-armed. Recovery still reports onTransportClose({recoverable:true}); no path claims the PTY exited.
This commit is contained in:
@@ -84,20 +84,36 @@ export abstract class RemoteRuntimeTerminalFlowController extends RemoteRuntimeT
|
||||
if (stream.pendingAckBytes >= TERMINAL_MULTIPLEX_ACK_BATCH_BYTES) {
|
||||
return this.flushOutputAcknowledgement(stream)
|
||||
}
|
||||
if (stream.ackFlushTimer === null) {
|
||||
stream.ackFlushTimer = setTimeout(() => {
|
||||
stream.ackFlushTimer = null
|
||||
this.flushOutputAcknowledgement(stream)
|
||||
}, TERMINAL_MULTIPLEX_ACK_FLUSH_MS)
|
||||
}
|
||||
this.scheduleOutputAcknowledgementFlush(stream)
|
||||
return true
|
||||
}
|
||||
|
||||
private scheduleOutputAcknowledgementFlush(stream: RemoteRuntimeMultiplexedTerminalState): void {
|
||||
if (stream.ackFlushTimer !== null) {
|
||||
return
|
||||
}
|
||||
stream.ackFlushTimer = setTimeout(() => {
|
||||
stream.ackFlushTimer = null
|
||||
this.flushOutputAcknowledgement(stream)
|
||||
}, TERMINAL_MULTIPLEX_ACK_FLUSH_MS)
|
||||
}
|
||||
|
||||
private flushOutputAcknowledgement(stream: RemoteRuntimeMultiplexedTerminalState): boolean {
|
||||
clearAckFlushTimer(stream)
|
||||
const bytes = stream.pendingAckBytes
|
||||
if (bytes <= 0) {
|
||||
return true
|
||||
}
|
||||
stream.pendingAckBytes = 0
|
||||
return bytes <= 0 || this.acknowledgeOutput(stream, bytes)
|
||||
if (this.acknowledgeOutput(stream, bytes)) {
|
||||
// Why: only a frame the transport took reopens the host's send window.
|
||||
stream.watchdog.recordOutputAcknowledged(bytes)
|
||||
return true
|
||||
}
|
||||
// Why re-charged: dropping an unsent ack shrinks the host window for the stream's life, and the only retry trigger is the output that shrunken window blocks.
|
||||
stream.pendingAckBytes += bytes
|
||||
this.scheduleOutputAcknowledgementFlush(stream)
|
||||
return false
|
||||
}
|
||||
|
||||
getStreamsForE2e(): Iterable<RemoteRuntimeMultiplexedTerminalState> {
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
REMOTE_TERMINAL_DELIVERY_STALL_TIMEOUT_MS,
|
||||
createRemoteTerminalStreamWatchdog
|
||||
} from './remote-terminal-stream-watchdog'
|
||||
|
||||
describe('remote terminal stream watchdog delivery deadline', () => {
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers()
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('anchors the deadline to the oldest unsettled delivery instead of the last settled sibling', () => {
|
||||
const onStall = vi.fn()
|
||||
const watchdog = createRemoteTerminalStreamWatchdog(onStall)
|
||||
|
||||
watchdog.beginOutputDelivery(100)
|
||||
for (let tick = 0; tick < 3; tick += 1) {
|
||||
vi.advanceTimersByTime(9_000)
|
||||
const settle = watchdog.beginOutputDelivery(10)
|
||||
settle()
|
||||
}
|
||||
expect(onStall).not.toHaveBeenCalled()
|
||||
|
||||
vi.advanceTimersByTime(REMOTE_TERMINAL_DELIVERY_STALL_TIMEOUT_MS - 27_000)
|
||||
|
||||
expect(onStall).toHaveBeenCalledTimes(1)
|
||||
expect(onStall.mock.calls[0]?.[0]).toMatchObject({ reason: 'delivery-credit-timeout' })
|
||||
})
|
||||
|
||||
it('stays armed while parsed bytes remain unacknowledged to the host', () => {
|
||||
const onStall = vi.fn()
|
||||
const watchdog = createRemoteTerminalStreamWatchdog(onStall)
|
||||
|
||||
watchdog.beginOutputDelivery(100)()
|
||||
vi.advanceTimersByTime(REMOTE_TERMINAL_DELIVERY_STALL_TIMEOUT_MS)
|
||||
|
||||
expect(onStall).toHaveBeenCalledTimes(1)
|
||||
expect(onStall.mock.calls[0]?.[0]).toMatchObject({
|
||||
outstandingDeliveryBytes: 0,
|
||||
reason: 'delivery-credit-timeout'
|
||||
})
|
||||
})
|
||||
|
||||
it('disarms once the acknowledgement reaches the transport', () => {
|
||||
const onStall = vi.fn()
|
||||
const watchdog = createRemoteTerminalStreamWatchdog(onStall)
|
||||
|
||||
watchdog.beginOutputDelivery(100)()
|
||||
watchdog.recordOutputAcknowledged(100)
|
||||
vi.advanceTimersByTime(REMOTE_TERMINAL_DELIVERY_STALL_TIMEOUT_MS * 2)
|
||||
|
||||
expect(onStall).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
@@ -9,6 +9,8 @@ export type RemoteTerminalStreamStall = {
|
||||
|
||||
export type RemoteTerminalStreamWatchdog = {
|
||||
beginOutputDelivery: (bytes: number) => () => void
|
||||
/** Bytes whose ACK frame reached the transport, releasing the host's window. */
|
||||
recordOutputAcknowledged: (bytes: number) => void
|
||||
completeCommandResponseProbe: () => void
|
||||
recordCommandInput: (text: string) => void
|
||||
recordInbound: () => void
|
||||
@@ -21,6 +23,9 @@ export function createRemoteTerminalStreamWatchdog(
|
||||
let responseTimer: ReturnType<typeof setTimeout> | null = null
|
||||
let deliveryTimer: ReturnType<typeof setTimeout> | null = null
|
||||
let outstandingDeliveryBytes = 0
|
||||
// Why separate from parse credit: the host reopens its window on ACK frames, so bytes parsed but not yet ACKed are still the credit whose loss stops output.
|
||||
let unacknowledgedBytes = 0
|
||||
let deliveryPendingSinceMs: number | null = null
|
||||
let lastInboundAtMs = Date.now()
|
||||
let commandResponseProbePending = false
|
||||
let disposed = false
|
||||
@@ -54,23 +59,29 @@ export function createRemoteTerminalStreamWatchdog(
|
||||
reason
|
||||
})
|
||||
}
|
||||
const armDeliveryTimer = (): void => {
|
||||
clearDeliveryTimer()
|
||||
if (outstandingDeliveryBytes <= 0 || disposed) {
|
||||
// Why anchored, never restarted: a deadline re-armed by sibling settles is postponed forever, and one cleared at zero parse credit can only re-arm from inbound output — which is what the stall stops.
|
||||
const syncDeliveryTimer = (): void => {
|
||||
if (disposed || outstandingDeliveryBytes + unacknowledgedBytes <= 0) {
|
||||
clearDeliveryTimer()
|
||||
deliveryPendingSinceMs = null
|
||||
return
|
||||
}
|
||||
deliveryTimer = setTimeout(
|
||||
() => trip('delivery-credit-timeout'),
|
||||
REMOTE_TERMINAL_DELIVERY_STALL_TIMEOUT_MS
|
||||
deliveryPendingSinceMs ??= Date.now()
|
||||
if (deliveryTimer) {
|
||||
return
|
||||
}
|
||||
const remainingMs = Math.max(
|
||||
0,
|
||||
REMOTE_TERMINAL_DELIVERY_STALL_TIMEOUT_MS - (Date.now() - deliveryPendingSinceMs)
|
||||
)
|
||||
deliveryTimer = setTimeout(() => trip('delivery-credit-timeout'), remainingMs)
|
||||
}
|
||||
|
||||
return {
|
||||
beginOutputDelivery(bytes) {
|
||||
outstandingDeliveryBytes += bytes
|
||||
if (!deliveryTimer) {
|
||||
armDeliveryTimer()
|
||||
}
|
||||
unacknowledgedBytes += bytes
|
||||
syncDeliveryTimer()
|
||||
let settled = false
|
||||
return () => {
|
||||
if (settled || disposed) {
|
||||
@@ -78,9 +89,16 @@ export function createRemoteTerminalStreamWatchdog(
|
||||
}
|
||||
settled = true
|
||||
outstandingDeliveryBytes = Math.max(0, outstandingDeliveryBytes - bytes)
|
||||
armDeliveryTimer()
|
||||
syncDeliveryTimer()
|
||||
}
|
||||
},
|
||||
recordOutputAcknowledged(bytes) {
|
||||
if (disposed) {
|
||||
return
|
||||
}
|
||||
unacknowledgedBytes = Math.max(0, unacknowledgedBytes - bytes)
|
||||
syncDeliveryTimer()
|
||||
},
|
||||
completeCommandResponseProbe() {
|
||||
commandResponseProbePending = false
|
||||
},
|
||||
@@ -103,6 +121,8 @@ export function createRemoteTerminalStreamWatchdog(
|
||||
clearResponseTimer()
|
||||
clearDeliveryTimer()
|
||||
outstandingDeliveryBytes = 0
|
||||
unacknowledgedBytes = 0
|
||||
deliveryPendingSinceMs = null
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user