diff --git a/CHANGELOG.md b/CHANGELOG.md index fbd429a96d..0f953af88f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,10 @@ - **ci:** heavy-pipeline dedup ([#6215](https://github.com/diegosouzapw/OmniRoute/pull/6215)) — the release-PR pipeline ran the unit suite 4× per sync (95 jobs, 208 machine-min; the v3.8.44 cycle fired 123 such runs, 88 cancelled). Now: Node 24/26 compat matrices move to a daily `nightly-compat.yml` (−28%/run; resolves the active release branch, opens a tracking issue on failure), coverage is collected inside the unit shards themselves via c8/`NODE_V8_COVERAGE` (−18%/run; the Coverage Shard ×8 matrix is gone — nodejs/node's own CI pattern), the ~40-job per-language i18n matrix becomes 1 job (the account has 20 concurrent-job slots total), and heavy jobs skip **draft** PRs — paired with `/generate-release` now opening the living release PR as draft (flipped ready at the new Phase 0a.0a), killing the per-merge churn for the whole cycle. Validated by a full `workflow_dispatch` of the new pipeline: 35 jobs, 0 failures, 23 min, merged coverage 80.16% (> ratchet baseline). - **feat(quality):** no-new-warnings per PR ([#6218](https://github.com/diegosouzapw/OmniRoute/pull/6218)) — native ESLint bulk suppressions (≥9.24) freeze the pre-existing debt (476 files / 4,273 violations in `config/quality/eslint-suppressions.json`); `npm run lint`, lint-staged (pre-commit) and a new fork-aware `lint-guard` job in quality.yml all run suppressions-aware, so a NEW warning goes red in the PR that introduces it instead of accruing invisibly (+41/+88 per cycle) and being blind-rebaselined at release. 3 warn rules promoted to error in `src/**` (`react-hooks/exhaustive-deps`, `@next/next/no-img-element`, `import/no-anonymous-default-export`); `collect-metrics` measures under the frozen baseline (ratchet metric = net-NEW debt; baseline tightened 4,279→0 in-PR per require-tighten); fork PRs run report-only (contributors are never blocked — the maintainer campaigns fix via co-authorship). Baseline stock shrinks via `--prune-suppressions` at release reconciliation. +### 🔧 Bug Fixes + +- fix(sse): Responses API passthrough now drops internal commentary-phase output before forwarding to clients (gated by RESPONSES_PASSTHROUGH_DROP_COMMENTARY, default on) (#6199) + --- ## [3.8.44] — TBD diff --git a/open-sse/handlers/responseSanitizer.ts b/open-sse/handlers/responseSanitizer.ts index 04d25871a6..f45ae5d3fe 100644 --- a/open-sse/handlers/responseSanitizer.ts +++ b/open-sse/handlers/responseSanitizer.ts @@ -554,6 +554,23 @@ function normalizeResponsesId(id: unknown): string { return `resp_${id}`; } +/** + * True when a Responses output item is an assistant `message` in the internal + * `commentary` phase — i.e. reasoning/scratchpad text that must never reach the + * client. Streaming `response.output_text.delta` events do not carry the `phase` + * themselves, so the passthrough path uses this on the `response.output_item.added` + * item to decide which subsequent deltas/dones to drop statefully (#6199). + */ +export function isResponsesCommentaryMessageItem(item: unknown): boolean { + const itemRecord = toRecord(item); + if (!itemRecord) return false; + const type = toString(itemRecord.type) || "message"; + if (type !== "message") return false; + const role = toString(itemRecord.role) || "assistant"; + const phase = toString(itemRecord.phase); + return role === "assistant" && phase === "commentary"; +} + function sanitizeResponsesStreamingOutputItem(item: unknown): JsonRecord | null { const itemRecord = toRecord(item); if (!itemRecord) return null; @@ -562,8 +579,7 @@ function sanitizeResponsesStreamingOutputItem(item: unknown): JsonRecord | null if (type === "message") { const role = toString(itemRecord.role) || "assistant"; - const phase = toString(itemRecord.phase); - if (role === "assistant" && phase === "commentary") { + if (isResponsesCommentaryMessageItem(itemRecord)) { return null; } diff --git a/open-sse/utils/stream.ts b/open-sse/utils/stream.ts index 8357ce2506..120644ac9d 100644 --- a/open-sse/utils/stream.ts +++ b/open-sse/utils/stream.ts @@ -30,7 +30,9 @@ import { STREAM_IDLE_TIMEOUT_MS, FETCH_BODY_TIMEOUT_MS, HTTP_STATUS } from "../c import { OMIT_STREAMING_CHUNK_MARKER, sanitizeStreamingChunk, + isResponsesCommentaryMessageItem, } from "../handlers/responseSanitizer.ts"; +import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags"; import { buildErrorBody } from "./error.ts"; import { parseTextualToolCallCandidate, isValidToolCallHeaderPrefix } from "./textualToolCall.ts"; import { recordToolLatency } from "../services/toolLatencyTracker.ts"; @@ -119,6 +121,12 @@ type StreamOptions = { copilotCompatibleReasoning?: boolean; /** Suppress the `` close marker for clients that render it verbatim (#5245). */ suppressThinkClose?: boolean; + /** + * Drop internal commentary-phase output items from Responses API passthrough + * streams before forwarding (#6199). When omitted, falls back to the + * `RESPONSES_PASSTHROUGH_DROP_COMMENTARY` feature flag (default on). + */ + dropResponsesCommentary?: boolean; provider?: string | null; reqLogger?: StreamLogger | null; toolNameMap?: unknown; @@ -621,9 +629,16 @@ export function createSSEStream(options: StreamOptions = {}) { body = null, onComplete = null, onFailure = null, + dropResponsesCommentary, } = options; const signatureNamespace = connectionId; + // Drop internal commentary-phase Responses output before forwarding (#6199). + // Explicit option wins; otherwise read the feature flag (default on). Resolved + // once per stream — never on the hot per-chunk path. + const shouldDropResponsesCommentary = + dropResponsesCommentary ?? isFeatureFlagEnabled("RESPONSES_PASSTHROUGH_DROP_COMMENTARY"); + const clientExpectsResponsesStream = (mode === STREAM_MODE.PASSTHROUGH ? clientResponseFormat === FORMATS.OPENAI_RESPONSES @@ -684,6 +699,12 @@ export function createSSEStream(options: StreamOptions = {}) { let passthroughResponsesId: string | null = null; let passthroughResponsesCurrentFunctionCallKey: string | null = null; const passthroughResponsesReasoningSummarySeen = new Set(); + // #6199 — commentary-phase items announced via `response.output_item.added` are + // internal. Their `response.output_text.delta`/`response.output_text.done`/ + // `response.output_item.done` events do not carry the `phase`, so we remember the + // item id + output_index here and drop every matching follow-up event. + const passthroughResponsesCommentaryItemIds = new Set(); + const passthroughResponsesCommentaryIndexes = new Set(); // #5786 — highest Responses-API `sequence_number` already forwarded on this stream. // The Responses API guarantees a strictly increasing sequence_number, so any event at // or below this watermark is an upstream reconnect/retry replay and must be dropped — @@ -1287,6 +1308,50 @@ export function createSSEStream(options: StreamOptions = {}) { parsed.type === "error"); if (isResponsesSSE) { + // #6199 — statefully drop internal commentary-phase output. The + // `response.output_item.added` announces the phase; the follow-up + // delta/done events only carry `item_id`/`output_index`, so we key + // off those. Happy-path (non-commentary) events are untouched. + if (shouldDropResponsesCommentary) { + const responsesEventType = parsed.type as string; + const eventOutputIndex = + typeof parsed.output_index === "number" ? parsed.output_index : null; + const eventItem = + parsed.item && typeof parsed.item === "object" && !Array.isArray(parsed.item) + ? (parsed.item as JsonRecord) + : null; + const eventItemId = + typeof parsed.item_id === "string" + ? parsed.item_id + : eventItem && typeof eventItem.id === "string" + ? eventItem.id + : null; + + if ( + responsesEventType === "response.output_item.added" && + isResponsesCommentaryMessageItem(parsed.item) + ) { + if (eventItemId) passthroughResponsesCommentaryItemIds.add(eventItemId); + if (eventOutputIndex !== null) + passthroughResponsesCommentaryIndexes.add(eventOutputIndex); + continue; + } + + const belongsToCommentary = + (eventItemId !== null && + passthroughResponsesCommentaryItemIds.has(eventItemId)) || + (eventOutputIndex !== null && + passthroughResponsesCommentaryIndexes.has(eventOutputIndex)); + if (belongsToCommentary) { + if (responsesEventType === "response.output_item.done") { + if (eventItemId) passthroughResponsesCommentaryItemIds.delete(eventItemId); + if (eventOutputIndex !== null) + passthroughResponsesCommentaryIndexes.delete(eventOutputIndex); + } + continue; + } + } + const responsesIdsNormalized = normalizeResponsesSseIds(parsed as JsonRecord); const parsedResponse = parsed.response && diff --git a/src/shared/constants/featureFlagDefinitions.ts b/src/shared/constants/featureFlagDefinitions.ts index 01be15921a..b42c26506a 100644 --- a/src/shared/constants/featureFlagDefinitions.ts +++ b/src/shared/constants/featureFlagDefinitions.ts @@ -234,7 +234,19 @@ export const FEATURE_FLAG_DEFINITIONS: FeatureFlagDefinition[] = [ warningLevel: "info", }, - // ──────────────── Runtime (12) ──────────────── + // ──────────────── Runtime (13) ──────────────── + { + key: "RESPONSES_PASSTHROUGH_DROP_COMMENTARY", + label: "Drop Responses Commentary", + description: + "Drop internal commentary-phase output items from Responses API passthrough streams before forwarding to clients. Disable to receive raw upstream commentary.", + descriptionI18nKey: "featureFlagResponsesPassthroughDropCommentaryDescription", + category: "runtime", + defaultValue: "true", + type: "boolean", + requiresRestart: false, + warningLevel: "info", + }, { key: "OMNIROUTE_MCP_ENFORCE_SCOPES", label: "MCP Enforce Scopes", diff --git a/tests/unit/feature-flags-settings.test.ts b/tests/unit/feature-flags-settings.test.ts index fccb08ee55..cfd7af226c 100644 --- a/tests/unit/feature-flags-settings.test.ts +++ b/tests/unit/feature-flags-settings.test.ts @@ -30,13 +30,13 @@ const { isControlPlaneProxyDirectFallbackEnabled, } = await import("../../src/shared/utils/featureFlags.ts"); -const EXPECTED_FEATURE_FLAG_COUNT = 40; +const EXPECTED_FEATURE_FLAG_COUNT = 41; // ────────────────────────────────────────────────────── // Test group 1 — Flag definitions registry // ────────────────────────────────────────────────────── describe("featureFlagDefinitions", () => { - it("has exactly 40 flag definitions", () => { + it("has exactly 41 flag definitions", () => { assert.strictEqual(FEATURE_FLAG_DEFINITIONS.length, EXPECTED_FEATURE_FLAG_COUNT); }); @@ -312,7 +312,7 @@ describe("resolveFeatureFlag", () => { }); describe("resolveAllFeatureFlags", () => { - it("returns all 40 flags", () => { + it("returns all 41 flags", () => { const all = resolveAllFeatureFlags(); assert.strictEqual(all.length, EXPECTED_FEATURE_FLAG_COUNT); }); diff --git a/tests/unit/responses-commentary-passthrough-6199.test.ts b/tests/unit/responses-commentary-passthrough-6199.test.ts new file mode 100644 index 0000000000..5346f46615 --- /dev/null +++ b/tests/unit/responses-commentary-passthrough-6199.test.ts @@ -0,0 +1,196 @@ +/** + * TDD test for fix(sse) #6199: Responses API passthrough leaks commentary-phase + * output text to clients. + * + * Background: #186 made the passthrough sanitizer format-aware and started SKIPPING + * the chat sanitizer for `response.*` events. Side effect: the streaming passthrough + * path never applied the commentary filter, so an assistant message item announced + * with `phase: "commentary"` (internal-only) had its `response.output_text.delta` + * chunks forwarded straight to the client. + * + * The commentary drop is STATEFUL: a `response.output_item.added` announcing a + * commentary item records its `output_index` / item id, then the matching + * `response.output_text.delta` / `response.output_item.done` events are dropped + * together (the delta events do not carry the `phase` themselves). + * + * Gated by the RESPONSES_PASSTHROUGH_DROP_COMMENTARY feature flag (default ON). The + * transform accepts an explicit `dropResponsesCommentary` boolean option so this test + * can exercise both the flag-on (drop) and flag-off (passthrough) behavior without + * touching env/DB state. + */ + +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-commentary-6199-")); +process.env.DATA_DIR = TEST_DATA_DIR; +const core = await import("../../src/lib/db/core.ts"); + +const { createSSEStream } = await import("../../open-sse/utils/stream.ts"); + +const textEncoder = new TextEncoder(); + +async function readTransformed(chunks: string[], options: object): Promise { + 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(); +} + +test.after(() => { + core.resetDbInstance(); + if (fs.existsSync(TEST_DATA_DIR)) { + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); + } +}); + +const COMMENTARY_TEXT = "internal chain-of-thought commentary that must stay hidden"; +const FINAL_TEXT = "The final answer visible to the user."; + +function sse(event: object): string { + return `data: ${JSON.stringify(event)}\n\n`; +} + +// A realistic Responses SSE sequence: a commentary item (index 0) followed by a +// real assistant answer item (index 1). +function buildResponsesStream(): string[] { + return [ + sse({ type: "response.created", response: { id: "resp_6199", output: [] } }), + // --- commentary item (internal, must be dropped when filtering) --- + sse({ + type: "response.output_item.added", + output_index: 0, + item: { + id: "msg_commentary", + type: "message", + role: "assistant", + phase: "commentary", + content: [], + }, + }), + sse({ + type: "response.output_text.delta", + output_index: 0, + item_id: "msg_commentary", + content_index: 0, + delta: COMMENTARY_TEXT, + }), + sse({ + type: "response.output_text.done", + output_index: 0, + item_id: "msg_commentary", + content_index: 0, + text: COMMENTARY_TEXT, + }), + sse({ + type: "response.output_item.done", + output_index: 0, + item: { + id: "msg_commentary", + type: "message", + role: "assistant", + phase: "commentary", + content: [{ type: "output_text", text: COMMENTARY_TEXT }], + }, + }), + // --- final answer item (must always be forwarded) --- + sse({ + type: "response.output_item.added", + output_index: 1, + item: { + id: "msg_final", + type: "message", + role: "assistant", + phase: "final", + content: [], + }, + }), + sse({ + type: "response.output_text.delta", + output_index: 1, + item_id: "msg_final", + content_index: 0, + delta: FINAL_TEXT, + }), + sse({ + type: "response.output_text.done", + output_index: 1, + item_id: "msg_final", + content_index: 0, + text: FINAL_TEXT, + }), + sse({ + type: "response.output_item.done", + output_index: 1, + item: { + id: "msg_final", + type: "message", + role: "assistant", + phase: "final", + content: [{ type: "output_text", text: FINAL_TEXT }], + }, + }), + sse({ + type: "response.completed", + response: { + id: "resp_6199", + output: [ + { + id: "msg_final", + type: "message", + role: "assistant", + phase: "final", + content: [{ type: "output_text", text: FINAL_TEXT }], + }, + ], + }, + }), + ]; +} + +const PASSTHROUGH_RESPONSES_OPTIONS = { + mode: "passthrough", + provider: "openai", + clientResponseFormat: "openai-responses", +}; + +test("commentary-phase output text is NOT forwarded when dropping is enabled (#6199)", async () => { + const output = await readTransformed(buildResponsesStream(), { + ...PASSTHROUGH_RESPONSES_OPTIONS, + dropResponsesCommentary: true, + }); + + assert.ok( + !output.includes(COMMENTARY_TEXT), + "commentary-phase text must be dropped from the passthrough stream" + ); + // The commentary item announcement / completion must not leak either. + assert.ok( + !output.includes("msg_commentary"), + "commentary item events must be dropped entirely" + ); + // The real answer must always be forwarded. + assert.ok(output.includes(FINAL_TEXT), "the final answer text must be forwarded"); + assert.ok(output.includes("msg_final"), "the final answer item must be forwarded"); +}); + +test("commentary passes through when dropping is disabled (gate/regression) (#6199)", async () => { + const output = await readTransformed(buildResponsesStream(), { + ...PASSTHROUGH_RESPONSES_OPTIONS, + dropResponsesCommentary: false, + }); + + assert.ok( + output.includes(COMMENTARY_TEXT), + "with the flag disabled, commentary text must pass through untouched" + ); + assert.ok(output.includes(FINAL_TEXT), "the final answer text must still be forwarded"); +});