fix(codex): settle a structured send on admission, and stop minting a colliding identity

Two sends could be written into the journal under one durable identity.

Codex coalesces a mid-turn `turn/start` into the running turn rather than
refusing it -- measured against real `codex app-server` builds 0.147.0,
0.150.1 and 0.153.4, none of which refuse and none of which fire a second
`turn/started`. The dispatch path read the turn id from the turn/start
response and stamped every accepted send `ordinal: 0`. Since a coalesced
send gets the running turn's id back, two submissions persisted the same
`providerItemId`. That string is durable, and it is the key a restore uses
to match a submission against provider history, so the second message's real
history row matched nothing and rendered as an extra bubble on replay.

On 0.147.0 it is worse than a collision: the coalesced response returns a
turn id that never starts and never completes, so the persisted key named a
turn absent from history and NEITHER message could match.

Identity is now minted from the echoed user message at `identityFor` -- the
single point that mints the journal row's own identity -- so the settled key
is by construction the one replay computes, rather than a parallel
calculation that can drift.

Dispatch returns `admitted` when the transport write completes; identity
settles on the echo through a channel that did not previously exist for
Codex. Waiters are keyed by client message id instead of being shifted off
the front of an array by arrival order, and they are cleared on session
close and child exit -- previously a timeout was the only thing that ever
ended one.

`TURN_ID_WAIT_MS` is deleted. It was never reachable on any build measured:
`readCodexTurnId` returns non-null on all three, so the 10s wait never
fired. The comment justifying it claimed older builds acknowledge before the
id exists, which no tested build does.

Three comments asserting Codex answers a mid-turn send with `turn already
running` are corrected. Their only backing was a test fixture inventing that
error string. The correction is factual only -- every changed line in
`src/main/runtime/orchestration/` is a comment, and mid-turn delivery is
still refused for both providers. Whether that policy is right is a separate
question; it was resting on a false premise.

Known gap, stated rather than implied: this prevents new collisions and does
not repair journals already written with a colliding or phantom key. Those
conversations keep duplicating on restore. Repairing them means re-matching
persisted submissions against provider history and rewriting
`providerItemId` -- which is what `journal-submission-reconciler.ts` is
written for, and it still has no production caller.
This commit is contained in:
Merge Sim
2026-09-11 11:31:17 -07:00
parent ec9c3e0550
commit 698eb3a74d
24 changed files with 578 additions and 138 deletions
@@ -1,3 +1,4 @@
import { createCodexDispatchEchoes } from './codex-structured-dispatch-echo'
import { afterEach, describe, expect, it, vi } from 'vitest'
import type { AgentJournalItemBody } from '../../shared/agent-session-journal-types'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
@@ -57,6 +58,7 @@ describe('requested-close durable turn timing', () => {
acquisitionGeneration: 'generation-1',
threadId: 'thread-1',
prompts: { clear: vi.fn() },
dispatchEchoes: createCodexDispatchEchoes(),
translator
} as unknown as CodexSession
const sessions = new Map([['session-1', session]])
@@ -0,0 +1,152 @@
import { describe, expect, it } from 'vitest'
import {
acquiredCodexAdapter,
echoUserMessage,
fakeCodexAppServer,
startTurn,
CODEX_TEST_THREAD_ID,
CODEX_TEST_USER_MESSAGE,
type LateSettlement
} from './codex-structured-dispatch-test-support'
function send(
adapter: Awaited<ReturnType<typeof acquiredCodexAdapter>>,
clientMessageId: string
): Promise<unknown> {
return adapter.dispatch({
sessionId: 'session-1',
clientMessageId,
body: CODEX_TEST_USER_MESSAGE,
fence: 7
})
}
describe('codex dispatch admission', () => {
it('admits a send queued behind a running turn and settles it when Codex echoes it', async () => {
// Measured on codex-cli 0.153.4: a `turn/start` issued while a turn runs is
// COALESCED into it -- same turn id back, no second `turn/started`, and the
// user message echoed only once the running turn reaches it.
const codex = fakeCodexAppServer({
'turn/start': () => ({ turn: { id: 'turn-1', status: 'inProgress' } })
})
const settlements: LateSettlement[] = []
const adapter = await acquiredCodexAdapter({ codex, settlements })
const connection = codex.connections[0]!
startTurn(connection, 'turn-1')
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u1', clientId: 'client-1' })
const outcome = await send(adapter, 'client-2')
// No doubt: elapsed time is not evidence, so nothing invites a Retry.
expect(outcome).toEqual({ state: 'admitted' })
expect(settlements).toEqual([])
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u2', clientId: 'client-2' })
// Ordinal 1, not 0: the queued send is the SECOND user message of the turn
// it was coalesced into, which is the key a history replay computes for it.
expect(settlements).toEqual([
{
sessionId: 'session-1',
clientMessageId: 'client-2',
providerIdentity: {
provider: 'codex',
threadId: CODEX_TEST_THREAD_ID,
turnId: 'turn-1',
ordinal: 1
}
}
])
})
it('correlates each send by client message id, not queue order', 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')
await send(adapter, 'client-1')
await send(adapter, 'client-2')
// The echoes arrive in the opposite order to the sends.
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u2', clientId: 'client-2' })
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u1', clientId: 'client-1' })
expect(
settlements.map((settlement) => [
settlement.clientMessageId,
(settlement.providerIdentity as { ordinal: number }).ordinal
])
).toEqual([
['client-2', 0],
['client-1', 1]
])
})
it('settles nothing for a user message this session never sent', 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')
await send(adapter, 'client-1')
// A message another client sent on the same thread, and one Codex did not
// correlate at all.
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-x', clientId: 'someone-else' })
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-y' })
expect(settlements).toEqual([])
})
it('rejects only when Codex answered and declined, and arms nothing for it', async () => {
const { CodexAppServerRequestError } = await import('./codex-app-server-connection')
const codex = fakeCodexAppServer({
'turn/start': () => {
throw new CodexAppServerRequestError('turn/start', -32602, 'thread not found')
}
})
const settlements: LateSettlement[] = []
const adapter = await acquiredCodexAdapter({ codex, settlements })
const connection = codex.connections[0]!
startTurn(connection, 'turn-1')
expect(await send(adapter, 'client-1')).toEqual({
state: 'rejected',
reason: 'thread not found'
})
// A refused write is disarmed, so a later echo of that id settles nothing.
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u1', clientId: 'client-1' })
expect(settlements).toEqual([])
})
it('leaves no waiter behind when the session closes', 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')
await send(adapter, 'client-1')
await adapter.closeSession('session-1')
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u1', clientId: 'client-1' })
expect(settlements).toEqual([])
})
it('leaves no waiter behind when the child exits', 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')
await send(adapter, 'client-1')
connection.handlers.onExit?.(new Error('codex app-server exited'))
echoUserMessage(connection, { turnId: 'turn-1', itemId: 'item-u1', clientId: 'client-1' })
expect(settlements).toEqual([])
})
})
@@ -0,0 +1,107 @@
import { describe, expect, it } from 'vitest'
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
import {
createCodexDispatchEchoes,
readCodexDispatchEcho,
MAX_CODEX_PENDING_DISPATCH_ECHOES
} from './codex-structured-dispatch-echo'
const CODEX_IDENTITY: AgentJournalItemIdentity = {
provider: 'codex',
threadId: 'thread-1',
turnId: 'turn-1',
ordinal: 3
}
describe('codex dispatch echoes', () => {
it('settles by client message id rather than arrival order', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.arm('client-2')
// Codex coalesces both sends into one turn, and the second can be echoed
// first. Queue position would settle the wrong submission here.
expect(echoes.settle('client-2')).toBe(true)
expect(echoes.settle('client-1')).toBe(true)
expect(echoes.size).toBe(0)
})
it('refuses an echo this session never armed', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
expect(echoes.settle('client-from-history')).toBe(false)
expect(echoes.size).toBe(1)
})
it('settles a send exactly once', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
expect(echoes.settle('client-1')).toBe(true)
expect(echoes.settle('client-1')).toBe(false)
})
it('drops a send whose write never reached the provider', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.disarm('client-1')
expect(echoes.settle('client-1')).toBe(false)
})
it('clears every armed send', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.arm('client-2')
echoes.clear()
expect(echoes.size).toBe(0)
expect(echoes.settle('client-1')).toBe(false)
})
it('retains a bounded window, dropping the oldest first', () => {
const echoes = createCodexDispatchEchoes()
for (let index = 0; index <= MAX_CODEX_PENDING_DISPATCH_ECHOES; index += 1) {
echoes.arm(`client-${index}`)
}
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)
})
})
describe('readCodexDispatchEcho', () => {
it('reads the client message id off a user message', () => {
expect(
readCodexDispatchEcho(
{ type: 'userMessage', id: 'item-1', clientId: 'client-1' },
CODEX_IDENTITY
)
).toEqual({ clientMessageId: 'client-1', providerIdentity: CODEX_IDENTITY })
})
it('ignores an item that is not a user message', () => {
expect(
readCodexDispatchEcho(
{ type: 'agentMessage', id: 'item-1', clientId: 'client-1' },
CODEX_IDENTITY
)
).toBeNull()
})
it('ignores a user message Codex did not correlate', () => {
expect(readCodexDispatchEcho({ type: 'userMessage', id: 'item-1' }, CODEX_IDENTITY)).toBeNull()
})
it('ignores an item with no durable Codex identity', () => {
expect(
readCodexDispatchEcho(
{ type: 'userMessage', id: 'item-1', clientId: 'client-1' },
{ provider: 'orca', clientMessageId: 'codex-item:thread-1:item-1' }
)
).toBeNull()
})
})
@@ -0,0 +1,61 @@
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
/** Sends awaiting their echo, oldest first. A send whose echo never arrives is
* retired by the journal's pending-submission recovery on exit, not from here. */
export const MAX_CODEX_PENDING_DISPATCH_ECHOES = 256
/**
* Which sends this session is still waiting to hear back about, keyed by the
* client message id Codex echoes on the user message.
*
* Keyed rather than ordered on purpose: Codex coalesces a `turn/start` issued
* while a turn is running into that turn, so two sends can share one turn id and
* 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
/** 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. */
disarm: (clientMessageId: string) => void
clear: () => void
readonly size: number
}
export function createCodexDispatchEchoes(): CodexDispatchEchoes {
const armed = new Set<string>()
return {
arm(clientMessageId) {
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)
}
},
settle: (clientMessageId) => armed.delete(clientMessageId),
disarm: (clientMessageId) => void armed.delete(clientMessageId),
clear: () => armed.clear(),
get size() {
return armed.size
}
}
}
/** The user-message echo a settlement is read off, or null for any other item. */
export function readCodexDispatchEcho(
item: { type: string; id: string } & Record<string, unknown>,
identity: AgentJournalItemIdentity
): { clientMessageId: string; providerIdentity: AgentJournalItemIdentity } | null {
if (item.type !== 'userMessage' || identity.provider !== 'codex') {
return null
}
const clientMessageId = item.clientId
return typeof clientMessageId === 'string' && clientMessageId.length > 0
? { clientMessageId, providerIdentity: identity }
: null
}
@@ -0,0 +1,138 @@
import type {
AgentJournalItemIdentity,
AgentJournalMessageItem,
AgentSessionJournalIdentity
} from '../../shared/agent-session-journal-types'
import type {
CodexAppServerConnection,
CodexAppServerConnectionHandlers,
CodexAppServerLaunch,
openCodexAppServerConnection
} from './codex-app-server-connection'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { CodexStructuredSessionAdapter } from './codex-structured-session-adapter'
export const CODEX_TEST_THREAD_ID = 'thread-abc'
export const CODEX_TEST_USER_MESSAGE: AgentJournalMessageItem = {
kind: 'message',
role: 'user',
blocks: [{ type: 'text', text: 'ship it' }]
}
export type CodexTestRoute = (params: Record<string, unknown> | undefined) => unknown
type FakeConnection = Omit<CodexAppServerConnection, 'closed'> & {
closed: boolean
launch: CodexAppServerLaunch
handlers: CodexAppServerConnectionHandlers
calls: { method: string; params?: Record<string, unknown> }[]
}
export type LateSettlement = {
sessionId: string
clientMessageId: string
providerIdentity: AgentJournalItemIdentity
}
/** A `codex app-server` whose turn traffic the test drives by hand. */
export function fakeCodexAppServer(routes: Record<string, CodexTestRoute> = {}): {
connections: FakeConnection[]
openConnection: typeof openCodexAppServerConnection
routes: Record<string, CodexTestRoute>
} {
const connections: FakeConnection[] = []
const openConnection = (async (launch, handlers = {}) => {
const connection: FakeConnection = {
launch,
handlers,
calls: [],
pid: 4321,
closed: false,
request: async (method, params) => {
connection.calls.push({ method, params })
return routes[method]?.(params) ?? {}
},
notify: () => {},
respond: () => {},
respondWithError: () => {},
close: async () => {
connection.closed = true
return true
}
}
connections.push(connection)
return connection
}) as typeof openCodexAppServerConnection
routes['thread/start'] ??= () => ({
thread: { id: CODEX_TEST_THREAD_ID, path: '/rollouts/abc.jsonl' }
})
return { connections, openConnection, routes }
}
/** A sink that records nothing but keeps the translator alive, which is what
* mints the identities a late settlement carries. */
export function recordingSink(): StructuredAgentSessionEventSink {
return {
appendItem: () => {},
appendTombstone: () => {},
publish: () => {}
} as unknown as StructuredAgentSessionEventSink
}
export async function acquiredCodexAdapter(input: {
codex: ReturnType<typeof fakeCodexAppServer>
settlements: LateSettlement[]
sink?: StructuredAgentSessionEventSink
}): Promise<CodexStructuredSessionAdapter> {
const adapter = new CodexStructuredSessionAdapter({
resolveLaunch: async () => ({
command: 'codex',
args: ['app-server'],
cwd: '/work/repo',
codexHome: null,
resumeThreadId: null
}),
openConnection: input.codex.openConnection,
readProcessStartTime: async () => 1_700_000_000_000,
now: () => 1_700_000_000_500,
onDispatchSettledLate: (settlement) => input.settlements.push(settlement)
})
const identity: AgentSessionJournalIdentity = {
sessionId: 'session-1',
workspaceId: 'ws-1',
hostId: 'host-1',
agent: 'codex',
providerHandle: { kind: 'codex', threadId: CODEX_TEST_THREAD_ID }
}
await adapter.acquire({
identity,
fence: 7,
spawnToken: 'spawn-9',
events: input.sink ?? recordingSink()
})
return adapter
}
/** Codex's own echo of a user message Orca sent, inside `turnId`. */
export function echoUserMessage(
connection: FakeConnection,
input: { turnId: string; itemId: string; clientId?: string; threadId?: string }
): void {
connection.handlers.onNotification?.('item/started', {
threadId: input.threadId ?? CODEX_TEST_THREAD_ID,
turn: { id: input.turnId },
item: {
type: 'userMessage',
id: input.itemId,
...(input.clientId ? { clientId: input.clientId } : {})
}
})
}
export function startTurn(connection: FakeConnection, turnId: string): void {
connection.handlers.onNotification?.('turn/started', {
threadId: CODEX_TEST_THREAD_ID,
turn: { id: turnId }
})
}
@@ -1,3 +1,4 @@
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
import type { AgentSessionDeltaCoalescerDeps } from '../native-chat/agent-session-wire/agent-session-delta-coalescer'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type { CodexStructuredSessionEvent } from './codex-structured-session-adapter'
@@ -9,6 +10,9 @@ export type CodexJournalTranslatorDeps = {
sessionId?: string
now?: () => number
bindPromptItemId?: (journalItemId: string, threadId: string, promptKey: string) => void
/** Settles a send's identity off the echoed user message, using the very
* identity the journal row carries so a replay computes the same key. */
onUserMessageEcho?: (clientMessageId: string, identity: AgentJournalItemIdentity) => void
primaryThreadId?: () => string | null
subagentExecutions?: CodexSubagentExecutions
coalesceMs?: number
@@ -29,6 +29,7 @@ import { appendCodexLifecycleItem, publishCodexLifecycle } from './codex-structu
import type { CodexActiveJournalItem } from './codex-structured-journal-settlement'
import { readCodexJournalString } from './codex-structured-journal-translation-values'
import { readCodexTurnId } from './codex-structured-thread-facts'
import { readCodexDispatchEcho } from './codex-structured-dispatch-echo'
export class CodexJournalItems {
readonly ordinals = new CodexTurnOrdinals()
@@ -40,7 +41,7 @@ export class CodexJournalItems {
constructor(
private readonly deps: Pick<
CodexJournalTranslatorDeps,
'sink' | 'coalesceMs' | 'maxRetainedBytes' | 'schedule'
'sink' | 'coalesceMs' | 'maxRetainedBytes' | 'schedule' | 'onUserMessageEcho'
> & { maxMetadataBytes?: number },
private readonly activeTurn: (threadId: string) => string | null,
private readonly suppress: (threadId: string, turnId: string) => void
@@ -78,6 +79,10 @@ export class CodexJournalItems {
const identity = this.identityFor(event.threadId, turnId, item)
// Count echoes for stable resume ordinals, but user bubbles come from submissions.
if (source === 'live' && item.type === 'userMessage') {
const echo = readCodexDispatchEcho(item, identity)
if (echo) {
this.deps.onUserMessageEcho?.(echo.clientMessageId, echo.providerIdentity)
}
return { handled: true, admission: CODEX_JOURNAL_ADMITTED }
}
if (item.type === 'contextCompaction' && event.method === 'item/started') {
@@ -6,7 +6,7 @@ import type {
} from '../../shared/agent-session-journal-types'
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import { agentJournalTurnBody } from '../../shared/agent-session-turn-record'
import { CODEX_USER_MESSAGE_ORDINAL } from './codex-structured-turn-start'
import { CODEX_USER_MESSAGE_ORDINAL } from './codex-turn-ordinals'
import type {
StructuredAgentSessionEventSink,
StructuredAgentSessionSinkAdmission
@@ -3,7 +3,7 @@ import { disposeCodexServerRequest } from './codex-server-request-disposition'
import type { CodexJournalTranslationAdmission } from './codex-structured-journal-translation'
import * as codexRewind from './codex-structured-rewind'
import type { CodexSession, CodexStructuredSessionEvent } from './codex-structured-session-state'
import { readCodexThreadId, readCodexTurnId } from './codex-structured-thread-facts'
import { readCodexThreadId } from './codex-structured-thread-facts'
import type { CodexStructuredTurnCancellation } from './codex-structured-turn-cancellation'
type EmitCodexEvent = (
@@ -41,10 +41,9 @@ export function deliverCodexNotification(
return { accepted: true }
}
const threadId = readCodexThreadId(params) ?? session.threadId
const turnId =
method === 'turn/started' && threadId === session.threadId ? readCodexTurnId(params) : null
const turnWaiter = turnId ? session.turnIdWaiters[0] : undefined
const admission = emit(session, {
// Dispatch identity settles on the user-message echo inside the translator,
// which is where the ordinal a replay will compute is minted.
return emit(session, {
type: 'notification',
sessionId,
threadId,
@@ -52,13 +51,6 @@ export function deliverCodexNotification(
params,
...(observedAt !== undefined ? { observedAt } : {})
})
if (method === 'turn/started' && threadId === session.threadId) {
if (admission.accepted && turnId && session.turnIdWaiters[0] === turnWaiter) {
session.turnIdWaiters.shift()
turnWaiter?.(turnId)
}
}
return admission
}
export function deliverCodexServerRequest(
@@ -10,6 +10,7 @@ import {
} from './codex-structured-acquisition-lifecycle'
import { CodexBackgroundTaskTracker } from './codex-background-task-tracker'
import { CodexSubagentExecutions } from './codex-subagent-executions'
import { createCodexDispatchEchoes } from './codex-structured-dispatch-echo'
import { createCodexJournalTranslator } from './codex-structured-journal-translation'
import { openCodexAppServerConnection } from './codex-app-server-connection'
import { codexProcessIdentity, codexProviderHandleLink } from './codex-structured-owner-identity'
@@ -77,6 +78,7 @@ export async function acquireCodexStructuredSession(input: {
? acquireInput.identity.providerHandle.threadId
: null
const subagentExecutions = new CodexSubagentExecutions()
const dispatchEchoes = createCodexDispatchEchoes()
const translator = acquireInput.events
? createCodexJournalTranslator({
sink: acquireInput.events,
@@ -85,7 +87,14 @@ export async function acquireCodexStructuredSession(input: {
primaryThreadId: () => primaryThreadId,
subagentExecutions,
bindPromptItemId: (journalItemId, threadId, promptKey) =>
acquisition.prompts.bindJournalItemId(journalItemId, threadId, promptKey)
acquisition.prompts.bindJournalItemId(journalItemId, threadId, promptKey),
onUserMessageEcho: (clientMessageId, providerIdentity) => {
// Only a send THIS session admitted; an echo from history restore or
// another client names no submission of ours to settle.
if (dispatchEchoes.settle(clientMessageId)) {
deps.onDispatchSettledLate?.({ sessionId, clientMessageId, providerIdentity })
}
}
})
: null
const open = deps.openConnection ?? openCodexAppServerConnection
@@ -207,7 +216,7 @@ export async function acquireCodexStructuredSession(input: {
prompts: acquisition.prompts,
options: restoredCodexSessionOptions(acquireInput.options),
reportedOptions: reportedCodexThreadOptions(opened),
turnIdWaiters: [],
dispatchEchoes,
translator,
backgroundTasks: new CodexBackgroundTaskTracker(opened.threadId, subagentExecutions),
forceCloseUnexpected: (reason) =>
@@ -432,7 +432,7 @@ describe('CodexStructuredSessionAdapter.acquire', () => {
})
describe('CodexStructuredSessionAdapter.dispatch', () => {
it('accepts a turn Codex names in its response', async () => {
it('admits a send as soon as Codex owns it', async () => {
const codex = fakeCodex({ 'turn/start': () => ({ turn: { id: 'turn-1' } }) })
const adapter = await acquired(codex)
@@ -451,10 +451,9 @@ describe('CodexStructuredSessionAdapter.dispatch', () => {
fence: 7
})
expect(outcome).toEqual({
state: 'accepted',
providerIdentity: { provider: 'codex', threadId: THREAD_ID, turnId: 'turn-1', ordinal: 0 }
})
// Identity is not knowable here: a send coalesced into a running turn shares
// that turn's id, so the echo settles which message landed where.
expect(outcome).toEqual({ state: 'admitted' })
expect(codex.connections[0].calls[1].params).toEqual({
threadId: THREAD_ID,
clientUserMessageId: 'client-1',
@@ -466,7 +465,7 @@ describe('CodexStructuredSessionAdapter.dispatch', () => {
})
})
it('accepts a turn named only by the notification that raced the ack', async () => {
it('admits a send on a build whose turn/start answers before the turn is named', async () => {
const codex = fakeCodex()
const events: CodexStructuredSessionEvent[] = []
const adapter = await acquired(codex, {}, events)
@@ -485,8 +484,7 @@ describe('CodexStructuredSessionAdapter.dispatch', () => {
fence: 7
})
expect(outcome).toMatchObject({ state: 'accepted' })
expect(outcome).toMatchObject({ providerIdentity: { turnId: 'turn-late' } })
expect(outcome).toEqual({ state: 'admitted' })
expect(events.at(-1)).toMatchObject({ type: 'notification', method: 'turn/started' })
})
@@ -510,10 +508,7 @@ describe('CodexStructuredSessionAdapter.dispatch', () => {
fence: 7
})
expect(outcome).toEqual({
state: 'accepted',
providerIdentity: { provider: 'codex', threadId: THREAD_ID, turnId: 'turn-root', ordinal: 0 }
})
expect(outcome).toEqual({ state: 'admitted' })
// Each event carries the thread it actually came from, so the journal can
// keep a subagent's turn out of the root conversation.
expect(events.map((event) => (event.type === 'notification' ? event.threadId : null))).toEqual([
@@ -522,29 +517,6 @@ describe('CodexStructuredSessionAdapter.dispatch', () => {
])
})
it('settles unknown rather than failed when Codex never names the turn', async () => {
vi.useFakeTimers()
try {
const codex = fakeCodex()
const adapter = await acquired(codex)
const dispatching = adapter.dispatch({
sessionId: 'session-1',
clientMessageId: 'client-1',
body: USER_MESSAGE,
fence: 7
})
await vi.advanceTimersByTimeAsync(10_000)
expect(await dispatching).toEqual({
state: 'unknown',
reason: 'codex app-server started a turn it did not name in time'
})
} finally {
vi.useRealTimers()
}
})
it('rejects only when Codex answered and declined', async () => {
const codex = fakeCodex({
'turn/start': () => {
@@ -257,9 +257,8 @@ describe('CodexStructuredSessionAdapter.cancelTurn', () => {
body: USER_MESSAGE,
fence: 7
})
).resolves.toMatchObject({
state: 'accepted',
providerIdentity: { turnId: 'turn-2' }
).resolves.toEqual({
state: 'admitted'
})
})
@@ -1,3 +1,4 @@
import { createCodexDispatchEchoes } from './codex-structured-dispatch-echo'
import { describe, expect, it, vi } from 'vitest'
import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types'
import type {
@@ -101,7 +102,7 @@ describe('Codex structured session close lifecycle', () => {
prompts,
options: new Map(),
reportedOptions: {},
turnIdWaiters: [],
dispatchEchoes: createCodexDispatchEchoes(),
translator
} as CodexSession
const sessions = new Map([['session-1', session]])
@@ -47,6 +47,9 @@ export function handleCodexSessionExit(input: {
event.settlementRetryRequired = true
}
session.ended = true
// Nothing can echo for this child any more; the journal's pending-submission
// recovery is what settles the sends these were armed for.
session.dispatchEchoes.clear()
session.backgroundTasks.clear()
input.onBackgroundTasksChanged?.(input.sessionId, null)
session.unbindReadingControl?.()
@@ -1,3 +1,4 @@
import { createCodexDispatchEchoes } from './codex-structured-dispatch-echo'
import { describe, expect, it, vi } from 'vitest'
import type { CodexAppServerConnection } from './codex-app-server-connection'
import { CodexAcquisitionWindow } from './codex-structured-acquisition-window'
@@ -31,7 +32,7 @@ function optionSession(request: CodexAppServerConnection['request']): CodexSessi
prompts: new CodexAcquisitionWindow().prompts,
options: new Map(),
reportedOptions: { model: 'gpt-live', effort: 'high' },
turnIdWaiters: [],
dispatchEchoes: createCodexDispatchEchoes(),
translator: null
}
}
@@ -1,4 +1,7 @@
import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types'
import type {
AgentJournalItemIdentity,
AgentSessionJournalIdentity
} from '../../shared/agent-session-journal-types'
import { randomUUID } from 'node:crypto'
import { cancelProcessAcquisition } from '../../shared/child-process/cancel-process-acquisition'
import type {
@@ -6,6 +9,7 @@ import type {
openCodexAppServerConnection
} from './codex-app-server-connection'
import { CodexAcquisitionWindow } from './codex-structured-acquisition-window'
import type { CodexDispatchEchoes } from './codex-structured-dispatch-echo'
import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire'
import type { CodexBackgroundTaskTracker } from './codex-background-task-tracker'
import type { CodexJournalTranslator } from './codex-structured-journal-translation'
@@ -58,6 +62,12 @@ export type CodexStructuredSessionAdapterDeps = {
sessionId: string,
state: AgentSessionBackgroundTaskState | null
) => void
/** Identity for a send admitted earlier, once Codex echoes the user message. */
onDispatchSettledLate?: (input: {
sessionId: string
clientMessageId: string
providerIdentity: AgentJournalItemIdentity
}) => void
openConnection?: typeof openCodexAppServerConnection
readProcessStartTime?: (pid: number) => Promise<number | null>
mintLinkId?: () => string
@@ -87,7 +97,8 @@ export type CodexSession = {
prompts: CodexAcquisitionWindow['prompts']
options: Map<string, string>
reportedOptions: { model?: string; effort?: string }
turnIdWaiters: ((turnId: string) => void)[]
/** Sends whose identity is still to be settled by the provider echo. */
dispatchEchoes: CodexDispatchEchoes
translator: CodexJournalTranslator | null
/** Ephemeral roster behind the background-tasks strip; never durable state. */
backgroundTasks: CodexBackgroundTaskTracker
+24 -51
View File
@@ -6,20 +6,14 @@ import {
type CodexAppServerConnection
} from './codex-app-server-connection'
import { isCodexAppServerUnsupportedError } from './codex-app-server-session'
import { readCodexTurnId } from './codex-structured-thread-facts'
import { DISPATCH_DOUBT_CODEX_TURN_UNNAMED } from '../native-chat/agent-session-journal/journal-dispatch-doubt-reasons'
import type { CodexDispatchEchoes } from './codex-structured-dispatch-echo'
// Starting a Codex turn and learning its id, which are not the same event:
// `turn/start` returns the id on newer builds and acks before it exists on
// older ones, where it arrives as a `turn/started` notification instead.
/** Codex records the user message first in a turn, so the submission Orca just
* accepted is ordinal 0 of `(threadId, turnId)`. */
export const CODEX_USER_MESSAGE_ORDINAL = 0
/** Past this the turn is real but unnameable, which the journal renders as
* delivery unconfirmed rather than failure. */
const TURN_ID_WAIT_MS = 10_000
// Writing a Codex turn and learning which message landed where, which are not
// the same event. `turn/start` answers as soon as Codex owns the message, but a
// message issued while a turn is running is COALESCED into that turn: the same
// turn id comes back, no second `turn/started` fires, and the user message is
// echoed only when the running turn reaches it. So the response proves
// admission and nothing about identity, which the echo settles later.
/** Keys Codex accepts as per-turn overrides. An unlisted key would otherwise
* become an arbitrary client-controlled `turn/start` parameter. */
@@ -36,14 +30,12 @@ export function isCodexTurnOptionKey(key: string): boolean {
return CODEX_TURN_OPTION_KEYS.has(key)
}
/** The session state one turn needs. `turnIdWaiters` is shared with the
* notification handler, which resolves the head of the queue — correct because
* Codex runs one turn per thread, so starts and `turn/started` share an order. */
/** The session state one turn needs. */
export type CodexTurnHost = {
connection: Pick<CodexAppServerConnection, 'request'>
threadId: string
options: Map<string, string>
turnIdWaiters: ((turnId: string) => void)[]
dispatchEchoes: CodexDispatchEchoes
}
function turnInputFor(body: AgentJournalMessageItem): Record<string, unknown>[] {
@@ -61,23 +53,17 @@ function turnInputFor(body: AgentJournalMessageItem): Record<string, unknown>[]
}
/**
* Resolves the turn id, or null when Codex owns a turn it never named. Throws
* only for outcomes the wire must not read as acceptance.
* Hands one submission to Codex. Resolves when Codex has taken it; throws only
* for outcomes the wire must not read as acceptance.
*/
export async function startCodexTurn(
host: CodexTurnHost,
input: { clientMessageId: string; body: AgentJournalMessageItem; timeoutMs?: number }
): Promise<string | null> {
// Registered BEFORE the call: on builds that ack first, `turn/started` can
// land while the response is still in flight.
let notified: ((turnId: string) => void) | null = null
const fromNotification = new Promise<string | null>((resolve) => {
notified = resolve
host.turnIdWaiters.push(resolve)
setTimeout(() => resolve(null), TURN_ID_WAIT_MS).unref?.()
})
): Promise<void> {
// Armed before the write: the echo can land while the response is in flight.
host.dispatchEchoes.arm(input.clientMessageId)
try {
const started = await host.connection.request(
await host.connection.request(
'turn/start',
{
threadId: host.threadId,
@@ -87,43 +73,30 @@ export async function startCodexTurn(
},
{ timeoutMs: input.timeoutMs }
)
return readCodexTurnId(started) ?? (await fromNotification)
} finally {
const index = notified ? host.turnIdWaiters.indexOf(notified) : -1
if (index !== -1) {
host.turnIdWaiters.splice(index, 1)
}
} catch (error) {
host.dispatchEchoes.disarm(input.clientMessageId)
throw error
}
}
/**
* One submission's outcome as the wire must read it: accepted names the turn,
* rejected is Codex answering and declining, and unknown covers a turn that is
* real but unnameable — never a failure the user is told their message hit.
* One submission's outcome as the wire must read it: admitted means Codex owns
* the message and its identity settles on the echo, rejected is Codex answering
* and declining. Elapsed time is never evidence here, because the wait a
* coalesced send would face is bounded only by the running turn.
*/
export async function dispatchCodexTurn(
session: CodexTurnHost,
input: { clientMessageId: string; body: AgentJournalMessageItem },
timeoutMs: number | undefined
): Promise<AgentSessionDispatchOutcome> {
let turnId: string | null
try {
turnId = await startCodexTurn(session, { ...input, timeoutMs })
await startCodexTurn(session, { ...input, timeoutMs })
} catch (error) {
if (isCodexAppServerRequestError(error) || isCodexAppServerUnsupportedError(error)) {
return { state: 'rejected', reason: (error as Error).message }
}
throw error
}
return turnId === null
? { state: 'unknown', reason: DISPATCH_DOUBT_CODEX_TURN_UNNAMED }
: {
state: 'accepted',
providerIdentity: {
provider: 'codex',
threadId: session.threadId,
turnId,
ordinal: CODEX_USER_MESSAGE_ORDINAL
}
}
return { state: 'admitted' }
}
+4
View File
@@ -3,6 +3,10 @@ import {
digestPayload
} from '../native-chat/agent-session-journal/journal-payload-bounds'
/** Codex records the user message first in a turn, so a restored submission is
* ordinal 0 of `(threadId, turnId)`. */
export const CODEX_USER_MESSAGE_ORDINAL = 0
/** Maximum forgotten turn keys retained for late-frame reconciliation. */
export const MAX_CODEX_TURN_ORDINAL_ENTRIES = 256
export const MAX_CODEX_TURN_ORDINAL_BYTES = 512 * 1024
@@ -25,13 +25,6 @@ export const DISPATCH_DOUBT_PERSISTENCE_FAILED = 'dispatch_result_persistence_fa
/** A retry was durably armed but had not yet recorded its dispatch outcome. */
export const DISPATCH_DOUBT_RETRY_IN_PROGRESS = 'dispatch_retry_in_progress'
/** Codex owns a turn it started but did not name, because its turn-start still
* settles on a deadline. Delete this once Codex settles on the app-server's
* turn-start response instead; until then this reason is never re-delivered,
* which is what the allowlist below already does by omitting it. */
export const DISPATCH_DOUBT_CODEX_TURN_UNNAMED =
'codex app-server started a turn it did not name in time'
/** The transport refused the frame; the underlying error follows the colon. */
export const DISPATCH_DOUBT_WRITE_FAILED = 'provider_write_failed'
@@ -1,6 +1,7 @@
// What one `agentSession.send` writes, and when a user's Retry is allowed to
// put the same message on the wire a second time.
import { DISPATCH_DOUBT_PROVIDER_EXITED } from '../agent-session-journal/journal-dispatch-doubt-reasons'
import { beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
import type { StructuredAgentSessionHost } from './structured-agent-session-host'
@@ -137,13 +138,15 @@ describe('send', () => {
expect(dispatch).toHaveBeenCalledTimes(2)
})
it('refuses to redeliver a retry for a turn the provider already owns', async () => {
it('refuses to redeliver a retry for a message the provider may already hold', async () => {
await attach()
// A dead child ends the wait without proving non-delivery: the message was
// already written to that child's stdin.
dispatch.mockImplementationOnce(async () => ({
state: 'unknown' as const,
reason: 'codex app-server started a turn it did not name in time'
reason: DISPATCH_DOUBT_PROVIDER_EXITED
}))
const body = hostTestMessage('a turn codex owns but did not name')
const body = hostTestMessage('a message the provider may already hold')
const params = { envelope: envelope('agentSession.send', { body }), body }
const first = await host.send(CALLER, params)
@@ -151,8 +154,8 @@ describe('send', () => {
ok: true,
value: { submission: { dispatchState: 'unknown' } }
})
// The turn is running; a second delivery would be a duplicate, so Retry
// replays the recorded outcome instead of re-sending.
// The reason is not on the fail-closed allowlist, so Retry replays the
// recorded outcome instead of putting the message on the wire again.
await expect(host.send(CALLER, { ...params, retryUnknown: true })).resolves.toMatchObject({
ok: true,
value: { submission: { dispatchState: 'unknown' } }
@@ -42,8 +42,8 @@ describe('structured mailbox pointer host', () => {
// The defect this pins: a running turn is announced by ONE lifecycle item, and settlement
// tombstones it rather than rewriting it. A long tool-calling turn pushes that item arbitrarily
// far from the tail, so any page-sized read reports a busy worker as idle — and the pointer is
// then delivered mid-turn, which Codex answers with `turn already running` and Claude settles
// `unknown` while the message is really queued.
// then delivered mid-turn, which Codex coalesces into the running turn and Claude queues behind
// it -- either way folded into work already in flight rather than read as a new instruction.
const items = [runningTurn(), ...transcript(500)]
hostRef.current = { journalSnapshot: () => ({ items }) }
expect(createStructuredMailboxPointerHost().readGateFacts('s1')).toEqual({
@@ -81,11 +81,14 @@ export function structuredSessionGateFacts(
* Decide whether the nudge may be sent right now.
*
* Mid-turn delivery is refused for both providers rather than delegated to
* them: Codex answers a mid-turn `turn/start` with `turn already running`, and
* Claude accepts the frame but cannot acknowledge it inside the dispatch ack
* window, settling `unknown` while the message is really queued. Waiting for
* the turn to settle is the one contract that holds for both, and it preserves
* orchestration's existing idle-edge-only delivery policy.
* them. Neither refuses the frame: Codex COALESCES a mid-turn `turn/start` into
* the running turn -- measured on codex-cli 0.147.0, 0.150.1 and 0.153.4, none
* of which refuse it and none of which fire a second `turn/started` -- and
* Claude queues it behind the turn. Both therefore
* fold the nudge into work already in flight, where it reads as part of the
* running turn rather than a new instruction. Waiting for the turn to settle is
* the one contract that holds for both, and it preserves orchestration's
* existing idle-edge-only delivery policy.
*/
export function decideStructuredPointerDelivery(input: {
refusal: AgentSessionPtyWriteRefusal
@@ -49,8 +49,8 @@ export function listAddressableStructuredWorkers(): OrchestrationAddressableAgen
* A structured worker's agent status, in the vocabulary `@idle` already matches on.
*
* Null when the session cannot be read: unknown must not read as idle, or a broadcast to `@idle`
* would wake a worker mid-turn — which Codex answers with `turn already running` and Claude queues
* behind the running turn.
* would wake a worker mid-turn — which Codex coalesces into the running turn and Claude queues
* behind it.
*/
export function structuredWorkerAgentStatus(sessionId: string): string | null {
const facts = readStructuredSessionGateFacts(sessionId)
@@ -7,6 +7,7 @@
// reads is module-level for the same reason the registry is — the runtime
// service is already far past its size budget.
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
import { existsSync } from 'node:fs'
import { join } from 'node:path'
import type { AgentSessionRecord } from '../../shared/agent-session-record'
@@ -218,6 +219,18 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
try {
let host: StructuredAgentSessionHost | null = null
let recoveryChain = Promise.resolve()
const onDispatchSettledLate = (settlement: {
sessionId: string
clientMessageId: string
providerIdentity: AgentJournalItemIdentity
}): void => {
void host?.settleLateDispatch(settlement).catch((error) =>
deps.onError?.({
scope: `structured-agent-session-late-settlement:${settlement.sessionId}`,
error
})
)
}
const codex = new CodexStructuredSessionAdapter({
resolveLaunch: createCodexStructuredLaunchResolver({
store,
@@ -229,6 +242,7 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
...(deps.readProcessStartTime ? { readProcessStartTime: deps.readProcessStartTime } : {}),
onBackgroundTasksChanged: (sessionId, state) =>
host?.publishBackgroundTaskState(sessionId, state),
onDispatchSettledLate,
onEvent: (event) => {
if (event.type !== 'ended' || !('cause' in event) || event.cause !== 'unexpected-exit') {
return
@@ -270,14 +284,7 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
},
onBackgroundTasksChanged: (sessionId, state) =>
host?.publishBackgroundTaskState(sessionId, state),
onDispatchSettledLate: (settlement) => {
void host?.settleLateDispatch(settlement).catch((error) =>
deps.onError?.({
scope: `structured-agent-session-late-settlement:${settlement.sessionId}`,
error
})
)
},
onDispatchSettledLate,
...(deps.openClaudeConnection ? { openClaudeConnection: deps.openClaudeConnection } : {}),
...(deps.readProcessStartTime ? { readProcessStartTime: deps.readProcessStartTime } : {})
})