diff --git a/packages/capture-kit/src/capture-admission/__tests__/audio-probe-session-resource.test.ts b/packages/capture-kit/src/capture-admission/__tests__/audio-probe-session-resource.test.ts index 073c33419b..359fb5f545 100644 --- a/packages/capture-kit/src/capture-admission/__tests__/audio-probe-session-resource.test.ts +++ b/packages/capture-kit/src/capture-admission/__tests__/audio-probe-session-resource.test.ts @@ -1,3 +1,4 @@ +import { makeCaptureSessionBinding } from '../../durable-capture/session-binding.fixtures.ts'; import fs from 'node:fs/promises'; import path from 'node:path'; import { expect, test, vi } from 'vitest'; @@ -44,6 +45,10 @@ test('audio-probe disposes on a failed finish because terminating the helper is ); const session: DurableCaptureSessionState = {}; sessionStore.set(sessionName, session); + const binding = makeCaptureSessionBinding(sessionStore, sessionName, { + read: (session) => session.audioProbe, + replace: (session, audioProbe) => ({ ...session, audioProbe }), + }); const statusPath = path.join(sessionStore.resolveSessionDir(sessionName), 'audio-probe.json'); const terminate = vi.fn(async () => {}); let resolveExit!: (result: HostCommandResult) => void; @@ -85,9 +90,7 @@ test('audio-probe disposes on a failed finish because terminating the helper is }); await adoptStartedAudioProbe({ admissionLedger: createAudioProbeAdmissionLedger(), - session, - sessionName, - sessionStore, + binding, device, owner: localRuntimeOwner('apple'), fence, diff --git a/packages/capture-kit/src/capture-admission/__tests__/durable-capture-resource.fixtures.ts b/packages/capture-kit/src/capture-admission/__tests__/durable-capture-resource.fixtures.ts index 0cef989b60..e022419770 100644 --- a/packages/capture-kit/src/capture-admission/__tests__/durable-capture-resource.fixtures.ts +++ b/packages/capture-kit/src/capture-admission/__tests__/durable-capture-resource.fixtures.ts @@ -1,3 +1,4 @@ +import { makeCaptureSessionBinding } from '../../durable-capture/session-binding.fixtures.ts'; import { vi } from 'vitest'; import type { AppLogCompletion, AppLogLiveHandle } from '@agent-device/contracts/app-log-runtime'; import type { CleanupOutcome, FinishOutcome } from '@agent-device/contracts/durable-resource'; @@ -72,6 +73,10 @@ export function makeDurableCaptureContext( const session: TestCaptureSession = {}; sessionStore.set(sessionName, session); return { + binding: makeCaptureSessionBinding(sessionStore, sessionName, { + read: (session) => session.appLog, + replace: (session, appLog) => ({ ...session, appLog, appLogFailure: undefined }), + }), admissionLedger: createDurableCaptureAdmissionLedger({ displayName: 'test capture' }), session, sessionName, diff --git a/packages/capture-kit/src/capture-admission/__tests__/screen-recording-boundary-faults.test.ts b/packages/capture-kit/src/capture-admission/__tests__/screen-recording-boundary-faults.test.ts index 775f5a3fae..dfab3f92e8 100644 --- a/packages/capture-kit/src/capture-admission/__tests__/screen-recording-boundary-faults.test.ts +++ b/packages/capture-kit/src/capture-admission/__tests__/screen-recording-boundary-faults.test.ts @@ -1,3 +1,4 @@ +import { makeCaptureSessionBinding } from '../../durable-capture/session-binding.fixtures.ts'; import path from 'node:path'; import { expect, test, vi } from 'vitest'; import { createDurableResourceEnvelope } from '../../durable-resource-envelope.ts'; @@ -217,9 +218,14 @@ function makeContext(resource: ReturnType session.screenRecording, + replace: (session, screenRecording) => ({ ...session, screenRecording }), + }); return { admissionLedger: createDurableCaptureAdmissionLedger({ displayName: 'screen recording' }), session, + binding, sessionName, sessionStore, device, diff --git a/packages/capture-kit/src/capture-admission/__tests__/screen-recording-session-resource.test.ts b/packages/capture-kit/src/capture-admission/__tests__/screen-recording-session-resource.test.ts index b048db0b3c..a64497c19d 100644 --- a/packages/capture-kit/src/capture-admission/__tests__/screen-recording-session-resource.test.ts +++ b/packages/capture-kit/src/capture-admission/__tests__/screen-recording-session-resource.test.ts @@ -1,3 +1,4 @@ +import { makeCaptureSessionBinding } from '../../durable-capture/session-binding.fixtures.ts'; import { expect, test, vi } from 'vitest'; import { PendingTransferGuard } from '@agent-device/contracts/async-lifecycle'; import { localRuntimeOwner } from '@agent-device/contracts/platform-runtime'; @@ -35,6 +36,10 @@ test('screen recording persists durable truth before adopting only handle and en const sessionName = 'recording'; const session: TestRecordingSession = { name: sessionName, device }; sessionStore.set(sessionName, session); + const binding = makeCaptureSessionBinding(sessionStore, sessionName, { + read: (session) => session.screenRecording, + replace: (session, screenRecording) => ({ ...session, screenRecording }), + }); const owner = localRuntimeOwner('android'); const fence = { token: 'recording-fence', generation: 1 } as const; const finish = vi.fn(async () => ({ @@ -79,9 +84,7 @@ test('screen recording persists durable truth before adopting only handle and en await adoptStartedScreenRecording({ admissionLedger: createScreenRecordingAdmissionLedger(), - session, - sessionName, - sessionStore, + binding, device: session.device, owner, fence, @@ -103,7 +106,10 @@ test('screen recording persists durable truth before adopting only handle and en if (!active) throw new Error('Expected screen-recording session'); await expect( finishLiveScreenRecording({ intent: 'capture', session: active, sessionName, sessionStore }), - ).resolves.toMatchObject({ backend: 'android', outPath: '/tmp/recording.mp4' }); + ).resolves.toMatchObject({ + backend: 'android', + outPath: '/tmp/recording.mp4', + }); expect(finish).toHaveBeenCalledOnce(); expect(sessionStore.get(sessionName)?.screenRecording).toBeUndefined(); }); @@ -115,6 +121,10 @@ test('a failed recording finish keeps the record open and never disposes the rec const sessionName = 'recording'; const session: TestRecordingSession = { name: sessionName, device }; sessionStore.set(sessionName, session); + const binding = makeCaptureSessionBinding(sessionStore, sessionName, { + read: (session) => session.screenRecording, + replace: (session, screenRecording) => ({ ...session, screenRecording }), + }); const owner = localRuntimeOwner('android'); const fence = { token: 'recording-fence', generation: 1 } as const; const finishError = new Error('failed to retrieve playable Android recording'); @@ -150,9 +160,7 @@ test('a failed recording finish keeps the record open and never disposes the rec }); await adoptStartedScreenRecording({ admissionLedger: createScreenRecordingAdmissionLedger(), - session, - sessionName, - sessionStore, + binding, device: session.device, owner, fence, @@ -182,6 +190,10 @@ test('a record stop that fails after collecting resumes through the fence withou const sessionName = 'recording'; const session: TestRecordingSession = { name: sessionName, device }; sessionStore.set(sessionName, session); + const binding = makeCaptureSessionBinding(sessionStore, sessionName, { + read: (session) => session.screenRecording, + replace: (session, screenRecording) => ({ ...session, screenRecording }), + }); const owner = localRuntimeOwner('android'); const fence = { token: 'recording-fence', generation: 1 } as const; const signals = vi.fn(async () => ({ observation: { recorder: 'confirmed' as const } })); @@ -216,9 +228,7 @@ test('a record stop that fails after collecting resumes through the fence withou ); await adoptStartedScreenRecording({ admissionLedger: createScreenRecordingAdmissionLedger(), - session, - sessionName, - sessionStore, + binding, device: session.device, owner, fence, @@ -239,7 +249,7 @@ test('a record stop that fails after collecting resumes through the fence withou if (!active) throw new Error('Expected screen-recording session'); return finishLiveScreenRecording({ intent: 'capture', - session: active, + session: sessionStore.get(sessionName) ?? session, sessionName, sessionStore, }); diff --git a/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts b/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts index dee8763861..ab0ccd73e0 100644 --- a/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts +++ b/packages/capture-kit/src/capture-admission/__tests__/session-store.fixtures.ts @@ -1,13 +1,13 @@ import path from 'node:path'; import { safeSessionName } from '@agent-device/host-kit/session-paths'; -import type { DurableCaptureSessionStore } from '../../durable-capture/index.ts'; import { mkdtempForTestSync } from '../../tmp-dir.fixtures.ts'; -export type CaptureAdmissionSessionStore = DurableCaptureSessionStore & - Readonly<{ - get(name: string): S | undefined; - sessionsDir: string; - }>; +export type CaptureAdmissionSessionStore = Readonly<{ + set(name: string, session: S): void; + resolveSessionDir(name: string): string; + get(name: string): S | undefined; + sessionsDir: string; +}>; /** * The whole of the daemon `SessionStore` these admission modules ever address โ€” `set`, diff --git a/packages/capture-kit/src/capture-admission/audio-probe-session-resource.ts b/packages/capture-kit/src/capture-admission/audio-probe-session-resource.ts index 83863d3336..ddcd8037b0 100644 --- a/packages/capture-kit/src/capture-admission/audio-probe-session-resource.ts +++ b/packages/capture-kit/src/capture-admission/audio-probe-session-resource.ts @@ -9,7 +9,10 @@ import type { RuntimeOwnerRef, } from '@agent-device/contracts/platform-runtime'; import type { DeviceInfo } from '@agent-device/kernel/device'; -import type { DurableCaptureSessionStore } from '../durable-capture/index.ts'; +import type { + DurableCaptureSessionBinding, + DurableCaptureSessionStore, +} from '../durable-capture/index.ts'; import { createDurableCaptureResource } from './durable-capture-resource.ts'; import type { DurableCaptureFinishIntent } from './durable-capture-resource.ts'; import type { AudioProbeAdmissionLedger } from './audio-probe-admission-ledger.ts'; @@ -48,9 +51,7 @@ export const audioProbeDurableResource = createDurableCaptureResource< export function adoptStartedAudioProbe(params: { admissionLedger: AudioProbeAdmissionLedger; - session: DurableCaptureSessionState; - sessionName: string; - sessionStore: DurableCaptureSessionStore; + binding: DurableCaptureSessionBinding<'audio-probe', AudioProbeLiveHandle>; device: DeviceInfo; owner: RuntimeOwnerRef; fence: ResourceOwnershipFence; diff --git a/packages/capture-kit/src/capture-admission/durable-capture-resource.ts b/packages/capture-kit/src/capture-admission/durable-capture-resource.ts index a03751be14..4dd4636244 100644 --- a/packages/capture-kit/src/capture-admission/durable-capture-resource.ts +++ b/packages/capture-kit/src/capture-admission/durable-capture-resource.ts @@ -23,8 +23,8 @@ import type { DurableSessionResourceKind } from './durable-session-resource-kind export type { DurableCaptureFinishIntent, DurableSessionResourceKind }; -type AdoptStartedSessionCaptureParams = Omit< - AdoptStartedDurableCaptureParams, +type AdoptStartedSessionCaptureParams = Omit< + AdoptStartedDurableCaptureParams, 'reportUndurableCleanup' > & Readonly<{ admissionLedger: DurableCaptureAdmissionLedger }>; @@ -48,7 +48,7 @@ export function createDurableCaptureResource< S, >(definition: DurableCaptureResourceDefinition) { const sessionResourcePath = ( - sessionStore: DurableCaptureSessionStore, + sessionStore: Readonly<{ resolveSessionDir(name: string): string }>, sessionName: string, ): string => definition.store.resolvePath(sessionStore.resolveSessionDir(sessionName)); const recoveryParams = ( @@ -70,7 +70,7 @@ export function createDurableCaptureResource< }): ResourceOwnershipFence { return createNextDurableCaptureFence(definition, params); }, - adoptStarted(params: AdoptStartedSessionCaptureParams): Promise { + adoptStarted(params: AdoptStartedSessionCaptureParams): Promise { return adoptStartedDurableCapture( definition, { @@ -80,7 +80,7 @@ export function createDurableCaptureResource< else params.admissionLedger.blockUndurableCleanup(device, outcome.reason); }, }, - sessionResourcePath(params.sessionStore, params.sessionName), + definition.store.resolvePath(params.binding.sessionDir), ); }, finishLive(params: { diff --git a/packages/capture-kit/src/capture-admission/perf-capture-session-resource.ts b/packages/capture-kit/src/capture-admission/perf-capture-session-resource.ts index 6e6e14104b..41546e034f 100644 --- a/packages/capture-kit/src/capture-admission/perf-capture-session-resource.ts +++ b/packages/capture-kit/src/capture-admission/perf-capture-session-resource.ts @@ -9,7 +9,10 @@ import type { RuntimeOwnerRef, } from '@agent-device/contracts/platform-runtime'; import type { DeviceInfo } from '@agent-device/kernel/device'; -import type { DurableCaptureSessionStore } from '../durable-capture/index.ts'; +import type { + DurableCaptureSessionBinding, + DurableCaptureSessionStore, +} from '../durable-capture/index.ts'; import { createDurableCaptureResource } from './durable-capture-resource.ts'; import type { DurableCaptureFinishIntent } from './durable-capture-resource.ts'; import type { PerfCaptureAdmissionLedger } from './perf-capture-admission-ledger.ts'; @@ -46,9 +49,7 @@ export const perfCaptureDurableResource = createDurableCaptureResource< export function adoptStartedPerfCapture(params: { admissionLedger: PerfCaptureAdmissionLedger; - session: DurableCaptureSessionState; - sessionName: string; - sessionStore: DurableCaptureSessionStore; + binding: DurableCaptureSessionBinding<'perf-capture', PerfNativeCaptureLiveHandle>; device: DeviceInfo; owner: RuntimeOwnerRef; fence: ResourceOwnershipFence; diff --git a/packages/capture-kit/src/capture-admission/screen-recording-session-resource.ts b/packages/capture-kit/src/capture-admission/screen-recording-session-resource.ts index 66b54f4bcc..2286e814e9 100644 --- a/packages/capture-kit/src/capture-admission/screen-recording-session-resource.ts +++ b/packages/capture-kit/src/capture-admission/screen-recording-session-resource.ts @@ -16,6 +16,7 @@ import type { StopObservation } from '@agent-device/contracts/recording-stop-obs import type { DeviceInfo } from '@agent-device/kernel/device'; import type { DurableCaptureRecoveryControl, + DurableCaptureSessionBinding, DurableCaptureSessionStore, } from '../durable-capture/index.ts'; import { createDurableCaptureResource } from './durable-capture-resource.ts'; @@ -50,9 +51,7 @@ export const screenRecordingDurableResource = createDurableCaptureResource< export function adoptStartedScreenRecording(params: { admissionLedger: ScreenRecordingAdmissionLedger; - session: DurableCaptureSessionState; - sessionName: string; - sessionStore: DurableCaptureSessionStore; + binding: DurableCaptureSessionBinding<'screen-recording', ScreenRecordingLiveHandle>; device: DeviceInfo; owner: RuntimeOwnerRef; fence: ResourceOwnershipFence; diff --git a/packages/capture-kit/src/durable-capture/adoption.ts b/packages/capture-kit/src/durable-capture/adoption.ts index 2390539458..eba37b023e 100644 --- a/packages/capture-kit/src/durable-capture/adoption.ts +++ b/packages/capture-kit/src/durable-capture/adoption.ts @@ -31,43 +31,41 @@ export async function adoptStartedDurableCapture< S, >( definition: DurableCaptureResourceDefinition, - params: AdoptStartedDurableCaptureParams, + params: AdoptStartedDurableCaptureParams, resourcePath: string, ): Promise { let state: AdoptionState = { kind: 'pending' }; try { + params.binding.assertAdoptable(); 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 }), - ); + params.binding.adopt({ handle, envelope }); } catch (error) { await recoverFailedAdoption(definition, params, resourcePath, state, error); throw error; } } -async function recoverFailedAdoption, C, S>( - definition: DurableCaptureResourceDefinition, - params: AdoptStartedDurableCaptureParams, +async function recoverFailedAdoption, C>( + definition: DurableCaptureRecordDefinition, + params: AdoptStartedDurableCaptureParams, resourcePath: string, state: AdoptionState, primaryError: unknown, ): Promise { + const mayPersist = params.binding.canPersist(); const persisted = - state.kind === 'pending' ? persistRecoveryTombstone(definition, params, resourcePath) : true; + state.kind === 'pending' + ? mayPersist && persistRecoveryTombstone(definition, params, resourcePath) + : true; const initialCleanupError = await disposeFailedAdoption(params, state); - const transition = confirmFailedAdoptionTransition( - definition, - params, - resourcePath, - initialCleanupError, - ); + const transition = !params.binding.canPersist() + ? { confirmed: false, cleanupError: initialCleanupError } + : confirmFailedAdoptionTransition(definition, params, resourcePath, initialCleanupError); params.reportUndurableCleanup( params.device, (!persisted && transition.cleanupError === undefined) || transition.confirmed @@ -84,9 +82,9 @@ async function recoverFailedAdoption, C, S>( +function confirmFailedAdoptionTransition, C>( definition: DurableCaptureRecordDefinition, - params: Pick, 'sessionName' | 'fence'>, + params: Pick, 'binding' | 'fence'>, resourcePath: string, cleanupError: unknown | undefined, ): { confirmed: boolean; cleanupError: unknown | undefined } { @@ -100,7 +98,7 @@ function confirmFailedAdoptionTransition( - params: AdoptStartedDurableCaptureParams, +async function disposeFailedAdoption( + params: AdoptStartedDurableCaptureParams, state: AdoptionState, ): Promise { try { @@ -122,16 +120,13 @@ async function disposeFailedAdoption, C, S>( +function persistRecoveryTombstone, C>( definition: DurableCaptureRecordDefinition, - params: AdoptStartedDurableCaptureParams, + params: AdoptStartedDurableCaptureParams, resourcePath: string, ): boolean { try { - definition.store.write( - resourcePath, - createExpectedEnvelope(definition, params, params.envelope.descriptor), - ); + definition.store.write(resourcePath, createRecoveryEnvelope(definition, params)); return true; } catch (descriptorError) { try { @@ -148,7 +143,7 @@ function persistRecoveryTombstone, C, S>( +function createRecoveryEnvelope, C>( definition: DurableCaptureRecordDefinition, - params: Pick< - AdoptStartedDurableCaptureParams, - 'sessionName' | 'device' | 'owner' | 'fence' - >, + params: AdoptStartedDurableCaptureParams, +): DurableResourceEnvelope { + try { + return withPhase(validateStartedEnvelope(definition, params), 'active'); + } catch { + return createExpectedEnvelope(definition, params, params.envelope.descriptor); + } +} + +function createExpectedEnvelope, C>( + definition: DurableCaptureRecordDefinition, + params: Pick, 'binding' | 'device' | 'owner' | 'fence'>, descriptor: DurableResourceEnvelope['descriptor'], ): DurableResourceEnvelope { return createDurableResourceEnvelope({ resourceKind: definition.resourceKind, - sessionId: params.sessionName, + sessionId: params.binding.address, device: deviceIdentity(params.device), owner: params.owner, fence: params.fence, @@ -180,11 +183,11 @@ function createExpectedEnvelope, C, S>( +function validateStartedEnvelope, C>( definition: DurableCaptureRecordDefinition, params: Pick< - AdoptStartedDurableCaptureParams, - 'sessionName' | 'device' | 'owner' | 'fence' | 'envelope' + AdoptStartedDurableCaptureParams, + 'binding' | 'device' | 'owner' | 'fence' | 'envelope' >, ): DurableResourceEnvelope { const decoded = decodeDurableResourceEnvelope(params.envelope); @@ -207,7 +210,7 @@ function validateStartedEnvelope( return true; } -function emitCleanupDiagnostic( +function emitCleanupDiagnostic( definition: DurableCaptureRecordDefinition, - params: Pick, 'sessionName'>, + params: Pick, 'binding'>, primaryError: unknown, cleanupError: unknown, ): void { @@ -260,7 +263,7 @@ function emitCleanupDiagnostic = Readonly<{ - set(name: string, session: S): void; - resolveSessionDir(name: string): string; -}>; - export type DurableCaptureSessionResource = Readonly<{ handle: H; envelope: DurableResourceEnvelope; }>; +export type DurableCaptureSessionBinding = Readonly<{ + address: string; + sessionDir: string; + read(): DurableCaptureSessionResource | undefined; + assertAdoptable(): void; + canPersist(): boolean; + adopt(resource: DurableCaptureSessionResource): void; + clear(expected: DurableCaptureSessionResource): 'cleared' | 'retired' | 'resource-changed'; +}>; + +export type DurableCaptureSessionStore = Readonly<{ + set(name: string, session: S): void; + resolveSessionDir(name: string): string; +}>; + export type DurableCaptureSessionSlot = Readonly<{ read(session: S): DurableCaptureSessionResource | undefined; replace(session: S, resource: DurableCaptureSessionResource | undefined): S; @@ -86,11 +91,9 @@ export type DurableCaptureCleanupOutcome = | { confirmed: true } | { confirmed: false; reason: string }; -export type AdoptStartedDurableCaptureParams = { +export type AdoptStartedDurableCaptureParams = { reportUndurableCleanup(device: DeviceInfo, outcome: DurableCaptureCleanupOutcome): void; - session: S; - sessionName: string; - sessionStore: DurableCaptureSessionStore; + binding: DurableCaptureSessionBinding; device: DeviceInfo; owner: RuntimeOwnerRef; fence: ResourceOwnershipFence; diff --git a/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts b/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts index 61bdf5d280..fcb88e857c 100644 --- a/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts +++ b/packages/capture-kit/src/durable-capture/durable-capture.fixtures.ts @@ -1,3 +1,4 @@ +import { makeCaptureSessionBinding } from './session-binding.fixtures.ts'; import path from 'node:path'; import { vi, type Mock } from 'vitest'; import { @@ -15,7 +16,6 @@ import type { DurableCaptureFailedFinishPolicy, DurableCaptureResourceDefinition, DurableCaptureSessionResource, - DurableCaptureSessionStore, } from './definition.ts'; import { createDurableCaptureResourceStore, type DurableCaptureResourceStore } from './store.ts'; @@ -91,8 +91,9 @@ export function makeDurableCaptureContext( const session: TestCaptureSession = { name: sessionName }; sessions.set(sessionName, session); const resolveSessionDir = (name: string): string => path.join(sessionsDir, name); - const sessionStore: DurableCaptureSessionStore = { - set: (name, next) => void sessions.set(name, next), + const sessionStore = { + get: (name: string) => sessions.get(name), + set: (name: string, next: TestCaptureSession) => void sessions.set(name, next), resolveSessionDir, }; const reportUndurableCleanup: Mock< @@ -100,6 +101,11 @@ export function makeDurableCaptureContext( > = vi.fn(); return { reportUndurableCleanup, + binding: makeCaptureSessionBinding( + sessionStore, + sessionName, + testCaptureDefinition.sessionSlot, + ), sessions, sessionsDir, resolveSessionDir, diff --git a/packages/capture-kit/src/durable-capture/index.ts b/packages/capture-kit/src/durable-capture/index.ts index 09013e180a..dc4349e202 100644 --- a/packages/capture-kit/src/durable-capture/index.ts +++ b/packages/capture-kit/src/durable-capture/index.ts @@ -14,6 +14,7 @@ export type { DurableCaptureRecordDefinition, DurableCaptureResourceDefinition, DurableCaptureSessionResource, + DurableCaptureSessionBinding, DurableCaptureSessionStore, } from './definition.ts'; export type { FinishRecoveredDurableCaptureParams } from './finish-recovered.ts'; diff --git a/packages/capture-kit/src/durable-capture/session-binding.fixtures.ts b/packages/capture-kit/src/durable-capture/session-binding.fixtures.ts new file mode 100644 index 0000000000..611d9b9108 --- /dev/null +++ b/packages/capture-kit/src/durable-capture/session-binding.fixtures.ts @@ -0,0 +1,45 @@ +import { AppError } from '@agent-device/kernel/errors'; +import type { DurableCaptureSessionBinding, DurableCaptureSessionSlot } from './definition.ts'; + +export function makeCaptureSessionBinding( + store: Readonly<{ + get(name: string): S | undefined; + set(name: string, session: S): void; + resolveSessionDir(name: string): string; + }>, + address: string, + slot: DurableCaptureSessionSlot, +): DurableCaptureSessionBinding { + const requireSession = (): S => { + const session = store.get(address); + if (!session) throw new AppError('COMMAND_FAILED', 'Test session retired'); + return session; + }; + const assertAdoptable = (): void => { + if (slot.read(requireSession())) throw new AppError('COMMAND_FAILED', 'Test resource changed'); + }; + return Object.freeze({ + address, + sessionDir: store.resolveSessionDir(address), + read: () => { + const session = store.get(address); + return session === undefined ? undefined : slot.read(session); + }, + assertAdoptable, + canPersist: () => { + const session = store.get(address); + return session !== undefined && slot.read(session) === undefined; + }, + adopt: (resource) => { + assertAdoptable(); + store.set(address, slot.replace(requireSession(), resource)); + }, + clear: (expected) => { + const current = store.get(address); + if (!current) return 'retired'; + if (slot.read(current)?.handle !== expected.handle) return 'resource-changed'; + store.set(address, slot.replace(current, undefined)); + return 'cleared'; + }, + }); +} diff --git a/packages/platform-android/src/recording/failed-finish.test.ts b/packages/platform-android/src/recording/failed-finish.test.ts index 4c4984020c..d3f7afb0aa 100644 --- a/packages/platform-android/src/recording/failed-finish.test.ts +++ b/packages/platform-android/src/recording/failed-finish.test.ts @@ -13,7 +13,6 @@ import { createDurableCaptureResourceStore, finishLiveDurableCapture, type DurableCaptureResourceDefinition, - type DurableCaptureSessionStore, } from '@agent-device/capture-kit/durable-capture'; import { mkdtempForTestSync } from '../__tests__/test-utils/tmp-dir.ts'; import { androidRecordingDevice, recordingHost, recordingInput } from './fixtures.ts'; @@ -150,20 +149,36 @@ async function adoptAndroidRecording(params: { }; const sessionsDir = mkdtempForTestSync('agent-device-android-failed-finish-session-'); let session: AndroidRecordingSession = {}; - const sessionStore: DurableCaptureSessionStore = { - set: (_name, next) => { + const sessionStore = { + get: () => session, + set: (_name: string, next: AndroidRecordingSession) => { session = next; }, - resolveSessionDir: (name) => path.join(sessionsDir, name), + resolveSessionDir: (name: string) => path.join(sessionsDir, name), }; - const resourcePath = store.resolvePath(sessionStore.resolveSessionDir(params.sessionName)); + const binding = { + address: params.sessionName, + sessionDir: sessionStore.resolveSessionDir(params.sessionName), + read: () => session.recording, + assertAdoptable: () => { + if (session.recording) throw new Error('Already recording'); + }, + canPersist: () => !session.recording, + adopt: (recording: AndroidRecordingSession['recording']) => { + session = { ...session, recording }; + }, + clear: (expected: NonNullable) => { + if (session.recording?.handle !== expected.handle) return 'resource-changed' as const; + session = { ...session, recording: undefined }; + return 'cleared' as const; + }, + }; + const resourcePath = store.resolvePath(binding.sessionDir); await adoptStartedDurableCapture( definition, { reportUndurableCleanup: () => {}, - session, - sessionName: params.sessionName, - sessionStore, + binding, device: androidRecordingDevice, owner: params.owner, fence: params.envelope.fence, diff --git a/scripts/layering/session-resource-ownership.test.ts b/scripts/layering/session-resource-ownership.test.ts index ffac560315..4b2e7fc2b8 100644 --- a/scripts/layering/session-resource-ownership.test.ts +++ b/scripts/layering/session-resource-ownership.test.ts @@ -19,6 +19,7 @@ test('session resources are constructed only by their durable domain owners', () appLogFailure: failure, audioProbe: audio, perfCapture: perf, + screenRecording: recording, });`, ], [ @@ -39,6 +40,7 @@ test('session resources are constructed only by their durable domain owners', () 'src/daemon/handlers/planted.ts: session appLogFailure record constructed outside its owner', 'src/daemon/handlers/planted.ts: session audioProbe record constructed outside its owner', 'src/daemon/handlers/planted.ts: session perfCapture record constructed outside its owner', + 'src/daemon/handlers/planted.ts: session screenRecording record constructed outside its owner', ], ); }); diff --git a/scripts/layering/session-resource-ownership.ts b/scripts/layering/session-resource-ownership.ts index 4990fdef9e..579dd3d637 100644 --- a/scripts/layering/session-resource-ownership.ts +++ b/scripts/layering/session-resource-ownership.ts @@ -29,10 +29,17 @@ const RESOURCE_OWNERS: Readonly>> = { appLogFailure: new Set(['src/daemon/app-log-session-resource.ts', 'src/daemon/session-state.ts']), audioProbe: new Set([ 'packages/capture-kit/src/capture-admission/audio-probe-session-resource.ts', + 'src/daemon/audio-probe-session-binding.ts', + 'src/daemon/session-state.ts', + ]), + screenRecording: new Set([ + 'packages/capture-kit/src/capture-admission/screen-recording-session-resource.ts', + 'src/daemon/screen-recording-session-binding.ts', 'src/daemon/session-state.ts', ]), perfCapture: new Set([ 'packages/capture-kit/src/capture-admission/perf-capture-session-resource.ts', + 'src/daemon/perf-capture-session-binding.ts', 'src/daemon/session-state.ts', ]), }; diff --git a/src/daemon/__tests__/app-log-session-resource.test.ts b/src/daemon/__tests__/app-log-session-resource.test.ts index 84a3add4f2..eb9402d3ba 100644 --- a/src/daemon/__tests__/app-log-session-resource.test.ts +++ b/src/daemon/__tests__/app-log-session-resource.test.ts @@ -88,7 +88,7 @@ test('SessionStore failure after transfer disposes the transferred handle and pr const context = makeContext(); const runtime = makeStartResult(context); const primary = new Error('store adoption failed'); - vi.spyOn(context.sessionStore, 'set').mockImplementationOnce(() => { + vi.spyOn(context.sessionStore, 'update').mockImplementationOnce(() => { throw primary; }); await expect( @@ -315,6 +315,77 @@ test('app-log disposes on a failed finish because its retry is that same finish ).not.toThrow(); }); +test('shutdown admission rejects a late start while its existing session still occupies the address', async () => { + const context = makeContext(); + const runtime = makeStartResult(context); + context.sessionStore.closeAdmission(); + await expect( + adoptStartedSessionAppLog({ ...context, ...runtime.result, throwIfCanceled: () => {} }), + ).rejects.toMatchObject({ details: { reason: 'daemon_shutting_down' } }); + expect(runtime.forceCleanup).toHaveBeenCalledOnce(); + expect(context.sessionStore.requireCurrent(context.ref).appLog).toBeUndefined(); + expect(appLogResourceStore.read(context.resourcePath)).toMatchObject({ + status: 'decoded', + envelope: { lifecycle: 'completed' }, + }); +}); + +test('failed adoption cannot terminalize successor evidence after its cleanup yields', async () => { + const context = makeContext(); + const runtime = makeStartResult(context); + let release!: () => void; + let entered!: () => void; + const held = new Promise((resolve) => { + release = resolve; + }); + const cleaning = new Promise((resolve) => { + entered = resolve; + }); + runtime.forceCleanup.mockImplementationOnce(async () => { + entered(); + await held; + return { status: 'cleaned' }; + }); + const canceled = new AppError('CANCELED', 'canceled'); + const adoption = adoptStartedSessionAppLog({ + ...context, + ...runtime.result, + throwIfCanceled: () => { + throw canceled; + }, + }); + const rejected = expect(adoption).rejects.toBe(canceled); + await cleaning; + context.sessionStore.retire(context.ref); + const successor = context.sessionStore.publish(context.sessionName, { ...context.session }); + appLogResourceStore.write(context.resourcePath, runtime.result.envelope); + release(); + await rejected; + expect(context.sessionStore.requireCurrent(successor).appLog).toBeUndefined(); + expect(appLogResourceStore.read(context.resourcePath)).toMatchObject({ + status: 'decoded', + envelope: { lifecycle: 'open' }, + }); +}); + +test('late adoption disposes its pending handle without overwriting a successor manifest', async () => { + const context = makeContext(); + const runtime = makeStartResult(context); + context.sessionStore.retire(context.ref); + const successor = context.sessionStore.publish(context.sessionName, { ...context.session }); + const envelope = { ...runtime.result.envelope, fence: { token: 'successor', generation: 2 } }; + appLogResourceStore.write(context.resourcePath, envelope); + await expect( + adoptStartedSessionAppLog({ ...context, ...runtime.result, throwIfCanceled: () => {} }), + ).rejects.toMatchObject({ details: { reason: 'session_lifetime_ended' } }); + expect(runtime.forceCleanup).toHaveBeenCalledOnce(); + expect(context.sessionStore.requireCurrent(successor).appLog).toBeUndefined(); + expect(appLogResourceStore.read(context.resourcePath)).toMatchObject({ + status: 'decoded', + envelope, + }); +}); + function makeContext( device: DeviceInfo = { platform: 'android', @@ -335,6 +406,7 @@ function makeContext( const resourcePath = appLogResourceStore.resolvePath(sessionStore.resolveSessionDir(sessionName)); return { admissionLedger: createAppLogAdmissionLedger(), + ref: sessionStore.lookup(sessionName)!, session, sessionName, sessionStore, diff --git a/src/daemon/__tests__/perf-capture-session-resource.test.ts b/src/daemon/__tests__/perf-capture-session-resource.test.ts index cd4010c381..6d7d4e8841 100644 --- a/src/daemon/__tests__/perf-capture-session-resource.test.ts +++ b/src/daemon/__tests__/perf-capture-session-resource.test.ts @@ -1,3 +1,4 @@ +import { bindSessionPerfCapture } from '../perf-capture-session-binding.ts'; import { beforeEach, expect, test, vi } from 'vitest'; import { localRuntimeOwner } from '@agent-device/contracts/platform-runtime'; import { AppError } from '@agent-device/kernel/errors'; @@ -112,9 +113,7 @@ test('a perf stop whose pull failed re-collects the device-side trace the first }); await adoptStartedPerfCapture({ admissionLedger: createPerfCaptureAdmissionLedger(), - session, - sessionName, - sessionStore, + binding: bindSessionPerfCapture(sessionStore, sessionStore.lookup(sessionName)!), device, owner: localRuntimeOwner('android'), fence, diff --git a/src/daemon/__tests__/screen-recording-session-binding.test.ts b/src/daemon/__tests__/screen-recording-session-binding.test.ts new file mode 100644 index 0000000000..520b6b3452 --- /dev/null +++ b/src/daemon/__tests__/screen-recording-session-binding.test.ts @@ -0,0 +1,60 @@ +import { expect, test, vi } from 'vitest'; +import { createScreenRecordingLiveHandle } from '@agent-device/capture-kit'; +import { PendingTransferGuard } from '@agent-device/contracts/async-lifecycle'; +import { makeSessionStore } from '../../__tests__/test-utils/store-factory.ts'; +import { makeRecordingSession } from './session-teardown.fixtures.ts'; +import { bindRecordOnlyScreenRecording } from '../screen-recording-session-binding.ts'; +import { createScreenRecordingAdmissionLedger } from '@agent-device/capture-kit/screen-recording-admission-ledger'; +import { + adoptStartedScreenRecording, + screenRecordingDurableResource, +} from '@agent-device/capture-kit/screen-recording-session-resource'; + +test('shutdown refuses draft publication while retaining unconfirmed recording cleanup evidence', async () => { + const store = makeSessionStore(); + const session = makeRecordingSession({ + name: 'draft', + sessionStore: store, + finish: async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }), + }); + const { handle: initial, envelope } = session.screenRecording!; + const cleanup = vi.fn( + async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }) as const, + ); + const handle = createScreenRecordingLiveHandle(initial.inspect(), { + finish: async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }), + forceCleanup: cleanup, + }); + const draft = bindRecordOnlyScreenRecording(store, 'draft', { + ...session, + screenRecording: undefined, + }); + store.closeAdmission(); + await expect( + adoptStartedScreenRecording({ + binding: draft.binding, + admissionLedger: createScreenRecordingAdmissionLedger(), + device: session.device, + owner: envelope.owner, + fence: envelope.fence, + pendingHandle: new PendingTransferGuard(handle), + envelope, + throwIfCanceled: () => {}, + }), + ).rejects.toMatchObject({ details: { reason: 'daemon_shutting_down' } }); + expect(cleanup).toHaveBeenCalledOnce(); + expect(store.lookup('draft')).toBeUndefined(); + const record = screenRecordingDurableResource.store.read( + screenRecordingDurableResource.store.resolvePath(draft.binding.sessionDir), + ); + expect(record).toMatchObject({ + status: 'decoded', + envelope: { + lifecycle: 'open', + descriptor: envelope.descriptor, + metadata: { phase: 'cleanup-pending' }, + }, + }); + if (record.status !== 'decoded') throw new Error('Expected recovery evidence'); + expect(record.envelope.metadata?.runtimeContractInvalid).toBeUndefined(); +}); diff --git a/src/daemon/__tests__/session-capture-binding.test.ts b/src/daemon/__tests__/session-capture-binding.test.ts new file mode 100644 index 0000000000..41f5322d5b --- /dev/null +++ b/src/daemon/__tests__/session-capture-binding.test.ts @@ -0,0 +1,72 @@ +import { expect, test } from 'vitest'; +import { makeSessionStore } from '../../__tests__/test-utils/store-factory.ts'; +import { makeRecordingSession } from './session-teardown.fixtures.ts'; +import { bindSessionScreenRecording } from '../screen-recording-session-binding.ts'; + +test('clearing a capture refreshes a rebuilt record without losing its other changes', () => { + const store = makeSessionStore(); + const session = makeRecordingSession({ + name: 'capture', + sessionStore: store, + finish: async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }), + }); + const ref = store.publish('capture', session); + const binding = bindSessionScreenRecording(store, ref); + const active = binding.read()!; + store.update(ref, { appName: 'updated', screenRecording: { ...active } }); + expect(binding.clear(active)).toBe('cleared'); + expect(store.requireCurrent(ref)).toMatchObject({ + appName: 'updated', + screenRecording: undefined, + }); +}); + +test('clearing an older handle or fence leaves a replacement capture intact', () => { + const store = makeSessionStore(); + const session = makeRecordingSession({ + name: 'capture', + sessionStore: store, + finish: async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }), + }); + const ref = store.publish('capture', session); + const binding = bindSessionScreenRecording(store, ref); + const active = binding.read()!; + const replacement = makeRecordingSession({ + name: 'other', + sessionStore: store, + finish: async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }), + }).screenRecording!; + store.update(ref, { screenRecording: replacement }); + expect(binding.clear(active)).toBe('resource-changed'); + expect(binding.read()).toBe(replacement); + const newerFence = { + ...active, + envelope: { ...active.envelope, fence: { token: 'next', generation: 2 } }, + }; + store.update(ref, { screenRecording: newerFence }); + expect(binding.clear(active)).toBe('resource-changed'); + expect(binding.read()).toBe(newerFence); +}); + +test('a retired binding retains its old resource but cannot write into the next lifetime', () => { + const store = makeSessionStore(); + const session = makeRecordingSession({ + name: 'capture', + sessionStore: store, + finish: async () => ({ status: 'cleanup-pending', reason: 'cleanup-unconfirmed' }), + }); + const ref = store.publish('capture', session); + const binding = bindSessionScreenRecording(store, ref); + const active = binding.read()!; + store.retire(ref); + const successor = store.publish('capture', session); + expect(binding.read()).toBe(active); + expect(binding.clear(active)).toBe('retired'); + expect(binding.canPersist()).toBe(false); + expect(() => binding.adopt(active)).toThrow( + expect.objectContaining({ + details: expect.objectContaining({ reason: 'session_lifetime_ended' }), + }), + ); + expect(store.requireCurrent(successor).screenRecording).toBe(active); +}); diff --git a/src/daemon/app-log-session-resource.ts b/src/daemon/app-log-session-resource.ts index e86d50f165..5c3993e19b 100644 --- a/src/daemon/app-log-session-resource.ts +++ b/src/daemon/app-log-session-resource.ts @@ -15,7 +15,8 @@ import { } from '@agent-device/capture-kit/durable-capture-resource'; import { appLogResourceStore } from './app-log-resource-store.ts'; import type { SessionStore } from './session-store.ts'; -import type { SessionState } from './session-state.ts'; +import type { SessionRef, SessionState } from './session-state.ts'; +import { bindSessionCapture } from './session-capture-binding.ts'; export type AppLogSessionSnapshot = Readonly<{ active: boolean; @@ -76,10 +77,8 @@ export function inspectSessionAppLog(session: SessionState): AppLogSessionSnapsh export function adoptStartedSessionAppLog(params: { admissionLedger: AppLogAdmissionLedger; - session: SessionState; - sessionName: string; + ref: SessionRef; sessionStore: SessionStore; - resourcePath: string; device: DeviceInfo; owner: RuntimeOwnerRef; fence: ResourceOwnershipFence; @@ -87,7 +86,10 @@ export function adoptStartedSessionAppLog(params: { envelope: DurableResourceEnvelope<'app-log'>; throwIfCanceled(): void; }): Promise { - return appLogDurableResource.adoptStarted(params); + return appLogDurableResource.adoptStarted({ + ...params, + binding: bindSessionAppLog(params.sessionStore, params.ref), + }); } export function finishSessionAppLog(params: { @@ -110,15 +112,15 @@ export function forceCleanupSessionAppLog(params: { } export function recordSessionAppLogFailure(params: { - session: SessionState; - sessionName: string; + ref: SessionRef; sessionStore: SessionStore; error: unknown; backend?: LogBackend; }): ReturnType { const normalized = normalizeError(params.error); - params.sessionStore.set(params.sessionName, { - ...params.session, + const current = params.sessionStore.resolveCurrent(params.ref); + if (!current || current.appLog) return normalized; + params.sessionStore.update(params.ref, { appLog: undefined, appLogFailure: { backend: params.backend, @@ -131,12 +133,19 @@ export function recordSessionAppLogFailure(params: { } export function clearSessionAppLogFailure(params: { - session: SessionState; - sessionName: string; + ref: SessionRef; sessionStore: SessionStore; }): void { - params.sessionStore.set(params.sessionName, { - ...params.session, + params.sessionStore.update(params.ref, { appLogFailure: undefined, }); } + +function bindSessionAppLog(sessionStore: SessionStore, ref: SessionRef) { + return bindSessionCapture(sessionStore, ref, { + read: (session) => session.appLog, + write: (appLog) => { + sessionStore.update(ref, { appLog, appLogFailure: undefined }); + }, + }); +} diff --git a/src/daemon/audio-probe-session-binding.ts b/src/daemon/audio-probe-session-binding.ts new file mode 100644 index 0000000000..b4eecb43c7 --- /dev/null +++ b/src/daemon/audio-probe-session-binding.ts @@ -0,0 +1,12 @@ +import { bindSessionCapture } from './session-capture-binding.ts'; +import type { SessionRef } from './session-state.ts'; +import type { SessionStore } from './session-store.ts'; + +export function bindSessionAudioProbe(sessionStore: SessionStore, ref: SessionRef) { + return bindSessionCapture(sessionStore, ref, { + read: (session) => session.audioProbe, + write: (audioProbe) => { + sessionStore.update(ref, { audioProbe }); + }, + }); +} diff --git a/src/daemon/handlers/record-runtime.ts b/src/daemon/handlers/record-runtime.ts index 758bd95bc4..1e63e684cb 100644 --- a/src/daemon/handlers/record-runtime.ts +++ b/src/daemon/handlers/record-runtime.ts @@ -30,7 +30,11 @@ import { resolveSessionScope } from '../session-routing.ts'; import type { SessionStore } from '../session-store.ts'; import type { BindDeviceRuntime, BindExactDeviceRuntime } from '../request-runtime-binding.ts'; import type { DaemonRequest, DaemonResponse } from '../daemon-request.ts'; -import type { SessionState } from '../session-state.ts'; +import type { SessionRef, SessionState } from '../session-state.ts'; +import { + bindRecordOnlyScreenRecording, + bindSessionScreenRecording, +} from '../screen-recording-session-binding.ts'; import { recordSessionAction } from '../session-action-recorder.ts'; import { missingAppSessionResponse, @@ -80,17 +84,19 @@ async function handleRecordCommandUnsafe( params: RecordRuntimeHandlerParams, ): Promise { const { req, sessionName, sessionStore } = params; - const existingSession = sessionStore.get(sessionName); + const existingRef = sessionStore.lookup(sessionName); + const existingSession = existingRef?.session; const { plan, scope } = resolveRecordPlan(req, existingSession); if (plan.kind === 'start' && !isWholeScreenRecordingScope(scope) && !existingSession) { return missingAppSessionResponse(req); } - const resolvedSession = await resolveRecordingSession(params, existingSession); + const resolvedSession = await resolveRecordingSession(params, existingRef); const { session } = resolvedSession; if (plan.kind === 'start') { return await startRecording( params, session, + resolvedSession.ref, prepareRecordingRequest(req), plan.use, resolvedSession.needsReadiness, @@ -113,17 +119,18 @@ function resolveRecordPlan(req: DaemonRequest, session: SessionState | undefined async function resolveRecordingSession( params: RecordRuntimeHandlerParams, - existing: SessionState | undefined, -): Promise> { - const device = existing?.device ?? (await resolveTargetDevice(params.req.flags ?? {})); + ref: SessionRef | undefined, +): Promise> { + const device = ref?.session.device ?? (await resolveTargetDevice(params.req.flags ?? {})); await params.retainDeviceExecutionLock(device.id); - if (existing) return { session: existing, needsReadiness: false }; + if (ref) return { session: params.sessionStore.requireCurrent(ref), ref, needsReadiness: false }; return { session: createRecordOnlySession(params, device), needsReadiness: true }; } async function startRecording( params: RecordRuntimeHandlerParams, session: SessionState, + ref: SessionRef | undefined, prepared: ReturnType, use: typeof screenRecordingStartUse, needsReadiness: boolean, @@ -131,6 +138,11 @@ async function startRecording( if (session.screenRecording) { return { ok: false, error: { code: 'INVALID_ARGS', message: 'recording already in progress' } }; } + const draft = ref + ? undefined + : bindRecordOnlyScreenRecording(params.sessionStore, params.sessionName, session); + const binding = ref ? bindSessionScreenRecording(params.sessionStore, ref) : draft!.binding; + binding.assertAdoptable(); const admission = await params.bindDevice(session.device, screenRecordingAdmissionUse); if (needsReadiness) await ensureBoundDeviceReady(admission); const startFact = admission.facts.screenRecordingStart; @@ -142,21 +154,20 @@ async function startRecording( ); await adoptStartedScreenRecording({ admissionLedger: params.admissionLedger, - session, - sessionName: params.sessionName, - sessionStore: params.sessionStore, + binding, device: session.device, owner: runtime.owner, fence, ...started, throwIfCanceled: params.throwIfCanceled, }); - const adopted = params.sessionStore.get(params.sessionName)?.screenRecording; + const adoptedRef = ref ?? draft!.requireRef(); + const adopted = binding.read(); if (!adopted) throw new TypeError('Screen recording adoption did not publish a live handle'); const snapshot = adopted.handle.inspect(); recordSessionAction( params.sessionStore, - session, + params.sessionStore.requireCurrent(adoptedRef), params.req, params.req.command, buildRecordingStartedAction(snapshot), diff --git a/src/daemon/perf-capture-session-binding.ts b/src/daemon/perf-capture-session-binding.ts new file mode 100644 index 0000000000..53463cd8ab --- /dev/null +++ b/src/daemon/perf-capture-session-binding.ts @@ -0,0 +1,12 @@ +import { bindSessionCapture } from './session-capture-binding.ts'; +import type { SessionRef } from './session-state.ts'; +import type { SessionStore } from './session-store.ts'; + +export function bindSessionPerfCapture(sessionStore: SessionStore, ref: SessionRef) { + return bindSessionCapture(sessionStore, ref, { + read: (session) => session.perfCapture, + write: (perfCapture) => { + sessionStore.update(ref, { perfCapture }); + }, + }); +} diff --git a/src/daemon/screen-recording-session-binding.ts b/src/daemon/screen-recording-session-binding.ts new file mode 100644 index 0000000000..d032f819ba --- /dev/null +++ b/src/daemon/screen-recording-session-binding.ts @@ -0,0 +1,42 @@ +import { bindSessionCapture } from './session-capture-binding.ts'; +import type { SessionRef } from './session-state.ts'; +import type { SessionStore } from './session-store.ts'; + +export function bindSessionScreenRecording(sessionStore: SessionStore, ref: SessionRef) { + return bindSessionCapture(sessionStore, ref, { + read: (session) => session.screenRecording, + write: (screenRecording) => { + sessionStore.update(ref, { screenRecording }); + }, + }); +} + +export function bindRecordOnlyScreenRecording( + sessionStore: SessionStore, + address: string, + draft: SessionRef['session'], +) { + let published: SessionRef | undefined; + const binding: ReturnType = Object.freeze({ + address, + sessionDir: sessionStore.resolveSessionDir(address), + read: () => + published ? bindSessionScreenRecording(sessionStore, published).read() : undefined, + assertAdoptable: () => sessionStore.assertPublishable(address), + canPersist: () => !published && sessionStore.lookup(address) === undefined, + adopt: (screenRecording) => { + sessionStore.assertPublishable(address); + draft.screenRecording = screenRecording; + published = sessionStore.publish(address, draft); + }, + clear: (expected) => + published ? bindSessionScreenRecording(sessionStore, published).clear(expected) : 'retired', + }); + return Object.freeze({ + binding, + requireRef: (): SessionRef => { + if (!published) throw new TypeError('Screen recording did not publish its session'); + return published; + }, + }); +} diff --git a/src/daemon/session-capture-binding.ts b/src/daemon/session-capture-binding.ts new file mode 100644 index 0000000000..c3f809cf37 --- /dev/null +++ b/src/daemon/session-capture-binding.ts @@ -0,0 +1,53 @@ +import type { + DurableCaptureSessionBinding, + DurableCaptureSessionResource, +} from '@agent-device/capture-kit/durable-capture'; +import { AppError } from '@agent-device/kernel/errors'; +import type { SessionRef, SessionState } from './session-state.ts'; +import type { SessionStore } from './session-store.ts'; + +export function bindSessionCapture( + sessionStore: SessionStore, + ref: SessionRef, + slot: Readonly<{ + read(session: SessionState): DurableCaptureSessionResource | undefined; + write(resource: DurableCaptureSessionResource | undefined): void; + }>, +): DurableCaptureSessionBinding { + const assertAdoptable = (): void => { + sessionStore.assertAdmissionOpen(ref.address); + if (slot.read(sessionStore.requireCurrent(ref))) { + throw new AppError('COMMAND_FAILED', 'Session capture resource has changed', { + reason: 'session_resource_changed', + session: ref.address, + }); + } + }; + return Object.freeze({ + address: ref.address, + sessionDir: sessionStore.resolveSessionDir(ref.address), + read: () => slot.read(sessionStore.resolveCurrent(ref) ?? ref.session), + assertAdoptable, + canPersist: () => { + const current = sessionStore.resolveCurrent(ref); + return current !== undefined && slot.read(current) === undefined; + }, + adopt: (resource) => { + assertAdoptable(); + slot.write(resource); + }, + clear: (expected) => { + const current = sessionStore.resolveCurrent(ref); + if (!current) return 'retired'; + const active = slot.read(current); + if ( + active?.handle !== expected.handle || + active.envelope.fence.token !== expected.envelope.fence.token || + active.envelope.fence.generation !== expected.envelope.fence.generation + ) + return 'resource-changed'; + slot.write(undefined); + return 'cleared'; + }, + }); +} diff --git a/src/daemon/session-observability/internal/__tests__/session-logs.test.ts b/src/daemon/session-observability/internal/__tests__/session-logs.test.ts index fb40412fe7..f356a1c7a2 100644 --- a/src/daemon/session-observability/internal/__tests__/session-logs.test.ts +++ b/src/daemon/session-observability/internal/__tests__/session-logs.test.ts @@ -269,7 +269,7 @@ test('rejected pending cleanup retains cleanup-pending record and blocks replace test('post-transfer SessionStore failure disposes the transferred handle and preserves primary error', async () => { const { sessionStore, sessionName } = openSession(); const primary = new Error('store adoption failed'); - vi.spyOn(sessionStore, 'set').mockImplementationOnce(() => { + vi.spyOn(sessionStore, 'update').mockImplementationOnce(() => { throw primary; }); const response = await runLogs(sessionStore, sessionName, ['start'], {}, runtime.bindDevice); diff --git a/src/daemon/session-observability/internal/session-audio.ts b/src/daemon/session-observability/internal/session-audio.ts index a7e60b98c7..6d439e65b5 100644 --- a/src/daemon/session-observability/internal/session-audio.ts +++ b/src/daemon/session-observability/internal/session-audio.ts @@ -20,7 +20,8 @@ import type { } from '../../request-runtime-binding.ts'; import type { SessionStore } from '../../session-store.ts'; import type { DaemonRequest, DaemonResponse } from '../../daemon-request.ts'; -import type { SessionState } from '../../session-state.ts'; +import type { SessionRef, SessionState } from '../../session-state.ts'; +import { bindSessionAudioProbe } from '../../audio-probe-session-binding.ts'; import { type DaemonFailureResponse, errorResponse } from '@agent-device/kernel/contracts'; type AudioParams = { @@ -44,7 +45,8 @@ export async function handleAudioCommand(params: AudioParams): Promise { const sessionResult = resolveAudioSession(params); if (!sessionResult.ok) return sessionResult; - const session = sessionResult.session; + const ref = sessionResult.ref; + const session = params.sessionStore.requireCurrent(ref); const request = parseAudioProbeRequest(params.req.positionals); // Facts, not a capability bucket, decide which of the two owner paths this device has โ€” // side-effect-free per ADR 0019 ยง9; the one bind below uses the plan's own use. @@ -67,7 +69,7 @@ async function handleAudioCommandUnsafe(params: AudioParams): Promise { + let session = params.sessionStore.requireCurrent(ref); + const binding = bindSessionAudioProbe(params.sessionStore, ref); // Start restarts an already-running probe (legacy parity), completing it through the durable // coordinator so the previous envelope terminalizes before a new fence is minted. Nobody reads // that completion, so the previous probe is being handed back rather than captured. @@ -110,12 +114,10 @@ async function startAudioProbe( await finishLiveAudioProbe({ intent: 'disposal', session, - sessionName: params.sessionName, + sessionName: ref.address, sessionStore: params.sessionStore, }); - const refreshed = params.sessionStore.get(params.sessionName); - if (!refreshed) return errorResponse('SESSION_NOT_FOUND', 'audio requires an active session'); - session = refreshed; + session = params.sessionStore.requireCurrent(ref); } const runtime = await params.bindDevice(session.device, use); const resourcePath = audioProbeDurableResource.store.resolvePath( @@ -139,16 +141,14 @@ async function startAudioProbe( }); await adoptStartedAudioProbe({ admissionLedger: params.audioProbeAdmissionLedger, - session, - sessionName: params.sessionName, - sessionStore: params.sessionStore, + binding, device: session.device, owner: runtime.owner, fence, ...started, throwIfCanceled: params.throwIfCanceled, }); - const adopted = params.sessionStore.get(params.sessionName)?.audioProbe; + const adopted = binding.read(); if (!adopted) throw new TypeError('Audio probe adoption did not publish a live handle'); return { ok: true, data: await adopted.handle.status() }; } diff --git a/src/daemon/session-observability/internal/session-observability.ts b/src/daemon/session-observability/internal/session-observability.ts index 2171842c39..2e98bd25de 100644 --- a/src/daemon/session-observability/internal/session-observability.ts +++ b/src/daemon/session-observability/internal/session-observability.ts @@ -31,7 +31,7 @@ import type { } from '../../request-runtime-binding.ts'; import type { SessionStore } from '../../session-store.ts'; import type { DaemonRequest, DaemonResponse } from '../../daemon-request.ts'; -import type { SessionState } from '../../session-state.ts'; +import type { SessionRef, SessionState } from '../../session-state.ts'; import { handleAudioCommand } from './session-audio.ts'; import { handlePerfRuntimeCommand } from './session-perf-runtime.ts'; import { handleNetworkCommand } from './session-network.ts'; @@ -51,6 +51,7 @@ export type SessionObservabilityCommandInput = { type ObservabilityInput = SessionObservabilityCommandInput; type LogsHandlerParams = Omit & { session: SessionState; + ref: SessionRef; bindDevice: BindDeviceRuntime; appLogAdmissionLedger: AppLogAdmissionLedger; }; @@ -143,12 +144,13 @@ async function handleEventsCommand(params: ObservabilityInput): Promise { const { req, sessionName, sessionStore } = params; - const session = sessionStore.get(sessionName); - if (!session) { + const ref = sessionStore.lookup(sessionName); + if (!ref) { return errorResponse('SESSION_NOT_FOUND', 'logs requires an active session'); } try { - const logsParams = requireLogsHandlerParams({ ...params, session }); + const session = sessionStore.requireCurrent(ref); + const logsParams = requireLogsHandlerParams({ ...params, session, ref }); const admission = await logsParams.bindDevice(session.device, appLogAdmissionUse); const inspectFact = admission.facts.appLogInspect; if (!inspectFact.available) { @@ -305,7 +307,7 @@ function handleLogsClear(params: LogsHandlerParams): DaemonResponse { } const logPath = sessionStore.resolveAppLogPath(sessionName); const cleared = clearAppLogFiles(logPath); - clearSessionAppLogFailure({ session, sessionName, sessionStore }); + clearSessionAppLogFailure({ ref: params.ref, sessionStore }); return { ok: true, data: cleared }; } @@ -386,10 +388,8 @@ async function startSessionAppLog( }); await adoptStartedSessionAppLog({ admissionLedger: params.appLogAdmissionLedger, - session, - sessionName, + ref: params.ref, sessionStore, - resourcePath, device: session.device, owner, fence, @@ -400,8 +400,7 @@ async function startSessionAppLog( return { ok: true, data: { path: outputPath, started: true } }; } catch (error) { const normalized = recordSessionAppLogFailure({ - session, - sessionName, + ref: params.ref, sessionStore, error, }); @@ -431,7 +430,7 @@ function requireAudioSeams(params: ObservabilityInput): Parameters { - const session = params.sessionStore.get(params.sessionName); - if (!session) { + const ref = params.sessionStore.lookup(params.sessionName); + if (!ref) { return errorResponse('SESSION_NOT_FOUND', 'perf requires an active session. Run open first.'); } + const session = params.sessionStore.requireCurrent(ref); + const bound = { ...params, ref }; try { if (isRemovedAggregatePerfToken(params.req.positionals?.[0])) { throw new AppError('INVALID_ARGS', PERF_AGGREGATE_REMOVED_ERROR_MESSAGE); @@ -89,7 +92,7 @@ export async function handlePerfRuntimeCommand( } return recordSuccessfulPerfResponse( params, - await executeAdmittedPerfPlan(params, session, admitted), + await executeAdmittedPerfPlan(bound, session, admitted), ); } catch (error) { return { ok: false, error: normalizeError(error) }; @@ -116,7 +119,7 @@ function recordSuccessfulPerfResponse( // the admission/runtime join this handler is meant to keep singular. // fallow-ignore-next-line complexity async function executeAdmittedPerfPlan( - params: PerfRuntimeHandlerParams, + params: PerfRuntimeHandlerParams & { ref: SessionRef }, session: SessionState, admission: AdmittedRuntimePlan>, ): Promise { @@ -183,7 +186,7 @@ async function executeAdmittedPerfPlan( } async function startPerfCapture( - params: PerfRuntimeHandlerParams, + params: PerfRuntimeHandlerParams & { ref: SessionRef }, session: SessionState, runtime: Readonly<{ owner: Parameters[0]['owner']; @@ -225,9 +228,7 @@ async function startPerfCapture( }); await adoptStartedPerfCapture({ admissionLedger: requirePerfCaptureAdmissionLedger(params), - session, - sessionName: params.sessionName, - sessionStore: params.sessionStore, + binding: bindSessionPerfCapture(params.sessionStore, params.ref), device: session.device, owner: runtime.owner, fence, diff --git a/src/daemon/session-store.ts b/src/daemon/session-store.ts index 9e7b1c2d8d..97e2d6dc4b 100644 --- a/src/daemon/session-store.ts +++ b/src/daemon/session-store.ts @@ -73,19 +73,27 @@ export class SessionStore { this.acceptingSessions = false; } - publish(address: string, session: SessionState): SessionRef { + assertAdmissionOpen(address: string): void { if (!this.acceptingSessions) { throw new AppError('COMMAND_FAILED', 'Daemon is shutting down', { reason: 'daemon_shutting_down', session: address, }); } + } + + assertPublishable(address: string): void { + this.assertAdmissionOpen(address); if (this.sessions.has(address)) { throw new AppError('COMMAND_FAILED', 'Session address is already occupied', { reason: 'session_address_occupied', session: address, }); } + } + + publish(address: string, session: SessionState): SessionRef { + this.assertPublishable(address); const entry = { current: session }; this.sessions.set(address, entry); this.clearIdleExpiryTombstone(address); diff --git a/src/daemon/session-teardown.ts b/src/daemon/session-teardown.ts index 377be3fb2b..2f051b2057 100644 --- a/src/daemon/session-teardown.ts +++ b/src/daemon/session-teardown.ts @@ -3,7 +3,6 @@ import { emitDiagnostic } from '@agent-device/host-kit/diagnostics'; import { cleanupRetainedMaterializedPathsForSession } from './materialized-path-registry.ts'; import type { SessionState } from './session-state.ts'; import type { SessionStore } from './session-store.ts'; -import { forceCleanupSessionAppLog } from './app-log-session-resource.ts'; import { appLogResourceStore } from './app-log-resource-store.ts'; import { finishLiveAudioProbe } from '@agent-device/capture-kit/audio-probe-session-resource'; import { finishLivePerfCapture } from '@agent-device/capture-kit/perf-capture-session-resource'; @@ -18,6 +17,7 @@ export async function stopSessionAppLog(params: { }): Promise { const { session, sessionName, sessionStore } = params; if (!session.appLog) return; + const { forceCleanupSessionAppLog } = await import('./app-log-session-resource.ts'); await forceCleanupSessionAppLog({ session, sessionName,