mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
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:
parent
9ab81b535c
commit
57f912b5f7
15 changed files with 99 additions and 123 deletions
|
|
@ -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.
|
||||
- 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.
|
||||
- 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.
|
||||
- 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.
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
||||
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
|
||||
`scripts/check-native-state-schema-version.mjs`, which fails the build when the
|
||||
|
|
|
|||
|
|
@ -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
|
||||
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
|
||||
event-loop turn for cache freshness. Canonical session validation uses the same
|
||||
schema revision. A migration by another process is detected on the next turn.
|
||||
Shared-state worker admission waits for a preceding probe to expire before
|
||||
starting a new operation, so requests do not reuse an earlier freshness result.
|
||||
Migration and snapshot before/after consistency checks remain fresh reads.
|
||||
rollback, or a foreign commit. A fresh `PRAGMA data_version` probe observes foreign
|
||||
commits on the next unpinned read, even within the same event-loop turn. Actual
|
||||
SQLite read snapshots retain their view until they end; the next read then observes
|
||||
committed changes. Canonical session validation uses the same schema revision.
|
||||
Unchanged versions reuse parsed schema facts and prepared statements without
|
||||
repeating schema scans. Migration and snapshot consistency checks remain fresh reads.
|
||||
This changes no stored schema, migration, durability, or update behavior.
|
||||
|
||||
The nullable requester-authority columns on GitHub publication lifecycle and
|
||||
|
|
|
|||
|
|
@ -573,8 +573,7 @@ describe("SQLite session entry cache", () => {
|
|||
expect(parseSessionEntryCalls).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("fully reloads on the next turn after another connection commits", async () => {
|
||||
vi.useFakeTimers({ toFake: ["setImmediate"] });
|
||||
it("fully reloads on the next read after another connection commits", async () => {
|
||||
const scope = createSessionScope("external-write");
|
||||
const siblingScope = { ...scope, sessionKey: "agent:main:external-write-sibling" };
|
||||
await upsertSessionEntryCore(scope, {
|
||||
|
|
@ -610,7 +609,6 @@ describe("SQLite session entry cache", () => {
|
|||
.run(JSON.stringify(updated), updated.label, updated.updatedAt, scope.sessionKey);
|
||||
|
||||
parseSessionEntryCalls.mockClear();
|
||||
vi.runOnlyPendingTimers();
|
||||
expect(
|
||||
listSessionEntriesCore({ ...scope, clone: false, projection: "list" })[0]?.entry.label,
|
||||
).toBe("projection-probe-after");
|
||||
|
|
@ -622,7 +620,6 @@ describe("SQLite session entry cache", () => {
|
|||
});
|
||||
|
||||
it("fully reloads a cross-connection same-millisecond entry rewrite", async () => {
|
||||
vi.useFakeTimers({ toFake: ["setImmediate"] });
|
||||
const scope = createSessionScope("external-same-ms");
|
||||
const siblingScope = { ...scope, sessionKey: "agent:main:external-same-ms-sibling" };
|
||||
await upsertSessionEntryCore(scope, {
|
||||
|
|
@ -656,7 +653,6 @@ describe("SQLite session entry cache", () => {
|
|||
.run(JSON.stringify(updated), updated.label, scope.sessionKey);
|
||||
|
||||
parseSessionEntryCalls.mockClear();
|
||||
vi.runOnlyPendingTimers();
|
||||
const after = listSessionEntriesCore({ ...scope, clone: false, projection: "list" });
|
||||
|
||||
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 () => {
|
||||
vi.useFakeTimers({ toFake: ["setImmediate"] });
|
||||
it("observes a commit during a listing on the next read", async () => {
|
||||
const scope = createSessionScope("external-race");
|
||||
const siblingScope = { ...scope, sessionKey: "agent:main:external-race-sibling" };
|
||||
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-sibling")?.label).toBe("external-before");
|
||||
vi.runOnlyPendingTimers();
|
||||
expect(
|
||||
listSessionEntriesCore(scope).find(
|
||||
({ entry }) => entry.sessionId === "external-race-sibling",
|
||||
|
|
|
|||
|
|
@ -60,7 +60,11 @@ it("bounds schema and freshness probes across admitted session reader entry poin
|
|||
expect(result.admitted).toBe(true);
|
||||
expect(result.schemaVersion).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") {
|
||||
let allowed = true;
|
||||
|
|
|
|||
|
|
@ -154,6 +154,7 @@ it("measures 100 composed catalog lists against real session and plugin stores",
|
|||
const currentIo = counters.snapshot();
|
||||
workPerList.push({
|
||||
sqliteReadCalls: currentIo.sqliteReadCalls - previousIo.sqliteReadCalls,
|
||||
sqliteFreshnessReads: currentIo.sqliteFreshnessReads - previousIo.sqliteFreshnessReads,
|
||||
bindingAuthorityReads:
|
||||
currentIo.bindingAuthorityReads - previousIo.bindingAuthorityReads,
|
||||
pluginStateWorkerOperations:
|
||||
|
|
@ -228,10 +229,12 @@ it("measures 100 composed catalog lists against real session and plugin stores",
|
|||
expect(io.pluginStateWorkerReadOperations).toBe(0);
|
||||
expect(io.sessionEntryReads).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) {
|
||||
expect(work).toEqual({
|
||||
sqliteReadCalls: 2,
|
||||
sqliteReadCalls: 5,
|
||||
sqliteFreshnessReads: 4,
|
||||
bindingAuthorityReads: 1,
|
||||
pluginStateWorkerOperations: 0,
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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", () => {
|
||||
const filename = path.join(tempDirs.make("openclaw-schema-siblings-"), "state.sqlite");
|
||||
const writer = openDatabase(undefined, true, filename);
|
||||
|
|
|
|||
|
|
@ -24,7 +24,6 @@ type SchemaOwner = {
|
|||
revision: number;
|
||||
facts?: SqliteSchemaFacts;
|
||||
dataVersion?: number;
|
||||
probeTurn?: Promise<void>;
|
||||
transactionalSchema: boolean;
|
||||
transactionalFacts: boolean;
|
||||
snapshot?: object;
|
||||
|
|
@ -239,7 +238,6 @@ function trackSchemaChanges(
|
|||
registerNodeSqliteDisposeCallback(database, () => {
|
||||
invalidate(owner);
|
||||
owner.dataVersion = undefined;
|
||||
owner.probeTurn = undefined;
|
||||
// Native close can still fail; transaction settlement retains pending DDL publication.
|
||||
if (owner.scope) {
|
||||
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 {
|
||||
const tracked = owners.get(database);
|
||||
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) =>
|
||||
statement.get(),
|
||||
);
|
||||
|
|
@ -268,27 +263,10 @@ export function readSqliteCacheDataVersion(database: DatabaseSync): number {
|
|||
invalidate(owner);
|
||||
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;
|
||||
}
|
||||
|
||||
/** 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. */
|
||||
export function trackSqliteSchema(database: DatabaseSync, native: NativeSqlite): void {
|
||||
if (!owners.has(database)) {
|
||||
|
|
|
|||
|
|
@ -119,18 +119,18 @@ beforeAll(async () => {
|
|||
});
|
||||
|
||||
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) => ({
|
||||
owner,
|
||||
userVersion: 0,
|
||||
sqliteMaster: 0,
|
||||
dataVersion: expect.toBeOneOf([0, 1]),
|
||||
})),
|
||||
);
|
||||
});
|
||||
|
||||
it("refuses schemas migrated by another process on the next turn", () => {
|
||||
vi.useFakeTimers({ toFake: ["setImmediate"] });
|
||||
it("refuses schemas migrated by another process on the next read", () => {
|
||||
const scope = {
|
||||
agentId: "main",
|
||||
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" },
|
||||
);
|
||||
vi.runOnlyPendingTimers();
|
||||
expect(() => withOpenClawAgentDatabaseReadOnly(() => undefined, scope)).toThrow(
|
||||
/uses newer schema version/,
|
||||
);
|
||||
|
|
@ -174,6 +173,5 @@ it("refuses schemas migrated by another process on the next turn", () => {
|
|||
}
|
||||
}
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
|
|
|||
|
|
@ -137,13 +137,12 @@ it.each([
|
|||
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 }) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
|
||||
const options = { agentId: "main", env: state.env };
|
||||
const { path } = openOpenClawAgentDatabase(options);
|
||||
await closeOpenClawAgentDatabaseByPathAsync(path);
|
||||
vi.useFakeTimers({ toFake: ["setImmediate"] });
|
||||
const target = { agentId: "main", path };
|
||||
const scope = new OpenClawAgentDatabaseReadOnlyScope();
|
||||
const read = () =>
|
||||
|
|
@ -156,11 +155,9 @@ it.each([
|
|||
} finally {
|
||||
writer.close();
|
||||
}
|
||||
vi.runOnlyPendingTimers();
|
||||
expect(read).toThrow(error);
|
||||
} finally {
|
||||
scope.close();
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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 }) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async ({ env }) => {
|
||||
const options = { agentId: "main", env };
|
||||
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 {
|
||||
const { path: databasePath } = openOpenClawAgentDatabase(options);
|
||||
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);
|
||||
writer.exec(`BEGIN IMMEDIATE; PRAGMA user_version = ${version}; COMMIT;`);
|
||||
} 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/,
|
||||
},
|
||||
])(
|
||||
"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 }) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async ({ env }) => {
|
||||
const options = { agentId: "main", env };
|
||||
const owner = openOpenClawAgentDatabase(options);
|
||||
vi.useFakeTimers({ toFake: ["setImmediate"] });
|
||||
try {
|
||||
inWriterTransaction(owner.db, () => expect(readStamp(options).db.isOpen).toBe(true));
|
||||
owner.db.exec(sql);
|
||||
vi.runOnlyPendingTimers();
|
||||
let invoked = false;
|
||||
inWriterTransaction(owner.db, () => {
|
||||
expect(() =>
|
||||
withOpenClawAgentDatabaseReadOnly(() => {
|
||||
invoked = true;
|
||||
}, options),
|
||||
).toThrow(error);
|
||||
});
|
||||
expect(invoked).toBe(false);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
inWriterTransaction(owner.db, () => expect(readStamp(options).db.isOpen).toBe(true));
|
||||
owner.db.exec(sql);
|
||||
let invoked = false;
|
||||
inWriterTransaction(owner.db, () => {
|
||||
expect(() =>
|
||||
withOpenClawAgentDatabaseReadOnly(() => {
|
||||
invoked = true;
|
||||
}, options),
|
||||
).toThrow(error);
|
||||
});
|
||||
expect(invoked).toBe(false);
|
||||
});
|
||||
},
|
||||
);
|
||||
|
|
|
|||
|
|
@ -15,7 +15,7 @@ import {
|
|||
confirmSqliteFileIntegrity,
|
||||
type SqliteIntegrityConfirmation,
|
||||
} 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 { registerSqliteCacheExitClose, type SqliteWalHealth } from "../infra/sqlite-wal.js";
|
||||
import type { DatabasePathIdentity } from "../infra/sqlite-worker-identity.js";
|
||||
|
|
@ -311,11 +311,6 @@ function publishOpenClawStateDatabase(database: OpenClawStateDatabase): OpenClaw
|
|||
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 {
|
||||
getOpenClawDatabaseMaintenanceScope()?.assertAdmission();
|
||||
assertExistingOpenClawStateSchemaCacheAdmission(pathname, stateDatabaseLifecycle);
|
||||
|
|
@ -628,7 +623,6 @@ export const openClawStateDatabaseCache = {
|
|||
evictCachedOpenClawStateDatabase,
|
||||
evictOpenClawStateDatabaseAfterCorruption,
|
||||
getCachedOpenClawStateDatabase,
|
||||
waitForCachedOpenClawStateSchemaProbe,
|
||||
getOpenClawStateDatabaseRuntimeFailure: runtimeFailures.get,
|
||||
getOpenClawStateDatabaseRecordedFailure: terminalOpenLatch.peek,
|
||||
getOpenClawStateDatabaseIfOpenAtPath,
|
||||
|
|
|
|||
|
|
@ -49,7 +49,7 @@ export function createOpenClawStateDatabaseRuntimeFailureOwner(owner: FailureOwn
|
|||
return undefined;
|
||||
}
|
||||
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);
|
||||
return undefined;
|
||||
} catch (error) {
|
||||
|
|
|
|||
|
|
@ -1,5 +1,5 @@
|
|||
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 { requireNodeSqlite } from "../infra/node-sqlite.js";
|
||||
import { runWithSqliteWorkerStateContext } from "../infra/sqlite-worker-state-context.js";
|
||||
|
|
@ -26,8 +26,7 @@ afterEach(async () => {
|
|||
await state.cleanup();
|
||||
});
|
||||
|
||||
it("admits a worker operation after the previous turn's schema probe expires", async () => {
|
||||
vi.useFakeTimers({ toFake: ["setImmediate"] });
|
||||
it("rechecks a foreign commit before the next worker operation", async () => {
|
||||
const context = captureOpenClawStateWorkerContext();
|
||||
const backend = runWithSqliteWorkerStateContext(context, () =>
|
||||
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 {
|
||||
peer.exec(`PRAGMA user_version = ${OPENCLAW_STATE_SCHEMA_VERSION + 1}`);
|
||||
const command = { type: "database.inspectIdle" as const, input: undefined };
|
||||
const result = Promise.resolve(backend.prepare?.(command))
|
||||
.then(() => runWithSqliteWorkerStateContext(context, () => backend.execute(command)))
|
||||
.then(
|
||||
() => undefined,
|
||||
(error: unknown) => error,
|
||||
);
|
||||
await Promise.resolve();
|
||||
vi.runOnlyPendingTimers();
|
||||
expect(await result).toMatchObject({
|
||||
message: expect.stringContaining("newer schema version"),
|
||||
});
|
||||
expect(() => runWithSqliteWorkerStateContext(context, () => backend.execute(command))).toThrow(
|
||||
"newer schema version",
|
||||
);
|
||||
} finally {
|
||||
peer.exec(`PRAGMA user_version = ${OPENCLAW_STATE_SCHEMA_VERSION}`);
|
||||
peer.close();
|
||||
vi.runOnlyPendingTimers();
|
||||
vi.useRealTimers();
|
||||
await backend.close();
|
||||
}
|
||||
});
|
||||
|
|
|
|||
|
|
@ -116,9 +116,6 @@ function createSharedStateWorkerBackend(
|
|||
return runtime.prepareSharedStateCommand(commandType);
|
||||
});
|
||||
},
|
||||
prepare() {
|
||||
return openClawStateDatabaseCache.waitForCachedOpenClawStateSchemaProbe(context.databasePath);
|
||||
},
|
||||
execute(command) {
|
||||
if (closed) {
|
||||
throw new Error("Shared-state worker is closed");
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue