From e2b3abe9e454f1b2edb041fce7b6cb93ad41d76f Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Wed, 9 Sep 2026 22:09:25 -0700 Subject: [PATCH] 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. --- .../use-mobile-structured-agent-session.ts | 2 +- .../use-mobile-structured-agent-state.ts | 2 +- ...bile-structured-agent-turn-timing.test.tsx | 30 +++++++- ...use-mobile-structured-agent-turn-timing.ts | 32 ++++++-- .../journal-epoch-rollover.ts | 5 +- .../journal-row-schema-version.test.ts | 66 ++++++++++++++++ .../agent-session-empty-batch.ts | 7 ++ ...d-agent-session-background-task-channel.ts | 12 ++- .../structured-agent-session-host.test.ts | 2 + .../structured-agent-session-host.ts | 3 +- .../structured-agent-session-rewind.test.ts | 50 +++++++++++++ .../structured-agent-session-rewind.ts | 27 ++++--- ...ructured-agent-session-subscribers.test.ts | 50 +++++++++++++ .../structured-agent-session-subscribers.ts | 75 +++++++++---------- .../structured-rewind-journal-body.test.ts | 17 +++++ .../structured-rewind-recovery.ts | 19 +++-- .../structured-rewind-retained-turns.ts | 34 +++++++++ .../structured-agent-session-read-owner.ts | 2 +- .../use-structured-agent-session.ts | 2 +- .../use-structured-agent-turn-timing.test.tsx | 28 +++++-- .../use-structured-agent-turn-timing.ts | 32 ++++++-- ...-session-background-task-state-equality.ts | 32 ++++++++ src/shared/agent-session-wire.ts | 18 +++-- .../structured-agent-session-reducer.test.ts | 73 ++++++++++++++++++ .../structured-agent-session-reducer.ts | 63 ++++++++-------- ...ructured-agent-session-turn-timing.test.ts | 66 ++++++++++++++++ .../structured-agent-session-turn-timing.ts | 22 ++++-- 27 files changed, 639 insertions(+), 132 deletions(-) create mode 100644 src/main/native-chat/agent-session-journal/journal-row-schema-version.test.ts create mode 100644 src/main/native-chat/agent-session-wire/agent-session-empty-batch.ts create mode 100644 src/main/native-chat/agent-session-wire/structured-rewind-retained-turns.ts create mode 100644 src/shared/agent-session-background-task-state-equality.ts diff --git a/mobile/src/session/use-mobile-structured-agent-session.ts b/mobile/src/session/use-mobile-structured-agent-session.ts index b079b1073ba..35291d4dd08 100644 --- a/mobile/src/session/use-mobile-structured-agent-session.ts +++ b/mobile/src/session/use-mobile-structured-agent-session.ts @@ -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, diff --git a/mobile/src/session/use-mobile-structured-agent-state.ts b/mobile/src/session/use-mobile-structured-agent-state.ts index 52aefab24aa..8d71c7fde29 100644 --- a/mobile/src/session/use-mobile-structured-agent-state.ts +++ b/mobile/src/session/use-mobile-structured-agent-state.ts @@ -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 } diff --git a/mobile/src/session/use-mobile-structured-agent-turn-timing.test.tsx b/mobile/src/session/use-mobile-structured-agent-turn-timing.test.tsx index ad231659d93..c35975a50eb 100644 --- a/mobile/src/session/use-mobile-structured-agent-turn-timing.test.tsx +++ b/mobile/src/session/use-mobile-structured-agent-turn-timing.test.tsx @@ -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', () => { diff --git a/mobile/src/session/use-mobile-structured-agent-turn-timing.ts b/mobile/src/session/use-mobile-structured-agent-turn-timing.ts index 7c05063a95b..8e6c9e33aae 100644 --- a/mobile/src/session/use-mobile-structured-agent-turn-timing.ts +++ b/mobile/src/session/use-mobile-structured-agent-turn-timing.ts @@ -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; 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 } } diff --git a/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts b/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts index e95ff821058..5a5f39e5a62 100644 --- a/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts +++ b/src/main/native-chat/agent-session-journal/journal-epoch-rollover.ts @@ -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, diff --git a/src/main/native-chat/agent-session-journal/journal-row-schema-version.test.ts b/src/main/native-chat/agent-session-journal/journal-row-schema-version.test.ts new file mode 100644 index 00000000000..b764ebd5ec9 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-row-schema-version.test.ts @@ -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() + } + }) +}) diff --git a/src/main/native-chat/agent-session-wire/agent-session-empty-batch.ts b/src/main/native-chat/agent-session-wire/agent-session-empty-batch.ts new file mode 100644 index 00000000000..bc713f92171 --- /dev/null +++ b/src/main/native-chat/agent-session-wire/agent-session-empty-batch.ts @@ -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: [] } +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts index 4592dea26e9..0096286ca18 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts @@ -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 { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts index 5354670fa0b..e7d389633f3 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host.test.ts @@ -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 }) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts index 1f11555b021..3cdda607ff4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts @@ -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 diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts index 554b8d34395..c75301c6461 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts @@ -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) => ({ diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.ts index f9b9554a86c..417a42fe414 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.ts @@ -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') > diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts index a32db8786d4..4f22b56397c 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts @@ -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' } }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts index a5062373dad..24415df1d3c 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts @@ -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 diff --git a/src/main/native-chat/agent-session-wire/structured-rewind-journal-body.test.ts b/src/main/native-chat/agent-session-wire/structured-rewind-journal-body.test.ts index 167c67c5ecf..e93aef352fe 100644 --- a/src/main/native-chat/agent-session-wire/structured-rewind-journal-body.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-rewind-journal-body.test.ts @@ -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({ diff --git a/src/main/native-chat/agent-session-wire/structured-rewind-recovery.ts b/src/main/native-chat/agent-session-wire/structured-rewind-recovery.ts index 41ccab4ae28..8f5a8942a86 100644 --- a/src/main/native-chat/agent-session-wire/structured-rewind-recovery.ts +++ b/src/main/native-chat/agent-session-wire/structured-rewind-recovery.ts @@ -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() 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 diff --git a/src/main/native-chat/agent-session-wire/structured-rewind-retained-turns.ts b/src/main/native-chat/agent-session-wire/structured-rewind-retained-turns.ts new file mode 100644 index 00000000000..aa0753b502c --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-rewind-retained-turns.ts @@ -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): 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() + 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 +} diff --git a/src/renderer/src/components/native-chat/structured-agent-session-read-owner.ts b/src/renderer/src/components/native-chat/structured-agent-session-read-owner.ts index 4c43326acb1..b36d1256787 100644 --- a/src/renderer/src/components/native-chat/structured-agent-session-read-owner.ts +++ b/src/renderer/src/components/native-chat/structured-agent-session-read-owner.ts @@ -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 }) } diff --git a/src/renderer/src/components/native-chat/use-structured-agent-session.ts b/src/renderer/src/components/native-chat/use-structured-agent-session.ts index f26e3f01825..8190750d3c5 100644 --- a/src/renderer/src/components/native-chat/use-structured-agent-session.ts +++ b/src/renderer/src/components/native-chat/use-structured-agent-session.ts @@ -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(() => { diff --git a/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.test.tsx b/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.test.tsx index 8ee622f0b17..e24416ccbcc 100644 --- a/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.test.tsx +++ b/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.test.tsx @@ -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) }) diff --git a/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.ts b/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.ts index 4552b750342..0ca196f93d6 100644 --- a/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.ts +++ b/src/renderer/src/components/native-chat/use-structured-agent-turn-timing.ts @@ -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; 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 } } diff --git a/src/shared/agent-session-background-task-state-equality.ts b/src/shared/agent-session-background-task-state-equality.ts new file mode 100644 index 00000000000..b852873b5ad --- /dev/null +++ b/src/shared/agent-session-background-task-state-equality.ts @@ -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 + ) +} diff --git a/src/shared/agent-session-wire.ts b/src/shared/agent-session-wire.ts index 2b5c0e5ae07..bd635ed2b1c 100644 --- a/src/shared/agent-session-wire.ts +++ b/src/shared/agent-session-wire.ts @@ -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 ──────────────────────────────────────────────────────────── diff --git a/src/shared/structured-agent-session-reducer.test.ts b/src/shared/structured-agent-session-reducer.test.ts index 9aec45c8bf1..89428bdab7e 100644 --- a/src/shared/structured-agent-session-reducer.test.ts +++ b/src/shared/structured-agent-session-reducer.test.ts @@ -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', diff --git a/src/shared/structured-agent-session-reducer.ts b/src/shared/structured-agent-session-reducer.ts index 075e4d8bb09..935d86409a1 100644 --- a/src/shared/structured-agent-session-reducer.ts +++ b/src/shared/structured-agent-session-reducer.ts @@ -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) } } diff --git a/src/shared/structured-agent-session-turn-timing.test.ts b/src/shared/structured-agent-session-turn-timing.test.ts index e23182a7bb0..f1487f0ef42 100644 --- a/src/shared/structured-agent-session-turn-timing.test.ts +++ b/src/shared/structured-agent-session-turn-timing.test.ts @@ -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) + }) +}) diff --git a/src/shared/structured-agent-session-turn-timing.ts b/src/shared/structured-agent-session-turn-timing.ts index a0fdb4abb50..fca083e418e 100644 --- a/src/shared/structured-agent-session-turn-timing.ts +++ b/src/shared/structured-agent-session-turn-timing.ts @@ -63,8 +63,10 @@ export function selectStructuredAgentTurnTimings( ): ReadonlyMap { const itemIds = new Set(items.map((item) => item.itemId)) const aliases = new Map() + // 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. */