fix(test): drain Codex fixture workers and diagnose teardown stalls (#151225)

Join canonical agent and shared-state SQLite cleanup in Codex fixtures and drain reusable disk-budget workers at file teardown. Concurrent drains share one completion so accepted scans finish and later worker reuse stays intact. Runtime lifecycle cleanup uses the same owner.

Add bounded, redacted native-worker diagnostics for extension-fork teardown stalls while preserving the existing failure deadline and Vitest termination and exit-joining ownership. No dependency patch, retry, schema, or public configuration change.

Exact-head CI run 35289664323 passed with zero pending checks. Native-resource regression, real-process shutdown scenarios, and 20 Linux pressure runs verify cleanup and unchanged failure outcomes. The original intermittent timeout was not reproduced verbatim.

Related: #150819. The voice-call fixture repair remains in #151187.
This commit is contained in:
Peter Steinberger 2026-09-17 17:45:52 -07:00 • committed by GitHub
parent 1e6ca9fe6e
commit db98ebc771
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 709 additions and 93 deletions

View file

@ -74,6 +74,8 @@ const config = {
".agents/skills/**/scripts/**/*.{js,mjs,cjs,ts,mts,cts}!",
"scripts/**/*.{test,spec}.{js,mjs,cjs,ts,mts,cts}!",
"test/**/*.{test,spec}.{js,mjs,cjs,ts,mts,cts}!",
// CLI subprocess fixtures consume the shared native-report collector.
"src/cli/cli-process-child.test-helpers.test.ts!",
// Core bootstrap packaging consumes the scripts' dist-import scanner.
"src/gateway/worker-environments/node-bootstrap-artifact.ts!",
"src/plugin-sdk/api-baseline.ts!",
@ -86,6 +88,7 @@ const config = {
"skills/**/*.{js,mjs,cjs,ts,mts,cts}!",
"scripts/**/*.{js,mjs,cjs,ts,mts,cts}!",
"test/**/*.{js,mjs,cjs,ts,mts,cts}!",
"src/cli/cli-process-child.test-helpers{,.test}.ts!",
"src/gateway/worker-environments/node-bootstrap-artifact.ts!",
"src/plugin-sdk/api-baseline.ts!",
],

View file

@ -178,6 +178,10 @@ This is home isolation, not a filesystem sandbox: explicit absolute paths,
`os.userInfo()` account lookup, children with stripped or replaced home variables,
and intentionally real-home live execution remain outside its protection.
Codex app-server fixtures await agent and shared-state SQLite drainage between
cases. Their file teardown drains the shared disk-budget scan worker, preserving
reuse during the file and releasing it before isolated fork shutdown.
- `src/test-utils/openclaw-test-state.ts`: use from Vitest when a test needs an isolated `HOME`, `OPENCLAW_STATE_DIR`, `OPENCLAW_CONFIG_PATH`, config fixture, workspace, agent dir, or auth-profile store.
- `pnpm test:env-mutations:report`: non-blocking report of tests/harnesses that mutate `HOME`, `OPENCLAW_STATE_DIR`, `OPENCLAW_CONFIG_PATH`, `OPENCLAW_WORKSPACE_DIR`, or related env keys directly. Use it to find migration candidates for the shared test-state helper.
- `test/helpers/openclaw-test-instance.ts`: process-level E2E tests needing a running Gateway, CLI env, log capture, and cleanup in one place.
@ -207,6 +211,19 @@ Redaction is unconditional and affects diagnostic output, not assertion behavior
Test console capture is outside this boundary; tests must still avoid logging
credentials directly.
Configured extension fork projects use the `openclaw-forks` diagnostic adapter
around Vitest's native fork transport. If the existing stop deadline fails while
the child remains alive, the adapter spends at most two additional seconds
collecting a Node report before native termination and pipe cleanup. The timeout
remains a test failure. The report distinguishes a missing stop acknowledgement
from a stall after acknowledgement and includes native stacks, libuv handles, and
worker-thread reports. Environment variables, command arguments, and socket
endpoints are omitted. A blocked event loop can prevent signal reporting; that
case explicitly reports that no complete report was captured.
Successful shutdown remains quiet. An explicit `--pool=forks` selects Vitest's
built-in pool and bypasses this adapter.
## JSON reports across native processes
For a multi-project or chunked run, explicitly request native JSON with an output

View file

@ -23,9 +23,14 @@ import {
resolveStorePath,
upsertSessionEntry,
} from "openclaw/plugin-sdk/session-store-runtime";
import { closeOpenClawAgentDatabasesForTest } from "openclaw/plugin-sdk/sqlite-runtime-testing";
import {
closeOpenClawAgentDatabasesAsync,
closeOpenClawAgentDatabasesForTest,
closeOpenClawStateDatabaseAsync,
drainSessionDiskBudgetWorkers,
} from "openclaw/plugin-sdk/sqlite-runtime-testing";
import { resolvePreferredOpenClawTmpDir } from "openclaw/plugin-sdk/temp-path";
import { afterEach, beforeEach, expect, vi } from "vitest";
import { afterAll, afterEach, beforeEach, expect, vi } from "vitest";
import { defaultCodexAppInventoryCache } from "./app-inventory-cache.js";
import { CodexAppServerClient } from "./client.js";
import {
@ -687,6 +692,8 @@ export function createRuntimeDynamicTool(name: string): RuntimeDynamicToolForTes
}
export function setupRunAttemptTestHooks(): void {
afterAll(drainSessionDiskBudgetWorkers);
beforeEach(async () => {
// Direct runtime tests supply the plugin root normally owned by loader registration.
setManagedCodexPluginRoot(fileURLToPath(new URL("../../", import.meta.url)));
@ -741,7 +748,9 @@ export function setupRunAttemptTestHooks(): void {
for (const owner of seededSessionOwnersForTest.splice(0)) {
await deleteSessionEntry(owner);
}
await closeOpenClawAgentDatabasesAsync();
closeOpenClawAgentDatabasesForTest();
await closeOpenClawStateDatabaseAsync();
await fs.rm(tempDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 50 });
});
}

View file

@ -1,6 +1,13 @@
import { createHook } from "node:async_hooks";
import { setImmediate } from "node:timers/promises";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import type { ModelCompatConfig } from "openclaw/plugin-sdk/provider-model-types";
import { describe, expect, it } from "vitest";
import {
closeOpenClawAgentDatabasesAsync,
closeOpenClawStateDatabaseAsync,
drainSessionDiskBudgetWorkers,
} from "openclaw/plugin-sdk/sqlite-runtime-testing";
import { beforeAll, describe, expect, it } from "vitest";
import {
createStartedThreadHarness,
createTestParams,
@ -10,6 +17,38 @@ import {
turnStartResult,
} from "./run-attempt-test-harness.js";
beforeAll(() => {
let allocatedWorkers = 0;
const pendingWorkers = new Map<number, string>();
const observer = createHook({
init(id, type) {
if (type === "WORKER") {
allocatedWorkers++;
pendingWorkers.set(id, new Error("WORKER allocated here").stack ?? "WORKER");
}
},
destroy(id) {
pendingWorkers.delete(id);
},
}).enable();
// The scan pool is shared across cases; verify its file owner after all fixture cleanup.
return async () => {
try {
await setImmediate();
expect(allocatedWorkers).toBeGreaterThan(0);
expect(
pendingWorkers.size,
`WORKER resources surviving Codex fixture teardown:\n${[...pendingWorkers.values()].join("\n")}`,
).toBe(0);
} finally {
observer.disable();
await drainSessionDiskBudgetWorkers();
await closeOpenClawAgentDatabasesAsync();
await closeOpenClawStateDatabaseAsync();
}
};
});
setupRunAttemptTestHooks();
describe("Codex reasoning effort across completed turns", () => {

View file

@ -0,0 +1,72 @@
import fs from "node:fs";
import { isRecord } from "../../packages/normalization-core/src/record-coerce.ts";
export const NODE_DIAGNOSTIC_REPORT_GRACE_MS = 2_000;
type NodeDiagnosticReport = {
threadId?: number;
javascriptStack: Record<string, unknown>;
nativeStack: unknown[];
libuv: unknown[];
workers: NodeDiagnosticReport[];
};
function projectDiagnosticReport(report: unknown): NodeDiagnosticReport | undefined {
if (
!isRecord(report) ||
!isRecord(report.javascriptStack) ||
!Array.isArray(report.nativeStack) ||
!Array.isArray(report.libuv)
) {
return undefined;
}
// Reports also contain argv, environment, host, and network metadata; never log those sections.
const threadId = isRecord(report.header) ? report.header.threadId : undefined;
return {
...(typeof threadId === "number" ? { threadId } : {}),
javascriptStack: report.javascriptStack,
nativeStack: report.nativeStack,
libuv: report.libuv.map((handle) => {
if (!isRecord(handle)) {
return handle;
}
// Node's network exclusion flag retains socket and named-pipe endpoints.
const { localEndpoint: _local, remoteEndpoint: _remote, ...execution } = handle;
return execution;
}),
workers: Array.isArray(report.workers)
? report.workers.map(projectDiagnosticReport).filter((worker) => worker !== undefined)
: [],
};
}
export function collectNodeDiagnosticReport(
reportPath: string,
minimumWaitMs = 0,
): Promise<string> {
const startedAt = performance.now();
return new Promise((resolve) => {
const finish = (report: string) => {
clearInterval(poll);
clearTimeout(deadline);
resolve(report);
};
const poll = setInterval(() => {
try {
const report = projectDiagnosticReport(JSON.parse(fs.readFileSync(reportPath, "utf8")));
if (report && performance.now() - startedAt >= minimumWaitMs) {
finish(JSON.stringify(report, null, 2));
}
} catch {
// Node writes directly to the report file; it may not exist or be complete yet.
}
}, 50);
const deadline = setTimeout(
() =>
finish(
`No complete Node diagnostic report captured within ${NODE_DIAGNOSTIC_REPORT_GRACE_MS}ms.`,
),
NODE_DIAGNOSTIC_REPORT_GRACE_MS,
);
});
}

View file

@ -2,19 +2,20 @@
// deadlock guard each, and failures that always carry the child's own output.
import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
import { once } from "node:events";
import fs from "node:fs";
import path from "node:path";
import { fileURLToPath } from "node:url";
import { toErrorObject } from "@openclaw/normalization-core/error-coercion";
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import { afterEach } from "vitest";
import {
collectNodeDiagnosticReport,
NODE_DIAGNOSTIC_REPORT_GRACE_MS as REPORT_GRACE_MS,
} from "../../scripts/lib/node-diagnostic-report.mts";
import { resolveVitestNodeArgs } from "../../scripts/lib/vitest-process-env.mts";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import { DEFAULT_VITEST_TEST_TIMEOUT_MS } from "../../test/vitest/vitest.timeouts.js";
const OUTPUT_TAIL_CHARS = 8_000;
const DIAGNOSTIC_GRACE_MS = 200;
const REPORT_GRACE_MS = 2_000;
const reportDirs = useAutoCleanupTempDirTracker(afterEach);
const diagnosticPreload = fileURLToPath(
new URL("./cli-process-diagnostics.test-support.cjs", import.meta.url),
@ -24,69 +25,6 @@ function withoutDiagnosticReadiness(stderr: string): string {
return stderr.replace(/^\[cli-process-diagnostics\] ready pid=\d+\r?\n/gmu, "");
}
type CliProcessReport = {
threadId?: number;
javascriptStack: Record<string, unknown>;
nativeStack: unknown[];
libuv: unknown[];
workers: CliProcessReport[];
};
function projectDiagnosticReport(report: unknown): CliProcessReport | undefined {
if (
!isRecord(report) ||
!isRecord(report.javascriptStack) ||
!Array.isArray(report.nativeStack) ||
!Array.isArray(report.libuv)
) {
return undefined;
}
// Reports also contain argv, environment, host, and network metadata; never log those sections.
const threadId = isRecord(report.header) ? report.header.threadId : undefined;
return {
...(typeof threadId === "number" ? { threadId } : {}),
javascriptStack: report.javascriptStack,
nativeStack: report.nativeStack,
libuv: report.libuv.map((handle) => {
if (!isRecord(handle)) {
return handle;
}
// Node's network exclusion flag retains socket and named-pipe endpoints.
const { localEndpoint: _local, remoteEndpoint: _remote, ...execution } = handle;
return execution;
}),
workers: Array.isArray(report.workers)
? report.workers.map(projectDiagnosticReport).filter((worker) => worker !== undefined)
: [],
};
}
function collectDiagnosticReport(reportPath: string): Promise<string> {
const startedAt = performance.now();
return new Promise((resolve) => {
const finish = (report: string) => {
clearInterval(poll);
clearTimeout(deadline);
resolve(report);
};
const poll = setInterval(() => {
try {
const report = projectDiagnosticReport(JSON.parse(fs.readFileSync(reportPath, "utf8")));
// Preserve the existing grace for the child's JS diagnostic and trailing pipe output.
if (report && performance.now() - startedAt >= DIAGNOSTIC_GRACE_MS) {
finish(JSON.stringify(report, null, 2));
}
} catch {
// Node writes directly to the report file; it may not exist or be complete yet.
}
}, 50);
const deadline = setTimeout(
() => finish(`No complete Node diagnostic report captured within ${REPORT_GRACE_MS}ms.`),
REPORT_GRACE_MS,
);
});
}
function releaseCliProcessChild(child: ChildProcessWithoutNullStreams): string[] {
const failures: string[] = [];
for (const release of [
@ -304,7 +242,11 @@ export async function runCliProcessChild(params: {
try {
if (child.kill("SIGUSR2")) {
diagnosticRequest = `SIGUSR2 requested; report grace<=${REPORT_GRACE_MS}ms`;
void collectDiagnosticReport(path.join(reportDir, "diagnostic.json")).then(finish);
// Preserve the child's JS diagnostic and trailing pipe output grace.
void collectNodeDiagnosticReport(
path.join(reportDir, "diagnostic.json"),
DIAGNOSTIC_GRACE_MS,
).then(finish);
return;
}
} catch (error) {

View file

@ -1,22 +1,55 @@
import path from "node:path";
import { resolveRuntimeWorkerUrl } from "../../infra/runtime-worker-url.js";
import { WorkerTaskPool } from "../../infra/worker-task-pool.js";
import { resolveGlobalSingleton } from "../../shared/global-singleton.js";
import type { SessionPhysicalDiskUsage } from "./disk-budget-files.js";
const measurements = new WorkerTaskPool<string, SessionPhysicalDiskUsage>({
workerUrl: resolveRuntimeWorkerUrl({
currentModuleUrl: import.meta.url,
sourceWorkerName: "disk-budget.worker",
distWorkerPath: "config/sessions/disk-budget.worker.js",
const measurements = resolveGlobalSingleton<{
pool: WorkerTaskPool<string, SessionPhysicalDiskUsage>;
pending: Set<Promise<SessionPhysicalDiskUsage>>;
draining?: Promise<void>;
}>(
Symbol.for("openclaw.sessionDiskBudgetWorkers"),
() => ({
pool: new WorkerTaskPool<string, SessionPhysicalDiskUsage>({
workerUrl: resolveRuntimeWorkerUrl({
currentModuleUrl: import.meta.url,
sourceWorkerName: "disk-budget.worker",
distWorkerPath: "config/sessions/disk-budget.worker.js",
}),
// Share one scan worker so concurrent stores cannot multiply filesystem scan heaps.
maxWorkers: 1,
}),
pending: new Set<Promise<SessionPhysicalDiskUsage>>(),
}),
// Share one scan worker so concurrent stores cannot multiply filesystem scan heaps.
maxWorkers: 1,
});
() => drainSessionDiskBudgetWorkers(),
);
/** Join admitted scans before retiring workers; later measurements reuse the pool. */
export function drainSessionDiskBudgetWorkers(): Promise<void> {
// A second teardown must join this rotation, not retire its successor worker.
return (measurements.draining ??= Promise.resolve()
.then(async () => {
while (measurements.pending.size > 0) {
await Promise.allSettled(measurements.pending);
}
await measurements.pool.rotate();
})
.finally(() => {
measurements.draining = undefined;
}));
}
/** Measures physical session artifacts without running per-file synchronous work on the caller. */
export async function measureSessionPhysicalDiskUsage(
storePath: string,
): Promise<SessionPhysicalDiskUsage> {
// Capture the locator before queueing; only four totals cross back to the caller.
return measurements.run(path.resolve(storePath), {});
const pending = measurements.pool.run(path.resolve(storePath), {});
measurements.pending.add(pending);
try {
return await pending;
} finally {
measurements.pending.delete(pending);
}
}

View file

@ -6,7 +6,9 @@ import { Worker } from "node:worker_threads";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { WorkerTaskPool } from "../../infra/worker-task-pool.js";
import { createDeferredCore } from "../../shared/deferred.js";
import { drainGlobalSingletonLifecycleState } from "../../shared/global-singleton.js";
import { withTestDir } from "../../test-helpers/temp-dir.js";
import { drainSessionDiskBudgetWorkers } from "./disk-budget-runtime.js";
import {
hasRetainedSessionTranscriptArchives,
measureSessionPhysicalDiskUsage,
@ -32,7 +34,94 @@ async function addSessionArtifacts(directory: string, index: number): Promise<vo
}
describe("physical session disk usage", () => {
it("reports scan overload before pruning archives and recovers after queued scans drain", async () => {
it("joins a measurement admitted after drainage starts before retiring its worker", async () => {
await withTestDir({ prefix: "openclaw-disk-usage-drain-admission-" }, async (directory) => {
const storePath = path.join(directory, "openclaw-agent.sqlite");
await fs.writeFile(storePath, Buffer.alloc(321));
const release = createDeferredCore();
const spy = vi
.spyOn(WorkerTaskPool.prototype, "run")
.mockImplementationOnce(function (this: WorkerTaskPool<unknown, unknown>, input, options) {
spy.mockRestore();
return this.run(async () => {
await release.promise;
return input;
}, options);
});
let completed = 0;
const first = measureSessionPhysicalDiskUsage(storePath).then(() => completed++);
const drainage = drainSessionDiskBudgetWorkers();
const late = measureSessionPhysicalDiskUsage(storePath).then(() => completed++);
try {
release.resolve();
await drainage;
expect(completed).toBe(2);
expect(workers).toHaveLength(1);
expect(workers[0]?.threadId).toBe(-1);
} finally {
release.resolve();
spy.mockRestore();
await Promise.allSettled([first, late, drainage]);
await drainSessionDiskBudgetWorkers();
}
});
});
it.each([
{ owner: "file teardown", drain: drainSessionDiskBudgetWorkers },
{ owner: "runtime lifecycle cleanup", drain: drainGlobalSingletonLifecycleState },
])(
"coalesces $owner with another teardown while a runtime scan awaits retirement",
async ({ drain }) => {
await withTestDir({ prefix: "openclaw-disk-usage-concurrent-drain-" }, async (directory) => {
const storePath = path.join(directory, "openclaw-agent.sqlite");
await fs.writeFile(storePath, Buffer.alloc(321));
await measureSessionPhysicalDiskUsage(storePath);
const retiringWorker = workers[0]!;
const entered = createDeferredCore();
const release = createDeferredCore();
const terminate = retiringWorker.terminate.bind(retiringWorker);
const retirement = vi
.spyOn(Worker.prototype, "terminate")
.mockImplementationOnce(async () => {
const code = await terminate();
entered.resolve();
await release.promise;
return code;
});
const firstDrain = drain();
let secondDrain: Promise<void> | undefined;
let measurement: ReturnType<typeof measureSessionPhysicalDiskUsage> | undefined;
try {
await entered.promise;
let measured = false;
measurement = measureSessionPhysicalDiskUsage(storePath).then((usage) => {
measured = true;
return usage;
});
secondDrain = drainSessionDiskBudgetWorkers();
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
expect(measured).toBe(false);
expect(retiringWorker.threadId).toBe(-1);
release.resolve();
await Promise.all([firstDrain, secondDrain]);
await expect(measurement).resolves.toMatchObject({ totalBytes: 321 });
expect(retirement).toHaveBeenCalledTimes(1);
expect(workers).toHaveLength(2);
expect(workers[1]?.threadId).toBeGreaterThan(0);
} finally {
release.resolve();
await Promise.allSettled([firstDrain, secondDrain, measurement]);
retirement.mockRestore();
await drainSessionDiskBudgetWorkers();
}
});
},
);
it("reports scan overload and retires queued scans before reusing workers", async () => {
await withTestDir({ prefix: "openclaw-disk-usage-pressure-" }, async (directory) => {
const storePath = path.join(directory, "openclaw-agent.sqlite");
const archivePath = path.join(directory, "old.jsonl.deleted.2026-01-01T00-00-00.000Z.zst");
@ -49,9 +138,14 @@ describe("physical session disk usage", () => {
return input;
}, options);
});
let settledScans = 0;
const accepted = Array.from({ length: 128 }, () =>
measureSessionPhysicalDiskUsage(storePath),
measureSessionPhysicalDiskUsage(storePath).then((usage) => {
settledScans++;
return usage;
}),
);
let drainage: Promise<void> | undefined;
let reported: unknown;
const excess = measureSessionPhysicalDiskUsage(storePath).catch((error: unknown) => {
reported = error;
@ -65,7 +159,19 @@ describe("physical session disk usage", () => {
pruneSessionTranscriptArchivesToHighWater({ storePath, highWaterBytes: 321 }),
).rejects.toMatchObject({ code: "overloaded" });
expect((await fs.stat(archivePath)).size).toBe(100);
let drained = false;
drainage = drainGlobalSingletonLifecycleState().then(() => {
drained = true;
});
await new Promise<void>((resolve) => {
setImmediate(resolve);
});
expect(drained).toBe(false);
release.resolve();
await drainage;
expect(settledScans).toBe(128);
expect(workers).toHaveLength(1);
expect(workers[0]?.threadId).toBe(-1);
const usage = {
databaseMainBytes: 321,
databaseWalBytes: 0,
@ -77,11 +183,13 @@ describe("physical session disk usage", () => {
pruneSessionTranscriptArchivesToHighWater({ storePath, highWaterBytes: 321 }),
).resolves.toMatchObject({ removedFiles: 1, usage: { totalBytes: 321 } });
await expect(fs.stat(archivePath)).rejects.toMatchObject({ code: "ENOENT" });
expect(workers).toHaveLength(1);
expect(workers).toHaveLength(2);
expect(workers[1]?.threadId).toBeGreaterThan(0);
} finally {
release.resolve();
spy.mockRestore();
await Promise.allSettled([...accepted, excess]);
await drainage;
}
});
});

View file

@ -13,6 +13,7 @@ export async function appendSqliteSessionTranscriptEventForTest(
await appendTranscriptEvent(params, params.event);
}
export { drainSessionDiskBudgetWorkers } from "../config/sessions/disk-budget-runtime.js";
export { formatSqliteSessionFileMarker } from "../config/sessions/legacy-sqlite-marker.js";
export {
appendSqliteTrajectoryRuntimeEvents,

View file

@ -0,0 +1,258 @@
import fs from "node:fs";
import path from "node:path";
import { afterEach, expect, it } from "vitest";
import { useAutoCleanupTempDirTracker } from "../helpers/temp-dir.js";
import { runVitestShutdownCommand } from "../helpers/vitest-shutdown-command.js";
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
const repoRoot = path.resolve(import.meta.dirname, "../..");
const posixNodeIt = it.skipIf(process.platform === "win32" || Boolean(process.versions.bun));
const teardownTimeoutError = "[vitest-pool-runner]: Timeout waiting for worker to respond";
posixNodeIt.for(["normal", "missing-ack", "after-ack", "blocked-after-ack"] as const)(
"retains native fork cleanup and captures only stalled teardown (%s)",
{ timeout: 180_000 },
async (mode, { signal }) => {
const root = tempDirs.make("openclaw-pool-diagnostics-");
const home = path.join(root, "home");
const tmp = path.join(root, "tmp");
fs.mkdirSync(home);
fs.mkdirSync(tmp);
fs.symlinkSync(
path.join(repoRoot, "node_modules"),
path.join(root, "node_modules"),
"junction",
);
fs.writeFileSync(path.join(root, "package.json"), '{"type":"module","private":true}');
const receipt = path.join(root, "deadline.json");
const diagnosticReceipt = path.join(root, "diagnostic-deadline.json");
const preload = path.join(root, "hold-teardown.cjs");
fs.writeFileSync(
preload,
`
const { subscribe } = require("node:diagnostics_channel");
const fs = require("node:fs");
const mode = ${JSON.stringify(mode)};
const schedule = globalThis.setTimeout;
const cancel = globalThis.clearTimeout;
const deadlines = new Map();
let stopDeadlineInvoked = false;
const diagnosticDeadline = { delay: 0, scheduled: 0, fired: 0 };
const recordDiagnosticDeadline = () => fs.writeFileSync(${JSON.stringify(diagnosticReceipt)}, JSON.stringify(diagnosticDeadline));
globalThis.setTimeout = (callback, delay, ...args) => {
if (stopDeadlineInvoked && delay === 2000) {
diagnosticDeadline.delay = delay;
diagnosticDeadline.scheduled++;
recordDiagnosticDeadline();
return schedule(() => {
diagnosticDeadline.fired++;
recordDiagnosticDeadline();
callback(...args);
}, delay);
}
if (delay !== 60000) return schedule(callback, delay, ...args);
const invoke = () => callback(...args);
const timer = schedule(() => { deadlines.delete(timer); invoke(); }, delay);
deadlines.set(timer, invoke);
return timer;
};
globalThis.clearTimeout = timer => { deadlines.delete(timer); return cancel(timer); };
const isFork = arg => typeof arg === "string" && arg.replaceAll("\\\\", "/").endsWith("/vitest/dist/workers/forks.js");
if (isFork(process.argv[1]) && process.send) {
const send = process.send;
const held = () => send.call(process, { fixtureTeardownHeld: true }, () => {
if (mode === "blocked-after-ack") Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0);
});
if (mode === "missing-ack") {
const emit = process.emit;
process.emit = function(event, message, ...args) {
if (event === "message" && message?.__vitest_worker_request__ === true && message.type === "stop") {
held();
return true;
}
return emit.call(this, event, message, ...args);
};
} else if (mode === "after-ack" || mode === "blocked-after-ack") {
process.send = function(message, ...args) {
if (message?.__vitest_worker_response__ === true && message.type === "stopped" && message.willExit === true) {
// Hold the transport's explicit process.exit callback after its acknowledgement flushes.
args[args.length - 1] = error => { if (error) throw error; held(); };
}
return send.call(this, message, ...args);
};
}
}
subscribe("child_process", ({ process: child }) => {
let selected = false;
child.once("spawn", () => { selected = child.spawnargs.some(isFork); });
child.on("message", message => {
if (!selected || message?.fixtureTeardownHeld !== true) return;
setImmediate(() => {
fs.writeFileSync(${JSON.stringify(receipt)}, JSON.stringify({ liveDeadlines: deadlines.size, delay: 60000 }));
if (deadlines.size !== 1) throw new Error("expected one live Vitest stop deadline");
const [timer, invoke] = deadlines.entries().next().value;
cancel(timer);
deadlines.delete(timer);
stopDeadlineInvoked = true;
invoke();
});
});
});
`,
);
const workerReceipts = path.join(root, "workers.jsonl");
for (const filename of ["first.test.ts", "second.test.ts"]) {
fs.writeFileSync(
path.join(root, filename),
`
import fs from "node:fs";
import { once } from "node:events";
import { createServer } from "node:net";
import { Worker, isMainThread } from "node:worker_threads";
import { expect, it } from "vitest";
it("runs on the fork main thread with ready native handles", async () => {
expect(isMainThread).toBe(true);
const server = createServer();
server.listen(0, "127.0.0.1");
await once(server, "listening");
const worker = new Worker("require('node:worker_threads').parentPort.postMessage('ready'); setInterval(() => {}, 1000)", { eval: true });
await once(worker, "message");
fs.appendFileSync(${JSON.stringify(workerReceipts)}, JSON.stringify({ pid: process.pid, reportDirectory: process.report.directory }) + "\\n");
if (${JSON.stringify(mode)} === "normal") {
await worker.terminate();
await new Promise(resolve => server.close(resolve));
}
});
`,
);
}
const config = path.join(root, "vitest.config.ts");
const outcomeFile = path.join(root, "outcome.json");
const outcomes: unknown[] = [];
for (const useAdapter of [false, true]) {
fs.writeFileSync(workerReceipts, "");
for (const file of [receipt, diagnosticReceipt, outcomeFile]) {
fs.rmSync(file, { force: true });
}
fs.writeFileSync(
config,
`
import fs from "node:fs";
import { createExtensionDatabaseWorkersVitestConfig } from ${JSON.stringify(path.join(repoRoot, "test/vitest/vitest.extension-database-workers.config.ts"))};
const extension = createExtensionDatabaseWorkersVitestConfig({});
export default {
root: ${JSON.stringify(root)},
test: {
pool: ${useAdapter ? "extension.test.pool" : '"forks"'},
include: ["*.test.ts"],
isolate: false,
maxWorkers: 1,
fileParallelism: false,
fsModuleCache: false,
reporters: ["default", {
onTestRunEnd(modules, errors, reason) {
fs.writeFileSync(${JSON.stringify(outcomeFile)}, JSON.stringify({
reason,
errors: errors.map(error => error.message),
}));
},
}],
},
};
`,
);
const result = await runVitestShutdownCommand({
args: [
path.join(repoRoot, "scripts/run-vitest.mjs"),
"run",
"--config",
config,
"--root",
root,
"--configLoader",
"native",
],
cwd: root,
env: {
...process.env,
HOME: home,
USERPROFILE: home,
TMPDIR: tmp,
TMP: tmp,
TEMP: tmp,
CI: "1",
NODE_OPTIONS: `--require=${preload}`,
OPENCLAW_VITEST_FS_MODULE_CACHE_PATH: path.join(root, "cache"),
POOL_DIAGNOSTIC_FIXTURE_SECRET: "fixture-env-value-do-not-print",
},
signal,
});
const output = `${result.stdout}\n${result.stderr}`;
expect(output).toMatch(/2 passed/u);
const outcome = JSON.parse(fs.readFileSync(outcomeFile, "utf8"));
expect(outcome, output).toEqual({
// Vitest's reason reflects test assertions; unhandled teardown errors set the CLI exit.
reason: "passed",
errors: mode === "normal" ? [] : [teardownTimeoutError],
});
outcomes.push({ code: result.code, outcome });
const workers = fs
.readFileSync(workerReceipts, "utf8")
.trim()
.split("\n")
.map((line) => JSON.parse(line) as { pid: number; reportDirectory: string });
expect(new Set(workers.map(({ pid }) => pid)).size).toBe(1);
for (const { reportDirectory } of workers) {
if (reportDirectory) {
expect(fs.existsSync(reportDirectory)).toBe(false);
}
}
if (mode === "normal") {
expect(result.code, output).toBe(0);
expect(output).not.toContain("vitest-pool-diagnostics");
expect(output).not.toContain("Writing Node.js report");
continue;
}
expect(result.code, output).toBe(1);
expect(JSON.parse(fs.readFileSync(receipt, "utf8"))).toEqual({
liveDeadlines: 1,
delay: 60_000,
});
expect(output).toContain(teardownTimeoutError);
if (!useAdapter) {
expect(output).not.toContain("vitest-pool-diagnostics");
continue;
}
const report = output.match(
/\[vitest-pool-diagnostics\][^\n]*\n([\s\S]*?)\n\[\/vitest-pool-diagnostics\]/u,
)?.[1];
expect(report, output).toBeDefined();
expect(output).toContain(`stopAcknowledged=${mode !== "missing-ack"}`);
expect(output).not.toMatch(
/fixture-env-value-do-not-print|127\.0\.0\.1|localEndpoint|remoteEndpoint/u,
);
if (mode === "blocked-after-ack") {
expect(report).toBe("No complete Node diagnostic report captured within 2000ms.");
expect(JSON.parse(fs.readFileSync(diagnosticReceipt, "utf8"))).toEqual({
delay: 2_000,
scheduled: 1,
fired: 1,
});
continue;
}
expect(JSON.parse(report!)).toMatchObject({
nativeStack: expect.any(Array),
libuv: expect.arrayContaining([
expect.objectContaining({ type: "tcp", is_active: true, is_referenced: true }),
]),
workers: expect.arrayContaining([
expect.objectContaining({
threadId: expect.any(Number),
libuv: expect.arrayContaining([expect.objectContaining({ type: "timer" })]),
}),
]),
});
}
expect(outcomes[1]).toEqual(outcomes[0]);
},
);

View file

@ -43,6 +43,7 @@ import { createExtensionDatabaseWorkersVitestConfig } from "./vitest/vitest.exte
import { createExtensionImessageVitestConfig } from "./vitest/vitest.extension-imessage.config.ts";
import { createExtensionSlackVitestConfig } from "./vitest/vitest.extension-slack.config.ts";
import { createExtensionsVitestConfig } from "./vitest/vitest.extensions.config.ts";
import { diagnosticForksPool } from "./vitest/vitest.forks-pool.ts";
import { createGatewayMethodsIsolatedVitestConfig } from "./vitest/vitest.gateway-methods-isolated.config.ts";
import { createGatewayMethodsVitestConfig } from "./vitest/vitest.gateway-methods.config.ts";
import { createGatewayServerIsolatedVitestConfig } from "./vitest/vitest.gateway-server-isolated.config.ts";
@ -715,7 +716,7 @@ describe("projects vitest config", () => {
it("keeps Slack's real cooldown store in its forked project", () => {
const project = "test/vitest/vitest.extension-slack.config.ts";
expect(requireTestConfig(createExtensionSlackVitestConfig({})).pool).toBe("forks");
expect(requireTestConfig(createExtensionSlackVitestConfig({})).pool).toBe(diagnosticForksPool);
expect(
buildVitestRunPlans(["extensions/slack/src/monitor/presence-cooldown-store.test.ts"]).map(
(plan) => plan.config,
@ -740,7 +741,7 @@ describe("projects vitest config", () => {
expect(
fullSuiteVitestShards.find((shard) => shard.name === "extensions")?.projects,
).toContain(project);
expect(testConfig.pool).toBe("forks");
expect(testConfig.pool).toBe(diagnosticForksPool);
expect(testConfig.isolate).toBe(true);
expect(testConfig.include).toEqual([
...databaseWorkerExtensionTestRoots.map(
@ -761,7 +762,7 @@ describe("projects vitest config", () => {
const config = requireTestConfig(createExtensionDatabaseWorkersVitestConfig({}));
expect(buildVitestRunPlans([file]).map((plan) => plan.config)).toEqual([project]);
expect(config.include).toContain(file.replace(/^extensions\//u, ""));
expect(config.pool).toBe("forks");
expect(config.pool).toBe(diagnosticForksPool);
expect(config.isolate).toBe(true);
},
);
@ -801,7 +802,7 @@ describe("projects vitest config", () => {
}
const workerConfig = requireTestConfig(createExtensionDatabaseWorkersVitestConfig({}));
expect(workerConfig.include).toContain(`imessage/src/${basename}`);
expect(workerConfig.pool).toBe("forks");
expect(workerConfig.pool).toBe(diagnosticForksPool);
expect(workerConfig.isolate).toBe(true);
expect(requireTestConfig(createExtensionImessageVitestConfig({})).exclude).toContain(
`imessage/src/${basename}`,

View file

@ -56,6 +56,7 @@ import { createExtensionVoiceCallVitestConfig } from "./vitest/vitest.extension-
import { createExtensionWhatsAppVitestConfig } from "./vitest/vitest.extension-whatsapp.config.ts";
import { createExtensionZaloVitestConfig } from "./vitest/vitest.extension-zalo.config.ts";
import { createExtensionsVitestConfig } from "./vitest/vitest.extensions.config.ts";
import { diagnosticForksPool } from "./vitest/vitest.forks-pool.ts";
import { createGatewayClientVitestConfig } from "./vitest/vitest.gateway-client.config.ts";
import { createGatewayCoreVitestConfig } from "./vitest/vitest.gateway-core.config.ts";
import { createGatewayMethodsVitestConfig } from "./vitest/vitest.gateway-methods.config.ts";
@ -148,11 +149,12 @@ function expectThreadedIsolatedRunner(config: {
expect(testConfig.isolate).toBe(true);
expect(testConfig.runner).toBeUndefined();
}
function expectForkedNonIsolatedRunner(config: {
test?: { pool?: unknown; isolate?: unknown; runner?: unknown };
}) {
function expectForkedNonIsolatedRunner(
config: { test?: { pool?: unknown; isolate?: unknown; runner?: unknown } },
pool: "forks" | typeof diagnosticForksPool = "forks",
) {
const testConfig = requireTestConfig(config);
expect(testConfig.pool).toBe("forks");
expect(testConfig.pool).toBe(pool);
expect(testConfig.isolate).toBe(false);
expect(normalizeConfigPath(testConfig.runner)).toBe("test/non-isolated-runner.ts");
}
@ -763,7 +765,7 @@ describe("scoped vitest configs", () => {
});
it("serializes Slack extension files that share process globals", () => {
expectForkedNonIsolatedRunner(defaultExtensionSlackConfig);
expectForkedNonIsolatedRunner(defaultExtensionSlackConfig, diagnosticForksPool);
expect(requireTestConfig(defaultExtensionSlackConfig).fileParallelism).toBe(false);
});
@ -1084,7 +1086,10 @@ describe("scoped vitest configs", () => {
.every((pattern) => /\.test\.[cm]?[jt]sx?$/u.test(pattern)),
).toBe(true);
expect(projects.map((project) => project.name)).toEqual(names);
expect(projects.map((project) => project.pool)).toEqual(["threads", "forks"]);
expect(projects.map((project) => project.pool)).toEqual([
"threads",
root.startsWith("extensions/") ? diagnosticForksPool.name : "forks",
]);
expect(projects[0]?.setupFiles).toEqual(owner.test?.setupFiles);
expect(projects[0]?.maxWorkers).toBe(owner.test?.maxWorkers);
const workerConfig = await resolveConfig(

View file

@ -0,0 +1,120 @@
import { ChildProcess } from "node:child_process";
import { subscribe, unsubscribe } from "node:diagnostics_channel";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import { ForksPoolWorker, type PoolOptions, type WorkerRequest } from "vitest/node";
import { isRecord } from "../../packages/normalization-core/src/record-coerce.ts";
import { collectNodeDiagnosticReport } from "../../scripts/lib/node-diagnostic-report.mts";
const POOL_NAME = "openclaw-forks";
const MAX_REPORT_CHARS = 64 * 1_024;
class DiagnosticForksPoolWorker extends ForksPoolWorker {
override readonly name = POOL_NAME;
private child?: ChildProcess;
private reportDir?: string;
private stopRequested = false;
private stopAcknowledged = false;
private files: string[] = [];
private readonly project: PoolOptions["project"];
constructor(options: PoolOptions) {
super(options);
this.project = options.project;
}
override async start(): Promise<void> {
if (process.platform !== "win32" && !process.versions.bun) {
try {
// openclaw-temp-dir: allow pool-owned diagnostics outlive individual test hooks.
this.reportDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-vitest-report-"));
this.execArgv = [
...this.execArgv,
"--report-on-signal",
"--report-signal=SIGUSR2",
`--report-directory=${this.reportDir}`,
"--report-filename=diagnostic.json",
"--report-exclude-env",
"--report-exclude-network",
];
} catch {
// Missing diagnostic storage must not prevent a worker from starting.
}
}
const children: ChildProcess[] = [];
const observe = (message: unknown) => {
if (isRecord(message) && message.process instanceof ChildProcess) {
children.push(message.process);
}
};
// Vitest 5.0.0 + patches/vitest@5.0.0.patch (pnpm hash 5e7c1655) starts
// forks synchronously and calls stop() after its deadline or joined exit.
// Recheck that contract on upgrades; the transport remains Vitest-owned.
subscribe("child_process", observe);
let started: Promise<void>;
try {
started = super.start();
} finally {
unsubscribe("child_process", observe);
}
this.child = children.find((child) => child.spawnargs?.includes(this.entrypoint));
await started;
}
override send(message: WorkerRequest): void {
if (message.type === "stop") {
this.stopRequested = true;
} else if (message.type === "run" || message.type === "collect") {
this.files = message.context.files.map(({ filepath }) =>
path.relative(this.project.config.root, filepath),
);
}
super.send(message);
}
override waitForExit(): Promise<void> {
// The public pool hook runs after the transport acknowledges its graceful exit.
this.stopAcknowledged = true;
return super.waitForExit();
}
override async stop(): Promise<void> {
try {
const child = this.child;
// The runner calls stop after its existing deadline, or after a graceful
// child exit. Only a still-live transport needs a dump before forced cleanup.
if (this.stopRequested && child?.exitCode === null && child.signalCode === null) {
let report = "Node diagnostic report unavailable on this runtime or host.";
if (this.reportDir && child.kill("SIGUSR2")) {
report = await collectNodeDiagnosticReport(path.join(this.reportDir, "diagnostic.json"));
}
const boundedReport =
report.length > MAX_REPORT_CHARS
? `${report.slice(0, MAX_REPORT_CHARS)}\n[Node diagnostic report truncated]`
: report;
this.project.vitest.logger.error(
`[vitest-pool-diagnostics] pid=${child.pid} project=${JSON.stringify(this.project.name)} stopAcknowledged=${this.stopAcknowledged} files=${JSON.stringify(this.files)}\n${boundedReport}\n[/vitest-pool-diagnostics]`,
);
}
} catch {
this.project.vitest.logger.error("[vitest-pool-diagnostics] Node report capture failed.");
} finally {
await super.stop();
if (this.reportDir) {
try {
fs.rmSync(this.reportDir, { recursive: true, force: true });
} catch {
this.project.vitest.logger.error(
"[vitest-pool-diagnostics] Report directory cleanup failed.",
);
}
}
}
}
}
export const diagnosticForksPool = {
name: POOL_NAME,
createPoolWorker: (options: PoolOptions) => new DiagnosticForksPoolWorker(options),
};

View file

@ -1,6 +1,7 @@
// Vitest scoped config helper builds test configs for scoped file patterns.
import path from "node:path";
import { defineConfig, type ViteUserConfig } from "vitest/config";
import { diagnosticForksPool } from "./vitest.forks-pool.ts";
import { intersectIncludePatterns } from "./vitest.include-patterns.ts";
import {
loadPatternListFromEnv,
@ -250,7 +251,14 @@ export function createScopedVitestConfig(
...(resolvedScopedDir ? { dir: resolvedScopedDir } : {}),
include: scopedInclude,
exclude,
...(options?.pool ? { pool: options.pool } : {}),
...(options?.pool
? {
pool:
options.pool === "forks" && scopedDir === "extensions"
? diagnosticForksPool
: options.pool,
}
: {}),
...(options?.fileParallelism === undefined
? {}
: { fileParallelism: options.fileParallelism }),