diff --git a/docs-site/src/content/docs/guides/sidecars.md b/docs-site/src/content/docs/guides/sidecars.md index 006287ab9..2287b71b1 100644 --- a/docs-site/src/content/docs/guides/sidecars.md +++ b/docs-site/src/content/docs/guides/sidecars.md @@ -29,10 +29,12 @@ When Codex requests hosted `web_search` for a non-passthrough routed model, open (default 3), then removes the search tool and forces a final answer. Real client tools such as `apply_patch` or shell finalize the turn so those calls reach Codex. -Every routed-model iteration requests upstream `stream: true`, but opencodex fully buffers semantic -events internally before deciding whether to search or return the final answer. Only the first -iteration's final headers/status and 429 key rotations are acquired eagerly. Thus synthetic search -calls and preliminary output are never exposed as client-visible model output. +Every routed-model iteration follows the route's effective upstream streaming policy. Streaming +routes use the streaming parser; compatibility routes use the bounded non-streaming parser while +the client-facing response stays Responses SSE. opencodex fully buffers semantic events internally +before deciding whether to search or return the final answer. Only the first iteration's final +headers/status and 429 key rotations are acquired eagerly. Thus synthetic search calls and +preliminary output are never exposed as client-visible model output. The injected result is wrapped in an untrusted-data boundary, length-capped, and de-duplicated by source URL. In structured-output turns (`json_schema` / `json_object`) it is handed over as compact diff --git a/docs-site/src/content/docs/ja/guides/sidecars.md b/docs-site/src/content/docs/ja/guides/sidecars.md index cc64fc789..31f283720 100644 --- a/docs-site/src/content/docs/ja/guides/sidecars.md +++ b/docs-site/src/content/docs/ja/guides/sidecars.md @@ -30,9 +30,11 @@ Codex がパススルーでないルーティングモデルにホスト型 `web **反復**します。限度に達すると検索ツールを削除し最終回答を強制します。`apply_patch` や shell のような実際のクライアントツールが出たらターンを終了し該当呼び出しが Codex に渡るようにします。 -ルーティングモデルのすべての反復は上流に `stream: true` を要求しますが、opencodex は検索可否や最終 -回答を決める前に意味のある event を内部ですべてバッファリングします。最初の反復の最終 -header/status と 429 キーローテーションのみ先行取得します。したがって合成検索呼び出しと中間出力はクライアントに +ルーティングモデルの各反復は、ルートの有効な上流ストリーミングポリシーに従います。ストリーミング +ルートはストリームパーサーを使い、互換性ルートはサイズ制限付きの非ストリームパーサーを使います。 +どちらの場合もクライアント向け応答は Responses SSE のままです。opencodex は検索可否や最終回答を +決める前に意味のある event を内部ですべてバッファリングします。最初の反復の最終 header/status と +429 キーローテーションのみ先行取得します。したがって合成検索呼び出しと中間出力はクライアントに モデル出力として公開されません。 注入結果は信頼できないデータ境界で囲んで長さを制限し、ソース URL 基準で重複を除去します。構造化出力ターン(`json_schema` / `json_object`)では散文ではなく簡潔な JSON で diff --git a/docs-site/src/content/docs/ko/guides/sidecars.md b/docs-site/src/content/docs/ko/guides/sidecars.md index 1ea7596bd..00803ff74 100644 --- a/docs-site/src/content/docs/ko/guides/sidecars.md +++ b/docs-site/src/content/docs/ko/guides/sidecars.md @@ -30,10 +30,11 @@ Codex가 패스스루가 아닌 라우팅 모델에 호스팅 `web_search`를 **반복**합니다. 한도에 닿으면 검색 도구를 제거하고 최종 답변을 강제합니다. `apply_patch`나 shell 같은 실제 클라이언트 도구가 나오면 턴을 끝내 해당 호출이 Codex에 전달되게 합니다. -라우팅 모델의 모든 반복은 업스트림에 `stream: true`를 요청하지만, opencodex는 검색 여부나 최종 -답변을 결정하기 전에 의미 있는 event를 내부에서 전부 버퍼링합니다. 첫 번째 반복의 최종 -header/status와 429 key rotation만 미리 가져옵니다. 따라서 합성 검색 호출과 중간 출력은 클라이언트에 -모델 출력으로 노출되지 않습니다. +라우팅 모델의 각 반복은 해당 경로에 적용된 업스트림 스트리밍 정책을 따릅니다. 스트리밍 경로는 +스트림 parser를 사용하고, 호환성 경로는 크기 제한이 있는 비스트리밍 parser를 사용합니다. 두 경우 모두 +클라이언트 응답은 Responses SSE로 유지됩니다. opencodex는 검색 여부나 최종 답변을 결정하기 전에 +의미 있는 event를 내부에서 전부 버퍼링합니다. 첫 번째 반복의 최종 header/status와 429 key rotation만 +미리 가져옵니다. 따라서 합성 검색 호출과 중간 출력은 클라이언트에 모델 출력으로 노출되지 않습니다. 주입 결과는 신뢰할 수 없는 데이터 경계로 감싸고 길이를 제한하며, 소스 URL 기준으로 중복을 제거합니다. 구조화된 출력 턴(`json_schema` / `json_object`)에서는 산문이 아니라 간결한 JSON으로 diff --git a/docs-site/src/content/docs/ru/guides/sidecars.md b/docs-site/src/content/docs/ru/guides/sidecars.md index 861d4bb05..a6077828b 100644 --- a/docs-site/src/content/docs/ru/guides/sidecars.md +++ b/docs-site/src/content/docs/ru/guides/sidecars.md @@ -34,11 +34,13 @@ opencodex: финальному ответу. Настоящие клиентские инструменты вроде `apply_patch` или shell завершают ход, чтобы эти вызовы дошли до Codex. -Каждая итерация маршрутизируемой модели запрашивает у вышестоящего провайдера `stream: true`, но -opencodex полностью буферизует семантические события внутри, прежде чем решить, искать дальше или -вернуть финальный ответ. Заранее получаются только финальные заголовки/статус первой итерации и -ротации ключей по 429. Поэтому синтетические поисковые вызовы и промежуточный вывод никогда не -попадают к клиенту как видимый вывод модели. +Каждая итерация маршрутизируемой модели следует действующей политике потоковой передачи upstream. +Потоковые маршруты используют потоковый парсер, а маршруты совместимости — ограниченный +непотоковый парсер; ответ клиенту в обоих случаях остаётся Responses SSE. opencodex полностью +буферизует семантические события внутри, прежде чем решить, искать дальше или вернуть финальный +ответ. Заранее получаются только финальные заголовки/статус первой итерации и ротации ключей по 429. +Поэтому синтетические поисковые вызовы и промежуточный вывод никогда не попадают к клиенту как +видимый вывод модели. Внедряемый результат оборачивается в границу недоверенных данных, ограничивается по длине и дедуплицируется по URL источника. В ходах со структурированным выводом (`json_schema` / diff --git a/docs-site/src/content/docs/zh-cn/guides/sidecars.md b/docs-site/src/content/docs/zh-cn/guides/sidecars.md index 1b49584b1..8219859ef 100644 --- a/docs-site/src/content/docs/zh-cn/guides/sidecars.md +++ b/docs-site/src/content/docs/zh-cn/guides/sidecars.md @@ -27,9 +27,11 @@ Anthropic OAuth provider。Sidecar 错误会转换成长度受限的工具结果 search 工具并强制生成最终答案。如果模型调用 `apply_patch` 或 shell 等真实客户端工具,当前 turn 会结束,以便这些调用到达 Codex。 -路由模型的每次迭代都会向上游请求 `stream: true`,但 opencodex 会在决定搜索还是返回最终答案前, -在内部完整缓冲所有语义 event。只有第一次迭代的最终 header/status 和 429 key rotation 会被提前 -取得。因此,合成搜索调用和中间输出不会作为模型输出暴露给客户端。 +路由模型的每次迭代都会遵循该路由生效的上游 streaming 策略。Streaming 路由使用流式 parser; +兼容性路由使用有大小限制的非流式 parser,而面向客户端的响应在两种情况下都保持为 Responses SSE。 +opencodex 会在决定搜索还是返回最终答案前,在内部完整缓冲所有语义 event。只有第一次迭代的最终 +header/status 和 429 key rotation 会被提前取得。因此,合成搜索调用和中间输出不会作为模型输出 +暴露给客户端。 注入结果会包裹在不可信数据边界中,限制长度,并按来源 URL 去重。在结构化输出 turn (`json_schema` / `json_object`)中,结果会以紧凑 JSON 而不是普通文本传入。若路由模型是纯文本 diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 2d00ab0c4..fce81cae8 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -2464,6 +2464,7 @@ async function handleResponsesInner( const wsResponse = await runWithWebSearch({ parsed, adapter, incomingMeta: { headers: selectedForwardHeaders, abortSignal: options.abortSignal, translatorBudget }, + upstreamStreaming: parsed.stream, backend: wsPlan.backend, forwardProvider: wsPlan.forwardSidecar?.provider, anthropicSidecar: wsPlan.anthropicSidecar, diff --git a/src/web-search/loop.ts b/src/web-search/loop.ts index ce4e4eb45..0bd7d8a42 100644 --- a/src/web-search/loop.ts +++ b/src/web-search/loop.ts @@ -244,6 +244,8 @@ export interface WebSearchLoopDeps { parsed: OcxParsedRequest; adapter: ProviderAdapter; incomingMeta: IncomingMeta; + /** Effective routed-model upstream transport. Downstream web-search output remains Responses SSE. */ + upstreamStreaming?: boolean; /** Which executor runs searches. Defaults to "openai" so existing callers keep the ChatGPT path (audit F4). */ backend?: "openai" | "anthropic"; /** Required for the openai backend; unused (and typically undefined) for the anthropic backend. */ @@ -285,14 +287,16 @@ export interface WebSearchLoopDeps { } /** - * Run the main (non-OpenAI) model in a small agentic loop. Each upstream iteration is streamed and - * fully buffered internally so raw byte progress is observable without leaking a synthetic tool or - * preliminary assistant output. If the model invokes web_search, run it via the hosted sidecar, - * inject the answer as a tool_result, and loop (bounded by `maxSearches`). + * Run the main (non-OpenAI) model in a small agentic loop. Each upstream iteration follows the + * route's effective streaming policy and is fully buffered internally so raw byte progress is + * observable without leaking a synthetic tool or preliminary assistant output. If the model + * invokes web_search, run it via the hosted sidecar, inject the answer as a tool_result, and loop + * (bounded by `maxSearches`). */ export async function runWithWebSearch(deps: WebSearchLoopDeps): Promise { const translatorBudget = deps.incomingMeta.translatorBudget; const { parsed, selectedForwardHeaders, forwardProvider, hostedTool, settings, maxSearches, abortSignal, recordSidecarOutcome } = deps; + const upstreamStreaming = deps.upstreamStreaming ?? true; const backend = deps.backend ?? "openai"; const anthropicSidecar = deps.anthropicSidecar; // Mutable: 429 key-failover (deps.on429) can swap in a rebuilt adapter mid-loop. @@ -362,7 +366,7 @@ export async function runWithWebSearch(deps: WebSearchLoopDeps): Promise { const events: AdapterEvent[] = []; try { - const parse = prepared.responseAdapter.parseStream.bind(prepared.responseAdapter); + let parse: ProviderAdapter["parseStream"]; + if (upstreamStreaming) { + parse = prepared.responseAdapter.parseStream.bind(prepared.responseAdapter); + } else { + const parseResponse = prepared.responseAdapter.parseResponse?.bind(prepared.responseAdapter); + if (!parseResponse) { + throw new LoopError(502, `Provider adapter ${prepared.responseAdapter.name} does not support buffered responses`); + } + parse = async function* (response, budget) { + for (const event of await parseResponse(response, budget)) yield event; + }; + } for await (const event of parseStreamWithProgress(prepared.response, parse, { signal, inactivityTimeoutMs: routedModelStallTimeoutMs, translatorBudget, + ...(upstreamStreaming ? {} : { maxBodyBytes: TRANSLATOR_MAX_TURN_BYTES }), })) { if (event.type === "heartbeat") yield event; // Kiro's explicit-completion protocol marks ordinary assistant text as commentary while @@ -555,9 +571,10 @@ export async function runWithWebSearch(deps: WebSearchLoopDeps): Promise diff --git a/src/web-search/progress-stream.ts b/src/web-search/progress-stream.ts index f51a55efc..148d6e263 100644 --- a/src/web-search/progress-stream.ts +++ b/src/web-search/progress-stream.ts @@ -30,6 +30,8 @@ export interface ParseStreamWithProgressOptions { inactivityTimeoutMs: number; translatorBudget: TranslatorBudget; signal?: AbortSignal; + /** Optional hard cap for the raw upstream body before adapter parsing. */ + maxBodyBytes?: number; /** Kept configurable for focused tests; production callers should use the 5 second default. */ postTerminalDrainTimeoutMs?: number; } @@ -142,6 +144,9 @@ export async function* parseStreamWithProgress( options: ParseStreamWithProgressOptions, ): AsyncGenerator { const inactivityTimeoutMs = normalizeTimeout(options.inactivityTimeoutMs, "inactivityTimeoutMs"); + const maxBodyBytes = options.maxBodyBytes === undefined + ? undefined + : normalizeTimeout(options.maxBodyBytes, "maxBodyBytes"); const postTerminalDrainTimeoutMs = normalizeTimeout( options.postTerminalDrainTimeoutMs ?? DEFAULT_POST_TERMINAL_DRAIN_TIMEOUT_MS, "postTerminalDrainTimeoutMs", @@ -151,6 +156,7 @@ export async function* parseStreamWithProgress( if (!reader) throw new WebSearchStreamProtocolError("upstream response has no body"); let inactivityTimer: ReturnType | undefined; + let bodyBytes = 0; let tappedController: ReadableStreamDefaultController | undefined; let settled = false; let iterator: AsyncGenerator | undefined; @@ -224,6 +230,11 @@ export async function* parseStreamWithProgress( return; } if (result.value.byteLength === 0) continue; + bodyBytes += result.value.byteLength; + if (maxBodyBytes !== undefined && bodyBytes > maxBodyBytes) { + fail(new TranslatorBudgetExceededError("live_transient", maxBodyBytes)); + return; + } resetInactivity(); handoff.offerProgress(); controller.enqueue(result.value); diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index ed401de19..d1cd60fe4 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -228,8 +228,11 @@ closes the stream with `response.incomplete` / `upstream_stall_timeout` and canc request if no real adapter events arrive. Adapter-yielded `{ type: "heartbeat" }` events DO reset the watchdog. -The web-search loop requests `stream: true` for every routed-model iteration, but buffers the events -needed to decide whether to intercept a synthetic search call. Text explicitly phased as +The web-search loop follows the final route's effective upstream streaming policy for every +routed-model iteration. Streaming routes use `parseStream`; compatibility routes use bounded +`parseResponse` while raw-byte progress, cancellation, inactivity limits, and the 32 MiB turn cap +remain enforced. The client-facing response stays Responses SSE in both cases. The loop buffers the +events needed to decide whether to intercept a synthetic search call. Text explicitly phased as `commentary` is safe to forward live because it cannot terminate the turn; this keeps Kiro's progress visible. A Kiro stream EOF after user-facing text or reasoning gets one bounded completion retry, because neither the upstream text event nor `END_TURN` / `STOP_SEQUENCE` reliably distinguishes diff --git a/tests/web-search-progress-stream.test.ts b/tests/web-search-progress-stream.test.ts index e994459be..546b8a195 100644 --- a/tests/web-search-progress-stream.test.ts +++ b/tests/web-search-progress-stream.test.ts @@ -6,6 +6,7 @@ import { WebSearchStreamProtocolError, } from "../src/web-search/progress-stream"; import type { AdapterEvent } from "../src/types"; +import { TranslatorBudgetExceededError } from "../src/lib/translator-budget"; type ParseStream = ProviderAdapter["parseStream"]; @@ -121,6 +122,28 @@ describe("web-search streamed-body progress collector", () => { expect(events.some(event => event.type === "heartbeat")).toBe(true); }); + test("an optional raw-body cap aborts before the adapter can buffer an oversized response", async () => { + let cancelled = false; + const response = new Response(new ReadableStream({ + pull(controller) { + controller.enqueue(bytes("four")); + }, + cancel() { + cancelled = true; + }, + }, { highWaterMark: 0 })); + + const error = await collect(parseStreamWithProgress(response, drainThenDone, { + inactivityTimeoutMs: 200, + maxBodyBytes: 3, + })).then(() => undefined, reason => reason); + + expect(error).toBeInstanceOf(TranslatorBudgetExceededError); + expect(error.code).toBe("translation_buffer_limit"); + await waitFor(() => cancelled); + expect(cancelled).toBe(true); + }); + test("continuous raw-byte silence raises the exact typed inactivity error", async () => { const response = new Response(new ReadableStream({ pull() { /* never resolves */ } }, { highWaterMark: 0 })); const error = await collect(parseStreamWithProgress(response, drainThenDone, { inactivityTimeoutMs: 20 })) diff --git a/tests/web-search.test.ts b/tests/web-search.test.ts index 60ded4b15..e7e388be5 100644 --- a/tests/web-search.test.ts +++ b/tests/web-search.test.ts @@ -557,6 +557,103 @@ describe("BUG-R86 routed web-search timeout semantics", () => { expect(frames.some(frame => frame.event === "response.completed")).toBe(true); }); + test("buffered routed iterations keep upstream JSON while downstream remains SSE", async () => { + const upstreamBodies: Record[] = []; + let routedPass = 0; + const upstream = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + async fetch(request) { + upstreamBodies.push(await request.json() as Record); + routedPass += 1; + return Response.json({ + id: `chatcmpl_buffered_${routedPass}`, + object: "chat.completion", + created: 1, + model: "routed/model", + choices: [{ + index: 0, + message: routedPass === 1 + ? { + role: "assistant", + content: null, + tool_calls: [{ + id: "call_buffered_search", + type: "function", + function: { name: "web_search", arguments: JSON.stringify({ query: "current docs" }) }, + }], + } + : { role: "assistant", content: "healthy buffered answer" }, + finish_reason: routedPass === 1 ? "tool_calls" : "stop", + }], + usage: { prompt_tokens: 3, completion_tokens: 4, total_tokens: 7 }, + }); + }, + }); + let sidecarCalls = 0; + const sidecar = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + fetch() { + sidecarCalls += 1; + return new Response( + 'event: response.output_text.delta\ndata: {"type":"response.output_text.delta","delta":"docs say X"}\n\n' + + 'event: response.completed\ndata: {"type":"response.completed"}\n\n', + { headers: { "Content-Type": "text/event-stream" } }, + ); + }, + }); + const fetchStreaming: Array = []; + try { + const baseUrl = `${upstream.url.toString().replace(/\/$/, "")}/v1`; + const sidecarBaseUrl = sidecar.url.toString().replace(/\/$/, ""); + const realAdapter = createOpenAIChatAdapter({ + adapter: "openai-chat", + baseUrl, + authMode: "key", + apiKey: "local-test-key", + }); + const adapter: ProviderAdapter = { + ...realAdapter, + async fetchResponse(request, context) { + fetchStreaming.push(context?.stream); + return fetch(request.url, { + method: request.method, + headers: request.headers, + body: request.body, + signal: context?.abortSignal, + }); + }, + }; + + const response = await runWithWebSearch({ + parsed: parseRequest({ model: "routed/model", input: "hi", stream: true, tools: [{ type: "web_search" }] }), + adapter, + upstreamStreaming: false, + forwardProvider: { ...forwardProvider, baseUrl: sidecarBaseUrl }, + hostedTool: { type: "web_search" }, + selectedForwardHeaders: new Headers({ authorization: "Bearer token" }), + settings: { model: "gpt-5.4-mini", reasoning: "low", timeoutMs: 30_000 }, + maxSearches: 1, + }); + + expect(response.status).toBe(200); + expect(response.headers.get("content-type")).toContain("text/event-stream"); + const frames = await collectSse(response.body!); + expect(upstreamBodies).toHaveLength(2); + expect(upstreamBodies.every(body => body.stream === false)).toBe(true); + expect(upstreamBodies.every(body => !("stream_options" in body))).toBe(true); + expect(fetchStreaming).toEqual([false, false]); + expect(sidecarCalls).toBe(1); + expect(frames.some(frame => frame.event === "response.output_text.delta" + && frame.data.delta === "healthy buffered answer")).toBe(true); + expect(frames.some(frame => frame.event === "response.completed")).toBe(true); + } finally { + upstream.stop(true); + sidecar.stop(true); + } + }); + test("fast headers plus raw byte progress can outlive connectTimeoutMs", async () => { const delay = (ms: number) => new Promise(resolve => setTimeout(resolve, ms)); let bodyCancelled = 0;