diff --git a/docs/reference/database-schemas/storage-changes.md b/docs/reference/database-schemas/storage-changes.md index 042177aa9641..4e06f27f9e71 100644 --- a/docs/reference/database-schemas/storage-changes.md +++ b/docs/reference/database-schemas/storage-changes.md @@ -24,19 +24,38 @@ and publishes the result. Avoid exposing a generic SQL callback to application code or adding an asynchronous wrapper around an existing asynchronous facade. The plugin KV API already has asynchronous methods over its SQLite owner. -Shared-state operations that request host transaction or commit admission retain -lifecycle coordinator custody before worker dispatch. Native host writers can -then borrow that same owner while servicing the worker's grants, avoiding a -coordinator wait that blocks the grant handler. A foreign coordinator owner is -waited out asynchronously before dispatch, within the existing SQLite lock budget. -Grant services reuse the physical worker owner's admitted path aliases, so native -handles opened through symlinked directories reach the same service without new -filesystem lookups. -The waiting job retains its FIFO position and capacity reservation; cancellation -or worker exit wakes the wait without replaying a dispatched write. Source -authority and persisted transaction checks remain unchanged. Coordinator acquisition -and release still perform control SQL on the host; this does not complete the migration of native -state writers. Database schemas, retention, and update behavior are unchanged. +Shared-state operations that request host transaction or commit admission acquire +fresh lifecycle coordinator custody on their executing SQLite worker. A live +parent-owned maintenance or native lease still delegates its existing custody. +The waiting job keeps its FIFO position and capacity reservation. Native attempts +use zero busy timeout and asynchronous backoff within the original captured lock +budget; only acquisition retries. The host rechecks current authority during +preparation, before native execution, and at the existing transaction and commit +grants. Cancellation before native execution joins coordinator cleanup without +replaying the command. + +Legacy native host writers service the same job's preparation and authority ports +between short coordinator-lock attempts, including path aliases. This lets the +worker finish while the host is inside a synchronous native caller. Successful +worker execution releases its physical coordinator after native settlement and +before result framing; the broker retains operation admission and transport +credits through the complete result. Unsettled native work retains custody until +worker exit. If native coordinator cleanup fails after the command settles, the +same per-job port services bounded result frames while a native host writer waits. +The broker receives the complete outcome, joins worker exit, and reports cleanup +separately without discarding that outcome or replaying the write. Incomplete +result delivery retains its existing unknown-outcome handling. Final publication +and follower dispatch run after synchronous native wait servicing returns. +Preparation refusals retain that port through terminal cleanup replies too. +The shared-state owner retires the exact unavailable actor after its accepted +callbacks finish, so the next call opens a usable actor without retrying the prior write. +Nested callbacks return their completed outcomes while further commands on the +failed actor refuse without waiting for the enclosing callback to close itself. +Retirement cleanup failures retain canonical retry custody and report separately +from the completed outcome. +Native host writers, Gateway lifecycle ownership, and source-handle +preparation retain their existing owners. Schemas, retention, and update behavior +are unchanged. Managed outgoing image metadata lookups and cleanup inventories read through the shared-state worker, retaining their writable, creating database-open behavior. @@ -285,8 +304,8 @@ comparison fences, and transactions. Read-only clients retain artifact-preservin reads and never create missing state. Reconnect waits for accepted persistence, and client shutdown drains it before returning; supplied cancellation and owner guards are checked again at worker admission. Device identity creation and the -compound pairing recovery transaction retain their existing owners. Host admission -still uses the synchronous lifecycle coordinator; token-data SQL runs in the worker. +compound pairing recovery transaction retain their existing owners. Fresh token +mutations acquire and release lifecycle custody on the same worker as token-data SQL. One-shot calls initialize that actor during request preparation, before starting the RPC timeout, without reading or caching token facts. diff --git a/scripts/pr-lib/wrapper-components.txt b/scripts/pr-lib/wrapper-components.txt index f3ea839eadf5..f5af39f9de1e 100644 --- a/scripts/pr-lib/wrapper-components.txt +++ b/scripts/pr-lib/wrapper-components.txt @@ -695,6 +695,7 @@ src/infra/sqlite-worker-broker.ts src/infra/sqlite-worker-client.ts src/infra/sqlite-worker-contract.ts src/infra/sqlite-worker-identity.ts +src/infra/sqlite-worker-lifecycle-preparation.ts src/infra/sqlite-worker-operation-admission.ts src/infra/sqlite-worker-store.ts src/infra/sqlite-worker-transfer.ts diff --git a/src/agents/harness/native-hook-relay-store.test.ts b/src/agents/harness/native-hook-relay-store.test.ts index 2ccc5681d943..be46c95472cc 100644 --- a/src/agents/harness/native-hook-relay-store.test.ts +++ b/src/agents/harness/native-hook-relay-store.test.ts @@ -1,13 +1,13 @@ import fs from "node:fs"; import path from "node:path"; -import { DatabaseSync, StatementSync } from "node:sqlite"; +import { DatabaseSync } from "node:sqlite"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; import { - resolveStateDatabaseCoordinatorPath, - resolveStateLifecycleRuntimeDirectory, -} from "../../infra/state-database-coordinator.js"; -import { closeOpenClawStateDatabaseAsync } from "../../state/openclaw-state-db-cache.js"; + closeOpenClawStateDatabaseAsync, + closeOpenClawStateDatabaseByPathAsync, +} from "../../state/openclaw-state-db-cache.js"; +import { observeMainThreadSql } from "../../test-utils/main-thread-sql-spies.js"; import { deleteNativeHookRelayBridgeRecordIfOwned, pruneNativeHookRelayBridgeRecords, @@ -50,30 +50,9 @@ function bridgeRecord( } describe("native hook relay store", () => { - it("persists bridge records off thread while retaining host coordinator SQL", async () => { - type ObservedSql = { databasePath: string | null; sql: string }; - const observed: ObservedSql[] = []; - const prepared = new WeakMap(); - // oxlint-disable-next-line typescript/unbound-method -- The observation wrapper forwards its database receiver with .call. - const prepare = DatabaseSync.prototype.prepare; - // oxlint-disable-next-line typescript/unbound-method -- The observation wrapper forwards its database receiver with .call. - const exec = DatabaseSync.prototype.exec; - vi.spyOn(DatabaseSync.prototype, "prepare").mockImplementation( - function (this: DatabaseSync, sql) { - const observation = { databasePath: this.location(), sql }; - observed.push(observation); - const statement = prepare.call(this, sql); - prepared.set(statement, observation); - return statement; - }, - ); - vi.spyOn(DatabaseSync.prototype, "exec").mockImplementation(function (this: DatabaseSync, sql) { - observed.push({ databasePath: this.location(), sql }); - return exec.call(this, sql); - }); - const statements = (["get", "all", "run", "iterate"] as const).map((method) => - vi.spyOn(StatementSync.prototype, method), - ); + it("persists bridge records and closes their worker without caller-thread SQLite", async () => { + const sql = observeMainThreadSql(); + const close = vi.spyOn(DatabaseSync.prototype, "close"); expect( await readNativeHookRelayBridgeRecord({ relayId: "absent", stateDbPath: primaryStateDbPath }), ).toBeUndefined(); @@ -144,39 +123,9 @@ describe("native hook relay store", () => { stateDbPath: primaryStateDbPath, }), ).toBeUndefined(); - for (const counter of statements) { - for (const receiver of counter.mock.contexts) { - const observation = receiver instanceof StatementSync ? prepared.get(receiver) : undefined; - if (!observation) { - throw new Error("Unattributed caller-thread SQLite statement"); - } - observed.push(observation); - } - } - const coordinatorPath = fs.realpathSync( - resolveStateDatabaseCoordinatorPath({ - databasePath: primaryStateDbPath, - runtimeDirectory: resolveStateLifecycleRuntimeDirectory(), - uid: process.getuid?.(), - }), - ); - const probes = new Set([ - "SELECT sqlite_version() AS version", - "SELECT sqlite_compileoption_used('OMIT_LOAD_EXTENSION') AS omitted", - ]); - const control = new Set([ - "PRAGMA busy_timeout = 0; PRAGMA journal_mode = MEMORY; BEGIN EXCLUSIVE;", - "ROLLBACK", - ]); - expect(observed.some(({ databasePath }) => databasePath !== null)).toBe(true); - for (const { databasePath, sql } of observed) { - if (databasePath === null) { - expect(probes.has(sql), sql).toBe(true); - } else { - expect(fs.realpathSync(databasePath)).toBe(coordinatorPath); - expect(control.has(sql), sql).toBe(true); - } - } + await closeOpenClawStateDatabaseByPathAsync(primaryStateDbPath); + sql.expectIdle(); + expect(close).not.toHaveBeenCalled(); }); it("requires matching token and pid to renew or delete a bridge", async () => { diff --git a/src/gateway/agent-turn/agent-run-dispatch.sqlite.test.ts b/src/gateway/agent-turn/agent-run-dispatch.sqlite.test.ts index 2ed865348e1f..16b064da846f 100644 --- a/src/gateway/agent-turn/agent-run-dispatch.sqlite.test.ts +++ b/src/gateway/agent-turn/agent-run-dispatch.sqlite.test.ts @@ -5,6 +5,7 @@ import { expect, it, vi } from "vitest"; import { ErrorCodes } from "../../../packages/gateway-protocol/src/index.js"; import { createDeferred } from "../../../test/helpers/promise.js"; import { trackSqliteStatementExecutions } from "../../../test/helpers/sqlite-statement-execution-counter.js"; +import { observeDeviceAuthHostSql } from "../../infra/device-auth-store.sql.test-support.js"; import { createEmptyPluginRegistry } from "../../plugins/registry-empty.js"; import { markPluginRegistryActive, @@ -76,6 +77,8 @@ it("creates a durable Gateway task before provider entry and settles its exact u }, ); const workerMessages = vi.spyOn(Worker.prototype, "postMessage"); + const hostSql = observeDeviceAuthHostSql(state.statePath("state", "openclaw.sqlite")); + let creationSql: ReturnType | undefined; const noWrites: HostWrites = { task: 0, delivery: 0, flow: 0 }; let creation: { taskId: string; writes: HostWrites } | undefined; configureTaskRegistryRuntime({ @@ -83,6 +86,7 @@ it("creates a durable Gateway task before provider entry and settles its exact u onEvent(event) { if (!creation && event.kind === "upserted" && event.task.runId === runId) { creation = { taskId: event.task.taskId, writes: { ...tracker.counts } }; + creationSql = hostSql.counts(); } }, }, @@ -154,6 +158,9 @@ it("creates a durable Gateway task before provider entry and settles its exact u notifyPolicy: "silent", }); expect(creation).toEqual({ taskId: running.taskId, writes: noWrites }); + expect(Object.values(creationSql ?? {}).flatMap((counts) => Object.values(counts))).toEqual( + Array(28).fill(0), + ); expect(running.parentFlowId).toBeUndefined(); expect(observed.flowCount).toBe(0); // Live run-owner binding is still a separate native writer in this initial-creation slice. @@ -164,8 +171,15 @@ it("creates a durable Gateway task before provider entry and settles its exact u expect(getTaskRunOwner(running)).toBeDefined(); expect(emitFinal).not.toHaveBeenCalled(); const beforeSettlement = { ...tracker.counts }; + const beforeSettlementSql = hostSql.counts(); releaseProvider.resolve(); await execution; + expect(hostSql.counts()).toEqual(beforeSettlementSql); + console.info("Gateway task host SQL", { + creation: creationSql, + beforeSettlement: beforeSettlementSql, + afterSettlement: hostSql.counts(), + }); const completed = loadTaskRegistryStateFromSqlite(); expect([...completed.tasks.keys()]).toEqual([running.taskId]); @@ -199,6 +213,7 @@ it("creates a durable Gateway task before provider entry and settles its exact u } finally { releaseProvider.resolve(); await execution; + hostSql.restore(); tracker.restore(); workerMessages.mockRestore(); provider.execute.mockReset(); diff --git a/src/infra/acquire-with-wait.ts b/src/infra/acquire-with-wait.ts index 9631407840a7..1431546e2d23 100644 --- a/src/infra/acquire-with-wait.ts +++ b/src/infra/acquire-with-wait.ts @@ -2,7 +2,7 @@ import { setTimeout as sleep } from "node:timers/promises"; /** Retry only acquisition, keeping the deadline independent of wall-clock changes. */ export async function acquireWithWait(params: { - acquire: () => T; + acquire: () => T | Promise; shouldRetry: (error: unknown) => boolean; deadlineMs: number; pollIntervalMs: number; @@ -14,7 +14,7 @@ export async function acquireWithWait(params: { let delayMs = params.pollIntervalMs; for (;;) { try { - return params.acquire(); + return await params.acquire(); } catch (error) { if (!params.shouldRetry(error)) { throw error; diff --git a/src/infra/device-auth-store.worker.test.ts b/src/infra/device-auth-store.worker.test.ts index a7f9292cd278..33f8d014d148 100644 --- a/src/infra/device-auth-store.worker.test.ts +++ b/src/infra/device-auth-store.worker.test.ts @@ -69,6 +69,8 @@ it("keeps cold, warm, read-only, ordered token-data operations and cleanup off t expect(await tokens.clearOriginDeviceToken(origin)).toBe(true); await closeOpenClawStateDatabaseAsync(); expect(Object.values(sql.counts().data)).toEqual(Array(7).fill(0)); + expect(Object.values(sql.counts().coordinator)).toEqual(Array(7).fill(0)); + expect(Object.values(sql.counts().runtimeInitialization)).toEqual(Array(7).fill(0)); expect(Object.values(sql.counts().unknown)).toEqual(Array(7).fill(0)); await expect(fs.stat(state.path("changed-state"))).rejects.toMatchObject({ code: "ENOENT" }); } finally { diff --git a/src/infra/sqlite-coordinator.ts b/src/infra/sqlite-coordinator.ts index 54f1d909638c..c7096369b037 100644 --- a/src/infra/sqlite-coordinator.ts +++ b/src/infra/sqlite-coordinator.ts @@ -7,6 +7,7 @@ import { sameFileIdentity } from "./fs-safe-advanced.js"; import { openNodeSqliteDatabase } from "./node-sqlite.js"; import { applyPrivateModeSync } from "./private-mode.js"; import { isSqliteLockError } from "./sqlite-error-diagnostics.js"; +import { sqliteWriteAdmissionServicesForLocation } from "./sqlite-transaction.js"; export const SqliteCoordinatorError = resolveGlobalSingleton( Symbol.for("openclaw.sqliteCoordinatorError"), @@ -278,13 +279,31 @@ function tryAcquireSqliteCoordinator( // Kysely transaction callbacks cannot own a lock beyond their synchronous commit section. // This handle never writes or commits data. Keep the empty database's initial // journal in memory so acquiring a lock does not create filesystem artifacts. - database.exec( - `PRAGMA busy_timeout = ${busyTimeoutMs}; PRAGMA journal_mode = MEMORY; ${ - mode === "exclusive" - ? "BEGIN EXCLUSIVE;" - : "BEGIN; SELECT rootpage FROM sqlite_schema LIMIT 1;" - }`, - ); + const services = + mode === "exclusive" ? sqliteWriteAdmissionServicesForLocation(location) : undefined; + const deadline = performance.now() + busyTimeoutMs; + for (;;) { + const attemptTimeout = services + ? Math.min(25, Math.max(0, Math.ceil(deadline - performance.now()))) + : busyTimeoutMs; + try { + database.exec( + `PRAGMA busy_timeout = ${attemptTimeout}; PRAGMA journal_mode = MEMORY; ${ + mode === "exclusive" + ? "BEGIN EXCLUSIVE;" + : "BEGIN; SELECT rootpage FROM sqlite_schema LIMIT 1;" + }`, + ); + break; + } catch (error) { + if (!services || !isSqliteLockError(error) || performance.now() >= deadline) { + throw error; + } + for (const service of services) { + service(); + } + } + } if (poolLocation && before) { const current = readCoordinatorIdentity(poolLocation); if (matchesCoordinatorIdentity(before, current)) { diff --git a/src/infra/sqlite-store.worker.ts b/src/infra/sqlite-store.worker.ts index c823883a65ac..74ac98f2b12d 100644 --- a/src/infra/sqlite-store.worker.ts +++ b/src/infra/sqlite-store.worker.ts @@ -4,7 +4,10 @@ import { parentPort, type MessagePort } from "node:worker_threads"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { routeLogsToStderr } from "../logging/console.js"; import { drainProcessOutput } from "../process/output-drain.js"; -import { encodeOpenClawStateWorkerError } from "../state/openclaw-state-worker-error.js"; +import { + encodeOpenClawStateWorkerError, + type OpenClawStateWorkerErrorPayload, +} from "../state/openclaw-state-worker-error.js"; import { withSqliteReaderOwner } from "./sqlite-reader-lifecycle.js"; import { SQLITE_WORKER_MAX_RESULT_BYTES, @@ -15,6 +18,7 @@ import { type SqliteWorkerRequest, } from "./sqlite-worker-contract.js"; import { assertExistingDatabaseIdentity } from "./sqlite-worker-identity.js"; +import { acquireSqliteWorkerLifecycle } from "./sqlite-worker-lifecycle-preparation.js"; import { withSqliteWorkerOperationAdmission, requestSqliteWorkerOperationAdmission, @@ -33,6 +37,7 @@ import { attachStateLifecycleDelegate, withStateDatabaseCoordinatorRuntimeDirectory, } from "./state-database-coordinator.js"; +import type { acquireStateDatabaseCoordinator } from "./state-database-coordinator.js"; import { ownedWorkerBytes } from "./worker-transfer-bytes.js"; const port = parentPort; @@ -58,6 +63,10 @@ const gatewayFences = new Map< Awaited> >(); let sourceLoaderRegistered = false; +let preparedGatewayActor: number | undefined; +let lifecycleReply: { actor: number; port: MessagePort } | undefined; +let nativeCleanupFailure: OpenClawStateWorkerErrorPayload | undefined; +let lifecyclePreparation: { actor: number; port: MessagePort; deadlineNs: bigint } | undefined; let operationAdmission: { actor: number; port: MessagePort } | undefined; // Input and result continuations retain the original job's delegation. let lifecycle: @@ -102,6 +111,16 @@ async function receive(request: SqliteWorkerRequest): Promise { try { let value: unknown; if (request.type !== "result-next" && request.type !== "execute-frame") { + if (request.lifecyclePreparation) { + if (lifecyclePreparation || !request.workerStateLifecycle) { + throw new Error("SQLite lifecycle preparation differs from its job"); + } + lifecyclePreparation = { + actor: request.actor, + port: request.lifecyclePreparation, + deadlineNs: request.workerStateLifecycle.deadlineNs, + }; + } if (request.operationAdmission) { if (operationAdmission) { throw new Error("SQLite operation admission still belongs to the preceding operation"); @@ -147,6 +166,9 @@ async function receive(request: SqliteWorkerRequest): Promise { actorId: String(request.actor), }), ); + if (request.workerStateLifecycle) { + preparedGatewayActor = request.actor; + } retire = false; } if (request.maintenanceSchemaFence) { @@ -168,11 +190,36 @@ async function receive(request: SqliteWorkerRequest): Promise { retire = false; } } - const executeCommand = (command: unknown) => { + const executeCommand = async (command: unknown) => { + let coordinator: ReturnType | undefined; + if (lifecyclePreparation) { + const context = stateContexts.get(request.actor); + const databasePath = actorPaths.get(request.actor); + if (lifecyclePreparation.actor !== request.actor || !context || !databasePath) { + throw new Error("SQLite lifecycle preparation lost its captured actor"); + } + const preparation = lifecyclePreparation; + lifecyclePreparation = undefined; + lifecycleReply = { actor: request.actor, port: preparation.port }; + const prepared = await acquireSqliteWorkerLifecycle({ + port: preparation.port, + databasePath, + deadlineNs: preparation.deadlineNs, + runtime: context.coordinatorRuntime, + onUnsettled: () => { + retire = true; + }, + }); + coordinator = prepared.coordinator; + if (prepared.admission) { + operationAdmission = { actor: request.actor, port: prepared.admission }; + } + } const backend = actors.get(request.actor); if (!backend) { throw new Error("SQLite worker actor is closed"); } + preparedGatewayActor = undefined; const assertSettled = (failure?: { error: unknown }) => { try { const settlement: unknown = runInActorContext(request.actor, () => @@ -201,36 +248,53 @@ async function receive(request: SqliteWorkerRequest): Promise { } }; try { - // SAFETY: The broker serialized a command from this actor's typed store contract. - const typedCommand = command as SqliteWorkerCommand; - value = runInActorContext(request.actor, () => - withSqliteReaderOwner( - { - operation: typedCommand.type, - ownerKind: "worker", - actorId: request.actor, - }, - () => ({ - // SAFETY: The typed host command is serialized once; framing validates complete reconstruction. - result: backend.execute(typedCommand), - }), - ), - ).result; - } catch (error) { - assertSettled({ error }); - throw error; - } - executed = true; - completeResult = true; - if (isPromise(value) || (isRecord(value) && typeof value.then === "function")) { - retire = true; - if (isPromise(value)) { - // Retirement owns the failure; consume rejection while native exit is joined. - void value.catch(() => {}); + try { + // SAFETY: The broker serialized a command from this actor's typed store contract. + const typedCommand = command as SqliteWorkerCommand; + value = runInActorContext(request.actor, () => + withSqliteReaderOwner( + { + operation: typedCommand.type, + ownerKind: "worker", + actorId: request.actor, + }, + () => ({ + // SAFETY: The typed host command is serialized once; framing validates complete reconstruction. + result: backend.execute(typedCommand), + }), + ), + ).result; + } catch (error) { + assertSettled({ error }); + throw error; + } + executed = true; + completeResult = true; + if (isPromise(value) || (isRecord(value) && typeof value.then === "function")) { + retire = true; + if (isPromise(value)) { + // Retirement owns the failure; consume rejection while native exit is joined. + void value.catch(() => {}); + } + throw new Error("SQLite worker operations must remain synchronous"); + } + assertSettled(); + } finally { + // No write-capable continuation may outlive this lease. A failed + // settlement retains the native owner until the broker joins worker exit. + if (coordinator && !retire) { + try { + coordinator.release(); + } catch (error) { + // The operation already settled. Retain its result and this native + // custody until the host receives the complete reply and joins exit. + const failure = error instanceof Error ? error : new Error(String(error)); + nativeCleanupFailure = encodeOpenClawStateWorkerError(failure, { + includeOrdinary: true, + }); + } } - throw new Error("SQLite worker operations must remain synchronous"); } - assertSettled(); }; if (request.type === "result-next") { if ( @@ -285,7 +349,7 @@ async function receive(request: SqliteWorkerRequest): Promise { } pendingInput = undefined; retire = false; - executeCommand(input.command); + await executeCommand(input.command); } else { inputNext = true; retire = false; @@ -355,7 +419,7 @@ async function receive(request: SqliteWorkerRequest): Promise { gatewayFences.get(request.actor)?.close(); gatewayFences.delete(request.actor); } else { - executeCommand(deserialize(request.input)); + await executeCommand(deserialize(request.input)); } const serialized = serialize(value); if (serialized.byteLength > SQLITE_WORKER_MAX_RESULT_BYTES) { @@ -377,6 +441,11 @@ async function receive(request: SqliteWorkerRequest): Promise { }; } } catch (error) { + if (preparedGatewayActor !== undefined) { + gatewayFences.get(preparedGatewayActor)?.close(); + gatewayFences.delete(preparedGatewayActor); + preparedGatewayActor = undefined; + } transfers.cancel(); pendingResult = undefined; pendingInput = undefined; @@ -390,7 +459,7 @@ async function receive(request: SqliteWorkerRequest): Promise { reply = { id: request.id, ok: false, - ...(retire ? { retire: true } : {}), + ...(retire || (nativeCleanupFailure && executed) ? { retire: true } : {}), ...(openNotEntered ? { openNotEntered: true } : {}), error: { name: executed ? "SqliteWorkerError" : failure.name, @@ -405,6 +474,8 @@ async function receive(request: SqliteWorkerRequest): Promise { maintenanceFence = undefined; lifecycle?.delegate.close(); lifecycle = undefined; + lifecyclePreparation?.port.close(); + lifecyclePreparation = undefined; operationAdmission?.port.close(); operationAdmission = undefined; } @@ -414,12 +485,29 @@ async function receive(request: SqliteWorkerRequest): Promise { drainProcessOutput(resolve); }); } + const complete = !reply.ok || (!pendingInput && !pendingResult); + if (complete && nativeCleanupFailure) { + reply.cleanupFailure = nativeCleanupFailure; + nativeCleanupFailure = undefined; + } + const resultPort = lifecycleReply?.actor === request.actor ? lifecycleReply.port : undefined; if (reply.ok) { const bytes = ownedWorkerBytes(reply.value); - port!.postMessage({ ...reply, value: bytes }, [bytes.buffer]); + const outgoing = { ...reply, value: bytes }; + if (resultPort) { + resultPort.postMessage({ type: "result", reply: outgoing }, [bytes.buffer]); + } else { + port!.postMessage(outgoing, [bytes.buffer]); + } + } else if (resultPort) { + resultPort.postMessage({ type: "result", reply }, []); } else { port!.postMessage(reply, []); } + if (complete) { + resultPort?.close(); + lifecycleReply = undefined; + } } // The broker sends one request at a time, including module initialization. diff --git a/src/infra/sqlite-transaction.ts b/src/infra/sqlite-transaction.ts index 4ed37219c562..92e2fbcd349a 100644 --- a/src/infra/sqlite-transaction.ts +++ b/src/infra/sqlite-transaction.ts @@ -107,6 +107,13 @@ export function retainSqliteWriteAdmissionService( }; } +/** Native coordinator waits must keep the same worker's current-authority grants serviceable. */ +export function sqliteWriteAdmissionServicesForLocation( + location: string, +): ReadonlySet<() => void> | undefined { + return writeAdmissionServices.get(normalizeWriteAdmissionLocation(location)); +} + type SqliteBeginAdmissionDiagnostics = { nativeAttempts: number; nativeMs: number; diff --git a/src/infra/sqlite-worker-broker-admission.ts b/src/infra/sqlite-worker-broker-admission.ts index 360f419115d3..cc704056889c 100644 --- a/src/infra/sqlite-worker-broker-admission.ts +++ b/src/infra/sqlite-worker-broker-admission.ts @@ -5,9 +5,6 @@ import { fileURLToPath, pathToFileURL } from "node:url"; import { serialize } from "node:v8"; import { INCOGNITO_AGENT_SQLITE_BASENAME } from "../state/openclaw-agent-db.paths.js"; import { OPENCLAW_SQLITE_BUSY_TIMEOUT_MS } from "../state/openclaw-state-db-contract.js"; -import { acquireWithWait } from "./acquire-with-wait.js"; -import { sleepWithAbort } from "./backoff.js"; -import { runWithSqliteCoordinator } from "./sqlite-coordinator.js"; import type { PreparedSqliteWorkerOpen, SqliteWorkerStoreOptions, @@ -18,8 +15,6 @@ import { readDatabasePathIdentity, type DatabasePathIdentity } from "./sqlite-wo import { createSqliteWorkerOperationAdmission } from "./sqlite-worker-operation-admission.js"; import type { SqliteWorkerStateContext } from "./sqlite-worker-state-context.js"; import { - acquireStateDatabaseCoordinator, - StateDatabaseCoordinatorContentionError, tryCreateGatewaySchemaFenceDelegate, tryCreateStateLifecycleDelegate, withStateDatabaseCoordinatorRuntimeDirectory, @@ -264,8 +259,7 @@ export function prepareSqliteWorkerLifecycle( job: Job, actor: Actor | undefined, assertDispatchable: () => void, - signal?: AbortSignal, -): void | Promise { +): void { const context = job.request.stateContext ?? actor?.stateContext; if (!actor || !context) { if (job.requireStateLifecycle) { @@ -298,48 +292,17 @@ export function prepareSqliteWorkerLifecycle( actorId: `${actor.id}:${job.request.id}`, }); if (!delegate && job.requireStateLifecycle) { - throw new Error("SQLite worker did not retain its required lifecycle custody"); + job.request.workerStateLifecycle = { + deadlineNs: + process.hrtime.bigint() + BigInt(OPENCLAW_SQLITE_BUSY_TIMEOUT_MS) * 1_000_000n, + }; } if (delegate) { job.stateLifecycle = { actor, delegate }; job.request.stateLifecycle = delegate.port; } }; - if (!job.requireStateLifecycle) { - prepare(); - return undefined; - } - let acquisitionError: unknown; - return acquireWithWait({ - deadlineMs: performance.now() + OPENCLAW_SQLITE_BUSY_TIMEOUT_MS, - pollIntervalMs: 25, - maxPollIntervalMs: 250, - sleep: (ms) => - sleepWithAbort(ms, signal).catch((error: unknown) => { - signal?.throwIfAborted(); - throw error; - }), - shouldRetry: (error) => - error === acquisitionError && - error instanceof StateDatabaseCoordinatorContentionError && - error.family === "state-lifecycle", - acquire: () => { - acquisitionError = undefined; - assertDispatchable(); - let coordinator: ReturnType; - try { - coordinator = acquireStateDatabaseCoordinator({ - databasePath: actor.databasePath, - busyTimeoutMs: 0, - }); - } catch (error) { - acquisitionError = error; - throw error; - } - // Only acquisition retries. Install delegates once, before releasing this real reference. - runWithSqliteCoordinator(coordinator, "SQLite worker lifecycle delegation", prepare); - }, - }); + prepare(); }); } diff --git a/src/infra/sqlite-worker-broker-reply.ts b/src/infra/sqlite-worker-broker-reply.ts index 87bda5d68d1d..f0125a7ba487 100644 --- a/src/infra/sqlite-worker-broker-reply.ts +++ b/src/infra/sqlite-worker-broker-reply.ts @@ -1,8 +1,12 @@ import { deserialize, serialize } from "node:v8"; -import type { Worker } from "node:worker_threads"; import { toErrorObject } from "@openclaw/normalization-core/error-coercion"; +import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { createDeferredCore } from "../shared/deferred.js"; -import { retainOpenClawStateWorkerErrorPayload } from "../state/openclaw-state-worker-error.js"; +import { + retainOpenClawStateWorkerErrorPayload, + hydrateOpenClawStateWorkerError, + type OpenClawStateWorkerErrorPayload, +} from "../state/openclaw-state-worker-error.js"; import { SqliteCoordinatorError } from "./sqlite-coordinator.js"; import { retainSqliteWriteAdmissionService } from "./sqlite-transaction.js"; import { @@ -17,6 +21,7 @@ import { type SqliteWorkerReply, type SqliteWorkerRequest, } from "./sqlite-worker-contract.js"; +import { createSqliteWorkerLifecyclePreparation } from "./sqlite-worker-lifecycle-preparation.js"; import type { SqliteWorkerOperationSettlement } from "./sqlite-worker-operation-settlement.js"; import { createSqliteWorkerTransferOwner, @@ -24,23 +29,29 @@ import { type SqliteWorkerTransferFrame, type SqliteWorkerTransferHandle, } from "./sqlite-worker-transfer.js"; +import { resolveStateDatabaseCoordinatorPath } from "./state-database-coordinator.js"; export function dispatchSqliteWorkerJob( slot: Slot, job: Job, onRejected: (error: unknown, retire: boolean) => void, ): void { - const reject = (error: unknown) => { + const reject = (error: unknown, preparedNotEntered = false) => { let failure = error; let retire = job.preparation - ? job.nativeDispatched === true + ? job.nativeDispatched === true || (job.requestPosted === true && !preparedNotEntered) : Boolean( job.request.gatewaySchemaFence || job.request.maintenanceSchemaFence || job.request.stateLifecycle || job.request.operationAdmission, ); - if (job.preparation && !job.nativeDispatched && slot.current === job) { + if ( + job.preparation && + !job.nativeDispatched && + (!job.requestPosted || preparedNotEntered) && + slot.current === job + ) { try { // No port reached native code. Release prepared custody before a follower can dispatch. releaseSqliteWorkerLifecycle(job); @@ -55,6 +66,7 @@ export function dispatchSqliteWorkerJob( } onRejected(failure, retire); }; + job.rejectPreparation = (error) => reject(error, true); if (job.requireStateLifecycle) { job.cancelPreparation = new AbortController(); } @@ -73,21 +85,15 @@ export function dispatchSqliteWorkerJob( const dispatch = () => { try { assertDispatchable(); - postSqliteWorkerJob(slot.worker, job, assertDispatchable, actor); - job.detach(); + postSqliteWorkerJob(slot, job, assertDispatchable, actor); } catch (error) { reject(error); } }; - const preparation = prepareSqliteWorkerLifecycle( - job, - actor, - assertDispatchable, - job.cancelPreparation?.signal, - ); - if (preparation) { - job.preparation = preparation; - void preparation.then(dispatch, reject); + prepareSqliteWorkerLifecycle(job, actor, assertDispatchable); + if (job.requireStateLifecycle && !job.request.workerStateLifecycle) { + job.preparation = Promise.resolve(); + void job.preparation.then(dispatch, reject); } else { dispatch(); } @@ -97,11 +103,83 @@ export function dispatchSqliteWorkerJob( } function postSqliteWorkerJob( - worker: Worker, + slot: Slot, job: Job, assertDispatchable: () => void, actor: Actor | undefined, ): void { + const dispatched = () => { + job.nativeDispatched = true; + job.detach(); + if (job.dispatchState) { + job.dispatchState.dispatched = true; + } + }; + if (job.request.workerStateLifecycle) { + const context = job.request.stateContext; + if (!actor || !context || !job.cancelPreparation) { + throw new Error("Worker lifecycle preparation requires its captured owner"); + } + const preparation = createSqliteWorkerLifecyclePreparation({ + assertCurrent: assertDispatchable, + signal: job.cancelPreparation.signal, + admit: () => prepareSqliteWorkerOperationAdmission(job, actor), + dispatch: dispatched, + receiveResult(reply, pumping) { + if (!isRecord(reply) || typeof reply.id !== "number" || typeof reply.ok !== "boolean") { + throw new Error("SQLite lifecycle reply is invalid"); + } + // SAFETY: This private port carries the same trusted worker reply as its message event. + slot.receiveReply(reply as SqliteWorkerReply, pumping); + }, + }); + const releaseService = retainSqliteWriteAdmissionService( + [ + resolveStateDatabaseCoordinatorPath({ + databasePath: actor.databasePath, + runtimeDirectory: context.coordinatorRuntime.directory, + uid: typeof process.getuid === "function" ? process.getuid() : undefined, + }), + ], + () => { + preparation.service(); + job.operationAdmission?.admission.service(); + }, + ); + job.lifecyclePreparation = { + get failure() { + return preparation.failure; + }, + finish() { + releaseService(); + preparation.finish(); + }, + }; + job.preparation = preparation.prepared; + job.request.lifecyclePreparation = preparation.port; + } else { + job.request.operationAdmission = prepareSqliteWorkerOperationAdmission(job, actor); + } + const request = prepareSqliteWorkerRequest(job); + assertDispatchable(); + // A throwing transfer may still have reached the worker; failure joins its exit. + if (!job.request.workerStateLifecycle) { + dispatched(); + } + job.requestPosted = true; + slot.worker.postMessage( + request, + [ + request.gatewaySchemaFence, + request.maintenanceSchemaFence, + request.stateLifecycle, + request.operationAdmission, + request.lifecyclePreparation, + ].filter((port) => port !== undefined), + ); +} + +function prepareSqliteWorkerOperationAdmission(job: Job, actor: Actor | undefined) { if (job.createAdmission) { const settlement = createDeferredCore(); job.settleNative = settlement.resolve; @@ -115,24 +193,9 @@ function postSqliteWorkerJob( () => retained.admission.service(), ), }; - job.request.operationAdmission = retained.admission.port; - } - const request = prepareSqliteWorkerRequest(job); - assertDispatchable(); - // A throwing transfer may still have reached the worker; failure joins its exit. - job.nativeDispatched = true; - worker.postMessage( - request, - [ - request.gatewaySchemaFence, - request.maintenanceSchemaFence, - request.stateLifecycle, - request.operationAdmission, - ].filter((port) => port !== undefined), - ); - if (job.dispatchState) { - job.dispatchState.dispatched = true; + return retained.admission.port; } + return undefined; } function prepareSqliteWorkerRequest(job: Job): SqliteWorkerRequest { @@ -152,7 +215,7 @@ function prepareSqliteWorkerRequest(job: Job): SqliteWorkerRequest { return { ...request, type: "execute-start", transfer }; } -export function decodeSqliteWorkerReplyValue( +function decodeSqliteWorkerReplyValue( job: Job, reply: Extract, ): @@ -237,7 +300,7 @@ export function decodeSqliteWorkerReplyValue( : { type: "complete", value }; } -export function decodeSqliteWorkerReplyError( +function decodeSqliteWorkerReplyError( job: Job, error: Extract["error"], ): Error { @@ -251,6 +314,116 @@ export function decodeSqliteWorkerReplyError( return failure; } +function decodeSqliteWorkerCleanupError(job: Job, payload: OpenClawStateWorkerErrorPayload): Error { + return hydrateOpenClawStateWorkerError( + decodeSqliteWorkerReplyError(job, { + name: "SqliteCoordinatorError", + message: "SQLite coordinator cleanup failed", + sharedState: payload, + }), + ); +} + +export function receiveSqliteWorkerReply( + slot: Pick & { worker: Pick }, + reply: SqliteWorkerReply, + owner: { + fail(reason: unknown, currentError?: Error, completed?: CompletedSqliteWorkerOutcome): void; + finish( + job: Job, + error?: unknown, + value?: unknown, + settlement?: SqliteWorkerOperationSettlement, + ): void; + dispatch(): void; + }, + pumping = false, +): void { + const job = slot.current; + if (!job || reply.id !== job.request.id) { + owner.fail(new Error("SQLite worker returned an unexpected response")); + return; + } + const settle = (operation: () => void) => { + if (pumping) { + queueMicrotask(() => { + if (slot.current === job && !slot.failed) { + operation(); + } + }); + } else { + operation(); + } + }; + if (!reply.ok) { + if (reply.cleanupFailure && job.nativeDispatched && !reply.retire) { + const original = + job.operationAdmission?.admission.failure ?? decodeSqliteWorkerReplyError(job, reply.error); + owner.fail(decodeSqliteWorkerCleanupError(job, reply.cleanupFailure), undefined, { + error: original, + }); + return; + } + if (reply.openNotEntered && job.request.type === "open" && job.dispatchState) { + job.dispatchState.openNotEntered = true; + } + const error = decodeSqliteWorkerReplyError(job, reply.error); + if (job.request.type === "open" && reply.openNotEntered && !reply.retire) { + settle(() => { + slot.current = undefined; + const refusal = job.operationAdmission?.admission.failure ?? error; + owner.finish(job, refusal, undefined, { kind: "not-entered", error: refusal }); + owner.dispatch(); + }); + return; + } + if (job.request.type !== "execute" || reply.retire) { + owner.fail(error, job.request.type !== "execute" ? error : undefined); + return; + } + if (job.lifecyclePreparation && !job.nativeDispatched) { + settle(() => { + job.lifecyclePreparation?.finish(); + job.rejectPreparation?.(job.lifecyclePreparation?.failure ?? error); + }); + return; + } + settle(() => { + slot.current = undefined; + owner.finish( + job, + job.lifecyclePreparation?.failure ?? job.operationAdmission?.admission.failure ?? error, + ); + owner.dispatch(); + }); + return; + } + let value: unknown; + try { + const result = decodeSqliteWorkerReplyValue(job, reply); + if (result.type === "continue") { + // Continuations retain the current job and its reserved transport credits through drain. + slot.worker.postMessage(result.request, []); + return; + } + value = result.value; + } catch (error) { + owner.fail(error); + return; + } + if (reply.cleanupFailure) { + owner.fail(decodeSqliteWorkerCleanupError(job, reply.cleanupFailure), undefined, { + value, + }); + return; + } + settle(() => { + slot.current = undefined; + owner.finish(job, undefined, value); + owner.dispatch(); + }); +} + /** Keep the original failure and outcome classification when retirement also fails. */ export function withSqliteWorkerCleanupFailure(failure: Error, cleanupError: unknown): Error { if (cleanupError === undefined) { @@ -264,12 +437,15 @@ export function withSqliteWorkerCleanupFailure(failure: Error, cleanupError: unk return retainSqliteWorkerErrorCode(combined, failure); } +export type CompletedSqliteWorkerOutcome = { value: unknown } | { error: unknown }; + export function settleFailedSqliteWorkerJobs({ queuedError, current, queued, error, currentError, + completed, retire, finish, }: { @@ -278,17 +454,33 @@ export function settleFailedSqliteWorkerJobs({ queued: Job[]; error: Error; currentError?: Error; + completed?: CompletedSqliteWorkerOutcome; retire: () => Promise; finish: typeof settleSqliteWorkerJob; }): void { current?.cancelPreparation?.abort(error); // Failed preparation must settle before retirement can release any actor custody. - const retirement = current?.preparation - ? current.preparation.catch(() => undefined).then(retire) - : retire(); + const retirement = current?.lifecyclePreparation + ? retire().finally(() => current.lifecyclePreparation?.finish()) + : current?.preparation + ? current.preparation.catch(() => undefined).then(retire) + : retire(); // Join native exit before releasing any operation that might have touched SQLite. const finishFailed = (cleanupError?: unknown) => { - if (current) { + if (current && completed) { + process.emitWarning( + new SqliteCoordinatorError( + "SQLite worker operation completed before coordinator cleanup failed", + withSqliteWorkerCleanupFailure(error, cleanupError), + ), + ); + finish( + current, + "error" in completed ? completed.error : undefined, + "value" in completed ? completed.value : undefined, + { kind: "completed" }, + ); + } else if (current) { finish( current, withSqliteWorkerCleanupFailure( @@ -325,6 +517,7 @@ export function settleSqliteWorkerJob( ); job.operationAdmission?.admission.finish(); job.operationAdmission?.releaseService(); + job.lifecyclePreparation?.finish(); let failure = error; const admissionCleanupFailures = job.operationAdmission?.admission.cleanupFailures ?? []; if (admissionCleanupFailures.length > 0) { diff --git a/src/infra/sqlite-worker-broker.ts b/src/infra/sqlite-worker-broker.ts index d789ebdeae52..5bad8579fa91 100644 --- a/src/infra/sqlite-worker-broker.ts +++ b/src/infra/sqlite-worker-broker.ts @@ -20,10 +20,10 @@ import { } from "./sqlite-worker-broker-admission.js"; import { settleSqliteWorkerJob, - decodeSqliteWorkerReplyError, - decodeSqliteWorkerReplyValue, dispatchSqliteWorkerJob, settleFailedSqliteWorkerJobs, + receiveSqliteWorkerReply, + type CompletedSqliteWorkerOutcome, } from "./sqlite-worker-broker-reply.js"; import type { Actor, @@ -388,8 +388,20 @@ export class SqliteWorkerBroker { }), ); const exited = createDeferredCore(); + const replyOwner = { + fail: (reason: unknown, currentError?: Error, completed?: CompletedSqliteWorkerOutcome) => + this.fail(slot, reason, currentError, completed), + finish: ( + job: Job, + error?: unknown, + value?: unknown, + settlement?: SqliteWorkerOperationSettlement, + ) => this.finish(job, error, value, settlement), + dispatch: () => this.dispatch(slot), + }; const slot: Slot = { worker, + receiveReply: (reply, pumping) => receiveSqliteWorkerReply(slot, reply, replyOwner, pumping), actors: new Set(), queue: [], exit: exited.promise, @@ -397,50 +409,7 @@ export class SqliteWorkerBroker { pendingOpens: 1, }; this.slots.add(slot); - worker.on("message", (reply: SqliteWorkerReply) => { - const job = slot.current; - if (!job || reply.id !== job.request.id) { - this.fail(slot, new Error("SQLite worker returned an unexpected response")); - return; - } - if (!reply.ok) { - if (reply.openNotEntered && job.request.type === "open" && job.dispatchState) { - job.dispatchState.openNotEntered = true; - } - const error = decodeSqliteWorkerReplyError(job, reply.error); - if (job.request.type === "open" && reply.openNotEntered && !reply.retire) { - slot.current = undefined; - const refusal = job.operationAdmission?.admission.failure ?? error; - this.finish(job, refusal, undefined, { kind: "not-entered", error: refusal }); - this.dispatch(slot); - return; - } - if (job.request.type !== "execute" || reply.retire) { - this.fail(slot, error, job.request.type !== "execute" ? error : undefined); - return; - } - slot.current = undefined; - this.finish(job, job.operationAdmission?.admission.failure ?? error); - this.dispatch(slot); - return; - } - let value: unknown; - try { - const result = decodeSqliteWorkerReplyValue(job, reply); - if (result.type === "continue") { - // Continuations retain the current job and its reserved transport credits through drain. - slot.worker.postMessage(result.request, []); - return; - } - value = result.value; - } catch (error) { - this.fail(slot, error); - return; - } - slot.current = undefined; - this.finish(job, undefined, value); - this.dispatch(slot); - }); + worker.on("message", (reply: SqliteWorkerReply) => slot.receiveReply(reply)); worker.on("error", (error) => this.fail(slot, error)); worker.on("messageerror", (error) => this.fail(slot, error)); worker.once("exit", (code) => { @@ -575,12 +544,20 @@ export class SqliteWorkerBroker { settleSqliteWorkerJob(job, error, value, settlement); } - private fail(slot: Slot, reason: unknown, currentError?: Error): void { + private fail( + slot: Slot, + reason: unknown, + currentError?: Error, + completed?: CompletedSqliteWorkerOutcome, + ): void { if (slot.failed) { return; } const error = toErrorObject(reason, "SQLite worker failed"); slot.failed = new SqliteWorkerError(error.message, "unavailable"); + if (completed) { + slot.retiredAfterCompletion = true; + } const current = slot.current; slot.current = undefined; if (current) { @@ -595,6 +572,7 @@ export class SqliteWorkerBroker { queued, error, currentError, + completed, retire: () => this.retire(slot), finish: (job, failure, value, settlement) => this.finish(job, failure, value, settlement), }); @@ -625,7 +603,7 @@ export class SqliteWorkerBroker { this.fail(actor.slot, error instanceof Error ? error : new Error(String(error))); await actor.slot.exit; } - } else if (firstAttempt && actor.slot.failed) { + } else if (firstAttempt && actor.slot.failed && !actor.slot.retiredAfterCompletion) { errors.push(actor.slot.failed); } try { diff --git a/src/infra/sqlite-worker-broker.types.ts b/src/infra/sqlite-worker-broker.types.ts index 45426d6e41cc..4ba8aa9a7e92 100644 --- a/src/infra/sqlite-worker-broker.types.ts +++ b/src/infra/sqlite-worker-broker.types.ts @@ -1,6 +1,6 @@ import type { Worker } from "node:worker_threads"; import type { OpenClawDatabaseMaintenanceScope } from "../state/openclaw-state-db-async-lifecycle.js"; -import type { SqliteWorkerRequest } from "./sqlite-worker-contract.js"; +import type { SqliteWorkerRequest, SqliteWorkerReply } from "./sqlite-worker-contract.js"; import type { SqliteWorkerAdmissionFactory, SqliteWorkerOperationAdmission, @@ -33,7 +33,10 @@ export type Job = { operationAdmission?: { admission: SqliteWorkerOperationAdmission; releaseService(): void }; settleNative?: (settlement: SqliteWorkerOperationSettlement) => void; nativeDispatched?: boolean; + requestPosted?: boolean; + rejectPreparation?: (error: unknown) => void; preparation?: Promise; + lifecyclePreparation?: { failure: unknown; finish(): void }; cancelPreparation?: AbortController; stateLifecycle?: { actor: Actor; delegate: StateLifecycleDelegate }; assertCurrent?: () => void; @@ -55,10 +58,12 @@ export type Job = { }; export type Slot = { worker: Worker; + receiveReply(reply: SqliteWorkerReply, pumping?: boolean): void; actors: Set; queue: Job[]; current?: Job; failed?: Error; + retiredAfterCompletion?: true; retiring?: Promise; exit: Promise; exited: boolean; diff --git a/src/infra/sqlite-worker-contract.ts b/src/infra/sqlite-worker-contract.ts index 3a1b9fc3a3fc..46762de8f4fb 100644 --- a/src/infra/sqlite-worker-contract.ts +++ b/src/infra/sqlite-worker-contract.ts @@ -30,6 +30,8 @@ export type SqliteWorkerRequest = { gatewaySchemaFence?: MessagePort; maintenanceSchemaFence?: MessagePort; stateLifecycle?: MessagePort; + workerStateLifecycle?: { deadlineNs: bigint }; + lifecyclePreparation?: MessagePort; operationAdmission?: MessagePort; } & ( | { @@ -49,6 +51,7 @@ export type SqliteWorkerRequest = { export type SqliteWorkerReply = { id: number; + cleanupFailure?: OpenClawStateWorkerErrorPayload; } & ( | { ok: true; value: Uint8Array; transfer?: "start" | "frame"; input?: "next" } | { diff --git a/src/infra/sqlite-worker-lifecycle-cleanup.test.ts b/src/infra/sqlite-worker-lifecycle-cleanup.test.ts new file mode 100644 index 000000000000..dd9fd11956c2 --- /dev/null +++ b/src/infra/sqlite-worker-lifecycle-cleanup.test.ts @@ -0,0 +1,253 @@ +import { createHash } from "node:crypto"; +import fs from "node:fs"; +import { MessagePort, Worker } from "node:worker_threads"; +import { isRecord } from "@openclaw/normalization-core/record-coerce"; +import { afterEach, expect, it, vi } from "vitest"; +import { createDeferredCore } from "../shared/deferred.js"; +import { + closeOpenClawStateDatabaseAsync, + runOpenClawStateWriteTransaction, +} from "../state/openclaw-state-db.js"; +import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.js"; +import { runOpenClawStateWorkerOperation } from "../state/openclaw-state-worker-store.js"; +import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js"; +import * as tokens from "./device-auth-store.js"; +import { storeDeviceAuthTokenInDatabase } from "./device-auth-store.kernel.js"; +import { sqliteWorkerPreloadEnv } from "./sqlite-worker-preload.test-support.js"; +import { + captureStateDatabaseCoordinatorRuntime, + resolveStateDatabaseCoordinatorPath, +} from "./state-database-coordinator.js"; + +afterEach(() => { + vi.restoreAllMocks(); + vi.unstubAllEnvs(); +}); + +it.each([ + { length: 32, ownerCurrent: true, queuedFollower: true, nested: false, preparation: false }, + { + length: 65 * 1024 * 1024, + ownerCurrent: true, + queuedFollower: true, + nested: false, + preparation: false, + }, + { length: 32, ownerCurrent: false, queuedFollower: true, nested: false, preparation: false }, + { length: 32, ownerCurrent: true, queuedFollower: false, nested: false, preparation: false }, + { length: 32, ownerCurrent: true, queuedFollower: false, nested: true, preparation: false }, + { length: 32, ownerCurrent: false, queuedFollower: false, nested: false, preparation: true }, +])( + "preserves a settled $length-byte token outcome through cleanup failure (owner current: $ownerCurrent, queued follower: $queuedFollower, nested: $nested, preparation: $preparation)", + async ({ length, ownerCurrent, queuedFollower, nested, preparation }) => { + await withOpenClawTestState({ label: "worker-lifecycle-cleanup" }, async (state) => { + const databasePath = state.statePath("state", "openclaw.sqlite"); + const coordinatorPath = resolveStateDatabaseCoordinatorPath({ + databasePath, + runtimeDirectory: captureStateDatabaseCoordinatorRuntime().directory, + uid: process.getuid?.(), + }); + const entered = state.path("transaction-entered"); + const failed = state.path("cleanup-failed"); + const preload = state.path("cleanup-preload.cjs"); + fs.writeFileSync( + preload, + ` +const { MessagePort, isMainThread } = require("node:worker_threads"); +if (!isMainThread) { + const fs = require("node:fs"); + const { DatabaseSync } = require("node:sqlite"); + let transactions = 0; + let armed = false; + let failedDatabase; + const post = MessagePort.prototype.postMessage; + MessagePort.prototype.postMessage = function(message, ...rest) { + if (message?.stage === "transaction") transactions += 1; + if (${preparation} + ? message?.type === "acquired" && transactions === 1 + : message?.stage === "transaction" && transactions === 2) { + armed = true; + fs.writeFileSync(${JSON.stringify(entered)}, "entered"); + } + return Reflect.apply(post, this, [message, ...rest]); + }; + const exec = DatabaseSync.prototype.exec; + DatabaseSync.prototype.exec = function(sql) { + if (armed && sql === "ROLLBACK" && this.location() === ${JSON.stringify(coordinatorPath)}) { + armed = false; + failedDatabase = this; + fs.writeFileSync(${JSON.stringify(failed)}, "one cleanup failure"); + throw new Error("Synthetic coordinator rollback failure"); + } + return Reflect.apply(exec, this, [sql]); + }; + const close = DatabaseSync.prototype.close; + DatabaseSync.prototype.close = function(...args) { + if (this === failedDatabase) { + failedDatabase = undefined; + throw new Error("Synthetic coordinator close failure"); + } + return Reflect.apply(close, this, args); + }; +} +`, + ); + for (const [key, value] of Object.entries(sqliteWorkerPreloadEnv(preload))) { + vi.stubEnv(key, value); + } + const lookup = { deviceId: "synthetic-cleanup-device", role: "operator", env: state.env }; + const posts = vi.spyOn(Worker.prototype, "postMessage"); + await tokens.storeDeviceAuthToken({ ...lookup, token: "synthetic-before" }); + const worker = posts.mock.contexts[0]; + posts.mockRestore(); + if (!(worker instanceof Worker)) { + throw new Error("Expected the shared-state worker"); + } + let exited = false; + worker.once("exit", () => { + exited = true; + }); + const warnings = vi.spyOn(process, "emitWarning").mockImplementation(() => {}); + // Keep the native implementation callable with the actual sending port. + const nativePost = vi.spyOn(MessagePort.prototype, "postMessage"); + nativePost.mockRestore(); + let nativeWrites = 0; + let current = true; + const refused = new Error("Synthetic token authority revoked"); + const dispatch = vi + .spyOn(MessagePort.prototype, "postMessage") + .mockImplementation(function (this: MessagePort, message, transferList) { + const result = nativePost.call(this, message, transferList); + if ( + !isRecord(message) || + message.type !== "accepted" || + (preparation + ? message.admission !== undefined + : !(message.admission instanceof MessagePort)) + ) { + return result; + } + dispatch.mockRestore(); + const deadline = Date.now() + 5_000; + const pause = new Int32Array(new SharedArrayBuffer(4)); + while (!fs.existsSync(entered)) { + if (Date.now() >= deadline) { + throw new Error("Worker did not reach its transaction grant"); + } + Atomics.wait(pause, 0, 0, 1); + } + current = ownerCurrent; + runOpenClawStateWriteTransaction( + ({ db }) => { + nativeWrites += 1; + storeDeviceAuthTokenInDatabase(db, { + deviceId: "synthetic-native-device", + role: "operator", + token: "synthetic-native-token", + }); + }, + { env: state.env }, + ); + return result; + }); + const token = "x".repeat(length); + const mutate = () => + tokens.storeDeviceAuthToken({ + ...lookup, + token, + expectedToken: "synthetic-before", + assertCurrent() { + if (!current) { + throw refused; + } + }, + }); + const escape = createDeferredCore(); + void escape.promise.catch(() => {}); + const mutation = nested + ? runOpenClawStateWorkerOperation( + captureOpenClawStateWorkerContext({ env: state.env }), + () => + Promise.race([ + (async () => { + const result = await mutate(); + await expect(tokens.loadDeviceAuthToken(lookup)).rejects.toMatchObject({ + code: "unavailable", + }); + return result; + })(), + escape.promise, + ]), + ) + : mutate(); + let completed = false; + void mutation.then( + () => { + completed = true; + }, + () => { + completed = true; + }, + ); + const follower = queuedFollower + ? tokens.loadDeviceAuthToken(lookup).then( + () => ({ failed: false, exited }), + (error: unknown) => ({ failed: true, exited, error }), + ) + : undefined; + try { + if (nested) { + await expect.poll(() => completed, { timeout: 5_000 }).toBe(true); + } + const hash = (value: string) => createHash("sha256").update(value).digest("hex"); + if (ownerCurrent) { + const result = await mutation; + expect(result?.token.length).toBe(length); + expect(hash(result?.token ?? "")).toBe(hash(token)); + } else if (preparation) { + const failure: unknown = await mutation.catch((error: unknown) => error); + expect(failure).toBeInstanceOf(Error); + expect(failure instanceof AggregateError ? failure.errors[0] : failure).toMatchObject({ + code: "unavailable", + }); + } else { + await expect(mutation).rejects.toBe(refused); + } + expect(exited).toBe(true); + expect(fs.readFileSync(failed, "utf8")).toBe("one cleanup failure"); + expect(nativeWrites).toBe(1); + if (follower) { + expect(await follower).toMatchObject({ + failed: true, + exited: true, + error: { code: "unavailable" }, + }); + } + if (!preparation) { + expect(warnings).toHaveBeenCalledWith( + expect.objectContaining({ + message: "SQLite worker operation completed before coordinator cleanup failed", + }), + ); + } + expect(hash((await tokens.loadDeviceAuthToken(lookup))?.token ?? "")).toBe( + hash(ownerCurrent ? token : "synthetic-before"), + ); + await closeOpenClawStateDatabaseAsync(); + expect(hash((await tokens.loadDeviceAuthToken(lookup))?.token ?? "")).toBe( + hash(ownerCurrent ? token : "synthetic-before"), + ); + expect( + await tokens.loadDeviceAuthToken({ ...lookup, deviceId: "synthetic-native-device" }), + ).toMatchObject({ token: "synthetic-native-token" }); + } finally { + escape.reject(new Error("Release nested fixture after observation")); + dispatch.mockRestore(); + await Promise.allSettled([mutation, follower]); + await closeOpenClawStateDatabaseAsync(); + warnings.mockRestore(); + vi.unstubAllEnvs(); + } + }); + }, +); diff --git a/src/infra/sqlite-worker-lifecycle-preparation.ts b/src/infra/sqlite-worker-lifecycle-preparation.ts new file mode 100644 index 000000000000..d6747643842c --- /dev/null +++ b/src/infra/sqlite-worker-lifecycle-preparation.ts @@ -0,0 +1,174 @@ +import { MessageChannel, MessagePort, receiveMessageOnPort } from "node:worker_threads"; +import { isRecord } from "@openclaw/normalization-core/record-coerce"; +import { createDeferredCore } from "../shared/deferred.js"; +import { acquireWithWait } from "./acquire-with-wait.js"; +import { sleepWithAbort } from "./backoff.js"; +import { + SqliteCoordinatorError, + createSqliteLifecycleAggregateError, +} from "./sqlite-coordinator.js"; +import { + acquireStateDatabaseCoordinator, + StateDatabaseCoordinatorContentionError, + withStateDatabaseCoordinatorRuntimeDirectory, + type StateDatabaseCoordinatorRuntime, +} from "./state-database-coordinator.js"; + +/** Service only this job's preparation while a legacy native caller owns the host turn. */ +export function createSqliteWorkerLifecyclePreparation(params: { + assertCurrent(): void; + admit(): MessagePort | undefined; + dispatch(): void; + receiveResult(reply: unknown, pumping: boolean): void; + signal: AbortSignal; +}) { + const { port1, port2 } = new MessageChannel(); + const prepared = createDeferredCore(); + let failure: unknown; + let finished = false; + let admitted = false; + const refuse = (error: unknown) => { + failure ??= error; + port1.postMessage({ type: "cancel" }); + }; + const receive = (request: unknown, pumping = false) => { + if (finished) { + return; + } + if (isRecord(request) && request.type === "result") { + params.receiveResult(request.reply, pumping); + return; + } + if (!isRecord(request) || (request.type !== "check" && request.type !== "acquired")) { + refuse(new SqliteCoordinatorError("SQLite lifecycle preparation request is invalid")); + return; + } + try { + params.signal.throwIfAborted(); + params.assertCurrent(); + if (admitted) { + throw new SqliteCoordinatorError("SQLite lifecycle preparation already admitted its job"); + } + if (request.type === "check") { + port1.postMessage({ type: "accepted" }); + return; + } + const admission = params.admit(); + params.signal.throwIfAborted(); + params.assertCurrent(); + params.dispatch(); + admitted = true; + port1.postMessage({ type: "accepted", admission }, admission ? [admission] : []); + prepared.resolve(); + } catch (error) { + refuse(error); + } + }; + const abort = () => refuse(params.signal.reason); + port1.on("message", receive); + port1.once("close", () => prepared.resolve()); + port1.unref(); + params.signal.addEventListener("abort", abort, { once: true }); + return { + port: port2, + prepared: prepared.promise, + get failure() { + return failure; + }, + service() { + for (let queued = receiveMessageOnPort(port1); queued; queued = receiveMessageOnPort(port1)) { + receive(queued.message, true); + } + }, + finish() { + finished = true; + params.signal.removeEventListener("abort", abort); + prepared.resolve(); + port1.close(); + port2.close(); + }, + }; +} + +/** The executing data worker owns acquisition, its synchronous command, and native release. */ +export async function acquireSqliteWorkerLifecycle(params: { + port: MessagePort; + databasePath: string; + deadlineNs: bigint; + runtime: StateDatabaseCoordinatorRuntime; + onUnsettled(): void; +}) { + const controller = new AbortController(); + let waiting: ReturnType> | undefined; + const cancel = () => { + const error = new SqliteCoordinatorError("SQLite lifecycle preparation was canceled"); + controller.abort(error); + waiting?.reject(error); + }; + const receive = (reply: unknown) => { + if (!isRecord(reply) || reply.type !== "accepted") { + cancel(); + return; + } + if (!waiting) { + cancel(); + return; + } + const pending = waiting; + waiting = undefined; + if (reply.admission !== undefined && !(reply.admission instanceof MessagePort)) { + pending.reject(new SqliteCoordinatorError("SQLite lifecycle admission port is invalid")); + return; + } + pending.resolve(reply.admission); + }; + params.port.on("message", receive); + params.port.once("close", cancel); + const check = (type: "check" | "acquired") => { + controller.signal.throwIfAborted(); + waiting = createDeferredCore(); + params.port.postMessage({ type }, []); + return waiting.promise; + }; + let coordinator: ReturnType | undefined; + try { + const held = await acquireWithWait({ + deadlineMs: + performance.now() + Number(params.deadlineNs - process.hrtime.bigint()) / 1_000_000, + pollIntervalMs: 25, + maxPollIntervalMs: 250, + sleep: (ms) => sleepWithAbort(ms, controller.signal), + shouldRetry: (error) => + error instanceof StateDatabaseCoordinatorContentionError && + error.family === "state-lifecycle", + acquire: async () => { + await check("check"); + controller.signal.throwIfAborted(); + return withStateDatabaseCoordinatorRuntimeDirectory(params.runtime, () => + acquireStateDatabaseCoordinator({ databasePath: params.databasePath, busyTimeoutMs: 0 }), + ); + }, + }); + coordinator = held; + const admission = await check("acquired"); + controller.signal.throwIfAborted(); + coordinator = undefined; + return { coordinator: held, admission }; + } catch (error) { + try { + coordinator?.release(); + } catch (cleanupError) { + params.onUnsettled(); + throw createSqliteLifecycleAggregateError( + [error, cleanupError], + "SQLite lifecycle preparation cleanup failed", + error, + ); + } + throw error; + } finally { + params.port.off("message", receive); + params.port.off("close", cancel); + // The caller retains this port through terminal replies, including failed native cleanup. + } +} diff --git a/src/infra/sqlite-worker-shared-state-admission.test-support.ts b/src/infra/sqlite-worker-shared-state-admission.test-support.ts index dcf50b57dfec..4893c1e22145 100644 --- a/src/infra/sqlite-worker-shared-state-admission.test-support.ts +++ b/src/infra/sqlite-worker-shared-state-admission.test-support.ts @@ -1,7 +1,7 @@ import { once } from "node:events"; import type { DatabaseSync } from "node:sqlite"; import { deserialize } from "node:v8"; -import { Worker, type Transferable } from "node:worker_threads"; +import { Worker } from "node:worker_threads"; import { expect, it, vi } from "vitest"; import type { NativeHookRelayBridgeRecord } from "../agents/harness/native-hook-relay-bridge-record.js"; import { closeOpenClawStateDatabaseAsync } from "../state/openclaw-state-db-cache.js"; @@ -11,12 +11,10 @@ import { runOpenClawStateWorkerOperation, } from "../state/openclaw-state-worker-store.js"; import * as nodeSqlite from "./node-sqlite.js"; -import type { SqliteWorkerRequest } from "./sqlite-worker-contract.js"; import { createSqliteWorkerOperationAdmission } from "./sqlite-worker-operation-admission.js"; import { acquireStateDatabaseCoordinator, acquireGatewayLifecycleCoordinator, - retainHeldStateDatabaseCoordinator, resolveStateDatabaseCoordinatorPath, withStateDatabaseCoordinatorRuntimeDirectory, } from "./state-database-coordinator.js"; @@ -25,7 +23,7 @@ export function registerSharedStateWorkerAdmissionTests( createContext: () => OpenClawStateWorkerContext, ): void { it.each([undefined, false])( - "lets host admission share physical lifecycle custody (explicit=%s)", + "borrows an already-held parent lifecycle owner for nested native admission (explicit=%s)", async (requireStateLifecycle) => { const captured = createContext(); const record: NativeHookRelayBridgeRecord = { @@ -36,35 +34,42 @@ export function registerSharedStateWorkerAdmissionTests( token: "synthetic-test-token", expiresAtMs: 20000, }; - let grants = 0; - await runOpenClawStateWorkerOperation( - captured, - (scope) => - scope.execute({ type: "nativeHookRelay.write", input: { record, updatedAtMs: 1 } }), - { - requireStateLifecycle, - createAdmission: () => ({ - nativeLocations: [captured.admission.databasePath], - admission: createSqliteWorkerOperationAdmission((request, grant) => { - expect(request.stage).toBe("transaction"); - // The worker already holds its write transaction and awaits this grant. - // An independent native acquisition must borrow the same physical owner. - const nested = withStateDatabaseCoordinatorRuntimeDirectory( - captured.coordinatorRuntime, - () => - acquireStateDatabaseCoordinator({ - databasePath: captured.admission.databasePath, - busyTimeoutMs: 0, - }), - ); - nested.release(); - captured.admission.assertCurrent(); - expect(grant()).toBe(true); - grants += 1; - }), - }), - }, + const parent = withStateDatabaseCoordinatorRuntimeDirectory(captured.coordinatorRuntime, () => + acquireStateDatabaseCoordinator({ databasePath: captured.admission.databasePath }), ); + let grants = 0; + try { + await runOpenClawStateWorkerOperation( + captured, + (scope) => + scope.execute({ type: "nativeHookRelay.write", input: { record, updatedAtMs: 1 } }), + { + requireStateLifecycle, + createAdmission: () => ({ + nativeLocations: [captured.admission.databasePath], + admission: createSqliteWorkerOperationAdmission((request, grant) => { + expect(request.stage).toBe("transaction"); + // The worker already holds its write transaction and awaits this grant. + // An independent native acquisition must borrow the same physical owner. + const nested = withStateDatabaseCoordinatorRuntimeDirectory( + captured.coordinatorRuntime, + () => + acquireStateDatabaseCoordinator({ + databasePath: captured.admission.databasePath, + busyTimeoutMs: 0, + }), + ); + nested.release(); + captured.admission.assertCurrent(); + expect(grant()).toBe(true); + grants += 1; + }), + }), + }, + ); + } finally { + parent.release(); + } expect(grants).toBe(1); expect( await executeOpenClawStateWorker(captured, { @@ -75,44 +80,37 @@ export function registerSharedStateWorkerAdmissionTests( }, ); - it("retains explicitly requested lifecycle custody without a host grant factory", async () => { + it("fences explicitly requested worker execution without a host grant factory", async () => { const captured = createContext(); - const messages = vi.spyOn(Worker.prototype, "postMessage"); await executeOpenClawStateWorker(captured, { type: "flows.list", input: { ownerKey: "agent:main:main" }, }); - const worker = messages.mock.contexts[0]; - messages.mockRestore(); - if (!(worker instanceof Worker)) { - throw new Error("Expected the shared-state worker"); - } - const postMessage = worker.postMessage.bind(worker); - let observed = false; - vi.spyOn(worker, "postMessage").mockImplementation( - (request: SqliteWorkerRequest, transferList?: readonly Transferable[]) => { - if (request.type === "execute") { - const held = withStateDatabaseCoordinatorRuntimeDirectory( - captured.coordinatorRuntime, - () => retainHeldStateDatabaseCoordinator(captured.admission.databasePath), - ); - expect(held).toBeDefined(); - held?.release(); - observed = true; - } - return postMessage(request, transferList); - }, - ); - await runOpenClawStateWorkerOperation( + const foreign = await holdForeignLifecycle(captured); + let checks = 0; + let completed = false; + const result = runOpenClawStateWorkerOperation( captured, - (scope) => - scope.execute({ - type: "flows.list", - input: { ownerKey: "agent:main:main" }, - }), - { requireStateLifecycle: true }, - ); - expect(observed).toBe(true); + (scope) => scope.execute({ type: "flows.list", input: { ownerKey: "agent:main:main" } }), + { + requireStateLifecycle: true, + assertCurrent: () => { + checks += 1; + }, + }, + ).then((value) => { + completed = true; + return value; + }); + try { + await vi.waitFor(() => expect(checks).toBeGreaterThan(4)); + expect(completed).toBe(false); + foreign.release(); + await expect(result).resolves.toEqual([]); + } finally { + await foreign.close(); + await result; + } }); it("waits for a bounded foreign lifecycle owner before an admitted write", async () => { @@ -225,7 +223,9 @@ export function registerSharedStateWorkerAdmissionTests( }, ); expect(createAdmission).not.toHaveBeenCalled(); - expect(posts.mock.calls.filter(([request]) => request.type === "execute")).toHaveLength(0); + expect(posts.mock.calls.some(([request]) => request.operationAdmission !== undefined)).toBe( + false, + ); } finally { posts.mockRestore(); await foreign.close(); @@ -248,9 +248,13 @@ export function registerSharedStateWorkerAdmissionTests( }, ); - it.each(["after-acquire", "admission-factory", "cleanup-failure"] as const)( - "settles a canceled head and its follower at %s before posting", - async (timing) => { + it.each( + (["after-acquire", "admission-factory", "cleanup-failure"] as const).flatMap((timing) => + [false, true].map((parentHeld) => ({ timing, parentHeld })), + ), + )( + "settles a canceled head and its follower at $timing before native execution (parent held: $parentHeld)", + async ({ timing, parentHeld }) => { const captured = createContext(); await executeOpenClawStateWorker(captured, { type: "flows.list", @@ -275,6 +279,11 @@ export function registerSharedStateWorkerAdmissionTests( gateway.release(); throw new Error("Expected the native Gateway coordinator"); } + const parent = parentHeld + ? withStateDatabaseCoordinatorRuntimeDirectory(captured.coordinatorRuntime, () => + acquireStateDatabaseCoordinator({ databasePath: captured.admission.databasePath }), + ) + : undefined; const canceled = new AbortController(); const stopped = new Error("synthetic cancellation before post"); let factories = 0; @@ -340,11 +349,13 @@ export function registerSharedStateWorkerAdmissionTests( }, ); const writes = posts.mock.calls.filter(([request]) => request.type === "execute"); - expect(writes).toHaveLength(timing === "cleanup-failure" ? 0 : 1); + expect(factories).toBe( + timing === "after-acquire" ? 1 : timing === "cleanup-failure" ? 1 : 2, + ); if (timing === "cleanup-failure") { return; } - expect(writes[0]?.[0].gatewaySchemaFence).toBeDefined(); + expect(writes.some(([request]) => request.gatewaySchemaFence !== undefined)).toBe(true); expect( await executeOpenClawStateWorker(captured, { type: "nativeHookRelay.read", @@ -359,7 +370,11 @@ export function registerSharedStateWorkerAdmissionTests( gateway.release(); } finally { closes.mockRestore(); - gateway.release(); + try { + gateway.release(); + } finally { + parent?.release(); + } } } }, @@ -437,7 +452,9 @@ export function registerSharedStateWorkerAdmissionTests( type: string; input: { record: NativeHookRelayBridgeRecord }; }; - return command.type === "nativeHookRelay.write" ? [command.input.record.pid] : []; + return command.type === "nativeHookRelay.write" && command.input.record.pid !== 0 + ? [command.input.record.pid] + : []; }); expect(dispatched).toEqual(Array.from({ length: 256 }, (_, index) => index + 1)); expect(factories).toBe(256); diff --git a/src/infra/sqlite-worker-shared-state-admission.test.ts b/src/infra/sqlite-worker-shared-state-admission.test.ts index ad8ab851a622..8f9c1e4c8429 100644 --- a/src/infra/sqlite-worker-shared-state-admission.test.ts +++ b/src/infra/sqlite-worker-shared-state-admission.test.ts @@ -1,7 +1,6 @@ import fs from "node:fs"; import path from "node:path"; -import { deserialize } from "node:v8"; -import { Worker } from "node:worker_threads"; +import { MessagePort, Worker } from "node:worker_threads"; import { asOptionalRecord } from "@openclaw/normalization-core/record-coerce"; import { afterEach, expect, it, vi } from "vitest"; import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; @@ -76,33 +75,32 @@ if (!isMainThread) { fs.rmSync(marker); let current = true; const mutation = vi.fn(() => "committed"); - const nativePost = worker.postMessage.bind(worker); - const dispatch = vi.spyOn(worker, "postMessage").mockImplementation((message, transferList) => { - const result = nativePost(message, transferList); - const request = asOptionalRecord(message); - if ( - request?.type !== "execute" || - !request.operationAdmission || - !(request.input instanceof Uint8Array) || - asOptionalRecord(deserialize(request.input))?.type !== "nativeHookRelay.renew" - ) { - return result; - } - dispatch.mockRestore(); - // Keep the host listener queued until the worker owns its transaction, - // then enter the competing synchronous write on the same event-loop turn. - const deadline = Date.now() + 5_000; - const pause = new Int32Array(new SharedArrayBuffer(4)); - while (!fs.existsSync(marker)) { - if (Date.now() >= deadline) { - throw new Error("Worker did not reach its transaction admission"); + // Keep the native implementation callable with the actual sending port. + const nativePost = vi.spyOn(MessagePort.prototype, "postMessage"); + nativePost.mockRestore(); + const dispatch = vi + .spyOn(MessagePort.prototype, "postMessage") + .mockImplementation(function (this: MessagePort, message, transferList) { + const result = nativePost.call(this, message, transferList); + const request = asOptionalRecord(message); + if (request?.type !== "accepted" || !(request.admission instanceof MessagePort)) { + return result; } - Atomics.wait(pause, 0, 0, 1); - } - current = ownerCurrent; - runOpenClawStateWriteTransaction(mutation, { path: stateDbPath }); - return result; - }); + dispatch.mockRestore(); + // Keep the host listener queued until the worker owns its transaction, + // then enter the competing synchronous write on the same event-loop turn. + const deadline = Date.now() + 5_000; + const pause = new Int32Array(new SharedArrayBuffer(4)); + while (!fs.existsSync(marker)) { + if (Date.now() >= deadline) { + throw new Error("Worker did not reach its transaction admission"); + } + Atomics.wait(pause, 0, 0, 1); + } + current = ownerCurrent; + runOpenClawStateWriteTransaction(mutation, { path: stateDbPath }); + return result; + }); const expiresAtMs = record.expiresAtMs + 1; const renewal = renewOrRestoreNativeHookRelayBridgeRecord({ record: { ...record, expiresAtMs }, diff --git a/src/infra/sqlite-worker-shared-state-idle-fixture.test-support.ts b/src/infra/sqlite-worker-shared-state-idle-fixture.test-support.ts index 23187ed5c2ba..2eded883a120 100644 --- a/src/infra/sqlite-worker-shared-state-idle-fixture.test-support.ts +++ b/src/infra/sqlite-worker-shared-state-idle-fixture.test-support.ts @@ -1,6 +1,7 @@ import { openOpenClawStateDatabase } from "../state/openclaw-state-db.js"; import { createSqliteWorkerBackend as createCanonicalBackend } from "../state/openclaw-state.worker.js"; import { getSqliteWorkerStateContext } from "./sqlite-worker-state-context.js"; +import { retainHeldStateDatabaseCoordinator } from "./state-database-coordinator.js"; /** Exercise canonical actor retirement with real native reader and transaction faults. */ export function createSqliteWorkerBackend(input: undefined, context: { databasePath: string }) { @@ -14,6 +15,13 @@ export function createSqliteWorkerBackend(input: undefined, context: { databaseP return { ...backend, execute(command: Parameters[0]) { + if (command.type === "database.inspectIdle") { + const held = retainHeldStateDatabaseCoordinator(context.databasePath); + if (!held) { + throw new Error("Idle inspection requires its executing worker's lifecycle custody"); + } + held.release(); + } if (command.type === "database.inspectIdle" && unsettleInspection) { db.exec("BEGIN"); } diff --git a/src/infra/sqlite-worker-shared-state-idle.test.ts b/src/infra/sqlite-worker-shared-state-idle.test.ts index b32512ffd484..352204b28d73 100644 --- a/src/infra/sqlite-worker-shared-state-idle.test.ts +++ b/src/infra/sqlite-worker-shared-state-idle.test.ts @@ -2,7 +2,8 @@ import { once } from "node:events"; import { performance } from "node:perf_hooks"; import { setImmediate as nextTurn } from "node:timers/promises"; import { deserialize } from "node:v8"; -import { Worker, type Transferable } from "node:worker_threads"; +import { MessagePort, Worker, type Transferable } from "node:worker_threads"; +import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { afterEach, expect, it, vi } from "vitest"; import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; import { createDeferredCore } from "../shared/deferred.js"; @@ -16,10 +17,7 @@ import { requireNodeSqlite } from "./node-sqlite.js"; import { runtimeProcessEntrypoints } from "./runtime-process-entrypoints.js"; import * as runtimeWorker from "./runtime-worker-url.js"; import type { SqliteWorkerRequest } from "./sqlite-worker-contract.js"; -import { - retainHeldStateDatabaseCoordinator, - withStateDatabaseCoordinatorRuntimeDirectory, -} from "./state-database-coordinator.js"; +import * as sqliteWorkers from "./sqlite-worker-store.js"; const dirs = useAutoCleanupTempDirTracker((cleanup) => afterEach(async () => { @@ -39,19 +37,13 @@ async function fixture(mode: "healthy" | "local-reader" | "unsettled-inspection" vi.spyOn(performance, "now").mockImplementation(() => now() + elapsed); const timers = vi.spyOn(globalThis, "setTimeout"); const resolveWorker = runtimeWorker.resolveRuntimeWorkerUrl; - const resolver = - mode !== "healthy" - ? vi - .spyOn(runtimeWorker, "resolveRuntimeWorkerUrl") - .mockImplementation((params) => - params.sourceWorkerName === runtimeProcessEntrypoints.sharedStateStore.sourceWorkerName - ? new URL( - "./sqlite-worker-shared-state-idle-fixture.test-support.ts", - import.meta.url, - ) - : resolveWorker(params), - ) - : undefined; + const resolver = vi + .spyOn(runtimeWorker, "resolveRuntimeWorkerUrl") + .mockImplementation((params) => + params.sourceWorkerName === runtimeProcessEntrypoints.sharedStateStore.sourceWorkerName + ? new URL("./sqlite-worker-shared-state-idle-fixture.test-support.ts", import.meta.url) + : resolveWorker(params), + ); const messages = vi.spyOn(Worker.prototype, "postMessage"); const read = () => executeOpenClawStateWorker(context, { @@ -61,7 +53,7 @@ async function fixture(mode: "healthy" | "local-reader" | "unsettled-inspection" expect(await read()).toEqual([]); const worker = messages.mock.contexts[0]; messages.mockRestore(); - resolver?.mockRestore(); + resolver.mockRestore(); if (!(worker instanceof Worker)) { throw new Error("Expected the canonical shared-state worker"); } @@ -99,29 +91,9 @@ it("retains the original healthy worker after one minute and closes it after 30 // Joining a real call also joins the original owner's retirement, if it retired at one minute. expect(await f.read()).toEqual([]); expect(f.worker.threadId).not.toBe(-1); - const postMessage = f.worker.postMessage.bind(f.worker); - let inspected = false; - vi.spyOn(f.worker, "postMessage").mockImplementation( - (request: SqliteWorkerRequest, transfers?: readonly Transferable[]) => { - if ( - request.type === "execute" && - deserialize(request.input).type === "database.inspectIdle" - ) { - const held = withStateDatabaseCoordinatorRuntimeDirectory( - f.context.coordinatorRuntime, - () => retainHeldStateDatabaseCoordinator(f.context.admission.databasePath), - ); - expect(held).toBeDefined(); - held?.release(); - inspected = true; - } - return postMessage(request, transfers); - }, - ); f.advance(minute); f.scheduled(minute)(); await vi.waitFor(() => expect(f.scheduled(29 * minute)).toBeTypeOf("function")); - expect(inspected).toBe(true); expect(f.worker.threadId).not.toBe(-1); const exited = once(f.worker, "exit"); f.advance(29 * minute); @@ -130,6 +102,21 @@ it("retains the original healthy worker after one minute and closes it after 30 expect(await f.read()).toEqual([]); }); +it("retires an unavailable actor even when its completed idle result is healthy", async () => { + const f = await fixture(); + const available = vi + .spyOn(sqliteWorkers, "isSqliteWorkerStoreAvailable") + .mockReturnValueOnce(false); + try { + f.advance(minute); + f.scheduled(minute)(); + await vi.waitFor(() => expect(f.worker.threadId).toBe(-1)); + expect(await f.read()).toEqual([]); + } finally { + available.mockRestore(); + } +}); + it("retires a worker with an untracked local reader and releases its actual WAL pin", async () => { const f = await fixture("local-reader"); const { DatabaseSync } = requireNodeSqlite(); @@ -223,49 +210,70 @@ it("ignores an inspection result and old expiry when real work resumes", async ( expect(f.worker.threadId).not.toBe(-1); }); -it("replaces a failed idle actor even when overlapping activity sends no command", async () => { +it("replaces a failed idle actor after an enclosing callback settles", async () => { const f = await fixture("unsettled-inspection"); - const postMessage = f.worker.postMessage.bind(f.worker); - const dispatched = createDeferredCore(); + const nativePost = vi.spyOn(MessagePort.prototype, "postMessage"); + nativePost.mockRestore(); let resume: (() => void) | undefined; const send = vi - .spyOn(f.worker, "postMessage") - .mockImplementation((request: SqliteWorkerRequest, transfers?: readonly Transferable[]) => { - if ( - request.type === "execute" && - deserialize(request.input).type === "database.inspectIdle" - ) { - resume = () => postMessage(request, transfers); - dispatched.resolve(); + .spyOn(MessagePort.prototype, "postMessage") + .mockImplementation(function (this: MessagePort, message, transfers) { + if (isRecord(message) && message.type === "accepted" && Object.hasOwn(message, "admission")) { + // Hold the acquired-custody grant, after live authority has accepted this inspection. + send.mockRestore(); + resume = () => nativePost.call(this, message, transfers); return; } - return postMessage(request, transfers); + return nativePost.call(this, message, transfers); }); f.advance(minute); f.scheduled(minute)(); - await dispatched.promise; const entered = createDeferredCore(); const finish = createDeferredCore(); - const active = runOpenClawStateWorkerOperation(f.context, async () => { - entered.resolve(); - await finish.promise; - return "completed without dispatch"; - }); - await entered.promise; - send.mockRestore(); + const escape = createDeferredCore(); + void escape.promise.catch(() => {}); + let active: Promise | undefined; try { + await vi.waitFor(() => expect(resume).toBeTypeOf("function")); if (!resume) { - throw new Error("Expected a held native inspection request"); + throw new Error("Expected an admitted native inspection"); } + active = runOpenClawStateWorkerOperation(f.context, () => + Promise.race([ + (async () => { + entered.resolve(); + await finish.promise; + await expect(f.read()).rejects.toMatchObject({ code: "unavailable" }); + return "completed without dispatch"; + })(), + escape.promise, + ]), + ); + await entered.promise; const exited = once(f.worker, "exit"); resume(); + resume = undefined; await exited; await nextTurn(); - } finally { finish.resolve(); + let settled = false; + void active.then( + () => { + settled = true; + }, + () => { + settled = true; + }, + ); + await vi.waitFor(() => expect(settled).toBe(true)); + expect(await active).toBe("completed without dispatch"); + await nextTurn(); + expect(await f.read()).toEqual([]); + } finally { + send.mockRestore(); + resume?.(); + finish.resolve(); + escape.reject(new Error("Release the enclosing fixture after failed observation")); + await active?.catch(() => {}); } - expect(await active).toBe("completed without dispatch"); - await nextTurn(); - // A stale unavailable store would make this first post-failure call fail once. - expect(await f.read()).toEqual([]); }); diff --git a/src/state/openclaw-state-worker-error.test.ts b/src/state/openclaw-state-worker-error.test.ts index 0e13cae0ec2d..d28d37d8d1d5 100644 --- a/src/state/openclaw-state-worker-error.test.ts +++ b/src/state/openclaw-state-worker-error.test.ts @@ -1,7 +1,8 @@ import { describe, expect, it, vi } from "vitest"; import { SqliteCoordinatorError } from "../infra/sqlite-coordinator.js"; import { SqliteSchemaVersionError } from "../infra/sqlite-user-version.js"; -import { decodeSqliteWorkerReplyError } from "../infra/sqlite-worker-broker-reply.js"; +import { receiveSqliteWorkerReply } from "../infra/sqlite-worker-broker-reply.js"; +import type { Job } from "../infra/sqlite-worker-broker.types.js"; import { findStartupMaintenanceRequiredError, StartupMaintenanceRequiredError, @@ -122,30 +123,55 @@ describe("shared-state worker error transport", () => { if (!payload) { throw new Error("Expected canonical payload"); } - const failure = decodeSqliteWorkerReplyError( + const job: Job = { + request: { + type: "execute", + id: 1, + actor: 1, + input: new Uint8Array(), + stateContext: { + environment: { OPENCLAW_STATE_DIR: "/fixture" }, + coordinatorRuntime: { directory: "/fixture/coordinator", keepAlive: false }, + }, + }, + bytes: 0, + resolve: () => undefined, + reject: () => undefined, + detach: () => undefined, + }; + let failure: unknown; + receiveSqliteWorkerReply( { - request: { - type: "execute", - id: 1, - actor: 1, - input: new Uint8Array(), - stateContext: { - environment: { OPENCLAW_STATE_DIR: "/fixture" }, - coordinatorRuntime: { directory: "/fixture/coordinator", keepAlive: false }, + current: job, + worker: { + postMessage: () => { + throw new Error("Unexpected native dispatch"); }, }, - bytes: 0, - resolve: () => undefined, - reject: () => undefined, - detach: () => undefined, }, { - name: "SqliteWorkerError", - message: "write outcome unknown", - code: "outcome-unknown", - sharedState: payload, + id: 1, + ok: false, + error: { + name: "SqliteWorkerError", + message: "write outcome unknown", + code: "outcome-unknown", + sharedState: payload, + }, + }, + { + fail(error) { + throw error; + }, + finish(_job, error) { + failure = error; + }, + dispatch() {}, }, ); + if (!(failure instanceof Error)) { + throw new Error("Expected the broker to settle the original failure"); + } expect(hydrateOpenClawStateWorkerError(failure)).toBe(failure); expect(failure).toMatchObject({ code: "outcome-unknown" }); expect(findStartupMaintenanceRequiredError(failure)).toBeUndefined(); diff --git a/src/state/openclaw-state-worker-store.ts b/src/state/openclaw-state-worker-store.ts index 27afb1478f86..19943e5c8d10 100644 --- a/src/state/openclaw-state-worker-store.ts +++ b/src/state/openclaw-state-worker-store.ts @@ -38,7 +38,7 @@ type StoreOperations = OpenClawStateWorkerOperations & OpenClawStateWorkerInspec type Store = SqliteWorkerStore; type DomainScope = Pick, "execute">; type OperationOptions = { - /** Acquire matching host lifecycle custody for each dispatched command. */ + /** Acquire matching lifecycle custody for each dispatched command. */ requireStateLifecycle?: boolean; assertCurrent?: (commandType?: PropertyKey) => void; createAdmission?: SqliteWorkerAdmissionFactory; @@ -131,13 +131,15 @@ function createSharedStateWorkerOwner() { let healthy = false; if (inspect) { try { + // This idle generation owns retirement while foreground callbacks may remain active. healthy = (await runWithCapturedWorkerContext(entry.context, () => - runWithOpenClawStateWorkerStore( + runSqliteWorkerStoreOperation( store, - entry.context, (scope) => scope.execute({ type: "database.inspectIdle", input: undefined }), + entry.context, () => { + entry.context.admission.assertCurrent(); if (!isCurrentIdle()) { throw new Error("Shared-state worker resumed before idle inspection"); } @@ -145,7 +147,7 @@ function createSharedStateWorkerOwner() { undefined, true, ), - )) === "healthy"; + )) === "healthy" && isSqliteWorkerStoreAvailable(store); entry.context.admission.assertCurrent(); } catch { // Unavailable inspection cannot justify retaining a potentially pinned native reader. @@ -192,6 +194,19 @@ function createSharedStateWorkerOwner() { } released = true; entry.activeOperations -= 1; + if ( + entry.activeOperations === 0 && + stores.has(entry) && + !isSqliteWorkerStoreAvailable(store) + ) { + void retire(entry).catch((error: unknown) => { + log.warn("Shared-state worker retirement failed", { + path: entry.context.admission.databasePath, + error, + }); + }); + return; + } scheduleIdleRetirement(entry); }; }; @@ -227,7 +242,12 @@ function createSharedStateWorkerOwner() { return { close, async retireFailedStore(store: Store): Promise { - await Promise.all([...stores.values()].filter((entry) => entry.store === store).map(retire)); + // An enclosing callback may await this operation; its final release owns retirement. + await Promise.all( + [...stores.values()] + .filter((entry) => entry.store === store && entry.activeOperations <= 1) + .map(retire), + ); }, retainOperation, async open( @@ -425,7 +445,7 @@ async function runAdmittedOpenClawStateWorkerOperation( operation, options?.assertCurrent, options?.createAdmission, - // Host-held lifecycle custody keeps native writers from blocking the worker's grant handler. + // Commands with live admission retain lifecycle custody through native settlement. options?.requireStateLifecycle === true || options?.createAdmission !== undefined, ); } finally { @@ -494,8 +514,9 @@ async function runWithOpenClawStateWorkerStore( requireStateLifecycle = false, ): Promise { const { admission } = context; + let result: T; try { - return await runSqliteWorkerStoreOperation( + result = await runSqliteWorkerStoreOperation( store, operation, context, @@ -520,6 +541,18 @@ async function runWithOpenClawStateWorkerStore( } throw error; } + if (!isSqliteWorkerStoreAvailable(store)) { + try { + await owner().retireFailedStore(store); + } catch (error) { + // Keep a completed outcome while canonical cleanup retains the failed entry. + log.warn("Completed shared-state operation retirement failed", { + path: admission.databasePath, + error, + }); + } + } + return result; } /** Retain the actual lease until every admitted worker transaction has settled. */ diff --git a/src/tasks/task-registry-live-flow.worker.test.ts b/src/tasks/task-registry-live-flow.worker.test.ts index 7f10e07f8a4f..3bc7d942ec26 100644 --- a/src/tasks/task-registry-live-flow.worker.test.ts +++ b/src/tasks/task-registry-live-flow.worker.test.ts @@ -1,5 +1,6 @@ import { afterEach, beforeEach, expect, it, vi } from "vitest"; import { createDeferred } from "../../test/helpers/promise.js"; +import { observeDeviceAuthHostSql } from "../infra/device-auth-store.sql.test-support.js"; import { requireNodeSqlite } from "../infra/node-sqlite.js"; import { getActiveGatewayRootWorkCount } from "../process/gateway-work-admission.js"; import { closeOpenClawStateDatabaseAsync } from "../state/openclaw-state-db.js"; @@ -119,6 +120,7 @@ it("retries the live equal-time winner through the worker and canonical close", for (const counter of counters) { counter.mockClear(); } + const hostSql = observeDeviceAuthHostSql(state.statePath("state", "openclaw.sqlite")); const startedAt = performance.now(); await vi.advanceTimersByTimeAsync(1_000); expect(await published.promise).toMatchObject({ @@ -127,7 +129,12 @@ it("retries the live equal-time winner through the worker and canonical close", status: "succeeded", }); await vi.waitFor(() => expect(getActiveGatewayRootWorkCount()).toBe(0)); + const beforeCloseSql = hostSql.counts(); + expect(Object.values(beforeCloseSql).flatMap((counts) => Object.values(counts))).toEqual( + Array(28).fill(0), + ); await closeOpenClawStateDatabaseAsync(); + console.info("Live flow host SQL", { beforeClose: beforeCloseSql, afterClose: hostSql.counts() }); rememberPreparedSql(); let unattributedStatements = 0; const sql = [ @@ -153,6 +160,7 @@ it("retries the live equal-time winner through the worker and canonical close", }); expect(unattributedStatements).toBe(0); expect(domainSql).toEqual([]); + hostSql.restore(); vi.restoreAllMocks(); expect(loadTaskFlowRegistryStateFromSqliteReadOnly().flows.get(flow.flowId)).toMatchObject({ goal: "Live insertion winner",