From 54885d54a0385b6137847a1b29ba14c4cf5a8474 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 1 Oct 2026 08:16:07 -0700 Subject: [PATCH] fix(state): keep agent database admission open across transient lease contention (#162750) A stalled cron sub-agent tool call drove sustained SQLite write contention; a state.lease operation with busy_timeout 0 failed with lock contention, and the per-agent database execution admission then closed permanently until restart because the execution owner retired logical admission after recoverable native-open failures and the worker transport dropped SQLite contention codes. Admission is now preserved across recoverable native-open failures, contention retries share the 25 ms cadence introduced for the lease heartbeat in #160702 and run only before application work starts, and shutdown, FIFO ownership and unknown-outcome safeguards are unchanged. Doctor preserves nested lease-loss diagnostics instead of flattening them; a nine-transaction Doctor probe confirms #160702 closes the heartbeat loss reported in #162083. Closes #162579 Refs #162083 --- .../integrity-and-recovery.md | 8 ++ .../state-migrations.plugin-doctor.test.ts | 14 +- src/infra/state-migrations.plugin-doctor.ts | 6 +- ...enclaw-agent-execution.close-wedge.test.ts | 61 ++++++++- src/state/openclaw-agent-execution.ts | 46 ++++++- src/state/openclaw-state-lease-acquisition.ts | 21 +++ .../openclaw-state-lease-contention.test.ts | 48 ++++--- .../openclaw-state-lease-heartbeat-shared.ts | 2 + .../openclaw-state-lease-heartbeat.worker.ts | 16 ++- src/state/openclaw-state-lease-storage.ts | 9 +- src/state/openclaw-state-worker-error.test.ts | 23 +++- src/state/openclaw-state-worker-error.ts | 3 + test/doctor-lease-contention.e2e.test.ts | 125 ++++++++++++++++++ 13 files changed, 348 insertions(+), 34 deletions(-) create mode 100644 test/doctor-lease-contention.e2e.test.ts diff --git a/docs/reference/database-schemas/integrity-and-recovery.md b/docs/reference/database-schemas/integrity-and-recovery.md index e7a7c416a4b7..9d322031e201 100644 --- a/docs/reference/database-schemas/integrity-and-recovery.md +++ b/docs/reference/database-schemas/integrity-and-recovery.md @@ -321,6 +321,14 @@ IPC error remains the reported failure even if termination also fails. Agent database maintenance fences other writers with a 60-second lease in the shared state database. A dedicated worker renews that lease during synchronous integrity scans and migration phases. Maintenance still checks the exact persisted owner before mutations and commit, and stops if the heartbeat fails or ownership expires or changes. Finishing or cancelling maintenance stops renewal before releasing the lease; process death leaves at most the remaining lease duration. +SQLite lock contention retries at 25 ms intervals within the lease acquisition +budget or the heartbeat's durable expiry. Agent execution admission also retries +contention before entering application work, for up to two seconds after its first +failed preparation settles. An exhausted retry returns the contention error and +leaves admission available for the next request. Shutdown still revokes admission; +completed or entered application work is never replayed. Doctor's plugin session +repair warning includes the nested lease-loss cause when maintenance cannot settle. + Before draining heartbeats for a file capture, each state-lease owner attempts a final ordinary renewal. Capture remains bounded by the shortest durable expiry read after drainage; it cannot renew while files are excluded or revive an expired owner. Asynchronous agent-database admission runs the first full-file integrity check in a read-only child process when that check is outside a write transaction. Later ordinary opens reuse remembered verification. Maintenance retains its independent full check. The connection and owning scope remain held until the native reader closes; cancellation and timeout wait for process exit. Schema changes, index repairs, and compaction retain their synchronous phases. diff --git a/src/infra/state-migrations.plugin-doctor.test.ts b/src/infra/state-migrations.plugin-doctor.test.ts index 6010de30acee..7ac02fd783c6 100644 --- a/src/infra/state-migrations.plugin-doctor.test.ts +++ b/src/infra/state-migrations.plugin-doctor.test.ts @@ -5,6 +5,7 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import type { listPluginDoctorStateMigrationEntries } from "../plugins/doctor-contract-registry.js"; import { closeOpenClawAgentDatabasesForTest } from "../state/openclaw-agent-db.js"; import { closeOpenClawStateDatabaseForTest } from "../state/openclaw-state-db.js"; +import { createOpenClawStateLeaseLostError } from "../state/openclaw-state-lease-error.js"; import { createTrackedTempDirs } from "../test-utils/tracked-temp-dirs.js"; import { autoMigrateLegacyPluginDoctorState, @@ -32,7 +33,16 @@ vi.mock("../plugins/plugin-lifecycle-lease.js", async (importOriginal) => { // The lease owner validates again after the callback returns. Model a lost // lease at that boundary, after migrations have already committed. if (controls.failSettlement) { - throw new Error("lease settlement failed"); + throw createOpenClawStateLeaseLostError( + { + scope: "core:agent-database-maintenance", + key: "global", + leaseLabel: "agent database maintenance lease", + }, + new Error( + "state lease heartbeat exited: lease expired or ownership lost (exitCode=0, acquiredAt=1800000000000, lastRenewedAt=1800000020000)", + ), + ); } return result; })) satisfies typeof actual.withPluginLifecycleLease, @@ -264,7 +274,7 @@ describe("plugin Doctor migrations", () => { ? "refusal" : failure === "detector" ? "second detector failed" - : "lease settlement failed", + : "lease expired or ownership lost (exitCode=0, acquiredAt=1800000000000, lastRenewedAt=1800000020000)", ); } diff --git a/src/infra/state-migrations.plugin-doctor.ts b/src/infra/state-migrations.plugin-doctor.ts index 5e3a3722f8e0..b035fa7c60b4 100644 --- a/src/infra/state-migrations.plugin-doctor.ts +++ b/src/infra/state-migrations.plugin-doctor.ts @@ -14,6 +14,7 @@ import { preparePluginDoctorMigrationResources } from "../plugins/doctor-migrati import { withPluginLifecycleLease } from "../plugins/plugin-lifecycle-lease.js"; import { withAgentDatabaseMaintenanceLease } from "../state/openclaw-agent-db.js"; import { prepareOpenClawStateDatabaseSchema } from "../state/openclaw-state-db.js"; +import { formatErrorMessage } from "./errors.js"; import { acquireGatewayLock } from "./gateway-lock.js"; import { formatStartupMigrationFailure } from "./state-migrations.messages.js"; import { createPluginDoctorStateMigrationContext } from "./state-migrations.plugin-doctor-context.js"; @@ -562,7 +563,10 @@ export async function runPostSessionPluginDoctorStateRepairs(params: { return { ...completed, completedPluginIds: undefined, - warnings: [...completed.warnings, `Plugin session repair did not settle: ${String(error)}.`], + warnings: [ + ...completed.warnings, + `Plugin session repair did not settle: ${formatErrorMessage(error)}.`, + ], warningDisposition: undefined, }; } diff --git a/src/state/openclaw-agent-execution.close-wedge.test.ts b/src/state/openclaw-agent-execution.close-wedge.test.ts index 3409646bc7ef..f9ffb855eea8 100644 --- a/src/state/openclaw-agent-execution.close-wedge.test.ts +++ b/src/state/openclaw-agent-execution.close-wedge.test.ts @@ -2,6 +2,7 @@ import fs from "node:fs"; import type { Worker } from "node:worker_threads"; import { afterEach, expect, it, vi } from "vitest"; import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js"; +import * as backoff from "../infra/backoff.js"; import { formatErrorMessageWithCode } from "../infra/errors.js"; import { createSqliteWorkerOperationAdmission } from "../infra/sqlite-worker-operation-admission.js"; import { closeOpenClawAgentDatabasesAsync } from "./openclaw-agent-db.js"; @@ -13,7 +14,7 @@ import { closeOpenClawStateDatabaseAsync } from "./openclaw-state-db.js"; // retain cleanup custody without permanently refusing later requests (#159438). const fault = vi.hoisted(() => ({ marker: "close-wedge-kill-marker", - enabled: new SharedArrayBuffer(2 * Int32Array.BYTES_PER_ELEMENT), + enabled: new SharedArrayBuffer(3 * Int32Array.BYTES_PER_ELEMENT), workers: new Set(), })); @@ -27,6 +28,17 @@ vi.mock("../infra/worker-cpu.js", async (importOriginal) => { const prepare = DatabaseSync.prototype.prepare; DatabaseSync.prototype.prepare = function (sql) { const statement = prepare.call(this, sql); + if (/insert into "?agent_database_leases"?/i.test(sql)) { + const run = statement.run.bind(statement); + statement.run = (...args) => { + const fault = new Int32Array(workerData.closeWedgeEnabled); + if (Atomics.compareExchange(fault, 0, 3, 0) === 3 || Atomics.load(fault, 0) === 4) { + Atomics.add(fault, 2, 1); + throw Object.assign(new Error("database is locked"), { code: "ERR_SQLITE_ERROR", errcode: 5 }); + } + return run(...args); + }; + } if (/delete from "?agent_database_leases"?/i.test(sql)) { const run = statement.run.bind(statement); statement.run = (...args) => { @@ -90,12 +102,59 @@ vi.mock("../infra/worker-cpu.js", async (importOriginal) => { const tempDirs = useAutoCleanupTempDirTracker((cleanup) => afterEach(async () => { Atomics.store(new Int32Array(fault.enabled), 0, 0); + vi.restoreAllMocks(); await closeOpenClawAgentDatabasesAsync(); await closeOpenClawStateDatabaseAsync(); cleanup(); }), ); +it.each(["retry", "deadline", "shutdown"] as const)( + "keeps contention recoverable while preserving the %s boundary", + async (ending) => { + const env = { OPENCLAW_STATE_DIR: fs.realpathSync(tempDirs.make("agent-open-contention-")) }; + const execution = captureOpenClawAgentDatabaseExecution({ agentId: "first", env }); + const attempts = Atomics.load(new Int32Array(fault.enabled), 2); + Atomics.store(new Int32Array(fault.enabled), 0, ending === "deadline" ? 4 : 3); + const clock = vi.spyOn(performance, "now"); + vi.spyOn(backoff, "sleepWithAbort").mockImplementation(async () => { + if (ending === "shutdown") { + await closeOpenClawAgentDatabasesAsync(); + } else if (ending === "deadline") { + clock.mockReturnValue(performance.now() + 2_000); + } + }); + try { + const preparing = execution.prepare(source); + if (ending === "retry") { + await preparing; + } else { + await expect(preparing).rejects.toThrow(ending === "shutdown" ? /closed/ : /locked/); + } + expect(Atomics.load(new Int32Array(fault.enabled), 2)).toBeGreaterThan(attempts); + Atomics.store(new Int32Array(fault.enabled), 0, 0); + if (ending !== "shutdown") { + execution.assertCurrent(); + await execution.prepare(source); + const operation = vi.fn(async () => "recovered"); + expect(await execution.runExisting(source, operation)).toBe("recovered"); + expect(operation).toHaveBeenCalledOnce(); + const failure = Object.assign(new Error("database is locked"), { errcode: 5 }); + const rejected = vi.fn(async () => { + throw failure; + }); + await expect(execution.runExisting(source, rejected)).rejects.toBe(failure); + expect(rejected).toHaveBeenCalledOnce(); + await closeOpenClawAgentDatabasesAsync(); + } + expect(() => execution.assertCurrent()).toThrow(/closed/); + } finally { + Atomics.store(new Int32Array(fault.enabled), 0, 0); + await execution.release(); + } + }, +); + const source: AgentDatabaseRequestExecutionSource = { assertCurrent: () => undefined, createAdmission(binding) { diff --git a/src/state/openclaw-agent-execution.ts b/src/state/openclaw-agent-execution.ts index 2735b35b2193..6cbd7513655f 100644 --- a/src/state/openclaw-agent-execution.ts +++ b/src/state/openclaw-agent-execution.ts @@ -2,7 +2,9 @@ import { AsyncLocalStorage } from "node:async_hooks"; import { addAbortListener } from "node:events"; import path from "node:path"; import { cloneEnvWithPlatformSemantics } from "../config/config-env-vars.js"; +import { sleepWithAbort } from "../infra/backoff.js"; import { formatErrorMessage } from "../infra/errors.js"; +import { isSqliteLockError } from "../infra/sqlite-error-diagnostics.js"; import { SQLITE_IDLE_HANDLE_TTL_MS } from "../infra/sqlite-handle-lifecycle.js"; import { retainSqliteWorkerErrorCode } from "../infra/sqlite-worker-contract.js"; import { @@ -41,6 +43,10 @@ import { } from "./openclaw-state-db-async-lifecycle.js"; import { registerOpenClawStateDatabaseAsyncResource } from "./openclaw-state-db-cache.js"; import { resolveOpenClawStateSqlitePath } from "./openclaw-state-db.paths.js"; +import { + LEASE_CONTENTION_RETRY_MS, + LEASE_CONTENTION_RETRY_TIMEOUT_MS, +} from "./openclaw-state-lease-heartbeat-shared.js"; import { captureOpenClawStateReadContext, captureOpenClawStateWorkerContext, @@ -292,6 +298,7 @@ function createAgentDatabaseExecution( createIfMissing = false, creatingTarget?: DatabasePathIdentity, signal?: AbortSignal, + contentionDeadline?: number, ): Promise { const pending = agentDatabaseLifecycle.pending.get(pathname); if (pending) { @@ -352,10 +359,14 @@ function createAgentDatabaseExecution( } } const current = generation; + let entered = false; try { const result = await current.run( source, - operation, + (scope) => { + entered = true; + return operation(scope); + }, assertCallerCurrent, createIfMissing, signal, @@ -370,9 +381,10 @@ function createAgentDatabaseExecution( return result; } catch (error) { const nativeFailed = current.failed(); + const contended = !entered && isSqliteLockError(error); if (generation === current && (nativeFailed || retireNativeOnFailure)) { try { - if (nativeFailed) { + if (nativeFailed && !contended) { await owner.close(); } else { // The rejected broker scope has settled; only its captured native owner is retired. @@ -387,6 +399,36 @@ function createAgentDatabaseExecution( ); } } + if (contended) { + const deadline = + contentionDeadline ?? performance.now() + LEASE_CONTENTION_RETRY_TIMEOUT_MS; + if (contentionDeadline === undefined) { + log.warn( + "Agent database execution admission delayed by SQLite lock contention; retrying before execution.", + ); + } + const remaining = deadline - performance.now(); + if (remaining > 0) { + await sleepWithAbort(Math.min(LEASE_CONTENTION_RETRY_MS, remaining), signal); + assertCurrent(); + source.assertCurrent(); + assertCallerCurrent(); + if (performance.now() >= deadline) { + throw error; + } + return run( + source, + operation, + assertCallerCurrent, + expectedIdentity, + retireNativeOnFailure, + createIfMissing, + creatingTarget, + signal, + deadline, + ); + } + } throw error; } } diff --git a/src/state/openclaw-state-lease-acquisition.ts b/src/state/openclaw-state-lease-acquisition.ts index a7657872ad23..27ba1438afa2 100644 --- a/src/state/openclaw-state-lease-acquisition.ts +++ b/src/state/openclaw-state-lease-acquisition.ts @@ -10,6 +10,7 @@ import { OpenClawStateLeaseAcquisitionError, OpenClawStateLeaseError, } from "./openclaw-state-lease-error.js"; +import { LEASE_CONTENTION_RETRY_MS } from "./openclaw-state-lease-heartbeat-shared.js"; import { STATE_LEASE_WRITE_BACKOFF } from "./openclaw-state-lease-storage.js"; import type { OpenClawStateLeaseAcquisition } from "./openclaw-state-lease.types.js"; @@ -30,6 +31,7 @@ export async function acquireOpenClawStateLease(params: { let preparation = params.prepare; let attempt = 0; let lastReportedHolder: string | undefined; + let contentionReported = false; const cancellation = params.signal ? new AbortController() : undefined; let aborted: OpenClawStateLeaseAcquisitionError | undefined; const abort = () => { @@ -56,6 +58,7 @@ export async function acquireOpenClawStateLease(params: { while (true) { assertCurrent(); let outcome: OpenClawStateLeaseAcquisition; + let acquiring = false; try { if (preparation) { const prepare = preparation; @@ -63,6 +66,7 @@ export async function acquireOpenClawStateLease(params: { prepare(); deadline = performance.now() + params.waitMs; } + acquiring = true; outcome = await params.acquire(assertCurrent, cancellation?.signal); } catch (error) { if ( @@ -82,6 +86,23 @@ export async function acquireOpenClawStateLease(params: { const failure = error instanceof OpenClawStateLeaseError ? error.cause : error; if (isSqliteLockError(failure)) { assertCurrent(); + const remainingMs = deadline - performance.now(); + if (acquiring && remainingMs > 0) { + if (!contentionReported) { + log.warn(`Waiting for ${params.label} after SQLite lock contention.`); + contentionReported = true; + } + try { + await sleepWithAbort(Math.min(remainingMs, LEASE_CONTENTION_RETRY_MS), params.signal); + } catch (sleepError) { + assertCurrent(); + throw sleepError; + } + assertCurrent(); + if (performance.now() < deadline) { + continue; + } + } } throw new OpenClawStateLeaseAcquisitionError( params.label, diff --git a/src/state/openclaw-state-lease-contention.test.ts b/src/state/openclaw-state-lease-contention.test.ts index ea1a60e1de94..dbd7bd9e740b 100644 --- a/src/state/openclaw-state-lease-contention.test.ts +++ b/src/state/openclaw-state-lease-contention.test.ts @@ -1,6 +1,7 @@ import { DatabaseSync } from "node:sqlite"; import { setImmediate as yieldImmediate } from "node:timers/promises"; import { describe, expect, it, vi } from "vitest"; +import * as backoff from "../infra/backoff.js"; import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js"; import { openOpenClawStateDatabase } from "./openclaw-state-db.js"; import { releaseOpenClawStateLeaseBestEffort } from "./openclaw-state-lease-storage.js"; @@ -9,7 +10,7 @@ import { withOpenClawStateLease } from "./openclaw-state-lease.js"; describe.each([undefined, "existing"] as const)( "lease SQLite contention (%s schema)", (schemaPolicy) => { - it.each(["release", "timeout", "abort", "cleanup"] as const)( + it.each(["release", "timeout", "deadline", "abort", "cleanup"] as const)( "preserves the %s contract while an independent writer holds a native write transaction", async (ending) => { await withOpenClawTestState({ label: "lease-native-contention" }, async (state) => { @@ -28,8 +29,20 @@ describe.each([undefined, "existing"] as const)( }; let writer = ending === "cleanup" ? undefined : takeWriter(); const controller = new AbortController(); - let entered = false; + let entered = 0; let cleanupRelease: Promise | undefined; + const clock = vi.spyOn(performance, "now").mockReturnValue(0); + const retry = vi.spyOn(backoff, "sleepWithAbort").mockImplementation(async () => { + if (ending === "release") { + writer?.release(); + } else if (ending === "deadline") { + clock.mockReturnValue(5_000); + } else if (ending === "abort") { + controller.abort(new Error("cancel waiting acquisition")); + } else { + throw new Error(`Unexpected acquisition retry for ${ending}`); + } + }); try { const operation = withOpenClawStateLease( { @@ -41,7 +54,7 @@ describe.each([undefined, "existing"] as const)( signal: controller.signal, }, async (lease) => { - entered = true; + entered += 1; lease.assertOwned(); if (ending === "cleanup") { vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] }); @@ -50,37 +63,38 @@ describe.each([undefined, "existing"] as const)( } }, ); - if (ending === "timeout" || ending === "abort") { + if (ending === "timeout" || ending === "deadline" || ending === "abort") { const rejected = expect(operation).rejects.toMatchObject({ code: - ending === "timeout" - ? "OPENCLAW_STATE_LEASE_STORAGE_FAILED" - : "OPENCLAW_STATE_LEASE_ABORTED", + ending === "abort" + ? "OPENCLAW_STATE_LEASE_ABORTED" + : "OPENCLAW_STATE_LEASE_STORAGE_FAILED", outcome: - ending === "timeout" - ? { kind: "store-unavailable", reason: "sqlite-busy" } - : { kind: "aborted", reason: "caller-signal", elapsedMs: expect.any(Number) }, + ending === "abort" + ? { kind: "aborted", reason: "caller-signal", elapsedMs: expect.any(Number) } + : { kind: "store-unavailable", reason: "sqlite-busy" }, }); - if (ending === "abort") { - controller.abort(new Error("cancel waiting acquisition")); - } await rejected; - expect(entered).toBe(false); + expect(entered).toBe(0); } else { if (ending === "release") { - expect(entered).toBe(false); - writer?.release(); + expect(entered).toBe(0); } await operation; await cleanupRelease; - expect(entered).toBe(true); + expect(entered).toBe(1); } + expect(retry).toHaveBeenCalledTimes( + ending === "timeout" || ending === "cleanup" ? 0 : 1, + ); expect( database.db .prepare("SELECT owner FROM state_leases WHERE scope = ? AND lease_key = ?") .all("core:test", "contending-writer"), ).toEqual([]); } finally { + retry.mockRestore(); + clock.mockRestore(); vi.useRealTimers(); writer?.release(); await cleanupRelease; diff --git a/src/state/openclaw-state-lease-heartbeat-shared.ts b/src/state/openclaw-state-lease-heartbeat-shared.ts index f9f3cd9ee3c9..4b817f3df88d 100644 --- a/src/state/openclaw-state-lease-heartbeat-shared.ts +++ b/src/state/openclaw-state-lease-heartbeat-shared.ts @@ -4,6 +4,8 @@ import type { OpenClawStateWorkerErrorPayload } from "./openclaw-state-worker-er // Allow headroom over observed 38 s cold Gateway boots under load; committed lease expiry still bounds startup. export const LEASE_HEARTBEAT_START_TIMEOUT_MS = 60_000; +export const LEASE_CONTENTION_RETRY_MS = 25; +export const LEASE_CONTENTION_RETRY_TIMEOUT_MS = 2_000; export const leaseHeartbeatState = { status: 0, diff --git a/src/state/openclaw-state-lease-heartbeat.worker.ts b/src/state/openclaw-state-lease-heartbeat.worker.ts index a485610b3258..fd641d1a7577 100644 --- a/src/state/openclaw-state-lease-heartbeat.worker.ts +++ b/src/state/openclaw-state-lease-heartbeat.worker.ts @@ -15,6 +15,7 @@ import { leaseHeartbeatState as state, leaseHeartbeatStartupPhase as startupPhase, LEASE_HEARTBEAT_START_TIMEOUT_MS, + LEASE_CONTENTION_RETRY_MS, type LeaseHeartbeatRenewalFailure, type LeaseHeartbeatReply, type LeaseHeartbeatParentMessage, @@ -52,12 +53,16 @@ function openHeartbeatDatabase() { throw error; } } - Atomics.wait(shared, state.status, state.starting, Math.max(1, Math.min(25, remaining()))); + Atomics.wait( + shared, + state.status, + state.starting, + Math.max(1, Math.min(LEASE_CONTENTION_RETRY_MS, remaining())), + ); } throw new Error("state lease heartbeat startup deadline expired or owner stopped"); } const db = openHeartbeatDatabase(); -const CONTENTION_RETRY_MS = 25; Atomics.store(shared, state.startupPhase, startupPhase["open-complete"]); let processOwner = params.processOwner; let heartbeat: ReturnType | undefined; @@ -159,7 +164,7 @@ const renewInWorker = (explicit: boolean): number | undefined => { Math.max( 1, Math.min( - contentionError === undefined ? params.heartbeatMs : CONTENTION_RETRY_MS, + contentionError === undefined ? params.heartbeatMs : LEASE_CONTENTION_RETRY_MS, expiresAt - Date.now(), ), ), @@ -197,7 +202,10 @@ function activateHeartbeat(): void { activateHeartbeat, Math.max( 1, - Math.min(params.heartbeatMs, Number(Atomics.load(shared, state.expiresAt)) - Date.now()), + Math.min( + LEASE_CONTENTION_RETRY_MS, + Number(Atomics.load(shared, state.expiresAt)) - Date.now(), + ), ), ); return; diff --git a/src/state/openclaw-state-lease-storage.ts b/src/state/openclaw-state-lease-storage.ts index 1b18c2c3a489..ade0db7d64c8 100644 --- a/src/state/openclaw-state-lease-storage.ts +++ b/src/state/openclaw-state-lease-storage.ts @@ -18,6 +18,10 @@ import { createOpenClawStateLeaseLostError, toOpenClawStateLeaseVerificationError, } from "./openclaw-state-lease-error.js"; +import { + LEASE_CONTENTION_RETRY_MS, + LEASE_CONTENTION_RETRY_TIMEOUT_MS, +} from "./openclaw-state-lease-heartbeat-shared.js"; import { readOpenClawStateLeaseExpiry, releaseOpenClawStateLeaseInTransaction, @@ -85,12 +89,11 @@ export function withLeaseWriteTransaction( } export const STATE_LEASE_WRITE_BACKOFF = { - initialMs: 25, + initialMs: LEASE_CONTENTION_RETRY_MS, maxMs: 250, factor: 1.5, jitter: 0.25, } as const; -const RELEASE_RETRY_TIMEOUT_MS = 2_000; export type OpenClawStateLeaseOwnerIdentity = OpenClawStateLeaseIdentity & { leaseLabel: string }; @@ -157,7 +160,7 @@ export async function releaseOpenClawStateLeaseBestEffort( params: Parameters[0], execute?: () => Promise, ): Promise { - const deadline = performance.now() + RELEASE_RETRY_TIMEOUT_MS; + const deadline = performance.now() + LEASE_CONTENTION_RETRY_TIMEOUT_MS; let attempt = 0; while (true) { try { diff --git a/src/state/openclaw-state-worker-error.test.ts b/src/state/openclaw-state-worker-error.test.ts index 499c0d62da0b..a3b4ded4dd7c 100644 --- a/src/state/openclaw-state-worker-error.test.ts +++ b/src/state/openclaw-state-worker-error.test.ts @@ -98,6 +98,18 @@ describe("shared-state worker error transport", () => { }, ); + it.each([ + { code: "ERR_SQLITE_ERROR", errcode: 5 }, + { code: "ERR_SQLITE_ERROR", errcode: 6 }, + { code: "ERR_SQLITE_ERROR", errcode: 517 }, + { code: "ERR_SQLITE_ERROR", errcode: 262 }, + { code: "SQLITE_BUSY" }, + { code: "SQLITE_LOCKED" }, + ])("preserves native SQLite contention without an enclosing domain error (%j)", (fields) => { + const original = Object.assign(new Error("native lock contention"), fields); + expect(roundTrip(original)).toMatchObject(fields); + }); + it.each( [RangeError, SyntaxError, TypeError, SkillUploadRequestError].flatMap((ErrorType) => [false, true].map((aggregate) => ({ ErrorType, name: ErrorType.name, aggregate })), @@ -218,8 +230,11 @@ describe("shared-state worker error transport", () => { expect(findStartupMaintenanceRequiredError(hydrated)).toBeInstanceOf(SqliteSchemaVersionError); }); - it("keeps outcome-unknown explicit instead of hydrating a maintenance payload", () => { - const payload = encodeOpenClawStateWorkerError(new SqliteSchemaVersionError("newer schema")); + it.each([ + new SqliteSchemaVersionError("newer schema"), + Object.assign(new Error("native lock contention"), { code: "ERR_SQLITE_ERROR", errcode: 5 }), + ])("keeps outcome-unknown explicit instead of hydrating %s", (original) => { + const payload = encodeOpenClawStateWorkerError(original); assert(payload); const job: Job = { request: { @@ -336,7 +351,7 @@ describe("shared-state worker error transport", () => { }); it("encodes and hydrates ordinary error graphs only with an explicit opt-in", () => { - const cause = Object.assign(new Error("native failure"), { code: "SQLITE_BUSY" }); + const cause = Object.assign(new Error("native failure"), { code: "SQLITE_IOERR" }); const original = new AggregateError([cause], "load and cleanup", { cause }); cause.cause = original; expect(encodeOpenClawStateWorkerError(original)).toBeUndefined(); @@ -345,7 +360,7 @@ describe("shared-state worker error transport", () => { const retained = remoteError(structuredClone(payload)); expect(hydrateOpenClawStateWorkerError(retained)).toBe(retained); const decoded = hydrateOpenClawStateWorkerError(retained, { includeOrdinary: true }); - expect(decoded.cause).toMatchObject({ message: "native failure", code: "SQLITE_BUSY" }); + expect(decoded.cause).toMatchObject({ message: "native failure", code: "SQLITE_IOERR" }); assert(decoded instanceof AggregateError && decoded.cause instanceof Error); expect(decoded.errors[0]).toBe(decoded.cause); expect(decoded.cause.cause).toBe(decoded); diff --git a/src/state/openclaw-state-worker-error.ts b/src/state/openclaw-state-worker-error.ts index c9c4145b595b..1ee15b8c3a79 100644 --- a/src/state/openclaw-state-worker-error.ts +++ b/src/state/openclaw-state-worker-error.ts @@ -1,5 +1,6 @@ import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { + isSqliteLockError, isSqliteNativeOpenFailure, markSqliteNativeOpenFailure, } from "../infra/sqlite-error-diagnostics.js"; @@ -91,6 +92,7 @@ export function encodeOpenClawStateWorkerError( canonical ||= stateDatabasePath !== undefined || nativeOpen || + isSqliteLockError(current) || current instanceof OpenClawQuarantineReadCleanupError || (identity.type !== "error" && identity.type !== "aggregate"); const code = "code" in current ? current.code : undefined; @@ -234,6 +236,7 @@ function decodeErrorGraph( canonical ||= node.stateDatabasePath !== undefined || node.nativeOpen === true || + isSqliteLockError(node) || (node.type === "aggregate" && node.name === DATABASE_QUARANTINE_READ_CLEANUP_ERROR_NAME) || (node.type !== "error" && node.type !== "aggregate"); for (const edge of [...(node.cause ? [node.cause] : []), ...(node.errors ?? [])]) { diff --git a/test/doctor-lease-contention.e2e.test.ts b/test/doctor-lease-contention.e2e.test.ts new file mode 100644 index 000000000000..971bb02a1968 --- /dev/null +++ b/test/doctor-lease-contention.e2e.test.ts @@ -0,0 +1,125 @@ +import { isMainThread } from "node:worker_threads"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { acquireGatewayLock } from "../src/infra/gateway-lock.js"; +import { runPostSessionPluginDoctorStateRepairs } from "../src/infra/state-migrations.plugin-doctor.js"; +import type { listPluginDoctorStateMigrationEntries } from "../src/plugins/doctor-contract-registry.js"; +import { AGENT_DATABASE_MAINTENANCE_LEASE } from "../src/state/openclaw-agent-db-lease.js"; +import { + openOpenClawStateDatabase, + runOpenClawStateWriteTransaction, +} from "../src/state/openclaw-state-db.js"; +import { withOpenClawTestState } from "../src/test-utils/openclaw-test-state.js"; + +const controls = vi.hoisted(() => ({ + entries: [] as ReturnType, +})); + +vi.mock("../src/plugins/doctor-contract-registry.js", async (importOriginal) => ({ + ...(await importOriginal()), + listPluginDoctorStateMigrationEntries: () => controls.entries, +})); + +afterEach(() => { + controls.entries = []; +}); + +// Opt-in release proof: the reporter's 61.748 s of real SQLite lock occupancy +// must span the production 60 s lease while its independent worker keeps time. +describe.runIf(process.env.OPENCLAW_DOCTOR_LEASE_CONTENTION_PROOF === "1")( + "Doctor plugin session repair under sustained SQLite contention", + () => { + it("settles after nine main-thread immediate transactions outlive the original maintenance lease", async () => { + await withOpenClawTestState({ label: "doctor-lease-contention" }, async (state) => { + const durations = [6_392, 8_648, 6_980, 5_296, 3_889, 5_700, 7_325, 9_312, 8_206]; + const blocker = new Int32Array(new SharedArrayBuffer(4)); + const observations: { heldMs: number; heartbeatAt: number; expiresAt: number }[] = []; + let initialExpiry = 0; + let repaired = false; + controls.entries = [ + { + pluginId: "contention-proof", + channelIds: [], + trustedForDurableStores: false, + migration: { + id: "repair-session-state", + label: "Contended plugin session repair", + phase: "after-session-repair", + detectLegacyState: () => ({ preview: ["repair synthetic plugin session state"] }), + migrateLegacyState: () => { + expect(isMainThread).toBe(true); + const { db } = openOpenClawStateDatabase({ env: state.env }); + const readLease = () => { + const row = db + .prepare( + "SELECT expires_at, heartbeat_at FROM state_leases WHERE scope = ? AND lease_key = ?", + ) + .get( + AGENT_DATABASE_MAINTENANCE_LEASE.scope, + AGENT_DATABASE_MAINTENANCE_LEASE.key, + ); + expect(row).toBeDefined(); + return { + expiresAt: Number(row?.expires_at), + heartbeatAt: Number(row?.heartbeat_at), + }; + }; + initialExpiry = readLease().expiresAt; + for (const [index, duration] of durations.entries()) { + const started = performance.now(); + runOpenClawStateWriteTransaction( + () => Atomics.wait(blocker, 0, 0, duration), + { env: state.env }, + { operationLabel: "agent.database.maintenance.admission" }, + ); + observations.push({ heldMs: performance.now() - started, ...readLease() }); + if (index < durations.length - 1) { + // Keep the parent blocked during a known unlocked window. The + // 25 ms worker retry can renew; a 20 s retry misses these gaps. + Atomics.wait(blocker, 0, 0, 100); + } + } + repaired = true; + return { changes: ["Repaired synthetic plugin session state"], warnings: [] }; + }, + }, + }, + ]; + const lock = await acquireGatewayLock({ + env: state.env, + role: "sqlite-maintenance", + allowInTests: true, + }); + expect(lock).not.toBeNull(); + if (!lock) { + throw new Error("Doctor did not acquire isolated maintenance ownership"); + } + try { + const result = await lock.run(() => + runPostSessionPluginDoctorStateRepairs({ + config: {}, + env: state.env, + maintenanceAuthority: lock, + plannedActions: [{ pluginId: "contention-proof", id: "repair-session-state" }], + }), + ); + console.info( + JSON.stringify({ + step: "plugin-doctor-post-session-state", + initialExpiry, + observations, + result, + }), + ); + expect(result.warnings).toEqual([]); + expect(result.changes).toEqual(["Repaired synthetic plugin session state"]); + expect(repaired).toBe(true); + expect(observations).toHaveLength(9); + expect(Date.now()).toBeGreaterThan(initialExpiry); + expect(observations.at(-1)?.expiresAt).toBeGreaterThan(initialExpiry); + } finally { + await lock.release(); + } + }); + }, 120_000); + }, +);