From c8baf4eaaa336f63896f4689cf2bb8c1b6111224 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Wed, 26 Aug 2026 13:41:24 -0700 Subject: [PATCH] fix(native-chat): close stale turns and retry rejected sends --- ...dex-structured-journal-translation.test.ts | 73 +++++++- .../codex-structured-journal-translation.ts | 43 +++-- ...e-structured-agent-session-outbox.test.tsx | 167 ++++++++++++++++++ .../use-structured-agent-session-outbox.ts | 33 +++- 4 files changed, 299 insertions(+), 17 deletions(-) diff --git a/src/main/codex/codex-structured-journal-translation.test.ts b/src/main/codex/codex-structured-journal-translation.test.ts index 464c53f41a2..84497c50c79 100644 --- a/src/main/codex/codex-structured-journal-translation.test.ts +++ b/src/main/codex/codex-structured-journal-translation.test.ts @@ -4,7 +4,10 @@ import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types' import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key' -import { projectStructuredItemsToNativeChat } from '../../shared/structured-agent-session-projection' +import { + projectStructuredAgentSessionStatus, + projectStructuredItemsToNativeChat +} from '../../shared/structured-agent-session-projection' import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' import { CodexTurnOrdinals } from './codex-structured-item-translation' import { @@ -138,6 +141,74 @@ describe('codex journal translation', () => { expect(tap.tombstones).toEqual(['legacy:codex:session-1:turn-lifecycle%3Aturn-1']) }) + it('closes every active turn when the provider session ends after a later turn starts', () => { + const tap = recorder() + const translator = createCodexJournalTranslator({ + sink: tap.sink, + primaryThreadId: () => THREAD_ID + }) + + translator.handle(notification('turn/started', { turn: { id: 'turn-stale' } })) + translator.handle(notification('turn/started', { turn: { id: 'turn-later' } })) + translator.handle({ type: 'ended', sessionId: SESSION_ID, reason: 'app-server exited' }) + + expect(tap.rows.filter((row) => row.body.kind === 'status')).toHaveLength(2) + expect(tap.rows.map((row) => row.body)).toEqual([ + expect.objectContaining({ turnLifecycle: { turnId: 'turn-stale', state: 'running' } }), + expect.objectContaining({ turnLifecycle: { turnId: 'turn-later', state: 'running' } }) + ]) + expect(tap.tombstones).toEqual([ + 'legacy:codex:session-1:turn-lifecycle%3Aturn-stale', + 'legacy:codex:session-1:turn-lifecycle%3Aturn-later' + ]) + // The tombstones remove both running rows from the reduced journal; no + // lifecycle identity remains live after a session end. + expect( + projectStructuredAgentSessionStatus( + tap.rows + .filter((row) => !tap.tombstones.includes(row.key)) + .map((row, sequence) => ({ + itemId: row.key, + revision: 1, + sequence: sequence + 1, + observedAt: sequence + 1, + body: row.body + })) + ) + ).toBe('idle') + }) + + it('matches out-of-order completions to each turn identity', () => { + const tap = recorder() + const translator = createCodexJournalTranslator({ + sink: tap.sink, + primaryThreadId: () => THREAD_ID + }) + + translator.handle(notification('turn/started', { turn: { id: 'turn-stale' } })) + translator.handle(notification('turn/started', { turn: { id: 'turn-later' } })) + translator.handle(notification('turn/completed', { turn: { id: 'turn-stale' } })) + translator.handle(notification('turn/completed', { turn: { id: 'turn-later' } })) + + expect(tap.tombstones).toEqual([ + 'legacy:codex:session-1:turn-lifecycle%3Aturn-stale', + 'legacy:codex:session-1:turn-lifecycle%3Aturn-later' + ]) + expect( + projectStructuredAgentSessionStatus( + tap.rows + .filter((row) => !tap.tombstones.includes(row.key)) + .map((row, sequence) => ({ + itemId: row.key, + revision: 1, + sequence: sequence + 1, + observedAt: sequence + 1, + body: row.body + })) + ) + ).toBe('idle') + }) + it('journals a user turn and the assistant answer under durable codex keys', () => { const { translator, tap } = translatorWith() diff --git a/src/main/codex/codex-structured-journal-translation.ts b/src/main/codex/codex-structured-journal-translation.ts index f31a2718758..10bfcb9ef98 100644 --- a/src/main/codex/codex-structured-journal-translation.ts +++ b/src/main/codex/codex-structured-journal-translation.ts @@ -65,17 +65,33 @@ export function createCodexJournalTranslator( const identities = new Map() /** What each announced item is, so an approval can name what it approves. */ const details = new Map() - const currentTurnIds = new Map() + /** Turns announced by the provider and not yet closed. */ + const currentTurnIds = new Map>() const genericRowsByTurn = new Map() const suppressedRowsByTurn = new Map() let fallbackSequence = 0 + const currentTurnIdFor = (threadId: string): string | null => + [...(currentTurnIds.get(threadId) ?? [])].at(-1) ?? null + + const rememberTurn = (threadId: string, turnId: string): void => { + currentTurnIds.set(threadId, new Set([...(currentTurnIds.get(threadId) ?? []), turnId])) + } + + const forgetTurn = (threadId: string, turnId: string): void => { + const active = currentTurnIds.get(threadId) + active?.delete(turnId) + if (!active?.size) { + currentTurnIds.delete(threadId) + } + } + const appendUnhandled = (kind: string, payload: unknown, threadId = 'session'): void => { const translated = unhandledProviderFrameJournalItem('codex', kind, payload) if (!translated) { return } - const turnId = readCodexTurnId(payload) ?? currentTurnIds.get(threadId) ?? 'outside-turn' + const turnId = readCodexTurnId(payload) ?? currentTurnIdFor(threadId) ?? 'outside-turn' const bucket = `${encodeURIComponent(threadId)}:${encodeURIComponent(turnId)}` const rowCount = genericRowsByTurn.get(bucket) ?? 0 // The cap bounds noise, never evidence: an error frame is always journaled, @@ -152,7 +168,7 @@ export function createCodexJournalTranslator( coalesceMs: deps.coalesceMs, schedule: deps.schedule, identityFor: (threadId, params, item) => { - const turnId = readCodexTurnId(params) ?? currentTurnIds.get(threadId) ?? null + const turnId = readCodexTurnId(params) ?? currentTurnIdFor(threadId) return identityFor(threadId, turnId, item) } }) @@ -167,7 +183,7 @@ export function createCodexJournalTranslator( if (!item) { return false } - const turnId = readCodexTurnId(event.params) ?? currentTurnIds.get(event.threadId) ?? null + const turnId = readCodexTurnId(event.params) ?? currentTurnIdFor(event.threadId) const identity = identityFor(event.threadId, turnId, item) const translated = codexJournalItem(item) const command = readString(item, 'command') @@ -238,7 +254,7 @@ export function createCodexJournalTranslator( if (!turnId) { continue } - currentTurnIds.set(threadId, turnId) + currentTurnIds.set(threadId, new Set([turnId])) for (const item of Array.isArray(turn.items) ? turn.items : []) { handleItemEvent({ threadId, method: 'item/completed', params: { turnId, item } }) } @@ -250,8 +266,11 @@ export function createCodexJournalTranslator( handle: (event) => { if (event.type === 'ended') { streams.flush() - for (const [threadId, turnId] of currentTurnIds) { - publishTurnLifecycle(event.sessionId, threadId, turnId, 'completed') + for (const [threadId, turnIds] of currentTurnIds) { + for (const turnId of turnIds) { + publishTurnLifecycle(event.sessionId, threadId, turnId, 'completed') + ordinals.forgetTurn(threadId, turnId) + } } currentTurnIds.clear() return @@ -279,20 +298,20 @@ export function createCodexJournalTranslator( if (event.method === 'turn/started') { const turnId = readCodexTurnId(event.params) if (turnId) { - currentTurnIds.set(event.threadId, turnId) + rememberTurn(event.threadId, turnId) publishTurnLifecycle(event.sessionId, event.threadId, turnId, 'running') } return } if (event.method === 'turn/completed') { - const turnId = readCodexTurnId(event.params) ?? currentTurnIds.get(event.threadId) + const turnId = readCodexTurnId(event.params) ?? currentTurnIdFor(event.threadId) if (turnId) { publishTurnLifecycle(event.sessionId, event.threadId, turnId, 'completed') ordinals.forgetTurn(event.threadId, turnId) + forgetTurn(event.threadId, turnId) } - // A later item with no turn of its own belongs to no turn, not to the - // one that just ended. - currentTurnIds.delete(event.threadId) + // A later item without its own turn id falls back to another active + // turn, if one exists; completed turns are never adopted again. return } if (event.method === 'item/started' || event.method === 'item/completed') { diff --git a/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.test.tsx b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.test.tsx index 331948ad9bf..d6b3eece713 100644 --- a/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.test.tsx +++ b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.test.tsx @@ -2,6 +2,7 @@ import { act, renderHook, waitFor } from '@testing-library/react' import { beforeEach, describe, expect, it, vi } from 'vitest' +import type { AgentJournalSubmission } from '../../../../shared/agent-session-journal-types' import type { AgentSessionWireRefusalCode } from '../../../../shared/agent-session-wire' const mocks = vi.hoisted(() => ({ @@ -46,6 +47,50 @@ function acceptedResult(fence: number) { } } +function acceptedResultFor(clientMessageId: string, fence: number) { + return { + ok: true, + replayed: false, + fence, + cursor: { epoch: 'epoch-1', sequence: fence }, + value: { + clientMessageId, + submission: { + clientMessageId, + fence, + payloadFingerprint: 'fingerprint', + dispatchState: 'accepted', + providerItemId: `provider-${clientMessageId}`, + reason: null, + submittedAt: fence, + resolvedAt: fence + } + } + } +} + +function unknownResultFor(clientMessageId: string, submittedAt: number) { + return { + ok: true, + replayed: false, + fence: 1, + cursor: { epoch: 'epoch-1', sequence: submittedAt }, + value: { + clientMessageId, + submission: { + clientMessageId, + fence: 1, + payloadFingerprint: 'fingerprint', + dispatchState: 'unknown' as const, + providerItemId: null, + reason: 'socket closed', + submittedAt, + resolvedAt: submittedAt + } + } + } +} + function refusedResult(code: AgentSessionWireRefusalCode) { return { ok: false, refusal: { code, message: code } } } @@ -175,4 +220,126 @@ describe('useStructuredAgentSessionOutbox', () => { } }) }) + + it('retries an unknown head and advances a queued tail', async () => { + vi.mocked(globalThis.crypto.randomUUID) + .mockReturnValueOnce('aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa') + .mockReturnValueOnce('bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb') + mocks.call + .mockImplementationOnce(async (_target, _method, params) => { + const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope + .clientOperationId + return unknownResultFor(clientMessageId, 10) + }) + .mockImplementationOnce(async (_target, _method, params) => { + const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope + .clientOperationId + return acceptedResultFor(clientMessageId, 11) + }) + .mockImplementationOnce(async (_target, _method, params) => { + const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope + .clientOperationId + return acceptedResultFor(clientMessageId, 12) + }) + const { result, rerender } = renderHook( + ({ submissions }: { submissions: readonly AgentJournalSubmission[] }) => + useStructuredAgentSessionOutbox({ + sessionId: 'session-1', + target: LOCAL_TARGET, + fence: 1, + submissions + }), + { initialProps: { submissions: [] as readonly AgentJournalSubmission[] } } + ) + + act(() => { + expect(result.current.send('first')).toBe(true) + }) + await waitFor(() => expect(result.current.outbox[0]?.state).toBe('unconfirmed')) + const firstId = result.current.outbox[0]!.clientMessageId + rerender({ + submissions: [ + { + clientMessageId: firstId, + fence: 1, + payloadFingerprint: 'fingerprint', + dispatchState: 'unknown', + providerItemId: null, + reason: 'socket closed', + submittedAt: 10, + resolvedAt: 10 + } + ] + }) + act(() => { + expect(result.current.send('second')).toBe(true) + }) + expect(result.current.outbox).toHaveLength(2) + + act(() => result.current.retry(firstId)) + await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(3)) + await waitFor(() => expect(result.current.outbox).toHaveLength(0)) + const retryParams = mocks.call.mock.calls[1]?.[2] as { retryUnknown?: true } | undefined + expect(retryParams?.retryUnknown).toBe(true) + }) + + it('rotates a history-rejected unknown head so the queued tail can advance', async () => { + vi.mocked(globalThis.crypto.randomUUID) + .mockReturnValueOnce('aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa') + .mockReturnValueOnce('bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb') + .mockReturnValueOnce('cccccccc-cccc-4ccc-8ccc-cccccccccccc') + mocks.call + .mockImplementationOnce(async (_target, _method, params) => { + const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope + .clientOperationId + return unknownResultFor(clientMessageId, 10) + }) + .mockImplementationOnce(async (_target, _method, params) => { + const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope + .clientOperationId + return acceptedResultFor(clientMessageId, 11) + }) + .mockImplementationOnce(async (_target, _method, params) => { + const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope + .clientOperationId + return acceptedResultFor(clientMessageId, 12) + }) + const { result, rerender } = renderHook( + ({ submissions }: { submissions: readonly AgentJournalSubmission[] }) => + useStructuredAgentSessionOutbox({ + sessionId: 'session-1', + target: LOCAL_TARGET, + fence: 1, + submissions + }), + { initialProps: { submissions: [] as readonly AgentJournalSubmission[] } } + ) + + act(() => expect(result.current.send('first')).toBe(true)) + await waitFor(() => expect(result.current.outbox[0]?.state).toBe('unconfirmed')) + const firstId = result.current.outbox[0]!.clientMessageId + act(() => expect(result.current.send('second')).toBe(true)) + rerender({ + submissions: [ + { + clientMessageId: firstId, + fence: 1, + payloadFingerprint: 'fingerprint', + dispatchState: 'rejected', + providerItemId: null, + reason: 'not_delivered', + submittedAt: 10, + resolvedAt: 10 + } + ] + }) + + act(() => result.current.retry(firstId)) + await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(3)) + await waitFor(() => expect(result.current.outbox).toHaveLength(0)) + const retryParams = mocks.call.mock.calls[1]?.[2] as + | { envelope: { clientOperationId: string } } + | undefined + expect(retryParams?.envelope.clientOperationId).not.toBe(firstId) + }) }) diff --git a/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts index 457fcad4b77..d8b9dd61dfa 100644 --- a/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts +++ b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts @@ -247,13 +247,38 @@ export function useStructuredAgentSessionOutbox(args: { const retry = (clientMessageId: string): void => { blockedIdRef.current = null setError(null) - const unknown = submissions.find( - (submission) => - submission.clientMessageId === clientMessageId && submission.dispatchState === 'unknown' + const submission = submissions.find( + (candidate) => candidate.clientMessageId === clientMessageId ) const current = outboxRef.current.find((entry) => entry.clientMessageId === clientMessageId) + // A provider-history reconciliation can settle an earlier unknown as + // rejected before the user presses Retry. Reusing that operation id only + // replays the settled rejection forever, so rotate the id for a safe resend. + if (current && submission?.dispatchState === 'rejected') { + const rotated = outboxRef.current.map((entry) => + entry.clientMessageId === clientMessageId + ? { + ...entry, + clientMessageId: structuredSessionOperationId(), + state: 'queued' as const, + retryAfterUnknownSubmittedAt: null + } + : entry + ) + if (!writeOutbox(sessionId, rotated)) { + setError('Message could not be saved to the outbox') + return + } + outboxRef.current = rotated + setOutbox(rotated) + return + } const retryAfterUnknownSubmittedAt = - unknown?.submittedAt ?? (current?.state === 'unconfirmed' ? -1 : null) + submission?.dispatchState === 'unknown' + ? submission.submittedAt + : current?.state === 'unconfirmed' + ? -1 + : null const next = outboxRef.current.map((entry) => entry.clientMessageId === clientMessageId ? {