diff --git a/docs/adr/0019-request-bound-platform-runtime.md b/docs/adr/0019-request-bound-platform-runtime.md index e88bf19e6..75ce8787f 100644 --- a/docs/adr/0019-request-bound-platform-runtime.md +++ b/docs/adr/0019-request-bound-platform-runtime.md @@ -304,6 +304,15 @@ contract handle beside its descriptor and metadata; the field has one R7 transit descriptor and neutral metadata enter the authoritative persisted recovery record. Concrete platform classes, provider clients, child handles, timers, transports, and wait promises enter neither store. +The daemon shares the lifecycle mechanics of those facet-specific resources through one +`DurableCaptureResource` coordinator. It owns bounded manifest I/O, per-resource fence serialization, +start/persist/adopt compensation, terminal transitions, and deadline-bounded exact-owner recovery. +The coordinator receives a facet's resource kind and manifest store, neutral session slot, completion +metadata projection, failure wording, and exact-owner recovery adapter; it does not select platforms, +interpret command flags, or form a generic runtime facet. App-log and screen recording retain distinct handles, +descriptors, facts, native finalization semantics, and admission policy. Reusable codec/live-handle +mechanics below the daemon remain in `@agent-device/capture-kit`, preserving the package direction. + A daemon-owned, process-lifetime admission ledger may retain bounded cleanup uncertainty that has no honest durable representation, but it is not a second live-resource store: it contains no handle or descriptor, never supersedes the persisted manifest, and is keyed by canonical device identity when diff --git a/src/daemon/__tests__/app-log-resource-store.test.ts b/src/daemon/__tests__/app-log-resource-store.test.ts index 358bf35b8..f452af557 100644 --- a/src/daemon/__tests__/app-log-resource-store.test.ts +++ b/src/daemon/__tests__/app-log-resource-store.test.ts @@ -39,6 +39,25 @@ test.each([ expect(fs.readFileSync(resourcePath, 'utf8')).toBe(body); }); +test('malformed app-log records preserve the underlying decoder message', () => { + const resourcePath = resolveAppLogResourcePath( + path.join(mkdtempForTestSync('app-log-malformed-message-'), 'session'), + ); + fs.mkdirSync(path.dirname(resourcePath), { recursive: true }); + fs.writeFileSync(resourcePath, '{'); + let expectedMessage = ''; + try { + JSON.parse('{'); + } catch (error) { + expectedMessage = error instanceof Error ? error.message : ''; + } + + expect(readAppLogResourceRecord(resourcePath)).toMatchObject({ + status: 'unreattachable', + message: expectedMessage, + }); +}); + test('app-log resource listing is deterministic and ignores unrelated artifacts', () => { const sessionsDir = mkdtempForTestSync('app-log-record-list-'); for (const sessionName of ['zeta', 'alpha']) { diff --git a/src/daemon/__tests__/durable-capture-admission-ledger.test.ts b/src/daemon/__tests__/durable-capture-admission-ledger.test.ts new file mode 100644 index 000000000..20d21dca4 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-admission-ledger.test.ts @@ -0,0 +1,54 @@ +import { expect, test } from 'vitest'; +import type { DeviceInfo } from '@agent-device/kernel/device'; +import { createDurableCaptureAdmissionLedger } from '../durable-capture-admission-ledger.ts'; + +const device: DeviceInfo = { + platform: 'android', + id: 'ledger-device', + name: 'Pixel', + kind: 'emulator', +}; +const otherDevice: DeviceInfo = { + platform: 'android', + id: 'other-device', + name: 'Other Pixel', + kind: 'emulator', +}; + +test('blocks only the affected device in this process-lifetime ledger', () => { + const ledger = createDurableCaptureAdmissionLedger({ displayName: 'Screen recording' }); + ledger.blockUndurableCleanup(device, 'cleanup could not be confirmed'); + expect(() => ledger.assertStartAllowed(device)).toThrow(/process-local/); + expect(() => ledger.assertStartAllowed(otherDevice)).not.toThrow(); + expect(() => + createDurableCaptureAdmissionLedger({ displayName: 'Screen recording' }).assertStartAllowed( + device, + ), + ).not.toThrow(); +}); + +test('expires bounded cleanup uncertainty and reports it once', () => { + let now = 1_000; + const expired: string[] = []; + const ledger = createDurableCaptureAdmissionLedger({ + displayName: 'Screen recording', + now: () => now, + undurableCleanupTtlMs: 500, + onUndurableCleanupExpired: ({ reason }) => expired.push(reason), + }); + ledger.blockUndurableCleanup(device, 'cleanup could not be confirmed'); + expect(() => ledger.assertStartAllowed(device)).toThrow(/process-local/); + now += 501; + expect(() => ledger.assertStartAllowed(device)).not.toThrow(); + expect(expired).toEqual(['cleanup could not be confirmed']); + expect(() => ledger.assertStartAllowed(device)).not.toThrow(); +}); + +test('clearing a device releases only its recorded uncertainty', () => { + const ledger = createDurableCaptureAdmissionLedger({ displayName: 'Screen recording' }); + ledger.blockUndurableCleanup(device, 'first'); + ledger.blockUndurableCleanup(otherDevice, 'second'); + ledger.clearUndurableCleanup(device); + expect(() => ledger.assertStartAllowed(device)).not.toThrow(); + expect(() => ledger.assertStartAllowed(otherDevice)).toThrow(/process-local/); +}); diff --git a/src/daemon/__tests__/durable-capture-recovery-authority.test.ts b/src/daemon/__tests__/durable-capture-recovery-authority.test.ts new file mode 100644 index 000000000..98f971863 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-recovery-authority.test.ts @@ -0,0 +1,73 @@ +import { expect, test, vi } from 'vitest'; +import { createDurableResourceEnvelope } from '@agent-device/capture-kit'; +import { localRuntimeOwner, type AppLogLiveHandle } from '@agent-device/contracts/platform'; +import { createTestAppLogLiveHandle } from '../../__tests__/test-utils/app-log-live-handle.ts'; +import { + acquireDurableCaptureRecoveryAuthorityBeforeDeadline, + DurableCaptureRecoveryDeadlineError, +} from '../durable-capture-recovery-authority.ts'; + +const envelope = createDurableResourceEnvelope({ + resourceKind: 'app-log', + sessionId: 'session', + device: { id: 'emulator-5554', family: 'android', kind: 'emulator' }, + owner: localRuntimeOwner('android'), + fence: { token: 'fence', generation: 1 }, + lifecycle: 'open', + descriptor: { version: 1, body: {} }, +}); + +test('deadline abort disposes authority that becomes active after the caller has timed out', async () => { + vi.useFakeTimers(); + try { + let resolveReattach!: (outcome: { status: 'active'; handle: AppLogLiveHandle }) => void; + const forceCleanup = vi.fn(async () => ({ status: 'cleaned' }) as const); + const disposeControl = vi.fn(async () => {}); + const cleanupFailures = vi.fn(); + const acquisition = acquireDurableCaptureRecoveryAuthorityBeforeDeadline({ + displayName: 'app-log', + envelope, + scope: { + signal: new AbortController().signal, + diagnostics: { emit: () => {} }, + progress: { report: () => {} }, + }, + deadlineMs: 25, + acquireControl: async () => ({ + reattach: async () => + await new Promise<{ status: 'active'; handle: AppLogLiveHandle }>((resolve) => { + resolveReattach = resolve; + }), + cleanup: async () => ({ status: 'already-missing' }), + [Symbol.asyncDispose]: disposeControl, + }), + onLateCleanupFailure: cleanupFailures, + }); + + const timedOut = expect(acquisition).rejects.toBeInstanceOf( + DurableCaptureRecoveryDeadlineError, + ); + await vi.advanceTimersByTimeAsync(25); + await timedOut; + + resolveReattach({ + status: 'active', + handle: createTestAppLogLiveHandle({ + inspect: () => ({ backend: 'android', state: 'active', startedAt: 1 }), + finish: async () => ({ + status: 'completed', + alreadyCompleted: true, + result: { backend: 'android', outputPath: '/tmp/app.log', completedAt: 1 }, + }), + forceCleanup, + }), + }); + await vi.advanceTimersByTimeAsync(0); + + expect(forceCleanup).toHaveBeenCalledOnce(); + expect(disposeControl).toHaveBeenCalledOnce(); + expect(cleanupFailures).not.toHaveBeenCalled(); + } finally { + vi.useRealTimers(); + } +}); diff --git a/src/daemon/__tests__/durable-capture-resource-adoption.test.ts b/src/daemon/__tests__/durable-capture-resource-adoption.test.ts new file mode 100644 index 000000000..290f6b888 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-resource-adoption.test.ts @@ -0,0 +1,57 @@ +import fs from 'node:fs'; +import path from 'node:path'; +import { expect, test } from 'vitest'; +import { AppError } from '@agent-device/kernel/errors'; +import { + makeDurableCaptureContext, + makeDurableCaptureStartResult, + testCaptureResource, + testCaptureStore, +} from './durable-capture-resource.fixtures.ts'; + +test('canceled adoption cleans the pending handle before terminalizing its manifest', async () => { + const context = makeDurableCaptureContext(); + const start = makeDurableCaptureStartResult(context); + const cancellation = new AppError('CANCELED', 'request canceled'); + + await expect( + testCaptureResource.adoptStarted({ + ...context, + ...start, + throwIfCanceled: () => { + throw cancellation; + }, + }), + ).rejects.toBe(cancellation); + expect(start.forceCleanup).toHaveBeenCalledOnce(); + expect(testCaptureStore.read(context.resourcePath)).toMatchObject({ + status: 'decoded', + envelope: { lifecycle: 'completed', metadata: { phase: 'completed' } }, + }); +}); + +test('a failed terminal transition preserves the primary error and blocks replacement', async () => { + const context = makeDurableCaptureContext(); + const start = makeDurableCaptureStartResult(context, { + cleanup: { status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }, + }); + const primary = new AppError('CANCELED', 'request canceled'); + const resourceDir = path.dirname(context.resourcePath); + + try { + await expect( + testCaptureResource.adoptStarted({ + ...context, + ...start, + throwIfCanceled: () => { + fs.chmodSync(resourceDir, 0o500); + throw primary; + }, + }), + ).rejects.toBe(primary); + } finally { + fs.chmodSync(resourceDir, 0o700); + } + + expect(() => context.admissionLedger.assertStartAllowed(context.device)).toThrow(/process-local/); +}); diff --git a/src/daemon/__tests__/app-log-resource-fence.test.ts b/src/daemon/__tests__/durable-capture-resource-fence.test.ts similarity index 62% rename from src/daemon/__tests__/app-log-resource-fence.test.ts rename to src/daemon/__tests__/durable-capture-resource-fence.test.ts index 2e2f60f82..edba2e525 100644 --- a/src/daemon/__tests__/app-log-resource-fence.test.ts +++ b/src/daemon/__tests__/durable-capture-resource-fence.test.ts @@ -3,18 +3,21 @@ import { expect, test, vi } from 'vitest'; import { localRuntimeOwner } from '@agent-device/contracts/platform'; import { createDurableResourceEnvelope } from '@agent-device/capture-kit'; import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; -import { withAppLogResourceFence } from '../app-log-resource-fence.ts'; -import { - readAppLogResourceRecord, - resolveAppLogResourcePath, - writeAppLogResourceRecord, -} from '../app-log-resource-store.ts'; +import { withDurableCaptureResourceFence } from '../durable-capture-resource-fence.ts'; +import { createDurableCaptureResourceStore } from '../durable-capture-resource-store.ts'; -test('ownership fence rejects a stale token before the native side effect', async () => { +const store = createDurableCaptureResourceStore({ + resourceKind: 'screen-recording', + fileName: 'screen-recording.resource.json', + displayName: 'Screen recording', +}); + +test('rejects a stale fence before its side effect', async () => { const resourcePath = makeRecord(); const sideEffect = vi.fn(async () => {}); await expect( - withAppLogResourceFence({ + withDurableCaptureResourceFence({ + store, resourcePath, expected: { token: 'stale', generation: 1 }, run: sideEffect, @@ -23,29 +26,31 @@ test('ownership fence rejects a stale token before the native side effect', asyn expect(sideEffect).not.toHaveBeenCalled(); }); -test('ownership fence serializes validation, side effect, and persisted transition', async () => { +test('serializes validation, native work, and transition for one resource path', async () => { const resourcePath = makeRecord(); const order: string[] = []; - let releaseFirst!: () => void; - let markFirstStarted!: () => void; + let release!: () => void; + let started!: () => void; const firstStarted = new Promise((resolve) => { - markFirstStarted = resolve; + started = resolve; }); const firstReleased = new Promise((resolve) => { - releaseFirst = resolve; + release = resolve; }); - const first = withAppLogResourceFence({ + const first = withDurableCaptureResourceFence({ + store, resourcePath, expected: { token: 'current', generation: 1 }, run: async (lease) => { order.push('first-start'); - markFirstStarted(); + started(); await firstReleased; lease.transition('open', { metadata: { phase: 'completing' } }); order.push('first-end'); }, }); - const second = withAppLogResourceFence({ + const second = withDurableCaptureResourceFence({ + store, resourcePath, expected: { token: 'current', generation: 1 }, run: async () => { @@ -54,23 +59,23 @@ test('ownership fence serializes validation, side effect, and persisted transiti }); await firstStarted; expect(order).toEqual(['first-start']); - releaseFirst(); + release(); await Promise.all([first, second]); expect(order).toEqual(['first-start', 'first-end', 'second']); - expect(readAppLogResourceRecord(resourcePath)).toMatchObject({ + expect(store.read(resourcePath)).toMatchObject({ status: 'decoded', - envelope: { lifecycle: 'open', metadata: { phase: 'completing' } }, + envelope: { metadata: { phase: 'completing' } }, }); }); function makeRecord(): string { - const resourcePath = resolveAppLogResourcePath( - path.join(mkdtempForTestSync('app-log-fence-'), 'session'), + const resourcePath = store.resolvePath( + path.join(mkdtempForTestSync('capture-fence-'), 'session'), ); - writeAppLogResourceRecord( + store.write( resourcePath, createDurableResourceEnvelope({ - resourceKind: 'app-log', + resourceKind: 'screen-recording', sessionId: 'session', device: { id: 'emulator-5554', family: 'android', kind: 'emulator' }, owner: localRuntimeOwner('android'), diff --git a/src/daemon/__tests__/durable-capture-resource-recovery.test.ts b/src/daemon/__tests__/durable-capture-resource-recovery.test.ts new file mode 100644 index 000000000..8d6d5d530 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-resource-recovery.test.ts @@ -0,0 +1,51 @@ +import path from 'node:path'; +import { expect, test, vi } from 'vitest'; +import { createDurableResourceEnvelope } from '@agent-device/capture-kit'; +import { localRuntimeOwner } from '@agent-device/contracts/platform'; +import { makeSessionStore } from '../../__tests__/test-utils/store-factory.ts'; +import { appLogDurableResource } from '../app-log-session-resource.ts'; + +test('generic recovery terminalizes missing native authority through one exact control', async () => { + const sessionStore = makeSessionStore('durable-capture-recovery-'); + const sessionsDir = path.dirname(sessionStore.resolveSessionDir('session')); + const resourcePath = appLogDurableResource.store.resolvePath( + sessionStore.resolveSessionDir('session'), + ); + appLogDurableResource.store.write( + resourcePath, + createDurableResourceEnvelope({ + resourceKind: 'app-log', + sessionId: 'session', + device: { id: 'emulator-5554', family: 'android', kind: 'emulator' }, + owner: localRuntimeOwner('android'), + fence: { token: 'fence', generation: 1 }, + lifecycle: 'open', + descriptor: { version: 1, body: { pid: 123 } }, + }), + ); + const dispose = vi.fn(async () => {}); + + await expect( + appLogDurableResource.recoverAll({ + sessionsDir, + scope: { + signal: new AbortController().signal, + diagnostics: { emit: () => {} }, + progress: { report: () => {} }, + }, + acquireControl: async () => ({ + reattach: async () => ({ status: 'missing' }), + cleanup: async () => ({ status: 'already-missing' }), + [Symbol.asyncDispose]: dispose, + }), + }), + ).resolves.toEqual({ scanned: 1, recovered: 1, retained: 0 }); + expect(dispose).toHaveBeenCalledOnce(); + expect(appLogDurableResource.store.read(resourcePath)).toMatchObject({ + status: 'decoded', + envelope: { + lifecycle: 'completed', + metadata: { phase: 'completed', recoveryStatus: 'already-missing' }, + }, + }); +}); diff --git a/src/daemon/__tests__/durable-capture-resource-store.test.ts b/src/daemon/__tests__/durable-capture-resource-store.test.ts new file mode 100644 index 000000000..bd17dc714 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-resource-store.test.ts @@ -0,0 +1,68 @@ +import fs from 'node:fs'; +import path from 'node:path'; +import { expect, test } from 'vitest'; +import { localRuntimeOwner } from '@agent-device/contracts/platform'; +import { createDurableResourceEnvelope } from '@agent-device/capture-kit'; +import { mkdtempForTestSync } from '../../__tests__/test-utils/tmp-dir.ts'; +import { createDurableCaptureResourceStore } from '../durable-capture-resource-store.ts'; + +const store = createDurableCaptureResourceStore({ + resourceKind: 'screen-recording', + fileName: 'screen-recording.resource.json', + displayName: 'Screen recording', +}); + +test('publishes typed resource records atomically with owner-only permissions', () => { + const resourcePath = store.resolvePath( + path.join(mkdtempForTestSync('capture-store-'), 'session'), + ); + store.write(resourcePath, envelope('session')); + + expect(store.read(resourcePath)).toMatchObject({ + status: 'decoded', + envelope: { resourceKind: 'screen-recording', lifecycle: 'open' }, + }); + expect(fs.statSync(resourcePath).mode & 0o777).toBe(0o600); + expect(fs.readdirSync(path.dirname(resourcePath))).toEqual(['screen-recording.resource.json']); +}); + +test('retains invalid evidence and never follows or replaces a symbolic link', () => { + const root = mkdtempForTestSync('capture-store-link-'); + const outsidePath = path.join(root, 'outside.json'); + fs.writeFileSync(outsidePath, `${JSON.stringify(envelope('outside'))}\n`); + const resourcePath = store.resolvePath(path.join(root, 'sessions', 'session')); + fs.mkdirSync(path.dirname(resourcePath), { recursive: true }); + fs.symlinkSync(outsidePath, resourcePath); + + expect(store.read(resourcePath)).toMatchObject({ + status: 'unreattachable', + reason: 'descriptor-invalid', + }); + expect(() => store.write(resourcePath, envelope('replacement'))).toThrow(/symbolic link/); + expect(fs.lstatSync(resourcePath).isSymbolicLink()).toBe(true); +}); + +test('lists only this resource kind in deterministic session order', () => { + const sessionsDir = mkdtempForTestSync('capture-store-list-'); + for (const name of ['zeta', 'alpha']) { + const resourcePath = store.resolvePath(path.join(sessionsDir, name)); + store.write(resourcePath, envelope(name)); + } + fs.writeFileSync(path.join(sessionsDir, 'unrelated.txt'), 'ignored'); + expect(store.list(sessionsDir)).toEqual([ + store.resolvePath(path.join(sessionsDir, 'alpha')), + store.resolvePath(path.join(sessionsDir, 'zeta')), + ]); +}); + +function envelope(sessionId: string) { + return createDurableResourceEnvelope({ + resourceKind: 'screen-recording', + sessionId, + device: { id: 'emulator-5554', family: 'android', kind: 'emulator' }, + owner: localRuntimeOwner('android'), + fence: { token: `fence-${sessionId}`, generation: 1 }, + lifecycle: 'open', + descriptor: { version: 1, body: {} }, + }); +} diff --git a/src/daemon/__tests__/durable-capture-resource-transitions.test.ts b/src/daemon/__tests__/durable-capture-resource-transitions.test.ts new file mode 100644 index 000000000..730b4b6e7 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-resource-transitions.test.ts @@ -0,0 +1,34 @@ +import { expect, test } from 'vitest'; +import { + makeDurableCaptureContext, + makeDurableCaptureStartResult, + testCaptureResource, + testCaptureStore, +} from './durable-capture-resource.fixtures.ts'; + +test('an uncertain finish retains both the live slot and cleanup-pending durable truth', async () => { + const context = makeDurableCaptureContext(); + const start = makeDurableCaptureStartResult(context, { + finish: { status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }, + }); + await testCaptureResource.adoptStarted({ + ...context, + ...start, + throwIfCanceled: () => {}, + }); + const active = context.sessionStore.get(context.sessionName); + if (!active) throw new Error('Expected adopted test capture session'); + + await expect( + testCaptureResource.finishLive({ + session: active, + sessionName: context.sessionName, + sessionStore: context.sessionStore, + }), + ).rejects.toMatchObject({ details: { reason: 'cleanup-unconfirmed' } }); + expect(context.sessionStore.get(context.sessionName)?.appLog?.handle).toBe(start.handle); + expect(testCaptureStore.read(context.resourcePath)).toMatchObject({ + status: 'decoded', + envelope: { lifecycle: 'open', metadata: { phase: 'cleanup-pending' } }, + }); +}); diff --git a/src/daemon/__tests__/durable-capture-resource.fixtures.ts b/src/daemon/__tests__/durable-capture-resource.fixtures.ts new file mode 100644 index 000000000..80f52d29a --- /dev/null +++ b/src/daemon/__tests__/durable-capture-resource.fixtures.ts @@ -0,0 +1,120 @@ +import { vi } from 'vitest'; +import { + localRuntimeOwner, + type AppLogCompletion, + type AppLogLiveHandle, + type CleanupOutcome, + type FinishOutcome, +} from '@agent-device/contracts/platform'; +import { createAppLogStartResult, createDurableResourceEnvelope } from '@agent-device/capture-kit'; +import type { DeviceInfo } from '@agent-device/kernel/device'; +import { makeSessionStore } from '../../__tests__/test-utils/store-factory.ts'; +import { createTestAppLogLiveHandle } from '../../__tests__/test-utils/app-log-live-handle.ts'; +import { createDurableCaptureAdmissionLedger } from '../durable-capture-admission-ledger.ts'; +import { createDurableCaptureResource } from '../durable-capture-resource.ts'; +import { createDurableCaptureResourceStore } from '../durable-capture-resource-store.ts'; +import type { SessionState } from '../types.ts'; + +export const testCaptureStore = createDurableCaptureResourceStore({ + resourceKind: 'app-log', + fileName: 'test-capture.resource.json', + displayName: 'test capture', +}); + +export const testCaptureResource = createDurableCaptureResource< + 'app-log', + AppLogLiveHandle, + AppLogCompletion +>({ + resourceKind: 'app-log', + displayName: 'test capture', + store: testCaptureStore, + sessionSlot: { + read: (session) => session.appLog, + replace: (session, appLog) => ({ ...session, appLog, appLogFailure: undefined }), + }, + completionMetadata: (completion) => ({ + outputPath: completion.outputPath, + completedAt: completion.completedAt, + }), + messages: { + noActive: 'no test capture active', + cleanupPendingHint: 'Keep the test capture manifest for exact-owner recovery.', + }, +}); + +export function makeDurableCaptureContext( + device: DeviceInfo = { + platform: 'android', + id: 'emulator-5554', + name: 'Pixel', + kind: 'emulator', + }, +) { + const sessionStore = makeSessionStore('durable-capture-resource-'); + const sessionName = 'session'; + const session: SessionState = { + name: sessionName, + device, + createdAt: 1, + actions: [], + }; + sessionStore.set(sessionName, session); + return { + admissionLedger: createDurableCaptureAdmissionLedger({ displayName: 'test capture' }), + session, + sessionName, + sessionStore, + device, + owner: localRuntimeOwner(device.platform), + fence: { token: 'fence', generation: 1 } as const, + resourcePath: testCaptureStore.resolvePath(sessionStore.resolveSessionDir(sessionName)), + }; +} + +export function makeDurableCaptureStartResult( + context: ReturnType, + options: { + cleanup?: CleanupOutcome; + finish?: FinishOutcome; + } = {}, +) { + const forceCleanup = vi.fn(async () => options.cleanup ?? ({ status: 'cleaned' } as const)); + const finish = vi.fn( + async () => + options.finish ?? + ({ + status: 'completed', + result: { backend: 'android', outputPath: '/tmp/app.log', completedAt: 2 }, + } as const), + ); + const handle = createTestAppLogLiveHandle({ + inspect: () => ({ backend: 'android', state: 'active', startedAt: 1 }), + finish, + forceCleanup, + }); + const envelope = createDurableResourceEnvelope({ + resourceKind: 'app-log', + sessionId: context.sessionName, + device: { + id: context.device.id, + family: context.device.platform, + ...(context.device.appleOs === undefined ? {} : { appleOs: context.device.appleOs }), + kind: context.device.kind, + ...(context.device.target === undefined ? {} : { target: context.device.target }), + ...(context.device.iosPhysicalDeviceBackend === undefined + ? {} + : { iosPhysicalDeviceBackend: context.device.iosPhysicalDeviceBackend }), + }, + owner: context.owner, + fence: context.fence, + lifecycle: 'open', + descriptor: { version: 1, body: { pid: 123 } }, + }); + return { + forceCleanup, + finish, + handle, + ...createAppLogStartResult(handle, envelope), + }; +} diff --git a/src/daemon/__tests__/durable-capture-resource.test.ts b/src/daemon/__tests__/durable-capture-resource.test.ts new file mode 100644 index 000000000..7a49dbc90 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-resource.test.ts @@ -0,0 +1,29 @@ +import { expect, test } from 'vitest'; +import { + makeDurableCaptureContext, + makeDurableCaptureStartResult, + testCaptureResource, +} from './durable-capture-resource.fixtures.ts'; + +test('one coordinator exposes the typed manifest and all lifecycle entrypoints', async () => { + const context = makeDurableCaptureContext(); + const start = makeDurableCaptureStartResult(context); + + await testCaptureResource.adoptStarted({ + ...context, + ...start, + throwIfCanceled: () => {}, + }); + const active = context.sessionStore.get(context.sessionName); + expect(active?.appLog?.handle).toBe(start.handle); + if (!active) throw new Error('Expected adopted test capture session'); + + await expect( + testCaptureResource.finishLive({ + session: active, + sessionName: context.sessionName, + sessionStore: context.sessionStore, + }), + ).resolves.toMatchObject({ outputPath: '/tmp/app.log', completedAt: 2 }); + expect(context.sessionStore.get(context.sessionName)?.appLog).toBeUndefined(); +}); diff --git a/src/daemon/__tests__/durable-capture-start-preflight.test.ts b/src/daemon/__tests__/durable-capture-start-preflight.test.ts new file mode 100644 index 000000000..de73644a8 --- /dev/null +++ b/src/daemon/__tests__/durable-capture-start-preflight.test.ts @@ -0,0 +1,34 @@ +import { expect, test } from 'vitest'; +import { createDurableResourceEnvelope } from '@agent-device/capture-kit'; +import { + makeDurableCaptureContext, + testCaptureResource, + testCaptureStore, +} from './durable-capture-resource.fixtures.ts'; + +test('a nonterminal capture blocks a replacement for the same device in another session', () => { + const context = makeDurableCaptureContext(); + testCaptureStore.write( + context.resourcePath, + createDurableResourceEnvelope({ + resourceKind: 'app-log', + sessionId: context.sessionName, + device: { id: context.device.id, family: 'android', kind: 'emulator' }, + owner: context.owner, + fence: context.fence, + lifecycle: 'open', + descriptor: { version: 1, body: { pid: 123 } }, + }), + ); + const replacementPath = testCaptureStore.resolvePath( + context.sessionStore.resolveSessionDir('replacement'), + ); + + expect(() => + testCaptureResource.createNextFence({ + admissionLedger: context.admissionLedger, + resourcePath: replacementPath, + device: context.device, + }), + ).toThrow(/terminal state/); +}); diff --git a/src/daemon/app-log-admission-ledger.ts b/src/daemon/app-log-admission-ledger.ts index baf7b918f..7cd2f7f4b 100644 --- a/src/daemon/app-log-admission-ledger.ts +++ b/src/daemon/app-log-admission-ledger.ts @@ -6,8 +6,10 @@ import { type DeviceInfo, } from '@agent-device/kernel/device'; import { AppError } from '@agent-device/kernel/errors'; - -const DEFAULT_UNDURABLE_CLEANUP_TTL_MS = 5 * 60_000; +import { + createDurableCaptureAdmissionLedger, + type DurableCaptureAdmissionLedger, +} from './durable-capture-admission-ledger.ts'; export type RetainedLegacyAppLogMarker = Readonly<{ markerPath: string; @@ -22,12 +24,10 @@ export type AppLogAdmissionLedgerOptions = Readonly<{ }>; /** Process-lifetime admission evidence that is intentionally reset by daemon restart. */ -export type AppLogAdmissionLedger = Readonly<{ - retainLegacyMarkers(markers: readonly RetainedLegacyAppLogMarker[]): void; - blockUndurableCleanup(device: DeviceInfo, reason: string): void; - clearUndurableCleanup(device: DeviceInfo): void; - assertStartAllowed(device: DeviceInfo): void; -}>; +export type AppLogAdmissionLedger = DurableCaptureAdmissionLedger & + Readonly<{ + retainLegacyMarkers(markers: readonly RetainedLegacyAppLogMarker[]): void; + }>; function findRetainedLegacyMarker( retainedMarkers: Map, @@ -51,32 +51,21 @@ function findRetainedLegacyMarker( export function createAppLogAdmissionLedger( options: AppLogAdmissionLedgerOptions = {}, ): AppLogAdmissionLedger { - const now = options.now ?? Date.now; const markerExists = options.markerExists ?? fs.existsSync; - const ttlMs = options.undurableCleanupTtlMs ?? DEFAULT_UNDURABLE_CLEANUP_TTL_MS; - const undurableCleanupBlocks = new Map< - string, - Readonly<{ device: DeviceIdentity; reason: string; expiresAt: number }> - >(); + const durableLedger = createDurableCaptureAdmissionLedger({ + displayName: 'app-log', + now: options.now, + undurableCleanupTtlMs: options.undurableCleanupTtlMs, + onUndurableCleanupExpired: options.onUndurableCleanupExpired, + }); const retainedLegacyMarkers = new Map(); return Object.freeze({ + ...durableLedger, retainLegacyMarkers(markers: readonly RetainedLegacyAppLogMarker[]): void { for (const marker of markers) retainedLegacyMarkers.set(marker.markerPath, marker); }, - blockUndurableCleanup(device: DeviceInfo, reason: string): void { - const identity = deviceIdentity(device); - undurableCleanupBlocks.set(deviceIdentityKey(identity), { - device: identity, - reason, - expiresAt: now() + ttlMs, - }); - }, - clearUndurableCleanup(device: DeviceInfo): void { - undurableCleanupBlocks.delete(deviceIdentityKey(deviceIdentity(device))); - }, assertStartAllowed(device: DeviceInfo): void { - const identity = deviceIdentity(device); - const identityKey = deviceIdentityKey(identity); + const identityKey = deviceIdentityKey(deviceIdentity(device)); const retained = findRetainedLegacyMarker(retainedLegacyMarkers, identityKey, markerExists); if (retained) { throw new AppError( @@ -88,22 +77,7 @@ export function createAppLogAdmissionLedger( }, ); } - - const block = undurableCleanupBlocks.get(identityKey); - if (!block) return; - if (now() > block.expiresAt) { - undurableCleanupBlocks.delete(identityKey); - options.onUndurableCleanupExpired?.({ device: block.device, reason: block.reason }); - return; - } - throw new AppError( - 'COMMAND_FAILED', - 'The existing app-log resource has process-local unconfirmed ownership', - { - reason: 'cleanup-unconfirmed', - hint: `Do not start a replacement in this daemon process: ${block.reason}`, - }, - ); + durableLedger.assertStartAllowed(device); }, }); } diff --git a/src/daemon/app-log-resource-recovery.ts b/src/daemon/app-log-resource-recovery.ts index 299e644de..d03cb61be 100644 --- a/src/daemon/app-log-resource-recovery.ts +++ b/src/daemon/app-log-resource-recovery.ts @@ -1,381 +1,74 @@ -import path from 'node:path'; import { - isConfirmedCleanup, narrowDeviceBinding, runtimeUse, type AppLogCompletion, type AppLogLiveHandle, - type CleanupOutcome, type DeviceBinding, type DeviceRuntimeGateway, type DurableResourceEnvelope, - type PlatformRuntimeOperations, type PlatformRequestScope, - type ReattachOutcome, + type PlatformRuntimeOperations, } from '@agent-device/contracts/platform'; import type { DeviceInfo } from '@agent-device/kernel/device'; -import { emitDiagnostic } from '../utils/diagnostics.ts'; -import { withAppLogResourceFence } from './app-log-resource-fence.ts'; -import { - listAppLogResourcePaths, - readAppLogResourceRecord, - resolveAppLogResourcePath, -} from './app-log-resource-store.ts'; -import { safeSessionName } from './session-paths.ts'; +import { appLogDurableResource } from './app-log-session-resource.ts'; +import type { DurableCaptureRecoveryControl } from './durable-capture-recovery-authority.ts'; +import type { + DurableCaptureRecoveryDiagnostic, + DurableCaptureRecoverySummary, +} from './durable-capture-resource-recovery.ts'; const appLogRecoveryUse = runtimeUse()({ required: ['appLogReattach', 'appLogCleanup'], }); -const DEFAULT_APP_LOG_RECOVERY_DEADLINE_MS = 5_000; - -export type AppLogRecoverySummary = Readonly<{ - scanned: number; - recovered: number; - retained: number; -}>; - -export type AppLogRecoveryDiagnostic = Readonly<{ - phase: string; - resourcePath: string; - data: Readonly>; -}>; +export type AppLogRecoverySummary = DurableCaptureRecoverySummary; +export type AppLogRecoveryDiagnostic = DurableCaptureRecoveryDiagnostic; -type AppLogRecoveryParams = { +export function recoverAppLogResourcesAfterDaemonLock(params: { sessionsDir: string; gateway: DeviceRuntimeGateway; scope: PlatformRequestScope; perRecordDeadlineMs?: number; onDiagnostic?: (diagnostic: AppLogRecoveryDiagnostic) => void; -}; - -type AppLogRecoveryPathOutcome = 'ignored' | 'recovered' | 'retained'; - -/** - * Recovers only persisted app-log resources and never creates SessionState. - * The caller must own the daemon lock for the entire invocation. - */ -export async function recoverAppLogResourcesAfterDaemonLock( - params: AppLogRecoveryParams, -): Promise { - const paths = listAppLogResourcePaths(params.sessionsDir); - const outcomes: AppLogRecoveryPathOutcome[] = []; - for (const resourcePath of paths) { - outcomes.push(await recoverAppLogResourcePath(params, resourcePath)); - } - return { - scanned: paths.length, - recovered: outcomes.filter((outcome) => outcome === 'recovered').length, - retained: outcomes.filter((outcome) => outcome === 'retained').length, - }; -} - -async function recoverAppLogResourcePath( - params: AppLogRecoveryParams, - resourcePath: string, -): Promise { - const record = readAppLogResourceRecord(resourcePath); - if (record.status === 'unreattachable') { - emitRecoveryDiagnostic( - 'app_log_recovery_record_unreattachable', - resourcePath, - { reason: record.reason, message: record.message, version: record.version }, - params.onDiagnostic, - ); - return 'retained'; - } - if (record.status === 'missing') return 'ignored'; - const canonicalResourcePath = canonicalAppLogResourcePath( - params.sessionsDir, - record.envelope.sessionId, - ); - if (path.resolve(resourcePath) !== path.resolve(canonicalResourcePath)) { - emitRecoveryDiagnostic( - 'app_log_recovery_session_path_mismatch', - resourcePath, - { envelopeSessionId: record.envelope.sessionId, canonicalResourcePath }, - params.onDiagnostic, - ); - return 'retained'; - } - if (record.envelope.lifecycle === 'completed') return 'ignored'; - return await recoverDecodedAppLogResource(params, resourcePath, record.envelope); -} - -async function recoverDecodedAppLogResource( - params: AppLogRecoveryParams, - resourcePath: string, - envelope: DurableResourceEnvelope<'app-log'>, -): Promise { - try { - const recovered = await recoverOneBeforeDeadline({ - resourcePath, - envelope, - gateway: params.gateway, - scope: params.scope, - deadlineMs: params.perRecordDeadlineMs ?? DEFAULT_APP_LOG_RECOVERY_DEADLINE_MS, - onDiagnostic: params.onDiagnostic, - }); - return recovered ? 'recovered' : 'retained'; - } catch (error) { - emitRecoveryFailure(resourcePath, error, params.onDiagnostic); - return 'retained'; - } -} - -function canonicalAppLogResourcePath(sessionsDir: string, sessionId: string): string { - return resolveAppLogResourcePath(path.join(sessionsDir, safeSessionName(sessionId))); -} - -function emitRecoveryFailure( - resourcePath: string, - error: unknown, - onDiagnostic?: (diagnostic: AppLogRecoveryDiagnostic) => void, -): void { - emitRecoveryDiagnostic( - error instanceof AppLogRecoveryDeadlineError - ? 'app_log_recovery_timed_out' - : 'app_log_recovery_failed', - resourcePath, - { - error: error instanceof Error ? error.message : String(error), - ...(error instanceof AppLogRecoveryDeadlineError ? { deadlineMs: error.deadlineMs } : {}), - }, - onDiagnostic, - ); -} - -class AppLogRecoveryDeadlineError extends Error { - readonly deadlineMs: number; - - constructor(deadlineMs: number) { - super(`App-log recovery exceeded its ${deadlineMs}ms deadline`); - this.name = 'AppLogRecoveryDeadlineError'; - this.deadlineMs = deadlineMs; - } -} - -async function recoverOneBeforeDeadline(params: { - resourcePath: string; - envelope: DurableResourceEnvelope<'app-log'>; - gateway: DeviceRuntimeGateway; - scope: PlatformRequestScope; - deadlineMs: number; - onDiagnostic?: (diagnostic: AppLogRecoveryDiagnostic) => void; -}): Promise { - const deadlineController = new AbortController(); - const scope: PlatformRequestScope = { - ...params.scope, - signal: AbortSignal.any([params.scope.signal, deadlineController.signal]), - }; - let rejectDeadline: (error: AppLogRecoveryDeadlineError) => void = () => {}; - const deadline = new Promise((_resolve, reject) => { - rejectDeadline = reject; +}): Promise { + return appLogDurableResource.recoverAll({ + sessionsDir: params.sessionsDir, + scope: params.scope, + perRecordDeadlineMs: params.perRecordDeadlineMs, + onDiagnostic: params.onDiagnostic, + acquireControl: async (envelope, scope) => + await acquireAppLogRecoveryControl(params.gateway, envelope, scope), }); - const timer = setTimeout(() => { - const error = new AppLogRecoveryDeadlineError(params.deadlineMs); - deadlineController.abort(error); - rejectDeadline(error); - }, params.deadlineMs); - timer.unref?.(); - - const acquisition = acquireRecoveryAuthority({ ...params, scope }); - let acquired: AppLogRecoveryAuthority; - try { - acquired = await Promise.race([acquisition, deadline]); - } finally { - clearTimeout(timer); - // A runtime that ignores cancellation is quarantined here. Acquisition checks - // the aborted request scope before returning any authority to clean or persist. - void acquisition.catch(() => {}); - } - - // Once exact-owner authority is acquired, settle its bounded cleanup before the - // daemon lock may be released. The deadline intentionally no longer races this phase. - return settleRecoveryAuthority({ ...params, ...acquired }); } -type AppLogRecoveryAuthority = Readonly<{ - binding: DeviceBinding; - reattached: ReattachOutcome; -}>; - -async function acquireRecoveryAuthority(params: { - resourcePath: string; - envelope: DurableResourceEnvelope<'app-log'>; - gateway: DeviceRuntimeGateway; - scope: PlatformRequestScope; - onDiagnostic?: (diagnostic: AppLogRecoveryDiagnostic) => void; -}): Promise { +async function acquireAppLogRecoveryControl( + gateway: DeviceRuntimeGateway, + envelope: DurableResourceEnvelope<'app-log'>, + scope: PlatformRequestScope, +): Promise> { let binding: DeviceBinding | undefined; - let reattached: ReattachOutcome | undefined; try { - binding = await params.gateway.bind({ - device: deviceFromEnvelope(params.envelope), - intent: { - kind: 'exact-owner', - owner: params.envelope.owner, - fence: params.envelope.fence, - }, - scope: params.scope, + binding = await gateway.bind({ + device: deviceFromEnvelope(envelope), + intent: { kind: 'exact-owner', owner: envelope.owner, fence: envelope.fence }, + scope, }); - params.scope.signal.throwIfAborted(); + scope.signal.throwIfAborted(); const runtime = narrowDeviceBinding(binding, appLogRecoveryUse); - reattached = await runtime.operations.appLogReattach({ envelope: params.envelope }); - params.scope.signal.throwIfAborted(); - return { binding, reattached }; + const bound = binding; + return Object.freeze({ + reattach: async (resource: DurableResourceEnvelope<'app-log'>) => + await runtime.operations.appLogReattach({ envelope: resource }), + cleanup: async (resource: DurableResourceEnvelope<'app-log'>) => + await runtime.operations.appLogCleanup({ envelope: resource }), + [Symbol.asyncDispose]: async () => await bound[Symbol.asyncDispose](), + }); } catch (error) { - await disposeAbortedRecoveryAuthority(params, error, binding, reattached); + if (binding) await binding[Symbol.asyncDispose](); throw error; } } -async function disposeAbortedRecoveryAuthority( - params: Pick[0], 'resourcePath' | 'onDiagnostic'>, - primaryError: unknown, - binding: DeviceBinding | undefined, - reattached: ReattachOutcome | undefined, -): Promise { - if (reattached?.status === 'active') { - await disposeRecoveryValue( - reattached.handle, - 'app_log_recovery_late_handle_cleanup_failed', - params, - primaryError, - ); - } - if (binding) { - await disposeRecoveryValue( - binding, - 'app_log_recovery_late_binding_cleanup_failed', - params, - primaryError, - ); - } -} - -async function disposeRecoveryValue( - value: AsyncDisposable, - phase: string, - params: Pick[0], 'resourcePath' | 'onDiagnostic'>, - primaryError: unknown, -): Promise { - try { - await value[Symbol.asyncDispose](); - } catch (cleanupError) { - emitRecoveryDiagnostic( - phase, - params.resourcePath, - { - error: cleanupError instanceof Error ? cleanupError.message : String(cleanupError), - primaryError: primaryError instanceof Error ? primaryError.message : String(primaryError), - }, - params.onDiagnostic, - ); - } -} - -async function settleRecoveryAuthority(params: { - resourcePath: string; - envelope: DurableResourceEnvelope<'app-log'>; - binding: DeviceBinding; - reattached: ReattachOutcome; - onDiagnostic?: (diagnostic: AppLogRecoveryDiagnostic) => void; -}): Promise { - try { - const runtime = narrowDeviceBinding(params.binding, appLogRecoveryUse); - const { reattached } = params; - switch (reattached.status) { - case 'active': { - const cleanup = await withAppLogResourceFence({ - resourcePath: params.resourcePath, - expected: params.envelope.fence, - run: async (lease) => { - lease.transition('open', { - metadata: { ...(lease.envelope.metadata ?? {}), phase: 'completing' }, - }); - const outcome = await reattached.handle.forceCleanup(); - transitionCleanupOutcome(lease, outcome); - return outcome; - }, - }); - return isConfirmedCleanup(cleanup); - } - case 'completed': - await transitionRecoveredTerminal(params, { - backend: reattached.result.backend, - outputPath: reattached.result.outputPath, - completedAt: reattached.result.completedAt, - }); - return true; - case 'missing': - await transitionRecoveredTerminal(params, { recoveryStatus: 'already-missing' }); - return true; - case 'unreattachable': { - if ( - reattached.reason === 'descriptor-invalid' || - reattached.reason === 'descriptor-version-unsupported' || - reattached.reason === 'ownership-fence-lost' - ) { - emitRecoveryDiagnostic( - 'app_log_recovery_retained', - params.resourcePath, - { reason: reattached.reason, message: reattached.message }, - params.onDiagnostic, - ); - return false; - } - const cleanup = await runtime.operations.appLogCleanup({ envelope: params.envelope }); - await withAppLogResourceFence({ - resourcePath: params.resourcePath, - expected: params.envelope.fence, - run: async (lease) => transitionCleanupOutcome(lease, cleanup), - }); - return isConfirmedCleanup(cleanup); - } - } - } finally { - await params.binding[Symbol.asyncDispose](); - } -} - -async function transitionRecoveredTerminal( - params: { - resourcePath: string; - envelope: DurableResourceEnvelope<'app-log'>; - }, - metadata: Record, -): Promise { - await withAppLogResourceFence({ - resourcePath: params.resourcePath, - expected: params.envelope.fence, - run: async (lease) => { - lease.transition('completed', { - metadata: { ...(lease.envelope.metadata ?? {}), phase: 'completed', ...metadata }, - }); - }, - }); -} - -function transitionCleanupOutcome( - lease: Parameters[0]['run']>[0], - outcome: CleanupOutcome, -): void { - lease.transition(isConfirmedCleanup(outcome) ? 'completed' : 'open', { - metadata: { - ...(lease.envelope.metadata ?? {}), - phase: isConfirmedCleanup(outcome) ? 'completed' : 'cleanup-pending', - cleanupStatus: outcome.status, - ...(outcome.status === 'cleanup-pending' - ? { - cleanupPendingReason: outcome.reason, - ...(outcome.message ? { cleanupPendingMessage: outcome.message } : {}), - } - : {}), - }, - }); -} - function deviceFromEnvelope(envelope: DurableResourceEnvelope<'app-log'>): DeviceInfo { return { platform: envelope.device.family, @@ -389,20 +82,3 @@ function deviceFromEnvelope(envelope: DurableResourceEnvelope<'app-log'>): Devic : { iosPhysicalDeviceBackend: envelope.device.iosPhysicalDeviceBackend }), }; } - -function emitRecoveryDiagnostic( - phase: string, - resourcePath: string, - data: Record, - onDiagnostic?: (diagnostic: AppLogRecoveryDiagnostic) => void, -): void { - if (onDiagnostic) { - onDiagnostic({ phase, resourcePath, data }); - return; - } - emitDiagnostic({ - level: 'warn', - phase, - data: { resourcePath, ...data }, - }); -} diff --git a/src/daemon/app-log-resource-store.ts b/src/daemon/app-log-resource-store.ts index 4f1a0d1d0..29601f62f 100644 --- a/src/daemon/app-log-resource-store.ts +++ b/src/daemon/app-log-resource-store.ts @@ -1,158 +1,12 @@ -import crypto from 'node:crypto'; -import fs from 'node:fs'; -import path from 'node:path'; -import { - type DurableEnvelopeDecodeOutcome, - type DurableResourceEnvelope, -} from '@agent-device/contracts/platform'; -import { decodeDurableResourceEnvelope } from '@agent-device/capture-kit'; -import { openVerifiedFileForRead } from '../utils/verified-file.ts'; - -const APP_LOG_RESOURCE_FILENAME = 'app-log.resource.json'; - -export type AppLogResourceRecordRead = - | Readonly<{ status: 'missing' }> - | Readonly<{ status: 'decoded'; envelope: DurableResourceEnvelope<'app-log'> }> - | Readonly<{ - status: 'unreattachable'; - reason: 'descriptor-invalid' | 'descriptor-version-unsupported'; - message: string; - version?: number; - }>; - -export function resolveAppLogResourcePath(sessionDir: string): string { - return path.join(sessionDir, APP_LOG_RESOURCE_FILENAME); -} - -export function readAppLogResourceRecord(resourcePath: string): AppLogResourceRecordRead { - let fd: number | undefined; - let value: unknown; - try { - fd = openVerifiedFileForRead(resourcePath); - if (fd === undefined) return { status: 'missing' }; - value = JSON.parse(fs.readFileSync(fd, 'utf8')) as unknown; - } catch (error) { - if (isMissingFile(error)) return { status: 'missing' }; - return invalidResourceRecord( - error instanceof Error ? error.message : 'App-log resource record is invalid', - ); - } finally { - if (fd !== undefined) fs.closeSync(fd); - } - - const decoded = decodeDurableResourceEnvelope(value); - return narrowAppLogEnvelope(decoded); -} - -export function writeAppLogResourceRecord( - resourcePath: string, - envelope: DurableResourceEnvelope<'app-log'>, -): void { - const dir = path.dirname(resourcePath); - fs.mkdirSync(dir, { recursive: true }); - const temporaryPath = path.join( - dir, - `.${path.basename(resourcePath)}.${process.pid}.${crypto.randomUUID()}.tmp`, - ); - let fd: number | undefined; - try { - assertSafeResourceDestination(resourcePath); - fd = fs.openSync(temporaryPath, 'wx', 0o600); - fs.writeFileSync(fd, `${JSON.stringify(envelope)}\n`, 'utf8'); - fs.fsyncSync(fd); - fs.closeSync(fd); - fd = undefined; - assertSafeResourceDestination(resourcePath); - fs.renameSync(temporaryPath, resourcePath); - syncDirectoryBestEffort(dir); - } finally { - if (fd !== undefined) fs.closeSync(fd); - try { - fs.rmSync(temporaryPath, { force: true }); - } catch {} - } -} - -export function listAppLogResourcePaths(sessionsDir: string): string[] { - let entries: fs.Dirent[]; - try { - entries = fs.readdirSync(sessionsDir, { withFileTypes: true }); - } catch { - return []; - } - return entries - .filter((entry) => entry.isDirectory()) - .map((entry) => resolveAppLogResourcePath(path.join(sessionsDir, entry.name))) - .filter(pathEntryExistsWithoutFollowing) - .sort(); -} - -function assertSafeResourceDestination(resourcePath: string): void { - let destination: fs.Stats; - try { - destination = fs.lstatSync(resourcePath); - } catch (error) { - if (isMissingFile(error)) return; - throw error; - } - if (destination.isSymbolicLink()) { - throw new Error('Refusing to replace an app-log resource symbolic link'); - } - if (!destination.isFile()) { - throw new Error('Refusing to replace an app-log resource that is not a regular file'); - } -} - -function pathEntryExistsWithoutFollowing(resourcePath: string): boolean { - try { - fs.lstatSync(resourcePath); - return true; - } catch (error) { - return !isMissingFile(error); - } -} - -function invalidResourceRecord(message: string): AppLogResourceRecordRead { - return { - status: 'unreattachable', - reason: 'descriptor-invalid', - message, - }; -} - -function narrowAppLogEnvelope(decoded: DurableEnvelopeDecodeOutcome): AppLogResourceRecordRead { - if (decoded.status !== 'decoded') return decoded; - if (decoded.envelope.resourceKind !== 'app-log') { - return { - status: 'unreattachable', - reason: 'descriptor-invalid', - message: `Expected app-log resource record, received ${decoded.envelope.resourceKind}`, - }; - } - return { - status: 'decoded', - envelope: decoded.envelope as DurableResourceEnvelope<'app-log'>, - }; -} - -function syncDirectoryBestEffort(dir: string): void { - let fd: number | undefined; - try { - fd = fs.openSync(dir, 'r'); - fs.fsyncSync(fd); - } catch { - // Atomic rename is the correctness boundary. Directory fsync support differs - // by host filesystem, so durability hardening remains best effort here. - } finally { - if (fd !== undefined) fs.closeSync(fd); - } -} - -function isMissingFile(error: unknown): boolean { - return ( - error !== null && - typeof error === 'object' && - 'code' in error && - (error as { code?: unknown }).code === 'ENOENT' - ); -} +import { createDurableCaptureResourceStore } from './durable-capture-resource-store.ts'; + +export const appLogResourceStore = createDurableCaptureResourceStore({ + resourceKind: 'app-log', + fileName: 'app-log.resource.json', + displayName: 'App-log', +}); + +export const resolveAppLogResourcePath = appLogResourceStore.resolvePath; +export const readAppLogResourceRecord = appLogResourceStore.read; +export const writeAppLogResourceRecord = appLogResourceStore.write; +export const listAppLogResourcePaths = appLogResourceStore.list; diff --git a/src/daemon/app-log-session-resource.ts b/src/daemon/app-log-session-resource.ts index 334204f52..aac782b59 100644 --- a/src/daemon/app-log-session-resource.ts +++ b/src/daemon/app-log-session-resource.ts @@ -1,30 +1,19 @@ import type { LogBackend } from '@agent-device/contracts/observability'; -import { - createDurableResourceEnvelope, - decodeDurableResourceEnvelope, -} from '@agent-device/capture-kit'; -import { - isConfirmedCleanup, - runtimeOwnerKey, - type AppLogCompletion, - type AppLogLiveHandle, - type CleanupOutcome, - type DurableResourceEnvelope, - type PendingTransferGuard, - type ResourceOwnershipFence, - type RuntimeOwnerRef, +import type { + AppLogCompletion, + AppLogLiveHandle, + DurableResourceEnvelope, + PendingTransferGuard, + ResourceOwnershipFence, + RuntimeOwnerRef, } from '@agent-device/contracts/platform'; -import { deviceIdentity, sameDeviceIdentity, type DeviceInfo } from '@agent-device/kernel/device'; -import { AppError, normalizeError } from '@agent-device/kernel/errors'; -import { emitDiagnostic } from '../utils/diagnostics.ts'; +import type { DeviceInfo } from '@agent-device/kernel/device'; +import { normalizeError } from '@agent-device/kernel/errors'; import type { AppLogAdmissionLedger } from './app-log-admission-ledger.ts'; +import { createDurableCaptureResource } from './durable-capture-resource.ts'; +import { appLogResourceStore } from './app-log-resource-store.ts'; import type { SessionStore } from './session-store.ts'; import type { SessionState } from './types.ts'; -import { - withAppLogResourceFence, - type AppLogResourceFenceLease, -} from './app-log-resource-fence.ts'; -import { readAppLogResourceRecord, writeAppLogResourceRecord } from './app-log-resource-store.ts'; export type AppLogSessionSnapshot = Readonly<{ active: boolean; @@ -36,6 +25,30 @@ export type AppLogSessionSnapshot = Readonly<{ hint?: string; }>; +export const appLogDurableResource = createDurableCaptureResource< + 'app-log', + AppLogLiveHandle, + AppLogCompletion +>({ + resourceKind: 'app-log', + displayName: 'app-log', + store: appLogResourceStore, + sessionSlot: { + read: (session) => session.appLog, + replace: (session, appLog) => ({ ...session, appLog, appLogFailure: undefined }), + }, + completionMetadata: (completion) => ({ + backend: completion.backend, + outputPath: completion.outputPath, + completedAt: completion.completedAt, + }), + messages: { + noActive: 'no app log stream active', + cleanupPendingHint: + 'Keep app-log.resource.json and retry cleanup through its exact runtime owner.', + }, +}); + export function inspectSessionAppLog(session: SessionState): AppLogSessionSnapshot { if (session.appLog) { const snapshot = session.appLog.handle.inspect(); @@ -57,7 +70,7 @@ export function inspectSessionAppLog(session: SessionState): AppLogSessionSnapsh return { active: false, state: 'inactive' }; } -type AdoptStartedSessionAppLogParams = { +export function adoptStartedSessionAppLog(params: { admissionLedger: AppLogAdmissionLedger; session: SessionState; sessionName: string; @@ -69,188 +82,26 @@ type AdoptStartedSessionAppLogParams = { pendingHandle: PendingTransferGuard; envelope: DurableResourceEnvelope<'app-log'>; throwIfCanceled(): void; -}; - -type AppLogAdoptionState = - | { kind: 'pending' } - | { kind: 'persisted' } - | { kind: 'transferred'; handle: AppLogLiveHandle }; - -export async function adoptStartedSessionAppLog( - params: AdoptStartedSessionAppLogParams, -): Promise { - let state: AppLogAdoptionState = { kind: 'pending' }; - try { - const envelope = withAppLogPhase(validateStartedEnvelope(params), 'active'); - writeAppLogResourceRecord(params.resourcePath, envelope); - state = { kind: 'persisted' }; - params.throwIfCanceled(); - const handle = params.pendingHandle.transfer(); - state = { kind: 'transferred', handle }; - params.sessionStore.set(params.sessionName, { - ...params.session, - appLog: { handle, envelope }, - appLogFailure: undefined, - }); - } catch (error) { - await recoverFailedAppLogAdoption(params, state, error); - throw error; - } -} - -async function recoverFailedAppLogAdoption( - params: AdoptStartedSessionAppLogParams, - state: AppLogAdoptionState, - primaryError: unknown, -): Promise { - const persisted = state.kind === 'pending' ? persistIncoherentRuntimeEnvelope(params) : true; - let cleanupError = await disposeFailedAppLogAdoption(params.pendingHandle, state); - const transition = confirmFailedAdoptionTransition(params, persisted, cleanupError); - cleanupError = transition.cleanupError; - updateUndurableCleanupBlock( - params.admissionLedger, - params.device, - persisted, - cleanupError, - transition.confirmed, - ); - emitFailedAdoptionCleanupDiagnostic(params.sessionName, primaryError, cleanupError); -} - -async function disposeFailedAppLogAdoption( - pendingHandle: PendingTransferGuard, - state: AppLogAdoptionState, -): Promise { - try { - if (state.kind === 'transferred') await state.handle[Symbol.asyncDispose](); - else await pendingHandle[Symbol.asyncDispose](); - return undefined; - } catch (error) { - return error; - } -} - -function confirmFailedAdoptionTransition( - params: Pick, - persisted: boolean, - cleanupError: unknown | undefined, -): { confirmed: boolean; cleanupError: unknown | undefined } { - if (!persisted) return { confirmed: false, cleanupError }; - try { - return { - confirmed: markCleanupAfterFailedAdoption( - params.resourcePath, - params.fence, - cleanupError === undefined, - ), - cleanupError, - }; - } catch (transitionError) { - emitDiagnostic({ - level: 'error', - phase: 'app_log_pending_adoption_transition_failed', - data: { - session: params.sessionName, - transitionError: - transitionError instanceof Error ? transitionError.message : String(transitionError), - }, - }); - return { confirmed: false, cleanupError: cleanupError ?? transitionError }; - } -} - -function updateUndurableCleanupBlock( - ledger: AppLogAdmissionLedger, - device: DeviceInfo, - persisted: boolean, - cleanupError: unknown | undefined, - transitionConfirmed: boolean, -): void { - if ((!persisted && cleanupError === undefined) || transitionConfirmed) { - ledger.clearUndurableCleanup(device); - return; - } - ledger.blockUndurableCleanup( - device, - cleanupError instanceof Error - ? cleanupError.message - : 'The durable cleanup transition could not be confirmed', - ); -} - -function emitFailedAdoptionCleanupDiagnostic( - sessionName: string, - primaryError: unknown, - cleanupError: unknown | undefined, -): void { - if (cleanupError === undefined) return; - emitDiagnostic({ - level: 'error', - phase: 'app_log_pending_adoption_cleanup_failed', - data: { - session: sessionName, - primaryError: primaryError instanceof Error ? primaryError.message : String(primaryError), - cleanupError: cleanupError instanceof Error ? cleanupError.message : String(cleanupError), - }, - }); +}): Promise { + return appLogDurableResource.adoptStarted(params); } -function persistIncoherentRuntimeEnvelope(params: { +export function finishSessionAppLog(params: { + session: SessionState; sessionName: string; + sessionStore: SessionStore; resourcePath: string; - device: DeviceInfo; - owner: RuntimeOwnerRef; - fence: ResourceOwnershipFence; - envelope: DurableResourceEnvelope<'app-log'>; -}): boolean { - try { - const envelope = createExpectedRecoveryEnvelope(params, params.envelope.descriptor); - writeAppLogResourceRecord(params.resourcePath, envelope); - return true; - } catch (descriptorError) { - try { - const envelope = createExpectedRecoveryEnvelope(params, { - version: 0, - body: { reason: 'runtime-contract-invalid' }, - }); - writeAppLogResourceRecord(params.resourcePath, envelope); - return true; - } catch (persistenceError) { - emitDiagnostic({ - level: 'error', - phase: 'app_log_runtime_contract_tombstone_failed', - data: { - session: params.sessionName, - descriptorError: - descriptorError instanceof Error ? descriptorError.message : String(descriptorError), - persistenceError: - persistenceError instanceof Error ? persistenceError.message : String(persistenceError), - }, - }); - return false; - } - } +}): Promise { + return appLogDurableResource.finishLive(params); } -function createExpectedRecoveryEnvelope( - params: { - sessionName: string; - device: DeviceInfo; - owner: RuntimeOwnerRef; - fence: ResourceOwnershipFence; - }, - descriptor: DurableResourceEnvelope<'app-log'>['descriptor'], -): DurableResourceEnvelope<'app-log'> { - return createDurableResourceEnvelope({ - resourceKind: 'app-log', - sessionId: params.sessionName, - device: deviceIdentity(params.device), - owner: params.owner, - fence: params.fence, - lifecycle: 'open', - descriptor, - metadata: { phase: 'runtime-contract-invalid', runtimeContractInvalid: true }, - }); +export function forceCleanupSessionAppLog(params: { + session: SessionState; + sessionName?: string; + sessionStore?: SessionStore; + resourcePath: string; +}): Promise { + return appLogDurableResource.forceCleanupLive(params); } export function recordSessionAppLogFailure(params: { @@ -284,204 +135,3 @@ export function clearSessionAppLogFailure(params: { appLogFailure: undefined, }); } - -export async function finishSessionAppLog(params: { - session: SessionState; - sessionName: string; - sessionStore: SessionStore; - resourcePath: string; -}): Promise { - const resource = params.session.appLog; - if (!resource) { - throw new AppError('INVALID_ARGS', 'no app log stream active'); - } - const outcome = await withAppLogResourceFence({ - resourcePath: params.resourcePath, - expected: resource.envelope.fence, - run: async (lease) => { - markAppLogResourceCompleting(lease); - const result = await resource.handle.finish(); - if (result.status === 'completed') { - lease.transition('completed', { - metadata: { - ...(lease.envelope.metadata ?? {}), - backend: result.result.backend, - outputPath: result.result.outputPath, - completedAt: result.result.completedAt, - phase: 'completed', - }, - }); - } else { - lease.transition('open', { - metadata: { - ...(lease.envelope.metadata ?? {}), - phase: 'cleanup-pending', - cleanupPendingReason: result.reason, - ...(result.message ? { cleanupPendingMessage: result.message } : {}), - }, - }); - } - return result; - }, - }); - if (outcome.status === 'cleanup-pending') throw cleanupPendingError(outcome); - params.sessionStore.set(params.sessionName, { - ...params.session, - appLog: undefined, - appLogFailure: undefined, - }); - return outcome.result; -} - -/** Generic close/teardown cleanup; it never re-selects a platform implementation. */ -export async function forceCleanupSessionAppLog(params: { - session: SessionState; - sessionName?: string; - sessionStore?: SessionStore; - resourcePath: string; -}): Promise { - const resource = params.session.appLog; - if (!resource) return; - const outcome = await withAppLogResourceFence({ - resourcePath: params.resourcePath, - expected: resource.envelope.fence, - run: async (lease) => { - markAppLogResourceCompleting(lease); - const result = await resource.handle.forceCleanup(); - lease.transition(isConfirmedCleanup(result) ? 'completed' : 'open', { - metadata: { - ...(lease.envelope.metadata ?? {}), - phase: isConfirmedCleanup(result) ? 'completed' : 'cleanup-pending', - cleanupStatus: result.status, - ...(result.status === 'cleanup-pending' - ? { - cleanupPendingReason: result.reason, - ...(result.message ? { cleanupPendingMessage: result.message } : {}), - } - : {}), - }, - }); - return result; - }, - }); - if (!isConfirmedCleanup(outcome)) throw cleanupPendingError(outcome); - if (params.sessionStore && params.sessionName) { - params.sessionStore.set(params.sessionName, { - ...params.session, - appLog: undefined, - appLogFailure: undefined, - }); - } -} - -function markAppLogResourceCompleting(lease: AppLogResourceFenceLease): void { - lease.transition('open', { - metadata: { ...(lease.envelope.metadata ?? {}), phase: 'completing' }, - }); -} - -function validateStartedEnvelope(params: { - sessionName: string; - device: DeviceInfo; - owner: RuntimeOwnerRef; - fence: ResourceOwnershipFence; - envelope: DurableResourceEnvelope<'app-log'>; -}): DurableResourceEnvelope<'app-log'> { - const decoded = decodeDurableResourceEnvelope(params.envelope); - if (decoded.status !== 'decoded' || decoded.envelope.resourceKind !== 'app-log') { - throw invalidStartedEnvelope(); - } - const { envelope } = decoded; - if (!matchesStartedEnvelopeAuthority(envelope, params)) throw invalidStartedEnvelope(); - return createDurableResourceEnvelope({ - resourceKind: 'app-log', - sessionId: envelope.sessionId, - device: envelope.device, - owner: envelope.owner, - fence: envelope.fence, - lifecycle: envelope.lifecycle, - descriptor: envelope.descriptor, - ...(envelope.metadata === undefined ? {} : { metadata: envelope.metadata }), - }); -} - -function matchesStartedEnvelopeAuthority( - envelope: DurableResourceEnvelope, - expected: Pick, -): boolean { - return ( - envelope.sessionId === expected.sessionName && - sameDeviceIdentity(envelope.device, deviceIdentity(expected.device)) && - runtimeOwnerKey(envelope.owner) === runtimeOwnerKey(expected.owner) && - envelope.fence.token === expected.fence.token && - envelope.fence.generation === expected.fence.generation && - isStartedLifecycle(envelope.lifecycle) - ); -} - -function isStartedLifecycle(lifecycle: DurableResourceEnvelope['lifecycle']): boolean { - return lifecycle === 'open'; -} - -function invalidStartedEnvelope(): AppError { - return new AppError('COMMAND_FAILED', 'App-log runtime returned an incoherent durable envelope', { - reason: 'runtime-contract-invalid', - hint: 'Retain the runtime diagnostics and report the selected device and owner.', - }); -} - -function markCleanupAfterFailedAdoption( - resourcePath: string, - expected: ResourceOwnershipFence, - confirmed: boolean, -): boolean { - const record = readAppLogResourceRecord(resourcePath); - if ( - record.status !== 'decoded' || - record.envelope.fence.token !== expected.token || - record.envelope.fence.generation !== expected.generation - ) { - return false; - } - writeAppLogResourceRecord(resourcePath, { - ...record.envelope, - lifecycle: confirmed ? 'completed' : 'open', - metadata: { - ...(record.envelope.metadata ?? {}), - phase: confirmed ? 'completed' : 'cleanup-pending', - cleanupStatus: confirmed ? 'cleaned' : 'cleanup-pending', - }, - }); - return true; -} - -function withAppLogPhase( - envelope: DurableResourceEnvelope<'app-log'>, - phase: string, -): DurableResourceEnvelope<'app-log'> { - return Object.freeze({ - ...envelope, - lifecycle: 'open', - metadata: { ...(envelope.metadata ?? {}), phase }, - }); -} - -function cleanupPendingError( - outcome: - | Extract - | { - status: 'cleanup-pending'; - reason: string; - message?: string; - }, -): AppError { - return new AppError( - 'COMMAND_FAILED', - outcome.message ?? 'App-log cleanup could not be confirmed', - { - reason: outcome.reason, - retriable: outcome.reason !== 'ownership-fence-lost', - hint: 'Keep app-log.resource.json and retry cleanup through its exact runtime owner.', - }, - ); -} diff --git a/src/daemon/app-log-start-preflight.ts b/src/daemon/app-log-start-preflight.ts index f56cefdb7..6fc98de01 100644 --- a/src/daemon/app-log-start-preflight.ts +++ b/src/daemon/app-log-start-preflight.ts @@ -1,61 +1,16 @@ -import crypto from 'node:crypto'; -import path from 'node:path'; import type { ResourceOwnershipFence } from '@agent-device/contracts/platform'; -import { deviceIdentity, deviceIdentityKey, type DeviceInfo } from '@agent-device/kernel/device'; -import { AppError } from '@agent-device/kernel/errors'; +import type { DeviceInfo } from '@agent-device/kernel/device'; import type { AppLogAdmissionLedger } from './app-log-admission-ledger.ts'; -import { listAppLogResourcePaths, readAppLogResourceRecord } from './app-log-resource-store.ts'; +import { appLogDurableResource } from './app-log-session-resource.ts'; export function createNextAppLogFence(params: { ledger: AppLogAdmissionLedger; resourcePath: string; device: DeviceInfo; }): ResourceOwnershipFence { - const { ledger, resourcePath, device } = params; - ledger.assertStartAllowed(device); - const selectedDeviceKey = deviceIdentityKey(deviceIdentity(device)); - assertNoConflictingManifest(resourcePath, selectedDeviceKey); - - const record = readAppLogResourceRecord(resourcePath); - return Object.freeze({ - token: crypto.randomUUID(), - generation: record.status === 'decoded' ? record.envelope.fence.generation + 1 : 1, + return appLogDurableResource.createNextFence({ + admissionLedger: params.ledger, + resourcePath: params.resourcePath, + device: params.device, }); } - -function assertNoConflictingManifest(resourcePath: string, selectedDeviceKey: string): void { - const sessionsDir = path.dirname(path.dirname(resourcePath)); - for (const existingPath of listAppLogResourcePaths(sessionsDir)) { - const existing = readAppLogResourceRecord(existingPath); - if (existing.status === 'unreattachable') throw unreattachableManifest(existing); - if ( - existing.status === 'decoded' && - existing.envelope.lifecycle !== 'completed' && - (existingPath === resourcePath || - deviceIdentityKey(existing.envelope.device) === selectedDeviceKey) - ) { - throw new AppError( - 'COMMAND_FAILED', - 'An app-log resource for this device has not reached a confirmed terminal state', - { - reason: 'cleanup-unconfirmed', - hint: 'Retry exact-owner cleanup using the existing app-log.resource.json before starting a replacement.', - }, - ); - } - } -} - -function unreattachableManifest( - record: Extract, { status: 'unreattachable' }>, -): AppError { - return new AppError( - 'COMMAND_FAILED', - `An app-log recovery record is unreattachable: ${record.message}`, - { - reason: record.reason, - retriable: false, - hint: 'Retain the corrupt or future-version app-log.resource.json for manual recovery; no replacement capture is safe.', - }, - ); -} diff --git a/src/daemon/durable-capture-admission-ledger.ts b/src/daemon/durable-capture-admission-ledger.ts new file mode 100644 index 000000000..a6e176b0e --- /dev/null +++ b/src/daemon/durable-capture-admission-ledger.ts @@ -0,0 +1,64 @@ +import { + deviceIdentity, + deviceIdentityKey, + type DeviceIdentity, + type DeviceInfo, +} from '@agent-device/kernel/device'; +import { AppError } from '@agent-device/kernel/errors'; + +const DEFAULT_UNDURABLE_CLEANUP_TTL_MS = 5 * 60_000; + +export type DurableCaptureAdmissionLedger = Readonly<{ + blockUndurableCleanup(device: DeviceInfo, reason: string): void; + clearUndurableCleanup(device: DeviceInfo): void; + assertStartAllowed(device: DeviceInfo): void; +}>; + +export function createDurableCaptureAdmissionLedger( + options: Readonly<{ + displayName: string; + now?: () => number; + undurableCleanupTtlMs?: number; + onUndurableCleanupExpired?: ( + block: Readonly<{ device: DeviceIdentity; reason: string }>, + ) => void; + }>, +): DurableCaptureAdmissionLedger { + const now = options.now ?? Date.now; + const ttlMs = options.undurableCleanupTtlMs ?? DEFAULT_UNDURABLE_CLEANUP_TTL_MS; + const blocks = new Map< + string, + Readonly<{ device: DeviceIdentity; reason: string; expiresAt: number }> + >(); + return Object.freeze({ + blockUndurableCleanup(device: DeviceInfo, reason: string): void { + const identity = deviceIdentity(device); + blocks.set(deviceIdentityKey(identity), { + device: identity, + reason, + expiresAt: now() + ttlMs, + }); + }, + clearUndurableCleanup(device: DeviceInfo): void { + blocks.delete(deviceIdentityKey(deviceIdentity(device))); + }, + assertStartAllowed(device: DeviceInfo): void { + const identityKey = deviceIdentityKey(deviceIdentity(device)); + const block = blocks.get(identityKey); + if (!block) return; + if (now() > block.expiresAt) { + blocks.delete(identityKey); + options.onUndurableCleanupExpired?.({ device: block.device, reason: block.reason }); + return; + } + throw new AppError( + 'COMMAND_FAILED', + `The existing ${options.displayName.toLowerCase()} resource has process-local unconfirmed ownership`, + { + reason: 'cleanup-unconfirmed', + hint: `Do not start a replacement in this daemon process: ${block.reason}`, + }, + ); + }, + }); +} diff --git a/src/daemon/durable-capture-recovery-authority.ts b/src/daemon/durable-capture-recovery-authority.ts new file mode 100644 index 000000000..d0e26af8f --- /dev/null +++ b/src/daemon/durable-capture-recovery-authority.ts @@ -0,0 +1,123 @@ +import type { + CleanupOutcome, + DurableResourceEnvelope, + LiveResourceHandle, + PlatformRequestScope, + ReattachOutcome, +} from '@agent-device/contracts/platform'; + +export type DurableCaptureRecoveryControl< + K extends string, + H extends LiveResourceHandle, + C, +> = AsyncDisposable & + Readonly<{ + reattach(envelope: DurableResourceEnvelope): Promise>; + cleanup(envelope: DurableResourceEnvelope): Promise; + }>; + +export type DurableCaptureRecoveryAuthority< + K extends string, + H extends LiveResourceHandle, + C, +> = Readonly<{ + control: DurableCaptureRecoveryControl; + reattached: ReattachOutcome; +}>; + +export type DurableCaptureRecoveryAuthorityParams< + K extends string, + H extends LiveResourceHandle, + C, +> = Readonly<{ + displayName: string; + envelope: DurableResourceEnvelope; + scope: PlatformRequestScope; + deadlineMs: number; + acquireControl( + envelope: DurableResourceEnvelope, + scope: PlatformRequestScope, + ): Promise>; + onLateCleanupFailure(phase: string, cleanupError: unknown, primaryError: unknown): void; +}>; + +export async function acquireDurableCaptureRecoveryAuthorityBeforeDeadline< + K extends string, + H extends LiveResourceHandle, + C, +>( + params: DurableCaptureRecoveryAuthorityParams, +): Promise> { + const controller = new AbortController(); + const scope = { + ...params.scope, + signal: AbortSignal.any([params.scope.signal, controller.signal]), + }; + let rejectDeadline: (error: DurableCaptureRecoveryDeadlineError) => void = () => {}; + const deadline = new Promise((_resolve, reject) => { + rejectDeadline = reject; + }); + const timer = setTimeout(() => { + const error = new DurableCaptureRecoveryDeadlineError(params.displayName, params.deadlineMs); + controller.abort(error); + rejectDeadline(error); + }, params.deadlineMs); + timer.unref?.(); + const acquisition = acquireRecoveryAuthority(params, scope); + try { + return await Promise.race([acquisition, deadline]); + } finally { + clearTimeout(timer); + void acquisition.catch(() => {}); + } +} + +async function acquireRecoveryAuthority, C>( + params: DurableCaptureRecoveryAuthorityParams, + scope: PlatformRequestScope, +): Promise> { + let control: DurableCaptureRecoveryControl | undefined; + let reattached: ReattachOutcome | undefined; + try { + control = await params.acquireControl(params.envelope, scope); + scope.signal.throwIfAborted(); + reattached = await control.reattach(params.envelope); + scope.signal.throwIfAborted(); + return { control, reattached }; + } catch (error) { + if (reattached?.status === 'active') { + await disposeLateAuthority(params, reattached.handle, 'late_handle_cleanup_failed', error); + } + if (control) { + await disposeLateAuthority(params, control, 'late_control_cleanup_failed', error); + } + throw error; + } +} + +async function disposeLateAuthority, C>( + params: DurableCaptureRecoveryAuthorityParams, + value: AsyncDisposable, + phase: string, + primaryError: unknown, +): Promise { + try { + await value[Symbol.asyncDispose](); + } catch (cleanupError) { + params.onLateCleanupFailure(phase, cleanupError, primaryError); + } +} + +export class DurableCaptureRecoveryDeadlineError extends Error { + readonly deadlineMs: number; + + constructor(displayName: string, deadlineMs: number) { + super(`${capitalize(displayName)} recovery exceeded its ${deadlineMs}ms deadline`); + this.name = 'DurableCaptureRecoveryDeadlineError'; + this.deadlineMs = deadlineMs; + } +} + +function capitalize(value: string): string { + return value.length === 0 ? value : value[0]!.toUpperCase() + value.slice(1); +} diff --git a/src/daemon/durable-capture-resource-adoption.ts b/src/daemon/durable-capture-resource-adoption.ts new file mode 100644 index 000000000..e5996350a --- /dev/null +++ b/src/daemon/durable-capture-resource-adoption.ts @@ -0,0 +1,282 @@ +import { + runtimeOwnerKey, + type DurableResourceEnvelope, + type LiveResourceHandle, + type ResourceOwnershipFence, + type RuntimeOwnerRef, +} from '@agent-device/contracts/platform'; +import { + createDurableResourceEnvelope, + decodeDurableResourceEnvelope, +} from '@agent-device/capture-kit'; +import { deviceIdentity, sameDeviceIdentity, type DeviceInfo } from '@agent-device/kernel/device'; +import { AppError } from '@agent-device/kernel/errors'; +import { emitDiagnostic } from '../utils/diagnostics.ts'; +import type { + AdoptStartedDurableCaptureParams, + DurableCaptureResourceDefinition, +} from './durable-capture-resource.ts'; + +type AdoptionState = + | { kind: 'pending' } + | { kind: 'persisted' } + | { kind: 'transferred'; handle: H }; + +export async function adoptStartedDurableCapture< + K extends string, + H extends LiveResourceHandle, + C, +>( + definition: DurableCaptureResourceDefinition, + params: AdoptStartedDurableCaptureParams, + resourcePath: string, +): Promise { + let state: AdoptionState = { kind: 'pending' }; + try { + const envelope = withPhase(validateStartedEnvelope(definition, params), 'active'); + definition.store.write(resourcePath, envelope); + state = { kind: 'persisted' }; + params.throwIfCanceled(); + const handle = params.pendingHandle.transfer(); + state = { kind: 'transferred', handle }; + params.sessionStore.set( + params.sessionName, + definition.sessionSlot.replace(params.session, { handle, envelope }), + ); + } catch (error) { + await recoverFailedAdoption(definition, params, resourcePath, state, error); + throw error; + } +} + +async function recoverFailedAdoption, C>( + definition: DurableCaptureResourceDefinition, + params: AdoptStartedDurableCaptureParams, + resourcePath: string, + state: AdoptionState, + primaryError: unknown, +): Promise { + const persisted = + state.kind === 'pending' ? persistRecoveryTombstone(definition, params, resourcePath) : true; + const initialCleanupError = await disposeFailedAdoption(params, state); + const transition = confirmFailedAdoptionTransition( + definition, + params, + resourcePath, + initialCleanupError, + ); + if ((!persisted && transition.cleanupError === undefined) || transition.confirmed) { + params.admissionLedger.clearUndurableCleanup(params.device); + } else { + params.admissionLedger.blockUndurableCleanup( + params.device, + transition.cleanupError instanceof Error + ? transition.cleanupError.message + : 'The durable cleanup transition could not be confirmed', + ); + } + if (transition.cleanupError !== undefined) + emitCleanupDiagnostic(definition, params, primaryError, transition.cleanupError); +} + +function confirmFailedAdoptionTransition, C>( + definition: DurableCaptureResourceDefinition, + params: Pick, 'sessionName' | 'fence'>, + resourcePath: string, + cleanupError: unknown | undefined, +): { confirmed: boolean; cleanupError: unknown | undefined } { + try { + return { + confirmed: confirmFailedAdoption(definition, resourcePath, params.fence, cleanupError), + cleanupError, + }; + } catch (transitionError) { + emitDiagnostic({ + level: 'error', + phase: `${diagnosticPrefix(definition.resourceKind)}_pending_adoption_transition_failed`, + data: { + session: params.sessionName, + transitionError: + transitionError instanceof Error ? transitionError.message : String(transitionError), + }, + }); + return { confirmed: false, cleanupError: cleanupError ?? transitionError }; + } +} + +async function disposeFailedAdoption( + params: AdoptStartedDurableCaptureParams, + state: AdoptionState, +): Promise { + try { + if (state.kind === 'transferred') await state.handle[Symbol.asyncDispose](); + else await params.pendingHandle[Symbol.asyncDispose](); + return undefined; + } catch (error) { + return error; + } +} + +function persistRecoveryTombstone, C>( + definition: DurableCaptureResourceDefinition, + params: AdoptStartedDurableCaptureParams, + resourcePath: string, +): boolean { + try { + definition.store.write( + resourcePath, + createExpectedEnvelope(definition, params, params.envelope.descriptor), + ); + return true; + } catch (descriptorError) { + try { + definition.store.write( + resourcePath, + createExpectedEnvelope(definition, params, { + version: 0, + body: { reason: 'runtime-contract-invalid' }, + }), + ); + return true; + } catch (persistenceError) { + emitDiagnostic({ + level: 'error', + phase: `${diagnosticPrefix(definition.resourceKind)}_runtime_contract_tombstone_failed`, + data: { + session: params.sessionName, + descriptorError: + descriptorError instanceof Error ? descriptorError.message : String(descriptorError), + persistenceError: + persistenceError instanceof Error ? persistenceError.message : String(persistenceError), + }, + }); + return false; + } + } +} + +function createExpectedEnvelope, C>( + definition: DurableCaptureResourceDefinition, + params: Pick< + AdoptStartedDurableCaptureParams, + 'sessionName' | 'device' | 'owner' | 'fence' + >, + descriptor: DurableResourceEnvelope['descriptor'], +): DurableResourceEnvelope { + return createDurableResourceEnvelope({ + resourceKind: definition.resourceKind, + sessionId: params.sessionName, + device: deviceIdentity(params.device), + owner: params.owner, + fence: params.fence, + lifecycle: 'open', + descriptor, + metadata: { phase: 'runtime-contract-invalid', runtimeContractInvalid: true }, + }); +} + +function validateStartedEnvelope, C>( + definition: DurableCaptureResourceDefinition, + params: Pick< + AdoptStartedDurableCaptureParams, + 'sessionName' | 'device' | 'owner' | 'fence' | 'envelope' + >, +): DurableResourceEnvelope { + const decoded = decodeDurableResourceEnvelope(params.envelope); + if ( + decoded.status === 'decoded' && + matchesAuthority(decoded.envelope, params, definition.resourceKind) + ) { + return decoded.envelope as DurableResourceEnvelope; + } + throw new AppError( + 'COMMAND_FAILED', + `${capitalize(definition.displayName)} runtime returned an incoherent durable envelope`, + { + reason: 'runtime-contract-invalid', + hint: 'Retain the runtime diagnostics and report the selected device and owner.', + }, + ); +} + +function matchesAuthority( + envelope: DurableResourceEnvelope, + expected: { + sessionName: string; + device: DeviceInfo; + owner: RuntimeOwnerRef; + fence: ResourceOwnershipFence; + }, + resourceKind: string, +): boolean { + return ( + envelope.resourceKind === resourceKind && + envelope.sessionId === expected.sessionName && + sameDeviceIdentity(envelope.device, deviceIdentity(expected.device)) && + runtimeOwnerKey(envelope.owner) === runtimeOwnerKey(expected.owner) && + envelope.fence.token === expected.fence.token && + envelope.fence.generation === expected.fence.generation && + envelope.lifecycle === 'open' + ); +} + +function confirmFailedAdoption, C>( + definition: DurableCaptureResourceDefinition, + resourcePath: string, + expected: ResourceOwnershipFence, + cleanupError: unknown | undefined, +): boolean { + const record = definition.store.read(resourcePath); + if ( + record.status !== 'decoded' || + record.envelope.fence.token !== expected.token || + record.envelope.fence.generation !== expected.generation + ) + return false; + definition.store.write(resourcePath, { + ...record.envelope, + lifecycle: cleanupError === undefined ? 'completed' : 'open', + metadata: { + ...(record.envelope.metadata ?? {}), + phase: cleanupError === undefined ? 'completed' : 'cleanup-pending', + cleanupStatus: cleanupError === undefined ? 'cleaned' : 'cleanup-pending', + }, + }); + return true; +} + +function emitCleanupDiagnostic, C>( + definition: DurableCaptureResourceDefinition, + params: Pick, 'sessionName'>, + primaryError: unknown, + cleanupError: unknown, +): void { + emitDiagnostic({ + level: 'error', + phase: `${diagnosticPrefix(definition.resourceKind)}_pending_adoption_cleanup_failed`, + data: { + session: params.sessionName, + primaryError: primaryError instanceof Error ? primaryError.message : String(primaryError), + cleanupError: cleanupError instanceof Error ? cleanupError.message : String(cleanupError), + }, + }); +} + +function withPhase( + envelope: DurableResourceEnvelope, + phase: string, +): DurableResourceEnvelope { + return Object.freeze({ + ...envelope, + lifecycle: 'open', + metadata: { ...(envelope.metadata ?? {}), phase }, + }); +} + +function capitalize(value: string): string { + return value.length === 0 ? value : value[0]!.toUpperCase() + value.slice(1); +} + +function diagnosticPrefix(resourceKind: string): string { + return resourceKind.replaceAll('-', '_'); +} diff --git a/src/daemon/app-log-resource-fence.ts b/src/daemon/durable-capture-resource-fence.ts similarity index 52% rename from src/daemon/app-log-resource-fence.ts rename to src/daemon/durable-capture-resource-fence.ts index 72efc3410..9bb1986a8 100644 --- a/src/daemon/app-log-resource-fence.ts +++ b/src/daemon/durable-capture-resource-fence.ts @@ -4,58 +4,59 @@ import type { DurableResourceLifecycleState, ResourceOwnershipFence, } from '@agent-device/contracts/platform'; -import { readAppLogResourceRecord, writeAppLogResourceRecord } from './app-log-resource-store.ts'; +import type { DurableCaptureResourceStore } from './durable-capture-resource-store.ts'; const resourceFenceTails = new Map>(); -export type AppLogResourceFenceLease = Readonly<{ - envelope: DurableResourceEnvelope<'app-log'>; +export type DurableCaptureResourceFenceLease = Readonly<{ + envelope: DurableResourceEnvelope; transition( lifecycle: DurableResourceLifecycleState, update?: Readonly<{ - descriptor?: DurableResourceEnvelope<'app-log'>['descriptor']; - metadata?: DurableResourceEnvelope<'app-log'>['metadata']; + descriptor?: DurableResourceEnvelope['descriptor']; + metadata?: DurableResourceEnvelope['metadata']; }>, - ): DurableResourceEnvelope<'app-log'>; + ): DurableResourceEnvelope; }>; -/** - * Serializes one app-log resource from persisted fence validation through the - * native side effect and its persisted transition. Production invokes this - * only while the process owns the daemon lock; the per-record queue supplies - * the narrower request/startup mutual exclusion inside that owner. - */ -export async function withAppLogResourceFence(params: { - resourcePath: string; - expected: ResourceOwnershipFence; - run(lease: AppLogResourceFenceLease): Promise; -}): Promise { +/** Holds one persisted ownership fence from validation through side effect and transition. */ +export async function withDurableCaptureResourceFence( + params: Readonly<{ + store: DurableCaptureResourceStore; + resourcePath: string; + expected: ResourceOwnershipFence; + run(lease: DurableCaptureResourceFenceLease): Promise; + }>, +): Promise { return await serializeResource(params.resourcePath, async () => { - const record = readAppLogResourceRecord(params.resourcePath); + const record = params.store.read(params.resourcePath); if (record.status !== 'decoded') { - throw resourceFenceError( + throw fenceError( + params.store, record.status === 'missing' - ? 'App-log resource record is missing' - : `App-log resource record is unreattachable: ${record.message}`, + ? `${params.store.displayName} resource record is missing` + : `${params.store.displayName} resource record is unreattachable: ${record.message}`, ); } if (!sameFence(record.envelope.fence, params.expected)) { - throw resourceFenceError('App-log resource ownership fence was lost'); + throw fenceError( + params.store, + `${params.store.displayName} resource ownership fence was lost`, + ); } - let current = record.envelope; return await params.run({ get envelope() { return current; }, - transition: (lifecycle, update = {}) => { + transition(lifecycle, update = {}) { current = Object.freeze({ ...current, lifecycle, ...(update.descriptor === undefined ? {} : { descriptor: update.descriptor }), ...(update.metadata === undefined ? {} : { metadata: update.metadata }), }); - writeAppLogResourceRecord(params.resourcePath, current); + params.store.write(params.resourcePath, current); return current; }, }); @@ -66,7 +67,10 @@ function sameFence(left: ResourceOwnershipFence, right: ResourceOwnershipFence): return left.token === right.token && left.generation === right.generation; } -async function serializeResource(resourcePath: string, task: () => Promise): Promise { +async function serializeResource( + resourcePath: string, + task: () => Promise, +): Promise { const previous = resourceFenceTails.get(resourcePath) ?? Promise.resolve(); let release!: () => void; const current = new Promise((resolve) => { @@ -83,10 +87,10 @@ async function serializeResource(resourcePath: string, task: () => Promise } } -function resourceFenceError(message: string): AppError { +function fenceError(store: DurableCaptureResourceStore, message: string): AppError { return new AppError('COMMAND_FAILED', message, { reason: 'ownership-fence-lost', retriable: false, - hint: 'Use the current app-log resource owner or retain the recovery record for manual recovery.', + hint: `Use the current ${store.displayName.toLowerCase()} resource owner or retain the recovery record for manual recovery.`, }); } diff --git a/src/daemon/durable-capture-resource-recovery.ts b/src/daemon/durable-capture-resource-recovery.ts new file mode 100644 index 000000000..7a605e73b --- /dev/null +++ b/src/daemon/durable-capture-resource-recovery.ts @@ -0,0 +1,258 @@ +import path from 'node:path'; +import type { + DurableResourceEnvelope, + LiveResourceHandle, + PlatformRequestScope, + ResourceUnreattachableReason, +} from '@agent-device/contracts/platform'; +import { isConfirmedCleanup } from '@agent-device/contracts/platform'; +import { emitDiagnostic } from '../utils/diagnostics.ts'; +import { + acquireDurableCaptureRecoveryAuthorityBeforeDeadline, + DurableCaptureRecoveryDeadlineError, + type DurableCaptureRecoveryAuthority, + type DurableCaptureRecoveryControl, +} from './durable-capture-recovery-authority.ts'; +import { + withDurableCaptureResourceFence, + type DurableCaptureResourceFenceLease, +} from './durable-capture-resource-fence.ts'; +import type { DurableCaptureResourceDefinition } from './durable-capture-resource.ts'; +import { transitionCleanupOutcome } from './durable-capture-resource-transitions.ts'; +import { safeSessionName } from './session-paths.ts'; + +const DEFAULT_RECOVERY_DEADLINE_MS = 5_000; + +export type DurableCaptureRecoverySummary = Readonly<{ + scanned: number; + recovered: number; + retained: number; +}>; + +export type DurableCaptureRecoveryDiagnostic = Readonly<{ + phase: string; + resourcePath: string; + data: Readonly>; +}>; + +export type DurableCaptureRecoveryParams, C> = { + definition: DurableCaptureResourceDefinition; + sessionsDir: string; + scope: PlatformRequestScope; + acquireControl( + envelope: DurableResourceEnvelope, + scope: PlatformRequestScope, + ): Promise>; + perRecordDeadlineMs?: number; + onDiagnostic?: (diagnostic: DurableCaptureRecoveryDiagnostic) => void; +}; + +export async function recoverDurableCaptureResourcesAfterDaemonLock< + K extends string, + H extends LiveResourceHandle, + C, +>(params: DurableCaptureRecoveryParams): Promise { + const paths = params.definition.store.list(params.sessionsDir); + const outcomes: Array<'ignored' | 'recovered' | 'retained'> = []; + for (const resourcePath of paths) { + outcomes.push(await recoverResourcePath(params, resourcePath)); + } + return { + scanned: paths.length, + recovered: outcomes.filter((outcome) => outcome === 'recovered').length, + retained: outcomes.filter((outcome) => outcome === 'retained').length, + }; +} + +async function recoverResourcePath, C>( + params: DurableCaptureRecoveryParams, + resourcePath: string, +): Promise<'ignored' | 'recovered' | 'retained'> { + const candidate = readRecoveryCandidate(params, resourcePath); + if (candidate.status !== 'recoverable') return candidate.status; + try { + const recovered = await recoverBeforeDeadline( + params, + resourcePath, + candidate.envelope, + params.perRecordDeadlineMs ?? DEFAULT_RECOVERY_DEADLINE_MS, + ); + return recovered ? 'recovered' : 'retained'; + } catch (error) { + report( + params, + error instanceof DurableCaptureRecoveryDeadlineError ? 'timed_out' : 'failed', + resourcePath, + { + error: error instanceof Error ? error.message : String(error), + ...(error instanceof DurableCaptureRecoveryDeadlineError + ? { deadlineMs: error.deadlineMs } + : {}), + }, + ); + return 'retained'; + } +} + +function readRecoveryCandidate, C>( + params: DurableCaptureRecoveryParams, + resourcePath: string, +): + | Readonly<{ status: 'ignored' | 'retained' }> + | Readonly<{ status: 'recoverable'; envelope: DurableResourceEnvelope }> { + const record = params.definition.store.read(resourcePath); + if (record.status === 'unreattachable') { + report(params, 'record_unreattachable', resourcePath, { + reason: record.reason, + message: record.message, + version: record.version, + }); + return { status: 'retained' }; + } + if (record.status === 'missing' || record.envelope.lifecycle === 'completed') { + return { status: 'ignored' }; + } + const canonicalPath = params.definition.store.resolvePath( + path.join(params.sessionsDir, safeSessionName(record.envelope.sessionId)), + ); + if (path.resolve(resourcePath) !== path.resolve(canonicalPath)) { + report(params, 'session_path_mismatch', resourcePath, { + envelopeSessionId: record.envelope.sessionId, + canonicalPath, + }); + return { status: 'retained' }; + } + return { status: 'recoverable', envelope: record.envelope }; +} + +async function recoverBeforeDeadline, C>( + params: DurableCaptureRecoveryParams, + resourcePath: string, + envelope: DurableResourceEnvelope, + deadlineMs: number, +): Promise { + const authority = await acquireDurableCaptureRecoveryAuthorityBeforeDeadline({ + displayName: params.definition.displayName, + envelope, + scope: params.scope, + deadlineMs, + acquireControl: params.acquireControl, + onLateCleanupFailure: (phase, cleanupError, primaryError) => + report(params, phase, resourcePath, { + error: cleanupError instanceof Error ? cleanupError.message : String(cleanupError), + primaryError: primaryError instanceof Error ? primaryError.message : String(primaryError), + }), + }); + return await settleRecoveryAuthority(params, resourcePath, envelope, authority); +} + +async function settleRecoveryAuthority, C>( + params: DurableCaptureRecoveryParams, + resourcePath: string, + envelope: DurableResourceEnvelope, + authority: DurableCaptureRecoveryAuthority, +): Promise { + try { + switch (authority.reattached.status) { + case 'active': { + const handle = authority.reattached.handle; + const cleanup = await withDurableCaptureResourceFence({ + store: params.definition.store, + resourcePath, + expected: envelope.fence, + run: async (lease) => { + markCompleting(lease); + const outcome = await handle.forceCleanup(); + transitionCleanupOutcome(lease, outcome); + return outcome; + }, + }); + return isConfirmedCleanup(cleanup); + } + case 'completed': + await transitionTerminal(params, resourcePath, envelope, { + ...params.definition.completionMetadata(authority.reattached.result), + }); + return true; + case 'missing': + await transitionTerminal(params, resourcePath, envelope, { + recoveryStatus: 'already-missing', + }); + return true; + case 'unreattachable': { + if (retainsWithoutCleanup(authority.reattached.reason)) { + report(params, 'retained', resourcePath, { + reason: authority.reattached.reason, + message: authority.reattached.message, + }); + return false; + } + const cleanup = await authority.control.cleanup(envelope); + await withDurableCaptureResourceFence({ + store: params.definition.store, + resourcePath, + expected: envelope.fence, + run: async (lease) => transitionCleanupOutcome(lease, cleanup), + }); + return isConfirmedCleanup(cleanup); + } + } + } finally { + await authority.control[Symbol.asyncDispose](); + } +} + +async function transitionTerminal, C>( + params: DurableCaptureRecoveryParams, + resourcePath: string, + envelope: DurableResourceEnvelope, + metadata: Record, +): Promise { + await withDurableCaptureResourceFence({ + store: params.definition.store, + resourcePath, + expected: envelope.fence, + run: async (lease) => { + lease.transition('completed', { + metadata: { ...(lease.envelope.metadata ?? {}), phase: 'completed', ...metadata }, + }); + }, + }); +} + +function markCompleting(lease: DurableCaptureResourceFenceLease): void { + lease.transition('open', { + metadata: { ...(lease.envelope.metadata ?? {}), phase: 'completing' }, + }); +} + +function retainsWithoutCleanup(reason: ResourceUnreattachableReason): boolean { + switch (reason) { + case 'descriptor-invalid': + case 'descriptor-version-unsupported': + case 'ownership-fence-lost': + return true; + case 'owner-unavailable': + case 'transport-not-reattachable': + return false; + } +} + +function report, C>( + params: DurableCaptureRecoveryParams, + suffix: string, + resourcePath: string, + data: Record, +): void { + const diagnostic = { + phase: `${diagnosticPrefix(params.definition.resourceKind)}_recovery_${suffix}`, + resourcePath, + data, + }; + if (params.onDiagnostic) params.onDiagnostic(diagnostic); + else emitDiagnostic({ level: 'warn', phase: diagnostic.phase, data: { resourcePath, ...data } }); +} + +function diagnosticPrefix(resourceKind: string): string { + return resourceKind.replaceAll('-', '_'); +} diff --git a/src/daemon/durable-capture-resource-store.ts b/src/daemon/durable-capture-resource-store.ts new file mode 100644 index 000000000..02de64935 --- /dev/null +++ b/src/daemon/durable-capture-resource-store.ts @@ -0,0 +1,172 @@ +import crypto from 'node:crypto'; +import fs from 'node:fs'; +import path from 'node:path'; +import type { + DurableEnvelopeDecodeOutcome, + DurableResourceEnvelope, +} from '@agent-device/contracts/platform'; +import { decodeDurableResourceEnvelope } from '@agent-device/capture-kit'; +import { openVerifiedFileForRead } from '../utils/verified-file.ts'; + +export type DurableCaptureResourceRecord = + | Readonly<{ status: 'missing' }> + | Readonly<{ status: 'decoded'; envelope: DurableResourceEnvelope }> + | Readonly<{ + status: 'unreattachable'; + reason: 'descriptor-invalid' | 'descriptor-version-unsupported'; + message: string; + version?: number; + }>; + +export type DurableCaptureResourceStore = Readonly<{ + resourceKind: K; + displayName: string; + resolvePath(sessionDir: string): string; + read(resourcePath: string): DurableCaptureResourceRecord; + write(resourcePath: string, envelope: DurableResourceEnvelope): void; + list(sessionsDir: string): string[]; +}>; + +export function createDurableCaptureResourceStore( + options: Readonly<{ + resourceKind: K; + fileName: string; + displayName: string; + }>, +): DurableCaptureResourceStore { + assertFileName(options.fileName); + const resolvePath = (sessionDir: string): string => path.join(sessionDir, options.fileName); + return Object.freeze({ + resourceKind: options.resourceKind, + displayName: options.displayName, + resolvePath, + read(resourcePath: string): DurableCaptureResourceRecord { + let descriptor: number | undefined; + let value: unknown; + try { + descriptor = openVerifiedFileForRead(resourcePath); + if (descriptor === undefined) return { status: 'missing' }; + value = JSON.parse(fs.readFileSync(descriptor, 'utf8')) as unknown; + } catch (error) { + if (isMissingFile(error)) return { status: 'missing' }; + return invalidRecord(errorMessage(error, options.displayName)); + } finally { + if (descriptor !== undefined) fs.closeSync(descriptor); + } + return narrowEnvelope(decodeDurableResourceEnvelope(value), options.resourceKind); + }, + write(resourcePath: string, envelope: DurableResourceEnvelope): void { + const directory = path.dirname(resourcePath); + fs.mkdirSync(directory, { recursive: true }); + const temporaryPath = path.join( + directory, + `.${path.basename(resourcePath)}.${process.pid}.${crypto.randomUUID()}.tmp`, + ); + let descriptor: number | undefined; + try { + assertSafeDestination(resourcePath, options.displayName); + descriptor = fs.openSync(temporaryPath, 'wx', 0o600); + fs.writeFileSync(descriptor, `${JSON.stringify(envelope)}\n`, 'utf8'); + fs.fsyncSync(descriptor); + fs.closeSync(descriptor); + descriptor = undefined; + assertSafeDestination(resourcePath, options.displayName); + fs.renameSync(temporaryPath, resourcePath); + syncDirectoryBestEffort(directory); + } finally { + if (descriptor !== undefined) fs.closeSync(descriptor); + try { + fs.rmSync(temporaryPath, { force: true }); + } catch {} + } + }, + list(sessionsDir: string): string[] { + let entries: fs.Dirent[]; + try { + entries = fs.readdirSync(sessionsDir, { withFileTypes: true }); + } catch { + return []; + } + return entries + .filter((entry) => entry.isDirectory()) + .map((entry) => resolvePath(path.join(sessionsDir, entry.name))) + .filter(pathEntryExistsWithoutFollowing) + .sort(); + }, + }); +} + +function narrowEnvelope( + decoded: DurableEnvelopeDecodeOutcome, + resourceKind: K, +): DurableCaptureResourceRecord { + if (decoded.status !== 'decoded') return decoded; + if (decoded.envelope.resourceKind !== resourceKind) { + return invalidRecord( + `Expected ${resourceKind} resource record, received ${decoded.envelope.resourceKind}`, + ); + } + return { status: 'decoded', envelope: decoded.envelope as DurableResourceEnvelope }; +} + +function assertFileName(fileName: string): void { + if (fileName.length === 0 || path.basename(fileName) !== fileName) { + throw new TypeError('Durable capture resource fileName must be one file name'); + } +} + +function assertSafeDestination(resourcePath: string, displayName: string): void { + let destination: fs.Stats; + try { + destination = fs.lstatSync(resourcePath); + } catch (error) { + if (isMissingFile(error)) return; + throw error; + } + if (destination.isSymbolicLink()) { + throw new Error(`Refusing to replace a ${displayName.toLowerCase()} resource symbolic link`); + } + if (!destination.isFile()) { + throw new Error( + `Refusing to replace a ${displayName.toLowerCase()} resource that is not a regular file`, + ); + } +} + +function pathEntryExistsWithoutFollowing(resourcePath: string): boolean { + try { + fs.lstatSync(resourcePath); + return true; + } catch (error) { + return !isMissingFile(error); + } +} + +function invalidRecord(message: string): DurableCaptureResourceRecord { + return { status: 'unreattachable', reason: 'descriptor-invalid', message }; +} + +function syncDirectoryBestEffort(directory: string): void { + let descriptor: number | undefined; + try { + descriptor = fs.openSync(directory, 'r'); + fs.fsyncSync(descriptor); + } catch { + // Atomic rename is the correctness boundary; directory fsync support varies by filesystem. + } finally { + if (descriptor !== undefined) fs.closeSync(descriptor); + } +} + +function isMissingFile(error: unknown): boolean { + return ( + error !== null && + typeof error === 'object' && + 'code' in error && + (error as { code?: unknown }).code === 'ENOENT' + ); +} + +function errorMessage(error: unknown, displayName: string): string { + return error instanceof Error ? error.message : `${displayName} resource record is invalid`; +} diff --git a/src/daemon/durable-capture-resource-transitions.ts b/src/daemon/durable-capture-resource-transitions.ts new file mode 100644 index 000000000..850294c5d --- /dev/null +++ b/src/daemon/durable-capture-resource-transitions.ts @@ -0,0 +1,144 @@ +import { + isConfirmedCleanup, + type CleanupOutcome, + type FinishOutcome, + type LiveResourceHandle, +} from '@agent-device/contracts/platform'; +import { AppError } from '@agent-device/kernel/errors'; +import { + withDurableCaptureResourceFence, + type DurableCaptureResourceFenceLease, +} from './durable-capture-resource-fence.ts'; +import type { DurableCaptureResourceDefinition } from './durable-capture-resource.ts'; +import type { SessionStore } from './session-store.ts'; +import type { SessionState } from './types.ts'; + +export async function finishLiveDurableCapture< + K extends string, + H extends LiveResourceHandle, + C, +>( + definition: DurableCaptureResourceDefinition, + params: { session: SessionState; sessionName: string; sessionStore: SessionStore }, + resourcePath: string, +): Promise { + const active = definition.sessionSlot.read(params.session); + if (!active) throw new AppError('INVALID_ARGS', definition.messages.noActive); + const outcome = await withDurableCaptureResourceFence({ + store: definition.store, + resourcePath, + expected: active.envelope.fence, + run: async (lease) => { + markCompleting(lease); + const result = await active.handle.finish(); + transitionFinishOutcome(definition, lease, result); + return result; + }, + }); + if (outcome.status === 'cleanup-pending') throw cleanupPendingError(definition, outcome); + params.sessionStore.set( + params.sessionName, + definition.sessionSlot.replace(params.session, undefined), + ); + return outcome.result; +} + +export async function forceCleanupLiveDurableCapture< + K extends string, + H extends LiveResourceHandle, + C, +>( + definition: DurableCaptureResourceDefinition, + params: { + session: SessionState; + sessionName?: string; + sessionStore?: SessionStore; + resourcePath: string; + }, +): Promise { + const active = definition.sessionSlot.read(params.session); + if (!active) return; + const outcome = await withDurableCaptureResourceFence({ + store: definition.store, + resourcePath: params.resourcePath, + expected: active.envelope.fence, + run: async (lease) => { + markCompleting(lease); + const result = await active.handle.forceCleanup(); + transitionCleanupOutcome(lease, result); + return result; + }, + }); + if (!isConfirmedCleanup(outcome)) throw cleanupPendingError(definition, outcome); + if (params.sessionStore && params.sessionName) { + params.sessionStore.set( + params.sessionName, + definition.sessionSlot.replace(params.session, undefined), + ); + } +} + +export function transitionCleanupOutcome( + lease: DurableCaptureResourceFenceLease, + outcome: CleanupOutcome, +): void { + lease.transition(isConfirmedCleanup(outcome) ? 'completed' : 'open', { + metadata: { + ...(lease.envelope.metadata ?? {}), + phase: isConfirmedCleanup(outcome) ? 'completed' : 'cleanup-pending', + cleanupStatus: outcome.status, + ...(outcome.status === 'cleanup-pending' + ? { + cleanupPendingReason: outcome.reason, + ...(outcome.message ? { cleanupPendingMessage: outcome.message } : {}), + } + : {}), + }, + }); +} + +function cleanupPendingError( + definition: Pick< + DurableCaptureResourceDefinition, unknown>, + 'displayName' | 'messages' + >, + outcome: Extract, +): AppError { + return new AppError( + 'COMMAND_FAILED', + outcome.message ?? `${capitalize(definition.displayName)} cleanup could not be confirmed`, + { + reason: outcome.reason, + retriable: outcome.reason !== 'ownership-fence-lost', + hint: definition.messages.cleanupPendingHint, + }, + ); +} + +function transitionFinishOutcome, C>( + definition: DurableCaptureResourceDefinition, + lease: DurableCaptureResourceFenceLease, + outcome: FinishOutcome, +): void { + if (outcome.status === 'completed') { + lease.transition('completed', { + metadata: { + ...(lease.envelope.metadata ?? {}), + ...definition.completionMetadata(outcome.result), + phase: 'completed', + }, + }); + return; + } + transitionCleanupOutcome(lease, outcome); +} + +function markCompleting(lease: DurableCaptureResourceFenceLease): void { + lease.transition('open', { + metadata: { ...(lease.envelope.metadata ?? {}), phase: 'completing' }, + }); +} + +function capitalize(value: string): string { + return value.length === 0 ? value : value[0]!.toUpperCase() + value.slice(1); +} diff --git a/src/daemon/durable-capture-resource.ts b/src/daemon/durable-capture-resource.ts new file mode 100644 index 000000000..d156e5f11 --- /dev/null +++ b/src/daemon/durable-capture-resource.ts @@ -0,0 +1,112 @@ +import type { JsonObject } from '@agent-device/contracts/client'; +import type { + DurableResourceEnvelope, + LiveResourceHandle, + PendingTransferGuard, + ResourceOwnershipFence, + RuntimeOwnerRef, +} from '@agent-device/contracts/platform'; +import type { DeviceInfo } from '@agent-device/kernel/device'; +import type { DurableCaptureAdmissionLedger } from './durable-capture-admission-ledger.ts'; +import { adoptStartedDurableCapture } from './durable-capture-resource-adoption.ts'; +import { + finishLiveDurableCapture, + forceCleanupLiveDurableCapture, +} from './durable-capture-resource-transitions.ts'; +import type { DurableCaptureResourceStore } from './durable-capture-resource-store.ts'; +import { createNextDurableCaptureFence } from './durable-capture-start-preflight.ts'; +import { + recoverDurableCaptureResourcesAfterDaemonLock, + type DurableCaptureRecoveryParams, +} from './durable-capture-resource-recovery.ts'; +import type { SessionStore } from './session-store.ts'; +import type { SessionState } from './types.ts'; + +export type DurableCaptureSessionResource = Readonly<{ + handle: H; + envelope: DurableResourceEnvelope; +}>; + +export type DurableCaptureSessionSlot = Readonly<{ + read(session: SessionState): DurableCaptureSessionResource | undefined; + replace( + session: SessionState, + resource: DurableCaptureSessionResource | undefined, + ): SessionState; +}>; + +export type DurableCaptureResourceDefinition< + K extends string, + H extends LiveResourceHandle, + C, +> = Readonly<{ + resourceKind: K; + displayName: string; + store: DurableCaptureResourceStore; + sessionSlot: DurableCaptureSessionSlot; + completionMetadata(result: C): JsonObject; + messages: Readonly<{ + noActive: string; + cleanupPendingHint: string; + }>; +}>; + +export type AdoptStartedDurableCaptureParams = { + admissionLedger: DurableCaptureAdmissionLedger; + session: SessionState; + sessionName: string; + sessionStore: SessionStore; + device: DeviceInfo; + owner: RuntimeOwnerRef; + fence: ResourceOwnershipFence; + pendingHandle: PendingTransferGuard; + envelope: DurableResourceEnvelope; + throwIfCanceled(): void; +}; + +export function createDurableCaptureResource, C>( + definition: DurableCaptureResourceDefinition, +) { + const resourcePath = (sessionStore: SessionStore, sessionName: string): string => + definition.store.resolvePath(sessionStore.resolveSessionDir(sessionName)); + + return Object.freeze({ + store: definition.store, + createNextFence(params: { + admissionLedger: DurableCaptureAdmissionLedger; + resourcePath: string; + device: DeviceInfo; + }): ResourceOwnershipFence { + return createNextDurableCaptureFence(definition, params); + }, + adoptStarted(params: AdoptStartedDurableCaptureParams): Promise { + return adoptStartedDurableCapture( + definition, + params, + resourcePath(params.sessionStore, params.sessionName), + ); + }, + finishLive(params: { + session: SessionState; + sessionName: string; + sessionStore: SessionStore; + }): Promise { + return finishLiveDurableCapture( + definition, + params, + resourcePath(params.sessionStore, params.sessionName), + ); + }, + forceCleanupLive(params: { + session: SessionState; + sessionName?: string; + sessionStore?: SessionStore; + resourcePath: string; + }): Promise { + return forceCleanupLiveDurableCapture(definition, params); + }, + recoverAll(params: Omit, 'definition'>) { + return recoverDurableCaptureResourcesAfterDaemonLock({ definition, ...params }); + }, + }); +} diff --git a/src/daemon/durable-capture-start-preflight.ts b/src/daemon/durable-capture-start-preflight.ts new file mode 100644 index 000000000..ffc97ba9c --- /dev/null +++ b/src/daemon/durable-capture-start-preflight.ts @@ -0,0 +1,69 @@ +import crypto from 'node:crypto'; +import path from 'node:path'; +import type { LiveResourceHandle, ResourceOwnershipFence } from '@agent-device/contracts/platform'; +import { deviceIdentity, deviceIdentityKey, type DeviceInfo } from '@agent-device/kernel/device'; +import { AppError } from '@agent-device/kernel/errors'; +import type { DurableCaptureAdmissionLedger } from './durable-capture-admission-ledger.ts'; +import type { DurableCaptureResourceDefinition } from './durable-capture-resource.ts'; + +export function createNextDurableCaptureFence, C>( + definition: DurableCaptureResourceDefinition, + params: { + admissionLedger: DurableCaptureAdmissionLedger; + resourcePath: string; + device: DeviceInfo; + }, +): ResourceOwnershipFence { + params.admissionLedger.assertStartAllowed(params.device); + assertNoConflictingManifest(definition, params.resourcePath, params.device); + const current = definition.store.read(params.resourcePath); + return Object.freeze({ + token: crypto.randomUUID(), + generation: current.status === 'decoded' ? current.envelope.fence.generation + 1 : 1, + }); +} + +function assertNoConflictingManifest, C>( + definition: DurableCaptureResourceDefinition, + resourcePath: string, + device: DeviceInfo, +): void { + const selectedDeviceKey = deviceIdentityKey(deviceIdentity(device)); + for (const existingPath of definition.store.list(path.dirname(path.dirname(resourcePath)))) { + const existing = definition.store.read(existingPath); + if (existing.status === 'unreattachable') { + throw new AppError( + 'COMMAND_FAILED', + `${capitalize(articleFor(definition.displayName))} ${definition.displayName} recovery record is unreattachable: ${existing.message}`, + { + reason: existing.reason, + retriable: false, + hint: `Retain the corrupt or future-version ${definition.displayName} manifest for manual recovery; no replacement capture is safe.`, + }, + ); + } + if ( + existing.status === 'decoded' && + existing.envelope.lifecycle !== 'completed' && + (existingPath === resourcePath || + deviceIdentityKey(existing.envelope.device) === selectedDeviceKey) + ) { + throw new AppError( + 'COMMAND_FAILED', + `${capitalize(articleFor(definition.displayName))} ${definition.displayName} resource for this device has not reached a confirmed terminal state`, + { + reason: 'cleanup-unconfirmed', + hint: `Retry exact-owner cleanup using the existing ${definition.displayName} manifest before starting a replacement.`, + }, + ); + } + } +} + +function capitalize(value: string): string { + return value.length === 0 ? value : value[0]!.toUpperCase() + value.slice(1); +} + +function articleFor(value: string): 'a' | 'an' { + return /^[aeiou]/i.test(value) ? 'an' : 'a'; +}