diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts index 80ea85e8683..f69bc3fb06a 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts @@ -120,6 +120,10 @@ export type StructuredAgentSessionHostDeps = { /** Whether an orchestration dispatch still owns this session's worker; absent answers no. */ hasOpenDispatch?: (record: AgentSessionRecord) => boolean onEventSinkError?: (input: { sessionId: string; error: unknown }) => void + /** Lease bookkeeping run for startup or a read (the reconcile, or resolving a chat's recovery) + * that refused or threw, once per distinct failure. Startup and the read carry on: the next + * attach or send reconciles and resolves recovery again before it acts. */ + onLeaseReconcileFailure?: (failure: unknown) => void /** Every status projection this host publishes. `replay` marks a re-projection of state the host * already knew (restore, an arriving subscriber) rather than a fresh journal edge. */ onSessionStatusChanged?: ( diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts index 79a129118ee..bcda340ee79 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts @@ -135,7 +135,7 @@ export class StructuredAgentSessionHost { flushStreamedEvents: (sessionId) => this.flushStreamedEvents(sessionId) }) this.restore = createStructuredAgentSessionHostRestore(deps, { - reconcile: this.reconcileLeases, + reconcileLeases: this.reconcileLeases, resolveRecovery: (sessionId) => this.runtimeState.resolveRecovery(sessionId), serialize: (sessionId, task) => this.serialize(sessionId, task), hasSession: this.hasSession, @@ -215,6 +215,7 @@ export class StructuredAgentSessionHost { listSessionTabs = () => sessionTabs.listStructuredAgentSessionTabs(this.sessions) getPersistedVisibleSessionTabIndex = () => this.deps.store.getVisibleSessionTabIndex() getSessionTabId = (sessionId: string): string | null => this.deps.store.getSessionTabId(sessionId) + showSessionTabs = (sessionIds: readonly string[]) => this.deps.store.showSessionTabs(sessionIds) setSessionTabVisibility = async ( sessionId: string, @@ -228,12 +229,7 @@ export class StructuredAgentSessionHost { } } - reconcileRestartLeases = async (): Promise => { - const refusal = await this.reconcileLeases('startup') - if (refusal) { - throw new Error(refusal.code) - } - } + reconcileRestartLeases = (): Promise => this.restore.reconcileRestartLeases() restoreReadableSessions = (sessionIds?: readonly string[]): Promise => this.restore.restoreReadableSessions(sessionIds) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-readable-restorer.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-readable-restorer.test.ts index 5eaedfffb38..09cab618ab9 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-readable-restorer.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-readable-restorer.test.ts @@ -39,8 +39,8 @@ describe('StructuredAgentSessionReadableRestorer', () => { adapter: {} }, supportsRecord: () => true, - reconcile: async () => null, - resolveRecovery: async () => undefined, + reconcile: async () => true, + resolveRecovery: async () => true, serialize: async (_sessionId, task) => task(), hasSession: () => false, onReadable: () => undefined diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-restart-reconcile.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-restart-reconcile.ts index 3b6cc30271e..f349f252ebe 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-restart-reconcile.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-restart-reconcile.ts @@ -45,6 +45,76 @@ export function createRestartReconciler(deps: { } } +/** Reports each distinct failure once, until `clear` says the bookkeeping settled again: the + * startup check, the restore pass and a restore retried after a failed journal open would each + * log the same store failure. */ +export type ReaderBookkeepingFailures = { + report: (failure: unknown) => void + clear: () => void +} + +export function reportEachFailureOnce( + onFailure: ((failure: unknown) => void) | undefined +): ReaderBookkeepingFailures { + let reported: string | null = null + return { + report: (failure) => { + const key = failureKey(failure) + if (key !== reported) { + reported = key + try { + onFailure?.(failure) + } catch (sinkError) { + // A throwing sink must not turn the reported failure back into a failed read. + console.warn('[structured-agent-session] reporting a lease bookkeeping failure failed', { + failure, + sinkError + }) + } + } + }, + clear: () => { + reported = null + } + } +} + +/** The reconcile a reader runs, at startup and before each restored read: it never throws, since + * an unreconciled lease grants no writer and the next send reconciles again before it acts. + * Answers whether every lease is settled. */ +export function createReaderReconcile( + reconcile: (sessionId: string) => Promise, + failures: ReaderBookkeepingFailures +): (sessionId: string) => Promise { + return async (sessionId) => { + let failure: unknown + try { + const refusal = await reconcile(sessionId) + if (!refusal) { + failures.clear() + return true + } + failure = refusal + } catch (error) { + failure = error + } + failures.report(failure) + return false + } +} + +function failureKey(failure: unknown): string { + if (typeof failure === 'object' && failure !== null) { + if ('code' in failure && failure.code) { + return String(failure.code) + } + if ('message' in failure) { + return String(failure.message) + } + } + return String(failure) +} + async function reconcileCurrentLeases(deps: { store: AgentSessionRecordStore probe: (record: AgentSessionRecord) => Promise diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.test.ts index c7d04ea64cc..0cffa1e32c7 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.test.ts @@ -62,8 +62,8 @@ describe('restart journal restoration', () => { const restoration = restoreStructuredAgentSessionsOnRestart({ openDeps: NO_OPEN_DEPS, records, - reconcile: async () => null, - resolveRecovery: async () => undefined, + reconcile: async () => true, + resolveRecovery: async () => true, serialize: async (_sessionId, task) => task(), hasSession: () => false, onReadable: () => undefined @@ -103,8 +103,8 @@ describe('restart journal restoration', () => { await restoreStructuredAgentSessionsOnRestart({ openDeps: NO_OPEN_DEPS, records, - reconcile: async () => null, - resolveRecovery: async () => undefined, + reconcile: async () => true, + resolveRecovery: async () => true, serialize: async (_sessionId, task) => task(), hasSession: () => false, onReadable: () => undefined @@ -150,9 +150,10 @@ describe('restart journal restoration', () => { // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the restore reads only the record's session id here. records: [{ sessionId: 'session-1' } as AgentSessionRecord], openDeps: NO_OPEN_DEPS, - reconcile: async () => null, + reconcile: async () => true, resolveRecovery: async () => { calls.push('resolveRecovery') + return true }, serialize: async (_sessionId, task) => task(), hasSession: () => false, @@ -164,6 +165,58 @@ describe('restart journal restoration', () => { expect(calls).toEqual(['resolveRecovery', 'open', 'onReadable:restored']) }) + // Each failed bookkeeping call stands for one wait on a held store lock. + describe('once lease bookkeeping fails in a pass', () => { + const records = Array.from( + { length: 8 }, + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the restore reads only the record's session id here. + (_, index) => ({ sessionId: `session-${index}` }) as AgentSessionRecord + ) + const restore = ( + bookkeeping: Pick< + Parameters[0], + 'reconcile' | 'resolveRecovery' + > + ) => + restoreStructuredAgentSessionsOnRestart({ + openDeps: NO_OPEN_DEPS, + records, + ...bookkeeping, + serialize: async (_sessionId, task) => task(), + hasSession: () => false, + onReadable: () => undefined + }) + + /** A failure that takes a while, as a lock wait does, so the chats open at once overlap it. */ + const slowFailure = async (): Promise => { + await new Promise((resolve) => setTimeout(resolve, 5)) + return false + } + + beforeEach(() => restoreRead.mockResolvedValue(null)) + + it('skips it for every chat when the pass check fails, and still opens them all', async () => { + const reconcile = vi.fn(slowFailure) + const resolveRecovery = vi.fn(async () => true) + + await restore({ reconcile, resolveRecovery }) + + expect(reconcile).toHaveBeenCalledOnce() + expect(resolveRecovery).not.toHaveBeenCalled() + expect(restoreRead).toHaveBeenCalledTimes(records.length) + }) + + it('starts no more after the first failed recovery, and still opens every chat', async () => { + const resolveRecovery = vi.fn(slowFailure) + + await restore({ reconcile: async () => true, resolveRecovery }) + + // Only those already started when the first failed: at most one per chat open at once. + expect(resolveRecovery.mock.calls.length).toBeLessThanOrEqual(4) + expect(restoreRead).toHaveBeenCalledTimes(records.length) + }) + }) + it('does not settle again when a second restore finds the session already open', async () => { restoreRead.mockResolvedValue({ session: { journal: {}, params: {}, child: null }, @@ -174,8 +227,8 @@ describe('restart journal restoration', () => { // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the restore reads only the record's session id here. records: [{ sessionId: 'session-1' } as AgentSessionRecord], openDeps: NO_OPEN_DEPS, - reconcile: async () => null, - resolveRecovery: async () => undefined, + reconcile: async () => true, + resolveRecovery: async () => true, serialize: async (_sessionId, task) => task(), hasSession: () => true, onReadable: () => undefined diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.ts index 4e4b7f1e5dd..29ea17845d8 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-restart-restore.ts @@ -13,7 +13,6 @@ import { setImmediate as yieldToEventLoop } from 'node:timers/promises' import type { AgentSessionRecord } from '../../../shared/agent-session-record' -import type { AgentSessionWireRefusal } from '../../../shared/agent-session-wire' import { mapWithConcurrency } from '../../../shared/map-with-concurrency' import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' import type { @@ -28,8 +27,11 @@ export type StructuredAgentSessionReadRestoreDeps = { openDeps: StructuredAgentSessionConversationOpenDeps & { store: Pick } - reconcile: (sessionId: string) => Promise - resolveRecovery: (sessionId: string) => Promise + // Lease bookkeeping. Neither throws: a read grants no writer, so bookkeeping must not block it. + /** Whether every lease is settled. */ + reconcile: (sessionId: string) => Promise + /** False when its store write failed; the next attach or send resolves it again. */ + resolveRecovery: (sessionId: string) => Promise serialize: (sessionId: string, task: () => Promise) => Promise hasSession: (sessionId: string) => boolean onReadable: ( @@ -41,13 +43,10 @@ export type StructuredAgentSessionReadRestoreDeps = { /** One session's share of the restart restore. Startup maps this over every supported record. */ async function restoreOneStructuredAgentSessionRead( input: StructuredAgentSessionReadRestoreDeps, - sessionId: string + sessionId: string, + settleLeases: (sessionId: string) => Promise ): Promise { - const unreconciled = await input.reconcile(sessionId) - if (!unreconciled) { - // A session latched in recovery exits here at startup, without waiting for a client. - await input.resolveRecovery(sessionId) - } + await settleLeases(sessionId) await input.serialize(sessionId, () => restoreOneStructuredAgentSessionReadUnderSerialize(input, sessionId) ) @@ -73,9 +72,26 @@ async function restoreOneStructuredAgentSessionReadUnderSerialize( export async function restoreStructuredAgentSessionsOnRestart( input: StructuredAgentSessionReadRestoreDeps & { records: AgentSessionRecord[] } ): Promise { + const [first] = input.records + if (!first) { + return + } + // One check for the pass. Each chat checks again while it holds, since another writer can mark + // leases unreconciled mid-pass; after the first failure, retrying per chat only waits on the + // same store again, and the next attach or send settles those chats instead. + let settled = await input.reconcile(first.sessionId) + const settleLeases = async (sessionId: string): Promise => { + // A session latched in recovery exits here at startup, without waiting for a client. + if ( + settled && + !((await input.reconcile(sessionId)) && (await input.resolveRecovery(sessionId))) + ) { + settled = false + } + } await mapWithConcurrency(input.records, JOURNAL_RESTORE_CONCURRENCY, async ({ sessionId }) => { // A journal open is synchronous SQLite: without a macrotask per chat the restore is one long task. await yieldToEventLoop() - await restoreOneStructuredAgentSessionRead(input, sessionId) + await restoreOneStructuredAgentSessionRead(input, sessionId, settleLeases) }) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts index 226e6536c9c..3290ba05f0d 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-restore-without-import.test.ts @@ -232,8 +232,8 @@ async function restore(sessionIds: readonly string[]) { await restoreStructuredAgentSessionsOnRestart({ openDeps: deps, records: sessionIds.map(recordFor), - reconcile: async () => null, - resolveRecovery: async () => undefined, + reconcile: async () => true, + resolveRecovery: async () => true, serialize: async (_sessionId, task) => task(), hasSession: (sessionId) => sessions.has(sessionId), onReadable: (sessionId, opened) => { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-reveal.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-reveal.ts index 2a1899ab1cc..262c9e18d3c 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-reveal.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-reveal.ts @@ -10,10 +10,15 @@ // a send does. And a journal it cannot open is not a refusal — the chat shows that failure with a // Retry, so the tab is worth publishing either way. +import type { AgentSessionWireRefusal } from '../../../shared/agent-session-wire' import { agentSessionRefusalError } from '../../../shared/agent-session-wire-refusals' import { adapterSupportsRecord } from './structured-agent-session-provider-support' import { StructuredAgentSessionReadableRestorer } from './structured-agent-session-readable-restorer' import { StructuredAgentSessionRestartRestoreGate } from './structured-agent-session-restart-restore-gate' +import { + createReaderReconcile, + reportEachFailureOnce +} from './structured-agent-session-restart-reconcile' import type { StructuredAgentSessionHostDeps, StructuredAgentSessionReveal @@ -51,23 +56,44 @@ export async function revealStructuredAgentSession( } } -/** The host's startup readable-restore sweep: reconcile, resolve, then open each chat's journal. */ +/** The host's startup readable-restore sweep: reconcile, resolve, then open each chat's journal. + * Its lease bookkeeping is a reader's, which never fails a read or startup; startup shares it. */ export function createStructuredAgentSessionHostRestore( deps: StructuredAgentSessionHostDeps, wiring: Omit< ConstructorParameters[0], - 'openDeps' | 'supportsRecord' - > + 'openDeps' | 'supportsRecord' | 'reconcile' | 'resolveRecovery' + > & { + reconcileLeases: (sessionId: string) => Promise + resolveRecovery: (sessionId: string) => Promise + } ): { + reconcileRestartLeases: () => Promise restoreReadableSessions: (sessionIds?: readonly string[]) => Promise } { + const { reconcileLeases, resolveRecovery, ...rest } = wiring + const failures = reportEachFailureOnce(deps.onLeaseReconcileFailure) + const reconcile = createReaderReconcile(reconcileLeases, failures) const restorer = new StructuredAgentSessionReadableRestorer({ openDeps: deps, supportsRecord: (record) => adapterSupportsRecord(deps.adapter, record), - ...wiring + reconcile, + // The next attach or send resolves recovery again, strictly, before it acts. + resolveRecovery: (sessionId) => + resolveRecovery(sessionId).then( + () => true, + (error: unknown) => { + failures.report(error) + return false + } + ), + ...rest }) const gate = new StructuredAgentSessionRestartRestoreGate() return { + reconcileRestartLeases: async () => { + await reconcile('startup') + }, restoreReadableSessions: (sessionIds) => gate.run(() => restorer.restore(sessionIds)) } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-startup-reconcile-failure.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-startup-reconcile-failure.test.ts new file mode 100644 index 00000000000..1d7b6621b3e --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-startup-reconcile-failure.test.ts @@ -0,0 +1,220 @@ +// The lease reconcile is bookkeeping: a lock that gives up or a store this build may not write is +// reported, and startup and every read carry on. Nothing is owed after it, because an unreconciled +// lease grants no writer and the next send reconciles every lease again before it acts. + +import { cp, readdir, readFile, rm } from 'node:fs/promises' +import { dirname } from 'node:path' +import { afterEach, expect, it, vi } from 'vitest' +import type * as FileTransactionLock from '../../file-transaction-lock' +import { AGENT_SESSION_STORE_SCHEMA_VERSION } from '../../runtime/agent-session-record-store-file' +import { + editPersistedTestAgentSessionStore, + openTestAgentSessionRecordStore, + type PersistedTestAgentSessionStore, + testAgentSessionStoreFilePath +} from '../../runtime/agent-session-record-store-test-harness' +import { + StructuredAgentSessionHost, + type StructuredAgentSessionHostDeps +} from './structured-agent-session-host' +import { + adapter, + attach, + CALLER, + envelope, + hostTestState, + replaceHostTestState +} from './structured-agent-session-host-test-harness' +import { + HOST_TEST_NOW as NOW, + HOST_TEST_SESSION as SESSION, + hostTestMessage +} from './structured-agent-session-host-test-data' +import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support' + +const lock = vi.hoisted(() => ({ failing: false })) + +vi.mock('../../file-transaction-lock', async (importOriginal) => { + const actual = await importOriginal() + return { + ...actual, + withFileTransactionLock: (...args: Parameters) => + lock.failing + ? // What proper-lockfile throws once its retries give up. + Promise.reject( + Object.assign(new Error('Lock file is already being held'), { code: 'ELOCKED' }) + ) + : actual.withFileTransactionLock(...args) + } +}) + +const relaunchedRoots: string[] = [] + +afterEach(async () => { + lock.failing = false + await Promise.all( + relaunchedRoots.splice(0).map((dir) => rm(dir, { recursive: true, force: true })) + ) +}) + +/** A host with one chat, relaunched over a copy of its files; `rewrite` edits the copied store. */ +async function relaunch( + rewrite?: (persisted: PersistedTestAgentSessionStore) => void, + probeOwner: StructuredAgentSessionHostDeps['probeOwner'] = async () => ({ + outcome: 'pid-absent' + }) +) { + const dying = hostTestState() + await attach() + // An empty renewal queues behind every record write, so they are on disk. + await dying.store.renewLeases([]) + const relaunched = `${dying.root}-relaunched` + relaunchedRoots.push(relaunched) + // A dead process holds no lock. + await cp(dying.root, relaunched, { + recursive: true, + filter: (source) => !source.includes('.lock') + }) + const storePath = testAgentSessionStoreFilePath(relaunched) + if (rewrite) { + await editPersistedTestAgentSessionStore(relaunched, rewrite) + } + const store = await openTestAgentSessionRecordStore(relaunched) + const onLeaseReconcileFailure = vi.fn() + const host = new StructuredAgentSessionHost({ + store, + adapter: adapter(), + journalDatabase: openTestJournalHostDatabase(relaunched), + claimKeyId: 'key-1', + mintSpawnToken: () => 'spawn-next', + probeOwner, + stopOwnerProcess: () => { + throw new Error('an owner not proven alive must not be stopped') + }, + now: () => NOW, + onLeaseReconcileFailure + }) + replaceHostTestState({ store, host }) + return { host, store, storePath, onLeaseReconcileFailure } +} + +it('reports a startup reconcile whose store write fails, and does not reject', async () => { + const { host, store, onLeaseReconcileFailure } = await relaunch() + + lock.failing = true + await expect(host.reconcileRestartLeases()).resolves.toBeUndefined() + + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + expect(onLeaseReconcileFailure).toHaveBeenCalledWith(expect.objectContaining({ code: 'ELOCKED' })) + // Nothing was adjudicated, so the lease still grants no writer. + expect(store.getRecord(SESSION)?.lease.unreconciled).toBe(true) +}) + +it('reconciles the chat on its next send once the store can be written again', async () => { + const { host, store, onLeaseReconcileFailure } = await relaunch() + lock.failing = true + await host.reconcileRestartLeases() + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + lock.failing = false + + const body = hostTestMessage('sent after a startup reconcile failed') + await expect( + host.send(CALLER, { envelope: envelope('agentSession.send', { body }), body }) + ).resolves.toMatchObject({ ok: true }) + + // The send is queued; delivering it starts the agent, and that start reconciles the lease. + await vi.waitFor(() => expect(hostTestState().dispatch).toHaveBeenCalledOnce(), { + timeout: 10_000 + }) + expect(store.getRecord(SESSION)?.lease.unreconciled).toBe(false) + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + // Before the relaunched directory is removed, so the child's wind-down can write its lease. + await host.flushAllStreamedEvents() +}) + +it('reports a store a newer Orca wrote without writing it, and does not reject', async () => { + const { host, store, storePath, onLeaseReconcileFailure } = await relaunch((persisted) => { + persisted.schemaVersion = AGENT_SESSION_STORE_SCHEMA_VERSION + 1 + }) + expect(store.readOnly).toBe(true) + const bytes = await readFile(storePath) + const files = await readdir(dirname(storePath)) + + await expect(host.reconcileRestartLeases()).resolves.toBeUndefined() + + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + expect(onLeaseReconcileFailure).toHaveBeenCalledWith( + expect.objectContaining({ message: 'agent_session_legacy_required' }) + ) + expect(store.listRecords().map((record) => record.sessionId)).toEqual([SESSION]) + expect(await readFile(storePath)).toEqual(bytes) + expect(await readdir(dirname(storePath))).toEqual(files) +}) + +it('restores a chat for reading while the reconcile keeps failing, and reports it once', async () => { + const { host, store, onLeaseReconcileFailure } = await relaunch() + lock.failing = true + await host.reconcileRestartLeases() + + await expect(host.restoreReadableSessions([SESSION])).resolves.toBeUndefined() + + expect(host.hasSession(SESSION)).toBe(true) + expect(store.getRecord(SESSION)?.lease.unreconciled).toBe(true) + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + expect(onLeaseReconcileFailure).toHaveBeenCalledWith(expect.objectContaining({ code: 'ELOCKED' })) +}) + +it('restores a chat for reading from a store a newer Orca wrote', async () => { + const { host, storePath, onLeaseReconcileFailure } = await relaunch((persisted) => { + persisted.schemaVersion = AGENT_SESSION_STORE_SCHEMA_VERSION + 1 + }) + const bytes = await readFile(storePath) + + await expect(host.restoreReadableSessions([SESSION])).resolves.toBeUndefined() + + expect(host.hasSession(SESSION)).toBe(true) + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + expect(await readFile(storePath)).toEqual(bytes) +}) + +// A chat whose owner could not be proven gone is left recovering; the next attach or send retries it. +it('restores a chat for reading when resolving its recovery cannot write the store', async () => { + const { host, store, onLeaseReconcileFailure } = await relaunch(undefined, async () => ({ + outcome: 'indeterminate', + reason: 'probe' + })) + await host.reconcileRestartLeases() + expect(store.getRecord(SESSION)?.lease.handoffStage).toBe('recovering') + lock.failing = true + + await expect(host.restoreReadableSessions([SESSION])).resolves.toBeUndefined() + + expect(host.hasSession(SESSION)).toBe(true) + expect(store.getRecord(SESSION)?.lease.handoffStage).toBe('recovering') + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + expect(onLeaseReconcileFailure).toHaveBeenCalledWith(expect.objectContaining({ code: 'ELOCKED' })) +}) + +it.each([ + ['startup reconcile', (host: StructuredAgentSessionHost) => host.reconcileRestartLeases()], + ['read restore', (host: StructuredAgentSessionHost) => host.restoreReadableSessions([SESSION])] +])('keeps the %s resolving when the failure sink throws', async (_step, read) => { + const { host, onLeaseReconcileFailure } = await relaunch() + const sinkError = new Error('error sink failed') + onLeaseReconcileFailure.mockImplementation(() => { + throw sinkError + }) + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}) + lock.failing = true + try { + await expect(read(host)).resolves.toBeUndefined() + + expect(onLeaseReconcileFailure).toHaveBeenCalledOnce() + expect(warn).toHaveBeenCalledWith( + expect.stringContaining('reporting a lease bookkeeping failure failed'), + expect.objectContaining({ failure: expect.objectContaining({ code: 'ELOCKED' }), sinkError }) + ) + } finally { + warn.mockRestore() + } +}) diff --git a/src/main/runtime/agent-session-record-store.ts b/src/main/runtime/agent-session-record-store.ts index 925fa38d8d1..603b5ed02fa 100644 --- a/src/main/runtime/agent-session-record-store.ts +++ b/src/main/runtime/agent-session-record-store.ts @@ -66,7 +66,7 @@ import { agentSessionStorePath, type AgentSessionStoreState } from './agent-session-record-store-file' -import { setAgentSessionTabVisibility } from './agent-session-tab-table' +import { setAgentSessionTabVisibility, showAgentSessionTabs } from './agent-session-tab-table' import { loadProtectedAgentSessionStore } from './agent-session-record-store-security' import { AgentSessionStoreTransactionQueue, @@ -146,6 +146,12 @@ export class AgentSessionRecordStore { return this.transact(() => setAgentSessionTabVisibility(this.state, sessionId, visible, tabId)) } + /** Shows each session that still has a record, in one write: an index written part way would + * read as complete at the next launch and drop the rest. */ + showSessionTabs(sessionIds: readonly string[]): Promise { + return this.transact(() => showAgentSessionTabs(this.state, sessionIds)) + } + listByScope(location: AgentSessionExecutionLocation): AgentSessionRecord[] { const scope = agentSessionScopeKey(location) return this.listRecords().filter((record) => agentSessionScopeKey(record.location) === scope) diff --git a/src/main/runtime/agent-session-tab-table.ts b/src/main/runtime/agent-session-tab-table.ts index 6022aac8bf2..72551d563ca 100644 --- a/src/main/runtime/agent-session-tab-table.ts +++ b/src/main/runtime/agent-session-tab-table.ts @@ -131,6 +131,17 @@ export function setAgentSessionTabVisibility( } } +export function showAgentSessionTabs( + state: AgentSessionStoreState, + sessionIds: readonly string[] +): void { + for (const sessionId of sessionIds) { + if (state.records.has(sessionId)) { + setAgentSessionTabVisibility(state, sessionId, true) + } + } +} + export type PersistedAgentSessionTab = { tabId: string; sessionId: string } export function serializeAgentSessionTabTable(table: AgentSessionTabTable): { diff --git a/src/main/runtime/orca-runtime-restore-structured-agent-session-tabs-once.ts b/src/main/runtime/orca-runtime-restore-structured-agent-session-tabs-once.ts index 8d29a0d360b..10dd920b253 100644 --- a/src/main/runtime/orca-runtime-restore-structured-agent-session-tabs-once.ts +++ b/src/main/runtime/orca-runtime-restore-structured-agent-session-tabs-once.ts @@ -5,6 +5,7 @@ import { getStructuredAgentSessionHost } from '../native-chat/agent-session-wire import { replaceConversationInSnapshot } from './structured-conversation-tab-replacement' import type { ConversationReplacement } from '../native-chat/agent-session-wire/structured-conversation-command' import { collectSavedStructuredAgentSessionIds } from './saved-structured-agent-session-restoration' +import { seedStructuredAgentSessionTabIndex } from './structured-agent-session-tab-index-seed' import { LOCAL_EXECUTION_HOST_ID } from '../../shared/execution-host' import type { RuntimeMobileSessionAgentTab, @@ -28,7 +29,9 @@ import { parseAppSshPtyId } from '../../shared/ssh-pty-id' import type { PtyProcessInspection } from '../providers/pty-process-inspection' export class OrcaRuntimeWithRestoreStructuredAgentSessionTabsOnce extends OrcaRuntimeWithGetStructuredAgentSessionCreateSupport { - async replaceStructuredAgentSessionTab(replacement: ConversationReplacement): Promise { + /** Projects only: a replacement's chat already has its tab in the store, which the /clear commit + * moved in the same write. */ + replaceStructuredAgentSessionTab(replacement: ConversationReplacement): void { const prior = this.mobileSessionTabsByWorktree.get(replacement.workspaceId) const next = prior ? replaceConversationInSnapshot(prior, replacement) : null if (next && next !== prior) { @@ -39,7 +42,7 @@ export class OrcaRuntimeWithRestoreStructuredAgentSessionTabsOnce extends OrcaRu (tab) => tab.type === 'agent-session' && tab.sessionId === replacement.sessionId ) ) { - await this.publishStructuredAgentSessionTab({ + this.projectStructuredAgentSessionTab({ ...replacement, replacesSessionId: replacement.sourceSessionId, activate: false @@ -57,9 +60,8 @@ export class OrcaRuntimeWithRestoreStructuredAgentSessionTabsOnce extends OrcaRu const profileIds = collectSavedStructuredAgentSessionIds( this.store?.getWorkspaceSession?.(LOCAL_EXECUTION_HOST_ID) ?? null ) - await host?.restoreReadableSessions( - persistedVisibleIndex.present ? persistedVisibleIndex.sessionIds : profileIds - ) + const targets = persistedVisibleIndex.present ? persistedVisibleIndex.sessionIds : profileIds + await host?.restoreReadableSessions(targets) for (const worktreeId of this.getKnownWorkspaceSessionWorktreeIds()) { this.hydrateHeadlessMobileSessionTabsFromWorkspaceSession(worktreeId, { allowAttachedWindow: true, @@ -67,24 +69,27 @@ export class OrcaRuntimeWithRestoreStructuredAgentSessionTabsOnce extends OrcaRu }) } this.hydrateHeadlessMobileSessionTabsFromWorkspaceSession() - for (const replacement of host?.conversationReplacements?.() ?? []) { - await this.replaceStructuredAgentSessionTab(replacement) - } - for (const session of host?.listSessionTabs() ?? []) { + const restored = (host?.listSessionTabs() ?? []).flatMap((session) => { if (session.agent !== 'codex' && session.agent !== 'claude') { - continue + return [] } let sessionId = session.sessionId while (sessionId.startsWith('agent-session:')) { sessionId = sessionId.slice('agent-session:'.length) } - await this.publishStructuredAgentSessionTab({ - ...session, - agent: session.agent, - sessionId, - activate: false, - notify: false - }) + return [{ ...session, agent: session.agent, sessionId }] + }) + await seedStructuredAgentSessionTabIndex( + host, + targets, + restored.map((session) => session.sessionId) + ) + // Past the seed, projecting records nothing. + for (const replacement of host?.conversationReplacements?.() ?? []) { + this.replaceStructuredAgentSessionTab(replacement) + } + for (const session of restored) { + this.projectStructuredAgentSessionTab({ ...session, activate: false, notify: false }) } const wasUnverifiable = this.structuredAgentSessionInventoryUnverifiable // No host means no one can say which chats exist; with none on disk, empty is the answer. @@ -114,6 +119,18 @@ export class OrcaRuntimeWithRestoreStructuredAgentSessionTabsOnce extends OrcaRu ...(input.tabId ? [input.tabId] : []) ) } + this.projectStructuredAgentSessionTab(input) + } + + /** The runtime's own snapshot of a chat tab. Records nothing: the caller owns the store write. */ + projectStructuredAgentSessionTab(input: { + workspaceId: string + sessionId: string + agent: 'claude' | 'codex' + activate: boolean + notify?: boolean + replacesSessionId?: string + }): void { const existing = this.mobileSessionTabsByWorktree.get(input.workspaceId) const id = `agent-session:${input.sessionId}` if (existing?.tabs.some((tab) => tab.id === id)) { diff --git a/src/main/runtime/orca-runtime-structured-session-restore.test.ts b/src/main/runtime/orca-runtime-structured-session-restore.test.ts index 1d8105d07f8..929ee82b6c8 100644 --- a/src/main/runtime/orca-runtime-structured-session-restore.test.ts +++ b/src/main/runtime/orca-runtime-structured-session-restore.test.ts @@ -324,7 +324,7 @@ describe('structured session cold restoration', () => { it('publishes restored Claude tabs with the Claude title', async () => { const runtime = new OrcaRuntimeService() - const publish = vi.spyOn(runtime, 'publishStructuredAgentSessionTab') + const project = vi.spyOn(runtime, 'projectStructuredAgentSessionTab') const internal = runtime as unknown as { hasPersistedStructuredAgentSessionStore(): boolean getKnownWorkspaceSessionWorktreeIds(): Set @@ -351,7 +351,7 @@ describe('structured session cold restoration', () => { await runtime.restoreStructuredAgentSessionTabs() - expect(publish).toHaveBeenCalledWith({ + expect(project).toHaveBeenCalledWith({ workspaceId: 'workspace-1', sessionId: 'restored-claude', agent: 'claude', @@ -372,6 +372,30 @@ describe('structured session cold restoration', () => { ) }) + // The /clear commit moved the tab in the store already, so replacing it only projects. + it('replaces a cleared conversation tab without a store write', async () => { + const runtime = new OrcaRuntimeService() + const setSessionTabVisibility = vi.fn(async () => undefined) + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: replacing a tab reaches the host only through setSessionTabVisibility, which must stay uncalled. + setStructuredAgentSessionHost({ setSessionTabVisibility } as never) + + runtime.replaceStructuredAgentSessionTab({ + sourceSessionId: 'cleared-session', + sessionId: 'replacement-session', + workspaceId: 'workspace-1', + agent: 'codex' + }) + + const snapshot = await runtime.listMobileSessionTabs('id:workspace-1') + expect(snapshot.tabs).toEqual([ + expect.objectContaining({ + id: 'agent-session:replacement-session', + replacesSessionId: 'cleared-session' + }) + ]) + expect(setSessionTabVisibility).not.toHaveBeenCalled() + }) + it('commits the host close when the renderer already removed the structured tab', async () => { const runtime = new OrcaRuntimeService() runtime.setNotifier({ diff --git a/src/main/runtime/rpc/methods/structured-agent-session.ts b/src/main/runtime/rpc/methods/structured-agent-session.ts index 534d4967779..463772d3337 100644 --- a/src/main/runtime/rpc/methods/structured-agent-session.ts +++ b/src/main/runtime/rpc/methods/structured-agent-session.ts @@ -123,7 +123,7 @@ export const STRUCTURED_AGENT_SESSION_METHODS = [ .conversationReplacements() .find((entry) => entry.sourceSessionId === params.envelope.sessionId) if (replacement) { - await ctx.runtime.replaceStructuredAgentSessionTab(replacement) + ctx.runtime.replaceStructuredAgentSessionTab(replacement) } await host.close(params.envelope.sessionId) } diff --git a/src/main/runtime/structured-agent-session-runtime.ts b/src/main/runtime/structured-agent-session-runtime.ts index 82c9803e2a7..bb82e070e12 100644 --- a/src/main/runtime/structured-agent-session-runtime.ts +++ b/src/main/runtime/structured-agent-session-runtime.ts @@ -317,6 +317,10 @@ async function installOnJournal( : {}), onEventSinkError: ({ sessionId, error }) => deps.onError?.({ scope: `structured-agent-session-journal:${sessionId}`, error }), + onLeaseReconcileFailure: (error) => + deps.onError + ? deps.onError({ scope: 'structured-agent-session-lease-reconcile', error }) + : console.warn('[structured-agent-session] reconciling chat leases failed', error), ...(deps.onSessionStatusChanged ? { onSessionStatusChanged: deps.onSessionStatusChanged } : {}), ...(deps.statusSink ? { statusSink: deps.statusSink } : {}), ...(deps.hasOpenDispatch ? { hasOpenDispatch: deps.hasOpenDispatch } : {}), diff --git a/src/main/runtime/structured-agent-session-startup-reconcile.test.ts b/src/main/runtime/structured-agent-session-startup-reconcile.test.ts new file mode 100644 index 00000000000..33b26f7569b --- /dev/null +++ b/src/main/runtime/structured-agent-session-startup-reconcile.test.ts @@ -0,0 +1,110 @@ +// A profile whose chat records a newer Orca wrote opens read-only. The startup reconcile cannot +// write there, which is bookkeeping: startup must still finish rather than put the whole app into +// its degraded "Session restore failed" mode, and the file must be left exactly as it was. + +import { mkdir, mkdtemp, readdir, readFile, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { agentSessionRecordFixture } from '../../shared/agent-session-record.test-fixture' +import { getStructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-registry' +import { + AGENT_SESSION_STORE_SCHEMA_VERSION, + agentSessionStorePath +} from './agent-session-record-store-file' +import { OrcaRuntimeService } from './orca-runtime' +import { + ensureStructuredAgentSessionHost, + stopStructuredAgentSessionRuntime +} from './structured-agent-session-runtime' + +let root: string + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-startup-reconcile-')) +}) + +afterEach(async () => { + await stopStructuredAgentSessionRuntime().catch(() => undefined) + vi.restoreAllMocks() + await rm(root, { recursive: true, force: true }) +}) + +/** A profile whose record store a newer Orca wrote, with one chat in it. */ +async function seedNewerStore() { + const storeDirectory = join(root, 'agent-sessions') + const path = agentSessionStorePath(storeDirectory) + const record = agentSessionRecordFixture() + await mkdir(storeDirectory, { recursive: true }) + await writeFile( + path, + JSON.stringify({ + schemaVersion: AGENT_SESSION_STORE_SCHEMA_VERSION + 1, + hostId: 'local', + records: { [record.sessionId]: record }, + operations: {}, + retiredClaimKeys: [], + unusableRecords: {} + }) + ) + return { storeDirectory, path, sessionId: record.sessionId } +} + +function startupRuntime(onError?: (input: { scope: string; error: unknown }) => void) { + const runtime = new OrcaRuntimeService() + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: these are the runtime's own protected members; the test roots the host at `root` and stubs the PTY daemon. + const internal = runtime as unknown as { + hasPersistedStructuredAgentSessionStore(): boolean + ensureStructuredAgentSessionHost(): Promise + refreshMobileSessionPtyRecords(): Promise | null> + } + internal.hasPersistedStructuredAgentSessionStore = () => true + internal.ensureStructuredAgentSessionHost = () => + ensureStructuredAgentSessionHost({ + stateDirectory: root, + hostId: 'local', + claimKeyId: 'key-1', + resolveWorkspacePath: async () => root, + resolveEnvironment: async () => ({}), + resolveClaudeAuthPolicy: () => ({ stripAuthEnv: true }), + ...(onError ? { onError } : {}) + }) + internal.refreshMobileSessionPtyRecords = async () => new Set() + return runtime +} + +it('finishes startup over records a newer Orca wrote, reports it, and writes nothing', async () => { + const { storeDirectory, path, sessionId } = await seedNewerStore() + const bytes = await readFile(path) + const files = await readdir(storeDirectory) + const onError = vi.fn() + const runtime = startupRuntime(onError) + + // What the renderer's startup awaits through `app:prepareTerminalStartupRestoration`. + await expect(runtime.prepareStructuredAgentSessionStartupRestoration()).resolves.toBeUndefined() + + expect(onError).toHaveBeenCalledOnce() + expect(onError).toHaveBeenCalledWith({ + scope: 'structured-agent-session-lease-reconcile', + error: expect.objectContaining({ message: 'agent_session_legacy_required' }) + }) + expect(getStructuredAgentSessionHost()?.sessionAgent(sessionId)).toBe('claude') + await stopStructuredAgentSessionRuntime() + expect(await readFile(path)).toEqual(bytes) + expect(await readdir(storeDirectory)).toEqual(files) +}) + +// The desktop installs its host with no error sink, so the failure is logged rather than dropped. +it('logs the failure when the host has no error sink', async () => { + await seedNewerStore() + const warn = vi.spyOn(console, 'warn').mockImplementation(() => undefined) + + await expect( + startupRuntime().prepareStructuredAgentSessionStartupRestoration() + ).resolves.toBeUndefined() + + expect(warn).toHaveBeenCalledWith( + '[structured-agent-session] reconciling chat leases failed', + expect.objectContaining({ message: 'agent_session_legacy_required' }) + ) +}) diff --git a/src/main/runtime/structured-agent-session-startup-tab-restore.test.ts b/src/main/runtime/structured-agent-session-startup-tab-restore.test.ts new file mode 100644 index 00000000000..61fa020e427 --- /dev/null +++ b/src/main/runtime/structured-agent-session-startup-tab-restore.test.ts @@ -0,0 +1,378 @@ +// With native chat on, the renderer's startup also awaits the chat tab restore (`session.tabs.listAll`), +// which reads every chat whose tab was open at quit. That read must not wait on record-store +// bookkeeping: with a saved tab index it writes nothing, and a store that cannot be written costs +// a bounded number of lock waits, not one per chat. + +import { mkdir, mkdtemp, readdir, readFile, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { AgentSessionRecord } from '../../shared/agent-session-record' +import type { RuntimeMobileSessionTabsSnapshot } from '../../shared/runtime-types' +import { + agentSessionLeaseFixture, + agentSessionRecordFixture +} from '../../shared/agent-session-record.test-fixture' +import type * as FileTransactionLock from '../file-transaction-lock' +import { JournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database' +import { openAgentSessionJournal } from '../native-chat/agent-session-journal/journal-store-factory' +import { journalIdentityFor } from '../native-chat/agent-session-wire/structured-agent-session-attach' +import { attachParamsForRecord } from '../native-chat/agent-session-wire/structured-agent-session-conversation-open' +import { getStructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-registry' +import { AgentSessionRecordStore } from './agent-session-record-store' +import { + AGENT_SESSION_STORE_SCHEMA_VERSION, + agentSessionStorePath +} from './agent-session-record-store-file' +import { OrcaRuntimeService } from './orca-runtime' +import { + ensureStructuredAgentSessionHost, + stopStructuredAgentSessionRuntime +} from './structured-agent-session-runtime' + +// `failing` refuses every take; `grants` lets that many more through, then refuses. +const lock = vi.hoisted(() => ({ failing: false, grants: Infinity, refused: 0 })) + +vi.mock('../file-transaction-lock', async (importOriginal) => { + const actual = await importOriginal() + return { + ...actual, + withFileTransactionLock: (...args: Parameters) => { + if (lock.failing || lock.grants <= 0) { + lock.refused += 1 + // What proper-lockfile throws once its retries give up, about 3 s later. + return Promise.reject( + Object.assign(new Error('Lock file is already being held'), { code: 'ELOCKED' }) + ) + } + lock.grants -= 1 + return actual.withFileTransactionLock(...args) + } + } +}) + +const PROMPT = 'add a retry' +const CHAT_A = 'chat-a-0001' +const CHAT_B = 'chat-b-0002' +const CLEARED = 'chat-s-0003' + +let root: string + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-startup-tab-restore-')) +}) + +afterEach(async () => { + Object.assign(lock, { failing: false, grants: Infinity, refused: 0 }) + await stopStructuredAgentSessionRuntime().catch(() => undefined) + vi.restoreAllMocks() + await rm(root, { recursive: true, force: true }) +}) + +/** A released chat, so no startup probe looks for a live owner. */ +function chatRecord( + sessionId: string, + options: { codex?: boolean; clearedInto?: string } = {} +): AgentSessionRecord { + const record = agentSessionRecordFixture( + agentSessionLeaseFixture({ + sessionId, + ownerProcess: null, + reservedSpawnToken: null, + claimStatus: 'released' + }) + ) + const codex = options.codex + ? { + provider: 'codex' as const, + providerHandleChain: record.providerHandleChain.map((link) => ({ + ...link, + handle: { provider: 'codex' as const, threadId: `thread-${sessionId}` } + })), + accountHome: { variable: 'CODEX_HOME' as const, path: join(root, 'codex-home') } + } + : {} + const clear = options.clearedInto + ? { + conversationCommand: { + command: 'clear' as const, + state: 'completed' as const, + phase: 'committed' as const, + operationId: `clear-${sessionId}`, + callerKey: 'client-1', + replacementSessionId: options.clearedInto + } + } + : {} + return { ...record, ...codex, ...clear } +} + +/** Writes the record store; `visible` is the saved tab index, absent on a legacy profile. */ +async function seedStore( + records: AgentSessionRecord[], + options: { newer?: boolean; visible?: string[] } = {} +) { + const storeDirectory = join(root, 'agent-sessions') + const path = agentSessionStorePath(storeDirectory) + await mkdir(storeDirectory, { recursive: true }) + await writeFile( + path, + JSON.stringify({ + schemaVersion: AGENT_SESSION_STORE_SCHEMA_VERSION + (options.newer ? 1 : 0), + hostId: 'local', + records: Object.fromEntries(records.map((record) => [record.sessionId, record])), + operations: {}, + retiredClaimKeys: [], + unusableRecords: {}, + ...(options.visible ? { visibleSessionIds: options.visible } : {}) + }) + ) + return { storeDirectory, path } +} + +/** Each chat's history as the last run left it: one prompt, accepted, so nothing is left to send. */ +async function seedHistory(records: AgentSessionRecord[]): Promise { + const database = JournalHostDatabase.open(root) + for (const record of records) { + const fence = record.lease.runtimeFence + const journal = await openAgentSessionJournal({ + identity: journalIdentityFor( + record, + attachParamsForRecord(record, { clientOperationId: 'seed', expectedRuntimeFence: fence }) + ), + database + }) + await journal.appendSubmission({ + clientMessageId: 'client-1', + payloadFingerprint: 'fp-1', + body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: PROMPT }] }, + fence, + handoverRecorded: true + }) + await journal.resolveDispatch({ + clientMessageId: 'client-1', + fence, + state: 'accepted', + providerIdentity: null + }) + await journal.close() + } + database.close() +} + +function startupRuntime(options: { afterInstall?: () => void; profileChats?: string[] } = {}) { + const onError = vi.fn() + const runtime = new OrcaRuntimeService() + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: these are the runtime's own protected members; the test roots the host at `root`, stubs the PTY daemon and gives it a profile. + const internal = runtime as unknown as { + hasPersistedStructuredAgentSessionStore(): boolean + ensureStructuredAgentSessionHost(): Promise + refreshMobileSessionPtyRecords(): Promise | null> + mobileSessionTabsByWorktree: Map + store: { getWorkspaceSession: () => unknown } + } + internal.hasPersistedStructuredAgentSessionStore = () => true + internal.ensureStructuredAgentSessionHost = async () => { + const installed = await ensureStructuredAgentSessionHost({ + stateDirectory: root, + hostId: 'local', + claimKeyId: 'key-1', + resolveWorkspacePath: async () => root, + resolveEnvironment: async () => ({}), + resolveClaudeAuthPolicy: () => ({ stripAuthEnv: true }), + onError + }) + options.afterInstall?.() + return installed + } + internal.refreshMobileSessionPtyRecords = async () => new Set() + // A legacy profile's saved chat tabs, which a store with no tab index restores from. + internal.store = { + getWorkspaceSession: () => ({ + activeRepoId: null, + activeWorktreeId: 'workspace-1', + activeTabId: null, + tabsByWorktree: {}, + terminalLayoutsByTabId: {}, + unifiedTabs: { + 'workspace-1': (options.profileChats ?? []).map((sessionId, index) => ({ + id: `agent-session:${sessionId}`, + entityId: sessionId, + groupId: 'group-1', + worktreeId: 'workspace-1', + contentType: 'agent-session', + label: 'Codex Chat', + customLabel: null, + color: null, + sortOrder: index, + createdAt: 1 + })) + } + }) + } + return { + runtime, + onError, + published: () => internal.mobileSessionTabsByWorktree.get('workspace-1')?.tabs ?? [] + } +} + +async function expectHistory(sessionId: string): Promise { + const host = getStructuredAgentSessionHost() + expect(JSON.stringify((await host!.journalSnapshot(sessionId)).items)).toContain(PROMPT) +} + +function spyOnTabWrites() { + return { + visibility: vi.spyOn(AgentSessionRecordStore.prototype, 'setSessionTabVisibility'), + seed: vi.spyOn(AgentSessionRecordStore.prototype, 'showSessionTabs') + } +} + +async function tabIndexOnDisk(path: string): Promise { + return JSON.parse(await readFile(path, 'utf-8')).visibleSessionIds +} + +describe('restoring the chat tabs open at quit', () => { + it.each([ + { store: 'records a newer Orca wrote', newer: true, lockFails: false }, + { store: 'a lock that keeps failing', newer: false, lockFails: true } + ])('lists and reads every chat from $store, and writes nothing', async ({ newer, lockFails }) => { + const records = [ + chatRecord(CHAT_A), + chatRecord(CHAT_B), + chatRecord(CLEARED, { clearedInto: CHAT_A }) + ] + const { storeDirectory, path } = await seedStore(records, { newer, visible: [CHAT_A, CHAT_B] }) + await seedHistory(records.slice(0, 2)) + const bytes = await readFile(path) + const files = await readdir(storeDirectory) + const writes = spyOnTabWrites() + vi.spyOn(console, 'warn').mockImplementation(() => undefined) + // Free when the store opens, then held for good. + const { runtime, onError, published } = startupRuntime({ + afterInstall: () => { + lock.failing = lockFails + } + }) + + await expect(runtime.restoreStructuredAgentSessionTabs()).resolves.toBeUndefined() + + expect(published().map((tab) => tab.id)).toEqual([ + `agent-session:${CHAT_A}`, + `agent-session:${CHAT_B}` + ]) + expect(published()[0]).toMatchObject({ replacesSessionId: CLEARED }) + await expectHistory(CHAT_A) + await expectHistory(CHAT_B) + expect(onError).toHaveBeenCalledOnce() + expect(onError).toHaveBeenCalledWith({ + scope: 'structured-agent-session-lease-reconcile', + error: expect.objectContaining( + newer ? { message: 'agent_session_legacy_required' } : { code: 'ELOCKED' } + ) + }) + expect(writes.visibility).not.toHaveBeenCalled() + expect(writes.seed).not.toHaveBeenCalled() + await stopStructuredAgentSessionRuntime() + if (newer) { + expect(await readFile(path)).toEqual(bytes) + expect(await readdir(storeDirectory)).toEqual(files) + } + }) + + // Each refused take stands for one lock wait of about 3 s. + it.each([1, 4, 8])( + 'waits on a held lock once per startup step with %i chats open', + async (count) => { + const chats = Array.from({ length: count }, (_, index) => `chat-${index}-000${index}`) + const records = chats.map((sessionId) => chatRecord(sessionId)) + await seedStore(records, { visible: chats }) + await seedHistory(records) + vi.spyOn(console, 'warn').mockImplementation(() => undefined) + const { runtime, published } = startupRuntime({ + afterInstall: () => { + lock.failing = true + } + }) + + await runtime.prepareStructuredAgentSessionStartupRestoration() + const prepared = lock.refused + await runtime.restoreStructuredAgentSessionTabs() + + expect(published()).toHaveLength(count) + expect(prepared).toBe(1) + expect(lock.refused - prepared).toBe(1) + } + ) + + describe('on a legacy profile, with no tab index yet', () => { + const legacyChats = () => [ + chatRecord(CHAT_A, { codex: true }), + chatRecord(CHAT_B, { codex: true }) + ] + + it('records every restored chat in one write', async () => { + const records = legacyChats() + const { path } = await seedStore(records) + await seedHistory(records) + const writes = spyOnTabWrites() + const { runtime, published } = startupRuntime({ profileChats: [CHAT_A, CHAT_B] }) + await runtime.prepareStructuredAgentSessionStartupRestoration() + // One more take, then held: a seed written chat by chat would stop part way. + lock.grants = 1 + + await runtime.restoreStructuredAgentSessionTabs() + + expect(published().map((tab) => tab.id)).toEqual([ + `agent-session:${CHAT_A}`, + `agent-session:${CHAT_B}` + ]) + expect(await tabIndexOnDisk(path)).toEqual([CHAT_A, CHAT_B]) + expect(writes.seed).toHaveBeenCalledOnce() + expect(writes.visibility).not.toHaveBeenCalled() + }) + + it('waits on a held lock once more, for that one write', async () => { + const records = legacyChats() + await seedStore(records) + await seedHistory(records) + vi.spyOn(console, 'warn').mockImplementation(() => undefined) + const { runtime, published } = startupRuntime({ + profileChats: [CHAT_A, CHAT_B], + afterInstall: () => { + lock.failing = true + } + }) + + await runtime.prepareStructuredAgentSessionStartupRestoration() + const prepared = lock.refused + await runtime.restoreStructuredAgentSessionTabs() + + expect(published()).toHaveLength(2) + expect(prepared).toBe(1) + // The restore's lease check and the seed. + expect(lock.refused - prepared).toBe(2) + }) + + it('still lists the chats when that write fails, and leaves the index absent', async () => { + const records = legacyChats() + const { path } = await seedStore(records) + await seedHistory(records) + const warn = vi.spyOn(console, 'warn').mockImplementation(() => undefined) + const { runtime, published } = startupRuntime({ profileChats: [CHAT_A, CHAT_B] }) + await runtime.prepareStructuredAgentSessionStartupRestoration() + lock.failing = true + + await expect(runtime.restoreStructuredAgentSessionTabs()).resolves.toBeUndefined() + + expect(published()).toHaveLength(2) + await expectHistory(CHAT_A) + expect(warn).toHaveBeenCalledWith( + '[structured-agent-session] recording restored chat tabs failed', + { sessionIds: [CHAT_A, CHAT_B], error: expect.objectContaining({ code: 'ELOCKED' }) } + ) + expect(await tabIndexOnDisk(path)).toBeUndefined() + }) + }) +}) diff --git a/src/main/runtime/structured-agent-session-tab-index-seed.ts b/src/main/runtime/structured-agent-session-tab-index-seed.ts new file mode 100644 index 00000000000..181d67128a1 --- /dev/null +++ b/src/main/runtime/structured-agent-session-tab-index-seed.ts @@ -0,0 +1,35 @@ +// A profile saved before the store kept a tab index restores its chats from the profile's own tabs. +// Those chats are recorded once, all together: an index written part way would read as complete at +// the next launch and drop the rest. With an index present, every restored chat is already listed. + +import type { StructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-host' + +type TabIndexHost = Pick< + StructuredAgentSessionHost, + 'getPersistedVisibleSessionTabIndex' | 'showSessionTabs' +> + +/** Best-effort: the tabs publish either way, and an index still absent is seeded again next launch. */ +export async function seedStructuredAgentSessionTabIndex( + host: Partial | null | undefined, + targets: readonly string[], + restored: readonly string[] +): Promise { + if (!host?.getPersistedVisibleSessionTabIndex || !host.showSessionTabs) { + return + } + const listed = new Set(host.getPersistedVisibleSessionTabIndex().sessionIds) + const opened = new Set(restored) + // In the restore's own order, so the seeded tabs keep the order the profile gave them. + const unlisted = [...new Set([...targets, ...restored])].filter( + (sessionId) => opened.has(sessionId) && !listed.has(sessionId) + ) + if (unlisted.length > 0) { + await host.showSessionTabs(unlisted).catch((error: unknown) => + console.warn('[structured-agent-session] recording restored chat tabs failed', { + sessionIds: unlisted, + error + }) + ) + } +}