refactor(gateway): centralize OpenAI-compatible HTTP contracts (#130321)

This commit is contained in:
Peter Steinberger
2026-08-26 13:25:41 -07:00
committed by GitHub
parent fc2724d831
commit 38966daa33
4 changed files with 123 additions and 240 deletions
+15 -38
View File
@@ -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/<agentId>`.",
type: "invalid_request_error",
},
});
sendInvalidRequest(res, "Invalid `model`. Use `openclaw` or `openclaw/<agentId>`.");
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]);
+18
View File
@@ -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<T extends z.ZodType>(
res: ServerResponse,
body: unknown,
schema: T,
): z.output<T> | 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[],
+33 -84
View File
@@ -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);
+57 -118
View File
@@ -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<typeof resolveGatewayRequestContext>;
@@ -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.