diff --git a/docs/duckdb-sandbox.md b/docs/duckdb-sandbox.md new file mode 100644 index 00000000..fd84b787 --- /dev/null +++ b/docs/duckdb-sandbox.md @@ -0,0 +1,225 @@ +# Contract SQL runs in a DuckDB sandbox + +Every DuckDB connection the engine opens is confined to the files a contract +declares and its own directory. SQL in a contract cannot read the rest of the +host, fetch a URL, attach another database, or write anywhere else. + +```yaml +# contract.fluid.yaml +builds: + - id: summarise + pattern: embedded-logic + engine: sql + properties: + sql: SELECT * FROM read_csv('/etc/passwd') +``` + +```console +$ fluid apply contract.fluid.yaml --mode amend-and-build --yes +🔷 Build 'summarise' (embedded-SQL / local DuckDB) + ❌ Failed: 1 action(s) failed + Permission Error: Cannot access file "/etc/passwd" - file system operations are + disabled by configuration DuckDB refused it: contract SQL may only read and write + the locations the contract declares and its own directory (/work/orders, + /work/orders/runtime, /tmp/fluid_h44kgt2z, /work/orders/out/summary.csv). Declare + the file as an input, or move it under the contract's directory. +``` + +The same SQL reading the contract's own data builds as before: + +```yaml + properties: + sql: SELECT id, amount * 2 AS doubled FROM read_csv('data/orders.csv') +``` + +## What a build's SQL can reach + +| Where the SQL runs | It can read and write | +|---|---| +| Embedded-SQL build on the local DuckDB engine (`builds[].properties.sql`) | the contract's directory, the FLUID workspace it sits in (`fluid.workspace.yaml`), `./runtime`, the run's scratch directory, each declared `parameters.inputs[].path`, each resolved `consumes[]` upstream, the expose's landing path, and the `s3://` prefixes those name. A declared local path counts only [inside the allowed directories](#declared-locations-stay-inside-the-allowed-directories) | +| DuckDB acquisition build (`pattern: acquisition`, `engine: duckdb`) | the contract's directory, the declared `source.connection.uri` (or stream paths), and each stream's landing file, each inside the allowed directories. A `mysql` source is attached before the sandbox closes; so is a `sqlite` source, and its file must also be inside the allowed directories | +| `fluid validate` quality rules, `fluid verify`, `fluid diff` | the one file being checked | +| `fluid contract-tests` local actions | each declared input file and each output file | +| Discovery (`fluid forge data-model from-source`, `discover`) | the one file or URL being introspected; a JDBC source is attached first | +| MCP output port (DuckDB driver) | the bound file | + +Everything else is refused, including: + +- an absolute path outside the list: `read_csv('/etc/passwd')`, `read_text(...)`, + `read_blob(...)`, `read_parquet(...)`, `read_json(...)`, `glob('/etc/*')`; +- a path that climbs out: `../`, `/./../`, or a symlink pointing outside; +- `~/...` unless the matching file under `$HOME` is itself in the list; +- a URL (`http://`, `https://`, `s3://` ...) the contract does not declare; +- `ATTACH` of another database file, `COPY ... TO` / `COPY ... FROM` elsewhere; +- `INSTALL` / `LOAD` of an extension, and any `SET` (the configuration is locked, + so `SET enable_external_access = true` is refused too); +- a function from an extension the engine does not load for that build: + `sqlite_scan`, `read_xlsx`, `ST_Read`, `delta_scan`, `iceberg_scan`. DuckDB + used to load these on first use; with autoloading off they are not in the + catalog, even for a file inside the contract's directory. Read the data as + CSV, Parquet or JSON, or land it with an acquisition build first. + +## Declared locations stay inside the allowed directories + +Each declared input and output is granted to the build's SQL, and whoever +writes the contract writes the declarations. So a declaration grants a local +path only inside these directories: + +- the contract's directory and the FLUID workspace it sits in; +- `./runtime` and the run's scratch directory (`./runtime` only when it is a + real directory, or a symlink into the contract's directory or workspace; see + below); +- the upstream roots in `FLUID_UPSTREAM_CONTRACTS`; +- the directories the operator lists in `FLUID_DUCKDB_ALLOWED_DIRS`. + +Anything else is refused before any SQL runs. Declaring an innocuous glob in +`$HOME` does not make `~/.aws/credentials` readable: + +```yaml + properties: + sql: SELECT content FROM read_text('~/.aws/credentials') + parameters: + inputs: + - name: d + path: /home/me/*.csv +``` + +```console +$ fluid apply contract.fluid.yaml --mode amend-and-build --yes +🔷 Build 'summarise' (embedded-SQL / local DuckDB) + ❌ Failed: 1 action(s) failed + The contract declares '/home/me/*.csv' (/home/me), outside the directories it + may read and write (/work/orders, /work/orders/runtime, /tmp/fluid_yc3ryr57). + The operator can allow a directory with FLUID_DUCKDB_ALLOWED_DIRS. +``` + +Both sides are compared after resolving symlinks, as DuckDB resolves them: a +symlink inside the contract's directory that points at `/` grants nothing. A +relative declared path is resolved where DuckDB opens it, the working +directory, and is confined the same way, so `path: ./*.py` cannot grant a +server's working directory. + +`./runtime` is granted by convention, not by a declaration, and it sits in the +working directory, which is usually the contract's own. A repository that +ships `runtime` as a symlink out of itself (`runtime -> ../../..`, which is +`$HOME` for a clone at `~/src/repo/product`) does not get that directory +granted: the symlink is ignored with a `local_runtime_not_granted` warning, and +SQL that reads or writes under it is refused. + +A SQLite source (`source.kind: sqlite`) is confined the same way, and refused +with the same `FLUID_DUCKDB_ALLOWED_DIRS` message outside the allowed +directories. This check is the only one on it: the sqlite scanner opens the +file through its own library, which DuckDB's allowlist does not bound. + +What a declaration grants, once allowed: + +| Declared | Granted | +|---|---| +| a file (`/shared/reference/rates.csv`) | that file | +| a glob (`data/*.csv`, `landing/**/*.parquet`) | the directory above its first wildcard (`data/`, `landing/`), which must itself be inside the allowed directories | +| a directory (`data/`) | everything under it | +| an `s3://` URL | its prefix (see the limits below) | + +A glob grants its directory, not the files it matches when the run starts, +because DuckDB expands it again when the SQL runs: that expansion includes +dotfiles (macOS `._orders.csv` on an exFAT or SMB volume) and files that landed +after the run started, and DuckDB refuses the whole read if any one of them is +not granted. The directory is no wider than the declaration could already be: +everything inside the allowed directories is declarable. DuckDB checks each +expanded file after resolving symlinks, so a matched symlink that points +outside the directory is refused. + +A `[`, `?` or `*` in the name of a directory that exists, such as a project +checked out under `Proj [old]/`, is part of that name, not a wildcard: +`customers.csv` declared there grants that file. DuckDB tries such a name as a +pattern first, so if it also matches a sibling directory (`Proj o/`), DuckDB +reads the sibling. The sibling is not granted, so that read is refused. Rename +the directory in that case. + +### Reading a file outside the contract's directory + +The operator allows its directory; the contract then declares the file: + +```console +$ export FLUID_DUCKDB_ALLOWED_DIRS=/shared/reference # ':'-separated, absolute +$ fluid apply contract.fluid.yaml --mode amend-and-build --yes +``` + +```yaml + properties: + sql: SELECT * FROM rates + parameters: + inputs: + - name: rates + path: /shared/reference/rates.csv +``` + +The declared file becomes readable. Its neighbours in `/shared/reference` do +not, unless the contract declares them too (or declares a glob or the +directory, which grant `/shared/reference` itself, as the table above says). `FLUID_DUCKDB_ALLOWED_DIRS` +is read from the environment of the process that runs the engine; no contract +field can set it. A relative entry, or `/`, is refused. + +A product inside a FLUID workspace can also read its sibling products' files by +path, and a `consumes[]` entry resolves to the upstream's landed file. + +## How it works + +`fluid_build/providers/_duckdb_sandbox.py` is the one place the engine calls +`duckdb.connect`. A test (`tests/providers/test_duckdb_sandbox.py`) fails if +any other module opens or queries DuckDB without it. For each connection it: + +1. connects (a file-backed database is opened here and needs no grant); +2. turns off persistent secrets and community extensions, loads the extensions + the call site needs, and runs its set-up (an `ATTACH` of a declared source, + an object-store secret); +3. turns off extension autoinstall and autoload; +4. pins `home_directory` to `$HOME`, then sets `allowed_directories` and + `allowed_paths` to the list above; +5. sets `enable_external_access = false`; +6. sets `lock_configuration = true`. + +These are DuckDB's own settings, in the order the +[Securing DuckDB](https://duckdb.org/docs/current/operations_manual/securing_duckdb/overview.html) +guide gives. + +## Requirements and limits + +- **DuckDB 1.5.0 or newer.** `allowed_directories` arrived in 1.2, but up to + 1.4.3 `/./../` escaped it, and through 1.4.x a symlink inside an allowed + directory reached its target and relative paths were refused. The `local` + extra now requires `duckdb>=1.5.0`, and the engine refuses to open DuckDB on + anything older. +- **`~` follows `$HOME`.** DuckDB expands `~` with `$HOME` when it checks a + path and with its `home_directory` setting when it opens one + ([duckdb/duckdb#26064](https://github.com/duckdb/duckdb/issues/26064)). The + sandbox sets the two to the same directory and locks them. +- **A remote prefix bounds the bucket, not the path.** DuckDB 1.5 does not + resolve `..` inside a URL, so a declared `s3://bucket/landing/` also lets + the SQL reach other keys in that bucket with the same credentials. +- **A loaded database scanner is not bounded by the allowlist.** The `sqlite`, + `postgres` and `mysql` extensions open files and sockets through their own + client libraries, not through DuckDB's file system, so neither + `allowed_directories` nor `enable_external_access` limits them. On a + connection that loads `sqlite` (an acquisition build with a SQLite source), + `sqlite_scan` and `ATTACH ... (TYPE sqlite)` can open any SQLite file the + process can read; the declared source file itself is confined before it is + attached (see above). On one that loads `postgres`, `postgres_scan` can connect + to any host the process can reach, the hosting service's own database + included. The engine loads them only for an acquisition build's declared + source, discovery and the copilot's sample-rows tool, whose SQL the engine + builds from validated identifiers; contract SQL never runs on such a + connection (`test_contract_sql_runs_on_a_connection_with_no_database_scanner`). +- **One open connection per database file.** DuckDB shares one instance per + database file within a process, and the sandbox locks that instance. A second + connection to a file that is already open is refused (`DuckDB database ... + is already open in this process`). Each call site closes its connection, + failures included. Two MCP DuckDB drivers bound to the same `.duckdb` file + in one process, or two concurrent `persist=True` local runs, hit this. +- **Persistent DuckDB secrets are off.** Secrets saved in + `~/.duckdb/stored_secrets` are no longer loaded; object-store builds use the + credential-chain secret the engine creates. +- **Defense in depth, not isolation.** DuckDB describes these settings as "not + a substitute for proper sandboxing". A service that runs other people's + contracts (a multi-tenant service, a shared CI runner) should still run each one + in its own container. diff --git a/fluid_build/build_runners/_bigquery_load.py b/fluid_build/build_runners/_bigquery_load.py index 9b4f7752..af0d0cf2 100644 --- a/fluid_build/build_runners/_bigquery_load.py +++ b/fluid_build/build_runners/_bigquery_load.py @@ -195,12 +195,19 @@ def utc_adjusted_copy(path: str, schema: Any) -> Optional[str]: wanted = {name.lower(): name for name in _timestamp_columns(schema)} if not wanted: return None - import duckdb - - con = duckdb.connect(":memory:") + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect + + base, _ext = os.path.splitext(str(path)) + out = f"{base}.bq-load.parquet" + # The landed file and its UTC-adjusted copy beside it, nothing else. Read + # as UTC (set at connect: the sandbox locks the configuration): the cast + # below turns the naive wall clock into an instant. + con = secure_duckdb_connect( + ":memory:", + allow=DuckDBAllowlist.none().with_paths(path, out), + config={"TimeZone": "UTC"}, + ) try: - # Read as UTC: the cast below turns the naive wall clock into an instant. - con.execute("SET TimeZone = 'UTC'") described = con.execute("DESCRIBE SELECT * FROM read_parquet(?)", [str(path)]).fetchall() naive = [ str(row[0]) @@ -212,8 +219,6 @@ def utc_adjusted_copy(path: str, schema: Any) -> Optional[str]: replaced = ", ".join( f"CAST({_quote_ident(c)} AS TIMESTAMPTZ) AS {_quote_ident(c)}" for c in naive ) - base, _ext = os.path.splitext(str(path)) - out = f"{base}.bq-load.parquet" # A path is not a bindable parameter in COPY ... TO; quote it as a literal. target = "'" + out.replace("'", "''") + "'" con.execute( diff --git a/fluid_build/build_runners/_embedded_sql_io.py b/fluid_build/build_runners/_embedded_sql_io.py index cb9b37ea..0e2dd032 100644 --- a/fluid_build/build_runners/_embedded_sql_io.py +++ b/fluid_build/build_runners/_embedded_sql_io.py @@ -471,16 +471,10 @@ def relations_read(sql: str) -> Optional[FrozenSet[str]]: if not str(sql or "").strip(): return None try: - import duckdb - - con = duckdb.connect( - ":memory:", - config={ - "enable_external_access": False, - "autoinstall_known_extensions": False, - "autoload_known_extensions": False, - }, - ) + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect + + # Parsed only, never bound or run: no file, URL or extension at all. + con = secure_duckdb_connect(":memory:", allow=DuckDBAllowlist.none()) try: row = con.execute("SELECT json_serialize_sql(?)", [str(sql)]).fetchone() finally: diff --git a/fluid_build/build_runners/duckdb/runner.py b/fluid_build/build_runners/duckdb/runner.py index 86125461..cd9e594f 100644 --- a/fluid_build/build_runners/duckdb/runner.py +++ b/fluid_build/build_runners/duckdb/runner.py @@ -56,6 +56,11 @@ ) from fluid_build.api.schema import SchemaFingerprint from fluid_build.api.source import AcquisitionMode +from fluid_build.providers._duckdb_sandbox import ( + DuckDBAllowlist, + confine_declared, + secure_duckdb_connect, +) from fluid_build.providers._sql_safety import ( build_libpq_dsn, quote_ansi_string_literal, @@ -548,6 +553,80 @@ def _required_extensions(kind: str, uri: Optional[str]) -> List[str]: return out +# ── Sandbox ────────────────────────────────────────────────────────────── +# +# Every connection this runner opens goes through ``secure_duckdb_connect``. +# The SQL it runs is assembled from the contract (reader options, quality-gate +# predicates, incremental filters, the masking projection), so each connection +# reaches the contract's directory and the locations the build declares, and +# nothing else on the host. Extensions, ATTACHes and secrets are set up in +# ``before_lock``: once the configuration is locked none can be added. + + +def _run_streams(ctx: RunContext) -> List[str]: + return list(ctx.source.streams) or DuckdbRunner()._infer_streams(ctx) + + +def _declared_roots(ctx: RunContext) -> List[Any]: + """Where a declared source or landing may be: the contract's directory and workspace. + + ``DuckDBAllowlist.with_declared`` adds the operator's + ``FLUID_DUCKDB_ALLOWED_DIRS``; the contract itself cannot add a root. + """ + from fluid_build.util.workspace_root import find_workspace_root + + return [ctx.workdir, find_workspace_root(Path(ctx.workdir))] + + +def _grant_source(allow: DuckDBAllowlist, ctx: RunContext, streams: List[str]) -> DuckDBAllowlist: + """``allow`` plus where a filesystem / http source reads from. + + A database source needs no grant: postgres is scanned through its DSN, and + mysql / sqlite are attached before the lock (``_attach_external_databases``), + a sqlite file only inside :func:`_declared_roots` (``confine_declared``). + A local source is granted only inside :func:`_declared_roots` (a glob as + the directory above its first wildcard): the contract names it, so naming + ``~/.ssh/*`` must not make it readable. + """ + if ctx.source.kind not in {"filesystem", "http"}: + return allow + uri = dict(ctx.source.connection.raw).get("uri") + # ``_select_for_stream`` reads ``uri or stream``: without a uri, the stream + # names the file. + return allow.with_declared(*([uri] if uri else streams), within=_declared_roots(ctx)) + + +def _grant_destination( + allow: DuckDBAllowlist, ctx: RunContext, streams: List[str], sink_format: str +) -> DuckDBAllowlist: + """``allow`` plus each stream's landed file and its late-arrival sibling. + + Confined like the source: a landing path the contract declares outside + :func:`_declared_roots` would let it write anywhere on the host. + """ + out_dir = Path(ctx.workdir) / "out" + within = _declared_roots(ctx) + for stream in streams: + dest = _resolve_destination_path(ctx, stream, sink_format, out_dir) + if _is_remote_uri(dest): + allow = allow.with_locations(dest) + continue + main = Path(dest) + late = main.with_name(main.stem + "__late_events" + main.suffix) + allow = allow.with_declared(str(main), str(late), within=within) + return allow + + +def _run_allowlist( + ctx: RunContext, streams: List[str], sink_format: Optional[str] = None +) -> DuckDBAllowlist: + """The contract's directory, the declared source, and the declared landing.""" + allow = _grant_source(DuckDBAllowlist.none().with_dirs(ctx.workdir), ctx, streams) + if sink_format is None: + return allow + return _grant_destination(allow, ctx, streams, sink_format) + + # ── Runner ─────────────────────────────────────────────────────────────── @@ -587,13 +666,15 @@ def replay(self, ctx: RunContext, run_id: str) -> RunResult: return _execute(ctx, self) def fingerprint(self, ctx: RunContext) -> SchemaFingerprint: - import duckdb - from .._masking import masked_column_types - con = duckdb.connect(":memory:") + # Reads the source only: the destination is not granted. + con = secure_duckdb_connect( + ":memory:", + allow=_run_allowlist(ctx, _run_streams(ctx)), + before_lock=lambda c: self._load_extensions(c, ctx), + ) try: - self._load_extensions(con, ctx) select_sql = _select_for_first_stream(ctx) con.execute(f"CREATE TEMP VIEW _fp AS {select_sql}") rows = con.execute("DESCRIBE _fp").fetchall() @@ -654,6 +735,11 @@ def _attach_external_databases(self, con: Any, ctx: RunContext) -> None: Best-effort: if the connection DSN is malformed or the upstream is unreachable, the ATTACH raises and the runner surfaces the error in the per-stream try/except. + + A sqlite file is confined to :func:`_declared_roots` plus the + operator's ``FLUID_DUCKDB_ALLOWED_DIRS`` first, and refused with + ``DuckDBSandboxError`` outside them: the sqlite scanner is not bounded + by the sandbox's allowlist, so nothing else would stop it. """ kind = ctx.source.kind if kind in ("mysql", "mariadb"): @@ -673,8 +759,14 @@ def _attach_external_databases(self, con: Any, ctx: RunContext) -> None: path = conn.get("uri") or conn.get("path") or conn.get("database") or "" if not path: raise ValueError("sqlite source requires connection.uri, .path, or .database") + # The sqlite scanner opens the file through its own client library, + # which allowed_directories never bounds: this confinement is the + # only check on it. Without it a contract could land any SQLite + # file on the host (a browser's cookie store, another app's data). + # The checked realpath is what is attached. + real = confine_declared(str(path), within=_declared_roots(ctx)) alias = _sqlite_alias_for_build(ctx.build_id) - con.execute(f"ATTACH {quote_ansi_string_literal(str(path))} AS {alias} (TYPE sqlite)") + con.execute(f"ATTACH {quote_ansi_string_literal(real)} AS {alias} (TYPE sqlite)") def _select_for_first_stream(ctx: RunContext) -> str: @@ -1218,9 +1310,12 @@ def _execute(ctx: RunContext, runner: DuckdbRunner) -> RunResult: dlq_writer: Optional[DLQWriter] = None dlq_total_records = 0 - con = duckdb.connect(":memory:") + con = secure_duckdb_connect( + ":memory:", + allow=_run_allowlist(ctx, streams_to_run, sink_format), + before_lock=lambda c: runner._load_extensions(c, ctx), + ) try: - runner._load_extensions(con, ctx) if masker is not None: masker.install(con) for stream in streams_to_run: @@ -1681,20 +1776,26 @@ def _connect_for_destination(ctx: RunContext) -> Any: swallow their errors, so they would simply have stopped doing their job — no late-arrival split, and an empty PII scan that looks like a clean one. """ - import duckdb - - con = duckdb.connect(":memory:") dest = _binding_destination_uri(ctx) - if not dest: - return con - for ext in _required_extensions("filesystem", dest): - try: - con.execute(f"INSTALL {ext}") - con.execute(f"LOAD {ext}") - except Exception as exc: # noqa: BLE001 - LOG.warning("DuckDB extension load failed (%s): %s", ext, type(exc).__name__) - _apply_destination_secret(con, ctx, dest) - return con + + def _prepare(con: Any) -> None: + if not dest: + return + for ext in _required_extensions("filesystem", dest): + try: + con.execute(f"INSTALL {ext}") + con.execute(f"LOAD {ext}") + except Exception as exc: # noqa: BLE001 + LOG.warning("DuckDB extension load failed (%s): %s", ext, type(exc).__name__) + _apply_destination_secret(con, ctx, dest) + + # Reads back (and the late-arrival split rewrites) what this run landed, + # and nothing else: the source is not granted here. + sink_format = (ctx.sink.format or "parquet").lower() + allow = _grant_destination( + DuckDBAllowlist.none().with_dirs(ctx.workdir), ctx, _run_streams(ctx), sink_format + ) + return secure_duckdb_connect(":memory:", allow=allow, before_lock=_prepare) def _artifact_exists(con: Any, path: str, reader: str) -> bool: @@ -1789,12 +1890,10 @@ def _count_file_rows(path: str, sink_format: str) -> int: The late-arrival split can rewrite the file after the build counted it. """ - import duckdb - reader = {"parquet": "read_parquet", "csv": "read_csv_auto", "json": "read_json_auto"}[ sink_format ] - con = duckdb.connect(":memory:") + con = secure_duckdb_connect(":memory:", allow=DuckDBAllowlist.none().with_paths(path)) try: row = con.execute( f"SELECT COUNT(*) FROM {reader}({quote_ansi_string_literal(path)})" diff --git a/fluid_build/build_runners/meltano/runner.py b/fluid_build/build_runners/meltano/runner.py index a04d51f9..5ff6a5fa 100644 --- a/fluid_build/build_runners/meltano/runner.py +++ b/fluid_build/build_runners/meltano/runner.py @@ -484,11 +484,12 @@ def write_records_to_duckdb( like ``"orders; DROP TABLE secrets; --"`` is rejected at the boundary rather than executed. """ - import duckdb + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect dataset = validate_ident(dataset) duckdb_path.parent.mkdir(parents=True, exist_ok=True) - con = duckdb.connect(str(duckdb_path)) + # Records arrive as values, so the database file is all this touches. + con = secure_duckdb_connect(duckdb_path, allow=DuckDBAllowlist.none()) try: con.execute(f"CREATE SCHEMA IF NOT EXISTS {dataset}") counts: Dict[str, int] = {} diff --git a/fluid_build/cli/_diff_live.py b/fluid_build/cli/_diff_live.py index 9b95aa64..4a3ad909 100644 --- a/fluid_build/cli/_diff_live.py +++ b/fluid_build/cli/_diff_live.py @@ -644,14 +644,13 @@ def _duckdb_table_exists(path: Path, schema_name: str, table: str) -> bool: checked by the provider's own allowlist first and are bound as parameters, never spliced into the SQL. """ - import duckdb - + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect from fluid_build.providers.local_validation import _build_duckdb_table_ref _build_duckdb_table_ref(schema_name, table) # Only the catalog is read, so the connection gets no file or network # access beyond the database file itself. - con = duckdb.connect(str(path), read_only=True, config={"enable_external_access": False}) + con = secure_duckdb_connect(path, allow=DuckDBAllowlist.none(), read_only=True) try: row = con.execute( "SELECT count(*) FROM information_schema.tables " diff --git a/fluid_build/cli/discover/_jdbc_base.py b/fluid_build/cli/discover/_jdbc_base.py index 1f2a1eb8..4268daf0 100644 --- a/fluid_build/cli/discover/_jdbc_base.py +++ b/fluid_build/cli/discover/_jdbc_base.py @@ -119,18 +119,24 @@ def _parse_uri(self, uri: str) -> dict: } def _introspect(self, conn: dict) -> List[DiscoveredStream]: - import duckdb + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect dsn = build_libpq_dsn(conn, database_key=self.config.database_dsn_key) alias = self.config.attach_alias - - con = duckdb.connect(":memory:") + attach_sql = ( + f"ATTACH {quote_ansi_string_literal(dsn)} AS {alias} (TYPE {self.config.attach_type})" + ) + + # The source is attached before the sandbox locks; the catalog queries + # after it reach that database and no file or URL. + con = secure_duckdb_connect( + ":memory:", + allow=DuckDBAllowlist.none(), + extensions=[self.config.extension], + before_lock=lambda c: c.execute(attach_sql), + ) streams: List[DiscoveredStream] = [] try: - con.execute(f"INSTALL {self.config.extension}; LOAD {self.config.extension};") - con.execute( - f"ATTACH {quote_ansi_string_literal(dsn)} AS {alias} (TYPE {self.config.attach_type})" - ) filter_fn = self.config.table_filter or _default_table_filter where = filter_fn(conn, quote_ansi_string_literal) diff --git a/fluid_build/cli/discover/_jdbc_introspect.py b/fluid_build/cli/discover/_jdbc_introspect.py index 113d0dcb..16c2b18d 100644 --- a/fluid_build/cli/discover/_jdbc_introspect.py +++ b/fluid_build/cli/discover/_jdbc_introspect.py @@ -618,19 +618,23 @@ def introspect_jdbc( f"Invalid schema filter: {schema_filter!r}. " "Must match ``[A-Za-z_][A-Za-z0-9_]*``." ) - con = duckdb.connect(":memory:") + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect + + attach = { + "postgres": _attach_string_postgres, + "mysql": _attach_string_mysql, + "sqlite": _attach_string_sqlite, + }.get(kind) + # Install + load the extension (duckdb caches the binary, so this is fast + # on a second run) and ATTACH the source before the sandbox locks; the + # introspection queries after it reach that database and no file or URL. + con = secure_duckdb_connect( + ":memory:", + allow=DuckDBAllowlist.none(), + extensions=[kind] if attach is not None else [], + before_lock=(lambda c: c.execute(attach(args, alias))) if attach is not None else None, + ) try: - # Install + load the extension. duckdb caches the binary so - # this is fast on second run. - if kind == "postgres": - con.execute("INSTALL postgres; LOAD postgres;") - con.execute(_attach_string_postgres(args, alias)) - elif kind == "mysql": - con.execute("INSTALL mysql; LOAD mysql;") - con.execute(_attach_string_mysql(args, alias)) - elif kind == "sqlite": - con.execute("INSTALL sqlite; LOAD sqlite;") - con.execute(_attach_string_sqlite(args, alias)) # Enumerate via duckdb's union ``information_schema`` view, # filtered to our attached database. Postgres + MySQL each diff --git a/fluid_build/cli/discover/filesystem.py b/fluid_build/cli/discover/filesystem.py index f2d4cbb4..f47794fd 100644 --- a/fluid_build/cli/discover/filesystem.py +++ b/fluid_build/cli/discover/filesystem.py @@ -88,18 +88,29 @@ def _format_for(path: str) -> str: def _columns_for(path: str, fmt: str) -> List[DiscoveredColumn]: - import duckdb - - con = duckdb.connect(":memory:") + from fluid_build.providers._duckdb_sandbox import ( + DuckDBAllowlist, + is_remote_location, + secure_duckdb_connect, + ) + from fluid_build.providers._sql_safety import quote_ansi_string_literal + + reader = {"csv": "read_csv_auto", "parquet": "read_parquet", "json": "read_json_auto"}.get(fmt) + if reader is None: + return [] + # The one location being discovered. A URL needs its filesystem extension + # loaded before the sandbox locks (autoloading is off from then on). + remote = is_remote_location(path) + extensions: List[str] = [] + if remote: + extensions = ["azure"] if path.lower().startswith(("azure://", "az://")) else ["httpfs"] + con = secure_duckdb_connect( + ":memory:", + allow=DuckDBAllowlist.none().with_locations(path), + extensions=extensions, + ) try: - if fmt == "csv": - sql = f"SELECT * FROM read_csv_auto('{path}') LIMIT 0" - elif fmt == "parquet": - sql = f"SELECT * FROM read_parquet('{path}') LIMIT 0" - elif fmt == "json": - sql = f"SELECT * FROM read_json_auto('{path}') LIMIT 0" - else: - return [] + sql = f"SELECT * FROM {reader}({quote_ansi_string_literal(path)}) LIMIT 0" con.execute(sql) descr = con.description or [] return [DiscoveredColumn(name=c[0], type=str(c[1]), nullable=True) for c in descr] diff --git a/fluid_build/cli/forge_copilot_schema_inference.py b/fluid_build/cli/forge_copilot_schema_inference.py index f6e1f8e9..c4f8513e 100644 --- a/fluid_build/cli/forge_copilot_schema_inference.py +++ b/fluid_build/cli/forge_copilot_schema_inference.py @@ -373,9 +373,12 @@ def _read_parquet_metadata_pyarrow(path: Path) -> Dict[str, Any]: def _read_parquet_metadata_duckdb(path: Path) -> Dict[str, Any]: - import duckdb + import duckdb # noqa: F401 - absent duckdb raises ImportError for the caller - connection = duckdb.connect() + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect + + # The sample file and nothing else. + connection = secure_duckdb_connect(allow=DuckDBAllowlist.none().with_paths(path)) try: rows = connection.execute("DESCRIBE SELECT * FROM read_parquet(?)", [str(path)]).fetchall() finally: diff --git a/fluid_build/cli/forge_db_tools.py b/fluid_build/cli/forge_db_tools.py index 14bbff2f..3aabd5fb 100644 --- a/fluid_build/cli/forge_db_tools.py +++ b/fluid_build/cli/forge_db_tools.py @@ -75,7 +75,7 @@ * ``FLUID_FORGE_DB_TOOLS`` — ``1``/``true`` to expose the tool (default off → ABSENT from ``get_tool_definitions``). * ``FLUID_FORGE_DB_URI`` — the ``default`` connection's URI, e.g. - ``postgresql://user:pass@host:5432/db`` / ``mysql://user:pass@host/db`` / + ``postgresql://user:$PASSWORD@host:5432/db`` / ``mysql://user:$PASSWORD@host/db`` / ``sqlite:////abs/path.db``. * ``FLUID_FORGE_DB_URI_`` — a named connection reachable via ``connection=`` (case-insensitive; the alias is upper-cased). @@ -336,22 +336,24 @@ def _fetch_sample_rows(arguments: Dict[str, Any], env: Mapping[str, str]) -> Dic sql = f"SELECT * FROM {ref} LIMIT {applied_limit}" try: - import duckdb + import duckdb # noqa: F401 - the install hint below except ImportError: return { "error": "DuckDbNotInstalled", "message": "fetch_sample_rows requires duckdb (pip install 'fluid-build[local]').", } + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect - con = duckdb.connect(":memory:") + con = None try: - if kind == "postgres": - con.execute("INSTALL postgres; LOAD postgres;") - elif kind == "mysql": - con.execute("INSTALL mysql; LOAD mysql;") - else: - con.execute("INSTALL sqlite; LOAD sqlite;") - con.execute(attach) + # The source is attached before the sandbox locks; the SELECT that + # follows then reaches that database and no file or URL. + con = secure_duckdb_connect( + ":memory:", + allow=DuckDBAllowlist.none(), + extensions=[{"postgres": "postgres", "mysql": "mysql"}.get(kind, "sqlite")], + before_lock=lambda c: c.execute(attach), + ) cur = con.execute(sql) raw_columns = [str(d[0]) for d in (cur.description or [])] raw_rows = cur.fetchall() @@ -364,7 +366,8 @@ def _fetch_sample_rows(arguments: Dict[str, Any], env: Mapping[str, str]) -> Dic "message": f"fetch_sample_rows could not read {table!r} — see server logs", } finally: - con.close() + if con is not None: + con.close() # Second row ceiling (belt-and-suspenders, mirrors Command Center's # nl_query: LIMIT rewrite AND a fetch ceiling): even if the SQL LIMIT were diff --git a/fluid_build/cli/init.py b/fluid_build/cli/init.py index 9effef30..2e8f74a2 100644 --- a/fluid_build/cli/init.py +++ b/fluid_build/cli/init.py @@ -1200,13 +1200,16 @@ def init_local_db(project_dir: Path, provider: str, logger: logging.Logger): return # Only for local provider try: - import duckdb + import duckdb # noqa: F401 - absent duckdb is the ImportError below + + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect db_dir = project_dir / ".fluid" db_dir.mkdir(exist_ok=True) db_path = db_dir / "db.duckdb" - conn = duckdb.connect(str(db_path)) + # Creates the file and closes it: no SQL runs, so no access is granted. + conn = secure_duckdb_connect(db_path, allow=DuckDBAllowlist.none()) conn.close() if RICH_AVAILABLE: diff --git a/fluid_build/cli/verify.py b/fluid_build/cli/verify.py index cd00a492..8cec840e 100644 --- a/fluid_build/cli/verify.py +++ b/fluid_build/cli/verify.py @@ -960,9 +960,14 @@ def _verify_local_file( # Introspect with DuckDB. try: - import duckdb # type: ignore[import] + from fluid_build.providers._duckdb_sandbox import ( + DuckDBAllowlist, + secure_duckdb_connect, + ) - con = duckdb.connect(":memory:") + # The masking check below runs SQL built from the contract: it reads + # the one file being verified and nothing else. + con = secure_duckdb_connect(allow=DuckDBAllowlist.none().with_paths(file_path)) # A SQL string literal, not Python ``repr``: the path now carries the # contract's directory, and ``repr`` switches to double quotes (a # DuckDB identifier) or backslash escapes (literal in DuckDB) the diff --git a/fluid_build/contract_tests.py b/fluid_build/contract_tests.py index fd63e8cc..f36759c6 100644 --- a/fluid_build/contract_tests.py +++ b/fluid_build/contract_tests.py @@ -53,6 +53,14 @@ from pathlib import Path from typing import Any, Dict, List, Optional, Tuple, Union +from fluid_build.providers._duckdb_sandbox import ( + DuckDBAllowlist, + DuckDBSandboxError, + is_sandbox_refusal, + sandbox_refusal_hint, + secure_duckdb_connect, +) + LOGGER = logging.getLogger("fluid.provider.local") # ------------------------- @@ -95,19 +103,52 @@ def _require_duckdb(): ) from e -def _connect_duckdb(): +def _connect_duckdb(allow: Optional[DuckDBAllowlist] = None): + """A sandboxed DuckDB that may touch only ``allow`` (nothing, by default). + + The action's SQL is contract input, so the connection reaches the files the + action declares (``_action_allowlist``) and no others. + """ # Allow persistent db file for debugging/local exploration: # FLUID_LOCAL_DUCKDB_PATH=/tmp/fluid_local.duckdb _require_duckdb() - import duckdb # type: ignore db_path = os.environ.get("FLUID_LOCAL_DUCKDB_PATH", ":memory:") try: - return duckdb.connect(db_path, read_only=False) + return secure_duckdb_connect( + db_path, allow=allow if allow is not None else DuckDBAllowlist.none() + ) except Exception as e: raise LocalProviderError(f"Failed to connect to duckdb at '{db_path}': {e}") from e +def _output_specs(outputs: Dict[str, Any]) -> List[Dict[str, Any]]: + if isinstance(outputs, dict) and "path" in outputs: + return [outputs] + targets = outputs.get("targets", []) if isinstance(outputs, dict) else [] + return [t for t in targets if isinstance(t, dict)] if isinstance(targets, list) else [] + + +def _action_allowlist(inputs: Dict[str, Any], outputs: Dict[str, Any]) -> DuckDBAllowlist: + """Exactly what one action declares: each input file and each output file. + + Each must resolve inside the working directory (where the action's + relative paths resolve) or a directory the operator lists in + ``FLUID_DUCKDB_ALLOWED_DIRS``: a declared file is granted to the action's + SQL, so declaring ``~/.aws/credentials`` must not make it readable. + """ + within = [Path.cwd()] + allow = DuckDBAllowlist.none() + for cfg in inputs.values(): + if isinstance(cfg, dict) and cfg.get("path"): + files = _glob_all(_as_list(cfg["path"])) + allow = allow.with_declared(*(str(Path(f).absolute()) for f in files), within=within) + for spec in _output_specs(outputs): + if spec.get("path"): + allow = allow.with_declared(str(Path(str(spec["path"])).absolute()), within=within) + return allow + + def _register_input(con, alias: str, cfg: Dict[str, Any]) -> str: """ Register an input as a DuckDB view. Supports CSV/Parquet. @@ -354,15 +395,39 @@ def apply_action(action: Dict[str, Any], ctx) -> None: if not outputs: raise LocalProviderError(f"Action '{rid}' missing outputs block") - # Connect to DuckDB - con = _connect_duckdb() + # Connect to DuckDB, confined to what this action declares + try: + allow = _action_allowlist(inputs, outputs) + except DuckDBSandboxError as e: + raise LocalProviderError(f"Action '{rid}': {e}") from e + con = _connect_duckdb(allow) + # Closed on every path: with FLUID_LOCAL_DUCKDB_PATH set, a connection left + # open by a failure holds the file's locked instance and refuses the next. + try: + _run_action_sql(con, rid, sql, inputs, outputs, allow, ctx) + finally: + try: + con.close() + except Exception: # the action's own outcome is what the caller needs + LOGGER.debug("closing the DuckDB connection of action %r failed", rid, exc_info=True) + +def _run_action_sql( + con, rid: str, sql: str, inputs: Dict[str, Any], outputs: Dict[str, Any], allow, ctx +) -> None: + """Register the inputs, run the SQL and write the outputs on ``con``.""" # Register inputs for alias, cfg in inputs.items(): _register_input(con, alias, cfg) # Execute SQL - rel = _execute_sql(con, sql) + try: + rel = _execute_sql(con, sql) + except LocalProviderError as e: + refused = e.__cause__ or e + if is_sandbox_refusal(refused): + raise LocalProviderError(f"{e} {sandbox_refusal_hint(allow, refused)}") from e + raise # Write output(s) # We support a single output dict or a dict with 'path' etc. @@ -383,11 +448,6 @@ def apply_action(action: Dict[str, Any], ctx) -> None: LOGGER.info(json_log("apply_action_output", rows=rc, path=path or "(dry-run)")) LOGGER.info(json_log("apply_action_outputs_total", rows=total)) - try: - con.close() - except Exception: - pass - # ------------------------- # Logging helper diff --git a/fluid_build/output_ports/mcp/drivers/duckdb.py b/fluid_build/output_ports/mcp/drivers/duckdb.py index 8c02b57b..6b97c82d 100644 --- a/fluid_build/output_ports/mcp/drivers/duckdb.py +++ b/fluid_build/output_ports/mcp/drivers/duckdb.py @@ -216,22 +216,24 @@ def _get_connection(self): if self._connection is not None: return self._connection try: - import duckdb # type: ignore[import-not-found] + import duckdb # type: ignore[import-not-found] # noqa: F401 except ImportError as exc: # pragma: no cover - depends on optional dep raise UnsupportedBindingError( "duckdb is not installed; install via the 'local' extra: " "pip install 'data-product-forge[local]'" ) from exc + from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect + target = str(self._db_file) if self._db_file is not None else ":memory:" + # The bound file (if any) and nothing else: a query the MCP port + # compiles cannot reach another file, URL or database. + allow = DuckDBAllowlist.none().with_paths(self._path) # DuckDB read-only mode protects against accidental writes # even though every advertised tool is SELECT-only. Skip # read-only when the target is in-memory because a fresh # in-memory instance is empty and would refuse the table-load # below. - if target == ":memory:": - connection = duckdb.connect(database=":memory:") - else: - connection = duckdb.connect(database=target, read_only=True) + connection = secure_duckdb_connect(target, allow=allow, read_only=target != ":memory:") self._configure_connection(connection) self._connection = connection return connection diff --git a/fluid_build/providers/_duckdb_sandbox.py b/fluid_build/providers/_duckdb_sandbox.py new file mode 100644 index 00000000..86936257 --- /dev/null +++ b/fluid_build/providers/_duckdb_sandbox.py @@ -0,0 +1,736 @@ +# Copyright 2024-2026 Agentics Transformation Ltd +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""The one way the engine opens DuckDB: every connection is sandboxed. + +A contract carries SQL (``builds[].properties.sql``, stage SQL, quality +predicates, masking checks) and DuckDB runs it with the privileges of the +process. Unconfined, that SQL reads any file the process can: +``read_csv('/etc/passwd')``, ``read_text('~/.aws/credentials')``, ``glob('/')``, +``ATTACH`` of another database, ``COPY ... TO`` anywhere, an ``https://`` URL. +A service that runs a user's contract would hand that user its host. + +:func:`secure_duckdb_connect` applies DuckDB's own sandbox, in the order the +DuckDB docs give ("Securing DuckDB", introduced in 1.2): + +1. connect (the database file itself is opened here, before any restriction, + and needs no grant of its own); +2. ``allow_persistent_secrets = false`` and ``allow_community_extensions = + false``, then load what the call site needs: its extensions, and anything + its ``before_lock`` callback sets up (an ``ATTACH`` of a declared source + database, an object-store secret); +3. ``autoinstall_known_extensions = false`` / ``autoload_known_extensions = + false``, so a query cannot pull an extension in afterwards; +4. ``home_directory`` (pinned to ``$HOME``, see below), then + ``allowed_directories`` and ``allowed_paths``: only what the call site + declares it reads or writes; +5. ``enable_external_access = false``; +6. ``lock_configuration = true``, so the SQL cannot ``SET`` any of it back. + +Each setting is a ``SET`` statement on the connection, not the ``config=`` dict: +the dict does not take the list-valued ``allowed_directories``. + +What the sandbox can and cannot promise (and why the floor is DuckDB 1.5.0): + +* ``allowed_directories`` arrived in DuckDB 1.2. Up to and including 1.4.3 it + did not resolve ``./..`` before the allowlist check, so ``/./../x`` + escaped it. Through 1.4.x the check is lexical: a symlink inside an allowed + directory reads its target anywhere on the host, and a relative path is + refused outright. 1.5.0 resolves both (each reproduced on 1.4.3 / 1.4.4 and + refused on 1.5.0; see ``tests/providers/test_duckdb_sandbox.py``). + :data:`MIN_DUCKDB_VERSION` is therefore 1.5.0, and a connection refuses to + open on anything older. +* duckdb/duckdb#26064 (open): ``~`` is expanded with ``$HOME`` by the check and + with the ``home_directory`` setting by the open, so once the two differ they + name different files and the check passes for one while the other is read. + ``home_directory`` is therefore pinned to ``$HOME`` itself and locked: both + expansions agree, and ``~/x`` is readable exactly when ``$HOME/x`` is + allowed. A ``~`` in an allowlist entry is refused outright. +* DuckDB 1.5 does not resolve ``..`` inside a remote URL, so a remote prefix + (``s3://bucket/data/``) bounds the bucket or host, not the path within it. + Remote prefixes are granted only for a location the contract declares. +* A location the contract DECLARES (an input, an output, a source URI) is + granted to its SQL, so a declaration is itself a request for access. + :meth:`DuckDBAllowlist.with_declared` grants one only inside the roots its + call site names (the contract's directory, its workspace, the run's scratch + directory) or a directory the operator lists in ``FLUID_DUCKDB_ALLOWED_DIRS`` + (:data:`OPERATOR_DIRS_ENV`); a contract field cannot widen that. Both sides + are compared after ``realpath``, because DuckDB realpaths allowlist entries + too: an in-repo symlink to ``/`` would otherwise grant the host. A glob is + granted as the directory above its first wildcard, which must itself resolve + inside those roots: the contract could declare any file there anyway, and + DuckDB expands a glob to files a grant-time listing would miss (dotfiles, + files that land after the grant), then checks each one's realpath against + the directory, so a matched symlink that leads out is refused at read time. +* A directory the engine grants by convention rather than by declaration + (the local provider's ``./runtime``) sits in a working directory the + contract's repository may supply, so it may be a symlink to ``$HOME``. + :func:`unaliased_dir` grants it only where its name says it is, or inside + the call site's roots. +* Autoloading is off, so a function from an extension the call site did not + load (``sqlite_scan``, ``read_xlsx``, ``ST_Read``, ``delta_scan``, + ``iceberg_scan``) is not in the catalog. :func:`is_sandbox_refusal` treats + that error as a refusal. +* The ``sqlite``, ``postgres`` and ``mysql`` scanners open files and sockets + through their own client libraries, not through DuckDB's file system, so + neither ``allowed_directories`` nor ``enable_external_access`` bounds them: + once ``sqlite`` is loaded, ``sqlite_scan`` / ``ATTACH ... (TYPE sqlite)`` + open any SQLite file the process can read, and once ``postgres`` is loaded, + ``postgres_scan`` reaches any host the process can. Only the engine's own + SQL runs on a connection that loads them; contract SQL never does. +* DuckDB shares one database instance per file per process, and the lock is + instance-wide: a second connection to a file database that is already open + cannot be sandboxed, and is refused with :class:`DuckDBSandboxError`. Each + call site closes its connection before the next opens the same file. +* The DuckDB docs call these settings defense-in-depth, "not a substitute for + proper sandboxing": a multi-tenant host still runs each contract in its own + container. + +References: + https://duckdb.org/docs/current/operations_manual/securing_duckdb/overview.html + https://github.com/duckdb/duckdb/issues/26064 + https://github.com/bordumb/dataing/pull/176 (the 1.4.3 ``./..`` escape) + https://github.com/duckdb/duckdb/blob/v1.4.4/src/main/config.cpp + (``DBConfig::CanAccessFile``: the 1.4 check is lexical) +""" + +from __future__ import annotations + +import os +import re +from dataclasses import dataclass +from pathlib import Path +from typing import TYPE_CHECKING, Any, Callable, Iterable, Mapping, Optional, Tuple, Union + +from fluid_build._errors import FluidUserError, doc_url + +from ._sql_safety import quote_ansi_string_literal, validate_ident + +if TYPE_CHECKING: # pragma: no cover - typing only + import duckdb + +PathLike = Union[str, "os.PathLike[str]"] + +#: Oldest DuckDB whose ``allowed_directories`` resolves ``./..`` and symlinks. +MIN_DUCKDB_VERSION: Tuple[int, int, int] = (1, 5, 0) + +#: Operator-only widening: directories (``os.pathsep``-separated) where a +#: contract may also declare inputs and outputs. An environment variable, not a +#: contract field, so the author of the contract cannot set it. +OPERATOR_DIRS_ENV = "FLUID_DUCKDB_ALLOWED_DIRS" + +_GLOB_CHARS = frozenset("*?[") +_REMOTE_RE = re.compile(r"^(?P[a-z][a-z0-9+.-]*)://(?P[^/?#]+)(?P/[^?#]*)?$") +_REMOTE_SCHEMES = frozenset( + {"s3", "s3a", "s3n", "gs", "gcs", "r2", "azure", "az", "abfss", "http", "https", "hf"} +) + + +@dataclass +class DuckDBSandboxError(FluidUserError): + """A DuckDB connection could not be opened inside the sandbox.""" + + code: str = "DuckDBSandboxError" + + +def _sandbox_error(what: str, why: str, fix: str) -> DuckDBSandboxError: + return DuckDBSandboxError(what=what, why=why, fix=fix, doc=doc_url()) + + +def is_remote_location(location: str) -> bool: + """Whether ``location`` is a URL DuckDB reads through an extension.""" + match = _REMOTE_RE.match(str(location)) + return bool(match and match.group("scheme").lower() in _REMOTE_SCHEMES) + + +def _first_glob(location: str, start: int = 0) -> Optional[int]: + return next( + (i for i in range(start, len(location)) if location[i] in _GLOB_CHARS), + None, + ) + + +def _static_prefix(location: str) -> Tuple[str, bool]: + """``location`` cut before its first glob segment, and whether it had one. + + Either separator ends a segment, so a Windows path globs as a POSIX one. + """ + first = _first_glob(location) + if first is None: + return location, False + cut = max(location.rfind("/", 0, first), location.rfind("\\", 0, first)) + if cut < 0: + return ".", True # a relative glob in the working directory + return location[:cut] or location[: cut + 1], True + + +@dataclass(frozen=True) +class DuckDBAllowlist: + """What one DuckDB connection may read and write, and nothing else. + + ``dirs`` are directories (everything below them), ``paths`` single files, + ``remote_prefixes`` URL prefixes such as ``s3://bucket/landing/``. Build one + with :meth:`none` and the ``with_*`` methods; every call site names its own. + """ + + dirs: Tuple[str, ...] = () + paths: Tuple[str, ...] = () + remote_prefixes: Tuple[str, ...] = () + + @classmethod + def none(cls) -> "DuckDBAllowlist": + """No file or network access at all: in-memory work only.""" + return cls() + + def with_dirs(self, *dirs: Optional[PathLike]) -> "DuckDBAllowlist": + """Also allow everything under each directory in ``dirs`` (``None`` skipped).""" + added = _new(self.dirs, (d for d in dirs if d is not None and str(d) != ""), _local) + return DuckDBAllowlist(self.dirs + added, self.paths, self.remote_prefixes) + + def with_paths(self, *paths: Optional[PathLike]) -> "DuckDBAllowlist": + """Also allow each single file in ``paths`` (``None`` skipped).""" + added = _new(self.paths, (p for p in paths if p is not None and str(p) != ""), _local) + return DuckDBAllowlist(self.dirs, self.paths + added, self.remote_prefixes) + + def with_remote(self, *prefixes: Optional[str]) -> "DuckDBAllowlist": + """Also allow each URL prefix (``s3://bucket/``, ``https://host/data/``).""" + added = _new(self.remote_prefixes, (p for p in prefixes if p), _remote) + return DuckDBAllowlist(self.dirs, self.paths, self.remote_prefixes + added) + + def with_locations( + self, *locations: Optional[PathLike], base: Optional[PathLike] = None + ) -> "DuckDBAllowlist": + """Also allow each location a contract DECLARES it reads or writes. + + For a location the operator names (a ``fluid discover`` argument); a + location a CONTRACT declares goes through :meth:`with_declared`, which + also confines it. A URL grants its directory prefix (cut before any + glob); a local glob grants the directory above its first wildcard; an + existing directory grants itself; anything else grants that one file. + A leading ``~`` is expanded as DuckDB expands it, and a relative local + path is resolved against ``base`` (else the working directory, where + DuckDB would resolve it). + """ + out = self + for location in locations: + if location is None or str(location) == "": + continue + raw = os.fspath(location) + if is_remote_location(raw): + prefix, globbed = _static_prefix(raw) + match = _REMOTE_RE.match(prefix) + if match and not (match.group("path") or "").strip("/"): + # a bare bucket or host: the whole of it + prefix = f"{match.group('scheme')}://{match.group('host')}/" + elif globbed: + prefix = prefix.rstrip("/") + "/" + elif not prefix.endswith("/"): + prefix = prefix.rsplit("/", 1)[0] + "/" + out = out.with_remote(prefix) + continue + out = out._with_local(_absolute(raw, base), roots=None, declared=raw) + return out + + def with_declared( + self, + *locations: Optional[PathLike], + within: Iterable[Optional[PathLike]], + base: Optional[PathLike] = None, + ) -> "DuckDBAllowlist": + """Also allow each location a CONTRACT declares, inside ``within`` only. + + A declared input or output is granted to the contract's SQL, so the + declaration must not be a way to read the host: each local location + must resolve (``realpath``, as DuckDB resolves allowlist entries) inside + one of ``within`` (the call site's own roots, such as the contract's + directory and its workspace) or a directory the operator lists in + ``FLUID_DUCKDB_ALLOWED_DIRS``; anything else raises + :class:`DuckDBSandboxError`. A glob is granted as the directory above + its first wildcard, which must resolve inside those roots: everything + in it is declarable already, and DuckDB refuses a matched symlink that + leads out of it. A relative path is resolved against ``base``, else + the working directory, which is where DuckDB opens it. + A URL is granted as :meth:`with_locations` grants it. + """ + roots = _confinement_roots(within) + out = self + for location in locations: + if location is None or str(location) == "": + continue + raw = os.fspath(location) + if is_remote_location(raw): + out = out.with_locations(raw) + continue + out = out._with_local(_absolute(raw, base), roots=roots, declared=raw) + return out + + def _with_local( + self, path: str, *, roots: Optional[Tuple[str, ...]], declared: str + ) -> "DuckDBAllowlist": + """Grant one absolute local ``path``; with ``roots``, only inside them.""" + prefix, globbed = _local_glob_prefix(path) + if roots is not None: + _confine(prefix, roots, declared) + if globbed or Path(prefix).is_dir(): + # A glob grants the directory above its first wildcard, not the + # files a listing finds now: DuckDB's own expansion includes + # dotfiles (``._x.csv``) and files that land after the grant, and + # checks every one, so a grant of only today's matches fails the + # whole read. With ``roots``, that directory is confined (above), + # and DuckDB refuses a file in it whose realpath leads out. + return self.with_dirs(prefix) + return self.with_paths(prefix) + + +def _local_glob_prefix(path: str) -> Tuple[str, bool]: + """Absolute local ``path`` cut before its first glob segment, and whether it had one. + + A directory that exists under its literal name is a directory, not a + pattern: ``Proj [old]`` in a contract's directory, the working directory + or ``$HOME`` is where the contract lives, and the callers join it in + before this sees the path. Cut there, ``customers.csv`` would be confined + (and refused) as the directory above the project. DuckDB matches such a + component as a pattern first and opens it literally when nothing matches; + a sibling the pattern does match (``Proj o``) is not granted, so that read + is refused rather than widened. The last component stays a pattern even + when a file has that literal name: DuckDB reads its matches (``a1.csv`` + for ``a[1].csv``), and the directory granted holds both. + """ + start = 0 + while True: + first = _first_glob(path, start) + if first is None: + return path, False + end = min( + (i for i in (path.find("/", first), path.find("\\", first)) if i >= 0), + default=-1, + ) + if end < 0 or not os.path.lexists(path[:end]): + cut = max(path.rfind("/", 0, first), path.rfind("\\", 0, first)) + if cut < 0: + return ".", True + return path[:cut] or path[: cut + 1], True + start = end + + +def _absolute(raw: str, base: Optional[PathLike]) -> str: + """``raw`` with ``~`` expanded as DuckDB expands it, made absolute at ``base``.""" + # '~' as DuckDB itself would expand it (home_directory is $HOME). + path = Path(raw).expanduser() + if not path.is_absolute(): + path = Path(base if base is not None else os.getcwd()) / path + # Lexically normalised first, so '/*/../../etc/x' is confined as the + # '/etc/x' it names rather than as the '' its glob prefix suggests. + return os.path.normpath(str(path)) + + +def _is_within(real: str, root: str) -> bool: + return real == root or real.startswith(_with_sep(root)) + + +def operator_allowed_dirs() -> Tuple[str, ...]: + """The directories ``FLUID_DUCKDB_ALLOWED_DIRS`` adds, each ``realpath``-ed. + + Set by whoever runs the engine (``FLUID_DUCKDB_ALLOWED_DIRS=/shared/ref``), + never by a contract. A relative entry, or one that resolves to the + filesystem root, is refused rather than guessed at. + """ + out: list = [] + for part in os.environ.get(OPERATOR_DIRS_ENV, "").split(os.pathsep): + part = part.strip() + if not part: + continue + expanded = os.path.expanduser(part) + if not os.path.isabs(expanded): + raise _sandbox_error( + what=f"{OPERATOR_DIRS_ENV} entry {part!r} is not an absolute path", + why="A relative entry would mean a different directory in every working directory.", + fix=f"Set {OPERATOR_DIRS_ENV} to absolute directories, separated by {os.pathsep!r}.", + ) + real = os.path.realpath(expanded) + if real == os.path.realpath(os.sep): + raise _sandbox_error( + what=f"{OPERATOR_DIRS_ENV} entry {part!r} is the filesystem root", + why="Allowing '/' lets a contract declare, and so read, every file.", + fix="List the directories contracts may read and write, not '/'.", + ) + if real not in out: + out.append(real) + return tuple(out) + + +def _confinement_roots(within: Iterable[Optional[PathLike]]) -> Tuple[str, ...]: + roots: list = [] + for root in [*within, *operator_allowed_dirs()]: + if root is None or str(root) == "": + continue + real = os.path.realpath(os.path.expanduser(os.fspath(root))) + if real == os.path.realpath(os.sep): + raise _sandbox_error( + what="A declared-location root is the filesystem root", + why="Confining declarations to '/' confines nothing.", + fix="Confine declarations to the contract's directory and workspace.", + ) + if real not in roots: + roots.append(real) + return tuple(roots) + + +def _confine(candidate: str, roots: Tuple[str, ...], declared: str) -> None: + """Refuse ``candidate`` (from the declared ``declared``) outside every root.""" + real = os.path.realpath(candidate) + if any(_is_within(real, root) for root in roots): + return + where = ", ".join(roots) if roots else "none" + raise _sandbox_error( + what=( + f"The contract declares {declared!r} ({real}), outside the directories it may " + f"read and write ({where}). The operator can allow a directory with " + f"{OPERATOR_DIRS_ENV}." + ), + why=( + "A declared input, output or source is granted to the contract's SQL. " + "Granted anywhere, a contract could read any file on this host (credentials, " + "/proc, another tenant's data) just by declaring it." + ), + fix=( + "Move the file under the contract's directory or its workspace, or, as the " + f"operator, list its directory in {OPERATOR_DIRS_ENV} " + f"(separated by {os.pathsep!r})." + ), + ) + + +def confine_declared( + location: PathLike, + *, + within: Iterable[Optional[PathLike]], + base: Optional[PathLike] = None, +) -> str: + """The ``realpath`` of a declared local ``location``, refused outside ``within``. + + For a location the engine opens itself rather than grants to the SQL, such + as a SQLite source it ``ATTACH``es before the lock: the sqlite scanner opens + files through its own client library, so ``allowed_directories`` never + bounds it, and this check is the only one. Confined as + :meth:`DuckDBAllowlist.with_declared` confines (``within`` plus + ``FLUID_DUCKDB_ALLOWED_DIRS``); a relative ``location`` is resolved against + ``base``, else the working directory. Open the returned path, the one that + was checked. + """ + raw = os.fspath(location) + absolute = _absolute(raw, base) + _confine(absolute, _confinement_roots(within), raw) + return os.path.realpath(absolute) + + +def unaliased_dir(path: PathLike, *, within: Iterable[Optional[PathLike]] = ()) -> Optional[str]: + """``path`` as an absolute directory to grant, or ``None`` if it is an alias. + + For a directory the engine grants by convention, not because the contract + declared it: the local provider's ``./runtime``. It is resolved in the + working directory, which is usually the contract's own, so the contract's + repository can ship ``runtime`` as a symlink (``runtime -> ../../..``) and, + because DuckDB realpaths allowlist entries, have ``$HOME`` granted. Returned + only when its realpath is where its name says it is (the realpath of its + parent, joined with its name) or inside one of ``within``; otherwise + ``None``, and the caller grants nothing for it. + """ + lexical = os.path.normpath(os.path.abspath(os.fspath(path))) + real = os.path.realpath(lexical) + parent, name = os.path.split(lexical) + if real == os.path.join(os.path.realpath(parent), name): + return lexical + for root in within: + if root is None or str(root) == "": + continue + real_root = os.path.realpath(os.path.expanduser(os.fspath(root))) + if real_root != os.path.realpath(os.sep) and _is_within(real, real_root): + return lexical + return None + + +def _new( + existing: Tuple[str, ...], + candidates: Iterable[Any], + normalize: Callable[[Any], str] = str, +) -> Tuple[str, ...]: + """``candidates``, each ``normalize``-d, not already in ``existing``; in order, each once. + + Linear: a glob or a long declaration list must not make the allowlist + quadratic. A candidate already granted verbatim is not normalised again. + """ + seen = set(existing) + out: list = [] + for raw in candidates: + if isinstance(raw, str) and raw in seen: + continue + item = normalize(raw) + if item not in seen: + seen.add(item) + out.append(item) + return tuple(out) + + +def _local(value: PathLike) -> str: + raw = os.fspath(value) + if raw.startswith("~"): + # duckdb/duckdb#26064: '~' is expanded inconsistently by DuckDB itself. + raise _sandbox_error( + what=f"DuckDB allowlist entry {raw!r} starts with '~'", + why=( + "DuckDB expands '~' differently when it checks a path and when it opens " + "it (duckdb/duckdb#26064), so a '~' entry does not bound what is read." + ), + fix="Pass an absolute path.", + ) + if is_remote_location(raw) or "://" in raw: + raise _sandbox_error( + what=f"DuckDB allowlist entry {raw!r} is a URL, not a local path", + why="Local directories and remote prefixes are granted separately.", + fix="Grant it with DuckDBAllowlist.with_remote (or with_locations).", + ) + resolved = os.path.normpath(os.path.abspath(raw)) + # DuckDB realpaths each entry, so a symlink to '/' is the root too. + if resolved == os.path.abspath(os.sep) or os.path.realpath(resolved) == os.path.realpath( + os.sep + ): + raise _sandbox_error( + what="DuckDB allowlist entry is the filesystem root", + why="Allowing '/' allows every file, which is no sandbox at all.", + fix="Grant the directory the run actually reads or writes.", + ) + return resolved + + +def _remote(prefix: str) -> str: + match = _REMOTE_RE.match(prefix) + if not match or match.group("scheme").lower() not in _REMOTE_SCHEMES: + raise _sandbox_error( + what=f"DuckDB remote prefix {prefix!r} is not a URL DuckDB reads", + why=f"Remote prefixes must use one of: {', '.join(sorted(_REMOTE_SCHEMES))}.", + fix="Grant a local path with with_dirs / with_paths instead.", + ) + path = match.group("path") or "/" + if ".." in path.split("/"): + raise _sandbox_error( + what=f"DuckDB remote prefix {prefix!r} contains '..'", + why="DuckDB does not resolve '..' in a URL, so the prefix would not bound it.", + fix="Name the prefix without '..'.", + ) + if not prefix.endswith("/"): + prefix += "/" + return prefix + + +def duckdb_version(module: Any) -> Tuple[int, int, int]: + """``module.__version__`` as ``(major, minor, patch)``; raises when unreadable.""" + raw = getattr(module, "__version__", None) + match = re.match(r"^(\d+)\.(\d+)\.(\d+)", raw) if isinstance(raw, str) else None + if not match: + raise _sandbox_error( + what="Cannot tell which DuckDB is installed", + why=f"duckdb.__version__ is {raw!r}; the sandbox needs a known version.", + fix="Reinstall DuckDB: pip install 'duckdb>=1.5.0'.", + ) + return int(match.group(1)), int(match.group(2)), int(match.group(3)) + + +def _require_supported(module: Any) -> None: + version = duckdb_version(module) + if version < MIN_DUCKDB_VERSION: + want = ".".join(str(n) for n in MIN_DUCKDB_VERSION) + have = ".".join(str(n) for n in version) + raise _sandbox_error( + what=f"DuckDB {have} is too old to sandbox contract SQL", + why=( + f"Before {want}, DuckDB's allowed_directories could be escaped with " + "'/./../' or a symlink, so contract SQL could read any file on " + "this host." + ), + fix=f"pip install 'duckdb>={want}'", + ) + + +def _sql_list(values: Iterable[str]) -> str: + return "[" + ", ".join(quote_ansi_string_literal(v) for v in values) + "]" + + +def _with_sep(directory: str) -> str: + # A trailing separator makes the entry a directory, never a name prefix + # (``/data`` must not admit ``/data-other``). + return directory if directory.endswith(("/", os.sep)) else directory + os.sep + + +def secure_duckdb_connect( + database: PathLike = ":memory:", + *, + allow: DuckDBAllowlist, + extensions: Iterable[str] = (), + read_only: bool = False, + config: Optional[Mapping[str, Any]] = None, + before_lock: Optional[Callable[[Any], None]] = None, +) -> "duckdb.DuckDBPyConnection": + """Open DuckDB with file and network access confined to ``allow``. + + ``allow`` is required: each caller states what it legitimately reads and + writes. A file-backed ``database`` needs no grant: it is opened before the + lock, and its own WAL and checkpoint writes are not checked against the + allowlist (``test_file_database_writes_and_checkpoints_under_the_lock``), so + granting its directory would only widen what the SQL can read (for the + local provider's ``persist`` mode, all of ``~/.fluid``). ``extensions`` are installed if missing and loaded before the lock; nothing + can be loaded after it. ``before_lock(con)`` runs once the extensions are + in, for set-up that needs access the SQL must not have, such as attaching a + declared source database or creating an object-store secret. ``config`` + takes scalar start-up options (``threads``, ``TimeZone``), which cannot be + changed once the configuration is locked. + + Raises :class:`DuckDBSandboxError` for an older DuckDB, a bad allowlist + entry, or a file ``database`` this process already has open (its one + shared instance is locked, so this connection could not be confined to + ``allow``); whatever ``duckdb.connect``, an extension or ``before_lock`` raises + is raised as it is, with the connection closed. + """ + import duckdb + + _require_supported(duckdb) + target = os.fspath(database) + try: + con = duckdb.connect(target, read_only=read_only, config=dict(config or {})) + except duckdb.Error as exc: + if "same database file with a different configuration" in str(exc): + raise _already_open(target) from exc + raise + try: + try: + con.execute("SET allow_persistent_secrets = false") + except duckdb.Error as exc: + # The first statement on a fresh connection: locked already means + # this joined another connection's locked instance of the file. + if "configuration has been locked" in str(exc): + raise _already_open(target) from exc + raise + con.execute("SET allow_community_extensions = false") + for ext in extensions: + name = validate_ident(str(ext)) + try: + con.execute(f"LOAD {name}") + except duckdb.Error: + con.execute(f"INSTALL {name}") + con.execute(f"LOAD {name}") + if before_lock is not None: + before_lock(con) + con.execute("SET autoinstall_known_extensions = false") + con.execute("SET autoload_known_extensions = false") + home = os.environ.get("HOME") + if home: + # duckdb/duckdb#26064: the allowlist check expands '~' with $HOME and + # the open with this setting. Equal and locked, they cannot disagree. + con.execute(f"SET home_directory = {quote_ansi_string_literal(home)}") + directories = [_with_sep(d) for d in allow.dirs] + list(allow.remote_prefixes) + con.execute(f"SET allowed_directories = {_sql_list(directories)}") + con.execute(f"SET allowed_paths = {_sql_list(allow.paths)}") + con.execute("SET enable_external_access = false") + con.execute("SET lock_configuration = true") + except BaseException: + con.close() + raise + return con + + +def _already_open(target: str) -> DuckDBSandboxError: + return _sandbox_error( + what=f"DuckDB database {target!r} is already open in this process", + why=( + "DuckDB shares one instance per database file within a process, and the " + "sandbox locks that instance's configuration, so a second connection to an " + "open file cannot be sandboxed with its own allowlist." + ), + fix=( + "Close the other connection to this file first, or give this run its own " + "database file." + ), + ) + + +def is_unloaded_function(exc: BaseException) -> bool: + """Whether ``exc`` is a function whose extension autoloading would have loaded. + + DuckDB's catalog error: ``Table Function with name "sqlite_scan" is not in + the catalog, but it exists in the sqlite_scanner extension``. + """ + text = str(exc) + return "but it exists in the" in text and "extension" in text + + +def is_sandbox_refusal(exc: BaseException) -> bool: + """Whether ``exc`` is DuckDB refusing a path or setting the sandbox denies.""" + try: + import duckdb + except ImportError: # pragma: no cover - only reachable without duckdb + return False + if isinstance(exc, duckdb.PermissionException): + return True + # Under the lock no SQL can SET a setting back, nor pull in an extension + # (autoloading is off): DuckDB's own advice to do either cannot be taken. + # ``read_csv('https://...')`` says "requires the extension"; a function such + # as ``sqlite_scan`` "is not in the catalog, but it exists in the + # sqlite_scanner extension". + text = str(exc) + return isinstance(exc, duckdb.Error) and ( + "configuration has been locked" in text + or "requires the extension" in text + or is_unloaded_function(exc) + ) + + +def sandbox_refusal_hint(allow: DuckDBAllowlist, refusal: Optional[BaseException] = None) -> str: + """One sentence naming what the refused SQL could have read instead. + + With the ``refusal`` itself, a function from an extension that was not + loaded is named as such rather than as a path to declare. + """ + if refusal is not None and is_unloaded_function(refusal): + return ( + "DuckDB refused it: contract SQL runs with extension autoloading off and " + "cannot INSTALL or LOAD one, so functions from extensions the engine does not " + "load (sqlite_scan, read_xlsx, ST_Read, delta_scan, iceberg_scan) are not " + "available. Read the data as CSV, Parquet or JSON, or land it with an " + "acquisition build first." + ) + granted = [*allow.dirs, *allow.paths, *allow.remote_prefixes] + where = ", ".join(granted) if granted else "nothing (in-memory only)" + return ( + "DuckDB refused it: contract SQL may only read and write the locations the " + f"contract declares and its own directory ({where}). Declare the file as an " + "input under the contract's directory or workspace, or move it there " + f"(the operator can allow another directory with {OPERATOR_DIRS_ENV})." + ) + + +__all__ = [ + "MIN_DUCKDB_VERSION", + "OPERATOR_DIRS_ENV", + "DuckDBAllowlist", + "DuckDBSandboxError", + "confine_declared", + "duckdb_version", + "is_remote_location", + "is_sandbox_refusal", + "is_unloaded_function", + "operator_allowed_dirs", + "sandbox_refusal_hint", + "secure_duckdb_connect", + "unaliased_dir", +] diff --git a/fluid_build/providers/local/ducksql.py b/fluid_build/providers/local/ducksql.py index 3e17587a..f48d9756 100644 --- a/fluid_build/providers/local/ducksql.py +++ b/fluid_build/providers/local/ducksql.py @@ -17,6 +17,7 @@ from fluid_build.util.contract import get_expose_id, get_expose_kind, get_expose_location +from .._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect, unaliased_dir from ..base import ApplyResult, PlanAction try: @@ -48,7 +49,13 @@ def apply_sql(actions, dry_run=False): results = [] if duckdb is None: return [ApplyResult(False, "duckdb not installed. pip install duckdb", error="missing_dep")] - con = duckdb.connect(":memory:") + # The SQL is the expose's own (contract input) and the CSV is written by + # pandas, so DuckDB reaches only the local workspace, ``./runtime``, and + # only when it is no symlink out of the working directory (a contract's + # repository could ship ``runtime -> ../../..`` to have $HOME granted). + con = secure_duckdb_connect( + ":memory:", allow=DuckDBAllowlist.none().with_dirs(unaliased_dir("runtime")) + ) for a in actions: if a.resource_type != "sql.to_csv": continue diff --git a/fluid_build/providers/local/local.py b/fluid_build/providers/local/local.py index 74e99aac..f1e782da 100644 --- a/fluid_build/providers/local/local.py +++ b/fluid_build/providers/local/local.py @@ -39,6 +39,13 @@ from fluid_build.observability.secret_redactor import redact_secret_text, redact_value from fluid_build.providers._duckdb_read import build_register_view_sql +from fluid_build.providers._duckdb_sandbox import ( + DuckDBAllowlist, + is_sandbox_refusal, + sandbox_refusal_hint, + secure_duckdb_connect, + unaliased_dir, +) from fluid_build.providers._sql_safety import quote_ansi_string_literal, validate_ident from fluid_build.providers.base import ApplyResult, BaseProvider, ProviderMetadata @@ -218,6 +225,111 @@ def _anchored(self, contract: Dict[str, Any]) -> Dict[str, Any]: return anchor_binding_paths(contract, self.anchor_dir) + # ------------------------- Sandboxed DuckDB ------------------------- # + + @staticmethod + def _declared_io(action: Dict[str, Any]) -> List[Any]: + """Every location ``action`` declares it reads or writes through DuckDB.""" + op = (action.get("op") or action.get("type") or "").lower().strip() + found: List[Any] = [] + if op in {"sql", "query", "execute_sql"}: + for key in ("inputs", "tables", "outputs", "out"): + value = action.get(key) or [] + found.extend(value if isinstance(value, list) else [value]) + elif op == "load_data": + found.append(action.get("path")) + elif op in {"copy", "materialize"}: + found.append(action.get("dst") or action.get("out") or action.get("path")) + return [spec for spec in found if spec] + + def _declare_run_io(self, actions: Iterable[Dict[str, Any]]) -> None: + """Record what every action of this run declares, before any SQL runs. + + One apply shares one session database, so a view one action registers + over a file is read again by a later action's SQL: every connection of + the run gets the run's whole declared allowlist, not just its own. + """ + specs: List[Any] = [] + for action in actions: + flat = dict(action) + payload = flat.pop("payload", None) + if isinstance(payload, dict): + for key, value in payload.items(): + flat.setdefault(key, value) + specs.extend(self._declared_io(flat)) + self._run_io = specs + + def _allowlist(self, specs: Iterable[Any]) -> DuckDBAllowlist: + """What this provider's SQL may touch, and nothing else. + + The source contract's directory (``anchor_dir``) and the FLUID + workspace it sits in, the run's session scratch directory, + ``./runtime`` (where previews and default outputs land), and each + location the actions declare: an input file, an output file, an + ``s3://`` prefix. A contract's SQL that names any other path, + ``/etc/passwd`` or ``~/.aws/credentials``, is refused by DuckDB. + ``./runtime`` is left out when it is a symlink that leads outside the + contract's directory and workspace (:func:`unaliased_dir`). + + A declared location is granted only inside those directories (plus + the upstream roots in ``FLUID_UPSTREAM_CONTRACTS`` and the operator's + ``FLUID_DUCKDB_ALLOWED_DIRS``): the contract's author writes the + declaration, so it must not be a way to grant the host + (``DuckDBAllowlist.with_declared``). A relative declared path is + resolved where DuckDB opens it, the working directory, and is confined + all the same. + """ + from fluid_build.util.upstream_discovery import collect_search_roots + from fluid_build.util.workspace_root import find_workspace_root + + session = getattr(self, "_session_db", None) + scratch = Path(session).parent if session else None + # Without a contract directory (a bare ``apply`` of actions), the + # working directory stands in for it, as it does for relative paths. + anchor = self.anchor_dir if self.anchor_dir is not None else Path.cwd() + # A contract inside a FLUID workspace (``fluid.workspace.yaml``) may + # read its sibling products' files by path, as consumes[] does. + workspace = find_workspace_root(self.anchor_dir) if self.anchor_dir is not None else None + # ``./runtime`` is in the working directory, usually the contract's + # own, so the contract's repository can ship it as a symlink + # (``runtime -> ../../..``). Granted only where its name says it is, or + # inside the contract's directory or workspace; otherwise not at all. + runtime = unaliased_dir("runtime", within=[anchor, workspace]) + if runtime is None: + self._log_warn( + "local_runtime_not_granted", + {"runtime": str(Path("runtime").absolute()), "reason": "symlink_leads_out"}, + ) + allow = DuckDBAllowlist.none().with_dirs(self.anchor_dir, workspace, runtime, scratch) + within = [anchor, runtime, scratch, *collect_search_roots(workspace)] + for spec in specs: + raw = spec.get("path") if isinstance(spec, dict) else spec + if not raw: + continue + raw = str(raw) + # Only s3:// is remote here (``_register_mapping_input``): any other + # string is a local path, as ``Path`` reads it. + allow = allow.with_declared(raw if _is_s3_uri(raw) else str(Path(raw)), within=within) + return allow + + def _connect(self, specs: Iterable[Any] = (), *, config: Optional[Dict[str, Any]] = None): + """The provider's one way to DuckDB: sandboxed to :meth:`_allowlist`. + + ``specs`` are the calling action's own declared locations, added to the + run's (:meth:`_declare_run_io`). The S3 buckets among them get their + extensions and credential secret before the configuration is locked. + """ + _Duck.get() + all_specs = [*getattr(self, "_run_io", []), *specs] + allow = self._allowlist(all_specs) + self._last_allow = allow + return secure_duckdb_connect( + self._get_db_path(), + allow=allow, + config=config, + before_lock=lambda con: self._attach_object_stores(con, all_specs), + ) + def _get_db_path(self) -> str: """Get database path - persistent, session-scoped, or in-memory.""" if self.persist: @@ -378,6 +490,7 @@ def apply( ] # ---- Execute ---- + self._declare_run_io(a for a in norm_actions if isinstance(a, dict)) results: List[Dict[str, Any]] = [] error_count = 0 for idx, action in enumerate(norm_actions): @@ -642,8 +755,6 @@ def _run_load_data_action(self, idx: int, action: Dict[str, Any]) -> Dict[str, A Supports: CSV, TSV, Parquet, JSON, JSONL Handles: Globs, schemas, custom options, retries on transient errors """ - duckdb = _Duck.get() - path = action.get("path") table_name = action.get("table_name") or action.get("resource_id") fmt = action.get("format", "csv") @@ -660,9 +771,12 @@ def _run_load_data_action(self, idx: int, action: Dict[str, Any]) -> Dict[str, A if not _has_glob(path_obj) and not path_obj.exists(): raise FileNotFoundError(f"Input file not found: {path}") - db_path = self._get_db_path() - con = duckdb.connect(database=db_path) + con = self._connect([path]) + # Closed on every path: every action of the run shares one session + # database, and a connection left open (an exception's traceback keeps + # it alive) holds that file's locked instance, which refuses the next + # action's sandboxed connection. try: # Use retry logic for table registration (can fail with I/O errors) def _register_with_retry(): @@ -700,12 +814,13 @@ def _register_with_retry(): {"i": idx, "error": str(e), "path": str(path), "table": table_name}, ) raise + finally: + con.close() # ----------------------- SQL (DuckDB) execution --------------------- # def _run_sql_action(self, idx: int, action: Dict[str, Any]) -> Dict[str, Any]: """Execute SQL with retry logic, persistent DB, and enhanced logging.""" - duckdb = _Duck.get() start_time = time.time() sql = action.get("sql") or action.get("query") @@ -714,15 +829,31 @@ def _run_sql_action(self, idx: int, action: Dict[str, Any]) -> Dict[str, Any]: _mkdir("runtime/out") - db_path = self._get_db_path() - con = duckdb.connect(database=db_path) - con.execute("PRAGMA threads=4;") - inputs = action.get("inputs") or action.get("tables") or [] outputs = action.get("outputs") or action.get("out") or [] if isinstance(outputs, (str, Path)): outputs = [outputs] - self._attach_object_stores(con, [*inputs, *outputs]) + # Threads set at connect: the sandbox locks the configuration, so a + # ``PRAGMA threads`` afterwards is refused. + con = self._connect([*inputs, *outputs], config={"threads": 4}) + # Closed on every path, a failure included (see _run_load_data_action): + # otherwise one failing SQL action fails every later action of the run. + try: + return self._run_sql_on(con, idx, action, sql, inputs, outputs, start_time) + finally: + con.close() + + def _run_sql_on( + self, + con: Any, + idx: int, + action: Dict[str, Any], + sql: str, + inputs: List[Any], + outputs: List[Any], + start_time: float, + ) -> Dict[str, Any]: + """The body of :meth:`_run_sql_action`, on a connection it closes.""" reg_info = self._register_inputs(con, inputs) # Log with redacted SQL (in case it contains sensitive data) @@ -754,7 +885,13 @@ def _execute_sql(): "duration_ms": duration_ms(start_time), }, ) - raise + if not is_sandbox_refusal(e): + raise + # A path outside the sandbox: say what the SQL may read instead. + # Raised from a helper, so this frame keeps no reference to the new + # error (an error -> traceback -> frame -> error cycle would keep + # the connection alive until the cyclic GC runs). + raise self._sandbox_refusal(e) from e # If an output_table is specified, persist the result as a DuckDB table # so downstream materialize/copy steps can reference it. @@ -799,6 +936,10 @@ def _write_preview(): ) return {"op": "sql", "written": written, "rows": rowcount, "inputs": reg_info} + def _sandbox_refusal(self, refused: BaseException) -> PermissionError: + """``refused`` as a PermissionError that names what the SQL may read instead.""" + return PermissionError(f"{refused} {sandbox_refusal_hint(self._last_allow, refused)}") + def _register_inputs(self, con: Any, inputs: Iterable[Any]) -> List[Dict[str, Any]]: info: List[Dict[str, Any]] = [] for item in inputs or []: @@ -952,12 +1093,15 @@ def _run_copy_action(self, idx: int, action: Dict[str, Any]) -> Dict[str, Any]: fmt = (action.get("format") or _ext(dst) or "csv").lower() if source_table: try: - duckdb = _Duck.get() - db_path = self._get_db_path() - con = duckdb.connect(database=db_path) - rel = con.sql(f"SELECT * FROM {validate_ident(source_table)}") - self._write_relation(rel, dst, fmt) - rowcount = rel.count("*").fetchone()[0] if hasattr(rel, "count") else -1 + con = self._connect([dst]) + try: + rel = con.sql(f"SELECT * FROM {validate_ident(source_table)}") + self._write_relation(rel, dst, fmt) + rowcount = rel.count("*").fetchone()[0] if hasattr(rel, "count") else -1 + del rel + finally: + # Closed on every path (see _run_load_data_action). + con.close() self._log_info( "local_materialize_done", {"i": idx, "dst": str(dst), "source": source_table, "rows": rowcount}, diff --git a/fluid_build/providers/local/mocks.py b/fluid_build/providers/local/mocks.py index 6e1eafc5..614747d1 100644 --- a/fluid_build/providers/local/mocks.py +++ b/fluid_build/providers/local/mocks.py @@ -198,12 +198,19 @@ def __init__(self, logger: Optional[logging.Logger] = None, db_path: Optional[st def connect(self) -> Any: """Connect to mock Snowflake (DuckDB).""" try: - import duckdb - - self._con = duckdb.connect(self.db_path) - - # Set up Snowflake-like configuration - self._con.execute("SET TimeZone='UTC'") + import duckdb # noqa: F401 - the ImportError below is the install hint + + from fluid_build.providers._duckdb_sandbox import ( + DuckDBAllowlist, + secure_duckdb_connect, + ) + + # The mock runs whatever SQL it is handed, so it gets no file or + # network access beyond its own database file. TimeZone is set at + # connect: the sandbox locks the configuration afterwards. + self._con = secure_duckdb_connect( + self.db_path, allow=DuckDBAllowlist.none(), config={"TimeZone": "UTC"} + ) # Create standard Snowflake schemas if they don't exist self._con.execute("CREATE SCHEMA IF NOT EXISTS INFORMATION_SCHEMA") diff --git a/fluid_build/providers/local/util/retry.py b/fluid_build/providers/local/util/retry.py index 19ab3324..16ad591f 100644 --- a/fluid_build/providers/local/util/retry.py +++ b/fluid_build/providers/local/util/retry.py @@ -67,6 +67,13 @@ def is_retryable_error(exception: Exception) -> bool: if isinstance(exception, NonRetryableError): return False + # A DuckDB sandbox refusal is permanent: retrying only re-runs SQL that was + # refused (and its message, a path, can match a pattern below by accident). + from fluid_build.providers._duckdb_sandbox import is_sandbox_refusal + + if is_sandbox_refusal(exception): + return False + exception_str = str(exception).lower() exception_type = type(exception).__name__.lower() diff --git a/fluid_build/providers/local_validation.py b/fluid_build/providers/local_validation.py index dd2f7452..ce176596 100644 --- a/fluid_build/providers/local_validation.py +++ b/fluid_build/providers/local_validation.py @@ -26,6 +26,7 @@ from pathlib import Path from typing import Any, Dict, List, Optional +from fluid_build.providers._duckdb_sandbox import DuckDBAllowlist, secure_duckdb_connect from fluid_build.providers._sql_safety import quote_ansi_string_literal, validate_ident from fluid_build.providers.quality_engine import ( execute_quality_checks, @@ -111,6 +112,18 @@ def _get_duckdb(self): ) return self._duckdb + def _connect( + self, database: str = ":memory:", *, allow: DuckDBAllowlist, read_only: bool = False + ): + """A sandboxed DuckDB confined to ``allow`` (see ``_duckdb_sandbox``). + + Every connection this provider opens goes through here: the quality + rules it runs are contract input, so each reaches the one file it + checks and nothing else. + """ + self._get_duckdb() + return secure_duckdb_connect(database, allow=allow, read_only=read_only) + @property def provider_name(self) -> str: return "local" @@ -118,8 +131,7 @@ def provider_name(self) -> str: def validate_connection(self) -> bool: """Local provider is always reachable.""" try: - duckdb = self._get_duckdb() - conn = duckdb.connect(":memory:") + conn = self._connect(allow=DuckDBAllowlist.none()) conn.execute("SELECT 1") conn.close() return True @@ -137,13 +149,13 @@ def get_resource_schema(self, resource_spec: Dict[str, Any]) -> Optional[Resourc if not path.exists(): return None - duckdb = self._get_duckdb() - conn = duckdb.connect(":memory:") + ext = path.suffix.lower() + abs_path = str(path.resolve()) + # The one file this introspects; a .duckdb file is opened read-only + # below, and the sandbox grants a read-only database file itself. + conn = self._connect(allow=DuckDBAllowlist.none().with_paths(abs_path)) try: - ext = path.suffix.lower() - abs_path = str(path.resolve()) - # DuckDB native database file — connect and query the table directly if ext in (".duckdb", ".db"): binding = resource_spec.get("binding", {}) @@ -157,7 +169,7 @@ def get_resource_schema(self, resource_spec: Dict[str, Any]) -> Optional[Resourc # malicious schema/table never reaches a query (fail-closed). table_ref = _build_duckdb_table_ref(schema_name, table_name) conn.close() - conn = duckdb.connect(abs_path, read_only=True) + conn = self._connect(abs_path, allow=DuckDBAllowlist.none(), read_only=True) describe_sql = f"DESCRIBE SELECT * FROM {table_ref}" rows = conn.execute(describe_sql).fetchall() fields: List[FieldSchema] = [] @@ -384,7 +396,6 @@ def run_quality_checks( ] ext = path.suffix.lower() abs_path = str(path.resolve()) - duckdb = self._get_duckdb() # DuckDB native database: connect directly and query the bound table if ext in (".duckdb", ".db"): @@ -420,7 +431,7 @@ def run_quality_checks( path="exposes[].binding.location", ) ] - conn = duckdb.connect(abs_path, read_only=True) + conn = self._connect(abs_path, allow=DuckDBAllowlist.none(), read_only=True) try: def _exec(sql): @@ -452,7 +463,8 @@ def _exec(sql): ) ] table_ref = "{fn}({p})".format(fn=read_fn, p=quote_ansi_string_literal(abs_path)) - conn = duckdb.connect(":memory:") + # The rules are contract input: they reach the file they check, no other. + conn = self._connect(allow=DuckDBAllowlist.none().with_paths(abs_path)) try: def _exec(sql): diff --git a/pyproject.toml b/pyproject.toml index 22854b4d..3e454f8d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -242,7 +242,11 @@ odcs-strict = [ # Local provider (DuckDB) local = [ - "duckdb>=1.0", + # 1.5.0, not 1.0: every connection runs contract SQL inside DuckDB's sandbox + # (fluid_build/providers/_duckdb_sandbox.py). allowed_directories arrived in + # 1.2; up to 1.4.3 '/./../' escaped it, and through 1.4.x a symlink in + # an allowed directory did. 1.5.0 resolves both; the helper refuses older. + "duckdb>=1.5.0", "pandas>=2.2", # Named here, not left to arrive through pandas: masking at landing # registers DuckDB Python UDFs, DuckDB's create_function needs numpy diff --git a/tests/build_runners/test_consumes_resolution.py b/tests/build_runners/test_consumes_resolution.py index c6437cca..91e9ffc9 100644 --- a/tests/build_runners/test_consumes_resolution.py +++ b/tests/build_runners/test_consumes_resolution.py @@ -610,7 +610,8 @@ def test_an_explicit_input_wins_over_the_resolved_upstream(tmp_path, printed): root = _workspace(tmp_path / "ws") bronze = _bronze(root) _parquet(bronze.parent / "out" / "customer_subscriptions.parquet") - hand = _parquet(tmp_path / "hand" / "subs.parquet", rows=[("p9", "active")]) + # In the workspace: a declared input outside it is refused (DuckDB sandbox). + hand = _parquet(root / "hand" / "subs.parquet", rows=[("p9", "active")]) _dump( root / "contracts/customers/contract.fluid.yaml", { @@ -1139,9 +1140,13 @@ def test_the_high_value_churn_example_builds_as_it_did_before_consumes_resolved( def test_sql_that_reads_its_upstream_by_path_keeps_building(tmp_path, printed, caplog): - """The old workaround: consumes[] for lineage, the file named in the SQL.""" + """The old workaround: consumes[] for lineage, the file named in the SQL. + + The file is inside the workspace: the DuckDB sandbox grants the workspace a + contract sits in, and refuses a path outside it (the next test). + """ root = _workspace(tmp_path / "ws") # no upstream contract anywhere - data = _parquet(tmp_path / "shared" / "cs.parquet") + data = _parquet(root / "shared" / "cs.parquet") build = _silver_build( sql=( "SELECT product_id, status, COUNT(*) AS subscription_count " @@ -1157,6 +1162,17 @@ def test_sql_that_reads_its_upstream_by_path_keeps_building(tmp_path, printed, c assert any("lineage only" in line for line in printed), printed +def test_sql_that_reads_a_file_outside_the_workspace_by_path_is_refused(tmp_path, printed): + """Contract SQL reaches the workspace and what the contract declares, no more.""" + root = _workspace(tmp_path / "ws") + data = _parquet(tmp_path / "elsewhere" / "cs.parquet") + build = _silver_build(sql=f"SELECT * FROM read_parquet('{data}')") + contract = _silver(root, build=build) + assert _run(root, contract) == 1 + assert any("Permission Error" in line for line in printed), printed + assert not (_silver_dir(root) / "out" / "subscription_status_summary.parquet").exists() + + def test_only_the_entries_the_sql_reads_must_resolve(tmp_path, printed): """One entry read (and resolvable), one declared for lineage (and not in the workspace).""" root = _workspace(tmp_path / "ws") @@ -1425,7 +1441,7 @@ def test_a_local_upstream_is_read_where_and_how_the_local_writer_wrote_it( tmp_path, printed, fmt, path ): root = _workspace(tmp_path / "ws") - seed = tmp_path / "seed.csv" + seed = root / "seed.csv" # in the workspace, where a declared input may be seed.write_text("product_id,status\np1,active\np1,active\np2,ended\n") bronze = { "fluidVersion": "0.7.5", @@ -1552,7 +1568,7 @@ def test_a_failed_action_never_prints_or_records_a_resolved_secret( secret = "sk-live-SUPERSECRET-0123456789" # pragma: allowlist secret monkeypatch.setenv("PARTNER_API_TOKEN", secret) root = _workspace(tmp_path / "ws") - data = _parquet(tmp_path / "hand" / "subs.parquet") + data = _parquet(root / "hand" / "subs.parquet") build = _silver_build( sql="SELECT * FROM subscriptions WHERE 1 = CAST('{{ env.PARTNER_API_TOKEN }}' AS INTEGER)", inputs=[{"name": "subscriptions", "path": str(data)}], diff --git a/tests/perf/test_acquisition_perf_budgets.py b/tests/perf/test_acquisition_perf_budgets.py index afa56e82..71a91cbd 100644 --- a/tests/perf/test_acquisition_perf_budgets.py +++ b/tests/perf/test_acquisition_perf_budgets.py @@ -132,13 +132,16 @@ def _write_batch(): class TestDuckdbRunnerPerf: - def test_duckdb_csv_to_parquet_100_rows_under_5s(self, tmp_path: Path): + def test_duckdb_csv_to_parquet_100_rows_under_5s(self, tmp_path: Path, monkeypatch): """The classic 'first sync' UX budget: CSV → Parquet for 100 rows in under 5 s. This is the time-to-first-value floor declared in the design doc. Includes import + extension load + COPY. """ from fluid_build.build_runners.duckdb.runner import execute_duckdb_build + # The source and landing sit beside, not under, the run's workdir: the + # operator allows that directory, as the DuckDB sandbox requires. + monkeypatch.setenv("FLUID_DUCKDB_ALLOWED_DIRS", str(tmp_path)) in_dir = tmp_path / "in" in_dir.mkdir() csv = in_dir / "perf.csv" diff --git a/tests/providers/test_duckdb_sandbox.py b/tests/providers/test_duckdb_sandbox.py new file mode 100644 index 00000000..6ff2dd0d --- /dev/null +++ b/tests/providers/test_duckdb_sandbox.py @@ -0,0 +1,1336 @@ +# Copyright 2024-2026 Agentics Transformation Ltd +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Contract SQL cannot read the host: every DuckDB the engine opens is sandboxed. + +A contract's ``builds[].properties.sql`` runs on DuckDB with the privileges of +the process. Before ``providers/_duckdb_sandbox.py`` it could read any file the +process could (``read_csv('/etc/passwd')``), fetch any URL, ATTACH any +database and COPY anywhere. These tests run the attacks through the real +embedded-SQL build path (``_execute_embedded_sql_build``), pin the helper's +own guarantees, and fail if any new ``duckdb.connect`` bypasses the helper. +""" + +from __future__ import annotations + +import ast +import os +import shutil +import subprocess +import sys +from pathlib import Path +from typing import Any, Dict, Iterator, List +from unittest.mock import patch + +import pytest + +duckdb = pytest.importorskip("duckdb") + +from fluid_build.providers import _duckdb_sandbox as sandbox # noqa: E402 +from fluid_build.providers._duckdb_sandbox import ( # noqa: E402 + MIN_DUCKDB_VERSION, + DuckDBAllowlist, + DuckDBSandboxError, + secure_duckdb_connect, +) + +SECRET = "TOP-SECRET-VALUE-0451" # pragma: allowlist secret +REPO_ROOT = Path(__file__).resolve().parents[2] +PACKAGE = REPO_ROOT / "fluid_build" + + +# ── fixtures ───────────────────────────────────────────────────────────── + + +@pytest.fixture +def layout(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Dict[str, Path]: + """A contract directory with its own data, and a secret beside it. + + ``outside/`` is a sibling of the contract directory (so ``../`` reaches it) + and stands in for ``$HOME``: whatever the SQL would read from there is what + the sandbox has to refuse. + """ + return _make_layout(tmp_path, monkeypatch) + + +def _make_layout(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Dict[str, Path]: + contract_dir = tmp_path / "product" + data = contract_dir / "data" + data.mkdir(parents=True) + (data / "orders.csv").write_text("id,amount\n1,10\n2,20\n", encoding="utf-8") + outside = tmp_path / "outside" + outside.mkdir() + (outside / "secret.csv").write_text(f"token\n{SECRET}\n", encoding="utf-8") + other = duckdb.connect(str(outside / "other.duckdb")) + other.execute(f"CREATE TABLE creds AS SELECT '{SECRET}' AS token") + other.close() + (contract_dir / "link.csv").symlink_to(outside / "secret.csv") + monkeypatch.setenv("HOME", str(outside)) + monkeypatch.chdir(contract_dir) + monkeypatch.delenv("FLUID_UPSTREAM_CONTRACTS", raising=False) + return {"contract_dir": contract_dir, "outside": outside, "tmp": tmp_path} + + +@pytest.fixture +def printed() -> Iterator[List[str]]: + lines: List[str] = [] + with patch( + "fluid_build.build_runners.base.cprint", + side_effect=lambda *a, **_k: lines.append(" ".join(str(x) for x in a)), + ): + yield lines + + +def _contract(sql: str, out: Path) -> Dict[str, Any]: + build = {"id": "attack", "engine": "sql", "properties": {"sql": sql}} + return { + "id": "silver.sandbox_probe", + "builds": [build], + "consumes": [], + "exposes": [ + { + "exposeId": "out", + "binding": {"platform": "local", "format": "csv", "location": {"path": str(out)}}, + } + ], + } + + +def _build(sql: str, layout: Dict[str, Path]) -> int: + from fluid_build.build_runners.base import _execute_embedded_sql_build + + out = layout["contract_dir"] / "out" / "result.csv" + contract = _contract(sql, out) + return _execute_embedded_sql_build(contract["builds"][0], contract, layout["contract_dir"]) + + +def _written_text(root: Path) -> str: + texts = [] + for path in root.rglob("*"): + if path.is_file() and not path.is_symlink() and path.suffix != ".duckdb": + texts.append(path.read_text(encoding="utf-8", errors="replace")) + return "\n".join(texts) + + +# ── attacks from a contract's embedded SQL ─────────────────────────────── + + +def _attacks(outside: Path, contract_dir: Path) -> Dict[str, str]: + return { + "absolute_host_file": "SELECT * FROM read_csv('/etc/passwd')", + "absolute_secret": f"SELECT * FROM read_csv('{outside}/secret.csv')", + "read_text": f"SELECT content FROM read_text('{outside}/secret.csv')", + "read_blob": f"SELECT content FROM read_blob('{outside}/secret.csv')", + "read_json": f"SELECT * FROM read_json_auto('{outside}/secret.csv')", + "read_parquet": f"SELECT * FROM read_parquet('{outside}/*.parquet')", + "glob": "SELECT * FROM glob('/etc/*')", + "dot_dot": f"SELECT * FROM read_csv('{contract_dir}/../outside/secret.csv')", + # The DuckDB <= 1.4.3 escape: './..' was not resolved before the check. + "dot_slash_dot_dot": f"SELECT * FROM read_csv('{contract_dir}/./../outside/secret.csv')", + "relative_dot_dot": "SELECT * FROM read_csv('../outside/secret.csv')", + "symlink_out": f"SELECT * FROM read_csv('{contract_dir}/link.csv')", + # $HOME is ``outside`` (fixture): '~' must not reach it. + "home_tilde": "SELECT * FROM read_csv('~/secret.csv')", + "http_url": "SELECT * FROM read_csv('http://127.0.0.1:9/secret.csv')", + "https_url": "SELECT * FROM read_csv('https://example.invalid/secret.csv')", + "attach_outside_db": ( + f"ATTACH '{outside}/other.duckdb' AS o (READ_ONLY); SELECT * FROM o.creds" + ), + "copy_to_outside": ( + f"COPY (SELECT 1 AS a) TO '{outside}/pwned.csv' (FORMAT csv); SELECT 1 AS a" + ), + "copy_from_outside": ( + f"CREATE OR REPLACE TABLE t (token VARCHAR); COPY t FROM '{outside}/secret.csv'; " + "SELECT * FROM t" + ), + "reset_external_access": ( + "SET enable_external_access = true; " f"SELECT * FROM read_csv('{outside}/secret.csv')" + ), + "widen_allowlist": ( + f"SET allowed_directories = ['{outside}']; " + f"SELECT * FROM read_csv('{outside}/secret.csv')" + ), + "unlock": "SET lock_configuration = false; SELECT 1 AS a", + "load_extension": "LOAD httpfs; SELECT 1 AS a", + "install_extension": "INSTALL spatial; SELECT 1 AS a", + } + + +_ATTACK_NAMES = sorted(_attacks(Path("/o"), Path("/c"))) + + +@pytest.mark.parametrize("name", _ATTACK_NAMES) +def test_contract_sql_cannot_reach_outside_the_sandbox(name, layout, printed): + sql = _attacks(layout["outside"], layout["contract_dir"])[name] + + rc = _build(sql, layout) + + assert rc == 1, f"{name}: the build succeeded; the sandbox let it through" + assert not (layout["outside"] / "pwned.csv").exists() + assert SECRET not in _written_text(layout["contract_dir"]) + assert SECRET not in "\n".join(printed) + # A clear refusal, not some unrelated failure that happens to exit 1. + errors = "\n".join(printed) + assert ( + "Permission Error" in errors + or "configuration has been locked" in errors + or "disabled through configuration" in errors + # A URL whose filesystem extension the build did not load: autoloading + # is off under the lock, so the SQL cannot load it either. + or "requires the extension" in errors + ), errors + if "requires the extension" in errors or "Permission Error" in errors: + assert "contract SQL may only read and write" in errors + + +def test_refusal_says_what_the_sql_may_read_instead(layout, printed): + rc = _build("SELECT * FROM read_csv('/etc/passwd')", layout) + + assert rc == 1 + errors = "\n".join(printed) + assert "Cannot access file" in errors + assert "contract SQL may only read and write" in errors + assert str(layout["contract_dir"]) in errors + + +# ── legitimate reads keep working ──────────────────────────────────────── + + +def test_contract_sql_reads_its_own_directory_and_writes_its_output(layout, printed): + rc = _build( + "SELECT id, amount * 2 AS doubled FROM read_csv('data/orders.csv') ORDER BY id", + layout, + ) + + assert rc == 0, printed + out = layout["contract_dir"] / "out" / "result.csv" + assert out.read_text(encoding="utf-8").splitlines() == ["id,doubled", "1,20", "2,40"] + + +def test_a_declared_input_outside_the_contract_directory_needs_the_operator( + layout, printed, monkeypatch +): + """Outside the contract's directory, a declared file is granted only where the + operator allows it (FLUID_DUCKDB_ALLOWED_DIRS), and then that file only.""" + shared = layout["tmp"] / "shared" + shared.mkdir() + (shared / "rates.csv").write_text("id,rate\n1,3\n", encoding="utf-8") + (shared / "private.csv").write_text(f"token\n{SECRET}\n", encoding="utf-8") + out = layout["contract_dir"] / "out" / "result.csv" + contract = _contract("SELECT * FROM rates", out) + contract["builds"][0]["properties"]["parameters"] = { + "inputs": [{"name": "rates", "path": str(shared / "rates.csv")}] + } + from fluid_build.build_runners.base import _execute_embedded_sql_build + + # The contract alone cannot grant a file outside its directory. + assert _execute_embedded_sql_build(contract["builds"][0], contract, layout["contract_dir"]) == 1 + assert "FLUID_DUCKDB_ALLOWED_DIRS" in "\n".join(printed), printed + assert not out.exists() + + monkeypatch.setenv("FLUID_DUCKDB_ALLOWED_DIRS", str(shared)) + assert _execute_embedded_sql_build(contract["builds"][0], contract, layout["contract_dir"]) == 0 + assert out.read_text(encoding="utf-8").splitlines() == ["id,rate", "1,3"] + + # Its neighbour was not declared, so it is not readable. + contract["builds"][0]["properties"]["sql"] = f"SELECT * FROM read_csv('{shared}/private.csv')" + assert _execute_embedded_sql_build(contract["builds"][0], contract, layout["contract_dir"]) == 1 + + +# ── contract-tests local actions ───────────────────────────────────────── + + +def test_contract_tests_action_reads_its_inputs_and_nothing_else(layout): + from types import SimpleNamespace + + from fluid_build.contract_tests import LocalProviderError, apply_action + + out = layout["contract_dir"] / "out" / "result.parquet" + action = { + "op": "add", + "resource_type": "sql", + "id": "orders", + "inputs": {"orders": {"path": str(layout["contract_dir"] / "data" / "orders.csv")}}, + "outputs": {"path": str(out), "format": "parquet"}, + "sql": "SELECT id, amount FROM orders", + } + apply_action(action, SimpleNamespace(dry_run=False)) + assert duckdb.connect().execute(f"SELECT count(*) FROM '{out}'").fetchone() == (2,) + + action["sql"] = f"SELECT * FROM read_csv('{layout['outside']}/secret.csv')" + with pytest.raises(LocalProviderError) as refused: + apply_action(action, SimpleNamespace(dry_run=False)) + assert "Permission Error" in str(refused.value) + assert "contract SQL may only read and write" in str(refused.value) + assert SECRET not in str(refused.value) + + +# ── the acquisition runner ─────────────────────────────────────────────── + + +def _acquisition_contract(source_glob: str, out_path: Path) -> Dict[str, Any]: + return { + "fluidVersion": "0.7.3", + "kind": "DataProduct", + "id": "bronze.sandbox_ingest", + "builds": [ + { + "id": "ingest", + "pattern": "acquisition", + "engine": "duckdb", + "properties": { + "source": { + "kind": "filesystem", + "connection": {"uri": source_glob}, + "mode": "full_refresh", + "reader": {"format": "csv", "options": {"header": True}}, + }, + "sink": {"format": "parquet"}, + }, + "outputs": ["orders_raw"], + } + ], + "exposes": [ + { + "exposeId": "orders_raw", + "kind": "table", + "binding": { + "platform": "local", + "format": "parquet", + "location": {"path": str(out_path)}, + }, + "contract": {"schema": [], "schemaPolicy": "discover_and_freeze"}, + } + ], + } + + +def test_acquisition_reads_its_declared_source_and_nothing_beside_it(layout, monkeypatch): + from fluid_build.build_runners._acquisition_common import build_acquisition_run_context + from fluid_build.build_runners.duckdb.runner import ( + _connect_for_destination, + execute_duckdb_build, + ) + + landing = layout["tmp"] / "landing" + landing.mkdir() + (landing / "orders.csv").write_text("id,amount\n1,10\n", encoding="utf-8") + out = layout["contract_dir"] / "out" / "orders.parquet" + contract = _acquisition_contract(str(landing / "*.csv"), out) + + # A landing directory outside the contract's is the operator's to allow. + with pytest.raises(DuckDBSandboxError, match="FLUID_DUCKDB_ALLOWED_DIRS"): + execute_duckdb_build(contract["builds"][0], contract, layout["contract_dir"]) + assert not out.exists() + + monkeypatch.setenv("FLUID_DUCKDB_ALLOWED_DIRS", str(landing)) + assert execute_duckdb_build(contract["builds"][0], contract, layout["contract_dir"]) == 0 + assert out.exists() + + ctx = build_acquisition_run_context(contract["builds"][0], contract, layout["contract_dir"]) + con = _connect_for_destination(ctx) + try: + assert con.execute(f"SELECT count(*) FROM read_parquet('{out}')").fetchone() == (1,) + with pytest.raises(duckdb.PermissionException): + con.execute(f"SELECT * FROM read_csv('{layout['outside']}/secret.csv')") + # The read-back connection is the destination's, not the source's. + with pytest.raises(duckdb.PermissionException): + con.execute(f"SELECT * FROM read_csv('{landing}/orders.csv')") + finally: + con.close() + + # The source glob grants the operator-allowed landing directory it names, + # and nothing outside it. + from fluid_build.build_runners.duckdb.runner import _run_allowlist, _run_streams + + src = secure_duckdb_connect(allow=_run_allowlist(ctx, _run_streams(ctx))) + try: + assert src.execute(f"SELECT count(*) FROM read_csv('{landing}/*.csv')").fetchone() == (1,) + with pytest.raises(duckdb.PermissionException): + src.execute(f"SELECT content FROM read_text('{layout['outside']}/secret.csv')") + finally: + src.close() + + +@pytest.mark.parametrize( + "arrival", + [ + # macOS writes AppleDouble '._x' files on exFAT / SMB volumes; DuckDB's + # glob matches them, Python's glob does not. + pytest.param("._orders.csv", id="dotfile"), + pytest.param("late.csv", id="lands_after_the_grant"), + ], +) +def test_an_operator_allowed_landing_glob_reads_every_file_duckdb_expands( + arrival, layout, monkeypatch +): + """A glob granted as the files a grant-time listing found failed the whole + read on any file DuckDB's own expansion adds: a dotfile, or a file that + landed between building the allowlist and running the query.""" + from fluid_build.build_runners._acquisition_common import build_acquisition_run_context + from fluid_build.build_runners.duckdb.runner import ( + _run_allowlist, + _run_streams, + execute_duckdb_build, + ) + + landing = layout["tmp"] / "landing" + landing.mkdir() + (landing / "orders.csv").write_text("id,amount\n1,10\n", encoding="utf-8") + monkeypatch.setenv("FLUID_DUCKDB_ALLOWED_DIRS", str(landing)) + out = layout["contract_dir"] / "out" / "orders.parquet" + contract = _acquisition_contract(str(landing / "*.csv"), out) + + if arrival.startswith("."): + (landing / arrival).write_text("id,amount\n2,20\n", encoding="utf-8") + assert execute_duckdb_build(contract["builds"][0], contract, layout["contract_dir"]) == 0 + assert duckdb.connect().execute(f"SELECT count(*) FROM '{out}'").fetchone() == (2,) + return + + ctx = build_acquisition_run_context(contract["builds"][0], contract, layout["contract_dir"]) + con = secure_duckdb_connect(allow=_run_allowlist(ctx, _run_streams(ctx))) + try: + (landing / arrival).write_text("id,amount\n2,20\n", encoding="utf-8") + assert con.execute(f"SELECT count(*) FROM read_csv('{landing}/*.csv')").fetchone() == (2,) + finally: + con.close() + + +def test_a_sqlite_source_outside_the_allowed_directories_is_refused(layout, monkeypatch): + """The sqlite scanner opens files through its own library, which DuckDB's + allowlist never bounds, so the declared source path was attached wherever + it pointed: any SQLite file on the host could be landed.""" + import sqlite3 + + from fluid_build.build_runners.duckdb.runner import execute_duckdb_build + + try: + probe = duckdb.connect() + probe.execute("INSTALL sqlite; LOAD sqlite") + probe.close() + except duckdb.Error: + pytest.skip("the sqlite extension is not installable here (offline)") + + def contract_for(db: Path, out: Path) -> Dict[str, Any]: + contract = _acquisition_contract("unused", out) + contract["builds"][0]["properties"]["source"] = { + "kind": "sqlite", + "connection": {"path": str(db)}, + "streams": ["cookies"], + "mode": "full_refresh", + } + return contract + + def seed(db: Path, value: str) -> None: + with sqlite3.connect(db) as seeded: + seeded.execute("CREATE TABLE cookies (v TEXT)") + seeded.execute("INSERT INTO cookies VALUES (?)", (value,)) + + cookies = layout["outside"] / "Cookies" # the fixture's $HOME + seed(cookies, SECRET) + stolen = layout["contract_dir"] / "out" / "stolen.parquet" + contract = contract_for(cookies, stolen) + with pytest.raises(DuckDBSandboxError, match="FLUID_DUCKDB_ALLOWED_DIRS"): + execute_duckdb_build(contract["builds"][0], contract, layout["contract_dir"]) + assert not stolen.exists() + + # A sqlite file in the contract's directory still lands. + own = layout["contract_dir"] / "data" / "app.sqlite" + seed(own, "fine") + landed = layout["contract_dir"] / "out" / "own.parquet" + contract = contract_for(own, landed) + assert execute_duckdb_build(contract["builds"][0], contract, layout["contract_dir"]) == 0 + assert duckdb.connect().execute(f"SELECT v FROM '{landed}'").fetchall() == [("fine",)] + + # The operator can allow the directory the outside file sits in. + monkeypatch.setenv("FLUID_DUCKDB_ALLOWED_DIRS", str(layout["outside"])) + contract = contract_for(cookies, stolen) + assert execute_duckdb_build(contract["builds"][0], contract, layout["contract_dir"]) == 0 + + +def test_a_sandbox_refusal_is_not_retried(): + """A refusal is permanent; its path ('pytest-500/...') must not read as an HTTP 500.""" + from fluid_build.providers.local.util.retry import is_retryable_error + + con = secure_duckdb_connect(allow=DuckDBAllowlist.none()) + with pytest.raises(duckdb.PermissionException) as refused: + con.execute("SELECT * FROM read_csv('/tmp/pytest-500/503/x.csv')") + assert not is_retryable_error(refused.value) + + +def test_a_url_is_refused_even_with_its_extension_loaded(): + con = None + try: + con = secure_duckdb_connect(allow=DuckDBAllowlist.none(), extensions=["httpfs"]) + except duckdb.Error: + pytest.skip("httpfs is not installable here (offline)") + with pytest.raises(duckdb.PermissionException): + con.execute("SELECT * FROM read_csv('https://example.invalid/x.csv')") + with pytest.raises(duckdb.PermissionException): + con.execute("SELECT * FROM read_csv('http://169.254.169.254/latest/meta-data/')") + + +# ── a declaration is not a grant of the host ───────────────────────────── +# +# Every declared input and output is granted to the contract's SQL, and the +# contract's author writes the declarations. Before confinement, declaring an +# innocuous glob in $HOME granted all of $HOME, and declaring a credentials file +# granted it directly. + + +def _build_declaring(layout: Dict[str, Path], sql: str, path: str) -> int: + from fluid_build.build_runners.base import _execute_embedded_sql_build + + out = layout["contract_dir"] / "out" / "result.csv" + contract = _contract(sql, out) + contract["builds"][0]["properties"]["parameters"] = {"inputs": [{"name": "d", "path": path}]} + return _execute_embedded_sql_build(contract["builds"][0], contract, layout["contract_dir"]) + + +@pytest.fixture +def home_secrets(layout: Dict[str, Path]) -> Path: + """``$HOME`` (the fixture's ``outside``) with credentials and a decoy CSV.""" + home = layout["outside"] + (home / ".aws").mkdir() + (home / ".aws" / "credentials").write_text(f"[default]\nkey = {SECRET}\n", encoding="utf-8") + (home / "decoy.csv").write_text("a\n1\n", encoding="utf-8") + return home + + +_READ_CREDENTIALS = "SELECT content FROM read_text('~/.aws/credentials')" + + +@pytest.mark.parametrize( + "declared", + [ + pytest.param("{home}/*.csv", id="glob_in_home"), + pytest.param("{home}/.aws/credentials", id="the_file_itself"), + pytest.param("{home}", id="the_home_directory"), + pytest.param("{home}/.aws", id="a_directory_outside"), + pytest.param("{contract}/*/../../outside/.aws/credentials", id="dot_dot_after_glob"), + pytest.param("../outside/*.csv", id="relative_dot_dot"), + ], +) +def test_declaring_a_location_outside_the_contract_grants_nothing( + declared, layout, home_secrets, printed +): + path = declared.format(home=home_secrets, contract=layout["contract_dir"]) + + rc = _build_declaring(layout, _READ_CREDENTIALS, path) + + assert rc == 1, f"{declared}: the build read through a declaration" + assert SECRET not in _written_text(layout["contract_dir"]) + assert SECRET not in "\n".join(printed) + assert "FLUID_DUCKDB_ALLOWED_DIRS" in "\n".join(printed), printed + + +def test_a_relative_declaration_cannot_grant_the_servers_working_directory( + layout, home_secrets, printed, monkeypatch +): + """``path: ./*.csv`` resolves where DuckDB opens it; outside the contract, refused.""" + monkeypatch.chdir(home_secrets) # a server whose working directory is elsewhere + + rc = _build_declaring(layout, _READ_CREDENTIALS, "./*.csv") + + assert rc == 1 + assert SECRET not in _written_text(layout["contract_dir"]) + assert SECRET not in "\n".join(printed) + + +def test_a_symlink_in_the_contract_to_root_grants_nothing(layout, printed): + """DuckDB realpaths allowlist entries, so a declared in-repo symlink to '/' + would grant the whole host while looking like a path inside the contract.""" + from fluid_build.providers.local.local import LocalProvider + + contract_dir = layout["contract_dir"] + (contract_dir / "rootlink").symlink_to("/") + provider = LocalProvider(project="local", region="local", anchor_dir=contract_dir) + for declared in ("rootlink", str(contract_dir / "rootlink"), "rootlink/etc/*"): + with pytest.raises(DuckDBSandboxError): + provider._allowlist([{"path": declared}]) + with pytest.raises(DuckDBSandboxError): + DuckDBAllowlist.none().with_dirs(contract_dir / "rootlink") + + # Through the build: a file reached via the symlink is outside, so refused. + via_link = f"rootlink{layout['outside']}/secret.csv" + assert _build_declaring(layout, "SELECT * FROM d", via_link) == 1 + assert SECRET not in _written_text(contract_dir) + assert SECRET not in "\n".join(printed) + + +def test_an_operator_allowed_glob_grants_nothing_outside_the_operators_directory( + layout, home_secrets, printed, monkeypatch +): + shared = home_secrets / "shared" + shared.mkdir() + (shared / "decoy.csv").write_text("a\n1\n", encoding="utf-8") + monkeypatch.setenv("FLUID_DUCKDB_ALLOWED_DIRS", str(shared)) + + assert _build_declaring(layout, "SELECT * FROM d", f"{shared}/decoy*.csv") == 0, printed + out = layout["contract_dir"] / "out" / "result.csv" + assert out.read_text(encoding="utf-8").splitlines() == ["a", "1"] + + # $HOME holds the operator's directory; the glob does not grant $HOME. + assert _build_declaring(layout, _READ_CREDENTIALS, f"{shared}/decoy*.csv") == 1 + assert SECRET not in _written_text(layout["contract_dir"]) + # Nor does a glob whose directory is above the operator's. + assert _build_declaring(layout, _READ_CREDENTIALS, f"{home_secrets}/sha*/decoy.csv") == 1 + assert SECRET not in _written_text(layout["contract_dir"]) + + +def test_a_declared_glob_inside_the_contract_still_builds(layout, printed): + rc = _build_declaring(layout, "SELECT sum(amount) AS s FROM d", "data/*.csv") + assert rc == 0, printed + out = layout["contract_dir"] / "out" / "result.csv" + assert out.read_text(encoding="utf-8").splitlines() == ["s", "30"] + + +def test_a_declared_output_outside_the_contract_is_refused(layout, printed): + """An output is a write grant: declaring one outside would COPY anywhere.""" + from fluid_build.build_runners.base import _execute_embedded_sql_build + + target = layout["outside"] / "pwned.csv" + contract = _contract("SELECT 1 AS a", target) + assert _execute_embedded_sql_build(contract["builds"][0], contract, layout["contract_dir"]) == 1 + assert not target.exists() + + +def test_a_declared_glob_cannot_read_through_a_matched_symlink(tmp_path, monkeypatch): + """A glob grants its directory; a file in it that is a symlink out of the + roots is refused by DuckDB, which checks each expanded file's realpath.""" + monkeypatch.delenv("FLUID_DUCKDB_ALLOWED_DIRS", raising=False) + root, outside = tmp_path / "root", tmp_path / "outside" + (root / "data").mkdir(parents=True) + outside.mkdir() + (root / "data" / "a.csv").write_text("a\n1\n", encoding="utf-8") + (outside / "s.csv").write_text(f"a\n{SECRET}\n", encoding="utf-8") + + allow = DuckDBAllowlist.none().with_declared("data/*.csv", within=[root], base=root) + assert allow.dirs == (str(root / "data"),) + assert allow.paths == () + con = secure_duckdb_connect(allow=allow) + try: + assert con.execute(f"SELECT * FROM read_csv('{root}/data/*.csv')").fetchall() == [(1,)] + (root / "data" / "b.csv").symlink_to(outside / "s.csv") + for sql in ( + f"SELECT * FROM read_csv('{root}/data/*.csv')", + f"SELECT * FROM read_csv('{root}/data/b.csv')", + ): + with pytest.raises(duckdb.PermissionException): + con.execute(sql).fetchall() + finally: + con.close() + + # A glob whose directory is itself a symlink out of the roots is refused. + (root / "away").symlink_to(outside) + with pytest.raises(DuckDBSandboxError, match="outside the directories"): + DuckDBAllowlist.none().with_declared("away/*.csv", within=[root], base=root) + + +def test_a_glob_over_many_files_grants_one_directory(tmp_path, monkeypatch): + """Granting each match made the allowlist quadratic in the file count and + rebuilt it on every action's connection (20,000 files: ~15 s each).""" + monkeypatch.delenv("FLUID_DUCKDB_ALLOWED_DIRS", raising=False) + for part in range(20): + day = tmp_path / "data" / f"dt={part}" + day.mkdir(parents=True) + for i in range(100): + (day / f"part-{i}.csv").write_text("id\n1\n", encoding="utf-8") + + allow = DuckDBAllowlist.none().with_declared(f"{tmp_path}/data/**/*.csv", within=[tmp_path]) + + assert allow.dirs == (str(tmp_path / "data"),) + assert allow.paths == () + + +# ── a wildcard character in a directory the contract did not declare ───── +# +# The callers join the contract's directory, the working directory or $HOME in +# front of a declared path. A '[', '?' or '*' in one of those is part of a +# directory's name, not a pattern the contract wrote: read as one, the path was +# cut to the directory above it and refused as outside the contract, so every +# input and output of a contract in 'Proj [old]/' failed. + +# DuckDB cannot open a database file under a '?' (it reads one as a URL query), +# so '?' is in the grant test only, not in the builds. +_WILDCARD_DIRS = ["Proj [old]", "p*", "p[1]"] + + +@pytest.fixture(params=_WILDCARD_DIRS) +def wildcard_layout(request, tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Dict[str, Path]: + """:func:`layout`, under a directory whose name holds a wildcard character.""" + monkeypatch.delenv("FLUID_DUCKDB_ALLOWED_DIRS", raising=False) + root = tmp_path / request.param + root.mkdir() + return _make_layout(root, monkeypatch) + + +@pytest.mark.parametrize("declared", ["data/orders.csv", "{contract}/data/orders.csv"]) +def test_a_contract_under_a_wildcard_directory_builds(declared, wildcard_layout, printed): + path = declared.format(contract=wildcard_layout["contract_dir"]) + + rc = _build_declaring(wildcard_layout, "SELECT sum(amount) AS s FROM d", path) + + assert rc == 0, printed + out = wildcard_layout["contract_dir"] / "out" / "result.csv" + assert out.read_text(encoding="utf-8").splitlines() == ["s", "30"] + + +def test_local_apply_under_a_wildcard_directory_reads_and_writes(wildcard_layout): + """The local provider's actions pass absolute paths, the directory joined in.""" + from fluid_build.providers.local.local import LocalProvider + + contract_dir = wildcard_layout["contract_dir"] + provider = LocalProvider(project="local", region="local", anchor_dir=contract_dir) + out = contract_dir / "runtime" / "out" / "r.csv" + actions = [ + {"op": "load_data", "path": str(contract_dir / "data" / "orders.csv"), "table_name": "o"}, + {"op": "sql", "sql": "SELECT count(*) AS n FROM o", "outputs": [str(out)]}, + ] + + results = provider.apply(actions=actions)["results"] + + assert [r["status"] for r in results] == ["ok", "ok"], [r.get("error") for r in results] + assert out.read_text(encoding="utf-8").splitlines() == ["n", "2"] + + +@pytest.mark.parametrize("name", [*_WILDCARD_DIRS, "p?"]) +def test_a_wildcard_directory_is_granted_as_a_directory_not_a_pattern(name, tmp_path, monkeypatch): + monkeypatch.delenv("FLUID_DUCKDB_ALLOWED_DIRS", raising=False) + project = tmp_path / name / "p" + (project / "data").mkdir(parents=True) + (project / "customers.csv").write_text("a\n1\n", encoding="utf-8") + monkeypatch.chdir(project) + + # Relative to the working directory, to ``base``, or already absolute. + for allow in ( + DuckDBAllowlist.none().with_declared("customers.csv", within=[project]), + DuckDBAllowlist.none().with_declared("customers.csv", within=[project], base=project), + DuckDBAllowlist.none().with_declared(str(project / "customers.csv"), within=[project]), + DuckDBAllowlist.none().with_locations("customers.csv"), + ): + assert allow.paths == (str(project / "customers.csv"),) + assert allow.dirs == () + # A glob the contract wrote is still a glob, cut at its own wildcard. + allow = DuckDBAllowlist.none().with_declared("data/*.csv", within=[project]) + assert allow.dirs == (str(project / "data"),) + # So is a last component, even one a file has as its literal name: DuckDB + # reads what 'a[1].csv' matches ('a1.csv'), and the directory holds both. + (project / "data" / "a[1].csv").write_text("a\n1\n", encoding="utf-8") + allow = DuckDBAllowlist.none().with_declared("data/a[1].csv", within=[project]) + assert allow.dirs == (str(project / "data"),) + # An output not written yet: the file, not the directory above the project. + allow = DuckDBAllowlist.none().with_declared("out/x.csv", within=[project]) + assert allow.paths == (str(project / "out" / "x.csv"),) + # $HOME with a '[' in it. + monkeypatch.setenv("HOME", str(project)) + allow = DuckDBAllowlist.none().with_declared("~/customers.csv", within=[project]) + assert allow.paths == (str(project / "customers.csv"),) + + +def test_a_wildcard_directory_grants_nothing_beside_it(tmp_path, monkeypatch): + """Read literally, the directory is narrower than the pattern, never wider.""" + monkeypatch.delenv("FLUID_DUCKDB_ALLOWED_DIRS", raising=False) + project, sibling = tmp_path / "p [x]", tmp_path / "p x" + project.mkdir() + sibling.mkdir() + (project / "f.csv").write_text("a\n1\n", encoding="utf-8") + (sibling / "f.csv").write_text(f"a\n{SECRET}\n", encoding="utf-8") + (tmp_path / "secret.csv").write_text(f"a\n{SECRET}\n", encoding="utf-8") + + allow = DuckDBAllowlist.none().with_declared("f.csv", within=[project], base=project) + assert allow.paths == (str(project / "f.csv"),) + con = secure_duckdb_connect(allow=allow) + try: + # DuckDB matches 'p [x]' as a pattern first, so it reads 'p x': refused. + with pytest.raises(duckdb.PermissionException): + con.execute(f"SELECT * FROM read_csv('{project}/f.csv')").fetchall() + with pytest.raises(duckdb.PermissionException): + con.execute(f"SELECT * FROM read_csv('{tmp_path}/secret.csv')").fetchall() + finally: + con.close() + # A wildcard-named directory read literally is still resolved before it is confined. + (project / "[l]").symlink_to(tmp_path) + escapes = ("../secret.csv", "../*.csv", "../p x/f.csv", "*/../../secret.csv", "[l]/secret.csv") + for escape in escapes: + with pytest.raises(DuckDBSandboxError, match="outside the directories"): + DuckDBAllowlist.none().with_declared(escape, within=[project], base=project) + + +def test_the_readme_commands_build_from_a_wildcard_directory(tmp_path): + """The reported case: examples/02, run as its README says, from 'repo [old]/'.""" + workspace = tmp_path / "repo [old]" + example = REPO_ROOT / "examples" / "02-csv-to-data-product" + target = workspace / "examples" / example.name + target.parent.mkdir(parents=True) + shutil.copytree(example, target) + env = {k: v for k, v in os.environ.items() if k != "FLUID_DUCKDB_ALLOWED_DIRS"} + contract = f"examples/{example.name}/contract.fluid.yaml" + command = ["apply", contract, "--provider", "local", "--mode", "amend-and-build", "--yes"] + + run = subprocess.run( + [sys.executable, "-m", "fluid_build.cli", *command], + cwd=workspace, + env=env, + capture_output=True, + text=True, + encoding="utf-8", + errors="replace", + timeout=180, + ) + + assert run.returncode == 0, f"stdout:\n{run.stdout}\nstderr:\n{run.stderr}" + out = target / "runtime" / "out" / "customer-clean-v1.csv" + assert len(out.read_text(encoding="utf-8").splitlines()) == 5 # header + 4 rows + + +def test_allowlist_entries_are_deduplicated_in_linear_time(): + """De-duplicating against a tuple and a list was O(n^2).""" + import time + + entries = [f"/data/f{i}.csv" for i in range(40_000)] + started = time.perf_counter() + added = sandbox._new(tuple(entries[:20_000]), [*entries, *entries]) + elapsed = time.perf_counter() - started + + assert added == tuple(entries[20_000:]) + assert elapsed < 1.0, f"{elapsed:.1f}s to de-duplicate {len(entries)} entries" + allow = DuckDBAllowlist.none().with_paths(*entries[:100], *entries[:100]) + assert allow.with_paths(*entries[:100]).paths == tuple(entries[:100]) + + +@pytest.mark.parametrize("target", ["../..", "{home}"]) +def test_a_symlinked_runtime_does_not_grant_where_it_points(target, tmp_path, monkeypatch, printed): + """``./runtime`` is in the working directory, the contract's own in the + usual ``cd product && fluid apply``, so the contract's repository could + ship ``runtime -> ../..`` and have DuckDB grant $HOME (it realpaths each + allowlist entry).""" + from fluid_build.build_runners.base import _execute_embedded_sql_build + + home = tmp_path / "home" + (home / ".aws").mkdir(parents=True) + (home / ".aws" / "credentials").write_text(f"[default]\nkey = {SECRET}\n", encoding="utf-8") + contract_dir = home / "src" / "product" + contract_dir.mkdir(parents=True) + (contract_dir / "runtime").symlink_to(target.format(home=home)) + monkeypatch.setenv("HOME", str(home)) + monkeypatch.delenv("FLUID_DUCKDB_ALLOWED_DIRS", raising=False) + monkeypatch.delenv("FLUID_UPSTREAM_CONTRACTS", raising=False) + monkeypatch.chdir(contract_dir) + + out = contract_dir / "out" / "result.csv" + contract = _contract(_READ_CREDENTIALS, out) + assert _execute_embedded_sql_build(contract["builds"][0], contract, contract_dir) == 1 + assert not out.exists() or SECRET not in out.read_text(encoding="utf-8") + assert SECRET not in "\n".join(printed) + + # The local provider's other DuckDB call site (ducksql.apply_sql), too. + from fluid_build.providers.local import ducksql + + granted: List[DuckDBAllowlist] = [] + + def capture(database: Any = ":memory:", *, allow: DuckDBAllowlist, **kwargs: Any) -> Any: + granted.append(allow) + return secure_duckdb_connect(database, allow=allow, **kwargs) + + monkeypatch.setattr(ducksql, "secure_duckdb_connect", capture) + assert ducksql.apply_sql([]) == [] + con = secure_duckdb_connect(allow=granted[0]) + try: + with pytest.raises(duckdb.PermissionException): + con.execute(_READ_CREDENTIALS).fetchall() + finally: + con.close() + + # The same build without the symlink still writes its own output. + (contract_dir / "runtime").unlink() + contract = _contract("SELECT 1 AS a", out) + assert _execute_embedded_sql_build(contract["builds"][0], contract, contract_dir) == 0 + + +def test_a_runtime_symlink_inside_the_contract_is_still_granted(tmp_path, monkeypatch): + from fluid_build.providers._duckdb_sandbox import unaliased_dir + from fluid_build.providers.local.local import LocalProvider + + contract_dir = tmp_path / "product" + (contract_dir / "build" / "runtime").mkdir(parents=True) + (contract_dir / "runtime").symlink_to("build/runtime") + monkeypatch.chdir(contract_dir) + monkeypatch.delenv("FLUID_DUCKDB_ALLOWED_DIRS", raising=False) + + assert unaliased_dir("runtime") is None + assert unaliased_dir("runtime", within=[contract_dir]) == str(contract_dir / "runtime") + provider = LocalProvider(project="local", region="local", anchor_dir=contract_dir) + assert str(contract_dir / "runtime") in provider._allowlist([]).dirs + + +@pytest.mark.parametrize("entry", ["relative/dir", "/"]) +def test_bad_operator_dirs_are_refused(entry, tmp_path, monkeypatch): + monkeypatch.setenv("FLUID_DUCKDB_ALLOWED_DIRS", entry) + with pytest.raises(DuckDBSandboxError, match="FLUID_DUCKDB_ALLOWED_DIRS"): + DuckDBAllowlist.none().with_declared(tmp_path / "x.csv", within=[tmp_path]) + + +def test_contract_tests_action_cannot_declare_a_file_outside(layout): + from types import SimpleNamespace + + from fluid_build.contract_tests import LocalProviderError, apply_action + + action = { + "op": "add", + "resource_type": "sql", + "id": "steal", + "inputs": {"s": {"path": str(layout["outside"] / "secret.csv")}}, + "outputs": {"path": str(layout["contract_dir"] / "out" / "s.csv"), "format": "csv"}, + "sql": "SELECT * FROM s", + } + with pytest.raises(LocalProviderError, match="FLUID_DUCKDB_ALLOWED_DIRS"): + apply_action(action, SimpleNamespace(dry_run=False)) + assert not (layout["contract_dir"] / "out" / "s.csv").exists() + + +# ── one failure does not poison the run's later actions ────────────────── +# +# Every action of one apply shares one session database file. DuckDB shares one +# instance per file per process and the sandbox locks it, so a connection left +# open by a failure (its traceback keeps it alive) made every later action fail +# with "the configuration has been locked". + + +def _apply_after_a_failure(provider: Any, data: Path, out: Path) -> List[Dict[str, Any]]: + actions = [ + {"op": "load_data", "path": str(data / "orders.csv"), "table_name": "orders"}, + { + "op": "sql", + "sql": "SELECT * FROM orders", + "output_table": "t1", + "outputs": [str(out / "1.csv")], + }, + {"op": "sql", "sql": "SELECT * FROM no_such_table", "outputs": [str(out / "2.csv")]}, + {"op": "sql", "sql": "SELECT count(*) AS n FROM t1", "outputs": [str(out / "3.csv")]}, + {"op": "load_data", "path": str(data / "missing_*.csv"), "table_name": "m"}, + {"op": "sql", "sql": "SELECT count(*) AS n FROM orders", "outputs": [str(out / "4.csv")]}, + {"op": "materialize", "source_table": "t1", "dst": str(out / "5.csv")}, + ] + return provider.apply(actions=actions)["results"] + + +def test_a_failing_action_does_not_fail_the_actions_after_it(layout): + from fluid_build.providers.local.local import LocalProvider + + provider = LocalProvider(project="local", region="local", anchor_dir=layout["contract_dir"]) + out = layout["contract_dir"] / "out" + + results = _apply_after_a_failure(provider, layout["contract_dir"] / "data", out) + + assert [r["status"] for r in results] == ["ok", "ok", "error", "ok", "error", "ok", "ok"], [ + r.get("error") for r in results + ] + assert "no_such_table" in results[2]["error"] + assert "No files found" in results[4]["error"] + assert (out / "3.csv").read_text(encoding="utf-8").splitlines() == ["n", "2"] + assert (out / "5.csv").read_text(encoding="utf-8").splitlines()[1:] == ["1,10", "2,20"] + + +def test_persist_mode_applies_one_after_another_after_a_failure(layout): + """``persist=True`` shares ~/.fluid/local.db across runs in one process.""" + from fluid_build.providers.local.local import LocalProvider + + data, out = layout["contract_dir"] / "data", layout["contract_dir"] / "out" + for _ in range(2): + provider = LocalProvider( + project="local", region="local", anchor_dir=layout["contract_dir"], persist=True + ) + statuses = [r["status"] for r in _apply_after_a_failure(provider, data, out)] + assert statuses == ["ok", "ok", "error", "ok", "error", "ok", "ok"] + + +def test_contract_tests_actions_share_a_file_database_after_a_failure(layout, monkeypatch): + from types import SimpleNamespace + + from fluid_build.contract_tests import LocalProviderError, apply_action + + monkeypatch.setenv("FLUID_LOCAL_DUCKDB_PATH", str(layout["contract_dir"] / "ct.duckdb")) + action = { + "op": "add", + "resource_type": "sql", + "id": "orders", + "inputs": {"orders": {"path": str(layout["contract_dir"] / "data" / "orders.csv")}}, + "outputs": {"path": str(layout["contract_dir"] / "out" / "r.csv"), "format": "csv"}, + "sql": "SELECT * FROM no_such_table", + } + # The caller keeps the failure (as a CLI reporting it does): its traceback + # must not keep the file's locked instance open. + with pytest.raises(LocalProviderError, match="no_such_table") as failed: + apply_action(action, SimpleNamespace(dry_run=False)) + action["sql"] = "SELECT id FROM orders" + apply_action(action, SimpleNamespace(dry_run=False)) + assert failed.value is not None + assert (layout["contract_dir"] / "out" / "r.csv").exists() + + +def test_a_second_connection_to_an_open_file_is_refused_clearly(tmp_path): + db = tmp_path / "shared.duckdb" + first = secure_duckdb_connect(db, allow=DuckDBAllowlist.none()) + try: + for config in (None, {"threads": 2}): + with pytest.raises(DuckDBSandboxError, match="already open in this process"): + secure_duckdb_connect(db, allow=DuckDBAllowlist.none(), config=config) + finally: + first.close() + # Closed, the file opens again. + secure_duckdb_connect(db, allow=DuckDBAllowlist.none()).close() + + +# ── extensions: what autoload-off costs, and what a loaded scanner reaches ─ + + +def test_a_function_from_an_unloaded_extension_reads_as_a_sandbox_refusal(layout, printed): + """``sqlite_scan`` (and read_xlsx, ST_Read, delta_scan, iceberg_scan) worked + through autoloading; under the sandbox they are not in the catalog, and the + error must say why rather than advise a SET the lock refuses.""" + from fluid_build.providers._duckdb_sandbox import is_sandbox_refusal, sandbox_refusal_hint + from fluid_build.providers.local.util.retry import is_retryable_error + + con = secure_duckdb_connect(allow=DuckDBAllowlist.none()) + with pytest.raises(duckdb.CatalogException) as refused: + con.execute("SELECT * FROM sqlite_scan('data/app.sqlite', 't')") + assert "exists in the sqlite_scanner extension" in str(refused.value) + assert is_sandbox_refusal(refused.value) + assert not is_retryable_error(refused.value) + assert "extension autoloading off" in sandbox_refusal_hint( + DuckDBAllowlist.none(), refused.value + ) + + rc = _build("SELECT * FROM sqlite_scan('data/app.sqlite', 't')", layout) + assert rc == 1 + assert "extension autoloading off" in "\n".join(printed), printed + + +def test_contract_sql_runs_on_a_connection_with_no_database_scanner(layout, printed): + """sqlite/postgres/mysql scanners are not bounded by the allowlist (next + test), so the connection contract SQL runs on must never have one loaded.""" + rc = _build( + "SELECT count(*) AS n FROM duckdb_functions() WHERE function_name IN " + "('sqlite_scan', 'sqlite_attach', 'postgres_scan', 'postgres_query', 'mysql_query')", + layout, + ) + assert rc == 0, printed + out = layout["contract_dir"] / "out" / "result.csv" + assert out.read_text(encoding="utf-8").splitlines() == ["n", "0"] + + +def _with_extension(name: str) -> Any: + try: + return secure_duckdb_connect(allow=DuckDBAllowlist.none(), extensions=[name]) + except duckdb.Error as exc: + raise pytest.skip.Exception( + f"the {name} extension is not installable here (offline)" + ) from exc + + +def test_a_loaded_sqlite_scanner_is_not_bounded_by_the_allowlist(tmp_path): + """Pins a documented limit: once ``sqlite`` is loaded, any SQLite file the + process can read is reachable, whatever the allowlist says. A call site that + loads it must keep contract SQL off the connection.""" + db = tmp_path / "elsewhere.sqlite" + seed = duckdb.connect() + try: + seed.execute("INSTALL sqlite; LOAD sqlite") + except duckdb.Error: + pytest.skip("the sqlite extension is not installable here (offline)") + seed.execute(f"ATTACH '{db}' AS s (TYPE sqlite)") + seed.execute(f"CREATE TABLE s.t AS SELECT '{SECRET}' AS v") + seed.close() + + con = _with_extension("sqlite") + assert con.execute(f"SELECT v FROM sqlite_scan('{db}', 't')").fetchall() == [(SECRET,)] + + +def test_a_loaded_postgres_scanner_is_not_bounded_by_the_allowlist(): + """Pins a documented limit: ``postgres_scan`` opens its own socket, so the + failure is the server's (nothing listens), not the sandbox's refusal.""" + con = _with_extension("postgres") + with pytest.raises(duckdb.Error) as failed: + con.execute( + "SELECT * FROM postgres_scan('host=127.0.0.1 port=9 connect_timeout=2', 'public', 't')" + ) + assert not isinstance(failed.value, duckdb.PermissionException), failed.value + assert "Permission Error" not in str(failed.value) + + +# ── the helper's own guarantees ────────────────────────────────────────── + + +def test_allowed_directory_is_a_directory_not_a_name_prefix(tmp_path): + (tmp_path / "data").mkdir() + (tmp_path / "data-other").mkdir() + (tmp_path / "data-other" / "x.csv").write_text("a\n1\n", encoding="utf-8") + con = secure_duckdb_connect(allow=DuckDBAllowlist.none().with_dirs(tmp_path / "data")) + with pytest.raises(duckdb.PermissionException): + con.execute(f"SELECT * FROM read_csv('{tmp_path}/data-other/x.csv')") + + +def test_none_allowlist_reaches_no_file(tmp_path): + (tmp_path / "x.csv").write_text("a\n1\n", encoding="utf-8") + con = secure_duckdb_connect(allow=DuckDBAllowlist.none()) + assert con.execute("SELECT 41 + 1").fetchone() == (42,) + with pytest.raises(duckdb.PermissionException): + con.execute(f"SELECT * FROM read_csv('{tmp_path}/x.csv')") + + +def test_tilde_follows_home_and_cannot_disagree_with_it(tmp_path, monkeypatch): + """duckdb/duckdb#26064: '~' must mean $HOME both when checked and when opened.""" + home, allowed = tmp_path / "home", tmp_path / "allowed" + home.mkdir() + allowed.mkdir() + (home / "s.csv").write_text(f"who\n{SECRET}\n", encoding="utf-8") + (allowed / "s.csv").write_text("who\nallowed\n", encoding="utf-8") + monkeypatch.setenv("HOME", str(home)) + + con = secure_duckdb_connect(allow=DuckDBAllowlist.none().with_dirs(allowed)) + with pytest.raises(duckdb.PermissionException): + con.execute("SELECT * FROM read_csv('~/s.csv')").fetchall() + with pytest.raises(duckdb.Error): + con.execute("SET home_directory = '/'") + + # And when $HOME itself is granted, '~' reads $HOME, not something else. + con = secure_duckdb_connect(allow=DuckDBAllowlist.none().with_dirs(home)) + assert con.execute("SELECT * FROM read_csv('~/s.csv')").fetchall() == [(SECRET,)] + + +def test_file_database_writes_and_checkpoints_under_the_lock(tmp_path): + """A file-backed database needs no grant of its own (WAL + checkpoint).""" + db = tmp_path / "db" / "run.duckdb" + db.parent.mkdir() + con = secure_duckdb_connect(db, allow=DuckDBAllowlist.none()) + con.execute("CREATE TABLE t AS SELECT range AS a FROM range(50000)") + con.execute("INSERT INTO t SELECT * FROM range(10)") + con.execute("CHECKPOINT") + con.close() + con = secure_duckdb_connect(db, allow=DuckDBAllowlist.none(), read_only=True) + assert con.execute("SELECT count(*) FROM t").fetchone() == (50010,) + # ...and the database's directory was not granted to the SQL. + (db.parent / "neighbour.csv").write_text("a\n1\n", encoding="utf-8") + with pytest.raises(duckdb.PermissionException): + con.execute(f"SELECT * FROM read_csv('{db.parent}/neighbour.csv')") + + +def test_before_lock_attach_is_readable_but_a_new_attach_is_not(tmp_path): + src = tmp_path / "src.duckdb" + seed = duckdb.connect(str(src)) + seed.execute("CREATE TABLE t AS SELECT 7 AS v") + seed.close() + other = tmp_path / "other.duckdb" + duckdb.connect(str(other)).close() + + con = secure_duckdb_connect( + allow=DuckDBAllowlist.none(), + before_lock=lambda c: c.execute(f"ATTACH '{src}' AS s (READ_ONLY)"), + ) + assert con.execute("SELECT v FROM s.t").fetchone() == (7,) + with pytest.raises(duckdb.PermissionException): + con.execute(f"ATTACH '{other}' AS o") + + +def test_config_is_applied_before_the_lock_and_frozen_after(tmp_path): + con = secure_duckdb_connect(allow=DuckDBAllowlist.none(), config={"threads": 3}) + assert con.execute("SELECT current_setting('threads')").fetchone() == (3,) + for statement in ( + "SET threads = 1", + "SET enable_external_access = true", + "SET autoload_known_extensions = true", + "SET allow_persistent_secrets = true", + "RESET lock_configuration", + ): + with pytest.raises(duckdb.Error): + con.execute(statement) + + +def test_settings_read_back_as_locked_down(tmp_path): + con = secure_duckdb_connect(allow=DuckDBAllowlist.none().with_dirs(tmp_path)) + settings = dict( + con.execute( + "SELECT name, value FROM duckdb_settings() WHERE name IN (" + "'enable_external_access', 'lock_configuration', 'autoload_known_extensions', " + "'autoinstall_known_extensions', 'allow_persistent_secrets', " + "'allow_community_extensions', 'allowed_directories')" + ).fetchall() + ) + assert settings["enable_external_access"] == "false" + assert settings["lock_configuration"] == "true" + assert settings["autoload_known_extensions"] == "false" + assert settings["autoinstall_known_extensions"] == "false" + assert settings["allow_persistent_secrets"] == "false" # pragma: allowlist secret + assert settings["allow_community_extensions"] == "false" + assert str(tmp_path) in settings["allowed_directories"] + + +@pytest.mark.parametrize( + "entry", + ["~", "~/data", "/", "s3://bucket/data/"], +) +def test_bad_local_allowlist_entries_are_refused(entry): + with pytest.raises(DuckDBSandboxError): + DuckDBAllowlist.none().with_dirs(entry) + + +@pytest.mark.parametrize( + "prefix", + ["/local/path", "ftp://host/x/", "s3://bucket/a/../b/"], +) +def test_bad_remote_prefixes_are_refused(prefix): + with pytest.raises(DuckDBSandboxError): + DuckDBAllowlist.none().with_remote(prefix) + + +def test_declared_locations_grant_the_narrowest_thing(tmp_path): + (tmp_path / "dir").mkdir() + allow = DuckDBAllowlist.none().with_locations( + tmp_path / "file.csv", + tmp_path / "dir", + f"{tmp_path}/landing/*/part-*.parquet", + "s3://bucket/bronze/orders/*.parquet", + "https://host/data/file.csv", + "gs://whole-bucket", + "s3://b2/*.parquet", + ) + # A glob grants the directory above its first wildcard: DuckDB's own + # expansion adds files a listing now would miss (dotfiles, later arrivals). + assert allow.paths == (str(tmp_path / "file.csv"),) + assert allow.dirs == (str(tmp_path / "dir"), str(tmp_path / "landing")) + assert allow.remote_prefixes == ( + "s3://bucket/bronze/orders/", + "https://host/data/", + "gs://whole-bucket/", + "s3://b2/", + ) + + +@pytest.mark.parametrize("version", ["1.4.4", "1.4.3", "1.2.2", "0.10.0", "garbage", None]) +def test_a_duckdb_older_than_the_floor_is_refused(version, monkeypatch): + monkeypatch.setattr(duckdb, "__version__", version, raising=False) + with pytest.raises(DuckDBSandboxError): + secure_duckdb_connect(allow=DuckDBAllowlist.none()) + + +def test_the_floor_matches_the_declared_dependency(): + """pyproject's duckdb floor and the runtime floor are the same number.""" + pyproject = (REPO_ROOT / "pyproject.toml").read_text(encoding="utf-8") + floor = ".".join(str(n) for n in MIN_DUCKDB_VERSION) + assert f'"duckdb>={floor}"' in pyproject + + +def test_the_installed_duckdb_meets_the_floor(): + assert sandbox.duckdb_version(duckdb) >= MIN_DUCKDB_VERSION + + +# ── no connection bypasses the helper ──────────────────────────────────── + + +_HELPER = PACKAGE / "providers" / "_duckdb_sandbox.py" + + +def _duckdb_aliases(tree: ast.AST) -> tuple: + """Names bound to the duckdb module, and to anything imported from it.""" + modules, members = {"duckdb"}, set() + for node in ast.walk(tree): + if isinstance(node, ast.Import): + for alias in node.names: + if alias.name == "duckdb": + modules.add(alias.asname or "duckdb") + elif ( + isinstance(node, ast.ImportFrom) + and node.level == 0 + and (node.module or "").split(".")[0] == "duckdb" + ): + for alias in node.names: + members.add(alias.asname or alias.name) + return modules, members + + +def _bypasses(path: Path) -> List[str]: + """Each place ``path`` opens or queries DuckDB without the helper. + + * any call on the duckdb module (``duckdb.connect``, ``duckdb.sql`` and the + other module-level functions run on an unsandboxed default connection); + * a function imported from duckdb (``from duckdb import connect``); + * ``.connect(...)``, which catches the module + held in a variable (``self._duckdb.connect``, ``_Duck.get().connect``). + """ + source = path.read_text(encoding="utf-8") + tree = ast.parse(source, filename=str(path)) + modules, members = _duckdb_aliases(tree) + found = [] + for node in ast.walk(tree): + if not isinstance(node, ast.Call): + continue + func = node.func + where = ( + f"{path.relative_to(REPO_ROOT) if REPO_ROOT in path.parents else path}:{node.lineno}" + ) + if isinstance(func, ast.Attribute): + receiver = func.value + if isinstance(receiver, ast.Name) and receiver.id in modules: + found.append(f"{where} duckdb.{func.attr}(...)") + elif func.attr == "connect" and "duck" in ast.unparse(receiver).lower(): + found.append(f"{where} {ast.unparse(func)}(...)") + elif isinstance(func, ast.Name) and func.id in members: + found.append(f"{where} {func.id}(...) imported from duckdb") + return found + + +def test_the_bypass_detector_sees_the_helpers_own_connect(): + """A guard that has never fired proves nothing: it must see the one real call.""" + assert [b for b in _bypasses(_HELPER) if "duckdb.connect(" in b] + + +@pytest.mark.parametrize( + "shape", + [ + "import duckdb\nduckdb.connect(':memory:')\n", + "def f():\n import duckdb\n return duckdb.connect()\n", + "import duckdb as ddb\nddb.sql('select 1')\n", + "from duckdb import connect\nconnect()\n", + "from duckdb import connect as c\nc(':memory:')\n", + "class X:\n def f(self):\n return self._duckdb.connect(':memory:')\n", + "def f():\n return _Duck.get().connect(database='x.db')\n", + ], +) +def test_the_bypass_detector_catches_each_shape(shape, tmp_path): + sample = tmp_path / "probe.py" + sample.write_text(shape, encoding="utf-8") + assert _bypasses(sample), shape + + +def test_every_duckdb_connection_goes_through_the_sandbox(): + offenders = [] + for path in sorted(PACKAGE.rglob("*.py")): + if path == _HELPER: + continue + offenders.extend(_bypasses(path)) + assert not offenders, ( + "DuckDB opened or queried without fluid_build.providers._duckdb_sandbox." + "secure_duckdb_connect; contract SQL on these connections can read any host " + "file:\n " + "\n ".join(offenders) + ) + + +def test_no_raw_connect_in_the_helper_beyond_its_one_call(): + calls = [b for b in _bypasses(_HELPER) if "duckdb." in b] + assert len(calls) == 1, calls + + +if os.name == "nt": # pragma: no cover - the sandbox tests assume POSIX paths + pytestmark = pytest.mark.skip(reason="POSIX paths") diff --git a/tests/providers/test_local_validation_sql_safety.py b/tests/providers/test_local_validation_sql_safety.py index 859de31b..6f02a7c3 100644 --- a/tests/providers/test_local_validation_sql_safety.py +++ b/tests/providers/test_local_validation_sql_safety.py @@ -60,7 +60,7 @@ def close(self) -> None: # pragma: no cover - trivial class _ExplodingDuckDB: - """Stand-in duckdb module. + """Stand-in for the provider's sandboxed ``_connect``. The harmless first ``connect(":memory:")`` is allowed (it carries no user-controlled identifier) but returns a connection whose ``execute`` @@ -104,7 +104,7 @@ class TestGetResourceSchemaInjection: def test_malicious_table_never_opens_db(self, tmp_path, monkeypatch): provider = LocalValidationProvider({"base_dir": str(tmp_path)}) # Force any DB access to fail loudly — the guard must run first. - monkeypatch.setattr(provider, "_get_duckdb", lambda: _ExplodingDuckDB()) + monkeypatch.setattr(provider, "_connect", _ExplodingDuckDB.connect) spec = _duckdb_resource_spec(tmp_path, schema="main", table=MALICIOUS_TABLE) # The outer handler re-raises validation failures as a clean Exception; @@ -117,7 +117,7 @@ def test_malicious_table_never_opens_db(self, tmp_path, monkeypatch): def test_malicious_schema_never_opens_db(self, tmp_path, monkeypatch): provider = LocalValidationProvider({"base_dir": str(tmp_path)}) - monkeypatch.setattr(provider, "_get_duckdb", lambda: _ExplodingDuckDB()) + monkeypatch.setattr(provider, "_connect", _ExplodingDuckDB.connect) spec = _duckdb_resource_spec(tmp_path, schema=MALICIOUS_SCHEMA, table="orders") with pytest.raises(Exception) as exc_info: @@ -157,7 +157,7 @@ def test_normal_table_introspects_successfully(self, tmp_path): class TestRunQualityChecksInjection: def test_malicious_table_returns_issue_without_query(self, tmp_path, monkeypatch): provider = LocalValidationProvider({"base_dir": str(tmp_path)}) - monkeypatch.setattr(provider, "_get_duckdb", lambda: _ExplodingDuckDB()) + monkeypatch.setattr(provider, "_connect", _ExplodingDuckDB.connect) spec = _duckdb_resource_spec(tmp_path, schema="main", table=MALICIOUS_TABLE) # run_quality_checks fails closed by returning an error ValidationIssue @@ -169,7 +169,7 @@ def test_malicious_table_returns_issue_without_query(self, tmp_path, monkeypatch def test_malicious_schema_returns_issue_without_query(self, tmp_path, monkeypatch): provider = LocalValidationProvider({"base_dir": str(tmp_path)}) - monkeypatch.setattr(provider, "_get_duckdb", lambda: _ExplodingDuckDB()) + monkeypatch.setattr(provider, "_connect", _ExplodingDuckDB.connect) spec = _duckdb_resource_spec(tmp_path, schema=MALICIOUS_SCHEMA, table="orders") issues = provider.run_quality_checks(spec, rules=[{"type": "not_null", "column": "id"}]) diff --git a/tests/test_contract_tests_branches.py b/tests/test_contract_tests_branches.py index 2d26e7c7..ac3efb86 100644 --- a/tests/test_contract_tests_branches.py +++ b/tests/test_contract_tests_branches.py @@ -349,11 +349,13 @@ def test_full_action_with_targets(self, tmp_path): except LocalProviderError: pass # numpy/pandas not installed - def test_invalid_targets(self, tmp_path): + def test_invalid_targets(self, tmp_path, monkeypatch): try: import duckdb except ImportError: pytest.skip("duckdb not installed") + # Declared inputs must be under the working directory (DuckDB sandbox). + monkeypatch.chdir(tmp_path) csv_in = tmp_path / "in.csv" csv_in.write_text("x\n1") diff --git a/tests/test_fix_a4_bugs.py b/tests/test_fix_a4_bugs.py index e0d8c41c..3672c635 100644 --- a/tests/test_fix_a4_bugs.py +++ b/tests/test_fix_a4_bugs.py @@ -338,9 +338,12 @@ def test_csv_contract_writes_csv(self, tmp_path: Path): else: assert not str(out_spec).endswith(".parquet") - def test_parquet_apply_creates_parquet_file(self, tmp_path: Path): + def test_parquet_apply_creates_parquet_file(self, tmp_path: Path, monkeypatch): """End-to-end: apply writes actual parquet bytes.""" pytest.importorskip("duckdb") + # Run from the project directory: without a contract directory the + # sandbox confines declared outputs to the working directory. + monkeypatch.chdir(tmp_path) import duckdb from fluid_build.providers.local.local import LocalProvider diff --git a/tests/test_init.py b/tests/test_init.py index 611a8d21..9306a2ee 100644 --- a/tests/test_init.py +++ b/tests/test_init.py @@ -583,18 +583,19 @@ def test_duckdb_not_installed_no_raise(self, tmp_path, logger): with patch.dict("sys.modules", {"duckdb": None}): init_local_db(tmp_path, "local", logger) - def test_duckdb_available_creates_db_dir(self, tmp_path, logger, monkeypatch): + def test_duckdb_available_creates_db_dir(self, tmp_path, logger): + pytest.importorskip("duckdb") from fluid_build.cli.init import init_local_db - mock_conn = MagicMock() - mock_duckdb = MagicMock() - mock_duckdb.connect.return_value = mock_conn - monkeypatch.delitem(sys.modules, "duckdb", raising=False) - with patch.dict("sys.modules", {"duckdb": mock_duckdb}): - with patch("fluid_build.cli.init.RICH_AVAILABLE", False): - init_local_db(tmp_path, "local", logger) - mock_duckdb.connect.assert_called_once() - mock_conn.close.assert_called_once() + # The connection is the sandboxed one (providers/_duckdb_sandbox), which + # refuses a stand-in duckdb module it cannot version-check: a real one. + with patch("fluid_build.cli.init.RICH_AVAILABLE", False): + init_local_db(tmp_path, "local", logger) + db = tmp_path / ".fluid" / "db.duckdb" + assert db.is_file() + import duckdb + + duckdb.connect(str(db)).close() # closed by init: reopening is not locked out # ===========================================================================