diff --git a/src/main/claude/claude-structured-dispatch.test.ts b/src/main/claude/claude-structured-dispatch.test.ts index d66a64f82eb..9756a4a94ae 100644 --- a/src/main/claude/claude-structured-dispatch.test.ts +++ b/src/main/claude/claude-structured-dispatch.test.ts @@ -65,6 +65,49 @@ describe('Claude structured dispatch image limits', () => { expect(session.activeTurnSequence).toBe(session.dispatchSequence) }) + it('reports the settled submission when a timed-out replay arrives late', async () => { + const session = sessionFor() + const accepted: { clientMessageId: string; uuid: string }[] = [] + session.onLateDispatchAccepted = (input) => accepted.push(input) + const dispatched = dispatchClaudeTurn( + session, + { clientMessageId: 'client-1', body: userMessage([{ type: 'text', text: 'one' }]) }, + 500 + ) + await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1)) + const sentUuid = (session.dispatchWaiters[0] as { sentUuid?: string }).sentUuid + // The send gave up waiting, so the host has already recorded this submission `unknown`. + await expect(dispatched).resolves.toMatchObject({ state: 'unknown' }) + + resolveClaudeReplayWaiter(session, userReplayFrame(sentUuid!, 'one')) + + expect(accepted).toEqual([{ clientMessageId: 'client-1', uuid: sentUuid }]) + }) + + it('reports a late replay even after a newer dispatch has started', async () => { + const session = sessionFor() + const accepted: { clientMessageId: string; uuid: string }[] = [] + session.onLateDispatchAccepted = (input) => accepted.push(input) + const first = dispatchClaudeTurn( + session, + { clientMessageId: 'client-1', body: userMessage([{ type: 'text', text: 'one' }]) }, + 500 + ) + await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1)) + const firstUuid = (session.dispatchWaiters[0] as { sentUuid?: string }).sentUuid + await expect(first).resolves.toMatchObject({ state: 'unknown' }) + dispatchClaudeTurn( + session, + { clientMessageId: 'client-2', body: userMessage([{ type: 'text', text: 'two' }]) }, + 100 + ) + await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1)) + + // Turn identity must stay with dispatch B, but message A still provably landed. + expect(resolveClaudeReplayWaiter(session, userReplayFrame(firstUuid!, 'one'))).toBe(false) + expect(accepted).toEqual([{ clientMessageId: 'client-1', uuid: firstUuid }]) + }) + it('never lets a late replay for dispatch A resolve dispatch B', async () => { const session = sessionFor() const first = dispatchClaudeTurn( diff --git a/src/main/claude/claude-structured-dispatch.ts b/src/main/claude/claude-structured-dispatch.ts index 96271e41d71..bef0b9f3f39 100644 --- a/src/main/claude/claude-structured-dispatch.ts +++ b/src/main/claude/claude-structured-dispatch.ts @@ -147,6 +147,12 @@ function recoverLateIdentity( session.activeTurnId = uuid session.activeTurnSequence = waiter.dispatchSequence } + // The waiter is retired, so its send already returned `unknown`. A user replay is Claude + // echoing the message, which proves delivery — regardless of whether a newer dispatch has + // since started, so this is deliberately not gated on `dispatchSequence`. + if (isUserReplay) { + session.onLateDispatchAccepted?.({ clientMessageId: waiter.clientMessageId, uuid }) + } return isUserReplay && waiter.dispatchSequence === session.dispatchSequence } @@ -155,12 +161,14 @@ function waitForReplay( timeoutMs: number, acceptsResult: boolean, sentUuid: string, - replayContentKey: string + replayContentKey: string, + clientMessageId: string ): { waiter: ClaudeDispatchWaiter; promise: Promise } { let waiter!: ClaudeDispatchWaiter const promise = new Promise((resolve) => { waiter = { acceptsResult, + clientMessageId, sentUuid, dispatchSequence: session.dispatchSequence, replayContentKey, @@ -220,7 +228,8 @@ export async function dispatchClaudeTurn( timeoutMs, acceptsResult, sentUuid, - claudeDispatchContentKey(content) + claudeDispatchContentKey(content), + input.clientMessageId ) const replayed = replay.promise try { diff --git a/src/main/claude/claude-structured-session-acquisition.ts b/src/main/claude/claude-structured-session-acquisition.ts index e8d09bd78d9..77d0f55a5f4 100644 --- a/src/main/claude/claude-structured-session-acquisition.ts +++ b/src/main/claude/claude-structured-session-acquisition.ts @@ -265,6 +265,12 @@ export async function acquireClaudeSession({ }) const acquired: AgentSessionAcquisition = publication.acquisition liveSession = publication.session + if (deps.onLateDispatchAccepted) { + const notify = deps.onLateDispatchAccepted + const bound = publication.session + bound.onLateDispatchAccepted = ({ clientMessageId, uuid }) => + notify({ sessionId, clientMessageId, uuid, providerSessionId: bound.providerSessionId }) + } await restoreClaudeStructuredSessionOptions(liveSession, deps.requestTimeoutMs) acquisitions.assertCurrent(sessionId, attempt) acquisitions.deleteIfCurrent(sessionId, attempt) diff --git a/src/main/claude/claude-structured-session-state.ts b/src/main/claude/claude-structured-session-state.ts index 1c0b1862913..651588edd86 100644 --- a/src/main/claude/claude-structured-session-state.ts +++ b/src/main/claude/claude-structured-session-state.ts @@ -60,6 +60,15 @@ export type ClaudeStructuredSessionAdapterDeps = { sessionId: string, state: AgentSessionBackgroundTaskState | null ) => void + /** A replay that lands after its dispatch timed out proves the message reached Claude, so the + * submission the host already recorded `unknown` can still be settled `accepted`. */ + onLateDispatchAccepted?: (input: { + sessionId: string + clientMessageId: string + uuid: string + /** Claude's own session id — the journal identity is provider-scoped, not Orca-scoped. */ + providerSessionId: string + }) => void openConnection?: typeof openClaudeStreamJsonConnection readProcessStartTime?: (pid: number) => Promise mintLinkId?: () => string @@ -87,6 +96,8 @@ export type ClaudeDispatchWaiter = { resolve: (uuid: string | null) => void timer: ReturnType acceptsResult: boolean + /** The submission this dispatch belongs to, so a late replay can settle its journal row. */ + clientMessageId: string /** Client uuid echoed by Claude so a replay is tied to its own dispatch. */ sentUuid: string /** Sequence used to fence a late identity from a newer dispatch. */ @@ -100,6 +111,8 @@ export type ClaudeDispatchWaiter = { } export type ClaudeSession = { + /** Set by the adapter; see `onLateDispatchAccepted` on the adapter deps. */ + onLateDispatchAccepted?: (input: { clientMessageId: string; uuid: string }) => void connection: ClaudeStreamJsonConnection providerSessionId: string /** Durable transcript files live under this account's `projects` directory. */ diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts index 5354670fa0b..c6fb0880d87 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts @@ -338,6 +338,54 @@ describe('send', () => { expect(result).toMatchObject({ ok: true, value: { submission: { dispatchState: 'unknown' } } }) }) + it('settles an unknown submission accepted when its replay lands late', async () => { + await attach() + dispatch.mockResolvedValueOnce({ state: 'unknown', reason: 'no replay in time' }) + const body = hostTestMessage('add a retry') + const result = await host.send(CALLER, { + envelope: envelope('agentSession.send', { body }), + body + }) + if (!result.ok) { + throw new Error(`expected a send, got ${result.refusal.code}`) + } + expect(result.value.submission.dispatchState).toBe('unknown') + + await host.settleLateDispatch({ + sessionId: SESSION, + clientMessageId: result.value.clientMessageId, + providerIdentity: { provider: 'codex', threadId: THREAD, turnId: 'turn-late', ordinal: 0 } + }) + + const page = host.history({ sessionId: SESSION, direction: 'tail' }) + const settled = (page.ok ? page.page.submissions : []).find( + (entry) => entry.clientMessageId === result.value.clientMessageId + ) + expect(settled?.dispatchState).toBe('accepted') + // The message was already delivered; a late settlement must never re-dispatch it. + expect(dispatch).toHaveBeenCalledTimes(1) + }) + + it('leaves an already accepted submission alone when a late replay arrives', async () => { + await attach() + const body = hostTestMessage('add a retry') + const result = await host.send(CALLER, { + envelope: envelope('agentSession.send', { body }), + body + }) + if (!result.ok) { + throw new Error(`expected a send, got ${result.refusal.code}`) + } + const before = host.history({ sessionId: SESSION, direction: 'tail' }) + await host.settleLateDispatch({ + sessionId: SESSION, + clientMessageId: result.value.clientMessageId, + providerIdentity: { provider: 'codex', threadId: THREAD, turnId: 'turn-late', ordinal: 0 } + }) + const after = host.history({ sessionId: SESSION, direction: 'tail' }) + expect(after.ok && after.page.fence).toBe(before.ok && before.page.fence) + }) + it('replays a retried send from the journal without dispatching twice', async () => { await attach() const body = hostTestMessage('add a retry') 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 13b7d7b441e..b7c0ad1fe5a 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 @@ -1,6 +1,7 @@ // Structured agent-session host: where the lease, journal, and provider adapter meet. // Mutations share one durable admission path and serialize per session. +import type { AgentJournalItemIdentity } from '../../../shared/agent-session-journal-types' import type { AgentSessionExecutionLocation } from '../../../shared/agent-session-record' import type * as SessionWire from '../../../shared/agent-session-wire' import type { AgentSessionAttachParams } from './structured-agent-session-attach' @@ -337,6 +338,37 @@ export class StructuredAgentSessionHost { ) => this.backgroundTasks.publish(sessionId, state) unsubscribe = (sessionId: string, id: string): void => this.subscribers.close(sessionId, id) + /** A provider replay that arrived after its dispatch timed out. The send already recorded + * `unknown`, which is terminal for the outbox — the client shows "delivery is unconfirmed" + * forever and its Retry redelivers a message the agent already has. The replay proves + * delivery, so settle the submission the host was unsure about. */ + settleLateDispatch = (input: { + sessionId: string + clientMessageId: string + providerIdentity: AgentJournalItemIdentity + }): Promise => + this.serialize(input.sessionId, async () => { + const session = this.sessions.get(input.sessionId) + if (!session) { + return + } + // Only an `unknown` row is ours to move; anything else is already the truth. + const submission = session.journal + .submissions() + .find((entry) => entry.clientMessageId === input.clientMessageId) + if (submission?.dispatchState !== 'unknown') { + return + } + await session.journal.resolveDispatch({ + clientMessageId: input.clientMessageId, + state: 'accepted', + providerIdentity: input.providerIdentity, + fence: session.fence, + recovered: true + }) + this.subscribers.publish(input.sessionId, session.journal) + }) + /** Every session's projected status for session lists; unlike `subscribe`, retains nothing. */ subscribeStatus = ( subscriber: Parameters[0] diff --git a/src/main/runtime/structured-agent-session-runtime.ts b/src/main/runtime/structured-agent-session-runtime.ts index d670c16f47d..ecfcf1b59cd 100644 --- a/src/main/runtime/structured-agent-session-runtime.ts +++ b/src/main/runtime/structured-agent-session-runtime.ts @@ -260,6 +260,17 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise host?.publishBackgroundTaskState(sessionId, state), + onLateDispatchAccepted: ({ sessionId, clientMessageId, uuid, providerSessionId }) => { + void host + ?.settleLateDispatch({ + sessionId, + clientMessageId, + providerIdentity: { provider: 'claude', sessionId: providerSessionId, uuid } + }) + .catch((error) => + deps.onError?.({ scope: `structured-agent-session-late-dispatch:${sessionId}`, error }) + ) + }, ...(deps.openClaudeConnection ? { openClaudeConnection: deps.openClaudeConnection } : {}), ...(deps.readProcessStartTime ? { readProcessStartTime: deps.readProcessStartTime } : {}) }) diff --git a/src/main/runtime/structured-claude-runtime-adapter.ts b/src/main/runtime/structured-claude-runtime-adapter.ts index 95151f1dae9..a595b4ff5f9 100644 --- a/src/main/runtime/structured-claude-runtime-adapter.ts +++ b/src/main/runtime/structured-claude-runtime-adapter.ts @@ -34,6 +34,7 @@ export type StructuredClaudeRuntimeAdapterDeps = { sessionId: string, state: AgentSessionBackgroundTaskState | null ) => void + onLateDispatchAccepted?: ClaudeStructuredSessionAdapterDeps['onLateDispatchAccepted'] } export function createStructuredClaudeRuntimeAdapter( @@ -100,6 +101,7 @@ export function createStructuredClaudeRuntimeAdapter( ...(deps.onBackgroundTasksChanged ? { onBackgroundTasksChanged: deps.onBackgroundTasksChanged } : {}), + ...(deps.onLateDispatchAccepted ? { onLateDispatchAccepted: deps.onLateDispatchAccepted } : {}), ...(deps.openClaudeConnection ? { openConnection: deps.openClaudeConnection } : {}), ...(deps.readProcessStartTime ? { readProcessStartTime: deps.readProcessStartTime } : {}) })