fix(commandcode): replay raw byte chunks to preserve all NDJSON lines

This commit is contained in:
Christian Gennari
2026-09-26 11:53:44 +07:00
committed by decolua
parent 5d2cfbf3c5
commit c2148179c0
2 changed files with 33 additions and 25 deletions

View File

@@ -141,7 +141,7 @@ export async function inspectAndWrapCommandCodeResponse(originalResponse, model)
const reader = originalResponse.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
const bufferedLines = [];
const rawChunks = [];
let detectedError = null;
try {
@@ -155,16 +155,15 @@ export async function inspectAndWrapCommandCodeResponse(originalResponse, model)
const parsed = JSON.parse(jsonStr);
if (parsed?.type === "error") {
detectedError = parsed;
} else {
bufferedLines.push(trimmed);
}
} catch {
bufferedLines.push(trimmed);
/* ignore */
}
}
break;
}
rawChunks.push(value);
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split("\n");
buffer = lines.pop() || "";
@@ -175,7 +174,6 @@ export async function inspectAndWrapCommandCodeResponse(originalResponse, model)
if (!trimmed) continue;
const jsonStr = trimmed.startsWith("data:") ? trimmed.slice(5).trim() : trimmed;
if (!jsonStr || jsonStr === "[DONE]") {
bufferedLines.push(trimmed);
stopLoop = true;
break;
}
@@ -184,7 +182,6 @@ export async function inspectAndWrapCommandCodeResponse(originalResponse, model)
try {
event = JSON.parse(jsonStr);
} catch {
bufferedLines.push(trimmed);
continue;
}
@@ -194,8 +191,6 @@ export async function inspectAndWrapCommandCodeResponse(originalResponse, model)
break;
}
bufferedLines.push(trimmed);
if (
event?.type === "text-delta" ||
event?.type === "reasoning-delta" ||
@@ -238,29 +233,18 @@ export async function inspectAndWrapCommandCodeResponse(originalResponse, model)
);
}
const combinedStream = createReplayedStream(bufferedLines, buffer, reader);
const combinedStream = createRawReplayedStream(rawChunks, reader);
return wrapNdjsonAsOpenAISse(combinedStream, model, originalResponse);
}
function createReplayedStream(bufferedLines, remainingBuffer, reader) {
const encoder = new TextEncoder();
let replayed = false;
function createRawReplayedStream(rawChunks, reader) {
let chunkIndex = 0;
return new ReadableStream({
async pull(controller) {
if (!replayed) {
replayed = true;
let prefix = bufferedLines.join("\n");
if (prefix && remainingBuffer) {
prefix += "\n" + remainingBuffer;
} else if (remainingBuffer) {
prefix = remainingBuffer;
} else if (prefix) {
prefix += "\n";
}
if (prefix) {
controller.enqueue(encoder.encode(prefix));
}
if (chunkIndex < rawChunks.length) {
controller.enqueue(rawChunks[chunkIndex++]);
return;
}
try {

View File

@@ -133,6 +133,30 @@ describe("inspectAndWrapCommandCodeResponse", () => {
expect(text).toContain("data: [DONE]");
});
it("preserves all lines in a multi-line packet when inspecting tool-input-start", async () => {
const packet = [
JSON.stringify({ type: "start" }),
JSON.stringify({ type: "start-step" }),
JSON.stringify({ type: "tool-input-start", id: "call_1", toolName: "terminal" }),
JSON.stringify({ type: "tool-input-delta", id: "call_1", delta: '{"command": "ls"}' }),
JSON.stringify({ type: "finish-step", finishReason: "tool-calls" }),
JSON.stringify({ type: "finish", finishReason: "tool-calls" }),
].join("\n") + "\n";
const ndjsonBody = createNdjsonStream([packet]);
const fakeResponse = new Response(ndjsonBody, {
status: 200,
headers: { "Content-Type": "text/event-stream" },
});
const result = await inspectAndWrapCommandCodeResponse(fakeResponse, "cmc/deepseek/deepseek-v4.1-flash");
expect(result.ok).toBe(true);
const text = await result.text();
expect(text).toContain('"name":"terminal"');
expect(text).toContain('"arguments":"{\\"command\\": \\"ls\\"}"');
});
it("retries when initial stream yields an error and succeeds on second attempt", async () => {
let callCount = 0;
const executor = new CommandCodeExecutor();