From 6049b27bb487ca22a339a1364641bb2499568bc9 Mon Sep 17 00:00:00 2001 From: Sergey Popov Date: Wed, 26 Aug 2026 19:57:56 +0300 Subject: [PATCH] feat: preload hybrid tables concurrently --- .changeset/green-indexes-rest.md | 2 + README.md | 5 + .../src/content/docs/runtime/db.md | 12 +- .../src/content/docs/start/llm-cheat-sheet.md | 10 +- packages/hyperdb/README.md | 5 + .../drivers/idb/idb-driver.scan-all.test.ts | 45 +++++ .../runtime/preloaded-hybrid-db.test.ts | 171 +++++++++++++++++- .../hyperdb/runtime/preloaded-hybrid-db.ts | 119 ++++++++++-- 8 files changed, 353 insertions(+), 16 deletions(-) diff --git a/.changeset/green-indexes-rest.md b/.changeset/green-indexes-rest.md index 62a6a64..8a2cce2 100644 --- a/.changeset/green-indexes-rest.md +++ b/.changeset/green-indexes-rest.md @@ -9,3 +9,5 @@ hash-index transactions provide copy-on-write commit and rollback behavior. Add `externalStorageMergeTrait` for changesets already persisted by another runtime sharing the primary. Their normal merge operations update the preloaded snapshot, notify subscribers, and persist external inserts idempotently. +Preloaded tables now load concurrently by default, with a generic +`preloadConcurrency` option accepting a positive bound or `"whole"`. diff --git a/README.md b/README.md index a7723dd..3e6e2ee 100644 --- a/README.md +++ b/README.md @@ -187,6 +187,11 @@ batch-fetch only unresolved rows through the built-in `byId` entity index. No explicit `preloadTables` call or extra B-tree `byIds` index is needed. Repeated `loadTables` calls are incremental: previously loaded tables remain available while supplied table definitions are added or refreshed. +Supplied tables preload concurrently by default. Pass +`{ preloadConcurrency: n }` to bound concurrent whole-table reads, or +`{ preloadConcurrency: "whole" }` to state the default explicitly. Drivers may +serialize internally; higher concurrency can temporarily retain several decoded +tables before their rows are released. Exact `byId` misses reconcile rows added by another connected runtime. Apply changesets already persisted by another runtime sharing the primary with `externalStorageMergeTrait`. Their normal merge operations update the preloaded diff --git a/packages/hyperdb-doc/src/content/docs/runtime/db.md b/packages/hyperdb-doc/src/content/docs/runtime/db.md index 7de0057..034bf4e 100644 --- a/packages/hyperdb-doc/src/content/docs/runtime/db.md +++ b/packages/hyperdb-doc/src/content/docs/runtime/db.md @@ -255,7 +255,9 @@ import { import { openIndexedDBDriver } from "@will-be-done/hyperdb/drivers/idb"; const primary = new DB(await openIndexedDBDriver("my-app")); -const db = new SubscribableDB(new PreloadedHybridDB(primary)); +const db = new SubscribableDB( + new PreloadedHybridDB(primary, { preloadConcurrency: "whole" }), +); // Automatically preloads all indexes on both tables, including byId. await execAsync(db.loadTables([tasksTable, projectsTable])); @@ -264,6 +266,14 @@ await execAsync(db.loadTables([tasksTable, projectsTable])); `loadTables` is incremental. A later call keeps previously loaded tables and adds or refreshes only the table definitions passed to that call. +Tables supplied to one `loadTables` call preload concurrently by default. Set +`preloadConcurrency` to a positive integer to bound active whole-table reads; +`1` is sequential, while `"whole"` starts every supplied table and is the +default. Scheduling is driver-agnostic: IndexedDB can overlap readonly cursor +transactions, while a SQLite driver may serialize access to its connection. +Each completed scan immediately builds ID-only indexes so its decoded rows can +be released, but higher concurrency can still increase peak startup memory. + A scan first reads its bounds from the in-memory ID-only index. This includes non-ID `uniqhash` indexes, which are preloaded as unique value-to-ID pointers. The built-in `byId` `uniqhash` is the canonical entity store: it begins with diff --git a/packages/hyperdb-doc/src/content/docs/start/llm-cheat-sheet.md b/packages/hyperdb-doc/src/content/docs/start/llm-cheat-sheet.md index 86a8db0..6686207 100644 --- a/packages/hyperdb-doc/src/content/docs/start/llm-cheat-sheet.md +++ b/packages/hyperdb-doc/src/content/docs/start/llm-cheat-sheet.md @@ -332,14 +332,20 @@ import { execAsync, } from "@will-be-done/hyperdb"; -const db = new SubscribableDB(new PreloadedHybridDB(primary)); +const db = new SubscribableDB( + new PreloadedHybridDB(primary, { preloadConcurrency: "whole" }), +); await execAsync(db.loadTables([tasksTable, projectsTable])); ``` `loadTables` automatically preloads all declared indexes with entity IDs as their leaves, including non-ID `uniqhash` value-to-ID pointers. Calls are incremental: existing tables remain loaded while supplied definitions are added -or refreshed. A bounded scan resolves IDs from memory and dereferences them +or refreshed. Supplied tables preload concurrently by default; use a positive +numeric `preloadConcurrency` to cap active reads (`1` is sequential), or +`"whole"` to start all supplied tables. Drivers may serialize internally, and +greater concurrency raises temporary startup memory. A bounded scan resolves +IDs from memory and dereferences them through the built-in `byId` `uniqhash`. Unresolved `byId` entries are batch-loaded from the primary and then reused by reads through every index. Do not add a `byIds` B-tree or call `preloadTables` for this runtime. Reads remain async, and diff --git a/packages/hyperdb/README.md b/packages/hyperdb/README.md index de868bb..237c7e1 100644 --- a/packages/hyperdb/README.md +++ b/packages/hyperdb/README.md @@ -184,6 +184,11 @@ resolves ordered IDs in memory, including through non-ID `uniqhash` pointers. It then uses the built-in `byId` `uniqhash` as the canonical entity cache and batch-loads only unresolved IDs. This mode does not need an explicit `preloadTables` call or a separate B-tree `byIds` index. +Tables supplied to `loadTables` preload concurrently by default. Use +`new PreloadedHybridDB(primary, { preloadConcurrency: n })` to cap concurrent +whole-table reads, or `"whole"` to make the default explicit. The storage driver +may still serialize reads internally, and higher concurrency increases temporary +startup memory. Exact `byId` misses check the primary, allowing rows added by another connected runtime to be incorporated after startup. Apply changesets already persisted by another runtime sharing the primary with diff --git a/packages/hyperdb/src/hyperdb/drivers/idb/idb-driver.scan-all.test.ts b/packages/hyperdb/src/hyperdb/drivers/idb/idb-driver.scan-all.test.ts index 2511a5c..faea2a4 100644 --- a/packages/hyperdb/src/hyperdb/drivers/idb/idb-driver.scan-all.test.ts +++ b/packages/hyperdb/src/hyperdb/drivers/idb/idb-driver.scan-all.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from "vitest"; import { execAsync } from "../../core/executor"; import { DB } from "../../runtime/db"; +import { PreloadedHybridDB } from "../../runtime/preloaded-hybrid-db"; import { defineTable } from "../../schema/table"; import { v } from "../../schema/values"; import { openIndexedDBDriver } from "./idb-driver"; @@ -10,6 +11,11 @@ const bulkRowsTable = defineTable("idbScanAllBulkRows", { value: v.number(), }); +const bulkGroupsTable = defineTable("idbScanAllBulkGroups", { + id: v.string(), + title: v.string(), +}); + let databaseCounter = 0; function deleteDatabase(databaseName: string): Promise { @@ -45,4 +51,43 @@ describe("IdbDriver scanAll", () => { await deleteDatabase(databaseName); } }); + + it("preloads multiple IndexedDB tables with whole concurrency", async () => { + databaseCounter += 1; + const databaseName = `hyperdb-idb-preloaded-${Date.now().toString(36)}-${databaseCounter}`; + await deleteDatabase(databaseName); + const driver = await openIndexedDBDriver(databaseName); + const primary = new DB(driver); + const row = { id: "row-1", value: 1 }; + const group = { id: "group-1", title: "Group" }; + + try { + await execAsync(primary.loadTables([bulkRowsTable, bulkGroupsTable])); + await execAsync(primary.insert(bulkRowsTable, [row])); + await execAsync(primary.insert(bulkGroupsTable, [group])); + + const preloaded = new PreloadedHybridDB(primary, { + preloadConcurrency: "whole", + }); + await execAsync(preloaded.loadTables([bulkRowsTable, bulkGroupsTable])); + + await expect( + execAsync( + preloaded.intervalScan(bulkRowsTable, "byId", [ + { eq: [{ col: "id", val: row.id }] }, + ]), + ), + ).resolves.toEqual([row]); + await expect( + execAsync( + preloaded.intervalScan(bulkGroupsTable, "byId", [ + { eq: [{ col: "id", val: group.id }] }, + ]), + ), + ).resolves.toEqual([group]); + } finally { + driver.close(); + await deleteDatabase(databaseName); + } + }); }); diff --git a/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.test.ts b/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.test.ts index 9234010..dc37ea5 100644 --- a/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.test.ts +++ b/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.test.ts @@ -7,10 +7,14 @@ import { import { SubscribableDB, type Op } from "./subscribable-db"; import { PreloadedTableIndexes } from "./preloaded-hybrid-db-indexes"; import { AsyncDB } from "../test-utils/async-db"; -import { createSqlJsDriver } from "../test-utils/sql-js-driver"; +import { + createSqlJsAsyncDriver, + createSqlJsDriver, +} from "../test-utils/sql-js-driver"; import { defineTable, type TableDefinition } from "../schema/table"; import { v } from "../schema/values"; -import { execAsync } from "../core/executor"; +import { execAsync, execSync } from "../core/executor"; +import { unwrap } from "../commands/async"; const tasksTable = defineTable("preloadedHybridTasks", { id: v.string(), @@ -27,6 +31,11 @@ const projectsTable = defineTable("preloadedHybridProjects", { title: v.string(), }); +const labelsTable = defineTable("preloadedHybridLabels", { + id: v.string(), + title: v.string(), +}); + const schemalessHashTable = { tableName: "preloadedHybridSchemalessHash", schema: {}, @@ -51,6 +60,52 @@ const task = (value: number, title = `Task ${value}`): Task => ({ slug: `task-${value}`, }); +type Deferred = { + promise: Promise; + resolve: (value: T) => void; + reject: (reason: unknown) => void; +}; + +function deferred(): Deferred { + let resolve!: (value: T) => void; + let reject!: (reason: unknown) => void; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, resolve, reject }; +} + +function controlTableScans(primary: DB, tables: TableDefinition[]) { + const gates = new Map( + tables.map((table) => [table.tableName, deferred()]), + ); + const started: string[] = []; + let active = 0; + let maxActive = 0; + + vi.spyOn(primary, "scanAll").mockImplementation(function* ( + table: TableDefinition, + ) { + const gate = gates.get(table.tableName); + if (!gate) throw new Error(`Missing scan gate for ${table.tableName}`); + started.push(table.tableName); + active++; + maxActive = Math.max(maxActive, active); + try { + return yield* unwrap(gate.promise); + } finally { + active--; + } + } as typeof primary.scanAll); + + return { + gates, + started, + maxActive: () => maxActive, + }; +} + async function createRuntime(rows: Task[]) { const primary = new DB(await createSqlJsDriver()); const primaryDB = new AsyncDB(primary); @@ -68,6 +123,118 @@ async function createRuntime(rows: Task[]) { } describe("PreloadedHybridDB", () => { + it("preloads all supplied tables concurrently by default", async () => { + const primary = new DB(await createSqlJsDriver()); + const tables = [tasksTable, projectsTable, labelsTable]; + const controlled = controlTableScans(primary, tables); + const loading = new AsyncDB(new PreloadedHybridDB(primary)).loadTables( + tables, + ); + + expect(controlled.started).toEqual(tables.map((table) => table.tableName)); + expect(controlled.maxActive()).toBe(3); + for (const gate of controlled.gates.values()) gate.resolve([]); + await loading; + }); + + it("bounds concurrent table preloads and preserves the option through traits", async () => { + const primary = new DB(await createSqlJsDriver()); + const tables = [tasksTable, projectsTable, labelsTable]; + const controlled = controlTableScans(primary, tables); + const runtime = new PreloadedHybridDB(primary, { preloadConcurrency: 2 }); + const loading = new AsyncDB( + runtime.withTraits({ type: "test" }), + ).loadTables(tables); + + expect(controlled.started).toEqual( + tables.slice(0, 2).map((table) => table.tableName), + ); + controlled.gates.get(tasksTable.tableName)!.resolve([]); + await vi.waitFor(() => { + expect(controlled.started).toHaveLength(3); + }); + expect(controlled.maxActive()).toBe(2); + controlled.gates.get(projectsTable.tableName)!.resolve([]); + controlled.gates.get(labelsTable.tableName)!.resolve([]); + await loading; + }); + + it("keeps successful table preloads when another table fails", async () => { + const primary = new DB(await createSqlJsDriver()); + const tables = [tasksTable, projectsTable, labelsTable]; + const controlled = controlTableScans(primary, tables); + const runtime = new PreloadedHybridDB(primary, { + preloadConcurrency: "whole", + }); + const db = new AsyncDB(runtime); + const loading = db.loadTables(tables); + const failure = new Error("projects preload failed"); + + controlled.gates.get(tasksTable.tableName)!.resolve([task(1)]); + controlled.gates.get(projectsTable.tableName)!.reject(failure); + controlled.gates + .get(labelsTable.tableName)! + .resolve([{ id: "label-1", title: "Label" }]); + + await expect(loading).rejects.toBe(failure); + await expect( + db.preloadTables([ + { table: tasksTable, scanIndex: "byId" }, + { table: labelsTable, scanIndex: "byId" }, + ]), + ).resolves.toBeUndefined(); + await expect( + db.preloadTables([{ table: projectsTable, scanIndex: "byId" }]), + ).rejects.toThrow(`Table ${projectsTable.tableName} not found`); + }); + + it("supports sequential synchronous preloading with concurrency one", async () => { + const primary = new DB(await createSqlJsDriver()); + const runtime = new PreloadedHybridDB(primary, { preloadConcurrency: 1 }); + + expect(() => + execSync(runtime.loadTables([tasksTable, projectsTable])), + ).not.toThrow(); + }); + + it.each([0, -1, 1.5, Number.NaN, Number.POSITIVE_INFINITY, "all"])( + "rejects invalid preload concurrency %s", + async (preloadConcurrency) => { + const primary = new DB(await createSqlJsDriver()); + expect( + () => + new PreloadedHybridDB(primary, { + preloadConcurrency: preloadConcurrency as 1, + }), + ).toThrow('preloadConcurrency must be "whole" or a positive integer'); + }, + ); + + it("loads through an async SQLite primary with whole concurrency", async () => { + const primary = new DB(await createSqlJsAsyncDriver()); + const primaryDB = new AsyncDB(primary); + const row = task(1); + await primaryDB.loadTables([tasksTable, projectsTable]); + await primaryDB.insert(tasksTable, [row]); + await primaryDB.insert(projectsTable, [{ id: "p1", title: "Project" }]); + + const db = new AsyncDB( + new PreloadedHybridDB(primary, { preloadConcurrency: "whole" }), + ); + await db.loadTables([tasksTable, projectsTable]); + + await expect( + db.intervalScan(tasksTable, "byValue", [ + { eq: [{ col: "value", val: row.value }] }, + ]), + ).resolves.toEqual([row]); + await expect( + db.intervalScan(projectsTable, "byId", [ + { eq: [{ col: "id", val: "p1" }] }, + ]), + ).resolves.toEqual([{ id: "p1", title: "Project" }]); + }); + it("preloads every table index without retaining entity rows", async () => { const rows = [task(1, "same"), task(2, "same"), task(3, "other")]; const { db, scanAllSpy, intervalScanSpy } = await createRuntime(rows); diff --git a/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.ts b/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.ts index 0903377..ca8743f 100644 --- a/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.ts +++ b/packages/hyperdb/src/hyperdb/runtime/preloaded-hybrid-db.ts @@ -1,5 +1,6 @@ import type { DBCmd } from "../commands/async"; import { unwrap } from "../commands/async"; +import { execMaybeAsync } from "../core/executor"; import type { HyperDB, HyperDBTx, @@ -139,8 +140,75 @@ const createState = (): PreloadedHybridDBState => ({ export type PreloadedHybridDBOptions = { traits?: Trait[]; + preloadConcurrency?: number | "whole"; }; +type TablePreloadResult = + | { + status: "fulfilled"; + table: TableDefinition; + state: PreloadedTableState; + } + | { status: "rejected"; table: TableDefinition; reason: unknown }; + +function normalizePreloadConcurrency( + concurrency: number | "whole" | undefined, +): number | "whole" { + if (concurrency === undefined || concurrency === "whole") return "whole"; + if (!Number.isSafeInteger(concurrency) || concurrency <= 0) { + throw new Error('preloadConcurrency must be "whole" or a positive integer'); + } + return concurrency; +} + +function createPreloadedTableState( + table: TableDefinition, + rows: Row[], +): PreloadedTableState { + const byId = new HashIndex({ + name: table.idIndexName, + unique: true, + }); + byId.insert(entityPointers(rows)); + return { + indexes: new PreloadedTableIndexes(table, rows), + byId, + }; +} + +async function preloadTablesConcurrently( + primary: DB, + tables: TableDefinition[], + concurrency: number, +): Promise { + const results = new Array(tables.length); + let nextTableIndex = 0; + + const worker = async (): Promise => { + while (nextTableIndex < tables.length) { + const tableIndex = nextTableIndex++; + const table = tables[tableIndex]!; + try { + const pendingRows = execMaybeAsync(primary.scanAll(table)); + const rows = + pendingRows instanceof Promise ? await pendingRows : pendingRows; + results[tableIndex] = { + status: "fulfilled", + table, + state: createPreloadedTableState(table, rows as Row[]), + }; + } catch (reason) { + results[tableIndex] = { status: "rejected", table, reason }; + } + } + }; + + await Promise.all( + Array.from({ length: Math.min(concurrency, tables.length) }, worker), + ); + return results; +} + function* acquireLock(lock: AwaitLock): Generator void> { if (!lock.tryAcquire()) yield* unwrap(lock.acquireAsync()); let released = false; @@ -239,17 +307,22 @@ function* scanPreloaded( export class PreloadedHybridDB implements HyperDB { readonly primary: DB; traits: Trait[]; + private readonly preloadConcurrency: number | "whole"; private state: PreloadedHybridDBState; constructor(primary: DB, options: PreloadedHybridDBOptions = {}) { this.primary = primary; this.traits = options.traits ?? []; + this.preloadConcurrency = normalizePreloadConcurrency( + options.preloadConcurrency, + ); this.state = createState(); } withTraits(...traits: Trait[]): HyperDB { const db = new PreloadedHybridDB(this.primary, { traits: [...this.traits, ...traits], + preloadConcurrency: this.preloadConcurrency, }); db.state = this.state; return db; @@ -290,17 +363,41 @@ export class PreloadedHybridDB implements HyperDB { try { yield* this.delegatePrimary().loadTables(tables); const data = this.state.data; - for (const table of tables) { - const rows = yield* this.primary.scanAll(table); - const byId = new HashIndex({ - name: table.idIndexName, - unique: true, - }); - byId.insert(entityPointers(rows)); - data.tables.set(table.tableName, { - indexes: new PreloadedTableIndexes(table, rows as Row[]), - byId, - }); + const concurrency = + this.preloadConcurrency === "whole" + ? tables.length + : Math.min(this.preloadConcurrency, tables.length); + + if (concurrency <= 1) { + let hasError = false; + let firstError: unknown; + for (const table of tables) { + try { + const rows = yield* this.primary.scanAll(table); + data.tables.set( + table.tableName, + createPreloadedTableState(table, rows as Row[]), + ); + } catch (error) { + if (!hasError) firstError = error; + hasError = true; + } + } + if (hasError) throw firstError; + return; + } + + const results = yield* unwrap( + preloadTablesConcurrently(this.primary, tables, concurrency), + ); + for (const result of results) { + if (result.status === "fulfilled") { + data.tables.set(result.table.tableName, result.state); + } + } + const failed = results.find((result) => result.status === "rejected"); + if (failed?.status === "rejected") { + throw failed.reason; } } finally { release();