use crate::Result;
use crate::error::{TypeChangeOperation, TypeChangeWriteMode};
use crate::metadata_writer::{
ColumnDef, ColumnStat, CommitIds, DataFileInfo, ExistingCatalogColumn, MetadataWriter,
SnapshotCommitMetadata, WriteMode, WriteSetupResult, assign_column_ids, catalog_column_defs,
catalog_column_type_equal, catalog_column_type_requires_migration, catalog_columns_differ,
quote_snapshot_name, quote_snapshot_table, table_write_changes, top_level_column_ids,
validate_name,
};
use crate::partition::PartitionTransform;
use duckdb::{Connection, OptionalExt, Transaction, params};
use std::sync::{Arc, Mutex, MutexGuard};
const SQL_CREATE_SCHEMA: &str = r#"
CREATE SEQUENCE IF NOT EXISTS ducklake_schema_id_seq START 1;
CREATE SEQUENCE IF NOT EXISTS ducklake_table_id_seq START 1;
CREATE SEQUENCE IF NOT EXISTS ducklake_data_file_id_seq START 1;
CREATE SEQUENCE IF NOT EXISTS ducklake_delete_file_id_seq START 1;
CREATE SEQUENCE IF NOT EXISTS ducklake_partition_id_seq START 1;
CREATE SEQUENCE IF NOT EXISTS ducklake_sort_id_seq START 1;
CREATE TABLE IF NOT EXISTS ducklake_metadata (
key VARCHAR NOT NULL,
value VARCHAR NOT NULL,
scope VARCHAR
);
CREATE TABLE IF NOT EXISTS ducklake_snapshot (
snapshot_id BIGINT PRIMARY KEY,
snapshot_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
schema_version BIGINT NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS ducklake_snapshot_changes (
snapshot_id BIGINT PRIMARY KEY,
changes_made VARCHAR,
author VARCHAR,
commit_message VARCHAR,
commit_extra_info VARCHAR
);
CREATE TABLE IF NOT EXISTS ducklake_schema_versions (
begin_snapshot BIGINT NOT NULL,
schema_version BIGINT NOT NULL,
table_id BIGINT NOT NULL,
UNIQUE (table_id, begin_snapshot)
);
CREATE TABLE IF NOT EXISTS ducklake_schema (
schema_id BIGINT PRIMARY KEY DEFAULT nextval('ducklake_schema_id_seq'),
schema_name VARCHAR NOT NULL,
path VARCHAR NOT NULL DEFAULT '',
path_is_relative BOOLEAN NOT NULL DEFAULT true,
begin_snapshot BIGINT NOT NULL,
end_snapshot BIGINT
);
CREATE TABLE IF NOT EXISTS ducklake_table (
table_id BIGINT PRIMARY KEY DEFAULT nextval('ducklake_table_id_seq'),
schema_id BIGINT NOT NULL,
table_name VARCHAR NOT NULL,
path VARCHAR NOT NULL DEFAULT '',
path_is_relative BOOLEAN NOT NULL DEFAULT true,
begin_snapshot BIGINT NOT NULL,
end_snapshot BIGINT
);
-- Bare table (no PK, no NOT NULL), upstream's exact column order, so a column
-- can be versioned by [begin_snapshot, end_snapshot) and a type promotion can
-- write a second row sharing the same column_id. Mirrors the SQLite writer's
-- `ducklake_column`. The four `*default*` columns + `parent_column` are left
-- NULL (no nested-type / column-default writes here).
CREATE TABLE IF NOT EXISTS ducklake_column (
column_id BIGINT,
begin_snapshot BIGINT,
end_snapshot BIGINT,
table_id BIGINT,
column_order BIGINT,
column_name VARCHAR,
column_type VARCHAR,
initial_default VARCHAR,
default_value VARCHAR,
nulls_allowed BOOLEAN,
parent_column BIGINT,
default_value_type VARCHAR,
default_value_dialect VARCHAR
);
CREATE TABLE IF NOT EXISTS ducklake_data_file (
data_file_id BIGINT PRIMARY KEY DEFAULT nextval('ducklake_data_file_id_seq'),
table_id BIGINT NOT NULL,
path VARCHAR NOT NULL,
path_is_relative BOOLEAN NOT NULL DEFAULT true,
file_size_bytes BIGINT NOT NULL,
footer_size BIGINT,
encryption_key VARCHAR,
record_count BIGINT,
row_id_start BIGINT,
mapping_id BIGINT,
begin_snapshot BIGINT NOT NULL,
end_snapshot BIGINT,
-- References ducklake_partition_info.partition_id; NULL when unpartitioned.
partition_id BIGINT
);
CREATE TABLE IF NOT EXISTS ducklake_table_stats (
table_id BIGINT PRIMARY KEY,
record_count BIGINT NOT NULL DEFAULT 0,
next_row_id BIGINT NOT NULL DEFAULT 0,
file_size_bytes BIGINT NOT NULL DEFAULT 0
);
-- Per-file, per-column zone maps (DuckLake spec) — powers file pruning.
-- Column set mirrors the official extension and the other backends.
CREATE TABLE IF NOT EXISTS ducklake_file_column_stats (
data_file_id BIGINT NOT NULL,
table_id BIGINT NOT NULL,
column_id BIGINT NOT NULL,
column_size_bytes BIGINT,
value_count BIGINT,
null_count BIGINT,
min_value VARCHAR,
max_value VARCHAR,
contains_nan BOOLEAN,
extra_stats VARCHAR
);
-- Table-wide per-column roll-up (DuckLake spec) — feeds the optimizer.
CREATE TABLE IF NOT EXISTS ducklake_table_column_stats (
table_id BIGINT NOT NULL,
column_id BIGINT NOT NULL,
contains_null BOOLEAN,
contains_nan BOOLEAN,
min_value VARCHAR,
max_value VARCHAR,
extra_stats VARCHAR
);
CREATE TABLE IF NOT EXISTS ducklake_delete_file (
delete_file_id BIGINT PRIMARY KEY DEFAULT nextval('ducklake_delete_file_id_seq'),
data_file_id BIGINT NOT NULL,
table_id BIGINT NOT NULL,
path VARCHAR NOT NULL,
path_is_relative BOOLEAN NOT NULL DEFAULT true,
file_size_bytes BIGINT NOT NULL,
footer_size BIGINT,
encryption_key VARCHAR,
delete_count BIGINT,
begin_snapshot BIGINT NOT NULL,
end_snapshot BIGINT
);
CREATE TABLE IF NOT EXISTS ducklake_files_scheduled_for_deletion (
data_file_id BIGINT NOT NULL,
path VARCHAR NOT NULL,
path_is_relative BOOLEAN NOT NULL DEFAULT true,
schedule_start TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- Partition spec generations (DuckLake spec); end_snapshot NULL == active.
CREATE TABLE IF NOT EXISTS ducklake_partition_info (
partition_id BIGINT NOT NULL,
table_id BIGINT NOT NULL,
begin_snapshot BIGINT NOT NULL,
end_snapshot BIGINT
);
-- Partition-key columns for a spec (DuckLake spec), ordered by partition_key_index.
CREATE TABLE IF NOT EXISTS ducklake_partition_column (
partition_id BIGINT NOT NULL,
table_id BIGINT NOT NULL,
partition_key_index BIGINT NOT NULL,
column_id BIGINT NOT NULL,
transform VARCHAR NOT NULL
);
-- Per-file partition values (DuckLake spec): the value every row in the file
-- shares for a partition key, DuckDB-canonical VARCHAR (NULL is legal).
CREATE TABLE IF NOT EXISTS ducklake_file_partition_value (
data_file_id BIGINT NOT NULL,
table_id BIGINT NOT NULL,
partition_key_index BIGINT NOT NULL,
partition_value VARCHAR
);
-- Sort spec generations (DuckLake spec); end_snapshot NULL == active. sort_id is
-- allocated from the next_sort_id counter (like partition_id).
CREATE TABLE IF NOT EXISTS ducklake_sort_info (
sort_id BIGINT NOT NULL,
table_id BIGINT NOT NULL,
begin_snapshot BIGINT NOT NULL,
end_snapshot BIGINT
);
-- Sort-key expressions for a spec (DuckLake spec), ordered by sort_key_index.
-- expression is a sort expression in `dialect` (this crate produces bare column
-- names under `duckdb`); sort_direction ASC/DESC; null_order NULLS_FIRST/LAST.
CREATE TABLE IF NOT EXISTS ducklake_sort_expression (
sort_id BIGINT NOT NULL,
table_id BIGINT NOT NULL,
sort_key_index BIGINT NOT NULL,
expression VARCHAR NOT NULL,
dialect VARCHAR NOT NULL,
sort_direction VARCHAR NOT NULL,
null_order VARCHAR NOT NULL
);
"#;
#[derive(Debug, Clone)]
pub struct DuckdbMetadataWriter {
conn: Arc<Mutex<Connection>>,
#[allow(dead_code)]
catalog_path: String,
}
impl DuckdbMetadataWriter {
pub fn new(path: impl Into<String>) -> Result<Self> {
let catalog_path = path.into();
let conn = Connection::open(&catalog_path)?;
Ok(Self {
conn: Arc::new(Mutex::new(conn)),
catalog_path,
})
}
pub fn new_with_init(path: impl Into<String>) -> Result<Self> {
let writer = Self::new(path)?;
writer.initialize_schema()?;
Ok(writer)
}
fn connection(&self) -> MutexGuard<'_, Connection> {
self.conn.lock().expect("DuckDB connection mutex poisoned")
}
}
fn reserve_ids(tx: &Transaction<'_>, key: &str, n: i64) -> Result<i64> {
let last: i64 = tx.query_row(
"UPDATE ducklake_metadata
SET value = CAST(CAST(value AS BIGINT) + ? AS VARCHAR)
WHERE key = ? AND scope IS NULL
RETURNING CAST(value AS BIGINT)",
params![n, key],
|row| row.get(0),
)?;
Ok(last)
}
fn insert_snapshot(tx: &Transaction<'_>) -> Result<(i64, i64)> {
let (snapshot_id, schema_version): (i64, i64) = tx.query_row(
"SELECT COALESCE(MAX(snapshot_id), 0) + 1, COALESCE(MAX(schema_version), 0)
FROM ducklake_snapshot",
[],
|row| Ok((row.get(0)?, row.get(1)?)),
)?;
tx.execute(
"INSERT INTO ducklake_snapshot (snapshot_id, snapshot_time, schema_version)
VALUES (?, CURRENT_TIMESTAMP, ?)",
params![snapshot_id, schema_version],
)?;
tx.execute(
"INSERT INTO ducklake_snapshot_changes (snapshot_id, changes_made)
VALUES (?, NULL)",
params![snapshot_id],
)?;
Ok((snapshot_id, schema_version))
}
fn record_snapshot_changes(
tx: &Transaction<'_>,
snapshot_id: i64,
changes_made: &str,
commit_metadata: &SnapshotCommitMetadata,
) -> Result<()> {
let changes_made = (!changes_made.is_empty()).then_some(changes_made);
tx.execute(
"UPDATE ducklake_snapshot_changes
SET changes_made = CASE
WHEN changes_made IS NULL THEN ?
WHEN ? IS NULL THEN changes_made
ELSE changes_made || ',' || ?
END,
author = ?,
commit_message = ?,
commit_extra_info = ?
WHERE snapshot_id = ?",
params![
changes_made,
changes_made,
changes_made,
commit_metadata.author(),
commit_metadata.message(),
commit_metadata.extra_info(),
snapshot_id,
],
)?;
Ok(())
}
fn record_table_write_changes(
tx: &Transaction<'_>,
snapshot_id: i64,
table_id: i64,
schema_name: &str,
table_name: &str,
mode: WriteMode,
commit_metadata: &SnapshotCommitMetadata,
) -> Result<()> {
let (schema_begin_snapshot, table_begin_snapshot): (i64, i64) = tx.query_row(
"SELECT s.begin_snapshot, t.begin_snapshot
FROM ducklake_table t
JOIN ducklake_schema s ON s.schema_id = t.schema_id
WHERE t.table_id = ?",
params![table_id],
|row| Ok((row.get(0)?, row.get(1)?)),
)?;
let altered: bool = tx.query_row(
"SELECT EXISTS(
SELECT 1 FROM ducklake_schema_versions
WHERE table_id = ? AND begin_snapshot = ?
)",
params![table_id, snapshot_id],
|row| row.get(0),
)?;
let replaced_existing_data: bool = tx.query_row(
"SELECT EXISTS(
SELECT 1 FROM ducklake_data_file
WHERE table_id = ? AND end_snapshot = ?
)",
params![table_id, snapshot_id],
|row| row.get(0),
)?;
let mut changes = Vec::new();
if schema_begin_snapshot == snapshot_id {
changes.push(format!(
"created_schema:{}",
quote_snapshot_name(schema_name)
));
}
if table_begin_snapshot == snapshot_id {
changes.push(format!(
"created_table:{}",
quote_snapshot_table(schema_name, table_name)
));
} else if altered {
changes.push(format!("altered_table:{table_id}"));
}
changes.push(table_write_changes(
table_id,
mode,
false,
replaced_existing_data,
));
record_snapshot_changes(tx, snapshot_id, &changes.join(","), commit_metadata)
}
fn bump_schema_version(tx: &Transaction<'_>, snapshot_id: i64) -> Result<i64> {
let prev_max: i64 = tx.query_row(
"SELECT COALESCE(MAX(schema_version), 0) FROM ducklake_snapshot WHERE snapshot_id <> ?",
params![snapshot_id],
|row| row.get(0),
)?;
let new_version = prev_max + 1;
tx.execute(
"UPDATE ducklake_snapshot SET schema_version = ? WHERE snapshot_id = ?",
params![new_version, snapshot_id],
)?;
Ok(new_version)
}
fn record_schema_version(
tx: &Transaction<'_>,
snapshot_id: i64,
schema_version: i64,
table_id: i64,
) -> Result<()> {
tx.execute(
"INSERT INTO ducklake_schema_versions (begin_snapshot, schema_version, table_id)
VALUES (?, ?, ?)",
params![snapshot_id, schema_version, table_id],
)?;
Ok(())
}
fn seed_stats_if_missing(tx: &Transaction<'_>, table_id: i64) -> Result<()> {
let exists: Option<i64> = tx
.query_row(
"SELECT 1 FROM ducklake_table_stats WHERE table_id = ?",
params![table_id],
|row| row.get(0),
)
.optional()?;
if exists.is_none() {
tx.execute(
"INSERT INTO ducklake_table_stats
(table_id, record_count, next_row_id, file_size_bytes)
VALUES (?, 0, 0, 0)",
params![table_id],
)?;
}
Ok(())
}
fn detect_replace_conflict(tx: &Transaction<'_>, table_id: i64, base_snapshot: i64) -> Result<()> {
let conflicts: i64 = tx.query_row(
"SELECT COUNT(*) FROM ducklake_data_file
WHERE table_id = ? AND (begin_snapshot > ? OR end_snapshot > ?)",
params![table_id, base_snapshot, base_snapshot],
|row| row.get(0),
)?;
if conflicts > 0 {
return Err(crate::DuckLakeError::Conflict(format!(
"Replace on table {table_id} conflicts with a concurrent write committed since \
snapshot {base_snapshot}; aborting (retry the write against the new generation)"
)));
}
Ok(())
}
fn retire_prior_generation(tx: &Transaction<'_>, table_id: i64, snapshot_id: i64) -> Result<()> {
tx.execute(
"UPDATE ducklake_data_file SET end_snapshot = ?
WHERE table_id = ? AND end_snapshot IS NULL AND begin_snapshot < ?",
params![snapshot_id, table_id, snapshot_id],
)?;
tx.execute(
"UPDATE ducklake_table_stats SET record_count = 0, file_size_bytes = 0 WHERE table_id = ?",
params![table_id],
)?;
Ok(())
}
fn insert_file_column_stats(
tx: &Transaction<'_>,
table_id: i64,
data_file_id: i64,
column_stats: &[ColumnStat],
) -> Result<()> {
for stat in column_stats {
tx.execute(
"INSERT INTO ducklake_file_column_stats
(data_file_id, table_id, column_id, column_size_bytes,
value_count, null_count, min_value, max_value, contains_nan, extra_stats)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, NULL)",
params![
data_file_id,
table_id,
stat.column_id,
stat.column_size_bytes,
stat.value_count,
stat.null_count,
stat.min_value.as_deref(),
stat.max_value.as_deref(),
stat.contains_nan,
],
)?;
}
Ok(())
}
fn insert_partition_metadata(
tx: &Transaction<'_>,
table_id: i64,
data_file_id: i64,
file: &DataFileInfo,
) -> Result<()> {
if let Some(partition_id) = file.partition_id {
tx.execute(
"UPDATE ducklake_data_file SET partition_id = ? WHERE data_file_id = ?",
params![partition_id, data_file_id],
)?;
}
for (key_index, value) in &file.partition_values {
tx.execute(
"INSERT INTO ducklake_file_partition_value
(data_file_id, table_id, partition_key_index, partition_value)
VALUES (?, ?, ?, ?)",
params![data_file_id, table_id, i64::from(*key_index), value.as_deref()],
)?;
}
Ok(())
}
fn recompute_table_column_stats(
tx: &Transaction<'_>,
table_id: i64,
columns: &[ColumnDef],
column_ids: &[i64],
) -> Result<()> {
use crate::stats_encode::{FileColumnStat, aggregate_global_column_stats};
let catalog_columns = catalog_column_defs(columns)?;
let column_ids = top_level_column_ids(&catalog_columns, column_ids)?;
let live_file_count: i64 = tx.query_row(
"SELECT COUNT(*) FROM ducklake_data_file WHERE table_id = ? AND end_snapshot IS NULL",
params![table_id],
|row| row.get(0),
)?;
let per_file: Vec<FileColumnStat> = {
let mut stmt = tx.prepare(
"SELECT s.column_id, s.min_value, s.max_value, s.null_count, s.contains_nan
FROM ducklake_file_column_stats s
JOIN ducklake_data_file d ON d.data_file_id = s.data_file_id
WHERE d.table_id = ? AND d.end_snapshot IS NULL",
)?;
let mapped = stmt.query_map(params![table_id], |row| {
Ok(FileColumnStat {
column_id: row.get::<_, i64>(0)?,
min_value: row.get::<_, Option<String>>(1)?,
max_value: row.get::<_, Option<String>>(2)?,
null_count: row.get::<_, Option<i64>>(3)?,
contains_nan: row.get::<_, Option<bool>>(4)?,
})
})?;
mapped.collect::<duckdb::Result<Vec<_>>>()?
};
let numeric_of = |column_id: i64| -> bool {
column_ids
.iter()
.position(|id| *id == column_id)
.and_then(|i| columns.get(i))
.map(|c| crate::stats_encode::is_numeric_ducklake_type(c.ducklake_type()))
.unwrap_or(false)
};
let globals = aggregate_global_column_stats(&per_file, live_file_count, numeric_of);
tx.execute(
"DELETE FROM ducklake_table_column_stats WHERE table_id = ?",
params![table_id],
)?;
for g in globals {
tx.execute(
"INSERT INTO ducklake_table_column_stats
(table_id, column_id, contains_null, contains_nan, min_value, max_value, extra_stats)
VALUES (?, ?, ?, ?, ?, ?, NULL)",
params![
table_id,
g.column_id,
g.contains_null,
g.contains_nan,
g.min_value,
g.max_value,
],
)?;
}
Ok(())
}
fn finalize_snapshot(
tx: &Transaction<'_>,
table_id: i64,
columns: &[ColumnDef],
column_ids: &[i64],
mode: WriteMode,
base_snapshot: i64,
) -> Result<i64> {
use std::collections::{HashMap, HashSet};
let proposed = catalog_column_defs(columns)?;
if proposed.len() != column_ids.len() {
return Err(crate::DuckLakeError::InvalidConfig(format!(
"column_ids has {} entries for {} catalog column nodes",
column_ids.len(),
proposed.len()
)));
}
let (snapshot_id, mut schema_version) = insert_snapshot(tx)?;
let current: Vec<(i64, String, String, i64, bool, Option<i64>)> = {
let mut stmt = tx.prepare(
"SELECT column_id, column_name, column_type, column_order, nulls_allowed, parent_column
FROM ducklake_column
WHERE table_id = ? AND end_snapshot IS NULL
ORDER BY column_order",
)?;
let mapped = stmt.query_map(params![table_id], |row| {
let column_id: i64 = row.get(0)?;
let name: String = row.get(1)?;
let ty: String = row.get(2)?;
let order: i64 = row.get(3)?;
let nullable: Option<bool> = row.get(4)?;
let parent_column: Option<i64> = row.get(5)?;
Ok((
column_id,
name,
ty,
order,
nullable.unwrap_or(true),
parent_column,
))
})?;
mapped.collect::<std::result::Result<Vec<_>, duckdb::Error>>()?
};
let existing_catalog_columns = current
.iter()
.map(|column| ExistingCatalogColumn {
column_id: column.0,
name: column.1.clone(),
ducklake_type: column.2.clone(),
parent_column: column.5,
})
.collect::<Vec<_>>();
let existing_nullability = current.iter().map(|column| column.4).collect::<Vec<_>>();
let committed_ids = assign_column_ids(&proposed, &existing_catalog_columns, column_ids)?;
if committed_ids != column_ids {
return Err(crate::DuckLakeError::Conflict(
"table columns were created concurrently with different field ids; retry the write"
.to_string(),
));
}
let is_ddl = current.is_empty()
|| catalog_columns_differ(
&existing_catalog_columns,
&existing_nullability,
&proposed,
column_ids,
);
if is_ddl {
schema_version = bump_schema_version(tx, snapshot_id)?;
}
let proposed_ids = column_ids.iter().copied().collect::<HashSet<_>>();
let mut current_by_id: HashMap<i64, (i64, bool, String)> = HashMap::new();
for (column_id, _name, ty, order, nullable, _parent_id) in ¤t {
if !proposed_ids.contains(column_id) {
tx.execute(
"UPDATE ducklake_column SET end_snapshot = ?
WHERE table_id = ? AND column_id = ? AND end_snapshot IS NULL",
params![snapshot_id, table_id, column_id],
)?;
}
current_by_id.insert(*column_id, (*order, *nullable, ty.clone()));
}
for (order, (column, column_id)) in proposed.iter().zip(column_ids).enumerate() {
let parent_id = column.parent_index.map(|index| column_ids[index]);
match current_by_id.get(column_id) {
Some((cur_order, cur_nullable, cur_type)) => {
let migrate_type = catalog_column_type_requires_migration(cur_type, column);
if migrate_type {
tx.execute(
"UPDATE ducklake_column SET end_snapshot = ?
WHERE table_id = ? AND column_id = ? AND end_snapshot IS NULL",
params![snapshot_id, table_id, column_id],
)?;
tx.execute(
"INSERT INTO ducklake_column
(column_id, table_id, column_name, column_type, column_order,
nulls_allowed, parent_column, begin_snapshot)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
params![
column_id,
table_id,
column.name.as_str(),
column.ducklake_type.as_str(),
order as i64,
column.is_nullable,
parent_id,
snapshot_id
],
)?;
} else if *cur_order != order as i64 || *cur_nullable != column.is_nullable {
tx.execute(
"UPDATE ducklake_column
SET column_order = ?, nulls_allowed = ?
WHERE table_id = ? AND column_id = ? AND end_snapshot IS NULL",
params![order as i64, column.is_nullable, table_id, column_id],
)?;
}
},
None => {
tx.execute(
"INSERT INTO ducklake_column
(column_id, table_id, column_name, column_type, column_order,
nulls_allowed, parent_column, begin_snapshot)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
params![
column_id,
table_id,
column.name.as_str(),
column.ducklake_type.as_str(),
order as i64,
column.is_nullable,
parent_id,
snapshot_id
],
)?;
},
}
}
if mode == WriteMode::Replace {
detect_replace_conflict(tx, table_id, base_snapshot)?;
seed_stats_if_missing(tx, table_id)?;
retire_prior_generation(tx, table_id, snapshot_id)?;
}
if is_ddl {
record_schema_version(tx, snapshot_id, schema_version, table_id)?;
}
Ok(snapshot_id)
}
impl MetadataWriter for DuckdbMetadataWriter {
fn create_snapshot(&self) -> Result<i64> {
let mut conn = self.connection();
let tx = conn.transaction()?;
let (snapshot_id, _schema_version) = insert_snapshot(&tx)?;
tx.commit()?;
Ok(snapshot_id)
}
fn get_or_create_schema(
&self,
name: &str,
path: Option<&str>,
snapshot_id: i64,
) -> Result<(i64, bool)> {
validate_name(name, "Schema")?;
let mut conn = self.connection();
let tx = conn.transaction()?;
let existing: Option<i64> = tx
.query_row(
"SELECT schema_id FROM ducklake_schema
WHERE schema_name = ? AND end_snapshot IS NULL",
params![name],
|row| row.get(0),
)
.optional()?;
if let Some(schema_id) = existing {
tx.commit()?;
return Ok((schema_id, false));
}
let schema_path = path.unwrap_or(name);
let schema_id: i64 = tx.query_row(
"INSERT INTO ducklake_schema (schema_name, path, path_is_relative, begin_snapshot)
VALUES (?, ?, true, ?) RETURNING schema_id",
params![name, schema_path, snapshot_id],
|row| row.get(0),
)?;
record_snapshot_changes(
&tx,
snapshot_id,
&format!("created_schema:{}", quote_snapshot_name(name)),
&SnapshotCommitMetadata::default(),
)?;
tx.commit()?;
Ok((schema_id, true))
}
fn get_or_create_table(
&self,
schema_id: i64,
name: &str,
path: Option<&str>,
snapshot_id: i64,
) -> Result<(i64, bool)> {
validate_name(name, "Table")?;
let mut conn = self.connection();
let tx = conn.transaction()?;
let existing: Option<i64> = tx
.query_row(
"SELECT table_id FROM ducklake_table
WHERE schema_id = ? AND table_name = ? AND end_snapshot IS NULL",
params![schema_id, name],
|row| row.get(0),
)
.optional()?;
if let Some(table_id) = existing {
tx.commit()?;
return Ok((table_id, false));
}
let schema_name: String = tx.query_row(
"SELECT schema_name FROM ducklake_schema WHERE schema_id = ?",
params![schema_id],
|row| row.get(0),
)?;
let table_path = path.unwrap_or(name);
let table_id: i64 = tx.query_row(
"INSERT INTO ducklake_table (schema_id, table_name, path, path_is_relative, begin_snapshot)
VALUES (?, ?, ?, true, ?) RETURNING table_id",
params![schema_id, name, table_path, snapshot_id],
|row| row.get(0),
)?;
record_snapshot_changes(
&tx,
snapshot_id,
&format!("created_table:{}", quote_snapshot_table(&schema_name, name)),
&SnapshotCommitMetadata::default(),
)?;
tx.commit()?;
Ok((table_id, true))
}
fn set_columns(
&self,
table_id: i64,
columns: &[ColumnDef],
snapshot_id: i64,
) -> Result<Vec<i64>> {
if columns.is_empty() {
return Err(crate::DuckLakeError::InvalidConfig(
"Table must have at least one column".to_string(),
));
}
let mut conn = self.connection();
let tx = conn.transaction()?;
tx.execute(
"UPDATE ducklake_column SET end_snapshot = ?
WHERE table_id = ? AND end_snapshot IS NULL",
params![snapshot_id, table_id],
)?;
let catalog_columns = catalog_column_defs(columns)?;
let n = catalog_columns.len() as i64;
let last_column_id = reserve_ids(&tx, "next_column_id", n)?;
let first_column_id = last_column_id - n + 1;
let column_ids = (first_column_id..=last_column_id).collect::<Vec<_>>();
for (order, (column, column_id)) in
catalog_columns.iter().zip(column_ids.iter()).enumerate()
{
let parent_id = column.parent_index.map(|index| column_ids[index]);
tx.execute(
"INSERT INTO ducklake_column
(column_id, table_id, column_name, column_type, column_order,
nulls_allowed, parent_column, begin_snapshot)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
params![
column_id,
table_id,
column.name.as_str(),
column.ducklake_type.as_str(),
order as i64,
column.is_nullable,
parent_id,
snapshot_id
],
)?;
}
let table_begin_snapshot: i64 = tx.query_row(
"SELECT begin_snapshot FROM ducklake_table WHERE table_id = ?",
params![table_id],
|row| row.get(0),
)?;
if table_begin_snapshot != snapshot_id {
record_snapshot_changes(
&tx,
snapshot_id,
&format!("altered_table:{table_id}"),
&SnapshotCommitMetadata::default(),
)?;
}
tx.commit()?;
top_level_column_ids(&catalog_columns, &column_ids)
}
fn register_data_file(
&self,
table_id: i64,
schema_name: &str,
table_name: &str,
snapshot_id: i64,
file: &DataFileInfo,
mode: WriteMode,
base_snapshot: i64,
columns: &[ColumnDef],
column_ids: &[i64],
) -> Result<CommitIds> {
self.register_data_file_with_commit_metadata(
table_id,
schema_name,
table_name,
snapshot_id,
file,
mode,
base_snapshot,
columns,
column_ids,
&SnapshotCommitMetadata::default(),
None,
)
}
fn register_data_file_with_commit_metadata(
&self,
table_id: i64,
schema_name: &str,
table_name: &str,
_snapshot_id: i64,
file: &DataFileInfo,
mode: WriteMode,
base_snapshot: i64,
columns: &[ColumnDef],
column_ids: &[i64],
commit_metadata: &SnapshotCommitMetadata,
expected_base_snapshot_id: Option<i64>,
) -> Result<CommitIds> {
if expected_base_snapshot_id.is_some() {
return Err(crate::DuckLakeError::InvalidConfig(
"conditional writes are not supported by the DuckDB metadata writer".to_string(),
));
}
let mut conn = self.connection();
let tx = conn.transaction()?;
let snapshot_id =
finalize_snapshot(&tx, table_id, columns, column_ids, mode, base_snapshot)?;
let live_partition_id: Option<i64> = tx.query_row(
"SELECT (SELECT partition_id FROM ducklake_partition_info
WHERE table_id = ? AND end_snapshot IS NULL LIMIT 1)",
params![table_id],
|row| row.get::<_, Option<i64>>(0),
)?;
crate::metadata_writer::enforce_partition_fence(table_id, live_partition_id, file)?;
seed_stats_if_missing(&tx, table_id)?;
let row_id_start: i64 = tx.query_row(
"SELECT next_row_id FROM ducklake_table_stats WHERE table_id = ?",
params![table_id],
|row| row.get(0),
)?;
let data_file_id: i64 = tx.query_row(
"INSERT INTO ducklake_data_file
(table_id, path, path_is_relative, file_size_bytes,
footer_size, record_count, row_id_start, begin_snapshot)
VALUES (?, ?, ?, ?, ?, ?, ?, ?) RETURNING data_file_id",
params![
table_id,
file.path.as_str(),
file.path_is_relative,
file.file_size_bytes,
file.footer_size,
file.record_count,
row_id_start,
snapshot_id
],
|row| row.get(0),
)?;
insert_file_column_stats(&tx, table_id, data_file_id, &file.column_stats)?;
insert_partition_metadata(&tx, table_id, data_file_id, file)?;
recompute_table_column_stats(&tx, table_id, columns, column_ids)?;
tx.execute(
"UPDATE ducklake_table_stats
SET next_row_id = next_row_id + ?,
record_count = record_count + ?,
file_size_bytes = file_size_bytes + ?
WHERE table_id = ?",
params![file.record_count, file.record_count, file.file_size_bytes, table_id],
)?;
record_table_write_changes(
&tx,
snapshot_id,
table_id,
schema_name,
table_name,
mode,
commit_metadata,
)?;
let schema_id: i64 = tx.query_row(
"SELECT schema_id FROM ducklake_table WHERE table_id = ?",
params![table_id],
|row| row.get(0),
)?;
tx.commit()?;
Ok(CommitIds {
snapshot_id,
schema_id,
table_id,
})
}
#[allow(clippy::too_many_arguments)]
fn register_data_files(
&self,
table_id: i64,
schema_name: &str,
table_name: &str,
snapshot_id: i64,
files: &[DataFileInfo],
mode: WriteMode,
base_snapshot: i64,
columns: &[ColumnDef],
column_ids: &[i64],
) -> Result<CommitIds> {
self.register_data_files_with_commit_metadata(
table_id,
schema_name,
table_name,
snapshot_id,
files,
mode,
base_snapshot,
columns,
column_ids,
&SnapshotCommitMetadata::default(),
None,
)
}
fn register_data_files_with_commit_metadata(
&self,
table_id: i64,
schema_name: &str,
table_name: &str,
_snapshot_id: i64,
files: &[DataFileInfo],
mode: WriteMode,
base_snapshot: i64,
columns: &[ColumnDef],
column_ids: &[i64],
commit_metadata: &SnapshotCommitMetadata,
expected_base_snapshot_id: Option<i64>,
) -> Result<CommitIds> {
if expected_base_snapshot_id.is_some() {
return Err(crate::DuckLakeError::InvalidConfig(
"conditional multi-file writes are not supported by the DuckDB metadata writer"
.to_string(),
));
}
if files.is_empty() {
return Err(crate::DuckLakeError::InvalidConfig(
"register_data_files: files must be non-empty".to_string(),
));
}
let mut conn = self.connection();
let tx = conn.transaction()?;
let snapshot_id =
finalize_snapshot(&tx, table_id, columns, column_ids, mode, base_snapshot)?;
let live_partition_id: Option<i64> = tx.query_row(
"SELECT (SELECT partition_id FROM ducklake_partition_info
WHERE table_id = ? AND end_snapshot IS NULL LIMIT 1)",
params![table_id],
|row| row.get::<_, Option<i64>>(0),
)?;
for file in files {
crate::metadata_writer::enforce_partition_fence(table_id, live_partition_id, file)?;
}
seed_stats_if_missing(&tx, table_id)?;
let mut next_row_id: i64 = tx.query_row(
"SELECT next_row_id FROM ducklake_table_stats WHERE table_id = ?",
params![table_id],
|row| row.get(0),
)?;
let mut total_records: i64 = 0;
let mut total_bytes: i64 = 0;
for file in files {
let data_file_id: i64 = tx.query_row(
"INSERT INTO ducklake_data_file
(table_id, path, path_is_relative, file_size_bytes,
footer_size, record_count, row_id_start, begin_snapshot)
VALUES (?, ?, ?, ?, ?, ?, ?, ?) RETURNING data_file_id",
params![
table_id,
file.path.as_str(),
file.path_is_relative,
file.file_size_bytes,
file.footer_size,
file.record_count,
next_row_id,
snapshot_id
],
|row| row.get(0),
)?;
insert_file_column_stats(&tx, table_id, data_file_id, &file.column_stats)?;
insert_partition_metadata(&tx, table_id, data_file_id, file)?;
next_row_id += file.record_count;
total_records += file.record_count;
total_bytes += file.file_size_bytes;
}
recompute_table_column_stats(&tx, table_id, columns, column_ids)?;
tx.execute(
"UPDATE ducklake_table_stats
SET next_row_id = next_row_id + ?,
record_count = record_count + ?,
file_size_bytes = file_size_bytes + ?
WHERE table_id = ?",
params![total_records, total_records, total_bytes, table_id],
)?;
record_table_write_changes(
&tx,
snapshot_id,
table_id,
schema_name,
table_name,
mode,
commit_metadata,
)?;
let schema_id: i64 = tx.query_row(
"SELECT schema_id FROM ducklake_table WHERE table_id = ?",
params![table_id],
|row| row.get(0),
)?;
tx.commit()?;
Ok(CommitIds {
snapshot_id,
schema_id,
table_id,
})
}
fn set_partition_spec(
&self,
table_id: i64,
columns: &[(String, PartitionTransform)],
) -> Result<i64> {
if columns.is_empty() {
return Err(crate::DuckLakeError::InvalidConfig(
"set_partition_spec: partition spec must have at least one column; \
use reset_partition_spec to remove partitioning"
.to_string(),
));
}
let mut conn = self.connection();
let tx = conn.transaction()?;
let partition_id: i64 =
tx.query_row("SELECT nextval('ducklake_partition_id_seq')", [], |row| {
row.get(0)
})?;
let (new_snapshot, _carried) = insert_snapshot(&tx)?;
let mut column_ids: Vec<i64> = Vec::with_capacity(columns.len());
for (name, _transform) in columns {
let column_id: i64 = tx
.query_row(
"SELECT column_id FROM ducklake_column
WHERE table_id = ? AND column_name = ? AND end_snapshot IS NULL
AND parent_column IS NULL",
params![table_id, name.as_str()],
|row| row.get(0),
)
.optional()?
.ok_or_else(|| {
crate::DuckLakeError::InvalidConfig(format!(
"set_partition_spec: no live column '{name}' in table {table_id}"
))
})?;
column_ids.push(column_id);
}
tx.execute(
"UPDATE ducklake_partition_info SET end_snapshot = ?
WHERE table_id = ? AND end_snapshot IS NULL",
params![new_snapshot, table_id],
)?;
tx.execute(
"INSERT INTO ducklake_partition_info
(partition_id, table_id, begin_snapshot, end_snapshot)
VALUES (?, ?, ?, NULL)",
params![partition_id, table_id, new_snapshot],
)?;
for (key_index, column_id) in column_ids.iter().enumerate() {
tx.execute(
"INSERT INTO ducklake_partition_column
(partition_id, table_id, partition_key_index, column_id, transform)
VALUES (?, ?, ?, ?, ?)",
params![
partition_id,
table_id,
key_index as i64,
*column_id,
columns[key_index].1.to_catalog_string()
],
)?;
}
let new_schema_version = bump_schema_version(&tx, new_snapshot)?;
record_schema_version(&tx, new_snapshot, new_schema_version, table_id)?;
record_snapshot_changes(
&tx,
new_snapshot,
&format!("altered_table:{table_id}"),
&SnapshotCommitMetadata::default(),
)?;
tx.commit()?;
Ok(new_snapshot)
}
fn live_partition_spec(
&self,
table_id: i64,
) -> Result<Option<crate::partition::PartitionSpec>> {
let conn = self.connection();
let mut stmt = conn.prepare(
"SELECT pi.partition_id, pc.partition_key_index, pc.column_id, pc.transform
FROM ducklake_partition_info AS pi
JOIN ducklake_partition_column AS pc
ON pc.partition_id = pi.partition_id AND pc.table_id = pi.table_id
WHERE pi.table_id = ? AND pi.end_snapshot IS NULL
ORDER BY pc.partition_key_index",
)?;
let rows = stmt
.query_map(duckdb::params![table_id], |row| {
Ok((
row.get::<_, i64>(0)?,
i32::try_from(row.get::<_, i64>(1)?).unwrap_or(0),
row.get::<_, i64>(2)?,
row.get::<_, String>(3)?,
))
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(crate::partition::PartitionSpec::from_rows(rows, false))
}
fn reset_partition_spec(&self, table_id: i64) -> Result<i64> {
let mut conn = self.connection();
let tx = conn.transaction()?;
let (new_snapshot, _carried) = insert_snapshot(&tx)?;
let ended = tx.execute(
"UPDATE ducklake_partition_info SET end_snapshot = ?
WHERE table_id = ? AND end_snapshot IS NULL",
params![new_snapshot, table_id],
)?;
if ended == 0 {
drop(tx);
let head: i64 = conn.query_row(
"SELECT COALESCE(MAX(snapshot_id), 0) FROM ducklake_snapshot",
[],
|row| row.get(0),
)?;
return Ok(head);
}
let new_schema_version = bump_schema_version(&tx, new_snapshot)?;
record_schema_version(&tx, new_snapshot, new_schema_version, table_id)?;
record_snapshot_changes(
&tx,
new_snapshot,
&format!("altered_table:{table_id}"),
&SnapshotCommitMetadata::default(),
)?;
tx.commit()?;
Ok(new_snapshot)
}
fn live_sort_spec(&self, table_id: i64) -> Result<Option<crate::sort::SortSpec>> {
let conn = self.connection();
let mut stmt = conn.prepare(
"SELECT si.sort_id, se.sort_key_index, se.expression, se.dialect,
se.sort_direction, se.null_order
FROM ducklake_sort_info AS si
JOIN ducklake_sort_expression AS se
ON se.sort_id = si.sort_id AND se.table_id = si.table_id
WHERE si.table_id = ? AND si.end_snapshot IS NULL
ORDER BY se.sort_key_index",
)?;
let rows = stmt
.query_map(duckdb::params![table_id], |row| {
Ok((
row.get::<_, i64>(0)?,
i32::try_from(row.get::<_, i64>(1)?).unwrap_or(0),
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, String>(4)?,
row.get::<_, String>(5)?,
))
})?
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(crate::sort::SortSpec::from_rows(rows))
}
fn set_sort_spec(&self, table_id: i64, fields: &[crate::sort::SortField]) -> Result<i64> {
if fields.is_empty() {
return Err(crate::DuckLakeError::InvalidConfig(
"set_sort_spec: at least one sort key is required (use reset_sort_spec to clear)"
.to_string(),
));
}
let mut conn = self.connection();
let tx = conn.transaction()?;
let sort_id: i64 = tx.query_row("SELECT nextval('ducklake_sort_id_seq')", [], |row| {
row.get(0)
})?;
let (new_snapshot, _carried) = insert_snapshot(&tx)?;
for field in fields {
let column = field.column_candidate().ok_or_else(|| {
crate::DuckLakeError::InvalidConfig(format!(
"set_sort_spec: sort key '{}' is not a bare column; only column \
sort keys are supported",
field.expression
))
})?;
let exists: Option<i64> = tx
.query_row(
"SELECT column_id FROM ducklake_column
WHERE table_id = ? AND column_name = ? AND end_snapshot IS NULL
AND parent_column IS NULL",
params![table_id, column.as_str()],
|row| row.get(0),
)
.optional()?;
if exists.is_none() {
return Err(crate::DuckLakeError::InvalidConfig(format!(
"set_sort_spec: no live column '{column}' in table {table_id}"
)));
}
}
tx.execute(
"UPDATE ducklake_sort_info SET end_snapshot = ?
WHERE table_id = ? AND end_snapshot IS NULL",
params![new_snapshot, table_id],
)?;
tx.execute(
"INSERT INTO ducklake_sort_info
(sort_id, table_id, begin_snapshot, end_snapshot)
VALUES (?, ?, ?, NULL)",
params![sort_id, table_id, new_snapshot],
)?;
for field in fields {
tx.execute(
"INSERT INTO ducklake_sort_expression
(sort_id, table_id, sort_key_index, expression, dialect,
sort_direction, null_order)
VALUES (?, ?, ?, ?, ?, ?, ?)",
params![
sort_id,
table_id,
field.sort_key_index as i64,
field.expression.as_str(),
field.dialect.as_str(),
field.direction.to_catalog_string(),
field.null_order.to_catalog_string()
],
)?;
}
record_snapshot_changes(
&tx,
new_snapshot,
&format!("altered_table:{table_id}"),
&SnapshotCommitMetadata::default(),
)?;
tx.commit()?;
Ok(new_snapshot)
}
fn reset_sort_spec(&self, table_id: i64) -> Result<i64> {
let mut conn = self.connection();
let tx = conn.transaction()?;
let (new_snapshot, _carried) = insert_snapshot(&tx)?;
let ended = tx.execute(
"UPDATE ducklake_sort_info SET end_snapshot = ?
WHERE table_id = ? AND end_snapshot IS NULL",
params![new_snapshot, table_id],
)?;
if ended == 0 {
drop(tx);
let head: i64 = conn.query_row(
"SELECT COALESCE(MAX(snapshot_id), 0) FROM ducklake_snapshot",
[],
|row| row.get(0),
)?;
return Ok(head);
}
record_snapshot_changes(
&tx,
new_snapshot,
&format!("altered_table:{table_id}"),
&SnapshotCommitMetadata::default(),
)?;
tx.commit()?;
Ok(new_snapshot)
}
fn publish_snapshot(
&self,
table_id: i64,
schema_name: &str,
table_name: &str,
_snapshot_id: i64,
mode: WriteMode,
base_snapshot: i64,
columns: &[ColumnDef],
column_ids: &[i64],
) -> Result<CommitIds> {
let mut conn = self.connection();
let tx = conn.transaction()?;
let snapshot_id =
finalize_snapshot(&tx, table_id, columns, column_ids, mode, base_snapshot)?;
record_table_write_changes(
&tx,
snapshot_id,
table_id,
schema_name,
table_name,
mode,
&SnapshotCommitMetadata::default(),
)?;
let schema_id: i64 = tx.query_row(
"SELECT schema_id FROM ducklake_table WHERE table_id = ?",
params![table_id],
|row| row.get(0),
)?;
tx.commit()?;
Ok(CommitIds {
snapshot_id,
schema_id,
table_id,
})
}
fn end_table_files(&self, table_id: i64, snapshot_id: i64) -> Result<u64> {
let mut conn = self.connection();
let tx = conn.transaction()?;
let rows_affected = tx.execute(
"UPDATE ducklake_data_file SET end_snapshot = ?
WHERE table_id = ? AND end_snapshot IS NULL",
params![snapshot_id, table_id],
)? as u64;
tx.execute(
"UPDATE ducklake_table_stats
SET record_count = 0, file_size_bytes = 0
WHERE table_id = ?",
params![table_id],
)?;
tx.commit()?;
Ok(rows_affected)
}
fn get_data_path(&self) -> Result<String> {
let conn = self.connection();
let row: Option<String> = conn
.query_row(
"SELECT value FROM ducklake_metadata WHERE key = ? AND scope IS NULL",
params!["data_path"],
|row| row.get(0),
)
.optional()?;
match row {
Some(path) => Ok(path),
None => Err(crate::error::DuckLakeError::InvalidConfig(
"Missing required catalog metadata: 'data_path' not configured.".to_string(),
)),
}
}
fn set_data_path(&self, path: &str) -> Result<()> {
let mut conn = self.connection();
let tx = conn.transaction()?;
tx.execute(
"DELETE FROM ducklake_metadata WHERE key = 'data_path' AND scope IS NULL",
[],
)?;
tx.execute(
"INSERT INTO ducklake_metadata (key, value, scope) VALUES ('data_path', ?, NULL)",
params![path],
)?;
tx.commit()?;
Ok(())
}
fn initialize_schema(&self) -> Result<()> {
let conn = self.connection();
conn.execute_batch(SQL_CREATE_SCHEMA)?;
conn.execute_batch(
"ALTER TABLE ducklake_data_file ADD COLUMN IF NOT EXISTS partition_id BIGINT",
)?;
conn.execute_batch(
"ALTER TABLE ducklake_snapshot_changes ALTER COLUMN changes_made DROP NOT NULL",
)?;
conn.execute_batch(
"INSERT INTO ducklake_metadata (key, value, scope)
SELECT 'next_column_id',
CAST(COALESCE((SELECT MAX(column_id) FROM ducklake_column), 0) AS VARCHAR),
NULL
WHERE NOT EXISTS (
SELECT 1 FROM ducklake_metadata WHERE key = 'next_column_id' AND scope IS NULL
)",
)?;
Ok(())
}
fn begin_write_transaction(
&self,
schema_name: &str,
table_name: &str,
columns: &[ColumnDef],
mode: WriteMode,
) -> Result<WriteSetupResult> {
validate_name(schema_name, "Schema")?;
validate_name(table_name, "Table")?;
if columns.is_empty() {
return Err(crate::DuckLakeError::InvalidConfig(
"Table must have at least one column".to_string(),
));
}
let mut conn = self.connection();
let tx = conn.transaction()?;
let catalog_columns = catalog_column_defs(columns)?;
let n = catalog_columns.len() as i64;
let last_column_id = reserve_ids(&tx, "next_column_id", n)?;
let fresh_ids: Vec<i64> = ((last_column_id - n + 1)..=last_column_id).collect();
let base_snapshot_id: i64 = tx.query_row(
"SELECT COALESCE(MAX(snapshot_id), 0) FROM ducklake_snapshot",
[],
|row| row.get(0),
)?;
let snapshot_id: i64 = base_snapshot_id + 1;
let schema_id: i64 = {
let existing: Option<i64> = tx
.query_row(
"SELECT schema_id FROM ducklake_schema
WHERE schema_name = ? AND end_snapshot IS NULL",
params![schema_name],
|row| row.get(0),
)
.optional()?;
match existing {
Some(id) => id,
None => tx.query_row(
"INSERT INTO ducklake_schema (schema_name, path, path_is_relative, begin_snapshot)
VALUES (?, ?, true, ?) RETURNING schema_id",
params![schema_name, schema_name, snapshot_id],
|row| row.get(0),
)?,
}
};
let table_id: i64 = {
let existing: Option<i64> = tx
.query_row(
"SELECT table_id FROM ducklake_table
WHERE schema_id = ? AND table_name = ? AND end_snapshot IS NULL",
params![schema_id, table_name],
|row| row.get(0),
)
.optional()?;
match existing {
Some(id) => id,
None => tx.query_row(
"INSERT INTO ducklake_table (schema_id, table_name, path, path_is_relative, begin_snapshot)
VALUES (?, ?, ?, true, ?) RETURNING table_id",
params![schema_id, table_name, table_name, snapshot_id],
|row| row.get(0),
)?,
}
};
let existing_rows: Vec<(String, String, i64, Option<i64>)> = {
let mut stmt = tx.prepare(
"SELECT column_name, column_type, column_id, parent_column
FROM ducklake_column
WHERE table_id = ? AND end_snapshot IS NULL
ORDER BY column_order",
)?;
let mapped = stmt.query_map(params![table_id], |row| {
let name: String = row.get(0)?;
let col_type: String = row.get(1)?;
let cid: i64 = row.get(2)?;
let parent_column: Option<i64> = row.get(3)?;
Ok((name, col_type, cid, parent_column))
})?;
mapped.collect::<std::result::Result<Vec<_>, duckdb::Error>>()?
};
let mut existing_catalog_columns = Vec::with_capacity(existing_rows.len());
for (name, ducklake_type, column_id, parent_column) in existing_rows {
existing_catalog_columns.push(ExistingCatalogColumn {
column_id,
name,
ducklake_type,
parent_column,
});
}
let field_ids = assign_column_ids(&catalog_columns, &existing_catalog_columns, &fresh_ids)?;
if !existing_catalog_columns.is_empty() {
use std::collections::HashMap;
let existing_map: HashMap<i64, &ExistingCatalogColumn> = existing_catalog_columns
.iter()
.map(|column| (column.column_id, column))
.collect();
for (new_column, column_id) in catalog_columns.iter().zip(&field_ids) {
if let Some(existing_column) = existing_map.get(column_id) {
let same_type =
catalog_column_type_equal(&existing_column.ducklake_type, new_column);
if !same_type {
return Err(crate::error::DuckLakeError::UnsupportedTypeChange {
operation: TypeChangeOperation::DataWrite {
mode: match mode {
WriteMode::Replace => TypeChangeWriteMode::Replace,
WriteMode::Append => TypeChangeWriteMode::Append,
},
},
column: new_column.name.clone(),
from: existing_column.ducklake_type.clone(),
to: new_column.ducklake_type.clone(),
});
}
} else if mode == WriteMode::Append
&& new_column.parent_index.is_none()
&& !new_column.is_nullable
{
return Err(crate::error::DuckLakeError::InvalidConfig(format!(
"Schema evolution error: new column '{}' must be nullable. Adding non-nullable columns is not allowed.",
new_column.name
)));
}
}
}
tx.commit()?;
Ok(WriteSetupResult {
snapshot_id,
base_snapshot_id,
schema_id,
table_id,
column_ids: top_level_column_ids(&catalog_columns, &field_ids)?,
field_ids,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::DuckdbMetadataProvider;
use crate::metadata_provider::MetadataProvider;
use arrow::datatypes::{DataType, Field};
use std::sync::Arc;
use tempfile::TempDir;
#[test]
fn duckdb_writer_persists_recursive_columns() {
let temp = TempDir::new().unwrap();
let db_path = temp.path().join("nested.ducklake");
let writer = DuckdbMetadataWriter::new_with_init(db_path.to_str().unwrap()).unwrap();
let levels = DataType::List(Arc::new(Field::new(
"item",
DataType::Struct(
vec![
Arc::new(Field::new("price", DataType::Decimal128(38, 16), false)),
Arc::new(Field::new("count", DataType::UInt32, false)),
]
.into(),
),
false,
)));
let columns = vec![ColumnDef::from_arrow("bids", &levels, false).unwrap()];
let setup = writer
.begin_write_transaction("main", "depths", &columns, WriteMode::Replace)
.unwrap();
assert_eq!(setup.column_ids, vec![setup.field_ids[0]]);
assert_eq!(setup.field_ids.len(), 4);
writer
.register_data_file(
setup.table_id,
"main",
"depths",
setup.snapshot_id,
&DataFileInfo::new("depths.parquet", 1, 1),
WriteMode::Replace,
setup.base_snapshot_id,
&columns,
&setup.field_ids,
)
.unwrap();
let mut connection = writer.connection();
let transaction = connection.transaction().unwrap();
let mut statement = transaction
.prepare(
"SELECT column_id, column_name, column_type, parent_column
FROM ducklake_column
WHERE table_id = ? AND end_snapshot IS NULL
ORDER BY column_order",
)
.unwrap();
let rows = statement
.query_map(params![setup.table_id], |row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, Option<i64>>(3)?,
))
})
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap();
assert_eq!(
rows,
vec![
(setup.field_ids[0], "bids".into(), "list".into(), None),
(
setup.field_ids[1],
"element".into(),
"struct".into(),
Some(setup.field_ids[0]),
),
(
setup.field_ids[2],
"price".into(),
"decimal(38, 16)".into(),
Some(setup.field_ids[1]),
),
(
setup.field_ids[3],
"count".into(),
"uint32".into(),
Some(setup.field_ids[1]),
),
]
);
}
#[test]
fn duckdb_write_then_read_back_via_provider() {
let temp = TempDir::new().unwrap();
let db_path = temp.path().join("catalog.ducklake");
let db_path_str = db_path.to_str().unwrap().to_string();
let data_path = temp.path().join("data");
let data_path_str = data_path.to_str().unwrap().to_string();
{
let writer = DuckdbMetadataWriter::new_with_init(&db_path_str).unwrap();
writer.set_data_path(&data_path_str).unwrap();
let columns = vec![
ColumnDef::new("id", "int64", false).unwrap(),
ColumnDef::new("name", "varchar", true).unwrap(),
];
let setup = writer
.begin_write_transaction("main", "users", &columns, WriteMode::Append)
.unwrap();
assert_eq!(setup.column_ids.len(), 2);
assert_eq!(setup.base_snapshot_id, 0, "fresh catalog: no prior head");
let file = DataFileInfo::new("users/data-0.parquet", 1024, 3).with_footer_size(256);
let commit = writer
.register_data_file(
setup.table_id,
"main",
"users",
setup.snapshot_id,
&file,
WriteMode::Append,
setup.base_snapshot_id,
&columns,
&setup.column_ids,
)
.unwrap();
assert_eq!(commit.snapshot_id, 1, "first write commits snapshot 1");
}
let provider = DuckdbMetadataProvider::new(&db_path_str).unwrap();
let snapshot = provider.get_current_snapshot().unwrap();
assert_eq!(snapshot, 1, "committed head is snapshot 1");
assert_eq!(provider.get_data_path().unwrap(), data_path_str);
let schema = provider
.get_schema_by_name("main", snapshot)
.unwrap()
.expect("schema 'main' must exist");
let table = provider
.get_table_by_name(schema.schema_id, "users", snapshot)
.unwrap()
.expect("table 'users' must exist");
assert_eq!(table.table_name, "users");
let cols = provider
.get_table_structure(table.table_id, snapshot)
.unwrap();
assert_eq!(cols.len(), 2, "both columns must read back");
assert_eq!(cols[0].column_name, "id");
assert_eq!(cols[0].column_type, "int64");
assert!(!cols[0].is_nullable);
assert_eq!(cols[1].column_name, "name");
assert_eq!(cols[1].column_type, "varchar");
assert!(cols[1].is_nullable);
let files = provider
.get_table_files_for_select(table.table_id, snapshot)
.unwrap();
assert_eq!(files.len(), 1, "exactly one data file");
assert_eq!(files[0].file.path, "users/data-0.parquet");
assert!(files[0].file.path_is_relative);
assert_eq!(files[0].file.file_size_bytes, 1024);
assert_eq!(files[0].file.footer_size, Some(256));
assert_eq!(files[0].max_row_count, Some(3));
assert_eq!(
files[0].row_id_start,
Some(0),
"first file starts at rowid 0"
);
}
#[test]
fn duckdb_row_id_start_advances_across_appends() {
let temp = TempDir::new().unwrap();
let db_path = temp.path().join("catalog.ducklake");
let db_path_str = db_path.to_str().unwrap().to_string();
let writer = DuckdbMetadataWriter::new_with_init(&db_path_str).unwrap();
writer.set_data_path("/tmp/does-not-matter").unwrap();
let columns = vec![ColumnDef::new("id", "int64", false).unwrap()];
let setup1 = writer
.begin_write_transaction("main", "t", &columns, WriteMode::Append)
.unwrap();
writer
.register_data_file(
setup1.table_id,
"main",
"t",
setup1.snapshot_id,
&DataFileInfo::new("a.parquet", 100, 3),
WriteMode::Append,
setup1.base_snapshot_id,
&columns,
&setup1.column_ids,
)
.unwrap();
let setup2 = writer
.begin_write_transaction("main", "t", &columns, WriteMode::Append)
.unwrap();
let commit2 = writer
.register_data_file(
setup2.table_id,
"main",
"t",
setup2.snapshot_id,
&DataFileInfo::new("b.parquet", 250, 7),
WriteMode::Append,
setup2.base_snapshot_id,
&columns,
&setup2.column_ids,
)
.unwrap();
assert_eq!(commit2.snapshot_id, 2, "second write commits snapshot 2");
let mut conn = writer.connection();
let tx = conn.transaction().unwrap();
let a_start: i64 = tx
.query_row(
"SELECT row_id_start FROM ducklake_data_file WHERE path = 'a.parquet'",
[],
|r| r.get(0),
)
.unwrap();
let b_start: i64 = tx
.query_row(
"SELECT row_id_start FROM ducklake_data_file WHERE path = 'b.parquet'",
[],
|r| r.get(0),
)
.unwrap();
let (records, next, bytes): (i64, i64, i64) = tx
.query_row(
"SELECT record_count, next_row_id, file_size_bytes
FROM ducklake_table_stats WHERE table_id = ?",
params![setup1.table_id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.unwrap();
tx.commit().unwrap();
assert_eq!(a_start, 0, "first file starts at 0");
assert_eq!(b_start, 3, "second file starts after the first file's rows");
assert_eq!(records, 10, "record_count = 3 + 7");
assert_eq!(next, 10, "next_row_id advances by sum of record_counts");
assert_eq!(bytes, 350, "file_size_bytes accumulates");
}
}