diff --git a/src/main/codex/codex-conversation-name-generation.test.ts b/src/main/codex/codex-conversation-name-generation.test.ts index 31a417030f9..ae19e2312d9 100644 --- a/src/main/codex/codex-conversation-name-generation.test.ts +++ b/src/main/codex/codex-conversation-name-generation.test.ts @@ -72,10 +72,10 @@ function run( // `null` leaves the turn unanswered entirely; 'DECLINE' completes it with // no message, which is a different fact the collector must distinguish. if (answer !== null && answer !== 'DECLINE') { - collector?.handle('item/completed', { item: { type: 'agentMessage', text: answer } }) + collector?.handle('item/completed', { item: { type: 'agentMessage', text: answer } }, true) } if (answer !== null) { - collector?.handle('turn/completed', {}) + collector?.handle('turn/completed', {}, true) } }) }) @@ -168,7 +168,7 @@ describe('createCodexNamingTurnCollector', () => { it('reports a completed turn that said nothing as a DECLINE', async () => { const collector = createCodexNamingTurnCollector(60_000) - collector.handle('turn/completed', {}) + collector.handle('turn/completed', {}, true) // A model that completed and said nothing has answered. Classifying this as // a host failure would re-ask, and pay, on every future acquisition. @@ -178,10 +178,12 @@ describe('createCodexNamingTurnCollector', () => { it('reports the message a completed turn produced', async () => { const collector = createCodexNamingTurnCollector(60_000) - collector.handle('item/completed', { - item: { type: 'agentMessage', text: '{"title":"Fix probe"}' } - }) - collector.handle('turn/completed', {}) + collector.handle( + 'item/completed', + { item: { type: 'agentMessage', text: '{"title":"Fix probe"}' } }, + true + ) + collector.handle('turn/completed', {}, true) await expect(collector.answer).resolves.toEqual({ outcome: 'answered', @@ -194,7 +196,7 @@ describe('createCodexNamingTurnCollector', () => { // There is no `turn/failed` notification; a rate-limited or rejected turn // arrives as `error`. Without it this would hold for the whole timeout. - collector.handle('error', { message: 'rate limit exceeded' }) + collector.handle('error', { message: 'rate limit exceeded' }, true) await expect(collector.answer).resolves.toEqual({ outcome: 'failed' }) }) @@ -202,16 +204,33 @@ describe('createCodexNamingTurnCollector', () => { it('reports a failure even when the turn had already said something', async () => { const collector = createCodexNamingTurnCollector(60_000) - collector.handle('item/completed', { - item: { type: 'agentMessage', text: '{"title":"Fix probe"}' } - }) - collector.handle('error', { message: 'stream closed' }) + collector.handle( + 'item/completed', + { item: { type: 'agentMessage', text: '{"title":"Fix probe"}' } }, + true + ) + collector.handle('error', { message: 'stream closed' }, true) // The host is why there is no title; a partial answer does not make it a // decline the conversation should be marked for. await expect(collector.answer).resolves.toEqual({ outcome: 'failed' }) }) + it.each([ + ['a bare completion', 'turn/completed', {}], + ['an agent message', 'item/completed', { item: { type: 'agentMessage', text: 'hi' } }], + ['a terminal error', 'error', { message: 'rate limit exceeded' }] + ])('ignores %s that could not be attributed to the naming thread', async (_l, method, params) => { + const collector = createCodexNamingTurnCollector(1) + + collector.handle(method, params, false) + + // Inside the `thread/start` window the broad divert rule can catch a + // SUB-AGENT frame. Settling on a bare completion would report a DECLINE, + // which is durable — the conversation could never be named again. + await expect(collector.answer).resolves.toEqual({ outcome: 'timed-out' }) + }) + it('reports a turn that never answered as TIMED OUT', async () => { const collector = createCodexNamingTurnCollector(1) @@ -227,7 +246,7 @@ describe('createCodexNamingTurnCollector', () => { const collector = createCodexNamingTurnCollector(60_000) expect(vi.getTimerCount()).toBe(1) - collector.handle(method, params) + collector.handle(method, params, true) await collector.answer // A pending timer holds the collector's closure for the whole timeout diff --git a/src/main/codex/codex-conversation-name-generation.ts b/src/main/codex/codex-conversation-name-generation.ts index 57f35c56cc0..d38a813ec45 100644 --- a/src/main/codex/codex-conversation-name-generation.ts +++ b/src/main/codex/codex-conversation-name-generation.ts @@ -98,10 +98,12 @@ export type CodexNamingState = { * * That window is one `thread/start` round trip on a healthy app-server, but it is * bounded by the request deadline, not instantaneous — a hung server stretches - * it. A sub-agent frame arriving inside it is still diverted and can still be - * recorded as the naming answer; that costs one wasted naming attempt rather - * than a permanent forfeiture, because unparseable output classifies as - * unsettled. + * it. A sub-agent frame arriving inside it is diverted from the journal and then + * DROPPED: the collector acts only on a frame matched to the retained throwaway + * id, so nothing this rule alone caught can become the naming answer or settle + * the turn. The cost is the sub-agent's own rows missing from the transcript for + * that window — never a forfeited name, which a bare `turn/completed` diverted + * here would otherwise cause by settling as a decline. * * Exact matching means leak protection now RELIES ON ID STABILITY: an app-server * that tagged naming frames with any id other than the one `thread/start` @@ -162,7 +164,9 @@ export type CodexNamingTurnResult = /** Collects one ephemeral naming turn's frames and reports how it ended. */ export type CodexNamingTurnCollector = { - handle: (method: string, params: unknown) => void + /** `attributed` is false for a frame only the broad pre-id divert rule caught; + * it is kept out of the journal but may never speak for this turn. */ + handle: (method: string, params: unknown, attributed: boolean) => void answer: Promise /** For a turn abandoned before `answer` is awaited, which would otherwise hold * the timer — and the closure behind it — for the whole timeout. */ @@ -218,7 +222,13 @@ export function createCodexNamingTurnCollector(timeoutMs: number): CodexNamingTu }) let latest: string | null = null return { - handle: (method, params) => { + handle: (method, params, attributed) => { + // Only a frame matched to the retained throwaway id may speak for this + // turn. A bare `turn/completed` from a sub-agent caught by the broad pre-id + // rule would otherwise settle as a decline and durably forfeit naming. + if (!attributed) { + return + } if (method === 'item/completed') { latest = agentMessageText(params) ?? latest } else if (method === 'turn/completed') { diff --git a/src/main/codex/codex-structured-session-adapter.ts b/src/main/codex/codex-structured-session-adapter.ts index 847dbf578e1..17e709f7a66 100644 --- a/src/main/codex/codex-structured-session-adapter.ts +++ b/src/main/codex/codex-structured-session-adapter.ts @@ -139,8 +139,10 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap // Before anything can journal it: a naming turn runs on a throwaway thread // over this same connection, and the item translator journals items from ANY // thread. Routed here, its prompt and its JSON answer never reach the chat. - if (isCodexNamingFrame(session, readCodexThreadId(params))) { - session.naming?.handle(method, params) + const frameThreadId = readCodexThreadId(params) + if (isCodexNamingFrame(session, frameThreadId)) { + // Diverted either way; only the exact-id half may settle the naming turn. + session.naming?.handle(method, params, isCodexNamingThread(session, frameThreadId)) return { accepted: true } } captureCodexConversationName(sessionId, session, method, params, this.deps) @@ -194,8 +196,9 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap private handleUnhandledFrame(sessionId: string, kind: string, params: unknown): void { const session = this.sessions.get(sessionId) - if (session && isCodexNamingFrame(session, readCodexThreadId(params))) { - session.naming?.handle(kind, params) + const frameThreadId = readCodexThreadId(params) + if (session && isCodexNamingFrame(session, frameThreadId)) { + session.naming?.handle(kind, params, isCodexNamingThread(session, frameThreadId)) return } deliverCodexUnhandledFrame(sessionId, session, kind, params, (current, event) => diff --git a/src/main/codex/codex-structured-session-close.ts b/src/main/codex/codex-structured-session-close.ts index bc33422d366..a34ce19026f 100644 --- a/src/main/codex/codex-structured-session-close.ts +++ b/src/main/codex/codex-structured-session-close.ts @@ -46,7 +46,8 @@ export function handleCodexSessionExit(input: { // A naming turn in flight otherwise holds its collector until the 60s deadline // and then runs its cleanup against a dead connection. Settling it as a host // failure — not a decline — leaves the conversation askable on reacquisition. - session.naming?.handle('error', {}) + // Attributed: the session owning the turn is what died, not a foreign thread. + session.naming?.handle('error', {}, true) session.naming = null session.unbindReadingControl?.() input.onEvent?.(event) diff --git a/src/main/codex/codex-structured-session-conversation-name.test.ts b/src/main/codex/codex-structured-session-conversation-name.test.ts index 205a9133f5f..d2ff0ca4498 100644 --- a/src/main/codex/codex-structured-session-conversation-name.test.ts +++ b/src/main/codex/codex-structured-session-conversation-name.test.ts @@ -132,6 +132,9 @@ function namingCodex( /** Never answers the ephemeral `thread/start`, so naming is in flight with * NO naming thread id known: the window the broad frame rule covers. */ hangNamingThreadStart?: boolean + /** Holds the ephemeral `thread/start` open until the test releases it, so a + * frame can arrive inside that window and the flow still runs to the end. */ + holdNamingThreadStart?: boolean /** Completes the naming turn having said nothing: a genuine model decline, * which is a different fact from prose that ignored the schema. */ declineNamingTurn?: boolean @@ -140,6 +143,10 @@ function namingCodex( const connections: FakeConnection[] = [] const calls: { method: string; params: Record }[] = [] const replies: { id: number | string; result?: unknown; code?: number; message?: string }[] = [] + let releaseNamingThreadStart = (): void => {} + const namingThreadStartGate = new Promise((resolve) => { + releaseNamingThreadStart = resolve + }) const openConnection = (async ( _launch: CodexAppServerLaunch, handlers: CodexAppServerConnectionHandlers = {} @@ -154,24 +161,14 @@ function namingCodex( if (options.hangNamingThreadStart) { return await new Promise(() => {}) } + if (options.holdNamingThreadStart) { + await namingThreadStartGate + } // Leaves the naming turn in flight: the thread id is known, but nothing // ever settles the collector, which is the window sub-agents run in. if (options.hangNamingTurn) { return { thread: { id: NAMING_THREAD, ephemeral: true } } } - // The naming turn's frames arrive on this same connection. - queueMicrotask(() => { - if (!options.declineNamingTurn) { - handlers.onNotification?.('item/completed', { - threadId: NAMING_THREAD, - item: { - type: 'agentMessage', - text: options.answer ?? '{"title":"Fix lease probe"}' - } - }) - } - handlers.onNotification?.('turn/completed', { threadId: NAMING_THREAD }) - }) return { thread: { id: NAMING_THREAD } } } if (method === 'thread/start') { @@ -186,6 +183,23 @@ function namingCodex( } } if (method === 'turn/start') { + // The naming turn's frames arrive on this same connection, and only + // once the turn exists — the app-server cannot emit for a turn that + // was never started, which is what makes them attributable. + if (params?.threadId === NAMING_THREAD) { + queueMicrotask(() => { + if (!options.declineNamingTurn) { + handlers.onNotification?.('item/completed', { + threadId: NAMING_THREAD, + item: { + type: 'agentMessage', + text: options.answer ?? '{"title":"Fix lease probe"}' + } + }) + } + handlers.onNotification?.('turn/completed', { threadId: NAMING_THREAD }) + }) + } return { turn: { id: 'turn-1' } } } return {} @@ -202,7 +216,7 @@ function namingCodex( connections.push(connection) return connection }) as typeof openCodexAppServerConnection - return { connections, openConnection, calls, replies } + return { connections, openConnection, calls, replies, releaseNamingThreadStart } } const NAMING_THREAD = 'thread-naming' @@ -566,6 +580,31 @@ describe('Codex sub-agent threads survive the naming window', () => { ]) }) + it('does not let a foreign turn/completed inside the window forfeit naming', async () => { + const codex = namingCodex({ holdNamingThreadStart: true }) + const markNamingAttempted = vi.fn() + const { onConversationName } = await dispatchedAdapter(codex, { markNamingAttempted }) + await settle() + + // A sub-agent's BARE completion, arriving before the throwaway thread has an + // id. Settling on it reports a decline, which is durable: the conversation + // would carry `conversationNamingAttempted` with no name and could never be + // named again. + codex.connections[0]!.handlers.onNotification?.('turn/completed', { + threadId: SUBAGENT_THREAD + }) + await settle() + + codex.releaseNamingThreadStart() + await settle() + + expect(codex.calls.find((call) => call.method === 'thread/name/set')?.params).toEqual({ + threadId: THREAD_ID, + name: 'Fix lease probe' + }) + expect(onConversationName).toHaveBeenCalledWith(SESSION, 'Fix lease probe') + }) + it('still keeps the naming thread out once its id is known', async () => { const codex = namingCodex({ hangNamingTurn: true }) const { events } = await dispatchedAdapter(codex)