diff --git a/CHANGELOG.md b/CHANGELOG.md index c3ddfcb2..9962eab9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,6 +3,10 @@ ## Features - **Qoder CN**: add the `qoder-cn` provider for qoder.com.cn — same device/OAuth flow and COSY-signed chat protocol as intl Qoder, served from the CN gateway (`gateway.qoder.com.cn` / `openapi.qoder.com.cn`). Region resolution is derived from the provider id (`qoder-cn` → CN) across executor, model catalog, quota usage, validation, and the dashboard. +## Fixes +- **Qoder**: prevent proxy failures from silently replaying signed inference requests over a direct connection with the same COSY request ID (`403/103 Duplicate request`). +- **Qoder**: return first-frame upstream errors, including `403/103 Duplicate request`, as HTTP failures instead of assistant text; handle fragmented frames and heartbeat prefixes while preserving billing error mapping. + # v0.5.81 (2026-09-18) ## Features diff --git a/open-sse/executors/qoder.js b/open-sse/executors/qoder.js index da6844f4..f9e7d769 100644 --- a/open-sse/executors/qoder.js +++ b/open-sse/executors/qoder.js @@ -29,7 +29,7 @@ import { BaseExecutor } from "./base.js"; import { PROVIDERS } from "../config/providers.js"; import { proxyAwareFetch } from "../utils/proxyFetch.js"; import { SSE_DONE } from "../utils/sseConstants.js"; -import { FETCH_CONNECT_TIMEOUT_MS } from "../config/runtimeConfig.js"; +import { FETCH_CONNECT_TIMEOUT_MS, HTTP_STATUS } from "../config/runtimeConfig.js"; import { QODER_CHAT_SIG_PATH, QODER_CONTEXT_TIER_ENV, @@ -354,38 +354,47 @@ function isBillingBlock(inner) { } /** - * Peek the first SSE frame to detect billing errors before piping. - * Returns { isBilling, statusVal, message, consumed } — `consumed` is every + * Peek the first SSE data line to detect upstream errors before piping. + * Returns { isError, isBilling, statusVal, message, consumed } — `consumed` is every * byte read so far (including the peeked line) so the caller can re-process * it and nothing is dropped from the stream. */ async function peekFirstQoderFrame(reader, decoder) { let consumed = ""; + let offset = 0; + let upstreamDone = false; while (true) { - const { done, value } = await reader.read(); - if (done) return { isBilling: false, consumed, upstreamDone: true }; + let nl = consumed.indexOf("\n", offset); + if (nl === -1 && !upstreamDone) { + const { done, value } = await reader.read(); + upstreamDone = done; + consumed += done ? decoder.decode() : decoder.decode(value, { stream: true }); + continue; + } + if (offset >= consumed.length) return { isError: false, consumed, upstreamDone }; + if (nl === -1) nl = consumed.length; - consumed += decoder.decode(value, { stream: true }); - const nl = consumed.indexOf("\n"); - if (nl === -1) continue; // need a full line first - - const line = consumed.slice(0, nl).replace(/\r$/, "").trim(); + const line = consumed.slice(offset, nl).replace(/\r$/, "").trim(); + offset = nl + 1; if (!line.startsWith("data:")) continue; const data = line.slice(5).trimStart(); - if (data === "[DONE]") return { isBilling: false, consumed }; + if (data === "[DONE]") return { isError: false, consumed, upstreamDone }; let envelope; - try { envelope = JSON.parse(data); } catch { return { isBilling: false, consumed }; } + try { envelope = JSON.parse(data); } catch { return { isError: false, consumed, upstreamDone }; } + // statusCodeValue is documented numeric, but accept numeric strings defensively. - const statusVal = Number(envelope.statusCodeValue) || 200; - const inner = typeof envelope.body === "string" + const raw = Number(envelope?.statusCodeValue); + const statusVal = Number.isNaN(raw) ? 200 : raw; + const inner = typeof envelope?.body === "string" ? envelope.body - : envelope.body != null ? JSON.stringify(envelope.body) : ""; - if (statusVal !== 200 && isBillingBlock(inner)) { - return { isBilling: true, statusVal, message: inner || `qoder billing block (${statusVal})` }; + : envelope?.body != null ? JSON.stringify(envelope.body) : ""; + + if (statusVal !== 200) { + return { isError: true, isBilling: isBillingBlock(inner), statusVal, message: inner || `upstream status ${statusVal}` }; } - return { isBilling: false, consumed }; + return { isError: false, consumed, upstreamDone }; } } @@ -396,8 +405,8 @@ async function peekFirstQoderFrame(reader, decoder) { * Each upstream line looks like: * data: {"statusCodeValue":200,"body":"{\"choices\":[{\"delta\":{...}}]}"} * The inner body is an OpenAI streaming chunk (or "[DONE]"). We unwrap it - * and re-emit as `data: \n\n`. Errors become a synthetic OpenAI error - * chunk + [DONE]. + * and re-emit as `data: \n\n`. First-frame errors become HTTP errors; + * errors after streaming starts retain the synthetic chunk + [DONE] path. * * Critical: Qoder's SSE often keeps the socket open after the terminal * [DONE]/error frame (agent keepalive). Non-streaming clients drain via @@ -409,9 +418,10 @@ async function peekFirstQoderFrame(reader, decoder) { * usage from the finish chunk, so we coalesce those two frames (see * createQoderSseCoalescer) before forwarding. * - * NEW: Peek first frame to detect billing blocks (code 112/10605/pricingUrl). - * If detected, return 403 response so chatCore marks connection unavailable - * and triggers combo fallback instead of leaking error text into chat. + * Peek the first frame for errors before committing to HTTP 200. Preserve + * upstream error statuses so chatCore can handle failures instead of recording + * error text as a successful completion. Billing blocks retain the existing + * 403 mapping for quota/account fallback. */ async function wrapQoderSSE(response, model, log = null) { if (!response.ok || !response.body) return response; @@ -419,14 +429,17 @@ async function wrapQoderSSE(response, model, log = null) { const decoder = new TextDecoder(); const reader = response.body.getReader(); - // Peek first frame to detect billing block + // Detect errors before returning a successful streaming response. const peek = await peekFirstQoderFrame(reader, decoder); - if (peek?.isBilling) { - // Billing block detected — return 403 so chatCore fails this connection + if (peek.isError) { await reader.cancel().catch(() => {}); + const status = peek.isBilling + ? HTTP_STATUS.FORBIDDEN + : Number.isInteger(peek.statusVal) && peek.statusVal >= HTTP_STATUS.BAD_REQUEST && peek.statusVal <= 599 + ? peek.statusVal : HTTP_STATUS.BAD_GATEWAY; return new Response( JSON.stringify({ error: { message: peek.message, code: peek.statusVal } }), - { status: 403, headers: { "Content-Type": "application/json" } } + { status, headers: { "Content-Type": "application/json" } } ); } @@ -699,8 +712,15 @@ export class QoderExecutor extends BaseExecutor { response = await proxyAwareFetch( url, { method: "POST", headers, body: encodedBodyBuf, signal: mergedSignal }, - proxyOptions, + // A failed proxy request may already have reached Qoder. Replaying + // the same COSY signature directly reuses its requestId and returns + // 403/code 103. Let the caller retry through execute() with fresh signing. + { ...proxyOptions, strictProxy: true }, ); + } catch (err) { + // strictProxy wraps transport errors; retain caller cancellation semantics. + if (mergedSignal.aborted) throw mergedSignal.reason; + throw err; } finally { clearTimeout(connectTimer); } diff --git a/tests/unit/qoder-billing.test.js b/tests/unit/qoder-billing.test.js index d7c5482e..4900f2df 100644 --- a/tests/unit/qoder-billing.test.js +++ b/tests/unit/qoder-billing.test.js @@ -236,7 +236,7 @@ describe("wrapQoderSSE billing detection", () => { expect(wrapped.status).toBe(403); }); - it("passes through normal errors (non-billing) as wrapped SSE", async () => { + it("returns non-billing errors with their upstream HTTP status", async () => { const errorEnv = JSON.stringify({ statusCodeValue: 500, body: "Internal server error", @@ -245,22 +245,11 @@ describe("wrapQoderSSE billing detection", () => { const wrapped = await wrapQoderSSE(makeResponse([upstream]), "qoder/ultimate"); - // Normal error: still 200 response, error text in SSE body - expect(wrapped.status).toBe(200); - expect(wrapped.ok).toBe(true); - - const reader = wrapped.body.getReader(); - const decoder = new TextDecoder(); - let buf = ""; - while (true) { - const { done, value } = await reader.read(); - if (done) break; - buf += decoder.decode(value, { stream: true }); - } - buf += decoder.decode(); - - expect(buf).toContain("[qoder error 500"); - expect(buf).toContain("data: [DONE]"); + expect(wrapped.status).toBe(500); + expect(wrapped.ok).toBe(false); + expect(await wrapped.json()).toEqual({ + error: { message: "Internal server error", code: 500 }, + }); }); it("passes through successful responses unchanged", async () => { diff --git a/tests/unit/qoder-proxy-replay.test.js b/tests/unit/qoder-proxy-replay.test.js new file mode 100644 index 00000000..1a3ff0c6 --- /dev/null +++ b/tests/unit/qoder-proxy-replay.test.js @@ -0,0 +1,97 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; + +vi.mock("../../open-sse/services/qoderModels.js", () => ({ + getQoderModelConfig: vi.fn(async () => ({ key: "auto", max_output_tokens: 32 })), + resolveQoderModels: vi.fn(), + isQoderPat: () => false, + resolveQoderCredentials: vi.fn(), +})); + +const request = { + model: "auto", + body: { messages: [{ role: "user", content: "hello" }], max_tokens: 32 }, + stream: true, + credentials: { + accessToken: "dt-test-token", + providerSpecificData: { userId: "test-user", machineId: "test-machine" }, + }, +}; + +function success() { + return new Response('data: {"statusCodeValue":200,"body":"[DONE]"}\n\n', { + headers: { "Content-Type": "text/event-stream" }, + }); +} + +async function loadExecutor(fetchMock, useProxy = true) { + vi.resetModules(); + for (const key of ["HTTP_PROXY", "HTTPS_PROXY", "ALL_PROXY", "NO_PROXY", "http_proxy", "https_proxy", "all_proxy", "no_proxy"]) { + vi.stubEnv(key, ""); + } + if (useProxy) vi.stubEnv("HTTPS_PROXY", "http://proxy.test:3128"); + // Exercise the real proxyAwareFetch: it captures fetch when imported. + vi.stubGlobal("fetch", fetchMock); + const { QoderExecutor } = await import("../../open-sse/executors/qoder.js"); + return new QoderExecutor(); +} + +afterEach(() => { + vi.unstubAllGlobals(); + vi.unstubAllEnvs(); +}); + +describe("Qoder signed inference transport", () => { + it.each([null, { strictProxy: false }])("does not replay a signed POST after proxy response loss (%j)", async (proxyOptions) => { + const seen = new Set(); + const fetchMock = vi.fn(async (_url, options) => { + const authorization = options.headers.Authorization; + if (seen.has(authorization)) { + return new Response('data: {"statusCodeValue":403,"body":"{\\"code\\":\\"103\\",\\"message\\":\\"Duplicate request\\"}"}\n\n'); + } + seen.add(authorization); + throw new TypeError("response lost after upstream accepted request"); + }); + const executor = await loadExecutor(fetchMock); + await expect(executor.execute({ ...request, proxyOptions })).rejects.toThrow("response lost"); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(fetchMock.mock.calls[0][1].dispatcher).toBeDefined(); + if (proxyOptions) expect(proxyOptions.strictProxy).toBe(false); + }); + + it("generates a fresh COSY identity when the caller retries after transport failure", async () => { + const fetchMock = vi.fn() + .mockRejectedValueOnce(new TypeError("response lost")) + .mockResolvedValueOnce(success()); + const executor = await loadExecutor(fetchMock); + await expect(executor.execute(request)).rejects.toThrow("response lost"); + const result = await executor.execute(request); + expect(result.response.ok).toBe(true); + await result.response.text(); + expect(fetchMock).toHaveBeenCalledTimes(2); + const ids = fetchMock.mock.calls.map(([, options]) => JSON.parse( + Buffer.from(options.headers.Authorization.split(".")[1], "base64").toString(), + ).requestId); + expect(ids[0]).not.toBe(ids[1]); + }); + + it.each([true, false])("still supports successful inference with proxy=%s", async (useProxy) => { + const fetchMock = vi.fn(async () => success()); + const executor = await loadExecutor(fetchMock, useProxy); + const result = await executor.execute(request); + expect(result.response.ok).toBe(true); + await result.response.text(); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(!!fetchMock.mock.calls[0][1].dispatcher).toBe(useProxy); + }); + + it("preserves caller cancellation without replaying the request", async () => { + const controller = new AbortController(); + const fetchMock = vi.fn(async (_url, options) => { + controller.abort(); + throw options.signal.reason; + }); + const executor = await loadExecutor(fetchMock); + await expect(executor.execute({ ...request, signal: controller.signal })).rejects.toMatchObject({ name: "AbortError" }); + expect(fetchMock).toHaveBeenCalledTimes(1); + }); +}); diff --git a/tests/unit/qoder-stream-errors.test.js b/tests/unit/qoder-stream-errors.test.js new file mode 100644 index 00000000..60b316b0 --- /dev/null +++ b/tests/unit/qoder-stream-errors.test.js @@ -0,0 +1,81 @@ +import { describe, it, expect, vi } from "vitest"; +import { __test__ } from "../../open-sse/executors/qoder.js"; + +const { wrapQoderSSE } = __test__; +const duplicate = '{"code":"103","message":"Duplicate request"}'; +const frame = (statusCodeValue, body) => `data: ${JSON.stringify({ statusCodeValue, body })}\n\n`; + +function upstream(chunks, { keepOpen = false } = {}) { + const cancel = vi.fn(); + const response = new Response(new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(new TextEncoder().encode(chunk)); + if (!keepOpen) controller.close(); + }, + cancel, + })); + return { response, cancel }; +} + +describe("Qoder first-frame errors", () => { + it.each([ + ["one chunk", [frame(403, duplicate)]], + ["fragmented frame", [frame(403, duplicate).slice(0, 35), frame(403, duplicate).slice(35)]], + ["heartbeat prefix", [": keepalive\r\n\r\n", frame(403, duplicate)]], + ["prefix and frame in one chunk", [": keepalive\n\nevent: message\n" + frame(403, duplicate)]], + ["EOF without newline", [frame(403, duplicate).trimEnd()]], + ["object body", [frame(403, JSON.parse(duplicate))]], + ])("surfaces duplicate-request errors as HTTP 403: %s", async (_name, chunks) => { + const { response } = upstream(chunks); + const wrapped = await wrapQoderSSE(response, "qoder/kmodel_latest"); + expect(wrapped.status).toBe(403); + expect(wrapped.ok).toBe(false); + expect(wrapped.headers.get("content-type")).toBe("application/json"); + const body = await wrapped.json(); + expect(body.error.message).toBe(duplicate); + expect(body).not.toHaveProperty("choices"); + }); + + it("cancels the upstream keepalive immediately after an error", async () => { + const { response, cancel } = upstream([": keepalive\n\n", frame(403, duplicate)], { keepOpen: true }); + const wrapped = await wrapQoderSSE(response, "qoder/kmodel_latest"); + expect(wrapped.status).toBe(403); + expect(cancel).toHaveBeenCalledOnce(); + }); + + it.each([401, 429, 500, 503])("preserves non-billing HTTP status %s", async (status) => { + const { response } = upstream([frame(status, "upstream failure")]); + const wrapped = await wrapQoderSSE(response, "qoder/auto"); + expect(wrapped.status).toBe(status); + expect((await wrapped.json()).error.message).toBe("upstream failure"); + }); + + it.each([0, 302, 600, 403.5])("maps invalid error status %s to 502", async (status) => { + const { response } = upstream([frame(status, "invalid upstream status")]); + const wrapped = await wrapQoderSSE(response, "qoder/auto"); + expect(wrapped.status).toBe(502); + }); + + it("replays successful frames after a heartbeat without losing or duplicating content", async () => { + const first = JSON.stringify({ choices: [{ delta: { content: "hello" } }] }); + const second = JSON.stringify({ choices: [{ delta: { content: "world" } }] }); + const { response } = upstream([": keepalive\n\n", frame(200, first) + frame(200, second) + "data: [DONE]\n\n"]); + const wrapped = await wrapQoderSSE(response, "qoder/auto"); + expect(wrapped.status).toBe(200); + expect(await wrapped.text()).toBe(`data: ${first}\n\ndata: ${second}\n\ndata: [DONE]\n\n`); + }); + + it("starts forwarding success without waiting for the upstream to close", async () => { + const inner = JSON.stringify({ choices: [{ delta: { content: "hello" } }] }); + const { response, cancel } = upstream([": keepalive\n\n", frame(200, inner)], { keepOpen: true }); + const wrapped = await wrapQoderSSE(response, "qoder/auto"); + const reader = wrapped.body.getReader(); + try { + const { value } = await reader.read(); + expect(new TextDecoder().decode(value)).toBe(`data: ${inner}\n\n`); + } finally { + await reader.cancel(); + } + expect(cancel).toHaveBeenCalledOnce(); + }); +}); diff --git a/tests/unit/qoder.test.js b/tests/unit/qoder.test.js index 5f8246e8..ce7edd67 100644 --- a/tests/unit/qoder.test.js +++ b/tests/unit/qoder.test.js @@ -506,14 +506,14 @@ describe("wrapQoderSSE", () => { // Regression for review finding #3: chunks could leak past [DONE] when // the success branch had no doneEmitted guard. We synthesize an error - // envelope (which sets doneEmitted=true) followed by a valid envelope + // envelope after content (which sets doneEmitted=true), followed by a valid envelope // and assert the second envelope is NOT forwarded. it("does not forward chunks after [DONE] has been emitted", async () => { const errorEnv = JSON.stringify({ statusCodeValue: 500, body: "boom" }); const validInner = JSON.stringify({ choices: [{ delta: { content: "leak" } }] }); const validEnv = JSON.stringify({ statusCodeValue: 200, body: validInner }); const wrapped = await wrapQoderSSE( - makeResponse([`data: ${errorEnv}\n\ndata: ${validEnv}\n\n`]), + makeResponse([envelope(JSON.stringify({ choices: [{ delta: { content: "hi" } }] })) + `data: ${errorEnv}\n\ndata: ${validEnv}\n\n`]), "qoder/auto", ); const out = await drain(wrapped); @@ -539,12 +539,13 @@ describe("wrapQoderSSE", () => { expect(() => JSON.parse(dataLine.slice("data: ".length))).not.toThrow(); }); - it("upstream error envelope produces an error chunk + [DONE]", async () => { + it("upstream first-frame error envelope produces an HTTP error", async () => { const env = JSON.stringify({ statusCodeValue: 503, body: "service unavailable" }); const wrapped = await wrapQoderSSE(makeResponse([`data: ${env}\n\n`]), "qoder/lite"); - const out = await drain(wrapped); - expect(out).toContain("[qoder error 503"); - expect(out).toContain("data: [DONE]\n\n"); + expect(wrapped.status).toBe(503); + expect(await wrapped.json()).toEqual({ + error: { message: "service unavailable", code: 503 }, + }); }); it("non-ok responses are returned unchanged (no transform)", async () => {