diff --git a/packages/host-kit/src/session-paths.test.ts b/packages/host-kit/src/session-paths.test.ts new file mode 100644 index 0000000000..be0c8d3116 --- /dev/null +++ b/packages/host-kit/src/session-paths.test.ts @@ -0,0 +1,19 @@ +import { test } from 'vitest'; +import assert from 'node:assert/strict'; +// oxlint-disable-next-line no-restricted-imports -- asserts a path under os.homedir +import os from 'node:os'; +import path from 'node:path'; +import { expandSessionPath } from './session-paths.ts'; + +test('expandSessionPath resolves tilde, relative-with-cwd, and absolute paths', () => { + const homePath = expandSessionPath('~/flows/replay.ad'); + assert.equal(homePath.startsWith(os.homedir()), true); + assert.equal(homePath.endsWith(path.join('flows', 'replay.ad')), true); + + const relativePath = expandSessionPath('workflows/replay.ad', '/tmp/agent-device-cwd'); + assert.equal(relativePath, path.resolve('/tmp/agent-device-cwd', 'workflows/replay.ad')); + + const absoluteInput = path.resolve('/tmp', 'agent-device-absolute.ad'); + const absolutePath = expandSessionPath(absoluteInput, '/tmp/ignored-cwd'); + assert.equal(absolutePath, absoluteInput); +}); diff --git a/src/__tests__/daemon-registration-owner.test.ts b/src/__tests__/daemon-registration-owner.test.ts index edb42abced..7f4795f6d8 100644 --- a/src/__tests__/daemon-registration-owner.test.ts +++ b/src/__tests__/daemon-registration-owner.test.ts @@ -1,16 +1,23 @@ import assert from 'node:assert/strict'; import fs from 'node:fs'; import { afterEach, test, vi } from 'vitest'; -import { readCurrentOwnerIdentity } from '@agent-device/host-kit/process'; +import { readCurrentOwnerIdentity, isProcessAlive } from '@agent-device/host-kit/process'; import { tryAcquireDaemonRegistration, stopAndRetireDaemon, recoverAbandonedDaemonRegistration, + createOwnedReplayStateDir, + DAEMON_STARTUP_EXIT_CODES, + launchDaemonProcess, + type OwnedReplayStateDir, } from '../daemon-registration-owner.ts'; import { resolveDaemonPaths, type DaemonPaths } from '../daemon-resolution.ts'; import { readRegisteredDaemonOwnership } from '../daemon-registration.ts'; import { readDaemonShutdownReport } from '../daemon-shutdown-report.ts'; import { mkdtempForTestSync } from './test-utils/tmp-dir.ts'; +import { registeredDaemonFixtureArgs } from './test-utils/registered-daemon-fixture.ts'; +import { sleep } from '@agent-device/host-kit/retry'; +import { stopDaemonProcess } from '../daemon-process.ts'; const fields = { socketPort: 4210, @@ -305,3 +312,156 @@ test.skipIf(process.getuid?.() === 0)( } }, ); + +test('a forged private-directory capability cannot authorize even matching dead metadata removal', async () => { + const paths = resolveDaemonPaths(mkdtempForTestSync('agent-device-private-forgery-')); + replaceInfo(paths, deadIdentity.pid, deadIdentity.startTime); + const before = fs.readFileSync(paths.infoPath, 'utf8'); + const result = await stopAndRetireDaemon({ + paths, + observed: deadIdentity, + mode: 'force', + ownedStateDir: Object.freeze({ paths }) as OwnedReplayStateDir, + }); + assert.equal(result.status, 'retained'); + assert.equal(fs.readFileSync(paths.infoPath, 'utf8'), before); +}); + +test('private retirement closes startup admission, joins the actual child and never recreates a removed directory', async () => { + const ownedStateDir = createOwnedReplayStateDir(); + const paths = ownedStateDir.paths; + const args = registeredDaemonFixtureArgs(paths, fields); + const launch = launchDaemonProcess({ paths, args, serverMode: 'socket', ownedStateDir }); + let contender: ReturnType | undefined; + try { + assert.ok(launch.startTime); + await waitForFixtureFile(paths.infoPath); + contender = launchDaemonProcess({ paths, args, serverMode: 'socket', ownedStateDir }); + assert.equal((await contender.exited).exitCode, DAEMON_STARTUP_EXIT_CODES.busy); + const input = { + paths, + observed: { pid: launch.pid, startTime: launch.startTime }, + mode: 'graceful' as const, + ownedStateDir, + }; + const pending = stopAndRetireDaemon(input); + assert.throws( + () => launchDaemonProcess({ paths, args, serverMode: 'socket', ownedStateDir }), + (error: { details?: { reason?: string } }) => + error.details?.reason === 'daemon_startup_admission_closed', + ); + const result = await pending; + assert.equal(result.status, 'retired', JSON.stringify(result)); + if (result.status !== 'retired') assert.fail('retirement not confirmed'); + assert.equal(result.removedStateDir, true); + assert.equal((await launch.exited).exitCode, 0); + assert.equal(fs.existsSync(paths.baseDir), false); + assert.deepEqual(await stopAndRetireDaemon(input), result); + assert.equal(fs.existsSync(paths.baseDir), false); + } finally { + await finishPrivateTestDaemons(paths, launch, contender); + } +}); + +test('private retirement retains the directory while an earlier actual startup child is still paused', async () => { + const ownedStateDir = createOwnedReplayStateDir(); + const paths = ownedStateDir.paths; + const args = registeredDaemonFixtureArgs(paths, fields); + const entry = args[1]!; + const ready = `${paths.baseDir}/paused-startup.ready`; + fs.writeFileSync( + entry, + `import fs from 'node:fs'; fs.writeFileSync(${JSON.stringify(ready)}, 'ready'); setInterval(() => {}, 1000);`, + ); + const first = launchDaemonProcess({ paths, args, serverMode: 'socket', ownedStateDir }); + let second: ReturnType | undefined; + try { + await waitForFixtureFile(ready); + second = launchDaemonProcess({ + paths, + args: registeredDaemonFixtureArgs(paths, fields), + serverMode: 'socket', + ownedStateDir, + }); + await waitForFixtureFile(paths.infoPath); + const result = await stopAndRetireDaemon({ + paths, + observed: { pid: second.pid, startTime: second.startTime ?? null }, + mode: 'graceful', + ownedStateDir, + startupJoinTimeoutMs: 0, + }); + assert.equal(result.status, 'retained', JSON.stringify(result)); + if (result.status !== 'retained') assert.fail('private state unexpectedly retired'); + assert.equal(result.reason, 'startup-unconfirmed'); + assert.equal(result.termination?.status, 'exited'); + assert.equal(fs.existsSync(paths.baseDir), true); + assert.equal(isProcessAlive(first.pid), true); + assert.equal((await second.exited).exitCode, 0); + } finally { + await finishPrivateTestDaemons(paths, first, second); + } +}); + +for (const contents of [ + '{broken', + '{"owner":"default","expiresAt":1e400}', + '{"owner":"default","expiresAt":-1e400}', + ...[false, null, 0, {}].map((commitFailure) => + JSON.stringify({ owner: 'default', expiresAt: Date.now() + 60_000, commitFailure }), + ), +]) { + test(`malformed repair evidence (${contents}) retains private state after the actual child has exited`, async () => { + const ownedStateDir = createOwnedReplayStateDir(); + const paths = ownedStateDir.paths; + const sessionDir = `${paths.sessionsDir}/default`; + fs.mkdirSync(sessionDir, { recursive: true }); + const evidencePath = `${sessionDir}/repair-tombstone.json`; + fs.writeFileSync(evidencePath, contents); + const launch = launchDaemonProcess({ + paths, + args: registeredDaemonFixtureArgs(paths, fields), + serverMode: 'socket', + ownedStateDir, + }); + try { + await waitForFixtureFile(paths.infoPath); + const result = await stopAndRetireDaemon({ + paths, + observed: { pid: launch.pid, startTime: launch.startTime ?? null }, + mode: 'graceful', + ownedStateDir, + }); + assert.equal(result.status, 'retained', JSON.stringify(result)); + if (result.status !== 'retained') assert.fail('repair evidence unexpectedly discarded'); + assert.equal(result.termination?.status, 'exited'); + assert.equal(result.error?.details?.reason, 'repair_evidence_invalid'); + assert.equal(fs.readFileSync(evidencePath, 'utf8'), contents); + await launch.exited; + } finally { + await finishPrivateTestDaemons(paths, launch); + } + }); +} + +async function waitForFixtureFile(filePath: string): Promise { + const deadline = Date.now() + 2_000; + while (!fs.existsSync(filePath) && Date.now() < deadline) await sleep(20); + assert.equal(fs.existsSync(filePath), true); +} + +async function finishPrivateTestDaemons( + paths: DaemonPaths, + ...launches: (ReturnType | undefined)[] +): Promise { + for (const launch of launches) { + if (!launch) continue; + const termination = await stopDaemonProcess( + { pid: launch.pid, startTime: launch.startTime ?? null }, + { mode: 'force', termTimeoutMs: 0, killTimeoutMs: 2_000 }, + ); + assert.notEqual(termination.status, 'retained', JSON.stringify(termination)); + await launch.exited; + } + fs.rmSync(paths.baseDir, { recursive: true, force: true }); +} diff --git a/src/__tests__/test-utils/registered-daemon-fixture.ts b/src/__tests__/test-utils/registered-daemon-fixture.ts new file mode 100644 index 0000000000..63ee16a7bc --- /dev/null +++ b/src/__tests__/test-utils/registered-daemon-fixture.ts @@ -0,0 +1,84 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import path from 'node:path'; +import { vi } from 'vitest'; +import type { runCmdDetachedMonitored } from '@agent-device/host-kit/command'; +import { readProcessStartTime } from '@agent-device/host-kit/process'; +import { stopDaemonProcess } from '../../daemon-process.ts'; +import type { DaemonPaths } from '../../daemon-resolution.ts'; +import type { DaemonRegistrationFields } from '../../daemon-registration-owner.ts'; + +const actualCommand = await vi.importActual( + '@agent-device/host-kit/command', +); +const children = new Map< + string, + { launch: ReturnType; startTime: string | null } +>(); + +/** A real registration owner, advertising the caller's HTTP fixture and joining before deletion. */ +export function registeredDaemonFixtureArgs( + paths: DaemonPaths, + fields: DaemonRegistrationFields, +): string[] { + const entry = path.join(paths.baseDir, 'dist', 'src', 'internal', 'daemon.js'); + fs.mkdirSync(path.dirname(entry), { recursive: true }); + fs.writeFileSync(path.join(paths.baseDir, 'package.json'), '{"type":"module"}'); + const registrationUrl = new URL('../../daemon-registration-owner.ts', import.meta.url).href; + fs.writeFileSync( + entry, + `import fs from 'node:fs'; +import path from 'node:path'; +import { tryAcquireDaemonRegistration } from ${JSON.stringify(registrationUrl)}; +const paths = ${JSON.stringify(paths)}; +const acquired = await tryAcquireDaemonRegistration(paths); +if (acquired.status !== 'acquired') process.exit(75); +process.on('SIGTERM', async () => { + const deferred = path.join(paths.baseDir, 'repair-on-shutdown.json'); + if (fs.existsSync(deferred)) { + const dir = path.join(paths.sessionsDir, 'default'); + fs.mkdirSync(dir, { recursive: true }); + fs.copyFileSync(deferred, path.join(dir, 'repair-tombstone.json')); + } + await acquired.owner.finish(); + process.exit(0); +}); +acquired.owner.publish(${JSON.stringify(fields)}); +setInterval(() => {}, 1000); +`, + ); + return ['--experimental-strip-types', entry]; +} + +export function spawnRegisteredDaemonFixture( + paths: DaemonPaths, + fields: DaemonRegistrationFields, + options: Parameters[2], +): ReturnType { + const child = actualCommand.runCmdDetachedMonitored( + process.execPath, + registeredDaemonFixtureArgs(paths, fields), + options, + ); + children.set(paths.baseDir, { launch: child, startTime: readProcessStartTime(child.pid) }); + return child; +} + +export async function finishRegisteredDaemonFixture(stateDir: string): Promise { + const owned = children.get(stateDir); + if (owned) { + const child = owned.launch; + const termination = await stopDaemonProcess( + { pid: child.pid, startTime: owned.startTime }, + { mode: 'force', termTimeoutMs: 0, killTimeoutMs: 2_000 }, + ); + assert.notEqual(termination.status, 'retained', JSON.stringify(termination)); + await child.exited; + children.delete(stateDir); + } + fs.rmSync(stateDir, { recursive: true, force: true }); +} + +export async function finishRegisteredDaemonFixtures(): Promise { + for (const stateDir of children.keys()) await finishRegisteredDaemonFixture(stateDir); +} diff --git a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts index 91f4da210a..d4d5ae9880 100644 --- a/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts +++ b/src/daemon-client/__tests__/daemon-client-lifecycle.test.ts @@ -6,6 +6,11 @@ import net from 'node:net'; import path from 'node:path'; import { afterEach, test, vi } from 'vitest'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; +import { + spawnRegisteredDaemonFixture, + finishRegisteredDaemonFixture, + finishRegisteredDaemonFixtures, +} from '../../__tests__/test-utils/registered-daemon-fixture.ts'; vi.mock('@agent-device/host-kit/command', async (importOriginal) => ({ ...(await importOriginal()), @@ -51,14 +56,19 @@ type DaemonInfoFixture = { processStartTime?: string; }; +const actualRetry = await vi.importActual( + '@agent-device/host-kit/retry', +); const mockRunCmdDetached = vi.mocked(runCmdDetachedMonitored); const mockRunCmdSync = vi.mocked(runCmdSync); const mockSleep = vi.mocked(sleep); -afterEach(() => { +afterEach(async () => { + await finishRegisteredDaemonFixtures(); mockRunCmdDetached.mockReset(); mockRunCmdSync.mockClear(); - mockSleep.mockClear(); + mockSleep.mockReset(); + mockSleep.mockImplementation(async () => {}); vi.unstubAllEnvs(); }); @@ -147,16 +157,22 @@ function installSpawnedHttpDaemonAtOwnedStateDir( httpPort: number, onStateDir: (stateDir: string) => void, ): void { + mockSleep.mockImplementation(actualRetry.sleep); mockRunCmdDetached.mockImplementation((_command, _args, options) => { const ownedStateDir = String(options?.env?.AGENT_DEVICE_STATE_DIR); onStateDir(ownedStateDir); const ownedPaths = resolveDaemonPaths(ownedStateDir); - writeDaemonInfo(ownedPaths, { httpPort, transport: 'http' }); - writeDaemonLock(ownedPaths, { - pid: process.pid, - processStartTime: readProcessStartTime(process.pid) ?? undefined, - }); - return { pid: process.pid, exited: new Promise(() => {}) }; + return spawnRegisteredDaemonFixture( + ownedPaths, + { + httpPort, + token: 'local-secret', + version: readVersion(), + codeOrigin: 'checkout', + codeSignature: currentDaemonCodeSignature(), + }, + options, + ); }); } @@ -298,7 +314,7 @@ function mockSocketErrorAfterWrite(failingPort: number): { }; } -test('sendToDaemon retries daemon spawn failures and cleans partial metadata on terminal failure', async () => { +test('sendToDaemon retains unknown metadata after a spawn failure', async () => { const stateDir = makeTempStateDir('agent-device-daemon-spawn-retry-'); const paths = resolveDaemonPaths(stateDir); vi.stubEnv('AGENT_DEVICE_STATE_DIR', stateDir); @@ -329,25 +345,17 @@ test('sendToDaemon retries daemon spawn failures and cleans partial metadata on assert.ok(thrown instanceof AppError); assert.equal(thrown.message, 'Failed to start daemon'); - assert.equal(thrown.details?.startError, 'spawn failed 2'); - assert.equal(thrown.details?.startupAttempts, 2); - const cleanupResults = thrown.details?.cleanupResults; - assert.ok(Array.isArray(cleanupResults)); - assert.deepEqual( - cleanupResults.map((result) => ({ - reason: result.reason, - removedInfo: result.removedInfo, - removedLock: result.removedLock, - })), - [ - { reason: 'start_error', removedInfo: true, removedLock: true }, - { reason: 'start_error', removedInfo: true, removedLock: true }, - ], - ); - assert.equal(attempts, 2); - assert.equal(mockSleep.mock.calls[0]?.[0], 150); - assert.equal(fs.existsSync(paths.infoPath), false); - assert.equal(fs.existsSync(paths.lockPath), false); + assert.equal(fs.readFileSync(paths.infoPath, 'utf8'), '{"partial":true}\n'); + assert.equal(fs.readFileSync(paths.lockPath, 'utf8'), 'not-json\n'); + assert.equal(thrown.details?.startError, 'spawn failed 1'); + assert.equal(thrown.details?.startupAttempts, 1); + const results = thrown.details?.cleanupResults as Array<{ + status: string; + removedInfo: boolean; + }>; + assert.equal(results[0]?.status, 'retained'); + assert.equal(results[0]?.removedInfo, false); + assert.equal(attempts, 1); } finally { fs.rmSync(stateDir, { recursive: true, force: true }); } @@ -516,7 +524,6 @@ test('sendToDaemon prints a takeover notice before replacing an unreachable daem const stateDir = makeTempStateDir('agent-device-daemon-unreachable-takeover-'); const paths = resolveDaemonPaths(stateDir); - // Bind fresh BEFORE freeing the port below: a later bind can reclaim it and skip the takeover. const freshDaemon = await startHttpDaemonFixture({ via: 'fresh-daemon' }); const unreachable = await startHttpDaemonFixture({ via: 'unused' }); await closeLoopbackServer(unreachable.server); @@ -738,12 +745,6 @@ test('sendToDaemon does not replay over HTTP after the socket request is written } }); -// --- ADR 0012 decision 6, R7 (Fix 1, C1): a repair-armed `replay --save-script` -// that comes back as a HELD divergence (the daemon's `resume.repairSessionHeld` -// signal) must keep its owning (owned/ephemeral) daemon alive and addressable. -// The keep-alive keys on that signal — the REPAIR-ARMED condition — NOT on -// `resume.allowed`, which reports only plan-resumability. --- - function heldDivergenceError( resume: Record = { allowed: true, from: 3, planDigest: 'digest-abc' }, ): Record { @@ -805,15 +806,13 @@ test('sendToDaemon keeps an owned ephemeral daemon alive and hints its --state-d assert.match(String(response.error.hint), /--state-dir/); assert.ok(String(response.error.hint).includes(ownedStateDir)); - // The daemon was NOT torn down: metadata and the owned state dir itself - // are still on disk, addressable by a follow-up command's --state-dir. const ownedPaths = resolveDaemonPaths(ownedStateDir); assert.equal(fs.existsSync(ownedPaths.infoPath), true); assert.equal(fs.existsSync(ownedPaths.lockPath), true); assert.equal(fs.existsSync(ownedStateDir), true); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); @@ -823,8 +822,6 @@ test('C1: keep-alive keys on repairSessionHeld, NOT resume.allowed — a HELD di return; } - // resume.allowed:false (plan not resumable), but the daemon still HELD the - // repair session — the agent must be able to reach it to close/inspect. const daemon = await startHttpDaemonErrorFixture( heldDivergenceError({ allowed: false, @@ -854,7 +851,7 @@ test('C1: keep-alive keys on repairSessionHeld, NOT resume.allowed — a HELD di assert.equal(fs.existsSync(ownedStateDir), true); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); @@ -883,92 +880,76 @@ test('sendToDaemon tears down an owned ephemeral daemon on an UNHELD divergence if (response.ok) return; assert.equal(response.error.hint, undefined); assert.ok(ownedStateDir.length > 0); - // No held signal (`resume.allowed:true` alone is not the keep-alive key) — - // ordinary one-shot teardown still applies. assert.equal(fs.existsSync(ownedStateDir), false); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); -// --- ADR 0012 decision 6 (BLOCKER 2, third follow-up): a one-shot -// `replay --save-script` that completes with no divergence returns SUCCESS -// immediately — the actual healed-script commit is deferred to daemon -// teardown. If that deferred commit then fails, the daemon leaves a -// REPAIR_COMMIT_FAILED tombstone in the owned state dir before exiting. The -// client cleanup must discover it (after waiting for the daemon to actually -// exit) BEFORE deleting the owned state dir, and must surface it in the -// response the caller receives — never silently delete the only evidence of -// the failure while reporting the success already computed for the replay -// itself. --- - test('BLOCKER 2 (third follow-up): a shutdown-time repair commit failure is surfaced and the owned state dir survives', async (t) => { if (!(await supportsLoopbackBind())) { t.skip('loopback listeners are not permitted in this environment'); return; } - // The daemon's RPC response for the replay itself is a plain SUCCESS (the - // plan completed with no divergence) — exactly what a real daemon would - // return before its deferred, teardown-time commit has even attempted. - const daemon = await startHttpDaemonFixture({ session: 'default' }); - let ownedStateDir = ''; - installSpawnedHttpDaemonAtOwnedStateDir(daemon.port, (dir) => { - ownedStateDir = dir; - // Simulate the daemon's OWN shutdown handler (`finalizeRepairTeardown`) - // having already run and left a commit-failure tombstone before this - // fake process "exits" — the real ordering `stopDaemonProcessForTakeover` - // depends on (it waits for the process to exit, and the real daemon only - // exits after teardown finishes writing this file). - const ownedPaths = resolveDaemonPaths(dir); - const sessionDir = path.join(ownedPaths.sessionsDir, 'default'); - fs.mkdirSync(sessionDir, { recursive: true }); - fs.writeFileSync( - path.join(sessionDir, 'repair-tombstone.json'), - `${JSON.stringify({ - owner: 'default', - reapedAt: Date.now(), - expiresAt: Date.now() + 60_000, - sourcePath: '/tmp/flow.ad', - commitFailure: { - code: 'COMMAND_FAILED', - message: 'a prior healed script already exists at /tmp/flow.healed.ad', - }, - })}\n`, - ); - }); - - try { - const response = await sendToDaemon({ - session: 'default', - command: 'replay', - positionals: ['flow.ad'], - flags: { saveScript: true, daemonTransport: 'http' }, - meta: { requestId: 'req-repair-commit-fail-teardown' }, + for (const failRelease of [false, true]) { + const daemon = await startHttpDaemonFixture({ session: 'default' }); + let ownedStateDir = ''; + installSpawnedHttpDaemonAtOwnedStateDir(daemon.port, (dir) => { + ownedStateDir = dir; + fs.writeFileSync( + path.join(dir, 'repair-on-shutdown.json'), + `${JSON.stringify({ + owner: 'default', + reapedAt: Date.now(), + expiresAt: Date.now() + 60_000, + sourcePath: '/tmp/flow.ad', + commitFailure: { + code: 'COMMAND_FAILED', + message: 'a prior healed script already exists at /tmp/flow.healed.ad', + }, + })}\n`, + ); }); - // The client-visible response must surface the deferred commit failure — - // never the raw success the daemon returned for the replay itself, and - // never silently swallowed by cleanup. - assert.equal(response.ok, false); - if (response.ok) return; - assert.equal(response.error.code, 'REPAIR_COMMIT_FAILED'); - assert.match(response.error.message, /a prior healed script already exists/); - assert.ok(response.error.message.includes('replay /tmp/flow.ad --save-script')); + const originalRmdir = fs.rmdirSync; + const releaseSpy = vi.spyOn(fs, 'rmdirSync').mockImplementation((target, options) => { + if (failRelease && target === resolveDaemonPaths(ownedStateDir).lockPath) + throw Object.assign(new Error('release failed'), { code: 'EBUSY' }); + return originalRmdir(target, options); + }); + try { + const response = await sendToDaemon({ + session: 'default', + command: 'replay', + positionals: ['flow.ad'], + flags: { saveScript: true, daemonTransport: 'http' }, + meta: { requestId: 'req-repair-commit-fail-teardown' }, + }); - // The owned state dir — and the tombstone evidence inside it — must - // survive: never rmSync'd while an unrecovered commit failure is on record. - assert.ok(ownedStateDir.length > 0); - assert.equal(fs.existsSync(ownedStateDir), true); - const ownedPaths = resolveDaemonPaths(ownedStateDir); - assert.equal( - fs.existsSync(path.join(ownedPaths.sessionsDir, 'default', 'repair-tombstone.json')), - true, - ); - } finally { - await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + assert.ok(!response.ok); + assert.equal(response.error.code, 'REPAIR_COMMIT_FAILED'); + const secondary = response.error.details?.cleanupFailure as + | { details?: { ownerReleaseUnverified?: boolean }; hint?: string } + | undefined; + assert.equal(Boolean(secondary?.details?.ownerReleaseUnverified), failRelease); + assert.equal(Boolean(secondary?.hint?.startsWith('Restore process inspection')), failRelease); + assert.match(response.error.message, /a prior healed script already exists/); + assert.ok(response.error.message.includes('replay /tmp/flow.ad --save-script')); + + assert.ok(ownedStateDir.length > 0); + assert.equal(fs.existsSync(ownedStateDir), true); + const ownedPaths = resolveDaemonPaths(ownedStateDir); + assert.equal( + fs.existsSync(path.join(ownedPaths.sessionsDir, 'default', 'repair-tombstone.json')), + true, + ); + } finally { + releaseSpy.mockRestore(); + await closeLoopbackServer(daemon.server); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); + } } }); @@ -1003,7 +984,7 @@ test('continuation: sendToDaemon keeps the daemon alive on a held divergence eve assert.equal(fs.existsSync(ownedStateDir), true); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); @@ -1131,7 +1112,7 @@ test('sendToDaemon keeps an owned ephemeral daemon alive and hints its --state-d assert.equal(fs.existsSync(ownedStateDir), true); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); @@ -1170,7 +1151,7 @@ test('closes the loop: a follow-up sendToDaemon using the hinted --state-dir/--s assert.equal(daemon.rpcRequests[1]?.params?.command, 'press'); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); @@ -1202,7 +1183,7 @@ test('sendToDaemon tears down an owned ephemeral daemon when replay reports the assert.equal(fs.existsSync(ownedStateDir), false); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); @@ -1248,7 +1229,7 @@ test('ADR 0012 R7 x ADR 0016: a completed --save-script repair also keeps its ow assert.equal(fs.existsSync(ownedStateDir), true); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); @@ -1328,6 +1309,6 @@ test('sendToDaemon still tears down a `test` command owned ephemeral daemon even assert.equal(fs.existsSync(ownedStateDir), false); } finally { await closeLoopbackServer(daemon.server); - if (ownedStateDir) fs.rmSync(ownedStateDir, { recursive: true, force: true }); + if (ownedStateDir) await finishRegisteredDaemonFixture(ownedStateDir); } }); diff --git a/src/daemon-client/__tests__/daemon-client-metadata.test.ts b/src/daemon-client/__tests__/daemon-client-metadata.test.ts index 5b46a91e8d..6bfec46b6e 100644 --- a/src/daemon-client/__tests__/daemon-client-metadata.test.ts +++ b/src/daemon-client/__tests__/daemon-client-metadata.test.ts @@ -5,10 +5,12 @@ import path from 'node:path'; import { afterEach, test, vi } from 'vitest'; import type { DaemonCodeOrigin } from '@agent-device/host-kit/code-signature'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; -import { tryAcquireDaemonRegistration } from '../../daemon-registration-owner.ts'; +import { + stopAndRetireDaemon, + tryAcquireDaemonRegistration, +} from '../../daemon-registration-owner.ts'; import { readDaemonInfo, - cleanupFailedDaemonStartupMetadata, stopDaemonProcessForTakeover, type DaemonInfo, } from '../daemon-client-metadata.ts'; @@ -81,13 +83,15 @@ for (const artifact of ['daemon.json', 'daemon.lock']) { fs.writeFileSync(file, contents); vi.mocked(isAgentDeviceDaemonProcess).mockReturnValue(true); vi.mocked(stopDaemonProcess).mockResolvedValue({ status: 'retained', reason: 'exit-timeout' }); - const result = await cleanupFailedDaemonStartupMetadata(paths, 'start_error'); + const result = await stopAndRetireDaemon({ + paths, + observed: { pid: 7, startTime: 'start' }, + mode: 'graceful', + }); assert.equal(fs.readFileSync(file, 'utf8'), contents); assert.equal(result.removedInfo, false); - assert.equal(result.removedLock, false); - assert.equal(result.stoppedInfoProcess, false); - assert.equal(result.stoppedLockProcess, false); - assert.match(result.error ?? '', /exit could not be confirmed/); + assert.equal(result.status, 'retained'); + if (result.status === 'retained') assert.equal(result.reason, 'exit-unconfirmed'); }); } diff --git a/src/daemon-client/__tests__/daemon-client-startup-race.test.ts b/src/daemon-client/__tests__/daemon-client-startup-race.test.ts index d52ba0e6c2..fa87f3f65b 100644 --- a/src/daemon-client/__tests__/daemon-client-startup-race.test.ts +++ b/src/daemon-client/__tests__/daemon-client-startup-race.test.ts @@ -200,7 +200,7 @@ test('a start race won by an older daemon replaces it instead of adopting it', a }), ); assert.equal(fixture.rpcRequests.length, 0); - assert.equal(launches, 2); + assert.equal(launches, 1); assert.match( String(stderr.mock.calls.flat().join('')), /Replacing daemon \(pid 43300, v0\.0\.1\)/, diff --git a/src/daemon-client/__tests__/daemon-client.test.ts b/src/daemon-client/__tests__/daemon-client.test.ts index f579a3c4a5..accebf0782 100644 --- a/src/daemon-client/__tests__/daemon-client.test.ts +++ b/src/daemon-client/__tests__/daemon-client.test.ts @@ -1,5 +1,5 @@ import type { RequestProgressEvent } from '@agent-device/contracts/progress'; -import { test, vi } from 'vitest'; +import { test } from 'vitest'; import assert from 'node:assert/strict'; import http from 'node:http'; import net from 'node:net'; @@ -11,57 +11,18 @@ import { listenOnLoopback, supportsLoopbackBind, } from '../../__tests__/test-utils/loopback.ts'; -import { runCmdBackground } from '@agent-device/host-kit/command'; -import { - isProcessAlive, - readProcessCommand, - readProcessStartTime, - waitForProcessExit, -} from '@agent-device/host-kit/process'; +import { readProcessStartTime } from '@agent-device/host-kit/process'; import { sendToDaemon } from '../daemon-client.ts'; import { currentDaemonCodeSignature } from '../../__tests__/test-utils/daemon-http-fixture.ts'; import { computeDaemonCodeSignature } from '@agent-device/host-kit/code-signature'; import { downloadRemoteArtifact } from '../../remote/daemon-artifacts.ts'; -import { - cleanupFailedDaemonStartupMetadata, - resolveDaemonStartupHint, -} from '../daemon-client-metadata.ts'; +import { resolveDaemonStartupHint } from '../daemon-client-metadata.ts'; import { canConnectSocket } from '../daemon-client-transport.ts'; import { DAEMON_RPC_PROTOCOL_VERSION } from '@agent-device/contracts/daemon-http'; import { resolveDaemonPaths } from '../../daemon-resolution.ts'; import { readVersion } from '@agent-device/host-kit/version'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; -// readProcessStartTime/readProcessCommand shell out to `ps` with a 1s -// timeout (see host-process.ts). isAgentDeviceDaemonProcess re-reads both for -// every liveness check, so a spawned-daemon fixture that is proven live once -// (a real read, right after the process starts) can still be misclassified -// as dead later if a *subsequent* `ps` call happens to miss its deadline -// under full-suite CPU contention. mockReadProcessStartTime/mockReadProcessCommand -// default to `undefined`, which falls through to the real implementation for -// every pid in every test in this file; only the one test below that needs a -// stable answer for its spawned pid configures an override, and clears it -// afterward. -const { mockReadProcessStartTime, mockReadProcessCommand } = vi.hoisted(() => ({ - mockReadProcessStartTime: vi.fn<(pid: number) => string | null | undefined>(), - mockReadProcessCommand: vi.fn<(pid: number) => string | null | undefined>(), -})); - -vi.mock('@agent-device/host-kit/process', async (importOriginal) => { - const actual = await importOriginal(); - return { - ...actual, - readProcessStartTime: (pid: number) => { - const overridden = mockReadProcessStartTime(pid); - return overridden !== undefined ? overridden : actual.readProcessStartTime(pid); - }, - readProcessCommand: (pid: number) => { - const overridden = mockReadProcessCommand(pid); - return overridden !== undefined ? overridden : actual.readProcessCommand(pid); - }, - }; -}); - type MockHttpResponse = EventEmitter & { headers?: Record; statusCode?: number; @@ -219,144 +180,6 @@ test('resolveDaemonStartupHint shell-quotes cleanup paths', () => { ); }); -test('cleanupFailedDaemonStartupMetadata removes partial startup metadata', async () => { - const stateDir = mkdtempForTestSync('agent-device-daemon-cleanup-'); - const paths = resolveDaemonPaths(stateDir); - try { - fs.mkdirSync(paths.baseDir, { recursive: true }); - fs.writeFileSync(paths.infoPath, '{"invalid":true}\n', 'utf8'); - fs.writeFileSync(paths.lockPath, 'not-json\n', 'utf8'); - - const result = await cleanupFailedDaemonStartupMetadata(paths, 'startup_timeout'); - - assert.deepEqual(result, { - reason: 'startup_timeout', - removedInfo: true, - removedLock: true, - stoppedInfoProcess: false, - stoppedLockProcess: false, - }); - assert.equal(fs.existsSync(paths.infoPath), false); - assert.equal(fs.existsSync(paths.lockPath), false); - } finally { - fs.rmSync(stateDir, { recursive: true, force: true }); - } -}); - -test('cleanupFailedDaemonStartupMetadata retains live startup daemon on timeout', async (t) => { - const stateDir = mkdtempForTestSync('agent-device-daemon-live-cleanup-'); - const root = mkdtempForTestSync('agent-device-live-daemon-'); - const daemonDir = path.join(root, 'agent-device', 'dist', 'src', 'internal'); - const daemonScriptPath = path.join(daemonDir, 'daemon.js'); - fs.mkdirSync(daemonDir, { recursive: true }); - fs.writeFileSync(daemonScriptPath, 'setInterval(() => {}, 1000);\n', 'utf8'); - const daemonProcess = runCmdBackground(process.execPath, [daemonScriptPath], { - stdio: 'ignore', - allowFailure: true, - captureOutput: false, - }); - void daemonProcess.wait.catch(() => {}); - const pid = daemonProcess.child.pid; - assert.ok(pid, 'spawned child should have a pid'); - - try { - await new Promise((resolve) => setTimeout(resolve, 50)); - // Read the spawned daemon's real identity once (ground truth: it is - // genuinely alive, with this real start time and command line), then - // pin readProcessStartTime/readProcessCommand to keep returning these - // same proven-real values for this pid. isAgentDeviceDaemonProcess reads - // both again internally on every call inside cleanupFailedDaemonStartupMetadata; - // without pinning, a second real `ps` call could miss its 1s timeout - // under load and misclassify this genuinely-live daemon as dead. - const processStartTime = readProcessStartTime(pid) ?? undefined; - const command = readProcessCommand(pid); - if (command === null || processStartTime === undefined) { - t.skip('process command/start inspection is unavailable in this environment'); - return; - } - mockReadProcessStartTime.mockImplementation((queriedPid: number) => - queriedPid === pid ? processStartTime : undefined, - ); - mockReadProcessCommand.mockImplementation((queriedPid: number) => - queriedPid === pid ? command : undefined, - ); - - const paths = resolveDaemonPaths(stateDir); - fs.mkdirSync(paths.baseDir, { recursive: true }); - fs.writeFileSync( - paths.infoPath, - `${JSON.stringify({ - token: 'startup-secret', - port: 65530, - transport: 'socket', - pid, - processStartTime, - })}\n`, - 'utf8', - ); - fs.writeFileSync( - paths.lockPath, - `${JSON.stringify({ pid, processStartTime, startedAt: Date.now() })}\n`, - 'utf8', - ); - - const result = await cleanupFailedDaemonStartupMetadata(paths, 'startup_timeout', { - stopLiveProcesses: false, - }); - - assert.equal(result.retainedInfoProcess, true); - assert.equal(result.retainedLockProcess, true); - assert.equal(result.removedInfo, false); - assert.equal(result.removedLock, false); - assert.equal(isProcessAlive(pid), true); - assert.equal(fs.existsSync(paths.infoPath), true); - assert.equal(fs.existsSync(paths.lockPath), true); - } finally { - mockReadProcessStartTime.mockReset(); - mockReadProcessCommand.mockReset(); - if (isProcessAlive(pid)) { - process.kill(pid, 'SIGKILL'); - await waitForProcessExit(pid, 1_500); - } - fs.rmSync(stateDir, { recursive: true, force: true }); - fs.rmSync(root, { recursive: true, force: true }); - } -}); - -test('cleanupFailedDaemonStartupMetadata removes stale daemon metadata on timeout', async () => { - const stateDir = mkdtempForTestSync('agent-device-daemon-stale-cleanup-'); - const paths = resolveDaemonPaths(stateDir); - try { - fs.mkdirSync(paths.baseDir, { recursive: true }); - fs.writeFileSync( - paths.infoPath, - `${JSON.stringify({ - token: 'startup-secret', - port: 65530, - transport: 'socket', - pid: 999_999, - })}\n`, - 'utf8', - ); - fs.writeFileSync( - paths.lockPath, - `${JSON.stringify({ pid: 999_999, startedAt: Date.now() })}\n`, - 'utf8', - ); - - const result = await cleanupFailedDaemonStartupMetadata(paths, 'startup_timeout', { - stopLiveProcesses: false, - }); - - assert.equal(result.removedInfo, true); - assert.equal(result.removedLock, true); - assert.equal(fs.existsSync(paths.infoPath), false); - assert.equal(fs.existsSync(paths.lockPath), false); - } finally { - fs.rmSync(stateDir, { recursive: true, force: true }); - } -}); - test('canConnectSocket times out stalled local daemon probes', async () => { const originalCreateConnection = net.createConnection; let timeoutMs: number | undefined; diff --git a/src/daemon-client/daemon-client-lifecycle.ts b/src/daemon-client/daemon-client-lifecycle.ts index 43a6556781..7e75749e61 100644 --- a/src/daemon-client/daemon-client-lifecycle.ts +++ b/src/daemon-client/daemon-client-lifecycle.ts @@ -1,17 +1,25 @@ import fs from 'node:fs'; import net from 'node:net'; -import os from 'node:os'; -import path from 'node:path'; -import { AppError, normalizeError } from '@agent-device/kernel/errors'; +import { AppError, normalizeError, type NormalizedError } from '@agent-device/kernel/errors'; import { readReplayDivergenceResume } from '@agent-device/ad-replay/divergence'; import type { DaemonRequest, DaemonResponse } from '../daemon/daemon-request.ts'; -import { runCmdDetachedMonitored, type ExecDetachedExit } from '@agent-device/host-kit/command'; +import { type ExecDetachedExit } from '@agent-device/host-kit/command'; import { shellQuoteIfNeeded } from '@agent-device/kernel/device-shell'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; -import { isProcessAlive, readProcessStartTime } from '@agent-device/host-kit/process'; +import { isProcessAlive } from '@agent-device/host-kit/process'; import { sleep } from '@agent-device/host-kit/retry'; +import { inspectProcessLock } from '@agent-device/host-kit/file'; -import { findUnrecoveredRepairCommitFailure } from '../session-repair-tombstone.ts'; +import type { findUnrecoveredRepairCommitFailure } from '../session-repair-tombstone.ts'; +import { + createOwnedReplayStateDir, + recoverAbandonedDaemonRegistration, + type DaemonRetirementResult, + launchDaemonProcess, + stopAndRetireDaemon, + type OwnedReplayStateDir, + type DaemonStartupLaunch, +} from '../daemon-registration-owner.ts'; import { resolveDaemonPaths, resolveDaemonServerMode, @@ -28,19 +36,15 @@ import { import { PUBLIC_COMMANDS } from '@agent-device/command-registry/catalog'; import { - cleanupFailedDaemonStartupMetadata, cleanupStaleDaemonLockIfSafe, getDaemonMetadataState, isDaemonLockHeldByAnotherDaemon, isRemoteDaemon, readDaemonInfo, - recoverDaemonLockHolder, removeDaemonInfo, - removeDaemonLock, resolveDaemonStartupHint, stopDaemonProcessForTakeover, type DaemonInfo, - type DaemonStartupCleanupResult, } from './daemon-client-metadata.ts'; import { canConnect, @@ -52,7 +56,7 @@ export type DaemonClientSettings = { paths: DaemonPaths; transportPreference: DaemonTransportPreference; serverMode: DaemonServerMode; - ownedStateDir?: boolean; + ownedStateDir?: OwnedReplayStateDir; remoteBaseUrl?: string; remoteAuthToken?: string; }; @@ -62,13 +66,6 @@ export type EnsuredDaemon = { startedByClient: boolean; }; -type DaemonStartupLaunch = { - pid: number; - /** The launched process's start time, so a reused pid is never taken for it. */ - startTime?: string; - exited: Promise; -}; - type DaemonStartupWaitResult = | { kind: 'ready'; daemon: EnsuredDaemon } | { kind: 'early_exit'; exit: ExecDetachedExit } @@ -91,10 +88,11 @@ export function resolveClientSettings( const explicitStateDir = resolveExplicitStateDir(req); const remote = resolveRemoteClientSettings(req, suppliedAuthToken); const transport = resolveTransportClientSettings(req, remote.remoteBaseUrl); - const ownedStateDir = shouldUseOwnedReplayStateDir(req, explicitStateDir, remote.rawBaseUrl); - const stateDir = ownedStateDir ? createOwnedReplayStateDir() : explicitStateDir; + const ownedStateDir = shouldUseOwnedReplayStateDir(req, explicitStateDir, remote.rawBaseUrl) + ? createOwnedReplayStateDir() + : undefined; return { - paths: resolveDaemonPaths(stateDir), + paths: ownedStateDir?.paths ?? resolveDaemonPaths(explicitStateDir), transportPreference: transport.preference, serverMode: transport.serverMode, ownedStateDir, @@ -153,10 +151,6 @@ function shouldUseOwnedReplayStateDir( return isOneShotReplayCommand(req.command) && !explicitStateDir && !rawRemoteBaseUrl; } -function createOwnedReplayStateDir(): string { - return fs.mkdtempSync(path.join(os.tmpdir(), 'agent-device-replay-daemon-')); -} - export async function ensureDaemon(settings: DaemonClientSettings): Promise { if (settings.remoteBaseUrl) { return await ensureRemoteDaemon(settings); @@ -303,65 +297,27 @@ function emitDaemonTakeoverNotice(info: DaemonInfo, reason: string, stateDir: st } } -async function startLocalDaemon(settings: DaemonClientSettings): Promise { - let lockRecoveryCount = 0; - const cleanupResults: DaemonStartupCleanupResult[] = []; - let startError: string | undefined; - let daemonProcess: ExecDetachedExit | { pid: number } | undefined; - for (let attempt = 1; attempt <= DAEMON_STARTUP_ATTEMPTS; attempt += 1) { - let launch: DaemonStartupLaunch; - try { - launch = startDaemon(settings); - daemonProcess = { pid: launch.pid }; - } catch (error) { - startError = error instanceof Error ? error.message : String(error); - cleanupResults.push(await cleanupFailedDaemonStartupMetadata(settings.paths, 'start_error')); - if (attempt < DAEMON_STARTUP_ATTEMPTS) { - await sleep(150); - continue; - } - break; - } - - const startup = await waitForDaemonStartup(DAEMON_STARTUP_TIMEOUT_MS, settings, launch); - if (startup.kind === 'ready') return startup.daemon; - if (startup.kind === 'early_exit') { - daemonProcess = startup.exit; - startError = describeDaemonEarlyExit(startup.exit); - cleanupResults.push(await cleanupFailedDaemonStartupMetadata(settings.paths, 'start_error')); - if (attempt < DAEMON_STARTUP_ATTEMPTS) { - await sleep(150); - continue; - } - break; - } - - if (await recoverDaemonLockHolder(settings.paths)) { - lockRecoveryCount += 1; - continue; - } - - const metadataState = getDaemonMetadataState(settings.paths); - const hasAnotherAttempt = attempt < DAEMON_STARTUP_ATTEMPTS; - const cleanup = await cleanupFailedDaemonStartupMetadata(settings.paths, 'startup_timeout', { - stopLiveProcesses: false, - }); - cleanupResults.push(cleanup); - if (cleanup.retainedInfoProcess || cleanup.retainedLockProcess) { - const extended = await waitForDaemonStartup(DAEMON_STARTUP_TIMEOUT_MS, settings, launch); - if (extended.kind === 'ready') return extended.daemon; - if (extended.kind === 'early_exit') { - daemonProcess = extended.exit; - startError = describeDaemonEarlyExit(extended.exit); - } - break; - } - if (!hasAnotherAttempt) break; +type FailedDaemonStartup = { + cleanup: DaemonRetirementResult; + startError?: string; + daemonProcess?: ExecDetachedExit | { pid: number }; + retry: boolean; +}; - // Detached daemon startup can race on busy CI hosts; retry when no metadata exists yet. - if (!metadataState.hasInfo && !metadataState.hasLock) await sleep(150); +async function startLocalDaemon(settings: DaemonClientSettings): Promise { + const deadline = Date.now() + DAEMON_STARTUP_TIMEOUT_MS; + const cleanupResults: DaemonRetirementResult[] = []; + let failure: FailedDaemonStartup | undefined; + let attempts = 0; + while (attempts < DAEMON_STARTUP_ATTEMPTS && Date.now() < deadline) { + attempts += 1; + const result = await attemptLocalDaemonStartup(settings, deadline); + if ('daemon' in result) return result.daemon; + failure = result; + cleanupResults.push(result.cleanup); + if (!result.retry) break; + await sleep(Math.min(150, Math.max(0, deadline - Date.now()))); } - const state = getDaemonMetadataState(settings.paths); const daemonLogTail = readRecentLogTail(settings.paths.logPath); throw new AppError('COMMAND_FAILED', 'Failed to start daemon', { @@ -371,17 +327,57 @@ async function startLocalDaemon(settings: DaemonClientSettings): Promise { + let launch: DaemonStartupLaunch; + try { + launch = startDaemon(settings); + } catch (error) { + const cleanup = await recoverAbandonedDaemonRegistration({ + paths: settings.paths, + observed: null, + lockTimeoutMs: 0, + }); + return { + cleanup, + startError: normalizeError(error).message, + retry: cleanup.status !== 'retained', + }; + } + const startup = await waitForDaemonStartup(Math.max(0, deadline - Date.now()), settings, launch); + if (startup.kind === 'ready') return { daemon: startup.daemon }; + const cleanup = await stopAndRetireDaemon({ + paths: settings.paths, + observed: { pid: launch.pid, startTime: launch.startTime ?? null }, + mode: 'graceful', + lockTimeoutMs: 0, + }); + const inspection = inspectProcessLock(settings.paths.lockPath); + const available = + inspection.state === 'absent' || + (inspection.state === 'held' && + (inspection.liveness === 'owner-process-dead' || + inspection.liveness === 'owner-process-reused')); + return { + cleanup, + retry: startup.kind === 'early_exit' && available, + startError: startup.kind === 'early_exit' ? describeDaemonEarlyExit(startup.exit) : undefined, + daemonProcess: startup.kind === 'early_exit' ? startup.exit : { pid: launch.pid }, + }; +} + /** * ADR 0012 decision 6 (BLOCKER 2, third follow-up): a one-shot repair * (`replay --save-script`) that COMPLETES without diverging returns SUCCESS @@ -435,49 +431,44 @@ export async function cleanupDaemonAfterRequest( return response; } - const result = { - pid: daemon.info.pid, - removedInfo: false, - removedLock: false, - removedStateDir: false, - error: undefined as string | undefined, - }; - let surfacedResponse = response; - - try { - await stopDaemonProcessForTakeover(daemon.info); - } catch (error) { - result.error = error instanceof Error ? error.message : String(error); - } finally { - const infoExists = fs.existsSync(settings.paths.infoPath); - removeDaemonInfo(settings.paths.infoPath); - result.removedInfo = infoExists && !fs.existsSync(settings.paths.infoPath); - - const lockExists = fs.existsSync(settings.paths.lockPath); - removeDaemonLock(settings.paths.lockPath); - result.removedLock = lockExists && !fs.existsSync(settings.paths.lockPath); - - if (settings.ownedStateDir) { - // `stopDaemonProcessForTakeover` above waits for the (real) daemon - // process to actually exit, which only happens AFTER its shutdown - // handler finishes `finalizeRepairTeardown` for every session — so by - // now any commit-failure tombstone it would leave is already on disk. - const unrecovered = findUnrecoveredRepairCommitFailure(settings.paths.sessionsDir); - if (unrecovered) { - surfacedResponse = surfaceUnrecoveredRepairCommitFailure(response, unrecovered); - } else { - fs.rmSync(settings.paths.baseDir, { recursive: true, force: true }); - result.removedStateDir = !fs.existsSync(settings.paths.baseDir); - } - } - } - + const result = await stopAndRetireDaemon({ + paths: settings.paths, + observed: { pid: daemon.info.pid, startTime: daemon.info.processStartTime ?? null }, + mode: 'graceful', + ownedStateDir: settings.ownedStateDir, + }); emitDiagnostic({ - level: result.error ? 'warn' : 'info', + level: result.status === 'retained' ? 'warn' : 'info', phase: 'daemon_replay_cleanup', - data: result, + data: { pid: daemon.info.pid, ...result }, }); - return surfacedResponse; + if (result.status !== 'absent' && result.repairCommitFailure) { + return surfaceUnrecoveredRepairCommitFailure( + response, + result.repairCommitFailure, + result.status === 'retained' ? result.error : undefined, + ); + } + if (result.status === 'retained' && response?.ok) { + return { + ok: false, + error: normalizeError( + new AppError( + 'COMMAND_FAILED', + 'Replay completed, but daemon cleanup could not be confirmed.', + { + reason: 'daemon_retirement_unconfirmed', + retirement: result, + stateDir: settings.paths.baseDir, + hint: + result.error?.hint ?? + `State and diagnostics were retained at ${settings.paths.baseDir}. Resolve the reported cleanup failure before retrying.`, + }, + ), + ), + }; + } + return response; } /** @@ -498,6 +489,7 @@ export async function cleanupDaemonAfterRequest( function surfaceUnrecoveredRepairCommitFailure( response: DaemonResponse | undefined, unrecovered: NonNullable>, + cleanupFailure?: NormalizedError, ): DaemonResponse { if (response && !response.ok) return response; const { sessionName, tombstone } = unrecovered; @@ -507,7 +499,16 @@ function surfaceUnrecoveredRepairCommitFailure( const message = `The repair transaction for session "${sessionName}" completed, but committing its ` + `healed script failed at teardown: ${tombstone.commitFailure.message}. ${reRun}.`; - return { ok: false, error: normalizeError(new AppError('REPAIR_COMMIT_FAILED', message)) }; + return { + ok: false, + error: normalizeError( + new AppError( + 'REPAIR_COMMIT_FAILED', + message, + cleanupFailure ? { cleanupFailure } : undefined, + ), + ), + }; } /** @@ -524,7 +525,7 @@ function surfaceUnrecoveredRepairCommitFailure( * `resume.allowed` (plan-resumability): a held divergence with `allowed: false` * still holds the session so the agent can inspect and `close` cleanly. */ -export function isHeldRepairDivergence(response: DaemonResponse | undefined): boolean { +function isHeldRepairDivergence(response: DaemonResponse | undefined): boolean { if (!response || response.ok) return false; if (response.error.code !== 'REPLAY_DIVERGENCE') return false; const resume = readReplayDivergenceResume(response.error.details?.divergence); @@ -539,7 +540,7 @@ export function isHeldRepairDivergence(response: DaemonResponse | undefined): bo * selector-miss's own guidance) so the agent's next command knows to target * the SAME daemon instead of resolving to the default one. */ -export function attachRepairSessionAddressHint( +function attachRepairSessionAddressHint( response: Extract, stateDir: string, ): Extract { @@ -571,7 +572,7 @@ function isOneShotReplayCommand(command: string | undefined): boolean { * anyway, but the explicit command check keeps that carve-out a decision * rather than an accident of the response shape. */ -export function isActiveReplaySessionResponse( +function isActiveReplaySessionResponse( req: Omit, response: DaemonResponse | undefined, ): boolean { @@ -669,28 +670,14 @@ function isLaunchedDaemon(info: DaemonInfo, launch: DaemonStartupLaunch): boolea function startDaemon(settings: DaemonClientSettings): DaemonStartupLaunch { const launchSpec = resolveDaemonLaunchSpec(); - const args = launchSpec.useSrc - ? ['--experimental-strip-types', launchSpec.srcPath] - : [launchSpec.distPath]; - const env = { - ...process.env, - AGENT_DEVICE_STATE_DIR: settings.paths.baseDir, - AGENT_DEVICE_DAEMON_SERVER_MODE: settings.serverMode, - }; - - fs.mkdirSync(settings.paths.baseDir, { recursive: true }); - const stdoutFd = fs.openSync(settings.paths.logPath, 'a'); - const stderrFd = fs.openSync(settings.paths.logPath, 'a'); - try { - const launched = runCmdDetachedMonitored(process.execPath, args, { - env, - stdio: ['ignore', stdoutFd, stderrFd], - }); - return { ...launched, startTime: readProcessStartTime(launched.pid) ?? undefined }; - } finally { - fs.closeSync(stdoutFd); - fs.closeSync(stderrFd); - } + return launchDaemonProcess({ + paths: settings.paths, + serverMode: settings.serverMode, + ownedStateDir: settings.ownedStateDir, + args: launchSpec.useSrc + ? ['--experimental-strip-types', launchSpec.srcPath] + : [launchSpec.distPath], + }); } function describeDaemonEarlyExit(exit: ExecDetachedExit): string { @@ -771,3 +758,21 @@ function isLoopbackHostname(hostname: string): boolean { if (net.isIPv6(normalized)) return LOOPBACK_BLOCK_LIST.check(normalized, 'ipv6'); return false; } + +export function attachSessionAddressHints( + response: DaemonResponse, + req: Omit, + settings: DaemonClientSettings, +): DaemonResponse { + if (!response.ok) { + return settings.ownedStateDir && isHeldRepairDivergence(response) + ? attachRepairSessionAddressHint(response, settings.paths.baseDir) + : response; + } + return isActiveReplaySessionResponse(req, response) + ? attachActiveSessionAddressHint( + response, + settings.ownedStateDir ? settings.paths.baseDir : undefined, + ) + : response; +} diff --git a/src/daemon-client/daemon-client-metadata.ts b/src/daemon-client/daemon-client-metadata.ts index 3a448b099f..bf886cef5c 100644 --- a/src/daemon-client/daemon-client-metadata.ts +++ b/src/daemon-client/daemon-client-metadata.ts @@ -1,7 +1,6 @@ import fs from 'node:fs'; import { AppError } from '@agent-device/kernel/errors'; import { shellQuote } from '@agent-device/kernel/device-shell'; -import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; import { isAgentDeviceDaemonProcess, stopDaemonProcess, @@ -44,19 +43,6 @@ export type DaemonMetadataState = { hasLock: boolean; }; -type DaemonStartupCleanupReason = 'start_error' | 'startup_timeout'; - -export type DaemonStartupCleanupResult = { - reason: DaemonStartupCleanupReason; - removedInfo: boolean; - removedLock: boolean; - stoppedInfoProcess: boolean; - stoppedLockProcess: boolean; - retainedInfoProcess?: boolean; - retainedLockProcess?: boolean; - error?: string; -}; - const DAEMON_TAKEOVER_TERM_TIMEOUT_MS = 3000; const DAEMON_TAKEOVER_KILL_TIMEOUT_MS = 1000; @@ -161,78 +147,6 @@ export function cleanupStaleDaemonLockIfSafe(paths: DaemonPaths): void { removeDaemonLock(paths.lockPath); } -export async function cleanupFailedDaemonStartupMetadata( - paths: DaemonPaths, - reason: DaemonStartupCleanupReason, - options: { stopLiveProcesses?: boolean } = {}, -): Promise { - const stopLiveProcesses = options.stopLiveProcesses ?? true; - const result: DaemonStartupCleanupResult = { - reason, - removedInfo: false, - removedLock: false, - stoppedInfoProcess: false, - stoppedLockProcess: false, - }; - - try { - const infoExists = fs.existsSync(paths.infoPath); - const info = readDaemonInfo(paths.infoPath); - if (info) { - const liveInfoProcess = isAgentDeviceDaemonProcess(info.pid, info.processStartTime); - if (liveInfoProcess && !stopLiveProcesses) { - result.retainedInfoProcess = true; - } else { - if (liveInfoProcess) { - await stopDaemonProcessForTakeover(info); - result.stoppedInfoProcess = true; - } - removeDaemonInfo(paths.infoPath); - result.removedInfo = true; - } - } else if (infoExists) { - removeDaemonInfo(paths.infoPath); - result.removedInfo = true; - } - - const lockExists = fs.existsSync(paths.lockPath); - const lockInfo = readDaemonLockInfo(paths.lockPath); - if (lockInfo) { - const liveLockProcess = isAgentDeviceDaemonProcess(lockInfo.pid, lockInfo.processStartTime); - if (liveLockProcess && !stopLiveProcesses) { - result.retainedLockProcess = true; - } else { - if (liveLockProcess) { - const termination = await stopDaemonProcess( - { pid: lockInfo.pid, startTime: lockInfo.processStartTime ?? null }, - { - mode: 'graceful', - termTimeoutMs: DAEMON_TAKEOVER_TERM_TIMEOUT_MS, - killTimeoutMs: DAEMON_TAKEOVER_KILL_TIMEOUT_MS, - }, - ); - requireDaemonExit(termination); - result.stoppedLockProcess = true; - } - removeDaemonLock(paths.lockPath); - result.removedLock = true; - } - } else if (lockExists) { - removeDaemonLock(paths.lockPath); - result.removedLock = true; - } - } catch (error) { - result.error = error instanceof Error ? error.message : String(error); - } - - emitDiagnostic({ - level: result.error ? 'warn' : 'info', - phase: 'daemon_startup_metadata_cleanup', - data: result, - }); - return result; -} - export function getDaemonMetadataState(paths: DaemonPaths): DaemonMetadataState { return { hasInfo: fs.existsSync(paths.infoPath), @@ -240,21 +154,6 @@ export function getDaemonMetadataState(paths: DaemonPaths): DaemonMetadataState }; } -export async function recoverDaemonLockHolder(paths: DaemonPaths): Promise { - const state = getDaemonMetadataState(paths); - if (!state.hasLock || state.hasInfo) return false; - const lockInfo = readDaemonLockInfo(paths.lockPath); - if (!lockInfo) { - removeDaemonLock(paths.lockPath); - return true; - } - if (!isAgentDeviceDaemonProcess(lockInfo.pid, lockInfo.processStartTime)) { - removeDaemonLock(paths.lockPath); - return true; - } - return false; -} - export async function stopDaemonProcessForTakeover( info: DaemonInfo, ): Promise { diff --git a/src/daemon-client/daemon-client.ts b/src/daemon-client/daemon-client.ts index 7d777254fb..4d5e97e561 100644 --- a/src/daemon-client/daemon-client.ts +++ b/src/daemon-client/daemon-client.ts @@ -17,17 +17,7 @@ import { prepareRemoteRequestArtifacts, type PreparedRemoteRequest, } from '../remote/daemon-artifacts.ts'; -import { - attachActiveSessionAddressHint, - attachRepairSessionAddressHint, - cleanupDaemonAfterRequest, - ensureDaemon, - isActiveReplaySessionResponse, - isHeldRepairDivergence, - resolveClientSettings, - type DaemonClientSettings, - type EnsuredDaemon, -} from './daemon-client-lifecycle.ts'; +import type { DaemonClientSettings, EnsuredDaemon } from './daemon-client-lifecycle.ts'; import { sendRequest } from './daemon-client-transport.ts'; import { isRemoteDaemon, type DaemonInfo } from './daemon-client-metadata.ts'; import { leaseScopeFromRequest } from '@agent-device/contracts/lease-scope'; @@ -42,6 +32,8 @@ export async function sendToDaemon( req: Omit, options: DaemonTransportOptions = {}, ): Promise { + const { resolveClientSettings, ensureDaemon, attachSessionAddressHints } = + await import('./daemon-client-lifecycle.ts'); const requestId = req.meta?.requestId ?? createRequestId(); const debug = Boolean(req.meta?.debug || req.flags?.verbose); // A few internal callers build DaemonRequest directly instead of using the @@ -107,11 +99,7 @@ export async function sendToDaemon( ), { requestId, command: req.command }, ); - return withActiveSessionAddressHint( - withRepairSessionAddressHintIfOwned(response, settings), - requestWithoutAuthFlag, - settings, - ); + return attachSessionAddressHints(response, requestWithoutAuthFlag, settings); }, ); } @@ -235,6 +223,7 @@ async function performDaemonRequestWithCleanup( requestFailed = true; requestError = error; } + const { cleanupDaemonAfterRequest } = await import('./daemon-client-lifecycle.ts'); const finalResponse = await cleanupDaemonAfterRequest(req, daemon, settings, response); if (requestFailed) throw requestError; if (!finalResponse) { @@ -246,48 +235,6 @@ async function performDaemonRequestWithCleanup( return finalResponse; } -/** - * ADR 0012 decision 6 (Fix 1): the owned ephemeral state dir this daemon was - * started at is otherwise unaddressable by a later invocation — hint it here, - * only when the daemon is actually being kept alive for it - * (`settings.ownedStateDir` means `daemon.startedByClient` is also true). - */ -function withRepairSessionAddressHintIfOwned( - response: DaemonResponse, - settings: DaemonClientSettings, -): DaemonResponse { - if (response.ok || !settings.ownedStateDir || !isHeldRepairDivergence(response)) { - return response; - } - return attachRepairSessionAddressHint(response, settings.paths.baseDir); -} - -/** - * ADR 0016 counterpart to `withRepairSessionAddressHintIfOwned` — but unlike - * that one, NOT gated on `settings.ownedStateDir`. An owned ephemeral state - * dir is unaddressable by a later invocation either way, so it's included - * when owned; an explicit `--state-dir`/`AGENT_DEVICE_STATE_DIR` caller - * already knows their own dir, so it's omitted then. But the session's own - * name is cwd-qualified and, per #1394, `session list` cannot rediscover it - * either — so `--session` is still worth hinting even at an explicit state - * dir, which is why this runs for every active-session response regardless - * of `ownedStateDir` (`attachActiveSessionAddressHint` itself decides what, - * if anything, is worth attaching). - */ -function withActiveSessionAddressHint( - response: DaemonResponse, - req: Omit, - settings: DaemonClientSettings, -): DaemonResponse { - if (!response.ok || !isActiveReplaySessionResponse(req, response)) { - return response; - } - return attachActiveSessionAddressHint( - response, - settings.ownedStateDir ? settings.paths.baseDir : undefined, - ); -} - function writeInstallInProgressNotice(command: string | undefined): void { if (!isInstallLikeCommand(command) || process.stderr.isTTY !== true || process.env.CI) return; process.stderr.write( diff --git a/src/daemon-registration-owner.ts b/src/daemon-registration-owner.ts index 0092f087af..d2afde67bb 100644 --- a/src/daemon-registration-owner.ts +++ b/src/daemon-registration-owner.ts @@ -1,11 +1,18 @@ import fs from 'node:fs'; -import { normalizeError, type NormalizedError } from '@agent-device/kernel/errors'; +import os from 'node:os'; +import path from 'node:path'; +import { runCmdDetachedMonitored, type ExecDetachedExit } from '@agent-device/host-kit/command'; +import { AppError, normalizeError, type NormalizedError } from '@agent-device/kernel/errors'; import { stopDaemonProcess, waitForDaemonExit, type DaemonTerminationResult, } from './daemon-process.ts'; -import { readCurrentOwnerIdentity, type OwnerIdentity } from '@agent-device/host-kit/process'; +import { + readCurrentOwnerIdentity, + readProcessStartTime, + type OwnerIdentity, +} from '@agent-device/host-kit/process'; import { publishFileSync, tryAcquireProcessLock, @@ -15,7 +22,12 @@ import { } from '@agent-device/host-kit/file'; import { emitDiagnostic, withDiagnosticsScope } from '@agent-device/host-kit/diagnostics'; import type { DaemonCodeOrigin } from '@agent-device/host-kit/code-signature'; -import type { DaemonPaths } from './daemon-resolution.ts'; +import { + resolveDaemonPaths, + type DaemonPaths, + type DaemonServerMode, +} from './daemon-resolution.ts'; +import { findUnrecoveredRepairCommitFailure } from './session-repair-tombstone.ts'; import { readRegisteredDaemonOwnership, type RegisteredDaemonOwnership, @@ -161,6 +173,95 @@ function truncateDaemonLog(logPath: string): void { } } +declare const privateReplayState: unique symbol; +export type OwnedReplayStateDir = Readonly<{ + paths: Readonly; + [privateReplayState]: true; +}>; +export type DaemonStartupLaunch = Readonly<{ + pid: number; + startTime?: string; + exited: Promise; +}>; +type OwnedStartup = { launch: DaemonStartupLaunch; joined: boolean }; +type PrivateReplayState = { + paths: Readonly; + startups: OwnedStartup[]; + sealed: boolean; + retirement?: Promise; +}; +const privateReplayStates = new WeakMap(); + +/** Creates deletion authority only for a fresh private replay directory. */ +export function createOwnedReplayStateDir(): OwnedReplayStateDir { + const paths = Object.freeze( + resolveDaemonPaths(fs.mkdtempSync(path.join(os.tmpdir(), 'agent-device-replay-daemon-'))), + ); + const owned = Object.freeze({ paths }) as OwnedReplayStateDir; + privateReplayStates.set(owned, { paths, startups: [], sealed: false }); + return owned; +} + +/** Launches and monitors the actual child before recording its private-directory authority. */ +export function launchDaemonProcess( + input: Readonly<{ + paths: DaemonPaths; + args: string[]; + serverMode: DaemonServerMode; + ownedStateDir?: OwnedReplayStateDir; + }>, +): DaemonStartupLaunch { + const owned = input.ownedStateDir && requirePrivateReplayState(input.ownedStateDir, input.paths); + if (owned?.sealed) + throw new AppError('COMMAND_FAILED', 'Replay daemon startup admission is closed.', { + reason: 'daemon_startup_admission_closed', + }); + fs.mkdirSync(input.paths.baseDir, { recursive: true }); + const logFd = fs.openSync(input.paths.logPath, 'a'); + try { + const monitored = runCmdDetachedMonitored(process.execPath, input.args, { + env: { + ...process.env, + AGENT_DEVICE_STATE_DIR: input.paths.baseDir, + AGENT_DEVICE_DAEMON_SERVER_MODE: input.serverMode, + }, + stdio: ['ignore', logFd, logFd], + }); + const startup: OwnedStartup = { launch: Object.freeze({ ...monitored }), joined: false }; + if (owned) { + owned.startups.push(startup); + void monitored.exited.then(() => { + startup.joined = true; + }); + } + startup.launch = Object.freeze({ + ...startup.launch, + startTime: readProcessStartTime(monitored.pid) ?? undefined, + }); + return startup.launch; + } finally { + fs.closeSync(logFd); + } +} + +function requirePrivateReplayState( + capability: OwnedReplayStateDir, + paths: DaemonPaths, +): PrivateReplayState { + const owned = privateReplayStates.get(capability); + if ( + !owned || + Object.entries(owned.paths).some(([key, value]) => paths[key as keyof DaemonPaths] !== value) + ) { + throw new AppError( + 'COMMAND_FAILED', + 'Private replay directory ownership could not be verified.', + { reason: 'daemon_private_state_unowned' }, + ); + } + return owned; +} + export type DaemonRetirementInput = Readonly<{ paths: DaemonPaths; observed: OwnerIdentity | null; @@ -169,16 +270,25 @@ export type DaemonRetirementInput = Readonly<{ lockTimeoutMs?: number; }>; +type RepairCommitFailure = NonNullable>; type ConfirmedDaemonTermination = Extract; export type DaemonRetirementResult = - | Readonly<{ status: 'retired'; termination: ConfirmedDaemonTermination; removedInfo: boolean }> + | Readonly<{ + status: 'retired'; + termination: ConfirmedDaemonTermination; + removedInfo: boolean; + removedStateDir?: boolean; + repairCommitFailure?: RepairCommitFailure; + }> | Readonly<{ status: 'absent'; removedInfo: false }> | Readonly<{ status: 'retained'; termination?: DaemonTerminationResult; removedInfo: boolean; + removedStateDir?: boolean; reason: | 'ownership-unproven' + | 'startup-unconfirmed' | 'exit-unconfirmed' | 'stop-failed' | 'lock-busy' @@ -186,19 +296,59 @@ export type DaemonRetirementResult = | 'metadata-unreadable' | 'retirement-unconfirmed'; error?: NormalizedError; + repairCommitFailure?: RepairCommitFailure; }>; /** Stops only the captured daemon lifetime, then retires its registration under the startup lock. */ export async function stopAndRetireDaemon( - input: DaemonRetirementInput & Readonly<{ mode: 'graceful' | 'force' }>, + input: DaemonRetirementInput & + Readonly<{ + mode: 'graceful' | 'force'; + ownedStateDir?: OwnedReplayStateDir; + startupJoinTimeoutMs?: number; + }>, ): Promise { - return await retireObservedDaemon(input, (identity) => - stopDaemonProcess(identity, { - mode: input.mode, - termTimeoutMs: input.termTimeoutMs ?? 3_000, - killTimeoutMs: input.killTimeoutMs ?? 1_000, - }), + let owned: PrivateReplayState | undefined; + try { + if (input.ownedStateDir) { + owned = requirePrivateReplayState(input.ownedStateDir, input.paths); + owned.sealed = true; + const launch = owned.startups.find( + ({ launch }) => + launch.pid === input.observed?.pid && launch.startTime === input.observed?.startTime, + )?.launch; + if (!launch?.startTime) + throw new AppError( + 'COMMAND_FAILED', + 'The observed daemon is not an owned startup lifetime.', + { reason: 'daemon_private_startup_unowned' }, + ); + if (owned.retirement) return await owned.retirement; + } + } catch (error) { + return { + status: 'retained', + reason: 'ownership-unproven', + removedInfo: false, + error: normalizeError(error), + }; + } + const retirement = retireObservedDaemon( + input, + (identity) => + stopDaemonProcess(identity, { + mode: input.mode, + termTimeoutMs: input.termTimeoutMs ?? 3_000, + killTimeoutMs: input.killTimeoutMs ?? 1_000, + }), + owned, + input.startupJoinTimeoutMs ?? 1_000, ); + if (owned) owned.retirement = retirement; + const result = await retirement; + if (owned && (result.status === 'absent' || !result.removedStateDir)) + owned.retirement = undefined; + return result; } /** Recovers a confirmed abandoned registration without signaling a live process. */ @@ -217,6 +367,8 @@ export async function recoverAbandonedDaemonRegistration( async function retireObservedDaemon( input: DaemonRetirementInput, terminate: (identity: OwnerIdentity) => Promise, + owned?: PrivateReplayState, + startupJoinTimeoutMs = 1_000, ): Promise { const paths = { ...input.paths }; const observed = input.observed && { ...input.observed }; @@ -241,12 +393,16 @@ async function retireObservedDaemon( }; } } - return await retireDaemonRegistration({ ...input, paths }, termination); + if (owned && !(await joinOwnedStartups(owned, startupJoinTimeoutMs))) { + return { status: 'retained', reason: 'startup-unconfirmed', termination, removedInfo: false }; + } + return await retireDaemonRegistration({ ...input, paths }, termination, owned); } async function retireDaemonRegistration( input: DaemonRetirementInput, termination: ConfirmedDaemonTermination | undefined, + owned?: PrivateReplayState, ): Promise { const paths = input.paths; let acquisition: ProcessLockAcquisition; @@ -268,7 +424,11 @@ async function retireDaemonRegistration( error: failure, }; } - let result: DaemonRetirementResult; + let result: DaemonRetirementResult = { + status: 'retained', + reason: 'retirement-unconfirmed', + removedInfo: false, + }; try { const removal = removeRegistrationUnderLock( paths.infoPath, @@ -276,12 +436,13 @@ async function retireDaemonRegistration( acquisition, ); result = retirementAfterRemoval(removal, termination); + result = retirePrivateStateIfEligible(owned, acquisition, result); } catch (error) { result = { status: 'retained', reason: 'retirement-unconfirmed', termination, - removedInfo: false, + removedInfo: result.removedInfo, error: normalizeError(error), }; } @@ -339,3 +500,42 @@ async function recordRegistrationWarning( emitDiagnostic({ level: 'warn', phase, data: { error: normalizeError(error) } }); }); } + +async function joinOwnedStartups(owned: PrivateReplayState, timeoutMs: number): Promise { + if (owned.startups.every((startup) => startup.joined)) return true; + let timer: ReturnType | undefined; + try { + return await Promise.race([ + Promise.all(owned.startups.map(({ launch }) => launch.exited)).then(() => true), + new Promise((resolve) => { + timer = setTimeout(() => resolve(false), timeoutMs); + }), + ]); + } finally { + if (timer) clearTimeout(timer); + } +} + +function retirePrivateStateIfEligible( + owned: PrivateReplayState | undefined, + acquisition: ProcessLockAcquisition, + result: DaemonRetirementResult, +): DaemonRetirementResult { + if (!owned || result.status !== 'retired') return result; + try { + acquisition.assertHeld(); + const failure = findUnrecoveredRepairCommitFailure(owned.paths.sessionsDir); + if (failure) return { ...result, removedStateDir: false, repairCommitFailure: failure }; + acquisition.assertHeld(); + fs.rmSync(owned.paths.baseDir, { recursive: true, force: true }); + return { ...result, removedStateDir: true }; + } catch (error) { + return { + ...result, + status: 'retained', + reason: 'retirement-unconfirmed', + removedStateDir: false, + error: normalizeError(error), + }; + } +} diff --git a/src/daemon/__tests__/session-artifact-paths.test.ts b/src/daemon/__tests__/session-artifact-paths.test.ts new file mode 100644 index 0000000000..c76caa9ffb --- /dev/null +++ b/src/daemon/__tests__/session-artifact-paths.test.ts @@ -0,0 +1,29 @@ +import { test } from 'vitest'; +import assert from 'node:assert/strict'; +import path from 'node:path'; +import { AppError } from '@agent-device/kernel/errors'; +import { resolveSessionDir } from '../session-artifact-paths.ts'; +import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; + +test('resolveSessionDir keeps every session dir beneath the sessions dir', () => { + const sessionsDir = path.join( + mkdtempForTestSync('agent-device-tests'), + 'agent-device-tests', + 'sessions', + ); + assert.equal(resolveSessionDir(sessionsDir, 'a/b:c d'), path.join(sessionsDir, 'a_b_c_d')); + // `.` and `..` survive `safeSessionName` unchanged, so without an explicit + // refusal `path.join` resolves them to the sessions dir itself and its parent + // (the daemon state dir): a remote caller's `--session ..` would then land + // app.log / runner.log / requests/*.ndjson outside the sessions tree. + for (const name of ['.', '..', '']) { + assert.throws( + () => resolveSessionDir(sessionsDir, name), + (error: unknown) => + error instanceof AppError && + error.code === 'INVALID_ARGS' && + /session name/i.test(error.message), + `expected resolveSessionDir(${JSON.stringify(name)}) to reject`, + ); + } +}); diff --git a/src/daemon/__tests__/session-store.test.ts b/src/daemon/__tests__/session-store.test.ts index 05fe386324..5eddff721d 100644 --- a/src/daemon/__tests__/session-store.test.ts +++ b/src/daemon/__tests__/session-store.test.ts @@ -1,10 +1,7 @@ import { test } from 'vitest'; import assert from 'node:assert/strict'; import fs from 'node:fs'; -// oxlint-disable-next-line no-restricted-imports -- asserts a path under os.homedir -import os from 'node:os'; import path from 'node:path'; -import { AppError } from '@agent-device/kernel/errors'; import { SessionStore } from '../session-store.ts'; import type { SessionState } from '../session-state.ts'; import { buildRequestFinishedEvent } from '@agent-device/session-journal/session-event-log'; @@ -96,19 +93,6 @@ function assertScriptMatches(script: string, patterns: RegExp[]): void { } } -test('expandHome resolves tilde, relative-with-cwd, and absolute paths', () => { - const homePath = SessionStore.expandHome('~/flows/replay.ad'); - assert.equal(homePath.startsWith(os.homedir()), true); - assert.equal(homePath.endsWith(path.join('flows', 'replay.ad')), true); - - const relativePath = SessionStore.expandHome('workflows/replay.ad', '/tmp/agent-device-cwd'); - assert.equal(relativePath, path.resolve('/tmp/agent-device-cwd', 'workflows/replay.ad')); - - const absoluteInput = path.resolve('/tmp', 'agent-device-absolute.ad'); - const absolutePath = SessionStore.expandHome(absoluteInput, '/tmp/ignored-cwd'); - assert.equal(absolutePath, absoluteInput); -}); - test('defaultTracePath sanitizes session name', () => { const store = new SessionStore( path.join(mkdtempForTestSync('agent-device-tests'), 'agent-device-tests'), @@ -119,30 +103,6 @@ test('defaultTracePath sanitizes session name', () => { assert.match(tracePath, /\.trace\.log$/); }); -test('resolveSessionDir keeps every session dir beneath the sessions dir', () => { - const sessionsDir = path.join( - mkdtempForTestSync('agent-device-tests'), - 'agent-device-tests', - 'sessions', - ); - const store = new SessionStore(sessionsDir); - assert.equal(store.resolveSessionDir('a/b:c d'), path.join(sessionsDir, 'a_b_c_d')); - // `.` and `..` survive `safeSessionName` unchanged, so without an explicit - // refusal `path.join` resolves them to the sessions dir itself and its parent - // (the daemon state dir): a remote caller's `--session ..` would then land - // app.log / runner.log / requests/*.ndjson outside the sessions tree. - for (const name of ['.', '..', '']) { - assert.throws( - () => store.resolveSessionDir(name), - (error: unknown) => - error instanceof AppError && - error.code === 'INVALID_ARGS' && - /session name/i.test(error.message), - `expected resolveSessionDir(${JSON.stringify(name)}) to reject`, - ); - } -}); - test('session lease metadata round-trips through the store', () => { const { store, session } = makeFixture('agent-device-session-lease-'); session.lease = { diff --git a/src/daemon/device/device-claim-owner-recovery.ts b/src/daemon/device/device-claim-owner-recovery.ts index c8b60069ee..f57a6c641f 100644 --- a/src/daemon/device/device-claim-owner-recovery.ts +++ b/src/daemon/device/device-claim-owner-recovery.ts @@ -4,7 +4,11 @@ import { createPlatformRuntimeGateway } from '../../platform-runtime.ts'; import { resolveDaemonPaths } from '../../daemon-resolution.ts'; import { createDeviceClaimReconciler } from './device-claim-reconciliation.ts'; import type { DeviceClaimReconciler } from './device-claims.ts'; -import { SessionStore } from '../session-store.ts'; +import { + resolveSessionDir, + resolveSessionAppLogPath, + resolveSessionAppLogPidPath, +} from '../session-artifact-paths.ts'; export type OwnerScopedClaimRecovery = { reconcile: DeviceClaimReconciler; @@ -49,17 +53,16 @@ function composeOwnerScopedClaimRecovery( scope: PlatformRequestScope, ): OwnerScopedClaimRecovery { const daemonPaths = resolveDaemonPaths(stateDir); - const sessionStore = new SessionStore(daemonPaths.sessionsDir); const gateway = createPlatformRuntimeGateway({ sessionsDir: daemonPaths.sessionsDir, ownedProcesses: createOwnedProcessRecordStore({ stateDir: daemonPaths.baseDir, sessionsDir: daemonPaths.sessionsDir, - resolveSessionDir: (sessionId) => sessionStore.resolveSessionDir(sessionId), + resolveSessionDir: (sessionId) => resolveSessionDir(daemonPaths.sessionsDir, sessionId), }), resolveSessionArtifacts: (sessionId) => ({ - outputPath: sessionStore.resolveAppLogPath(sessionId), - pidPath: sessionStore.resolveAppLogPidPath(sessionId), + outputPath: resolveSessionAppLogPath(daemonPaths.sessionsDir, sessionId), + pidPath: resolveSessionAppLogPidPath(daemonPaths.sessionsDir, sessionId), }), }); return { diff --git a/src/daemon/handlers/session-app-deployment.ts b/src/daemon/handlers/session-app-deployment.ts index edf468e306..79768d7158 100644 --- a/src/daemon/handlers/session-app-deployment.ts +++ b/src/daemon/handlers/session-app-deployment.ts @@ -1,3 +1,4 @@ +import { expandSessionPath } from '@agent-device/host-kit/session-paths'; import fs from 'node:fs'; import type { AppDeploymentResult } from '@agent-device/contracts/app-deployment-runtime'; import { @@ -9,7 +10,7 @@ import { readNotificationPayload } from '../dispatch-payload.ts'; import { cleanupUploadedArtifact, prepareUploadedArtifact } from '../artifact-tracking.ts'; import { expireRefFrame } from '../ref-frame.ts'; import type { BindDeviceRuntime, InspectDeviceRuntimeFacts } from '../request-runtime-binding.ts'; -import { SessionStore } from '../session-store.ts'; +import type { SessionStore } from '../session-store.ts'; import type { DaemonRequest, DaemonResponse } from '../daemon-request.ts'; import type { SessionState } from '../session-state.ts'; import { resolvePayloadInput } from '../payload-input.ts'; @@ -62,7 +63,7 @@ export async function handleAppDeploymentCommand(params: { try { const appPath = uploadedArtifactId ? prepareUploadedArtifact(uploadedArtifactId, req.meta?.tenantId) - : SessionStore.expandHome(target.appPathInput); + : expandSessionPath(target.appPathInput); if (!fs.existsSync(appPath)) { return errorResponse('INVALID_ARGS', `App binary not found: ${appPath}`); } @@ -242,7 +243,7 @@ function resolvePushPayload(payloadArg: string, cwd?: string): string { const resolved = resolvePayloadInput(payloadArg, { subject: 'Push payload', cwd, - expandPath: (value, currentCwd) => SessionStore.expandHome(value, currentCwd), + expandPath: (value, currentCwd) => expandSessionPath(value, currentCwd), }); return resolved.kind === 'file' ? resolved.path : resolved.text; } diff --git a/src/daemon/handlers/trace-runtime.ts b/src/daemon/handlers/trace-runtime.ts index 43b8ee86c6..64f804f0b5 100644 --- a/src/daemon/handlers/trace-runtime.ts +++ b/src/daemon/handlers/trace-runtime.ts @@ -1,7 +1,8 @@ +import { expandSessionPath } from '@agent-device/host-kit/session-paths'; import fs from 'node:fs'; import path from 'node:path'; import type { TraceCommandResult } from '@agent-device/contracts/recording'; -import { SessionStore } from '../session-store.ts'; +import type { SessionStore } from '../session-store.ts'; import type { DaemonRequest, DaemonResponse } from '../daemon-request.ts'; import type { SessionState } from '../session-state.ts'; import { recordSessionAction } from '../session-action-recorder.ts'; @@ -29,9 +30,7 @@ function startTrace( session: SessionState, ): DaemonResponse { if (session.trace) return errorResponse('INVALID_ARGS', 'trace already in progress'); - const outPath = SessionStore.expandHome( - req.positionals?.[1] ?? sessionStore.defaultTracePath(session), - ); + const outPath = expandSessionPath(req.positionals?.[1] ?? sessionStore.defaultTracePath(session)); fs.mkdirSync(path.dirname(outPath), { recursive: true }); fs.appendFileSync(outPath, ''); session.trace = { outPath, startedAt: Date.now() }; @@ -72,7 +71,7 @@ function stopTrace( function relocateTraceOutput(currentPath: string, requestedPath: string | undefined): string { if (!requestedPath) return currentPath; - const resolved = SessionStore.expandHome(requestedPath); + const resolved = expandSessionPath(requestedPath); fs.mkdirSync(path.dirname(resolved), { recursive: true }); if (fs.existsSync(currentPath)) fs.renameSync(currentPath, resolved); else fs.appendFileSync(resolved, ''); diff --git a/src/daemon/screenshot-runtime.ts b/src/daemon/screenshot-runtime.ts index f0e4c76477..f10f44ab42 100644 --- a/src/daemon/screenshot-runtime.ts +++ b/src/daemon/screenshot-runtime.ts @@ -1,3 +1,4 @@ +import { expandSessionPath } from '@agent-device/host-kit/session-paths'; import type { CommandFlags } from '@agent-device/contracts/command'; import { retiredScreenshotMaxSizeFlagError, @@ -38,7 +39,6 @@ import { type ScreenshotRuntimeBindings, } from './screenshot-runtime-binding.ts'; import { setSessionSnapshot } from './session-snapshot.ts'; -import { SessionStore } from './session-store.ts'; import type { DaemonRequest } from './daemon-request.ts'; import type { SessionState } from './session-state.ts'; @@ -321,7 +321,7 @@ function readScreenshotRequest( const positionals = req.positionals ?? []; const flags = req.flags ?? {}; const expand = (value: string | undefined) => - value === undefined ? undefined : SessionStore.expandHome(value, req.meta?.cwd); + value === undefined ? undefined : expandSessionPath(value, req.meta?.cwd); const positionalPath = expand(positionals[0]); const outFlag = expand(flags.out); return { diff --git a/src/daemon/session-artifact-paths.ts b/src/daemon/session-artifact-paths.ts index d456d4d199..b445ed1ec1 100644 --- a/src/daemon/session-artifact-paths.ts +++ b/src/daemon/session-artifact-paths.ts @@ -1,6 +1,7 @@ import path from 'node:path'; import type { DiagnosticsRecordRef } from '@agent-device/kernel/errors'; -import { safeSessionName } from '@agent-device/host-kit/session-paths'; +import { AppError } from '@agent-device/kernel/errors'; +import { isSafeSessionSegment, safeSessionName } from '@agent-device/host-kit/session-paths'; /** Path to session-scoped platform subprocess output, such as Apple runner xcodebuild logs. */ export function resolveSessionRunnerLogPath(sessionDir: string): string { @@ -49,3 +50,27 @@ export function resolveRemoteRequestDiagnosticsPath( ref.requestId, ); } + +/** + * The one place a session name becomes a directory, so the invariant that every + * session dir lies beneath `sessionsDir` is enforced here rather than by each + * caller: `.` and `..` survive `safeSessionName` and would resolve to the + * sessions dir itself or the daemon state dir above it. + */ +export function resolveSessionDir(sessionsDir: string, sessionName: string): string { + if (!isSafeSessionSegment(sessionName)) { + throw new AppError( + 'INVALID_ARGS', + `Invalid session name ${JSON.stringify(sessionName)}: a session name cannot be empty, ".", or "..".`, + ); + } + return path.join(sessionsDir, safeSessionName(sessionName)); +} + +export function resolveSessionAppLogPath(sessionsDir: string, address: string): string { + return path.join(resolveSessionDir(sessionsDir, address), 'app.log'); +} + +export function resolveSessionAppLogPidPath(sessionsDir: string, address: string): string { + return path.join(resolveSessionDir(sessionsDir, address), 'app-log.pid'); +} diff --git a/src/daemon/session-observability/internal/session-perf-runtime.ts b/src/daemon/session-observability/internal/session-perf-runtime.ts index 210528eea6..cb6acdb0a7 100644 --- a/src/daemon/session-observability/internal/session-perf-runtime.ts +++ b/src/daemon/session-observability/internal/session-perf-runtime.ts @@ -1,3 +1,4 @@ +import { expandSessionPath } from '@agent-device/host-kit/session-paths'; import path from 'node:path'; import { type PerfCaptureAdmissionLedger } from '@agent-device/capture-kit/perf-capture-admission-ledger'; import { @@ -28,7 +29,7 @@ import type { BindDeviceRuntime, InspectDeviceRuntimeFacts, } from '../../request-runtime-binding.ts'; -import { SessionStore } from '../../session-store.ts'; +import type { SessionStore } from '../../session-store.ts'; import type { DaemonRequest, DaemonResponse } from '../../daemon-request.ts'; import type { SessionState } from '../../session-state.ts'; import { recordSessionAction } from '../../session-action-recorder.ts'; @@ -138,7 +139,7 @@ async function executeAdmittedPerfPlan( appId: session.appBundleId, kind: plan.request.kind, outPath: plan.request.outPath - ? SessionStore.expandHome(plan.request.outPath, params.req.meta?.cwd) + ? expandSessionPath(plan.request.outPath, params.req.meta?.cwd) : undefined, artifactsDir: path.join( params.sessionStore.ensureSessionDir(params.sessionName), @@ -171,7 +172,7 @@ async function executeAdmittedPerfPlan( const data = await runtime.operations.perfProfileReport({ appId: session.appBundleId, kind: plan.request.kind, - tracePath: SessionStore.expandHome(tracePath, params.req.meta?.cwd), + tracePath: expandSessionPath(tracePath, params.req.meta?.cwd), outPath, template: plan.request.template ?? (last?.kind === 'xctrace' ? last.template : undefined), profile: last, @@ -257,7 +258,7 @@ async function stopPerfCapture( const mismatchMessage = perfCaptureStopMismatch(snapshot, request); if (mismatchMessage) return errorResponse('INVALID_ARGS', mismatchMessage); if (request.outPath) { - capture.handle.setOutputPath(SessionStore.expandHome(request.outPath, params.req.meta?.cwd)); + capture.handle.setOutputPath(expandSessionPath(request.outPath, params.req.meta?.cwd)); } const completion = await finishLivePerfCapture({ intent: 'capture', @@ -395,7 +396,7 @@ function resolveNativeOutPath( requestedPath: string | undefined, fallbackFileName: string, ): string { - if (requestedPath) return SessionStore.expandHome(requestedPath, params.req.meta?.cwd); + if (requestedPath) return expandSessionPath(requestedPath, params.req.meta?.cwd); return path.join( params.sessionStore.ensureSessionDir(params.sessionName), `${timestampToken()}-${fallbackFileName}`, diff --git a/src/daemon/session-store.ts b/src/daemon/session-store.ts index 1598dd9414..e31c66403a 100644 --- a/src/daemon/session-store.ts +++ b/src/daemon/session-store.ts @@ -1,14 +1,14 @@ import path from 'node:path'; import fs from 'node:fs'; -import { AppError } from '@agent-device/kernel/errors'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; import type { SessionRef, SessionRuntimeHints, SessionState } from './session-state.ts'; import { recordActionEntry, type RecordActionEntry } from './session-action-recorder.ts'; +import { isSafeSessionSegment, safeSessionName } from '@agent-device/host-kit/session-paths'; import { - expandSessionPath, - isSafeSessionSegment, - safeSessionName, -} from '@agent-device/host-kit/session-paths'; + resolveSessionDir, + resolveSessionAppLogPath, + resolveSessionAppLogPidPath, +} from './session-artifact-paths.ts'; import { readRepairTombstoneFile, resolveRepairTombstonePath, @@ -380,20 +380,8 @@ export class SessionStore { return path.join(this.sessionsDir, `${safeName}-${timestamp}.trace.log`); } - /** - * The one place a session name becomes a directory, so the invariant that every - * session dir lies beneath `sessionsDir` is enforced here rather than by each - * caller: `.` and `..` survive `safeSessionName` and would resolve to the - * sessions dir itself or the daemon state dir above it. - */ resolveSessionDir(sessionName: string): string { - if (!isSafeSessionSegment(sessionName)) { - throw new AppError( - 'INVALID_ARGS', - `Invalid session name ${JSON.stringify(sessionName)}: a session name cannot be empty, ".", or "..".`, - ); - } - return path.join(this.sessionsDir, safeSessionName(sessionName)); + return resolveSessionDir(this.sessionsDir, sessionName); } // Daemon state dir (parent of the `sessions/` dir), matching daemonPaths.baseDir. Called via @@ -411,21 +399,17 @@ export class SessionStore { /** Path to session-scoped app log file. Agent can grep this for token-efficient debugging. */ resolveAppLogPath(sessionName: string): string { - return path.join(this.resolveSessionDir(sessionName), 'app.log'); + return resolveSessionAppLogPath(this.sessionsDir, sessionName); } resolveAppLogPidPath(sessionName: string): string { - return path.join(this.resolveSessionDir(sessionName), 'app-log.pid'); + return resolveSessionAppLogPidPath(this.sessionsDir, sessionName); } resolveEventLogPath(sessionName: string): string { return resolveSessionEventLogPath(this.resolveSessionDir(sessionName)); } - static expandHome(filePath: string, cwd?: string): string { - return expandSessionPath(filePath, cwd); - } - /** * Resolve the map key for a live session object. SessionState.name is the * public session name, while the map key may include cwd/tenant isolation. diff --git a/src/session-repair-tombstone.ts b/src/session-repair-tombstone.ts index 59c2b60a9f..a0fae0461b 100644 --- a/src/session-repair-tombstone.ts +++ b/src/session-repair-tombstone.ts @@ -1,5 +1,6 @@ import path from 'node:path'; import fs from 'node:fs'; +import { AppError } from '@agent-device/kernel/errors'; /** * ADR 0012 decision 6, R7 (C5a): a reaped repair session leaves this bounded @@ -31,20 +32,51 @@ export function resolveRepairTombstonePath(sessionDir: string): string { /** Parses/validates a tombstone file at `tombstonePath`; `undefined` if missing, malformed, or expired. */ export function readRepairTombstoneFile(tombstonePath: string): RepairSessionTombstone | undefined { - let raw: string; try { - raw = fs.readFileSync(tombstonePath, 'utf8'); + return readRepairTombstoneForCleanup(tombstonePath); } catch { return undefined; } - let parsed: RepairSessionTombstone; +} + +function readRepairTombstoneForCleanup(tombstonePath: string): RepairSessionTombstone | undefined { + let raw: string; try { - parsed = JSON.parse(raw) as RepairSessionTombstone; - } catch { - return undefined; + raw = fs.readFileSync(tombstonePath, 'utf8'); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return undefined; + throw error; } - if (typeof parsed?.expiresAt !== 'number' || parsed.expiresAt <= Date.now()) return undefined; - return parsed; + const parsed = parseRepairTombstone(raw, tombstonePath); + return parsed.expiresAt > Date.now() ? parsed : undefined; +} + +function parseRepairTombstone(raw: string, tombstonePath: string): RepairSessionTombstone { + try { + const parsed = JSON.parse(raw) as RepairSessionTombstone; + if ( + !Number.isFinite(parsed?.expiresAt) || + typeof parsed?.owner !== 'string' || + !validRepairCommitFailure(parsed.commitFailure) + ) + throw new Error('Invalid repair tombstone fields'); + return parsed; + } catch (error) { + throw new AppError( + 'COMMAND_FAILED', + 'Repair evidence could not be inspected.', + { reason: 'repair_evidence_invalid', path: tombstonePath }, + error instanceof Error ? error : undefined, + ); + } +} + +function validRepairCommitFailure(value: unknown): boolean { + const failure = value as RepairSessionTombstone['commitFailure'] | null; + return ( + value === undefined || + (typeof failure?.code === 'string' && typeof failure?.message === 'string') + ); } /** @@ -54,6 +86,7 @@ export function readRepairTombstoneFile(tombstonePath: string): RepairSessionTom * CLIENT side of the daemon boundary (`cleanupDaemonAfterRequest` in * `daemon-client-lifecycle.ts`), which has no live `SessionStore`/session name * to key off of, only the filesystem path an owned ephemeral daemon was given. + * Unreadable or malformed evidence throws so cleanup retains the directory. * An owned ephemeral state dir services exactly one repair transaction at a * time, so the first match found is returned. * @@ -71,12 +104,13 @@ export function findUnrecoveredRepairCommitFailure(sessionsDir: string): let entries: fs.Dirent[]; try { entries = fs.readdirSync(sessionsDir, { withFileTypes: true }); - } catch { - return undefined; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return undefined; + throw error; } for (const entry of entries) { if (!entry.isDirectory()) continue; - const tombstone = readRepairTombstoneFile( + const tombstone = readRepairTombstoneForCleanup( resolveRepairTombstonePath(path.join(sessionsDir, entry.name)), ); if (tombstone?.commitFailure) {