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 } } }