fix(native-chat): a long reply keeps its Working bar and Stop after a reload, from the host's turn (#25751)

* fix(native-chat): the chat reads whether a turn runs from the host, not the loaded history page

* fix(native-chat): the Working bar, card cancel and phone Stop follow the host's turn when its rows are not loaded

* fix(native-chat): keep the phone's Stop as it was, naming the host's running turn

Drops the conversation-wide phone Stop: the phone shows Stop only once a turn
is running, so a Stop naming no turn could not be reached and only changed the
ordinary Stop's request. Stop and card cancel still name the turn the host says
is running.

* fix(native-chat): the host's turn answer never outlives its host, and a half-loaded turn shows no edit totals

- A batch carrying rows restates the host's newest turn record, so one without
  it (an older host after a downgrade) drops back to the loaded rows instead of
  holding the last answer; the event coalescer merges by the same rule.
- A turn whose record is above the loaded rows is grouped only while it runs,
  and its loaded tail gets no "N changed files" total, which would undercount.
- "Thinking" reads the host's running turn too, so it shows when the turn's
  record is not loaded.
- The wire doc for observedAt says what the journal stores: creation time.

* refactor(native-chat): totals per turn come from their own hook, keeping the transcript list under its line limit

Also drops the chat zoom prop main removed from the unloaded live turn test.
This commit is contained in:
Brennan Benson
2026-10-06 10:57:15 -07:00
committed by GitHub
parent 842c6c667a
commit ccab2a3110
33 changed files with 1021 additions and 162 deletions
@@ -1,7 +1,7 @@
import type { AgentSessionCancelResult } from '../../../src/shared/agent-session-wire'
import type { AgentJournalRenderItem } from '../../../src/shared/agent-session-journal-types'
import type { StructuredAgentSessionState } from '../../../src/shared/structured-agent-session-reducer'
import { activeStructuredAgentSessionTurnId } from '../../../src/shared/structured-agent-session-live-turn'
import { runningStructuredAgentSessionTurnId } from '../../../src/shared/structured-agent-session-live-turn'
import type { RpcClient } from '../transport/rpc-client'
import {
requestStructuredAgentSessionMutation,
@@ -37,7 +37,7 @@ export async function requestMobileStructuredAgentSessionCancel(args: {
}): Promise<boolean> {
const { client, enabled, inFlight, onSendError, sessionId, stateRef } = args
const current = stateRef.current
const turnId = activeStructuredAgentSessionTurnId(current.items)
const turnId = runningStructuredAgentSessionTurnId(current)
if (!client || !sessionId || !enabled || current.fence === null || !turnId) {
onSendError('Stop not sent')
return false
@@ -767,6 +767,49 @@ describe('useMobileNativeChatTurnDisclosure', () => {
expect(rows.map((row) => row.activeTurnIsWorking)).toEqual([false, false, true])
expect(rows[0].turnStatus?.workedSeconds).toBe(4)
})
it('draws the live bar on a turn whose record and opening message are not loaded', () => {
// The newest rows of a long turn; its record is above the page, named only by the host.
const items: AgentJournalRenderItem[] = ['a300', 'a301'].map((itemId, index) => ({
itemId,
revision: 0,
sequence: 300 + index,
observedAt: 300 + index,
body: { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: itemId }] },
turnScope: { kind: 'turn', turnItemId: 'turn-record' }
}))
const messages: NativeChatMessage[] = items.map((item) => ({
...userMessage(item.itemId),
role: 'assistant'
}))
const rowsWith = (latestTurn: NativeChatTurnJournal['latestTurn']) => {
act(() => {
renderer = create(
createElement(Harness, {
messages,
enabled: true,
workingStartedAt: 1_000,
turnJournal: { items, submissions: [], latestTurn }
})
)
})
const disclosure = renderer!.root.findByType('result').props.disclosure
const rows = messages.map((message, index) => disclosure.resolveRow(index, message))
act(() => renderer?.unmount())
return rows
}
const hosted = rowsWith({
itemId: 'turn-record',
observedAt: 1,
turn: { turnId: 'turn-1', state: 'running', startedAt: 1_000, userItemId: 'user-1' }
})
expect(hosted[0]?.turnStatus).toMatchObject({ startedAt: 1_000 })
expect(hosted.map((row) => row.activeTurnIsWorking)).toEqual([true, true])
// From the loaded rows alone, nothing names the turn, so no bar draws.
expect(rowsWith(undefined).map((row) => row.turnStatus)).toEqual([null, null])
})
})
describe('the open reasoning block the live line discloses', () => {
@@ -224,6 +224,49 @@ describe('mobile structured prompt cancellation', () => {
)
})
it("stops the host's running turn when its record is not loaded", async () => {
state = {
...state,
hasOlder: true,
items: [pendingApproval()],
latestTurn: {
itemId: 'turn-status',
observedAt: 1,
turn: { turnId: 'turn-1', state: 'running', startedAt: 1 }
}
}
act(() => {
renderer = create(createElement(Harness, { promptCancelSupported: true }))
})
expect(hook.isWorking).toBe(true)
await act(async () => {
expect(await hook.cancelPrompt()).toBe(true)
})
expect(mocks.sendRequest).toHaveBeenCalledWith(
'agentSession.cancel',
expect.objectContaining({
turnId: 'turn-1',
prompt: { itemId: 'approval-1', expectedRevision: 4 }
}),
expect.any(Object)
)
})
it('reads a turn the host ended as ended, though its running row is still loaded', () => {
state = {
...state,
latestTurn: {
itemId: 'turn-status',
observedAt: 1,
turn: { turnId: 'turn-1', state: 'completed', startedAt: 1, completedAt: 2 }
}
}
act(() => {
renderer = create(createElement(Harness, { promptCancelSupported: true }))
})
expect(hook.turnId).toBeNull()
})
it('uses the rendered prompt identity when the journal changes before tap', async () => {
act(() => {
renderer = create(createElement(Harness, { promptCancelSupported: true }))
@@ -243,4 +286,28 @@ describe('mobile structured prompt cancellation', () => {
expect.any(Object)
)
})
it("names the host's running turn for a plain Stop when its record is not loaded", async () => {
state = {
...state,
hasOlder: true,
items: [],
latestTurn: {
itemId: 'turn-status',
observedAt: 1,
turn: { turnId: 'turn-1', state: 'running' }
}
}
act(() => {
renderer = create(createElement(Harness, { promptCancelSupported: true }))
})
expect(hook.turnId).toBe('turn-1')
await act(async () => {
hook.cancel()
})
const call = mocks.sendRequest.mock.calls.find(([method]) => method === 'agentSession.cancel')
expect(call).toBeDefined()
expect(call?.[1]).toMatchObject({ turnId: 'turn-1' })
expect(call?.[1]).not.toHaveProperty('prompt')
})
})
@@ -7,8 +7,8 @@ import { TUI_AGENT_DISPLAY_NAMES } from '../../../src/shared/tui-agent-display-n
import { isStructuredAgentSessionMainAgentWorking } from '../../../src/shared/structured-agent-session-main-agent-working'
import { isFinalAgentSessionReadRefusal } from '../../../src/shared/structured-agent-session-read-refusal'
import {
activeStructuredAgentSessionTurnId,
isStructuredAgentSessionThinking
isStructuredAgentSessionThinking,
runningStructuredAgentSessionTurnId
} from '../../../src/shared/structured-agent-session-live-turn'
import { selectStructuredAgentTurnActivity } from '../../../src/shared/native-chat-turn-activity'
import {
@@ -177,14 +177,14 @@ export function useMobileStructuredAgentSession(args: {
}),
[transcriptItems, state.submissions]
)
const turnId = activeStructuredAgentSessionTurnId(state.items)
const turnId = runningStructuredAgentSessionTurnId(state)
const turnTiming = useMobileStructuredAgentTurnTiming(
{ ...state, items: transcriptItems },
turnId
)
const activityText =
selectStructuredAgentTurnActivity(state.items, turnId, state.activity)?.text ?? null
const thinking = isStructuredAgentSessionThinking(state.items)
const thinking = isStructuredAgentSessionThinking(state)
const turnIndicator = useMemo(() => ({ thinking, activityText }), [thinking, activityText])
const status = state.status === 'idle' ? 'idle' : state.status
const approvalPrompt = useMemo(
@@ -4,6 +4,7 @@ import type {
AgentJournalSubmission
} from '../../../src/shared/agent-session-journal-types'
import type { NativeChatSettledTurns } from '../../../src/shared/native-chat-turn-status'
import type { AgentSessionLatestTurn } from '../../../src/shared/agent-session-wire'
import type { NativeChatTurnJournal } from '../../../src/shared/native-chat-turn-membership'
import type { StructuredAgentHostClock } from '../../../src/shared/structured-agent-session-reducer'
import { selectStructuredAgentTurnBars } from '../../../src/shared/structured-agent-session-turn-timing'
@@ -19,10 +20,12 @@ export function useMobileStructuredAgentTurnTiming(
{
items,
submissions,
latestTurn,
hostClock
}: {
items: readonly AgentJournalRenderItem[]
submissions: readonly AgentJournalSubmission[]
latestTurn?: AgentSessionLatestTurn | null
hostClock?: StructuredAgentHostClock | null
},
turnId: string | null
@@ -33,10 +36,13 @@ export function useMobileStructuredAgentTurnTiming(
workingStartedAt: number | null
} {
const { settledTurns, runningTiming } = useMemo(
() => selectStructuredAgentTurnBars(items, submissions, turnId),
[items, submissions, turnId]
() => selectStructuredAgentTurnBars(items, submissions, turnId, latestTurn),
[items, submissions, turnId, latestTurn]
)
const turnJournal = useMemo(
() => ({ items, submissions, latestTurn }),
[items, submissions, latestTurn]
)
const turnJournal = useMemo(() => ({ items, submissions }), [items, submissions])
const [latch, setLatch] = useState<StructuredAgentTurnClockLatch | null>(null)
// Stamp during render (React's derive-from-props pattern) so the first paint of
// a new turn already counts from the right instant.
@@ -4,7 +4,7 @@
// host advertises the queue and no pending prompt is one this build cannot answer.
import { useCallback } from 'react'
import { activeStructuredAgentSessionTurnId } from '../../../src/shared/structured-agent-session-live-turn'
import { runningStructuredAgentSessionTurnId } from '../../../src/shared/structured-agent-session-live-turn'
import { pendingPromptsAllUnanswerableHere } from '../../../src/shared/agent-session-approval-subject'
import {
structuredAgentSessionSendBody,
@@ -95,7 +95,7 @@ export function useMobileStructuredSendWithOutcome(args: {
...controller
},
canRun: () =>
!activeStructuredAgentSessionTurnId(stateRef.current.items) &&
!runningStructuredAgentSessionTurnId(stateRef.current) &&
!stateRef.current.items.some(
(item) => pendingStructuredApproval(item) || pendingStructuredQuestion(item)
),
@@ -163,7 +163,7 @@ describe("a Codex subagent's rows on the parent's surfaces", () => {
summary: ['Reading the diff']
})
expect(isStructuredAgentSessionThinking(await items())).toBe(false)
expect(isStructuredAgentSessionThinking({ items: await items() })).toBe(false)
})
it("does not show the child's compaction as the parent's activity line", async () => {
@@ -9,6 +9,7 @@
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
import { isRootAgentJournalItem } from '../../../shared/agent-session-journal-producer'
import { latestStructuredAgentSessionTurn } from '../../../shared/structured-agent-session-live-turn'
import type {
AgentJournalCursor,
AgentJournalRenderItem,
@@ -21,6 +22,7 @@ import {
type AgentSessionHistoryPage,
type AgentSessionHistoryRequest,
type AgentSessionHistoryResult,
type AgentSessionLatestTurn,
type AgentSessionSubagentRosterEntry
} from '../../../shared/agent-session-wire'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
@@ -308,6 +310,20 @@ function buildPage(input: {
liveCursor: input.snapshot.cursor,
hasOlder: input.hasOlder,
hasNewer: input.hasNewer,
...(input.subagentRoster === undefined ? {} : { subagentRoster: input.subagentRoster })
...(input.subagentRoster === undefined ? {} : { subagentRoster: input.subagentRoster }),
// From the whole timeline, never the page: a long turn's record sits below any window.
latestTurn: snapshotLatestTurn(input.snapshot)
}
}
// A catch-up run pages one snapshot many times; its scan back to the turn record is paid once.
const latestTurnBySnapshot = new WeakMap<AgentJournalSnapshot, AgentSessionLatestTurn | null>()
function snapshotLatestTurn(snapshot: AgentJournalSnapshot): AgentSessionLatestTurn | null {
let latest = latestTurnBySnapshot.get(snapshot)
if (latest === undefined) {
latest = latestStructuredAgentSessionTurn(snapshot.items)
latestTurnBySnapshot.set(snapshot, latest)
}
return latest
}
@@ -93,6 +93,7 @@ export function deliverToSubscriber(
submissions: page.submissions
},
fence: subscriber.fence,
...(page.latestTurn !== undefined ? { latestTurn: page.latestTurn } : {}),
...shared
},
// On a multi-page catch-up the draft list rides only the final page, or a
@@ -0,0 +1,214 @@
// A long turn's record keeps the place it opened at when revised, so a bounded page leaves it out.
// Whether a turn runs comes from the host's whole-journal record, never the loaded rows.
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 { codexProviderHandle } from '../../../shared/agent-session-provider-handle-encoding'
import { agentJournalItemKey } from '../../../shared/agent-session-journal-item-key'
import type {
AgentJournalItemIdentity,
AgentJournalTurnLifecycle,
AgentSessionJournalIdentity
} from '../../../shared/agent-session-journal-types'
import type { AgentSessionSubscribeEvent } from '../../../shared/agent-session-wire'
import {
activeStructuredAgentSessionTurnId,
runningStructuredAgentSessionTurnId
} from '../../../shared/structured-agent-session-live-turn'
import {
EMPTY_STRUCTURED_AGENT_SESSION,
reduceStructuredAgentSession,
type StructuredAgentSessionState
} from '../../../shared/structured-agent-session-reducer'
import { selectStructuredAgentTurnBars } from '../../../shared/structured-agent-session-turn-timing'
import { createTrackedJournalOpener } from '../agent-session-journal/journal-host-database-test-support'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import {
readAgentSessionHistory,
readAgentSessionHydrationPage
} from './agent-session-history-page'
import { AgentSessionSubscribers } from './structured-agent-session-subscribers'
const IDENTITY: AgentSessionJournalIdentity = {
sessionId: 'session-1',
workspaceId: 'ws-1',
hostId: 'host-1',
agent: 'codex',
providerHandle: codexProviderHandle('thread-1')
}
const journals = createTrackedJournalOpener()
let root: string
let journal: AgentSessionJournal
beforeEach(async () => {
root = await mkdtemp(join(tmpdir(), 'orca-live-turn-window-'))
journal = await journals.open({ identity: IDENTITY, stateDirectory: root })
})
afterEach(async () => {
await journals.closeAll()
await rm(root, { recursive: true, force: true })
})
function turnRecord(turnId: string, ordinal = 0): AgentJournalItemIdentity {
return { provider: 'codex', threadId: 'thread-1', turnId, ordinal }
}
async function writeTurn(turnId: string, turn: Omit<AgentJournalTurnLifecycle, 'turnId'>) {
await journal.appendItem(
turnRecord(turnId),
{ kind: 'turn', turnId, ...turn },
{ fence: 1, turnScope: { kind: 'turn', turnItemId: agentJournalItemKey(turnRecord(turnId)) } }
)
}
/** A running turn whose record sits `rows` rows above the journal head. */
async function runLongTurn(turnId: string, rows: number): Promise<void> {
await writeTurn(turnId, { state: 'running', startedAt: 1_000, requestedAt: 900 })
const turnScope = { kind: 'turn' as const, turnItemId: agentJournalItemKey(turnRecord(turnId)) }
for (let ordinal = 1; ordinal <= rows; ordinal += 1) {
await journal.appendItem(
turnRecord(turnId, ordinal),
{ kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: `row ${ordinal}` }] },
{ fence: 1, turnScope }
)
}
}
function hydrate(): StructuredAgentSessionState {
const page = readAgentSessionHydrationPage(journal, 1)
return reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
type: 'event',
event: { type: 'snapshot', sessionId: 'session-1', page, fence: 1 }
})
}
/** A live subscriber attached at `state`'s cursor; returns a pump that applies what it was sent. */
function subscribeFrom(state: StructuredAgentSessionState) {
const events: AgentSessionSubscribeEvent[] = []
const subscribers = new AgentSessionSubscribers()
subscribers.open({
id: 'client-1',
sessionId: 'session-1',
journal,
fence: 1,
emit: (event) => events.push(event),
...(state.cursor ? { cursor: state.cursor } : {})
})
let current = state
return {
subscribers,
events,
apply(): StructuredAgentSessionState {
for (const event of events.splice(0)) {
current = reduceStructuredAgentSession(current, { type: 'event', event })
}
return current
}
}
}
describe('a running turn whose record is older than the loaded page', () => {
it('is running in the chat, as it is on the host', async () => {
await runLongTurn('turn-1', 250)
const state = hydrate()
expect(journal.activeTurnId()).toBe('turn-1')
expect(state.hasOlder).toBe(true)
// What every client read before: the record is off the page, so the loaded rows said idle.
expect(activeStructuredAgentSessionTurnId(state.items)).toBeNull()
expect(runningStructuredAgentSessionTurnId(state)).toBe('turn-1')
})
it('counts Working from the host record when the record is not loaded', async () => {
await runLongTurn('turn-1', 250)
const state = hydrate()
const turnId = runningStructuredAgentSessionTurnId(state)
const { runningTiming } = selectStructuredAgentTurnBars(
state.items,
state.submissions,
turnId,
state.latestTurn
)
expect(runningTiming).toMatchObject({ state: 'running', startedAt: 1_000, requestedAt: 900 })
})
it('stops reading as running when the turn ends off the page', async () => {
await runLongTurn('turn-1', 250)
const live = subscribeFrom(hydrate())
await writeTurn('turn-1', { state: 'completed', startedAt: 1_000, completedAt: 2_000 })
live.subscribers.publish('session-1', journal)
const state = live.apply()
// The completing revision keeps its old position, so the window never admits it.
expect(
state.items.some((item) => item.itemId === agentJournalItemKey(turnRecord('turn-1')))
).toBe(false)
expect(runningStructuredAgentSessionTurnId(state)).toBeNull()
expect(state.latestTurn?.turn).toMatchObject({ turnId: 'turn-1', state: 'completed' })
})
it('follows the next turn as it opens', async () => {
await runLongTurn('turn-1', 250)
await writeTurn('turn-1', { state: 'completed', startedAt: 1_000, completedAt: 2_000 })
const live = subscribeFrom(hydrate())
await writeTurn('turn-2', { state: 'running', startedAt: 3_000 })
live.subscribers.publish('session-1', journal)
expect(runningStructuredAgentSessionTurnId(live.apply())).toBe('turn-2')
})
it('is not undone by an older page read while live frames moved on', async () => {
await runLongTurn('turn-1', 250)
const before = hydrate()
const older = readAgentSessionHistory(journal, {
sessionId: 'session-1',
direction: 'before',
cursor: { epoch: before.epoch ?? '', sequence: before.items[0]?.sequence ?? 0 },
limit: 40
})
const live = subscribeFrom(before)
await writeTurn('turn-1', { state: 'completed', startedAt: 1_000, completedAt: 2_000 })
live.subscribers.publish('session-1', journal)
const ended = live.apply()
if (!older.ok) {
throw new Error(`expected an older page, got reset ${older.reset}`)
}
// The older page was read while the turn still ran.
expect(older.page.latestTurn?.turn.state).toBe('running')
const after = reduceStructuredAgentSession(ended, {
type: 'older-page',
requestedCursor: { epoch: ended.epoch ?? '', sequence: ended.items[0]?.sequence ?? 0 },
page: older.page
})
expect(runningStructuredAgentSessionTurnId(after)).toBeNull()
})
})
describe('an older host, which publishes no latest turn', () => {
it('still reads the turn off the loaded rows', async () => {
await writeTurn('turn-1', { state: 'running', startedAt: 1_000 })
const { latestTurn: _omitted, ...page } = readAgentSessionHydrationPage(journal, 1)
const state = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
type: 'event',
event: { type: 'snapshot', sessionId: 'session-1', page, fence: 1 }
})
expect(state.latestTurn).toBeUndefined()
expect(runningStructuredAgentSessionTurnId(state)).toBe('turn-1')
})
it('reads a journal with no turn as idle, not unknown', async () => {
const state = hydrate()
expect(state.latestTurn).toBeNull()
expect(runningStructuredAgentSessionTurnId(state)).toBeNull()
})
})
@@ -34,7 +34,6 @@ import {
} from './native-chat-transcript-slots'
import { useNativeChatTranscriptSlots } from './use-native-chat-transcript-slots'
import { useNativeChatTranscriptWindow } from './use-native-chat-transcript-window'
import { nativeChatRowsInTranscriptOrder } from './native-chat-subagent-sections'
import { useNativeChatSubagentSections } from './use-native-chat-subagent-sections'
import { toggleNativeChatExpandedKey } from './native-chat-expanded-keys'
import { useNativeChatTurnMembership } from './use-native-chat-turn-membership'
@@ -54,14 +53,11 @@ import type {
AgentJournalRenderItem,
AgentJournalSubmission
} from '../../../../shared/agent-session-journal-types'
import type { AgentSessionLatestTurn } from '../../../../shared/agent-session-wire'
import { isStructuredAgentSessionThinking } from '../../../../shared/structured-agent-session-live-turn'
import type { NativeChatSettledTurns } from '../../../../shared/native-chat-turn-status'
import {
nativeChatTurnDiffs,
type NativeChatDiffReveal,
type NativeChatDiffTarget,
type NativeChatTurnDiff
} from './native-chat-turn-diffs'
import type { NativeChatDiffReveal, NativeChatDiffTarget } from './native-chat-turn-diffs'
import { useNativeChatTurnDiffs } from './use-native-chat-turn-diffs'
/** The turn is blocked on the reader. `shown`: the pane draws the prompt itself, as a card;
* `unshown`: it cannot (the prompt is only in the agent's terminal). */
@@ -75,6 +71,7 @@ export function NativeChatMessageList({
session,
journalItems,
journalSubmissions,
journalLatestTurn: latestTurn,
subagentRoster,
railOutline = null,
isVisible = true,
@@ -93,6 +90,8 @@ export function NativeChatMessageList({
journalItems?: readonly AgentJournalRenderItem[]
/** With the items, what places each row in its turn (structured lane). */
journalSubmissions?: readonly AgentJournalSubmission[]
/** The host's newest turn record, which places a live turn whose record is not loaded. */
journalLatestTurn?: AgentSessionLatestTurn | null
/** Every subagent the session's rosters named, whether or not its roster row is loaded. */
subagentRoster?: Parameters<typeof useNativeChatSubagentSections>[2]
/** User messages older than the loaded window, from the host's outline. */
@@ -139,7 +138,8 @@ export function NativeChatMessageList({
const { messages, subagentRows } = useNativeChatTranscriptProjection(
session,
journalItems,
journalSubmissions
journalSubmissions,
latestTurn
)
const {
sections: subagentSections,
@@ -151,23 +151,23 @@ export function NativeChatMessageList({
const taskListPredecessors = useMemo(() => nativeChatTaskListPredecessors(messages), [messages])
const taskListState = useMemo(() => nativeChatTaskListState(messages), [messages])
// Each row's turn, which turn is live, and the order the rows draw in, resolved once.
const {
messages: rows,
turnKeys,
liveTurnKey
} = useNativeChatTurnMembership(messages, journalItems, journalSubmissions)
const turnDiffs = useMemo(() => {
if (!journalItems) {
return new Map<string, NativeChatTurnDiff>()
}
const merged = nativeChatRowsInTranscriptOrder(rows, turnKeys, subagentRowsInOrder)
return nativeChatTurnDiffs(merged.messages, merged.turnKeys, subagentSections.pathOf)
}, [journalItems, rows, subagentRowsInOrder, subagentSections.pathOf, turnKeys])
const turnRows = useNativeChatTurnMembership(
messages,
journalItems,
journalSubmissions,
latestTurn
)
const { messages: rows, turnKeys, liveTurnKey } = turnRows
const turnDiffs = useNativeChatTurnDiffs(
journalItems ? turnRows : null,
subagentRowsInOrder,
subagentSections.pathOf
)
// "Thinking" is real reasoning content at the tail of the turn, not the absence
// of output — the latter reports thinking while the request is merely in flight.
const thinking = useMemo(
() => (journalItems ? isStructuredAgentSessionThinking(journalItems) : false),
[journalItems]
() => isStructuredAgentSessionThinking({ items: journalItems ?? [], latestTurn }),
[journalItems, latestTurn]
)
const turnStatuses = useNativeChatTurnStatus({
turnKeys,
@@ -0,0 +1,119 @@
// @vitest-environment happy-dom
// A long turn's record and opening message sit above the loaded page. The host's newest turn record
// still names the turn, so its loaded rows draw under the live Working bar.
import '@testing-library/jest-dom/vitest'
import { cleanup, render, screen } from '@testing-library/react'
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'
import type {
AgentJournalRenderItem,
AgentJournalTurnScope
} from '../../../../shared/agent-session-journal-types'
import type { AgentSessionLatestTurn } from '../../../../shared/agent-session-wire'
import { projectStructuredAgentSessionMessages } from '../../../../shared/structured-agent-session-message-projection'
import { NativeChatMessageList } from './NativeChatMessageList'
import { installNativeChatMessageListTestViewport } from './native-chat-message-list-test-viewport'
let restoreViewport = (): void => {}
beforeAll(() => {
restoreViewport = installNativeChatMessageListTestViewport()
})
afterAll(() => restoreViewport())
afterEach(cleanup)
const TURN_RECORD = 'turn-record-1'
const IN_TURN: AgentJournalTurnScope = { kind: 'turn', turnItemId: TURN_RECORD }
/** The newest rows of the turn: everything from its 300th row on. */
const loaded: AgentJournalRenderItem[] = [300, 301, 302].map((sequence) => ({
itemId: `row-${sequence}`,
revision: 0,
sequence,
observedAt: sequence,
body: {
kind: 'message',
role: 'assistant',
blocks: [{ type: 'text', text: `Step ${sequence}` }]
},
turnScope: IN_TURN
}))
const patch = '@@ -1 +1 @@\n-before\n+after'
/** An edit among the loaded rows; the turn's earlier edits are above them. */
const edit: AgentJournalRenderItem = {
itemId: 'row-303',
revision: 0,
sequence: 303,
observedAt: 303,
body: {
kind: 'diff',
path: 'src/a.ts',
patch: { head: patch, truncated: false, digest: 'fixture', byteLength: patch.length }
},
turnScope: IN_TURN
}
const running: AgentSessionLatestTurn = {
itemId: TURN_RECORD,
observedAt: 1,
turn: { turnId: 'turn-1', state: 'running', startedAt: 1_000, userItemId: 'user-1' }
}
function list(
latestTurn: AgentSessionLatestTurn | null | undefined,
items: AgentJournalRenderItem[] = loaded
): React.JSX.Element {
return (
<NativeChatMessageList
session={{
messages: projectStructuredAgentSessionMessages(items, [], [], { rejectedInPlace: true }),
status: 'ready',
sessionId: 'session-1',
agent: 'codex',
hasMore: true,
loadingEarlier: false,
olderHistoryGeneration: 0,
loadEarlier: vi.fn(),
readPhase: 'ready'
}}
journalItems={items}
journalSubmissions={[]}
journalLatestTurn={latestTurn}
isWorking
workingStartedAt={1_000}
expandSignal={false}
/>
)
}
describe('a live turn whose record and opening message are not loaded', () => {
it('draws its Working bar with the running clock', () => {
const now = vi.spyOn(Date, 'now').mockReturnValue(64_000)
try {
render(list(running))
expect(screen.getByText('Step 302')).toBeInTheDocument()
expect(screen.getByText(/Working for 1m 3s/)).toBeInTheDocument()
} finally {
now.mockRestore()
}
})
it('totals no edits, since only the end of the turn is loaded', () => {
render(list(running, [...loaded, edit]))
expect(screen.getByText(/Working for/)).toBeInTheDocument()
expect(screen.queryByRole('button', { name: /changed file/ })).toBeNull()
})
it('had no bar to draw from the loaded rows alone', () => {
const now = vi.spyOn(Date, 'now').mockReturnValue(64_000)
try {
render(list(undefined))
expect(screen.getByText('Step 302')).toBeInTheDocument()
expect(screen.queryByText(/Working for/)).not.toBeInTheDocument()
} finally {
now.mockRestore()
}
})
})
@@ -286,6 +286,7 @@ export function NativeChatStructuredSession(
session={session}
journalItems={controller.journalItems}
journalSubmissions={controller.submissions}
journalLatestTurn={controller.latestTurn}
subagentRoster={controller.subagentRoster}
railOutline={controller.railOutline}
isVisible={props.isVisible}
@@ -30,12 +30,14 @@ export type NativeChatTurnDiff = {
export function nativeChatTurnDiffs(
messages: readonly NativeChatMessage[],
turnKeys: readonly (string | undefined)[],
subagentSectionsOf?: ReadonlyMap<string, readonly string[]>
subagentSectionsOf?: ReadonlyMap<string, readonly string[]>,
/** A turn only partly loaded, whose totals would read as the whole turn's. */
partialTurnKey?: string
): Map<string, NativeChatTurnDiff> {
const turns = new Map<string, Map<string, NativeChatTurnDiffFile>>()
for (const [index, message] of messages.entries()) {
const turnKey = turnKeys[index]
if (!turnKey) {
if (!turnKey || turnKey === partialTurnKey) {
continue
}
const sections = subagentSectionsOf?.get(message.id)
@@ -5,6 +5,7 @@ import type {
} from '../../../../shared/agent-session-journal-types'
import type { NativeChatSubagentRow } from '../../../../shared/native-chat-transcript-projection'
import type { NativeChatMessage } from '../../../../shared/native-chat-types'
import type { AgentSessionLatestTurn } from '../../../../shared/agent-session-wire'
import { createNativeChatMessageListProjection } from './native-chat-message-list-projection'
import { projectNativeChatTaskListFrames } from './native-chat-task-list-frames'
import { omitNativeChatThreadGoalRows } from './native-chat-thread-goal-rows'
@@ -15,7 +16,8 @@ import type { NativeChatLiveSession } from './use-native-chat-live-session'
export function useNativeChatTranscriptProjection(
session: NativeChatLiveSession,
journalItems: readonly AgentJournalRenderItem[] | undefined,
journalSubmissions: readonly AgentJournalSubmission[] | undefined
journalSubmissions: readonly AgentJournalSubmission[] | undefined,
latestTurn?: AgentSessionLatestTurn | null
): {
messages: NativeChatMessage[]
subagentRows: ReadonlyMap<string, readonly NativeChatSubagentRow[]>
@@ -30,9 +32,11 @@ export function useNativeChatTranscriptProjection(
() =>
projectMessages(
session.messages,
journalItems ? { items: journalItems, submissions: journalSubmissions ?? [] } : null
journalItems
? { items: journalItems, submissions: journalSubmissions ?? [], latestTurn }
: null
),
[journalItems, journalSubmissions, projectMessages, session.messages]
[journalItems, journalSubmissions, latestTurn, projectMessages, session.messages]
)
const messages = useMemo(() => {
const projected = projectNativeChatTaskListFrames(projection.conversation)
@@ -0,0 +1,28 @@
import { useMemo } from 'react'
import type { NativeChatSubagentRow } from '../../../../shared/native-chat-transcript-projection'
import { nativeChatRowsInTranscriptOrder } from './native-chat-subagent-sections'
import { nativeChatTurnDiffs, type NativeChatTurnDiff } from './native-chat-turn-diffs'
import type { NativeChatTurnRows } from './use-native-chat-turn-membership'
const NO_TURN_DIFFS: ReadonlyMap<string, NativeChatTurnDiff> = new Map()
/** Each turn's recorded edit totals, a subagent's edits counted in the turn they happened. None
* without a journal (`turns` null), and none for a turn only partly loaded. */
export function useNativeChatTurnDiffs(
turns: NativeChatTurnRows | null,
subagentRows: readonly NativeChatSubagentRow[],
subagentSectionsOf: ReadonlyMap<string, readonly string[]>
): ReadonlyMap<string, NativeChatTurnDiff> {
return useMemo(() => {
if (!turns) {
return NO_TURN_DIFFS
}
const merged = nativeChatRowsInTranscriptOrder(turns.messages, turns.turnKeys, subagentRows)
return nativeChatTurnDiffs(
merged.messages,
merged.turnKeys,
subagentSectionsOf,
turns.partialTurnKey
)
}, [subagentRows, subagentSectionsOf, turns])
}
@@ -4,6 +4,7 @@ import type {
AgentJournalSubmission
} from '../../../../shared/agent-session-journal-types'
import type { NativeChatMessage } from '../../../../shared/native-chat-types'
import type { AgentSessionLatestTurn } from '../../../../shared/agent-session-wire'
import { nativeChatTurnMembership } from '../../../../shared/native-chat-turn-membership'
import { nativeChatRowsInDrawOrder } from '../../../../shared/native-chat-turn-grouping'
@@ -13,6 +14,8 @@ export type NativeChatTurnRows = {
/** Each drawn row's turn, by index into `messages`. */
turnKeys: readonly (string | undefined)[]
liveTurnKey: string | undefined
/** The live turn, when only its tail is loaded (`NativeChatTurnMembership.partialTurnKey`). */
partialTurnKey?: string
}
/** Each row's turn, which turn is live, and the order the rows draw in, resolved once: from the
@@ -20,17 +23,21 @@ export type NativeChatTurnRows = {
export function useNativeChatTurnMembership(
messages: readonly NativeChatMessage[],
journalItems: readonly AgentJournalRenderItem[] | undefined,
journalSubmissions: readonly AgentJournalSubmission[] | undefined
journalSubmissions: readonly AgentJournalSubmission[] | undefined,
latestTurn?: AgentSessionLatestTurn | null
): NativeChatTurnRows {
return useMemo(() => {
const { turnKeys, liveTurnKey, drawOrder } = nativeChatTurnMembership(
const { turnKeys, liveTurnKey, drawOrder, partialTurnKey } = nativeChatTurnMembership(
messages,
journalItems ? { items: journalItems, submissions: journalSubmissions ?? [] } : null
journalItems
? { items: journalItems, submissions: journalSubmissions ?? [], latestTurn }
: null
)
return {
messages: nativeChatRowsInDrawOrder(messages, drawOrder),
turnKeys: nativeChatRowsInDrawOrder(turnKeys, drawOrder),
liveTurnKey
liveTurnKey,
...(partialTurnKey !== undefined ? { partialTurnKey } : {})
}
}, [journalItems, journalSubmissions, messages])
}, [journalItems, journalSubmissions, latestTurn, messages])
}
@@ -11,6 +11,7 @@ import type {
AgentJournalSubmission
} from '../../../../shared/agent-session-journal-types'
import type { StructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox'
import type { AgentSessionLatestTurn } from '../../../../shared/agent-session-wire'
const mocks = vi.hoisted(() => ({
call: vi.fn(),
@@ -21,6 +22,7 @@ let items: AgentJournalRenderItem[] = []
let submissions: AgentJournalSubmission[] = []
let outbox: StructuredAgentSessionOutboxEntry[] = []
let fence = 3
let latestTurn: AgentSessionLatestTurn | null | undefined
vi.mock('@/runtime/structured-agent-session-client', () => ({
callStructuredAgentSession: mocks.call,
@@ -30,7 +32,15 @@ vi.mock('@/runtime/structured-agent-session-client', () => ({
vi.mock('./use-structured-agent-session-read', () => ({
useStructuredAgentSessionRead: () => ({
state: { fence, items, submissions, status: 'ready', error: null, hasOlder: false },
state: {
fence,
items,
submissions,
status: 'ready',
error: null,
hasOlder: false,
...(latestTurn !== undefined ? { latestTurn } : {})
},
loadingOlder: false,
loadOlder: vi.fn()
})
@@ -139,6 +149,7 @@ beforeEach(() => {
submissions = []
outbox = []
fence = 3
latestTurn = undefined
})
describe('Stop against a host that stops the conversation', () => {
@@ -542,3 +553,42 @@ describe.each([
expect(mocks.withdrawUnsent).not.toHaveBeenCalled()
})
})
// A long turn's record sits above the loaded page; the host's record is what says it runs.
describe('a running turn whose record is not loaded', () => {
beforeEach(() => {
items = []
latestTurn = {
itemId: 'turn-1',
observedAt: 1,
turn: { turnId: 'provider-turn', state: 'running', startedAt: 1 }
}
})
it('shows Working and Stop, and an older host gets the turn by name', async () => {
setLocalRuntimeCapabilitiesForTests([])
const { result } = render()
expect(result.current.isWorking).toBe(true)
expect(result.current.turnId).toBe('provider-turn')
expect(result.current.canStop).toBe(true)
await act(async () => {
await result.current.stop()
})
expect(cancels()).toEqual([expect.objectContaining({ turnId: 'provider-turn' })])
})
it('reads idle once the host ends it, though a running record is still loaded', () => {
items = [RUNNING_TURN]
latestTurn = {
itemId: 'turn-1',
observedAt: 2,
turn: { turnId: 'provider-turn', state: 'completed', startedAt: 1, completedAt: 2 }
}
setLocalRuntimeCapabilitiesForTests([])
const { result } = render()
expect(result.current.isWorking).toBe(false)
expect(result.current.canStop).toBe(false)
})
})
@@ -1,5 +1,5 @@
import { useMemo } from 'react'
import { activeStructuredAgentSessionTurnId } from '../../../../shared/structured-agent-session-projection'
import { runningStructuredAgentSessionTurnId } from '../../../../shared/structured-agent-session-live-turn'
import { isStructuredAgentSessionMainAgentWorking } from '../../../../shared/structured-agent-session-main-agent-working'
import type { StructuredAgentSessionState } from '../../../../shared/structured-agent-session-reducer'
import type { StructuredAgentSubagentRoster } from '../../../../shared/structured-agent-session-subagent-roster'
@@ -19,7 +19,9 @@ export function useStructuredAgentSessionTransportState(
const submissions = enabled ? state.submissions : NO_SUBMISSIONS
const subagentRoster = (enabled ? state.subagentRoster : undefined) ?? NO_SUBAGENT_ROSTER
const fence = enabled ? state.fence : null
const turnId = activeStructuredAgentSessionTurnId(journalItems)
const latestTurn = enabled ? state.latestTurn : undefined
// The host's whole-journal turn, never the loaded rows': a long turn's record is off the page.
const turnId = runningStructuredAgentSessionTurnId({ items: journalItems, latestTurn })
// The rule the host projects every session list's Working from, so this chat cannot disagree.
const isWorking = isStructuredAgentSessionMainAgentWorking(turnId, submissions, fence)
const turnActivity = useMemo(
@@ -30,12 +32,14 @@ export function useStructuredAgentSessionTransportState(
{
items: journalItems,
submissions,
latestTurn,
...(enabled ? { hostClock: state.hostClock } : {})
},
turnId
)
return {
journalItems,
latestTurn,
subagentRoster,
submissions,
fence,
@@ -233,6 +233,8 @@ export function useStructuredAgentSession(args: {
)
}),
journalItems: transcriptItems,
/** The host's newest turn record, which places a live turn whose record is not loaded. */
latestTurn: transportState.latestTurn,
subagentRoster: transportState.subagentRoster,
messages,
status: transportEnabled ? state.status : 'ready',
@@ -4,6 +4,7 @@ import type {
AgentJournalSubmission
} from '../../../../shared/agent-session-journal-types'
import type { NativeChatSettledTurns } from '../../../../shared/native-chat-turn-status'
import type { AgentSessionLatestTurn } from '../../../../shared/agent-session-wire'
import type { StructuredAgentHostClock } from '../../../../shared/structured-agent-session-reducer'
import { selectStructuredAgentTurnBars } from '../../../../shared/structured-agent-session-turn-timing'
import {
@@ -18,17 +19,19 @@ export function useStructuredAgentTurnTiming(
{
items,
submissions,
latestTurn,
hostClock
}: {
items: readonly AgentJournalRenderItem[]
submissions: readonly AgentJournalSubmission[]
latestTurn?: AgentSessionLatestTurn | null
hostClock?: StructuredAgentHostClock | null
},
turnId: string | null
): { settledTurns: NativeChatSettledTurns; workingStartedAt: number | null } {
const { settledTurns, runningTiming } = useMemo(
() => selectStructuredAgentTurnBars(items, submissions, turnId),
[items, submissions, turnId]
() => selectStructuredAgentTurnBars(items, submissions, turnId, latestTurn),
[items, submissions, turnId, latestTurn]
)
const [latch, setLatch] = useState<StructuredAgentTurnClockLatch | null>(null)
// Stamp during render (React's derive-from-props pattern) so the first paint of
@@ -0,0 +1,49 @@
import type { AgentJournalTurnOutcome } from './agent-session-journal-types'
import { agentSessionScopeKey, type AgentSessionExecutionLocation } from './agent-session-record'
// Turn completion feed: the per-turn EDGE beside the status feed's STATE.
/**
* The session's latest request reaching a terminal outcome — a root turn, or a send the agent or
* its start refused — derived by the EXECUTION HOST at journal commit.
*
* This is the EDGE, with turn identity; `AgentSessionStatusSummary.turnOutcome` is the STATE.
* The summary carries the verdict only while the session is idle, as a fact about the main agent's
* last turn that a status reader may act on (attention alerts, the `mainAgent.outcome` row field),
* and never a turn id: a reader that needs to know WHICH turn finished, or to react exactly once
* per finish, subscribes here. Re-broadcasting the summary on every status change therefore
* repeats a state, not a completion.
*
* `outcome` is the journal's recorded verdict (the provider's, a stop, or the host's supersede) and
* is never inferred — a turn the host only observed ending carries no outcome and produces no event
* at all, because absent means UNKNOWN, not success.
*/
export type AgentSessionTurnCompletion = {
/** Host-and-workspace scope; a bare provider turn id is not globally unique. */
scope: AgentSessionExecutionLocation
sessionId: string
/** The request's identity: the root turn's id, or for a send refused before any turn, that
* send's journal item key. Neither is minted here. */
turnId: string
outcome: AgentJournalTurnOutcome
/** Execution host's clock at journal commit. */
completedAt: number
/** The request settled while a prompt waits on the user. Absent otherwise, and from older hosts. */
awaitingUser?: true
}
/**
* LIVE-ONLY: there is no snapshot arm and no replay arm, by decision. A subscriber is told what
* completes while it is subscribed and nothing else; completions that land while it is away are
* dropped rather than queued, so nothing durable can strand. On reconnect the client baselines.
*/
export type AgentSessionTurnCompletionEvent =
| { type: 'completion'; completion: AgentSessionTurnCompletion }
| { type: 'end' }
/** Delivery dedupe address. Unread is idempotent and does not need it; mobile fanout does. */
export function agentSessionTurnCompletionKey(completion: AgentSessionTurnCompletion): string {
return [agentSessionScopeKey(completion.scope), completion.sessionId, completion.turnId].join(
'\u0000'
)
}
+20 -54
View File
@@ -12,6 +12,7 @@ import type {
export * from './agent-session-wire-refusals'
export * from './agent-session-queued-message-wire'
export * from './agent-session-turn-completion-wire'
import type { AgentSessionConversationCommand } from './agent-session-conversation-command'
import type { AgentSessionContextUsage } from './agent-session-context-usage'
// ─── Structured agent-session wire contract ─────────────────────────────────
@@ -28,15 +29,10 @@ import type {
AgentJournalResolution,
AgentJournalSubmission,
AgentJournalThreadGoal,
AgentJournalTurnOutcome
AgentJournalTurnLifecycle
} from './agent-session-journal-types'
import type { AgentTurnOutcome } from './agent-turn-outcome'
import {
agentSessionScopeKey,
type AgentSessionExecutionLocation,
type AgentSessionHandoffStage,
type AgentSessionRecord
} from './agent-session-record'
import type { AgentSessionHandoffStage, AgentSessionRecord } from './agent-session-record'
import type { AgentProviderSessionMetadata } from './agent-session-resume'
import type { NativeChatSubagentEntry } from './native-chat-types'
import type { StructuredAgentSessionProjectedStatus } from './structured-agent-session-projection'
@@ -64,6 +60,18 @@ export type AgentSessionTurnActivity = {
text: string
}
/** The session's newest turn record over the WHOLE journal. A page windows the timeline and a
* turn's record keeps the place it opened at, so a long turn's record falls off the page; this is
* what tells a client a turn is running. Present null: the journal records no turn. Absent: an
* older host, whose clients still read the loaded rows. */
export type AgentSessionLatestTurn = {
/** The record's journal key, which rows of the turn name as their scope. */
itemId: string
/** Host clock at the record's creation, as on its own row; a revision does not move it. */
observedAt: number
turn: AgentJournalTurnLifecycle
}
export const AGENT_SESSION_ID_MAX_LENGTH = 512
/** Backward paging is the client's normal read; 40 matches the page size the
@@ -133,6 +141,8 @@ export type AgentSessionHistoryPage = {
/** Names the subagents with rows on the page whose roster row is older than it; bounded.
* Absent from older hosts, and when every such roster row is on the page. */
subagentRoster?: AgentSessionSubagentRosterEntry[]
/** As of the page's read; a client applies it only from a page that replaces its state. */
latestTurn?: AgentSessionLatestTurn | null
}
export type AgentSessionHistoryResult =
@@ -192,6 +202,9 @@ export type AgentSessionSubscribeEvent =
commands?: AgentSessionSlashCommand[] | null
/** Additive ephemeral state; it never creates or advances journal rows. */
activity?: AgentSessionTurnActivity | null
/** Rides every batch that carries rows, removals or submissions, so absent there means an
* older host; absent on one that carries none, which changes no turn. */
latestTurn?: AgentSessionLatestTurn | null
} & AgentSessionHostClockField)
| ({
type: 'reset'
@@ -267,53 +280,6 @@ export type AgentSessionStatusEvent =
| { type: 'status'; session: AgentSessionStatusSummary }
| { type: 'end' }
// ─── Turn completion feed ───────────────────────────────────────────────────
/**
* The session's latest request reaching a terminal outcome — a root turn, or a send the agent or
* its start refused — derived by the EXECUTION HOST at journal commit.
*
* This is the EDGE, with turn identity; `AgentSessionStatusSummary.turnOutcome` is the STATE.
* The summary carries the verdict only while the session is idle, as a fact about the main agent's
* last turn that a status reader may act on (attention alerts, the `mainAgent.outcome` row field),
* and never a turn id: a reader that needs to know WHICH turn finished, or to react exactly once
* per finish, subscribes here. Re-broadcasting the summary on every status change therefore
* repeats a state, not a completion.
*
* `outcome` is the journal's recorded verdict (the provider's, a stop, or the host's supersede) and
* is never inferred — a turn the host only observed ending carries no outcome and produces no event
* at all, because absent means UNKNOWN, not success.
*/
export type AgentSessionTurnCompletion = {
/** Host-and-workspace scope; a bare provider turn id is not globally unique. */
scope: AgentSessionExecutionLocation
sessionId: string
/** The request's identity: the root turn's id, or for a send refused before any turn, that
* send's journal item key. Neither is minted here. */
turnId: string
outcome: AgentJournalTurnOutcome
/** Execution host's clock at journal commit. */
completedAt: number
/** The request settled while a prompt waits on the user. Absent otherwise, and from older hosts. */
awaitingUser?: true
}
/**
* LIVE-ONLY: there is no snapshot arm and no replay arm, by decision. A subscriber is told what
* completes while it is subscribed and nothing else; completions that land while it is away are
* dropped rather than queued, so nothing durable can strand. On reconnect the client baselines.
*/
export type AgentSessionTurnCompletionEvent =
| { type: 'completion'; completion: AgentSessionTurnCompletion }
| { type: 'end' }
/** Delivery dedupe address. Unread is idempotent and does not need it; mobile fanout does. */
export function agentSessionTurnCompletionKey(completion: AgentSessionTurnCompletion): string {
return [agentSessionScopeKey(completion.scope), completion.sessionId, completion.turnId].join(
'\u0000'
)
}
// ─── Mutation envelope ──────────────────────────────────────────────────────
/**
@@ -200,6 +200,39 @@ describe('the live turn', () => {
})
})
describe('a turn whose record is above the loaded rows', () => {
const tail = () => [assistant('a300', inTurn('t1')), assistant('a301', inTurn('t1'))]
const hostTurn = (state: 'running' | 'completed', itemId = 't1') => ({
itemId,
observedAt: 1,
turn: { turnId: itemId, state, userItemId: 'u1' }
})
const membership = (
items: readonly AgentJournalRenderItem[],
latestTurn: ReturnType<typeof hostTurn>
) => nativeChatTurnMembership(rows(items), { items, submissions: [], latestTurn })
it('groups its tail under the live bar while it runs, marked partial', () => {
const live = membership(tail(), hostTurn('running'))
expect(live.turnKeys).toEqual(['t1', 't1'])
expect(live.liveTurnKey).toBe('t1')
expect(live.partialTurnKey).toBe('t1')
})
it('leaves its tail ungrouped once it ends, whichever turn is newest', () => {
expect(membership(tail(), hostTurn('completed')).turnKeys).toEqual([undefined, undefined])
const next = membership([...tail(), user('u2')], hostTurn('running', 't2'))
expect(next.turnKeys.slice(0, 2)).toEqual([undefined, undefined])
})
it('is whole once its record loads', () => {
const items = [user('u1'), turn('t1', 'u1', THREAD, 'running'), ...tail()]
const loaded = membership(items, hostTurn('running'))
expect(loaded.turnKeys).toEqual(['u1', 'u1', 'u1'])
expect(loaded.partialTurnKey).toBeUndefined()
})
})
// A message sent while A runs, which the provider queued behind A: its turn opens only after A's
// remaining rows, which the journal wrote after it.
describe('a message the provider answered after the running turn', () => {
+27 -5
View File
@@ -26,7 +26,11 @@ import {
} from './native-chat-turn-grouping'
import type { NativeChatRole } from './native-chat-types'
import { isStructuredAgentSessionCommandTurn } from './structured-agent-session-command-entry'
import { liveStructuredAgentSessionTurnScope } from './structured-agent-session-live-turn'
import {
liveStructuredAgentSessionTurnScope,
runningStructuredAgentSessionTurnScope
} from './structured-agent-session-live-turn'
import type { AgentSessionLatestTurn } from './agent-session-wire'
/** Whether the host writing this journal states each row's turn. Only a host that runs `/compact`
* as a turn of the send path does, so this is also how a client tells that host from an older one. */
@@ -37,6 +41,10 @@ export function hostStatesTurnScopes(items: readonly AgentJournalRenderItem[]):
export type NativeChatTurnJournal = {
items: readonly AgentJournalRenderItem[]
submissions: readonly AgentJournalSubmission[]
/** The host's newest turn record: it names the live turn, and anchors it while it runs when the
* record is not loaded, so the turn's loaded rows still draw under its bar. Absent from older
* hosts. */
latestTurn?: AgentSessionLatestTurn | null
}
/**
@@ -48,7 +56,9 @@ export type NativeChatTurnJournal = {
*/
export function structuredAgentTurnAnchors(
items: readonly AgentJournalRenderItem[],
submissions: readonly AgentJournalSubmission[] = []
submissions: readonly AgentJournalSubmission[] = [],
/** Anchored too while it runs with its record above the loaded rows; nothing loaded precedes it. */
latestTurn?: AgentSessionLatestTurn | null
): ReadonlyMap<string, string> {
const userItemIds = new Set(
items.flatMap((item) =>
@@ -89,6 +99,13 @@ export function structuredAgentTurnAnchors(
)
inFlightSinceLastTurn = null
}
// Only while running: an ended turn's tail would otherwise regroup when the next turn opens.
if (latestTurn?.turn.state === 'running' && !anchors.has(latestTurn.itemId)) {
anchors.set(
latestTurn.itemId,
anchorOf(latestTurn.itemId, latestTurn.turn, userItemIds, aliases, null, null)
)
}
return anchors
}
@@ -123,6 +140,9 @@ export type NativeChatTurnMembership = {
/** Row indexes in the order the transcript draws them (`nativeChatTurnDrawOrder`), or null when
* that is the order given. */
drawOrder: readonly number[] | null
/** The live turn's key while its record is above the loaded rows: those rows are only its tail,
* so nothing may total them as the whole turn. */
partialTurnKey?: string
}
type NativeChatTurnMember = { id: string; role: NativeChatRole; unsent?: true }
@@ -144,8 +164,8 @@ export function nativeChatTurnMembership(
const turnKeys = withoutUnsent(messages, nativeChatRowTurnKeys(messages, null, opens))
return { turnKeys, liveTurnKey: newestUserTurnKey(messages, turnKeys), drawOrder: null }
}
const anchors = structuredAgentTurnAnchors(journal.items, journal.submissions)
const running = liveStructuredAgentSessionTurnScope(journal.items)
const anchors = structuredAgentTurnAnchors(journal.items, journal.submissions, journal.latestTurn)
const running = runningStructuredAgentSessionTurnScope(journal)
if (!hostStatesTurnScopes(journal.items)) {
const recordKeys = namedRecordKeys(journal.items, anchors)
const turnKeys = withoutUnsent(
@@ -166,6 +186,7 @@ export function nativeChatTurnMembership(
const runningKey = running.kind === 'turn' ? anchors.get(running.turnItemId) : undefined
const anchoring = new Set(anchors.values())
const scopes = new Map(journal.items.map((item) => [item.itemId, item.turnScope]))
const runningRecordLoaded = running.kind === 'turn' && scopes.has(running.turnItemId)
const turnKeys = messages.map((message) => {
if (message.unsent === true) {
return undefined
@@ -181,7 +202,8 @@ export function nativeChatTurnMembership(
return {
turnKeys,
liveTurnKey: runningKey ?? newestUserTurnKey(messages, turnKeys),
drawOrder: nativeChatTurnDrawOrder(messages, turnKeys, anchoring)
drawOrder: nativeChatTurnDrawOrder(messages, turnKeys, anchoring),
...(runningKey !== undefined && !runningRecordLoaded ? { partialTurnKey: runningKey } : {})
}
}
@@ -125,4 +125,58 @@ describe('structured agent session event coalescer', () => {
coalescer.flush()
expect(events[1]).toMatchObject({ queuePause: { reason: 'stopped' } })
})
it('keeps the latest turn a coalesced frame carried, and a null one as an answer', () => {
const events: AgentSessionSubscribeEvent[] = []
const coalescer = createStructuredAgentSessionEventCoalescer((event) => events.push(event))
const running = {
itemId: 'turn-1',
observedAt: 1,
turn: { turnId: 'turn-1', state: 'running' as const }
}
coalescer.push({ ...batch(1), latestTurn: running })
coalescer.push(batch(2))
coalescer.flush()
coalescer.push({ ...batch(3), latestTurn: running })
coalescer.push({ ...batch(4), latestTurn: null })
coalescer.flush()
expect(events.map((event) => (event.type === 'batch' ? event.latestTurn : 'other'))).toEqual([
running,
null
])
})
it('drops it when a later frame carries rows without it, as applying both would', () => {
const events: AgentSessionSubscribeEvent[] = []
const coalescer = createStructuredAgentSessionEventCoalescer((event) => events.push(event))
const token = (sequence: number) => ({
...batch(sequence),
batch: {
...batch(sequence).batch,
items: [
{
itemId: `token-${sequence}`,
revision: 1,
sequence,
observedAt: sequence,
body: { kind: 'message' as const, role: 'assistant' as const, blocks: [] }
}
]
}
})
const running = {
itemId: 'turn-1',
observedAt: 1,
turn: { turnId: 'turn-1', state: 'running' as const }
}
// The second frame is an older host's: rows, and no answer to keep the first one alive.
coalescer.push({ ...token(1), latestTurn: running })
coalescer.push(token(2))
coalescer.flush()
expect(events).toHaveLength(1)
expect(events[0]).not.toHaveProperty('latestTurn')
})
})
@@ -1,4 +1,5 @@
import type { AgentSessionSubscribeEvent } from './agent-session-wire'
import { latestTurnAfterStructuredAgentSessionBatch } from './structured-agent-session-live-turn'
export const STRUCTURED_AGENT_SESSION_CLIENT_COALESCE_MS = 48
@@ -23,6 +24,8 @@ function mergeBatch(
for (const submission of right.batch.submissions) {
submissions.set(submission.clientMessageId, submission)
}
// As applying both in turn would leave it, so an older host's rows still drop a stale claim.
const latestTurn = latestTurnAfterStructuredAgentSessionBatch(left.latestTurn, right)
return {
type: 'batch',
...(right.commands !== undefined || left.commands !== undefined
@@ -55,7 +58,8 @@ function mergeBatch(
: {}),
...(right.activity !== undefined || left.activity !== undefined
? { activity: right.activity !== undefined ? right.activity : (left.activity ?? null) }
: {})
: {}),
...(latestTurn !== undefined ? { latestTurn } : {})
}
}
@@ -1,6 +1,7 @@
// How much of a session's journal a client keeps in memory.
import type { AgentJournalRenderItem } from './agent-session-journal-types'
import type { AgentJournalRenderItem, AgentJournalSubmission } from './agent-session-journal-types'
import { agentJournalSubmissionKey } from './agent-session-journal-item-key'
import { isRootAgentJournalItem } from './agent-session-journal-producer'
// Well above the renderer's initial read window (300) plus a page, so only genuinely
@@ -37,3 +38,28 @@ export function trimRetainedItems(
}
return start === 0 ? items : items.slice(start)
}
const MAX_RETAINED_SUBMISSIONS = 256
export function mergeSubmissions(
current: readonly AgentJournalSubmission[],
incoming: readonly AgentJournalSubmission[],
items: readonly AgentJournalRenderItem[]
): AgentJournalSubmission[] {
const byId = new Map(current.map((submission) => [submission.clientMessageId, submission]))
for (const submission of incoming) {
byId.set(submission.clientMessageId, submission)
}
const sorted = [...byId.values()].sort((left, right) => left.submittedAt - right.submittedAt)
const itemIds = new Set(
items
.filter((item) => item.body.kind === 'message' && item.body.role === 'user')
.map((item) => item.itemId)
)
// Loaded user messages need their provider alias for durable turn attribution.
return sorted.filter(
(submission, index) =>
index >= sorted.length - MAX_RETAINED_SUBMISSIONS ||
itemIds.has(agentJournalSubmissionKey(submission.clientMessageId))
)
}
@@ -27,7 +27,7 @@ describe('isStructuredAgentSessionThinking', () => {
})
it('is true while reasoning is the newest thing the turn produced', () => {
expect(isStructuredAgentSessionThinking([turnStart, reasoning(2)])).toBe(true)
expect(isStructuredAgentSessionThinking({ items: [turnStart, reasoning(2)] })).toBe(true)
})
it("reads a reasoning row's own state when its host keeps one", () => {
@@ -46,7 +46,7 @@ describe('isStructuredAgentSessionThinking', () => {
it('is false once a tool call, a message or a diff lands after the reasoning', () => {
const after = (body: AgentJournalRenderItem['body']): boolean =>
isStructuredAgentSessionThinking([turnStart, reasoning(2), item('after', 3, body)])
isStructuredAgentSessionThinking({ items: [turnStart, reasoning(2), item('after', 3, body)] })
expect(after({ kind: 'tool-call', name: 'shell', input: null, state: 'running' })).toBe(false)
expect(
after({ kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'Here you go' }] })
@@ -61,8 +61,8 @@ describe('isStructuredAgentSessionThinking', () => {
})
it('is false when a turn produced no reasoning at all', () => {
expect(isStructuredAgentSessionThinking([turnStart])).toBe(false)
expect(isStructuredAgentSessionThinking([])).toBe(false)
expect(isStructuredAgentSessionThinking({ items: [turnStart] })).toBe(false)
expect(isStructuredAgentSessionThinking({ items: [] })).toBe(false)
})
it('does not read an earlier turn as this one reasoning', () => {
@@ -73,7 +73,7 @@ describe('isStructuredAgentSessionThinking', () => {
text: 'Working',
turnLifecycle: { turnId: 'turn-2', state: 'running' }
})
expect(isStructuredAgentSessionThinking([reasoning(1), newTurn])).toBe(false)
expect(isStructuredAgentSessionThinking({ items: [reasoning(1), newTurn] })).toBe(false)
})
it('does not read a completed turn as reasoning during the next pending dispatch', () => {
@@ -82,20 +82,47 @@ describe('isStructuredAgentSessionThinking', () => {
turnId: 'turn-1',
state: 'completed'
})
expect(isStructuredAgentSessionThinking([completedTurn, reasoning(2)])).toBe(false)
expect(isStructuredAgentSessionThinking({ items: [completedTurn, reasoning(2)] })).toBe(false)
})
it('stops at a typed turn item, the carrier this host writes', () => {
const typedTurn = (sequence: number, turnId: string): AgentJournalRenderItem =>
item(`turn-${turnId}`, sequence, { kind: 'turn', turnId, state: 'running' })
expect(isStructuredAgentSessionThinking([typedTurn(1, 'turn-1'), reasoning(2)])).toBe(true)
expect(isStructuredAgentSessionThinking([reasoning(1), typedTurn(2, 'turn-2')])).toBe(false)
expect(
isStructuredAgentSessionThinking({ items: [typedTurn(1, 'turn-1'), reasoning(2)] })
).toBe(true)
expect(
isStructuredAgentSessionThinking({ items: [reasoning(1), typedTurn(2, 'turn-2')] })
).toBe(false)
})
it("reads the host's running turn when the turn's record is above the loaded rows", () => {
const hostTurn = (state: 'running' | 'completed') => ({
itemId: 'turn-start',
observedAt: 1,
turn: { turnId: 'turn-1', state }
})
expect(
isStructuredAgentSessionThinking({ items: [reasoning(2)], latestTurn: hostTurn('running') })
).toBe(true)
expect(
isStructuredAgentSessionThinking({ items: [reasoning(2)], latestTurn: hostTurn('completed') })
).toBe(false)
// An older host states nothing, and the loaded rows alone name no turn.
expect(isStructuredAgentSessionThinking({ items: [reasoning(2)] })).toBe(false)
// The host's answer outranks a loaded record whose ending revision is off the window.
expect(
isStructuredAgentSessionThinking({
items: [turnStart, reasoning(2)],
latestTurn: hostTurn('completed')
})
).toBe(false)
})
it('lets an unmarked status stay transparent to the latest reasoning state', () => {
const plan = item('plan', 3, { kind: 'status', text: 'Step 1. Read the file' })
expect(isStructuredAgentSessionThinking([turnStart, plan])).toBe(false)
expect(isStructuredAgentSessionThinking([turnStart, reasoning(2), plan])).toBe(true)
expect(isStructuredAgentSessionThinking({ items: [turnStart, plan] })).toBe(false)
expect(isStructuredAgentSessionThinking({ items: [turnStart, reasoning(2), plan] })).toBe(true)
})
it.each([
@@ -124,7 +151,9 @@ describe('isStructuredAgentSessionThinking', () => {
}
])('stops thinking when the turn is waiting on a $kind', (prompt) => {
expect(
isStructuredAgentSessionThinking([turnStart, reasoning(2), item('prompt', 3, prompt)])
isStructuredAgentSessionThinking({
items: [turnStart, reasoning(2), item('prompt', 3, prompt)]
})
).toBe(false)
})
})
@@ -154,7 +183,9 @@ describe("the live-turn readers answer for the session's own agent", () => {
role: 'reasoning',
blocks: [{ type: 'text', text: 'Weighing two approaches' }]
})
expect(isStructuredAgentSessionThinking([turnStart, spawnCall, childReasoning])).toBe(false)
expect(
isStructuredAgentSessionThinking({ items: [turnStart, spawnCall, childReasoning] })
).toBe(false)
})
it('still reports the parent as thinking when the parent itself is reasoning', () => {
@@ -163,7 +194,9 @@ describe("the live-turn readers answer for the session's own agent", () => {
role: 'reasoning',
blocks: [{ type: 'text', text: 'Weighing two approaches' }]
})
expect(isStructuredAgentSessionThinking([turnStart, spawnCall, ownReasoning])).toBe(true)
expect(isStructuredAgentSessionThinking({ items: [turnStart, spawnCall, ownReasoning] })).toBe(
true
)
})
it("reports the parent's own running call while a subagent runs its own", () => {
@@ -24,6 +24,7 @@ import {
type AgentJournalTurnScope
} from './agent-session-journal-types'
import { isRootAgentJournalItem } from './agent-session-journal-producer'
import type { AgentSessionLatestTurn, AgentSessionSubscribeEvent } from './agent-session-wire'
import { readAgentJournalTurn } from './agent-session-turn-record'
import type { NativeChatToolCallBlock } from './native-chat-types'
import {
@@ -101,15 +102,61 @@ export function activeStructuredAgentSessionTurnIdBySequence(
export function newestStructuredAgentSessionTurn(
items: readonly AgentJournalRenderItem[]
): AgentJournalTurnLifecycle | null {
return latestStructuredAgentSessionTurn(items)?.turn ?? null
}
/** The same record as a page publishes it, with the identity a client keys the turn by. */
export function latestStructuredAgentSessionTurn(
items: readonly AgentJournalRenderItem[]
): AgentSessionLatestTurn | null {
for (let index = items.length - 1; index >= 0; index -= 1) {
const turn = readAgentJournalTurn(items[index]?.body)
if (turn) {
return turn
const item = items[index]
const turn = readAgentJournalTurn(item?.body)
if (item && turn) {
return { itemId: item.itemId, observedAt: item.observedAt, turn }
}
}
return null
}
/** The turn the session's own agent is running, for a client: the host's answer over the whole
* journal, which no loaded window can hide. Temporary: an older host sends none, so its clients
* still read the loaded rows' newest record; delete that arm once such hosts age out. */
export function runningStructuredAgentSessionTurnId(state: HostTurnSource): string | null {
if (state.latestTurn === undefined) {
return activeStructuredAgentSessionTurnId(state.items)
}
return state.latestTurn?.turn.state === 'running' ? state.latestTurn.turn.turnId : null
}
/** The scope a row of that running turn names, read the same way. */
export function runningStructuredAgentSessionTurnScope(
state: HostTurnSource
): AgentJournalTurnScope {
if (state.latestTurn === undefined) {
return liveStructuredAgentSessionTurnScope(state.items)
}
return state.latestTurn?.turn.state === 'running'
? { kind: 'turn', turnItemId: state.latestTurn.itemId }
: AGENT_JOURNAL_THREAD_SCOPE
}
type HostTurnSource = {
items: readonly AgentJournalRenderItem[]
latestTurn?: AgentSessionLatestTurn | null
}
/** The host's answer once `event` applies. A batch carrying rows restates it, so one without it
* came from an older host and falls back to the rows rather than keep a claim nothing renews. */
export function latestTurnAfterStructuredAgentSessionBatch(
previous: AgentSessionLatestTurn | null | undefined,
event: Extract<AgentSessionSubscribeEvent, { type: 'batch' }>
): AgentSessionLatestTurn | null | undefined {
const { items, removedItemIds, submissions } = event.batch
const carriesRows = items.length > 0 || removedItemIds.length > 0 || submissions.length > 0
return carriesRows || event.latestTurn !== undefined ? event.latestTurn : previous
}
/**
* Whether the newest thing the active turn produced is the model's own reasoning.
*
@@ -118,16 +165,16 @@ export function newestStructuredAgentSessionTurn(
* thinking while the request is merely in flight, and stops reporting it the moment a tool call
* lands, which is usually when reasoning actually starts.
*/
export function isStructuredAgentSessionThinking(
items: readonly AgentJournalRenderItem[]
): boolean {
export function isStructuredAgentSessionThinking({ items, latestTurn }: HostTurnSource): boolean {
// The host's answer outranks a loaded record, whose newest revision may be off the window.
const hostRunning = latestTurn === undefined ? null : latestTurn?.turn.state === 'running'
let newestContentIsReasoning: boolean | null = null
for (let index = items.length - 1; index >= 0; index -= 1) {
const item = items[index]
const body = item?.body
const turn = readAgentJournalTurn(body)
if (turn) {
return turn.state === 'running' && newestContentIsReasoning === true
return (hostRunning ?? turn.state === 'running') && newestContentIsReasoning === true
}
if (newestContentIsReasoning !== null || !isRootAgentJournalItem(item)) {
continue
@@ -146,7 +193,8 @@ export function isStructuredAgentSessionThinking(
}
// A status row is a notice, not newer transcript content.
}
return false
// The record is above the loaded rows, so every loaded root row is newer than it.
return hostRunning === true && newestContentIsReasoning === true
}
/** The tool the status row names for the SESSION'S OWN agent, as the chat draws it: the running
@@ -1,6 +1,7 @@
import { describe, expect, it } from 'vitest'
import type { AgentJournalRenderItem, AgentJournalSubmission } from './agent-session-journal-types'
import type { AgentSessionHistoryPage } from './agent-session-wire'
import type { AgentSessionHistoryPage, AgentSessionLatestTurn } from './agent-session-wire'
import { runningStructuredAgentSessionTurnId } from './structured-agent-session-live-turn'
import {
EMPTY_STRUCTURED_AGENT_SESSION,
reduceStructuredAgentSession
@@ -534,6 +535,66 @@ describe('structured agent session reducer', () => {
expect(cleared.items).toBe(active.items)
})
describe("the host's newest turn record", () => {
const turnRow = (state: 'running' | 'completed'): AgentJournalRenderItem => ({
itemId: 'turn-record',
revision: state === 'running' ? 1 : 2,
sequence: 1,
observedAt: 1,
body: { kind: 'turn', turnId: 'turn-1', state }
})
const hostTurn = (state: 'running' | 'completed'): AgentSessionLatestTurn => ({
itemId: 'turn-record',
observedAt: 1,
turn: { turnId: 'turn-1', state }
})
const opened = () =>
reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
type: 'event',
event: {
type: 'snapshot',
sessionId: 'session-a',
fence: 1,
page: { ...hydrationPage([turnRow('running')]), latestTurn: hostTurn('running') }
}
})
const batch = (
state: ReturnType<typeof opened>,
items: AgentJournalRenderItem[],
latestTurn?: AgentSessionLatestTurn | null
) =>
reduceStructuredAgentSession(state, {
type: 'event',
event: {
type: 'batch',
sessionId: 'session-a',
batch: {
cursor: { epoch: 'epoch-a', sequence: state.cursor!.sequence + items.length },
items,
removedItemIds: [],
submissions: []
},
activity: { turnId: 'turn-1', text: 'Reading' },
...(latestTurn !== undefined ? { latestTurn } : {})
}
})
it('keeps it across a frame that carries no rows, and takes the next one rows carry', () => {
const quiet = batch(opened(), [])
expect(quiet.latestTurn).toEqual(hostTurn('running'))
const ended = batch(quiet, [turnRow('completed')], hostTurn('completed'))
expect(ended.latestTurn).toEqual(hostTurn('completed'))
expect(runningStructuredAgentSessionTurnId(ended)).toBeNull()
})
it('falls back to the loaded rows when rows arrive without it, as from an older host', () => {
// A cursor resume against a downgraded host: its rows end the turn but restate nothing.
const downgraded = batch(opened(), [turnRow('completed')])
expect(downgraded.latestTurn).toBeUndefined()
expect(runningStructuredAgentSessionTurnId(downgraded)).toBeNull()
})
})
it('records the host clock from frames that carry it and keeps it otherwise', () => {
const snapshot = reduceStructuredAgentSession(
EMPTY_STRUCTURED_AGENT_SESSION,
+8 -26
View File
@@ -7,6 +7,7 @@ import type {
AgentSessionBackgroundTaskState,
AgentSessionSlashCommand,
AgentSessionHistoryPage,
AgentSessionLatestTurn,
AgentSessionQueuedMessage,
AgentSessionQueuePause,
AgentSessionSubscribeEvent,
@@ -15,15 +16,16 @@ import type {
import type { AgentSessionRefusalReference } from './agent-session-wire-refusals'
import { backgroundTaskStatesEqual } from './agent-session-background-task-state-equality'
import { admitAgentSessionBackgroundTaskState } from './agent-session-background-task-state-admission'
import { agentJournalSubmissionKey } from './agent-session-journal-item-key'
import {
MAX_RETAINED_ITEMS,
MAX_RETAINED_OWN_ITEMS,
mergeSubmissions,
ownItemCount,
trimRetainedItems
} from './structured-agent-session-item-retention'
import { compareAgentJournalItems } from './agent-session-journal-position'
import { readAgentJournalTurn } from './agent-session-turn-record'
import { latestTurnAfterStructuredAgentSessionBatch } from './structured-agent-session-live-turn'
import {
foldStructuredAgentSubagentRoster,
foldStructuredAgentSubagentRosterPage,
@@ -67,6 +69,9 @@ export type StructuredAgentSessionState = {
/** Every subagent a roster row this client received named, by agent id; not trimmed with
* `items`. Absent until a page has been applied. */
subagentRoster?: StructuredAgentSubagentRoster
/** The host's newest turn record over the whole journal, which says whether a turn runs; absent
* from an older host, whose answer is read off `items` instead. */
latestTurn?: AgentSessionLatestTurn | null
/** Bumped per live batch that leaves a turn row's newest revision outside the window
* (dropped or trimmed), so a whole-journal answer derived from turn rows is asked for again. */
unloadedTurnRevisions?: number
@@ -81,8 +86,6 @@ export type StructuredAgentSessionAction =
| { type: 'history-page'; page: AgentSessionHistoryPage }
| { type: 'older-page'; requestedCursor: AgentJournalCursor; page: AgentSessionHistoryPage }
const MAX_RETAINED_SUBMISSIONS = 256
export const EMPTY_STRUCTURED_AGENT_SESSION: StructuredAgentSessionState = {
epoch: null,
cursor: null,
@@ -136,6 +139,7 @@ function replacePage(
status: 'ready',
subagentRoster: foldStructuredAgentSubagentRosterPage(undefined, page),
activity: activity ?? null,
...(page.latestTurn !== undefined ? { latestTurn: page.latestTurn } : {}),
...(backgroundTasks !== undefined
? { backgroundTasks }
: page.backgroundTasks !== undefined
@@ -182,29 +186,6 @@ function liveItemsWithinWindow(
return incoming.filter((item) => item.sequence >= head.sequence)
}
function mergeSubmissions(
current: readonly AgentJournalSubmission[],
incoming: readonly AgentJournalSubmission[],
items: readonly AgentJournalRenderItem[]
): AgentJournalSubmission[] {
const byId = new Map(current.map((submission) => [submission.clientMessageId, submission]))
for (const submission of incoming) {
byId.set(submission.clientMessageId, submission)
}
const sorted = [...byId.values()].sort((left, right) => left.submittedAt - right.submittedAt)
const itemIds = new Set(
items
.filter((item) => item.body.kind === 'message' && item.body.role === 'user')
.map((item) => item.itemId)
)
// Loaded user messages need their provider alias for durable turn attribution.
return sorted.filter(
(submission, index) =>
index >= sorted.length - MAX_RETAINED_SUBMISSIONS ||
itemIds.has(agentJournalSubmissionKey(submission.clientMessageId))
)
}
/** `receivedAt` is the client clock at apply time; callers pass it so the reducer stays pure. */
export function reduceStructuredAgentSession(
state: StructuredAgentSessionState,
@@ -336,6 +317,7 @@ export function reduceStructuredAgentSession(
error: undefined,
readRefusal: undefined,
commands: event.commands !== undefined ? event.commands : state.commands,
latestTurn: latestTurnAfterStructuredAgentSessionBatch(state.latestTurn, event),
...queuePublicationField(event, state),
backgroundTasks,
...(activity !== undefined ? { activity } : {}),
@@ -13,6 +13,7 @@ import { readAgentJournalTurn, readAgentJournalTurnOutcome } from './agent-sessi
import { agentTurnVerdict, type AgentTurnOutcome } from './agent-turn-outcome'
import { structuredAgentTurnAnchors } from './native-chat-turn-membership'
import type { NativeChatSettledTurn, NativeChatSettledTurns } from './native-chat-turn-status'
import type { AgentSessionLatestTurn } from './agent-session-wire'
export type StructuredAgentTurnTiming = {
state: AgentJournalTurnLifecycleState
@@ -36,7 +37,7 @@ export type StructuredAgentTurnTiming = {
}
function readTiming(
item: AgentJournalRenderItem,
item: Pick<AgentJournalRenderItem, 'body' | 'observedAt'>,
precedingTurnEndedAt: number | undefined
): StructuredAgentTurnTiming | null {
const turn = readAgentJournalTurn(item.body)
@@ -220,14 +221,27 @@ function settledTurnsOf(
export function selectStructuredAgentTurnBars(
items: readonly AgentJournalRenderItem[],
submissions: readonly AgentJournalSubmission[],
turnId: string | null
turnId: string | null,
/** The host's newest turn record, which times a running turn whose row is not loaded. */
latestTurn?: AgentSessionLatestTurn | null
): {
settledTurns: NativeChatSettledTurns
runningTiming: StructuredAgentTurnTiming | null
} {
const turns = readStructuredAgentJournalTurns(items, submissions)
const unloaded =
turnId !== null && !turns.byTurnId.has(turnId) && latestTurn?.turn.turnId === turnId
? latestTurn
: null
return {
settledTurns: settledTurnsOf(turns, submissions),
runningTiming: turnId === null ? null : (turns.byTurnId.get(turnId) ?? null)
runningTiming: unloaded
? readTiming(
{ body: { kind: 'turn', ...unloaded.turn }, observedAt: unloaded.observedAt },
undefined
)
: turnId === null
? null
: (turns.byTurnId.get(turnId) ?? null)
}
}