diff --git a/src/server/index.ts b/src/server/index.ts index 7d08fc422..52466cbd3 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -872,7 +872,6 @@ export function startServer(port?: number, deps: StartServerDeps = {}) { } if (url.pathname === "/v1/responses" && req.method === "POST") { - disableResponsesRequestTimeout(req, requestServer); if (isDraining()) { return drainingResponse(req); } @@ -901,6 +900,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}) { return runAdmittedHttpTurn(req, async turnAdmissionLease => { const response = await handleResponses(req, config, logCtx, { turnAdmissionLease, + onRequestBodyRead: () => disableResponsesRequestTimeout(req, requestServer), abortSignal: req.signal, onFirstOutput: () => recordFirstOutput(logCtx, start), onNativePassthroughTerminal: status => { diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 2d00ab0c4..ac318f29b 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -567,6 +567,8 @@ export interface ConsumedComboFailure { export interface HandleResponsesOptions { turnAdmissionLease?: AdmissionLease; + /** Called only after the complete inbound body has been read and parsed successfully. */ + onRequestBodyRead?: () => void; forceEmptyResponseId?: boolean; abortSignal?: AbortSignal; /** One-shot TTFT callback: first non-empty model output observed (WP4). */ @@ -1366,6 +1368,7 @@ async function handleResponsesInner( } return formatErrorResponse(400, "invalid_request_error", err instanceof Error ? err.message : String(err)); } + options.onRequestBodyRead?.(); // Prefer a pre-populated id (routed Claude) over Responses headers that may be // absent or synthetically injected (session_id from prompt_cache_key). if (!logCtx.conversationId) { diff --git a/tests/server-auth.test.ts b/tests/server-auth.test.ts index 849a5a204..9d1496e59 100644 --- a/tests/server-auth.test.ts +++ b/tests/server-auth.test.ts @@ -35,6 +35,7 @@ import { } from "../src/server"; import { clearRequestLogsForTests, getRequestLogEntries } from "../src/server/request-log"; import { handleManagementAPI } from "../src/server/management-api"; +import { handleResponses } from "../src/server/responses"; import type { OcxConfig } from "../src/types"; import { fakeChatGptJwt } from "./helpers/fake-chatgpt-jwt"; import { installIsolatedCodexHome, type IsolatedCodexHome } from "./helpers/isolated-codex-home"; @@ -345,6 +346,40 @@ describe("server local API auth", () => { })).toBe(false); }); + test("responses handler keeps the request timeout until the body is fully parsed", async () => { + let controller!: ReadableStreamDefaultController; + const body = new ReadableStream({ + start(value) { + controller = value; + }, + }); + const req = new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body, + }); + let bodyRead = false; + const responsePromise = handleResponses(req, config(), { + model: "unknown", + provider: "unknown", + }, { + onRequestBodyRead: () => { + bodyRead = true; + }, + }); + + controller.enqueue(new TextEncoder().encode('{"model":"missing/provider","input":"hello"')); + await Bun.sleep(10); + expect(bodyRead).toBe(false); + + controller.enqueue(new TextEncoder().encode("}")); + controller.close(); + const response = await responsePromise; + expect(bodyRead).toBe(true); + expect(response.status).toBeGreaterThanOrEqual(400); + expect(response.status).toBeLessThan(500); + }); + test("loopback hostnames do not require opencodex API auth", () => { expect(isLoopbackHostname(undefined)).toBe(true); expect(isLoopbackHostname("")).toBe(true);