diff --git a/packages/client/src/lib/__tests__/event-reducer.test.ts b/packages/client/src/lib/__tests__/event-reducer.test.ts index 525872475..d547f3ee1 100644 --- a/packages/client/src/lib/__tests__/event-reducer.test.ts +++ b/packages/client/src/lib/__tests__/event-reducer.test.ts @@ -3354,3 +3354,188 @@ describe("advisor message_end rows", () => { expect(unrelated.messages).toHaveLength(0); }); }); +describe("issue #126: live thinking flush on mid-turn tool execution", () => { + it("flushes streamingThinking into a thinking row when tool_execution_start occurs", () => { + let state = createInitialState(); + const t0 = 1000; + + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0, + data: { assistantMessageEvent: { type: "thinking_start" } }, + }, + { isLive: true }, + ); + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0 + 10, + data: { assistantMessageEvent: { type: "thinking_delta", delta: "Analyzing code..." } }, + }, + { isLive: true }, + ); + + expect(state.streamingThinking).toBe("Analyzing code..."); + expect(state.messages).toHaveLength(0); + + state = reduceEvent( + state, + { + eventType: "tool_execution_start", + timestamp: t0 + 20, + data: { toolCallId: "call_sub126", toolName: "agent" }, + }, + { isLive: true }, + ); + + expect(state.streamingThinking).toBe(""); + expect(state.thinkingStartedAt).toBeUndefined(); + + expect(state.messages).toHaveLength(2); + expect(state.messages[0]).toMatchObject({ + role: "thinking", + content: "Analyzing code...", + streamedLive: true, + startedAt: t0, + duration: 20, + }); + expect(state.messages[1]).toMatchObject({ + role: "toolResult", + toolCallId: "call_sub126", + toolName: "agent", + }); + }); + + it("does not create a duplicate thinking row when thinking_end fires after flush", () => { + let state = createInitialState(); + const t0 = 1000; + + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0, + data: { assistantMessageEvent: { type: "thinking_start" } }, + }, + { isLive: true }, + ); + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0 + 10, + data: { assistantMessageEvent: { type: "thinking_delta", delta: "Subagent prompt" } }, + }, + { isLive: true }, + ); + state = reduceEvent( + state, + { + eventType: "tool_execution_start", + timestamp: t0 + 20, + data: { toolCallId: "call_1", toolName: "task" }, + }, + { isLive: true }, + ); + + expect(state.messages.filter((m) => m.role === "thinking")).toHaveLength(1); + + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0 + 30, + data: { assistantMessageEvent: { type: "thinking_end" } }, + }, + { isLive: true }, + ); + + expect(state.messages.filter((m) => m.role === "thinking")).toHaveLength(1); + }); + + it("flushes previous thinking block when a new thinking_start arrives before thinking_end", () => { + let state = createInitialState(); + const t0 = 1000; + + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0, + data: { assistantMessageEvent: { type: "thinking_start" } }, + }, + { isLive: true }, + ); + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0 + 10, + data: { assistantMessageEvent: { type: "thinking_delta", delta: "First block" } }, + }, + { isLive: true }, + ); + + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0 + 50, + data: { assistantMessageEvent: { type: "thinking_start" } }, + }, + { isLive: true }, + ); + + expect(state.messages).toHaveLength(1); + expect(state.messages[0]).toMatchObject({ + role: "thinking", + content: "First block", + }); + expect(state.streamingThinking).toBe(""); + expect(state.thinkingStartedAt).toBe(t0 + 50); + }); + + it("respects isLive flag for streamedLive and isolates subagent state", () => { + let state = createInitialState(); + const t0 = 1000; + + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0, + data: { assistantMessageEvent: { type: "thinking_start" } }, + }, + { isLive: false }, + ); + state = reduceEvent( + state, + { + eventType: "message_update", + timestamp: t0 + 10, + data: { assistantMessageEvent: { type: "thinking_delta", delta: "Cold replay thinking" } }, + }, + { isLive: false }, + ); + state = reduceEvent( + state, + { + eventType: "tool_execution_start", + timestamp: t0 + 20, + data: { toolCallId: "call_sub", toolName: "agent" }, + }, + { isLive: false }, + ); + + expect(state.messages[0]).toMatchObject({ + role: "thinking", + content: "Cold replay thinking", + streamedLive: false, + }); + + expect(state.subagents.size).toBe(0); + }); +}); diff --git a/packages/client/src/lib/event-reducer.ts b/packages/client/src/lib/event-reducer.ts index d0df2c33a..c09d4cf70 100644 --- a/packages/client/src/lib/event-reducer.ts +++ b/packages/client/src/lib/event-reducer.ts @@ -727,6 +727,59 @@ const TURN_BOUNDARY_ROLES: ReadonlySet = new Set([ * @param timestamp Event timestamp (used as the row's `timestamp`) * @param toolCallId Id of the upcoming tool — used as the row's stable id anchor */ +/** + * Flush the current `streamingThinking` into a permanent `thinking` ChatMessage + * row. Called when a tool/subagent is executed mid-thought or when a new thinking + * block begins before `thinking_end` fires. + */ +export function flushStreamingThinkingAsRow( + state: SessionState, + timestamp: number, + isLive: boolean, + seq?: number, + toolCallId?: string, +): SessionState { + if (!state.streamingThinking) return state; + + const startedAt = state.thinkingStartedAt; + const duration = startedAt ? timestamp - startedAt : undefined; + const id = toolCallId + ? seq !== undefined + ? `thinking-${seq}-flush-${toolCallId}` + : `thinking-${state.messages.length}-flush-${toolCallId}` + : seq !== undefined + ? `thinking-${seq}-0` + : `thinking-${state.messages.length}`; + + if (state.messages.some((m) => m.id === id)) { + return { + ...state, + streamingThinking: "", + streamingThinkingCollapsed: false, + thinkingStartedAt: undefined, + }; + } + + const thinkingRow: ChatMessage = { + id, + role: "thinking", + content: state.streamingThinking, + timestamp, + startedAt, + duration, + streamedLive: state.streamingThinkingCollapsed ? false : isLive, + seq, + }; + + return { + ...state, + messages: [...state.messages, thinkingRow], + streamingThinking: "", + streamingThinkingCollapsed: false, + thinkingStartedAt: undefined, + }; +} + export function flushStreamingTextAsAssistantRow( state: SessionState, timestamp: number, @@ -1589,6 +1642,12 @@ export function reduceEvent( // Handle thinking events from assistantMessageEvent if (assistantEvent) { if (assistantEvent.type === "thinking_start") { + if (next.streamingThinking) { + Object.assign( + next, + flushStreamingThinkingAsRow(next, event.timestamp, isLive, seq), + ); + } next.pendingThinking = ""; next.streamingThinking = ""; next.streamingThinkingCollapsed = false; @@ -1848,6 +1907,12 @@ export function reduceEvent( // See change: fix-reducer-crash-undefined-toolname. const toolName = typeof data.toolName === "string" ? data.toolName : "unknown"; + if (next.streamingThinking) { + Object.assign( + next, + flushStreamingThinkingAsRow(next, event.timestamp, isLive, seq, toolCallId), + ); + } // Flush any pending streamingText into a permanent assistant row // BEFORE pushing the new toolResult, so the message's content-array // order is preserved in messages[] for the entire tool runtime —