diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-agent-start.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-agent-start.ts index 6fbcca0dce2..d2ef237aa1b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-agent-start.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-agent-start.ts @@ -28,6 +28,10 @@ import type { StructuredAgentSessionAttachContext } from './structured-agent-ses import { attachStructuredAgentSessionUnderSerialize } from './structured-agent-session-attach-orchestration' import { failedCreateRefusal } from './structured-agent-session-failed-create-refusal' import { adapterSupportsRecord } from './structured-agent-session-provider-support' +import { + finishOwedStructuredAgentSessionWindDownUnderSerialize, + type StructuredAgentSessionLifetimeContext +} from './structured-agent-session-host-lifetime' import { structuredAgentSessionResumeOperationId, structuredAgentSessionResumeParams @@ -48,6 +52,48 @@ export type StructuredAgentSessionResumeOutcome = /** The attach's caller key: the ledger row a start settles is Orca's own. */ const AGENT_START_CALLER_KEY = 'trusted-local:agent-start' +/** + * Every operation that reaches the provider finishes a stop an earlier attempt left owed first: the + * child that stop could not prove gone takes no input, and none may start beside it. Still + * unproven, the operation is refused with the exit `unverifiable`, never assumed `exited`. + */ +export async function finishOwedStructuredAgentSessionStop( + context: StructuredAgentSessionLifetimeContext, + sessionId: string +): Promise { + if (await finishOwedStructuredAgentSessionWindDownUnderSerialize(context, sessionId)) { + return { ok: true } + } + return { + ok: false, + refusal: refuse( + 'agent_session_ownership_unknown', + { reason: 'previousExitUnverifiable', ownerVerdict: 'unverifiable' }, + "Orca could not prove this chat's previous agent process exited." + ) + } +} + +/** The same, for an operation the running child performs, which starts none: only an exit still + * unproven, its child still on record, holds it back. A proven exit owes only bookkeeping, which + * gates no such operation: its retry was reported, and the operation goes on as with none owed. */ +export async function finishOwedStructuredAgentSessionStopForProviderWrite( + context: StructuredAgentSessionLifetimeContext, + sessionId: string +): Promise { + const settled = await finishOwedStructuredAgentSessionStop(context, sessionId) + return settled.ok || context.sessions.get(sessionId)?.child ? settled : { ok: true } +} + +export function isStructuredAgentSessionPreviousExitUnverifiable( + refusal: AgentSessionWireRefusal +): boolean { + return ( + refusal.code === 'agent_session_ownership_unknown' && + refusal.details?.reason === 'previousExitUnverifiable' + ) +} + /** * Gives the session a provider child if it has none, for a caller inside its serialize — which is * what makes "if it has none" exact: two askers run this in turn, and the second finds the first @@ -58,6 +104,10 @@ export async function ensureStructuredAgentSessionAgent( sessionId: string, startedFor?: string ): Promise { + const settled = await finishOwedStructuredAgentSessionStop(context, sessionId) + if (!settled.ok) { + return settled + } if (context.sessions.get(sessionId)?.child) { return { ok: true } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-claude-unproven-stop-send.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-claude-unproven-stop-send.test.ts new file mode 100644 index 00000000000..b86e2567464 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-claude-unproven-stop-send.test.ts @@ -0,0 +1,569 @@ +// A Claude Stop whose close could not prove the child gone, then the user's next message, on the +// shipping adapter and host. That child takes no input, so the send retries the stop first; still +// unproven, the message waits with its reason until a later retry proves the exit. + +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' +import { AGENT_JOURNAL_THREAD_SCOPE } from '../../../shared/agent-session-journal-types' +import { activeStructuredAgentSessionTurnId } from '../../../shared/structured-agent-session-live-turn' +import { claudeUnwrittenUserMessageError } from '../../claude/claude-agent-sdk-user-message-queue' +import { ClaudeStructuredSessionAdapter } from '../../claude/claude-structured-session-adapter' +import { + fakeClaude, + PROVIDER_SESSION_ID, + type FakeConnection +} from '../../claude/claude-structured-session-test-support' +import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' +import { openTestAgentSessionRecordStore } from '../../runtime/agent-session-record-store-test-harness' +import { structuredClaudeLifecycleEvent } from '../../runtime/structured-claude-runtime-adapter' +import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support' +import { StructuredAgentSessionHost } from './structured-agent-session-host' +import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support' +import { + HOST_TEST_NOW as NOW, + HOST_TEST_SESSION as SESSION, + hostTestAttachParams, + hostTestMessage, + hostTestOperationId, + resetHostTestOperationIds +} from './structured-agent-session-host-test-data' + +const CALLER = { callerKey: 'client-1' } +const CAPABILITIES = ['interrupt_receipt_v1', 'interrupt_cancel_queued_v1', 'msg_lifecycle_v1'] + +let root: string +let host: StructuredAgentSessionHost +let adapter: ClaudeStructuredSessionAdapter +let store: AgentSessionRecordStore +let claude: ReturnType +let log: ReturnType + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-claude-unproven-stop-send-')) + resetHostTestOperationIds() + log = recordingStructuredAgentSessionLogger() + claude = fakeClaude({ + replayUuid: null, + routes: { interrupt: () => ({ still_queued: [], cancelled: [] }) } + }) + const lifecycle: Promise[] = [] + adapter = new ClaudeStructuredSessionAdapter({ + resolveLaunch: async () => ({ + pathToClaudeCodeExecutable: 'claude', + options: {}, + cwd: root, + claudeConfigDir: join(root, 'claude-home'), + providerSessionId: PROVIDER_SESSION_ID, + resumeLeafUuid: null, + resumesTranscript: (store.getRecord(SESSION)?.providerHandleChain.length ?? 0) > 0, + continuesChain: (store.getRecord(SESSION)?.providerHandleChain.length ?? 0) > 0 + }), + onEvent: (event) => { + const mapped = structuredClaudeLifecycleEvent(event) + if (mapped) { + lifecycle.push(host.handleAdapterEvent(mapped)) + } + }, + onDispatchSettledLate: (settlement) => void host.settleLateDispatch(settlement), + openConnection: claude.openConnection, + readProcessStartTime: async () => 1_700_000_000_000, + now: () => NOW + }) + store = await openTestAgentSessionRecordStore(root) + host = new StructuredAgentSessionHost({ + store, + adapter: Object.assign(adapter, { supportsCreate: () => true }), + journalDatabase: openTestJournalHostDatabase(root), + claimKeyId: 'key-1', + mintSpawnToken: () => 'spawn-a', + logger: log.logger, + now: () => NOW + }) + const params = hostTestAttachParams(null, { + provider: 'claude', + agent: 'claude', + accountHome: { variable: 'CLAUDE_CONFIG_DIR', path: join(root, 'claude-home') }, + providerHandle: { kind: 'claude', sessionId: PROVIDER_SESSION_ID, leafUuid: null } + }) + expect(await host.attach(CALLER, params)).toMatchObject({ ok: true }) + await adapter.awaitStarted(SESSION) + await Promise.all(lifecycle) +}) + +afterEach(async () => { + await adapter.closeAll() + await host.flushAllStreamedEvents() + await rm(root, { recursive: true, force: true }) +}) + +function eventually(assertion: () => T | Promise): Promise { + return vi.waitFor(assertion, { timeout: 10_000 }) +} + +function envelope( + method: + | 'agentSession.send' + | 'agentSession.cancel' + | 'agentSession.setOption' + | 'agentSession.queuedMessageSend', + // The fingerprint's own field shape, as the sibling host tests type it. + fields: Parameters[0]['fields'] +) { + return { + sessionId: SESSION, + clientOperationId: hostTestOperationId(), + expectedRuntimeFence: store.getRecord(SESSION)!.lease.runtimeFence, + payloadFingerprint: computeAgentSessionPayloadFingerprint({ + method, + sessionId: SESSION, + fields + }) + } +} + +async function send(text: string): Promise { + const body = hostTestMessage(text) + const sent = await host.send(CALLER, { envelope: envelope('agentSession.send', { body }), body }) + if (!sent.ok) { + throw new Error(`send refused: ${JSON.stringify(sent.refusal)}`) + } + return sent.value.clientMessageId +} + +async function submission(clientMessageId: string) { + await host.flushStreamedEvents(SESSION) + return (await host.journalSnapshot(SESSION)).submissions.find( + (entry) => entry.clientMessageId === clientMessageId + ) +} + +function frame(connection: FakeConnection, message: Record): void { + connection.handlers.onMessage?.({ session_id: PROVIDER_SESSION_ID, ...message }) +} + +function wrote(connection: FakeConnection, text: string): boolean { + return connection.sent.some((message) => JSON.stringify(message).includes(text)) +} + +async function openTurn(connection: FakeConnection): Promise { + const text = 'Write a long reply.' + const clientMessageId = await send(text) + await eventually(() => expect(wrote(connection, text)).toBe(true)) + frame(connection, { + type: 'system', + subtype: 'init', + uuid: 'init-1', + model: 'claude-sonnet-5', + capabilities: CAPABILITIES + }) + const written = connection.sent.at(-1)! + frame(connection, { ...written, uuid: written.uuid }) + frame(connection, { + type: 'assistant', + uuid: 'stopped-turn-leaf', + parent_tool_use_id: null, + message: { id: 'msg-1', role: 'assistant', content: [{ type: 'text', text: 'Working on' }] } + }) + await eventually(async () => + expect((await submission(clientMessageId))?.dispatchState).toBe('accepted') + ) + expect( + activeStructuredAgentSessionTurnId((await host.journalSnapshot(SESSION)).items) + ).not.toBeNull() +} + +function stop() { + return host.cancel(CALLER, { envelope: envelope('agentSession.cancel', {}) }) +} + +function laneDrained(): Promise { + return host['tasks'].serialize(SESSION, async () => {}) +} + +/** A commit queues a serialized wake, and the wake queues its delivery step behind that. */ +async function commitSettled(): Promise { + await laneDrained() + await laneDrained() +} + +/** As the real connection: once a close begins it refuses every write, proven or not. */ +function closeUnprovenFor(connection: FakeConnection, failures: number): void { + const close = connection.close + let left = failures + connection.close = async () => { + if (left === 0) { + return close() + } + left -= 1 + connection.closeCount += 1 + connection.closed = true + return false + } + const write = connection.send + connection.send = (message, beforeDispatch) => + connection.closed + ? Promise.reject( + claudeUnwrittenUserMessageError(new Error('claude stream-json connection is closed')) + ) + : write(message, beforeDispatch) +} + +const INTERRUPTED_RESULT = { + type: 'result', + subtype: 'error_during_execution', + is_error: true, + terminal_reason: 'aborted_streaming', + uuid: 'interrupted-result' +} + +/** Stops a turn whose child's close cannot prove the exit `failures` times, and lets the Stop's + * second step fail on it. */ +async function stopWithUnprovenClose(failures: number): Promise { + const connection = claude.connections[0]! + await openTurn(connection) + closeUnprovenFor(connection, failures) + await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } }) + frame(connection, INTERRUPTED_RESULT) + await laneDrained() + expect(connection.closeCount).toBe(1) + expect(owedWindDown()).toBeDefined() + return connection +} + +function owedWindDown() { + return host['sessions'].get(SESSION)?.owesProviderChildWindDown +} + +async function waitRows() { + await host.flushStreamedEvents(SESSION) + return (await host.journalSnapshot(SESSION)).items.flatMap((item) => + item.body.kind === 'status' && item.body.failure?.kind === 'previousExitUnverifiable' + ? [item.body] + : [] + ) +} + +async function resumedWith(connection: FakeConnection, text: string): Promise { + return eventually(() => { + const started = claude.connections.at(-1)! + expect(started).not.toBe(connection) + expect(wrote(started, text)).toBe(true) + return started + }) +} + +it('retries the close a Stop could not prove before the next message, then sends it to a resumed child', async () => { + const connection = await stopWithUnprovenClose(1) + + const next = await send('Carry on.') + const resumed = await resumedWith(connection, 'Carry on.') + + expect(connection.closeCount).toBe(2) + expect(wrote(connection, 'Carry on.')).toBe(false) + expect(resumed.closed).toBe(false) + expect(owedWindDown()).toBeUndefined() + expect((await submission(next))?.dispatchState).toBe('pending') + expect(await waitRows()).toEqual([]) +}) + +it('holds the message with its reason while the exit stays unverifiable, and sends it once a later retry proves it', async () => { + const connection = await stopWithUnprovenClose(3) + + // Never refused: the send is accepted and returns at once. + const next = await send('Carry on.') + await eventually(async () => expect(await waitRows()).toHaveLength(1)) + await laneDrained() + expect(await waitRows()).toEqual([ + { + kind: 'status', + tone: 'warning', + text: "Orca couldn't confirm Claude's previous process ended. Messages wait to be sent until Orca confirms it has ended.", + failure: { kind: 'previousExitUnverifiable' } + } + ]) + // Still queued, drawn below the chat; never written to the child the Stop could not end. + expect(await submission(next)).toMatchObject({ dispatchState: 'pending' }) + expect((await submission(next))?.handedOverAt).toBeUndefined() + expect(claude.connections).toHaveLength(1) + expect(wrote(connection, 'Carry on.')).toBe(false) + expect(owedWindDown()).toBeDefined() + // The Stop's attempt and the send's one retry: the row's own commit retries nothing. + expect(connection.closeCount).toBe(2) + + // A second message retries once more, and waits under the same row. + const second = await send('And this.') + await eventually(() => expect(connection.closeCount).toBe(3)) + await laneDrained() + expect(await submission(second)).toMatchObject({ dispatchState: 'pending' }) + expect(await waitRows()).toHaveLength(1) + + // The sweep's next tick retries the stop; the exit is proven and both messages go out. + await host['lifetime'].idleSweep.tick() + const resumed = await resumedWith(connection, 'Carry on.') + await eventually(() => expect(wrote(resumed, 'And this.')).toBe(true)) + expect(connection.closeCount).toBe(4) + expect(owedWindDown()).toBeUndefined() + expect(await waitRows()).toHaveLength(1) +}) + +it('lets a held message be withdrawn with Stop, and the owed stop still ends on the next retry', async () => { + const connection = await stopWithUnprovenClose(2) + const next = await send('Carry on.') + await eventually(async () => expect(await waitRows()).toHaveLength(1)) + await laneDrained() + + await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } }) + expect(await submission(next)).toMatchObject({ dispatchState: 'rejected' }) + expect(owedWindDown()).toBeDefined() + + await host['lifetime'].idleSweep.tick() + expect(connection.closeCount).toBe(3) + expect(owedWindDown()).toBeUndefined() + expect(store.getRecord(SESSION)?.lease.claimStatus).toBe('released') + // Nothing was waiting, so no child starts. + expect(claude.connections).toHaveLength(1) + expect(wrote(connection, 'Carry on.')).toBe(false) +}) + +function setModel(model: string) { + const fields = { key: 'model', value: model } + return host.setOption(CALLER, { envelope: envelope('agentSession.setOption', fields), ...fields }) +} + +it('refuses an option change with the exit unverifiable after retrying the stop, never writing to the old child', async () => { + const connection = await stopWithUnprovenClose(2) + + await expect(setModel('claude-opus-5')).resolves.toMatchObject({ + ok: false, + refusal: { + code: 'agent_session_ownership_unknown', + details: { reason: 'previousExitUnverifiable', ownerVerdict: 'unverifiable' } + } + }) + expect(connection.closeCount).toBe(2) + expect(connection.calls.some((call) => call.subtype === 'set_model')).toBe(false) + expect(owedWindDown()).toBeDefined() + + // Trying again retries the stop; proven, the pick is the chat's at rest, for the next start. + await expect(setModel('claude-opus-5')).resolves.toMatchObject({ ok: true }) + expect(connection.closeCount).toBe(3) + expect(owedWindDown()).toBeUndefined() + expect(connection.calls.some((call) => call.subtype === 'set_model')).toBe(false) + expect(store.getRecord(SESSION)?.options).toMatchObject({ model: 'claude-opus-5' }) + expect(claude.connections).toHaveLength(1) +}) + +function stopBackgroundTasks() { + const fields = { scope: 'background-tasks' as const } + return host.cancel(CALLER, { envelope: envelope('agentSession.cancel', fields), ...fields }) +} + +/** Holds the user's next message on a Stop whose exit stays unproven through the send's retry. */ +async function heldAfterStop(): Promise { + const connection = await stopWithUnprovenClose(2) + await send('Carry on.') + await eventually(async () => expect(await waitRows()).toHaveLength(1)) + await laneDrained() + expect(connection.closeCount).toBe(2) + return connection +} + +it('sends a held message once an option change proves the exit: the stop that lands hands it over', async () => { + const connection = await heldAfterStop() + + await expect(setModel('claude-opus-5')).resolves.toMatchObject({ ok: true }) + expect(connection.closeCount).toBe(3) + const resumed = await resumedWith(connection, 'Carry on.') + expect(resumed.closed).toBe(false) + expect(owedWindDown()).toBeUndefined() + expect(store.getRecord(SESSION)?.options).toMatchObject({ model: 'claude-opus-5' }) +}) + +it('sends a held message once a background-task stop proves the exit, though that stop finds no agent', async () => { + const connection = await heldAfterStop() + + // Proven gone, the old agent has no tasks left to stop, and the message still goes out. + await expect(stopBackgroundTasks()).resolves.toMatchObject({ + ok: false, + refusal: { details: { reason: 'noLiveOwner' } } + }) + expect(connection.closeCount).toBe(3) + await resumedWith(connection, 'Carry on.') + expect(owedWindDown()).toBeUndefined() +}) + +it('sends a message held after a tab close whose exit was unproven, rather than rejecting it as closed', async () => { + const connection = claude.connections[0]! + closeUnprovenFor(connection, 2) + await expect(host.close(SESSION, 'user-close')).rejects.toThrow() + expect(owedWindDown()).toMatchObject({ cause: 'user-close' }) + + // The chat stays open on the host; the user sends again, and the message waits on that close. + const next = await send('Carry on.') + await eventually(async () => expect(await waitRows()).toHaveLength(1)) + await laneDrained() + expect(connection.closeCount).toBe(2) + + // The close was asked for before this message, so the retry that lands ends nothing it waited on. + await host['lifetime'].idleSweep.tick() + await laneDrained() + expect(await submission(next)).not.toMatchObject({ dispatchState: 'rejected' }) + await resumedWith(connection, 'Carry on.') + expect(owedWindDown()).toBeUndefined() +}) + +it('never stops a live child for a wind-down another, earlier child still owes', async () => { + const connection = claude.connections[0]! + const session = host['sessions'].get(SESSION)! + session.owesProviderChildWindDown = { + generation: 'an-earlier-child', + fence: session.child!.fence, + cause: 'user-stop', + requestedAt: session.journal.cursor() + } + + await send('Carry on.') + await eventually(() => expect(wrote(connection, 'Carry on.')).toBe(true)) + await laneDrained() + await host['lifetime'].idleSweep.tick() + expect(connection.closeCount).toBe(0) + expect(claude.connections).toHaveLength(1) +}) + +it('retries once per new message: commits after a second waiting message retry nothing', async () => { + const connection = await stopWithUnprovenClose(6) + await send('Carry on.') + await eventually(async () => expect(await waitRows()).toHaveLength(1)) + await laneDrained() + await send('And this.') + await eventually(() => expect(connection.closeCount).toBe(3)) + await laneDrained() + + // Any journal commit wakes the delivery loop while a message is queued. + const session = host['sessions'].get(SESSION)! + for (const n of [1, 2, 3]) { + await session.journal.appendItem( + { provider: 'orca', clientMessageId: `unrelated-${n}` }, + { kind: 'status', text: `Unrelated ${n}.` }, + { fence: store.getRecord(SESSION)!.lease.runtimeFence, turnScope: AGENT_JOURNAL_THREAD_SCOPE } + ) + await commitSettled() + } + expect(connection.closeCount).toBe(3) + expect(claude.connections).toHaveLength(1) + + // The sweep still retries, one pass a tick, until the exit is proven; then both go out. + for (const _tick of [1, 2, 3, 4]) { + await host['lifetime'].idleSweep.tick() + } + expect(connection.closeCount).toBe(7) + const resumed = await resumedWith(connection, 'Carry on.') + await eventually(() => expect(wrote(resumed, 'And this.')).toBe(true)) +}) + +it('queues a follow-up as a draft by default while a message waits, and Steer retries the stop once for it', async () => { + const connection = await stopWithUnprovenClose(3) + await send('Carry on.') + await eventually(async () => expect(await waitRows()).toHaveLength(1)) + await laneDrained() + + // Queueing follow-ups is the default: the waiting message reads as working, so this is a draft. + const body = hostTestMessage('And this.') + const delivery = 'queue-if-active' as const + const queued = await host.send(CALLER, { + envelope: envelope('agentSession.send', { body, delivery }), + body, + delivery + }) + expect(queued).toMatchObject({ ok: true, value: { queued: { state: 'waiting' } } }) + await commitSettled() + expect(connection.closeCount).toBe(2) + + // Steer makes it a waiting message, which retries once and waits under the same note. + const messageId = queued.ok && 'queued' in queued.value ? queued.value.queued.messageId : '' + await expect( + host.queuedMessageSend(CALLER, { + envelope: envelope('agentSession.queuedMessageSend', { messageId }), + messageId + }) + ).resolves.toMatchObject({ ok: true }) + await eventually(() => expect(connection.closeCount).toBe(3)) + await commitSettled() + expect(connection.closeCount).toBe(3) + expect(claude.connections).toHaveLength(1) + expect(await waitRows()).toHaveLength(1) + + await host['lifetime'].idleSweep.tick() + const resumed = await resumedWith(connection, 'Carry on.') + await eventually(() => expect(wrote(resumed, 'And this.')).toBe(true)) +}) + +it('rejects as closed a message a second tab close closed, though that close could not reject it itself', async () => { + const connection = claude.connections[0]! + closeUnprovenFor(connection, 3) + await expect(host.close(SESSION, 'user-close')).rejects.toThrow() + const next = await send('Carry on.') + await eventually(async () => expect(await waitRows()).toHaveLength(1)) + await laneDrained() + expect(connection.closeCount).toBe(2) + + // The second close's own rejection of what is queued fails; its stop is a new ask all the same. + const session = host['sessions'].get(SESSION)! + vi.spyOn(session.journal, 'rejectQueuedSubmissions').mockRejectedValueOnce(new Error('disk full')) + await expect(host.close(SESSION, 'user-close')).rejects.toThrow() + expect(connection.closeCount).toBe(3) + expect(await submission(next)).toMatchObject({ dispatchState: 'pending' }) + + // The retry that proves the exit closes what that second close closed. + await host['lifetime'].idleSweep.tick() + await commitSettled() + expect(connection.closeCount).toBe(4) + expect(await submission(next)).toMatchObject({ dispatchState: 'rejected' }) + expect(claude.connections).toHaveLength(1) +}) + +it('notes why a message waits when another operation failed its retry before the message was delivered', async () => { + const connection = await stopWithUnprovenClose(2) + + // The send is accepted, then the option change's failed retry runs before the send's delivery. + const sent = send('Carry on.') + const option = setModel('claude-opus-5') + await expect(option).resolves.toMatchObject({ + ok: false, + refusal: { details: { reason: 'previousExitUnverifiable' } } + }) + await sent + await commitSettled() + expect(connection.closeCount).toBe(2) + expect(await waitRows()).toHaveLength(1) + expect(claude.connections).toHaveLength(1) +}) + +it('keeps an option change at rest when the Stop proved the exit and only its bookkeeping keeps failing', async () => { + const connection = claude.connections[0]! + await openTurn(connection) + // The close proves the exit; draining what the old agent wrote fails for the Stop. + const barrierLost = { ok: false as const, error: new Error('drain barrier lost') } + const drained = vi + .spyOn(host['runtimeState'].eventSinkFor(SESSION), 'drained') + .mockResolvedValueOnce(barrierLost) + .mockResolvedValueOnce(barrierLost) + await expect(stop()).resolves.toMatchObject({ ok: true, value: { cancelled: true } }) + frame(connection, INTERRUPTED_RESULT) + await laneDrained() + expect(connection.closeCount).toBe(1) + expect(host['sessions'].get(SESSION)?.child).toBeNull() + expect(owedWindDown()).toBeDefined() + + // The option change's retry of that bookkeeping fails again: reported, and the pick still lands. + drained.mockResolvedValueOnce(barrierLost) + await expect(setModel('claude-opus-5')).resolves.toMatchObject({ ok: true }) + expect(owedWindDown()).toBeDefined() + expect(log.scopes().filter((scope) => scope === 'owed-stop-retry')).toHaveLength(1) + expect(store.getRecord(SESSION)?.options).toMatchObject({ model: 'claude-opus-5' }) + expect(connection.closeCount).toBe(1) + expect(claude.connections).toHaveLength(1) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts index d6a6819679e..e02d605fa73 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-codex-stop-row.test.ts @@ -36,6 +36,8 @@ let host: StructuredAgentSessionHost let turns: ReturnType let codex: ReturnType let notify: (method: string, params: unknown) => void +/** Read at each start, so a test can say what the next start resumes. */ +let launch: { resumeThreadId?: string | null } let disposeSession: MockInstance> beforeEach(async () => { @@ -57,7 +59,8 @@ beforeEach(async () => { } const store = await openTestAgentSessionRecordStore(root) // The runtime's wiring: an echo accepts its send, and an exit reaches the host. - const adapter = adapterFor(codex, {}, [], { + launch = {} + const adapter = adapterFor(codex, launch, [], { onDispatchSettledLate: (settlement) => void host.settleLateDispatch(settlement), onEvent: (event) => { if (event.type === 'ended' && 'cause' in event && event.cause === 'unexpected-exit') { @@ -442,3 +445,52 @@ describe('a Codex Stop whose interrupt failed', () => { expect(childEndedByStop()).toBe(false) }) }) + +describe('a message after a Codex Stop whose exit was unproven', () => { + it('retries that stop first, then goes to a fresh Codex, never to the old one', async () => { + await runningTurn() + codex.routes['turn/interrupt'] = () => { + throw interruptFailure('internal error') + } + // As the real connection: once a close begins it refuses every request, proven or not. + const old = codex.connections.at(-1)! + const close = old.close + const request = old.request + let unproven = 1 + old.close = async () => { + old.closed = true + if (unproven === 0) { + return close() + } + unproven -= 1 + return false + } + old.request = (method, params) => + old.closed + ? Promise.reject(new Error('codex app-server is closing')) + : request(method, params) + + await stop() + await host.flushStreamedEvents(SESSION) + expect(host['sessions'].get(SESSION)?.owesProviderChildWindDown).toBeDefined() + // The next start resumes the chat's thread, as the runtime's launch resolves it from the record. + launch.resumeThreadId = THREAD + + const sent = await send('carry on') + expect(sent).toMatchObject({ ok: true }) + await vi.waitFor(() => expect(codex.connections).toHaveLength(2)) + await vi.waitFor(() => + expect( + codex.connections[1]!.calls.some( + (call) => call.method === 'turn/start' && JSON.stringify(call.params).includes('carry on') + ) + ).toBe(true) + ) + expect( + old.calls.some( + (call) => call.method === 'turn/start' && JSON.stringify(call.params).includes('carry on') + ) + ).toBe(false) + expect(host['sessions'].get(SESSION)?.owesProviderChildWindDown).toBeUndefined() + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts index 499c5f3ff48..5b88983d872 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-close.test.ts @@ -301,6 +301,7 @@ describe('the wind-down retry with a message queued (P2-31)', () => { owesProviderChildWindDown: { generation: 'generation-1', fence: 1 } } const stopAgent = vi.fn(async () => undefined) + const finishOwedWindDown = vi.fn(async () => true) const sweep = new StructuredAgentSessionIdleSweep({ // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: a session fixture carrying only the journal and child facts the sweep reads. sessions: Object.assign(new Map([[SESSION, session as never]]), { @@ -316,6 +317,7 @@ describe('the wind-down retry with a message queued (P2-31)', () => { providerHoldsDispatch: () => false, stopAgent, stopStartingAgent: stopAgent, + finishOwedWindDown, closeConversation: vi.fn(async () => false), // A failed step fails the test. logger: { @@ -329,5 +331,6 @@ describe('the wind-down retry with a message queued (P2-31)', () => { }) await sweep.tick() expect(stopAgent).not.toHaveBeenCalled() + expect(finishOwedWindDown).not.toHaveBeenCalled() }) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-lifetime.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-lifetime.ts index 491021dc1c1..c79d53dc23f 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-lifetime.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-conversation-lifetime.ts @@ -15,6 +15,7 @@ import type { StructuredAgentSessionConversations } from './structured-agent-ses import { abandonQueuedStructuredAgentSessionMessages, closeStructuredAgentSessionConversationUnderSerialize, + finishOwedStructuredAgentSessionWindDownUnderSerialize, stopStructuredAgentSessionAgentUnderSerialize, type StructuredAgentSessionCloseCause, type StructuredAgentSessionLifetimeContext @@ -56,6 +57,8 @@ export function createStructuredAgentSessionConversationLifetime(host: { }) const stopAgent = (sessionId: string, cause: StructuredAgentSessionStopCause) => stopStructuredAgentSessionAgentUnderSerialize(host.context(), sessionId, { cause }) + const finishOwedWindDown = (sessionId: string) => + finishOwedStructuredAgentSessionWindDownUnderSerialize(host.context(), sessionId) const closeConversation = (sessionId: string): Promise => closeStructuredAgentSessionConversationUnderSerialize( @@ -84,6 +87,7 @@ export function createStructuredAgentSessionConversationLifetime(host: { providerHoldsDispatch: (sessionId) => deps().adapter.holdsDispatch?.(sessionId) === true, // The host puts an idle agent to rest: a turn it cuts short is news, not the user's Stop. stopAgent: (sessionId) => stopAgent(sessionId, 'evict'), + finishOwedWindDown, // A host stop: the delivery loop waiting on this child writes the one error row and rejects // what is queued with it, both worded from the hostStopped fact. stopStartingAgent: (sessionId) => diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-delivery-loop.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-delivery-loop.ts index 3727040e3c0..0414836f4a5 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-delivery-loop.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-delivery-loop.ts @@ -26,7 +26,10 @@ import { structuredAgentSessionStartFailure, type StructuredAgentSessionStartFailureCause } from './structured-agent-session-failure-text' -import type { StructuredAgentSessionResumeOutcome } from './structured-agent-session-agent-start' +import { + isStructuredAgentSessionPreviousExitUnverifiable, + type StructuredAgentSessionResumeOutcome +} from './structured-agent-session-agent-start' import type { StructuredAgentSessionChildEndCause, StructuredAgentSessionEndedChild, @@ -39,6 +42,10 @@ import { } from './structured-agent-session-start-failure-row' import { failedProviderChildStart } from './structured-agent-session-provider-child' import { handOverSubmission } from './structured-agent-session-turns' +import { + recordStructuredAgentSessionWindDownWait, + structuredAgentSessionWindDownWaitHolds +} from './structured-agent-session-wind-down-wait-row' import { structuredAgentSessionCommandRunning } from './structured-agent-session-command-turn' import type { StructuredAgentSessionLogger } from './structured-agent-session-logger' @@ -182,11 +189,24 @@ export class StructuredAgentSessionDeliveryLoop { if (!oldest || (session.child && structuredAgentSessionCommandRunning(session.journal))) { return this.stop(sessionId) } + // Already waiting on a stop that could not prove its child gone: a new message retries it, and + // any other retry that lands wakes this loop itself, so the waiting row's own commit does not. + // Another operation's retry may have failed first, so the row is made sure of here too. + if (structuredAgentSessionWindDownWaitHolds(session)) { + await recordStructuredAgentSessionWindDownWait(session, sessionId, this.deps) + return this.stop(sessionId) + } const failedStart = startThatFailedWhileQueued(session, oldest) if (failedStart) { return this.fail(sessionId, failedStart) } const ready = await this.deps.ensureProviderChild(sessionId, oldest.clientMessageId) + if (!ready.ok && isStructuredAgentSessionPreviousExitUnverifiable(ready.refusal)) { + // The start retried that stop first and still could not prove the exit: the message waits, + // saying why, rather than being refused. The row is written once per unproven child. + await recordStructuredAgentSessionWindDownWait(session, sessionId, this.deps) + return this.stop(sessionId) + } if (!ready.ok) { return ready } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts index a2f38c2e6bb..fa90e015169 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts @@ -22,10 +22,13 @@ import type { StructuredAgentSessionHostRuntimeState } from './structured-agent- import type { StructuredAgentSessionHostDeps, StructuredAgentSessionHostSession, + StructuredAgentSessionOwedWindDown, StructuredAgentSessionProviderChildIdentity } from './structured-agent-session-host-types' import { endProviderChild, + pendingProviderChildWindDown, + sameProviderChild, structuredAgentSessionConversationFence } from './structured-agent-session-provider-child' import { releaseStoredStructuredAgentSessionOwner } from './structured-agent-session-lease-release' @@ -40,6 +43,8 @@ export type StructuredAgentSessionLifetimeContext = { now: () => number /** Re-projects the session's status after its agent stopped and the chat stays. */ publishStatus?: (sessionId: string) => void + /** Hands the delivery loop what is queued; for a caller inside the session's serialize. */ + wakeDelivery?: (sessionId: string) => void /** Quit-only snapshot taken immediately before the provider child is stopped. */ restartWitness?: { beforeStop: (sessionId: string) => void @@ -91,6 +96,27 @@ function owedProviderChildWindDown( : session.owesProviderChildWindDown } +/** The stop this pass owes. A retry continues the one already asked for, keeping where it was asked; + * any other stop is a new ask, even one with the same cause: a second close closes what came since. */ +function owedStop( + session: StructuredAgentSessionHostSession, + cause: StructuredAgentSessionStopCause, + retry: boolean +): StructuredAgentSessionOwedWindDown | undefined { + const owed = owedProviderChildWindDown(session) + if (!owed) { + return undefined + } + const asked = session.owesProviderChildWindDown + const continues = retry && asked !== undefined && sameProviderChild(asked, owed) + return { + generation: owed.generation, + fence: owed.fence, + cause, + requestedAt: continues ? asked.requestedAt : session.journal.cursor() + } +} + /** * The agent goes to rest; the conversation stays. Runs the eviction steps under a deadline. A step * that fails — or runs out of time — aborts the rest and leaves the wind-down owed, so the next @@ -100,8 +126,9 @@ function owedProviderChildWindDown( export async function stopStructuredAgentSessionAgentUnderSerialize( context: StructuredAgentSessionLifetimeContext, sessionId: string, - // Required: an omitted cause must not default to the user's cancellation. - ending: { cause: StructuredAgentSessionStopCause; reason?: string } + // Required: an omitted cause must not default to the user's cancellation. `retry` is set only by + // the retry of a stop already owed. + ending: { cause: StructuredAgentSessionStopCause; reason?: string; retry?: true } ): Promise { const session = context.sessions.get(sessionId) if (!session) { @@ -114,8 +141,8 @@ export async function stopStructuredAgentSessionAgentUnderSerialize( // The obligation OUTLIVES the child. `child` is ended the instant the adapter proves the exit, // so a step that aborts after that point would otherwise leave the retry reading "no child // here" and skipping the settlement and the lease release it still owes. - const owed = owedProviderChildWindDown(session) - session.owesProviderChildWindDown = owed ? { ...owed, cause } : undefined + const owed = owedStop(session, cause, ending.retry === true) + session.owesProviderChildWindDown = owed const stopping = session.child const eviction: StructuredAgentSessionEvictionContext = { sessionId, @@ -139,6 +166,8 @@ export async function stopStructuredAgentSessionAgentUnderSerialize( cause, reason: ending.reason ?? null, duringStartup: stopping.phase === 'starting', + // A later retry that proves the exit still ends the child at the Stop it finishes. + ...(owed ? { endedAt: owed.requestedAt } : {}), ...verdict }) } @@ -183,12 +212,53 @@ export async function stopStructuredAgentSessionAgentUnderSerialize( // Whatever ended the child, the row belongs to the conversation: it shows not-running, and // only the conversation's close forgets it. context.publishStatus?.(sessionId) + // The stop's own end hands over what waited on it, whichever caller's retry landed. + context.wakeDelivery?.(sessionId) } } - await evictStructuredAgentSession( - eviction, - withStructuredAgentSessionEvictionDeadline(STRUCTURED_AGENT_SESSION_EVICTION_STEPS) - ) + try { + await evictStructuredAgentSession( + eviction, + withStructuredAgentSessionEvictionDeadline(STRUCTURED_AGENT_SESSION_EVICTION_STEPS) + ) + } catch (error) { + if (session.owesProviderChildWindDown === owed && owed) { + session.owesProviderChildWindDown = { ...owed, failedAt: session.journal.cursor() } + } + throw error + } +} + +/** + * Retries the wind-down an earlier stop left owed, with that stop's own cause: a child it could + * not prove gone takes no input, so nothing may write to it or start beside it until this lands. + * One pass of the stop, bounded by its own step deadline (10 s): a shorter bound would cut a + * supervised Claude's exit proof (up to about 7 s) short. Resolves whether nothing is owed now; a + * failure is reported, never thrown, and leaves the exit unverifiable, never exited. Landing, the + * stop itself hands over what waited on it. + */ +export async function finishOwedStructuredAgentSessionWindDownUnderSerialize( + context: StructuredAgentSessionLifetimeContext, + sessionId: string +): Promise { + const session = context.sessions.get(sessionId) + const owed = session && pendingProviderChildWindDown(session) + if (!owed) { + return true + } + try { + await stopStructuredAgentSessionAgentUnderSerialize(context, sessionId, { + cause: owed.cause, + retry: true + }) + } catch (error) { + context.deps.logger.warn('retrying an unfinished agent stop failed', { + scope: 'owed-stop-retry', + sessionId, + error + }) + } + return context.sessions.get(sessionId)?.owesProviderChildWindDown === undefined } /** A close's cause: the user closing this chat, or the host evicting it (quit, idle, teardown). */ 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 5067bcb3fb2..3fbb96ddaf8 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 @@ -25,7 +25,7 @@ import { type StructuredAgentSessionMutationContext } from './structured-agent-session-mutation-context' import { - openForWrite, + openForProviderWrite, openWithAgent, sendPreparation, structuredAgentSessionSendBlock @@ -101,7 +101,7 @@ export function cancelStructuredAgentSessionTurn( caller, params.envelope, cancelPlan({ ...params, childWork: () => context.readChildWork(params.envelope.sessionId) }), - openForWrite(context, params.envelope) + openForProviderWrite(context, params.envelope) ) } const plan = cancelPlan(params) @@ -130,7 +130,7 @@ export function respondToStructuredAgentSessionPrompt( caller, params.envelope, promptPlan(params), - openForWrite(context, params.envelope) + openForProviderWrite(context, params.envelope) ) } @@ -159,7 +159,7 @@ export async function setStructuredAgentSessionOption( ? recordStructuredAgentSessionOptionIntent(context.deps.store, ctx, params) : plan.run(ctx) }, - openForWrite(context, params.envelope) + openForProviderWrite(context, params.envelope) ) } 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 dcf708643e9..6552f8a2601 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 @@ -38,9 +38,13 @@ export type StructuredAgentSessionProviderChildIdentity = { readonly fence: number } -/** A wind-down still owed, with the cause of the stop that owes it: a retry finishes that stop. */ +/** A wind-down still owed, with the stop that owes it: a retry finishes that stop. */ export type StructuredAgentSessionOwedWindDown = StructuredAgentSessionProviderChildIdentity & { readonly cause: StructuredAgentSessionStopCause + /** Where the journal stood when the stop was asked for; the child's end is ordered there. */ + readonly requestedAt: AgentJournalCursor + /** Where it stood once the newest pass failed: a message accepted by then waited through a retry. */ + readonly failedAt?: AgentJournalCursor } /** The provider process behind a conversation. Written only in @@ -74,7 +78,8 @@ export type StructuredAgentSessionEndedChild = StructuredAgentSessionProviderChi duringStartup: boolean startedFor?: string /** Where the conversation's journal stood when the child ended, to order the end against a - * message's acceptance. */ + * message's acceptance. A stop's end stands where it was asked for: a message accepted while + * retries proved the exit waited on it, and came after it. */ endedAt: AgentJournalCursor } 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 95cc822bba1..e11e1dbd04e 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 @@ -22,10 +22,7 @@ import { structuredAgentSessionOwnerStatus } from './structured-agent-session-ow import { StructuredAgentSessionHostRuntimeState } from './structured-agent-session-host-runtime-state' import { attachStructuredAgentSession } from './structured-agent-session-attach-orchestration' import type { StructuredAgentSessionLifetimeContext } from './structured-agent-session-host-lifetime' -import { - ensureStructuredAgentSessionAgent, - ensureStructuredAgentSessionAgentForOperation -} from './structured-agent-session-agent-start' +import * as agentStart from './structured-agent-session-agent-start' import { createStructuredAgentSessionConversationLifetime, type StructuredAgentSessionConversationLifetime @@ -128,7 +125,7 @@ export class StructuredAgentSessionHost { // Quit drains a delivery start before it evicts, so the child it produces is stopped. trackStart: (start) => this.tasks.trackAttach(start), ensureProviderChild: (sessionId, startedFor) => - ensureStructuredAgentSessionAgent(this.attachContext(), sessionId, startedFor), + agentStart.ensureStructuredAgentSessionAgent(this.attachContext(), sessionId, startedFor), reset: (sessionId, journal, reset) => this.subscribers.reset( sessionId, @@ -195,7 +192,8 @@ export class StructuredAgentSessionHost { runtimeState: this.runtimeState, sessions: this.sessions, now: () => this.now(), - publishStatus: this.clientDelivery.publishStatus + publishStatus: this.clientDelivery.publishStatus, + wakeDelivery: (sessionId: string) => this.conversationDelivery.loop.wake(sessionId) } satisfies StructuredAgentSessionLifetimeContext } @@ -283,7 +281,12 @@ export class StructuredAgentSessionHost { serialize: (sessionId, task) => this.serialize(sessionId, task), openConversation: this.conversationDelivery.open, ensureAgent: (sessionId) => - ensureStructuredAgentSessionAgentForOperation(this.attachContext(), sessionId), + agentStart.ensureStructuredAgentSessionAgentForOperation(this.attachContext(), sessionId), + finishOwedStop: (sessionId) => + agentStart.finishOwedStructuredAgentSessionStopForProviderWrite( + this.attachContext(), + sessionId + ), wakeDelivery: (sessionId) => this.conversationDelivery.loop.wake(sessionId), stopAgent: (sessionId) => this.lifetime.stopAgent(sessionId, 'user-stop'), wakeQueuedDrain: (sessionId) => this.queued.drain.schedule(sessionId), diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.test.ts index 8ff4cdbeb42..fba7abcfdeb 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.test.ts @@ -319,6 +319,7 @@ describe('the idle sweep with no child running (P2-22 ii)', () => { providerHoldsDispatch: () => false, stopAgent, stopStartingAgent: stopAgent, + finishOwedWindDown: vi.fn(async () => true), closeConversation, // A failed step fails the test. logger: { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.ts index e9b0e1b45c6..96a5f59cf84 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-idle-sweep.ts @@ -15,6 +15,7 @@ import type { AgentJournalRenderItem } from '../../../shared/agent-session-journ import type { AgentChildWorkView } from '../../../shared/agent-status-child-work-view' import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types' import type { StructuredAgentSessionLogger } from './structured-agent-session-logger' +import { pendingProviderChildWindDown } from './structured-agent-session-provider-child' export const STRUCTURED_AGENT_SESSION_IDLE_SWEEP_INTERVAL_MS = 5 * 60_000 export const STRUCTURED_AGENT_SESSION_IDLE_MS = 30 * 60_000 @@ -36,6 +37,8 @@ export type StructuredAgentSessionIdleSweepDeps = { providerHoldsDispatch: (sessionId: string) => boolean /** Each of these runs inside the session's serialize and never takes it again. */ stopAgent: (sessionId: string) => Promise + /** Retries a stop that did not finish; landing, it hands over what waited on it. */ + finishOwedWindDown: (sessionId: string) => Promise stopStartingAgent: (sessionId: string) => Promise closeConversation: (sessionId: string) => Promise logger: StructuredAgentSessionLogger @@ -108,14 +111,10 @@ export class StructuredAgentSessionIdleSweep { if (!session || this.deps.isDisposed()) { return } - // A stop that failed after the child was proven gone: finish it now, before the idle test, so - // the rows its settlement wrote cannot push the retry out. A message accepted since goes first. - if ( - session.owesProviderChildWindDown !== undefined && - !session.child && - !this.queuedOrDelivering(sessionId, session) - ) { - await this.deps.stopAgent(sessionId) + // A stop that did not finish: retry it now, before the idle test, so the rows its settlement + // wrote cannot push the retry out. A running delivery step retries it itself. + if (pendingProviderChildWindDown(session) && !this.deps.deliveryActive(sessionId)) { + await this.deps.finishOwedWindDown(sessionId) return } // Owed work is activity, read every tick, so the agent gets a full window once it ends: a child diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-context.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-context.ts index 01059d67ceb..f02a982024b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-context.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-context.ts @@ -32,6 +32,9 @@ export type StructuredAgentSessionMutationContext = { openConversation: (sessionId: string) => Promise /** Gives the session a provider child; inside the caller's serialize. */ ensureAgent: (sessionId: string) => Promise + /** Finishes a stop an earlier attempt left owed, for an operation that starts no child; inside + * the caller's serialize. */ + finishOwedStop: (sessionId: string) => Promise /** A message was accepted: the session's delivery loop hands it over. */ wakeDelivery: (sessionId: string) => void /** Stops the session's provider child, keeping its conversation; inside the caller's serialize. */ diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-provider-child.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-provider-child.ts index da5e491e1d1..5785cfc53e1 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-provider-child.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-provider-child.ts @@ -5,11 +5,13 @@ // an exit, a failed re-attach, a Stop and an eviction. Each is matched on the child's generation // and fence, so an ending that arrives late for an older child cannot end a newer one. +import type { AgentJournalCursor } from '../../../shared/agent-session-journal-types' import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import type { StructuredAgentSessionEndedChild, StructuredAgentSessionHostSession, + StructuredAgentSessionOwedWindDown, StructuredAgentSessionProviderChild, StructuredAgentSessionProviderChildIdentity } from './structured-agent-session-host-types' @@ -46,9 +48,12 @@ export function markProviderChildStarted( return child !== null } +/** `endedAt` is a stop's ask; an exit ends where the journal stands. */ export function endProviderChild( session: ChildBearer, - ended: Omit + ended: Omit & { + endedAt?: AgentJournalCursor + } ): boolean { const child = matchingChild(session, ended) if (!child) { @@ -58,7 +63,7 @@ export function endProviderChild( session.lastEndedChild = { ...ended, ...(child.startedFor === undefined ? {} : { startedFor: child.startedFor }), - endedAt: session.journal.cursor() + endedAt: ended.endedAt ?? session.journal.cursor() } return true } @@ -72,12 +77,26 @@ export function failedProviderChildStart( return !session.child && ended?.duringStartup && ended.cause !== 'user-stop' ? ended : null } +/** The owed wind-down an operation reaching the provider finishes first. One owed for another child + * never outranks the child in front of it, which that child's own stop finishes. */ +export function pendingProviderChildWindDown( + session: Pick +): StructuredAgentSessionOwedWindDown | undefined { + const { child, owesProviderChildWindDown: owed } = session + return owed && (!child || sameProviderChild(child, owed)) ? owed : undefined +} + +export function sameProviderChild( + a: StructuredAgentSessionProviderChildIdentity, + b: StructuredAgentSessionProviderChildIdentity +): boolean { + return a.generation === b.generation && a.fence === b.fence +} + function matchingChild( session: ChildBearer, identity: StructuredAgentSessionProviderChildIdentity ): StructuredAgentSessionProviderChild | null { const { child } = session - return child && child.generation === identity.generation && child.fence === identity.fence - ? child - : null + return child && sameProviderChild(child, identity) ? child : null } 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 4b2011a02ee..7c6511dd2f2 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 @@ -83,6 +83,25 @@ export function openForWrite( return () => openConversationForWrite(context.openConversation, envelope, context.deps.logger) } +/** For an operation the running child performs, which starts none: the conversation, then any + * stop an earlier attempt left owed, so it never reaches a child that takes no input. */ +export function openForProviderWrite( + context: Pick< + StructuredAgentSessionMutationContext, + 'openConversation' | 'finishOwedStop' | 'deps' + >, + envelope: AgentSessionMutationEnvelope +): () => Promise { + return async () => { + const opened = await openConversationForWrite( + context.openConversation, + envelope, + context.deps.logger + ) + return opened.ok ? context.finishOwedStop(envelope.sessionId) : opened + } +} + /** For an operation only the provider can perform: the conversation, then its agent. */ export function openWithAgent( context: Pick, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts index 30809bb5cd2..48500afa12f 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-turns-cancel.ts @@ -80,8 +80,8 @@ export type StructuredAgentSessionStopWindDown = { waitsForProvider: boolean; st /** * A session-ending Stop's second step, queued behind its first in the same tick so nothing sent * meanwhile reaches the child it ends. The Stop has answered: a failure here is reported. The next - * Stop retries the wind-down it leaves owed, and so does the idle sweep: at its next tick once the - * child is proven gone, else only after the chat idles with no child work left. + * operation that reaches the agent retries the wind-down it leaves owed, and so does the idle + * sweep's next tick. */ export async function endStoppedStructuredAgentSession( ctx: Pick, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-wind-down-wait-row.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-wind-down-wait-row.ts new file mode 100644 index 00000000000..5f1d03a7ead --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-wind-down-wait-row.ts @@ -0,0 +1,83 @@ +// The row a message waiting on an unfinished stop leaves in the chat: why it has not gone out. + +import { agentSessionFailureFact } from '../../../shared/agent-session-failure' +import { + agentSessionFailureWords, + type AgentSessionFailureWordsContext +} from '../../../shared/agent-session-failure-words' +import { agentJournalItemKey } from '../../../shared/agent-session-journal-item-key' +import { isQueuedAgentJournalSubmission } from '../../../shared/agent-session-queued-submission' +import { + AGENT_JOURNAL_THREAD_SCOPE, + type AgentJournalItemIdentity +} from '../../../shared/agent-session-journal-types' +import type { + StructuredAgentSessionHostSession, + StructuredAgentSessionProviderChildIdentity +} from './structured-agent-session-host-types' +import { pendingProviderChildWindDown } from './structured-agent-session-provider-child' + +/** Keyed by the child the stop could not prove gone, so every send that waits on it shares a row. */ +export function structuredAgentSessionWindDownWaitIdentity( + owed: StructuredAgentSessionProviderChildIdentity +): AgentJournalItemIdentity { + return { + provider: 'orca', + clientMessageId: `wind-down-wait:${owed.fence}:${owed.generation ?? 'unknown'}` + } +} + +type WaitingSession = Pick< + StructuredAgentSessionHostSession, + 'journal' | 'child' | 'owesProviderChildWindDown' +> + +function windDownWaitRow( + session: WaitingSession, + owed: StructuredAgentSessionProviderChildIdentity +) { + const itemId = agentJournalItemKey(structuredAgentSessionWindDownWaitIdentity(owed)) + return session.journal.snapshot().items.find((item) => item.itemId === itemId) +} + +/** Every queued message already waited through a retry of the owed stop: it was accepted before the + * newest pass failed. Only a new message retries again, and a retry that lands wakes the loop + * itself, so no other commit, each of which wakes the delivery loop, retries. */ +export function structuredAgentSessionWindDownWaitHolds(session: WaitingSession): boolean { + const failedAt = pendingProviderChildWindDown(session)?.failedAt + if (!failedAt || failedAt.epoch !== session.journal.cursor().epoch) { + return false + } + return session.journal + .submissions() + .every( + (submission) => + !isQueuedAgentJournalSubmission(submission) || + (submission.acceptedSequence ?? 0) <= failedAt.sequence + ) +} + +/** Written once per child, so every message that waits on it shares the row. */ +export async function recordStructuredAgentSessionWindDownWait( + session: WaitingSession, + sessionId: string, + input: { + conversationFence: (sessionId: string) => number + failureTextContext: (sessionId: string) => AgentSessionFailureWordsContext + } +): Promise { + const owed = pendingProviderChildWindDown(session) + if (!owed || windDownWaitRow(session, owed)) { + return + } + const identity = structuredAgentSessionWindDownWaitIdentity(owed) + const words = agentSessionFailureWords(agentSessionFailureFact('previousExitUnverifiable'), { + ...input.failureTextContext(sessionId), + surface: 'row' + }) + await session.journal.appendItem( + identity, + { kind: 'status', tone: 'warning', ...words }, + { fence: input.conversationFence(sessionId), turnScope: AGENT_JOURNAL_THREAD_SCOPE } + ) +} diff --git a/src/renderer/src/components/native-chat/agent-session-failure-words-text.ts b/src/renderer/src/components/native-chat/agent-session-failure-words-text.ts index 26833f720ee..34dab4ab94f 100644 --- a/src/renderer/src/components/native-chat/agent-session-failure-words-text.ts +++ b/src/renderer/src/components/native-chat/agent-session-failure-words-text.ts @@ -227,6 +227,12 @@ const PIECES: Record + translate( + 'components.native-chat.failureWords.previousExitUnverifiable', + COPY.previousExitUnverifiable, + values ) } diff --git a/src/renderer/src/i18n/en-runtime-required.json b/src/renderer/src/i18n/en-runtime-required.json index 5dae9e98fb1..72118c9ab68 100644 --- a/src/renderer/src/i18n/en-runtime-required.json +++ b/src/renderer/src/i18n/en-runtime-required.json @@ -2759,6 +2759,7 @@ "notDelivered": "This message was not delivered.", "notDeliveredSendAgain": "This message was not delivered. Send it again to continue.", "notSignedIn": "{{agent}} is not signed in for the selected account.", + "previousExitUnverifiable": "Orca couldn't confirm {{agent}}'s previous process ended. Messages wait to be sent until Orca confirms it has ended.", "providerExitedRejection": "{{agent}} stopped before this message was sent.", "providerExitedRow": "{{agent}} stopped while this response was in progress. You can continue in this conversation.", "providerRateLimited": "{{agent}} is rate-limited and retrying.", diff --git a/src/renderer/src/i18n/locales/en.json b/src/renderer/src/i18n/locales/en.json index 3e930124282..f38bb0588ee 100644 --- a/src/renderer/src/i18n/locales/en.json +++ b/src/renderer/src/i18n/locales/en.json @@ -17545,7 +17545,8 @@ "hostStopped": "{{agent}} never finished starting, so Orca stopped it.", "providerRateLimited": "{{agent}} is rate-limited and retrying.", "providerRetrying": "{{agent}} hit a temporary problem and is retrying.", - "providerRetryingQuoted": "{{agent}} is retrying: {{detail}}." + "providerRetryingQuoted": "{{agent}} is retrying: {{detail}}.", + "previousExitUnverifiable": "Orca couldn't confirm {{agent}}'s previous process ended. Messages wait to be sent until Orca confirms it has ended." }, "writeNotice": { "notDoneReadHistory": "This chat's history couldn't be loaded.", diff --git a/src/renderer/src/i18n/locales/es.json b/src/renderer/src/i18n/locales/es.json index 03b3f528bad..ae52313f4ec 100644 --- a/src/renderer/src/i18n/locales/es.json +++ b/src/renderer/src/i18n/locales/es.json @@ -14986,7 +14986,8 @@ "hostStopped": "{{agent}} no terminó de iniciarse, así que Orca lo detuvo.", "providerRateLimited": "{{agent}} alcanzó un límite de solicitudes y está reintentando.", "providerRetrying": "{{agent}} tuvo un problema temporal y está reintentando.", - "providerRetryingQuoted": "{{agent}} está reintentando: {{detail}}." + "providerRetryingQuoted": "{{agent}} está reintentando: {{detail}}.", + "previousExitUnverifiable": "Orca no pudo confirmar que el proceso anterior de {{agent}} terminó. Los mensajes esperan para enviarse hasta que Orca confirme que terminó." }, "writeNotice": { "notDoneReadHistory": "No se pudo cargar el historial de este chat.", diff --git a/src/renderer/src/i18n/locales/fr.json b/src/renderer/src/i18n/locales/fr.json index 8b01f3b4c25..f3b7f238079 100644 --- a/src/renderer/src/i18n/locales/fr.json +++ b/src/renderer/src/i18n/locales/fr.json @@ -17452,7 +17452,8 @@ "hostStopped": "{{agent}} n'a jamais fini de démarrer, Orca l'a donc arrêté.", "providerRateLimited": "{{agent}} est limité en débit et réessaie.", "providerRetrying": "{{agent}} a rencontré un problème temporaire et réessaie.", - "providerRetryingQuoted": "{{agent}} réessaie : {{detail}}." + "providerRetryingQuoted": "{{agent}} réessaie : {{detail}}.", + "previousExitUnverifiable": "Orca n'a pas pu confirmer que le processus précédent de {{agent}} s'est arrêté. Les messages attendent d'être envoyés jusqu'à ce qu'Orca confirme son arrêt." }, "writeNotice": { "notDoneReadHistory": "L'historique de ce chat n'a pas pu être chargé.", diff --git a/src/renderer/src/i18n/locales/ja.json b/src/renderer/src/i18n/locales/ja.json index 8ad615691e8..354b147116b 100644 --- a/src/renderer/src/i18n/locales/ja.json +++ b/src/renderer/src/i18n/locales/ja.json @@ -17388,7 +17388,8 @@ "hostStopped": "{{agent}} の起動が完了しなかったため、Orca が停止しました。", "providerRateLimited": "{{agent}} はレート制限を受けているため、再試行しています。", "providerRetrying": "{{agent}} で一時的な問題が発生したため、再試行しています。", - "providerRetryingQuoted": "{{agent}} は再試行しています: {{detail}}。" + "providerRetryingQuoted": "{{agent}} は再試行しています: {{detail}}。", + "previousExitUnverifiable": "Orca は {{agent}} の前のプロセスが終了したことを確認できませんでした。Orca が終了を確認するまで、メッセージは送信を待機します。" }, "writeNotice": { "notDoneReadHistory": "このチャットの履歴を読み込めませんでした。", diff --git a/src/renderer/src/i18n/locales/ko.json b/src/renderer/src/i18n/locales/ko.json index 0cc7eaff0d4..58c01a08b59 100644 --- a/src/renderer/src/i18n/locales/ko.json +++ b/src/renderer/src/i18n/locales/ko.json @@ -17388,7 +17388,8 @@ "hostStopped": "{{agent}}이(가) 시작을 마치지 못해 Orca가 중지했습니다.", "providerRateLimited": "{{agent}}이(가) 속도 제한에 걸려 다시 시도하고 있습니다.", "providerRetrying": "{{agent}}에 일시적인 문제가 발생하여 다시 시도하고 있습니다.", - "providerRetryingQuoted": "{{agent}}이(가) 다시 시도하고 있습니다: {{detail}}." + "providerRetryingQuoted": "{{agent}}이(가) 다시 시도하고 있습니다: {{detail}}.", + "previousExitUnverifiable": "Orca가 {{agent}}의 이전 프로세스가 종료되었는지 확인하지 못했습니다. Orca가 종료를 확인할 때까지 메시지는 전송을 기다립니다." }, "writeNotice": { "notDoneReadHistory": "이 채팅의 기록을 불러오지 못했습니다.", diff --git a/src/renderer/src/i18n/locales/zh.json b/src/renderer/src/i18n/locales/zh.json index f54fbab8102..afdfc6a1cb7 100644 --- a/src/renderer/src/i18n/locales/zh.json +++ b/src/renderer/src/i18n/locales/zh.json @@ -17353,7 +17353,8 @@ "hostStopped": "{{agent}} 始终未完成启动,因此 Orca 已将其停止。", "providerRateLimited": "{{agent}} 已被限流,正在重试。", "providerRetrying": "{{agent}} 遇到临时问题,正在重试。", - "providerRetryingQuoted": "{{agent}} 正在重试:{{detail}}。" + "providerRetryingQuoted": "{{agent}} 正在重试:{{detail}}。", + "previousExitUnverifiable": "Orca 无法确认 {{agent}} 的上一个进程已结束。在 Orca 确认其已结束之前,消息会等待发送。" }, "writeNotice": { "notDoneReadHistory": "无法加载此聊天的历史记录。", diff --git a/src/shared/agent-session-failure-copy.ts b/src/shared/agent-session-failure-copy.ts index 20d6a69e112..e6eeab28222 100644 --- a/src/shared/agent-session-failure-copy.ts +++ b/src/shared/agent-session-failure-copy.ts @@ -80,7 +80,9 @@ export const AGENT_SESSION_FAILURE_COPY = { hostStopped: '{{agent}} never finished starting, so Orca stopped it.', providerRateLimited: '{{agent}} is rate-limited and retrying.', providerRetrying: '{{agent}} hit a temporary problem and is retrying.', - providerRetryingQuoted: '{{agent}} is retrying: {{detail}}.' + providerRetryingQuoted: '{{agent}} is retrying: {{detail}}.', + previousExitUnverifiable: + "Orca couldn't confirm {{agent}}'s previous process ended. Messages wait to be sent until Orca confirms it has ended." } as const export type AgentSessionFailureCopyId = keyof typeof AGENT_SESSION_FAILURE_COPY diff --git a/src/shared/agent-session-failure-words.ts b/src/shared/agent-session-failure-words.ts index 47c32fa4599..f64ba5cb335 100644 --- a/src/shared/agent-session-failure-words.ts +++ b/src/shared/agent-session-failure-words.ts @@ -267,7 +267,9 @@ const FAILURE_SENTENCES = { agent(say, context) ), retry?.cause - ) + ), + previousExitUnverifiable: (context, _fact, _surface, say) => + say('previousExitUnverifiable', agent(say, context)) } satisfies Record /** The sentence a person reads for this fact on this surface; never a marker. */ diff --git a/src/shared/agent-session-failure.ts b/src/shared/agent-session-failure.ts index 14bc8eb3a1e..b08c5250a15 100644 --- a/src/shared/agent-session-failure.ts +++ b/src/shared/agent-session-failure.ts @@ -47,7 +47,9 @@ export const AGENT_SESSION_FAILURE_KINDS = [ /** Orca stopped an agent whose start never finished. */ 'hostStopped', /** The provider is retrying a request its API refused; not a failure yet. */ - 'providerRetrying' + 'providerRetrying', + /** A message waits on a child a Stop could not prove gone: its exit is unverifiable. */ + 'previousExitUnverifiable' ] as const export type AgentSessionFailureKind = (typeof AGENT_SESSION_FAILURE_KINDS)[number] @@ -62,7 +64,8 @@ const STATUS_ROW_ONLY_FAILURE_KINDS = [ 'cancelUnconfirmed', 'stopRefused', 'answerUnconfirmed', - 'providerRetrying' + 'providerRetrying', + 'previousExitUnverifiable' ] as const satisfies readonly AgentSessionFailureKind[] /** Why a message was not sent. A new failure kind is one of these until listed above. */ diff --git a/src/shared/agent-session-refusal-details.ts b/src/shared/agent-session-refusal-details.ts index 9e697c10b53..70c6121d03f 100644 --- a/src/shared/agent-session-refusal-details.ts +++ b/src/shared/agent-session-refusal-details.ts @@ -66,7 +66,10 @@ export const AGENT_SESSION_REFUSAL_REASONS = { 'spawnIdentityMismatch', 'notResumable', 'noProviderChild', - 'conversationHeldElsewhere' + 'conversationHeldElsewhere', + /** A Stop could not prove its child gone, and a retry could not either: that child takes no + * input and none starts beside it. Sent with `ownerVerdict: 'unverifiable'`. */ + 'previousExitUnverifiable' ], agent_session_conflict: [ 'chatStarting', diff --git a/src/shared/agent-session-refusal-notice.ts b/src/shared/agent-session-refusal-notice.ts index b2021797bed..0f9c47ded03 100644 --- a/src/shared/agent-session-refusal-notice.ts +++ b/src/shared/agent-session-refusal-notice.ts @@ -156,7 +156,9 @@ const REASON_WORDS = { spawnIdentityMismatch: codeWords('hostFinding'), notResumable: codeWords('retry'), noProviderChild: codeWords('retry'), - conversationHeldElsewhere: codeWords('retry') + conversationHeldElsewhere: codeWords('retry'), + // Trying again retries the stop that could not prove the exit. + previousExitUnverifiable: causeWords('ownerUnproven', 'retry') }, agent_session_conflict: { chatStarting: AGENT_STARTING,