mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
Fix post-heal terminal parser credit and SSH snapshot replay
This commit is contained in:
@@ -261,6 +261,110 @@ describe('registerPtyHandlers', () => {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
it.each(['cumulative ACK', 'health repair', 'legacy ACK'])(
|
||||
'holds fresh parser bytes after writing off old loss until %s',
|
||||
(creditPath) => {
|
||||
vi.useFakeTimers()
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
try {
|
||||
const provider = installObservableDaemonTestProvider()
|
||||
registerPtyHandlers(mainWindow as never)
|
||||
const id = 'post-heal-pty'
|
||||
provider.emitData(id, 'x'.repeat(201))
|
||||
vi.advanceTimersByTime(100)
|
||||
const lost = reportRendererDeliveryState({
|
||||
receivedCharsByPty: {},
|
||||
processedCharsByPty: {},
|
||||
heal: true
|
||||
})
|
||||
expect(lost.writtenOff).toEqual([{ id, writtenOffChars: 201 }])
|
||||
provider.emitData(id, 'y'.repeat(199))
|
||||
vi.advanceTimersByTime(10_000)
|
||||
const held = reportRendererDeliveryState({
|
||||
receivedCharsByPty: { [id]: 199 },
|
||||
processedCharsByPty: {},
|
||||
heal: true
|
||||
})
|
||||
expect(held.writtenOff).toBeUndefined()
|
||||
expect(held.inFlightTotalChars).toBe(199)
|
||||
if (creditPath === 'health repair') {
|
||||
reportRendererDeliveryState({
|
||||
receivedCharsByPty: { [id]: 199 },
|
||||
processedCharsByPty: { [id]: 199 }
|
||||
})
|
||||
} else {
|
||||
getPtyAckDataListener()(null, {
|
||||
id,
|
||||
...(creditPath === 'legacy ACK' ? { charCount: 199 } : { processedChars: 199 })
|
||||
})
|
||||
}
|
||||
expect(getPtyRendererDeliveryDebugSnapshot().rendererInFlightChars).toBe(0)
|
||||
provider.emitData(id, 'z'.repeat(17))
|
||||
vi.advanceTimersByTime(10_000)
|
||||
// Repeated reports must neither credit fresh loss nor count the old writeoff twice.
|
||||
const nextLoss = reportRendererDeliveryState({
|
||||
receivedCharsByPty: { [id]: 199 },
|
||||
processedCharsByPty: { [id]: 199 },
|
||||
heal: true
|
||||
})
|
||||
expect(nextLoss.writtenOff).toEqual([{ id, writtenOffChars: 17 }])
|
||||
provider.emitData(id, 'fresh')
|
||||
vi.advanceTimersByTime(10_000)
|
||||
expect(
|
||||
reportRendererDeliveryState({
|
||||
receivedCharsByPty: { [id]: 204 },
|
||||
processedCharsByPty: { [id]: 199 },
|
||||
heal: true
|
||||
})
|
||||
).toMatchObject({ inFlightTotalChars: 5 })
|
||||
} finally {
|
||||
warn.mockRestore()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
}
|
||||
)
|
||||
|
||||
it.each(['exit', 'navigation'])(
|
||||
'discards writeoff coordinates on %s before the PTY id is reused',
|
||||
(lifecycle) => {
|
||||
vi.useFakeTimers()
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
try {
|
||||
const provider = installObservableDaemonTestProvider()
|
||||
registerPtyHandlers(mainWindow as never)
|
||||
const id = 'reused-healed-pty'
|
||||
provider.emitData(id, 'x'.repeat(201))
|
||||
vi.advanceTimersByTime(100)
|
||||
expect(
|
||||
reportRendererDeliveryState({
|
||||
receivedCharsByPty: {},
|
||||
processedCharsByPty: {},
|
||||
heal: true
|
||||
}).writtenOff
|
||||
).toEqual([{ id, writtenOffChars: 201 }])
|
||||
if (lifecycle === 'exit') {
|
||||
provider.emitExit(id)
|
||||
} else {
|
||||
getMainFrameNavigationListener()()
|
||||
getPtyRendererDispatcherReadyListener()()
|
||||
}
|
||||
provider.emitData(id, 'fresh')
|
||||
vi.advanceTimersByTime(100)
|
||||
expect(
|
||||
reportRendererDeliveryState({
|
||||
receivedCharsByPty: { [id]: 5 },
|
||||
processedCharsByPty: {}
|
||||
}).inFlightTotalChars
|
||||
).toBe(5)
|
||||
getPtyAckDataListener()(null, { id, processedChars: 1 })
|
||||
expect(getPtyRendererDeliveryDebugSnapshot().rendererInFlightChars).toBe(4)
|
||||
} finally {
|
||||
warn.mockRestore()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
}
|
||||
)
|
||||
|
||||
it('refuses a heal while main has seen a recent ACK', async () => {
|
||||
vi.useFakeTimers()
|
||||
const mockProc = createMockProc()
|
||||
|
||||
@@ -17,6 +17,7 @@ describe('bounded delivery health recovery candidates', () => {
|
||||
accounting.set(id, {
|
||||
sentChars: 2000,
|
||||
ackedChars: 1000,
|
||||
writtenOffChars: 0,
|
||||
lastSendAtMs: now,
|
||||
lastAckAtMs: siblingState === 'streaming' ? now : now - 60_000
|
||||
})
|
||||
@@ -25,6 +26,7 @@ describe('bounded delivery health recovery candidates', () => {
|
||||
accounting.set('lost', {
|
||||
sentChars: 100,
|
||||
ackedChars: 0,
|
||||
writtenOffChars: 0,
|
||||
lastSendAtMs: now - 60_000,
|
||||
lastAckAtMs: null
|
||||
})
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import { createPtyIpcSession } from './pty/session'
|
||||
import { PtyPendingDataDrainQueue } from './pty-pending-data-drain-queue'
|
||||
import { sendPtyDataToRenderer } from './pty/delivery/payload'
|
||||
import { writeOffLostRendererDelivery } from './pty/delivery/accounting'
|
||||
import { handleRendererDeliveryStateReport } from './pty/delivery/renderer-delivery-state-report'
|
||||
import {
|
||||
createSshPtyOutputIntakeHarness,
|
||||
sshPtyOutputEvent
|
||||
} from './ssh-pty-output-intake-test-harness'
|
||||
|
||||
vi.mock('./pty/provider/registry', () => ({ tryGetProviderForPty: () => undefined }))
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks()
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
describe('SSH delivery heal parser credit', () => {
|
||||
it('settles old 201 loss but holds the fresh 199 projection until parser completion', async () => {
|
||||
vi.useFakeTimers()
|
||||
vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
const session = createPtyIpcSession({
|
||||
mainWindow: { webContents: { send: vi.fn() } } as never
|
||||
})
|
||||
session.pendingData = new PtyPendingDataDrainQueue(() => 'active')
|
||||
session.readCurrentPtyRendererDeliveryDebugSnapshot = () => ({}) as never
|
||||
session.schedulePendingDataAfterCreditReport = vi.fn()
|
||||
session.writeOffLostRendererDelivery = (report, ids) =>
|
||||
writeOffLostRendererDelivery(session, report, ids)
|
||||
const id = 'ssh:connection@@pty-1'
|
||||
const harness = createSshPtyOutputIntakeHarness({
|
||||
project: (event, projection) => {
|
||||
sendPtyDataToRenderer(session, event.id, event, [projection.identity.projectionSemanticsId])
|
||||
}
|
||||
})
|
||||
session.sshOutputIntake = harness.intake
|
||||
try {
|
||||
const lost = harness.intake.acceptData(
|
||||
sshPtyOutputEvent({
|
||||
id,
|
||||
data: 'x'.repeat(201),
|
||||
rawLength: 201
|
||||
})
|
||||
)
|
||||
harness.completions[0]!.resolve()
|
||||
await lost
|
||||
expect(
|
||||
handleRendererDeliveryStateReport(session, {
|
||||
receivedCharsByPty: {},
|
||||
processedCharsByPty: {},
|
||||
heal: true
|
||||
}).writtenOff
|
||||
).toEqual([{ id, writtenOffChars: 201 }])
|
||||
expect(harness.intake.getDebugSnapshot().projection.records).toBe(0)
|
||||
|
||||
const fresh = harness.intake.acceptData(
|
||||
sshPtyOutputEvent({
|
||||
id,
|
||||
data: 'y'.repeat(199),
|
||||
rawLength: 199
|
||||
})
|
||||
)
|
||||
harness.completions[1]!.resolve()
|
||||
await fresh
|
||||
vi.advanceTimersByTime(10_000)
|
||||
const held = handleRendererDeliveryStateReport(session, {
|
||||
receivedCharsByPty: { [id]: 199 },
|
||||
processedCharsByPty: {},
|
||||
heal: true
|
||||
})
|
||||
expect(held.writtenOff).toBeUndefined()
|
||||
expect(held.inFlightTotalChars).toBe(199)
|
||||
expect(harness.intake.getDebugSnapshot().projection.records).toBe(1)
|
||||
const settled = handleRendererDeliveryStateReport(session, {
|
||||
receivedCharsByPty: { [id]: 199 },
|
||||
processedCharsByPty: { [id]: 199 }
|
||||
})
|
||||
expect(settled.inFlightTotalChars).toBe(0)
|
||||
expect(harness.intake.getDebugSnapshot().projection.records).toBe(0)
|
||||
} finally {
|
||||
harness.intake.dispose()
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -160,6 +160,16 @@ export function applyCumulativeAck(
|
||||
return acknowledged
|
||||
}
|
||||
|
||||
// Renderer totals exclude lost pushes; main totals retain them as written-off credit.
|
||||
export function applyRendererCumulativeAck(
|
||||
session: PtyIpcSession,
|
||||
id: string,
|
||||
processedChars: number
|
||||
): number {
|
||||
const writtenOffChars = session.rendererDeliveryAccountingByPty.get(id)?.writtenOffChars ?? 0
|
||||
return applyCumulativeAck(session, id, processedChars + writtenOffChars)
|
||||
}
|
||||
|
||||
export function schedulePendingDataAfterCreditReport(
|
||||
session: PtyIpcSession,
|
||||
creditedAny: boolean
|
||||
@@ -217,7 +227,7 @@ export function hasUnreceivedRendererDelivery(
|
||||
): boolean {
|
||||
return (
|
||||
accounting.sentChars > accounting.ackedChars &&
|
||||
sanitizeReportedChars(receivedChars) <= accounting.ackedChars
|
||||
sanitizeReportedChars(receivedChars) + accounting.writtenOffChars <= accounting.ackedChars
|
||||
)
|
||||
}
|
||||
|
||||
@@ -239,6 +249,7 @@ export function writeOffLostRendererDelivery(
|
||||
if (acknowledged <= 0) {
|
||||
continue
|
||||
}
|
||||
accounting.writtenOffChars += acknowledged
|
||||
tryGetProviderForPty(id)?.acknowledgeDataEvent(id, acknowledged)
|
||||
// Why drop pending: everything at/before markerSeq comes from the snapshot, so flushing pre-marker bytes would double-paint the restore.
|
||||
const pending = session.pendingData.get(id)
|
||||
|
||||
@@ -71,6 +71,7 @@ export function sendPtyDataToRenderer(
|
||||
session.rendererDeliveryAccountingByPty.set(id, {
|
||||
sentChars: charCount,
|
||||
ackedChars: 0,
|
||||
writtenOffChars: 0,
|
||||
lastSendAtMs: Date.now(),
|
||||
lastAckAtMs: null
|
||||
})
|
||||
|
||||
@@ -10,7 +10,7 @@ import type {
|
||||
} from '../../../../shared/pty-renderer-delivery-health'
|
||||
import { tryGetProviderForPty } from '../provider/registry'
|
||||
import {
|
||||
applyCumulativeAck,
|
||||
applyRendererCumulativeAck,
|
||||
collectAckSilentPtyIdsForHeal,
|
||||
hasAckSilentRendererDeliveryDebt,
|
||||
hasUnreceivedRendererDelivery,
|
||||
@@ -32,7 +32,7 @@ export function applyRendererProcessedCharTotals(
|
||||
if (typeof processedChars !== 'number' || !Number.isFinite(processedChars)) {
|
||||
continue
|
||||
}
|
||||
const acknowledged = applyCumulativeAck(session, id, Math.max(0, processedChars))
|
||||
const acknowledged = applyRendererCumulativeAck(session, id, Math.max(0, processedChars))
|
||||
if (acknowledged > 0) {
|
||||
creditedAny = true
|
||||
tryGetProviderForPty(id)?.acknowledgeDataEvent(id, acknowledged)
|
||||
|
||||
@@ -26,7 +26,7 @@ import {
|
||||
mainDeliveryBreadcrumbs,
|
||||
resetRendererDeliveryAccountingForLifecycleReset
|
||||
} from '../delivery/debug'
|
||||
import { applyCumulativeAck } from '../delivery/accounting'
|
||||
import { applyCumulativeAck, applyRendererCumulativeAck } from '../delivery/accounting'
|
||||
import {
|
||||
applyRendererProcessedCharTotals,
|
||||
handleRendererDeliveryStateReport
|
||||
@@ -105,7 +105,11 @@ export function installPtyResizeVisibilityIpc(session: PtyIpcSession): void {
|
||||
session.deliveryResyncUnansweredWarnLogged = false
|
||||
let acknowledged = 0
|
||||
if (typeof args.processedChars === 'number' && Number.isFinite(args.processedChars)) {
|
||||
acknowledged = applyCumulativeAck(session, args.id, Math.max(0, args.processedChars))
|
||||
acknowledged = applyRendererCumulativeAck(
|
||||
session,
|
||||
args.id,
|
||||
Math.max(0, args.processedChars)
|
||||
)
|
||||
} else {
|
||||
// Why: tolerate legacy per-chunk delta payloads — dev hot-reload can pair an old renderer with a new main.
|
||||
const accounting = session.rendererDeliveryAccountingByPty.get(args.id)
|
||||
|
||||
@@ -34,6 +34,7 @@ export type PtyDataPayload = {
|
||||
export type RendererPtyDeliveryAccounting = {
|
||||
sentChars: number
|
||||
ackedChars: number
|
||||
writtenOffChars: number
|
||||
lastSendAtMs: number
|
||||
lastAckAtMs: number | null
|
||||
}
|
||||
|
||||
+52
-4
@@ -213,6 +213,54 @@ describe('connectPanePty', () => {
|
||||
expect(api.pty.signal).toHaveBeenCalledWith('leaf-session', 'SIGWINCH')
|
||||
})
|
||||
|
||||
it('paints a bound parked SSH normal-buffer snapshot once after reattach', async () => {
|
||||
const { connectPanePty } = await import('./pty-connection')
|
||||
const sshPtyId = toAppSshPtyId('conn-1', 'relay-pty-1')
|
||||
const reattach = createDeferred<{ id: string; isReattach: true }>()
|
||||
const transport = createMockTransport(sshPtyId)
|
||||
transport.connect.mockReturnValue(reattach.promise)
|
||||
transportFactoryQueue.push(transport)
|
||||
vi.mocked(window.api.pty.getMainBufferSnapshot).mockResolvedValue({
|
||||
data: 'PARKED-NORMAL-HISTORY\r\nPARKED-NORMAL-SCREEN\r\n',
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
seq: 123,
|
||||
source: 'headless',
|
||||
alternateScreen: false
|
||||
})
|
||||
mockStoreState = {
|
||||
...mockStoreState,
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: sshPtyId }] },
|
||||
ptyIdsByTabId: { 'tab-1': [sshPtyId] },
|
||||
repos: [{ id: 'repo1', connectionId: 'conn-1' }],
|
||||
sshConnectionStates: new Map([['conn-1', { status: 'connected' }]]),
|
||||
deferredSshSessionIdsByTabId: { 'tab-1': sshPtyId }
|
||||
}
|
||||
const pane = createPane(1)
|
||||
const { writes, parseCallbacks } = captureCallbackTerminalWrites(pane)
|
||||
const binding = connectPanePty(
|
||||
pane as never,
|
||||
createManager(1) as never,
|
||||
createDeps({
|
||||
mountFollowsTerminalPark: true,
|
||||
restoredLeafId: LEAF_1,
|
||||
restoredPtyIdByLeafId: { [LEAF_1]: sshPtyId }
|
||||
}) as never
|
||||
)
|
||||
await flushAsyncTicks(20)
|
||||
expect(transport.connect).toHaveBeenCalledOnce()
|
||||
expect(window.api.pty.getMainBufferSnapshot).not.toHaveBeenCalled()
|
||||
reattach.resolve({ id: sshPtyId, isReattach: true })
|
||||
for (let step = 0; step < 40; step += 1) {
|
||||
parseCallbacks.shift()?.()
|
||||
await flushAsyncTicks(2)
|
||||
}
|
||||
expect(window.api.pty.getMainBufferSnapshot).toHaveBeenCalledOnce()
|
||||
expect(writes.join('').match(/PARKED-NORMAL-HISTORY/g)).toHaveLength(1)
|
||||
expect(writes.join('').match(/PARKED-NORMAL-SCREEN/g)).toHaveLength(1)
|
||||
binding.dispose()
|
||||
})
|
||||
|
||||
it('keeps a too-wide parked SSH alt frame while no live process can repaint it', async () => {
|
||||
const { connectPanePty } = await import('./pty-connection')
|
||||
const sshPtyId = toAppSshPtyId('conn-1', 'relay-pty-1')
|
||||
@@ -231,7 +279,7 @@ describe('connectPanePty', () => {
|
||||
})
|
||||
mockStoreState = {
|
||||
...mockStoreState,
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: sshPtyId }] },
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: null }] },
|
||||
ptyIdsByTabId: { 'tab-1': [sshPtyId] },
|
||||
repos: [{ id: 'repo1', connectionId: 'conn-1' }],
|
||||
sshConnectionStates: new Map([['conn-1', { status: 'disconnected' }]]),
|
||||
@@ -294,7 +342,7 @@ describe('connectPanePty', () => {
|
||||
vi.mocked(window.api.ssh.connect).mockReturnValue(sshConnect.promise)
|
||||
mockStoreState = {
|
||||
...mockStoreState,
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: sshPtyId, generation: 7 }] },
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: null, generation: 7 }] },
|
||||
ptyIdsByTabId: { 'tab-1': [sshPtyId] },
|
||||
repos: [{ id: 'repo1', connectionId: 'conn-1' }],
|
||||
sshConnectionStates: new Map([
|
||||
@@ -373,7 +421,7 @@ describe('connectPanePty', () => {
|
||||
})
|
||||
mockStoreState = {
|
||||
...mockStoreState,
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: foreignPtyId }] },
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: null }] },
|
||||
ptyIdsByTabId: { 'tab-1': [foreignPtyId] },
|
||||
repos: [{ id: 'repo1', connectionId: 'conn-1' }],
|
||||
sshConnectionStates: new Map([['conn-1', { status: 'disconnected' }]]),
|
||||
@@ -425,7 +473,7 @@ describe('connectPanePty', () => {
|
||||
vi.mocked(window.api.pty.getMainBufferSnapshot).mockReturnValue(snapshot.promise)
|
||||
mockStoreState = {
|
||||
...mockStoreState,
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: sshPtyId }] },
|
||||
tabsByWorktree: { 'wt-1': [{ id: 'tab-1', ptyId: null }] },
|
||||
ptyIdsByTabId: { 'tab-1': [sshPtyId] },
|
||||
repos: [{ id: 'repo1', connectionId: 'conn-1' }],
|
||||
sshConnectionStates: new Map([['conn-1', { status: 'connected' }]]),
|
||||
|
||||
@@ -63,8 +63,10 @@ export function runDeferredSessionAttach(session: ConnectPanePtySession): void {
|
||||
)
|
||||
const legacyWorkerOwnsPane = session.isLegacyWorkerAutomaticResumeBlocked()
|
||||
if (gate.enterDeferredFlow && (!legacyWorkerOwnsPane || !gate.sshConnected)) {
|
||||
// Paint main's parked model while SSH recovery continues off the render path.
|
||||
session.prepaintParkedSshSnapshot(pendingSessionId)
|
||||
// Bound parked panes paint within reattach; only unbound deferred sessions need early paint.
|
||||
if (tabPtyId !== pendingSessionId) {
|
||||
session.prepaintParkedSshSnapshot(pendingSessionId)
|
||||
}
|
||||
void (async () => {
|
||||
// Why: for a passphrase target with no cached credential, don't auto-fire ssh.connect — a prompt popping just from focusing a tab / Cmd+J would surprise the user.
|
||||
// Wait for a user-initiated connect first; no-passphrase targets return false here and auto-connect as before.
|
||||
|
||||
Reference in New Issue
Block a user