Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 39 additions & 8 deletions src/app/v1/_lib/proxy/stream-gate/frame-classifier.ts
Original file line number Diff line number Diff line change
Expand Up @@ -167,12 +167,6 @@ const STREAM_SIGNALS: Record<ProtocolFamily, StreamSignal> = {
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 机会。
Expand Down Expand Up @@ -356,7 +350,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;
}
Expand All @@ -365,11 +359,12 @@ function classifyFrameInner(
// 只有外层中性时才解包,供所有门控与 observer 共用同一分类结果。
const response = (parsed as Record<string, unknown>).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
Expand All @@ -385,6 +380,9 @@ function classifyParsedFrame(
for (const rule of signal.errorRules) {
if (frameRuleMatches(rule, effective, parsed)) return "error";
}
if (family === "openai-responses" && isResponsesCompactionContent(effective, parsed)) {
return "content";
}
for (const rule of signal.contentRules) {
if (frameRuleMatches(rule, effective, parsed)) return "content";
}
Expand All @@ -397,6 +395,39 @@ function classifyParsedFrame(
return "neutral";
}

/**
* Remote/server-side compaction 的 opaque state 是完整协议 payload。部分上游只在
* response.completed 中返回 output, 不会先发送 response.output_item.done。
*/
function isResponsesCompactionContent(eventType: string, parsed: object): boolean {
if (Array.isArray(parsed)) return false;

const record = parsed as Record<string, unknown>;
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<string, unknown>).output;
if (!Array.isArray(output)) return false;

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<string, unknown>;
return (
record.type === "compaction" &&
typeof record.encrypted_content === "string" &&
record.encrypted_content !== ""
);
}

/** 单条帧规则 AND 语义;空规则永不命中(防目录笔误把所有帧判成内容/错误)。 */
function frameRuleMatches(rule: FrameRule, eventType: string, parsed: unknown): boolean {
if (rule.eventTypes && rule.eventTypes.length > 0 && !rule.eventTypes.includes(eventType)) {
Expand Down
36 changes: 36 additions & 0 deletions tests/unit/proxy/stream-gate-content-gate.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down
76 changes: 76 additions & 0 deletions tests/unit/proxy/stream-gate-forwarder-integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"];
Expand Down Expand Up @@ -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" });
Expand Down Expand Up @@ -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();
Expand Down
96 changes: 96 additions & 0 deletions tests/unit/proxy/stream-gate-frame-classifier.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -300,6 +334,68 @@ 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("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(
Expand Down
Loading