From b3eded1744206bc0b0a2cec57ed5629e669d71fc Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 27 Aug 2026 11:34:07 -0700 Subject: [PATCH] fix(hooks): run plugin triggers with HTTP hooks disabled (#131059) * fix(hooks): schedule plugin turns independently of HTTP hooks * docs(hooks): clarify shared HTTP capacity reservation --- docs/plugins/sdk-runtime.md | 6 +- .../server/hooks.early-failure.test.ts | 205 +++++++++++------- src/gateway/server/hooks.ts | 8 +- 3 files changed, 136 insertions(+), 83 deletions(-) diff --git a/docs/plugins/sdk-runtime.md b/docs/plugins/sdk-runtime.md index 3ef6b7921281..e6955a98e864 100644 --- a/docs/plugins/sdk-runtime.md +++ b/docs/plugins/sdk-runtime.md @@ -416,8 +416,10 @@ snapshots; OpenClaw owns all persistence and lifecycle coordination. Dispatch isolated agent turns for untrusted external-content triggers, such as an email watcher. Unlike `api.runtime.subagent.run(...)`, hook dispatch - wraps external content, serializes runs for the same session, and uses the - Gateway hook execution lane and completion reporting. + wraps external content, serializes runs for the same session, and reports + completion through the Gateway. Plugin turns share the cron execution + budget without requiring the HTTP hooks endpoint. When HTTP hooks are + enabled, one slot in that shared budget remains reserved for HTTP work. ```typescript const result = await api.runtime.hooks.dispatchHookAgentTurn({ diff --git a/src/gateway/server/hooks.early-failure.test.ts b/src/gateway/server/hooks.early-failure.test.ts index c3ef494c3e64..2b6ef02b1710 100644 --- a/src/gateway/server/hooks.early-failure.test.ts +++ b/src/gateway/server/hooks.early-failure.test.ts @@ -5,6 +5,7 @@ import type { AcpRuntime, AcpRuntimeTurnInput } from "@openclaw/acp-core/runtime import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; import { createDeferred } from "../../../test/helpers/promise.js"; import { consumeAcpTurnStream } from "../../acp/control-plane/manager.turn-stream.js"; +import { DEFAULT_CRON_MAX_CONCURRENT_RUNS } from "../../config/cron-limits.js"; import type { HookMappingConfig } from "../../config/types.hooks.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import type { RunCronAgentTurnResult } from "../../cron/isolated-agent/run.types.js"; @@ -19,12 +20,15 @@ import { import { parseLogLine } from "../../logging/parse-log-line.js"; import { loggingState } from "../../logging/state.js"; import { createSubsystemLogger } from "../../logging/subsystem.js"; +import { enqueueCommandInLane, getCommandLaneSnapshot } from "../../process/command-queue.js"; +import { resetCommandQueueStateForTest } from "../../process/command-queue.test-support.js"; import { getActiveGatewayRootWorkCount, resetGatewayWorkAdmission, } from "../../process/gateway-work-admission.js"; import { CommandLane } from "../../process/lanes.js"; import { resolveHooksConfig } from "../hooks.js"; +import { applyGatewayLaneConcurrency, resolveGatewayLaneConcurrency } from "../server-lanes.js"; const mocks = vi.hoisted(() => ({ enqueueSystemEvent: vi.fn(), @@ -68,6 +72,18 @@ function createPluginHookDispatcher(options: { admissionTimeoutMs?: number } = { return { dispatcher, logHooks }; } +function queueHookRunner(onStart = vi.fn()) { + mocks.runCronIsolatedAgentTurn.mockImplementationOnce( + async (params: { lane: string; onExecutionStarted?: () => void }) => + await enqueueCommandInLane(params.lane, async () => { + params.onExecutionStarted?.(); + onStart(); + return { status: "ok", summary: "done" }; + }), + ); + return onStart; +} + function createConfig(global: boolean): OpenClawConfig { return { agents: { entries: { main: { default: true }, hooks: {} } }, @@ -148,6 +164,8 @@ describe("gateway hook early-failure recovery", () => { beforeEach(() => { resetGatewayWorkAdmission(); + resetCommandQueueStateForTest(); + applyGatewayLaneConcurrency(resolveGatewayLaneConcurrency({})); vi.clearAllMocks(); }); @@ -156,6 +174,7 @@ describe("gateway hook early-failure recovery", () => { resetGatewayWorkAdmission(); loggingState.rawConsole = null; resetLogger(); + resetCommandQueueStateForTest(); }); afterAll(async () => { @@ -589,65 +608,69 @@ describe("gateway hook early-failure recovery", () => { }, ); - it("contains plugin email turns without enabling the HTTP hook surface", async () => { - const config: OpenClawConfig = { - agents: { entries: { main: { default: true }, hooks: {} } }, - hooks: { - allowedAgentIds: ["main"], - allowedSessionKeyPrefixes: ["hook:http:"], - }, - }; - mocks.getRuntimeConfig.mockReturnValue(config); - mocks.runCronIsolatedAgentTurn.mockImplementationOnce( - async (params: { onExecutionStarted?: () => void }) => { - params.onExecutionStarted?.(); - return { status: "ok", summary: "done" }; - }, - ); - const { dispatcher, logHooks } = createPluginHookDispatcher(); - const unsafePluginTurn = { - ...pluginHookTurn, - allowUnsafeExternalContent: true, - sessionMode: "persistent", - }; - - const result = await dispatcher.dispatchHookAgentTurn(unsafePluginTurn, "imap"); - - expect(result).toEqual({ ok: true, runId: expect.any(String) }); - expect(mocks.runCronIsolatedAgentTurn).toHaveBeenCalledWith( - expect.objectContaining({ - agentId: "hooks", - sessionKey: pluginHookTurn.sessionKey, - lane: CommandLane.HookDispatch, - job: expect.objectContaining({ - name: "IMAP fastmail", - agentId: "hooks", - sessionTarget: "isolated", - payload: expect.objectContaining({ - kind: "agentTurn", - message: pluginHookTurn.message, - externalContentSource: "email", - allowUnsafeExternalContent: undefined, - }), - delivery: { mode: "none" }, - }), - executionIdentity: { - ingress: { - kind: "webhook", - boundary: "gateway.hooks.plugin", - state: "present", - rawSourceRef: "imap:IMAP fastmail", - }, + it.each([undefined, false, true])( + "contains plugin email turns with HTTP hooks enabled=%s", + async (enabled) => { + const config: OpenClawConfig = { + agents: { entries: { main: { default: true }, hooks: {} } }, + hooks: { + enabled, + allowedAgentIds: ["main"], + allowedSessionKeyPrefixes: ["hook:http:"], }, - }), - ); - await vi.waitFor(() => - expect(logHooks.info).toHaveBeenCalledWith( - expect.stringMatching(/^hook agent run completed /), - expect.objectContaining({ name: "IMAP fastmail" }), - ), - ); - }); + }; + mocks.getRuntimeConfig.mockReturnValue(config); + applyGatewayLaneConcurrency(resolveGatewayLaneConcurrency(config)); + const onStart = queueHookRunner(); + const { dispatcher, logHooks } = createPluginHookDispatcher({ admissionTimeoutMs: 100 }); + const unsafePluginTurn = { + ...pluginHookTurn, + allowUnsafeExternalContent: true, + sessionMode: "persistent", + }; + + const result = await dispatcher.dispatchHookAgentTurn(unsafePluginTurn, "imap"); + // Drain a closed lane after the regression fails; expired admission still fences execution. + applyGatewayLaneConcurrency(resolveGatewayLaneConcurrency(createConfig(false))); + await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(0)); + + expect(result).toEqual({ ok: true, runId: expect.any(String) }); + expect(onStart).toHaveBeenCalledOnce(); + expect(mocks.runCronIsolatedAgentTurn).toHaveBeenCalledWith( + expect.objectContaining({ + agentId: "hooks", + sessionKey: pluginHookTurn.sessionKey, + lane: CommandLane.CronNested, + job: expect.objectContaining({ + name: "IMAP fastmail", + agentId: "hooks", + sessionTarget: "isolated", + payload: expect.objectContaining({ + kind: "agentTurn", + message: pluginHookTurn.message, + externalContentSource: "email", + allowUnsafeExternalContent: undefined, + }), + delivery: { mode: "none" }, + }), + executionIdentity: { + ingress: { + kind: "webhook", + boundary: "gateway.hooks.plugin", + state: "present", + rawSourceRef: "imap:IMAP fastmail", + }, + }, + }), + ); + await vi.waitFor(() => + expect(logHooks.info).toHaveBeenCalledWith( + expect.stringMatching(/^hook agent run completed /), + expect.objectContaining({ name: "IMAP fastmail" }), + ), + ); + }, + ); it("announces successful plugin hook turns through the existing heartbeat path", async () => { mocks.getRuntimeConfig.mockReturnValue(createConfig(false)); @@ -757,25 +780,55 @@ describe("gateway hook early-failure recovery", () => { }); }); - it("preserves plugin hook admission timeout and fences late execution", async () => { - const releasePreparation = createDeferred(); - mocks.getRuntimeConfig.mockReturnValue(createConfig(false)); - mocks.runCronIsolatedAgentTurn.mockImplementationOnce(async () => { - await releasePreparation.promise; - return { status: "ok", summary: "done" }; - }); - const { dispatcher } = createPluginHookDispatcher({ admissionTimeoutMs: 10 }); - - try { - await expect(dispatcher.dispatchHookAgentTurn(pluginHookTurn, "imap")).resolves.toEqual({ - ok: false, - reason: "hook agent run did not start before admission timeout", + it.each([false, true])( + "bounds plugin work by cron capacity and fences expired admission=%s", + async (expire) => { + const config = { agents: createConfig(false).agents }; + mocks.getRuntimeConfig.mockReturnValue(config); + applyGatewayLaneConcurrency(resolveGatewayLaneConcurrency(config)); + const releaseCron = createDeferred(); + const cronRuns = Array.from({ length: DEFAULT_CRON_MAX_CONCURRENT_RUNS }, () => + enqueueCommandInLane(CommandLane.CronNested, async () => await releaseCron.promise), + ); + const onStart = queueHookRunner(); + const { dispatcher } = createPluginHookDispatcher({ + admissionTimeoutMs: expire ? 100 : 5_000, }); - } finally { - releasePreparation.resolve(); - } - await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(0)); - }); + const admission = dispatcher.dispatchHookAgentTurn(pluginHookTurn, "imap"); + + try { + await vi.waitFor(() => + expect(getCommandLaneSnapshot(CommandLane.CronNested)).toMatchObject({ + activeCount: DEFAULT_CRON_MAX_CONCURRENT_RUNS, + queuedCount: 1, + }), + ); + expect(onStart).not.toHaveBeenCalled(); + if (!expire) { + releaseCron.resolve(); + } + await expect(admission).resolves.toEqual( + expire + ? { + ok: false, + reason: "hook agent run did not start before admission timeout", + } + : { ok: true, runId: expect.any(String) }, + ); + } finally { + releaseCron.resolve(); + applyGatewayLaneConcurrency(resolveGatewayLaneConcurrency(createConfig(false))); + await Promise.all(cronRuns); + await admission; + await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(0)); + } + expect(onStart).toHaveBeenCalledTimes(expire ? 0 : 1); + expect(getCommandLaneSnapshot(CommandLane.CronNested)).toMatchObject({ + activeCount: 0, + queuedCount: 0, + }); + }, + ); it("serializes HTTP and plugin turns together while replaying plugin idempotency keys", async () => { const releaseHttpRun = createDeferred(); diff --git a/src/gateway/server/hooks.ts b/src/gateway/server/hooks.ts index 7a0b1ae44eef..a23208f7e54a 100644 --- a/src/gateway/server/hooks.ts +++ b/src/gateway/server/hooks.ts @@ -518,11 +518,9 @@ export function createGatewayHookDispatcher(params: { // Isolated runs derive their lifecycle key from random jobId (or an // already-stable cron: key), so accepted agentId closes reload drift. agentId, - // Hook agent runs get their own lane rather than sharing - // `cron-nested` with cron inner work, so a saturated cron budget - // cannot starve them. Aggregate capacity stays bounded by the lane - // group that owns both lanes. - lane: CommandLane.HookDispatch, + // Only HTTP hooks own the opt-in reserved lane. Trusted plugin + // triggers share cron capacity even when the HTTP surface is off. + lane: pluginId ? CommandLane.CronNested : CommandLane.HookDispatch, executionIdentity: { ingress: pluginId ? {