diff --git a/packages/app/src/contexts/session-stream-reducers.test.ts b/packages/app/src/contexts/session-stream-reducers.test.ts index 5a65c597a..ee4ffff4b 100644 --- a/packages/app/src/contexts/session-stream-reducers.test.ts +++ b/packages/app/src/contexts/session-stream-reducers.test.ts @@ -910,6 +910,67 @@ describe("processAgentStreamEvents", () => { lastActivityAt: new Date(3000), }); }); + + it("keeps a live Claude assistant paragraph contiguous when init tail hydration lands mid-stream", () => { + const seq186Text = "Now let me write the updated tool-call-detail-parser.ts with all the sub"; + const seq187Text = "agent additions."; + + const liveSeq186 = processAgentStreamEvent({ + ...baseStreamInput, + event: makeTimelineEvent(seq186Text), + seq: 186, + epoch: "epoch-1", + currentTail: [], + currentHead: [], + currentCursor: undefined, + timestamp: new Date("2026-05-02T10:00:00.186Z"), + }); + + expect(getAssistantTexts(liveSeq186.head)).toEqual([seq186Text]); + expect(getAssistantTexts(liveSeq186.tail)).toEqual([]); + + const initTailHydration = processTimelineResponse({ + ...baseTimelineInput, + currentTail: liveSeq186.tail, + currentHead: liveSeq186.head, + currentCursor: liveSeq186.cursor ?? undefined, + isInitializing: true, + hasActiveInitDeferred: true, + initRequestDirection: "tail", + payload: { + agentId: "agent-1", + direction: "tail", + reset: false, + epoch: "epoch-1", + startCursor: { seq: 186 }, + endCursor: { seq: 186 }, + entries: [makeTimelineEntry(186, seq186Text)], + error: null, + }, + }); + + expect(getAssistantTexts(initTailHydration.tail)).toEqual([]); + expect(getAssistantTexts(initTailHydration.head)).toEqual([seq186Text]); + + const liveSeq187 = processAgentStreamEvent({ + ...baseStreamInput, + event: makeTimelineEvent(seq187Text), + seq: 187, + epoch: "epoch-1", + currentTail: initTailHydration.tail, + currentHead: initTailHydration.head, + currentCursor: initTailHydration.cursor ?? undefined, + timestamp: new Date("2026-05-02T10:00:00.187Z"), + }); + + const finalAssistantItems = [...liveSeq187.tail, ...liveSeq187.head].filter( + (item): item is Extract => + item.kind === "assistant_message", + ); + + expect(finalAssistantItems).toHaveLength(1); + expect(finalAssistantItems[0]?.text).toBe(`${seq186Text}${seq187Text}`); + }); }); describe("createAgentStreamReducerQueue", () => { diff --git a/packages/app/src/contexts/session-stream-reducers.ts b/packages/app/src/contexts/session-stream-reducers.ts index 1617cd7a6..b8aa311ad 100644 --- a/packages/app/src/contexts/session-stream-reducers.ts +++ b/packages/app/src/contexts/session-stream-reducers.ts @@ -99,16 +99,49 @@ interface TimelinePathResult { sideEffects: TimelineReducerSideEffect[]; } +function preserveReplacePathAssistantHead(params: { + tail: StreamItem[]; + currentHead: StreamItem[]; +}): { + tail: StreamItem[]; + head: StreamItem[]; +} { + const { tail, currentHead } = params; + const liveAssistant = currentHead.findLast( + (item): item is Extract => + item.kind === "assistant_message", + ); + if (!liveAssistant) { + return { tail, head: [] }; + } + const tailAssistant = tail.at(-1); + if (!tailAssistant || tailAssistant.kind !== "assistant_message") { + return { tail, head: currentHead }; + } + if (!liveAssistant.text.startsWith(tailAssistant.text)) { + return { tail, head: [] }; + } + return { + tail: tail.slice(0, -1), + head: [{ ...liveAssistant, text: tailAssistant.text }], + }; +} + function applyTimelineReplacePath(args: { timelineUnits: TimelineUnit[]; payload: ProcessTimelineResponseInput["payload"]; bootstrapPolicy: ReturnType; + currentHead: StreamItem[]; toHydratedEvents: ( units: TimelineUnit[], ) => Array<{ event: AgentStreamEventPayload; timestamp: Date }>; }): TimelinePathResult { - const { timelineUnits, payload, bootstrapPolicy, toHydratedEvents } = args; - const tail = hydrateStreamState(toHydratedEvents(timelineUnits), { source: "canonical" }); + const { timelineUnits, payload, bootstrapPolicy, currentHead, toHydratedEvents } = args; + const hydratedTail = hydrateStreamState(toHydratedEvents(timelineUnits), { source: "canonical" }); + const { tail, head } = preserveReplacePathAssistantHead({ + tail: hydratedTail, + currentHead, + }); const cursor: TimelineCursor | null = payload.startCursor && payload.endCursor ? { @@ -121,7 +154,7 @@ function applyTimelineReplacePath(args: { if (bootstrapPolicy.catchUpCursor) { sideEffects.push({ type: "catch_up", cursor: bootstrapPolicy.catchUpCursor }); } - return { tail, head: [], cursor, cursorChanged: true, sideEffects }; + return { tail, head, cursor, cursorChanged: true, sideEffects }; } interface IncrementalAcceptResult { @@ -307,6 +340,7 @@ export function processTimelineResponse( timelineUnits, payload, bootstrapPolicy, + currentHead, toHydratedEvents, }) : applyTimelineIncrementalPath({