use std::sync::Arc;
use delta_kernel::arrow::array::{ArrayRef, Int32Array};
use delta_kernel::committer::{Committer, FileSystemCommitter};
use delta_kernel::schema::{schema, schema_ref};
use delta_kernel::snapshot::{
CheckpointWriteResult, ChecksumWriteResult, IncrementalReplay, SnapshotBuilder,
};
use delta_kernel::transaction::create_table::create_table;
use delta_kernel::{DeltaResult, Error, Snapshot, Version};
use rstest::rstest;
use serde_json::json;
use test_utils::delta_kernel_default_engine::executor::TaskExecutor;
use test_utils::delta_kernel_default_engine::DefaultEngine;
use test_utils::{
add_commit, assert_result_error_with_message, engine_store_setup, insert_data_with,
test_table_setup_mt, TestCatalogCommitter,
};
const LATEST_VERSION: Version = 3;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum TableKind {
FileSystem,
CatalogManaged,
}
impl TableKind {
fn committer(self) -> Box<dyn Committer> {
match self {
TableKind::FileSystem => Box::new(FileSystemCommitter::new()),
TableKind::CatalogManaged => Box::new(TestCatalogCommitter),
}
}
}
fn maybe_attach_max_catalog_version(
builder: SnapshotBuilder,
max_catalog_version: Version,
kind: TableKind,
) -> SnapshotBuilder {
match kind {
TableKind::FileSystem => builder,
TableKind::CatalogManaged => builder.with_max_catalog_version(max_catalog_version),
}
}
async fn append_row<E: TaskExecutor>(
snapshot: Arc<Snapshot>,
engine: &Arc<DefaultEngine<E>>,
kind: TableKind,
value: i32,
) -> DeltaResult<Arc<Snapshot>> {
let column: ArrayRef = Arc::new(Int32Array::from(vec![value]));
Ok(insert_data_with(
snapshot,
engine,
vec![column],
kind.committer(),
"WRITE",
true,
false,
)
.await?
.unwrap_post_commit_snapshot())
}
async fn setup_multi_version_table<E: TaskExecutor>(
engine: &Arc<DefaultEngine<E>>,
table_path: &str,
kind: TableKind,
) -> DeltaResult<()> {
let schema = schema_ref! { nullable "id": INTEGER };
let builder = create_table(table_path, schema, "test_engine");
let builder = match kind {
TableKind::FileSystem => builder,
TableKind::CatalogManaged => builder.with_table_properties([
("delta.feature.catalogManaged", "supported"),
("io.unitycatalog.tableId", "snapshot-build-test"),
]),
};
let create_snapshot = builder
.build(engine.as_ref(), kind.committer())?
.commit(engine.as_ref())?
.unwrap_post_commit_snapshot();
assert!(create_snapshot.is_built_as_latest());
create_snapshot.write_checksum(engine.as_ref())?;
let mut snap = maybe_attach_max_catalog_version(Snapshot::builder_for(table_path), 0, kind)
.build(engine.as_ref())?;
for value in 1..=LATEST_VERSION as i32 {
snap = append_row(snap, engine, kind, value).await?;
}
Ok(())
}
#[tokio::test]
async fn deeply_nested_schema_snapshot_load_returns_schema_error(
) -> Result<(), Box<dyn std::error::Error>> {
let deeply_nested_schema = (0..42).fold(
schema! { nullable "leaf": INTEGER },
|nested, depth| schema! { nullable (format!("level_{depth}")): (nested) },
);
let (store, engine, table_url) = engine_store_setup("deeply_nested_schema", None);
let commit = [
json!({"protocol": {"minReaderVersion": 1, "minWriterVersion": 2}}),
json!({
"metaData": {
"id": "deeply-nested-schema",
"format": {"provider": "parquet", "options": {}},
"schemaString": serde_json::to_string(&deeply_nested_schema)?,
"partitionColumns": [],
"configuration": {},
"createdTime": 0
}
}),
]
.map(|action| action.to_string())
.join("\n");
add_commit(table_url.as_str(), store.as_ref(), 0, commit).await?;
let result = Snapshot::builder_for(table_url).build(&engine);
assert_result_error_with_message(
result.as_ref(),
concat!(
"Schema error: Table schema is too deeply nested: decoding ",
"metaData.schemaString exceeded serde_json's ",
"recursion limit: recursion limit exceeded"
),
);
let error = match result.unwrap_err() {
Error::Backtraced { source, .. } => *source,
error => error,
};
assert!(matches!(error, Error::Schema(_)));
Ok(())
}
#[rstest]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn built_as_latest_is_inherited_by_derived_snapshots(
#[values(true, false)] base_snap_time_travel_to_latest: bool,
#[values(TableKind::FileSystem, TableKind::CatalogManaged)] kind: TableKind,
) -> DeltaResult<()> {
let (_temp_dir, table_path, engine) = test_table_setup_mt()?;
setup_multi_version_table(&engine, &table_path, kind).await?;
let mut builder =
maybe_attach_max_catalog_version(Snapshot::builder_for(&table_path), LATEST_VERSION, kind)
.with_incremental_crc_replay(IncrementalReplay::Unlimited);
if base_snap_time_travel_to_latest {
builder = builder.at_version(LATEST_VERSION);
}
let base = builder.build(engine.as_ref())?;
assert_eq!(base.version(), LATEST_VERSION);
let base_is_latest = !base_snap_time_travel_to_latest || kind == TableKind::CatalogManaged;
assert_eq!(base.is_built_as_latest(), base_is_latest);
let (ckpt_result, after_checkpoint) = base.checkpoint(engine.as_ref(), None)?;
assert_eq!(ckpt_result, CheckpointWriteResult::Written);
assert_eq!(after_checkpoint.is_built_as_latest(), base_is_latest);
let (crc_result, after_checksum) = after_checkpoint.write_checksum(engine.as_ref())?;
assert_eq!(crc_result, ChecksumWriteResult::Written);
assert_eq!(after_checksum.is_built_as_latest(), base_is_latest);
if kind == TableKind::CatalogManaged {
let published = base.publish(engine.as_ref(), &TestCatalogCommitter)?;
assert_eq!(published.is_built_as_latest(), base_is_latest);
}
let post_commit = append_row(base, &engine, kind, 4).await?;
assert_eq!(post_commit.version(), 4);
assert!(post_commit.is_built_as_latest());
Ok(())
}
#[rstest]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn built_as_latest_on_fresh_and_incremental_build(
#[values(None, Some(1), Some(3))] time_travel_version: Option<Version>,
#[values(TableKind::FileSystem, TableKind::CatalogManaged)] kind: TableKind,
) -> DeltaResult<()> {
let (_temp_dir, table_path, engine) = test_table_setup_mt()?;
setup_multi_version_table(&engine, &table_path, kind).await?;
let pins_catalog_latest =
kind == TableKind::CatalogManaged && time_travel_version == Some(LATEST_VERSION);
let mut builder =
maybe_attach_max_catalog_version(Snapshot::builder_for(&table_path), LATEST_VERSION, kind);
if let Some(v) = time_travel_version {
builder = builder.at_version(v);
}
let base = builder.build(engine.as_ref())?;
assert_eq!(
base.version(),
time_travel_version.unwrap_or(LATEST_VERSION)
);
let base_is_latest = time_travel_version.is_none() || pins_catalog_latest;
assert_eq!(base.is_built_as_latest(), base_is_latest);
let refreshed = maybe_attach_max_catalog_version(
Snapshot::builder_from(base.clone()),
LATEST_VERSION,
kind,
)
.build(engine.as_ref())?;
assert_eq!(refreshed.version(), LATEST_VERSION);
assert!(refreshed.is_built_as_latest());
let pinned =
maybe_attach_max_catalog_version(Snapshot::builder_from(base), LATEST_VERSION, kind)
.at_version(LATEST_VERSION)
.build(engine.as_ref())?;
assert_eq!(pinned.version(), LATEST_VERSION);
assert_eq!(
pinned.is_built_as_latest(),
base_is_latest || kind == TableKind::CatalogManaged
);
Ok(())
}