diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json index cef1865d17e2..47e28e540f31 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json @@ -1 +1 @@ -{"contentHash":"e2a5ce75757e7a275312c91ce95bdb4f13091a3c613dd4c66fdda0b2e4d3bf83","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"} +{"contentHash":"57cb6a98df58f04f4c31d37b35aa857e5329eb29560f272381be592e39946d5a","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json index be3be85887a4..c76f918f806d 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json @@ -1 +1 @@ -{"contentHash":"0e4494129d83b508c82788383f8dfead578976ed474be2cfbec557aa6caf88e7","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} +{"contentHash":"a3e367817d1aa078cbd028aec5b080e73f7b67a85f6b91e4cfdfe32a09e6b509","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-core.json b/docs/.generated/plugin-sdk-api-baseline/channel-core.json index f09a6a6b79c0..a523521b8e3c 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-core.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-core.json @@ -1 +1 @@ -{"contentHash":"8643a3553553855a519e825fcaa483024ecd3538675c5a2d474af841fb5070d7","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} +{"contentHash":"355f3a91dabce0393a540e1d9f476eba6926cff82153f41e98c5a01d4e011ea2","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json index 75f1fa2425cb..ac28841df7cb 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json @@ -1 +1 @@ -{"contentHash":"13ec7276e7c37ae4e7589e67302c2bf531674f2c439b70d4a8a44a4ee136f685","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"} +{"contentHash":"695b06456e6e2bb7c3ca5891619683bf14e5cb5939e4a44bd907b2c798910cae","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json b/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json index cef050dd62be..0279888b253e 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json @@ -1 +1 @@ -{"contentHash":"ed0a5bf19bbe6e9641c8de12ea94105865850343c319a58953a2dce21e6d3f49","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"} +{"contentHash":"cf487c3dd656a5f9d8c4a927bf0cb8ff27b1fb29c09c466c8685482c647df12d","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-message.json b/docs/.generated/plugin-sdk-api-baseline/channel-message.json index 5c94495737f6..ad991daf3eae 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-message.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-message.json @@ -1 +1 @@ -{"contentHash":"feb6187469b8f0e88c0445ba025cc557c94871db55f39075eb63044d8ba48cd7","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} +{"contentHash":"f9a6c7bd0ded857948b59ec8a8e2e8cc36204888e76cd5cfb50133439be59b1a","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json index 825f8710f7d7..c3bc84b19510 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json @@ -1 +1 @@ -{"contentHash":"97d39817d24eaf5b2808960770447188daa0b2959514efbd44db007797e40e3b","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} +{"contentHash":"dfebc8eae36256aedf44b8e4ea9a137f3b35946cc624d09246ec8d614223d752","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json index 3ef00e4dac54..17512d7ba4f2 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json @@ -1 +1 @@ -{"contentHash":"e824840c22ccade58891ad2952a6aea4aafa8c163adb4383e9a5aa8b47cce44b","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"} +{"contentHash":"524471a854edf593a8f9f3b80eaedf4492e9bc60cea40e314ad7efbbbbf46d32","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"} diff --git a/docs/.generated/plugin-sdk-api-baseline/core.json b/docs/.generated/plugin-sdk-api-baseline/core.json index f4b0c317ba1b..0bba4e2c152d 100644 --- a/docs/.generated/plugin-sdk-api-baseline/core.json +++ b/docs/.generated/plugin-sdk-api-baseline/core.json @@ -1 +1 @@ -{"contentHash":"08bba8059fb44659737936c0ef4c163cfff29872db3088b72d2fff82f4c4ac85","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} +{"contentHash":"f0a95c43b07dffdfbd05ac63a1b1dfd4b4828fab8d3613ca51465a7bbcae8cdf","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/discord.json b/docs/.generated/plugin-sdk-api-baseline/discord.json index 8e24587c0713..3776e2eb546d 100644 --- a/docs/.generated/plugin-sdk-api-baseline/discord.json +++ b/docs/.generated/plugin-sdk-api-baseline/discord.json @@ -1 +1 @@ -{"contentHash":"db6ac3119377920f375bdb1cd01082815022f2ad49de3dc5d7ac258a6a864ca6","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} +{"contentHash":"8c243343f62835f8570cf85a7c58523d76d5fed4f80704a179417d993a3ffc5d","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} diff --git a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json index e5f37fac7973..e6f08fd922d1 100644 --- a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json +++ b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json @@ -1 +1 @@ -{"contentHash":"7bb8bb13f101a7018aea253debd9c3611635541b8ad2d847e827300b9f94991f","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"} +{"contentHash":"a78162d9df85c5e96ef95b985a036af55110afe5768174c761d6fa0a485aa621","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"} diff --git a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json index 1620befd327d..01360aa32342 100644 --- a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json @@ -1 +1 @@ -{"contentHash":"e5c950c85caf1fe7edf4d672f95afc074e926e300981d96707ed904304e001b4","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} +{"contentHash":"46ce609eeb7ac8b23cdfb90469301cb988d4c18fdecfe5aa9aa2e444a52cb345","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json index 077b526765f1..824c124c059c 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json @@ -1 +1 @@ -{"contentHash":"f6867c932f2632409f933d8e158d01f303679950b4e6e3a01ba28b0ea8b4d7da","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} +{"contentHash":"fd5a6f903c1d177ccb090d0507bdec51d9494f094c8219d56334553044fff7cf","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json index 41fcb539bf51..cd28c10daf0e 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json @@ -1 +1 @@ -{"contentHash":"53a44cf7b1b8866054e905dd7a34623e75dbf4be6da5b58e7d47fd638e2a92f5","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} +{"contentHash":"bbd1fa1f3afa8c8699a69a17d33d4ec9e6db46c5b94636f19e0cc5fe31b0b125","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json index 824af6e0846e..7afbc163161e 100644 --- a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json @@ -1 +1 @@ -{"contentHash":"03aeade34faf6a1c6bfee3e88f684ccb4abe657ea8f3475c96162698baeeb314","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"} +{"contentHash":"9bab151043bafb665d20beb0f3e2685bbb9ac8fd7a19316cc695dafa7f977a5b","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json index 750aa90a0ebb..da2b38e190cc 100644 --- a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json +++ b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json @@ -1 +1 @@ -{"contentHash":"83795fac5e17a90840255de8446aaf23b988a14d19bc76a327b0b74ef5f9ee8a","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} +{"contentHash":"6e2ee162d479ea014a08c64d6b85276e60d099137baf348f5d193bc8ed1cf162","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} diff --git a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json index 50c7b7d50551..f83a0576dc7e 100644 --- a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json +++ b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json @@ -1 +1 @@ -{"contentHash":"9511576a6c8d6631cbe81dc7aca200f670feaf7fc03d6b3308066b5caa3ab50c","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} +{"contentHash":"267759cc867869904787ed33483da8c59eba7f290427500836d98247ed135d91","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} diff --git a/docs/plan/runners.md b/docs/plan/runners.md index d08f53f1639c..e1bd5270b417 100644 --- a/docs/plan/runners.md +++ b/docs/plan/runners.md @@ -13,22 +13,22 @@ Proposal, revision 2. Supersedes revision 1 in place (2026-08-11, operator decision). Implementation in progress; update this table in every PR that advances a milestone. -| # | Milestone | Status | PRs | -| --- | ---------------------------------------------------------- | ----------- | ------------------------- | -| 0 | This plan (revision 2) | landed | #122454 | -| 1a | Naming: session copy revert | landed | #120667 | -| 1b | Naming: devices consolidation | landed | #120689 | -| 1c | Cleanup: node-pairing → device-pairing merge | landed | #120726 | -| 2 | `openclaw resume` + web Continue in terminal | in progress | #120664 | -| 3 | `openclaw connect` one-paste onboarding + `/j/` join route | in progress | #120768, #122499 | -| 4 | Picker: grouping, placement, liveness, enrichment | in progress | #120804, #122531, #122635 | -| F | Real-wire session boundary harness | landed | #121212 | -| 5 | Public worker ingress path | landed | #122578, #122643 | -| 6 | Node worker provider (device runners) | in progress | #122683 | -| 7 | Bundle push consent + runner updates | not started | — | -| 8 | Stop-and-continue moves | not started | — | -| 9 | Deletions (ssh sandbox, openshell, exec-host clones, …) | not started | — | -| 10 | Cloud convergence (provisioners run `openclaw connect`) | not started | — | +| # | Milestone | Status | PRs | +| --- | ---------------------------------------------------------- | ----------- | ---------------------------------- | +| 0 | This plan (revision 2) | landed | #122454 | +| 1a | Naming: session copy revert | landed | #120667 | +| 1b | Naming: devices consolidation | landed | #120689 | +| 1c | Cleanup: node-pairing → device-pairing merge | landed | #120726 | +| 2 | `openclaw resume` + web Continue in terminal | in progress | #120664 | +| 3 | `openclaw connect` one-paste onboarding + `/j/` join route | in progress | #120768, #122499 | +| 4 | Picker: grouping, placement, liveness, enrichment | in progress | #120804, #122531, #122635, #122774 | +| F | Real-wire session boundary harness | landed | #121212 | +| 5 | Public worker ingress path | landed | #122578, #122643 | +| 6 | Node worker provider (device runners) | in progress | #122683, #122829 | +| 7 | Bundle push consent + runner updates | not started | — | +| 8 | Stop-and-continue moves | not started | — | +| 9 | Deletions (ssh sandbox, openshell, exec-host clones, …) | not started | — | +| 10 | Cloud convergence (provisioners run `openclaw connect`) | not started | — | Revision history: revision 1 (2026-08-08) established the session/runner vocabulary, the naming rulings, and the milestone skeleton after a diff --git a/src/cli/worker-cli.ts b/src/cli/worker-cli.ts index 3aa455868f64..c9d86b570142 100644 --- a/src/cli/worker-cli.ts +++ b/src/cli/worker-cli.ts @@ -1,14 +1,107 @@ -import type { Command } from "commander"; +import { Option, type Command } from "commander"; +import { signalProcessTree } from "../process/kill-tree.js"; +import type { WorkerCommandLifetime } from "../worker/worker-command.runtime.js"; + +const WORKER_START_MESSAGE_TYPE = "openclaw-worker-start-v1"; + +function isWorkerStartMessage(value: unknown): boolean { + return ( + typeof value === "object" && + value !== null && + !Array.isArray(value) && + Object.keys(value).length === 1 && + (value as { type?: unknown }).type === WORKER_START_MESSAGE_TYPE + ); +} + +function createWorkerIpcLifetime(): WorkerCommandLifetime { + if (!process.connected || !process.channel || typeof process.send !== "function") { + throw new Error("internal worker IPC mode requires a connected Node IPC channel"); + } + const abortController = new AbortController(); + let disposed = false; + let started = false; + let settled = false; + let resolveStarted!: (started: boolean) => void; + let rejectStarted!: (error: Error) => void; + const startedPromise = new Promise((resolve, reject) => { + resolveStarted = resolve; + rejectStarted = reject; + }); + const rejectOrAbort = (error: Error) => { + if (!settled) { + settled = true; + rejectStarted(error); + return; + } + abortController.abort(error); + }; + const onMessage = (message: unknown) => { + if (disposed) { + return; + } + if (!isWorkerStartMessage(message) || settled) { + rejectOrAbort(new Error("invalid internal worker IPC start message")); + return; + } + started = true; + settled = true; + resolveStarted(true); + }; + const onDisconnect = () => { + if (disposed) { + return; + } + if (!settled) { + settled = true; + resolveStarted(false); + return; + } + if (started) { + abortController.abort(new Error("worker supervisor lifetime ended")); + } + }; + process.on("message", onMessage); + process.once("disconnect", onDisconnect); + return { + started: startedPromise, + signal: abortController.signal, + terminateOwnedTree: () => { + signalProcessTree(process.pid, "SIGKILL", { + detached: process.platform !== "win32", + }); + }, + dispose: () => { + if (disposed) { + return; + } + disposed = true; + process.off("message", onMessage); + process.off("disconnect", onDisconnect); + if (process.connected) { + try { + process.disconnect?.(); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ERR_IPC_DISCONNECTED") { + throw error; + } + } + } + }, + }; +} /** Register the restricted cloud worker runtime entry point. */ export function registerWorkerCli(program: Command): void { program .command("worker") .description("Run the restricted cloud worker runtime") - .action(async () => { + .addOption(new Option("--internal-worker-ipc").hideHelp()) + .action(async (options: { internalWorkerIpc?: boolean }) => { const { runWorkerCommand } = await import("../worker/worker-command.runtime.js"); await runWorkerCommand({ input: process.stdin, + ...(options.internalWorkerIpc ? { lifetime: createWorkerIpcLifetime() } : {}), output: process.stdout, }); }); diff --git a/src/node-host/node-worker-environment.ts b/src/node-host/node-worker-environment.ts new file mode 100644 index 000000000000..777b2ae2894d --- /dev/null +++ b/src/node-host/node-worker-environment.ts @@ -0,0 +1,51 @@ +const POSIX_WORKER_ENV_KEYS = new Set([ + "PATH", + "HOME", + "TMPDIR", + "TMP", + "TEMP", + "LANG", + "LANGUAGE", + "TZ", + "NODE_EXTRA_CA_CERTS", + "NODE_USE_SYSTEM_CA", + "OPENCLAW_ALLOW_INSECURE_PRIVATE_WS", +]); +const WINDOWS_WORKER_ENV_KEYS = new Set([ + ...POSIX_WORKER_ENV_KEYS, + "USERPROFILE", + "HOMEDRIVE", + "HOMEPATH", + "SYSTEMROOT", + "WINDIR", + "COMSPEC", + "PATHEXT", +]); + +/** Freeze the minimal non-secret environment inherited by node-host workers. */ +export function snapshotNodeWorkerEnv(source: NodeJS.ProcessEnv): NodeJS.ProcessEnv { + const windows = process.platform === "win32"; + const snapshot: NodeJS.ProcessEnv = {}; + const retainedWindowsKeys = new Map(); + for (const [key, value] of Object.entries(source)) { + if (value === undefined) { + continue; + } + const normalized = windows ? key.toUpperCase() : key; + const allowed = + (windows ? WINDOWS_WORKER_ENV_KEYS : POSIX_WORKER_ENV_KEYS).has(normalized) || + normalized.startsWith("LC_"); + if (!allowed) { + continue; + } + if (windows) { + const previousKey = retainedWindowsKeys.get(normalized); + if (previousKey) { + delete snapshot[previousKey]; + } + retainedWindowsKeys.set(normalized, key); + } + snapshot[key] = value; + } + return snapshot; +} diff --git a/src/node-host/node-worker-launch-store.ts b/src/node-host/node-worker-launch-store.ts new file mode 100644 index 000000000000..0ae63e902927 --- /dev/null +++ b/src/node-host/node-worker-launch-store.ts @@ -0,0 +1,434 @@ +import type { DatabaseSync } from "node:sqlite"; +import type { Selectable } from "kysely"; +import { + executeSqliteQuerySync, + executeSqliteQueryTakeFirstSync, + getNodeSqliteKysely, +} from "../infra/kysely-sync.js"; +import type { DB as OpenClawStateDatabase } from "../state/openclaw-state-db.generated.js"; +import { + runOpenClawStateWriteTransaction, + type OpenClawStateDatabaseOptions, +} from "../state/openclaw-state-db.js"; +import { OPENCLAW_STATE_SCHEMA_SQL } from "../state/openclaw-state-schema.js"; +import { + inspectNodeWorkerProcessIdentity, + type NodeWorkerProcessIdentity, +} from "./node-worker-process-identity.js"; + +type NodeWorkerLaunchState = + | "pending" + | "running" + | "completed" + | "failed" + | "interrupted" + | "cancelled"; +export type NodeWorkerTerminalState = Exclude; + +type NodeWorkerLaunchDatabase = Pick; +type NodeWorkerLaunchRow = Selectable; + +export type NodeWorkerLaunchReceipt = { + launchId: string; + planHash: string; + gatewayNamespace: string; + environmentId: string; + sessionId: string; + ownerEpoch: number; + placementGeneration: number; + runId: string; + state: NodeWorkerLaunchState; + supervisor: NodeWorkerProcessIdentity; + worker: NodeWorkerProcessIdentity | null; + resultJson: string | null; + errorText: string | null; + completedAtMs: number | null; + createdAtMs: number; + updatedAtMs: number; +}; + +type NodeWorkerLaunchClaim = Pick< + NodeWorkerLaunchReceipt, + | "environmentId" + | "gatewayNamespace" + | "launchId" + | "ownerEpoch" + | "placementGeneration" + | "planHash" + | "runId" + | "sessionId" +>; + +type NodeWorkerLaunchClaimResult = { + action: "start" | "replay" | "recover"; + receipt: NodeWorkerLaunchReceipt; +}; + +const NODE_WORKER_LAUNCH_SCHEMA_START = "CREATE TABLE IF NOT EXISTS node_worker_launches ("; +const NODE_WORKER_LAUNCH_SCHEMA_END = "\n) STRICT;"; +const initializedDatabases = new WeakSet(); +const TERMINAL_STATES: ReadonlySet = new Set([ + "completed", + "failed", + "interrupted", + "cancelled", +]); + +function ensureNodeWorkerLaunchSchema(database: DatabaseSync): void { + const start = OPENCLAW_STATE_SCHEMA_SQL.indexOf(NODE_WORKER_LAUNCH_SCHEMA_START); + const end = + start >= 0 ? OPENCLAW_STATE_SCHEMA_SQL.indexOf(NODE_WORKER_LAUNCH_SCHEMA_END, start) : -1; + if (start < 0 || end < start) { + throw new Error("OpenClaw node worker launch schema marker is missing."); + } + database.exec(OPENCLAW_STATE_SCHEMA_SQL.slice(start, end + NODE_WORKER_LAUNCH_SCHEMA_END.length)); // sqlite-allow-raw -- Canonical feature-local additive DDL only. +} + +function query(database: DatabaseSync) { + return getNodeSqliteKysely(database); +} + +function readRow(database: DatabaseSync, launchId: string): NodeWorkerLaunchRow | undefined { + return executeSqliteQueryTakeFirstSync( + database, + query(database) + .selectFrom("node_worker_launches") + .selectAll() + .where("launch_id", "=", launchId), + ); +} + +function processIdentity(pid: number, startTime: number): NodeWorkerProcessIdentity { + return { pid, startTime }; +} + +function receiptFromRow(row: NodeWorkerLaunchRow): NodeWorkerLaunchReceipt { + if (!isNodeWorkerLaunchState(row.state)) { + throw new Error(`invalid node worker launch state ${row.state}`); + } + return { + launchId: row.launch_id, + planHash: row.plan_hash, + gatewayNamespace: row.gateway_namespace, + environmentId: row.environment_id, + sessionId: row.session_id, + ownerEpoch: row.owner_epoch, + placementGeneration: row.placement_generation, + runId: row.run_id, + state: row.state, + supervisor: processIdentity(row.supervisor_pid, row.supervisor_start_time), + worker: + row.worker_pid === null || row.worker_start_time === null + ? null + : processIdentity(row.worker_pid, row.worker_start_time), + resultJson: row.result_json, + errorText: row.error_text, + completedAtMs: row.completed_at_ms, + createdAtMs: row.created_at_ms, + updatedAtMs: row.updated_at_ms, + }; +} + +function isNodeWorkerLaunchState(value: string): value is NodeWorkerLaunchState { + return value === "pending" || value === "running" || TERMINAL_STATES.has(value); +} + +function validateIdentifier(value: string, label: string): void { + if (!value || value.trim() !== value || value.length > 256 || value.includes("\0")) { + throw new Error(`${label} must be a bounded non-empty identifier`); + } +} + +function validatePlanHash(value: string): void { + if (!/^[a-f0-9]{64}$/u.test(value)) { + throw new Error("node worker plan hash must be 64 lowercase hexadecimal characters"); + } +} + +function validateTimestamp(value: number): void { + if (!Number.isSafeInteger(value) || value < 0) { + throw new Error("node worker launch timestamp must be a non-negative safe integer"); + } +} + +function validateProcessIdentity(identity: NodeWorkerProcessIdentity): void { + if ( + !Number.isSafeInteger(identity.pid) || + identity.pid <= 0 || + identity.pid > 2_147_483_647 || + !Number.isSafeInteger(identity.startTime) || + identity.startTime < 0 + ) { + throw new Error("node worker process identity must contain a bounded pid and start time"); + } +} + +function requireMatchingRow( + database: DatabaseSync, + launchId: string, + planHash: string, +): NodeWorkerLaunchRow { + const row = readRow(database, launchId); + if (!row) { + throw new Error(`node worker launch ${launchId} does not exist`); + } + if (row.plan_hash !== planHash) { + throw new Error(`node worker launch ${launchId} was replayed with a different plan`); + } + return row; +} + +function rowHasSupervisor(row: NodeWorkerLaunchRow, identity: NodeWorkerProcessIdentity): boolean { + return row.supervisor_pid === identity.pid && row.supervisor_start_time === identity.startTime; +} + +function rowHasWorker( + row: NodeWorkerLaunchRow, + identity: NodeWorkerProcessIdentity | null, +): boolean { + return identity === null + ? row.worker_pid === null && row.worker_start_time === null + : row.worker_pid === identity.pid && row.worker_start_time === identity.startTime; +} + +function sameObservedOwner(current: NodeWorkerLaunchRow, observed: NodeWorkerLaunchRow): boolean { + return ( + current.state === observed.state && + current.supervisor_pid === observed.supervisor_pid && + current.supervisor_start_time === observed.supervisor_start_time && + current.worker_pid === observed.worker_pid && + current.worker_start_time === observed.worker_start_time + ); +} + +/** Synchronous shared-state owner for durable node worker launch supervision. */ +export class NodeWorkerLaunchStore { + private readonly databaseOptions: OpenClawStateDatabaseOptions; + + constructor(options: { env?: NodeJS.ProcessEnv } = {}) { + this.databaseOptions = options.env ? { env: options.env } : {}; + } + + private write(operationLabel: string, operation: (database: DatabaseSync) => T): T { + let initializedDatabase: DatabaseSync | undefined; + const result = runOpenClawStateWriteTransaction( + ({ db }) => { + if (!initializedDatabases.has(db)) { + ensureNodeWorkerLaunchSchema(db); + initializedDatabase = db; + } + return operation(db); + }, + this.databaseOptions, + { operationLabel }, + ); + if (initializedDatabase) { + initializedDatabases.add(initializedDatabase); + } + return result; + } + + claim( + claim: NodeWorkerLaunchClaim, + supervisor: NodeWorkerProcessIdentity, + nowMs = Date.now(), + ): NodeWorkerLaunchClaimResult { + validateIdentifier(claim.launchId, "node worker launch id"); + validatePlanHash(claim.planHash); + validateTimestamp(nowMs); + validateProcessIdentity(supervisor); + + // Process inspection is intentionally outside SQLite. The second transaction + // re-reads the exact owner tuple before an adoption or recovery decision. + const observed = this.write("node-worker-launch.claim-inspect", (database) => + readRow(database, claim.launchId), + ); + if (observed && observed.plan_hash !== claim.planHash) { + throw new Error(`node worker launch ${claim.launchId} was replayed with a different plan`); + } + const observedSupervisorState = observed + ? inspectNodeWorkerProcessIdentity( + processIdentity(observed.supervisor_pid, observed.supervisor_start_time), + ) + : undefined; + + return this.write("node-worker-launch.claim", (database) => { + let current = readRow(database, claim.launchId); + if (!current) { + executeSqliteQuerySync( + database, + query(database).insertInto("node_worker_launches").values({ + launch_id: claim.launchId, + plan_hash: claim.planHash, + gateway_namespace: claim.gatewayNamespace, + environment_id: claim.environmentId, + session_id: claim.sessionId, + owner_epoch: claim.ownerEpoch, + placement_generation: claim.placementGeneration, + run_id: claim.runId, + state: "pending", + supervisor_pid: supervisor.pid, + supervisor_start_time: supervisor.startTime, + worker_pid: null, + worker_start_time: null, + result_json: null, + error_text: null, + completed_at_ms: null, + created_at_ms: nowMs, + updated_at_ms: nowMs, + }), + ); + return { + action: "start", + receipt: receiptFromRow(requireMatchingRow(database, claim.launchId, claim.planHash)), + }; + } + if (current.plan_hash !== claim.planHash) { + throw new Error(`node worker launch ${claim.launchId} was replayed with a different plan`); + } + const previousOwnerDefinitelyStale = + observedSupervisorState === "dead" || observedSupervisorState === "reused"; + if ( + current.state === "pending" && + observed && + sameObservedOwner(current, observed) && + previousOwnerDefinitelyStale + ) { + const updatedAtMs = Math.max(nowMs, current.created_at_ms, current.updated_at_ms); + executeSqliteQuerySync( + database, + query(database) + .updateTable("node_worker_launches") + .set({ + supervisor_pid: supervisor.pid, + supervisor_start_time: supervisor.startTime, + updated_at_ms: updatedAtMs, + }) + .where("launch_id", "=", claim.launchId) + .where("plan_hash", "=", claim.planHash) + .where("state", "=", "pending") + .where("supervisor_pid", "=", observed.supervisor_pid) + .where("supervisor_start_time", "=", observed.supervisor_start_time) + .where("worker_pid", "is", null) + .where("worker_start_time", "is", null), + ); + current = requireMatchingRow(database, claim.launchId, claim.planHash); + return { + action: rowHasSupervisor(current, supervisor) ? "start" : "replay", + receipt: receiptFromRow(current), + }; + } + if ( + current.state === "running" && + observed && + sameObservedOwner(current, observed) && + previousOwnerDefinitelyStale + ) { + return { action: "recover", receipt: receiptFromRow(current) }; + } + return { action: "replay", receipt: receiptFromRow(current) }; + }); + } + + get(launchId: string): NodeWorkerLaunchReceipt | undefined { + validateIdentifier(launchId, "node worker launch id"); + return this.write("node-worker-launch.get", (database) => { + const row = readRow(database, launchId); + return row ? receiptFromRow(row) : undefined; + }); + } + + markRunning(params: { + launchId: string; + planHash: string; + supervisor: NodeWorkerProcessIdentity; + worker: NodeWorkerProcessIdentity; + nowMs?: number; + }): NodeWorkerLaunchReceipt { + const nowMs = params.nowMs ?? Date.now(); + validateTimestamp(nowMs); + validateProcessIdentity(params.supervisor); + validateProcessIdentity(params.worker); + return this.write("node-worker-launch.mark-running", (database) => { + const current = requireMatchingRow(database, params.launchId, params.planHash); + if (TERMINAL_STATES.has(current.state)) { + return receiptFromRow(current); + } + if (current.state === "running") { + return receiptFromRow(current); + } + if (!rowHasSupervisor(current, params.supervisor) || !rowHasWorker(current, null)) { + return receiptFromRow(current); + } + const updatedAtMs = Math.max(nowMs, current.created_at_ms, current.updated_at_ms); + executeSqliteQuerySync( + database, + query(database) + .updateTable("node_worker_launches") + .set({ + state: "running", + worker_pid: params.worker.pid, + worker_start_time: params.worker.startTime, + updated_at_ms: updatedAtMs, + }) + .where("launch_id", "=", params.launchId) + .where("plan_hash", "=", params.planHash) + .where("state", "=", "pending") + .where("supervisor_pid", "=", params.supervisor.pid) + .where("supervisor_start_time", "=", params.supervisor.startTime) + .where("worker_pid", "is", null) + .where("worker_start_time", "is", null), + ); + return receiptFromRow(requireMatchingRow(database, params.launchId, params.planHash)); + }); + } + + finish(params: { + launchId: string; + planHash: string; + supervisor: NodeWorkerProcessIdentity; + worker: NodeWorkerProcessIdentity | null; + state: NodeWorkerTerminalState; + resultJson?: string; + errorText?: string; + nowMs?: number; + }): NodeWorkerLaunchReceipt { + const nowMs = params.nowMs ?? Date.now(); + validateTimestamp(nowMs); + validateProcessIdentity(params.supervisor); + if (params.worker) { + validateProcessIdentity(params.worker); + } + return this.write("node-worker-launch.finish", (database) => { + const current = requireMatchingRow(database, params.launchId, params.planHash); + if (TERMINAL_STATES.has(current.state)) { + return receiptFromRow(current); + } + if (!rowHasSupervisor(current, params.supervisor) || !rowHasWorker(current, params.worker)) { + return receiptFromRow(current); + } + const completedAtMs = Math.max(nowMs, current.created_at_ms, current.updated_at_ms); + let update = query(database) + .updateTable("node_worker_launches") + .set({ + state: params.state, + result_json: params.state === "completed" ? (params.resultJson ?? null) : null, + error_text: params.state === "completed" ? null : (params.errorText ?? null), + completed_at_ms: completedAtMs, + updated_at_ms: completedAtMs, + }) + .where("launch_id", "=", params.launchId) + .where("plan_hash", "=", params.planHash) + .where("state", "in", ["pending", "running"]) + .where("supervisor_pid", "=", params.supervisor.pid) + .where("supervisor_start_time", "=", params.supervisor.startTime); + update = params.worker + ? update + .where("worker_pid", "=", params.worker.pid) + .where("worker_start_time", "=", params.worker.startTime) + : update.where("worker_pid", "is", null).where("worker_start_time", "is", null); + executeSqliteQuerySync(database, update); + return receiptFromRow(requireMatchingRow(database, params.launchId, params.planHash)); + }); + } +} diff --git a/src/node-host/node-worker-process-identity.ts b/src/node-host/node-worker-process-identity.ts new file mode 100644 index 000000000000..c69b08c033d7 --- /dev/null +++ b/src/node-host/node-worker-process-identity.ts @@ -0,0 +1,36 @@ +import { readWindowsProcessStartTimeSync } from "../infra/windows-port-pids.js"; +import { getFileLockProcessStartTime, isPidDefinitelyDead } from "../shared/pid-alive.js"; + +export type NodeWorkerProcessIdentity = { + pid: number; + startTime: number; +}; + +type NodeWorkerProcessIdentityState = "live" | "dead" | "reused" | "unknown"; + +function readNodeWorkerProcessStartTime(pid: number): number | null { + return process.platform === "win32" + ? readWindowsProcessStartTimeSync(pid) + : getFileLockProcessStartTime(pid); +} + +export function requireNodeWorkerProcessIdentity(pid: number): NodeWorkerProcessIdentity { + const startTime = readNodeWorkerProcessStartTime(pid); + if (startTime === null) { + throw new Error(`cannot establish PID-reuse-safe identity for process ${pid}`); + } + return { pid, startTime }; +} + +export function inspectNodeWorkerProcessIdentity( + identity: NodeWorkerProcessIdentity, +): NodeWorkerProcessIdentityState { + const observedStartTime = readNodeWorkerProcessStartTime(identity.pid); + if (observedStartTime !== null) { + if (observedStartTime !== identity.startTime) { + return "reused"; + } + return isPidDefinitelyDead(identity.pid) ? "dead" : "live"; + } + return isPidDefinitelyDead(identity.pid) ? "dead" : "unknown"; +} diff --git a/src/node-host/node-worker-supervisor.recovery.test.ts b/src/node-host/node-worker-supervisor.recovery.test.ts new file mode 100644 index 000000000000..d31dff1e78ae --- /dev/null +++ b/src/node-host/node-worker-supervisor.recovery.test.ts @@ -0,0 +1,362 @@ +import { spawn, type ChildProcess } from "node:child_process"; +import { createHash } from "node:crypto"; +import fs from "node:fs"; +import path from "node:path"; +import { pathToFileURL } from "node:url"; +import { stableStringify } from "@openclaw/normalization-core"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; +import { + closeOpenClawStateDatabaseForTest, + openOpenClawStateDatabase, +} from "../state/openclaw-state-db.js"; +import type { NodeWorkerLaunchReceipt } from "./node-worker-launch-store.js"; +import { + inspectNodeWorkerProcessIdentity, + requireNodeWorkerProcessIdentity, + type NodeWorkerProcessIdentity, +} from "./node-worker-process-identity.js"; +import { createNodeWorkerSupervisor } from "./node-worker-supervisor.js"; +import { + testWorkerLaunchInput, + writeNodeWorkerFixture, +} from "./node-worker-supervisor.test-support.js"; + +const tempDirs = useAutoCleanupTempDirTracker(afterEach); +const spawned = new Set(); +const ownedProcessGroups: NodeWorkerProcessIdentity[] = []; + +afterEach(async () => { + for (const child of spawned) { + if (child.exitCode === null && child.signalCode === null) { + child.kill("SIGKILL"); + } + } + if (process.platform !== "win32") { + for (const identity of ownedProcessGroups) { + if (inspectNodeWorkerProcessIdentity(identity) === "reused") { + continue; + } + try { + process.kill(-identity.pid, "SIGKILL"); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ESRCH") { + throw error; + } + } + } + } + spawned.clear(); + ownedProcessGroups.length = 0; + closeOpenClawStateDatabaseForTest(); +}); + +function fixture(label: string) { + return writeNodeWorkerFixture(tempDirs.make(label)); +} + +function planHash(input: ReturnType): string { + return createHash("sha256") + .update( + stableStringify({ + bundleHash: input.bundleHash, + descriptor: input.descriptor, + gatewayNamespace: input.gatewayNamespace, + placementGeneration: input.placementGeneration, + }), + ) + .digest("hex"); +} + +function insertLaunch(params: { + env: NodeJS.ProcessEnv; + input: ReturnType; + state: "pending" | "running"; + supervisor: NodeWorkerProcessIdentity; + worker?: NodeWorkerProcessIdentity; +}) { + const database = openOpenClawStateDatabase({ env: params.env }).db; + database + .prepare( + `INSERT INTO node_worker_launches ( + launch_id, plan_hash, gateway_namespace, environment_id, session_id, + owner_epoch, placement_generation, run_id, state, + supervisor_pid, supervisor_start_time, worker_pid, worker_start_time, + result_json, error_text, completed_at_ms, created_at_ms, updated_at_ms + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, NULL, NULL, 1, 1)`, + ) + .run( + params.input.launchId, + planHash(params.input), + params.input.gatewayNamespace, + params.input.descriptor.admission.environmentId, + params.input.descriptor.admission.sessionId, + params.input.descriptor.admission.ownerEpoch, + params.input.placementGeneration, + params.input.descriptor.assignment.runId, + params.state, + params.supervisor.pid, + params.supervisor.startTime, + params.worker?.pid ?? null, + params.worker?.startTime ?? null, + ); +} + +function waitForChildLine(child: ChildProcess): Promise { + return new Promise((resolve, reject) => { + let stdout = ""; + let stderr = ""; + const onData = (chunk: Buffer) => { + stdout += chunk.toString("utf8"); + const newline = stdout.indexOf("\n"); + if (newline >= 0) { + resolve(stdout.slice(0, newline)); + } + }; + child.stdout?.on("data", onData); + child.stderr?.on("data", (chunk: Buffer) => { + stderr += chunk.toString("utf8"); + }); + child.once("error", reject); + child.once("close", (code, signal) => { + reject(new Error(`owner exited before ready (${code ?? signal}): ${stderr}`)); + }); + }); +} + +function waitForChildExit(child: ChildProcess): Promise { + if (child.exitCode !== null || child.signalCode !== null) { + return Promise.resolve(); + } + return new Promise((resolve, reject) => { + child.once("error", reject); + child.once("close", () => resolve()); + }); +} + +function writeSupervisorOwnerScript(root: string): string { + const supervisorUrl = pathToFileURL(path.resolve("src/node-host/node-worker-supervisor.ts")).href; + const scriptPath = path.join(root, "supervisor-owner.mts"); + fs.writeFileSync( + scriptPath, + ` + import fs from "node:fs"; + import { createNodeWorkerSupervisor } from ${JSON.stringify(supervisorUrl)}; + const [bundleRoot, stateDir, inputPath] = process.argv.slice(2); + const supervisor = createNodeWorkerSupervisor({ + bundleRoot, + env: { ...process.env, OPENCLAW_STATE_DIR: stateDir }, + }); + const shutdown = async () => { + await supervisor.close(); + process.exit(0); + }; + process.once("SIGTERM", () => void shutdown()); + const input = JSON.parse(fs.readFileSync(inputPath, "utf8")); + const receipt = await supervisor.launch(input); + process.stdout.write(JSON.stringify(receipt) + "\\n"); + setInterval(() => {}, 1000); + `, + ); + return scriptPath; +} + +function spawnSupervisorOwner(params: { + bundleRoot: string; + env: NodeJS.ProcessEnv; + input: ReturnType; + root: string; +}): ChildProcess { + const inputPath = path.join(params.root, `${params.input.launchId}.json`); + fs.writeFileSync(inputPath, JSON.stringify(params.input)); + const child = spawn( + process.execPath, + [ + "--import", + "tsx", + writeSupervisorOwnerScript(params.root), + params.bundleRoot, + params.env.OPENCLAW_STATE_DIR!, + inputPath, + ], + { stdio: ["ignore", "pipe", "pipe"] }, + ); + spawned.add(child); + return child; +} + +async function waitForIdentityDeath(identity: NodeWorkerProcessIdentity) { + await vi.waitFor(() => expect(inspectNodeWorkerProcessIdentity(identity)).not.toBe("live"), { + timeout: 5_000, + }); +} + +describe("node worker supervisor recovery", () => { + it("atomically adopts pending work only after the previous supervisor is stale", async () => { + const { bundleRoot, env, workspaceDir } = fixture("node-worker-stale-pending-"); + const supervisor = createNodeWorkerSupervisor({ bundleRoot, env }); + await supervisor.status("schema-probe"); + const input = testWorkerLaunchInput(workspaceDir, "stale-pending-launch"); + insertLaunch({ + env, + input, + state: "pending", + supervisor: { pid: 2_147_483_647, startTime: 1 }, + }); + + const running = await supervisor.launch(input); + + expect(running).toMatchObject({ + state: "running", + supervisor: requireNodeWorkerProcessIdentity(process.pid), + worker: { pid: expect.any(Number), startTime: expect.any(Number) }, + }); + await supervisor.close(); + }); + + it.runIf(process.platform !== "win32")( + "kills the exact stale-owner worker group before marking it interrupted", + async () => { + const { bundleRoot, env, root, workspaceDir } = fixture("node-worker-stale-running-"); + const marker = path.join(root, "recovery-grandchild.pid"); + const workerSource = ` + const { spawn } = require("node:child_process"); + const fs = require("node:fs"); + const child = spawn(process.execPath, ["-e", "setInterval(() => {}, 1000)"], { stdio: "ignore" }); + fs.writeFileSync(process.argv[1], String(child.pid)); + setInterval(() => {}, 1000); + `; + const workerProcess = spawn(process.execPath, ["-e", workerSource, marker], { + detached: true, + stdio: "ignore", + }); + spawned.add(workerProcess); + const worker = requireNodeWorkerProcessIdentity(workerProcess.pid!); + ownedProcessGroups.push(worker); + await vi.waitFor(() => expect(fs.existsSync(marker)).toBe(true)); + const grandchild = requireNodeWorkerProcessIdentity(Number(fs.readFileSync(marker, "utf8"))); + const input = testWorkerLaunchInput(workspaceDir, "stale-running-launch", "wait"); + const supervisor = createNodeWorkerSupervisor({ bundleRoot, env }); + await supervisor.status("schema-probe"); + insertLaunch({ + env, + input, + state: "running", + supervisor: { pid: 2_147_483_647, startTime: 1 }, + worker, + }); + + const recovered = await supervisor.launch(input); + + expect(recovered).toMatchObject({ state: "interrupted", worker }); + await waitForIdentityDeath(worker); + await waitForIdentityDeath(grandchild); + expect((await supervisor.status(input.launchId))?.worker).toEqual(worker); + await supervisor.close(); + }, + ); + + it("returns a live foreign running receipt from a real second process without mutation", async () => { + const { bundleRoot, env, root, workspaceDir } = fixture("node-worker-live-replay-"); + const input = testWorkerLaunchInput(workspaceDir, "live-running-launch", "wait"); + const owner = spawnSupervisorOwner({ bundleRoot, env, input, root }); + const owned = JSON.parse(await waitForChildLine(owner)) as NodeWorkerLaunchReceipt; + if (owned.worker) { + ownedProcessGroups.push(owned.worker); + } + const second = createNodeWorkerSupervisor({ bundleRoot, env }); + + const replay = await second.launch(input); + + expect(replay).toEqual(owned); + expect(inspectNodeWorkerProcessIdentity(owned.supervisor)).toBe("live"); + expect(inspectNodeWorkerProcessIdentity(owned.worker!)).toBe("live"); + owner.kill("SIGTERM"); + await waitForChildExit(owner); + await second.close(); + }); + + it.runIf(process.platform !== "win32")( + "uses IPC disconnect after owner SIGKILL, then reconciles only after exact tree death", + async () => { + const { bundleRoot, env, root, workspaceDir } = fixture("node-worker-owner-kill-"); + const input = testWorkerLaunchInput(workspaceDir, "owner-kill-launch", "tree"); + const owner = spawnSupervisorOwner({ bundleRoot, env, input, root }); + const owned = JSON.parse(await waitForChildLine(owner)) as NodeWorkerLaunchReceipt; + ownedProcessGroups.push(owned.worker!); + const grandchildPath = path.join(workspaceDir, "grandchild.pid"); + await vi.waitFor(() => expect(fs.existsSync(grandchildPath)).toBe(true)); + const grandchild = requireNodeWorkerProcessIdentity( + Number(fs.readFileSync(grandchildPath, "utf8")), + ); + + owner.kill("SIGKILL"); + await waitForChildExit(owner); + await waitForIdentityDeath(owned.supervisor); + await waitForIdentityDeath(owned.worker!); + await waitForIdentityDeath(grandchild); + + const restarted = createNodeWorkerSupervisor({ bundleRoot, env }); + const reconciled = await restarted.launch(input); + expect(reconciled).toMatchObject({ + state: "interrupted", + supervisor: owned.supervisor, + worker: owned.worker, + }); + await restarted.close(); + }, + ); + + it("keeps a live foreign pending claim unchanged across real processes", async () => { + const { bundleRoot, env, root, workspaceDir } = fixture("node-worker-live-pending-"); + const input = testWorkerLaunchInput(workspaceDir, "live-pending-launch", "wait"); + const claim = { + launchId: input.launchId, + planHash: planHash(input), + gatewayNamespace: input.gatewayNamespace, + environmentId: input.descriptor.admission.environmentId, + sessionId: input.descriptor.admission.sessionId, + ownerEpoch: input.descriptor.admission.ownerEpoch, + placementGeneration: input.placementGeneration, + runId: input.descriptor.assignment.runId, + }; + const storeUrl = pathToFileURL(path.resolve("src/node-host/node-worker-launch-store.ts")).href; + const identityUrl = pathToFileURL( + path.resolve("src/node-host/node-worker-process-identity.ts"), + ).href; + const claimPath = path.join(root, "claim.json"); + const scriptPath = path.join(root, "pending-owner.mts"); + fs.writeFileSync(claimPath, JSON.stringify(claim)); + fs.writeFileSync( + scriptPath, + ` + import fs from "node:fs"; + import { NodeWorkerLaunchStore } from ${JSON.stringify(storeUrl)}; + import { requireNodeWorkerProcessIdentity } from ${JSON.stringify(identityUrl)}; + const [stateDir, claimPath] = process.argv.slice(2); + const store = new NodeWorkerLaunchStore({ env: { ...process.env, OPENCLAW_STATE_DIR: stateDir } }); + const result = store.claim( + JSON.parse(fs.readFileSync(claimPath, "utf8")), + requireNodeWorkerProcessIdentity(process.pid), + ); + process.stdout.write(JSON.stringify(result.receipt) + "\\n"); + setInterval(() => {}, 1000); + `, + ); + const owner = spawn( + process.execPath, + ["--import", "tsx", scriptPath, env.OPENCLAW_STATE_DIR!, claimPath], + { stdio: ["ignore", "pipe", "pipe"] }, + ); + spawned.add(owner); + const owned = JSON.parse(await waitForChildLine(owner)) as NodeWorkerLaunchReceipt; + const second = createNodeWorkerSupervisor({ bundleRoot, env }); + + const replay = await second.launch(input); + + expect(replay).toEqual(owned); + owner.kill("SIGKILL"); + await waitForChildExit(owner); + await second.close(); + }); +}); diff --git a/src/node-host/node-worker-supervisor.test-support.ts b/src/node-host/node-worker-supervisor.test-support.ts new file mode 100644 index 000000000000..0124b4d0f3a4 --- /dev/null +++ b/src/node-host/node-worker-supervisor.test-support.ts @@ -0,0 +1,184 @@ +import fs from "node:fs"; +import path from "node:path"; +import { + WORKER_PROTOCOL_FEATURES, + WORKER_RPC_SET_VERSION, +} from "../../packages/gateway-protocol/src/schema/worker-admission.js"; +import type { WorkerLaunchDescriptor } from "../worker/launch-descriptor.js"; + +const TEST_BUNDLE_HASH = "a".repeat(64); +export const TEST_WORKER_CREDENTIAL = 'node worker/"credential\\secret?'; + +export const TEST_WORKER_SOURCE = String.raw` +import fs from "node:fs"; +import path from "node:path"; +import { spawn } from "node:child_process"; +let input = ""; +for await (const chunk of process.stdin) input += chunk; +const descriptor = JSON.parse(input); +if (descriptor.assignment.prompt === "exit-before-start") { + fs.writeFileSync(path.join(descriptor.assignment.workspaceDir, "prestart-exited"), "exited"); + process.exit(23); +} +if (!process.connected || !process.channel || !process.argv.includes("--internal-worker-ipc")) { + process.exit(24); +} +let grandchild; +let disposed = false; +let started = false; +let resolveStart; +const start = new Promise((resolve) => { resolveStart = resolve; }); +const hardTerminate = () => { + if (process.platform === "win32") { + spawn("taskkill", ["/F", "/T", "/PID", String(process.pid)], { + detached: true, + stdio: "ignore", + windowsHide: true, + }); + return; + } + process.kill(-process.pid, "SIGKILL"); +}; +const onMessage = (message) => { + if ( + started || + typeof message !== "object" || + message === null || + Array.isArray(message) || + Object.keys(message).length !== 1 || + message.type !== "openclaw-worker-start-v1" + ) { + hardTerminate(); + return; + } + started = true; + resolveStart(); +}; +const onDisconnect = () => { + if (disposed) return; + if (!started) process.exit(0); + hardTerminate(); +}; +process.on("message", onMessage); +process.once("disconnect", onDisconnect); +await start; +const exitWorker = (code) => { + disposed = true; + process.off("message", onMessage); + process.off("disconnect", onDisconnect); + if (process.connected) process.disconnect(); + process.exit(code); +}; +const writeResultAndExit = (value) => { + fs.writeSync(1, value); + exitWorker(0); +}; +const mode = descriptor.assignment.prompt; +if (mode === "wait") { + setInterval(() => {}, 1000); +} else if (mode === "tree") { + grandchild = spawn(process.execPath, ["-e", "setInterval(() => {}, 1000)"], { stdio: "ignore" }); + fs.writeFileSync(path.join(descriptor.assignment.workspaceDir, "grandchild.pid"), String(grandchild.pid)); + setInterval(() => {}, 1000); +} else if (mode === "secret-fail") { + await new Promise((resolve) => setTimeout(resolve, 500)); + const credential = descriptor.admission.credential; + const escaped = JSON.stringify(credential).slice(1, -1); + process.stderr.write( + "failure " + "x".repeat(5000) + " " + credential + " " + encodeURIComponent(credential) + " " + escaped, + ); + exitWorker(7); +} else if (mode.startsWith("secret-cutoff-")) { + const credential = descriptor.admission.credential; + const representations = { + "secret-cutoff-raw": credential, + "secret-cutoff-url": encodeURIComponent(credential), + "secret-cutoff-json": JSON.stringify(credential).slice(1, -1), + }; + const representation = representations[mode]; + const suffixBytes = 4096 - Math.floor(Buffer.byteLength(representation, "utf8") / 2); + process.stderr.write("x".repeat(5000) + representation + "y".repeat(suffixBytes)); + exitWorker(7); +} else if (mode === "secret-success") { + await new Promise((resolve) => setTimeout(resolve, 500)); + const credential = descriptor.admission.credential; + writeResultAndExit( + JSON.stringify({ raw: credential, encoded: encodeURIComponent(credential), status: "completed" }) + "\n", + ); +} else if (mode === "overflow") { + writeResultAndExit("x".repeat(70 * 1024)); +} else if (mode === "fast-terminal") { + const marker = path.join(descriptor.assignment.workspaceDir, "fast-terminal-marker"); + process.once("SIGTERM", () => { + fs.writeFileSync(marker, "signal"); + process.exit(143); + }); + await new Promise((resolve) => setTimeout(resolve, 100)); + fs.writeFileSync(marker, "normal"); + writeResultAndExit(JSON.stringify({ status: "completed" }) + "\n"); +} else if (mode === "env") { + writeResultAndExit(JSON.stringify(process.env) + "\n"); +} else { + await new Promise((resolve) => setTimeout(resolve, 25)); + writeResultAndExit(JSON.stringify({ argv: process.argv.slice(2), status: "completed" }) + "\n"); +} +`; + +export function testWorkerDescriptor( + workspaceDir: string, + prompt = "success", +): WorkerLaunchDescriptor { + return { + version: 3, + connectionEndpoint: { kind: "unix", socketPath: "/tmp/openclaw-worker/gateway.sock" }, + admission: { + environmentId: "environment-1", + credential: TEST_WORKER_CREDENTIAL, + sessionId: "session-1", + ownerEpoch: 3, + rpcSetVersion: WORKER_RPC_SET_VERSION, + handshake: { + bundleHash: TEST_BUNDLE_HASH, + openclawVersion: "2026.8.1", + protocolFeatures: [...WORKER_PROTOCOL_FEATURES], + }, + }, + assignment: { + agentId: "agent-1", + operationalRunInstance: { instanceId: "instance-1", runId: "run-1" }, + agentRuntimeIdentityToken: "signed-runtime-token", + runId: "run-1", + turnId: "turn-1", + prompt, + suppressPromptTranscript: false, + workspaceDir, + modelRef: { provider: "provider-1", model: "model-1" }, + inferenceOptions: {}, + initialMessages: [], + transcript: { baseLeafId: null, nextSeq: 1 }, + liveEvents: { ackedSeq: 0, nextSeq: 1 }, + toolAuthority: { allowedToolNames: [] }, + }, + }; +} + +export function writeNodeWorkerFixture(root: string) { + const stateDir = path.join(root, "state-root"); + const bundleRoot = path.join(root, "bundles-root"); + const workspaceDir = path.join(root, "workspace"); + const bundleDir = path.join(bundleRoot, "gateway-1", "bundles", TEST_BUNDLE_HASH); + fs.mkdirSync(bundleDir, { recursive: true }); + fs.mkdirSync(workspaceDir, { recursive: true }); + fs.writeFileSync(path.join(bundleDir, "openclaw.mjs"), TEST_WORKER_SOURCE); + return { bundleRoot, env: { OPENCLAW_STATE_DIR: stateDir }, root, stateDir, workspaceDir }; +} + +export function testWorkerLaunchInput(workspaceDir: string, launchId: string, prompt = "success") { + return { + launchId, + gatewayNamespace: "gateway-1", + bundleHash: TEST_BUNDLE_HASH, + placementGeneration: 4, + descriptor: testWorkerDescriptor(workspaceDir, prompt), + }; +} diff --git a/src/node-host/node-worker-supervisor.test.ts b/src/node-host/node-worker-supervisor.test.ts new file mode 100644 index 000000000000..a96f9b983400 --- /dev/null +++ b/src/node-host/node-worker-supervisor.test.ts @@ -0,0 +1,513 @@ +import childProcess from "node:child_process"; +import fs from "node:fs"; +import path from "node:path"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; +import { registerSecretValueForRedaction } from "../logging/secret-redaction-registry.js"; +import { resetSecretRedactionRegistryForTest } from "../logging/secret-redaction-registry.test-support.js"; +import { + closeOpenClawStateDatabaseForTest, + openOpenClawStateDatabase, +} from "../state/openclaw-state-db.js"; +import { withEnvAsync } from "../test-utils/env.js"; +import { NodeWorkerLaunchStore } from "./node-worker-launch-store.js"; +import { + inspectNodeWorkerProcessIdentity, + requireNodeWorkerProcessIdentity, +} from "./node-worker-process-identity.js"; +import { createNodeWorkerSupervisor } from "./node-worker-supervisor.js"; +import { + TEST_WORKER_CREDENTIAL, + TEST_WORKER_SOURCE, + testWorkerDescriptor, + testWorkerLaunchInput, + writeNodeWorkerFixture, +} from "./node-worker-supervisor.test-support.js"; + +type NodeWorkerSupervisor = ReturnType; + +const tempDirs = useAutoCleanupTempDirTracker(afterEach); + +afterEach(() => { + vi.restoreAllMocks(); + resetSecretRedactionRegistryForTest(); + closeOpenClawStateDatabaseForTest(); +}); + +function fixture() { + const root = tempDirs.make("node-worker-supervisor-"); + const { bundleRoot, env, stateDir, workspaceDir } = writeNodeWorkerFixture(root); + const supervisor = createNodeWorkerSupervisor({ bundleRoot, env }); + return { bundleRoot, env, root, stateDir, supervisor, workspaceDir }; +} + +function launchInput(workspaceDir: string, launchId: string, prompt = "success") { + return testWorkerLaunchInput(workspaceDir, launchId, prompt); +} + +async function waitForTerminal(supervisor: NodeWorkerSupervisor, launchId: string) { + await vi.waitFor( + async () => { + expect((await supervisor.status(launchId))?.state).not.toMatch(/^(?:pending|running)$/u); + }, + { timeout: 5_000 }, + ); + const receipt = await supervisor.status(launchId); + if (!receipt) { + throw new Error(`missing launch receipt ${launchId}`); + } + return receipt; +} + +describe("node worker supervisor", () => { + it("keeps construction and close inert without resolving process identity", async () => { + const root = tempDirs.make("node-worker-inert-"); + const { bundleRoot, env } = writeNodeWorkerFixture(root); + const originalPlatform = Object.getOwnPropertyDescriptor(process, "platform"); + const spawnSync = vi.spyOn(childProcess, "spawnSync"); + const execFileSync = vi.spyOn(childProcess, "execFileSync"); + Object.defineProperty(process, "platform", { configurable: true, value: "win32" }); + try { + const supervisor = createNodeWorkerSupervisor({ bundleRoot, env }); + await supervisor.close(); + expect(spawnSync).not.toHaveBeenCalled(); + expect(execFileSync).not.toHaveBeenCalled(); + } finally { + if (originalPlatform) { + Object.defineProperty(process, "platform", originalPlatform); + } + } + }); + + it("keeps the additive table absent until the first stateful operation", async () => { + const { bundleRoot, env, supervisor } = fixture(); + const database = openOpenClawStateDatabase({ env }); + const findTable = () => + database.db + .prepare("SELECT name FROM sqlite_schema WHERE type = 'table' AND name = ?") + .get("node_worker_launches"); + + expect(findTable()).toBeUndefined(); + await supervisor.close(); + expect(findTable()).toBeUndefined(); + + const active = createNodeWorkerSupervisor({ bundleRoot, env }); + expect(await active.status("missing-launch")).toBeUndefined(); + expect( + database.db + .prepare("SELECT strict FROM pragma_table_list WHERE name = ?") + .get("node_worker_launches"), + ).toEqual({ strict: 1 }); + await active.close(); + }); + + it("keeps pending and running launches owned by a live supervisor unchanged", async () => { + const { bundleRoot, env, supervisor } = fixture(); + await supervisor.status("schema-probe"); + const database = openOpenClawStateDatabase({ env }).db; + const supervisorIdentity = requireNodeWorkerProcessIdentity(process.pid); + const insert = database.prepare(` + INSERT INTO node_worker_launches ( + launch_id, plan_hash, gateway_namespace, environment_id, session_id, + owner_epoch, placement_generation, run_id, state, + supervisor_pid, supervisor_start_time, worker_pid, worker_start_time, + result_json, error_text, completed_at_ms, created_at_ms, updated_at_ms + ) VALUES (?, ?, 'gateway-1', 'environment-1', 'session-1', 3, 4, 'run-1', ?, ?, ?, ?, ?, NULL, NULL, NULL, 1, 1) + `); + insert.run( + "pending-launch", + "b".repeat(64), + "pending", + supervisorIdentity.pid, + supervisorIdentity.startTime, + null, + null, + ); + insert.run( + "running-launch", + "c".repeat(64), + "running", + supervisorIdentity.pid, + supervisorIdentity.startTime, + process.pid, + supervisorIdentity.startTime, + ); + + const sameHandle = createNodeWorkerSupervisor({ bundleRoot, env }); + expect(await sameHandle.status("pending-launch")).toMatchObject({ + state: "pending", + worker: null, + }); + expect(await sameHandle.status("running-launch")).toMatchObject({ + state: "running", + worker: supervisorIdentity, + }); + await supervisor.close(); + await sameHandle.close(); + closeOpenClawStateDatabaseForTest(); + + openOpenClawStateDatabase({ env }); + const recovered = createNodeWorkerSupervisor({ bundleRoot, env }); + expect(await recovered.status("pending-launch")).toMatchObject({ + state: "pending", + worker: null, + }); + expect(await recovered.status("running-launch")).toMatchObject({ + state: "running", + worker: supervisorIdentity, + }); + await recovered.close(); + }); + + it("launches idempotently and persists only bounded non-secret facts", async () => { + const { env, supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, "success-launch"); + + expect(await supervisor.launch(input)).toMatchObject({ + launchId: "success-launch", + state: "running", + environmentId: "environment-1", + sessionId: "session-1", + ownerEpoch: 3, + placementGeneration: 4, + runId: "run-1", + }); + const completed = await waitForTerminal(supervisor, input.launchId); + expect(completed).toMatchObject({ state: "completed", errorText: null }); + expect(JSON.parse(completed.resultJson ?? "null")).toEqual({ + argv: ["worker", "--internal-worker-ipc"], + status: "completed", + }); + expect(await supervisor.launch(input)).toEqual(completed); + await expect( + supervisor.launch({ + ...input, + descriptor: testWorkerDescriptor(workspaceDir, "different-plan"), + }), + ).rejects.toThrow("replayed with a different plan"); + + const row = openOpenClawStateDatabase({ env }) + .db.prepare("SELECT * FROM node_worker_launches WHERE launch_id = ?") + .get(input.launchId); + expect(JSON.stringify(row)).not.toContain(TEST_WORKER_CREDENTIAL); + await supervisor.close(); + }); + + it.each(["status", "launch", "cancel", "close"] as const)( + "retains an observed terminal outcome when %s reconciliation keeps failing", + async (operation) => { + const { env, supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, `finish-failure-${operation}`); + const store = (supervisor as unknown as { store: NodeWorkerLaunchStore }).store; + const originalFinish = store.finish.bind(store); + let persistenceUnavailable = true; + const finish = vi.spyOn(store, "finish").mockImplementation((params) => { + if (persistenceUnavailable) { + throw new Error("injected finish failure"); + } + return originalFinish(params); + }); + const invoke = async () => { + switch (operation) { + case "status": + return await supervisor.status(input.launchId); + case "launch": + return await supervisor.launch(input); + case "cancel": + return await supervisor.cancel(input.launchId); + case "close": + await supervisor.close(); + return new NodeWorkerLaunchStore({ env }).get(input.launchId); + default: + throw new Error("unsupported reconciliation operation"); + } + }; + + expect(await supervisor.launch(input)).toMatchObject({ state: "running" }); + await vi.waitFor(() => expect(finish).toHaveBeenCalled(), { timeout: 5_000 }); + expect(new NodeWorkerLaunchStore({ env }).get(input.launchId)?.state).toBe("running"); + + await expect(invoke()).rejects.toThrow("injected finish failure"); + expect(new NodeWorkerLaunchStore({ env }).get(input.launchId)?.state).toBe("running"); + + persistenceUnavailable = false; + const completed = await invoke(); + expect(completed).toMatchObject({ + state: "completed", + resultJson: expect.stringContaining('"status":"completed"'), + }); + expect(new NodeWorkerLaunchStore({ env }).get(input.launchId)?.state).toBe("completed"); + await supervisor.close(); + }, + ); + + it("spawns workers with only supplied runtime essentials", async () => { + const root = tempDirs.make("node-worker-env-"); + const { bundleRoot, env, workspaceDir } = writeNodeWorkerFixture(root); + const suppliedPathKey = process.platform === "win32" ? "Path" : "PATH"; + const suppliedEnv: NodeJS.ProcessEnv = { + ...env, + [suppliedPathKey]: process.env.PATH, + HOME: path.join(root, "worker-home"), + LANG: "en_US.UTF-8", + LC_TIME: "de_DE.UTF-8", + NODE_EXTRA_CA_CERTS: path.join(root, "private-ca.pem"), + NODE_USE_SYSTEM_CA: "1", + OPENCLAW_ALLOW_INSECURE_PRIVATE_WS: "1", + OPENCLAW_SUPPLIED_SECRET: "supplied-openclaw-secret", + NODE_OPTIONS: "--title=forbidden-worker-title", + BASH_ENV: path.join(root, "forbidden-shell-init"), + DYLD_INSERT_LIBRARIES: path.join(root, "forbidden-runtime-injection"), + HTTPS_PROXY: "http://supplied-proxy.invalid", + SUPPLIED_SECRET: "supplied-secret", + }; + + await withEnvAsync( + { + AMBIENT_SECRET: "ambient-secret", + OPENCLAW_AMBIENT_SECRET: "ambient-openclaw-secret", + HTTP_PROXY: "http://ambient-proxy.invalid", + NODE_OPTIONS: undefined, + }, + async () => { + const expectedWorkerEnv: NodeJS.ProcessEnv = { + HOME: suppliedEnv.HOME, + LANG: suppliedEnv.LANG, + LC_TIME: suppliedEnv.LC_TIME, + NODE_EXTRA_CA_CERTS: suppliedEnv.NODE_EXTRA_CA_CERTS, + NODE_USE_SYSTEM_CA: suppliedEnv.NODE_USE_SYSTEM_CA, + OPENCLAW_ALLOW_INSECURE_PRIVATE_WS: suppliedEnv.OPENCLAW_ALLOW_INSECURE_PRIVATE_WS, + [suppliedPathKey]: suppliedEnv[suppliedPathKey], + }; + const supervisor = createNodeWorkerSupervisor({ bundleRoot, env: suppliedEnv }); + suppliedEnv.HOME = path.join(root, "mutated-home"); + suppliedEnv.LANG = "mutated-locale"; + const input = launchInput(workspaceDir, "env-launch", "env"); + await supervisor.launch(input); + const completed = await waitForTerminal(supervisor, input.launchId); + const workerEnv = JSON.parse(completed.resultJson ?? "null") as Record; + + expect(workerEnv).toMatchObject(expectedWorkerEnv); + expect(workerEnv).not.toHaveProperty("AMBIENT_SECRET"); + expect(workerEnv).not.toHaveProperty("OPENCLAW_AMBIENT_SECRET"); + expect(workerEnv).not.toHaveProperty("OPENCLAW_STATE_DIR"); + expect(workerEnv).not.toHaveProperty("OPENCLAW_SUPPLIED_SECRET"); + expect(workerEnv).not.toHaveProperty("NODE_OPTIONS"); + expect(workerEnv).not.toHaveProperty("BASH_ENV"); + expect(workerEnv).not.toHaveProperty("DYLD_INSERT_LIBRARIES"); + expect(workerEnv).not.toHaveProperty("HTTP_PROXY"); + expect(workerEnv).not.toHaveProperty("HTTPS_PROXY"); + expect(workerEnv).not.toHaveProperty("SUPPLIED_SECRET"); + expect(JSON.stringify(workerEnv)).not.toContain(TEST_WORKER_CREDENTIAL); + const platformInjectedKeys = + process.platform === "darwin" ? ["__CF_USER_TEXT_ENCODING"] : []; + expect(Object.keys(workerEnv).toSorted()).toEqual( + [...Object.keys(expectedWorkerEnv), ...platformInjectedKeys] + .filter( + (key) => expectedWorkerEnv[key] !== undefined || platformInjectedKeys.includes(key), + ) + .toSorted(), + ); + await supervisor.close(); + }, + ); + }); + + it("bounds output and scrubs launch credentials after registry eviction", async () => { + const { supervisor, workspaceDir } = fixture(); + const successInput = launchInput(workspaceDir, "secret-success-launch", "secret-success"); + const failureInput = launchInput(workspaceDir, "failure-launch", "secret-fail"); + const overflowInput = launchInput(workspaceDir, "overflow-launch", "overflow"); + + await supervisor.launch(successInput); + await supervisor.launch(failureInput); + await supervisor.launch(overflowInput); + for (let index = 0; index < 600; index += 1) { + registerSecretValueForRedaction(`eviction-secret-${index}`); + } + const success = await waitForTerminal(supervisor, successInput.launchId); + const failure = await waitForTerminal(supervisor, failureInput.launchId); + const overflow = await waitForTerminal(supervisor, overflowInput.launchId); + const representations = [ + TEST_WORKER_CREDENTIAL, + encodeURIComponent(TEST_WORKER_CREDENTIAL), + JSON.stringify(TEST_WORKER_CREDENTIAL).slice(1, -1), + ]; + expect(success.state).toBe("completed"); + expect(JSON.parse(success.resultJson ?? "null")).toEqual({ + raw: "[REDACTED]", + encoded: "[REDACTED]", + status: "completed", + }); + expect(failure.state).toBe("failed"); + expect(Buffer.byteLength(failure.errorText ?? "", "utf8")).toBeLessThanOrEqual(4 * 1024); + for (const representation of representations) { + expect(success.resultJson).not.toContain(representation); + expect(failure.errorText).not.toContain(representation); + } + expect(overflow).toMatchObject({ + state: "failed", + errorText: expect.stringContaining("stdout exceeded 65536 bytes"), + }); + await supervisor.close(); + }); + + it.each([ + ["raw", "secret-cutoff-raw", TEST_WORKER_CREDENTIAL], + ["URL", "secret-cutoff-url", encodeURIComponent(TEST_WORKER_CREDENTIAL)], + ["JSON-escaped", "secret-cutoff-json", JSON.stringify(TEST_WORKER_CREDENTIAL).slice(1, -1)], + ])( + "redacts a %s credential representation across the stderr cutoff", + async (_, prompt, representation) => { + const { supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, `cutoff-${prompt}`, prompt); + + await supervisor.launch(input); + const failure = await waitForTerminal(supervisor, input.launchId); + + expect(failure.state).toBe("failed"); + expect(Buffer.byteLength(failure.errorText ?? "", "utf8")).toBeLessThanOrEqual(4 * 1024); + expect(failure.errorText).not.toContain(representation); + expect(failure.errorText).not.toContain(representation.slice(-8)); + await supervisor.close(); + }, + ); + + it("does not open or signal a child after markRunning observes its terminal receipt", async () => { + const { supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, "fast-terminal-launch", "fast-terminal"); + vi.spyOn(NodeWorkerLaunchStore.prototype, "markRunning").mockImplementation( + function (this: NodeWorkerLaunchStore, params) { + return this.finish({ + launchId: params.launchId, + planHash: params.planHash, + supervisor: params.supervisor, + worker: null, + state: "completed", + resultJson: '{"status":"completed"}', + }); + }, + ); + + expect(await supervisor.launch(input)).toMatchObject({ state: "completed" }); + const marker = path.join(workspaceDir, "fast-terminal-marker"); + await new Promise((resolve) => { + setTimeout(resolve, 150); + }); + expect(fs.existsSync(marker)).toBe(false); + await supervisor.close(); + }); + + it("records a gated child that exits before journal readiness as terminal", async () => { + const { supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, "prestart-exit-launch", "exit-before-start"); + const exitedPath = path.join(workspaceDir, "prestart-exited"); + + await supervisor.launch(input); + const terminal = await waitForTerminal(supervisor, input.launchId); + + expect(fs.existsSync(exitedPath)).toBe(true); + expect(terminal.state).toBe("failed"); + await supervisor.close(); + }); + + it.each([ + ["cancel", "cancelled"], + ["close", "interrupted"], + ] as const)("records %s while awaiting the owned child", async (operation, state) => { + const { supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, `${operation}-launch`, "wait"); + expect(await supervisor.launch(input)).toMatchObject({ state: "running" }); + + if (operation === "cancel") { + await supervisor.cancel(input.launchId); + } else { + await supervisor.close(); + } + + expect(await supervisor.status(input.launchId)).toMatchObject({ + state, + worker: { pid: expect.any(Number), startTime: expect.any(Number) }, + }); + await supervisor.close(); + }); + + it.each([ + ["cancel", "cancelled"], + ["close", "interrupted"], + ] as const)( + "%s during startup closes the gate before worker code runs", + async (operation, state) => { + const { supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, `${operation}-startup-launch`, "tree"); + const originalMarkRunning = Object.getOwnPropertyDescriptor( + NodeWorkerLaunchStore.prototype, + "markRunning", + )?.value as NodeWorkerLaunchStore["markRunning"]; + let stopping: Promise | undefined; + vi.spyOn(NodeWorkerLaunchStore.prototype, "markRunning").mockImplementation( + function (this: NodeWorkerLaunchStore, params) { + const receipt = Reflect.apply(originalMarkRunning, this, [params]); + stopping = + operation === "cancel" ? supervisor.cancel(input.launchId) : supervisor.close(); + return receipt; + }, + ); + + await supervisor.launch(input); + await stopping; + + expect((await supervisor.status(input.launchId))?.state).toBe(state); + expect(fs.existsSync(path.join(workspaceDir, "grandchild.pid"))).toBe(false); + await supervisor.close(); + }, + ); + + it.each([ + ["cancel", "cancelled"], + ["close", "interrupted"], + ] as const)("%s terminates the worker-owned grandchild", async (operation, state) => { + const { supervisor, workspaceDir } = fixture(); + const input = launchInput(workspaceDir, `${operation}-tree-launch`, "tree"); + const running = await supervisor.launch(input); + expect(running.state).toBe("running"); + const grandchildPath = path.join(workspaceDir, "grandchild.pid"); + await vi.waitFor(() => expect(fs.existsSync(grandchildPath)).toBe(true)); + const grandchildPid = Number(fs.readFileSync(grandchildPath, "utf8")); + const grandchild = requireNodeWorkerProcessIdentity(grandchildPid); + expect(inspectNodeWorkerProcessIdentity(grandchild)).toBe("live"); + + if (operation === "cancel") { + await supervisor.cancel(input.launchId); + } else { + await supervisor.close(); + } + + const terminal = await supervisor.status(input.launchId); + expect(terminal).toMatchObject({ state, worker: running.worker }); + await vi.waitFor(() => { + expect(inspectNodeWorkerProcessIdentity(running.worker!)).not.toBe("live"); + expect(inspectNodeWorkerProcessIdentity(grandchild)).not.toBe("live"); + }); + await supervisor.close(); + }); + + it("fails closed when the bundle entry resolves outside its namespaced bundle", async () => { + const { bundleRoot, root, supervisor, workspaceDir } = fixture(); + const escapedHash = "b".repeat(64); + const escapedBundle = path.join(bundleRoot, "gateway-1", "bundles", escapedHash); + const outsideEntry = path.join(root, "outside.mjs"); + fs.mkdirSync(escapedBundle, { recursive: true }); + fs.writeFileSync(outsideEntry, TEST_WORKER_SOURCE); + fs.symlinkSync(outsideEntry, path.join(escapedBundle, "openclaw.mjs")); + const input = launchInput(workspaceDir, "escaped-entry"); + input.bundleHash = escapedHash; + input.descriptor.admission.handshake.bundleHash = escapedHash; + + expect(await supervisor.launch(input)).toMatchObject({ + state: "failed", + errorText: expect.stringContaining("inside its bundle"), + }); + await supervisor.close(); + }); +}); diff --git a/src/node-host/node-worker-supervisor.ts b/src/node-host/node-worker-supervisor.ts new file mode 100644 index 000000000000..61fca11e0d73 --- /dev/null +++ b/src/node-host/node-worker-supervisor.ts @@ -0,0 +1,692 @@ +import { createHash } from "node:crypto"; +import fs from "node:fs"; +import path from "node:path"; +import { setTimeout as delay } from "node:timers/promises"; +import { stableStringify } from "@openclaw/normalization-core"; +import { resolveStateDir } from "../config/paths.js"; +import { formatErrorMessage } from "../infra/errors.js"; +import { isPathInside } from "../infra/path-guards.js"; +import { redactToolPayloadText } from "../logging/redact.js"; +import { + redactRegisteredSecretValues, + registerSecretValueForRedaction, +} from "../logging/secret-redaction-registry.js"; +import { + appendCapturedOutput, + createCapturedOutputBuffers, + finalizeCapturedOutput, +} from "../process/exec-output.js"; +import { signalProcessTree } from "../process/kill-tree.js"; +import { createChildAdapter } from "../process/supervisor/adapters/child.js"; +import { truncateUtf8Suffix } from "../utils/utf8-truncate.js"; +import { + parseWorkerLaunchDescriptor, + type WorkerLaunchDescriptor, +} from "../worker/launch-descriptor.js"; +import { snapshotNodeWorkerEnv } from "./node-worker-environment.js"; +import { + NodeWorkerLaunchStore, + type NodeWorkerLaunchReceipt, + type NodeWorkerTerminalState, +} from "./node-worker-launch-store.js"; +import { + inspectNodeWorkerProcessIdentity, + requireNodeWorkerProcessIdentity, + type NodeWorkerProcessIdentity, +} from "./node-worker-process-identity.js"; + +const STDOUT_MAX_BYTES = 64 * 1024; +const STDERR_MAX_BYTES = 4 * 1024; +const STOP_GRACE_MS = 1_000; +const FORCE_STOP_WAIT_MS = 4_000; +const RECOVERY_POLL_MS = 25; +const GATEWAY_NAMESPACE_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$/u; +const BUNDLE_HASH_PATTERN = /^[a-f0-9]{64}$/u; + +type NodeWorkerLaunchInput = { + launchId: string; + gatewayNamespace: string; + bundleHash: string; + placementGeneration: number; + descriptor: WorkerLaunchDescriptor; +}; + +type ChildAdapter = Awaited>; +type StopState = Extract; +type OwnedTreeState = "live" | "dead" | "unknown"; +type CredentialScrubber = { + maxRepresentationBytes: number; + scrub: (text: string) => string; +}; +type ActiveBase = { + launchId: string; + planHash: string; + supervisor: NodeWorkerProcessIdentity; + worker: NodeWorkerProcessIdentity; +}; +type RunningChild = ActiveBase & { + state: "running"; + adapter: ChildAdapter; + done: Promise; + journalReady: Promise; + releaseJournal: () => void; + scrubber: CredentialScrubber; + stopState?: StopState; +}; +type TerminalOutcome = Readonly<{ + state: NodeWorkerTerminalState; + resultJson?: string; + errorText?: string; +}>; +type ObservedTerminal = ActiveBase & { + state: "observed"; + outcome: TerminalOutcome; + persistenceError?: unknown; +}; +type ActiveOwnership = RunningChild | ObservedTerminal; + +function nodeWorkerPlanHash(params: { + bundleHash: string; + descriptor: WorkerLaunchDescriptor; + gatewayNamespace: string; + placementGeneration: number; +}): string { + return createHash("sha256").update(stableStringify(params)).digest("hex"); +} + +function resolveWorkerEntry(params: { + bundleRoot: string; + bundleHash: string; + gatewayNamespace: string; +}): string { + const root = fs.realpathSync.native(params.bundleRoot); + const bundle = fs.realpathSync.native( + path.join(root, params.gatewayNamespace, "bundles", params.bundleHash), + ); + if (!isPathInside(root, bundle)) { + throw new Error("node worker bundle resolves outside its configured root"); + } + const entry = fs.realpathSync.native(path.join(bundle, "openclaw.mjs")); + if (!isPathInside(bundle, entry) || !fs.statSync(entry).isFile()) { + throw new Error("node worker entry must be a regular file inside its bundle"); + } + return entry; +} + +function createCredentialScrubber(credential: string): CredentialScrubber { + const representations = new Set([ + credential, + encodeURIComponent(credential), + JSON.stringify(credential).slice(1, -1), + ]); + const ordered = [...representations].toSorted((left, right) => right.length - left.length); + return { + maxRepresentationBytes: Math.max( + ...ordered.map((representation) => Buffer.byteLength(representation, "utf8")), + ), + scrub: (text) => { + let scrubbed = text; + for (const representation of ordered) { + scrubbed = scrubbed.replaceAll(representation, "[REDACTED]"); + } + return scrubbed; + }, + }; +} + +function redactLaunchText(value: string, scrubCredential: (text: string) => string): string { + const launchRedacted = scrubCredential(value); + const exactRedacted = redactRegisteredSecretValues(launchRedacted, () => "[REDACTED]"); + return redactToolPayloadText(exactRedacted); +} + +function sanitizeDiagnostic( + value: string, + fallback: string, + scrubCredential: (text: string) => string, +): string { + const oneLine = redactLaunchText(value, scrubCredential).replace(/\s+/gu, " ").trim(); + return truncateUtf8Suffix(oneLine || fallback, STDERR_MAX_BYTES); +} + +function successfulResult( + stdout: ReturnType, + scrubCredential: (text: string) => string, +): string { + if (stdout.truncatedBytes > 0) { + throw new Error(`worker stdout exceeded ${STDOUT_MAX_BYTES} bytes`); + } + const raw = finalizeCapturedOutput(stdout, "head", true).toString("utf8").trim(); + const redacted = redactLaunchText(raw, scrubCredential); + let parsed: unknown; + try { + parsed = JSON.parse(redacted) as unknown; + } catch (error) { + throw new Error("worker returned invalid JSON output", { cause: error }); + } + const result = JSON.stringify(parsed); + if (Buffer.byteLength(result, "utf8") > STDOUT_MAX_BYTES) { + throw new Error(`worker result exceeded ${STDOUT_MAX_BYTES} bytes`); + } + return result; +} + +function inspectPosixProcessGroup(pid: number): OwnedTreeState { + try { + process.kill(-pid, 0); + return "live"; + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + return code === "ESRCH" ? "dead" : "unknown"; + } +} + +function inspectOwnedWorkerTree(worker: NodeWorkerProcessIdentity): OwnedTreeState { + const root = inspectNodeWorkerProcessIdentity(worker); + if (root === "reused") { + return "dead"; + } + if (root === "live") { + return "live"; + } + if (root === "unknown") { + return "unknown"; + } + return process.platform === "win32" ? "dead" : inspectPosixProcessGroup(worker.pid); +} + +async function signalOwnedWorkerTree( + worker: NodeWorkerProcessIdentity, + signal: "SIGTERM" | "SIGKILL", +): Promise { + const root = inspectNodeWorkerProcessIdentity(worker); + if (root === "reused" || root === "unknown") { + return; + } + await new Promise((resolve) => { + signalProcessTree(worker.pid, signal, { detached: true, onComplete: resolve }); + }); +} + +async function waitForOwnedWorkerTreeDeath( + worker: NodeWorkerProcessIdentity, + timeoutMs: number, +): Promise { + const deadline = Date.now() + timeoutMs; + let state = inspectOwnedWorkerTree(worker); + while (state === "live" && Date.now() < deadline) { + await delay(RECOVERY_POLL_MS); + state = inspectOwnedWorkerTree(worker); + } + return state; +} + +/** Owns worker process groups, lifetime gates, and the durable node-host launch journal. */ +class NodeWorkerSupervisor { + private readonly active = new Map(); + private readonly starting = new Map>(); + private readonly bundleRoot: string; + private readonly store: NodeWorkerLaunchStore; + private readonly workerEnv: NodeJS.ProcessEnv; + private supervisorIdentity?: NodeWorkerProcessIdentity; + private closed = false; + private closePromise?: Promise; + + constructor(options: { bundleRoot?: string; env?: NodeJS.ProcessEnv } = {}) { + const env = options.env ?? process.env; + this.bundleRoot = path.resolve( + options.bundleRoot ?? path.join(resolveStateDir(env), "node-host"), + ); + this.store = new NodeWorkerLaunchStore({ env }); + this.workerEnv = snapshotNodeWorkerEnv(env); + } + + private requireSupervisorIdentity(): NodeWorkerProcessIdentity { + return (this.supervisorIdentity ??= requireNodeWorkerProcessIdentity(process.pid)); + } + + async launch(input: NodeWorkerLaunchInput): Promise { + if (!GATEWAY_NAMESPACE_PATTERN.test(input.gatewayNamespace)) { + throw new Error("gateway namespace must be a safe bounded path component"); + } + if (!BUNDLE_HASH_PATTERN.test(input.bundleHash)) { + throw new Error("node worker bundle hash must be 64 lowercase hexadecimal characters"); + } + if (!Number.isSafeInteger(input.placementGeneration) || input.placementGeneration < 0) { + throw new Error("node worker placement generation must be a non-negative safe integer"); + } + const descriptor = parseWorkerLaunchDescriptor(structuredClone(input.descriptor)); + if (descriptor.admission.handshake.bundleHash !== input.bundleHash) { + throw new Error("node worker descriptor bundle hash does not match the launch bundle"); + } + const planHash = nodeWorkerPlanHash({ + bundleHash: input.bundleHash, + descriptor, + gatewayNamespace: input.gatewayNamespace, + placementGeneration: input.placementGeneration, + }); + const local = this.active.get(input.launchId); + if (local) { + if (local.planHash !== planHash) { + throw new Error(`node worker launch ${input.launchId} was replayed with a different plan`); + } + if (local.state === "observed") { + return this.reconcileActiveTerminal(local); + } + const receipt = this.store.get(input.launchId); + if (receipt) { + return receipt; + } + } + if (this.closed) { + throw new Error("node worker supervisor is closed"); + } + const supervisor = this.requireSupervisorIdentity(); + const claim = this.store.claim( + { + launchId: input.launchId, + planHash, + gatewayNamespace: input.gatewayNamespace, + environmentId: descriptor.admission.environmentId, + sessionId: descriptor.admission.sessionId, + ownerEpoch: descriptor.admission.ownerEpoch, + placementGeneration: input.placementGeneration, + runId: descriptor.assignment.runId, + }, + supervisor, + ); + if (claim.action === "recover") { + return await this.recoverRunning(claim.receipt); + } + if (claim.action === "replay") { + const replay = this.active.get(input.launchId); + if (replay?.planHash === planHash && replay.state === "observed") { + return this.reconcileActiveTerminal(replay); + } + const startup = this.starting.get(input.launchId); + return startup && claim.receipt.state === "pending" ? await startup : claim.receipt; + } + const startup = this.startClaimed({ input, descriptor, planHash, supervisor }); + this.starting.set(input.launchId, startup); + try { + return await startup; + } finally { + if (this.starting.get(input.launchId) === startup) { + this.starting.delete(input.launchId); + } + } + } + + async status(launchId: string): Promise { + const active = this.active.get(launchId); + if (active?.state === "observed") { + return this.reconcileActiveTerminal(active); + } + return this.store.get(launchId); + } + + async cancel(launchId: string): Promise { + const active = this.active.get(launchId); + if (active) { + if (active.state === "running") { + await this.stopChild(active, "cancelled"); + } + const observed = this.active.get(launchId); + if (observed?.state === "observed") { + return this.reconcileActiveTerminal(observed); + } + return this.store.get(launchId); + } + const startup = this.starting.get(launchId); + const receipt = this.store.get(launchId); + if (!receipt || receipt.state === "completed" || receipt.state === "failed") { + return receipt; + } + if (receipt.state === "interrupted" || receipt.state === "cancelled") { + return receipt; + } + if (!startup || receipt.state !== "pending" || receipt.supervisor.pid !== process.pid) { + return receipt; + } + const cancelled = this.store.finish({ + launchId, + planHash: receipt.planHash, + supervisor: this.requireSupervisorIdentity(), + worker: null, + state: "cancelled", + errorText: "node worker launch cancelled", + }); + await startup; + return this.store.get(launchId) ?? cancelled; + } + + close(): Promise { + if (this.closePromise) { + return this.closePromise; + } + this.closed = true; + const operation = (async () => { + await Promise.allSettled(this.starting.values()); + await Promise.all( + [...this.active.values()] + .filter((active): active is RunningChild => active.state === "running") + .map(async (active) => await this.stopChild(active, "interrupted")), + ); + const errors: unknown[] = []; + for (const active of this.active.values()) { + if (active.state !== "observed") { + continue; + } + try { + this.reconcileActiveTerminal(active); + } catch (error) { + errors.push(error); + } + } + if (errors.length === 1) { + throw errors[0]; + } + if (errors.length > 1) { + throw new AggregateError(errors, "node worker terminal reconciliation failed"); + } + })(); + const closePromise = operation.finally(() => { + if (this.closePromise === closePromise) { + this.closePromise = undefined; + } + }); + this.closePromise = closePromise; + return closePromise; + } + + private reconcileActiveTerminal(active: ObservedTerminal): NodeWorkerLaunchReceipt { + try { + const receipt = this.store.finish({ + launchId: active.launchId, + planHash: active.planHash, + supervisor: active.supervisor, + worker: active.worker, + ...active.outcome, + }); + if (receipt.state === "pending" || receipt.state === "running") { + throw new Error(`node worker launch ${active.launchId} terminal state was not persisted`); + } + if (this.active.get(active.launchId) === active) { + this.active.delete(active.launchId); + } + return receipt; + } catch (error) { + active.persistenceError = error; + throw error; + } + } + + private async recoverRunning(receipt: NodeWorkerLaunchReceipt): Promise { + if (receipt.state !== "running" || !receipt.worker) { + return receipt; + } + const previousSupervisor = inspectNodeWorkerProcessIdentity(receipt.supervisor); + if (previousSupervisor !== "dead" && previousSupervisor !== "reused") { + return this.store.get(receipt.launchId) ?? receipt; + } + let workerState = inspectOwnedWorkerTree(receipt.worker); + if (workerState === "unknown") { + return this.store.get(receipt.launchId) ?? receipt; + } + if (workerState === "live") { + await signalOwnedWorkerTree(receipt.worker, "SIGTERM"); + workerState = await waitForOwnedWorkerTreeDeath(receipt.worker, STOP_GRACE_MS); + } + if (workerState === "live") { + await signalOwnedWorkerTree(receipt.worker, "SIGKILL"); + workerState = await waitForOwnedWorkerTreeDeath(receipt.worker, FORCE_STOP_WAIT_MS); + } + if (workerState !== "dead") { + return this.store.get(receipt.launchId) ?? receipt; + } + return this.store.finish({ + launchId: receipt.launchId, + planHash: receipt.planHash, + supervisor: receipt.supervisor, + worker: receipt.worker, + state: "interrupted", + errorText: "node host stopped before the worker launch completed", + }); + } + + private async startClaimed(params: { + input: NodeWorkerLaunchInput; + descriptor: WorkerLaunchDescriptor; + planHash: string; + supervisor: NodeWorkerProcessIdentity; + }): Promise { + const credential = params.descriptor.admission.credential; + const scrubber = createCredentialScrubber(credential); + registerSecretValueForRedaction(credential); + let adapter: ChildAdapter; + try { + const entry = resolveWorkerEntry({ + bundleRoot: this.bundleRoot, + bundleHash: params.input.bundleHash, + gatewayNamespace: params.input.gatewayNamespace, + }); + adapter = await createChildAdapter({ + argv: [process.execPath, entry, "worker", "--internal-worker-ipc"], + env: this.workerEnv, + exactEnv: true, + ownedWorker: true, + input: JSON.stringify(params.descriptor), + }); + } catch (error) { + return this.store.finish({ + launchId: params.input.launchId, + planHash: params.planHash, + supervisor: params.supervisor, + worker: null, + state: "failed", + errorText: sanitizeDiagnostic( + formatErrorMessage(error), + "node worker spawn failed", + scrubber.scrub, + ), + }); + } + if (!adapter.pid) { + adapter.kill("SIGKILL"); + adapter.dispose(); + return this.store.finish({ + launchId: params.input.launchId, + planHash: params.planHash, + supervisor: params.supervisor, + worker: null, + state: "failed", + errorText: "node worker spawn did not return a process id", + }); + } + let worker: NodeWorkerProcessIdentity; + try { + worker = requireNodeWorkerProcessIdentity(adapter.pid); + } catch (error) { + adapter.kill("SIGKILL"); + await adapter.wait().catch(() => undefined); + adapter.dispose(); + return this.store.finish({ + launchId: params.input.launchId, + planHash: params.planHash, + supervisor: params.supervisor, + worker: null, + state: "failed", + errorText: sanitizeDiagnostic( + formatErrorMessage(error), + "node worker process identity unavailable", + scrubber.scrub, + ), + }); + } + let journalReleased = false; + let releaseJournalPromise!: () => void; + const journalReady = new Promise((resolve) => { + releaseJournalPromise = resolve; + }); + const releaseJournal = () => { + if (!journalReleased) { + journalReleased = true; + releaseJournalPromise(); + } + }; + const active = { + state: "running", + adapter, + journalReady, + launchId: params.input.launchId, + planHash: params.planHash, + releaseJournal, + scrubber, + supervisor: params.supervisor, + worker, + } as RunningChild; + active.done = this.observeChild(active); + this.active.set(active.launchId, active); + void active.done.catch(() => undefined); + let running: NodeWorkerLaunchReceipt; + try { + running = this.store.markRunning({ + launchId: active.launchId, + planHash: active.planHash, + supervisor: params.supervisor, + worker, + }); + } catch (error) { + active.releaseJournal(); + await this.stopChild(active, "interrupted").catch(() => undefined); + throw error; + } + active.releaseJournal(); + if (running.state === "cancelled" || running.state === "interrupted") { + await this.stopChild(active, running.state); + return this.store.get(active.launchId) ?? running; + } + if (running.state !== "running") { + adapter.closeStartGate?.(); + return running; + } + if (this.closed) { + await this.stopChild(active, "interrupted"); + return this.store.get(active.launchId) ?? running; + } + try { + await adapter.openStartGate?.(); + } catch { + await this.stopChild(active, "interrupted"); + return this.store.get(active.launchId) ?? running; + } + return running; + } + + private async observeChild(active: RunningChild): Promise { + const stdout = createCapturedOutputBuffers(); + const stderr = createCapturedOutputBuffers(); + active.adapter.onStdout((chunk) => + appendCapturedOutput(stdout, chunk, STDOUT_MAX_BYTES, "head"), + ); + active.adapter.onStderr((chunk) => + appendCapturedOutput( + stderr, + chunk, + STDERR_MAX_BYTES + active.scrubber.maxRepresentationBytes, + "tail", + ), + ); + let outcome: TerminalOutcome; + try { + const exit = await active.adapter.wait(); + await active.journalReady; + if (active.stopState) { + outcome = Object.freeze({ + state: active.stopState, + errorText: + active.stopState === "cancelled" + ? "node worker launch cancelled" + : "node worker launch interrupted during node-host shutdown", + }); + } else if (exit.code === 0 && exit.signal === null) { + try { + outcome = Object.freeze({ + state: "completed", + resultJson: successfulResult(stdout, active.scrubber.scrub), + }); + } catch (error) { + outcome = Object.freeze({ + state: "failed", + errorText: sanitizeDiagnostic( + formatErrorMessage(error), + "invalid worker result", + active.scrubber.scrub, + ), + }); + } + } else { + const detail = finalizeCapturedOutput(stderr, "tail", true).toString("utf8"); + const exitLabel = exit.signal ? `signal ${exit.signal}` : `exit code ${String(exit.code)}`; + outcome = Object.freeze({ + state: "failed", + errorText: sanitizeDiagnostic( + `node worker failed with ${exitLabel}${detail ? `: ${detail}` : ""}`, + "node worker failed", + active.scrubber.scrub, + ), + }); + } + } catch (error) { + await active.journalReady; + outcome = Object.freeze({ + state: active.stopState ?? "failed", + errorText: sanitizeDiagnostic( + formatErrorMessage(error), + "node worker wait failed", + active.scrubber.scrub, + ), + }); + } finally { + active.adapter.dispose(); + } + const observed: ObservedTerminal = { + state: "observed", + launchId: active.launchId, + planHash: active.planHash, + supervisor: active.supervisor, + worker: active.worker, + outcome, + }; + if (this.active.get(active.launchId) !== active) { + return; + } + this.active.set(active.launchId, observed); + try { + this.reconcileActiveTerminal(observed); + } catch { + // The observed outcome stays owned in memory for the next supervisor operation. + } + } + + private async stopChild(active: RunningChild, state: StopState): Promise { + active.stopState ??= state; + active.adapter.kill("SIGTERM"); + const forceKill = setTimeout(() => active.adapter.kill("SIGKILL"), STOP_GRACE_MS); + forceKill.unref?.(); + try { + await active.done; + } finally { + clearTimeout(forceKill); + } + } +} + +export function createNodeWorkerSupervisor( + options: { + bundleRoot?: string; + env?: NodeJS.ProcessEnv; + } = {}, +): NodeWorkerSupervisor { + return new NodeWorkerSupervisor(options); +} diff --git a/src/node-host/runner.test.ts b/src/node-host/runner.test.ts index 7ecb43b28b8e..65ece1715fcf 100644 --- a/src/node-host/runner.test.ts +++ b/src/node-host/runner.test.ts @@ -545,6 +545,8 @@ describe("runNodeHost", () => { await vi.waitFor(() => expect(mocks.capturedGatewayClients[0]?.stop).toHaveBeenCalledOnce()); expect(clearIntervalSpy).not.toHaveBeenCalled(); + await vi.waitFor(() => expect(mocks.closeMcpManager).toHaveBeenCalledOnce()); + expect(resolveCloseMcp).toBeTypeOf("function"); resolveCloseMcp?.(); await running; diff --git a/src/node-host/runtime.test.ts b/src/node-host/runtime.test.ts index bf3b1e82dce1..5310a6d29778 100644 --- a/src/node-host/runtime.test.ts +++ b/src/node-host/runtime.test.ts @@ -8,6 +8,7 @@ import { prepareNodeHostRuntime } from "./runtime.js"; const mocks = vi.hoisted(() => ({ closeMcp: vi.fn(async () => undefined), + closeWorkerSupervisor: vi.fn(async () => undefined), handleInvoke: vi.fn(async () => undefined), progressStartHeartbeats: vi.fn(), progressWrite: vi.fn(async () => undefined), @@ -39,6 +40,10 @@ vi.mock("./node-invoke-progress.js", () => ({ })), })); +vi.mock("./node-worker-supervisor.js", () => ({ + createNodeWorkerSupervisor: vi.fn(() => ({ close: mocks.closeWorkerSupervisor })), +})); + vi.mock("./plugin-node-host.js", () => ({ ensureNodeHostPluginRegistry: vi.fn(async () => undefined), isRegisteredNodeHostCommandDuplex: vi.fn((command: string) => command === "test.duplex"), @@ -178,6 +183,7 @@ describe("node-host invocation cancellation", () => { await runtime.close(); expect(held.signal?.aborted).toBe(true); + expect(mocks.closeWorkerSupervisor).toHaveBeenCalledOnce(); held.release(); await invoking; }); diff --git a/src/node-host/runtime.ts b/src/node-host/runtime.ts index e36e86967a28..403f59b50836 100644 --- a/src/node-host/runtime.ts +++ b/src/node-host/runtime.ts @@ -26,6 +26,7 @@ import { handleInvoke, type NodeInvokeRequestPayload, type SkillBinsProvider } f import { startNodeHostMcpManager, type NodeHostMcpManager } from "./mcp.js"; import { buildNodeEventParams } from "./node-event-params.js"; import { createNodeInvokeProgressWriter } from "./node-invoke-progress.js"; +import { createNodeWorkerSupervisor } from "./node-worker-supervisor.js"; import { ensureNodeHostPluginRegistry, isRegisteredNodeHostCommandDuplex, @@ -305,6 +306,7 @@ export async function prepareNodeHostRuntime(params?: { initialInventory, start({ client, onInventoryChanged, onManifestChanged }) { const mcpAbort = new AbortController(); + const workerSupervisor = createNodeWorkerSupervisor({ env }); const skillBins = new SkillBinsCache(client, pathEnv); const activeInvokes = new Map(); const pluginCommandContext: OpenClawPluginNodeHostCommandContext = { @@ -448,6 +450,7 @@ export async function prepareNodeHostRuntime(params?: { async close() { this.cancelAll(); stopAvailabilityWatch(); + await workerSupervisor.close(); mcpAbort.abort(); const resolved = manager ?? (await startup.catch(() => undefined)); await resolved?.close(); diff --git a/src/process/spawn-secret-input.ts b/src/process/spawn-secret-input.ts index a1e6a95313f9..01b2608456de 100644 --- a/src/process/spawn-secret-input.ts +++ b/src/process/spawn-secret-input.ts @@ -2,7 +2,7 @@ import type { ChildProcess } from "node:child_process"; import type { Writable } from "node:stream"; import type { SpawnSecretInput } from "./supervisor/types.js"; -export type SpawnStdioEntry = "ignore" | "inherit" | "overlapped" | "pipe"; +export type SpawnStdioEntry = "ignore" | "inherit" | "ipc" | "overlapped" | "pipe"; export function addSecretInputStdio( stdio: SpawnStdioEntry[], diff --git a/src/process/supervisor/adapters/child.test.ts b/src/process/supervisor/adapters/child.test.ts index 2bd099634d4c..a041803faace 100644 --- a/src/process/supervisor/adapters/child.test.ts +++ b/src/process/supervisor/adapters/child.test.ts @@ -57,8 +57,23 @@ function createStubChild(pid = 1234) { Object.defineProperty(child, "killed", { value: false, configurable: true, writable: true }); Object.defineProperty(child, "exitCode", { value: null, configurable: true, writable: true }); Object.defineProperty(child, "signalCode", { value: null, configurable: true, writable: true }); + Object.defineProperty(child, "channel", { value: {}, configurable: true }); + Object.defineProperty(child, "connected", { value: true, configurable: true, writable: true }); const killMock = vi.fn(() => true); + const sendMock = vi.fn((_message: unknown, ...args: unknown[]) => { + const callback = args.findLast((value) => typeof value === "function") as + | ((error: Error | null) => void) + | undefined; + callback?.(null); + return true; + }); + const disconnectMock = vi.fn(() => { + Object.defineProperty(child, "connected", { value: false, configurable: true, writable: true }); + child.emit("disconnect"); + }); child.kill = killMock as ChildProcess["kill"]; + child.send = sendMock as ChildProcess["send"]; + child.disconnect = disconnectMock as ChildProcess["disconnect"]; const emitClose = (code: number | null, signal: NodeJS.Signals | null = null) => { child.emit("close", code, signal); }; @@ -71,7 +86,7 @@ function createStubChild(pid = 1234) { }); child.emit("exit", code, signal); }; - return { child, killMock, emitClose, emitExit }; + return { child, disconnectMock, killMock, sendMock, emitClose, emitExit }; } async function createAdapterHarness(params?: { @@ -224,6 +239,30 @@ describe("createChildAdapter", () => { expect(killMock).toHaveBeenCalledWith("SIGKILL"); }); + it("creates owned worker trees in a dedicated POSIX process group without fallback", async () => { + process.env.OPENCLAW_SERVICE_MARKER = "service-managed"; + const { child, disconnectMock, sendMock } = createStubChild(); + spawnWithFallbackMock.mockResolvedValue({ child, usedFallback: false }); + + const adapter = await createChildAdapter({ + argv: ["node", "worker"], + ownedWorker: true, + input: "{}", + }); + + expect(firstSpawnWithFallbackParams().options?.detached).toBe(process.platform !== "win32"); + expect(firstSpawnWithFallbackParams().fallbacks).toEqual([]); + expect(firstSpawnWithFallbackParams().options?.stdio).toEqual(["pipe", "pipe", "pipe", "ipc"]); + + await adapter.openStartGate?.(); + expect(sendMock).toHaveBeenCalledWith( + { type: "openclaw-worker-start-v1" }, + expect.any(Function), + ); + adapter.closeStartGate?.(); + expect(disconnectMock).toHaveBeenCalledOnce(); + }); + it("writes secret input to an extra descriptor and zeroes the transient buffer", async () => { const { child } = createStubChild(); const secretStream = new PassThrough(); @@ -788,6 +827,28 @@ describe("createChildAdapter", () => { expect(spawnArgs.options.env.CDPATH).toBeUndefined(); }); + it("keeps an exact Linux child environment out of the OOM shell wrapper", async () => { + setPlatform("linux"); + const restoreLinuxShell = mockLinuxOomWrapperShell(); + const { child } = createStubChild(3335); + spawnWithFallbackMock.mockResolvedValue({ child, usedFallback: false }); + try { + const adapter = await createChildAdapter({ + argv: ["/usr/bin/node", "-e", "process.exit(0)"], + env: { HOME: "/worker-home", PATH: "/usr/bin" }, + exactEnv: true, + stdinMode: "pipe-open", + }); + expect(adapter.oomScoreWrapperSelected).toBe(false); + } finally { + restoreLinuxShell(); + } + + const spawnArgs = firstSpawnWithFallbackParams(); + expect(spawnArgs.argv).toEqual(["/usr/bin/node", "-e", "process.exit(0)"]); + expect(spawnArgs.options?.env).toEqual({ HOME: "/worker-home", PATH: "/usr/bin" }); + }); + it("passes explicit env overrides as strings", async () => { await createAdapterHarness({ pid: 4444, diff --git a/src/process/supervisor/adapters/child.ts b/src/process/supervisor/adapters/child.ts index 7e08e4b18e98..839eefaf39c4 100644 --- a/src/process/supervisor/adapters/child.ts +++ b/src/process/supervisor/adapters/child.ts @@ -71,6 +71,12 @@ function resolveChildInvocation(params: { } type ChildAdapter = SpawnProcessAdapter; +type WorkerChildAdapter = ChildAdapter & { + closeStartGate?: () => void; + openStartGate?: () => Promise; +}; + +const WORKER_START_MESSAGE = { type: "openclaw-worker-start-v1" } as const; function isServiceManagedRuntime(): boolean { return Boolean(process.env.OPENCLAW_SERVICE_MARKER?.trim()); @@ -78,32 +84,40 @@ function isServiceManagedRuntime(): boolean { export async function createChildAdapter(params: { argv: string[]; + /** Own a separately signalable tree whose private IPC channel gates worker startup. */ + ownedWorker?: true; + /** Preserve the supplied environment exactly by skipping environment-mutating spawn wrappers. */ + exactEnv?: true; cwd?: string; env?: NodeJS.ProcessEnv; windowsVerbatimArguments?: boolean; input?: string; stdinMode?: "inherit" | "pipe-open" | "pipe-closed"; secretInput?: SpawnSecretInput; -}): Promise { +}): Promise { const baseEnv = params.env ? toStringEnv(params.env) : undefined; const invocation = resolveChildInvocation({ argv: params.argv, env: baseEnv, windowsVerbatimArguments: params.windowsVerbatimArguments, }); - const preparedSpawn = prepareOomScoreAdjustedSpawn(invocation.command, invocation.args, { - env: baseEnv, - }); + const preparedSpawn = params.exactEnv + ? { command: invocation.command, args: invocation.args, env: baseEnv, wrapped: false } + : prepareOomScoreAdjustedSpawn(invocation.command, invocation.args, { env: baseEnv }); const stdinMode = params.stdinMode ?? (params.input !== undefined ? "pipe-closed" : "inherit"); - // In service-managed mode keep children attached so systemd/launchd can - // stop the full process tree reliably. Outside service mode preserve the - // existing POSIX detached behavior. - const useDetached = process.platform !== "win32" && !isServiceManagedRuntime(); + // A detached POSIX child is still a descendant in the service cgroup/job, but + // owns a process group that can be killed without touching the node host. + const useDetached = + process.platform !== "win32" && + (params.ownedWorker !== undefined || !isServiceManagedRuntime()); const stdio: SpawnStdioEntry[] = [stdinMode === "inherit" ? "inherit" : "pipe", "pipe", "pipe"]; addSecretInputStdio(stdio, params.secretInput); + if (params.ownedWorker !== undefined) { + stdio.push("ipc"); + } const options: SpawnOptions = { cwd: params.cwd, @@ -117,17 +131,34 @@ export async function createChildAdapter(params: { const spawned = await spawnWithFallback({ argv: [preparedSpawn.command, ...preparedSpawn.args], options, - fallbacks: useDetached - ? [ - { - label: "no-detach", - options: { detached: false }, - }, - ] - : [], + fallbacks: + useDetached && params.ownedWorker === undefined + ? [ + { + label: "no-detach", + options: { detached: false }, + }, + ] + : [], }); const child = spawned.child as ChildProcessWithoutNullStreams; + if (params.ownedWorker !== undefined && (!child.connected || !child.channel)) { + spawned.child.kill("SIGKILL"); + throw new Error("worker lifecycle IPC channel was not created"); + } + const disconnectWorkerIpc = () => { + if (!child.connected) { + return; + } + try { + child.disconnect(); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ERR_IPC_DISCONNECTED") { + throw error; + } + } + }; // Pipe errors can arrive before output subscribers attach. Close remains // responsible for decoder flush and Windows drain completion. const ignoreOutputStreamError = () => {}; @@ -489,9 +520,41 @@ export async function createChildAdapter(params: { const dispose = () => { clearForceKillWaitFallback(); clearForcedWindowsCloseTimer(); + if (params.ownedWorker !== undefined) { + disconnectWorkerIpc(); + } child.removeAllListeners(); }; + const closeStartGate = params.ownedWorker ? disconnectWorkerIpc : undefined; + + let startGateOpened = false; + const openStartGate = params.ownedWorker + ? async () => { + if (startGateOpened) { + return; + } + startGateOpened = true; + await new Promise((resolve, reject) => { + if (!child.connected) { + reject(new Error("worker lifecycle IPC channel closed before startup")); + return; + } + try { + child.send(WORKER_START_MESSAGE, (error) => { + if (error) { + reject(error); + return; + } + resolve(); + }); + } catch (error) { + reject(toErrorObject(error, "worker lifecycle IPC send failed")); + } + }); + } + : undefined; + return { pid: child.pid ?? undefined, stdin, @@ -501,5 +564,7 @@ export async function createChildAdapter(params: { wait, kill, dispose, + closeStartGate, + openStartGate, }; } diff --git a/src/state/openclaw-database-maintenance.test.ts b/src/state/openclaw-database-maintenance.test.ts index ad2e154cbda8..869865edc2c4 100644 --- a/src/state/openclaw-database-maintenance.test.ts +++ b/src/state/openclaw-database-maintenance.test.ts @@ -311,39 +311,39 @@ describe("OpenClaw database maintenance schema validation", () => { } }); - it("allows the lazy worker SSH fallback table to be absent but rejects drift", () => { - const database = createGlobalDatabase(); - try { - const canonicalTable = database - .prepare("SELECT sql FROM sqlite_schema WHERE type = 'table' AND name = ?") - .get("worker_environment_ssh_fallback_ports") as { sql?: unknown } | undefined; - if (typeof canonicalTable?.sql !== "string") { - throw new Error("missing canonical worker SSH fallback port table"); + it.each(["node_worker_launches", "worker_environment_ssh_fallback_ports"])( + "allows lazy table %s to be absent but rejects drift", + (tableName) => { + const database = createGlobalDatabase(); + try { + const canonicalTable = database + .prepare("SELECT sql FROM sqlite_schema WHERE type = 'table' AND name = ?") + .get(tableName) as { sql?: unknown } | undefined; + if (typeof canonicalTable?.sql !== "string") { + throw new Error(`missing canonical ${tableName} table`); + } + database.exec(`DROP TABLE ${tableName};`); + + expect(() => + assertOpenClawStateDatabaseForMaintenance(database, { + pathname: "global.sqlite", + }), + ).not.toThrow(); + + const driftedTableSql = canonicalTable.sql.replace("(\n", "(\n unexpected TEXT,\n"); + expect(driftedTableSql).not.toBe(canonicalTable.sql); + database.exec(driftedTableSql); + + expect(() => + assertOpenClawStateDatabaseForMaintenance(database, { + pathname: "global.sqlite", + }), + ).toThrow(`column definitions differ for ${tableName}`); + } finally { + database.close(); } - database.exec("DROP TABLE worker_environment_ssh_fallback_ports;"); - - expect(() => - assertOpenClawStateDatabaseForMaintenance(database, { - pathname: "global.sqlite", - }), - ).not.toThrow(); - - const driftedTableSql = canonicalTable.sql.replace( - " PRIMARY KEY (environment_id, position)", - " unexpected TEXT,\n PRIMARY KEY (environment_id, position)", - ); - expect(driftedTableSql).not.toBe(canonicalTable.sql); - database.exec(driftedTableSql); - - expect(() => - assertOpenClawStateDatabaseForMaintenance(database, { - pathname: "global.sqlite", - }), - ).toThrow("column definitions differ for worker_environment_ssh_fallback_ports"); - } finally { - database.close(); - } - }); + }, + ); it("rejects a current agent database with a missing canonical table", () => { const database = createAgentDatabase(); diff --git a/src/state/openclaw-state-db-contract.ts b/src/state/openclaw-state-db-contract.ts index 53e984a6dbd6..9c71a98ed588 100644 --- a/src/state/openclaw-state-db-contract.ts +++ b/src/state/openclaw-state-db-contract.ts @@ -12,6 +12,7 @@ export const FIRST_USE_STATE_TABLES = [ "cron_job_runtime_authorities", "execution_identity_contexts", "mcp_oauth_pending_authorizations", + "node_worker_launches", "operator_approval_execution_identities", "execution_decision_facts", ] as const; diff --git a/src/state/openclaw-state-db.generated.d.ts b/src/state/openclaw-state-db.generated.d.ts index cc8ab3e44855..ac017ab08e82 100644 --- a/src/state/openclaw-state-db.generated.d.ts +++ b/src/state/openclaw-state-db.generated.d.ts @@ -941,6 +941,27 @@ export interface NodeHostConfig { version: number; } +export interface NodeWorkerLaunches { + completed_at_ms: number | null; + created_at_ms: number; + environment_id: string; + error_text: string | null; + gateway_namespace: string; + launch_id: string; + owner_epoch: number; + placement_generation: number; + plan_hash: string; + result_json: string | null; + run_id: string; + session_id: string; + state: string; + supervisor_pid: number; + supervisor_start_time: number; + updated_at_ms: number; + worker_pid: number | null; + worker_start_time: number | null; +} + export interface OfficialExternalPluginCatalogSnapshots { body: string; checksum: string; @@ -1717,6 +1738,7 @@ export interface DB { model_catalog_remote: ModelCatalogRemote; native_hook_relay_bridges: NativeHookRelayBridges; node_host_config: NodeHostConfig; + node_worker_launches: NodeWorkerLaunches; official_external_plugin_catalog_snapshots: OfficialExternalPluginCatalogSnapshots; onboarding_recommendations: OnboardingRecommendations; operator_approval_execution_identities: OperatorApprovalExecutionIdentities; diff --git a/src/state/openclaw-state-schema.sql b/src/state/openclaw-state-schema.sql index bb21ce4682eb..0a83a2a9ce3a 100644 --- a/src/state/openclaw-state-schema.sql +++ b/src/state/openclaw-state-schema.sql @@ -854,6 +854,81 @@ CREATE TABLE IF NOT EXISTS node_host_config ( updated_at_ms INTEGER NOT NULL ) STRICT; +-- Node-host-owned launch journal. The descriptor and its credential remain +-- process memory only; this table records bounded supervision facts. +CREATE TABLE IF NOT EXISTS node_worker_launches ( + launch_id TEXT NOT NULL PRIMARY KEY + CHECK (length(launch_id) BETWEEN 1 AND 256 AND instr(launch_id, char(0)) = 0), + plan_hash TEXT NOT NULL + CHECK (length(plan_hash) = 64 AND plan_hash NOT GLOB '*[^0-9a-f]*'), + gateway_namespace TEXT NOT NULL + CHECK ( + length(gateway_namespace) BETWEEN 1 AND 128 + AND gateway_namespace NOT GLOB '*[^A-Za-z0-9._-]*' + AND gateway_namespace GLOB '[A-Za-z0-9]*' + ), + environment_id TEXT NOT NULL + CHECK (length(environment_id) BETWEEN 1 AND 256 AND instr(environment_id, char(0)) = 0), + session_id TEXT NOT NULL + CHECK (length(session_id) BETWEEN 1 AND 256 AND instr(session_id, char(0)) = 0), + owner_epoch INTEGER NOT NULL CHECK (owner_epoch BETWEEN 1 AND 9007199254740991), + placement_generation INTEGER NOT NULL + CHECK (placement_generation BETWEEN 0 AND 9007199254740991), + run_id TEXT NOT NULL + CHECK (length(run_id) BETWEEN 1 AND 256 AND instr(run_id, char(0)) = 0), + state TEXT NOT NULL + CHECK (state IN ('pending', 'running', 'completed', 'failed', 'interrupted', 'cancelled')), + supervisor_pid INTEGER NOT NULL CHECK (supervisor_pid BETWEEN 1 AND 2147483647), + supervisor_start_time INTEGER NOT NULL + CHECK (supervisor_start_time BETWEEN 0 AND 9007199254740991), + worker_pid INTEGER CHECK (worker_pid IS NULL OR worker_pid BETWEEN 1 AND 2147483647), + worker_start_time INTEGER CHECK ( + worker_start_time IS NULL OR worker_start_time BETWEEN 0 AND 9007199254740991 + ), + result_json TEXT CHECK ( + result_json IS NULL + OR ( + length(CAST(result_json AS BLOB)) BETWEEN 1 AND 65536 + AND instr(result_json, char(0)) = 0 + AND json_valid(result_json) + ) + ), + error_text TEXT CHECK ( + error_text IS NULL + OR ( + length(CAST(error_text AS BLOB)) BETWEEN 1 AND 4096 + AND instr(error_text, char(0)) = 0 + AND instr(error_text, char(10)) = 0 + AND instr(error_text, char(13)) = 0 + ) + ), + completed_at_ms INTEGER CHECK ( + completed_at_ms IS NULL OR completed_at_ms BETWEEN 0 AND 9007199254740991 + ), + created_at_ms INTEGER NOT NULL CHECK (created_at_ms BETWEEN 0 AND 9007199254740991), + updated_at_ms INTEGER NOT NULL CHECK ( + updated_at_ms BETWEEN created_at_ms AND 9007199254740991 + ), + CHECK ((worker_pid IS NULL) = (worker_start_time IS NULL)), + CHECK ( + (state = 'pending' + AND worker_pid IS NULL AND result_json IS NULL AND error_text IS NULL + AND completed_at_ms IS NULL) + OR + (state = 'running' + AND worker_pid IS NOT NULL AND result_json IS NULL AND error_text IS NULL + AND completed_at_ms IS NULL) + OR + (state = 'completed' + AND result_json IS NOT NULL AND error_text IS NULL + AND completed_at_ms BETWEEN created_at_ms AND updated_at_ms) + OR + (state IN ('failed', 'interrupted', 'cancelled') + AND result_json IS NULL AND error_text IS NOT NULL + AND completed_at_ms BETWEEN created_at_ms AND updated_at_ms) + ) +) STRICT; + CREATE TABLE IF NOT EXISTS voicewake_triggers ( config_key TEXT NOT NULL, position INTEGER NOT NULL, diff --git a/src/worker/worker-command.runtime.test.ts b/src/worker/worker-command.runtime.test.ts new file mode 100644 index 000000000000..b3d7d86a077b --- /dev/null +++ b/src/worker/worker-command.runtime.test.ts @@ -0,0 +1,177 @@ +import { PassThrough } from "node:stream"; +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { + WORKER_PROTOCOL_FEATURES, + WORKER_RPC_SET_VERSION, +} from "../../packages/gateway-protocol/src/schema/worker-admission.js"; +import type { WorkerLaunchDescriptor } from "./launch-descriptor.js"; +import { runWorkerCommand } from "./worker-command.runtime.js"; +import { runWorkerDescriptor } from "./worker.runtime.js"; + +vi.mock("./worker.runtime.js", () => ({ + runWorkerDescriptor: vi.fn(), +})); + +const descriptor = { + version: 3, + connectionEndpoint: { kind: "unix", socketPath: "/tmp/openclaw-worker/gateway.sock" }, + admission: { + environmentId: "environment-1", + credential: ["worker", "fixture", "value"].join("-"), + sessionId: "session-1", + ownerEpoch: 1, + rpcSetVersion: WORKER_RPC_SET_VERSION, + handshake: { + bundleHash: "a".repeat(64), + openclawVersion: "2026.7.12", + protocolFeatures: [...WORKER_PROTOCOL_FEATURES], + }, + }, + assignment: { + agentId: "agent-1", + operationalRunInstance: { instanceId: "instance-run-1", runId: "run-1" }, + agentRuntimeIdentityToken: "signed-runtime-token", + runId: "run-1", + turnId: "turn-1", + prompt: "Inspect the workspace.", + suppressPromptTranscript: false, + workspaceDir: "/tmp/openclaw-worker/workspace", + modelRef: { provider: "provider-1", model: "model-1" }, + inferenceOptions: { reasoning: "medium", maxTokens: 512 }, + initialMessages: [ + { + role: "user", + content: [{ type: "text", text: "Earlier context." }], + timestamp: 1, + }, + ], + transcript: { baseLeafId: "leaf-7", nextSeq: 8 }, + liveEvents: { ackedSeq: 12, nextSeq: 13 }, + toolAuthority: { allowedToolNames: ["read", "exec"] }, + }, +} satisfies WorkerLaunchDescriptor; + +function commandInput() { + const input = new PassThrough(); + input.end(JSON.stringify(descriptor)); + return input; +} + +function lifetimeHarness() { + const controller = new AbortController(); + let resolveStarted!: (started: boolean) => void; + const started = new Promise((resolve) => { + resolveStarted = resolve; + }); + const dispose = vi.fn(); + const terminateOwnedTree = vi.fn(); + return { + contract: { dispose, signal: controller.signal, started, terminateOwnedTree }, + disconnectAfterStart: () => controller.abort(new Error("worker supervisor lifetime ended")), + disconnectBeforeStart: () => resolveStarted(false), + dispose, + open: () => resolveStarted(true), + terminateOwnedTree, + }; +} + +describe("worker command lifetime gate", () => { + beforeEach(() => { + vi.mocked(runWorkerDescriptor).mockReset(); + vi.mocked(runWorkerDescriptor).mockResolvedValue({ + status: "completed", + transcriptLeafId: null, + transcriptNextSeq: 1, + }); + }); + + it("keeps the ordinary worker command path ungated", async () => { + const output = new PassThrough(); + const chunks: Buffer[] = []; + output.on("data", (chunk: Buffer) => chunks.push(Buffer.from(chunk))); + + await runWorkerCommand({ input: commandInput(), output }); + + expect(runWorkerDescriptor).toHaveBeenCalledOnce(); + expect(JSON.parse(Buffer.concat(chunks).toString("utf8"))).toMatchObject({ + status: "completed", + }); + }); + + it("does not enter the worker runtime before the explicit start message", async () => { + const output = new PassThrough(); + const lifetime = lifetimeHarness(); + const running = runWorkerCommand({ + input: commandInput(), + output, + lifetime: lifetime.contract, + }); + + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(runWorkerDescriptor).not.toHaveBeenCalled(); + lifetime.open(); + + await running; + expect(runWorkerDescriptor).toHaveBeenCalledOnce(); + expect(lifetime.terminateOwnedTree).not.toHaveBeenCalled(); + expect(lifetime.dispose).toHaveBeenCalledOnce(); + }); + + it("exits without starting when IPC disconnects before the start message", async () => { + const output = new PassThrough(); + const lifetime = lifetimeHarness(); + const running = runWorkerCommand({ + input: commandInput(), + output, + lifetime: lifetime.contract, + }); + + lifetime.disconnectBeforeStart(); + + await running; + expect(runWorkerDescriptor).not.toHaveBeenCalled(); + expect(lifetime.terminateOwnedTree).not.toHaveBeenCalled(); + expect(lifetime.dispose).toHaveBeenCalledOnce(); + }); + + it("aborts the real worker path and terminates its owned tree on IPC disconnect", async () => { + const output = new PassThrough(); + let runtimeSignal: AbortSignal | undefined; + vi.mocked(runWorkerDescriptor).mockImplementation(async (_descriptor, options) => { + const signal = options?.signal; + if (!signal) { + throw new Error("expected worker lifetime abort signal"); + } + runtimeSignal = signal; + return await new Promise((_, reject) => { + signal.addEventListener( + "abort", + () => { + const reason = signal.reason; + reject(reason instanceof Error ? reason : new Error("worker interrupted")); + }, + { once: true }, + ); + }); + }); + const lifetime = lifetimeHarness(); + lifetime.terminateOwnedTree.mockImplementation(() => { + expect(runtimeSignal?.aborted).toBe(true); + }); + const running = runWorkerCommand({ + input: commandInput(), + output, + lifetime: lifetime.contract, + }); + lifetime.open(); + await vi.waitFor(() => expect(runWorkerDescriptor).toHaveBeenCalledOnce()); + + lifetime.disconnectAfterStart(); + + await expect(running).rejects.toThrow("worker supervisor lifetime ended"); + expect(lifetime.terminateOwnedTree).toHaveBeenCalledOnce(); + expect(lifetime.dispose).toHaveBeenCalledOnce(); + }); +}); diff --git a/src/worker/worker-command.runtime.ts b/src/worker/worker-command.runtime.ts index 21499d6811f4..6b86ffaf2d8f 100644 --- a/src/worker/worker-command.runtime.ts +++ b/src/worker/worker-command.runtime.ts @@ -5,9 +5,17 @@ import { runWorkerDescriptor } from "./worker.runtime.js"; type RunWorkerCommandOptions = { input: Readable; + lifetime?: WorkerCommandLifetime; output: Writable; }; +export type WorkerCommandLifetime = { + dispose: () => void; + signal: AbortSignal; + started: Promise; + terminateOwnedTree: () => void; +}; + async function readLaunchDescriptor(input: Readable): Promise { const chunks: Buffer[] = []; let byteLength = 0; @@ -41,16 +49,39 @@ async function readLaunchDescriptor(input: Readable): Promise { - const descriptor = await readLaunchDescriptor(options.input); const abortController = new AbortController(); const stop = () => abortController.abort(new Error("worker interrupted")); - process.once("SIGINT", stop); - process.once("SIGTERM", stop); + let lifetimeEnded = false; + const stopForLifetime = () => { + if (lifetimeEnded || !options.lifetime) { + return; + } + lifetimeEnded = true; + abortController.abort( + options.lifetime.signal.reason ?? new Error("worker supervisor lifetime ended"), + ); + options.lifetime.terminateOwnedTree(); + }; try { + const [descriptor, started] = await Promise.all([ + readLaunchDescriptor(options.input), + options.lifetime?.started ?? Promise.resolve(true), + ]); + if (!started) { + return; + } + options.lifetime?.signal.addEventListener("abort", stopForLifetime, { once: true }); + if (options.lifetime?.signal.aborted) { + stopForLifetime(); + } + process.once("SIGINT", stop); + process.once("SIGTERM", stop); const result = await runWorkerDescriptor(descriptor, { signal: abortController.signal }); const encoded = `${JSON.stringify(result)}\n`; options.output.write(encoded); } finally { + options.lifetime?.signal.removeEventListener("abort", stopForLifetime); + options.lifetime?.dispose(); process.off("SIGINT", stop); process.off("SIGTERM", stop); }