diff --git a/src/adapters/run-turn-queue.ts b/src/adapters/run-turn-queue.ts index 23c1cd50f18..dc43e6ea622 100644 --- a/src/adapters/run-turn-queue.ts +++ b/src/adapters/run-turn-queue.ts @@ -152,6 +152,7 @@ async function* replay( export async function preflightAdapterEvents( source: AsyncIterable, + classifyFirstEvent?: (event: AdapterEvent) => Extract | undefined, ): Promise { const iterator = source[Symbol.asyncIterator](); const buffered: AdapterEvent[] = []; @@ -165,6 +166,12 @@ export async function preflightAdapterEvents( if (buffered.length > PREFLIGHT_HEARTBEAT_RETAIN_LIMIT) buffered.shift(); continue; } + const classifiedError = replayUnsafe ? undefined : classifyFirstEvent?.(next.value); + if (classifiedError) { + buffered.push(classifiedError); + await iterator.return?.(); + return { stream: replay(buffered, iterator), error: classifiedError, empty: false, replayUnsafe }; + } buffered.push(next.value); if (next.value.type === "error") { await iterator.return?.(); diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 02f7b7e567c..9cc1cd6a46a 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -18,7 +18,7 @@ import type { AdapterEventQueue } from "../../adapters/run-turn-queue"; import type { AttemptRecoveryKind } from "../../usage/log"; import { providerFetch } from "./fetch-helpers"; import { normalizeLogConversationId } from "../request-log-conversation"; -import type { AdapterEvent, OcxProviderContinuationState } from "../../types"; +import { normalizeDeclaredToolName, type AdapterEvent, type OcxProviderContinuationState } from "../../types"; import { adapterFailureFromMessage, SEND_BUDGET_EXHAUSTED_CODE } from "../../lib/errors"; import { SendBudgetExhaustedError } from "../../lib/upstream-retry"; import { @@ -38,6 +38,7 @@ import { import { rememberResponseState } from "../../responses/state"; import { trackStreamLifetime } from "../lifecycle"; import { awaitThoughtSignatureDurability } from "../../responses/thought-signature-replay"; +import { undeclaredToolCallMessage } from "../responses-undeclared-tool-guard"; /** One responsibility of the Responses request pipeline; state owners are explicit. */ export async function executeResponsesRunTurn( @@ -365,6 +366,20 @@ export async function executeResponsesRunTurn( }; const { toolNsMap, declaredToolNames, toolParameterSchemas, freeformToolNames, toolSearchToolNames } = toolBridgeMaps; + const enforceDeclaredToolNames = inboundWire !== "chat" && inboundWire !== "anthropic"; + const classifyUndeclaredFirstTool = ( + event: AdapterEvent, + ): Extract | undefined => { + if (!enforceDeclaredToolNames || event.type !== "tool_call_start") return undefined; + const effectiveName = normalizeDeclaredToolName(event.name, declaredToolNames); + if (declaredToolNames.has(effectiveName)) return undefined; + return { + type: "error", + status: 502, + errorType: "upstream_error", + message: undeclaredToolCallMessage(effectiveName), + }; + }; if (parsed.stream) { void runTurn(); let eventSource: AsyncIterable = queue.stream(); @@ -374,7 +389,7 @@ export async function executeResponsesRunTurn( eventSource = await preflightRunTurnFailover(eventSource); } if (options.comboAttempt) { - const preflight = await preflightAdapterEvents(eventSource); + const preflight = await preflightAdapterEvents(eventSource, classifyUndeclaredFirstTool); if (preflight.error || preflight.empty) { runTurnAbort.abort(); queue.close(); @@ -411,7 +426,7 @@ export async function executeResponsesRunTurn( stallTimeoutSec: config.stallTimeoutSec, hideThinkingSummary: parsed.options.hideThinkingSummary, declaredToolNames, - enforceDeclaredToolNames: inboundWire !== "chat" && inboundWire !== "anthropic", + enforceDeclaredToolNames, toolParameterSchemas, ...(options.onFirstOutput ? { onFirstOutput: options.onFirstOutput } : {}), ...(routedCompaction ? { compaction: true } : {}), @@ -468,10 +483,13 @@ export async function executeResponsesRunTurn( } if (options.comboAttempt) { const firstMeaningful = events.find(event => event.type !== "heartbeat"); - if (!firstMeaningful || firstMeaningful.type === "error") { - const message = firstMeaningful?.type === "error" + const classifiedError = firstMeaningful + ? classifyUndeclaredFirstTool(firstMeaningful) + : undefined; + if (!firstMeaningful || firstMeaningful.type === "error" || classifiedError) { + const message = classifiedError?.message ?? (firstMeaningful?.type === "error" ? firstMeaningful.message - : "Adapter ended before producing a response"; + : "Adapter ended before producing a response"); return formatErrorResponse(502, "upstream_error", redactSecretString(message)); } } @@ -482,7 +500,7 @@ export async function executeResponsesRunTurn( hideThinkingSummary: parsed.options.hideThinkingSummary, toolNsMap, declaredToolNames, - enforceDeclaredToolNames: inboundWire !== "chat" && inboundWire !== "anthropic", + enforceDeclaredToolNames, toolParameterSchemas, freeformToolNames, toolSearchToolNames, diff --git a/structure/runtime.md b/structure/runtime.md index 26f5b2a43f8..02851294dcc 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -498,7 +498,7 @@ The combo may advance to its next eligible unattempted target before output comm A definite context-window overflow is the fourth request-local verdict. A heterogeneous combo mixes windows, so "this turn does not fit THIS model" is not "this turn is impossible", and stopping at the first undersized target burned the ladder on turns a later target could hold. Evidence must come from the innermost provider message: `classifyError` remaps any occurrence of `context window`, `context length`, `maximum context` or `too many tokens` anywhere in the blob, and inheriting that looseness would let a `context_length_exceeded` token sitting in a `code` field beside `Unsupported parameter: user` authorize a replay. `src/combos/failover.ts` therefore unwraps only the exact proxy wrapper, within four envelopes and 16,384 characters, and reads the leaf message. A JSON-shaped body that does not parse fails closed, because `normalizeUpstreamErrorText` caps `classificationText` at 500 characters and a long envelope arrives here as a prefix. The verdict is admitted only for statuses that speak about the request — 400, 413, 422 and 5xx — so a 401/403 body that merely quotes context prose keeps its provider-wide cooldown instead of being rescored as request-shaped. Structured `origin_rejected`, cyber policy and the non-replayable post-send codes are all tested before it. -This is also why the classifier cannot duplicate visible output. A streaming child reaches combo classification only through `preflightComboStreamResponse`, which commits the child on any text, tool call or unknown event and synthesizes a failure envelope only for a zero-output terminal, so a turn whose text or tool call the client already saw is never reclassified as a hop. +This is also why the classifier cannot duplicate visible output. Native byte streams reach combo classification only through `preflightComboStreamResponse`, which commits the child on any text, tool call or unknown event and synthesizes a failure envelope only for a zero-output terminal. A `runTurn` adapter has the equivalent boundary in `preflightAdapterEvents`: when the first meaningful event is an undeclared tool call and no replay-unsafe heartbeat recorded a side effect, `src/server/responses/run-turn-execution.ts` checks it against the exact current request catalog and projects the existing fail-closed refusal as a pre-commit 502 so failover can continue without changing the catalog. Any earlier text, tool call, control boundary, unknown event or replay-unsafe heartbeat commits that child, so a turn whose output the client may have seen or whose side effect may have run is never replayed. Regression coverage: `tests/responses/responses-forward-prompt-envelope.test.ts`, `tests/routing/router-combo-failover-classification.test.ts`, `tests/routing/routing-policy-fallback.test.ts`, `tests/helpers/combo-context-overflow-cases.ts`, and `tests/server/server-combo-failover-e2e.test.ts`. diff --git a/tests/adapters/run-turn-queue.test.ts b/tests/adapters/run-turn-queue.test.ts index b1f821d4818..b7515820937 100644 --- a/tests/adapters/run-turn-queue.test.ts +++ b/tests/adapters/run-turn-queue.test.ts @@ -317,6 +317,37 @@ describe("run-turn adapter event preflight", () => { expect(await collect(preflight.stream)).toEqual(values); }); + test("first-event classifier replaces only the first meaningful event", async () => { + const values: AdapterEvent[] = [ + heartbeat, + { type: "tool_call_start", id: "call_stale", name: "stale_tool" }, + text("must not run"), + ]; + const classified: Extract = { + type: "error", + status: 502, + message: "undeclared tool", + }; + const preflight = await preflightAdapterEvents(events(values), event => + event.type === "tool_call_start" ? classified : undefined); + expect(preflight.error).toEqual(classified); + expect(preflight.empty).toBe(false); + expect(await collect(preflight.stream)).toEqual([heartbeat, classified]); + }); + + test("first-event classifier cannot replace after a replay-unsafe heartbeat", async () => { + const tool: AdapterEvent = { type: "tool_call_start", id: "call_stale", name: "stale_tool" }; + const values: AdapterEvent[] = [{ type: "heartbeat", replayUnsafe: true }, tool]; + const preflight = await preflightAdapterEvents(events(values), () => ({ + type: "error", + status: 502, + message: "must not replace", + })); + expect(preflight.error).toBeUndefined(); + expect(preflight.replayUnsafe).toBe(true); + expect(await collect(preflight.stream)).toEqual(values); + }); + test("immediate done is a commit", async () => { const preflight = await preflightAdapterEvents(events([done])); expect(preflight.error).toBeUndefined(); diff --git a/tests/server/server-combo-zero-output-failover.test.ts b/tests/server/server-combo-zero-output-failover.test.ts index a5a7a9fa292..306ee73a9f1 100644 --- a/tests/server/server-combo-zero-output-failover.test.ts +++ b/tests/server/server-combo-zero-output-failover.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, setDefaultTimeout, test } from "bun:test"; +import { afterEach, beforeEach, describe, expect, mock, setDefaultTimeout, test } from "bun:test"; import { mkdtempSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -17,8 +17,34 @@ import { flushResponseState, responseStatePersistPendingForTests, } from "../../src/responses/state"; -import { handleResponses } from "../../src/server/responses"; -import type { OcxConfig } from "../../src/types"; +import type { ProviderAdapter } from "../../src/adapters/base"; +import type { AdapterEvent, OcxConfig, OcxProviderConfig } from "../../src/types"; + +const actualResolver = await import("../../src/server/adapter-resolve"); +const actualResolveAdapter = actualResolver.resolveAdapter; +let customRunTurn: NonNullable | undefined; + +mock.module("../../src/server/adapter-resolve", () => ({ + ...actualResolver, + resolveAdapter(provider: OcxProviderConfig, cacheRetention?: "none" | "short" | "long") { + if (provider.adapter !== "test-run-turn") { + return actualResolveAdapter(provider, cacheRetention); + } + return { + name: "test-run-turn", + buildRequest: () => ({ url: provider.baseUrl, method: "POST", headers: {}, body: "" }), + async *parseStream(): AsyncGenerator { + yield { type: "error", message: "test runTurn adapter does not use parseStream" }; + }, + async runTurn(parsed, incoming, emit) { + if (!customRunTurn) throw new Error("custom runTurn not installed"); + await customRunTurn(parsed, incoming, emit); + }, + } satisfies ProviderAdapter; + }, +})); + +const { handleResponses } = await import("../../src/server/responses"); /** * Zero-output combo failover driven by a bare Responses SSE `error` event. @@ -29,10 +55,10 @@ import type { OcxConfig } from "../../src/types"; * lowers a cap, so the way back under it is to hold new cases in a sibling file rather * than to raise the number. * - * The harness below is the subset of that file's fixture this case actually uses: real + * The harness below is the subset of that file's fixture these cases actually use: real * loopback upstreams, an isolated home, and the combo/request-log state that leaks - * between tests. No module is mocked here, because this case drives the real - * `openai-responses` adapter. + * between tests. The loopback cases drive real adapters; the runTurn case uses the same + * narrow resolver seam as the parent file to emit deterministic adapter events. */ // The parent file raises this for the same reason: a real loopback server plus combo @@ -63,6 +89,7 @@ beforeEach(() => { }); afterEach(async () => { + customRunTurn = undefined; // Release before home teardown to prevent Windows removal failures and a live unlinked database. releaseSpendHome?.(); releaseSpendHome = undefined; @@ -186,4 +213,114 @@ describe("combo zero-output bare Responses error failover", () => { ], }); }); + + test("undeclared first adapter tool call hops without changing the request catalog", async () => { + const requests: Record[] = []; + const upstream = (content: string) => serve(async request => { + requests.push(await request.json() as Record); + return new Response([ + `data: ${JSON.stringify({ choices: [{ delta: JSON.parse(content) }] })}`, + "data: [DONE]", + "", + ].join("\n"), { headers: { "content-type": "text/event-stream" } }); + }); + const a = upstream(JSON.stringify({ + tool_calls: [{ index: 0, id: "call_stale", function: { name: "stale_tool", arguments: "{}" } }], + })); + const b = upstream(JSON.stringify({ content: "tool-safe backup" })); + const config = comboConfig({ + a: provider("openai-chat", baseUrl(a), "key-a"), + b: provider("openai-chat", baseUrl(b), "key-b"), + }); + const tools = [{ type: "function", name: "current_tool", parameters: { type: "object" } }]; + + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "combo/free", input: "hello", stream: true, tools }), + }), config, { model: "", provider: "" }); + const body = await response.text(); + expect(body).toContain("tool-safe backup"); + expect(body).not.toContain("stale_tool"); + expect(requests).toHaveLength(2); + expect(requests[0]?.tools).toEqual(requests[1]?.tools); + expect(JSON.stringify(requests[0]?.tools)).toContain("current_tool"); + }); + + test("non-streaming undeclared adapter tool call hops without changing the request catalog", async () => { + const catalogs: unknown[] = []; + const hits: string[] = []; + customRunTurn = async (parsed, _incoming, emit) => { + hits.push(parsed.modelId); + catalogs.push(structuredClone((parsed._rawBody as { tools?: unknown }).tools)); + if (parsed.modelId === "m1") { + emit({ type: "tool_call_start", id: "call_stale", name: "stale_tool" }); + emit({ type: "tool_call_delta", arguments: "{}" }); + emit({ type: "tool_call_end" }); + emit({ type: "done" }); + return; + } + emit({ type: "text_delta", text: "non-streaming tool-safe backup" }); + emit({ type: "done" }); + }; + const config = comboConfig({ + a: provider("test-run-turn", "https://a.test/v1", "key-a"), + b: provider("test-run-turn", "https://b.test/v1", "key-b"), + }); + const tools = [{ type: "function", name: "current_tool", parameters: { type: "object" } }]; + + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "combo/free", input: "hello", stream: false, tools }), + }), config, { model: "", provider: "" }); + const body = await response.text(); + expect(response.status).toBe(200); + expect(body).toContain("non-streaming tool-safe backup"); + expect(body).not.toContain("stale_tool"); + expect(hits).toEqual(["m1", "m2"]); + expect(catalogs).toHaveLength(2); + expect(catalogs[0]).toEqual(catalogs[1]); + expect(JSON.stringify(catalogs[0])).toContain("current_tool"); + }); + + test("undeclared adapter tool call after text never replays on backup", async () => { + const hits: string[] = []; + const a = serve(() => { + hits.push("a"); + return new Response([ + `data: ${JSON.stringify({ choices: [{ delta: { content: "already visible" } }] })}`, + `data: ${JSON.stringify({ choices: [{ delta: { + tool_calls: [{ index: 0, id: "call_stale", function: { name: "stale_tool", arguments: "{}" } }], + } }] })}`, + "data: [DONE]", + "", + ].join("\n"), { headers: { "content-type": "text/event-stream" } }); + }); + const b = serve(() => { + hits.push("b"); + return new Response("data: [DONE]\n\n", { + headers: { "content-type": "text/event-stream" }, + }); + }); + const config = comboConfig({ + a: provider("openai-chat", baseUrl(a), "key-a"), + b: provider("openai-chat", baseUrl(b), "key-b"), + }); + + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + model: "combo/free", + input: "hello", + stream: true, + tools: [{ type: "function", name: "current_tool", parameters: { type: "object" } }], + }), + }), config, { model: "", provider: "" }); + const body = await response.text(); + expect(body).toContain("already visible"); + expect(body).toContain("response.failed"); + expect(hits).toEqual(["a"]); + }); });