From 404cee070e0cd3a12aca69a0decb349dd6b2e06b Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Mon, 7 Sep 2026 18:37:34 -0700 Subject: [PATCH] refactor(terminal): separate parked pane health from producer credit --- src/main/ipc/pty-delivery-health-heal.test.ts | 60 +------ .../ipc/pty-ipc-producer-flow-control.test.ts | 8 +- src/main/ipc/pty/delivery/accounting.ts | 6 +- ...pty-dispatcher-parked-delivery-ack.test.ts | 69 +++---- .../terminal-pane/pty-dispatcher.ts | 23 +-- .../terminal-pane/pty-parked-delivery-debt.ts | 105 ----------- .../terminal-pane/pty-pre-handler-buffer.ts | 108 +++++------ ...-pre-handler-parked-debt-invariant.test.ts | 170 ------------------ ...nal-delivery-watchdog-parked-stall.test.ts | 85 +++++---- .../terminal-delivery-watchdog.test.ts | 16 +- .../terminal-delivery-watchdog.ts | 82 +++------ .../terminal-pane/terminal-pane-recovery.ts | 2 +- .../terminal-parked-delivery-stall.ts | 14 +- .../terminal-parked-pane-ownership.test.ts | 93 ---------- .../terminal-parked-pane-recovery.ts | 30 +--- src/shared/pty-renderer-delivery-health.ts | 9 - 16 files changed, 179 insertions(+), 701 deletions(-) delete mode 100644 src/renderer/src/components/terminal-pane/pty-parked-delivery-debt.ts delete mode 100644 src/renderer/src/components/terminal-pane/pty-pre-handler-parked-debt-invariant.test.ts delete mode 100644 src/renderer/src/components/terminal-pane/terminal-parked-pane-ownership.test.ts diff --git a/src/main/ipc/pty-delivery-health-heal.test.ts b/src/main/ipc/pty-delivery-health-heal.test.ts index 0c013ecd259..9dcfda35d8d 100644 --- a/src/main/ipc/pty-delivery-health-heal.test.ts +++ b/src/main/ipc/pty-delivery-health-heal.test.ts @@ -294,32 +294,6 @@ describe('registerPtyHandlers', () => { vi.useRealTimers() } }) - it('writes off parked bytes the renderer reported as having no consumer', async () => { - vi.useFakeTimers() - const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - const mockProc = createMockProc() - spawnMock.mockReturnValue(mockProc.proc) - - try { - const spawnResult = await spawnAndSaturateRendererDeliveryGate(mockProc) - - // The renderer received every byte but has no handler for them, so they sit in the - // pre-handler buffer un-ACKed. Received-minus-parked is what can still repay itself. - const healed = reportRendererDeliveryState({ - receivedCharsByPty: { [spawnResult.id]: 512 * 1024 }, - processedCharsByPty: {}, - parkedCharsByPty: { [spawnResult.id]: 512 * 1024 }, - heal: true, - rendererPtyDataListenerCount: 1 - }) - - expect(healed.writtenOff).toEqual([{ id: spawnResult.id, writtenOffChars: 512 * 1024 }]) - expect(getPtyRendererDeliveryDebugSnapshot()).toMatchObject({ rendererInFlightChars: 0 }) - } finally { - warnSpy.mockRestore() - vi.useRealTimers() - } - }) it('still heals when the same report repairs a lost ACK for the wedged PTY', async () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) @@ -334,9 +308,8 @@ describe('registerPtyHandlers', () => { // lastAckAtMs, and reading that stamp back would veto the heal the report asked for — // a recovered ACK is evidence of a LOST one, never of a live consumer. const healed = reportRendererDeliveryState({ - receivedCharsByPty: { [spawnResult.id]: 512 * 1024 }, + receivedCharsByPty: { [spawnResult.id]: 1 }, processedCharsByPty: { [spawnResult.id]: 1 }, - parkedCharsByPty: { [spawnResult.id]: 512 * 1024 }, heal: true, rendererPtyDataListenerCount: 1 }) @@ -348,31 +321,6 @@ describe('registerPtyHandlers', () => { vi.useRealTimers() } }) - it('still skips a PTY whose received bytes are only partly parked', async () => { - vi.useFakeTimers() - const mockProc = createMockProc() - spawnMock.mockReturnValue(mockProc.proc) - - try { - const spawnResult = await spawnAndSaturateRendererDeliveryGate(mockProc) - - // Half the bytes are still in the live parse path; their deferred ACK repays that half, - // so nothing here is provably lost and the write-off must stay out of it. - const health = reportRendererDeliveryState({ - receivedCharsByPty: { [spawnResult.id]: 512 * 1024 }, - processedCharsByPty: {}, - parkedCharsByPty: { [spawnResult.id]: 256 * 1024 }, - heal: true - }) - - expect(health.writtenOff).toBeUndefined() - expect(getPtyRendererDeliveryDebugSnapshot()).toMatchObject({ - rendererInFlightChars: 512 * 1024 - }) - } finally { - vi.useRealTimers() - } - }) it('heals an ACK-silent PTY while a sibling PTY keeps ACKing', () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) @@ -391,9 +339,8 @@ describe('registerPtyHandlers', () => { getPtyAckDataListener()(null, { id: 'pty-live', processedChars: 512 * 1024 }) const healed = reportRendererDeliveryState({ - receivedCharsByPty: { 'pty-wedged': 512 * 1024 }, + receivedCharsByPty: {}, processedCharsByPty: {}, - parkedCharsByPty: { 'pty-wedged': 512 * 1024 }, heal: true, rendererPtyDataListenerCount: 1 }) @@ -426,9 +373,8 @@ describe('registerPtyHandlers', () => { getPtyAckDataListener()(null, { id: 'pty-live', processedChars: 256 * 1024 }) const healed = reportRendererDeliveryState({ - receivedCharsByPty: { 'pty-wedged': 512 * 1024 }, + receivedCharsByPty: {}, processedCharsByPty: {}, - parkedCharsByPty: { 'pty-wedged': 512 * 1024 }, heal: true, rendererPtyDataListenerCount: 1 }) diff --git a/src/main/ipc/pty-ipc-producer-flow-control.test.ts b/src/main/ipc/pty-ipc-producer-flow-control.test.ts index ca607848ba5..f71a8991383 100644 --- a/src/main/ipc/pty-ipc-producer-flow-control.test.ts +++ b/src/main/ipc/pty-ipc-producer-flow-control.test.ts @@ -188,7 +188,7 @@ describe('registerPtyHandlers', () => { } ]) }) - it('pauses the shell for a pane whose bytes are parked, and releases it on the write-off', () => { + it('pauses the shell for lost push delivery and releases it on write-off', () => { vi.useFakeTimers() const warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) try { @@ -196,8 +196,7 @@ describe('registerPtyHandlers', () => { registerPtyHandlers(mainWindow as never) mainWindow.webContents.send.mockClear() - // Backpressure for free: the renderer withholds the ACK for bytes nothing consumes, and - // main's existing in-flight window turns that into a paused producer. No new mechanism. + // No push bytes arrive, so main retains credit until the invoke heal recovers delivery. provider.emitData('parked-pty', 'x'.repeat(900 * 1024)) vi.advanceTimersByTime(2) for (let index = 0; index < 400; index++) { @@ -211,9 +210,8 @@ describe('registerPtyHandlers', () => { expect(getPtyDataSendCalls()).toHaveLength(sendsWhileParked) const healed = reportRendererDeliveryState({ - receivedCharsByPty: { 'parked-pty': 512 * 1024 }, + receivedCharsByPty: {}, processedCharsByPty: {}, - parkedCharsByPty: { 'parked-pty': 512 * 1024 }, heal: true }) diff --git a/src/main/ipc/pty/delivery/accounting.ts b/src/main/ipc/pty/delivery/accounting.ts index 009b93ebcf5..4075be3501e 100644 --- a/src/main/ipc/pty/delivery/accounting.ts +++ b/src/main/ipc/pty/delivery/accounting.ts @@ -225,12 +225,8 @@ export function writeOffLostRendererDelivery( continue } const receivedChars = sanitizeReportedChars(report.receivedCharsByPty?.[id]) - // Parked bytes sit in the renderer's pre-handler buffer with no handler to consume them - // and a withheld ACK, so they can never repay themselves; only what is left in the live - // parse path does. Absent field ⇒ 0 parked ⇒ the predicate this replaced, byte for byte. - const parkedChars = sanitizeReportedChars(report.parkedCharsByPty?.[id]) // Why skip: received-but-unparsed bytes are alive in the renderer write queue; their deferred ACK still repays this debt. - if (receivedChars - parkedChars > accounting.ackedChars) { + if (receivedChars > accounting.ackedChars) { continue } const acknowledged = applyCumulativeAck(session, id, accounting.sentChars) diff --git a/src/renderer/src/components/terminal-pane/pty-dispatcher-parked-delivery-ack.test.ts b/src/renderer/src/components/terminal-pane/pty-dispatcher-parked-delivery-ack.test.ts index 81d1c866162..b4d163acac7 100644 --- a/src/renderer/src/components/terminal-pane/pty-dispatcher-parked-delivery-ack.test.ts +++ b/src/renderer/src/components/terminal-pane/pty-dispatcher-parked-delivery-ack.test.ts @@ -1,11 +1,3 @@ -/** - * The held-ACK contract, driven through the real dispatcher rather than the buffer helper. - * - * Why it matters: `pty:data` for a PTY with no registered handler used to be parked AND - * ACKed. Debt for a dead pane was therefore zero, and the delivery watchdog's stall - * predicate needs `inFlightTotalChars > 0` — so it was blind by construction, main's flow - * control read healthy, and the shell kept flooding a pane that rendered nothing. - */ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' vi.mock('@/lib/e2e-config', () => ({ e2eConfig: { exposeStore: false } })) @@ -37,10 +29,6 @@ describe('parked pty:data delivery credit', () => { return () => {} }, ackData, - // Partial test APIs omit it, and the watchdog refuses to start without it, so the - // dispatcher keeps today's ACK-at-return. Note the web remote client is NOT this - // case — it stubs a real one and does arm the watchdog. Nothing parks there because - // its `onData` never emits, so this branch is not what makes that surface safe. ...(options.deliveryWatchdog ? { reportRendererDeliveryState: vi.fn(async () => null) } : {}) @@ -64,19 +52,20 @@ describe('parked pty:data delivery credit', () => { } }) - it('holds the ACK for bytes parked with no handler, and repays it exactly once on bind', async () => { + it('credits parked bytes immediately and does not credit them again on bind', async () => { installWindow({ deliveryWatchdog: true }) const { ensurePtyDispatcher, registerEagerPtyBuffer } = await import('./pty-dispatcher') - const { getParkedPreHandlerCharsByPty } = await import('./pty-parked-delivery-debt') + const { getParkedPreHandlerCharsByPty } = await import('./pty-pre-handler-buffer') ensurePtyDispatcher() emitPtyData({ id: PTY_ID, data: BOOT_OUTPUT }) - expect(ackData).not.toHaveBeenCalled() + expect(ackData.mock.calls).toEqual([[PTY_ID, BOOT_OUTPUT.length, BOOT_OUTPUT.length]]) expect(getParkedPreHandlerCharsByPty()).toEqual({ [PTY_ID]: BOOT_OUTPUT.length }) // The real bind seam: registering the eager buffer drains the parked chunks. - registerEagerPtyBuffer(PTY_ID, () => {}) + const handle = registerEagerPtyBuffer(PTY_ID, () => {}) + expect(handle.flush()).toBe(BOOT_OUTPUT) expect(ackData.mock.calls).toEqual([[PTY_ID, BOOT_OUTPUT.length, BOOT_OUTPUT.length]]) expect(getParkedPreHandlerCharsByPty()).toEqual({}) @@ -85,48 +74,48 @@ describe('parked pty:data delivery credit', () => { it('credits the ACK at return on a surface with no delivery watchdog', async () => { installWindow({ deliveryWatchdog: false }) const { ensurePtyDispatcher } = await import('./pty-dispatcher') - const { getParkedPreHandlerCharsByPty } = await import('./pty-parked-delivery-debt') + const { getParkedPreHandlerCharsByPty } = await import('./pty-pre-handler-buffer') const { clearPreHandlerPtyState } = await import('./pty-pre-handler-buffer') ensurePtyDispatcher() emitPtyData({ id: PTY_ID, data: BOOT_OUTPUT }) - // Held debt with no heal lane would be a paused shell nobody can unstick. + // Consumer health remains observable without the optional invoke API. expect(ackData.mock.calls).toEqual([[PTY_ID, BOOT_OUTPUT.length, BOOT_OUTPUT.length]]) - expect(getParkedPreHandlerCharsByPty()).toEqual({}) + expect(getParkedPreHandlerCharsByPty()).toEqual({ [PTY_ID]: BOOT_OUTPUT.length }) clearPreHandlerPtyState(PTY_ID) }) - it('leaves a bound pane on the parse-deferred ACK it already had', async () => { + it('leaves a bound pane on parse-deferred ACK credit', async () => { installWindow({ deliveryWatchdog: true }) const { ensurePtyDispatcher, ptyDataHandlers } = await import('./pty-dispatcher') - const { getParkedPreHandlerCharsByPty } = await import('./pty-parked-delivery-debt') + const { takeCurrentPtyDeliveryAckCredit } = await import('./terminal-pty-ack-gate') ensurePtyDispatcher() - ptyDataHandlers.set(PTY_ID, () => {}) + let consumed: (() => void) | null = null + ptyDataHandlers.set(PTY_ID, () => { + consumed = takeCurrentPtyDeliveryAckCredit() + }) emitPtyData({ id: PTY_ID, data: BOOT_OUTPUT }) - + expect(ackData).not.toHaveBeenCalled() + expect(consumed).toBeTypeOf('function') + ;(consumed as unknown as () => void)() expect(ackData.mock.calls).toEqual([[PTY_ID, BOOT_OUTPUT.length, BOOT_OUTPUT.length]]) - expect(getParkedPreHandlerCharsByPty()).toEqual({}) - ptyDataHandlers.delete(PTY_ID) }) - it('repays parked credit on exit, when no handler will ever drain it', async () => { + it('retains startup output on exit without issuing another credit', async () => { installWindow({ deliveryWatchdog: true }) const { ensurePtyDispatcher } = await import('./pty-dispatcher') - const { getParkedPreHandlerCharsByPty } = await import('./pty-parked-delivery-debt') + const { getParkedPreHandlerCharsByPty, drainPreHandlerPtyData } = + await import('./pty-pre-handler-buffer') ensurePtyDispatcher() - emitPtyData({ id: PTY_ID, data: BOOT_OUTPUT }) - expect(getParkedPreHandlerCharsByPty()).toEqual({ [PTY_ID]: BOOT_OUTPUT.length }) - - // No primary exit handler, so the exit is only buffered and no drain ever runs. Main - // deletes this pty's accounting on exit, so the write-off lane cannot forgive it either: - // un-repaid here, the debt pins the session in-flight total until the window reloads. emitPtyExit({ id: PTY_ID, code: 0 }) - - expect(getParkedPreHandlerCharsByPty()).toEqual({}) - expect(ackData).toHaveBeenCalledWith(PTY_ID, BOOT_OUTPUT.length, BOOT_OUTPUT.length) + expect(getParkedPreHandlerCharsByPty()).toEqual({ [PTY_ID]: BOOT_OUTPUT.length }) + const data: string[] = [] + drainPreHandlerPtyData(PTY_ID, (chunk) => data.push(chunk)) + expect(data).toEqual([BOOT_OUTPUT]) + expect(ackData.mock.calls).toEqual([[PTY_ID, BOOT_OUTPUT.length, BOOT_OUTPUT.length]]) }) it('leaves no processed-char total behind for the next incarnation of a reused id', async () => { @@ -161,18 +150,18 @@ describe('parked pty:data delivery credit', () => { ptyExitHandlers.delete(PTY_ID) }) - it('holds the credit in rawLength chars, not the UTF-8 byte count', async () => { + it('credits rawLength while reporting the actual retained characters', async () => { installWindow({ deliveryWatchdog: true }) const { ensurePtyDispatcher } = await import('./pty-dispatcher') - const { getParkedPreHandlerCharsByPty } = await import('./pty-parked-delivery-debt') + const { getParkedPreHandlerCharsByPty } = await import('./pty-pre-handler-buffer') const { clearPreHandlerPtyState } = await import('./pty-pre-handler-buffer') ensurePtyDispatcher() // Main counts UTF-16 chars and says so via rawLength; the buffer counts UTF-8 bytes. - emitPtyData({ id: PTY_ID, data: 'ééé', rawLength: 3 }) + emitPtyData({ id: PTY_ID, data: 'ééé', rawLength: 9 }) expect(getParkedPreHandlerCharsByPty()).toEqual({ [PTY_ID]: 3 }) clearPreHandlerPtyState(PTY_ID) - expect(ackData.mock.calls).toEqual([[PTY_ID, 3, 3]]) + expect(ackData.mock.calls).toEqual([[PTY_ID, 9, 9]]) }) }) diff --git a/src/renderer/src/components/terminal-pane/pty-dispatcher.ts b/src/renderer/src/components/terminal-pane/pty-dispatcher.ts index 361ce64149c..fdfd10f23f4 100644 --- a/src/renderer/src/components/terminal-pane/pty-dispatcher.ts +++ b/src/renderer/src/components/terminal-pane/pty-dispatcher.ts @@ -4,8 +4,7 @@ import { clearProcessedPtyCharTotal, deliverPtyDataWithDeferredAck, exposeE2eTerminalPtyAckGate, - getProcessedPtyCharTotals, - takeCurrentPtyDeliveryAckCredit + getProcessedPtyCharTotals } from './terminal-pty-ack-gate' import { clampUtf8Tail, type EagerBufferChunk } from './pty-eager-buffer-clamp' import { @@ -16,11 +15,9 @@ import { } from './pty-pre-handler-buffer' import { buildPtyDataMeta, type PtyDataPayload } from './pty-data-meta' import { deliverPtyExitToHandlers } from './pty-exit-delivery' -import { settleParkedPtyDeliveryDebtsForPty } from './pty-parked-delivery-debt' import { clearReceivedPtyCharTotal, isPtyPushDeliveryBlackholed, - isTerminalDeliveryWatchdogArmed, recordPtyDataReceived, startTerminalDeliveryWatchdog } from './terminal-delivery-watchdog' @@ -97,8 +94,8 @@ export function ensurePtyDispatcher(): void { import('./terminal-parked-pane-recovery') .then((module) => module.recoverParkedPanes(ptyIds)) // Why swallowed: a failed chunk load must not reject out of the tick before the heal - // runs. Claiming no owner is the safe answer — the write-off lane still forgives. - .catch(() => []) + // runs. A later tick retries ownership without holding producer credit. + .catch(() => {}) }) } @@ -129,13 +126,8 @@ function handleDispatchedPtyData(payload: PtyDataPayload): void { if (handler) { handler(payload.data, meta) } else { - // Why hold the ACK: crediting bytes no handler will consume left main's in-flight window - // empty for a dead pane, so its flow control read healthy and it kept flooding. Gated on - // the watchdog because held debt with no heal lane is a paused shell nobody can unstick. - bufferPreHandlerPtyData(payload.id, payload.data, meta, { - chars, - settle: isTerminalDeliveryWatchdogArmed() ? takeCurrentPtyDeliveryAckCredit() : null - }) + // No consumer owns parse credit; retain a bounded tail and report occupancy separately. + bufferPreHandlerPtyData(payload.id, payload.data, meta) } const sidecars = ptyDataSidecars.get(payload.id) if (sidecars && sidecars.size > 0) { @@ -172,11 +164,6 @@ function attachPtySecondaryPushListeners(unsubscribes: (() => void)[]): void { // Why: host-initiated remote sleep has no requester transaction in this renderer; classify its ordered exit before pane cleanup runs. markCommittedPtyShutdowns([payload.id]) } - // Why before the totals are cleared: main deletes this pty's accounting on exit, so - // parked credit held past here can be repaid by no drain and forgiven by no write-off. - // An exit with no primary handler never reaches a drain at all, which pinned the debt — - // and with it the session in-flight total — until the window reloaded. - settleParkedPtyDeliveryDebtsForPty(payload.id) const sidecars = ptyExitSidecars.get(payload.id) if (sidecars) { ptyExitSidecars.delete(payload.id) diff --git a/src/renderer/src/components/terminal-pane/pty-parked-delivery-debt.ts b/src/renderer/src/components/terminal-pane/pty-parked-delivery-debt.ts deleted file mode 100644 index 89e10bdfbc3..00000000000 --- a/src/renderer/src/components/terminal-pane/pty-parked-delivery-debt.ts +++ /dev/null @@ -1,105 +0,0 @@ -/** - * Ledger for the pty:data delivery credit held while bytes sit in the pre-handler buffer. - * - * Why hold it at all: ACKing bytes that no handler will ever consume made main's in-flight - * window read empty for a pane rendering nothing. The delivery watchdog needs debt > 0 to - * call a wedge, so it was blind by construction, and the producer was never paused — main - * kept flooding a pane nobody could see. Holding the credit turns parked bytes into real - * debt, which is what makes main's existing per-PTY window pause the shell, the watchdog - * see the pane, and the write-off lane able to forgive it. The pause is not guaranteed: this - * buffer caps at 512KB of UTF-8 BYTES while main's window counts UTF-16 CHARS, so on - * multi-byte output eviction settles the credit before main's window ever fills. - - * - * Unit is delivery-credit CHARS (`rawLength ?? data.length`), never the buffer's UTF-8 byte - * count: the two diverge on multi-byte output and main's accounting is in chars. - */ - -export type ParkedPtyDeliveryDebt = { - /** Idempotent — a second settle ACKs once and decrements the ledger once. */ - settle: () => void -} - -type OpenParkedDebt = { settle: () => void } - -const parkedCharsByPty = new Map() -const openDebtsByPty = new Map>() - -/** Park `chars` of credit for `ptyId`. `settleAck` is the open delivery credit claimed by the - * dispatcher; a null one means this surface keeps today's ACK-at-return behaviour. */ -export function openParkedPtyDeliveryDebt( - ptyId: string, - chars: number, - settleAck: (() => void) | null -): ParkedPtyDeliveryDebt | null { - if (!settleAck) { - return null - } - if (!Number.isFinite(chars) || chars <= 0) { - settleAck() - return null - } - parkedCharsByPty.set(ptyId, (parkedCharsByPty.get(ptyId) ?? 0) + chars) - let settled = false - const open: OpenParkedDebt = { - settle: () => { - if (settled) { - return - } - settled = true - const debts = openDebtsByPty.get(ptyId) - if (debts) { - debts.delete(open) - if (debts.size === 0) { - openDebtsByPty.delete(ptyId) - } - } - const remaining = (parkedCharsByPty.get(ptyId) ?? 0) - chars - if (remaining > 0) { - parkedCharsByPty.set(ptyId, remaining) - } else { - parkedCharsByPty.delete(ptyId) - } - settleAck() - } - } - let debts = openDebtsByPty.get(ptyId) - if (!debts) { - debts = new Set() - openDebtsByPty.set(ptyId, debts) - } - debts.add(open) - return { settle: open.settle } -} - -/** Repay every credit a buffer state still holds. Safe to call twice, and on nothing. */ -export function settleParkedPtyDeliveryDebts( - state: { chunks: { debt?: ParkedPtyDeliveryDebt }[]; head: number } | undefined -): void { - if (!state) { - return - } - for (let index = state.head; index < state.chunks.length; index += 1) { - state.chunks[index].debt?.settle() - } -} - -/** Repay every credit still held for one PTY while leaving its bytes buffered for a late - * bind. Main deletes a PTY's accounting on exit, so debt still held past that point can be - * repaid by no drain and forgiven by no write-off — it would pin the session's in-flight - * total until the window reloaded. */ -export function settleParkedPtyDeliveryDebtsForPty(ptyId: string): void { - const debts = openDebtsByPty.get(ptyId) - if (!debts) { - return - } - for (const debt of Array.from(debts)) { - debt.settle() - } -} - -/** Chars parked per PTY — the discriminator main's write-off skip was missing: bytes with a - * consumer repay themselves, these have none. */ -export function getParkedPreHandlerCharsByPty(): Record { - return Object.fromEntries(parkedCharsByPty) -} diff --git a/src/renderer/src/components/terminal-pane/pty-pre-handler-buffer.ts b/src/renderer/src/components/terminal-pane/pty-pre-handler-buffer.ts index 8017af24c3d..2566f3b556c 100644 --- a/src/renderer/src/components/terminal-pane/pty-pre-handler-buffer.ts +++ b/src/renderer/src/components/terminal-pane/pty-pre-handler-buffer.ts @@ -1,25 +1,18 @@ import { isPtyIncarnationId } from '../../../../shared/pty-incarnation' import { clampUtf8Tail } from './pty-eager-buffer-clamp' -import { - openParkedPtyDeliveryDebt, - settleParkedPtyDeliveryDebts, - type ParkedPtyDeliveryDebt -} from './pty-parked-delivery-debt' import type { PtyDataMeta } from './pty-dispatcher' type BufferedPreHandlerPtyData = { data: string bytes: number meta?: PtyDataMeta - /** Withheld ACK for this chunk. Settled exactly once, on whichever of drain, eviction, - * clear or write-off discard takes the chunk out of the buffer. */ - debt?: ParkedPtyDeliveryDebt } type BufferedPreHandlerPtyState = { chunks: BufferedPreHandlerPtyData[] head: number bytes: number + chars: number /** Sequence of the newest chunk, used to fence state left by a prior incarnation of a reused id. */ sequence: number } @@ -72,22 +65,17 @@ export function currentPreHandlerPtySequence(): number { } /** Map preserves insertion order, so the first key is the least recently admitted id. - * Returns the evicted entry so the caller can drop state keyed alongside it — and, for - * buffered data, repay the credit those now-unreachable chunks were holding. */ -function evictOldestPtyIfAtCap( - map: Map, - ptyId: string, - cap: number -): [string, V] | null { + * Returns the evicted id so the caller can drop state keyed alongside the entry. */ +function evictOldestPtyIfAtCap(map: Map, ptyId: string, cap: number): string | null { if (map.has(ptyId) || map.size < cap) { return null } - const oldest = map.entries().next().value - if (!oldest) { - return null + const oldestPtyId = map.keys().next().value + if (typeof oldestPtyId === 'string') { + map.delete(oldestPtyId) + return oldestPtyId } - map.delete(oldest[0]) - return oldest + return null } /** Drop a buffered exit proven to describe a different lifetime of `ptyId` than the one now @@ -137,29 +125,22 @@ function retainPreHandlerPtyExits( preHandlerPtyExit.set(ptyId, kept) } -/** `ack` carries the dispatcher's claim on this chunk's delivery credit. Every return path - * here must repay it: a claim that never settles is permanent debt, and main answers - * permanent debt by pausing a healthy shell forever. */ -export function bufferPreHandlerPtyData( - ptyId: string, - data: string, - meta?: PtyDataMeta, - ack?: { chars: number; settle: (() => void) | null } -): void { +export function bufferPreHandlerPtyData(ptyId: string, data: string, meta?: PtyDataMeta): void { if (discardedPreHandlerPtyStates.has(ptyId)) { - ack?.settle?.() return } const chunk = clampUtf8Tail(data, PRE_HANDLER_PTY_DATA_MAX_BYTES) if (!chunk.data) { - ack?.settle?.() return } - const evicted = evictOldestPtyIfAtCap(preHandlerPtyData, ptyId, PRE_HANDLER_PTY_DATA_MAX_PTYS) - if (evicted !== null) { + const evictedPtyId = evictOldestPtyIfAtCap( + preHandlerPtyData, + ptyId, + PRE_HANDLER_PTY_DATA_MAX_PTYS + ) + if (evictedPtyId !== null) { // The warn breadcrumb describes buffered bytes that no longer exist. - warnedLostHandlerPtyIds.delete(evicted[0]) - settleParkedPtyDeliveryDebts(evicted[1]) + warnedLostHandlerPtyIds.delete(evictedPtyId) } const bufferedMeta = meta && chunk.data.length !== data.length && typeof meta.rawLength === 'number' @@ -167,26 +148,22 @@ export function bufferPreHandlerPtyData( : meta let state = preHandlerPtyData.get(ptyId) if (!state) { - state = { chunks: [], head: 0, bytes: 0, sequence: 0 } + state = { chunks: [], head: 0, bytes: 0, chars: 0, sequence: 0 } preHandlerPtyData.set(ptyId, state) } state.sequence = nextPreHandlerPtySequence() - const debt = ack ? openParkedPtyDeliveryDebt(ptyId, ack.chars, ack.settle) : null state.chunks.push({ data: chunk.data, bytes: chunk.bytes, - ...(bufferedMeta ? { meta: bufferedMeta } : {}), - ...(debt ? { debt } : {}) + ...(bufferedMeta ? { meta: bufferedMeta } : {}) }) state.bytes += chunk.bytes + state.chars += chunk.data.length // Why: a missing handler can accumulate many small chunks; a stored total // and head index keep that failure path linear instead of rescanning/shifting. while (state.bytes > PRE_HANDLER_PTY_DATA_MAX_BYTES && state.head < state.chunks.length - 1) { state.bytes -= state.chunks[state.head].bytes - // Why settle on evict: this buffer caps UTF-8 BYTES while main's window counts UTF-16 - // chars, so eviction can precede main's pause. An evicted chunk that kept its credit - // would leave debt larger than the parked bytes, which no later drain can repay. - state.chunks[state.head].debt?.settle() + state.chars -= state.chunks[state.head].data.length state.chunks[state.head] = { data: '', bytes: 0 } state.head += 1 } @@ -213,32 +190,12 @@ export function drainPreHandlerPtyData( return } preHandlerPtyData.delete(ptyId) - try { - for (let index = state.head; index < state.chunks.length; index += 1) { - const chunk = state.chunks[index] - handler(chunk.data, chunk.meta) - } - } finally { - // Handing the bytes to the pane's write path is the consume point; settle even when the - // handler throws, because the chunks are already out of the buffer and their debt would - // otherwise have no payer left. - settleParkedPtyDeliveryDebts(state) + for (let index = state.head; index < state.chunks.length; index += 1) { + const chunk = state.chunks[index] + handler(chunk.data, chunk.meta) } } -/** Drop parked bytes for PTYs main just wrote off. Their credit is forgiven (a post-write-off - * ACK is a clamped no-op main-side) and a very late bind must not paint bytes the restore - * marker already superseded. */ -export function discardParkedPtyDataAfterWriteOff(ptyIds: string[]): void { - ptyIds.forEach(dropParkedPreHandlerPtyData) -} - -function dropParkedPreHandlerPtyData(ptyId: string): void { - settleParkedPtyDeliveryDebts(preHandlerPtyData.get(ptyId)) - preHandlerPtyData.delete(ptyId) - warnedLostHandlerPtyIds.delete(ptyId) -} - /** Replay buffered startup bytes without taking them from the future primary handler. */ export function replayPreHandlerPtyData(ptyId: string, observer: (data: string) => void): void { const state = preHandlerPtyData.get(ptyId) @@ -302,7 +259,8 @@ export function discardPreHandlerPtyStateFromPriorIncarnation( retainPreHandlerPtyExits(ptyId, (exit) => exit.sequence > fenceSequence) const data = preHandlerPtyData.get(ptyId) if (data && data.sequence <= fenceSequence) { - dropParkedPreHandlerPtyData(ptyId) + preHandlerPtyData.delete(ptyId) + warnedLostHandlerPtyIds.delete(ptyId) } // Why: the id now names a different, live PTY, so a prior incarnation's consumed/discarded marks // must not suppress this one's own exit — the same admission boundary a same-id reattach gets. @@ -407,7 +365,7 @@ export function drainPreHandlerPtyExit( } export function clearPreHandlerPtyState(ptyId: string): void { - dropParkedPreHandlerPtyData(ptyId) + preHandlerPtyData.delete(ptyId) preHandlerPtyExit.delete(ptyId) consumedPreHandlerPtyExits.delete(ptyId) const discardTimer = discardedPreHandlerPtyStates.get(ptyId) @@ -415,4 +373,18 @@ export function clearPreHandlerPtyState(ptyId: string): void { clearTimeout(discardTimer) } discardedPreHandlerPtyStates.delete(ptyId) + warnedLostHandlerPtyIds.delete(ptyId) +} + +/** Buffer occupancy is a pane-health signal, independent of producer delivery credit. */ +export function getParkedPreHandlerCharsByPty(): Record { + return Object.fromEntries(Array.from(preHandlerPtyData, ([id, state]) => [id, state.chars])) +} + +/** A restore supersedes these bytes, including for a pane that binds later. */ +export function discardParkedPtyDataAfterWriteOff(ptyIds: string[]): void { + for (const id of ptyIds) { + preHandlerPtyData.delete(id) + warnedLostHandlerPtyIds.delete(id) + } } diff --git a/src/renderer/src/components/terminal-pane/pty-pre-handler-parked-debt-invariant.test.ts b/src/renderer/src/components/terminal-pane/pty-pre-handler-parked-debt-invariant.test.ts deleted file mode 100644 index 2242df981c0..00000000000 --- a/src/renderer/src/components/terminal-pane/pty-pre-handler-parked-debt-invariant.test.ts +++ /dev/null @@ -1,170 +0,0 @@ -/** - * The ratchet for the inverse bug. Holding the ACK for parked bytes is only safe while every - * exit from this buffer repays it: a claim that never settles is permanent renderer-held - * debt, and main answers permanent debt by pausing a healthy shell forever. - * - * So: for every sequence, total settled credit == total claimed credit, and nothing stays - * parked. Any future exit path that forgets to settle turns this red. - */ -import { afterEach, describe, expect, it } from 'vitest' -import { - bufferPreHandlerPtyData, - clearConsumedPreHandlerPtyExit, - clearPreHandlerPtyState, - currentPreHandlerPtySequence, - discardParkedPtyDataAfterWriteOff, - discardPreHandlerPtyState, - discardPreHandlerPtyStateFromPriorIncarnation, - drainPreHandlerPtyData -} from './pty-pre-handler-buffer' -import { getParkedPreHandlerCharsByPty } from './pty-parked-delivery-debt' - -const PRE_HANDLER_PTY_DATA_MAX_BYTES = 512 * 1024 -const PRE_HANDLER_PTY_DATA_MAX_PTYS = 64 -const PTY_ID = 'pty-parked-debt' -const LRU_PTY_IDS = Array.from( - { length: PRE_HANDLER_PTY_DATA_MAX_PTYS + 1 }, - (_, index) => `pty-parked-lru-${index}` -) - -/** Stands in for the dispatcher's claim on the open delivery credit: one settle per claim, - * idempotent, exactly like `takeCurrentTerminalDeliveryCredit`. */ -function createDeliveryCreditLedger(): { - claimedChars: () => number - settledChars: () => number - park: (ptyId: string, data: string) => void -} { - let claimedChars = 0 - let settledChars = 0 - return { - claimedChars: () => claimedChars, - settledChars: () => settledChars, - park: (ptyId, data) => { - claimedChars += data.length - let settled = false - bufferPreHandlerPtyData(ptyId, data, undefined, { - chars: data.length, - settle: () => { - if (settled) { - return - } - settled = true - settledChars += data.length - } - }) - } - } -} - -function expectNoOutstandingDebt(ledger: ReturnType): void { - expect(ledger.settledChars()).toBe(ledger.claimedChars()) - expect(getParkedPreHandlerCharsByPty()).toEqual({}) -} - -describe('parked pre-handler delivery debt always settles', () => { - afterEach(() => { - clearConsumedPreHandlerPtyExit(PTY_ID) - clearPreHandlerPtyState(PTY_ID) - for (const ptyId of LRU_PTY_IDS) { - clearPreHandlerPtyState(ptyId) - } - }) - - it('holds the credit while parked and repays it on drain', () => { - const ledger = createDeliveryCreditLedger() - - ledger.park(PTY_ID, 'prompt') - ledger.park(PTY_ID, ' and motd') - expect(getParkedPreHandlerCharsByPty()).toEqual({ - [PTY_ID]: 'prompt'.length + ' and motd'.length - }) - expect(ledger.settledChars()).toBe(0) - - const drained: string[] = [] - drainPreHandlerPtyData(PTY_ID, (data) => drained.push(data)) - - expect(drained).toEqual(['prompt', ' and motd']) - expectNoOutstandingDebt(ledger) - }) - - it('repays evicted chunks: the byte cap can bite before main pauses the producer', () => { - const ledger = createDeliveryCreditLedger() - - // The buffer caps UTF-8 BYTES while main's in-flight window counts UTF-16 chars, so an - // eviction can happen with main still crediting; the evicted credit must not be stranded. - ledger.park(PTY_ID, 'x'.repeat(PRE_HANDLER_PTY_DATA_MAX_BYTES)) - ledger.park(PTY_ID, 'y'.repeat(1024)) - expect(ledger.settledChars()).toBe(PRE_HANDLER_PTY_DATA_MAX_BYTES) - - drainPreHandlerPtyData(PTY_ID, () => {}) - - expectNoOutstandingDebt(ledger) - }) - - it('repays a whole PTY evicted by the per-PTY LRU cap', () => { - const ledger = createDeliveryCreditLedger() - - for (const ptyId of LRU_PTY_IDS) { - ledger.park(ptyId, 'startup') - } - - // The oldest id was pushed out of the map entirely; its bytes will never reach a pane. - expect(getParkedPreHandlerCharsByPty()[LRU_PTY_IDS[0]]).toBeUndefined() - expect(ledger.settledChars()).toBe('startup'.length) - - for (const ptyId of LRU_PTY_IDS) { - drainPreHandlerPtyData(ptyId, () => {}) - } - expectNoOutstandingDebt(ledger) - }) - - it('repays bytes dropped because the PTY state is discarded', () => { - const ledger = createDeliveryCreditLedger() - - discardPreHandlerPtyState(PTY_ID) - ledger.park(PTY_ID, 'kill-flush output') - - // Dropped, exactly as before — but dropping still has to ACK. - expectNoOutstandingDebt(ledger) - }) - - it('repays an empty chunk that never becomes a buffered record', () => { - const ledger = createDeliveryCreditLedger() - - ledger.park(PTY_ID, '') - - expectNoOutstandingDebt(ledger) - }) - - it('repays on clear, on write-off discard, and on the prior-incarnation fence', () => { - const ledger = createDeliveryCreditLedger() - - ledger.park(PTY_ID, 'cleared') - clearPreHandlerPtyState(PTY_ID) - expectNoOutstandingDebt(ledger) - - ledger.park(PTY_ID, 'written-off') - discardParkedPtyDataAfterWriteOff([PTY_ID, 'pty-never-seen']) - expectNoOutstandingDebt(ledger) - - ledger.park(PTY_ID, 'previous incarnation') - discardPreHandlerPtyStateFromPriorIncarnation(PTY_ID, currentPreHandlerPtySequence()) - expectNoOutstandingDebt(ledger) - }) - - it('repays the whole PTY when a draining handler throws', () => { - const ledger = createDeliveryCreditLedger() - - ledger.park(PTY_ID, 'first') - ledger.park(PTY_ID, 'second') - - expect(() => - drainPreHandlerPtyData(PTY_ID, () => { - throw new Error('xterm write threw') - }) - ).toThrow('xterm write threw') - - // The chunks are already out of the buffer, so their debt would have no payer left. - expectNoOutstandingDebt(ledger) - }) -}) diff --git a/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog-parked-stall.test.ts b/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog-parked-stall.test.ts index 9b99bedda62..7e8d5402e4d 100644 --- a/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog-parked-stall.test.ts +++ b/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog-parked-stall.test.ts @@ -25,12 +25,11 @@ const LIVE_PTY_ID = 'pty-live' const WEDGED_TAB_ID = 'tab-wedged' const WEDGED_OUTPUT = 'output nobody renders' -/** Main is healthy overall — some other pane ACKed a moment ago — but holds this pane's debt. */ +/** Parked output has already returned all producer credit. */ const BUSY_MAIN: PtyRendererDeliveryHealthReply = { - inFlightTotalChars: WEDGED_OUTPUT.length, - inFlightPtyCount: 1, - msSinceLastAck: 200, - stalledPtys: [{ id: WEDGED_PTY_ID, inFlightChars: WEDGED_OUTPUT.length, msSinceLastAck: null }] + inFlightTotalChars: 0, + inFlightPtyCount: 0, + msSinceLastAck: 200 } type PtyDataPayload = { id: string; data: string } @@ -96,7 +95,7 @@ describe('per-PTY parked delivery stall', () => { const dispatcher = await import('./pty-dispatcher') dispatcher.ptyDataHandlers.set(LIVE_PTY_ID, () => {}) dispatcher.ensurePtyDispatcher() - const { getParkedPreHandlerCharsByPty } = await import('./pty-parked-delivery-debt') + const { getParkedPreHandlerCharsByPty } = await import('./pty-pre-handler-buffer') return { parkedCharsByPty: getParkedPreHandlerCharsByPty, streamLiveOutput: () => emitPtyData({ id: LIVE_PTY_ID, data: 'still streaming' }) @@ -120,53 +119,69 @@ describe('per-PTY parked delivery stall', () => { streamLiveOutput() await vi.advanceTimersByTimeAsync(INTERVAL_MS) expect(storeState.remountTerminalTabForRecovery).not.toHaveBeenCalled() - expect(reportMock.mock.calls[0]![0]).toMatchObject({ - parkedCharsByPty: { [WEDGED_PTY_ID]: WEDGED_OUTPUT.length } - }) + expect(reportMock.mock.calls[0]![0]).not.toHaveProperty('parkedCharsByPty') streamLiveOutput() await vi.advanceTimersByTimeAsync(INTERVAL_MS) - // The remount is the heal for an owned pane: rebinding drains the parked bytes and - // their held ACK repays the debt. No push listener churn for one pane. + // Local occupancy drives recovery even with zero main debt. expect(storeState.remountTerminalTabForRecovery).toHaveBeenCalledWith(WEDGED_TAB_ID) expect(reattachMock).not.toHaveBeenCalled() expect(healCalls()).toHaveLength(0) }) - it('writes off an orphan-parked pane and drops the superseded bytes', async () => { - reportMock.mockImplementation((args) => - Promise.resolve( - (args as { heal?: boolean }).heal - ? { - inFlightTotalChars: 0, - inFlightPtyCount: 0, - msSinceLastAck: 0, - writtenOff: [{ id: WEDGED_PTY_ID, writtenOffChars: WEDGED_OUTPUT.length }] - } - : BUSY_MAIN - ) - ) + it('keeps unowned output bounded without requesting a credit write-off', async () => { + reportMock.mockResolvedValue(BUSY_MAIN) const { parkedCharsByPty, streamLiveOutput } = await startDispatcherAndWatchdog() - - // No tab owns this pty, so there is nothing to remount; only a write-off frees the debt. - emitPtyData({ id: WEDGED_PTY_ID, data: WEDGED_OUTPUT }) + emitPtyData({ id: WEDGED_PTY_ID, data: 'x'.repeat(512 * 1024) }) streamLiveOutput() await vi.advanceTimersByTimeAsync(INTERVAL_MS) + emitPtyData({ id: WEDGED_PTY_ID, data: 'tail' }) streamLiveOutput() await vi.advanceTimersByTimeAsync(INTERVAL_MS) - + expect(parkedCharsByPty()).toEqual({ [WEDGED_PTY_ID]: 4 }) expect(storeState.remountTerminalTabForRecovery).not.toHaveBeenCalled() - expect(healCalls()).toHaveLength(1) - expect(healCalls()[0]).toMatchObject({ - heal: true, - parkedCharsByPty: { [WEDGED_PTY_ID]: WEDGED_OUTPUT.length } - }) - // A very late bind must not paint bytes the restore marker already superseded. - expect(parkedCharsByPty()).toEqual({}) + expect(healCalls()).toHaveLength(0) expect(reattachMock).not.toHaveBeenCalled() }) + it('recovers local occupancy even when the main health reply is unavailable', async () => { + storeState.ptyIdsByTabId = { [WEDGED_TAB_ID]: [WEDGED_PTY_ID] } + reportMock.mockResolvedValue(null) + await startDispatcherAndWatchdog() + emitPtyData({ id: WEDGED_PTY_ID, data: WEDGED_OUTPUT }) + await vi.advanceTimersByTimeAsync(INTERVAL_MS * 2) + expect(storeState.remountTerminalTabForRecovery).toHaveBeenCalledWith(WEDGED_TAB_ID) + }) + + it('does not mistake byte-cap eviction for consumer progress', async () => { + storeState.ptyIdsByTabId = { [WEDGED_TAB_ID]: [WEDGED_PTY_ID] } + reportMock.mockResolvedValue(BUSY_MAIN) + await startDispatcherAndWatchdog() + emitPtyData({ id: WEDGED_PTY_ID, data: 'x'.repeat(512 * 1024) }) + await vi.advanceTimersByTimeAsync(INTERVAL_MS) + emitPtyData({ id: WEDGED_PTY_ID, data: 'tail' }) + await vi.advanceTimersByTimeAsync(INTERVAL_MS) + expect(storeState.remountTerminalTabForRecovery).toHaveBeenCalledWith(WEDGED_TAB_ID) + }) + + it('heals an unreceived PTY while siblings stream and nothing is parked', async () => { + reportMock.mockResolvedValue({ + ...BUSY_MAIN, + inFlightTotalChars: 100, + inFlightPtyCount: 1, + stalledPtys: [{ id: WEDGED_PTY_ID, inFlightChars: 100, msSinceLastAck: null }] + }) + const { streamLiveOutput } = await startDispatcherAndWatchdog() + streamLiveOutput() + await vi.advanceTimersByTimeAsync(INTERVAL_MS) + streamLiveOutput() + await vi.advanceTimersByTimeAsync(INTERVAL_MS) + expect(healCalls()).toHaveLength(1) + expect(reattachMock).not.toHaveBeenCalled() + expect(storeState.remountTerminalTabForRecovery).not.toHaveBeenCalled() + }) + it('leaves the ordinary pre-attach race alone: parked bytes that drain never heal', async () => { storeState.ptyIdsByTabId = { [WEDGED_TAB_ID]: [WEDGED_PTY_ID] } reportMock.mockResolvedValue(BUSY_MAIN) diff --git a/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.test.ts b/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.test.ts index d76fd90be06..b8035b809f5 100644 --- a/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.test.ts +++ b/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.test.ts @@ -71,7 +71,7 @@ describe('terminal delivery watchdog', () => { watchdog.startTerminalDeliveryWatchdog({ reattachPushListeners: reattachMock, hasAttachedPtys: () => true, - recoverParkedPanes: async () => [] + recoverParkedPanes: async () => {} }) return { recordPtyDataReceived: watchdog.recordPtyDataReceived, @@ -81,7 +81,8 @@ describe('terminal delivery watchdog', () => { } } - it('does zero IPC while pty output is flowing', async () => { + it('probes per-PTY health while output flows without healing', async () => { + reportMock.mockResolvedValue(HEALTHY) const { recordPtyDataReceived } = await startWatchdog() for (let tick = 0; tick < 8; tick++) { @@ -89,7 +90,8 @@ describe('terminal delivery watchdog', () => { await vi.advanceTimersByTimeAsync(INTERVAL_MS) } - expect(reportMock).not.toHaveBeenCalled() + expect(reportMock).toHaveBeenCalledTimes(8) + expect(reattachMock).not.toHaveBeenCalled() }) it('reports during silence but never heals a healthy-idle main', async () => { @@ -125,11 +127,11 @@ describe('terminal delivery watchdog', () => { // Bytes flowed once, then the push channel died: the field shape. recordPtyDataReceived('pty-1', 128) await vi.advanceTimersByTimeAsync(INTERVAL_MS) - expect(reportMock).not.toHaveBeenCalled() + expect(reportMock).toHaveBeenCalledTimes(1) // First silent tick: report only, no heal yet. await vi.advanceTimersByTimeAsync(INTERVAL_MS) - expect(reportMock).toHaveBeenCalledTimes(1) + expect(reportMock).toHaveBeenCalledTimes(2) expect(reattachMock).not.toHaveBeenCalled() // Second silent tick confirms: re-attach precedes the heal report, the @@ -178,7 +180,7 @@ describe('terminal delivery watchdog', () => { watchdog.startTerminalDeliveryWatchdog({ reattachPushListeners: reattachMock, hasAttachedPtys: () => false, - recoverParkedPanes: async () => [] + recoverParkedPanes: async () => {} }) await vi.advanceTimersByTimeAsync(INTERVAL_MS * 3) @@ -193,7 +195,7 @@ describe('terminal delivery watchdog', () => { watchdog.startTerminalDeliveryWatchdog({ reattachPushListeners: reattachMock, hasAttachedPtys: () => true, - recoverParkedPanes: async () => [] + recoverParkedPanes: async () => {} }) await vi.advanceTimersByTimeAsync(INTERVAL_MS * 3) diff --git a/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.ts b/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.ts index a428abcf83c..ef7f3ffaae9 100644 --- a/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.ts +++ b/src/renderer/src/components/terminal-pane/terminal-delivery-watchdog.ts @@ -11,15 +11,17 @@ * channel that is dead (upstream precedent: electron#37067, one-directional * Mojo IPC death). This watchdog is the missing lane: it detects the wedge * and heals over invoke — the direction proven alive — with zero cost on the - * data hot path (one Map upsert per received chunk; a tick does no IPC while - * output flows or while no PTY delivery is expected). + * data hot path (one Map upsert per received chunk; probes run every 15 seconds + * while PTY delivery is expected). */ import { e2eConfig } from '@/lib/e2e-config' import type { PtyRendererDeliveryHealthReply } from '../../../../shared/pty-renderer-delivery-health' import { redactPtyIdForDiagnostics } from '../../../../shared/pty-delivery-diagnostics' import { deliverPulledPtyModelRestoreMarkers } from './pty-model-restore-channel' -import { getParkedPreHandlerCharsByPty } from './pty-parked-delivery-debt' -import { discardParkedPtyDataAfterWriteOff } from './pty-pre-handler-buffer' +import { + getParkedPreHandlerCharsByPty, + discardParkedPtyDataAfterWriteOff +} from './pty-pre-handler-buffer' import { advanceParkedDeliveryStallStreaks, advanceStalledPtyStreaks, @@ -48,11 +50,11 @@ type TerminalDeliveryWatchdogDeps = { reattachPushListeners: () => void /** True while any PTY handler or eager buffer expects push delivery. */ hasAttachedPtys: () => boolean - /** Remount the tabs owning these parked ptys; answers which ids a tab actually owned. + /** Remount the tabs owning these parked ptys. * Injected, and loaded lazily by the dispatcher, because resolving ownership reads the app * store: a static edge would close an import cycle and drag the store into the * freeze-report graph, which runs in bare-node contexts. */ - recoverParkedPanes: (ptyIds: string[]) => Promise + recoverParkedPanes: (ptyIds: string[]) => Promise } const receivedPtyCharTotals = new Map() @@ -90,10 +92,6 @@ export function isPtyPushDeliveryBlackholed(): boolean { return blackholePtyPushDelivery } -export function isTerminalDeliveryWatchdogArmed(): boolean { - return watchdogDeps !== null -} - function isMainDeliveryStalled(health: PtyRendererDeliveryHealthReply): boolean { // Why msSinceLastAck may be null: a wedged-from-first-byte session (the // field case's brand-new terminal) never ACKs; in-flight debt alone is the @@ -117,19 +115,7 @@ async function runWatchdogTick(): Promise { eventCountAtLastTick = receivedPtyDataEventCount const parkedCharsByPty = getParkedPreHandlerCharsByPty() const hasParked = Object.keys(parkedCharsByPty).length > 0 - // Keeps "no IPC while output flows": an empty parked map is the cheap proof there is no - // per-pane debt to ask main about. - if (!hasParked && (globalEventsMoved || !deps.hasAttachedPtys())) { - stallStreakTicks = 0 - resetParkedDeliveryStallStreaks() - return - } - const health = await report({ - receivedCharsByPty: Object.fromEntries(receivedPtyCharTotals), - processedCharsByPty: getProcessedPtyCharTotals(), - ...(hasParked ? { parkedCharsByPty } : {}) - }) - if (!health) { + if (!hasParked && !deps.hasAttachedPtys()) { stallStreakTicks = 0 resetParkedDeliveryStallStreaks() return @@ -138,17 +124,22 @@ async function runWatchdogTick(): Promise { parkedCharsByPty, watchdogConfig.stallTicksToHeal ) - // A remount is the heal for an owned pane: rebinding drains the parked bytes and their held - // ACK repays the debt. Its own budget/cooldown is the churn control, so it is not gated on - // the write-off cooldown below. - const ownedParkedPtyIds = new Set( - stalledParkedPtyIds.length > 0 ? await deps.recoverParkedPanes(stalledParkedPtyIds) : [] - ) - const orphanParkedPtyIds = stalledParkedPtyIds.filter((id) => !ownedParkedPtyIds.has(id)) - // An owned parked pane is already routed to the remount; letting it also drive a write-off - // would forgive the very bytes the rebind is about to drain. + // Local consumer recovery does not depend on main's debt or a successful health reply. + if (stalledParkedPtyIds.length > 0) { + await deps.recoverParkedPanes(stalledParkedPtyIds) + } + // Probe attached panes even while a sibling streams: global movement cannot rule out + // a single PTY whose push events never arrive. + const health = await report({ + receivedCharsByPty: Object.fromEntries(receivedPtyCharTotals), + processedCharsByPty: getProcessedPtyCharTotals() + }) + if (!health) { + stallStreakTicks = 0 + return + } const perPtyStalled = advanceStalledPtyStreaks( - health.stalledPtys?.filter((entry) => !ownedParkedPtyIds.has(entry.id)), + health.stalledPtys, receivedPtyCharTotals, watchdogConfig.stallTicksToHeal ) @@ -165,32 +156,20 @@ async function runWatchdogTick(): Promise { }) } const globalWedge = stallStreakTicks >= watchdogConfig.stallTicksToHeal - if (!globalWedge && !perPtyStalled && orphanParkedPtyIds.length === 0) { + if (!globalWedge && !perPtyStalled) { return } if (lastHealAtMs !== null && Date.now() - lastHealAtMs < watchdogConfig.healCooldownMs) { return } - await healDeadPushDelivery(deps, report, health, { - parkedCharsByPty, - // Only the session-wide wedge implicates the push listeners; a single parked pane - // must not churn every pane's subscription. - reattachPushListeners: globalWedge, - ownedParkedPtyCount: ownedParkedPtyIds.size, - orphanParkedPtyCount: orphanParkedPtyIds.length - }) + await healDeadPushDelivery(deps, report, health, globalWedge) } async function healDeadPushDelivery( deps: TerminalDeliveryWatchdogDeps, report: NonNullable, stalled: PtyRendererDeliveryHealthReply, - context: { - parkedCharsByPty: Record - reattachPushListeners: boolean - ownedParkedPtyCount: number - orphanParkedPtyCount: number - } + reattachPushListeners: boolean ): Promise { lastHealAtMs = Date.now() stallStreakTicks = 0 @@ -200,15 +179,12 @@ async function healDeadPushDelivery( // bug to hunt); ≥1 = events are being dropped below the emitter (channel // dead, platform-level). The single most valuable field discriminator. const listenerCountBeforeReattach = window.api?.pty?.getPtyDataListenerCount?.() ?? null - if (context.reattachPushListeners) { + if (reattachPushListeners) { deps.reattachPushListeners() } const healed = await report({ receivedCharsByPty: Object.fromEntries(receivedPtyCharTotals), processedCharsByPty: getProcessedPtyCharTotals(), - ...(Object.keys(context.parkedCharsByPty).length > 0 - ? { parkedCharsByPty: context.parkedCharsByPty } - : {}), heal: true, rendererPtyDataListenerCount: listenerCountBeforeReattach }) @@ -226,8 +202,6 @@ async function healDeadPushDelivery( ) } recordTerminalFreezeBreadcrumb('watchdog-heal', { - orphanParkedPtyCount: context.orphanParkedPtyCount, - ownedParkedPtyCount: context.ownedParkedPtyCount, listenerCountBeforeReattach, writtenOffPtyCount: writtenOff.length, writtenOffChars: writtenOff.reduce((sum, entry) => sum + entry.writtenOffChars, 0) diff --git a/src/renderer/src/components/terminal-pane/terminal-pane-recovery.ts b/src/renderer/src/components/terminal-pane/terminal-pane-recovery.ts index 1be9d267897..f504ab35d36 100644 --- a/src/renderer/src/components/terminal-pane/terminal-pane-recovery.ts +++ b/src/renderer/src/components/terminal-pane/terminal-pane-recovery.ts @@ -34,7 +34,7 @@ export type TerminalPaneRecoveryReason = // binding. pty:data for the old id then lands in the pre-handler buffer, which // ACKs it — main's delivery health stays green while the pane shows nothing. | 'spawn-left-pane-unbound' - // pty:data kept arriving for a pane with no handler and sat parked, un-ACKed, for + // pty:data kept arriving for a pane with no handler and sat parked for // two watchdog ticks. Skips the liveness probe like 'input-rejected-by-host', but for a // simpler reason: it renders no verdict on the PTY at all. A remount preserves the PTY // whether or not it is still alive, so nothing here reads silence as death — which is what diff --git a/src/renderer/src/components/terminal-pane/terminal-parked-delivery-stall.ts b/src/renderer/src/components/terminal-pane/terminal-parked-delivery-stall.ts index ed84620e32c..e9dae6a26a6 100644 --- a/src/renderer/src/components/terminal-pane/terminal-parked-delivery-stall.ts +++ b/src/renderer/src/components/terminal-pane/terminal-parked-delivery-stall.ts @@ -18,8 +18,7 @@ type StallStreak = { previous: number; ticks: number } const parkedStreakByPty = new Map() const wedgedStreakByPty = new Map() -/** Advance one id's streak. Progress — parked chars falling, or received chars moving — is - * evidence a consumer exists, and resets it. */ +/** Advance a sampled stall streak, restarting when its progress predicate changes. */ function advanceStreak( streaks: Map, id: string, @@ -44,9 +43,9 @@ function retainStreaks(streaks: Map, liveIds: Set): } } -/** Ids whose held debt has not shrunk for `stallTicksToHeal` ticks. A drain zeroes a pty's - * parked total, so "still parked, no smaller" is the honest evidence that nothing consumed - * it — new bytes arriving for the same dead pane do not make it healthier. */ +/** Ids whose buffer has remained occupied for `stallTicksToHeal` ticks. A drain zeroes a pty's + * parked total, so "still parked" is the honest evidence that nothing consumed + * it — byte-cap eviction is not consumer progress. */ export function advanceParkedDeliveryStallStreaks( parkedCharsByPty: Record, stallTicksToHeal: number @@ -57,10 +56,7 @@ export function advanceParkedDeliveryStallStreaks( if (chars <= 0) { continue } - if ( - advanceStreak(parkedStreakByPty, ptyId, chars, (previous) => chars >= previous) >= - stallTicksToHeal - ) { + if (advanceStreak(parkedStreakByPty, ptyId, chars, () => true) >= stallTicksToHeal) { stalled.push(ptyId) } } diff --git a/src/renderer/src/components/terminal-pane/terminal-parked-pane-ownership.test.ts b/src/renderer/src/components/terminal-pane/terminal-parked-pane-ownership.test.ts deleted file mode 100644 index 47091f20450..00000000000 --- a/src/renderer/src/components/terminal-pane/terminal-parked-pane-ownership.test.ts +++ /dev/null @@ -1,93 +0,0 @@ -/** - * "Owned" has to mean the held ACK has a payer. - * - * The watchdog excludes owned ids from BOTH heal lanes — the write-off lane and the per-PTY - * stalled lane — on the grounds that a remount will drain their bytes and repay the credit. - * So an id reported as owned when no remount can happen is excluded forever: main keeps that - * pty's in-flight window full and pauses a perfectly healthy shell, with no path back. - * - * Ownership therefore tracks the recovery's own verdict, not the store's `ptyIdsByTabId`, - * which only says a tab once listed the id. The converse matters just as much: a request that - * was declined but re-queued IS owned, because its bytes are about to be drained — reporting - * it unowned would hand the write-off lane output the imminent rebind was going to render. - */ -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' - -const PTY_ID = 'pty-parked-owner' -const TAB_ID = 'tab-parked-owner' - -const storeState = { - ptyIdsByTabId: {} as Record, - getTab: vi.fn((tabId: string) => ({ id: tabId }) as { id: string; viewMode?: string }), - remountTerminalTabForRecovery: vi.fn((_tabId: string) => true) -} - -vi.mock('@/store', () => ({ useAppStore: { getState: () => storeState } })) -vi.mock('@/lib/crash-breadcrumb-recorder', () => ({ recordRendererCrashBreadcrumb: vi.fn() })) - -describe('parked pane ownership', () => { - let warnSpy: ReturnType - - beforeEach(() => { - vi.resetModules() - vi.useFakeTimers() - storeState.ptyIdsByTabId = { [TAB_ID]: [PTY_ID] } - storeState.getTab = vi.fn((tabId: string) => ({ id: tabId })) - storeState.remountTerminalTabForRecovery = vi.fn(() => true) - warnSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}) - }) - - afterEach(() => { - warnSpy.mockRestore() - vi.useRealTimers() - }) - - async function recover(): Promise { - const { recoverParkedPanes } = await import('./terminal-parked-pane-recovery') - return recoverParkedPanes([PTY_ID]) - } - - it('claims a pty whose tab actually remounted', async () => { - expect(await recover()).toEqual([PTY_ID]) - expect(storeState.remountTerminalTabForRecovery).toHaveBeenCalledWith(TAB_ID) - }) - - it('disclaims a pty whose tab left the store, so the write-off lane can forgive it', async () => { - // remountTerminalTabForRecovery answers false when the tab is gone. Recovery does not - // even re-queue this one — its own comment says retrying is pointless — so if ownership - // still claimed it, nothing anywhere would ever repay the credit. - storeState.remountTerminalTabForRecovery = vi.fn(() => false) - - expect(await recover()).toEqual([]) - }) - - it('disclaims a pty whose tab is in chat view, which refuses recovery unconditionally', async () => { - storeState.getTab = vi.fn((tabId: string) => ({ id: tabId, viewMode: 'chat' as const })) - - expect(await recover()).toEqual([]) - expect(storeState.remountTerminalTabForRecovery).not.toHaveBeenCalled() - }) - - it('keeps claiming a pty whose remount was declined but re-queued', async () => { - // The first call consumes the tab's recovery budget; the second lands inside the cooldown - // and is re-queued rather than refused. Its bytes are about to be drained by the pending - // retry, so writing them off here would discard output, not rescue it. - const { recoverParkedPanes } = await import('./terminal-parked-pane-recovery') - const { hasPendingTerminalPaneRecovery } = await import('./terminal-pane-recovery') - expect(await recoverParkedPanes([PTY_ID])).toEqual([PTY_ID]) - - storeState.remountTerminalTabForRecovery = vi.fn(() => true) - const claimed = await recoverParkedPanes([PTY_ID]) - - expect(storeState.remountTerminalTabForRecovery).not.toHaveBeenCalled() - expect(hasPendingTerminalPaneRecovery(TAB_ID)).toBe(true) - expect(claimed).toEqual([PTY_ID]) - }) - - it('disclaims a pty no tab lists at all', async () => { - storeState.ptyIdsByTabId = {} - - expect(await recover()).toEqual([]) - expect(storeState.remountTerminalTabForRecovery).not.toHaveBeenCalled() - }) -}) diff --git a/src/renderer/src/components/terminal-pane/terminal-parked-pane-recovery.ts b/src/renderer/src/components/terminal-pane/terminal-parked-pane-recovery.ts index 90af61d585f..7b2cc273d86 100644 --- a/src/renderer/src/components/terminal-pane/terminal-parked-pane-recovery.ts +++ b/src/renderer/src/components/terminal-pane/terminal-parked-pane-recovery.ts @@ -1,35 +1,19 @@ /** - * Ownership and remount for PTYs whose bytes have been parked, un-ACKed, across the watchdog's + * Ownership and remount for PTYs whose bytes have been parked across the watchdog's * stall streak. Wired into the watchdog as a dep by `pty-dispatcher.ts` rather than imported * by it: the watchdog is on the freeze-report path and must stay clear of the app store. */ import { useAppStore } from '@/store' import { captureTerminalPaneRecoveryGeneration, - hasPendingTerminalPaneRecovery, requestTerminalPaneRecovery } from './terminal-pane-recovery' -/** Remount the tabs owning `ptyIds` and return the ids whose debt a remount will actually - * repay. The rest belong to the write-off lane. - * - * Why the recovery's own verdict decides this, not the store: "some tab lists this id" is a - * store fact, while whether a remount can happen is a runtime one — the tab may have left - * `tabsByWorktree`, or sit in chat view, which refuses unconditionally. Reporting the store - * fact as ownership excluded such ids from BOTH heal lanes on every tick, so their held ACK - * had no payer at all and main kept a healthy shell paused. A retry still queued counts as - * owned: its bytes are about to be drained, not lost. - * - * No liveness probe, and none is needed: this infers nothing about whether the PTY is alive. - * The stall predicate is "held debt has not shrunk", which is evidence about the past — the - * bytes arrived over a live path — not a claim about the present. The action is a - * renderer-local remount that preserves the PTY either way, so nothing here reads silence as - * death, which is what keeps it honest across the SSH execution boundary. The recovery - * budget/cooldown is the anti-churn control. */ -export async function recoverParkedPanes(ptyIds: string[]): Promise { +/** Remount owned panes without inferring PTY liveness or changing producer flow control. + * Recovery's budget/cooldown bounds churn; unowned bytes remain in the bounded buffer. */ +export async function recoverParkedPanes(ptyIds: string[]): Promise { // Bounded: this scan runs only for ids stalled across two ticks, never per render. const ptyIdsByTabId = useAppStore.getState().ptyIdsByTabId ?? {} - const owned: string[] = [] for (const ptyId of ptyIds) { const tabId = Object.keys(ptyIdsByTabId).find((candidate) => ptyIdsByTabId[candidate]?.includes(ptyId) @@ -37,15 +21,11 @@ export async function recoverParkedPanes(ptyIds: string[]): Promise { if (tabId === undefined) { continue } - const remounted = await requestTerminalPaneRecovery({ + await requestTerminalPaneRecovery({ tabId, ptyId, reason: 'delivery-parked', terminalRecoveryGeneration: captureTerminalPaneRecoveryGeneration(tabId) }) - if (remounted || hasPendingTerminalPaneRecovery(tabId)) { - owned.push(ptyId) - } } - return owned } diff --git a/src/shared/pty-renderer-delivery-health.ts b/src/shared/pty-renderer-delivery-health.ts index a644ce096bf..c1ccc8d6e5b 100644 --- a/src/shared/pty-renderer-delivery-health.ts +++ b/src/shared/pty-renderer-delivery-health.ts @@ -19,15 +19,6 @@ export type PtyRendererDeliveryStateReport = { * ACK path and resync response carry; merging them here is a free extra * repair lane for the lost-ACK variant. */ processedCharsByPty: Record - /** Chars CURRENTLY parked for a PTY with no registered data handler — a live balance, - * decremented as each chunk settles, not a cumulative total like the fields either side - * of it. `writeOffLostRendererDelivery` subtracts it from a cumulative `receivedChars` - * precisely because of that: what is still parked is what cannot repay itself. Their ACK - * is withheld, so — - * unlike received-but-unparsed bytes, which their own deferred ACK repays — - * this debt has no consumer to repay it and only a write-off or a bind clears - * it. Absent means "none parked", which is exactly how an older renderer read. */ - parkedCharsByPty?: Record /** Set on the confirming tick: main may write off provably-lost bytes and * answer with restore markers for the renderer to route locally. */ heal?: boolean