diff --git a/open-sse/executors/commandcode.js b/open-sse/executors/commandcode.js index aad40439..f694e61b 100644 --- a/open-sse/executors/commandcode.js +++ b/open-sse/executors/commandcode.js @@ -42,12 +42,235 @@ export class CommandCodeExecutor extends BaseExecutor { async execute(opts) { const result = await super.execute(opts); if (!result?.response?.ok || !result.response.body) return result; - result.response = wrapNdjsonAsOpenAISse(result.response, opts.model); + result.response = await inspectAndWrapCommandCodeResponse(result.response, opts.model); return result; } + + parseError(response, bodyText) { + let parsed = null; + try { + parsed = JSON.parse(bodyText || "{}"); + } catch { + parsed = null; + } + const errObj = parsed?.error || parsed; + const msg = errObj?.message || parsed?.message || bodyText || response.statusText; + const status = Number(errObj?.code || errObj?.statusCode || response.status) || response.status; + return { + status, + message: msg || `CommandCode upstream error: ${response.status}`, + }; + } } -function wrapNdjsonAsOpenAISse(originalResponse, model) { +export function parseCommandCodeError(event) { + if (!event || typeof event !== "object") { + return { + statusCode: 503, + message: "CommandCode upstream error", + type: "server_error", + }; + } + + const errVal = event.error ?? event.message ?? "unknown"; + let message = ""; + let statusCode = null; + let type = "server_error"; + + if (typeof errVal === "object" && errVal !== null) { + message = errVal.message || errVal.error || JSON.stringify(errVal); + if (errVal.statusCode && Number.isInteger(Number(errVal.statusCode))) { + statusCode = Number(errVal.statusCode); + } else if (errVal.status && Number.isInteger(Number(errVal.status))) { + statusCode = Number(errVal.status); + } + if (errVal.type) type = errVal.type; + } else if (typeof errVal === "string") { + message = errVal; + } else { + message = JSON.stringify(errVal); + } + + if (event.statusCode && Number.isInteger(Number(event.statusCode))) { + statusCode = Number(event.statusCode); + } + + if (!statusCode || statusCode < 400 || statusCode > 599) { + const lower = message.toLowerCase(); + if (lower.includes("rate limit") || lower.includes("too many requests")) { + statusCode = 429; + type = "rate_limit_error"; + } else if (lower.includes("unauthorized") || lower.includes("invalid api key") || lower.includes("authentication")) { + statusCode = 401; + type = "authentication_error"; + } else if (lower.includes("payment required") || lower.includes("billing")) { + statusCode = 402; + type = "billing_error"; + } else if (lower.includes("quota") || lower.includes("forbidden") || lower.includes("permission")) { + statusCode = 403; + type = "permission_error"; + } else if (lower.includes("not found")) { + statusCode = 404; + type = "invalid_request_error"; + } else if (lower.includes("unavailable") || lower.includes("overloaded") || lower.includes("server error")) { + statusCode = 503; + type = "server_error"; + } else { + statusCode = 503; + } + } + + return { statusCode, message, type }; +} + +export async function inspectAndWrapCommandCodeResponse(originalResponse, model) { + const reader = originalResponse.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + const bufferedLines = []; + let detectedError = null; + + try { + while (true) { + const { value, done } = await reader.read(); + if (done) { + const trimmed = buffer.trim(); + if (trimmed) { + try { + const jsonStr = trimmed.startsWith("data:") ? trimmed.slice(5).trim() : trimmed; + const parsed = JSON.parse(jsonStr); + if (parsed?.type === "error") { + detectedError = parsed; + } else { + bufferedLines.push(trimmed); + } + } catch { + bufferedLines.push(trimmed); + } + } + break; + } + + buffer += decoder.decode(value, { stream: true }); + const lines = buffer.split("\n"); + buffer = lines.pop() || ""; + + let stopLoop = false; + for (const line of lines) { + const trimmed = line.trim(); + if (!trimmed) continue; + const jsonStr = trimmed.startsWith("data:") ? trimmed.slice(5).trim() : trimmed; + if (!jsonStr || jsonStr === "[DONE]") { + bufferedLines.push(trimmed); + stopLoop = true; + break; + } + + let event; + try { + event = JSON.parse(jsonStr); + } catch { + bufferedLines.push(trimmed); + continue; + } + + if (event?.type === "error") { + detectedError = event; + stopLoop = true; + break; + } + + bufferedLines.push(trimmed); + + if ( + event?.type === "text-delta" || + event?.type === "reasoning-delta" || + event?.type === "tool-input-start" || + event?.type === "tool-call" || + event?.type === "finish" || + event?.type === "finish-step" + ) { + stopLoop = true; + break; + } + } + + if (stopLoop) break; + } + } catch { + try { reader.releaseLock(); } catch { /* ignore */ } + return originalResponse; + } + + if (detectedError) { + try { await reader.cancel(); } catch { /* ignore */ } + const { statusCode, message, type } = parseCommandCodeError(detectedError); + return new Response( + JSON.stringify({ + error: { + message: `[CommandCode error: ${message}]`, + type, + code: statusCode, + }, + }), + { + status: statusCode, + statusText: statusCode === 503 ? "Service Unavailable" : (statusCode === 429 ? "Too Many Requests" : "Bad Gateway"), + headers: { + "Content-Type": "application/json", + "Access-Control-Allow-Origin": "*", + }, + } + ); + } + + const combinedStream = createReplayedStream(bufferedLines, buffer, reader); + return wrapNdjsonAsOpenAISse(combinedStream, model, originalResponse); +} + +function createReplayedStream(bufferedLines, remainingBuffer, reader) { + const encoder = new TextEncoder(); + let replayed = false; + + 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)); + } + } + + try { + const { value, done } = await reader.read(); + if (done) { + controller.close(); + } else { + controller.enqueue(value); + } + } catch (err) { + controller.error(err); + } + }, + async cancel(reason) { + try { + await reader.cancel(reason); + } catch { + /* ignore */ + } + }, + }); +} + +function wrapNdjsonAsOpenAISse(streamBody, model, originalResponse = null) { const decoder = new TextDecoder(); const encoder = new TextEncoder(); let buffer = ""; @@ -70,7 +293,6 @@ function wrapNdjsonAsOpenAISse(originalResponse, model) { for (const line of lines) { const trimmed = line.trim(); if (!trimmed) continue; - // Translate AI SDK v5 NDJSON line to one or more OpenAI chunks emitChunks(commandCodeToOpenAIResponse(trimmed, state), controller); } }, @@ -83,11 +305,17 @@ function wrapNdjsonAsOpenAISse(originalResponse, model) { }, }); - const newBody = originalResponse.body.pipeThrough(transform); + const newBody = streamBody.pipeThrough(transform); return new Response(newBody, { - status: originalResponse.status, - statusText: originalResponse.statusText, - headers: originalResponse.headers, + status: originalResponse?.status || 200, + statusText: originalResponse?.statusText || "OK", + headers: { + "Content-Type": "text/event-stream", + "Cache-Control": "no-cache", + "Connection": "keep-alive", + ...(originalResponse?.headers ? Object.fromEntries(originalResponse.headers.entries()) : {}), + "content-type": "text/event-stream", + }, }); } diff --git a/tests/unit/commandcode-executor.test.js b/tests/unit/commandcode-executor.test.js new file mode 100644 index 00000000..bd0a23cd --- /dev/null +++ b/tests/unit/commandcode-executor.test.js @@ -0,0 +1,192 @@ +import { describe, it, expect, vi } from "vitest"; +import { + parseCommandCodeError, + inspectAndWrapCommandCodeResponse, + CommandCodeExecutor, +} from "../../open-sse/executors/commandcode.js"; +import { handleComboChat } from "../../open-sse/services/combo.js"; + +function createNdjsonStream(lines) { + const encoder = new TextEncoder(); + return new ReadableStream({ + start(controller) { + for (const line of lines) { + controller.enqueue(encoder.encode(typeof line === "string" ? line : JSON.stringify(line) + "\n")); + } + controller.close(); + }, + }); +} + +describe("parseCommandCodeError", () => { + it("parses user exact error payload with statusCode 503 and isRetryable", () => { + const event = { + type: "error", + error: { + type: "server_error", + message: "Service temporarily unavailable. Please try again shortly.", + statusCode: 503, + isRetryable: true, + }, + }; + const parsed = parseCommandCodeError(event); + expect(parsed.statusCode).toBe(503); + expect(parsed.message).toBe("Service temporarily unavailable. Please try again shortly."); + expect(parsed.type).toBe("server_error"); + }); + + it("handles string error message", () => { + const event = { + type: "error", + message: "Rate limit exceeded. Please wait 30s.", + }; + const parsed = parseCommandCodeError(event); + expect(parsed.statusCode).toBe(429); + expect(parsed.message).toBe("Rate limit exceeded. Please wait 30s."); + }); + + it("handles plain error string in error property", () => { + const event = { + type: "error", + error: "Unauthorized access", + }; + const parsed = parseCommandCodeError(event); + expect(parsed.statusCode).toBe(401); + expect(parsed.message).toBe("Unauthorized access"); + }); +}); + +describe("inspectAndWrapCommandCodeResponse", () => { + it("converts initial upstream 200 with error event to 503 Response", async () => { + const ndjsonBody = createNdjsonStream([ + JSON.stringify({ + type: "error", + error: { + type: "server_error", + message: "Service temporarily unavailable. Please try again shortly.", + statusCode: 503, + isRetryable: true, + }, + }) + "\n", + ]); + + const fakeResponse = new Response(ndjsonBody, { + status: 200, + headers: { "Content-Type": "text/event-stream" }, + }); + + const result = await inspectAndWrapCommandCodeResponse(fakeResponse, "poolside/laguna-s-2.1-free"); + expect(result.ok).toBe(false); + expect(result.status).toBe(503); + + const body = await result.json(); + expect(body.error.message).toContain("Service temporarily unavailable"); + expect(body.error.code).toBe(503); + }); + + it("converts initial upstream 200 with start/start-step followed by error to 503 Response", async () => { + const ndjsonBody = createNdjsonStream([ + JSON.stringify({ type: "start" }) + "\n", + JSON.stringify({ type: "start-step" }) + "\n", + JSON.stringify({ + type: "error", + error: { + type: "server_error", + message: "Service temporarily unavailable. Please try again shortly.", + statusCode: 503, + isRetryable: true, + }, + }) + "\n", + ]); + + const fakeResponse = new Response(ndjsonBody, { + status: 200, + headers: { "Content-Type": "text/event-stream" }, + }); + + const result = await inspectAndWrapCommandCodeResponse(fakeResponse, "poolside/laguna-s-2.1-free"); + expect(result.ok).toBe(false); + expect(result.status).toBe(503); + + const body = await result.json(); + expect(body.error.message).toContain("Service temporarily unavailable"); + }); + + it("streams successful responses when content is emitted", async () => { + const ndjsonBody = createNdjsonStream([ + JSON.stringify({ type: "start" }) + "\n", + JSON.stringify({ type: "text-delta", text: "Hello from Laguna" }) + "\n", + JSON.stringify({ type: "finish" }) + "\n", + ]); + + const fakeResponse = new Response(ndjsonBody, { + status: 200, + headers: { "Content-Type": "text/event-stream" }, + }); + + const result = await inspectAndWrapCommandCodeResponse(fakeResponse, "poolside/laguna-s-2.1-free"); + expect(result.ok).toBe(true); + expect(result.status).toBe(200); + + const text = await result.text(); + expect(text).toContain("Hello from Laguna"); + expect(text).toContain("data: [DONE]"); + }); +}); + +describe("CommandCode in Combo Fallback", () => { + it("automatically falls back to next model when commandcode returns 503 error", async () => { + const log = { + info: vi.fn(), + warn: vi.fn(), + debug: vi.fn(), + }; + + const handleSingleModel = vi.fn(async (body, modelStr) => { + if (modelStr === "commandcode/poolside/laguna-s-2.1-free") { + // Simulated failed CommandCode response + return new Response( + JSON.stringify({ + error: { + message: "Service temporarily unavailable. Please try again shortly.", + type: "server_error", + code: 503, + }, + }), + { status: 503, headers: { "Content-Type": "application/json" } } + ); + } + + if (modelStr === "openai/gpt-4o-mini") { + // Fallback model succeeds + return new Response( + JSON.stringify({ + id: "chatcmpl-test", + choices: [{ message: { role: "assistant", content: "Fallback success!" } }], + }), + { status: 200, headers: { "Content-Type": "application/json" } } + ); + } + + return new Response("Not found", { status: 404 }); + }); + + const comboResponse = await handleComboChat({ + body: { messages: [{ role: "user", content: "Hello" }] }, + models: ["commandcode/poolside/laguna-s-2.1-free", "openai/gpt-4o-mini"], + handleSingleModel, + log, + comboName: "test-combo", + comboStrategy: "fallback", + }); + + expect(comboResponse.ok).toBe(true); + expect(comboResponse.status).toBe(200); + + const data = await comboResponse.json(); + expect(data.choices[0].message.content).toBe("Fallback success!"); + expect(handleSingleModel).toHaveBeenCalledTimes(2); + expect(handleSingleModel).toHaveBeenNthCalledWith(1, expect.anything(), "commandcode/poolside/laguna-s-2.1-free"); + expect(handleSingleModel).toHaveBeenNthCalledWith(2, expect.anything(), "openai/gpt-4o-mini"); + }); +});