diff --git a/docs/plugins/architecture.md b/docs/plugins/architecture.md index 9aedd56ecd95..a5dbf6e5f68e 100644 --- a/docs/plugins/architecture.md +++ b/docs/plugins/architecture.md @@ -210,6 +210,13 @@ changing that digest. Invalid optional package metadata fails only when selected Module acquisition uses the instance's current admission, and disposal closes further capture. +Model-catalog workers keep their captured plugin files in a directory owned by +one worker. The parent removes any remaining captures after that worker exits, +including cancellation and crashes. Files remain available while the worker is +running, and retiring one worker does not remove another generation's captures. +Cancellation releases compute capacity after the worker exits; terminal shutdown +also waits for file cleanup. Failed file removal is reported as a cleanup warning. + Loading metadata alone does not execute every plugin, and registration remains synchronous. Synchronously loaded TypeScript entries and their synchronous TypeScript imports retain Jiti's CommonJS compilation behavior, including `.mts` diff --git a/docs/plugins/sdk-overview/infrastructure.md b/docs/plugins/sdk-overview/infrastructure.md index 47d7f3361fab..582ab7125f31 100644 --- a/docs/plugins/sdk-overview/infrastructure.md +++ b/docs/plugins/sdk-overview/infrastructure.md @@ -100,6 +100,20 @@ For stateless computation, `sharedCompute: true` also shares an aggregate pools in the same isolate. Dedicated ordered pools retain their own execution capacity and still enforce their individual admission limits. +Pass static Node.js Worker settings in `workerOptions`. For per-worker settings, +`prepareWorker()` runs once per Worker creation attempt and returns +`{ options, temporaryDirectory? }`. Its `options` shallowly override +`workerOptions`: properties such as `env`, `workerData`, and `resourceLimits` +replace the whole static property rather than merging nested values. + +A returned `temporaryDirectory` transfers a newly allocated disposable directory +to the pool. Preparation owns cleanup if it fails before returning. The pool +removes the directory only after that Worker exits, including startup failure or +cancellation, and reports deletion failures without replacing the task outcome. +Worker exit releases execution capacity; `close()` also waits for pending file +cleanup. Keep persistent data and files borrowed outside the Worker out of this +directory. + ### SQLite worker stores Use `openSqliteWorkerStore` from diff --git a/src/agents/prepared-model-catalog-worker.integration.test.ts b/src/agents/prepared-model-catalog-worker.integration.test.ts index e9d9c9a9f5a7..c3595ce7bc0a 100644 --- a/src/agents/prepared-model-catalog-worker.integration.test.ts +++ b/src/agents/prepared-model-catalog-worker.integration.test.ts @@ -3,6 +3,7 @@ import fs from "node:fs"; import { createRequire } from "node:module"; import path from "node:path"; import { setImmediate as nextTurn } from "node:timers/promises"; +import { threadId } from "node:worker_threads"; import { beforeEach, describe, expect, it, vi } from "vitest"; import type { OpenClawConfig } from "../config/types.openclaw.js"; import { buildModelsListResult } from "../gateway/server-methods/models-list-result.js"; @@ -53,6 +54,8 @@ import { retainPreparedPluginGeneration } from "./prepared-model-runtime.plugin- import { AuthStorage } from "./sessions/auth-storage.js"; import { markPluginMetadataSnapshotProvided, + loadCompletedFullCatalog, + readCatalogDiscoveryCaptures, usePreparedCatalogWorkerFixtures, } from "./test-helpers/prepared-model-catalog-worker-fixture.js"; @@ -161,16 +164,19 @@ describe("prepared model catalog worker boundary", () => { prepareInboundPluginRegistry: selection.prepareInboundPluginRegistry, }, ); - const catalog = await fixture.snapshot.loadFullModelCatalog!(); - expect.soft(catalog.entries).toContainEqual( - expect.objectContaining({ - provider: PROVIDER_ID, - id: `plugin-generation-${selection.version}`, - }), - ); + await fixture.snapshot.loadFullModelCatalog!(); const auth = await loadPreparedModelRuntimeAuth(fixture.snapshot, { providerIds: [PROVIDER_ID], }); + // The foreground read can return starter rows while cold discovery continues. + await expect + .poll(() => fixture.snapshot.readFullModelCatalog?.()?.entries, { timeout: 30_000 }) + .toContainEqual( + expect.objectContaining({ + provider: PROVIDER_ID, + id: `plugin-generation-${selection.version}`, + }), + ); expect(auth?.authStore.profiles[EXTERNAL_AUTH_PROFILE_ID]).toMatchObject({ access: `${selection.version}:A`, }); @@ -182,13 +188,25 @@ describe("prepared model catalog worker boundary", () => { .split("\n"), ), ).toEqual(new Set([selection.version])); + const captures = readCatalogDiscoveryCaptures(fixture.root); + const workerCaptures = captures.filter((capture) => capture.threadId !== threadId); + const parentCaptures = captures.filter((capture) => capture.threadId === threadId); + expect(workerCaptures.length).toBeGreaterThan(0); + expect(parentCaptures.length).toBeGreaterThan(0); + expect(captures.every((capture) => fs.existsSync(capture.filename))).toBe(true); + fixture.supersede(); + await waitForWorkers(); + await expect + .poll(() => workerCaptures.filter((capture) => fs.existsSync(capture.filename))) + .toEqual([]); + expect(parentCaptures.every((capture) => fs.existsSync(capture.filename))).toBe(true); }); it("keeps explicit read-only full inventories discoverable without a runtime registry", async () => { const fixture = await createStaticSnapshot(0, {}, { readOnly: true }); expect(fixture.snapshot.pluginRegistry).toBeUndefined(); - const catalog = await fixture.snapshot.loadFullModelCatalog!(); + const catalog = await loadCompletedFullCatalog(fixture.snapshot); expect(catalog.entries).toContainEqual( expect.objectContaining({ provider: PROVIDER_ID, id: "plugin-generation-v1" }), ); @@ -242,7 +260,7 @@ describe("prepared model catalog worker boundary", () => { const result = (await build.pending)[0]!; retireAfterTest(retainPreparedPluginGeneration(result.pluginGeneration)); snapshot = result.snapshot; - const modelCatalog = await snapshot.loadFullModelCatalog!(); + const modelCatalog = await loadCompletedFullCatalog(snapshot); expect(modelCatalog.entries).toContainEqual( expect.objectContaining({ provider: PROVIDER_ID, id: "plugin-generation-v1" }), ); @@ -318,7 +336,7 @@ describe("prepared model catalog worker boundary", () => { await refreshPreparedModelRuntimeSnapshots(initialConfig, options); const main = getPreparedModelRuntimeSnapshot(mainInput)!; const sibling = getPreparedModelRuntimeSnapshot(siblingInput)!; - const catalog = await main.loadFullModelCatalog!(); + const catalog = await loadCompletedFullCatalog(main); expect(catalog.entries).toContainEqual( expect.objectContaining({ provider: PROVIDER_ID, id: "plugin-generation-v1" }), ); @@ -326,13 +344,13 @@ describe("prepared model catalog worker boundary", () => { await expect(loadPreparedModelRuntimeAuth(sibling, authScope)).resolves.toMatchObject({ authStore: { profiles: { [EXTERNAL_AUTH_PROFILE_ID]: { access: "v1:A" } } }, }); - await sibling.loadFullModelCatalog!(); + await loadCompletedFullCatalog(sibling); const completed = fs.readFileSync(fixture.marker, "utf8"); const nextInvocation = () => fs.readFileSync(fixture.marker, "utf8").split("start\n").length; const heldInvocation = nextInvocation(); fs.writeFileSync(`${fixture.marker}.hold`, "", "utf8"); - const inFlight = main.loadFullModelCatalog!({ refresh: true }); + const inFlight = loadCompletedFullCatalog(main, { refresh: true }); void inFlight.catch(() => undefined); try { await expect.poll(() => fs.readFileSync(fixture.marker, "utf8")).toBe(`${completed}start\n`); @@ -371,13 +389,13 @@ describe("prepared model catalog worker boundary", () => { }), ); const replaced = getPreparedModelRuntimeSnapshot({ ...siblingInput, config: nextConfig })!; - await replaced.loadFullModelCatalog!(); + await loadCompletedFullCatalog(replaced); fs.writeFileSync(fixture.externalAuthPath, "B", "utf8"); await expect(loadPreparedModelRuntimeAuth(retained, authScope)).resolves.toMatchObject({ authStore: { profiles: { [EXTERNAL_AUTH_PROFILE_ID]: { access: "v1:B" } } }, }); const refreshedInvocation = nextInvocation(); - await expect(retained.loadFullModelCatalog!({ refresh: true })).resolves.toMatchObject({ + await expect(loadCompletedFullCatalog(retained, { refresh: true })).resolves.toMatchObject({ entries: expect.arrayContaining([ expect.objectContaining({ id: `proof-refresh-${refreshedInvocation}-sqlite-true-shared-true-unrelated-true`, @@ -412,7 +430,7 @@ describe("prepared model catalog worker boundary", () => { access: "v1:A", }); } - const catalog = await fixture.snapshot.loadFullModelCatalog?.(); + const catalog = await loadCompletedFullCatalog(fixture.snapshot); expect(catalog?.entries).toContainEqual( expect.objectContaining({ provider: PROVIDER_ID, id: "plugin-generation-v1" }), ); @@ -643,7 +661,7 @@ describe("prepared model catalog worker boundary", () => { modelCatalog: { entries: [route], routeVariants: [route] }, }); const project = async () => { - const fullCatalog = await fixture.snapshot.loadFullModelCatalog?.({ refresh: true }); + const fullCatalog = await loadCompletedFullCatalog(fixture.snapshot, { refresh: true }); const fullAuth = fullCatalog && getPreparedModelFullCatalogAuth(fullCatalog); if (!fullAuth) { throw new Error("full catalog omitted prepared auth"); @@ -973,10 +991,16 @@ describe("prepared model catalog worker boundary", () => { void catalog?.catch(() => {}); try { await waitForMarker(fixture.marker); + const captures = readCatalogDiscoveryCaptures(fixture.root).filter( + (capture) => capture.threadId !== threadId, + ); + expect(captures.length).toBeGreaterThan(0); + expect(captures.every((capture) => fs.existsSync(capture.filename))).toBe(true); fixture.supersede(); await expect(catalog).rejects.toThrow("superseded"); await waitForWorkers(); + expect(captures.filter((capture) => fs.existsSync(capture.filename))).toEqual([]); expect(fs.readFileSync(fixture.marker, "utf8")).toBe("start\n"); } finally { fixture.supersede(); diff --git a/src/agents/prepared-model-catalog-worker.ts b/src/agents/prepared-model-catalog-worker.ts index af6401a10c33..f9547ab8d8b0 100644 --- a/src/agents/prepared-model-catalog-worker.ts +++ b/src/agents/prepared-model-catalog-worker.ts @@ -1,4 +1,7 @@ /** Runs complete model-catalog discovery outside the Gateway event loop. */ +import fs from "node:fs"; +import { tmpdir } from "node:os"; +import path from "node:path"; import { getConfigResolutionFacts, serializeConfigResolutionFacts, @@ -50,6 +53,10 @@ export type PreparedModelCatalogWorkerInput = Readonly<{ pluginMetadataSnapshot: Omit; }>; +export type PreparedModelCatalogWorkerData = PreparedModelCatalogWorkerInput & { + sourceCaptureDirectory: string; +}; + type PreparedModelWorkerCommand = | Readonly<{ kind: "catalog"; providerIds?: readonly string[] }> | Readonly<{ @@ -258,10 +265,19 @@ export function createPreparedModelCatalogWorker( // Only the lifecycle owner may retire it; crashes close the generation permanently. idleTimeoutMs: 0, restartOnError: false, - workerOptions: { - workerData: workerInput, - // Establish state/config environment before worker module initialization reads process.env. - env: workerInput.input.env, + prepareWorker: () => { + const directory = fs.mkdtempSync(path.join(tmpdir(), "openclaw-model-catalog-")); + return { + temporaryDirectory: directory, + options: { + workerData: { + ...workerInput, + sourceCaptureDirectory: directory, + } satisfies PreparedModelCatalogWorkerData, + // Establish state/config environment before module initialization reads process.env. + env: workerInput.input.env, + }, + }; }, validateResult: (message) => { assertCurrent(); diff --git a/src/agents/prepared-model-catalog.worker.ts b/src/agents/prepared-model-catalog.worker.ts index df69728cad2e..fbc1a69c8e5d 100644 --- a/src/agents/prepared-model-catalog.worker.ts +++ b/src/agents/prepared-model-catalog.worker.ts @@ -13,6 +13,7 @@ import { listRuntimePluginIdsFromRegistry } from "../plugins/active-runtime-regi import { normalizePluginsConfig } from "../plugins/config-state.js"; import { isManifestPluginAvailableForControlPlane } from "../plugins/manifest-contract-eligibility.js"; import { restorePluginMetadataSnapshot } from "../plugins/plugin-metadata-snapshot.js"; +import { withPluginSourceCaptureDirectory } from "../plugins/plugin-package-metadata-capture.js"; import { captureProviderCatalogExpiries } from "../plugins/provider-catalog-expiry.js"; import { planRuntimePluginDiscovery } from "../plugins/provider-discovery.js"; import { restorePreparedSyntheticAuthFacts } from "../plugins/provider-synthetic-auth.js"; @@ -37,6 +38,7 @@ import { fingerprintPreparedModelCatalogGeneration, fingerprintPreparedModelWorkerRequest, type PreparedModelCatalogWorkerInput, + type PreparedModelCatalogWorkerData, type PreparedModelWorkerRequest, type PreparedModelWorkerResult, } from "./prepared-model-catalog-worker.js"; @@ -407,16 +409,18 @@ function isWorkerRequest(value: unknown): value is PreparedModelWorkerRequest { } if (parentPort) { - const value = workerData as PreparedModelCatalogWorkerInput; + const value = workerData as PreparedModelCatalogWorkerData; let preparedGeneration: ReturnType | undefined; serveWorkerTasks((request) => { if (!isWorkerRequest(request)) { throw new Error("invalid prepared model catalog worker request"); } - return runPreparedModelCatalogWorkerRequest( - value, - request, - () => (preparedGeneration ??= prepareWorkerGeneration(value)), + return withPluginSourceCaptureDirectory(value.sourceCaptureDirectory, () => + runPreparedModelCatalogWorkerRequest( + value, + request, + () => (preparedGeneration ??= prepareWorkerGeneration(value)), + ), ); }); } diff --git a/src/agents/test-helpers/prepared-model-catalog-worker-fixture.ts b/src/agents/test-helpers/prepared-model-catalog-worker-fixture.ts index ada259d9de38..ab947a08a393 100644 --- a/src/agents/test-helpers/prepared-model-catalog-worker-fixture.ts +++ b/src/agents/test-helpers/prepared-model-catalog-worker-fixture.ts @@ -8,10 +8,14 @@ import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js" import type { PluginMetadataSnapshot } from "../../plugins/plugin-metadata-snapshot.types.js"; import { closeOpenClawAgentDatabasesForTest } from "../../state/openclaw-agent-db.js"; import { clearRuntimeAuthProfileStoreSnapshots } from "../auth-profiles/runtime-snapshots.js"; +import type { ModelCatalogSnapshot } from "../model-catalog.types.js"; +import { isPreparedModelCatalogFull } from "../prepared-model-runtime.full-catalog.js"; import { resetPreparedModelRuntimeSnapshotsForTest } from "../prepared-model-runtime.test-support.js"; +import type { PreparedModelRuntimeSnapshot } from "../prepared-model-runtime.types.js"; + +const waitTimeoutMs = 30_000; export function usePreparedCatalogWorkerFixtures() { - const waitTimeoutMs = 30_000; const retirements = new Set<() => void | Promise>(); const workers = new Set(); // Synchronous capture also covers failures before a worker request is awaited. @@ -78,6 +82,7 @@ export function writeSyntheticAuthDiscoveryFixture(params: { fs.writeFileSync( path.join(params.pluginDir, "provider-discovery.cjs"), `const fs = require("node:fs"); +fs.appendFileSync(${JSON.stringify(path.join(params.root, "discovery-artifact-paths.jsonl"))}, JSON.stringify({ threadId: require("node:worker_threads").threadId, filename: __filename }) + "\\n"); fs.appendFileSync(${JSON.stringify(path.join(params.root, "discovery-artifacts.txt"))}, ${JSON.stringify(params.pluginVersion)} + "\\n"); module.exports = { id: ${JSON.stringify(params.harnessId)}, @@ -125,3 +130,41 @@ export function markPluginMetadataSnapshotProvided( ): PluginMetadataSnapshot { return { ...snapshot, registrySource: "provided", registryDiagnostics: [] }; } + +export function readCatalogDiscoveryCaptures(root: string) { + return fs + .readFileSync(path.join(root, "discovery-artifact-paths.jsonl"), "utf8") + .trim() + .split("\n") + .map((line) => JSON.parse(line) as { threadId: number; filename: string }); +} + +/** Full-result assertions follow publication after the bounded foreground read returns. */ +export async function loadCompletedFullCatalog( + snapshot: PreparedModelRuntimeSnapshot, + options?: { refresh?: boolean }, +): Promise { + const previous = snapshot.readFullModelCatalog!(); + await snapshot.loadFullModelCatalog!(options); + let completed: ModelCatalogSnapshot | undefined; + await expect + .poll( + () => { + const catalog = snapshot.readFullModelCatalog!(); + if ( + catalog && + isPreparedModelCatalogFull(catalog) && + !catalog.pendingProviders?.length && + !catalog.refreshFailed && + (!options?.refresh || catalog !== previous) + ) { + completed = catalog; + return true; + } + return false; + }, + { timeout: waitTimeoutMs }, + ) + .toBe(true); + return completed!; +} diff --git a/src/infra/non-fatal-cleanup.test.ts b/src/infra/non-fatal-cleanup.test.ts index 58a9b783144d..9117445a4a30 100644 --- a/src/infra/non-fatal-cleanup.test.ts +++ b/src/infra/non-fatal-cleanup.test.ts @@ -11,19 +11,26 @@ describe("runBestEffortCleanup", () => { ).resolves.toBe(7); }); - it("swallows cleanup failures and reports them through onError", async () => { - const onError = vi.fn(); - const error = new Error("cleanup failed"); + it.each([false, true])( + "preserves the primary result when cleanup fails (reporter throws: %s)", + async (reporterThrows) => { + const onError = vi.fn(() => { + if (reporterThrows) { + throw new Error("cleanup warning failed"); + } + }); + const error = new Error("cleanup failed"); - await expect( - runBestEffortCleanup({ - cleanup: async () => { - throw error; - }, - onError, - }), - ).resolves.toBeUndefined(); + await expect( + runBestEffortCleanup({ + cleanup: async () => { + throw error; + }, + onError, + }), + ).resolves.toBeUndefined(); - expect(onError).toHaveBeenCalledWith(error); - }); + expect(onError).toHaveBeenCalledWith(error); + }, + ); }); diff --git a/src/infra/non-fatal-cleanup.ts b/src/infra/non-fatal-cleanup.ts index d8cbaf2586be..4aa1cdb29855 100644 --- a/src/infra/non-fatal-cleanup.ts +++ b/src/infra/non-fatal-cleanup.ts @@ -8,7 +8,11 @@ export async function runBestEffortCleanup(params: { try { return await params.cleanup(); } catch (error) { - params.onError?.(error); + try { + params.onError?.(error); + } catch { + // A failed warning sink must not replace the result that cleanup preserves. + } return undefined; } } diff --git a/src/infra/worker-task-pool.cleanup.test.ts b/src/infra/worker-task-pool.cleanup.test.ts new file mode 100644 index 000000000000..4ea547338a81 --- /dev/null +++ b/src/infra/worker-task-pool.cleanup.test.ts @@ -0,0 +1,77 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { describe, expect, it, vi } from "vitest"; +import { createDeferredCore } from "../shared/deferred.js"; +import { WorkerTaskPool } from "./worker-task-pool.js"; +import type { PoolFixtureInput, PoolFixtureResult } from "./worker-task-pool.test-support.js"; + +const cleanup = vi.hoisted(() => vi.fn<() => Promise>()); +vi.mock("./temp-artifact-cleanup.js", () => ({ removeTemporaryArtifacts: cleanup })); + +const workerUrl = new URL("./worker-task-pool.test-support.ts", import.meta.url); + +describe("worker task artifact lifetime", () => { + it("releases stopped execution before disposal while every close joins pending cleanup", async () => { + const directory = fs.mkdtempSync(path.join(os.tmpdir(), "worker-cleanup-owner-")); + const gate = createDeferredCore(); + const entered = createDeferredCore(); + const context = new AsyncLocalStorage(); + let cleanupContext: string | undefined; + cleanup + .mockImplementationOnce(() => { + cleanupContext = context.getStore(); + entered.resolve(); + return gate.promise; + }) + .mockResolvedValue(undefined); + const roots: string[] = []; + const pool = new WorkerTaskPool({ + workerUrl, + maxWorkers: 1, + maxPendingTasks: 1, + prepareWorker: () => { + const owned = fs.mkdtempSync(path.join(directory, "generation-")); + roots.push(owned); + return { options: {}, temporaryDirectory: owned }; + }, + }); + try { + const controller = new AbortController(); + const counters = new SharedArrayBuffer(8); + const active = context.run("request", () => + pool.run({ label: "held", counters, wait: true }, { signal: controller.signal }), + ); + let taskSettled = false; + void active + .finally(() => { + taskSettled = true; + }) + .catch(() => {}); + await expect.poll(() => Atomics.load(new Int32Array(counters), 0)).toBe(1); + context.run("request", () => controller.abort(new Error("canceled owner"))); + await entered.promise; + expect(cleanupContext).toBeUndefined(); + await expect.poll(() => taskSettled).toBe(true); + await expect(active).rejects.toThrow("canceled owner"); + await expect(pool.run({ label: "replacement" }, {})).resolves.toMatchObject({ + label: "replacement", + }); + expect(new Set(roots).size).toBe(2); + let closed = false; + const closing = Promise.all([pool.close(), pool.close()]).then(() => { + closed = true; + }); + await expect.poll(() => cleanup.mock.calls.length).toBe(2); + expect(closed).toBe(false); + gate.resolve(); + await closing; + expect(closed).toBe(true); + } finally { + gate.resolve(); + await pool.close(); + fs.rmSync(directory, { recursive: true, force: true }); + } + }); +}); diff --git a/src/infra/worker-task-pool.test-support.ts b/src/infra/worker-task-pool.test-support.ts index 8b87e7e44c4f..4c3605511dcb 100644 --- a/src/infra/worker-task-pool.test-support.ts +++ b/src/infra/worker-task-pool.test-support.ts @@ -1,10 +1,11 @@ import assert from "node:assert/strict"; -import { threadId } from "node:worker_threads"; +import { threadId, workerData } from "node:worker_threads"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { serveWorkerTasks } from "./worker-task-pool.js"; export type PoolFixtureInput = { label: string; + readStartupOptions?: boolean; exchanges?: number; counters?: SharedArrayBuffer; wait?: boolean; @@ -18,6 +19,7 @@ export type PoolFixtureResult = { buffer?: ArrayBuffer; previousBufferBytes?: number; relayedBufferBytes?: number; + startupOptions?: { data: unknown; argv: string[] }; }; let previousBuffer: ArrayBuffer | undefined; @@ -66,6 +68,9 @@ serveWorkerTasks( buffer: input.buffer, previousBufferBytes, relayedBufferBytes, + ...(input.readStartupOptions + ? { startupOptions: { data: workerData, argv: process.argv.slice(2) } } + : {}), }; }, { transferList: (value) => (value.buffer ? [value.buffer] : []) }, diff --git a/src/infra/worker-task-pool.test.ts b/src/infra/worker-task-pool.test.ts index e5676dba390c..82b078ab438e 100644 --- a/src/infra/worker-task-pool.test.ts +++ b/src/infra/worker-task-pool.test.ts @@ -1,11 +1,15 @@ import assert from "node:assert/strict"; import { AsyncLocalStorage } from "node:async_hooks"; import { execFile } from "node:child_process"; +import { channel } from "node:diagnostics_channel"; +import fs from "node:fs"; import { availableParallelism } from "node:os"; +import path from "node:path"; import { fileURLToPath } from "node:url"; import { promisify } from "node:util"; import type { Worker } from "node:worker_threads"; import { afterEach, describe, expect, it, vi } from "vitest"; +import { createTempDirTracker } from "../../test/helpers/temp-dir.js"; import { createDeferredCore } from "../shared/deferred.js"; import { WorkerTaskPool } from "./worker-task-pool.js"; import type { PoolFixtureInput, PoolFixtureResult } from "./worker-task-pool.test-support.js"; @@ -13,6 +17,7 @@ import type { PoolFixtureInput, PoolFixtureResult } from "./worker-task-pool.tes const workerUrl = new URL("./worker-task-pool.test-support.ts", import.meta.url); const pools: WorkerTaskPool[] = []; const workers = vi.hoisted(() => [] as Worker[]); +const directories = createTempDirTracker(); vi.mock("node:os", async (importOriginal) => ({ ...(await importOriginal()), @@ -51,9 +56,131 @@ afterEach(async () => { for (const worker of workers.splice(0)) { expect(worker.threadId).toBe(-1); } + directories.cleanup(); }); describe("worker task pool", () => { + it.each(["factory", "options", "constructor"] as const)( + "joins cancellation during worker %s preparation before removing scratch", + async (phase) => { + const directory = directories.make("worker-reentrant-preparation-"); + const controller = new AbortController(); + const reason = new Error("canceled during worker preparation"); + const createdBefore = workers.length; + const workerChannel = channel("worker_threads"); + const cancel = () => controller.abort(reason); + if (phase === "constructor") { + workerChannel.subscribe(cancel); + } + const pool = createPool({ + workerUrl, + workerOptions: { + get workerData() { + if (phase === "options") { + cancel(); + } + return { prepared: true }; + }, + }, + prepareWorker: () => { + if (phase === "factory") { + cancel(); + } + return { options: {}, temporaryDirectory: directory }; + }, + }); + try { + await expect(pool.run({ label: "canceled" }, { signal: controller.signal })).rejects.toBe( + reason, + ); + await pool.close(); + const created = workers.slice(createdBefore); + expect(created).toHaveLength(phase === "constructor" ? 1 : 0); + expect(created.map((worker) => worker.threadId)).toEqual( + phase === "constructor" ? [-1] : [], + ); + expect(fs.existsSync(directory)).toBe(false); + } finally { + workerChannel.unsubscribe(cancel); + // A failed regression must still join any Worker created after cancellation. + await Promise.all(workers.slice(createdBefore).map((worker) => worker.terminate())); + await pool.close(); + } + }, + ); + + it.each([false, true])( + "preserves static Worker options with prepared overrides: %s", + async (prepared) => { + const pool = createPool({ + workerUrl, + workerOptions: { + argv: ["shared-argument"], + workerData: { source: "static", retained: true }, + }, + ...(prepared + ? { prepareWorker: () => ({ options: { workerData: { source: "prepared" } } }) } + : {}), + }); + const result = await pool.run({ label: "options", readStartupOptions: true }, {}); + expect(result.startupOptions).toEqual({ + argv: ["shared-argument"], + data: prepared ? { source: "prepared" } : { source: "static", retained: true }, + }); + }, + ); + + it.each(["close", "abort", "exit", "startup-error", "clone-error"] as const)( + "reclaims only its exited worker's scratch after %s", + async (ending) => { + const directory = directories.make("worker-owned-scratch-"); + const unrelated = directories.make("worker-unrelated-scratch-"); + fs.writeFileSync(path.join(directory, "captured-module.js"), "synthetic capture"); + fs.writeFileSync(path.join(unrelated, "retained-module.js"), "unrelated capture"); + const controller = new AbortController(); + const pool = createPool({ + workerUrl: + ending === "startup-error" ? new URL("./missing-worker.mjs", import.meta.url) : workerUrl, + restartOnError: false, + prepareWorker: () => ({ + temporaryDirectory: directory, + options: ending === "clone-error" ? { workerData: () => {} } : {}, + }), + }); + if (ending === "startup-error" || ending === "clone-error") { + await expect(pool.run({ label: ending }, {})).rejects.toMatchObject({ + code: "unavailable", + }); + } else { + await pool.run({ label: "warm" }, {}); + expect(fs.existsSync(directory)).toBe(true); + const worker = workers.at(-1)!; + if (ending === "close") { + await pool.close(); + } else if (ending === "exit") { + await worker.terminate(); + await pool.close(); + } else { + const counters = new SharedArrayBuffer(Int32Array.BYTES_PER_ELEMENT * 2); + const active = pool.run( + { label: "blocked", counters, wait: true }, + { signal: controller.signal }, + ); + void active.catch(() => {}); + await expect.poll(() => Atomics.load(new Int32Array(counters), 0)).toBe(1); + expect(fs.existsSync(directory)).toBe(true); + controller.abort(new Error("scratch canceled")); + await expect(active).rejects.toThrow("scratch canceled"); + } + expect(worker.threadId).toBe(-1); + } + await pool.close(); + expect(fs.existsSync(directory)).toBe(false); + expect(fs.readFileSync(path.join(unrelated, "retained-module.js"), "utf8")).toBe( + "unrelated capture", + ); + }, + ); it("keeps canceled preparation charged until its retained input is released", async () => { const pool = createPool({ workerUrl, maxPendingTasks: 1 }); const gate = createDeferredCore(); diff --git a/src/infra/worker-task-pool.ts b/src/infra/worker-task-pool.ts index d3747e1d0a28..b6af1efdde99 100644 --- a/src/infra/worker-task-pool.ts +++ b/src/infra/worker-task-pool.ts @@ -6,6 +6,7 @@ import { toErrorObject } from "@openclaw/normalization-core/error-coercion"; import { resolveTimerTimeoutMs } from "@openclaw/normalization-core/number-coercion"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { createDeferredCore, type Deferred } from "../shared/deferred.js"; +import { runBestEffortCleanup } from "./non-fatal-cleanup.js"; import { DEFAULT_WORKER_PENDING_BYTES, DEFAULT_WORKER_PENDING_TASKS, @@ -82,6 +83,7 @@ type Task = Deferred & { }; type Slot = { worker?: Worker; + temporaryDirectory?: string; task?: Task; idleTimer?: NodeJS.Timeout; retiring?: Promise; @@ -100,6 +102,7 @@ export class WorkerTaskError extends Error { /** Bounded execution workers; each worker accepts one task at a time. */ export class WorkerTaskPool { private readonly slots = new Set>(); + private readonly artifactCleanups = new Set>(); private readonly queue: Task[] = []; private readonly maxWorkers: number; private readonly maxPendingTasks: number; @@ -121,6 +124,11 @@ export class WorkerTaskPool { private readonly options: { workerUrl: URL; workerOptions?: Omit; + /** Shallow per-Worker overrides; returned scratch stays owned until Worker exit. */ + prepareWorker?: () => { + options: Omit; + temporaryDirectory?: string; + }; maxWorkers?: number; /** Share CPU admission with other stateless compute pools in this isolate. */ sharedCompute?: boolean; @@ -210,7 +218,9 @@ export class WorkerTaskPool { this.finish(slot.task, this.closedError, undefined, true); } } - return Promise.all([...this.slots].map((slot) => this.retire(slot))).then(() => undefined); + return Promise.all([...this.slots].map((slot) => this.retire(slot))) + .then(() => Promise.all(this.artifactCleanups)) + .then(() => undefined); } private dispatch(): void { @@ -267,14 +277,22 @@ export class WorkerTaskPool { // Worker listeners outlive tasks; their creation scope must not retain an async task frame. private createWorker(slot: Slot): Worker { - const worker = runInWorkerPoolContext( - () => - new Worker(this.options.workerUrl, { - // Preserve native require(ESM) and its transitive import-only exports. - execArgv: this.options.workerUrl.pathname.endsWith(".ts") ? ["--import", "tsx/esm"] : [], - ...this.options.workerOptions, - }), - ); + const worker = runInWorkerPoolContext(() => { + const prepared = this.options.prepareWorker?.(); + slot.temporaryDirectory = prepared?.temporaryDirectory; + const workerUrl = this.options.workerUrl; + const workerOptions = { + // Preserve native require(ESM) and its transitive import-only exports. + execArgv: workerUrl.pathname.endsWith(".ts") ? ["--import", "tsx/esm"] : [], + ...this.options.workerOptions, + ...prepared?.options, + }; + // Preparation and option getters can synchronously close the task. + if (slot.retiring) { + throw new WorkerTaskError("worker creation closed during preparation", "unavailable"); + } + return new Worker(workerUrl, workerOptions); + }); slot.worker = worker; worker.on("message", (message: unknown) => { const task = slot.task; @@ -602,11 +620,32 @@ export class WorkerTaskPool { private retire(slot: Slot): Promise { this.clearTimeoutFn(slot.idleTimer); // Retain error listeners until exit: termination can race a worker startup error. - return (slot.retiring ??= (slot.worker?.terminate() ?? Promise.resolve()).then(() => { - slot.worker?.removeAllListeners(); - this.slots.delete(slot); - this.dispatch(); - })); + // Constructor observers can retire this slot before its Worker is assigned. + return (slot.retiring ??= Promise.resolve() + .then(() => slot.worker?.terminate()) + .then(() => { + const directory = slot.temporaryDirectory; + if (directory) { + runInWorkerPoolContext(() => { + const cleanup = runBestEffortCleanup({ + cleanup: async () => { + const { removeTemporaryArtifacts } = await import("./temp-artifact-cleanup.js"); + await removeTemporaryArtifacts(directory, "Worker task"); + }, + onError: (error) => + process.emitWarning( + `Worker task cleanup could not load for ${directory}: ${String(error)}`, + ), + }); + // Release execution capacity at exit; terminal close still joins disposable files. + this.artifactCleanups.add(cleanup); + void cleanup.then(() => this.artifactCleanups.delete(cleanup)); + }); + } + slot.worker?.removeAllListeners(); + this.slots.delete(slot); + this.dispatch(); + })); } } diff --git a/src/plugins/plugin-package-metadata-capture.ts b/src/plugins/plugin-package-metadata-capture.ts index 7c4a80c5ed05..cf1c820521f9 100644 --- a/src/plugins/plugin-package-metadata-capture.ts +++ b/src/plugins/plugin-package-metadata-capture.ts @@ -1,3 +1,4 @@ +import { AsyncLocalStorage } from "node:async_hooks"; import { createHash, randomUUID } from "node:crypto"; import fs from "node:fs"; import { createRequire, isBuiltin } from "node:module"; @@ -627,9 +628,20 @@ export function createPluginPackageMetadataCapture(params: { }; } +const sourceCaptureDirectory = new AsyncLocalStorage(); + +/** A compute worker's parent reclaims this scratch directory after confirmed exit. */ +export function withPluginSourceCaptureDirectory(directory: string, run: () => T): T { + return sourceCaptureDirectory.run(directory, run); +} + /** Admissions and failed-input receipts belong to one source acquisition lifetime. */ export function createPluginSourceCapture(execute?: (run: () => T) => T) { - const directory = fs.realpathSync(fs.mkdtempSync(path.join(tmpdir(), "openclaw-plugin-build-"))); + const directory = fs.realpathSync( + fs.mkdtempSync( + path.join(sourceCaptureDirectory.getStore() ?? tmpdir(), "openclaw-plugin-build-"), + ), + ); fs.chmodSync(directory, 0o700); const inputs = new Map(); const pendingInputs = new Set();