From 005d0c5a48525d7b8a530d36f89e59da044145b6 Mon Sep 17 00:00:00 2001 From: Neil <4138956+nwparker@users.noreply.github.com> Date: Tue, 1 Sep 2026 00:12:27 -0700 Subject: [PATCH] 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. --- ...remote-runtime-terminal-flow-controller.ts | 30 +++++++--- .../remote-terminal-stream-watchdog.test.ts | 58 +++++++++++++++++++ .../remote-terminal-stream-watchdog.ts | 40 +++++++++---- 3 files changed, 111 insertions(+), 17 deletions(-) create mode 100644 src/renderer/src/runtime/remote-terminal-stream-watchdog.test.ts diff --git a/src/renderer/src/runtime/remote-runtime-terminal-flow-controller.ts b/src/renderer/src/runtime/remote-runtime-terminal-flow-controller.ts index 640ec25933e..f4fc1a5c0cd 100644 --- a/src/renderer/src/runtime/remote-runtime-terminal-flow-controller.ts +++ b/src/renderer/src/runtime/remote-runtime-terminal-flow-controller.ts @@ -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 { diff --git a/src/renderer/src/runtime/remote-terminal-stream-watchdog.test.ts b/src/renderer/src/runtime/remote-terminal-stream-watchdog.test.ts new file mode 100644 index 00000000000..66d65f9bcbc --- /dev/null +++ b/src/renderer/src/runtime/remote-terminal-stream-watchdog.test.ts @@ -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() + }) +}) diff --git a/src/renderer/src/runtime/remote-terminal-stream-watchdog.ts b/src/renderer/src/runtime/remote-terminal-stream-watchdog.ts index 4cb6b081c5f..be9e1b9d404 100644 --- a/src/renderer/src/runtime/remote-terminal-stream-watchdog.ts +++ b/src/renderer/src/runtime/remote-terminal-stream-watchdog.ts @@ -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 | null = null let deliveryTimer: ReturnType | 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 } } }