Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions services/api/src/ai/sse-mapper.ts
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,7 @@ export class ChannelTokenRouter {
}


export function mapEventToSSE(event: Record<string, unknown>): SSEEvent | null {
export function mapEventToSSE(event: Record<string, unknown>, options: { preserveThinkingMarkers?: boolean } = {}): SSEEvent | null {
const type = event.type as string;
const data = event.data as Record<string, unknown> | undefined;

Expand All @@ -182,14 +182,14 @@ export function mapEventToSSE(event: Record<string, unknown>): 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) ─
Expand All @@ -199,7 +199,7 @@ export function mapEventToSSE(event: Record<string, unknown>): 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) ──
Expand Down
12 changes: 7 additions & 5 deletions services/api/src/ai/tool-messages.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -303,10 +303,12 @@ export function sanitizeText(text: string): string {
// <|channel>thought, <channel>, <|channel|> — Gemma 4
// <rationale>, </rationale> — Claude (when prompted)
// <answer>, </answer> — 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);
Expand Down
59 changes: 14 additions & 45 deletions services/api/src/routes/chat/event-processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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";
Expand Down Expand Up @@ -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(
Expand Down
150 changes: 150 additions & 0 deletions services/api/src/routes/chat/final-response.integration.test.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown>) =>
process({ type, data } as Parameters<typeof process>[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: "<think>Internal reasoning.</think>Final answer.",
});
emit("assistant.message", {
messageId: "m1",
content: "<think>Internal reasoning.</think>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 ["<thi", "nk>Reasoning", "</thi", "nk>Answer."]) {
emit("assistant.message_delta", { messageId: "final", deltaContent });
}
emit("assistant.message", {
messageId: "final",
content: "<think>Reasoning</think>Answer.",
});
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.",
);
});
73 changes: 73 additions & 0 deletions services/api/src/routes/chat/final-response.test.ts
Original file line number Diff line number Diff line change
@@ -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, "");
});
Loading