refactor(terminal): separate parked pane health from producer credit

This commit is contained in:
Merge Sim
2026-09-07 18:37:34 -07:00
parent 0a704b0ced
commit 404cee070e
16 changed files with 179 additions and 701 deletions
+3 -57
View File
@@ -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
})
@@ -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
})
+1 -5
View File
@@ -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)
@@ -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]])
})
})
@@ -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)
@@ -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<string, number>()
const openDebtsByPty = new Map<string, Set<OpenParkedDebt>>()
/** 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<string, number> {
return Object.fromEntries(parkedCharsByPty)
}
@@ -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<V>(
map: Map<string, V>,
ptyId: string,
cap: number
): [string, V] | null {
* Returns the evicted id so the caller can drop state keyed alongside the entry. */
function evictOldestPtyIfAtCap<V>(map: Map<string, V>, 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<string, number> {
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)
}
}
@@ -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<typeof createDeliveryCreditLedger>): 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)
})
})
@@ -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)
@@ -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)
@@ -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<string[]>
recoverParkedPanes: (ptyIds: string[]) => Promise<void>
}
const receivedPtyCharTotals = new Map<string, number>()
@@ -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<void> {
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<void> {
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<void> {
})
}
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<Window['api']['pty']['reportRendererDeliveryState']>,
stalled: PtyRendererDeliveryHealthReply,
context: {
parkedCharsByPty: Record<string, number>
reattachPushListeners: boolean
ownedParkedPtyCount: number
orphanParkedPtyCount: number
}
reattachPushListeners: boolean
): Promise<void> {
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)
@@ -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
@@ -18,8 +18,7 @@ type StallStreak = { previous: number; ticks: number }
const parkedStreakByPty = new Map<string, StallStreak>()
const wedgedStreakByPty = new Map<string, StallStreak>()
/** 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<string, StallStreak>,
id: string,
@@ -44,9 +43,9 @@ function retainStreaks(streaks: Map<string, StallStreak>, liveIds: Set<string>):
}
}
/** 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<string, number>,
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)
}
}
@@ -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<string, string[]>,
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<typeof vi.spyOn>
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<string[]> {
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()
})
})
@@ -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<string[]> {
/** 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<void> {
// 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<string[]> {
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
}
@@ -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<string, number>
/** 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<string, number>
/** 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