fix(channels): release rejected webhook connections after answering them

Deliver webhook body-limit and timeout responses before closing their request connections.

Use the shared HTTP rejection lifecycle across channel webhook handlers while preserving each endpoint's status and response body. Add raw-socket regressions for response bytes, close headers, and connection cleanup.

Closes #126808

Co-authored-by: Ayaan Zaidi <hi@obviy.us>
This commit is contained in:
Eden 2026-09-02 10:41:28 +08:00 • committed by GitHub
parent e77e16c085
commit 32b4abb341
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
41 changed files with 1544 additions and 336 deletions

View file

@ -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<typeof import("node:timers")>();
return {
...actual,
setTimeout: ((callback: (...args: unknown[]) => void, delay?: number, ...args: unknown[]) =>
globalThis.setTimeout(callback, delay, ...args)) as typeof actual.setTimeout,
clearTimeout: ((timer: ReturnType<typeof globalThis.setTimeout> | undefined) =>
globalThis.clearTimeout(timer)) as typeof actual.clearTimeout,
};
});
const activeStores = new Set<A2aTaskStore>();
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<void>((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<void>((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<void>((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<void>((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<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}
});
});

View file

@ -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(

View file

@ -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<void>((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;

View file

@ -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);

View file

@ -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 &&

View file

@ -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)}`));

View file

@ -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<void>;
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<typeof import("openclaw/plugin-sdk/webhook-ingress")>(
"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<void>) => {
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<void>((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<void>((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" });
});
});
});

View file

@ -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<string>;
/**
* 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<boolean> {
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<ReturnType<typeof createLineBot>, "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)}`));

View file

@ -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<void>((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<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}
});
});

View file

@ -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<string> {
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" }));

View file

@ -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";

View file

@ -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,

View file

@ -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<void>((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<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}
});
it("rejects the startup token when Mattermost has rotated the current command token", async () => {

View file

@ -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;
}

View file

@ -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;
}

View file

@ -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<string, string>;
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",

View file

@ -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<void> {
if (res.headersSent) {
return;
}
await sendHttpRequestRejection(req, res, status, JSON.stringify({ error }), "application/json");
}
export function createNextcloudTalkWebhookServer(opts: NextcloudTalkWebhookServerOptions): {
server: Server;
start: () => Promise<void>;
@ -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) {

View file

@ -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<boolean>;
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<boolean> {
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;
}

View file

@ -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<number> {
await new Promise<void>((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<void> {
if (!server.listening) {
return;
}
await new Promise<void>((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<Record<string, string>> => {
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();
}
});
});

View file

@ -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 {

View file

@ -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();
}
},
);
});

View file

@ -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<v
});
}
export function writeQaRequestBodyLimitError(res: ServerResponse, error: unknown): boolean {
export async function writeQaRequestBodyLimitError(
req: IncomingMessage,
res: ServerResponse,
error: unknown,
): Promise<boolean> {
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);

View file

@ -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<typeof writeError>[0], error: unknown): void {
if (writeQaRequestBodyLimitError(res, error)) {
async function writeQaLabServerError(
req: IncomingMessage,
res: Parameters<typeof writeError>[0],
error: unknown,
): Promise<void> {
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);
}
});
});

View file

@ -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 });

View file

@ -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) => {

View file

@ -309,13 +309,17 @@ export async function readTwilioWebhookForm(req: IncomingMessage): Promise<Recor
const body = await readRequestBodyWithLimit(req, {
maxBytes: WEBHOOK_BODY_LIMIT_BYTES,
timeoutMs: WEBHOOK_BODY_TIMEOUT_MS,
// Defer destruction so the webhook can answer 413 before the connection closes.
destroyOnLimit: false,
});
return parseTwilioFormBody(body);
}
export const TWIML_CONTENT_TYPE = "text/xml; charset=utf-8";
export function respondTwiml(res: ServerResponse, statusCode: number, body = ""): void {
res.statusCode = statusCode;
res.setHeader("content-type", "text/xml; charset=utf-8");
res.setHeader("content-type", TWIML_CONTENT_TYPE);
res.end(body || "<Response></Response>");
}

View file

@ -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> = {},
): 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<SmsDeliveryRecorder["record"]>(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 };
}

View file

@ -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> = {}): 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<string, string> } {
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<SmsDeliveryRecorder["record"]>(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<SmsDeliveryRecorder["record"]>(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<SmsDeliveryRecorder["record"]>(async ({ form }) => {
await new Promise<void>((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<string, string> = {
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,

View file

@ -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;

View file

@ -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<typeof import("node:timers")>();
return {
...actual,
setTimeout: ((callback: (...args: unknown[]) => void, delay?: number, ...args: unknown[]) =>
globalThis.setTimeout(callback, delay, ...args)) as typeof actual.setTimeout,
clearTimeout: ((timer: ReturnType<typeof globalThis.setTimeout> | 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<void>((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<void>((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<void>((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<void>((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<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}
});
});

View file

@ -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<void>((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<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}
});
it("rejects excess concurrent pre-auth body reads from the same remote IP", async () => {

View file

@ -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 };
}

View file

@ -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();
},
);

View file

@ -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") {

View file

@ -9,6 +9,7 @@ export {
type OpenClawPluginApi,
readRequestBodyWithLimit,
requestBodyErrorToText,
sendHttpRequestRejection,
type SessionEntry,
sleep,
TtsAutoSchema,

View file

@ -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";

View file

@ -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();
}
});
});

View file

@ -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<void> {
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<WebhookResponsePayload> {
res: http.ServerResponse,
): Promise<WebhookResponsePayload | null> {
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<string> {
return readRequestBodyWithLimit(req, { maxBytes, timeoutMs });
// Defer destruction so a limit rejection can be answered before the close.
return readRequestBodyWithLimit(req, { maxBytes, timeoutMs, destroyOnLimit: false });
}
/**

View file

@ -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";

View file

@ -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<string, string>;
/** 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<string, string>;
/** 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<RawHttpResult> {
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<RawHttpResult>((resolve) => {
const socket = net.connect(port, target.hostname);
const received: Buffer[] = [];
let settled = false;
let timer: ReturnType<typeof setTimeout> | 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));
});
}

View file

@ -302,7 +302,7 @@ const RUNTIME_API_EXPORT_GUARDS: Record<string, readonly string[]> = {
'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";',