Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110import assert from "node:assert/strict";import { test } from "node:test";import { MockLanguageModelV3 } from "ai/test";import { streamText, stepCountIs, tool, jsonSchema } from "ai";import { Secret } from "../configuration/secrets.ts";import { MODEL_PROVIDERS, MODEL_EFFORTS, EFFORT_LABELS, modelLabel,} from "../shared/model-providers.ts";import { createConfiguredModel, modelProviderOptions, protectModel,} from "../worker/model-provider.ts";import { DEFAULT_MODEL, parseModelConfiguration, parseProviderKey, encryptProviderKey, decryptProviderKey,} from "../worker/model-settings.ts";
const credential = "sk-ant-test-secret-that-must-never-escape";const secret = new Secret("session-signing-key-with-at-least-32-bytes");const installation = "a".repeat(32);const anthropic = { provider: "anthropic", model: "claude-haiku-4-5-20251001" };const params = { prompt: [{ role: "user", content: [{ type: "text", text: "Hi" }] }], maxOutputTokens: 10,};
const openrouter = { provider: "openrouter", model: "openai/gpt-5.6-luna",};
test("effort validation is exhaustive per model and never accepts budgets or toggles", () => { for (const [provider, definition] of Object.entries(MODEL_PROVIDERS)) { assert.deepEqual( Object.keys(MODEL_EFFORTS[provider]).sort(), Object.keys(definition.models).sort(), ); for (const model of Object.keys(definition.models)) { const configuration = { provider, model }; assert.deepEqual(parseModelConfiguration(configuration), configuration); assert.equal(modelProviderOptions(configuration), undefined); for (const effort of [ ...Object.keys(EFFORT_LABELS), "default", "off", "on", "none", "adaptive", 2048, null, undefined, ]) { if (MODEL_EFFORTS[provider][model].includes(effort)) assert.deepEqual( parseModelConfiguration({ ...configuration, effort }), { ...configuration, effort }, ); else assert.throws( () => parseModelConfiguration({ ...configuration, effort }), /supported effort/, ); } for (const extra of [ { budget: 2048 }, { thinking: true }, { providerOptions: {} }, ]) assert.throws( () => parseModelConfiguration({ ...configuration, ...extra }), /supported provider/, ); } } assert.deepEqual(MODEL_EFFORTS.openrouter["tencent/hy4-preview"], [ "low", "high", ]); assert.deepEqual(MODEL_EFFORTS.openrouter["deepseek/deepseek-v4-flash"], [ "high", "xhigh", ]); assert.deepEqual(MODEL_EFFORTS.anthropic[anthropic.model], []);});
test("native SDKs serialize effort on both stream and generate without budget mappings", async (t) => { let captured; const response = () => Response.json({ error: { message: "wire captured" } }, { status: 400 }); t.mock.method(globalThis, "fetch", async (_url, init) => { captured = JSON.parse(init.body); return response(); }); const binding = { run: async (_model, body) => { captured = body; return response(); }, gateway: () => ({ run: async (request) => { captured = request.query; return response(); }, }), }; for (const [provider, models] of Object.entries(MODEL_EFFORTS)) { for (const [model, efforts] of Object.entries(models)) { // Workers AI's ordinary non-reasoning transport is covered separately. if (provider === "workers-ai" && !efforts.length) continue; for (const effort of [undefined, ...efforts]) { const configuration = parseModelConfiguration({ provider, model, ...(effort ? { effort } : {}), }); const languageModel = createConfiguredModel( binding, configuration, new Secret(credential), ); for (const operation of ["doGenerate", "doStream"]) { captured = undefined; await assert.rejects( languageModel[operation]({ ...params, providerOptions: modelProviderOptions(configuration), }), ); assert.ok( captured, `${provider}/${model}/${operation} must reach transport`, ); assert.equal(captured.reasoning_effort, undefined); if (provider === "openrouter") { assert.deepEqual( captured.reasoning, effort ? { effort } : undefined, ); assert.equal(captured.output_config, undefined); } else { assert.deepEqual( captured.output_config, effort ? { effort } : undefined, ); assert.deepEqual( captured.thinking, effort && provider === "anthropic" ? { type: "adaptive" } : undefined, ); } assert.ok(!JSON.stringify(captured).includes("budget_tokens")); } } } }});
test("adaptive reasoning preserves signed and redacted blocks across a real SDK tool round-trip", async (t) => { const configuration = { provider: "anthropic", model: "claude-sonnet-5", effort: "high", }; const requests = []; t.mock.method(globalThis, "fetch", async (_url, init) => { const body = JSON.parse(init.body); requests.push(body); const first = requests.length === 1; const blocks = first ? [ { type: "thinking", thinking: "", signature: "" }, { type: "redacted_thinking", data: "opaque-redacted-block" }, { type: "tool_use", id: "call-1", name: "echo", input: {} }, ] : [{ type: "text", text: "Done" }]; const events = [ { type: "message_start", message: { id: `msg-${requests.length}`, type: "message", role: "assistant", model: configuration.model, content: [], stop_reason: null, stop_sequence: null, usage: { input_tokens: 1, output_tokens: 0 }, }, }, ]; for (const [index, block] of blocks.entries()) { events.push({ type: "content_block_start", index, content_block: block }); if (block.type === "thinking") events.push({ type: "content_block_delta", index, delta: { type: "signature_delta", signature: "opaque-thinking-signature", }, }); events.push({ type: "content_block_stop", index }); } events.push( { type: "message_delta", delta: { stop_reason: first ? "tool_use" : "end_turn", stop_sequence: null, }, usage: { output_tokens: 1 }, }, { type: "message_stop" }, ); return new Response(anthropicMessagesStream(events), { headers: { "content-type": "text/event-stream" }, }); }); const result = streamText({ model: protectModel( createConfiguredModel({}, configuration, new Secret(credential)), "anthropic", ), providerOptions: modelProviderOptions(configuration), prompt: "Use the echo tool", tools: { echo: tool({ inputSchema: jsonSchema({ type: "object", properties: {} }), execute: async () => "ok", }), }, stopWhen: stepCountIs(2), }); await result.consumeStream(); assert.equal(requests.length, 2); for (const request of requests) assert.deepEqual(request.output_config, { effort: "high" }); const assistant = requests[1].messages.find( (message) => message.role === "assistant", ); assert.deepEqual(assistant.content.slice(0, 2), [ { type: "thinking", thinking: "", signature: "opaque-thinking-signature" }, { type: "redacted_thinking", data: "opaque-redacted-block" }, ]);});
test("reasoning protocol metadata is allowlisted, not arbitrary provider diagnostics", async () => { const metadata = { anthropic: { signature: "", redactedData: "opaque", secret: credential }, other: { secret: credential }, }; const model = protectModel( new MockLanguageModelV3({ doGenerate: async () => ({ content: [{ type: "reasoning", text: "", providerMetadata: metadata }], finishReason: { unified: "stop" }, usage: { inputTokens: {}, outputTokens: {} }, }), doStream: async () => ({ stream: new ReadableStream({ start(controller) { controller.enqueue({ type: "reasoning-start", id: "r", providerMetadata: metadata, }); controller.enqueue({ type: "reasoning-delta", id: "r", delta: "", providerMetadata: metadata, }); controller.enqueue({ type: "reasoning-end", id: "r", providerMetadata: metadata, }); controller.close(); }, }), }), }), "anthropic", ); const generated = await model.doGenerate(params); assert.deepEqual(generated.content[0].providerMetadata, { anthropic: { signature: "", redactedData: "opaque" }, }); const chunks = []; for await (const chunk of (await model.doStream(params)).stream) { chunks.push(chunk); assert.deepEqual( chunk.providerMetadata, generated.content[0].providerMetadata, ); } assert.equal(chunks.length, 3); assert.ok(!JSON.stringify({ generated, chunks }).includes(credential));});
test("OpenRouter model selection and credentials are bounded by the provider registry", () => { for (const model of Object.keys(MODEL_PROVIDERS.openrouter.models)) assert.deepEqual(parseModelConfiguration({ ...openrouter, model }), { ...openrouter, model, }); for (const value of [ { provider: "opencode-go", model: "glm-5.3" }, { ...openrouter, model: anthropic.model }, { ...openrouter, baseURL: "https://attacker.example" }, { ...openrouter, gatewayId: "other" }, { provider: "__proto__", model: "constructor" }, { provider: "constructor", model: "name" }, ]) assert.throws(() => parseModelConfiguration(value), /supported provider/); assert.equal(parseProviderKey("openrouter", credential).reveal(), credential); assert.equal(parseProviderKey("openrouter", null), null); assert.throws( () => parseProviderKey("opencode-go", credential), /Unsupported/, ); assert.throws( () => parseProviderKey("workers-ai", credential), /Unsupported/, ); assert.throws( () => createConfiguredModel({}, openrouter), /Add an OpenRouter API key/, );});
test("OpenRouter streams text and tool arguments through the native gateway with cancellation", async (t) => { t.mock.method(globalThis, "fetch", () => { throw new Error("Direct provider fetch is forbidden"); }); const requests = []; const signal = new AbortController().signal; const binding = { gateway(id) { assert.equal(id, "default"); return { run: async (request, options) => { requests.push({ request, options }); return new Response( workersAIStream([ { id: "chat_openrouter", choices: [ { index: 0, delta: { role: "assistant", content: "Hello" }, finish_reason: null, }, ], }, { id: "chat_openrouter", choices: [ { index: 0, delta: { tool_calls: [ { index: 0, id: "call_openrouter", type: "function", function: { name: "shell", arguments: '{"command":' }, }, ], }, finish_reason: null, }, ], }, { id: "chat_openrouter", choices: [ { index: 0, delta: { tool_calls: [ { index: 0, function: { arguments: '"printf 0"}' } }, ], }, finish_reason: null, }, ], }, { id: "chat_openrouter", choices: [{ index: 0, delta: {}, finish_reason: "tool_calls" }], usage: { prompt_tokens: 2, completion_tokens: 3, total_tokens: 5, }, }, ]), { headers: { "content-type": "text/event-stream" } }, ); }, }; }, }; const model = protectModel( createConfiguredModel( binding, openrouter, new Secret(credential), "conversation-openrouter", ), openrouter.provider, ); assert.equal(model.provider, "openrouter"); for (let turn = 0; turn < 2; turn++) { const result = await model.doStream({ ...params, abortSignal: signal, headers: { authorization: "override" }, }); const chunks = []; for await (const chunk of result.stream) chunks.push(chunk); assert.equal( chunks .filter((chunk) => chunk.type === "text-delta") .map((chunk) => chunk.delta) .join(""), "Hello", ); assert.equal( chunks.find((chunk) => chunk.type === "tool-call").input, '{"command":"printf 0"}', ); assert.ok(!JSON.stringify({ ...result, chunks }).includes(credential)); } for (const { request, options } of requests) { assert.equal(request.provider, "openrouter"); assert.equal( request.endpoint, "https://openrouter.ai/api/v1/chat/completions", ); assert.equal(request.query.model, openrouter.model); assert.equal(request.query.stream, true); assert.equal(request.headers.authorization, `Bearer ${credential}`); assert.deepEqual( Object.keys(request.headers).filter((name) => name.includes("session")), [], ); assert.match(request.headers["user-agent"], /^flarebot\//); assert.equal(request.headers["cf-aig-collect-log"], "false"); assert.equal(request.headers["cf-aig-skip-cache"], "true"); assert.equal(options.signal, signal); }});
test("OpenRouter errors distinguish billing, authentication and rate limits without exposing upstream data", async () => { for (const [status, message] of [ [401, /OpenRouter rejected authentication/], [402, /OpenRouter allowance or balance/], [429, /OpenRouter rate limit/], [500, /OpenRouter request failed/], ]) { const binding = { gateway: () => ({ run: async () => Response.json({ error: { message: credential } }, { status }), }), }; const model = protectModel( createConfiguredModel(binding, openrouter, new Secret(credential)), openrouter.provider, ); for (const operation of ["doGenerate", "doStream"]) await assert.rejects(model[operation](params), (error) => { assert.match(error.message, message); assert.ok(!JSON.stringify(error).includes(credential)); assert.equal(error.cause, undefined); return true; }); } let called = false; const controller = new AbortController(); controller.abort(); const binding = { gateway: () => ({ run: async () => { called = true; throw new Error(credential); }, }), }; const model = protectModel( createConfiguredModel(binding, openrouter, new Secret(credential)), openrouter.provider, ); await assert.rejects( model.doStream({ ...params, abortSignal: controller.signal }), { name: "AbortError" }, ); assert.equal(called, false);});
function workersAIStream(events) { const encoder = new TextEncoder(); return new ReadableStream({ start(controller) { for (const event of events) controller.enqueue( encoder.encode(`data: ${JSON.stringify(event)}\n\n`), ); controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.close(); }, });}
function anthropicMessagesStream(events) { const encoder = new TextEncoder(); return new ReadableStream({ start(controller) { for (const event of events) controller.enqueue( encoder.encode( `event: ${event.type}\ndata: ${JSON.stringify(event)}\n\n`, ), ); controller.close(); }, });}
test("configuration is strictly bounded and credential encryption is installation bound", async () => { assert.deepEqual(parseModelConfiguration(DEFAULT_MODEL), DEFAULT_MODEL); for (const value of [ null, [], {}, { ...anthropic, apiKey: credential }, { provider: "anthropic", model: "../../host" }, { provider: "workers-ai", model: anthropic.model }, ]) assert.throws( () => parseModelConfiguration(value), /supported provider and model/, ); for (const key of ["", "x".repeat(513), "x".repeat(20) + "\n", 123, {}]) assert.throws(() => parseProviderKey("anthropic", key), /API key/); assert.throws(() => parseProviderKey("other", credential), /Unsupported/); const key = parseProviderKey("anthropic", credential); assert.equal(JSON.stringify(key), '"[REDACTED]"'); const ciphertext = await encryptProviderKey(key, secret, installation); assert.ok(!ciphertext.includes(credential)); assert.notEqual( await encryptProviderKey(key, secret, installation), ciphertext, ); assert.equal( (await decryptProviderKey(ciphertext, secret, installation)).reveal(), credential, ); for (const [value, signingSecret, id] of [ [ciphertext, new Secret("rotated-session-signing-key"), installation], [ciphertext, secret, "b".repeat(32)], ["invalid-envelope", secret, installation], ]) await assert.rejects( decryptProviderKey(value, signingSecret, id), /replace it in settings/, );});
test("official Anthropic adapter sends the key only to its fixed provider endpoint and sanitizes failures", async (t) => { let calls = 0; t.mock.method(globalThis, "fetch", async (url, init) => { calls++; assert.equal(String(url), "https://api.anthropic.com/v1/messages"); assert.equal(new Headers(init.headers).get("x-api-key"), credential); const body = JSON.parse(init.body); assert.equal(body.model, anthropic.model); assert.ok(!init.body.includes(credential)); return Response.json( { type: "error", error: { type: "authentication_error", message: credential }, }, { status: 401 }, ); }); assert.throws( () => createConfiguredModel({}, anthropic), /Add an Anthropic API key/, ); const model = protectModel( createConfiguredModel({}, anthropic, new Secret(credential)), "anthropic", ); for (const operation of ["doStream", "doGenerate"]) await assert.rejects(model[operation](params), (error) => { assert.match(error.message, /Anthropic rejected authentication/); assert.ok(!JSON.stringify(error).includes(credential)); assert.equal(error.cause, undefined); return true; }); assert.equal(calls, 2);});
test("Workers AI catalog labels and routes text models through the default gateway", async () => { for (const [model, protocol] of Object.entries( MODEL_PROVIDERS["workers-ai"].models, )) { const configuration = { provider: "workers-ai", model }; assert.deepEqual(parseModelConfiguration(configuration), configuration); assert.match(modelLabel("workers-ai", model), / · (Free allowance|Paid)$/); if (protocol !== "workers-ai") continue; let request; const binding = { run: async (...args) => { request = args; return workersAIStream([{ response: "Hello" }]); }, }; const result = await createConfiguredModel(binding, configuration).doStream( params, ); for await (const chunk of result.stream) assert.notEqual(chunk.type, "error"); assert.equal(request[0], model); assert.deepEqual(request[2].gateway, { id: "default" }); } assert.equal( modelLabel("workers-ai", "@cf/zai-org/glm-5.3"), "@cf/zai-org/glm-5.3 · Paid", ); assert.equal( modelLabel("openrouter", "minimax/minimax-m3:free"), "minimax/minimax-m3:free · Free", ); assert.equal( modelLabel("anthropic", anthropic.model), `${anthropic.model} · Paid`, );});
test("Workers AI mixed-format stream events retain numeric tool arguments once", async () => { const argumentChunks = ['{"command":"printf ', 0, '"}']; const events = [ { response: "Hello World", choices: [{ delta: { content: "Hello World" } }], }, ...argumentChunks.map((argumentsDelta, index) => { const nativeToolCall = { index: 0, ...(index === 0 ? { id: "call", type: "function" } : {}), ...(index === 0 ? { name: "shell" } : {}), arguments: argumentsDelta, }; const openAIToolCall = { index: 0, ...(index === 0 ? { id: "call", type: "function" } : {}), function: { ...(index === 0 ? { name: "shell" } : {}), arguments: String(argumentsDelta), }, }; return { tool_calls: [nativeToolCall], choices: [{ delta: { tool_calls: [openAIToolCall] } }], }; }), ]; const binding = { run: async () => workersAIStream(events) }; const result = await createConfiguredModel(binding, DEFAULT_MODEL).doStream( params, ); const chunks = []; for await (const chunk of result.stream) chunks.push(chunk); assert.equal( chunks .filter((chunk) => chunk.type === "text-delta") .map((chunk) => chunk.delta) .join(""), "Hello World", ); assert.equal( chunks.find((chunk) => chunk.type === "tool-call").input, '{"command":"printf 0"}', );});
test("Inkling uses Cloudflare's Anthropic Messages wire format", async () => { let request; const events = [ { type: "message_start", message: { id: "msg_inkling", type: "message", role: "assistant", content: [], model: "thinkingmachines/Inkling", stop_reason: null, stop_sequence: null, usage: { input_tokens: 1, output_tokens: 0 }, }, }, { type: "content_block_start", index: 0, content_block: { type: "text", text: "" }, }, { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "Inkling" }, }, { type: "content_block_stop", index: 0 }, { type: "message_delta", delta: { stop_reason: "end_turn", stop_sequence: null }, usage: { output_tokens: 1 }, }, { type: "message_stop" }, ]; const binding = { run: async (model, input, options) => { request = { model, input, options }; const stream = anthropicMessagesStream(events); return options?.returnRawResponse ? new Response(stream, { headers: { "content-type": "text/event-stream" }, }) : stream; }, }; const model = createConfiguredModel( binding, { provider: "workers-ai", model: "thinkingmachines/inkling-256k", }, undefined, "conversation-id", ); const chunks = []; for await (const chunk of (await model.doStream(params)).stream) chunks.push(chunk); assert.equal( chunks .filter((chunk) => chunk.type === "text-delta") .map((chunk) => chunk.delta) .join(""), "Inkling", ); assert.equal(request.model, "thinkingmachines/inkling-256k"); assert.equal(request.input.model, undefined); assert.equal(request.input.stream, true); assert.deepEqual(request.options.gateway, { id: "default" }); assert.equal(request.options.returnRawResponse, true); assert.deepEqual(request.options.extraHeaders, { "x-session-affinity": "conversation-id", });});
test("Inkling reports insufficient Cloudflare AI Gateway balance", async () => { const binding = { run: async () => Response.json( { error: [ { code: 2021, message: "Insufficient balance; add money to your gateway or use BYOK", }, ], }, { status: 402 }, ), }; const model = protectModel( createConfiguredModel(binding, { provider: "workers-ai", model: "thinkingmachines/inkling-256k", }), "workers-ai", ); await assert.rejects(model.doStream(params), /AI Gateway balance/);});
test("model boundary skips raw data, strips diagnostics, sanitizes stream errors and retains cancellation", async () => { const model = protectModel( new MockLanguageModelV3({ doStream: async () => ({ request: { body: credential }, response: { headers: { secret: credential } }, stream: new ReadableStream({ start(controller) { for (const value of [ { type: "raw", rawValue: credential }, { type: "stream-start", warnings: [{ type: "other", message: credential }], }, { type: "text-delta", id: "a", delta: "safe", providerMetadata: { secret: credential }, }, { type: "error", error: new Error(credential) }, ]) controller.enqueue(value); controller.close(); }, }), }), }), "anthropic", ); const result = await model.doStream(params); const chunks = []; for await (const chunk of result.stream) chunks.push(chunk); assert.deepEqual( chunks.map((c) => c.type), ["stream-start", "text-delta", "error"], ); assert.equal(chunks[1].delta, "safe"); assert.match(chunks[2].error.message, /Anthropic request failed/); assert.ok(!JSON.stringify({ ...result, chunks }).includes(credential));
const controller = new AbortController(); const aborted = protectModel( new MockLanguageModelV3({ doStream: async () => { controller.abort(); throw new Error(credential); }, }), "anthropic", ); await assert.rejects( aborted.doStream({ ...params, abortSignal: controller.signal }), { name: "AbortError", message: "Model request cancelled" }, );
const streamFailure = protectModel( new MockLanguageModelV3({ doStream: async () => ({ stream: new ReadableStream({ pull(controller) { controller.error(new Error(credential)); }, }), }), }), "anthropic", ); await assert.rejects(async () => { for await (const _chunk of (await streamFailure.doStream(params)).stream) { /* drain */ } }, /Anthropic request failed/);});
test("optional diagnostics observe one terminal result per provider attempt without content", async () => { const events = []; const successful = protectModel( new MockLanguageModelV3({ doGenerate: async () => ({ content: [{ type: "text", text: credential }], finishReason: { unified: "stop", raw: undefined }, usage: { inputTokens: { total: 1 }, outputTokens: { total: 1 } }, warnings: [], providerMetadata: { secret: credential }, }), }), "anthropic", (event) => events.push(event), ); await successful.doGenerate(params); assert.deepEqual( events.map((e) => e.status), ["started", "completed"], ); assert.equal(events[0].attemptId, events[1].attemptId); assert.ok(events[1].durationMs >= 0); assert.equal(events[0].durationMs, null); assert.ok(!JSON.stringify(events).includes(credential)); for (const mode of ["throw", "error-chunk", "read-error", "abort"]) { const observed = []; const abort = new AbortController(); const model = protectModel( new MockLanguageModelV3({ doStream: async () => { if (mode === "throw") throw new Error(credential); if (mode === "abort") { abort.abort(); throw new Error(credential); } return { stream: new ReadableStream({ start(controller) { if (mode === "read-error") { controller.error(new Error(credential)); return; } controller.enqueue({ type: "error", error: new Error(credential), }); controller.close(); }, }), }; }, }), "anthropic", (event) => observed.push(event), ); try { for await (const _ of ( await model.doStream({ ...params, abortSignal: abort.signal }) ).stream) { } } catch {} assert.deepEqual( observed.map((e) => e.status), ["started", mode === "abort" ? "aborted" : "error"], ); assert.equal(observed[0].attemptId, observed[1].attemptId); assert.ok(!JSON.stringify(observed).includes(credential)); }});
test("failed model diagnostics retain only bounded HTTP status and failure stage", async () => { for (const stage of ["request", "stream-read", "stream-event"]) { for (const statusCode of [ 401, 402, 404, 429, 502, credential, 999, 401.5, undefined, ]) { const events = []; const error = Object.assign(new Error(credential), { statusCode, responseBody: credential, requestHeaders: { authorization: credential }, }); const model = protectModel( new MockLanguageModelV3({ doStream: async () => { if (stage === "request") throw error; return { stream: new ReadableStream({ start(controller) { if (stage === "stream-read") controller.error(error); else { controller.enqueue({ type: "error", error }); controller.close(); } }, }), }; }, }), "openrouter", (event) => events.push(event), ); try { for await (const _ of (await model.doStream(params)).stream) { } } catch {} assert.equal(events.at(-1).failureStage, stage); assert.equal( events.at(-1).httpStatus, Number.isInteger(statusCode) && statusCode >= 400 && statusCode <= 599 ? statusCode : null, ); assert.ok(!JSON.stringify(events).includes(credential)); } }});
test("gateway failures expose safe codes and correlation IDs without raw error content", async () => { for (const body of [ JSON.stringify({ name: "AiGatewayError", internalCode: 2006, message: credential, }), credential, "x".repeat(8193), JSON.stringify({ name: "AiGatewayError", internalCode: credential }), ]) { const events = []; const error = Object.assign(new Error(credential), { statusCode: 502, responseBody: body, responseHeaders: { "cf-ray": "0123456789abcdef-WAW", authorization: credential, }, }); const model = protectModel( new MockLanguageModelV3({ doStream: async () => { throw error; }, }), "openrouter", (e) => events.push(e), ); await assert.rejects(model.doStream(params)); const failure = events.at(-1); assert.equal(failure.errorCategory, "upstream-unavailable"); assert.equal( failure.gatewayErrorCode, body.includes('"internalCode":2006') ? 2006 : null, ); assert.equal(failure.cfRay, "0123456789abcdef-WAW"); assert.ok(!JSON.stringify(events).includes(credential)); }});
test("stream attempts stay open until consumption and retain backpressure, cancellation and observer isolation", async () => { const observed = []; let pulls = 0; let cancelled; const model = protectModel( new MockLanguageModelV3({ doStream: async () => ({ stream: new ReadableStream( { pull(controller) { pulls++; controller.enqueue({ type: "text-delta", id: "text", delta: "safe", }); }, cancel(reason) { cancelled = reason; }, }, { highWaterMark: 0 }, ), }), }), "anthropic", (event) => { observed.push(event); throw new Error(credential); }, ); const { stream } = await model.doStream(params); await new Promise((resolve) => setTimeout(resolve, 15)); assert.equal( observed.length, 1, "response headers must not finish the attempt", ); assert.ok(pulls <= 1, "no eager diagnostic drain"); const reader = stream.getReader(); assert.equal((await reader.read()).value.delta, "safe"); await reader.cancel(credential); assert.equal(cancelled, credential); assert.deepEqual( observed.map((e) => e.status), ["started", "aborted"], ); assert.ok(!JSON.stringify(observed).includes(credential)); const completed = []; const finite = protectModel( new MockLanguageModelV3({ doStream: async () => ({ stream: new ReadableStream({ start(controller) { controller.enqueue({ type: "raw", rawValue: credential }); controller.close(); }, }), }), }), "anthropic", (event) => completed.push(event), ); for await (const _ of (await finite.doStream(params)).stream) { } assert.deepEqual( completed.map((e) => e.status), ["started", "completed"], );});