feat(native-chat): say how a withdrawn Codex send joined the turn it names

The answered turn becomes { turnItemId, via }: 'start' when Codex
answered the send's turn/start with that turn, 'steer' when Orca steered
it into that running turn. A client can then tell the turn's opener from
a send steered into it even when the opener is on an older page.
This commit is contained in:
Brennan Benson
2026-10-05 09:02:33 -07:00
parent 93749070e4
commit ab2826ee18
17 changed files with 171 additions and 86 deletions
@@ -133,9 +133,9 @@ describe('the turn Codex answered a send into but has not opened', () => {
it('is the turn the latest armed send was answered into', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
echoes.arm('client-2')
echoes.bindTurn('client-2', 'thread-1', 'turn-2')
echoes.bindTurn('client-2', 'thread-1', 'turn-2', 'start')
expect(echoes.answeredUnopenedTurn('thread-1', NONE_OPEN)).toBe('turn-2')
})
@@ -143,7 +143,7 @@ describe('the turn Codex answered a send into but has not opened', () => {
it('is none once Codex opened that turn', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
expect(echoes.answeredUnopenedTurn('thread-1', new Set(['turn-1']))).toBeNull()
})
@@ -151,7 +151,7 @@ describe('the turn Codex answered a send into but has not opened', () => {
it('is none once that turn ended, even with its send still armed for an echo', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
// A completed end leaves its unechoed send armed.
expect(echoes.endTurn('thread-1', 'turn-1', { status: 'completed' })).toEqual([])
@@ -161,9 +161,9 @@ describe('the turn Codex answered a send into but has not opened', () => {
it('skips a turn a wait left unopened, and still names an earlier one', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
echoes.arm('client-2')
echoes.bindTurn('client-2', 'thread-1', 'turn-2')
echoes.bindTurn('client-2', 'thread-1', 'turn-2', 'start')
echoes.leftUnopened('thread-1', 'turn-2')
@@ -174,7 +174,7 @@ describe('the turn Codex answered a send into but has not opened', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.arm('client-2')
echoes.bindTurn('client-2', 'thread-2', 'turn-2')
echoes.bindTurn('client-2', 'thread-2', 'turn-2', 'start')
expect(echoes.answeredUnopenedTurn('thread-1', NONE_OPEN)).toBeNull()
})
@@ -1,5 +1,8 @@
import type { ProviderDiagnostic } from '../../shared/agent-session-failure'
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
import type {
AgentJournalItemIdentity,
AgentJournalTurnJoin
} from '../../shared/agent-session-journal-types'
/** Sends awaiting their echo. One bound to a turn that ended without taking it settles from that
* end; any other whose echo never arrives is retired by the journal's recovery on exit. */
@@ -35,10 +38,15 @@ export type CodexDispatchEchoes = {
/** Drops an armed send whose write never reached the provider. */
disarm: (clientMessageId: string) => void
/**
* Binds a send to the turn Codex answered it into. Returns that turn's end when the answer is
* read after it; a send that end settles is no longer armed.
* Binds a send to the turn Codex answered it into, and how it joined that turn. Returns that
* turn's end when the answer is read after it; a send that end settles is no longer armed.
*/
bindTurn: (clientMessageId: string, threadId: string, turnId: string) => CodexTurnEnd | null
bindTurn: (
clientMessageId: string,
threadId: string,
turnId: string,
via: AgentJournalTurnJoin
) => CodexTurnEnd | null
/** The turn the latest armed send was answered into that is neither in `openTurnIds`, ended, nor
* left unopened through a wait: one Codex has picked for the send but not opened. */
answeredUnopenedTurn: (threadId: string, openTurnIds: ReadonlySet<string>) => string | null
@@ -49,7 +57,11 @@ export type CodexDispatchEchoes = {
* completed, which echoes its pending input first, so one it never echoed waits for recovery. An
* interrupt withdraws an un-echoed send, steered or the turn's own input: neither reached history.
*/
endTurn: (threadId: string, turnId: string, end: CodexTurnEnd) => string[]
endTurn: (
threadId: string,
turnId: string,
end: CodexTurnEnd
) => { clientMessageId: string; via: AgentJournalTurnJoin }[]
/** Submission origin for this exact send, retained until its echo settles it. */
requestOrigin: (clientMessageId: string) => CodexDispatchRequestOrigin | null
/** Highest causal sequence assigned to a dispatch in this session. */
@@ -61,7 +73,11 @@ export type CodexDispatchEchoes = {
export function createCodexDispatchEchoes(): CodexDispatchEchoes {
const armed = new Map<
string,
{ requestedAt: number | null; sequence: number; turn?: { threadId: string; turnId: string } }
{
requestedAt: number | null
sequence: number
turn?: { threadId: string; turnId: string; via: AgentJournalTurnJoin }
}
>()
const endedTurns = new Map<string, CodexTurnEnd>()
const unopenedTurns = new Set<string>()
@@ -85,12 +101,12 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes {
},
settle: (clientMessageId) => armed.delete(clientMessageId),
disarm: (clientMessageId) => void armed.delete(clientMessageId),
bindTurn: (clientMessageId, threadId, turnId) => {
bindTurn: (clientMessageId, threadId, turnId, via) => {
const entry = armed.get(clientMessageId)
if (!entry) {
return null
}
entry.turn = { threadId, turnId }
entry.turn = { threadId, turnId, via }
const end = endedTurns.get(turnKey(threadId, turnId)) ?? null
if (end && settles(end)) {
armed.delete(clientMessageId)
@@ -132,10 +148,10 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes {
}
const settled = [...armed].flatMap(([clientMessageId, entry]) =>
entry.turn && turnKey(entry.turn.threadId, entry.turn.turnId) === turn
? [clientMessageId]
? [{ clientMessageId, via: entry.turn.via }]
: []
)
for (const clientMessageId of settled) {
for (const { clientMessageId } of settled) {
armed.delete(clientMessageId)
}
return settled
@@ -1,4 +1,5 @@
import type {
AgentJournalAnsweredTurnIdentity,
AgentJournalItemIdentity,
AgentSessionJournalIdentity
} from '../../shared/agent-session-journal-types'
@@ -86,7 +87,7 @@ export type CodexStructuredSessionAdapterDeps = {
| { providerIdentity: AgentJournalItemIdentity }
| ({
state: 'rejected'
answeredInTurn: AgentJournalItemIdentity
answeredInTurn: AgentJournalAnsweredTurnIdentity
} & AgentJournalDispatchRejection)
)
) => void
@@ -317,7 +317,10 @@ describe('the turn a withdrawn Codex send names', () => {
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] })
expect.objectContaining({
clientMessageId: 'client-1',
answeredInTurn: { turn: rig.turnRecords[0], via: 'start' }
})
])
})
@@ -329,8 +332,14 @@ describe('the turn a withdrawn Codex send names', () => {
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] })]
[
'client-1',
expect.objectContaining({ answeredInTurn: { turn: rig.turnRecords[0], via: 'start' } })
],
[
'client-2',
expect.objectContaining({ answeredInTurn: { turn: rig.turnRecords[0], via: 'steer' } })
]
])
})
@@ -345,7 +354,7 @@ describe('the turn a withdrawn Codex send names', () => {
await expect(sending).resolves.toMatchObject({
state: 'rejected',
answeredInTurn: rig.turnRecords[0]
answeredInTurn: { turn: rig.turnRecords[0], via: 'start' }
})
})
@@ -369,9 +378,11 @@ describe('a send bound to a turn', () => {
it('dies with the settlement its turn end makes', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
expect(echoes.endTurn('thread-1', 'turn-1', { status: 'interrupted' })).toEqual(['client-1'])
expect(echoes.endTurn('thread-1', 'turn-1', { status: 'interrupted' })).toEqual([
{ clientMessageId: 'client-1', via: 'start' }
])
expect(echoes.size).toBe(0)
expect(echoes.settle('client-1')).toBe(false)
})
@@ -379,20 +390,20 @@ describe('a send bound to a turn', () => {
it('dies with its child, which forgets recorded turn ends too', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
echoes.endTurn('thread-2', 'turn-2', { status: 'interrupted' })
echoes.clear()
expect(echoes.size).toBe(0)
echoes.arm('client-2')
expect(echoes.bindTurn('client-2', 'thread-2', 'turn-2')).toBeNull()
expect(echoes.bindTurn('client-2', 'thread-2', 'turn-2', 'start')).toBeNull()
})
it('is matched by thread as well as turn id', () => {
const echoes = createCodexDispatchEchoes()
echoes.arm('client-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
expect(echoes.endTurn('thread-2', 'turn-1', { status: 'interrupted' })).toEqual([])
expect(echoes.size).toBe(1)
@@ -17,7 +17,7 @@ import {
agentSessionFailureWords,
type AgentJournalDispatchRejection
} from '../../shared/agent-session-failure-words'
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
import type { AgentJournalAnsweredTurnIdentity } 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'
@@ -43,8 +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
/** The turn Codex answered the send into, whose end settled it, and how the send joined it. */
answeredInTurn: AgentJournalAnsweredTurnIdentity
} & AgentJournalDispatchRejection
function errorDetail(params: unknown): ProviderDiagnostic | undefined {
@@ -97,10 +97,14 @@ export function settleCodexSendsInEndedTurn(
return
}
const rejection = codexTurnEndRejection(end)
const answeredInTurn = codexTurnLifecycleIdentity(frame.sessionId, turnId)
for (const clientMessageId of session.dispatchEchoes.endTurn(session.threadId, turnId, end)) {
const turn = codexTurnLifecycleIdentity(frame.sessionId, turnId)
for (const { clientMessageId, via } of session.dispatchEchoes.endTurn(
session.threadId,
turnId,
end
)) {
if (rejection) {
settle({ clientMessageId, state: 'rejected', answeredInTurn, ...rejection })
settle({ clientMessageId, state: 'rejected', answeredInTurn: { turn, via }, ...rejection })
}
}
}
+19 -8
View File
@@ -1,5 +1,8 @@
import { agentSessionFailureFact, providerDiagnosticOf } from '../../shared/agent-session-failure'
import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
import type {
AgentJournalMessageItem,
AgentJournalTurnJoin
} from '../../shared/agent-session-journal-types'
import type { NativeChatBlock } from '../../shared/native-chat-types'
import type { AgentSessionDispatchOutcome } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import {
@@ -108,7 +111,7 @@ async function steerCodexTurn(
host: CodexTurnHost,
expectedTurnId: string,
input: { clientMessageId: string; body: AgentJournalMessageItem; timeoutMs?: number }
): Promise<{ turnId: string } | null> {
): Promise<{ turnId: string; via: 'steer' } | null> {
try {
const answer = await host.connection.request(
'turn/steer',
@@ -120,7 +123,7 @@ async function steerCodexTurn(
},
{ timeoutMs: input.timeoutMs }
)
return { turnId: readCodexTurnId(answer) ?? expectedTurnId }
return { turnId: readCodexTurnId(answer) ?? expectedTurnId, via: 'steer' }
} catch (error) {
if (isCodexAppServerRequestError(error) || isCodexAppServerUnsupportedError(error)) {
return null
@@ -142,7 +145,7 @@ export async function startCodexTurn(
requestedAt?: number
timeoutMs?: number
}
): Promise<{ turnId: string | null } | false> {
): Promise<{ turnId: string | null; via: AgentJournalTurnJoin } | false> {
// Armed before the write: the echo and `turn/started` can both land while the
// response is in flight, and the start must snapshot this send in its frontier.
if (!host.dispatchEchoes.arm(input.clientMessageId, input.requestedAt)) {
@@ -168,7 +171,7 @@ export async function startCodexTurn(
},
{ timeoutMs: input.timeoutMs }
)
return { turnId: readCodexTurnId(answer) }
return { turnId: readCodexTurnId(answer), via: 'start' }
}
/**
@@ -187,7 +190,7 @@ export async function dispatchCodexTurn(
},
timeoutMs: number | undefined
): Promise<AgentSessionDispatchOutcome> {
let answer: { turnId: string | null } | false
let answer: { turnId: string | null; via: AgentJournalTurnJoin } | false
try {
answer = await startCodexTurn(session, { ...input, timeoutMs })
} catch (error) {
@@ -211,13 +214,21 @@ export async function dispatchCodexTurn(
}
// An answer read after the turn it names already ended is settled by that end.
const endedFirst = answer.turnId
? session.dispatchEchoes.bindTurn(input.clientMessageId, session.threadId, answer.turnId)
? session.dispatchEchoes.bindTurn(
input.clientMessageId,
session.threadId,
answer.turnId,
answer.via
)
: null
const rejection = endedFirst ? codexTurnEndRejection(endedFirst) : null
return rejection && answer.turnId
? {
state: 'rejected',
answeredInTurn: codexTurnLifecycleIdentity(input.sessionId, answer.turnId),
answeredInTurn: {
turn: codexTurnLifecycleIdentity(input.sessionId, answer.turnId),
via: answer.via
},
...rejection
}
: { state: 'admitted' }
@@ -6,6 +6,7 @@ import {
type UnreadAgentSessionFailureFact
} from '../../../shared/agent-session-failure'
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
import type { AgentJournalAnsweredTurn } from '../../../shared/agent-session-journal-types'
import { journalDispatchRowApplies } from './journal-dispatch-settlement'
import type { JournalReducerState } from './journal-reducer'
import {
@@ -34,14 +35,11 @@ export function applyJournalDispatchRow(
} else {
delete submission.rejection
}
if (
row.state === 'rejected' &&
typeof row.answeredInTurnItemId === 'string' &&
row.answeredInTurnItemId
) {
submission.answeredInTurnItemId = row.answeredInTurnItemId
const answeredInTurn = row.state === 'rejected' ? readAnsweredTurn(row.answeredInTurn) : undefined
if (answeredInTurn) {
submission.answeredInTurn = answeredInTurn
} else {
delete submission.answeredInTurnItemId
delete submission.answeredInTurn
}
submission.resolvedAt = row.state === 'pending' ? null : row.ts
if (row.state === 'pending') {
@@ -70,6 +68,19 @@ export function applyJournalDispatchRow(
})
}
/** A stored answered turn; one malformed, or naming a way of joining this build does not know, is
* dropped, never the row. */
function readAnsweredTurn(value: unknown): AgentJournalAnsweredTurn | undefined {
if (typeof value !== 'object' || value === null) {
return undefined
}
const turnItemId = 'turnItemId' in value ? value.turnItemId : undefined
const via = 'via' in value ? value.via : undefined
return typeof turnItemId === 'string' && turnItemId && (via === 'start' || via === 'steer')
? { turnItemId, via }
: undefined
}
/** A stored rejection fact, read where it can be placed; a kind it cannot place is kept as
* written, so the classifier still knows a fact was there without this build claiming what it
* says. Shared with the queued-draft table, whose returned card mirrors its submission. */
@@ -111,7 +111,12 @@ export function journalDispatchRowBuilder(
reason: boundedDispatchReason(input),
...(input.state === 'rejected' ? { rejection: input.rejection } : {}),
...(input.state === 'rejected' && input.answeredInTurn
? { answeredInTurnItemId: agentJournalItemKey(input.answeredInTurn) }
? {
answeredInTurn: {
turnItemId: agentJournalItemKey(input.answeredInTurn.turn),
via: input.answeredInTurn.via
}
}
: {}),
...journalRowBase(state().epoch, seq, input.fence, ts),
...(input.recovered ? { recovered: input.recovered } : {}),
@@ -16,6 +16,7 @@
import type { AgentSessionFailureFact } from '../../../shared/agent-session-failure'
import {
AGENT_SESSION_JOURNAL_SCHEMA_VERSION,
type AgentJournalAnsweredTurn,
type AgentJournalDispatchState,
type AgentJournalItemBody,
type AgentJournalMessageItem,
@@ -150,10 +151,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
/** On `rejected`: the turn a Codex send was answered into, and how it joined it, when that
* turn's end settled the send. Absent on every other row. Older readers keep the key and ignore
* it; a malformed one is dropped when read, never the row. */
answeredInTurn?: AgentJournalAnsweredTurn
}
/** An item mutation may name its own producer, because one batch can CREATE
@@ -1,5 +1,6 @@
import type { AgentJournalDispatchRejection } from '../../../shared/agent-session-failure-words'
import type {
AgentJournalAnsweredTurnIdentity,
AgentJournalCursor,
AgentJournalItemBody,
AgentJournalItemIdentity,
@@ -42,7 +43,7 @@ export type ResolveDispatchInput = {
* `agentSessionFailureWords`, never written by hand. */
| ({
state: 'rejected'
answeredInTurn?: AgentJournalItemIdentity
answeredInTurn?: AgentJournalAnsweredTurnIdentity
} & AgentJournalDispatchRejection)
| { state: 'unknown'; reason?: string | null }
)
@@ -94,7 +94,7 @@ async function sendWithdrawnByItsTurnEnd() {
clientMessageId: 'send-1',
state: 'rejected',
...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }),
answeredInTurn: TURN,
answeredInTurn: { turn: TURN, via: 'start' },
fence: 1
})
return journal
@@ -107,23 +107,25 @@ describe('the turn a rejected submission was answered into', () => {
expect(journal.submission('send-1')).toMatchObject({
dispatchState: 'rejected',
answeredInTurnItemId: turnItemId
answeredInTurn: { turnItemId, via: 'start' }
})
expect(readAgentSessionHydrationPage(journal).submissions).toEqual([
expect.objectContaining({ answeredInTurnItemId: turnItemId })
expect.objectContaining({ answeredInTurn: { turnItemId, via: 'start' } })
])
await journals.closeAll()
const replayed = await journals.open({ identity: IDENTITY, stateDirectory: root! })
expect(replayed.submission('send-1')).toMatchObject({ answeredInTurnItemId: turnItemId })
expect(replayed.submission('send-1')).toMatchObject({
answeredInTurn: { turnItemId, via: 'start' }
})
})
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(journal.submission('send-1')).not.toHaveProperty('answeredInTurn')
expect(readAgentSessionHydrationPage(journal).submissions[0]).not.toHaveProperty(
'answeredInTurnItemId'
'answeredInTurn'
)
})
})
@@ -12,6 +12,7 @@ import type {
import type { AgentSessionBackgroundTaskStops } from '../../../shared/agent-child-work-stop-targets'
import type {
AgentJournalAnsweredTurnIdentity,
AgentJournalItemIdentity,
AgentJournalItemBody,
AgentJournalMessageItem,
@@ -180,7 +181,7 @@ export type AgentSessionDispatchOutcome =
* provider answered the send into, which ended before the answer was read. */
| ({
state: 'rejected'
answeredInTurn?: AgentJournalItemIdentity
answeredInTurn?: AgentJournalAnsweredTurnIdentity
} & AgentJournalDispatchRejection)
/** The call did not settle. Never re-send on the user's behalf. */
| { state: 'unknown'; reason: string }
@@ -1,4 +1,7 @@
import type { AgentJournalItemIdentity } from '../../../shared/agent-session-journal-types'
import type {
AgentJournalAnsweredTurnIdentity,
AgentJournalItemIdentity
} from '../../../shared/agent-session-journal-types'
import type { AgentJournalDispatchRejection } from '../../../shared/agent-session-failure-words'
import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations'
import { structuredAgentSessionConversationFence } from './structured-agent-session-provider-child'
@@ -14,7 +17,7 @@ export async function settleStructuredAgentSessionLateDispatch(
| { providerIdentity: AgentJournalItemIdentity }
| ({
state: 'rejected'
answeredInTurn?: AgentJournalItemIdentity
answeredInTurn?: AgentJournalAnsweredTurnIdentity
} & AgentJournalDispatchRejection)
| { state: 'unknown'; reason: string }
)
@@ -317,12 +317,12 @@ describe('the turn a withdrawn Codex send was answered into', () => {
item.body.kind === 'turn' ? [item.itemId] : []
),
named: snapshot.submissions.find((entry) => entry.clientMessageId === clientMessageId)
?.answeredInTurnItemId,
onPage: onPage?.answeredInTurnItemId
?.answeredInTurn,
onPage: onPage?.answeredInTurn
}
}
it("is that turn's record, for the send that opened it and for one steered into it", async () => {
it("is that turn's record, started by the send that opened it and steered by a later one", async () => {
const opening = await send('look around')
await vi.waitFor(() => expect(answers).toBe(1))
turns.start()
@@ -336,13 +336,13 @@ describe('the turn a withdrawn Codex send was answered into', () => {
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]
const started = { turnItemId: turnRecords[0], via: 'start' }
const steeredIn = { turnItemId: turnRecords[0], via: 'steer' }
expect(await answeredInto(opening)).toEqual({ turnRecords, named: started, onPage: started })
expect(await answeredInto(steered)).toEqual({
turnRecords,
named: steeredIn,
onPage: steeredIn
})
})
@@ -373,7 +373,7 @@ describe('the turn a withdrawn Codex send was answered into', () => {
)
const { turnRecords, named } = await answeredInto(sent)
expect(turnRecords).toHaveLength(1)
expect(named).toBe(turnRecords[0])
expect(named).toEqual({ turnItemId: turnRecords[0], via: 'start' })
})
})
+2 -1
View File
@@ -319,7 +319,8 @@ export const AgentJournalSubmissionSchema = z.object({
submittedAt: z.number(),
resolvedAt: z.number().nullable(),
submittedSequence: z.number().int().optional(),
answeredInTurnItemId: z.string().min(1).optional(),
// Open like `dispatchState`: a way of joining a newer host names must not drop the submission.
answeredInTurn: z.object({ turnItemId: z.string().min(1), via: z.string().min(1) }).optional(),
recovered: z.literal(true).optional(),
handoverRecorded: z.literal(true).optional(),
handedOverAt: z.number().optional(),
+16 -4
View File
@@ -391,6 +391,18 @@ export type AgentJournalRenderItem = AgentJournalProducerLinkage & {
// ─── Submissions ────────────────────────────────────────────────────────────
/** The turn a send was answered into: its record's item id, and how the send joined it. `start`:
* the provider answered the send's start request with that turn; `steer`: Orca steered it into
* that running turn. Known limit: a start the provider silently folds into a running turn reads
* as `start`. Open: a newer host may name another way, which a reader leaves unclaimed. */
export type AgentJournalAnsweredTurn = { turnItemId: string; via: AgentJournalTurnJoin }
export type AgentJournalTurnJoin = 'start' | 'steer'
/** The same, as a writer names it: the turn record's identity, keyed when the row is written. */
export type AgentJournalAnsweredTurnIdentity = {
turn: AgentJournalItemIdentity
via: AgentJournalTurnJoin
}
export const AGENT_JOURNAL_DISPATCH_STATES = ['pending', 'accepted', 'rejected', 'unknown'] as const
export type AgentJournalDispatchState = (typeof AGENT_JOURNAL_DISPATCH_STATES)[number]
@@ -415,10 +427,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
/** On `rejected`: the turn 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. */
answeredInTurn?: AgentJournalAnsweredTurn
/** 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
@@ -40,7 +40,7 @@ function withoutPosition(page: AgentSessionHistoryPage): AgentSessionHistoryPage
return {
...page,
submissions: page.submissions.map(
({ submittedSequence: _submitted, answeredInTurnItemId: _turn, ...rest }) => rest
({ submittedSequence: _submitted, answeredInTurn: _turn, ...rest }) => rest
)
}
}
@@ -78,10 +78,13 @@ test('a released client folds and draws a page whose submissions carry their jou
state: 'rejected',
...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }),
answeredInTurn: {
provider: 'legacy',
agent: 'codex',
sessionId: IDENTITY.sessionId,
recordId: 'turn-lifecycle:turn-1'
turn: {
provider: 'legacy',
agent: 'codex',
sessionId: IDENTITY.sessionId,
recordId: 'turn-lifecycle:turn-1'
},
via: 'start'
},
fence: 1
})
@@ -93,7 +96,9 @@ test('a released client folds and draws a page whose submissions carry their jou
})
)
expect(page.submissions.find((entry) => entry.clientMessageId === 'turn-ended')).toEqual(
expect.objectContaining({ answeredInTurnItemId: expect.any(String) })
expect.objectContaining({
answeredInTurn: { turnItemId: expect.any(String), via: 'start' }
})
)
const checkout = await materializeReleaseCheckout(BASELINE_REF)
@@ -117,7 +122,7 @@ test('a released client folds and draws a page whose submissions carry their jou
const state = reduce(empty, { type: 'history-page', page: from })
return {
submissions: state.submissions.map(
({ submittedSequence: _s, answeredInTurnItemId: _t, ...rest }) => rest
({ submittedSequence: _s, answeredInTurn: _t, ...rest }) => rest
),
messages: project(state.items, [], state.submissions)
}