Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
354 changes: 124 additions & 230 deletions Cargo.lock

Large diffs are not rendered by default.

63 changes: 62 additions & 1 deletion docs/catalog/duckdb.md
Original file line number Diff line number Diff line change
Expand Up @@ -506,6 +506,67 @@ A `create server` statement example used to access local Iceberg service:
);
```

#### DuckLake

This is to access an existing [DuckLake](https://ducklake.select/) catalog backed by PostgreSQL, with data files on S3-compatible storage. The wrapper attaches the catalog read-only and does not create or migrate it. Local catalog and data files are disabled, as with the other remote server types.

| Server Option | Description | Required | Default |
| ------------- | ----------- | :------: | ------- |
| type | Server type, must be `ducklake` | Y | |
| metadata_path | Catalog connection string, including the `postgres:` prefix | Y | |
| metadata_schema | Schema containing the DuckLake metadata tables in PostgreSQL | | DuckLake default (`main`) |
| data_path | Optional data location override; omit to use the catalog's stored location | | Stored catalog path |
| key_id | S3 access key ID | For private S3 storage | |
| secret | S3 secret access key | For private S3 storage | |

The S3 options listed above, including `region`, `endpoint`, `url_style`, `use_ssl` and `session_token`, are also supported. Use these for services such as MinIO, Supabase Storage, or Cloudflare R2's S3 endpoint. The `ducklake`, `postgres`, `httpfs` and `parquet` extensions are installed and loaded automatically.

Store the complete catalog connection string in Vault to protect its credentials:

```sql
select vault.create_secret(
'postgres:host=catalog.example.com port=5432 dbname=ducklake user=reader password=<password> sslmode=require',
'ducklake_catalog'
);

create server ducklake_server
foreign data wrapper duckdb_wrapper
options (
type 'ducklake',
vault_metadata_path '<catalog_secret_id>',
vault_key_id '<s3_key_secret_id>',
vault_secret '<s3_secret_secret_id>',
region 'us-east-1'
);
```

Alternatively, set `metadata_path`, `key_id` and `secret` directly. The metadata catalog must already be initialized and accessible to the PostgreSQL user in the connection string.

Import tables from a logical DuckLake schema (this is separate from `metadata_schema`):

```sql
create schema if not exists ducklake;

import foreign schema "main"
limit to (products)
from server ducklake_server into ducklake;

select * from ducklake.products;
```

Imported tables retain their original names. `EXCEPT` and importing the whole schema are also supported. To define a table manually, qualify its source with the `ducklake` catalog alias:

```sql
create foreign table ducklake.products (
id bigint,
name text
)
server ducklake_server
options (table 'ducklake.main.products');
```

See [DuckLake connection options](https://ducklake.select/docs/stable/duckdb/usage/connecting) for details about metadata and data paths.

#### MotherDuck

This is to access [MotherDuck](https://motherduck.com/), a cloud-hosted DuckDB service.
Expand Down Expand Up @@ -605,7 +666,7 @@ import foreign schema "docs_example"
from server duckdb_server into duckdb;
```

Currently only MotherDuck and Iceberg-like servers, such as S3 Tables, R2 Data Catalog and etc., support `import foreign schema` without specifying source tables. For other types of servers, source tables must be explicitly specified in options. For example,
MotherDuck, DuckLake and Iceberg-like servers, such as S3 Tables and R2 Data Catalog, support `import foreign schema` without specifying source tables. For other types of servers, source tables must be explicitly specified in options. For example,

```sql
-- 'duckdb_server_md' server type is 'md', all tables under 'main' schema
Expand Down
28 changes: 27 additions & 1 deletion supabase-wrappers/src/utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,13 +97,37 @@ pub fn is_sensitive_option(option_name: &str) -> bool {
SENSITIVE_OPTION_NAMES.iter().any(|&s| lower.contains(s))
}

fn redact_postgres_urls(message: &str) -> String {
let mut result = String::with_capacity(message.len());
for token in message.split_inclusive(char::is_whitespace) {
// ASCII case folding keeps byte offsets valid for UTF-8 messages.
let lower = token.to_ascii_lowercase();
let start = ["postgres://", "postgresql://"]
.iter()
.filter_map(|scheme| lower.find(scheme))
.min();
if let Some(start) = start {
result.push_str(&token[..start]);
result.push_str("[REDACTED]");
// Conservatively hide the whole URL, including query credentials
// and trailing punctuation. Valid URLs encode embedded whitespace.
let end = token.trim_end_matches(char::is_whitespace).len();
result.push_str(&token[end..]);
} else {
result.push_str(token);
}
}
result
}

/// Masks credential values in an error message string.
/// Scans for patterns like `key = 'value'` or `key: value` and masks sensitive values.
///
/// This function handles multiple common formats:
/// - SQL-style: `secret = 'value'`
/// - JSON-style: `"secret": "value"`
/// - URL parameters: `secret=value`
/// - PostgreSQL connection URLs: the complete URL is hidden
///
/// # Examples
/// ```
Expand All @@ -114,7 +138,9 @@ pub fn is_sensitive_option(option_name: &str) -> bool {
/// assert!(masked.contains("wJal***"));
/// ```
pub fn mask_credentials_in_message(message: &str) -> String {
let mut result = message.to_string();
// Do this before named-field masking, which could otherwise modify a URL
// containing query parameters before its embedded credentials are removed.
let mut result = redact_postgres_urls(message);

for sensitive_name in SENSITIVE_OPTION_NAMES {
let lower_name = sensitive_name.to_lowercase();
Expand Down
16 changes: 15 additions & 1 deletion wrappers/.ci/docker-compose-native.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -272,4 +272,18 @@ services:
test: ["CMD", "sh", "-c", "echo 'db.runCommand({ ping: 1 })' | mongosh --quiet"]
interval: 10s
timeout: 5s
retries: 5
retries: 5

ducklake-postgres:
image: postgres:17
environment:
POSTGRES_USER: ducklake
POSTGRES_PASSWORD: ducklake
POSTGRES_DB: ducklake
ports:
- "5433:5432"
healthcheck:
test: ["CMD-SHELL", "pg_isready -U ducklake -d ducklake"]
interval: 5s
timeout: 5s
retries: 10
16 changes: 6 additions & 10 deletions wrappers/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -185,9 +185,9 @@ iceberg_fdw = [
"uuid",
]
duckdb_fdw = [
"arrow-array-compat",
"arrow-schema-compat",
"arrow-json-compat",
"arrow-array",
"arrow-schema",
"arrow-json",
"chrono",
"duckdb",
"regex",
Expand Down Expand Up @@ -356,13 +356,9 @@ iceberg-catalog-rest = { version = "0.8.0", optional = true }
rust_decimal = { version = "1.37.1", optional = true }

# for duckdb_fdw
duckdb = { version = "=1.3.2", features = ["bundled"], optional = true }
# Version-pinned arrow compat deps matching duckdb's bundled arrow (55.x).
# The workspace uses arrow 57.x for iceberg_fdw; using 57.x here causes type
# incompatibilities because duckdb::arrow::* types implement the 55.x traits.
arrow-array-compat = { package = "arrow-array", version = "55", optional = true }
arrow-json-compat = { package = "arrow-json", version = "55", optional = true }
arrow-schema-compat = { package = "arrow-schema", version = "55", optional = true }
# Bundles DuckDB 1.5.0 and shares Arrow 57 with iceberg_fdw and s3_fdw.
# Keep pinned: duckdb-rs 1.10501.0 moves to Arrow 58.
duckdb = { version = "=1.10500.0", features = ["bundled"], optional = true }

[dev-dependencies]
pgrx-tests = "=0.19.2"
11 changes: 6 additions & 5 deletions wrappers/src/fdw/duckdb_fdw/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,9 @@ This is a foreign data wrapper for [DuckDB](https://duckdb.org/). It is develope

## Changelog

| Version | Date | Notes |
| ------- | ---------- | ---------------------------------------------------- |
| 0.1.2 | 2025-10-16 | Add MotherDuck support |
| 0.1.1 | 2025-08-15 | Replace execute_batch() with execute() |
| 0.1.0 | 2024-10-31 | Initial version |
| Version | Date | Notes |
| ------- | ---------- | ---------------------------------------------------------- |
| 0.1.3 | 2026-09-29 | Upgrade DuckDB to 1.5.0 and add read-only DuckLake support |
| 0.1.2 | 2025-10-16 | Add MotherDuck support |
| 0.1.1 | 2025-08-15 | Replace execute_batch() with execute() |
| 0.1.0 | 2024-10-31 | Initial version |
28 changes: 13 additions & 15 deletions wrappers/src/fdw/duckdb_fdw/duckdb_fdw.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ use supabase_wrappers::prelude::*;
use super::{DuckdbFdwError, DuckdbFdwResult, mapper, server_type::ServerType};

#[wrappers_fdw(
version = "0.1.2",
version = "0.1.3",
author = "Supabase",
website = "https://github.com/supabase/wrappers/tree/main/wrappers/src/fdw/duckdb_fdw",
error_type = "DuckdbFdwError"
Expand All @@ -27,20 +27,18 @@ impl DuckdbFdw {
const FDW_NAME: &'static str = "DuckdbFdw";

fn init_duckdb(&self) -> DuckdbFdwResult<()> {
let sql_batch = String::default()
+ self.svr_type.get_duckdb_extension_sql()
+ &self.svr_type.get_settings_sql(&self.svr_opts)
+ &self.svr_type.get_create_secret_sql(&self.svr_opts)
+ &self.svr_type.get_attach_sql(&self.svr_opts)?;

// execute_batch() won't raise error when one of the statements failed,
// so we execute each sql separately
for sql in sql_batch
.split(";")
.map(|s| s.trim())
.filter(|s| !s.is_empty())
{
self.conn.execute(sql, [])?;
let statements = self
.svr_type
.get_duckdb_extension_sql()
.into_iter()
.chain(self.svr_type.get_settings_sql(&self.svr_opts))
.chain(self.svr_type.get_create_secret_sql(&self.svr_opts))
.chain([self.svr_type.get_attach_sql(&self.svr_opts)?]);

// Execute each statement separately to propagate errors, preserving
// semicolons inside quoted credentials and connection strings.
for sql in statements.filter(|sql| !sql.is_empty()) {
self.conn.execute(&sql, [])?;
}

Ok(())
Expand Down
4 changes: 2 additions & 2 deletions wrappers/src/fdw/duckdb_fdw/mapper.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
use arrow_array_compat::{Array, RecordBatch, array::ArrayRef};
use arrow_json_compat::ArrayWriter;
use arrow_array::{Array, RecordBatch, array::ArrayRef};
use arrow_json::ArrayWriter;
use duckdb::{
self,
types::{EnumType, ListType, ValueRef},
Expand Down
2 changes: 1 addition & 1 deletion wrappers/src/fdw/duckdb_fdw/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ enum DuckdbFdwError {
NumericError(#[from] pgrx::datum::numeric_support::error::Error),

#[error("arrow error: {0}")]
ArrowError(#[from] arrow_schema_compat::ArrowError),
ArrowError(#[from] arrow_schema::ArrowError),

#[error("uuid error: {0}")]
UuidConversionError(#[from] uuid::Error),
Expand Down
66 changes: 48 additions & 18 deletions wrappers/src/fdw/duckdb_fdw/server_type.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ pub(super) enum ServerType {

// SQL-like Remotes
MotherDuck,
DuckLake,
}

impl ServerType {
Expand All @@ -36,6 +37,7 @@ impl ServerType {
"polaris" => Self::Polaris,
"lakekeeper" => Self::Lakekeeper,
"md" => Self::MotherDuck,
"ducklake" => Self::DuckLake,
_ => return Err(DuckdbFdwError::InvalidServerType(svr_type.to_owned())),
};
Ok(ret)
Expand All @@ -51,6 +53,7 @@ impl ServerType {
Self::Polaris => "polaris",
Self::Lakekeeper => "lakekeeper",
Self::MotherDuck => "md",
Self::DuckLake => "ducklake",
}
}

Expand All @@ -62,17 +65,23 @@ impl ServerType {
}

pub(super) fn is_sql_like(&self) -> bool {
matches!(self, Self::MotherDuck)
matches!(self, Self::MotherDuck | Self::DuckLake)
}

pub(super) fn get_duckdb_extension_sql(&self) -> &'static str {
match self {
pub(super) fn get_duckdb_extension_sql(&self) -> Vec<String> {
let extensions = match self {
Self::Iceberg | Self::S3Tables | Self::R2Catalog | Self::Polaris | Self::Lakekeeper => {
"install iceberg;load iceberg;"
vec!["iceberg"]
}
Self::MotherDuck => "install md;load md;",
_ => "",
}
Self::MotherDuck => vec!["md"],
// Load dependencies before disabling local filesystem access.
Self::DuckLake => vec!["ducklake", "postgres", "httpfs", "parquet"],
_ => vec![],
};
extensions
.into_iter()
.flat_map(|ext| [format!("install {ext}"), format!("load {ext}")])
.collect()
}

fn allowed_secret_params(&self) -> Vec<&'static str> {
Expand Down Expand Up @@ -100,7 +109,7 @@ impl ServerType {
"oauth2_scope",
"oauth2_server_uri",
],
Self::MotherDuck => vec![],
Self::MotherDuck | Self::DuckLake => vec![],
}
}

Expand All @@ -126,10 +135,11 @@ impl ServerType {
}

// make 'create secret' sql for DuckDB from server options
pub(super) fn get_create_secret_sql(&self, svr_opts: &ServerOptions) -> String {
pub(super) fn get_create_secret_sql(&self, svr_opts: &ServerOptions) -> Vec<String> {
let secrets: Vec<(&str, Vec<&str>)> = match self {
Self::S3 | Self::S3Tables => vec![("s3", self.allowed_secret_params())],
Self::R2 => vec![("r2", self.allowed_secret_params())],
Self::DuckLake => vec![("s3", Self::S3.allowed_secret_params())],

// note: for generic Iceberg, we only support S3 compatible storage for now,
// so we need to create 2 secrets: one for S3 and one for Iceberg
Expand All @@ -147,16 +157,19 @@ impl ServerType {
_ => vec![],
};

let mut ret = String::default();
let mut ret = Vec::new();
for (typ, params) in secrets {
let opts = self.format_options(svr_opts, &params);
ret.push_str(&format!("create or replace secret (type {typ}, {opts});"));
// Public storage and catalogs with inlined data need no S3 secret.
if !opts.is_empty() {
ret.push(format!("create or replace secret (type {typ}, {opts})"));
}
}

ret
}

pub(super) fn get_settings_sql(&self, svr_opts: &ServerOptions) -> String {
pub(super) fn get_settings_sql(&self, svr_opts: &ServerOptions) -> Vec<String> {
let settings: Vec<(&str, String)> = match self {
Self::MotherDuck => {
let token = if svr_opts.contains_key("vault_motherduck_token") {
Expand Down Expand Up @@ -186,16 +199,33 @@ impl ServerType {
]
}
};
let mut ret = String::default();
for (key, value) in settings {
ret.push_str(&format!("set {key}={value};"));
}

ret
settings
.into_iter()
.map(|(key, value)| format!("set {key}={value}"))
.collect()
}

pub(super) fn get_attach_sql(&self, svr_opts: &ServerOptions) -> DuckdbFdwResult<String> {
let ret = match self {
Self::DuckLake => {
let metadata_path = if let Some(secret_id) = svr_opts.get("vault_metadata_path") {
get_vault_secret(secret_id).ok_or_else(|| {
OptionsError::OptionNameNotFound("vault_metadata_path".to_string())
})?
} else {
require_option("metadata_path", svr_opts)?.to_string()
};
let opts = self.format_options(svr_opts, &["data_path", "metadata_schema"]);
let opts = if opts.is_empty() {
String::new()
} else {
format!(", {opts}")
};
format!(
"attach 'ducklake:{}' as ducklake (read_only, create_if_not_exists false{opts})",
metadata_path.replace("'", "''")
)
}
Self::S3Tables => {
let arn = require_option("s3_tables_arn", svr_opts)?;
let db_name = self.as_str();
Expand Down
Loading