diff --git a/services/api/src/ai/sse-mapper.ts b/services/api/src/ai/sse-mapper.ts index 219a5e52..e8bc48d9 100644 --- a/services/api/src/ai/sse-mapper.ts +++ b/services/api/src/ai/sse-mapper.ts @@ -173,7 +173,7 @@ export class ChannelTokenRouter { } -export function mapEventToSSE(event: Record): SSEEvent | null { +export function mapEventToSSE(event: Record, options: { preserveThinkingMarkers?: boolean } = {}): SSEEvent | null { const type = event.type as string; const data = event.data as Record | undefined; @@ -182,14 +182,14 @@ export function mapEventToSSE(event: Record): SSEEvent | null { case "assistant.message_delta": { const delta = (data?.deltaContent ?? "") as string; if (!delta) return null; - return { type: "text_delta", data: sanitizeText(delta) }; + return { type: "text_delta", data: sanitizeText(delta, options) }; } // ─── SDK v0.2.0 streaming delta (raw text chunks) ──── case "assistant.streaming_delta": { const streamDelta = (data?.deltaContent ?? data?.content ?? data?.delta ?? "") as string; if (!streamDelta) return null; - return { type: "text_delta", data: sanitizeText(streamDelta) }; + return { type: "text_delta", data: sanitizeText(streamDelta, options) }; } // ─── Final complete message (sent after streaming ends) ─ @@ -199,7 +199,7 @@ export function mapEventToSSE(event: Record): SSEEvent | null { // ─── Legacy / direct provider text events ───────────── case "text_delta": { const raw = (data?.content ?? data ?? "") as string; - return { type: "text_delta", data: sanitizeText(String(raw)) }; + return { type: "text_delta", data: sanitizeText(String(raw), options) }; } // ─── Streaming reasoning deltas (token-by-token thinking) ── diff --git a/services/api/src/ai/tool-messages.ts b/services/api/src/ai/tool-messages.ts index 1a90d49c..077a7a71 100644 --- a/services/api/src/ai/tool-messages.ts +++ b/services/api/src/ai/tool-messages.ts @@ -293,7 +293,7 @@ export function sanitizeCommand(cmd: string): string { return result; } -export function sanitizeText(text: string): string { +export function sanitizeText(text: string, options: { preserveThinkingMarkers?: boolean } = {}): string { if (!text) return text; let result = text; @@ -303,10 +303,12 @@ export function sanitizeText(text: string): string { // <|channel>thought, , <|channel|> — Gemma 4 // , — Claude (when prompted) // , — DeepSeek (post-thinking answer marker) - result = result.replace(/<\/?think>/gi, ""); - result = result.replace(/<\|?channel\|?>(?:thought)?/gi, ""); - result = result.replace(/<\/?rationale>/gi, ""); - result = result.replace(/<\/?answer>/gi, ""); + if (!options.preserveThinkingMarkers) { + result = result.replace(/<\/?think>/gi, ""); + result = result.replace(/<\|?channel\|?>(?:thought)?/gi, ""); + result = result.replace(/<\/?rationale>/gi, ""); + result = result.replace(/<\/?answer>/gi, ""); + } // 1. Strip absolute server paths result = stripServerPaths(result); diff --git a/services/api/src/routes/chat/event-processor.ts b/services/api/src/routes/chat/event-processor.ts index c2009f2a..8d28d40c 100644 --- a/services/api/src/routes/chat/event-processor.ts +++ b/services/api/src/routes/chat/event-processor.ts @@ -80,8 +80,11 @@ export function createProcessEvent( } state.lastCapturedMsgId = deltaMessageId; state.msgIdDeltaStart = state.assistantContent.length; + state.currentMessageTextLength = 0; state.lastMsgIdSepEmitted = true; } + const rawText = String(evtData?.deltaContent ?? evtData?.content ?? evtData?.delta ?? ""); + state.currentMessageTextLength += rawText.length; } // assistant.message catch-up (BUG-119) @@ -95,7 +98,7 @@ export function createProcessEvent( // later only if recovery fails (see send-handler.ts). // EXCEPTION: Rate limit errors are sent immediately — they are not // transient and the user needs to know why generation stopped. - const sseData = mapEventToSSE(event); + const sseData = mapEventToSSE(event, { preserveThinkingMarkers: true }); if (sseData) { if (evtType === "session.error" && sseData.type === "error") { const errMsg = typeof sseData.data === "string" ? sseData.data : "Unknown error"; @@ -145,55 +148,21 @@ function handleAssistantMessageCatchUp( } state.lastCapturedMsgId = msgId; state.msgIdDeltaStart = state.assistantContent.length; + state.currentMessageTextLength = 0; } // Reset the flag after catch-up so the next transition works fresh state.lastMsgIdSepEmitted = false; if (!content) return; - const sanitizedContent = sanitizeText(content); - const deltasSoFar = state.assistantContent.slice(state.msgIdDeltaStart); - // Account for text we classified as thinking via the leading-text buffer. - // The SDK's assistant.message includes ALL text (reasoning + content), but - // our delta handler split it into assistantContent and assistantThinking. - // Without this, the catch-up would see thinking text as "missing" and - // re-emit it as text_delta, leaking reasoning into the chat. - const totalProcessed = deltasSoFar.length + state.assistantThinking.length; - if (sanitizedContent.length > totalProcessed) { - const missing = sanitizedContent.slice(totalProcessed); - console.log(`[Chat][${projectId.slice(0, 8)}] catch-up: msg=${sanitizedContent.length} processed=${totalProcessed} (content=${deltasSoFar.length} thinking=${state.assistantThinking.length}) missing=${missing.length}`); - let visibleText = ""; - for (const chunk of channelRouter.process(missing)) { - if (!chunk.content) continue; - if (chunk.type === "text") { - visibleText += chunk.content; - stream.writeSSE({ data: JSON.stringify({ type: "text_delta", data: chunk.content }) }).catch(() => {}); - } else if (chunk.type === "thinking") { - state.assistantThinking += chunk.content; - stream.writeSSE({ data: JSON.stringify({ type: "thinking", data: stripServerPaths(chunk.content) }) }).catch(() => {}); - } else if (chunk.type === "tool") { - state.sawToolDelta = true; - stream.writeSSE({ data: JSON.stringify({ type: "tool_delta", data: chunk.content }) }).catch(() => {}); - } - } - if (visibleText) { - state.assistantContent = state.assistantContent.slice(0, state.msgIdDeltaStart) + deltasSoFar + visibleText; - } - } else if (!totalProcessed && !state.assistantContent) { - let visibleText = ""; - for (const chunk of channelRouter.process(sanitizedContent)) { - if (!chunk.content) continue; - if (chunk.type === "text") { - visibleText += chunk.content; - stream.writeSSE({ data: JSON.stringify({ type: "text_delta", data: chunk.content }) }).catch(() => {}); - } else if (chunk.type === "thinking") { - state.assistantThinking += chunk.content; - stream.writeSSE({ data: JSON.stringify({ type: "thinking", data: stripServerPaths(chunk.content) }) }).catch(() => {}); - } else if (chunk.type === "tool") { - state.sawToolDelta = true; - stream.writeSSE({ data: JSON.stringify({ type: "tool_delta", data: chunk.content }) }).catch(() => {}); - } - } - state.assistantContent = visibleText; + // Count only deltas from THIS message, including markers. Session-wide + // thinking length can hide a complete final answer from a non-streaming + // provider. Slice raw text BEFORE sanitizing: jargon replacements can change + // lengths differently for partial deltas and complete messages. + const missing = sanitizeText(content.slice(state.currentMessageTextLength), { preserveThinkingMarkers: true }); + if (missing) { + routeSseEvent(stream, state, channelRouter, { type: "text_delta", data: missing }, evtData, projectId, userId, messageId); + state.currentMessageTextLength = content.length; } + } function routeSseEvent( diff --git a/services/api/src/routes/chat/final-response.integration.test.ts b/services/api/src/routes/chat/final-response.integration.test.ts new file mode 100644 index 00000000..b465b503 --- /dev/null +++ b/services/api/src/routes/chat/final-response.integration.test.ts @@ -0,0 +1,150 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { createInitialState } from "./types.js"; +import { createProcessEvent } from "./event-processor.js"; +import { ChannelTokenRouter } from "../../ai/sse-mapper.js"; +import { finalizeLeadingResponse } from "./final-response.js"; + +function fixture() { + const state = createInitialState(); + const frames: Array<{ type: string; data: unknown }> = []; + const stream = { + writeSSE: async ({ data }: { data: string }) => { + frames.push(JSON.parse(data)); + }, + } as unknown as import("hono/streaming").SSEStreamingApi; + const process = createProcessEvent( + stream, + state, + new ChannelTokenRouter(), + "regression-test", + "test-user", + "test-message", + "build", + () => {}, + () => {}, + () => "test-session", + ); + const emit = (type: string, data: Record) => + process({ type, data } as Parameters[0]); + return { state, frames, emit }; +} + +test("SDK tool round-trip and final deltas preserve answer and native reasoning", () => { + const { state, frames, emit } = fixture(); + emit("assistant.reasoning_delta", { deltaContent: "Native reasoning." }); + emit("assistant.message_delta", { + messageId: "m1", + deltaContent: "I will inspect the files.", + }); + emit("tool.execution_start", { + toolName: "read_file", + toolCallId: "t1", + arguments: { path: "package.json" }, + }); + emit("tool.execution_complete", { + toolName: "read_file", + toolCallId: "t1", + success: true, + }); + emit("assistant.message_delta", { + messageId: "m2", + deltaContent: "这是一个 React 项目。", + }); + emit("assistant.message", { + messageId: "m2", + content: "这是一个 React 项目。", + }); + emit("session.idle", {}); + assert.equal(state.hadToolCalls, true); + assert.ok(frames.some((frame) => frame.type === "tool_result")); + assert.equal(finalizeLeadingResponse(state), "这是一个 React 项目。"); + assert.equal(state.assistantContent, "这是一个 React 项目。"); + assert.equal( + state.assistantThinking, + "Native reasoning.I will inspect the files.", + ); +}); + +test("tagged reasoning after a tool is not promoted with the final answer", () => { + const { state, emit } = fixture(); + emit("tool.execution_complete", { toolName: "read_file", success: true }); + emit("assistant.message_delta", { + messageId: "m1", + deltaContent: "Internal reasoning.Final answer.", + }); + emit("assistant.message", { + messageId: "m1", + content: "Internal reasoning.Final answer.", + }); + assert.equal(finalizeLeadingResponse(state), "Final answer."); + assert.equal(state.assistantThinking, "Internal reasoning."); +}); + +test("text followed by another tool call is not a final answer", () => { + const { state, emit } = fixture(); + emit("assistant.message_delta", { + messageId: "m1", + deltaContent: "I need another file.", + }); + emit("tool.execution_start", { toolName: "read_file", toolCallId: "t1" }); + emit("tool.execution_complete", { + toolName: "read_file", + toolCallId: "t1", + success: true, + }); + assert.equal(finalizeLeadingResponse(state), ""); + assert.equal(state.assistantContent, ""); + assert.equal(state.assistantThinking, "I need another file."); +}); + +test("complete-only final SDK message is not hidden by prior reasoning length", () => { + const { state, emit } = fixture(); + emit("assistant.reasoning_delta", { + deltaContent: "Long reasoning from an earlier tool round. ".repeat(20), + }); + emit("tool.execution_complete", { toolName: "read_file", success: true }); + emit("assistant.message", { messageId: "final", content: "Final answer." }); + assert.equal(finalizeLeadingResponse(state), "Final answer."); + assert.equal(state.assistantContent, "Final answer."); +}); + +test("thinking markers split across deltas stay out of the answer", () => { + const { state, emit } = fixture(); + for (const deltaContent of ["Reasoning", "Answer."]) { + emit("assistant.message_delta", { messageId: "final", deltaContent }); + } + emit("assistant.message", { + messageId: "final", + content: "ReasoningAnswer.", + }); + assert.equal(finalizeLeadingResponse(state), "Answer."); + assert.equal(state.assistantContent, "Answer."); + assert.equal(state.assistantThinking, "Reasoning"); +}); + +test("catch-up uses raw offsets when sanitization expands split jargon", () => { + const { state, emit } = fixture(); + const content = "Read package.json. Final answer."; + for (const deltaContent of ["Read pack", "age.json. Final answer."]) { + emit("assistant.message_delta", { messageId: "final", deltaContent }); + } + emit("assistant.message", { messageId: "final", content }); + assert.equal(finalizeLeadingResponse(state), content); +}); + +test("catch-up sanitizes only a genuinely missing raw suffix", () => { + const { state, emit } = fixture(); + emit("assistant.message_delta", { + messageId: "final", + deltaContent: "Read package.json.", + }); + emit("assistant.message", { + messageId: "final", + content: "Read package.json. Final answer.", + }); + assert.equal( + finalizeLeadingResponse(state), + "Read project configuration. Final answer.", + ); +}); diff --git a/services/api/src/routes/chat/final-response.test.ts b/services/api/src/routes/chat/final-response.test.ts new file mode 100644 index 00000000..d31f2e97 --- /dev/null +++ b/services/api/src/routes/chat/final-response.test.ts @@ -0,0 +1,73 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { createInitialState } from "./types.js"; +import { finalizeLeadingResponse } from "./final-response.js"; + +const answer = "这是一个空白 React 项目。可以添加业务页面。没有修改文件。"; + +test("final answer after read-only tools becomes persisted content", () => { + const state = createInitialState(); + state.hadToolCalls = true; + state.assistantThinking = "Earlier tool planning.\n" + answer; + state.leadingTextBuffer = answer; + assert.equal(finalizeLeadingResponse(state), answer); + assert.equal(state.assistantContent, answer); + assert.equal(state.assistantThinking, "Earlier tool planning.\n"); + assert.equal(state.leadingTextBuffer, ""); + assert.equal( + finalizeLeadingResponse(state), + "", + "finalization is idempotent", + ); +}); + +test("direct answers still work without tools", () => { + const state = createInitialState(); + state.assistantThinking = state.leadingTextBuffer = answer; + assert.equal(finalizeLeadingResponse(state), answer); + assert.equal(state.assistantThinking, ""); +}); + +test("explicit reasoning without a plain-text final segment stays private", () => { + const state = createInitialState(); + state.hadToolCalls = true; + state.assistantThinking = "Native reasoning only"; + assert.equal(finalizeLeadingResponse(state), ""); + assert.equal(state.assistantContent, ""); + assert.equal(state.assistantThinking, "Native reasoning only"); +}); + +test("a deferred stream error does not promote an unfinished segment", () => { + const state = createInitialState(); + state.assistantThinking = state.leadingTextBuffer = "I still need to inspect"; + state.deferredError = "Request timed out"; + assert.equal(finalizeLeadingResponse(state), ""); + assert.equal(state.assistantContent, ""); + assert.equal(state.assistantThinking, "I still need to inspect"); +}); + +test("later native reasoning is preserved rather than removed as a suffix", () => { + const state = createInitialState(); + state.hadToolCalls = true; + state.leadingTextBuffer = answer; + state.assistantThinking = + "Earlier reasoning. " + answer + " Later reasoning."; + assert.equal(finalizeLeadingResponse(state), answer); + assert.equal(state.assistantThinking, "Earlier reasoning. Later reasoning."); +}); + +test("previous visible content is retained", () => { + const state = createInitialState(); + state.assistantContent = "Existing answer. "; + state.assistantThinking = state.leadingTextBuffer = answer; + assert.equal(finalizeLeadingResponse(state), answer); + assert.equal(state.assistantContent, "Existing answer. " + answer); +}); + +test("unmatched/interleaved reasoning is conservatively left untouched", () => { + const state = createInitialState(); + state.leadingTextBuffer = "Hello world"; + state.assistantThinking = "Hello native reasoning world"; + assert.equal(finalizeLeadingResponse(state), ""); + assert.equal(state.assistantContent, ""); +}); diff --git a/services/api/src/routes/chat/final-response.ts b/services/api/src/routes/chat/final-response.ts new file mode 100644 index 00000000..1669af3b --- /dev/null +++ b/services/api/src/routes/chat/final-response.ts @@ -0,0 +1,27 @@ +import type { ChatStreamState } from "./types.js"; + +/** + * Ordinary text is provisionally routed to thinking while waiting to see + * whether another tool call follows. A successfully completed final segment + * is the answer, even when earlier segments invoked tools. Explicit reasoning + * events and text confirmed by a subsequent tool call are never promoted. + */ +export function finalizeLeadingResponse(state: ChatStreamState): string { + const buffered = state.leadingTextBuffer; + state.leadingTextBuffer = ""; + state.leadingTextFlushed = true; + if (!buffered || state.deferredError) return ""; + + // Remove only the exact provisional segment, not an arbitrary suffix: a + // native reasoning event may have arrived after the ordinary text. + const index = state.assistantThinking.lastIndexOf(buffered); + if (index < 0) return ""; + const visible = buffered.replace(/[\s\S]*?<\/think>\s*/gi, "").trim(); + if (!visible) return ""; + + state.assistantThinking = + state.assistantThinking.slice(0, index) + + state.assistantThinking.slice(index + buffered.length); + state.assistantContent += visible; + return visible; +} diff --git a/services/api/src/routes/chat/send-handler.ts b/services/api/src/routes/chat/send-handler.ts index 13067c7d..dbb8809f 100644 --- a/services/api/src/routes/chat/send-handler.ts +++ b/services/api/src/routes/chat/send-handler.ts @@ -34,6 +34,7 @@ import { projectSessions, activeRequests } from "./session-state.js"; import { buildSystemPrompt } from "./system-prompts.js"; import { createRecordAssistantToolCall, createToolProgressCallbacks } from "./tool-callbacks.js"; import { createProcessEvent } from "./event-processor.js"; +import { finalizeLeadingResponse } from "./final-response.js"; import { popArtifacts } from "./artifact-stash.js"; import { scaffoldAndStartDev, emitConfigTraces, logToolManifest, handleToolEndEvent } from "./send-helpers.js"; import { checkAndEvictOnModeChange, checkAndEvictOnProviderChange, resolveSession, persistSessionToDb, filterToolsForMode, recreateSession } from "./session-manager.js"; @@ -736,38 +737,14 @@ export function registerSendHandler(app: Hono) { } } - // ── Finalize leading-text buffer at stream end ── - // Post-tool text stays as thinking — tool results (file changes, - // build cards, MCP UI resources) provide all the visible UI the - // user needs. Converting the buffer to content leaks internal - // reasoning (BUG-119: MiniMax emits untagged reasoning as text). - // However, if no tool calls occurred AND no content was emitted, - // the buffer is the actual response (e.g. simple chat greeting). - if (state.leadingTextBuffer) { - const bufLen = state.leadingTextBuffer.length; - if (!state.hadToolCalls && state.assistantContent.length === 0) { - // No tool calls, no content — this IS the response, not reasoning - const buffered = state.leadingTextBuffer; - state.leadingTextBuffer = ""; - state.leadingTextFlushed = true; - // Strip any ... blocks that the channel router - // didn't catch during streaming (token boundary issue) - const visibleContent = buffered.replace(/[\s\S]*?<\/think>\s*/gi, "").trim(); - if (visibleContent) { - // Move buffer from thinking to content (only the visible portion) - state.assistantThinking = state.assistantThinking.slice(0, state.assistantThinking.length - buffered.length); - state.assistantContent += visibleContent; - console.log(`[Chat][${projectId.slice(0, 8)}] Flushing ${visibleContent.length} chars as content (stripped from ${bufLen} buffer, no tools)`); - broadcastToRoom(projectId, { type: "ai:stream-chunk", chunk: visibleContent, messageId, isThinking: false }, userId).catch(() => {}); - await stream.writeSSE({ data: JSON.stringify({ type: "thinking_to_text", data: visibleContent }) }); - } else { - console.log(`[Chat][${projectId.slice(0, 8)}] Keeping ${bufLen} chars as thinking (all thinking, no visible content)`); - } - } else { - state.leadingTextBuffer = ""; - state.leadingTextFlushed = true; - console.log(`[Chat][${projectId.slice(0, 8)}] Keeping ${bufLen} chars as thinking (stream end, hadTools=${state.hadToolCalls})`); - } + // A tool call confirms preceding text as intermediate reasoning. + // The remaining ordinary-text segment at successful stream end is + // the final answer, regardless of tools in earlier turns. + const finalResponse = finalizeLeadingResponse(state); + if (finalResponse) { + state.traceCollector?.onTextDelta(finalResponse); + broadcastToRoom(projectId, { type: "ai:stream-chunk", chunk: finalResponse, messageId, isThinking: false }, userId).catch(() => {}); + await stream.writeSSE({ data: JSON.stringify({ type: "thinking_to_text", data: finalResponse }) }); } console.log(`[Chat][${projectId.slice(0, 8)}] stream done — content: ${state.assistantContent.length}, thinking: ${state.assistantThinking.length}, tools: ${state.hadToolCalls}`); diff --git a/services/api/src/routes/chat/types.ts b/services/api/src/routes/chat/types.ts index 8c76d00d..6d57fbb9 100644 --- a/services/api/src/routes/chat/types.ts +++ b/services/api/src/routes/chat/types.ts @@ -18,6 +18,8 @@ export interface ChatStreamState { lastCapturedMsgId: string | undefined; lastMsgIdSepEmitted: boolean; msgIdDeltaStart: number; + /** Raw ordinary-text delta length for the current SDK message, before routing. */ + currentMessageTextLength: number; assistantMessageId: string | undefined; lastFlushLen: number; /** Track last thinking_content flush length — mirrors lastFlushLen so a @@ -109,6 +111,7 @@ export function createInitialState(): ChatStreamState { lastCapturedMsgId: undefined, lastMsgIdSepEmitted: false, msgIdDeltaStart: 0, + currentMessageTextLength: 0, assistantMessageId: undefined, lastFlushLen: 0, lastThinkingFlushLen: 0,