Retain turn attribution for loaded chat history

This commit is contained in:
Merge Sim
2026-09-09 21:29:30 -07:00
parent 1274bdac40
commit c6caab5fb7
2 changed files with 187 additions and 13 deletions
+27 -13
View File
@@ -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,
@@ -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)
})
})