From 6fa64d9f7182f405759ccddb427e35192a4b01b2 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 17 Sep 2026 01:33:59 -0700 Subject: [PATCH] fix(memory): reduce CPU work during index sync (#150602) --- config/assertion-safety-baseline.txt | 2 +- docs/plugins/sdk-subpaths.md | 3 +- .../imessage/src/send-receipt-db.worker.ts | 2 +- extensions/logbook/src/store-frame-queries.ts | 2 +- extensions/logbook/src/store-prune.test.ts | 6 +- .../logbook/src/store-read-budgets.test.ts | 6 +- extensions/logbook/src/store-schema.ts | 2 +- extensions/logbook/src/store.worker.ts | 2 +- .../src/memory-session-tombstones.ts | 2 +- .../src/memory/manager-chunk-writer.ts | 4 +- .../src/memory/manager-database-context.ts | 30 ++- .../src/memory/manager-db-kernel.ts | 191 ++++++++++++++++++ .../memory-core/src/memory/manager-db.test.ts | 8 +- .../memory-core/src/memory/manager-db.ts | 191 +----------------- .../src/memory/manager-embedding-ops.ts | 2 +- .../manager-publication-transfer.test.ts | 22 +- .../src/memory/manager-publication.worker.ts | 6 +- .../src/memory/manager-search-maintenance.ts | 2 +- .../manager-search-orchestration.test.ts | 2 +- .../src/memory/manager-source-index-kernel.ts | 7 +- .../src/memory/manager-source-sync-ops.ts | 2 +- .../src/memory/manager-sync-ops.ts | 8 +- .../memory/manager-vector-rebuild-state.ts | 6 +- .../manager.borrowed-connection.test.ts | 2 +- extensions/memory-core/src/memory/manager.ts | 26 +-- .../src/tools.index-diagnostic.test.ts | 2 +- .../src/tools.index-upgrade.test.ts | 2 +- .../src/tools.real-manager.test.ts | 2 +- ...tion-identity-storage-inspection.worker.ts | 2 +- extensions/team-reports/src/store.worker.ts | 2 +- .../tsconfig.package-boundary.paths.json | 3 + .../src/sqlite-store-batch-read.test.ts | 4 +- .../workboard/src/sqlite-store-kernel.ts | 2 +- .../workboard/src/sqlite-store-policy.test.ts | 2 +- .../workboard/src/sqlite-store-records.ts | 2 +- .../workboard/src/sqlite-store-schema.ts | 2 +- .../workboard/src/sqlite-store-write.ts | 2 +- .../workboard/src/sqlite-store.worker.ts | 5 +- extensions/xai/tsconfig.json | 3 + package.json | 6 +- scripts/lib/plugin-sdk-entrypoints.json | 1 + ...lugin-sdk-private-local-only-subpaths.json | 1 + .../memory-core-host-engine-schema.ts | 8 +- src/plugin-sdk/sqlite-worker-runtime.ts | 28 +++ 44 files changed, 345 insertions(+), 270 deletions(-) create mode 100644 extensions/memory-core/src/memory/manager-db-kernel.ts create mode 100644 src/plugin-sdk/sqlite-worker-runtime.ts diff --git a/config/assertion-safety-baseline.txt b/config/assertion-safety-baseline.txt index 8f256cffaf43..6004b7463c7c 100644 --- a/config/assertion-safety-baseline.txt +++ b/config/assertion-safety-baseline.txt @@ -710,7 +710,7 @@ extensions/memory-core/src/dreaming-phases.ts 2 extensions/memory-core/src/dreaming.ts 1 extensions/memory-core/src/memory/embeddings.ts 2 extensions/memory-core/src/memory/hybrid.ts 1 -extensions/memory-core/src/memory/manager-db.ts 5 +extensions/memory-core/src/memory/manager-db.ts 2 extensions/memory-core/src/memory/manager-embedding-errors.ts 2 extensions/memory-core/src/memory/manager-embedding-ops.ts 1 extensions/memory-core/src/memory/manager-keyword-retrieval.ts 1 diff --git a/docs/plugins/sdk-subpaths.md b/docs/plugins/sdk-subpaths.md index 6c236e423458..87baf3a08c09 100644 --- a/docs/plugins/sdk-subpaths.md +++ b/docs/plugins/sdk-subpaths.md @@ -374,6 +374,7 @@ Use `isLoopbackHost(host)` when a plugin must accept only the local machine. It | `plugin-sdk/session-discussion` | External session discussion provider contracts, registration, and canonical Control UI session path building | | `plugin-sdk/session-transcript-runtime` | Private-local after July 2026; Transcript identity, bounded raw and visible cursors, scoped target/read/write helpers, read-only native transcript catalog pages with portable sender attribution, visible message-entry projection, update publishing, write locks, and transcript memory hit keys | | `plugin-sdk/sqlite-runtime` | Private-local after July 2026; SQLite agent-schema, path, transaction, and shared-handle borrowing helpers for first-party runtime. Type-only `Generated` and `Selectable` model generated columns and selected rows in Kysely table definitions. `compileSqliteQueryBindings` compiles fixed Kysely SQL with fresh bindings for caller-owned native statements; statement lifetime stays with the caller. `iterateSqliteQuerySync` streams Kysely query rows for incremental decoding; consume the iterator before closing its database. `sqliteStringSet` binds string membership as one SQLite JSON table-valued query parameter, preserving one query snapshot without a variable-count placeholder list. `borrowOpenClawAgentDatabase` returns `{ db, release }`; active borrows prevent cache eviction, `release()` does not close the handle, and explicit owner disposal still revokes it. `withOpenClawAgentDatabaseAsync` admits a connection asynchronously and retains its original identity through the operation’s settlement; it does not move callback execution off the caller thread. | + | `plugin-sdk/sqlite-worker-runtime` | Private-local native SQLite, Kysely query, synchronous transaction, and operation-admission primitives for worker backends and their shared helpers. Imports existing owners directly so workers do not load host database lifecycle code. Host store creation and lifecycle coordination remain on `sqlite-runtime`. | | `plugin-sdk/cron-store-runtime` | Private-local after July 2026; Cron store path/load/save helpers | | `plugin-sdk/state-paths` | State/OAuth dir path helpers | | `plugin-sdk/plugin-state-runtime` | Private-local after July 2026; Plugin-scoped keyed-state and BLOB contracts plus connection pragma, verified WAL maintenance, and atomic STRICT-schema migration helpers. Plugin-state leases were removed; use SQLite transactions and keyed stores instead | @@ -521,7 +522,7 @@ Use `isLoopbackHost(host)` when a plugin must accept only the local machine. It | `plugin-sdk/memory-core-host-engine-fs` | Private-local focused filesystem and user-path helpers for doctor migrations | | `plugin-sdk/memory-core-host-engine-embeddings` | Private-local after July 2026; Memory host embedding contracts and batch/remote helpers. Providers register through the generic embedding provider API. | | `plugin-sdk/memory-core-host-engine-sessions` | Private-local after July 2026; Memory session transcript and query helpers | - | `plugin-sdk/memory-core-host-engine-schema` | Private-local focused memory index schema and sqlite-vec helpers for doctor migrations | + | `plugin-sdk/memory-core-host-engine-schema` | Private-local memory index schema and sqlite-vec operations shared by Doctor, host maintenance, and native publication workers | | `plugin-sdk/memory-core-host-engine-indexing` | Private-local immutable chunk preparation, annotations, hashes, and embedding input limits for indexing workers | | `plugin-sdk/memory-core-host-engine-knn` | Private-local read-only SQLite ownership checks, sqlite-vec, and text/vector primitives for isolated retrieval workers and children | | `plugin-sdk/memory-core-host-engine-storage` | Private-local after July 2026; Memory host storage engine exports | diff --git a/extensions/imessage/src/send-receipt-db.worker.ts b/extensions/imessage/src/send-receipt-db.worker.ts index 1498bc77712c..fda45645437b 100644 --- a/extensions/imessage/src/send-receipt-db.worker.ts +++ b/extensions/imessage/src/send-receipt-db.worker.ts @@ -4,7 +4,7 @@ import { getNodeSqliteKysely, openNodeSqliteDatabase, type SqliteWorkerBackend, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import type { IMessageReceiptDbOperations } from "./send-receipt-db.js"; import { normalizeIMessageHandle } from "./targets.js"; diff --git a/extensions/logbook/src/store-frame-queries.ts b/extensions/logbook/src/store-frame-queries.ts index cb95600ec445..970576fcc9a1 100644 --- a/extensions/logbook/src/store-frame-queries.ts +++ b/extensions/logbook/src/store-frame-queries.ts @@ -2,7 +2,7 @@ import type { DatabaseSync } from "node:sqlite"; import { prepareSqliteQuerySync, type getNodeSqliteKysely, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { MAX_FRAMES_PER_CALL } from "./analyze.js"; import type { LogbookDatabase, toFrame } from "./store-schema.js"; diff --git a/extensions/logbook/src/store-prune.test.ts b/extensions/logbook/src/store-prune.test.ts index 790065eaf506..7daeae25c0f0 100644 --- a/extensions/logbook/src/store-prune.test.ts +++ b/extensions/logbook/src/store-prune.test.ts @@ -1,7 +1,7 @@ import * as fs from "node:fs"; import path from "node:path"; import { DatabaseSync } from "node:sqlite"; -import type { SqliteWorkerBackend } from "openclaw/plugin-sdk/sqlite-runtime"; +import type { SqliteWorkerBackend } from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { useAutoCleanupTempDirTracker } from "openclaw/plugin-sdk/test-env"; import { afterEach, describe, expect, it, vi } from "vitest"; import type { LogbookOperations } from "./store-contract.js"; @@ -12,8 +12,8 @@ vi.mock("node:fs", async (importOriginal) => { const actual = await importOriginal(); return { ...actual, rmSync: vi.fn(actual.rmSync) }; }); -vi.mock("openclaw/plugin-sdk/sqlite-runtime", async (importOriginal) => { - const actual = await importOriginal(); +vi.mock("openclaw/plugin-sdk/sqlite-worker-runtime", async (importOriginal) => { + const actual = await importOriginal(); return { ...actual, openNodeSqliteDatabase: (...args: Parameters) => { diff --git a/extensions/logbook/src/store-read-budgets.test.ts b/extensions/logbook/src/store-read-budgets.test.ts index 91dd9365dea1..0306f3de58b1 100644 --- a/extensions/logbook/src/store-read-budgets.test.ts +++ b/extensions/logbook/src/store-read-budgets.test.ts @@ -1,6 +1,6 @@ import path from "node:path"; import { DatabaseSync } from "node:sqlite"; -import type { SqliteWorkerBackend } from "openclaw/plugin-sdk/sqlite-runtime"; +import type { SqliteWorkerBackend } from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { useAutoCleanupTempDirTracker } from "openclaw/plugin-sdk/test-env"; import { afterEach, describe, expect, it, vi } from "vitest"; import type { LogbookOperations } from "./store-contract.js"; @@ -14,8 +14,8 @@ const reads = vi.hoisted(() => ({ frameTextBytes: 0, })); const preparations = vi.hoisted(() => new Map()); -vi.mock("openclaw/plugin-sdk/sqlite-runtime", async (importOriginal) => { - const actual = await importOriginal(); +vi.mock("openclaw/plugin-sdk/sqlite-worker-runtime", async (importOriginal) => { + const actual = await importOriginal(); return { ...actual, openNodeSqliteDatabase: (...args: Parameters) => { diff --git a/extensions/logbook/src/store-schema.ts b/extensions/logbook/src/store-schema.ts index e24884784fab..e2cb67b31e3b 100644 --- a/extensions/logbook/src/store-schema.ts +++ b/extensions/logbook/src/store-schema.ts @@ -1,4 +1,4 @@ -import type { Generated, Selectable } from "openclaw/plugin-sdk/sqlite-runtime"; +import type { Generated, Selectable } from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { asOptionalObjectRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; import type { LogbookBatch, diff --git a/extensions/logbook/src/store.worker.ts b/extensions/logbook/src/store.worker.ts index b6e466719fa4..a3d9fb756327 100644 --- a/extensions/logbook/src/store.worker.ts +++ b/extensions/logbook/src/store.worker.ts @@ -14,7 +14,7 @@ import { prepareSqliteQuerySync, runSqliteImmediateTransactionSync, type SqliteWorkerBackend, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { pickKeyframeId } from "./analyze.js"; import type { LogbookBatchInput, diff --git a/extensions/memory-core/src/memory-session-tombstones.ts b/extensions/memory-core/src/memory-session-tombstones.ts index 3678b3b995cd..cb6d36858d7c 100644 --- a/extensions/memory-core/src/memory-session-tombstones.ts +++ b/extensions/memory-core/src/memory-session-tombstones.ts @@ -3,7 +3,7 @@ import { executeSqliteQuerySync, getNodeSqliteKysely, tableExists, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; type TombstoneDatabase = { memory_session_tombstones: { session_id: string; agent_id: string }; diff --git a/extensions/memory-core/src/memory/manager-chunk-writer.ts b/extensions/memory-core/src/memory/manager-chunk-writer.ts index ae28db553656..952b10b82b3c 100644 --- a/extensions/memory-core/src/memory/manager-chunk-writer.ts +++ b/extensions/memory-core/src/memory/manager-chunk-writer.ts @@ -3,11 +3,11 @@ import type { MemoryChunk, MemoryEntryProvenance, MemorySource, -} from "openclaw/plugin-sdk/memory-core-host-engine-storage"; +} from "openclaw/plugin-sdk/memory-core-host-engine-indexing"; import { compileSqliteQueryBindings, getNodeSqliteKysely, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; export type IndexedMemoryChunk = MemoryChunk & { importance: number | null; diff --git a/extensions/memory-core/src/memory/manager-database-context.ts b/extensions/memory-core/src/memory/manager-database-context.ts index 75fe265ad76a..03a3726feb24 100644 --- a/extensions/memory-core/src/memory/manager-database-context.ts +++ b/extensions/memory-core/src/memory/manager-database-context.ts @@ -9,6 +9,7 @@ import { } from "openclaw/plugin-sdk/memory-core-host-engine-foundation"; import { resolveRuntimeWorkerUrl } from "openclaw/plugin-sdk/process-runtime"; import { + borrowOpenClawAgentDatabase, openSqliteWorkerStore, openOpenClawAgentSqliteWorkerStore, runSqliteWorkerStoreWrite, @@ -19,10 +20,11 @@ import { type StoreWriterQueue, } from "openclaw/plugin-sdk/sqlite-runtime"; import { memoryCpuProcessEntrypoints } from "./manager-cpu-entrypoints.js"; +import { MemoryIndexRevisionConflictError } from "./manager-db-kernel.js"; import { - MemoryIndexRevisionConflictError, closeMemoryDatabase, openMemoryDatabaseAtPath, + openMemoryDatabaseReadOnlyAtPath, } from "./manager-db.js"; import type { MemoryPublicationConnection, @@ -53,6 +55,32 @@ export class MemoryIndexDatabase { private releaseInProgress = false; shadowReleased = false; + static openPublished(params: { + agentId: string; + writeOptions: Parameters[0] & { path: string }; + readOnly: boolean; + allowExtension: boolean; + maintenanceSource?: MemoryIndexDatabase; + }): MemoryIndexDatabase { + const connection = params.readOnly + ? openMemoryDatabaseReadOnlyAtPath( + params.writeOptions.path, + params.allowExtension, + params.agentId, + ) + : borrowOpenClawAgentDatabase(params.writeOptions); + if (params.maintenanceSource && connection.db !== params.maintenanceSource.db) { + connection.release(); + throw new Error("Memory maintenance source connection changed"); + } + return new MemoryIndexDatabase( + connection.db, + connection.release, + params.readOnly, + params.writeOptions, + ); + } + static openShadow(filename: string, allowExtension: boolean): MemoryIndexDatabase { let database: MemoryIndexDatabase | undefined; const db = openMemoryDatabaseAtPath(filename, allowExtension, (operation) => diff --git a/extensions/memory-core/src/memory/manager-db-kernel.ts b/extensions/memory-core/src/memory/manager-db-kernel.ts new file mode 100644 index 000000000000..fdad71b33879 --- /dev/null +++ b/extensions/memory-core/src/memory/manager-db-kernel.ts @@ -0,0 +1,191 @@ +import type { DatabaseSync } from "node:sqlite"; +import { + dropMemoryPathFtsTriggers, + ensureMemoryChunkProvenance, + ensureMemoryRecallMetadataSchema, + ensureMemoryPathFtsTriggers, + MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE, + MEMORY_INDEX_PATHS_FTS_TABLE, +} from "openclaw/plugin-sdk/memory-core-host-engine-schema"; +import { runSqliteImmediateTransactionSync } from "openclaw/plugin-sdk/sqlite-worker-runtime"; +import { markMemoryVectorIndexClean } from "./manager-vector-rebuild-state.js"; + +const MEMORY_REINDEX_SCHEMA = "memory_reindex"; +export const MEMORY_INDEX_STATE_ID = 1; + +function tableExists(db: DatabaseSync, schema: string, tableName: string): boolean { + const row = db + .prepare(`SELECT 1 AS ok FROM ${schema}.sqlite_master WHERE type = 'table' AND name = ?`) + .get(tableName); + return row?.ok === 1; +} + +export { tableExists as memoryDatabaseTableExists }; + +function readTableSql(db: DatabaseSync, schema: string, tableName: string): string | null { + const row = db + .prepare(`SELECT sql FROM ${schema}.sqlite_master WHERE type = 'table' AND name = ?`) + .get(tableName); + return typeof row?.sql === "string" && row.sql.trim() ? row.sql : null; +} + +export function readMemoryDatabaseRevision(db: DatabaseSync): number { + const row = db + .prepare("SELECT revision FROM memory_index_state WHERE id = ?") + .get(MEMORY_INDEX_STATE_ID); + if (typeof row?.revision !== "number" || !Number.isSafeInteger(row.revision)) { + throw new Error("Memory index revision is missing or invalid"); + } + return row.revision; +} + +export class MemoryIndexRevisionConflictError extends Error { + override name = "MemoryIndexRevisionConflictError"; +} + +function replaceVirtualTable(params: { + db: DatabaseSync; + tableName: "memory_index_chunks_fts" | "memory_index_chunks_vec"; + columns: string; + ignoreDropErrorWhenSourceMissing?: boolean; +}): void { + const { db, tableName, columns } = params; + const createSql = readTableSql(db, MEMORY_REINDEX_SCHEMA, tableName); + if (!createSql) { + try { + db.exec(`DROP TABLE IF EXISTS main.${tableName}`); + } catch (err) { + if (!params.ignoreDropErrorWhenSourceMissing) { + throw err; + } + } + return; + } + db.exec(`DROP TABLE IF EXISTS main.${tableName}`); + db.exec(createSql); + db.exec( + `INSERT INTO main.${tableName} (${columns}) ` + + `SELECT ${columns} FROM ${MEMORY_REINDEX_SCHEMA}.${tableName}`, + ); +} + +function replaceMemoryPathFtsTable(db: DatabaseSync): void { + const createSql = readTableSql(db, MEMORY_REINDEX_SCHEMA, MEMORY_INDEX_PATHS_FTS_TABLE); + db.exec(`DROP TABLE IF EXISTS main.${MEMORY_INDEX_PATHS_FTS_TABLE}`); + if (!createSql) { + return; + } + db.exec(createSql); + // Bulk publication already suspends row triggers. Rebuild from the copied + // stable source ids so later singleton deletes remain direct rowid lookups. + db.exec( + `INSERT INTO main.${MEMORY_INDEX_PATHS_FTS_TABLE} (rowid, path, source) ` + + `SELECT id, path, source FROM main.memory_index_sources`, + ); +} + +/** The native publication owner receives prepared connection and source facts. */ +type MemoryDatabasePublication = { + targetDb: DatabaseSync; + sourcePath: string; + metaKey: string; + expectedRevision: number; + onBegin?: () => void; + withCommit?: (commit: () => void) => void; + vectorIndexComplete?: boolean; +}; + +/** The admitted connection owns ATTACH, atomic replacement, COMMIT and DETACH. */ +export function publishMemoryDatabaseTables(params: MemoryDatabasePublication): void { + ensureMemoryRecallMetadataSchema(params.targetDb); + // Existing pre-provenance databases need this before the publication writes it. + ensureMemoryChunkProvenance(params.targetDb); + // Admission precedes ATTACH; no shadow attachment or transaction crosses an await. + params.targetDb.prepare(`ATTACH DATABASE ? AS ${MEMORY_REINDEX_SCHEMA}`).run(params.sourcePath); + try { + runSqliteImmediateTransactionSync( + params.targetDb, + () => { + params.onBegin?.(); + const liveRevision = readMemoryDatabaseRevision(params.targetDb); + if (liveRevision !== params.expectedRevision) { + throw new MemoryIndexRevisionConflictError( + `Memory index changed while full reindex was building ` + + `(expected revision ${params.expectedRevision}, found ${liveRevision}); retry the full reindex.`, + ); + } + const publishesPathFts = tableExists( + params.targetDb, + MEMORY_REINDEX_SCHEMA, + MEMORY_INDEX_PATHS_FTS_TABLE, + ); + // Bulk source replacement must not fire one FTS5 scan per old row. + // Restore the schema-owned triggers only after the derived table is replaced. + dropMemoryPathFtsTriggers(params.targetDb); + params.targetDb + .prepare("DELETE FROM main.memory_index_meta WHERE key = ?") + .run(params.metaKey); + params.targetDb + .prepare( + `INSERT INTO main.memory_index_meta (key, value) + SELECT key, value FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_meta WHERE key = ?`, + ) + .run(params.metaKey); + + params.targetDb.exec(` + DELETE FROM main.memory_index_sources; + INSERT INTO main.memory_index_sources (id, path, source, hash, mtime, size) + SELECT id, path, source, hash, mtime, size + FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_sources; + + DELETE FROM main.memory_index_chunks; + INSERT INTO main.memory_index_chunks ( + id, path, source, start_line, end_line, hash, model, text, embedding, updated_at + ) + SELECT + id, path, source, start_line, end_line, hash, model, text, embedding, updated_at + FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_chunks; + + DELETE FROM main.${MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE}; + INSERT INTO main.${MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE} ( + chunk_id, importance, triggers, project_key + ) + SELECT chunk_id, importance, triggers, project_key + FROM ${MEMORY_REINDEX_SCHEMA}.${MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE}; + + DELETE FROM main.memory_index_chunk_provenance; + INSERT INTO main.memory_index_chunk_provenance ( + chunk_id, origin_class, session_kind, observed_at, supersedes_key + ) + SELECT chunk_id, origin_class, session_kind, observed_at, supersedes_key + FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_chunk_provenance; + `); + + replaceVirtualTable({ + db: params.targetDb, + tableName: "memory_index_chunks_fts", + columns: "text, id, path, source, model, start_line, end_line", + }); + replaceMemoryPathFtsTable(params.targetDb); + if (publishesPathFts) { + ensureMemoryPathFtsTriggers(params.targetDb); + } + replaceVirtualTable({ + db: params.targetDb, + tableName: "memory_index_chunks_vec", + columns: "id, embedding", + // A vector-disabled connection may not have sqlite-vec loaded and cannot + // drop an old virtual table. Missing vector metadata forces a strict + // rebuild before that table can be queried again. + ignoreDropErrorWhenSourceMissing: true, + }); + if (params.vectorIndexComplete) { + markMemoryVectorIndexClean(params.targetDb); + } + }, + { withCommit: params.withCommit }, + ); + } finally { + params.targetDb.exec(`DETACH DATABASE ${MEMORY_REINDEX_SCHEMA}`); + } +} diff --git a/extensions/memory-core/src/memory/manager-db.test.ts b/extensions/memory-core/src/memory/manager-db.test.ts index 57b5930e853d..0604ea51b49d 100644 --- a/extensions/memory-core/src/memory/manager-db.test.ts +++ b/extensions/memory-core/src/memory/manager-db.test.ts @@ -13,12 +13,14 @@ import { resetMemoryCoreDreamingStateForTests, } from "../test-helpers.js"; import { - cleanupAgedMemoryReindexTempFiles, - closeMemoryDatabase, - openMemoryDatabaseAtPath, publishMemoryDatabaseTables, readMemoryDatabaseRevision, MemoryIndexRevisionConflictError, +} from "./manager-db-kernel.js"; +import { + cleanupAgedMemoryReindexTempFiles, + closeMemoryDatabase, + openMemoryDatabaseAtPath, resetMemoryDatabase, } from "./manager-db.js"; import { waitForMemoryReindexLock } from "./manager-reindex-lock.js"; diff --git a/extensions/memory-core/src/memory/manager-db.ts b/extensions/memory-core/src/memory/manager-db.ts index 966c78807cea..daaf94eb2b18 100644 --- a/extensions/memory-core/src/memory/manager-db.ts +++ b/extensions/memory-core/src/memory/manager-db.ts @@ -6,14 +6,8 @@ import type { DatabaseSync } from "node:sqlite"; import { closeMemorySqliteWalMaintenance, configureMemorySqliteWalMaintenance, - dropMemoryPathFtsTriggers, - ensureMemoryChunkProvenance, ensureMemoryIndexSchema, - ensureMemoryRecallMetadataSchema, - ensureMemoryPathFtsTriggers, loadSqliteVecExtension, - MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE, - MEMORY_INDEX_PATHS_FTS_TABLE, MEMORY_INDEX_DERIVED_TABLES, MEMORY_INDEX_STATE_TABLE, MEMORY_INDEX_VECTOR_TABLE, @@ -24,12 +18,14 @@ import { runSqliteImmediateTransactionSync, } from "openclaw/plugin-sdk/sqlite-runtime"; import { withMemoryWorkspaceLock } from "../memory-workspace-lock.js"; +import { + MEMORY_INDEX_STATE_ID, + memoryDatabaseTableExists as tableExists, + readMemoryDatabaseRevision, +} from "./manager-db-kernel.js"; import { withMemoryIndexPublishGeneration } from "./manager-index-generation-lease.js"; import { waitForMemoryReindexLock } from "./manager-reindex-lock.js"; -import { markMemoryVectorIndexClean } from "./manager-vector-rebuild-state.js"; -const MEMORY_REINDEX_SCHEMA = "memory_reindex"; -const MEMORY_INDEX_STATE_ID = 1; const MEMORY_DATABASE_FILE_SUFFIXES = ["", "-wal", "-shm", "-journal"] as const; const MEMORY_REINDEX_ENTRY_SUFFIXES = ["-wal", "-shm", "-journal", ""] as const; const MEMORY_REINDEX_UUID_PATTERN = @@ -64,22 +60,6 @@ async function isRegularFile(filePath: string): Promise { } } -function tableExists(db: DatabaseSync, schema: string, tableName: string): boolean { - const row = db - .prepare(`SELECT 1 AS ok FROM ${schema}.sqlite_master WHERE type = 'table' AND name = ?`) - .get(tableName) as { ok?: unknown } | undefined; - return row?.ok === 1; -} - -export { tableExists as memoryDatabaseTableExists }; - -function readTableSql(db: DatabaseSync, schema: string, tableName: string): string | null { - const row = db - .prepare(`SELECT sql FROM ${schema}.sqlite_master WHERE type = 'table' AND name = ?`) - .get(tableName) as { sql?: unknown } | undefined; - return typeof row?.sql === "string" && row.sql.trim() ? row.sql : null; -} - function hasSqliteVecExtension(db: DatabaseSync): boolean { try { const row = db.prepare("SELECT vec_version() AS version").get() as @@ -91,20 +71,6 @@ function hasSqliteVecExtension(db: DatabaseSync): boolean { } } -export function readMemoryDatabaseRevision(db: DatabaseSync): number { - const row = db - .prepare("SELECT revision FROM memory_index_state WHERE id = ?") - .get(MEMORY_INDEX_STATE_ID) as { revision?: unknown } | undefined; - if (typeof row?.revision !== "number" || !Number.isSafeInteger(row.revision)) { - throw new Error("Memory index revision is missing or invalid"); - } - return row.revision; -} - -export class MemoryIndexRevisionConflictError extends Error { - override name = "MemoryIndexRevisionConflictError"; -} - /** Reset derived content without replacing the shared agent database or its schema. */ export async function resetMemoryDatabase(params: { targetDb: DatabaseSync; @@ -178,153 +144,6 @@ export async function resetMemoryDatabase(params: { } } -function replaceVirtualTable(params: { - db: DatabaseSync; - tableName: "memory_index_chunks_fts" | "memory_index_chunks_vec"; - columns: string; - ignoreDropErrorWhenSourceMissing?: boolean; -}): void { - const { db, tableName, columns } = params; - const createSql = readTableSql(db, MEMORY_REINDEX_SCHEMA, tableName); - if (!createSql) { - try { - db.exec(`DROP TABLE IF EXISTS main.${tableName}`); - } catch (err) { - if (!params.ignoreDropErrorWhenSourceMissing) { - throw err; - } - } - return; - } - db.exec(`DROP TABLE IF EXISTS main.${tableName}`); - db.exec(createSql); - db.exec( - `INSERT INTO main.${tableName} (${columns}) ` + - `SELECT ${columns} FROM ${MEMORY_REINDEX_SCHEMA}.${tableName}`, - ); -} - -function replaceMemoryPathFtsTable(db: DatabaseSync): void { - const createSql = readTableSql(db, MEMORY_REINDEX_SCHEMA, MEMORY_INDEX_PATHS_FTS_TABLE); - db.exec(`DROP TABLE IF EXISTS main.${MEMORY_INDEX_PATHS_FTS_TABLE}`); - if (!createSql) { - return; - } - db.exec(createSql); - // Bulk publication already suspends row triggers. Rebuild from the copied - // stable source ids so later singleton deletes remain direct rowid lookups. - db.exec( - `INSERT INTO main.${MEMORY_INDEX_PATHS_FTS_TABLE} (rowid, path, source) ` + - `SELECT id, path, source FROM main.memory_index_sources`, - ); -} - -/** The native publication owner receives prepared connection and source facts. */ -type MemoryDatabasePublication = { - targetDb: DatabaseSync; - sourcePath: string; - metaKey: string; - expectedRevision: number; - onBegin?: () => void; - withCommit?: (commit: () => void) => void; - vectorIndexComplete?: boolean; -}; - -/** The admitted connection owns ATTACH, atomic replacement, COMMIT and DETACH. */ -export function publishMemoryDatabaseTables(params: MemoryDatabasePublication): void { - ensureMemoryRecallMetadataSchema(params.targetDb); - // Existing pre-provenance databases need this before the publication writes it. - ensureMemoryChunkProvenance(params.targetDb); - // Admission precedes ATTACH; no shadow attachment or transaction crosses an await. - params.targetDb.prepare(`ATTACH DATABASE ? AS ${MEMORY_REINDEX_SCHEMA}`).run(params.sourcePath); - try { - runSqliteImmediateTransactionSync( - params.targetDb, - () => { - params.onBegin?.(); - const liveRevision = readMemoryDatabaseRevision(params.targetDb); - if (liveRevision !== params.expectedRevision) { - throw new MemoryIndexRevisionConflictError( - `Memory index changed while full reindex was building ` + - `(expected revision ${params.expectedRevision}, found ${liveRevision}); retry the full reindex.`, - ); - } - const publishesPathFts = tableExists( - params.targetDb, - MEMORY_REINDEX_SCHEMA, - MEMORY_INDEX_PATHS_FTS_TABLE, - ); - // Bulk source replacement must not fire one FTS5 scan per old row. - // Restore the schema-owned triggers only after the derived table is replaced. - dropMemoryPathFtsTriggers(params.targetDb); - params.targetDb - .prepare("DELETE FROM main.memory_index_meta WHERE key = ?") - .run(params.metaKey); - params.targetDb - .prepare( - `INSERT INTO main.memory_index_meta (key, value) - SELECT key, value FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_meta WHERE key = ?`, - ) - .run(params.metaKey); - - params.targetDb.exec(` - DELETE FROM main.memory_index_sources; - INSERT INTO main.memory_index_sources (id, path, source, hash, mtime, size) - SELECT id, path, source, hash, mtime, size - FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_sources; - - DELETE FROM main.memory_index_chunks; - INSERT INTO main.memory_index_chunks ( - id, path, source, start_line, end_line, hash, model, text, embedding, updated_at - ) - SELECT - id, path, source, start_line, end_line, hash, model, text, embedding, updated_at - FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_chunks; - - DELETE FROM main.${MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE}; - INSERT INTO main.${MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE} ( - chunk_id, importance, triggers, project_key - ) - SELECT chunk_id, importance, triggers, project_key - FROM ${MEMORY_REINDEX_SCHEMA}.${MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE}; - - DELETE FROM main.memory_index_chunk_provenance; - INSERT INTO main.memory_index_chunk_provenance ( - chunk_id, origin_class, session_kind, observed_at, supersedes_key - ) - SELECT chunk_id, origin_class, session_kind, observed_at, supersedes_key - FROM ${MEMORY_REINDEX_SCHEMA}.memory_index_chunk_provenance; - `); - - replaceVirtualTable({ - db: params.targetDb, - tableName: "memory_index_chunks_fts", - columns: "text, id, path, source, model, start_line, end_line", - }); - replaceMemoryPathFtsTable(params.targetDb); - if (publishesPathFts) { - ensureMemoryPathFtsTriggers(params.targetDb); - } - replaceVirtualTable({ - db: params.targetDb, - tableName: "memory_index_chunks_vec", - columns: "id, embedding", - // A vector-disabled connection may not have sqlite-vec loaded and cannot - // drop an old virtual table. Missing vector metadata forces a strict - // rebuild before that table can be queried again. - ignoreDropErrorWhenSourceMissing: true, - }); - if (params.vectorIndexComplete) { - markMemoryVectorIndexClean(params.targetDb); - } - }, - { withCommit: params.withCommit }, - ); - } finally { - params.targetDb.exec(`DETACH DATABASE ${MEMORY_REINDEX_SCHEMA}`); - } -} - /** Remove one closed shadow memory database and its journal-mode sidecars. */ export async function removeMemoryDatabaseFiles(dbPath: string): Promise { for (const suffix of MEMORY_DATABASE_FILE_SUFFIXES) { diff --git a/extensions/memory-core/src/memory/manager-embedding-ops.ts b/extensions/memory-core/src/memory/manager-embedding-ops.ts index 09131fd63620..3f18f43bc523 100644 --- a/extensions/memory-core/src/memory/manager-embedding-ops.ts +++ b/extensions/memory-core/src/memory/manager-embedding-ops.ts @@ -37,7 +37,7 @@ import { readSessionResetRecallCutoffMetadata } from "../session-reset-recall-me import type { EmbeddingProvider } from "./embeddings.js"; import type { IndexedMemoryChunk } from "./manager-chunk-writer.js"; import { prepareMemoryIndexInWorker } from "./manager-cpu-worker-runtime.js"; -import { readMemoryDatabaseRevision } from "./manager-db.js"; +import { readMemoryDatabaseRevision } from "./manager-db-kernel.js"; import { clearMemoryEmbeddingCacheIdentities, collectMemoryCachedEmbeddings, diff --git a/extensions/memory-core/src/memory/manager-publication-transfer.test.ts b/extensions/memory-core/src/memory/manager-publication-transfer.test.ts index ac92b1420e8c..2c4d05f5a705 100644 --- a/extensions/memory-core/src/memory/manager-publication-transfer.test.ts +++ b/extensions/memory-core/src/memory/manager-publication-transfer.test.ts @@ -1,12 +1,12 @@ import path from "node:path"; import { serialize } from "node:v8"; -import * as sqliteCapabilities from "openclaw/plugin-sdk/memory-core-host-engine-knn"; import { ensureMemoryIndexSchema } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; import * as sqliteRuntime from "openclaw/plugin-sdk/sqlite-runtime"; +import * as sqliteWorkerRuntime from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { useAutoCleanupTempDirTracker } from "openclaw/plugin-sdk/test-env"; import { afterEach, describe, expect, it, vi } from "vitest"; import { MemoryIndexDatabase } from "./manager-database-context.js"; -import { readMemoryDatabaseRevision } from "./manager-db.js"; +import { readMemoryDatabaseRevision } from "./manager-db-kernel.js"; import { memoryPublicationBatches } from "./manager-publication-transfer.js"; import { openExistingSqliteWorkerBackend } from "./manager-publication.worker.js"; import { readMemoryShadowIdentity } from "./manager-shadow-task.js"; @@ -100,14 +100,16 @@ describe("bounded memory publication transfer", () => { it("opens keyword publication when SQLite extension loading is unavailable", () => { const owner = createOwner(); - vi.spyOn(sqliteCapabilities, "supportsNodeSqliteExtensionLoading").mockReturnValue(false); - const open = sqliteRuntime.openNodeSqliteDatabase; - vi.spyOn(sqliteRuntime, "openNodeSqliteDatabase").mockImplementation((location, options) => { - if (options?.allowExtension) { - throw new Error("SQLite extension loading is unavailable"); - } - return open(location, options); - }); + vi.spyOn(sqliteWorkerRuntime, "supportsNodeSqliteExtensionLoading").mockReturnValue(false); + const open = sqliteWorkerRuntime.openNodeSqliteDatabase; + vi.spyOn(sqliteWorkerRuntime, "openNodeSqliteDatabase").mockImplementation( + (location, options) => { + if (options?.allowExtension) { + throw new Error("SQLite extension loading is unavailable"); + } + return open(location, options); + }, + ); const backend = createBackend(owner); const { chunks, embeddings: _embeddings, ...header } = replacement(); backend.execute({ diff --git a/extensions/memory-core/src/memory/manager-publication.worker.ts b/extensions/memory-core/src/memory/manager-publication.worker.ts index f122056898e8..55d76d43f94e 100644 --- a/extensions/memory-core/src/memory/manager-publication.worker.ts +++ b/extensions/memory-core/src/memory/manager-publication.worker.ts @@ -1,15 +1,15 @@ import type { DatabaseSync } from "node:sqlite"; -import { supportsNodeSqliteExtensionLoading } from "openclaw/plugin-sdk/memory-core-host-engine-knn"; import { assertTransactionUsable, openNodeSqliteDatabase, resolveExistingSqliteFileUri, requestSqliteWorkerOperationAdmission, runSqliteImmediateTransactionSync, + supportsNodeSqliteExtensionLoading, type SqliteWorkerBackend, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { hasMemorySessionTombstone } from "../memory-session-tombstones.js"; -import { publishMemoryDatabaseTables, readMemoryDatabaseRevision } from "./manager-db.js"; +import { publishMemoryDatabaseTables, readMemoryDatabaseRevision } from "./manager-db-kernel.js"; import type { MemoryPublicationConnection, MemoryPublicationOperations, diff --git a/extensions/memory-core/src/memory/manager-search-maintenance.ts b/extensions/memory-core/src/memory/manager-search-maintenance.ts index 49bcbe3d40f1..c4ecd353aa7e 100644 --- a/extensions/memory-core/src/memory/manager-search-maintenance.ts +++ b/extensions/memory-core/src/memory/manager-search-maintenance.ts @@ -1,6 +1,6 @@ // Memory Core owns detached search-time index maintenance lifecycle. import { toErrorObject } from "openclaw/plugin-sdk/error-runtime"; -import { MemoryIndexRevisionConflictError } from "./manager-db.js"; +import { MemoryIndexRevisionConflictError } from "./manager-db-kernel.js"; type MemorySearchMaintenanceManager = { adoptReindexRetryState(generation: DirtyGeneration): void; diff --git a/extensions/memory-core/src/memory/manager-search-orchestration.test.ts b/extensions/memory-core/src/memory/manager-search-orchestration.test.ts index 1a572b79805a..f0c58a40f4d2 100644 --- a/extensions/memory-core/src/memory/manager-search-orchestration.test.ts +++ b/extensions/memory-core/src/memory/manager-search-orchestration.test.ts @@ -11,7 +11,7 @@ import { forgetMemoryEntries } from "../memory-forget.js"; import type { EmbeddingProvider } from "./embeddings.js"; import { memoryCpuProcessEntrypoints } from "./manager-cpu-entrypoints.js"; import * as memoryCpuWorkerRuntime from "./manager-cpu-worker-runtime.js"; -import { MemoryIndexRevisionConflictError } from "./manager-db.js"; +import { MemoryIndexRevisionConflictError } from "./manager-db-kernel.js"; import { createManagerIndexFixture } from "./manager-index.test-support.js"; const { closeAllMemorySearchManagers, getMemorySearchManager } = await import("./index.js"); diff --git a/extensions/memory-core/src/memory/manager-source-index-kernel.ts b/extensions/memory-core/src/memory/manager-source-index-kernel.ts index 59347f550898..b67136277940 100644 --- a/extensions/memory-core/src/memory/manager-source-index-kernel.ts +++ b/extensions/memory-core/src/memory/manager-source-index-kernel.ts @@ -1,16 +1,15 @@ import type { DatabaseSync, StatementSync } from "node:sqlite"; +import { hashText, type MemorySource } from "openclaw/plugin-sdk/memory-core-host-engine-indexing"; import { - hashText, MEMORY_INDEX_FTS_TABLE, MEMORY_INDEX_VECTOR_TABLE, - type MemorySource, -} from "openclaw/plugin-sdk/memory-core-host-engine-storage"; +} from "openclaw/plugin-sdk/memory-core-host-engine-schema"; import { executeSqliteQuerySync, executeSqliteQueryTakeFirstSync, getNodeSqliteKysely, runSqliteImmediateTransactionSync, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { createMemoryChunkWriter, type IndexedMemoryChunk } from "./manager-chunk-writer.js"; import { markMemoryVectorRebuildRequired, diff --git a/extensions/memory-core/src/memory/manager-source-sync-ops.ts b/extensions/memory-core/src/memory/manager-source-sync-ops.ts index b9d714618281..bf80559b7661 100644 --- a/extensions/memory-core/src/memory/manager-source-sync-ops.ts +++ b/extensions/memory-core/src/memory/manager-source-sync-ops.ts @@ -11,7 +11,7 @@ import { } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; import { runSqliteImmediateTransaction } from "openclaw/plugin-sdk/sqlite-runtime"; import { withMemoryWorkspaceLock } from "../memory-workspace-lock.js"; -import { MemoryIndexRevisionConflictError } from "./manager-db.js"; +import { MemoryIndexRevisionConflictError } from "./manager-db-kernel.js"; import type { MemoryIndexEntry } from "./manager-index-preparation.js"; import { MemoryManagerSessionSyncOps } from "./manager-session-sync-ops.js"; import { diff --git a/extensions/memory-core/src/memory/manager-sync-ops.ts b/extensions/memory-core/src/memory/manager-sync-ops.ts index b21dac98a69e..b1fa7c1b1c00 100644 --- a/extensions/memory-core/src/memory/manager-sync-ops.ts +++ b/extensions/memory-core/src/memory/manager-sync-ops.ts @@ -21,12 +21,8 @@ import { type EmbeddingProviderRuntime, } from "./embeddings.js"; import { MemoryIndexDatabase } from "./manager-database-context.js"; -import { - cleanupAgedMemoryReindexTempFiles, - memoryDatabaseTableExists, - readMemoryDatabaseRevision, - removeMemoryDatabaseFiles, -} from "./manager-db.js"; +import { memoryDatabaseTableExists, readMemoryDatabaseRevision } from "./manager-db-kernel.js"; +import { cleanupAgedMemoryReindexTempFiles, removeMemoryDatabaseFiles } from "./manager-db.js"; import { isMemoryEmbeddingOperationError } from "./manager-embedding-errors.js"; import { withMemoryIndexPublishGeneration } from "./manager-index-generation-lease.js"; import { diff --git a/extensions/memory-core/src/memory/manager-vector-rebuild-state.ts b/extensions/memory-core/src/memory/manager-vector-rebuild-state.ts index cdb095a4bb74..289e0609f8d1 100644 --- a/extensions/memory-core/src/memory/manager-vector-rebuild-state.ts +++ b/extensions/memory-core/src/memory/manager-vector-rebuild-state.ts @@ -1,9 +1,7 @@ // Memory Core plugin module owns persisted vector completeness state. import type { DatabaseSync } from "node:sqlite"; -import { - MEMORY_INDEX_META_TABLE, - type MemoryVectorIndexState, -} from "openclaw/plugin-sdk/memory-core-host-engine-storage"; +import { MEMORY_INDEX_META_TABLE } from "openclaw/plugin-sdk/memory-core-host-engine-schema"; +import type { MemoryVectorIndexState } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; const VECTOR_REBUILD_META_KEY = "memory_vector_rebuild_v1"; diff --git a/extensions/memory-core/src/memory/manager.borrowed-connection.test.ts b/extensions/memory-core/src/memory/manager.borrowed-connection.test.ts index bd9e5b35fd92..d1096f058a01 100644 --- a/extensions/memory-core/src/memory/manager.borrowed-connection.test.ts +++ b/extensions/memory-core/src/memory/manager.borrowed-connection.test.ts @@ -19,7 +19,7 @@ import * as sqliteRuntime from "openclaw/plugin-sdk/sqlite-runtime"; import { closeOpenClawAgentDatabasesForTest } from "openclaw/plugin-sdk/sqlite-runtime-testing"; import { afterEach, describe, expect, it, vi } from "vitest"; import { recordMemorySessionTombstones } from "../memory-entry-origins.js"; -import { MemoryIndexRevisionConflictError } from "./manager-db.js"; +import { MemoryIndexRevisionConflictError } from "./manager-db-kernel.js"; import { createManagerIndexFixture } from "./manager-index.test-support.js"; import { closeAllMemoryIndexManagers } from "./manager-runtime.js"; import { MemoryIndexManager } from "./manager.js"; diff --git a/extensions/memory-core/src/memory/manager.ts b/extensions/memory-core/src/memory/manager.ts index a81536c46ed0..af39f471ab96 100644 --- a/extensions/memory-core/src/memory/manager.ts +++ b/extensions/memory-core/src/memory/manager.ts @@ -19,16 +19,13 @@ import { } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; import { normalizeAgentId } from "openclaw/plugin-sdk/routing"; import { createPluginRuntimeStore } from "openclaw/plugin-sdk/runtime-store"; -import { - borrowOpenClawAgentDatabase, - withOpenClawAgentDatabaseWrite, -} from "openclaw/plugin-sdk/sqlite-runtime"; +import { withOpenClawAgentDatabaseWrite } from "openclaw/plugin-sdk/sqlite-runtime"; import { runInMemoryBackgroundContext } from "./background-context.js"; import type { MemoryCoreAcquireLocalService } from "./embedding-local-service.js"; import type { EmbeddingProvider, EmbeddingProviderRequest } from "./embeddings.js"; import { getMemoryManagerLifecycle } from "./lifecycle.js"; import { MemoryIndexDatabase } from "./manager-database-context.js"; -import { memoryDatabaseTableExists, openMemoryDatabaseReadOnlyAtPath } from "./manager-db.js"; +import { memoryDatabaseTableExists } from "./manager-db-kernel.js"; import { resolveEffectiveMemorySearchSettings, resolveMemoryEmbeddingProviderRequirement, @@ -249,24 +246,17 @@ export class MemoryIndexManager extends MemorySearchOrchestration implements Mem for (const memorySource of effectiveSettings.sources) { this.sources.add(memorySource); } - const vectorEnabled = effectiveSettings.store.vector.enabled; const readOnly = this.purpose === "status"; if (source && (!source.publishedDatabase.db.isOpen || this.purpose !== "maintenance")) { throw new Error("Memory maintenance source connection is unavailable"); } - const connection = readOnly - ? openMemoryDatabaseReadOnlyAtPath(dbPath, vectorEnabled, this.agentId) - : borrowOpenClawAgentDatabase(params.databaseOptions); - if (source && connection.db !== source.publishedDatabase.db) { - connection.release(); - throw new Error("Memory maintenance source connection changed"); - } - this.publishedDatabase = new MemoryIndexDatabase( - connection.db, - connection.release, + this.publishedDatabase = MemoryIndexDatabase.openPublished({ + agentId: this.agentId, + writeOptions: params.databaseOptions, readOnly, - params.databaseOptions, - ); + allowExtension: effectiveSettings.store.vector.enabled, + maintenanceSource: source?.publishedDatabase, + }); try { this.providerKey = this.computeProviderKey(); this.cache = { diff --git a/extensions/memory-core/src/tools.index-diagnostic.test.ts b/extensions/memory-core/src/tools.index-diagnostic.test.ts index dbdbb850cedc..8b342b21627f 100644 --- a/extensions/memory-core/src/tools.index-diagnostic.test.ts +++ b/extensions/memory-core/src/tools.index-diagnostic.test.ts @@ -1,6 +1,6 @@ import { openOpenClawAgentDatabase } from "openclaw/plugin-sdk/sqlite-runtime"; import { describe, expect, it, vi } from "vitest"; -import { readMemoryDatabaseRevision } from "./memory/manager-db.js"; +import { readMemoryDatabaseRevision } from "./memory/manager-db-kernel.js"; import { createManagerIndexFixture } from "./memory/manager-index.test-support.js"; import { MEMORY_INDEX_PROVENANCE_VERSION } from "./memory/manager-reindex-state.js"; import { createMemorySearchTool, testing } from "./tools.js"; diff --git a/extensions/memory-core/src/tools.index-upgrade.test.ts b/extensions/memory-core/src/tools.index-upgrade.test.ts index 4dd8bf764ff1..2ad16269cd0b 100644 --- a/extensions/memory-core/src/tools.index-upgrade.test.ts +++ b/extensions/memory-core/src/tools.index-upgrade.test.ts @@ -2,7 +2,7 @@ import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; import { MEMORY_CHUNKING_VERSION } from "openclaw/plugin-sdk/memory-core-host-engine-storage"; import { openOpenClawAgentDatabase } from "openclaw/plugin-sdk/sqlite-runtime"; import { beforeEach, describe, expect, it, vi } from "vitest"; -import { readMemoryDatabaseRevision } from "./memory/manager-db.js"; +import { readMemoryDatabaseRevision } from "./memory/manager-db-kernel.js"; import { createManagerIndexFixture } from "./memory/manager-index.test-support.js"; import { MEMORY_INDEX_PROVENANCE_VERSION } from "./memory/manager-reindex-state.js"; import { createMemorySearchTool, testing } from "./tools.js"; diff --git a/extensions/memory-core/src/tools.real-manager.test.ts b/extensions/memory-core/src/tools.real-manager.test.ts index c206b5407b77..6eea03009e70 100644 --- a/extensions/memory-core/src/tools.real-manager.test.ts +++ b/extensions/memory-core/src/tools.real-manager.test.ts @@ -12,7 +12,7 @@ import { openOpenClawAgentDatabase } from "openclaw/plugin-sdk/sqlite-runtime"; import { closeOpenClawAgentDatabasesForTest } from "openclaw/plugin-sdk/sqlite-runtime-testing"; import { beforeEach, describe, expect, it, vi } from "vitest"; import type { EmbeddingProvider } from "./memory/embeddings.js"; -import { readMemoryDatabaseRevision } from "./memory/manager-db.js"; +import { readMemoryDatabaseRevision } from "./memory/manager-db-kernel.js"; import * as generationLease from "./memory/manager-index-generation-lease.js"; import { createManagerIndexFixture, diff --git a/extensions/qa-lab/src/execution-identity-storage-inspection.worker.ts b/extensions/qa-lab/src/execution-identity-storage-inspection.worker.ts index 44752057ffcc..f381548a2447 100644 --- a/extensions/qa-lab/src/execution-identity-storage-inspection.worker.ts +++ b/extensions/qa-lab/src/execution-identity-storage-inspection.worker.ts @@ -4,7 +4,7 @@ import { openNodeSqliteDatabase, tableExists, type SqliteWorkerBackend, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import type { QaExecutionIdentityStorageOperations } from "./execution-identity-storage-inspection.js"; type QaExecutionIdentityDatabase = { diff --git a/extensions/team-reports/src/store.worker.ts b/extensions/team-reports/src/store.worker.ts index e80d10be43f5..8fe918cb85d2 100644 --- a/extensions/team-reports/src/store.worker.ts +++ b/extensions/team-reports/src/store.worker.ts @@ -13,7 +13,7 @@ import { openNodeSqliteDatabase, runSqliteImmediateTransactionSync, type SqliteWorkerCommand, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { MAX_REPORT_BYTES } from "./limits.js"; import { DAY_MS } from "./periods.js"; import type { diff --git a/extensions/tsconfig.package-boundary.paths.json b/extensions/tsconfig.package-boundary.paths.json index 72e9f636c59c..117f839dad74 100644 --- a/extensions/tsconfig.package-boundary.paths.json +++ b/extensions/tsconfig.package-boundary.paths.json @@ -933,6 +933,9 @@ ], "openclaw/plugin-sdk/memory-core-host-engine-indexing": [ "../packages/plugin-sdk/dist/src/plugin-sdk/memory-core-host-engine-indexing.d.ts" + ], + "openclaw/plugin-sdk/sqlite-worker-runtime": [ + "../packages/plugin-sdk/dist/src/plugin-sdk/sqlite-worker-runtime.d.ts" ] } } diff --git a/extensions/workboard/src/sqlite-store-batch-read.test.ts b/extensions/workboard/src/sqlite-store-batch-read.test.ts index 8d2048c11b1f..fcc0befb22e9 100644 --- a/extensions/workboard/src/sqlite-store-batch-read.test.ts +++ b/extensions/workboard/src/sqlite-store-batch-read.test.ts @@ -19,8 +19,8 @@ const workerModuleUrl = resolveRuntimeWorkerUrl(workboardSqliteBackendEntrypoint const sqliteStatements = vi.hoisted(() => ({ count: 0 })); -vi.mock("openclaw/plugin-sdk/sqlite-runtime", async (importOriginal) => { - const actual = await importOriginal(); +vi.mock("openclaw/plugin-sdk/sqlite-worker-runtime", async (importOriginal) => { + const actual = await importOriginal(); return { ...actual, openNodeSqliteDatabase: (...args: Parameters) => { diff --git a/extensions/workboard/src/sqlite-store-kernel.ts b/extensions/workboard/src/sqlite-store-kernel.ts index 2703b1e3bbda..7f163536d4e0 100644 --- a/extensions/workboard/src/sqlite-store-kernel.ts +++ b/extensions/workboard/src/sqlite-store-kernel.ts @@ -11,7 +11,7 @@ import { iterateSqliteQuerySync, runSqliteImmediateTransactionSync, sqliteStringSet, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; import type { PersistedWorkboardAttachment, diff --git a/extensions/workboard/src/sqlite-store-policy.test.ts b/extensions/workboard/src/sqlite-store-policy.test.ts index 9d7b248fe5da..c820b8cbfc35 100644 --- a/extensions/workboard/src/sqlite-store-policy.test.ts +++ b/extensions/workboard/src/sqlite-store-policy.test.ts @@ -8,7 +8,7 @@ const { close, configureSqliteConnectionPragmas } = vi.hoisted(() => ({ configureSqliteConnectionPragmas: vi.fn(), })); -vi.mock("openclaw/plugin-sdk/sqlite-runtime", () => ({ +vi.mock("openclaw/plugin-sdk/sqlite-worker-runtime", () => ({ openNodeSqliteDatabase: vi.fn(() => ({ close })), })); vi.mock("openclaw/plugin-sdk/plugin-state-runtime", () => ({ diff --git a/extensions/workboard/src/sqlite-store-records.ts b/extensions/workboard/src/sqlite-store-records.ts index b5513d066634..9706805bd959 100644 --- a/extensions/workboard/src/sqlite-store-records.ts +++ b/extensions/workboard/src/sqlite-store-records.ts @@ -18,7 +18,7 @@ import { getNodeSqliteKysely, iterateSqliteQuerySync, sqliteStringSet, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; export type Row = Record; export function jsonValue(value: unknown): string | null { diff --git a/extensions/workboard/src/sqlite-store-schema.ts b/extensions/workboard/src/sqlite-store-schema.ts index d402d357a3a6..02b90b8ffaa8 100644 --- a/extensions/workboard/src/sqlite-store-schema.ts +++ b/extensions/workboard/src/sqlite-store-schema.ts @@ -5,7 +5,7 @@ import { configureSqliteConnectionPragmas, migrateSqliteSchemaToStrict, } from "openclaw/plugin-sdk/plugin-state-runtime"; -import { openNodeSqliteDatabase } from "openclaw/plugin-sdk/sqlite-runtime"; +import { openNodeSqliteDatabase } from "openclaw/plugin-sdk/sqlite-worker-runtime"; const SCHEMA_VERSION = 3; const WORKBOARD_SQLITE_BUSY_TIMEOUT_MS = 5000; const WORKBOARD_SQLITE_DIR_MODE = 0o700; diff --git a/extensions/workboard/src/sqlite-store-write.ts b/extensions/workboard/src/sqlite-store-write.ts index 7e978394e999..e9e505443171 100644 --- a/extensions/workboard/src/sqlite-store-write.ts +++ b/extensions/workboard/src/sqlite-store-write.ts @@ -3,7 +3,7 @@ import type { WorkboardCard } from "@openclaw/workboard-contract"; import { compileSqliteQueryBindings, getNodeSqliteKysely, -} from "openclaw/plugin-sdk/sqlite-runtime"; +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import { jsonValue, type CARD_CHILD_TABLES, diff --git a/extensions/workboard/src/sqlite-store.worker.ts b/extensions/workboard/src/sqlite-store.worker.ts index 612de4279b02..3a45c0952074 100644 --- a/extensions/workboard/src/sqlite-store.worker.ts +++ b/extensions/workboard/src/sqlite-store.worker.ts @@ -1,4 +1,7 @@ -import type { SqliteWorkerBackend, SqliteWorkerCommand } from "openclaw/plugin-sdk/sqlite-runtime"; +import type { + SqliteWorkerBackend, + SqliteWorkerCommand, +} from "openclaw/plugin-sdk/sqlite-worker-runtime"; import type { WorkboardSqliteOperations, WorkboardSqliteWorkerOperations, diff --git a/extensions/xai/tsconfig.json b/extensions/xai/tsconfig.json index 78d87cd6c26d..626e8fe3aaf5 100644 --- a/extensions/xai/tsconfig.json +++ b/extensions/xai/tsconfig.json @@ -917,6 +917,9 @@ ], "openclaw/plugin-sdk/memory-core-host-engine-indexing": [ "../../packages/plugin-sdk/dist/src/plugin-sdk/memory-core-host-engine-indexing.d.ts" + ], + "openclaw/plugin-sdk/sqlite-worker-runtime": [ + "../../packages/plugin-sdk/dist/src/plugin-sdk/sqlite-worker-runtime.d.ts" ] } } diff --git a/package.json b/package.json index 5fa2514cd250..42b90f901a95 100644 --- a/package.json +++ b/package.json @@ -423,7 +423,8 @@ "!dist/plugin-sdk/provider-auth-aliases.d.ts", "!dist/plugin-sdk/provider-model-metadata.d.ts", "!dist/plugin-sdk/memory-core-host-engine-knn.d.ts", - "!dist/plugin-sdk/memory-core-host-engine-indexing.d.ts" + "!dist/plugin-sdk/memory-core-host-engine-indexing.d.ts", + "!dist/plugin-sdk/sqlite-worker-runtime.d.ts" ], "type": "module", "main": "dist/index.js", @@ -1198,6 +1199,9 @@ "./plugin-sdk/sqlite-runtime": { "default": "./dist/plugin-sdk/sqlite-runtime.js" }, + "./plugin-sdk/sqlite-worker-runtime": { + "default": "./dist/plugin-sdk/sqlite-worker-runtime.js" + }, "./plugin-sdk/session-transcript-hit": { "default": "./dist/plugin-sdk/session-transcript-hit.js" }, diff --git a/scripts/lib/plugin-sdk-entrypoints.json b/scripts/lib/plugin-sdk-entrypoints.json index 497aacc22a07..59e239e27f11 100644 --- a/scripts/lib/plugin-sdk-entrypoints.json +++ b/scripts/lib/plugin-sdk-entrypoints.json @@ -238,6 +238,7 @@ "session-transcript-runtime", "sqlite-runtime", "sqlite-runtime-testing", + "sqlite-worker-runtime", "session-transcript-hit", "session-visibility", "ssrf-dispatcher", diff --git a/scripts/lib/plugin-sdk-private-local-only-subpaths.json b/scripts/lib/plugin-sdk-private-local-only-subpaths.json index 22126c04f535..fbb02b2fa895 100644 --- a/scripts/lib/plugin-sdk-private-local-only-subpaths.json +++ b/scripts/lib/plugin-sdk-private-local-only-subpaths.json @@ -181,6 +181,7 @@ "speech-provider", "sqlite-runtime", "sqlite-runtime-testing", + "sqlite-worker-runtime", "ssrf-dispatcher", "ssrf-runtime-internal", "string-normalization-runtime", diff --git a/src/plugin-sdk/memory-core-host-engine-schema.ts b/src/plugin-sdk/memory-core-host-engine-schema.ts index 2611e386ce6e..17f9f8858882 100644 --- a/src/plugin-sdk/memory-core-host-engine-schema.ts +++ b/src/plugin-sdk/memory-core-host-engine-schema.ts @@ -1,10 +1,16 @@ -// Focused memory host schema helpers for doctor and migration control-plane paths. +// Memory schema operations shared by host maintenance and native publication workers. export { + dropMemoryPathFtsTriggers, + ensureMemoryChunkProvenance, ensureMemoryIndexSchema, + ensureMemoryPathFtsTriggers, + ensureMemoryRecallMetadataSchema, MEMORY_EMBEDDING_CACHE_TABLE, + MEMORY_INDEX_CHUNK_RECALL_METADATA_TABLE, MEMORY_INDEX_CHUNKS_TABLE, MEMORY_INDEX_FTS_TABLE, MEMORY_INDEX_META_TABLE, + MEMORY_INDEX_PATHS_FTS_TABLE, MEMORY_INDEX_SOURCES_TABLE, MEMORY_INDEX_VECTOR_TABLE, } from "../../packages/memory-host-sdk/src/host/memory-schema.js"; diff --git a/src/plugin-sdk/sqlite-worker-runtime.ts b/src/plugin-sdk/sqlite-worker-runtime.ts new file mode 100644 index 000000000000..8a2dfcf8d0f8 --- /dev/null +++ b/src/plugin-sdk/sqlite-worker-runtime.ts @@ -0,0 +1,28 @@ +// SQLite backends use native primitives without loading host database lifecycle owners. +export type { Generated, Selectable } from "kysely"; +export { + compileSqliteQueryBindings, + enableNodeSqliteKyselyStatementCache, + executeSqliteQuerySync, + executeSqliteQueryTakeFirstSync, + getNodeSqliteKysely, + iterateSqliteQuerySync, + prepareSqliteQuerySync, + sqliteStringSet, +} from "../infra/kysely-sync.js"; +export { + openNodeSqliteDatabase, + resolveExistingSqliteFileUri, + supportsNodeSqliteExtensionLoading, +} from "../infra/node-sqlite.js"; +export { + assertTransactionUsable, + runSqliteImmediateTransactionSync, +} from "../infra/sqlite-transaction.js"; +export type { + SqliteWorkerBackend, + SqliteWorkerCommand, + SqliteWorkerOperations, +} from "../infra/sqlite-worker-contract.js"; +export { requestSqliteWorkerOperationAdmission } from "../infra/sqlite-worker-operation-admission.js"; +export { tableExists } from "../state/openclaw-state-db-schema-helpers.js";