mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
fix(state): avoid host SQLite waits during async state operations (#152377)
* fix(state): move shared-state lifecycle custody off-thread Acquire and release fresh required per-command lifecycle leases on the existing executing SQLite worker. Keep actual parent-held delegation, FIFO reservations, bounded transfers, lock budgets, and current transaction/commit authority. Retain the per-job port through preparation refusals and result framing so native host writers can service progress. Preserve complete outcomes through failed coordinator cleanup, joining native exit before settlement and credit release. Retire exact unavailable shared-state entries without making nested callbacks wait for their own close. Never replay a mutation. Token lifecycle coordinator calls drop from ten to zero; actual Gateway task creation, unstarted settlement, and live-flow retry each drop from two to zero. Storage APIs, schemas, retention, and remaining native owners stay unchanged. Validation: 178 tests in 13 files, normal build, full changed-file gates, and independent P0-P2 review. Final compiled Node 24.21.0 and Bun 1.4.2 token flows record zero host calls across all 28 SQLite counters, including initialization, with natural exit and unchanged runtime artifacts. * fix(state): include lifecycle worker in wrapper source bundles Native PR wrappers eagerly import the worker broker from their extracted source inventory. Include its lifecycle preparation dependency so cold provisioning and extracted package validation can load the complete import closure. Replace the hook-relay test's retired coordinator SQL allowance with strict caller-thread SQL and close-zero observation through path-specific cleanup. Preserve every record persistence and ownership assertion. Validation: 36 tests passed across three existing files, with one existing platform skip; changed-file gates and independent P0-P2 review passed. Isolated provisioning fixtures used an adequately sized APFS temp directory with the existing capacity guard unchanged. Application runtime source is unchanged. * fix(state): avoid idle cleanup waiting on active callbacks Keep background inspection retirement with the existing idle-generation owner. Use the existing broker operation with captured schema context, current admission, and required lifecycle custody, without foreground completion cleanup that can wait on an enclosing callback. Retain a healthy idle actor only while the broker also records it as available. Observe lifecycle custody in the executing canonical worker fixture and pause failed inspection after its real admission boundary. Cover nested callback recovery and a healthy result from an unavailable actor without changing the one-minute inspection or thirty-minute healthy retirement policy. Validation: 186 focused tests in 14 files, normal build, changed-file gates, and independent P0-P2 review passed. Fresh compiled Node 24.21.0 and Bun 1.4.2 token flows each record zero across all 28 host SQLite counters, with natural exit.
This commit is contained in:
parent
c09841409e
commit
ed2192d22d
23 changed files with 1209 additions and 442 deletions
|
|
@ -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.
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<object, ObservedSql>();
|
||||
// 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 () => {
|
||||
|
|
|
|||
|
|
@ -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<typeof hostSql.counts> | 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();
|
||||
|
|
|
|||
|
|
@ -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<T>(params: {
|
||||
acquire: () => T;
|
||||
acquire: () => T | Promise<T>;
|
||||
shouldRetry: (error: unknown) => boolean;
|
||||
deadlineMs: number;
|
||||
pollIntervalMs: number;
|
||||
|
|
@ -14,7 +14,7 @@ export async function acquireWithWait<T>(params: {
|
|||
let delayMs = params.pollIntervalMs;
|
||||
for (;;) {
|
||||
try {
|
||||
return params.acquire();
|
||||
return await params.acquire();
|
||||
} catch (error) {
|
||||
if (!params.shouldRetry(error)) {
|
||||
throw error;
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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)) {
|
||||
|
|
|
|||
|
|
@ -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<ReturnType<typeof attachGatewaySchemaFenceDelegate>>
|
||||
>();
|
||||
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<void> {
|
|||
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<void> {
|
|||
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<void> {
|
|||
retire = false;
|
||||
}
|
||||
}
|
||||
const executeCommand = (command: unknown) => {
|
||||
const executeCommand = async (command: unknown) => {
|
||||
let coordinator: ReturnType<typeof acquireStateDatabaseCoordinator> | 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<void> {
|
|||
}
|
||||
};
|
||||
try {
|
||||
// SAFETY: The broker serialized a command from this actor's typed store contract.
|
||||
const typedCommand = command as SqliteWorkerCommand<SqliteWorkerOperations>;
|
||||
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<SqliteWorkerOperations>;
|
||||
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<void> {
|
|||
}
|
||||
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<void> {
|
|||
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<void> {
|
|||
};
|
||||
}
|
||||
} 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<void> {
|
|||
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<void> {
|
|||
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<void> {
|
|||
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.
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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> {
|
||||
): 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<typeof acquireStateDatabaseCoordinator>;
|
||||
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();
|
||||
});
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<SqliteWorkerOperationSettlement>();
|
||||
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<SqliteWorkerReply, { ok: true }>,
|
||||
):
|
||||
|
|
@ -237,7 +300,7 @@ export function decodeSqliteWorkerReplyValue(
|
|||
: { type: "complete", value };
|
||||
}
|
||||
|
||||
export function decodeSqliteWorkerReplyError(
|
||||
function decodeSqliteWorkerReplyError(
|
||||
job: Job,
|
||||
error: Extract<SqliteWorkerReply, { ok: false }>["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<Slot, "current" | "failed"> & { worker: Pick<Slot["worker"], "postMessage"> },
|
||||
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<void>;
|
||||
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) {
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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<void>;
|
||||
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<Actor>;
|
||||
queue: Job[];
|
||||
current?: Job;
|
||||
failed?: Error;
|
||||
retiredAfterCompletion?: true;
|
||||
retiring?: Promise<void>;
|
||||
exit: Promise<void>;
|
||||
exited: boolean;
|
||||
|
|
|
|||
|
|
@ -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" }
|
||||
| {
|
||||
|
|
|
|||
253
src/infra/sqlite-worker-lifecycle-cleanup.test.ts
Normal file
253
src/infra/sqlite-worker-lifecycle-cleanup.test.ts
Normal file
|
|
@ -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<never>();
|
||||
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();
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
174
src/infra/sqlite-worker-lifecycle-preparation.ts
Normal file
174
src/infra/sqlite-worker-lifecycle-preparation.ts
Normal file
|
|
@ -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<typeof createDeferredCore<MessagePort | undefined>> | 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<MessagePort | undefined>();
|
||||
params.port.postMessage({ type }, []);
|
||||
return waiting.promise;
|
||||
};
|
||||
let coordinator: ReturnType<typeof acquireStateDatabaseCoordinator> | 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.
|
||||
}
|
||||
}
|
||||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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 },
|
||||
|
|
|
|||
|
|
@ -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<typeof backend.execute>[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");
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<never>();
|
||||
void escape.promise.catch(() => {});
|
||||
let active: Promise<string> | 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([]);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -38,7 +38,7 @@ type StoreOperations = OpenClawStateWorkerOperations & OpenClawStateWorkerInspec
|
|||
type Store = SqliteWorkerStore<StoreOperations>;
|
||||
type DomainScope = Pick<SqliteWorkerStore<OpenClawStateWorkerOperations>, "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<void> {
|
||||
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<T>(
|
|||
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<T>(
|
|||
requireStateLifecycle = false,
|
||||
): Promise<T> {
|
||||
const { admission } = context;
|
||||
let result: T;
|
||||
try {
|
||||
return await runSqliteWorkerStoreOperation<StoreOperations, T>(
|
||||
result = await runSqliteWorkerStoreOperation<StoreOperations, T>(
|
||||
store,
|
||||
operation,
|
||||
context,
|
||||
|
|
@ -520,6 +541,18 @@ async function runWithOpenClawStateWorkerStore<T>(
|
|||
}
|
||||
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. */
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue