From 9469fa5f594593d0320974248a7939e4b8aebf86 Mon Sep 17 00:00:00 2001 From: Diego Rodrigues de Sa e Souza Date: Thu, 3 Sep 2026 20:58:56 -0300 Subject: [PATCH] fix(streaming): sanitize generic stream failure boundaries (#12457) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Validado em lote numa worktree combinada com os 14 PRs desta campanha de error-boundary sobre o tip de `release/v3.8.51`: `typecheck:core` limpo e **120/120** nos 23 arquivos de teste que os PRs trazem. Um ponto que só apareceu no tree combinado: **#12465 e #12466 criam o mesmo arquivo novo** `open-sse/utils/streamReadiness.ts` (que não existe no tip) com desenhos divergentes de cancelamento — `cancelled` + `releaseLock` imediato num, `readInFlight`/`cancelRequested` com `cancelReader` fire-and-forget no outro. Adotei a versão do #12466, que difere e defere o release do lock para quando a leitura em voo termina, e validei a escolha rodando as suítes dos **dois** PRs contra ela: 21/21 no readiness compartilhado e 22/22 incluindo o boundary do Perplexity. --- CHANGELOG.md | 4 + open-sse/utils/streamHandler.ts | 14 +- ...m-handler-public-error-boundary.fixture.ts | 211 ++++++++++++++++++ ...ream-handler-public-error-boundary.test.ts | 62 +++++ tests/unit/stream-handler.test.ts | 17 +- 5 files changed, 296 insertions(+), 12 deletions(-) create mode 100644 tests/fixtures/stream-handler-public-error-boundary.fixture.ts create mode 100644 tests/unit/stream-handler-public-error-boundary.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index bfcad514f4..80abce64dc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -98,6 +98,10 @@ _Living section — cycle opened at the v3.8.50 freeze (parallel-cycle model). B ### 🐛 Bug Fixes +- **security(streaming):** sanitize generic mid-stream error messages before emitting OpenAI, + Responses, or Claude SSE failure frames and before diagnostic logging, while preserving raw + failures for internal classification and keeping client disconnects out of provider failure state. + ### 📝 Maintenance --- diff --git a/open-sse/utils/streamHandler.ts b/open-sse/utils/streamHandler.ts index 7776f2e5e9..841346bcfe 100644 --- a/open-sse/utils/streamHandler.ts +++ b/open-sse/utils/streamHandler.ts @@ -1,6 +1,7 @@ import { trackPendingRequest } from "@/lib/usageDb"; import { STREAM_IDLE_TIMEOUT_MS } from "../config/constants.ts"; import { FORMATS } from "../translator/formats.ts"; +import { buildErrorBody } from "./error.ts"; import { PENDING_REQUEST_CLEARED_MARKER } from "./stream.ts"; import { createCompletedResponsesToolHandoffWatcher } from "./responsesToolHandoff.ts"; import { createStreamContentWatcher, type StreamContentWatcher } from "./streamReadiness.ts"; @@ -187,6 +188,10 @@ function getErrorStatusCode(error: unknown): number { return 502; } +function getPublicErrorMessage(errorMsg: string, statusCode: number): string { + return buildErrorBody(statusCode, errorMsg).error.message; +} + function isDeadlineAbortReason(reason: unknown): reason is Error { return ( reason instanceof Error && @@ -406,7 +411,7 @@ export function createStreamController({ } if (error instanceof Error) { - logStream(`error: ${error.message}`); + logStream(`error: ${getPublicErrorMessage(error.message, getErrorStatusCode(error))}`); return; } logStream("error: unknown"); @@ -452,6 +457,7 @@ export function buildStreamErrorChunks( clientResponseFormat?: string | null ) { const statusMapping = getStreamErrorStatusMapping(statusCode); + const publicErrorMessage = getPublicErrorMessage(errorMsg, statusCode); if (isResponsesClientFormat(clientResponseFormat)) { const errorEvent = { @@ -460,7 +466,7 @@ export function buildStreamErrorChunks( id: null, status: "failed", error: { - message: errorMsg, + message: publicErrorMessage, type: statusMapping.responses.type, code: statusMapping.responses.code, }, @@ -475,7 +481,7 @@ export function buildStreamErrorChunks( type: "error", error: { type: statusMapping.claude.type, - message: errorMsg, + message: publicErrorMessage, }, }; @@ -498,7 +504,7 @@ export function buildStreamErrorChunks( }, ], error: { - message: errorMsg, + message: publicErrorMessage, type: statusMapping.responses.type, code: statusMapping.responses.code, }, 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 new file mode 100644 index 0000000000..41aa796b19 --- /dev/null +++ b/tests/unit/stream-handler-public-error-boundary.test.ts @@ -0,0 +1,62 @@ +import assert from "node:assert/strict"; +import { spawnSync } from "node:child_process"; +import { fileURLToPath } from "node:url"; +import test from "node:test"; + +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) +); + +const CHILD_RUNTIME_ENV_KEYS = [ + "PATH", + "TMPDIR", + "TMP", + "TEMP", + "SystemRoot", + "ComSpec", + "PATHEXT", + "LANG", + "LC_ALL", + "TZ", +] as const; + +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", + }; + + for (const key of CHILD_RUNTIME_ENV_KEYS) { + const value = process.env[key]; + if (value !== undefined) env[key] = value; + } + + // Nested test runners must not inherit the parent runner's recursion marker. + delete env.NODE_TEST_CONTEXT; + return env; +} + +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}`; + + 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); +}); diff --git a/tests/unit/stream-handler.test.ts b/tests/unit/stream-handler.test.ts index 8c19c08802..c358ff7c2a 100644 --- a/tests/unit/stream-handler.test.ts +++ b/tests/unit/stream-handler.test.ts @@ -256,7 +256,8 @@ test("createDisconnectAwareStream emits Responses API failure events for Respons assert.match(text, /event: response\.failed/); assert.match(text, /"type":"response\.failed"/); - assert.match(text, /"message":"responses stream\\ndied"/); + assert.match(text, /"message":"responses stream"/); + assert.doesNotMatch(text, /died/); assert.match(text, /"type":"server_error"/); assert.match(text, /"code":"server_error"/); assert.doesNotMatch(text, /chat\.completion\.chunk/); @@ -264,7 +265,7 @@ test("createDisconnectAwareStream emits Responses API failure events for Respons assert.doesNotMatch(text, /\[DONE\]/); }); -test("createDisconnectAwareStream keeps newlines escaped inside SSE data fields", async () => { +test("createDisconnectAwareStream strips multiline diagnostic tails from Responses errors", async () => { const upstreamError = Object.assign(new Error("line one\nline two\rline three"), { statusCode: 400, }); @@ -290,9 +291,9 @@ test("createDisconnectAwareStream keeps newlines escaped inside SSE data fields" const text = await readStreamText(stream); assert.match(text, /^event: response\.failed\ndata: \{"type":"response\.failed"/); - assert.match(text, /"message":"line one\\nline two\\rline three"/); - assert.doesNotMatch(text, /^line two/m); - assert.doesNotMatch(text, /^line three/m); + assert.match(text, /"message":"line one"/); + assert.doesNotMatch(text, /line two/); + assert.doesNotMatch(text, /line three/); }); test("createDisconnectAwareStream treats legacy OpenAI response format alias as Responses", async () => { @@ -360,7 +361,7 @@ test("createDisconnectAwareStream emits Claude SSE errors for Claude clients", a assert.doesNotMatch(text, /\[DONE\]/); }); -test("createDisconnectAwareStream keeps newlines escaped for Claude SSE errors", async () => { +test("createDisconnectAwareStream strips multiline diagnostic tails from Claude errors", async () => { const upstreamError = Object.assign(new Error("claude line one\nclaude line two"), { statusCode: 502, }); @@ -386,8 +387,8 @@ test("createDisconnectAwareStream keeps newlines escaped for Claude SSE errors", const text = await readStreamText(stream); assert.match(text, /^event: error\ndata: \{"type":"error"/); - assert.match(text, /"message":"claude line one\\nclaude line two"/); - assert.doesNotMatch(text, /^claude line two/m); + assert.match(text, /"message":"claude line one"/); + assert.doesNotMatch(text, /claude line two/); }); // #7699/#7816 — heuristic is scoped to FORMATS.CLAUDE (/v1/messages); a