mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
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
This commit is contained in:
parent
43873ce187
commit
54885d54a0
13 changed files with 348 additions and 34 deletions
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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)",
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<Worker>(),
|
||||
}));
|
||||
|
||||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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<T | undefined> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<void> | 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;
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<typeof setTimeout> | 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;
|
||||
|
|
|
|||
|
|
@ -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<T>(
|
|||
}
|
||||
|
||||
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<typeof releaseOpenClawStateLease>[0],
|
||||
execute?: () => Promise<void>,
|
||||
): Promise<void> {
|
||||
const deadline = performance.now() + RELEASE_RETRY_TIMEOUT_MS;
|
||||
const deadline = performance.now() + LEASE_CONTENTION_RETRY_TIMEOUT_MS;
|
||||
let attempt = 0;
|
||||
while (true) {
|
||||
try {
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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 ?? [])]) {
|
||||
|
|
|
|||
125
test/doctor-lease-contention.e2e.test.ts
Normal file
125
test/doctor-lease-contention.e2e.test.ts
Normal file
|
|
@ -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<typeof listPluginDoctorStateMigrationEntries>,
|
||||
}));
|
||||
|
||||
vi.mock("../src/plugins/doctor-contract-registry.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../src/plugins/doctor-contract-registry.js")>()),
|
||||
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);
|
||||
},
|
||||
);
|
||||
Loading…
Add table
Add a link
Reference in a new issue