Keep earlier turns through a Codex rewind and count a mid-turn attach from the real start

Findings from an independent adversarial review of the typed turn record:

- A Codex rewind adopted the provider's item list as the new epoch, and the
  provider never returns the host's own turn rows, so every duration before
  the rewind point vanished. The host's turn rows are now spliced back beside
  the item each followed, and recovery no longer expects the provider to
  prove rows it never owned.
- The epoch row was stamped with the current schema version, so an older host
  latched read-only at row 1 of every new session, defeating the mixed
  version design. It carries no body and stays at v2; a stored-row test now
  reads SQLite directly, because the reader upcasts every row on read.
- A send Codex folds into a running turn shares the opening prompt's provider
  key, and the alias map credited the duration to the later prompt. The
  earliest submission naming a key now wins.
- The live counter anchored on first sight, so a client attaching mid-turn
  counted from zero. Published frames now carry the host's clock, the reducer
  keeps the last sample with its local receipt time, and both clients anchor
  on how long the host says the turn has run.
This commit is contained in:
Merge Sim
2026-09-09 22:09:25 -07:00
parent 4ae26deb11
commit e2b3abe9e4
27 changed files with 639 additions and 132 deletions
@@ -272,7 +272,7 @@ export function useMobileStructuredAgentSession(args: {
[state.items, state.submissions]
)
const turnId = activeStructuredAgentSessionTurnId(state.items)
const turnTiming = useMobileStructuredAgentTurnTiming(state.items, state.submissions, turnId)
const turnTiming = useMobileStructuredAgentTurnTiming(state, turnId)
const status = state.status === 'idle' ? 'idle' : state.status
const approvalPrompt = useMemo(
() => state.items.find(pendingStructuredApproval) ?? null,
@@ -65,7 +65,7 @@ export function useMobileStructuredAgentState(args: {
}
setSessionStates((current) => {
const previous = current.get(sessionKey) ?? EMPTY_STRUCTURED_AGENT_SESSION
const next = reduceStructuredAgentSession(previous, action)
const next = reduceStructuredAgentSession(previous, action, Date.now())
if (next === previous) {
return current
}
@@ -63,13 +63,15 @@ describe('useMobileStructuredAgentTurnTiming', () => {
function Harness({
items,
submissions = NO_SUBMISSIONS,
turnId
turnId,
hostClock
}: {
items: readonly AgentJournalRenderItem[]
submissions?: readonly AgentJournalSubmission[]
turnId: string | null
hostClock?: { hostNow: number; receivedAt: number }
}): null {
timing = useMobileStructuredAgentTurnTiming(items, submissions, turnId)
timing = useMobileStructuredAgentTurnTiming({ items, submissions, hostClock }, turnId)
return null
}
@@ -123,6 +125,30 @@ describe('useMobileStructuredAgentTurnTiming', () => {
act(() => renderer?.update(createElement(Harness, { items, turnId: null })))
expect(timing?.workingStartedAt).toBeNull()
// With a host clock that said the turn was 35s old 5s ago, the anchor sits
// 40s before first sight, wherever the client's absolute clock is.
vi.setSystemTime(CLIENT_NOW + 60_000)
const next = [
...items,
user('u3', 5),
lifecycle(
't3',
6,
{ state: 'running', startedAt: HOST_START + 150_000 },
HOST_START + 150_100
)
]
act(() =>
renderer?.update(
createElement(Harness, {
items: next,
turnId: 't3',
hostClock: { hostNow: HOST_START + 185_000, receivedAt: CLIENT_NOW + 55_000 }
})
)
)
expect(timing?.workingStartedAt).toBe(CLIENT_NOW + 60_000 - 40_000)
})
it('leaves the anchor null when an older host records no start', () => {
@@ -12,22 +12,40 @@ import {
type TurnAnchor = { turnId: string; startedAt: number | null }
/** The host's clock as last published, paired with the client clock at receipt. */
type HostClock = { hostNow: number; receivedAt: number }
/** The live turn's local-clock anchor. Null when its row carries no host start
* (older hosts), so local observation applies. */
function anchorRunningTurn(items: readonly AgentJournalRenderItem[], turnId: string): TurnAnchor {
function anchorRunningTurn(
items: readonly AgentJournalRenderItem[],
turnId: string,
hostClock: HostClock | null | undefined
): TurnAnchor {
const timing = selectStructuredAgentRunningTurnTiming(items, turnId)
return {
turnId,
startedAt: timing ? structuredAgentTurnLocalStartedAt(timing, Date.now()) : null
if (!timing) {
return { turnId, startedAt: null }
}
const now = Date.now()
// Advance the published host clock by the client time since receipt; both
// terms stay single-clock, so a mid-turn attach counts from the real start.
const hostNow = hostClock ? hostClock.hostNow + (now - hostClock.receivedAt) : undefined
return { turnId, startedAt: structuredAgentTurnLocalStartedAt(timing, now, hostNow) }
}
/** Host-recorded turn timing for the structured lane: settled durations straight
* off the journal, and a skew-free start for the live counter stamped once per
* turn so re-renders never move it. */
export function useMobileStructuredAgentTurnTiming(
items: readonly AgentJournalRenderItem[],
submissions: readonly AgentJournalSubmission[],
{
items,
submissions,
hostClock
}: {
items: readonly AgentJournalRenderItem[]
submissions: readonly AgentJournalSubmission[]
hostClock?: HostClock | null
},
turnId: string | null
): { settledTurns: ReadonlyMap<string, NativeChatSettledTurn>; workingStartedAt: number | null } {
const settledTurns = useMemo(
@@ -44,7 +62,7 @@ export function useMobileStructuredAgentTurnTiming(
return { settledTurns, workingStartedAt: null }
}
if (anchor?.turnId !== turnId) {
const next = anchorRunningTurn(items, turnId)
const next = anchorRunningTurn(items, turnId, hostClock)
setAnchor(next)
return { settledTurns, workingStartedAt: next.startedAt }
}
@@ -5,7 +5,7 @@
// repair marker the superseded epoch was carrying. Superseded rows are DELETED
// rather than retained — nothing would ever shed them.
import { AGENT_SESSION_JOURNAL_SCHEMA_VERSION } from '../../../shared/agent-session-journal-types'
import { journalRowSchemaVersion } from '../../../shared/agent-session-journal-types'
import type { AgentSessionProviderHandle } from '../../../shared/agent-session-journal-types'
import type Database from '../../sqlite/sync-database'
import type { JournalLoad } from './journal-open'
@@ -33,7 +33,8 @@ export function publishNewEpoch(input: {
kind: 'epoch',
reason: input.reason,
providerHandle: input.providerHandle,
v: AGENT_SESSION_JOURNAL_SCHEMA_VERSION,
// Carries no body: an older host must keep reading a turn-free session past row 1.
v: journalRowSchemaVersion([]),
epoch: input.epoch,
seq: 1,
fence: input.fence,
@@ -0,0 +1,66 @@
import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
import { agentJournalTurnBody } from '../../../shared/agent-session-turn-record'
import { openJournalDatabase } from './journal-database'
import { journalDatabaseFile } from './journal-paths'
import { createTrackedJournalOpener } from './journal-store-test-open'
// Which rows an older host can still read: only rows that carry a turn item
// are stamped with the version it does not know, and the epoch row never is.
// Read raw: the reader upcasts every row to the current version, so only the
// stored row_json says what an older build would see.
describe('journal row schema versions', () => {
let root = ''
const opener = createTrackedJournalOpener()
beforeEach(async () => {
root = await mkdtemp(join(tmpdir(), 'orca-journal-row-version-'))
})
afterEach(async () => {
await opener.closeAll()
await rm(root, { recursive: true, force: true })
})
it('stamps v3 only on rows that carry a turn item', async () => {
const journal = await opener.open({
identity: {
sessionId: 'session-1',
workspaceId: 'workspace-1',
hostId: 'local',
agent: 'codex',
providerHandle: { kind: 'codex', threadId: 'thread-1' }
},
now: () => 1_000,
journalDir: join(root, 'session-1')
})
const identity = { provider: 'orca' as const, clientMessageId: 'm1' }
await journal.appendItem(
identity,
{ kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'hi' }] },
{ fence: 1 }
)
await journal.appendItem(
{ provider: 'legacy', agent: 'codex', sessionId: 'session-1', recordId: 'turn-lifecycle:t1' },
agentJournalTurnBody({ turnId: 't1', state: 'running', startedAt: 1_000 }),
{ fence: 1 }
)
await journal.close()
const opened = openJournalDatabase(journalDatabaseFile(join(root, 'session-1')))
try {
const stored = opened.db
.prepare('SELECT row_json FROM journal_rows ORDER BY seq')
.all()
.map((row) => JSON.parse(String((row as { row_json: string }).row_json)))
.map((row: { kind: string; v: number }) => [row.kind, row.v])
expect(stored).toEqual([
['epoch', 2],
['item', 2],
['item', 3]
])
} finally {
opened.db.close()
}
})
})
@@ -0,0 +1,7 @@
import type { AgentJournalCursor } from '../../../shared/agent-session-journal-types'
import type { AgentSessionJournalBatch } from '../../../shared/agent-session-wire'
/** A batch that advances nothing: the carrier for fence, handoff, roster, and clock updates. */
export function emptyAgentSessionBatch(cursor: AgentJournalCursor): AgentSessionJournalBatch {
return { cursor, items: [], removedItemIds: [], submissions: [] }
}
@@ -31,9 +31,15 @@ export class StructuredAgentSessionBackgroundTaskChannel {
request
})
const backgroundTasks = this.state(request.sessionId)
return backgroundTasks === undefined
? result
: { ...result, page: { ...result.page, backgroundTasks } }
const hostNow = this.deps.now?.() ?? Date.now()
return {
...result,
page: {
...result.page,
hostNow,
...(backgroundTasks !== undefined ? { backgroundTasks } : {})
}
}
}
subscribe(input: AgentSessionSubscribeInput): () => void {
@@ -324,6 +324,8 @@ describe('send', () => {
const page = host.history({ sessionId: SESSION, direction: 'tail' })
expect(page.ok && page.page.items).toHaveLength(1)
expect(page.ok && page.page.fence).toBe(1)
// The injected host clock, so a client can anchor a live counter on it.
expect(page.page.hostNow).toBe(NOW)
expect(page.providerSession).toEqual({ key: 'session_id', id: THREAD })
})
@@ -72,7 +72,8 @@ export class StructuredAgentSessionHost {
})
private readonly subscribers = new AgentSessionSubscribers({
readCommands: (sessionId) => this.deps.adapter.readCommands?.(sessionId),
onJournalPublished: (sessionId, journal) => this.statusFeed.publish(sessionId, journal)
onJournalPublished: (sessionId, journal) => this.statusFeed.publish(sessionId, journal),
now: () => this.now()
})
private readonly tasks = new StructuredAgentSessionTaskQueue()
private readonly runtimeState: StructuredAgentSessionHostRuntimeState
@@ -420,6 +420,56 @@ describe('host rewind', () => {
expect(await host.rewind(caller, params(target))).toMatchObject({ ok: true })
})
it('keeps the host-stamped turn rows before the boundary through a Codex provider hydration', async () => {
expect(await host.attach(caller, hostTestAttachParams(null))).toMatchObject({ ok: true })
const message = (turnId: string) => ({
provider: 'codex' as const,
threadId: HOST_TEST_THREAD,
turnId,
ordinal: 0
})
const turnRow = (turnId: string) => ({
provider: 'legacy' as const,
agent: 'codex',
sessionId: HOST_TEST_SESSION,
recordId: `turn-lifecycle:${turnId}`
})
const keptTurn = {
kind: 'turn' as const,
turnId: 'kept',
state: 'completed' as const,
userItemId: agentJournalItemKey(message('kept')),
startedAt: HOST_TEST_NOW - 9_000,
completedAt: HOST_TEST_NOW - 4_000,
durationMs: 5_000
}
sink.appendItem(message('kept'), hostTestMessage('kept'))
sink.appendItem(turnRow('kept'), keptTurn)
sink.appendItem(message('drop'), hostTestMessage('drop'))
sink.appendItem(turnRow('drop'), { ...keptTurn, turnId: 'drop', durationMs: 1_000 })
sink.appendItem(message('tip'), { ...hostTestMessage('tip'), role: 'assistant' })
await host.flushStreamedEvents(HOST_TEST_SESSION)
// The provider preflight knows only its own items, never the host's turn rows.
const items = [{ identity: message('kept'), body: hostTestMessage('kept from provider') }]
rewind.mockImplementationOnce(async (input) => {
await input.onPrepared?.(items)
await input.onReverted?.()
return { ok: true, items }
})
expect(await host.rewind(caller, params(agentJournalItemKey(message('drop'))))).toMatchObject({
ok: true
})
expect(
host.journalSnapshot(HOST_TEST_SESSION).items.map(({ itemId, body }) => ({ itemId, body }))
).toEqual([
{ itemId: agentJournalItemKey(message('kept')), body: hostTestMessage('kept from provider') },
{ itemId: agentJournalItemKey(turnRow('kept')), body: keptTurn }
])
expect(store.getRecord(HOST_TEST_SESSION)?.rewind?.phase).toBe('completed')
})
it('recovers against the complete provider preflight when the local journal omitted an older turn', async () => {
const target = await seed()
const items = ['older', 'kept'].map((turnId) => ({
@@ -20,6 +20,7 @@ import { conversationCommandBlocked } from './structured-conversation-command-ad
import { rewindRefusal } from './structured-rewind-refusal'
import { persistRewindRecord, recoverStructuredRewind } from './structured-rewind-recovery'
import { replaceClaudeRewindOwner } from './structured-rewind-claude-owner'
import { mergeRetainedTurnRows } from './structured-rewind-retained-turns'
export async function rewindStructuredAgentSession(
context: StructuredAgentSessionMutationContext,
@@ -173,11 +174,14 @@ export async function rewindStructuredAgentSession(
fence: ctx.fence,
beforeTurnId: key.provider === 'codex' ? key.turnId : '',
onPrepared: async (items) => {
const retained = items.map(({ identity, body }) => ({
itemId: agentJournalItemKey(identity),
body,
observedAt: ctx.now()
}))
const retained = mergeRetainedTurnRows(
prepared.retained,
items.map(({ identity, body }) => ({
itemId: agentJournalItemKey(identity),
body,
observedAt: ctx.now()
}))
)
if (
retained.length > 10_000 ||
Buffer.byteLength(JSON.stringify(retained), 'utf8') >
@@ -216,11 +220,14 @@ export async function rewindStructuredAgentSession(
return rewindRefusal(reason)
}
const confirmed = provider.items
? provider.items.map(({ identity, body }) => ({
itemId: agentJournalItemKey(identity),
body,
observedAt: ctx.now()
}))
? mergeRetainedTurnRows(
prepared.retained,
provider.items.map(({ identity, body }) => ({
itemId: agentJournalItemKey(identity),
body,
observedAt: ctx.now()
}))
)
: prepared.retained
if (
Buffer.byteLength(JSON.stringify(confirmed), 'utf8') >
@@ -68,11 +68,57 @@ describe('AgentSessionSubscribers', () => {
submissions: []
},
fence: 7,
hostNow: expect.any(Number),
activity: null
}
])
})
it('stamps the host clock once per published frame', async () => {
const journal = await journals.open({
identity: {
sessionId: SESSION,
workspaceId: 'workspace-1',
hostId: 'local',
agent: 'codex',
providerHandle: { kind: 'codex', threadId: 'thread-1' }
},
journalDir: join(root, 'clock-journal')
})
let now = 1_000
const events: AgentSessionSubscribeEvent[] = []
const subscribers = new AgentSessionSubscribers({ now: () => (now += 1) })
const emit = (event: AgentSessionSubscribeEvent): void => {
events.push(event)
}
subscribers.open({ id: 'one', sessionId: SESSION, journal, fence: 1, emit })
subscribers.open({ id: 'two', sessionId: SESSION, journal, fence: 1, emit })
await journal.appendItem(
{ provider: 'orca', clientMessageId: 'clocked' },
{ kind: 'status', text: 'Clocked' },
{ fence: 1 }
)
subscribers.publish(SESSION, journal)
subscribers.handoff(SESSION, 1, { owner: 'native' } as AgentSessionHandoffStatus)
subscribers.reset(SESSION, journal, 'epoch_changed', 1)
expect(events.map((event) => ('hostNow' in event ? event.hostNow : null))).toEqual([
1_001, 1_002,
// Both subscribers of one publication read the same clock sample.
1_003, 1_003, 1_004, 1_004, 1_005, 1_005
])
expect(events.map((event) => event.type)).toEqual([
'snapshot',
'snapshot',
'batch',
'batch',
'batch',
'batch',
'reset',
'reset'
])
})
it('includes catalogs on reconnect and sends an idle checkpoint without journal work', async () => {
const journal = await journals.open({
identity: {
@@ -101,6 +147,7 @@ describe('AgentSessionSubscribers', () => {
type: 'batch',
sessionId: SESSION,
fence: 7,
hostNow: expect.any(Number),
commands,
batch: { cursor: journal.cursor(), items: [], removedItemIds: [], submissions: [] }
})
@@ -245,6 +292,7 @@ describe('AgentSessionSubscribers', () => {
submissions: []
},
fence: 2,
hostNow: expect.any(Number),
handoff
})
})
@@ -284,6 +332,7 @@ describe('AgentSessionSubscribers', () => {
sessionId: SESSION,
batch: { cursor, items: [], removedItemIds: [], submissions: [] },
fence: 2,
hostNow: expect.any(Number),
backgroundTasks
})
@@ -330,6 +379,7 @@ describe('AgentSessionSubscribers', () => {
sessionId: SESSION,
batch: { cursor, items: [], removedItemIds: [], submissions: [] },
fence: 1,
hostNow: expect.any(Number),
activity: { turnId: 'turn-1', text: 'Inspecting the session wire' }
})
@@ -17,6 +17,7 @@ import {
type AgentSessionTurnActivity
} from '../../../shared/agent-session-wire'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import { emptyAgentSessionBatch } from './agent-session-empty-batch'
import {
createAgentSessionCatchUpReader,
readAgentSessionHydrationPage
@@ -44,6 +45,8 @@ export type AgentSessionSubscribersHooks = {
/** Fires after any publication that can change journal content, whether or not anyone
* is subscribed to the transcript: session lists project status from this same edge. */
onJournalPublished?: (sessionId: string, journal: AgentSessionJournal) => void
/** Host wall clock, stamped once per published frame as `hostNow`. */
now?: () => number
}
export class AgentSessionSubscribers {
@@ -76,8 +79,9 @@ export class AgentSessionSubscribers {
session.set(input.id, subscriber)
this.bySession.set(input.sessionId, session)
const hostNow = this.now()
if (input.cursor) {
this.deliver(subscriber, input.journal, input.handoff, true, input.backgroundTasks)
this.deliver(subscriber, input.journal, hostNow, input.handoff, true, input.backgroundTasks)
} else {
const page = readAgentSessionHydrationPage(input.journal, input.fence)
this.emit(subscriber, {
@@ -85,6 +89,7 @@ export class AgentSessionSubscribers {
sessionId: input.sessionId,
page,
fence: input.fence,
hostNow,
...(input.handoff ? { handoff: input.handoff } : {}),
...(input.backgroundTasks !== undefined ? { backgroundTasks: input.backgroundTasks } : {}),
...this.activityField(input.sessionId)
@@ -121,8 +126,9 @@ export class AgentSessionSubscribers {
this.activityBySession.delete(sessionId)
}
}
const hostNow = this.now()
for (const subscriber of this.subscribers(sessionId)) {
this.deliver(subscriber, journal, undefined, false, undefined, activity)
this.deliver(subscriber, journal, hostNow, undefined, false, undefined, activity)
}
if (activity === undefined) {
this.hooks.onJournalPublished?.(sessionId, journal)
@@ -138,21 +144,7 @@ export class AgentSessionSubscribers {
fence: number,
backgroundTasks?: AgentSessionBackgroundTaskState | null
): void {
const page = readAgentSessionHydrationPage(journal, fence)
for (const subscriber of this.subscribers(sessionId)) {
this.emit(subscriber, {
type: 'reset',
sessionId,
reset: reason,
page,
fence,
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
...this.activityField(sessionId)
})
subscriber.cursor = page.liveCursor ?? page.window.nextCursor
subscriber.fence = fence
}
this.hooks.onJournalPublished?.(sessionId, journal)
this.replay(sessionId, journal, fence, backgroundTasks, { type: 'reset', reset: reason })
}
snapshot(
@@ -160,14 +152,26 @@ export class AgentSessionSubscribers {
journal: AgentSessionJournal,
fence: number,
backgroundTasks?: AgentSessionBackgroundTaskState | null
): void {
this.replay(sessionId, journal, fence, backgroundTasks, { type: 'snapshot' })
}
private replay(
sessionId: string,
journal: AgentSessionJournal,
fence: number,
backgroundTasks: AgentSessionBackgroundTaskState | null | undefined,
frame: { type: 'snapshot' } | { type: 'reset'; reset: AgentJournalResetReason }
): void {
const page = readAgentSessionHydrationPage(journal, fence)
const hostNow = this.now()
for (const subscriber of this.subscribers(sessionId)) {
this.emit(subscriber, {
type: 'snapshot',
...frame,
sessionId,
page,
fence,
hostNow,
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
...this.activityField(sessionId)
})
@@ -178,18 +182,15 @@ export class AgentSessionSubscribers {
}
handoff(sessionId: string, fence: number, handoff: AgentSessionHandoffStatus): void {
const hostNow = this.now()
for (const subscriber of this.subscribers(sessionId)) {
this.emit(subscriber, {
type: 'batch',
sessionId,
batch: {
cursor: subscriber.cursor,
items: [],
removedItemIds: [],
submissions: []
},
batch: emptyAgentSessionBatch(subscriber.cursor),
fence,
handoff
handoff,
hostNow
})
subscriber.fence = fence
}
@@ -200,18 +201,15 @@ export class AgentSessionSubscribers {
state: AgentSessionBackgroundTaskState | null,
fence: number
): void {
const hostNow = this.now()
for (const subscriber of this.subscribers(sessionId)) {
this.emit(subscriber, {
type: 'batch',
sessionId,
batch: {
cursor: subscriber.cursor,
items: [],
removedItemIds: [],
submissions: []
},
batch: emptyAgentSessionBatch(subscriber.cursor),
fence,
backgroundTasks: state
backgroundTasks: state,
hostNow
})
subscriber.fence = fence
}
@@ -224,6 +222,7 @@ export class AgentSessionSubscribers {
private deliver(
subscriber: Subscriber,
journal: AgentSessionJournal,
hostNow: number,
handoff?: AgentSessionHandoffStatus,
emitCheckpoint = false,
backgroundTasks?: AgentSessionBackgroundTaskState | null,
@@ -249,6 +248,7 @@ export class AgentSessionSubscribers {
reset: result.reset,
page,
fence: subscriber.fence,
hostNow,
...(handoff ? { handoff } : {}),
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
...(publishedActivity !== undefined ? { activity: publishedActivity } : {})
@@ -266,13 +266,9 @@ export class AgentSessionSubscribers {
this.emit(subscriber, {
type: 'batch',
sessionId: subscriber.sessionId,
batch: {
cursor: page.window.nextCursor,
items: [],
removedItemIds: [],
submissions: []
},
batch: emptyAgentSessionBatch(page.window.nextCursor),
fence: subscriber.fence,
hostNow,
...(handoff ? { handoff } : {}),
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
...(publishedActivity !== undefined ? { activity: publishedActivity } : {})
@@ -290,6 +286,7 @@ export class AgentSessionSubscribers {
submissions: page.submissions
},
fence: subscriber.fence,
hostNow,
...(handoff ? { handoff } : {}),
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
...(publishedActivity !== undefined ? { activity: publishedActivity } : {})
@@ -301,6 +298,8 @@ export class AgentSessionSubscribers {
}
}
private now = (): number => this.hooks.now?.() ?? Date.now()
private isActive = (subscriber: Subscriber): boolean =>
this.bySession.get(subscriber.sessionId)?.get(subscriber.id) === subscriber
@@ -50,6 +50,23 @@ describe('rewind recovery of newer durable records', () => {
expect(restoreRewindJournalBody(status)).toEqual(status)
}
)
it('accepts a canonical turn body with a known state and keeps an unknown one as evidence', () => {
const turn = {
kind: 'turn' as const,
turnId: 'turn',
state: 'completed',
userItemId: 'codex:thread:turn:0',
startedAt: 10,
completedAt: 20,
durationMs: 10
}
expect(restoreRewindJournalBody(turn)).toEqual(turn)
const unknown = { ...turn, state: 'future-state' }
expect(restoreRewindJournalBody(unknown)).toEqual({
kind: 'status',
text: JSON.stringify(unknown)
})
})
it('does not reject a saved recovery prefix over a newer refusal reason', () => {
expect(
AgentSessionRewindRecordSchema.safeParse({
@@ -1,4 +1,5 @@
import { restoreRewindJournalBody } from './structured-rewind-journal-body'
import { isRetainedTurnRow, mergeRetainedTurnRows } from './structured-rewind-retained-turns'
import { isDeepStrictEqual } from 'node:util'
import {
agentJournalItemKey,
@@ -60,7 +61,10 @@ export async function recoverStructuredRewind(
}
throw new Error(`agent_session_rewind:${recovered?.reason ?? 'outcome-unknown'}`)
}
const expectedItems = new Set(rewind.retained.map((item) => item.itemId))
// Turn rows are the host's, never the provider's; the proof covers provider items only.
const expectedItems = new Set(
rewind.retained.filter((item) => !isRetainedTurnRow(item)).map((item) => item.itemId)
)
const observedItems = new Set<string>()
for (const { identity } of recovered.items) {
const itemId = agentJournalItemKey(identity)
@@ -76,11 +80,14 @@ export async function recoverStructuredRewind(
if (observedItems.size !== expectedItems.size) {
throw new Error('agent_session_rewind:proof-mismatch')
}
const retained = recovered.items.map(({ identity, body }) => ({
itemId: agentJournalItemKey(identity),
body,
observedAt: now()
}))
const retained = mergeRetainedTurnRows(
rewind.retained,
recovered.items.map(({ identity, body }) => ({
itemId: agentJournalItemKey(identity),
body,
observedAt: now()
}))
)
if (
retained.length > 10_000 ||
Buffer.byteLength(JSON.stringify(retained), 'utf8') > AGENT_SESSION_HISTORY_MAX_PAGE_BYTES
@@ -0,0 +1,34 @@
// The Codex preflight returns provider items only. The host's turn rows are its own record, so a
// rewind that takes the provider's list as the new epoch would drop every duration before the
// boundary unless those rows are spliced back beside the item each one followed.
import type { AgentJournalItemBody } from '../../../shared/agent-session-journal-types'
import type { AgentSessionRewindRecord } from '../../../shared/agent-session-rewind'
import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record'
type RetainedRow = AgentSessionRewindRecord['retained'][number]
export function isRetainedTurnRow(item: Pick<RetainedRow, 'body'>): boolean {
return readAgentJournalTurn(item.body as AgentJournalItemBody) !== null
}
/** `reference` fixes where each turn row sits; the provider items are the spine and keep their
* own order, including turns the local journal never saw. */
export function mergeRetainedTurnRows(
reference: readonly RetainedRow[],
providerItems: readonly RetainedRow[]
): RetainedRow[] {
const spineIndex = new Map(providerItems.map((item, index) => [item.itemId, index]))
const rowsAfter = new Map<number, RetainedRow[]>()
let anchor = -1
for (const item of reference) {
if (!isRetainedTurnRow(item)) {
anchor = spineIndex.get(item.itemId) ?? anchor
} else if (!spineIndex.has(item.itemId)) {
rowsAfter.set(anchor, [...(rowsAfter.get(anchor) ?? []), item])
}
}
const merged = [...(rowsAfter.get(-1) ?? [])]
providerItems.forEach((item, index) => merged.push(item, ...(rowsAfter.get(index) ?? [])))
return merged
}
@@ -72,7 +72,7 @@ function createReadOwner(
emit()
}
const apply = (action: StructuredAgentSessionAction): void => {
const state = reduceStructuredAgentSession(snapshot.state, action)
const state = reduceStructuredAgentSession(snapshot.state, action, Date.now())
if (state !== snapshot.state) {
setSnapshot({ ...snapshot, state })
}
@@ -85,7 +85,7 @@ export function useStructuredAgentSession(args: {
() => selectStructuredAgentTurnActivity(state.items, turnId, state.activity),
[state.activity, state.items, turnId]
)
const turnTiming = useStructuredAgentTurnTiming(state.items, state.submissions, turnId)
const turnTiming = useStructuredAgentTurnTiming(state, turnId)
const backgroundTasksView = structuredSessionBackgroundTasksView(state.backgroundTasks, turnId)
useEffect(() => {
@@ -97,11 +97,17 @@ describe('useStructuredAgentTurnTiming', () => {
HOST_START + 102_500
)
]
type Props = {
items: AgentJournalRenderItem[]
turnId: string | null
hostClock?: { hostNow: number; receivedAt: number }
}
const { result, rerender } = renderHook(
({ items, turnId }: { items: AgentJournalRenderItem[]; turnId: string | null }) =>
useStructuredAgentTurnTiming(items, SUBMISSIONS, turnId),
{ initialProps: { items: running, turnId: 't2' as string | null } }
({ items, turnId, hostClock }: Props) =>
useStructuredAgentTurnTiming({ items, submissions: SUBMISSIONS, hostClock }, turnId),
{ initialProps: { items: running, turnId: 't2' } as Props }
)
// Without a host clock the counter starts at first sight, less the append lag.
expect(result.current.workingStartedAt).toBe(CLIENT_NOW - 2_500)
// The row's provider key resolves through the submission alias, not journal order.
expect([...result.current.settledTurns]).toEqual([
@@ -116,7 +122,9 @@ describe('useStructuredAgentTurnTiming', () => {
expect(result.current.workingStartedAt).toBeNull()
vi.setSystemTime(CLIENT_NOW + 60_000)
// An older host's status carrier still anchors the counter.
// An older host's status carrier still anchors the counter. With a host clock
// that said the turn was 35s old 5s ago, the anchor sits 40s before first
// sight, wherever the client's absolute clock is.
const next = [
...running,
user('u3', 5),
@@ -127,15 +135,21 @@ describe('useStructuredAgentTurnTiming', () => {
HOST_START + 150_100
)
]
rerender({ items: next, turnId: 't3' })
expect(result.current.workingStartedAt).toBe(CLIENT_NOW + 60_000 - 100)
rerender({
items: next,
turnId: 't3',
hostClock: { hostNow: HOST_START + 185_000, receivedAt: CLIENT_NOW + 55_000 }
})
expect(result.current.workingStartedAt).toBe(CLIENT_NOW + 60_000 - 40_000)
})
it('leaves the anchor null when an older host records no start', () => {
vi.useFakeTimers()
vi.setSystemTime(CLIENT_NOW)
const items = [user('u1', 1), lifecycle('t1', 2, { state: 'running' }, HOST_START)]
const { result } = renderHook(() => useStructuredAgentTurnTiming(items, [], 't1'))
const { result } = renderHook(() =>
useStructuredAgentTurnTiming({ items, submissions: [] }, 't1')
)
expect(result.current.workingStartedAt).toBeNull()
expect(result.current.settledTurns.size).toBe(0)
})
@@ -12,22 +12,40 @@ import {
type TurnAnchor = { turnId: string; startedAt: number | null }
/** The host's clock as last published, paired with the client clock at receipt. */
type HostClock = { hostNow: number; receivedAt: number }
/** The live turn's local-clock anchor. Null when its row carries no host start
* (older hosts), so local observation applies. */
function anchorRunningTurn(items: readonly AgentJournalRenderItem[], turnId: string): TurnAnchor {
function anchorRunningTurn(
items: readonly AgentJournalRenderItem[],
turnId: string,
hostClock: HostClock | null | undefined
): TurnAnchor {
const timing = selectStructuredAgentRunningTurnTiming(items, turnId)
return {
turnId,
startedAt: timing ? structuredAgentTurnLocalStartedAt(timing, Date.now()) : null
if (!timing) {
return { turnId, startedAt: null }
}
const now = Date.now()
// Advance the published host clock by the client time since receipt; both
// terms stay single-clock, so a mid-turn attach counts from the real start.
const hostNow = hostClock ? hostClock.hostNow + (now - hostClock.receivedAt) : undefined
return { turnId, startedAt: structuredAgentTurnLocalStartedAt(timing, now, hostNow) }
}
/** Host-recorded turn timing for the structured lane: settled durations straight
* off the journal, and a skew-free start for the live counter stamped once per
* turn so re-renders never move it. */
export function useStructuredAgentTurnTiming(
items: readonly AgentJournalRenderItem[],
submissions: readonly AgentJournalSubmission[],
{
items,
submissions,
hostClock
}: {
items: readonly AgentJournalRenderItem[]
submissions: readonly AgentJournalSubmission[]
hostClock?: HostClock | null
},
turnId: string | null
): { settledTurns: ReadonlyMap<string, NativeChatSettledTurn>; workingStartedAt: number | null } {
const settledTurns = useMemo(
@@ -44,7 +62,7 @@ export function useStructuredAgentTurnTiming(
return { settledTurns, workingStartedAt: null }
}
if (anchor?.turnId !== turnId) {
const next = anchorRunningTurn(items, turnId)
const next = anchorRunningTurn(items, turnId, hostClock)
setAnchor(next)
return { settledTurns, workingStartedAt: next.startedAt }
}
@@ -0,0 +1,32 @@
import type { AgentSessionBackgroundTaskState } from './agent-session-wire'
/** Structural equality so a republished roster never churns transcript identity. */
export function backgroundTaskStatesEqual(
left: AgentSessionBackgroundTaskState | null | undefined,
right: AgentSessionBackgroundTaskState | null | undefined
): boolean {
if (left === right) {
return true
}
if (
!left ||
!right ||
left.state !== right.state ||
left.supportsTaskStop !== right.supportsTaskStop ||
left.supportsStopAll !== right.supportsStopAll
) {
return false
}
if (left.tasks === right.tasks) {
return true
}
if (!left.tasks || !right.tasks || left.tasks.length !== right.tasks.length) {
return false
}
return left.tasks.every(
(task, index) =>
task.id === right.tasks?.[index]?.id &&
task.kind === right.tasks[index]?.kind &&
task.description === right.tasks[index]?.description
)
}
+12 -6
View File
@@ -125,6 +125,9 @@ export type AgentSessionHistoryPage = {
hasNewer: boolean
/** Present on hosts that expose provider-owned background task lifecycle. */
backgroundTasks?: AgentSessionBackgroundTaskState | null
/** Host wall clock (ms epoch) when the page was read, so a client attaching mid-turn
* can anchor a live counter on the real start. Absent from older hosts. */
hostNow?: number
}
export type AgentSessionHistoryResult =
@@ -149,8 +152,11 @@ export type AgentSessionJournalBatch = {
submissions: AgentJournalSubmission[]
}
/** Host wall clock (ms epoch) stamped once per published frame; see `AgentSessionHistoryPage`. */
type AgentSessionHostClockField = { hostNow?: number }
export type AgentSessionSubscribeEvent =
| {
| ({
type: 'snapshot'
sessionId: string
page: AgentSessionHistoryPage
@@ -161,8 +167,8 @@ export type AgentSessionSubscribeEvent =
commands?: AgentSessionSlashCommand[] | null
/** Latest provider-authored turn activity; optional for mixed-version hosts. */
activity?: AgentSessionTurnActivity | null
}
| {
} & AgentSessionHostClockField)
| ({
type: 'batch'
sessionId: string
batch: AgentSessionJournalBatch
@@ -174,8 +180,8 @@ export type AgentSessionSubscribeEvent =
commands?: AgentSessionSlashCommand[] | null
/** Additive ephemeral state; it never creates or advances journal rows. */
activity?: AgentSessionTurnActivity | null
}
| {
} & AgentSessionHostClockField)
| ({
type: 'reset'
sessionId: string
reset: AgentJournalResetReason
@@ -186,7 +192,7 @@ export type AgentSessionSubscribeEvent =
/** Omitted when unchanged; null clears a previous provider catalog. */
commands?: AgentSessionSlashCommand[] | null
activity?: AgentSessionTurnActivity | null
}
} & AgentSessionHostClockField)
| { type: 'end' }
// ─── Status feed ────────────────────────────────────────────────────────────
@@ -456,6 +456,79 @@ describe('structured agent session reducer', () => {
expect(cleared.items).toBe(active.items)
})
it('records the host clock from frames that carry it and keeps it otherwise', () => {
const snapshot = reduceStructuredAgentSession(
EMPTY_STRUCTURED_AGENT_SESSION,
{
type: 'event',
event: {
type: 'snapshot',
sessionId: 'session-a',
fence: 1,
page: hydrationPage([item('first', 1)]),
hostNow: 5_000
}
},
9_000
)
expect(snapshot.hostClock).toEqual({ hostNow: 5_000, receivedAt: 9_000 })
const batch = reduceStructuredAgentSession(
snapshot,
{
type: 'event',
event: {
type: 'batch',
sessionId: 'session-a',
fence: 1,
hostNow: 5_400,
batch: {
cursor: { epoch: 'epoch-a', sequence: 2 },
items: [item('second', 2)],
removedItemIds: [],
submissions: []
}
}
},
9_400
)
expect(batch.hostClock).toEqual({ hostNow: 5_400, receivedAt: 9_400 })
// An older host stamps nothing; the last sample stays usable.
const unstamped = reduceStructuredAgentSession(
batch,
{
type: 'event',
event: {
type: 'batch',
sessionId: 'session-a',
fence: 1,
batch: {
cursor: { epoch: 'epoch-a', sequence: 3 },
items: [item('third', 3)],
removedItemIds: [],
submissions: []
}
}
},
9_800
)
expect(unstamped.hostClock).toEqual({ hostNow: 5_400, receivedAt: 9_400 })
const paged = reduceStructuredAgentSession(
unstamped,
{ type: 'tail-page', page: { ...hydrationPage([item('fourth', 4)]), hostNow: 6_000 } },
10_000
)
expect(paged.hostClock).toEqual({ hostNow: 6_000, receivedAt: 10_000 })
expect(
reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
type: 'tail-page',
page: hydrationPage([item('first', 1)])
}).hostClock
).toBeUndefined()
})
it('retains same-epoch activity across a newer journal tail refresh', () => {
const active = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
type: 'event',
+30 -33
View File
@@ -11,8 +11,16 @@ import type {
AgentSessionSubscribeEvent,
AgentSessionTurnActivity
} from './agent-session-wire'
import { backgroundTaskStatesEqual } from './agent-session-background-task-state-equality'
import { agentJournalSubmissionKey } from './agent-session-journal-item-key'
/** The last host clock sample: `hostNow - receivedAt` is the client's skew from the host,
* which is what lets a client attaching mid-turn anchor its live counter on the real start. */
export type StructuredAgentHostClock = {
hostNow: number
receivedAt: number
}
export type StructuredAgentSessionState = {
epoch: string | null
cursor: AgentJournalCursor | null
@@ -26,6 +34,8 @@ export type StructuredAgentSessionState = {
backgroundTasks?: AgentSessionBackgroundTaskState | null
commands?: AgentSessionSlashCommand[] | null
activity?: AgentSessionTurnActivity | null
/** Absent until a frame from a host that stamps `hostNow` has been applied. */
hostClock?: StructuredAgentHostClock
}
export type StructuredAgentSessionAction =
@@ -49,34 +59,14 @@ export const EMPTY_STRUCTURED_AGENT_SESSION: StructuredAgentSessionState = {
const MAX_RETAINED_SUBMISSIONS = 256
function backgroundTaskStatesEqual(
left: AgentSessionBackgroundTaskState | null | undefined,
right: AgentSessionBackgroundTaskState | null | undefined
): boolean {
if (left === right) {
return true
}
if (
!left ||
!right ||
left.state !== right.state ||
left.supportsTaskStop !== right.supportsTaskStop ||
left.supportsStopAll !== right.supportsStopAll
) {
return false
}
if (left.tasks === right.tasks) {
return true
}
if (!left.tasks || !right.tasks || left.tasks.length !== right.tasks.length) {
return false
}
return left.tasks.every(
(task, index) =>
task.id === right.tasks?.[index]?.id &&
task.kind === right.tasks[index]?.kind &&
task.description === right.tasks[index]?.description
)
/** A frame without `hostNow` (older host) leaves the previous sample in place. */
function hostClockField(
hostNow: number | undefined,
receivedAt: number,
previous: StructuredAgentHostClock | undefined
): { hostClock?: StructuredAgentHostClock } {
const hostClock = hostNow !== undefined ? { hostNow, receivedAt } : previous
return hostClock ? { hostClock } : {}
}
function replacePage(
@@ -145,9 +135,11 @@ function mergeSubmissions(
)
}
/** `receivedAt` is the client clock at apply time; callers pass it so the reducer stays pure. */
export function reduceStructuredAgentSession(
state: StructuredAgentSessionState,
action: StructuredAgentSessionAction
action: StructuredAgentSessionAction,
receivedAt: number = Date.now()
): StructuredAgentSessionState {
if (action.type === 'loading') {
// Keep the last transcript visible while a reconnect rehydrates the stream.
@@ -182,6 +174,7 @@ export function reduceStructuredAgentSession(
...(action.page.backgroundTasks !== undefined
? { backgroundTasks: action.page.backgroundTasks }
: {}),
...hostClockField(action.page.hostNow, receivedAt, state.hostClock),
status: 'ready',
error: undefined
}
@@ -206,7 +199,8 @@ export function reduceStructuredAgentSession(
? { backgroundTasks: action.page.backgroundTasks }
: state.backgroundTasks !== undefined
? { backgroundTasks: state.backgroundTasks }
: {})
: {}),
...hostClockField(action.page.hostNow, receivedAt, state.hostClock)
}
}
if (action.type === 'older-page') {
@@ -218,7 +212,8 @@ export function reduceStructuredAgentSession(
...state,
items,
submissions: mergeSubmissions(state.submissions, action.page.submissions, items),
hasOlder: action.page.hasOlder
hasOlder: action.page.hasOlder,
...hostClockField(action.page.hostNow, receivedAt, state.hostClock)
}
}
const event = action.event
@@ -228,7 +223,8 @@ export function reduceStructuredAgentSession(
if (event.type === 'snapshot' || event.type === 'reset') {
return {
...replacePage(event.page, event.fence, event.handoff, event.backgroundTasks, event.activity),
commands: event.commands
commands: event.commands,
...hostClockField(event.hostNow, receivedAt, state.hostClock)
}
}
if (state.epoch !== event.batch.cursor.epoch) {
@@ -275,7 +271,8 @@ export function reduceStructuredAgentSession(
handoff: event.handoff ?? state.handoff,
commands: event.commands !== undefined ? event.commands : state.commands,
...(backgroundTasks !== undefined ? { backgroundTasks } : {}),
...(activity !== undefined ? { activity } : {})
...(activity !== undefined ? { activity } : {}),
...hostClockField(event.hostNow, receivedAt, state.hostClock)
}
}
@@ -272,3 +272,69 @@ describe('provider-measured duration', () => {
).toBeNull()
})
})
describe('coalesced sends and canonical rows', () => {
const accepted = (clientMessageId: string, providerItemId: string) => ({
clientMessageId,
fence: 1,
payloadFingerprint: 'fp',
dispatchState: 'accepted' as const,
providerItemId,
reason: null,
submittedAt: 1,
resolvedAt: 2
})
it('gives a turn shared by two accepted sends to the prompt that opened it', () => {
const items = [
user('orca:first'),
user('orca:second'),
lifecycle('t1', {
state: 'completed',
userItemId: 'codex:thread:t1:0',
startedAt: 1_000,
completedAt: 5_000
})
]
const timings = selectStructuredAgentTurnTimings(items, [
accepted('first', 'codex:thread:t1:0'),
accepted('second', 'codex:thread:t1:0')
])
expect([...timings.keys()]).toEqual(['orca:first'])
})
it('reads a canonical turn item exactly like the legacy carrier', () => {
sequence += 1
const canonical: AgentJournalRenderItem = {
itemId: 'legacy:codex:s:turn-lifecycle%3At9',
revision: 2,
sequence,
observedAt: 1_000,
body: {
kind: 'turn',
turnId: 't9',
state: 'completed',
userItemId: 'orca:u9',
startedAt: 1_000,
completedAt: 9_000,
durationMs: 7_172
}
}
const timings = selectStructuredAgentTurnTimings([user('orca:u9'), canonical])
expect(timings.get('orca:u9')).toMatchObject({ state: 'completed', durationMs: 7_172 })
expect(selectStructuredAgentSettledTurns([user('orca:u9'), canonical]).get('orca:u9')).toEqual({
startedAt: 1_000,
workedSeconds: 7
})
expect(selectStructuredAgentRunningTurnTiming([canonical], 't9')?.startedAt).toBe(1_000)
})
})
describe('structuredAgentTurnLocalStartedAt with the host clock', () => {
it('counts a mid-turn attach from the real start, not from first sight', () => {
const timing = { state: 'running' as const, startedAt: 50_000, observedAt: 50_000 }
// Host says the turn has run 40s; client clock is arbitrary.
expect(structuredAgentTurnLocalStartedAt(timing, 3_600_000, 90_000)).toBe(3_600_000 - 40_000)
expect(structuredAgentTurnLocalStartedAt(timing, 3_600_000, 40_000)).toBe(3_600_000)
})
})
@@ -63,8 +63,10 @@ export function selectStructuredAgentTurnTimings(
): ReadonlyMap<string, StructuredAgentTurnTiming> {
const itemIds = new Set(items.map((item) => item.itemId))
const aliases = new Map<string, string>()
// Codex folds a send issued mid-turn into the running turn under the SAME provider
// key, so the earliest submission that names a key is the prompt that opened the turn.
for (const submission of submissions) {
if (submission.providerItemId) {
if (submission.providerItemId && !aliases.has(submission.providerItemId)) {
aliases.set(submission.providerItemId, agentJournalSubmissionKey(submission.clientMessageId))
}
}
@@ -119,14 +121,22 @@ export function completedStructuredAgentTurnSeconds(
: null
}
/** A local-clock anchor for the live counter that carries no host/client skew:
* the client's first sighting of the running row, moved back by the host-side
* lag between turn-start receipt and the row's append. Both terms are single-clock. */
/** A local-clock anchor for the live counter that carries no host/client skew.
* With the host's own clock at publish time, the anchor is the client's first
* sighting moved back by how long the host says the turn has already run, so a
* client attaching mid-turn counts from the real start. Without it, only the
* host-side lag between turn-start receipt and the row's append is known, and
* the counter starts at first sight. Every difference is single-clock. */
export function structuredAgentTurnLocalStartedAt(
timing: StructuredAgentTurnTiming,
firstSeenAt: number
firstSeenAt: number,
hostNow?: number
): number {
return firstSeenAt - Math.max(0, timing.observedAt - timing.startedAt)
const hostElapsed =
hostNow !== undefined && Number.isFinite(hostNow)
? hostNow - timing.startedAt
: timing.observedAt - timing.startedAt
return firstSeenAt - Math.max(0, hostElapsed)
}
/** The settled turns a chat surface hands to the shared turn-status selector. */