diff --git a/config/knip.scripts-exports.config.ts b/config/knip.scripts-exports.config.ts index 9d3da48c71f1..7d52731c2ec2 100644 --- a/config/knip.scripts-exports.config.ts +++ b/config/knip.scripts-exports.config.ts @@ -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!", ], diff --git a/docs/reference/test/runner-internals.md b/docs/reference/test/runner-internals.md index f2e0a9b22562..aa80e4081731 100644 --- a/docs/reference/test/runner-internals.md +++ b/docs/reference/test/runner-internals.md @@ -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 diff --git a/extensions/codex/src/app-server/run-attempt-test-harness.ts b/extensions/codex/src/app-server/run-attempt-test-harness.ts index 2f565ef82ab9..19990ee10878 100644 --- a/extensions/codex/src/app-server/run-attempt-test-harness.ts +++ b/extensions/codex/src/app-server/run-attempt-test-harness.ts @@ -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 }); }); } diff --git a/extensions/codex/src/app-server/run-attempt.reasoning-effort.test.ts b/extensions/codex/src/app-server/run-attempt.reasoning-effort.test.ts index 374ba48fb12a..980c41254c53 100644 --- a/extensions/codex/src/app-server/run-attempt.reasoning-effort.test.ts +++ b/extensions/codex/src/app-server/run-attempt.reasoning-effort.test.ts @@ -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(); + 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", () => { diff --git a/scripts/lib/node-diagnostic-report.mts b/scripts/lib/node-diagnostic-report.mts new file mode 100644 index 000000000000..cd63a6a36990 --- /dev/null +++ b/scripts/lib/node-diagnostic-report.mts @@ -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; + 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 { + 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, + ); + }); +} diff --git a/src/cli/cli-process-child.test-helpers.ts b/src/cli/cli-process-child.test-helpers.ts index c8c00e22eace..bdf5659fe560 100644 --- a/src/cli/cli-process-child.test-helpers.ts +++ b/src/cli/cli-process-child.test-helpers.ts @@ -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; - 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 { - 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) { diff --git a/src/config/sessions/disk-budget-runtime.ts b/src/config/sessions/disk-budget-runtime.ts index b5934ae1b80d..a65488e9ae42 100644 --- a/src/config/sessions/disk-budget-runtime.ts +++ b/src/config/sessions/disk-budget-runtime.ts @@ -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({ - workerUrl: resolveRuntimeWorkerUrl({ - currentModuleUrl: import.meta.url, - sourceWorkerName: "disk-budget.worker", - distWorkerPath: "config/sessions/disk-budget.worker.js", +const measurements = resolveGlobalSingleton<{ + pool: WorkerTaskPool; + pending: Set>; + draining?: Promise; +}>( + Symbol.for("openclaw.sessionDiskBudgetWorkers"), + () => ({ + pool: new WorkerTaskPool({ + 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>(), }), - // 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 { + // 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 { // 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); + } } diff --git a/src/config/sessions/disk-budget.physical-usage.test.ts b/src/config/sessions/disk-budget.physical-usage.test.ts index e36e23a4c570..828fd1053f56 100644 --- a/src/config/sessions/disk-budget.physical-usage.test.ts +++ b/src/config/sessions/disk-budget.physical-usage.test.ts @@ -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 { - 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, 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 | undefined; + let measurement: ReturnType | undefined; + try { + await entered.promise; + let measured = false; + measurement = measureSessionPhysicalDiskUsage(storePath).then((usage) => { + measured = true; + return usage; + }); + secondDrain = drainSessionDiskBudgetWorkers(); + await new Promise((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 | 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((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; } }); }); diff --git a/src/plugin-sdk/sqlite-runtime-testing.ts b/src/plugin-sdk/sqlite-runtime-testing.ts index 499532acbdce..61e7e057616d 100644 --- a/src/plugin-sdk/sqlite-runtime-testing.ts +++ b/src/plugin-sdk/sqlite-runtime-testing.ts @@ -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, diff --git a/test/scripts/vitest-forks-pool.test.ts b/test/scripts/vitest-forks-pool.test.ts new file mode 100644 index 000000000000..20cbbc2b7c8c --- /dev/null +++ b/test/scripts/vitest-forks-pool.test.ts @@ -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]); + }, +); diff --git a/test/vitest-projects-config.test.ts b/test/vitest-projects-config.test.ts index f1cc58698d5b..5707b9db1396 100644 --- a/test/vitest-projects-config.test.ts +++ b/test/vitest-projects-config.test.ts @@ -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}`, diff --git a/test/vitest-scoped-config.test.ts b/test/vitest-scoped-config.test.ts index dccfd57d83c8..85413070f534 100644 --- a/test/vitest-scoped-config.test.ts +++ b/test/vitest-scoped-config.test.ts @@ -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( diff --git a/test/vitest/vitest.forks-pool.ts b/test/vitest/vitest.forks-pool.ts new file mode 100644 index 000000000000..8455b190f124 --- /dev/null +++ b/test/vitest/vitest.forks-pool.ts @@ -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 { + 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; + 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 { + // The public pool hook runs after the transport acknowledges its graceful exit. + this.stopAcknowledged = true; + return super.waitForExit(); + } + + override async stop(): Promise { + 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), +}; diff --git a/test/vitest/vitest.scoped-config.ts b/test/vitest/vitest.scoped-config.ts index 7f8a67129f16..606c446c62db 100644 --- a/test/vitest/vitest.scoped-config.ts +++ b/test/vitest/vitest.scoped-config.ts @@ -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 }),