diff --git a/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts b/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts index ab0ccd73e0..02cc494180 100644 --- a/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts +++ b/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts @@ -1,28 +1,16 @@ import path from 'node:path'; import { safeSessionName } from '@agent-device/host-kit/session-paths'; import { mkdtempForTestSync } from '../../tmp-dir.fixtures.ts'; +import { makeCaptureFixtureStore } from '../../durable-capture/session-binding.fixtures.ts'; -export type CaptureAdmissionSessionStore = Readonly<{ - set(name: string, session: S): void; - resolveSessionDir(name: string): string; - get(name: string): S | undefined; - sessionsDir: string; -}>; +export type CaptureAdmissionSessionStore = ReturnType< + typeof makeCaptureAdmissionSessionStore +>; -/** - * The whole of the daemon `SessionStore` these admission modules ever address — `set`, - * `resolveSessionDir`, and the read-back a test asserts on — over a fresh temp directory, so a - * test of this family needs no session record, store class, or daemon import. - */ -export function makeCaptureAdmissionSessionStore( - prefix: string, -): CaptureAdmissionSessionStore { +export function makeCaptureAdmissionSessionStore(prefix: string) { const sessionsDir = mkdtempForTestSync(prefix); - const sessions = new Map(); - return { - set: (name, session) => void sessions.set(name, session), - get: (name) => sessions.get(name), - resolveSessionDir: (name) => path.join(sessionsDir, safeSessionName(name)), + return Object.freeze({ + ...makeCaptureFixtureStore((name) => path.join(sessionsDir, safeSessionName(name))), sessionsDir, - }; + }); } diff --git a/packages/capture-kit/src/durable-capture/adoption.test.ts b/packages/capture-kit/src/durable-capture/adoption.test.ts index ac2cbff303..f3fb4a521d 100644 --- a/packages/capture-kit/src/durable-capture/adoption.test.ts +++ b/packages/capture-kit/src/durable-capture/adoption.test.ts @@ -68,3 +68,25 @@ test('a failed terminal transition preserves the primary error and reports it un reason: expect.stringMatching(/./), }); }); + +test('adoption refuses a retired lifetime even when its address has a vacant successor', async () => { + const context = makeDurableCaptureContext(); + const start = makeDurableCaptureStartResult(context); + context.sessionStore.retire(context.sessionStore.lookup(context.sessionName)); + const successor = { name: 'successor' }; + context.sessionStore.set(context.sessionName, successor); + await expect( + adoptStartedDurableCapture( + testCaptureDefinition, + { + ...context, + ...start, + throwIfCanceled: () => {}, + }, + context.resourcePath, + ), + ).rejects.toMatchObject({ code: 'COMMAND_FAILED' }); + expect(start.forceCleanup).toHaveBeenCalledOnce(); + expect(context.sessionStore.get(context.sessionName)).toBe(successor); + expect(testCaptureStore.read(context.resourcePath).status).toBe('missing'); +}); diff --git a/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts b/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts index 363b83a45c..9496e6845f 100644 --- a/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts +++ b/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts @@ -1,4 +1,4 @@ -import { makeCaptureSessionBinding } from './session-binding.fixtures.ts'; +import { makeCaptureSessionBinding, makeCaptureFixtureStore } from './session-binding.fixtures.ts'; import path from 'node:path'; import { vi, type Mock } from 'vitest'; import { @@ -78,15 +78,10 @@ export function makeDurableCaptureContext( ) { const sessionsDir = mkdtempForTestSync('durable-capture-resource-'); const sessionName = 'session'; - const sessions = new Map(); const session: TestCaptureSession = { name: sessionName }; - sessions.set(sessionName, session); const resolveSessionDir = (name: string): string => path.join(sessionsDir, name); - const sessionStore = { - get: (name: string) => sessions.get(name), - set: (name: string, next: TestCaptureSession) => void sessions.set(name, next), - resolveSessionDir, - }; + const sessionStore = makeCaptureFixtureStore(resolveSessionDir); + sessionStore.set(sessionName, session); const reportUndurableCleanup: Mock< (device: DeviceInfo, outcome: DurableCaptureCleanupOutcome) => void > = vi.fn(); @@ -96,7 +91,7 @@ export function makeDurableCaptureContext( read: (session) => session.capture, replace: (session, capture) => ({ ...session, capture }), }), - sessions, + sessions: sessionStore, sessionsDir, resolveSessionDir, session, diff --git a/packages/capture-kit/src/durable-capture/session-binding.fixtures.ts b/packages/capture-kit/src/durable-capture/session-binding.fixtures.ts index 8c37b7be49..e84740ee55 100644 --- a/packages/capture-kit/src/durable-capture/session-binding.fixtures.ts +++ b/packages/capture-kit/src/durable-capture/session-binding.fixtures.ts @@ -1,21 +1,51 @@ import { AppError } from '@agent-device/kernel/errors'; import type { DurableCaptureSessionBinding, DurableCaptureSessionResource } from './definition.ts'; +type FixtureSessionRef = Readonly<{ address: string; session: S; lifetime: object }>; + +export function makeCaptureFixtureStore(resolveSessionDir: (address: string) => string) { + const entries = new Map(); + const resolveCurrent = (ref: FixtureSessionRef): S | undefined => { + const entry = entries.get(ref.address); + return entry === ref.lifetime ? entry.current : undefined; + }; + return Object.freeze({ + resolveSessionDir, + get: (address: string): S | undefined => entries.get(address)?.current, + set: (address: string, session: S): void => { + const entry = entries.get(address); + if (entry) entry.current = session; + else entries.set(address, { current: session }); + }, + lookup: (address: string): FixtureSessionRef => { + const entry = entries.get(address); + if (!entry) throw new AppError('COMMAND_FAILED', 'Test session retired'); + return Object.freeze({ address, session: entry.current, lifetime: entry }); + }, + resolveCurrent, + update: (ref: FixtureSessionRef, rebuild: (current: S) => S): void => { + const current = resolveCurrent(ref); + if (current === undefined) throw new AppError('COMMAND_FAILED', 'Test session retired'); + entries.get(ref.address)!.current = rebuild(current); + }, + retire: (ref: FixtureSessionRef): boolean => + resolveCurrent(ref) !== undefined && entries.delete(ref.address), + }); +} + export function makeCaptureSessionBinding( - store: Readonly<{ - get(name: string): S | undefined; - set(name: string, session: S): void; - resolveSessionDir(name: string): string; - }>, + store: ReturnType>, address: string, slot: Readonly<{ read(session: S): DurableCaptureSessionResource | undefined; replace(session: S, resource: DurableCaptureSessionResource | undefined): S; }>, ): DurableCaptureSessionBinding { + const ref = store.lookup(address); + let retained = slot.read(ref.session); const requireSession = (): S => { - const session = store.get(address); - if (!session) throw new AppError('COMMAND_FAILED', 'Test session retired'); + const session = store.resolveCurrent(ref); + if (session === undefined) throw new AppError('COMMAND_FAILED', 'Test session retired'); return session; }; const assertAdoptable = (): void => { @@ -25,23 +55,32 @@ export function makeCaptureSessionBinding { - const session = store.get(address); - return session === undefined ? undefined : slot.read(session); + const session = store.resolveCurrent(ref); + if (session !== undefined) retained = slot.read(session); + return retained; }, assertAdoptable, canPersist: () => { - const session = store.get(address); + const session = store.resolveCurrent(ref); return session !== undefined && slot.read(session) === undefined; }, adopt: (resource) => { assertAdoptable(); - store.set(address, slot.replace(requireSession(), resource)); + store.update(ref, (current) => slot.replace(current, resource)); + retained = resource; }, clear: (expected) => { - const current = store.get(address); - if (!current) return 'retired'; - if (slot.read(current)?.handle !== expected.handle) return 'resource-changed'; - store.set(address, slot.replace(current, undefined)); + const current = store.resolveCurrent(ref); + if (current === undefined) return 'retired'; + const active = slot.read(current); + if ( + active?.handle !== expected.handle || + active.envelope.fence.token !== expected.envelope.fence.token || + active.envelope.fence.generation !== expected.envelope.fence.generation + ) + return 'resource-changed'; + store.update(ref, (session) => slot.replace(session, undefined)); + retained = undefined; return 'cleared'; }, }); diff --git a/packages/capture-kit/src/durable-capture/transitions.test.ts b/packages/capture-kit/src/durable-capture/transitions.test.ts index f691674c53..0ff01c1062 100644 --- a/packages/capture-kit/src/durable-capture/transitions.test.ts +++ b/packages/capture-kit/src/durable-capture/transitions.test.ts @@ -219,3 +219,67 @@ test('a disposal finish disposes a preserving kind’s material too', async () = envelope: { lifecycle: 'completed', metadata: { phase: 'completed' } }, }); }); + +test.each(['rebuild', 'retire', 'token', 'generation'] as const)( + 'a held finish after %s clears only its matching lifetime, handle and fence', + async (change) => { + const context = makeDurableCaptureContext(); + const start = makeDurableCaptureStartResult(context); + let enter!: () => void; + let release!: () => void; + const entered = new Promise((resolve) => { + enter = resolve; + }); + const resumed = new Promise((resolve) => { + release = resolve; + }); + start.finish.mockImplementationOnce(async () => { + enter(); + await resumed; + return { status: 'completed', result: { outputPath: '/tmp/capture', completedAt: 2 } }; + }); + await adoptStartedDurableCapture( + testCaptureDefinition, + { + ...context, + ...start, + throwIfCanceled: () => {}, + }, + context.resourcePath, + ); + const finishing = finishLiveDurableCapture( + testCaptureDefinition, + { + binding: context.binding, + intent: 'capture', + }, + context.resourcePath, + ); + await entered; + const ref = context.sessionStore.lookup(context.sessionName); + const active = context.sessionStore.get(context.sessionName)!.capture!; + if (change === 'retire') { + context.sessionStore.retire(ref); + context.sessionStore.set(context.sessionName, { name: 'successor', capture: active }); + } else { + const fence = { + ...active.envelope.fence, + ...(change === 'token' ? { token: 'replacement' } : {}), + ...(change === 'generation' ? { generation: active.envelope.fence.generation + 1 } : {}), + }; + context.sessionStore.update(ref, (current) => ({ + ...current, + name: 'updated', + capture: { ...active, envelope: { ...active.envelope, fence } }, + })); + } + const before = context.sessionStore.get(context.sessionName)!; + release(); + await finishing; + const current = context.sessionStore.get(context.sessionName)!; + expect(start.finish).toHaveBeenCalledOnce(); + expect(current.name).toBe(change === 'retire' ? 'successor' : 'updated'); + if (change === 'rebuild') expect(current.capture).toBeUndefined(); + else expect(current).toBe(before); + }, +); diff --git a/packages/host-kit/src/internal/owner-identity-liveness.test.ts b/packages/host-kit/src/internal/owner-identity-liveness.test.ts index 772d29cac5..511da64644 100644 --- a/packages/host-kit/src/internal/owner-identity-liveness.test.ts +++ b/packages/host-kit/src/internal/owner-identity-liveness.test.ts @@ -66,4 +66,6 @@ test('a pid outside the native range is unknown without a liveness probe', () => 'unknown', ); assert.equal(mockIsProcessAlive.mock.calls.length, 0); + assert.equal(mockIsProcessZombie.mock.calls.length, 0); + assert.equal(mockReadProcessStartTime.mock.calls.length, 0); }); diff --git a/packages/host-kit/src/session-paths.test.ts b/packages/host-kit/src/session-paths.test.ts index be0c8d3116..0b443d3452 100644 --- a/packages/host-kit/src/session-paths.test.ts +++ b/packages/host-kit/src/session-paths.test.ts @@ -7,8 +7,7 @@ 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); + assert.equal(homePath, path.join(os.homedir(), 'flows', 'replay.ad')); const relativePath = expandSessionPath('workflows/replay.ad', '/tmp/agent-device-cwd'); assert.equal(relativePath, path.resolve('/tmp/agent-device-cwd', 'workflows/replay.ad')); diff --git a/scripts/layering/session-resource-ownership.test.ts b/scripts/layering/session-resource-ownership.test.ts index 4fd9fd7045..8a1453e6cb 100644 --- a/scripts/layering/session-resource-ownership.test.ts +++ b/scripts/layering/session-resource-ownership.test.ts @@ -27,16 +27,10 @@ test('session resources are constructed only by their durable domain owners', () `sessionStore.update(ref, { appLog: log, appLogFailure: undefined });`, ], [ - 'src/daemon/audio-probe-session-binding.ts', - `sessionStore.update(ref, { audioProbe: audio });`, - ], - [ - 'src/daemon/perf-capture-session-binding.ts', - `sessionStore.update(ref, { perfCapture: perf });`, - ], - [ - 'src/daemon/screen-recording-session-binding.ts', - `sessionStore.update(ref, { screenRecording: recording });`, + 'src/daemon/session-capture-binding.ts', + `sessionStore.update(ref, { audioProbe: audio }); + sessionStore.update(ref, { perfCapture: perf }); + sessionStore.update(ref, { screenRecording: recording });`, ], [ 'packages/capture-kit/src/capture-admission/audio-probe-session-resource.ts', diff --git a/scripts/layering/session-resource-ownership.ts b/scripts/layering/session-resource-ownership.ts index 24408f22b4..323d8af3be 100644 --- a/scripts/layering/session-resource-ownership.ts +++ b/scripts/layering/session-resource-ownership.ts @@ -1,4 +1,4 @@ -// Catches: a session resource field (appLog, appLogFailure, audioProbe, perfCapture) written +// Catches: a session resource field (appLog, appLogFailure, audioProbe, perfCapture, screenRecording) written // from outside its declared owner module — R7's session-state-ownership shape applied to the // narrower set of per-resource fields these session-scoped runtimes carry, where the same // aliasing hazard (get()/set() hand back and re-put the live reference) applies. @@ -27,15 +27,12 @@ const SCANNED_ROOTS = ['src/daemon/', 'packages/capture-kit/src/capture-admissio const RESOURCE_OWNERS: Readonly>> = { appLog: new Set(['src/daemon/app-log-session-resource.ts', 'src/daemon/session-state.ts']), appLogFailure: new Set(['src/daemon/app-log-session-resource.ts', 'src/daemon/session-state.ts']), - audioProbe: new Set(['src/daemon/audio-probe-session-binding.ts', 'src/daemon/session-state.ts']), + audioProbe: new Set(['src/daemon/session-capture-binding.ts', 'src/daemon/session-state.ts']), screenRecording: new Set([ - 'src/daemon/screen-recording-session-binding.ts', - 'src/daemon/session-state.ts', - ]), - perfCapture: new Set([ - 'src/daemon/perf-capture-session-binding.ts', + 'src/daemon/session-capture-binding.ts', 'src/daemon/session-state.ts', ]), + perfCapture: new Set(['src/daemon/session-capture-binding.ts', 'src/daemon/session-state.ts']), }; /** Durable session-resource records have one whole-record construction owner per domain. */ diff --git a/src/__tests__/test-utils/registered-daemon-fixture.test.ts b/src/__tests__/test-utils/registered-daemon-fixture.test.ts new file mode 100644 index 0000000000..e4e265f4af --- /dev/null +++ b/src/__tests__/test-utils/registered-daemon-fixture.test.ts @@ -0,0 +1,43 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import { test } from 'vitest'; +import { resolveDaemonPaths } from '../../daemon-resolution.ts'; +import { stopDaemonProcess } from '../../daemon-process.ts'; +import { mkdtempForTestSync } from './tmp-dir.ts'; +import { + spawnRegisteredDaemonFixture, + waitForRegisteredDaemonFixture, + finishRegisteredDaemonFixture, +} from './registered-daemon-fixture.ts'; + +test('a joined fixture exit refuses its remaining registration metadata', async () => { + const paths = resolveDaemonPaths(mkdtempForTestSync('daemon-fixture-exited-publication-')); + const child = spawnRegisteredDaemonFixture( + paths, + { + httpPort: 4210, + token: 'fixture', + version: 'test', + codeOrigin: 'checkout', + codeSignature: 'fixture', + }, + undefined, + ); + try { + const observed = await waitForRegisteredDaemonFixture(paths, child); + assert.equal( + ( + await stopDaemonProcess( + { pid: child.pid, startTime: observed.processStartTime ?? null }, + { mode: 'force', termTimeoutMs: 0, killTimeoutMs: 1_000 }, + ) + ).status, + 'exited', + ); + await child.exited; + assert.equal(fs.existsSync(paths.infoPath), true); + await assert.rejects(waitForRegisteredDaemonFixture(paths, child), /exited before publication/); + } finally { + await finishRegisteredDaemonFixture(paths.baseDir); + } +}); diff --git a/src/__tests__/test-utils/registered-daemon-fixture.ts b/src/__tests__/test-utils/registered-daemon-fixture.ts index 6c99455a4e..9131617410 100644 --- a/src/__tests__/test-utils/registered-daemon-fixture.ts +++ b/src/__tests__/test-utils/registered-daemon-fixture.ts @@ -85,6 +85,7 @@ export async function waitForRegisteredDaemonFixture( void child.exited.then((result) => { exit = result; }); + await Promise.resolve(); for (let attempt = 0; attempt < 400; attempt += 1) { if (exit) throw new Error(`Registered child exited before publication: ${JSON.stringify(exit)}`); 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 8788a8c887..60eb5a9366 100644 --- a/src/daemon-client/__tests__/daemon-client-startup-race.test.ts +++ b/src/daemon-client/__tests__/daemon-client-startup-race.test.ts @@ -212,23 +212,22 @@ test('a joined busy contender waits for a published winner to become ready witho fs.writeFileSync(deferred, 'wait'); const winner = spawnRegisteredDaemonFixture(paths, fields(http.port), { stdio: 'ignore' }); await awaitFile(path.join(paths.baseDir, 'registration-held')); - let joined = false; + let contenderExit: ExecDetachedExit | undefined; spawn.mockImplementation((_command, _args, options) => { const child = spawnRegisteredDaemonFixture(paths, fields(http.port), options); void child.exited.then((exit) => { - assert.equal(exit.exitCode, DAEMON_STARTUP_EXIT_CODES.busy); - joined = true; + contenderExit = exit; }); return child; }); pause.mockImplementation(async () => { - if (joined) fs.rmSync(deferred, { force: true }); + if (contenderExit) fs.rmSync(deferred, { force: true }); await actualRetry.sleep(10); }); const signal = vi.spyOn(process, 'kill'); try { assert.equal((await sendToDaemon(request(paths))).ok, true); - assert.equal(joined, true); + assert.equal(contenderExit?.exitCode, DAEMON_STARTUP_EXIT_CODES.busy); assert.ok(probes >= 12); assert.equal(spawn.mock.calls.length, 1); assert.equal(http.rpcRequests.length, 1); diff --git a/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts b/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts index c28d397d1d..88415d1a23 100644 --- a/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts +++ b/src/daemon-client/__tests__/daemon-client-timeout-route.test.ts @@ -311,89 +311,115 @@ function timeoutDiagnostic(paths: DaemonPaths): Record { return event.data; } +function assertForcedRetirement( + retirement: DaemonRetirementResult | undefined, + removalFails: boolean, +): void { + if (removalFails) { + assert.ok(retirement?.status === 'retained'); + assert.equal(retirement.reason, 'retirement-unconfirmed'); + assert.ok(retirement.termination?.status === 'exited'); + assert.equal(retirement.termination.mode, 'forced'); + } else { + assert.ok(retirement?.status === 'retired'); + assert.equal(retirement.termination.mode, 'forced'); + } +} + for (const transport of ['socket', 'http'] as const) { - test(`${transport} timeout waits for force retirement before reporting reset`, async (t) => { - if (await skipWhenLoopbackUnavailable(t)) return; - mockRunCmdSync.mockReturnValue({ exitCode: 1, stdout: '', stderr: '' }); - const endpoint = await (transport === 'socket' - ? startHangingSocketServer() - : startHangingHttpServer()); - const paths = resolveDaemonPaths(mkdtempForTestSync('agent-device-timeout-owner-')); - const child = spawnRegisteredDaemonFixture( - paths, - { - ...(transport === 'socket' ? { socketPort: endpoint.port } : { httpPort: endpoint.port }), - token: 'test-token', - version: 'test', - codeOrigin: 'checkout', - codeSignature: 'test', - }, - undefined, - ); - let exited = false; - void child.exited.then(() => { - exited = true; - }); - let killRequested!: () => void; - const requested = new Promise((resolve) => { - killRequested = resolve; - }); - const actualKill = process.kill.bind(process); - let settled = false; - let outcome: Promise | undefined; - const kill = vi.spyOn(process, 'kill').mockImplementation((pid, signal) => { - if (pid === child.pid && signal === 'SIGKILL') { - killRequested(); - return true; - } - return actualKill(pid, signal); - }); - try { - const info = await waitForRegisteredDaemonFixture(paths, child); - outcome = withDiagnosticsScope( - { debug: true, logPath: path.join(paths.baseDir, 'timeout-diagnostics.ndjson') }, - () => - sendRequest( - info, - { ...buildRequest(undefined), command: 'open' }, - transport, - paths, - TIMEOUT_MS, - ), - ).then( - () => assert.fail('hanging request unexpectedly succeeded'), - (error: unknown) => { - settled = true; - return error; + for (const removalFails of [false, true]) { + test(`${transport} timeout awaits force exit when metadata removal ${removalFails ? 'fails' : 'succeeds'}`, async (t) => { + if (await skipWhenLoopbackUnavailable(t)) return; + mockRunCmdSync.mockReturnValue({ exitCode: 1, stdout: '', stderr: '' }); + const endpoint = await (transport === 'socket' + ? startHangingSocketServer() + : startHangingHttpServer()); + const paths = resolveDaemonPaths(mkdtempForTestSync('agent-device-timeout-owner-')); + const child = spawnRegisteredDaemonFixture( + paths, + { + ...(transport === 'socket' ? { socketPort: endpoint.port } : { httpPort: endpoint.port }), + token: 'test-token', + version: 'test', + codeOrigin: 'checkout', + codeSignature: 'test', }, + undefined, ); - await waitForForceStop(requested); - await sleep(30); - assert.equal(actualKill(child.pid, 0), true); - assert.equal(settled, false, 'request must remain pending while the daemon is alive'); - kill.mockRestore(); - actualKill(child.pid, 'SIGKILL'); - await child.exited; - const error = await outcome; - assert.ok(error instanceof AppError); - assert.equal(normalizeError(error).details?.reason, 'daemon_transport_timeout'); - const retirement = error.details?.retirement as DaemonRetirementResult | undefined; - assert.ok(retirement?.status === 'retired'); - assert.equal(retirement.termination.mode, 'forced'); - assert.equal(fs.existsSync(paths.infoPath), false); - assert.equal(fs.existsSync(paths.lockPath), false); - assert.equal(fs.existsSync(paths.baseDir), true); - assert.equal(timeoutDiagnostic(paths).daemonPreservedAfterTimeout, false); - assert.equal(mockRunCmdSync.mock.calls.filter(([cmd]) => cmd === 'pkill').length, 3); - } finally { - kill.mockRestore(); - if (!exited) actualKill(child.pid, 'SIGKILL'); - await child.exited; - await outcome; - await finishRegisteredDaemonFixture(paths.baseDir); - await closeLoopbackServer(endpoint.server); - } - }); + let exited = false; + void child.exited.then(() => { + exited = true; + }); + let killRequested!: () => void; + const requested = new Promise((resolve) => { + killRequested = resolve; + }); + const actualKill = process.kill.bind(process); + let settled = false; + let outcome: Promise | undefined; + const kill = vi.spyOn(process, 'kill').mockImplementation((pid, signal) => { + if (pid === child.pid && signal === 'SIGKILL') { + killRequested(); + return true; + } + return actualKill(pid, signal); + }); + const actualUnlink = fs.unlinkSync.bind(fs); + const remove = vi.spyOn(fs, 'unlinkSync').mockImplementation((file) => { + if (removalFails && file === paths.infoPath) + throw Object.assign(new Error('retained registration control'), { code: 'EACCES' }); + actualUnlink(file); + }); + try { + const info = await waitForRegisteredDaemonFixture(paths, child); + outcome = withDiagnosticsScope( + { debug: true, logPath: path.join(paths.baseDir, 'timeout-diagnostics.ndjson') }, + () => + sendRequest( + info, + { ...buildRequest(undefined), command: 'open' }, + transport, + paths, + TIMEOUT_MS, + ), + ).then( + () => assert.fail('hanging request unexpectedly succeeded'), + (error: unknown) => { + settled = true; + return error; + }, + ); + await waitForForceStop(requested); + await sleep(30); + assert.equal(actualKill(child.pid, 0), true); + assert.equal(settled, false, 'request must remain pending while the daemon is alive'); + kill.mockRestore(); + actualKill(child.pid, 'SIGKILL'); + await child.exited; + const error = await outcome; + assert.ok(error instanceof AppError); + assert.equal(normalizeError(error).details?.reason, 'daemon_transport_timeout'); + assertForcedRetirement( + error.details?.retirement as DaemonRetirementResult | undefined, + removalFails, + ); + assert.equal(fs.existsSync(paths.infoPath), removalFails); + assert.equal(fs.existsSync(paths.lockPath), false); + assert.equal(fs.existsSync(paths.baseDir), true); + assert.equal(timeoutDiagnostic(paths).daemonPreservedAfterTimeout, false); + assert.equal(timeoutDiagnostic(paths).daemonPidForceKilled, true); + assert.equal(mockRunCmdSync.mock.calls.filter(([cmd]) => cmd === 'pkill').length, 3); + } finally { + remove.mockRestore(); + kill.mockRestore(); + if (!exited) actualKill(child.pid, 'SIGKILL'); + await child.exited; + await outcome; + await finishRegisteredDaemonFixture(paths.baseDir); + await closeLoopbackServer(endpoint.server); + } + }); + } } test('timeout retains a live registration without captured birth proof and reports that outcome', async (t) => { diff --git a/src/daemon-client/daemon-client-timeout.ts b/src/daemon-client/daemon-client-timeout.ts index f7132c4f87..12c87c0d93 100644 --- a/src/daemon-client/daemon-client-timeout.ts +++ b/src/daemon-client/daemon-client-timeout.ts @@ -73,7 +73,8 @@ export async function handleRequestTimeout( mode: 'force', }); } - const retained = retirement?.status === 'retained'; + const preserved = + retirement?.status === 'retained' && retirement.termination?.status !== 'exited'; // The HINT, unlike cleanup, may only name Apple-runner involvement on // evidence this call site actually has: an explicitly declared Apple // platform selector, or the cleanup itself having terminated a matching @@ -94,7 +95,7 @@ export async function handleRequestTimeout( daemonPidReset: retirement?.status === 'retired' ? info.pid : undefined, daemonPidForceKilled: resetDaemon ? daemonWasForceKilled(retirement) : undefined, daemonRetirement: retirement, - daemonPreservedAfterTimeout: retained || (!remote && !resetDaemon), + daemonPreservedAfterTimeout: preserved || (!remote && !resetDaemon), daemonBaseUrl: info.baseUrl, }, }); diff --git a/src/daemon/__tests__/app-log-session-resource.test.ts b/src/daemon/__tests__/app-log-session-resource.test.ts index 2a20ed6a0c..bf65e67c88 100644 --- a/src/daemon/__tests__/app-log-session-resource.test.ts +++ b/src/daemon/__tests__/app-log-session-resource.test.ts @@ -19,6 +19,28 @@ import { adoptStartedSessionAppLog, finishSessionAppLog } from '../app-log-sessi import { createNextAppLogFence } from '../app-log-start-preflight.ts'; import { appLogResourceStore } from '../app-log-resource-store.ts'; import type { SessionState } from '../session-state.ts'; +import { stopSessionAppLog } from '../session-teardown.ts'; + +test('teardown captures an adopted app log before its lazy import can cross retirement', async () => { + const context = makeContext(); + const runtime = makeStartResult(context); + await adoptStartedSessionAppLog({ + ...context, + ...runtime.result, + throwIfCanceled: () => {}, + }); + expect(context.ref.session.appLog).toBeUndefined(); + const stopping = stopSessionAppLog(context); + context.sessionStore.retire(context.ref); + const successor = context.sessionStore.publish(context.sessionName, { + ...context.session, + appName: 'successor', + }); + await stopping; + expect(runtime.forceCleanup).toHaveBeenCalledOnce(); + expect(context.sessionStore.requireCurrent(successor)).toBe(successor.session); + expect(context.sessionStore.requireCurrent(successor).appLog).toBeUndefined(); +}); test('start persists open recovery truth before adopting the live handle', async () => { const context = makeContext(); diff --git a/src/daemon/__tests__/perf-capture-session-resource.test.ts b/src/daemon/__tests__/perf-capture-session-resource.test.ts index ecd14eb954..7441e03774 100644 --- a/src/daemon/__tests__/perf-capture-session-resource.test.ts +++ b/src/daemon/__tests__/perf-capture-session-resource.test.ts @@ -1,4 +1,4 @@ -import { bindSessionPerfCapture } from '../perf-capture-session-binding.ts'; +import { bindSessionPerfCapture } from '../session-capture-binding.ts'; import { beforeEach, expect, test, vi } from 'vitest'; import { localRuntimeOwner } from '@agent-device/contracts/platform-runtime'; import { AppError } from '@agent-device/kernel/errors'; @@ -103,6 +103,7 @@ test('a perf stop whose pull failed re-collects the device-side trace the first }, ); + const binding = bindSessionPerfCapture(sessionStore, sessionStore.lookup(sessionName)!); const started = await startAndroidPerfCapture(device, localRuntimeOwner('android'), { sessionId: sessionName, appId: capture.packageName, @@ -113,7 +114,7 @@ test('a perf stop whose pull failed re-collects the device-side trace the first }); await adoptStartedPerfCapture({ admissionLedger: createPerfCaptureAdmissionLedger(), - binding: bindSessionPerfCapture(sessionStore, sessionStore.lookup(sessionName)!), + binding, device, owner: localRuntimeOwner('android'), fence, @@ -127,7 +128,7 @@ test('a perf stop whose pull failed re-collects the device-side trace the first const stop = () => finishLivePerfCapture({ intent: 'capture', - binding: bindSessionPerfCapture(sessionStore, sessionStore.lookup(sessionName)!), + binding, }); await expect(stop()).rejects.toBe(pullFailure); diff --git a/src/daemon/__tests__/session-capture-binding.test.ts b/src/daemon/__tests__/session-capture-binding.test.ts index 41f5322d5b..1058b0f354 100644 --- a/src/daemon/__tests__/session-capture-binding.test.ts +++ b/src/daemon/__tests__/session-capture-binding.test.ts @@ -1,7 +1,44 @@ import { expect, test } from 'vitest'; import { makeSessionStore } from '../../__tests__/test-utils/store-factory.ts'; import { makeRecordingSession } from './session-teardown.fixtures.ts'; -import { bindSessionScreenRecording } from '../screen-recording-session-binding.ts'; +import { bindSessionScreenRecording } from '../session-capture-binding.ts'; + +test('adoption owns a vacant slot and retains its adopted handle after retirement', () => { + const store = makeSessionStore(); + const recorded = makeRecordingSession({ + name: 'capture', + sessionStore: store, + finish: async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }), + }); + const resource = recorded.screenRecording!; + const ref = store.publish('capture', { ...recorded, screenRecording: undefined }); + const binding = bindSessionScreenRecording(store, ref); + expect(binding.canPersist()).toBe(true); + binding.adopt(resource); + expect(store.requireCurrent(ref).screenRecording).toBe(resource); + expect(binding.canPersist()).toBe(false); + expect(() => binding.adopt(resource)).toThrow( + expect.objectContaining({ + details: expect.objectContaining({ reason: 'session_resource_changed' }), + }), + ); + store.retire(ref); + const successor = store.publish('capture', { + ...recorded, + screenRecording: undefined, + appName: 'successor', + }); + expect(binding.read()).toBe(resource); + expect(binding.clear(resource)).toBe('retired'); + expect(binding.canPersist()).toBe(false); + expect(() => binding.adopt(resource)).toThrow( + expect.objectContaining({ + details: expect.objectContaining({ reason: 'session_lifetime_ended' }), + }), + ); + expect(store.requireCurrent(successor)).toBe(successor.session); + expect(store.requireCurrent(successor).screenRecording).toBeUndefined(); +}); test('clearing a capture refreshes a rebuilt record without losing its other changes', () => { const store = makeSessionStore(); @@ -39,13 +76,15 @@ test('clearing an older handle or fence leaves a replacement capture intact', () store.update(ref, { screenRecording: replacement }); expect(binding.clear(active)).toBe('resource-changed'); expect(binding.read()).toBe(replacement); - const newerFence = { - ...active, - envelope: { ...active.envelope, fence: { token: 'next', generation: 2 } }, - }; - store.update(ref, { screenRecording: newerFence }); - expect(binding.clear(active)).toBe('resource-changed'); - expect(binding.read()).toBe(newerFence); + for (const fence of [ + { ...active.envelope.fence, token: 'next' }, + { ...active.envelope.fence, generation: active.envelope.fence.generation + 1 }, + ]) { + const newerFence = { ...active, envelope: { ...active.envelope, fence } }; + store.update(ref, { screenRecording: newerFence }); + expect(binding.clear(active)).toBe('resource-changed'); + expect(binding.read()).toBe(newerFence); + } }); test('a retired binding retains its old resource but cannot write into the next lifetime', () => { diff --git a/src/daemon/__tests__/session-store-lifetime.test.ts b/src/daemon/__tests__/session-store-lifetime.test.ts index d3b7ae7e65..d0e6a84cf0 100644 --- a/src/daemon/__tests__/session-store-lifetime.test.ts +++ b/src/daemon/__tests__/session-store-lifetime.test.ts @@ -30,6 +30,10 @@ test('refs capture records while resolving rebuilds from the same lifetime', () assert.equal(store.resolveCurrent(ref), rebuilt); } assert.equal(store.lookup(ADDRESS)?.session, rebuilt); + const refreshed = store.refresh(initial); + assert.equal(refreshed.lifetime, initial.lifetime); + assert.equal(refreshed.session, rebuilt); + assert.equal(initial.session, session); assert.equal(store.get('default'), undefined); }); @@ -56,6 +60,7 @@ test('address reuse with the same record still starts a different lifetime', () store.setRuntimeHints(ADDRESS, { metroPort: 8082 }); assert.notEqual(successor.lifetime, old.lifetime); assert.equal(store.resolveCurrent(old), undefined); + assert.equal(store.refresh(old), old); assert.throws(() => store.requireCurrent(old), ended); assert.throws(() => store.update(old, { appName: 'Stale' }), ended); assert.equal(store.retire(old), false); diff --git a/src/daemon/app-log-session-resource.ts b/src/daemon/app-log-session-resource.ts index 89335eebcb..1099c5811e 100644 --- a/src/daemon/app-log-session-resource.ts +++ b/src/daemon/app-log-session-resource.ts @@ -137,7 +137,7 @@ export function clearSessionAppLogFailure(params: { }); } -function bindSessionAppLog(sessionStore: SessionStore, ref: SessionRef) { +export function bindSessionAppLog(sessionStore: SessionStore, ref: SessionRef) { return bindSessionCapture(sessionStore, ref, { read: (session) => session.appLog, write: (appLog) => { diff --git a/src/daemon/audio-probe-session-binding.ts b/src/daemon/audio-probe-session-binding.ts deleted file mode 100644 index b4eecb43c7..0000000000 --- a/src/daemon/audio-probe-session-binding.ts +++ /dev/null @@ -1,12 +0,0 @@ -import { bindSessionCapture } from './session-capture-binding.ts'; -import type { SessionRef } from './session-state.ts'; -import type { SessionStore } from './session-store.ts'; - -export function bindSessionAudioProbe(sessionStore: SessionStore, ref: SessionRef) { - return bindSessionCapture(sessionStore, ref, { - read: (session) => session.audioProbe, - write: (audioProbe) => { - sessionStore.update(ref, { audioProbe }); - }, - }); -} diff --git a/src/daemon/handlers/record-runtime.ts b/src/daemon/handlers/record-runtime.ts index ad8cfc5b16..12cf18cad1 100644 --- a/src/daemon/handlers/record-runtime.ts +++ b/src/daemon/handlers/record-runtime.ts @@ -4,6 +4,7 @@ import type { ScreenRecordingCompletion, ScreenRecordingStartInput, } from '@agent-device/contracts/screen-recording-runtime'; +import { bindSessionScreenRecording } from '../session-capture-binding.ts'; import { resolveScreenRecordingRuntimePlan, screenRecordingAdmissionUse, @@ -31,10 +32,7 @@ import type { SessionStore } from '../session-store.ts'; import type { BindDeviceRuntime, BindExactDeviceRuntime } from '../request-runtime-binding.ts'; import type { DaemonRequest, DaemonResponse } from '../daemon-request.ts'; import type { SessionRef, SessionState } from '../session-state.ts'; -import { - bindRecordOnlyScreenRecording, - bindSessionScreenRecording, -} from '../screen-recording-session-binding.ts'; +import { bindRecordOnlyScreenRecording } from '../screen-recording-session-binding.ts'; import { recordSessionAction } from '../session-action-recorder.ts'; import { missingAppSessionResponse, @@ -155,6 +153,7 @@ async function startRecording( if (!startFact.available) return buildRecordingUnsupportedResponse(startFact); const runtime = await params.bindDevice(session.device, use); const { fence, outputPaths } = prepareRecordingStart(params, session); + binding.assertAdoptable(); const started = await runtime.operations.screenRecordingStart( screenRecordingStartInput(params, session, prepared, fence, outputPaths.outputPath), ); diff --git a/src/daemon/perf-capture-session-binding.ts b/src/daemon/perf-capture-session-binding.ts deleted file mode 100644 index 53463cd8ab..0000000000 --- a/src/daemon/perf-capture-session-binding.ts +++ /dev/null @@ -1,12 +0,0 @@ -import { bindSessionCapture } from './session-capture-binding.ts'; -import type { SessionRef } from './session-state.ts'; -import type { SessionStore } from './session-store.ts'; - -export function bindSessionPerfCapture(sessionStore: SessionStore, ref: SessionRef) { - return bindSessionCapture(sessionStore, ref, { - read: (session) => session.perfCapture, - write: (perfCapture) => { - sessionStore.update(ref, { perfCapture }); - }, - }); -} diff --git a/src/daemon/screen-recording-session-binding.ts b/src/daemon/screen-recording-session-binding.ts index d032f819ba..b5d56813d0 100644 --- a/src/daemon/screen-recording-session-binding.ts +++ b/src/daemon/screen-recording-session-binding.ts @@ -1,16 +1,7 @@ -import { bindSessionCapture } from './session-capture-binding.ts'; +import { bindSessionScreenRecording } from './session-capture-binding.ts'; import type { SessionRef } from './session-state.ts'; import type { SessionStore } from './session-store.ts'; -export function bindSessionScreenRecording(sessionStore: SessionStore, ref: SessionRef) { - return bindSessionCapture(sessionStore, ref, { - read: (session) => session.screenRecording, - write: (screenRecording) => { - sessionStore.update(ref, { screenRecording }); - }, - }); -} - export function bindRecordOnlyScreenRecording( sessionStore: SessionStore, address: string, diff --git a/src/daemon/session-capture-binding.ts b/src/daemon/session-capture-binding.ts index c3f809cf37..97555f0b93 100644 --- a/src/daemon/session-capture-binding.ts +++ b/src/daemon/session-capture-binding.ts @@ -14,6 +14,7 @@ export function bindSessionCapture( write(resource: DurableCaptureSessionResource | undefined): void; }>, ): DurableCaptureSessionBinding { + let retained = slot.read(sessionStore.resolveCurrent(ref) ?? ref.session); const assertAdoptable = (): void => { sessionStore.assertAdmissionOpen(ref.address); if (slot.read(sessionStore.requireCurrent(ref))) { @@ -26,7 +27,11 @@ export function bindSessionCapture( return Object.freeze({ address: ref.address, sessionDir: sessionStore.resolveSessionDir(ref.address), - read: () => slot.read(sessionStore.resolveCurrent(ref) ?? ref.session), + read: () => { + const current = sessionStore.resolveCurrent(ref); + if (current) retained = slot.read(current); + return retained; + }, assertAdoptable, canPersist: () => { const current = sessionStore.resolveCurrent(ref); @@ -35,6 +40,7 @@ export function bindSessionCapture( adopt: (resource) => { assertAdoptable(); slot.write(resource); + retained = resource; }, clear: (expected) => { const current = sessionStore.resolveCurrent(ref); @@ -47,7 +53,35 @@ export function bindSessionCapture( ) return 'resource-changed'; slot.write(undefined); + retained = undefined; return 'cleared'; }, }); } + +export function bindSessionAudioProbe(sessionStore: SessionStore, ref: SessionRef) { + return bindSessionCapture(sessionStore, ref, { + read: (session) => session.audioProbe, + write: (audioProbe) => { + sessionStore.update(ref, { audioProbe }); + }, + }); +} + +export function bindSessionPerfCapture(sessionStore: SessionStore, ref: SessionRef) { + return bindSessionCapture(sessionStore, ref, { + read: (session) => session.perfCapture, + write: (perfCapture) => { + sessionStore.update(ref, { perfCapture }); + }, + }); +} + +export function bindSessionScreenRecording(sessionStore: SessionStore, ref: SessionRef) { + return bindSessionCapture(sessionStore, ref, { + read: (session) => session.screenRecording, + write: (screenRecording) => { + sessionStore.update(ref, { screenRecording }); + }, + }); +} diff --git a/src/daemon/session-observability/internal/__tests__/session-perf-runtime.test.ts b/src/daemon/session-observability/internal/__tests__/session-perf-runtime.test.ts index 6547e94231..ed08e60858 100644 --- a/src/daemon/session-observability/internal/__tests__/session-perf-runtime.test.ts +++ b/src/daemon/session-observability/internal/__tests__/session-perf-runtime.test.ts @@ -201,6 +201,42 @@ test('perf native capture is adopted durably and stop uses the live handle witho ); }); +test.each(['shutdown', 'retire'] as const)( + 'perf refuses native startup after %s during binding', + async (change) => { + const sessionStore = makeStore(); + const ref = sessionStore.lookup('android')!; + const start = vi.fn(); + const runtime = createPerfRuntime({ perfNativeCaptureStart: start }); + const bindDevice: BindDeviceRuntime = async (device, use) => { + const bound = await runtime.bindDevice(device, use); + if (change === 'shutdown') sessionStore.closeAdmission(); + else sessionStore.retire(ref); + return bound; + }; + const response = await handleSessionObservabilityCommands({ + req: { + token: 't', + session: 'android', + command: 'perf', + positionals: ['trace', 'start', 'xctrace'], + }, + sessionName: 'android', + sessionStore, + inspectFacts: runtime.inspectFacts, + bindDevice, + perfCaptureAdmissionLedger: createPerfCaptureAdmissionLedger(), + }); + assert.equal(response?.ok, false); + if (response && !response.ok) + assert.equal( + response.error.details?.reason, + change === 'shutdown' ? 'daemon_shutting_down' : 'session_lifetime_ended', + ); + assert.equal(start.mock.calls.length, 0); + }, +); + function makeStore() { const sessionStore = makeSessionStore('agent-device-perf-runtime-'); sessionStore.set('android', makeAndroidSession('android', { appBundleId: 'com.example.app' })); diff --git a/src/daemon/session-observability/internal/session-audio.ts b/src/daemon/session-observability/internal/session-audio.ts index 2d5292c1ea..9a947fa1ce 100644 --- a/src/daemon/session-observability/internal/session-audio.ts +++ b/src/daemon/session-observability/internal/session-audio.ts @@ -21,7 +21,7 @@ import type { import type { SessionStore } from '../../session-store.ts'; import type { DaemonRequest, DaemonResponse } from '../../daemon-request.ts'; import type { SessionRef } from '../../session-state.ts'; -import { bindSessionAudioProbe } from '../../audio-probe-session-binding.ts'; +import { bindSessionAudioProbe } from '../../session-capture-binding.ts'; import { type DaemonFailureResponse, errorResponse } from '@agent-device/kernel/contracts'; type AudioParams = { @@ -130,6 +130,7 @@ async function startAudioProbe( params.sessionStore.ensureSessionDir(params.sessionName), 'audio-probe.json', ); + binding.assertAdoptable(); const started = await runtime.operations.audioProbeStart({ sessionId: params.sessionName, statusPath, diff --git a/src/daemon/session-observability/internal/session-observability.ts b/src/daemon/session-observability/internal/session-observability.ts index 7d00c3cba8..8b4302a297 100644 --- a/src/daemon/session-observability/internal/session-observability.ts +++ b/src/daemon/session-observability/internal/session-observability.ts @@ -18,6 +18,7 @@ import { type PerfCaptureAdmissionLedger } from '@agent-device/capture-kit/perf- import { appLogResourceStore } from '../../app-log-resource-store.ts'; import { adoptStartedSessionAppLog, + bindSessionAppLog, clearSessionAppLogFailure, finishSessionAppLog, inspectSessionAppLog, @@ -375,6 +376,7 @@ async function startSessionAppLog( resourcePath, device: session.device, }); + bindSessionAppLog(sessionStore, params.ref).assertAdoptable(); const result = await start({ sessionId: sessionName, appBundleId, diff --git a/src/daemon/session-observability/internal/session-perf-runtime.ts b/src/daemon/session-observability/internal/session-perf-runtime.ts index 4c8b4142e0..f8761f6231 100644 --- a/src/daemon/session-observability/internal/session-perf-runtime.ts +++ b/src/daemon/session-observability/internal/session-perf-runtime.ts @@ -32,7 +32,7 @@ import type { import type { SessionStore } from '../../session-store.ts'; import type { DaemonRequest, DaemonResponse } from '../../daemon-request.ts'; import type { SessionRef, SessionState } from '../../session-state.ts'; -import { bindSessionPerfCapture } from '../../perf-capture-session-binding.ts'; +import { bindSessionPerfCapture } from '../../session-capture-binding.ts'; import { recordSessionAction } from '../../session-action-recorder.ts'; import { admitRuntimePlan, @@ -217,6 +217,8 @@ async function startPerfCapture( resourcePath, device: session.device, }); + const binding = bindSessionPerfCapture(params.sessionStore, params.ref); + binding.assertAdoptable(); const started = await runtime.operations.perfNativeCaptureStart({ sessionId: params.sessionName, appId: session.appBundleId, @@ -228,7 +230,7 @@ async function startPerfCapture( }); await adoptStartedPerfCapture({ admissionLedger: requirePerfCaptureAdmissionLedger(params), - binding: bindSessionPerfCapture(params.sessionStore, params.ref), + binding, device: session.device, owner: runtime.owner, fence, diff --git a/src/daemon/session-store.ts b/src/daemon/session-store.ts index 97e2d6dc4b..217b8b8dac 100644 --- a/src/daemon/session-store.ts +++ b/src/daemon/session-store.ts @@ -105,6 +105,11 @@ export class SessionStore { return entry === ref.lifetime ? entry.current : undefined; } + refresh(ref: SessionRef): SessionRef { + const entry = this.sessions.get(ref.address); + return entry === ref.lifetime ? this.captureRef(ref.address, entry) : ref; + } + requireCurrent(ref: SessionRef): SessionState { const session = this.resolveCurrent(ref); if (!session) { diff --git a/src/daemon/session-teardown.ts b/src/daemon/session-teardown.ts index a5a4639bd6..e9e684fbef 100644 --- a/src/daemon/session-teardown.ts +++ b/src/daemon/session-teardown.ts @@ -2,11 +2,12 @@ import { AppError } from '@agent-device/kernel/errors'; import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; import { cleanupRetainedMaterializedPathsForSession } from './materialized-path-registry.ts'; import type { SessionRef, SessionState } from './session-state.ts'; -import { bindSessionAudioProbe } from './audio-probe-session-binding.ts'; -import { bindSessionPerfCapture } from './perf-capture-session-binding.ts'; -import { bindSessionScreenRecording } from './screen-recording-session-binding.ts'; +import { + bindSessionAudioProbe, + bindSessionPerfCapture, + bindSessionScreenRecording, +} from './session-capture-binding.ts'; import type { SessionStore } from './session-store.ts'; -import { forceCleanupSessionAppLog } from './app-log-session-resource.ts'; import { finishLiveAudioProbe } from '@agent-device/capture-kit/audio-probe-session-resource'; import { finishLivePerfCapture } from '@agent-device/capture-kit/perf-capture-session-resource'; import { finishLiveScreenRecording } from '@agent-device/capture-kit/screen-recording-session-resource'; @@ -17,7 +18,10 @@ export async function stopSessionAppLog(params: { ref: SessionRef; sessionStore: SessionStore; }): Promise { - await forceCleanupSessionAppLog(params); + const ref = params.sessionStore.refresh(params.ref); + if (!ref.session.appLog) return; + const { forceCleanupSessionAppLog } = await import('./app-log-session-resource.ts'); + await forceCleanupSessionAppLog({ ...params, ref }); } export async function stopSessionPerfCapture(params: { diff --git a/test/integration/support/daemon-test-cleanup.test.ts b/test/integration/support/daemon-test-cleanup.test.ts index a83780c9d5..baba38e0a6 100644 --- a/test/integration/support/daemon-test-cleanup.test.ts +++ b/test/integration/support/daemon-test-cleanup.test.ts @@ -47,31 +47,99 @@ test.each(['string', 'out-of-range'])( }, ); -test('cleanup stops the registered replacement when its observation still names the exited original', async () => { - const paths = resolveDaemonPaths(mkdtempForTestSync('daemon-test-replaced-registration-')); +test.each([false, true])( + 'cleanup joins a registered successor when the old observation lacks birth proof: %s', + async (missingBirth) => { + const paths = resolveDaemonPaths(mkdtempForTestSync('daemon-test-replaced-registration-')); + const original = spawnRegisteredDaemonFixture(paths, fields, undefined); + try { + const observed = await waitForRegisteredDaemonFixture(paths, original); + const stopped = await stopDaemonProcess( + { pid: original.pid, startTime: observed.processStartTime ?? null }, + { mode: 'force', termTimeoutMs: 0, killTimeoutMs: 1_000 }, + ); + assert.equal(stopped.status, 'exited'); + await original.exited; + const replacement = spawnRegisteredDaemonFixture(paths, fields, undefined); + await waitForRegisteredDaemonFixture(paths, replacement); + await cleanupDaemonTestState(paths.baseDir, { + ...observed, + processStartTime: missingBirth ? undefined : observed.processStartTime, + }); + assert.ok( + !isProcessAlive(replacement.pid) || + readHostProcessIdentityObservations([replacement.pid]) + .get(replacement.pid) + ?.state.startsWith('Z'), + 'the registered replacement must be dead before cleanup returns', + ); + await replacement.exited; + assert.equal(isProcessAlive(replacement.pid), false); + assert.equal(fs.existsSync(paths.baseDir), false); + } finally { + await finishRegisteredDaemonFixture(paths.baseDir); + } + }, +); + +test.each(['invalid-json', 'ownerless'])( + 'cleanup joins its observed child and retains %s metadata', + async (kind) => { + const paths = resolveDaemonPaths(mkdtempForTestSync('daemon-test-corrupt-observed-')); + const child = spawnRegisteredDaemonFixture(paths, fields, undefined); + const warnings = vi.spyOn(console, 'warn').mockImplementation(() => {}); + try { + const observed = await waitForRegisteredDaemonFixture(paths, child); + const corrupt = kind === 'invalid-json' ? '{invalid' : '{"pid": "unknown"}'; + fs.writeFileSync(paths.infoPath, corrupt); + await cleanupDaemonTestState(paths.baseDir, observed); + assert.ok( + !isProcessAlive(child.pid) || + readHostProcessIdentityObservations([child.pid]).get(child.pid)?.state.startsWith('Z'), + 'the observed child must be terminated before cleanup returns', + ); + await child.exited; + assert.equal(isProcessAlive(child.pid), false); + assert.equal(fs.readFileSync(paths.infoPath, 'utf8'), corrupt); + assert.equal(warnings.mock.calls.length, 1); + } finally { + warnings.mockRestore(); + await finishRegisteredDaemonFixture(paths.baseDir); + } + }, +); + +test('cleanup retains an unpublished successor holding the registration lock', async () => { + const paths = resolveDaemonPaths(mkdtempForTestSync('daemon-test-unpublished-successor-')); const original = spawnRegisteredDaemonFixture(paths, fields, undefined); + const warnings = vi.spyOn(console, 'warn').mockImplementation(() => {}); try { const observed = await waitForRegisteredDaemonFixture(paths, original); - const stopped = await stopDaemonProcess( - { pid: original.pid, startTime: observed.processStartTime ?? null }, - { mode: 'force', termTimeoutMs: 0, killTimeoutMs: 1_000 }, + assert.equal( + ( + await stopDaemonProcess( + { pid: original.pid, startTime: observed.processStartTime ?? null }, + { mode: 'force', termTimeoutMs: 0, killTimeoutMs: 1_000 }, + ) + ).status, + 'exited', ); - assert.equal(stopped.status, 'exited'); await original.exited; - const replacement = spawnRegisteredDaemonFixture(paths, fields, undefined); - await waitForRegisteredDaemonFixture(paths, replacement); - await cleanupDaemonTestState(paths.baseDir, observed); - assert.ok( - !isProcessAlive(replacement.pid) || - readHostProcessIdentityObservations([replacement.pid]) - .get(replacement.pid) - ?.state.startsWith('Z'), - 'the registered replacement must be dead before cleanup returns', + fs.rmSync(paths.infoPath); + fs.rmSync(paths.baseDir + '/registration-held'); + fs.writeFileSync(paths.baseDir + '/defer-publication', 'wait'); + const successor = spawnRegisteredDaemonFixture(paths, fields, undefined); + await vi.waitFor( + () => assert.equal(fs.existsSync(paths.baseDir + '/registration-held'), true), + { timeout: 4_000, interval: 10 }, ); - await replacement.exited; - assert.equal(isProcessAlive(replacement.pid), false); - assert.equal(fs.existsSync(paths.baseDir), false); + await cleanupDaemonTestState(paths.baseDir, observed); + assert.equal(fs.existsSync(paths.baseDir), true); + assert.equal(fs.existsSync(paths.infoPath), false); + assert.equal(isProcessAlive(successor.pid), true); + assert.equal(warnings.mock.calls.length, 1); } finally { + warnings.mockRestore(); await finishRegisteredDaemonFixture(paths.baseDir); } }); diff --git a/test/integration/support/daemon-test-cleanup.ts b/test/integration/support/daemon-test-cleanup.ts index 848a90d450..aa5c64fe75 100644 --- a/test/integration/support/daemon-test-cleanup.ts +++ b/test/integration/support/daemon-test-cleanup.ts @@ -1,12 +1,17 @@ import fs from 'node:fs'; -import path from 'node:path'; import { normalizeError } from '@agent-device/kernel/errors'; +import { tryAcquireProcessLock } from '@agent-device/host-kit/file'; +import { + readCurrentOwnerIdentity, + ownerIdentityMatches, + type OwnerIdentity, +} from '@agent-device/host-kit/process'; import { stopDaemonProcess } from '../../../src/daemon-process.ts'; +import { resolveDaemonPaths } from '../../../src/daemon-resolution.ts'; import { readRegisteredDaemonIdentity, readRegisteredDaemonOwnership, } from '../../../src/daemon-registration.ts'; -import { ownerIdentityMatches, type OwnerIdentity } from '@agent-device/host-kit/process'; type TestDaemonIdentity = { pid: number; processStartTime?: string }; @@ -16,33 +21,55 @@ export async function cleanupDaemonTestState( observed: TestDaemonIdentity | null, ): Promise { try { - const identities = [ - observed ? { pid: observed.pid, startTime: observed.processStartTime ?? null } : null, - readIdentity(stateDir), - ].filter((identity): identity is OwnerIdentity => identity !== null); - if (identities.length === 0) throw new Error('No daemon lifetime was observed'); + const paths = resolveDaemonPaths(stateDir); + const identities: OwnerIdentity[] = observed + ? [{ pid: observed.pid, startTime: observed.processStartTime ?? null }] + : []; + let registrationFailure: unknown; + try { + const registered = readIdentity(paths.infoPath); + if (registered) identities.push(registered); + } catch (error) { + registrationFailure = error; + } + const confirmed: OwnerIdentity[] = []; + let retained = false; for (const identity of identities) { const termination = await stopDaemonProcess(identity, { mode: 'graceful', termTimeoutMs: 1_500, killTimeoutMs: 1_500, }); - if (termination.status !== 'exited') { - console.warn('Daemon test cleanup retained state:', stateDir, termination); - return; - } + if (termination.status === 'exited') confirmed.push(identity); + else if (termination.status === 'retained') retained = true; + } + if (registrationFailure) throw registrationFailure; + if (retained || confirmed.length === 0) + throw new Error('Daemon termination could not be confirmed'); + const attempt = tryAcquireProcessLock({ + lockDirPath: paths.lockPath, + owner: { ...readCurrentOwnerIdentity(), acquiredAtMs: Date.now() }, + description: 'daemon test cleanup', + }); + if (attempt.status !== 'acquired') + throw new Error('Daemon registration is held or unproven during test cleanup'); + const { acquisition } = attempt; + try { + acquisition.assertHeld(); + const current = readIdentity(paths.infoPath); + if (current && !confirmed.some((identity) => ownerIdentityMatches(identity, current))) + throw new Error('Daemon registration changed during test cleanup'); + acquisition.assertHeld(); + fs.rmSync(stateDir, { recursive: true, force: true }); + } finally { + await acquisition.release(); } - const current = readIdentity(stateDir); - if (current && !identities.some((identity) => ownerIdentityMatches(identity, current))) - throw new Error('Daemon registration changed during test cleanup'); - fs.rmSync(stateDir, { recursive: true, force: true }); } catch (error) { console.warn('Daemon test cleanup retained state:', stateDir, normalizeError(error)); } } -function readIdentity(stateDir: string): OwnerIdentity | null { - const infoPath = path.join(stateDir, 'daemon.json'); +function readIdentity(infoPath: string): OwnerIdentity | null { if (readRegisteredDaemonOwnership(infoPath, null).state === 'absent') return null; const identity = readRegisteredDaemonIdentity(infoPath); if (!identity) throw new Error('Daemon registration identity is invalid or unreadable'); diff --git a/vitest.config.ts b/vitest.config.ts index bb9a0a0046..33237c9ecf 100644 --- a/vitest.config.ts +++ b/vitest.config.ts @@ -178,6 +178,7 @@ export default defineConfig({ // decisions over fixture state-dir listings, so they need no daemon, // device, or subprocess. 'test/integration/support/daemon-leak-model.test.ts', + // Cleanup uses real detached registration owners and process-exit proof; no device is needed. 'test/integration/support/daemon-test-cleanup.test.ts', // The Android failed-step evidence reader: it replays adb output through the probe // seam, so the crash/process/activity selectors need no emulator to be pinned.