use std::collections::HashMap;
use std::sync::{Arc, LazyLock};
use delta_kernel_derive::internal_api;
use itertools::Itertools;
use super::log_replay::TableChangesScanMetadata;
use crate::actions::deletion_vector::DeletionVectorDescriptor;
use crate::actions::visitors::visit_deletion_vector_at;
use crate::engine_data::{GetData, TypedGetData};
use crate::expressions::{col, lit, Expression};
use crate::scan::state::DvInfo;
use crate::schema::{
ColumnName, ColumnNamesAndTypes, DataType, MapType, SchemaRef, StructField, StructType,
};
use crate::utils::require;
use crate::{DeltaResult, Error, RowVisitor};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum CdfScanFileType {
Add,
Remove,
Cdc,
}
impl CdfScanFileType {
pub(crate) fn get_cdf_string_value(&self) -> &str {
match self {
CdfScanFileType::Add => super::ADD_CHANGE_TYPE,
CdfScanFileType::Remove => super::REMOVE_CHANGE_TYPE,
CdfScanFileType::Cdc => "not-expected",
}
}
}
#[derive(Debug, PartialEq, Clone)]
pub(crate) struct CdfScanFile {
pub scan_type: CdfScanFileType,
pub path: String,
pub dv_info: DvInfo,
pub remove_dv: Option<DvInfo>,
pub partition_values: HashMap<String, String>,
pub commit_version: i64,
pub commit_timestamp: i64,
pub size: Option<i64>,
pub base_row_id: Option<i64>,
pub default_row_commit_version: Option<i64>,
}
pub(crate) type CdfScanCallback<T> = fn(context: &mut T, scan_file: CdfScanFile);
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
#[internal_api]
pub(crate) struct TableChangesScanFile {
pub path: String,
pub deletion_vector: Option<DeletionVectorDescriptor>,
pub partition_values: HashMap<String, String>,
pub size: Option<i64>,
pub commit_version: i64,
pub commit_timestamp: i64,
pub base_row_id: Option<i64>,
pub default_row_commit_version: Option<i64>,
}
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
#[internal_api]
pub(crate) struct TableChangesFileAction {
pub add: Option<TableChangesScanFile>,
pub remove: Option<TableChangesScanFile>,
}
impl TableChangesFileAction {
pub(crate) fn try_from_scan_file(scan_file: CdfScanFile) -> DeltaResult<Self> {
match scan_file.scan_type {
CdfScanFileType::Add => {
require!(
scan_file.base_row_id.is_some(),
Error::missing_data(format!(
"baseRowId for row-tracking add action at path {} in version {}",
scan_file.path, scan_file.commit_version
))
);
require!(
scan_file.default_row_commit_version.is_some(),
Error::missing_data(format!(
"defaultRowCommitVersion for row-tracking add action at path {} in \
version {}",
scan_file.path, scan_file.commit_version
))
);
let remove = scan_file.remove_dv.as_ref().map(|dv| {
TableChangesScanFile::from_scan_file(&scan_file, dv.deletion_vector.clone())
});
Ok(TableChangesFileAction {
add: Some(TableChangesScanFile::from_scan_file(
&scan_file,
scan_file.dv_info.deletion_vector.clone(),
)),
remove,
})
}
CdfScanFileType::Remove => Ok(TableChangesFileAction {
add: None,
remove: Some(TableChangesScanFile::from_scan_file(
&scan_file,
scan_file.dv_info.deletion_vector.clone(),
)),
}),
CdfScanFileType::Cdc => Err(Error::internal_error(format!(
"Row-tracking change feed listing unexpectedly produced a cdc scan file: \
path={}, version={}",
scan_file.path, scan_file.commit_version
))),
}
}
}
impl TableChangesScanFile {
fn from_scan_file(
scan_file: &CdfScanFile,
deletion_vector: Option<DeletionVectorDescriptor>,
) -> Self {
TableChangesScanFile {
path: scan_file.path.clone(),
deletion_vector,
partition_values: scan_file.partition_values.clone(),
size: scan_file.size,
commit_version: scan_file.commit_version,
commit_timestamp: scan_file.commit_timestamp,
base_row_id: scan_file.base_row_id,
default_row_commit_version: scan_file.default_row_commit_version,
}
}
}
pub(crate) fn scan_metadata_to_scan_file(
scan_metadata: impl Iterator<Item = DeltaResult<TableChangesScanMetadata>>,
) -> impl Iterator<Item = DeltaResult<CdfScanFile>> {
scan_metadata
.map(|scan_metadata| -> DeltaResult<_> {
let scan_metadata = scan_metadata?;
let callback: CdfScanCallback<Vec<CdfScanFile>> =
|context, scan_file| context.push(scan_file);
Ok(visit_cdf_scan_files(&scan_metadata, vec![], callback)?.into_iter())
}) .flatten_ok() }
pub(crate) fn visit_cdf_scan_files<T>(
scan_metadata: &TableChangesScanMetadata,
context: T,
callback: CdfScanCallback<T>,
) -> DeltaResult<T> {
let mut visitor = CdfScanFileVisitor {
callback,
context,
selection_vector: &scan_metadata.selection_vector,
remove_dvs: scan_metadata.remove_dvs.as_ref(),
};
visitor.visit_rows_of(scan_metadata.scan_metadata.as_ref())?;
Ok(visitor.context)
}
struct CdfScanFileVisitor<'a, T> {
callback: CdfScanCallback<T>,
selection_vector: &'a [bool],
remove_dvs: &'a HashMap<String, DvInfo>,
context: T,
}
struct FileSideSpec {
scan_type: CdfScanFileType,
start_index: usize,
path_field: &'static str,
partition_values_field: &'static str,
size_field: &'static str,
base_row_id_field: &'static str,
default_row_commit_version_field: &'static str,
}
const ADD_FILE_SIDE: FileSideSpec = FileSideSpec {
scan_type: CdfScanFileType::Add,
start_index: 0,
path_field: "scanFile.add.path",
partition_values_field: "scanFile.add.fileConstantValues.partitionValues",
size_field: "scanFile.add.size",
base_row_id_field: "scanFile.add.baseRowId",
default_row_commit_version_field: "scanFile.add.defaultRowCommitVersion",
};
const REMOVE_FILE_SIDE: FileSideSpec = FileSideSpec {
scan_type: CdfScanFileType::Remove,
start_index: 10,
path_field: "scanFile.remove.path",
partition_values_field: "scanFile.remove.fileConstantValues.partitionValues",
size_field: "scanFile.remove.size",
base_row_id_field: "scanFile.remove.baseRowId",
default_row_commit_version_field: "scanFile.remove.defaultRowCommitVersion",
};
const CDC_PATH_INDEX: usize = 20;
const CDC_PARTITION_VALUES_INDEX: usize = 21;
const CDC_SIZE_INDEX: usize = 22;
const COMMIT_TIMESTAMP_INDEX: usize = 23;
const COMMIT_VERSION_INDEX: usize = 24;
const CDF_SCAN_FILE_GETTER_COUNT: usize = 25;
struct FileSide {
scan_type: CdfScanFileType,
path: String,
deletion_vector: Option<DeletionVectorDescriptor>,
partition_values: Option<HashMap<String, String>>,
size: Option<i64>,
base_row_id: Option<i64>,
default_row_commit_version: Option<i64>,
}
fn read_file_side<'a>(
row_index: usize,
getters: &[&'a dyn GetData<'a>],
spec: &FileSideSpec,
) -> DeltaResult<Option<FileSide>> {
let Some(path) = getters[spec.start_index].get_opt(row_index, spec.path_field)? else {
return Ok(None);
};
let deletion_vector = visit_deletion_vector_at(
row_index,
&getters[spec.start_index + 1..=spec.start_index + 5],
)?;
let partition_values =
getters[spec.start_index + 6].get_opt(row_index, spec.partition_values_field)?;
let size = getters[spec.start_index + 7].get_opt(row_index, spec.size_field)?;
let base_row_id = getters[spec.start_index + 8].get_opt(row_index, spec.base_row_id_field)?;
let default_row_commit_version =
getters[spec.start_index + 9].get_opt(row_index, spec.default_row_commit_version_field)?;
Ok(Some(FileSide {
scan_type: spec.scan_type,
path,
deletion_vector,
partition_values,
size,
base_row_id,
default_row_commit_version,
}))
}
impl<T> RowVisitor for CdfScanFileVisitor<'_, T> {
fn visit<'a>(&mut self, row_count: usize, getters: &[&'a dyn GetData<'a>]) -> DeltaResult<()> {
require!(
getters.len() == CDF_SCAN_FILE_GETTER_COUNT,
Error::InternalError(format!(
"Wrong number of CdfScanFileVisitor getters: {}",
getters.len()
))
);
for row_index in 0..row_count {
if !self.selection_vector[row_index] {
continue;
}
let file_side = if let Some(side) = read_file_side(row_index, getters, &ADD_FILE_SIDE)?
{
side
} else if let Some(side) = read_file_side(row_index, getters, &REMOVE_FILE_SIDE)? {
side
} else if let Some(path) =
getters[CDC_PATH_INDEX].get_opt(row_index, "scanFile.cdc.path")?
{
let partition_values = getters[CDC_PARTITION_VALUES_INDEX]
.get_opt(row_index, "scanFile.cdc.fileConstantValues.partitionValues")?;
let size = getters[CDC_SIZE_INDEX].get_opt(row_index, "scanFile.cdc.size")?;
FileSide {
scan_type: CdfScanFileType::Cdc,
path,
deletion_vector: None,
partition_values,
size,
base_row_id: None,
default_row_commit_version: None,
}
} else {
continue;
};
let scan_file = CdfScanFile {
remove_dv: self.remove_dvs.get(&file_side.path).cloned(),
scan_type: file_side.scan_type,
path: file_side.path,
dv_info: DvInfo {
deletion_vector: file_side.deletion_vector,
},
partition_values: file_side.partition_values.unwrap_or_default(),
commit_timestamp: getters[COMMIT_TIMESTAMP_INDEX]
.get(row_index, "scanFile.timestamp")?,
commit_version: getters[COMMIT_VERSION_INDEX]
.get(row_index, "scanFile.commit_version")?,
size: file_side.size,
base_row_id: file_side.base_row_id,
default_row_commit_version: file_side.default_row_commit_version,
};
(self.callback)(&mut self.context, scan_file)
}
Ok(())
}
fn selected_column_names_and_types(&self) -> (&'static [ColumnName], &'static [DataType]) {
static NAMES_AND_TYPES: LazyLock<ColumnNamesAndTypes> =
LazyLock::new(|| cdf_scan_row_schema().leaves(None));
NAMES_AND_TYPES.as_ref()
}
}
pub(crate) fn cdf_scan_row_schema() -> SchemaRef {
static CDF_SCAN_ROW_SCHEMA: LazyLock<Arc<StructType>> = LazyLock::new(|| {
let deletion_vector = StructType::new_unchecked([
StructField::nullable("storageType", DataType::STRING),
StructField::nullable("pathOrInlineDv", DataType::STRING),
StructField::nullable("offset", DataType::INTEGER),
StructField::nullable("sizeInBytes", DataType::INTEGER),
StructField::nullable("cardinality", DataType::LONG),
]);
let partition_values = MapType::new(DataType::STRING, DataType::STRING, true);
let file_constant_values =
StructType::new_unchecked([StructField::nullable("partitionValues", partition_values)]);
let add = StructType::new_unchecked([
StructField::nullable("path", DataType::STRING),
StructField::nullable("deletionVector", deletion_vector.clone()),
StructField::nullable("fileConstantValues", file_constant_values.clone()),
StructField::nullable("size", DataType::LONG),
StructField::nullable("baseRowId", DataType::LONG),
StructField::nullable("defaultRowCommitVersion", DataType::LONG),
]);
let remove = StructType::new_unchecked([
StructField::nullable("path", DataType::STRING),
StructField::nullable("deletionVector", deletion_vector),
StructField::nullable("fileConstantValues", file_constant_values.clone()),
StructField::nullable("size", DataType::LONG),
StructField::nullable("baseRowId", DataType::LONG),
StructField::nullable("defaultRowCommitVersion", DataType::LONG),
]);
let cdc = StructType::new_unchecked([
StructField::nullable("path", DataType::STRING),
StructField::nullable("fileConstantValues", file_constant_values),
StructField::nullable("size", DataType::LONG),
]);
Arc::new(StructType::new_unchecked([
StructField::nullable("add", add),
StructField::nullable("remove", remove),
StructField::nullable("cdc", cdc),
StructField::not_null("timestamp", DataType::LONG),
StructField::not_null("commit_version", DataType::LONG),
]))
});
CDF_SCAN_ROW_SCHEMA.clone()
}
pub(crate) fn cdf_scan_row_expression(commit_timestamp: i64, commit_number: i64) -> Expression {
Expression::struct_from([
Expression::struct_from([
col!("add.path"),
col!("add.deletionVector"),
Expression::struct_from([col!("add.partitionValues")]),
col!("add.size"),
col!("add.baseRowId"),
col!("add.defaultRowCommitVersion"),
]),
Expression::struct_from([
col!("remove.path"),
col!("remove.deletionVector"),
Expression::struct_from([col!("remove.partitionValues")]),
col!("remove.size"),
col!("remove.baseRowId"),
col!("remove.defaultRowCommitVersion"),
]),
Expression::struct_from([
col!("cdc.path"),
Expression::struct_from([col!("cdc.partitionValues")]),
col!("cdc.size"),
]),
lit(commit_timestamp),
lit(commit_number),
])
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::Arc;
use itertools::Itertools;
use rstest::rstest;
use super::{scan_metadata_to_scan_file, CdfScanFile, CdfScanFileType, TableChangesFileAction};
use crate::actions::deletion_vector::{DeletionVectorDescriptor, DeletionVectorStorageType};
use crate::actions::{Add, Cdc, Metadata, Protocol, Remove};
use crate::engine::sync::SyncEngine;
use crate::log_segment::LogSegment;
use crate::scan::state::DvInfo;
use crate::schema::{DataType, StructField, StructType};
use crate::table_changes::log_replay::{
table_changes_action_iter, table_changes_action_iter_with_mode,
};
use crate::table_changes::test_utils::{row_tracking_table_config, test_deletion_vector};
use crate::table_changes::CdfMode;
use crate::table_configuration::TableConfiguration;
use crate::table_properties::{COLUMN_MAPPING_MODE, ENABLE_CHANGE_DATA_FEED};
use crate::unit_test_utils::{assert_result_error_with_message, Action, LocalMockTable};
use crate::Engine as _;
fn row_tracking_add(
path: &str,
deletion_vector: Option<DeletionVectorDescriptor>,
data_change: bool,
base_row_id: i64,
default_row_commit_version: i64,
) -> Add {
Add {
path: path.into(),
deletion_vector,
partition_values: HashMap::new(),
data_change,
size: 100,
base_row_id: Some(base_row_id),
default_row_commit_version: Some(default_row_commit_version),
..Default::default()
}
}
fn row_tracking_remove(
path: &str,
deletion_vector: Option<DeletionVectorDescriptor>,
base_row_id: Option<i64>,
default_row_commit_version: Option<i64>,
) -> Remove {
Remove {
path: path.into(),
deletion_vector,
data_change: true,
base_row_id,
default_row_commit_version,
..Default::default()
}
}
fn row_tracking_scan_files(
engine: Arc<SyncEngine>,
mock_table: &LocalMockTable,
) -> Vec<CdfScanFile> {
let table_root = url::Url::from_directory_path(mock_table.table_root()).unwrap();
let log_segment = LogSegment::for_table_changes(
engine.storage_handler().as_ref(),
table_root.join("_delta_log/").unwrap(),
0,
None,
)
.unwrap();
let table_schema = Arc::new(StructType::new_unchecked([
StructField::nullable("id", DataType::INTEGER),
StructField::nullable("value", DataType::STRING),
]));
let table_config = row_tracking_table_config(table_root, table_schema.clone());
let scan_metadata = table_changes_action_iter_with_mode(
engine,
&table_config,
log_segment.listed.ascending_commit_files,
table_schema,
None,
CdfMode::RowTracking,
)
.unwrap();
scan_metadata_to_scan_file(scan_metadata)
.try_collect()
.unwrap()
}
#[tokio::test]
async fn test_scan_file_visiting() {
let engine = SyncEngine::new();
let mut mock_table = LocalMockTable::new();
let dv_info = DeletionVectorDescriptor {
storage_type: DeletionVectorStorageType::PersistedRelative,
path_or_inline_dv: "vBn[lx{q8@P<9BNH/isA".to_string(),
offset: Some(1),
size_in_bytes: 36,
cardinality: 2,
};
let add_partition_values = HashMap::from([("a".to_string(), "b".to_string())]);
let add_paired = Add {
path: "fake_path_1".into(),
deletion_vector: Some(dv_info.clone()),
partition_values: add_partition_values,
data_change: true,
size: 100i64,
..Default::default()
};
let paired_remove = Remove {
path: "fake_path_1".into(),
deletion_vector: None,
partition_values: None,
data_change: true,
size: Some(200i64),
..Default::default()
};
let rm_dv = DeletionVectorDescriptor {
storage_type: DeletionVectorStorageType::PersistedRelative,
path_or_inline_dv: "U5OWRz5k%CFT.Td}yCPW".to_string(),
offset: Some(1),
size_in_bytes: 38,
cardinality: 3,
};
let rm_partition_values = Some(HashMap::from([("c".to_string(), "d".to_string())]));
let remove = Remove {
path: "fake_path_2".into(),
deletion_vector: Some(rm_dv),
partition_values: rm_partition_values,
data_change: true,
size: None,
..Default::default()
};
let cdc_partition_values = HashMap::from([("x".to_string(), "y".to_string())]);
let cdc = Cdc {
path: "fake_path_3".into(),
partition_values: cdc_partition_values,
..Default::default()
};
let remove_no_partition = Remove {
path: "fake_path_2".into(),
deletion_vector: None,
partition_values: None,
data_change: true,
size: None,
..Default::default()
};
mock_table
.commit([
Action::Remove(paired_remove.clone()),
Action::Add(add_paired.clone()),
Action::Remove(remove.clone()),
])
.await;
mock_table.commit([Action::Cdc(cdc.clone())]).await;
mock_table
.commit([Action::Remove(remove_no_partition.clone())])
.await;
let table_root = url::Url::from_directory_path(mock_table.table_root()).unwrap();
let log_root = table_root.join("_delta_log/").unwrap();
let log_segment =
LogSegment::for_table_changes(engine.storage_handler().as_ref(), log_root, 0, None)
.unwrap();
let table_schema = Arc::new(StructType::new_unchecked([
StructField::nullable("id", DataType::INTEGER),
StructField::nullable("value", DataType::STRING),
]));
let metadata = Metadata::try_new(
None,
None,
table_schema.clone(),
vec![],
0,
HashMap::from([
(ENABLE_CHANGE_DATA_FEED.to_string(), "true".to_string()),
(COLUMN_MAPPING_MODE.to_string(), "none".to_string()),
]),
)
.unwrap();
let protocol = Protocol::try_new_legacy(1, 4).unwrap();
let table_config =
TableConfiguration::try_new(metadata, protocol, table_root.clone(), 0).unwrap();
let scan_metadata = table_changes_action_iter(
Arc::new(engine),
&table_config,
log_segment.listed.ascending_commit_files.clone(),
table_schema,
None,
)
.unwrap();
let scan_files: Vec<_> = scan_metadata_to_scan_file(scan_metadata)
.try_collect()
.unwrap();
let timestamps = log_segment
.listed
.ascending_commit_files
.iter()
.map(|commit| commit.location.last_modified)
.collect_vec();
let expected_remove_dv = DvInfo {
deletion_vector: None,
};
let expected_scan_files = vec![
CdfScanFile {
scan_type: CdfScanFileType::Add,
path: add_paired.path,
dv_info: DvInfo {
deletion_vector: add_paired.deletion_vector,
},
partition_values: add_paired.partition_values,
commit_version: 0,
commit_timestamp: timestamps[0],
remove_dv: Some(expected_remove_dv),
size: Some(add_paired.size),
base_row_id: None,
default_row_commit_version: None,
},
CdfScanFile {
scan_type: CdfScanFileType::Remove,
path: remove.path,
dv_info: DvInfo {
deletion_vector: remove.deletion_vector,
},
partition_values: remove.partition_values.unwrap(),
commit_version: 0,
commit_timestamp: timestamps[0],
remove_dv: None,
size: remove.size,
base_row_id: None,
default_row_commit_version: None,
},
CdfScanFile {
scan_type: CdfScanFileType::Cdc,
path: cdc.path,
dv_info: DvInfo {
deletion_vector: None,
},
partition_values: cdc.partition_values,
commit_version: 1,
commit_timestamp: timestamps[1],
remove_dv: None,
size: Some(cdc.size),
base_row_id: None,
default_row_commit_version: None,
},
CdfScanFile {
scan_type: CdfScanFileType::Remove,
path: remove_no_partition.path,
dv_info: DvInfo {
deletion_vector: None,
},
partition_values: HashMap::new(),
commit_version: 2,
commit_timestamp: timestamps[2],
remove_dv: None,
size: remove_no_partition.size,
base_row_id: None,
default_row_commit_version: None,
},
];
assert_eq!(scan_files, expected_scan_files);
}
#[tokio::test]
async fn test_row_tracking_listing_ignores_cdc_and_surfaces_row_tracking_fields() {
let engine = Arc::new(SyncEngine::new());
let mut mock_table = LocalMockTable::new();
let post_dv = test_deletion_vector("vBn[lx{q8@P<9BNH/isA", 2);
let baseline_dv = test_deletion_vector("U5OWRz5k%CFT.Td}yCPW", 1);
let add_paired = row_tracking_add("path_1", Some(post_dv.clone()), true, 10, 0);
let paired_remove =
row_tracking_remove("path_1", Some(baseline_dv.clone()), Some(10), Some(0));
let remove_unpaired = row_tracking_remove("path_2", None, None, None);
let cdc = Cdc {
path: "cdc_path".into(),
partition_values: HashMap::new(),
..Default::default()
};
let add_insert = row_tracking_add("path_4", None, true, 20, 2);
let add_no_data_change = row_tracking_add("path_6", None, false, 40, 3);
mock_table
.commit([
Action::Remove(paired_remove),
Action::Add(add_paired),
Action::Remove(remove_unpaired),
])
.await;
mock_table.commit([Action::Cdc(cdc)]).await;
mock_table.commit([Action::Add(add_insert)]).await;
mock_table.commit([Action::Add(add_no_data_change)]).await;
let scan_files = row_tracking_scan_files(engine, &mock_table);
assert!(scan_files
.iter()
.all(|f| f.scan_type != CdfScanFileType::Cdc));
assert!(
!scan_files.iter().any(|f| f.path == "path_6"),
"metadata-only add (data_change=false) must be excluded"
);
assert_eq!(scan_files.len(), 3);
let paired = scan_files
.iter()
.find(|f| f.path == "path_1")
.expect("paired add present");
assert_eq!(paired.scan_type, CdfScanFileType::Add);
assert_eq!(paired.base_row_id, Some(10));
assert_eq!(paired.default_row_commit_version, Some(0));
assert_eq!(paired.dv_info.deletion_vector, Some(post_dv));
assert_eq!(
paired
.remove_dv
.as_ref()
.and_then(|dv| dv.deletion_vector.clone()),
Some(baseline_dv)
);
let unpaired = scan_files
.iter()
.find(|f| f.path == "path_2")
.expect("unpaired remove present");
assert_eq!(unpaired.scan_type, CdfScanFileType::Remove);
assert_eq!(unpaired.remove_dv, None);
assert_eq!(unpaired.base_row_id, None);
assert_eq!(unpaired.default_row_commit_version, None);
let insert = scan_files
.iter()
.find(|f| f.path == "path_4")
.expect("insert add present");
assert_eq!(insert.scan_type, CdfScanFileType::Add);
assert_eq!(insert.base_row_id, Some(20));
let listing: Vec<TableChangesFileAction> = scan_files
.into_iter()
.map(TableChangesFileAction::try_from_scan_file)
.try_collect()
.unwrap();
let path_of = |fa: &TableChangesFileAction| {
fa.add
.as_ref()
.or(fa.remove.as_ref())
.map(|s| s.path.clone())
};
let paired = listing
.iter()
.find(|fa| path_of(fa).as_deref() == Some("path_1"))
.expect("paired listing present");
let paired_add = paired.add.as_ref().expect("paired add side");
let paired_remove = paired.remove.as_ref().expect("paired remove side");
assert_eq!(paired_add.base_row_id, Some(10));
assert!(paired_add.deletion_vector.is_some());
assert!(paired_remove.deletion_vector.is_some());
assert_eq!(paired_remove.base_row_id, Some(10));
let unpaired_delete = listing
.iter()
.find(|fa| path_of(fa).as_deref() == Some("path_2"))
.expect("unpaired delete present");
assert!(unpaired_delete.add.is_none());
assert!(unpaired_delete.remove.is_some());
let insert = listing
.iter()
.find(|fa| path_of(fa).as_deref() == Some("path_4"))
.expect("insert present");
assert!(insert.remove.is_none());
assert_eq!(
insert.add.as_ref().expect("insert add side").base_row_id,
Some(20)
);
}
#[rstest]
#[case::with_deletion_vector(Some(test_deletion_vector("baseline", 1)))]
#[case::without_deletion_vector(None)]
fn cdf_file_action_preserves_paired_remove_deletion_vector(
#[case] deletion_vector: Option<DeletionVectorDescriptor>,
) {
let scan_file = CdfScanFile {
scan_type: CdfScanFileType::Add,
path: "path".into(),
dv_info: DvInfo {
deletion_vector: Some(test_deletion_vector("current", 2)),
},
remove_dv: Some(DvInfo {
deletion_vector: deletion_vector.clone(),
}),
partition_values: HashMap::new(),
commit_version: 1,
commit_timestamp: 2,
size: Some(3),
base_row_id: Some(4),
default_row_commit_version: Some(5),
};
let action = TableChangesFileAction::try_from_scan_file(scan_file).unwrap();
assert_eq!(action.remove.unwrap().deletion_vector, deletion_vector);
assert!(action.add.is_some());
}
#[rstest]
#[case::missing_base_row_id(None, Some(5), "baseRowId")]
#[case::missing_default_row_commit_version(Some(4), None, "defaultRowCommitVersion")]
fn cdf_file_action_rejects_add_missing_row_tracking_fields(
#[case] base_row_id: Option<i64>,
#[case] default_row_commit_version: Option<i64>,
#[case] expected: &str,
) {
let scan_file = CdfScanFile {
scan_type: CdfScanFileType::Add,
path: "path".into(),
dv_info: DvInfo {
deletion_vector: None,
},
remove_dv: None,
partition_values: HashMap::new(),
commit_version: 1,
commit_timestamp: 2,
size: Some(3),
base_row_id,
default_row_commit_version,
};
assert_result_error_with_message(
TableChangesFileAction::try_from_scan_file(scan_file),
expected,
);
}
#[test]
fn cdf_file_action_rejects_cdc_scan_file() {
let cdc_scan_file = CdfScanFile {
scan_type: CdfScanFileType::Cdc,
path: "cdc".into(),
dv_info: DvInfo {
deletion_vector: None,
},
remove_dv: None,
partition_values: HashMap::new(),
commit_version: 0,
commit_timestamp: 0,
size: None,
base_row_id: None,
default_row_commit_version: None,
};
assert!(TableChangesFileAction::try_from_scan_file(cdc_scan_file).is_err());
}
}