diff --git a/packages/app/src/timeline/session-stream-reducers.test.ts b/packages/app/src/timeline/session-stream-reducers.test.ts index 6c2ce8d90..848c34982 100644 --- a/packages/app/src/timeline/session-stream-reducers.test.ts +++ b/packages/app/src/timeline/session-stream-reducers.test.ts @@ -69,6 +69,14 @@ function makeTimelineEvent( } as AgentStreamEventPayload; } +function makeAssistantTimelineEvent(text: string, messageId?: string): AgentStreamEventPayload { + return { + type: "timeline", + provider: "claude", + item: { type: "assistant_message", text, ...(messageId ? { messageId } : {}) }, + } as AgentStreamEventPayload; +} + function makeToolCallTimelineEvent(callId: string): AgentStreamEventPayload { return { type: "timeline", @@ -1523,6 +1531,65 @@ describe("processAgentStreamEvents", () => { expect(result.sideEffects).toEqual([]); }); + it("keeps matching assistant message ids in the live head", () => { + const result = processAgentStreamEvents({ + events: [ + makeStreamReducerEvent(makeAssistantTimelineEvent("Hel", "assistant-one"), 1), + makeStreamReducerEvent(makeAssistantTimelineEvent("lo", "assistant-one"), 2), + ], + currentTail: [], + currentHead: [], + currentCursor: undefined, + currentAgent: null, + }); + + expect(result.changedTail).toBe(false); + expect(result.changedHead).toBe(true); + expect(result.tail).toEqual([]); + expect(result.head).toHaveLength(1); + expect(result.head[0]).toMatchObject({ + kind: "assistant_message", + text: "Hello", + messageId: "assistant-one", + }); + }); + + it("flushes the live assistant head before starting a different assistant message id", () => { + const result = processAgentStreamEvents({ + events: [ + makeStreamReducerEvent(makeAssistantTimelineEvent("First", "assistant-one"), 1), + makeStreamReducerEvent(makeAssistantTimelineEvent("Second", "assistant-two"), 2), + ], + currentTail: [], + currentHead: [], + currentCursor: undefined, + currentAgent: null, + }); + + expect(result.changedTail).toBe(true); + expect(result.changedHead).toBe(true); + expect(getAssistantTexts(result.tail)).toEqual(["First"]); + expect(getAssistantTexts(result.head)).toEqual(["Second"]); + }); + + it("flushes an anonymous assistant head before starting an identified assistant message", () => { + const result = processAgentStreamEvents({ + events: [ + makeStreamReducerEvent(makeAssistantTimelineEvent("Anonymous"), 1), + makeStreamReducerEvent(makeAssistantTimelineEvent("Identified", "assistant-two"), 2), + ], + currentTail: [], + currentHead: [], + currentCursor: undefined, + currentAgent: null, + }); + + expect(result.changedTail).toBe(true); + expect(result.changedHead).toBe(true); + expect(getAssistantTexts(result.tail)).toEqual(["Anonymous"]); + expect(getAssistantTexts(result.head)).toEqual(["Identified"]); + }); + it("promotes completed assistant markdown blocks to tail while keeping the live block in head", () => { const result = processAgentStreamEvents({ events: [ @@ -1632,6 +1699,92 @@ describe("processAgentStreamEvents", () => { } satisfies TimelineCursor); }); + it("keeps Claude image tool-result output before following assistant blocks while text streams", () => { + const imageMarkdown = + "![Image](/tmp/paseo-attachments/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.png)"; + const result = processAgentStreamEvents({ + events: [ + makeStreamReducerEvent( + { + type: "timeline", + provider: "claude", + item: { + type: "assistant_message", + text: "ABC", + messageId: "assistant-before", + }, + } as AgentStreamEventPayload, + 1, + ), + makeStreamReducerEvent(makeToolCallTimelineEvent("toolu_read_png"), 2), + makeStreamReducerEvent( + { + type: "timeline", + provider: "claude", + item: { + type: "tool_call", + callId: "toolu_read_png", + name: "Read", + status: "completed", + detail: { + type: "read", + filePath: "/tmp/image.png", + }, + error: null, + }, + } as AgentStreamEventPayload, + 3, + ), + makeStreamReducerEvent( + { + type: "timeline", + provider: "claude", + item: { + type: "assistant_message", + text: imageMarkdown, + }, + } as AgentStreamEventPayload, + 4, + ), + makeStreamReducerEvent( + { + type: "timeline", + provider: "claude", + item: { + type: "assistant_message", + text: "D", + messageId: "assistant-after", + }, + } as AgentStreamEventPayload, + 5, + ), + makeStreamReducerEvent( + { + type: "timeline", + provider: "claude", + item: { + type: "assistant_message", + text: "\n\nE", + messageId: "assistant-after", + }, + } as AgentStreamEventPayload, + 6, + ), + ], + currentTail: [], + currentHead: [], + currentCursor: undefined, + currentAgent: null, + }); + + expect(getAssistantTexts([...result.tail, ...result.head])).toEqual([ + "ABC", + imageMarkdown, + "D", + "E", + ]); + }); + it("returns the final optimistic lifecycle patch across a batch", () => { const result = processAgentStreamEvents({ events: [ diff --git a/packages/app/src/types/stream.ts b/packages/app/src/types/stream.ts index 67c7480d0..fb711fd2a 100644 --- a/packages/app/src/types/stream.ts +++ b/packages/app/src/types/stream.ts @@ -888,6 +888,13 @@ function getEventItemKind(event: AgentStreamEventPayload): StreamItem["kind"] | } } +function getIncomingAssistantMessageId(event: AgentStreamEventPayload): string | undefined { + if (event.type !== "timeline" || event.item.type !== "assistant_message") { + return undefined; + } + return event.item.messageId; +} + /** * Finalize head items before flushing to tail. * Marks thoughts as "ready" since they're no longer being streamed. @@ -939,10 +946,7 @@ function getTailAssistantToResume(params: { if (params.tailAssistant?.kind !== "assistant_message") { return null; } - const incomingMessageId = - params.event.type === "timeline" && params.event.item.type === "assistant_message" - ? params.event.item.messageId - : undefined; + const incomingMessageId = getIncomingAssistantMessageId(params.event); if (incomingMessageId !== undefined && params.tailAssistant.messageId !== incomingMessageId) { return null; } @@ -1036,9 +1040,14 @@ export function flushHeadToTail(tail: StreamItem[], head: StreamItem[]): StreamI /** * Determine if the head should be flushed based on incoming event kind. - * Flush when a different kind arrives or when the incoming kind is not streamable. + * Flush when a different streamable lane starts, including a new identified assistant message. */ -function shouldFlushHead(head: StreamItem[], incomingKind: StreamItem["kind"] | null): boolean { +function shouldFlushHead(input: { + head: StreamItem[]; + incomingKind: StreamItem["kind"] | null; + event: AgentStreamEventPayload; +}): boolean { + const { head, incomingKind, event } = input; if (head.length === 0) { return false; } @@ -1072,6 +1081,11 @@ function shouldFlushHead(head: StreamItem[], incomingKind: StreamItem["kind"] | return true; } + if (incomingKind === "assistant_message" && lastStreamable.kind === "assistant_message") { + const incomingMessageId = getIncomingAssistantMessageId(event); + return incomingMessageId !== undefined && lastStreamable.messageId !== incomingMessageId; + } + return false; } @@ -1134,7 +1148,13 @@ export function applyStreamEvent(params: { const incomingKind = getEventItemKind(event); // Check if we need to flush head before processing this event - if (shouldFlushHead(nextHead, incomingKind)) { + if ( + shouldFlushHead({ + head: nextHead, + incomingKind, + event, + }) + ) { flushHead(); }