test(state): complete reader pool lifecycle mocks (#162711)

State-reader test doubles omitted the pool rotation lifecycle, causing Bun's conservative SQLite cleanup to fail and contaminate later tests.

Share a typed, complete worker-pool fake across the reader fixtures and verify both targeted resource closure and pool rotation. Production LOC is 0.

Proof: 52 files pass on Node 24 and the requested Bun fork (314 tests, one existing skip each); changed-file checks, import-cycle checks and P2 review pass.
This commit is contained in:
Peter Steinberger 2026-10-01 06:19:46 -07:00 • committed by GitHub
parent 1c95e958d9
commit 6b8a374c18
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 132 additions and 65 deletions

View file

@ -0,0 +1,32 @@
import { createRetainedOperation } from "./retained-operation.js";
import type { createOwnedWorkerTaskPool } from "./worker-task-pool.js";
type Pool<Input, Output> = ReturnType<typeof createOwnedWorkerTaskPool<Input, Output>>;
function unexpectedTask(): never {
throw new Error("Worker pool fixture requires an explicit task implementation");
}
function closedResources() {
const closed = createRetainedOperation<void>(() => {});
closed.resolve(undefined);
return closed.operation;
}
/** No native resources are allocated; both resource-retirement paths remain available. */
export function createOwnedWorkerTaskPoolMock<Input, Output>(
overrides: Partial<Pool<Input, Output>>,
): Pool<Input, Output> {
return {
run: unexpectedTask,
runTask: unexpectedTask,
startTask: unexpectedTask,
getSnapshot: unexpectedTask,
rotate: async () => {},
startRotate: closedResources,
closeResources: async () => {},
startCloseResources: closedResources,
close: async () => {},
...overrides,
};
}

View file

@ -2,7 +2,9 @@ import fs from "node:fs";
import path from "node:path";
import { afterEach, beforeEach, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import * as sqliteRuntime from "../infra/bun-sqlite-library.js";
import { createRetainedOperation, type RetainedOperation } from "../infra/retained-operation.js";
import { createOwnedWorkerTaskPoolMock } from "../infra/worker-task-pool.mock.test-support.js";
import type { OwnedWorkerTask, RetainedWorkerTask } from "../infra/worker-task-pool.types.js";
import { PluginBlobStoreError } from "../plugin-state/plugin-blob-store.types.js";
import { createDeferredCore } from "../shared/deferred.js";
@ -15,6 +17,7 @@ import { withExistingOpenClawStateSchema } from "./openclaw-state-db-schema-poli
import type {
OpenClawStateReadPhase,
OpenClawStateReadReply,
OpenClawStateReadRequest,
} from "./openclaw-state-read.types.js";
import { captureOpenClawStateReadWorkerContext } from "./openclaw-state-worker-context.js";
import { encodeOpenClawStateWorkerError } from "./openclaw-state-worker-error.js";
@ -35,6 +38,7 @@ const mock = vi.hoisted(() => ({
close: vi.fn<OwnedWorkerTask<OpenClawStateReadReply>["close"]>(),
closePool: vi.fn<() => Promise<void>>(),
closeResources: vi.fn<(key?: string) => Promise<void>>(),
rotate: vi.fn<() => Promise<void>>(),
}));
vi.mock("./openclaw-state-worker-context.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("./openclaw-state-worker-context.js")>();
@ -45,14 +49,16 @@ vi.mock("./openclaw-state-worker-context.js", async (importOriginal) => {
});
vi.mock("../infra/worker-task-pool.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../infra/worker-task-pool.js")>()),
createOwnedWorkerTaskPool: () => ({
startTask: (): RetainedWorkerTask<OpenClawStateReadReply> => ({
...observeAsyncFixture(mock.run),
release: (options) => observeAsyncFixture(() => mock.close(options)),
createOwnedWorkerTaskPool: () =>
createOwnedWorkerTaskPoolMock<OpenClawStateReadRequest, OpenClawStateReadReply>({
startTask: (): RetainedWorkerTask<OpenClawStateReadReply> => ({
...observeAsyncFixture(mock.run),
release: (options) => observeAsyncFixture(() => mock.close(options)),
}),
close: mock.closePool,
closeResources: mock.closeResources,
rotate: mock.rotate,
}),
close: mock.closePool,
closeResources: mock.closeResources,
}),
}));
const tempDirs = useAutoCleanupTempDirTracker((cleanup) =>
@ -60,6 +66,7 @@ const tempDirs = useAutoCleanupTempDirTracker((cleanup) =>
mock.close.mockResolvedValue();
mock.closePool.mockResolvedValue();
mock.closeResources.mockResolvedValue();
mock.rotate.mockResolvedValue();
await closeOpenClawStateDatabaseAsync();
cleanup();
}),
@ -75,6 +82,7 @@ beforeEach(() => {
mock.close.mockReset().mockResolvedValue();
mock.closePool.mockReset().mockResolvedValue();
mock.closeResources.mockReset().mockResolvedValue();
mock.rotate.mockReset().mockResolvedValue();
});
function source() {
const root = tempDirs.make("state-read-error-phase-");
@ -88,6 +96,29 @@ function mapper() {
return { mapped, mapError: vi.fn((_error: unknown, _phase: OpenClawStateReadPhase) => mapped) };
}
it.each([false, true])(
"closes a mapped reader with explicit SQLite close capability %s",
async (explicitClose) => {
const capabilities = vi.spyOn(sqliteRuntime, "getSqliteRuntimeCapabilities").mockReturnValue({
explicitSqliteCloseReleasesNativeResources: explicitClose,
decided: true,
reason: "test policy",
});
try {
const options = source();
await executeExistingOpenClawStateRead(options, { type: "fleet.list" });
await closeOpenClawStateDatabaseByPathAsync(options.path);
expect(explicitClose ? mock.closeResources : mock.rotate).toHaveBeenCalledOnce();
expect(explicitClose ? mock.rotate : mock.closeResources).not.toHaveBeenCalled();
expect(mock.closePool).not.toHaveBeenCalled();
await closeOpenClawStateDatabaseAsync();
expect(mock.closePool).toHaveBeenCalledOnce();
} finally {
capabilities.mockRestore();
}
},
);
it.each(["retired", "different-source"] as const)(
"maps %s captured authority before dispatching a read",
async (kind) => {

View file

@ -5,6 +5,7 @@ import { afterEach, beforeEach, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import { createRetainedOperation } from "../infra/retained-operation.js";
import type { RetainedPreparedSqliteReadOnlyLocation } from "../infra/sqlite-readonly-location.types.js";
import { createOwnedWorkerTaskPoolMock } from "../infra/worker-task-pool.mock.test-support.js";
import type { RetainedWorkerTask, WorkerTaskInput } from "../infra/worker-task-pool.types.js";
import { createDeferredCore } from "../shared/deferred.js";
import type {
@ -87,59 +88,59 @@ beforeEach(() => {
mock.prepareNative.mockReset();
mock.prepareFresh.mockReset();
mock.prepareInherited.mockReset();
mock.pool.mockReset().mockImplementation(() => ({
startTask(input: WorkerTaskInput<OpenClawStateReadRequest>) {
const name = owner.getStore();
const inContext = AsyncLocalStorage.snapshot();
let admitted = false;
let released = false;
let taskReply: OpenClawStateReadReply = reply;
const completion = createRetainedOperation<OpenClawStateReadReply>(() => {
if (completion.operation.read().status !== "pending") {
return;
}
if (!admitted && occupied < 2) {
const request = typeof input === "function" ? input() : input;
if (request instanceof Promise) {
throw new Error("This worker fixture requires synchronous admission");
mock.pool.mockReset().mockImplementation(() =>
createOwnedWorkerTaskPoolMock<OpenClawStateReadRequest, OpenClawStateReadReply>({
startTask(input: WorkerTaskInput<OpenClawStateReadRequest>) {
const name = owner.getStore();
const inContext = AsyncLocalStorage.snapshot();
let admitted = false;
let released = false;
let taskReply: OpenClawStateReadReply = reply;
const completion = createRetainedOperation<OpenClawStateReadReply>(() => {
if (completion.operation.read().status !== "pending") {
return;
}
expect(owner.getStore()).toBe(name);
taskReply = request.command.type === "admit" ? { ok: true, type: "admit" } : reply;
occupied++;
admitted = true;
events.push(`admit ${name}`);
}
if (admitted && ready) {
completion.resolve(taskReply);
}
});
const cleanup = createRetainedOperation<void>(() => {
completion.operation.service();
if (completion.operation.read().status === "pending") {
return;
}
if (!released) {
released = true;
occupied--;
expect(owner.getStore()).toBe(name);
events.push(`release ${name}`);
}
cleanup.resolve(undefined);
});
const task = { ...completion.operation, release: () => cleanup.operation };
tasks.push(task);
releaseFixtures.push(() =>
inContext(() => {
task.service();
task.release().service();
}),
);
submitted.resolve();
return task;
},
closeResources: async () => {},
close: async () => {},
}));
if (!admitted && occupied < 2) {
const request = typeof input === "function" ? input() : input;
if (request instanceof Promise) {
throw new Error("This worker fixture requires synchronous admission");
}
expect(owner.getStore()).toBe(name);
taskReply = request.command.type === "admit" ? { ok: true, type: "admit" } : reply;
occupied++;
admitted = true;
events.push(`admit ${name}`);
}
if (admitted && ready) {
completion.resolve(taskReply);
}
});
const cleanup = createRetainedOperation<void>(() => {
completion.operation.service();
if (completion.operation.read().status === "pending") {
return;
}
if (!released) {
released = true;
occupied--;
expect(owner.getStore()).toBe(name);
events.push(`release ${name}`);
}
cleanup.resolve(undefined);
});
const task = { ...completion.operation, release: () => cleanup.operation };
tasks.push(task);
releaseFixtures.push(() =>
inContext(() => {
task.service();
task.release().service();
}),
);
submitted.resolve();
return task;
},
}),
);
});
function source() {

View file

@ -3,6 +3,7 @@ import path from "node:path";
import { afterEach, beforeEach, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import { createRetainedOperation, type RetainedOperation } from "../infra/retained-operation.js";
import { createOwnedWorkerTaskPoolMock } from "../infra/worker-task-pool.mock.test-support.js";
import type {
OwnedWorkerTask,
WorkerTaskInput,
@ -73,12 +74,14 @@ beforeEach(() => {
mock.closePool.mockReset().mockResolvedValue();
mock.closeResources.mockReset().mockResolvedValue();
mock.rotate.mockReset().mockResolvedValue();
mock.create.mockReset().mockImplementation(() => ({
startTask: mock.runTask,
close: mock.closePool,
closeResources: mock.closeResources,
rotate: mock.rotate,
}));
mock.create.mockReset().mockImplementation(() =>
createOwnedWorkerTaskPoolMock<OpenClawStateReadRequest, OpenClawStateReadReply>({
startTask: mock.runTask,
close: mock.closePool,
closeResources: mock.closeResources,
rotate: mock.rotate,
}),
);
});
export function source(name = "source.sqlite") {