fix(state): observe foreign commits on the next database read (#156824)

* fix(state): observe foreign commits on the next database read

* test(sqlite): account for fresh version probes

* test(sqlite): bound freshness probes per use instead of per turn

* test(sqlite): cover the worker entry point probe count

---------

Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
Vincent Koc 2026-09-24 09:15:39 +08:00 • committed by GitHub
parent 9ab81b535c
commit 57f912b5f7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
15 changed files with 99 additions and 123 deletions

View file

@ -69,7 +69,7 @@ lightweight artifact, not load the plugin's execution runtime.
- OpenClaw state and caches use SQLite, not new JSON/JSONL/sidecar stores. Files are for named user artifacts, imports/exports, attachments, logs, backups, or external-tool contracts. - OpenClaw state and caches use SQLite, not new JSON/JSONL/sidecar stores. Files are for named user artifacts, imports/exports, attachments, logs, backups, or external-tool contracts.
- Use Kysely for ordinary SQLite access; raw SQL is limited to schema, migrations, bootstrap, and justified primitives. Write transactions are synchronous: finish asynchronous planning first, then reread authoritative rows before writing. No Promise or `await` in a transaction callback. - Use Kysely for ordinary SQLite access; raw SQL is limited to schema, migrations, bootstrap, and justified primitives. Write transactions are synchronous: finish asynchronous planning first, then reread authoritative rows before writing. No Promise or `await` in a transaction callback.
- Database access runs in worker threads, never on the Gateway main thread: readers use the read-only worker scope, writers the SQLite worker broker, and the main thread only awaits results and installs published facts. Boot admission, migrations, Doctor/CLI one-shots, and lock primitives are the only synchronous exceptions. Existing synchronous main-thread access is legacy: never add more, and migrate any such path you touch. - Database access runs in worker threads, never on the Gateway main thread: readers use the read-only worker scope, writers the SQLite worker broker, and the main thread only awaits results and installs published facts. Boot admission, migrations, Doctor/CLI one-shots, and lock primitives are the only synchronous exceptions. Existing synchronous main-thread access is legacy: never add more, and migrate any such path you touch.
- Schema-version, integrity, canonical-index, and table-existence checks belong to database open/admission and the migration owner after migrations; runtime paths carry admitted schema facts with the handle and never re-query them. Bound `PRAGMA data_version` cache probes to one per event-loop turn per handle. Existing per-call checks are legacy: never add more, and migrate any you touch. - Schema-version, integrity, canonical-index, and table-existence checks belong to database open/admission and the migration owner after migrations; runtime paths carry admitted schema facts with the handle and never re-query them. Cache freshness probes must reflect foreign commits on the next unpinned use while preserving actual SQLite snapshot semantics and retained schema facts. Existing per-call checks are legacy: never add more, and migrate any you touch.
- Privileged actions require current owner-held authority. Revalidate after awaited work and immediately before side effects; tokens, signatures, expiry, and matching IDs alone do not prove live authority. - Privileged actions require current owner-held authority. Revalidate after awaited work and immediately before side effects; tokens, signatures, expiry, and matching IDs alone do not prove live authority.
- Core owns shared message tools, action vocabulary, and dispatch. Channels own their account, security, conversation, and transport contracts. Preserve typed command/approval/URL/action distinctions until encoding; never infer product commands from raw strings. - Core owns shared message tools, action vocabulary, and dispatch. Channels own their account, security, conversation, and transport contracts. Preserve typed command/approval/URL/action distinctions until encoding; never infer product commands from raw strings.
- Carry prepared facts through hot paths. Reuse process-stable plugin metadata and lifecycle-owned caches; do not repeatedly load registries or freshness-poll files. Preserve lazy module boundaries and verify relevant builds on the authorized host. - Carry prepared facts through hot paths. Reuse process-stable plugin metadata and lifecycle-owned caches; do not repeatedly load registries or freshness-poll files. Preserve lazy module boundaries and verify relevant builds on the authorized host.

View file

@ -11,7 +11,7 @@ title: "Database schemas"
OpenClaw stores control-plane state in the shared state database and agent data in one SQLite database per agent. Schema migrations run forward when a database opens. Older OpenClaw builds refuse databases written by a newer schema. OpenClaw stores control-plane state in the shared state database and agent data in one SQLite database per agent. Schema migrations run forward when a database opens. Older OpenClaw builds refuse databases written by a newer schema.
Schema-version, integrity, canonical-index, and table-existence checks belong to open/admission and the migration owner after migrations; runtime paths must carry admitted schema facts with the handle, never re-query them, and bound `PRAGMA data_version` cache probes to one per event-loop turn per handle. Existing per-call checks are legacy and must be migrated when touched. Schema-version, integrity, canonical-index, and table-existence checks belong to open/admission and the migration owner after migrations; runtime paths must carry admitted schema facts with the handle, never re-query them, and use fresh `PRAGMA data_version` probes to observe foreign commits on the next unpinned read while preserving active SQLite snapshots. Existing per-call checks are legacy and must be migrated when touched.
Two mechanisms back that contract. CI runs Two mechanisms back that contract. CI runs
`scripts/check-native-state-schema-version.mjs`, which fails the build when the `scripts/check-native-state-schema-version.mjs`, which fails the build when the

View file

@ -23,12 +23,12 @@ Matching numeric versions are necessary but not sufficient. A release can add a
Admitted agent and cached shared-state handles retain their schema version and Admitted agent and cached shared-state handles retain their schema version and
table facts. The handle owner revokes these facts after local DDL, transaction table facts. The handle owner revokes these facts after local DDL, transaction
rollback, or a foreign commit; it checks `PRAGMA data_version` at most once per rollback, or a foreign commit. A fresh `PRAGMA data_version` probe observes foreign
event-loop turn for cache freshness. Canonical session validation uses the same commits on the next unpinned read, even within the same event-loop turn. Actual
schema revision. A migration by another process is detected on the next turn. SQLite read snapshots retain their view until they end; the next read then observes
Shared-state worker admission waits for a preceding probe to expire before committed changes. Canonical session validation uses the same schema revision.
starting a new operation, so requests do not reuse an earlier freshness result. Unchanged versions reuse parsed schema facts and prepared statements without
Migration and snapshot before/after consistency checks remain fresh reads. repeating schema scans. Migration and snapshot consistency checks remain fresh reads.
This changes no stored schema, migration, durability, or update behavior. This changes no stored schema, migration, durability, or update behavior.
The nullable requester-authority columns on GitHub publication lifecycle and The nullable requester-authority columns on GitHub publication lifecycle and

View file

@ -573,8 +573,7 @@ describe("SQLite session entry cache", () => {
expect(parseSessionEntryCalls).not.toHaveBeenCalled(); expect(parseSessionEntryCalls).not.toHaveBeenCalled();
}); });
it("fully reloads on the next turn after another connection commits", async () => { it("fully reloads on the next read after another connection commits", async () => {
vi.useFakeTimers({ toFake: ["setImmediate"] });
const scope = createSessionScope("external-write"); const scope = createSessionScope("external-write");
const siblingScope = { ...scope, sessionKey: "agent:main:external-write-sibling" }; const siblingScope = { ...scope, sessionKey: "agent:main:external-write-sibling" };
await upsertSessionEntryCore(scope, { await upsertSessionEntryCore(scope, {
@ -610,7 +609,6 @@ describe("SQLite session entry cache", () => {
.run(JSON.stringify(updated), updated.label, updated.updatedAt, scope.sessionKey); .run(JSON.stringify(updated), updated.label, updated.updatedAt, scope.sessionKey);
parseSessionEntryCalls.mockClear(); parseSessionEntryCalls.mockClear();
vi.runOnlyPendingTimers();
expect( expect(
listSessionEntriesCore({ ...scope, clone: false, projection: "list" })[0]?.entry.label, listSessionEntriesCore({ ...scope, clone: false, projection: "list" })[0]?.entry.label,
).toBe("projection-probe-after"); ).toBe("projection-probe-after");
@ -622,7 +620,6 @@ describe("SQLite session entry cache", () => {
}); });
it("fully reloads a cross-connection same-millisecond entry rewrite", async () => { it("fully reloads a cross-connection same-millisecond entry rewrite", async () => {
vi.useFakeTimers({ toFake: ["setImmediate"] });
const scope = createSessionScope("external-same-ms"); const scope = createSessionScope("external-same-ms");
const siblingScope = { ...scope, sessionKey: "agent:main:external-same-ms-sibling" }; const siblingScope = { ...scope, sessionKey: "agent:main:external-same-ms-sibling" };
await upsertSessionEntryCore(scope, { await upsertSessionEntryCore(scope, {
@ -656,7 +653,6 @@ describe("SQLite session entry cache", () => {
.run(JSON.stringify(updated), updated.label, scope.sessionKey); .run(JSON.stringify(updated), updated.label, scope.sessionKey);
parseSessionEntryCalls.mockClear(); parseSessionEntryCalls.mockClear();
vi.runOnlyPendingTimers();
const after = listSessionEntriesCore({ ...scope, clone: false, projection: "list" }); const after = listSessionEntriesCore({ ...scope, clone: false, projection: "list" });
expect(after[0]?.entry.label).toBe("projection-probe-after"); expect(after[0]?.entry.label).toBe("projection-probe-after");
@ -667,8 +663,7 @@ describe("SQLite session entry cache", () => {
} }
}); });
it("observes a commit during a listing on the next turn", async () => { it("observes a commit during a listing on the next read", async () => {
vi.useFakeTimers({ toFake: ["setImmediate"] });
const scope = createSessionScope("external-race"); const scope = createSessionScope("external-race");
const siblingScope = { ...scope, sessionKey: "agent:main:external-race-sibling" }; const siblingScope = { ...scope, sessionKey: "agent:main:external-race-sibling" };
await upsertSessionEntryCore(scope, { await upsertSessionEntryCore(scope, {
@ -718,7 +713,6 @@ describe("SQLite session entry cache", () => {
expect(byId.get("external-race-local")?.label).toBe("local-after"); expect(byId.get("external-race-local")?.label).toBe("local-after");
expect(byId.get("external-race-sibling")?.label).toBe("external-before"); expect(byId.get("external-race-sibling")?.label).toBe("external-before");
vi.runOnlyPendingTimers();
expect( expect(
listSessionEntriesCore(scope).find( listSessionEntriesCore(scope).find(
({ entry }) => entry.sessionId === "external-race-sibling", ({ entry }) => entry.sessionId === "external-race-sibling",

View file

@ -60,7 +60,11 @@ it("bounds schema and freshness probes across admitted session reader entry poin
expect(result.admitted).toBe(true); expect(result.admitted).toBe(true);
expect(result.schemaVersion).toBe(0); expect(result.schemaVersion).toBe(0);
expect(result.userVersion).toBe(0); expect(result.userVersion).toBe(0);
expect(result.dataVersion).toBeLessThanOrEqual(1); // Freshness probes run on every use so foreign commits are visible on the next
// read; today one session read passes through 9-11 layered uses depending on the
// entry point. The bound holds that ceiling until the layers share one probe per
// operation.
expect(result.dataVersion).toBeLessThanOrEqual(12 * 100);
} }
if (typeof writer.db.setAuthorizer === "function") { if (typeof writer.db.setAuthorizer === "function") {
let allowed = true; let allowed = true;

View file

@ -154,6 +154,7 @@ it("measures 100 composed catalog lists against real session and plugin stores",
const currentIo = counters.snapshot(); const currentIo = counters.snapshot();
workPerList.push({ workPerList.push({
sqliteReadCalls: currentIo.sqliteReadCalls - previousIo.sqliteReadCalls, sqliteReadCalls: currentIo.sqliteReadCalls - previousIo.sqliteReadCalls,
sqliteFreshnessReads: currentIo.sqliteFreshnessReads - previousIo.sqliteFreshnessReads,
bindingAuthorityReads: bindingAuthorityReads:
currentIo.bindingAuthorityReads - previousIo.bindingAuthorityReads, currentIo.bindingAuthorityReads - previousIo.bindingAuthorityReads,
pluginStateWorkerOperations: pluginStateWorkerOperations:
@ -228,10 +229,12 @@ it("measures 100 composed catalog lists against real session and plugin stores",
expect(io.pluginStateWorkerReadOperations).toBe(0); expect(io.pluginStateWorkerReadOperations).toBe(0);
expect(io.sessionEntryReads).toBe(0); expect(io.sessionEntryReads).toBe(0);
expect(io.sessionPayloadReads).toBe(0); expect(io.sessionPayloadReads).toBe(0);
// The adopted cohort shares bounded freshness and authority reads with admitted schema facts. // Cached-handle and reused-read admission each check published/content freshness.
// The adopted cohort still shares one bulk binding query without rescanning rows.
for (const work of workPerList) { for (const work of workPerList) {
expect(work).toEqual({ expect(work).toEqual({
sqliteReadCalls: 2, sqliteReadCalls: 5,
sqliteFreshnessReads: 4,
bindingAuthorityReads: 1, bindingAuthorityReads: 1,
pluginStateWorkerOperations: 0, pluginStateWorkerOperations: 0,
}); });

View file

@ -39,6 +39,42 @@ describe("admitted SQLite schema facts", () => {
} }
}); });
it.each(["transaction", "implicit snapshot"])(
"observes foreign commits on the next read while preserving an active %s",
(pin) => {
const filename = path.join(tempDirs.make("openclaw-schema-foreign-"), "state.sqlite");
const reader = openDatabase(undefined, true, filename);
reader.exec("PRAGMA journal_mode=WAL");
// Bypass local schema publications, as a worker or another process does.
const writer = new DatabaseSync(filename);
databases.push(writer);
expect(tableExists(reader, "committed")).toBe(false);
writer.exec("BEGIN; CREATE TABLE committed (id); PRAGMA user_version = 2; COMMIT;");
expect(tableExists(reader, "committed")).toBe(true);
expect(assertSupportedAgentSchemaVersion(reader, filename)).toBe(2);
const readSnapshot = () => {
expect(tableExists(reader, "later")).toBe(false);
writer.exec("BEGIN; CREATE TABLE later (id); PRAGMA user_version = 3; COMMIT;");
expect(tableExists(reader, "later")).toBe(false);
expect(assertSupportedAgentSchemaVersion(reader, filename)).toBe(2);
};
if (pin === "transaction") {
reader.exec("BEGIN");
try {
reader.prepare("SELECT id FROM original").all();
readSnapshot();
} finally {
reader.exec("COMMIT");
}
} else {
runSqlitePinnedReadSnapshotSync(reader, readSnapshot);
}
expect(tableExists(reader, "later")).toBe(true);
expect(assertSupportedAgentSchemaVersion(reader, filename)).toBe(3);
},
);
it("publishes local DDL to sibling handles while preserving their active snapshots", () => { it("publishes local DDL to sibling handles while preserving their active snapshots", () => {
const filename = path.join(tempDirs.make("openclaw-schema-siblings-"), "state.sqlite"); const filename = path.join(tempDirs.make("openclaw-schema-siblings-"), "state.sqlite");
const writer = openDatabase(undefined, true, filename); const writer = openDatabase(undefined, true, filename);

View file

@ -24,7 +24,6 @@ type SchemaOwner = {
revision: number; revision: number;
facts?: SqliteSchemaFacts; facts?: SqliteSchemaFacts;
dataVersion?: number; dataVersion?: number;
probeTurn?: Promise<void>;
transactionalSchema: boolean; transactionalSchema: boolean;
transactionalFacts: boolean; transactionalFacts: boolean;
snapshot?: object; snapshot?: object;
@ -239,7 +238,6 @@ function trackSchemaChanges(
registerNodeSqliteDisposeCallback(database, () => { registerNodeSqliteDisposeCallback(database, () => {
invalidate(owner); invalidate(owner);
owner.dataVersion = undefined; owner.dataVersion = undefined;
owner.probeTurn = undefined;
// Native close can still fail; transaction settlement retains pending DDL publication. // Native close can still fail; transaction settlement retains pending DDL publication.
if (owner.scope) { if (owner.scope) {
scopes.finalizer.unregister(owner); scopes.finalizer.unregister(owner);
@ -250,13 +248,10 @@ function trackSchemaChanges(
}); });
} }
/** Cache freshness only: migration and snapshot before/after probes must remain uncached. */ /** Foreign commits can occur within one JS turn; only SQLite owns snapshot visibility. */
export function readSqliteCacheDataVersion(database: DatabaseSync): number { export function readSqliteCacheDataVersion(database: DatabaseSync): number {
const tracked = owners.get(database); const tracked = owners.get(database);
const owner = tracked?.admitted ? tracked : undefined; const owner = tracked?.admitted ? tracked : undefined;
if (owner?.probeTurn && !owner.authorizerActive && owner.dataVersion !== undefined) {
return owner.dataVersion;
}
const row = executeWithCachedStatement(database, "PRAGMA data_version", [], (statement) => const row = executeWithCachedStatement(database, "PRAGMA data_version", [], (statement) =>
statement.get(), statement.get(),
); );
@ -268,27 +263,10 @@ export function readSqliteCacheDataVersion(database: DatabaseSync): number {
invalidate(owner); invalidate(owner);
owner.dataVersion = row.data_version; owner.dataVersion = row.data_version;
} }
if (!owner.probeTurn) {
// Retain only the facts, never a native handle or statement, until the next turn.
const turn = new Promise<void>((resolve) => {
setImmediate(() => {
if (owner.probeTurn === turn) {
owner.probeTurn = undefined;
}
resolve();
}).unref();
});
owner.probeTurn = turn;
}
} }
return row.data_version; return row.data_version;
} }
/** A new asynchronous operation must not inherit a preceding operation's freshness probe. */
export function waitForSqliteSchemaProbeTurn(database: DatabaseSync): Promise<void> | undefined {
return owners.get(database)?.probeTurn;
}
/** Install at native open, before callers can retain statements or install an authorizer. */ /** Install at native open, before callers can retain statements or install an authorizer. */
export function trackSqliteSchema(database: DatabaseSync, native: NativeSqlite): void { export function trackSqliteSchema(database: DatabaseSync, native: NativeSqlite): void {
if (!owners.has(database)) { if (!owners.has(database)) {

View file

@ -119,18 +119,18 @@ beforeAll(async () => {
}); });
it("keeps admitted reads within the schema-query budget", () => { it("keeps admitted reads within the schema-query budget", () => {
expect(counts).toEqual( expect(
counts.map(({ owner, userVersion, sqliteMaster }) => ({ owner, userVersion, sqliteMaster })),
).toEqual(
["agent", "state"].map((owner) => ({ ["agent", "state"].map((owner) => ({
owner, owner,
userVersion: 0, userVersion: 0,
sqliteMaster: 0, sqliteMaster: 0,
dataVersion: expect.toBeOneOf([0, 1]),
})), })),
); );
}); });
it("refuses schemas migrated by another process on the next turn", () => { it("refuses schemas migrated by another process on the next read", () => {
vi.useFakeTimers({ toFake: ["setImmediate"] });
const scope = { const scope = {
agentId: "main", agentId: "main",
env: { ...process.env, OPENCLAW_STATE_DIR: tempDirs.make("openclaw-schema-migration-") }, env: { ...process.env, OPENCLAW_STATE_DIR: tempDirs.make("openclaw-schema-migration-") },
@ -158,7 +158,6 @@ it("refuses schemas migrated by another process on the next turn", () => {
], ],
{ stdio: "pipe" }, { stdio: "pipe" },
); );
vi.runOnlyPendingTimers();
expect(() => withOpenClawAgentDatabaseReadOnly(() => undefined, scope)).toThrow( expect(() => withOpenClawAgentDatabaseReadOnly(() => undefined, scope)).toThrow(
/uses newer schema version/, /uses newer schema version/,
); );
@ -174,6 +173,5 @@ it("refuses schemas migrated by another process on the next turn", () => {
} }
} }
closeOpenClawStateDatabaseForTest(); closeOpenClawStateDatabaseForTest();
vi.useRealTimers();
} }
}); });

View file

@ -137,13 +137,12 @@ it.each([
error: "has schema role state", error: "has schema role state",
}, },
])( ])(
"revalidates retained read admission on the next turn after a commit: $sql", "revalidates retained read admission on the next read after a commit: $sql",
async ({ sql, error }) => { async ({ sql, error }) => {
await withOpenClawTestState({ scenario: "minimal" }, async (state) => { await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
const options = { agentId: "main", env: state.env }; const options = { agentId: "main", env: state.env };
const { path } = openOpenClawAgentDatabase(options); const { path } = openOpenClawAgentDatabase(options);
await closeOpenClawAgentDatabaseByPathAsync(path); await closeOpenClawAgentDatabaseByPathAsync(path);
vi.useFakeTimers({ toFake: ["setImmediate"] });
const target = { agentId: "main", path }; const target = { agentId: "main", path };
const scope = new OpenClawAgentDatabaseReadOnlyScope(); const scope = new OpenClawAgentDatabaseReadOnlyScope();
const read = () => const read = () =>
@ -156,11 +155,9 @@ it.each([
} finally { } finally {
writer.close(); writer.close();
} }
vi.runOnlyPendingTimers();
expect(read).toThrow(error); expect(read).toThrow(error);
} finally { } finally {
scope.close(); scope.close();
vi.useRealTimers();
} }
}); });
}, },

View file

@ -152,39 +152,31 @@ describe("committed agent database reads", () => {
}, },
}, },
])( ])(
"rejects a committed version change to $version on the next turn before borrowed and fresh reads", "rejects a committed version change to $version on the next read before borrowed and fresh reads",
async ({ version, expectedError }) => { async ({ version, expectedError }) => {
await withOpenClawTestState({ scenario: "minimal" }, async ({ env }) => { await withOpenClawTestState({ scenario: "minimal" }, async ({ env }) => {
const options = { agentId: "main", env }; const options = { agentId: "main", env };
closeOpenClawAgentDatabaseByPath(resolveOpenClawAgentSqlitePath(options)); closeOpenClawAgentDatabaseByPath(resolveOpenClawAgentSqlitePath(options));
vi.useFakeTimers({ toFake: ["setImmediate"] }); const { path: databasePath } = openOpenClawAgentDatabase(options);
let admitted: boolean;
const read = () =>
withOpenClawAgentDatabaseReadOnly(({ db }) => {
admitted = true;
return db.prepare("SELECT agent_id FROM schema_meta WHERE meta_key = 'primary'").get();
}, options);
expect(read()).toEqual({ found: true, value: { agent_id: "main" } });
admitted = false;
const writer = new DatabaseSync(databasePath);
try { try {
const { path: databasePath } = openOpenClawAgentDatabase(options); writer.exec(`BEGIN IMMEDIATE; PRAGMA user_version = ${version}; COMMIT;`);
let admitted = false;
const read = () =>
withOpenClawAgentDatabaseReadOnly(({ db }) => {
admitted = true;
return db
.prepare("SELECT agent_id FROM schema_meta WHERE meta_key = 'primary'")
.get();
}, options);
expect(read()).toEqual({ found: true, value: { agent_id: "main" } });
admitted = false;
const writer = new DatabaseSync(databasePath);
try {
writer.exec(`BEGIN IMMEDIATE; PRAGMA user_version = ${version}; COMMIT;`);
} finally {
writer.close();
}
vi.runOnlyPendingTimers();
expect(read).toThrow(expect.objectContaining(expectedError));
expect(admitted).toBe(false);
expect(closeOpenClawAgentDatabaseByPath(databasePath)).toBe(true);
expect(read).toThrow(expect.objectContaining(expectedError));
expect(admitted).toBe(false);
} finally { } finally {
vi.useRealTimers(); writer.close();
} }
expect(read).toThrow(expect.objectContaining(expectedError));
expect(admitted).toBe(false);
expect(closeOpenClawAgentDatabaseByPath(databasePath)).toBe(true);
expect(read).toThrow(expect.objectContaining(expectedError));
expect(admitted).toBe(false);
}); });
}, },
); );
@ -206,28 +198,22 @@ describe("committed agent database reads", () => {
error: /belongs to agent other.*requested agent main/, error: /belongs to agent other.*requested agent main/,
}, },
])( ])(
"rechecks committed $name on the next turn before invoking a retained reader", "rechecks committed $name on the next read before invoking a retained reader",
async ({ sql, error }) => { async ({ sql, error }) => {
await withOpenClawTestState({ scenario: "minimal" }, async ({ env }) => { await withOpenClawTestState({ scenario: "minimal" }, async ({ env }) => {
const options = { agentId: "main", env }; const options = { agentId: "main", env };
const owner = openOpenClawAgentDatabase(options); const owner = openOpenClawAgentDatabase(options);
vi.useFakeTimers({ toFake: ["setImmediate"] }); inWriterTransaction(owner.db, () => expect(readStamp(options).db.isOpen).toBe(true));
try { owner.db.exec(sql);
inWriterTransaction(owner.db, () => expect(readStamp(options).db.isOpen).toBe(true)); let invoked = false;
owner.db.exec(sql); inWriterTransaction(owner.db, () => {
vi.runOnlyPendingTimers(); expect(() =>
let invoked = false; withOpenClawAgentDatabaseReadOnly(() => {
inWriterTransaction(owner.db, () => { invoked = true;
expect(() => }, options),
withOpenClawAgentDatabaseReadOnly(() => { ).toThrow(error);
invoked = true; });
}, options), expect(invoked).toBe(false);
).toThrow(error);
});
expect(invoked).toBe(false);
} finally {
vi.useRealTimers();
}
}); });
}, },
); );

View file

@ -15,7 +15,7 @@ import {
confirmSqliteFileIntegrity, confirmSqliteFileIntegrity,
type SqliteIntegrityConfirmation, type SqliteIntegrityConfirmation,
} from "../infra/sqlite-integrity.js"; } from "../infra/sqlite-integrity.js";
import { admitSqliteSchema, waitForSqliteSchemaProbeTurn } from "../infra/sqlite-schema-facts.js"; import { admitSqliteSchema } from "../infra/sqlite-schema-facts.js";
import { createSqliteTerminalOpenLatch } from "../infra/sqlite-terminal-open-latch.js"; import { createSqliteTerminalOpenLatch } from "../infra/sqlite-terminal-open-latch.js";
import { registerSqliteCacheExitClose, type SqliteWalHealth } from "../infra/sqlite-wal.js"; import { registerSqliteCacheExitClose, type SqliteWalHealth } from "../infra/sqlite-wal.js";
import type { DatabasePathIdentity } from "../infra/sqlite-worker-identity.js"; import type { DatabasePathIdentity } from "../infra/sqlite-worker-identity.js";
@ -311,11 +311,6 @@ function publishOpenClawStateDatabase(database: OpenClawStateDatabase): OpenClaw
return database; return database;
} }
function waitForCachedOpenClawStateSchemaProbe(pathname: string): Promise<void> | undefined {
const database = cachedDatabases.get(pathname);
return database ? waitForSqliteSchemaProbeTurn(database.db) : undefined;
}
function getCachedOpenClawStateDatabase(pathname: string): OpenClawStateDatabase | undefined { function getCachedOpenClawStateDatabase(pathname: string): OpenClawStateDatabase | undefined {
getOpenClawDatabaseMaintenanceScope()?.assertAdmission(); getOpenClawDatabaseMaintenanceScope()?.assertAdmission();
assertExistingOpenClawStateSchemaCacheAdmission(pathname, stateDatabaseLifecycle); assertExistingOpenClawStateSchemaCacheAdmission(pathname, stateDatabaseLifecycle);
@ -628,7 +623,6 @@ export const openClawStateDatabaseCache = {
evictCachedOpenClawStateDatabase, evictCachedOpenClawStateDatabase,
evictOpenClawStateDatabaseAfterCorruption, evictOpenClawStateDatabaseAfterCorruption,
getCachedOpenClawStateDatabase, getCachedOpenClawStateDatabase,
waitForCachedOpenClawStateSchemaProbe,
getOpenClawStateDatabaseRuntimeFailure: runtimeFailures.get, getOpenClawStateDatabaseRuntimeFailure: runtimeFailures.get,
getOpenClawStateDatabaseRecordedFailure: terminalOpenLatch.peek, getOpenClawStateDatabaseRecordedFailure: terminalOpenLatch.peek,
getOpenClawStateDatabaseIfOpenAtPath, getOpenClawStateDatabaseIfOpenAtPath,

View file

@ -49,7 +49,7 @@ export function createOpenClawStateDatabaseRuntimeFailureOwner(owner: FailureOwn
return undefined; return undefined;
} }
try { try {
// Admission owns schema facts and bounds the foreign-commit probe to one per turn. // Admission retains schema facts but checks foreign commits before reusing them.
assertSupportedStateSchemaVersion(cached.db, resolvedPath); assertSupportedStateSchemaVersion(cached.db, resolvedPath);
return undefined; return undefined;
} catch (error) { } catch (error) {

View file

@ -1,5 +1,5 @@
import { existsSync } from "node:fs"; import { existsSync } from "node:fs";
import { afterEach, beforeEach, expect, it, vi } from "vitest"; import { afterEach, beforeEach, expect, it } from "vitest";
import { getNodeSqliteKysely, iterateSqliteQuerySync } from "../infra/kysely-sync.js"; import { getNodeSqliteKysely, iterateSqliteQuerySync } from "../infra/kysely-sync.js";
import { requireNodeSqlite } from "../infra/node-sqlite.js"; import { requireNodeSqlite } from "../infra/node-sqlite.js";
import { runWithSqliteWorkerStateContext } from "../infra/sqlite-worker-state-context.js"; import { runWithSqliteWorkerStateContext } from "../infra/sqlite-worker-state-context.js";
@ -26,8 +26,7 @@ afterEach(async () => {
await state.cleanup(); await state.cleanup();
}); });
it("admits a worker operation after the previous turn's schema probe expires", async () => { it("rechecks a foreign commit before the next worker operation", async () => {
vi.useFakeTimers({ toFake: ["setImmediate"] });
const context = captureOpenClawStateWorkerContext(); const context = captureOpenClawStateWorkerContext();
const backend = runWithSqliteWorkerStateContext(context, () => const backend = runWithSqliteWorkerStateContext(context, () =>
createSqliteWorkerBackend(undefined, { databasePath: context.admission.databasePath }), createSqliteWorkerBackend(undefined, { databasePath: context.admission.databasePath }),
@ -36,22 +35,12 @@ it("admits a worker operation after the previous turn's schema probe expires", a
try { try {
peer.exec(`PRAGMA user_version = ${OPENCLAW_STATE_SCHEMA_VERSION + 1}`); peer.exec(`PRAGMA user_version = ${OPENCLAW_STATE_SCHEMA_VERSION + 1}`);
const command = { type: "database.inspectIdle" as const, input: undefined }; const command = { type: "database.inspectIdle" as const, input: undefined };
const result = Promise.resolve(backend.prepare?.(command)) expect(() => runWithSqliteWorkerStateContext(context, () => backend.execute(command))).toThrow(
.then(() => runWithSqliteWorkerStateContext(context, () => backend.execute(command))) "newer schema version",
.then( );
() => undefined,
(error: unknown) => error,
);
await Promise.resolve();
vi.runOnlyPendingTimers();
expect(await result).toMatchObject({
message: expect.stringContaining("newer schema version"),
});
} finally { } finally {
peer.exec(`PRAGMA user_version = ${OPENCLAW_STATE_SCHEMA_VERSION}`); peer.exec(`PRAGMA user_version = ${OPENCLAW_STATE_SCHEMA_VERSION}`);
peer.close(); peer.close();
vi.runOnlyPendingTimers();
vi.useRealTimers();
await backend.close(); await backend.close();
} }
}); });

View file

@ -116,9 +116,6 @@ function createSharedStateWorkerBackend(
return runtime.prepareSharedStateCommand(commandType); return runtime.prepareSharedStateCommand(commandType);
}); });
}, },
prepare() {
return openClawStateDatabaseCache.waitForCachedOpenClawStateSchemaProbe(context.databasePath);
},
execute(command) { execute(command) {
if (closed) { if (closed) {
throw new Error("Shared-state worker is closed"); throw new Error("Shared-state worker is closed");