diff --git a/docs/capabilities.md b/docs/capabilities.md index 30a31b5..70a2db6 100644 --- a/docs/capabilities.md +++ b/docs/capabilities.md @@ -75,7 +75,7 @@ visible in review rather than only in production. | Chains per project | 1 | `schemas/project.schema.json` (`maxItems`) | | Contract addresses | 20 | `schemas/project.schema.json` (`maxItems`) | | Blocks per approved run | 100_000, `policy.block_budget` to change. A block count, not a duration: 100k blocks is about 14 days on Ethereum (12 s blocks), 2.3 days on Base (2 s), 7 hours on Arbitrum One (0.25 s) — set it for the chain you index | `src/plan/generate.ts` (`DEFAULT_BLOCK_BUDGET`) | -| Query deadline | 60 s | `src/query/runQuery.ts` (`DEADLINE_MS`, SIGKILL) | +| Query deadline | 60 s per query, and 60 s for loading snapshots and building models, which a build does once for all its queries | `src/query/runQuery.ts` (`DEADLINE_MS`, SIGKILL) | | DuckDB memory | 1 GiB, spills to a temp dir | `src/query/workerMain.ts` (`MEMORY_LIMIT`) | | Returned rows | 10_000, `policy.row_limit` to change | `src/project/limits.ts`, enforced in `src/query/workerMain.ts` (the reader stops at the limit) | | RPC job wall clock | 30 min, resumable | `src/ingest/rindexer/runBounded.ts` | @@ -86,8 +86,11 @@ visible in review rather than only in production. | `fork` per-request timeout | 30 s | `src/fork/fetchGuard.ts` (`FORK_LIMITS`) | | `fork` redirect hops | 0 | `src/fork/fetchGuard.ts` (refused outright) | -`CHAINPLOT_QUERY_MEMORY_LIMIT` overrides the memory figure; the query still -spills to disk rather than failing when it goes over. +`CHAINPLOT_QUERY_MEMORY_LIMIT` overrides the memory figure, as a size such as +`2GB` or `1536MB`; anything else is refused by name. The query still spills to +disk rather than failing when it goes over. The default stays small because a +forked recipe is untrusted; a project with millions of rows behind its models +builds far faster with 2–3 GB, where the memory is there to give. The row limit bounds the *download*, not the rendering. The table is virtualised — 5,000 rows put 26 in the DOM — so a wide result no longer diff --git a/docs/security.md b/docs/security.md index dd7b08b..bfb97fb 100644 --- a/docs/security.md +++ b/docs/security.md @@ -16,24 +16,32 @@ than the file format. ## Containment Query execution happens in a forked child process (`src/query/workerMain.ts`), -never in the CLI process: +never in the CLI process. `build` runs all of a release's queries in one such +process: snapshots load and models build once, then each query runs in turn +against the same session. | Control | Where | |---|---| | Separate process, env stripped to `PATH`/`HOME`/`LANG` | `src/query/runQuery.ts` | | In-memory DuckDB; no database file on disk | `workerMain.ts` | +| File reads allowed for exactly the declared snapshots (`allowed_paths`), not their directories | `workerMain.ts` | | Extension autoinstall and autoload disabled | `workerMain.ts` | | `enable_external_access=false` **before any project SQL runs** | `workerMain.ts` | | Single-SELECT admission control, via DuckDB's parser | `src/query/sqlGuard.ts` | -| 60 s deadline, SIGKILL on expiry | `runQuery.ts` | +| 60 s deadline for loading and models, then 60 s per query; SIGKILL on expiry, naming the query | `runQuery.ts` | | Row limit enforced by stopping the reader, not by truncating after | `workerMain.ts` | -Ordering matters and is the part that was wrong before 2026-09-15. Snapshots -are read first, because `read_parquet` needs filesystem access. External -access is then disabled, and only after that are models materialized and the -query run. Models are project-supplied SQL like any other, so they must land +Ordering matters and is the part that was wrong before 2026-09-15. The +snapshot files are allowlisted first, by exact path, and external access is +then disabled; DuckDB refuses both to widen that list and to re-enable access +afterwards. Each snapshot is a view read in place, so a model scans only the +columns it uses rather than a copy of every column held in memory. Only after +that are models materialized and the queries run. Models are project-supplied SQL like any other, so they must land on the closed side of that door; DuckDB does not allow external access to be -re-enabled within a session. +re-enabled within a session, so the door stays shut for every query in the +batch, not only the first. Sharing the session gives one query nothing over +another: each is still admitted only as a single SELECT, which cannot change +the session the next one runs in, and all of them come from the same recipe. ### Admission control @@ -70,7 +78,8 @@ so a value cannot carry markup or a scheme into an `href`. to the bucket can serve a consistent, hostile release. **Fork only from buckets you would trust with the data itself.** - **Denial of service by a hostile recipe.** Bounded, not eliminated: a forked - query gets 60 s, a row limit, and a 1 GiB memory cap that spills to a temp + query gets 60 s, as does loading the snapshots and building its models, plus a + row limit and a 1 GiB memory cap that spills to a temp directory rather than failing. A release can still make your build slow, and can still fill that temp directory. - **Secrets you place inside the recipe directories.** The source bundle is an diff --git a/src/publish/writeRelease.ts b/src/publish/writeRelease.ts index 295fe0c..7dbd0a4 100644 --- a/src/publish/writeRelease.ts +++ b/src/publish/writeRelease.ts @@ -6,7 +6,7 @@ import type { CommandError } from "../cli/envelope.js"; import { loadProject } from "../project/load.js"; import { validateProject } from "../project/validate.js"; import { topoSortModels } from "../project/modelGraph.js"; -import { runQuery } from "../query/runQuery.js"; +import { runQueries } from "../query/runQuery.js"; import { isComplete, requiredEnd } from "../ingest/coverage.js"; import { readCoverageFile, segmentsFor } from "../ingest/coverageStore.js"; import { lastProvenCompleteBlock } from "../ingest/coverage.js"; @@ -178,7 +178,7 @@ export async function buildRelease( ); } - // Models materialize in dependency order for every query. + // Models materialize in dependency order, once per build. const modelOrder = topoSortModels(models); const modelSql: { id: string; sql: string }[] = modelOrder.map((id) => { const model = models.find((m) => m.id === id)!; @@ -198,8 +198,13 @@ export async function buildRelease( const files: string[] = []; try { - // Queries (with models materialized first). - for (const query of queries) { + // Queries: one isolated session for the whole release. Snapshots load and + // models build once, then every query runs against them in turn. + // + // Every dataset is in scope, not just the declared one, so a model over + // dataset A can feed a query on dataset B and a query can join across + // datasets. `query.dataset` still names the provenance. + const planned = queries.map((query) => { const dataset = datasetById.get(query.dataset); if (!dataset) { throw error("validation", `unknown dataset: ${query.dataset}`, { @@ -207,20 +212,22 @@ export async function buildRelease( pointer: "/queries", }); } - const sqlPath = path.resolve(projectDir, query.file); - const data = await runQuery({ - sql: fs.readFileSync(sqlPath, "utf8"), - // Every dataset is in scope, not just the declared one. Models are - // materialized into each query's session, so loading one table meant a - // model over dataset A failed every query on dataset B — which made - // models unusable in any multi-dataset project. It also lets a query - // join across datasets. `query.dataset` still names the provenance. - tables: allTables, + const sql = fs.readFileSync(path.resolve(projectDir, query.file), "utf8"); + return { query, dataset, sql }; + }); + const results = await runQueries({ + tables: allTables, + models: modelSql, + queries: planned.map(({ query, sql }) => ({ + id: query.id, + sql, rawAmountColumns: rawAmountNames(query.raw_amount_columns), rowLimit: rowLimitFor(project), - models: modelSql, - }); + })), + }); + for (const [i, { query, dataset, sql }] of planned.entries()) { + const data = results[i]!; const rel = path.join("results", `${query.id}.json`); writeJson(path.join(staging, rel), { schema_version: 1, @@ -231,7 +238,7 @@ export async function buildRelease( columns: decorateColumns(data.columns, query.raw_amount_columns), rows: data.rows, snapshot: dataset.snapshot, - query_digest: sha256(fs.readFileSync(sqlPath, "utf8")), + query_digest: sha256(sql), raw_amount_columns: rawAmountNames(query.raw_amount_columns), }); files.push(rel); diff --git a/src/query/runQuery.ts b/src/query/runQuery.ts index eddd707..f270115 100644 --- a/src/query/runQuery.ts +++ b/src/query/runQuery.ts @@ -21,23 +21,65 @@ export interface ParquetColumn { logical_type: string; } +export interface BatchQuery { + id: string; + sql: string; + rawAmountColumns: string[]; + rowLimit: number; +} + +export interface BatchRequest { + tables: Record; + models?: { id: string; sql: string }[]; + queries: BatchQuery[]; +} + +export interface BatchResult { + id: string; + columns: { name: string; logical_type: string }[]; + rows: unknown[][]; +} + +// Per step, not per batch: loading the snapshots and building the models get +// one deadline, and each query gets its own after that. A batch of nineteen +// queries is therefore held to what nineteen separate workers were. const DEADLINE_MS = 60_000; function error( code: CommandError["code"], message: string, - opts: { retryable?: boolean } = {}, + opts: { retryable?: boolean; resource_id?: string | null } = {}, ): CommandError { return { code, message, - resource_id: null, + resource_id: opts.resource_id ?? null, pointer: null, retryable: opts.retryable ?? false, suggested_next: null, }; } +// A DuckDB size: a number and a unit, as in 2GB or 1.5GiB. +const MEMORY_SIZE = /^\d+(\.\d+)?\s*(B|KB|MB|GB|TB|KiB|MiB|GiB|TiB)$/i; + +/** + * The worker's memory cap, from CHAINPLOT_QUERY_MEMORY_LIMIT when set. Read + * here because the worker starts with a stripped environment, and checked + * here so a typo is refused by name instead of surfacing as a DuckDB error. + */ +function memoryLimit(): string | undefined { + const raw = process.env.CHAINPLOT_QUERY_MEMORY_LIMIT?.trim(); + if (!raw) return undefined; + if (!MEMORY_SIZE.test(raw)) { + throw error( + "validation", + `CHAINPLOT_QUERY_MEMORY_LIMIT must be a size such as 2GB or 1536MB, not ${JSON.stringify(raw)}`, + ); + } + return raw; +} + function workerLaunch(): { modulePath: string; execArgv: string[] } { const self = fileURLToPath(import.meta.url); const isTs = self.endsWith(".ts"); @@ -60,63 +102,41 @@ function strippedEnv(): NodeJS.ProcessEnv { return env; } -interface WorkerSuccess { - columns: { name: string; logical_type: string }[]; - rows: unknown[][]; - truncated: boolean; -} +/** One line the worker wrote: setup done, a query's rows, or a failure. */ +type WorkerLine = + | { ready: true } + | { ok: true; id: string; columns: BatchResult["columns"]; rows: unknown[][]; truncated: boolean } + | { ok: false; code?: string; message?: string; model?: string; query?: string }; -function parseWorkerPayload( - line: string, -): - | { ok: true; value: WorkerSuccess } - | { ok: false; code: CommandError["code"]; message: string } - | null { - let parsed: unknown; +const named = (id: string): string => (id ? `query ${id}` : "query"); + +function parseLine(line: string): WorkerLine | null { try { - parsed = JSON.parse(line); + const parsed: unknown = JSON.parse(line); + return parsed !== null && typeof parsed === "object" ? (parsed as WorkerLine) : null; } catch { return null; } - if (parsed === null || typeof parsed !== "object" || !("ok" in parsed)) { - return null; - } - const body = parsed as { - ok: unknown; - columns?: WorkerSuccess["columns"]; - rows?: WorkerSuccess["rows"]; - truncated?: unknown; - code?: unknown; - message?: unknown; - }; - if (body.ok === true && body.columns !== undefined && body.rows !== undefined) { - return { - ok: true, - value: { - columns: body.columns, - rows: body.rows, - truncated: body.truncated === true, - }, - }; - } - if (body.ok === false) { - const code = typeof body.code === "string" ? body.code : "validation"; - return { - ok: false, - code: code as CommandError["code"], - message: String(body.message ?? "worker failed"), - }; - } - return null; } -function invokeWorker(req: { - sql: string; - tables: Record; - rawAmountColumns?: string[]; - rowLimit?: number; - models?: { id: string; sql: string }[]; -}): Promise { +/** + * Run a batch of queries in one isolated worker: snapshots load and models + * build once, then each query runs in turn. Resolves with every result in + * order, or rejects with the first failure, naming the model or query that + * caused it. + */ +export function runQueries( + req: BatchRequest, + opts: { deadlineMs?: number } = {}, +): Promise { + if (req.queries.length === 0) return Promise.resolve([]); + let limit: string | undefined; + try { + limit = memoryLimit(); + } catch (err) { + return Promise.reject(err); + } + const deadlineMs = opts.deadlineMs ?? DEADLINE_MS; const { modulePath, execArgv } = workerLaunch(); return new Promise((resolve, reject) => { const child = fork(modulePath, [], { @@ -124,64 +144,105 @@ function invokeWorker(req: { env: strippedEnv(), stdio: ["pipe", "pipe", "pipe", "ipc"], }); - let stdout = ""; + const results: BatchResult[] = []; + let buffered = ""; let stderr = ""; + let ready = false; let settled = false; + let timer: NodeJS.Timeout | undefined; - const timer = setTimeout(() => { - child.kill("SIGKILL"); - finish( - error("transient_dependency", "query deadline exceeded", { - retryable: true, - }), - ); - }, DEADLINE_MS); + // What the worker is doing right now, for the deadline's message. A single + // query from runQuery has an empty id, and is simply "query". + const waitingOn = (): string | null => + ready ? (req.queries[results.length]?.id ?? null) : null; - function finish(err: CommandError | null, value?: WorkerSuccess): void { - if (settled) { - return; - } + const arm = (): void => { + clearTimeout(timer); + timer = setTimeout(() => { + child.kill("SIGKILL"); + const id = waitingOn(); + finish( + error( + "transient_dependency", + id === null + ? `loading snapshots and building models took longer than ${deadlineMs / 1000} s` + : `${named(id)} took longer than ${deadlineMs / 1000} s`, + { retryable: true, resource_id: id || null }, + ), + ); + }, deadlineMs); + }; + + function finish(err: CommandError | null): void { + if (settled) return; settled = true; clearTimeout(timer); if (err) { + child.kill("SIGKILL"); reject(err); } else { - resolve(value as WorkerSuccess); + resolve(results); } } + function handle(line: WorkerLine): void { + if ("ready" in line) { + ready = true; + arm(); + return; + } + if (line.ok === true) { + const query = req.queries[results.length]; + if (line.truncated && query) { + finish( + error( + "policy_refused", + `${named(query.id)} returned more than ${query.rowLimit} rows; add a LIMIT or aggregate instead`, + { resource_id: query.id || null }, + ), + ); + return; + } + results.push({ id: line.id, columns: line.columns, rows: line.rows }); + if (results.length === req.queries.length) finish(null); + else arm(); + return; + } + finish( + error((line.code ?? "validation") as CommandError["code"], String(line.message ?? "worker failed"), { + resource_id: line.model ?? (line.query || null), + }), + ); + } + + arm(); child.stdout?.setEncoding("utf8"); child.stdout?.on("data", (chunk: string) => { - stdout += chunk; + buffered += chunk; + let nl: number; + while ((nl = buffered.indexOf("\n")) !== -1) { + const parsed = parseLine(buffered.slice(0, nl)); + buffered = buffered.slice(nl + 1); + if (parsed) handle(parsed); + } }); child.stderr?.setEncoding("utf8"); child.stderr?.on("data", (chunk: string) => { stderr += chunk; }); - child.on("error", (err) => { - finish(error("internal", err.message)); - }); + child.on("error", (err) => finish(error("internal", err.message))); child.on("exit", (code) => { - const payload = parseWorkerPayload(stdout.trim().split("\n").pop() ?? ""); - if (payload?.ok === true) { - finish(null, payload.value); - return; - } - if (payload?.ok === false) { - finish(error(payload.code, payload.message)); - return; - } + if (settled) return; const detail = stderr.trim() || `worker exited with code ${code ?? "unknown"}`; - finish(error("internal", detail)); + finish(error("internal", detail, { resource_id: waitingOn() || null })); }); child.stdin?.write( JSON.stringify({ - sql: req.sql, tables: req.tables, - rawAmountColumns: req.rawAmountColumns ?? [], - rowLimit: req.rowLimit, models: req.models ?? [], + queries: req.queries, + memoryLimit: limit, }) + "\n", ); child.stdin?.end(); @@ -189,24 +250,24 @@ function invokeWorker(req: { } export async function runQuery(req: QueryRequest): Promise { - const { columns, rows, truncated } = await invokeWorker(req); - if (truncated) { - throw error( - "policy_refused", - `query returned more than ${req.rowLimit} rows; add a LIMIT or aggregate instead`, - ); - } + const [result] = await runQueries({ + tables: req.tables, + models: req.models, + queries: [{ id: "", sql: req.sql, rawAmountColumns: req.rawAmountColumns, rowLimit: req.rowLimit }], + }); return { - columns, - rows, + columns: result!.columns, + rows: result!.rows, snapshot: Object.values(req.tables)[0] ?? "", }; } export async function describeParquet(parquetPath: string): Promise { - const { rows } = await invokeWorker({ + const { rows } = await runQuery({ sql: "DESCRIBE SELECT * FROM snapshot", tables: { snapshot: parquetPath }, + rawAmountColumns: [], + rowLimit: 10_000, }); return rows.map((row) => ({ name: String(row[0]), diff --git a/src/query/workerMain.ts b/src/query/workerMain.ts index 5999ac5..8e374ec 100644 --- a/src/query/workerMain.ts +++ b/src/query/workerMain.ts @@ -15,12 +15,24 @@ const { inspectSerializedSql } = (await import( ).href )) as typeof import("./sqlGuard.js"); -interface WorkerRequest { +interface WorkerQuery { + id: string; sql: string; - tables: Record; rawAmountColumns?: string[]; rowLimit?: number; +} + +/** + * One session for a whole batch: snapshots load and models build once, then + * every query runs against them in turn. A build used to fork a worker per + * query, each repeating the load and the models before its one SELECT. + */ +interface WorkerRequest { + tables: Record; models?: { id: string; sql: string }[]; + queries: WorkerQuery[]; + /** Validated by the parent; the worker's own environment is stripped. */ + memoryLimit?: string; } const MODEL_ID = /^[a-z0-9_]+$/; @@ -30,8 +42,9 @@ const DEFAULT_ROW_LIMIT = 10_000; // A forked recipe runs here, so an unbounded query is the host's problem. // DuckDB spills past this rather than failing, provided a temp directory -// exists — without one it raises an out-of-memory error instead. -const MEMORY_LIMIT = process.env.CHAINPLOT_QUERY_MEMORY_LIMIT ?? "1GB"; +// exists — without one it raises an out-of-memory error instead. The parent +// passes CHAINPLOT_QUERY_MEMORY_LIMIT in the request when it is set. +const DEFAULT_MEMORY_LIMIT = "1GB"; /** * Canonical sort key for uint256/int256 amounts carried as decimal strings. @@ -65,6 +78,14 @@ function issueError(issue: SqlIssue): Error { return err; } +/** One protocol line on stdout. The parent reads them as they arrive. */ +function emit(payload: unknown): void { + fs.writeSync(1, JSON.stringify(payload) + "\n"); +} + +/** What the session is busy with, so a failure can say which step it was. */ +let current: { model: string } | { query: string } | null = null; + function quoteIdent(name: string): string { return `"${name.replaceAll('"', '""')}"`; } @@ -131,27 +152,29 @@ async function readRequest(): Promise { return JSON.parse(buf) as WorkerRequest; } -async function execute(req: WorkerRequest): Promise<{ - columns: { name: string; logical_type: string }[]; - rows: unknown[][]; - truncated: boolean; -}> { - const rowLimit = req.rowLimit ?? DEFAULT_ROW_LIMIT; +async function execute(req: WorkerRequest): Promise { const spillDir = fs.mkdtempSync(path.join(os.tmpdir(), "chainplot-duckdb-")); const instance = await DuckDBInstance.create(":memory:", { autoinstall_known_extensions: "false", autoload_known_extensions: "false", - memory_limit: MEMORY_LIMIT, + memory_limit: req.memoryLimit ?? DEFAULT_MEMORY_LIMIT, temp_directory: spillDir, }); try { const conn = (await instance.connect()) as unknown as Conn; try { // Snapshots are the only filesystem reads this process is allowed to - // make, so they happen first... - for (const [name, parquetPath] of Object.entries(req.tables)) { + // make: exactly those files, not their directories, so a model cannot + // reach a file sitting beside one. The list is fixed before the door + // shuts, and DuckDB refuses to widen it or to reopen the door after. + // Each snapshot is then a view read in place, so a model scans only + // the columns it uses instead of a full copy held in memory. + const snapshots = Object.entries(req.tables).map( + ([name, file]) => [name, path.resolve(file)] as const, + ); + if (snapshots.length > 0) { await conn.run( - `CREATE TABLE ${quoteIdent(name)} AS SELECT * FROM read_parquet(${quoteString(parquetPath)})`, + `SET allowed_paths = [${snapshots.map(([, file]) => quoteString(file)).join(", ")}]`, ); } @@ -160,6 +183,11 @@ async function execute(req: WorkerRequest): Promise<{ // trusted than the query itself, so they must land on this side of it. // DuckDB does not allow re-enabling external access in a session. await conn.run("SET enable_external_access=false"); + for (const [name, file] of snapshots) { + await conn.run( + `CREATE VIEW ${quoteIdent(name)} AS SELECT * FROM read_parquet(${quoteString(file)})`, + ); + } // rindexer exports block_timestamp as TIMESTAMP WITH TIME ZONE, and // DuckDB renders, casts and buckets that type in the session's // TimeZone, which defaults to the machine's. Pinning UTC is what makes @@ -181,6 +209,7 @@ async function execute(req: WorkerRequest): Promise<{ await conn.run(SORT_KEY_MACRO); for (const model of req.models ?? []) { + current = { model: model.id }; if (!MODEL_ID.test(model.id)) { throw issueError({ code: "validation", @@ -198,24 +227,32 @@ async function execute(req: WorkerRequest): Promise<{ } } - await assertAdmissible(conn, req.sql, { - label: "query", - rawAmountColumns: req.rawAmountColumns ?? [], - }); - - // Stop reading at the limit instead of materializing everything and - // rejecting afterwards; `done` tells us whether more rows existed. - const reader = await conn.streamAndReadUntil(req.sql, rowLimit + 1); - const columns: { name: string; logical_type: string }[] = []; - for (let i = 0; i < reader.columnCount; i++) { - columns.push({ - name: reader.columnName(i), - logical_type: reader.columnType(i).toString(), + current = null; + emit({ ready: true }); + + for (const query of req.queries) { + current = { query: query.id }; + const rowLimit = query.rowLimit ?? DEFAULT_ROW_LIMIT; + await assertAdmissible(conn, query.sql, { + label: "query", + rawAmountColumns: query.rawAmountColumns ?? [], }); + + // Stop reading at the limit instead of materializing everything and + // rejecting afterwards; `done` tells us whether more rows existed. + const reader = await conn.streamAndReadUntil(query.sql, rowLimit + 1); + const columns: { name: string; logical_type: string }[] = []; + for (let i = 0; i < reader.columnCount; i++) { + columns.push({ + name: reader.columnName(i), + logical_type: reader.columnType(i).toString(), + }); + } + const all = reader.getRowsJson(); + const truncated = all.length > rowLimit || !reader.done; + emit({ ok: true, id: query.id, columns, rows: all.slice(0, rowLimit), truncated }); } - const all = reader.getRowsJson(); - const truncated = all.length > rowLimit || !reader.done; - return { columns, rows: all.slice(0, rowLimit), truncated }; + current = null; } finally { (conn as unknown as { closeSync(): void }).closeSync(); } @@ -232,13 +269,13 @@ function reply(payload: unknown, exitCode: number): void { try { const req = await readRequest(); - const { columns, rows, truncated } = await execute(req); - reply({ ok: true, columns, rows, truncated }, 0); + await execute(req); + process.exit(0); } catch (err) { const code = err !== null && typeof err === "object" && "chainplotCode" in err ? String((err as { chainplotCode: unknown }).chainplotCode) : "validation"; const message = err instanceof Error ? err.message : String(err); - reply({ ok: false, code, message }, 1); + reply({ ok: false, code, message, ...(current ?? {}) }, 1); } diff --git a/tests/query/batch.test.ts b/tests/query/batch.test.ts new file mode 100644 index 0000000..a7049c5 --- /dev/null +++ b/tests/query/batch.test.ts @@ -0,0 +1,162 @@ +import { describe, expect, it } from "vitest"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { runQueries } from "../../src/query/runQuery.js"; + +function tempCopyOfFixture(): string { + const src = path.resolve( + path.dirname(fileURLToPath(import.meta.url)), + "../../templates/fixture-transfers/snapshots/amounts.parquet", + ); + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "chainplot-batch-")); + const dest = path.join(dir, "amounts.parquet"); + fs.copyFileSync(src, dest); + return dest; +} + +const query = (id: string, sql: string, rowLimit = 10_000) => ({ + id, + sql, + rawAmountColumns: [], + rowLimit, +}); + +// A build used to fork one worker per query, and each worker loaded every +// snapshot and rebuilt every model before running its one SELECT. A batch +// loads and builds once, then runs the queries in turn. +describe("query batches", () => { + it("runs every query against one build of the models, in order", async () => { + const results = await runQueries({ + tables: { amounts: tempCopyOfFixture() }, + // random() is drawn once per build of the model, so two queries + // reading the same value prove they shared one build. + models: [{ id: "drawn", sql: "SELECT random() AS r, (SELECT COUNT(*) FROM amounts) AS n" }], + queries: [query("first", "SELECT r, n FROM drawn"), query("second", "SELECT r FROM drawn")], + }); + expect(results.map((r) => r.id)).toEqual(["first", "second"]); + expect(results[0]!.rows[0]![1]).toBe("8"); + expect(results[1]!.rows[0]![0]).toBe(results[0]!.rows[0]![0]); + }); + + it("names the query that failed", async () => { + const run = runQueries({ + tables: { amounts: tempCopyOfFixture() }, + queries: [ + query("fine", "SELECT COUNT(*) FROM amounts"), + query("broken", "SELECT no_such_column FROM amounts"), + ], + }); + await expect(run).rejects.toMatchObject({ code: "validation", resource_id: "broken" }); + }); + + it("names a failing model rather than the query that would have used it", async () => { + const run = runQueries({ + tables: { amounts: tempCopyOfFixture() }, + models: [{ id: "bad", sql: "SELECT nope FROM amounts" }], + queries: [query("uses_bad", "SELECT * FROM bad")], + }); + await expect(run).rejects.toMatchObject({ code: "validation", resource_id: "bad" }); + }); + + it("holds each query to its own row limit, and says which one overran", async () => { + const run = runQueries({ + tables: { amounts: tempCopyOfFixture() }, + queries: [ + query("small", "SELECT 1"), + query("too_big", "SELECT * FROM range(50)", 10), + ], + }); + await expect(run).rejects.toMatchObject({ code: "policy_refused", resource_id: "too_big" }); + }); + + it("keeps external access shut for every query, not just the first", async () => { + const run = runQueries({ + tables: { amounts: tempCopyOfFixture() }, + queries: [ + query("fine", "SELECT 1"), + query("reads_disk", "SELECT * FROM read_csv('/etc/passwd')"), + ], + }); + await expect(run).rejects.toMatchObject({ resource_id: "reads_disk" }); + }); + + it("gives each query its own deadline, and names the one that ran out", async () => { + const run = runQueries( + { + tables: { amounts: tempCopyOfFixture() }, + queries: [ + query("quick", "SELECT 1"), + // About 10^10 rows: far longer than the deadline below. + query("slow", "SELECT SUM(a.range * b.range) FROM range(100000) a, range(100000) b"), + ], + }, + { deadlineMs: 1_500 }, + ); + await expect(run).rejects.toMatchObject({ + code: "transient_dependency", + resource_id: "slow", + }); + }); + + // Snapshots are read in place through views, with access allowed to exactly + // the declared files. Another parquet beside one of them stays out of reach. + it("reads the declared snapshot and nothing beside it", async () => { + const parquet = tempCopyOfFixture(); + const sibling = path.join(path.dirname(parquet), "other.parquet"); + fs.copyFileSync(parquet, sibling); + const ok = await runQueries({ + tables: { amounts: parquet }, + queries: [query("declared", "SELECT COUNT(*) FROM amounts")], + }); + expect(ok[0]!.rows[0]![0]).toBe("8"); + const run = runQueries({ + tables: { amounts: parquet }, + queries: [query("sibling", `SELECT COUNT(*) FROM read_parquet('${sibling}')`)], + }); + await expect(run).rejects.toThrow(/file system operations are disabled/); + }); + + describe("CHAINPLOT_QUERY_MEMORY_LIMIT", () => { + const saved = process.env.CHAINPLOT_QUERY_MEMORY_LIMIT; + const restore = () => { + if (saved === undefined) delete process.env.CHAINPLOT_QUERY_MEMORY_LIMIT; + else process.env.CHAINPLOT_QUERY_MEMORY_LIMIT = saved; + }; + + // The worker's environment is stripped to PATH/HOME/LANG, so the override + // documented in capabilities.md used to stop at the parent. + it("reaches the worker", async () => { + process.env.CHAINPLOT_QUERY_MEMORY_LIMIT = "2GB"; + try { + const [result] = await runQueries({ + tables: { amounts: tempCopyOfFixture() }, + queries: [query("limit", "SELECT current_setting('memory_limit')")], + }); + // 2GB is decimal; DuckDB reports it in binary units. + expect(result!.rows[0]![0]).toBe("1.8 GiB"); + } finally { + restore(); + } + }); + + it("refuses a value that is not a size", async () => { + process.env.CHAINPLOT_QUERY_MEMORY_LIMIT = "lots"; + try { + const run = runQueries({ + tables: { amounts: tempCopyOfFixture() }, + queries: [query("limit", "SELECT 1")], + }); + await expect(run).rejects.toMatchObject({ code: "validation" }); + await expect(run).rejects.toThrow(/CHAINPLOT_QUERY_MEMORY_LIMIT/); + } finally { + restore(); + } + }); + }); + + it("runs nothing for an empty batch", async () => { + expect(await runQueries({ tables: {}, queries: [] })).toEqual([]); + }); +});