diff --git a/open-sse/executors/commandcode.js b/open-sse/executors/commandcode.js index d923fd50..c1af7e69 100644 --- a/open-sse/executors/commandcode.js +++ b/open-sse/executors/commandcode.js @@ -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 { diff --git a/tests/unit/commandcode-executor.test.js b/tests/unit/commandcode-executor.test.js index f498b39c..3c97e27a 100644 --- a/tests/unit/commandcode-executor.test.js +++ b/tests/unit/commandcode-executor.test.js @@ -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();