Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291import * as Effect from "effect/Effect";import 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("native attachment images and PDFs reach each provider transport", async (t) => { const png = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAwMCAO+jRZkAAAAASUVORK5CYII="; const pdf = Buffer.from("%PDF-1.4\nfixture\n%%EOF").toString("base64"); let captured; t.mock.method(globalThis, "fetch", async (_url, init) => { captured = JSON.parse(init.body); return new Response("captured", { status: 400 }); }); const binding = { gateway: () => ({run: async ({query}) => {captured = query; throw new Error("captured");}}), run: async (_model, body) => { captured = body; throw new Error("captured"); }, }; const configurations = [ { provider: "workers-ai", model: "@cf/meta/llama-4-scout-17b-16e-instruct", }, anthropic, {provider: "workers-ai", model: "@cf/qwen/qwen3.8-27b"}, {provider: "openrouter", model: "openai/gpt-5.6-luna"}, ]; for (const configuration of configurations) { for (const operation of ["doGenerate", "doStream"]) { captured = undefined; const files = [ { type: "file", mediaType: "image/png", data: { type: "data", data: png }, }, ]; if (configuration.provider !== "workers-ai") files.push({ type: "file", mediaType: "application/pdf", data: { type: "data", data: pdf }, }); const model = protectModel( createConfiguredModel(binding, configuration, new Secret(credential)), configuration.provider, ); await assert.rejects( model[operation]({ prompt: [ { role: "user", content: [{ type: "text", text: "Inspect the attachment" }], }, { role: "assistant", content: [ { type: "tool-call", toolCallId: "call_1", toolName: "read_attachment", input: { id: "fixture" }, }, ], }, { role: "tool", content: [ { type: "tool-result", toolCallId: "call_1", toolName: "read_attachment", output: { type: "content", value: [{ type: "text", text: "Attachment" }, ...files], }, }, ], }, ], }), ); assert.ok( captured, `${configuration.provider}/${operation} reaches the transport`, ); if (configuration.provider !== "anthropic") { assert.equal( captured.messages.at(-1).content.find(p => p.type === "image_url").image_url.url, `data:image/png;base64,${png}`, ); assert.equal(captured.messages.at(-2).role, "tool"); if (configuration.provider === "openrouter") assert.equal(captured.messages.at(-1).content.find(p => p.type === "file").file.file_data, `data:application/pdf;base64,${pdf}`); } else { const result = captured.messages .at(-1) .content.find((part) => part.type === "tool_result"); assert.equal( result.content.find((part) => part.type === "image").source.data, png, ); assert.equal( result.content.find((part) => part.type === "document").source.data, pdf, ); } } }});
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; throw new Error("wire captured"); }, 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`, ); if ( provider === "workers-ai" && model !== "thinkingmachines/inkling-256k" ) { assert.equal(captured.reasoning_effort, effort); assert.equal(captured.output_config, undefined); assert.equal(captured.thinking, undefined); } else if (provider === "openrouter") { assert.equal(captured.reasoning_effort, undefined); assert.deepEqual( captured.reasoning, effort ? { effort } : undefined, ); assert.equal(captured.output_config, undefined); } else { assert.equal(captured.reasoning_effort, undefined); 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 Effect.runPromise( encryptProviderKey(key, secret, installation), ); assert.ok(!ciphertext.includes(credential)); assert.notEqual( await Effect.runPromise(encryptProviderKey(key, secret, installation)), ciphertext, ); assert.equal( ( await Effect.runPromise( 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( Effect.runPromise(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"], );});
test("provider setup and normalization failures stay safe and finish their attempts", async () => { const locked = new ReadableStream(); const lock = locked.getReader(); const observations = []; try { const model = protectModel( new MockLanguageModelV3({ doStream: async () => ({ stream: locked }) }), "workers-ai", (event) => observations.push(event), ); await assert.rejects(model.doStream(params), /Workers AI request failed/); assert.deepEqual( observations.map((event) => event.status), ["started", "error"], ); assert.equal(observations.at(-1).failureStage, "request"); } finally { lock.releaseLock(); }
for (const streaming of [false, true]) { const events = []; const part = { type: "reasoning", get providerMetadata() { throw new Error(credential); }, }; const source = new ReadableStream({ start(controller) { controller.enqueue(part); }, }); const model = protectModel( new MockLanguageModelV3({ doGenerate: async () => ({ content: [part] }), doStream: async () => ({ stream: source }), }), "workers-ai", (event) => events.push(event), ); await assert.rejects(async () => { if (streaming) for await (const _ of (await model.doStream(params)).stream) { } else await model.doGenerate(params); }, /Workers AI request failed/); assert.deepEqual( events.map((event) => event.status), ["started", "error"], ); assert.ok(!JSON.stringify(events).includes(credential)); if (streaming) assert.equal(source.locked, false); }});