diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts index 0248799980a7..061b66eeb369 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.test.ts @@ -15,6 +15,7 @@ import { performAttach, type AttachFlowInput } from './structured-agent-session- import { AgentSessionJournal } from '../agent-session-journal/journal-store' import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' import * as legacyImport from '../agent-session-journal/journal-legacy-import' +import { restoreStructuredAgentSessionRead } from './structured-agent-session-read-restore' const NOW = 1_800_000_000_000 const SESSION = 'codex_adopting_session' @@ -210,7 +211,7 @@ describe('adopting a provider conversation on create', () => { } ) - it('still releases acquisition and closes the provisional journal on an import write failure', async () => { + it('fails an import write before any spawn and closes the founding journal', async () => { root = await mkdtemp(join(tmpdir(), 'orca-adopt-write-failure-')) const transcriptPath = join(root, 'rollout.jsonl') await writeCodexRollout(transcriptPath, 'valid source') @@ -220,9 +221,29 @@ describe('adopting a provider conversation on create', () => { const close = vi.spyOn(agentSessionJournalCloseRetries, 'closeOrRetain') const sessionAdapter = adapter() await expect(attach(transcriptPath, sessionAdapter)).rejects.toThrow('disk write failed') - expect(sessionAdapter.acquire).toHaveBeenCalledTimes(1) - expect(sessionAdapter.releaseAcquisition).toHaveBeenCalledTimes(1) + // The conversation is founded before a child exists, so nothing was spawned to release. + expect(sessionAdapter.acquire).not.toHaveBeenCalled() + expect(sessionAdapter.releaseAcquisition).not.toHaveBeenCalled() expect(close).toHaveBeenCalledTimes(1) + expect(store?.getRecord(SESSION)?.lease.claimStatus).toBe('released') + }) + + it('keeps the adopted history readable when the first start fails', async () => { + root = await mkdtemp(join(tmpdir(), 'orca-adopt-failed-start-')) + const transcriptPath = join(root, 'rollout.jsonl') + await writeCodexRollout(transcriptPath, 'history before the failed start') + const sessionAdapter = adapter() + vi.mocked(sessionAdapter.acquire).mockRejectedValueOnce(new Error('codex exited (code 1)')) + + await expect(attach(transcriptPath, sessionAdapter)).resolves.toMatchObject({ ok: false }) + + // A send restarts from the record alone, with no adopt source, so the history must already + // be in the conversation the create founded. + const restored = await restoreStructuredAgentSessionRead(store!, root, SESSION) + expect(JSON.stringify(restored?.journal.snapshot().items)).toContain( + 'history before the failed start' + ) + await restored?.journal.close() }) it('prepares a valid source once before acquisition and imports those exact items', async () => { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts index 459d21b1f3b7..922aa9810bf4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-adopted-import.ts @@ -1,7 +1,7 @@ import type { AgentSessionWireRefusal } from '../../../shared/agent-session-wire' import type { AgentSessionRecord } from '../../../shared/agent-session-record' -import type { AgentSessionAttachParams, AttachedJournal } from './structured-agent-session-attach' -import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' +import type { AgentSessionAttachParams } from './structured-agent-session-attach' +import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement' import { importLegacyTranscriptIntoJournal, @@ -58,39 +58,24 @@ async function readAdoptedTranscript( // Import before publication so the first visible chat agrees with the provider's resumed context. export async function importAdoptedTranscript( params: AgentSessionAttachParams, - attached: AttachedJournal, - record: AgentSessionRecord, - prepared: JournalReplacementItem[] | null -): Promise { - try { - await applyAdoptedTranscript(params, attached, record, prepared) - } catch (error) { - // Publication has not taken ownership of this provisional journal yet. - await agentSessionJournalCloseRetries.closeOrRetain(attached.journal) - throw error - } -} - -async function applyAdoptedTranscript( - params: AgentSessionAttachParams, - attached: AttachedJournal, + journal: AgentSessionJournal, record: AgentSessionRecord, prepared: JournalReplacementItem[] | null ): Promise { const adopt = params.adopt // A new journal contains only its epoch row; replay must preserve subsequent durable writes. - if (!adopt || attached.journal.cursor().sequence > 1) { + if (!adopt || journal.cursor().sequence > 1) { return } if (prepared) { - await attached.journal.replaceEpochItems('legacy_import', record.lease.runtimeFence, prepared) + await journal.replaceEpochItems('legacy_import', record.lease.runtimeFence, prepared) return } if (!adopt.transcriptPath) { throw new Error('agent_session_identity_required') } const imported = await importLegacyTranscriptIntoJournal({ - journal: attached.journal, + journal, agent: params.agent, sessionId: adopt.providerHandle.kind === 'claude' diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts index ab9c44949e00..5dc8e0620ac2 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-flow.ts @@ -3,9 +3,10 @@ import { failedAcquisitionRefusal, failedAcquisitionSettlement } from './structured-agent-session-failed-create-refusal' -import type { - StructuredAgentSessionAdapter, - StructuredAgentSessionProviderChildPhase +import { + AgentSessionPreSpawnError, + type StructuredAgentSessionAdapter, + type StructuredAgentSessionProviderChildPhase } from './structured-agent-session-adapter' // The host supplies owner authority; this flow reserves, proves, and publishes the session. @@ -32,15 +33,14 @@ import type { StructuredAgentSessionEventSink } from './structured-agent-session import { resolveAgentSessionReplayOutcome } from './structured-agent-session-replay-outcome' import { readAgentSessionHydrationPage } from './agent-session-history-page' import { acquireOwner } from './structured-agent-session-acquisition' -import { - importAdoptedTranscript, - prepareAdoptedTranscript -} from './structured-agent-session-adopted-import' +import { prepareAdoptedTranscript } from './structured-agent-session-adopted-import' +import { foundAgentSessionConversation } from './structured-agent-session-conversation-founding' import { withAgentSessionCreatePhase, type AgentSessionCreatePhaseRecorder } from '../../observability/agent-session-instrumentation' import type { ProviderHistoryWindow } from '../agent-session-journal/journal-submission-reconciler' +import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement' export type AttachFlowInput = { store: AgentSessionRecordStore @@ -159,6 +159,7 @@ export async function performAttach( ownerAlreadyAdmitted: agentSessionLeaseAdmitsWriter(record.lease) }) if (!agentSessionLeaseAdmitsWriter(record.lease)) { + await foundConversationBeforeSpawn(input, record, preparedTranscript.items) const acquired = await withAgentSessionCreatePhase('acquire_owner', input.recordPhase, () => acquireOwner(input, record) ) @@ -210,7 +211,6 @@ export async function performAttach( adapter: input.adapter, providerHistoryWindow }) - await importAdoptedTranscript(params, attached, record, preparedTranscript.items) await input.onAttached(attached, acquisitionGeneration, acquiredOwner, providerChildPhase) await store.recordOperationOutcome({ callerKey: input.callerKey, @@ -237,6 +237,24 @@ export async function performAttach( } } +/** Nothing has spawned yet, so a failure here settles the reservation as processless. */ +async function foundConversationBeforeSpawn( + input: AttachFlowInput, + record: AgentSessionRecord, + adoptedItems: JournalReplacementItem[] | null +): Promise { + try { + await foundAgentSessionConversation({ + record, + params: input.params, + journalRoot: input.journalRoot, + adoptedItems + }) + } catch (error) { + throw new AgentSessionPreSpawnError(error) + } +} + async function readProviderHistoryWindow(input: { adapter: StructuredAgentSessionAdapter identity: AgentSessionJournalIdentity diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-founding.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-founding.ts new file mode 100644 index 000000000000..90ba0b46a25d --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-founding.ts @@ -0,0 +1,44 @@ +// A conversation exists from its reservation, not from its first live provider child. +// +// The journal used to be opened only once acquisition succeeded, so a create whose child died at +// startup left a record and nothing to read: no chat to publish, and nothing a later send could be +// admitted into. Founding it here makes that create an ordinary readable session with a released +// lease, which a send restarts like any other. + +import { existsSync } from 'node:fs' +import type { AgentSessionRecord } from '../../../shared/agent-session-record' +import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry' +import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement' +import { journalDatabaseFile, journalDirectoryFor } from '../agent-session-journal/journal-paths' +import { openAgentSessionJournal } from '../agent-session-journal/journal-store-factory' +import { importAdoptedTranscript } from './structured-agent-session-adopted-import' +import { + journalIdentityFor, + type AgentSessionAttachParams +} from './structured-agent-session-attach' + +export async function foundAgentSessionConversation(input: { + record: AgentSessionRecord + params: AgentSessionAttachParams + journalRoot: string + adoptedItems: JournalReplacementItem[] | null +}): Promise { + const journalDir = journalDirectoryFor(input.journalRoot, { + workspaceId: input.params.location.workspaceId, + sessionId: input.record.sessionId + }) + // An adopted conversation whose import never landed is still owed it; the import skips a + // journal that already holds more than its epoch. + if (existsSync(journalDatabaseFile(journalDir)) && !input.params.adopt) { + return + } + const journal = await openAgentSessionJournal({ + identity: journalIdentityFor(input.record, input.params), + journalDir + }) + try { + await importAdoptedTranscript(input.params, journal, input.record, input.adoptedItems) + } finally { + await agentSessionJournalCloseRetries.closeOrRetain(journal) + } +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-failed-create-readable.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-failed-create-readable.test.ts new file mode 100644 index 000000000000..fe7046469084 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-failed-create-readable.test.ts @@ -0,0 +1,187 @@ +// A create whose child dies while being acquired still made its conversation. The create answers it +// as a readable chat that says why the agent stopped, and a send restarts the agent through the +// host's ensure-owner step and delivers — the same path as any chat whose child ended. + +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest' +import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' +import type { + AgentSessionAttachResult, + AgentSessionMutationResult +} from '../../../shared/agent-session-wire' +import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' +import { + AgentSessionAcquisitionExitUnprovenError, + AgentSessionAcquisitionRefusal, + type StructuredAgentSessionAdapter +} from './structured-agent-session-adapter' +import { StructuredAgentSessionHost } from './structured-agent-session-host' +import { + HOST_TEST_NOW as NOW, + HOST_TEST_SESSION as SESSION, + HOST_TEST_THREAD as THREAD, + hostTestAttachParams, + hostTestMessage, + hostTestOperationId, + resetHostTestOperationIds +} from './structured-agent-session-host-test-data' + +const CALLER = { callerKey: 'client-1' } +const EXIT_REASON = 'Claude Code is not signed in. Sign in with the Claude CLI' +const STARTUP_ROW = `The provider stopped before it finished starting: ${EXIT_REASON}.` + +let root: string +let store: AgentSessionRecordStore +let host: StructuredAgentSessionHost +let acquire: Mock +let dispatch: Mock +let supported = true + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-failed-create-readable-')) + resetHostTestOperationIds() + supported = true + acquire = vi.fn(async ({ fence, spawnToken }) => ({ + process: { hostId: 'local', pid: 4242, processStartTimeMs: 1_700_000_000_000, spawnToken }, + link: { + linkId: `link-${fence}`, + handle: { provider: 'codex', threadId: THREAD }, + origin: 'created' as const, + mintedAtFence: fence, + observedAt: NOW + } + })) + dispatch = vi.fn(async () => ({ state: 'admitted' as const })) + store = await AgentSessionRecordStore.open({ directory: join(root, 'store'), hostId: 'local' }) + let spawns = 0 + host = new StructuredAgentSessionHost({ + store, + adapter: { + supportsCreate: () => supported, + acquire, + releaseAcquisition: vi.fn(async () => true), + closeSession: vi.fn(async () => true), + dispatch, + cancelTurn: vi.fn(async () => ({ cancelled: true })), + answerPrompt: vi.fn(async () => undefined), + setOption: vi.fn(async () => undefined) + }, + journalRoot: root, + claimKeyId: 'key-1', + mintSpawnToken: () => `spawn-${++spawns}`, + now: () => NOW + }) +}) + +afterEach(async () => { + await host.flushAllStreamedEvents() + await rm(root, { recursive: true, force: true }) +}) + +function statusRows(created: AgentSessionMutationResult): string[] { + return created.ok + ? created.value.page.items.flatMap((item) => + item.body.kind === 'status' ? [item.body.text] : [] + ) + : [] +} + +function send(text: string, fence: number) { + const body = hostTestMessage(text) + return host.send(CALLER, { + envelope: { + sessionId: SESSION, + clientOperationId: hostTestOperationId(), + expectedRuntimeFence: fence, + payloadFingerprint: computeAgentSessionPayloadFingerprint({ + method: 'agentSession.send', + sessionId: SESSION, + fields: { body } + }) + }, + body + }) +} + +describe('a create whose child dies while being acquired', () => { + it('answers the conversation, with the cause in it, and a send restarts the agent and delivers', async () => { + acquire.mockRejectedValueOnce(new Error(EXIT_REASON)) + const params = hostTestAttachParams(null) + + const created = await host.create(CALLER, params) + + expect(created).toMatchObject({ ok: true, value: { sessionId: SESSION } }) + expect(statusRows(created)).toEqual([STARTUP_ROW]) + expect(store.getRecord(SESSION)?.lease.claimStatus).toBe('released') + expect(host.listSessionTabs()).toEqual([ + { sessionId: SESSION, workspaceId: 'workspace-1', agent: 'codex' } + ]) + // A lost reply's replay answers the same conversation and writes no second row. + const replayed = await host.create(CALLER, params) + expect(replayed).toMatchObject({ ok: true, fence: created.ok ? created.fence : -1 }) + expect(statusRows(replayed)).toEqual([STARTUP_ROW]) + expect(acquire).toHaveBeenCalledOnce() + + const sent = await send('hello', created.ok ? created.fence : -1) + + expect(sent).toMatchObject({ ok: true, value: { submission: { dispatchState: 'pending' } } }) + expect(acquire).toHaveBeenCalledTimes(2) + expect(dispatch).toHaveBeenCalledOnce() + expect(store.getRecord(SESSION)?.lease.claimStatus).toBe('live') + }) + + it('keeps a sign-in failure retryable by sending, and delivers once the user has signed in', async () => { + acquire.mockRejectedValueOnce(new AgentSessionAcquisitionRefusal(EXIT_REASON)) + const created = await host.create(CALLER, hostTestAttachParams(null)) + expect(statusRows(created)).toEqual([STARTUP_ROW]) + const fence = created.ok ? created.fence : -1 + + // Still signed out: the restart fails with its cause, and nothing is admitted. + acquire.mockRejectedValueOnce(new AgentSessionAcquisitionRefusal(EXIT_REASON)) + await expect(send('are you there?', fence)).resolves.toMatchObject({ + ok: false, + refusal: { + code: 'agent_session_owner_restart_failed', + message: expect.stringContaining(EXIT_REASON) + } + }) + expect(dispatch).not.toHaveBeenCalled() + + // Signed in: the same kind of send restarts the agent and delivers. + const sent = await send( + 'are you there now?', + store.getRecord(SESSION)?.lease.runtimeFence ?? -1 + ) + expect(sent).toMatchObject({ ok: true }) + expect(acquire).toHaveBeenCalledTimes(3) + expect(dispatch).toHaveBeenCalledOnce() + }) + + it('leaves an unverifiable exit unknown rather than claiming a readable chat', async () => { + acquire.mockRejectedValueOnce(new AgentSessionAcquisitionExitUnprovenError(new Error('hung'))) + + await expect(host.create(CALLER, hostTestAttachParams(null))).rejects.toThrow() + expect(host.hasSession(SESSION)).toBe(false) + }) +}) + +describe('a create refused before its reservation', () => { + it('stays a refusal with no conversation, so the client offers a fresh create', async () => { + supported = false + + await expect(host.create(CALLER, hostTestAttachParams(null))).resolves.toEqual({ + ok: false, + refusal: expect.objectContaining({ code: 'structured_agent_session_unsupported' }) + }) + expect(store.getRecord(SESSION)).toBeNull() + expect(host.hasSession(SESSION)).toBe(false) + + supported = true + await expect(host.create(CALLER, hostTestAttachParams(null))).resolves.toMatchObject({ + ok: true + }) + expect(acquire).toHaveBeenCalledOnce() + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-failed-create.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-failed-create.ts new file mode 100644 index 000000000000..63e64b0cc6b4 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-failed-create.ts @@ -0,0 +1,91 @@ +// A create makes a conversation and tries to start its agent. When only the start fails, the +// conversation still exists: its journal was founded at reservation. The create answers it as a +// readable chat whose first row says why the agent stopped, and a send restarts the agent through +// the host's ensure-owner step like any chat whose child ended — there is no second restart path. +// +// Only a start whose child is proven gone qualifies. A definitive refusal still answers as one, so +// the client falls back instead; an unverifiable one still answers as unknown, so the client +// reconciles instead. + +import type { + AgentSessionAttachResult, + AgentSessionMutationResult, + AgentSessionWireRefusal +} from '../../../shared/agent-session-wire' +import { isDefinitiveAgentSessionCreateRefusal } from '../../../shared/agent-session-definitive-refusal' +import { readAgentSessionHydrationPage } from './agent-session-history-page' +import type { AgentSessionAttachParams } from './structured-agent-session-attach' +import { providerStartupFailureOutcome } from './structured-agent-session-dead-generation-settlement' +import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' +import { recordStructuredAgentSessionStartFailure } from './structured-agent-session-start-failure-row' + +type FailedCreateContext = Pick< + StructuredAgentSessionMutationContext, + 'deps' | 'sessions' | 'restoreReadable' | 'publish' +> +type CreateResult = AgentSessionMutationResult + +export async function answerStructuredAgentSessionCreate( + attached: Promise, + params: AgentSessionAttachParams, + /** Runs the read inside the session's serialize. */ + serialized: ( + read: (context: FailedCreateContext) => Promise + ) => Promise +): Promise { + const result = await attached + if (result.ok) { + return result + } + const { sessionId, clientOperationId } = params.envelope + const created = await serialized((context) => + readFailedCreate(context, { + sessionId, + operationId: clientOperationId, + refusal: result.refusal + }) + ) + return created ?? result +} + +/** Null leaves the create's refusal as its answer. */ +async function readFailedCreate( + context: FailedCreateContext, + input: { sessionId: string; operationId: string; refusal: AgentSessionWireRefusal } +): Promise { + const { sessionId, refusal } = input + if (refusal.ownerVerdict !== 'exited' || isDefinitiveAgentSessionCreateRefusal(refusal.code)) { + return null + } + let session: Awaited> + try { + session = await recordStructuredAgentSessionStartFailure( + context, + sessionId, + `failed-start:${input.operationId}`, + providerStartupFailureOutcome(refusal.message) + ) + } catch (error) { + // Bookkeeping: the refusal still reaches the user, whose Retry is a fresh create. + context.deps.onEventSinkError?.({ sessionId, error }) + return null + } + if (!session) { + return null + } + const { fence, journal } = session + const tabId = context.deps.store.getRecord(sessionId)?.surfaceTabId + return { + ok: true, + replayed: false, + fence, + cursor: journal.cursor(), + value: { + sessionId, + fence, + page: readAgentSessionHydrationPage(journal, fence), + unconfirmedClientMessageIds: [], + ...(tabId ? { tabId } : {}) + } + } +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts index 4773ea5eab8d..e035c836b73d 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts @@ -1,5 +1,5 @@ // Everything a client can ask an ATTACHED session to do: send a turn, cancel one, answer a prompt, -// change an option, read the options back. +// change an option, read the options back — and the answer a create gives when only its start failed. // // They share one shape — admit the envelope against the lease, run a plan, publish the journal — so // they share one path here rather than five copies in the host. The host keeps attach, holds and @@ -23,6 +23,8 @@ import type { } from '../../../shared/agent-session-wire' import type { StructuredAgentSessionHolds } from './structured-agent-session-holds' import { threadGoalPlan } from './structured-agent-session-thread-goal' +import type { AgentSessionAttachParams } from './structured-agent-session-attach' +import { answerStructuredAgentSessionCreate } from './structured-agent-session-failed-create' import { admitAndRunAgentSessionMutation, type AgentSessionMutationRequest @@ -288,6 +290,18 @@ export function structuredAgentSessionMutationDelegates( caller: StructuredAgentSessionCaller, params: Parameters[2] ) => changeStructuredAgentSessionThreadGoal(context(), caller, params), - readOptions: (sessionId: string) => readStructuredAgentSessionOptions(context(), sessionId) + readOptions: (sessionId: string) => readStructuredAgentSessionOptions(context(), sessionId), + settleLateDispatch: (input: Parameters[1]) => + settleStructuredAgentSessionLateDispatch(context(), input), + releaseUnansweredDispatches: ( + input: Parameters[1] + ) => releaseStructuredAgentSessionUnansweredDispatches(context(), input), + answerCreate: ( + attached: Parameters[0], + params: AgentSessionAttachParams + ) => + answerStructuredAgentSessionCreate(attached, params, (read) => + context().serialize(params.envelope.sessionId, () => read(context())) + ) } } 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 a2b03e120cac..bd9c3d14ceb4 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 @@ -34,9 +34,7 @@ import type { StructuredAgentSessionAttachContext } from './structured-agent-ses import { listStructuredAgentSessionTabs } from './structured-agent-session-host-tabs' import { structuredAgentSessionMutationDelegates, - settleStructuredAgentSessionLateDispatch, - type StructuredAgentSessionMutationContext, - releaseStructuredAgentSessionUnansweredDispatches + type StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' import { flushStructuredAgentSessionHost } from './structured-agent-session-host-teardown' import type { @@ -251,6 +249,10 @@ export class StructuredAgentSessionHost { return attachStructuredAgentSession(this.attachContext(), caller.callerKey, params) } + /** A create's answer: the attach's, or the conversation it made when only its start failed. */ + create = (caller: StructuredAgentSessionCaller, params: AgentSessionAttachParams) => + this.mutations.answerCreate(this.attach(caller, params), params) + flushStreamedEvents = (sessionId: string): Promise => this.runtimeState.flushEventSink(sessionId) @@ -329,12 +331,8 @@ export class StructuredAgentSessionHost { subscribe = (input: AgentSessionSubscribeInput): (() => void) => this.backgroundTasks.subscribe(input) - settleLateDispatch = (input: Parameters[1]) => - settleStructuredAgentSessionLateDispatch(this.mutationContext(), input) - - releaseUnansweredDispatches = ( - input: Parameters[1] - ) => releaseStructuredAgentSessionUnansweredDispatches(this.mutationContext(), input) + settleLateDispatch = this.mutations.settleLateDispatch + releaseUnansweredDispatches = this.mutations.releaseUnansweredDispatches publishBackgroundTaskState: StructuredAgentSessionBackgroundTaskChannel['publish'] = (...args) => this.backgroundTasks.publish(...args) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts index e4a332f7acf8..e68e73394b9b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts @@ -30,7 +30,6 @@ import type { AgentSessionWireRefusal } from '../../../shared/agent-session-wire' import type { AgentSessionWireRefusalCode } from '../../../shared/agent-session-wire-refusals' -import { boundJournalStatusText } from '../agent-session-journal/journal-prompt-body-bounds' import { TUI_AGENT_DISPLAY_NAMES } from '../../../shared/tui-agent-display-names' import { ownerRestartFailedOutcome, @@ -41,6 +40,7 @@ import type { StructuredAgentSessionMutationContext } from './structured-agent-s import type { AgentSessionMutationSessionPreparation } from './structured-agent-session-mutation-admission' import { isResumableStructuredAgentSessionRecord } from './structured-agent-session-resume-eligibility' import { rewindRefusal } from './structured-rewind-refusal' +import { recordStructuredAgentSessionStartFailure } from './structured-agent-session-start-failure-row' /** * What a refused resume means for the send that ran it. `transient`: the resume met a lease @@ -194,9 +194,7 @@ function ownerRestartFailedRefusal( } } -/** The same status row a start that failed leaves in the chat, so the reason outlives the error - * strip. The journal is made readable for it when the failed attach left none behind. Keyed by - * the send, not the clock: a resend of the same id that fails again adds no second row. */ +/** The same status row a start that failed leaves in the chat, keyed by the send. */ async function recordFailedRestart( context: SendPreparationContext, envelope: AgentSessionMutationEnvelope, @@ -204,27 +202,12 @@ async function recordFailedRestart( ): Promise { const { sessionId } = envelope try { - if (!context.sessions.has(sessionId)) { - await context.restoreReadable(sessionId) - } - const session = context.sessions.get(sessionId) - if (!session) { - return - } - const settlementId = `failed-restart:${envelope.clientOperationId}` - await session.journal.appendLifecycleBatch({ - settlementId, - fence: session.fence, - recovered: true, - mutations: [ - { - kind: 'item', - identity: { provider: 'orca', clientMessageId: settlementId }, - body: { kind: 'status', text: boundJournalStatusText(text) } - } - ] - }) - context.publish(sessionId, session.journal) + await recordStructuredAgentSessionStartFailure( + context, + sessionId, + `failed-restart:${envelope.clientOperationId}`, + text + ) } catch (error) { context.deps.onEventSinkError?.({ sessionId, error }) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-start-failure-row.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-start-failure-row.ts new file mode 100644 index 000000000000..167441f65fef --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-start-failure-row.ts @@ -0,0 +1,38 @@ +import { boundJournalStatusText } from '../agent-session-journal/journal-prompt-body-bounds' +import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types' +import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' + +/** + * The status row a start that failed leaves in the chat, so the reason outlives any error strip. + * The journal is made readable for it when no session holds it. Keyed by the operation that met + * the failure, not the clock, so a replay of that operation adds no second row. Answers the session + * it wrote to, or null when the conversation has nothing to read. + */ +export async function recordStructuredAgentSessionStartFailure( + context: Pick, + sessionId: string, + settlementId: string, + text: string +): Promise { + if (!context.sessions.has(sessionId)) { + await context.restoreReadable(sessionId) + } + const session = context.sessions.get(sessionId) + if (!session) { + return null + } + await session.journal.appendLifecycleBatch({ + settlementId, + fence: session.fence, + recovered: true, + mutations: [ + { + kind: 'item', + identity: { provider: 'orca', clientMessageId: settlementId }, + body: { kind: 'status', text: boundJournalStatusText(text) } + } + ] + }) + context.publish(sessionId, session.journal) + return session +} diff --git a/src/main/runtime/rpc/methods/orchestration-structured-worker-failed-start.test.ts b/src/main/runtime/rpc/methods/orchestration-structured-worker-failed-start.test.ts new file mode 100644 index 000000000000..ec54c92161f0 --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration-structured-worker-failed-start.test.ts @@ -0,0 +1,87 @@ +// A chat surface gets a readable conversation when its agent fails to start; a worker launch needs a +// running agent, so its create still answers the start's refusal with the reason, after one spawn. + +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest' +import { AgentSessionRecordStore } from '../../agent-session-record-store' +import type { StructuredAgentSessionAdapter } from '../../../native-chat/agent-session-wire/structured-agent-session-adapter' +import { StructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-host' +import { setStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry' +import { + HOST_TEST_NOW as NOW, + hostTestAttachParams +} from '../../../native-chat/agent-session-wire/structured-agent-session-host-test-data' +import type { OrcaRuntimeService } from '../../orca-runtime' +import { createStructuredWorkerSession } from './orchestration-structured-worker-session' + +const EXIT_REASON = 'Codex is not signed in. Run codex login' + +let root: string +let host: StructuredAgentSessionHost +let acquire: Mock + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-worker-failed-start-')) + // The worker mints its operation id from `Date.now()`; the host expires ids against its own clock. + vi.spyOn(Date, 'now').mockReturnValue(NOW) + acquire = vi.fn(async () => { + throw new Error(EXIT_REASON) + }) + const store = await AgentSessionRecordStore.open({ + directory: join(root, 'store'), + hostId: 'local' + }) + host = new StructuredAgentSessionHost({ + store, + adapter: { + acquire, + releaseAcquisition: vi.fn(async () => true), + closeSession: vi.fn(async () => true), + dispatch: vi.fn(async () => ({ state: 'admitted' as const })), + cancelTurn: vi.fn(async () => ({ cancelled: true })), + answerPrompt: vi.fn(async () => undefined), + setOption: vi.fn(async () => undefined) + }, + journalRoot: root, + claimKeyId: 'key-1', + mintSpawnToken: () => 'spawn-1', + now: () => NOW + }) + setStructuredAgentSessionHost(host) +}) + +afterEach(async () => { + vi.restoreAllMocks() + setStructuredAgentSessionHost(null) + await host.flushAllStreamedEvents() + await rm(root, { recursive: true, force: true }) +}) + +describe('a structured worker whose agent fails to start', () => { + it('is refused with the start’s own reason, after one spawn and no readable chat', async () => { + const { envelope: _envelope, ...resolved } = hostTestAttachParams(null) + const publishStructuredAgentSessionTab = vi.fn(async () => undefined) + const runtime = { + ensureStructuredAgentSessionHost: async () => {}, + resolveStructuredAgentSessionCreateIntent: async () => resolved, + publishStructuredAgentSessionTab, + retireStructuredAgentSessionTabFromSnapshot: () => {} + } + + await expect( + createStructuredWorkerSession({ + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the create path reaches only the four runtime members stubbed above; a missing one fails the call. + runtime: runtime as unknown as OrcaRuntimeService, + worktreeId: 'wt_1', + agent: 'codex', + dispatchId: 'd_failed_start', + onJournalActivity: () => {} + }) + ).rejects.toThrow(`The structured codex session for this worker was refused: ${EXIT_REASON}`) + expect(acquire).toHaveBeenCalledOnce() + expect(publishStructuredAgentSessionTab).not.toHaveBeenCalled() + expect(host.listSessionTabs()).toEqual([]) + }) +}) diff --git a/src/main/runtime/rpc/methods/structured-agent-session-create.ts b/src/main/runtime/rpc/methods/structured-agent-session-create.ts index ab52c2aa855c..6b830511c5e4 100644 --- a/src/main/runtime/rpc/methods/structured-agent-session-create.ts +++ b/src/main/runtime/rpc/methods/structured-agent-session-create.ts @@ -125,9 +125,15 @@ export async function commitStructuredAgentSessionCreate(args: { caller: StructuredAgentSessionCaller prepared: PreparedStructuredAgentSessionCreate activate: boolean + /** A chat surface gets a readable conversation even when its agent failed to start, and restarts + * it by sending; an in-process launch needs a running agent, so a failed start stays a refusal. */ + answer: 'conversation' | 'running-agent' }): Promise> { const { prepared } = args - const result = await prepared.host.attach(args.caller, prepared.attachParams) + const result = + args.answer === 'conversation' + ? await prepared.host.create(args.caller, prepared.attachParams) + : await prepared.host.attach(args.caller, prepared.attachParams) if (!result.ok || !prepared.tab) { return result } @@ -173,6 +179,7 @@ export async function createStructuredAgentSessionForWorktree(args: { runtime: args.runtime, caller: args.caller, prepared, - activate: args.activate + activate: args.activate, + answer: 'running-agent' }) } diff --git a/src/main/runtime/rpc/methods/structured-agent-session-precommit-refusal.test.ts b/src/main/runtime/rpc/methods/structured-agent-session-precommit-refusal.test.ts index f34ecb5d6cd8..353599c4c141 100644 --- a/src/main/runtime/rpc/methods/structured-agent-session-precommit-refusal.test.ts +++ b/src/main/runtime/rpc/methods/structured-agent-session-precommit-refusal.test.ts @@ -49,7 +49,9 @@ function hostStub(): StructuredAgentSessionHost { cursor: { epoch: 'epoch-a', sequence: 0 }, value: { sessionId: SESSION, fence: 1, page: {}, unconfirmedClientMessageIds: [] } })) - return { attach } as unknown as StructuredAgentSessionHost + // A create answers its attach here; only a failed start's readable chat differs, host-side. + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the create route reaches only `attach`/`create`, both stubbed; a missing member fails the call. + return { attach, create: attach } as unknown as StructuredAgentSessionHost } const resolvedIntent = { diff --git a/src/main/runtime/rpc/methods/structured-agent-session-rpc.test-fixture.ts b/src/main/runtime/rpc/methods/structured-agent-session-rpc.test-fixture.ts index 9f698a447418..67b9fe61d876 100644 --- a/src/main/runtime/rpc/methods/structured-agent-session-rpc.test-fixture.ts +++ b/src/main/runtime/rpc/methods/structured-agent-session-rpc.test-fixture.ts @@ -123,34 +123,37 @@ function statusFeed(): StructuredAgentSessionStatusFeed { export function hostStub(): StructuredAgentSessionHost { reset(hostCalls) - Object.assign(hostCalls, { - attach: vi.fn(async () => ({ - ok: true, - replayed: false, + const attach = vi.fn(async (..._args: unknown[]) => ({ + ok: true, + replayed: false, + fence: 1, + cursor: { epoch: 'epoch-a', sequence: 0 }, + value: { + sessionId: SESSION, fence: 1, - cursor: { epoch: 'epoch-a', sequence: 0 }, - value: { + page: { sessionId: SESSION, - fence: 1, - page: { - sessionId: SESSION, - epoch: 'epoch-a', - direction: 'tail', - items: [], - removedItemIds: [], - submissions: [], - window: { - oldest: null, - newest: null, - nextCursor: { epoch: 'epoch-a', sequence: 0 } - }, - liveCursor: { epoch: 'epoch-a', sequence: 0 }, - hasOlder: false, - hasNewer: false + epoch: 'epoch-a', + direction: 'tail', + items: [], + removedItemIds: [], + submissions: [], + window: { + oldest: null, + newest: null, + nextCursor: { epoch: 'epoch-a', sequence: 0 } }, - unconfirmedClientMessageIds: [] - } - })), + liveCursor: { epoch: 'epoch-a', sequence: 0 }, + hasOlder: false, + hasNewer: false + }, + unconfirmedClientMessageIds: [] + } + })) + Object.assign(hostCalls, { + attach, + // A create whose attach succeeded answers the attach; the failed-start answer is the host's. + create: vi.fn((...args: unknown[]) => attach(...args)), rewind: vi.fn(async () => ({ ok: true, value: { itemId: 'chosen', epoch: 'next' } })), send: vi.fn(async () => ({ ok: true, diff --git a/src/main/runtime/rpc/methods/structured-agent-session.ts b/src/main/runtime/rpc/methods/structured-agent-session.ts index 3b2fde22b771..23b92819a7c2 100644 --- a/src/main/runtime/rpc/methods/structured-agent-session.ts +++ b/src/main/runtime/rpc/methods/structured-agent-session.ts @@ -177,7 +177,8 @@ export const STRUCTURED_AGENT_SESSION_METHODS = [ runtime: ctx.runtime, caller: callerFor(ctx), prepared, - activate: true + activate: true, + answer: 'conversation' }) } }), diff --git a/src/renderer/src/components/native-chat/NativeChatStructuredSession.launch-lifecycle.test.tsx b/src/renderer/src/components/native-chat/NativeChatStructuredSession.launch-lifecycle.test.tsx index 859b93824f32..895f16d16e9e 100644 --- a/src/renderer/src/components/native-chat/NativeChatStructuredSession.launch-lifecycle.test.tsx +++ b/src/renderer/src/components/native-chat/NativeChatStructuredSession.launch-lifecycle.test.tsx @@ -121,7 +121,10 @@ describe('NativeChatStructuredSession launch lifecycle', () => { ) }) - it('relaunches a failed start on send, then delivers the message once it publishes', async () => { + // A host answers a create whose agent died starting as a published chat, and a send restarts + // the agent there. A launch that still reads failed was refused before the host made anything + // (or by a host that predates that answer), so only launch Retry, a fresh create, recovers it. + it('parks a send into a failed launch until launch Retry publishes it, without relaunching', async () => { mocks.mode = 'outbox' mocks.launchLifecycle = 'failed' mocks.call.mockResolvedValue({ @@ -130,11 +133,12 @@ describe('NativeChatStructuredSession launch lifecycle', () => { }) const { rerender } = render(sessionView()) - expect(composerSend()('restart and say hi', [])).toBe(true) - // The relaunch is launch Retry's own: a new create operation under the same session. - expect(mocks.retryLaunch).toHaveBeenCalledExactlyOnceWith('wt-1', 'session-1') + expect(composerSend()('say hi once it starts', [])).toBe(true) + expect(mocks.retryLaunch).not.toHaveBeenCalled() expect(mocks.call).not.toHaveBeenCalled() + fireEvent.click(screen.getByRole('button', { name: 'Retry' })) + expect(mocks.retryLaunch).toHaveBeenCalledExactlyOnceWith('wt-1', 'session-1') mocks.launchLifecycle = 'published' rerender(sessionView()) await waitFor(() => expect(mocks.call).toHaveBeenCalledOnce()) @@ -145,7 +149,7 @@ describe('NativeChatStructuredSession launch lifecycle', () => { ) }) - it('keeps the message queued with the reason shown when the relaunch fails again', async () => { + it('keeps the message queued with the reason shown while the launch stays failed', async () => { mocks.mode = 'outbox' mocks.launchLifecycle = 'failed' mocks.call.mockResolvedValue({ diff --git a/src/renderer/src/components/native-chat/NativeChatStructuredSession.tsx b/src/renderer/src/components/native-chat/NativeChatStructuredSession.tsx index 6191febf5e3a..20e2a8f866c3 100644 --- a/src/renderer/src/components/native-chat/NativeChatStructuredSession.tsx +++ b/src/renderer/src/components/native-chat/NativeChatStructuredSession.tsx @@ -39,7 +39,6 @@ export function NativeChatStructuredSession( fileLinkContext?.worktreeId, props.sessionId ) - const { sendThroughRelaunch } = provisionalLaunch // The host's own word on whether the provider child has answered startup yet. const startupPhase = useStructuredAgentSessionHostExecutionPhase(props.sessionId, props.target) const controller = useStructuredAgentSession({ @@ -166,14 +165,12 @@ export function NativeChatStructuredSession( : null return { send: (text: string, attachments: readonly { id: string; path: string }[]): boolean => - sendThroughRelaunch(() => - controller.send( - text, - attachments.map((attachment) => ({ - path: attachment.path, - previewUri: attachment.path - })) - ) + controller.send( + text, + attachments.map((attachment) => ({ + path: attachment.path, + previewUri: attachment.path + })) ), dispatchCommand: (text: string) => dispatchStructuredAgentSessionComposerCommand(text, { @@ -208,8 +205,7 @@ export function NativeChatStructuredSession( optionPickerRequest, props.agent, props.sessionId, - props.target, - sendThroughRelaunch + props.target ]) return ( diff --git a/src/renderer/src/components/native-chat/use-native-chat-provisional-launch.ts b/src/renderer/src/components/native-chat/use-native-chat-provisional-launch.ts index 800e56476f02..ce790cf7c6c2 100644 --- a/src/renderer/src/components/native-chat/use-native-chat-provisional-launch.ts +++ b/src/renderer/src/components/native-chat/use-native-chat-provisional-launch.ts @@ -1,6 +1,5 @@ import { useCallback } from 'react' import { - getStructuredAgentSessionLaunchLifecycle, retryStructuredAgentSessionLaunch, useStructuredAgentSessionLaunchFailureReason, useStructuredAgentSessionLaunchLifecycle @@ -17,26 +16,10 @@ export function useNativeChatProvisionalLaunch( retryStructuredAgentSessionLaunch(worktreeId, sessionId) } }, [sessionId, worktreeId]) - // A send into a start that never published relaunches it; the queued message goes out on publish. - const sendThroughRelaunch = useCallback( - (send: () => boolean): boolean => { - const accepted = send() - if ( - accepted && - worktreeId && - getStructuredAgentSessionLaunchLifecycle(worktreeId, sessionId) === 'failed' - ) { - retryStructuredAgentSessionLaunch(worktreeId, sessionId) - } - return accepted - }, - [sessionId, worktreeId] - ) return { lifecycle, failureReason, retry, - sendThroughRelaunch, transportEnabled: lifecycle === null || lifecycle === 'published' } } diff --git a/src/renderer/src/lib/structured-agent-session-launch-exited-owner.test.ts b/src/renderer/src/lib/structured-agent-session-launch-exited-owner.test.ts index 1662e4c04300..1b9d38d4fd9c 100644 --- a/src/renderer/src/lib/structured-agent-session-launch-exited-owner.test.ts +++ b/src/renderer/src/lib/structured-agent-session-launch-exited-owner.test.ts @@ -1,6 +1,8 @@ // @vitest-environment happy-dom // The launch client against the create wire contract: a failed create whose provider is proven // gone reads as failed (with its reason) and Retry starts it again; any other verdict stays unknown. +// A current host answers such a create as a readable chat instead; this refusal is what a host that +// predates that answer sends, and launch Retry is how a newer client recovers from it. import { beforeEach, describe, expect, it, vi } from 'vitest' import type { diff --git a/tests/e2e/cross-version-wire/structured-agent-session-host-fixture.ts b/tests/e2e/cross-version-wire/structured-agent-session-host-fixture.ts index 0de322c1b1d3..15dc434f6527 100644 --- a/tests/e2e/cross-version-wire/structured-agent-session-host-fixture.ts +++ b/tests/e2e/cross-version-wire/structured-agent-session-host-fixture.ts @@ -102,6 +102,9 @@ export function installableHost( ): StructuredAgentSessionHost { const host = { ...hostCalls, + // A create answers its attach; the real host adds only a failed start's readable chat, on a + // refusal. Reassembled like `restartResume`, so the manifest still names `attach`. + create: (...args: unknown[]) => hostCalls.attach(...args), restartResume: { list: hostCalls.restartResumableList, listFailures: hostCalls.restartResumableFailures,