From 7d3c57c8deac108c378bd9967c5e45f83fa95b46 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Sat, 12 Sep 2026 01:01:49 -0700 Subject: [PATCH] fix(codex): preserve unsettled dispatch correlations --- ...odex-structured-dispatch-admission.test.ts | 48 +++++++++++++++++++ .../codex-structured-dispatch-echo.test.ts | 11 +++-- .../codex/codex-structured-dispatch-echo.ts | 15 +++--- .../codex-structured-dispatch-test-support.ts | 1 + src/main/codex/codex-structured-turn-start.ts | 42 ++++++++-------- 5 files changed, 84 insertions(+), 33 deletions(-) diff --git a/src/main/codex/codex-structured-dispatch-admission.test.ts b/src/main/codex/codex-structured-dispatch-admission.test.ts index 70878434641..7d9c94efe39 100644 --- a/src/main/codex/codex-structured-dispatch-admission.test.ts +++ b/src/main/codex/codex-structured-dispatch-admission.test.ts @@ -1,4 +1,5 @@ import { describe, expect, it } from 'vitest' +import { MAX_CODEX_PENDING_DISPATCH_ECHOES } from './codex-structured-dispatch-echo' import { acquiredCodexAdapter, echoUserMessage, @@ -122,6 +123,53 @@ describe('codex dispatch admission', () => { expect(settlements).toEqual([]) }) + it('retains correlation when a request fails after its write may have landed', async () => { + const codex = fakeCodexAppServer({ + 'turn/start': () => { + throw new Error('request timed out after write') + } + }) + const settlements: LateSettlement[] = [] + const adapter = await acquiredCodexAdapter({ codex, settlements }) + const connection = codex.connections[0]! + startTurn(connection, 'turn-1') + + await expect(send(adapter, 'client-1')).rejects.toThrow('request timed out after write') + echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u1', clientId: 'client-1' }) + + expect(settlements).toEqual([ + { + sessionId: 'session-1', + clientMessageId: 'client-1', + providerIdentity: { + provider: 'codex', + threadId: CODEX_TEST_THREAD_ID, + turnId: 'turn-1', + ordinal: 0 + } + } + ]) + }) + + it('refuses overflow without discarding an older accepted send', async () => { + const codex = fakeCodexAppServer({ 'turn/start': () => ({ turn: { id: 'turn-1' } }) }) + const settlements: LateSettlement[] = [] + const adapter = await acquiredCodexAdapter({ codex, settlements }) + const connection = codex.connections[0]! + startTurn(connection, 'turn-1') + + for (let index = 0; index < MAX_CODEX_PENDING_DISPATCH_ECHOES; index += 1) { + expect(await send(adapter, `client-${index}`)).toEqual({ state: 'admitted' }) + } + expect(await send(adapter, 'client-overflow')).toEqual({ + state: 'rejected', + reason: 'codex structured dispatch queue is full' + }) + + echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u0', clientId: 'client-0' }) + expect(settlements.map(({ clientMessageId }) => clientMessageId)).toEqual(['client-0']) + }) + it('leaves no waiter behind when the session closes', async () => { const codex = fakeCodexAppServer({ 'turn/start': () => ({ turn: { id: 'turn-1' } }) }) const settlements: LateSettlement[] = [] diff --git a/src/main/codex/codex-structured-dispatch-echo.test.ts b/src/main/codex/codex-structured-dispatch-echo.test.ts index b53cd0d1a3e..58203b26964 100644 --- a/src/main/codex/codex-structured-dispatch-echo.test.ts +++ b/src/main/codex/codex-structured-dispatch-echo.test.ts @@ -61,15 +61,16 @@ describe('codex dispatch echoes', () => { expect(echoes.settle('client-1')).toBe(false) }) - it('retains a bounded window, dropping the oldest first', () => { + it('refuses new correlations at capacity without dropping an older send', () => { const echoes = createCodexDispatchEchoes() - for (let index = 0; index <= MAX_CODEX_PENDING_DISPATCH_ECHOES; index += 1) { - echoes.arm(`client-${index}`) + for (let index = 0; index < MAX_CODEX_PENDING_DISPATCH_ECHOES; index += 1) { + expect(echoes.arm(`client-${index}`)).toBe(true) } + expect(echoes.arm(`client-${MAX_CODEX_PENDING_DISPATCH_ECHOES}`)).toBe(false) expect(echoes.size).toBe(MAX_CODEX_PENDING_DISPATCH_ECHOES) - expect(echoes.settle('client-0')).toBe(false) - expect(echoes.settle(`client-${MAX_CODEX_PENDING_DISPATCH_ECHOES}`)).toBe(true) + expect(echoes.settle('client-0')).toBe(true) + expect(echoes.settle(`client-${MAX_CODEX_PENDING_DISPATCH_ECHOES}`)).toBe(false) }) }) diff --git a/src/main/codex/codex-structured-dispatch-echo.ts b/src/main/codex/codex-structured-dispatch-echo.ts index 80c682477fb..8ea97561c59 100644 --- a/src/main/codex/codex-structured-dispatch-echo.ts +++ b/src/main/codex/codex-structured-dispatch-echo.ts @@ -13,8 +13,8 @@ export const MAX_CODEX_PENDING_DISPATCH_ECHOES = 256 * their echoes arrive far apart. Queue position identifies neither. */ export type CodexDispatchEchoes = { - /** Arms settlement for a send about to be written. */ - arm: (clientMessageId: string) => void + /** Arms settlement for a send about to be written; false preserves older waits at capacity. */ + arm: (clientMessageId: string) => boolean /** True once, for a send this session armed and has not yet settled. */ settle: (clientMessageId: string) => boolean /** Drops an armed send whose write never reached the provider. */ @@ -27,15 +27,12 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes { const armed = new Set() return { arm(clientMessageId) { + if (!armed.has(clientMessageId) && armed.size >= MAX_CODEX_PENDING_DISPATCH_ECHOES) { + return false + } armed.delete(clientMessageId) armed.add(clientMessageId) - while (armed.size > MAX_CODEX_PENDING_DISPATCH_ECHOES) { - const oldest = armed.values().next().value - if (typeof oldest !== 'string') { - break - } - armed.delete(oldest) - } + return true }, settle: (clientMessageId) => armed.delete(clientMessageId), disarm: (clientMessageId) => void armed.delete(clientMessageId), diff --git a/src/main/codex/codex-structured-dispatch-test-support.ts b/src/main/codex/codex-structured-dispatch-test-support.ts index ddfca7f8c08..94891838863 100644 --- a/src/main/codex/codex-structured-dispatch-test-support.ts +++ b/src/main/codex/codex-structured-dispatch-test-support.ts @@ -95,6 +95,7 @@ export async function acquiredCodexAdapter(input: { }), openConnection: input.codex.openConnection, readProcessStartTime: async () => 1_700_000_000_000, + captureTurnProcesses: async () => null, now: () => 1_700_000_000_500, onDispatchSettledLate: (settlement) => input.settlements.push(settlement) }) diff --git a/src/main/codex/codex-structured-turn-start.ts b/src/main/codex/codex-structured-turn-start.ts index cb110ece982..92d7d217fef 100644 --- a/src/main/codex/codex-structured-turn-start.ts +++ b/src/main/codex/codex-structured-turn-start.ts @@ -53,30 +53,28 @@ function turnInputFor(body: AgentJournalMessageItem): Record[] } /** - * Hands one submission to Codex. Resolves when Codex has taken it; throws only - * for outcomes the wire must not read as acceptance. + * Hands one submission to Codex. False means the bounded correlation window + * refused it before the write; otherwise resolves when Codex has taken it. */ export async function startCodexTurn( host: CodexTurnHost, input: { clientMessageId: string; body: AgentJournalMessageItem; timeoutMs?: number } -): Promise { +): Promise { // Armed before the write: the echo can land while the response is in flight. - host.dispatchEchoes.arm(input.clientMessageId) - try { - await host.connection.request( - 'turn/start', - { - threadId: host.threadId, - clientUserMessageId: input.clientMessageId, - input: turnInputFor(input.body), - ...Object.fromEntries(host.options) - }, - { timeoutMs: input.timeoutMs } - ) - } catch (error) { - host.dispatchEchoes.disarm(input.clientMessageId) - throw error + if (!host.dispatchEchoes.arm(input.clientMessageId)) { + return false } + await host.connection.request( + 'turn/start', + { + threadId: host.threadId, + clientUserMessageId: input.clientMessageId, + input: turnInputFor(input.body), + ...Object.fromEntries(host.options) + }, + { timeoutMs: input.timeoutMs } + ) + return true } /** @@ -91,11 +89,17 @@ export async function dispatchCodexTurn( timeoutMs: number | undefined ): Promise { try { - await startCodexTurn(session, { ...input, timeoutMs }) + if (!(await startCodexTurn(session, { ...input, timeoutMs }))) { + return { state: 'rejected', reason: 'codex structured dispatch queue is full' } + } } catch (error) { if (isCodexAppServerRequestError(error) || isCodexAppServerUnsupportedError(error)) { + // Codex answered and declined, so no echo for this write can arrive. + session.dispatchEchoes.disarm(input.clientMessageId) return { state: 'rejected', reason: (error as Error).message } } + // A timeout or transport failure can happen after the frame was written. + // Keep the correlation armed so a later echo can prove delivery. throw error } return { state: 'admitted' }