mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
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.
This commit is contained in:
parent
388d856006
commit
2ce0573b02
38 changed files with 1108 additions and 400 deletions
|
|
@ -76,6 +76,7 @@ COPY node-version.mjs ./
|
||||||
COPY node-sqlite.mjs ./
|
COPY node-sqlite.mjs ./
|
||||||
COPY node-runtime-update.mjs ./
|
COPY node-runtime-update.mjs ./
|
||||||
COPY node-runtime-recovery.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 node-host-launcher.mjs ./
|
||||||
COPY openclaw.mjs ./
|
COPY openclaw.mjs ./
|
||||||
COPY ui/package.json ./ui/package.json
|
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-sqlite.mjs .
|
||||||
COPY --from=runtime-assets --chown=node:node /app/node-runtime-update.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/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/node-host-launcher.mjs .
|
||||||
COPY --from=runtime-assets --chown=node:node /app/openclaw.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}
|
COPY --from=runtime-assets --chown=node:node /app/${OPENCLAW_BUNDLED_PLUGIN_DIR} ./${OPENCLAW_BUNDLED_PLUGIN_DIR}
|
||||||
|
|
|
||||||
23
cli-root-options.d.mts
Normal file
23
cli-root-options.d.mts
Normal file
|
|
@ -0,0 +1,23 @@
|
||||||
|
export const FLAG_TERMINATOR: "--";
|
||||||
|
|
||||||
|
export function isValueToken(arg: string | undefined): boolean;
|
||||||
|
export function consumeRootOptionToken(args: ReadonlyArray<string>, 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<string>;
|
||||||
|
booleanFlags?: ReadonlyArray<string>;
|
||||||
|
valueFlags?: ReadonlyArray<string>;
|
||||||
|
maxPositionals?: number;
|
||||||
|
mode?: "route" | "command-path";
|
||||||
|
};
|
||||||
|
|
||||||
|
export function getCommandPositionalsWithRootOptions(
|
||||||
|
argv: readonly string[],
|
||||||
|
options: CommandPositionalsParseOptions,
|
||||||
|
): string[] | null;
|
||||||
|
export function getCommandArgsWithRootOptions(
|
||||||
|
argv: readonly string[],
|
||||||
|
options: Omit<CommandPositionalsParseOptions, "maxPositionals">,
|
||||||
|
): string[] | null;
|
||||||
176
cli-root-options.mjs
Normal file
176
cli-root-options.mjs
Normal file
|
|
@ -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;
|
||||||
|
}
|
||||||
|
|
@ -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
|
background tasks to finish, up to a drain budget (5 minutes by default). Most
|
||||||
restarts therefore interrupt nothing at all.
|
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
|
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
|
reaches only the Gateway. Systemd still kills remaining child processes when the
|
||||||
Gateway exits or its stop deadline expires. Older `KillMode=control-group` units
|
Gateway exits or its stop deadline expires. Older `KillMode=control-group` units
|
||||||
|
|
|
||||||
3
gateway-run-argv.d.mts
Normal file
3
gateway-run-argv.d.mts
Normal file
|
|
@ -0,0 +1,3 @@
|
||||||
|
export const GATEWAY_RUN_VALUE_FLAGS: ReadonlySet<string>;
|
||||||
|
export const GATEWAY_RUN_BOOLEAN_FLAGS: ReadonlySet<string>;
|
||||||
|
export function isForegroundGatewayRunArgv(argv: string[]): boolean;
|
||||||
46
gateway-run-argv.mjs
Normal file
46
gateway-run-argv.mjs
Normal file
|
|
@ -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");
|
||||||
|
}
|
||||||
5
gateway-shutdown-budget.d.mts
Normal file
5
gateway-shutdown-budget.d.mts
Normal file
|
|
@ -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;
|
||||||
10
gateway-shutdown-budget.mjs
Normal file
10
gateway-shutdown-budget.mjs
Normal file
|
|
@ -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;
|
||||||
|
|
@ -3,47 +3,21 @@ import { spawn, spawnSync } from "node:child_process";
|
||||||
import { lstatSync, readFileSync, readdirSync, realpathSync, statSync } from "node:fs";
|
import { lstatSync, readFileSync, readdirSync, realpathSync, statSync } from "node:fs";
|
||||||
import os from "node:os";
|
import os from "node:os";
|
||||||
import path from "node:path";
|
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 {
|
import {
|
||||||
detectCurrentSqliteCapabilities,
|
detectCurrentSqliteCapabilities,
|
||||||
nodeRuntimeFailure,
|
nodeRuntimeFailure,
|
||||||
SQLITE_CAPABILITY_PROBE,
|
SQLITE_CAPABILITY_PROBE,
|
||||||
} from "./node-sqlite.mjs";
|
} from "./node-sqlite.mjs";
|
||||||
|
|
||||||
const LAUNCHER_ROOT_BOOLEAN_FLAGS = new Set(["--dev", "--no-color"]);
|
export { consumeLauncherRootOptionToken };
|
||||||
const LAUNCHER_ROOT_VALUE_FLAGS = new Set(["--profile", "--log-level", "--container"]);
|
|
||||||
export const isNativeHookRelayInvocation = (argv) => argv[2] === "hooks" && argv[3] === "relay";
|
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.
|
// Mirror the entry's foreground Gmail policy: a wrapper would kill that run before descendant cleanup finishes.
|
||||||
export const isForegroundGmailRunInvocation = (argv) => {
|
export const isForegroundGmailRunInvocation = (argv) => {
|
||||||
const args = argv.slice(2);
|
const args = argv.slice(2);
|
||||||
|
|
@ -70,6 +44,17 @@ const respawnSignalForceKillGraceMs = 1_000;
|
||||||
const respawnSignalHardExitGraceMs = 1_000;
|
const respawnSignalHardExitGraceMs = 1_000;
|
||||||
|
|
||||||
export const runRespawnedChild = (command, args, env) => {
|
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 stdioIsTerminal = process.stdin.isTTY || process.stdout.isTTY;
|
||||||
const child = spawn(command, args, {
|
const child = spawn(command, args, {
|
||||||
stdio: "inherit",
|
stdio: "inherit",
|
||||||
|
|
@ -131,7 +116,7 @@ export const runRespawnedChild = (command, args, env) => {
|
||||||
}
|
}
|
||||||
signalExitTimer = setTimeout(() => {
|
signalExitTimer = setTimeout(() => {
|
||||||
requestChildTermination();
|
requestChildTermination();
|
||||||
}, respawnSignalExitGraceMs);
|
}, signalExitGraceMs);
|
||||||
signalExitTimer.unref?.();
|
signalExitTimer.unref?.();
|
||||||
};
|
};
|
||||||
for (const signal of respawnSignals) {
|
for (const signal of respawnSignals) {
|
||||||
|
|
|
||||||
|
|
@ -37,6 +37,9 @@
|
||||||
"files": [
|
"files": [
|
||||||
"CHANGELOG.md",
|
"CHANGELOG.md",
|
||||||
"LICENSE",
|
"LICENSE",
|
||||||
|
"cli-root-options.mjs",
|
||||||
|
"gateway-run-argv.mjs",
|
||||||
|
"gateway-shutdown-budget.mjs",
|
||||||
"node-sqlite.mjs",
|
"node-sqlite.mjs",
|
||||||
"node-version.mjs",
|
"node-version.mjs",
|
||||||
"node-runtime-update.mjs",
|
"node-runtime-update.mjs",
|
||||||
|
|
|
||||||
|
|
@ -22,6 +22,9 @@ const targets = [
|
||||||
"test",
|
"test",
|
||||||
"skills",
|
"skills",
|
||||||
"config",
|
"config",
|
||||||
|
"cli-root-options.mjs",
|
||||||
|
"gateway-run-argv.mjs",
|
||||||
|
"gateway-shutdown-budget.mjs",
|
||||||
"node-host-launcher.mjs",
|
"node-host-launcher.mjs",
|
||||||
"node-runtime-update.mjs",
|
"node-runtime-update.mjs",
|
||||||
"node-runtime-recovery.mjs",
|
"node-runtime-recovery.mjs",
|
||||||
|
|
|
||||||
|
|
@ -19,6 +19,7 @@ COPY node-version.mjs ./
|
||||||
COPY node-sqlite.mjs ./
|
COPY node-sqlite.mjs ./
|
||||||
COPY node-runtime-update.mjs ./
|
COPY node-runtime-update.mjs ./
|
||||||
COPY node-runtime-recovery.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 node-host-launcher.mjs ./
|
||||||
COPY ui/package.json ./ui/package.json
|
COPY ui/package.json ./ui/package.json
|
||||||
COPY packages ./packages
|
COPY packages ./packages
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,5 @@
|
||||||
// Daemon lifecycle core tests cover service lifecycle transitions and platform adapters.
|
// 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 { OpenClawConfig } from "../../config/config.js";
|
||||||
import type { GatewayServiceControlArgs } from "../../daemon/service-types.js";
|
import type { GatewayServiceControlArgs } from "../../daemon/service-types.js";
|
||||||
import type { GatewayService } from "../../daemon/service.js";
|
import type { GatewayService } from "../../daemon/service.js";
|
||||||
|
|
@ -55,6 +55,7 @@ vi.mock("../../runtime.js", () => ({
|
||||||
vi.mock("../../infra/restart-intent.js", () => ({
|
vi.mock("../../infra/restart-intent.js", () => ({
|
||||||
clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(),
|
clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(),
|
||||||
writeGatewayRestartIntentSync: (opts: unknown) => writeGatewayRestartIntentSync(opts),
|
writeGatewayRestartIntentSync: (opts: unknown) => writeGatewayRestartIntentSync(opts),
|
||||||
|
writeGatewayServiceRestartIntentSync: (opts: unknown) => writeGatewayRestartIntentSync(opts),
|
||||||
}));
|
}));
|
||||||
|
|
||||||
vi.mock("./lifecycle-audit.js", () => ({
|
vi.mock("./lifecycle-audit.js", () => ({
|
||||||
|
|
@ -79,10 +80,8 @@ vi.mock("./lifecycle-audit.js", () => ({
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
let runServiceRestart: typeof import("./lifecycle-core.js").runServiceRestart;
|
const { runServiceRestart, runServiceStart, runServiceStop, runServiceUninstall } =
|
||||||
let runServiceStart: typeof import("./lifecycle-core.js").runServiceStart;
|
await import("./lifecycle-core.js");
|
||||||
let runServiceStop: typeof import("./lifecycle-core.js").runServiceStop;
|
|
||||||
let runServiceUninstall: typeof import("./lifecycle-core.js").runServiceUninstall;
|
|
||||||
|
|
||||||
// oxlint-disable-next-line typescript/no-unnecessary-type-parameters -- Test helper lets assertions ascribe logged JSON shape.
|
// oxlint-disable-next-line typescript/no-unnecessary-type-parameters -- Test helper lets assertions ascribe logged JSON shape.
|
||||||
function readJsonLog<T extends object>() {
|
function readJsonLog<T extends object>() {
|
||||||
|
|
@ -141,11 +140,6 @@ function expectUnsupportedServiceCheckFailure() {
|
||||||
}
|
}
|
||||||
|
|
||||||
describe("runServiceRestart token drift", () => {
|
describe("runServiceRestart token drift", () => {
|
||||||
beforeAll(async () => {
|
|
||||||
({ runServiceRestart, runServiceStart, runServiceStop, runServiceUninstall } =
|
|
||||||
await import("./lifecycle-core.js"));
|
|
||||||
});
|
|
||||||
|
|
||||||
afterEach(() => {
|
afterEach(() => {
|
||||||
vi.restoreAllMocks();
|
vi.restoreAllMocks();
|
||||||
});
|
});
|
||||||
|
|
@ -479,11 +473,14 @@ describe("runServiceRestart token drift", () => {
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
expect(service.restart).not.toHaveBeenCalled();
|
expect(service.restart).not.toHaveBeenCalled();
|
||||||
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({
|
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith(
|
||||||
targetPid: 1234,
|
expect.objectContaining({
|
||||||
reason: "gateway.restart",
|
env: expect.any(Object),
|
||||||
intent: { waitMs: 2_500 },
|
targetPid: 1234,
|
||||||
});
|
reason: "gateway.restart",
|
||||||
|
intent: { waitMs: 2_500 },
|
||||||
|
}),
|
||||||
|
);
|
||||||
expect(readJsonLog<{ result?: string; message?: string }>()).toMatchObject({
|
expect(readJsonLog<{ result?: string; message?: string }>()).toMatchObject({
|
||||||
result: "restarted",
|
result: "restarted",
|
||||||
message: "Gateway service definition repaired and restarted.",
|
message: "Gateway service definition repaired and restarted.",
|
||||||
|
|
@ -729,10 +726,13 @@ describe("runServiceRestart token drift", () => {
|
||||||
|
|
||||||
await runServiceRestart(createServiceRunArgs());
|
await runServiceRestart(createServiceRunArgs());
|
||||||
|
|
||||||
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({
|
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith(
|
||||||
targetPid: 1234,
|
expect.objectContaining({
|
||||||
reason: "gateway.restart",
|
env: expect.any(Object),
|
||||||
});
|
targetPid: 1234,
|
||||||
|
reason: "gateway.restart",
|
||||||
|
}),
|
||||||
|
);
|
||||||
expect(clearGatewayRestartIntentSync).not.toHaveBeenCalled();
|
expect(clearGatewayRestartIntentSync).not.toHaveBeenCalled();
|
||||||
expect(service.restart).toHaveBeenCalledTimes(1);
|
expect(service.restart).toHaveBeenCalledTimes(1);
|
||||||
});
|
});
|
||||||
|
|
@ -769,13 +769,16 @@ describe("runServiceRestart token drift", () => {
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
|
||||||
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({
|
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith(
|
||||||
targetPid: 1234,
|
expect.objectContaining({
|
||||||
reason: "gateway.restart",
|
env: expect.any(Object),
|
||||||
intent: {
|
targetPid: 1234,
|
||||||
waitMs: 2_500,
|
reason: "gateway.restart",
|
||||||
},
|
intent: {
|
||||||
});
|
waitMs: 2_500,
|
||||||
|
},
|
||||||
|
}),
|
||||||
|
);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("clears restart intent when service-manager restart fails before signaling", async () => {
|
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");
|
await expect(runServiceRestart(createServiceRunArgs())).rejects.toThrow("__exit__:1");
|
||||||
|
|
||||||
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith({
|
expect(writeGatewayRestartIntentSync).toHaveBeenCalledWith(
|
||||||
targetPid: 1234,
|
expect.objectContaining({
|
||||||
reason: "gateway.restart",
|
env: expect.any(Object),
|
||||||
});
|
targetPid: 1234,
|
||||||
|
reason: "gateway.restart",
|
||||||
|
}),
|
||||||
|
);
|
||||||
expect(clearGatewayRestartIntentSync).toHaveBeenCalledOnce();
|
expect(clearGatewayRestartIntentSync).toHaveBeenCalledOnce();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,6 @@ import { readBestEffortConfig } from "../../config/config.js";
|
||||||
import { resolveIsNixMode } from "../../config/paths.js";
|
import { resolveIsNixMode } from "../../config/paths.js";
|
||||||
import { checkTokenDrift } from "../../daemon/service-audit.js";
|
import { checkTokenDrift } from "../../daemon/service-audit.js";
|
||||||
import type { GatewayServiceRestartResult } from "../../daemon/service-types.js";
|
import type { GatewayServiceRestartResult } from "../../daemon/service-types.js";
|
||||||
import { assertGatewayServiceUpdateCurrent } from "../../daemon/service-update-authority.js";
|
|
||||||
import type {
|
import type {
|
||||||
GatewayServiceStartRepairIssue,
|
GatewayServiceStartRepairIssue,
|
||||||
GatewayServiceState,
|
GatewayServiceState,
|
||||||
|
|
@ -19,11 +18,7 @@ import {
|
||||||
import { renderSystemdUnavailableHints } from "../../daemon/systemd-hints.js";
|
import { renderSystemdUnavailableHints } from "../../daemon/systemd-hints.js";
|
||||||
import { isSystemdUserServiceAvailable } from "../../daemon/systemd.js";
|
import { isSystemdUserServiceAvailable } from "../../daemon/systemd.js";
|
||||||
import { isGatewaySecretRefUnavailableError } from "../../gateway/credentials.js";
|
import { isGatewaySecretRefUnavailableError } from "../../gateway/credentials.js";
|
||||||
import {
|
import type { GatewayRestartIntent } from "../../infra/restart-intent.js";
|
||||||
clearGatewayRestartIntentSync,
|
|
||||||
type GatewayRestartIntent,
|
|
||||||
writeGatewayRestartIntentSync,
|
|
||||||
} from "../../infra/restart-intent.js";
|
|
||||||
import { isWSL } from "../../infra/wsl.js";
|
import { isWSL } from "../../infra/wsl.js";
|
||||||
import { defaultRuntime } from "../../runtime.js";
|
import { defaultRuntime } from "../../runtime.js";
|
||||||
import { formatCliCommand } from "../command-format.js";
|
import { formatCliCommand } from "../command-format.js";
|
||||||
|
|
@ -34,6 +29,7 @@ import {
|
||||||
appendServiceLifecycleRepairAudit,
|
appendServiceLifecycleRepairAudit,
|
||||||
createServiceLifecycleMutationAudit,
|
createServiceLifecycleMutationAudit,
|
||||||
} from "./lifecycle-audit.js";
|
} from "./lifecycle-audit.js";
|
||||||
|
import { createServiceRestartIntent } from "./lifecycle-restart-intent.js";
|
||||||
import {
|
import {
|
||||||
buildDaemonServiceSnapshot,
|
buildDaemonServiceSnapshot,
|
||||||
createDaemonActionContext,
|
createDaemonActionContext,
|
||||||
|
|
@ -498,26 +494,18 @@ export async function runServiceRestart(params: {
|
||||||
let handledRecovery: ServiceRecoveryResult<"restarted"> | null = null;
|
let handledRecovery: ServiceRecoveryResult<"restarted"> | null = null;
|
||||||
let handledRepair: ServiceRecoveryResult<"restarted"> | null = null;
|
let handledRepair: ServiceRecoveryResult<"restarted"> | null = null;
|
||||||
let recoveredLoadedState: boolean | null = null;
|
let recoveredLoadedState: boolean | null = null;
|
||||||
let wroteRestartIntent = false;
|
const { prepare: prepareGatewayRestartIntent, clear: clearPreparedRestartIntent } =
|
||||||
const prepareGatewayRestartIntent = async () => {
|
createServiceRestartIntent({
|
||||||
if (params.serviceNoun !== "Gateway" || wroteRestartIntent) {
|
serviceNoun: params.serviceNoun,
|
||||||
return;
|
service: params.service,
|
||||||
}
|
intent: restartIntent,
|
||||||
const runtime = await params.service.readRuntime(process.env).catch(() => null);
|
warn: (message) => {
|
||||||
assertGatewayServiceUpdateCurrent();
|
warnings.push(message);
|
||||||
wroteRestartIntent = writeGatewayRestartIntentSync({
|
if (!json) {
|
||||||
targetPid: runtime?.pid,
|
defaultRuntime.log(message);
|
||||||
reason: "gateway.restart",
|
}
|
||||||
...(restartIntent ? { intent: restartIntent } : {}),
|
},
|
||||||
});
|
});
|
||||||
};
|
|
||||||
const clearPreparedRestartIntent = () => {
|
|
||||||
if (wroteRestartIntent) {
|
|
||||||
assertGatewayServiceUpdateCurrent();
|
|
||||||
clearGatewayRestartIntentSync();
|
|
||||||
wroteRestartIntent = false;
|
|
||||||
}
|
|
||||||
};
|
|
||||||
const emitScheduledRestart = (
|
const emitScheduledRestart = (
|
||||||
restartStatus: ReturnType<typeof describeGatewayServiceRestart>,
|
restartStatus: ReturnType<typeof describeGatewayServiceRestart>,
|
||||||
serviceLoaded: boolean,
|
serviceLoaded: boolean,
|
||||||
|
|
|
||||||
73
src/cli/daemon-cli/lifecycle-restart-intent.ts
Normal file
73
src/cli/daemon-cli/lifecycle-restart-intent.ts
Normal file
|
|
@ -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;
|
||||||
|
}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
@ -60,6 +60,7 @@ vi.mock("../../infra/gateway-lock.js", async (importOriginal) => {
|
||||||
|
|
||||||
vi.mock("../../infra/restart-intent.js", () => ({
|
vi.mock("../../infra/restart-intent.js", () => ({
|
||||||
writeGatewayRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params),
|
writeGatewayRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params),
|
||||||
|
writeGatewayServiceRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params),
|
||||||
clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(),
|
clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(),
|
||||||
}));
|
}));
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -79,6 +79,7 @@ vi.mock("./lifecycle-audit.js", () => ({
|
||||||
vi.mock("../../infra/restart-intent.js", async (original) => ({
|
vi.mock("../../infra/restart-intent.js", async (original) => ({
|
||||||
...(await original<typeof import("../../infra/restart-intent.js")>()),
|
...(await original<typeof import("../../infra/restart-intent.js")>()),
|
||||||
writeGatewayRestartIntentSync: () => true,
|
writeGatewayRestartIntentSync: () => true,
|
||||||
|
writeGatewayServiceRestartIntentSync: () => true,
|
||||||
clearGatewayRestartIntentSync: vi.fn(),
|
clearGatewayRestartIntentSync: vi.fn(),
|
||||||
}));
|
}));
|
||||||
|
|
||||||
|
|
|
||||||
289
src/cli/daemon-cli/lifecycle.restart-intent.test.ts
Normal file
289
src/cli/daemon-cli/lifecycle.restart-intent.test.ts
Normal file
|
|
@ -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<typeof consumeGatewayRestartIntentPayloadSync> = 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<typeof consumeGatewayRestartIntentPayloadSync> = 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,
|
||||||
|
});
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
@ -1,5 +1,5 @@
|
||||||
// Daemon lifecycle tests cover CLI service lifecycle orchestration and cleanup.
|
// 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 { mockSystemAccountHome } from "../../daemon/service.test-helpers.js";
|
||||||
import { captureEnv } from "../../test-utils/env.js";
|
import { captureEnv } from "../../test-utils/env.js";
|
||||||
import {
|
import {
|
||||||
|
|
@ -115,6 +115,7 @@ vi.mock("../../infra/gateway-owner-lease.js", () => ({ readGatewayOwnerLease }))
|
||||||
|
|
||||||
vi.mock("../../infra/restart-intent.js", () => ({
|
vi.mock("../../infra/restart-intent.js", () => ({
|
||||||
writeGatewayRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params),
|
writeGatewayRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params),
|
||||||
|
writeGatewayServiceRestartIntentSync: (params: unknown) => writeGatewayRestartIntentSync(params),
|
||||||
clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(),
|
clearGatewayRestartIntentSync: () => clearGatewayRestartIntentSync(),
|
||||||
}));
|
}));
|
||||||
|
|
||||||
|
|
@ -181,10 +182,9 @@ vi.mock("./lifecycle-core.js", () => ({
|
||||||
runServiceUninstall: vi.fn(),
|
runServiceUninstall: vi.fn(),
|
||||||
}));
|
}));
|
||||||
|
|
||||||
|
const { runDaemonStart, runDaemonRestart, runDaemonStop } = await import("./lifecycle.js");
|
||||||
|
|
||||||
describe("runDaemonRestart health checks", () => {
|
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<typeof captureEnv>;
|
let envSnapshot: ReturnType<typeof captureEnv>;
|
||||||
|
|
||||||
function mockUnmanagedRestart({
|
function mockUnmanagedRestart({
|
||||||
|
|
@ -222,10 +222,6 @@ describe("runDaemonRestart health checks", () => {
|
||||||
return outcome;
|
return outcome;
|
||||||
}
|
}
|
||||||
|
|
||||||
beforeAll(async () => {
|
|
||||||
({ runDaemonStart, runDaemonRestart, runDaemonStop } = await import("./lifecycle.js"));
|
|
||||||
});
|
|
||||||
|
|
||||||
beforeEach(() => {
|
beforeEach(() => {
|
||||||
envSnapshot = captureEnv([
|
envSnapshot = captureEnv([
|
||||||
"OPENCLAW_CONTAINER_HINT",
|
"OPENCLAW_CONTAINER_HINT",
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
// Fast-path argv parser for `openclaw gateway ...` without full Commander registration.
|
// 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 {
|
import {
|
||||||
WINDOWS_TASK_SUPERVISOR_CHILD_FLAG,
|
WINDOWS_TASK_SUPERVISOR_CHILD_FLAG,
|
||||||
WINDOWS_TASK_SUPERVISOR_FLAG,
|
WINDOWS_TASK_SUPERVISOR_FLAG,
|
||||||
|
|
@ -10,49 +11,7 @@ import {
|
||||||
isValueToken,
|
isValueToken,
|
||||||
} from "../infra/cli-root-options.js";
|
} from "../infra/cli-root-options.js";
|
||||||
|
|
||||||
const GATEWAY_RUN_VALUE_FLAGS = new Set([
|
export { isForegroundGatewayRunArgv } from "../../gateway-run-argv.mjs";
|
||||||
"--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");
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Return how many argv tokens a gateway-run option consumes, or 0 when not recognized. */
|
/** Return how many argv tokens a gateway-run option consumes, or 0 when not recognized. */
|
||||||
export function consumeGatewayRunOptionToken(args: ReadonlyArray<string>, index: number): number {
|
export function consumeGatewayRunOptionToken(args: ReadonlyArray<string>, index: number): number {
|
||||||
|
|
|
||||||
|
|
@ -146,6 +146,9 @@ export async function writeRepairCandidate(candidate: string, configChange: bool
|
||||||
"node-sqlite.mjs",
|
"node-sqlite.mjs",
|
||||||
"node-runtime-update.mjs",
|
"node-runtime-update.mjs",
|
||||||
"node-runtime-recovery.mjs",
|
"node-runtime-recovery.mjs",
|
||||||
|
"cli-root-options.mjs",
|
||||||
|
"gateway-run-argv.mjs",
|
||||||
|
"gateway-shutdown-budget.mjs",
|
||||||
"package.json",
|
"package.json",
|
||||||
]) {
|
]) {
|
||||||
await fs.copyFile(path.join(process.cwd(), file), path.join(candidate, file));
|
await fs.copyFile(path.join(process.cwd(), file), path.join(candidate, file));
|
||||||
|
|
|
||||||
|
|
@ -113,6 +113,9 @@ export function createSourceRuntime(root: string): string {
|
||||||
"node-sqlite.mjs",
|
"node-sqlite.mjs",
|
||||||
"node-runtime-update.mjs",
|
"node-runtime-update.mjs",
|
||||||
"node-runtime-recovery.mjs",
|
"node-runtime-recovery.mjs",
|
||||||
|
"cli-root-options.mjs",
|
||||||
|
"gateway-run-argv.mjs",
|
||||||
|
"gateway-shutdown-budget.mjs",
|
||||||
"package.json",
|
"package.json",
|
||||||
"tsconfig.json",
|
"tsconfig.json",
|
||||||
]) {
|
]) {
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@
|
||||||
import fs from "node:fs/promises";
|
import fs from "node:fs/promises";
|
||||||
import { asOptionalRecord, isStringRecord } from "@openclaw/normalization-core/record-coerce";
|
import { asOptionalRecord, isStringRecord } from "@openclaw/normalization-core/record-coerce";
|
||||||
import { hasErrnoCode } from "../infra/errno.js";
|
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 { runExec } from "../process/exec.js";
|
||||||
import type {
|
import type {
|
||||||
GatewayServiceCommandConfig,
|
GatewayServiceCommandConfig,
|
||||||
|
|
@ -9,10 +10,11 @@ import type {
|
||||||
GatewayServiceReadOptions,
|
GatewayServiceReadOptions,
|
||||||
} from "./service-types.js";
|
} from "./service-types.js";
|
||||||
|
|
||||||
|
export { LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS };
|
||||||
|
|
||||||
// launchd defaults to a 10s spawn throttle. Keep that default explicitly so
|
// launchd defaults to a 10s spawn throttle. Keep that default explicitly so
|
||||||
// crash loops back off instead of respawning every second while still allowing
|
// crash loops back off instead of respawning every second while still allowing
|
||||||
// explicit kickstart restarts to take effect.
|
// 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).
|
// launchd stores plist integer values in decimal; 0o077 renders as 63 (owner-only files).
|
||||||
export const LAUNCH_AGENT_POLICY = {
|
export const LAUNCH_AGENT_POLICY = {
|
||||||
RunAtLoad: true,
|
RunAtLoad: true,
|
||||||
|
|
|
||||||
|
|
@ -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 () => {
|
it("keeps interactive no-cache respawns attached to the terminal", async () => {
|
||||||
await markSourceCheckout();
|
await markSourceCheckout();
|
||||||
argv = [process.execPath, entryFile, "tui"];
|
argv = [process.execPath, entryFile, "tui"];
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,7 @@ import { enableCompileCache, getCompileCacheDir } from "node:module";
|
||||||
import os from "node:os";
|
import os from "node:os";
|
||||||
import path from "node:path";
|
import path from "node:path";
|
||||||
import process from "node:process";
|
import process from "node:process";
|
||||||
|
import { isForegroundGatewayRunArgv } from "./cli/gateway-run-argv.js";
|
||||||
import {
|
import {
|
||||||
isForegroundGmailRunArgv,
|
isForegroundGmailRunArgv,
|
||||||
isTerminalInteractiveRespawnArgv,
|
isTerminalInteractiveRespawnArgv,
|
||||||
|
|
@ -117,6 +118,10 @@ function buildOpenClawCompileCacheRespawnPlan(params: {
|
||||||
const env = params.env ?? process.env;
|
const env = params.env ?? process.env;
|
||||||
const argv = process.argv;
|
const argv = process.argv;
|
||||||
const platform = process.platform;
|
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)) {
|
if (isForegroundGmailRunArgv(argv) || shouldKeepNativeHookRelayInProcess(argv, platform)) {
|
||||||
return undefined;
|
return undefined;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -89,6 +89,9 @@ export function useNodeBootstrapArtifactFixtures() {
|
||||||
await write(packageRoot, "node-sqlite.mjs", "export const probe = true;");
|
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-update.mjs", "export const update = true;");
|
||||||
await write(packageRoot, "node-runtime-recovery.mjs", "export const recovery = 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, "node-host-launcher.mjs", "export const launcher = true;");
|
||||||
await write(packageRoot, "scripts/preinstall.mjs", "export {};\n");
|
await write(packageRoot, "scripts/preinstall.mjs", "export {};\n");
|
||||||
await write(
|
await write(
|
||||||
|
|
|
||||||
|
|
@ -51,6 +51,9 @@ const BOOTSTRAP_LAUNCHER_FILES = [
|
||||||
"node-sqlite.mjs",
|
"node-sqlite.mjs",
|
||||||
"node-runtime-update.mjs",
|
"node-runtime-update.mjs",
|
||||||
"node-runtime-recovery.mjs",
|
"node-runtime-recovery.mjs",
|
||||||
|
"cli-root-options.mjs",
|
||||||
|
"gateway-run-argv.mjs",
|
||||||
|
"gateway-shutdown-budget.mjs",
|
||||||
"node-host-launcher.mjs",
|
"node-host-launcher.mjs",
|
||||||
];
|
];
|
||||||
const READ_CONCURRENCY = 16;
|
const READ_CONCURRENCY = 16;
|
||||||
|
|
|
||||||
|
|
@ -130,6 +130,9 @@ describe("worker node enrollment", () => {
|
||||||
path.join(packageRoot, "node-runtime-recovery.mjs"),
|
path.join(packageRoot, "node-runtime-recovery.mjs"),
|
||||||
"export const recovery = true;",
|
"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/entry.js"), "export const ready = true;"),
|
||||||
fs.writeFile(
|
fs.writeFile(
|
||||||
path.join(packageRoot, "dist/build-info.json"),
|
path.join(packageRoot, "dist/build-info.json"),
|
||||||
|
|
|
||||||
|
|
@ -1,200 +1,9 @@
|
||||||
/** CLI token that stops root option scanning and leaves following args positional. */
|
export {
|
||||||
export const FLAG_TERMINATOR = "--";
|
FLAG_TERMINATOR,
|
||||||
|
consumeRootOptionToken,
|
||||||
const ROOT_BOOLEAN_FLAGS = new Set(["--dev", "--no-color"]);
|
consumeRootCommandOptionToken,
|
||||||
const ROOT_VALUE_FLAGS = new Set(["--profile", "--log-level", "--container"]);
|
getRootOptionAwareCommandPath,
|
||||||
|
getCommandPositionalsWithRootOptions,
|
||||||
/** Returns whether a token can be consumed as a root option value. */
|
getCommandArgsWithRootOptions,
|
||||||
export function isValueToken(arg: string | undefined): boolean {
|
isValueToken,
|
||||||
if (!arg || arg === FLAG_TERMINATOR) {
|
} from "../../cli-root-options.mjs";
|
||||||
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<string>, 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<string>;
|
|
||||||
booleanFlags?: ReadonlyArray<string>;
|
|
||||||
valueFlags?: ReadonlyArray<string>;
|
|
||||||
maxPositionals?: number;
|
|
||||||
mode?: "route" | "command-path";
|
|
||||||
};
|
|
||||||
|
|
||||||
function consumeKnownOptionToken(
|
|
||||||
args: ReadonlyArray<string>,
|
|
||||||
index: number,
|
|
||||||
booleanFlags: ReadonlySet<string>,
|
|
||||||
valueFlags: ReadonlySet<string>,
|
|
||||||
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<CommandPositionalsParseOptions, "maxPositionals">,
|
|
||||||
): 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;
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,6 @@
|
||||||
import { randomUUID } from "node:crypto";
|
import { randomUUID } from "node:crypto";
|
||||||
import { hostname } from "node:os";
|
import { hostname } from "node:os";
|
||||||
|
import type { DatabaseSync } from "node:sqlite";
|
||||||
import { isRecord } from "@openclaw/normalization-core/record-coerce";
|
import { isRecord } from "@openclaw/normalization-core/record-coerce";
|
||||||
import { createSubsystemLogger } from "../logging/subsystem.js";
|
import { createSubsystemLogger } from "../logging/subsystem.js";
|
||||||
import { getFileLockProcessStartTime } from "../shared/pid-alive.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 };
|
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(
|
export function readGatewayOwnerLease(
|
||||||
params: {
|
params: {
|
||||||
env?: NodeJS.ProcessEnv;
|
env?: NodeJS.ProcessEnv;
|
||||||
|
|
@ -77,54 +127,8 @@ export function readGatewayOwnerLease(
|
||||||
openStateSchemaReadAdmission?: OpenClawStateSchemaReadAdmission;
|
openStateSchemaReadAdmission?: OpenClawStateSchemaReadAdmission;
|
||||||
} = {},
|
} = {},
|
||||||
): GatewayOwnerLeaseIdentity | undefined {
|
): GatewayOwnerLeaseIdentity | undefined {
|
||||||
const operation = ({
|
const operation = ({ db }: { db: DatabaseSync }) =>
|
||||||
db,
|
readGatewayOwnerLeaseFromDatabase(db, params.port);
|
||||||
}: {
|
|
||||||
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(),
|
|
||||||
};
|
|
||||||
};
|
|
||||||
return params.current || params.openStateSchemaReadAdmission
|
return params.current || params.openStateSchemaReadAdmission
|
||||||
? withExistingOpenClawStateDatabaseCurrentReadOnly(
|
? withExistingOpenClawStateDatabaseCurrentReadOnly(
|
||||||
operation,
|
operation,
|
||||||
|
|
|
||||||
|
|
@ -1,8 +1,7 @@
|
||||||
// Shared Gateway stop policy used by the run loop and native service definitions.
|
export {
|
||||||
const GATEWAY_SHUTDOWN_DRAIN_TIMEOUT_MS = 315_000;
|
GATEWAY_SHUTDOWN_RESERVE_MS,
|
||||||
export const GATEWAY_SHUTDOWN_RESERVE_MS = 10_000;
|
GATEWAY_SUPERVISOR_EXIT_MARGIN_MS,
|
||||||
export const GATEWAY_SUPERVISOR_EXIT_MARGIN_MS = 5_000;
|
GATEWAY_SHUTDOWN_TIMEOUT_MS,
|
||||||
export const GATEWAY_SHUTDOWN_TIMEOUT_MS =
|
GATEWAY_SERVICE_STOP_TIMEOUT_MS,
|
||||||
GATEWAY_SHUTDOWN_DRAIN_TIMEOUT_MS + GATEWAY_SHUTDOWN_RESERVE_MS;
|
LAUNCH_AGENT_EXIT_TIMEOUT_SECONDS,
|
||||||
export const GATEWAY_SERVICE_STOP_TIMEOUT_MS =
|
} from "../../gateway-shutdown-budget.mjs";
|
||||||
GATEWAY_SHUTDOWN_TIMEOUT_MS + GATEWAY_SUPERVISOR_EXIT_MARGIN_MS;
|
|
||||||
|
|
|
||||||
93
src/infra/node-runtime-recovery.gateway-drain.test.ts
Normal file
93
src/infra/node-runtime-recovery.gateway-drain.test.ts
Normal file
|
|
@ -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);
|
||||||
|
});
|
||||||
85
src/infra/node-runtime-recovery.shutdown.test.ts
Normal file
85
src/infra/node-runtime-recovery.shutdown.test.ts
Normal file
|
|
@ -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<typeof import("node:child_process")>()),
|
||||||
|
spawn,
|
||||||
|
}));
|
||||||
|
|
||||||
|
let child: ChildProcess;
|
||||||
|
let kill: MockInstance<ChildProcess["kill"]>;
|
||||||
|
let exit: MockInstance<typeof process.exit>;
|
||||||
|
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<typeof process.exit>());
|
||||||
|
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]);
|
||||||
|
});
|
||||||
|
|
@ -1,13 +1,16 @@
|
||||||
import { existsSync } from "node:fs";
|
import { existsSync } from "node:fs";
|
||||||
|
import type { DatabaseSync } from "node:sqlite";
|
||||||
// Persists short-lived gateway restart intent for supervisor SIGTERM handoff.
|
// Persists short-lived gateway restart intent for supervisor SIGTERM handoff.
|
||||||
import { asPositiveSafeInteger } from "@openclaw/normalization-core/number-coercion";
|
import { asPositiveSafeInteger } from "@openclaw/normalization-core/number-coercion";
|
||||||
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
|
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
|
||||||
|
import { resolveSystemdServiceName } from "../daemon/systemd-service-files.js";
|
||||||
import { createSubsystemLogger } from "../logging/subsystem.js";
|
import { createSubsystemLogger } from "../logging/subsystem.js";
|
||||||
import { runExistingOpenClawStateWriteTransaction } from "../state/openclaw-state-db-existing-write.js";
|
import { runExistingOpenClawStateWriteTransaction } from "../state/openclaw-state-db-existing-write.js";
|
||||||
import type { DB as OpenClawStateKyselyDatabase } from "../state/openclaw-state-db.generated.js";
|
import type { DB as OpenClawStateKyselyDatabase } from "../state/openclaw-state-db.generated.js";
|
||||||
import { openOpenClawStateDatabase } from "../state/openclaw-state-db.js";
|
import { openOpenClawStateDatabase } from "../state/openclaw-state-db.js";
|
||||||
import { resolveOpenClawStateSqlitePath } from "../state/openclaw-state-db.paths.js";
|
import { resolveOpenClawStateSqlitePath } from "../state/openclaw-state-db.paths.js";
|
||||||
import { OPENCLAW_STATE_SCHEMA_SQL } from "../state/openclaw-state-schema.js";
|
import { OPENCLAW_STATE_SCHEMA_SQL } from "../state/openclaw-state-schema.js";
|
||||||
|
import { readGatewayOwnerLeaseFromDatabase } from "./gateway-owner-lease.js";
|
||||||
import {
|
import {
|
||||||
executeSqliteQuerySync,
|
executeSqliteQuerySync,
|
||||||
executeSqliteQueryTakeFirstSync,
|
executeSqliteQueryTakeFirstSync,
|
||||||
|
|
@ -65,6 +68,67 @@ export function writeGatewayRestartIntentSync(opts: {
|
||||||
if (targetPid === null) {
|
if (targetPid === null) {
|
||||||
return false;
|
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;
|
const env = opts.env ?? process.env;
|
||||||
try {
|
try {
|
||||||
if (!existsSync(resolveOpenClawStateSqlitePath(env))) {
|
if (!existsSync(resolveOpenClawStateSqlitePath(env))) {
|
||||||
|
|
@ -78,10 +142,17 @@ export function writeGatewayRestartIntentSync(opts: {
|
||||||
opts.intent.waitMs >= 0
|
opts.intent.waitMs >= 0
|
||||||
? Math.floor(opts.intent.waitMs)
|
? Math.floor(opts.intent.waitMs)
|
||||||
: null;
|
: null;
|
||||||
const createdAt = Date.now();
|
|
||||||
// The old Gateway still owns the schema until the restart hands off.
|
// The old Gateway still owns the schema until the restart hands off.
|
||||||
runExistingOpenClawStateWriteTransaction(
|
return runExistingOpenClawStateWriteTransaction(
|
||||||
({ db }) => {
|
({ 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<GatewayRestartIntentDatabase>(db);
|
const stateDb = getNodeSqliteKysely<GatewayRestartIntentDatabase>(db);
|
||||||
executeSqliteQuerySync(
|
executeSqliteQuerySync(
|
||||||
db,
|
db,
|
||||||
|
|
@ -109,12 +180,14 @@ export function writeGatewayRestartIntentSync(opts: {
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
);
|
);
|
||||||
|
return true;
|
||||||
},
|
},
|
||||||
{ env },
|
{ env },
|
||||||
{ schemaSql: schema, operationLabel: "gateway.restart-intent.write" },
|
{ schemaSql: schema, operationLabel: "gateway.restart-intent.write" },
|
||||||
);
|
);
|
||||||
return true;
|
|
||||||
} catch (err) {
|
} 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)}`);
|
restartLog.warn(`failed to write gateway restart intent: ${String(err)}`);
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -66,6 +66,9 @@ async function writePackage(root: string, version: string, body: string, schema
|
||||||
"node-sqlite.mjs",
|
"node-sqlite.mjs",
|
||||||
"node-runtime-update.mjs",
|
"node-runtime-update.mjs",
|
||||||
"node-runtime-recovery.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));
|
await fs.copyFile(path.resolve(name), path.join(root, name));
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -46,6 +46,18 @@ async function makeLauncherVersionFixture(
|
||||||
path.resolve(process.cwd(), "node-runtime-recovery.mjs"),
|
path.resolve(process.cwd(), "node-runtime-recovery.mjs"),
|
||||||
path.join(fixtureRoot, "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.mkdir(path.join(fixtureRoot, "dist"), { recursive: true });
|
||||||
await fs.writeFile(
|
await fs.writeFile(
|
||||||
path.join(fixtureRoot, "package.json"),
|
path.join(fixtureRoot, "package.json"),
|
||||||
|
|
|
||||||
|
|
@ -36,6 +36,18 @@ async function makeLauncherFixture(fixtureRoots: string[]): Promise<string> {
|
||||||
path.resolve(process.cwd(), "node-runtime-recovery.mjs"),
|
path.resolve(process.cwd(), "node-runtime-recovery.mjs"),
|
||||||
path.join(fixtureRoot, "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(
|
await fs.copyFile(
|
||||||
path.resolve(process.cwd(), "node-sqlite.mjs"),
|
path.resolve(process.cwd(), "node-sqlite.mjs"),
|
||||||
path.join(fixtureRoot, "node-sqlite.mjs"),
|
path.join(fixtureRoot, "node-sqlite.mjs"),
|
||||||
|
|
|
||||||
|
|
@ -34,6 +34,9 @@ it.runIf(process.platform !== "win32")(
|
||||||
"node-version.mjs",
|
"node-version.mjs",
|
||||||
"node-runtime-update.mjs",
|
"node-runtime-update.mjs",
|
||||||
"node-runtime-recovery.mjs",
|
"node-runtime-recovery.mjs",
|
||||||
|
"cli-root-options.mjs",
|
||||||
|
"gateway-run-argv.mjs",
|
||||||
|
"gateway-shutdown-budget.mjs",
|
||||||
"node-sqlite.mjs",
|
"node-sqlite.mjs",
|
||||||
]) {
|
]) {
|
||||||
copyFileSync(filename, path.join(checkoutRoot, filename));
|
copyFileSync(filename, path.join(checkoutRoot, filename));
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue