fix(codex): preserve unsettled dispatch correlations

This commit is contained in:
Merge Sim
2026-09-12 01:01:49 -07:00
parent a60b714e0f
commit 7d3c57c8de
5 changed files with 84 additions and 33 deletions
@@ -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[] = []
@@ -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)
})
})
@@ -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<string>()
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),
@@ -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)
})
+23 -19
View File
@@ -53,30 +53,28 @@ function turnInputFor(body: AgentJournalMessageItem): Record<string, unknown>[]
}
/**
* 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<void> {
): Promise<boolean> {
// 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<AgentSessionDispatchOutcome> {
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' }