diff --git a/packages/ui/src/__tests__/chat-turn-steering-order.test.ts b/packages/ui/src/__tests__/chat-turn-steering-order.test.ts index f4a2f6e25c..756032d42a 100644 --- a/packages/ui/src/__tests__/chat-turn-steering-order.test.ts +++ b/packages/ui/src/__tests__/chat-turn-steering-order.test.ts @@ -24,13 +24,14 @@ import { renderToStaticMarkup } from 'react-dom/server'; import { parseHTML } from 'linkedom'; import { TurnView } from '../chat-turn.js'; import { LocaleProvider } from '../locale-context.js'; -import { materializeTurns, type TurnViewModel } from '../materialize.js'; +import { foldTimeline } from '../timeline-fold.js'; +import { materializeTurns, overlayLiveTurn, type TurnTimelineItem, type TurnViewModel } from '../materialize.js'; import { createTranscriptProjection } from '../transcript-projection.js'; import { ChatView } from '../chat-view.js'; import { Composer } from '../composer.js'; import { renderTranscriptMarkup } from './transcript-test-dom.js'; import { ChatSurfaceLayout } from '../chat-surface-layout.js'; -import { armLiveTurn } from '../live-turn-projection.js'; +import { armLiveTurn, type LiveTurnProjection } from '../live-turn-projection.js'; import { applyLiveTurnEvent } from './live-turn-zh.js'; import type { SessionSummary, StoredMessage } from '@maka/core/session'; @@ -125,3 +126,220 @@ test('holds steering above the composer while old output continues, then renders assert.deepEqual(timeline(), ['old answer continues', pending.text, 'reply to new instruction']); assert.equal(text.split(pending.text).length - 1, 1); }); + +const timelineOrder = (timeline: readonly TurnTimelineItem[]): string[] => + timeline.map((item) => + item.kind === 'user' + ? `user:${item.message.text}` + : item.kind === 'tools' + ? `tools:${item.items.map((tool) => tool.toolUseId).join('+')}` + : `${item.kind}:${item.text}`, + ); + +test('keeps post-steering work below the steering row even when it continues the same step', () => { + let live: LiveTurnProjection | undefined = applyLiveTurnEvent(armLiveTurn('turn-1'), { + type: 'thinking_delta', id: 'e1', turnId: 'turn-1', messageId: 'm1', ts: 1, text: 'pre-steer reasoning', + }); + live = applyLiveTurnEvent(live, { + type: 'tool_start', id: 'e2', turnId: 'turn-1', stepId: 'm1', toolUseId: 'tool-1', toolName: 'Read', args: {}, ts: 2, + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 'steer-event', turnId: 'turn-1', messageId: 'steer-1', ts: 3, content: { text: 'steer' }, + }); + // A late result for a row that already exists updates that row in place; it + // does not claim the steering because its position was fixed at tool_start. + live = applyLiveTurnEvent(live, { + type: 'tool_result', id: 'e3', turnId: 'turn-1', toolUseId: 'tool-1', isError: false, ts: 4, + content: { kind: 'text', text: 'done' }, + }); + // While the steering awaits its boundary it must already render as one. + assert.deepEqual(timelineOrder(overlayLiveTurn([], live, 'en')[0]!.timeline), [ + 'thinking:pre-steer reasoning', 'tools:tool-1', 'user:steer', + ]); + live = applyLiveTurnEvent(live, { + type: 'thinking_delta', id: 'e4', turnId: 'turn-1', messageId: 'm1', ts: 5, text: 'post-steer reasoning', + }); + live = applyLiveTurnEvent(live, { + type: 'tool_start', id: 'e5', turnId: 'turn-1', stepId: 'm1', toolUseId: 'tool-2', toolName: 'Bash', args: {}, ts: 6, + }); + live = applyLiveTurnEvent(live, { + type: 'text_delta', id: 'e6', turnId: 'turn-1', messageId: 'm1', ts: 7, text: 'answer continues', + }); + + const timeline = overlayLiveTurn([], live, 'en')[0]!.timeline; + assert.deepEqual(timelineOrder(timeline), [ + 'thinking:pre-steer reasoning', + 'tools:tool-1', + 'user:steer', + 'thinking:post-steer reasoning', + 'tools:tool-2', + 'text:answer continues', + ]); + const folded = foldTimeline(timeline).entries.map((entry) => + entry.kind === 'processing' + ? `fold:${entry.children.map((child) => child.kind).join('+')}` + : entry.kind === 'user' + ? `user:${entry.message.text}` + : `text:${entry.text}`, + ); + assert.deepEqual(folded, [ + 'fold:thinking+tools', + 'user:steer', + 'fold:thinking+tools', + 'text:answer continues', + ]); +}); + +test('anchors a persisted steering row ahead of live work the stream seeded after it', () => { + const settled = materializeTurns([ + { type: 'user', id: 'original', turnId: 't1', ts: 1, text: 'request' }, + { type: 'user', id: 'steer-1', turnId: 't1', ts: 3, text: 'steer', steeringEventId: 'steer-event' }, + ], 'en'); + const live = applyLiveTurnEvent(armLiveTurn('t1'), { + type: 'thinking_delta', id: 'e1', turnId: 't1', messageId: 'm2', ts: 5, text: 'in-flight after the steer', + }); + + const [overlaid] = overlayLiveTurn(settled, live, 'en'); + assert.deepEqual(timelineOrder(overlaid!.timeline), [ + 'user:steer', + 'thinking:in-flight after the steer', + ]); +}); + +test('keeps pre-steering content in place when its completion lands after the boundary', () => { + let live: LiveTurnProjection | undefined = applyLiveTurnEvent(armLiveTurn('turn-1'), { + type: 'text_delta', id: 'e1', turnId: 'turn-1', messageId: 'm1', ts: 1, text: 'answer before', + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 'steer-event', turnId: 'turn-1', messageId: 'steer-1', ts: 2, content: { text: 'steer' }, + }); + live = applyLiveTurnEvent(live, { + type: 'text_complete', id: 'e2', turnId: 'turn-1', messageId: 'm1', ts: 3, text: 'answer before', + }); + + assert.deepEqual(timelineOrder(overlayLiveTurn([], live, 'en')[0]!.timeline), [ + 'text:answer before', + 'user:steer', + ]); +}); + +test('splits a completion across the boundary instead of duplicating the sealed portion', () => { + let live: LiveTurnProjection | undefined = applyLiveTurnEvent(armLiveTurn('turn-1'), { + type: 'thinking_delta', id: 'e1', turnId: 'turn-1', messageId: 'm1', ts: 1, text: 'pre-', + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 'steer-event', turnId: 'turn-1', messageId: 'steer-1', ts: 2, content: { text: 'steer' }, + }); + live = applyLiveTurnEvent(live, { + type: 'thinking_delta', id: 'e2', turnId: 'turn-1', messageId: 'm1', ts: 3, text: 'post', + }); + live = applyLiveTurnEvent(live, { + type: 'thinking_complete', id: 'e3', turnId: 'turn-1', messageId: 'm1', ts: 4, text: 'pre-post', + }); + + assert.deepEqual(timelineOrder(overlayLiveTurn([], live, 'en')[0]!.timeline), [ + 'thinking:pre-', + 'user:steer', + 'thinking:post', + ]); +}); + +test('lands a divergent completion whole instead of slicing at the delta offset', () => { + // thinking_complete may carry a provider summary that replaces the streamed + // deltas outright — the accumulated delta length is not a safe cut point. + let live: LiveTurnProjection | undefined = applyLiveTurnEvent(armLiveTurn('turn-1'), { + type: 'thinking_delta', id: 'e1', turnId: 'turn-1', messageId: 'm1', ts: 1, text: 'AAAA', + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 's1', turnId: 'turn-1', messageId: 'steer-1', ts: 2, content: { text: 'steer' }, + }); + live = applyLiveTurnEvent(live, { + type: 'thinking_delta', id: 'e2', turnId: 'turn-1', messageId: 'm1', ts: 3, text: 'BBBB', + }); + live = applyLiveTurnEvent(live, { + type: 'thinking_complete', id: 'e3', turnId: 'turn-1', messageId: 'm1', ts: 4, text: 'Short summary.', + }); + + assert.deepEqual(timelineOrder(overlayLiveTurn([], live, 'en')[0]!.timeline), [ + 'thinking:AAAA', + 'user:steer', + 'thinking:Short summary.', + ]); +}); + +test('lands a completion shorter than the delta offset whole', () => { + let live: LiveTurnProjection | undefined = applyLiveTurnEvent(armLiveTurn('turn-1'), { + type: 'thinking_delta', id: 'e1', turnId: 'turn-1', messageId: 'm1', ts: 1, text: 'AAAA', + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 's1', turnId: 'turn-1', messageId: 'steer-1', ts: 2, content: { text: 'steer' }, + }); + live = applyLiveTurnEvent(live, { + type: 'thinking_delta', id: 'e2', turnId: 'turn-1', messageId: 'm1', ts: 3, text: 'BBBB', + }); + live = applyLiveTurnEvent(live, { + type: 'thinking_complete', id: 'e3', turnId: 'turn-1', messageId: 'm1', ts: 4, text: 'ABC', + }); + + assert.deepEqual(timelineOrder(overlayLiveTurn([], live, 'en')[0]!.timeline), [ + 'thinking:AAAA', + 'user:steer', + 'thinking:ABC', + ]); +}); + +test('keeps consecutive steering rows in arrival order', () => { + let live: LiveTurnProjection | undefined = applyLiveTurnEvent(armLiveTurn('turn-1'), { + type: 'text_delta', id: 'e1', turnId: 'turn-1', messageId: 'm1', ts: 1, text: 'before', + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 's1', turnId: 'turn-1', messageId: 'steer-1', ts: 2, content: { text: 'one' }, + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 's2', turnId: 'turn-1', messageId: 'steer-2', ts: 3, content: { text: 'two' }, + }); + live = applyLiveTurnEvent(live, { + type: 'text_delta', id: 'e2', turnId: 'turn-1', messageId: 'm2', ts: 4, text: 'after', + }); + + assert.deepEqual(timelineOrder(overlayLiveTurn([], live, 'en')[0]!.timeline), [ + 'text:before', + 'user:one', + 'user:two', + 'text:after', + ]); +}); + +test('renders a repeated steering event once', () => { + let live: LiveTurnProjection | undefined = applyLiveTurnEvent(armLiveTurn('turn-1'), { + type: 'steering_message', id: 's1', turnId: 'turn-1', messageId: 'steer-1', ts: 1, content: { text: 'steer' }, + }); + live = applyLiveTurnEvent(live, { + type: 'steering_message', id: 's1-echo', turnId: 'turn-1', messageId: 'steer-1', ts: 2, content: { text: 'steer' }, + }); + + assert.deepEqual(timelineOrder(overlayLiveTurn([], live, 'en')[0]!.timeline), ['user:steer']); +}); + +test('trails a settled steering that carries no ts behind live work', () => { + const turn: TurnViewModel = { + turnId: 't1', + status: 'running', + tools: [], + notes: [], + startedAt: 1, + timeline: [ + { kind: 'text', text: 'persisted answer', messageId: 'm1', ts: 1 }, + { kind: 'user', message: { id: 'steer-1', role: 'user', text: 'steer' }, messageId: 'steer-1', steeringEventId: 'steer-event' }, + ], + }; + const live = applyLiveTurnEvent(armLiveTurn('t1'), { + type: 'text_delta', id: 'e1', turnId: 't1', messageId: 'm2', ts: 5, text: 'in-flight', + }); + + assert.deepEqual(timelineOrder(overlayLiveTurn([turn], live, 'en')[0]!.timeline), [ + 'text:persisted answer', + 'text:in-flight', + 'user:steer', + ]); +}); diff --git a/packages/ui/src/__tests__/live-turn-projection.test.ts b/packages/ui/src/__tests__/live-turn-projection.test.ts index 250f6adb22..1bb3511b87 100644 --- a/packages/ui/src/__tests__/live-turn-projection.test.ts +++ b/packages/ui/src/__tests__/live-turn-projection.test.ts @@ -282,6 +282,7 @@ describe('applyLiveTurnEvent', () => { text: '完整思考', truncated: false, complete: true, + sourceEndOffset: 4, }); }); @@ -710,13 +711,14 @@ describe('reconcileTerminalLiveTurn', () => { it('keeps terminal live steering until the terminal transcript catches up', () => { const message = { id: 'steer-1', content: { text: 'change direction' }, ts: 2 }; + const boundary = { stepId: 'steering:steer-1', tools: [], steering: message }; const withSteering: LiveTurnProjection = { ...toolOnly, - steps: [{ ...toolOnly.steps[0]!, leadingSteering: [message] }], + steps: [...toolOnly.steps, boundary], }; assert.equal(reconcileTerminalLiveTurn(withSteering, []), withSteering); - const steeringOnly = { ...withSteering, steps: [] }; + const steeringOnly = { ...toolOnly, steps: [boundary] }; assert.equal(reconcileTerminalLiveTurn(steeringOnly, []), steeringOnly); assert.deepEqual(reconcileTerminalLiveTurn(withSteering, [{ type: 'turn_state', id: 'state-1', turnId: 'turn-1', ts: 3, @@ -735,7 +737,7 @@ describe('reconcileTerminalLiveTurn', () => { }); assert.equal(aborted?.terminal, true); - assert.deepEqual(aborted?.pendingSteering, [message]); + assert.deepEqual(aborted?.steps[0]?.steering, message); }); it('retains interrupted live output until a persisted result covers it', () => { diff --git a/packages/ui/src/live-turn-projection.ts b/packages/ui/src/live-turn-projection.ts index 9f0130dd53..0c6bf2a188 100644 --- a/packages/ui/src/live-turn-projection.ts +++ b/packages/ui/src/live-turn-projection.ts @@ -63,9 +63,22 @@ export interface LiveThinkingProjection { export interface LiveTurnStepProjection { stepId: string; + /** Event ts of the first word that opened this step or slice. */ + startedAt?: number; contentOrder?: LiveTurnStepContentKind[]; - /** Steering drained immediately before this provider step began. */ - leadingSteering?: LiveSteeringProjection[]; + /** + * A steering boundary slice: the steering row is emitted at this position + * and the slice holds no content. Content events never resolve to it + * (`steering:`-prefixed stepIds collide with no real stepId). + */ + steering?: LiveSteeringProjection; + /** + * A steering boundary can split a step mid-flight; the continuation slice + * keeps the same durable stepId. Replayed deltas trim against the earlier + * slice's source length through these baselines instead of re-appending it. + */ + continuedThinkingEndOffset?: number; + continuedTextEndOffset?: number; thinking?: LiveThinkingProjection; text?: LiveTextProjection; tools: ToolActivityItem[]; @@ -103,8 +116,6 @@ export interface LiveTurnProjection { /** Event ts of the first authority word about this Turn, so a Turn the * transcript has not reached yet still has a stable start. */ startedAt?: number; - /** Steering acknowledged after the current content and awaiting its next provider step. */ - pendingSteering?: LiveSteeringProjection[]; /** * Set by `armLiveTurn` and cleared by the first word the authority says about * this turn (any event carrying the same turnId). @@ -219,14 +230,21 @@ function projectLiveTurnEvent( if (liveSteeringMessages(prior).some((message) => message.id === event.messageId)) { return confirmed(prior); } + // A steering's position is fixed by stream order: it becomes a boundary + // slice immediately instead of parking until a later event claims it. return { ...confirmed(prior), - pendingSteering: [ - ...(prior.pendingSteering ?? []), + steps: [ + ...prior.steps, { - id: event.messageId, - content: structuredClone(event.content), - ts: event.ts, + stepId: `steering:${event.messageId}`, + startedAt: event.ts, + tools: [], + steering: { + id: event.messageId, + content: structuredClone(event.content), + ts: event.ts, + }, }, ], }; @@ -296,23 +314,41 @@ function projectLiveTurnEvent( : event.type === 'tool_start' ? event.stepId ?? existingToolStep?.stepId ?? `tool:${event.toolUseId}` : existingToolStep?.stepId ?? `tool:${event.toolUseId}`; - const stepIndex = prior.steps.findIndex((step) => step.stepId === stepId); - const isNewStep = stepIndex < 0; - const claimsPendingSteering = isNewStep - && existingToolStep === undefined - && (prior.pendingSteering?.length ?? 0) > 0; - const step: LiveTurnStepProjection = isNewStep + // A steering boundary freezes the positions before it: same-stepId deltas + // arriving after one continue the step in a fresh slice, so a stepId can + // repeat across the array. Completions and events for an existing tool row + // are updates to a row whose position is already fixed — they resolve to + // the row's last slice wherever it sits, on either side of a boundary. + const boundaryIndex = prior.steps.findLastIndex((candidate) => candidate.steering !== undefined); + const sameStepIndex = prior.steps.findLastIndex((candidate) => candidate.stepId === stepId); + const stepIndex = existingToolStep === undefined + ? event.type === 'thinking_complete' || event.type === 'text_complete' + ? sameStepIndex + : sameStepIndex > boundaryIndex ? sameStepIndex : -1 + : event.type !== 'tool_start' + || event.stepId === undefined + || event.stepId === existingToolStep.stepId + ? prior.steps.indexOf(existingToolStep) + : sameStepIndex; + const continuedFrom = stepIndex < 0 && sameStepIndex >= 0 ? prior.steps[sameStepIndex]! : undefined; + const continuedThinkingEnd = continuedFrom?.thinking?.sourceEndOffset ?? continuedFrom?.continuedThinkingEndOffset; + const continuedTextEnd = continuedFrom?.text?.sourceEndOffset ?? continuedFrom?.continuedTextEndOffset; + const step: LiveTurnStepProjection = stepIndex < 0 ? { stepId, + startedAt: event.ts, tools: [], - ...(claimsPendingSteering - ? { leadingSteering: prior.pendingSteering } - : {}), + ...(continuedThinkingEnd === undefined + ? {} + : { continuedThinkingEndOffset: continuedThinkingEnd }), + ...(continuedTextEnd === undefined + ? {} + : { continuedTextEndOffset: continuedTextEnd }), } : prior.steps[stepIndex]!; let nextStep: LiveTurnStepProjection; if (event.type === 'thinking_delta') { - const delta = replaySafeDelta(step.thinking?.sourceEndOffset, event); + const delta = replaySafeDelta(step.thinking?.sourceEndOffset ?? step.continuedThinkingEndOffset, event); const applied = applyThinkingDelta(step.thinking?.text ?? '', delta.text, { locale, ...(step.thinking?.redactionState === undefined @@ -334,20 +370,23 @@ function projectLiveTurnEvent( }, }; } else if (event.type === 'thinking_complete') { - const applied = applyThinkingComplete(event.text, { locale }); + const applied = applyThinkingComplete( + completionRemainder(prior, step, 'thinking', event.text), + { locale }, + ); nextStep = { ...step, thinking: { text: applied.text, truncated: applied.truncated, complete: true, - ...(step.thinking?.sourceEndOffset === undefined + ...((step.thinking?.sourceEndOffset ?? step.continuedThinkingEndOffset) === undefined ? {} : { sourceEndOffset: event.text.length }), }, }; } else if (event.type === 'text_delta') { - const delta = replaySafeDelta(step.text?.sourceEndOffset, event); + const delta = replaySafeDelta(step.text?.sourceEndOffset ?? step.continuedTextEndOffset, event); const applied = applyAssistantDelta(step.text?.text ?? '', delta.text, { locale, ...(step.text?.redactionState === undefined @@ -369,7 +408,10 @@ function projectLiveTurnEvent( }, }; } else if (event.type === 'text_complete') { - const applied = applyAssistantComplete(event.text, { locale }); + const applied = applyAssistantComplete( + completionRemainder(prior, step, 'text', event.text), + { locale }, + ); nextStep = { ...step, text: { @@ -377,7 +419,7 @@ function projectLiveTurnEvent( text: applied.text, truncated: applied.truncated, complete: true, - ...(step.text?.sourceEndOffset === undefined + ...((step.text?.sourceEndOffset ?? step.continuedTextEndOffset) === undefined ? {} : { sourceEndOffset: event.text.length }), }, @@ -496,7 +538,7 @@ function projectLiveTurnEvent( }; let steps: LiveTurnStepProjection[]; if (existingToolStep && existingToolStep.stepId !== stepId && !messageEvent) { - const sourceIndex = prior.steps.findIndex((candidate) => candidate.stepId === existingToolStep.stepId); + const sourceIndex = prior.steps.indexOf(existingToolStep); const sourceWithoutTool = { ...existingToolStep, tools: existingToolStep.tools.filter((tool) => tool.toolUseId !== event.toolUseId), @@ -507,7 +549,7 @@ function projectLiveTurnEvent( const sourceIsEmpty = !sourceWithoutTool.thinking && !sourceWithoutTool.text && sourceWithoutTool.tools.length === 0 - && (sourceWithoutTool.leadingSteering?.length ?? 0) === 0; + && sourceWithoutTool.steering === undefined; steps = []; for (let index = 0; index < prior.steps.length; index += 1) { const candidate = prior.steps[index]!; @@ -526,18 +568,47 @@ function projectLiveTurnEvent( ? prior.steps.map((candidate, index) => index === stepIndex ? nextStep : candidate) : [...prior.steps, nextStep]; } - const { pendingSteering: _pendingSteering, ...withoutPendingSteering } = priorWithoutRetry; - return { - ...(claimsPendingSteering ? withoutPendingSteering : priorWithoutRetry), - steps, - }; + // A completion finalizes the message across every slice it occupies: an + // earlier slice keeps the portion it rendered — marked complete — so a + // steering boundary never relocates pre-steering content into the full text + // a later slice finalizes. + const finalizedKind = event.type === 'thinking_complete' + ? 'thinking' + : event.type === 'text_complete' ? 'text' : undefined; + if (finalizedKind !== undefined) { + steps = steps.map((candidate) => + candidate !== nextStep + && candidate.stepId === stepId + && candidate[finalizedKind] !== undefined + ? { ...candidate, [finalizedKind]: { ...candidate[finalizedKind]!, complete: true } } + : candidate); + } + return { ...priorWithoutRetry, steps }; } function liveSteeringMessages(current: LiveTurnProjection): LiveSteeringProjection[] { - return [ - ...(current.pendingSteering ?? []), - ...current.steps.flatMap((step) => step.leadingSteering ?? []), - ]; + return current.steps.flatMap((step) => (step.steering ? [step.steering] : [])); +} + +/** + * A completion's full text shares the delta stream's coordinates only when it + * extends the already-rendered prefix — a provider summary replaces the + * streamed text outright (`reasoningSummaryText` adoption), so a bare offset + * would cut real content. Trim only on a verified prefix; otherwise land the + * payload whole. + */ +function completionRemainder( + prior: LiveTurnProjection, + step: LiveTurnStepProjection, + kind: 'thinking' | 'text', + fullText: string, +): string { + const rendered = prior.steps.flatMap((candidate) => + candidate !== step && candidate.stepId === step.stepId && candidate[kind] + ? [candidate[kind]!.text] + : []); + const prefix = rendered.join(''); + return fullText.startsWith(prefix) ? fullText.slice(prefix.length) : fullText; } function replaySafeDelta( @@ -547,9 +618,7 @@ function replaySafeDelta( if (event.startOffset === undefined) { return { text: event.text, - ...(currentEndOffset === undefined - ? {} - : { sourceEndOffset: currentEndOffset + event.text.length }), + sourceEndOffset: (currentEndOffset ?? 0) + event.text.length, }; } const endOffset = event.startOffset + event.text.length; @@ -567,28 +636,30 @@ function replaySafeDelta( * Streaming display handoff: drop the committed text/thinking slots for `stepId`. * Tools that still carry live stream evidence (outputChunks) stay — empty * shell_run durable results do not cover them, and co-located Bash+answer - * steps must not lose pre-handoff output when the answer settles. + * steps must not lose pre-handoff output when the answer settles. `stepId` + * can match multiple slices once a steering boundary split the step; + * steering boundary slices carry their own namespaced id and never match. */ export function settleLiveTurnStep( current: LiveTurnProjection, stepId: string, ): LiveTurnProjection | undefined { - const stepIndex = current.steps.findIndex((step) => step.stepId === stepId); - if (stepIndex < 0) return current; - const step = current.steps[stepIndex]!; - const retainedTools = step.tools.filter((tool) => (tool.outputChunks?.length ?? 0) > 0); - const steps = retainedTools.length > 0 - ? current.steps.map((candidate, index) => ( - index === stepIndex - ? { - stepId: candidate.stepId, - tools: retainedTools, - contentOrder: ['tools' as const], - } - : candidate - )) - : current.steps.filter((candidate) => candidate.stepId !== stepId); - if (steps.length === current.steps.length && retainedTools.length === 0) return current; + let found = false; + const steps = current.steps.flatMap((step) => { + if (step.stepId !== stepId) return [step]; + found = true; + const retainedTools = step.tools.filter((tool) => (tool.outputChunks?.length ?? 0) > 0); + if (retainedTools.length === 0) { + return []; + } + return [{ + stepId, + tools: retainedTools, + ...(retainedTools.length > 0 ? { contentOrder: ['tools' as const] } : {}), + ...(step.startedAt !== undefined ? { startedAt: step.startedAt } : {}), + }]; + }); + if (!found) return current; if (steps.length === 0 && current.terminal) return undefined; return { ...current, steps }; } @@ -655,6 +726,9 @@ export function reconcileTerminalLiveTurn( const toolCallIds = new Set(turnMessages.flatMap((message) => message.type === 'tool_call' ? [message.id] : [])); const toolResultIds = new Set(turnMessages.flatMap((message) => message.type === 'tool_result' ? [message.toolUseId] : [])); let steps = projection.steps.filter((step) => { + // Steering boundary slices hold no durable-comparable content; the + // overlay dedupes against the persisted user row by id. + if (step.steering !== undefined) return true; if (step.text?.text.length) return true; if (step.thinking && !assistantIds.has(step.stepId)) return true; const toolsCovered = step.tools.every((tool) => { @@ -676,11 +750,7 @@ export function reconcileTerminalLiveTurn( && transcriptReachedTerminal && liveSteeringMessages(projection).length > 0; if (steeringSettled) { - steps = steps.map((step) => { - if (!step.leadingSteering) return step; - const { leadingSteering: _leadingSteering, ...withoutSteering } = step; - return withoutSteering; - }); + steps = steps.filter((step) => step.steering === undefined); } if ( steps.length === 0 @@ -692,7 +762,5 @@ export function reconcileTerminalLiveTurn( ) ) return undefined; if (steps.length === projection.steps.length && !steeringSettled) return projection; - if (!steeringSettled) return { ...projection, steps }; - const { pendingSteering: _pendingSteering, ...withoutSteering } = projection; - return { ...withoutSteering, steps }; + return { ...projection, steps }; } diff --git a/packages/ui/src/materialize.ts b/packages/ui/src/materialize.ts index 3178a5f9da..24fce88a1d 100644 --- a/packages/ui/src/materialize.ts +++ b/packages/ui/src/materialize.ts @@ -504,22 +504,14 @@ export function overlayLiveTurn( } satisfies TurnViewModel, ]; } - if ( - targetIndex >= 0 - && liveTurn.steps.length === 0 - && (liveTurn.pendingSteering?.length ?? 0) === 0 - ) { + if (targetIndex >= 0 && liveTurn.steps.length === 0) { return turns; } // A send arm is only a presentation claim that the next message may still // arrive. It is not a Turn record and must not manufacture one while the // canonical transcript is catching up. A real live step (or steering // message) is sufficient evidence to project a missing external Turn. - if ( - targetIndex < 0 - && liveTurn.steps.length === 0 - && (liveTurn.pendingSteering?.length ?? 0) === 0 - ) { + if (targetIndex < 0 && liveTurn.steps.length === 0) { return turns; } const current = @@ -560,22 +552,20 @@ export function overlayLiveTurn( } } const liveTimeline: TurnTimelineItem[] = []; - const emittedSteeringIds = new Set(); - const appendLiveSteering = ( - messages: readonly LiveSteeringProjection[], - ): void => { - for (const message of messages) { - if (emittedSteeringIds.has(message.id)) continue; - emittedSteeringIds.add(message.id); - liveTimeline.push({ - kind: "user", - message: chatItemFromContent(message.id, message.ts, message.content), - messageId: message.id, - }); - } + const liveItemTs: number[] = []; + const pushLive = (item: TurnTimelineItem, ts: number | undefined): void => { + liveTimeline.push(item); + liveItemTs.push(ts ?? 0); + }; + const appendLiveSteering = (message: LiveSteeringProjection): void => { + pushLive({ + kind: "user", + message: chatItemFromContent(message.id, message.ts, message.content), + messageId: message.id, + }, message.ts); }; for (const step of liveTurn.steps) { - appendLiveSteering(step.leadingSteering ?? []); + if (step.steering !== undefined) appendLiveSteering(step.steering); const contentOrder = step.contentOrder ?? [ ...(step.thinking ? ["thinking" as const] : []), ...(step.text ? ["text" as const] : []), @@ -583,15 +573,15 @@ export function overlayLiveTurn( ]; for (const kind of contentOrder) { if (kind === "thinking" && step.thinking?.text) { - liveTimeline.push({ + pushLive({ kind: "thinking", text: step.thinking.text, messageId: step.stepId, live: step.thinking.complete !== true, truncated: step.thinking.truncated, - }); + }, step.startedAt); } else if (kind === "text" && step.text && (step.text.text || step.text.interrupted)) { - liveTimeline.push({ + pushLive({ kind: "text", text: step.text.text, ...(step.text.interrupted ? { interrupted: true } : {}), @@ -599,36 +589,48 @@ export function overlayLiveTurn( live: true, complete: step.text.complete, truncated: step.text.truncated, - }); + }, step.startedAt); } else if (kind === "tools") { const stepTools = step.tools.flatMap((tool) => { const projected = toolByUseId.get(tool.toolUseId); return projected ? [projected] : []; }); if (stepTools.length > 0) - liveTimeline.push({ kind: "tools", items: stepTools }); + pushLive({ kind: "tools", items: stepTools }, step.startedAt); } } } - appendLiveSteering(liveTurn.pendingSteering ?? []); // Shared entries are handoff points: replace them in place while preserving // live production order. Appending all live content after settled rows moved // an earlier answer (and its steering anchor) behind later persisted steps. - const liveEntries = flattenTimelineTools(liveTimeline); - const liveIndex = new Map(liveEntries.map((item, index) => [timelineItemKey(item), index])); + const liveEntries = liveTimeline.flatMap((item, index) => + flattenTimelineTools([item]).map((entry) => ({ entry, ts: liveItemTs[index]! }))); + const lastSettledContentIndex = current.timeline.findLastIndex((item) => item.kind !== 'user'); + const liveKeys = new Set(liveEntries.map(({ entry }) => timelineItemKey(entry))); + // A steering the live stream missed is position-ambiguous once it sits at + // the settled tail; splice it into the live order where its own ts falls + // instead of blindly trailing content that followed it. Missing ts sorts + // conservatively on both sides: a live entry without one keeps the deferred + // steering behind it, and a settled steering without one trails the live + // order rather than leaping ahead of known content. + const deferredSettled = new Set(); + for (const [index, item] of current.timeline.entries()) { + if (item.kind !== 'user' || item.steeringEventId === undefined + || index <= lastSettledContentIndex || liveKeys.has(timelineItemKey(item))) continue; + deferredSettled.add(index); + const ts = item.message.ts ?? Number.MAX_SAFE_INTEGER; + const at = liveEntries.findIndex((entry) => entry.ts > ts); + if (at < 0) liveEntries.push({ entry: item, ts }); + else liveEntries.splice(at, 0, { entry: item, ts }); + } + const liveIndex = new Map(liveEntries.map(({ entry }, index) => [timelineItemKey(entry), index])); const timeline: TurnTimelineItem[] = []; let nextLive = 0; const appendLiveThrough = (index: number) => { - while (nextLive <= index) timeline.push(liveEntries[nextLive++]!); + while (nextLive <= index) timeline.push(liveEntries[nextLive++]!.entry); }; - const lastSettledContentIndex = current.timeline.findLastIndex((item) => item.kind !== 'user'); - const deferredSteering: TurnTimelineItem[] = []; for (const [index, item] of current.timeline.entries()) { - if (item.kind === 'user' && item.steeringEventId !== undefined - && !liveIndex.has(timelineItemKey(item)) && index > lastSettledContentIndex) { - deferredSteering.push(item); - continue; - } + if (deferredSettled.has(index)) continue; for (const entry of flattenTimelineTools([item])) { const key = timelineItemKey(entry); const livePosition = liveIndex.get(key); @@ -637,7 +639,6 @@ export function overlayLiveTurn( } } appendLiveThrough(liveEntries.length - 1); - timeline.push(...deferredSteering); const mergedTimeline = mergeAdjacentTimeline(timeline); const next = { ...current,