From 2ce0573b02eae69ca68d7a65b85534bbda1f0960 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 19 Sep 2026 22:41:17 -0700 Subject: [PATCH] fix(gateway): preserve active work during recovered Node restarts (#153435) * fix(gateway): preserve work during recovered Node restarts Let Unix Node recovery wrappers respect the existing Gateway stop budget. Share startup-safe command classification and deadline facts, and resolve managed restart intent against the live serving owner in its write transaction with current update authority. Preserve the existing public PID intent API, record format, and Windows behavior. Keep cold lifecycle test preparation outside timed hooks. Inspired by @ly85206559's wrapper-grace approach in #147054; this repair is independently authored from main. Related: #146956. Windows cooperative shutdown remains a follow-up. * refactor(gateway): stabilize restart callback and test setup Declare receiver-independent restart intent callbacks as arrows and retain typed signal spy handles for lint-safe assertions. Move cold lifecycle imports into test collection so worker preparation cannot consume timed setup hooks. --- Dockerfile | 2 + cli-root-options.d.mts | 23 ++ cli-root-options.mjs | 176 +++++++++++ docs/gateway/restart-recovery.md | 9 + gateway-run-argv.d.mts | 3 + gateway-run-argv.mjs | 46 +++ gateway-shutdown-budget.d.mts | 5 + gateway-shutdown-budget.mjs | 10 + node-runtime-recovery.mjs | 53 ++-- package.json | 3 + scripts/check-duplicates.mts | 3 + scripts/docker/cleanup-smoke/Dockerfile | 1 + src/cli/daemon-cli/lifecycle-core.test.ts | 66 ++-- src/cli/daemon-cli/lifecycle-core.ts | 38 +-- .../daemon-cli/lifecycle-restart-intent.ts | 73 +++++ .../lifecycle.external-supervision.test.ts | 1 + .../lifecycle.gateway-owner.test.ts | 1 + .../lifecycle.restart-intent.test.ts | 289 ++++++++++++++++++ src/cli/daemon-cli/lifecycle.test.ts | 12 +- src/cli/gateway-run-argv.ts | 45 +-- ...e-command-repair-isolation.test-support.ts | 3 + ...r-config-preflight.process.test-support.ts | 3 + src/daemon/launchd-plist.ts | 4 +- src/entry.compile-cache.test.ts | 14 + src/entry.compile-cache.ts | 5 + .../node-bootstrap-artifact.test-support.ts | 3 + .../node-bootstrap-artifact.ts | 3 + .../node-enrollment.test.ts | 3 + src/infra/cli-root-options.ts | 209 +------------ src/infra/gateway-owner-lease.ts | 100 +++--- src/infra/gateway-shutdown-budget.ts | 15 +- ...ode-runtime-recovery.gateway-drain.test.ts | 93 ++++++ .../node-runtime-recovery.shutdown.test.ts | 85 ++++++ src/infra/restart-intent.ts | 79 ++++- test/node-host-launcher.test.ts | 3 + test/openclaw-launcher-version.e2e.test.ts | 12 + test/openclaw-launcher.e2e.test.ts | 12 + test/scripts/run-node-lifecycle.test.ts | 3 + 38 files changed, 1108 insertions(+), 400 deletions(-) create mode 100644 cli-root-options.d.mts create mode 100644 cli-root-options.mjs create mode 100644 gateway-run-argv.d.mts create mode 100644 gateway-run-argv.mjs create mode 100644 gateway-shutdown-budget.d.mts create mode 100644 gateway-shutdown-budget.mjs create mode 100644 src/cli/daemon-cli/lifecycle-restart-intent.ts create mode 100644 src/cli/daemon-cli/lifecycle.restart-intent.test.ts create mode 100644 src/infra/node-runtime-recovery.gateway-drain.test.ts create mode 100644 src/infra/node-runtime-recovery.shutdown.test.ts diff --git a/Dockerfile b/Dockerfile index be632fc81d5b..00003badd624 100644 --- a/Dockerfile +++ b/Dockerfile @@ -76,6 +76,7 @@ COPY node-version.mjs ./ COPY node-sqlite.mjs ./ COPY node-runtime-update.mjs ./ COPY node-runtime-recovery.mjs ./ +COPY cli-root-options.mjs gateway-run-argv.mjs gateway-shutdown-budget.mjs ./ COPY node-host-launcher.mjs ./ COPY openclaw.mjs ./ COPY ui/package.json ./ui/package.json @@ -281,6 +282,7 @@ COPY --from=runtime-assets --chown=node:node /app/node-version.mjs . COPY --from=runtime-assets --chown=node:node /app/node-sqlite.mjs . COPY --from=runtime-assets --chown=node:node /app/node-runtime-update.mjs . COPY --from=runtime-assets --chown=node:node /app/node-runtime-recovery.mjs . +COPY --from=runtime-assets --chown=node:node /app/cli-root-options.mjs /app/gateway-run-argv.mjs /app/gateway-shutdown-budget.mjs ./ COPY --from=runtime-assets --chown=node:node /app/node-host-launcher.mjs . COPY --from=runtime-assets --chown=node:node /app/openclaw.mjs . COPY --from=runtime-assets --chown=node:node /app/${OPENCLAW_BUNDLED_PLUGIN_DIR} ./${OPENCLAW_BUNDLED_PLUGIN_DIR} diff --git a/cli-root-options.d.mts b/cli-root-options.d.mts new file mode 100644 index 000000000000..ddfc59d61d2f --- /dev/null +++ b/cli-root-options.d.mts @@ -0,0 +1,23 @@ +export const FLAG_TERMINATOR: "--"; + +export function isValueToken(arg: string | undefined): boolean; +export function consumeRootOptionToken(args: ReadonlyArray, index: number): number; +export function consumeRootCommandOptionToken(args: readonly string[], index: number): number; +export function getRootOptionAwareCommandPath(argv: readonly string[], depth: number): string[]; + +type CommandPositionalsParseOptions = { + commandPath: ReadonlyArray; + booleanFlags?: ReadonlyArray; + valueFlags?: ReadonlyArray; + maxPositionals?: number; + mode?: "route" | "command-path"; +}; + +export function getCommandPositionalsWithRootOptions( + argv: readonly string[], + options: CommandPositionalsParseOptions, +): string[] | null; +export function getCommandArgsWithRootOptions( + argv: readonly string[], + options: Omit, +): string[] | null; diff --git a/cli-root-options.mjs b/cli-root-options.mjs new file mode 100644 index 000000000000..21958ee5b197 --- /dev/null +++ b/cli-root-options.mjs @@ -0,0 +1,176 @@ +/** CLI token that stops root option scanning and leaves following args positional. */ +export const FLAG_TERMINATOR = "--"; + +const ROOT_BOOLEAN_FLAGS = new Set(["--dev", "--no-color"]); +const ROOT_VALUE_FLAGS = new Set(["--profile", "--log-level", "--container"]); + +/** Returns whether a token can be consumed as a root option value. */ +export function isValueToken(arg) { + if (!arg || arg === FLAG_TERMINATOR) { + return false; + } + if (!arg.startsWith("-")) { + return true; + } + return /^-\d+(?:\.\d+)?$/.test(arg); +} + +/** Count root-option tokens conservatively for route matching. */ +export function consumeRootOptionToken(args, index) { + const arg = args[index]; + if (!arg) { + return 0; + } + if (ROOT_BOOLEAN_FLAGS.has(arg)) { + return 1; + } + if ( + arg.startsWith("--profile=") || + arg.startsWith("--log-level=") || + arg.startsWith("--container=") + ) { + return 1; + } + if (ROOT_VALUE_FLAGS.has(arg)) { + return isValueToken(args[index + 1]) ? 2 : 1; + } + return 0; +} + +/** Consume required root values by their Commander role before startup policy is selected. */ +export function consumeRootCommandOptionToken(args, index) { + return consumeKnownOptionToken(args, index, ROOT_BOOLEAN_FLAGS, ROOT_VALUE_FLAGS, "command-path"); +} + +/** Read positional command tokens while accepting root options at any pre-terminator position. */ +export function getRootOptionAwareCommandPath(argv, depth) { + const args = argv.slice(2); + const path = []; + let literal = false; + for (let index = 0; index < args.length; index += 1) { + const arg = args[index]; + if (!arg) { + break; + } + if (!literal && arg === FLAG_TERMINATOR) { + // A leading terminator still leaves a command to discover; later operands belong to callers. + if (path.length > 0) { + break; + } + literal = true; + continue; + } + const consumed = literal ? 0 : consumeRootCommandOptionToken(args, index); + if (consumed > 0) { + index += consumed - 1; + continue; + } + if (!literal && arg.startsWith("-")) { + continue; + } + path.push(arg); + if (path.length >= depth) { + break; + } + } + return path; +} + +function consumeKnownOptionToken(args, index, booleanFlags, valueFlags, mode) { + const arg = args[index]; + if (!arg || arg === FLAG_TERMINATOR || !arg.startsWith("-")) { + return 0; + } + + const equalsIndex = arg.indexOf("="); + const flag = equalsIndex === -1 ? arg : arg.slice(0, equalsIndex); + if (booleanFlags.has(flag)) { + return equalsIndex === -1 ? 1 : 0; + } + if (!valueFlags.has(flag)) { + return 0; + } + if (equalsIndex !== -1) { + return mode === "command-path" || arg.slice(equalsIndex + 1).trim() ? 1 : 0; + } + // Required Commander values include empty strings, flag-looking tokens, and `--`. + // Discovery must consume them before choosing startup policy; routes still validate values. + if (mode === "command-path") { + return args[index + 1] !== undefined ? 2 : 0; + } + return isValueToken(args[index + 1]) ? 2 : 0; +} + +/** Parse command positionals while consuming known root and command options. */ +export function getCommandPositionalsWithRootOptions(argv, options) { + return parseCommandArgsWithRootOptions(argv, options, false); +} + +/** Preserve the leaf's raw arguments after consuming its root and parent options. */ +export function getCommandArgsWithRootOptions(argv, options) { + return parseCommandArgsWithRootOptions(argv, options, true); +} + +function parseCommandArgsWithRootOptions(argv, options, returnTail) { + const args = argv.slice(2); + const booleanFlags = new Set(options.booleanFlags ?? []); + const valueFlags = new Set(options.valueFlags ?? []); + const positionals = []; + let commandIndex = 0; + let literal = false; + + for (let index = 0; index < args.length; index += 1) { + const arg = args[index]; + if (arg === undefined) { + break; + } + if (!literal && arg === FLAG_TERMINATOR) { + if (options.mode !== "command-path") { + break; + } + literal = true; + continue; + } + const rootConsumed = literal + ? 0 + : options.mode === "command-path" + ? consumeRootCommandOptionToken(args, index) + : consumeRootOptionToken(args, index); + if (rootConsumed > 0) { + index += rootConsumed - 1; + continue; + } + if (!literal && arg.startsWith("-")) { + const optionConsumed = consumeKnownOptionToken( + args, + index, + booleanFlags, + valueFlags, + options.mode, + ); + if (optionConsumed === 0 || commandIndex === 0) { + return null; + } + index += optionConsumed - 1; + continue; + } + if (commandIndex < options.commandPath.length) { + if (arg !== options.commandPath[commandIndex]) { + return null; + } + commandIndex += 1; + if (returnTail && commandIndex === options.commandPath.length) { + const tail = args.slice(index + 1); + // A downstream parser must not reactivate flags after an earlier literal boundary. + return literal ? [FLAG_TERMINATOR, ...tail] : tail; + } + continue; + } + positionals.push(arg); + if (positionals.length === options.maxPositionals) { + return positionals; + } + } + + return commandIndex < options.commandPath.length ? null : positionals; +} diff --git a/docs/gateway/restart-recovery.md b/docs/gateway/restart-recovery.md index 1cd242afcd85..d1c96eb74a6e 100644 --- a/docs/gateway/restart-recovery.md +++ b/docs/gateway/restart-recovery.md @@ -123,6 +123,15 @@ gateway stops accepting new work, then waits for active agent turns and background tasks to finish, up to a drain budget (5 minutes by default). Most restarts therefore interrupt nothing at all. +On Linux and macOS, this also applies when startup recovers from an unsupported +Node version and the service manager tracks a launcher parent. The launcher +forwards the stop signal and waits for the serving Gateway to drain within the +shared service budget. Managed restart intent targets the live serving owner, +so unfinished work still follows restart recovery when its drain budget expires. +The launchd stop budget remains 20 seconds; Linux units use the deadlines below. +This requires a Gateway started with the updated launcher: replacing files cannot +change a launcher that is already running. + On Linux, the systemd unit must use `KillMode=mixed` so the initial stop signal reaches only the Gateway. Systemd still kills remaining child processes when the Gateway exits or its stop deadline expires. Older `KillMode=control-group` units diff --git a/gateway-run-argv.d.mts b/gateway-run-argv.d.mts new file mode 100644 index 000000000000..d3237ca0488d --- /dev/null +++ b/gateway-run-argv.d.mts @@ -0,0 +1,3 @@ +export const GATEWAY_RUN_VALUE_FLAGS: ReadonlySet; +export const GATEWAY_RUN_BOOLEAN_FLAGS: ReadonlySet; +export function isForegroundGatewayRunArgv(argv: string[]): boolean; diff --git a/gateway-run-argv.mjs b/gateway-run-argv.mjs new file mode 100644 index 000000000000..b98823063aff --- /dev/null +++ b/gateway-run-argv.mjs @@ -0,0 +1,46 @@ +// Startup-safe Gateway command-role parsing, shared with the CLI. +import { getCommandPositionalsWithRootOptions } from "./cli-root-options.mjs"; + +export const GATEWAY_RUN_VALUE_FLAGS = new Set([ + "--port", + "--bind", + "--token", + "--token-file", + "--auth", + "--password", + "--password-file", + "--tailscale", + "--ws-log", + "--raw-stream-path", +]); + +export const GATEWAY_RUN_BOOLEAN_FLAGS = new Set([ + "--tailscale-reset-on-exit", + "--allow-unconfigured", + "--dev", + "--ambient-channels", + "--dev-ambient-channels", + "--reset", + "--update-canary", + "--force", + "--verbose", + "--cli-backend-logs", + "--claude-cli-logs", + "--compact", + "--raw-stream", +]); + +export function isForegroundGatewayRunArgv(argv) { + const positionals = getCommandPositionalsWithRootOptions(argv, { + commandPath: ["gateway"], + booleanFlags: [...GATEWAY_RUN_BOOLEAN_FLAGS], + valueFlags: [...GATEWAY_RUN_VALUE_FLAGS], + mode: "command-path", + }); + if (!positionals) { + return false; + } + // Foreground gateway owns the terminal/process environment itself; respawning would + // add an extra parent process around the long-lived server. + return positionals.length === 0 || (positionals.length === 1 && positionals[0] === "run"); +} diff --git a/gateway-shutdown-budget.d.mts b/gateway-shutdown-budget.d.mts new file mode 100644 index 000000000000..03637034718f --- /dev/null +++ b/gateway-shutdown-budget.d.mts @@ -0,0 +1,5 @@ +export const GATEWAY_SHUTDOWN_RESERVE_MS: number; +export const GATEWAY_SUPERVISOR_EXIT_MARGIN_MS: number; +export const GATEWAY_SHUTDOWN_TIMEOUT_MS: number; +export const GATEWAY_SERVICE_STOP_TIMEOUT_MS: number; +export const LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS: 20; diff --git a/gateway-shutdown-budget.mjs b/gateway-shutdown-budget.mjs new file mode 100644 index 000000000000..b8aec66c937f --- /dev/null +++ b/gateway-shutdown-budget.mjs @@ -0,0 +1,10 @@ +// Shared Gateway stop policy used by the run loop and native service definitions. +const GATEWAY_SHUTDOWN_DRAIN_TIMEOUT_MS = 315_000; +export const GATEWAY_SHUTDOWN_RESERVE_MS = 10_000; +export const GATEWAY_SUPERVISOR_EXIT_MARGIN_MS = 5_000; +export const GATEWAY_SHUTDOWN_TIMEOUT_MS = + GATEWAY_SHUTDOWN_DRAIN_TIMEOUT_MS + GATEWAY_SHUTDOWN_RESERVE_MS; +export const GATEWAY_SERVICE_STOP_TIMEOUT_MS = + GATEWAY_SHUTDOWN_TIMEOUT_MS + GATEWAY_SUPERVISOR_EXIT_MARGIN_MS; + +export const LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS = 20; diff --git a/node-runtime-recovery.mjs b/node-runtime-recovery.mjs index ab6bd5049358..91c7f8c895c5 100644 --- a/node-runtime-recovery.mjs +++ b/node-runtime-recovery.mjs @@ -3,47 +3,21 @@ import { spawn, spawnSync } from "node:child_process"; import { lstatSync, readFileSync, readdirSync, realpathSync, statSync } from "node:fs"; import os from "node:os"; import path from "node:path"; +import { consumeRootOptionToken as consumeLauncherRootOptionToken } from "./cli-root-options.mjs"; +import { isForegroundGatewayRunArgv } from "./gateway-run-argv.mjs"; +import { + GATEWAY_SERVICE_STOP_TIMEOUT_MS, + LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS, +} from "./gateway-shutdown-budget.mjs"; import { detectCurrentSqliteCapabilities, nodeRuntimeFailure, SQLITE_CAPABILITY_PROBE, } from "./node-sqlite.mjs"; -const LAUNCHER_ROOT_BOOLEAN_FLAGS = new Set(["--dev", "--no-color"]); -const LAUNCHER_ROOT_VALUE_FLAGS = new Set(["--profile", "--log-level", "--container"]); +export { consumeLauncherRootOptionToken }; export const isNativeHookRelayInvocation = (argv) => argv[2] === "hooks" && argv[3] === "relay"; -const isLauncherRootOptionValueToken = (arg) => { - if (!arg || arg === "--") { - return false; - } - if (!arg.startsWith("-")) { - return true; - } - return /^-\d+(?:\.\d+)?$/.test(arg); -}; - -export const consumeLauncherRootOptionToken = (args, index) => { - const arg = args[index]; - if (!arg) { - return 0; - } - if (LAUNCHER_ROOT_BOOLEAN_FLAGS.has(arg)) { - return 1; - } - if ( - arg.startsWith("--profile=") || - arg.startsWith("--log-level=") || - arg.startsWith("--container=") - ) { - return 1; - } - if (LAUNCHER_ROOT_VALUE_FLAGS.has(arg)) { - return isLauncherRootOptionValueToken(args[index + 1]) ? 2 : 1; - } - return 0; -}; - // Mirror the entry's foreground Gmail policy: a wrapper would kill that run before descendant cleanup finishes. export const isForegroundGmailRunInvocation = (argv) => { const args = argv.slice(2); @@ -70,6 +44,17 @@ const respawnSignalForceKillGraceMs = 1_000; const respawnSignalHardExitGraceMs = 1_000; export const runRespawnedChild = (command, args, env) => { + const launchdService = env.OPENCLAW_LAUNCHD_LABEL?.trim(); + const serviceStopTimeoutMs = + process.platform === "darwin" && launchdService && env.XPC_SERVICE_NAME === launchdService + ? LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS * 1_000 + : GATEWAY_SERVICE_STOP_TIMEOUT_MS; + // The serving Gateway owns drain and cleanup. Reap a stuck child only in the + // supervisor's exit margin, after that owner has had its full shutdown budget. + const signalExitGraceMs = + process.platform !== "win32" && isForegroundGatewayRunArgv(process.argv) + ? serviceStopTimeoutMs - respawnSignalForceKillGraceMs - respawnSignalHardExitGraceMs + : respawnSignalExitGraceMs; const stdioIsTerminal = process.stdin.isTTY || process.stdout.isTTY; const child = spawn(command, args, { stdio: "inherit", @@ -131,7 +116,7 @@ export const runRespawnedChild = (command, args, env) => { } signalExitTimer = setTimeout(() => { requestChildTermination(); - }, respawnSignalExitGraceMs); + }, signalExitGraceMs); signalExitTimer.unref?.(); }; for (const signal of respawnSignals) { diff --git a/package.json b/package.json index d7ff67a2402b..9a9f9654a03b 100644 --- a/package.json +++ b/package.json @@ -37,6 +37,9 @@ "files": [ "CHANGELOG.md", "LICENSE", + "cli-root-options.mjs", + "gateway-run-argv.mjs", + "gateway-shutdown-budget.mjs", "node-sqlite.mjs", "node-version.mjs", "node-runtime-update.mjs", diff --git a/scripts/check-duplicates.mts b/scripts/check-duplicates.mts index b091b468cf9c..4f1bc4e38fbd 100644 --- a/scripts/check-duplicates.mts +++ b/scripts/check-duplicates.mts @@ -22,6 +22,9 @@ const targets = [ "test", "skills", "config", + "cli-root-options.mjs", + "gateway-run-argv.mjs", + "gateway-shutdown-budget.mjs", "node-host-launcher.mjs", "node-runtime-update.mjs", "node-runtime-recovery.mjs", diff --git a/scripts/docker/cleanup-smoke/Dockerfile b/scripts/docker/cleanup-smoke/Dockerfile index 50e27d6972f5..2388182f183e 100644 --- a/scripts/docker/cleanup-smoke/Dockerfile +++ b/scripts/docker/cleanup-smoke/Dockerfile @@ -19,6 +19,7 @@ COPY node-version.mjs ./ COPY node-sqlite.mjs ./ COPY node-runtime-update.mjs ./ COPY node-runtime-recovery.mjs ./ +COPY cli-root-options.mjs gateway-run-argv.mjs gateway-shutdown-budget.mjs ./ COPY node-host-launcher.mjs ./ COPY ui/package.json ./ui/package.json COPY packages ./packages diff --git a/src/cli/daemon-cli/lifecycle-core.test.ts b/src/cli/daemon-cli/lifecycle-core.test.ts index a14752a92749..16f7dc4a8f06 100644 --- a/src/cli/daemon-cli/lifecycle-core.test.ts +++ b/src/cli/daemon-cli/lifecycle-core.test.ts @@ -1,5 +1,5 @@ // Daemon lifecycle core tests cover service lifecycle transitions and platform adapters. -import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import type { OpenClawConfig } from "../../config/config.js"; import type { GatewayServiceControlArgs } from "../../daemon/service-types.js"; import type { GatewayService } from "../../daemon/service.js"; @@ -55,6 +55,7 @@ vi.mock("../../runtime.js", () => ({ vi.mock("../../infra/restart-intent.js", () => ({ clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(), writeGatewayRestartIntentSync: (opts: unknown) => writeGatewayRestartIntentSync(opts), + writeGatewayServiceRestartIntentSync: (opts: unknown) => writeGatewayRestartIntentSync(opts), })); vi.mock("./lifecycle-audit.js", () => ({ @@ -79,10 +80,8 @@ vi.mock("./lifecycle-audit.js", () => ({ }, })); -let runServiceRestart: typeof import("./lifecycle-core.js").runServiceRestart; -let runServiceStart: typeof import("./lifecycle-core.js").runServiceStart; -let runServiceStop: typeof import("./lifecycle-core.js").runServiceStop; -let runServiceUninstall: typeof import("./lifecycle-core.js").runServiceUninstall; +const { runServiceRestart, runServiceStart, runServiceStop, runServiceUninstall } = + await import("./lifecycle-core.js"); // oxlint-disable-next-line typescript/no-unnecessary-type-parameters -- Test helper lets assertions ascribe logged JSON shape. function readJsonLog() { @@ -141,11 +140,6 @@ function expectUnsupportedServiceCheckFailure() { } describe("runServiceRestart token drift", () => { - beforeAll(async () => { - ({ runServiceRestart, runServiceStart, runServiceStop, runServiceUninstall } = - await import("./lifecycle-core.js")); - }); - afterEach(() => { vi.restoreAllMocks(); }); @@ -479,11 +473,14 @@ describe("runServiceRestart token drift", () => { }), ); expect(service.restart).not.toHaveBeenCalled(); - expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({ - targetPid: 1234, - reason: "gateway.restart", - intent: { waitMs: 2_500 }, - }); + expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith( + expect.objectContaining({ + env: expect.any(Object), + targetPid: 1234, + reason: "gateway.restart", + intent: { waitMs: 2_500 }, + }), + ); expect(readJsonLog<{ result?: string; message?: string }>()).toMatchObject({ result: "restarted", message: "Gateway service definition repaired and restarted.", @@ -729,10 +726,13 @@ describe("runServiceRestart token drift", () => { await runServiceRestart(createServiceRunArgs()); - expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({ - targetPid: 1234, - reason: "gateway.restart", - }); + expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith( + expect.objectContaining({ + env: expect.any(Object), + targetPid: 1234, + reason: "gateway.restart", + }), + ); expect(clearGatewayRestartIntentSync).not.toHaveBeenCalled(); expect(service.restart).toHaveBeenCalledTimes(1); }); @@ -769,13 +769,16 @@ describe("runServiceRestart token drift", () => { }, }); - expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({ - targetPid: 1234, - reason: "gateway.restart", - intent: { - waitMs: 2_500, - }, - }); + expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith( + expect.objectContaining({ + env: expect.any(Object), + targetPid: 1234, + reason: "gateway.restart", + intent: { + waitMs: 2_500, + }, + }), + ); }); it("clears restart intent when service-manager restart fails before signaling", async () => { @@ -785,10 +788,13 @@ describe("runServiceRestart token drift", () => { await expect(runServiceRestart(createServiceRunArgs())).rejects.toThrow("__exit__:1"); - expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({ - targetPid: 1234, - reason: "gateway.restart", - }); + expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith( + expect.objectContaining({ + env: expect.any(Object), + targetPid: 1234, + reason: "gateway.restart", + }), + ); expect(clearGatewayRestartIntentSync).toHaveBeenCalledOnce(); }); diff --git a/src/cli/daemon-cli/lifecycle-core.ts b/src/cli/daemon-cli/lifecycle-core.ts index 7c9a97434cdc..4f3ef8495092 100644 --- a/src/cli/daemon-cli/lifecycle-core.ts +++ b/src/cli/daemon-cli/lifecycle-core.ts @@ -4,7 +4,6 @@ import { readBestEffortConfig } from "../../config/config.js"; import { resolveIsNixMode } from "../../config/paths.js"; import { checkTokenDrift } from "../../daemon/service-audit.js"; import type { GatewayServiceRestartResult } from "../../daemon/service-types.js"; -import { assertGatewayServiceUpdateCurrent } from "../../daemon/service-update-authority.js"; import type { GatewayServiceStartRepairIssue, GatewayServiceState, @@ -19,11 +18,7 @@ import { import { renderSystemdUnavailableHints } from "../../daemon/systemd-hints.js"; import { isSystemdUserServiceAvailable } from "../../daemon/systemd.js"; import { isGatewaySecretRefUnavailableError } from "../../gateway/credentials.js"; -import { - clearGatewayRestartIntentSync, - type GatewayRestartIntent, - writeGatewayRestartIntentSync, -} from "../../infra/restart-intent.js"; +import type { GatewayRestartIntent } from "../../infra/restart-intent.js"; import { isWSL } from "../../infra/wsl.js"; import { defaultRuntime } from "../../runtime.js"; import { formatCliCommand } from "../command-format.js"; @@ -34,6 +29,7 @@ import { appendServiceLifecycleRepairAudit, createServiceLifecycleMutationAudit, } from "./lifecycle-audit.js"; +import { createServiceRestartIntent } from "./lifecycle-restart-intent.js"; import { buildDaemonServiceSnapshot, createDaemonActionContext, @@ -498,26 +494,18 @@ export async function runServiceRestart(params: { let handledRecovery: ServiceRecoveryResult<"restarted"> | null = null; let handledRepair: ServiceRecoveryResult<"restarted"> | null = null; let recoveredLoadedState: boolean | null = null; - let wroteRestartIntent = false; - const prepareGatewayRestartIntent = async () => { - if (params.serviceNoun !== "Gateway" || wroteRestartIntent) { - return; - } - const runtime = await params.service.readRuntime(process.env).catch(() => null); - assertGatewayServiceUpdateCurrent(); - wroteRestartIntent = writeGatewayRestartIntentSync({ - targetPid: runtime?.pid, - reason: "gateway.restart", - ...(restartIntent ? { intent: restartIntent } : {}), + const { prepare: prepareGatewayRestartIntent, clear: clearPreparedRestartIntent } = + createServiceRestartIntent({ + serviceNoun: params.serviceNoun, + service: params.service, + intent: restartIntent, + warn: (message) => { + warnings.push(message); + if (!json) { + defaultRuntime.log(message); + } + }, }); - }; - const clearPreparedRestartIntent = () => { - if (wroteRestartIntent) { - assertGatewayServiceUpdateCurrent(); - clearGatewayRestartIntentSync(); - wroteRestartIntent = false; - } - }; const emitScheduledRestart = ( restartStatus: ReturnType, serviceLoaded: boolean, diff --git a/src/cli/daemon-cli/lifecycle-restart-intent.ts b/src/cli/daemon-cli/lifecycle-restart-intent.ts new file mode 100644 index 000000000000..f0e6df907114 --- /dev/null +++ b/src/cli/daemon-cli/lifecycle-restart-intent.ts @@ -0,0 +1,73 @@ +import { resolveLaunchAgentLabel } from "../../daemon/launchd-label.js"; +import { mergeGatewayServiceEnv } from "../../daemon/service-env-merge.js"; +import { assertGatewayServiceUpdateCurrent } from "../../daemon/service-update-authority.js"; +import type { GatewayService } from "../../daemon/service.js"; +import { resolveSystemdServiceName } from "../../daemon/systemd-service-files.js"; +import { + clearGatewayRestartIntentSync, + type GatewayRestartIntent, + type GatewayRestartIntentService, + writeGatewayRestartIntentSync, + writeGatewayServiceRestartIntentSync, +} from "../../infra/restart-intent.js"; + +export function createServiceRestartIntent(params: { + serviceNoun: string; + service: GatewayService; + intent?: GatewayRestartIntent; + warn: (message: string) => void; +}) { + let recorded = false; + let env = process.env; + return { + prepare: async () => { + if (params.serviceNoun !== "Gateway" || recorded) { + return; + } + const runtime = await params.service.readRuntime(process.env).catch(() => null); + assertGatewayServiceUpdateCurrent(); + const nativeService = process.platform === "linux" || process.platform === "darwin"; + let service: GatewayRestartIntentService | undefined; + if (nativeService) { + try { + const command = await params.service.readCommand(process.env, { requireEffective: true }); + assertGatewayServiceUpdateCurrent(); + env = mergeGatewayServiceEnv(process.env, command); + service = + process.platform === "linux" + ? { + kind: "systemd", + name: runtime?.systemd?.unit ?? resolveSystemdServiceName(process.env), + } + : { kind: "launchd", name: resolveLaunchAgentLabel(process.env) }; + } catch { + params.warn( + "Could not verify the serving Gateway owner; using native service status for restart intent.", + ); + } + } + assertGatewayServiceUpdateCurrent(); + const options = { + env, + targetPid: runtime?.pid, + reason: "gateway.restart", + ...(params.intent ? { intent: params.intent } : {}), + }; + recorded = nativeService + ? writeGatewayServiceRestartIntentSync({ + ...options, + service, + assertCurrent: assertGatewayServiceUpdateCurrent, + warn: params.warn, + }) + : writeGatewayRestartIntentSync(options); + }, + clear: () => { + if (recorded) { + assertGatewayServiceUpdateCurrent(); + clearGatewayRestartIntentSync(env); + recorded = false; + } + }, + }; +} diff --git a/src/cli/daemon-cli/lifecycle.external-supervision.test.ts b/src/cli/daemon-cli/lifecycle.external-supervision.test.ts index 310dd5a8dc06..6163af5af8c8 100644 --- a/src/cli/daemon-cli/lifecycle.external-supervision.test.ts +++ b/src/cli/daemon-cli/lifecycle.external-supervision.test.ts @@ -60,6 +60,7 @@ vi.mock("../../infra/gateway-lock.js", async (importOriginal) => { vi.mock("../../infra/restart-intent.js", () => ({ writeGatewayRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params), + writeGatewayServiceRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params), clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(), })); diff --git a/src/cli/daemon-cli/lifecycle.gateway-owner.test.ts b/src/cli/daemon-cli/lifecycle.gateway-owner.test.ts index dc2c338189dc..cb9a254287b1 100644 --- a/src/cli/daemon-cli/lifecycle.gateway-owner.test.ts +++ b/src/cli/daemon-cli/lifecycle.gateway-owner.test.ts @@ -79,6 +79,7 @@ vi.mock("./lifecycle-audit.js", () => ({ vi.mock("../../infra/restart-intent.js", async (original) => ({ ...(await original()), writeGatewayRestartIntentSync: () => true, + writeGatewayServiceRestartIntentSync: () => true, clearGatewayRestartIntentSync: vi.fn(), })); diff --git a/src/cli/daemon-cli/lifecycle.restart-intent.test.ts b/src/cli/daemon-cli/lifecycle.restart-intent.test.ts new file mode 100644 index 000000000000..71596cca97e7 --- /dev/null +++ b/src/cli/daemon-cli/lifecycle.restart-intent.test.ts @@ -0,0 +1,289 @@ +import { hostname } from "node:os"; +import { afterEach, beforeEach, expect, it, vi } from "vitest"; +import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; +import { withGatewayServiceUpdateAuthority } from "../../daemon/service-update-authority.js"; +import { + acquireGatewayOwnerLease, + type GatewayOwnerSupervisor, +} from "../../infra/gateway-owner-lease.js"; +import { consumeGatewayRestartIntentPayloadSync } from "../../infra/restart-intent.js"; +import { acquireGatewayLifecycleCoordinator } from "../../infra/state-database-coordinator.js"; +import * as processOwners from "../../infra/state-lease-process-owner.js"; +import * as existingWrites from "../../state/openclaw-state-db-existing-write.js"; +import { + closeOpenClawStateDatabaseForTest, + openOpenClawStateDatabase, +} from "../../state/openclaw-state-db.js"; +import { resolveOpenClawStateSqlitePath } from "../../state/openclaw-state-db.paths.js"; +import { + createGatewayServiceRunArgs, + lifecycleTestRuntime, + lifecycleRuntimeLogs, + resetLifecycleRuntimeLogs, + resetLifecycleServiceMocks, + service, +} from "./test-helpers/lifecycle-core-harness.js"; + +vi.mock("../../runtime.js", () => ({ defaultRuntime: lifecycleTestRuntime })); +vi.mock("./lifecycle-action-preflight.js", () => ({ + getServiceActionPreflightFailure: async () => null, +})); +vi.mock("./lifecycle-audit.js", () => ({ + createServiceLifecycleMutationAudit: () => undefined, + appendServiceLifecycleRepairAudit: vi.fn(), +})); + +const tempDirs = useAutoCleanupTempDirTracker(afterEach); +const { runServiceRestart } = await import("./lifecycle-core.js"); + +function beforeIntentWriteAdmission(operation: () => void) { + const write = existingWrites.runExistingOpenClawStateWriteTransaction; + vi.spyOn(existingWrites, "runExistingOpenClawStateWriteTransaction").mockImplementation( + (mutate, options, contract) => { + if (contract.operationLabel === "gateway.restart-intent.write") { + operation(); + } + return write(mutate, options, contract); + }, + ); +} + +function publishServingOwner( + pid: number, + mode: "supervised" | "foreground" = "supervised", + supervisor: GatewayOwnerSupervisor | null = { + kind: "systemd", + name: "openclaw-gateway.service", + }, +) { + const { db } = openOpenClawStateDatabase(); + const now = Date.now(); + db.prepare( + `INSERT INTO state_leases + (scope, lease_key, owner, expires_at, heartbeat_at, payload_json, created_at, updated_at) + VALUES ('gateway-owner', 'global', ?, ?, ?, ?, ?, ?) + ON CONFLICT(scope, lease_key) DO UPDATE SET + owner = excluded.owner, payload_json = excluded.payload_json`, + ).run( + `generation-${pid}`, + now + 60_000, + now, + JSON.stringify({ + owner: { pid, host: hostname(), startedAt: 1 }, + port: 18789, + mode, + supervisor, + }), + now, + now, + ); +} + +beforeEach(() => { + resetLifecycleRuntimeLogs(); + resetLifecycleServiceMocks(); + vi.stubEnv("OPENCLAW_STATE_DIR", tempDirs.make("openclaw-restart-intent-cli-")); + vi.stubEnv("OPENCLAW_PROFILE", "default"); + vi.stubEnv("OPENCLAW_SYSTEMD_UNIT", ""); + vi.stubEnv("OPENCLAW_LAUNCHD_LABEL", ""); + service.readRuntime.mockResolvedValue({ status: "running", pid: process.pid + 1 }); +}); + +afterEach(() => { + closeOpenClawStateDatabaseForTest(); + vi.restoreAllMocks(); + vi.unstubAllEnvs(); +}); + +it.each([ + { + platform: "linux", + supervisor: { kind: "systemd", name: "openclaw-gateway.service" }, + serviceState: false, + }, + { + platform: "darwin", + supervisor: { kind: "launchd", name: "ai.openclaw.gateway" }, + serviceState: false, + }, + { + platform: "linux", + supervisor: { kind: "systemd", name: "openclaw-gateway.service" }, + serviceState: true, + }, + { + platform: "linux", + supervisor: { kind: "systemd", name: "openclaw-gateway" }, + serviceState: false, + }, +] as const)( + "delivers the managed $platform/$supervisor.name restart intent to the serving owner (service state=$serviceState)", + async ({ platform, supervisor, serviceState }) => { + const stateDir = tempDirs.make("openclaw-serving-state-"); + const env = { ...process.env, OPENCLAW_STATE_DIR: stateDir }; + if (!serviceState) { + vi.stubEnv("OPENCLAW_STATE_DIR", stateDir); + } + const coordinator = acquireGatewayLifecycleCoordinator({ + databasePath: resolveOpenClawStateSqlitePath(env), + }); + const lease = acquireGatewayOwnerLease({ + env, + port: 18789, + mode: "supervised", + supervisor, + }); + try { + await lease.ready; + vi.spyOn(process, "platform", "get").mockReturnValue(platform); + service.readRuntime.mockResolvedValue({ + status: "running", + pid: process.pid + 1, + ...(platform === "linux" ? { systemd: { unit: "openclaw-gateway.service" } } : {}), + }); + service.readCommand.mockResolvedValue({ + programArguments: ["node", "openclaw.mjs", "gateway", "run"], + environment: { OPENCLAW_STATE_DIR: stateDir }, + }); + let consumed: ReturnType = null; + service.restart.mockImplementationOnce(async () => { + consumed = consumeGatewayRestartIntentPayloadSync(env); + return { outcome: "completed" }; + }); + + await expect( + runServiceRestart({ + ...createGatewayServiceRunArgs(), + opts: { json: true, preserveDefinition: true, restartIntent: { waitMs: 30_000 } }, + }), + ).resolves.toBe(true); + + expect(consumed).toEqual({ reason: "gateway.restart", waitMs: 30_000 }); + expect(consumeGatewayRestartIntentPayloadSync(env)).toBeNull(); + } finally { + await lease.release(); + coordinator.release(); + } + }, +); + +it("targets the replacement published before restart-intent write admission", async () => { + vi.spyOn(process, "platform", "get").mockReturnValue("linux"); + // Synthetic process identities isolate the race; lease storage and intent consumption are real. + vi.spyOn(processOwners, "readStateLeaseProcessOwnerStatus").mockReturnValue("live"); + const coordinator = acquireGatewayLifecycleCoordinator({ + databasePath: resolveOpenClawStateSqlitePath(process.env), + }); + try { + publishServingOwner(process.pid + 2); + let publications = 0; + beforeIntentWriteAdmission(() => { + publishServingOwner(process.pid); + publications += 1; + }); + let consumed: ReturnType = null; + service.restart.mockImplementationOnce(async () => { + consumed = consumeGatewayRestartIntentPayloadSync(); + return { outcome: "completed" }; + }); + + await expect(runServiceRestart(createGatewayServiceRunArgs())).resolves.toBe(true); + + expect(publications).toBe(1); + expect(consumed).toEqual({ reason: "gateway.restart" }); + } finally { + coordinator.release(); + } +}); + +it.each([ + { state: "dead", mode: "supervised", name: "openclaw-gateway.service" }, + { state: "unknown", mode: "supervised", name: "openclaw-gateway.service" }, + { state: "live", mode: "foreground", name: "openclaw-gateway.service" }, + { state: "live", mode: "supervised", name: "another-gateway.service" }, + { state: "live", mode: "supervised", name: null }, +] as const)( + "keeps native targeting for an unrelated or unverifiable owner ($state/$mode/$name)", + async ({ state, mode, name }) => { + vi.spyOn(process, "platform", "get").mockReturnValue("linux"); + vi.spyOn(processOwners, "readStateLeaseProcessOwnerStatus").mockReturnValue(state); + const coordinator = acquireGatewayLifecycleCoordinator({ + databasePath: resolveOpenClawStateSqlitePath(process.env), + }); + try { + publishServingOwner( + process.pid + 2, + mode, + mode === "supervised" ? { kind: "systemd", name } : null, + ); + service.readRuntime.mockResolvedValue({ status: "running", pid: process.pid }); + + await expect(runServiceRestart(createGatewayServiceRunArgs())).resolves.toBe(true); + + expect(consumeGatewayRestartIntentPayloadSync()).toEqual({ reason: "gateway.restart" }); + expect(service.restart).toHaveBeenCalledOnce(); + } finally { + coordinator.release(); + } + }, +); + +it.each([true, false])( + "warns and preserves native restart when serving ownership cannot be inspected (json=%s)", + async (json) => { + vi.spyOn(process, "platform", "get").mockReturnValue("linux"); + openOpenClawStateDatabase(); + service.readRuntime.mockResolvedValue({ status: "running", pid: process.pid }); + service.readCommand.mockRejectedValue(new Error("native command inspection unavailable")); + + await expect( + runServiceRestart({ ...createGatewayServiceRunArgs(), opts: { json } }), + ).resolves.toBe(true); + + expect(consumeGatewayRestartIntentPayloadSync()).toEqual({ reason: "gateway.restart" }); + expect(lifecycleRuntimeLogs.join("\n")).toContain( + "Could not verify the serving Gateway owner; using native service status for restart intent.", + ); + expect(service.restart).toHaveBeenCalledOnce(); + }, +); + +it.each(["runtime", "command", "write admission"])( + "revalidates update authority after native %s inspection before recording intent", + async (inspection) => { + vi.spyOn(process, "platform", "get").mockReturnValue("linux"); + const { db } = openOpenClawStateDatabase(); + let current = true; + if (inspection === "runtime") { + service.readRuntime.mockImplementationOnce(async () => { + current = false; + return { status: "running", pid: process.pid }; + }); + } else if (inspection === "command") { + service.readCommand.mockImplementationOnce(async () => { + current = false; + return { programArguments: [] }; + }); + } else { + beforeIntentWriteAdmission(() => { + current = false; + }); + } + + await expect( + withGatewayServiceUpdateAuthority( + () => { + if (!current) { + throw new Error("update owner revoked"); + } + }, + () => runServiceRestart(createGatewayServiceRunArgs()), + ), + ).rejects.toThrow(); + + expect(service.restart).not.toHaveBeenCalled(); + expect(db.prepare("SELECT count(*) AS count FROM gateway_restart_intent").get()).toEqual({ + count: 0, + }); + }, +); diff --git a/src/cli/daemon-cli/lifecycle.test.ts b/src/cli/daemon-cli/lifecycle.test.ts index 590d1d27245a..863e9e6952da 100644 --- a/src/cli/daemon-cli/lifecycle.test.ts +++ b/src/cli/daemon-cli/lifecycle.test.ts @@ -1,5 +1,5 @@ // Daemon lifecycle tests cover CLI service lifecycle orchestration and cleanup. -import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { mockSystemAccountHome } from "../../daemon/service.test-helpers.js"; import { captureEnv } from "../../test-utils/env.js"; import { @@ -115,6 +115,7 @@ vi.mock("../../infra/gateway-owner-lease.js", () => ({ readGatewayOwnerLease })) vi.mock("../../infra/restart-intent.js", () => ({ writeGatewayRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params), + writeGatewayServiceRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params), clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(), })); @@ -181,10 +182,9 @@ vi.mock("./lifecycle-core.js", () => ({ runServiceUninstall: vi.fn(), })); +const { runDaemonStart, runDaemonRestart, runDaemonStop } = await import("./lifecycle.js"); + describe("runDaemonRestart health checks", () => { - let runDaemonStart: typeof import("./lifecycle.js").runDaemonStart; - let runDaemonRestart: typeof import("./lifecycle.js").runDaemonRestart; - let runDaemonStop: typeof import("./lifecycle.js").runDaemonStop; let envSnapshot: ReturnType; function mockUnmanagedRestart({ @@ -222,10 +222,6 @@ describe("runDaemonRestart health checks", () => { return outcome; } - beforeAll(async () => { - ({ runDaemonStart, runDaemonRestart, runDaemonStop } = await import("./lifecycle.js")); - }); - beforeEach(() => { envSnapshot = captureEnv([ "OPENCLAW_CONTAINER_HINT", diff --git a/src/cli/gateway-run-argv.ts b/src/cli/gateway-run-argv.ts index d651a0f9a3ad..6c5e6343584a 100644 --- a/src/cli/gateway-run-argv.ts +++ b/src/cli/gateway-run-argv.ts @@ -1,4 +1,5 @@ // Fast-path argv parser for `openclaw gateway ...` without full Commander registration. +import { GATEWAY_RUN_BOOLEAN_FLAGS, GATEWAY_RUN_VALUE_FLAGS } from "../../gateway-run-argv.mjs"; import { WINDOWS_TASK_SUPERVISOR_CHILD_FLAG, WINDOWS_TASK_SUPERVISOR_FLAG, @@ -10,49 +11,7 @@ import { isValueToken, } from "../infra/cli-root-options.js"; -const GATEWAY_RUN_VALUE_FLAGS = new Set([ - "--port", - "--bind", - "--token", - "--token-file", - "--auth", - "--password", - "--password-file", - "--tailscale", - "--ws-log", - "--raw-stream-path", -]); - -const GATEWAY_RUN_BOOLEAN_FLAGS = new Set([ - "--tailscale-reset-on-exit", - "--allow-unconfigured", - "--dev", - "--ambient-channels", - "--dev-ambient-channels", - "--reset", - "--update-canary", - "--force", - "--verbose", - "--cli-backend-logs", - "--claude-cli-logs", - "--compact", - "--raw-stream", -]); - -export function isForegroundGatewayRunArgv(argv: string[]): boolean { - const positionals = getCommandPositionalsWithRootOptions(argv, { - commandPath: ["gateway"], - booleanFlags: [...GATEWAY_RUN_BOOLEAN_FLAGS], - valueFlags: [...GATEWAY_RUN_VALUE_FLAGS], - mode: "command-path", - }); - if (!positionals) { - return false; - } - // Foreground gateway owns the terminal/process environment itself; respawning would - // add an extra parent process around the long-lived server. - return positionals.length === 0 || (positionals.length === 1 && positionals[0] === "run"); -} +export { isForegroundGatewayRunArgv } from "../../gateway-run-argv.mjs"; /** Return how many argv tokens a gateway-run option consumes, or 0 when not recognized. */ export function consumeGatewayRunOptionToken(args: ReadonlyArray, index: number): number { diff --git a/src/cli/update-cli/update-command-repair-isolation.test-support.ts b/src/cli/update-cli/update-command-repair-isolation.test-support.ts index c8cd76b9e75b..49b925cfa5f7 100644 --- a/src/cli/update-cli/update-command-repair-isolation.test-support.ts +++ b/src/cli/update-cli/update-command-repair-isolation.test-support.ts @@ -146,6 +146,9 @@ export async function writeRepairCandidate(candidate: string, configChange: bool "node-sqlite.mjs", "node-runtime-update.mjs", "node-runtime-recovery.mjs", + "cli-root-options.mjs", + "gateway-run-argv.mjs", + "gateway-shutdown-budget.mjs", "package.json", ]) { await fs.copyFile(path.join(process.cwd(), file), path.join(candidate, file)); diff --git a/src/commands/doctor-config-preflight.process.test-support.ts b/src/commands/doctor-config-preflight.process.test-support.ts index 97b3f65f4203..1f5d57199da3 100644 --- a/src/commands/doctor-config-preflight.process.test-support.ts +++ b/src/commands/doctor-config-preflight.process.test-support.ts @@ -113,6 +113,9 @@ export function createSourceRuntime(root: string): string { "node-sqlite.mjs", "node-runtime-update.mjs", "node-runtime-recovery.mjs", + "cli-root-options.mjs", + "gateway-run-argv.mjs", + "gateway-shutdown-budget.mjs", "package.json", "tsconfig.json", ]) { diff --git a/src/daemon/launchd-plist.ts b/src/daemon/launchd-plist.ts index 162b29854f03..b69f6ed05669 100644 --- a/src/daemon/launchd-plist.ts +++ b/src/daemon/launchd-plist.ts @@ -2,6 +2,7 @@ import fs from "node:fs/promises"; import { asOptionalRecord, isStringRecord } from "@openclaw/normalization-core/record-coerce"; import { hasErrnoCode } from "../infra/errno.js"; +import { LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS } from "../infra/gateway-shutdown-budget.js"; import { runExec } from "../process/exec.js"; import type { GatewayServiceCommandConfig, @@ -9,10 +10,11 @@ import type { GatewayServiceReadOptions, } from "./service-types.js"; +export { LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS }; + // launchd defaults to a 10s spawn throttle. Keep that default explicitly so // crash loops back off instead of respawning every second while still allowing // explicit kickstart restarts to take effect. -export const LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS = 20; // launchd stores plist integer values in decimal; 0o077 renders as 63 (owner-only files). export const LAUNCH_AGENT_POLICY = { RunAtLoad: true, diff --git a/src/entry.compile-cache.test.ts b/src/entry.compile-cache.test.ts index b9f80ffaa0ed..de96d7bf5c45 100644 --- a/src/entry.compile-cache.test.ts +++ b/src/entry.compile-cache.test.ts @@ -219,6 +219,20 @@ describe("entry compile cache", () => { }, ); + it.each(["linux", "darwin"] as const)( + "keeps the serving Gateway in process with inherited compile cache on %s", + async (platform) => { + await markSourceCheckout(); + argv = [process.execPath, entryFile, "--profile=fixture", "gateway", "run"]; + await withMockedPlatform(platform, async () => { + await expect( + respawnWithoutOpenClawCompileCacheIfNeeded({ currentFile: entryFile, installRoot: root }), + ).resolves.toBe(false); + expect(spawn).not.toHaveBeenCalled(); + }); + }, + ); + it("keeps interactive no-cache respawns attached to the terminal", async () => { await markSourceCheckout(); argv = [process.execPath, entryFile, "tui"]; diff --git a/src/entry.compile-cache.ts b/src/entry.compile-cache.ts index d13fb75f5b8b..d3daad7d30d1 100644 --- a/src/entry.compile-cache.ts +++ b/src/entry.compile-cache.ts @@ -5,6 +5,7 @@ import { enableCompileCache, getCompileCacheDir } from "node:module"; import os from "node:os"; import path from "node:path"; import process from "node:process"; +import { isForegroundGatewayRunArgv } from "./cli/gateway-run-argv.js"; import { isForegroundGmailRunArgv, isTerminalInteractiveRespawnArgv, @@ -117,6 +118,10 @@ function buildOpenClawCompileCacheRespawnPlan(params: { const env = params.env ?? process.env; const argv = process.argv; const platform = process.platform; + // A recovered Unix Gateway must not acquire another short-lived stop wrapper. + if (platform !== "win32" && isForegroundGatewayRunArgv(argv)) { + return undefined; + } if (isForegroundGmailRunArgv(argv) || shouldKeepNativeHookRelayInProcess(argv, platform)) { return undefined; } diff --git a/src/gateway/worker-environments/node-bootstrap-artifact.test-support.ts b/src/gateway/worker-environments/node-bootstrap-artifact.test-support.ts index cc2f44d44937..c174beabc102 100644 --- a/src/gateway/worker-environments/node-bootstrap-artifact.test-support.ts +++ b/src/gateway/worker-environments/node-bootstrap-artifact.test-support.ts @@ -89,6 +89,9 @@ export function useNodeBootstrapArtifactFixtures() { await write(packageRoot, "node-sqlite.mjs", "export const probe = true;"); await write(packageRoot, "node-runtime-update.mjs", "export const update = true;"); await write(packageRoot, "node-runtime-recovery.mjs", "export const recovery = true;"); + await write(packageRoot, "cli-root-options.mjs", "export {};"); + await write(packageRoot, "gateway-run-argv.mjs", "export {};"); + await write(packageRoot, "gateway-shutdown-budget.mjs", "export {};"); await write(packageRoot, "node-host-launcher.mjs", "export const launcher = true;"); await write(packageRoot, "scripts/preinstall.mjs", "export {};\n"); await write( diff --git a/src/gateway/worker-environments/node-bootstrap-artifact.ts b/src/gateway/worker-environments/node-bootstrap-artifact.ts index 083564d88005..ef6ce09d6324 100644 --- a/src/gateway/worker-environments/node-bootstrap-artifact.ts +++ b/src/gateway/worker-environments/node-bootstrap-artifact.ts @@ -51,6 +51,9 @@ const BOOTSTRAP_LAUNCHER_FILES = [ "node-sqlite.mjs", "node-runtime-update.mjs", "node-runtime-recovery.mjs", + "cli-root-options.mjs", + "gateway-run-argv.mjs", + "gateway-shutdown-budget.mjs", "node-host-launcher.mjs", ]; const READ_CONCURRENCY = 16; diff --git a/src/gateway/worker-environments/node-enrollment.test.ts b/src/gateway/worker-environments/node-enrollment.test.ts index 7dabea92c12b..c534748bc3e7 100644 --- a/src/gateway/worker-environments/node-enrollment.test.ts +++ b/src/gateway/worker-environments/node-enrollment.test.ts @@ -130,6 +130,9 @@ describe("worker node enrollment", () => { path.join(packageRoot, "node-runtime-recovery.mjs"), "export const recovery = true;", ), + fs.writeFile(path.join(packageRoot, "cli-root-options.mjs"), "export {};"), + fs.writeFile(path.join(packageRoot, "gateway-run-argv.mjs"), "export {};"), + fs.writeFile(path.join(packageRoot, "gateway-shutdown-budget.mjs"), "export {};"), fs.writeFile(path.join(packageRoot, "dist/entry.js"), "export const ready = true;"), fs.writeFile( path.join(packageRoot, "dist/build-info.json"), diff --git a/src/infra/cli-root-options.ts b/src/infra/cli-root-options.ts index 0509255e82e7..bdc17fad807e 100644 --- a/src/infra/cli-root-options.ts +++ b/src/infra/cli-root-options.ts @@ -1,200 +1,9 @@ -/** CLI token that stops root option scanning and leaves following args positional. */ -export const FLAG_TERMINATOR = "--"; - -const ROOT_BOOLEAN_FLAGS = new Set(["--dev", "--no-color"]); -const ROOT_VALUE_FLAGS = new Set(["--profile", "--log-level", "--container"]); - -/** Returns whether a token can be consumed as a root option value. */ -export function isValueToken(arg: string | undefined): boolean { - if (!arg || arg === FLAG_TERMINATOR) { - return false; - } - if (!arg.startsWith("-")) { - return true; - } - return /^-\d+(?:\.\d+)?$/.test(arg); -} - -/** Count root-option tokens conservatively for route matching. */ -export function consumeRootOptionToken(args: ReadonlyArray, index: number): number { - const arg = args[index]; - if (!arg) { - return 0; - } - if (ROOT_BOOLEAN_FLAGS.has(arg)) { - return 1; - } - if ( - arg.startsWith("--profile=") || - arg.startsWith("--log-level=") || - arg.startsWith("--container=") - ) { - return 1; - } - if (ROOT_VALUE_FLAGS.has(arg)) { - return isValueToken(args[index + 1]) ? 2 : 1; - } - return 0; -} - -/** Consume required root values by their Commander role before startup policy is selected. */ -export function consumeRootCommandOptionToken(args: readonly string[], index: number): number { - return consumeKnownOptionToken(args, index, ROOT_BOOLEAN_FLAGS, ROOT_VALUE_FLAGS, "command-path"); -} - -/** Read positional command tokens while accepting root options at any pre-terminator position. */ -export function getRootOptionAwareCommandPath(argv: readonly string[], depth: number): string[] { - const args = argv.slice(2); - const path: string[] = []; - let literal = false; - for (let index = 0; index < args.length; index += 1) { - const arg = args[index]; - if (!arg) { - break; - } - if (!literal && arg === FLAG_TERMINATOR) { - // A leading terminator still leaves a command to discover; later operands belong to callers. - if (path.length > 0) { - break; - } - literal = true; - continue; - } - const consumed = literal ? 0 : consumeRootCommandOptionToken(args, index); - if (consumed > 0) { - index += consumed - 1; - continue; - } - if (!literal && arg.startsWith("-")) { - continue; - } - path.push(arg); - if (path.length >= depth) { - break; - } - } - return path; -} - -type CommandPositionalsParseOptions = { - commandPath: ReadonlyArray; - booleanFlags?: ReadonlyArray; - valueFlags?: ReadonlyArray; - maxPositionals?: number; - mode?: "route" | "command-path"; -}; - -function consumeKnownOptionToken( - args: ReadonlyArray, - index: number, - booleanFlags: ReadonlySet, - valueFlags: ReadonlySet, - mode: CommandPositionalsParseOptions["mode"], -): number { - const arg = args[index]; - if (!arg || arg === FLAG_TERMINATOR || !arg.startsWith("-")) { - return 0; - } - - const equalsIndex = arg.indexOf("="); - const flag = equalsIndex === -1 ? arg : arg.slice(0, equalsIndex); - if (booleanFlags.has(flag)) { - return equalsIndex === -1 ? 1 : 0; - } - if (!valueFlags.has(flag)) { - return 0; - } - if (equalsIndex !== -1) { - return mode === "command-path" || arg.slice(equalsIndex + 1).trim() ? 1 : 0; - } - // Required Commander values include empty strings, flag-looking tokens, and `--`. - // Discovery must consume them before choosing startup policy; routes still validate values. - if (mode === "command-path") { - return args[index + 1] !== undefined ? 2 : 0; - } - return isValueToken(args[index + 1]) ? 2 : 0; -} - -/** Parse command positionals while consuming known root and command options. */ -export function getCommandPositionalsWithRootOptions( - argv: readonly string[], - options: CommandPositionalsParseOptions, -): string[] | null { - return parseCommandArgsWithRootOptions(argv, options, false); -} - -/** Preserve the leaf's raw arguments after consuming its root and parent options. */ -export function getCommandArgsWithRootOptions( - argv: readonly string[], - options: Omit, -): string[] | null { - return parseCommandArgsWithRootOptions(argv, options, true); -} - -function parseCommandArgsWithRootOptions( - argv: readonly string[], - options: CommandPositionalsParseOptions, - returnTail: boolean, -): string[] | null { - const args = argv.slice(2); - const booleanFlags = new Set(options.booleanFlags ?? []); - const valueFlags = new Set(options.valueFlags ?? []); - const positionals: string[] = []; - let commandIndex = 0; - let literal = false; - - for (let index = 0; index < args.length; index += 1) { - const arg = args[index]; - if (arg === undefined) { - break; - } - if (!literal && arg === FLAG_TERMINATOR) { - if (options.mode !== "command-path") { - break; - } - literal = true; - continue; - } - const rootConsumed = literal - ? 0 - : options.mode === "command-path" - ? consumeRootCommandOptionToken(args, index) - : consumeRootOptionToken(args, index); - if (rootConsumed > 0) { - index += rootConsumed - 1; - continue; - } - if (!literal && arg.startsWith("-")) { - const optionConsumed = consumeKnownOptionToken( - args, - index, - booleanFlags, - valueFlags, - options.mode, - ); - if (optionConsumed === 0 || commandIndex === 0) { - return null; - } - index += optionConsumed - 1; - continue; - } - if (commandIndex < options.commandPath.length) { - if (arg !== options.commandPath[commandIndex]) { - return null; - } - commandIndex += 1; - if (returnTail && commandIndex === options.commandPath.length) { - const tail = args.slice(index + 1); - // A downstream parser must not reactivate flags after an earlier literal boundary. - return literal ? [FLAG_TERMINATOR, ...tail] : tail; - } - continue; - } - positionals.push(arg); - if (positionals.length === options.maxPositionals) { - return positionals; - } - } - - return commandIndex < options.commandPath.length ? null : positionals; -} +export { + FLAG_TERMINATOR, + consumeRootOptionToken, + consumeRootCommandOptionToken, + getRootOptionAwareCommandPath, + getCommandPositionalsWithRootOptions, + getCommandArgsWithRootOptions, + isValueToken, +} from "../../cli-root-options.mjs"; diff --git a/src/infra/gateway-owner-lease.ts b/src/infra/gateway-owner-lease.ts index 1f6a3b2f14d8..bb759be3fb9d 100644 --- a/src/infra/gateway-owner-lease.ts +++ b/src/infra/gateway-owner-lease.ts @@ -1,5 +1,6 @@ import { randomUUID } from "node:crypto"; import { hostname } from "node:os"; +import type { DatabaseSync } from "node:sqlite"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { createSubsystemLogger } from "../logging/subsystem.js"; import { getFileLockProcessStartTime } from "../shared/pid-alive.js"; @@ -68,6 +69,55 @@ function parseSupervisor(value: unknown): GatewayOwnerSupervisor | null { return { kind: value.kind, name: value.name }; } +/** Read through the caller's admitted connection when ownership guards a write. */ +export function readGatewayOwnerLeaseFromDatabase( + db: DatabaseSync, + port?: number, +): GatewayOwnerLeaseIdentity | undefined { + if (!tableExists(db, "state_leases")) { + return undefined; + } + const row = readOpenClawStateLease(db, gatewayOwnerKey); + if (!row) { + return undefined; + } + const processOwner = parseStateLeaseProcessOwner(row.payloadJson); + let payload: unknown; + try { + payload = row.payloadJson ? JSON.parse(row.payloadJson) : null; + } catch { + payload = null; + } + if ( + !processOwner || + !isRecord(payload) || + typeof payload.port !== "number" || + !Number.isInteger(payload.port) || + payload.port <= 0 || + payload.port > 65535 || + (payload.mode !== "foreground" && payload.mode !== "supervised") + ) { + throw new Error("Gateway owner lease identity could not be verified"); + } + if (port !== undefined && payload.port !== port) { + return undefined; + } + const supervisor = parseSupervisor(payload.supervisor); + if ((payload.mode === "foreground") !== (supervisor === null)) { + throw new Error("Gateway owner lease supervisor does not match its listener mode"); + } + return { + ...processOwner, + owner: row.owner, + port: payload.port, + mode: payload.mode, + supervisor, + // Expiry cannot revoke the separate physical Gateway coordinator. + state: readStateLeaseProcessOwnerStatus(processOwner), + expired: row.expiresAt === null || row.expiresAt <= Date.now(), + }; +} + export function readGatewayOwnerLease( params: { env?: NodeJS.ProcessEnv; @@ -77,54 +127,8 @@ export function readGatewayOwnerLease( openStateSchemaReadAdmission?: OpenClawStateSchemaReadAdmission; } = {}, ): GatewayOwnerLeaseIdentity | undefined { - const operation = ({ - db, - }: { - db: import("node:sqlite").DatabaseSync; - }): GatewayOwnerLeaseIdentity | undefined => { - if (!tableExists(db, "state_leases")) { - return undefined; - } - const row = readOpenClawStateLease(db, gatewayOwnerKey); - if (!row) { - return undefined; - } - const processOwner = parseStateLeaseProcessOwner(row.payloadJson); - let payload: unknown; - try { - payload = row.payloadJson ? JSON.parse(row.payloadJson) : null; - } catch { - payload = null; - } - if ( - !processOwner || - !isRecord(payload) || - typeof payload.port !== "number" || - !Number.isInteger(payload.port) || - payload.port <= 0 || - payload.port > 65535 || - (payload.mode !== "foreground" && payload.mode !== "supervised") - ) { - throw new Error("Gateway owner lease identity could not be verified"); - } - if (params.port !== undefined && payload.port !== params.port) { - return undefined; - } - const supervisor = parseSupervisor(payload.supervisor); - if ((payload.mode === "foreground") !== (supervisor === null)) { - throw new Error("Gateway owner lease supervisor does not match its listener mode"); - } - return { - ...processOwner, - owner: row.owner, - port: payload.port, - mode: payload.mode, - supervisor, - // Expiry cannot revoke the separate physical Gateway coordinator. - state: readStateLeaseProcessOwnerStatus(processOwner), - expired: row.expiresAt === null || row.expiresAt <= Date.now(), - }; - }; + const operation = ({ db }: { db: DatabaseSync }) => + readGatewayOwnerLeaseFromDatabase(db, params.port); return params.current || params.openStateSchemaReadAdmission ? withExistingOpenClawStateDatabaseCurrentReadOnly( operation, diff --git a/src/infra/gateway-shutdown-budget.ts b/src/infra/gateway-shutdown-budget.ts index 6d194a38bbc9..0cde0821380a 100644 --- a/src/infra/gateway-shutdown-budget.ts +++ b/src/infra/gateway-shutdown-budget.ts @@ -1,8 +1,7 @@ -// Shared Gateway stop policy used by the run loop and native service definitions. -const GATEWAY_SHUTDOWN_DRAIN_TIMEOUT_MS = 315_000; -export const GATEWAY_SHUTDOWN_RESERVE_MS = 10_000; -export const GATEWAY_SUPERVISOR_EXIT_MARGIN_MS = 5_000; -export const GATEWAY_SHUTDOWN_TIMEOUT_MS = - GATEWAY_SHUTDOWN_DRAIN_TIMEOUT_MS + GATEWAY_SHUTDOWN_RESERVE_MS; -export const GATEWAY_SERVICE_STOP_TIMEOUT_MS = - GATEWAY_SHUTDOWN_TIMEOUT_MS + GATEWAY_SUPERVISOR_EXIT_MARGIN_MS; +export { + GATEWAY_SHUTDOWN_RESERVE_MS, + GATEWAY_SUPERVISOR_EXIT_MARGIN_MS, + GATEWAY_SHUTDOWN_TIMEOUT_MS, + GATEWAY_SERVICE_STOP_TIMEOUT_MS, + LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS, +} from "../../gateway-shutdown-budget.mjs"; diff --git a/src/infra/node-runtime-recovery.gateway-drain.test.ts b/src/infra/node-runtime-recovery.gateway-drain.test.ts new file mode 100644 index 000000000000..abc8c3882b59 --- /dev/null +++ b/src/infra/node-runtime-recovery.gateway-drain.test.ts @@ -0,0 +1,93 @@ +import { spawn } from "node:child_process"; +import { once } from "node:events"; +import fs from "node:fs/promises"; +import path from "node:path"; +import { afterEach, describe, expect, it } from "vitest"; +import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; +import { resolveTestNodeExecPath } from "../test-utils/node-process.js"; + +describe.skipIf(process.platform === "win32")("recovered Gateway shutdown", () => { + const tempDirs = useAutoCleanupTempDirTracker(afterEach); + + it("keeps admitted work alive until the serving child drains", async () => { + const root = tempDirs.make("openclaw-recovery-drain-"); + const worker = path.join(root, "worker.mjs"); + const wrapper = path.join(root, "wrapper.mjs"); + await fs.writeFile( + worker, + `import { createServer } from "node:http"; +let pending; +let stopping = false; +const server = createServer((request, response) => { + if (stopping) { response.writeHead(503); response.end("draining"); return; } + pending = response; + response.writeHead(200); + response.write("admitted\\n"); +}); +process.on("SIGTERM", () => { + if (stopping) return; + stopping = true; + process.stdout.write("draining\\n"); + setTimeout(() => { + pending.end("completed\\n"); + server.close(() => process.exit(0)); + }, 3_000); +}); +server.listen(0, "127.0.0.1", () => process.stdout.write(JSON.stringify({ + port: server.address().port, pid: process.pid, parent: process.ppid, +}) + "\\n")); +`, + ); + await fs.writeFile( + wrapper, + `import { runRespawnedChild } from ${JSON.stringify(new URL("../../node-runtime-recovery.mjs", import.meta.url).href)}; +runRespawnedChild(process.execPath, [${JSON.stringify(worker)}], process.env); +`, + ); + const child = spawn( + resolveTestNodeExecPath(), + [wrapper, "--profile=fixture", "gateway", "run", "--bind", "loopback"], + { env: { PATH: process.env.PATH, HOME: root }, stdio: ["ignore", "pipe", "pipe"] }, + ); + let stdout = ""; + let stderr = ""; + child.stdout.on("data", (data: Buffer) => { + stdout += data.toString(); + }); + child.stderr.on("data", (data: Buffer) => { + stderr += data.toString(); + }); + const exited = once(child, "exit"); + let servingPid: number | undefined; + try { + await expect.poll(() => stdout, { timeout: 10_000 }).toContain("\n"); + const ready = JSON.parse(stdout.split("\n")[0]!) as { + port: number; + pid: number; + parent: number; + }; + servingPid = ready.pid; + expect(ready.parent).toBe(child.pid); + expect(servingPid).not.toBe(child.pid); + const response = await fetch(`http://127.0.0.1:${ready.port}/work`); + const body = response.text().catch(() => "connection terminated before completion"); + child.kill("SIGTERM"); + await expect.poll(() => stdout).toContain("draining\n"); + const denied = await fetch(`http://127.0.0.1:${ready.port}/new-work`); + expect(denied.status).toBe(503); + await denied.text(); + expect(await body).toBe("admitted\ncompleted\n"); + expect(await exited, stderr).toEqual([0, null]); + } finally { + if (child.exitCode === null && child.signalCode === null) { + if (servingPid) { + try { + process.kill(servingPid, "SIGKILL"); + } catch {} + } + child.kill("SIGKILL"); + } + await exited; + } + }, 15_000); +}); diff --git a/src/infra/node-runtime-recovery.shutdown.test.ts b/src/infra/node-runtime-recovery.shutdown.test.ts new file mode 100644 index 000000000000..7bb2b8b6cf68 --- /dev/null +++ b/src/infra/node-runtime-recovery.shutdown.test.ts @@ -0,0 +1,85 @@ +import { ChildProcess, type SpawnOptions } from "node:child_process"; +import { afterEach, beforeEach, expect, it, vi, type MockInstance } from "vitest"; +import { runRespawnedChild } from "../../node-runtime-recovery.mjs"; + +const spawn = vi.hoisted(() => + vi.fn<(command: string, args: string[], options: SpawnOptions) => ChildProcess>(), +); +vi.mock("node:child_process", async (importOriginal) => ({ + ...(await importOriginal()), + spawn, +})); + +let child: ChildProcess; +let kill: MockInstance; +let exit: MockInstance; +let detach: (() => void) | undefined; +beforeEach(() => { + vi.useFakeTimers(); + child = new ChildProcess(); + kill = vi.spyOn(child, "kill").mockReturnValue(true); + spawn.mockReturnValue(child); + exit = vi.spyOn(process, "exit").mockImplementation(vi.fn()); + vi.spyOn(process, "kill").mockReturnValue(true); +}); +afterEach(() => { + detach?.(); + detach = undefined; + vi.useRealTimers(); + vi.restoreAllMocks(); +}); + +it.each([ + { platform: "linux", args: ["gateway", "run"], nativeBudgetMs: 330_000 }, + { platform: "darwin", args: ["gateway"], nativeBudgetMs: 20_000 }, + { platform: "linux", args: ["gateway", "status"], nativeBudgetMs: 3_000 }, + { platform: "win32", args: ["gateway", "run"], nativeBudgetMs: 3_000 }, +] as const)( + "bounds $platform $args shutdown without preempting the serving owner", + ({ platform, args, nativeBudgetMs }) => { + vi.spyOn(process, "platform", "get").mockReturnValue(platform); + vi.spyOn(process, "argv", "get").mockReturnValue([ + "node", + "openclaw.mjs", + "--profile=fixture", + ...args, + ]); + const previous = new Set(process.listeners("SIGTERM")); + runRespawnedChild("node", ["child.mjs"], { + OPENCLAW_LAUNCHD_LABEL: "ai.openclaw.fixture", + XPC_SERVICE_NAME: "ai.openclaw.fixture", + }); + detach = () => child.emit("exit", 0, null); + const signal = process.listeners("SIGTERM").find((listener) => !previous.has(listener)); + expect(signal).toBeDefined(); + signal!("SIGTERM"); + expect(kill).toHaveBeenCalledExactlyOnceWith("SIGTERM"); + // Reserve the final two seconds for escalation; all earlier time belongs to the child. + vi.advanceTimersByTime(nativeBudgetMs - 2_001); + expect(kill).toHaveBeenCalledTimes(1); + signal!("SIGTERM"); + expect(kill).toHaveBeenCalledTimes(2); + vi.advanceTimersByTime(1); + expect(kill).toHaveBeenCalledTimes(3); + vi.advanceTimersByTime(1_000); + expect(kill).toHaveBeenLastCalledWith(platform === "win32" ? "SIGTERM" : "SIGKILL"); + expect(exit).not.toHaveBeenCalled(); + vi.advanceTimersByTime(1_000); + expect(exit).toHaveBeenCalledExactlyOnceWith(1); + }, +); + +it("removes the shutdown deadline when the child exits cooperatively", () => { + vi.spyOn(process, "platform", "get").mockReturnValue("linux"); + vi.spyOn(process, "argv", "get").mockReturnValue(["node", "openclaw.mjs", "gateway"]); + const previous = new Set(process.listeners("SIGTERM")); + runRespawnedChild("node", ["child.mjs"], {}); + detach = () => child.emit("exit", 0, null); + process.listeners("SIGTERM").find((listener) => !previous.has(listener))!("SIGTERM"); + vi.advanceTimersByTime(3_000); + child.emit("exit", 0, null); + vi.advanceTimersByTime(330_000); + expect(kill).toHaveBeenCalledExactlyOnceWith("SIGTERM"); + expect(exit).toHaveBeenCalledExactlyOnceWith(0); + expect(process.listeners("SIGTERM")).toEqual([...previous]); +}); diff --git a/src/infra/restart-intent.ts b/src/infra/restart-intent.ts index f28d98041e8c..e07e2305b06b 100644 --- a/src/infra/restart-intent.ts +++ b/src/infra/restart-intent.ts @@ -1,13 +1,16 @@ import { existsSync } from "node:fs"; +import type { DatabaseSync } from "node:sqlite"; // Persists short-lived gateway restart intent for supervisor SIGTERM handoff. import { asPositiveSafeInteger } from "@openclaw/normalization-core/number-coercion"; import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice"; +import { resolveSystemdServiceName } from "../daemon/systemd-service-files.js"; import { createSubsystemLogger } from "../logging/subsystem.js"; import { runExistingOpenClawStateWriteTransaction } from "../state/openclaw-state-db-existing-write.js"; import type { DB as OpenClawStateKyselyDatabase } from "../state/openclaw-state-db.generated.js"; import { openOpenClawStateDatabase } from "../state/openclaw-state-db.js"; import { resolveOpenClawStateSqlitePath } from "../state/openclaw-state-db.paths.js"; import { OPENCLAW_STATE_SCHEMA_SQL } from "../state/openclaw-state-schema.js"; +import { readGatewayOwnerLeaseFromDatabase } from "./gateway-owner-lease.js"; import { executeSqliteQuerySync, executeSqliteQueryTakeFirstSync, @@ -65,6 +68,67 @@ export function writeGatewayRestartIntentSync(opts: { if (targetPid === null) { return false; } + return writeGatewayRestartIntentForTargetSync(opts, () => targetPid); +} + +export type GatewayRestartIntentService = { + kind: "systemd" | "launchd"; + name: string; +}; + +/** Native service control keeps its selected service; resolve its serving process at admission. */ +export function writeGatewayServiceRestartIntentSync(opts: { + env?: NodeJS.ProcessEnv; + targetPid?: number; + service?: GatewayRestartIntentService; + intent?: GatewayRestartIntent; + reason?: string; + assertCurrent: () => void; + warn: (message: string) => void; +}): boolean { + let ownershipUnverified = false; + const written = writeGatewayRestartIntentForTargetSync( + opts, + (db) => { + if (!opts.service) { + return opts.targetPid; + } + try { + const owner = readGatewayOwnerLeaseFromDatabase(db); + ownershipUnverified = owner?.state === "unknown"; + const supervisor = owner?.supervisor; + if ( + owner?.state === "live" && + owner.mode === "supervised" && + supervisor?.kind === opts.service.kind && + supervisor.name !== null && + (supervisor.kind === "systemd" + ? resolveSystemdServiceName({ OPENCLAW_SYSTEMD_UNIT: supervisor.name }) === + resolveSystemdServiceName({ OPENCLAW_SYSTEMD_UNIT: opts.service.name }) + : supervisor.name === opts.service.name) + ) { + return owner.pid; + } + } catch { + ownershipUnverified = true; + } + return opts.targetPid; + }, + opts.assertCurrent, + ); + if (ownershipUnverified) { + opts.warn( + "Could not verify the serving Gateway owner; using native service status for restart intent.", + ); + } + return written; +} + +function writeGatewayRestartIntentForTargetSync( + opts: { env?: NodeJS.ProcessEnv; intent?: GatewayRestartIntent; reason?: string }, + resolveTargetPid: (db: DatabaseSync) => number | undefined, + assertCurrent?: () => void, +): boolean { const env = opts.env ?? process.env; try { if (!existsSync(resolveOpenClawStateSqlitePath(env))) { @@ -78,10 +142,17 @@ export function writeGatewayRestartIntentSync(opts: { opts.intent.waitMs >= 0 ? Math.floor(opts.intent.waitMs) : null; - const createdAt = Date.now(); // The old Gateway still owns the schema until the restart hands off. - runExistingOpenClawStateWriteTransaction( + return runExistingOpenClawStateWriteTransaction( ({ db }) => { + // Coordinator/BEGIN admission can block while the supervised owner changes. + assertCurrent?.(); + const targetPid = asPositiveSafeInteger(resolveTargetPid(db)) ?? null; + assertCurrent?.(); + if (targetPid === null) { + return false; + } + const createdAt = Date.now(); const stateDb = getNodeSqliteKysely(db); executeSqliteQuerySync( db, @@ -109,12 +180,14 @@ export function writeGatewayRestartIntentSync(opts: { }), ), ); + return true; }, { env }, { schemaSql: schema, operationLabel: "gateway.restart-intent.write" }, ); - return true; } catch (err) { + // Revoked native control authority must not become a best-effort storage warning. + assertCurrent?.(); restartLog.warn(`failed to write gateway restart intent: ${String(err)}`); return false; } diff --git a/test/node-host-launcher.test.ts b/test/node-host-launcher.test.ts index 11464dd5176c..90607f67e23e 100644 --- a/test/node-host-launcher.test.ts +++ b/test/node-host-launcher.test.ts @@ -66,6 +66,9 @@ async function writePackage(root: string, version: string, body: string, schema "node-sqlite.mjs", "node-runtime-update.mjs", "node-runtime-recovery.mjs", + "cli-root-options.mjs", + "gateway-run-argv.mjs", + "gateway-shutdown-budget.mjs", ]) { await fs.copyFile(path.resolve(name), path.join(root, name)); } diff --git a/test/openclaw-launcher-version.e2e.test.ts b/test/openclaw-launcher-version.e2e.test.ts index 05fa688d2010..be3624f4c678 100644 --- a/test/openclaw-launcher-version.e2e.test.ts +++ b/test/openclaw-launcher-version.e2e.test.ts @@ -46,6 +46,18 @@ async function makeLauncherVersionFixture( path.resolve(process.cwd(), "node-runtime-recovery.mjs"), path.join(fixtureRoot, "node-runtime-recovery.mjs"), ); + await fs.copyFile( + path.resolve(process.cwd(), "cli-root-options.mjs"), + path.join(fixtureRoot, "cli-root-options.mjs"), + ); + await fs.copyFile( + path.resolve(process.cwd(), "gateway-run-argv.mjs"), + path.join(fixtureRoot, "gateway-run-argv.mjs"), + ); + await fs.copyFile( + path.resolve(process.cwd(), "gateway-shutdown-budget.mjs"), + path.join(fixtureRoot, "gateway-shutdown-budget.mjs"), + ); await fs.mkdir(path.join(fixtureRoot, "dist"), { recursive: true }); await fs.writeFile( path.join(fixtureRoot, "package.json"), diff --git a/test/openclaw-launcher.e2e.test.ts b/test/openclaw-launcher.e2e.test.ts index 5f5ab5163112..8fd0e5c3f7e5 100644 --- a/test/openclaw-launcher.e2e.test.ts +++ b/test/openclaw-launcher.e2e.test.ts @@ -36,6 +36,18 @@ async function makeLauncherFixture(fixtureRoots: string[]): Promise { path.resolve(process.cwd(), "node-runtime-recovery.mjs"), path.join(fixtureRoot, "node-runtime-recovery.mjs"), ); + await fs.copyFile( + path.resolve(process.cwd(), "cli-root-options.mjs"), + path.join(fixtureRoot, "cli-root-options.mjs"), + ); + await fs.copyFile( + path.resolve(process.cwd(), "gateway-run-argv.mjs"), + path.join(fixtureRoot, "gateway-run-argv.mjs"), + ); + await fs.copyFile( + path.resolve(process.cwd(), "gateway-shutdown-budget.mjs"), + path.join(fixtureRoot, "gateway-shutdown-budget.mjs"), + ); await fs.copyFile( path.resolve(process.cwd(), "node-sqlite.mjs"), path.join(fixtureRoot, "node-sqlite.mjs"), diff --git a/test/scripts/run-node-lifecycle.test.ts b/test/scripts/run-node-lifecycle.test.ts index a38ac4632140..5d57c2e6075f 100644 --- a/test/scripts/run-node-lifecycle.test.ts +++ b/test/scripts/run-node-lifecycle.test.ts @@ -34,6 +34,9 @@ it.runIf(process.platform !== "win32")( "node-version.mjs", "node-runtime-update.mjs", "node-runtime-recovery.mjs", + "cli-root-options.mjs", + "gateway-run-argv.mjs", + "gateway-shutdown-budget.mjs", "node-sqlite.mjs", ]) { copyFileSync(filename, path.join(checkoutRoot, filename));