From 38966daa33163759e6f9fbfabe949dc6fc7d1cac Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 26 Aug 2026 13:25:41 -0700 Subject: [PATCH] refactor(gateway): centralize OpenAI-compatible HTTP contracts (#130321) --- src/gateway/embeddings-http.ts | 53 +++------ src/gateway/http-common.ts | 18 +++ src/gateway/openai-http.ts | 117 ++++++-------------- src/gateway/openresponses-http.ts | 175 ++++++++++-------------------- 4 files changed, 123 insertions(+), 240 deletions(-) diff --git a/src/gateway/embeddings-http.ts b/src/gateway/embeddings-http.ts index 5cbbbebf66b6..becc8f0c88be 100644 --- a/src/gateway/embeddings-http.ts +++ b/src/gateway/embeddings-http.ts @@ -18,7 +18,13 @@ import { getMemoryEmbeddingProvider } from "../plugins/memory-embedding-provider import type { MemoryEmbeddingProvider } from "../plugins/memory-embedding-providers.js"; import type { AuthRateLimiter } from "./auth-rate-limit.js"; import type { ResolvedGatewayAuth } from "./auth.js"; -import { sendJson, sendMissingScopeForbidden, watchClientDisconnect } from "./http-common.js"; +import { + parseGatewayJsonRequest, + sendInvalidRequest, + sendJson, + sendMissingScopeForbidden, + watchClientDisconnect, +} from "./http-common.js"; import { handleGatewayPostJsonEndpoint } from "./http-endpoint-helpers.js"; import { authorizeOpenAiCompatibleHttpModelOverride, @@ -328,52 +334,30 @@ export async function handleOpenAiEmbeddingsHttpRequest( return true; } - const parsed = EmbeddingsRequestSchema.safeParse(handled.body); - if (!parsed.success) { - const issue = parsed.error.issues[0]; - sendJson(res, 400, { - error: { - message: issue ? `${issue.path.join(".")}: ${issue.message}` : "Invalid request body", - type: "invalid_request_error", - }, - }); + const payload = parseGatewayJsonRequest(res, handled.body, EmbeddingsRequestSchema); + if (!payload) { return true; } - const payload = parsed.data; const requestModel = normalizeOptionalString(payload.model) ?? ""; if (!requestModel) { - sendJson(res, 400, { - error: { message: "Missing `model`.", type: "invalid_request_error" }, - }); + sendInvalidRequest(res, "Missing `model`."); return true; } const cfg = getRuntimeConfig(); if (!isOpenClawAgentModelId(requestModel)) { - sendJson(res, 400, { - error: { - message: "Invalid `model`. Use `openclaw` or `openclaw/`.", - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, "Invalid `model`. Use `openclaw` or `openclaw/`."); return true; } const texts = resolveInputTexts(payload.input); if (!texts) { - sendJson(res, 400, { - error: { - message: "`input` must be a string or an array of strings.", - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, "`input` must be a string or an array of strings."); return true; } const inputError = validateInputTexts(texts); if (inputError) { - sendJson(res, 400, { - error: { message: inputError, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, inputError); return true; } @@ -382,9 +366,7 @@ export async function handleOpenAiEmbeddingsHttpRequest( agentId = resolveAgentIdForRequest({ req, model: requestModel }); } catch (err) { if (isAgentSelectionRequiredError(err) || isUnknownGatewayAgentError(err)) { - sendJson(res, 400, { - error: { message: err.message, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, err.message); return true; } throw err; @@ -401,12 +383,7 @@ export async function handleOpenAiEmbeddingsHttpRequest( configuredProvider, }); if ("errorMessage" in target) { - sendJson(res, 400, { - error: { - message: target.errorMessage, - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, target.errorMessage); return true; } const providerScopeKey = JSON.stringify([agentId, target.provider]); diff --git a/src/gateway/http-common.ts b/src/gateway/http-common.ts index ad04d455337a..206dddaac0bb 100644 --- a/src/gateway/http-common.ts +++ b/src/gateway/http-common.ts @@ -1,6 +1,7 @@ // Shared Gateway HTTP helpers handle small JSON/text responses, SSE headers, // body-size errors, and client disconnect aborts. import type { IncomingMessage, ServerResponse } from "node:http"; +import type { z } from "zod"; import { buildMissingScopeErrorDetails } from "../../packages/gateway-protocol/src/index.js"; import { closeRequestAfterResponse } from "../infra/http-body.js"; import { @@ -105,6 +106,23 @@ export function sendInvalidRequest(res: ServerResponse, message: string) { }); } +export function parseGatewayJsonRequest( + res: ServerResponse, + body: unknown, + schema: T, +): z.output | undefined { + const parsed = schema.safeParse(body); + if (parsed.success) { + return parsed.data; + } + const issue = parsed.error.issues[0]; + sendInvalidRequest( + res, + issue ? `${issue.path.join(".")}: ${issue.message}` : "Invalid request body", + ); + return undefined; +} + function buildMissingScopeForbiddenBody( missingScope: string | undefined, requiredScopes?: readonly string[], diff --git a/src/gateway/openai-http.ts b/src/gateway/openai-http.ts index 0e86d179d644..7e064e2d126d 100644 --- a/src/gateway/openai-http.ts +++ b/src/gateway/openai-http.ts @@ -49,6 +49,8 @@ import { import type { AuthRateLimiter } from "./auth-rate-limit.js"; import type { ResolvedGatewayAuth } from "./auth.js"; import { + parseGatewayJsonRequest, + sendInvalidRequest, sendJson, sendMissingScopeForbidden, setSseHeaders, @@ -285,12 +287,22 @@ function applyChatToolChoice(params: { tools: ClientToolDefinition[]; toolChoice type ChatCompletionStreamIdentity = { runId: string; model: string; created: number }; -function writeAssistantRoleChunk(res: ServerResponse, params: ChatCompletionStreamIdentity) { +function writeChatCompletionChunk( + res: ServerResponse, + identity: ChatCompletionStreamIdentity, + chunk: { choices: unknown[]; usage?: OpenAiChatCompletionsUsage }, +) { writeSse(res, { - id: params.runId, + id: identity.runId, object: "chat.completion.chunk", - created: params.created, - model: params.model, + created: identity.created, + model: identity.model, + ...chunk, + }); +} + +function writeAssistantRoleChunk(res: ServerResponse, params: ChatCompletionStreamIdentity) { + writeChatCompletionChunk(res, params, { choices: [{ index: 0, delta: { role: "assistant" }, finish_reason: null }], }); } @@ -299,11 +311,7 @@ function writeAssistantContentChunk( res: ServerResponse, params: ChatCompletionStreamIdentity & { content: string }, ) { - writeSse(res, { - id: params.runId, - object: "chat.completion.chunk", - created: params.created, - model: params.model, + writeChatCompletionChunk(res, params, { choices: [ { index: 0, @@ -318,11 +326,7 @@ function writeAssistantFinishChunk( res: ServerResponse, params: ChatCompletionStreamIdentity & { finishReason: "stop" | "tool_calls" }, ) { - writeSse(res, { - id: params.runId, - object: "chat.completion.chunk", - created: params.created, - model: params.model, + writeChatCompletionChunk(res, params, { choices: [ { index: 0, @@ -358,11 +362,7 @@ function writeAssistantToolCallsIncrementalChunks( }, ) { for (const [index, call] of params.toolCalls.entries()) { - writeSse(res, { - id: params.runId, - object: "chat.completion.chunk", - created: params.created, - model: params.model, + writeChatCompletionChunk(res, params, { choices: [ { index: 0, @@ -382,11 +382,7 @@ function writeAssistantToolCallsIncrementalChunks( }); for (const argsDelta of splitArgumentsForStreaming(call.arguments)) { - writeSse(res, { - id: params.runId, - object: "chat.completion.chunk", - created: params.created, - model: params.model, + writeChatCompletionChunk(res, params, { choices: [ { index: 0, @@ -412,11 +408,7 @@ function writeUsageChunk( usage: OpenAiChatCompletionsUsage; }, ) { - writeSse(res, { - id: params.runId, - object: "chat.completion.chunk", - created: params.created, - model: params.model, + writeChatCompletionChunk(res, params, { choices: [], usage: params.usage, }); @@ -892,18 +884,10 @@ export async function handleOpenAiHttpRequest( return true; } const senderIsOwner = resolveOpenAiCompatibleHttpSenderIsOwner(req, handled.requestAuth); - const parsed = OpenAiChatCompletionRequestSchema.safeParse(handled.body); - if (!parsed.success) { - const issue = parsed.error.issues[0]; - sendJson(res, 400, { - error: { - message: issue ? `${issue.path.join(".")}: ${issue.message}` : "Invalid request body", - type: "invalid_request_error", - }, - }); + const payload = parseGatewayJsonRequest(res, handled.body, OpenAiChatCompletionRequestSchema); + if (!payload) { return true; } - const payload = parsed.data; const stream = payload.stream === true; const streamIncludeUsage = stream && resolveIncludeUsageForStreaming(payload); const model = typeof payload.model === "string" ? payload.model : "openclaw"; @@ -917,9 +901,7 @@ export async function handleOpenAiHttpRequest( const legacyMaxTokens = resolveChatCompletionTokenCap(payload.max_tokens, "max_tokens"); maxTokens = maxCompletionTokens ?? legacyMaxTokens; } catch (err) { - sendJson(res, 400, { - error: { message: formatErrorMessage(err).trim(), type: "invalid_request_error" }, - }); + sendInvalidRequest(res, formatErrorMessage(err).trim()); return true; } const temperature = typeof payload.temperature === "number" ? payload.temperature : undefined; @@ -933,24 +915,14 @@ export async function handleOpenAiHttpRequest( try { responseFormat = resolveResponseFormat(payload.response_format); } catch (err) { - sendJson(res, 400, { - error: { - message: `Invalid response_format: ${formatErrorMessage(err).trim()}`, - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, `Invalid response_format: ${formatErrorMessage(err).trim()}`); return true; } let stop: string[] | undefined; try { stop = resolveStopSequences(payload.stop); } catch (err) { - sendJson(res, 400, { - error: { - message: `Invalid stop: ${formatErrorMessage(err).trim()}`, - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, `Invalid stop: ${formatErrorMessage(err).trim()}`); return true; } const samplingError = validateOpenAiSamplingParams({ @@ -961,9 +933,7 @@ export async function handleOpenAiHttpRequest( seed: payload.seed, }); if (samplingError) { - sendJson(res, 400, { - error: { message: samplingError, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, samplingError); return true; } const streamParams = @@ -1006,9 +976,7 @@ export async function handleOpenAiHttpRequest( isInvalidGatewayModelError(err) || isGatewaySessionKeyOverrideError(err) ) { - sendJson(res, 400, { - error: { message: err.message, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, err.message); return true; } throw err; @@ -1042,9 +1010,7 @@ export async function handleOpenAiHttpRequest( model, }); if (modelError) { - sendJson(res, 400, { - error: { message: modelError, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, modelError); return true; } const activeTurnContext = resolveActiveTurnContext(payload.messages); @@ -1062,12 +1028,7 @@ export async function handleOpenAiHttpRequest( toolChoicePrompt = toolChoiceResult.extraSystemPrompt; toolChoiceConstraint = toolChoiceResult.constraint; } catch (err) { - sendJson(res, 400, { - error: { - message: `Invalid tools/tool_choice: ${formatErrorMessage(err).trim()}`, - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, `Invalid tools/tool_choice: ${formatErrorMessage(err).trim()}`); return true; } let images: ImageContent[]; @@ -1075,22 +1036,12 @@ export async function handleOpenAiHttpRequest( images = await resolveImagesForRequest(activeTurnContext, limits); } catch (err) { logWarn(`openai-compat: invalid image_url content: ${String(err)}`); - sendJson(res, 400, { - error: { - message: "Invalid image_url content in `messages`.", - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, "Invalid image_url content in `messages`."); return true; } if (!prompt.message && images.length === 0) { - sendJson(res, 400, { - error: { - message: "Missing user message in `messages`.", - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, "Missing user message in `messages`."); return true; } @@ -1213,9 +1164,7 @@ export async function handleOpenAiHttpRequest( } logWarn(`openai-compat: chat completion failed: ${String(err)}`); if (isClientToolNameConflictError(err)) { - sendJson(res, 400, { - error: { message: "invalid tool configuration", type: "invalid_request_error" }, - }); + sendInvalidRequest(res, "invalid tool configuration"); return true; } const mapped = resolveOpenAiCompatError(err); diff --git a/src/gateway/openresponses-http.ts b/src/gateway/openresponses-http.ts index fad5d8de1425..fde495938022 100644 --- a/src/gateway/openresponses-http.ts +++ b/src/gateway/openresponses-http.ts @@ -47,6 +47,8 @@ import { import type { AuthRateLimiter } from "./auth-rate-limit.js"; import type { ResolvedGatewayAuth } from "./auth.js"; import { + parseGatewayJsonRequest, + sendInvalidRequest, sendJson, sendMissingScopeForbidden, setSseHeaders, @@ -474,18 +476,10 @@ export async function handleOpenResponsesHttpRequest( return true; } const senderIsOwner = resolveOpenAiCompatibleHttpSenderIsOwner(req, handled.requestAuth); - // Validate request body with Zod - const parseResult = CreateResponseBodySchema.safeParse(handled.body); - if (!parseResult.success) { - const issue = parseResult.error.issues[0]; - const message = issue ? `${issue.path.join(".")}: ${issue.message}` : "Invalid request body"; - sendJson(res, 400, { - error: { message, type: "invalid_request_error" }, - }); + const payload = parseGatewayJsonRequest(res, handled.body, CreateResponseBodySchema); + if (!payload) { return true; } - - const payload: CreateResponseBody = parseResult.data; const stream = Boolean(payload.stream); const model = payload.model; const user = payload.user; @@ -498,9 +492,7 @@ export async function handleOpenResponsesHttpRequest( isInvalidGatewayModelError(err) || isUnknownGatewayAgentError(err) ) { - sendJson(res, 400, { - error: { message: err.message, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, err.message); return true; } throw err; @@ -524,9 +516,7 @@ export async function handleOpenResponsesHttpRequest( model, }); if (modelError) { - sendJson(res, 400, { - error: { message: modelError, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, modelError); return true; } @@ -620,9 +610,7 @@ export async function handleOpenResponsesHttpRequest( } } catch (err) { logWarn(`openresponses: request parsing failed: ${String(err)}`); - sendJson(res, 400, { - error: { message: "invalid request", type: "invalid_request_error" }, - }); + sendInvalidRequest(res, "invalid request"); return true; } @@ -640,9 +628,7 @@ export async function handleOpenResponsesHttpRequest( toolChoiceConstraint = toolChoiceResult.constraint; } catch (err) { logWarn(`openresponses: tool configuration failed: ${String(err)}`); - sendJson(res, 400, { - error: { message: "invalid tool configuration", type: "invalid_request_error" }, - }); + sendInvalidRequest(res, "invalid tool configuration"); return true; } let resolved: ReturnType; @@ -662,9 +648,7 @@ export async function handleOpenResponsesHttpRequest( isInvalidGatewayModelError(err) || isGatewaySessionKeyOverrideError(err) ) { - sendJson(res, 400, { - error: { message: err.message, type: "invalid_request_error" }, - }); + sendInvalidRequest(res, err.message); return true; } throw err; @@ -708,17 +692,24 @@ export async function handleOpenResponsesHttpRequest( .join("\n\n"); if (!prompt.message) { - sendJson(res, 400, { - error: { - message: "Missing user message in `input`.", - type: "invalid_request_error", - }, - }); + sendInvalidRequest(res, "Missing user message in `input`."); return true; } const responseId = `resp_${randomUUID()}`; const responseIdentity = { id: responseId, createdAt: Math.floor(Date.now() / 1000) }; + const createFailedResponse = ( + error: { code: string; message: string }, + usage?: Usage, + ): ResponseResource => + createResponseResource({ + ...responseIdentity, + model, + status: "failed", + output: [], + error, + usage, + }); const rememberResponseSession = () => storeResponseSession(responseId, sessionKey, responseSessionScope); const outputItemId = `msg_${randomUUID()}`; @@ -776,17 +767,13 @@ export async function handleOpenResponsesHttpRequest( toolChoiceConstraint && !isToolChoiceConstraintSatisfied({ constraint: toolChoiceConstraint, pendingToolCalls }) ) { - const failed = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { + const failed = createFailedResponse( + { code: "api_error", message: resolveUnsatisfiedToolChoiceMessage(toolChoiceConstraint), }, usage, - }); + ); rememberResponseSession(); sendJson(res, 502, failed); return true; @@ -855,41 +842,25 @@ export async function handleOpenResponsesHttpRequest( } logWarn(`openresponses: non-stream response failed: ${String(err)}`); if (isClientToolNameConflictError(err)) { - const response = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { code: "invalid_request_error", message: "invalid tool configuration" }, + const response = createFailedResponse({ + code: "invalid_request_error", + message: "invalid tool configuration", }); sendJson(res, 400, response); return true; } - const response = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { code: "api_error", message: "internal error" }, - }); const mapped = resolveOpenAiCompatError(err); if (mapped) { - const mappedResponse = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { - code: mapped.error.type, - message: mapped.error.message, - }, + const mappedResponse = createFailedResponse({ + code: mapped.error.type, + message: mapped.error.message, }); rememberResponseSession(); sendJson(res, mapped.status, mappedResponse); return true; } rememberResponseSession(); - sendJson(res, 500, response); + sendJson(res, 500, createFailedResponse({ code: "api_error", message: "internal error" })); } finally { stopWatchingDisconnect(); } @@ -1035,17 +1006,13 @@ export async function handleOpenResponsesHttpRequest( } rememberResponseSession(); finalizeFailedResponse( - createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { + createFailedResponse( + { code: "server_error", message: "Assistant output cannot be represented as an append-only response stream.", }, usage, - }), + ), ); }; @@ -1219,14 +1186,10 @@ export async function handleOpenResponsesHttpRequest( terminalLifecyclePhase = "error"; rememberResponseSession(); finalizeFailedResponse( - createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { code: "api_error", message: "internal error" }, - usage: extractUsageFromResult(result), - }), + createFailedResponse( + { code: "api_error", message: "internal error" }, + extractUsageFromResult(result), + ), ); return; } @@ -1253,24 +1216,15 @@ export async function handleOpenResponsesHttpRequest( toolChoiceConstraint && !isToolChoiceConstraintSatisfied({ constraint: toolChoiceConstraint, pendingToolCalls }) ) { - const failed = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { + const failed = createFailedResponse( + { code: "api_error", message: resolveUnsatisfiedToolChoiceMessage(toolChoiceConstraint), }, - usage: finalUsage ?? createEmptyUsage(), - }); - closed = true; - stopWatchingDisconnect(); - unsubscribe(); + finalUsage ?? createEmptyUsage(), + ); rememberResponseSession(); - writeSseEvent(res, { type: "response.failed", response: failed }); - writeDone(res); - res.end(); + finalizeFailedResponse(failed); return; } @@ -1406,46 +1360,31 @@ export async function handleOpenResponsesHttpRequest( finalUsage = finalUsage ?? createEmptyUsage(); if (isClientToolNameConflictError(err)) { - const errorResponse = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { code: "invalid_request_error", message: "invalid tool configuration" }, - usage: finalUsage, - }); - - finalizeFailedResponse(errorResponse); + finalizeFailedResponse( + createFailedResponse( + { code: "invalid_request_error", message: "invalid tool configuration" }, + finalUsage, + ), + ); return; } - const errorResponse = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { code: "api_error", message: "internal error" }, - usage: finalUsage, - }); - const mapped = resolveOpenAiCompatError(err); if (mapped) { - const mappedResponse = createResponseResource({ - ...responseIdentity, - model, - status: "failed", - output: [], - error: { + const mappedResponse = createFailedResponse( + { code: mapped.error.type, message: mapped.error.message, }, - usage: finalUsage, - }); + finalUsage, + ); rememberResponseSession(); finalizeFailedResponse(mappedResponse); return; } rememberResponseSession(); - finalizeFailedResponse(errorResponse); + finalizeFailedResponse( + createFailedResponse({ code: "api_error", message: "internal error" }, finalUsage), + ); } finally { releaseAgentRootWork?.(); // Existing provider terminals must not be replaced or emitted twice.