From 54e2cc7aa09ff80cc5fe2ce843b6d2903cb0c7f3 Mon Sep 17 00:00:00 2001 From: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com> Date: Wed, 2 Sep 2026 08:20:20 -0300 Subject: [PATCH] test(streaming): isolate public error boundary fixture --- ...m-handler-public-error-boundary.fixture.ts | 211 +++++++++++++++++ ...ream-handler-public-error-boundary.test.ts | 219 ++++-------------- 2 files changed, 260 insertions(+), 170 deletions(-) create mode 100644 tests/fixtures/stream-handler-public-error-boundary.fixture.ts diff --git a/tests/fixtures/stream-handler-public-error-boundary.fixture.ts b/tests/fixtures/stream-handler-public-error-boundary.fixture.ts new file mode 100644 index 0000000000..de48e9bdf4 --- /dev/null +++ b/tests/fixtures/stream-handler-public-error-boundary.fixture.ts @@ -0,0 +1,211 @@ +// This suite owns process-wide DATA_DIR, plugin, logger, and DB state. It must run only inside +// the subprocess launched by tests/unit/stream-handler-public-error-boundary.test.ts. +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import test from "node:test"; + +const originalDataDir = process.env.DATA_DIR; +const originalPluginsDir = process.env.OMNIROUTE_PLUGINS_DIR; +const testRoot = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-stream-public-error-")); +const TEST_DATA_DIR = path.join(testRoot, "data"); +const TEST_PLUGINS_DIR = path.join(testRoot, "plugins"); +fs.mkdirSync(TEST_DATA_DIR, { recursive: true }); +fs.mkdirSync(TEST_PLUGINS_DIR, { recursive: true }); +process.env.DATA_DIR = TEST_DATA_DIR; +process.env.OMNIROUTE_PLUGINS_DIR = TEST_PLUGINS_DIR; + +const [core, callLogs, artifactWriter, loggerResource, streamHandler, { FORMATS }] = + await Promise.all([ + import("../../src/lib/db/core.ts"), + import("../../src/lib/usage/callLogs.ts"), + import("../../src/lib/usage/callLogArtifactWriter.ts"), + import("../../src/shared/utils/loggerResource.ts"), + import("../../open-sse/utils/streamHandler.ts"), + import("../../open-sse/translator/formats.ts"), + ]); +const { createStreamController, pipeWithDisconnect } = streamHandler; + +const SECRET = "sk-live-streamhandler-secret-123456"; +const API_KEY = "provider-key-streamhandler-654321"; +const PRIVATE_PATH = "/srv/omniroute/private/provider.ts:42:9"; +const RAW_MESSAGE = + `Upstream failed at ${PRIVATE_PATH} Authorization: Bearer ${SECRET} api_key=${API_KEY}` + + `\n at dispatch (/srv/omniroute/private/dispatcher.ts:88:3)`; + +test.after(async () => { + assert.equal(await callLogs.waitForCallLogSaves(3_000), true); + await artifactWriter.closeCallLogArtifactWriter(); + core.resetDbInstance(); + await loggerResource.closeSharedLoggerResource(); + + if (originalDataDir === undefined) delete process.env.DATA_DIR; + else process.env.DATA_DIR = originalDataDir; + if (originalPluginsDir === undefined) delete process.env.OMNIROUTE_PLUGINS_DIR; + else process.env.OMNIROUTE_PLUGINS_DIR = originalPluginsDir; + + fs.rmSync(testRoot, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); +}); + +test("fixture binds all persistent state to its process-owned directories", () => { + assert.equal(core.DATA_DIR, TEST_DATA_DIR); + assert.equal(core.SQLITE_FILE, path.join(TEST_DATA_DIR, "storage.sqlite")); + assert.equal(process.env.DATA_DIR, TEST_DATA_DIR); + assert.equal(process.env.OMNIROUTE_PLUGINS_DIR, TEST_PLUGINS_DIR); + assert.equal(fs.existsSync(TEST_DATA_DIR), true); + assert.equal(fs.existsSync(TEST_PLUGINS_DIR), true); +}); + +test("OpenAI stream failures keep raw diagnostics internal and sanitize the public wire", async () => { + const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 502 }); + const source = new ReadableStream({ + start(controller) { + controller.error(upstreamError); + }, + }); + let internalMessage = ""; + + const stream = pipeWithDisconnect( + new Response(source), + new TransformStream(), + createStreamController({ + clientResponseFormat: FORMATS.OPENAI, + onError(event) { + internalMessage = event.message; + return true; + }, + }), + { stallTimeoutMs: 0 } + ); + const publicWire = await new Response(stream).text(); + + assert.equal(internalMessage, RAW_MESSAGE, "failure classification must retain the raw message"); + assert.match(publicWire, /"finish_reason":"error"/); + assert.match(publicWire, /"code":"server_error"/); + assert.match(publicWire, /\[DONE\]/); + assert.doesNotMatch(publicWire, new RegExp(SECRET)); + assert.doesNotMatch(publicWire, new RegExp(API_KEY)); + assert.doesNotMatch(publicWire, /\/srv\/omniroute\/private/); + assert.doesNotMatch(publicWire, /dispatcher\.ts/); + assert.match(publicWire, /Authorization: \[REDACTED\]/); + assert.match(publicWire, //); +}); + +test("Responses stream failures preserve the failure event shape without leaking diagnostics", async () => { + const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 429 }); + const source = new ReadableStream({ + start(controller) { + controller.error(upstreamError); + }, + }); + let internalError: unknown; + + const stream = pipeWithDisconnect( + new Response(source), + new TransformStream(), + createStreamController({ + clientResponseFormat: FORMATS.OPENAI_RESPONSES, + onError(event) { + internalError = event.error; + return true; + }, + }), + { stallTimeoutMs: 0 } + ); + const publicWire = await new Response(stream).text(); + + assert.equal(internalError, upstreamError, "the original error object must reach classification"); + assert.match(publicWire, /event: response\.failed/); + assert.match(publicWire, /"type":"response\.failed"/); + assert.match(publicWire, /"type":"rate_limit_error"/); + assert.match(publicWire, /"code":"rate_limit_exceeded"/); + assert.doesNotMatch(publicWire, new RegExp(SECRET)); + assert.doesNotMatch(publicWire, new RegExp(API_KEY)); + assert.doesNotMatch(publicWire, /\/srv\/omniroute\/private/); + assert.doesNotMatch(publicWire, /dispatcher\.ts/); + assert.match(publicWire, /Authorization: \[REDACTED\]/); + assert.match(publicWire, //); +}); + +test("Claude stream failures preserve error and stop events without leaking diagnostics", async () => { + const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 403 }); + const source = new ReadableStream({ + start(controller) { + controller.error(upstreamError); + }, + }); + let internalStatusCode = 0; + + const stream = pipeWithDisconnect( + new Response(source), + new TransformStream(), + createStreamController({ + clientResponseFormat: FORMATS.CLAUDE, + onError(event) { + internalStatusCode = event.statusCode; + return true; + }, + }), + { stallTimeoutMs: 0 } + ); + const publicWire = await new Response(stream).text(); + + assert.equal(internalStatusCode, 403); + assert.match(publicWire, /event: error/); + assert.match(publicWire, /"type":"permission_error"/); + assert.match(publicWire, /event: message_stop/); + assert.doesNotMatch(publicWire, new RegExp(SECRET)); + assert.doesNotMatch(publicWire, new RegExp(API_KEY)); + assert.doesNotMatch(publicWire, /\/srv\/omniroute\/private/); + assert.doesNotMatch(publicWire, /dispatcher\.ts/); + assert.match(publicWire, /Authorization: \[REDACTED\]/); + assert.match(publicWire, //); +}); + +test("stream diagnostics sanitize logs while callbacks retain the original failure", () => { + const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 502 }); + const originalLog = console.log; + const logLines: string[] = []; + let internalError: unknown; + console.log = (...args: unknown[]) => { + logLines.push(args.map(String).join(" ")); + }; + + try { + createStreamController({ + provider: "test-provider", + model: "test-model", + onError(event) { + internalError = event.error; + return true; + }, + }).handleError(upstreamError); + } finally { + console.log = originalLog; + } + + const logs = logLines.join("\n"); + assert.equal(internalError, upstreamError); + assert.match(logs, /error: Upstream failed at /); + assert.match(logs, /Authorization: \[REDACTED\]/); + assert.doesNotMatch(logs, new RegExp(SECRET)); + assert.doesNotMatch(logs, new RegExp(API_KEY)); + assert.doesNotMatch(logs, /\/srv\/omniroute\/private/); + assert.doesNotMatch(logs, /dispatcher\.ts/); +}); + +test("client disconnects stay outside the provider-failure callback", () => { + let providerFailureRecorded = false; + const controller = createStreamController({ + onError() { + providerFailureRecorded = true; + return true; + }, + }); + + controller.handleError(new DOMException("request_signal_aborted", "AbortError")); + + assert.equal(providerFailureRecorded, false); + assert.equal(controller.signal.aborted, false); +}); diff --git a/tests/unit/stream-handler-public-error-boundary.test.ts b/tests/unit/stream-handler-public-error-boundary.test.ts index 607c1068a3..41aa796b19 100644 --- a/tests/unit/stream-handler-public-error-boundary.test.ts +++ b/tests/unit/stream-handler-public-error-boundary.test.ts @@ -1,183 +1,62 @@ import assert from "node:assert/strict"; -import fs from "node:fs"; -import os from "node:os"; -import path from "node:path"; +import { spawnSync } from "node:child_process"; +import { fileURLToPath } from "node:url"; import test from "node:test"; -const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-stream-public-error-")); -const TEST_PLUGINS_DIR = fs.mkdtempSync( - path.join(os.tmpdir(), "omniroute-stream-public-error-plugins-") +const REPO_ROOT = fileURLToPath(new URL("../..", import.meta.url)); +const FIXTURE = fileURLToPath( + new URL("../fixtures/stream-handler-public-error-boundary.fixture.ts", import.meta.url) ); -process.env.DATA_DIR = TEST_DATA_DIR; -process.env.OMNIROUTE_PLUGINS_DIR = TEST_PLUGINS_DIR; -const core = await import("../../src/lib/db/core.ts"); -const { createStreamController, pipeWithDisconnect } = - await import("../../open-sse/utils/streamHandler.ts"); -const { FORMATS } = await import("../../open-sse/translator/formats.ts"); +const CHILD_RUNTIME_ENV_KEYS = [ + "PATH", + "TMPDIR", + "TMP", + "TEMP", + "SystemRoot", + "ComSpec", + "PATHEXT", + "LANG", + "LC_ALL", + "TZ", +] as const; -const SECRET = "sk-live-streamhandler-secret-123456"; -const API_KEY = "provider-key-streamhandler-654321"; -const PRIVATE_PATH = "/srv/omniroute/private/provider.ts:42:9"; -const RAW_MESSAGE = - `Upstream failed at ${PRIVATE_PATH} Authorization: Bearer ${SECRET} api_key=${API_KEY}` + - `\n at dispatch (/srv/omniroute/private/dispatcher.ts:88:3)`; - -test.after(() => { - core.resetDbInstance(); - fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); - fs.rmSync(TEST_PLUGINS_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); -}); - -test("OpenAI stream failures keep raw diagnostics internal and sanitize the public wire", async () => { - const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 502 }); - const source = new ReadableStream({ - start(controller) { - controller.error(upstreamError); - }, - }); - let internalMessage = ""; - - const stream = pipeWithDisconnect( - new Response(source), - new TransformStream(), - createStreamController({ - clientResponseFormat: FORMATS.OPENAI, - onError(event) { - internalMessage = event.message; - return true; - }, - }), - { stallTimeoutMs: 0 } - ); - const publicWire = await new Response(stream).text(); - - assert.equal(internalMessage, RAW_MESSAGE, "failure classification must retain the raw message"); - assert.match(publicWire, /"finish_reason":"error"/); - assert.match(publicWire, /"code":"server_error"/); - assert.match(publicWire, /\[DONE\]/); - assert.doesNotMatch(publicWire, new RegExp(SECRET)); - assert.doesNotMatch(publicWire, new RegExp(API_KEY)); - assert.doesNotMatch(publicWire, /\/srv\/omniroute\/private/); - assert.doesNotMatch(publicWire, /dispatcher\.ts/); - assert.match(publicWire, /Authorization: \[REDACTED\]/); - assert.match(publicWire, //); -}); - -test("Responses stream failures preserve the failure event shape without leaking diagnostics", async () => { - const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 429 }); - const source = new ReadableStream({ - start(controller) { - controller.error(upstreamError); - }, - }); - let internalError: unknown; - - const stream = pipeWithDisconnect( - new Response(source), - new TransformStream(), - createStreamController({ - clientResponseFormat: FORMATS.OPENAI_RESPONSES, - onError(event) { - internalError = event.error; - return true; - }, - }), - { stallTimeoutMs: 0 } - ); - const publicWire = await new Response(stream).text(); - - assert.equal(internalError, upstreamError, "the original error object must reach classification"); - assert.match(publicWire, /event: response\.failed/); - assert.match(publicWire, /"type":"response\.failed"/); - assert.match(publicWire, /"type":"rate_limit_error"/); - assert.match(publicWire, /"code":"rate_limit_exceeded"/); - assert.doesNotMatch(publicWire, new RegExp(SECRET)); - assert.doesNotMatch(publicWire, new RegExp(API_KEY)); - assert.doesNotMatch(publicWire, /\/srv\/omniroute\/private/); - assert.doesNotMatch(publicWire, /dispatcher\.ts/); - assert.match(publicWire, /Authorization: \[REDACTED\]/); - assert.match(publicWire, //); -}); - -test("Claude stream failures preserve error and stop events without leaking diagnostics", async () => { - const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 403 }); - const source = new ReadableStream({ - start(controller) { - controller.error(upstreamError); - }, - }); - let internalStatusCode = 0; - - const stream = pipeWithDisconnect( - new Response(source), - new TransformStream(), - createStreamController({ - clientResponseFormat: FORMATS.CLAUDE, - onError(event) { - internalStatusCode = event.statusCode; - return true; - }, - }), - { stallTimeoutMs: 0 } - ); - const publicWire = await new Response(stream).text(); - - assert.equal(internalStatusCode, 403); - assert.match(publicWire, /event: error/); - assert.match(publicWire, /"type":"permission_error"/); - assert.match(publicWire, /event: message_stop/); - assert.doesNotMatch(publicWire, new RegExp(SECRET)); - assert.doesNotMatch(publicWire, new RegExp(API_KEY)); - assert.doesNotMatch(publicWire, /\/srv\/omniroute\/private/); - assert.doesNotMatch(publicWire, /dispatcher\.ts/); - assert.match(publicWire, /Authorization: \[REDACTED\]/); - assert.match(publicWire, //); -}); - -test("stream diagnostics sanitize logs while callbacks retain the original failure", () => { - const upstreamError = Object.assign(new Error(RAW_MESSAGE), { statusCode: 502 }); - const originalLog = console.log; - const logLines: string[] = []; - let internalError: unknown; - console.log = (...args: unknown[]) => { - logLines.push(args.map(String).join(" ")); +function buildFixtureEnv(): NodeJS.ProcessEnv { + const env: NodeJS.ProcessEnv = { + NODE_ENV: "test", + APP_LOG_TO_FILE: "false", + API_KEY_SECRET: "stream-handler-boundary-fixture-secret-20260902", + DISABLE_SQLITE_AUTO_BACKUP: "true", + NO_COLOR: "1", }; - try { - createStreamController({ - provider: "test-provider", - model: "test-model", - onError(event) { - internalError = event.error; - return true; - }, - }).handleError(upstreamError); - } finally { - console.log = originalLog; + for (const key of CHILD_RUNTIME_ENV_KEYS) { + const value = process.env[key]; + if (value !== undefined) env[key] = value; } - const logs = logLines.join("\n"); - assert.equal(internalError, upstreamError); - assert.match(logs, /error: Upstream failed at /); - assert.match(logs, /Authorization: \[REDACTED\]/); - assert.doesNotMatch(logs, new RegExp(SECRET)); - assert.doesNotMatch(logs, new RegExp(API_KEY)); - assert.doesNotMatch(logs, /\/srv\/omniroute\/private/); - assert.doesNotMatch(logs, /dispatcher\.ts/); -}); + // Nested test runners must not inherit the parent runner's recursion marker. + delete env.NODE_TEST_CONTEXT; + return env; +} -test("client disconnects stay outside the provider-failure callback", () => { - let providerFailureRecorded = false; - const controller = createStreamController({ - onError() { - providerFailureRecorded = true; - return true; - }, - }); +test("generic stream public error boundaries pass in an isolated process", () => { + const result = spawnSync( + process.execPath, + ["--import", "tsx/esm", "--import", "./open-sse/utils/setupPolyfill.ts", "--test", FIXTURE], + { + cwd: REPO_ROOT, + encoding: "utf8", + env: buildFixtureEnv(), + timeout: 120_000, + } + ); + const output = `${result.stdout}\n${result.stderr}`; - controller.handleError(new DOMException("request_signal_aborted", "AbortError")); - - assert.equal(providerFailureRecorded, false); - assert.equal(controller.signal.aborted, false); + assert.ifError(result.error); + assert.equal(result.signal, null, output.slice(-12_000)); + assert.equal(result.status, 0, output.slice(-12_000)); + assert.match(output, /(?:^|\s)tests\s+6(?:\s|$)/m); + assert.match(output, /(?:^|\s)pass\s+6(?:\s|$)/m); + assert.match(output, /(?:^|\s)fail\s+0(?:\s|$)/m); });