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));