diff --git a/extensions/a2a/src/http.test.ts b/extensions/a2a/src/http.test.ts index 7e796fb8d4e5..f75020c4b310 100644 --- a/extensions/a2a/src/http.test.ts +++ b/extensions/a2a/src/http.test.ts @@ -1,13 +1,28 @@ import { EventEmitter } from "node:events"; -import type { ServerResponse } from "node:http"; +import { createServer, type ServerResponse } from "node:http"; import { VERSION } from "openclaw/plugin-sdk/cli-runtime"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; -import { createMockIncomingRequest, createMockServerResponse } from "openclaw/plugin-sdk/test-env"; +import { + createMockIncomingRequest, + createMockServerResponse, + postRawWebhook, +} from "openclaw/plugin-sdk/test-env"; import { afterEach, describe, expect, it, vi } from "vitest"; import { createA2aHttpHandler } from "./http.js"; import { A2aTaskStore } from "./task-store.js"; import type { A2aChannelConfig } from "./types.js"; +vi.mock("node:timers", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + setTimeout: ((callback: (...args: unknown[]) => void, delay?: number, ...args: unknown[]) => + globalThis.setTimeout(callback, delay, ...args)) as typeof actual.setTimeout, + clearTimeout: ((timer: ReturnType | undefined) => + globalThis.clearTimeout(timer)) as typeof actual.clearTimeout, + }; +}); + const activeStores = new Set(); afterEach(() => { @@ -101,6 +116,7 @@ async function startHttpHarness(options?: { return { baseUrl, taskStore, + handler, async get(endpoint: string) { return await dispatchRequest({ method: "GET", endpoint }); }, @@ -291,14 +307,99 @@ describe("A2A HTTP authentication and request limits", () => { } }); - it("returns a readable HTTP 413 response for request bodies above 1 MiB", async () => { + it("delivers HTTP 413 over the wire and closes for request bodies above 1 MiB", async () => { const harness = await startHttpHarness(); - const response = await harness.post("x".repeat(1024 * 1024 + 1)); - - expect(response.status).toBe(413); - await expect(response.json()).resolves.toMatchObject({ - error: expect.stringContaining("1 MiB"), + const server = createServer((req, res) => { + void harness.handler(req, res); }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected the A2A test server to have a TCP address"); + } + + // Declared and sent in one write: the shape whose rejection used to race the flush. + const result = await postRawWebhook({ + url: `http://127.0.0.1:${address.port}/a2a/v1`, + body: "x".repeat(1024 * 1024 + 1), + headers: { + "content-type": "application/json", + authorization: "Bearer alpha-secret", + }, + }); + + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(result.headers.connection).toBe("close"); + expect(JSON.parse(result.body)).toEqual({ + error: "Request body exceeds the 1 MiB limit", + }); + expect(result.closedByServer).toBe(true); + } finally { + server.closeAllConnections(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } + }); + + it("delivers the JSON-RPC timeout response before closing a partial upload", async () => { + vi.useFakeTimers(); + const harness = await startHttpHarness(); + const server = createServer((req, res) => { + void harness.handler(req, res); + }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected the A2A test server to have a TCP address"); + } + const requestReceived = new Promise((resolve) => { + server.once("connection", (socket) => socket.once("data", () => resolve())); + }); + const resultPromise = postRawWebhook({ + url: `http://127.0.0.1:${address.port}/a2a/v1`, + body: "{", + contentLength: 2, + idleTimeoutMs: 60_000, + headers: { + "content-type": "application/json", + authorization: "Bearer alpha-secret", + }, + }); + + await requestReceived; + await vi.advanceTimersByTimeAsync(31_000); + const result = await resultPromise; + + expect(result.statusLine).toBe("HTTP/1.1 200 OK"); + expect(result.headers.connection).toBe("close"); + expect(JSON.parse(result.body)).toEqual({ + jsonrpc: "2.0", + id: null, + error: { code: -32000, message: "Request body could not be read" }, + }); + expect(result.closedByServer).toBe(true); + } finally { + vi.useRealTimers(); + server.closeAllConnections(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } }); }); diff --git a/extensions/a2a/src/http.ts b/extensions/a2a/src/http.ts index bf30885524ac..27ecd9e49174 100644 --- a/extensions/a2a/src/http.ts +++ b/extensions/a2a/src/http.ts @@ -7,7 +7,10 @@ import { isRequestBodyLimitError, readRequestBodyWithLimit, } from "openclaw/plugin-sdk/webhook-ingress"; -import { runDetachedWebhookWork } from "openclaw/plugin-sdk/webhook-request-guards"; +import { + runDetachedWebhookWork, + sendHttpRequestRejection, +} from "openclaw/plugin-sdk/webhook-request-guards"; import { A2aProtocolError, A2aRpcRequestSchema, @@ -270,13 +273,27 @@ export function createA2aHttpHandler(params: A2aHttpHandlerParams) { try { body = await readRequestBodyWithLimit(request, { maxBytes: MAX_REQUEST_BODY_BYTES, + // Defer destruction so the rejection below reaches the peer before the close. destroyOnLimit: false, }); } catch (error) { - if (isRequestBodyLimitError(error, "PAYLOAD_TOO_LARGE")) { - response.setHeader("connection", "close"); - response.once("finish", () => request.destroy()); - writeJsonResponse(response, 413, { error: "Request body exceeds the 1 MiB limit" }); + const bodyRejection = isRequestBodyLimitError(error, "PAYLOAD_TOO_LARGE") + ? { statusCode: 413, body: { error: "Request body exceeds the 1 MiB limit" } } + : isRequestBodyLimitError(error, "REQUEST_BODY_TIMEOUT") + ? { + statusCode: 200, + body: createRpcError(null, -32000, "Request body could not be read"), + } + : undefined; + if (bodyRejection) { + response.setHeader("cache-control", "no-store"); + await sendHttpRequestRejection( + request, + response, + bodyRejection.statusCode, + JSON.stringify(bodyRejection.body), + "application/json; charset=utf-8", + ); return true; } writeJsonResponse( diff --git a/extensions/admin-http-rpc/src/handler.test.ts b/extensions/admin-http-rpc/src/handler.test.ts index c8e99c2ae824..eb5168a91581 100644 --- a/extensions/admin-http-rpc/src/handler.test.ts +++ b/extensions/admin-http-rpc/src/handler.test.ts @@ -2,6 +2,7 @@ import { createServer, type Server } from "node:http"; import { connect, Socket } from "node:net"; import { Readable } from "node:stream"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { beforeEach, describe, expect, it, vi } from "vitest"; import { handleAdminHttpRpcRequest } from "./handler.js"; import { listAdminHttpRpcAllowedMethods } from "./methods.js"; @@ -357,6 +358,41 @@ describe("admin-http-rpc plugin handler", () => { } }); + it("flushes the 413 when the whole oversized body is already in flight", async () => { + // The declared-length case above answers before any body arrives. A sender that + // writes the whole over-cap body first leaves the rejection queued behind bytes the + // server has stopped reading, which is the case a hard teardown discards. + const server = createServer((req, res) => { + void handleAdminHttpRpcRequest(req, res); + }); + try { + const port = await listen(server); + const result = await postRawWebhook({ + url: `http://127.0.0.1:${port}/api/v1/admin/rpc`, + body: "x".repeat(1024 * 1024 + 128 * 1024), + headers: { "content-type": "application/json" }, + // The transport owner half-closes first and destroys on its bounded deadline, so + // observing the actual close needs a window wider than that deadline. + idleTimeoutMs: 3_000, + }); + + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(JSON.parse(result.body) as unknown).toEqual({ + ok: false, + error: { + type: "invalid_request", + message: "Payload too large", + }, + }); + expect(result.closedByServer).toBe(true); + expect(dispatchGatewayMethod).not.toHaveBeenCalled(); + } finally { + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } + }); + it("delivers a real HTTP 408 response before closing timed-out partial bodies", async () => { vi.useFakeTimers(); let markRequestStarted: (() => void) | undefined; diff --git a/extensions/discord/src/activities/http.test.ts b/extensions/discord/src/activities/http.test.ts index b8814864ac19..1fcb06f6be5e 100644 --- a/extensions/discord/src/activities/http.test.ts +++ b/extensions/discord/src/activities/http.test.ts @@ -5,6 +5,7 @@ import os from "node:os"; import path from "node:path"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import type { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { afterEach, describe, expect, it, vi } from "vitest"; import { buildDiscordActivityCustomId } from "../component-custom-id.js"; import { createDiscordActivityHttpHandler } from "./http.js"; @@ -281,6 +282,42 @@ describe("Discord Activity HTTP OAuth", () => { await expect(observeStalledTokenRequest(base, 1_000)).resolves.toBe("server-terminated"); }); + it.each([ + { + name: "413 when the token body exceeds its limit", + bodyTimeoutMs: 5_000, + // Declared and sent in one write: the shape whose rejection used to race the flush. + body: JSON.stringify({ code: "x".repeat(8 * 1024) }), + contentLength: undefined, + statusLine: "HTTP/1.1 413 Payload Too Large", + error: "request body too large", + }, + { + name: "408 when the sender stalls mid-upload", + bodyTimeoutMs: 50, + // Promises more than is ever sent, so the read deadline fires with the request open. + body: "{", + contentLength: 4 * 1024, + statusLine: "HTTP/1.1 408 Request Timeout", + error: "request body timeout", + }, + ])("delivers $name and then closes the connection", async (scenario) => { + const base = await startServer(createActivityTestRuntime(), { + bodyTimeoutMs: scenario.bodyTimeoutMs, + }); + + const result = await postRawWebhook({ + url: `${base}/discord/activity/api/token`, + body: scenario.body, + contentLength: scenario.contentLength, + headers: { "content-type": "application/json" }, + }); + + expect(result.statusLine).toBe(scenario.statusLine); + expect(JSON.parse(result.body)).toEqual({ error: scenario.error }); + expect(result.closedByServer).toBe(true); + }); + it("exchanges a code, creates a session, and uses it on the widget endpoint", async () => { const runtime = createActivityTestRuntime(); const widgetId = await createWidget(runtime); diff --git a/extensions/discord/src/activities/http.ts b/extensions/discord/src/activities/http.ts index 105ee2286839..57d545619188 100644 --- a/extensions/discord/src/activities/http.ts +++ b/extensions/discord/src/activities/http.ts @@ -4,6 +4,7 @@ import { logError } from "openclaw/plugin-sdk/logging-core"; import { resolveRequestClientIp } from "openclaw/plugin-sdk/webhook-ingress"; import { readJsonBodyWithLimit, + sendHttpRequestRejection, WEBHOOK_BODY_READ_DEFAULTS, } from "openclaw/plugin-sdk/webhook-request-guards"; import { parseDiscordActivityCustomId } from "../component-custom-id.js"; @@ -26,6 +27,7 @@ import { } from "./shell.js"; const BODY_MAX_BYTES = 8 * 1024; +const JSON_CONTENT_TYPE = "application/json; charset=utf-8"; const WIDGET_ID_PATTERN = /^[A-Za-z0-9_-]{22}$/; const DOC_TOKEN_PATTERN = /^[A-Za-z0-9_-]{43}$/; @@ -69,8 +71,12 @@ function respond( return true; } +function jsonBody(body: unknown): string { + return `${JSON.stringify(body)}\n`; +} + function respondJson(res: ServerResponse, statusCode: number, body: unknown): true { - return respond(res, statusCode, `${JSON.stringify(body)}\n`, "application/json; charset=utf-8"); + return respond(res, statusCode, jsonBody(body), JSON_CONTENT_TYPE); } function notFound(res: ServerResponse, widgetDocument = false): true { @@ -158,9 +164,28 @@ export function createDiscordActivityHttpHandler(deps: DiscordActivityHttpDeps): maxBytes: BODY_MAX_BYTES, timeoutMs: bodyTimeoutMs, emptyObjectOnEmpty: true, + // Defer destruction so the rejections below reach the client before the close. + destroyOnLimit: false, }); if (!bodyResult.ok && bodyResult.code === "REQUEST_BODY_TIMEOUT") { - return respondJson(res, 408, { error: "request body timeout" }); + await sendHttpRequestRejection( + req, + res, + 408, + jsonBody({ error: "request body timeout" }), + JSON_CONTENT_TYPE, + ); + return true; + } + if (!bodyResult.ok && bodyResult.code === "PAYLOAD_TOO_LARGE") { + await sendHttpRequestRejection( + req, + res, + 413, + jsonBody({ error: "request body too large" }), + JSON_CONTENT_TYPE, + ); + return true; } const body = bodyResult.ok && diff --git a/extensions/line/src/monitor.ts b/extensions/line/src/monitor.ts index 83d0036d72f2..589185052987 100644 --- a/extensions/line/src/monitor.ts +++ b/extensions/line/src/monitor.ts @@ -14,11 +14,9 @@ import { } from "openclaw/plugin-sdk/runtime-env"; import { canonicalizeWebhookRouteKey, - isRequestBodyLimitError, normalizePluginHttpPath, normalizeWebhookPath, registerWebhookTargetWithPluginRoute, - requestBodyErrorToText, resolveSingleWebhookTarget, } from "openclaw/plugin-sdk/webhook-ingress"; import { @@ -42,7 +40,11 @@ import { } from "./send.js"; import { buildTemplateMessageFromPayload } from "./template-messages.js"; import type { LineChannelData, ResolvedLineAccount } from "./types.js"; -import { createLineNodeWebhookHandler, readLineWebhookRequestBody } from "./webhook-node.js"; +import { + createLineNodeWebhookHandler, + readLineWebhookRequestBody, + rejectLineWebhookRequest, +} from "./webhook-node.js"; import { LineWebhookTerminalDeliveryError } from "./webhook-spool.js"; import { parseLineWebhookBody, validateLineSignature } from "./webhook-utils.js"; @@ -426,16 +428,7 @@ export async function monitorLineProvider( res.setHeader("Content-Type", "application/json"); res.end(JSON.stringify({ status: "ok" })); } catch (err) { - if (isRequestBodyLimitError(err, "PAYLOAD_TOO_LARGE")) { - res.statusCode = 413; - res.setHeader("Content-Type", "application/json"); - res.end(JSON.stringify({ error: "Payload too large" })); - return; - } - if (isRequestBodyLimitError(err, "REQUEST_BODY_TIMEOUT")) { - res.statusCode = 408; - res.setHeader("Content-Type", "application/json"); - res.end(JSON.stringify({ error: requestBodyErrorToText("REQUEST_BODY_TIMEOUT") })); + if (await rejectLineWebhookRequest(req, res, err)) { return; } runtime.error?.(danger(`line webhook error: ${formatErrorMessage(err)}`)); diff --git a/extensions/line/src/monitor.webhook-wire.test.ts b/extensions/line/src/monitor.webhook-wire.test.ts new file mode 100644 index 000000000000..4586040e5060 --- /dev/null +++ b/extensions/line/src/monitor.webhook-wire.test.ts @@ -0,0 +1,193 @@ +// Line tests cover webhook body-limit answers as the sender receives them on the wire. +import crypto from "node:crypto"; +import { createServer, type IncomingMessage, type ServerResponse } from "node:http"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import type { RuntimeEnv } from "openclaw/plugin-sdk/runtime-env"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; +import { afterAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { monitorLineProvider } from "./monitor.js"; + +type LineNodeWebhookHandler = (req: IncomingMessage, res: ServerResponse) => Promise; + +const { createLineBotMock, registerWebhookTargetWithPluginRouteMock } = vi.hoisted(() => ({ + createLineBotMock: vi.fn(), + registerWebhookTargetWithPluginRouteMock: vi.fn(), +})); + +vi.mock("./bot.js", () => ({ createLineBot: createLineBotMock })); + +vi.mock("openclaw/plugin-sdk/webhook-ingress", async () => { + const actual = await vi.importActual( + "openclaw/plugin-sdk/webhook-ingress", + ); + return { + ...actual, + normalizePluginHttpPath: (path: string | undefined, fallback: string) => path ?? fallback, + registerWebhookTargetWithPluginRoute: registerWebhookTargetWithPluginRouteMock, + }; +}); + +function requireRegisteredHandler(): LineNodeWebhookHandler { + const registration = registerWebhookTargetWithPluginRouteMock.mock.calls[0]?.[0] as + | { route?: { handler?: LineNodeWebhookHandler } } + | undefined; + const handler = registration?.route?.handler; + if (!handler) { + throw new Error("expected registered LINE webhook route"); + } + return handler; +} + +function signLineWebhook(body: string): string { + return crypto.createHmac("SHA256", "secret").update(body).digest("base64"); +} + +describe("monitorLineProvider webhook body limits over a real connection", () => { + beforeEach(() => { + createLineBotMock.mockReset().mockImplementation(() => ({ + account: { accountId: "default" }, + handleWebhook: vi.fn(async () => "durable" as const), + stop: vi.fn(async () => undefined), + })); + // The monitor matches an inbound signature against the targets it registered, so the + // registration seam has to keep that map for the accepted case to reach an account. + registerWebhookTargetWithPluginRouteMock.mockReset().mockImplementation((params) => { + const key = params.target.path.toLowerCase(); + const normalizedTarget = { ...params.target, path: key }; + params.targetsByPath.set(key, [...(params.targetsByPath.get(key) ?? []), normalizedTarget]); + return { + target: normalizedTarget, + unregister: () => params.targetsByPath.delete(key), + }; + }); + }); + + afterAll(() => { + vi.doUnmock("./bot.js"); + vi.doUnmock("openclaw/plugin-sdk/webhook-ingress"); + vi.resetModules(); + }); + + // A real socket is the only place both halves of the contract are observable: the handler + // answers while the sender is still uploading, and then closes the connection. + const withLineWebhookWire = async (run: (webhookUrl: string) => Promise) => { + const monitor = await monitorLineProvider({ + channelAccessToken: "token", + channelSecret: "secret", // pragma: allowlist secret + accountId: "default", + config: {} as OpenClawConfig, + runtime: {} as RuntimeEnv, + }); + const handler = requireRegisteredHandler(); + const server = createServer((req, res) => { + void handler(req, res); + }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected LINE webhook test server to have a TCP address"); + } + await run(`http://127.0.0.1:${address.port}/line/webhook`); + } finally { + try { + if (server.listening) { + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } + } finally { + await monitor.stop(); + } + } + }; + + it("answers an over-limit webhook with 413 and then closes the connection", async () => { + await withLineWebhookWire(async (webhookUrl) => { + const oversizedPayload = JSON.stringify({ + events: [{ type: "message" }], + padding: "x".repeat(70 * 1024), + }); + const oversized = await postRawWebhook({ + url: webhookUrl, + body: oversizedPayload, + headers: { + "content-type": "application/json", + "x-line-signature": signLineWebhook(oversizedPayload), + }, + }); + expect(oversized.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(JSON.parse(oversized.body)).toEqual({ error: "Payload too large" }); + expect(oversized.closedByServer).toBe(true); + + // Same route: an in-limit webhook is still admitted and answered normally. + const acceptedPayload = JSON.stringify({ events: [{ type: "message" }] }); + const accepted = await fetch(webhookUrl, { + method: "POST", + headers: { + "content-type": "application/json", + "x-line-signature": signLineWebhook(acceptedPayload), + }, + body: acceptedPayload, + }); + expect(accepted.status).toBe(200); + expect(await accepted.json()).toEqual({ status: "ok" }); + }); + }); + + it("counts UTF-8 bytes, not characters, for a chunked non-ASCII webhook", async () => { + // LINE payloads are routinely non-ASCII, and chunk sizes are byte counts. Sending the + // framing and the body as one string would re-encode every byte above 0x7f after the + // declared length was computed, so the under-cap control below only reaches the handler + // when the driver is byte-exact. + await withLineWebhookWire(async (webhookUrl) => { + // 30,000 characters is well under the 65,536-byte cap as characters, and well over it + // as UTF-8. A character-counting limit would admit this body. + const overCap = JSON.stringify({ + events: [{ type: "message" }], + padding: "壓".repeat(30_000), + }); + expect(overCap.length).toBeLessThan(64 * 1024); + expect(Buffer.byteLength(overCap, "utf-8")).toBeGreaterThan(64 * 1024); + const rejected = await postRawWebhook({ + url: webhookUrl, + body: overCap, + chunkedEncoding: true, + chunk: { bytes: 8 * 1024, intervalMs: 1 }, + headers: { + "content-type": "application/json", + "x-line-signature": signLineWebhook(overCap), + }, + }); + expect(rejected.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(JSON.parse(rejected.body)).toEqual({ error: "Payload too large" }); + expect(rejected.closedByServer).toBe(true); + + // Control: the same non-ASCII chunked shape under the cap is signed, parsed, and + // answered, so the rejection above is attributable to the byte count rather than to + // chunked framing the server could not read. + const underCap = JSON.stringify({ + events: [{ type: "message" }], + padding: "壓".repeat(1_000), + }); + const accepted = await postRawWebhook({ + url: webhookUrl, + body: underCap, + chunkedEncoding: true, + chunk: { bytes: 1024, intervalMs: 1 }, + headers: { + "content-type": "application/json", + "x-line-signature": signLineWebhook(underCap), + }, + }); + expect(accepted.statusLine).toBe("HTTP/1.1 200 OK"); + expect(JSON.parse(accepted.body)).toEqual({ status: "ok" }); + }); + }); +}); diff --git a/extensions/line/src/webhook-node.ts b/extensions/line/src/webhook-node.ts index 214ab3012a19..73ef22e3a6fa 100644 --- a/extensions/line/src/webhook-node.ts +++ b/extensions/line/src/webhook-node.ts @@ -6,6 +6,7 @@ import { isRequestBodyLimitError, readRequestBodyWithLimit, requestBodyErrorToText, + sendHttpRequestRejection, } from "openclaw/plugin-sdk/webhook-request-guards"; import type { createLineBot } from "./bot.js"; import { parseLineWebhookBody, validateLineSignature } from "./webhook-utils.js"; @@ -22,11 +23,41 @@ export async function readLineWebhookRequestBody( return await readRequestBodyWithLimit(req, { maxBytes, timeoutMs, + // Defer destruction so the caller can answer 413/408 before the connection closes. + destroyOnLimit: false, }); } type ReadBodyFn = (req: IncomingMessage, maxBytes: number, timeoutMs?: number) => Promise; +/** + * Answer a body-limit failure through the connection owner. + * + * The reader defers destruction for these two codes, so the connection is already fenced + * and only the owner can still write: responding directly would race the teardown and LINE + * would see a reset instead of the status. + */ +export async function rejectLineWebhookRequest( + req: IncomingMessage, + res: ServerResponse, + error: unknown, +): Promise { + if ( + !isRequestBodyLimitError(error, "PAYLOAD_TOO_LARGE") && + !isRequestBodyLimitError(error, "REQUEST_BODY_TIMEOUT") + ) { + return false; + } + await sendHttpRequestRejection( + req, + res, + error.statusCode, + JSON.stringify({ error: requestBodyErrorToText(error.code) }), + "application/json", + ); + return true; +} + export function createLineNodeWebhookHandler(params: { channelSecret: string; bot: Pick, "handleWebhook">; @@ -108,16 +139,7 @@ export function createLineNodeWebhookHandler(params: { res.setHeader("Content-Type", "application/json"); res.end(JSON.stringify({ status: "ok" })); } catch (err) { - if (isRequestBodyLimitError(err, "PAYLOAD_TOO_LARGE")) { - res.statusCode = 413; - res.setHeader("Content-Type", "application/json"); - res.end(JSON.stringify({ error: "Payload too large" })); - return; - } - if (isRequestBodyLimitError(err, "REQUEST_BODY_TIMEOUT")) { - res.statusCode = 408; - res.setHeader("Content-Type", "application/json"); - res.end(JSON.stringify({ error: requestBodyErrorToText("REQUEST_BODY_TIMEOUT") })); + if (await rejectLineWebhookRequest(req, res, err)) { return; } params.runtime.error?.(danger(`line webhook error: ${formatErrorMessage(err)}`)); diff --git a/extensions/mattermost/src/mattermost/interactions.test.ts b/extensions/mattermost/src/mattermost/interactions.test.ts index 46ecede6af68..618ae0728781 100644 --- a/extensions/mattermost/src/mattermost/interactions.test.ts +++ b/extensions/mattermost/src/mattermost/interactions.test.ts @@ -1,5 +1,6 @@ // Mattermost tests cover interactions plugin behavior. -import type { IncomingMessage, ServerResponse } from "node:http"; +import { createServer, type IncomingMessage, type ServerResponse } from "node:http"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { describe, expect, it, beforeEach, afterEach, vi } from "vitest"; import type { PluginRuntime } from "../../runtime-api.js"; import { setMattermostRuntime } from "../runtime.js"; @@ -981,3 +982,50 @@ describe("createMattermostInteractionHandler", () => { expect(dispatchButtonClick).not.toHaveBeenCalled(); }); }); + +describe("createMattermostInteractionHandler body limits", () => { + // The 10s read deadline makes the sibling 408 case too slow to drive here; slash-http + // and Synology Chat cover that half of the same branch. + it("delivers 413 for an over-limit callback body and then closes the connection", async () => { + const handleInteraction = vi.fn(); + const handler = createMattermostInteractionHandler({ + client: {} as MattermostClient, + botUserId: "bot", + accountId: "acct", + handleInteraction, + }); + const server = createServer((req, res) => { + void handler(req, res); + }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected the interaction test server to have a TCP address"); + } + + // Declared and sent in one write: the shape whose rejection used to race the flush. + const result = await postRawWebhook({ + url: `http://127.0.0.1:${address.port}/mattermost/interactions`, + body: "x".repeat(64 * 1024 + 1), + headers: { "content-type": "application/json" }, + }); + + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(JSON.parse(result.body)).toEqual({ error: "Payload too large" }); + expect(result.closedByServer).toBe(true); + expect(handleInteraction).not.toHaveBeenCalled(); + } finally { + server.closeAllConnections(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } + }); +}); diff --git a/extensions/mattermost/src/mattermost/interactions.ts b/extensions/mattermost/src/mattermost/interactions.ts index 96ad9d530db7..a26ff51cc394 100644 --- a/extensions/mattermost/src/mattermost/interactions.ts +++ b/extensions/mattermost/src/mattermost/interactions.ts @@ -11,9 +11,11 @@ import { getMattermostRuntime } from "../runtime.js"; import { isWildcardBindHost } from "./callback-host.js"; import { updateMattermostPost, type MattermostClient, type MattermostPost } from "./client.js"; import { + isRequestBodyLimitError, isTrustedProxyAddress, readRequestBodyWithLimit, resolveClientIp, + sendHttpRequestRejection, type OpenClawConfig, } from "./runtime-api.js"; @@ -348,6 +350,8 @@ function readInteractionBody(req: IncomingMessage): Promise { return readRequestBodyWithLimit(req, { maxBytes: INTERACTION_MAX_BODY_BYTES, timeoutMs: INTERACTION_BODY_TIMEOUT_MS, + // Defer destruction so the rejection below reaches Mattermost before the close. + destroyOnLimit: false, }); } @@ -433,6 +437,26 @@ export function createMattermostInteractionHandler(params: { payload = parseInteractionPayload(raw); } catch (err) { log?.(`mattermost interaction: failed to parse body: ${String(err)}`); + if (isRequestBodyLimitError(err, "PAYLOAD_TOO_LARGE")) { + await sendHttpRequestRejection( + req, + res, + 413, + JSON.stringify({ error: "Payload too large" }), + "application/json", + ); + return; + } + if (isRequestBodyLimitError(err, "REQUEST_BODY_TIMEOUT")) { + await sendHttpRequestRejection( + req, + res, + 408, + JSON.stringify({ error: "Request body timeout" }), + "application/json", + ); + return; + } res.statusCode = 400; res.setHeader("Content-Type", "application/json"); res.end(JSON.stringify({ error: "Invalid request body" })); diff --git a/extensions/mattermost/src/mattermost/runtime-api.ts b/extensions/mattermost/src/mattermost/runtime-api.ts index 344ff18eaa66..731825386944 100644 --- a/extensions/mattermost/src/mattermost/runtime-api.ts +++ b/extensions/mattermost/src/mattermost/runtime-api.ts @@ -32,8 +32,9 @@ export { createChannelHistoryWindow, } from "openclaw/plugin-sdk/reply-history"; export { registerPluginHttpRoute } from "openclaw/plugin-sdk/webhook-targets"; +export { isRequestBodyLimitError } from "openclaw/plugin-sdk/webhook-ingress"; export { - isRequestBodyLimitError, readRequestBodyWithLimit, -} from "openclaw/plugin-sdk/webhook-ingress"; + sendHttpRequestRejection, +} from "openclaw/plugin-sdk/webhook-request-guards"; export { isTrustedProxyAddress, resolveClientIp } from "openclaw/plugin-sdk/core"; diff --git a/extensions/mattermost/src/mattermost/slash-http.send-config.test.ts b/extensions/mattermost/src/mattermost/slash-http.send-config.test.ts index f2530fd4adfd..891f2c8cc1c7 100644 --- a/extensions/mattermost/src/mattermost/slash-http.send-config.test.ts +++ b/extensions/mattermost/src/mattermost/slash-http.send-config.test.ts @@ -69,6 +69,7 @@ vi.mock("./runtime-api.js", () => { formatInboundFromLabel: vi.fn(() => ""), rawDataToString: vi.fn((value: unknown) => (typeof value === "string" ? value : "")), readRequestBodyWithLimit: mockState.readRequestBodyWithLimit, + sendHttpRequestRejection: vi.fn(async () => undefined), resolveThreadSessionKeys: vi.fn((params: { baseSessionKey: string }) => ({ sessionKey: params.baseSessionKey, parentSessionKey: undefined, diff --git a/extensions/mattermost/src/mattermost/slash-http.test.ts b/extensions/mattermost/src/mattermost/slash-http.test.ts index 0f48e6fe11e8..5b7bce6545dc 100644 --- a/extensions/mattermost/src/mattermost/slash-http.test.ts +++ b/extensions/mattermost/src/mattermost/slash-http.test.ts @@ -1,6 +1,7 @@ // Mattermost tests cover slash http plugin behavior. -import { IncomingMessage, type ServerResponse } from "node:http"; +import { createServer, IncomingMessage, type ServerResponse } from "node:http"; import { Socket } from "node:net"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { beforeEach, describe, expect, it, vi } from "vitest"; import type { OpenClawConfig, RuntimeEnv } from "../../runtime-api.js"; import type { ResolvedMattermostAccount } from "./accounts.js"; @@ -330,21 +331,65 @@ describe("slash-http", () => { expect(response.getBody()).toContain("Unauthorized: invalid command token."); }); - it("returns 408 when the request body stalls", async () => { + it.each([ + { + name: "413 when the upload exceeds the body limit", + bodyTimeoutMs: 5_000, + // Declared and sent in one write: the shape whose rejection used to race the flush. + body: "x".repeat(64 * 1024 + 1), + contentLength: undefined, + statusLine: "HTTP/1.1 413 Payload Too Large", + responseBody: "Payload Too Large", + }, + { + name: "408 when the sender stalls mid-upload", + bodyTimeoutMs: 50, + // Promises more than is ever sent, so the read deadline fires with the request open. + body: "x".repeat(16), + contentLength: 64 * 1024, + statusLine: "HTTP/1.1 408 Request Timeout", + responseBody: "Request body timeout", + }, + ])("delivers $name and then closes the connection", async (scenario) => { const handler = createSlashCommandHttpHandler({ account: accountFixture, cfg: {} as OpenClawConfig, runtime: {} as RuntimeEnv, registeredCommands: [createRegisteredCommand()], - bodyTimeoutMs: 1, + bodyTimeoutMs: scenario.bodyTimeoutMs, }); - const req = createRequest({ autoEnd: false }); - const response = createResponse(); + const server = createServer((req, res) => { + void handler(req, res); + }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected the slash-command test server to have a TCP address"); + } - await handler(req, response.res); + const result = await postRawWebhook({ + url: `http://127.0.0.1:${address.port}/slash`, + body: scenario.body, + contentLength: scenario.contentLength, + headers: { "content-type": "application/x-www-form-urlencoded" }, + }); - expect(response.res.statusCode).toBe(408); - expect(response.getBody()).toBe("Request body timeout"); + expect(result.statusLine).toBe(scenario.statusLine); + expect(result.body).toBe(scenario.responseBody); + expect(result.closedByServer).toBe(true); + } finally { + server.closeAllConnections(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } }); it("rejects the startup token when Mattermost has rotated the current command token", async () => { diff --git a/extensions/mattermost/src/mattermost/slash-http.ts b/extensions/mattermost/src/mattermost/slash-http.ts index cfa41b6f56f3..264988f99b44 100644 --- a/extensions/mattermost/src/mattermost/slash-http.ts +++ b/extensions/mattermost/src/mattermost/slash-http.ts @@ -40,6 +40,7 @@ import { isRequestBodyLimitError, logTypingFailure, readRequestBodyWithLimit, + sendHttpRequestRejection, type OpenClawConfig, type RuntimeEnv, } from "./runtime-api.js"; @@ -111,6 +112,8 @@ function readBody( return readRequestBodyWithLimit(req, { maxBytes, timeoutMs, + // Defer destruction so the rejections below reach Mattermost before the close. + destroyOnLimit: false, }); } @@ -593,12 +596,10 @@ export function createSlashCommandHttpHandler(params: SlashHttpHandlerParams) { body = bufferedBody ?? (await readBody(req, MAX_BODY_BYTES, bodyTimeoutMs)); } catch (error) { if (isRequestBodyLimitError(error, "REQUEST_BODY_TIMEOUT")) { - res.statusCode = 408; - res.end("Request body timeout"); + await sendHttpRequestRejection(req, res, 408, "Request body timeout"); return; } - res.statusCode = 413; - res.end("Payload Too Large"); + await sendHttpRequestRejection(req, res, 413, "Payload Too Large"); return; } diff --git a/extensions/mattermost/src/mattermost/slash-state.ts b/extensions/mattermost/src/mattermost/slash-state.ts index 179140674443..79e5ef44409c 100644 --- a/extensions/mattermost/src/mattermost/slash-state.ts +++ b/extensions/mattermost/src/mattermost/slash-state.ts @@ -16,6 +16,7 @@ import type { ResolvedMattermostAccount } from "./accounts.js"; import { isRequestBodyLimitError, readRequestBodyWithLimit, + sendHttpRequestRejection, type OpenClawPluginApi, } from "./runtime-api.js"; import { @@ -334,15 +335,15 @@ export function registerSlashCommandRoute(api: OpenClawPluginApi) { bodyStr = await readRequestBodyWithLimit(req, { maxBytes: MULTI_ACCOUNT_BODY_MAX_BYTES, timeoutMs: MULTI_ACCOUNT_BODY_TIMEOUT_MS, + // Defer destruction so the rejections below reach Mattermost before the close. + destroyOnLimit: false, }); } catch (error) { if (isRequestBodyLimitError(error, "REQUEST_BODY_TIMEOUT")) { - res.statusCode = 408; - res.end("Request body timeout"); + await sendHttpRequestRejection(req, res, 408, "Request body timeout"); return; } - res.statusCode = 413; - res.end("Payload Too Large"); + await sendHttpRequestRejection(req, res, 413, "Payload Too Large"); return; } diff --git a/extensions/nextcloud-talk/src/monitor.replay.test.ts b/extensions/nextcloud-talk/src/monitor.replay.test.ts index 087b2de0642d..2466ab7f978b 100644 --- a/extensions/nextcloud-talk/src/monitor.replay.test.ts +++ b/extensions/nextcloud-talk/src/monitor.replay.test.ts @@ -1,6 +1,6 @@ // Nextcloud Talk tests cover monitor.replay plugin behavior. import type { IncomingMessage, ServerResponse } from "node:http"; -import { createMockIncomingRequest } from "openclaw/plugin-sdk/test-env"; +import { createMockIncomingRequest, postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { describe, expect, it, vi } from "vitest"; import { createNextcloudTalkWebhookServer as createRawNextcloudTalkWebhookServer } from "./monitor.js"; import { createSignedCreateMessageRequest } from "./monitor.test-fixtures.js"; @@ -59,38 +59,6 @@ async function invokeWebhookRequestListener(params: { }); } -async function invokeWebhookServerRequest(params: { - body: string; - headers: Record; - maxBodyBytes: number; -}) { - const { server, stop } = createNextcloudTalkWebhookServer({ - host: "127.0.0.1", - port: 0, - path: "/nextcloud-body-limit", - secret: "nextcloud-secret", // pragma: allowlist secret - maxBodyBytes: params.maxBodyBytes, - onMessage: vi.fn(), - }); - try { - const listener = server.listeners("request")[0] as - | ((req: IncomingMessage, res: ServerResponse) => void) - | undefined; - if (!listener) { - throw new Error("expected Nextcloud Talk request listener"); - } - return await invokeWebhookRequestListener({ - listener, - path: "/nextcloud-body-limit", - body: params.body, - headers: params.headers, - remoteAddress: "127.0.0.1", - }); - } finally { - await stop(); - } -} - describe("createNextcloudTalkWebhookServer auth order", () => { it("closes when abort races with listener startup", async () => { const abortController = new AbortController(); @@ -134,19 +102,6 @@ describe("createNextcloudTalkWebhookServer auth order", () => { expect(await response.json()).toEqual({ error: "Missing signature headers" }); expect(readBody).not.toHaveBeenCalled(); }); - - it("rejects signed payloads over the configured body limit", async () => { - const { body, headers } = createSignedCreateMessageRequest(); - - const response = await invokeWebhookServerRequest({ - body, - headers, - maxBodyBytes: 128, - }); - - expect(response.status).toBe(413); - expect(JSON.parse(response.body)).toEqual({ error: "Payload too large" }); - }); }); describe("createNextcloudTalkWebhookServer backend allowlist", () => { @@ -249,6 +204,39 @@ describe("createNextcloudTalkWebhookServer payload validation", () => { expect(onMessage).not.toHaveBeenCalled(); }); + it("answers an over-limit webhook with 413 and then closes the connection", async () => { + // Driven over a raw socket rather than fetch: the server answers while the sender is + // still uploading and then closes, so both halves of the contract - the status is + // delivered, and the rejected request does not stay open - have to be observed on the + // wire. A mocked response records status(413) either way and proves neither half. + const body = JSON.stringify({ type: "Create", padding: "x".repeat(70 * 1024) }); + const { random, signature } = generateNextcloudTalkSignature({ + body, + secret: "nextcloud-secret", // pragma: allowlist secret + }); + const onMessage = vi.fn(); + const harness = await startWebhookServer({ + path: "/nextcloud-oversized-body", + onMessage, + }); + + const result = await postRawWebhook({ + url: harness.webhookUrl, + body, + headers: { + "content-type": "application/json", + "x-nextcloud-talk-random": random, + "x-nextcloud-talk-signature": signature, + "x-nextcloud-talk-backend": "https://nextcloud.example", + }, + }); + + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(result.body).toBe(JSON.stringify({ error: "Payload too large" })); + expect(result.closedByServer).toBe(true); + expect(onMessage).not.toHaveBeenCalled(); + }); + it("acknowledges signed non-Create Talk events instead of rejecting them", async () => { const payload = { type: "Join", diff --git a/extensions/nextcloud-talk/src/monitor.ts b/extensions/nextcloud-talk/src/monitor.ts index 6036b82bc4bd..c4c245a234ca 100644 --- a/extensions/nextcloud-talk/src/monitor.ts +++ b/extensions/nextcloud-talk/src/monitor.ts @@ -9,6 +9,7 @@ import { resolveRequestClientIp, requestBodyErrorToText, } from "openclaw/plugin-sdk/webhook-ingress"; +import { sendHttpRequestRejection } from "openclaw/plugin-sdk/webhook-request-guards"; import { extractNextcloudTalkHeaders, verifyNextcloudTalkSignature } from "./signature.js"; import type { NextcloudTalkWebhookHeaders, NextcloudTalkWebhookServerOptions } from "./types.js"; import { NextcloudTalkWebhookPayloadError } from "./webhook-spool-state.js"; @@ -98,9 +99,23 @@ function readNextcloudTalkWebhookBody(req: IncomingMessage, maxBodyBytes: number // body budget bounded even if the operator-configured post-parse limit is larger. maxBytes: Math.min(maxBodyBytes, PREAUTH_WEBHOOK_MAX_BODY_BYTES), timeoutMs: PREAUTH_WEBHOOK_BODY_TIMEOUT_MS, + // Defer destruction so the rejections below reach the backend before the close. + destroyOnLimit: false, }); } +async function rejectWebhookRequest( + req: IncomingMessage, + res: ServerResponse, + status: number, + error: string, +): Promise { + if (res.headersSent) { + return; + } + await sendHttpRequestRejection(req, res, status, JSON.stringify({ error }), "application/json"); +} + export function createNextcloudTalkWebhookServer(opts: NextcloudTalkWebhookServerOptions): { server: Server; start: () => Promise; @@ -193,11 +208,11 @@ export function createNextcloudTalkWebhookServer(opts: NextcloudTalkWebhookServe writeJsonResponse(res, 200); } catch (err) { if (isRequestBodyLimitError(err, "PAYLOAD_TOO_LARGE")) { - writeWebhookError(res, 413, WEBHOOK_ERRORS.payloadTooLarge); + await rejectWebhookRequest(req, res, 413, WEBHOOK_ERRORS.payloadTooLarge); return; } if (isRequestBodyLimitError(err, "REQUEST_BODY_TIMEOUT")) { - writeWebhookError(res, 408, requestBodyErrorToText("REQUEST_BODY_TIMEOUT")); + await rejectWebhookRequest(req, res, 408, requestBodyErrorToText("REQUEST_BODY_TIMEOUT")); return; } if (err instanceof NextcloudTalkWebhookPayloadError) { diff --git a/extensions/openai/realtime-quicksilver-offer-http.ts b/extensions/openai/realtime-quicksilver-offer-http.ts index 1ba8180e4550..f7f4c4b6a2df 100644 --- a/extensions/openai/realtime-quicksilver-offer-http.ts +++ b/extensions/openai/realtime-quicksilver-offer-http.ts @@ -1,29 +1,18 @@ +// Transport helpers for the GPT-Live browser offer endpoint. import type { IncomingMessage, ServerResponse } from "node:http"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; -import { resolveAcceptedBrowserOrigin } from "openclaw/plugin-sdk/webhook-request-guards"; +import { + isRequestBodyLimitError, + requestBodyErrorToText, + resolveAcceptedBrowserOrigin, + sendHttpRequestRejection, +} from "openclaw/plugin-sdk/webhook-request-guards"; type ResponseDeliveryWaiter = { result: Promise; cancel: () => void; }; -export function applyRealtimeOfferCorsHeaders( - req: IncomingMessage, - res: ServerResponse, - cfg: OpenClawConfig | undefined, -): boolean { - if (!req.headers.origin) { - return true; - } - const origin = resolveAcceptedBrowserOrigin({ req, cfg }); - if (!origin) { - return false; - } - res.setHeader("Access-Control-Allow-Origin", origin); - res.setHeader("Vary", "Origin"); - return true; -} - export function createResponseDeliveryWaiter( res: ServerResponse, onDelivered: () => void, @@ -46,11 +35,17 @@ export function createResponseDeliveryWaiter( return { result, cancel: () => settle(false) }; } +const OFFER_TEXT_CONTENT_TYPE = "text/plain; charset=utf-8"; + +/** + * The one writer for every offer-endpoint response, error text and SDP answer + * alike, so no caller sets the security headers by hand and forgets one. + */ export function respondRealtimeOffer( res: ServerResponse, statusCode: number, body: string, - contentType = "text/plain; charset=utf-8", + contentType = OFFER_TEXT_CONTENT_TYPE, ): void { res.statusCode = statusCode; res.setHeader("cache-control", "no-store"); @@ -58,3 +53,52 @@ export function respondRealtimeOffer( res.setHeader("x-content-type-options", "nosniff"); res.end(body); } + +export function applyRealtimeOfferCorsHeaders( + req: IncomingMessage, + res: ServerResponse, + cfg: OpenClawConfig | undefined, +): boolean { + if (!req.headers.origin) { + return true; + } + const origin = resolveAcceptedBrowserOrigin({ req, cfg }); + if (!origin) { + return false; + } + res.setHeader("Access-Control-Allow-Origin", origin); + res.setHeader("Vary", "Origin"); + return true; +} + +export function readOfferBearerToken(req: IncomingMessage): string | undefined { + const authorization = req.headers.authorization?.trim(); + return authorization?.match(/^Bearer\s+([^\s]+)$/i)?.[1]; +} + +/** + * Answer an over-limit or stalled offer through the connection owner. + * + * The reader defers destruction for those two codes, so the connection is already fenced + * and only the owner can still write; responding directly would race the teardown and the + * browser would see a reset instead of the status. + */ +export async function rejectOversizedOffer( + req: IncomingMessage, + res: ServerResponse, + error: unknown, +): Promise { + if (!isRequestBodyLimitError(error)) { + return false; + } + res.setHeader("cache-control", "no-store"); + res.setHeader("x-content-type-options", "nosniff"); + await sendHttpRequestRejection( + req, + res, + error.statusCode, + requestBodyErrorToText(error.code), + OFFER_TEXT_CONTENT_TYPE, + ); + return true; +} diff --git a/extensions/openai/realtime-quicksilver-session.transport.test.ts b/extensions/openai/realtime-quicksilver-session.transport.test.ts new file mode 100644 index 000000000000..d3f72fd7da9f --- /dev/null +++ b/extensions/openai/realtime-quicksilver-session.transport.test.ts @@ -0,0 +1,118 @@ +/** + * Transport-level coverage for the GPT-Live offer endpoint. + * + * A mocked response records a status even when the socket died first, so the + * answered-then-released contract is only observable over a real socket. + */ +import { createServer, type Server } from "node:http"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { OPENAI_QUICKSILVER_OFFER_PATH } from "./realtime-quicksilver-session.js"; +import { createBroker } from "./realtime-quicksilver.test-helpers.js"; + +afterEach(() => { + vi.restoreAllMocks(); + vi.unstubAllEnvs(); +}); + +async function listenOnLoopback(server: Server): Promise { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("Expected the realtime offer test server to have a TCP address"); + } + return address.port; +} + +async function closeServer(server: Server): Promise { + if (!server.listening) { + return; + } + await new Promise((resolve, reject) => { + server.close((error) => { + if (error) { + reject(error); + return; + } + resolve(); + }); + }); +} + +describe("GPT-Live offer transport", () => { + it("delivers a rejected offer's response before releasing the connection", async () => { + const { realtime } = createBroker(); + const server = createServer((req, res) => { + void realtime.handler(req, res); + }); + try { + const port = await listenOnLoopback(server); + const url = `http://127.0.0.1:${port}${OPENAI_QUICKSILVER_OFFER_PATH}`; + // Offer credentials are single-use, so every request needs its own reservation. + const offerHeaders = async (): Promise> => { + const reservation = await realtime.broker.createBrowserSession( + { + providerConfig: {}, + model: "gpt-live-1", + runAgentConsult: vi.fn(async () => ({ text: "Done" })), + }, + { type: "api-key", token: "platform-key" }, + ); + if (reservation.transport !== "webrtc") { + throw new Error("Expected WebRTC reservation"); + } + return { + authorization: `Bearer ${reservation.clientSecret}`, + "content-type": "application/sdp", + }; + }; + + // Declared over-cap length: rejected from the header alone, while the browser is + // still mid-upload and can only learn the outcome from what reaches it. + const declared = await postRawWebhook({ + url, + body: "v=offer", + contentLength: 300 * 1024, + headers: await offerHeaders(), + idleTimeoutMs: 500, + }); + expect(declared.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(declared.body).toBe("Payload too large"); + expect(declared.closedByServer).toBe(true); + + // Chunked, so there is no declared length and the cap can only be hit by counting + // the bytes that actually arrive. + const streamed = await postRawWebhook({ + url, + body: "v".repeat(300 * 1024), + headers: await offerHeaders(), + chunkedEncoding: true, + chunk: { bytes: 32 * 1024, intervalMs: 5 }, + idleTimeoutMs: 500, + }); + expect(streamed.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(streamed.closedByServer).toBe(true); + + // Control: an under-cap offer is still read and answered on a retained connection, + // so the teardown above is attributable to the rejection, not to every offer. + const withinCap = await postRawWebhook({ + url, + body: " ", + headers: await offerHeaders(), + idleTimeoutMs: 500, + }); + expect(withinCap.statusLine).toBe("HTTP/1.1 400 Bad Request"); + expect(withinCap.body).toContain("SDP offer is required"); + expect(withinCap.closedByServer).toBe(false); + } finally { + await closeServer(server); + await realtime.cleanup(); + } + }); +}); diff --git a/extensions/openai/realtime-quicksilver-session.ts b/extensions/openai/realtime-quicksilver-session.ts index 0f053ac1bf40..6669dc50d9f1 100644 --- a/extensions/openai/realtime-quicksilver-session.ts +++ b/extensions/openai/realtime-quicksilver-session.ts @@ -20,6 +20,8 @@ import { OpenAIQuicksilverDelegationController } from "./realtime-quicksilver-de import { applyRealtimeOfferCorsHeaders, createResponseDeliveryWaiter, + readOfferBearerToken, + rejectOversizedOffer, respondRealtimeOffer, } from "./realtime-quicksilver-offer-http.js"; import { @@ -97,11 +99,6 @@ type OpenAIRealtimeOfferMetrics = { totalOfferMs: number; }; -function readBearerToken(req: IncomingMessage): string | undefined { - const authorization = req.headers.authorization?.trim(); - return authorization?.match(/^Bearer\s+([^\s]+)$/i)?.[1]; -} - export async function resolveOpenAIChatGptSubscriptionAuth(params: { cfg?: OpenClawConfig; agentDir?: string; @@ -368,7 +365,7 @@ export function createOpenAIQuicksilverBrowserSessionBroker(params: { return true; } prunePendingOffers(); - const token = readBearerToken(req); + const token = readOfferBearerToken(req); const offer = token ? pendingOffers.get(token) : undefined; if (!token || !offer || offer.expiresAt <= Date.now()) { respondRealtimeOffer(res, 401, "Invalid or expired realtime session token"); @@ -426,6 +423,8 @@ export function createOpenAIQuicksilverBrowserSessionBroker(params: { const sdp = await readRequestBodyWithLimit(req, { maxBytes: OPENAI_QUICKSILVER_MAX_SDP_BYTES, timeoutMs: 15_000, + // Defer destruction so the rejection below reaches the browser before the close. + destroyOnLimit: false, }); if (!sdp.trim()) { respondRealtimeOffer(res, 400, "SDP offer is required"); @@ -656,6 +655,9 @@ export function createOpenAIQuicksilverBrowserSessionBroker(params: { if (browserDisconnected) { return true; } + if (await rejectOversizedOffer(req, res, error)) { + return true; + } respondRealtimeOffer(res, 502, sessionError.message); return true; } finally { diff --git a/extensions/qa-lab/src/bus-server.test.ts b/extensions/qa-lab/src/bus-server.test.ts index e073ee37da5c..15993b268495 100644 --- a/extensions/qa-lab/src/bus-server.test.ts +++ b/extensions/qa-lab/src/bus-server.test.ts @@ -1,9 +1,9 @@ // Qa Lab tests cover bus server plugin behavior. import { Agent, createServer, request } from "node:http"; import { setTimeout as sleep } from "node:timers/promises"; -import { createMockIncomingRequest } from "openclaw/plugin-sdk/test-env"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import { closeQaHttpServer, handleQaBusRequest, startQaBusServer } from "./bus-server.js"; +import { closeQaHttpServer, startQaBusServer } from "./bus-server.js"; import { createQaBusState } from "./bus-state.js"; import type { QaBusPollResult } from "./runtime-api.js"; @@ -529,34 +529,25 @@ describe("handleQaBusRequest", () => { { pathname: "/v1/outbound/message", limit: 16 * 1024 * 1024 }, { pathname: "/v1/reset", limit: 1024 * 1024 }, ])( - "returns a controlled error when the $pathname body exceeds its limit", + "delivers 413 over the wire and closes when the $pathname body exceeds its limit", async ({ pathname, limit }) => { - const req = Object.assign(createMockIncomingRequest([]), { - method: "POST", - url: pathname, - headers: { "content-length": String(limit + 1) }, - }); - const res = { - statusCode: 0, - body: "", - writeHead(statusCode: number) { - this.statusCode = statusCode; - }, - end(payload: string) { - this.body = payload; - }, - }; + const bus = await startQaBusServer({ state: createQaBusState() }); + try { + // Declared over-cap length: rejected from the header alone, while the sender is + // still mid-upload and can only learn the outcome from what reaches it. + const result = await postRawWebhook({ + url: `${bus.baseUrl}${pathname}`, + body: "{}", + contentLength: limit + 1, + headers: { "content-type": "application/json" }, + }); - const handled = await handleQaBusRequest({ - req, - res: res as never, - state: createQaBusState(), - }); - - expect(handled).toBe(true); - expect(req.destroyed).toBe(true); - expect(res.statusCode).toBe(413); - expect(JSON.parse(res.body)).toEqual({ error: "Payload too large" }); + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(JSON.parse(result.body)).toEqual({ error: "Payload too large" }); + expect(result.closedByServer).toBe(true); + } finally { + await bus.stop(); + } }, ); }); diff --git a/extensions/qa-lab/src/bus-server.ts b/extensions/qa-lab/src/bus-server.ts index 84a53d398c0b..d2f03798761a 100644 --- a/extensions/qa-lab/src/bus-server.ts +++ b/extensions/qa-lab/src/bus-server.ts @@ -6,6 +6,7 @@ import { readRequestBodyWithLimit, requestBodyErrorToText, } from "openclaw/plugin-sdk/webhook-ingress"; +import { sendHttpRequestRejection } from "openclaw/plugin-sdk/webhook-request-guards"; import { z } from "zod"; import { normalizeAccountId, resolveQaBusPollStartCursor } from "./bus-queries.js"; import type { QaBusState } from "./bus-state.js"; @@ -187,6 +188,8 @@ export async function readQaJsonBody( await readRequestBodyWithLimit(req, { maxBytes: options?.maxBytes ?? QA_HTTP_JSON_MAX_BODY_BYTES, timeoutMs: QA_HTTP_JSON_BODY_TIMEOUT_MS, + // Defer destruction so writeQaRequestBodyLimitError can answer before the close. + destroyOnLimit: false, }) ).trim(); if (!text) { @@ -226,11 +229,21 @@ export function dispatchQaHttpRequest(res: ServerResponse, task: () => Promise { if (!isRequestBodyLimitError(error)) { return false; } - writeError(res, error.statusCode, requestBodyErrorToText(error.code)); + await sendHttpRequestRejection( + req, + res, + error.statusCode, + JSON.stringify({ error: requestBodyErrorToText(error.code) }), + "application/json; charset=utf-8", + ); return true; } @@ -453,7 +466,7 @@ export async function handleQaBusRequest(params: { return true; } } catch (error) { - if (writeQaRequestBodyLimitError(params.res, error)) { + if (await writeQaRequestBodyLimitError(params.req, params.res, error)) { return true; } writeError(params.res, 400, error); diff --git a/extensions/qa-lab/src/lab-server.ts b/extensions/qa-lab/src/lab-server.ts index 78a10dbd73d2..2b7699a5ced0 100644 --- a/extensions/qa-lab/src/lab-server.ts +++ b/extensions/qa-lab/src/lab-server.ts @@ -1,6 +1,6 @@ // Qa Lab plugin module implements lab server behavior. import fs from "node:fs"; -import { createServer } from "node:http"; +import { createServer, type IncomingMessage } from "node:http"; import path from "node:path"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime"; @@ -85,8 +85,12 @@ export type { QaLabServerStartParams, } from "./lab-server.types.js"; -function writeQaLabServerError(res: Parameters[0], error: unknown): void { - if (writeQaRequestBodyLimitError(res, error)) { +async function writeQaLabServerError( + req: IncomingMessage, + res: Parameters[0], + error: unknown, +): Promise { + if (await writeQaRequestBodyLimitError(req, res, error)) { return; } if (isQaMalformedJsonBodyError(error)) { @@ -928,7 +932,7 @@ export async function startQaLabServer( } res.end(body); } catch (error) { - writeQaLabServerError(res, error); + await writeQaLabServerError(req, res, error); } }); }); diff --git a/extensions/raft/src/gateway.test.ts b/extensions/raft/src/gateway.test.ts index 2ac21a08fc42..ccce8ed2689f 100644 --- a/extensions/raft/src/gateway.test.ts +++ b/extensions/raft/src/gateway.test.ts @@ -8,6 +8,7 @@ import { tempWorkspaceSync, type TempWorkspaceSync, } from "openclaw/plugin-sdk/temp-path"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { withTimeout } from "openclaw/plugin-sdk/text-utility-runtime"; import { afterEach, describe, expect, it, vi } from "vitest"; import type { ResolvedRaftAccount } from "./accounts.js"; @@ -186,6 +187,46 @@ describe("Raft wake gateway", () => { } }); + // Raft already answered this case through its own close-after-response teardown; the + // wire behavior must survive replacing that teardown with the shared transport owner. + it("keeps delivering 413 for an over-limit wake payload and closing the connection", async () => { + const { ctx, controller, wakeDedupe } = createContext(); + Object.defineProperty(ctx, "abortSignal", { value: controller.signal }); + const bridge = new FakeBridge(); + const start = startRaftGatewayAccount(ctx, { + spawnBridge: bridge.spawn, + wakeDedupe, + }); + void start.catch(bridge.started.reject); + + try { + const { endpoint: wakeEndpoint, token: bridgeToken } = await withTimeout( + bridge.started.promise, + 500, + "Raft bridge startup", + ); + + // Declared and sent in one write: the shape whose rejection used to race the flush. + const result = await postRawWebhook({ + url: wakeEndpoint, + body: JSON.stringify({ deliveryId: "x".repeat(16 * 1024) }), + headers: { + "content-type": "application/json", + "x-raft-bridge-token": bridgeToken, + }, + }); + + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(JSON.parse(result.body)).toEqual({ + error: "Wake payload exceeds the 16 KiB limit.", + }); + expect(result.closedByServer).toBe(true); + } finally { + controller.abort(); + await start; + } + }); + it("accepts authenticated content-free wake hints and dedupes retry delivery ids", async () => { const { ctx, controller, run, wakeDedupe } = createContext(); Object.defineProperty(ctx, "abortSignal", { value: controller.signal }); diff --git a/extensions/raft/src/gateway.ts b/extensions/raft/src/gateway.ts index e418922f835c..b28485215dad 100644 --- a/extensions/raft/src/gateway.ts +++ b/extensions/raft/src/gateway.ts @@ -14,6 +14,7 @@ import { killProcessTree } from "openclaw/plugin-sdk/process-runtime"; import { safeEqualSecret } from "openclaw/plugin-sdk/security-runtime"; import { readJsonBodyWithLimit, + sendHttpRequestRejection, WEBHOOK_BODY_READ_DEFAULTS, } from "openclaw/plugin-sdk/webhook-request-guards"; import { RAFT_CHANNEL_ID, type ResolvedRaftAccount } from "./accounts.js"; @@ -82,6 +83,7 @@ class WakeRequestError extends Error { constructor( readonly statusCode: number, message: string, + /** Body-limit rejections own the connection, so their answer goes through the transport. */ readonly closeAfterResponse = false, ) { super(message); @@ -337,23 +339,26 @@ export async function startRaftGatewayAccount( runtimeSession, ...(dispatched ? {} : { duplicate: true }), }); - })().catch((error: unknown) => { + })().catch(async (error: unknown) => { const statusCode = error instanceof WakeRequestError ? error.statusCode : 500; const message = error instanceof WakeRequestError ? error.message : "Internal server error."; ctx.log?.warn?.(`Raft wake request rejected: ${message}`); - if (!response.headersSent) { - if (error instanceof WakeRequestError && error.closeAfterResponse) { - response.setHeader("Connection", "close"); - response.once("close", () => { - if (!request.destroyed) { - request.destroy(); - } - }); - } - sendJson(response, statusCode, { error: message }); - } else { + if (response.headersSent) { response.destroy(); + return; } + if (error instanceof WakeRequestError && error.closeAfterResponse) { + response.setHeader("cache-control", "no-store"); + await sendHttpRequestRejection( + request, + response, + statusCode, + JSON.stringify({ error: message }), + "application/json; charset=utf-8", + ); + return; + } + sendJson(response, statusCode, { error: message }); }); }); server.on("connection", (socket) => { diff --git a/extensions/sms/src/twilio.ts b/extensions/sms/src/twilio.ts index d8c3fe21c613..b919fdc15239 100644 --- a/extensions/sms/src/twilio.ts +++ b/extensions/sms/src/twilio.ts @@ -309,13 +309,17 @@ export async function readTwilioWebhookForm(req: IncomingMessage): Promise"); } diff --git a/extensions/sms/src/webhook.test-support.ts b/extensions/sms/src/webhook.test-support.ts new file mode 100644 index 000000000000..52ce58e607ce --- /dev/null +++ b/extensions/sms/src/webhook.test-support.ts @@ -0,0 +1,52 @@ +// Sms test support shares webhook fixtures between the unit and raw-wire suites. +import { vi } from "vitest"; +import type { SmsDeliveryRecorder } from "./delivery-observations.js"; +import type { ResolvedSmsAccount } from "./types.js"; + +let testAccountSequence = 0; +let activeAccountId = "test-0"; + +// Each test gets its own account id so the handler's per-account rate-limiter +// buckets never carry a previous test's counters into the next one. +export function advanceSmsTestAccountId(): string { + activeAccountId = `test-${++testAccountSequence}`; + return activeAccountId; +} + +export function createSmsTestAccount( + overrides: Partial = {}, +): ResolvedSmsAccount { + return { + accountId: activeAccountId, + enabled: true, + accountSid: "AC123", + authToken: "secret", + fromNumber: "+15557654321", + messagingServiceSid: "", + defaultTo: "", + webhookPath: "/webhooks/sms", + publicWebhookUrl: "https://gateway.example.com/webhooks/sms", + dangerouslyDisableSignatureValidation: false, + dmPolicy: "pairing", + allowFrom: [], + textChunkLimit: 1500, + ...overrides, + }; +} + +export function createSmsTestDeliveryRecorder( + record = vi.fn(async ({ account, form }) => ({ + duplicate: false, + record: { + accountId: account.accountId, + accountSidHash: "account-sid-hash", + messageSid: form.MessageSid ?? form.SmsSid ?? form.SmsMessageSid ?? "", + status: form.MessageStatus ?? form.SmsStatus ?? "", + firstObservedAt: 1, + lastObservedAt: 1, + observations: [], + }, + })), +): SmsDeliveryRecorder & { record: typeof record } { + return { record }; +} diff --git a/extensions/sms/src/webhook.test.ts b/extensions/sms/src/webhook.test.ts index b0461d9cf33d..4b5993200a20 100644 --- a/extensions/sms/src/webhook.test.ts +++ b/extensions/sms/src/webhook.test.ts @@ -1,11 +1,17 @@ // Sms tests cover webhook plugin behavior. import { createHmac } from "node:crypto"; import type { IncomingMessage, ServerResponse } from "node:http"; +import { Socket } from "node:net"; import { Readable } from "node:stream"; import { beforeEach, describe, expect, it, vi } from "vitest"; import type { SmsDeliveryRecorder } from "./delivery-observations.js"; import type { ResolvedSmsAccount } from "./types.js"; import { createSmsWebhookHandler } from "./webhook.js"; +import { + advanceSmsTestAccountId, + createSmsTestAccount, + createSmsTestDeliveryRecorder, +} from "./webhook.test-support.js"; const assertSmsCredentialOwnerAvailable = vi.hoisted(() => vi.fn()); const enqueueSmsIngress = vi.hoisted(() => @@ -14,7 +20,6 @@ const enqueueSmsIngress = vi.hoisted(() => vi.mock("./credential-availability.js", () => ({ assertSmsCredentialOwnerAvailable })); -let testAccountSequence = 0; let activeAccountId = "test-0"; function createIngress() { @@ -41,31 +46,12 @@ function computeTestTwilioSignature(params: { return createHmac("sha1", params.authToken).update(data).digest("base64"); } -function createAccount(overrides: Partial = {}): ResolvedSmsAccount { - return { - accountId: activeAccountId, - enabled: true, - accountSid: "AC123", - authToken: "secret", - fromNumber: "+15557654321", - messagingServiceSid: "", - defaultTo: "", - webhookPath: "/webhooks/sms", - publicWebhookUrl: "https://gateway.example.com/webhooks/sms", - dangerouslyDisableSignatureValidation: false, - dmPolicy: "pairing", - allowFrom: [], - textChunkLimit: 1500, - ...overrides, - }; -} - function createSignedBody(params?: { account?: ResolvedSmsAccount; body?: string; messageSid?: string; }): { body: string; signature: string } { - const account = params?.account ?? createAccount(); + const account = params?.account ?? createSmsTestAccount(); const body = params?.body ?? `AccountSid=${encodeURIComponent(account.accountSid)}&From=%2B15551234567&To=%2B15557654321&Body=hello&MessageSid=${encodeURIComponent(params?.messageSid ?? "SM123")}`; @@ -79,6 +65,12 @@ function createSignedBody(params?: { }; } +function createLoopbackSocket(remoteAddress = "127.0.0.1"): Socket { + const socket = new Socket(); + Object.defineProperty(socket, "remoteAddress", { value: remoteAddress }); + return socket; +} + function createRequest( body: string, signature: string, @@ -92,7 +84,7 @@ function createRequest( ...options?.headers, }; Object.defineProperty(req, "socket", { - value: { remoteAddress: options?.remoteAddress ?? "127.0.0.1" }, + value: createLoopbackSocket(options?.remoteAddress), }); return req; } @@ -100,9 +92,9 @@ function createRequest( function configureRequest(req: IncomingMessage): IncomingMessage { req.method = "POST"; req.headers = {}; - Object.defineProperty(req, "socket", { - value: { remoteAddress: "127.0.0.1" }, - }); + // A real socket: the reader hands a limited request to the connection-level + // rejection owner, which subscribes to socket events before the caller answers. + Object.defineProperty(req, "socket", { value: createLoopbackSocket() }); return req; } @@ -168,7 +160,7 @@ function createSignedDeliveryPayload(params: { account?: ResolvedSmsAccount; accountSid?: string; }): { body: string; signature: string; form: Record } { - const account = params.account ?? createAccount(); + const account = params.account ?? createSmsTestAccount(); const form = { AccountSid: params.accountSid ?? account.accountSid, From: account.fromNumber, @@ -188,23 +180,6 @@ function createSignedDeliveryPayload(params: { }; } -function createDeliveryRecorder( - record = vi.fn(async ({ account, form }) => ({ - duplicate: false, - record: { - accountId: account.accountId, - accountSidHash: "account-sid-hash", - messageSid: form.MessageSid ?? form.SmsSid ?? form.SmsMessageSid ?? "", - status: form.MessageStatus ?? form.SmsStatus ?? "", - firstObservedAt: 1, - lastObservedAt: 1, - observations: [], - }, - })), -): SmsDeliveryRecorder & { record: typeof record } { - return { record }; -} - function createMessageSid(index: number): string { return `SM${index.toString(16).padStart(32, "0")}`; } @@ -214,7 +189,7 @@ describe("createSmsWebhookHandler", () => { assertSmsCredentialOwnerAvailable.mockReset(); enqueueSmsIngress.mockReset(); enqueueSmsIngress.mockResolvedValue({ kind: "accepted", duplicate: false }); - activeAccountId = `test-${++testAccountSequence}`; + activeAccountId = advanceSmsTestAccountId(); }); it("rechecks the owner after parsing and before authentication or durable admission", async () => { @@ -226,7 +201,7 @@ describe("createSmsWebhookHandler", () => { const { body, signature } = createSignedBody(); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -249,7 +224,7 @@ describe("createSmsWebhookHandler", () => { const { body, signature } = createSignedSmsPayload(createMessageSid(1)); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount({ + account: createSmsTestAccount({ publicWebhookUrl: "https://gateway.example.com/webhooks/sms#rp=4xx", }), ingress: createIngress(), @@ -263,57 +238,10 @@ describe("createSmsWebhookHandler", () => { expect(enqueueSmsIngress).toHaveBeenCalledWith(parseTestTwilioForm(body)); }); - it("returns terminal HTTP 413 for an oversized callback body", async () => { - const delivery = createDeliveryRecorder(); - const handler = createSmsWebhookHandler({ - cfg: {}, - account: createAccount(), - ingress: createIngress(), - delivery, - }); - const res = createResponse(); - - await handler( - createRequest("x", "unused", { - headers: { "content-length": String(32 * 1024 + 1) }, - }), - res, - ); - - expect(res.statusCode).toBe(413); - expect(res.body).toBe("Payload too large"); - expect(delivery.record).not.toHaveBeenCalled(); - expect(enqueueSmsIngress).not.toHaveBeenCalled(); - }); - - it("rethrows request body timeouts for Gateway-owned retry responses", async () => { - vi.useFakeTimers(); - try { - const handler = createSmsWebhookHandler({ - cfg: {}, - account: createAccount(), - ingress: createIngress(), - }); - const res = createResponse(); - const handling = handler(createPendingRequest(), res); - const expected = expect(handling).rejects.toMatchObject({ - code: "REQUEST_BODY_TIMEOUT", - }); - - await vi.advanceTimersByTimeAsync(5_000); - await expected; - - expect(res.endMock).not.toHaveBeenCalled(); - expect(enqueueSmsIngress).not.toHaveBeenCalled(); - } finally { - vi.useRealTimers(); - } - }); - it("rethrows unexpected request read failures for Gateway-owned retry responses", async () => { const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -328,7 +256,7 @@ describe("createSmsWebhookHandler", () => { it("rethrows a closed request body for Gateway-owned retry responses", async () => { const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -345,10 +273,10 @@ describe("createSmsWebhookHandler", () => { messageSid: createMessageSid(20), status: "delivered", }); - const delivery = createDeliveryRecorder(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), delivery, }); @@ -366,7 +294,7 @@ describe("createSmsWebhookHandler", () => { }); it("accepts legacy SmsSid and SmsStatus delivery callbacks", async () => { - const account = createAccount(); + const account = createSmsTestAccount(); const form = { AccountSid: account.accountSid, From: account.fromNumber, @@ -380,7 +308,7 @@ describe("createSmsWebhookHandler", () => { authToken: account.authToken, form, }); - const delivery = createDeliveryRecorder(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, account, @@ -399,7 +327,7 @@ describe("createSmsWebhookHandler", () => { it.each(["receiving", "received"])( "keeps legacy inbound SmsStatus=%s on the durable ingress path", async (status) => { - const account = createAccount(); + const account = createSmsTestAccount(); const form = { AccountSid: account.accountSid, From: "+15551234567", @@ -414,7 +342,7 @@ describe("createSmsWebhookHandler", () => { authToken: account.authToken, form, }); - const delivery = createDeliveryRecorder(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, account, @@ -436,10 +364,10 @@ describe("createSmsWebhookHandler", () => { messageSid: createMessageSid(26), status: "delivered", }); - const delivery = createDeliveryRecorder(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), delivery, }); @@ -457,14 +385,14 @@ describe("createSmsWebhookHandler", () => { messageSid: createMessageSid(21), status: "sent", }); - const delivery = createDeliveryRecorder( + const delivery = createSmsTestDeliveryRecorder( vi.fn(async () => { throw new Error("sqlite unavailable"); }), ); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), delivery, }); @@ -479,14 +407,14 @@ describe("createSmsWebhookHandler", () => { }); it("waits for the durable delivery commit before returning HTTP 200", async () => { - const account = createAccount(); + const account = createSmsTestAccount(); const payload = createSignedDeliveryPayload({ account, messageSid: createMessageSid(22), status: "sent", }); let releaseCommit: (() => void) | undefined; - const delivery = createDeliveryRecorder( + const delivery = createSmsTestDeliveryRecorder( vi.fn(async ({ form }) => { await new Promise((resolve) => { releaseCommit = resolve; @@ -535,10 +463,10 @@ describe("createSmsWebhookHandler", () => { status: "failed", accountSid: "AC-other", }); - const delivery = createDeliveryRecorder(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), delivery, }); @@ -558,7 +486,7 @@ describe("createSmsWebhookHandler", () => { ["whitespace", " "], ["padded", " AC123 "], ])("acknowledges but does not store a delivery callback with %s AccountSid", async (_, sid) => { - const account = createAccount(); + const account = createSmsTestAccount(); const form: Record = { MessageSid: createMessageSid(27), MessageStatus: "failed", @@ -572,7 +500,7 @@ describe("createSmsWebhookHandler", () => { authToken: account.authToken, form, }); - const delivery = createDeliveryRecorder(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, account, @@ -594,7 +522,7 @@ describe("createSmsWebhookHandler", () => { enqueueSmsIngress.mockRejectedValueOnce(new Error("sqlite unavailable")); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -618,7 +546,7 @@ describe("createSmsWebhookHandler", () => { ); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -642,7 +570,7 @@ describe("createSmsWebhookHandler", () => { enqueueSmsIngress.mockResolvedValueOnce({ kind: "accepted", duplicate: true }); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -662,7 +590,7 @@ describe("createSmsWebhookHandler", () => { }); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -683,7 +611,7 @@ describe("createSmsWebhookHandler", () => { }); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); const res = createResponse(); @@ -704,7 +632,7 @@ describe("createSmsWebhookHandler", () => { }); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); @@ -734,7 +662,7 @@ describe("createSmsWebhookHandler", () => { }); const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); @@ -748,7 +676,7 @@ describe("createSmsWebhookHandler", () => { }); it("does not let unsigned proxy traffic consume the same client's signed webhook rate limit", async () => { - const account = createAccount(); + const account = createSmsTestAccount(); const handler = createSmsWebhookHandler({ cfg: { gateway: { trustedProxies: ["127.0.0.1"] } }, account, @@ -789,13 +717,13 @@ describe("createSmsWebhookHandler", () => { }); it("scopes signed webhook rate limits to one SMS account and route", async () => { - const supportAccount = createAccount({ + const supportAccount = createSmsTestAccount({ accountId: "support", accountSid: "AC-support", webhookPath: "/webhooks/sms/support", publicWebhookUrl: "https://gateway.example.com/webhooks/sms/support", }); - const defaultAccount = createAccount(); + const defaultAccount = createSmsTestAccount(); const supportHandler = createSmsWebhookHandler({ cfg: {}, account: supportAccount, @@ -836,8 +764,8 @@ describe("createSmsWebhookHandler", () => { it("meters inbound dispatch per sender without throttling signed delivery callbacks", async () => { const warn = vi.fn(); - const account = createAccount(); - const delivery = createDeliveryRecorder(); + const account = createSmsTestAccount(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, account, @@ -900,8 +828,8 @@ describe("createSmsWebhookHandler", () => { }); it("bounds aggregate inbound fan-out without throttling signed delivery callbacks", async () => { - const account = createAccount(); - const delivery = createDeliveryRecorder(); + const account = createSmsTestAccount(); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: {}, account, @@ -946,7 +874,7 @@ describe("createSmsWebhookHandler", () => { try { const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); @@ -974,7 +902,7 @@ describe("createSmsWebhookHandler", () => { it("shares one quota for invalid signed senders without throttling a valid sender", async () => { const handler = createSmsWebhookHandler({ cfg: {}, - account: createAccount(), + account: createSmsTestAccount(), ingress: createIngress(), }); @@ -1005,7 +933,7 @@ describe("createSmsWebhookHandler", () => { }); it("keeps validation-disabled webhook dispatches on the stricter callback budget", async () => { - const account = createAccount({ dangerouslyDisableSignatureValidation: true }); + const account = createSmsTestAccount({ dangerouslyDisableSignatureValidation: true }); const handler = createSmsWebhookHandler({ cfg: { gateway: { trustedProxies: ["127.0.0.1"] } }, account, @@ -1042,8 +970,8 @@ describe("createSmsWebhookHandler", () => { }); it("rate limits unsigned delivery callbacks by client address before persistence", async () => { - const account = createAccount({ dangerouslyDisableSignatureValidation: true }); - const delivery = createDeliveryRecorder(); + const account = createSmsTestAccount({ dangerouslyDisableSignatureValidation: true }); + const delivery = createSmsTestDeliveryRecorder(); const handler = createSmsWebhookHandler({ cfg: { gateway: { trustedProxies: ["127.0.0.1"] } }, account, diff --git a/extensions/sms/src/webhook.ts b/extensions/sms/src/webhook.ts index 646b6dd86e63..8b514b82517c 100644 --- a/extensions/sms/src/webhook.ts +++ b/extensions/sms/src/webhook.ts @@ -6,6 +6,7 @@ import { isRequestBodyLimitError, resolveRequestClientIp, } from "openclaw/plugin-sdk/webhook-ingress"; +import { sendHttpRequestRejection } from "openclaw/plugin-sdk/webhook-request-guards"; import { assertSmsCredentialOwnerAvailable } from "./credential-availability.js"; import { createSmsDeliveryRecorder, @@ -18,6 +19,7 @@ import { resolveTwilioInboundSender, resolveTwilioMessageSid, resolveTwilioWebhookSignatureUrl, + TWIML_CONTENT_TYPE, verifyTwilioSignature, } from "./twilio.js"; import type { ResolvedSmsAccount } from "./types.js"; @@ -127,7 +129,19 @@ export function createSmsWebhookHandler(params: SmsWebhookHandlerParams) { form = await readTwilioWebhookForm(req); } catch (error) { if (isRequestBodyLimitError(error, "PAYLOAD_TOO_LARGE")) { - respondTwiml(res, 413, "Payload too large"); + await sendHttpRequestRejection(req, res, 413, "Payload too large", TWIML_CONTENT_TYPE); + return true; + } + if (isRequestBodyLimitError(error, "REQUEST_BODY_TIMEOUT")) { + // Twilio retries 5xx responses. Keep that outcome while the shared owner closes in order. + res.setHeader("cache-control", "no-store"); + await sendHttpRequestRejection( + req, + res, + 500, + "Internal Server Error", + "text/plain; charset=utf-8", + ); return true; } throw error; diff --git a/extensions/sms/src/webhook.wire.test.ts b/extensions/sms/src/webhook.wire.test.ts new file mode 100644 index 000000000000..19996fa9c1e2 --- /dev/null +++ b/extensions/sms/src/webhook.wire.test.ts @@ -0,0 +1,146 @@ +// Sms tests cover webhook responses as the sender receives them on the wire. +import { createServer } from "node:http"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { createSmsWebhookHandler } from "./webhook.js"; +import { + advanceSmsTestAccountId, + createSmsTestAccount, + createSmsTestDeliveryRecorder, +} from "./webhook.test-support.js"; + +vi.mock("node:timers", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + setTimeout: ((callback: (...args: unknown[]) => void, delay?: number, ...args: unknown[]) => + globalThis.setTimeout(callback, delay, ...args)) as typeof actual.setTimeout, + clearTimeout: ((timer: ReturnType | undefined) => + globalThis.clearTimeout(timer)) as typeof actual.clearTimeout, + }; +}); + +const assertSmsCredentialOwnerAvailable = vi.hoisted(() => vi.fn()); +const enqueueSmsIngress = vi.hoisted(() => + vi.fn(async () => ({ kind: "accepted" as const, duplicate: false })), +); + +vi.mock("./credential-availability.js", () => ({ assertSmsCredentialOwnerAvailable })); + +describe("createSmsWebhookHandler over a real connection", () => { + beforeEach(() => { + assertSmsCredentialOwnerAvailable.mockReset(); + enqueueSmsIngress.mockReset(); + enqueueSmsIngress.mockResolvedValue({ kind: "accepted", duplicate: false }); + advanceSmsTestAccountId(); + }); + + it("delivers HTTP 413 over the wire and closes for an oversized callback body", async () => { + const delivery = createSmsTestDeliveryRecorder(); + const handler = createSmsWebhookHandler({ + cfg: {}, + account: createSmsTestAccount(), + ingress: { enqueue: enqueueSmsIngress }, + delivery, + }); + const server = createServer((req, res) => { + void handler(req, res); + }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected the SMS webhook test server to have a TCP address"); + } + + // Declared and sent in one write: the shape whose rejection used to race the flush. + const result = await postRawWebhook({ + url: `http://127.0.0.1:${address.port}/sms`, + body: "x".repeat(32 * 1024 + 1), + headers: { + "content-type": "application/x-www-form-urlencoded", + "x-twilio-signature": "unused", + }, + }); + + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(result.headers.connection).toBe("close"); + expect(result.body).toBe("Payload too large"); + expect(result.closedByServer).toBe(true); + expect(delivery.record).not.toHaveBeenCalled(); + expect(enqueueSmsIngress).not.toHaveBeenCalled(); + } finally { + server.closeAllConnections(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } + }); + + it("delivers a retryable 500 before closing a timed-out callback upload", async () => { + vi.useFakeTimers(); + const handler = createSmsWebhookHandler({ + cfg: {}, + account: createSmsTestAccount(), + ingress: { enqueue: enqueueSmsIngress }, + }); + let routeError: unknown; + const server = createServer((req, res) => { + void handler(req, res).catch((error: unknown) => { + routeError = error; + res.statusCode = 500; + res.setHeader("content-type", "text/plain; charset=utf-8"); + res.end("Internal Server Error"); + }); + }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected the SMS webhook test server to have a TCP address"); + } + const requestReceived = new Promise((resolve) => { + server.once("connection", (socket) => socket.once("data", () => resolve())); + }); + const resultPromise = postRawWebhook({ + url: `http://127.0.0.1:${address.port}/sms`, + body: "x", + contentLength: 2, + idleTimeoutMs: 10_000, + headers: { + "content-type": "application/x-www-form-urlencoded", + "x-twilio-signature": "unused", + }, + }); + + await requestReceived; + await vi.advanceTimersByTimeAsync(6_000); + const result = await resultPromise; + + expect(routeError).toBeUndefined(); + expect(result.statusLine).toBe("HTTP/1.1 500 Internal Server Error"); + expect(result.headers.connection).toBe("close"); + expect(result.body).toBe("Internal Server Error"); + expect(result.closedByServer).toBe(true); + expect(enqueueSmsIngress).not.toHaveBeenCalled(); + } finally { + vi.useRealTimers(); + server.closeAllConnections(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } + }); +}); diff --git a/extensions/synology-chat/src/webhook-handler.test.ts b/extensions/synology-chat/src/webhook-handler.test.ts index caf6ba9bf926..5ac09cadc17e 100644 --- a/extensions/synology-chat/src/webhook-handler.test.ts +++ b/extensions/synology-chat/src/webhook-handler.test.ts @@ -1,5 +1,7 @@ // Synology Chat tests cover webhook handler plugin behavior. +import { createServer } from "node:http"; import { expectDefined } from "@openclaw/normalization-core"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { describe, it, expect, vi, beforeEach } from "vitest"; import { makeFormBody, makeReq, makeRes, makeStalledReq } from "./test-http-utils.js"; import type { ResolvedSynologyChatAccount } from "./types.js"; @@ -304,20 +306,66 @@ describe("createWebhookHandler", () => { expect(res.status).toBe(400); }); - it("returns 408 when request body times out", async () => { + it.each([ + { + name: "413 when the upload exceeds the pre-auth body limit", + bodyTimeoutMs: 5_000, + // Declared and sent in one write: the shape whose rejection used to race the flush. + body: "x".repeat(64 * 1024 + 1), + contentLength: undefined, + statusLine: "HTTP/1.1 413 Payload Too Large", + responseBody: JSON.stringify({ error: "Payload too large" }), + }, + { + name: "408 when the sender stalls mid-upload", + bodyTimeoutMs: 50, + // Promises more than is ever sent, so the read deadline fires with the request open. + body: "x".repeat(16), + contentLength: 64 * 1024, + statusLine: "HTTP/1.1 408 Request Timeout", + responseBody: JSON.stringify({ error: "Request body timeout" }), + }, + ])("delivers $name and then closes the connection", async (scenario) => { + const deliver = vi.fn(); const handler = createWebhookHandler({ account: makeAccount(), - deliver: vi.fn(), + deliver, log, - bodyTimeoutMs: 1, + bodyTimeoutMs: scenario.bodyTimeoutMs, }); + const server = createServer((req, res) => { + void handler(req, res); + }); + try { + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.removeListener("error", reject); + resolve(); + }); + }); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("expected the Synology webhook test server to have a TCP address"); + } - const req = makeStalledReq("POST"); - const res = makeRes(); - await handler(req, res); + const result = await postRawWebhook({ + url: `http://127.0.0.1:${address.port}/webhook/synology`, + body: scenario.body, + contentLength: scenario.contentLength, + headers: { "content-type": "application/x-www-form-urlencoded" }, + }); - expect(res.status).toBe(408); - expect(res.body).toContain("timeout"); + expect(result.statusLine).toBe(scenario.statusLine); + expect(result.body).toBe(scenario.responseBody); + expect(result.closedByServer).toBe(true); + expect(deliver).not.toHaveBeenCalled(); + } finally { + server.closeAllConnections(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } }); it("rejects excess concurrent pre-auth body reads from the same remote IP", async () => { diff --git a/extensions/synology-chat/src/webhook-handler.ts b/extensions/synology-chat/src/webhook-handler.ts index 9bd4858ab0fc..9f3f8012d02c 100644 --- a/extensions/synology-chat/src/webhook-handler.ts +++ b/extensions/synology-chat/src/webhook-handler.ts @@ -15,6 +15,7 @@ import { resolveRequestClientIp, requestBodyErrorToText, } from "openclaw/plugin-sdk/webhook-ingress"; +import { sendHttpRequestRejection } from "openclaw/plugin-sdk/webhook-request-guards"; import * as synologyClient from "./client.js"; import { validateToken, @@ -152,19 +153,23 @@ function getSynologyWebhookInFlightKey(account: ResolvedSynologyChatAccount): st /** Read the full request body as a string. */ async function readBody( req: IncomingMessage, - timeoutMs = PREAUTH_BODY_TIMEOUT_MS, + timeoutMs: number = PREAUTH_BODY_TIMEOUT_MS, ): Promise< | { ok: true; body: string } | { ok: false; statusCode: number; error: string; + /** Limit rejections own the connection, so their answer goes through the transport. */ + closeAfterResponse: boolean; } > { try { const body = await readRequestBodyWithLimit(req, { maxBytes: PREAUTH_MAX_BODY_BYTES, timeoutMs, + // Defer destruction so the caller can answer before the connection closes. + destroyOnLimit: false, }); return { ok: true, body }; } catch (err) { @@ -173,12 +178,14 @@ async function readBody( ok: false, statusCode: err.statusCode, error: requestBodyErrorToText(err.code), + closeAfterResponse: true, }; } return { ok: false, statusCode: 400, error: "Invalid request body", + closeAfterResponse: false, }; } } @@ -410,6 +417,16 @@ async function parseWebhookPayloadRequest(params: { const bodyResult = await readBody(params.req, params.bodyTimeoutMs); if (!bodyResult.ok) { params.log?.error("Failed to read request body", bodyResult.error); + if (bodyResult.closeAfterResponse) { + await sendHttpRequestRejection( + params.req, + params.res, + bodyResult.statusCode, + JSON.stringify({ error: bodyResult.error }), + "application/json", + ); + return { ok: false }; + } respondJson(params.res, bodyResult.statusCode, { error: bodyResult.error }); return { ok: false }; } diff --git a/extensions/telegram/src/webhook.test.ts b/extensions/telegram/src/webhook.test.ts index 42fdabf45497..926dff72687f 100644 --- a/extensions/telegram/src/webhook.test.ts +++ b/extensions/telegram/src/webhook.test.ts @@ -2949,12 +2949,11 @@ describe("startTelegramWebhook", () => { req.end("{}"); }); - if (responseOrError.kind === "response") { - expect(responseOrError.statusCode).toBe(413); - expect(responseOrError.body).toBe("Payload too large"); - } else { - expect(responseOrError.code).toBeOneOf(["ECONNRESET", "EPIPE"]); - } + expect(responseOrError).toEqual({ + kind: "response", + statusCode: 413, + body: "Payload too large", + }); expect(handleUpdateSpy).not.toHaveBeenCalled(); }, ); diff --git a/extensions/telegram/src/webhook.ts b/extensions/telegram/src/webhook.ts index 455ccfbbf0c5..d28e7fcbc6a6 100644 --- a/extensions/telegram/src/webhook.ts +++ b/extensions/telegram/src/webhook.ts @@ -29,7 +29,10 @@ import { createFixedWindowRateLimiter, WEBHOOK_RATE_LIMIT_DEFAULTS, } from "openclaw/plugin-sdk/webhook-ingress"; -import { readJsonBodyWithLimit } from "openclaw/plugin-sdk/webhook-request-guards"; +import { + readJsonBodyWithLimit, + sendHttpRequestRejection, +} from "openclaw/plugin-sdk/webhook-request-guards"; import { mergeTelegramAccountConfig } from "./account-config.js"; import { resolveTelegramAllowedUpdates } from "./allowed-updates.js"; import { withTelegramApiErrorLogging } from "./api-logging.js"; @@ -42,6 +45,7 @@ import { createTelegramWebhookStatusPublisher } from "./webhook-status.js"; const TELEGRAM_WEBHOOK_MAX_BODY_BYTES = 1024 * 1024; const TELEGRAM_WEBHOOK_BODY_TIMEOUT_MS = 30_000; +const TELEGRAM_WEBHOOK_TEXT_TYPE = "text/plain; charset=utf-8"; const TELEGRAM_WEBHOOK_ACCEPTED_HEADER = "x-openclaw-delivery-accepted"; const TELEGRAM_WEBHOOK_ACCEPTED_VALUE = "durable"; const TELEGRAM_WEBHOOK_SPOOLED_DRAIN_INTERVAL_MS = 500; @@ -472,7 +476,7 @@ export async function startTelegramWebhook(opts: { if (res.headersSent || res.writableEnded) { return; } - res.writeHead(statusCode, { "Content-Type": "text/plain; charset=utf-8" }); + res.writeHead(statusCode, { "Content-Type": TELEGRAM_WEBHOOK_TEXT_TYPE }); res.end(text); }; @@ -514,14 +518,16 @@ export async function startTelegramWebhook(opts: { maxBytes: TELEGRAM_WEBHOOK_MAX_BODY_BYTES, timeoutMs: TELEGRAM_WEBHOOK_BODY_TIMEOUT_MS, emptyObjectOnEmpty: false, + // Defer destruction so the rejections below reach Telegram before the close. + destroyOnLimit: false, }); if (!body.ok) { if (body.code === "PAYLOAD_TOO_LARGE") { - respondText(413, body.error); + await sendHttpRequestRejection(req, res, 413, body.error, TELEGRAM_WEBHOOK_TEXT_TYPE); return; } if (body.code === "REQUEST_BODY_TIMEOUT") { - respondText(408, body.error); + await sendHttpRequestRejection(req, res, 408, body.error, TELEGRAM_WEBHOOK_TEXT_TYPE); return; } if (body.code === "CONNECTION_CLOSED") { diff --git a/extensions/voice-call/api.ts b/extensions/voice-call/api.ts index ed30e9c130c5..4e5d09770916 100644 --- a/extensions/voice-call/api.ts +++ b/extensions/voice-call/api.ts @@ -9,6 +9,7 @@ export { type OpenClawPluginApi, readRequestBodyWithLimit, requestBodyErrorToText, + sendHttpRequestRejection, type SessionEntry, sleep, TtsAutoSchema, diff --git a/extensions/voice-call/runtime-api.ts b/extensions/voice-call/runtime-api.ts index fdf6e27aad1f..6a9065a4590d 100644 --- a/extensions/voice-call/runtime-api.ts +++ b/extensions/voice-call/runtime-api.ts @@ -8,6 +8,7 @@ export { isRequestBodyLimitError, readRequestBodyWithLimit, requestBodyErrorToText, + sendHttpRequestRejection, } from "openclaw/plugin-sdk/webhook-request-guards"; export { fetchWithSsrFGuard, isBlockedHostnameOrIp } from "openclaw/plugin-sdk/ssrf-runtime"; export type { SessionEntry } from "openclaw/plugin-sdk/session-store-runtime"; diff --git a/extensions/voice-call/src/webhook.hangup-once.lifecycle.test.ts b/extensions/voice-call/src/webhook.hangup-once.lifecycle.test.ts index ddbf518a8b84..fe538372ced4 100644 --- a/extensions/voice-call/src/webhook.hangup-once.lifecycle.test.ts +++ b/extensions/voice-call/src/webhook.hangup-once.lifecycle.test.ts @@ -6,6 +6,7 @@ import { createPluginStateSyncKeyedStoreForTests, resetPluginStateStoreForTests, } from "openclaw/plugin-sdk/plugin-state-test-runtime"; +import { postRawWebhook } from "openclaw/plugin-sdk/test-env"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { VoiceCallConfigSchema, type VoiceCallConfig } from "./config.js"; import { CallManager } from "./manager.js"; @@ -327,3 +328,42 @@ describe("Voice-call webhook hangup-once lifecycle", () => { expect(secondManager.getCallByProviderCallId("provider-inbound-1")).toBeUndefined(); }); }); + +describe("Voice-call webhook body limits", () => { + beforeEach(() => { + resetPluginStateStoreForTests(); + installStateRuntime(); + }); + + afterEach(() => { + resetPluginStateStoreForTests(); + }); + + it("answers an over-limit webhook with 413 and then closes the connection", async () => { + // Driven over a real socket: the server answers while the sender is still uploading + // and then closes, so a mocked response cannot show whether either half happened. + const provider = new FakeProvider(); + const config = createConfig(); + const manager = new CallManager(config, createTestStorePath()); + await manager.initialize(provider, "https://example.com/voice/webhook"); + const server = new VoiceCallWebhookServer(config, manager, provider); + + try { + const baseUrl = await server.start(); + const result = await postRawWebhook({ + url: baseUrl, + body: `CallSid=CA123&From=%2B15552222222&Padding=${"x".repeat(2 * 1024 * 1024)}`, + headers: { + "content-type": "application/x-www-form-urlencoded", + "x-plivo-signature-v2": "sig", + "x-plivo-signature-v2-nonce": "nonce", + }, + }); + + expect(result.statusLine).toBe("HTTP/1.1 413 Payload Too Large"); + expect(result.closedByServer).toBe(true); + } finally { + await server.stop(); + } + }); +}); diff --git a/extensions/voice-call/src/webhook.ts b/extensions/voice-call/src/webhook.ts index 60f07deac736..7fad4b2c1101 100644 --- a/extensions/voice-call/src/webhook.ts +++ b/extensions/voice-call/src/webhook.ts @@ -22,6 +22,7 @@ import { isRequestBodyLimitError, readRequestBodyWithLimit, requestBodyErrorToText, + sendHttpRequestRejection, } from "../api.js"; import type { OpenClawPluginApi } from "../api.js"; import { isAllowlistedCaller, normalizePhoneNumber } from "./allowlist.js"; @@ -672,14 +673,18 @@ export class VoiceCallWebhookServer { res: http.ServerResponse, webhookPath: string, ): Promise { - const payload = await this.runWebhookPipeline(req, webhookPath); - this.writeWebhookResponse(res, payload); + const payload = await this.runWebhookPipeline(req, webhookPath, res); + // A body-limit rejection already wrote its answer through the transport owner. + if (payload) { + this.writeWebhookResponse(res, payload); + } } private async runWebhookPipeline( req: http.IncomingMessage, webhookPath: string, - ): Promise { + res: http.ServerResponse, + ): Promise { const url = buildRequestUrl(req.url); if (url.pathname === "/voice/hold-music") { @@ -729,10 +734,17 @@ export class VoiceCallWebhookServer { body = await this.readBody(req, MAX_WEBHOOK_BODY_BYTES, WEBHOOK_BODY_TIMEOUT_MS); } catch (err) { if (isRequestBodyLimitError(err, "PAYLOAD_TOO_LARGE")) { - return { statusCode: 413, body: "Payload Too Large" }; + await sendHttpRequestRejection(req, res, 413, "Payload Too Large"); + return null; } if (isRequestBodyLimitError(err, "REQUEST_BODY_TIMEOUT")) { - return { statusCode: 408, body: requestBodyErrorToText("REQUEST_BODY_TIMEOUT") }; + await sendHttpRequestRejection( + req, + res, + 408, + requestBodyErrorToText("REQUEST_BODY_TIMEOUT"), + ); + return null; } throw err; } @@ -1047,9 +1059,10 @@ export class VoiceCallWebhookServer { private readBody( req: http.IncomingMessage, maxBytes: number, - timeoutMs = WEBHOOK_BODY_TIMEOUT_MS, + timeoutMs: number, ): Promise { - return readRequestBodyWithLimit(req, { maxBytes, timeoutMs }); + // Defer destruction so a limit rejection can be answered before the close. + return readRequestBodyWithLimit(req, { maxBytes, timeoutMs, destroyOnLimit: false }); } /** diff --git a/src/plugin-sdk/test-env.ts b/src/plugin-sdk/test-env.ts index d009bffa2172..bea0bdb5bae9 100644 --- a/src/plugin-sdk/test-env.ts +++ b/src/plugin-sdk/test-env.ts @@ -14,4 +14,5 @@ export { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; export { useFrozenTime, useRealTime } from "../test-utils/frozen-time.js"; export { withServer } from "./test-helpers/http-test-server.js"; export { createMockIncomingRequest } from "./test-helpers/mock-incoming-request.js"; +export { postRawWebhook, type RawHttpResult } from "./test-helpers/raw-http-request.js"; export { withTempHomeCore as withTempHome } from "./test-helpers/temp-home.js"; diff --git a/src/plugin-sdk/test-helpers/raw-http-request.ts b/src/plugin-sdk/test-helpers/raw-http-request.ts new file mode 100644 index 000000000000..db844c1f6e3f --- /dev/null +++ b/src/plugin-sdk/test-helpers/raw-http-request.ts @@ -0,0 +1,172 @@ +/** + * Raw-socket HTTP request driver for webhook tests. + * + * `fetch` cannot express a server that answers while the sender is still uploading and + * then closes: undici surfaces the failed upload and discards the response it already + * received. Body-limit rejections are exactly that shape, so asserting them needs a + * client that reports the status line and the connection outcome independently. + */ +import net from "node:net"; + +export type RawHttpResult = { + /** First line of the response, or an empty string when the server sent nothing. */ + statusLine: string; + /** Lowercase response header names and their final values. */ + headers: Record; + /** Response body after the header block. */ + body: string; + /** True when the server closed the connection rather than leaving it open. */ + closedByServer: boolean; +}; + +/** + * Reassembles a chunked response body so callers can assert the payload the server sent. + * + * Chunk sizes count bytes, so the scan stays on the Buffer: decoding first would make a + * multi-byte character consume one slice position and misplace every later boundary. + */ +function decodeChunkedBody(raw: Buffer): Buffer { + const parts: Buffer[] = []; + let offset = 0; + for (;;) { + const lineEnd = raw.indexOf("\r\n", offset, "latin1"); + if (lineEnd === -1) { + break; + } + const size = Number.parseInt(raw.toString("latin1", offset, lineEnd).trim(), 16); + if (!Number.isFinite(size) || size <= 0) { + break; + } + parts.push(raw.subarray(lineEnd + 2, lineEnd + 2 + size)); + offset = lineEnd + 2 + size + 2; + } + return Buffer.concat(parts); +} + +export async function postRawWebhook(params: { + /** Absolute URL of the webhook endpoint. */ + url: string; + /** Request body. */ + body: string; + headers?: Record; + /** How long to keep the socket open before reporting it as retained. */ + idleTimeoutMs?: number; + /** + * Send the body incrementally instead of in one write, so the server decides while the + * upload is still active rather than from the declared length alone. + */ + chunk?: { bytes: number; intervalMs: number }; + /** Content-Length to declare; defaults to the body's real length. */ + contentLength?: number; + /** + * Send with chunked transfer encoding and no Content-Length, so a size limit can only be + * detected from the bytes that actually arrive rather than from a declared length. + */ + chunkedEncoding?: boolean; +}): Promise { + const target = new URL(params.url); + const port = Number(target.port); + const payload = Buffer.from(params.body, "utf-8"); + const idleTimeoutMs = params.idleTimeoutMs ?? 2_000; + const headerLines = Object.entries(params.headers ?? {}) + .map(([name, value]) => `${name}: ${value}\r\n`) + .join(""); + const head = + `POST ${target.pathname}${target.search} HTTP/1.1\r\n` + + `Host: ${target.hostname}:${port}\r\n` + + headerLines + + (params.chunkedEncoding + ? `Transfer-Encoding: chunked\r\n\r\n` + : `Content-Length: ${params.contentLength ?? payload.length}\r\n\r\n`); + + return await new Promise((resolve) => { + const socket = net.connect(port, target.hostname); + const received: Buffer[] = []; + let settled = false; + let timer: ReturnType | undefined; + + const settle = (closedByServer: boolean) => { + if (settled) { + return; + } + settled = true; + if (timer) { + clearTimeout(timer); + } + socket.destroy(); + const raw = Buffer.concat(received); + const headerEnd = raw.indexOf("\r\n\r\n", 0, "latin1"); + const headBlock = + headerEnd === -1 ? raw.toString("latin1") : raw.toString("latin1", 0, headerEnd); + const rawBody = headerEnd === -1 ? Buffer.alloc(0) : raw.subarray(headerEnd + 4); + const chunked = /transfer-encoding:\s*chunked/i.test(headBlock); + const headers = Object.fromEntries( + headBlock + .split("\r\n") + .slice(1) + .map((line) => { + const separator = line.indexOf(":"); + return [line.slice(0, separator).toLowerCase(), line.slice(separator + 1).trim()]; + }), + ); + resolve({ + statusLine: (headBlock.split("\r\n")[0] ?? "").trim(), + headers, + body: (chunked ? decodeChunkedBody(rawBody) : rawBody).toString("utf-8"), + closedByServer, + }); + }; + + // Chunk framing is ASCII and the body is raw bytes. Writing them together through a + // string would re-encode every byte above 0x7f after its length was already declared. + const writeBody = (slice: Buffer) => { + if (!params.chunkedEncoding) { + socket.write(slice); + return; + } + socket.write(`${slice.length.toString(16)}\r\n`); + socket.write(slice); + socket.write("\r\n"); + }; + + socket.on("connect", () => { + socket.write(head); + if (!params.chunk) { + writeBody(payload); + if (params.chunkedEncoding) { + socket.write("0\r\n\r\n"); + } + } else { + const { bytes, intervalMs } = params.chunk; + let offset = 0; + const pump = () => { + if (settled || socket.destroyed) { + return; + } + if (offset >= payload.length) { + if (params.chunkedEncoding) { + socket.write(`0\r\n\r\n`); + } + return; + } + const end = Math.min(offset + bytes, payload.length); + writeBody(payload.subarray(offset, end)); + offset = end; + setTimeout(pump, intervalMs).unref?.(); + }; + pump(); + } + // Only the idle timeout reports a retained connection; a server that answers and + // closes always reaches "close" below. + timer = setTimeout(() => settle(false), idleTimeoutMs); + timer.unref?.(); + }); + socket.on("data", (chunk: Buffer) => { + received.push(chunk); + }); + // Rejecting mid-upload makes the write fail on this side; the response may already be + // buffered, so wait for "close" rather than settling on the write error. + socket.on("error", () => {}); + socket.on("close", () => settle(true)); + }); +} diff --git a/src/plugins/contracts/plugin-sdk-runtime-api-guardrails.test.ts b/src/plugins/contracts/plugin-sdk-runtime-api-guardrails.test.ts index 5dbeab94fafc..9bb780975f47 100644 --- a/src/plugins/contracts/plugin-sdk-runtime-api-guardrails.test.ts +++ b/src/plugins/contracts/plugin-sdk-runtime-api-guardrails.test.ts @@ -302,7 +302,7 @@ const RUNTIME_API_EXPORT_GUARDS: Record = { 'export { definePluginEntry } from "openclaw/plugin-sdk/plugin-entry";', 'export type { OpenClawPluginApi } from "openclaw/plugin-sdk/plugin-entry";', 'export type { GatewayRequestHandlerOptions } from "openclaw/plugin-sdk/gateway-runtime";', - 'export { isRequestBodyLimitError, readRequestBodyWithLimit, requestBodyErrorToText } from "openclaw/plugin-sdk/webhook-request-guards";', + 'export { isRequestBodyLimitError, readRequestBodyWithLimit, requestBodyErrorToText, sendHttpRequestRejection } from "openclaw/plugin-sdk/webhook-request-guards";', 'export { fetchWithSsrFGuard, isBlockedHostnameOrIp } from "openclaw/plugin-sdk/ssrf-runtime";', 'export type { SessionEntry } from "openclaw/plugin-sdk/session-store-runtime";', 'export { TtsAutoSchema, TtsConfigSchema, TtsModeSchema, TtsProviderSchema } from "openclaw/plugin-sdk/tts-runtime";',