fix(test): managed-process, Vitest shutdown, and worker-artifact tooling tests time out on loaded hosts while polling ready files (#162613)

Test-tooling suites for the managed child-process runner, check-changed cancellation, Vitest fork/worker shutdown, worker artifacts, and run-vitest progress waited for fixture wrappers, workers, and descendants by polling ready files, PID files, and the process table against 1-15 s wall-clock deadlines, which a loaded host outlasts.

Readiness now arrives on the owned channel each fixture already has: runManagedCommand's onReady stream and child events, stdout/IPC readiness written after the fact it announces (PID files are renamed into place first), relayed worker readiness for the Vitest fork/worker fixtures, and retained command outcomes that fail fast if the child settles first. Foreign-descendant extinction without an owner join keeps a deadline-free check bound to the test signal. check-changed-cancellation and run-vitest-progress drop out of the timeout-race baseline. The reverse-release wait in run-vitest-progress and the native Windows Job cases stay unchanged (follow-ups; the latter has no native CI execution). No product source, test timeout, product timeout, or assertion meaning changed.

Part of the polling audit from #162274. Proof on Blacksmith Testbox: delayed-readiness fault probes fail the original bytes and pass these bytes; forced-abort probes reach cleanup; 20/20 standalone runs per file; 3/3 replays of all seven owning tooling shards on the rebased head; tsgo; Codex autoreview and ClawSweeper clean.
This commit is contained in:
Peter Steinberger 2026-10-01 03:55:51 -07:00 • committed by GitHub
parent 8941a07d69
commit b9512e7883
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 805 additions and 360 deletions

View file

@ -137,7 +137,6 @@ test/helpers/openclaw-test-instance.acquisition.test.ts 1
test/helpers/qa-gateway-test-lifetime.test.ts 2
test/plugins/browser-session-authority.gateway.test.ts 2
test/plugins/team-reports-http.gateway.test.ts 1
test/scripts/check-changed-cancellation.test.ts 2
test/scripts/ci-run-node-test-shard.test.ts 3
test/scripts/compiler-input-snapshot.test.ts 1
test/scripts/nested-retention.test-support.ts 4
@ -147,7 +146,6 @@ test/scripts/process-fixture-cleanup.test.ts 3
test/scripts/profile-extension-memory.test.ts 2
test/scripts/run-tsgo.test.ts 2
test/scripts/run-vitest-bounded.test.ts 1
test/scripts/run-vitest-progress.test.ts 5
test/scripts/run-with-env.test.ts 1
test/scripts/test-projects-build-admission.test.ts 7
test/scripts/upstream-repository-advisories.test.ts 2

View file

@ -22,6 +22,15 @@ export function installVitestShutdownCancellation({ root, preload }) {
},
});
publish("shim.pid", child.pid);
let output = "";
const forwardReady = (chunk) => {
output += chunk.toString();
if (output.includes("shutdown-worker-ready\n")) {
child.stdout.off("data", forwardReady);
fs.writeSync(1, "shutdown-worker-ready\n");
}
};
child.stdout.on("data", forwardReady);
return child;
};
// Release after every TERM listener and its microtasks run. An immediate KILL
@ -47,6 +56,8 @@ export function installVitestShutdownCancellation({ root, preload }) {
if (target === file("receipt.json")) {
// Stop at the actual worker receipt, before any test result or teardown.
publish("worker.pid", process.pid);
// Flush to the already-owned pipe before suspending this event loop.
fs.writeSync(1, "shutdown-worker-ready\n");
process.kill(process.pid, "SIGSTOP");
}
return result;

View file

@ -1,20 +1,36 @@
import type { ChildProcess } from "node:child_process";
import fs from "node:fs";
import path from "node:path";
import { setTimeout as waitForReaper } from "node:timers/promises";
import { pathToFileURL } from "node:url";
import { afterEach, describe, expect, it } from "vitest";
import { hasErrnoCode } from "../../src/infra/errno.js";
import { createFixtureLifetime } from "../helpers/fixture-lifetime.js";
import { isProcessAlive, waitForDead, waitForFixtureFile } from "../helpers/process-wait.js";
import { withTestTimeout } from "../helpers/promise.js";
import { isProcessAlive } from "../helpers/process-wait.js";
import { awaitGateBeforeSettlement, createDeferred, withinTest } from "../helpers/promise.js";
import { runNodeScript } from "../helpers/run-node-script.js";
const fixture = createFixtureLifetime();
afterEach(() => fixture.cleanup());
// The detached implementation/leaf PIDs have no ChildProcess handles in this harness.
async function waitForRecordedPidsDead(pids: number[], signal: AbortSignal) {
for (const pid of pids) {
try {
while (isProcessAlive(pid)) {
await waitForReaper(10, undefined, { signal });
}
} catch (error) {
throw new Error(`process still alive: ${pid}`, { cause: error });
}
}
}
describe.skipIf(process.platform === "win32")("check-changed public wrapper cancellation", () => {
it.each(["resistant", "cooperative", "failure"] as const)(
it.for(["resistant", "cooperative", "failure"] as const)(
"joins a %s check and stops admitting commands",
async (mode) => {
{ timeout: 20_000 },
async (mode, { signal }) => {
await fixture.run(async () => {
const cwd = fixture.createTempDir("check-changed-cancellation-");
const wrapperPath = path.resolve("scripts/check-changed.mjs");
@ -23,7 +39,7 @@ describe.skipIf(process.platform === "win32")("check-changed public wrapper canc
const pidPaths = ["implementation", "command", "descendant"].map((name) =>
path.join(cwd, `${name}.pid`),
);
const readyPath = path.join(cwd, "ready.pid");
const ready = createDeferred();
const clockPath = path.join(cwd, "supervisor-clock.mjs");
const binDir = path.join(cwd, "bin");
fs.mkdirSync(binDir);
@ -87,7 +103,7 @@ child.once("close", () => {
});
setTimeout(() => process.exit(98), 15000);
child.once("message", () => {
fs.writeFileSync(${JSON.stringify(readyPath)}, String(process.pid));
fs.writeSync(1, "changed-check command ready\\n");
});
`,
{ mode: 0o755 },
@ -108,8 +124,13 @@ child.once("message", () => {
runNodeScript([wrapperPath, "--staged", "--", "README.md"], env, 10_000, {
cwd,
maxBuffer: 64 * 1024,
onReady(child) {
onReady(child, readOutput) {
wrapper = child;
child.stdout!.on("data", () => {
if (readOutput().stdout.includes("changed-check command ready\n")) {
ready.resolve();
}
});
},
}),
);
@ -126,16 +147,19 @@ child.once("message", () => {
});
try {
if (mode !== "failure") {
// The receipt is written only after both leaf signal handlers are installed.
await withTestTimeout(
waitForFixtureFile(readyPath, completion),
5_000,
"changed-check command did not become ready",
// The relay writes readiness only after both leaf signal handlers are installed.
await withinTest(
awaitGateBeforeSettlement(
ready.promise,
completion,
"changed-check command did not become ready",
),
signal,
);
expect(readOwnedPids()).toHaveLength(3);
expect(wrapper?.kill("SIGTERM")).toBe(true);
}
const result = await completion;
const result = await withinTest(completion, signal);
expect(result.error, result.stderr).toBeUndefined();
expect(result.status, result.stderr).toBe(mode === "failure" ? 7 : 143);
expect(result.stderr.trim().split("\n").at(-1)).toBe(
@ -163,20 +187,17 @@ child.once("message", () => {
try {
process.kill(pid, "SIGKILL");
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "ESRCH") {
if (!hasErrnoCode(error, "ESRCH")) {
throw error;
}
}
}
}
await Promise.all([
...pids.map((pid) => waitForDead(pid, 2_000)),
withTestTimeout(completion, 2_000, "wrapper output did not close during cleanup"),
]);
await completion;
await waitForRecordedPidsDead(pids, signal);
});
}
});
},
20_000,
);
});

View file

@ -2,11 +2,12 @@ import { execFileSync } from "node:child_process";
import { getEventListeners } from "node:events";
import fs from "node:fs";
import path from "node:path";
import { setTimeout as waitForReaper } from "node:timers/promises";
import { afterEach, beforeEach, expect, it, vi } from "vitest";
import { createVitestResourceOwner } from "../../scripts/lib/vitest-resource-ownership.mts";
import { runNodeStep } from "../../scripts/prepare-extension-package-boundary-artifacts.mts";
import { createFixtureLifetime } from "../helpers/fixture-lifetime.js";
import { isProcessAlive, waitForDead } from "../helpers/process-wait.js";
import { isProcessAlive } from "../helpers/process-wait.js";
import { createDeferred } from "../helpers/promise.js";
import { useAutoCleanupTempDirTracker } from "../helpers/temp-dir.js";
@ -25,6 +26,18 @@ afterEach(() => vi.unstubAllEnvs());
const fixture = createFixtureLifetime();
afterEach(() => fixture.cleanup());
// runNodeStep exposes its joined outcome, but no ChildProcess handle for rescue
// after an unverified join. Observe only the fixture's recorded PID.
async function waitForRescuedChild(pid: number, signal: AbortSignal) {
try {
while (isProcessAlive(pid)) {
await waitForReaper(10, undefined, { signal });
}
} catch (error) {
throw new Error(`process still alive: ${pid}`, { cause: error });
}
}
it("releases inputs and claims after a native execFileSync ENOENT error", async () => {
const lifetime = createFixtureLifetime();
const root = lifetime.createTempDir("fixture-missing-command-");
@ -361,7 +374,7 @@ it
try {
if (pid && isProcessAlive(pid)) {
process.kill(pid, "SIGKILL");
await waitForDead(pid, 2_000);
await waitForRescuedChild(pid, contextSignal);
}
} finally {
await rescue;

File diff suppressed because it is too large Load diff

View file

@ -1,40 +1,85 @@
import fs from "node:fs";
import path from "node:path";
import { setTimeout as delay } from "node:timers/promises";
import { pathToFileURL } from "node:url";
import { afterEach, describe, it, vi } from "vitest";
import { afterAll, afterEach, beforeAll, describe, it, vi, type TestContext } from "vitest";
import * as managedChild from "../../scripts/lib/managed-child-process.mts";
import { resolveVitestCliEntry } from "../../scripts/lib/vitest-build-prerequisites.mts";
import { resolveVitestNodeArgs } from "../../scripts/lib/vitest-process-env.mts";
import { createVitestWorkerRun } from "../../scripts/lib/vitest-worker-run.mts";
import { resolveVitestSpawnParams, spawnWatchedVitestProcess } from "../../scripts/run-vitest.mts";
import { forceKillVitestProcessGroup } from "../../scripts/vitest-process-group.mts";
import {
fixtureReceiptClientSource,
openFixtureReceiptChannel,
type FixtureReceiptChannel,
} from "../helpers/fixture-receipts.js";
import { isProcessAlive } from "../helpers/process-wait.js";
import { createDeferred, withTestTimeout } from "../helpers/promise.js";
import { awaitGateBeforeSettlement, createDeferred, withinTest } from "../helpers/promise.js";
import { useAutoCleanupTempDirTracker } from "../helpers/temp-dir.js";
import { createControlledWorkerCompiler } from "./vitest-worker-artifacts.test-support.js";
const repoRoot = path.resolve(import.meta.dirname, "../..");
const posixDescribe = process.platform === "win32" ? describe.skip : describe.concurrent;
const posixSerialDescribe = process.platform === "win32" ? describe.skip : describe;
const ioTimeoutMs = 15_000;
const silenceMs = 1_000;
let receipts: FixtureReceiptChannel;
const progressBodies = new WeakMap<TestContext, Promise<void>>();
async function waitForRealIo(ready: () => boolean, description: string) {
const deadline = Date.now() + ioTimeoutMs;
while (!ready()) {
if (Date.now() >= deadline) {
throw new Error(`Timed out waiting for ${description}`);
}
await delay(5);
}
beforeAll(async () => {
receipts = await openFixtureReceiptChannel();
});
afterAll(async () => {
await receipts?.close();
});
function joinedProgressTest(body: (context: TestContext) => Promise<void>) {
return (context: TestContext) => {
// Timeout aborts the wait; the finish hook also owns the asynchronous finally.
const run = Promise.resolve().then(() => {
context.signal.throwIfAborted();
return body(context);
});
progressBodies.set(context, run);
context.onTestFinished(() => run);
return run;
};
}
function cleanupAfterProgressBody(cleanup: () => void) {
afterEach(async (context) => {
// Vitest runs afterEach before onTestFinished; retained writers must join first.
await Promise.allSettled([progressBodies.get(context)]);
cleanup();
});
}
function fixtureReadyBeforeSettlement(
readyPath: string,
operation: PromiseLike<unknown>,
description: string,
) {
// The worker publishes its PID before reporting readiness. Receipt delivery
// can trail completion on the independently owned output pipes.
const settled = Promise.resolve(operation).then(
() => {
if (!fs.existsSync(readyPath)) {
throw new Error(`Timed out waiting for ${description}`);
}
},
(error: unknown) => {
if (!fs.existsSync(readyPath)) {
throw error;
}
},
);
return Promise.race([receipts.waitFor(readyPath, "ready"), settled]);
}
posixDescribe.each([false, true])(
"keeps real case progress alive across watchdog windows, then stall=%s",
(stall) => {
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
it("reports the expected outcome and stops its process group", async ({ expect }) => {
const tempDirs = useAutoCleanupTempDirTracker(cleanupAfterProgressBody);
const testProgress = joinedProgressTest(async ({ expect, signal }) => {
const root = tempDirs.make("oc-vt-progress-");
fs.symlinkSync(
path.join(repoRoot, "node_modules"),
@ -51,11 +96,13 @@ posixDescribe.each([false, true])(
`import fs from "node:fs";
import { expect, it } from "vitest";
import { waitForFile } from ${JSON.stringify(path.join(repoRoot, "test/helpers/process-wait.ts"))};
${fixtureReceiptClientSource(receipts.endpoint)}
const index = ${index};
it("real progress " + index, async () => {
const ready = ${JSON.stringify(root)} + "/ready-" + index;
fs.writeFileSync(ready + ".tmp", String(process.pid));
fs.renameSync(ready + ".tmp", ready);
sendReceipt(ready, "ready");
await waitForFile(${JSON.stringify(root)} + "/release-" + index, 15000);
expect(index).toBeLessThan(5);
});
@ -126,35 +173,55 @@ export default {
vi.useRealTimers();
}
let output = "";
watched.child.stdout!.on("data", (chunk: string) => {
output += chunk;
});
watched.child.stderr!.on("data", (chunk: string) => {
output += chunk;
});
const caseCompletions = Array.from({ length: 5 }, () => createDeferred());
const casePassed = (index: number) =>
output
.split("\n")
.some((line) => line.includes("✓") && line.includes(` > real progress ${index} `));
const observeOutput = (chunk: string) => {
output += chunk;
for (const [index, completion] of caseCompletions.entries()) {
if (casePassed(index)) {
completion.resolve();
}
}
};
watched.child.stdout!.on("data", observeOutput);
watched.child.stderr!.on("data", observeOutput);
let workerPid: number | undefined;
try {
for (let index = 0; index < 4; index++) {
const readyPath = path.join(root, `ready-${index}`);
await waitForRealIo(() => fs.existsSync(readyPath), `case ${index} readiness\n${output}`);
await withinTest(
fixtureReadyBeforeSettlement(
readyPath,
watched.completion,
`case ${index} readiness\n${output}`,
),
signal,
);
workerPid = Number(fs.readFileSync(readyPath, "utf8"));
clock.tick(600);
expect(onNoOutputTimeout, output).not.toHaveBeenCalled();
fs.writeFileSync(path.join(root, `release-${index}`), "");
// File barriers never enter the watched pipes. Only Vitest's completed
// case output can reset the watchdog before the next 600ms advance.
await waitForRealIo(
() => casePassed(index),
`Vitest completion for case ${index}\n${output}`,
await withinTest(
awaitGateBeforeSettlement(
caseCompletions[index]!.promise,
watched.completion,
`Timed out waiting for Vitest completion for case ${index}\n${output}`,
),
signal,
);
}
await waitForRealIo(
() => fs.existsSync(path.join(root, "ready-4")),
"final case readiness",
await withinTest(
fixtureReadyBeforeSettlement(
path.join(root, "ready-4"),
watched.completion,
"final case readiness",
),
signal,
);
expect(isProcessAlive(watched.child.pid!)).toBe(true);
expect(onNoOutputTimeout).not.toHaveBeenCalled();
@ -166,7 +233,7 @@ export default {
} else {
fs.writeFileSync(path.join(root, "release-4"), "");
}
const result = await watched.completion;
const result = await withinTest(watched.completion, signal);
// Vitest's logger handles SIGTERM and exits with 128 + 15, rather than
// leaving Node to report a signal-only exit (as a bare silent child does).
expect(result, output).toEqual({ code: stall ? 143 : 0, signal: null, groupJoined: true });
@ -196,18 +263,18 @@ export default {
} finally {
watched.teardown();
forceKillVitestProcessGroup(watched.child);
await withTestTimeout(watched.completion, ioTimeoutMs, "owned Vitest group did not stop");
await watched.completion;
}
});
it("reports the expected outcome and stops its process group", testProgress);
},
);
posixSerialDescribe("compiled subprocess preparation progress", { concurrent: false }, () => {
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
const tempDirs = useAutoCleanupTempDirTracker(cleanupAfterProgressBody);
it.for(["valid", "tampered"] as const)(
"counts accepted and verified work without hiding a later stall (%s)",
async (verification, { expect }) => {
const testPreparation = (verification: "valid" | "tampered") =>
joinedProgressTest(async ({ expect, signal }) => {
const directory = tempDirs.make("oc-vt-preparation-progress-");
const env = {
...process.env,
@ -291,8 +358,10 @@ process.send({event:"ready"});
vi.useRealTimers();
}
const handle = watched;
let ready = false;
const ready = createDeferred();
let duplicateAcknowledged = createDeferred();
let duplicates = 0;
const replied = createDeferred();
const replies: unknown[] = [];
let output = "";
handle.child.on("message", (message: unknown) => {
@ -300,13 +369,15 @@ process.send({event:"ready"});
return;
}
if (message.event === "ready") {
ready = true;
ready.resolve();
}
if (message.event === "duplicate") {
duplicates += 1;
duplicateAcknowledged.resolve();
}
if (message.event === "reply") {
replies.push(message);
replied.resolve();
}
});
handle.child.stdout!.on("data", (chunk: string) => {
@ -315,15 +386,26 @@ process.send({event:"ready"});
handle.child.stderr!.on("data", (chunk: string) => {
output += chunk;
});
const waitForBorrower = (gate: PromiseLike<unknown>, description: string) =>
withinTest(
awaitGateBeforeSettlement(
gate,
handle.completion,
`Timed out waiting for ${description}`,
),
signal,
);
const duplicate = async () => {
const expected = duplicates + 1;
duplicateAcknowledged = createDeferred();
handle.child.send("duplicate");
await waitForRealIo(() => duplicates === expected, "duplicate borrower IPC delivery");
await waitForBorrower(duplicateAcknowledged.promise, "duplicate borrower IPC delivery");
expect(duplicates).toBe(expected);
};
await waitForRealIo(() => ready, "native borrower readiness");
await waitForBorrower(ready.promise, "native borrower readiness");
clock.tick(100_000);
handle.child.send("request");
await withTestTimeout(compilerStarted.promise, ioTimeoutMs, "compiler admission");
await waitForBorrower(compilerStarted.promise, "compiler admission");
clock.tick(100_000);
// No watched pipe output exists: only accepted owner work can keep
// this real borrower alive past its original 120-second deadline.
@ -332,7 +414,7 @@ process.send({event:"ready"});
expect(isProcessAlive(handle.child.pid!)).toBe(true);
await duplicate();
releaseCompiler.resolve();
await withTestTimeout(verificationStarted.promise, ioTimeoutMs, "artifact verification");
await waitForBorrower(verificationStarted.promise, "artifact verification");
expect(controlled.read()).toHaveLength(1);
clock.tick(10_000);
expect(onNoOutputTimeout).not.toHaveBeenCalled();
@ -340,7 +422,8 @@ process.send({event:"ready"});
fs.appendFileSync(heldOutput, "\nchanged after compilation\n");
}
releaseVerification.resolve();
await waitForRealIo(() => replies.length === 1, "verified borrower reply");
await waitForBorrower(replied.promise, "verified borrower reply");
expect(replies).toHaveLength(1);
expect(output).toBe("");
if (verification === "valid") {
expect(replies[0]).toEqual({ event: "reply", ok: true });
@ -359,11 +442,7 @@ process.send({event:"ready"});
expect(fs.existsSync(generation)).toBe(true);
clock.tick(1);
expect(onNoOutputTimeout).toHaveBeenCalledOnce();
const result = await withTestTimeout(
handle.completion,
ioTimeoutMs,
"timed-out borrower join",
);
const result = await withinTest(handle.completion, signal);
expect(handle.child.exitCode).toBe(0);
expect(result).toEqual({ code: 1, signal: null, groupJoined: true });
expect(isProcessAlive(handle.child.pid!)).toBe(false);
@ -374,7 +453,7 @@ process.send({event:"ready"});
if (watched) {
watched.teardown();
forceKillVitestProcessGroup(watched.child);
await withTestTimeout(watched.completion, ioTimeoutMs, "borrower cleanup");
await watched.completion;
}
} finally {
try {
@ -395,6 +474,9 @@ process.send({event:"ready"});
expect(disposalFailure).toBeUndefined();
}
expect(fs.existsSync(generation)).toBe(false);
},
});
it.for(["valid", "tampered"] as const)(
"counts accepted and verified work without hiding a later stall (%s)",
(verification, context) => testPreparation(verification)(context),
);
});

View file

@ -1,15 +1,12 @@
import { execFileSync, type ChildProcess } from "node:child_process";
import fs from "node:fs";
import path from "node:path";
import { setTimeout as waitForReaper } from "node:timers/promises";
import { fileURLToPath } from "node:url";
import { expect, it, vi, type TestContext } from "vitest";
import { inspectManagedProcessGroup } from "../../scripts/lib/managed-child-process.mts";
import {
isProcessAlive,
waitForDead,
waitForFile,
waitForFixtureFile,
} from "../helpers/process-wait.js";
import { isProcessAlive } from "../helpers/process-wait.js";
import { awaitGateBeforeSettlement, createDeferred, withinTest } from "../helpers/promise.js";
import { createTempDirTracker } from "../helpers/temp-dir.js";
import { runVitestShutdownCommand } from "../helpers/vitest-shutdown-command.js";
import { fixturePreloadArgs } from "./fixtures/ci-fixture-runtime.cjs";
@ -21,6 +18,30 @@ const fixtureRoots = fileURLToPath(
);
fs.mkdirSync(fixtureRoots, { recursive: true });
function observeReadyLine(child: ChildProcess, line: string, ready: () => void) {
let output = "";
const observe = (chunk: Buffer) => {
output += chunk.toString();
if (output.includes(line)) {
child.stdout!.off("data", observe);
ready();
}
};
child.stdout!.on("data", observe);
}
async function waitForFixtureProcessesDead(pids: number[], signal: AbortSignal) {
// Group completion can precede Darwin's foreign-zombie reap. The harness
// has no child handles for those PIDs; only the owning test bounds observation.
try {
while (pids.some(isProcessAlive)) {
await waitForReaper(5, undefined, { signal });
}
} catch (cause) {
throw new Error(`process still alive: ${pids.filter(isProcessAlive).join(", ")}`, { cause });
}
}
function runJoinedShutdownTest(context: TestContext, body: () => Promise<void>) {
// Register before the body starts: outer cancellation must join every continuation and finally.
const run = Promise.resolve().then(() => {
@ -216,6 +237,7 @@ installVitestShutdownCancellation({root:${JSON.stringify(root)},preload:import.m
`,
);
let child!: ChildProcess;
const workerReady = createDeferred();
let fireDeadline: (() => void) | undefined;
const schedule = globalThis.setTimeout;
// Capture only this command's deadline; readiness and native cleanup keep real timers.
@ -249,6 +271,7 @@ installVitestShutdownCancellation({root:${JSON.stringify(root)},preload:import.m
{
onReady(owned) {
child = owned;
observeReadyLine(child, "shutdown-worker-ready\n", () => workerReady.resolve());
},
signal: context.signal,
},
@ -262,7 +285,14 @@ installVitestShutdownCancellation({root:${JSON.stringify(root)},preload:import.m
);
const pids: number[] = [];
try {
await waitForFixtureFile(path.join(root, "worker.pid"), invocation);
await withinTest(
awaitGateBeforeSettlement(
workerReady.promise,
invocation,
`Child exited before writing ${path.join(root, "worker.pid")}`,
),
context.signal,
);
const shim = Number(fs.readFileSync(path.join(root, "shim.pid"), "utf8"));
const worker = Number(fs.readFileSync(path.join(root, "worker.pid"), "utf8"));
context.signal.throwIfAborted();
@ -308,8 +338,9 @@ installVitestShutdownCancellation({root:${JSON.stringify(root)},preload:import.m
).toBeTypeOf("function");
fireDeadline!();
}
const result = await outcome;
await waitForFile(path.join(root, "term-received"), 1_000);
const result = await withinTest(outcome, context.signal);
// The fixture writes this synchronously in its TERM handler, before exit.
expect(fs.existsSync(path.join(root, "term-received"))).toBe(true);
console.log(
JSON.stringify({
mode,
@ -343,9 +374,7 @@ installVitestShutdownCancellation({root:${JSON.stringify(root)},preload:import.m
child.kill("SIGTERM");
}
await outcome;
for (const pid of pids) {
await waitForDead(pid, 5_000);
}
await waitForFixtureProcessesDead(pids, context.signal);
}
}),
);
@ -368,8 +397,10 @@ process.on("SIGTERM", () => {
setInterval(() => {}, 1000);
fs.writeFileSync(process.argv[1] + ".tmp", String(process.pid));
fs.renameSync(process.argv[1] + ".tmp", process.argv[1]);
process.stdout.write("descendant-ready\\n");
`;
const controller = new AbortController();
const descendantReady = createDeferred();
let child!: ChildProcess;
const invocation = runVitestShutdownCommand({
args: [
@ -391,6 +422,7 @@ spawn(process.execPath, ["-e", ${JSON.stringify(descendantSource)}, process.argv
signal: AbortSignal.any([context.signal, controller.signal]),
onReady(owned) {
child = owned;
observeReadyLine(child, "descendant-ready\n", () => descendantReady.resolve());
},
});
const outcome = invocation.then(
@ -398,7 +430,14 @@ spawn(process.execPath, ["-e", ${JSON.stringify(descendantSource)}, process.argv
(error: unknown) => ({ result: undefined, error }),
);
try {
await waitForFixtureFile(descendantPidPath, invocation);
await withinTest(
awaitGateBeforeSettlement(
descendantReady.promise,
invocation,
`Child exited before writing ${descendantPidPath}`,
),
context.signal,
);
const descendantPid = Number(fs.readFileSync(descendantPidPath, "utf8"));
if (!Number.isSafeInteger(descendantPid) || descendantPid <= 0) {
throw new Error("Invalid descendant PID receipt");
@ -426,7 +465,7 @@ spawn(process.execPath, ["-e", ${JSON.stringify(descendantSource)}, process.argv
expect(inspectManagedProcessGroup(child, { errorPolicy: "indeterminate" })).toBe("live");
controller.abort();
const result = await outcome;
const result = await withinTest(outcome, context.signal);
expect(result.error).toMatchObject({ code: "ABORT_ERR" });
for (const role of ["leader", "descendant"]) {
expect(result.error).toMatchObject({

View file

@ -1,17 +1,30 @@
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import { setTimeout as tick } from "node:timers/promises";
import { fileURLToPath, pathToFileURL } from "node:url";
import { expect } from "vitest";
import { isConstrainedCiCheckHost } from "../../scripts/lib/local-check-runtime.mts";
import { isProcessAlive, waitForDead, waitForFixtureFile } from "../helpers/process-wait.js";
import {
fixtureReceiptClientSource,
openFixtureReceiptChannel,
type FixtureReceiptChannel,
} from "../helpers/fixture-receipts.js";
import { isProcessAlive } from "../helpers/process-wait.js";
import { withinTest } from "../helpers/promise.js";
import {
createControlledWorkerCompiler,
createWorkerArtifactTest,
fixtureFileBeforeSettlement,
writeFixture,
} from "./vitest-worker-artifacts.test-support.js";
const it = createWorkerArtifactTest();
let receipts: FixtureReceiptChannel;
it.beforeAll(async () => {
receipts = await openFixtureReceiptChannel();
});
it.afterAll(() => receipts.close());
const root = process.cwd();
const command = ["--import", "tsx", "scripts/ci-run-node-test-shard.mts"];
type Observation = {
@ -25,6 +38,21 @@ type Observation = {
};
const generationDirectory = (generation: string) => fileURLToPath(new URL("../../", generation));
// Native grandchildren expose no harness-owned close event. Missing-claim errors can suppress
// group-end output, and Darwin managed joins can precede orphan-zombie reaping.
async function waitForBorrowerExit(pid: number, signal: AbortSignal): Promise<void> {
try {
while (isProcessAlive(pid)) {
await tick(10, undefined, { signal });
}
} catch (error) {
if (signal.aborted) {
throw new Error(`process still alive: ${pid}`, { cause: error });
}
throw error;
}
}
function createCiProbe(
directory: string,
retain = false,
@ -51,6 +79,7 @@ function createCiProbe(
directory,
"child.test.ts",
`
${fixtureReceiptClientSource(receipts.endpoint)}
import fs from 'node:fs';
import { createHash } from 'node:crypto';
import { fileURLToPath } from 'node:url';
@ -97,12 +126,14 @@ function createCiProbe(
}
} finally {
fs.writeFileSync(${JSON.stringify(firstReady)}, 'ready');
sendReceipt(${JSON.stringify(firstReady)}, 'written');
}
if (${generationClaim === "released"}) {
throw new Error('ordinary failure after releasing generation claim');
}
} else {
fs.writeFileSync(${JSON.stringify(ready)}, 'ready');
sendReceipt(${JSON.stringify(ready)}, 'written');
await waitForSignal(${JSON.stringify(release)});
fs.accessSync(generation);
fs.writeFileSync(${JSON.stringify(ready + ".read")}, 'read after sibling exit');
@ -159,7 +190,7 @@ it.runIf(process.platform !== "win32").for([
{ parallelism: 2, shared: false },
])(
"owns real CI group generations (parallelism=$parallelism, shared=$shared)",
({ parallelism, shared }, { workerArtifacts }) =>
({ parallelism, shared }, { workerArtifacts, signal }) =>
workerArtifacts.fixtureLifetime.run(async () => {
const { node } = workerArtifacts.createFixtureCommands();
const directory = workerArtifacts.fixtureDirectory();
@ -198,13 +229,15 @@ it.runIf(process.platform !== "win32").for([
);
expect(result.code, result.stderr + result.stdout).toBe(0);
if (controlled) {
const receipts = controlled.read();
console.log("Controlled compiler receipts", JSON.stringify(receipts));
expect(receipts).toHaveLength(shared ? 1 : 2);
const compilerReceipts = controlled.read();
console.log("Controlled compiler receipts", JSON.stringify(compilerReceipts));
expect(compilerReceipts).toHaveLength(shared ? 1 : 2);
expect(
new Set(receipts.map(({ pid, processStartTime }) => `${pid}:${processStartTime}`)).size,
).toBe(receipts.length);
for (const receipt of receipts) {
new Set(
compilerReceipts.map(({ pid, processStartTime }) => `${pid}:${processStartTime}`),
).size,
).toBe(compilerReceipts.length);
for (const receipt of compilerReceipts) {
expect(receipt).toMatchObject({
processStartTime: expect.any(Number),
isMainThread: true,
@ -242,12 +275,18 @@ it.runIf(process.platform !== "win32").for([
}
} finally {
const observations = fs.existsSync(fixture.observationsFile) ? fixture.read() : [];
await Promise.all(
observations.flatMap(({ pid, parent }) => [
waitForDead(pid, 5_000),
waitForDead(parent, 5_000),
]),
);
// Shared group joins prove extinction. The non-detached fixture instead uses node()'s
// managed join, which can accept Darwin zombies before the kernel reaps their PIDs.
for (const { pid, parent } of observations) {
if (!shared && process.platform === "darwin") {
await Promise.all([
waitForBorrowerExit(pid, signal),
waitForBorrowerExit(parent, signal),
]);
}
expect(isProcessAlive(pid), `process still alive: ${pid}`).toBe(false);
expect(isProcessAlive(parent), `process still alive: ${parent}`).toBe(false);
}
for (const run of new Set(observations.map(({ generation }) => generation))) {
fs.rmSync(generationDirectory(run), { recursive: true, force: true });
}
@ -281,7 +320,7 @@ it
name: "removes a shared generation after a borrower releases its generation claim and fails normally",
claim: "released",
},
] as const)("$name", ({ claim }, { workerArtifacts }) =>
] as const)("$name", ({ claim }, { workerArtifacts, signal }) =>
workerArtifacts.fixtureLifetime.run(async () => {
const { node } = workerArtifacts.createFixtureCommands();
const directory = workerArtifacts.fixtureDirectory();
@ -294,24 +333,27 @@ it
const controlled = createControlledWorkerCompiler(directory, env);
const running = node(command, root, controlled.env);
try {
await waitForFixtureFile(fixture.ready, running);
await withinTest(fixtureFileBeforeSettlement(receipts, fixture.ready, running), signal);
// Hold first-group until its sibling is borrowing, then join its own receipt.
fs.writeFileSync(fixture.startFirst, "start");
await waitForFixtureFile(fixture.firstReady, running);
await withinTest(fixtureFileBeforeSettlement(receipts, fixture.firstReady, running), signal);
const observations = fixture.read();
const first = observations.find(({ group }) => group === "first-group")!;
const second = observations.find(({ group }) => group === "second-group")!;
expect(observations).toHaveLength(2);
expect(new Set(observations.map(({ generation }) => generation)).size).toBe(1);
await Promise.all([waitForDead(first.pid, 5_000), waitForDead(first.parent, 5_000)]);
await Promise.all([
waitForBorrowerExit(first.pid, signal),
waitForBorrowerExit(first.parent, signal),
]);
expect(isProcessAlive(second.pid)).toBe(true);
expect(fs.existsSync(generationDirectory(first.generation))).toBe(true);
expect(fs.existsSync(first.includeFile)).toBe(true);
fs.writeFileSync(fixture.release, "finish");
const result = await running;
const receipts = controlled.read();
expect(receipts).toHaveLength(1);
console.log("Controlled compiler receipts", JSON.stringify(receipts));
const result = await withinTest(running, signal);
const compilerReceipts = controlled.read();
expect(compilerReceipts).toHaveLength(1);
console.log("Controlled compiler receipts", JSON.stringify(compilerReceipts));
expect(result.code).not.toBe(0);
if (claim === "temporary") {
expect(result.stdout + result.stderr).toContain("retained temporary namespace");
@ -340,12 +382,11 @@ it
fs.writeFileSync(fixture.release, "finish");
await running;
const observations = fs.existsSync(fixture.observationsFile) ? fixture.read() : [];
await Promise.all(
observations.flatMap(({ pid, parent }) => [
waitForDead(pid, 5_000),
waitForDead(parent, 5_000),
]),
);
// The CI command finishes only after both group completions and worker-run disposal.
for (const { pid, parent } of observations) {
expect(isProcessAlive(pid), `process still alive: ${pid}`).toBe(false);
expect(isProcessAlive(parent), `process still alive: ${parent}`).toBe(false);
}
for (const run of new Set(observations.map(({ generation }) => generation))) {
fs.rmSync(generationDirectory(run), { recursive: true, force: true });
}

View file

@ -12,6 +12,10 @@ import {
import { resolveVitestSpawnParams, spawnWatchedVitestProcess } from "../../scripts/run-vitest.mts";
import { createVitestProcessCompletion } from "../../scripts/vitest-process-group.mts";
import { createFixtureLifetime } from "../helpers/fixture-lifetime.js";
import {
fixtureReceiptClientSource,
type FixtureReceiptChannel,
} from "../helpers/fixture-receipts.js";
import { runNodeScript } from "../helpers/run-node-script.js";
import { fixturePreloadEnv } from "./fixtures/ci-fixture-runtime.cjs";
@ -240,6 +244,35 @@ export function writeFixture(directory: string, name: string, source: string) {
return filename;
}
// Fixture records precede receipts and child completion; separate pipes can deliver them out of order.
export function fixtureFileBeforeSettlement(
receipts: FixtureReceiptChannel,
filename: string,
completion: PromiseLike<unknown>,
expected?: string,
): Promise<void> {
const matches = () => {
if (!fs.existsSync(filename)) {
return false;
}
const text = fs.readFileSync(filename, "utf8");
return text.length > 0 && (expected === undefined || text === expected);
};
const settled = Promise.resolve(completion).then(
() => {
if (!matches()) {
throw new Error(`Child exited before writing ${filename}`);
}
},
(error: unknown) => {
if (!matches()) {
throw new Error(`Child failed before writing ${filename}`, { cause: error });
}
},
);
return Promise.race([receipts.waitFor(filename, expected ?? "written"), settled]);
}
export function workerBorrowingProbe(directory: string) {
const value = writeFixture(directory, "value.ts", 'export const value: string = "first";');
const test = writeFixture(
@ -281,12 +314,14 @@ export function workerProbe(
directory: string,
holdSecond = false,
mode: "compiled" | "source" = "compiled",
receiptEndpoint?: string,
) {
const value = writeFixture(directory, "value.ts", 'export const value: string = "first";');
const test = writeFixture(
directory,
"child.test.ts",
`
${receiptEndpoint ? fixtureReceiptClientSource(receiptEndpoint) : ""}
import * as cp from 'node:child_process';
import fs from 'node:fs';
import path from 'node:path';
@ -369,6 +404,7 @@ export function workerProbe(
expect(args[sourceLoader ? 2 : 0]).toMatch(sourceMode ? /\\.ts$/ : /\\.js$/);
fs.appendFileSync(${JSON.stringify(path.join(directory, "observations.jsonl"))}, JSON.stringify({args, tuiUrls, setupUrls, retentionUrl, value, configValue:inject('configValue'), knn:resolveRuntimeWorkerUrl(vectorKnnProcessEntrypoint).href})+'\\n');
fs.appendFileSync(${JSON.stringify(path.join(directory, "generations.jsonl"))}, JSON.stringify(generation)+'\\n');
${receiptEndpoint ? `sendReceipt(${JSON.stringify(path.join(directory, "generations.jsonl"))}, 'written');` : ""}
const release = inject('releaseFile');
if (release) await new Promise(resolve => {
const check = () => {if(fs.existsSync(release)){clearInterval(poll);resolve();}};

View file

@ -17,7 +17,12 @@ import { resolveVitestSpawnParams, spawnWatchedVitestProcess } from "../../scrip
import { createVitestProcessCompletion } from "../../scripts/vitest-process-group.mts";
import { resolveRuntimeWorkerArgv } from "../../src/infra/runtime-worker-url.js";
import { resolveTestNodeExecPath } from "../../src/test-utils/node-process.js";
import { waitForFixtureFile } from "../helpers/process-wait.js";
import {
fixtureReceiptClientSource,
openFixtureReceiptChannel,
type FixtureReceiptChannel,
} from "../helpers/fixture-receipts.js";
import { awaitGateBeforeSettlement, createDeferred, withinTest } from "../helpers/promise.js";
import { fixturePreloadArgs } from "./fixtures/ci-fixture-runtime.cjs";
import { copyCompiledFsSafeRuntimeFixture } from "./fs-safe-package.test-support.js";
import {
@ -26,6 +31,7 @@ import {
} from "./vitest-worker-artifacts.prepared.test-support.js";
import {
createWorkerArtifactTest,
fixtureFileBeforeSettlement,
preparationClient,
workerBorrowingProbe,
workerProbe,
@ -37,6 +43,11 @@ const preparedCompiler = createPreparedWorkerCompiler();
const it = createWorkerArtifactTest(preparedCompiler.env);
it.beforeAll(() => preparedCompiler.prepare());
it.afterAll(() => preparedCompiler.cleanup());
let receipts: FixtureReceiptChannel;
it.beforeAll(async () => {
receipts = await openFixtureReceiptChannel();
});
it.afterAll(() => receipts.close());
const compilerModule = "scripts/lib/vitest-worker-run.mts";
const compilerEntry = "scripts/lib/vitest-worker-compiler.mts";
@ -947,11 +958,11 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
it.for(["cancel", "owner disconnect"])(
"joins actual borrowers after %s before deleting artifacts",
(action, { workerArtifacts }) =>
(action, { workerArtifacts, signal }) =>
workerArtifacts.fixtureLifetime.run(async () => {
const { startBorrower } = workerArtifacts.createFixtureCommands();
const directory = workerArtifacts.fixtureDirectory();
const { config } = workerProbe(directory, true);
const { config } = workerProbe(directory, true, "compiled", receipts.endpoint);
const owner = workerArtifacts.createWorkerRun();
// Node parent-side child.disconnect() can omit ChildProcess.close. Close the
// fixture endpoint so the owner receives EOF and retains its real join contract.
@ -972,7 +983,10 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
);
try {
const observed = path.join(directory, "generations.jsonl");
await waitForFixtureFile(observed, handle.completion);
await withinTest(
fixtureFileBeforeSettlement(receipts, observed, handle.completion),
signal,
);
const generation = JSON.parse(fs.readFileSync(observed, "utf8").trim());
expect(fs.existsSync(new URL(generation))).toBe(true);
if (action === "cancel") {
@ -980,7 +994,7 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
} else {
handle.child.send("fixture-disconnect");
}
const result = await handle.result;
const result = await withinTest(handle.result, signal);
expect(result.code).not.toBe(0);
if (action === "owner disconnect") {
expect(result.stderr).toContain("owner disconnected");
@ -996,7 +1010,7 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
it.runIf(process.platform !== "win32")(
"retains artifacts after an uncertain join and waits for the surviving borrower",
({ workerArtifacts }) =>
({ workerArtifacts, signal }) =>
workerArtifacts.fixtureLifetime.run(async () => {
const { observeChild } = workerArtifacts.createFixtureCommands();
const directory = workerArtifacts.fixtureDirectory();
@ -1016,7 +1030,7 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
process.disconnect();
}});
process.channel.ref();
fs.writeFileSync(process.argv[2],'ready');
process.send('fixture-ready');
`,
);
const clients = ["first", "second"].map((name) => {
@ -1025,6 +1039,12 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
detached: true,
stdio: ["ignore", "pipe", "pipe", "ipc"],
});
const readySignal = createDeferred();
child.on("message", (message: unknown) => {
if (message === "fixture-ready") {
readySignal.resolve();
}
});
const closed = new Promise<void>((resolve) => {
child.once("close", () => resolve());
});
@ -1045,14 +1065,23 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
}),
),
);
return { child, ready, closed, completion };
return { child, ready, readySignal: readySignal.promise, closed, completion };
});
try {
await Promise.all(
clients.map((client) => waitForFixtureFile(client.ready, client.completion)),
await withinTest(
Promise.all(
clients.map((client) =>
awaitGateBeforeSettlement(
client.readySignal,
client.completion,
`Child exited before writing ${client.ready}`,
),
),
),
signal,
);
clients[0]!.child.send("finish");
await expect(clients[0]!.completion).rejects.toThrow(
await expect(withinTest(clients[0]!.completion, signal)).rejects.toThrow(
"injected process-group join failure",
);
let disposed = false;
@ -1063,9 +1092,11 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
expect(disposed).toBe(false);
expect(clients[1]!.child.exitCode).toBeNull();
clients[1]!.child.send("finish");
await clients[1]!.completion;
await withinTest(clients[1]!.completion, signal);
expect(fs.readFileSync(clients[1]!.ready + ".read", "utf8")).toBe("read");
await expect(disposal).rejects.toThrow("injected process-group join failure");
await expect(withinTest(disposal, signal)).rejects.toThrow(
"injected process-group join failure",
);
expect(fs.existsSync(artifact)).toBe(true);
} finally {
for (const { child } of clients) {
@ -1195,6 +1226,7 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
it("keeps watch launches on live source across dependency edits", ({
workerArtifacts,
onTestFinished,
signal,
}) =>
workerArtifacts.fixtureLifetime.run(async () => {
const directory = workerArtifacts.fixtureDirectory();
@ -1209,6 +1241,7 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
directory,
"watch.test.ts",
`
${fixtureReceiptClientSource(receipts.endpoint)}
import {execFileSync} from 'node:child_process';
import fs from 'node:fs';
import {it,expect} from 'vitest';
@ -1222,15 +1255,20 @@ if (process.argv[1]?.endsWith("vitest-worker-compiler.mts")) {
const actual=execFileSync(process.execPath,['--import','tsx',${JSON.stringify(dependency)}],{encoding:'utf8'}).trim();
expect(actual).toBe(value);
fs.writeFileSync(${JSON.stringify(observed)},actual);
sendReceipt(${JSON.stringify(observed)},actual);
});
`,
);
const reporter = writeFixture(
directory,
"watch-reporter.mjs",
`import fs from 'node:fs';
`${fixtureReceiptClientSource(receipts.endpoint)}
import fs from 'node:fs';
export default class {
onWatcherStart() { fs.writeFileSync(${JSON.stringify(watchReady)}, 'ready'); }
onWatcherStart() {
fs.writeFileSync(${JSON.stringify(watchReady)}, 'ready');
sendReceipt(${JSON.stringify(watchReady)}, 'written');
}
}`,
);
const config = writeFixture(
@ -1259,13 +1297,16 @@ export default class {
await handle.completion;
});
try {
await Promise.all([
waitForFixtureFile(observed, handle.completion, "first"),
waitForFixtureFile(watchReady, handle.completion),
]);
const rerun = waitForFixtureFile(observed, handle.completion, "second");
await withinTest(
Promise.all([
fixtureFileBeforeSettlement(receipts, observed, handle.completion, "first"),
fixtureFileBeforeSettlement(receipts, watchReady, handle.completion),
]),
signal,
);
const rerun = fixtureFileBeforeSettlement(receipts, observed, handle.completion, "second");
fs.writeFileSync(dependency, 'export const value: string = "second"; console.log(value);');
await rerun;
await withinTest(rerun, signal);
expect(output).not.toContain("[vitest-workers] prepared");
} catch (error) {
console.error(output);