From c6caab5fb76f0487de4451997cd41bb11dadcd75 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Wed, 9 Sep 2026 21:29:30 -0700 Subject: [PATCH] Retain turn attribution for loaded chat history --- .../structured-agent-session-reducer.ts | 40 +++-- ...gent-session-turn-timing-retention.test.ts | 160 ++++++++++++++++++ 2 files changed, 187 insertions(+), 13 deletions(-) create mode 100644 src/shared/structured-agent-session-turn-timing-retention.test.ts diff --git a/src/shared/structured-agent-session-reducer.ts b/src/shared/structured-agent-session-reducer.ts index 10897c3f61f..075e4d8bb09 100644 --- a/src/shared/structured-agent-session-reducer.ts +++ b/src/shared/structured-agent-session-reducer.ts @@ -11,6 +11,7 @@ import type { AgentSessionSubscribeEvent, AgentSessionTurnActivity } from './agent-session-wire' +import { agentJournalSubmissionKey } from './agent-session-journal-item-key' export type StructuredAgentSessionState = { epoch: string | null @@ -123,15 +124,25 @@ function mergeItems( function mergeSubmissions( current: readonly AgentJournalSubmission[], - incoming: readonly AgentJournalSubmission[] + incoming: readonly AgentJournalSubmission[], + items: readonly AgentJournalRenderItem[] ): AgentJournalSubmission[] { const byId = new Map(current.map((submission) => [submission.clientMessageId, submission])) for (const submission of incoming) { byId.set(submission.clientMessageId, submission) } - return [...byId.values()] - .sort((left, right) => left.submittedAt - right.submittedAt) - .slice(-MAX_RETAINED_SUBMISSIONS) + const sorted = [...byId.values()].sort((left, right) => left.submittedAt - right.submittedAt) + const itemIds = new Set( + items + .filter((item) => item.body.kind === 'message' && item.body.role === 'user') + .map((item) => item.itemId) + ) + // Loaded user messages need their provider alias for durable turn attribution. + return sorted.filter( + (submission, index) => + index >= sorted.length - MAX_RETAINED_SUBMISSIONS || + itemIds.has(agentJournalSubmissionKey(submission.clientMessageId)) + ) } export function reduceStructuredAgentSession( @@ -184,7 +195,7 @@ export function reduceStructuredAgentSession( fence: action.page.fence ?? null, items: action.page.items, submissions: sameEpoch - ? mergeSubmissions(state.submissions, action.page.submissions) + ? mergeSubmissions(state.submissions, action.page.submissions, action.page.items) : action.page.submissions, hasOlder: action.page.hasOlder, status: 'ready', @@ -202,10 +213,11 @@ export function reduceStructuredAgentSession( if (state.epoch !== action.requestedEpoch || action.page.epoch !== action.requestedEpoch) { return state } + const items = mergeItems(state.items, action.page.items, action.page.removedItemIds) return { ...state, - items: mergeItems(state.items, action.page.items, action.page.removedItemIds), - submissions: mergeSubmissions(state.submissions, action.page.submissions), + items, + submissions: mergeSubmissions(state.submissions, action.page.submissions, items), hasOlder: action.page.hasOlder } } @@ -246,16 +258,18 @@ export function reduceStructuredAgentSession( ) { return state } + const items = journalUnchanged + ? state.items + : mergeItems(state.items, event.batch.items, event.batch.removedItemIds) return { ...state, cursor: event.batch.cursor, fence: event.fence ?? state.fence, - items: journalUnchanged - ? state.items - : mergeItems(state.items, event.batch.items, event.batch.removedItemIds), - submissions: journalUnchanged - ? state.submissions - : mergeSubmissions(state.submissions, event.batch.submissions), + items, + submissions: + event.batch.submissions.length === 0 && event.batch.removedItemIds.length === 0 + ? state.submissions + : mergeSubmissions(state.submissions, event.batch.submissions, items), status: 'ready', error: undefined, handoff: event.handoff ?? state.handoff, diff --git a/src/shared/structured-agent-session-turn-timing-retention.test.ts b/src/shared/structured-agent-session-turn-timing-retention.test.ts new file mode 100644 index 00000000000..7b08e3ad691 --- /dev/null +++ b/src/shared/structured-agent-session-turn-timing-retention.test.ts @@ -0,0 +1,160 @@ +import { describe, expect, it } from 'vitest' +import type { AgentJournalRenderItem, AgentJournalSubmission } from './agent-session-journal-types' +import type { AgentSessionHistoryPage } from './agent-session-wire' +import { + EMPTY_STRUCTURED_AGENT_SESSION, + reduceStructuredAgentSession +} from './structured-agent-session-reducer' +import { selectStructuredAgentSettledTurns } from './structured-agent-session-turn-timing' + +function submission(index: number): AgentJournalSubmission { + return { + clientMessageId: `user-${index}`, + providerItemId: `codex:thread:turn-${index}:0`, + fence: 1, + payloadFingerprint: 'fingerprint', + dispatchState: 'accepted', + reason: null, + submittedAt: index, + resolvedAt: index + } +} + +function turnItems(index: number): AgentJournalRenderItem[] { + return [ + { + itemId: `orca:user-${index}`, + revision: 1, + sequence: index * 2 + 1, + observedAt: 1_000, + body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'hello' }] } + }, + { + itemId: `legacy:codex:session:turn-${index}`, + revision: 2, + sequence: index * 2 + 2, + observedAt: 1_000, + body: { + kind: 'turn', + turnId: `turn-${index}`, + userItemId: submission(index).providerItemId!, + state: 'completed', + startedAt: 1_000, + completedAt: 8_000 + } + } + ] +} + +function page(indices: number[]): AgentSessionHistoryPage { + return { + sessionId: 'session', + epoch: 'epoch', + direction: 'tail', + items: indices.flatMap(turnItems), + submissions: indices.map(submission), + removedItemIds: [], + window: { oldest: null, newest: null, nextCursor: { epoch: 'epoch', sequence: 1_000 } }, + hasOlder: false, + hasNewer: false + } +} + +function loadedHistory() { + let state = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, { + type: 'event', + event: { + type: 'snapshot', + sessionId: 'session', + fence: 1, + page: page(Array.from({ length: 64 }, (_, index) => index + 193)) + } + }) + for (const first of [129, 65, 1]) { + state = reduceStructuredAgentSession(state, { + type: 'older-page', + requestedEpoch: 'epoch', + page: page(Array.from({ length: 64 }, (_, index) => index + first)) + }) + } + return state +} + +describe('durable turn attribution across paginated history', () => { + it('keeps an older page duration after the recent submission budget fills', () => { + const state = reduceStructuredAgentSession(loadedHistory(), { + type: 'older-page', + requestedEpoch: 'epoch', + page: page([0]) + }) + expect( + selectStructuredAgentSettledTurns(state.items, state.submissions).get('orca:user-0') + ).toEqual({ startedAt: 1_000, workedSeconds: 7 }) + expect(state.submissions).toHaveLength(257) + + const next = reduceStructuredAgentSession(state, { + type: 'event', + event: { + type: 'batch', + sessionId: 'session', + batch: { + cursor: { epoch: 'epoch', sequence: 1_001 }, + items: turnItems(257), + submissions: [submission(257)], + removedItemIds: [] + } + } + }) + expect(selectStructuredAgentSettledTurns(next.items, next.submissions).size).toBe(258) + const streamed = reduceStructuredAgentSession(next, { + type: 'event', + event: { + type: 'batch', + sessionId: 'session', + batch: { + cursor: { epoch: 'epoch', sequence: 1_002 }, + items: [{ ...turnItems(257)[1]!, revision: 3 }], + submissions: [], + removedItemIds: [] + } + } + }) + expect(streamed.submissions).toBe(next.submissions) + }) + + it('drops an old alias once rewind removes its user item', () => { + const loaded = reduceStructuredAgentSession(loadedHistory(), { + type: 'older-page', + requestedEpoch: 'epoch', + page: page([0]) + }) + const removed = reduceStructuredAgentSession(loaded, { + type: 'event', + event: { + type: 'batch', + sessionId: 'session', + batch: { + cursor: { epoch: 'epoch', sequence: 1_001 }, + items: [], + submissions: [], + removedItemIds: turnItems(0).map((item) => item.itemId) + } + } + }) + expect(removed.submissions).toHaveLength(256) + expect(removed.submissions.some((entry) => entry.clientMessageId === 'user-0')).toBe(false) + + const reset = reduceStructuredAgentSession(loaded, { + type: 'event', + event: { + type: 'reset', + sessionId: 'session', + fence: 2, + page: page([256]), + reset: 'epoch_changed' + } + }) + expect(reset.submissions).toEqual([submission(256)]) + expect(selectStructuredAgentSettledTurns(reset.items, reset.submissions).size).toBe(1) + }) +})