fix(native-chat): keep a Codex ask's questions in the order it asked them (#23502)

* fix(native-chat): keep a Codex ask's questions in the order it asked them

Codex journals every question of one ask in a single write, so the
questions share a timestamp. The transcript list sorted rows by
timestamp and broke ties by row id, and a question's id ends in the
question id the model chose, so answered, cancelled and still-pending
rows of one ask came out in the alphabetical order of those ids.

The transcript projection now breaks timestamp ties by the order its
source gave: the journal's order for structured sessions. The session
assembler keeps its id tie-break, since the sources it merges share no
order of their own.

* refactor(native-chat): name the id tie-break comparator for what it does

Two shared comparators differed only in whether they break timestamp ties by id,
under near-identical names; spell the id tie-break in the name.

* fix(native-chat): order structured chat rows by their journal position

The desktop list and the host's conversation outline sorted structured rows
by timestamp. The journal's contract is that the sequence orders the
timeline and the timestamp is the provider's clock: a Codex ask writes all
of its questions in one write (one sequence, one timestamp), and a row
recovered after a crash carries an earlier clock at a later sequence.

The host now records each item's place within the write that created it,
keeps it across revisions like the sequence, and sends it as an optional
field. Rows projected from the journal carry that position, and both the
desktop list and the outline order journal rows by it. Rows outside the
journal keep their rules: rank for the streaming and pending tail, outbox
sends after every journal row, and terminal-backed chats keep time then id.

* fix(native-chat): keep a refused send at its journal place, and keep list positions off worker reads

A send the host journalled before the provider refused it is shown from the outbox, and it
sorted after every journal row, so it dropped below whatever the agent wrote after it. It now
takes the journal position of the submission the host recorded.

Structured worker reads and their archives projected journal rows through the same
projection, so they returned the list-only journal position on every message. The worker
payload bound now drops it.

* test(native-chat): a failed restart's row draws below the message it failed

Since a message is accepted before its delivery starts the agent, a restart that fails is
journalled after the message, and the refused message keeps that journal place in the chat.
Both host paths now pin it through the chat's own projection: a start the host could not make,
and a restarted child that exits before proving its start.

Folds the journal reducer's batch item write onto fewer lines, which the merge of main pushed
past the file's line limit.
This commit is contained in:
Brennan Benson
2026-09-27 23:44:36 -07:00
committed by GitHub
parent 56e691344f
commit 5e5f4f6603
26 changed files with 683 additions and 55 deletions
+2
View File
@@ -6,6 +6,8 @@
"./scripts/vitest-host-ports-setup.ts",
"../src/main/**/*",
"../src/renderer/src/lib/skill-freshness-display-status.ts",
"../src/renderer/src/components/native-chat/native-chat-resolution-receipt.ts",
"../src/renderer/src/components/native-chat/structured-agent-question-projection.ts",
"../src/preload/**/*",
"../src/shared/**/*",
"../src/relay/**/*",
@@ -0,0 +1,257 @@
// A Codex ask with several questions, driven from the real host journal through
// the client's session reducer to the rows the transcript list draws.
import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { computeAgentSessionPayloadFingerprint } from '../../shared/agent-session-mutation-envelope'
import type { AgentJournalRenderItem } from '../../shared/agent-session-journal-types'
import {
EMPTY_STRUCTURED_AGENT_SESSION,
reduceStructuredAgentSession,
type StructuredAgentSessionState
} from '../../shared/structured-agent-session-reducer'
import { projectStructuredAgentSessionMessages } from '../../shared/structured-agent-session-message-projection'
import { projectNativeChatTranscriptMessages } from '../../shared/native-chat-transcript-projection'
import { AgentSessionRecordStore } from '../runtime/agent-session-record-store'
import { CodexJournalPrompts } from './codex-structured-journal-prompts'
import { CODEX_USER_INPUT_METHOD } from './codex-structured-prompt-replies'
import type { StructuredAgentSessionAdapter } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { StructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-host'
import {
HOST_TEST_NOW,
HOST_TEST_SESSION as SESSION,
HOST_TEST_THREAD as THREAD,
hostTestAttachParams,
hostTestOperationId,
resetHostTestOperationIds
} from '../native-chat/agent-session-wire/structured-agent-session-host-test-data'
import {
projectStructuredQuestionMessages,
structuredQuestionTranscript
} from '../../renderer/src/components/native-chat/structured-agent-question-projection'
const CALLER = { callerKey: 'client-1' }
type Asked = readonly { id: string; question: string }[]
// Codex's question ids are the model's own words, so their text order is not
// the order it asked in. These two asks spell the two orders live QA saw.
const ASKED: Asked = [
{ id: 'scope', question: 'Which files are in scope?' },
{ id: 'priority', question: 'What matters most?' },
{ id: 'deadline', question: 'When is it due?' }
]
const ASKED_OUT_OF_ORDER: Asked = [
{ id: 'format', question: 'Which format?' },
{ id: 'audience', question: 'Who reads it?' },
{ id: 'length', question: 'How long?' }
]
let root: string
let store: AgentSessionRecordStore
let host: StructuredAgentSessionHost
let sink: StructuredAgentSessionEventSink | null
let client: StructuredAgentSessionState
let clock: number
beforeEach(async () => {
root = await mkdtemp(join(tmpdir(), 'orca-codex-question-order-'))
resetHostTestOperationIds()
sink = null
clock = HOST_TEST_NOW
const adapter: StructuredAgentSessionAdapter = {
acquire: vi.fn<StructuredAgentSessionAdapter['acquire']>(async ({ fence, events }) => {
sink = events ?? null
return {
process: {
hostId: 'local',
pid: 4242,
processStartTimeMs: 1_700_000_000_000,
spawnToken: store.getRecord(SESSION)?.lease.reservedSpawnToken ?? 'spawn-a'
},
link: {
linkId: `link-${fence}`,
handle: { provider: 'codex', threadId: THREAD },
origin: 'created',
mintedAtFence: fence,
observedAt: HOST_TEST_NOW
}
}
}),
releaseAcquisition: vi.fn(async () => true),
dispatch: vi.fn(),
cancelTurn: vi.fn(async () => ({ cancelled: true })),
answerPrompt: vi.fn(async ({ commit }) => commit()),
setOption: vi.fn(async () => undefined)
}
store = await AgentSessionRecordStore.open({ directory: join(root, 'store'), hostId: 'local' })
host = new StructuredAgentSessionHost({
store,
adapter,
journalRoot: root,
claimKeyId: 'key-1',
mintSpawnToken: () => 'spawn-a',
// Every write lands on its own millisecond, as it does live.
now: () => (clock += 1)
})
expect((await host.attach(CALLER, hostTestAttachParams(null))).ok).toBe(true)
const page = host.history({ sessionId: SESSION, direction: 'tail' })
if (!page.ok) {
throw new Error('no history page')
}
client = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
type: 'history-page',
page: page.page
})
host.subscribe({
id: 'client',
sessionId: SESSION,
cursor: page.page.liveCursor ?? page.page.window.nextCursor,
emit: (event) => {
client = reduceStructuredAgentSession(client, { type: 'event', event })
}
})
})
afterEach(async () => {
await host.flushAllStreamedEvents()
await rm(root, { recursive: true, force: true })
})
function codexPrompts(): CodexJournalPrompts {
if (!sink) {
throw new Error('session was never acquired')
}
return new CodexJournalPrompts(
{ sink, linkageFor: () => ({}) },
() => null,
() => 'turn-1'
)
}
async function ask(prompts: CodexJournalPrompts, asked: Asked = ASKED): Promise<void> {
prompts.handle({
threadId: THREAD,
method: CODEX_USER_INPUT_METHOD,
codexItemId: 'codex-item-1',
promptKey: '7',
params: {
threadId: THREAD,
turnId: 'turn-1',
questions: asked.map((question) => ({
...question,
options: [
{ label: 'Yes', description: '' },
{ label: 'No', description: '' }
]
}))
}
})
await host.flushStreamedEvents(SESSION)
}
function questionItem(question: string): AgentJournalRenderItem {
const item = client.items.find(
(candidate) => candidate.body.kind === 'question' && candidate.body.question === question
)
if (!item || item.body.kind !== 'question') {
throw new Error(`no journal item for ${question}`)
}
return item
}
async function answer(question: string): Promise<void> {
const item = questionItem(question)
if (item.body.kind !== 'question') {
return
}
const fields = {
itemId: item.itemId,
expectedRevision: item.revision,
optionId: item.body.options[0]!.id
}
const result = await host.respondToPrompt(CALLER, {
envelope: {
sessionId: SESSION,
clientOperationId: hostTestOperationId(),
expectedRuntimeFence: store.getRecord(SESSION)?.lease.runtimeFence ?? 1,
payloadFingerprint: computeAgentSessionPayloadFingerprint({
method: 'agentSession.respondTo:question',
sessionId: SESSION,
fields
})
},
kind: 'question',
...fields
})
expect(result.ok).toBe(true)
await host.flushStreamedEvents(SESSION)
}
/** The prompt rows the transcript list draws, top to bottom, as the questions each one shows. */
function drawnPromptRows(): string[][] {
const { receipts } = structuredQuestionTranscript(client.items)
// The desktop list's projection: its comparator adds only a rank for rows the host never writes.
const rows = projectNativeChatTranscriptMessages(
projectStructuredAgentSessionMessages(
client.items,
[],
client.submissions,
projectStructuredQuestionMessages
)
)
return rows.flatMap((row) => {
const prompt = receipts.get(row.id)
if (!prompt || prompt.kind !== 'question') {
return []
}
const questions = prompt.questions?.length ? prompt.questions : [prompt]
return [questions.map(({ question }) => `${question} (${prompt.resolution.state})`)]
})
}
describe('a Codex ask with several questions', () => {
it('keeps the pending rest of the ask below the question already answered', async () => {
await ask(codexPrompts())
await answer(ASKED[0]!.question)
expect(drawnPromptRows()).toEqual([
['Which files are in scope? (resolved)'],
['What matters most? (pending)', 'When is it due? (pending)']
])
})
it('draws every answered question in the order Codex asked it', async () => {
await ask(codexPrompts())
for (const { question } of ASKED) {
await answer(question)
}
expect(drawnPromptRows()).toEqual([
['Which files are in scope? (resolved)'],
['What matters most? (resolved)'],
['When is it due? (resolved)']
])
// Mobile draws the shared projection in journal order, one row per question.
expect(
projectStructuredAgentSessionMessages(client.items, [], client.submissions).map(
({ blocks }) => (blocks[0]?.type === 'text' ? blocks[0].text.split('\n')[0] : null)
)
).toEqual(ASKED.map(({ question }) => question))
})
it('draws a cancelled ask in the order Codex asked it', async () => {
const prompts = codexPrompts()
await ask(prompts, ASKED_OUT_OF_ORDER)
prompts.cancel(questionItem(ASKED_OUT_OF_ORDER[0]!.question).itemId)
await host.flushStreamedEvents(SESSION)
expect(drawnPromptRows()).toEqual([
['Which format? (cancelled)'],
['Who reads it? (cancelled)'],
['How long? (cancelled)']
])
})
})
@@ -118,6 +118,39 @@ describe('ordering', () => {
expect(renderJournalState(state).items.map((item) => item.itemId)).toEqual(['earlier', 'later'])
})
it("places a batch's writes by their order in it, and keeps that place on revision", () => {
// One Codex ask writes all its questions in one batch; their ids are not their order.
const state = fold([
{
kind: 'lifecycle-batch',
settlementId: 'ask',
mutations: [
{ kind: 'item', itemId: 'scope', revision: 1, body: text('first') },
{ kind: 'item', itemId: 'priority', revision: 1, body: text('second') },
{ kind: 'item', itemId: 'deadline', revision: 1, body: text('third') }
],
...base(1)
},
{
kind: 'lifecycle-batch',
settlementId: 'answer',
mutations: [{ kind: 'item', itemId: 'deadline', revision: 2, body: text('answered') }],
...base(2)
}
])
expect(
renderJournalState(state).items.map(({ itemId, sequence, sequenceIndex }) => ({
itemId,
sequence,
sequenceIndex
}))
).toEqual([
{ itemId: 'scope', sequence: 1, sequenceIndex: undefined },
{ itemId: 'priority', sequence: 1, sequenceIndex: 1 },
{ itemId: 'deadline', sequence: 1, sequenceIndex: 2 }
])
})
it('orders by sequence even when the observed timestamp runs backwards', () => {
const state = fold([
{ kind: 'item', itemId: 'late', revision: 1, body: text('late'), ...base(1), ts: 9_000 },
@@ -4,9 +4,10 @@
//
// Rules: highest revision wins, a tombstone removes, a late lower revision is
// dropped rather than resurrecting stale content, and ordering is by the
// sequence of the row that CREATED an item (a later revision updates the body,
// it does not move the bubble). Producer linkage is likewise the creating
// write's: a revision naming no producer keeps it, one naming any replaces it.
// position (sequence, then place in the row) of the write that CREATED an item
// (a later revision updates the body, it does not move the bubble). Producer
// linkage is likewise the creating write's: a revision naming no producer keeps
// it, one naming any replaces it.
import type {
AgentJournalAcceptanceReceipt,
@@ -15,6 +16,7 @@ import type {
AgentJournalSubmission
} from '../../../shared/agent-session-journal-types'
import { journalBatchMutationProducer, journalRenderItem } from './journal-render-item'
import { compareAgentJournalItems } from '../../../shared/agent-session-journal-position'
import {
agentJournalSubmissionKey,
parseAgentJournalItemKey
@@ -90,25 +92,17 @@ export function applyJournalRow(state: JournalReducerState, row: JournalRow): vo
if (state.appliedSettlementIds.has(row.settlementId)) {
return
}
for (const mutation of row.mutations) {
for (const [sequenceIndex, mutation] of row.mutations.entries()) {
if (mutation.kind === 'item') {
if (journalItemRevisionIsStale(state, mutation.itemId, mutation.revision)) {
continue
}
const itemId = resolveJournalItemId(state, mutation.itemId, mutation.body)
const { revision, body } = mutation
const itemId = resolveJournalItemId(state, mutation.itemId, body)
acceptSubmissionFromProviderItem(state, mutation.itemId, itemId, row)
upsertItem(
state,
itemId,
mutation.revision,
journalRenderItem(
itemId,
mutation.revision,
mutation.body,
row,
journalBatchMutationProducer(row, mutation)
)
)
const producer = journalBatchMutationProducer(row, mutation)
const item = journalRenderItem(itemId, revision, body, row, producer, sequenceIndex)
upsertItem(state, itemId, revision, item)
} else {
removeItem(state, resolveItemId(state, mutation.itemId), mutation.revision)
}
@@ -216,14 +210,16 @@ function upsertItem(
existing.body.kind === 'message' &&
existing.body.role === 'user' &&
parseAgentJournalItemKey(itemId)?.provider === 'orca'
const { sequenceIndex: _revisedAt, ...revised } = next
state.items.set(itemId, {
...next,
...revised,
// Settlements, prompt answers and reopen sweeps revise rows any agent wrote
// without naming one; each would otherwise hand a subagent's row to the session.
...(namesAgentJournalProducer(next) ? {} : agentJournalLinkageFields(existing)),
// Provider history may normalize text or omit local attachments from the original send.
body: submitted ? existing.body : next.body,
sequence: existing.sequence,
...(existing.sequenceIndex !== undefined ? { sequenceIndex: existing.sequenceIndex } : {}),
observedAt: existing.observedAt
})
state.tombstones.delete(itemId)
@@ -333,9 +329,9 @@ function acceptSubmissionFromProviderItem(
/** Project the folded state into the client-facing snapshot. */
export function renderJournalState(state: JournalReducerState): AgentJournalSnapshot {
// Sequence is the sole ordering key; map insertion order is not, because a
// re-created item re-enters the map after the items that followed it.
const items = [...state.items.values()].sort((a, b) => a.sequence - b.sequence)
// The journal position is the sole ordering key; map insertion order is not,
// because a re-created item re-enters the map after the items that followed it.
const items = [...state.items.values()].sort(compareAgentJournalItems)
return {
sessionId: state.sessionId,
cursor: { epoch: state.epoch, sequence: state.lastSequence },
@@ -19,13 +19,16 @@ export function journalRenderItem(
revision: number,
body: AgentJournalItemBody,
row: JournalRow,
producer: AgentJournalProducerLinkage = row
producer: AgentJournalProducerLinkage = row,
/** Which of the row's writes this is; only a lifecycle batch has more than one. */
sequenceIndex = 0
): AgentJournalRenderItem {
return {
itemId,
revision,
body,
sequence: row.seq,
...(sequenceIndex > 0 ? { sequenceIndex } : {}),
observedAt: row.ts,
...(row.recovered ? { recoveredAt: row.ts } : {}),
...(row.recovered ? { recovered: row.recovered } : {}),
@@ -6,6 +6,7 @@
// key instead of appearing as a second copy of the user's own message.
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
import { compareAgentJournalItems } from '../../../shared/agent-session-journal-position'
import type {
AgentJournalRenderItem,
AgentJournalSnapshot,
@@ -65,7 +66,7 @@ export function projectJournalBatch(input: {
const items = [...touchedItemIds]
.map((itemId) => live.get(itemId))
.filter((item) => item !== undefined)
.sort((a, b) => a.sequence - b.sequence)
.sort(compareAgentJournalItems)
return {
ok: true,
batch: {
@@ -7,8 +7,8 @@ import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
import type { AgentJournalSubmission } from '../../../shared/agent-session-journal-types'
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
import type { AgentJournalSubmission } from '../../../shared/agent-session-journal-types'
import type {
AgentSessionSubscribeEvent,
AgentSessionTurnCompletionEvent
@@ -35,6 +35,7 @@ import {
HOST_TEST_SESSION as SESSION,
HOST_TEST_THREAD as THREAD,
hostTestAttachParams,
hostTestDrawnRowIds,
hostTestMessage,
hostTestOperationId,
resetHostTestOperationIds
@@ -322,6 +323,25 @@ describe('a start the chat needed and did not get', () => {
expect(errorRows()).toHaveLength(1)
})
it('draws the messages it failed above the error row, since they were accepted first', async () => {
await host.close(SESSION)
acquire.mockRejectedValueOnce(new Error('spawn codex ENOENT'))
const first = await accept('first')
const second = await accept('second')
await eventually(() => expect(submission(second)?.dispatchState).toBe('rejected'))
const snapshot = host.journalSnapshot(SESSION)
const errorRow = snapshot.items.find(
(item) => item.body.kind === 'status' && item.body.tone === 'error'
)?.itemId
const shown = [agentJournalSubmissionKey(first), agentJournalSubmissionKey(second), errorRow]
const drawn = hostTestDrawnRowIds(snapshot, [
{ clientMessageId: first, text: 'first' },
{ clientMessageId: second, text: 'second' }
])
expect(drawn.filter((id) => shown.includes(id))).toEqual(shown)
})
it('notifies failed once for the queued messages one start failure refused', async () => {
await host.close(SESSION)
acquire.mockRejectedValueOnce(new Error('spawn codex ENOENT'))
@@ -1,6 +1,12 @@
import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types'
import type {
AgentJournalMessageItem,
AgentJournalSnapshot
} from '../../../shared/agent-session-journal-types'
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
import type { AgentSessionExecutionLocation } from '../../../shared/agent-session-record'
import { projectNativeChatTranscriptMessages } from '../../../shared/native-chat-transcript-projection'
import { projectStructuredAgentSessionMessages } from '../../../shared/structured-agent-session-message-projection'
import { createStructuredAgentSessionOutboxEntry } from '../../../shared/structured-agent-session-outbox'
import { attachFingerprintFields } from './structured-agent-session-attach'
import type { AgentSessionAttachParams } from './structured-agent-session-attach'
@@ -61,3 +67,23 @@ export function hostTestAttachParams(
}
}
}
/** The row ids a chat draws from `snapshot`, top to bottom, while its composer still holds `sent`
* as dispatched, as it does until the journal accepts them. */
export function hostTestDrawnRowIds(
snapshot: AgentJournalSnapshot,
sent: readonly { clientMessageId: string; text: string }[]
): string[] {
const outbox = sent.map((message) => ({
...createStructuredAgentSessionOutboxEntry({
...message,
sessionId: HOST_TEST_SESSION,
attachments: [],
queuedAt: HOST_TEST_NOW
}),
state: 'dispatching' as const
}))
return projectNativeChatTranscriptMessages(
projectStructuredAgentSessionMessages(snapshot.items, outbox, snapshot.submissions)
).map(({ id }) => id)
}
@@ -11,6 +11,7 @@ import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
import type { AgentSessionMutationEnvelope } from '../../../shared/agent-session-wire'
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
@@ -21,6 +22,7 @@ import {
HOST_TEST_SESSION as SESSION,
HOST_TEST_THREAD as THREAD,
hostTestAttachParams,
hostTestDrawnRowIds,
hostTestMessage,
hostTestOperationId,
resetHostTestOperationIds
@@ -206,6 +208,15 @@ describe('a send into a published session whose child ended before startup', ()
expect(journalStatuses().slice(rowsBefore)).toEqual([
expect.stringMatching(/stopped before it finished starting: .*not signed in/)
])
// Accepted before the restart it needed, so the chat draws it above the row naming the cause.
const snapshot = host.journalSnapshot(SESSION)
const causeRow = snapshot.items.findLast((item) => item.body.kind === 'status')?.itemId
const shown = [agentJournalSubmissionKey(held), causeRow]
expect(
hostTestDrawnRowIds(snapshot, [
{ clientMessageId: held, text: 'still not signed in' }
]).filter((id) => shown.includes(id))
).toEqual(shown)
// The failed restart moved the fence twice: the acquisition, and the exit that released it.
expect(store.getRecord(SESSION)?.lease.runtimeFence).toBe(releasedFence + 2)
expect(store.getRecord(SESSION)?.lease.claimStatus).toBe('released')
@@ -1,7 +1,9 @@
import { describe, expect, it } from 'vitest'
import { MAX_CODEX_SUBAGENTS_PER_GROUP } from '../../codex/codex-structured-journal-limits'
import { projectStructuredItemsToNativeChat } from '../../../shared/structured-agent-session-projection'
import {
boundWorkerTranscriptMessages,
boundWorkerTranscriptTail,
redactWorkerTerminalLines
} from './worker-transcript-payload'
@@ -177,6 +179,28 @@ describe('worker transcript wire bounds', () => {
expect(result).toMatchObject({ limited: false, warnings: [] })
})
it("serves a structured worker's journal rows without their list position", () => {
const [message] = projectStructuredItemsToNativeChat([
{
itemId: 'reply',
revision: 1,
sequence: 7,
sequenceIndex: 1,
observedAt: 1,
body: { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'done' }] }
}
])
// Anti-vacuous: the projection itself does position the row.
expect(message?.journalPosition).toEqual({ sequence: 7, index: 1 })
for (const served of [
boundWorkerTranscriptMessages([message!]).messages,
boundWorkerTranscriptTail([message!], 262_144).messages
]) {
expect(served).toHaveLength(1)
expect(served[0]).not.toHaveProperty('journalPosition')
}
})
it('keeps two roster ids sharing a 512-char prefix distinct', () => {
// The id is the roster key: a plain prefix clip would merge the two children.
const head = 'a'.repeat(512)
@@ -108,8 +108,10 @@ function boundMessage(
if (blocks.length < message.blocks.length) {
markClipped(state, 'Some transcript blocks were omitted from oversized messages.')
}
// The journal position only orders a live list; a worker read is already in order.
const { journalPosition: _journalPosition, ...served } = message
return {
...message,
...served,
id: boundIdentifier(message.id, transcriptPath, state),
...(message.turnId ? { turnId: boundIdentifier(message.turnId, transcriptPath, state) } : {}),
blocks: blocks.map((block) => boundBlock(block, state))
@@ -113,6 +113,8 @@ describe('live Codex checklist frames', () => {
role: 'assistant',
timestamp: 1,
source: 'transcript',
// Journalled like the frames it sits between.
journalPosition: { sequence: 1, index: 0 },
blocks: [
{
type: 'tool-call',
@@ -60,7 +60,8 @@ const JOURNAL: AgentJournalRenderItem[] = [
]),
row(10, { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'Done.' }] }),
user(11, [{ type: 'text', text: 'Thanks' }]),
// Journalled after `Thanks` but observed before `Done.`: the transcript orders by observation.
// Recovered after a crash: journalled after `Thanks`, but carrying the provider's clock from
// before `Done.`. The transcript orders by journal position, never observation.
{ ...user(12, [{ type: 'text', text: 'Observed earlier' }]), observedAt: 1_009.5 }
]
@@ -107,13 +108,13 @@ describe('conversation outline parity with the loaded rail', () => {
}))
).toEqual(loaded.map(({ id, text, hasImages }) => ({ id, text, hasImages })))
// Anti-vacuous: the folded tool result, the refused send, the harness turn and the empty
// prompt were all dropped, and the late-journalled row sits where it was observed.
// prompt were all dropped, and the recovered row sits where it was journalled.
expect(outline.map((entry) => entry.itemId)).toEqual([
'item-1',
'item-5',
'item-9',
'item-12',
'item-11'
'item-11',
'item-12'
])
})
@@ -123,8 +124,8 @@ describe('conversation outline parity with the loaded rail', () => {
{ itemId: 'item-1', sequence: 1, preview: 'Fix the parser', imageCount: 0 },
{ itemId: 'item-5', sequence: 5, preview: '', imageCount: 1 },
{ itemId: 'item-9', sequence: 9, preview: 'Compare these', imageCount: 2 },
{ itemId: 'item-12', sequence: 12, preview: 'Observed earlier', imageCount: 0 },
{ itemId: 'item-11', sequence: 11, preview: 'Thanks', imageCount: 0 }
{ itemId: 'item-11', sequence: 11, preview: 'Thanks', imageCount: 0 },
{ itemId: 'item-12', sequence: 12, preview: 'Observed earlier', imageCount: 0 }
])
})
})
@@ -7,7 +7,7 @@ import {
type NativeChatSessionStatus
} from '../../../../shared/native-chat-types'
import { NATIVE_CHAT_STREAMING_ID } from '../../../../shared/native-chat-streaming'
import { compareNativeChatMessagesByTime } from '../../../../shared/native-chat-transcript-projection'
import { compareNativeChatTranscriptMessages } from '../../../../shared/native-chat-transcript-projection'
import {
hasImagePromptMarker,
isImageSourceUserTurn,
@@ -105,14 +105,14 @@ function messageSortRank(message: NativeChatMessage): number {
return 0
}
// Rank first; within a tier the transcript's shared time order.
// Rank first; within a tier the transcript's shared order.
export function compareMessages(a: NativeChatMessage, b: NativeChatMessage): number {
const ar = messageSortRank(a)
const br = messageSortRank(b)
if (ar !== br) {
return ar - br
}
return compareNativeChatMessagesByTime(a, b)
return compareNativeChatTranscriptMessages(a, b)
}
/**
@@ -4,6 +4,7 @@ import type {
} from '../../../../shared/agent-session-journal-types'
import { isAskUserQuestionTool } from '../../../../shared/agent-question-answered-intent'
import { parseAskFromToolInput } from '../../../../shared/native-chat-ask'
import { agentJournalItemRowOrigin } from '../../../../shared/agent-session-journal-position'
import type { NativeChatMessage } from '../../../../shared/native-chat-types'
import { projectStructuredItemToNativeChat } from '../../../../shared/structured-agent-session-projection'
import { readAgentJournalTurn } from '../../../../shared/agent-session-turn-record'
@@ -69,10 +70,8 @@ function projectItem(item: AgentJournalRenderItem): Projection {
if (body.resolution.state === 'pending') {
// A system row preserves question identity through tool folding; the receipt renders its body.
message = {
id: item.itemId,
...agentJournalItemRowOrigin(item),
role: 'system',
timestamp: item.observedAt,
source: 'transcript',
blocks: [{ type: 'text', text: body.question }]
}
}
@@ -0,0 +1,130 @@
import { describe, expect, it } from 'vitest'
import type {
AgentJournalItemBody,
AgentJournalRenderItem,
AgentJournalSubmission
} from '../../../../shared/agent-session-journal-types'
import { projectAgentSessionConversationOutline } from '../../../../shared/agent-session-conversation-outline'
import { agentJournalSubmissionKey } from '../../../../shared/agent-session-journal-item-key'
import {
createStructuredAgentSessionOutboxEntry,
type StructuredAgentSessionOutboxEntry
} from '../../../../shared/structured-agent-session-outbox'
import { createNativeChatMessageListProjection } from './native-chat-message-list-projection'
import { projectStructuredAgentSessionMessages } from './structured-agent-session-message-projection'
function journalItem(
itemId: string,
sequence: number,
observedAt: number,
body: AgentJournalItemBody,
extra: Partial<AgentJournalRenderItem> = {}
): AgentJournalRenderItem {
return { itemId, revision: 1, sequence, observedAt, body, ...extra }
}
function said(role: 'user' | 'assistant', text: string): AgentJournalItemBody {
return { kind: 'message', role, blocks: [{ type: 'text', text }] }
}
function answered(question: string): AgentJournalItemBody {
return {
kind: 'question',
question,
options: [{ id: 'yes', label: 'Yes' }],
resolution: { state: 'resolved', selectedOptionId: 'yes', resolvedBy: 'client', resolvedAt: 9 }
}
}
/** The ids the desktop transcript list draws, top to bottom. */
function drawn(
items: AgentJournalRenderItem[],
outbox: StructuredAgentSessionOutboxEntry[] = [],
submissions: AgentJournalSubmission[] = []
): string[] {
return createNativeChatMessageListProjection()(
projectStructuredAgentSessionMessages(items, outbox, submissions)
).map(({ id }) => id)
}
function queued(clientMessageId: string, text: string, queuedAt: number) {
return createStructuredAgentSessionOutboxEntry({
clientMessageId,
sessionId: 'session',
text,
attachments: [],
queuedAt
})
}
describe('structured transcript order', () => {
it('draws a row recovered after a crash at its journal position, not its earlier clock', () => {
// The provider logged the last prompt before the crash; the host journalled it on recovery.
const items = [
journalItem('ask', 1, 100, said('user', 'Fix the parser')),
journalItem('reply', 2, 300, said('assistant', 'Working on it')),
journalItem('steer', 3, 310, said('user', 'Keep the old API')),
journalItem('recovered', 4, 200, said('user', 'Also add a test'), {
recovered: true,
recoveredAt: 400
})
]
expect(drawn(items)).toEqual(['ask', 'reply', 'steer', 'recovered'])
// The host's outline, published to every client, uses the same order.
expect(projectAgentSessionConversationOutline(items, []).map(({ itemId }) => itemId)).toEqual([
'ask',
'steer',
'recovered'
])
})
it("draws one write's questions by their place in it, whatever order they arrive in", () => {
// One Codex ask: a single journal write, one sequence and one timestamp. Neither
// the ids nor the arrival order below spell the order it was asked in.
const scope = journalItem('q:scope', 5, 500, answered('Which files?'))
const priority = journalItem('q:priority', 5, 500, answered('What matters?'), {
sequenceIndex: 1
})
const deadline = journalItem('q:deadline', 5, 500, answered('When?'), { sequenceIndex: 2 })
expect(drawn([priority, deadline, scope])).toEqual(['q:scope', 'q:priority', 'q:deadline'])
})
it('keeps a send the journal does not hold yet below every row it does', () => {
// The composer's clock can trail the host's; the unsent message still reads last.
const outbox = [queued('queued', 'One more thing', 150)]
const items = [
journalItem('ask', 1, 100, said('user', 'Fix the parser')),
journalItem('reply', 2, 200, said('assistant', 'Working on it'))
]
expect(drawn(items, outbox)).toEqual(['ask', 'reply', agentJournalSubmissionKey('queued')])
})
it('keeps a send the journal recorded and the provider refused at its journal place', () => {
// The agent kept writing after the refused steer; its Retry stays with the composer.
const refused = agentJournalSubmissionKey('steer')
const items = [
journalItem('ask', 1, 100, said('user', 'Fix the parser')),
journalItem(refused, 2, 200, said('user', 'Keep the old API')),
journalItem('reply', 3, 300, said('assistant', 'Done.'))
]
const submissions: AgentJournalSubmission[] = [
{
clientMessageId: 'steer',
fence: 1,
payloadFingerprint: 'fingerprint',
dispatchState: 'rejected',
providerItemId: null,
reason: 'provider_refused',
submittedAt: 200,
resolvedAt: 250
}
]
expect(
drawn(
items,
[queued('steer', 'Keep the old API', 190), queued('next', 'And docs', 400)],
submissions
)
).toEqual(['ask', refused, 'reply', agentJournalSubmissionKey('next')])
})
})
@@ -75,8 +75,8 @@ export function truncateOutlinePreview(text: string, maxChars: number): string {
/** User messages that draw a transcript row, in transcript order. Projected over the
* whole journal, not user items alone: whether a user row survives depends on its
* neighbours (a harness sidecar folds into the turn before it) and its order on
* when it was observed. Previews are uncut; the reply bound owns length. */
* neighbours (a harness sidecar folds into the turn before it), and its order is
* its journal position. Previews are uncut; the reply bound owns length. */
export function projectAgentSessionConversationOutline(
items: readonly AgentJournalRenderItem[],
submissions: readonly AgentJournalSubmission[]
@@ -0,0 +1,36 @@
// The journal's own order: the sequence of the row that created an item, then the
// item's place among that row's writes. Time never enters it — a row recovered
// after a crash carries the provider's earlier clock at a later sequence.
import type { AgentJournalPosition, AgentJournalRenderItem } from './agent-session-journal-types'
import type { NativeChatMessage } from './native-chat-types'
type PositionedItem = Pick<AgentJournalRenderItem, 'sequence' | 'sequenceIndex'>
/** What every transcript row projected from a journal item carries, so a row
* that changes shape (a question once answered) keeps its identity and place. */
export function agentJournalItemRowOrigin(
item: AgentJournalRenderItem
): Pick<NativeChatMessage, 'id' | 'timestamp' | 'source' | 'journalPosition'> {
return {
id: item.itemId,
timestamp: item.observedAt,
source: 'transcript',
journalPosition: agentJournalItemPosition(item)
}
}
export function agentJournalItemPosition(item: PositionedItem): AgentJournalPosition {
return { sequence: item.sequence, index: item.sequenceIndex ?? 0 }
}
export function compareAgentJournalPositions(
a: AgentJournalPosition,
b: AgentJournalPosition
): number {
return a.sequence - b.sequence || a.index - b.index
}
export function compareAgentJournalItems(a: PositionedItem, b: PositionedItem): number {
return a.sequence - b.sequence || (a.sequenceIndex ?? 0) - (b.sequenceIndex ?? 0)
}
@@ -289,6 +289,7 @@ export const AgentJournalRenderItemSchema = z.object({
revision: z.number().int(),
body: AgentJournalItemBodySchema,
sequence: z.number().int(),
sequenceIndex: z.number().int().nonnegative().optional(),
observedAt: z.number(),
recovered: z.literal(true).optional(),
recoveredAt: z.number().optional(),
+10
View File
@@ -327,6 +327,13 @@ export type AgentJournalProducerLinkage = {
attempt?: number
}
/** Where the journal placed an item: the sequence of the row that created it,
* then its place among that row's writes. The timeline's only ordering key. */
export type AgentJournalPosition = {
sequence: number
index: number
}
/** One reduced timeline entry. `sequence` orders the list; `observedAt` is the
* provider's own clock and may sort earlier than a later sequence when the row
* was recovered after a crash. */
@@ -335,6 +342,9 @@ export type AgentJournalRenderItem = AgentJournalProducerLinkage & {
revision: number
body: AgentJournalItemBody
sequence: number
/** Place among the writes of the row at `sequence`, which one lifecycle batch
* shares across every item it creates. Absent ⇒ 0, and on a host that predates it. */
sequenceIndex?: number
observedAt: number
/** Set when the row was appended by crash reconciliation rather than live. */
recovered?: true
@@ -4,6 +4,7 @@
// of exactly what the renderer's transcript runs, not a second reading of it.
import type { NativeChatMessage } from './native-chat-types'
import { compareAgentJournalPositions } from './agent-session-journal-position'
import { stripNoiseMessages } from './native-chat-noise'
import { foldToolMessages } from './native-chat-tool-fold'
@@ -27,11 +28,32 @@ export function compareNativeChatMessagesByTime(
return 0
}
/** Rows the journal holds read in the journal's own order, never its clock: a
* batch shares one timestamp, and a row recovered after a crash carries an
* earlier one. A row not in the journal yet — a send still in the outbox — was
* made after everything the journal holds, so it follows them; only such rows,
* and terminal-backed transcripts, which have no journal, order by time. */
export function compareNativeChatTranscriptMessages(
a: NativeChatMessage,
b: NativeChatMessage
): number {
if (a.journalPosition && b.journalPosition) {
return compareAgentJournalPositions(a.journalPosition, b.journalPosition)
}
if (a.journalPosition || b.journalPosition) {
return a.journalPosition ? -1 : 1
}
return compareNativeChatMessagesByTime(a, b)
}
/** `compare` lets the renderer order its own tail rows (streaming, optimistic
* sends), which never exist on the host. */
export function projectNativeChatTranscriptMessages(
messages: readonly NativeChatMessage[],
compare: (a: NativeChatMessage, b: NativeChatMessage) => number = compareNativeChatMessagesByTime
compare: (
a: NativeChatMessage,
b: NativeChatMessage
) => number = compareNativeChatTranscriptMessages
): NativeChatMessage[] {
// Not `toSorted`: mobile's Hermes lacks it, and src/shared must stay loadable there.
return stripNoiseMessages(foldToolMessages(Array.from(messages).sort(compare)))
+7 -1
View File
@@ -10,7 +10,10 @@ import type {
AgentSessionBackgroundTask,
AgentSessionBackgroundTaskRunState
} from './agent-session-background-task-wire'
import type { AgentJournalMessageSendMode } from './agent-session-journal-types'
import type {
AgentJournalMessageSendMode,
AgentJournalPosition
} from './agent-session-journal-types'
import type { AgentType } from './agent-status-types'
import type { NativeChatToolMetadata } from './native-chat-tool-identity'
@@ -199,6 +202,9 @@ export type NativeChatMessage = {
turnId?: string
/** How a user message was delivered when it was not an ordinary prompt. */
sentAs?: AgentJournalMessageSendMode
/** Set only by the structured projection, on rows the journal holds, and ranks
* them ahead of time. Terminal-backed messages never carry it, and worker reads strip it. */
journalPosition?: AgentJournalPosition
}
export const NATIVE_CHAT_TURN_LIFECYCLE_STATES = ['working', 'completed', 'interrupted'] as const
@@ -1,5 +1,6 @@
import type { AgentJournalRenderItem, AgentJournalSubmission } from './agent-session-journal-types'
import { agentJournalSubmissionKey } from './agent-session-journal-item-key'
import { agentJournalItemPosition } from './agent-session-journal-position'
import type { NativeChatMessage } from './native-chat-types'
import {
reconcileStructuredAgentSessionOutbox,
@@ -20,18 +21,32 @@ export function projectStructuredAgentSessionMessages(
.filter((submission) => submission.dispatchState === 'rejected')
.map((submission) => agentJournalSubmissionKey(submission.clientMessageId))
)
const visibleItems = items.filter((item) => !rejected.has(item.itemId))
const visibleItems: AgentJournalRenderItem[] = []
const refused = new Map<string, AgentJournalRenderItem>()
for (const item of items) {
if (rejected.has(item.itemId)) {
refused.set(item.itemId, item)
} else {
visibleItems.push(item)
}
}
const journalled = new Set(visibleItems.map((item) => item.itemId))
return [
...projectItems(visibleItems),
...optimistic
.filter((entry) => !journalled.has(agentJournalSubmissionKey(entry.clientMessageId)))
.map((entry): NativeChatMessage => ({
id: agentJournalSubmissionKey(entry.clientMessageId),
role: 'user',
source: 'transcript',
timestamp: entry.queuedAt,
blocks: entry.body.blocks
}))
.map((entry): NativeChatMessage => {
const id = agentJournalSubmissionKey(entry.clientMessageId)
const recorded = refused.get(id)
return {
id,
role: 'user',
source: 'transcript',
timestamp: entry.queuedAt,
blocks: entry.body.blocks,
// A send the journal recorded before refusing it keeps its place there.
...(recorded ? { journalPosition: agentJournalItemPosition(recorded) } : {})
}
})
]
}
@@ -11,6 +11,7 @@ import {
type AgentJournalTurnOutcome
} from './agent-session-journal-types'
import { isRootAgentJournalItem } from './agent-session-journal-producer'
import { agentJournalItemRowOrigin } from './agent-session-journal-position'
import {
AGENT_STATUS_TOOL_INPUT_MAX_LENGTH,
AGENT_STATUS_TOOL_NAME_MAX_LENGTH
@@ -170,11 +171,9 @@ export function projectStructuredItemToNativeChat(
const sentAs = item.body.kind === 'message' ? item.body.sentAs : undefined
const message: NativeChatMessage | null = projected
? {
id: item.itemId,
...agentJournalItemRowOrigin(item),
role: projected.role,
blocks: projected.blocks,
timestamp: item.observedAt,
source: 'transcript',
// A send mode this build cannot name renders as an ordinary message.
...(sentAs !== undefined && isAgentJournalMessageSendMode(sentAs) ? { sentAs } : {})
}
@@ -54,6 +54,37 @@ function hydrationPage(
}
describe('structured agent session reducer', () => {
it("orders one journal write's items by their place in it, whatever order they arrive in", () => {
const at = (id: string, sequence: number, sequenceIndex: number): AgentJournalRenderItem => ({
...item(id, sequence),
...(sequenceIndex > 0 ? { sequenceIndex } : {})
})
const paged = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
type: 'history-page',
page: hydrationPage([
item('before', 4),
at('third', 5, 2),
at('first', 5, 0),
at('second', 5, 1)
])
})
expect(paged.items.map(({ itemId }) => itemId)).toEqual(['before', 'first', 'second', 'third'])
const live = reduceStructuredAgentSession(paged, {
type: 'event',
event: {
type: 'batch',
sessionId: 'session-a',
batch: {
cursor: { epoch: 'epoch-a', sequence: 6 },
items: [at('next-b', 6, 1), at('next-a', 6, 0)],
removedItemIds: [],
submissions: []
}
}
})
expect(live.items.map(({ itemId }) => itemId).slice(-2)).toEqual(['next-a', 'next-b'])
})
it('applies an additive targeted-stop capability update without journal churn', () => {
const backgroundTasks = {
state: 'monitoring' as const,
@@ -12,6 +12,7 @@ import type {
} from './agent-session-wire'
import { backgroundTaskStatesEqual } from './agent-session-background-task-state-equality'
import { agentJournalSubmissionKey } from './agent-session-journal-item-key'
import { compareAgentJournalItems } from './agent-session-journal-position'
import { readAgentJournalTurn } from './agent-session-turn-record'
/** The last host clock sample: `hostNow - receivedAt` is the client's skew from the host,
@@ -85,7 +86,7 @@ function replacePage(
epoch: page.epoch,
cursor: page.liveCursor ?? page.window.nextCursor,
fence,
items: [...page.items].sort((left, right) => left.sequence - right.sequence),
items: [...page.items].sort(compareAgentJournalItems),
submissions: page.submissions,
retainedItemLimit: Math.max(MAX_RETAINED_ITEMS, page.items.length),
hasOlder: page.hasOlder,
@@ -114,7 +115,7 @@ function mergeItems(
byId.set(item.itemId, item)
}
}
return [...byId.values()].sort((left, right) => left.sequence - right.sequence)
return [...byId.values()].sort(compareAgentJournalItems)
}
/**