From 40dbb569398fcc1acc66487eccac156a0d6b2210 Mon Sep 17 00:00:00 2001 From: ding113 Date: Mon, 10 Aug 2026 17:01:03 +0800 Subject: [PATCH 1/2] fix(stream-gate): recognize Responses compaction in terminal frames Some Responses upstreams return compaction output only in response.completed without first emitting response.output_item.done, so the content gate treated the stream as empty and triggered false failover with circuit-breaking. classifyParsedFrame now accepts the protocol family and, for openai-responses, detects a compaction output item with non-empty encrypted_content inside response.completed, classifying the frame as content so the gate commits the stream. Regression tests cover compaction-only terminal frames and custom tool-call input deltas across the classifier, content gate, and forwarder integration paths. Fixes #1410 --- .../proxy/stream-gate/frame-classifier.ts | 29 ++++++- .../proxy/stream-gate-content-gate.test.ts | 36 +++++++++ .../stream-gate-forwarder-integration.test.ts | 76 +++++++++++++++++++ .../stream-gate-frame-classifier.test.ts | 68 +++++++++++++++++ 4 files changed, 207 insertions(+), 2 deletions(-) diff --git a/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts b/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts index b6dac5c55..bb9fe5869 100644 --- a/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts +++ b/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts @@ -356,7 +356,7 @@ function classifyFrameInner( return "malformed"; } - const outerVerdict = classifyParsedFrame(signal, eventName, parsed); + const outerVerdict = classifyParsedFrame(family, signal, eventName, parsed); if (outerVerdict !== "neutral" || family !== "gemini" || Array.isArray(parsed)) { return outerVerdict; } @@ -365,11 +365,12 @@ function classifyFrameInner( // 只有外层中性时才解包,供所有门控与 observer 共用同一分类结果。 const response = (parsed as Record).response; return response && typeof response === "object" && !Array.isArray(response) - ? classifyParsedFrame(signal, eventName, response) + ? classifyParsedFrame(family, signal, eventName, response) : outerVerdict; } function classifyParsedFrame( + family: ProtocolFamily, signal: StreamSignal, eventName: string | null, parsed: object @@ -385,6 +386,9 @@ function classifyParsedFrame( for (const rule of signal.errorRules) { if (frameRuleMatches(rule, effective, parsed)) return "error"; } + if (family === "openai-responses" && isCompletedResponsesCompaction(effective, parsed)) { + return "content"; + } for (const rule of signal.contentRules) { if (frameRuleMatches(rule, effective, parsed)) return "content"; } @@ -397,6 +401,27 @@ function classifyParsedFrame( return "neutral"; } +/** + * 部分 Responses 上游只在 response.completed 中返回 compaction output, + * 不会先发送 response.output_item.done。逐项关联 type 与 opaque state, + * 避免不同 output item 的字段组合造成误提交。 + */ +function isCompletedResponsesCompaction(eventType: string, parsed: object): boolean { + if (eventType !== "response.completed" || Array.isArray(parsed)) return false; + + const response = (parsed as Record).response; + if (response === null || typeof response !== "object" || Array.isArray(response)) return false; + + const output = (response as Record).output; + if (!Array.isArray(output)) return false; + + return output.some((item) => { + if (item === null || typeof item !== "object" || Array.isArray(item)) return false; + const record = item as Record; + return record.type === "compaction" && isNonEmptyValue(record.encrypted_content); + }); +} + /** 单条帧规则 AND 语义;空规则永不命中(防目录笔误把所有帧判成内容/错误)。 */ function frameRuleMatches(rule: FrameRule, eventType: string, parsed: unknown): boolean { if (rule.eventTypes && rule.eventTypes.length > 0 && !rule.eventTypes.includes(eventType)) { diff --git a/tests/unit/proxy/stream-gate-content-gate.test.ts b/tests/unit/proxy/stream-gate-content-gate.test.ts index bd97a540a..56140dd01 100644 --- a/tests/unit/proxy/stream-gate-content-gate.test.ts +++ b/tests/unit/proxy/stream-gate-content-gate.test.ts @@ -254,6 +254,42 @@ describe("runStreamContentGate", () => { expect(new TextDecoder().decode(rest.value)).toBe(completed); }); + it("openai-responses: commits compaction carried only by response.completed", async () => { + const completed = + 'event: response.completed\ndata: {"type":"response.completed","response":{"status":"completed","output":[{"type":"compaction","encrypted_content":"opaque-state"}]}}\n\n'; + const reader = readerFromChunks([completed]); + + const result = await runStreamContentGate(reader, { + ...GATE_OPTIONS, + family: "openai-responses", + }); + + expect(result.committed).toBe(true); + if (!result.committed) return; + expect(await drainPrefix(result.prefixChunks)).toBe(completed); + expect(result.readerDone).toBe(false); + }); + + it("openai-responses: commits custom tool-call input before response.completed", async () => { + const toolInput = + 'event: response.custom_tool_call_input.delta\ndata: {"type":"response.custom_tool_call_input.delta","delta":"{\\"path\\":\\"README.md\\"}"}\n\n'; + const completed = + 'event: response.completed\ndata: {"type":"response.completed","response":{"status":"completed"}}\n\n'; + const reader = readerFromChunks([toolInput, completed]); + + const result = await runStreamContentGate(reader, { + ...GATE_OPTIONS, + family: "openai-responses", + }); + + expect(result.committed).toBe(true); + if (!result.committed) return; + expect(await drainPrefix(result.prefixChunks)).toBe(toolInput); + expect(result.readerDone).toBe(false); + const rest = await reader.read(); + expect(new TextDecoder().decode(rest.value)).toBe(completed); + }); + it("gemini: usage-only chunks buffer until content commits", async () => { const reader = readerFromChunks([ 'data: {"usageMetadata":{"totalTokenCount":1}}\n\n', diff --git a/tests/unit/proxy/stream-gate-forwarder-integration.test.ts b/tests/unit/proxy/stream-gate-forwarder-integration.test.ts index 8ce3c6296..23d4beb42 100644 --- a/tests/unit/proxy/stream-gate-forwarder-integration.test.ts +++ b/tests/unit/proxy/stream-gate-forwarder-integration.test.ts @@ -224,6 +224,35 @@ const OPENAI_RESPONSES_WINNER_FRAMES = [ }), ]; +const VALID_OPENAI_RESPONSES_STREAMS = [ + { + name: "terminal compaction output", + frames: [ + sseFrame("response.completed", { + type: "response.completed", + response: { + id: "resp_compaction", + status: "completed", + output: [{ id: "cmp_1", type: "compaction", encrypted_content: "opaque-state" }], + }, + }), + ], + }, + { + name: "custom tool-call input deltas", + frames: [ + sseFrame("response.custom_tool_call_input.delta", { + type: "response.custom_tool_call_input.delta", + delta: '{"path":"README.md"}', + }), + sseFrame("response.completed", { + type: "response.completed", + response: { id: "resp_tool", status: "completed" }, + }), + ], + }, +] as const; + type ReplayGateCase = { name: string; providerType: Provider["providerType"]; @@ -516,6 +545,30 @@ describe("F1 stream content gate x ProxyForwarder sequential path", () => { expect(mocks.recordFailure).not.toHaveBeenCalled(); }); + test("Responses terminal compaction 在 enforce 模式下直接提交且不计入熔断", async () => { + const provider = createProvider({ id: 1, name: "compaction", providerType: "codex" }); + const session = createSession(); + session.setProvider(provider); + Object.assign(session, { + requestUrl: new URL("https://example.com/v1/responses"), + originalFormat: "response", + endpointPolicy: resolveEndpointPolicy("/v1/responses"), + }); + + const frames = VALID_OPENAI_RESPONSES_STREAMS[0].frames.slice(); + const doForward = spyOnDoForward(); + doForward.mockImplementationOnce(async () => createSseResponse(frames)); + + const response = await ProxyForwarder.send(session); + const text = await response.text(); + + expect(response.status).toBe(200); + expect(text).toBe(frames.join("")); + expect(doForward).toHaveBeenCalledTimes(1); + expect(mocks.pickRandomProviderWithExclusion).not.toHaveBeenCalled(); + expect(mocks.recordFailure).not.toHaveBeenCalled(); + }); + test("terminal-only 流(message_stop 即终止)按 empty_stream 失败并切换供应商", async () => { const provider1 = createProvider({ id: 1, name: "gate-p1" }); const provider2 = createProvider({ id: 2, name: "gate-p2" }); @@ -631,6 +684,29 @@ describe("F1 stream content gate x ProxyForwarder sequential path", () => { expect(mocks.recordFailure).not.toHaveBeenCalled(); }); + test.each(VALID_OPENAI_RESPONSES_STREAMS)( + "Replay owner 将 $name 视为有效内容,不触发 502/failover/熔断", + async ({ frames }) => { + const provider = createProvider({ id: 1, name: "responses-valid", providerType: "codex" }); + const session = createSession(); + session.setProvider(provider); + attachReplayOwner(session, REPLAY_GATE_CASES[0]); + + const streamFrames = frames.slice(); + const doForward = spyOnDoForward(); + doForward.mockImplementationOnce(async () => createSseResponse(streamFrames)); + + const response = await ProxyForwarder.send(session); + const text = await response.text(); + + expect(response.status).toBe(200); + expect(text).toBe(streamFrames.join("")); + expect(doForward).toHaveBeenCalledTimes(1); + expect(mocks.pickRandomProviderWithExclusion).not.toHaveBeenCalled(); + expect(mocks.recordFailure).not.toHaveBeenCalled(); + } + ); + test("Replay owner 在所有 precommit attempt 失败后立即释放所有权", async () => { const provider = createProvider({ id: 1, name: "replay-only", providerType: "codex" }); const session = createSession(); diff --git a/tests/unit/proxy/stream-gate-frame-classifier.test.ts b/tests/unit/proxy/stream-gate-frame-classifier.test.ts index 3e8fb3053..8866dca3b 100644 --- a/tests/unit/proxy/stream-gate-frame-classifier.test.ts +++ b/tests/unit/proxy/stream-gate-frame-classifier.test.ts @@ -270,6 +270,40 @@ describe("classifyFrame: openai-responses", () => { ).toBe("content"); }); + it("content: custom tool-call input delta and done payloads", () => { + expect( + classifyFrame( + "openai-responses", + "response.custom_tool_call_input.delta", + '{"type":"response.custom_tool_call_input.delta","delta":"{\\"path\\":\\"README.md\\"}"}' + ) + ).toBe("content"); + expect( + classifyFrame( + "openai-responses", + "response.custom_tool_call_input.done", + '{"type":"response.custom_tool_call_input.done","input":"{\\"path\\":\\"README.md\\"}"}' + ) + ).toBe("content"); + }); + + it("neutral: empty custom tool-call input delta and done payloads", () => { + expect( + classifyFrame( + "openai-responses", + "response.custom_tool_call_input.delta", + '{"type":"response.custom_tool_call_input.delta","delta":""}' + ) + ).toBe("neutral"); + expect( + classifyFrame( + "openai-responses", + "response.custom_tool_call_input.done", + '{"type":"response.custom_tool_call_input.done","input":""}' + ) + ).toBe("neutral"); + }); + it("neutral: output_item.added carrying only tool metadata", () => { expect( classifyFrame( @@ -300,6 +334,40 @@ describe("classifyFrame: openai-responses", () => { ).toBe("content"); }); + it("content: response.completed carrying compaction output with opaque state", () => { + expect( + classifyFrame( + "openai-responses", + "response.completed", + '{"type":"response.completed","response":{"status":"completed","output":[{"type":"compaction","encrypted_content":"opaque-state"}]}}' + ) + ).toBe("content"); + }); + + it("terminal: response.completed without a non-empty compaction output", () => { + expect( + classifyFrame( + "openai-responses", + "response.completed", + '{"type":"response.completed","response":{"status":"completed","output":[{"type":"compaction","encrypted_content":""}]}}' + ) + ).toBe("terminal"); + expect( + classifyFrame( + "openai-responses", + "response.completed", + '{"type":"response.completed","response":{"status":"completed","output":[{"type":"reasoning","encrypted_content":"opaque-state"}]}}' + ) + ).toBe("terminal"); + expect( + classifyFrame( + "openai-responses", + "response.completed", + '{"type":"response.completed","response":{"status":"completed","output":[{"type":"compaction","encrypted_content":""},{"type":"reasoning","encrypted_content":"opaque-state"}]}}' + ) + ).toBe("terminal"); + }); + it("neutral: empty or non-compaction encrypted output item", () => { expect( classifyFrame( From 2856cd74cf624c8cae03a32f8964878ac158e322 Mon Sep 17 00:00:00 2001 From: ding113 Date: Mon, 10 Aug 2026 17:27:20 +0800 Subject: [PATCH 2/2] fix(stream-gate): require non-empty string for compaction encrypted_content The compaction signal rule on response.output_item.done accepted any truthy encrypted_content value, allowing non-string types (booleans, numbers, objects) to be misclassified as content. Consolidate the per-item type guard into isNonEmptyCompactionItem and apply it to both response.output_item.done and response.completed paths so the opaque state must be a non-empty string before a frame is committed as content. Add regression tests covering malformed encrypted_content types across event variants. --- .../proxy/stream-gate/frame-classifier.ts | 42 +++++++++++-------- .../stream-gate-frame-classifier.test.ts | 28 +++++++++++++ 2 files changed, 52 insertions(+), 18 deletions(-) diff --git a/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts b/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts index bb9fe5869..609616b6a 100644 --- a/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts +++ b/src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts @@ -167,12 +167,6 @@ const STREAM_SIGNALS: Record = { eventTypes: ["response.code_interpreter_call_code.done"], anyPaths: ["code"], }, - { - // Remote/server-side compaction 的 opaque state 是完整协议 payload,只有非空且类型精确匹配才提交。 - eventTypes: ["response.output_item.done"], - anyPaths: ["item.encrypted_content"], - valueMatches: [{ path: "item.type", values: ["compaction"] }], - }, { // output_item.added 的 name/id/status 只是结构元数据;真实 payload 到达前不能提交, // 否则紧随其后的 response.failed / 断流将失去透明 fallback 机会。 @@ -386,7 +380,7 @@ function classifyParsedFrame( for (const rule of signal.errorRules) { if (frameRuleMatches(rule, effective, parsed)) return "error"; } - if (family === "openai-responses" && isCompletedResponsesCompaction(effective, parsed)) { + if (family === "openai-responses" && isResponsesCompactionContent(effective, parsed)) { return "content"; } for (const rule of signal.contentRules) { @@ -402,24 +396,36 @@ function classifyParsedFrame( } /** - * 部分 Responses 上游只在 response.completed 中返回 compaction output, - * 不会先发送 response.output_item.done。逐项关联 type 与 opaque state, - * 避免不同 output item 的字段组合造成误提交。 + * Remote/server-side compaction 的 opaque state 是完整协议 payload。部分上游只在 + * response.completed 中返回 output, 不会先发送 response.output_item.done。 */ -function isCompletedResponsesCompaction(eventType: string, parsed: object): boolean { - if (eventType !== "response.completed" || Array.isArray(parsed)) return false; +function isResponsesCompactionContent(eventType: string, parsed: object): boolean { + if (Array.isArray(parsed)) return false; - const response = (parsed as Record).response; + const record = parsed as Record; + if (eventType === "response.output_item.done") { + return isNonEmptyCompactionItem(record.item); + } + if (eventType !== "response.completed") return false; + + const response = record.response; if (response === null || typeof response !== "object" || Array.isArray(response)) return false; const output = (response as Record).output; if (!Array.isArray(output)) return false; - return output.some((item) => { - if (item === null || typeof item !== "object" || Array.isArray(item)) return false; - const record = item as Record; - return record.type === "compaction" && isNonEmptyValue(record.encrypted_content); - }); + return output.some(isNonEmptyCompactionItem); +} + +/** 同一 output item 内的 type 与 opaque state 必须同时满足协议类型约束。 */ +function isNonEmptyCompactionItem(item: unknown): boolean { + if (item === null || typeof item !== "object" || Array.isArray(item)) return false; + const record = item as Record; + return ( + record.type === "compaction" && + typeof record.encrypted_content === "string" && + record.encrypted_content !== "" + ); } /** 单条帧规则 AND 语义;空规则永不命中(防目录笔误把所有帧判成内容/错误)。 */ diff --git a/tests/unit/proxy/stream-gate-frame-classifier.test.ts b/tests/unit/proxy/stream-gate-frame-classifier.test.ts index 8866dca3b..50b8ecde8 100644 --- a/tests/unit/proxy/stream-gate-frame-classifier.test.ts +++ b/tests/unit/proxy/stream-gate-frame-classifier.test.ts @@ -368,6 +368,34 @@ describe("classifyFrame: openai-responses", () => { ).toBe("terminal"); }); + it("rejects non-string compaction encrypted content", () => { + for (const encryptedContent of [true, 42, { opaque: "state" }]) { + expect( + classifyFrame( + "openai-responses", + "response.output_item.done", + JSON.stringify({ + type: "response.output_item.done", + item: { type: "compaction", encrypted_content: encryptedContent }, + }) + ) + ).toBe("neutral"); + expect( + classifyFrame( + "openai-responses", + "response.completed", + JSON.stringify({ + type: "response.completed", + response: { + status: "completed", + output: [{ type: "compaction", encrypted_content: encryptedContent }], + }, + }) + ) + ).toBe("terminal"); + } + }); + it("neutral: empty or non-compaction encrypted output item", () => { expect( classifyFrame(