use std::collections::HashMap;
use std::sync::LazyLock;
use delta_kernel_derive::{
internal_api, IntoEngineData, IntoStructData, ToSchema, TryFromStructData,
};
use serde::{Deserialize, Serialize};
use tracing::warn;
use url::Url;
use visitors::{MetadataVisitor, ProtocolVisitor};
use self::deletion_vector::DeletionVectorDescriptor;
#[cfg(feature = "adaptive-metadata-in-dev")]
use crate::expressions::Scalar;
#[cfg(feature = "adaptive-metadata-in-dev")]
use crate::expressions::{ArrayData, StructData};
use crate::schema::{
is_unsupported_delta_type_error, lazy_schema_ref, schema_ref, SchemaRef, StructField,
StructType, ToSchema as _,
};
#[cfg(feature = "adaptive-metadata-in-dev")]
use crate::schema::{schema, ArrayType, DataType};
use crate::table_features::{
FeatureType, TableFeature, LEGACY_READER_FEATURES, MIN_VALID_RW_VERSION,
TABLE_FEATURES_MIN_READER_VERSION, TABLE_FEATURES_MIN_WRITER_VERSION,
};
use crate::table_properties::TableProperties;
use crate::utils::require;
use crate::{
DeltaResult, Engine, EngineData, Error, EvaluationHandlerExtension as _, FileMeta, FileSize,
IntoEngineData, RowVisitor as _,
};
const KERNEL_VERSION: &str = env!("CARGO_PKG_VERSION");
const SERDE_JSON_RECURSION_LIMIT_ERROR_PREFIX: &str = "recursion limit exceeded";
const UNKNOWN_OPERATION: &str = "UNKNOWN";
pub mod deletion_vector;
pub mod deletion_vector_writer;
pub mod set_transaction;
#[cfg(feature = "internal-api")]
pub mod visitors;
#[cfg(not(feature = "internal-api"))]
pub(crate) mod visitors;
#[internal_api]
pub(crate) const ADD_NAME: &str = "add";
#[internal_api]
pub(crate) const REMOVE_NAME: &str = "remove";
#[internal_api]
pub(crate) const METADATA_NAME: &str = "metaData";
#[internal_api]
pub(crate) const PROTOCOL_NAME: &str = "protocol";
#[internal_api]
pub(crate) const SET_TRANSACTION_NAME: &str = "txn";
#[internal_api]
pub(crate) const COMMIT_INFO_NAME: &str = "commitInfo";
#[internal_api]
pub(crate) const CDC_NAME: &str = "cdc";
#[internal_api]
pub(crate) const SIDECAR_NAME: &str = "sidecar";
#[internal_api]
pub(crate) const CHECKPOINT_METADATA_NAME: &str = "checkpointMetadata";
#[internal_api]
pub(crate) const DOMAIN_METADATA_NAME: &str = "domainMetadata";
#[cfg(feature = "adaptive-metadata-in-dev")]
#[internal_api]
pub(crate) const CHECKPOINT_ACTION_NAME: &str = "checkpoint";
#[cfg(feature = "adaptive-metadata-in-dev")]
#[internal_api]
pub(crate) const CONTENT_ROOT_NAME: &str = "contentRoot";
pub(crate) const INTERNAL_DOMAIN_PREFIX: &str = "delta.";
pub(crate) fn action_presence_leaf(action_name: &str) -> Option<&'static str> {
match action_name {
ADD_NAME | REMOVE_NAME | CDC_NAME | SIDECAR_NAME => Some("path"),
METADATA_NAME => Some("id"),
PROTOCOL_NAME => Some("minReaderVersion"),
SET_TRANSACTION_NAME => Some("appId"),
DOMAIN_METADATA_NAME => Some("domain"),
CHECKPOINT_METADATA_NAME => Some("version"),
_ => None,
}
}
#[internal_api]
pub(crate) const NUM_RECORDS: &str = "numRecords";
#[internal_api]
pub(crate) const NULL_COUNT: &str = "nullCount";
#[internal_api]
pub(crate) const MIN_VALUES: &str = "minValues";
#[internal_api]
pub(crate) const MAX_VALUES: &str = "maxValues";
#[internal_api]
pub(crate) const TIGHT_BOUNDS: &str = "tightBounds";
#[internal_api]
pub(crate) const STATS_PARSED: &str = "stats_parsed";
pub(crate) static ADD_SCHEMA: LazyLock<StructType> = LazyLock::new(Add::to_schema);
pub(crate) static ADD_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(ADD_NAME, ADD_SCHEMA.clone()));
pub(crate) static REMOVE_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(REMOVE_NAME, Remove::to_schema()));
pub(crate) static METADATA_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(METADATA_NAME, Metadata::to_schema()));
pub(crate) static PROTOCOL_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(PROTOCOL_NAME, Protocol::to_schema()));
pub(crate) static SET_TRANSACTION_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(SET_TRANSACTION_NAME, SetTransaction::to_schema()));
pub(crate) static COMMIT_INFO_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(COMMIT_INFO_NAME, CommitInfo::to_schema()));
pub(crate) static CDC_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(CDC_NAME, Cdc::to_schema()));
pub(crate) static DOMAIN_METADATA_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(DOMAIN_METADATA_NAME, DomainMetadata::to_schema()));
pub(crate) static CHECKPOINT_METADATA_FIELD: LazyLock<StructField> = LazyLock::new(|| {
StructField::nullable(CHECKPOINT_METADATA_NAME, CheckpointMetadata::to_schema())
});
pub(crate) static SIDECAR_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(SIDECAR_NAME, Sidecar::to_schema()));
#[cfg(feature = "adaptive-metadata-in-dev")]
pub(crate) static CONTENT_ROOT_FIELD: LazyLock<StructField> =
LazyLock::new(|| StructField::nullable(CONTENT_ROOT_NAME, ContentRoot::to_schema()));
#[cfg(feature = "adaptive-metadata-in-dev")]
static CONTENT_SIDECAR_FIELD: LazyLock<StructField> = LazyLock::new(|| {
StructField::nullable(
SIDECAR_NAME,
schema! {
not_null "type": STRING,
..(Sidecar::to_schema().into_fields()),
},
)
});
#[cfg(feature = "adaptive-metadata-in-dev")]
static CHECKPOINT_ACTION_ELEMENT_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! {
(&CHECKPOINT_METADATA_FIELD),
(&CONTENT_ROOT_FIELD),
(&PROTOCOL_FIELD),
(&METADATA_FIELD),
(&DOMAIN_METADATA_FIELD),
(&SET_TRANSACTION_FIELD),
(&CONTENT_SIDECAR_FIELD),
};
#[cfg(feature = "adaptive-metadata-in-dev")]
pub(crate) static CHECKPOINT_ACTION_FIELD: LazyLock<StructField> = LazyLock::new(|| {
StructField::nullable(
CHECKPOINT_ACTION_NAME,
ArrayType::new(CHECKPOINT_ACTION_ELEMENT_SCHEMA.clone(), false),
)
});
fn checkpoint_action_field() -> impl IntoIterator<Item = &'static StructField> {
#[cfg(feature = "adaptive-metadata-in-dev")]
{
Some(&*CHECKPOINT_ACTION_FIELD)
}
#[cfg(not(feature = "adaptive-metadata-in-dev"))]
{
None::<&'static StructField>
}
}
#[cfg(any(test, feature = "internal-api"))]
static COMMIT_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! {
(&ADD_FIELD),
(&REMOVE_FIELD),
(&METADATA_FIELD),
(&PROTOCOL_FIELD),
(&SET_TRANSACTION_FIELD),
(&COMMIT_INFO_FIELD),
(&CDC_FIELD),
(&DOMAIN_METADATA_FIELD),
..(checkpoint_action_field()),
};
static ALL_ACTIONS_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! {
(&ADD_FIELD),
(&REMOVE_FIELD),
(&METADATA_FIELD),
(&PROTOCOL_FIELD),
(&SET_TRANSACTION_FIELD),
(&COMMIT_INFO_FIELD),
(&CDC_FIELD),
(&DOMAIN_METADATA_FIELD),
..(checkpoint_action_field()),
(&CHECKPOINT_METADATA_FIELD),
(&SIDECAR_FIELD),
};
#[internal_api]
pub(crate) static LOG_ADD_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! { (&ADD_FIELD) };
#[internal_api]
pub(crate) static LOG_REMOVE_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! { (&REMOVE_FIELD) };
#[internal_api]
pub(crate) static LOG_METADATA_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! { (&METADATA_FIELD) };
#[cfg(feature = "adaptive-metadata-in-dev")]
#[internal_api]
pub(crate) static LOG_CHECKPOINT_SCHEMA: LazyLock<SchemaRef> =
lazy_schema_ref! { (&CHECKPOINT_ACTION_FIELD) };
#[internal_api]
pub(crate) static LOG_PROTOCOL_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! { (&PROTOCOL_FIELD) };
#[internal_api]
pub(crate) static LOG_COMMIT_INFO_SCHEMA: LazyLock<SchemaRef> =
lazy_schema_ref! { (&COMMIT_INFO_FIELD) };
#[internal_api]
pub(crate) static LOG_TXN_SCHEMA: LazyLock<SchemaRef> =
lazy_schema_ref! { (&SET_TRANSACTION_FIELD) };
#[internal_api]
pub(crate) static LOG_DOMAIN_METADATA_SCHEMA: LazyLock<SchemaRef> =
lazy_schema_ref! { (&DOMAIN_METADATA_FIELD) };
#[cfg(any(test, feature = "internal-api"))]
#[internal_api]
pub(crate) fn get_commit_schema() -> &'static SchemaRef {
&COMMIT_SCHEMA
}
#[internal_api]
#[allow(dead_code)]
pub(crate) fn get_all_actions_schema() -> &'static SchemaRef {
&ALL_ACTIONS_SCHEMA
}
#[internal_api]
pub(crate) fn schema_contains_file_actions(schema: &SchemaRef) -> bool {
schema.contains(ADD_NAME) || schema.contains(REMOVE_NAME)
}
pub(crate) fn as_log_add_schema(add_schema: SchemaRef) -> SchemaRef {
schema_ref! { nullable ADD_NAME: (add_schema) }
}
#[derive(
Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema, IntoStructData, TryFromStructData,
)]
#[serde(rename_all = "camelCase")]
#[internal_api]
pub(crate) struct Format {
pub(crate) provider: String,
pub(crate) options: HashMap<String, String>,
}
impl Default for Format {
fn default() -> Self {
Self {
provider: String::from("parquet"),
options: HashMap::new(),
}
}
}
#[derive(
Debug, Default, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema, IntoStructData,
)]
#[serde(rename_all = "camelCase")]
pub struct Metadata {
id: String,
name: Option<String>,
description: Option<String>,
format: Format,
schema_string: String,
partition_columns: Vec<String>,
created_time: Option<i64>,
configuration: HashMap<String, String>,
}
impl Metadata {
#[internal_api]
pub(crate) fn try_new(
name: Option<String>,
description: Option<String>,
schema: SchemaRef,
partition_columns: Vec<String>,
created_time: i64,
configuration: HashMap<String, String>,
) -> DeltaResult<Self> {
if let Some(metadata_field) = schema.fields().find(|field| field.is_metadata_column()) {
return Err(Error::Schema(format!(
"Table schema must not contain metadata columns. Found metadata column: '{}'",
metadata_field.name
)));
}
Ok(Self {
id: uuid::Uuid::new_v4().to_string(),
name,
description,
format: Format::default(),
schema_string: serde_json::to_string(&schema)?,
partition_columns,
created_time: Some(created_time),
configuration,
})
}
#[internal_api]
pub(crate) fn try_new_from_data(data: &dyn EngineData) -> DeltaResult<Option<Metadata>> {
let mut visitor = MetadataVisitor::default();
visitor.visit_rows_of(data)?;
Ok(visitor.metadata)
}
#[internal_api]
#[allow(dead_code)]
pub(crate) fn id(&self) -> &str {
&self.id
}
#[internal_api]
#[allow(dead_code)]
pub(crate) fn name(&self) -> Option<&str> {
self.name.as_deref()
}
#[internal_api]
#[allow(dead_code)]
pub(crate) fn description(&self) -> Option<&str> {
self.description.as_deref()
}
#[internal_api]
#[allow(dead_code)]
pub(crate) fn created_time(&self) -> Option<i64> {
self.created_time
}
#[internal_api]
pub(crate) fn configuration(&self) -> &HashMap<String, String> {
&self.configuration
}
#[internal_api]
#[allow(dead_code)]
pub(crate) fn format_provider(&self) -> &str {
&self.format.provider
}
#[internal_api]
pub(crate) fn schema_string(&self) -> &String {
&self.schema_string
}
#[internal_api]
pub(crate) fn parse_schema(&self) -> DeltaResult<StructType> {
serde_json::from_str(&self.schema_string).map_err(|error| {
if error.is_syntax()
&& error
.to_string()
.starts_with(SERDE_JSON_RECURSION_LIMIT_ERROR_PREFIX)
{
Error::schema(format!(
"Table schema is too deeply nested: decoding metaData.schemaString exceeded \
serde_json's recursion limit: {error}"
))
.with_backtrace()
} else if is_unsupported_delta_type_error(&error) {
Error::schema(error.to_string()).with_backtrace()
} else {
error.into()
}
})
}
#[internal_api]
pub(crate) fn partition_columns(&self) -> &[String] {
&self.partition_columns
}
#[internal_api]
pub(crate) fn parse_table_properties(&self) -> TableProperties {
TableProperties::from(self.configuration.iter())
}
pub(crate) fn with_schema(self, schema: SchemaRef) -> DeltaResult<Self> {
Ok(Self {
schema_string: serde_json::to_string(&schema)?,
..self
})
}
pub(crate) fn with_configuration_entry(
mut self,
key: impl Into<String>,
value: impl Into<String>,
) -> Self {
self.configuration.insert(key.into(), value.into());
self
}
#[cfg(test)]
#[allow(clippy::too_many_arguments)]
pub(crate) fn new_unchecked(
id: impl Into<String>,
name: Option<String>,
description: Option<String>,
format: Format,
schema_string: impl Into<String>,
partition_columns: Vec<String>,
created_time: Option<i64>,
configuration: HashMap<String, String>,
) -> Self {
Self {
id: id.into(),
name,
description,
format,
schema_string: schema_string.into(),
partition_columns,
created_time,
configuration,
}
}
}
impl IntoEngineData for Metadata {
fn into_engine_data(
self,
schema: SchemaRef,
engine: &dyn Engine,
) -> DeltaResult<Box<dyn EngineData>> {
let values = [
self.id.into(),
self.name.into(),
self.description.into(),
self.format.provider.into(),
self.format.options.into(),
self.schema_string.into(),
self.partition_columns.into(),
self.created_time.into(),
self.configuration.into(),
];
engine.evaluation_handler().create_one(schema, &values)
}
}
#[derive(
Default,
Debug,
Clone,
PartialEq,
Eq,
ToSchema,
IntoStructData,
Serialize,
Deserialize,
IntoEngineData,
)]
#[serde(rename_all = "camelCase", try_from = "ProtocolRaw")]
pub struct Protocol {
min_reader_version: i32,
min_writer_version: i32,
#[serde(skip_serializing_if = "Option::is_none")]
reader_features: Option<Vec<TableFeature>>,
#[serde(skip_serializing_if = "Option::is_none")]
writer_features: Option<Vec<TableFeature>>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct ProtocolRaw {
min_reader_version: i32,
min_writer_version: i32,
reader_features: Option<Vec<TableFeature>>,
writer_features: Option<Vec<TableFeature>>,
}
impl TryFrom<ProtocolRaw> for Protocol {
type Error = Error;
fn try_from(protocol: ProtocolRaw) -> DeltaResult<Self> {
Protocol::try_new(
protocol.min_reader_version,
protocol.min_writer_version,
protocol.reader_features,
protocol.writer_features,
)
}
}
fn parse_features(
features: Option<impl IntoIterator<Item = impl Into<TableFeature>>>,
) -> Option<Vec<TableFeature>> {
let features = features?.into_iter().map(Into::into);
Some(features.collect())
}
impl Protocol {
pub(crate) fn try_new_modern(
reader_features: impl IntoIterator<Item = impl Into<TableFeature>>,
writer_features: impl IntoIterator<Item = impl Into<TableFeature>>,
) -> DeltaResult<Self> {
Self::try_new(
TABLE_FEATURES_MIN_READER_VERSION,
TABLE_FEATURES_MIN_WRITER_VERSION,
Some(reader_features),
Some(writer_features),
)
}
#[cfg(test)]
pub(crate) fn try_new_legacy(
min_reader_version: i32,
min_writer_version: i32,
) -> DeltaResult<Self> {
Self::try_new(
min_reader_version,
min_writer_version,
TableFeature::NO_LIST,
TableFeature::NO_LIST,
)
}
pub(crate) fn try_new(
min_reader_version: i32,
min_writer_version: i32,
reader_features: Option<impl IntoIterator<Item = impl Into<TableFeature>>>,
writer_features: Option<impl IntoIterator<Item = impl Into<TableFeature>>>,
) -> DeltaResult<Self> {
require!(
min_reader_version >= MIN_VALID_RW_VERSION,
Error::InvalidProtocol(format!(
"min_reader_version must be >= {MIN_VALID_RW_VERSION}, got {min_reader_version}"
))
);
require!(
min_writer_version >= MIN_VALID_RW_VERSION,
Error::InvalidProtocol(format!(
"min_writer_version must be >= {MIN_VALID_RW_VERSION}, got {min_writer_version}"
))
);
let reader_features = parse_features(reader_features);
let writer_features = parse_features(writer_features);
if min_reader_version == TABLE_FEATURES_MIN_READER_VERSION {
require!(
reader_features.is_some(),
Error::invalid_protocol(
"Reader features must be present when minimum reader version = 3"
)
);
} else {
require!(
reader_features.is_none(),
Error::invalid_protocol(
"Reader features must not be present when minimum reader version != 3"
)
);
}
if min_writer_version == TABLE_FEATURES_MIN_WRITER_VERSION {
require!(
writer_features.is_some(),
Error::invalid_protocol(
"Writer features must be present when minimum writer version = 7"
)
);
} else {
require!(
writer_features.is_none(),
Error::invalid_protocol(
"Writer features must not be present when minimum writer version != 7"
)
);
}
match (&reader_features, &writer_features) {
(Some(reader_features), Some(writer_features)) => {
if let Some(offending) = reader_features.iter().find(|feature| {
!matches!(
feature.feature_type(),
FeatureType::ReaderWriter | FeatureType::Unknown
) || !writer_features.contains(*feature)
}) {
return Err(Error::invalid_protocol(format!(
"Reader features must contain only ReaderWriter features that are also \
listed in writer features, but {offending:?} is not \
(readerFeatures={reader_features:?}, writerFeatures={writer_features:?}, \
minReaderVersion={min_reader_version}, minWriterVersion={min_writer_version})"
)));
}
let mut legacy_orphans = Vec::new();
for feature in writer_features.iter() {
let orphaned_reader_writer_feature = feature.feature_type()
== FeatureType::ReaderWriter
&& !reader_features.contains(feature);
if !orphaned_reader_writer_feature {
continue;
}
if LEGACY_READER_FEATURES.contains(feature) {
legacy_orphans.push(feature);
} else {
return Err(Error::invalid_protocol(format!(
"Writer features must be Writer-only or also listed in reader features, \
but ReaderWriter feature {feature:?} is listed in writerFeatures and \
missing from readerFeatures \
(readerFeatures={reader_features:?}, \
writerFeatures={writer_features:?}, \
minReaderVersion={min_reader_version}, \
minWriterVersion={min_writer_version})"
)));
}
}
for feature in legacy_orphans {
warn!(
"ReaderWriter feature {feature:?} is listed in writerFeatures but \
missing from readerFeatures at minReaderVersion={min_reader_version}; \
treating it as reader-enabled (malformed protocol)"
);
}
Ok(())
}
(None, None) => Ok(()),
(None, Some(writer_features)) => {
if let Some(offending) = writer_features.iter().find(|feature| {
match feature.feature_type() {
FeatureType::WriterOnly | FeatureType::Unknown => false,
FeatureType::ReaderWriter => {
!(min_reader_version == 2 && *feature == &TableFeature::ColumnMapping)
}
}
}) {
return Err(Error::invalid_protocol(format!(
"Writer features must be Writer-only or also listed in reader features, \
but ReaderWriter feature {offending:?} is listed in writerFeatures with \
no reader features present \
(writerFeatures={writer_features:?}, minReaderVersion={min_reader_version}, \
minWriterVersion={min_writer_version})"
)));
}
Ok(())
}
(Some(_), None) => Err(Error::invalid_protocol(
"Reader features should be present in writer features",
)),
}?;
Ok(Protocol {
min_reader_version,
min_writer_version,
reader_features,
writer_features,
})
}
pub(crate) fn try_new_from_data(data: &dyn EngineData) -> DeltaResult<Option<Protocol>> {
let mut visitor = ProtocolVisitor::default();
visitor.visit_rows_of(data)?;
Ok(visitor.protocol)
}
#[internal_api]
pub(crate) fn min_reader_version(&self) -> i32 {
self.min_reader_version
}
#[internal_api]
pub(crate) fn min_writer_version(&self) -> i32 {
self.min_writer_version
}
#[internal_api]
pub(crate) fn reader_features(&self) -> Option<&[TableFeature]> {
self.reader_features.as_deref()
}
#[internal_api]
pub(crate) fn writer_features(&self) -> Option<&[TableFeature]> {
self.writer_features.as_deref()
}
pub(crate) fn has_table_feature(&self, feature: &TableFeature) -> bool {
self.writer_features()
.is_some_and(|features| features.contains(feature))
}
#[cfg(test)]
pub(crate) fn new_unchecked(
min_reader_version: i32,
min_writer_version: i32,
reader_features: Option<Vec<TableFeature>>,
writer_features: Option<Vec<TableFeature>>,
) -> Self {
Self {
min_reader_version,
min_writer_version,
reader_features,
writer_features,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, ToSchema, IntoEngineData)]
#[internal_api]
#[cfg_attr(test, derive(Serialize, Default), serde(rename_all = "camelCase"))]
pub(crate) struct CommitInfo {
pub(crate) timestamp: Option<i64>,
pub(crate) in_commit_timestamp: Option<i64>,
pub(crate) operation: Option<String>,
pub(crate) operation_parameters: Option<HashMap<String, Option<String>>>,
pub(crate) operation_metrics: Option<HashMap<String, Option<String>>>,
pub(crate) kernel_version: Option<String>,
pub(crate) is_blind_append: Option<bool>,
pub(crate) engine_info: Option<String>,
pub(crate) txn_id: Option<String>,
}
impl CommitInfo {
pub(crate) fn new(
timestamp: i64,
in_commit_timestamp: Option<i64>,
operation: Option<String>,
engine_info: Option<String>,
is_blind_append: bool,
) -> Self {
Self {
timestamp: Some(timestamp),
in_commit_timestamp,
operation: Some(operation.unwrap_or_else(|| UNKNOWN_OPERATION.to_string())),
operation_parameters: Some(HashMap::new()),
operation_metrics: None,
kernel_version: Some(format!("v{KERNEL_VERSION}")),
is_blind_append: is_blind_append.then_some(true),
engine_info,
txn_id: Some(uuid::Uuid::new_v4().to_string()),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, ToSchema)]
#[cfg_attr(
test,
derive(Serialize, Deserialize, Default),
serde(rename_all = "camelCase")
)]
#[internal_api]
pub(crate) struct Add {
pub(crate) path: String,
#[allow_null_container_values]
pub(crate) partition_values: HashMap<String, String>,
pub(crate) size: i64,
pub(crate) modification_time: i64,
pub(crate) data_change: bool,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub stats: Option<String>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub tags: Option<HashMap<String, Option<String>>>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub deletion_vector: Option<DeletionVectorDescriptor>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub base_row_id: Option<i64>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub default_row_commit_version: Option<i64>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub clustering_provider: Option<String>,
}
impl Add {
#[internal_api]
#[allow(dead_code)]
pub(crate) fn dv_unique_id(&self) -> Option<String> {
self.deletion_vector.as_ref().map(|dv| dv.unique_id())
}
}
#[derive(Debug, Clone, PartialEq, Eq, ToSchema)]
#[internal_api]
#[cfg_attr(test, derive(Serialize, Default), serde(rename_all = "camelCase"))]
pub(crate) struct Remove {
pub(crate) path: String,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) deletion_timestamp: Option<i64>,
pub(crate) data_change: bool,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) extended_file_metadata: Option<bool>,
#[allow_null_container_values]
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) partition_values: Option<HashMap<String, String>>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) size: Option<i64>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub stats: Option<String>,
#[allow_null_container_values]
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) tags: Option<HashMap<String, String>>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) deletion_vector: Option<DeletionVectorDescriptor>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) base_row_id: Option<i64>,
#[cfg_attr(test, serde(skip_serializing_if = "Option::is_none"))]
pub(crate) default_row_commit_version: Option<i64>,
}
#[derive(Debug, Clone, PartialEq, Eq, ToSchema)]
#[internal_api]
#[cfg_attr(test, derive(Serialize, Default), serde(rename_all = "camelCase"))]
pub(crate) struct Cdc {
pub path: String,
#[allow_null_container_values]
pub partition_values: HashMap<String, String>,
pub size: i64,
pub data_change: bool,
#[allow_null_container_values]
pub tags: Option<HashMap<String, String>>,
}
#[derive(
Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema, IntoStructData, IntoEngineData,
)]
#[serde(rename_all = "camelCase")]
#[internal_api]
pub(crate) struct SetTransaction {
pub(crate) app_id: String,
pub(crate) version: i64,
pub(crate) last_updated: Option<i64>,
}
impl SetTransaction {
pub(crate) fn new(app_id: String, version: i64, last_updated: Option<i64>) -> Self {
Self {
app_id,
version,
last_updated,
}
}
pub(crate) fn is_expired(&self, expiration_timestamp: Option<i64>) -> bool {
matches!(
(expiration_timestamp, self.last_updated),
(Some(exp_ts), Some(lu)) if lu <= exp_ts
)
}
pub(crate) fn non_expired_version(&self, expiration_timestamp: Option<i64>) -> Option<i64> {
(!self.is_expired(expiration_timestamp)).then_some(self.version)
}
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[derive(Debug, Clone, PartialEq, Eq, ToSchema, IntoStructData)]
#[internal_api]
#[cfg_attr(
test,
derive(Serialize, Deserialize, Default),
serde(rename_all = "camelCase")
)]
pub(crate) struct ContentRoot {
pub(crate) path: String,
size_in_bytes: i64,
version: i64,
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[derive(Debug, Clone, PartialEq, Eq)]
#[internal_api]
pub(crate) struct CheckpointAction {
pub(crate) version: i64,
pub(crate) content_root: ContentRoot,
pub(crate) protocol: Protocol,
pub(crate) metadata: Metadata,
pub(crate) transactions: Vec<SetTransaction>,
pub(crate) domain_metadata: Vec<DomainMetadata>,
pub(crate) txn_sidecars: Vec<Sidecar>,
pub(crate) domain_metadata_sidecars: Vec<Sidecar>,
}
#[cfg(feature = "adaptive-metadata-in-dev")]
fn content_sidecar_element(type_str: &str, sidecar: Sidecar) -> DeltaResult<Scalar> {
let sidecar: StructData = sidecar.into();
let fields = std::iter::once(StructField::not_null("type", DataType::STRING))
.chain(sidecar.fields().iter().cloned());
let values = std::iter::once(Scalar::from(type_str))
.chain(sidecar.values().iter().cloned())
.collect();
Ok(Scalar::Struct(StructData::from_values_unchecked(
StructType::try_new(fields)?,
values,
)))
}
#[cfg(feature = "adaptive-metadata-in-dev")]
fn checkpoint_action_union_element(field_name: &str, value: Scalar) -> DeltaResult<Scalar> {
let fields: Vec<StructField> = CHECKPOINT_ACTION_ELEMENT_SCHEMA.fields().cloned().collect();
require!(
fields.iter().any(|f| f.name() == field_name),
Error::generic(format!(
"checkpoint union element field {field_name:?} not found in element schema"
))
);
let values = fields
.iter()
.map(|field| {
if field.name() == field_name {
value.clone()
} else {
Scalar::null(field.data_type().clone())
}
})
.collect();
Ok(Scalar::Struct(StructData::try_new(fields, values)?))
}
#[cfg(feature = "adaptive-metadata-in-dev")]
impl IntoEngineData for CheckpointAction {
fn into_engine_data(
self,
schema: SchemaRef,
engine: &dyn Engine,
) -> DeltaResult<Box<dyn EngineData>> {
self.validate()?;
let checkpoint_metadata = CheckpointMetadata {
version: self.version,
tags: None,
};
let mut elements = vec![
checkpoint_action_union_element(CHECKPOINT_METADATA_NAME, checkpoint_metadata.into())?,
checkpoint_action_union_element(CONTENT_ROOT_NAME, self.content_root.into())?,
checkpoint_action_union_element(PROTOCOL_NAME, self.protocol.into())?,
checkpoint_action_union_element(METADATA_NAME, self.metadata.into())?,
];
for txn in self.transactions {
elements.push(checkpoint_action_union_element(
SET_TRANSACTION_NAME,
txn.into(),
)?);
}
for dm in self.domain_metadata {
elements.push(checkpoint_action_union_element(
DOMAIN_METADATA_NAME,
dm.into(),
)?);
}
for sidecar in self.txn_sidecars {
let element = content_sidecar_element(SET_TRANSACTION_NAME, sidecar)?;
elements.push(checkpoint_action_union_element(SIDECAR_NAME, element)?);
}
for sidecar in self.domain_metadata_sidecars {
let element = content_sidecar_element(DOMAIN_METADATA_NAME, sidecar)?;
elements.push(checkpoint_action_union_element(SIDECAR_NAME, element)?);
}
let array_type = ArrayType::new(CHECKPOINT_ACTION_ELEMENT_SCHEMA.clone(), false);
let array = Scalar::Array(ArrayData::try_new(array_type, elements)?);
engine.evaluation_handler().create_many(schema, &[&[array]])
}
}
#[cfg(feature = "adaptive-metadata-in-dev")]
fn has_scheme(location: &str) -> bool {
for (position, ch) in location.char_indices() {
if ch == ':' {
return position > 0;
}
if !is_scheme_char(ch, position) {
return false;
}
}
false
}
#[cfg(feature = "adaptive-metadata-in-dev")]
fn is_scheme_char(ch: char, position: usize) -> bool {
if ch.is_ascii_alphabetic() {
return true;
}
position > 0 && (ch.is_ascii_digit() || ch == '+' || ch == '-' || ch == '.')
}
#[cfg(feature = "adaptive-metadata-in-dev")]
impl ContentRoot {
#[internal_api]
pub(crate) fn to_filemeta(&self, table_root: &Url) -> DeltaResult<FileMeta> {
let path = &self.path;
let location = if has_scheme(path) {
Url::parse(path).map_err(|e| {
Error::generic(format!(
"Failed to parse absolute checkpoint contentRoot path {path:?}: {e}"
))
})?
} else {
let mut base = table_root.as_str().to_string();
if !base.ends_with('/') {
base.push('/');
}
Url::parse(&format!("{base}{path}")).map_err(|e| {
Error::generic(format!(
"Failed to resolve checkpoint contentRoot path {path:?} against table \
root {base}: {e}"
))
})?
};
Ok(FileMeta {
location,
last_modified: i64::MAX,
size: to_file_size(self.size_in_bytes, "checkpoint contentRoot")?,
})
}
}
#[cfg(feature = "adaptive-metadata-in-dev")]
impl CheckpointAction {
#[internal_api]
pub(crate) fn try_new_from_data(
data: &dyn EngineData,
) -> DeltaResult<Option<CheckpointAction>> {
let mut visitor = visitors::CheckpointVisitor::default();
visitor.visit_rows_of(data)?;
Ok(visitor.checkpoint)
}
fn validate(&self) -> DeltaResult<()> {
require!(
self.content_root.version <= self.version,
Error::generic(format!(
"checkpoint contentRoot.version {} exceeds checkpointMetadata.version {}",
self.content_root.version, self.version
))
);
Ok(())
}
#[internal_api]
pub(crate) fn path(&self) -> &str {
&self.content_root.path
}
#[internal_api]
pub(crate) fn version(&self) -> i64 {
self.version
}
#[internal_api]
pub(crate) fn root_filemeta(&self, table_root: &Url) -> DeltaResult<FileMeta> {
self.content_root.to_filemeta(table_root)
}
#[internal_api]
pub(crate) fn protocol(&self) -> &Protocol {
&self.protocol
}
#[internal_api]
pub(crate) fn metadata(&self) -> &Metadata {
&self.metadata
}
}
#[derive(ToSchema, IntoStructData, Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[internal_api]
pub(crate) struct Sidecar {
pub path: String,
pub size_in_bytes: i64,
pub modification_time: i64,
#[allow_null_container_values]
pub tags: Option<HashMap<String, String>>,
}
fn to_file_size(bytes: i64, context: &str) -> DeltaResult<FileSize> {
bytes.try_into().map_err(|_| {
Error::generic(format!(
"Failed to convert {context} size {bytes} to FileSize"
))
})
}
impl Sidecar {
pub(crate) fn to_filemeta(&self, log_root: &Url) -> DeltaResult<FileMeta> {
Ok(FileMeta {
location: log_root.join("_sidecars/")?.join(&self.path)?,
last_modified: self.modification_time,
size: to_file_size(self.size_in_bytes, "sidecar")?,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, ToSchema, IntoStructData, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[internal_api]
pub(crate) struct CheckpointMetadata {
pub(crate) version: i64,
#[allow_null_container_values]
pub(crate) tags: Option<HashMap<String, String>>,
}
#[derive(
Debug, Clone, PartialEq, Eq, Serialize, Deserialize, ToSchema, IntoStructData, IntoEngineData,
)]
pub struct DomainMetadata {
domain: String,
configuration: String,
removed: bool,
}
impl DomainMetadata {
pub(crate) fn new(domain: String, configuration: String) -> Self {
Self {
domain,
configuration,
removed: false,
}
}
pub(crate) fn remove(domain: String, configuration: String) -> Self {
Self {
domain,
configuration,
removed: true,
}
}
#[allow(unused)]
#[internal_api]
pub(crate) fn is_internal(&self) -> bool {
self.domain.starts_with(INTERNAL_DOMAIN_PREFIX)
}
pub fn domain(&self) -> &str {
&self.domain
}
pub fn configuration(&self) -> &str {
&self.configuration
}
pub fn is_removed(&self) -> bool {
self.removed
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use rstest::rstest;
use serde_json::json;
use super::*;
use crate::arrow::array::{
Array, BooleanArray, Int32Array, Int64Array, ListArray, ListBuilder, MapBuilder,
MapFieldNames, RecordBatch, StringArray, StringBuilder, StructArray,
};
use crate::arrow::datatypes::{DataType as ArrowDataType, Field, Schema};
use crate::arrow::json::ReaderBuilder;
use crate::engine::arrow_data::EngineDataArrowExt as _;
use crate::engine::arrow_expression::ArrowEvaluationHandler;
use crate::expressions::Scalar;
use crate::schema::{schema, schema_ref, DataType, MapType, StructField};
use crate::unit_test_utils::assert_result_error_with_message;
use crate::{
Engine, EvaluationHandler, IntoEngineData, JsonHandler, ParquetHandler, StorageHandler,
};
#[rstest]
#[case::add(ADD_NAME, Some("path"))]
#[case::remove(REMOVE_NAME, Some("path"))]
#[case::metadata(METADATA_NAME, Some("id"))]
#[case::protocol(PROTOCOL_NAME, Some("minReaderVersion"))]
#[case::transaction(SET_TRANSACTION_NAME, Some("appId"))]
#[case::cdc(CDC_NAME, Some("path"))]
#[case::domain_metadata(DOMAIN_METADATA_NAME, Some("domain"))]
#[case::checkpoint_metadata(CHECKPOINT_METADATA_NAME, Some("version"))]
#[case::sidecar(SIDECAR_NAME, Some("path"))]
#[case::witnessless_action(COMMIT_INFO_NAME, None)]
#[case::unknown_action("futureAction", None)]
fn test_action_presence_leaf(#[case] action_name: &str, #[case] expected_leaf: Option<&str>) {
assert_eq!(action_presence_leaf(action_name), expected_leaf);
}
struct ExprEngine(Arc<dyn EvaluationHandler>);
impl ExprEngine {
fn new() -> Self {
ExprEngine(Arc::new(ArrowEvaluationHandler))
}
}
impl Engine for ExprEngine {
fn evaluation_handler(&self) -> Arc<dyn EvaluationHandler> {
self.0.clone()
}
fn json_handler(&self) -> Arc<dyn JsonHandler> {
unimplemented!()
}
fn parquet_handler(&self) -> Arc<dyn ParquetHandler> {
unimplemented!()
}
fn storage_handler(&self) -> Arc<dyn StorageHandler> {
unimplemented!()
}
}
fn create_string_map_builder(
nullable_values: bool,
) -> MapBuilder<StringBuilder, StringBuilder> {
MapBuilder::new(
Some(MapFieldNames {
entry: "key_value".to_string(),
key: "key".to_string(),
value: "value".to_string(),
}),
StringBuilder::new(),
StringBuilder::new(),
)
.with_values_field(Field::new(
"value".to_string(),
ArrowDataType::Utf8,
nullable_values,
))
}
#[rstest]
#[case::no_expiration_configured(None, Some(1000), false)]
#[case::null_last_updated_never_expires(Some(5000), None, false)]
#[case::both_none(None, None, false)]
#[case::last_updated_before_expiration(Some(2000), Some(1000), true)]
#[case::last_updated_at_expiration(Some(1000), Some(1000), true)]
#[case::last_updated_after_expiration(Some(2000), Some(3000), false)]
fn test_set_transaction_expiration(
#[case] expiration_timestamp: Option<i64>,
#[case] last_updated: Option<i64>,
#[case] expired: bool,
) {
let txn = SetTransaction::new("app".to_string(), 7, last_updated);
assert_eq!(txn.is_expired(expiration_timestamp), expired);
assert_eq!(
txn.non_expired_version(expiration_timestamp),
(!expired).then_some(7)
);
}
#[test]
fn test_metadata_schema() {
let schema = get_commit_schema()
.project(&[METADATA_NAME])
.expect("Couldn't get metaData field");
let expected = schema_ref! {
nullable "metaData": {
not_null "id": STRING,
nullable "name": STRING,
nullable "description": STRING,
not_null "format": {
not_null "provider": STRING,
not_null "options": { STRING => not_null STRING },
},
not_null "schemaString": STRING,
not_null "partitionColumns": [ not_null STRING ],
nullable "createdTime": LONG,
not_null "configuration": { STRING => not_null STRING },
},
};
assert_eq!(schema, expected);
}
#[rstest]
#[case::supported(41, false)]
#[case::exceeded(42, true)]
fn parse_schema_nesting_boundary(#[case] depth: usize, #[case] exceeds_limit: bool) {
let metadata = Metadata {
schema_string: serde_json::to_string(&nested_schema(depth)).unwrap(),
..Default::default()
};
let result = metadata.parse_schema();
if exceeds_limit {
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(_)));
} else {
result.unwrap();
}
}
#[rstest]
#[case::malformed_syntax("{", "MalformedJson")]
#[case::malformed_bad_decimal(
r#"{"type":"struct","fields":[{"name":"t","type":"decimal(nope)","nullable":true,"metadata":{}}]}"#,
"MalformedJson"
)]
#[case::malformed_value_echoes_prefix(
r#"{"type":"struct","fields":[{"name":"t","type":"string","nullable":"Unsupported Delta table type","metadata":{}}]}"#,
"MalformedJson"
)]
#[case::unsupported_time(
r#"{"type":"struct","fields":[{"name":"t","type":"time(6)","nullable":true,"metadata":{}}]}"#,
"Schema"
)]
#[case::unsupported_interval(
r#"{"type":"struct","fields":[{"name":"t","type":"interval week","nullable":true,"metadata":{}}]}"#,
"Schema"
)]
#[case::unsupported_nested(
r#"{"type":"struct","fields":[{"name":"t","type":{"type":"struct","fields":[{"name":"inner","type":"time(6)","nullable":true,"metadata":{}}]},"nullable":true,"metadata":{}}]}"#,
"Schema"
)]
#[case::unsupported_complex(
r#"{"type":"struct","fields":[{"name":"t","type":{"type":"matrix"},"nullable":true,"metadata":{}}]}"#,
"Schema"
)]
fn parse_schema_error_classification(
#[case] schema_string: &str,
#[case] expected_error: &str,
) {
let metadata = Metadata {
schema_string: schema_string.to_string(),
..Default::default()
};
let error = match metadata.parse_schema().unwrap_err() {
Error::Backtraced { source, .. } => *source,
error => error,
};
match expected_error {
"MalformedJson" => {
assert!(matches!(error, Error::MalformedJson(_)), "got: {error:?}")
}
"Schema" => {
assert!(matches!(error, Error::Schema(_)), "got: {error:?}")
}
other => panic!("unknown expected_error discriminant: {other}"),
}
}
fn nested_schema(depth: usize) -> StructType {
(0..depth).fold(
schema! { nullable "leaf": INTEGER },
|nested, depth| schema! { nullable (format!("level_{depth}")): (nested) },
)
}
#[test]
fn test_add_schema() {
let schema = get_commit_schema()
.project(&[ADD_NAME])
.expect("Couldn't get add field");
let expected = schema_ref! {
nullable "add": {
not_null "path": STRING,
not_null "partitionValues": { STRING => nullable STRING },
not_null "size": LONG,
not_null "modificationTime": LONG,
not_null "dataChange": BOOLEAN,
nullable "stats": STRING,
nullable "tags": { STRING => nullable STRING },
(deletion_vector_field()),
nullable "baseRowId": LONG,
nullable "defaultRowCommitVersion": LONG,
nullable "clusteringProvider": STRING,
},
};
assert_eq!(schema, expected);
}
fn tags_field() -> StructField {
StructField::nullable(
"tags",
MapType::new(DataType::STRING, DataType::STRING, true),
)
}
fn partition_values_field() -> StructField {
StructField::nullable(
"partitionValues",
MapType::new(DataType::STRING, DataType::STRING, true),
)
}
fn deletion_vector_field() -> StructField {
StructField::nullable(
"deletionVector",
schema! {
not_null "storageType": STRING,
not_null "pathOrInlineDv": STRING,
nullable "offset": INTEGER,
not_null "sizeInBytes": INTEGER,
not_null "cardinality": LONG,
},
)
}
#[test]
fn test_remove_schema() {
let schema = get_commit_schema()
.project(&[REMOVE_NAME])
.expect("Couldn't get remove field");
let expected = schema_ref! {
nullable "remove": {
not_null "path": STRING,
nullable "deletionTimestamp": LONG,
not_null "dataChange": BOOLEAN,
nullable "extendedFileMetadata": BOOLEAN,
(partition_values_field()),
nullable "size": LONG,
nullable "stats": STRING,
(tags_field()),
(deletion_vector_field()),
nullable "baseRowId": LONG,
nullable "defaultRowCommitVersion": LONG,
},
};
assert_eq!(schema, expected);
}
#[test]
fn test_cdc_schema() {
let schema = get_commit_schema()
.project(&[CDC_NAME])
.expect("Couldn't get cdc field");
let expected = schema_ref! {
nullable "cdc": {
not_null "path": STRING,
not_null "partitionValues": { STRING => nullable STRING },
not_null "size": LONG,
not_null "dataChange": BOOLEAN,
(tags_field()),
},
};
assert_eq!(schema, expected);
}
#[test]
fn test_sidecar_schema() {
let schema = Sidecar::to_schema();
let expected = schema! {
not_null "path": STRING,
not_null "sizeInBytes": LONG,
not_null "modificationTime": LONG,
(tags_field()),
};
assert_eq!(schema, expected);
}
#[test]
fn test_checkpoint_metadata_schema() {
let schema = get_all_actions_schema()
.project(&[CHECKPOINT_METADATA_NAME])
.expect("Couldn't get checkpointMetadata field");
let expected = schema_ref! {
nullable "checkpointMetadata": {
not_null "version": LONG,
(tags_field()),
},
};
assert_eq!(schema, expected);
}
#[test]
fn test_transaction_schema() {
let schema = get_commit_schema()
.project(&["txn"])
.expect("Couldn't get transaction field");
let expected = schema_ref! {
nullable "txn": {
not_null "appId": STRING,
not_null "version": LONG,
nullable "lastUpdated": LONG,
},
};
assert_eq!(schema, expected);
}
#[test]
fn test_commit_info_schema() {
let schema = get_commit_schema()
.project(&["commitInfo"])
.expect("Couldn't get commitInfo field");
let expected = schema_ref! {
nullable "commitInfo": {
nullable "timestamp": LONG,
nullable "inCommitTimestamp": LONG,
nullable "operation": STRING,
nullable "operationParameters": { STRING => nullable STRING },
nullable "operationMetrics": { STRING => nullable STRING },
nullable "kernelVersion": STRING,
nullable "isBlindAppend": BOOLEAN,
nullable "engineInfo": STRING,
nullable "txnId": STRING,
},
};
assert_eq!(schema, expected);
}
#[test]
fn test_domain_metadata_schema() {
let schema = get_commit_schema()
.project(&[DOMAIN_METADATA_NAME])
.expect("Couldn't get domainMetadata field");
let expected = schema_ref! {
nullable "domainMetadata": {
not_null "domain": STRING,
not_null "configuration": STRING,
not_null "removed": BOOLEAN,
},
};
assert_eq!(schema, expected);
}
#[test]
fn test_validate_protocol() {
let invalid_protocols = [
Protocol {
min_reader_version: 3,
min_writer_version: 7,
reader_features: None,
writer_features: Some(vec![]),
},
Protocol {
min_reader_version: 3,
min_writer_version: 7,
reader_features: Some(vec![]),
writer_features: None,
},
Protocol {
min_reader_version: 3,
min_writer_version: 7,
reader_features: None,
writer_features: None,
},
];
for Protocol {
min_reader_version,
min_writer_version,
reader_features,
writer_features,
} in invalid_protocols
{
assert!(matches!(
Protocol::try_new(
min_reader_version,
min_writer_version,
reader_features,
writer_features
),
Err(Error::InvalidProtocol(_)),
));
}
}
#[rstest]
#[case(0, 1)]
#[case(1, 0)]
#[case(-1, 2)]
#[case(1, -1)]
fn reject_protocol_version_below_minimum(#[case] rv: i32, #[case] wv: i32) {
let expected = if rv < 1 {
format!("Invalid protocol action in the delta log: min_reader_version must be >= 1, got {rv}")
} else {
format!("Invalid protocol action in the delta log: min_writer_version must be >= 1, got {wv}")
};
assert_result_error_with_message(
Protocol::try_new(rv, wv, TableFeature::NO_LIST, TableFeature::NO_LIST),
&expected,
);
}
#[test]
fn accept_min_versions() {
let p = Protocol::try_new_legacy(1, 1).unwrap();
assert_eq!(p.min_reader_version(), 1);
assert_eq!(p.min_writer_version(), 1);
}
#[test]
fn test_validate_table_features_invalid() {
let invalid_features = [
(
vec![TableFeature::DeletionVectors],
vec![TableFeature::AppendOnly],
"Reader features must contain only ReaderWriter features that are also listed in writer features",
),
(
vec![TableFeature::DeletionVectors],
vec![],
"Reader features must contain only ReaderWriter features that are also listed in writer features",
),
(
vec![],
vec![TableFeature::DeletionVectors],
"Writer features must be Writer-only or also listed in reader features",
),
(
vec![TableFeature::VariantType],
vec![
TableFeature::VariantType,
TableFeature::DeletionVectors,
],
"Writer features must be Writer-only or also listed in reader features",
),
(
vec![TableFeature::AppendOnly],
vec![TableFeature::AppendOnly],
"Reader features must contain only ReaderWriter features that are also listed in writer features",
),
];
for (reader_features, writer_features, error_msg) in invalid_features {
let res = Protocol::try_new_modern(reader_features, writer_features);
assert!(
matches!(
&res,
Err(Error::InvalidProtocol(error)) if error.to_string().contains(error_msg)
),
"Expected message containing:\t{error_msg}\nBut got:{res:?}\n"
);
}
}
#[test]
fn test_validate_table_features_unknown() {
let protocol = Protocol::try_new_modern(
vec![TableFeature::Unknown("unknown_reader".to_string())],
vec![TableFeature::Unknown("unknown_reader".to_string())],
);
assert!(protocol.is_ok());
let protocol = Protocol::try_new_modern(
TableFeature::EMPTY_LIST,
vec![TableFeature::Unknown("unknown_writer".to_string())],
);
assert!(protocol.is_ok());
}
#[test]
fn test_validate_table_features_valid() {
let valid_features = [
(
vec![TableFeature::DeletionVectors],
vec![TableFeature::DeletionVectors],
),
(vec![], vec![TableFeature::AppendOnly]),
(
vec![TableFeature::VariantType],
vec![TableFeature::VariantType, TableFeature::AppendOnly],
),
(
vec![TableFeature::Unknown("rw".to_string())],
vec![
TableFeature::Unknown("rw".to_string()),
TableFeature::Unknown("w".to_string()),
],
),
(vec![], vec![]),
];
for (reader_features, writer_features) in valid_features {
assert!(Protocol::try_new_modern(reader_features, writer_features).is_ok());
}
}
#[test]
fn test_validate_legacy_column_mapping_valid() {
let protocol = Protocol::try_new(
2,
7,
TableFeature::NO_LIST,
Some(vec![TableFeature::ColumnMapping]),
);
assert!(protocol.is_ok());
}
#[test]
fn test_validate_legacy_writer_only_features_valid() {
let protocol = Protocol::try_new(
1,
7,
TableFeature::NO_LIST,
Some(vec![TableFeature::AppendOnly]),
);
assert!(protocol.is_ok());
}
#[test]
fn test_validate_legacy_column_mapping_with_writer_features_valid() {
let protocol = Protocol::try_new(
2,
7,
TableFeature::NO_LIST,
Some(vec![TableFeature::AppendOnly, TableFeature::ColumnMapping]),
);
assert!(protocol.is_ok());
}
#[test]
fn test_validate_column_mapping_reader_v1_invalid() {
let protocol = Protocol::try_new(
1,
7,
TableFeature::NO_LIST,
Some(vec![TableFeature::ColumnMapping]),
);
assert!(protocol.is_err());
}
#[test]
fn test_validate_multiple_readerwriter_features_reader_v2_invalid() {
let protocol = Protocol::try_new(
2,
7,
TableFeature::NO_LIST,
Some(vec![
TableFeature::ColumnMapping,
TableFeature::DeletionVectors,
]),
);
assert!(protocol.is_err());
}
#[test]
fn test_parse_table_feature_never_fails() {
let features = Some(["", "absurD_)(+13%^⚙️"]);
let expected = Some(FromIterator::from_iter([
TableFeature::unknown(""),
TableFeature::unknown("absurD_)(+13%^⚙️"),
]));
assert_eq!(parse_features(features), expected);
}
#[test]
fn test_into_engine_data() {
let engine = ExprEngine::new();
let set_transaction = SetTransaction {
app_id: "app_id".to_string(),
version: 0,
last_updated: None,
};
let engine_data =
set_transaction.into_engine_data(SetTransaction::to_schema().into(), &engine);
let record_batch = engine_data.try_into_record_batch().unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new("appId", ArrowDataType::Utf8, false),
Field::new("version", ArrowDataType::Int64, false),
Field::new("lastUpdated", ArrowDataType::Int64, true),
]));
let expected = RecordBatch::try_new(
schema,
vec![
Arc::new(StringArray::from(vec!["app_id"])),
Arc::new(Int64Array::from(vec![0_i64])),
Arc::new(Int64Array::from(vec![None::<i64>])),
],
)
.unwrap();
assert_eq!(record_batch, expected);
}
#[test]
fn test_commit_info_into_engine_data() {
let engine = ExprEngine::new();
let commit_info = CommitInfo::new(0, None, None, None, false);
let commit_info_txn_id = commit_info.txn_id.clone();
let engine_data = commit_info.into_engine_data(CommitInfo::to_schema().into(), &engine);
let record_batch = engine_data.try_into_record_batch().unwrap();
let mut map_builder = create_string_map_builder(true);
map_builder.append(true).unwrap();
let operation_parameters = Arc::new(map_builder.finish());
let mut map_builder = create_string_map_builder(true);
map_builder.append(false).unwrap();
let operation_metrics = Arc::new(map_builder.finish());
let expected = RecordBatch::try_new(
record_batch.schema(),
vec![
Arc::new(Int64Array::from(vec![Some(0)])),
Arc::new(Int64Array::from(vec![None::<i64>])),
Arc::new(StringArray::from(vec![Some("UNKNOWN")])),
operation_parameters,
operation_metrics,
Arc::new(StringArray::from(vec![Some(format!("v{KERNEL_VERSION}"))])),
Arc::new(BooleanArray::from(vec![None::<bool>])),
Arc::new(StringArray::from(vec![None::<String>])),
Arc::new(StringArray::from(vec![commit_info_txn_id])),
],
)
.unwrap();
assert_eq!(record_batch, expected);
}
#[test]
fn test_domain_metadata_into_engine_data() {
let engine = ExprEngine::new();
let domain_metadata = DomainMetadata {
domain: "my.domain".to_string(),
configuration: "config_value".to_string(),
removed: false,
};
let engine_data =
domain_metadata.into_engine_data(DomainMetadata::to_schema().into(), &engine);
let record_batch = engine_data.try_into_record_batch().unwrap();
let expected = RecordBatch::try_new(
record_batch.schema(),
vec![
Arc::new(StringArray::from(vec!["my.domain"])),
Arc::new(StringArray::from(vec!["config_value"])),
Arc::new(BooleanArray::from(vec![false])),
],
)
.unwrap();
assert_eq!(record_batch, expected);
}
#[test]
fn test_metadata_try_new() {
let schema = schema_ref! { not_null "id": INTEGER };
let config = HashMap::from([("key1".to_string(), "value1".to_string())]);
let metadata = Metadata::try_new(
Some("test_table".to_string()),
Some("description".to_string()),
schema.clone(),
vec!["year".to_string()],
1234567890,
config.clone(),
)
.unwrap();
assert!(!metadata.id.is_empty());
assert_eq!(metadata.name, Some("test_table".to_string()));
assert_eq!(
metadata.schema_string,
serde_json::to_string(&schema).unwrap()
);
assert_eq!(metadata.created_time, Some(1234567890));
assert_eq!(metadata.configuration, config);
}
#[test]
fn test_metadata_try_new_default() {
let schema = schema_ref! { not_null "id": INTEGER };
let metadata = Metadata::try_new(None, None, schema, vec![], 0, HashMap::new()).unwrap();
assert!(!metadata.id.is_empty());
assert_eq!(metadata.name, None);
assert_eq!(metadata.description, None);
}
#[test]
fn test_metadata_unique_ids() {
let schema = schema_ref! { not_null "id": INTEGER };
let m1 = Metadata::try_new(None, None, schema.clone(), vec![], 0, HashMap::new()).unwrap();
let m2 = Metadata::try_new(None, None, schema, vec![], 0, HashMap::new()).unwrap();
assert_ne!(m1.id, m2.id);
}
#[rstest]
#[case::typical(HashMap::from([
("path".to_string(), "/delta/table".to_string()),
("compressionType".to_string(), "snappy".to_string()),
]))]
#[case::empty(HashMap::new())]
#[case::special_characters(HashMap::from([
("path".to_string(), "/path/with spaces".to_string()),
("unicode".to_string(), "测试🎉".to_string()),
("empty".to_string(), String::new()),
]))]
fn test_format_scalar_round_trip(#[case] options: HashMap<String, String>) {
let format = Format {
provider: "parquet".to_string(),
options: options.clone(),
};
let scalar = Scalar::from(format.clone());
let Scalar::Struct(struct_data) = &scalar else {
panic!("Expected struct scalar, got {scalar}");
};
let field_names: Vec<_> = struct_data.fields().iter().map(|f| f.name()).collect();
assert_eq!(field_names, ["provider", "options"]);
assert_eq!(struct_data.values()[0], Scalar::from("parquet"));
let Scalar::Map(map_data) = &struct_data.values()[1] else {
panic!("Expected map options");
};
assert_eq!(map_data.pairs().len(), options.len());
assert_eq!(Format::try_from(scalar).unwrap(), format);
}
#[test]
fn test_format_default() {
let format = Format::default();
let expected = Format {
provider: "parquet".to_string(),
options: HashMap::new(),
};
assert_eq!(format, expected);
}
#[test]
fn test_metadata_into_engine_data() {
let engine = ExprEngine::new();
let schema = schema_ref! { not_null "id": INTEGER };
let test_metadata = Metadata::try_new(
Some("test".to_string()),
Some("my table".to_string()),
schema.clone(),
vec!["part".to_string()],
123,
HashMap::from([("k".to_string(), "v".to_string())]),
)
.unwrap();
let test_id = test_metadata.id.clone();
let actual = test_metadata
.into_engine_data(Metadata::to_schema().into(), &engine)
.unwrap()
.try_into_record_batch()
.unwrap();
let expected_json = json!({
"id": test_id,
"name": "test",
"description": "my table",
"format": {
"provider": "parquet",
"options": {}
},
"schemaString": "{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"integer\",\"nullable\":false,\"metadata\":{}}]}",
"partitionColumns": ["part"],
"createdTime": 123,
"configuration": {
"k": "v"
}
}).to_string();
let expected = ReaderBuilder::new(actual.schema())
.build(expected_json.as_bytes())
.unwrap()
.next()
.unwrap()
.unwrap();
assert_eq!(actual, expected);
}
#[test]
fn test_metadata_with_log_schema() {
let engine = ExprEngine::new();
let schema = schema_ref! { not_null "id": INTEGER };
let metadata = Metadata::try_new(
Some("table".to_string()),
None, schema,
vec![],
456,
HashMap::new(),
)
.unwrap();
let metadata_id = metadata.id.clone();
let commit_schema = LOG_METADATA_SCHEMA.clone();
let actual = metadata
.into_engine_data(commit_schema, &engine)
.unwrap()
.try_into_record_batch()
.unwrap();
let expected_json = json!({
"metaData": {
"id": metadata_id,
"name": "table",
"format": {
"provider": "parquet",
"options": {}
},
"schemaString": "{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"integer\",\"nullable\":false,\"metadata\":{}}]}",
"partitionColumns": [],
"createdTime": 456,
"configuration": {}
}
}).to_string();
let expected = ReaderBuilder::new(actual.schema())
.build(expected_json.as_bytes())
.unwrap()
.next()
.unwrap()
.unwrap();
assert_eq!(actual, expected);
}
#[test]
fn test_protocol_into_engine_data() {
let engine = ExprEngine::new();
let protocol = Protocol::try_new_modern(
[TableFeature::DeletionVectors, TableFeature::ColumnMapping],
[TableFeature::DeletionVectors, TableFeature::ColumnMapping],
)
.unwrap();
let engine_data = protocol
.clone()
.into_engine_data(Protocol::to_schema().into(), &engine);
let record_batch = engine_data.try_into_record_batch().unwrap();
let list_field = Arc::new(Field::new("element", ArrowDataType::Utf8, false));
let protocol_fields = vec![
Field::new("minReaderVersion", ArrowDataType::Int32, false),
Field::new("minWriterVersion", ArrowDataType::Int32, false),
Field::new(
"readerFeatures",
ArrowDataType::List(list_field.clone()),
true, ),
Field::new(
"writerFeatures",
ArrowDataType::List(list_field.clone()),
true, ),
];
let schema = Arc::new(Schema::new(protocol_fields.clone()));
let string_builder = StringBuilder::new();
let mut list_builder = ListBuilder::new(string_builder).with_field(list_field.clone());
list_builder.values().append_value("deletionVectors");
list_builder.values().append_value("columnMapping");
list_builder.append(true);
let reader_features_array = list_builder.finish();
let string_builder = StringBuilder::new();
let mut list_builder = ListBuilder::new(string_builder).with_field(list_field.clone());
list_builder.values().append_value("deletionVectors");
list_builder.values().append_value("columnMapping");
list_builder.append(true);
let writer_features_array = list_builder.finish();
let expected = RecordBatch::try_new(
schema,
vec![
Arc::new(Int32Array::from(vec![3])),
Arc::new(Int32Array::from(vec![7])),
Arc::new(reader_features_array.clone()),
Arc::new(writer_features_array.clone()),
],
)
.unwrap();
assert_eq!(record_batch, expected);
let commit_schema = LOG_PROTOCOL_SCHEMA.clone();
let engine_data = protocol.into_engine_data(commit_schema, &engine);
let schema = Arc::new(Schema::new(vec![Field::new(
"protocol",
ArrowDataType::Struct(protocol_fields.into()),
true,
)]));
let expected = RecordBatch::try_new(
schema,
vec![Arc::new(StructArray::from(vec![
(
Arc::new(Field::new("minReaderVersion", ArrowDataType::Int32, false)),
Arc::new(Int32Array::from(vec![3])) as Arc<dyn Array>,
),
(
Arc::new(Field::new("minWriterVersion", ArrowDataType::Int32, false)),
Arc::new(Int32Array::from(vec![7])) as Arc<dyn Array>,
),
(
Arc::new(Field::new(
"readerFeatures",
ArrowDataType::List(list_field.clone()),
true,
)),
Arc::new(reader_features_array) as Arc<dyn Array>,
),
(
Arc::new(Field::new(
"writerFeatures",
ArrowDataType::List(list_field),
true,
)),
Arc::new(writer_features_array) as Arc<dyn Array>,
),
]))],
)
.unwrap();
let record_batch = engine_data.try_into_record_batch().unwrap();
assert_eq!(record_batch, expected);
}
#[test]
fn test_protocol_into_engine_data_empty_features() {
let engine = ExprEngine::new();
let protocol =
Protocol::try_new_modern(TableFeature::EMPTY_LIST, TableFeature::EMPTY_LIST).unwrap();
let engine_data = protocol
.into_engine_data(Protocol::to_schema().into(), &engine)
.unwrap();
let record_batch = engine_data.try_into_record_batch().unwrap();
assert_eq!(record_batch.num_rows(), 1);
assert_eq!(record_batch.num_columns(), 4);
let reader_features_col = record_batch
.column(2)
.as_any()
.downcast_ref::<ListArray>()
.unwrap();
assert_eq!(reader_features_col.len(), 1);
assert_eq!(reader_features_col.value(0).len(), 0); let writer_features_col = record_batch
.column(3)
.as_any()
.downcast_ref::<ListArray>()
.unwrap();
assert_eq!(writer_features_col.len(), 1);
assert_eq!(writer_features_col.value(0).len(), 0); }
#[test]
fn test_protocol_into_engine_data_no_features() {
let engine = ExprEngine::new();
let protocol = Protocol::try_new_legacy(1, 2).unwrap();
let engine_data = protocol
.into_engine_data(Protocol::to_schema().into(), &engine)
.unwrap();
let record_batch = engine_data.try_into_record_batch().unwrap();
assert_eq!(record_batch.num_rows(), 1);
assert_eq!(record_batch.num_columns(), 4);
assert!(record_batch.column(2).is_null(0));
assert!(record_batch.column(3).is_null(0));
}
#[test]
fn test_schema_contains_file_actions_with_add() {
let schema = get_commit_schema()
.project(&[ADD_NAME, PROTOCOL_NAME])
.unwrap();
assert!(schema_contains_file_actions(&schema));
assert!(schema_contains_file_actions(
&schema.project(&[ADD_NAME]).unwrap()
));
}
#[test]
fn test_schema_contains_file_actions_with_remove() {
let schema = get_commit_schema()
.project(&[REMOVE_NAME, METADATA_NAME])
.unwrap();
assert!(schema_contains_file_actions(&schema));
assert!(schema_contains_file_actions(
&schema.project(&[REMOVE_NAME]).unwrap()
));
}
#[test]
fn test_schema_contains_file_actions_with_both() {
let schema = get_commit_schema()
.project(&[ADD_NAME, REMOVE_NAME])
.unwrap();
assert!(schema_contains_file_actions(&schema));
}
#[test]
fn test_schema_contains_file_actions_with_neither() {
let schema = get_commit_schema()
.project(&[PROTOCOL_NAME, METADATA_NAME])
.unwrap();
assert!(!schema_contains_file_actions(&schema));
}
#[test]
fn test_schema_contains_file_actions_empty_schema() {
let schema = schema_ref! {};
assert!(!schema_contains_file_actions(&schema));
}
#[test]
fn test_add_tags_deserialization_null_case() {
let json1 = r#"{"path":"file1.parquet","partitionValues":{},"size":100,"modificationTime":1234567890,"dataChange":true,"tags":null}"#;
let add1: Add = serde_json::from_str(json1).unwrap();
assert_eq!(add1.tags, None);
}
#[test]
fn test_add_tags_deserialization_nullable_values_case() {
let json2 = r#"{"path":"file2.parquet","partitionValues":{},"size":200,"modificationTime":1234567890,"dataChange":true,"tags":{"INSERTION_TIME":"1677811178336000","NULLABLE_TAG":null}}"#;
let add2: Add = serde_json::from_str(json2).unwrap();
assert!(add2.tags.is_some());
let tags = add2.tags.unwrap();
assert_eq!(tags.len(), 2);
assert_eq!(
tags.get("INSERTION_TIME"),
Some(&Some("1677811178336000".to_string()))
);
assert_eq!(tags.get("NULLABLE_TAG"), Some(&None));
}
#[test]
fn test_add_tags_deserialization_non_null_values_case() {
let json3 = r#"{"path":"file3.parquet","partitionValues":{},"size":300,"modificationTime":1234567890,"dataChange":true,"tags":{"INSERTION_TIME":"1677811178336000","MIN_INSERTION_TIME":"1677811178336000"}}"#;
let add3: Add = serde_json::from_str(json3).unwrap();
assert!(add3.tags.is_some());
let tags = add3.tags.unwrap();
assert_eq!(tags.len(), 2);
assert_eq!(
tags.get("INSERTION_TIME"),
Some(&Some("1677811178336000".to_string()))
);
assert_eq!(
tags.get("MIN_INSERTION_TIME"),
Some(&Some("1677811178336000".to_string()))
);
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_checkpoint_action_schema() {
let schema = get_commit_schema()
.project(&[CHECKPOINT_ACTION_NAME])
.unwrap();
let checkpoint_field = schema.field(CHECKPOINT_ACTION_NAME).unwrap();
assert!(checkpoint_field.is_nullable());
let array = match checkpoint_field.data_type() {
DataType::Array(array) => array,
other => panic!("Expected array, got {other:?}"),
};
assert!(!array.contains_null());
let element = match array.element_type() {
DataType::Struct(s) => s,
other => panic!("Expected struct element, got {other:?}"),
};
let field_names: Vec<&str> = element.fields().map(|f| f.name.as_str()).collect();
assert_eq!(
field_names,
vec![
CHECKPOINT_METADATA_NAME,
CONTENT_ROOT_NAME,
PROTOCOL_NAME,
METADATA_NAME,
DOMAIN_METADATA_NAME,
SET_TRANSACTION_NAME,
SIDECAR_NAME,
]
);
assert!(!field_names.contains(&COMMIT_INFO_NAME));
for field in element.fields() {
assert!(field.is_nullable(), "{} should be nullable", field.name);
assert!(matches!(field.data_type(), DataType::Struct(_)));
}
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_content_sidecar_is_type_prefixed_sidecar() {
let DataType::Struct(content_sidecar) = CONTENT_SIDECAR_FIELD.data_type() else {
panic!("content sidecar should be a struct");
};
let expected: Vec<StructField> =
std::iter::once(StructField::not_null("type", DataType::STRING))
.chain(Sidecar::to_schema().into_fields())
.collect();
let actual: Vec<StructField> = content_sidecar.fields().cloned().collect();
assert_eq!(actual, expected);
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[rstest]
#[case::relative_path(
"memory:///table/",
"metadata/root.parquet",
2048,
"memory:///table/metadata/root.parquet",
Ok(2048)
)]
#[case::absolute_path(
"memory:///table/",
"s3://bucket/table/metadata/root.parquet",
2048,
"s3://bucket/table/metadata/root.parquet",
Ok(2048)
)]
#[case::negative_size(
"memory:///table/",
"metadata/root.parquet",
-1,
"memory:///table/metadata/root.parquet",
Err("Failed to convert checkpoint contentRoot size -1")
)]
#[case::table_root_without_trailing_slash_gets_one(
"memory:///table",
"metadata/root.parquet",
2048,
"memory:///table/metadata/root.parquet",
Ok(2048)
)]
#[case::single_char_scheme_treated_as_absolute(
"memory:///table/",
"c:/foo/root.parquet",
2048,
"c:/foo/root.parquet",
Ok(2048)
)]
#[case::colon_in_relative_segment_stays_relative(
"memory:///table/",
"metadata/snap-123:456.parquet",
2048,
"memory:///table/metadata/snap-123:456.parquet",
Ok(2048)
)]
#[case::leading_digit_scheme_treated_as_relative(
"memory:///table/",
"3com/root.parquet",
2048,
"memory:///table/3com/root.parquet",
Ok(2048)
)]
#[case::non_ascii_scheme_treated_as_relative(
"memory:///table/",
"\u{03b1}scheme/root.parquet",
2048,
"memory:///table/%CE%B1scheme/root.parquet",
Ok(2048)
)]
#[case::compound_scheme_treated_as_absolute(
"memory:///table/",
"git+ssh://host/repo/root.parquet",
2048,
"git+ssh://host/repo/root.parquet",
Ok(2048)
)]
fn test_checkpoint_action_root_filemeta(
#[case] table_root: &str,
#[case] path: &str,
#[case] size_in_bytes: i64,
#[case] expected_location: &str,
#[case] expected: Result<FileSize, &str>,
) {
let table_root = Url::parse(table_root).unwrap();
let checkpoint_action = CheckpointAction {
version: 1,
content_root: ContentRoot {
path: path.to_string(),
size_in_bytes,
version: 1,
},
protocol: Protocol::new_unchecked(1, 2, None, None),
metadata: Metadata::default(),
transactions: Vec::new(),
domain_metadata: Vec::new(),
txn_sidecars: Vec::new(),
domain_metadata_sidecars: Vec::new(),
};
let result = checkpoint_action.root_filemeta(&table_root);
match expected {
Ok(expected_size) => {
let file_meta = result.unwrap();
assert_eq!(file_meta.location.as_str(), expected_location);
assert_eq!(file_meta.size, expected_size);
assert_eq!(file_meta.last_modified, i64::MAX);
}
Err(expected_message) => assert_result_error_with_message(result, expected_message),
}
}
#[cfg(feature = "adaptive-metadata-in-dev")]
fn sample_checkpoint_action() -> CheckpointAction {
let sidecar = |path: &str| Sidecar {
path: path.to_string(),
size_in_bytes: 100,
modification_time: 1,
tags: None,
};
CheckpointAction {
version: 42,
content_root: ContentRoot {
path: "s3://bucket/manifest".to_string(),
size_in_bytes: 1024,
version: 40,
},
protocol: Protocol::new_unchecked(1, 2, None, None),
metadata: Metadata::default(),
transactions: vec![SetTransaction {
app_id: "myApp".to_string(),
version: 3,
last_updated: None,
}],
domain_metadata: vec![DomainMetadata {
domain: "myDomain".to_string(),
configuration: "cfg".to_string(),
removed: false,
}],
txn_sidecars: vec![sidecar("txn-sidecar.parquet")],
domain_metadata_sidecars: vec![sidecar("dm-sidecar.parquet")],
}
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_checkpoint_action_into_engine_data_round_trip() -> DeltaResult<()> {
let engine = ExprEngine::new();
let action = sample_checkpoint_action();
let data = action
.clone()
.into_engine_data(LOG_CHECKPOINT_SCHEMA.clone(), &engine)?;
let back = CheckpointAction::try_new_from_data(data.as_ref())?
.expect("checkpoint action should round-trip");
assert_eq!(action, back);
Ok(())
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_checkpoint_action_into_engine_data_rejects_invalid_content_root_version() {
let base = sample_checkpoint_action();
let action = CheckpointAction {
content_root: ContentRoot {
version: base.version + 1,
..base.content_root
},
..sample_checkpoint_action()
};
let engine = ExprEngine::new();
let result = action.into_engine_data(LOG_CHECKPOINT_SCHEMA.clone(), &engine);
assert_result_error_with_message(result, "exceeds checkpointMetadata.version");
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_checkpoint_action_wire_format() -> DeltaResult<()> {
use crate::engine::to_json_bytes;
use crate::engine_data::FilteredEngineData;
let engine = ExprEngine::new();
let data =
sample_checkpoint_action().into_engine_data(LOG_CHECKPOINT_SCHEMA.clone(), &engine)?;
let filtered = FilteredEngineData::with_all_rows_selected(data);
let bytes = to_json_bytes(std::iter::once(Ok(filtered)))?;
let json: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
assert_eq!(
json,
json!({ "checkpoint": [
{ "checkpointMetadata": { "version": 42 } },
{ "contentRoot": { "path": "s3://bucket/manifest", "sizeInBytes": 1024, "version": 40 } },
{ "protocol": { "minReaderVersion": 1, "minWriterVersion": 2 } },
{ "metaData": {
"id": "",
"format": { "provider": "parquet", "options": {} },
"schemaString": "",
"partitionColumns": [],
"configuration": {},
} },
{ "txn": { "appId": "myApp", "version": 3 } },
{ "domainMetadata": { "domain": "myDomain", "configuration": "cfg", "removed": false } },
{ "sidecar": { "type": "txn", "path": "txn-sidecar.parquet", "sizeInBytes": 100, "modificationTime": 1 } },
{ "sidecar": { "type": "domainMetadata", "path": "dm-sidecar.parquet", "sizeInBytes": 100, "modificationTime": 1 } },
] })
);
Ok(())
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_try_new_from_data_returns_none_when_no_checkpoint_action() -> DeltaResult<()> {
let data = crate::unit_test_utils::action_batch();
assert!(CheckpointAction::try_new_from_data(data.as_ref())?.is_none());
Ok(())
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_checkpoint_action_round_trip_multiple_and_empty_collections() -> DeltaResult<()> {
let sidecar = |path: &str| Sidecar {
path: path.to_string(),
size_in_bytes: 1,
modification_time: 2,
tags: None,
};
let action = CheckpointAction {
version: 10,
content_root: ContentRoot {
path: "s3://bucket/manifest".to_string(),
size_in_bytes: 8,
version: 8,
},
protocol: Protocol::new_unchecked(1, 2, None, None),
metadata: Metadata::default(),
transactions: vec![
SetTransaction {
app_id: "a1".to_string(),
version: 1,
last_updated: None,
},
SetTransaction {
app_id: "a2".to_string(),
version: 2,
last_updated: None,
},
],
domain_metadata: vec![
DomainMetadata {
domain: "d1".to_string(),
configuration: "c1".to_string(),
removed: false,
},
DomainMetadata {
domain: "d2".to_string(),
configuration: "c2".to_string(),
removed: true,
},
],
txn_sidecars: vec![sidecar("t1.parquet"), sidecar("t2.parquet")],
domain_metadata_sidecars: vec![],
};
let engine = ExprEngine::new();
let data = action
.clone()
.into_engine_data(LOG_CHECKPOINT_SCHEMA.clone(), &engine)?;
let back = CheckpointAction::try_new_from_data(data.as_ref())?
.expect("checkpoint action should round-trip");
assert_eq!(action, back);
Ok(())
}
#[cfg(feature = "adaptive-metadata-in-dev")]
#[test]
fn test_checkpoint_action_round_trip_protocol_with_features() -> DeltaResult<()> {
let action = CheckpointAction {
protocol: Protocol::new_unchecked(
3,
7,
Some(vec![TableFeature::AdaptiveMetadataPreview]),
Some(vec![TableFeature::AdaptiveMetadataPreview]),
),
..sample_checkpoint_action()
};
let engine = ExprEngine::new();
let data = action
.clone()
.into_engine_data(LOG_CHECKPOINT_SCHEMA.clone(), &engine)?;
let back = CheckpointAction::try_new_from_data(data.as_ref())?
.expect("checkpoint action should round-trip");
assert_eq!(action, back);
Ok(())
}
}