mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
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.
This commit is contained in:
parent
e1ccea9127
commit
5d605627c5
8 changed files with 180 additions and 9 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
});
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ const removedAsyncTransactionNames = new Set([
|
|||
]);
|
||||
const synchronousTransactionCallbackIndexes = new Map([
|
||||
["runOpenClawAgentWriteTransaction", 0],
|
||||
["runOpenClawAgentWriteWithYieldingAdmission", 0],
|
||||
["runOpenClawStateWriteTransaction", 0],
|
||||
["runSqliteImmediateTransactionSync", 1],
|
||||
]);
|
||||
|
|
|
|||
|
|
@ -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<void> | 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<void> | 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, {
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
70
src/state/openclaw-agent-db-transaction.ts
Normal file
70
src/state/openclaw-agent-db-transaction.ts
Normal file
|
|
@ -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<T>(
|
||||
operation: (database: OpenClawAgentDatabase) => T,
|
||||
options: OpenClawAgentDatabaseOptions,
|
||||
transactionOptions: Pick<
|
||||
SqliteTransactionOptions,
|
||||
"operationLabel" | "slowTransactionHoldMs"
|
||||
> = {},
|
||||
): Promise<T | undefined> {
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
|
@ -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([]);
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue