diff --git a/src/infra/worker-task-pool.mock.test-support.ts b/src/infra/worker-task-pool.mock.test-support.ts new file mode 100644 index 000000000000..6dcd0c35b0e4 --- /dev/null +++ b/src/infra/worker-task-pool.mock.test-support.ts @@ -0,0 +1,32 @@ +import { createRetainedOperation } from "./retained-operation.js"; +import type { createOwnedWorkerTaskPool } from "./worker-task-pool.js"; + +type Pool = ReturnType>; + +function unexpectedTask(): never { + throw new Error("Worker pool fixture requires an explicit task implementation"); +} + +function closedResources() { + const closed = createRetainedOperation(() => {}); + closed.resolve(undefined); + return closed.operation; +} + +/** No native resources are allocated; both resource-retirement paths remain available. */ +export function createOwnedWorkerTaskPoolMock( + overrides: Partial>, +): Pool { + return { + run: unexpectedTask, + runTask: unexpectedTask, + startTask: unexpectedTask, + getSnapshot: unexpectedTask, + rotate: async () => {}, + startRotate: closedResources, + closeResources: async () => {}, + startCloseResources: closedResources, + close: async () => {}, + ...overrides, + }; +} diff --git a/src/state/openclaw-state-db-readonly.error-mapping.test.ts b/src/state/openclaw-state-db-readonly.error-mapping.test.ts index 8ddbb878ad45..20ef06dd4be8 100644 --- a/src/state/openclaw-state-db-readonly.error-mapping.test.ts +++ b/src/state/openclaw-state-db-readonly.error-mapping.test.ts @@ -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["close"]>(), closePool: vi.fn<() => Promise>(), closeResources: vi.fn<(key?: string) => Promise>(), + rotate: vi.fn<() => Promise>(), })); vi.mock("./openclaw-state-worker-context.js", async (importOriginal) => { const actual = await importOriginal(); @@ -45,14 +49,16 @@ vi.mock("./openclaw-state-worker-context.js", async (importOriginal) => { }); vi.mock("../infra/worker-task-pool.js", async (importOriginal) => ({ ...(await importOriginal()), - createOwnedWorkerTaskPool: () => ({ - startTask: (): RetainedWorkerTask => ({ - ...observeAsyncFixture(mock.run), - release: (options) => observeAsyncFixture(() => mock.close(options)), + createOwnedWorkerTaskPool: () => + createOwnedWorkerTaskPoolMock({ + startTask: (): RetainedWorkerTask => ({ + ...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) => { diff --git a/src/state/openclaw-state-db-readonly.retained.test.ts b/src/state/openclaw-state-db-readonly.retained.test.ts index 803d7dc99d7a..6fa8568132c1 100644 --- a/src/state/openclaw-state-db-readonly.retained.test.ts +++ b/src/state/openclaw-state-db-readonly.retained.test.ts @@ -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) { - const name = owner.getStore(); - const inContext = AsyncLocalStorage.snapshot(); - let admitted = false; - let released = false; - let taskReply: OpenClawStateReadReply = reply; - const completion = createRetainedOperation(() => { - 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({ + startTask(input: WorkerTaskInput) { + const name = owner.getStore(); + const inContext = AsyncLocalStorage.snapshot(); + let admitted = false; + let released = false; + let taskReply: OpenClawStateReadReply = reply; + const completion = createRetainedOperation(() => { + 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(() => { - 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(() => { + 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() { diff --git a/src/state/openclaw-state-read-worker.test-harness.ts b/src/state/openclaw-state-read-worker.test-harness.ts index 433e5d45f4a1..ea4371856743 100644 --- a/src/state/openclaw-state-read-worker.test-harness.ts +++ b/src/state/openclaw-state-read-worker.test-harness.ts @@ -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({ + startTask: mock.runTask, + close: mock.closePool, + closeResources: mock.closeResources, + rotate: mock.rotate, + }), + ); }); export function source(name = "source.sqlite") {