feat(native-chat): publish the turn a withdrawn Codex send was answered into

A send Codex answered into a turn that then ended without taking it is
rejected by that turn's end. The rejection row now names that turn's
record, and the submission carries it as answeredInTurnItemId, so a
client can tell which turn the send belonged to without guessing from
journal order or clocks. Absent on every other send.
This commit is contained in:
Brennan Benson
2026-10-05 08:50:55 -07:00
parent bd3789ef24
commit 93749070e4
17 changed files with 319 additions and 26 deletions
@@ -132,7 +132,7 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
}
if (event.type === 'notification') {
// Only an admitted turn end settles; a refused one settles on the retry that lands.
settleCodexSendsInEndedTurn(session, event.method, event.params, (settlement) =>
settleCodexSendsInEndedTurn(session, event, (settlement) =>
this.deps.onDispatchSettledLate?.({ sessionId: event.sessionId, ...settlement })
)
// After the admission check, so a refused frame is observed by the child records
@@ -84,7 +84,10 @@ export type CodexStructuredSessionAdapterDeps = {
onDispatchSettledLate?: (
input: { sessionId: string; clientMessageId: string } & (
| { providerIdentity: AgentJournalItemIdentity }
| ({ state: 'rejected' } & AgentJournalDispatchRejection)
| ({
state: 'rejected'
answeredInTurn: AgentJournalItemIdentity
} & AgentJournalDispatchRejection)
)
) => void
/** Codex reported its thread not running with no turn open: a send whose
@@ -1,5 +1,9 @@
import { describe, expect, it, vi } from 'vitest'
import type { AgentJournalItemBody } from '../../shared/agent-session-journal-types'
import type {
AgentJournalItemBody,
AgentJournalItemIdentity
} from '../../shared/agent-session-journal-types'
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import { classifyDispatchRejection } from '../../shared/structured-agent-session-dispatch-rejection'
import { createCodexDispatchEchoes } from './codex-structured-dispatch-echo'
import {
@@ -20,11 +24,17 @@ async function turnEndRig() {
Object.assign(codex.routes, turns.routes)
const settlements: LateSettlement[] = []
const bodies: AgentJournalItemBody[] = []
const turnRecords: AgentJournalItemIdentity[] = []
const adapter = await acquiredCodexAdapter({
codex,
settlements,
sink: {
appendItem: (_identity, body) => bodies.push(body),
appendItem: (identity, body) => {
bodies.push(body)
if (body.kind === 'turn') {
turnRecords.push(identity)
}
},
appendTombstone: () => {},
publish: () => {}
}
@@ -46,7 +56,18 @@ async function turnEndRig() {
const settledIds = () => settlements.map(({ clientMessageId }) => clientMessageId)
const categoryOf = (settlement: LateSettlement | undefined) =>
settlement && 'state' in settlement ? classifyDispatchRejection(settlement).category : null
return { codex, turns, adapter, send, sendAndOpen, settlements, settledIds, categoryOf, bodies }
return {
codex,
turns,
adapter,
send,
sendAndOpen,
settlements,
settledIds,
categoryOf,
bodies,
turnRecords
}
}
describe('a Codex send its turn ended without echoing', () => {
@@ -286,6 +307,64 @@ describe('a Codex send its turn ended without echoing', () => {
})
})
describe('the turn a withdrawn Codex send names', () => {
it('is the record of the turn it opened, when that turn is interrupted', async () => {
const rig = await turnEndRig()
await rig.sendAndOpen('client-1')
rig.turns.end('interrupted')
expect(rig.turnRecords.length).toBeGreaterThan(0)
expect(new Set(rig.turnRecords.map(agentJournalItemKey)).size).toBe(1)
expect(rig.settlements).toEqual([
expect.objectContaining({ clientMessageId: 'client-1', answeredInTurn: rig.turnRecords[0] })
])
})
it('is the running turn for a send steered into it', async () => {
const rig = await turnEndRig()
await rig.sendAndOpen('client-1')
await rig.send('client-2')
rig.turns.end('interrupted')
expect(rig.settlements.map((settlement) => [settlement.clientMessageId, settlement])).toEqual([
['client-1', expect.objectContaining({ answeredInTurn: rig.turnRecords[0] })],
['client-2', expect.objectContaining({ answeredInTurn: rig.turnRecords[0] })]
])
})
it('is the ended turn whose end was read before the answer', async () => {
const rig = await turnEndRig()
const release = rig.turns.holdNextAnswer()
const sending = rig.send('client-1')
await vi.waitFor(() => expect(rig.turns.turnId).toBe('turn-1'))
rig.turns.start()
rig.turns.end('interrupted')
release()
await expect(sending).resolves.toMatchObject({
state: 'rejected',
answeredInTurn: rig.turnRecords[0]
})
})
it('is not named when its turn completes without echoing it: the send stays pending', async () => {
const rig = await turnEndRig()
const release = rig.turns.holdNextAnswer()
const sending = rig.send('client-1')
await vi.waitFor(() => expect(rig.turns.turnId).toBe('turn-1'))
rig.turns.start()
rig.turns.end('completed')
release()
await expect(sending).resolves.toEqual({ state: 'admitted' })
await rig.sendAndOpen('client-2')
rig.turns.end('completed')
expect(rig.settlements).toEqual([])
})
})
describe('a send bound to a turn', () => {
it('dies with the settlement its turn end makes', () => {
const echoes = createCodexDispatchEchoes()
@@ -17,7 +17,9 @@ import {
agentSessionFailureWords,
type AgentJournalDispatchRejection
} from '../../shared/agent-session-failure-words'
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
import type { CodexTurnEnd } from './codex-structured-dispatch-echo'
import { codexTurnLifecycleIdentity } from './codex-structured-journal-translation-turns'
import type { CodexSession } from './codex-structured-session-state'
import {
readCodexThreadId,
@@ -41,6 +43,8 @@ export function codexDispatchRejection(
export type CodexTurnEndSettlement = {
clientMessageId: string
state: 'rejected'
/** The turn Codex answered the send into, whose end settled it. */
answeredInTurn: AgentJournalItemIdentity
} & AgentJournalDispatchRejection
function errorDetail(params: unknown): ProviderDiagnostic | undefined {
@@ -80,19 +84,23 @@ export function codexTurnEndRejection(end: CodexTurnEnd): AgentJournalDispatchRe
/** Settles the sends bound to the turn this admitted notification ended. */
export function settleCodexSendsInEndedTurn(
session: Pick<CodexSession, 'threadId' | 'dispatchEchoes'>,
method: string,
params: unknown,
frame: { sessionId: string; method: string; params: unknown },
settle: (settlement: CodexTurnEndSettlement) => void
): void {
const turnId = readCodexTurnId(params)
const end = readCodexTurnEnd(method, params)
if (!turnId || !end || (readCodexThreadId(params) ?? session.threadId) !== session.threadId) {
const turnId = readCodexTurnId(frame.params)
const end = readCodexTurnEnd(frame.method, frame.params)
if (
!turnId ||
!end ||
(readCodexThreadId(frame.params) ?? session.threadId) !== session.threadId
) {
return
}
const rejection = codexTurnEndRejection(end)
const answeredInTurn = codexTurnLifecycleIdentity(frame.sessionId, turnId)
for (const clientMessageId of session.dispatchEchoes.endTurn(session.threadId, turnId, end)) {
if (rejection) {
settle({ clientMessageId, state: 'rejected', ...rejection })
settle({ clientMessageId, state: 'rejected', answeredInTurn, ...rejection })
}
}
}
+14 -2
View File
@@ -8,6 +8,7 @@ import {
} from './codex-app-server-connection'
import { isCodexAppServerUnsupportedError } from './codex-app-server-session'
import type { CodexDispatchEchoes } from './codex-structured-dispatch-echo'
import { codexTurnLifecycleIdentity } from './codex-structured-journal-translation-turns'
import { readCodexTurnId } from './codex-structured-thread-facts'
import {
codexRunningOrOpeningTurn,
@@ -178,7 +179,12 @@ export async function startCodexTurn(
*/
export async function dispatchCodexTurn(
session: CodexTurnHost,
input: { clientMessageId: string; body: AgentJournalMessageItem; requestedAt?: number },
input: {
sessionId: string
clientMessageId: string
body: AgentJournalMessageItem
requestedAt?: number
},
timeoutMs: number | undefined
): Promise<AgentSessionDispatchOutcome> {
let answer: { turnId: string | null } | false
@@ -208,5 +214,11 @@ export async function dispatchCodexTurn(
? session.dispatchEchoes.bindTurn(input.clientMessageId, session.threadId, answer.turnId)
: null
const rejection = endedFirst ? codexTurnEndRejection(endedFirst) : null
return rejection ? { state: 'rejected', ...rejection } : { state: 'admitted' }
return rejection && answer.turnId
? {
state: 'rejected',
answeredInTurn: codexTurnLifecycleIdentity(input.sessionId, answer.turnId),
...rejection
}
: { state: 'admitted' }
}
@@ -34,6 +34,15 @@ export function applyJournalDispatchRow(
} else {
delete submission.rejection
}
if (
row.state === 'rejected' &&
typeof row.answeredInTurnItemId === 'string' &&
row.answeredInTurnItemId
) {
submission.answeredInTurnItemId = row.answeredInTurnItemId
} else {
delete submission.answeredInTurnItemId
}
submission.resolvedAt = row.state === 'pending' ? null : row.ts
if (row.state === 'pending') {
submission.handedOverAt = row.ts
@@ -110,6 +110,9 @@ export function journalDispatchRowBuilder(
providerItemId,
reason: boundedDispatchReason(input),
...(input.state === 'rejected' ? { rejection: input.rejection } : {}),
...(input.state === 'rejected' && input.answeredInTurn
? { answeredInTurnItemId: agentJournalItemKey(input.answeredInTurn) }
: {}),
...journalRowBase(state().epoch, seq, input.fence, ts),
...(input.recovered ? { recovered: input.recovered } : {}),
...(input.state === 'pending' ? { turnScope: input.turnScope } : {})
@@ -150,6 +150,10 @@ export type JournalDispatchRow = JournalRowBase & {
/** On `rejected`: why, typed. Older readers keep the key and ignore it; a malformed one is
* dropped when read, never the row. */
rejection?: AgentSessionFailureFact
/** On `rejected`: the turn record a Codex send was answered into (`turn/start` or `turn/steer`)
* when that turn's end settled the send. Absent on every other row. Older readers keep the key
* and ignore it. */
answeredInTurnItemId?: string
}
/** An item mutation may name its own producer, because one batch can CREATE
@@ -40,7 +40,10 @@ export type ResolveDispatchInput = {
| { state: 'pending'; turnScope: AgentJournalTurnScope }
/** `reason` is what released clients print, `rejection` what newer ones read: both from
* `agentSessionFailureWords`, never written by hand. */
| ({ state: 'rejected' } & AgentJournalDispatchRejection)
| ({
state: 'rejected'
answeredInTurn?: AgentJournalItemIdentity
} & AgentJournalDispatchRejection)
| { state: 'unknown'; reason?: string | null }
)
@@ -1,6 +1,7 @@
// Each submission carries where the journal wrote its row. A rejected send's own row moves to its
// rejection, so this is the only journal-order record of where it was sent. Derived on every fold: a
// replay of the same rows gives the same value, and a history page carries it.
// replay of the same rows gives the same value, and a history page carries it. A rejection a turn's
// end made also names that turn.
import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
@@ -8,7 +9,10 @@ import { join } from 'node:path'
import { afterEach, describe, expect, it } from 'vitest'
import { agentSessionFailureFact } from '../../../shared/agent-session-failure'
import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words'
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
import {
agentJournalItemKey,
agentJournalSubmissionKey
} from '../../../shared/agent-session-journal-item-key'
import {
AGENT_JOURNAL_THREAD_SCOPE,
type AgentJournalMessageItem,
@@ -69,6 +73,61 @@ async function sendHandedOverThenWithdrawn() {
return { journal, submitted, handedOver, pending, withdrawn }
}
const TURN = {
provider: 'legacy',
agent: 'codex',
sessionId: 'session-1',
recordId: 'turn-lifecycle:turn-1'
} as const
/** A send its turn ended without taking: the rejection names that turn. */
async function sendWithdrawnByItsTurnEnd() {
root = await mkdtemp(join(tmpdir(), 'orca-submission-positions-'))
const journal = await journals.open({ identity: IDENTITY, stateDirectory: root })
await journal.appendSubmission({
clientMessageId: 'send-1',
payloadFingerprint: 'send-1',
body: BODY,
fence: 1
})
await journal.resolveDispatch({
clientMessageId: 'send-1',
state: 'rejected',
...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }),
answeredInTurn: TURN,
fence: 1
})
return journal
}
describe('the turn a rejected submission was answered into', () => {
it('is the turn record its rejection named, on the snapshot, a replay, and the page', async () => {
const journal = await sendWithdrawnByItsTurnEnd()
const turnItemId = agentJournalItemKey(TURN)
expect(journal.submission('send-1')).toMatchObject({
dispatchState: 'rejected',
answeredInTurnItemId: turnItemId
})
expect(readAgentSessionHydrationPage(journal).submissions).toEqual([
expect.objectContaining({ answeredInTurnItemId: turnItemId })
])
await journals.closeAll()
const replayed = await journals.open({ identity: IDENTITY, stateDirectory: root! })
expect(replayed.submission('send-1')).toMatchObject({ answeredInTurnItemId: turnItemId })
})
it('is absent on a take-back that names no turn, as on rows from older hosts', async () => {
const { journal } = await sendHandedOverThenWithdrawn()
expect(journal.submission('send-1')).toMatchObject({ dispatchState: 'rejected' })
expect(journal.submission('send-1')).not.toHaveProperty('answeredInTurnItemId')
expect(readAgentSessionHydrationPage(journal).submissions[0]).not.toHaveProperty(
'answeredInTurnItemId'
)
})
})
describe("a submission's journal position", () => {
it('is its own row, which a take-back does not move', async () => {
const { journal, submitted, handedOver, pending, withdrawn } =
@@ -176,8 +176,12 @@ export type AgentSessionDispatchOutcome =
* anything and never promotes this to `unknown`.
*/
| { state: 'admitted' }
/** Words from `agentSessionFailureWords`, never written by hand. */
| ({ state: 'rejected' } & AgentJournalDispatchRejection)
/** Words from `agentSessionFailureWords`, never written by hand. `answeredInTurn`: the turn the
* provider answered the send into, which ended before the answer was read. */
| ({
state: 'rejected'
answeredInTurn?: AgentJournalItemIdentity
} & AgentJournalDispatchRejection)
/** The call did not settle. Never re-send on the user's behalf. */
| { state: 'unknown'; reason: string }
@@ -12,7 +12,10 @@ export async function settleStructuredAgentSessionLateDispatch(
clientMessageId: string
} & (
| { providerIdentity: AgentJournalItemIdentity }
| ({ state: 'rejected' } & AgentJournalDispatchRejection)
| ({
state: 'rejected'
answeredInTurn?: AgentJournalItemIdentity
} & AgentJournalDispatchRejection)
| { state: 'unknown'; reason: string }
)
): Promise<void> {
@@ -43,6 +46,7 @@ export async function settleStructuredAgentSessionLateDispatch(
state: 'rejected',
reason: input.reason,
rejection: input.rejection,
...(input.answeredInTurn ? { answeredInTurn: input.answeredInTurn } : {}),
fence
}
)
@@ -231,6 +231,7 @@ export async function handOverSubmission(
state: 'rejected',
reason: outcome.reason,
rejection: outcome.rejection,
...(outcome.answeredInTurn ? { answeredInTurn: outcome.answeredInTurn } : {}),
fence: ctx.fence
}
: { clientMessageId, state: 'unknown', reason: outcome.reason, fence: ctx.fence }
@@ -304,6 +304,79 @@ describe('a Codex send its turn ended without taking it', () => {
})
})
describe('the turn a withdrawn Codex send was answered into', () => {
async function answeredInto(clientMessageId: string) {
await host.flushStreamedEvents(SESSION)
const snapshot = await host.journalSnapshot(SESSION)
const page = await host.history({ sessionId: SESSION, direction: 'tail' })
const onPage = page.ok
? page.page.submissions.find((entry) => entry.clientMessageId === clientMessageId)
: undefined
return {
turnRecords: snapshot.items.flatMap((item) =>
item.body.kind === 'turn' ? [item.itemId] : []
),
named: snapshot.submissions.find((entry) => entry.clientMessageId === clientMessageId)
?.answeredInTurnItemId,
onPage: onPage?.answeredInTurnItemId
}
}
it("is that turn's record, for the send that opened it and for one steered into it", async () => {
const opening = await send('look around')
await vi.waitFor(() => expect(answers).toBe(1))
turns.start()
const steered = await send('and check the tests')
await vi.waitFor(() => expect(steers).toBe(1))
await stop('turn-1')
await vi.waitFor(async () =>
expect(verdictOf((await settled()).submissions, steered)).toBe('withdrawn')
)
const { turnRecords } = await answeredInto(opening)
expect(turnRecords).toHaveLength(1)
expect(await answeredInto(opening)).toMatchObject({
named: turnRecords[0],
onPage: turnRecords[0]
})
expect(await answeredInto(steered)).toMatchObject({
named: turnRecords[0],
onPage: turnRecords[0]
})
})
it('is not named on a send the turn took', async () => {
const opening = await send('look around')
await vi.waitFor(() => expect(answers).toBe(1))
turns.start()
turns.echo(opening)
await stop('turn-1')
await vi.waitFor(async () =>
expect(verdictOf((await settled()).submissions, opening)).toBe('accepted')
)
expect(await answeredInto(opening)).toMatchObject({ named: undefined, onPage: undefined })
})
it('is named when the answer is read after that turn ended', async () => {
const release = turns.holdNextAnswer()
const sent = await send('look around')
await vi.waitFor(() => expect(turns.turnId).toBe('turn-1'))
turns.start()
turns.end('interrupted')
release()
await vi.waitFor(async () =>
expect(verdictOf((await settled()).submissions, sent)).toBe('withdrawn')
)
const { turnRecords, named } = await answeredInto(sent)
expect(turnRecords).toHaveLength(1)
expect(named).toBe(turnRecords[0])
})
})
describe('a queued card sent now into the turn a Stop ends', () => {
async function handoffs(messageId: string): Promise<AgentJournalSubmission[]> {
return (await settled()).submissions.filter((entry) => entry.queuedMessageId === messageId)
@@ -319,6 +319,7 @@ export const AgentJournalSubmissionSchema = z.object({
submittedAt: z.number(),
resolvedAt: z.number().nullable(),
submittedSequence: z.number().int().optional(),
answeredInTurnItemId: z.string().min(1).optional(),
recovered: z.literal(true).optional(),
handoverRecorded: z.literal(true).optional(),
handedOverAt: z.number().optional(),
@@ -415,6 +415,10 @@ export type AgentJournalSubmission = {
* never stored. Needed because a rejected send's own row moves to its rejection, which erases
* where it was sent. Absent from hosts that predate it. */
submittedSequence?: number
/** On `rejected`: the item id of the turn record a Codex send was answered into, when that turn
* ended without taking it. Absent on every other send: accepted ones (the echo places them),
* Claude, queued take-backs, restart recovery, and rows from hosts that predate it. */
answeredInTurnItemId?: string
/** Set when crash reconciliation resolved the dispatch, not the provider. A live
* `unknown` is a send still outstanding; a recovered one outlived its writer. */
recovered?: true
@@ -10,8 +10,8 @@ import { createTrackedJournalOpener } from '../../../src/main/native-chat/agent-
import { readAgentSessionHydrationPage } from '../../../src/main/native-chat/agent-session-wire/agent-session-history-page'
import { importReleaseCheckoutModule, materializeReleaseCheckout } from './release-checkout'
// A released client that predates the published submission position: it must fold and draw a page
// carrying it exactly as it draws the same page without it.
// A released client that predates the published submission position and answered turn: it must fold
// and draw a page carrying them exactly as it draws the same page without them.
const BASELINE_REF = 'v1.4.219'
const IDENTITY: AgentSessionJournalIdentity = {
@@ -34,16 +34,19 @@ function releaseExport<T>(module: Record<string, unknown>, name: string): T {
type OldState = { submissions: Record<string, unknown>[]; items: unknown[] }
/** Strips the published position, as an older host's page would carry its submissions. */
/** Strips the published position and answered turn, as an older host's page would carry its
* submissions. */
function withoutPosition(page: AgentSessionHistoryPage): AgentSessionHistoryPage {
return {
...page,
submissions: page.submissions.map(({ submittedSequence: _submitted, ...rest }) => rest)
submissions: page.submissions.map(
({ submittedSequence: _submitted, answeredInTurnItemId: _turn, ...rest }) => rest
)
}
}
// Loads a real release checkout, cold extraction and transforms included.
test('a released client folds and draws a page whose submissions carry their journal position', async () => {
test('a released client folds and draws a page whose submissions carry their journal position and answered turn', async () => {
const directory = mkdtempSync(join(tmpdir(), 'orca-submission-positions-downgrade-'))
const journals = createTrackedJournalOpener()
try {
@@ -64,13 +67,34 @@ test('a released client folds and draws a page whose submissions carry their jou
fence: 1,
recovered: true
})
await journal.appendSubmission({
clientMessageId: 'turn-ended',
payloadFingerprint: 'turn-ended',
body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'turn-ended' }] },
fence: 1
})
await journal.resolveDispatch({
clientMessageId: 'turn-ended',
state: 'rejected',
...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }),
answeredInTurn: {
provider: 'legacy',
agent: 'codex',
sessionId: IDENTITY.sessionId,
recordId: 'turn-lifecycle:turn-1'
},
fence: 1
})
const page = readAgentSessionHydrationPage(journal, 1)
// Anti-vacuous: the page this build publishes does carry the position.
// Anti-vacuous: the page this build publishes does carry both.
expect(page.submissions.find((entry) => entry.clientMessageId === 'taken-back')).toEqual(
expect.objectContaining({
submittedSequence: expect.any(Number)
})
)
expect(page.submissions.find((entry) => entry.clientMessageId === 'turn-ended')).toEqual(
expect.objectContaining({ answeredInTurnItemId: expect.any(String) })
)
const checkout = await materializeReleaseCheckout(BASELINE_REF)
const reducer = await importReleaseCheckoutModule(
@@ -92,7 +116,9 @@ test('a released client folds and draws a page whose submissions carry their jou
const draw = (from: AgentSessionHistoryPage) => {
const state = reduce(empty, { type: 'history-page', page: from })
return {
submissions: state.submissions.map(({ submittedSequence: _s, ...rest }) => rest),
submissions: state.submissions.map(
({ submittedSequence: _s, answeredInTurnItemId: _t, ...rest }) => rest
),
messages: project(state.items, [], state.submissions)
}
}