fix(qoder): prevent signed request replay and surface upstream errors
This commit is contained in:
@@ -3,6 +3,10 @@
|
|||||||
## Features
|
## 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.
|
- **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)
|
# v0.5.81 (2026-09-18)
|
||||||
|
|
||||||
## Features
|
## Features
|
||||||
|
|||||||
@@ -29,7 +29,7 @@ import { BaseExecutor } from "./base.js";
|
|||||||
import { PROVIDERS } from "../config/providers.js";
|
import { PROVIDERS } from "../config/providers.js";
|
||||||
import { proxyAwareFetch } from "../utils/proxyFetch.js";
|
import { proxyAwareFetch } from "../utils/proxyFetch.js";
|
||||||
import { SSE_DONE } from "../utils/sseConstants.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 {
|
import {
|
||||||
QODER_CHAT_SIG_PATH,
|
QODER_CHAT_SIG_PATH,
|
||||||
QODER_CONTEXT_TIER_ENV,
|
QODER_CONTEXT_TIER_ENV,
|
||||||
@@ -354,38 +354,47 @@ function isBillingBlock(inner) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Peek the first SSE frame to detect billing errors before piping.
|
* Peek the first SSE data line to detect upstream errors before piping.
|
||||||
* Returns { isBilling, statusVal, message, consumed } — `consumed` is every
|
* Returns { isError, isBilling, statusVal, message, consumed } — `consumed` is every
|
||||||
* byte read so far (including the peeked line) so the caller can re-process
|
* byte read so far (including the peeked line) so the caller can re-process
|
||||||
* it and nothing is dropped from the stream.
|
* it and nothing is dropped from the stream.
|
||||||
*/
|
*/
|
||||||
async function peekFirstQoderFrame(reader, decoder) {
|
async function peekFirstQoderFrame(reader, decoder) {
|
||||||
let consumed = "";
|
let consumed = "";
|
||||||
|
let offset = 0;
|
||||||
|
let upstreamDone = false;
|
||||||
while (true) {
|
while (true) {
|
||||||
|
let nl = consumed.indexOf("\n", offset);
|
||||||
|
if (nl === -1 && !upstreamDone) {
|
||||||
const { done, value } = await reader.read();
|
const { done, value } = await reader.read();
|
||||||
if (done) return { isBilling: false, consumed, upstreamDone: true };
|
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 line = consumed.slice(offset, nl).replace(/\r$/, "").trim();
|
||||||
const nl = consumed.indexOf("\n");
|
offset = nl + 1;
|
||||||
if (nl === -1) continue; // need a full line first
|
|
||||||
|
|
||||||
const line = consumed.slice(0, nl).replace(/\r$/, "").trim();
|
|
||||||
if (!line.startsWith("data:")) continue;
|
if (!line.startsWith("data:")) continue;
|
||||||
|
|
||||||
const data = line.slice(5).trimStart();
|
const data = line.slice(5).trimStart();
|
||||||
if (data === "[DONE]") return { isBilling: false, consumed };
|
if (data === "[DONE]") return { isError: false, consumed, upstreamDone };
|
||||||
|
|
||||||
let envelope;
|
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.
|
// statusCodeValue is documented numeric, but accept numeric strings defensively.
|
||||||
const statusVal = Number(envelope.statusCodeValue) || 200;
|
const raw = Number(envelope?.statusCodeValue);
|
||||||
const inner = typeof envelope.body === "string"
|
const statusVal = Number.isNaN(raw) ? 200 : raw;
|
||||||
|
const inner = typeof envelope?.body === "string"
|
||||||
? envelope.body
|
? envelope.body
|
||||||
: envelope.body != null ? JSON.stringify(envelope.body) : "";
|
: envelope?.body != null ? JSON.stringify(envelope.body) : "";
|
||||||
if (statusVal !== 200 && isBillingBlock(inner)) {
|
|
||||||
return { isBilling: true, statusVal, message: inner || `qoder billing block (${statusVal})` };
|
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:
|
* Each upstream line looks like:
|
||||||
* data: {"statusCodeValue":200,"body":"{\"choices\":[{\"delta\":{...}}]}"}
|
* data: {"statusCodeValue":200,"body":"{\"choices\":[{\"delta\":{...}}]}"}
|
||||||
* The inner body is an OpenAI streaming chunk (or "[DONE]"). We unwrap it
|
* The inner body is an OpenAI streaming chunk (or "[DONE]"). We unwrap it
|
||||||
* and re-emit as `data: <inner>\n\n`. Errors become a synthetic OpenAI error
|
* and re-emit as `data: <inner>\n\n`. First-frame errors become HTTP errors;
|
||||||
* chunk + [DONE].
|
* errors after streaming starts retain the synthetic chunk + [DONE] path.
|
||||||
*
|
*
|
||||||
* Critical: Qoder's SSE often keeps the socket open after the terminal
|
* Critical: Qoder's SSE often keeps the socket open after the terminal
|
||||||
* [DONE]/error frame (agent keepalive). Non-streaming clients drain via
|
* [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
|
* usage from the finish chunk, so we coalesce those two frames (see
|
||||||
* createQoderSseCoalescer) before forwarding.
|
* createQoderSseCoalescer) before forwarding.
|
||||||
*
|
*
|
||||||
* NEW: Peek first frame to detect billing blocks (code 112/10605/pricingUrl).
|
* Peek the first frame for errors before committing to HTTP 200. Preserve
|
||||||
* If detected, return 403 response so chatCore marks connection unavailable
|
* upstream error statuses so chatCore can handle failures instead of recording
|
||||||
* and triggers combo fallback instead of leaking error text into chat.
|
* error text as a successful completion. Billing blocks retain the existing
|
||||||
|
* 403 mapping for quota/account fallback.
|
||||||
*/
|
*/
|
||||||
async function wrapQoderSSE(response, model, log = null) {
|
async function wrapQoderSSE(response, model, log = null) {
|
||||||
if (!response.ok || !response.body) return response;
|
if (!response.ok || !response.body) return response;
|
||||||
@@ -419,14 +429,17 @@ async function wrapQoderSSE(response, model, log = null) {
|
|||||||
const decoder = new TextDecoder();
|
const decoder = new TextDecoder();
|
||||||
const reader = response.body.getReader();
|
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);
|
const peek = await peekFirstQoderFrame(reader, decoder);
|
||||||
if (peek?.isBilling) {
|
if (peek.isError) {
|
||||||
// Billing block detected — return 403 so chatCore fails this connection
|
|
||||||
await reader.cancel().catch(() => {});
|
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(
|
return new Response(
|
||||||
JSON.stringify({ error: { message: peek.message, code: peek.statusVal } }),
|
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(
|
response = await proxyAwareFetch(
|
||||||
url,
|
url,
|
||||||
{ method: "POST", headers, body: encodedBodyBuf, signal: mergedSignal },
|
{ 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 {
|
} finally {
|
||||||
clearTimeout(connectTimer);
|
clearTimeout(connectTimer);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -236,7 +236,7 @@ describe("wrapQoderSSE billing detection", () => {
|
|||||||
expect(wrapped.status).toBe(403);
|
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({
|
const errorEnv = JSON.stringify({
|
||||||
statusCodeValue: 500,
|
statusCodeValue: 500,
|
||||||
body: "Internal server error",
|
body: "Internal server error",
|
||||||
@@ -245,22 +245,11 @@ describe("wrapQoderSSE billing detection", () => {
|
|||||||
|
|
||||||
const wrapped = await wrapQoderSSE(makeResponse([upstream]), "qoder/ultimate");
|
const wrapped = await wrapQoderSSE(makeResponse([upstream]), "qoder/ultimate");
|
||||||
|
|
||||||
// Normal error: still 200 response, error text in SSE body
|
expect(wrapped.status).toBe(500);
|
||||||
expect(wrapped.status).toBe(200);
|
expect(wrapped.ok).toBe(false);
|
||||||
expect(wrapped.ok).toBe(true);
|
expect(await wrapped.json()).toEqual({
|
||||||
|
error: { message: "Internal server error", code: 500 },
|
||||||
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]");
|
|
||||||
});
|
});
|
||||||
|
|
||||||
it("passes through successful responses unchanged", async () => {
|
it("passes through successful responses unchanged", async () => {
|
||||||
|
|||||||
97
tests/unit/qoder-proxy-replay.test.js
Normal file
97
tests/unit/qoder-proxy-replay.test.js
Normal file
@@ -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);
|
||||||
|
});
|
||||||
|
});
|
||||||
81
tests/unit/qoder-stream-errors.test.js
Normal file
81
tests/unit/qoder-stream-errors.test.js
Normal file
@@ -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();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -506,14 +506,14 @@ describe("wrapQoderSSE", () => {
|
|||||||
|
|
||||||
// Regression for review finding #3: chunks could leak past [DONE] when
|
// Regression for review finding #3: chunks could leak past [DONE] when
|
||||||
// the success branch had no doneEmitted guard. We synthesize an error
|
// 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.
|
// and assert the second envelope is NOT forwarded.
|
||||||
it("does not forward chunks after [DONE] has been emitted", async () => {
|
it("does not forward chunks after [DONE] has been emitted", async () => {
|
||||||
const errorEnv = JSON.stringify({ statusCodeValue: 500, body: "boom" });
|
const errorEnv = JSON.stringify({ statusCodeValue: 500, body: "boom" });
|
||||||
const validInner = JSON.stringify({ choices: [{ delta: { content: "leak" } }] });
|
const validInner = JSON.stringify({ choices: [{ delta: { content: "leak" } }] });
|
||||||
const validEnv = JSON.stringify({ statusCodeValue: 200, body: validInner });
|
const validEnv = JSON.stringify({ statusCodeValue: 200, body: validInner });
|
||||||
const wrapped = await wrapQoderSSE(
|
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",
|
"qoder/auto",
|
||||||
);
|
);
|
||||||
const out = await drain(wrapped);
|
const out = await drain(wrapped);
|
||||||
@@ -539,12 +539,13 @@ describe("wrapQoderSSE", () => {
|
|||||||
expect(() => JSON.parse(dataLine.slice("data: ".length))).not.toThrow();
|
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 env = JSON.stringify({ statusCodeValue: 503, body: "service unavailable" });
|
||||||
const wrapped = await wrapQoderSSE(makeResponse([`data: ${env}\n\n`]), "qoder/lite");
|
const wrapped = await wrapQoderSSE(makeResponse([`data: ${env}\n\n`]), "qoder/lite");
|
||||||
const out = await drain(wrapped);
|
expect(wrapped.status).toBe(503);
|
||||||
expect(out).toContain("[qoder error 503");
|
expect(await wrapped.json()).toEqual({
|
||||||
expect(out).toContain("data: [DONE]\n\n");
|
error: { message: "service unavailable", code: 503 },
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
it("non-ok responses are returned unchanged (no transform)", async () => {
|
it("non-ok responses are returned unchanged (no transform)", async () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user