From ba32236c69e4133f1db4bd3fa7a148600fd5cd49 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Sun, 6 Sep 2026 16:09:06 -0700 Subject: [PATCH] fix(native-chat): persist late dispatch receipts before session close --- ...structured-agent-session-host-mutations.ts | 29 ++----- ...ured-agent-session-late-settlement.test.ts | 80 ++++++++++++++++++- 2 files changed, 84 insertions(+), 25 deletions(-) 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 6f20a4d16b2..9a5f3e0c475 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 @@ -122,19 +122,7 @@ export function readStructuredAgentSessionOptions( }) } -/** - * Settle a send whose ack window expired but which the provider later proved it - * had received. - * - * Without this the submission stays `unknown` for the life of the session: the - * client renders an unconfirmed bubble whose Retry redispatches, so the user is - * invited to deliver the same message to the agent a second time. Every send made - * while a turn is already running takes this path, because the provider does not - * echo the new message until the running turn ends. - * - * Serialized with the session's other journal writes, and a no-op once the row is - * `accepted` or `rejected` — the reducer treats both as terminal. - */ +/** Settle provider-proven delivery independently of an in-flight client mutation. */ export async function settleStructuredAgentSessionLateDispatch( context: StructuredAgentSessionMutationContext, input: { @@ -147,13 +135,12 @@ export async function settleStructuredAgentSessionLateDispatch( if (!session) { return } - await context.serialize(input.sessionId, async () => { - await session.journal.resolveDispatch({ - clientMessageId: input.clientMessageId, - state: 'accepted', - providerIdentity: input.providerIdentity, - fence: session.fence - }) - context.publish(input.sessionId, session.journal) + // The journal queue drains before close; the host queue would defer this past teardown. + await session.journal.resolveDispatch({ + clientMessageId: input.clientMessageId, + state: 'accepted', + providerIdentity: input.providerIdentity, + fence: session.fence }) + context.publish(input.sessionId, session.journal) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-late-settlement.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-late-settlement.test.ts index e345a29f39e..a75d2ea6512 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-late-settlement.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-late-settlement.test.ts @@ -3,9 +3,11 @@ 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 { AgentSessionMutationEnvelope } from '../../../shared/agent-session-wire' +import type { + AgentSessionMutationEnvelope, + AgentSessionSubscribeEvent +} from '../../../shared/agent-session-wire' import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' -import { createTrackedJournalOpener } from '../agent-session-journal/journal-store-test-open' import type { AgentSessionDispatchOutcome, StructuredAgentSessionAdapter @@ -21,13 +23,13 @@ import { resetHostTestOperationIds } from './structured-agent-session-host-test-data' -const journals = createTrackedJournalOpener() const CALLER = { callerKey: 'client-1' } let root: string let store: AgentSessionRecordStore let host: StructuredAgentSessionHost let dispatch: Mock +let closeSession: Mock> function accepted(): AgentSessionDispatchOutcome { return { @@ -65,6 +67,7 @@ beforeEach(async () => { root = await mkdtemp(join(tmpdir(), 'orca-wire-late-settle-')) resetHostTestOperationIds() dispatch = vi.fn(async () => accepted()) + closeSession = vi.fn(async () => true) store = await AgentSessionRecordStore.open({ directory: join(root, 'store'), hostId: 'local' }) host = new StructuredAgentSessionHost({ store, @@ -86,6 +89,7 @@ beforeEach(async () => { })), releaseAcquisition: vi.fn(async () => true), dispatch, + closeSession, cancelTurn: vi.fn(async () => ({ cancelled: true })), answerPrompt: vi.fn(async () => undefined), setOption: vi.fn(async () => undefined) @@ -99,12 +103,80 @@ beforeEach(async () => { }) afterEach(async () => { - await journals.closeAll() await host.flushAllStreamedEvents() + await host.close(SESSION) await rm(root, { recursive: true, force: true }) }) describe('settling a send the provider proves it received after the ack window', () => { + it('publishes acceptance during a pending send and never reopens it for retry', async () => { + let finishDispatch!: (outcome: AgentSessionDispatchOutcome) => void + dispatch.mockImplementationOnce( + () => + new Promise((resolve) => { + finishDispatch = resolve + }) + ) + const events: AgentSessionSubscribeEvent[] = [] + const unsubscribe = host.subscribe({ + id: 'late-receipt', + sessionId: SESSION, + emit: (event) => events.push(event) + }) + const params = sendParams('echo before send completes') + const pending = host.send(CALLER, params) + await vi.waitFor(() => expect(dispatch).toHaveBeenCalledTimes(1)) + try { + await host.settleLateDispatch({ + sessionId: SESSION, + clientMessageId: params.envelope.clientOperationId, + providerIdentity: { provider: 'claude', sessionId: THREAD, uuid: 'early-echo' } + }) + expect(events.at(-1)).toMatchObject({ + type: 'batch', + batch: { + submissions: [ + { clientMessageId: params.envelope.clientOperationId, dispatchState: 'accepted' } + ] + } + }) + } finally { + finishDispatch({ state: 'unknown', reason: 'ack timeout' }) + unsubscribe() + } + await expect(pending).resolves.toMatchObject({ + ok: true, + value: { submission: { dispatchState: 'accepted' } } + }) + await expect(host.send(CALLER, { ...params, retryUnknown: true })).resolves.toMatchObject({ + ok: true, + value: { submission: { dispatchState: 'accepted' } } + }) + expect(dispatch).toHaveBeenCalledTimes(1) + }) + + it('persists an echo received while the provider is closing', async () => { + dispatch.mockResolvedValueOnce({ state: 'unknown', reason: 'ack timeout' }) + const params = sendParams('received just before shutdown') + await host.send(CALLER, params) + let settlement: Promise | undefined + closeSession.mockImplementationOnce(async () => { + settlement = host.settleLateDispatch({ + sessionId: SESSION, + clientMessageId: params.envelope.clientOperationId, + providerIdentity: { provider: 'claude', sessionId: THREAD, uuid: 'closing-echo' } + }) + void settlement.catch(() => undefined) + return true + }) + + await host.close(SESSION) + await expect(settlement).resolves.toBeUndefined() + await host.revealSession(SESSION) + expect(submissions()).toMatchObject([{ dispatchState: 'accepted' }]) + expect(dispatch).toHaveBeenCalledTimes(1) + }) + it('moves a durable unknown to accepted so nothing offers to send it again', async () => { dispatch.mockRejectedValueOnce(new Error('socket closed')) const params = sendParams('sent while a turn was running')