Files
OmniRoute/tests/unit/stream-utils.test.ts
Diego Rodrigues de Sa e Souza a25d5f1ef6 Release v3.8.13 (#3327)
* chore(release): open v3.8.13 development cycle

Bump 3.8.12 → 3.8.13 across package.json, lockfile, electron/, open-sse/, and
docs/reference/openapi.yaml; add the [3.8.13] cycle placeholder to the root
CHANGELOG and the 41 i18n mirrors. Integration branch for the v3.8.13 cycle —
fixes/features land here via per-issue PRs and it merges to main at release time.

* fix(ci): skip auto-deploy when VPS host is unreachable from the runner (#3299)

Integrated into release/v3.8.13

* fix(dev): auto-rebuild better-sqlite3 on Node ABI mismatch at dev startup (#3301)

Integrated into release/v3.8.13

* feat(api): accept path-scoped API keys on client API routes (#3300)

Integrated into release/v3.8.13

* fix(sse): harden against empty responses causing Copilot Chat failures (#3297)

Integrated into release/v3.8.13

* fix(api): remove Completions.me rickroll provider (discussion #3293) (#3302)

Integrated into release/v3.8.13

* fix(opencode-provider): extract contextLength from live model catalog (#3298)

Integrated into release/v3.8.13

* feat(web-cookie): self-service login infrastructure + auto-refresh daemon (#3292)

Integrated into release/v3.8.13

* docs(changelog): record the v3.8.13 PRs merged this round (#3292/#3300/#3297/#3298/#3301/#3302/#3299)

* fix(auth): harden URL token extraction — drop query-string fallback, gate to client routes (security follow-up to #3300) (#3309)

Security follow-up to #3300 — integrated into release/v3.8.13

* docs: rename resolve-issues → review-issues skill references

* fix(dashboard): keep no-auth providers visible under 'Show configured only' (#3290) (#3312)

no-auth providers (opencode, duckduckgo-web, theoldllm, veoaifree-web) never
create a DB connection row so stats.total stays 0, which the configured-only
filter treated as 'unconfigured' and hid them — even though they are always
usable and appear unconditionally in /v1/models. filterConfiguredProviderEntries
now treats displayAuthType === 'no-auth' as configured.

Co-authored-by: uniQta <uniQta@users.noreply.github.com>

* fix(cli): resolve update paths relative to script + recursive backup (#3295) (#3313)

omniroute update always failed on a global install:
- getCurrentVersion() read package.json from process.cwd(), which on a global
  npm/brew install is the user's working dir, not the package root → null →
  'Could not determine current version'.
- createBackup() resolved bin/ from cwd too, and passed the 'cli' directory to
  copyFileSync → EISDIR, swallowed by the catch → 'Failed to create backup'.

Both now resolve package.json/bin relative to the script via import.meta.url,
and the backup uses cpSync({recursive:true}) so the cli/ directory is copied.

Co-authored-by: uniQta <uniQta@users.noreply.github.com>

* fix(theoldllm): read upstream body once to avoid [502] body-already-read (#3296) (#3314)

On the cached-token path the executor never enters the refresh branch, so the
same upstream Response was read with .text() twice (token-rejection check +
final body). A Response body is single-use, so the second read threw
'Body is unusable: Body has already been read', caught and surfaced as [502].

Read the body once into finalBody and only re-read after a token-rejection
refetch.

Co-authored-by: onizukashonan14-png <onizukashonan14-png@users.noreply.github.com>

* fix(sse): strip leaked internal tool envelopes from streaming output (#3311)

Integrated into release/v3.8.13

* fix(sse): expose Claude + Gemini budget tiers in the antigravity catalog (#3184) (#3303)

Integrated into release/v3.8.13 (#3184)

* fix(catalog): compute combo context_length from known targets only (#3304)

Integrated into release/v3.8.13 — live contextLength + known-targets combo context (#3298 follow-up)

* chore(i18n): add message keys for proxy UI + vscode/ollama endpoint (#3307)

Integrated into release/v3.8.13 — i18n message keys for proxy UI + vscode/ollama

* feat(dashboard): i18n the proxy settings UI (#3310)

Integrated into release/v3.8.13 — i18n the proxy settings UI

* feat(api): model catalog enrichment + MCP model-catalog tools (#3306)

Integrated into release/v3.8.13 — model catalog enrichment + MCP model-catalog tools, reconciled with #3309 URL-token hardening

* test(catalog): align Antigravity preview-alias test with #3303 budget tiers

#3303 added the Gemini `-high`/`-low` budget tiers to ANTIGRAVITY_PUBLIC_MODELS
(user-callable on the Antigravity OAuth backend, verified via #3184), but did
not update the catalog-route test that asserted `antigravity/gemini-3.1-pro-high`
must NOT be exposed. The assertion now reflects the intended behavior — the
client-visible budget alias IS surfaced — while keeping the legacy
`gemini-claude-*` alias keys unexposed. Caught running the full catalog suite
on the merged release HEAD (the #3303 round only ran the antigravity-aliases
and usage-hardening files).

* docs(changelog): record the 6 PRs merged this review round into v3.8.13

#3306/#3307/#3310 (New Features — VS Code split: catalog+MCP, i18n keys, proxy
UI i18n), #3311/#3303/#3304 (Bug Fixes — SSE envelope sanitizer, antigravity
budget tiers, combo known-targets context_length).

* chore(release): finalize v3.8.13 changelog and cleanup

Finalize the v3.8.13 changelog with release date, maintenance notes,
and contributor credits. Update MCP docs to reference the correct tool
inventory diagram, exclude nested .claude worktrees from ESLint scans,
and tighten a response sanitizer type guard.

* fix(dashboard): refresh connections after provider auth import (#3320)

Integrated into release/v3.8.13 — refresh connections after provider auth import

* fix(codex): strip client-only params on native /responses passthrough (#3317) (#3325)

A /v1/responses request against the built-in codex/ provider does an
openai-responses -> openai-responses passthrough (CodexExecutor.transformRequest
returns the body early for _nativeCodexPassthrough). It forwarded client-only
fields verbatim and the Codex upstream rejected them with 400 Unsupported
parameter: prompt_cache_retention / safety_identifier / user — breaking Factory
Droid (which injects all three). The chat-completions path already strips these
(base.ts #1884, openai-responses translator #2770) but the passthrough skips
translation. Strip the three fields in the shared block before the passthrough
return; user is removed unconditionally since Codex /responses always rejects it.

Co-authored-by: tycronk20 <tycronk20@users.noreply.github.com>

* fix(dashboard): normalize agent-bridge /state response to stop page crash (#3318) (#3326)

The Agent Bridge page seeded a well-shaped initialData default then replaced it
wholesale with the raw /api/tools/agent-bridge/state response. The route returns
{ server, agents } but the UI reads { serverState, agentStates, bypassPatterns,
mappings }, so serverState became undefined and AgentBridgeServerCard crashed on
serverState.running — surfaced as the full-page 'Internal Server Error' boundary
(client render error, not a real 5xx).

Add a shared normalizeAgentBridgeState() that maps the route shape into the page
contract (server.running/certExists -> serverState) and always returns safe
defaults (never undefined serverState). Wired into both the SSR loader (page.tsx)
and the polling hook. The legacy 'agents' entry shape differs from AgentStateEntry
so it is not coerced; full route<->page contract reconciliation (port, upstreamCa,
bypassPatterns, mappings, agentStates) is a follow-up.

Co-authored-by: tycronk20 <tycronk20@users.noreply.github.com>

* docs: VS Code/Ollama endpoints + env & i18n tooling (#3319)

Integrated into release/v3.8.13 — VS Code/Ollama docs + env & i18n tooling

* feat(provider): test-all endpoint, rate-limit overrides, visibility f… (#3267)

Integrated into release/v3.8.13 — provider test-all endpoint, rate-limit overrides, model visibility

* feat: auto-combo optimization, playground model dropdown, only-configured toggle (#3322)

Integrated into release/v3.8.13 — auto-combo candidate expansion + playground dropdown + only-configured toggle

* feat(api): VS Code Copilot Ollama-compatible BYOK endpoint (#3316)

Integrated into release/v3.8.13 — VS Code Copilot Ollama-compatible BYOK endpoint (reconciled with #3306/#3309 auth hardening)

* chore(release): document #3320 in the v3.8.13 changelog + contributor credits

---------

Co-authored-by: Felipe Almeman <4226997+zhiru@users.noreply.github.com>
Co-authored-by: Wilson <pedbookmed@gmail.com>
Co-authored-by: Hernan Javier Ardila Sanchez <hjasgr@gmail.com>
Co-authored-by: Paijo <14921983+oyi77@users.noreply.github.com>
Co-authored-by: uniQta <uniQta@users.noreply.github.com>
Co-authored-by: onizukashonan14-png <onizukashonan14-png@users.noreply.github.com>
Co-authored-by: tycronk20 <tycronk20@users.noreply.github.com>
Co-authored-by: Vinayrnani <vinayrnani@gmail.com>
2026-06-06 19:13:11 -03:00

1661 lines
52 KiB
TypeScript

import test from "node:test";
import assert from "node:assert/strict";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-stream-utils-"));
process.env.DATA_DIR = TEST_DATA_DIR;
const core = await import("../../src/lib/db/core.ts");
const { createSSEStream, createSSETransformStreamWithLogger, createPassthroughStreamWithLogger } =
await import("../../open-sse/utils/stream.ts");
const {
buildStreamSummaryFromEvents,
compactStructuredStreamPayload,
createStructuredSSECollector,
} = await import("../../open-sse/utils/streamPayloadCollector.ts");
const { FORMATS } = await import("../../open-sse/translator/formats.ts");
const { createRequestLogger } = await import("../../open-sse/utils/requestLogger.ts");
const textEncoder = new TextEncoder();
const SYNTHETIC_CLAUDE_EMPTY_RESPONSE_TEXT =
"[Proxy Error] The upstream API returned an empty response. Please retry the request.";
async function readTransformed(chunks, options) {
const source = new ReadableStream({
start(controller) {
for (const chunk of chunks) {
controller.enqueue(textEncoder.encode(chunk));
}
controller.close();
},
});
return new Response(source.pipeThrough(createSSEStream(options))).text();
}
async function readWithTransform(chunks, transformStream) {
const source = new ReadableStream({
start(controller) {
for (const chunk of chunks) {
controller.enqueue(textEncoder.encode(chunk));
}
controller.close();
},
});
return new Response(source.pipeThrough(transformStream)).text();
}
test.after(() => {
core.resetDbInstance();
if (fs.existsSync(TEST_DATA_DIR)) {
for (const entry of fs.readdirSync(TEST_DATA_DIR)) {
fs.rmSync(path.join(TEST_DATA_DIR, entry), { recursive: true, force: true });
}
}
});
test("createSSEStream passthrough normalizes tool-call finishes and reports the assembled response", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_1",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: { role: "assistant", content: "Hello " } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_1",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [
{
index: 0,
delta: {
tool_calls: [
{
index: 0,
id: "call_1",
type: "function",
function: {
name: "read_file",
arguments: '{"path":"/tmp/a"}',
},
},
],
},
},
],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_1",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "openai",
model: "gpt-4.1-mini",
body: {
messages: [{ role: "user", content: "hello" }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.match(text, /"content":"Hello "/);
assert.match(text, /"name":"read_file"/);
assert.match(text, /"finish_reason":"tool_calls"/);
assert.equal(onCompletePayload.status, 200);
assert.equal(onCompletePayload.responseBody.choices[0].finish_reason, "tool_calls");
assert.equal(onCompletePayload.responseBody.choices[0].message.tool_calls[0].id, "call_1");
assert.equal(onCompletePayload.responseBody.choices[0].message.content, "Hello");
assert.equal(onCompletePayload.clientPayload._streamed, true);
});
test("createSSEStream passthrough converts textual tool-call content into structured call log tool_calls", async () => {
let onCompletePayload = null;
const toolArgs = JSON.stringify({
command: 'sqlite3 /root/.o\u200dmniroute/omniroute.db ".tables"',
});
const toolText = `[Tool call: terminal]\nArguments: ${toolArgs}`;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: { role: "assistant", content: toolText } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "antigravity",
model: "antigravity/gemini-3.5-flash-low",
body: {
messages: [{ role: "user", content: "inspect db" }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.equal(onCompletePayload.status, 200);
assert.match(text, /"tool_calls":\[/);
assert.match(text, /"name":"terminal"/);
assert.doesNotMatch(text, /"content":"\[Tool call: terminal/);
const choice = onCompletePayload.responseBody.choices[0];
assert.equal(choice.finish_reason, "tool_calls");
assert.equal(choice.message.content, null);
assert.equal(choice.message.tool_calls[0].function.name, "terminal");
assert.deepEqual(JSON.parse(choice.message.tool_calls[0].function.arguments), {
command: 'sqlite3 /root/.omniroute/omniroute.db ".tables"',
});
assert.doesNotMatch(text, /\[Tool call: terminal\]/);
});
test("createSSEStream passthrough converts split textual tool-call content at completion", async () => {
let onCompletePayload = null;
const splitToolArgs = JSON.stringify({
command: 'sqlite3 ~/.o\u200dmniroute/o\u200dmniroute.db ".tables"',
});
const chunks = ["[Tool call: terminal]\n", `Arguments: ${splitToolArgs}`];
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_split_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: { role: "assistant", content: chunks[0] } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_split_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: { content: chunks[1] } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_split_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "antigravity",
model: "antigravity/gemini-3.5-flash-low",
body: { messages: [{ role: "user", content: "inspect db" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.match(text, /"tool_calls":\[/);
assert.match(text, /"name":"terminal"/);
assert.doesNotMatch(text, /"content":"\[Tool call: terminal/);
const choice = onCompletePayload.responseBody.choices[0];
assert.equal(choice.finish_reason, "tool_calls");
assert.equal(choice.message.content, null);
assert.equal(choice.message.tool_calls[0].function.name, "terminal");
assert.deepEqual(JSON.parse(choice.message.tool_calls[0].function.arguments), {
command: 'sqlite3 ~/.omniroute/omniroute.db ".tables"',
});
assert.doesNotMatch(text, /\[Tool call: terminal\]/);
});
test("createSSEStream passthrough buffers fragmented textual tool-call JSON before emitting", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_fragmented_live_shape",
object: "chat.completion.chunk",
created: 1,
model: "MainAgent",
choices: [
{
index: 0,
delta: { role: "assistant", content: '[Tool call: terminal]\nArguments: {"' },
},
],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_fragmented_live_shape",
object: "chat.completion.chunk",
created: 1,
model: "MainAgent",
choices: [
{
index: 0,
delta: { content: 'command":"echo live_shape","timeout":10}' },
},
],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_fragmented_live_shape",
object: "chat.completion.chunk",
created: 1,
model: "MainAgent",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "omniroute",
model: "MainAgent",
body: { messages: [{ role: "user", content: "inspect" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.doesNotMatch(text, /\[Tool call:/);
assert.doesNotMatch(text, /Arguments:/);
assert.match(text, /"tool_calls":\[/);
assert.match(text, /"finish_reason":"tool_calls"/);
const choice = onCompletePayload.responseBody.choices[0];
assert.equal(choice.finish_reason, "tool_calls");
assert.equal(choice.message.content, null);
assert.equal(choice.message.tool_calls[0].function.name, "terminal");
assert.deepEqual(JSON.parse(choice.message.tool_calls[0].function.arguments), {
command: "echo live_shape",
timeout: 10,
});
});
test("createSSEStream passthrough suppresses trailing prose plus textual tool call", async () => {
let onCompletePayload = null;
const toolArgs = JSON.stringify({
command: "echo should_not_leak",
timeout: 10,
});
const toolText = `Вот оно! Статические файлы Next.js отдают 404. Чанки не найдены.\n\n[Tool call: terminal]\nArguments: ${toolArgs}`;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_trailing_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "MainAgent",
choices: [{ index: 0, delta: { role: "assistant", content: toolText } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_trailing_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "MainAgent",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "omniroute",
model: "MainAgent",
body: { messages: [{ role: "user", content: "inspect static files" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.equal(onCompletePayload.status, 200);
const choice = onCompletePayload.responseBody.choices[0];
assert.equal(choice.finish_reason, "tool_calls");
assert.equal(choice.message.content, null);
assert.equal(choice.message.tool_calls[0].function.name, "terminal");
assert.deepEqual(JSON.parse(choice.message.tool_calls[0].function.arguments), {
command: "echo should_not_leak",
timeout: 10,
});
assert.doesNotMatch(text, /\[Tool call: terminal\]/);
assert.doesNotMatch(text, /Arguments:/);
assert.doesNotMatch(JSON.stringify(onCompletePayload.responseBody), /\[Tool call: terminal\]/);
assert.doesNotMatch(JSON.stringify(onCompletePayload.responseBody), /Arguments:/);
});
test("createSSEStream passthrough suppresses textual tool calls for unknown tools", async () => {
let onCompletePayload = null;
const toolText = `[Tool call: search_files_ide]
Arguments: {"path":"/opt/OmniRoute/src","target":"files"}`;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_unknown_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: { role: "assistant", content: toolText } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_unknown_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "antigravity",
model: "antigravity/gemini-3.5-flash-low",
body: {
messages: [{ role: "user", content: "inspect files" }],
tools: [
{ type: "function", function: { name: "search_files", parameters: { type: "object" } } },
],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
const choice = onCompletePayload.responseBody.choices[0];
assert.equal(choice.finish_reason, "stop");
assert.equal(choice.message.content, null);
assert.equal(choice.message.tool_calls, undefined);
assert.doesNotMatch(text, /search_files_ide/);
assert.doesNotMatch(JSON.stringify(onCompletePayload.responseBody), /search_files_ide/);
});
test("createSSEStream passthrough suppresses malformed textual tool-call content", async () => {
let onCompletePayload = null;
const malformedToolText = `(empty)[Tool call: terminal]\nArguments: {"command":"sqlite3 /opt/O\u200dmniRoute/data/o\u200dmniroute.`;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_malformed_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: { role: "assistant", content: malformedToolText } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_malformed_textual_tool",
object: "chat.completion.chunk",
created: 1,
model: "antigravity/gemini-3.5-flash-low",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "antigravity",
model: "antigravity/gemini-3.5-flash-low",
body: { messages: [{ role: "user", content: "inspect db" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
const choice = onCompletePayload.responseBody.choices[0];
assert.equal(choice.finish_reason, "stop");
assert.equal(choice.message.content, null);
assert.equal(choice.message.tool_calls, undefined);
assert.doesNotMatch(text, /\[Tool call: terminal\]/);
assert.doesNotMatch(JSON.stringify(onCompletePayload.responseBody), /\[Tool call: terminal\]/);
});
test("createSSEStream suppresses malformed compact textual tool-call content", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
candidates: [
{
content: {
parts: [
{
text: "[Tool call: search_files_ide{file_glob:*combos*.ts,path:/opt/OmniRoute,target:files}]",
},
],
},
},
],
})}\n\n`,
`data: ${JSON.stringify({ candidates: [{ finishReason: "STOP" }] })}\n\n`,
],
{
mode: "translate",
targetFormat: FORMATS.ANTIGRAVITY,
sourceFormat: FORMATS.OPENAI,
provider: "antigravity",
model: "antigravity/gemini-3.5-flash-low",
body: { messages: [{ role: "user", content: "inspect files" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
const choice = onCompletePayload.responseBody.choices[0];
assert.equal(choice.finish_reason, "stop");
assert.equal(choice.message.content, null);
assert.equal(choice.message.tool_calls, undefined);
assert.doesNotMatch(JSON.stringify(onCompletePayload.responseBody), /\[Tool call:/);
});
test("createSSEStream passthrough flushes a buffered final line without a trailing newline", async () => {
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_2",
object: "chat.completion.chunk",
created: 2,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: { role: "assistant", content: "tail chunk" } }],
})}`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "openai",
model: "gpt-4.1-mini",
body: {
messages: [{ role: "user", content: "hello" }],
},
}
);
assert.match(text, /tail chunk/);
assert.equal(text.includes("data: "), true);
});
test("createSSEStream translate mode converts Claude SSE into OpenAI chunks and completion payload", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
type: "message_start",
message: {
id: "msg_1",
model: "claude-sonnet-4",
role: "assistant",
usage: { input_tokens: 3 },
},
})}\n\n`,
`data: ${JSON.stringify({
type: "content_block_start",
index: 0,
content_block: { type: "text", text: "" },
})}\n\n`,
`data: ${JSON.stringify({
type: "content_block_delta",
index: 0,
delta: { type: "text_delta", text: "Hello Claude" },
})}\n\n`,
`data: ${JSON.stringify({
type: "message_delta",
delta: { stop_reason: "end_turn" },
usage: { output_tokens: 4 },
})}\n\n`,
`data: ${JSON.stringify({
type: "message_stop",
})}\n\n`,
],
{
mode: "translate",
targetFormat: FORMATS.CLAUDE,
sourceFormat: FORMATS.OPENAI,
provider: "claude",
model: "claude-sonnet-4",
body: {
messages: [{ role: "user", content: "hello" }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.match(text, /"content":"Hello Claude"/);
assert.match(text, /\[DONE\]/);
assert.equal(onCompletePayload.status, 200);
assert.equal(onCompletePayload.responseBody.choices[0].message.content, "Hello Claude");
assert.equal(onCompletePayload.responseBody.usage.completion_tokens, 4);
assert.equal(onCompletePayload.responseBody.usage.total_tokens, 4);
});
test("createSSEStream Responses passthrough converts textual tool-call deltas before streaming", async () => {
let onCompletePayload = null;
const toolText = `[Tool call: terminal]
Arguments: {"command":"systemctl status omniroute"}`;
const text = await readTransformed(
[
`data: ${JSON.stringify({
type: "response.output_text.delta",
delta: toolText,
})}
`,
`data: ${JSON.stringify({
type: "response.completed",
response: {
id: "resp_textual_tool",
object: "response",
model: "antigravity/gemini-3.5-flash-low",
status: "completed",
output: [],
usage: { input_tokens: 10, output_tokens: 4, total_tokens: 14 },
},
})}
`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI_RESPONSES,
clientResponseFormat: FORMATS.OPENAI_RESPONSES,
provider: "antigravity",
model: "antigravity/gemini-3.5-flash-low",
body: {
input: "check service",
tools: [{ type: "function", name: "terminal", parameters: { type: "object" } }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.doesNotMatch(text, /\[Tool call: terminal\]/);
assert.doesNotMatch(text, /Arguments:/);
assert.match(text, /response.output_item.added/);
assert.match(text, /response.function_call_arguments.done/);
assert.match(text, /"name":"terminal"/);
assert.equal(onCompletePayload.responseBody.choices[0].finish_reason, "tool_calls");
assert.equal(onCompletePayload.responseBody.choices[0].message.content, null);
assert.equal(
onCompletePayload.responseBody.choices[0].message.tool_calls[0].function.name,
"terminal"
);
assert.doesNotMatch(JSON.stringify(onCompletePayload.clientPayload), /\[Tool call: terminal\]/);
});
test("createSSEStream passthrough preserves Responses API events and completion summaries", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
type: "response.output_text.delta",
delta: "Hello ",
})}\n\n`,
`data: ${JSON.stringify({
type: "response.output_text.delta",
delta: "world",
})}\n\n`,
`data: ${JSON.stringify({
type: "response.completed",
response: {
id: "resp_1",
object: "response",
model: "gpt-4.1-mini",
status: "completed",
usage: { input_tokens: 2, output_tokens: 3, total_tokens: 5 },
},
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI_RESPONSES,
provider: "openai",
model: "gpt-4.1-mini",
body: { input: "hello" },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.match(text, /response.output_text.delta/);
assert.match(text, /response.completed/);
assert.equal(onCompletePayload.responseBody.usage.total_tokens, 5);
assert.equal(onCompletePayload.providerPayload.summary.object, "response");
});
test("createSSEStream passthrough drops leaked empty chat bootstrap chunks for Responses clients", async () => {
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl-dummy",
object: "chat.completion.chunk",
created: 1,
model: "gpt-5.4",
choices: [
{ index: 0, delta: { role: "assistant", content: null, refusal: null }, finish_reason: null },
],
})}\n\n`,
`event: response.created\ndata: ${JSON.stringify({
type: "response.created",
response: {
id: "resp_1",
object: "response",
model: "gpt-5.4",
status: "in_progress",
output: [],
},
})}\n\n`,
`event: response.in_progress\ndata: ${JSON.stringify({
type: "response.in_progress",
response: {
id: "resp_1",
object: "response",
model: "gpt-5.4",
status: "in_progress",
output: [],
},
})}\n\n`,
`data: ${JSON.stringify({
type: "response.output_text.delta",
delta: "OK",
})}\n\n`,
`data: ${JSON.stringify({
type: "response.completed",
response: {
id: "resp_1",
object: "response",
model: "gpt-5.4",
status: "completed",
output: [],
usage: { input_tokens: 2, output_tokens: 1, total_tokens: 3 },
},
})}\n\n`,
`data: [DONE]\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI_RESPONSES,
clientResponseFormat: FORMATS.OPENAI_RESPONSES,
provider: "openai",
model: "gpt-5.4",
body: { input: "hello" },
}
);
assert.doesNotMatch(text, /chatcmpl-dummy/);
assert.match(text, /response\.created/);
assert.match(text, /response\.output_text\.delta/);
assert.match(text, /"delta":"OK"/);
});
test("buildStreamSummaryFromEvents falls back to response.output_text.delta when completed output is empty", () => {
const summary = buildStreamSummaryFromEvents(
[
{
index: 0,
data: {
type: "response.output_text.delta",
delta: "Hello ",
},
},
{
index: 1,
data: {
type: "response.output_text.delta",
delta: "world",
},
},
{
index: 2,
data: {
type: "response.completed",
response: {
id: "resp_fallback",
object: "response",
model: "gpt-5.4",
status: "completed",
output: [],
usage: { output_tokens: 2 },
},
},
},
],
FORMATS.OPENAI_RESPONSES,
"gpt-5.4"
);
assert.equal((summary as any).object, "response");
assert.equal((summary as any).output[0].type, "message");
assert.equal((summary as any).output[0].content[0].type, "output_text");
assert.equal((summary as any).output[0].content[0].text, "Hello world");
assert.equal((summary as any).usage.output_tokens, 2);
});
test("createSSEStream translate mode aborts on Responses failure with rate limit error", async () => {
let onCompletePayload = null;
await assert.rejects(
readTransformed(
[
`data: ${JSON.stringify({
type: "response.created",
response: {
id: "resp_fail",
object: "response",
model: "gpt-5.4",
status: "in_progress",
output: [],
},
})}\n\n`,
`data: ${JSON.stringify({
type: "response.failed",
response: {
id: "resp_fail",
object: "response",
model: "gpt-5.4",
status: "failed",
error: {
message: "Rate limit reached for gpt-5.4",
code: "rate_limit_exceeded",
},
},
})}\n\n`,
`data: [DONE]\n\n`,
],
{
mode: "translate",
targetFormat: FORMATS.OPENAI_RESPONSES,
sourceFormat: FORMATS.OPENAI,
provider: "codex",
model: "gpt-5.4",
body: { messages: [{ role: "user", content: "hello" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
),
/Rate limit reached for gpt-5\.4|Upstream failure/
);
assert.ok(onCompletePayload, "should capture completion payload before aborting");
assert.equal(onCompletePayload.status, 429);
assert.equal(onCompletePayload.responseBody.error.type, "rate_limit_error");
assert.equal(onCompletePayload.responseBody.error.code, "rate_limit_exceeded");
assert.match(onCompletePayload.responseBody.error.message, /Rate limit reached/);
});
test("createSSEStream passthrough restores Claude tool names from the mapping table", async () => {
const toolNameMap = new Map([["tool_alias", "read_file"]]);
const text = await readTransformed(
[
`data: ${JSON.stringify({
type: "content_block_start",
index: 0,
content_block: {
type: "tool_use",
id: "tool_1",
name: "tool_alias",
input: { path: "/tmp/a" },
},
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.CLAUDE,
provider: "claude",
model: "claude-sonnet-4",
toolNameMap,
body: { messages: [{ role: "user", content: "hello" }] },
}
);
assert.match(text, /"name":"read_file"/);
assert.equal(text.includes("tool_alias"), false);
});
test("createSSEStream passthrough fixes generic ids and normalizes reasoning aliases", async () => {
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chat",
object: "chat.completion.chunk",
created: 1,
model: "kimi-k2.5",
choices: [
{
index: 0,
delta: {
reasoning: "Let me think first",
},
},
],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "openai",
model: "kimi-k2.5",
body: { messages: [{ role: "user", content: "hello" }] },
}
);
assert.match(text, /"id":"chatcmpl-/);
assert.match(text, /"reasoning_content":"Let me think first"/);
assert.equal(text.includes('"reasoning":"Let me think first"'), false);
});
test("createSSEStream passthrough splits mixed reasoning and content deltas and estimates usage", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_reasoning",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [
{
index: 0,
delta: {
reasoning_content: "First think",
content: "Then answer",
},
},
],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_reasoning",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "openai",
model: "gpt-4.1-mini",
body: {
messages: [{ role: "user", content: "hello world" }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
const reasoningIndex = text.indexOf('"reasoning_content":"First think"');
const contentIndex = text.indexOf('"content":"Then answer"');
assert.ok(reasoningIndex >= 0);
assert.ok(contentIndex > reasoningIndex);
assert.match(text, /"total_tokens":\d+/);
assert.equal(onCompletePayload.responseBody.choices[0].message.reasoning_content, "First think");
assert.equal(onCompletePayload.responseBody.choices[0].message.content, "Then answer");
assert.ok(onCompletePayload.responseBody.usage.total_tokens > 0);
});
test("createSSEStream passthrough merges Claude usage chunks and restores mapped tool names", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
type: "message_start",
message: {
id: "msg_passthrough",
model: "claude-sonnet-4",
role: "assistant",
usage: { input_tokens: 6 },
},
})}\n\n`,
`data: ${JSON.stringify({
type: "content_block_start",
index: 0,
content_block: {
type: "tool_use",
id: "tool_1",
name: "tool_alias",
input: { path: "/tmp/a" },
},
})}\n\n`,
`data: ${JSON.stringify({
type: "content_block_delta",
index: 1,
delta: { text: "Claude says hi" },
})}\n\n`,
`data: ${JSON.stringify({
type: "message_delta",
delta: { stop_reason: "end_turn" },
usage: { output_tokens: 4 },
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.CLAUDE,
provider: "claude",
model: "claude-sonnet-4",
toolNameMap: new Map([["tool_alias", "read_file"]]),
body: {
messages: [{ role: "user", content: "hello" }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.match(text, /"name":"read_file"/);
assert.equal(text.includes('"name":"tool_alias"'), false);
assert.equal(onCompletePayload.responseBody.choices[0].message.content, "Claude says hi");
assert.equal(onCompletePayload.responseBody.usage.prompt_tokens, 6);
assert.equal(onCompletePayload.responseBody.usage.completion_tokens, 4);
assert.equal(onCompletePayload.responseBody.usage.total_tokens, 10);
});
test("createSSEStream passthrough injects a synthetic Claude text block for empty assistant SSE", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`event: message_start\ndata: ${JSON.stringify({
type: "message_start",
message: {
id: "msg_empty_passthrough",
type: "message",
role: "assistant",
model: "claude-sonnet-4",
content: [],
stop_reason: null,
stop_sequence: null,
usage: { input_tokens: 7, output_tokens: 0 },
},
})}\n\n`,
`event: message_stop\ndata: ${JSON.stringify({
type: "message_stop",
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.CLAUDE,
provider: "claude",
model: "claude-sonnet-4",
body: {
messages: [{ role: "user", content: "hello" }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.equal((text.match(/event: message_start/g) || []).length, 1);
assert.equal((text.match(/event: message_delta/g) || []).length, 1);
assert.match(text, /event: content_block_start/);
assert.match(text, /event: content_block_delta/);
assert.match(text, /event: message_stop/);
assert.match(text, /\[Proxy Error\] The upstream API returned an empty response/);
assert.ok(text.indexOf("event: content_block_start") > text.indexOf("event: message_start"));
assert.ok(text.indexOf("event: message_stop") > text.indexOf("event: content_block_stop"));
assert.equal(
onCompletePayload.responseBody.choices[0].message.content,
SYNTHETIC_CLAUDE_EMPTY_RESPONSE_TEXT
);
});
test("createSSEStream passthrough does not emit [DONE] for Claude SSE clients", async () => {
const text = await readTransformed(
[
`event: message_start\ndata: ${JSON.stringify({
type: "message_start",
message: {
id: "msg_claude_done_gate",
type: "message",
role: "assistant",
model: "claude-sonnet-4",
content: [],
stop_reason: null,
stop_sequence: null,
usage: { input_tokens: 3, output_tokens: 0 },
},
})}\n\n`,
`event: content_block_start\ndata: ${JSON.stringify({
type: "content_block_start",
index: 0,
content_block: { type: "text", text: "" },
})}\n\n`,
`event: content_block_delta\ndata: ${JSON.stringify({
type: "content_block_delta",
index: 0,
delta: { type: "text_delta", text: "Claude client stream" },
})}\n\n`,
`event: content_block_stop\ndata: ${JSON.stringify({
type: "content_block_stop",
index: 0,
})}\n\n`,
`event: message_delta\ndata: ${JSON.stringify({
type: "message_delta",
delta: { stop_reason: "end_turn", stop_sequence: null },
usage: { output_tokens: 3 },
})}\n\n`,
`event: message_stop\ndata: ${JSON.stringify({
type: "message_stop",
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.CLAUDE,
clientResponseFormat: FORMATS.CLAUDE,
provider: "claude",
model: "claude-sonnet-4",
body: {
messages: [{ role: "user", content: "hello" }],
},
}
);
assert.match(text, /event: message_stop/);
assert.match(text, /Claude client stream/);
assert.doesNotMatch(text, /\[DONE\]/);
});
test("createSSEStream translate mode injects a synthetic Claude text block when OpenAI finishes empty", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_empty_1",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: { role: "assistant" } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_empty_1",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
usage: { prompt_tokens: 3, completion_tokens: 0, total_tokens: 3 },
})}\n\n`,
],
{
mode: "translate",
targetFormat: FORMATS.OPENAI,
sourceFormat: FORMATS.CLAUDE,
provider: "openai",
model: "gpt-4.1-mini",
body: {
messages: [{ role: "user", content: "hello" }],
},
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.equal((text.match(/event: message_start/g) || []).length, 1);
assert.match(text, /event: content_block_start/);
assert.match(text, /event: content_block_delta/);
assert.match(text, /event: message_delta/);
assert.match(text, /event: message_stop/);
assert.match(text, /\[Proxy Error\] The upstream API returned an empty response/);
assert.ok(text.indexOf("event: content_block_start") > text.indexOf("event: message_start"));
assert.ok(text.indexOf("event: message_delta") > text.indexOf("event: content_block_stop"));
assert.equal(
onCompletePayload.responseBody.choices[0].message.content,
SYNTHETIC_CLAUDE_EMPTY_RESPONSE_TEXT
);
assert.equal(onCompletePayload.responseBody.usage.total_tokens, 3);
});
test("createSSETransformStreamWithLogger flushes a trailing Claude usage event without a newline", async () => {
let onCompletePayload = null;
const text = await readWithTransform(
[
`data: ${JSON.stringify({
type: "message_start",
message: {
id: "msg_tail",
model: "claude-sonnet-4",
role: "assistant",
usage: { input_tokens: 3 },
},
})}\n\n`,
`data: ${JSON.stringify({
type: "content_block_start",
index: 0,
content_block: { type: "text", text: "" },
})}\n\n`,
`data: ${JSON.stringify({
type: "content_block_delta",
index: 0,
delta: { type: "text_delta", text: "Buffered tail" },
})}\n\n`,
`data: ${JSON.stringify({
type: "message_delta",
delta: { stop_reason: "end_turn" },
usage: { output_tokens: 5 },
})}`,
],
createSSETransformStreamWithLogger(
FORMATS.CLAUDE,
FORMATS.OPENAI,
"claude",
null,
null,
"claude-sonnet-4",
null,
{ messages: [{ role: "user", content: "hello" }] },
(payload) => {
onCompletePayload = payload;
}
)
);
assert.match(text, /Buffered tail/);
assert.match(text, /\[DONE\]/);
assert.equal(onCompletePayload.responseBody.choices[0].message.content, "Buffered tail");
assert.equal(onCompletePayload.responseBody.usage.completion_tokens, 5);
assert.equal(onCompletePayload.responseBody.usage.total_tokens, 5);
});
test("buildStreamSummaryFromEvents compacts Responses API deltas into a synthetic response", () => {
const summary = buildStreamSummaryFromEvents(
[
{ index: 0, data: { type: "response.output_text.delta", delta: "Hello " } },
{ index: 1, data: { type: "response.output_text.delta", delta: "world" } },
{
index: 2,
data: {
type: "response.output_text.done",
usage: { input_tokens: 2, output_tokens: 3, total_tokens: 5 },
},
},
],
FORMATS.OPENAI_RESPONSES,
"gpt-4.1-mini"
);
assert.equal((summary as any).object, "response");
assert.equal((summary as any).model, "gpt-4.1-mini");
assert.equal((summary as any).output[0].content[0].text, "Hello world");
assert.deepEqual((summary as any).usage, { input_tokens: 2, output_tokens: 3, total_tokens: 5 });
});
test("buildStreamSummaryFromEvents preserves Gemini thought parts and function calls", () => {
const summary = buildStreamSummaryFromEvents(
[
{
index: 0,
data: {
modelVersion: "gemini-2.5-pro",
candidates: [
{
content: {
role: "model",
parts: [
{ text: "Thinking", thought: true },
{ text: " aloud", thought: true },
],
},
},
],
},
},
{
index: 1,
data: {
candidates: [
{
content: {
role: "model",
parts: [
{ text: "Done." },
{ functionCall: { name: "read_file", args: { path: "/tmp/a" } } },
],
},
finishReason: "STOP",
},
],
usageMetadata: {
promptTokenCount: 4,
candidatesTokenCount: 5,
totalTokenCount: 9,
},
},
},
],
FORMATS.GEMINI,
"gemini-2.5-pro"
);
assert.equal((summary as any).modelVersion, "gemini-2.5-pro");
assert.equal((summary as any).candidates[0].content.parts[0].text, "Thinking aloud");
assert.equal((summary as any).candidates[0].content.parts[0].thought, true);
assert.deepEqual((summary as any).candidates[0].content.parts[2], {
functionCall: { name: "read_file", args: { path: "/tmp/a" } },
});
assert.deepEqual((summary as any).usageMetadata, {
promptTokenCount: 4,
candidatesTokenCount: 5,
totalTokenCount: 9,
});
});
test("compactStructuredStreamPayload wraps primitive summaries with Omniroute stream metadata", () => {
const compact = compactStructuredStreamPayload({
_streamed: true,
_format: "sse-json",
_stage: "client_response",
_eventCount: 2,
summary: "done",
});
assert.deepEqual(compact, {
summary: "done",
_omniroute_stream: {
format: "sse-json",
stage: "client_response",
eventCount: 2,
},
});
});
test("createSSETransformStreamWithLogger flushes Responses API terminal events on stream end", async () => {
const text = await readWithTransform(
[
`data: ${JSON.stringify({
id: "chatcmpl_flush",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: { role: "assistant", content: "Hello" } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_flush",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
usage: { prompt_tokens: 2, completion_tokens: 3, total_tokens: 5 },
})}\n\n`,
],
createSSETransformStreamWithLogger(
FORMATS.OPENAI,
FORMATS.OPENAI_RESPONSES,
"openai",
null,
null,
"gpt-4.1-mini",
null,
{ messages: [{ role: "user", content: "hello" }] }
)
);
assert.match(text, /response\.created/);
assert.match(text, /response\.completed/);
assert.doesNotMatch(text, /\[DONE\]/);
});
test("createPassthroughStreamWithLogger reuses passthrough mode helpers", async () => {
const text = await readWithTransform(
[
`data: ${JSON.stringify({
id: "chatcmpl_passthrough",
object: "chat.completion.chunk",
created: 1,
model: "gpt-4.1-mini",
choices: [{ index: 0, delta: { role: "assistant", content: "Hello again" } }],
})}\n\n`,
"data: [DONE]\n\n",
],
createPassthroughStreamWithLogger("openai", null, null, "gpt-4.1-mini", null, {
messages: [{ role: "user", content: "hello" }],
})
);
assert.match(text, /Hello again/);
assert.match(text, /\[DONE\]/);
});
test("createStructuredSSECollector drops excess events and compactStructuredStreamPayload preserves metadata for object summaries", () => {
const collector = createStructuredSSECollector({
stage: "client_response",
maxEvents: 1,
maxBytes: 512,
});
collector.push({ type: "response.output_text.delta", delta: "one" });
collector.push({ type: "response.output_text.delta", delta: "two" });
const built = collector.build(
{
object: "response",
status: "completed",
},
{ includeEvents: false }
);
const compact = compactStructuredStreamPayload(built);
assert.equal(built._truncated, true);
assert.equal(built._droppedEvents, 1);
assert.equal(built._eventCount, 2);
assert.deepEqual(compact, {
object: "response",
status: "completed",
_omniroute_stream: {
format: "sse-json",
stage: "client_response",
eventCount: 2,
truncated: true,
droppedEvents: 1,
},
});
});
test("createSSEStream passthrough drops keepalive event blocks without losing Responses deltas", async () => {
const text = await readTransformed(
[
"event: keepalive\ndata:\n\n",
`data: ${JSON.stringify({
type: "response.output_text.delta",
delta: "Hello keepalive-safe",
})}\n\n`,
`data: ${JSON.stringify({
type: "response.completed",
response: {
id: "resp_keepalive",
object: "response",
model: "gpt-4.1-mini",
status: "completed",
usage: { input_tokens: 2, output_tokens: 1, total_tokens: 3 },
},
})}\n\n`,
"data: [DONE]\n\n",
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI_RESPONSES,
provider: "openai",
model: "gpt-4.1-mini",
body: { input: "hello" },
}
);
assert.equal(text.includes("event: keepalive"), false);
assert.equal(text.includes("data:\n\n"), false);
assert.match(text, /response\.output_text\.delta/);
assert.match(text, /Hello keepalive-safe/);
assert.match(text, /data: \[DONE\]/);
});
test("createSSEStream passthrough aborts on Responses usage-limit failures and reports 429", async () => {
let failurePayload = null;
await assert.rejects(
readTransformed(
[
`data: ${JSON.stringify({
type: "response.failed",
response: {
id: "resp_usage_limit",
object: "response",
model: "gpt-5.5",
status: "failed",
error: {
code: "usage_limit_reached",
message: "Your weekly usage limit has been reached",
},
},
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI_RESPONSES,
provider: "codex",
model: "gpt-5.5",
body: { input: "hello" },
onFailure(payload) {
failurePayload = payload;
},
}
),
/weekly usage limit|Upstream failure/
);
assert.ok(failurePayload, "should report the stream failure before aborting");
assert.equal(failurePayload.status, 429);
assert.equal(failurePayload.code, "usage_limit_reached");
});
test("createRequestLogger skips disabled logs and caps retained stream chunk bytes", async () => {
const disabled = await createRequestLogger("openai", "openai", "gpt-test", {
enabled: false,
});
disabled.logClientRawRequest("/v1/chat/completions", { prompt: "hello" });
disabled.appendProviderChunk("x".repeat(32));
assert.equal(disabled.getPipelinePayloads(), null);
const logger = await createRequestLogger("openai", "openai", "gpt-test", {
enabled: true,
captureStreamChunks: true,
maxStreamChunkBytes: 5,
});
logger.appendProviderChunk("abcdef");
logger.appendProviderChunk("ghijkl");
const payloads = logger.getPipelinePayloads();
assert.deepEqual(payloads.streamChunks.provider, [
"abcde",
"[stream chunk log truncated after 5 bytes]",
]);
});
test("createRequestLogger caps retained stream chunk item count", async () => {
const logger = await createRequestLogger("openai", "openai", "gpt-test", {
enabled: true,
captureStreamChunks: true,
maxStreamChunkBytes: 1024,
maxStreamChunkItems: 2,
});
logger.appendProviderChunk("one");
logger.appendProviderChunk("two");
logger.appendProviderChunk("three");
const payloads = logger.getPipelinePayloads();
assert.deepEqual(payloads.streamChunks.provider, [
"one",
"[stream chunk log truncated after 2 chunks]",
]);
});
// T-VERIFY: passthrough mode failure decrements pending requests
// Regression test for missing trackPendingRequest(false) on passthrough failure
import { getPendingRequests, clearPendingRequests } from "../../src/lib/usage/usageHistory.ts";
test("createSSEStream passthrough mode decrements pending requests on failure", async () => {
// Clear any existing pending requests first
clearPendingRequests();
const initial = getPendingRequests();
assert.equal(Object.keys(initial.byModel).length, 0, "should start with no pending requests");
let failurePayload = null;
const testProvider = "openai-compatible-test-failure";
const testModel = "gpt-test";
const testConnectionId = "test-conn-123";
await assert.rejects(
readTransformed(
[
`data: ${JSON.stringify({
type: "response.failed",
response: {
id: "resp_failed_test",
object: "response",
model: testModel,
status: "failed",
error: {
code: "test_failure",
message: "Test failure for pending request tracking",
},
},
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI_RESPONSES,
provider: testProvider,
model: testModel,
connectionId: testConnectionId,
body: { input: "hello" },
onFailure(payload) {
failurePayload = payload;
},
}
),
/Test failure|Upstream failure/
);
assert.ok(failurePayload, "should report the stream failure");
// Verify pending requests are properly decremented after failure
const pending = getPendingRequests();
const modelKey = `${testModel} (${testProvider})`;
const count = pending.byModel[modelKey] || 0;
assert.equal(
count,
0,
`pending request count for ${modelKey} should be 0 after failure, got ${count}`
);
});
test("createSSEStream passthrough emits synthetic error chunk for empty choices array", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_empty",
object: "chat.completion.chunk",
created: 1,
model: "kimi-k2.6",
choices: [],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_empty",
object: "chat.completion.chunk",
created: 1,
model: "kimi-k2.6",
choices: [{ index: 0, delta: { role: "assistant", content: "Hello" } }],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_empty",
object: "chat.completion.chunk",
created: 1,
model: "kimi-k2.6",
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "opencode-go",
model: "kimi-k2.6",
body: { messages: [{ role: "user", content: "hello" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
// The empty choices chunk should have been replaced with a synthetic error chunk
assert.match(text, /\[OmniRoute\] Upstream returned an empty response/);
assert.match(text, /"finish_reason":"stop"/);
// Subsequent valid chunks should still be present
assert.match(text, /"content":"Hello"/);
assert.equal(onCompletePayload.status, 200);
});
test("createSSEStream passthrough logs empty response after tool_calls completion", async () => {
let onCompletePayload = null;
const text = await readTransformed(
[
`data: ${JSON.stringify({
id: "chatcmpl_tool_then_empty",
object: "chat.completion.chunk",
created: 1,
model: "gpt-5.5-xhigh",
choices: [
{
index: 0,
delta: {
tool_calls: [
{
index: 0,
id: "call_tc",
type: "function",
function: { name: "task_complete", arguments: '{}' },
},
],
},
},
],
})}\n\n`,
`data: ${JSON.stringify({
id: "chatcmpl_tool_then_empty",
object: "chat.completion.chunk",
created: 1,
model: "gpt-5.5-xhigh",
choices: [{ index: 0, delta: {}, finish_reason: "tool_calls" }],
})}\n\n`,
],
{
mode: "passthrough",
sourceFormat: FORMATS.OPENAI,
provider: "codex",
model: "gpt-5.5-xhigh",
body: { messages: [{ role: "user", content: "do task" }] },
onComplete(payload) {
onCompletePayload = payload;
},
}
);
assert.match(text, /"finish_reason":"tool_calls"/);
assert.equal(onCompletePayload.status, 200);
assert.equal(onCompletePayload.responseBody.choices[0].finish_reason, "tool_calls");
assert.equal(onCompletePayload.responseBody.choices[0].message.tool_calls[0].function.name, "task_complete");
// Content should be null (empty) since no text was generated
assert.equal(onCompletePayload.responseBody.choices[0].message.content, null);
});