diff --git a/config/tsconfig.node.json b/config/tsconfig.node.json index c42b937a467..4f7cc41bcc1 100644 --- a/config/tsconfig.node.json +++ b/config/tsconfig.node.json @@ -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/**/*", diff --git a/src/main/codex/codex-structured-question-order.test.ts b/src/main/codex/codex-structured-question-order.test.ts new file mode 100644 index 00000000000..35758327076 --- /dev/null +++ b/src/main/codex/codex-structured-question-order.test.ts @@ -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(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 { + 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 { + 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)'] + ]) + }) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-reducer.test.ts b/src/main/native-chat/agent-session-journal/journal-reducer.test.ts index 825c5aaf347..eea740e754d 100644 --- a/src/main/native-chat/agent-session-journal/journal-reducer.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-reducer.test.ts @@ -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 }, diff --git a/src/main/native-chat/agent-session-journal/journal-reducer.ts b/src/main/native-chat/agent-session-journal/journal-reducer.ts index 3fbee85a417..cf0f3739acb 100644 --- a/src/main/native-chat/agent-session-journal/journal-reducer.ts +++ b/src/main/native-chat/agent-session-journal/journal-reducer.ts @@ -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 }, diff --git a/src/main/native-chat/agent-session-journal/journal-render-item.ts b/src/main/native-chat/agent-session-journal/journal-render-item.ts index 6f41ff3e862..516086eed3d 100644 --- a/src/main/native-chat/agent-session-journal/journal-render-item.ts +++ b/src/main/native-chat/agent-session-journal/journal-render-item.ts @@ -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 } : {}), diff --git a/src/main/native-chat/agent-session-wire/agent-session-journal-batch.ts b/src/main/native-chat/agent-session-wire/agent-session-journal-batch.ts index 57f2291e476..b70a18a8386 100644 --- a/src/main/native-chat/agent-session-wire/agent-session-journal-batch.ts +++ b/src/main/native-chat/agent-session-wire/agent-session-journal-batch.ts @@ -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: { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-accept-then-deliver.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-accept-then-deliver.test.ts index b9e1b49c213..c280d062488 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-accept-then-deliver.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-accept-then-deliver.test.ts @@ -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')) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts index 69e6d029d60..959a02919a2 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-test-data.ts @@ -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) +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-send-restarts-failed-start.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-send-restarts-failed-start.test.ts index 211d8684e15..c3f192ff3d3 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-send-restarts-failed-start.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-send-restarts-failed-start.test.ts @@ -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') diff --git a/src/main/runtime/orchestration/worker-transcript-payload.test.ts b/src/main/runtime/orchestration/worker-transcript-payload.test.ts index c843fe1cfa7..802b724a2aa 100644 --- a/src/main/runtime/orchestration/worker-transcript-payload.test.ts +++ b/src/main/runtime/orchestration/worker-transcript-payload.test.ts @@ -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) diff --git a/src/main/runtime/orchestration/worker-transcript-payload.ts b/src/main/runtime/orchestration/worker-transcript-payload.ts index b4365e259c2..684dfc771e2 100644 --- a/src/main/runtime/orchestration/worker-transcript-payload.ts +++ b/src/main/runtime/orchestration/worker-transcript-payload.ts @@ -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)) diff --git a/src/renderer/src/components/native-chat/NativeChatMessageList.task-list-frames.test.tsx b/src/renderer/src/components/native-chat/NativeChatMessageList.task-list-frames.test.tsx index 1408342336a..4be1ea7cdb3 100644 --- a/src/renderer/src/components/native-chat/NativeChatMessageList.task-list-frames.test.tsx +++ b/src/renderer/src/components/native-chat/NativeChatMessageList.task-list-frames.test.tsx @@ -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', diff --git a/src/renderer/src/components/native-chat/native-chat-rail-outline-parity.test.ts b/src/renderer/src/components/native-chat/native-chat-rail-outline-parity.test.ts index 3057be1f908..385899ae91f 100644 --- a/src/renderer/src/components/native-chat/native-chat-rail-outline-parity.test.ts +++ b/src/renderer/src/components/native-chat/native-chat-rail-outline-parity.test.ts @@ -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 } ]) }) }) diff --git a/src/renderer/src/components/native-chat/native-chat-session-assembler.ts b/src/renderer/src/components/native-chat/native-chat-session-assembler.ts index eb61380967f..674e14d6bf8 100644 --- a/src/renderer/src/components/native-chat/native-chat-session-assembler.ts +++ b/src/renderer/src/components/native-chat/native-chat-session-assembler.ts @@ -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) } /** diff --git a/src/renderer/src/components/native-chat/structured-agent-question-projection.ts b/src/renderer/src/components/native-chat/structured-agent-question-projection.ts index b56fa1c7ab0..9469c443d0a 100644 --- a/src/renderer/src/components/native-chat/structured-agent-question-projection.ts +++ b/src/renderer/src/components/native-chat/structured-agent-question-projection.ts @@ -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 }] } } diff --git a/src/renderer/src/components/native-chat/structured-agent-session-transcript-order.test.ts b/src/renderer/src/components/native-chat/structured-agent-session-transcript-order.test.ts new file mode 100644 index 00000000000..ea9b3a0cf3b --- /dev/null +++ b/src/renderer/src/components/native-chat/structured-agent-session-transcript-order.test.ts @@ -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 { + 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')]) + }) +}) diff --git a/src/shared/agent-session-conversation-outline.ts b/src/shared/agent-session-conversation-outline.ts index bc4fa3b40ad..2b6bb846adb 100644 --- a/src/shared/agent-session-conversation-outline.ts +++ b/src/shared/agent-session-conversation-outline.ts @@ -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[] diff --git a/src/shared/agent-session-journal-position.ts b/src/shared/agent-session-journal-position.ts new file mode 100644 index 00000000000..9e317f1ee71 --- /dev/null +++ b/src/shared/agent-session-journal-position.ts @@ -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 + +/** 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 { + 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) +} diff --git a/src/shared/agent-session-journal-schemas.ts b/src/shared/agent-session-journal-schemas.ts index 04f7183eb3f..438352b1ea3 100644 --- a/src/shared/agent-session-journal-schemas.ts +++ b/src/shared/agent-session-journal-schemas.ts @@ -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(), diff --git a/src/shared/agent-session-journal-types.ts b/src/shared/agent-session-journal-types.ts index 0a6cb655cd8..f68361fb13c 100644 --- a/src/shared/agent-session-journal-types.ts +++ b/src/shared/agent-session-journal-types.ts @@ -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 diff --git a/src/shared/native-chat-transcript-projection.ts b/src/shared/native-chat-transcript-projection.ts index 26b21824711..45f2488f4aa 100644 --- a/src/shared/native-chat-transcript-projection.ts +++ b/src/shared/native-chat-transcript-projection.ts @@ -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))) diff --git a/src/shared/native-chat-types.ts b/src/shared/native-chat-types.ts index 24bf8652fb2..45c1a3f6ed4 100644 --- a/src/shared/native-chat-types.ts +++ b/src/shared/native-chat-types.ts @@ -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 diff --git a/src/shared/structured-agent-session-message-projection.ts b/src/shared/structured-agent-session-message-projection.ts index 6219538fbec..334cde7cfe8 100644 --- a/src/shared/structured-agent-session-message-projection.ts +++ b/src/shared/structured-agent-session-message-projection.ts @@ -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() + 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) } : {}) + } + }) ] } diff --git a/src/shared/structured-agent-session-projection.ts b/src/shared/structured-agent-session-projection.ts index 023c9ef0970..9dcd952d271 100644 --- a/src/shared/structured-agent-session-projection.ts +++ b/src/shared/structured-agent-session-projection.ts @@ -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 } : {}) } diff --git a/src/shared/structured-agent-session-reducer.test.ts b/src/shared/structured-agent-session-reducer.test.ts index aa90749ffdf..a374c794d7f 100644 --- a/src/shared/structured-agent-session-reducer.test.ts +++ b/src/shared/structured-agent-session-reducer.test.ts @@ -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, diff --git a/src/shared/structured-agent-session-reducer.ts b/src/shared/structured-agent-session-reducer.ts index 2f839f93330..e47ed5ad57b 100644 --- a/src/shared/structured-agent-session-reducer.ts +++ b/src/shared/structured-agent-session-reducer.ts @@ -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) } /**