mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
fix(models): reclaim plugin captures after catalog workers stop (#146513)
Give each catalog Worker a parent-owned capture directory. Release execution at confirmed Worker exit and join detached advisory cleanup on pool close. Preserve static Worker options, and fence synchronous cancellation during preparation or construction before releasing ownership. Make full-catalog test assertions follow the existing completed publication after one bounded foreground request, preserving pending and cancellation coverage and all production deadlines.
This commit is contained in:
parent
fadcb209a2
commit
a84d90b9c2
13 changed files with 435 additions and 56 deletions
|
|
@ -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`
|
||||
|
|
|
|||
|
|
@ -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<Operations>` from
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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<PluginMetadataSnapshot, "normalizePluginId">;
|
||||
}>;
|
||||
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<typeof prepareWorkerGeneration> | 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)),
|
||||
),
|
||||
);
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<void>>();
|
||||
const workers = new Set<Worker>();
|
||||
// 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<ModelCatalogSnapshot> {
|
||||
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!;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
},
|
||||
);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -8,7 +8,11 @@ export async function runBestEffortCleanup<T>(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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
77
src/infra/worker-task-pool.cleanup.test.ts
Normal file
77
src/infra/worker-task-pool.cleanup.test.ts
Normal file
|
|
@ -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<void>>());
|
||||
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<string>();
|
||||
let cleanupContext: string | undefined;
|
||||
cleanup
|
||||
.mockImplementationOnce(() => {
|
||||
cleanupContext = context.getStore();
|
||||
entered.resolve();
|
||||
return gate.promise;
|
||||
})
|
||||
.mockResolvedValue(undefined);
|
||||
const roots: string[] = [];
|
||||
const pool = new WorkerTaskPool<PoolFixtureInput, PoolFixtureResult>({
|
||||
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 });
|
||||
}
|
||||
});
|
||||
});
|
||||
|
|
@ -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<PoolFixtureResult>(
|
|||
buffer: input.buffer,
|
||||
previousBufferBytes,
|
||||
relayedBufferBytes,
|
||||
...(input.readStartupOptions
|
||||
? { startupOptions: { data: workerData, argv: process.argv.slice(2) } }
|
||||
: {}),
|
||||
};
|
||||
},
|
||||
{ transferList: (value) => (value.buffer ? [value.buffer] : []) },
|
||||
|
|
|
|||
|
|
@ -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<PoolFixtureInput, PoolFixtureResult>[] = [];
|
||||
const workers = vi.hoisted(() => [] as Worker[]);
|
||||
const directories = createTempDirTracker();
|
||||
|
||||
vi.mock("node:os", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("node:os")>()),
|
||||
|
|
@ -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<PoolFixtureInput>();
|
||||
|
|
|
|||
|
|
@ -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<Input, Output> = Deferred<Output> & {
|
|||
};
|
||||
type Slot<Input, Output> = {
|
||||
worker?: Worker;
|
||||
temporaryDirectory?: string;
|
||||
task?: Task<Input, Output>;
|
||||
idleTimer?: NodeJS.Timeout;
|
||||
retiring?: Promise<void>;
|
||||
|
|
@ -100,6 +102,7 @@ export class WorkerTaskError extends Error {
|
|||
/** Bounded execution workers; each worker accepts one task at a time. */
|
||||
export class WorkerTaskPool<Input, Output> {
|
||||
private readonly slots = new Set<Slot<Input, Output>>();
|
||||
private readonly artifactCleanups = new Set<Promise<void>>();
|
||||
private readonly queue: Task<Input, Output>[] = [];
|
||||
private readonly maxWorkers: number;
|
||||
private readonly maxPendingTasks: number;
|
||||
|
|
@ -121,6 +124,11 @@ export class WorkerTaskPool<Input, Output> {
|
|||
private readonly options: {
|
||||
workerUrl: URL;
|
||||
workerOptions?: Omit<WorkerOptions, "eval">;
|
||||
/** Shallow per-Worker overrides; returned scratch stays owned until Worker exit. */
|
||||
prepareWorker?: () => {
|
||||
options: Omit<WorkerOptions, "eval">;
|
||||
temporaryDirectory?: string;
|
||||
};
|
||||
maxWorkers?: number;
|
||||
/** Share CPU admission with other stateless compute pools in this isolate. */
|
||||
sharedCompute?: boolean;
|
||||
|
|
@ -210,7 +218,9 @@ export class WorkerTaskPool<Input, Output> {
|
|||
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<Input, Output> {
|
|||
|
||||
// Worker listeners outlive tasks; their creation scope must not retain an async task frame.
|
||||
private createWorker(slot: Slot<Input, Output>): 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<Input, Output> {
|
|||
private retire(slot: Slot<Input, Output>): Promise<void> {
|
||||
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();
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<string>();
|
||||
|
||||
/** A compute worker's parent reclaims this scratch directory after confirmed exit. */
|
||||
export function withPluginSourceCaptureDirectory<T>(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?: <T>(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<string, PluginSourceInput>();
|
||||
const pendingInputs = new Set<string>();
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue