Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/main/daemon/daemon-checkpoint-file.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,5 +20,7 @@ export type TerminalCheckpointFile = {
/** Ties this checkpoint to the output.log whose header carries the same
* generation. Absent on checkpoints written before incremental logs. */
generation?: number
/** Last daemon pending-output batch included in this checkpoint. */
pendingOutputSeq?: number
checkpointedAt: string
}
159 changes: 159 additions & 0 deletions src/main/daemon/daemon-durable-history-snapshot.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
import { ColdRestoreReplayWriter } from './cold-restore-replay-writer'
import { DAEMON_RESTORE_SCROLLBACK_ROWS } from './daemon-restore-scrollback-depth'
import { HeadlessEmulator } from './headless-emulator'
import { isValidTerminalHistorySize } from './terminal-history-dimensions'
import type { ColdRestoreInfo } from './terminal-history-cold-restore-info'
import type { PendingOutputRecord, TerminalSnapshot } from './types'

type RestoreBase = {
scrollbackAnsi: string
rehydrateSequences: string
snapshotAnsi: string
pendingEscapeTailAnsi?: string
oscLinks?: TerminalSnapshot['oscLinks']
lastTitle?: string
cwd: string | null
cols: number
rows: number
}

export function terminalSnapshotFromColdRestore(
info: ColdRestoreInfo,
opts?: { outputSequence?: number; frameRestoreAnsi?: string }
): TerminalSnapshot {
return {
snapshotAnsi: info.snapshotAnsi,
scrollbackAnsi: info.modes.alternateScreen ? info.scrollbackAnsi : '',
oscLinks: info.oscLinks,
rehydrateSequences: info.rehydrateSequences,
...(info.pendingEscapeTailAnsi ? { pendingEscapeTailAnsi: info.pendingEscapeTailAnsi } : {}),
...(opts?.frameRestoreAnsi ? { frameRestoreAnsi: opts.frameRestoreAnsi } : {}),
cwd: info.cwd,
modes: info.modes,
cols: info.cols,
rows: info.rows,
scrollbackLines:
info.scrollbackLines ?? Math.max(0, countAnsiRows(info.scrollbackAnsi) - info.rows),
...(info.lastTitle ? { lastTitle: info.lastTitle } : {}),
...(opts?.outputSequence !== undefined ? { outputSequence: opts.outputSequence } : {})
}
}

export async function buildDurableCheckpointSnapshot(opts: {
liveSnapshot: TerminalSnapshot
restoreInfo: ColdRestoreInfo | null
pendingRecords?: readonly PendingOutputRecord[]
scrollbackRows?: number
}): Promise<TerminalSnapshot> {
const pendingRecords = opts.pendingRecords ?? []
if (!opts.restoreInfo && pendingRecords.length === 0) {
return opts.liveSnapshot
}
if (
opts.restoreInfo &&
pendingRecords.length === 0 &&
(opts.scrollbackRows === undefined || opts.scrollbackRows >= DAEMON_RESTORE_SCROLLBACK_ROWS)
) {
return terminalSnapshotFromColdRestore(opts.restoreInfo, {
outputSequence: opts.liveSnapshot.outputSequence,
frameRestoreAnsi: opts.liveSnapshot.frameRestoreAnsi
})
}

const emulator = new HeadlessEmulator({
cols: opts.restoreInfo?.cols ?? opts.liveSnapshot.cols,
rows: opts.restoreInfo?.rows ?? opts.liveSnapshot.rows,
scrollback: Math.min(
opts.scrollbackRows ?? DAEMON_RESTORE_SCROLLBACK_ROWS,
DAEMON_RESTORE_SCROLLBACK_ROWS
)
})
const replay = new ColdRestoreReplayWriter(emulator)
try {
// Why not seed the live window when there is no disk history: pending records
// are the raw stream. Replaying them on top of the already-truncated live
// snapshot would duplicate the newest rows and evict the older recoverable ones.
if (opts.restoreInfo) {
const base = restoreBaseFrom(opts.restoreInfo)
for (const segment of [
base.scrollbackAnsi,
base.rehydrateSequences,
base.snapshotAnsi,
base.pendingEscapeTailAnsi ?? ''
]) {
if (!(await replay.write(segment))) {
return opts.liveSnapshot
}
}
emulator.setRestoredOscLinks(base.oscLinks)
if (base.lastTitle) {
emulator.setLastTitle(base.lastTitle)
}
emulator.setCwd(base.cwd)
}
if (!(await replayPendingRecords(replay, pendingRecords))) {
return opts.liveSnapshot
}
const snapshot = emulator.getSnapshot()
return {
...snapshot,
...(opts.liveSnapshot.outputSequence !== undefined
? { outputSequence: opts.liveSnapshot.outputSequence }
: {}),
...(opts.liveSnapshot.frameRestoreAnsi && !snapshot.frameRestoreAnsi
? { frameRestoreAnsi: opts.liveSnapshot.frameRestoreAnsi }
: {})
}
} catch (error) {
console.warn('[history] durable snapshot rebuild failed:', error)
return opts.liveSnapshot
} finally {
emulator.dispose()
}
}

function restoreBaseFrom(restoreInfo: ColdRestoreInfo): RestoreBase {
return {
scrollbackAnsi: restoreInfo.modes.alternateScreen ? restoreInfo.scrollbackAnsi : '',
rehydrateSequences: restoreInfo.rehydrateSequences,
snapshotAnsi: restoreInfo.snapshotAnsi,
...(restoreInfo.pendingEscapeTailAnsi
? { pendingEscapeTailAnsi: restoreInfo.pendingEscapeTailAnsi }
: {}),
oscLinks: restoreInfo.oscLinks,
lastTitle: restoreInfo.lastTitle,
cwd: restoreInfo.cwd,
cols: restoreInfo.cols,
rows: restoreInfo.rows
}
}

async function replayPendingRecords(
replay: ColdRestoreReplayWriter,
records: readonly PendingOutputRecord[]
): Promise<boolean> {
for (const record of records) {
if (record.kind === 'output') {
if (!(await replay.write(record.data))) {
return false
}
continue
}
if (record.kind === 'resize') {
if (!isValidTerminalHistorySize(record.cols, record.rows)) {
return false
}
await replay.resize(record.cols, record.rows)
continue
}
await replay.clearScrollback()
}
return true
}

function countAnsiRows(ansi: string): number {
if (ansi.length === 0) {
return 0
}
return ansi.split(/\r\n|\n|\r/).filter((row) => row.length > 0).length
}
34 changes: 19 additions & 15 deletions src/main/daemon/daemon-pty-adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2717,7 +2717,8 @@ describe('DaemonPtyAdapter (IPtyProvider)', () => {

expect(checkpointSpy).toHaveBeenCalledWith(
id,
expect.objectContaining({ snapshotAnsi: expect.stringContaining('latest before sleep') })
expect.objectContaining({ snapshotAnsi: expect.stringContaining('latest before sleep') }),
{ pendingOutputSeq: expect.any(Number) }
)
expect(existsSync(join(historyDir, getHistorySessionDirName(id)))).toBe(true)

Expand Down Expand Up @@ -2850,7 +2851,8 @@ describe('DaemonPtyAdapter (IPtyProvider)', () => {

expect(checkpointSpy).toHaveBeenCalledWith(
id,
expect.objectContaining({ snapshotAnsi: expect.stringContaining('fresh output') })
expect.objectContaining({ snapshotAnsi: expect.stringContaining('fresh output') }),
{ pendingOutputSeq: expect.any(Number) }
)
expect(existsSync(join(historyDir, getHistorySessionDirName(id)))).toBe(true)
})
Expand All @@ -2866,14 +2868,18 @@ describe('DaemonPtyAdapter (IPtyProvider)', () => {
env: { SHELL: '/bin/zsh' },
sessionId: 'sleep-checkpoint-tail'
})
const appendSpy = vi.spyOn(historyAdapter.getHistoryManager()!, 'appendIncrements')
const checkpointSpy = vi.spyOn(historyAdapter.getHistoryManager()!, 'checkpoint')

lastSubprocess._simulateData('\x1b]777;orca-shell-ready')
await historyAdapter.shutdown(id, { immediate: true, keepHistory: true })

expect(appendSpy).toHaveBeenCalledWith(id, expect.any(Number), [
{ kind: 'output', data: '\x1b]777;orca-shell-ready' }
])
expect(checkpointSpy).toHaveBeenCalledWith(
id,
expect.objectContaining({
pendingEscapeTailAnsi: expect.stringContaining('\x1b]777;orca-shell-ready')
}),
{ pendingOutputSeq: expect.any(Number) }
)
})

it('returns cold restore data when disk history has unclean shutdown', async () => {
Expand Down Expand Up @@ -3261,7 +3267,8 @@ describe('DaemonPtyAdapter (IPtyProvider)', () => {
expect(appendSpy).not.toHaveBeenCalled()
expect(checkpointSpy).toHaveBeenCalledWith(
sessionId,
expect.objectContaining({ snapshotAnsi: expect.stringContaining('revived session') })
expect.objectContaining({ snapshotAnsi: expect.stringContaining('revived session') }),
{ pendingOutputSeq: expect.any(Number) }
)

// Subsequent ticks return to incremental appends.
Expand Down Expand Up @@ -3320,8 +3327,8 @@ describe('DaemonPtyAdapter (IPtyProvider)', () => {
}
const result = await historyAdapter.spawn({ cols: 80, rows: 24, sessionId })
expect(result.isReattach).toBe(true)
expect(internals.sessionsNeedingFullCheckpoint.has(sessionId)).toBe(true)
expect(internals.lastFullCheckpointAt.has(sessionId)).toBe(false)
expect(internals.sessionsNeedingFullCheckpoint.has(sessionId)).toBe(false)
expect(internals.lastFullCheckpointAt.has(sessionId)).toBe(true)
})

it('skips the cold-restore replay when the daemon session is still alive', async () => {
Expand All @@ -3332,20 +3339,17 @@ describe('DaemonPtyAdapter (IPtyProvider)', () => {
await first.disconnectOnly()

historyAdapter = new DaemonPtyAdapter({ socketPath, tokenPath, historyPath: historyDir })
const reader = (historyAdapter as unknown as { historyReader: HistoryReader }).historyReader
const detectSpy = vi.spyOn(reader, 'detectColdRestore')
const result = await historyAdapter.spawn({ cols: 80, rows: 24, sessionId })

expect(result.isReattach).toBe(true)
expect(result.coldRestore).toBeUndefined()
expect(detectSpy).not.toHaveBeenCalled()
// The unmanaged-reattach re-anchor must survive the skipped detect.
// The unmanaged-reattach re-anchor must survive a live remount overlay.
const internals = historyAdapter as unknown as {
sessionsNeedingFullCheckpoint: Set<string>
lastFullCheckpointAt: Map<string, number>
}
expect(internals.sessionsNeedingFullCheckpoint.has(sessionId)).toBe(true)
expect(internals.lastFullCheckpointAt.has(sessionId)).toBe(false)
expect(internals.sessionsNeedingFullCheckpoint.has(sessionId)).toBe(false)
expect(internals.lastFullCheckpointAt.has(sessionId)).toBe(true)
})

it('does not probe session aliveness when there is no restorable history', async () => {
Expand Down
Loading