diff --git a/.gitignore b/.gitignore index 7881bd9..5d0a844 100644 --- a/.gitignore +++ b/.gitignore @@ -9,3 +9,4 @@ coverage/ .alchemy/ playwright-report/ test-results/ +duckdb-benchmark*.json diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index e5373d7..8d262eb 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -121,6 +121,16 @@ hosts typed configuration as a modular layer over the same core. scalar-family joins and typed keyset ordering compare values consistently with the KV executor. A shared KV/SQLite/PostgreSQL conformance corpus is the regression boundary for this contract. - Backend packages construct the storage adapters and runtime layers for their platforms. +- `@bjacobso/triplex-duckdb` is a private analytical snapshot POC. It copies a migrated SQLite + source through `SqlClient` in one transaction, imports bounded batches into native DuckDB, + and reuses the shared Datalog compiler and row decoder at a pinned temporal/commit basis. + Its Effect service exposes read-only analytical queries, not the writable `Triples` contract. + Its federation service opens a trusted catalog of local SQLite sources through independent + read-only clients. Local worker processes filter pinned facts and virtual database membership, + then a DuckDB coordinator runs the shared compiler over source-qualified identities. Source + shortcuts lower to ordinary membership patterns without changing the core Datalog contract. + It does not implement distributed recursive rounds, incremental projection, or object storage. + See `packages/duckdb/README.md` for the opt-in comparison benchmark and limitations. - `@bjacobso/triplex-cloudflare` composes its synchronous Durable Object SQLite adapter with the shared SQL compiler/row decoder through the narrow `SqlStatementRunner` contract. Its public `CloudflareTriples.layer({ state, scope })` binds migrations, Datalog, writes, journal positions, diff --git a/docs/datalog-performance.md b/docs/datalog-performance.md index 6920ea7..f0f1ed6 100644 --- a/docs/datalog-performance.md +++ b/docs/datalog-performance.md @@ -1,5 +1,10 @@ # Datalog pagination and performance +A private [DuckDB analytical snapshot POC](https://github.com/bjacobso/triplex/tree/main/packages/duckdb) compares the shared +Datalog executor over SQLite and native DuckDB at the same commit and temporal cut. Its +benchmark is opt-in and reports import cost separately; it does not establish distributed +execution or replace SQLite's transactional backend. + `Triples.query` and `Triples.queryPage` return at most 100 bindings by default. Page sizes are validated in the shared service boundary and may not exceed 1,000. SQL fetches one extra row to determine whether `nextCursor` exists; only the requested page reaches the caller. Total counts diff --git a/packages/core/src/datalog/wrapper.ts b/packages/core/src/datalog/wrapper.ts index 46cca1f..e700972 100644 --- a/packages/core/src/datalog/wrapper.ts +++ b/packages/core/src/datalog/wrapper.ts @@ -84,12 +84,12 @@ const compileFilterOp = ( case "ilike": // SQLite: use LIKE with COLLATE NOCASE // PostgreSQL: use native ILIKE - if (dialect.name === "postgresql") { + if (dialect.name === "postgresql" || dialect.name === "duckdb") { return `${colName} ILIKE ${collector.add(value)}`; } return `${colName} LIKE ${collector.add(value)} COLLATE NOCASE`; case "not-ilike": - if (dialect.name === "postgresql") { + if (dialect.name === "postgresql" || dialect.name === "duckdb") { return `${colName} NOT ILIKE ${collector.add(value)}`; } return `${colName} NOT LIKE ${collector.add(value)} COLLATE NOCASE`; diff --git a/packages/core/src/dialects/types.ts b/packages/core/src/dialects/types.ts index 63fdf7d..0bc8de3 100644 --- a/packages/core/src/dialects/types.ts +++ b/packages/core/src/dialects/types.ts @@ -1,7 +1,7 @@ import { Context } from "effect"; export interface SqlDialect { - readonly name: "sqlite" | "postgresql" | "mysql"; + readonly name: "sqlite" | "postgresql" | "mysql" | "duckdb"; readonly limitOffset: (limit: number | undefined, offset: number | undefined) => string; readonly castAsText: (expression: string) => string; readonly booleanLiteral: (value: boolean) => string; diff --git a/packages/duckdb/LICENSE b/packages/duckdb/LICENSE new file mode 100644 index 0000000..ac737c0 --- /dev/null +++ b/packages/duckdb/LICENSE @@ -0,0 +1,21 @@ +MIT License + +Copyright (c) 2026 Ben Jacobson + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. diff --git a/packages/duckdb/README.md b/packages/duckdb/README.md new file mode 100644 index 0000000..9d07c82 --- /dev/null +++ b/packages/duckdb/README.md @@ -0,0 +1,254 @@ +# DuckDB snapshots and local federation POC + +Private, experimental Node.js package. Copies a migrated Triplex SQLite database into an +in-memory DuckDB columnar table and executes Triplex Datalog through the shared SQL compiler +and result decoder. SQLite remains the transactional system of record. + +This evaluates native DuckDB execution without an engine fork. It includes a local multi-process +federation POC: one immutable reader per SQLite database and a DuckDB query coordinator. +No DuckDB extensions or network downloads are required at runtime. + +## Run + +From the repository root: + +```sh +pnpm install +pnpm exec turbo run build --filter=@bjacobso/triplex-duckdb +pnpm --filter @bjacobso/triplex-duckdb test +pnpm --filter @bjacobso/triplex-duckdb --silent benchmark > duckdb-benchmark.json +``` + +The benchmark creates disposable in-memory databases. It runs only when explicitly invoked; +small correctness tests participate in `pnpm check`. + +```sh +DUCKDB_BENCH_ENTITIES=100000 DUCKDB_BENCH_ROUNDS=5 \ + pnpm --filter @bjacobso/triplex-duckdb --silent benchmark > duckdb-benchmark.json +``` + +`DUCKDB_BENCH_ENTITIES` defaults to 10,000 (maximum 1,000,000), `DUCKDB_BENCH_ROUNDS` to 5, +and `DUCKDB_BENCH_THREADS` to 4. `DUCKDB_BENCH_GRAPH_NODES` defaults to 128, capped by the +entity count. Increase it separately to stress recursive closure: the current SQLite plan +can become expensive even for modest graphs. + +## Query two local SQLite databases + +Run the disposable two-database example (including real worker processes): + +```sh +pnpm --filter @bjacobso/triplex-duckdb demo:federation +``` + +Both databases contain `person:1` with the same email. The output contains one cross-database +match with distinct, source-qualified person IDs, worker PIDs, per-source commit cuts, scanned +attributes, and transferred row counts. The script verifies the two query syntaxes below. + +```ts +import { Effect } from "effect"; +import { DuckdbFederation, SnapshotProvider } from "@bjacobso/triplex-duckdb"; + +const program = Effect.gen(function* () { + const db = yield* DuckdbFederation; + return yield* db.queryAll({ + sources: { $a: "a", $b: "b" }, + find: ["?personA", "?personB", "?email"], + where: [ + ["$a", "?personA", ":person/email", "?email"], + ["$b", "?personB", ":person/email", "?email"], + ], + }); +}); + +const result = await Effect.runPromise( + program.pipe( + Effect.provide(DuckdbFederation.layer()), + Effect.provide( + SnapshotProvider.local([ + { id: "a", tenant: "customer-a", filename: "/absolute/customer-a.sqlite" }, + { id: "b", tenant: "customer-b", filename: "/absolute/customer-b.sqlite" }, + ]), + ), + ), +); +``` + +Use existing Triplex SQLite files; federation opens them read-only and does not migrate or +create them. The catalog is trusted application configuration, not user query input. Aliases +bind catalog IDs, never filenames, and declaring an alias does not restrict unscoped patterns. +Only include databases the caller may query: this POC does not implement tenant authorization. + +The same query can use ordinary Datalog membership and catalog patterns: + +```ts +const query = { + find: ["?personA", "?personB", "?email"], + where: [ + ["?dbA", ":triplex/tenant", "customer-a"], + ["?personA", ":triplex/database", "?dbA"], + ["?dbB", ":triplex/tenant", "customer-b"], + ["?personB", ":triplex/database", "?dbB"], + ["?personA", ":person/email", "?email"], + ["?personB", ":person/email", "?email"], + ], +} as const; +``` + +`[$a, e, a, v, tx?]` lowers to the ordinary fact pattern plus database membership. The +`:triplex/database` ref links local entity, referenced, and transaction identities to their +database. `:triplex/tenant` links a database entity to its catalog tenant string (multiple +databases may share a tenant). Transaction provenance also works through +`["?e", ":person/email", "?email", "?tx"]`, `["?tx", ":_tx/database", "?db"]`, and +`["?db", ":triplex/tenant", "customer-a"]`. These three attributes are reserved virtual +facts; acquisition rejects visible source facts that use them. + +Use `sourceEntity(databaseId, localId)` for literal entity/transaction IDs and ref constants +in ordinary patterns; `decodeSourceEntity` reverses result IDs. `databaseEntity(id)` addresses +a catalog entity. Source shortcuts qualify literal identities and typed ref constants +automatically. Plain string values remain unchanged, and every stored ref is local to its +source. Equal local IDs in different databases never join as entities, refs, or transactions. +Cross-database joins use shared scalar values such as email. + +`makeDuckdbFederation(options)` requires `SnapshotProvider` and `Scope`. The layer handles +scope management. `mode` defaults to `"workers"`; `"in-process"` uses the identical source +reader without IPC and serves as a reference. `batchSize` defaults to 2,048 (maximum 65,536), +coordinator `threads` to 4 (maximum 64), and `workerTimeoutMs` to 30,000 per request. Each +worker uses one native DuckDB thread. At most 32 configured databases are accepted. + +Each source is copied once at acquisition in a SQLite read transaction, then filtered at its +own pinned commit position and `basis: { recordedAt?, validAt? }`. These cuts are independent, +not a globally atomic snapshot. Subsequent writes do not change results; reacquire federation +to refresh. `sources` reports the cuts even before a query. Scope release closes native +handles and workers; a worker crash or timeout fails the query without partial results or +silently acquiring a newer snapshot. + +`explain(query)` shows the normalized core query, selected databases, and projected attributes. +Flat source-constrained conjunctions can prune databases; complex clauses conservatively scan +all configured sources. Workers scan visible facts and requested membership attributes in +bounded IPC batches. The coordinator imports them into a fresh DuckDB table and executes joins, +negation, aggregation, and the existing bounded recursive rules. Results include `federation` +with per-source transferred row counts; `queryAll(query, true)` also exposes generated SQL. `queryWindow` +accepts a wrapped query with a federated `inner`, ordering, limit, and optional count. + +Current limits: + +- Source shortcuts work in top-level patterns and conjunctive `not`. Use ordinary core syntax + in `or` and rules; unsupported shortcuts fail explicitly. The core negation clause limit + applies after expansion. Bounded recursion runs at the coordinator, with no distributed + fixpoint rounds or cross-source stored refs. +- Pushdown currently filters attributes, not join keys, predicates, or partial aggregates. + Every configured source is snapshotted at acquisition, including subsequently pruned sources. + Snapshots and each coordinator query table are in memory; IPC batches bound transfer size, + not total memory. Queries are serialized within one federation instance. This establishes + semantics and process boundaries, not a horizontal scaling or performance claim. +- No S3 access, replication, remote RPC, authentication, incremental updates, or snapshot cache. + A future `SnapshotProvider` can restore a consistent, versioned SQLite snapshot locally. + Direct random reads of SQLite pages from S3 require a separate storage/VFS design; uploading + a SQLite file or its WAL does not make remote access equivalent to a local file. + +## Use a pinned snapshot + +```ts +import { Effect } from "effect"; +import { makeSqliteLayer } from "@bjacobso/triplex-sqlite"; +import { DuckdbSnapshot } from "@bjacobso/triplex-duckdb"; + +const program = Effect.gen(function* () { + const snapshot = yield* DuckdbSnapshot; + const result = yield* snapshot.queryAll({ + find: ["?group", "?count"], + where: [["?entity", ":person/group", "?group"]], + aggregate: [["count", "?entity", "?count"]], + }); + return { snapshot: snapshot.metadata, rows: result.results }; +}); + +const result = await Effect.runPromise( + program.pipe( + Effect.provide(DuckdbSnapshot.layer({ scope: "my-database" })), + Effect.provide(makeSqliteLayer("data.sqlite")), + ), +); +``` + +`makeDuckdbSnapshot(options)` is also available for creating a snapshot after seeding or +writing within an existing Effect program. It requires `SqlClient` and `Scope`; keep all +queries inside `Effect.scoped`. `DuckdbSnapshot.layer` manages that scope for its consumers. +Use an already migrated SQLite source; the example's SQLite layer applies migrations. + +Options include `basis: { recordedAt?, validAt? }`, `batchSize` (default 2,048; maximum +65,536), and native `threads` (default 4; maximum 64). Scope is a caller-supplied source label, +not an authorization boundary. + +The importer reads the commit counter and all fact batches in one SQLite transaction. +It retains assertion/retraction history and pins every query to that commit position and +the resolved temporal basis. Later source writes cannot change a snapshot. To observe new +commits or choose another business-time instant, acquire a new snapshot. Invoke acquisition +outside an ambient write transaction if the snapshot must contain only committed data. + +`queryAll(query, debug?)` completely materializes results. `queryWindow(wrappedQuery, debug?)` +supports wrapper filtering, explicit ordering, limits, and optional counts. Supply an explicit +`orderBy` with tie-breakers when comparing limited results. This low-level analytical service +does not expose `Triples` opaque cursors, subscriptions, transactions, or read-after-write +waiting. Debug output includes generated SQL, parameters, and compiler/execution metrics. + +## What the benchmark measures + +Both engines receive the same Datalog and exact temporal basis through the shared executor. +The fixture contains names, numeric scores, groups, bounded reference chains, and a retracted +score plus replacement for every tenth entity. It uses direct SQL seeding and a synthetic +commit counter, so it measures query execution rather than write throughput or journal cost. + +For both current and historical cuts, the benchmark compares attribute scans, point lookups, +joins, grouped counts, bounded recursion, ordered first pages, and pages with counts. Each +backend is warmed before measurement. Timings include compilation, execution, and decoding; +seeding, snapshot import, warm-up, and checksum calculation are excluded. Import and seed +times are reported separately. Results and counts must match SQLite or the command fails. +JSON includes versions, fixture dimensions, snapshot metadata, result checksums, generated +SQL, and median/minimum/maximum elapsed times. Progress goes to stderr. + +The default recursion fixture is intentionally small and remains fixed as the fact count +grows. It measures traversal amid unrelated facts, not distributed or large-graph scaling. +These are warm, single-client measurements in one process, with indexed SQLite and DuckDB's +shared-schema primary key; they are not concurrency, disk, or cold-start benchmarks. + +## Local measurement + +Measured September 9, 2026 in the development Linux x64 VM, using Node 24.14.1, +SQLite 3.51.2, DuckDB 1.5.5, and four DuckDB threads. This run used 100,000 entities, +310,124 stored facts including retractions, and 128 graph nodes. Values below are the +median of five warm executions at the current fixture cut (`recordedAt=25`, `validAt=100`, +`recordedPosition=2`). Builds and tests were idle during this run. + +| Query | SQLite | DuckDB | +| -------------------- | --------: | --------: | +| attribute-scan | 254.65 ms | 236.16 ms | +| point-lookup | 0.19 ms | 4.23 ms | +| join | 294.67 ms | 214.40 ms | +| grouped-count | 161.56 ms | 23.08 ms | +| bounded-recursion | 72.85 ms | 36.14 ms | +| first-100 | 108.41 ms | 21.51 ms | +| first-100-with-count | 202.47 ms | 39.02 ms | + +Snapshot acquisition and import cost 2.39 seconds, excluded from the query timings. +Both temporal cuts passed result and count comparisons. The gains are workload-dependent: +point lookups favor SQLite, while grouped counts and ordered pages benefit from DuckDB. +The small fixed graph does not demonstrate large-graph scaling. These local diagnostics +are not production latency guarantees. + +## Boundaries and next experiments + +The POC transfers bounded batches through JavaScript, but retains the entire imported +database in DuckDB and materializes complete query results in JavaScript. It holds the SQLite +read transaction for the duration of copying; long copies can retain WAL history. Native +queries finish before scope cleanup, so interruption waits for an in-flight native query. + +Focused tests cover typed projections, joins, negation, disjunction, aggregation, recursive +cycles, ordered windows/counts, recorded/valid time, and concurrent source writes during +batching. This is not the full writable-backend conformance suite. DuckDB uses native `ILIKE`; +cross-engine Unicode collation equivalence has not been established. + +The next experiments are direct SQLite scanning versus native columnar storage, immutable +Parquet segments with snapshot manifests, and incremental projection from the Triplex journal. +Cross-shard recursive exchange, global snapshot coordination, retraction maintenance, and +object-storage persistence still require separate implementations. diff --git a/packages/duckdb/package.json b/packages/duckdb/package.json new file mode 100644 index 0000000..e5aedef --- /dev/null +++ b/packages/duckdb/package.json @@ -0,0 +1,66 @@ +{ + "name": "@bjacobso/triplex-duckdb", + "version": "0.0.0", + "private": true, + "description": "Experimental DuckDB analytical snapshots for Triplex", + "keywords": [ + "database", + "effect", + "sql", + "triplestore" + ], + "homepage": "https://github.com/bjacobso/triplex#readme", + "bugs": { + "url": "https://github.com/bjacobso/triplex/issues" + }, + "license": "MIT", + "repository": { + "type": "git", + "url": "git+https://github.com/bjacobso/triplex.git", + "directory": "packages/duckdb" + }, + "files": [ + "dist", + "LICENSE", + "README.md" + ], + "type": "module", + "sideEffects": false, + "main": "./dist/index.js", + "types": "./dist/index.d.ts", + "exports": { + ".": { + "types": "./dist/index.d.ts", + "import": "./dist/index.js", + "default": "./dist/index.js" + } + }, + "scripts": { + "build": "tsdown && tsc --declaration --emitDeclarationOnly --rootDir src --outDir dist", + "dev": "tsdown --watch", + "test": "vitest --run --passWithNoTests", + "typecheck": "tsc --noEmit -p tsconfig.check.json", + "clean": "rm -rf dist", + "benchmark": "tsx scripts/benchmark.ts", + "demo:federation": "pnpm run build && tsx scripts/federation.ts" + }, + "dependencies": { + "@bjacobso/triplex": "workspace:^", + "@bjacobso/triplex-sql": "workspace:^", + "@duckdb/node-api": "catalog:", + "@effect/sql-sqlite-node": "catalog:" + }, + "devDependencies": { + "@bjacobso/triplex-sqlite": "workspace:^", + "@types/node": "catalog:", + "effect": "catalog:", + "tsx": "catalog:", + "vitest": "catalog:" + }, + "peerDependencies": { + "effect": "catalog:" + }, + "engines": { + "node": ">=22" + } +} diff --git a/packages/duckdb/scripts/benchmark.ts b/packages/duckdb/scripts/benchmark.ts new file mode 100644 index 0000000..5160f21 --- /dev/null +++ b/packages/duckdb/scripts/benchmark.ts @@ -0,0 +1,257 @@ +/** Opt-in, disposable in-memory analytical benchmark. No provider credentials. */ +import { createHash } from "node:crypto"; +import { Effect } from "effect"; +import { SqlClient } from "effect/unstable/sql"; +import { type DatalogQuery, type WrappedQuery } from "@bjacobso/triplex"; +import { QueryExecutor } from "@bjacobso/triplex/internal"; +import { SqliteTriples } from "@bjacobso/triplex-sqlite"; +import { makeDuckdbSnapshot } from "../src/index.js"; + +const positiveInteger = (name: string, fallback: number, max: number) => { + const value = Number(process.env[name] ?? fallback); + if (!Number.isSafeInteger(value) || value < 1 || value > max) + throw new Error(`${name} must be 1..${max}`); + return value; +}; +const entities = positiveInteger("DUCKDB_BENCH_ENTITIES", 10_000, 1_000_000); +const rounds = positiveInteger("DUCKDB_BENCH_ROUNDS", 5, 100); +const threads = positiveInteger("DUCKDB_BENCH_THREADS", 4, 64); +const graphNodes = Math.min(entities, positiveInteger("DUCKDB_BENCH_GRAPH_NODES", 128, 1_000_000)); +const names: DatalogQuery = { find: ["?e", "?name"], where: [["?e", ":bench/name", "?name"]] }; +const cases: readonly { name: string; query: DatalogQuery; window?: WrappedQuery }[] = [ + { name: "attribute-scan", query: names }, + { + name: "point-lookup", + query: { find: ["?name"], where: [["e:0000000", ":bench/name", "?name"]] }, + }, + { + name: "join", + query: { + find: ["?e", "?name", "?score"], + where: [ + ["?e", ":bench/name", "?name"], + ["?e", ":bench/score", "?score"], + [">=", "?score", 50], + ], + }, + }, + { + name: "grouped-count", + query: { + find: ["?group", "?count"], + where: [["?e", ":bench/group", "?group"]], + aggregate: [["count", "?e", "?count"]], + }, + }, + { + name: "bounded-recursion", + query: { + find: ["?target"], + where: [["reach", "e:0000000", "?target"]], + rules: [ + { name: "reach", body: [["?x", ":bench/next", "?y"]], maxDepth: 16 }, + { + name: "reach", + body: [ + ["?x", ":bench/next", "?z"], + ["reach", "?z", "?y"], + ], + maxDepth: 16, + }, + ], + }, + }, + { + name: "first-100", + query: names, + window: { inner: names, limit: 100, orderBy: [{ variable: "?e", direction: "asc" }] }, + }, + { + name: "first-100-with-count", + query: names, + window: { + inner: names, + limit: 100, + includeCount: true, + orderBy: [{ variable: "?e", direction: "asc" }], + }, + }, +]; + +const checksum = (rows: readonly Record[]) => + createHash("sha256") + .update( + JSON.stringify( + rows + .map((row) => + JSON.stringify( + Object.fromEntries(Object.entries(row).sort(([a], [b]) => a.localeCompare(b))), + ), + ) + .sort(), + ), + ) + .digest("hex"); + +const program = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const sqlite = yield* QueryExecutor; + const seedStart = performance.now(); + // Synthetic facts bypass the write API deliberately; these timings measure + // query execution, not transaction throughput. Include retraction history. + yield* sql.withTransaction( + Effect.gen(function* () { + for (let start = 0; start < entities; start += 100) { + const params: unknown[] = []; + const tuples: string[] = []; + const add = (row: unknown[]) => + tuples.push( + `(${row + .map((value) => { + params.push(value); + return `?${params.length}`; + }) + .join(",")})`, + ); + for (let i = start; i < Math.min(start + 100, entities); i++) { + const entity = `e:${String(i).padStart(7, "0")}`; + const updated = i % 10 === 0; + for (const [attribute, type, text, value] of [ + [":bench/name", "string", `Person ${i}`, null], + [":bench/score", "number", null, i % 100], + [":bench/group", "string", `group:${i % 250}`, null], + ]) { + const retracted = updated && attribute === ":bench/score"; + add([ + `${entity}:${attribute}`, + entity, + attribute, + type, + text, + value, + 10, + 1, + 0, + retracted ? 20 : null, + retracted ? 2 : null, + ]); + } + if (updated) + add([ + `${entity}:new-score`, + entity, + ":bench/score", + "number", + null, + 99, + 20, + 2, + 0, + null, + null, + ]); + // Disjoint chains keep closure size linear in entities at fixed depth. + if (i % 32 !== 31 && i + 1 < graphNodes) + add([ + `${entity}:next`, + entity, + ":bench/next", + "ref", + `e:${String(i + 1).padStart(7, "0")}`, + null, + 10, + 1, + 0, + null, + null, + ]); + } + yield* sql.unsafe( + `INSERT INTO triples + (id,entity_id,attribute,value_type,value_string,value_number,recorded_at,recorded_position,valid_from,retracted_at,retracted_position) + VALUES ${tuples.join(",")}`, + params, + ); + } + yield* sql.unsafe("INSERT INTO triplex_commit_position VALUES (1,2)"); + }), + ); + yield* sql.unsafe("ANALYZE triples"); + const seedTimeMs = performance.now() - seedStart; + const reports = []; + for (const recordedAt of [25, 15]) { + const report = yield* Effect.scoped( + Effect.gen(function* () { + const snapshot = yield* makeDuckdbSnapshot({ + scope: "benchmark", + basis: { recordedAt, validAt: 100 }, + threads, + }); + const measurements = []; + for (const entry of cases) { + process.stderr.write(`recordedAt=${recordedAt} query=${entry.name}\n`); + const requests = { + sqlite: entry.window + ? sqlite.executePage(entry.window, true, snapshot.metadata.basis) + : sqlite.execute(entry.query, true, snapshot.metadata.basis), + duckdb: entry.window + ? snapshot.queryWindow(entry.window, true) + : snapshot.queryAll(entry.query, true), + }; + const baseline = yield* requests.sqlite; + const expected = checksum(baseline.results); + const expectedCount = "totalCount" in baseline ? baseline.totalCount : undefined; + for (const backend of ["sqlite", "duckdb"] as const) { + yield* requests[backend]; // Warm up separately from measurements. + const elapsed: number[] = []; + let result = baseline; + for (let round = 0; round < rounds; round++) { + const start = performance.now(); + result = yield* requests[backend]; + elapsed.push(performance.now() - start); + } + const digest = checksum(result.results); + const totalCount = "totalCount" in result ? result.totalCount : undefined; + if (digest !== expected || totalCount !== expectedCount) + return yield* Effect.fail( + new Error(`Result mismatch: ${entry.name}, ${backend}, recordedAt=${recordedAt}`), + ); + elapsed.sort((a, b) => a - b); + const middle = Math.floor(elapsed.length / 2); + measurements.push({ + query: entry.name, + backend, + rows: result.results.length, + totalCount, + medianMs: + elapsed.length % 2 === 0 + ? (elapsed[middle - 1]! + elapsed[middle]!) / 2 + : elapsed[middle], + minMs: elapsed[0], + maxMs: elapsed[elapsed.length - 1], + checksum: digest, + sql: result.debug?.generatedSql, + params: result.debug?.params, + }); + } + } + return { snapshot: snapshot.metadata, measurements }; + }), + ); + reports.push(report); + } + const version = yield* sql.unsafe<{ version: string }>("SELECT sqlite_version() AS version"); + return { + entities, + graphNodes, + rounds, + seedTimeMs, + sqliteVersion: version[0]?.version, + node: process.version, + platform: `${process.platform}/${process.arch}`, + reports, + }; +}); + +const report = await Effect.runPromise(program.pipe(Effect.provide(SqliteTriples.layerMemory))); +process.stdout.write(`${JSON.stringify(report, null, 2)}\n`); diff --git a/packages/duckdb/scripts/federation.ts b/packages/duckdb/scripts/federation.ts new file mode 100644 index 0000000..ad41366 --- /dev/null +++ b/packages/duckdb/scripts/federation.ts @@ -0,0 +1,73 @@ +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { Effect } from "effect"; +import { EntityId, Triples, string } from "@bjacobso/triplex"; +import { SqliteTriples } from "@bjacobso/triplex-sqlite"; +import { makeDuckdbFederation, SnapshotProvider, type FederatedQuery } from "../dist/index.js"; + +const directory = await mkdtemp(join(tmpdir(), "triplex-duckdb-demo-")); +try { + const catalog = ["a", "b"].map((id) => ({ + id, + tenant: `customer-${id}`, + filename: join(directory, `${id}.sqlite`), + })); + for (const database of catalog) { + await Effect.runPromise( + Effect.gen(function* () { + const triples = yield* Triples; + yield* triples.assertBatch([ + { + entityId: EntityId.make("person:1"), + attribute: ":person/email", + value: string("shared@example.com"), + }, + { + entityId: EntityId.make("person:1"), + attribute: ":person/name", + value: string(database.id === "a" ? "Alice" : "Bob"), + }, + ]); + }).pipe(Effect.provide(SqliteTriples.layer({ filename: database.filename }))), + ); + } + const result = await Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const federation = yield* makeDuckdbFederation(); + const shortcut: FederatedQuery = { + sources: { $a: "a", $b: "b" }, + find: ["?personA", "?personB", "?email"], + where: [ + ["$a", "?personA", ":person/email", "?email"], + ["$b", "?personB", ":person/email", "?email"], + ], + }; + const membership: FederatedQuery = { + find: shortcut.find, + where: [ + ["?dbA", ":triplex/tenant", "customer-a"], + ["?personA", ":triplex/database", "?dbA"], + ["?dbB", ":triplex/tenant", "customer-b"], + ["?personB", ":triplex/database", "?dbB"], + ["?personA", ":person/email", "?email"], + ["?personB", ":person/email", "?email"], + ], + }; + const a = yield* federation.queryAll(shortcut); + const b = yield* federation.queryAll(membership); + if (JSON.stringify(a.results) !== JSON.stringify(b.results)) + throw new Error("Syntax results differ"); + return { + results: a.results, + federation: a.federation, + membershipPlan: yield* federation.explain(membership), + }; + }), + ).pipe(Effect.provide(SnapshotProvider.local(catalog))), + ); + console.log(JSON.stringify(result, null, 2)); +} finally { + await rm(directory, { recursive: true, force: true }); +} diff --git a/packages/duckdb/src/DuckdbFederation.ts b/packages/duckdb/src/DuckdbFederation.ts new file mode 100644 index 0000000..bf330ce --- /dev/null +++ b/packages/duckdb/src/DuckdbFederation.ts @@ -0,0 +1,232 @@ +import { isAbsolute } from "node:path"; +import { Context, Effect, Layer, Semaphore, type Scope } from "effect"; +import { ReadError, type DatalogQueryError, type WrappedQuery } from "@bjacobso/triplex"; +import type { QueryExecutorService } from "@bjacobso/triplex/internal"; +import { makeSqlQueryExecutor } from "@bjacobso/triplex-sql"; +import { DuckdbDialect } from "./dialect.js"; +import { makeTable } from "./federation/table.js"; +import { makeSource, type SourceWorker } from "./federation/source.js"; +import { makeProcessSource } from "./federation/worker-client.js"; +import { planFederation, type FederationPlan } from "./federation/planner.js"; +import { + SnapshotProvider, + DATABASE_ATTRIBUTE, + TENANT_ATTRIBUTE, + TX_DATABASE_ATTRIBUTE, + databaseEntity, + sourceEntity, + federationFailure, + type FactRow, + type FederatedQuery, +} from "./federation/types.js"; + +export interface DuckdbFederationOptions { + /** Separate local processes by default; in-process is useful as a reference. */ + readonly mode?: "workers" | "in-process"; + readonly batchSize?: number; + readonly threads?: number; + readonly workerTimeoutMs?: number; +} +export interface FederationSourceMetadata { + readonly database: string; + readonly tenant: string; + readonly pid: number; + readonly snapshot: SourceWorker["metadata"]; +} +export interface FederationExecution { + readonly mode: "workers" | "in-process"; + readonly sources: readonly FederationSourceMetadata[]; + /** Includes requested virtual membership facts; excludes coordinator catalog facts. */ + readonly rowsTransferred: number; + readonly fragments: readonly { + readonly database: string; + readonly attributes: readonly string[] | null; + readonly rows: number; + }[]; + readonly coordinator: "duckdb"; +} +type AllResult = Effect.Success>; +type WindowResult = Effect.Success>; +export type FederatedWrappedQuery = Omit & { + readonly inner: FederatedQuery; +}; +export interface DuckdbFederationService { + readonly sources: readonly FederationSourceMetadata[]; + readonly explain: (query: FederatedQuery) => Effect.Effect; + readonly queryAll: ( + query: FederatedQuery, + debug?: boolean, + ) => Effect.Effect< + AllResult & { readonly federation: FederationExecution }, + ReadError | DatalogQueryError + >; + readonly queryWindow: ( + query: FederatedWrappedQuery, + debug?: boolean, + ) => Effect.Effect< + WindowResult & { readonly federation: FederationExecution }, + ReadError | DatalogQueryError + >; +} + +const qualify = (database: string, row: FactRow): FactRow => ({ + ...row, + id: sourceEntity(database, String(row["id"])), + entity_id: sourceEntity(database, String(row["entity_id"])), + value_string: + row["value_type"] === "ref" && + row["attribute"] !== DATABASE_ATTRIBUTE && + row["attribute"] !== TX_DATABASE_ATTRIBUTE + ? sourceEntity(database, String(row["value_string"])) + : row["value_string"], + tx_id: row["tx_id"] == null ? null : sourceEntity(database, String(row["tx_id"])), + created_by: row["created_by"] == null ? null : sourceEntity(database, String(row["created_by"])), + // Source visibility is already resolved; a single global temporal cut would be incorrect. + retracted_at: null, + retracted_position: null, + retract_tx_id: null, +}); + +export const makeDuckdbFederation = ( + options: DuckdbFederationOptions = {}, +): Effect.Effect => + Effect.gen(function* () { + const provider = yield* SnapshotProvider; + const catalog = yield* provider.databases; + const mode = options.mode ?? "workers"; + const batchSize = options.batchSize ?? 2048; + const threads = options.threads ?? 4; + const timeoutMs = options.workerTimeoutMs ?? 30_000; + if ( + !["workers", "in-process"].includes(mode) || + !Number.isSafeInteger(batchSize) || + batchSize < 1 || + batchSize > 65_536 || + !Number.isSafeInteger(threads) || + threads < 1 || + threads > 64 || + !Number.isSafeInteger(timeoutMs) || + timeoutMs < 1 || + timeoutMs > 2_147_483_647 + ) + return yield* Effect.fail( + federationFailure( + "Invalid federation mode, batchSize (1..65536), threads (1..64), or positive workerTimeoutMs", + ), + ); + if ( + catalog.length === 0 || + catalog.length > 32 || + new Set(catalog.map((db) => db.id)).size !== catalog.length || + catalog.some((db) => !db.id.trim() || !db.tenant.trim() || !isAbsolute(db.filename)) + ) + return yield* Effect.fail( + federationFailure( + "Expected 1..32 unique database IDs, nonempty tenants, and absolute SQLite filenames", + ), + ); + const workers = new Map(); + for (const database of catalog) { + const worker = yield* mode === "workers" + ? makeProcessSource(database, batchSize, timeoutMs) + : makeSource(database, batchSize); + workers.set(database.id, worker); + } + const sources = Object.freeze( + catalog.map((db): FederationSourceMetadata => { + const worker = workers.get(db.id)!; + return Object.freeze({ + database: db.id, + tenant: db.tenant, + pid: worker.pid, + snapshot: Object.freeze({ + ...worker.metadata, + basis: Object.freeze({ ...worker.metadata.basis }), + }), + }); + }), + ); + const semaphore = yield* Semaphore.make(1); + const explain = (query: FederatedQuery) => + Effect.try({ try: () => planFederation(query, catalog), catch: federationFailure }); + const execute = ( + query: FederatedQuery, + run: ( + executor: QueryExecutorService, + plan: FederationPlan, + ) => Effect.Effect, + ) => + semaphore.withPermits(1)( + Effect.scoped( + Effect.gen(function* () { + const plan = yield* explain(query); + const table = yield* makeTable(threads); + const fragments: { + database: string; + attributes: readonly string[] | null; + rows: number; + }[] = []; + for (const id of plan.databases) { + const worker = workers.get(id)!; + let after: string | null = null; + let count = 0; + while (true) { + const rows: readonly FactRow[] = yield* worker.scan({ + attributes: plan.attributes, + after, + limit: batchSize, + }); + if (rows.length === 0) break; + yield* table.append(rows.map((row) => qualify(id, row))); + count += rows.length; + after = String(rows[rows.length - 1]!["id"]); + if (rows.length < batchSize) break; + } + fragments.push({ database: id, attributes: plan.attributes, rows: count }); + } + // Catalog facts are available even for an empty database. + yield* table.append( + catalog.map((db) => ({ + id: databaseEntity(db.id), + entity_id: databaseEntity(db.id), + attribute: TENANT_ATTRIBUTE, + value_type: "string", + value_string: db.tenant, + recorded_at: 0, + recorded_position: 0, + valid_from: 0, + schema_version: 1, + })), + ); + yield* table.runner.run("ANALYZE triples", []).pipe(Effect.mapError(federationFailure)); + const result = yield* run(makeSqlQueryExecutor(table.runner, DuckdbDialect), plan); + const federation: FederationExecution = { + mode, + sources: sources.filter((source) => plan.databases.includes(source.database)), + rowsTransferred: fragments.reduce((sum, fragment) => sum + fragment.rows, 0), + fragments, + coordinator: "duckdb", + }; + return { ...result, federation }; + }), + ), + ); + return { + sources, + explain, + queryAll: (query, debug = false) => + execute(query, (executor, plan) => executor.execute(plan.query, debug)), + queryWindow: (query, debug = false) => + execute(query.inner, (executor, plan) => + executor.executePage({ ...query, inner: plan.query }, debug), + ), + }; + }); + +export class DuckdbFederation extends Context.Service()( + "triplex/DuckdbFederation", +) { + static layer(options: DuckdbFederationOptions = {}) { + return Layer.effect(DuckdbFederation, makeDuckdbFederation(options)); + } +} diff --git a/packages/duckdb/src/DuckdbSnapshot.ts b/packages/duckdb/src/DuckdbSnapshot.ts new file mode 100644 index 0000000..472a2da --- /dev/null +++ b/packages/duckdb/src/DuckdbSnapshot.ts @@ -0,0 +1,222 @@ +import { Clock, Context, Effect, Layer, type Scope } from "effect"; +import { SqlClient } from "effect/unstable/sql"; +import { DuckDBInstance, type DuckDBValue } from "@duckdb/node-api"; +import { + ReadError, + resolveTemporalBasis, + type DatalogQuery, + type ResolvedTemporalBasis, + type TemporalBasis, + type WrappedQuery, +} from "@bjacobso/triplex"; +import { + makeSqlQueryExecutor, + TRIPLES_TABLE_DDL, + type SqlStatementRunner, +} from "@bjacobso/triplex-sql"; +import { DuckdbDialect } from "./dialect.js"; +import type { QueryExecutorService } from "@bjacobso/triplex/internal"; + +export interface DuckdbSnapshotOptions { + /** Identifies the source database in snapshot metadata. */ + readonly scope: string; + readonly basis?: TemporalBasis; + /** Bounds the rows transferred through JavaScript per batch. Default: 2048. */ + readonly batchSize?: number; + /** Native execution threads. Default: 4. */ + readonly threads?: number; +} + +export interface DuckdbSnapshotService { + readonly metadata: { + readonly scope: string; + readonly basis: ResolvedTemporalBasis; + readonly factCount: number; + readonly threads: number; + readonly importTimeMs: number; + readonly duckdbVersion: string; + }; + readonly queryAll: ( + query: DatalogQuery, + debug?: boolean, + ) => ReturnType; + readonly queryWindow: ( + query: WrappedQuery, + debug?: boolean, + ) => ReturnType; +} + +const readFailure = (cause: unknown) => + new ReadError({ message: `DuckDB snapshot: ${String(cause)}`, cause }); + +const scalar = (value: unknown): DuckDBValue => { + if ( + value === null || + typeof value === "string" || + typeof value === "number" || + typeof value === "bigint" || + typeof value === "boolean" + ) + return value; + throw new Error(`Unsupported SQL parameter: ${typeof value}`); +}; + +/** + * Copy a migrated Triplex SQLite source using one SQL transaction, including + * assertion/retraction history. The private runner also supports federation scans. + * Keep it inside the acquiring Effect scope so native handles remain alive. + */ +export const makeDuckdbSnapshotStore = ( + options: DuckdbSnapshotOptions, +): Effect.Effect< + DuckdbSnapshotService & { readonly runner: SqlStatementRunner }, + ReadError, + SqlClient.SqlClient | Scope.Scope +> => + Effect.gen(function* () { + const source = yield* SqlClient.SqlClient; + const batchSize = options.batchSize ?? 2048; + const threads = options.threads ?? 4; + if ( + !source.onDialectOrElse({ sqlite: () => true, orElse: () => false }) || + options.scope.trim() === "" || + !Number.isSafeInteger(batchSize) || + batchSize < 1 || + batchSize > 65_536 || + !Number.isSafeInteger(threads) || + threads < 1 || + threads > 64 + ) + return yield* Effect.fail( + readFailure("Expected SQLite, a nonempty scope, batchSize 1..65536 and threads 1..64"), + ); + + const now = yield* Clock.currentTimeMillis; + const temporal = yield* Effect.try({ + try: () => resolveTemporalBasis(options.basis, now), + catch: readFailure, + }); + const started = performance.now(); + const instance = yield* Effect.acquireRelease( + Effect.tryPromise({ + try: () => DuckDBInstance.create(":memory:", { threads: String(threads) }), + catch: readFailure, + }), + (db) => Effect.sync(() => db.closeSync()), + ); + const connection = yield* Effect.acquireRelease( + Effect.tryPromise({ try: () => instance.connect(), catch: readFailure }), + (conn) => Effect.sync(() => conn.closeSync()), + ); + // Wait for native operations before scope finalizers close the connection. + const runner: SqlStatementRunner = { + run: >(sql: string, params: readonly unknown[]) => + Effect.tryPromise({ + try: async () => { + const result = await connection.runAndReadAll(sql, params.map(scalar)); + return result.getRowObjects() as Row[]; + }, + catch: readFailure, + }).pipe(Effect.uninterruptible), + }; + yield* runner.run(TRIPLES_TABLE_DDL, []); + const columns = yield* runner.run<{ column_name: string; column_type: string }>( + "DESCRIBE triples", + [], + ); + const columnList = columns.map((column) => `"${column.column_name}"`).join(","); + let factCount = 0; + const recordedPosition = yield* source + .withTransaction( + Effect.gen(function* () { + // This first read establishes the same SQLite snapshot for the counter + // and every subsequent keyset batch, even if another writer commits. + const positions = yield* source.unsafe<{ position: number }>( + "SELECT position FROM triplex_commit_position WHERE singleton = 1", + ); + const position = positions[0]?.position ?? 0; + if (!Number.isSafeInteger(position) || position < 0) + return yield* Effect.fail(readFailure("Invalid source commit position")); + + yield* Effect.scoped( + Effect.gen(function* () { + const appender = yield* Effect.acquireRelease( + Effect.tryPromise({ + try: () => connection.createAppender("triples"), + catch: readFailure, + }), + (handle) => Effect.sync(() => handle.closeSync()), + ); + let lastId: string | undefined; + while (true) { + const rows = yield* source.unsafe>( + `SELECT ${columnList} FROM triples ${lastId === undefined ? "" : "WHERE id > ?1"} ORDER BY id LIMIT ${batchSize}`, + lastId === undefined ? [] : [lastId], + ); + if (rows.length === 0) break; + yield* Effect.try({ + try: () => { + for (const row of rows) { + for (const column of columns) { + const value = row[column.column_name]; + if (value === null) appender.appendNull(); + else if (column.column_type === "BIGINT") + appender.appendBigInt(BigInt(value as number)); + else if (column.column_type === "INTEGER") + appender.appendInteger(value as number); + else if (column.column_type === "DOUBLE") + appender.appendDouble(value as number); + else appender.appendVarchar(value as string); + } + appender.endRow(); + } + appender.flushSync(); + }, + catch: readFailure, + }); + factCount += rows.length; + lastId = rows[rows.length - 1]!["id"] as string; + } + }), + ); + return position; + }), + ) + .pipe(Effect.mapError(readFailure)); + + yield* runner.run("ANALYZE triples", []); + const versions = yield* runner.run<{ version: string }>("SELECT version() AS version", []); + const basis: ResolvedTemporalBasis = Object.freeze({ ...temporal, recordedPosition }); + const executor = makeSqlQueryExecutor(runner, DuckdbDialect); + const metadata = Object.freeze({ + scope: options.scope, + basis, + factCount, + threads, + importTimeMs: performance.now() - started, + duckdbVersion: versions[0]!.version, + }); + return { + runner, + metadata, + /** Complete materialization for trusted analytical workloads. */ + queryAll: (query: DatalogQuery, debug = false) => executor.execute(query, debug, basis), + /** SQL result window; this POC does not issue Triples opaque cursors. */ + queryWindow: (query: WrappedQuery, debug = false) => + executor.executePage(query, debug, basis), + }; + }).pipe(Effect.mapError(readFailure)); + +/** Acquire an immutable analytical copy of the current SQLite source. */ +export const makeDuckdbSnapshot = ( + options: DuckdbSnapshotOptions, +): Effect.Effect => + makeDuckdbSnapshotStore(options).pipe(Effect.map(({ runner: _runner, ...snapshot }) => snapshot)); + +export class DuckdbSnapshot extends Context.Service()( + "triplex/DuckdbSnapshot", +) { + static layer(options: DuckdbSnapshotOptions) { + return Layer.effect(DuckdbSnapshot, makeDuckdbSnapshot(options)); + } +} diff --git a/packages/duckdb/src/dialect.ts b/packages/duckdb/src/dialect.ts new file mode 100644 index 0000000..53dcba1 --- /dev/null +++ b/packages/duckdb/src/dialect.ts @@ -0,0 +1,14 @@ +import type { SqlDialect } from "@bjacobso/triplex/internal"; + +export const DuckdbDialect: SqlDialect = { + name: "duckdb", + limitOffset: (limit, offset) => + [ + ...(limit === undefined ? [] : [`LIMIT ${limit}`]), + ...(offset === undefined || offset === 0 ? [] : [`OFFSET ${offset}`]), + ].join(" "), + castAsText: (expression) => `CAST(${expression} AS VARCHAR)`, + booleanLiteral: (value) => (value ? "TRUE" : "FALSE"), + paramPlaceholder: (index) => `$${index + 1}`, + escapeLikePattern: (value) => value.replace(/[%_\\]/g, "\\$&"), +}; diff --git a/packages/duckdb/src/federation/planner.ts b/packages/duckdb/src/federation/planner.ts new file mode 100644 index 0000000..4d70df6 --- /dev/null +++ b/packages/duckdb/src/federation/planner.ts @@ -0,0 +1,136 @@ +import { assertDatalogQuery } from "@bjacobso/triplex/internal"; +import type { DatalogQuery } from "@bjacobso/triplex"; +import { + DATABASE_ATTRIBUTE, + TENANT_ATTRIBUTE, + databaseEntity, + sourceEntity, + type FederatedQuery, + type LocalDatabase, +} from "./types.js"; + +export interface FederationPlan { + readonly query: DatalogQuery; + readonly databases: readonly string[]; + readonly attributes: readonly string[] | null; +} + +/** Normalize only; all typing, joins, aggregation and recursion stay in the shared compiler. */ +export const planFederation = ( + input: FederatedQuery, + catalog: readonly LocalDatabase[], +): FederationPlan => { + const ids = new Set(catalog.map((db) => db.id)); + const sources = input.sources ?? {}; + for (const [alias, id] of Object.entries(sources)) { + if (!/^\$[a-zA-Z_][a-zA-Z0-9_]*$/.test(alias) || !ids.has(id)) + throw new Error(`Unknown database or invalid source alias: ${alias} -> ${id}`); + } + const expand = (clause: readonly unknown[]): readonly (readonly unknown[])[] => { + if (clause[0] === "not") + return [["not", ...clause.slice(1).flatMap((inner) => expand(inner as readonly unknown[]))]]; + if (typeof clause[0] !== "string" || !clause[0].startsWith("$")) return [clause]; + const id = sources[clause[0] as `$${string}`]; + if (id === undefined) throw new Error(`Unbound source: ${clause[0]}`); + if (clause.length !== 4 && clause.length !== 5) + throw new Error("Source patterns require four or five terms"); + const identity = (value: unknown): string => { + if (typeof value !== "string") throw new Error("Source identities must be strings"); + return value.startsWith("?") ? value : sourceEntity(id, value); + }; + const entity = identity(clause[1]); + const value = clause[3]; + const qualifiedValue = + typeof value === "object" && + value !== null && + "type" in value && + value.type === "ref" && + "value" in value + ? { type: "ref", value: sourceEntity(id, String(value.value)) } + : value; + return [ + [entity, clause[2], qualifiedValue, ...(clause.length === 5 ? [identity(clause[4])] : [])], + [entity, DATABASE_ATTRIBUTE, { type: "ref", value: databaseEntity(id) }], + ]; + }; + const { sources: _sources, ...rest } = input; + const normalized = { ...rest, where: input.where.flatMap(expand) }; + // Explicitly reject shortcuts where expansion would change core rule/OR semantics. + const rejectRemainingSources = (value: unknown): void => { + if (!Array.isArray(value)) return; + if (typeof value[0] === "string" && value[0].startsWith("$")) + throw new Error( + "Source shortcuts are supported in top-level patterns and not; use membership patterns in or and ordinary rules", + ); + value.forEach(rejectRemainingSources); + }; + rejectRemainingSources(normalized.where); + normalized.rules?.forEach((rule) => rejectRemainingSources(rule.body)); + const query = assertDatalogQuery(normalized); + const attributes = new Set(); + let wildcard = false; + const walk = (value: unknown): void => { + if (!Array.isArray(value)) return; + if ( + typeof value[0] === "string" && + value[0] !== "not" && + value[0] !== "or" && + typeof value[1] === "string" + ) { + if (value[1].startsWith(":")) attributes.add(value[1]); + else if (value[1].startsWith("?") && !["=", "!=", ">", ">=", "<", "<="].includes(value[0])) + wildcard = true; + } + value.forEach(walk); + }; + walk(query.where); + query.rules?.forEach((rule) => walk(rule.body)); + query.optionalProjection?.fields.forEach((field) => attributes.add(field.attribute)); + // Conservatively prune only flat conjunctions whose every entity has an explicit source. + const bindings = new Map>(); + const databaseBindings = new Map>(); + for (const clause of query.where) { + if ( + clause[1] === TENANT_ATTRIBUTE && + typeof clause[0] === "string" && + typeof clause[2] === "string" && + !clause[2].startsWith("?") + ) { + const matches = databaseBindings.get(clause[0]) ?? new Set(); + catalog.filter((db) => db.tenant === clause[2]).forEach((db) => matches.add(db.id)); + databaseBindings.set(clause[0], matches); + } + } + for (const clause of query.where) { + if (clause[1] === DATABASE_ATTRIBUTE && typeof clause[0] === "string") { + const value = clause[2]; + const matches = + typeof value === "string" + ? databaseBindings.get(value) + : typeof value === "object" && value !== null && "value" in value + ? new Set( + catalog.filter((db) => databaseEntity(db.id) === value.value).map((db) => db.id), + ) + : undefined; + if (matches !== undefined) { + const bound = bindings.get(clause[0]) ?? new Set(); + matches.forEach((id) => bound.add(id)); + bindings.set(clause[0], bound); + } + } + } + const used = new Set(); + let all = query.rules !== undefined || query.optionalProjection !== undefined; + for (const clause of query.where) { + if (["=", "!=", ">", ">=", "<", "<="].includes(String(clause[0]))) continue; + if (clause[1] === TENANT_ATTRIBUTE) continue; // The coordinator owns the complete catalog. + const bound = bindings.get(String(clause[0])); + if (!bound || typeof clause[1] !== "string" || !clause[1].startsWith(":")) all = true; + else bound.forEach((id) => used.add(id)); + } + return { + query, + databases: catalog.filter((db) => all || used.has(db.id)).map((db) => db.id), + attributes: wildcard ? null : [...attributes].sort(), + }; +}; diff --git a/packages/duckdb/src/federation/source.ts b/packages/duckdb/src/federation/source.ts new file mode 100644 index 0000000..3964214 --- /dev/null +++ b/packages/duckdb/src/federation/source.ts @@ -0,0 +1,122 @@ +import { Effect, type Scope } from "effect"; +import { SqliteClient } from "@effect/sql-sqlite-node"; +import { ReadError } from "@bjacobso/triplex"; +import { makeDuckdbSnapshotStore, type DuckdbSnapshotService } from "../DuckdbSnapshot.js"; +import { + DATABASE_ATTRIBUTE, + TENANT_ATTRIBUTE, + TX_DATABASE_ATTRIBUTE, + databaseEntity, + federationFailure, + type FactRow, + type LocalDatabase, + type ScanRequest, +} from "./types.js"; + +export interface SourceWorker { + readonly metadata: DuckdbSnapshotService["metadata"]; + readonly pid: number; + readonly scan: (request: ScanRequest) => Effect.Effect; +} + +/** The same immutable source implementation runs in-process or behind IPC. */ +export const makeSource = ( + database: LocalDatabase, + batchSize: number, +): Effect.Effect => + Effect.gen(function* () { + const snapshot = yield* makeDuckdbSnapshotStore({ + scope: database.id, + ...(database.basis === undefined ? {} : { basis: database.basis }), + batchSize, + threads: 1, + }).pipe( + Effect.provide( + SqliteClient.layer({ filename: database.filename, readonly: true, disableWAL: true }), + ), + ); + const { runner, metadata } = snapshot; + const basis = metadata.basis; + const params: unknown[] = [basis.recordedPosition, basis.validAt]; + const recorded = basis.recordedAt === undefined ? "" : " AND recorded_at <= $3"; + const retracted = basis.recordedAt === undefined ? "" : " OR retracted_at > $3"; + if (basis.recordedAt !== undefined) params.push(basis.recordedAt); + // Materialize visibility at this source's own commit cut, before federation. + yield* runner.run( + `CREATE TABLE visible AS SELECT * FROM triples WHERE + recorded_position <= $1 ${recorded} + AND (retracted_position IS NULL OR retracted_position > $1 ${retracted}) + AND valid_from <= $2 AND (valid_to IS NULL OR valid_to > $2)`, + params, + ); + const reserved = yield* runner.run( + `SELECT id FROM visible WHERE attribute IN ($1, $2, $3) LIMIT 1`, + [DATABASE_ATTRIBUTE, TENANT_ATTRIBUTE, TX_DATABASE_ATTRIBUTE], + ); + if (reserved.length) + return yield* Effect.fail( + federationFailure("Source contains reserved federation attributes"), + ); + const columns = yield* runner.run<{ column_name: string }>("DESCRIBE triples", []); + // Membership includes referenced identities and transactions, even if the requested + // ordinary attribute is absent. It must not depend on the current scan's projection. + yield* runner.run( + `CREATE VIEW identities AS + SELECT entity_id AS identity FROM visible UNION + SELECT value_string FROM visible WHERE value_type = 'ref' UNION + SELECT tx_id FROM visible WHERE tx_id IS NOT NULL`, + [], + ); + const literal = (value: string) => `'${value.replaceAll("'", "''")}'`; + const virtual = (attribute: string, relation: string, identity: string) => { + const expressions: Record = { + id: `${literal(`v:${attribute}:`)} || ${identity}`, + entity_id: identity, + attribute: literal(attribute), + value_type: "'ref'", + value_string: literal(databaseEntity(database.id)), + recorded_at: "0", + recorded_position: "0", + valid_from: "0", + schema_version: "1", + }; + return `SELECT ${columns.map(({ column_name }) => `${expressions[column_name] ?? "NULL"} AS "${column_name}"`).join(",")} FROM ${relation}`; + }; + yield* runner.run( + `CREATE VIEW federation_facts AS + SELECT ${columns.map(({ column_name }) => (column_name === "id" ? "'f:' || id AS id" : `"${column_name}"`)).join(",")} FROM visible + UNION ALL ${virtual(DATABASE_ATTRIBUTE, "identities", "identity")} + UNION ALL ${virtual(TX_DATABASE_ATTRIBUTE, "(SELECT DISTINCT tx_id FROM visible WHERE tx_id IS NOT NULL)", "tx_id")}`, + [], + ); + return { + metadata, + pid: process.pid, + scan: (request: ScanRequest) => { + const values: unknown[] = []; + const conditions: string[] = []; + if (request.after !== null) { + values.push(request.after); + conditions.push(`id > $${values.length}`); + } + if (request.attributes !== null) { + if (request.attributes.length === 0) return Effect.succeed([]); + conditions.push( + `attribute IN (${request.attributes + .map((attribute) => { + values.push(attribute); + return `$${values.length}`; + }) + .join(",")})`, + ); + } + values.push(request.limit); + return runner + .run( + `SELECT * FROM federation_facts ${conditions.length ? `WHERE ${conditions.join(" AND ")}` : ""} ORDER BY id LIMIT $${values.length}`, + values, + ) + .pipe(Effect.mapError(federationFailure)); + }, + }; + }).pipe(Effect.mapError(federationFailure)); diff --git a/packages/duckdb/src/federation/table.ts b/packages/duckdb/src/federation/table.ts new file mode 100644 index 0000000..bc76b2a --- /dev/null +++ b/packages/duckdb/src/federation/table.ts @@ -0,0 +1,66 @@ +import { Effect } from "effect"; +import { DuckDBInstance, type DuckDBValue } from "@duckdb/node-api"; +import { TRIPLES_TABLE_DDL, type SqlStatementRunner } from "@bjacobso/triplex-sql"; +import { federationFailure, type FactRow } from "./types.js"; + +/** A private native table for one coordinator query, released with its Effect scope. */ +export const makeTable = (threads: number) => + Effect.gen(function* () { + const instance = yield* Effect.acquireRelease( + Effect.tryPromise({ + try: () => DuckDBInstance.create(":memory:", { threads: String(threads) }), + catch: federationFailure, + }), + (db) => Effect.sync(() => db.closeSync()), + ); + const connection = yield* Effect.acquireRelease( + Effect.tryPromise({ try: () => instance.connect(), catch: federationFailure }), + (db) => Effect.sync(() => db.closeSync()), + ); + const runner: SqlStatementRunner = { + run: >(sql: string, parameters: readonly unknown[]) => + Effect.tryPromise({ + try: async () => + ( + await connection.runAndReadAll(sql, parameters as DuckDBValue[]) + ).getRowObjects() as Row[], + catch: federationFailure, + }).pipe(Effect.uninterruptible), + }; + yield* runner.run(TRIPLES_TABLE_DDL, []); + const columns = yield* runner.run<{ column_name: string; column_type: string }>( + "DESCRIBE triples", + [], + ); + const append = (rows: readonly FactRow[]) => + Effect.scoped( + Effect.gen(function* () { + const appender = yield* Effect.acquireRelease( + Effect.tryPromise({ + try: () => connection.createAppender("triples"), + catch: federationFailure, + }), + (handle) => Effect.sync(() => handle.closeSync()), + ); + yield* Effect.try({ + try: () => { + for (const row of rows) { + for (const column of columns) { + const value = row[column.column_name]; + if (value === null || value === undefined) appender.appendNull(); + else if (column.column_type === "BIGINT") + appender.appendBigInt(BigInt(value as bigint)); + else if (column.column_type === "INTEGER") appender.appendInteger(Number(value)); + else if (column.column_type === "DOUBLE") appender.appendDouble(Number(value)); + else appender.appendVarchar(String(value)); + } + appender.endRow(); + } + appender.flushSync(); + }, + catch: federationFailure, + }); + }), + ).pipe(Effect.mapError(federationFailure)); + return { runner, append }; + }).pipe(Effect.mapError(federationFailure)); diff --git a/packages/duckdb/src/federation/types.ts b/packages/duckdb/src/federation/types.ts new file mode 100644 index 0000000..0e939a9 --- /dev/null +++ b/packages/duckdb/src/federation/types.ts @@ -0,0 +1,77 @@ +import { Context, Effect, Layer } from "effect"; +import { ReadError, type DatalogQuery, type TemporalBasis } from "@bjacobso/triplex"; + +export interface LocalDatabase { + readonly id: string; + readonly tenant: string; + readonly filename: string; + readonly basis?: TemporalBasis; +} + +/** Trusted catalog. A future provider can restore a versioned remote snapshot here. */ +export class SnapshotProvider extends Context.Service< + SnapshotProvider, + { + readonly databases: Effect.Effect; + } +>()("triplex/duckdb/SnapshotProvider") { + static local(databases: readonly LocalDatabase[]) { + const catalog = databases.map((database) => + Object.freeze({ + ...database, + ...(database.basis === undefined ? {} : { basis: Object.freeze({ ...database.basis }) }), + }), + ); + return Layer.succeed(SnapshotProvider, { databases: Effect.succeed(Object.freeze(catalog)) }); + } +} + +type Clause = DatalogQuery["where"][number]; +type Term = string | number | boolean | { readonly type: "ref"; readonly value: string }; +export type SourcePattern = readonly [ + source: `$${string}`, + entity: string, + attribute: string, + value: Term, + tx?: string, +]; +export type FederatedClause = + | Clause + | SourcePattern + | readonly ["not", ...(Clause | SourcePattern)[]]; +export type FederatedQuery = Omit & { + /** Aliases reference catalog database IDs, never filenames. */ + readonly sources?: Readonly>; + readonly where: readonly FederatedClause[]; +}; + +export const DATABASE_ATTRIBUTE = ":triplex/database"; +export const TENANT_ATTRIBUTE = ":triplex/tenant"; +export const TX_DATABASE_ATTRIBUTE = ":_tx/database"; + +/** Length-independent, reversible identities. Plain string values are never rewritten. */ +export const sourceEntity = (database: string, entity: string): string => + `triplex:entity:${encodeURIComponent(database)}:${encodeURIComponent(entity)}`; +export const databaseEntity = (database: string): string => + `triplex:database:${encodeURIComponent(database)}`; +export const decodeSourceEntity = ( + identity: string, +): { readonly database: string; readonly entity: string } | undefined => { + const match = /^triplex:entity:([^:]*):([^:]*)$/.exec(identity); + if (!match) return undefined; + try { + return { database: decodeURIComponent(match[1]!), entity: decodeURIComponent(match[2]!) }; + } catch { + return undefined; + } +}; + +export const federationFailure = (cause: unknown): ReadError => + new ReadError({ message: `DuckDB federation: ${String(cause)}`, cause }); + +export type FactRow = Record; +export interface ScanRequest { + readonly attributes: readonly string[] | null; + readonly after: string | null; + readonly limit: number; +} diff --git a/packages/duckdb/src/federation/worker-client.ts b/packages/duckdb/src/federation/worker-client.ts new file mode 100644 index 0000000..2e48846 --- /dev/null +++ b/packages/duckdb/src/federation/worker-client.ts @@ -0,0 +1,119 @@ +import { fork } from "node:child_process"; +import { Effect } from "effect"; +import { federationFailure, type FactRow, type LocalDatabase, type ScanRequest } from "./types.js"; +import type { SourceWorker } from "./source.js"; +import type { WorkerRequest, WorkerResponse } from "./worker-process.js"; + +export const makeProcessSource = (database: LocalDatabase, batchSize: number, timeoutMs: number) => + Effect.gen(function* () { + const client = yield* Effect.acquireRelease( + Effect.try({ + try: () => { + const child = fork(new URL("./worker-process.js", import.meta.url), [], { + execArgv: [], + serialization: "advanced", + stdio: ["ignore", "ignore", "pipe", "ipc"], + }); + let stderr = ""; + child.stderr?.on("data", (data) => { + stderr = (stderr + String(data)).slice(-8192); + }); + let sequence = 0; + let failure: Error | undefined; + const pending = new Map< + number, + { resolve: (value: unknown) => void; reject: (error: Error) => void } + >(); + const fail = (error: Error) => { + failure ??= error; + for (const request of pending.values()) request.reject(failure); + pending.clear(); + }; + child.on("error", fail); + child.on("exit", (code, signal) => + fail(new Error(`Worker ${database.id} exited (${code ?? signal}): ${stderr}`)), + ); + child.on("message", (message: WorkerResponse) => { + const request = pending.get(message.id); + if (!request) return; + pending.delete(message.id); + if (message.ok) request.resolve(message.value); + else request.reject(new Error(message.error)); + }); + type Command = WorkerRequest extends infer R + ? R extends WorkerRequest + ? Omit + : never + : never; + const request = (command: Command) => + Effect.tryPromise({ + try: (signal) => + new Promise((resolve, reject) => { + if (failure) { + reject(failure); + return; + } + const id = ++sequence; + const abort = () => { + fail(new Error(`Worker ${database.id} request aborted`)); + child.kill(); + }; + const timer = setTimeout(() => { + fail(new Error(`Worker ${database.id} timed out after ${timeoutMs}ms`)); + child.kill(); + }, timeoutMs); + const cleanup = () => { + clearTimeout(timer); + signal.removeEventListener("abort", abort); + pending.delete(id); + }; + pending.set(id, { + resolve: (value) => { + cleanup(); + resolve(value as A); + }, + reject: (error) => { + cleanup(); + reject(error); + }, + }); + signal.addEventListener("abort", abort, { once: true }); + child.send({ ...command, id }, (error) => { + if (error) fail(error); + }); + }), + catch: federationFailure, + }); + return { child, request }; + }, + catch: federationFailure, + }), + ({ child }) => + Effect.promise( + () => + new Promise((resolve) => { + if (child.pid === undefined || child.exitCode !== null || child.signalCode !== null) { + resolve(); + return; + } + const timer = setTimeout(() => child.kill("SIGKILL"), 1000); + child.once("exit", () => { + clearTimeout(timer); + resolve(); + }); + if (child.connected) child.disconnect(); + else child.kill(); + }), + ), + ); + const opened = yield* client.request>({ + command: "open", + database, + batchSize, + }); + return { + ...opened, + scan: (request: ScanRequest) => + client.request({ command: "scan", request }), + } satisfies SourceWorker; + }); diff --git a/packages/duckdb/src/federation/worker-process.ts b/packages/duckdb/src/federation/worker-process.ts new file mode 100644 index 0000000..df8c157 --- /dev/null +++ b/packages/duckdb/src/federation/worker-process.ts @@ -0,0 +1,93 @@ +import { Effect } from "effect"; +import { makeSource } from "./source.js"; +import type { LocalDatabase, ScanRequest } from "./types.js"; + +export type WorkerRequest = + | { + readonly id: number; + readonly command: "open"; + readonly database: LocalDatabase; + readonly batchSize: number; + } + | { readonly id: number; readonly command: "scan"; readonly request: ScanRequest }; +export type WorkerResponse = + | { readonly id: number; readonly ok: true; readonly value: unknown } + | { readonly id: number; readonly ok: false; readonly error: string }; + +// This file is an executable package artifact, launched only through fork(). +const queue: WorkerRequest[] = []; +let wake: ((message: WorkerRequest | null) => void) | undefined; +let disconnected = false; +process.on("message", (message: WorkerRequest) => { + if (wake) { + const resolve = wake; + wake = undefined; + resolve(message); + } else queue.push(message); +}); +process.on("disconnect", () => { + disconnected = true; + wake?.(null); + wake = undefined; +}); +const next = () => + Effect.promise(() => { + const message = queue.shift(); + return message + ? Promise.resolve(message) + : disconnected + ? Promise.resolve(null) + : new Promise((resolve) => { + wake = resolve; + }); + }); +const send = (response: WorkerResponse) => + Effect.promise( + () => + new Promise((resolve, reject) => { + if (!process.connected || !process.send) { + resolve(); + return; + } + process.send(response, (error) => (error ? reject(error) : resolve())); + }), + ); + +const main = Effect.scoped( + Effect.gen(function* () { + const open = yield* next(); + if (!open || open.command !== "open") return; + const source = yield* makeSource(open.database, open.batchSize).pipe( + Effect.catch((error) => + send({ id: open.id, ok: false, error: String(error) }).pipe( + Effect.andThen(Effect.fail(error)), + ), + ), + ); + yield* send({ id: open.id, ok: true, value: { metadata: source.metadata, pid: source.pid } }); + while (true) { + const message = yield* next(); + if (message === null) break; + if (message.command !== "scan") { + yield* send({ id: message.id, ok: false, error: "Worker is already open" }); + continue; + } + const response = yield* source.scan(message.request).pipe( + Effect.map((value): WorkerResponse => ({ id: message.id, ok: true, value })), + Effect.catch((error) => + Effect.succeed({ id: message.id, ok: false, error: String(error) }), + ), + ); + yield* send(response); + } + }), +); +Effect.runPromise(main).then( + () => { + if (process.connected) process.disconnect?.(); + }, + () => { + process.exitCode = 1; + if (process.connected) process.disconnect?.(); + }, +); diff --git a/packages/duckdb/src/index.ts b/packages/duckdb/src/index.ts new file mode 100644 index 0000000..92e2531 --- /dev/null +++ b/packages/duckdb/src/index.ts @@ -0,0 +1,30 @@ +export { + DuckdbSnapshot, + makeDuckdbSnapshot, + type DuckdbSnapshotOptions, + type DuckdbSnapshotService, +} from "./DuckdbSnapshot.js"; +export { DuckdbDialect } from "./dialect.js"; +export { + DuckdbFederation, + makeDuckdbFederation, + type DuckdbFederationOptions, + type DuckdbFederationService, + type FederationExecution, + type FederationSourceMetadata, + type FederatedWrappedQuery, +} from "./DuckdbFederation.js"; +export { + SnapshotProvider, + sourceEntity, + decodeSourceEntity, + databaseEntity, + DATABASE_ATTRIBUTE, + TENANT_ATTRIBUTE, + TX_DATABASE_ATTRIBUTE, + type LocalDatabase, + type FederatedQuery, + type FederatedClause, + type SourcePattern, +} from "./federation/types.js"; +export type { FederationPlan } from "./federation/planner.js"; diff --git a/packages/duckdb/test/federation.test.ts b/packages/duckdb/test/federation.test.ts new file mode 100644 index 0000000..4d2c60b --- /dev/null +++ b/packages/duckdb/test/federation.test.ts @@ -0,0 +1,397 @@ +import { afterEach, describe, expect, it } from "vitest"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { DatabaseSync } from "node:sqlite"; +import { Effect, Result } from "effect"; +import { TRIPLES_TABLE_DDL, COMMIT_POSITION_TABLE_DDL } from "@bjacobso/triplex-sql"; +// Exercise the built worker entry and public API, including its package-relative fork path. +import { + makeDuckdbFederation, + SnapshotProvider, + sourceEntity, + decodeSourceEntity, + databaseEntity, + type FederatedQuery, + type LocalDatabase, +} from "../dist/index.js"; + +const directories: string[] = []; +afterEach(() => { + for (const directory of directories.splice(0)) + rmSync(directory, { recursive: true, force: true }); +}); +const fixture = (): readonly LocalDatabase[] => { + const directory = mkdtempSync(join(tmpdir(), "triplex-federation-")); + directories.push(directory); + return ["a", "b"].map((id) => { + const filename = join(directory, `${id}.sqlite`); + const db = new DatabaseSync(filename); + db.exec( + `${TRIPLES_TABLE_DDL}; ${COMMIT_POSITION_TABLE_DDL}; INSERT INTO triplex_commit_position VALUES (1, 7)`, + ); + const insert = + db.prepare(`INSERT INTO triples (id, entity_id, attribute, value_type, value_string, value_number, + recorded_at, recorded_position, valid_from, tx_id) VALUES (?, ?, ?, ?, ?, ?, 10, 1, 0, 'tx:1')`); + insert.run("1", "person:1", ":person/email", "string", "shared@example.com", null); + insert.run("2", "person:1", ":person/name", "string", id === "a" ? "Alice" : "Bob", null); + insert.run("3", "person:1", ":person/score", "number", null, id === "a" ? 7 : 11); + insert.run("4", "person:1", ":person/friend", "ref", "person:2", null); + insert.run("5", "person:2", ":person/name", "string", id === "a" ? "Ann" : "Ben", null); + insert.run("6", "person:2", ":person/friend", "ref", "person:1", null); + insert.run("7", "person:1", ":history", "string", "old", null); + db.exec("UPDATE triples SET retracted_at = 20, retracted_position = 2 WHERE id = '7'"); + insert.run("8", "person:1", ":future", "string", "future", null); + db.exec("UPDATE triples SET valid_from = 200 WHERE id = '8'"); + insert.run( + "9", + "person:1", + ":typed", + id === "a" ? "string" : "number", + id === "a" ? "7" : null, + id === "a" ? null : 7, + ); + db.close(); + return { id, tenant: `tenant-${id}`, filename, basis: { validAt: 100 } }; + }); +}; +const joinQuery: FederatedQuery = { + sources: { $a: "a", $b: "b" }, + find: ["?a", "?b", "?email"], + where: [ + ["$a", "?a", ":person/email", "?email"], + ["$b", "?b", ":person/email", "?email"], + ], +}; +const membershipQuery: FederatedQuery = { + find: ["?a", "?b", "?email"], + where: [ + ["?dbA", ":triplex/tenant", "tenant-a"], + ["?a", ":triplex/database", "?dbA"], + ["?dbB", ":triplex/tenant", "tenant-b"], + ["?b", ":triplex/database", "?dbB"], + ["?a", ":person/email", "?email"], + ["?b", ":person/email", "?email"], + ], +}; +const normalize = (rows: readonly unknown[]) => rows.map((row) => JSON.stringify(row)).sort(); +const run = ( + catalog: readonly LocalDatabase[], + effect: Effect.Effect, +) => Effect.runPromise(effect.pipe(Effect.provide(SnapshotProvider.local(catalog)))); + +describe("local DuckDB federation", () => { + it("joins two worker snapshots with either syntax and matches the in-process reference", async () => { + const catalog = fixture(); + let pids: readonly number[] = []; + await run( + catalog, + Effect.scoped( + Effect.gen(function* () { + const distributed = yield* makeDuckdbFederation({ batchSize: 2 }); + const local = yield* makeDuckdbFederation({ mode: "in-process", batchSize: 3 }); + pids = distributed.sources.map((source) => source.pid); + expect(new Set(pids).size).toBe(2); + expect(pids).not.toContain(process.pid); + const result = yield* distributed.queryAll(joinQuery, true); + expect(result.results).toEqual([ + { + "?a": sourceEntity("a", "person:1"), + "?b": sourceEntity("b", "person:1"), + "?email": "shared@example.com", + }, + ]); + expect(result.debug?.generatedSql).toContain("SELECT"); + expect((yield* distributed.queryAll(membershipQuery)).results).toEqual(result.results); + const corpus: readonly FederatedQuery[] = [ + joinQuery, + membershipQuery, + { + find: ["?sum"], + where: [["?p", ":person/score", "?score"]], + aggregate: [["sum", "?score", "?sum"]], + }, + { + sources: { $a: "a" }, + find: ["?p"], + where: [ + ["$a", "?p", ":person/name", "?n"], + ["not", ["$a", "?p", ":person/email", "?email"]], + ], + }, + { + find: ["?p"], + where: [ + ["?p", ":person/name", "?n"], + [ + "or", + [ + ["?p", ":person/name", "Alice"], + ["?p", ":person/name", "Ben"], + ], + ], + ], + }, + { + find: ["?to"], + where: [["reach", sourceEntity("a", "person:1"), "?to"]], + rules: [ + { name: "reach", body: [["?x", ":person/friend", "?y"]], maxDepth: 3 }, + { + name: "reach", + body: [ + ["?x", ":person/friend", "?z"], + ["reach", "?z", "?y"], + ], + maxDepth: 3, + }, + ], + }, + ]; + for (const query of corpus) + expect(normalize((yield* distributed.queryAll(query)).results)).toEqual( + normalize((yield* local.queryAll(query)).results), + ); + const total = yield* distributed.queryAll(corpus[2]!); + expect(total.results).toEqual([{ "?sum": 18 }]); + const reach = yield* distributed.queryAll(corpus[5]!); + expect(reach.results).toHaveLength(2); + expect( + reach.results.every((row) => decodeSourceEntity(String(row["?to"]))?.database === "a"), + ).toBe(true); + const window = yield* distributed.queryWindow({ + inner: membershipQuery, + orderBy: [{ variable: "?email", direction: "asc" }], + limit: 1, + includeCount: true, + }); + expect(window.results).toEqual(result.results); + }), + ), + ); + for (const pid of pids) expect(() => process.kill(pid, 0)).toThrow(); + }, 30_000); + + it( + "qualifies entity/ref/transaction identities and preserves typed scalar joins", + () => + run( + fixture(), + Effect.scoped( + Effect.gen(function* () { + const db = yield* makeDuckdbFederation({ batchSize: 1 }); + const collision = yield* db.queryAll({ + ...joinQuery, + find: ["?p"], + where: [ + ["$a", "?p", ":person/name", "?a"], + ["$b", "?p", ":person/name", "?b"], + ], + }); + expect(collision.results).toEqual([]); + const typed = yield* db.queryAll({ + ...joinQuery, + where: [ + ["$a", "?a", ":typed", "?email"], + ["$b", "?b", ":typed", "?email"], + ], + }); + expect(typed.results).toEqual([]); + const refs = yield* db.queryAll({ + sources: { $a: "a" }, + find: ["?friend"], + where: [["$a", "person:1", ":person/friend", "?friend"]], + }); + expect(refs.results).toEqual([{ "?friend": sourceEntity("a", "person:2") }]); + const tx = yield* db.queryAll({ + find: ["?tx", "?tenant"], + where: [ + ["?p", ":person/email", "?email", "?tx"], + ["?tx", ":_tx/database", "?db"], + ["?db", ":triplex/tenant", "?tenant"], + ], + }); + expect(normalize(tx.results)).toEqual( + normalize([ + { "?tx": sourceEntity("a", "tx:1"), "?tenant": "tenant-a" }, + { "?tx": sourceEntity("b", "tx:1"), "?tenant": "tenant-b" }, + ]), + ); + expect(decodeSourceEntity(sourceEntity("a:b%", "p:/%"))).toEqual({ + database: "a:b%", + entity: "p:/%", + }); + expect(sourceEntity("a", "b:c")).not.toBe(sourceEntity("a:b", "c")); + }), + ), + ), + 30_000, + ); + + it("pins per-source temporal cuts across later writes and filters transfer by source and attribute", () => { + const catalog = fixture(); + return run( + catalog, + Effect.scoped( + Effect.gen(function* () { + const db = yield* makeDuckdbFederation(); + const query: FederatedQuery = { + sources: { $a: "a" }, + find: ["?p", "?name"], + where: [["$a", "?p", ":person/name", "?name"]], + }; + const before = yield* db.queryAll(query); + expect(before.federation.fragments.map((fragment) => fragment.database)).toEqual(["a"]); + expect(before.federation.fragments[0]?.attributes).toEqual([ + ":person/name", + ":triplex/database", + ]); + expect(before.federation.rowsTransferred).toBe(5); // 2 names, 2 people + 1 transaction memberships + expect(db.sources.map((source) => source.snapshot.basis.recordedPosition)).toEqual([ + 7, 7, + ]); + yield* Effect.sync(() => { + const writer = new DatabaseSync(catalog[0]!.filename); + writer.exec( + "BEGIN; UPDATE triples SET value_string = 'Changed', recorded_position = 8 WHERE id = '2'; UPDATE triplex_commit_position SET position = 8; COMMIT", + ); + writer.close(); + }); + expect(normalize((yield* db.queryAll(query)).results)).toEqual(normalize(before.results)); + expect( + (yield* db.queryAll({ find: ["?p"], where: [["?p", ":future", "?v"]] })).results, + ).toEqual([]); + expect( + (yield* db.queryAll({ find: ["?p"], where: [["?p", ":history", "?v"]] })).results, + ).toEqual([]); + const next = yield* makeDuckdbFederation({ mode: "in-process" }); + expect((yield* next.queryAll(query)).results).toContainEqual({ + "?p": sourceEntity("a", "person:1"), + "?name": "Changed", + }); + expect(next.sources[0]?.snapshot.basis.recordedPosition).toBe(8); + }), + ), + ); + }, 30_000); + + it("supports distinct historical bases, metadata-only queries, and variable attributes", () => { + const catalog = fixture().map((db) => + db.id === "a" ? { ...db, basis: { validAt: 100, recordedAt: 15 } } : db, + ); + return run( + catalog, + Effect.scoped( + Effect.gen(function* () { + const db = yield* makeDuckdbFederation(); + expect( + (yield* db.queryAll({ find: ["?p"], where: [["?p", ":history", "old"]] })).results, + ).toEqual([{ "?p": sourceEntity("a", "person:1") }]); + const members = yield* db.queryAll({ + find: ["?p"], + where: [["?p", ":triplex/database", { type: "ref", value: databaseEntity("a") }]], + }); + expect(members.results).toHaveLength(3); + const wildcard = yield* db.queryAll({ + sources: { $a: "a" }, + find: ["?attr"], + where: [["$a", "person:1", "?attr", "?value"]], + }); + expect(wildcard.federation.fragments[0]?.attributes).toBeNull(); + expect(wildcard.results).toContainEqual({ "?attr": ":history" }); + expect( + (yield* db.queryAll({ find: ["?db"], where: [["?db", ":triplex/tenant", "tenant-b"]] })) + .results, + ).toEqual([{ "?db": databaseEntity("b") }]); + }), + ), + ); + }, 30_000); + + it( + "keeps unscoped and optional facts and prunes an explicit tenant selector", + () => + run( + fixture(), + Effect.scoped( + Effect.gen(function* () { + const db = yield* makeDuckdbFederation(); + const mixed = yield* db.queryAll({ + sources: { $a: "a" }, + find: ["?name"], + where: [ + ["$a", "?p", ":person/email", "?email"], + ["?other", ":person/email", "?email"], + ["?other", ":person/name", "?name"], + ], + }); + expect(normalize(mixed.results)).toEqual( + normalize([{ "?name": "Alice" }, { "?name": "Bob" }]), + ); + expect(mixed.federation.fragments).toHaveLength(2); + const optional = yield* db.queryAll({ + find: ["?name"], + where: [["?p", ":person/email", "?email"]], + optionalProjection: { + rowBinding: "?p", + fields: [{ attribute: ":person/name", variable: "?name" }], + }, + }); + expect(normalize(optional.results)).toEqual(normalize(mixed.results)); + const selected = yield* db.queryAll({ + find: ["?name"], + where: [ + ["?db", ":triplex/tenant", "tenant-a"], + ["?p", ":triplex/database", "?db"], + ["?p", ":person/name", "?name"], + ], + }); + expect(selected.federation.fragments.map((fragment) => fragment.database)).toEqual([ + "a", + ]); + expect(normalize(selected.results)).toEqual( + normalize([{ "?name": "Alice" }, { "?name": "Ann" }]), + ); + }), + ), + ), + 30_000, + ); + + it("rejects unknown aliases, reserved data, and unreadable sources; reports worker exits", async () => { + const catalog = fixture(); + await run( + catalog, + Effect.scoped( + Effect.gen(function* () { + const db = yield* makeDuckdbFederation(); + expect( + Result.isFailure( + yield* Effect.result( + db.queryAll({ ...joinQuery, sources: { $a: "missing", $b: "b" } }), + ), + ), + ).toBe(true); + yield* Effect.sync(() => process.kill(db.sources[0]!.pid, "SIGKILL")); + const result = yield* Effect.result(db.queryAll(joinQuery)); + expect(Result.isFailure(result)).toBe(true); + }), + ), + ); + await expect( + run( + [{ ...catalog[0]!, filename: join(directories[0]!, "missing.sqlite") }], + Effect.scoped(makeDuckdbFederation()), + ), + ).rejects.toThrow(); + await expect( + run(catalog, Effect.scoped(makeDuckdbFederation({ workerTimeoutMs: 1 }))), + ).rejects.toThrow(/timed out/); + const writer = new DatabaseSync(catalog[0]!.filename); + writer.exec("UPDATE triples SET attribute = ':triplex/database' WHERE id = '1'"); + writer.close(); + await expect( + run(catalog, Effect.scoped(makeDuckdbFederation({ mode: "in-process" }))), + ).rejects.toThrow(/reserved/); + }, 30_000); +}); diff --git a/packages/duckdb/test/snapshot.test.ts b/packages/duckdb/test/snapshot.test.ts new file mode 100644 index 0000000..ba7ffd4 --- /dev/null +++ b/packages/duckdb/test/snapshot.test.ts @@ -0,0 +1,321 @@ +import { describe, expect, it } from "vitest"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { Effect } from "effect"; +import { SqlClient } from "effect/unstable/sql"; +import { + EntityId, + Triples, + string, + number, + boolean, + datetime, + json, + ref, + type DatalogQuery, + type WrappedQuery, +} from "@bjacobso/triplex"; +import { QueryExecutor } from "@bjacobso/triplex/internal"; +import { SqliteTriples } from "@bjacobso/triplex-sqlite"; +import { DuckdbSnapshot, makeDuckdbSnapshot } from "../src/index.js"; + +const names: DatalogQuery = { find: ["?e", "?name"], where: [["?e", ":name", "?name"]] }; +const normalize = (rows: readonly unknown[]) => rows.map((row) => JSON.stringify(row)).sort(); +const run = (program: Effect.Effect) => + Effect.runPromise(program.pipe(Effect.provide(SqliteTriples.layerMemory))); + +const seed = Effect.gen(function* () { + const triples = yield* Triples; + yield* triples.assertBatch([ + { entityId: EntityId.make("a"), attribute: ":name", value: string("Alice") }, + { entityId: EntityId.make("b"), attribute: ":name", value: string("Bob") }, + { entityId: EntityId.make("c"), attribute: ":name", value: string("Carol") }, + { entityId: EntityId.make("a"), attribute: ":score", value: number(7.25) }, + { entityId: EntityId.make("b"), attribute: ":score", value: number(12) }, + { entityId: EntityId.make("a"), attribute: ":active", value: boolean(true) }, + { entityId: EntityId.make("b"), attribute: ":active", value: boolean(false) }, + { entityId: EntityId.make("a"), attribute: ":date", value: datetime(1_700_000_000_123) }, + { entityId: EntityId.make("a"), attribute: ":json", value: json({ test: true }) }, + { entityId: EntityId.make("a"), attribute: ":parent", value: ref(EntityId.make("b")) }, + { entityId: EntityId.make("b"), attribute: ":parent", value: ref(EntityId.make("c")) }, + { entityId: EntityId.make("c"), attribute: ":parent", value: ref(EntityId.make("a")) }, + { entityId: EntityId.make("a"), attribute: ":mixed", value: string("7") }, + { entityId: EntityId.make("b"), attribute: ":mixed", value: number(7) }, + { entityId: EntityId.make("c"), attribute: ":mixed", value: boolean(true) }, + ]); +}); + +const recursive: DatalogQuery = { + find: ["?descendant"], + where: [["reach", "a", "?descendant"]], + rules: [ + { name: "reach", body: [["?x", ":parent", "?y"]], maxDepth: 5 }, + { + name: "reach", + body: [ + ["?x", ":parent", "?z"], + ["reach", "?z", "?y"], + ], + maxDepth: 5, + }, + ], +}; +const corpus: readonly DatalogQuery[] = [ + names, + { + find: ["?name", "?score"], + where: [ + [">", "?score", 8], + ["?e", ":name", "?name"], + ["?e", ":score", "?score"], + ], + }, + { + find: ["?name"], + where: [ + ["?e", ":name", "?name"], + ["not", ["?e", ":score", "?score"]], + ], + }, + { + find: ["?e"], + where: [ + ["?e", ":name", "?name"], + [ + "or", + [ + ["?e", ":name", "Alice"], + ["?e", ":name", "Bob"], + ], + ], + ], + }, + { find: ["?v"], where: [["a", ":active", "?v"]] }, + { find: ["?v"], where: [["a", ":date", "?v"]] }, + { find: ["?v"], where: [["a", ":json", "?v"]] }, + { find: ["?v"], where: [["a", ":parent", "?v"]] }, + { + find: ["?v"], + where: [["?e", ":mixed", "?v"]], + orderBy: [{ variable: "?v", direction: "asc" }], + }, + { + find: ["?active", "?count"], + where: [["?e", ":active", "?active"]], + aggregate: [["count", "?e", "?count"]], + }, + { find: ["?sum"], where: [["?e", ":score", "?score"]], aggregate: [["sum", "?score", "?sum"]] }, + { find: ["?count"], where: [["?e", ":missing", "?v"]], aggregate: [["count", "?e", "?count"]] }, + { find: ["?name", 42, true], where: [["a", ":name", "?name"]] }, + recursive, +]; + +describe("DuckDB analytical snapshot", () => { + it("matches SQLite for typed joins, negation, disjunction, aggregation, and cyclic recursion", () => + run( + Effect.scoped( + Effect.gen(function* () { + yield* seed; + const sqlite = yield* QueryExecutor; + const snapshot = yield* makeDuckdbSnapshot({ scope: "test", batchSize: 3 }); + expect(snapshot.metadata.factCount).toBeGreaterThan(15); + expect(snapshot.metadata.basis.recordedPosition).toBeGreaterThan(0); + for (const query of corpus) { + const expected = yield* sqlite.execute(query, false, snapshot.metadata.basis); + const actual = yield* snapshot.queryAll(query, true); + expect(normalize(actual.results), JSON.stringify(query)).toEqual( + normalize(expected.results), + ); + expect(actual.debug?.queryPlan?.backend).toBe("duckdb"); + } + }), + ), + )); + + it("matches ordered windows, counts, boolean filters, and ILIKE", () => + run( + Effect.scoped( + Effect.gen(function* () { + yield* seed; + const sqlite = yield* QueryExecutor; + const snapshot = yield* makeDuckdbSnapshot({ scope: "test" }); + const windows: WrappedQuery[] = [ + { + inner: { ...names, offset: 1, orderBy: [{ variable: "?name", direction: "asc" }] }, + limit: 1, + includeCount: true, + orderBy: [{ variable: "?name", direction: "asc" }], + }, + { + inner: names, + filters: [{ column: "?name", op: "ilike", value: "a%" }], + includeCount: true, + }, + { inner: names, filters: [{ column: "?name", op: "not-ilike", value: "a%" }] }, + { + inner: { find: ["?active"], where: [["?e", ":active", "?active"]] }, + filters: [{ column: "?active", op: "=", value: true }], + }, + { + inner: recursive, + limit: 2, + includeCount: true, + orderBy: [{ variable: "?descendant", direction: "asc" }], + }, + ]; + for (const query of windows) { + const expected = yield* sqlite.executePage(query, false, snapshot.metadata.basis); + const actual = yield* snapshot.queryWindow(query, true); + expect(actual.results).toEqual(expected.results); + expect(actual.totalCount).toEqual(expected.totalCount); + } + }), + ), + )); + + it("keeps old facts after source retraction and excludes subsequent assertions", () => + run( + Effect.scoped( + Effect.gen(function* () { + const triples = yield* Triples; + const alice = yield* triples.assert({ + entityId: EntityId.make("a"), + attribute: ":name", + value: string("Alice"), + }); + const snapshot = yield* makeDuckdbSnapshot({ scope: "test", batchSize: 1 }); + yield* triples.retract(alice.id); + yield* triples.assert({ + entityId: EntityId.make("b"), + attribute: ":name", + value: string("Bob"), + }); + expect((yield* snapshot.queryAll(names)).results).toEqual([ + { "?e": "a", "?name": "Alice" }, + ]); + const next = yield* makeDuckdbSnapshot({ scope: "test" }); + expect(next.metadata.basis.recordedPosition).toBeGreaterThan( + snapshot.metadata.basis.recordedPosition!, + ); + expect((yield* next.queryAll(names)).results).toEqual([{ "?e": "b", "?name": "Bob" }]); + }), + ), + )); + + it("preserves recorded history and valid-time intervals", () => + run( + Effect.scoped( + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* sql.unsafe(`INSERT INTO triples + (id, entity_id, attribute, value_type, value_string, recorded_at, recorded_position, + retracted_at, retracted_position, valid_from, valid_to) + VALUES ('old','a',':name','string','Old',10,1,20,2,100,200), + ('new','a',':name','string','New',20,2,NULL,NULL,100,NULL)`); + yield* sql.unsafe("INSERT INTO triplex_commit_position VALUES (1,2)"); + for (const [basis, expected] of [ + [{ recordedAt: 15, validAt: 150 }, "Old"], + [{ recordedAt: 25, validAt: 150 }, "New"], + [{ recordedAt: 15, validAt: 200 }, undefined], + ] as const) { + yield* Effect.scoped( + Effect.gen(function* () { + const snapshot = yield* makeDuckdbSnapshot({ scope: "history", basis }); + const result = yield* snapshot.queryAll(names); + expect(result.results).toEqual( + expected === undefined ? [] : [{ "?e": "a", "?name": expected }], + ); + }), + ); + } + }), + ), + )); + + it( + "holds one SQLite read snapshot when another connection commits between batches", + () => + Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const directory = yield* Effect.acquireRelease( + Effect.tryPromise(() => mkdtemp(join(tmpdir(), "triplex-duckdb-"))), + (path) => Effect.promise(() => rm(path, { recursive: true, force: true })), + ); + const filename = join(directory, "source.sqlite"); + yield* Effect.scoped( + Effect.gen(function* () { + yield* seed; + const source = yield* SqlClient.SqlClient; + let wrote = false; + const writer = Effect.gen(function* () { + const triples = yield* Triples; + yield* triples.retractByPattern({ attribute: ":name" }); + yield* triples.assert({ + entityId: EntityId.make("d"), + attribute: ":name", + value: string("Dora"), + }); + }).pipe(Effect.provide(SqliteTriples.layer({ filename }))); + const intercepted = new Proxy(source, { + get(target, property, receiver) { + if (property !== "unsafe") return Reflect.get(target, property, receiver); + return (statement: string, params: readonly unknown[] = []) => + target.unsafe(statement, [...params]).pipe( + Effect.tap(() => { + if (wrote || !statement.startsWith('SELECT "id"')) return Effect.void; + wrote = true; + return writer; + }), + ); + }, + }); + const snapshot = yield* makeDuckdbSnapshot({ scope: filename, batchSize: 1 }).pipe( + Effect.provideService(SqlClient.SqlClient, intercepted), + ); + expect(wrote).toBe(true); + expect(normalize((yield* snapshot.queryAll(names)).results)).toEqual( + normalize([ + { "?e": "a", "?name": "Alice" }, + { "?e": "b", "?name": "Bob" }, + { "?e": "c", "?name": "Carol" }, + ]), + ); + const triples = yield* Triples; + expect((yield* triples.queryAll(names)).results).toEqual([ + { "?e": "d", "?name": "Dora" }, + ]); + }), + ).pipe(Effect.provide(SqliteTriples.layer({ filename }))); + }), + ), + ), + 15_000, + ); + + it("provides an empty snapshot through its Effect layer", () => + run( + Effect.gen(function* () { + const result = yield* Effect.gen(function* () { + const snapshot = yield* DuckdbSnapshot; + expect(snapshot.metadata.factCount).toBe(0); + return yield* snapshot.queryAll(names); + }).pipe(Effect.provide(DuckdbSnapshot.layer({ scope: "empty" }))); + expect(result.results).toEqual([]); + }), + )); + + it("returns typed failures for invalid configuration and queries", () => + run( + Effect.scoped( + Effect.gen(function* () { + const invalid = yield* Effect.result(makeDuckdbSnapshot({ scope: "test", batchSize: 0 })); + expect(invalid._tag).toBe("Failure"); + const snapshot = yield* makeDuckdbSnapshot({ scope: "test" }); + const result = yield* Effect.result(snapshot.queryAll({ find: ["?unbound"], where: [] })); + expect(result._tag).toBe("Failure"); + }), + ), + )); +}); diff --git a/packages/duckdb/tsconfig.check.json b/packages/duckdb/tsconfig.check.json new file mode 100644 index 0000000..cf96e88 --- /dev/null +++ b/packages/duckdb/tsconfig.check.json @@ -0,0 +1,5 @@ +{ + "extends": "../../tsconfig.node.json", + "include": ["src/**/*.ts", "test/**/*.ts", "scripts/**/*.ts"], + "exclude": ["node_modules", "dist"] +} diff --git a/packages/duckdb/tsconfig.json b/packages/duckdb/tsconfig.json new file mode 100644 index 0000000..fa20aed --- /dev/null +++ b/packages/duckdb/tsconfig.json @@ -0,0 +1,9 @@ +{ + "extends": "../../tsconfig.node.json", + "compilerOptions": { + "outDir": "dist", + "rootDir": "src" + }, + "include": ["src/**/*"], + "exclude": ["node_modules", "dist", "test"] +} diff --git a/packages/duckdb/tsdown.config.ts b/packages/duckdb/tsdown.config.ts new file mode 100644 index 0000000..0e634be --- /dev/null +++ b/packages/duckdb/tsdown.config.ts @@ -0,0 +1,13 @@ +import { defineConfig } from "tsdown"; + +export default defineConfig({ + entry: ["src/**/*.ts"], + format: "esm", + target: "esnext", + unbundle: true, + fixedExtension: false, + dts: false, + sourcemap: true, + clean: true, + outDir: "dist", +}); diff --git a/packages/duckdb/turbo.json b/packages/duckdb/turbo.json new file mode 100644 index 0000000..29a59c7 --- /dev/null +++ b/packages/duckdb/turbo.json @@ -0,0 +1,12 @@ +{ + "extends": ["//"], + "tasks": { + "typecheck": { + "inputs": ["src/**", "test/**", "scripts/**", "package.json", "tsconfig*.json"], + "dependsOn": ["build"] + }, + "test": { + "dependsOn": ["build"] + } + } +} diff --git a/packages/duckdb/vitest.config.ts b/packages/duckdb/vitest.config.ts new file mode 100644 index 0000000..6210d8b --- /dev/null +++ b/packages/duckdb/vitest.config.ts @@ -0,0 +1,8 @@ +import { defineConfig } from "vitest/config"; +import { workspaceAliases } from "../../vitest.workspace-aliases"; + +export default defineConfig({ + resolve: { + alias: workspaceAliases(), + }, +}); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 311a029..5a78e6b 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -15,6 +15,9 @@ catalogs: '@cloudflare/workers-types': specifier: 5.20260904.1 version: 5.20260904.1 + '@duckdb/node-api': + specifier: 1.5.5-r.4 + version: 1.5.5-r.4 '@effect/platform-node': specifier: 4.0.0-rc.112 version: 4.0.0-rc.112 @@ -426,6 +429,37 @@ importers: specifier: 'catalog:' version: 4.1.11(@types/node@25.9.1)(vite@8.2.2(@types/node@25.9.1)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.19.0)(yaml@2.9.0)) + packages/duckdb: + dependencies: + '@bjacobso/triplex': + specifier: workspace:^ + version: link:../core + '@bjacobso/triplex-sql': + specifier: workspace:^ + version: link:../sql + '@duckdb/node-api': + specifier: 'catalog:' + version: 1.5.5-r.4 + '@effect/sql-sqlite-node': + specifier: 'catalog:' + version: 4.0.0-rc.112(effect@4.0.0-rc.112) + devDependencies: + '@bjacobso/triplex-sqlite': + specifier: workspace:^ + version: link:../sqlite + '@types/node': + specifier: 'catalog:' + version: 25.9.1 + effect: + specifier: 'catalog:' + version: 4.0.0-rc.112 + tsx: + specifier: 'catalog:' + version: 4.19.0 + vitest: + specifier: 'catalog:' + version: 4.1.11(@types/node@25.9.1)(vite@8.2.2(@types/node@25.9.1)(esbuild@0.28.1)(jiti@2.7.0)(tsx@4.19.0)(yaml@2.9.0)) + packages/foundationdb: dependencies: '@bjacobso/triplex': @@ -1032,6 +1066,52 @@ packages: '@docsearch/sidepanel-js@4.7.0': resolution: {integrity: sha512-A8r34jCU8kcIk2viECEn2msA28ojUF1BLi/3v5OWWc5G2N3jOuuumBXoeYjfr8dA0UxgFSy5R2bt12dnFJQSyA==} + '@duckdb/node-api@1.5.5-r.4': + resolution: {integrity: sha512-8v0CZNo7aM6GQCNUHERGTrZWIfss8xKzZPF3ACNhPdMMrt7UjkNI5nRQ9FTFO7Yx8ghxYxVpdotmBBMQ5vXUOA==} + + '@duckdb/node-bindings-darwin-arm64@1.5.5-r.4': + resolution: {integrity: sha512-4OdO3pkoJzAZDnZ2iyehi67XL4+WqHPCGYWbyAsYyhGJS+WscGrAnFhMhHudMM2XF8pIEM2tuOKMmwWnTNN2wQ==} + cpu: [arm64] + os: [darwin] + + '@duckdb/node-bindings-darwin-x64@1.5.5-r.4': + resolution: {integrity: sha512-WHg3E+TupdujG31LGujMMcmtPIdW6rqRHeihlKgJR8EMetHycoYkwYxRud0niitJvS+uV0du/6pdemzs3HG9GQ==} + cpu: [x64] + os: [darwin] + + '@duckdb/node-bindings-linux-arm64-musl@1.5.5-r.4': + resolution: {integrity: sha512-Tswnf+/XWpOcFJNUUu1fkFIX/su4tMcnNz0ZV0Nru3/CkJKIMwBh+VL4SM9YV4zPK90GC4Lm8gzafjc1qDmfyQ==} + cpu: [arm64] + os: [linux] + + '@duckdb/node-bindings-linux-arm64@1.5.5-r.4': + resolution: {integrity: sha512-VIeHMpYAKpGiWZ4QYsOhODAHDEcXcn089aKty0MFsMEKaEfRJWixKwnwXy/ba2/wwFmnPiSFhKNqygj2BrQbqw==} + cpu: [arm64] + os: [linux] + + '@duckdb/node-bindings-linux-x64-musl@1.5.5-r.4': + resolution: {integrity: sha512-h0ixrgGHtHh+C/Fu1eAL9hX4iCf8yuyJN6Y24Gh8axXXP8UuSh0rQRAQUQTYarQ8pm2qSG/67ZYjyYYydD+J/w==} + cpu: [x64] + os: [linux] + + '@duckdb/node-bindings-linux-x64@1.5.5-r.4': + resolution: {integrity: sha512-EY+CL/4h8MQZd9MxTBq+98m3U9osmvHBwhE9b5fMWQGt2I6p1j8fvX0SXiwy9KBfQIbjVUOf5XvGdXncI54KSg==} + cpu: [x64] + os: [linux] + + '@duckdb/node-bindings-win32-arm64@1.5.5-r.4': + resolution: {integrity: sha512-uorAnySIMwWkRDuUhM1yTDxWIislnX8mndhFIw6r4VKhVW61HEOOMX5oAgkSQFII4RNWpC5PCcUAEhSATgNvVQ==} + cpu: [arm64] + os: [win32] + + '@duckdb/node-bindings-win32-x64@1.5.5-r.4': + resolution: {integrity: sha512-X9XGcWQ10P3mvUIaMXXk2bi94Cow7b/ziTPMKxJ0U8U3wQPxsLzEUP7D7C/KrBWINl5am2k+5SkDVyA/THUgPg==} + cpu: [x64] + os: [win32] + + '@duckdb/node-bindings@1.5.5-r.4': + resolution: {integrity: sha512-n+4hEfjp4vny3BuWn5p1Gh5CzHaPRxoI8TzTytxL1GMlIKKrXcg/o6sSAjrsROUXGAM4WQCiWPLKnPsjJJSggg==} + '@effect/platform-node-shared@4.0.0-rc.112': resolution: {integrity: sha512-ttjz0xKamFN7vL8pNDYVwddJLjZvqKePc05djlz2VcdaKbLsnYbtMnL1rbOfHgEnIUSHGh7FkjaN4DM1Ov81sQ==} engines: {node: '>=18.0.0'} @@ -5793,6 +5873,47 @@ snapshots: '@docsearch/sidepanel-js@4.7.0': {} + '@duckdb/node-api@1.5.5-r.4': + dependencies: + '@duckdb/node-bindings': 1.5.5-r.4 + + '@duckdb/node-bindings-darwin-arm64@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings-darwin-x64@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings-linux-arm64-musl@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings-linux-arm64@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings-linux-x64-musl@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings-linux-x64@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings-win32-arm64@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings-win32-x64@1.5.5-r.4': + optional: true + + '@duckdb/node-bindings@1.5.5-r.4': + dependencies: + detect-libc: 2.1.2 + optionalDependencies: + '@duckdb/node-bindings-darwin-arm64': 1.5.5-r.4 + '@duckdb/node-bindings-darwin-x64': 1.5.5-r.4 + '@duckdb/node-bindings-linux-arm64': 1.5.5-r.4 + '@duckdb/node-bindings-linux-arm64-musl': 1.5.5-r.4 + '@duckdb/node-bindings-linux-x64': 1.5.5-r.4 + '@duckdb/node-bindings-linux-x64-musl': 1.5.5-r.4 + '@duckdb/node-bindings-win32-arm64': 1.5.5-r.4 + '@duckdb/node-bindings-win32-x64': 1.5.5-r.4 + '@effect/platform-node-shared@4.0.0-rc.112(effect@4.0.0-rc.112)': dependencies: '@types/ws': 8.18.1 diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 7ec8d09..fab11c7 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -12,6 +12,7 @@ overrides: "better-sqlite3": "12.10.0" catalog: + "@duckdb/node-api": 1.5.5-r.4 "@changesets/cli": 3.0.1 "@cloudflare/vitest-plugin": 1.1.4 "@cloudflare/workers-types": 5.20260904.1 diff --git a/scripts/check-packages.mjs b/scripts/check-packages.mjs index e43a0a9..2b3eca0 100644 --- a/scripts/check-packages.mjs +++ b/scripts/check-packages.mjs @@ -10,6 +10,7 @@ const packageNames = [ "@bjacobso/triplex-sqlite", "@bjacobso/triplex-postgres", "@bjacobso/triplex-cloudflare", + "@bjacobso/triplex-duckdb", "@bjacobso/triplex-host", "@bjacobso/triplex-foundationdb", "@bjacobso/triplex-testkit", @@ -27,6 +28,7 @@ const publishPackageNames = new Set([ ]); const heldPackageNames = new Set([ "@bjacobso/triplex-cloudflare", + "@bjacobso/triplex-duckdb", "@bjacobso/triplex-host", "@bjacobso/triplex-foundationdb", ]); @@ -162,6 +164,7 @@ import { validateDatalogQuery, type DatalogQuery } from "@bjacobso/triplex/datal import { Attribute, ConfigRuntime, ConfigStore, EntityType, EntityValidation, Evaluate, GraphConstraint, TypeExpr } from "@bjacobso/triplex/config"; import * as Derivation from "@bjacobso/triplex/derivation"; import * as Cloudflare from "@bjacobso/triplex-cloudflare"; +import { DuckdbSnapshot } from "@bjacobso/triplex-duckdb"; import * as FoundationDb from "@bjacobso/triplex-foundationdb"; import * as Host from "@bjacobso/triplex-host"; import * as Postgres from "@bjacobso/triplex-postgres"; @@ -219,6 +222,7 @@ void Derivation.Materialization; void Derivation.Overlay; void makeSqliteLayer; void Cloudflare; +void DuckdbSnapshot; void FoundationDb; void Host; void Postgres; @@ -254,6 +258,8 @@ void HttpAuthorizationAllowAll; import { EntityId, Triples, string } from "@bjacobso/triplex"; import * as Derivation from "@bjacobso/triplex/derivation"; import { SqliteTriples } from "@bjacobso/triplex-sqlite"; +import { makeDuckdbSnapshot, makeDuckdbFederation, SnapshotProvider } from "@bjacobso/triplex-duckdb"; +import { resolve } from "node:path"; const result = await Effect.runPromise( Effect.gen(function* () { @@ -290,19 +296,22 @@ const result = await Effect.runPromise( }], }, }); + const analytical = yield* makeDuckdbSnapshot({ scope: "pack-smoke", basis: { validAt: 1 } }); return { + analytical: yield* analytical.queryAll(query), query: yield* triples.query(query, { basis: { validAt: 1 } }), preview, materialization: yield* Derivation.Materialization.current(triples, definition, { basis: { validAt: 1 }, }), }; - }).pipe(Effect.provide(SqliteTriples.layerMemory)), + }).pipe(Effect.scoped, Effect.provide(SqliteTriples.layerMemory)), ); if ( result.query.results.length !== 1 || result.query.results[0]?.["?name"] !== "Alice" || + result.analytical.results[0]?.["?name"] !== "Alice" || result.preview.candidates.length !== 2 || result.preview.nextTemporalBoundary !== 5 || !result.preview.candidates.some((candidate) => @@ -314,6 +323,25 @@ if ( ) { throw new Error("Unexpected packaged SQLite/Datalog result: " + JSON.stringify(result)); } + +const catalog = ["a", "b"].map((id) => ({ id, tenant: id, filename: resolve("federation-" + id + ".sqlite") })); +for (const source of catalog) { + await Effect.runPromise(Effect.gen(function* () { + const triples = yield* Triples; + yield* triples.assert({ entityId: EntityId.make("same-local-id"), attribute: ":email", value: string("shared@example.com") }); + }).pipe(Effect.provide(SqliteTriples.layer({ filename: source.filename })))); +} +const federated = await Effect.runPromise(Effect.scoped(Effect.gen(function* () { + const db = yield* makeDuckdbFederation(); + return yield* db.queryAll({ + sources: { $a: "a", $b: "b" }, find: ["?a", "?b"], + where: [["$a", "?a", ":email", "?email"], ["$b", "?b", ":email", "?email"]], + }); +})).pipe(Effect.provide(SnapshotProvider.local(catalog)))); +if (federated.results.length !== 1 || federated.results[0]["?a"] === federated.results[0]["?b"] || + federated.federation.sources.some((source) => source.pid === process.pid)) { + throw new Error("Unexpected packaged federation result: " + JSON.stringify(federated)); +} `, );