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. */