From 6e97fbf3406e703c7cf34211d0c9586751b4f3cc Mon Sep 17 00:00:00 2001 From: Jacky Lam Date: Sun, 16 Aug 2026 11:16:00 +0800 Subject: [PATCH] fix(sse): dedupe header-budget drop warns by drop-set fingerprint (#10397) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(sse): dedupe header-budget drop warns by drop-set fingerprint The 768-byte forwarded-header budget drop path emitted a full warn (with up to 20 dropped entries) on every SSE response whose headers exceeded the budget. The dropped set is usually identical across responses from the same upstream, so the repeats carried no new information — under Desktop multi-stream use this buried real errors and added event-loop serialization work. Warn once per unique drop fingerprint (sorted dropped-header names, capped at 1000 fingerprints) per process, then log at debug level. Fixes #10315 * changelog: fragment for #10397 --- .../fixes/10397-header-budget-warn-dedupe.md | 1 + open-sse/handlers/chatCore/responseHeaders.ts | 50 ++++++++- ...core-header-drop-warn-dedupe-10315.test.ts | 102 ++++++++++++++++++ 3 files changed, 150 insertions(+), 3 deletions(-) create mode 100644 changelog.d/fixes/10397-header-budget-warn-dedupe.md create mode 100644 tests/unit/chatcore-header-drop-warn-dedupe-10315.test.ts diff --git a/changelog.d/fixes/10397-header-budget-warn-dedupe.md b/changelog.d/fixes/10397-header-budget-warn-dedupe.md new file mode 100644 index 0000000000..d4d117b913 --- /dev/null +++ b/changelog.d/fixes/10397-header-budget-warn-dedupe.md @@ -0,0 +1 @@ +- **fix(sse):** the header-budget drop warning fires once per unique dropped-header set instead of on every SSE response (warn-storm fix) ([#10397](https://github.com/diegosouzapw/OmniRoute/pull/10397) — thanks @lamchun1110) diff --git a/open-sse/handlers/chatCore/responseHeaders.ts b/open-sse/handlers/chatCore/responseHeaders.ts index 43fdc5a88e..1a6602f506 100644 --- a/open-sse/handlers/chatCore/responseHeaders.ts +++ b/open-sse/handlers/chatCore/responseHeaders.ts @@ -40,7 +40,10 @@ const DEFAULT_FORWARDED_HEADER_BUDGET_BYTES = 768; * module-cache manipulation. */ export function resolveForwardedHeaderBudget(env?: string): number { - const parsed = Number.parseInt(String(env ?? process.env.OMNIROUTE_FORWARDING_HEADER_BUDGET_BYTES), 10); + const parsed = Number.parseInt( + String(env ?? process.env.OMNIROUTE_FORWARDING_HEADER_BUDGET_BYTES), + 10 + ); return Number.isFinite(parsed) && parsed > 0 ? parsed : DEFAULT_FORWARDED_HEADER_BUDGET_BYTES; } @@ -56,8 +59,31 @@ const responseHeaderEncoder = new TextEncoder(); type ResponseHeaderLogger = { warn?: (tag: string, message: string, data?: Record) => void; + debug?: (tag: string, message: string, data?: Record) => void; } | null; +/** + * #10315: the dropped-header set is usually identical across responses from the + * same upstream, so warn once per unique drop fingerprint per process, then log + * at debug level — a per-SSE-response warn storm buries real errors and adds + * event-loop serialization work. Fingerprints are dropped-header-name sets, so + * the set stays bounded by the distinct upstream header shapes in practice. + */ +const DROPPED_HEADER_WARN_FINGERPRINT_LIMIT = 1000; +const droppedHeaderWarnFingerprints = new Set(); + +export function fingerprintDroppedHeaders(dropped: Array<{ name: string; bytes: number }>): string { + return dropped + .map((header) => header.name.toLowerCase()) + .sort() + .join(","); +} + +/** Test hook: forget already-warned drop fingerprints. */ +export function resetDroppedHeaderWarnFingerprints(): void { + droppedHeaderWarnFingerprints.clear(); +} + function responseHeaderWireBytes(name: string, value: string): number { return responseHeaderEncoder.encode(`${name}: ${value}\r\n`).byteLength; } @@ -182,12 +208,30 @@ export function buildStreamingResponseHeaders( } if (droppedHeaders.length > 0) { - log?.warn?.("HTTP", "Dropped upstream response headers that exceeded forwarding budget", { + const dropPayload = { budgetBytes: MAX_FORWARDED_UPSTREAM_RESPONSE_HEADER_BYTES, forwardedBytes, droppedCount: droppedHeaders.length, droppedHeaders: droppedHeaders.slice(0, MAX_LOGGED_DROPPED_RESPONSE_HEADERS), - }); + }; + const fingerprint = fingerprintDroppedHeaders(droppedHeaders); + if (droppedHeaderWarnFingerprints.has(fingerprint)) { + log?.debug?.( + "HTTP", + "Dropped upstream response headers that exceeded forwarding budget (already warned once for this drop set)", + dropPayload + ); + } else { + if (droppedHeaderWarnFingerprints.size >= DROPPED_HEADER_WARN_FINGERPRINT_LIMIT) { + droppedHeaderWarnFingerprints.clear(); + } + droppedHeaderWarnFingerprints.add(fingerprint); + log?.warn?.( + "HTTP", + "Dropped upstream response headers that exceeded forwarding budget", + dropPayload + ); + } } const responseHeaders: Record = { diff --git a/tests/unit/chatcore-header-drop-warn-dedupe-10315.test.ts b/tests/unit/chatcore-header-drop-warn-dedupe-10315.test.ts new file mode 100644 index 0000000000..8485024b1e --- /dev/null +++ b/tests/unit/chatcore-header-drop-warn-dedupe-10315.test.ts @@ -0,0 +1,102 @@ +// #10315: the header-budget drop warn must not storm — identical dropped-header +// sets recur on every SSE response from the same upstream, so we warn once per +// unique drop fingerprint per process and fall back to debug afterwards. +import { test } from "node:test"; +import assert from "node:assert/strict"; + +const { + buildStreamingResponseHeaders, + fingerprintDroppedHeaders, + resetDroppedHeaderWarnFingerprints, +} = await import("../../open-sse/handlers/chatCore/responseHeaders.ts"); + +type DropPayload = { + budgetBytes: number; + forwardedBytes: number; + droppedCount: number; + droppedHeaders: Array<{ name: string; bytes: number }>; +}; + +function makeLogger() { + const warns: DropPayload[] = []; + const debugs: DropPayload[] = []; + return { + logger: { + warn: (_tag: string, _msg: string, data?: DropPayload) => warns.push(data as DropPayload), + debug: (_tag: string, _msg: string, data?: DropPayload) => debugs.push(data as DropPayload), + }, + warns, + debugs, + }; +} + +const meta = {} as Parameters[1]; + +// Small header + two ~600-byte headers: the first big one fits the 768-byte +// budget alongside the small one, the second is always dropped. +function oversizedProviderHeaders(): Headers { + return new Headers({ + "x-kept-small": "k".repeat(10), + "x-drop-alpha": "a".repeat(600), + "x-drop-beta": "b".repeat(600), + }); +} + +test("#10315: 100 identical oversized responses emit exactly one warn, the rest at debug", () => { + resetDroppedHeaderWarnFingerprints(); + const { logger, warns, debugs } = makeLogger(); + for (let i = 0; i < 100; i++) { + buildStreamingResponseHeaders(oversizedProviderHeaders(), meta, logger); + } + assert.equal(warns.length, 1); + assert.equal(debugs.length, 99); + assert.equal(warns[0].droppedCount, 1); +}); + +test("#10315: a different drop set warns again", () => { + resetDroppedHeaderWarnFingerprints(); + const { logger, warns } = makeLogger(); + buildStreamingResponseHeaders(oversizedProviderHeaders(), meta, logger); + assert.equal(warns.length, 1); + buildStreamingResponseHeaders( + new Headers({ + "x-kept-small": "k".repeat(10), + "x-drop-gamma": "g".repeat(600), + "x-drop-delta": "d".repeat(600), + }), + meta, + logger + ); + assert.equal(warns.length, 2); +}); + +test("#10315: fingerprint is order-insensitive to dropped header names", () => { + assert.equal( + fingerprintDroppedHeaders([ + { name: "X-Drop-Beta", bytes: 600 }, + { name: "x-drop-alpha", bytes: 600 }, + ]), + fingerprintDroppedHeaders([ + { name: "x-drop-alpha", bytes: 600 }, + { name: "X-Drop-Beta", bytes: 600 }, + ]) + ); +}); + +test("#10315: reset hook forgets fingerprints so the same drop set warns again", () => { + resetDroppedHeaderWarnFingerprints(); + const { logger, warns } = makeLogger(); + buildStreamingResponseHeaders(oversizedProviderHeaders(), meta, logger); + resetDroppedHeaderWarnFingerprints(); + buildStreamingResponseHeaders(oversizedProviderHeaders(), meta, logger); + assert.equal(warns.length, 2); +}); + +test("#10315: responses within budget never warn", () => { + resetDroppedHeaderWarnFingerprints(); + const { logger, warns, debugs } = makeLogger(); + const headers = buildStreamingResponseHeaders(new Headers({ "x-fits": "ok" }), meta, logger); + assert.equal(headers["x-fits"], "ok"); + assert.equal(warns.length, 0); + assert.equal(debugs.length, 0); +});