1use std::collections::{HashMap, HashSet};
2use std::iter;
3use std::marker::PhantomData;
4use std::ops::Deref;
5use std::sync::{Arc, LazyLock};
6use std::time::{Duration, Instant};
7
8use delta_kernel_derive::internal_api;
9use tracing::{info, instrument};
10
11use crate::actions::{
12 as_log_add_schema, CommitInfo, DomainMetadata, Metadata, Protocol, SetTransaction,
13 LOG_METADATA_SCHEMA, LOG_PROTOCOL_SCHEMA, LOG_REMOVE_SCHEMA, LOG_TXN_SCHEMA, MAX_VALUES,
14 MIN_VALUES, NULL_COUNT, NUM_RECORDS, TIGHT_BOUNDS,
15};
16use crate::committer::{
17 CommitMetadata, CommitProtocolMetadata, CommitResponse, CommitType, Committer,
18};
19use crate::crc::{is_incremental_safe_operation, CrcDelta, FileStatsDelta};
20use crate::engine_data::FilteredEngineData;
21use crate::error::Error;
22use crate::expressions::UnaryExpressionOp::ToJson;
23use crate::expressions::{
24 col, lit, ArrayData, ColumnName, ExpressionStructPatch, ExpressionStructPatchBuilder, Scalar,
25};
26use crate::log_segment::LogSegment;
27use crate::metrics::events::TRANSACTION_COMMIT_SPAN;
28use crate::metrics::{CommitFailureReason, MetricId};
29use crate::partition::serialization::serialize_partition_value;
30use crate::partition::validation::validate_partition_values;
31use crate::path::{LogRoot, ParsedLogPath};
32use crate::row_tracking::{RowTrackingDomainMetadata, RowTrackingVisitor};
33use crate::scan::data_skipping::stats_schema::schema_with_all_fields_nullable;
34use crate::scan::log_replay::{
35 BASE_ROW_ID_NAME, DEFAULT_ROW_COMMIT_VERSION_NAME, FILE_CONSTANT_VALUES_NAME,
36 PARTITION_VALUES_PARSED_NAME, STATS_PARSED_NAME, TAGS_NAME,
37};
38use crate::scan::scan_row_schema;
39use crate::schema::void_utils::{add_void_stripping, validate_schema_for_write};
40#[cfg(feature = "column-defaults-in-dev")]
41use crate::schema::ColumnDefault;
42use crate::schema::{
43 lazy_schema_ref, ArrayType, SchemaRef, SchemaStructPatchBuilder, StructField, StructType,
44};
45use crate::snapshot::{Snapshot, SnapshotRef};
46use crate::struct_patch::ProjectionStructPatchBuilder;
47use crate::table_configuration::TableConfiguration;
48use crate::table_features::TableFeature;
49use crate::utils::require;
50use crate::{
51 DataType, DeltaResult, Engine, EngineData, Expression, FileMeta, IntoEngineData, RowVisitor,
52 Version,
53};
54
55#[cfg(feature = "internal-api")]
56pub mod builder;
57#[cfg(not(feature = "internal-api"))]
58pub(crate) mod builder;
59
60#[cfg(feature = "internal-api")]
61pub mod create_table;
62#[cfg(not(feature = "internal-api"))]
63pub(crate) mod create_table;
64
65#[cfg(feature = "internal-api")]
66pub mod data_layout;
67#[cfg(not(feature = "internal-api"))]
68pub(crate) mod data_layout;
69
70pub(crate) mod alter_table;
71pub use alter_table::AlterTableTransaction;
72mod commit_info;
73mod domain_metadata;
74pub(crate) mod schema_evolution;
75#[cfg(feature = "internal-api")]
76pub mod stats_verifier;
77#[cfg(not(feature = "internal-api"))]
78mod stats_verifier;
79mod update;
80mod write_context;
81
82use stats_verifier::StatsColumnVerifier;
83use write_context::SharedWriteState;
84pub use write_context::WriteContext;
85
86pub(crate) type EngineDataResultIterator<'a> =
88 Box<dyn Iterator<Item = DeltaResult<Box<dyn EngineData>>> + Send + 'a>;
89
90pub(crate) static MANDATORY_ADD_FILE_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! {
93 not_null "path": STRING,
94 not_null "partitionValues": { STRING => nullable STRING },
95 not_null "size": LONG,
96 not_null "modificationTime": LONG,
97};
98
99pub(crate) fn mandatory_add_file_schema() -> &'static SchemaRef {
105 &MANDATORY_ADD_FILE_SCHEMA
106}
107
108pub(crate) static BASE_ADD_FILES_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! {
120 ..(mandatory_add_file_schema().fields().cloned()),
121 nullable "stats": {
122 nullable NUM_RECORDS: LONG,
123 nullable NULL_COUNT: {},
127 nullable MIN_VALUES: {},
128 nullable MAX_VALUES: {},
129 nullable TIGHT_BOUNDS: BOOLEAN,
130 },
131};
132
133static DATA_CHANGE_COLUMN: LazyLock<StructField> =
134 LazyLock::new(|| StructField::not_null("dataChange", DataType::BOOLEAN));
135
136fn with_row_tracking_cols(schema: &SchemaRef) -> DeltaResult<SchemaRef> {
140 let patch = SchemaStructPatchBuilder::new()
141 .append(StructField::nullable("baseRowId", DataType::LONG))
142 .append(StructField::nullable(
143 "defaultRowCommitVersion",
144 DataType::LONG,
145 ));
146 Ok(Arc::new(patch.build(schema)?))
147}
148
149#[derive(Debug)]
154pub struct ExistingTable;
155
156#[derive(Debug)]
161pub struct CreateTable;
162
163#[derive(Debug)]
168pub struct AlterTable;
169
170pub trait SupportsDataFiles {}
176impl SupportsDataFiles for ExistingTable {}
177impl SupportsDataFiles for CreateTable {}
178
179pub struct Transaction<S = ExistingTable> {
199 span: tracing::Span,
200 operation_id: MetricId,
202 correlation_id: Option<Arc<str>>,
205 read_snapshot_opt: Option<SnapshotRef>,
208 effective_table_config: TableConfiguration,
212 should_emit_protocol: bool,
214 should_emit_metadata: bool,
216 committer: Box<dyn Committer>,
217 operation: Option<String>,
218 engine_info: Option<String>,
219 engine_commit_info: Option<(Box<dyn EngineData>, SchemaRef)>,
220 add_files_metadata: Vec<Box<dyn EngineData>>,
221 remove_files_metadata: Vec<FilteredEngineData>,
222 set_transactions: Vec<SetTransaction>,
228 commit_timestamp: i64,
231 user_domain_metadata_additions: Vec<DomainMetadata>,
233 system_domain_metadata_additions: Vec<DomainMetadata>,
238 user_domain_removals: Vec<String>,
241 data_change: bool,
243 is_blind_append: bool,
245 dv_matched_files: Vec<FilteredEngineData>,
248 physical_clustering_columns: Option<Vec<ColumnName>>,
252 _state: PhantomData<S>,
255}
256
257impl<S> std::fmt::Debug for Transaction<S> {
258 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
259 let version_info = match &self.read_snapshot_opt {
260 Some(snap) => format!("{}", snap.version()),
261 None => "create_table".to_string(),
262 };
263 f.write_str(&format!(
264 "Transaction {{ read_snapshot version: {}, engine_info: {} }}",
265 version_info,
266 self.engine_info.is_some()
267 ))
268 }
269}
270
271fn build_add_action_projection(
273 input_schema: &StructType,
274 data_change: bool,
275) -> DeltaResult<(SchemaRef, Expression)> {
276 let (output_schema, patch) = ProjectionStructPatchBuilder::new(input_schema)
277 .insert_after(
278 "modificationTime",
279 DATA_CHANGE_COLUMN.clone(),
280 lit(data_change),
281 )
282 .replace(
283 "stats",
284 StructField::nullable("stats", DataType::STRING),
285 Expression::unary(ToJson, col!("stats")),
286 )
287 .build()?;
288 let patch = Expression::struct_from([patch]);
289 Ok((output_schema, patch))
290}
291
292fn build_add_actions<'a, I, T>(
295 engine: &dyn Engine,
296 add_files_metadata: I,
297 input_schema: SchemaRef,
298 data_change: bool,
299) -> DeltaResult<impl Iterator<Item = DeltaResult<Box<dyn EngineData>>> + 'a>
300where
301 I: Iterator<Item = DeltaResult<T>> + Send + 'a,
302 T: Deref<Target = dyn EngineData> + Send + 'a,
303{
304 let evaluation_handler = engine.evaluation_handler();
305 let (output_schema, adds_expr) = build_add_action_projection(&input_schema, data_change)?;
306 let adds_expr = Arc::new(adds_expr);
307 Ok(add_files_metadata.map(move |add_files_batch| {
308 let adds_evaluator = evaluation_handler.new_expression_evaluator(
309 input_schema.clone(),
310 adds_expr.clone(),
311 as_log_add_schema(output_schema.clone()).into(),
312 )?;
313 adds_evaluator.evaluate(add_files_batch?.deref())
314 }))
315}
316
317impl<S> Transaction<S> {
321 #[instrument(
327 parent = &self.span,
328 name = TRANSACTION_COMMIT_SPAN,
329 skip_all,
330 fields(
331 report,
332 operation_id = %self.operation_id,
333 is_catalog_managed = self.effective_table_config.is_catalog_managed(),
334 correlation_id = self.correlation_id.as_deref().unwrap_or(""),
335 commit_version = self.get_commit_version(),
336 num_add_files,
337 num_remove_files,
338 add_files_bytes,
339 remove_files_bytes,
340 is_blind_append,
341 data_change,
342 operation,
343 prepare_duration_ns,
344 committer_duration_ns,
345 failure_reason,
346 ),
347 err
348 )]
349 pub fn commit(self, engine: &dyn Engine) -> DeltaResult<CommitResult<S>> {
350 let commit_start = Instant::now();
351 info!(
355 num_add_files = self.add_files_metadata.len(),
356 num_remove_files = self.remove_files_metadata.len(),
357 num_dv_updates = self.dv_matched_files.len(),
358 );
359
360 if !self.remove_files_metadata.is_empty() {
362 self.effective_table_config
363 .validate_feature_support_for_remove()?;
364 }
365
366 let mut app_ids = HashSet::with_capacity(self.set_transactions.len());
372 if let Some(dup) = self
373 .set_transactions
374 .iter()
375 .find(|t| !app_ids.insert(&t.app_id))
376 {
377 return Err(Error::generic(format!(
378 "app_id {} already exists in transaction",
379 dup.app_id
380 )));
381 }
382
383 self.validate_blind_append_semantics()?;
384 self.ensure_schema_non_empty_for_data_writes()?;
385
386 if !self.add_files_metadata.is_empty() {
390 validate_schema_for_write(&self.effective_table_config.logical_schema())?;
391 }
392
393 if !self.is_create_table()
399 && !self.add_files_metadata.is_empty()
400 && !self.remove_files_metadata.is_empty()
401 && self.data_change
402 {
403 let cdf_enabled = self
404 .effective_table_config
405 .table_properties()
406 .enable_change_data_feed
407 .unwrap_or(false);
408 require!(
409 !cdf_enabled,
410 Error::generic(
411 "Cannot add and remove data in the same transaction when Change Data Feed is enabled (delta.enableChangeDataFeed = true). \
412 This would require writing CDC files for DML operations, which is not yet supported. \
413 Consider using separate transactions: one to add files, another to remove files."
414 )
415 );
416 }
417
418 self.validate_add_files_stats(&self.add_files_metadata)?;
420
421 let set_transaction_actions = self
423 .set_transactions
424 .clone()
425 .into_iter()
426 .map(|txn| txn.into_engine_data(LOG_TXN_SCHEMA.clone(), engine));
427
428 let in_commit_timestamp = self.get_in_commit_timestamp(engine)?;
430 let kernel_commit_info = CommitInfo::new(
431 self.commit_timestamp,
432 in_commit_timestamp,
433 self.operation.clone(),
434 self.engine_info.clone(),
435 self.is_blind_append,
436 );
437 let commit_info_action = self.generate_commit_info(engine, kernel_commit_info);
438
439 let (protocol_action, protocol) = if self.should_emit_protocol {
441 let protocol = self.effective_table_config.protocol().clone();
442 let schema = LOG_PROTOCOL_SCHEMA.clone();
443 let action = protocol.clone().into_engine_data(schema, engine)?;
444 (Some(action), Some(protocol))
445 } else {
446 (None, None)
447 };
448 let (metadata_action, metadata) = if self.should_emit_metadata {
449 let metadata = self.effective_table_config.metadata().clone();
450 let schema = LOG_METADATA_SCHEMA.clone();
451 let action = metadata.clone().into_engine_data(schema, engine)?;
452 (Some(action), Some(metadata))
453 } else {
454 (None, None)
455 };
456
457 let commit_version = self.get_commit_version();
460 let (add_actions, row_tracking_domain_metadata) =
461 self.generate_adds(engine, commit_version)?;
462
463 let (domain_metadata_actions, dm_changes) =
465 self.generate_domain_metadata_actions(engine, row_tracking_domain_metadata)?;
466
467 let dv_update_actions = self.generate_dv_update_actions(engine)?;
469
470 let remove_actions =
472 self.generate_remove_actions(engine, self.remove_files_metadata.iter(), &[])?;
473
474 let actions = iter::once(commit_info_action)
479 .chain(protocol_action.map(Ok))
480 .chain(metadata_action.map(Ok))
481 .chain(add_actions)
482 .chain(set_transaction_actions)
483 .chain(domain_metadata_actions);
484
485 let filtered_actions = actions
486 .map(|action_result| action_result.map(FilteredEngineData::with_all_rows_selected))
487 .chain(remove_actions)
488 .chain(dv_update_actions);
489
490 let commit_metadata = self.create_commit_metadata(
492 commit_version,
493 in_commit_timestamp,
494 protocol,
495 metadata,
496 dm_changes.clone(),
497 )?;
498 let prepare_duration = commit_start.elapsed();
499 let committer_start = Instant::now();
500 let commit_response =
501 self.committer
502 .commit(engine, Box::new(filtered_actions), commit_metadata);
503 let committer_duration = committer_start.elapsed();
504 match commit_response {
505 Ok(CommitResponse::Committed { file_meta }) => {
506 let bin_boundaries = self
509 .read_snapshot_opt
510 .as_ref()
511 .and_then(|snap| snap.get_file_stats_if_present())
512 .and_then(|s| s.file_size_histogram)
513 .map(|h| h.sorted_bin_boundaries);
514 let file_stats = FileStatsDelta::try_compute_for_txn(
515 &self.add_files_metadata,
516 &self.remove_files_metadata,
517 bin_boundaries.as_deref(),
518 )?;
519 self.record_commit_success_metrics(
520 &file_stats,
521 prepare_duration,
522 committer_duration,
523 );
524 let crc_delta =
525 self.build_crc_delta(file_stats, in_commit_timestamp, dm_changes)?;
526 Ok(CommitResult::CommittedTransaction(
527 self.into_committed(file_meta, crc_delta)?,
528 ))
529 }
530 Ok(CommitResponse::Conflict { version }) => {
531 tracing::Span::current()
533 .record("failure_reason", CommitFailureReason::Conflict.as_ref());
534 Ok(CommitResult::ConflictedTransaction(
535 self.into_conflicted(version),
536 ))
537 }
538 Err(e @ Error::IOError(_)) => {
541 tracing::Span::current()
543 .record("failure_reason", CommitFailureReason::RetryableIo.as_ref());
544 Ok(CommitResult::RetryableTransaction(self.into_retryable(e)))
545 }
546 Err(e) => Err(e),
547 }
548 }
549
550 fn record_commit_success_metrics(
551 &self,
552 file_stats: &FileStatsDelta,
553 prepare_duration: Duration,
554 committer_duration: Duration,
555 ) {
556 let span = tracing::Span::current();
557 span.record("num_add_files", file_stats.gross_add_files);
558 span.record("num_remove_files", file_stats.gross_remove_files);
559 span.record("add_files_bytes", file_stats.gross_add_bytes);
560 span.record("remove_files_bytes", file_stats.gross_remove_bytes);
561 span.record("is_blind_append", self.is_blind_append);
562 span.record("data_change", self.data_change);
563 if let Some(operation) = self.operation.as_deref() {
564 span.record("operation", operation);
565 }
566 span.record("prepare_duration_ns", prepare_duration.as_nanos() as u64);
567 span.record(
568 "committer_duration_ns",
569 committer_duration.as_nanos() as u64,
570 );
571 }
572
573 pub fn with_data_change(mut self, data_change: bool) -> Self {
584 self.data_change = data_change;
585 self
586 }
587
588 #[internal_api]
591 #[allow(dead_code)] pub(crate) fn set_data_change(&mut self, data_change: bool) {
593 self.data_change = data_change;
594 }
595
596 pub fn with_engine_info(mut self, engine_info: impl Into<String>) -> Self {
598 self.engine_info = Some(engine_info.into());
599 self
600 }
601
602 pub fn with_correlation_id(mut self, correlation_id: impl Into<Arc<str>>) -> Self {
606 self.correlation_id = Some(correlation_id.into()).filter(|id| !id.is_empty());
607 self
608 }
609
610 pub fn with_commit_info(
623 mut self,
624 engine_commit_info: Box<dyn EngineData>,
625 commit_info_schema: SchemaRef,
626 ) -> Self {
627 self.engine_commit_info = Some((engine_commit_info, commit_info_schema));
628 self
629 }
630
631 pub fn with_transaction_id(mut self, app_id: String, version: i64) -> Self {
637 let set_transaction = SetTransaction::new(app_id, version, Some(self.commit_timestamp));
638 self.set_transactions.push(set_transaction);
639 self
640 }
641
642 pub fn with_domain_metadata(mut self, domain: String, configuration: String) -> Self {
649 self.user_domain_metadata_additions
650 .push(DomainMetadata::new(domain, configuration));
651 self
652 }
653
654 fn determine_commit_type(
657 is_create: bool,
658 table_config: &crate::table_configuration::TableConfiguration,
659 ) -> CommitType {
660 let is_catalog_managed = table_config.is_catalog_managed();
661
662 match (is_create, is_catalog_managed) {
665 (true, true) => CommitType::CatalogManagedCreate,
666 (true, false) => CommitType::PathBasedCreate,
667 (false, true) => CommitType::CatalogManagedWrite,
668 (false, false) => CommitType::PathBasedWrite,
669 }
670 }
671
672 fn validate_commit_type(
675 is_catalog_committer: bool,
676 commit_type: &CommitType,
677 ) -> DeltaResult<()> {
678 match (
679 is_catalog_committer,
680 commit_type.requires_catalog_committer(),
681 ) {
682 (true, true) | (false, false) => Ok(()),
683 (false, true) => Err(Error::generic(
684 "This table is catalog-managed and requires a catalog committer. \
685 Please provide a catalog committer via Snapshot::transaction().",
686 )),
687 (true, false) => Err(Error::generic(
688 "This table is path-based and cannot be committed to with a catalog committer.",
689 )),
690 }
691 }
692
693 fn create_commit_metadata(
696 &self,
697 commit_version: Version,
698 in_commit_timestamp: Option<i64>,
699 new_protocol: Option<Protocol>,
700 new_metadata: Option<Metadata>,
701 domain_metadata_changes: Vec<crate::actions::DomainMetadata>,
702 ) -> DeltaResult<CommitMetadata> {
703 let log_root = LogRoot::new(self.effective_table_config.table_root().clone())?;
704 let is_create = self.is_create_table();
705 let commit_type = Self::determine_commit_type(is_create, &self.effective_table_config);
706 Self::validate_commit_type(self.committer.is_catalog_committer(), &commit_type)?;
707 let (read_protocol, read_metadata, max_published_version) = if is_create {
712 (None, None, None)
713 } else {
714 let snap = self.read_snapshot()?;
715 let read_config = snap.table_configuration();
716 (
717 Some(read_config.protocol().clone()),
718 Some(read_config.metadata().clone()),
719 snap.log_segment().listed.max_published_version,
720 )
721 };
722 let protocol_metadata = CommitProtocolMetadata::try_new(
723 read_protocol,
724 read_metadata,
725 new_protocol,
726 new_metadata,
727 )?;
728 Ok(CommitMetadata::new(
729 log_root,
730 commit_version,
731 commit_type,
732 in_commit_timestamp.unwrap_or(self.commit_timestamp),
733 max_published_version,
734 protocol_metadata,
735 domain_metadata_changes,
736 ))
737 }
738
739 fn validate_blind_append_semantics(&self) -> DeltaResult<()> {
745 if !self.is_blind_append {
746 return Ok(());
747 }
748 require!(
749 !self.is_create_table(),
750 Error::invalid_transaction_state(
751 "Blind append is not supported for create-table transactions",
752 )
753 );
754 require!(
755 !self.add_files_metadata.is_empty(),
756 Error::invalid_transaction_state("Blind append requires at least one added data file")
757 );
758 require!(
759 self.data_change,
760 Error::invalid_transaction_state("Blind append requires data_change to be true")
761 );
762 require!(
763 self.remove_files_metadata.is_empty(),
764 Error::invalid_transaction_state("Blind append cannot remove files")
765 );
766 require!(
767 self.dv_matched_files.is_empty(),
768 Error::invalid_transaction_state("Blind append cannot update deletion vectors")
769 );
770
771 Ok(())
772 }
773
774 fn ensure_schema_non_empty_for_data_writes(&self) -> DeltaResult<()> {
777 if self.is_create_table() {
778 return Ok(());
779 }
780 if self.has_data_file_actions() {
781 self.ensure_schema_non_empty_for_write_context()?;
782 }
783 Ok(())
784 }
785
786 fn ensure_schema_non_empty_for_write_context(&self) -> DeltaResult<()> {
789 if self.is_create_table() {
790 return Ok(());
791 }
792 if self.effective_table_config.logical_schema().num_fields() == 0 {
793 return Err(Error::generic(
794 "Cannot write data files to a Delta table with empty schema; \
795 use `snapshot.alter_table().add_column(...)` to add at least one \
796 column before writing data",
797 ));
798 }
799 Ok(())
800 }
801
802 fn is_create_table(&self) -> bool {
805 debug_assert!(
806 self.operation.as_deref() != Some("CREATE TABLE") || self.read_snapshot_opt.is_none(),
807 "CREATE TABLE operation should not have a read snapshot"
808 );
809 self.read_snapshot_opt.is_none()
810 }
811
812 fn has_data_file_actions(&self) -> bool {
814 !self.add_files_metadata.is_empty()
815 || !self.remove_files_metadata.is_empty()
816 || !self.dv_matched_files.is_empty()
817 }
818
819 fn read_snapshot(&self) -> DeltaResult<&Snapshot> {
822 self.read_snapshot_opt.as_deref().ok_or_else(|| {
823 Error::internal_error("read_snapshot() called on create-table transaction")
824 })
825 }
826
827 fn get_in_commit_timestamp(&self, engine: &dyn Engine) -> DeltaResult<Option<i64>> {
832 let has_ict = self
833 .effective_table_config
834 .is_feature_enabled(&TableFeature::InCommitTimestamp);
835
836 if !has_ict {
837 return Ok(None);
838 }
839
840 if self.is_create_table() {
841 return Ok(Some(self.commit_timestamp));
843 }
844
845 Ok(self
850 .read_snapshot()?
851 .get_in_commit_timestamp(engine)?
852 .map(|prev_ict| self.commit_timestamp.max(prev_ict + 1)))
853 }
854
855 fn get_commit_version(&self) -> Version {
859 match &self.read_snapshot_opt {
860 Some(snap) => snap.version() + 1,
861 None => 0,
862 }
863 }
864
865 pub fn add_files_schema(&self) -> &'static SchemaRef {
886 &BASE_ADD_FILES_SCHEMA
887 }
888}
889
890impl<S: SupportsDataFiles> Transaction<S> {
894 #[allow(unused)]
918 pub fn stats_schema(&self) -> DeltaResult<SchemaRef> {
919 let stats_schemas = self
920 .effective_table_config
921 .build_expected_stats_schemas(self.physical_clustering_columns.as_deref(), None)?;
922 Ok(stats_schemas.physical)
923 }
924
925 #[allow(unused)]
937 pub fn stats_columns(&self) -> Vec<ColumnName> {
938 self.effective_table_config
939 .physical_stats_column_names(self.physical_clustering_columns.as_deref())
940 }
941
942 fn generate_logical_to_physical(
945 &self,
946 partition_values: Option<&HashMap<String, Scalar>>,
947 ) -> DeltaResult<Expression> {
948 let logical_schema = self.effective_table_config.logical_schema();
949 let mut patch = ExpressionStructPatchBuilder::new();
950 if self
951 .effective_table_config
952 .should_materialize_partition_columns()
953 {
954 let partition_cols: HashSet<&str> = self
955 .effective_table_config
956 .partition_columns()
957 .iter()
958 .map(String::as_str)
959 .collect();
960 let mut predecessor: Option<&str> = None;
964 for field in logical_schema.fields() {
965 let name = field.name().as_str();
966 if partition_cols.contains(name) {
967 let value = partition_values.and_then(|m| m.get(name)).ok_or_else(|| {
968 Error::internal_error(format!(
969 "partition column '{name}' missing while building logical-to-physical \
970 expression"
971 ))
972 })?;
973 let literal = lit(value.clone());
974 patch = match predecessor {
975 Some(predecessor) => patch.insert_after(predecessor, literal),
976 None => patch.prepend(literal),
977 };
978 } else if *field.data_type() != DataType::VOID {
979 predecessor = Some(name);
980 }
981 }
982 }
983 let patch = add_void_stripping(patch, &logical_schema);
984 Expression::struct_patch(patch)
985 }
986
987 pub fn logical_partition_columns(&self) -> &[String] {
989 self.effective_table_config.partition_columns()
990 }
991
992 #[cfg(feature = "column-defaults-in-dev")]
1011 pub fn column_defaults(&self) -> DeltaResult<HashMap<String, ColumnDefault<'_>>> {
1012 let allow_column_defaults = self
1013 .effective_table_config
1014 .is_feature_enabled(&TableFeature::AllowColumnDefaults);
1015 let mut defaults = HashMap::new();
1016 for field in self.effective_table_config.logical_schema_ref().fields() {
1017 let Some(column_default) = field.column_default()? else {
1018 continue;
1019 };
1020 if !allow_column_defaults {
1021 return Err(Error::generic(format!(
1022 "Field '{}' declares a `CURRENT_DEFAULT` but the table does not enable the \
1023 `allowColumnDefaults` writer feature",
1024 field.name()
1025 )));
1026 }
1027 defaults.insert(field.name().clone(), column_default);
1028 }
1029 Ok(defaults)
1030 }
1031
1032 fn validate_for_data_write(&self) -> DeltaResult<()> {
1041 validate_schema_for_write(&self.effective_table_config.logical_schema())
1042 }
1043
1044 fn shared_write_state(&self) -> DeltaResult<Arc<SharedWriteState>> {
1046 let table_config = &self.effective_table_config;
1047 let props = table_config.table_properties();
1048 Ok(Arc::new(SharedWriteState {
1049 table_root: table_config.table_root().clone(),
1050 logical_schema: table_config.logical_schema_without_partition_columns(),
1051 physical_schema: table_config.physical_write_schema(),
1052 column_mapping_mode: table_config.column_mapping_mode(),
1053 stats_columns: self.stats_columns(),
1054 logical_partition_columns: table_config.partition_columns().to_vec(),
1055 randomize_file_prefixes: props.should_randomize_file_prefixes(),
1056 random_prefix_length: props.random_prefix_length(),
1057 }))
1058 }
1059
1060 pub fn partitioned_write_context(
1103 &self,
1104 partition_values: HashMap<String, Scalar>,
1105 ) -> DeltaResult<WriteContext> {
1106 self.ensure_schema_non_empty_for_write_context()?;
1107 self.validate_for_data_write()?;
1108 let shared = self.shared_write_state()?;
1109 require!(
1110 !shared.logical_partition_columns.is_empty(),
1111 Error::generic("table is not partitioned; use unpartitioned_write_context() instead")
1112 );
1113 let full_logical_schema = self.effective_table_config.logical_schema();
1116 let normalized = validate_partition_values(
1117 &shared.logical_partition_columns,
1118 &full_logical_schema,
1119 partition_values,
1120 )?;
1121
1122 let mut serialized = HashMap::with_capacity(normalized.len());
1124 for logical_name in &shared.logical_partition_columns {
1125 let scalar = normalized.get(logical_name).ok_or_else(|| {
1126 Error::internal_error(format!(
1127 "partition column '{logical_name}' missing after validation"
1128 ))
1129 })?;
1130 let value = serialize_partition_value(scalar)?;
1131 let physical_name = full_logical_schema
1132 .field(logical_name)
1133 .ok_or_else(|| {
1134 Error::internal_error(format!(
1135 "partition column '{logical_name}' not found in schema after validation"
1136 ))
1137 })?
1138 .physical_name(shared.column_mapping_mode)
1139 .to_string();
1140 serialized.insert(physical_name, value);
1141 }
1142 let logical_to_physical = Arc::new(self.generate_logical_to_physical(Some(&normalized))?);
1143
1144 Ok(WriteContext {
1145 shared,
1146 logical_to_physical,
1147 physical_partition_values: serialized,
1148 })
1149 }
1150
1151 pub fn unpartitioned_write_context(&self) -> DeltaResult<WriteContext> {
1156 self.ensure_schema_non_empty_for_write_context()?;
1157 self.validate_for_data_write()?;
1158 let shared = self.shared_write_state()?;
1159 require!(
1160 shared.logical_partition_columns.is_empty(),
1161 Error::generic("table is partitioned; use partitioned_write_context() instead")
1162 );
1163 let logical_to_physical = Arc::new(self.generate_logical_to_physical(None)?);
1164 Ok(WriteContext {
1165 shared,
1166 logical_to_physical,
1167 physical_partition_values: HashMap::new(),
1168 })
1169 }
1170
1171 pub fn add_files(&mut self, add_metadata: Box<dyn EngineData>) {
1177 self.add_files_metadata.push(add_metadata);
1178 }
1179}
1180
1181impl<S> Transaction<S> {
1185 fn validate_add_files_stats(&self, add_files: &[Box<dyn EngineData>]) -> DeltaResult<()> {
1199 if add_files.is_empty() {
1200 return Ok(());
1201 }
1202 if self.effective_table_config.requires_stats_num_records() {
1203 stats_verifier::verify_num_records_present(add_files)?;
1207 }
1208 if let Some(ref clustering_cols) = self.physical_clustering_columns {
1209 if !clustering_cols.is_empty() {
1210 let physical_schema = self.effective_table_config.physical_schema();
1211 let columns_with_types: Vec<(ColumnName, DataType)> = clustering_cols
1212 .iter()
1213 .map(|col| {
1214 let data_type = physical_schema
1215 .fields_of_path(col)?
1216 .last()
1217 .map(|field| field.data_type().clone())
1218 .ok_or_else(|| {
1219 Error::internal_error(format!(
1220 "Required column '{col}' not found in table schema"
1221 ))
1222 })?;
1223 Ok((col.clone(), data_type))
1224 })
1225 .collect::<DeltaResult<_>>()?;
1226 let verifier = StatsColumnVerifier::new(columns_with_types);
1227 verifier.verify(add_files)?;
1228 }
1229 }
1230 Ok(())
1231 }
1232
1233 #[instrument(name = "txn.gen_adds", skip_all, err)]
1235 fn generate_adds<'a>(
1236 &'a self,
1237 engine: &dyn Engine,
1238 commit_version: u64,
1239 ) -> DeltaResult<(
1240 EngineDataResultIterator<'a>,
1241 Option<RowTrackingDomainMetadata>,
1242 )> {
1243 let row_tracking_supported = self.effective_table_config.should_write_row_tracking();
1246
1247 if self.add_files_metadata.is_empty() {
1248 let row_tracking_dm = (row_tracking_supported && self.is_create_table())
1253 .then(RowTrackingDomainMetadata::initial);
1254 return Ok((Box::new(iter::empty()), row_tracking_dm));
1255 }
1256
1257 let commit_version = i64::try_from(commit_version)
1258 .map_err(|_| Error::generic("Commit version too large to fit in i64"))?;
1259
1260 if row_tracking_supported {
1261 self.generate_adds_with_row_tracking(engine, commit_version)
1262 } else {
1263 let add_actions = build_add_actions(
1264 engine,
1265 self.add_files_metadata.iter().map(|a| Ok(a.deref())),
1266 self.add_files_schema().clone(),
1267 self.data_change,
1268 )?;
1269 Ok((Box::new(add_actions), None))
1270 }
1271 }
1272
1273 fn generate_adds_with_row_tracking<'a>(
1281 &'a self,
1282 engine: &dyn Engine,
1283 commit_version: i64,
1284 ) -> DeltaResult<(
1285 EngineDataResultIterator<'a>,
1286 Option<RowTrackingDomainMetadata>,
1287 )> {
1288 let row_id_high_water_mark = if self.is_create_table() {
1289 None
1290 } else {
1291 RowTrackingDomainMetadata::get_high_water_mark(self.read_snapshot()?, engine)?
1292 };
1293
1294 let mut row_tracking_visitor =
1296 RowTrackingVisitor::new(row_id_high_water_mark, Some(self.add_files_metadata.len()));
1297
1298 for add_files_batch in &self.add_files_metadata {
1301 row_tracking_visitor.visit_rows_of(add_files_batch.deref())?;
1302 }
1303
1304 let RowTrackingVisitor {
1307 base_row_id_batches,
1308 row_id_high_water_mark,
1309 } = row_tracking_visitor;
1310
1311 let extended_add_files = self.add_files_metadata.iter().zip(base_row_id_batches).map(
1313 move |(add_files_batch, base_row_ids)| {
1314 let commit_versions = vec![commit_version; base_row_ids.len()];
1315 let base_row_ids_array =
1316 ArrayData::try_new(ArrayType::new(DataType::LONG, true), base_row_ids)?;
1317 let commit_versions_array =
1318 ArrayData::try_new(ArrayType::new(DataType::LONG, true), commit_versions)?;
1319
1320 let row_tracking_schema =
1321 with_row_tracking_cols(&Arc::new(StructType::new_unchecked(vec![])))?;
1322 add_files_batch.append_columns(
1323 row_tracking_schema,
1324 vec![base_row_ids_array, commit_versions_array],
1325 )
1326 },
1327 );
1328
1329 let add_actions = build_add_actions(
1331 engine,
1332 extended_add_files,
1333 with_row_tracking_cols(self.add_files_schema())?,
1334 self.data_change,
1335 )?;
1336
1337 let row_tracking_domain_metadata: RowTrackingDomainMetadata =
1339 RowTrackingDomainMetadata::new(row_id_high_water_mark);
1340
1341 Ok((Box::new(add_actions), Some(row_tracking_domain_metadata)))
1342 }
1343
1344 fn into_committed(
1345 self,
1346 file_meta: FileMeta,
1347 crc_delta: CrcDelta,
1348 ) -> DeltaResult<CommittedTransaction> {
1349 let parsed_commit = ParsedLogPath::parse_commit(file_meta)?;
1350
1351 let commit_version = parsed_commit.version;
1352
1353 let (post_commit_stats, post_commit_snapshot) = match &self.read_snapshot_opt {
1354 Some(snap) => {
1355 let stats = PostCommitStats {
1357 commits_since_checkpoint: snap.log_segment().commits_since_checkpoint() + 1,
1358 commits_since_log_compaction: snap
1359 .log_segment()
1360 .commits_since_log_compaction_or_checkpoint()
1361 + 1,
1362 };
1363 let snapshot = snap.new_post_commit(parsed_commit, crc_delta)?;
1364 (stats, Arc::new(snapshot))
1365 }
1366 None => {
1367 let log_root = self
1369 .effective_table_config
1370 .table_root()
1371 .join("_delta_log/")?;
1372 let log_segment = LogSegment::new_for_version_zero(log_root, parsed_commit)?;
1373 let crc = crc_delta.into_crc_for_version_zero().ok_or_else(|| {
1374 Error::internal_error("CREATE TABLE CRC delta is missing protocol or metadata")
1375 })?;
1376 let stats = PostCommitStats {
1377 commits_since_checkpoint: 1,
1378 commits_since_log_compaction: 1,
1379 };
1380 let snapshot = Snapshot::new_with_crc(
1381 log_segment,
1382 self.effective_table_config,
1383 Some(Arc::new(crc)),
1384 )?;
1385 (stats, Arc::new(snapshot))
1386 }
1387 };
1388
1389 Ok(CommittedTransaction {
1390 commit_version,
1391 post_commit_stats,
1392 post_commit_snapshot: Some(post_commit_snapshot),
1393 })
1394 }
1395
1396 fn build_crc_delta(
1399 &self,
1400 file_stats: FileStatsDelta,
1401 in_commit_timestamp: Option<i64>,
1402 dm_changes: Vec<DomainMetadata>,
1403 ) -> DeltaResult<CrcDelta> {
1404 let domain_metadata = dm_changes
1409 .into_iter()
1410 .map(|dm| (dm.domain().to_string(), dm))
1411 .collect();
1412 let set_transactions = self
1413 .set_transactions
1414 .iter()
1415 .map(|txn| (txn.app_id.clone(), txn.clone()))
1416 .collect();
1417 let is_incremental_safe = self
1423 .operation
1424 .as_deref()
1425 .is_some_and(is_incremental_safe_operation);
1426 Ok(CrcDelta {
1427 file_stats,
1428 protocol: self
1429 .should_emit_protocol
1430 .then(|| self.effective_table_config.protocol().clone()),
1431 metadata: self
1432 .should_emit_metadata
1433 .then(|| self.effective_table_config.metadata().clone()),
1434 domain_metadata,
1435 set_transactions,
1436 in_commit_timestamp,
1437 is_incremental_safe,
1438 })
1439 }
1440
1441 fn into_conflicted(self, conflict_version: Version) -> ConflictedTransaction<S> {
1442 ConflictedTransaction {
1443 transaction: self,
1444 conflict_version,
1445 }
1446 }
1447
1448 fn into_retryable(self, error: Error) -> RetryableTransaction<S> {
1449 RetryableTransaction {
1450 transaction: self,
1451 error,
1452 }
1453 }
1454
1455 #[instrument(name = "txn.gen_removes", skip_all, err)]
1476 fn generate_remove_actions<'a>(
1477 &'a self,
1478 engine: &dyn Engine,
1479 remove_files_metadata: impl Iterator<Item = &'a FilteredEngineData> + Send + 'a,
1480 columns_to_drop: &'a [&str],
1481 ) -> DeltaResult<impl Iterator<Item = DeltaResult<FilteredEngineData>> + Send + 'a> {
1482 if self.is_create_table() && !self.remove_files_metadata.is_empty() {
1485 return Err(Error::internal_error(
1486 "CREATE TABLE transaction cannot have remove actions",
1487 ));
1488 }
1489
1490 let input_schema = scan_row_schema();
1491 let target_schema = schema_with_all_fields_nullable(&LOG_REMOVE_SCHEMA);
1492 let evaluation_handler = engine.evaluation_handler();
1493
1494 let make_eval = |coalesce_stats_with_parsed: bool| -> DeltaResult<_> {
1495 let patch = build_remove_struct_patch(
1496 self.commit_timestamp,
1497 self.data_change,
1498 columns_to_drop,
1499 coalesce_stats_with_parsed,
1500 )?;
1501 let expr = Arc::new(Expression::struct_from([Expression::struct_patch(patch)?]));
1502 evaluation_handler.new_expression_evaluator(
1503 input_schema.clone(),
1504 expr,
1505 target_schema.clone().into(),
1506 )
1507 };
1508
1509 let base_eval = Arc::new(make_eval(false)?);
1515 let stats_parsed_eval = Arc::new(make_eval(true)?);
1516 let stats_parsed_col = ColumnName::new([STATS_PARSED_NAME]);
1517
1518 Ok(remove_files_metadata.map(move |file_metadata_batch| {
1519 let data = file_metadata_batch.data();
1520 let evaluator = if data.has_field(&stats_parsed_col) {
1521 &stats_parsed_eval
1522 } else {
1523 &base_eval
1524 };
1525 let updated_engine_data = evaluator.evaluate(data)?;
1526 FilteredEngineData::try_new(
1527 updated_engine_data,
1528 file_metadata_batch.selection_vector().to_vec(),
1529 )
1530 }))
1531 }
1532}
1533
1534fn build_remove_struct_patch(
1546 commit_timestamp: i64,
1547 data_change: bool,
1548 columns_to_drop: &[&str],
1549 coalesce_stats_with_parsed: bool,
1550) -> DeltaResult<ExpressionStructPatch> {
1551 let mut patch = ExpressionStructPatchBuilder::new()
1552 .insert_after("path", lit(commit_timestamp))
1554 .insert_after("path", lit(data_change))
1556 .insert_after("path", lit(true))
1558 .insert_after("path", col!(FILE_CONSTANT_VALUES_NAME, "partitionValues"));
1559
1560 if coalesce_stats_with_parsed {
1561 let coalesce_stats = Expression::coalesce([
1563 col!("stats"),
1564 Expression::unary(ToJson, col!(STATS_PARSED_NAME)),
1565 ]);
1566 patch = patch
1567 .replace("stats", coalesce_stats)
1568 .drop(STATS_PARSED_NAME);
1569 }
1570
1571 patch = patch
1572 .insert_after("stats", col!(FILE_CONSTANT_VALUES_NAME, TAGS_NAME))
1573 .insert_after(
1574 "deletionVector",
1575 col!(FILE_CONSTANT_VALUES_NAME, BASE_ROW_ID_NAME),
1576 )
1577 .insert_after(
1578 "deletionVector",
1579 col!(FILE_CONSTANT_VALUES_NAME, DEFAULT_ROW_COMMIT_VERSION_NAME),
1580 )
1581 .drop(FILE_CONSTANT_VALUES_NAME)
1582 .drop("modificationTime")
1583 .drop_if_exists(PARTITION_VALUES_PARSED_NAME);
1585
1586 for column_to_drop in columns_to_drop {
1587 patch = patch.drop(*column_to_drop);
1588 }
1589
1590 patch.build()
1591}
1592
1593#[derive(Debug)]
1596pub struct PostCommitStats {
1597 pub commits_since_checkpoint: u64,
1600 pub commits_since_log_compaction: u64,
1604}
1605
1606#[derive(Debug)]
1618#[must_use]
1619pub enum CommitResult<S = ExistingTable> {
1620 CommittedTransaction(CommittedTransaction),
1622 ConflictedTransaction(ConflictedTransaction<S>),
1629 RetryableTransaction(RetryableTransaction<S>),
1631}
1632
1633impl<S> CommitResult<S> {
1634 pub fn is_committed(&self) -> bool {
1636 matches!(self, CommitResult::CommittedTransaction(_))
1637 }
1638}
1639
1640impl<S: std::fmt::Debug> CommitResult<S> {
1641 #[cfg(any(test, feature = "test-utils"))]
1643 #[allow(clippy::panic)]
1644 pub fn unwrap_committed(self) -> CommittedTransaction {
1645 match self {
1646 CommitResult::CommittedTransaction(c) => c,
1647 other => panic!("Expected CommittedTransaction, got: {other:?}"),
1648 }
1649 }
1650
1651 #[cfg(any(test, feature = "test-utils"))]
1655 #[allow(clippy::panic, clippy::expect_used)]
1656 pub fn unwrap_post_commit_snapshot(self) -> SnapshotRef {
1657 self.unwrap_committed()
1658 .post_commit_snapshot()
1659 .expect("expected post-commit snapshot")
1660 .clone()
1661 }
1662}
1663
1664#[derive(Debug)]
1672pub struct CommittedTransaction {
1673 commit_version: Version,
1675 post_commit_stats: PostCommitStats,
1677 post_commit_snapshot: Option<SnapshotRef>,
1682}
1683
1684impl CommittedTransaction {
1685 pub fn commit_version(&self) -> Version {
1687 self.commit_version
1688 }
1689
1690 pub fn post_commit_stats(&self) -> &PostCommitStats {
1692 &self.post_commit_stats
1693 }
1694
1695 pub fn post_commit_snapshot(&self) -> Option<&SnapshotRef> {
1697 self.post_commit_snapshot.as_ref()
1698 }
1699}
1700
1701#[derive(Debug)]
1706pub struct ConflictedTransaction<S = ExistingTable> {
1707 #[allow(dead_code)]
1709 transaction: Transaction<S>,
1710 conflict_version: Version,
1711}
1712
1713impl<S> ConflictedTransaction<S> {
1714 pub fn conflict_version(&self) -> Version {
1716 self.conflict_version
1717 }
1718}
1719
1720#[derive(Debug)]
1724pub struct RetryableTransaction<S = ExistingTable> {
1725 pub transaction: Transaction<S>,
1727 pub error: Error,
1729}
1730
1731#[cfg(test)]
1732mod tests {
1733 use std::collections::HashMap;
1734 use std::path::PathBuf;
1735 use std::sync::Mutex;
1736
1737 use ::test_utils::get_column;
1738 use rstest::rstest;
1739 use url::Url;
1740
1741 use super::*;
1742 use crate::actions::deletion_vector::DeletionVectorDescriptor;
1743 use crate::actions::CommitInfo;
1744 use crate::arrow::array::{
1745 ArrayRef, Float64Array, Int32Array, Int64Array, NullArray, StringArray,
1746 };
1747 use crate::arrow::datatypes::{
1748 DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema,
1749 };
1750 use crate::arrow::record_batch::RecordBatch;
1751 use crate::committer::{FileSystemCommitter, PublishMetadata};
1752 use crate::engine::arrow_conversion::{TryFromArrow, TryIntoArrow};
1753 use crate::engine::arrow_data::ArrowEngineData;
1754 use crate::engine::arrow_expression::ArrowEvaluationHandler;
1755 use crate::engine::sync::SyncEngine;
1756 use crate::expressions::{MapData, Scalar, StructData};
1757 use crate::metrics::{MetricEvent, TableType, TransactionCommitFailure};
1758 use crate::object_store::memory::InMemory;
1759 use crate::object_store::path::Path;
1760 use crate::object_store::ObjectStoreExt as _;
1761 use crate::schema::{schema_ref, MapType};
1762 use crate::table_features::ColumnMappingMode;
1763 use crate::transaction::create_table::create_table;
1764 use crate::transaction::data_layout::DataLayout;
1765 use crate::utils::test_utils::{
1766 install_thread_local_metrics_reporter, load_test_table, string_array_to_engine_data,
1767 test_schema_flat, test_schema_nested, test_schema_with_array, test_schema_with_map,
1768 CapturingReporter,
1769 };
1770 use crate::{DeltaResultIterator, EvaluationHandler, Snapshot};
1771
1772 impl Transaction {
1773 fn with_clustering_columns_for_test(mut self, columns: Vec<ColumnName>) -> Self {
1776 self.physical_clustering_columns = Some(columns);
1777 self
1778 }
1779 }
1780
1781 struct IoErrorCommitter;
1783
1784 impl Committer for IoErrorCommitter {
1785 fn commit(
1786 &self,
1787 _engine: &dyn Engine,
1788 _actions: DeltaResultIterator<'_, FilteredEngineData>,
1789 _commit_metadata: CommitMetadata,
1790 ) -> DeltaResult<CommitResponse> {
1791 Err(Error::IOError(std::io::Error::other("simulated IO error")))
1792 }
1793 fn is_catalog_committer(&self) -> bool {
1794 false
1795 }
1796 fn publish(
1797 &self,
1798 _engine: &dyn Engine,
1799 _publish_metadata: PublishMetadata,
1800 ) -> DeltaResult<()> {
1801 Ok(())
1802 }
1803 }
1804
1805 struct GenericErrorCommitter;
1808
1809 impl Committer for GenericErrorCommitter {
1810 fn commit(
1811 &self,
1812 _engine: &dyn Engine,
1813 _actions: DeltaResultIterator<'_, FilteredEngineData>,
1814 _commit_metadata: CommitMetadata,
1815 ) -> DeltaResult<CommitResponse> {
1816 Err(Error::generic("simulated commit error"))
1817 }
1818 fn is_catalog_committer(&self) -> bool {
1819 false
1820 }
1821 fn publish(
1822 &self,
1823 _engine: &dyn Engine,
1824 _publish_metadata: PublishMetadata,
1825 ) -> DeltaResult<()> {
1826 Ok(())
1827 }
1828 }
1829
1830 struct MockCatalogCommitter;
1832
1833 impl Committer for MockCatalogCommitter {
1834 fn commit(
1835 &self,
1836 _engine: &dyn Engine,
1837 _actions: DeltaResultIterator<'_, FilteredEngineData>,
1838 _commit_metadata: CommitMetadata,
1839 ) -> DeltaResult<CommitResponse> {
1840 Ok(CommitResponse::Conflict { version: 0 })
1842 }
1843 fn is_catalog_committer(&self) -> bool {
1844 true
1845 }
1846 fn publish(
1847 &self,
1848 _engine: &dyn Engine,
1849 _publish_metadata: PublishMetadata,
1850 ) -> DeltaResult<()> {
1851 Ok(())
1852 }
1853 }
1854
1855 fn setup_dv_enabled_table() -> (SyncEngine, Arc<Snapshot>) {
1857 let engine = SyncEngine::new();
1858 let path =
1859 std::fs::canonicalize(PathBuf::from("./tests/data/table-with-dv-small/")).unwrap();
1860 let url = url::Url::from_directory_path(path).unwrap();
1861 let snapshot = Snapshot::builder_for(url)
1862 .at_version(1)
1863 .build(&engine)
1864 .unwrap();
1865 (engine, snapshot)
1866 }
1867
1868 fn setup_non_dv_table() -> (SyncEngine, Arc<Snapshot>) {
1869 let engine = SyncEngine::new();
1870 let path =
1871 std::fs::canonicalize(PathBuf::from("./tests/data/table-without-dv-small/")).unwrap();
1872 let url = url::Url::from_directory_path(path).unwrap();
1873 let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
1874 (engine, snapshot)
1875 }
1876
1877 fn setup_dv_supported_but_disabled_table() -> DeltaResult<(Arc<dyn Engine>, Arc<Snapshot>)> {
1878 let storage = Arc::new(InMemory::new());
1879 let table_root = url::Url::parse("memory:///").unwrap();
1880 let engine = Arc::new(SyncEngine::new_with_store(storage.clone()));
1881 let schema_json = serde_json::json!({
1882 "type": "struct",
1883 "fields": [{
1884 "name": "id",
1885 "type": "integer",
1886 "nullable": true,
1887 "metadata": {}
1888 }]
1889 });
1890 let actions = [
1891 r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["deletionVectors"],"writerFeatures":["deletionVectors"]}}"#.to_string(),
1892 serde_json::json!({
1893 "metaData": {
1894 "id": "test-id",
1895 "format": {"provider": "parquet", "options": {}},
1896 "schemaString": schema_json.to_string(),
1897 "partitionColumns": [],
1898 "configuration": {},
1899 "createdTime": 1234567890
1900 }
1901 })
1902 .to_string(),
1903 ]
1904 .join("\n");
1905
1906 let commit_path = Path::from("_delta_log/00000000000000000000.json");
1907 let rt = tokio::runtime::Runtime::new().unwrap();
1908 rt.block_on(storage.put(&commit_path, actions.into()))?;
1909 let engine: Arc<dyn Engine> = engine;
1910 let snapshot = Snapshot::builder_for(table_root).build(engine.as_ref())?;
1911 Ok((engine, snapshot))
1912 }
1913
1914 fn create_test_dv_descriptor(path_suffix: &str) -> DeletionVectorDescriptor {
1917 use crate::actions::deletion_vector::{
1918 DeletionVectorDescriptor, DeletionVectorStorageType,
1919 };
1920 DeletionVectorDescriptor {
1921 storage_type: DeletionVectorStorageType::PersistedRelative,
1922 path_or_inline_dv: format!("dv_{path_suffix}"),
1923 offset: Some(0),
1924 size_in_bytes: 100,
1925 cardinality: 1,
1926 }
1927 }
1928
1929 fn create_dv_transaction(
1930 snapshot: Arc<Snapshot>,
1931 engine: &dyn Engine,
1932 ) -> DeltaResult<Transaction> {
1933 Ok(snapshot
1934 .transaction(Box::new(FileSystemCommitter::new()), engine)?
1935 .with_operation("DELETE".to_string())
1936 .with_engine_info("test_engine"))
1937 }
1938
1939 #[test]
1941 fn test_add_files_schema() -> Result<(), Box<dyn std::error::Error>> {
1942 let engine = SyncEngine::new();
1943 let path =
1944 std::fs::canonicalize(PathBuf::from("./tests/data/table-with-dv-small/")).unwrap();
1945 let url = url::Url::from_directory_path(path).unwrap();
1946 let snapshot = Snapshot::builder_for(url)
1947 .at_version(1)
1948 .build(&engine)
1949 .unwrap();
1950 let txn = snapshot
1951 .transaction(Box::new(FileSystemCommitter::new()), &engine)?
1952 .with_engine_info("default engine");
1953
1954 let schema = txn.add_files_schema();
1955 let expected = StructType::new_unchecked(vec![
1956 StructField::not_null("path", DataType::STRING),
1957 StructField::not_null(
1958 "partitionValues",
1959 MapType::new(DataType::STRING, DataType::STRING, true),
1960 ),
1961 StructField::not_null("size", DataType::LONG),
1962 StructField::not_null("modificationTime", DataType::LONG),
1963 StructField::nullable(
1964 "stats",
1965 DataType::struct_type_unchecked(vec![
1966 StructField::nullable(NUM_RECORDS, DataType::LONG),
1967 StructField::nullable(NULL_COUNT, DataType::struct_type_unchecked(vec![])),
1968 StructField::nullable(MIN_VALUES, DataType::struct_type_unchecked(vec![])),
1969 StructField::nullable(MAX_VALUES, DataType::struct_type_unchecked(vec![])),
1970 StructField::nullable(TIGHT_BOUNDS, DataType::BOOLEAN),
1971 ]),
1972 ),
1973 ]);
1974 assert_eq!(*schema, expected.into());
1975 Ok(())
1976 }
1977
1978 #[rstest]
1979 #[case::base(false)]
1980 #[case::row_tracking(true)]
1981 fn test_add_action_projection_schema(#[case] row_tracking: bool) -> DeltaResult<()> {
1982 let input_schema = if row_tracking {
1983 with_row_tracking_cols(&BASE_ADD_FILES_SCHEMA)?
1984 } else {
1985 BASE_ADD_FILES_SCHEMA.clone()
1986 };
1987 let (schema, _) = build_add_action_projection(input_schema.as_ref(), true)?;
1988 let field_names: Vec<_> = schema.fields().map(|f| f.name().as_str()).collect();
1989 let expected_field_names = if row_tracking {
1990 vec![
1991 "path",
1992 "partitionValues",
1993 "size",
1994 "modificationTime",
1995 "dataChange",
1996 "stats",
1997 "baseRowId",
1998 "defaultRowCommitVersion",
1999 ]
2000 } else {
2001 vec![
2002 "path",
2003 "partitionValues",
2004 "size",
2005 "modificationTime",
2006 "dataChange",
2007 "stats",
2008 ]
2009 };
2010 assert_eq!(field_names, expected_field_names);
2011 assert_eq!(schema.field("dataChange"), Some(&*DATA_CHANGE_COLUMN));
2012 assert_eq!(
2013 schema.field("stats").unwrap().data_type(),
2014 &DataType::STRING
2015 );
2016 Ok(())
2017 }
2018
2019 #[test]
2020 fn test_new_deletion_vector_path() -> Result<(), Box<dyn std::error::Error>> {
2021 let engine = SyncEngine::new();
2022 let path =
2023 std::fs::canonicalize(PathBuf::from("./tests/data/table-with-dv-small/")).unwrap();
2024 let url = url::Url::from_directory_path(path).unwrap();
2025 let snapshot = Snapshot::builder_for(url.clone())
2026 .at_version(1)
2027 .build(&engine)
2028 .unwrap();
2029 let txn = snapshot
2030 .transaction(Box::new(FileSystemCommitter::new()), &engine)?
2031 .with_engine_info("default engine");
2032 let write_context = txn.unpartitioned_write_context().unwrap();
2033
2034 let dv_path1 = write_context.new_deletion_vector_path(String::from(""));
2036 let abs_path1 = dv_path1.absolute_path()?;
2037 assert!(abs_path1.as_str().contains(url.as_str()));
2038
2039 let prefix = String::from("dv_test");
2041 let dv_path2 = write_context.new_deletion_vector_path(prefix.clone());
2042 let abs_path2 = dv_path2.absolute_path()?;
2043 assert!(abs_path2.as_str().contains(url.as_str()));
2044 assert!(abs_path2.as_str().contains(&prefix));
2045
2046 let dv_path3 = write_context.new_deletion_vector_path(prefix.clone());
2048 let abs_path3 = dv_path3.absolute_path()?;
2049 assert_ne!(abs_path2, abs_path3);
2050
2051 Ok(())
2052 }
2053
2054 #[test]
2055 fn write_context_reflects_updated_effective_table_config(
2056 ) -> Result<(), Box<dyn std::error::Error>> {
2057 let (engine, snapshot) = setup_non_dv_table();
2058 let mut txn = snapshot
2059 .clone()
2060 .transaction(Box::new(FileSystemCommitter::new()), &engine)?
2061 .with_engine_info("default engine");
2062
2063 let initial_write_context = txn.unpartitioned_write_context()?;
2066 assert!(!initial_write_context
2067 .logical_schema()
2068 .contains("fresh_column"));
2069
2070 let mut evolved_fields: Vec<StructField> = txn
2071 .effective_table_config
2072 .logical_schema()
2073 .fields()
2074 .cloned()
2075 .collect();
2076 evolved_fields.push(StructField::nullable("fresh_column", DataType::INTEGER));
2077 let evolved_schema = Arc::new(StructType::new_unchecked(evolved_fields));
2078 let evolved_metadata = txn
2079 .effective_table_config
2080 .metadata()
2081 .clone()
2082 .with_schema(evolved_schema.clone())?;
2083 txn.effective_table_config = TableConfiguration::try_new_with_schema(
2084 &txn.effective_table_config,
2085 evolved_metadata,
2086 evolved_schema,
2087 )?;
2088
2089 let updated_write_context = txn.unpartitioned_write_context()?;
2090 assert!(updated_write_context
2091 .logical_schema()
2092 .contains("fresh_column"));
2093 assert!(updated_write_context
2094 .physical_schema()
2095 .contains("fresh_column"));
2096 assert!(!initial_write_context
2097 .logical_schema()
2098 .contains("fresh_column"));
2099
2100 Ok(())
2101 }
2102
2103 #[cfg(feature = "column-defaults-in-dev")]
2104 mod column_defaults {
2105 use super::*;
2106 use crate::schema::{ColumnMetadataKey, MetadataValue};
2107
2108 fn txn_with_schema(schema: StructType) -> Transaction {
2115 txn_with_schema_and_writer_features(schema, [TableFeature::AllowColumnDefaults])
2116 }
2117
2118 fn txn_with_schema_and_writer_features(
2123 schema: StructType,
2124 writer_features: impl IntoIterator<Item = TableFeature>,
2125 ) -> Transaction {
2126 let (engine, snapshot) = setup_non_dv_table();
2127 let mut txn = snapshot
2128 .transaction(Box::new(FileSystemCommitter::new()), &engine)
2129 .unwrap();
2130 let metadata = txn
2131 .effective_table_config
2132 .metadata()
2133 .clone()
2134 .with_schema(Arc::new(schema))
2135 .unwrap();
2136 let protocol =
2137 Protocol::try_new_modern(TableFeature::EMPTY_LIST, writer_features).unwrap();
2138 let version = txn.effective_table_config.version();
2139 txn.effective_table_config = TableConfiguration::try_new_from(
2140 &txn.effective_table_config,
2141 Some(metadata),
2142 Some(protocol),
2143 version,
2144 )
2145 .unwrap();
2146 txn
2147 }
2148
2149 fn field_with_default(name: &str, data_type: DataType, raw_sql: &str) -> StructField {
2151 StructField::nullable(name, data_type).add_metadata([(
2152 ColumnMetadataKey::CurrentDefault.as_ref().to_string(),
2153 MetadataValue::String(raw_sql.to_string()),
2154 )])
2155 }
2156
2157 #[test]
2158 fn collects_present_defaults_and_skips_columns_without_one() {
2159 let schema = StructType::try_new(vec![
2160 field_with_default("parsable", DataType::INTEGER, "42"),
2161 field_with_default("unparsable", DataType::TIMESTAMP, "current_timestamp()"),
2162 StructField::nullable("no_default", DataType::STRING),
2163 ])
2164 .unwrap();
2165 let txn = txn_with_schema(schema);
2166
2167 let defaults = txn.column_defaults().unwrap();
2168 assert_eq!(
2169 defaults.len(),
2170 2,
2171 "only columns with a default are returned"
2172 );
2173 assert!(!defaults.contains_key("no_default"));
2174
2175 let parsable = &defaults["parsable"];
2176 assert_eq!(parsable.raw_sql(), "42");
2177 assert!(parsable.to_scalar().unwrap().is_some());
2178
2179 let unparsable = &defaults["unparsable"];
2180 assert_eq!(unparsable.raw_sql(), "current_timestamp()");
2181 assert!(unparsable.to_scalar().unwrap().is_none());
2182 }
2183
2184 #[test]
2185 fn returns_empty_map_when_no_column_has_a_default() {
2186 let schema = StructType::try_new(vec![
2187 StructField::nullable("a", DataType::INTEGER),
2188 StructField::nullable("b", DataType::STRING),
2189 ])
2190 .unwrap();
2191 let txn = txn_with_schema(schema);
2192 assert!(txn.column_defaults().unwrap().is_empty());
2193 }
2194
2195 #[test]
2196 fn propagates_error_for_malformed_default() {
2197 let field = StructField::nullable("c", DataType::INTEGER).add_metadata([(
2198 ColumnMetadataKey::CurrentDefault.as_ref().to_string(),
2199 MetadataValue::Number(7),
2200 )]);
2201 let schema = StructType::try_new(vec![field]).unwrap();
2202 let txn = txn_with_schema(schema);
2203
2204 let err = txn
2205 .column_defaults()
2206 .expect_err("non-string CURRENT_DEFAULT must error")
2207 .to_string();
2208 assert!(err.contains("non-string"), "got: {err}");
2209 }
2210
2211 #[test]
2212 fn errors_when_default_present_but_feature_not_enabled() {
2213 let schema =
2214 StructType::try_new(vec![field_with_default("c", DataType::INTEGER, "42")])
2215 .unwrap();
2216 let txn = txn_with_schema_and_writer_features(schema, []);
2217
2218 let err = txn
2219 .column_defaults()
2220 .expect_err("a column default without the allowColumnDefaults feature must error")
2221 .to_string();
2222 assert!(err.contains("allowColumnDefaults"), "got: {err}");
2223 }
2224 }
2225
2226 #[test]
2227 fn test_write_context_schemas_exclude_partition_columns(
2228 ) -> Result<(), Box<dyn std::error::Error>> {
2229 let engine = SyncEngine::new();
2230 let path = std::fs::canonicalize(PathBuf::from("./tests/data/basic_partitioned/")).unwrap();
2231 let url = url::Url::from_directory_path(path).unwrap();
2232 let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
2233 let txn = snapshot
2234 .transaction(Box::new(FileSystemCommitter::new()), &engine)?
2235 .with_engine_info("default engine");
2236
2237 let write_context = txn.partitioned_write_context(HashMap::from([(
2238 "letter".to_string(),
2239 Scalar::String("a".into()),
2240 )]))?;
2241 let logical_schema = write_context.logical_schema();
2242 let physical_schema = write_context.physical_schema();
2243
2244 assert!(
2246 !logical_schema.contains("letter"),
2247 "Logical schema should not contain partition column 'letter'"
2248 );
2249 assert!(
2250 !physical_schema.contains("letter"),
2251 "Physical schema should not contain partition column 'letter' (stored in path)"
2252 );
2253
2254 assert!(
2256 logical_schema.contains("number"),
2257 "Logical schema should contain data column 'number'"
2258 );
2259
2260 assert!(
2261 physical_schema.contains("number"),
2262 "Physical schema should contain data column 'number'"
2263 );
2264
2265 Ok(())
2266 }
2267
2268 fn snapshot_and_partitioned_write_context(
2271 table_path: &str,
2272 partition_values: HashMap<String, Scalar>,
2273 ) -> Result<(Arc<Snapshot>, WriteContext), Box<dyn std::error::Error>> {
2274 let engine = SyncEngine::new();
2275 let path = std::fs::canonicalize(PathBuf::from(table_path)).unwrap();
2276 let url = url::Url::from_directory_path(path).unwrap();
2277 let snapshot = Snapshot::builder_for(url).build(&engine)?;
2278 let txn = snapshot
2279 .clone()
2280 .transaction(Box::new(FileSystemCommitter::new()), &engine)?;
2281 let wc = txn.partitioned_write_context(partition_values)?;
2282 Ok((snapshot, wc))
2283 }
2284
2285 fn eval_logical_to_physical(
2288 wc: &WriteContext,
2289 batch: RecordBatch,
2290 ) -> Result<RecordBatch, Box<dyn std::error::Error>> {
2291 let input_schema = StructType::try_from_arrow(batch.schema())?;
2292 let physical_schema = wc.physical_schema();
2293 let l2p = wc.logical_to_physical();
2294
2295 let handler = ArrowEvaluationHandler;
2296 let evaluator = handler.new_expression_evaluator(
2297 input_schema.into(),
2298 l2p,
2299 physical_schema.clone().into(),
2300 )?;
2301 let result = ArrowEngineData::try_from_engine_data(
2302 evaluator.evaluate(&ArrowEngineData::new(batch))?,
2303 )?;
2304 Ok(result.record_batch().clone())
2305 }
2306
2307 #[rstest]
2308 #[case::not_materialized("./tests/data/basic_partitioned/", false)]
2309 #[case::materialized("./tests/data/partitioned_with_materialize_feature/", true)]
2310 fn test_partition_columns_materialized_in_logical_to_physical(
2311 #[case] table_path: &str,
2312 #[case] materialized: bool,
2313 ) -> Result<(), Box<dyn std::error::Error>> {
2314 let (snapshot, wc) = snapshot_and_partitioned_write_context(
2315 table_path,
2316 HashMap::from([("letter".to_string(), Scalar::String("a".into()))]),
2317 )?;
2318 assert_eq!(
2319 snapshot
2320 .table_configuration()
2321 .protocol()
2322 .has_table_feature(&TableFeature::MaterializePartitionColumns),
2323 materialized
2324 );
2325
2326 let input_schema = Arc::new(ArrowSchema::new(vec![
2328 ArrowField::new("number", ArrowDataType::Int64, true),
2329 ArrowField::new("a_float", ArrowDataType::Float64, true),
2330 ]));
2331 let batch = RecordBatch::try_new(
2332 input_schema,
2333 vec![
2334 Arc::new(Int64Array::from(vec![42])) as ArrayRef,
2335 Arc::new(Float64Array::from(vec![1.5])),
2336 ],
2337 )?;
2338 let rb = eval_logical_to_physical(&wc, batch)?;
2339
2340 let rb_schema = rb.schema();
2341 let names: Vec<&str> = rb_schema
2342 .fields()
2343 .iter()
2344 .map(|f| f.name().as_str())
2345 .collect();
2346 if materialized {
2347 assert_eq!(names, vec!["letter", "number", "a_float"]);
2348 assert_eq!(get_column!(rb, "letter", StringArray).value(0), "a");
2349 } else {
2350 assert_eq!(names, vec!["number", "a_float"]);
2351 }
2352 Ok(())
2353 }
2354
2355 #[rstest]
2356 #[case::cm_none(ColumnMappingMode::None)]
2357 #[case::cm_name(ColumnMappingMode::Name)]
2358 #[case::cm_id(ColumnMappingMode::Id)]
2359 fn test_materialized_partition_column_insert(
2360 #[case] cm_mode: ColumnMappingMode,
2361 ) -> Result<(), Box<dyn std::error::Error>> {
2362 let cm = match cm_mode {
2363 ColumnMappingMode::None => "none",
2364 ColumnMappingMode::Name => "name",
2365 ColumnMappingMode::Id => "id",
2366 };
2367 let engine: Arc<dyn Engine> =
2368 Arc::new(SyncEngine::new_with_store(Arc::new(InMemory::new())));
2369 let schema = Arc::new(StructType::try_new(vec![
2371 StructField::nullable("p1", DataType::STRING),
2372 StructField::nullable("p2", DataType::INTEGER),
2373 StructField::nullable("d1", DataType::INTEGER),
2374 StructField::nullable("v", DataType::VOID),
2375 StructField::nullable("p3", DataType::STRING),
2376 StructField::nullable("p4", DataType::INTEGER),
2377 StructField::nullable("d2", DataType::INTEGER),
2378 ])?);
2379 let txn = create_table("memory:///t", schema, "DefaultEngine")
2380 .with_data_layout(DataLayout::partitioned(["p1", "p2", "p3", "p4"]))
2381 .with_table_properties([
2382 ("delta.feature.materializePartitionColumns", "supported"),
2383 ("delta.columnMapping.mode", cm),
2384 ])
2385 .build(engine.as_ref(), Box::new(FileSystemCommitter::new()))?;
2386
2387 let wc = txn.partitioned_write_context(HashMap::from([
2388 ("p1".to_string(), Scalar::String("aa".into())),
2389 ("p2".to_string(), Scalar::Integer(7)),
2390 ("p3".to_string(), Scalar::String("cc".into())),
2391 ("p4".to_string(), Scalar::Integer(9)),
2392 ]))?;
2393
2394 let input_schema = Arc::new(ArrowSchema::new(vec![
2397 ArrowField::new("d1", ArrowDataType::Int32, true),
2398 ArrowField::new("v", ArrowDataType::Null, true),
2399 ArrowField::new("d2", ArrowDataType::Int32, true),
2400 ]));
2401 let batch = RecordBatch::try_new(
2402 input_schema,
2403 vec![
2404 Arc::new(Int32Array::from(vec![10])) as ArrayRef,
2405 Arc::new(NullArray::new(1)),
2406 Arc::new(Int32Array::from(vec![20])),
2407 ],
2408 )?;
2409 let rb = eval_logical_to_physical(&wc, batch)?;
2410
2411 let rb_schema = rb.schema();
2414 let names: Vec<&str> = rb_schema
2415 .fields()
2416 .iter()
2417 .map(|f| f.name().as_str())
2418 .collect();
2419 let physical_schema = wc.physical_schema();
2420 let expected_names: Vec<&str> = physical_schema
2421 .fields()
2422 .map(|f| f.name().as_str())
2423 .collect();
2424 assert_eq!(names, expected_names);
2425
2426 assert_eq!(get_column!(rb, names[0], StringArray).value(0), "aa"); assert_eq!(get_column!(rb, names[1], Int32Array).value(0), 7); assert_eq!(get_column!(rb, names[2], Int32Array).value(0), 10); assert_eq!(get_column!(rb, names[3], StringArray).value(0), "cc"); assert_eq!(get_column!(rb, names[4], Int32Array).value(0), 9); assert_eq!(get_column!(rb, names[5], Int32Array).value(0), 20); Ok(())
2434 }
2435
2436 #[test]
2438 fn test_physical_schema_includes_partition_columns_when_materialized(
2439 ) -> Result<(), Box<dyn std::error::Error>> {
2440 let (_snapshot, write_context) = snapshot_and_partitioned_write_context(
2441 "./tests/data/partitioned_with_materialize_feature/",
2442 HashMap::from([("letter".to_string(), Scalar::String("a".into()))]),
2443 )?;
2444 let physical_schema = write_context.physical_schema();
2445
2446 assert!(
2447 physical_schema.contains("letter"),
2448 "Partition column 'letter' should be in physical schema when materialized"
2449 );
2450 assert!(
2451 physical_schema.contains("number"),
2452 "Non-partition column 'number' should be in physical schema"
2453 );
2454 Ok(())
2455 }
2456
2457 #[rstest]
2459 #[case::partitioned_on_unpartitioned(
2460 "./tests/data/table-without-dv-small/",
2461 true,
2462 "not partitioned"
2463 )]
2464 #[case::unpartitioned_on_partitioned(
2465 "./tests/data/basic_partitioned/",
2466 false,
2467 "table is partitioned"
2468 )]
2469 fn test_wrong_write_context_method_returns_error(
2470 #[case] table_path: &str,
2471 #[case] call_partitioned: bool,
2472 #[case] expected_msg: &str,
2473 ) -> Result<(), Box<dyn std::error::Error>> {
2474 let engine = SyncEngine::new();
2475 let path = std::fs::canonicalize(PathBuf::from(table_path)).unwrap();
2476 let url = url::Url::from_directory_path(path).unwrap();
2477 let snapshot = Snapshot::builder_for(url).build(&engine)?;
2478 let txn = snapshot.transaction(Box::new(FileSystemCommitter::new()), &engine)?;
2479 let result = if call_partitioned {
2480 txn.partitioned_write_context(HashMap::from([("x".to_string(), Scalar::Integer(1))]))
2481 } else {
2482 txn.unpartitioned_write_context()
2483 };
2484 let err = result.unwrap_err().to_string();
2485 assert!(
2486 err.contains(expected_msg),
2487 "expected '{expected_msg}' in error, got: {err}"
2488 );
2489 Ok(())
2490 }
2491
2492 #[test]
2495 fn test_update_deletion_vectors_unsupported_table() -> Result<(), Box<dyn std::error::Error>> {
2496 let (engine, snapshot) = setup_non_dv_table();
2497 let mut txn = create_dv_transaction(snapshot, &engine)?;
2498
2499 let dv_map = HashMap::new();
2500 let result = txn.update_deletion_vectors(dv_map, std::iter::empty());
2501
2502 let err = result.expect_err("Should fail on table without DV support");
2503 let err_msg = err.to_string();
2504 assert!(
2505 err_msg.contains("Deletion vector")
2506 && (err_msg.contains("require") || err_msg.contains("version")),
2507 "Expected protocol error about DV requirements, got: {err_msg}"
2508 );
2509 Ok(())
2510 }
2511
2512 #[test]
2513 fn test_update_deletion_vectors_requires_enablement_property(
2514 ) -> Result<(), Box<dyn std::error::Error>> {
2515 let (engine, snapshot) = setup_dv_supported_but_disabled_table()?;
2516 let mut txn = create_dv_transaction(snapshot, engine.as_ref())?;
2517
2518 let err = txn
2519 .update_deletion_vectors(HashMap::new(), std::iter::empty())
2520 .expect_err("DV updates should require delta.enableDeletionVectors=true");
2521
2522 assert!(
2523 matches!(err, Error::Unsupported(_)),
2524 "unexpected error: {err}"
2525 );
2526 assert!(
2527 err.to_string().contains("delta.enableDeletionVectors"),
2528 "error should mention the enablement property, got: {err}"
2529 );
2530 Ok(())
2531 }
2532
2533 #[test]
2536 fn test_update_deletion_vectors_mismatch_count() -> Result<(), Box<dyn std::error::Error>> {
2537 let (engine, snapshot) = setup_dv_enabled_table();
2538 let mut txn = create_dv_transaction(snapshot, &engine)?;
2539
2540 let mut dv_map = HashMap::new();
2541 let descriptor = create_test_dv_descriptor("non_existent");
2542 dv_map.insert("non_existent_file.parquet".to_string(), descriptor);
2543
2544 let result = txn.update_deletion_vectors(dv_map, std::iter::empty());
2545
2546 assert!(
2547 result.is_err(),
2548 "Should fail when DV descriptors don't match scan files"
2549 );
2550 let err_msg = result.unwrap_err().to_string();
2551 assert!(
2552 err_msg.contains("matched") && err_msg.contains("does not match"),
2553 "Expected error about mismatched count (expected 1 descriptor, 0 matched files), got: {err_msg}");
2554 Ok(())
2555 }
2556
2557 #[test]
2559 fn test_update_deletion_vectors_mismatch_does_not_mutate_transaction(
2560 ) -> Result<(), Box<dyn std::error::Error>> {
2561 let (engine, snapshot) = setup_dv_enabled_table();
2562 let mut txn = create_dv_transaction(snapshot.clone(), &engine)?;
2563 let scan = snapshot.scan_builder().build()?;
2564 let scan_metadata = scan
2565 .scan_metadata(&engine)?
2566 .collect::<DeltaResult<Vec<_>>>()?;
2567
2568 let mut paths = Vec::new();
2569 for metadata in &scan_metadata {
2570 paths =
2571 metadata.visit_scan_files(paths, |paths, scan_file| paths.push(scan_file.path))?;
2572 }
2573 let existing_path = paths
2574 .into_iter()
2575 .next()
2576 .ok_or_else(|| Error::generic("expected at least one scan file"))?;
2577
2578 let mut dv_map = HashMap::new();
2579 dv_map.insert(existing_path, create_test_dv_descriptor("matched"));
2580 dv_map.insert(
2581 "non_existent_file.parquet".to_string(),
2582 create_test_dv_descriptor("missing"),
2583 );
2584
2585 let result = txn.update_deletion_vectors(
2586 dv_map,
2587 scan_metadata
2588 .into_iter()
2589 .map(|metadata| Ok(metadata.scan_files)),
2590 );
2591
2592 assert!(
2593 result.is_err(),
2594 "Should fail when only some DV descriptors match scan files"
2595 );
2596 assert!(
2597 txn.dv_matched_files.is_empty(),
2598 "Failed DV update should not leave staged file updates"
2599 );
2600 Ok(())
2601 }
2602
2603 #[test]
2604 fn test_update_deletion_vectors_iter_error_does_not_mutate_transaction(
2605 ) -> Result<(), Box<dyn std::error::Error>> {
2606 let (engine, snapshot) = setup_dv_enabled_table();
2607 let mut txn = create_dv_transaction(snapshot.clone(), &engine)?;
2608 txn.dv_matched_files
2609 .push(FilteredEngineData::with_all_rows_selected(
2610 string_array_to_engine_data(StringArray::from(vec!["sentinel"])),
2611 ));
2612 let staged_len_before = txn.dv_matched_files.len();
2613 let scan = snapshot.scan_builder().build()?;
2614 let scan_metadata = scan
2615 .scan_metadata(&engine)?
2616 .collect::<DeltaResult<Vec<_>>>()?;
2617
2618 let mut paths = Vec::new();
2619 for metadata in &scan_metadata {
2620 paths =
2621 metadata.visit_scan_files(paths, |paths, scan_file| paths.push(scan_file.path))?;
2622 }
2623 let existing_path = paths
2624 .into_iter()
2625 .next()
2626 .ok_or_else(|| Error::generic("expected at least one scan file"))?;
2627
2628 let mut dv_map = HashMap::new();
2629 dv_map.insert(existing_path, create_test_dv_descriptor("matched"));
2630
2631 let result = txn.update_deletion_vectors(
2632 dv_map,
2633 scan_metadata
2634 .into_iter()
2635 .map(|metadata| Ok(metadata.scan_files))
2636 .chain(std::iter::once(Err(Error::generic(
2637 "simulated scan metadata failure",
2638 )))),
2639 );
2640
2641 assert!(result.is_err(), "iterator error should propagate");
2642 assert_eq!(
2643 txn.dv_matched_files.len(),
2644 staged_len_before,
2645 "Failed DV update should not stage additional file updates"
2646 );
2647 Ok(())
2648 }
2649
2650 #[test]
2653 fn test_update_deletion_vectors_empty_inputs() -> Result<(), Box<dyn std::error::Error>> {
2654 let (engine, snapshot) = setup_dv_enabled_table();
2655 let mut txn = create_dv_transaction(snapshot, &engine)?;
2656
2657 let dv_map = HashMap::new();
2658 let result = txn.update_deletion_vectors(dv_map, std::iter::empty());
2659
2660 assert!(
2661 result.is_ok(),
2662 "Empty DV updates should succeed as no-op, got error: {result:?}"
2663 );
2664
2665 Ok(())
2666 }
2667
2668 fn add_dummy_file<S: SupportsDataFiles>(txn: &mut Transaction<S>) {
2672 let data = string_array_to_engine_data(StringArray::from(vec!["dummy"]));
2673 txn.add_files(data);
2674 }
2675
2676 fn create_existing_table_txn(
2677 ) -> DeltaResult<(Arc<dyn Engine>, Transaction, Option<tempfile::TempDir>)> {
2678 let (engine, snapshot, tempdir) = load_test_table("table-without-dv-small")?;
2679 let txn = snapshot.transaction(Box::new(FileSystemCommitter::new()), engine.as_ref())?;
2680 Ok((engine, txn, tempdir))
2681 }
2682
2683 #[test]
2684 fn test_validate_blind_append_success() -> DeltaResult<()> {
2685 let (_engine, mut txn, _tempdir) = create_existing_table_txn()?;
2686 txn = txn.with_blind_append();
2687 add_dummy_file(&mut txn);
2688 txn.validate_blind_append_semantics()?;
2689 Ok(())
2690 }
2691
2692 #[test]
2693 fn test_validate_blind_append_requires_adds() -> DeltaResult<()> {
2694 let (_engine, mut txn, _tempdir) = create_existing_table_txn()?;
2695 txn = txn.with_blind_append();
2696 let result = txn.validate_blind_append_semantics();
2697 assert!(matches!(result, Err(Error::InvalidTransactionState(_))));
2698 Ok(())
2699 }
2700
2701 #[test]
2702 fn test_validate_blind_append_requires_data_change() -> DeltaResult<()> {
2703 let (_engine, mut txn, _tempdir) = create_existing_table_txn()?;
2704 txn = txn.with_blind_append();
2705 txn.set_data_change(false);
2706 add_dummy_file(&mut txn);
2707 let result = txn.validate_blind_append_semantics();
2708 assert!(matches!(result, Err(Error::InvalidTransactionState(_))));
2709 Ok(())
2710 }
2711
2712 #[test]
2713 fn test_validate_blind_append_rejects_removes() -> DeltaResult<()> {
2714 let (_engine, mut txn, _tempdir) = create_existing_table_txn()?;
2715 txn = txn.with_blind_append();
2716 add_dummy_file(&mut txn);
2717 let remove_data = FilteredEngineData::with_all_rows_selected(string_array_to_engine_data(
2718 StringArray::from(vec!["remove"]),
2719 ));
2720 txn.remove_files(remove_data);
2721 let result = txn.validate_blind_append_semantics();
2722 assert!(matches!(result, Err(Error::InvalidTransactionState(_))));
2723 Ok(())
2724 }
2725
2726 #[test]
2727 fn test_validate_blind_append_rejects_dv_updates() -> DeltaResult<()> {
2728 let (_engine, mut txn, _tempdir) = create_existing_table_txn()?;
2729 txn = txn.with_blind_append();
2730 add_dummy_file(&mut txn);
2731 let dv_data = FilteredEngineData::with_all_rows_selected(string_array_to_engine_data(
2732 StringArray::from(vec!["dv"]),
2733 ));
2734 txn.dv_matched_files.push(dv_data);
2735 let result = txn.validate_blind_append_semantics();
2736 assert!(matches!(result, Err(Error::InvalidTransactionState(_))));
2737 Ok(())
2738 }
2739
2740 #[test]
2741 fn test_validate_blind_append_rejects_create_table() -> DeltaResult<()> {
2742 let tempdir = tempfile::tempdir()?;
2743 let schema = schema_ref! { nullable "id": INTEGER };
2744 let engine = Arc::new(crate::engine::sync::SyncEngine::new());
2745 let mut txn = create_table(
2746 tempdir.path().to_str().expect("valid temp path"),
2747 schema,
2748 "test_engine",
2749 )
2750 .build(engine.as_ref(), Box::new(FileSystemCommitter::new()))?;
2751 txn.is_blind_append = true;
2754 add_dummy_file(&mut txn);
2755 let result = txn.validate_blind_append_semantics();
2756 assert!(matches!(result, Err(Error::InvalidTransactionState(_))));
2757 Ok(())
2758 }
2759
2760 #[test]
2761 fn test_blind_append_sets_commit_info_flag() -> Result<(), Box<dyn std::error::Error>> {
2762 let commit_info = CommitInfo::new(1, None, None, None, true);
2763 assert_eq!(commit_info.is_blind_append, Some(true));
2764
2765 let commit_info_false = CommitInfo::new(1, None, None, None, false);
2766 assert_eq!(commit_info_false.is_blind_append, None);
2767 Ok(())
2768 }
2769
2770 #[test]
2771 fn test_blind_append_commit_rejects_no_adds() -> DeltaResult<()> {
2772 let (_engine, mut txn, _tempdir) = create_existing_table_txn()?;
2773 txn = txn.with_blind_append();
2774 let err = txn
2776 .commit(_engine.as_ref())
2777 .expect_err("Blind append with no adds should fail");
2778 assert!(
2779 err.to_string()
2780 .contains("Blind append requires at least one added data file"),
2781 "Unexpected error: {err}"
2782 );
2783 Ok(())
2784 }
2785
2786 #[test]
2787 fn test_blind_append_commit_success() -> DeltaResult<()> {
2788 let (engine, mut txn, _tempdir) = create_existing_table_txn()?;
2789 txn = txn.with_blind_append();
2790 add_dummy_file(&mut txn);
2791 let result = txn.commit(engine.as_ref());
2795 if let Err(e) = result {
2797 assert!(
2798 !matches!(e, Error::InvalidTransactionState(_)),
2799 "Blind append validation should have passed, got: {e}"
2800 );
2801 }
2802 Ok(())
2803 }
2804
2805 #[test]
2811 fn test_commit_io_error_returns_retryable_transaction() -> DeltaResult<()> {
2812 let (engine, snapshot, _tempdir) = load_test_table("table-without-dv-small")?;
2813 let mut txn = snapshot.transaction(Box::new(IoErrorCommitter), engine.as_ref())?;
2814 add_dummy_file(&mut txn);
2815 let result = txn.commit(engine.as_ref())?;
2816 assert!(
2817 matches!(result, CommitResult::RetryableTransaction(_)),
2818 "Expected RetryableTransaction, got: {result:?}"
2819 );
2820 if let CommitResult::RetryableTransaction(retryable) = result {
2821 assert!(
2822 retryable.error.to_string().contains("simulated IO error"),
2823 "Unexpected error: {}",
2824 retryable.error
2825 );
2826 }
2827 Ok(())
2828 }
2829
2830 #[test]
2831 fn test_existing_table_txn_debug() -> DeltaResult<()> {
2832 let (_engine, txn, _tempdir) = create_existing_table_txn()?;
2833 let debug_str = format!("{txn:?}");
2834 assert!(
2836 debug_str.contains("Transaction") && debug_str.contains("read_snapshot version"),
2837 "Debug output should contain Transaction info: {debug_str}"
2838 );
2839 assert!(
2841 !debug_str.contains("create_table"),
2842 "Existing table debug should not contain create_table: {debug_str}"
2843 );
2844 Ok(())
2845 }
2846
2847 #[rstest]
2850 #[case::flat_none(test_schema_flat(), ColumnMappingMode::None)]
2851 #[case::flat_name(test_schema_flat(), ColumnMappingMode::Name)]
2852 #[case::flat_id(test_schema_flat(), ColumnMappingMode::Id)]
2853 #[case::nested_none(test_schema_nested(), ColumnMappingMode::None)]
2854 #[case::nested_name(test_schema_nested(), ColumnMappingMode::Name)]
2855 #[case::nested_id(test_schema_nested(), ColumnMappingMode::Id)]
2856 #[case::map_none(test_schema_with_map(), ColumnMappingMode::None)]
2857 #[case::map_name(test_schema_with_map(), ColumnMappingMode::Name)]
2858 #[case::map_id(test_schema_with_map(), ColumnMappingMode::Id)]
2859 #[case::array_none(test_schema_with_array(), ColumnMappingMode::None)]
2860 #[case::array_name(test_schema_with_array(), ColumnMappingMode::Name)]
2861 #[case::array_id(test_schema_with_array(), ColumnMappingMode::Id)]
2862 fn test_physical_schema_column_mapping(
2863 #[case] schema: SchemaRef,
2864 #[case] mode: ColumnMappingMode,
2865 ) -> DeltaResult<()> {
2866 let (_engine, txn) = crate::utils::test_utils::setup_column_mapping_txn(schema, mode)?;
2867 let write_context = txn.unpartitioned_write_context().unwrap();
2868 crate::utils::test_utils::validate_physical_schema_column_mapping(
2869 write_context.logical_schema(),
2870 write_context.physical_schema(),
2871 mode,
2872 );
2873 Ok(())
2874 }
2875
2876 fn build_test_record_batch() -> DeltaResult<Box<dyn EngineData>> {
2878 let schema = test_schema_nested();
2879 let tag_type = MapType::new(DataType::STRING, DataType::STRING, true);
2880 let score_type = ArrayType::new(DataType::INTEGER, true);
2881 let info_fields = vec![
2882 StructField::nullable("name", DataType::STRING),
2883 StructField::nullable("age", DataType::INTEGER),
2884 StructField::nullable("tags", tag_type.clone()),
2885 StructField::nullable("scores", score_type.clone()),
2886 ];
2887 let info1 = Scalar::Struct(StructData::try_new(
2888 info_fields.clone(),
2889 vec![
2890 "alice".into(),
2891 30i32.into(),
2892 Scalar::Map(MapData::try_new(tag_type.clone(), [("k1", "v1")])?),
2893 Scalar::Array(ArrayData::try_new(score_type.clone(), [10i32, 20i32])?),
2894 ],
2895 )?);
2896 let info2 = Scalar::Struct(StructData::try_new(
2897 info_fields,
2898 vec![
2899 "bob".into(),
2900 25i32.into(),
2901 Scalar::Map(MapData::try_new(tag_type, [("k2", "v2")])?),
2902 Scalar::Array(ArrayData::try_new(score_type, [30i32])?),
2903 ],
2904 )?);
2905 ArrowEvaluationHandler.create_many(schema, &[&[1i64.into(), info1], &[2i64.into(), info2]])
2906 }
2907
2908 fn validate_logical_to_physical_transform(mode: ColumnMappingMode) -> DeltaResult<()> {
2913 let schema = test_schema_nested();
2914 let (_engine, txn) = crate::utils::test_utils::setup_column_mapping_txn(schema, mode)?;
2915 let write_context = txn.unpartitioned_write_context().unwrap();
2916 let logical_schema = write_context.logical_schema();
2917 let physical_schema = write_context.physical_schema();
2918 let logical_to_physical_expression = write_context.logical_to_physical();
2919
2920 if mode != ColumnMappingMode::None {
2921 assert_ne!(
2922 logical_schema, physical_schema,
2923 "Physical schema should differ from logical schema when column mapping is enabled"
2924 );
2925 }
2926
2927 let data = build_test_record_batch()?;
2928
2929 let input_schema: SchemaRef = logical_schema.clone();
2931 let handler = ArrowEvaluationHandler;
2932 let evaluator = handler.new_expression_evaluator(
2933 input_schema,
2934 logical_to_physical_expression.clone(),
2935 physical_schema.clone().into(),
2936 )?;
2937 let result = evaluator.evaluate(data.as_ref())?;
2938 let result = ArrowEngineData::try_from_engine_data(result)?;
2939 let result_batch = result.record_batch();
2940
2941 let expected_arrow_schema: ArrowSchema = physical_schema.as_ref().try_into_arrow()?;
2943 assert_eq!(result_batch.schema().as_ref(), &expected_arrow_schema);
2944
2945 let id_col = result_batch
2947 .column(0)
2948 .as_any()
2949 .downcast_ref::<Int64Array>()
2950 .expect("id column should be Int64");
2951 assert_eq!(id_col.values(), &[1i64, 2]);
2952
2953 Ok(())
2954 }
2955
2956 #[rstest]
2957 #[case::name_mode(ColumnMappingMode::Name)]
2958 #[case::id_mode(ColumnMappingMode::Id)]
2959 #[case::none_mode(ColumnMappingMode::None)]
2960 fn test_logical_to_physical_transform(#[case] mode: ColumnMappingMode) -> DeltaResult<()> {
2961 validate_logical_to_physical_transform(mode)
2962 }
2963
2964 enum TestFileStats {
2970 None,
2972 Present,
2974 AllNull,
2976 }
2977
2978 fn create_test_add_files(paths: Vec<&str>, stats: Vec<TestFileStats>) -> Box<dyn EngineData> {
2980 let value_fields = vec![StructField::nullable("value", DataType::LONG)];
2981 let value_struct_type = DataType::struct_type_unchecked(value_fields.clone());
2982 let stats_type = DataType::struct_type_unchecked(vec![
2983 StructField::nullable(NUM_RECORDS, DataType::LONG),
2984 StructField::nullable(NULL_COUNT, value_struct_type.clone()),
2985 StructField::nullable(MIN_VALUES, value_struct_type.clone()),
2986 StructField::nullable(MAX_VALUES, value_struct_type.clone()),
2987 ]);
2988 let stats_fields = vec![
2989 StructField::nullable(NUM_RECORDS, DataType::LONG),
2990 StructField::nullable(NULL_COUNT, value_struct_type.clone()),
2991 StructField::nullable(MIN_VALUES, value_struct_type.clone()),
2992 StructField::nullable(MAX_VALUES, value_struct_type),
2993 ];
2994 let schema = Arc::new(StructType::new_unchecked(vec![
2995 StructField::not_null("path", DataType::STRING),
2996 StructField::not_null(
2997 "partitionValues",
2998 MapType::new(DataType::STRING, DataType::STRING, true),
2999 ),
3000 StructField::not_null("size", DataType::LONG),
3001 StructField::not_null("modificationTime", DataType::LONG),
3002 StructField::nullable("stats", stats_type.clone()),
3003 ]));
3004
3005 let empty_map = Scalar::Map(
3006 MapData::try_new(
3007 MapType::new(DataType::STRING, DataType::STRING, true),
3008 Vec::<(&str, &str)>::new(),
3009 )
3010 .unwrap(),
3011 );
3012
3013 let rows: Vec<Vec<Scalar>> = paths
3014 .iter()
3015 .zip(stats.iter())
3016 .map(|(path, stat)| {
3017 let stats_scalar = match stat {
3018 TestFileStats::None => Scalar::Null(stats_type.clone()),
3019 TestFileStats::Present | TestFileStats::AllNull => {
3020 let value_struct = |v: Option<i64>| {
3021 let scalar = v.map_or(Scalar::Null(DataType::LONG), |n| n.into());
3022 Scalar::Struct(
3023 StructData::try_new(value_fields.clone(), vec![scalar]).unwrap(),
3024 )
3025 };
3026 let (null_count, min, max) = match stat {
3027 TestFileStats::Present => (
3028 value_struct(Some(0)),
3029 value_struct(Some(1)),
3030 value_struct(Some(100)),
3031 ),
3032 _ => (
3033 value_struct(Some(100)),
3034 value_struct(None),
3035 value_struct(None),
3036 ),
3037 };
3038 Scalar::Struct(
3039 StructData::try_new(
3040 stats_fields.clone(),
3041 vec![100i64.into(), null_count, min, max],
3042 )
3043 .unwrap(),
3044 )
3045 }
3046 };
3047 vec![
3048 (*path).into(),
3049 empty_map.clone(),
3050 1024i64.into(),
3051 1000000i64.into(),
3052 stats_scalar,
3053 ]
3054 })
3055 .collect();
3056 let row_refs: Vec<&[Scalar]> = rows.iter().map(|r| r.as_slice()).collect();
3057 ArrowEvaluationHandler
3058 .create_many(schema, &row_refs)
3059 .unwrap()
3060 }
3061
3062 #[test]
3063 fn test_stats_validation_allows_all_null_clustering_column() {
3064 let (engine, snapshot) = setup_non_dv_table();
3065 let txn = snapshot
3066 .transaction(Box::new(FileSystemCommitter::new()), &engine)
3067 .unwrap()
3068 .with_operation("WRITE".to_string())
3069 .with_clustering_columns_for_test(vec![ColumnName::new(["value"])]);
3070
3071 let add_files = create_test_add_files(vec!["file1.parquet"], vec![TestFileStats::AllNull]);
3072
3073 let result = txn.validate_add_files_stats(&[add_files]);
3074
3075 assert!(
3076 result.is_ok(),
3077 "Stats validation should pass for all-null clustering columns, got: {result:?}",
3078 );
3079 }
3080
3081 #[test]
3082 fn test_stats_validation_when_clustering_cols_missing_stats() {
3083 let (engine, snapshot) = setup_non_dv_table();
3084 let txn = snapshot
3085 .transaction(Box::new(FileSystemCommitter::new()), &engine)
3086 .unwrap()
3087 .with_operation("WRITE".to_string())
3088 .with_clustering_columns_for_test(vec![ColumnName::new(["value"])]);
3090
3091 let add_files = create_test_add_files(vec!["file1.parquet"], vec![TestFileStats::None]);
3093
3094 let result = txn.validate_add_files_stats(&[add_files]);
3096
3097 assert!(
3098 result.is_err(),
3099 "Expected validation to fail when stats are missing for clustering columns"
3100 );
3101
3102 let err_msg = result.unwrap_err().to_string();
3103 assert!(
3104 err_msg.contains("Stats validation error") || err_msg.contains("no stats"),
3105 "Expected stats validation error, got: {err_msg}"
3106 );
3107 }
3108
3109 #[test]
3110 fn test_stats_validation_when_clustering_stats_present() {
3111 let (engine, snapshot) = setup_non_dv_table();
3112 let txn = snapshot
3113 .transaction(Box::new(FileSystemCommitter::new()), &engine)
3114 .unwrap()
3115 .with_operation("WRITE".to_string())
3116 .with_clustering_columns_for_test(vec![ColumnName::new(["value"])]);
3118
3119 let add_files = create_test_add_files(vec!["file1.parquet"], vec![TestFileStats::Present]);
3121
3122 let result = txn.validate_add_files_stats(&[add_files]);
3124
3125 assert!(
3126 result.is_ok(),
3127 "Stats validation should pass when stats are present, got: {result:?}"
3128 );
3129 }
3130
3131 #[test]
3132 fn test_stats_validation_skipped_without_clustering() {
3133 let (engine, snapshot) = setup_non_dv_table();
3134 let txn = snapshot
3135 .transaction(Box::new(FileSystemCommitter::new()), &engine)
3136 .unwrap()
3137 .with_operation("WRITE".to_string());
3138 let add_files = create_test_add_files(vec!["file1.parquet"], vec![TestFileStats::None]);
3142
3143 let result = txn.validate_add_files_stats(&[add_files]);
3145
3146 assert!(
3147 result.is_ok(),
3148 "Stats validation should be skipped without clustering, got: {result:?}"
3149 );
3150 }
3151
3152 #[test]
3153 fn disallow_catalog_committer_for_non_catalog_managed_table() {
3154 let storage = Arc::new(InMemory::new());
3155 let table_root = url::Url::parse("memory:///").unwrap();
3156 let engine = crate::engine::sync::SyncEngine::new_with_store(storage.clone());
3157
3158 let actions = [
3160 r#"{"commitInfo":{"timestamp":12345678900,"inCommitTimestamp":12345678900}}"#,
3161 r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":[],"writerFeatures":["inCommitTimestamp"]}}"#,
3162 r#"{"metaData":{"id":"test-id","format":{"provider":"parquet","options":{}},"schemaString":"{\"type\":\"struct\",\"fields\":[]}","partitionColumns":[],"configuration":{"delta.enableInCommitTimestamps":"true"},"createdTime":1234567890}}"#,
3163 ].join("\n");
3164
3165 let commit_path = Path::from("_delta_log/00000000000000000000.json");
3166 let rt = tokio::runtime::Runtime::new().unwrap();
3167 rt.block_on(storage.put(&commit_path, actions.into()))
3168 .unwrap();
3169
3170 let snapshot = Snapshot::builder_for(table_root).build(&engine).unwrap();
3171
3172 let committer = Box::new(MockCatalogCommitter);
3174 let err = snapshot
3175 .transaction(committer, &engine)
3176 .unwrap()
3177 .commit(&engine)
3178 .unwrap_err();
3179 assert!(matches!(
3180 err,
3181 crate::Error::Generic(e) if e.contains("This table is path-based and cannot be committed to with a catalog committer")
3182 ));
3183 }
3184
3185 #[test]
3186 fn disallow_catalog_committer_for_non_catalog_managed_create_table() {
3187 let storage = Arc::new(InMemory::new());
3188 let engine = crate::engine::sync::SyncEngine::new_with_store(storage);
3189
3190 let schema = Arc::new(crate::schema::StructType::new_unchecked(vec![
3192 crate::schema::StructField::new("id", crate::schema::DataType::INTEGER, true),
3193 ]));
3194 let committer = Box::new(MockCatalogCommitter);
3195 let err = create_table("memory:///", schema, "test-engine")
3196 .build(&engine, committer)
3197 .unwrap()
3198 .commit(&engine)
3199 .unwrap_err();
3200 assert!(matches!(
3201 err,
3202 crate::Error::Generic(e) if e.contains("This table is path-based and cannot be committed to with a catalog committer")
3203 ));
3204 }
3205
3206 struct CapturingCommitter {
3207 captured: Arc<Mutex<Option<i64>>>,
3208 }
3209
3210 impl CapturingCommitter {
3211 fn new() -> (Self, Arc<Mutex<Option<i64>>>) {
3212 let captured = Arc::new(Mutex::new(None));
3213 (
3214 Self {
3215 captured: captured.clone(),
3216 },
3217 captured,
3218 )
3219 }
3220 }
3221
3222 impl Committer for CapturingCommitter {
3223 fn commit(
3224 &self,
3225 _engine: &dyn Engine,
3226 _actions: DeltaResultIterator<'_, FilteredEngineData>,
3227 commit_metadata: CommitMetadata,
3228 ) -> DeltaResult<CommitResponse> {
3229 *self.captured.lock().unwrap() = Some(commit_metadata.in_commit_timestamp());
3230 Ok(CommitResponse::Conflict {
3231 version: commit_metadata.version(),
3232 })
3233 }
3234 fn is_catalog_committer(&self) -> bool {
3235 false
3236 }
3237 fn publish(
3238 &self,
3239 _engine: &dyn Engine,
3240 _publish_metadata: PublishMetadata,
3241 ) -> DeltaResult<()> {
3242 Ok(())
3243 }
3244 }
3245
3246 #[test]
3247 fn test_commit_metadata_receives_ict_not_wall_time() -> DeltaResult<()> {
3248 let tempdir = tempfile::tempdir().unwrap();
3253 let log_dir = tempdir.path().join("_delta_log");
3254 std::fs::create_dir_all(&log_dir).unwrap();
3255
3256 let future_ict: i64 = 9_999_999_999_999; let commit_info = serde_json::json!({
3258 "commitInfo": {
3259 "timestamp": 1000,
3260 "operation": "WRITE",
3261 "inCommitTimestamp": future_ict
3262 }
3263 });
3264 let protocol = serde_json::json!({
3265 "protocol": {
3266 "minReaderVersion": 3,
3267 "minWriterVersion": 7,
3268 "readerFeatures": [],
3269 "writerFeatures": ["inCommitTimestamp"]
3270 }
3271 });
3272 let schema_json = serde_json::json!({
3273 "type": "struct",
3274 "fields": [{
3275 "name": "id",
3276 "type": "integer",
3277 "nullable": true,
3278 "metadata": {}
3279 }]
3280 });
3281 let metadata = serde_json::json!({
3282 "metaData": {
3283 "id": "test-id",
3284 "format": {"provider": "parquet", "options": {}},
3285 "schemaString": schema_json.to_string(),
3286 "partitionColumns": [],
3287 "configuration": {
3288 "delta.enableInCommitTimestamps": "true"
3289 }
3290 }
3291 });
3292 let commit0 = format!("{commit_info}\n{protocol}\n{metadata}\n");
3293 std::fs::write(log_dir.join("00000000000000000000.json"), commit0).unwrap();
3294
3295 let table_url = Url::from_directory_path(tempdir.path()).unwrap();
3296 let engine = SyncEngine::new();
3297 let snapshot = Snapshot::builder_for(table_url).build(&engine)?;
3298
3299 let prev_ict = snapshot.get_in_commit_timestamp(&engine)?;
3300 assert_eq!(prev_ict, Some(future_ict));
3301
3302 let (committer, captured_ts) = CapturingCommitter::new();
3303 let mut txn = snapshot.transaction(Box::new(committer), &engine)?;
3304 add_dummy_file(&mut txn);
3305
3306 let result = txn.commit(&engine)?;
3307 assert!(
3308 matches!(result, CommitResult::ConflictedTransaction(_)),
3309 "Expected ConflictedTransaction from capturing committer"
3310 );
3311
3312 let captured = captured_ts
3314 .lock()
3315 .unwrap()
3316 .expect("should have captured a timestamp");
3317 assert_eq!(
3318 captured,
3319 future_ict + 1,
3320 "CommitMetadata.in_commit_timestamp should be the computed ICT (prev_ict + 1), \
3321 not the wall-clock time"
3322 );
3323 Ok(())
3324 }
3325
3326 fn commit_failure_event(reporter: &CapturingReporter) -> Option<TransactionCommitFailure> {
3329 reporter.events().into_iter().find_map(|event| match event {
3330 MetricEvent::TransactionCommitFailure(f) => Some(f),
3331 _ => None,
3332 })
3333 }
3334
3335 #[test]
3336 fn test_commit_io_error_emits_retryable_io_failure_metric() -> DeltaResult<()> {
3337 let (engine, snapshot, _tempdir) = load_test_table("table-without-dv-small")?;
3338 let reporter = Arc::new(CapturingReporter::default());
3339 let _guard = install_thread_local_metrics_reporter(reporter.clone());
3340 let mut txn = snapshot.transaction(Box::new(IoErrorCommitter), engine.as_ref())?;
3341 add_dummy_file(&mut txn);
3342 let result = txn.commit(engine.as_ref())?;
3343 assert!(matches!(result, CommitResult::RetryableTransaction(_)));
3344 let failure = commit_failure_event(&reporter).expect("commit failure event");
3345 assert_eq!(failure.reason, CommitFailureReason::RetryableIo);
3346 assert_eq!(failure.table_type, TableType::PathBased);
3347 Ok(())
3348 }
3349
3350 #[test]
3351 fn test_commit_terminal_error_emits_error_failure_metric() -> DeltaResult<()> {
3352 let (engine, snapshot, _tempdir) = load_test_table("table-without-dv-small")?;
3353 let reporter = Arc::new(CapturingReporter::default());
3354 let _guard = install_thread_local_metrics_reporter(reporter.clone());
3355 let mut txn = snapshot.transaction(Box::new(GenericErrorCommitter), engine.as_ref())?;
3356 add_dummy_file(&mut txn);
3357 assert!(txn.commit(engine.as_ref()).is_err());
3358 let failure = commit_failure_event(&reporter).expect("commit failure event");
3359 assert_eq!(failure.reason, CommitFailureReason::Error);
3360 assert_eq!(failure.table_type, TableType::PathBased);
3361 Ok(())
3362 }
3363}