From 0684f50136e77bb832417d292fa20ad62380527d Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Tue, 18 Aug 2026 10:39:56 -0700 Subject: [PATCH] fix(gateway): harden hook admission and shutdown (#125893) * fix(gateway): harden hook admission and shutdown Flush structured 413 and 408 hook responses before closing oversized or timed-out request sockets, and reject deferred wake requests that target an explicit session. Direct Gateway close now marks startup and readiness as draining immediately and attempts every HTTP listener shutdown before reporting an incomplete close. * test(gateway): model request destruction in fuzzing Give the HTTP request fuzz double the IncomingMessage destroy contract exercised after 413 and 408 responses finish. * fix(gateway): preserve response and shutdown bounds Route installed webhook body-limit failures through the shared response-first close owner and require complete Feishu 413/408 responses before terminal connection closure. Retain the current-main grace bound for agent-harness disposal so one stuck harness cannot block listener teardown. * test(telegram): model response-close request cleanup Keep the slow-body request alive through response finish, then assert cleanup when the mock response transport closes. --- config/assertion-safety-baseline.txt | 2 +- docs/automation/cron-jobs.md | 2 +- docs/gateway/configuration-reference.md | 2 +- .../src/monitor.webhook-security.test.ts | 53 +++++----- .../telegram/src/miniapp/routes.test.ts | 2 + src/gateway/hooks.test.ts | 8 +- src/gateway/hooks.ts | 9 +- src/gateway/http-common.fuzz.test.ts | 9 +- src/gateway/http-common.test.ts | 32 +++++- src/gateway/http-common.ts | 3 + src/gateway/server-close.test.ts | 30 ++++++ src/gateway/server-close.ts | 100 ++++++++++-------- .../server-http.hooks-request-timeout.test.ts | 25 ++++- src/gateway/server-kernel.test.ts | 57 ++++++++++ src/gateway/server-lifecycle.ts | 6 +- src/gateway/server-runtime-state-prepare.ts | 4 +- src/gateway/server.hooks-admission.test.ts | 93 ++++++++++++++++ src/gateway/server/hooks-request-handler.ts | 4 + src/infra/http-body.test.ts | 49 ++++++++- src/infra/http-body.ts | 22 +++- src/plugin-sdk/webhook-request-guards.ts | 14 +-- 21 files changed, 420 insertions(+), 106 deletions(-) diff --git a/config/assertion-safety-baseline.txt b/config/assertion-safety-baseline.txt index f93d69229c00..caf7f917f4ea 100644 --- a/config/assertion-safety-baseline.txt +++ b/config/assertion-safety-baseline.txt @@ -2965,7 +2965,7 @@ src/gateway/server-broadcast.ts 2 src/gateway/server-channels.ts 1 src/gateway/server-chat-state.ts 1 src/gateway/server-chat.ts 5 -src/gateway/server-close.ts 2 +src/gateway/server-close.ts 1 src/gateway/server-core-runtime.ts 2 src/gateway/server-cron-reconciled.ts 1 src/gateway/server-cron.ts 5 diff --git a/docs/automation/cron-jobs.md b/docs/automation/cron-jobs.md index 6a57daf6fe2b..d6c60d9f6bdf 100644 --- a/docs/automation/cron-jobs.md +++ b/docs/automation/cron-jobs.md @@ -578,7 +578,7 @@ Query-string tokens are rejected. Target agent. When supplied, it must name a configured agent. It is required when the configured agent fleet has no implicit or retained legacy owner. - Target session. Requires `hooks.allowRequestSessionKey: true` and must match `hooks.allowedSessionKeyPrefixes` when configured. Both wake modes can target an explicit session. + Target session. Requires `mode: "now"` and `hooks.allowRequestSessionKey: true`, and must match `hooks.allowedSessionKeyPrefixes` when configured. Deferred `next-heartbeat` wakes use the agent's main session. diff --git a/docs/gateway/configuration-reference.md b/docs/gateway/configuration-reference.md index b74afd610637..0271649f832e 100644 --- a/docs/gateway/configuration-reference.md +++ b/docs/gateway/configuration-reference.md @@ -1081,7 +1081,7 @@ Validation and safety notes: **Endpoints:** - `POST /hooks/wake` → `{ text, mode?: "now"|"next-heartbeat", agentId?, sessionKey? }` - - `sessionKey` is accepted only when `hooks.allowRequestSessionKey=true` (default: `false`) and must match `hooks.allowedSessionKeyPrefixes` when configured. + - `sessionKey` requires `mode: "now"`, is accepted only when `hooks.allowRequestSessionKey=true` (default: `false`), and must match `hooks.allowedSessionKeyPrefixes` when configured. - A supplied `agentId` must name a configured agent. - `POST /hooks/agent` → `{ message, name?, agentId?, sessionKey?, sessionMode?, wakeMode?, deliver?, channel?, to?, accountId?, model?, thinking?, timeoutSeconds? }` - A supplied `agentId` must name a configured agent. diff --git a/extensions/feishu/src/monitor.webhook-security.test.ts b/extensions/feishu/src/monitor.webhook-security.test.ts index de13e87b76d1..b08bddf85f68 100644 --- a/extensions/feishu/src/monitor.webhook-security.test.ts +++ b/extensions/feishu/src/monitor.webhook-security.test.ts @@ -60,6 +60,7 @@ async function waitForSlowBodyTimeoutResponse( const target = new URL(url); const startedAt = Date.now(); let response = ""; + let settled = false; const socket = createConnection( { host: target.hostname, @@ -76,17 +77,27 @@ async function waitForSlowBodyTimeoutResponse( ); socket.setEncoding("utf8"); - socket.on("error", () => {}); socket.on("data", (chunk) => { response += chunk.toString(); - if (response.includes("Request body timeout")) { + }); + socket.on("close", () => { + if (settled) { + return; + } + settled = true; + clearTimeout(failTimer); + resolve({ body: response, elapsedMs: Date.now() - startedAt }); + }); + socket.on("error", (error) => { + if (!settled) { + settled = true; clearTimeout(failTimer); - socket.destroy(); - resolve({ body: response, elapsedMs: Date.now() - startedAt }); + reject(error); } }); const failTimer = setTimeout(() => { + settled = true; socket.destroy(); reject(new Error(`timeout response did not arrive within ${timeoutMs}ms`)); }, timeoutMs); @@ -120,32 +131,26 @@ async function waitForOversizedBodyResponse(url: string): Promise { } settled = true; clearTimeout(failTimer); - socket.destroy(); resolve(result); }; socket.setEncoding("utf8"); socket.on("data", (chunk) => { response += chunk.toString(); - if (response.includes("Payload too large")) { - finish(response); - } }); socket.on("close", () => { - if (response.includes("Payload too large")) { - finish(response); - } + finish(response); }); socket.on("error", (error: NodeJS.ErrnoException) => { - if (response.includes("Payload too large")) { - finish(response); - return; + if (!settled) { + if (response.includes("Payload too large")) { + finish(response); + return; + } + settled = true; + clearTimeout(failTimer); + reject(new Error(`${error.message}; partial response: ${JSON.stringify(response)}`)); } - if (error.code === "ECONNRESET") { - finish("ECONNRESET"); - return; - } - reject(error); }); const failTimer = setTimeout(() => { @@ -267,12 +272,9 @@ describe("Feishu webhook security hardening", () => { async (url) => { const response = await waitForOversizedBodyResponse(url); - if (response === "ECONNRESET") { - expect(response).toBe("ECONNRESET"); - } else { - expect(response).toContain("413 Payload Too Large"); - expect(response).toContain("Payload too large"); - } + expect(response).toContain("413 Payload Too Large"); + expect(response).toContain("Payload too large"); + expect(response).toMatch(/connection: close/i); }, ); }); @@ -291,6 +293,7 @@ describe("Feishu webhook security hardening", () => { const result = await waitForSlowBodyTimeoutResponse(url, 1_000); expect(result.body).toContain("408 Request Timeout"); expect(result.body).toContain("Request body timeout"); + expect(result.body).toMatch(/connection: close/i); expect(result.elapsedMs).toBeLessThan(500); }, ); diff --git a/extensions/telegram/src/miniapp/routes.test.ts b/extensions/telegram/src/miniapp/routes.test.ts index 16277626f5f4..9dcfb66d270c 100644 --- a/extensions/telegram/src/miniapp/routes.test.ts +++ b/extensions/telegram/src/miniapp/routes.test.ts @@ -484,6 +484,8 @@ describe("registerTelegramMiniAppRoutes", () => { expect(res.statusCode).toBe(408); expect(res.body).toBe("Request body timeout"); + expect(req.destroyed).toBe(false); + res.emit("close"); expect(req.destroyed).toBe(true); expectBodyReadListenersCleaned(req); expect(issueDeviceBootstrapToken).not.toHaveBeenCalled(); diff --git a/src/gateway/hooks.test.ts b/src/gateway/hooks.test.ts index cbabf5cfbe26..8a6fe5187657 100644 --- a/src/gateway/hooks.test.ts +++ b/src/gateway/hooks.test.ts @@ -149,12 +149,8 @@ describe("gateway hooks helpers", () => { sessionKey: "hook:wake:later", }), ).toEqual({ - ok: true, - value: { - text: "wake later", - mode: "next-heartbeat", - sessionKey: "hook:wake:later", - }, + ok: false, + error: "sessionKey requires mode=now", }); expect(normalizeWakePayload({ text: " ", mode: "now" }).ok).toBe(false); expect(normalizeWakePayload({ text: "wake", agentId: 42 })).toEqual({ diff --git a/src/gateway/hooks.ts b/src/gateway/hooks.ts index 970772aeb33c..ec4155d3419c 100644 --- a/src/gateway/hooks.ts +++ b/src/gateway/hooks.ts @@ -196,7 +196,11 @@ export async function readJsonBody( req: IncomingMessage, maxBytes: number, ): Promise> { - const result = await readJsonBodyWithLimit(req, { maxBytes, emptyObjectOnEmpty: true }); + const result = await readJsonBodyWithLimit(req, { + maxBytes, + emptyObjectOnEmpty: true, + destroyOnLimit: false, + }); if (result.ok) { return result; } @@ -256,6 +260,9 @@ export function normalizeWakePayload( if (payload.sessionKey !== undefined && !sessionKey) { return { ok: false, error: "sessionKey must be a non-empty string" }; } + if (mode === "next-heartbeat" && sessionKey) { + return { ok: false, error: "sessionKey requires mode=now" }; + } return { ok: true, value: { diff --git a/src/gateway/http-common.fuzz.test.ts b/src/gateway/http-common.fuzz.test.ts index 79ba47b84d5b..e49d224756da 100644 --- a/src/gateway/http-common.fuzz.test.ts +++ b/src/gateway/http-common.fuzz.test.ts @@ -248,7 +248,14 @@ describe("fuzz: sendInvalidRequest", () => { }); describe("fuzz: readJsonBodyOrError", () => { - const makeRequest = () => ({}) as IncomingMessage; + const makeRequest = () => { + const req = { destroyed: false } as IncomingMessage; + req.destroy = vi.fn(() => { + req.destroyed = true; + return req; + }); + return req; + }; it("maps readJsonBody results to the documented status/body contract", async () => { const rng = makeRng(0xc0de); diff --git a/src/gateway/http-common.test.ts b/src/gateway/http-common.test.ts index 5df82e7e926c..e6493caf292f 100644 --- a/src/gateway/http-common.test.ts +++ b/src/gateway/http-common.test.ts @@ -227,7 +227,17 @@ describe("sendInvalidRequest", () => { }); describe("readJsonBodyOrError", () => { - const makeRequest = () => ({}) as IncomingMessage; + const makeRequest = (headers: Record = {}) => { + const req = Object.assign(new EventEmitter(), { + destroyed: false, + headers, + destroy: vi.fn(() => { + req.destroyed = true; + return req; + }), + }) as IncomingMessage & { destroy: ReturnType }; + return req; + }; it("returns the parsed body on success", async () => { readJsonBodyMock.mockResolvedValueOnce({ ok: true, value: { hello: "world" } }); @@ -243,7 +253,7 @@ describe("readJsonBodyOrError", () => { const events: DiagnosticEventPayload[] = []; const stop = onDiagnosticEvent((event) => events.push(event)); const { res, end } = makeMockHttpResponse(); - const req = { headers: { "content-length": "2048" } } as IncomingMessage; + const req = makeRequest({ "content-length": "2048" }); const result = await readJsonBodyOrError(req, res, 1024); stop(); expect(result).toBeUndefined(); @@ -259,12 +269,18 @@ describe("readJsonBodyOrError", () => { expect(event?.bytes).toBe(2048); expect(event?.limitBytes).toBe(1024); expect(event?.reason).toBe("json_body_limit"); + expect(req.destroy).not.toHaveBeenCalled(); + res.emit("finish"); + expect(req.destroy).not.toHaveBeenCalled(); + res.emit("close"); + expect(req.destroy).toHaveBeenCalledOnce(); }); it("responds with 408 when the request body times out", async () => { readJsonBodyMock.mockResolvedValueOnce({ ok: false, error: "request body timeout" }); const { res, end } = makeMockHttpResponse(); - const result = await readJsonBodyOrError(makeRequest(), res, 1024); + const req = makeRequest(); + const result = await readJsonBodyOrError(req, res, 1024); expect(result).toBeUndefined(); expect(res.statusCode).toBe(408); expect(end).toHaveBeenCalledWith( @@ -272,17 +288,25 @@ describe("readJsonBodyOrError", () => { error: { message: "Request body timeout", type: "invalid_request_error" }, }), ); + expect(req.destroy).not.toHaveBeenCalled(); + res.emit("finish"); + expect(req.destroy).not.toHaveBeenCalled(); + res.emit("close"); + expect(req.destroy).toHaveBeenCalledOnce(); }); it("responds with 400 for other parse failures", async () => { readJsonBodyMock.mockResolvedValueOnce({ ok: false, error: "bad json" }); const { res, end } = makeMockHttpResponse(); - const result = await readJsonBodyOrError(makeRequest(), res, 1024); + const req = makeRequest(); + const result = await readJsonBodyOrError(req, res, 1024); expect(result).toBeUndefined(); expect(res.statusCode).toBe(400); expect(end).toHaveBeenCalledWith( JSON.stringify({ error: { message: "bad json", type: "invalid_request_error" } }), ); + res.emit("finish"); + expect(req.destroy).not.toHaveBeenCalled(); }); }); diff --git a/src/gateway/http-common.ts b/src/gateway/http-common.ts index cbbb82cf737e..1473c7502b9c 100644 --- a/src/gateway/http-common.ts +++ b/src/gateway/http-common.ts @@ -2,6 +2,7 @@ // body-size errors, and client disconnect aborts. import type { IncomingMessage, ServerResponse } from "node:http"; import { buildMissingScopeErrorDetails } from "../../packages/gateway-protocol/src/index.js"; +import { closeRequestAfterResponse } from "../infra/http-body.js"; import { logRejectedLargePayload, parseContentLengthHeader, @@ -150,12 +151,14 @@ export async function readJsonBodyOrError( reason: "json_body_limit", ...(contentLength !== undefined ? { bytes: contentLength } : {}), }); + closeRequestAfterResponse(req, res); sendJson(res, 413, { error: { message: "Payload too large", type: "invalid_request_error" }, }); return undefined; } if (body.error === "request body timeout") { + closeRequestAfterResponse(req, res); sendJson(res, 408, { error: { message: "Request body timeout", type: "invalid_request_error" }, }); diff --git a/src/gateway/server-close.test.ts b/src/gateway/server-close.test.ts index 2a5e5f0ef46e..32bda0871dee 100644 --- a/src/gateway/server-close.test.ts +++ b/src/gateway/server-close.test.ts @@ -1960,6 +1960,36 @@ describe("createGatewayCloseHandler", () => { ).toBe(true); }); + it("attempts every HTTP listener before rejecting a stuck close", async () => { + vi.useFakeTimers(); + + const stuckServer = { + close: vi.fn(() => undefined), + closeAllConnections: vi.fn(), + closeIdleConnections: vi.fn(), + }; + const laterServer = { + close: vi.fn((cb: (err?: Error | null) => void) => cb(null)), + closeIdleConnections: vi.fn(), + }; + const close = createGatewayCloseHandler( + createGatewayCloseTestDeps({ + httpServers: [stuckServer as never, laterServer as never], + }), + ); + + const closePromise = close({ reason: "test shutdown" }); + const closeExpectation = expect(closePromise).rejects.toThrow( + "http-server[0] close still pending after forced connection shutdown (5000ms)", + ); + await vi.waitFor(() => expect(stuckServer.close).toHaveBeenCalledOnce()); + expect(laterServer.close).toHaveBeenCalledOnce(); + await vi.advanceTimersByTimeAsync(HTTP_CLOSE_GRACE_MS + HTTP_CLOSE_FORCE_WAIT_MS); + await closeExpectation; + + expect(stuckServer.closeAllConnections).toHaveBeenCalledOnce(); + }); + it("labels warnings for multiple HTTP servers with their index", async () => { const okServer = { close: (cb: (err?: Error | null) => void) => cb(null), diff --git a/src/gateway/server-close.ts b/src/gateway/server-close.ts index c6a0b9523207..56cc52067ad0 100644 --- a/src/gateway/server-close.ts +++ b/src/gateway/server-close.ts @@ -664,6 +664,50 @@ async function waitForHttpClose(params: { } } +async function closeHttpListener(params: { + server: HttpServer; + label: string; + warnings: string[]; +}): Promise { + const { server, label, warnings } = params; + server.closeIdleConnections?.(); + const closePromise = new Promise((resolve, reject) => { + server.close((err) => { + if (!err || isServerNotRunningError(err)) { + resolve(); + return; + } + reject(err); + }); + }); + void closePromise.catch(() => undefined); + const closedWithinGrace = await waitForHttpClose({ + closePromise, + timeoutMs: HTTP_CLOSE_GRACE_MS, + label, + warnings, + }); + if (closedWithinGrace) { + return; + } + shutdownLog.warn( + `${label} close exceeded ${HTTP_CLOSE_GRACE_MS}ms; forcing connection shutdown and waiting for close`, + ); + recordShutdownWarning(warnings, label); + server.closeAllConnections?.(); + const closedAfterForce = await waitForHttpClose({ + closePromise, + timeoutMs: HTTP_CLOSE_FORCE_WAIT_MS, + label, + warnings, + }); + if (!closedAfterForce) { + throw new Error( + `${label} close still pending after forced connection shutdown (${HTTP_CLOSE_FORCE_WAIT_MS}ms)`, + ); + } +} + export function createGatewayCloseHandler( params: { bonjourStop: (() => Promise) | null; @@ -1021,50 +1065,20 @@ export function createGatewayCloseHandler( try { if (transportServers.length > 0) { await measureCloseStep("http-server", async () => { - const servers = transportServers; - for (let i = 0; i < servers.length; i++) { - const httpServer = servers[i] as HttpServer & { - closeAllConnections?: () => void; - closeIdleConnections?: () => void; - }; - const label = servers.length > 1 ? `http-server[${i}]` : "http-server"; - if (typeof httpServer.closeIdleConnections === "function") { - httpServer.closeIdleConnections(); - } - const closePromise = new Promise((resolve, reject) => { - httpServer.close((err) => { - if (!err || isServerNotRunningError(err)) { - resolve(); - return; - } - reject(err); - }); - }); - void closePromise.catch(() => undefined); - const closedWithinGrace = await waitForHttpClose({ - closePromise, - timeoutMs: HTTP_CLOSE_GRACE_MS, - label, - warnings, - }); - if (!closedWithinGrace) { - shutdownLog.warn( - `${label} close exceeded ${HTTP_CLOSE_GRACE_MS}ms; forcing connection shutdown and waiting for close`, - ); - recordShutdownWarning(warnings, label); - httpServer.closeAllConnections?.(); - const closedAfterForce = await waitForHttpClose({ - closePromise, - timeoutMs: HTTP_CLOSE_FORCE_WAIT_MS, - label, + const results = await Promise.allSettled( + transportServers.map((server, index) => + closeHttpListener({ + server, + label: transportServers.length > 1 ? `http-server[${index}]` : "http-server", warnings, - }); - if (!closedAfterForce) { - throw new Error( - `${label} close still pending after forced connection shutdown (${HTTP_CLOSE_FORCE_WAIT_MS}ms)`, - ); - } - } + }), + ), + ); + const failure = results.find( + (result): result is PromiseRejectedResult => result.status === "rejected", + ); + if (failure) { + throw failure.reason; } }); } diff --git a/src/gateway/server-http.hooks-request-timeout.test.ts b/src/gateway/server-http.hooks-request-timeout.test.ts index c43d88698ea1..239b235a4ddf 100644 --- a/src/gateway/server-http.hooks-request-timeout.test.ts +++ b/src/gateway/server-http.hooks-request-timeout.test.ts @@ -1,6 +1,8 @@ /** * Tests timeout behavior for gateway HTTP hook request handling. */ +import { EventEmitter } from "node:events"; +import type { ServerResponse } from "node:http"; import { beforeEach, describe, expect, test, vi } from "vitest"; import { createHookRequest, @@ -40,14 +42,33 @@ describe("createHooksRequestHandler timeout status mapping", () => { const dispatchWakeHook = vi.fn(); const dispatchAgentHook = vi.fn(() => ({ ok: true as const, runId: "run-1" })); const handler = createHooksHandler({ dispatchWakeHook, dispatchAgentHook }); - const req = createHookRequest(); - const { res, end } = createResponse(); + const req = createHookRequest() as ReturnType & { + destroyed: boolean; + destroy: ReturnType; + }; + req.destroyed = false; + req.destroy = vi.fn(() => { + req.destroyed = true; + return req; + }); + const res = new EventEmitter() as ServerResponse; + res.statusCode = 200; + const setHeader = vi.fn(); + res.setHeader = setHeader; + const end = vi.fn(); + res.end = end; const handled = await handler(req, res); expect(handled).toBe(true); expect(res.statusCode).toBe(408); expect(end).toHaveBeenCalledWith(JSON.stringify({ ok: false, error: "request body timeout" })); + expect(setHeader).toHaveBeenCalledWith("Connection", "close"); + expect(req.destroy).not.toHaveBeenCalled(); + res.emit("finish"); + expect(req.destroy).not.toHaveBeenCalled(); + res.emit("close"); + expect(req.destroy).toHaveBeenCalledOnce(); expect(dispatchWakeHook).not.toHaveBeenCalled(); expect(dispatchAgentHook).not.toHaveBeenCalled(); }); diff --git a/src/gateway/server-kernel.test.ts b/src/gateway/server-kernel.test.ts index 4b5691a68a37..38235e01b60d 100644 --- a/src/gateway/server-kernel.test.ts +++ b/src/gateway/server-kernel.test.ts @@ -1,5 +1,6 @@ import fs from "node:fs/promises"; import { describe, expect, it, vi } from "vitest"; +import { createDeferred } from "../../test/helpers/promise.js"; import { createPluginRecord } from "../plugins/loader-records.js"; import { createEmptyPluginRegistry } from "../plugins/registry-empty.js"; import { @@ -18,6 +19,62 @@ import { createGatewayKernel } from "./server-kernel.js"; import { createSyntheticPluginRuntimeClient } from "./server-plugin-runtime-client.js"; describe("createGatewayKernel", () => { + it("reports startup and readiness as draining during a direct close", async () => { + const port = await getFreePort(); + const state = await createOpenClawTestState({ + label: "gateway-kernel-direct-close-readiness", + layout: "home", + env: { + OPENCLAW_GATEWAY_PASSWORD: undefined, + OPENCLAW_GATEWAY_TOKEN: undefined, + OPENCLAW_SKIP_BROWSER_CONTROL_SERVER: "1", + OPENCLAW_SKIP_CANVAS_HOST: "1", + OPENCLAW_SKIP_CHANNELS: "1", + OPENCLAW_SKIP_CRON: "1", + OPENCLAW_SKIP_GMAIL_WATCHER: "1", + OPENCLAW_SKIP_PROVIDERS: "1", + OPENCLAW_TEST_MINIMAL_GATEWAY: "1", + VITEST: "1", + }, + }); + const token = "gateway-kernel-direct-close-readiness-token"; + let kernel: Awaited> | undefined; + try { + await state.writeConfig({ + gateway: { auth: { mode: "token", token }, controlUi: { enabled: false }, port }, + }); + state.applyEnv(); + kernel = await createGatewayKernel(port, { + auth: { mode: "token", token }, + bind: "loopback", + controlUiEnabled: false, + sidecarStartup: "defer", + }); + kernel.kernel.unlockStartupMethods(); + kernel.kernel.markSidecarsReady(); + const { getStartup, getReadiness } = kernel.createHttpTransportOptions(); + expect(getStartup()).toMatchObject({ ok: true, status: "started" }); + expect(getReadiness()).toMatchObject({ ready: true, failing: [] }); + + const configReloaderStop = createDeferred(); + vi.spyOn(kernel.runtimeState.configReloader, "stop").mockReturnValue( + configReloaderStop.promise, + ); + const closing = kernel.createCloseHandler()({ reason: "direct close readiness test" }); + + expect(getStartup()).toMatchObject({ ok: false, status: "draining" }); + expect(getReadiness()).toMatchObject({ ready: false, failing: ["gateway-draining"] }); + configReloaderStop.resolve(); + await closing; + } finally { + try { + await kernel?.closeOnStartupFailure(); + } finally { + await state.cleanup(); + } + } + }); + it("keeps startup readiness and sidecar shutdown at their lifecycle boundaries", async () => { const port = await getFreePort(); const state = await createOpenClawTestState({ diff --git a/src/gateway/server-lifecycle.ts b/src/gateway/server-lifecycle.ts index 6f252f47472c..8e7aa5d0577b 100644 --- a/src/gateway/server-lifecycle.ts +++ b/src/gateway/server-lifecycle.ts @@ -86,6 +86,7 @@ export async function prepareGatewayLifecycle(params: { nodeDesktopStreamBroker, bindDeviceNodeControl, workerPlacementRuntime, + lifecycle, } = runtime; const subscribeSessionMessageEvents: GatewayRequestContext["subscribeSessionMessageEvents"] = ( connId, @@ -346,7 +347,6 @@ export async function prepareGatewayLifecycle(params: { }, }; - const lifecycle = { closePreludeStarted: false }; const cronReconciliation = createGatewayCronReconciliation({ port, workspaceDir: defaultWorkspaceDir, @@ -387,6 +387,9 @@ export async function prepareGatewayLifecycle(params: { return mediaCleanupStopPromise; }; const markClosePreludeStarted = () => { + if (lifecycle.closePreludeStarted) { + return; + } lifecycle.closePreludeStarted = true; postReadySidecarStopOwner.beginClose(); gatewayLifetimeSidecarStopOwner.beginClose(); @@ -487,6 +490,7 @@ export async function prepareGatewayLifecycle(params: { } }; const createCloseHandler = () => async (optsValue?: GatewayCloseOptions) => { + markClosePreludeStarted(); const channelIds = listLoadedChannelPlugins().map((plugin) => plugin.id as ChannelId); const transport = transportBridge.current(); await transport?.portalService.closeAll(); diff --git a/src/gateway/server-runtime-state-prepare.ts b/src/gateway/server-runtime-state-prepare.ts index 9ace248fb10f..58e7e26328cf 100644 --- a/src/gateway/server-runtime-state-prepare.ts +++ b/src/gateway/server-runtime-state-prepare.ts @@ -364,6 +364,7 @@ export async function prepareGatewayKernelState(params: { pendingReason: "startup-sidecars", dispatchReady: false, }; + const lifecycle = { closePreludeStarted: false }; let releaseStartupAccountStarts = () => {}; const startupAccountStartsReady = new Promise((resolve) => { releaseStartupAccountStarts = resolve; @@ -401,7 +402,7 @@ export async function prepareGatewayKernelState(params: { startedAt: serverStartedAt, getStartupPending: isGatewayStartupPending, getStartupPendingReason: () => startupState.pendingReason, - getGatewayDraining: isGatewayDraining, + getGatewayDraining: () => lifecycle.closePreludeStarted || isGatewayDraining(), }; const getStartup = createStartupChecker(startupCheckerDeps); const getReadiness = createReadinessChecker({ @@ -544,6 +545,7 @@ export async function prepareGatewayKernelState(params: { gatewayTls, readinessEventLoopHealth, startupState, + lifecycle, releaseStartupAccountStarts, gatewayInstanceRuntimeRef, channelManager, diff --git a/src/gateway/server.hooks-admission.test.ts b/src/gateway/server.hooks-admission.test.ts index b2017aa57b92..e5249104f13b 100644 --- a/src/gateway/server.hooks-admission.test.ts +++ b/src/gateway/server.hooks-admission.test.ts @@ -1,7 +1,9 @@ /** Focused HTTP coverage for hook admission feedback and pending replay behavior. */ +import { Agent, request as httpRequest } from "node:http"; import { afterEach, describe, expect, test, vi } from "vitest"; import { createDeferred } from "../../test/helpers/promise.js"; import { resolveMainSessionKeyFromConfig } from "../config/sessions.js"; +import { DEFAULT_WEBHOOK_MAX_BODY_BYTES } from "../infra/http-body.js"; import { drainSystemEvents } from "../infra/system-events.js"; import { cronIsolatedRun, @@ -50,7 +52,98 @@ async function waitForDuplicateRequest(): Promise { }); } +async function postOversizedChunkedHook(port: number): Promise<{ + statusCode: number | undefined; + body: string; + connection: string | undefined; + events: string[]; +}> { + const agent = new Agent({ keepAlive: true }); + const events: string[] = []; + try { + return await new Promise((resolve, reject) => { + const req = httpRequest( + { + agent, + host: "127.0.0.1", + port, + path: "/hooks/wake", + method: "POST", + headers: { + Authorization: `Bearer ${HOOK_TOKEN}`, + "Content-Type": "application/json", + }, + }, + (res) => { + const chunks: Buffer[] = []; + res.on("data", (chunk: Buffer) => chunks.push(chunk)); + res.on("error", reject); + res.on("end", () => { + events.push("response-end"); + void socketClosed.then(() => { + events.push("socket-close"); + resolve({ + statusCode: res.statusCode, + body: Buffer.concat(chunks).toString("utf8"), + connection: res.headers.connection, + events, + }); + }, reject); + }); + }, + ); + const socketClosed = new Promise((resolveClose) => { + req.once("socket", (socket) => socket.once("close", resolveClose)); + }); + req.on("error", reject); + req.setTimeout(5_000, () => req.destroy(new Error("chunked hook request timed out"))); + req.write('{"text":"'); + req.write("x".repeat(DEFAULT_WEBHOOK_MAX_BODY_BYTES + 1)); + req.end('"}'); + }); + } finally { + agent.destroy(); + } +} + describe("gateway hook admission", () => { + test("flushes an oversized chunked hook response before closing the socket", async () => { + testState.hooksConfig = { enabled: true, token: HOOK_TOKEN }; + await withGatewayServer(async ({ port }) => { + const response = await postOversizedChunkedHook(port); + + expect(response).toEqual({ + statusCode: 413, + body: JSON.stringify({ ok: false, error: "payload too large" }), + connection: "close", + events: ["response-end", "socket-close"], + }); + }); + }); + + test("rejects deferred wake delivery to an explicit session", async () => { + testState.hooksConfig = { + enabled: true, + token: HOOK_TOKEN, + allowRequestSessionKey: true, + allowedSessionKeyPrefixes: ["hook:"], + }; + await withGatewayServer(async ({ port }) => { + const response = await postHook( + port, + "/hooks/wake", + { text: "Wake later", mode: "next-heartbeat", sessionKey: "hook:wake:later" }, + "deferred-custom-wake", + ); + + expect(response.status).toBe(400); + await expect(response.json()).resolves.toEqual({ + ok: false, + error: "sessionKey requires mode=now", + }); + }); + }); + test("shares one pending persistent dispatch without losing its session target", async () => { testState.hooksConfig = { enabled: true, diff --git a/src/gateway/server/hooks-request-handler.ts b/src/gateway/server/hooks-request-handler.ts index e0cf640f3953..cd53bc8dd75f 100644 --- a/src/gateway/server/hooks-request-handler.ts +++ b/src/gateway/server/hooks-request-handler.ts @@ -1,6 +1,7 @@ // Hook request handler validates hook tokens, applies mappings, dedupes requests, and dispatches wake or agent work. import { createHash } from "node:crypto"; import type { IncomingMessage, ServerResponse } from "node:http"; +import { closeRequestAfterResponse } from "../../infra/http-body.js"; import { pruneMapToMaxSize } from "../../infra/map-size.js"; import type { createSubsystemLogger } from "../../logging/subsystem.js"; import { resolveHookExternalContentSource as resolveHookExternalContentSourceFromSession } from "../../security/external-content.js"; @@ -296,6 +297,9 @@ export function createHooksRequestHandler( : body.error === "request body timeout" ? 408 : 400; + if (status === 413 || status === 408) { + closeRequestAfterResponse(req, res); + } sendJson(res, status, { ok: false, error: body.error }); return true; } diff --git a/src/infra/http-body.test.ts b/src/infra/http-body.test.ts index fac35a2f47db..bbe567457485 100644 --- a/src/infra/http-body.test.ts +++ b/src/infra/http-body.test.ts @@ -1,10 +1,11 @@ // Tests HTTP body reading and size-limit handling. import { EventEmitter } from "node:events"; -import type { IncomingMessage } from "node:http"; +import type { IncomingMessage, ServerResponse } from "node:http"; import { MAX_TIMER_TIMEOUT_MS } from "@openclaw/normalization-core/number-coercion"; import { beforeEach, describe, expect, it, vi } from "vitest"; import { createMockServerResponse } from "../test-utils/mock-http-response.js"; import { + closeRequestAfterResponse, installRequestBodyLimitGuard, RequestBodyLimitError, type RequestBodyLimitErrorCode, @@ -330,4 +331,50 @@ describe("http body limits", () => { expect(req.destroyed).toBe(false); expect(pause).toHaveBeenCalledOnce(); }); + + it("closes a limited request only after its response transport closes", () => { + const req = createMockRequest({ emitEnd: false }); + const res = new EventEmitter() as ServerResponse; + const setHeader = vi.fn(); + res.setHeader = setHeader; + + closeRequestAfterResponse(req, res); + + expect(setHeader).toHaveBeenCalledWith("Connection", "close"); + expect(req.destroyed).toBe(false); + res.emit("finish"); + expect(req.destroyed).toBe(false); + res.emit("close"); + expect(req.destroyed).toBe(true); + }); + + it("flushes an installed guard response before destroying the request", async () => { + const req = createMockRequest({ chunks: ["oversized"], emitEnd: false }); + const res = new EventEmitter() as ServerResponse; + const setHeader = vi.fn(); + const end = vi.fn(); + res.setHeader = setHeader; + res.end = end; + + installRequestBodyLimitGuard(req, res, { maxBytes: 1 }); + await waitForMicrotaskTurn(); + + expect(res.statusCode).toBe(413); + expect(setHeader).toHaveBeenCalledWith("Connection", "close"); + expect(end).toHaveBeenCalledWith(JSON.stringify({ error: "Payload too large" })); + expect(req.destroyed).toBe(false); + res.emit("finish"); + expect(req.destroyed).toBe(false); + res.emit("close"); + expect(req.destroyed).toBe(true); + }); + + it("allows lightweight responses without finish listeners", () => { + const req = createMockRequest({ emitEnd: false }); + const res = createMockServerResponse(); + + expect(() => closeRequestAfterResponse(req, res)).not.toThrow(); + expect(res.getHeader("connection")).toBe("close"); + expect(req.destroyed).toBe(false); + }); }); diff --git a/src/infra/http-body.ts b/src/infra/http-body.ts index 2df03cf2bf37..4ab081283d38 100644 --- a/src/infra/http-body.ts +++ b/src/infra/http-body.ts @@ -155,6 +155,22 @@ function stopRequestBodyAfterLimit(req: IncomingMessage, destroyOnLimit: boolean req.pause(); } +/** Close a limited request only after its response transport has closed. */ +export function closeRequestAfterResponse(req: IncomingMessage, res: ServerResponse): void { + if (!res.headersSent) { + res.setHeader("Connection", "close"); + } + const once = Reflect.get(res, "once"); + if (typeof once !== "function") { + return; + } + once.call(res, "close", () => { + if (!req.destroyed) { + req.destroy(); + } + }); +} + type ReadResponsePrefixResult = { buffer: Buffer; size: number; @@ -529,6 +545,7 @@ export function installRequestBodyLimitGuard( }; const respond = (error: RequestBodyLimitError) => { + closeRequestAfterResponse(req, res); const text = customText[error.code] ?? requestBodyErrorToText(error.code); if (!res.headersSent) { res.statusCode = error.statusCode; @@ -550,11 +567,6 @@ export function installRequestBodyLimitGuard( reason = error.code; finish(); respond(error); - if (!req.destroyed) { - // Limit violations are expected user input; destroying with an Error causes - // an async 'error' event which can crash the process if no listener remains. - req.destroy(); - } }; const onData = (chunk: Buffer | string) => { diff --git a/src/plugin-sdk/webhook-request-guards.ts b/src/plugin-sdk/webhook-request-guards.ts index 0461f9886d53..8867c2d2f0aa 100644 --- a/src/plugin-sdk/webhook-request-guards.ts +++ b/src/plugin-sdk/webhook-request-guards.ts @@ -4,6 +4,7 @@ import { resolveIntegerOption } from "@openclaw/normalization-core/number-coerci import { normalizeOptionalLowercaseString } from "../../packages/normalization-core/src/string-coerce.js"; import { formatErrorMessage } from "../infra/errors.js"; import { + closeRequestAfterResponse, isRequestBodyLimitError, readJsonBodyWithLimit, readRequestBodyWithLimit, @@ -113,19 +114,6 @@ function respondWebhookBodyReadError(params: { return { ok: false }; } -function closeRequestAfterResponse(req: IncomingMessage, res: ServerResponse): void { - const once = Reflect.get(res, "once"); - if (typeof once !== "function") { - return; - } - res.setHeader("Connection", "close"); - once.call(res, "finish", () => { - if (!req.destroyed) { - req.destroy(); - } - }); -} - /** Create an in-memory limiter that caps concurrent webhook handlers per key. */ export function createWebhookInFlightLimiter(options?: { /** Maximum concurrent handlers allowed for one key. */