Skip to content
2 changes: 2 additions & 0 deletions config/tsconfig.node.json
Original file line number Diff line number Diff line change
Expand Up @@ -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/**/*",
Expand Down
257 changes: 257 additions & 0 deletions src/main/codex/codex-structured-question-order.test.ts
Original file line number Diff line number Diff line change
@@ -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)']
])
})
})
33 changes: 33 additions & 0 deletions src/main/native-chat/agent-session-journal/journal-reducer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
Expand Down
38 changes: 17 additions & 21 deletions src/main/native-chat/agent-session-journal/journal-reducer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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 },
Expand Down
Loading
Loading