mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
Package recovery treated checked Bun executables as Node and dropped configured macOS SQLite selection. Carry admitted runtime facts and the shared environment choice into durable standalone recovery, preserving version-1 journals, backup custody, and separately owned service restart. Verified focused suites on Node and Bun, real macOS custom-library recovery from a clean shell, published-driver updates and Node-free interrupted recovery on AWS, static/import checks, and independent review. CI fixture and cell-lifetime repairs retain the existing assertions and product timeout policies.
2829 lines
112 KiB
TypeScript
2829 lines
112 KiB
TypeScript
// OpenClaw test instance tests cover spawned test instance lifecycle.
|
||
import { AsyncLocalStorage, createHook } from "node:async_hooks";
|
||
import { execFile, spawn } from "node:child_process";
|
||
import { EventEmitter, once } from "node:events";
|
||
import { statSync } from "node:fs";
|
||
import fs from "node:fs/promises";
|
||
import { createServer } from "node:http";
|
||
import net, { type Socket } from "node:net";
|
||
import { tmpdir } from "node:os";
|
||
import path from "node:path";
|
||
import { PassThrough } from "node:stream";
|
||
import { setTimeout as delay } from "node:timers/promises";
|
||
import { promisify } from "node:util";
|
||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||
import {
|
||
inspectManagedProcessGroup,
|
||
terminateManagedChild,
|
||
} from "../../scripts/lib/managed-child-process.mts";
|
||
import { hasErrnoCode } from "../../src/infra/errno.js";
|
||
import { createFileLockManager } from "../../src/infra/file-lock-manager.js";
|
||
import {
|
||
drainFileLockStateForTest,
|
||
FILE_LOCK_TIMEOUT_ERROR_CODE,
|
||
resetFileLockStateForTest,
|
||
} from "../../src/plugin-sdk/file-lock.js";
|
||
import { resolveMaxOutputBytes } from "../../src/process/exec-output.js";
|
||
import { withEnvAsync } from "../../src/test-utils/env.js";
|
||
import { resolveTestNodeExecPath } from "../../src/test-utils/node-process.js";
|
||
import {
|
||
createOpenClawTestInstance,
|
||
formatGatewayReadinessDiagnostic,
|
||
GatewayStartupRefusedError,
|
||
type GatewayReadinessDiagnostic,
|
||
type OpenClawTestInstance,
|
||
testing,
|
||
} from "./openclaw-test-instance.js";
|
||
import { isProcessAlive, waitForDead } from "./process-wait.js";
|
||
import { awaitGateBeforeSettlement, createDeferred, withinTest } from "./promise.js";
|
||
import { runQaGatewayFixture } from "./qa-gateway-cleanup.js";
|
||
|
||
const MIGRATION_CONVERGENCE_REFUSAL =
|
||
"OpenClaw plugin migration inputs changed during startup convergence;";
|
||
const RESTART_MARKER =
|
||
"[openclaw-test-instance] restarting gateway after migration convergence refusal";
|
||
const LEGACY_STORE_PATH = "/fixture/sessions/sessions.json";
|
||
const LEGACY_MIGRATION_REFUSAL = `Legacy session store requires migration: ${LEGACY_STORE_PATH}. Run "openclaw doctor --fix" against the same state/config before starting OpenClaw.`;
|
||
const PROFILED_LEGACY_MIGRATION_REFUSAL = `Legacy session store requires migration: ${LEGACY_STORE_PATH}. Run "openclaw --profile qa-fixture doctor --fix" against the same state/config before starting OpenClaw.`;
|
||
const LEGACY_STARTUP_FAILURE = [
|
||
"OpenClaw startup migrations did not complete cleanly; refusing to report the gateway ready.",
|
||
`- Legacy sessions store unreadable; left in place at ${LEGACY_STORE_PATH}`,
|
||
'- Migration step "agent-dir" was not run because prior step "sessions" refused execution.',
|
||
'Run "openclaw doctor --fix" against the same state/config, then restart the gateway.',
|
||
].join("\n");
|
||
const CONTROL_REQUEST_SOURCE = `
|
||
import { get as getControl } from "node:http";
|
||
function waitForControl(url) {
|
||
return new Promise((resolve, reject) => {
|
||
getControl(url, { agent: false }, (response) => {
|
||
response.once("error", reject);
|
||
response.once("end", resolve);
|
||
response.resume();
|
||
}).once("error", reject);
|
||
});
|
||
}
|
||
`;
|
||
const fakeInstances: {
|
||
instance: Awaited<ReturnType<typeof createOpenClawTestInstance>>;
|
||
writerPidPath?: string;
|
||
}[] = [];
|
||
const fakeRoots: string[] = [];
|
||
const fakeOperations: Promise<unknown>[] = [];
|
||
const fakeControls: FakeGatewayControl[] = [];
|
||
|
||
type FakeGatewayControl = {
|
||
url: string;
|
||
reached: Promise<void>;
|
||
launches: number[];
|
||
parents: number[];
|
||
observers: { beforeRelease: () => void; onLaunch: () => void };
|
||
unblock: () => void;
|
||
release: () => Promise<void>;
|
||
close: () => Promise<void>;
|
||
};
|
||
|
||
type FakeGatewayAttempt = {
|
||
argv: string[];
|
||
config: unknown;
|
||
cwd: string;
|
||
env: Record<string, string | undefined>;
|
||
pid: number;
|
||
port: number;
|
||
};
|
||
|
||
afterEach(async () => {
|
||
const controls = fakeControls.splice(0);
|
||
for (const control of controls) {
|
||
control.unblock();
|
||
}
|
||
await Promise.allSettled(fakeOperations.splice(0));
|
||
const results = await Promise.allSettled(
|
||
fakeInstances.splice(0).map(async (owner) => {
|
||
const { instance, writerPidPath } = owner;
|
||
// Baseline failures can spawn after cleanup has already marked itself done.
|
||
try {
|
||
await runQaGatewayFixture(
|
||
() => instance.stopGateway(),
|
||
async () => {
|
||
if (writerPidPath) {
|
||
const pid = Number(await fs.readFile(writerPidPath, "utf8"));
|
||
expect(Number.isSafeInteger(pid) && pid > 1).toBe(true);
|
||
await waitForDead(pid, 5_000);
|
||
}
|
||
},
|
||
);
|
||
await instance.cleanup();
|
||
} catch (error) {
|
||
fakeInstances.push(owner);
|
||
throw error;
|
||
}
|
||
}),
|
||
);
|
||
const controlResults = await Promise.allSettled(
|
||
controls.map(async (control) => {
|
||
try {
|
||
await control.close();
|
||
} catch (error) {
|
||
fakeControls.push(control);
|
||
throw error;
|
||
}
|
||
}),
|
||
);
|
||
const failures = [...results, ...controlResults].flatMap((result) =>
|
||
result.status === "rejected" ? [result.reason] : [],
|
||
);
|
||
if (failures.length > 0) {
|
||
throw new AggregateError(failures, "fake Gateway cleanup failed; owners and roots retained");
|
||
}
|
||
await Promise.all(fakeRoots.splice(0).map((root) => fs.rm(root, { recursive: true })));
|
||
});
|
||
|
||
function trackOperation<T>(operation: Promise<T>): Promise<T> {
|
||
fakeOperations.push(operation.catch(() => undefined));
|
||
return operation;
|
||
}
|
||
|
||
async function observeCliExit() {
|
||
const managed = await import("../../scripts/lib/managed-child-process.mts");
|
||
const run = managed.runManagedCommand;
|
||
const exited = createDeferred();
|
||
// Observe the existing process handoff without replacing spawning or pipe ownership.
|
||
const observer = vi.spyOn(managed, "runManagedCommand").mockImplementation((options) =>
|
||
run({
|
||
...options,
|
||
onReady(child) {
|
||
child.once("exit", () => exited.resolve());
|
||
options.onReady?.(child);
|
||
},
|
||
}),
|
||
);
|
||
return { exited: exited.promise, restore: () => observer.mockRestore() };
|
||
}
|
||
|
||
// Detached inherited writers have no ChildProcess handle in this test process.
|
||
async function waitForProcessExit(pid: number, signal: AbortSignal): Promise<void> {
|
||
try {
|
||
for (;;) {
|
||
signal.throwIfAborted();
|
||
if (!isProcessAlive(pid)) {
|
||
return;
|
||
}
|
||
await delay(5, undefined, { signal });
|
||
}
|
||
} catch (error) {
|
||
if (signal.aborted) {
|
||
throw new Error(`process still alive: ${pid}`, { cause: error });
|
||
}
|
||
throw error;
|
||
}
|
||
}
|
||
|
||
async function createGatewayControl(): Promise<FakeGatewayControl> {
|
||
const reached = createDeferred();
|
||
const released = createDeferred();
|
||
const launches: number[] = [];
|
||
const parents: number[] = [];
|
||
const observers = { beforeRelease: () => {}, onLaunch: () => {} };
|
||
const sockets = new Set<Socket>();
|
||
const server = createServer((request, response) => {
|
||
const url = new URL(request.url ?? "/", "http://127.0.0.1");
|
||
if (url.pathname === "/wait") {
|
||
reached.resolve();
|
||
void released.promise.then(() => response.end("released"));
|
||
return;
|
||
}
|
||
if (url.pathname === "/release") {
|
||
observers.beforeRelease();
|
||
released.resolve();
|
||
} else if (url.pathname === "/launch") {
|
||
launches.push(Number(url.searchParams.get("pid")));
|
||
parents.push(Number(url.searchParams.get("parent")));
|
||
observers.onLaunch();
|
||
}
|
||
response.end("ok");
|
||
});
|
||
server.on("connection", (socket) => {
|
||
sockets.add(socket);
|
||
socket.once("close", () => sockets.delete(socket));
|
||
});
|
||
await new Promise<void>((resolve, reject) => {
|
||
server.once("error", reject);
|
||
server.listen(0, "127.0.0.1", resolve);
|
||
});
|
||
const address = server.address();
|
||
if (!address || typeof address === "string") {
|
||
throw new Error("control server has no port");
|
||
}
|
||
const url = `http://127.0.0.1:${address.port}`;
|
||
const control = {
|
||
url,
|
||
reached: reached.promise,
|
||
launches,
|
||
parents,
|
||
observers,
|
||
unblock: () => released.resolve(),
|
||
release: async () => {
|
||
const response = await fetch(`${url}/release`);
|
||
await response.text();
|
||
},
|
||
close: async () => {
|
||
released.resolve();
|
||
await new Promise<void>((resolve, reject) => {
|
||
server.close((error) => (error ? reject(error) : resolve()));
|
||
for (const socket of sockets) {
|
||
socket.destroy();
|
||
}
|
||
});
|
||
},
|
||
};
|
||
fakeControls.push(control);
|
||
return control;
|
||
}
|
||
|
||
async function createFakeGateway(
|
||
sequence: string,
|
||
startTimeoutMs = 1_000,
|
||
stopTimeoutMs = 1_500,
|
||
control?: { url: string; holdPreparation?: boolean },
|
||
ownerOptions: {
|
||
signal?: AbortSignal;
|
||
verifyCleanup?: (cleanup: () => Promise<void>) => Promise<void>;
|
||
} = {},
|
||
) {
|
||
const cwd = await fs.mkdtemp(path.join(tmpdir(), "openclaw-test-instance-gateway-"));
|
||
fakeRoots.push(cwd);
|
||
const distDir = path.join(cwd, "dist");
|
||
const tracePath = path.join(cwd, "attempts.jsonl");
|
||
// Diagnostic runs keep these receipts outside Vitest's disposable temp tree.
|
||
const processReceipt = `
|
||
const registry = ${JSON.stringify(process.env.OPENCLAW_HELPER_PROOF_PID_REGISTRY ?? null)};
|
||
function recordFixtureProcess(pid) {
|
||
if (!registry) return;
|
||
let identity;
|
||
try {
|
||
identity = execFileSync("/bin/ps", ["-p", String(pid), "-o", "pgid=", "-o", "lstart=", "-o", "command="], { encoding: "utf8", env: { PATH: "/usr/bin:/bin", LC_ALL: "C", TZ: "UTC" }, timeout: 1_000 }).trim();
|
||
} catch (error) {
|
||
if (error.status === 1) return;
|
||
throw error;
|
||
}
|
||
appendFileSync(registry, JSON.stringify({ pid, cwd: process.cwd(), identity }) + "\\n");
|
||
}
|
||
recordFixtureProcess(process.pid);
|
||
`;
|
||
await fs.mkdir(distDir);
|
||
await Promise.all([
|
||
...(control?.holdPreparation
|
||
? []
|
||
: [
|
||
fs.writeFile(path.join(distDir, ".buildstamp"), ""),
|
||
fs.writeFile(path.join(distDir, ".runtime-postbuildstamp"), ""),
|
||
]),
|
||
fs.writeFile(
|
||
path.join(distDir, "index.mjs"),
|
||
`
|
||
import { execFileSync, spawn } from "node:child_process";
|
||
import { appendFileSync, existsSync, readFileSync, writeFileSync } from "node:fs";
|
||
import { createServer } from "node:http";
|
||
${CONTROL_REQUEST_SOURCE}
|
||
${processReceipt}
|
||
const tracePath = process.env.OPENCLAW_FAKE_GATEWAY_TRACE;
|
||
const controlUrl = process.env.OPENCLAW_FAKE_GATEWAY_CONTROL;
|
||
function spawnInheritedWriter(stream, output) {
|
||
// Windows must keep the writer alive after its leader exits; the HTTP gate owns release.
|
||
const delayed = spawn(process.execPath, ["-e", 'require("node:http").get(process.argv[1] + "/wait", (response) => { response.resume(); response.on("end", () => process[process.argv[2]].write(process.argv[3], () => process.exit(0))); });', controlUrl, stream, output], { detached: process.platform === "win32", stdio: ["ignore", stream === "stdout" ? "inherit" : "ignore", stream === "stderr" ? "inherit" : "ignore"] });
|
||
recordFixtureProcess(delayed.pid);
|
||
writeFileSync(tracePath + ".writer-pid", String(delayed.pid));
|
||
}
|
||
if (controlUrl) await waitForControl(controlUrl + "/launch?pid=" + process.pid + "&parent=" + process.ppid);
|
||
const countPath = tracePath + ".count";
|
||
let attempt = 1;
|
||
try { attempt = Number(readFileSync(countPath, "utf8")) + 1; } catch {}
|
||
writeFileSync(countPath, String(attempt));
|
||
const argv = process.argv.slice(2);
|
||
const port = Number(argv[argv.indexOf("--port") + 1]);
|
||
const env = Object.fromEntries(["HOME", "OPENCLAW_CONFIG_PATH", "OPENCLAW_GATEWAY_TOKEN", "OPENCLAW_STATE_DIR"].map((key) => [key, process.env[key]]));
|
||
appendFileSync(tracePath, JSON.stringify({ argv, config: JSON.parse(readFileSync(process.env.OPENCLAW_CONFIG_PATH, "utf8")), cwd: process.cwd(), env, pid: process.pid, port }) + "\\n");
|
||
const kind = (process.env.OPENCLAW_FAKE_GATEWAY_SEQUENCE || "ready").split(",")[attempt - 1] || "ready";
|
||
if (kind === "cli-json") {
|
||
if (argv[1] === "overflow-close") {
|
||
const splitCodePoint = Buffer.from([0xf0, 0x9f, 0xa6, 0x8a]);
|
||
const first = Buffer.concat([
|
||
Buffer.alloc(Number(argv[0]) - 2, 0x61),
|
||
splitCodePoint.subarray(0, 2),
|
||
]);
|
||
await new Promise((resolve) => process.stdout.write(first, resolve));
|
||
await waitForControl(controlUrl + "/wait");
|
||
await new Promise((resolve) =>
|
||
process.stdout.end(
|
||
Buffer.concat([
|
||
splitCodePoint.subarray(2),
|
||
Buffer.from("\\ntrailing overflow output\\n"),
|
||
]),
|
||
resolve,
|
||
),
|
||
);
|
||
setInterval(() => {}, 1_000);
|
||
await new Promise(() => {});
|
||
}
|
||
const json = JSON.stringify({ first: "complete", payload: "é".repeat(Number(argv[0])), providerCredentialPresent: Object.hasOwn(process.env, "OPENAI_API_KEY"), last: "complete" });
|
||
await Promise.all([
|
||
new Promise((resolve) => process.stdout.write(json, resolve)),
|
||
new Promise((resolve) => process.stderr.write("discarded diagnostic " + "x".repeat(300 * 1024) + "\\nrecent cli diagnostic\\n", resolve)),
|
||
]);
|
||
// Keep the overflow producer alive so cancellation owns a live process.
|
||
if (argv[1] === "overflow") {
|
||
setInterval(() => {}, 1_000);
|
||
await new Promise(() => {});
|
||
}
|
||
process.exit(0);
|
||
}
|
||
process.stdout.write("fake gateway attempt " + attempt + "\\n");
|
||
if (kind === "cli" || kind === "cli-drain") {
|
||
process.stderr.write("cli diagnostic\\n");
|
||
if (argv[0] === "wait") {
|
||
setInterval(() => {}, 1_000);
|
||
if (controlUrl) await waitForControl(controlUrl + "/wait");
|
||
await new Promise(() => {});
|
||
}
|
||
if (kind === "cli-drain") {
|
||
spawnInheritedWriter("stdout", "drained cli output\\n");
|
||
process.exit(0);
|
||
}
|
||
if (argv[0] === "large") {
|
||
await new Promise((resolve) => process.stdout.write("x".repeat(300 * 1024), resolve));
|
||
process.exit(0);
|
||
}
|
||
process.exit(Number(argv[0]));
|
||
}
|
||
const refusal = ${JSON.stringify(MIGRATION_CONVERGENCE_REFUSAL)};
|
||
const legacyRefusal = kind.startsWith("startup-") ? ${JSON.stringify(LEGACY_STARTUP_FAILURE)} : kind.startsWith("profile-") ? ${JSON.stringify(PROFILED_LEGACY_MIGRATION_REFUSAL)} : ${JSON.stringify(LEGACY_MIGRATION_REFUSAL)};
|
||
if (kind === "legacy-refuse" || kind === "startup-legacy-refuse" || kind === "profile-legacy-refuse") { process.stderr.write(legacyRefusal + "\\n"); process.exit(78); }
|
||
if (kind === "late-legacy-refuse" || kind === "startup-late-legacy-refuse") {
|
||
spawnInheritedWriter("stderr", legacyRefusal + "\\n");
|
||
process.exit(78);
|
||
}
|
||
if (kind === "legacy-stdout") { process.stdout.write(legacyRefusal + "\\n"); process.exit(78); }
|
||
if (kind === "legacy-status1") { process.stderr.write(legacyRefusal + "\\n"); process.exit(1); }
|
||
if (kind === "legacy-no-advice") { process.stderr.write("Legacy session store requires migration: " + ${JSON.stringify(LEGACY_STORE_PATH)} + "\\n"); process.exit(78); }
|
||
if (kind === "startup-no-advice") { process.stderr.write(legacyRefusal.split("\\n").slice(0, -1).join("\\n") + "\\n"); process.exit(78); }
|
||
if (kind === "startup-warning") { process.stderr.write(legacyRefusal.split("\\n").slice(1).join("\\n") + "\\n"); process.exit(78); }
|
||
if (kind === "config-refuse") { process.stderr.write("unrelated configuration failure\\n"); process.exit(78); }
|
||
if (kind === "refuse") { process.stderr.write(refusal + " fixture\\n"); process.exit(1); }
|
||
if (kind === "late-refuse") {
|
||
spawnInheritedWriter("stderr", refusal + " delayed fixture\\n");
|
||
process.exit(1);
|
||
}
|
||
if (kind === "resist-after-exit" || kind === "resist-ignored-after-exit") {
|
||
const resistant = spawn(process.execPath, ["-e", 'const fs = require("node:fs");fs.writeFileSync(process.argv[1], String(process.pid));process.on("SIGTERM", () => { fs.appendFileSync(process.argv[2], "SIGTERM"); process.stderr.write("SIGTERM"); });process.send("ready");setInterval(() => {}, 1_000);', tracePath + ".resistant-pid", tracePath + ".signals"], { stdio: ["ignore", "ignore", kind === "resist-after-exit" ? "inherit" : "ignore", "ipc"] });
|
||
await new Promise((resolve) => resistant.once("message", resolve));
|
||
recordFixtureProcess(resistant.pid);
|
||
await waitForControl(controlUrl + "/wait");
|
||
process.exit(1);
|
||
}
|
||
if (kind === "terminal-drain" || kind === "refusal-drain") {
|
||
const draining = spawn(process.execPath, ["-e", 'const fs = require("node:fs");const release = process.argv[1];const deadline = Date.now() + 5_000;const timer = setInterval(() => { if (fs.existsSync(release) || Date.now() >= deadline) clearInterval(timer); }, 10);', tracePath + ".draining-release"], { detached: true, stdio: ["ignore", "ignore", "inherit"] });
|
||
draining.unref();
|
||
recordFixtureProcess(draining.pid);
|
||
writeFileSync(tracePath + ".draining-pid", String(draining.pid));
|
||
process.stderr.write(kind === "refusal-drain" ? refusal + " held fixture\\n" : "terminal startup failure\\n"); process.exit(kind === "refusal-drain" ? 1 : 7);
|
||
}
|
||
if (kind === "near") { process.stderr.write(refusal.slice(0, -1) + " fixture\\n"); process.exit(1); }
|
||
if (kind === "stdout") { process.stdout.write(refusal + " fixture\\n"); process.exit(1); }
|
||
if (kind === "status2") { process.stderr.write(refusal + " fixture\\n"); process.exit(2); }
|
||
if (kind === "signal") { process.stderr.write(refusal + " fixture\\n"); process.kill(process.pid, "SIGTERM"); }
|
||
if (kind === "late-unrelated") {
|
||
spawnInheritedWriter("stderr", "unrelated startup failure\\n");
|
||
process.exit(1);
|
||
}
|
||
if (kind === "held-unrelated") await waitForControl(controlUrl + "/wait");
|
||
if (kind === "unrelated" || kind === "held-unrelated") { process.stderr.write("unrelated startup failure\\n"); process.exit(1); }
|
||
const server = createServer(async (req, res) => {
|
||
if (req.url === "/readyz" && kind === "held-ready") await waitForControl(controlUrl + "/wait");
|
||
res.writeHead(req.url === "/readyz" ? 200 : 404, { "content-type": "application/json" });
|
||
res.end(JSON.stringify({ ready: req.url === "/readyz" && kind !== "never-ready" }));
|
||
});
|
||
process.on("SIGTERM", () => server.close(() => process.exit(0))); server.listen(port, "127.0.0.1");
|
||
`,
|
||
),
|
||
]);
|
||
if (control?.holdPreparation) {
|
||
await fs.mkdir(path.join(cwd, "scripts"));
|
||
// This is a fixture bootstrap, not the repository's build entrypoint.
|
||
await fs.writeFile(
|
||
path.join(cwd, "scripts", "run-node.mjs"),
|
||
`import { appendFileSync, writeFileSync } from "node:fs";
|
||
import { execFileSync } from "node:child_process";
|
||
${CONTROL_REQUEST_SOURCE}
|
||
${processReceipt}
|
||
await waitForControl(${JSON.stringify(`${control.url}/wait`)});
|
||
writeFileSync("dist/.buildstamp", "");
|
||
writeFileSync("dist/.runtime-postbuildstamp", "");
|
||
`,
|
||
);
|
||
}
|
||
const instance = await createOpenClawTestInstance({
|
||
name: `fake-gateway-${path.basename(cwd)}`,
|
||
cwd,
|
||
env: {
|
||
OPENCLAW_FAKE_GATEWAY_SEQUENCE: sequence,
|
||
OPENCLAW_FAKE_GATEWAY_TRACE: tracePath,
|
||
OPENCLAW_FAKE_GATEWAY_CONTROL: control?.url,
|
||
},
|
||
startTimeoutMs,
|
||
stopTimeoutMs,
|
||
...ownerOptions,
|
||
});
|
||
fakeInstances.push({
|
||
instance,
|
||
// Join inherited writers after releasing their HTTP gate, including failed commands/startup.
|
||
writerPidPath: sequence
|
||
.split(",")
|
||
.some(
|
||
(kind) =>
|
||
kind === "late-unrelated" ||
|
||
kind === "late-refuse" ||
|
||
kind === "late-legacy-refuse" ||
|
||
kind === "startup-late-legacy-refuse" ||
|
||
kind === "cli-drain",
|
||
)
|
||
? `${tracePath}.writer-pid`
|
||
: undefined,
|
||
});
|
||
return {
|
||
instance,
|
||
tracePath,
|
||
readAttempts: async (): Promise<FakeGatewayAttempt[]> =>
|
||
(await fs.readFile(tracePath, "utf8"))
|
||
.trim()
|
||
.split("\n")
|
||
.filter(Boolean)
|
||
.map((line) => JSON.parse(line) as FakeGatewayAttempt),
|
||
};
|
||
}
|
||
|
||
async function expectPathMissing(targetPath: string): Promise<void> {
|
||
try {
|
||
await fs.stat(targetPath);
|
||
} catch (error) {
|
||
expect((error as NodeJS.ErrnoException).code).toBe("ENOENT");
|
||
return;
|
||
}
|
||
throw new Error(`Expected missing path: ${targetPath}`);
|
||
}
|
||
|
||
async function isPortReserved(port: number): Promise<boolean> {
|
||
const competitor = net.createServer();
|
||
try {
|
||
await new Promise<void>((resolve, reject) => {
|
||
competitor.once("error", reject);
|
||
competitor.listen(port, "127.0.0.1", resolve);
|
||
});
|
||
return false;
|
||
} catch (error) {
|
||
if (!hasErrnoCode(error, "EADDRINUSE")) {
|
||
throw error;
|
||
}
|
||
return true;
|
||
} finally {
|
||
if (competitor.listening) {
|
||
await new Promise<void>((resolve, reject) => {
|
||
competitor.close((error) => (error ? reject(error) : resolve()));
|
||
});
|
||
}
|
||
}
|
||
}
|
||
|
||
async function startGatewayForPortLifecycle(
|
||
instance: OpenClawTestInstance,
|
||
signal: AbortSignal,
|
||
): Promise<void> {
|
||
// Port custody is ordered by real child exit/readiness, not native bootstrap speed.
|
||
// Freeze only startup policy time; HTTP, child events, and teardown stay real.
|
||
signal.throwIfAborted();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(Date.now());
|
||
const restoreClock = () => clock.mockRestore();
|
||
signal.addEventListener("abort", restoreClock, { once: true });
|
||
try {
|
||
await trackOperation(instance.startGateway());
|
||
} finally {
|
||
restoreClock();
|
||
signal.removeEventListener("abort", restoreClock);
|
||
}
|
||
}
|
||
|
||
function createGatewayProcessState(
|
||
overrides: Partial<{
|
||
pid: number;
|
||
exitCode: number | null;
|
||
signalCode: NodeJS.Signals | null;
|
||
}> = {},
|
||
) {
|
||
return Object.assign(new EventEmitter(), {
|
||
exitCode: null as number | null,
|
||
signalCode: null as NodeJS.Signals | null,
|
||
...overrides,
|
||
});
|
||
}
|
||
|
||
function readReadinessReceipt(error: unknown) {
|
||
if (!(error instanceof Error)) {
|
||
throw new Error("expected a readiness failure", { cause: error });
|
||
}
|
||
const prefix = "[openclaw-test-instance] readiness ";
|
||
const lines = error.message.split("\n").filter((line) => line.startsWith(prefix));
|
||
expect(lines).toHaveLength(1);
|
||
const line = lines[0]!;
|
||
expect(Buffer.byteLength(line, "utf8")).toBeLessThan(1_024);
|
||
return JSON.parse(line.slice(prefix.length));
|
||
}
|
||
|
||
function createStalledReadinessFetch(phase: "headers" | "body") {
|
||
return vi.fn<typeof fetch>((_url, init) => {
|
||
const signal = init?.signal;
|
||
if (!signal) {
|
||
throw new Error("readiness probe signal is required");
|
||
}
|
||
const stalled = (async (): Promise<never> => {
|
||
for (;;) {
|
||
signal.throwIfAborted();
|
||
await once(signal, "abort");
|
||
}
|
||
})();
|
||
if (phase === "headers") {
|
||
return stalled;
|
||
}
|
||
const response = new Response(null, { status: 503 });
|
||
response.json = () => stalled;
|
||
return Promise.resolve(response);
|
||
});
|
||
}
|
||
|
||
describe("openclaw test instance", () => {
|
||
it("reserves its idle port through refusal, CLI work, and stopped restarts", async ({
|
||
signal,
|
||
}) => {
|
||
const { instance, readAttempts } = await createFakeGateway("unrelated,cli,ready,ready");
|
||
const reserved = {
|
||
created: await isPortReserved(instance.port),
|
||
refused: false,
|
||
stopped: false,
|
||
};
|
||
await expect(startGatewayForPortLifecycle(instance, signal)).rejects.toThrow(
|
||
"unrelated startup failure",
|
||
);
|
||
expect(instance.readiness).toMatchObject([
|
||
{ outcome: "child-exit", child: { exitCode: 1, signalCode: null } },
|
||
]);
|
||
expect(instance.child).toBeUndefined();
|
||
reserved.refused = await isPortReserved(instance.port);
|
||
await expect(instance.cli(["0"])).resolves.toMatchObject({ code: 0, signal: null });
|
||
await startGatewayForPortLifecycle(instance, signal);
|
||
expect(instance.readiness).toMatchObject([
|
||
{ outcome: "ready", lastProbe: { phase: "complete", status: 200, ready: true } },
|
||
]);
|
||
await instance.stopGateway();
|
||
reserved.stopped = await isPortReserved(instance.port);
|
||
await startGatewayForPortLifecycle(instance, signal);
|
||
expect(instance.readiness).toMatchObject([
|
||
{ outcome: "ready", lastProbe: { phase: "complete", status: 200, ready: true } },
|
||
]);
|
||
const attempts = await readAttempts();
|
||
expect(attempts).toHaveLength(4);
|
||
expect(
|
||
attempts.filter((attempt) => attempt.argv[0] === "gateway").map((attempt) => attempt.port),
|
||
).toEqual([instance.port, instance.port, instance.port]);
|
||
const child = instance.child;
|
||
expect(child).toBeDefined();
|
||
const serverSpy = vi.spyOn(net, "createServer");
|
||
try {
|
||
await instance.cleanup();
|
||
expect(child?.exitCode !== null || child?.signalCode !== null).toBe(true);
|
||
expect(child?.stdout.closed).toBe(true);
|
||
expect(child?.stderr.closed).toBe(true);
|
||
await instance.cleanup();
|
||
await instance.stopGateway();
|
||
// Terminal calls must not leave a new idle reservation behind.
|
||
for (const result of serverSpy.mock.results) {
|
||
if (result.type === "return") {
|
||
expect(result.value.listening).toBe(false);
|
||
}
|
||
}
|
||
} finally {
|
||
serverSpy.mockRestore();
|
||
}
|
||
await expectPathMissing(instance.state.root);
|
||
expect(reserved).toEqual({ created: true, refused: true, stopped: true });
|
||
});
|
||
|
||
it("releases reservation probe connections before startup and terminal cleanup", async ({
|
||
signal,
|
||
}) => {
|
||
const { instance } = await createFakeGateway("ready");
|
||
const serverSpy = vi.spyOn(net, "createServer");
|
||
const probe = net.connect(instance.port, "127.0.0.1");
|
||
try {
|
||
await withinTest(once(probe, "close"), signal);
|
||
await startGatewayForPortLifecycle(instance, signal);
|
||
expect(instance.readiness).toMatchObject([
|
||
{ outcome: "ready", lastProbe: { phase: "complete", status: 200, ready: true } },
|
||
]);
|
||
await instance.stopGateway();
|
||
const reservation = serverSpy.mock.results.find(
|
||
(result) => result.type === "return" && result.value.listening,
|
||
)?.value;
|
||
expect(reservation?.listening).toBe(true);
|
||
const closed = vi.fn();
|
||
reservation?.once("close", closed);
|
||
await instance.cleanup();
|
||
expect(closed).toHaveBeenCalledOnce();
|
||
expect(reservation?.listening).toBe(false);
|
||
expect(reservation?.address()).toBeNull();
|
||
} finally {
|
||
probe.destroy();
|
||
serverSpy.mockRestore();
|
||
}
|
||
});
|
||
|
||
it.for(["held-unrelated", "late-unrelated"])(
|
||
"preserves the refusal when reacquiring the same port fails (%s)",
|
||
async (action, { signal }) => {
|
||
const control = await createGatewayControl();
|
||
const { instance } = await createFakeGateway(action, 1_000, 1_500, control, { signal });
|
||
const exited = createDeferred();
|
||
control.observers.onLaunch = () => {
|
||
instance.child?.once("exit", () => exited.resolve());
|
||
};
|
||
const competitor = net.createServer((socket) => socket.destroy());
|
||
const startup = startGatewayForPortLifecycle(instance, signal);
|
||
const outcome = startup.catch((error: unknown) => error);
|
||
try {
|
||
await Promise.race([control.reached, startup]);
|
||
await new Promise<void>((resolve, reject) => {
|
||
competitor.once("error", reject);
|
||
competitor.listen(instance.port, "127.0.0.1", resolve);
|
||
});
|
||
if (action === "late-unrelated") {
|
||
await Promise.race([exited.promise, outcome]);
|
||
expect(instance.child?.stderr.closed).toBe(false);
|
||
expect(instance.logs()).not.toContain("unrelated startup failure");
|
||
}
|
||
await control.release();
|
||
const error = await outcome;
|
||
expect(error).toBeInstanceOf(AggregateError);
|
||
expect((error as AggregateError).errors).toEqual([
|
||
expect.objectContaining({
|
||
message: expect.stringContaining("unrelated startup failure"),
|
||
}),
|
||
expect.objectContaining({ code: "EADDRINUSE" }),
|
||
]);
|
||
expect(instance.child).toBeUndefined();
|
||
await instance.cleanup();
|
||
await instance.stopGateway();
|
||
expect(competitor.listening).toBe(true);
|
||
await expectPathMissing(instance.state.root);
|
||
} finally {
|
||
control.unblock();
|
||
await outcome;
|
||
if (competitor.listening) {
|
||
await new Promise<void>((resolve, reject) => {
|
||
competitor.close((error) => (error ? reject(error) : resolve()));
|
||
});
|
||
}
|
||
}
|
||
},
|
||
);
|
||
|
||
it("leaves explicitly supplied ports owned by the caller", async () => {
|
||
const caller = net.createServer();
|
||
const closed = vi.fn();
|
||
caller.on("close", closed);
|
||
await new Promise<void>((resolve, reject) => {
|
||
caller.once("error", reject);
|
||
caller.listen(0, "127.0.0.1", resolve);
|
||
});
|
||
const address = caller.address();
|
||
if (!address || typeof address === "string") {
|
||
throw new Error("caller has no port");
|
||
}
|
||
try {
|
||
const instance = await createOpenClawTestInstance({
|
||
name: "caller-owned-port",
|
||
port: address.port,
|
||
});
|
||
fakeInstances.push({ instance });
|
||
await instance.stopGateway();
|
||
await instance.cleanup();
|
||
expect(caller.listening).toBe(true);
|
||
expect(closed).not.toHaveBeenCalled();
|
||
await expect(isPortReserved(address.port)).resolves.toBe(true);
|
||
} finally {
|
||
await new Promise<void>((resolve, reject) => {
|
||
caller.close((error) => (error ? reject(error) : resolve()));
|
||
});
|
||
}
|
||
expect(closed).toHaveBeenCalledOnce();
|
||
expect(caller.listening).toBe(false);
|
||
expect(caller.address()).toBeNull();
|
||
});
|
||
|
||
it.for(["complete", "overflow", "overflow-close"] as const)(
|
||
"owns complete CLI JSON and diagnostic tails (%s)",
|
||
async (mode, { signal }) => {
|
||
await withEnvAsync({ OPENAI_API_KEY: "ambient-provider-fixture" }, async () => {
|
||
const control = mode === "overflow-close" ? await createGatewayControl() : undefined;
|
||
const { instance, tracePath, readAttempts } = await createFakeGateway(
|
||
"cli-json",
|
||
1_000,
|
||
1_500,
|
||
control,
|
||
);
|
||
// The command must not merge a removed credential back from its parent.
|
||
delete instance.env.OPENAI_API_KEY;
|
||
const characters =
|
||
mode === "complete"
|
||
? 160 * 1024
|
||
: mode === "overflow-close"
|
||
? resolveMaxOutputBytes(undefined, "stdout")
|
||
: resolveMaxOutputBytes(undefined, "stdout") / 2;
|
||
const command = trackOperation(instance.cli([String(characters), mode]));
|
||
if (mode !== "complete") {
|
||
const outcomePromise = command.then(
|
||
(result) => ({
|
||
code: result.code,
|
||
stdoutBytes: Buffer.byteLength(result.stdout),
|
||
}),
|
||
(error: unknown) => error,
|
||
);
|
||
if (mode === "overflow-close") {
|
||
try {
|
||
// The fixture requests release only after its first stdout write completes.
|
||
await withinTest(
|
||
awaitGateBeforeSettlement(
|
||
control!.reached,
|
||
command,
|
||
`timeout waiting for ${tracePath}.overflow-ready`,
|
||
),
|
||
signal,
|
||
);
|
||
} finally {
|
||
control!.unblock();
|
||
await Promise.allSettled([command]);
|
||
}
|
||
}
|
||
const outcome = await outcomePromise;
|
||
if (!(outcome instanceof Error)) {
|
||
throw new Error(`Expected command output overflow failure: ${JSON.stringify(outcome)}`);
|
||
}
|
||
expect(outcome.message).toContain("command stdout exceeded capture limit");
|
||
if (mode === "overflow") {
|
||
expect(outcome.message).toContain('"last":"complete"');
|
||
} else {
|
||
expect(outcome.message).toContain(String.fromCodePoint(0x1f98a));
|
||
expect(outcome.message).toContain("trailing overflow output");
|
||
expect(outcome.message).not.toContain("\uFFFD");
|
||
}
|
||
expect(Buffer.byteLength(outcome.message)).toBeLessThan(600 * 1024);
|
||
} else {
|
||
const result = await command;
|
||
expect({ code: result.code, signal: result.signal }).toEqual({ code: 0, signal: null });
|
||
expect(JSON.parse(result.stdout)).toEqual({
|
||
first: "complete",
|
||
payload: "é".repeat(characters),
|
||
providerCredentialPresent: false,
|
||
last: "complete",
|
||
});
|
||
expect(result.stderr).toContain("[output truncated to last");
|
||
expect(result.stderr).toContain("recent cli diagnostic");
|
||
expect(result.stderr).not.toContain("discarded diagnostic");
|
||
expect(Buffer.byteLength(result.stderr)).toBeLessThan(300 * 1024);
|
||
}
|
||
const attempts = await readAttempts();
|
||
expect(attempts).toHaveLength(1);
|
||
expect(isProcessAlive(attempts[0]!.pid)).toBe(false);
|
||
});
|
||
},
|
||
);
|
||
|
||
it.for([
|
||
{ mode: "0", prepare: false },
|
||
{ mode: "7", prepare: false },
|
||
{ mode: "drain", prepare: false },
|
||
{ mode: "large", prepare: false },
|
||
{ mode: "wait", prepare: false },
|
||
{ mode: "0", prepare: true },
|
||
])("releases the CLI deadline after $mode (prepare=$prepare)", async (scenario, { signal }) => {
|
||
const { mode, prepare } = scenario;
|
||
const control = prepare || mode === "drain" ? await createGatewayControl() : undefined;
|
||
if (prepare) {
|
||
await control?.release();
|
||
}
|
||
const fixtureControl = control ? { url: control.url, holdPreparation: prepare } : undefined;
|
||
const { instance, readAttempts } = await createFakeGateway(
|
||
mode === "drain" ? "cli-drain" : "cli",
|
||
1_000,
|
||
1_500,
|
||
fixtureControl,
|
||
);
|
||
const scope = new AsyncLocalStorage<boolean>();
|
||
const timers = new Map<number, NodeJS.Timeout>();
|
||
const hook = createHook({
|
||
init(id, type, _trigger, resource) {
|
||
if (type === "Timeout" && scope.getStore()) {
|
||
// Node's Timeout async resource is the cancellable timer handle.
|
||
timers.set(id, resource as NodeJS.Timeout);
|
||
}
|
||
},
|
||
destroy(id) {
|
||
timers.delete(id);
|
||
},
|
||
});
|
||
hook.enable();
|
||
const observed = mode === "drain" ? await observeCliExit() : undefined;
|
||
try {
|
||
const timeoutMs = mode === "wait" ? 1_000 : 30_000;
|
||
const command = trackOperation(scope.run(true, () => instance.cli([mode], { timeoutMs })));
|
||
if (mode === "drain") {
|
||
await withinTest(
|
||
awaitGateBeforeSettlement(
|
||
control!.reached,
|
||
command,
|
||
"CLI stdout writer did not reach its gate",
|
||
),
|
||
signal,
|
||
);
|
||
const [attempt] = await readAttempts();
|
||
await withinTest(
|
||
awaitGateBeforeSettlement(
|
||
observed!.exited,
|
||
command,
|
||
`process still alive: ${attempt!.pid}`,
|
||
),
|
||
signal,
|
||
);
|
||
await control!.release();
|
||
}
|
||
if (mode === "wait") {
|
||
await expect(command).rejects.toThrow(`command timed out after ${timeoutMs}ms`);
|
||
} else {
|
||
const result = await command;
|
||
expect(result).toMatchObject({
|
||
code: mode === "drain" || mode === "large" ? 0 : Number(mode),
|
||
signal: null,
|
||
stderr: "cli diagnostic\n",
|
||
});
|
||
if (mode === "large") {
|
||
expect(result.stdout.startsWith("fake gateway attempt 1\n")).toBe(true);
|
||
expect(result.stdout.length).toBe("fake gateway attempt 1\n".length + 300 * 1024);
|
||
} else {
|
||
expect(result.stdout).toBe(
|
||
mode === "drain"
|
||
? "fake gateway attempt 1\ndrained cli output\n"
|
||
: "fake gateway attempt 1\n",
|
||
);
|
||
}
|
||
}
|
||
const attempts = await readAttempts();
|
||
expect(attempts).toHaveLength(1);
|
||
expect(isProcessAlive(attempts[0]!.pid)).toBe(false);
|
||
// Deliver Node's queued destroy hooks; elapsed wall time is not the oracle.
|
||
await new Promise<void>((resolve) => {
|
||
setImmediate(resolve);
|
||
});
|
||
expect(timers.size, "completed CLI invocation retained a deadline").toBe(0);
|
||
} finally {
|
||
control?.unblock();
|
||
observed?.restore();
|
||
hook.disable();
|
||
scope.disable();
|
||
// Retain the failing assertion while releasing only this invocation's timers.
|
||
for (const timer of timers.values()) {
|
||
clearTimeout(timer);
|
||
}
|
||
}
|
||
});
|
||
|
||
it("captures CLI stderr written after the process exits", async ({ signal }) => {
|
||
const control = await createGatewayControl();
|
||
const { instance, tracePath, readAttempts } = await createFakeGateway(
|
||
"late-refuse",
|
||
1_000,
|
||
1_500,
|
||
control,
|
||
);
|
||
const observed = await observeCliExit();
|
||
const command = trackOperation(instance.cli(["fixture"], { timeoutMs: 10_000 }));
|
||
let writerPid: number | undefined;
|
||
try {
|
||
await withinTest(
|
||
awaitGateBeforeSettlement(
|
||
control.reached,
|
||
command,
|
||
"CLI stderr fixture did not reach its release gate",
|
||
),
|
||
signal,
|
||
);
|
||
writerPid = Number(await fs.readFile(`${tracePath}.writer-pid`, "utf8"));
|
||
const [attempt] = await readAttempts();
|
||
await withinTest(
|
||
awaitGateBeforeSettlement(observed.exited, command, `process still alive: ${attempt!.pid}`),
|
||
signal,
|
||
);
|
||
// Release the inherited stderr writer only after the CLI leader has exited.
|
||
await control.release();
|
||
const result = await command;
|
||
expect(result).toEqual({
|
||
code: 1,
|
||
signal: null,
|
||
stdout: "fake gateway attempt 1\n",
|
||
stderr: `${MIGRATION_CONVERGENCE_REFUSAL} delayed fixture\n`,
|
||
});
|
||
} finally {
|
||
control.unblock();
|
||
observed.restore();
|
||
await Promise.allSettled([command]);
|
||
if (writerPid !== undefined) {
|
||
await waitForDead(writerPid, 5_000);
|
||
}
|
||
}
|
||
});
|
||
|
||
it.each([0, 1])(
|
||
"keeps claimed port offset %i unavailable while its child starts",
|
||
async (offset) => {
|
||
const control = await createGatewayControl();
|
||
const { instance, tracePath } = await createFakeGateway(
|
||
"held-unrelated",
|
||
10_000,
|
||
1_500,
|
||
control,
|
||
);
|
||
const starting = trackOperation(instance.startGateway());
|
||
await Promise.race([control.reached, starting]);
|
||
await drainFileLockStateForTest();
|
||
resetFileLockStateForTest();
|
||
const script = `
|
||
import { mock } from "node:test";
|
||
import { realpath } from "node:fs/promises";
|
||
import { tmpdir } from "node:os";
|
||
const ports = await import(${JSON.stringify(new URL("../../src/test-utils/ports.ts", import.meta.url).href)});
|
||
const attempted = [];
|
||
const candidates = ${JSON.stringify([instance.port + offset, instance.port + (1 - offset)])};
|
||
mock.module(${JSON.stringify(new URL("../../src/test-utils/ports.ts", import.meta.url).href)}, {
|
||
namedExports: { ...ports, getDeterministicFreePortBlock: async (options) => {
|
||
let candidate = candidates.shift() ?? await ports.getDeterministicFreePortBlock(options);
|
||
while (attempted.includes(candidate)) candidate = await ports.getDeterministicFreePortBlock(options);
|
||
attempted.push(candidate);
|
||
return candidate;
|
||
} },
|
||
});
|
||
const { createOpenClawTestInstance } = await import(${JSON.stringify(new URL("./openclaw-test-instance.ts", import.meta.url).href)});
|
||
const fixture = await createOpenClawTestInstance({ name: "port-claim-contender", cwd: ${JSON.stringify(path.dirname(tracePath))} });
|
||
try { console.log(JSON.stringify({ pid: process.pid, port: fixture.port, attempted, tempRoot: await realpath(tmpdir()) })); }
|
||
finally { await fixture.cleanup(); }
|
||
`;
|
||
const runContender = async (source: string) => {
|
||
const args = [
|
||
"--experimental-test-module-mocks",
|
||
"--import",
|
||
new URL("../../scripts/tsx.mjs", import.meta.url).href,
|
||
"--input-type=module",
|
||
"-e",
|
||
source,
|
||
];
|
||
const launcher = `
|
||
import { spawnOwnedVitestProcess } from ${JSON.stringify(new URL("../../scripts/lib/vitest-process.mts", import.meta.url).href)};
|
||
import { installVitestProcessGroupCleanup } from ${JSON.stringify(new URL("../../scripts/vitest-process-group.mts", import.meta.url).href)};
|
||
const { child, completion } = spawnOwnedVitestProcess({
|
||
command: process.execPath,
|
||
args: ${JSON.stringify(args)},
|
||
options: { cwd: process.cwd(), stdio: "inherit" },
|
||
homeMode: "tooling",
|
||
});
|
||
const cleanup = installVitestProcessGroupCleanup({ child, forceSignal: "SIGKILL" });
|
||
const result = await completion.finally(() => cleanup.teardown());
|
||
process.exitCode = result.code ?? 1;
|
||
`;
|
||
const result = await promisify(execFile)(
|
||
resolveTestNodeExecPath(),
|
||
[
|
||
"--import",
|
||
new URL("../../scripts/tsx.mjs", import.meta.url).href,
|
||
"--input-type=module",
|
||
"-e",
|
||
launcher,
|
||
],
|
||
{ cwd: process.cwd(), timeout: 20_000 },
|
||
);
|
||
return JSON.parse(result.stdout.trim());
|
||
};
|
||
const allocateContender = () => runContender(script);
|
||
const reserveInProcessPort = () =>
|
||
runContender(`
|
||
const { reserveGatewayTestListener } = await import(${JSON.stringify(new URL("../../src/gateway/test-helpers.listener.ts", import.meta.url).href)});
|
||
try {
|
||
const listener = await reserveGatewayTestListener(${instance.port + offset});
|
||
await listener.closeUnadopted();
|
||
console.log(JSON.stringify({ reserved: true }));
|
||
} catch (error) {
|
||
console.log(JSON.stringify({ reserved: false, code: error.code, causeCode: error.cause?.code }));
|
||
}
|
||
`);
|
||
const claimHandles: {
|
||
handle: Awaited<ReturnType<typeof fs.open>>;
|
||
lockPath: string;
|
||
dev: bigint;
|
||
ino: bigint;
|
||
}[] = [];
|
||
try {
|
||
const claims = createFileLockManager("openclaw.test-gateway-ports")
|
||
.heldEntries()
|
||
.filter((claim) =>
|
||
[instance.port, instance.port + 1].some(
|
||
(port) => path.basename(claim.normalizedTargetPath) === `openclaw-test-port-${port}`,
|
||
),
|
||
);
|
||
expect(claims).toHaveLength(2);
|
||
for (const claim of claims) {
|
||
const handle = await fs.open(claim.lockPath, "r");
|
||
const identity = await handle.stat({ bigint: true });
|
||
claimHandles.push({
|
||
handle,
|
||
lockPath: claim.lockPath,
|
||
dev: identity.dev,
|
||
ino: identity.ino,
|
||
});
|
||
}
|
||
expect(await reserveInProcessPort()).toEqual({
|
||
reserved: false,
|
||
code: "EADDRINUSE",
|
||
causeCode: FILE_LOCK_TIMEOUT_ERROR_CODE,
|
||
});
|
||
const contender = await allocateContender();
|
||
expect(contender.tempRoot).not.toBe(await fs.realpath(tmpdir()));
|
||
expect(contender.pid).not.toBe(process.pid);
|
||
expect([instance.port, instance.port + 1]).not.toContain(contender.port);
|
||
expect(contender.attempted.slice(0, 2)).toEqual([
|
||
instance.port + offset,
|
||
instance.port + (1 - offset),
|
||
]);
|
||
control.unblock();
|
||
await Promise.allSettled([starting]);
|
||
await instance.cleanup();
|
||
for (const claim of claimHandles) {
|
||
// Pin the old inode while checking release: another fixture may already
|
||
// own the same port and lock pathname, but cannot own this claim file.
|
||
let retained = false;
|
||
try {
|
||
const current = statSync(claim.lockPath, { bigint: true });
|
||
retained = current.dev === claim.dev && current.ino === claim.ino;
|
||
} catch (error) {
|
||
if (!hasErrnoCode(error, "ENOENT")) {
|
||
throw error;
|
||
}
|
||
}
|
||
expect(retained).toBe(false);
|
||
}
|
||
} finally {
|
||
control.unblock();
|
||
await Promise.allSettled([starting]);
|
||
await Promise.all(claimHandles.map((claim) => claim.handle.close()));
|
||
}
|
||
},
|
||
);
|
||
|
||
it("joins concurrent starts until the real readiness response arrives", async ({ signal }) => {
|
||
const control = await createGatewayControl();
|
||
const { instance } = await createFakeGateway("held-ready", 1_000, 1_500, control);
|
||
// Startup ordering must not depend on native bootstrap consuming the readiness budget.
|
||
signal.throwIfAborted();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(Date.now());
|
||
const restoreClock = () => clock.mockRestore();
|
||
signal.addEventListener("abort", restoreClock, { once: true });
|
||
let firstStart: Promise<void> | undefined;
|
||
let secondStart: Promise<void> | undefined;
|
||
try {
|
||
firstStart = trackOperation(instance.startGateway());
|
||
await Promise.race([control.reached, firstStart]);
|
||
|
||
let secondSettled = false;
|
||
let settledBeforeReady: boolean | undefined;
|
||
secondStart = trackOperation(
|
||
instance.startGateway().finally(() => {
|
||
secondSettled = true;
|
||
}),
|
||
);
|
||
// Observe at a real HTTP boundary before the child's held /readyz can reply.
|
||
control.observers.beforeRelease = () => {
|
||
settledBeforeReady = secondSettled;
|
||
};
|
||
await control.release();
|
||
await Promise.all([firstStart, secondStart]);
|
||
|
||
expect(settledBeforeReady).toBe(false);
|
||
expect(control.launches).toHaveLength(1);
|
||
expect(instance.child?.pid).toBe(
|
||
(process.platform === "win32" ? control.parents : control.launches)[0],
|
||
);
|
||
const response = await fetch(`http://127.0.0.1:${instance.port}/readyz`);
|
||
expect(await response.json()).toEqual({ ready: true });
|
||
} finally {
|
||
restoreClock();
|
||
signal.removeEventListener("abort", restoreClock);
|
||
control.unblock();
|
||
await Promise.allSettled([firstStart, secondStart]);
|
||
}
|
||
});
|
||
|
||
it("joins shared entrypoint preparation without launching the cancelled owner", async () => {
|
||
const controller = new AbortController();
|
||
const cancelled = new Error("instance owner cancelled during preparation");
|
||
const control = await createGatewayControl();
|
||
const { instance, tracePath, readAttempts } = await createFakeGateway(
|
||
"ready",
|
||
1_000,
|
||
1_500,
|
||
{
|
||
url: control.url,
|
||
holdPreparation: true,
|
||
},
|
||
{ signal: controller.signal },
|
||
);
|
||
const first = trackOperation(instance.startGateway());
|
||
const firstOutcome = first.catch((error: unknown) => error);
|
||
await Promise.race([control.reached, first]);
|
||
const borrower = await createOpenClawTestInstance({
|
||
name: "shared-entrypoint-borrower",
|
||
cwd: path.dirname(tracePath),
|
||
env: {
|
||
OPENCLAW_FAKE_GATEWAY_SEQUENCE: "ready",
|
||
OPENCLAW_FAKE_GATEWAY_TRACE: tracePath,
|
||
OPENCLAW_FAKE_GATEWAY_CONTROL: control.url,
|
||
},
|
||
});
|
||
fakeInstances.push({ instance: borrower });
|
||
const borrowedStart = trackOperation(borrower.startGateway());
|
||
const borrowedOutcome = borrowedStart.catch((error: unknown) => error);
|
||
controller.abort(cancelled);
|
||
const laterStart = trackOperation(instance.startGateway()).catch((error: unknown) => error);
|
||
// Observe cancellation only after the original shared preparation is joined.
|
||
await control.release();
|
||
expect(await firstOutcome).toBe(cancelled);
|
||
expect(await laterStart).toBeInstanceOf(Error);
|
||
expect(await borrowedOutcome).toBeUndefined();
|
||
expect(instance.child).toBeUndefined();
|
||
expect(control.launches).toHaveLength(1);
|
||
expect((await readAttempts()).map((attempt) => attempt.port)).toEqual([borrower.port]);
|
||
const response = await fetch(`http://127.0.0.1:${borrower.port}/readyz`);
|
||
expect(await response.json()).toEqual({ ready: true });
|
||
});
|
||
|
||
it("rolls back the real child when its owner aborts a pending readiness probe", async () => {
|
||
const controller = new AbortController();
|
||
const cancelled = new Error("instance owner cancelled during readiness");
|
||
const control = await createGatewayControl();
|
||
const { instance } = await createFakeGateway("held-ready", 1_000, 1_500, control, {
|
||
signal: controller.signal,
|
||
});
|
||
const startup = trackOperation(instance.startGateway());
|
||
const outcome = startup.catch((error: unknown) => error);
|
||
await Promise.race([control.reached, startup]);
|
||
const child = instance.child;
|
||
expect(child).toBeDefined();
|
||
controller.abort(cancelled);
|
||
expect(await outcome).toBe(cancelled);
|
||
expect(instance.child).toBeUndefined();
|
||
expect(child?.stdout.closed).toBe(true);
|
||
expect(child?.stderr.closed).toBe(true);
|
||
expect(inspectManagedProcessGroup(child!, { errorPolicy: "indeterminate" })).toBe("dead");
|
||
expect(isProcessAlive(child!.pid!)).toBe(false);
|
||
await expect(fs.stat(instance.state.root)).resolves.toBeDefined();
|
||
});
|
||
|
||
it("rolls back readiness that completes after owner cancellation", async () => {
|
||
const controller = new AbortController();
|
||
const cancelled = new Error("instance owner cancelled at readiness handoff");
|
||
const { instance } = await createFakeGateway("ready", 1_000, 1_500, undefined, {
|
||
signal: controller.signal,
|
||
});
|
||
const nativeFetch = globalThis.fetch;
|
||
let child: typeof instance.child;
|
||
const fetchSpy = vi.spyOn(globalThis, "fetch").mockImplementation(async (...args) => {
|
||
const response = await nativeFetch(...args);
|
||
if (args[0] === `http://127.0.0.1:${instance.port}/readyz`) {
|
||
const json = response.json.bind(response);
|
||
vi.spyOn(response, "json").mockImplementation(async () => {
|
||
const result: unknown = await json();
|
||
child = instance.child;
|
||
controller.abort(cancelled);
|
||
return result;
|
||
});
|
||
}
|
||
return response;
|
||
});
|
||
try {
|
||
const outcome = await trackOperation(instance.startGateway()).catch(
|
||
(error: unknown) => error,
|
||
);
|
||
expect(outcome).toBe(cancelled);
|
||
expect(child).toBeDefined();
|
||
expect(instance.child).toBeUndefined();
|
||
expect(child?.stdout.closed).toBe(true);
|
||
expect(child?.stderr.closed).toBe(true);
|
||
expect(isProcessAlive(child!.pid!)).toBe(false);
|
||
} finally {
|
||
fetchSpy.mockRestore();
|
||
}
|
||
});
|
||
|
||
it("joins a cancelled CLI command before releasing instance state", async ({ signal }) => {
|
||
const controller = new AbortController();
|
||
const control = await createGatewayControl();
|
||
const { instance, tracePath, readAttempts } = await createFakeGateway(
|
||
"cli",
|
||
1_000,
|
||
1_500,
|
||
control,
|
||
{
|
||
signal: controller.signal,
|
||
},
|
||
);
|
||
const command = trackOperation(instance.cli(["wait"], { timeoutMs: 1_000 }));
|
||
const outcome = command.catch((error: unknown) => error);
|
||
try {
|
||
await withinTest(
|
||
awaitGateBeforeSettlement(
|
||
control.reached,
|
||
command,
|
||
`Child exited before writing ${tracePath}.cli-ready`,
|
||
),
|
||
signal,
|
||
);
|
||
const [attempt] = await readAttempts();
|
||
expect(isProcessAlive(attempt!.pid)).toBe(true);
|
||
controller.abort(new Error("instance owner cancelled during CLI"));
|
||
expect(await outcome).toMatchObject({ message: expect.stringContaining("aborted") });
|
||
expect(isProcessAlive(attempt!.pid)).toBe(false);
|
||
await expect(fs.stat(instance.state.root)).resolves.toBeDefined();
|
||
await expect(instance.cli(["0"])).rejects.toThrow();
|
||
await instance.cleanup();
|
||
await expectPathMissing(instance.state.root);
|
||
} finally {
|
||
controller.abort();
|
||
control.unblock();
|
||
await outcome;
|
||
}
|
||
});
|
||
|
||
it("orders a new start after an intervening stop instead of joining the earlier start", async ({
|
||
signal,
|
||
}) => {
|
||
const control = await createGatewayControl();
|
||
const { instance } = await createFakeGateway("held-ready,ready", 1_000, 1_500, control);
|
||
// This case owns start/stop ordering, independently of native bootstrap speed.
|
||
signal.throwIfAborted();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(Date.now());
|
||
const restoreClock = () => clock.mockRestore();
|
||
signal.addEventListener("abort", restoreClock, { once: true });
|
||
let firstStart: Promise<void> | undefined;
|
||
let stopped: Promise<void> | undefined;
|
||
let secondStart: Promise<void> | undefined;
|
||
try {
|
||
firstStart = trackOperation(instance.startGateway());
|
||
await Promise.race([control.reached, firstStart]);
|
||
stopped = trackOperation(instance.stopGateway());
|
||
secondStart = trackOperation(instance.startGateway());
|
||
|
||
await control.release();
|
||
await Promise.allSettled([firstStart]);
|
||
await stopped;
|
||
await secondStart;
|
||
|
||
expect(control.launches).toHaveLength(2);
|
||
expect(control.launches[1]).not.toBe(control.launches[0]);
|
||
expect(isProcessAlive(control.launches[0]!)).toBe(false);
|
||
expect(instance.child?.pid).toBe(
|
||
(process.platform === "win32" ? control.parents : control.launches)[1],
|
||
);
|
||
const response = await fetch(`http://127.0.0.1:${instance.port}/readyz`);
|
||
expect(await response.json()).toEqual({ ready: true });
|
||
} finally {
|
||
restoreClock();
|
||
signal.removeEventListener("abort", restoreClock);
|
||
control.unblock();
|
||
await Promise.allSettled([firstStart, stopped, secondStart]);
|
||
}
|
||
});
|
||
|
||
it("starts a ready replacement after a real readiness deadline expires", async ({ signal }) => {
|
||
const control = await createGatewayControl();
|
||
const { instance } = await createFakeGateway("never-ready,ready", 1_000, 1_500, control);
|
||
let firstOwnerPid: number | undefined;
|
||
control.observers.onLaunch = () => {
|
||
firstOwnerPid = instance.child?.pid;
|
||
};
|
||
// Charge the unchanged real deadline only after the child proves it is running but unready.
|
||
signal.throwIfAborted();
|
||
const now = Date.now.bind(Date);
|
||
const fixtureTime = now();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(fixtureTime);
|
||
const restoreClock = () => clock.mockRestore();
|
||
signal.addEventListener("abort", restoreClock, { once: true });
|
||
let readinessStartedAt: number | undefined;
|
||
const nativeFetch = globalThis.fetch;
|
||
const fetchSpy = vi.spyOn(globalThis, "fetch").mockImplementation(async (...args) => {
|
||
const response = await nativeFetch(...args);
|
||
if (
|
||
args[0] === `http://127.0.0.1:${instance.port}/readyz` &&
|
||
readinessStartedAt === undefined
|
||
) {
|
||
const json = response.json.bind(response);
|
||
vi.spyOn(response, "json").mockImplementation(async () => {
|
||
const result: unknown = await json();
|
||
expect(response.status).toBe(200);
|
||
expect(result).toEqual({ ready: false });
|
||
const startedAt = now();
|
||
readinessStartedAt = startedAt;
|
||
clock.mockImplementation(() => fixtureTime + now() - startedAt);
|
||
return result;
|
||
});
|
||
}
|
||
return response;
|
||
});
|
||
let error: unknown;
|
||
try {
|
||
error = await trackOperation(instance.startGateway()).catch((failure: unknown) => failure);
|
||
} finally {
|
||
restoreClock();
|
||
signal.removeEventListener("abort", restoreClock);
|
||
fetchSpy.mockRestore();
|
||
}
|
||
expect(error).toBeInstanceOf(Error);
|
||
expect((error as Error).message).toContain("timeout waiting for gateway readiness");
|
||
expect(readinessStartedAt).toBeTypeOf("number");
|
||
const firstPid = control.launches[0];
|
||
expect(firstPid).toBeTypeOf("number");
|
||
expect(firstOwnerPid).toBeTypeOf("number");
|
||
const receipt = readReadinessReceipt(error);
|
||
expect(receipt).toMatchObject({
|
||
child: { pid: firstOwnerPid, exitCode: null, signalCode: null },
|
||
lastFailedResponse: { phase: "complete", status: 200, ready: false },
|
||
});
|
||
expect(receipt.elapsedMs).toBeGreaterThanOrEqual(1_000);
|
||
// Assert automatic failure cleanup before a replacement start can reap the owner.
|
||
expect(instance.child).toBeUndefined();
|
||
expect(isProcessAlive(firstPid as number)).toBe(false);
|
||
await expect(fs.stat(instance.state.root)).resolves.toBeDefined();
|
||
|
||
await startGatewayForPortLifecycle(instance, signal);
|
||
const response = await fetch(`http://127.0.0.1:${instance.port}/readyz`);
|
||
expect(await response.json()).toEqual({ ready: true });
|
||
expect(control.launches).toHaveLength(2);
|
||
expect(control.launches[1]).not.toBe(firstPid);
|
||
expect(instance.child?.pid).toBe(
|
||
(process.platform === "win32" ? control.parents : control.launches)[1],
|
||
);
|
||
expect(isProcessAlive(firstPid as number)).toBe(false);
|
||
await instance.stopGateway();
|
||
expect(instance.child).toBeUndefined();
|
||
expect(isProcessAlive(control.launches[1]!)).toBe(false);
|
||
await expect(fs.stat(instance.state.root)).resolves.toBeDefined();
|
||
await instance.cleanup();
|
||
await expectPathMissing(instance.state.root);
|
||
});
|
||
|
||
it.each(["stopGateway", "cleanup"] as const)(
|
||
"does not launch after %s settles during entrypoint preparation",
|
||
async (method) => {
|
||
const control = await createGatewayControl();
|
||
const { instance } = await createFakeGateway("ready", 1_000, 1_500, {
|
||
url: control.url,
|
||
holdPreparation: true,
|
||
});
|
||
const firstStart = trackOperation(instance.startGateway());
|
||
await Promise.race([control.reached, firstStart]);
|
||
await expect(isPortReserved(instance.port)).resolves.toBe(true);
|
||
let teardownSettled = false;
|
||
let launchedAfterTeardown = false;
|
||
control.observers.onLaunch = () => {
|
||
launchedAfterTeardown ||= teardownSettled;
|
||
};
|
||
const teardown = trackOperation(
|
||
instance[method]().finally(() => {
|
||
teardownSettled = true;
|
||
}),
|
||
);
|
||
|
||
// A valid owner may join startup or cancel this instance. Release preparation
|
||
// before joining teardown so either policy can complete without a deadlock.
|
||
await control.release();
|
||
const [, stopped] = await Promise.allSettled([firstStart, teardown]);
|
||
expect(stopped.status).toBe("fulfilled");
|
||
expect(launchedAfterTeardown).toBe(false);
|
||
expect(instance.child).toBeUndefined();
|
||
await instance.cleanup();
|
||
await expectPathMissing(instance.state.root);
|
||
for (const pid of control.launches) {
|
||
expect(isProcessAlive(pid)).toBe(false);
|
||
}
|
||
},
|
||
);
|
||
|
||
it("classifies only exact stderr convergence refusals with status 1", () => {
|
||
const classify = testing.isGatewayMigrationConvergenceRefusal;
|
||
expect(classify(1, null, `notice\n${MIGRATION_CONVERGENCE_REFUSAL} retry\n`)).toBe(true);
|
||
for (const candidate of [
|
||
[2, null, MIGRATION_CONVERGENCE_REFUSAL],
|
||
[1, "SIGTERM", MIGRATION_CONVERGENCE_REFUSAL],
|
||
[1, null, MIGRATION_CONVERGENCE_REFUSAL.slice(0, -1)],
|
||
[1, null, `prefix ${MIGRATION_CONVERGENCE_REFUSAL}`],
|
||
]) {
|
||
expect(classify(...(candidate as [number, NodeJS.Signals | null, string]))).toBe(false);
|
||
}
|
||
});
|
||
|
||
it.for([
|
||
"legacy-refuse",
|
||
"profile-legacy-refuse",
|
||
"late-legacy-refuse",
|
||
"startup-legacy-refuse",
|
||
"startup-late-legacy-refuse",
|
||
])("reports a typed %s only after the child's stderr closes", async (action) => {
|
||
const message = action.startsWith("startup-")
|
||
? LEGACY_STARTUP_FAILURE
|
||
: action === "profile-legacy-refuse"
|
||
? PROFILED_LEGACY_MIGRATION_REFUSAL
|
||
: LEGACY_MIGRATION_REFUSAL;
|
||
const control = action.includes("late-") ? await createGatewayControl() : undefined;
|
||
const { instance, readAttempts } = await createFakeGateway(
|
||
`${action},config-refuse`,
|
||
10_000,
|
||
1_500,
|
||
control,
|
||
);
|
||
if (action === "profile-legacy-refuse" || action.startsWith("startup-")) {
|
||
instance.env.OPENCLAW_PROFILE = "qa-fixture";
|
||
}
|
||
const exited = createDeferred();
|
||
if (control) {
|
||
control.observers.onLaunch = () => {
|
||
instance.child?.once("exit", () => exited.resolve());
|
||
};
|
||
}
|
||
const startup = trackOperation(instance.startGateway());
|
||
const outcome = startup.catch((error: unknown) => error);
|
||
if (control) {
|
||
await Promise.race([exited.promise, outcome]);
|
||
await Promise.race([control.reached, outcome]);
|
||
expect(instance.child?.stderr.closed).toBe(false);
|
||
expect(instance.logs()).not.toContain(message);
|
||
await control.release();
|
||
}
|
||
const error = await outcome;
|
||
expect(error).toMatchObject({
|
||
reason: "legacy-migration-required",
|
||
exitCode: 78,
|
||
signalCode: null,
|
||
legacyStorePath: LEGACY_STORE_PATH,
|
||
stderr: `${message}\n`,
|
||
cause: expect.any(Error),
|
||
});
|
||
expect(error).toBeInstanceOf(GatewayStartupRefusedError);
|
||
expect(instance.child).toBeUndefined();
|
||
expect(await readAttempts()).toHaveLength(1);
|
||
// Historical stderr must not classify the next child's unrelated EX_CONFIG.
|
||
const nextError = await instance.startGateway().catch((failure: unknown) => failure);
|
||
expect(nextError).toBeInstanceOf(Error);
|
||
expect(nextError).not.toBeInstanceOf(GatewayStartupRefusedError);
|
||
expect(await readAttempts()).toHaveLength(2);
|
||
});
|
||
|
||
it.for([
|
||
"legacy-stdout",
|
||
"legacy-status1",
|
||
"legacy-no-advice",
|
||
"profile-legacy-refuse",
|
||
"startup-no-advice",
|
||
"startup-warning",
|
||
])("does not classify %s as an explained legacy migration refusal", async (action) => {
|
||
const { instance, readAttempts } = await createFakeGateway(action);
|
||
if (action === "profile-legacy-refuse") {
|
||
instance.env.OPENCLAW_PROFILE = "different-fixture";
|
||
}
|
||
const error = await instance.startGateway().catch((failure: unknown) => failure);
|
||
expect(error).toBeInstanceOf(Error);
|
||
expect(error).not.toBeInstanceOf(GatewayStartupRefusedError);
|
||
expect(await readAttempts()).toHaveLength(1);
|
||
});
|
||
|
||
it("does not retry a drained migration refusal after owner cancellation", async () => {
|
||
const controller = new AbortController();
|
||
const control = await createGatewayControl();
|
||
const { instance, readAttempts } = await createFakeGateway(
|
||
"late-refuse,ready",
|
||
10_000,
|
||
1_500,
|
||
control,
|
||
{
|
||
signal: controller.signal,
|
||
},
|
||
);
|
||
const exited = createDeferred();
|
||
control.observers.onLaunch = () => {
|
||
instance.child?.once("exit", () => exited.resolve());
|
||
};
|
||
const startup = trackOperation(instance.startGateway());
|
||
const outcome = startup.catch((error: unknown) => error);
|
||
await Promise.race([exited.promise, startup]);
|
||
await Promise.race([control.reached, startup]);
|
||
expect(instance.child?.stderr.closed).toBe(false);
|
||
controller.abort(new Error("instance owner cancelled during refusal drain"));
|
||
await control.release();
|
||
const result = await outcome;
|
||
const attempts = await readAttempts();
|
||
expect(attempts).toHaveLength(1);
|
||
expect(result).toBeInstanceOf(Error);
|
||
expect(instance.logs()).not.toContain(RESTART_MARKER);
|
||
expect(instance.child).toBeUndefined();
|
||
expect(isProcessAlive(attempts[0]!.pid)).toBe(false);
|
||
});
|
||
|
||
it.for(["refuse", "late-refuse"])(
|
||
"restarts one %s refusal with identical launch state and owns the ready child",
|
||
async (refusalAction, { signal: testSignal }) => {
|
||
const control = refusalAction === "late-refuse" ? await createGatewayControl() : undefined;
|
||
const { instance, readAttempts } = await createFakeGateway(
|
||
`${refusalAction},ready`,
|
||
1_000,
|
||
1_500,
|
||
control,
|
||
);
|
||
// This case owns refusal/retry ordering, not deadline expiry. Keep native
|
||
// bootstrap and HTTP gates real without charging them to the policy clock.
|
||
testSignal.throwIfAborted();
|
||
const fixtureTime = Date.now();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(fixtureTime);
|
||
const restoreClock = () => clock.mockRestore();
|
||
testSignal.addEventListener("abort", restoreClock, { once: true });
|
||
const exited = createDeferred<{ code: number | null; signal: NodeJS.Signals | null }>();
|
||
if (control) {
|
||
control.observers.onLaunch = () => {
|
||
if (control.launches.length !== 1) {
|
||
return;
|
||
}
|
||
const leader = instance.child;
|
||
if (!leader) {
|
||
exited.reject(new Error("fixture launched without a process owner"));
|
||
return;
|
||
}
|
||
leader.once("exit", (code, signal) => exited.resolve({ code, signal }));
|
||
};
|
||
}
|
||
const startup = trackOperation(instance.startGateway());
|
||
try {
|
||
if (control) {
|
||
expect(await Promise.race([exited.promise, startup])).toEqual({ code: 1, signal: null });
|
||
await Promise.race([control.reached, startup]);
|
||
expect(instance.child?.stderr.closed).toBe(false);
|
||
expect(instance.logs()).not.toContain(MIGRATION_CONVERGENCE_REFUSAL);
|
||
// /launch installs the exit observer before the leader proceeds. Only
|
||
// release its waiting stderr writer after that native exit, never a timer.
|
||
await control.release();
|
||
}
|
||
await startup;
|
||
} finally {
|
||
restoreClock();
|
||
testSignal.removeEventListener("abort", restoreClock);
|
||
control?.unblock();
|
||
await Promise.allSettled([startup]);
|
||
}
|
||
const attempts = await readAttempts();
|
||
expect(attempts).toHaveLength(2);
|
||
expect(attempts[0]?.pid).not.toBe(attempts[1]?.pid);
|
||
expect({ ...attempts[0], pid: 0 }).toEqual({ ...attempts[1], pid: 0 });
|
||
expect(instance.logs()).toContain(MIGRATION_CONVERGENCE_REFUSAL);
|
||
expect(instance.logs()).toContain(RESTART_MARKER);
|
||
const readyPid = instance.child?.pid;
|
||
expect(readyPid).toBeTypeOf("number");
|
||
await instance.stopGateway();
|
||
expect(instance.child).toBeUndefined();
|
||
expect(isProcessAlive(readyPid as number)).toBe(false);
|
||
},
|
||
);
|
||
|
||
it.for(["near", "stdout", "status2", "signal", "unrelated"])(
|
||
"keeps %s convergence lookalikes terminal",
|
||
async (action, context) => {
|
||
// A Windows self-SIGTERM is status 1/null, already covered by the refusal case.
|
||
if (action === "signal" && process.platform === "win32") {
|
||
context.skip();
|
||
}
|
||
const { instance, readAttempts } = await createFakeGateway(`${action},ready`);
|
||
const startup = instance.startGateway();
|
||
try {
|
||
await expect(startup).rejects.toThrow("gateway exited before readiness");
|
||
} catch (error) {
|
||
console.error(
|
||
`Unexpected ${action} fake Gateway startup outcome`,
|
||
instance.logs(),
|
||
await startup.catch((startupError: unknown) => startupError),
|
||
);
|
||
throw error;
|
||
}
|
||
expect(await readAttempts()).toHaveLength(1);
|
||
expect(instance.logs()).not.toContain(RESTART_MARKER);
|
||
expect(instance.child).toBeUndefined();
|
||
await expect(fs.stat(instance.state.root)).resolves.toBeDefined();
|
||
await instance.stopGateway();
|
||
await expect(fs.stat(instance.state.root)).resolves.toBeDefined();
|
||
await instance.cleanup();
|
||
await expectPathMissing(instance.state.root);
|
||
},
|
||
);
|
||
|
||
it("preserves both refusals and never spawns a third gateway", async ({ signal }) => {
|
||
const control = await createGatewayControl();
|
||
const { instance, readAttempts } = await createFakeGateway(
|
||
"refuse,refuse,ready",
|
||
1_000,
|
||
1_500,
|
||
control,
|
||
);
|
||
const children: NonNullable<typeof instance.child>[] = [];
|
||
const previousOutputClosed: boolean[] = [];
|
||
control.observers.onLaunch = () => {
|
||
const previous = children.at(-1);
|
||
if (previous) {
|
||
previousOutputClosed.push(previous.stdout.closed && previous.stderr.closed);
|
||
}
|
||
if (instance.child) {
|
||
children.push(instance.child);
|
||
}
|
||
};
|
||
// Retry admission is the contract here; native process startup/exit must not
|
||
// spend the policy budget before both refusal facts reach the parent.
|
||
signal.throwIfAborted();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(Date.now());
|
||
const restoreClock = () => clock.mockRestore();
|
||
signal.addEventListener("abort", restoreClock, { once: true });
|
||
const servers = vi.spyOn(net, "createServer");
|
||
let reservation: net.Server | undefined;
|
||
const startup = trackOperation(instance.startGateway());
|
||
try {
|
||
await expect(startup).rejects.toThrow("gateway exited before readiness (code=1 signal=null)");
|
||
reservation = servers.mock.results.find(
|
||
(result) => result.type === "return" && result.value.listening,
|
||
)?.value;
|
||
expect(reservation).toBeDefined();
|
||
} finally {
|
||
restoreClock();
|
||
signal.removeEventListener("abort", restoreClock);
|
||
await Promise.allSettled([startup]);
|
||
servers.mockRestore();
|
||
}
|
||
const attempts = await readAttempts();
|
||
expect(attempts).toHaveLength(2);
|
||
expect(children).toHaveLength(2);
|
||
expect(previousOutputClosed).toEqual([true]);
|
||
expect(instance.readiness.map(({ outcome, child }) => ({ outcome, child }))).toEqual(
|
||
children.map(({ pid }) => ({
|
||
outcome: "child-exit",
|
||
child: { pid, exitCode: 1, signalCode: null },
|
||
})),
|
||
);
|
||
for (const child of children) {
|
||
expect(child.exitCode).toBe(1);
|
||
expect(child.signalCode).toBeNull();
|
||
expect(child.stdout.closed && child.stderr.closed).toBe(true);
|
||
}
|
||
expect(attempts.every(({ pid }) => !isProcessAlive(pid))).toBe(true);
|
||
expect(instance.child).toBeUndefined();
|
||
expect(instance.logs().split(MIGRATION_CONVERGENCE_REFUSAL)).toHaveLength(3);
|
||
expect(instance.logs().split(RESTART_MARKER)).toHaveLength(2);
|
||
await expect(isPortReserved(instance.port)).resolves.toBe(true);
|
||
await instance.cleanup();
|
||
// Once its claim is released, another worker may bind the numeric port.
|
||
expect(reservation?.listening).toBe(false);
|
||
expect(reservation?.address()).toBeNull();
|
||
});
|
||
|
||
it.runIf(process.platform !== "win32")(
|
||
"preserves an eligible refusal when its startup deadline expires during stdio drain",
|
||
async ({ signal: testSignal }) => {
|
||
const startupBudgetMs = 500;
|
||
const control = await createGatewayControl();
|
||
const { instance, tracePath, readAttempts } = await createFakeGateway(
|
||
"refusal-drain,ready",
|
||
startupBudgetMs,
|
||
1_500,
|
||
control,
|
||
);
|
||
// Start policy time only after the native leader has exited with stdio held.
|
||
// OS/module bootstrap must not turn this drain case into a readiness timeout.
|
||
const now = Date.now.bind(Date);
|
||
const fixtureTime = now();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(fixtureTime);
|
||
const restoreClock = () => clock.mockRestore();
|
||
testSignal.addEventListener("abort", restoreClock, { once: true });
|
||
const exited = createDeferred<{
|
||
code: number | null;
|
||
signal: NodeJS.Signals | null;
|
||
at: number;
|
||
}>();
|
||
control.observers.onLaunch = () => {
|
||
const leader = instance.child;
|
||
if (!leader) {
|
||
exited.reject(new Error("fixture launched without a process owner"));
|
||
return;
|
||
}
|
||
leader.once("exit", (code, signal) => exited.resolve({ code, signal, at: Date.now() }));
|
||
};
|
||
let startupSettled = false;
|
||
const startup = trackOperation(
|
||
instance.startGateway().finally(() => {
|
||
startupSettled = true;
|
||
}),
|
||
);
|
||
try {
|
||
const firstExit = await Promise.race([exited.promise, startup]);
|
||
expect(firstExit).toMatchObject({ code: 1, signal: null });
|
||
if (!firstExit) {
|
||
throw new Error("startup settled before the fixture leader exited");
|
||
}
|
||
const drainingPid = Number(await fs.readFile(`${tracePath}.draining-pid`, "utf8"));
|
||
expect(instance.child?.stderr.closed).toBe(false);
|
||
expect(isProcessAlive(drainingPid)).toBe(true);
|
||
const drainStartedAt = now();
|
||
clock.mockImplementation(() => fixtureTime + now() - drainStartedAt);
|
||
// The real drain now consumes the unchanged startup budget while cleanup
|
||
// retains its separate allowance and ownership of the inherited pipe.
|
||
const admissionExpiredAt = firstExit.at + startupBudgetMs;
|
||
while (Date.now() < admissionExpiredAt) {
|
||
await delay(admissionExpiredAt - Date.now());
|
||
}
|
||
expect(instance.child?.stderr.closed).toBe(false);
|
||
expect(isProcessAlive(drainingPid)).toBe(true);
|
||
expect(startupSettled).toBe(false);
|
||
|
||
await fs.writeFile(`${tracePath}.draining-release`, "");
|
||
await expect(startup).rejects.toThrow(
|
||
"gateway exited before readiness (code=1 signal=null)",
|
||
);
|
||
expect(instance.logs()).toContain(MIGRATION_CONVERGENCE_REFUSAL);
|
||
expect(instance.logs()).not.toContain(RESTART_MARKER);
|
||
expect(await readAttempts()).toHaveLength(1);
|
||
expect(instance.child).toBeUndefined();
|
||
await withinTest(waitForProcessExit(drainingPid, testSignal), testSignal);
|
||
expect(isProcessAlive(drainingPid)).toBe(false);
|
||
} finally {
|
||
restoreClock();
|
||
testSignal.removeEventListener("abort", restoreClock);
|
||
await fs.writeFile(`${tracePath}.draining-release`, "");
|
||
await Promise.allSettled([startup]);
|
||
}
|
||
},
|
||
);
|
||
|
||
it.each([
|
||
{ closePipes: true, groupError: "ESRCH", initiallyClosed: false, stopped: true },
|
||
{ closePipes: false, groupError: "ESRCH", initiallyClosed: false, stopped: false },
|
||
{ closePipes: true, groupError: "EPERM", initiallyClosed: false, stopped: false },
|
||
{ closePipes: true, groupError: "EPERM", initiallyClosed: true, stopped: false },
|
||
])(
|
||
"bounds TERM/KILL cleanup (close=$closePipes, group=$groupError, initiallyClosed=$initiallyClosed)",
|
||
async ({ closePipes, groupError, initiallyClosed, stopped: expectedStopped }) => {
|
||
const stdout = new PassThrough();
|
||
const stderr = new PassThrough();
|
||
const directKill = vi.fn(() => true);
|
||
const child = {
|
||
exitCode: 1,
|
||
kill: directKill,
|
||
pid: 12345,
|
||
signalCode: null,
|
||
stderr,
|
||
stdout,
|
||
} as unknown as Parameters<typeof testing.stopGatewayProcess>[0];
|
||
if (initiallyClosed) {
|
||
const closed = Promise.all([once(stdout, "close"), once(stderr, "close")]);
|
||
stdout.destroy();
|
||
stderr.destroy();
|
||
await closed;
|
||
}
|
||
const originalKill = process.kill.bind(process);
|
||
const signalTimes: number[] = [];
|
||
vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] });
|
||
const kill = vi.spyOn(process, "kill").mockImplementation((pid, signal) => {
|
||
if (pid !== -12345) {
|
||
return originalKill(pid, signal);
|
||
}
|
||
if (signal === 0) {
|
||
throw Object.assign(new Error("group lookup"), { code: groupError });
|
||
}
|
||
signalTimes.push(Date.now());
|
||
if (signal === "SIGKILL" && closePipes) {
|
||
stdout.destroy();
|
||
setTimeout(() => stderr.destroy(), 1);
|
||
}
|
||
return true;
|
||
});
|
||
try {
|
||
// Policy time is independent of OS scheduling; the real group proof below
|
||
// verifies native signal delivery and inherited-pipe closure separately.
|
||
const startedAt = Date.now();
|
||
const completion = testing
|
||
.stopGatewayProcess(child, startedAt + 80, 40, { platform: "linux" })
|
||
.then((stopped) => ({
|
||
stopped,
|
||
pipesClosed: stdout.closed && stderr.closed,
|
||
elapsedMs: Date.now() - startedAt,
|
||
}));
|
||
const [result] = await Promise.all([completion, vi.runAllTimersAsync()]);
|
||
expect(result.stopped).toBe(expectedStopped);
|
||
expect(result.pipesClosed).toBe(closePipes);
|
||
expect(result.elapsedMs).toBeLessThanOrEqual(80);
|
||
expect(kill.mock.calls.filter(([pid, signal]) => pid === -12345 && signal !== 0)).toEqual([
|
||
[-12345, "SIGTERM"],
|
||
[-12345, "SIGKILL"],
|
||
]);
|
||
const termGraceMs = signalTimes[1]! - signalTimes[0]!;
|
||
expect(termGraceMs).toBeGreaterThan(0);
|
||
expect(termGraceMs).toBeLessThanOrEqual(40);
|
||
expect(directKill).not.toHaveBeenCalled();
|
||
} finally {
|
||
stdout.destroy();
|
||
stderr.destroy();
|
||
vi.clearAllTimers();
|
||
kill.mockRestore();
|
||
vi.useRealTimers();
|
||
}
|
||
},
|
||
);
|
||
|
||
it.runIf(process.platform !== "win32").for(["inherit", "ignore"] as const)(
|
||
"certifies completion of a TERM-resistant group with %s stderr after leader exit",
|
||
async (stderrMode, { signal }) => {
|
||
const control = await createGatewayControl();
|
||
const { instance, tracePath } = await createFakeGateway(
|
||
stderrMode === "inherit" ? "resist-after-exit" : "resist-ignored-after-exit",
|
||
500,
|
||
40,
|
||
control,
|
||
);
|
||
const exerciseGroup = async () => {
|
||
// Preparation is outside policy time. The instance cases own slot/state
|
||
// assertions; this exercises the stopper's actual native signal primitive.
|
||
const leader = spawn(process.execPath, await instance.entrypoint(), {
|
||
cwd: path.dirname(tracePath),
|
||
env: instance.env,
|
||
detached: true,
|
||
stdio: ["ignore", "pipe", "pipe"],
|
||
});
|
||
const exited = once(leader, "exit");
|
||
const closed = once(leader, "close");
|
||
// A spawn error rejects both event promises; the joins below own it.
|
||
void closed.catch(() => undefined);
|
||
// A timeout must reap the raw group even before preparation reaches its gate.
|
||
const abortGroup = () => {
|
||
terminateManagedChild(leader, "SIGKILL");
|
||
};
|
||
signal.addEventListener("abort", abortGroup, { once: true });
|
||
if (signal.aborted) {
|
||
abortGroup();
|
||
}
|
||
let resistantPid: number | undefined;
|
||
let stopped: boolean | undefined;
|
||
const verifySignals = async () => {
|
||
await Promise.race([
|
||
control.reached,
|
||
exited.then(() => {
|
||
throw new Error("fixture exited before descendant readiness");
|
||
}),
|
||
]);
|
||
resistantPid = Number(await fs.readFile(`${tracePath}.resistant-pid`, "utf8"));
|
||
await control.release();
|
||
await exited;
|
||
if (stderrMode === "ignore") {
|
||
await closed;
|
||
}
|
||
expect(leader.exitCode).toBe(1);
|
||
expect(leader.signalCode).toBeNull();
|
||
expect(leader.stderr.closed).toBe(stderrMode === "ignore");
|
||
expect(isProcessAlive(resistantPid)).toBe(true);
|
||
|
||
terminateManagedChild(leader, "SIGTERM");
|
||
await expect
|
||
.poll(() => fs.readFile(`${tracePath}.signals`, "utf8"), { timeout: 500 })
|
||
.toBe("SIGTERM");
|
||
expect(isProcessAlive(resistantPid)).toBe(true);
|
||
expect(leader.stderr.closed).toBe(stderrMode === "ignore");
|
||
|
||
stopped = await testing.stopGatewayProcess(leader, Date.now() + 80, 40);
|
||
// An incomplete stop may exhaust its deadline before KILL. Only success
|
||
// certifies closure; the owned rescue below verifies extinction unconditionally.
|
||
if (stopped) {
|
||
expect(leader.stdout.closed).toBe(true);
|
||
expect(leader.stderr.closed).toBe(true);
|
||
expect(isProcessAlive(resistantPid)).toBe(false);
|
||
expect(inspectManagedProcessGroup(leader, { errorPolicy: "indeterminate" })).toBe(
|
||
"dead",
|
||
);
|
||
await withinTest(closed, signal);
|
||
expect(inspectManagedProcessGroup(leader, { errorPolicy: "indeterminate" })).toBe(
|
||
"dead",
|
||
);
|
||
}
|
||
};
|
||
const reapGroup = async () => {
|
||
terminateManagedChild(leader, "SIGKILL");
|
||
if (resistantPid && isProcessAlive(resistantPid)) {
|
||
try {
|
||
process.kill(resistantPid, "SIGKILL");
|
||
} catch (error) {
|
||
// Group termination can reap the descendant after the liveness probe.
|
||
if (!hasErrnoCode(error, "ESRCH")) {
|
||
throw error;
|
||
}
|
||
}
|
||
}
|
||
await closed;
|
||
if (resistantPid) {
|
||
await waitForDead(resistantPid, 500);
|
||
}
|
||
expect(inspectManagedProcessGroup(leader, { errorPolicy: "indeterminate" })).toBe("dead");
|
||
};
|
||
const [proof] = await Promise.allSettled([verifySignals()]);
|
||
const [cleanup] = await Promise.allSettled([reapGroup()]);
|
||
signal.removeEventListener("abort", abortGroup);
|
||
const registry = process.env.OPENCLAW_HELPER_PROOF_PID_REGISTRY;
|
||
if (registry) {
|
||
await fs.appendFile(
|
||
`${registry}.stops`,
|
||
`${JSON.stringify({ stderrMode, stopped, leaderPid: leader.pid, resistantPid })}\n`,
|
||
);
|
||
}
|
||
// Join before afterEach removes either root, and preserve both failures
|
||
// rather than letting last-resort cleanup hide the original regression.
|
||
if (cleanup.status === "rejected") {
|
||
const ownerIndex = fakeInstances.findIndex((owner) => owner.instance === instance);
|
||
expect(ownerIndex).toBeGreaterThanOrEqual(0);
|
||
fakeInstances.splice(ownerIndex, 1);
|
||
fakeRoots.splice(fakeRoots.indexOf(path.dirname(tracePath)), 1);
|
||
}
|
||
const failures = [proof, cleanup].flatMap((result) =>
|
||
result.status === "rejected" ? [result.reason] : [],
|
||
);
|
||
if (failures.length === 1) {
|
||
throw failures[0];
|
||
}
|
||
if (failures.length > 1) {
|
||
throw new AggregateError(failures, "native process proof and cleanup failed");
|
||
}
|
||
};
|
||
await trackOperation(exerciseGroup());
|
||
},
|
||
);
|
||
|
||
it.runIf(process.platform !== "win32")(
|
||
"reaps terminal children with inherited stdio before starting a new gateway",
|
||
async ({ signal: testSignal }) => {
|
||
const stopTimeoutMs = 100;
|
||
const control = await createGatewayControl();
|
||
const { instance, readAttempts, tracePath } = await createFakeGateway(
|
||
"terminal-drain,ready",
|
||
300,
|
||
stopTimeoutMs,
|
||
control,
|
||
);
|
||
// Bootstrap is not the teardown deadline. Start policy time at native exit,
|
||
// keeping the real inherited-pipe drain and failed restart budgets intact.
|
||
testSignal.throwIfAborted();
|
||
const now = Date.now.bind(Date);
|
||
const fixtureTime = now();
|
||
const clock = vi.spyOn(Date, "now").mockReturnValue(fixtureTime);
|
||
const restoreClock = () => clock.mockRestore();
|
||
testSignal.addEventListener("abort", restoreClock, { once: true });
|
||
const exited = createDeferred<{ code: number | null; signal: NodeJS.Signals | null }>();
|
||
control.observers.onLaunch = () => {
|
||
if (control.launches.length !== 1) {
|
||
return;
|
||
}
|
||
const leader = instance.child;
|
||
if (!leader) {
|
||
exited.reject(new Error("fixture launched without a process owner"));
|
||
return;
|
||
}
|
||
leader.once("exit", (code, signal) => exited.resolve({ code, signal }));
|
||
};
|
||
const startup = trackOperation(instance.startGateway());
|
||
try {
|
||
expect(await Promise.race([exited.promise, startup])).toEqual({ code: 7, signal: null });
|
||
const drainStartedAt = now();
|
||
clock.mockImplementation(() => fixtureTime + now() - drainStartedAt);
|
||
const startupError = await startup.catch((error: unknown) => error);
|
||
expect(startupError).toBeInstanceOf(Error);
|
||
expect((startupError as Error).message).toContain(
|
||
"gateway exited before readiness (code=7 signal=null)",
|
||
);
|
||
expect((startupError as Error).message).toContain("terminal startup failure");
|
||
const firstChild = instance.child;
|
||
expect(firstChild?.exitCode).toBe(7);
|
||
expect(firstChild?.stderr.closed).toBe(false);
|
||
if (!firstChild) {
|
||
throw new Error("terminal fixture lost its process owner");
|
||
}
|
||
const firstAttempt = (await readAttempts())[0];
|
||
const drainingPid = Number(await fs.readFile(`${tracePath}.draining-pid`, "utf8"));
|
||
expect(isProcessAlive(drainingPid)).toBe(true);
|
||
|
||
// Keep the inherited pipe held through a complete failed restart: eventual
|
||
// replacement after release alone cannot prove that stale ownership blocked it.
|
||
await expect(trackOperation(instance.startGateway())).rejects.toThrow(
|
||
new Error(
|
||
`gateway process cleanup could not verify termination and output closure\n${instance.logs()}`,
|
||
),
|
||
);
|
||
expect(instance.child).toBe(firstChild);
|
||
expect(firstChild.stderr.closed).toBe(false);
|
||
expect(isProcessAlive(drainingPid)).toBe(true);
|
||
expect(await readAttempts()).toHaveLength(1);
|
||
await expect(fs.stat(instance.state.root)).resolves.toBeDefined();
|
||
// A free socket is not evidence that the retained process owner has closed.
|
||
await expect(isPortReserved(instance.port)).resolves.toBe(false);
|
||
|
||
// Register before release so the native close cannot be missed.
|
||
const closed = trackOperation(once(firstChild, "close"));
|
||
await fs.writeFile(`${tracePath}.draining-release`, "");
|
||
await closed;
|
||
expect(firstChild.stdout.closed).toBe(true);
|
||
expect(firstChild.stderr.closed).toBe(true);
|
||
clock.mockReturnValue(Date.now());
|
||
await trackOperation(instance.startGateway());
|
||
restoreClock();
|
||
|
||
const attempts = await readAttempts();
|
||
expect(attempts).toHaveLength(2);
|
||
expect(attempts[1]?.pid).not.toBe(firstAttempt?.pid);
|
||
expect(instance.child?.pid).toBe(attempts[1]?.pid);
|
||
await instance.stopGateway();
|
||
expect(instance.child).toBeUndefined();
|
||
expect(isProcessAlive(attempts[1]?.pid as number)).toBe(false);
|
||
await withinTest(waitForProcessExit(drainingPid, testSignal), testSignal);
|
||
expect(isProcessAlive(drainingPid)).toBe(false);
|
||
} finally {
|
||
restoreClock();
|
||
testSignal.removeEventListener("abort", restoreClock);
|
||
await fs.writeFile(`${tracePath}.draining-release`, "");
|
||
await Promise.allSettled([startup]);
|
||
}
|
||
},
|
||
);
|
||
|
||
it.each([true, false])(
|
||
"does not target an exited Windows gateway (inherited pipes=%s)",
|
||
async (inheritedPipes) => {
|
||
const stdout = new PassThrough();
|
||
const stderr = new PassThrough();
|
||
const kill = vi.fn(() => true);
|
||
// SAFETY: The stub supplies every process and pipe member consumed by the stopper.
|
||
const child = {
|
||
exitCode: 1,
|
||
kill,
|
||
pid: 12345,
|
||
signalCode: null,
|
||
stderr,
|
||
stdout,
|
||
} as unknown as Parameters<typeof testing.stopGatewayProcess>[0];
|
||
const closePipes = () => {
|
||
stdout.destroy();
|
||
stderr.destroy();
|
||
};
|
||
const runTaskkill = vi.fn(() => {
|
||
closePipes();
|
||
return { status: 0 };
|
||
});
|
||
if (!inheritedPipes) {
|
||
const closed = Promise.all([once(stdout, "close"), once(stderr, "close")]);
|
||
closePipes();
|
||
await closed;
|
||
}
|
||
|
||
try {
|
||
await expect(
|
||
testing.stopGatewayProcess(child, Date.now() + 500, 250, {
|
||
forceWindowsTree: true,
|
||
platform: "win32",
|
||
runTaskkill,
|
||
}),
|
||
).rejects.toMatchObject({ processTreeState: "indeterminate" });
|
||
expect(runTaskkill).not.toHaveBeenCalled();
|
||
expect(kill).not.toHaveBeenCalled();
|
||
expect(stdout.closed).toBe(!inheritedPipes);
|
||
expect(stderr.closed).toBe(!inheritedPipes);
|
||
} finally {
|
||
const closed = Promise.all(
|
||
[stdout, stderr].map((pipe) => (pipe.closed ? Promise.resolve() : once(pipe, "close"))),
|
||
);
|
||
closePipes();
|
||
await closed;
|
||
}
|
||
},
|
||
);
|
||
|
||
it.each([
|
||
{ label: "joined closure", taskkillStatus: 0, closePipes: true, stopped: true },
|
||
{ label: "held pipe", taskkillStatus: 0, closePipes: false, stopped: false },
|
||
{ label: "unverified tree", taskkillStatus: 1, closePipes: true, stopped: false },
|
||
{ label: "unverified exited leader", taskkillStatus: 1, closePipes: true, stopped: false },
|
||
{ label: "taskkill timeout", taskkillStatus: 1, closePipes: true, stopped: false },
|
||
{ label: "taskkill exception", taskkillStatus: 1, closePipes: true, stopped: false },
|
||
])("observes Windows $label after blocking termination", async (scenario) => {
|
||
const stdout = new PassThrough();
|
||
const stderr = new PassThrough();
|
||
const stopLog = testing.createBoundedStringLog();
|
||
const exitedLeader = scenario.label === "unverified exited leader";
|
||
const heldPipe = scenario.label === "held pipe";
|
||
const processState = createGatewayProcessState({ exitCode: exitedLeader ? 7 : null });
|
||
// SAFETY: The stub supplies every process and pipe member consumed by the stopper.
|
||
const child = Object.assign(processState, {
|
||
pid: process.pid,
|
||
kill: vi.fn(() => true),
|
||
stdout,
|
||
stderr,
|
||
}) as unknown as Parameters<typeof testing.stopGatewayProcess>[0];
|
||
const observed = createDeferred();
|
||
const now = Date.now.bind(Date);
|
||
let offset = 0;
|
||
let scheduled = false;
|
||
const clock = vi.spyOn(Date, "now").mockImplementation(() => now() + offset);
|
||
const runTaskkill = vi.fn(() => {
|
||
// A synchronous taskkill consumes wall time before Node can deliver exit/close.
|
||
offset += 1_000;
|
||
if (!scheduled) {
|
||
scheduled = true;
|
||
setImmediate(() => {
|
||
processState.exitCode = heldPipe ? null : 0;
|
||
processState.signalCode = heldPipe ? "SIGTERM" : null;
|
||
stdout.destroy();
|
||
if (scenario.closePipes) {
|
||
stderr.destroy();
|
||
}
|
||
observed.resolve();
|
||
});
|
||
}
|
||
if (scenario.label === "taskkill exception") {
|
||
throw Object.assign(new Error("private exception text must not be logged"), {
|
||
code: "EACCES",
|
||
});
|
||
}
|
||
if (scenario.label === "taskkill timeout" && runTaskkill.mock.calls.length === 2) {
|
||
return {
|
||
status: null,
|
||
signal: "SIGKILL" as const,
|
||
error: Object.assign(new Error("private timeout text must not be logged"), {
|
||
code: "ETIMEDOUT",
|
||
}),
|
||
};
|
||
}
|
||
return { status: scenario.taskkillStatus, signal: null };
|
||
});
|
||
try {
|
||
const stop = testing.stopGatewayProcess(
|
||
child,
|
||
Date.now() + 500,
|
||
250,
|
||
{ platform: "win32", runTaskkill },
|
||
stopLog,
|
||
);
|
||
let stopped = false;
|
||
if (!scenario.stopped) {
|
||
await expect(stop).rejects.toBeInstanceOf(Error);
|
||
} else {
|
||
stopped = await stop;
|
||
expect(stopped).toBe(scenario.stopped);
|
||
}
|
||
const threw = scenario.label === "taskkill exception";
|
||
expect(runTaskkill).toHaveBeenCalledTimes(
|
||
exitedLeader ? 0 : scenario.taskkillStatus === 0 || threw ? 1 : 2,
|
||
);
|
||
if (!scenario.closePipes) {
|
||
expect(stderr.closed).toBe(false);
|
||
}
|
||
if (stopped) {
|
||
expect(child.exitCode).toBe(0);
|
||
expect(stdout.closed && stderr.closed).toBe(true);
|
||
}
|
||
expect(stopLog).toHaveLength(stopped ? 0 : 1);
|
||
if (!stopped) {
|
||
const prefix = "[openclaw-test-instance] Windows shutdown ";
|
||
expect(stopLog[0]?.startsWith(prefix)).toBe(true);
|
||
const diagnostic: { taskkill: Array<{ elapsedMs: number }> } = JSON.parse(
|
||
stopLog[0]!.slice(prefix.length),
|
||
);
|
||
const attempt = { force: false, elapsedMs: expect.any(Number) };
|
||
const taskkill: Record<string, unknown>[] = exitedLeader
|
||
? []
|
||
: threw
|
||
? [{ ...attempt, threw: true, errorCode: "EACCES" }]
|
||
: [{ ...attempt, status: scenario.taskkillStatus, signal: null }];
|
||
if (!exitedLeader && !heldPipe && !threw) {
|
||
taskkill.push({
|
||
...attempt,
|
||
force: true,
|
||
...(scenario.label === "taskkill timeout"
|
||
? { status: null, signal: "SIGKILL", errorCode: "ETIMEDOUT" }
|
||
: { status: 1, signal: null }),
|
||
});
|
||
}
|
||
expect(diagnostic).toEqual({
|
||
reason: "exception",
|
||
pid: process.pid,
|
||
exitCode: exitedLeader ? 7 : null,
|
||
signalCode: heldPipe ? "SIGTERM" : null,
|
||
stdoutClosed: heldPipe,
|
||
stderrClosed: false,
|
||
elapsedMs: expect.any(Number),
|
||
taskkill,
|
||
errorCode: threw ? "EACCES" : "EPROCESSGROUP_CLEANUP_FAILED",
|
||
});
|
||
for (const entry of diagnostic.taskkill) {
|
||
expect(entry.elapsedMs).toBeGreaterThanOrEqual(1_000);
|
||
}
|
||
expect(Buffer.byteLength(stopLog[0]!)).toBeLessThan(1_024);
|
||
}
|
||
} finally {
|
||
if (scheduled) {
|
||
await observed.promise;
|
||
}
|
||
const closed = Promise.all(
|
||
[stdout, stderr].map((pipe) => (pipe.closed ? Promise.resolve() : once(pipe, "close"))),
|
||
);
|
||
stdout.destroy();
|
||
stderr.destroy();
|
||
await closed.finally(() => clock.mockRestore());
|
||
}
|
||
});
|
||
|
||
it("keeps only bounded child output tails in helper logs", () => {
|
||
const stdout = testing.createBoundedStringLog(32);
|
||
const stderr = testing.createBoundedStringLog(32);
|
||
|
||
testing.appendLogChunk(stdout, `old stdout ${"x".repeat(64)}\n`);
|
||
testing.appendLogChunk(stdout, "recent stdout\n");
|
||
testing.appendLogChunk(stderr, `old stderr ${"y".repeat(64)}\n`);
|
||
testing.appendLogChunk(stderr, "recent stderr\n");
|
||
|
||
const logs = testing.formatLogs(stdout, stderr);
|
||
expect(logs).toContain("[output truncated to last 32 bytes]");
|
||
expect(logs).toContain("recent stdout");
|
||
expect(logs).toContain("recent stderr");
|
||
expect(logs).not.toContain("old stdout");
|
||
expect(logs).not.toContain("old stderr");
|
||
|
||
const exact = testing.createBoundedStringLog(32);
|
||
testing.appendLogChunk(exact, "x".repeat(32));
|
||
expect(testing.formatLogs(exact, [])).not.toContain("output truncated");
|
||
});
|
||
|
||
describe("UTF-8 log trimming", () => {
|
||
let exerciseTrimming: () => Promise<void>;
|
||
let stopTrimming: (() => Promise<void>) | undefined;
|
||
|
||
afterEach(async () => {
|
||
await stopTrimming?.();
|
||
});
|
||
|
||
// Use the fixture-hook budget for TS bootstrap before exercising trimming.
|
||
beforeEach(async ({ signal }) => {
|
||
const cases = [
|
||
{ chunks: ["€a", "b"], limit: 4, expected: "ab" },
|
||
{ chunks: ["old", "recent"], limit: 8, expected: "ldrecent" },
|
||
{ chunks: ["€abc"], limit: 4, expected: "abc" },
|
||
{ chunks: ["😀a", "b"], limit: 5, expected: "ab" },
|
||
{ chunks: ["😀a"], limit: 1, expected: "a" },
|
||
{ chunks: ["😀"], limit: 1, expected: "" },
|
||
{ chunks: ["😀"], limit: 3, expected: "" },
|
||
{ chunks: ["€"], limit: 2, expected: "" },
|
||
{ chunks: ["a", "€"], limit: 3, expected: "€" },
|
||
];
|
||
const control = await createGatewayControl();
|
||
// A synchronous regression must be killed and joined outside the Vitest event loop.
|
||
const script = `
|
||
import assert from "node:assert/strict";
|
||
import { testing } from ${JSON.stringify(new URL("./openclaw-test-instance.ts", import.meta.url).href)};
|
||
${CONTROL_REQUEST_SOURCE}
|
||
process.stderr.write("loaded actual log helper; waiting to start UTF-8 cases\\n");
|
||
await waitForControl(${JSON.stringify(`${control.url}/wait`)});
|
||
for (const { chunks, limit, expected } of JSON.parse(process.argv[1])) {
|
||
const log = testing.createBoundedStringLog(limit);
|
||
for (const chunk of chunks) {
|
||
testing.appendLogChunk(log, chunk);
|
||
assert.ok(Buffer.byteLength(log.join("")) <= limit);
|
||
}
|
||
assert.equal(log.join(""), expected);
|
||
assert.ok(!log.join("").includes("<22>"));
|
||
assert.match(testing.formatLogs(log, []), /output truncated to last/);
|
||
}
|
||
process.stdout.write("UTF-8 cases completed");
|
||
`;
|
||
const completed = promisify(execFile)(
|
||
process.execPath,
|
||
["--import", "tsx", "--input-type=module", "--eval", script, JSON.stringify(cases)],
|
||
{ signal, encoding: "utf8" },
|
||
);
|
||
const closed = new Promise<void>((resolve) => {
|
||
completed.child.once("close", () => resolve());
|
||
});
|
||
stopTrimming = async () => {
|
||
completed.child.kill("SIGKILL");
|
||
await closed;
|
||
};
|
||
exerciseTrimming = async () => {
|
||
await control.release();
|
||
const { stdout } = await completed;
|
||
expect(stdout).toBe("UTF-8 cases completed");
|
||
};
|
||
await trackOperation(
|
||
Promise.race([
|
||
control.reached,
|
||
completed.then(() => {
|
||
throw new Error("log helper child exited before reaching the trimming gate");
|
||
}),
|
||
]),
|
||
);
|
||
});
|
||
|
||
it("terminates UTF-8 log trimming within the byte cap", { timeout: 15_000 }, async () => {
|
||
await trackOperation(exerciseTrimming());
|
||
});
|
||
});
|
||
|
||
it("fails startup waits immediately after signaled gateway exits", async () => {
|
||
const error = await testing
|
||
.waitForGatewayReady(
|
||
createGatewayProcessState({ pid: 12345, signalCode: "SIGTERM" }),
|
||
[],
|
||
[],
|
||
1,
|
||
10_000,
|
||
)
|
||
.catch((failure: unknown) => failure);
|
||
expect(error).toBeInstanceOf(Error);
|
||
expect((error as Error).message).toContain("gateway exited before readiness");
|
||
expect(readReadinessReceipt(error)).toEqual({
|
||
attempts: 0,
|
||
elapsedMs: expect.any(Number),
|
||
lastProbe: null,
|
||
lastFailedResponse: null,
|
||
child: { pid: 12345, exitCode: null, signalCode: "SIGTERM" },
|
||
});
|
||
});
|
||
|
||
it("waits until the gateway readiness probe reports ready", async () => {
|
||
const fetchImpl = vi
|
||
.fn<typeof fetch>()
|
||
.mockResolvedValueOnce(
|
||
new Response('{"ready":false,"failing":["startup-sidecars"]}', { status: 503 }),
|
||
)
|
||
.mockResolvedValueOnce(new Response('{"ready":true,"failing":[]}', { status: 200 }));
|
||
|
||
const record = vi.fn<(diagnostic: GatewayReadinessDiagnostic) => void>();
|
||
await expect(
|
||
testing.waitForGatewayReady(
|
||
createGatewayProcessState(),
|
||
[],
|
||
[],
|
||
12345,
|
||
1_000,
|
||
fetchImpl,
|
||
undefined,
|
||
record,
|
||
),
|
||
).resolves.toBeUndefined();
|
||
expect(record).toHaveBeenCalledWith(
|
||
expect.objectContaining({
|
||
outcome: "ready",
|
||
lastProbe: expect.objectContaining({ attempt: 2, status: 200, ready: true }),
|
||
lastFailedResponse: expect.objectContaining({ attempt: 1, status: 503, ready: false }),
|
||
}),
|
||
);
|
||
|
||
expect(fetchImpl).toHaveBeenCalledTimes(2);
|
||
expect(fetchImpl.mock.calls[0]?.[0]).toBe("http://127.0.0.1:12345/readyz");
|
||
});
|
||
|
||
it("bounds not-ready diagnostics without exposing response details", async () => {
|
||
const privateDetail = `private-channel-${"x".repeat(2_048)}`;
|
||
const fetchImpl = vi.fn<typeof fetch>(async () => {
|
||
return new Response(
|
||
JSON.stringify({
|
||
ready: false,
|
||
failing: [
|
||
"startup-sidecars",
|
||
"gateway-draining",
|
||
"state-database",
|
||
"internal",
|
||
privateDetail,
|
||
null,
|
||
{ privateDetail },
|
||
false,
|
||
...Array<string>(12).fill(privateDetail),
|
||
],
|
||
config: { privateDetail },
|
||
error: privateDetail,
|
||
}),
|
||
{ status: 503, headers: { "x-private-detail": privateDetail } },
|
||
);
|
||
});
|
||
const error = await testing
|
||
.waitForGatewayReady(createGatewayProcessState(), [], [], 12345, 100, fetchImpl)
|
||
.catch((failure: unknown) => failure);
|
||
|
||
expect(error).toBeInstanceOf(Error);
|
||
expect((error as Error).message).toContain("timeout waiting for gateway readiness");
|
||
expect((error as Error).message).not.toContain(privateDetail);
|
||
expect((error as Error).message).not.toContain("x-private-detail");
|
||
expect((error as Error).message).not.toContain("config");
|
||
expect(readReadinessReceipt(error)).toEqual({
|
||
attempts: fetchImpl.mock.calls.length,
|
||
elapsedMs: expect.any(Number),
|
||
lastProbe: {
|
||
attempt: fetchImpl.mock.calls.length,
|
||
phase: "complete",
|
||
elapsedMs: expect.any(Number),
|
||
status: 503,
|
||
ready: false,
|
||
failing: [
|
||
"startup-sidecars",
|
||
"gateway-draining",
|
||
"state-database",
|
||
"internal",
|
||
"other",
|
||
"other",
|
||
"other",
|
||
"other",
|
||
],
|
||
omittedFailing: 12,
|
||
},
|
||
lastFailedResponse: {
|
||
attempt: fetchImpl.mock.calls.length,
|
||
phase: "complete",
|
||
elapsedMs: expect.any(Number),
|
||
status: 503,
|
||
ready: false,
|
||
},
|
||
child: { pid: null, exitCode: null, signalCode: null },
|
||
});
|
||
});
|
||
|
||
it.each([
|
||
{ status: 502, body: "proxy refusal", phase: "body", error: "invalid-json" },
|
||
{ status: 503, body: '{"ready":false}', phase: "complete", ready: false },
|
||
{ status: 200, body: '{"ready":false}', phase: "complete", ready: false },
|
||
{ status: 503, body: '{"ready":true}', phase: "complete", ready: true },
|
||
{ status: 200, body: "{}", phase: "complete" },
|
||
])("retains the last failed response ($status/$body) before a final timeout", async (failed) => {
|
||
vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] });
|
||
const privateDetail = "synthetic-proxy-credential";
|
||
const processState = createGatewayProcessState();
|
||
const fetchImpl = createStalledReadinessFetch("headers")
|
||
.mockResolvedValueOnce(new Response('{"ready":false}', { status: 429 }))
|
||
.mockResolvedValueOnce(
|
||
new Response(failed.body, {
|
||
status: failed.status,
|
||
headers: { "x-private": privateDetail },
|
||
}),
|
||
)
|
||
.mockRejectedValueOnce(new Error(privateDetail));
|
||
const record = vi.fn<(diagnostic: GatewayReadinessDiagnostic) => void>();
|
||
const completion = testing
|
||
.waitForGatewayReady(processState, [], [], 12345, 39, fetchImpl, undefined, record)
|
||
.catch((failure: unknown) => failure);
|
||
try {
|
||
await vi.advanceTimersByTimeAsync(39);
|
||
const failure = await completion;
|
||
const receipt = readReadinessReceipt(failure);
|
||
expect(receipt).toMatchObject({
|
||
attempts: 4,
|
||
elapsedMs: 39,
|
||
lastProbe: { attempt: 4, phase: "headers", error: "timeout", elapsedMs: 9 },
|
||
lastFailedResponse: {
|
||
attempt: 2,
|
||
phase: failed.phase,
|
||
elapsedMs: 0,
|
||
status: failed.status,
|
||
...(failed.ready === undefined ? {} : { ready: failed.ready }),
|
||
...(failed.error === undefined ? {} : { error: failed.error }),
|
||
},
|
||
});
|
||
expect(record).toHaveBeenCalledWith(
|
||
expect.objectContaining({
|
||
outcome: "timeout",
|
||
lastProbe: receipt.lastProbe,
|
||
lastFailedResponse: receipt.lastFailedResponse,
|
||
}),
|
||
);
|
||
const diagnostic: GatewayReadinessDiagnostic = record.mock.calls[0]![0];
|
||
expect(readReadinessReceipt(new Error(formatGatewayReadinessDiagnostic(diagnostic)))).toEqual(
|
||
receipt,
|
||
);
|
||
expect(JSON.stringify(record.mock.calls)).not.toContain(privateDetail);
|
||
expect((failure as Error).message).not.toContain(privateDetail);
|
||
expect(processState.listenerCount("exit")).toBe(0);
|
||
expect(vi.getTimerCount()).toBe(0);
|
||
} finally {
|
||
await vi.runAllTimersAsync();
|
||
await completion;
|
||
vi.useRealTimers();
|
||
}
|
||
});
|
||
|
||
it.each(["timeout", "child-exit", "aborted"] as const)(
|
||
"does not replace a failed response with an expired body on %s",
|
||
async (outcome) => {
|
||
vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] });
|
||
const processState = createGatewayProcessState();
|
||
const lateBody = createDeferred<unknown>();
|
||
const response = new Response(null, { status: 503 });
|
||
response.json = () => lateBody.promise;
|
||
const fetchImpl = vi
|
||
.fn<typeof fetch>()
|
||
.mockResolvedValueOnce(new Response("synthetic refusal", { status: 502 }))
|
||
.mockResolvedValueOnce(response);
|
||
const abort = new AbortController();
|
||
const record = vi.fn<(diagnostic: GatewayReadinessDiagnostic) => void>();
|
||
const completion = testing
|
||
.waitForGatewayReady(processState, [], [], 12345, 25, fetchImpl, abort.signal, record)
|
||
.catch((failure: unknown) => failure);
|
||
try {
|
||
await vi.advanceTimersByTimeAsync(10);
|
||
if (outcome === "child-exit") {
|
||
processState.exitCode = 7;
|
||
processState.emit("exit", 7, null);
|
||
} else if (outcome === "aborted") {
|
||
abort.abort(new Error("synthetic cancellation"));
|
||
}
|
||
await vi.advanceTimersByTimeAsync(15);
|
||
const failure = await completion;
|
||
const diagnostic = record.mock.calls[0]![0];
|
||
expect(diagnostic).toMatchObject({
|
||
outcome,
|
||
lastProbe: { attempt: 2, phase: "body", status: 503 },
|
||
lastFailedResponse: { attempt: 1, phase: "body", status: 502, error: "invalid-json" },
|
||
});
|
||
if (outcome !== "aborted") {
|
||
expect(readReadinessReceipt(failure).lastFailedResponse).toEqual(
|
||
diagnostic.lastFailedResponse,
|
||
);
|
||
}
|
||
const retained = JSON.stringify(diagnostic);
|
||
lateBody.resolve({ ready: false, failing: ["state-database"] });
|
||
await vi.advanceTimersByTimeAsync(0);
|
||
expect(JSON.stringify(diagnostic)).toBe(retained);
|
||
expect(processState.listenerCount("exit")).toBe(0);
|
||
expect(vi.getTimerCount()).toBe(0);
|
||
} finally {
|
||
lateBody.resolve({ ready: false });
|
||
await vi.runAllTimersAsync();
|
||
await completion;
|
||
vi.useRealTimers();
|
||
}
|
||
},
|
||
);
|
||
|
||
it("keeps the formatted receipt below 1024 bytes with maximal sanitized fields", () => {
|
||
const large = Number.MAX_SAFE_INTEGER;
|
||
const lastProbe = {
|
||
attempt: large,
|
||
phase: "complete" as const,
|
||
elapsedMs: large,
|
||
status: 503,
|
||
ready: false,
|
||
failing: Array<string>(8).fill("startup-sidecars"),
|
||
omittedFailing: large,
|
||
};
|
||
const receipt = readReadinessReceipt(
|
||
new Error(
|
||
formatGatewayReadinessDiagnostic({
|
||
attempts: large,
|
||
elapsedMs: large,
|
||
lastProbe,
|
||
lastFailedResponse: {
|
||
attempt: large,
|
||
phase: "body",
|
||
elapsedMs: large,
|
||
status: 502,
|
||
error: "invalid-json",
|
||
},
|
||
child: { pid: large, exitCode: large, signalCode: "SIGVTALRM" },
|
||
}),
|
||
),
|
||
);
|
||
expect(receipt.lastProbe.failing).toHaveLength(8);
|
||
expect(receipt.lastFailedResponse.status).toBe(502);
|
||
});
|
||
|
||
it.each(["fetch-failed", "invalid-json", "body-failed"] as const)(
|
||
"records only the fixed %s readiness category",
|
||
async (category) => {
|
||
const privateDetail = "private readiness failure details";
|
||
const fetchImpl = vi.fn<typeof fetch>(async () => {
|
||
if (category === "fetch-failed") {
|
||
throw new Error(privateDetail);
|
||
}
|
||
if (category === "invalid-json") {
|
||
return new Response(`{${privateDetail}`, { status: 503 });
|
||
}
|
||
const response = new Response(null, { status: 503 });
|
||
response.json = () => Promise.reject(new Error(privateDetail));
|
||
return response;
|
||
});
|
||
const error = await testing
|
||
.waitForGatewayReady(createGatewayProcessState(), [], [], 12345, 25, fetchImpl)
|
||
.catch((failure: unknown) => failure);
|
||
|
||
expect(error).toBeInstanceOf(Error);
|
||
expect((error as Error).message).toContain("timeout waiting for gateway readiness");
|
||
expect((error as Error).message).not.toContain(privateDetail);
|
||
expect(readReadinessReceipt(error)).toMatchObject({
|
||
lastProbe: {
|
||
phase: category === "fetch-failed" ? "headers" : "body",
|
||
error: category,
|
||
...(category === "fetch-failed" ? {} : { status: 503 }),
|
||
},
|
||
lastFailedResponse: category === "fetch-failed" ? null : { status: 503, error: category },
|
||
});
|
||
},
|
||
);
|
||
|
||
it.each(
|
||
(["headers", "body"] as const).flatMap((phase) => [
|
||
{ phase, exits: false, name: `keeps stalled readiness ${phase} inside the startup deadline` },
|
||
{ phase, exits: true, name: `aborts stalled readiness ${phase} when the gateway exits` },
|
||
]),
|
||
)("$name", async ({ phase, exits }) => {
|
||
vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] });
|
||
const processState = createGatewayProcessState(exits ? { pid: 12345 } : {});
|
||
const fetchImpl = createStalledReadinessFetch(phase);
|
||
const cleanup = new AbortController();
|
||
let settled = false;
|
||
const completion = testing
|
||
.waitForGatewayReady(
|
||
processState,
|
||
[],
|
||
[],
|
||
12345,
|
||
exits ? 5_000 : 25,
|
||
fetchImpl,
|
||
cleanup.signal,
|
||
)
|
||
.catch((failure: unknown) => failure)
|
||
.then((result) => {
|
||
settled = true;
|
||
return result;
|
||
});
|
||
const exitTimer = exits
|
||
? setTimeout(() => {
|
||
processState.signalCode = "SIGTERM";
|
||
processState.emit("exit", null, "SIGTERM");
|
||
}, 25)
|
||
: undefined;
|
||
try {
|
||
await vi.advanceTimersByTimeAsync(24);
|
||
expect(settled).toBe(false);
|
||
expect(fetchImpl).toHaveBeenCalledOnce();
|
||
expect(fetchImpl.mock.calls[0]?.[1]?.signal?.aborted).toBe(false);
|
||
expect(processState.listenerCount("exit")).toBe(1);
|
||
|
||
await vi.advanceTimersByTimeAsync(1);
|
||
expect(settled).toBe(true);
|
||
const error = await completion;
|
||
expect(error).toBeInstanceOf(Error);
|
||
expect((error as Error).message).toContain(
|
||
exits ? "gateway exited before readiness" : "timeout waiting for gateway readiness",
|
||
);
|
||
expect(readReadinessReceipt(error)).toMatchObject({
|
||
attempts: 1,
|
||
elapsedMs: 25,
|
||
lastProbe: { attempt: 1, phase, error: exits ? "child-exit" : "timeout", elapsedMs: 25 },
|
||
...(exits ? { child: { pid: 12345, exitCode: null, signalCode: "SIGTERM" } } : {}),
|
||
});
|
||
expect(fetchImpl).toHaveBeenCalledOnce();
|
||
expect(fetchImpl.mock.calls[0]?.[1]?.signal?.aborted).toBe(true);
|
||
expect(processState.listenerCount("exit")).toBe(0);
|
||
expect(vi.getTimerCount()).toBe(0);
|
||
} finally {
|
||
clearTimeout(exitTimer);
|
||
cleanup.abort();
|
||
await completion;
|
||
vi.useRealTimers();
|
||
}
|
||
});
|
||
|
||
it("retains an expired probe snapshot after abort-ignoring headers and body arrive", async () => {
|
||
const processState = createGatewayProcessState({ pid: 12345 });
|
||
const headers = createDeferred<Response>();
|
||
const lateJson = vi.fn(async () => {
|
||
processState.exitCode = 7;
|
||
return { ready: false, failing: ["state-database"] };
|
||
});
|
||
const response = new Response(null, { status: 503 });
|
||
response.json = lateJson;
|
||
const fetchImpl = vi.fn<typeof fetch>(() => headers.promise);
|
||
const releaseLateHeaders = (event: string | symbol) => {
|
||
if (event === "exit") {
|
||
expect(fetchImpl.mock.calls[0]?.[1]?.signal?.aborted).toBe(true);
|
||
headers.resolve(response);
|
||
}
|
||
};
|
||
// Removing the attempt's exit listener follows its snapshot; the old body
|
||
// then completes in microtasks before the existing poll delay can resume.
|
||
processState.on("removeListener", releaseLateHeaders);
|
||
try {
|
||
const error = await testing
|
||
.waitForGatewayReady(processState, [], [], 12345, 5_000, fetchImpl)
|
||
.catch((failure: unknown) => failure);
|
||
|
||
expect(error).toBeInstanceOf(Error);
|
||
expect((error as Error).message).toContain("gateway exited before readiness");
|
||
expect(lateJson).toHaveBeenCalledOnce();
|
||
expect(fetchImpl).toHaveBeenCalledOnce();
|
||
expect(readReadinessReceipt(error)).toEqual({
|
||
attempts: 1,
|
||
elapsedMs: expect.any(Number),
|
||
lastProbe: {
|
||
attempt: 1,
|
||
phase: "headers",
|
||
elapsedMs: expect.any(Number),
|
||
error: "timeout",
|
||
},
|
||
lastFailedResponse: null,
|
||
child: { pid: 12345, exitCode: 7, signalCode: null },
|
||
});
|
||
} finally {
|
||
processState.off("removeListener", releaseLateHeaders);
|
||
}
|
||
});
|
||
|
||
it.runIf(process.platform !== "win32")("preserves native CLI signal termination", async () => {
|
||
const { instance, readAttempts } = await createFakeGateway("signal");
|
||
await expect(trackOperation(instance.cli(["fixture"]))).resolves.toEqual({
|
||
code: null,
|
||
signal: "SIGTERM",
|
||
stdout: "fake gateway attempt 1\n",
|
||
stderr: `${MIGRATION_CONVERGENCE_REFUSAL} fixture\n`,
|
||
});
|
||
const [attempt] = await readAttempts();
|
||
expect(isProcessAlive(attempt!.pid)).toBe(false);
|
||
});
|
||
|
||
it("creates isolated config and spawn env without mutating process env", async () => {
|
||
const previousHome = process.env.HOME;
|
||
const inst = await createOpenClawTestInstance({
|
||
name: "instance-unit",
|
||
gatewayToken: "gateway-token",
|
||
hookToken: "hook-token",
|
||
config: {
|
||
gateway: {
|
||
bind: "loopback",
|
||
},
|
||
},
|
||
env: {
|
||
OPENCLAW_SKIP_CRON: "0",
|
||
},
|
||
});
|
||
|
||
try {
|
||
expect(process.env.HOME).toBe(previousHome);
|
||
expect(inst.homeDir).toBe(path.join(inst.state.root, "home"));
|
||
expect(inst.stateDir).toBe(path.join(inst.homeDir, ".openclaw"));
|
||
expect(inst.configPath).toBe(path.join(inst.stateDir, "openclaw.json"));
|
||
expect(inst.env.HOME).toBe(inst.homeDir);
|
||
expect(inst.env.OPENCLAW_STATE_DIR).toBe(inst.stateDir);
|
||
expect(inst.env.OPENCLAW_CONFIG_PATH).toBe(inst.configPath);
|
||
expect(inst.env.OPENCLAW_SKIP_CRON).toBe("0");
|
||
|
||
const config = JSON.parse(await fs.readFile(inst.configPath, "utf8"));
|
||
expect(config).toStrictEqual({
|
||
gateway: {
|
||
bind: "loopback",
|
||
port: inst.port,
|
||
auth: {
|
||
mode: "token",
|
||
token: "gateway-token",
|
||
},
|
||
controlUi: {
|
||
enabled: false,
|
||
},
|
||
},
|
||
hooks: {
|
||
enabled: true,
|
||
token: "hook-token",
|
||
path: "/hooks",
|
||
},
|
||
});
|
||
} finally {
|
||
await inst.cleanup();
|
||
}
|
||
|
||
await expectPathMissing(inst.state.root);
|
||
});
|
||
});
|