From 5d605627c5003c503f8b054438265fdebd26f412 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Mon, 28 Sep 2026 18:03:21 -0700 Subject: [PATCH] perf(sessions): yield during SQLite entry write contention (#160801) * perf(sessions): yield during SQLite entry write contention Keep the Gateway event loop responsive while session-entry patches wait for another SQLite writer. Reuse begin-only asynchronous admission with zero native busy waiting, retaining FIFO ordering, connection custody, commit revalidation, and non-replay of admitted writes. No schema or configuration changes. * fix(sessions): enforce synchronous callbacks during yielding admission Name yielding admission separately from retired async-transaction APIs and register its synchronous callback with the source guard. Make the Slack Stop-owner fixture await its competing committed session write instead of depending on updater microtask ordering. --- .../database-schemas/worker-access.md | 10 +++ .../slack/src/monitor/events/agent.test.ts | 11 +-- scripts/check-sqlite-transaction-boundary.mts | 1 + ...accessor.sqlite-entry-revalidation.test.ts | 76 +++++++++++++++++++ ...-accessor.sqlite-entry.env-capture.test.ts | 4 +- .../sessions/session-accessor.sqlite-entry.ts | 7 +- src/state/openclaw-agent-db-transaction.ts | 70 +++++++++++++++++ .../check-sqlite-transaction-boundary.test.ts | 10 +++ 8 files changed, 180 insertions(+), 9 deletions(-) create mode 100644 src/state/openclaw-agent-db-transaction.ts diff --git a/docs/reference/database-schemas/worker-access.md b/docs/reference/database-schemas/worker-access.md index 7f2ada181463..69dd49e57c02 100644 --- a/docs/reference/database-schemas/worker-access.md +++ b/docs/reference/database-schemas/worker-access.md @@ -48,6 +48,16 @@ Writes, schema transitions, lease grants, and all generic SQLite broker jobs retain immediate fresh ownership verification, including their transaction and commit grants. Schemas, retained data, and update behavior are unchanged. +Legacy session-entry patches yield while waiting for a competing SQLite writer. +Each native attempt uses a zero busy timeout through commit; only a failed +`BEGIN IMMEDIATE` can retry, within the connection's existing admission budget. +The session writer queue retains FIFO order, the captured connection stays +retained, and each attempt rechecks its owner. The admitted transaction revalidates +the prepared rows and caller authority before mutation. Its callback and committed +publications never replay. Entry reads and transaction bodies still execute on +the calling thread; this bounded cutover removes native lock waits without +changing schemas, durability, or update behavior. + Channel setup awaits a fresh policy read after the agent-selection prompt. Deferred plugin migration rows are read by the shared-state worker, and setup rechecks its config owner after the read before using the selected agent. Each diff --git a/extensions/slack/src/monitor/events/agent.test.ts b/extensions/slack/src/monitor/events/agent.test.ts index d330b52cb64b..a80725f080d2 100644 --- a/extensions/slack/src/monitor/events/agent.test.ts +++ b/extensions/slack/src/monitor/events/agent.test.ts @@ -17,11 +17,11 @@ import { patchSessionEntry as patchStoredSessionEntry, upsertSessionEntry, } from "openclaw/plugin-sdk/session-store-runtime"; -import * as sessionStoreRuntime from "openclaw/plugin-sdk/session-store-runtime"; // Slack tests cover Agent View lifecycle handling. import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { getSlackListenerWriteClient } from "../../client.js"; import { appendSlackStream, markSlackStreamsStopped, startSlackStream } from "../../streaming.js"; +import * as sessionEventRouting from "../message-handler/prepare-routing.js"; import { deliverSlackSlashReplies } from "../replies.js"; import { getSlackSessionRuns, registerSlackSessionRun } from "../session-run-targets.js"; import { getSlackSlashMocks, resetSlackSlashMocks } from "../slash.test-harness.js"; @@ -686,13 +686,14 @@ describe("registerSlackAgentEvents", () => { }, }); await moving.promise; - const readOwner = sessionStoreRuntime.getConversationSession; + const resolveRouting = sessionEventRouting.resolveSlackSessionEventRoutingContext; const lookup = vi - .spyOn(sessionStoreRuntime, "getConversationSession") - .mockImplementationOnce((params) => { - const owner = readOwner(params); + .spyOn(sessionEventRouting, "resolveSlackSessionEventRoutingContext") + .mockImplementationOnce(async (params) => { + const owner = await resolveRouting(params); if (phase === "admission") { releaseMove.resolve(); + await move; } return owner; }); diff --git a/scripts/check-sqlite-transaction-boundary.mts b/scripts/check-sqlite-transaction-boundary.mts index 6b6057c45487..cecb3545a5d1 100644 --- a/scripts/check-sqlite-transaction-boundary.mts +++ b/scripts/check-sqlite-transaction-boundary.mts @@ -16,6 +16,7 @@ const removedAsyncTransactionNames = new Set([ ]); const synchronousTransactionCallbackIndexes = new Map([ ["runOpenClawAgentWriteTransaction", 0], + ["runOpenClawAgentWriteWithYieldingAdmission", 0], ["runOpenClawStateWriteTransaction", 0], ["runSqliteImmediateTransactionSync", 1], ]); diff --git a/src/config/sessions/session-accessor.sqlite-entry-revalidation.test.ts b/src/config/sessions/session-accessor.sqlite-entry-revalidation.test.ts index d2087338ab41..e33d1a78c10b 100644 --- a/src/config/sessions/session-accessor.sqlite-entry-revalidation.test.ts +++ b/src/config/sessions/session-accessor.sqlite-entry-revalidation.test.ts @@ -1,5 +1,6 @@ import fs from "node:fs"; import { DatabaseSync } from "node:sqlite"; +import { setImmediate } from "node:timers/promises"; import { afterEach, beforeEach, describe, expect, it } from "vitest"; import { createTempDirTracker } from "../../../test/helpers/temp-dir.js"; import { @@ -97,6 +98,81 @@ describe("SQLite session entry patch commit revalidation", () => { ); } + it("yields to the event loop while another connection holds the entry write lock", async () => { + const other = new DatabaseSync(database.path); + // Bound the defective synchronous path without a timer or a timing assertion. + database.db.exec("PRAGMA busy_timeout = 250"); + let release: Promise | undefined; + try { + const result = await patchEntry("ordinary", () => { + other.exec("BEGIN IMMEDIATE"); + release = setImmediate().then(() => other.exec("ROLLBACK")); + return { label: "after contention" }; + }); + expect(result?.label).toBe("after contention"); + expect(loadExactSessionEntry(scope)?.entry.label).toBe("after contention"); + } finally { + await release; + other.close(); + } + }); + + it.each(["authority", "row", "connection"] as const)( + "rejects a changed %s after yielding for the entry write lock", + async (changed) => { + const other = new DatabaseSync(database.path); + let current = true; + let updates = 0; + let release: Promise | undefined; + try { + const patch = patchSessionEntryCore( + scope, + () => { + updates += 1; + other.exec("BEGIN IMMEDIATE"); + release = setImmediate().then(() => { + if (changed === "row") { + other + .prepare( + "UPDATE session_nodes SET entry_json = json_set(entry_json, '$.label', 'foreign') WHERE session_key = ?", + ) + .run(sessionKey); + } + other.exec("COMMIT"); + current = false; + if (changed === "connection") { + closeOpenClawAgentDatabaseByPath(database.path); + } + }); + return { label: "must not commit" }; + }, + { + skipMaintenance: true, + assertCommitAllowed: () => { + if (changed === "authority" && !current) { + throw new Error("patch authority revoked"); + } + }, + }, + ); + await expect(patch).rejects.toThrow( + changed === "authority" + ? "patch authority revoked" + : changed === "row" + ? /changed/ + : /closed|replaced|not open/, + ); + expect(updates).toBe(1); + expect(loadExactSessionEntry(scope)?.entry.label).toBe( + changed === "row" ? "foreign" : "original", + ); + } finally { + await release; + other.close(); + } + }, + ); + describe("prepared session mutation guard", () => { function ownerPredicate() { return createSessionTranscriptOwnerPredicate(database, { diff --git a/src/config/sessions/session-accessor.sqlite-entry.env-capture.test.ts b/src/config/sessions/session-accessor.sqlite-entry.env-capture.test.ts index c00386578d08..2f4a315f5656 100644 --- a/src/config/sessions/session-accessor.sqlite-entry.env-capture.test.ts +++ b/src/config/sessions/session-accessor.sqlite-entry.env-capture.test.ts @@ -34,13 +34,15 @@ vi.mock("../../infra/sqlite-number.js", () => ({})); vi.mock("../../state/openclaw-agent-db-identity.js", () => ({})); vi.mock("../../state/openclaw-agent-db-readonly-scope.js", () => ({})); vi.mock("../../state/openclaw-agent-db-readonly.js", () => ({})); +vi.mock("../../state/openclaw-agent-db-transaction.js", () => ({ + runOpenClawAgentWriteWithYieldingAdmission: boundary.commit, +})); vi.mock("../../state/openclaw-agent-db.js", () => ({ getOpenClawAgentDatabaseIfOpen: () => undefined, isIncognitoOpenClawAgentSqlitePath: () => false, openOpenClawAgentDatabase: boundary.open, resolveOpenClawAgentSqlitePath: (options: OpenClawAgentDatabaseOptions) => options.path ?? `${options.env?.OPENCLAW_STATE_DIR}/${options.agentId}.sqlite`, - runOpenClawAgentWriteTransaction: boundary.commit, withOpenClawAgentDatabaseAsync: async ( _options: OpenClawAgentDatabaseOptions, run: () => unknown, diff --git a/src/config/sessions/session-accessor.sqlite-entry.ts b/src/config/sessions/session-accessor.sqlite-entry.ts index ef438fca5402..5110f874f86c 100644 --- a/src/config/sessions/session-accessor.sqlite-entry.ts +++ b/src/config/sessions/session-accessor.sqlite-entry.ts @@ -6,6 +6,7 @@ import { import { coerceRequiredSqliteNumber as sqliteNumber } from "../../infra/sqlite-number.js"; import { OpenClawAgentDatabaseReadOnlyScope } from "../../state/openclaw-agent-db-readonly-scope.js"; import { withOpenClawAgentDatabaseReadOnly } from "../../state/openclaw-agent-db-readonly.js"; +import { runOpenClawAgentWriteWithYieldingAdmission } from "../../state/openclaw-agent-db-transaction.js"; import { getOpenClawAgentDatabaseIfOpen, isIncognitoOpenClawAgentSqlitePath, @@ -504,10 +505,10 @@ async function patchSqliteSessionEntrySnapshot( previous: writeBase, sessionKey, }); - // The updater may dispose the prepared handle; re-admit before the synchronous commit. - return withDatabase(() => { + // The updater may dispose the prepared handle; re-admit before waiting for the write lock. + return withDatabase(async () => { let result: SessionEntry | null = null; - const publish = runOpenClawAgentWriteTransaction( + const publish = await runOpenClawAgentWriteWithYieldingAdmission( (writeDatabase) => { assertCapturedSource(writeDatabase); if (options.shouldCommit?.() === false) { diff --git a/src/state/openclaw-agent-db-transaction.ts b/src/state/openclaw-agent-db-transaction.ts new file mode 100644 index 000000000000..33f2128c9d93 --- /dev/null +++ b/src/state/openclaw-agent-db-transaction.ts @@ -0,0 +1,70 @@ +import { cloneEnvWithPlatformSemantics } from "../config/config-env-vars.js"; +import { runWithSqliteBusyTimeout } from "../infra/sqlite-busy-timeout.js"; +import { withSqlitePostCommitPublications } from "../infra/sqlite-post-commit.js"; +import { + runSqliteImmediateTransaction, + type SqliteTransactionOptions, +} from "../infra/sqlite-transaction.js"; +import { + assertAgentDeletionDatabaseCleanupAccess, + getAgentDeletionDatabaseCleanup, +} from "./agent-deletion-cleanup.js"; +import type { + OpenClawAgentDatabase, + OpenClawAgentDatabaseOptions, +} from "./openclaw-agent-db-contract.js"; +import { + agentDatabaseLifecycle as cache, + retainAgentDatabase, +} from "./openclaw-agent-db-lifecycle.js"; +import { ensureOpenClawAgentDatabasePermissions } from "./openclaw-agent-db-permissions.js"; +import { getOpenClawAgentDatabaseIfOpen, openOpenClawAgentDatabase } from "./openclaw-agent-db.js"; + +/** Yield only for BEGIN admission; admitted writes and publications are never replayed. */ +export async function runOpenClawAgentWriteWithYieldingAdmission( + operation: (database: OpenClawAgentDatabase) => T, + options: OpenClawAgentDatabaseOptions, + transactionOptions: Pick< + SqliteTransactionOptions, + "operationLabel" | "slowTransactionHoldMs" + > = {}, +): Promise { + const captured = { + ...options, + env: cloneEnvWithPlatformSemantics(options.env ?? process.env), + }; + const database = openOpenClawAgentDatabase(captured); + captured.path = database.path; + const release = retainAgentDatabase(database.db); + try { + return await runSqliteImmediateTransaction( + database.db, + async () => () => { + assertAgentDeletionDatabaseCleanupAccess(database, captured); + const result = operation(database); + if (!cache.incognito.has(database)) { + ensureOpenClawAgentDatabasePermissions(database.path, captured); + } + return result; + }, + { + ...transactionOptions, + busyTimeoutMs: 0, + databaseLabel: database.path, + operationLabel: transactionOptions.operationLabel ?? "agent.write", + withCommit: getAgentDeletionDatabaseCleanup(captured)?.withCommit, + }, + (write) => { + if (getOpenClawAgentDatabaseIfOpen(captured) !== database) { + throw new Error(`Agent database closed or replaced before write: ${database.path}`); + } + // Keep zero native wait through COMMIT too; the helper retains the original BEGIN budget. + return runWithSqliteBusyTimeout(database.db, 0, () => + withSqlitePostCommitPublications(database.db, write), + ); + }, + ); + } finally { + release(); + } +} diff --git a/test/scripts/check-sqlite-transaction-boundary.test.ts b/test/scripts/check-sqlite-transaction-boundary.test.ts index d98d8a0edf45..5cb4936f0247 100644 --- a/test/scripts/check-sqlite-transaction-boundary.test.ts +++ b/test/scripts/check-sqlite-transaction-boundary.test.ts @@ -45,6 +45,7 @@ describe("SQLite transaction boundary guard", () => { runSqliteImmediateTransactionSync(db, async () => await prepare()); runOpenClawAgentWriteTransaction(async (database) => await write(database), options); runOpenClawStateWriteTransaction(async (database) => await write(database)); + await runOpenClawAgentWriteWithYieldingAdmission(async (database) => await write(database), options); `), ), ).toEqual([ @@ -63,6 +64,11 @@ describe("SQLite transaction boundary guard", () => { reason: 'passes an async callback to synchronous SQLite transaction helper "runOpenClawStateWriteTransaction"', }, + { + line: 5, + reason: + 'passes an async callback to synchronous SQLite transaction helper "runOpenClawAgentWriteWithYieldingAdmission"', + }, ]); }); @@ -116,6 +122,10 @@ describe("SQLite transaction boundary guard", () => { validate(database, prepared.expected); apply(database, prepared.patch); }, options); + await runOpenClawAgentWriteWithYieldingAdmission((database) => { + validate(database, prepared.expected); + apply(database, prepared.patch); + }, options); `), ), ).toEqual([]);