From ad505a7b55296f66596082966edf4bb323567f68 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 23 Jul 2026 06:26:31 -0700 Subject: [PATCH] fix(swarm): keep collector results reliable through races and restarts (#112989) * fix(swarm): harden collector lifecycle and dashboards * fix(swarm): initialize collector completion state * test(swarm): satisfy cross-environment type checks * test(codex): allow direct request handler calls * style(ui): avoid Swarm widget shadowing * test(swarm): keep internal helpers private * refactor(ui): own Swarm roster helpers in runtime --- docs/plugins/codex-harness-reference.md | 16 +- docs/plugins/codex-harness-runtime.md | 11 +- docs/plugins/codex-harness.md | 17 +- docs/tools/swarm.md | 42 ++- .../codex/src/app-server/client.test.ts | 7 +- extensions/codex/src/app-server/client.ts | 17 +- .../app-server/dynamic-tool-execution.test.ts | 21 +- .../src/app-server/dynamic-tool-execution.ts | 8 +- .../app-server/run-attempt-server-requests.ts | 8 +- .../codex/src/app-server/turn-router.test.ts | 27 ++ .../codex/src/app-server/turn-router.ts | 70 +++-- src/agents/subagent-registry-helpers.test.ts | 122 ++++++++ src/agents/subagent-registry-helpers.ts | 54 +++- .../subagent-registry-lifecycle.test.ts | 60 ++++ src/agents/subagent-registry-lifecycle.ts | 24 +- src/agents/subagent-registry-run-manager.ts | 105 +++++-- src/agents/subagent-registry.test.ts | 125 +++++++- src/agents/subagent-registry.ts | 77 ++--- src/agents/subagent-registry.types.ts | 1 + src/agents/subagent-run-timeout.test.ts | 27 ++ src/agents/subagent-run-timeout.ts | 8 +- .../subagent-spawn.in-process-gateway.test.ts | 182 +++++++++++- src/agents/subagent-spawn.ts | 33 ++- src/agents/swarm-collector.ts | 28 +- .../tools/structured-output-tool.test.ts | 21 ++ src/agents/tools/structured-output-tool.ts | 18 +- src/process/gateway-work-admission.ts | 6 +- .../app-sidebar-child-session-data.ts | 43 +-- ui/src/lib/board/builtin-dashboard.ts | 8 +- ui/src/lib/board/swarm-dashboard-roster.ts | 40 +++ ui/src/lib/board/swarm-dashboard.test.ts | 267 ++++++++++++++++++ ui/src/lib/board/swarm-dashboard.ts | 138 +++++++++ ui/src/lib/board/widgets/swarm.test.ts | 38 +++ ui/src/lib/board/widgets/swarm.ts | 25 +- ui/src/lib/sessions/child-session-data.ts | 60 ++++ ui/src/lib/sessions/index.ts | 5 +- .../swarm-activity.integration.test.ts | 5 + ui/src/lib/sessions/swarm-activity.ts | 20 +- ui/src/pages/chat/chat-pane.ts | 81 ++++-- ui/src/styles/board.css | 6 + 40 files changed, 1638 insertions(+), 233 deletions(-) create mode 100644 ui/src/lib/board/swarm-dashboard-roster.ts create mode 100644 ui/src/lib/board/swarm-dashboard.test.ts create mode 100644 ui/src/lib/sessions/child-session-data.ts diff --git a/docs/plugins/codex-harness-reference.md b/docs/plugins/codex-harness-reference.md index d69527f68181..2aafc1047e44 100644 --- a/docs/plugins/codex-harness-reference.md +++ b/docs/plugins/codex-harness-reference.md @@ -507,6 +507,11 @@ toward native `spawn_agent` for Codex-native subagent work, while Message-tool-only source replies also stay direct, since that is a turn-control contract. +Codex Code Mode projects generic OpenClaw dynamic-tool results as text. Parse a +JSON result before reading fields. Nested dynamic calls are serialized by the +Codex runtime, so `Promise.all` does not submit them concurrently; use a +bounded sequential launch loop when starting collector children. + Tools marked `catalogMode: "direct-only"`, including the OpenClaw `computer` tool, are grouped under `openclaw_direct`. OpenClaw adds that namespace to Codex's `code_mode.direct_only_tool_namespaces` list without replacing @@ -534,15 +539,16 @@ first available timeout in this order: converted to milliseconds, or the 60 second media default. For image understanding, this applies to the request itself and is not reduced by earlier preparation work. -- For the `message` tool, a fixed 120 second default. +- For the `message` tool, a fixed 600 second outer budget that covers Gateway delivery and bounded same-key reconciliation. - The 90 second dynamic-tool default. This watchdog is the outer dynamic `item/tool/call` budget. Provider-specific request timeouts run inside that call and keep their own timeout semantics. -Dynamic tool budgets are capped at 600000 ms. On timeout, OpenClaw aborts the -tool signal where supported and returns a failed dynamic-tool response to -Codex so the turn can continue instead of leaving the session in -`processing`. +Dynamic tool budgets are capped at 600000 ms. `agents_wait` adds 30000 ms of +outer completion grace, and the app-server client allows 660000 ms so that +structured wait result can reach Codex. On timeout, OpenClaw aborts the tool +signal where supported and returns a failed dynamic-tool response to Codex so +the turn can continue instead of leaving the session in `processing`. After Codex accepts a turn, and after OpenClaw responds to a turn-scoped app-server request, the harness expects Codex to make current-turn progress diff --git a/docs/plugins/codex-harness-runtime.md b/docs/plugins/codex-harness-runtime.md index 68805f69514e..8ddc0f5162c3 100644 --- a/docs/plugins/codex-harness-runtime.md +++ b/docs/plugins/codex-harness-runtime.md @@ -149,10 +149,13 @@ hooks such as `SessionStart` and `UserPromptSubmit` remain Codex-level controls; they are not exposed as OpenClaw plugin hooks in the v1 contract. For OpenClaw dynamic tools, OpenClaw executes the tool after Codex asks for -the call, so plugin and middleware behavior runs in the harness adapter. For -Codex-native tools, Codex owns the canonical tool record; OpenClaw can mirror -selected events but cannot rewrite the native thread unless Codex exposes that -through app-server or native hook callbacks. +the call, so plugin and middleware behavior runs in the harness adapter. Codex +Code Mode receives generic dynamic results as text and serializes nested +dynamic calls; callers must parse JSON-looking results and cannot rely on +`Promise.all` for concurrent submission. For Codex-native tools, Codex owns the +canonical tool record; OpenClaw can mirror selected events but cannot rewrite +the native thread unless Codex exposes that through app-server or native hook +callbacks. Codex app-server report-mode `PreToolUse` events defer plugin approval to the matching app-server approval. If an OpenClaw `before_tool_call` hook returns diff --git a/docs/plugins/codex-harness.md b/docs/plugins/codex-harness.md index d661f1c3cde2..177ca3c1764a 100644 --- a/docs/plugins/codex-harness.md +++ b/docs/plugins/codex-harness.md @@ -789,11 +789,18 @@ Persistent effective search-policy changes rotate the bound Codex thread before the next turn; transient per-turn restrictions use a temporary restricted thread and preserve the existing binding for later resume. -`sessions_yield` and message-tool-only source replies stay direct because -those are turn-control contracts. `sessions_spawn` stays searchable so -Codex's native `spawn_agent` remains the primary Codex subagent surface, -while explicit OpenClaw or ACP delegation is still available through the -`openclaw` dynamic tool namespace. Heartbeat collaboration instructions +`sessions_yield`, `sessions_spawn`, and message-tool-only source replies stay +direct because they are turn-control or delegation contracts. Guidance still +prefers Codex's native `spawn_agent` as the primary Codex subagent surface, +while explicit OpenClaw or ACP delegation remains directly callable through +`sessions_spawn`. In Codex Code Mode, generic OpenClaw +dynamic-tool results are JSON text rather than JavaScript objects, so parse +JSON-looking results before reading fields. Codex also serializes nested +dynamic calls; submit several `sessions_spawn` calls in a bounded loop rather +than expecting `Promise.all` to launch them concurrently. Already-accepted +children can still overlap while later calls are submitted. See +[Swarm](/tools/swarm#use-swarm-from-other-harnesses) for a complete pattern. +Heartbeat collaboration instructions tell Codex to search for `heartbeat_respond` before ending a heartbeat turn when the tool is not already loaded. diff --git a/docs/tools/swarm.md b/docs/tools/swarm.md index 2436576c8d05..d6bce3b9e5dd 100644 --- a/docs/tools/swarm.md +++ b/docs/tools/swarm.md @@ -172,6 +172,12 @@ return await agents.run( `maxConcurrent` children for the group and queues the rest in submission order. +Code Mode separately bounds concurrent guest bridge calls with +`tools.codeMode.maxPendingToolCalls` (default `16`, maximum `128`). For very +large groups, launch bounded batches below that limit and leave headroom for +`phase()`, `log()`, and child wait transitions. `maxConcurrent` limits running +children; it does not raise the guest bridge-call limit. + ### Loop on a decision gate Use a bounded `while` loop when each pass decides whether another pass is @@ -328,29 +334,39 @@ calls. Codex Code Mode automatically exposes eligible dynamic OpenClaw tools under `tools.*`. It does not use OpenClaw's QuickJS guest API or require `tools.codeMode`, but `tools.swarm` must still be enabled. Codex harness -`agents_wait` calls support the full 600-second timeout. Use this pattern: +`agents_wait` calls support the full 600-second timeout. + +With the currently supported Codex runtime, dynamic OpenClaw tool results reach +Code Mode as JSON text. Parse each result before reading fields. Codex also +serializes dynamic tool calls, so `Promise.all` does not submit several +`sessions_spawn` calls concurrently. Launch collectors in a bounded loop; +already-accepted children can still run while later launches are submitted. ```javascript +function parseToolResult(value) { + if (typeof value !== "string") return value; + return JSON.parse(value); +} + const tasks = [ "Check the authentication path.", "Check the storage path.", "Check the recovery path.", ]; +const launches = []; -const launches = await Promise.all( - tasks.map((task, index) => - tools.sessions_spawn({ +for (const [index, task] of tasks.entries()) { + const launch = parseToolResult( + await tools.sessions_spawn({ task, collect: true, label: `review-${index + 1}`, }), - ), -); - -for (const launch of launches) { + ); if (launch.status !== "accepted") { throw new Error(launch.error ?? "Collector spawn was not accepted."); } + launches.push(launch); } const pending = new Set(launches.map((launch) => launch.runId)); @@ -358,10 +374,12 @@ const completed = []; while (pending.size > 0) { const ids = [...pending].slice(0, 1000); - const batch = await tools.agents_wait({ - ids, - timeoutSeconds: 30, - }); + const batch = parseToolResult( + await tools.agents_wait({ + ids, + timeoutSeconds: 30, + }), + ); // Rotate this bounded window behind ids that have not been checked yet. for (const runId of ids) { diff --git a/extensions/codex/src/app-server/client.test.ts b/extensions/codex/src/app-server/client.test.ts index bc47dd1eab72..80ebb87ae853 100644 --- a/extensions/codex/src/app-server/client.test.ts +++ b/extensions/codex/src/app-server/client.test.ts @@ -10,7 +10,7 @@ import { resetSharedCodexAppServerClientForTests } from "./shared-client.js"; import { createClientHarness } from "./test-support.js"; import { MAX_CODEX_APP_SERVER_VERSION, MIN_CODEX_APP_SERVER_VERSION } from "./version.js"; -const CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS = 600_000; +const CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS = 660_000; describe("CodexAppServerClient", () => { const clients: CodexAppServerClient[] = []; @@ -585,8 +585,10 @@ describe("CodexAppServerClient", () => { const warn = vi.spyOn(embeddedAgentLog, "warn").mockImplementation(() => undefined); const harness = createClientHarness(); clients.push(harness.client); - harness.client.addRequestHandler((request) => { + let requestSignal: AbortSignal | undefined; + harness.client.addRequestHandler((request, signal) => { if (request.method === "item/tool/call") { + requestSignal = signal; return new Promise(() => {}); } return undefined; @@ -608,6 +610,7 @@ describe("CodexAppServerClient", () => { ], }, }); + expect(requestSignal?.aborted).toBe(true); expect(warn).toHaveBeenCalledWith("codex app-server server request timed out", { id: "srv-timeout", method: "item/tool/call", diff --git a/extensions/codex/src/app-server/client.ts b/extensions/codex/src/app-server/client.ts index b7787a139c6e..8f44396ca6e3 100644 --- a/extensions/codex/src/app-server/client.ts +++ b/extensions/codex/src/app-server/client.ts @@ -34,7 +34,9 @@ import { MAX_CODEX_APP_SERVER_VERSION, MIN_CODEX_APP_SERVER_VERSION } from "./ve const CODEX_APP_SERVER_PARSE_LOG_MAX = 500; const CODEX_APP_SERVER_PARSE_BUFFER_MAX = 8 * 1024 * 1024; const CODEX_APP_SERVER_PARSE_BUFFER_MAX_LINES = 1_000; -const CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS = 600_000; +// agents_wait can use a 600s inner budget plus 30s handler grace. Keep the +// app-server request guard outside that window so Codex receives the tool result. +const CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS = 660_000; const CODEX_APP_SERVER_STDERR_TAIL_MAX = 2_000; const CODEX_APP_SERVER_OVERLOADED_ERROR_CODE = -32_001; const CODEX_APP_SERVER_OVERLOAD_MAX_RETRIES = 3; @@ -204,6 +206,7 @@ export function isCodexAppServerConnectionClosedError(error: unknown): boolean { type CodexServerRequestHandler = ( request: Required> & { params?: JsonValue }, + signal?: AbortSignal, ) => Promise | JsonValue | undefined; /** Notification handler registered on a Codex app-server client. */ @@ -870,15 +873,16 @@ export class CodexAppServerClient { private async runServerRequestHandlers( request: Required> & { params?: JsonValue }, ): Promise { + const controller = new AbortController(); const timeoutResponse = timeoutServerRequestResponse(request); if (!timeoutResponse) { - return await this.runServerRequestHandlersWithoutTimeout(request); + return await this.runServerRequestHandlersWithoutTimeout(request, controller.signal); } let timeout: ReturnType | undefined; try { return await Promise.race([ - this.runServerRequestHandlersWithoutTimeout(request), + this.runServerRequestHandlersWithoutTimeout(request, controller.signal), new Promise((resolve) => { timeout = setTimeout(() => { embeddedAgentLog.warn("codex app-server server request timed out", { @@ -886,6 +890,7 @@ export class CodexAppServerClient { method: request.method, timeoutMs: CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS, }); + controller.abort(new Error("codex app-server server request timed out")); resolve(timeoutResponse); }, CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS); timeout.unref?.(); @@ -900,9 +905,13 @@ export class CodexAppServerClient { private async runServerRequestHandlersWithoutTimeout( request: Required> & { params?: JsonValue }, + signal: AbortSignal, ): Promise { for (const handler of this.requestHandlers) { - const result = await handler(request); + if (signal.aborted) { + return undefined; + } + const result = await handler(request, signal); if (result !== undefined) { return result; } diff --git a/extensions/codex/src/app-server/dynamic-tool-execution.test.ts b/extensions/codex/src/app-server/dynamic-tool-execution.test.ts index d26803007a85..22c2a6443674 100644 --- a/extensions/codex/src/app-server/dynamic-tool-execution.test.ts +++ b/extensions/codex/src/app-server/dynamic-tool-execution.test.ts @@ -19,6 +19,7 @@ const CODEX_DYNAMIC_TOOL_TIMEOUT_MS = 90_000; const CODEX_DYNAMIC_TOOL_MAX_TIMEOUT_MS = 600_000; const CODEX_DYNAMIC_IMAGE_TOOL_TIMEOUT_MS = 60_000; const CODEX_DYNAMIC_MESSAGE_TOOL_TIMEOUT_MS = CODEX_DYNAMIC_TOOL_MAX_TIMEOUT_MS; +const CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS = 660_000; describe("dynamic tool execution helpers", () => { afterEach(() => { @@ -138,9 +139,11 @@ describe("dynamic tool execution helpers", () => { config: { agents: { defaults: { - imageGenerationModel: { - primary: "openai/gpt-image-1", - timeoutMs: 180_000, + mediaModels: { + image: { + primary: "openai/gpt-image-1", + timeoutMs: 180_000, + }, }, }, }, @@ -361,12 +364,12 @@ describe("dynamic tool execution helpers", () => { config: undefined, }), ).toBe(150_000); - expect( - resolveDynamicToolCallTimeoutMs({ - call: { ...call, arguments: { ids: ["run-1"], timeoutSeconds: 600 } }, - config: undefined, - }), - ).toBe(630_000); + const fullWaitTimeoutMs = resolveDynamicToolCallTimeoutMs({ + call: { ...call, arguments: { ids: ["run-1"], timeoutSeconds: 600 } }, + config: undefined, + }); + expect(fullWaitTimeoutMs).toBe(630_000); + expect(CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS).toBeGreaterThan(fullWaitTimeoutMs); }); it("returns a failed dynamic tool response when an app-server tool call exceeds the deadline", async () => { diff --git a/extensions/codex/src/app-server/dynamic-tool-execution.ts b/extensions/codex/src/app-server/dynamic-tool-execution.ts index dc3dd0d93a86..f6de4cd840f0 100644 --- a/extensions/codex/src/app-server/dynamic-tool-execution.ts +++ b/extensions/codex/src/app-server/dynamic-tool-execution.ts @@ -47,7 +47,7 @@ const CODEX_DYNAMIC_COMPUTER_COMPLETION_GRACE_MS = 30_000; /** Timeout for image-understanding style dynamic tool calls. */ const CODEX_DYNAMIC_IMAGE_TOOL_TIMEOUT_MS = 60_000; /** Timeout for message-delivery dynamic tool calls. */ -const CODEX_DYNAMIC_MESSAGE_TOOL_TIMEOUT_MS = CODEX_DYNAMIC_TOOL_MAX_TIMEOUT_MS; +const CODEX_DYNAMIC_MESSAGE_TOOL_TIMEOUT_MS = 600_000; /** Outer default for collector waits: full swarm budget plus completion grace. */ const CODEX_DYNAMIC_AGENTS_WAIT_TOOL_TIMEOUT_MS = CODEX_DYNAMIC_TOOL_MAX_TIMEOUT_MS + CODEX_DYNAMIC_TOOL_TIMEOUT_SECONDS_GRACE_MS; @@ -554,12 +554,12 @@ function readConfiguredDynamicToolTimeoutMs( config: EmbeddedRunAttemptParams["config"], ): number | undefined { if (toolName === "image_generate") { - const imageGenerationModel = config?.agents?.defaults?.imageGenerationModel; - if (!imageGenerationModel || typeof imageGenerationModel !== "object") { + const imageModel = config?.agents?.defaults?.mediaModels?.image; + if (!imageModel || typeof imageModel !== "object") { return CODEX_DYNAMIC_IMAGE_GENERATION_TOOL_TIMEOUT_MS; } return ( - readPositiveFiniteTimeoutMs(imageGenerationModel.timeoutMs) ?? + readPositiveFiniteTimeoutMs(imageModel.timeoutMs) ?? CODEX_DYNAMIC_IMAGE_GENERATION_TOOL_TIMEOUT_MS ); } diff --git a/extensions/codex/src/app-server/run-attempt-server-requests.ts b/extensions/codex/src/app-server/run-attempt-server-requests.ts index 25816b83998f..5a4d52b2a39b 100644 --- a/extensions/codex/src/app-server/run-attempt-server-requests.ts +++ b/extensions/codex/src/app-server/run-attempt-server-requests.ts @@ -66,7 +66,9 @@ export function createCodexAttemptServerRequestController( const handleServerRequest = async ( request: CodexAppServerServerRequest, scope: CodexThreadRouteScope, + requestSignal: AbortSignal = new AbortController().signal, ) => { + const signal = AbortSignal.any([runAbortController.signal, requestSignal]); const turnId = turnIdRef.current; const projector = projectorRef.current; let armCompletionWatchOnResponse = false; @@ -96,7 +98,7 @@ export function createCodexAttemptServerRequestController( ...(computerUseConfig.enabled ? { computerUseMcpServerName: computerUseConfig.mcpServerName } : {}), - signal: runAbortController.signal, + signal, }); } if (request.method === "item/tool/requestUserInput") { @@ -123,7 +125,7 @@ export function createCodexAttemptServerRequestController( turnId, nativeHookRelay: resourceState.nativeHookRelay, autoApprove: shouldAutoApproveCodexAppServerApprovals(appServer), - signal: runAbortController.signal, + signal, onNativeToolFailureDisposition: (itemId, disposition) => projector?.recordNativeToolApprovalFailure(itemId, disposition), }); @@ -210,7 +212,7 @@ export function createCodexAttemptServerRequestController( handleDynamicToolCallWithTimeout({ call, toolBridge, - signal: runAbortController.signal, + signal, timeoutMs: dynamicToolTimeoutMs, toolMeta, toolCallOrdinal, diff --git a/extensions/codex/src/app-server/turn-router.test.ts b/extensions/codex/src/app-server/turn-router.test.ts index c847de8d8e53..c11e6803e7d0 100644 --- a/extensions/codex/src/app-server/turn-router.test.ts +++ b/extensions/codex/src/app-server/turn-router.test.ts @@ -5,6 +5,8 @@ import type { JsonValue } from "./protocol.js"; import { createClientHarness } from "./test-support.js"; import { getCodexAppServerTurnRouter, type CodexAppServerServerRequest } from "./turn-router.js"; +const CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS = 660_000; + type ClientHarness = ReturnType; type WireResponse = { @@ -46,6 +48,31 @@ describe("CodexAppServerTurnRouter", () => { expect(addCloseHandler).toHaveBeenCalledTimes(1); }); + it("does not dispatch a request that times out before route activation", async () => { + vi.useFakeTimers(); + vi.spyOn(embeddedAgentLog, "warn").mockImplementation(() => undefined); + const harness = createHarness(); + const requestHandler = vi.fn(() => ({ executed: true })); + const route = getCodexAppServerTurnRouter(harness.client).reserveThread({ + threadId: "thread-late", + }); + + harness.send({ + id: "request-late", + method: "item/tool/call", + params: { threadId: "thread-late", turnId: "turn-late", tool: "message" }, + }); + await vi.advanceTimersByTimeAsync(CODEX_DYNAMIC_TOOL_SERVER_REQUEST_TIMEOUT_MS); + expect(await waitForResponse(harness, "request-late")).toMatchObject({ + id: "request-late", + result: { success: false }, + }); + + await route.activate({ onRequest: requestHandler }); + + expect(requestHandler).not.toHaveBeenCalled(); + }); + it("routes concurrent traffic to the exact thread and turn", async () => { const harness = createHarness(); const router = getCodexAppServerTurnRouter(harness.client); diff --git a/extensions/codex/src/app-server/turn-router.ts b/extensions/codex/src/app-server/turn-router.ts index cf45679a30a6..6bc5822a22b4 100644 --- a/extensions/codex/src/app-server/turn-router.ts +++ b/extensions/codex/src/app-server/turn-router.ts @@ -26,6 +26,7 @@ export type CodexThreadRouteScope = { type CodexThreadRequestHandler = ( request: CodexAppServerServerRequest, scope: CodexThreadRouteScope, + signal: AbortSignal, ) => Promise | JsonValue | undefined; type CodexThreadNotificationHandler = ( notification: CodexServerNotification, @@ -127,7 +128,7 @@ class ClientTurnRouter implements CodexAppServerTurnRouter { constructor(client: CodexAppServerClient) { client.addNotificationHandler((notification) => this.routeNotification(notification)); - client.addRequestHandler((request) => this.routeRequest(request)); + client.addRequestHandler((request, signal) => this.routeRequest(request, signal)); client.addCloseHandler(() => this.dispose()); } @@ -353,8 +354,11 @@ class ClientTurnRouter implements CodexAppServerTurnRouter { return route.notificationTail; } - private async routeRequest(request: CodexAppServerServerRequest): Promise { - if (this.disposed) { + private async routeRequest( + request: CodexAppServerServerRequest, + signal: AbortSignal = new AbortController().signal, + ): Promise { + if (this.disposed || signal.aborted) { return undefined; } const scope = readScope(request.params); @@ -365,10 +369,10 @@ class ClientTurnRouter implements CodexAppServerTurnRouter { if (!route || route.released) { return undefined; } - if (!route.handlers) { - await route.activated.promise; + if (!route.handlers && !(await waitForPromiseOrAbort(route.activated.promise, signal))) { + return undefined; } - if (route.released || !route.handlers) { + if (signal.aborted || route.released || !route.handlers) { return undefined; } const handler = route.handlers.onRequest; @@ -378,8 +382,11 @@ class ClientTurnRouter implements CodexAppServerTurnRouter { // Open routes service a resumed native turn. Arming starts the handoff to a // new OpenClaw turn, whose requests must wait for its accepted turn id. while (route.gate === "armed") { - await route.binding?.promise; - if (route.released) { + const binding = route.binding?.promise; + if (!binding || !(await waitForPromiseOrAbort(binding, signal))) { + return undefined; + } + if (signal.aborted || route.released) { return undefined; } } @@ -391,18 +398,24 @@ class ClientTurnRouter implements CodexAppServerTurnRouter { return undefined; } } - await this.waitForNotifications(route); - if (route.released) { + if (!(await waitForPromiseOrAbort(this.waitForNotifications(route), signal))) { + return undefined; + } + if (signal.aborted || route.released) { return undefined; } try { - const result = await handler(request, { - threadId: scope.threadId, - ...(scope.turnId ? { turnId: scope.turnId } : {}), - }); - return route.released ? undefined : result; + const result = await handler( + request, + { + threadId: scope.threadId, + ...(scope.turnId ? { turnId: scope.turnId } : {}), + }, + signal, + ); + return signal.aborted || route.released ? undefined : result; } catch (error) { - if (route.released) { + if (signal.aborted || route.released) { return undefined; } throw error; @@ -591,6 +604,31 @@ function isCodexTerminalTurnNotification(notification: CodexServerNotification): ); } +async function waitForPromiseOrAbort( + promise: Promise, + signal: AbortSignal, +): Promise { + if (signal.aborted) { + return false; + } + let removeAbort: (() => void) | undefined; + try { + return await Promise.race([ + promise.then(() => true), + new Promise((resolve) => { + const onAbort = () => resolve(false); + signal.addEventListener("abort", onAbort, { once: true }); + removeAbort = () => signal.removeEventListener("abort", onAbort); + if (signal.aborted) { + onAbort(); + } + }), + ]); + } finally { + removeAbort?.(); + } +} + function deferred(): Deferred { let resolve!: () => void; const promise = new Promise((resolvePromise) => { diff --git a/src/agents/subagent-registry-helpers.test.ts b/src/agents/subagent-registry-helpers.test.ts index a1e08ca7ea2a..afac3d21b6c9 100644 --- a/src/agents/subagent-registry-helpers.test.ts +++ b/src/agents/subagent-registry-helpers.test.ts @@ -4,12 +4,16 @@ import { promises as fs } from "node:fs"; import { afterEach, describe, expect, it, vi } from "vitest"; import { defaultRuntime } from "../runtime.js"; import { + backfillCollectorArchiveAtMs, capFrozenResultText, logAnnounceGiveUp, + reconcileOrphanedRestoredRuns, reconcileOrphanedRun, + resolveSubagentArchiveAtMs, safeRemoveAttachmentsDir, } from "./subagent-registry-helpers.js"; import type { SubagentRunRecord } from "./subagent-registry.types.js"; +import { updateSwarmCollectorCompletion } from "./swarm-collector.js"; function createRunEntry(overrides: Partial = {}): SubagentRunRecord { return { @@ -36,6 +40,124 @@ describe("capFrozenResultText", () => { }); }); +describe("resolveSubagentArchiveAtMs", () => { + const cfg = { agents: { defaults: { subagents: { archiveAfterMinutes: 5 } } } }; + + it("defers collector retention until terminal completion", () => { + for (const cleanup of ["keep", "delete"] as const) { + expect( + resolveSubagentArchiveAtMs({ + cfg, + now: 1_000, + spawnMode: "run", + cleanup, + collect: true, + }), + ).toBeUndefined(); + } + }); + + it("starts collector retention when terminal completion is frozen", () => { + const entry = createRunEntry({ + collect: true, + endedAt: 2_000, + outcome: { status: "ok" }, + completion: { required: false, resultText: "done", capturedAt: 2_000 }, + }); + + expect(updateSwarmCollectorCompletion(entry, cfg)).toBe(true); + expect(entry.collectorCompletion).toEqual({ status: "done" }); + expect(entry.archiveAtMs).toBe(302_000); + }); + + it("starts retention when a delayed result first becomes waitable", () => { + vi.useFakeTimers(); + vi.setSystemTime(10_000); + const entry = createRunEntry({ + collect: true, + endedAt: 2_000, + outcome: { status: "ok" }, + completion: { required: false, resultText: "done" }, + }); + + expect(updateSwarmCollectorCompletion(entry, cfg)).toBe(true); + expect(entry.completion?.capturedAt).toBe(10_000); + expect(entry.archiveAtMs).toBe(310_000); + vi.useRealTimers(); + }); + + it("backfills legacy collectors from their terminal time", () => { + const entry = createRunEntry({ + collect: true, + endedAt: 2_000, + archiveAtMs: 10_000, + }); + + expect(backfillCollectorArchiveAtMs(entry, cfg)).toBe(true); + expect(entry.archiveAtMs).toBe(302_000); + expect(backfillCollectorArchiveAtMs(entry, cfg)).toBe(false); + }); + + it("clears stale deadlines from active, persistent, or retention-disabled collectors", () => { + const active = createRunEntry({ collect: true, archiveAtMs: 10_000 }); + expect(backfillCollectorArchiveAtMs(active, cfg)).toBe(true); + expect(active.archiveAtMs).toBeUndefined(); + + const persistent = createRunEntry({ + collect: true, + spawnMode: "session", + endedAt: 2_000, + archiveAtMs: 10_000, + }); + expect(backfillCollectorArchiveAtMs(persistent, cfg)).toBe(true); + expect(persistent.archiveAtMs).toBeUndefined(); + + const completed = createRunEntry({ collect: true, endedAt: 2_000, archiveAtMs: 10_000 }); + expect( + backfillCollectorArchiveAtMs(completed, { + agents: { defaults: { subagents: { archiveAfterMinutes: 0 } } }, + }), + ).toBe(true); + expect(completed.archiveAtMs).toBeUndefined(); + }); + + it("preserves ordinary keep and persistent session semantics", () => { + expect( + resolveSubagentArchiveAtMs({ + cfg, + now: 1_000, + spawnMode: "run", + cleanup: "keep", + }), + ).toBeUndefined(); + expect( + resolveSubagentArchiveAtMs({ + cfg, + now: 1_000, + spawnMode: "session", + cleanup: "delete", + collect: true, + }), + ).toBeUndefined(); + }); +}); + +describe("reconcileOrphanedRestoredRuns", () => { + it("keeps waitable collector tombstones after delete-mode sessions disappear", () => { + const entry = createRunEntry({ + collect: true, + cleanup: "delete", + endedAt: 2_000, + completion: { required: false, resultText: "done", capturedAt: 2_000 }, + collectorCompletion: { status: "done" }, + }); + const runs = new Map([[entry.runId, entry]]); + + expect(reconcileOrphanedRestoredRuns({ runs, resumedRuns: new Set() })).toBe(false); + expect(runs.get(entry.runId)).toBe(entry); + }); +}); + describe("safeRemoveAttachmentsDir", () => { it("reports non-ENOENT realpath failures instead of treating cleanup as complete", async () => { const realpathSpy = vi diff --git a/src/agents/subagent-registry-helpers.ts b/src/agents/subagent-registry-helpers.ts index 96da4ffa78fa..e211d190a592 100644 --- a/src/agents/subagent-registry-helpers.ts +++ b/src/agents/subagent-registry-helpers.ts @@ -316,6 +316,10 @@ export function reconcileOrphanedRestoredRuns(params: { const now = Date.now(); let changed = false; for (const [runId, entry] of params.runs.entries()) { + if (entry.collect && entry.collectorCompletion) { + // Waitable collector tombstones intentionally outlive delete-mode sessions. + continue; + } if (entry.requesterSettleWake) { // Requester-settle outbox rows can intentionally outlive delete-mode // child sessions. Restore replays the obligation before retiring them. @@ -351,7 +355,7 @@ export function reconcileOrphanedRestoredRuns(params: { } /** Resolves the completed subagent archive delay from config. */ -export function resolveArchiveAfterMs(cfg?: OpenClawConfig) { +function resolveArchiveAfterMs(cfg?: OpenClawConfig) { const config = cfg ?? getRuntimeConfig(); const minutes = config.agents?.defaults?.subagents?.archiveAfterMinutes ?? @@ -364,3 +368,51 @@ export function resolveArchiveAfterMs(cfg?: OpenClawConfig) { } return Math.max(1, Math.floor(minutes)) * 60_000; } + +/** Resolves the archive deadline for one newly registered run. */ +export function resolveSubagentArchiveAtMs(params: { + cfg?: OpenClawConfig; + now: number; + spawnMode: "run" | "session"; + cleanup: "keep" | "delete"; + collect?: boolean; +}): number | undefined { + if (params.spawnMode === "session" || params.collect || params.cleanup === "keep") { + return undefined; + } + const archiveAfterMs = resolveArchiveAfterMs(params.cfg); + return archiveAfterMs ? params.now + archiveAfterMs : undefined; +} + +/** Backfills the retention deadline added after collector groups first shipped. */ +export function backfillCollectorArchiveAtMs( + entry: SubagentRunRecord, + cfg?: OpenClawConfig, +): boolean { + if (!entry.collect) { + return false; + } + const endedAt = + typeof entry.endedAt === "number" && Number.isFinite(entry.endedAt) ? entry.endedAt : undefined; + const capturedAt = + endedAt === undefined && !entry.collectorCompletion + ? undefined + : typeof entry.completion?.capturedAt === "number" && + Number.isFinite(entry.completion.capturedAt) + ? entry.completion.capturedAt + : endedAt; + const archiveAfterMs = entry.spawnMode === "session" ? undefined : resolveArchiveAfterMs(cfg); + const expectedArchiveAt = + capturedAt !== undefined && archiveAfterMs !== undefined + ? capturedAt + archiveAfterMs + : undefined; + if (entry.archiveAtMs === expectedArchiveAt) { + return false; + } + if (expectedArchiveAt === undefined) { + delete entry.archiveAtMs; + } else { + entry.archiveAtMs = expectedArchiveAt; + } + return true; +} diff --git a/src/agents/subagent-registry-lifecycle.test.ts b/src/agents/subagent-registry-lifecycle.test.ts index a86876675d61..4aef89686edf 100644 --- a/src/agents/subagent-registry-lifecycle.test.ts +++ b/src/agents/subagent-registry-lifecycle.test.ts @@ -23,6 +23,7 @@ import { import { createSubagentRegistryLifecycleController } from "./subagent-registry-lifecycle.js"; import { markSubagentRunPausedAfterYield } from "./subagent-registry-run-manager.js"; import type { SubagentRunRecord } from "./subagent-registry.types.js"; +import { createStructuredOutputTool } from "./tools/structured-output-tool.js"; type LifecycleControllerParams = Parameters[0]; @@ -107,6 +108,7 @@ vi.mock("./subagent-registry-helpers.js", () => ({ MAX_ANNOUNCE_RETRY_COUNT: 3, MIN_ANNOUNCE_RETRY_DELAY_MS: 1_000, PROVISIONAL_KILL_RECONCILIATION_MS: 5 * 60_000, + backfillCollectorArchiveAtMs: () => false, capFrozenResultText: (text: string) => text.trim(), logAnnounceGiveUp: helperMocks.logAnnounceGiveUp, persistSubagentSessionTiming: helperMocks.persistSubagentSessionTiming, @@ -195,6 +197,7 @@ function createLifecycleController({ runs, resumedRuns: new Set(), subagentAnnounceTimeoutMs: 1_000, + getRuntimeConfig: () => ({}), persist: vi.fn(), persistOrThrow: vi.fn(), clearPendingLifecycleError: vi.fn(), @@ -2579,6 +2582,63 @@ describe("subagent registry lifecycle hardening", () => { expect(entry.collectorCompletion).toEqual({ status: "done" }); }); + it("treats accepted structured output as success for a tool-only collector turn", async () => { + const structured = { answer: "yes" }; + const entry = createRunEntry({ + expectsCompletionMessage: false, + collect: true, + outputSchema: { type: "object" }, + }); + const structuredOutput = createStructuredOutputTool({ + runId: entry.runId, + schema: { type: "object" }, + }); + await structuredOutput.execute("tool-call", { result: structured }); + const controller = createLifecycleController({ + entry, + captureSubagentCompletionReply: vi.fn(async () => ""), + }); + + await controller.completeSubagentRun({ + runId: entry.runId, + endedAt: 4_000, + outcome: { status: "error", error: "completed" }, + reason: SUBAGENT_ENDED_REASON_ERROR, + triggerCleanup: true, + }); + + await waitForLifecycleState(() => expect(entry.cleanupCompletedAt).toBeTypeOf("number")); + expect(entry.collectorCompletion).toEqual({ status: "done", structured }); + expect(entry.outcome).toMatchObject({ status: "ok" }); + expect(entry.execution).toMatchObject({ + status: "terminal", + outcome: expect.objectContaining({ status: "ok" }), + }); + expect(entry.endedReason).toBe(SUBAGENT_ENDED_REASON_COMPLETE); + }); + + it("preserves a real failure after structured output was accepted", async () => { + const structured = { answer: "yes" }; + const entry = createRunEntry({ + expectsCompletionMessage: false, + collect: true, + outputSchema: { type: "object" }, + structuredOutput: { structured, invalidAttempts: 0 }, + }); + const controller = createLifecycleController({ entry }); + + await controller.completeSubagentRun({ + runId: entry.runId, + endedAt: 4_000, + outcome: { status: "error", error: "provider failed after tool output" }, + reason: SUBAGENT_ENDED_REASON_ERROR, + triggerCleanup: true, + }); + + await waitForLifecycleState(() => expect(entry.cleanupCompletedAt).toBeTypeOf("number")); + expect(entry.collectorCompletion).toEqual({ status: "failed", structured }); + }); + it("marks a successful collector with invalid structured output failed", async () => { const entry = createRunEntry({ expectsCompletionMessage: false, diff --git a/src/agents/subagent-registry-lifecycle.ts b/src/agents/subagent-registry-lifecycle.ts index a5d32db2ca96..b775790c1b71 100644 --- a/src/agents/subagent-registry-lifecycle.ts +++ b/src/agents/subagent-registry-lifecycle.ts @@ -7,6 +7,7 @@ import { uniqueStrings } from "@openclaw/normalization-core/string-normalization import { isSilentReplyText, SILENT_REPLY_TOKEN } from "../auto-reply/tokens.js"; import type { cleanupBrowserSessionsForLifecycleEnd } from "../browser-lifecycle-cleanup.js"; import { formatSqliteSessionFileMarker } from "../config/sessions/sqlite-marker.js"; +import type { OpenClawConfig } from "../config/types.openclaw.js"; import type { callGateway as defaultCallGateway } from "../gateway/call.js"; import { formatErrorMessage, readErrorName } from "../infra/errors.js"; import { @@ -84,6 +85,7 @@ import { import { deleteSubagentSessionForCleanup } from "./subagent-session-cleanup.js"; import { updateSwarmCollectorCompletion } from "./swarm-collector.js"; import { releaseSwarmRun } from "./swarm-scheduler.js"; +import { peekSwarmStructuredOutput } from "./tools/structured-output-tool.js"; type CaptureSubagentCompletionReply = (typeof import("./subagent-announce.js"))["captureSubagentCompletionReply"]; @@ -174,6 +176,7 @@ export function createSubagentRegistryLifecycleController(params: { runs: Map; resumedRuns: Set; subagentAnnounceTimeoutMs: number; + getRuntimeConfig(): OpenClawConfig; persist(): void; persistOrThrow(): void; clearPendingLifecycleError(runId: string): void; @@ -1906,6 +1909,25 @@ export function createSubagentRegistryLifecycleController(params: { let endedAt = requestedEndedAt; let completionOutcome = shouldDrainExistingTerminal && entry.outcome ? entry.outcome : completeParams.outcome; + const liveStructuredOutput = entry.collect + ? (entry.structuredOutput ?? + peekSwarmStructuredOutput(entry.runId) ?? + (entry.swarmRunId ? peekSwarmStructuredOutput(entry.swarmRunId) : undefined)) + : undefined; + if (!entry.structuredOutput && liveStructuredOutput) { + entry.structuredOutput = liveStructuredOutput; + mutated = true; + } + if ( + liveStructuredOutput?.structured !== undefined && + completionOutcome.status === "error" && + completionOutcome.error === "completed" + ) { + // Tool-only collector turns use this runner sentinel after the result is + // durably recorded. Normalize before every task/session/hook projection. + completionOutcome = { status: "ok" }; + completionReason = SUBAGENT_ENDED_REASON_COMPLETE; + } const observedStartedAt = !shouldDrainExistingTerminal && typeof completeParams.startedAt === "number" && @@ -2085,7 +2107,7 @@ export function createSubagentRegistryLifecycleController(params: { mutated = true; } } - if (updateSwarmCollectorCompletion(entry)) { + if (updateSwarmCollectorCompletion(entry, params.getRuntimeConfig())) { mutated = true; } if (provisionalKillSnapshot) { diff --git a/src/agents/subagent-registry-run-manager.ts b/src/agents/subagent-registry-run-manager.ts index 3f76aff8f7dc..418edc46bf3e 100644 --- a/src/agents/subagent-registry-run-manager.ts +++ b/src/agents/subagent-registry-run-manager.ts @@ -43,7 +43,7 @@ import { } from "./subagent-registry-completion.js"; import { persistSubagentSessionTiming, - resolveArchiveAfterMs, + resolveSubagentArchiveAtMs, safeRemoveAttachmentsDir, } from "./subagent-registry-helpers.js"; import type { @@ -659,14 +659,14 @@ export function createSubagentRunManager(params: { source.childSessionKey, ); const cfg = params.getRuntimeConfig(); - const archiveAfterMs = resolveArchiveAfterMs(cfg); const spawnMode = source.spawnMode === "session" ? "session" : "run"; - const archiveAtMs = - spawnMode === "session" || source.cleanup === "keep" - ? undefined - : archiveAfterMs - ? now + archiveAfterMs - : undefined; + const archiveAtMs = resolveSubagentArchiveAtMs({ + cfg, + now, + spawnMode, + cleanup: source.cleanup, + collect: source.collect, + }); const runTimeoutSeconds = replaceParams.runTimeoutSeconds ?? source.runTimeoutSeconds ?? 0; const waitTimeoutMs = params.resolveSubagentWaitTimeoutMs(cfg, runTimeoutSeconds); const preserveFrozenResultFallback = replaceParams.preserveFrozenResultFallback === true; @@ -790,14 +790,14 @@ export function createSubagentRunManager(params: { const now = Date.now(); const generation = nextSubagentRunGeneration(params.runs.values(), childSessionKey); const cfg = params.getRuntimeConfig(); - const archiveAfterMs = resolveArchiveAfterMs(cfg); const spawnMode = registerParams.spawnMode === "session" ? "session" : "run"; - const archiveAtMs = - spawnMode === "session" || registerParams.cleanup === "keep" - ? undefined - : archiveAfterMs - ? now + archiveAfterMs - : undefined; + const archiveAtMs = resolveSubagentArchiveAtMs({ + cfg, + now, + spawnMode, + cleanup: registerParams.cleanup, + collect: registerParams.collect, + }); const runTimeoutSeconds = registerParams.runTimeoutSeconds ?? 0; const waitTimeoutMs = params.resolveSubagentWaitTimeoutMs(cfg, runTimeoutSeconds); const requesterOrigin = normalizeDeliveryContext(registerParams.requesterOrigin); @@ -922,11 +922,26 @@ export function createSubagentRunManager(params: { const entry = params.runs.get(key) ?? [...params.runs.values()].find((candidate) => candidate.swarmRunId === key); + const lifecycleStarted = + entry?.execution?.status === "running" && + typeof entry.execution.startedAt === "number" && + entry.swarmLaunchPending === true; + const provisionalTerminalBeforeAcceptance = + entry?.swarmLaunchPending === true && + typeof entry.endedAt === "number" && + entry.collectorCompletion === undefined; + if (provisionalTerminalBeforeAcceptance) { + // Cancellation won before Gateway acceptance. The caller must abort the + // newly accepted run before freezing completion or releasing the FIFO slot. + return false; + } + // Completion clears swarmLaunchPending, but queuedLaunch remains until the + // delayed acceptance response remaps the durable terminal row. + const terminalBeforeAcceptance = + entry?.collectorCompletion !== undefined && entry.queuedLaunch !== undefined; if ( !entry || - entry.execution?.status !== "queued" || - typeof entry.endedAt === "number" || - entry.collectorCompletion + (!terminalBeforeAcceptance && entry.execution?.status !== "queued" && !lifecycleStarted) ) { return false; } @@ -935,7 +950,7 @@ export function createSubagentRunManager(params: { if (conflicting && conflicting !== entry) { throw new Error(`collector gateway run id already exists: ${nextRunId}`); } - const startedAt = Date.now(); + const acceptedAt = Date.now(); const previousRunId = entry.runId; const previousStartedAt = entry.startedAt; const previousSessionStartedAt = entry.sessionStartedAt; @@ -951,9 +966,45 @@ export function createSubagentRunManager(params: { entry.runId = nextRunId; params.runs.set(nextRunId, entry); } - entry.startedAt = startedAt; - entry.sessionStartedAt ??= startedAt; - entry.execution = { ...entry.execution, status: "running", startedAt }; + if (terminalBeforeAcceptance) { + entry.swarmLaunchPending = false; + entry.queuedLaunch = undefined; + try { + params.persistOrThrow(); + return true; + } catch (error) { + if (previousRunId !== nextRunId) { + params.runs.delete(nextRunId); + entry.runId = previousRunId; + params.runs.set(previousRunId, entry); + } + entry.queuedLaunch = previousQueuedLaunch; + entry.swarmRunId = previousSwarmRunId; + entry.schedulerSlotId = previousSchedulerSlotId; + entry.swarmLaunchPending = previousSwarmLaunchPending; + throw error; + } + } + // Gateway acceptance only proves admission. Preserve a lifecycle start that + // raced ahead of this response; otherwise leave the run clock unset until + // preparation and lane dequeue emit the canonical start event. + const lifecycleStartedAt = + entry.execution?.status === "running" ? entry.execution.startedAt : undefined; + if (typeof lifecycleStartedAt === "number") { + entry.startedAt = lifecycleStartedAt; + entry.sessionStartedAt ??= lifecycleStartedAt; + entry.execution = { + ...entry.execution, + status: "running", + acceptedAt, + startedAt: lifecycleStartedAt, + }; + } else { + delete entry.startedAt; + delete entry.sessionStartedAt; + entry.execution = { ...entry.execution, status: "running", acceptedAt }; + delete entry.execution.startedAt; + } entry.swarmLaunchPending = false; entry.queuedLaunch = undefined; let persistedRunning = false; @@ -964,8 +1015,8 @@ export function createSubagentRunManager(params: { runId: entry.taskRunId ?? entry.runId, runtime: "subagent", sessionKey: entry.childSessionKey, - startedAt, - lastEventAt: startedAt, + startedAt: acceptedAt, + lastEventAt: acceptedAt, }); } catch (error) { if (previousRunId !== nextRunId) { @@ -1019,7 +1070,7 @@ export function createSubagentRunManager(params: { entry.queuedLaunch = undefined; entry.collectorLaunchCleanupPending = true; entry.completion = { required: false, resultText: error, capturedAt: endedAt }; - updateSwarmCollectorCompletion(entry); + updateSwarmCollectorCompletion(entry, params.getRuntimeConfig()); try { params.persistOrThrow(); } catch (persistError) { @@ -1080,7 +1131,7 @@ export function createSubagentRunManager(params: { resultText: entry.outcome?.status === "error" ? (entry.outcome.error ?? error) : error, capturedAt: entry.endedAt, }; - updateSwarmCollectorCompletion(entry); + updateSwarmCollectorCompletion(entry, params.getRuntimeConfig()); try { params.persistOrThrow(); } catch (persistError) { @@ -1238,7 +1289,7 @@ export function createSubagentRunManager(params: { supersededAt: existingKillReconciliation?.supersededAt, }; if (wasQueuedCollector && !collectorLaunchInFlight) { - updateSwarmCollectorCompletion(entry); + updateSwarmCollectorCompletion(entry, params.getRuntimeConfig()); } pendingTaskFinalizations.push({ entry, endedAt: taskEndedAt }); if (!entriesByChildSessionKey.has(entry.childSessionKey)) { diff --git a/src/agents/subagent-registry.test.ts b/src/agents/subagent-registry.test.ts index 38d7e556cecc..5e4438fc448c 100644 --- a/src/agents/subagent-registry.test.ts +++ b/src/agents/subagent-registry.test.ts @@ -17,6 +17,7 @@ import { getActiveGatewayRootWorkCount, markGatewayRestartDraining, resetGatewayWorkAdmission, + tryBeginGatewaySuspendAdmission, } from "../process/gateway-work-admission.js"; import { SUBAGENT_KILL_TASK_ERROR } from "../tasks/detached-task-runtime-contract.js"; import { @@ -449,7 +450,7 @@ describe("subagent registry seam flow", () => { runId: `run-collector-${suffix}`, childSessionKey: `agent:main:subagent:collector-${suffix}`, task: "retain lifetime group count", - cleanup: "delete", + cleanup: suffix === "one" ? "keep" : "delete", createdAt: now - 10_000, endedAt: now - 5_000, cleanupCompletedAt: now - 4_000, @@ -720,6 +721,8 @@ describe("subagent registry seam flow", () => { expect(mod.markSubagentRunTerminated({ runId, reason: "manual kill" })).toBe(1); expect(mod.getSubagentRunByRunId(runId)?.collectorCompletion).toBeUndefined(); + expect(mod.startQueuedSubagentRun(runId, "gateway-launch-kill")).toBe(false); + expect(mod.getSubagentRunByRunId("gateway-launch-kill")).toBeUndefined(); expect(mod.settleFailedQueuedSubagentLaunch(runId, "launch response lost")).toBe(true); expect(mod.getSubagentRunByRunId(runId)?.collectorCompletion).toMatchObject({ @@ -874,7 +877,7 @@ describe("subagent registry seam flow", () => { ); }); - it("rehydrates persisted collector FIFO queues after registry restore", async () => { + it("rehydrates persisted collector FIFO queues after admission reopens", async () => { const now = Date.now(); mocks.getRuntimeConfig.mockReturnValue({ tools: { swarm: { enabled: true, maxConcurrent: 1 } }, @@ -921,7 +924,15 @@ describe("subagent registry seam flow", () => { return request.method === "agent.wait" ? { status: "pending" } : {}; }); + const suspension = tryBeginGatewaySuspendAdmission(() => {}); + expect(suspension?.commit()).toBe(true); mod.initSubagentRegistry(); + await Promise.resolve(); + expect(mocks.callGateway.mock.calls.filter(([request]) => request.method === "agent")).toEqual( + [], + ); + + suspension?.release(); await waitForFast(() => { const agentCalls = mocks.callGateway.mock.calls.filter( ([request]) => request.method === "agent", @@ -937,14 +948,122 @@ describe("subagent registry seam flow", () => { }); }); expect(mod.getSubagentRunByRunId("run-queued-one")?.execution?.status).toBe("running"); - expect(mod.getSubagentRunByRunId("gateway-run-one")).toMatchObject({ + const acceptedRun = mod.getSubagentRunByRunId("gateway-run-one"); + expect(acceptedRun).toMatchObject({ runId: "gateway-run-one", swarmRunId: "run-queued-one", schedulerSlotId: "run-queued-one", + execution: { status: "running" }, }); + expect(acceptedRun).not.toHaveProperty("startedAt"); + expect(acceptedRun).not.toHaveProperty("sessionStartedAt"); + expect(acceptedRun?.execution).not.toHaveProperty("startedAt"); expect(mod.getSubagentRunByRunId("run-queued-two")?.execution?.status).toBe("queued"); }); + it("preserves a lifecycle start that arrives before collector acceptance returns", async () => { + const startedAt = 12_345; + mod.registerSubagentRun({ + runId: "run-start-race", + childSessionKey: "agent:main:subagent:start-race", + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + task: "start before acceptance", + cleanup: "keep", + collect: true, + groupId: "start-race", + queued: true, + expectsCompletionMessage: false, + }); + const lastOnAgentEventCall = mocks.onAgentEvent.mock.calls.at(-1) as unknown as + | [(event: AgentEventPayload) => void] + | undefined; + const lifecycleHandler = lastOnAgentEventCall?.[0]; + expect(lifecycleHandler).toBeTypeOf("function"); + + lifecycleHandler?.({ + runId: "run-start-race", + seq: 1, + stream: "lifecycle", + ts: startedAt, + data: { phase: "start", startedAt }, + }); + await waitForFast(() => + expect(mod.getSubagentRunByRunId("run-start-race")?.startedAt).toBe(startedAt), + ); + + expect(mod.startQueuedSubagentRun("run-start-race", "gateway-start-race")).toBe(true); + expect(mod.getSubagentRunByRunId("gateway-start-race")).toMatchObject({ + startedAt, + sessionStartedAt: startedAt, + execution: { status: "running", acceptedAt: expect.any(Number), startedAt }, + }); + }); + + it("remaps a collector that completed before its acceptance response", () => { + mod.addSubagentRunForTests({ + runId: "run-terminal-race", + childSessionKey: "agent:main:subagent:terminal-race", + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + task: "finish before acceptance", + cleanup: "keep", + collect: true, + swarmRunId: "run-terminal-race", + schedulerSlotId: "run-terminal-race", + swarmLaunchPending: false, + queuedLaunch: { + request: { sessionKey: "agent:main:subagent:terminal-race" }, + timeoutMs: 1_000, + schedulerGroupKey: "terminal-race", + maxConcurrent: 1, + }, + groupId: "terminal-race", + createdAt: 1_000, + endedAt: 2_000, + execution: { status: "terminal", endedAt: 2_000 }, + completion: { required: false, resultText: "done", capturedAt: 2_000 }, + collectorCompletion: { status: "done" }, + }); + + expect(mod.startQueuedSubagentRun("run-terminal-race", "gateway-terminal-race")).toBe(true); + const remapped = mod.getSubagentRunByRunId("gateway-terminal-race"); + expect(mod.getSubagentRunByRunId("run-terminal-race")).toBe(remapped); + expect(remapped).toMatchObject({ + runId: "gateway-terminal-race", + swarmRunId: "run-terminal-race", + collectorCompletion: { status: "done" }, + swarmLaunchPending: false, + }); + }); + + it("refuses to remap an unrelated terminal collector without a pending launch", () => { + mod.addSubagentRunForTests({ + runId: "run-terminal-stale", + childSessionKey: "agent:main:subagent:terminal-stale", + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + task: "stale acceptance callback", + cleanup: "keep", + collect: true, + swarmRunId: "run-terminal-stale", + schedulerSlotId: "run-terminal-stale", + groupId: "terminal-stale", + createdAt: 1_000, + endedAt: 2_000, + execution: { status: "terminal", endedAt: 2_000 }, + completion: { required: false, resultText: "done", capturedAt: 2_000 }, + collectorCompletion: { status: "done" }, + }); + + expect(mod.startQueuedSubagentRun("run-terminal-stale", "gateway-terminal-stale")).toBe(false); + expect(mod.getSubagentRunByRunId("run-terminal-stale")).toMatchObject({ + runId: "run-terminal-stale", + collectorCompletion: { status: "done" }, + }); + expect(mod.getSubagentRunByRunId("gateway-terminal-stale")).toBeUndefined(); + }); + it("holds a restored FIFO slot until an accepted collector is confirmed stopped", async () => { vi.useRealTimers(); const now = Date.now(); diff --git a/src/agents/subagent-registry.ts b/src/agents/subagent-registry.ts index c42dc93eeb83..66e3f1ccfe6a 100644 --- a/src/agents/subagent-registry.ts +++ b/src/agents/subagent-registry.ts @@ -16,6 +16,7 @@ import { getAgentRunContext, onAgentEvent } from "../infra/agent-events.js"; import { isFastTestRuntimeEnv } from "../infra/env.js"; import { createSubsystemLogger } from "../logging/subsystem.js"; import { + GatewayDrainingError, isGatewayRestartDraining, runWithGatewayIndependentRootWorkAdmission, } from "../process/gateway-work-admission.js"; @@ -67,6 +68,7 @@ import { } from "./subagent-registry-completion.js"; import { ANNOUNCE_EXPIRY_MS, + backfillCollectorArchiveAtMs, MAX_ANNOUNCE_RETRY_COUNT, PROVISIONAL_KILL_RECONCILIATION_MS, reconcileOrphanedRestoredRuns, @@ -865,6 +867,7 @@ const subagentLifecycleController = createSubagentRegistryLifecycleController({ runs: subagentRuns, resumedRuns, subagentAnnounceTimeoutMs: SUBAGENT_ANNOUNCE_TIMEOUT_MS, + getRuntimeConfig: () => subagentRegistryDeps.getRuntimeConfig(), persist: persistSubagentRuns, persistOrThrow: persistSubagentRunsOrThrow, clearPendingLifecycleError, @@ -1071,12 +1074,17 @@ function restoreSubagentRunsOnce() { if (restoredCount === 0) { return; } - if ( - reconcileOrphanedRestoredRuns({ - runs: subagentRuns, - resumedRuns, - }) - ) { + const cfg = subagentRegistryDeps.getRuntimeConfig(); + let restoredStateChanged = reconcileOrphanedRestoredRuns({ + runs: subagentRuns, + resumedRuns, + }); + for (const entry of subagentRuns.values()) { + if (backfillCollectorArchiveAtMs(entry, cfg)) { + restoredStateChanged = true; + } + } + if (restoredStateChanged) { persistSubagentRuns(); } const requesterTurns = new Map>(); @@ -1144,32 +1152,37 @@ function restoreSubagentRunsOnce() { .filter((candidate) => candidate.execution?.status === "running") .map((candidate) => candidate.schedulerSlotId ?? candidate.runId), start: async () => { - const response = await subagentRegistryDeps.callGateway({ - method: "agent", - params: applySubagentLaunchAuthorization(launch.request, launch.authorization), - // Restart replay must restore the trusted launch capability; otherwise - // the queued child silently falls back to its session/default route. - ...(launch.authorization ? { scopes: [ADMIN_SCOPE] } : {}), - timeoutMs: launch.timeoutMs, - }); - const gatewayRunId = readGatewayRunId(response) ?? runId; - try { - if (!startQueuedSubagentRun(runId, gatewayRunId)) { - throw new Error( - "collector registry row could not transition from queued to running", - ); - } - } catch (error) { - await terminateAcceptedRestoredCollectorRun({ - entry, - gatewayRunId, + await runWithGatewayIndependentRootWorkAdmission(async () => { + const response = await subagentRegistryDeps.callGateway({ + method: "agent", + params: applySubagentLaunchAuthorization(launch.request, launch.authorization), + // Restart replay must restore the trusted launch capability; otherwise + // the queued child silently falls back to its session/default route. + ...(launch.authorization ? { scopes: [ADMIN_SCOPE] } : {}), timeoutMs: launch.timeoutMs, }); - launchTerminationConfirmed = true; - throw error; - } + const gatewayRunId = readGatewayRunId(response) ?? runId; + try { + if (!startQueuedSubagentRun(runId, gatewayRunId)) { + throw new Error( + "collector registry row could not transition from queued to running", + ); + } + } catch (error) { + await terminateAcceptedRestoredCollectorRun({ + entry, + gatewayRunId, + timeoutMs: launch.timeoutMs, + }); + launchTerminationConfirmed = true; + throw error; + } + }); }, onStartFailure: (error) => { + if (error instanceof GatewayDrainingError) { + return false; + } return failAndCleanupRestoredQueuedRun( runId, entry, @@ -1827,12 +1840,9 @@ async function sweepSubagentRuns() { continue; } let deleteFailed = false; - // Lifecycle cleanup already attempted each delete-mode session. Retry - // here only so a transient cleanup failure cannot survive group archive. + // Group retention owns the final session archive for both keep- and + // delete-mode collectors. Retry every member so the batch is idempotent. for (const [candidateRunId, candidate] of groupEntries) { - if (candidate.cleanup !== "delete") { - continue; - } try { await subagentRegistryDeps.callGateway({ method: "sessions.delete", @@ -2027,6 +2037,7 @@ function ensureListener() { if (typeof entry.sessionStartedAt !== "number") { entry.sessionStartedAt = startedAt; } + entry.execution = { ...entry.execution, status: "running", startedAt }; persistSubagentRuns(); } return; diff --git a/src/agents/subagent-registry.types.ts b/src/agents/subagent-registry.types.ts index 752b1cc0f2c2..e3c9922502f2 100644 --- a/src/agents/subagent-registry.types.ts +++ b/src/agents/subagent-registry.types.ts @@ -39,6 +39,7 @@ export type PendingFinalDeliveryPayload = { export type SubagentExecutionState = { status: "queued" | "running" | "interrupted" | "terminal"; + acceptedAt?: number; startedAt?: number; endedAt?: number; outcome?: SubagentRunOutcome; diff --git a/src/agents/subagent-run-timeout.test.ts b/src/agents/subagent-run-timeout.test.ts index 52e8d931847d..46864654fce9 100644 --- a/src/agents/subagent-run-timeout.test.ts +++ b/src/agents/subagent-run-timeout.test.ts @@ -22,6 +22,33 @@ describe("subagent run timeout helpers", () => { ).toBe(2_592_001_000); }); + it("waits for the collector lifecycle start before setting its deadline", () => { + expect( + resolveSubagentRunDeadlineMs({ + collect: true, + createdAt: 1_000, + runTimeoutSeconds: 60, + }), + ).toBeUndefined(); + expect( + resolveSubagentRunDeadlineMs({ + collect: true, + createdAt: 1_000, + runTimeoutSeconds: 60, + }), + ).toBeUndefined(); + expect( + resolveSubagentRunDeadlineMs( + { + collect: true, + createdAt: 1_000, + runTimeoutSeconds: 60, + }, + 5_000, + ), + ).toBe(65_000); + }); + it("caps actual timer delays without shortening semantic durations", () => { // Long-lived subagent runs retain their requested deadline even though the // watchdog timer must be scheduled in bounded chunks. diff --git a/src/agents/subagent-run-timeout.ts b/src/agents/subagent-run-timeout.ts index 5bc593903ff4..b2cf3a85bdb5 100644 --- a/src/agents/subagent-run-timeout.ts +++ b/src/agents/subagent-run-timeout.ts @@ -29,7 +29,7 @@ export function resolveSubagentRunDurationMs(timeoutSeconds: unknown): number | /** Resolve the absolute timeout deadline for a subagent run. */ export function resolveSubagentRunDeadlineMs( - entry: Pick, + entry: Pick, observedStartedAt?: number, ): number | undefined { const durationMs = resolveSubagentRunDurationMs(entry.runTimeoutSeconds); @@ -41,7 +41,9 @@ export function resolveSubagentRunDeadlineMs( ? observedStartedAt : typeof entry.startedAt === "number" && Number.isFinite(entry.startedAt) ? entry.startedAt - : entry.createdAt; + : entry.collect + ? undefined + : entry.createdAt; const safeStartedAt = asDateTimestampMs(startedAt); if (safeStartedAt === undefined) { return undefined; @@ -54,7 +56,7 @@ export function resolveSubagentRunDeadlineMs( /** Clamp a reported terminal time to the run's explicit timeout deadline. */ export function resolveSubagentRunEffectiveEndedAt( - entry: Pick, + entry: Pick, endedAt: number, observedStartedAt?: number, ): number { diff --git a/src/agents/subagent-spawn.in-process-gateway.test.ts b/src/agents/subagent-spawn.in-process-gateway.test.ts index 41679c6bd51d..bfdb9b93724b 100644 --- a/src/agents/subagent-spawn.in-process-gateway.test.ts +++ b/src/agents/subagent-spawn.in-process-gateway.test.ts @@ -18,8 +18,14 @@ import { type dispatchGatewayMethodInProcess, } from "../gateway/server-plugins.js"; import { withPluginRuntimeGatewayRequestScope } from "../plugins/runtime/gateway-request-scope.js"; +import { + isGatewaySubordinateWorkAdmissionClosed, + resetGatewayWorkAdmission, + tryBeginGatewayRootWorkAdmission, +} from "../process/gateway-work-admission.js"; import { captureEnv, setTestEnvValue } from "../test-utils/env.js"; import { subagentRuns } from "./subagent-registry-memory.js"; +import { markSubagentRunTerminated } from "./subagent-registry.js"; import { resetSubagentRegistryForTests, testing as subagentRegistryTesting, @@ -90,6 +96,7 @@ async function waitForAssertion(assertion: () => void, timeoutMs = 2_000): Promi describe("spawnSubagentDirect in-process Gateway collector launch", () => { beforeEach(async () => { + resetGatewayWorkAdmission(); swarmSchedulerTesting.reset(); resetSubagentRegistryForTests({ persist: false }); clearFallbackGatewayContext(); @@ -109,7 +116,7 @@ describe("spawnSubagentDirect in-process Gateway collector launch", () => { path.join(stateDir, "openclaw.json"), `${JSON.stringify({ session: { mainKey: "main", scope: "per-sender" }, - tools: { swarm: true }, + tools: { swarm: { enabled: true, maxConcurrent: 1 } }, agents: { defaults: { workspace: stateDir }, entries: { main: { workspace: stateDir } }, @@ -121,6 +128,7 @@ describe("spawnSubagentDirect in-process Gateway collector launch", () => { afterEach(async () => { clearFallbackGatewayContext(); + resetGatewayWorkAdmission(); swarmSchedulerTesting.reset(); resetSubagentRegistryForTests({ persist: false }); subagentRegistryTesting.setDepsForTest(); @@ -134,6 +142,178 @@ describe("spawnSubagentDirect in-process Gateway collector launch", () => { } }); + it("launches queued collectors after the parent admission lease is released", async () => { + const gatewayContext = makeGatewayContext(); + let releaseFirstLaunch!: () => void; + const firstLaunchGate = new Promise((resolve) => { + releaseFirstLaunch = resolve; + }); + const subordinateAdmissionStates: boolean[] = []; + let launchCount = 0; + subagentSpawnTesting.setDepsForTest({ + dispatchGatewayMethodInProcess: async ( + _method: string, + params: Record, + ) => { + subordinateAdmissionStates.push(isGatewaySubordinateWorkAdmissionClosed()); + launchCount += 1; + if (launchCount === 1) { + await firstLaunchGate; + } + return { + runId: params.idempotencyKey as string, + status: "accepted", + } as T; + }, + }); + + const parentAdmission = tryBeginGatewayRootWorkAdmission(); + expect(parentAdmission).not.toBeNull(); + const results = await parentAdmission!.run(() => + withPluginRuntimeGatewayRequestScope( + { + context: gatewayContext, + client: externalCliClient(), + isWebchatConnect: () => false, + }, + () => + Promise.all([ + spawnSubagentDirect( + { + task: "first collector", + collect: true, + context: "isolated", + lightContext: true, + groupId: "swarm-queued-launch", + swarmLaunchReplayKey: "code-mode:agentSpawn:1", + }, + { + agentSessionKey: "agent:main:main", + requesterRunId: "parent-run", + }, + ), + spawnSubagentDirect( + { + task: "second collector", + collect: true, + context: "isolated", + lightContext: true, + groupId: "swarm-queued-launch", + swarmLaunchReplayKey: "code-mode:agentSpawn:2", + }, + { + agentSessionKey: "agent:main:main", + requesterRunId: "parent-run", + }, + ), + ]), + ), + ); + parentAdmission!.release(); + + expect(results.map((result) => result.status)).toEqual(["accepted", "accepted"]); + await waitForAssertion(() => { + expect(launchCount).toBe(1); + }); + releaseFirstLaunch(); + await waitForAssertion(() => { + expect(launchCount).toBe(2); + for (const result of results) { + expect(subagentRuns.get(result.runId!)).toMatchObject({ + collect: true, + swarmLaunchPending: false, + }); + } + }); + expect(subordinateAdmissionStates).toEqual([false, false]); + }); + + it("aborts a collector cancelled while Gateway acceptance is in flight", async () => { + const gatewayContext = makeGatewayContext(); + let releaseFirstLaunch!: () => void; + const firstLaunchGate = new Promise((resolve) => { + releaseFirstLaunch = resolve; + }); + const requests: Array<{ method: string; params: Record }> = []; + let launchCount = 0; + subagentSpawnTesting.setDepsForTest({ + dispatchGatewayMethodInProcess: async ( + method: string, + params: Record, + ) => { + requests.push({ method, params }); + if (method === "agent") { + launchCount += 1; + if (launchCount === 1) { + await firstLaunchGate; + } + return { runId: `gateway-run-${launchCount}`, status: "accepted" } as T; + } + return {} as T; + }, + }); + + const parentAdmission = tryBeginGatewayRootWorkAdmission(); + expect(parentAdmission).not.toBeNull(); + const results = await parentAdmission!.run(() => + withPluginRuntimeGatewayRequestScope( + { + context: gatewayContext, + client: externalCliClient(), + isWebchatConnect: () => false, + }, + () => + Promise.all([ + spawnSubagentDirect( + { + task: "cancelled collector", + collect: true, + context: "isolated", + lightContext: true, + groupId: "swarm-cancel-launch", + swarmLaunchReplayKey: "code-mode:agentSpawn:cancelled", + }, + { agentSessionKey: "agent:main:main", requesterRunId: "parent-run" }, + ), + spawnSubagentDirect( + { + task: "next collector", + collect: true, + context: "isolated", + lightContext: true, + groupId: "swarm-cancel-launch", + swarmLaunchReplayKey: "code-mode:agentSpawn:next", + }, + { agentSessionKey: "agent:main:main", requesterRunId: "parent-run" }, + ), + ]), + ), + ); + parentAdmission!.release(); + const firstRunId = results[0]?.runId; + expect(firstRunId).toBeTruthy(); + await waitForAssertion(() => expect(launchCount).toBe(1)); + + expect(markSubagentRunTerminated({ runId: firstRunId, reason: "manual kill" })).toBe(1); + releaseFirstLaunch(); + + await waitForAssertion(() => { + expect( + requests.some( + (request) => request.method === "chat.abort" && request.params.runId === "gateway-run-1", + ), + ).toBe(true); + expect(launchCount).toBe(2); + expect(subagentRuns.get(firstRunId!)).toMatchObject({ + collectorCompletion: { status: "killed" }, + }); + expect(subagentRuns.get("gateway-run-2")).toMatchObject({ + swarmRunId: results[1]!.runId, + swarmLaunchPending: false, + }); + }); + }); + it("hands a registered collector launch to Gateway as the host", async () => { const gatewayContext = makeGatewayContext(); const dispatchOptions: Array<{ method: string; forceSyntheticClient?: boolean }> = []; diff --git a/src/agents/subagent-spawn.ts b/src/agents/subagent-spawn.ts index a8bc6a2398d1..2e40f0ce1e1b 100644 --- a/src/agents/subagent-spawn.ts +++ b/src/agents/subagent-spawn.ts @@ -26,6 +26,10 @@ import { isFastTestRuntimeEnv } from "../infra/env.js"; import { stringifyRouteThreadId } from "../plugin-sdk/channel-route.js"; import { listRegisteredPluginAgentPromptGuidance } from "../plugins/command-registry-state.js"; import type { SubagentLifecycleHookRunner } from "../plugins/hooks.js"; +import { + GatewayDrainingError, + runWithGatewayIndependentRootWorkContinuation, +} from "../process/gateway-work-admission.js"; import { isValidAgentId, normalizeAgentId, parseAgentSessionKey } from "../routing/session-key.js"; import { recordSessionCreated, recordSubagentSpawned } from "../sessions/session-state-events.js"; import type { FastMode } from "../shared/fast-mode.js"; @@ -1846,20 +1850,27 @@ export async function spawnSubagentDirect( groupId: swarmSchedulerGroupKey, runId: childRunId, start: async () => { - const response = await launchChildRun(); - const gatewayRunId = readGatewayRunId(response) ?? childRunId; - try { - if (!startQueuedSubagentRun(childRunId, gatewayRunId)) { - throw new Error("collector registry row could not transition from queued to running"); + await runWithGatewayIndependentRootWorkContinuation(async () => { + const response = await launchChildRun(); + const gatewayRunId = readGatewayRunId(response) ?? childRunId; + try { + if (!startQueuedSubagentRun(childRunId, gatewayRunId)) { + throw new Error( + "collector registry row could not transition from queued to running", + ); + } + } catch (error) { + await terminateAcceptedCollectorRun({ childSessionKey, gatewayRunId }); + launchTerminationConfirmed = true; + throw error; } - } catch (error) { - await terminateAcceptedCollectorRun({ childSessionKey, gatewayRunId }); - launchTerminationConfirmed = true; - throw error; - } - await emitSpawnLifecycleHooks(gatewayRunId); + await emitSpawnLifecycleHooks(gatewayRunId); + }); }, onStartFailure: async (error) => { + if (error instanceof GatewayDrainingError) { + return false; + } const launchError = summarizeError(error); const [contextRollback, sessionCleanup] = await Promise.allSettled([ rollbackPreparedContextEngine(pipelineResult.state.contextEnginePreparation), diff --git a/src/agents/swarm-collector.ts b/src/agents/swarm-collector.ts index fa1abb1c3271..68b68513f8a4 100644 --- a/src/agents/swarm-collector.ts +++ b/src/agents/swarm-collector.ts @@ -1,27 +1,45 @@ +import type { OpenClawConfig } from "../config/types.openclaw.js"; +import { ensureCompletionState } from "./subagent-delivery-state.js"; import { SUBAGENT_ENDED_REASON_KILLED } from "./subagent-lifecycle-events.js"; +import { backfillCollectorArchiveAtMs } from "./subagent-registry-helpers.js"; import type { SubagentRunRecord, SwarmCollectorStatus } from "./subagent-registry.types.js"; import { loadSubagentSessionEntry } from "./subagent-session-reconciliation.js"; import { consumeSwarmStructuredOutput } from "./tools/structured-output-tool.js"; -function resolveStatus(entry: SubagentRunRecord): SwarmCollectorStatus { +function resolveStatus( + entry: SubagentRunRecord, + hasStructuredResult: boolean, +): SwarmCollectorStatus { if (entry.endedReason === SUBAGENT_ENDED_REASON_KILLED) { return "killed"; } if (entry.outcome?.status === "timeout") { return "timeout"; } - return entry.outcome?.status === "ok" ? "done" : "failed"; + if (entry.outcome?.status === "ok") { + return "done"; + } + // Tool-only structured turns can surface the runner's synthetic completion + // marker as an error despite having fulfilled the collector contract. + return hasStructuredResult && entry.outcome?.error === "completed" ? "done" : "failed"; } /** Freeze the waitable collector record after raw completion capture. */ -export function updateSwarmCollectorCompletion(entry: SubagentRunRecord): boolean { +export function updateSwarmCollectorCompletion( + entry: SubagentRunRecord, + cfg: OpenClawConfig, +): boolean { if (!entry.collect) { return false; } const clearedPendingLaunch = entry.swarmLaunchPending === true; entry.swarmLaunchPending = false; + const completion = ensureCompletionState(entry); + const capturedAtAdded = completion.capturedAt === undefined; + completion.capturedAt ??= Date.now(); + const archiveDeadlineAdded = backfillCollectorArchiveAtMs(entry, cfg); if (entry.collectorCompletion) { - return clearedPendingLaunch; + return clearedPendingLaunch || capturedAtAdded || archiveDeadlineAdded; } const executionCaptured = consumeSwarmStructuredOutput(entry.runId); const publicCaptured = @@ -42,7 +60,7 @@ export function updateSwarmCollectorCompletion(entry: SubagentRunRecord): boolea outputTokens: session.outputTokens ?? 0, } : undefined; - const resolvedStatus = resolveStatus(entry); + const resolvedStatus = resolveStatus(entry, captured?.structured !== undefined); const next = { status: schemaError && resolvedStatus === "done" ? ("failed" as const) : resolvedStatus, ...(captured?.structured !== undefined ? { structured: captured.structured } : {}), diff --git a/src/agents/tools/structured-output-tool.test.ts b/src/agents/tools/structured-output-tool.test.ts index 4c0fbb91bae0..10e037d5170e 100644 --- a/src/agents/tools/structured-output-tool.test.ts +++ b/src/agents/tools/structured-output-tool.test.ts @@ -21,6 +21,27 @@ describe("structured_output", () => { expect(testing.readSwarmStructuredOutput("run-1")?.structured).toEqual({ answer: "yes" }); }); + it("publishes a provider-valid schema while accepting any JSON result", () => { + const tool = createStructuredOutputTool({ + runId: "run-json-value", + schema: {}, + }); + expect(tool.parameters).toEqual({ + type: "object", + required: ["result"], + properties: { + result: { + type: ["object", "array", "string", "number", "boolean", "null"], + }, + }, + additionalProperties: false, + }); + for (const result of [{ answer: "yes" }, ["yes"], "yes", 1, true, null]) { + expect(Value.Check(tool.parameters, { result })).toBe(true); + } + expect(Value.Check(tool.parameters, { result: undefined })).toBe(false); + }); + it("nudges once then freezes schemaError", async () => { const tool = createStructuredOutputTool({ runId: "run-2", diff --git a/src/agents/tools/structured-output-tool.ts b/src/agents/tools/structured-output-tool.ts index 6d6bd446caaa..9b228b6ec0b4 100644 --- a/src/agents/tools/structured-output-tool.ts +++ b/src/agents/tools/structured-output-tool.ts @@ -14,7 +14,7 @@ function formatSchemaError(errors: Array<{ text: string }>): string { .join("; "); } -function readSwarmStructuredOutput(runId: string): SwarmStructuredOutputState | undefined { +export function peekSwarmStructuredOutput(runId: string): SwarmStructuredOutputState | undefined { const state = states.get(runId); return state ? structuredClone(state) : undefined; } @@ -22,7 +22,7 @@ function readSwarmStructuredOutput(runId: string): SwarmStructuredOutputState | export function consumeSwarmStructuredOutput( runId: string, ): SwarmStructuredOutputState | undefined { - const state = readSwarmStructuredOutput(runId); + const state = peekSwarmStructuredOutput(runId); states.delete(runId); return state; } @@ -59,8 +59,16 @@ export function createStructuredOutputTool(params: { displaySummary: "Record the collector result.", description: `Call exactly once as {"result": ...}, where result matches this JSON Schema: ${requestedSchema}`, // Runtime argument validation must reach execute so invalid attempts consume - // the durable one-retry budget. The requested schema remains model-visible above. - parameters: Type.Object({ result: Type.Unknown() }, { additionalProperties: false }), + // the durable one-retry budget. Providers still require every tool property to + // declare a JSON type before they will send the request. + parameters: Type.Object( + { + result: Type.Unsafe({ + type: ["object", "array", "string", "number", "boolean", "null"], + }), + }, + { additionalProperties: false }, + ), execute: async (_toolCallId, args) => { const prior = states.get(params.runId); if (prior?.structured !== undefined) { @@ -99,7 +107,7 @@ export function createStructuredOutputTool(params: { } const testing = { - readSwarmStructuredOutput, + readSwarmStructuredOutput: peekSwarmStructuredOutput, reset() { states.clear(); }, diff --git a/src/process/gateway-work-admission.ts b/src/process/gateway-work-admission.ts index 07f6b7e50abc..68c8a5a562f1 100644 --- a/src/process/gateway-work-admission.ts +++ b/src/process/gateway-work-admission.ts @@ -9,8 +9,8 @@ type AdmissionCloseReason = "restart-signal fence" | "restart drain" | "suspend type AdmissionReopenReason = "restart-signal fence" | "suspend phase"; export class GatewayDrainingError extends Error { - constructor() { - super("Gateway is draining; new tasks are not accepted"); + constructor(message = "Gateway is draining; new tasks are not accepted") { + super(message); this.name = "GatewayDrainingError"; } } @@ -302,7 +302,7 @@ export async function runWithGatewayIndependentRootWorkAdmission( ): Promise { while (true) { if (GATEWAY_WORK_ADMISSION_STATE.restartDraining) { - throw new Error("gateway is draining for restart"); + throw new GatewayDrainingError("gateway is draining for restart"); } const admission = tryBeginGatewayIndependentRootWorkAdmission(); if (admission) { diff --git a/ui/src/components/app-sidebar-child-session-data.ts b/ui/src/components/app-sidebar-child-session-data.ts index 355b6f27c24e..d6d219598e26 100644 --- a/ui/src/components/app-sidebar-child-session-data.ts +++ b/ui/src/components/app-sidebar-child-session-data.ts @@ -1,50 +1,9 @@ import type { GatewayBrowserClient } from "../api/gateway.ts"; import type { GatewaySessionRow } from "../api/types.ts"; -import type { SessionCapability } from "../lib/sessions/index.ts"; +export { fetchChildSessionRows } from "../lib/sessions/child-session-data.ts"; const MAX_SESSION_LINEAGE_DEPTH = 16; -export async function fetchChildSessionRows(params: { - sessions: SessionCapability; - parentKey: string; - isCurrent: () => boolean; -}): Promise { - const rows: GatewaySessionRow[] = []; - const seenOffsets = new Set(); - let offset = 0; - while (!seenOffsets.has(offset)) { - seenOffsets.add(offset); - const result = await params.sessions.list({ - spawnedBy: params.parentKey, - ...(offset > 0 ? { offset } : {}), - limit: 20, - includeGlobal: false, - includeUnknown: false, - configuredAgentsOnly: true, - }); - if (!params.isCurrent()) { - return null; - } - if (!result) { - throw new Error("child session list returned no result"); - } - const runtimeSampledAt = Date.now(); - for (const row of result.sessions) { - if (!rows.some((candidate) => candidate.key === row.key)) { - rows.push({ ...row, runtimeSampledAt }); - } - } - const hasMore = - result.hasMore ?? (typeof result.totalCount === "number" && rows.length < result.totalCount); - const nextOffset = result.nextOffset ?? rows.length; - if (!hasMore || nextOffset <= offset) { - break; - } - offset = nextOffset; - } - return rows; -} - export function collectKnownSessionRows( rootRows: readonly GatewaySessionRow[], childRowsByParent: Readonly>, diff --git a/ui/src/lib/board/builtin-dashboard.ts b/ui/src/lib/board/builtin-dashboard.ts index ad96121f4a35..2b658e9da621 100644 --- a/ui/src/lib/board/builtin-dashboard.ts +++ b/ui/src/lib/board/builtin-dashboard.ts @@ -1,7 +1,9 @@ import type { SessionObserverDigest } from "../../../../packages/gateway-protocol/src/schema/sessions.js"; import type { GatewaySessionRow } from "../../api/types.ts"; import { withObserverWidget } from "./observer-dashboard.ts"; -import { withSwarmWidget } from "./swarm-dashboard.ts"; +import { isSwarmEnabledInConfig, SwarmRosterHydrator, withSwarmWidget } from "./swarm-dashboard.ts"; + +export { isSwarmEnabledInConfig, SwarmRosterHydrator }; import type { BoardSnapshot } from "./types.ts"; import type { BoardViewSnapshot } from "./view-types.ts"; @@ -9,6 +11,8 @@ export function withBuiltinDashboardWidgets( snapshot: BoardSnapshot, sessions: readonly GatewaySessionRow[], observerDigests: readonly SessionObserverDigest[], + swarmEnabled = true, ): BoardViewSnapshot { - return withObserverWidget(withSwarmWidget(snapshot, sessions), observerDigests); + const withSwarm = swarmEnabled ? withSwarmWidget(snapshot, sessions) : snapshot; + return withObserverWidget(withSwarm, observerDigests); } diff --git a/ui/src/lib/board/swarm-dashboard-roster.ts b/ui/src/lib/board/swarm-dashboard-roster.ts new file mode 100644 index 000000000000..9b69e1a0b184 --- /dev/null +++ b/ui/src/lib/board/swarm-dashboard-roster.ts @@ -0,0 +1,40 @@ +import type { GatewaySessionRow } from "../../api/types.ts"; +import { fetchChildSessionRows } from "../sessions/child-session-data.ts"; +import type { SessionCapability } from "../sessions/index.ts"; + +const SWARM_SESSION_PAGE_SIZE = 10_000; + +function isNewerSessionRow(candidate: GatewaySessionRow, current: GatewaySessionRow): boolean { + // Callers pass hydrated rows first and the current lifecycle-decorated page + // second, so equal persisted timestamps intentionally prefer the latter. + return (candidate.updatedAt ?? 0) >= (current.updatedAt ?? 0); +} + +export function mergeSwarmSessionRows( + childRows: readonly GatewaySessionRow[], + currentRows: readonly GatewaySessionRow[], +): GatewaySessionRow[] { + const merged = new Map(); + for (const row of [...childRows, ...currentRows]) { + const current = merged.get(row.key); + if (!current || isNewerSessionRow(row, current)) { + merged.set(row.key, row); + } + } + return [...merged.values()]; +} + +export async function hydrateSwarmSessionRows(params: { + sessions: SessionCapability; + parentKey: string; + currentRows: readonly GatewaySessionRow[]; + isCurrent: () => boolean; +}): Promise { + const childRows = await fetchChildSessionRows({ + sessions: params.sessions, + parentKey: params.parentKey, + isCurrent: params.isCurrent, + pageSize: SWARM_SESSION_PAGE_SIZE, + }); + return childRows ? mergeSwarmSessionRows(params.currentRows, childRows) : null; +} diff --git a/ui/src/lib/board/swarm-dashboard.test.ts b/ui/src/lib/board/swarm-dashboard.test.ts new file mode 100644 index 000000000000..abc66459ef23 --- /dev/null +++ b/ui/src/lib/board/swarm-dashboard.test.ts @@ -0,0 +1,267 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import type { GatewaySessionRow, SessionsListResult } from "../../api/types.ts"; +import type { SessionCapability, SessionListOptions } from "../sessions/index.ts"; +import { hydrateSwarmSessionRows, mergeSwarmSessionRows } from "./swarm-dashboard-roster.ts"; +import { isSwarmEnabledInConfig, SwarmRosterHydrator } from "./swarm-dashboard.ts"; + +function row(index: number): GatewaySessionRow { + return { + key: `agent:worker:subagent:${index}`, + kind: "other", + updatedAt: index, + spawnedBy: "agent:main:parent", + swarmGroupId: "swarm:agent:main:parent:run-1", + } as unknown as GatewaySessionRow; +} + +function result(rows: GatewaySessionRow[], offset: number, totalCount: number): SessionsListResult { + const nextOffset = offset + rows.length; + return { + ts: Date.now(), + path: "state/openclaw.sqlite", + count: rows.length, + totalCount, + limitApplied: 10_000, + offset, + nextOffset: nextOffset < totalCount ? nextOffset : null, + hasMore: nextOffset < totalCount, + defaults: {} as SessionsListResult["defaults"], + sessions: rows, + }; +} + +afterEach(() => { + vi.useRealTimers(); +}); + +describe("isSwarmEnabledInConfig", () => { + it("accepts both the boolean and object configuration forms", () => { + expect(isSwarmEnabledInConfig({ tools: { swarm: true } })).toBe(true); + expect(isSwarmEnabledInConfig({ tools: { swarm: { enabled: true } } })).toBe(true); + expect(isSwarmEnabledInConfig({ tools: { swarm: false } })).toBe(false); + expect(isSwarmEnabledInConfig({ tools: { swarm: { enabled: false } } })).toBe(false); + expect( + isSwarmEnabledInConfig( + { + tools: { swarm: false }, + agents: { entries: { worker: { tools: { swarm: true } } } }, + }, + "worker", + ), + ).toBe(true); + expect( + isSwarmEnabledInConfig( + { + tools: { swarm: true }, + agents: { entries: { worker: { tools: { swarm: false } } } }, + }, + "worker", + ), + ).toBe(false); + expect( + isSwarmEnabledInConfig( + { + tools: { swarm: false }, + agents: { list: [{ id: "worker", tools: { swarm: true } }] }, + }, + "worker", + ), + ).toBe(true); + }); +}); + +describe("SwarmRosterHydrator", () => { + it("clears rows when the gateway source epoch changes", () => { + vi.useFakeTimers(); + const onRows = vi.fn(); + const hydrator = new SwarmRosterHydrator(); + const sessions = { + canonicalListRevision: 0, + list: vi.fn(async () => result([row(0)], 0, 1)), + } as unknown as SessionCapability; + + hydrator.update({ + sessions, + parentKey: "agent:main:parent", + sourceEpoch: 1, + currentRows: () => [row(0)], + onRows, + }); + expect(hydrator.rows).toHaveLength(1); + + hydrator.update({ + sessions, + parentKey: "agent:main:parent", + sourceEpoch: 2, + currentRows: () => [], + onRows, + }); + + expect(hydrator.rows).toEqual([]); + expect(onRows).toHaveBeenLastCalledWith([]); + hydrator.dispose(); + }); + + it("keeps a freshly fetched tie winner over an unchanged current page", async () => { + vi.useFakeTimers(); + const running = { ...row(0), status: "running" as const, updatedAt: 5 }; + const done = { ...row(0), status: "done" as const, updatedAt: 5 }; + const hydrator = new SwarmRosterHydrator(); + const sessions = { + canonicalListRevision: 0, + list: vi.fn(async () => result([done], 0, 1)), + } as unknown as SessionCapability; + + hydrator.update({ + sessions, + parentKey: "agent:main:parent", + sourceEpoch: 1, + currentRows: () => [running], + onRows: () => undefined, + }); + await vi.runAllTimersAsync(); + + expect(hydrator.rows).toEqual([expect.objectContaining({ status: "done" })]); + hydrator.dispose(); + }); + + it("keeps retrying at a bounded cadence after three transient failures", async () => { + vi.useFakeTimers(); + const onRows = vi.fn(); + const list = vi + .fn() + .mockRejectedValueOnce(new Error("offline")) + .mockRejectedValueOnce(new Error("offline")) + .mockRejectedValueOnce(new Error("offline")) + .mockResolvedValue(result([row(0)], 0, 1)); + const hydrator = new SwarmRosterHydrator(); + const sessions = { canonicalListRevision: 0, list } as unknown as SessionCapability; + + hydrator.update({ + sessions, + parentKey: "agent:main:parent", + sourceEpoch: 1, + currentRows: () => [], + onRows, + }); + await vi.runAllTimersAsync(); + + expect(list).toHaveBeenCalledTimes(4); + expect(hydrator.rows).toEqual([expect.objectContaining({ key: row(0).key })]); + hydrator.dispose(); + }); +}); + +describe("hydrateSwarmSessionRows", () => { + it("hydrates paginated cross-agent children outside the normal session page", async () => { + const children = Array.from({ length: 10_055 }, (_, index) => row(index)); + const list = vi.fn(async (options: SessionListOptions) => { + const offset = options.offset ?? 0; + return result(children.slice(offset, offset + 10_000), offset, children.length); + }); + const currentChild = { + ...row(0), + status: "running" as const, + updatedAt: 2_000, + } satisfies GatewaySessionRow; + const currentRows: GatewaySessionRow[] = [ + { + key: "agent:main:parent", + kind: "main", + updatedAt: 2_000, + } as unknown as GatewaySessionRow, + currentChild, + ]; + + const rows = await hydrateSwarmSessionRows({ + sessions: { list } as unknown as SessionCapability, + parentKey: "agent:main:parent", + currentRows, + isCurrent: () => true, + }); + + expect(rows).toHaveLength(10_056); + expect(rows?.find((candidate) => candidate.key === currentChild.key)?.status).toBe("running"); + expect(list).toHaveBeenCalledTimes(2); + expect(list).toHaveBeenNthCalledWith( + 1, + expect.objectContaining({ + spawnedBy: "agent:main:parent", + limit: 10_000, + includeGlobal: false, + configuredAgentsOnly: true, + }), + ); + expect(list).toHaveBeenNthCalledWith(2, expect.objectContaining({ offset: 10_000 })); + }); + + it("prefers the post-request server row when persisted timestamps tie", async () => { + const current = { ...row(0), status: "running" as const, updatedAt: 5 }; + const fetched = { ...row(0), status: "done" as const, updatedAt: 5 }; + + const rows = await hydrateSwarmSessionRows({ + sessions: { + list: vi.fn(async () => result([fetched], 0, 1)), + } as unknown as SessionCapability, + parentKey: "agent:main:parent", + currentRows: [current], + isCurrent: () => true, + }); + + expect(rows).toEqual([expect.objectContaining({ key: fetched.key, status: "done" })]); + }); + + it("restarts pagination when updated rows move across offset boundaries", async () => { + const running = { ...row(1), status: "running" as const }; + const done = { ...row(1), status: "done" as const, updatedAt: 10 }; + const pages = [ + [row(0), running], + [running, row(2)], + [row(3), row(0)], + [done, row(2)], + ]; + let callIndex = 0; + const list = vi.fn(async (options: SessionListOptions) => { + const rows = pages[callIndex] ?? []; + callIndex += 1; + return result(rows, options.offset ?? 0, 4); + }); + + const rows = await hydrateSwarmSessionRows({ + sessions: { list } as unknown as SessionCapability, + parentKey: "agent:main:parent", + currentRows: [], + isCurrent: () => true, + }); + + expect(rows?.map((candidate) => candidate.key).toSorted()).toEqual( + [row(0).key, row(1).key, row(2).key, row(3).key].toSorted(), + ); + expect(rows?.find((candidate) => candidate.key === done.key)?.status).toBe("done"); + expect(list).toHaveBeenCalledTimes(4); + }); + + it("keeps the freshest row when hydration overlaps a current-page snapshot", () => { + const stale = { ...row(0), status: "running" as const, updatedAt: 5, runtimeSampledAt: 10 }; + const fresh = { ...row(0), status: "done" as const, updatedAt: 6, runtimeSampledAt: 20 }; + + expect(mergeSwarmSessionRows([fresh], [stale])).toEqual([fresh]); + expect(mergeSwarmSessionRows([stale], [fresh])).toEqual([fresh]); + + const decorated = { ...stale, status: "done" as const }; + expect(mergeSwarmSessionRows([stale], [decorated])).toEqual([decorated]); + }); + + it("drops stale hydration results", async () => { + const rows = await hydrateSwarmSessionRows({ + sessions: { + list: vi.fn(async () => result([row(0)], 0, 1)), + } as unknown as SessionCapability, + parentKey: "agent:main:parent", + currentRows: [], + isCurrent: () => false, + }); + + expect(rows).toBeNull(); + }); +}); diff --git a/ui/src/lib/board/swarm-dashboard.ts b/ui/src/lib/board/swarm-dashboard.ts index 2af48db7cb1d..03527936b85b 100644 --- a/ui/src/lib/board/swarm-dashboard.ts +++ b/ui/src/lib/board/swarm-dashboard.ts @@ -1,12 +1,150 @@ +import { asNullableRecord } from "@openclaw/normalization-core/record-coerce"; import type { GatewaySessionRow } from "../../api/types.ts"; import { t } from "../../i18n/index.ts"; +import type { SessionCapability } from "../sessions/index.ts"; import { areUiSessionKeysEquivalent } from "../sessions/session-key.ts"; +import { hydrateSwarmSessionRows, mergeSwarmSessionRows } from "./swarm-dashboard-roster.ts"; import type { BoardSnapshot } from "./types.ts"; import type { BoardViewSnapshot } from "./view-types.ts"; const SWARM_TAB_ID = "builtin-swarm"; const SWARM_WIDGET_NAME = "builtin:swarm"; +function readSwarmEnabled(value: unknown): boolean | undefined { + if (typeof value === "boolean") { + return value; + } + const enabled = asNullableRecord(value)?.enabled; + return typeof enabled === "boolean" ? enabled : undefined; +} + +export function isSwarmEnabledInConfig(config: unknown, agentId?: string): boolean { + const root = asNullableRecord(config); + const globalEnabled = readSwarmEnabled(asNullableRecord(root?.tools)?.swarm); + const agents = asNullableRecord(root?.agents); + const entries = asNullableRecord(agents?.entries); + const listedEntries = Array.isArray(agents?.list) + ? agents.list + : Array.isArray(agents?.entries) + ? agents.entries + : []; + const listedAgent = agentId + ? listedEntries.map((entry) => asNullableRecord(entry)).find((entry) => entry?.id === agentId) + : undefined; + const agent = agentId ? (asNullableRecord(entries?.[agentId]) ?? listedAgent) : null; + const agentEnabled = readSwarmEnabled(asNullableRecord(agent?.tools)?.swarm); + return agentEnabled ?? globalEnabled ?? false; +} + +type SwarmHydrationParams = { + sessions: SessionCapability; + parentKey: string; + sourceEpoch: number; + currentRows: () => readonly GatewaySessionRow[]; + onRows: (rows: GatewaySessionRow[]) => void; +}; + +export class SwarmRosterHydrator { + rows: GatewaySessionRow[] = []; + private key = ""; + private revision = -1; + private generation = 0; + private attemptRevision = -1; + private attempts = 0; + private timer: ReturnType | null = null; + + update(params: SwarmHydrationParams): void { + const key = `${params.sourceEpoch}:${params.parentKey}`; + if (this.key !== key) { + this.reset(key); + } + this.rows = mergeSwarmSessionRows(this.rows, params.currentRows()); + params.onRows(this.rows); + const revision = params.sessions.canonicalListRevision; + if (this.attemptRevision !== revision) { + this.attemptRevision = revision; + this.attempts = 0; + } + if (this.revision === revision || this.timer !== null) { + return; + } + this.timer = setTimeout(() => this.hydrate(params), 250); + } + + dispose(): void { + this.reset(""); + } + + private hydrate(params: SwarmHydrationParams): void { + const generation = this.generation; + const revision = params.sessions.canonicalListRevision; + const key = `${params.sourceEpoch}:${params.parentKey}`; + const isCurrent = () => generation === this.generation && this.key === key; + const currentRowsAtStart = params.currentRows(); + const currentRowsAtStartByKey = new Map( + currentRowsAtStart.map((row) => [row.key, JSON.stringify(row)]), + ); + let hydrated = false; + let retrying = false; + this.attempts += 1; + void hydrateSwarmSessionRows({ + sessions: params.sessions, + parentKey: params.parentKey, + currentRows: currentRowsAtStart, + isCurrent, + }) + .then((rows) => { + if (!rows || !isCurrent()) { + return; + } + hydrated = true; + this.revision = revision; + const changedCurrentRows = params + .currentRows() + .filter((row) => currentRowsAtStartByKey.get(row.key) !== JSON.stringify(row)); + this.rows = mergeSwarmSessionRows(rows, changedCurrentRows); + params.onRows(this.rows); + }) + .catch(() => { + if (!isCurrent()) { + return; + } + retrying = true; + const retryDelayMs = Math.min(30_000, 1_000 * 2 ** Math.min(this.attempts - 1, 5)); + this.timer = setTimeout(() => { + this.timer = null; + if (isCurrent()) { + this.update(params); + } + }, retryDelayMs); + }) + .finally(() => { + if (!isCurrent()) { + return; + } + if (!retrying) { + this.timer = null; + } + if (hydrated && this.revision !== params.sessions.canonicalListRevision) { + this.update(params); + } + }); + } + + private reset(key: string): void { + if (this.timer !== null) { + clearTimeout(this.timer); + } + this.rows = []; + this.key = key; + this.revision = -1; + this.generation += 1; + this.attemptRevision = -1; + this.attempts = 0; + this.timer = null; + } +} + function hasSwarmRowsForSession( sessions: readonly GatewaySessionRow[], sessionKey: string, diff --git a/ui/src/lib/board/widgets/swarm.test.ts b/ui/src/lib/board/widgets/swarm.test.ts index 5f3737652208..485ef0db80ac 100644 --- a/ui/src/lib/board/widgets/swarm.test.ts +++ b/ui/src/lib/board/widgets/swarm.test.ts @@ -64,6 +64,44 @@ describe("swarm board widget", () => { ]); }); + it("renders every child beyond the ordinary 50-row session page", () => { + const container = document.createElement("div"); + document.body.append(container); + render( + renderSwarmWidget({ + sessionKey: parentSessionKey, + sessions: Array.from({ length: 55 }, (_, index) => + session({ key: `child-${index}`, status: "running" }), + ), + }), + container, + ); + + expect(container.querySelectorAll(".swarm-widget__dot")).toHaveLength(55); + }); + + it("caps historical dots while keeping active workers visible", () => { + const container = document.createElement("div"); + document.body.append(container); + render( + renderSwarmWidget({ + sessionKey: parentSessionKey, + sessions: [ + ...Array.from({ length: 300 }, (_, index) => + session({ key: `done-${index}`, status: "done" }), + ), + session({ key: "running", status: "running" }), + ], + }), + container, + ); + + expect(container.querySelectorAll(".swarm-widget__dot")).toHaveLength(256); + expect(container.querySelector(".swarm-widget__dot--running")).not.toBeNull(); + expect(container.querySelector(".swarm-widget__more")?.textContent?.trim()).toBe("+45"); + expect(container.textContent?.replace(/\s+/g, " ")).toContain("1 Running ยท 300 Done"); + }); + it("renders dot tooltips and keeps an empty state when no group is active", () => { const container = document.createElement("div"); document.body.append(container); diff --git a/ui/src/lib/board/widgets/swarm.ts b/ui/src/lib/board/widgets/swarm.ts index 3ca1db441856..fb3d230b94c0 100644 --- a/ui/src/lib/board/widgets/swarm.ts +++ b/ui/src/lib/board/widgets/swarm.ts @@ -5,15 +5,20 @@ import { areUiSessionKeysEquivalent } from "../../sessions/session-key.ts"; type SwarmDotStatus = "queued" | "running" | "done" | "failed"; +const SWARM_DOT_STATUS_RANK = { running: 0, queued: 1, failed: 2, done: 3 } as const; + type SwarmDot = { key: string; label: string; status: SwarmDotStatus; }; +const MAX_RENDERED_DOTS_PER_PHASE = 256; + type SwarmPhase = { title?: string; dots: SwarmDot[]; + hidden: number; }; type SwarmGroup = { @@ -140,7 +145,20 @@ function collectActiveSwarmGroups( narrator: entries.map((entry) => entry.log).find(Boolean), phases: [...phases.entries()] .toSorted((left, right) => left[1].rank - right[1].rank) - .map(([title, bucket]) => ({ title, dots: bucket.dots })), + .map(([title, bucket]) => { + const visibleFirst = + bucket.dots.length > MAX_RENDERED_DOTS_PER_PHASE + ? bucket.dots.toSorted( + (left, right) => + SWARM_DOT_STATUS_RANK[left.status] - SWARM_DOT_STATUS_RANK[right.status], + ) + : bucket.dots; + return { + title, + dots: visibleFirst.slice(0, MAX_RENDERED_DOTS_PER_PHASE), + hidden: Math.max(0, visibleFirst.length - MAX_RENDERED_DOTS_PER_PHASE), + }; + }), } satisfies SwarmGroup; }) .filter((group) => @@ -196,6 +214,11 @@ export function renderSwarmWidget({ > `, )} + ${phase.hidden > 0 + ? html`+${phase.hidden}` + : nothing} `, diff --git a/ui/src/lib/sessions/child-session-data.ts b/ui/src/lib/sessions/child-session-data.ts new file mode 100644 index 000000000000..0c051b35a553 --- /dev/null +++ b/ui/src/lib/sessions/child-session-data.ts @@ -0,0 +1,60 @@ +import type { GatewaySessionRow } from "../../api/types.ts"; +import type { SessionCapability } from "./index.ts"; + +const MAX_CHILD_SESSION_LIST_PASSES = 4; + +export async function fetchChildSessionRows(params: { + sessions: SessionCapability; + parentKey: string; + isCurrent: () => boolean; + pageSize?: number; +}): Promise { + const rowsByKey = new Map(); + const pageSize = params.pageSize ?? 20; + for (let pass = 0; pass < MAX_CHILD_SESSION_LIST_PASSES; pass += 1) { + const seenOffsets = new Set(); + const rowsBeforePass = rowsByKey.size; + let expectedTotal: number | undefined; + let offset = 0; + while (!seenOffsets.has(offset)) { + seenOffsets.add(offset); + const result = await params.sessions.list({ + spawnedBy: params.parentKey, + ...(offset > 0 ? { offset } : {}), + limit: pageSize, + includeGlobal: false, + includeUnknown: false, + configuredAgentsOnly: true, + }); + if (!params.isCurrent()) { + return null; + } + if (!result) { + throw new Error("child session list returned no result"); + } + expectedTotal = result.totalCount; + const runtimeSampledAt = Date.now(); + for (const row of result.sessions) { + // A later pass is a fresher server observation even when the row moved + // across an updatedAt-sorted offset boundary. + rowsByKey.set(row.key, { ...row, runtimeSampledAt }); + } + const hasMore = + result.hasMore ?? + (typeof result.totalCount === "number" && + offset + result.sessions.length < result.totalCount); + const nextOffset = result.nextOffset ?? offset + result.sessions.length; + if (!hasMore || nextOffset <= offset) { + break; + } + offset = nextOffset; + } + const addedThisPass = rowsByKey.size - rowsBeforePass; + if (addedThisPass === 0 || expectedTotal === undefined || rowsByKey.size >= expectedTotal) { + break; + } + // updatedAt ordering can move a child across an offset boundary while paging. + // Repeat from zero until the deduplicated roster reaches the latest total. + } + return [...rowsByKey.values()]; +} diff --git a/ui/src/lib/sessions/index.ts b/ui/src/lib/sessions/index.ts index ebb1a5f45cf5..88b3a21eceb7 100644 --- a/ui/src/lib/sessions/index.ts +++ b/ui/src/lib/sessions/index.ts @@ -189,7 +189,7 @@ export type SessionMessageSubscription = { export type SessionCapability = { readonly state: SessionState; - /** Advances only when a canonical sessions.list response is published. */ + /** Advances only when a canonical sessions.list result is published. */ readonly canonicalListRevision: number; list: (options?: SessionListOptions) => Promise; setCreatorFilter: (creatorId: string | null) => Promise; @@ -869,7 +869,7 @@ export function createSessionCapability(gateway: SessionGateway): SessionCapabil return null; } const result = await requestSessionList(scope.client, options); - return isCurrentConnection(scope) ? (result ?? null) : null; + return isCurrentConnection(scope) ? swarmActivity.decorate(result ?? null) : null; }; const publish = (next: SessionState) => { @@ -1419,7 +1419,6 @@ export function createSessionCapability(gateway: SessionGateway): SessionCapabil payload: unknown, options?: SessionReconcileOptions, ): SessionChangedResult => { - swarmActivity.observe(payload); const base = reconcileSessionChanged(state.result, payload, options); const result = swarmActivity.decorate(base.result); const reconciled = diff --git a/ui/src/lib/sessions/swarm-activity.integration.test.ts b/ui/src/lib/sessions/swarm-activity.integration.test.ts index ba3c88271fa2..d754414765b3 100644 --- a/ui/src/lib/sessions/swarm-activity.integration.test.ts +++ b/ui/src/lib/sessions/swarm-activity.integration.test.ts @@ -96,8 +96,13 @@ describe("session swarm activity", () => { emitEvent({ type: "event", event: "sessions.changed", payload }); await sessions.refresh({ force: true }); + const revisionBeforePhase = sessions.canonicalListRevision; emitChanged(note("phase", "Plan")); + expect(sessions.canonicalListRevision).toBe(revisionBeforePhase); await waitForFast(() => expect(request).toHaveBeenCalledTimes(2)); + await waitForFast(() => + expect(sessions.canonicalListRevision).toBeGreaterThan(revisionBeforePhase), + ); rows = [ ...rows, { diff --git a/ui/src/lib/sessions/swarm-activity.ts b/ui/src/lib/sessions/swarm-activity.ts index 9318f7bb277e..6b825957d8da 100644 --- a/ui/src/lib/sessions/swarm-activity.ts +++ b/ui/src/lib/sessions/swarm-activity.ts @@ -2,8 +2,10 @@ import { asNullableRecord } from "@openclaw/normalization-core/record-coerce"; import type { GatewaySessionRow, SessionsListResult } from "../../api/types.ts"; // Lifecycle notes are transient UI state, so bound them for long-lived board tabs. -const MAX_TRACKED_SWARM_GROUPS = 128; -const MAX_TRACKED_SWARM_CHILDREN = 2_048; +const MAX_TRACKED_SWARM_GROUPS = 10_000; +// Completed members stay visible while any group child is active, so retain the +// supported lifetime membership ceiling rather than only the live-child cap. +const MAX_TRACKED_SWARM_CHILDREN = 100_000; type SwarmDisplayCarrier = { swarmPhaseRank?: number; @@ -45,17 +47,16 @@ export class SwarmActivityTracker { this.phaseByChild.clear(); } - observe(payload: unknown): void { + observe(payload: unknown): boolean { const event = asNullableRecord(payload); if (!event) { - return; + return false; } const source = asNullableRecord(event.session) ?? event; const groupId = normalizedString(event.swarmGroupId) ?? normalizedString(source.swarmGroupId); if (!groupId) { - return; + return false; } - const kind = normalizedString(event.kind); const text = normalizedString(event.text); if ((kind === "phase" || kind === "log") && text) { @@ -77,17 +78,17 @@ export class SwarmActivityTracker { text, MAX_TRACKED_SWARM_GROUPS, ); - return; + return true; } const childKey = normalizedString(source.key) ?? normalizedString(event.sessionKey); if (!childKey) { - return; + return true; } const explicitPhase = normalizedString(source.swarmPhase) ?? normalizedString(event.swarmPhase); if (explicitPhase) { setBounded(this.phaseByChild, childKey, explicitPhase, MAX_TRACKED_SWARM_CHILDREN); - return; + return true; } // Implicit phase assignment is a creation-time fact: only a child ADMITTED // after phase('X') belongs to X. Status/completion updates for a child that @@ -98,6 +99,7 @@ export class SwarmActivityTracker { setBounded(this.phaseByChild, childKey, currentPhase, MAX_TRACKED_SWARM_CHILDREN); } } + return true; } decorate(result: SessionsListResult | null): SessionsListResult | null { diff --git a/ui/src/pages/chat/chat-pane.ts b/ui/src/pages/chat/chat-pane.ts index bbc580b6588d..9b34ade5376a 100644 --- a/ui/src/pages/chat/chat-pane.ts +++ b/ui/src/pages/chat/chat-pane.ts @@ -91,6 +91,7 @@ import { type BoardFace, type BoardSessionView, } from "../../lib/board/settings.ts"; +import type { SwarmRosterHydrator } from "../../lib/board/swarm-dashboard.ts"; import type { BoardSnapshot, BoardTab } from "../../lib/board/types.ts"; import type { BoardViewSnapshot } from "../../lib/board/view-types.ts"; import { @@ -492,7 +493,7 @@ class ChatPane extends OpenClawLightDomElement { private readonly observerDigestHistory = new ObserverDigestHistory(); private builtinBoardSnapshot: BoardViewSnapshot | null = null; private builtinBoardSnapshotBase: BoardSnapshot | null = null; - private builtinBoardSnapshotRequest = 0; + private swarmHydrator: SwarmRosterHydrator | null = null; private readonly sessionDiscussionStates = new Map(); private readonly sessionDiscussionOpenUrls = new Map(); private readonly sessionDiscussionProbes = new Set(); @@ -623,7 +624,11 @@ class ChatPane extends OpenClawLightDomElement { ) .watch( () => this.context?.runtimeConfig, - (runtimeConfig, notify) => runtimeConfig.subscribe(notify), + (runtimeConfig, notify) => + runtimeConfig.subscribe(() => { + this.refreshBuiltinBoardSnapshot(); + notify(); + }), ) .watch( () => this.resolveBoardProvider(), @@ -1825,22 +1830,54 @@ class ChatPane extends OpenClawLightDomElement { if (!state) { return; } - const request = ++this.builtinBoardSnapshotRequest; - const sessions = state.sessionsResult?.sessions ?? []; - void import("../../lib/board/builtin-dashboard.ts").then(({ withBuiltinDashboardWidgets }) => { - if (request !== this.builtinBoardSnapshotRequest) { - return; - } - const currentBase = this.resolveBoardProvider().snapshot$.value; - const sessionKey = this.resolveBoardSessionKey(currentBase.sessionKey); - this.builtinBoardSnapshotBase = currentBase; - this.builtinBoardSnapshot = withBuiltinDashboardWidgets( - currentBase, - sessions, - this.observerDigestHistory.get(sessionKey), - ); - this.requestUpdate(); - }); + const parentKey = this.resolveBoardSessionKey(); + const sourceEpoch = state.connectionEpoch; + void import("../../lib/board/builtin-dashboard.ts").then( + ({ isSwarmEnabledInConfig, SwarmRosterHydrator, withBuiltinDashboardWidgets }) => { + if ( + !this.state || + this.state.connectionEpoch !== sourceEpoch || + parentKey !== this.resolveBoardSessionKey() + ) { + return; + } + const swarmEnabled = + this.state.connected && + isSwarmEnabledInConfig( + this.context.runtimeConfig?.state.configSnapshot?.config, + resolveAgentIdFromSessionKey(parentKey), + ); + const applyRows = (rows: readonly GatewaySessionRow[], includeSwarm: boolean) => { + const base = this.resolveBoardProvider().snapshot$.value; + const sessionKey = this.resolveBoardSessionKey(base.sessionKey); + this.builtinBoardSnapshotBase = base; + this.builtinBoardSnapshot = withBuiltinDashboardWidgets( + base, + rows, + this.observerDigestHistory.get(sessionKey), + includeSwarm, + ); + this.requestUpdate(); + }; + if (!swarmEnabled) { + this.swarmHydrator?.dispose(); + this.swarmHydrator = null; + applyRows(this.state.sessionsResult?.sessions ?? [], false); + return; + } + this.swarmHydrator ??= new SwarmRosterHydrator(); + this.swarmHydrator.update({ + sessions: this.context.sessions, + parentKey, + sourceEpoch, + currentRows: () => + this.state?.connectionEpoch === sourceEpoch + ? (this.state.sessionsResult?.sessions ?? []) + : [], + onRows: (rows) => applyRows(rows, true), + }); + }, + ); } private recordObserverDigest(digest: SessionObserverDigest): void { @@ -2614,6 +2651,8 @@ class ChatPane extends OpenClawLightDomElement { window.clearTimeout(this.headerCopiedTimer); this.headerCopiedTimer = null; } + this.swarmHydrator?.dispose(); + this.swarmHydrator = null; this.headerWorktreePaths.clear(); this.headerBranches.clear(); this.presencePayload = undefined; @@ -2748,6 +2787,10 @@ class ChatPane extends OpenClawLightDomElement { // A reconnect can retain the browser client. Keep async ownership tied // to the logical connection, not only the transport object identity. this.connectionGeneration += 1; + this.swarmHydrator?.dispose(); + this.swarmHydrator = null; + this.builtinBoardSnapshot = null; + this.builtinBoardSnapshotBase = null; this.taskSuggestionsRequestVersion += 1; this.taskSuggestions = []; this.taskSuggestionBusyIds.clear(); @@ -3907,7 +3950,7 @@ class ChatPane extends OpenClawLightDomElement { board.hasBoard && board.face === "dashboard" ? renderBoardSessionSurface({ snapshot: board.snapshot, - sessions: state.sessionsResult?.sessions ?? [], + sessions: this.swarmHydrator?.rows ?? state.sessionsResult?.sessions ?? [], observer: { activeRunId: observerRunId, digests: this.observerDigestHistory.get( diff --git a/ui/src/styles/board.css b/ui/src/styles/board.css index f2d18bbd4542..6d539f101488 100644 --- a/ui/src/styles/board.css +++ b/ui/src/styles/board.css @@ -388,6 +388,12 @@ openclaw-board-widget-cell { width: 9px; } +.swarm-widget__more { + color: var(--muted, #8a919e); + font-size: 10px; + line-height: 9px; +} + .swarm-widget__dot--running { animation: swarm-widget-pulse 1.25s ease-in-out infinite; background: var(--accent, #ff5c5c);