Skip to main content

buoyant_kernel/transaction/
mod.rs

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
86/// Type alias for an iterator of [`EngineData`] results.
87pub(crate) type EngineDataResultIterator<'a> =
88    Box<dyn Iterator<Item = DeltaResult<Box<dyn EngineData>>> + Send + 'a>;
89
90/// The static instance referenced by [`add_files_schema`] that doesn't contain the dataChange
91/// column.
92pub(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
99/// Returns a reference to the mandatory fields in an add action.
100///
101/// Note this does not include "dataChange" which is a required field but
102/// but should be set on the transactoin level. Getting the full schema
103/// can be done with [`Transaction::add_files_schema`].
104pub(crate) fn mandatory_add_file_schema() -> &'static SchemaRef {
105    &MANDATORY_ADD_FILE_SCHEMA
106}
107
108/// The base schema for add file metadata, referenced by [`Transaction::add_files_schema`].
109///
110/// The `stats` field represents the minimum structure. The actual stats written by
111/// `DefaultEngine::write_parquet` include additional fields computed from the data:
112/// - `nullCount`: nested struct mirroring the data schema (all fields LONG)
113/// - `minValues`: nested struct with min/max eligible column types
114/// - `maxValues`: nested struct with min/max eligible column types
115///
116/// The nested structures within nullCount/minValues/maxValues depend on the table's data schema
117/// and which columns have statistics enabled. Use [`Transaction::stats_schema`] to get the
118/// expected stats schema for a specific table.
119pub(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        // nullCount, minValues, maxValues are dynamic based on data schema. Empty struct
124        // placeholders indicate these fields exist but their inner structure depends on the
125        // table schema and stats column configuration.
126        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
136/// Extend a schema with row tracking columns and return a new SchemaRef.
137///
138/// Note that this method is only useful to extend an Add action schema.
139fn 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/// Marker type for transactions on existing tables.
150///
151/// This is the default state for [`Transaction`] and provides the full set of operations
152/// including file removal, deletion vector updates, and blind append semantics.
153#[derive(Debug)]
154pub struct ExistingTable;
155
156/// Marker type for create-table transactions.
157///
158/// Transactions in this state have a restricted API surface — operations that are semantically
159/// invalid for table creation (e.g. file removal, domain metadata removal) are not available.
160#[derive(Debug)]
161pub struct CreateTable;
162
163/// Marker type for alter-table (schema evolution) transactions.
164///
165/// Transactions in this state perform metadata-only commits. Data file operations are not
166/// available at compile time because `AlterTable` does not implement [`SupportsDataFiles`].
167#[derive(Debug)]
168pub struct AlterTable;
169
170/// Marker trait for transaction states that support data file operations.
171///
172/// Only transaction types that implement this trait can access methods for adding, removing, or
173/// updating data files. This prevents compile-time misuse by states like `AlterTable` that
174/// only perform metadata-only commits.
175pub trait SupportsDataFiles {}
176impl SupportsDataFiles for ExistingTable {}
177impl SupportsDataFiles for CreateTable {}
178
179/// A transaction represents an in-progress write to a table. After creating a transaction, changes
180/// to the table may be staged via the transaction methods before calling `commit` to commit the
181/// changes to the table.
182///
183/// The type parameter `S` controls which operations are available:
184/// - [`ExistingTable`] (default): Full API for modifying existing tables.
185/// - [`CreateTable`]: Restricted API for table creation (see
186///   [`CreateTableTransaction`](create_table::CreateTableTransaction)).
187///
188/// # Examples
189///
190/// ```rust,ignore
191/// // create a transaction
192/// let mut txn = table.new_transaction(&engine)?;
193/// // stage table changes (right now only commit info)
194/// txn.commit_info(Box::new(ArrowEngineData::new(engine_commit_info)));
195/// // commit! (consume the transaction)
196/// txn.commit(&engine)?;
197/// ```
198pub struct Transaction<S = ExistingTable> {
199    span: tracing::Span,
200    // Correlates all metric events emitted by this transaction.
201    operation_id: MetricId,
202    // Opaque, caller-supplied id recorded on this transaction's commit metric events alongside
203    // `operation_id`. Set via `with_correlation_id`; not interpreted by kernel.
204    correlation_id: Option<Arc<str>>,
205    // The snapshot this transaction is based on. None for CREATE TABLE (no pre-existing table).
206    // Use `read_snapshot()` to access; it returns an error if None.
207    read_snapshot_opt: Option<SnapshotRef>,
208    // The table configuration that this commit will produce. For writes that don't change the
209    // config, this is cloned from the read snapshot; when the config changes (e.g. schema
210    // evolution), it is constructed separately with the new schema/protocol.
211    effective_table_config: TableConfiguration,
212    // Whether to emit a Protocol action. True for CREATE TABLE and ALTER TABLE, false otherwise.
213    should_emit_protocol: bool,
214    // Whether to emit a Metadata action. True for CREATE TABLE and ALTER TABLE, false otherwise.
215    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    // NB: hashmap would require either duplicating the appid or splitting SetTransaction
223    // key/payload. HashSet requires Borrow<&str> with matching Eq, Ord, and Hash. Plus,
224    // HashSet::insert drops the to-be-inserted value without returning the existing one, which
225    // would make error messaging unnecessarily difficult. Thus, we keep Vec here and deduplicate
226    // in the commit method.
227    set_transactions: Vec<SetTransaction>,
228    // commit-wide timestamp (in milliseconds since epoch) - used in ICT, `txn` action, etc. to
229    // keep all timestamps within the same commit consistent.
230    commit_timestamp: i64,
231    // User-provided domain metadata additions (via with_domain_metadata API).
232    user_domain_metadata_additions: Vec<DomainMetadata>,
233    // System-generated domain metadata (from transforms, e.g., clustering).
234    // TODO(#1779): Currently only populated during CREATE TABLE. For inserts, row tracking
235    // domain metadata is handled separately via `row_tracking_high_watermark` parameter in
236    // `generate_domain_metadata_actions`. Consider unifying system domain handling.
237    system_domain_metadata_additions: Vec<DomainMetadata>,
238    // Domain names to remove in this transaction. The configuration values are fetched during
239    // commit from the log to preserve the pre-image in tombstones.
240    user_domain_removals: Vec<String>,
241    // Whether this transaction contains any logical data changes.
242    data_change: bool,
243    // Whether this transaction should be marked as a blind append.
244    is_blind_append: bool,
245    // Files matched by update_deletion_vectors() with new DV descriptors appended. These are used
246    // to generate remove/add action pairs during commit, ensuring file statistics are preserved.
247    dv_matched_files: Vec<FilteredEngineData>,
248    // Clustering columns from domain metadata. Only populated if the ClusteredTable feature is
249    // enabled. Used for determining which columns require statistics collection. Expected to be
250    // physical column names.
251    physical_clustering_columns: Option<Vec<ColumnName>>,
252    // PhantomData marker for transaction state (ExistingTable or CreateTable).
253    // Zero-sized; only affects the type system.
254    _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
271/// Builds the projection for converting add file metadata into commit-ready Add actions.
272fn 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
292/// Transforms add file metadata into commit-ready add actions by converting stats to JSON and
293/// setting the `dataChange` field.
294fn 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
317// =============================================================================
318// Shared methods available on ALL transaction types
319// =============================================================================
320impl<S> Transaction<S> {
321    /// Consume the transaction and commit it to the table. The result is a result of
322    /// [CommitResult] with the following semantics:
323    /// - Ok(CommitResult) for either success or a recoverable error (includes the failed
324    ///   transaction in case of a conflict so the user can retry, etc.)
325    /// - Err(Error) indicates a non-retryable error (e.g. logic/validation error).
326    #[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        // Fields-only event: these feed the `txn.commit` metric via the layer's `on_event`
352        // channel. `num_dv_updates` has no other source (it is not a declared span field and
353        // gets no `span.record` below), so this event must keep its structured fields.
354        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        // Some table features don't yet support removeFiles. Reject here.
361        if !self.remove_files_metadata.is_empty() {
362            self.effective_table_config
363                .validate_feature_support_for_remove()?;
364        }
365
366        // Step 1: Check for duplicate app_ids and generate set transactions (`txn`)
367        // Note: The commit info must always be the first action in the commit but we generate it in
368        // step 2 to fail early on duplicate transaction appIds
369        // TODO(zach): we currently do this in two passes - can we do it in one and still keep refs
370        // in the HashSet?
371        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        // Validate that the schema supports data writes when files are being added.
387        // Void-in-array/map, all-void structs, and all-void tables cannot produce valid Parquet.
388        // Reads and metadata-only commits are always allowed.
389        if !self.add_files_metadata.is_empty() {
390            validate_schema_for_write(&self.effective_table_config.logical_schema())?;
391        }
392
393        // CDF check only applies to existing tables (not create table)
394        // If there are add and remove files with data change in the same transaction, we block it.
395        // This is because kernel does not yet have a way to discern DML operations. For DML
396        // operations that perform updates on rows, ChangeDataFeed requires that a `cdc` file be
397        // written to the delta log.
398        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        // Validate clustering column stats if ClusteredTable feature is enabled
419        self.validate_add_files_stats(&self.add_files_metadata)?;
420
421        // Step 1: Generate SetTransaction actions
422        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        // Step 2: Construct commit info with ICT if enabled
429        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        // Step 3: Generate Protocol and Metadata actions based on emit flags
440        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        // Step 4: Generate add actions and get data for domain metadata actions (e.g. row tracking
458        // high watermark)
459        let commit_version = self.get_commit_version();
460        let (add_actions, row_tracking_domain_metadata) =
461            self.generate_adds(engine, commit_version)?;
462
463        // Step 4b: Generate all domain metadata actions (user and system domains)
464        let (domain_metadata_actions, dm_changes) =
465            self.generate_domain_metadata_actions(engine, row_tracking_domain_metadata)?;
466
467        // Step 5: Generate DV update actions (remove/add pairs) if any DV updates are present
468        let dv_update_actions = self.generate_dv_update_actions(engine)?;
469
470        // Step 6: Generate remove actions (collect to avoid borrowing self)
471        let remove_actions =
472            self.generate_remove_actions(engine, self.remove_files_metadata.iter(), &[])?;
473
474        // Build the action chain
475        // For create-table: CommitInfo -> Protocol -> Metadata -> adds -> txns -> domain_metadata
476        // -> removes For existing table: CommitInfo -> adds -> txns -> domain_metadata ->
477        // removes
478        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        // Step 7: Commit via the committer
491        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                // TODO(#2717): the commit already succeeded atomically; the post-commit `?`
507                //              below must not fail the txn (and must not mislabel the metric).
508                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                // Flips the metric event from success -> failure.
532                tracing::Span::current()
533                    .record("failure_reason", CommitFailureReason::Conflict.as_ref());
534                Ok(CommitResult::ConflictedTransaction(
535                    self.into_conflicted(version),
536                ))
537            }
538            // TODO: we may want to be more or less selective about what is retryable (this is tied
539            // to the idea of "what kind of Errors should write_json_file return?")
540            Err(e @ Error::IOError(_)) => {
541                // Flips the metric event from success -> failure.
542                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    /// Set the data change flag.
574    ///
575    /// True indicates this commit is a "data changing" commit. False indicates table data was
576    /// reorganized but not materially modified.
577    ///
578    /// Data change might be set to false in the following scenarios:
579    /// 1. Operations that only change metadata (e.g. backfilling statistics)
580    /// 2. Operations that make no logical changes to the contents of the table (i.e. rows are only
581    ///    moved from old files to new ones.  OPTIMIZE commands is one example of this type of
582    ///    optimizaton).
583    pub fn with_data_change(mut self, data_change: bool) -> Self {
584        self.data_change = data_change;
585        self
586    }
587
588    /// Same as [`Transaction::with_data_change`] but set the value directly instead of
589    /// using a fluent API.
590    #[internal_api]
591    #[allow(dead_code)] // used in FFI
592    pub(crate) fn set_data_change(&mut self, data_change: bool) {
593        self.data_change = data_change;
594    }
595
596    /// Set the engine info field of this transaction's commit info action. This field is optional.
597    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    /// Attach an opaque, caller-supplied correlation id for joining this transaction's commit
603    /// metric events to the caller's own request or operation id. An empty id is treated as unset.
604    /// When unset, behavior is unchanged.
605    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    /// Set the content of the commitInfo action for this transaction. Note that kernel will
611    /// _always_ write a commitInfo, this function simply allows engines to add their own data
612    /// into that action if they wish. Note that the following fields in `engine_commit_info`
613    /// will be overridden by kernel if they are set (meaning you should not set them):
614    /// - timestamp
615    /// - inCommitTimestamp
616    /// - operation
617    /// - operationParameters
618    /// - kernelVersion
619    /// - isBlindAppend
620    /// - engineInfo
621    /// - txnId
622    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    /// Include a SetTransaction (app_id and version) action for this transaction (with an optional
632    /// `last_updated` timestamp).
633    /// Note that each app_id can only appear once per transaction. That is, multiple app_ids with
634    /// different versions are disallowed in a single transaction. If a duplicate app_id is
635    /// included, the `commit` will fail (that is, we don't eagerly check app_id validity here).
636    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    /// Set domain metadata to be written to the Delta log.
643    /// Note that each domain can only appear once per transaction. That is, multiple configurations
644    /// of the same domain are disallowed in a single transaction, as well as setting and removing
645    /// the same domain in a single transaction. If a duplicate domain is included, the commit will
646    /// fail (that is, we don't eagerly check domain validity here).
647    /// Setting metadata for multiple distinct domains is allowed.
648    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    /// Determines the commit type based on whether this is a create-table operation and whether
655    /// the table is catalog-managed.
656    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        // TODO: Handle UpgradeToCatalogManaged and DowngradeToPathBased when ALTER TABLE
663        // SET TBLPROPERTIES is supported.
664        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    /// Validates that the committer type matches the commit type. A catalog committer must be
673    /// used for catalog-managed operations, and a non-catalog committer for path-based operations.
674    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    /// Builds the [`CommitMetadata`] for this transaction. Determines the commit type,
694    /// validates the committer, and assembles the protocol/metadata state.
695    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        // For create-table: previous P&M is None (no prior table), new P&M is set.
708        // For existing table with metadata change: previous P&M is from snapshot, new P&M
709        // is from effective config.
710        // For existing table without metadata change: previous P&M is from snapshot, new is None.
711        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    /// Validate that the transaction is eligible to be marked as a blind append.
740    ///
741    /// Note: Domain metadata additions/removals are allowed; blind append only constrains
742    /// data-file operations and read predicates. Conflict resolution determines whether
743    /// metadata changes are problematic.
744    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    /// Reject data file writes (add/remove/DV) against an empty-schema table.
775    /// CREATE TABLE and metadata-only commits are exempt.
776    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    /// Reject `WriteContext` handouts on empty-schema tables, so engines fail
787    /// before staging any parquet. CREATE TABLE is exempt.
788    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    /// Returns true if this is a create-table transaction.
803    /// A create-table transaction has no read snapshot (no pre-existing table).
804    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    /// True iff this transaction stages any data-file action (add, remove, or DV update).
813    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    // Returns the read snapshot. Returns an error if this is a create-table transaction.
820    // To get the `Option<SnapshotRef>` directly, use the `read_snapshot_opt` field.
821    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    /// Computes the in-commit timestamp for this transaction if ICT is enabled.
828    /// Returns `None` if ICT is not enabled on the table. A feature being in the protocol
829    /// (`is_feature_supported`) is not sufficient -- the `delta.enableInCommitTimestamps`
830    /// property must also be `true` (`is_feature_enabled`).
831    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            // For CREATE TABLE there are no prior commits -- use the wall-clock time directly.
842            return Ok(Some(self.commit_timestamp));
843        }
844
845        // Existing table: enforce monotonicity per the Delta protocol. The timestamp
846        // must be the larger of:
847        // - The time at which the writer attempted the commit
848        // - One millisecond later than the previous commit's inCommitTimestamp
849        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    /// Returns the commit version for this transaction.
856    /// For existing table transactions, this is snapshot.version() + 1.
857    /// For create-table transactions, this is 0.
858    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    /// The schema that the [`Engine`]'s [`ParquetHandler`] is expected to use when reporting
866    /// information about a Parquet write operation back to Kernel.
867    ///
868    /// Concretely, it is the expected schema for [`EngineData`] passed to [`add_files`], as it is
869    /// the base for constructing an add_file. Each row represents metadata about a
870    /// file to be added to the table. Kernel takes this information and extends it to the full
871    /// add_file action schema, adding internal fields (e.g., baseRowID) as necessary.
872    ///
873    /// The `stats` field contains file-level statistics. The schema returned here shows the base
874    /// structure; the actual stats written by `DefaultEngine::write_parquet` include dynamically
875    /// computed fields (numRecords, nullCount, minValues, maxValues, tightBounds) based on the
876    /// data schema and table configuration. See [`stats_schema`] for the table-specific expected
877    /// stats schema.
878    ///
879    /// Note: While currently static, in the future the schema might change depending on
880    /// options set on the transaction or features enabled on the table.
881    ///
882    /// [`add_files`]: crate::transaction::Transaction::add_files
883    /// [`ParquetHandler`]: crate::ParquetHandler
884    /// [`stats_schema`]: Transaction::stats_schema
885    pub fn add_files_schema(&self) -> &'static SchemaRef {
886        &BASE_ADD_FILES_SCHEMA
887    }
888}
889
890// =============================================================================
891// Data file methods -- only available on transaction types that support data files
892// =============================================================================
893impl<S: SupportsDataFiles> Transaction<S> {
894    /// Returns the expected schema for file statistics.
895    ///
896    /// The schema structure is derived from table configuration:
897    /// - `delta.dataSkippingStatsColumns`: Explicit column list (if set)
898    /// - `delta.dataSkippingNumIndexedCols`: Column count limit (default 32)
899    /// - Partition columns: Always excluded
900    ///
901    /// The returned schema has the following structure:
902    /// ```ignore
903    /// {
904    ///   numRecords: long,
905    ///   nullCount: { ... },   // Nested struct mirroring data schema, all fields LONG
906    ///   minValues: { ... },   // Nested struct, only min/max eligible types
907    ///   maxValues: { ... },   // Nested struct, only min/max eligible types
908    ///   tightBounds: boolean,
909    /// }
910    /// ```
911    ///
912    /// Engines should collect statistics matching this schema structure when writing files.
913    ///
914    /// Per the Delta protocol, required columns (e.g. clustering columns) are always included
915    /// in statistics, regardless of `dataSkippingStatsColumns` or `dataSkippingNumIndexedCols`
916    /// settings.
917    #[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    /// Returns the list of column names that should have statistics collected.
926    ///
927    /// This returns leaf column paths as [`ColumnName`] objects. Each `ColumnName`
928    /// stores path components separately (e.g., `ColumnName::new(["nested", "field"])`).
929    /// See [`ColumnName`'s `Display` implementation][ColumnName#impl-Display-for-ColumnName]
930    /// for details on string formatting and escaping.
931    ///
932    /// Engines can use this to determine which columns need stats during writes.
933    ///
934    /// Per the Delta protocol, clustering columns are always included in statistics,
935    /// regardless of `dataSkippingStatsColumns` or `dataSkippingNumIndexedCols` settings.
936    #[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    // Generate the logical-to-physical expression which must be evaluated on every data chunk
943    // before writing.
944    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            // Insert each partition column after the nearest preceding surviving field
961            // (non-partition and non-void), in the order they appear in the logical schema.
962            // This keeps the post-transform data aligned with the physical schema.
963            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    /// Returns the logical partition column names for this table.
988    pub fn logical_partition_columns(&self) -> &[String] {
989        self.effective_table_config.partition_columns()
990    }
991
992    /// Returns the column default for every top-level column in this table's logical schema that
993    /// declares one, keyed by logical column name.
994    ///
995    /// Connectors use this to discover which columns have defaults, then call
996    /// [`ColumnDefault::to_scalar`] on each (or fall back to [`ColumnDefault::raw_sql`] when the
997    /// kernel cannot parse the default) to materialize the column before writing.
998    ///
999    /// Keys are `String` rather than [`ColumnName`] because the kernel currently surfaces defaults
1000    /// only for top-level columns, consistent with partition columns. This is a kernel limitation,
1001    /// not a protocol one.
1002    ///
1003    /// # Errors
1004    ///
1005    /// - A column declares a `CURRENT_DEFAULT` but the table does not enable the
1006    ///   `allowColumnDefaults` writer feature. The protocol only honors defaults "when enabled", so
1007    ///   such metadata is stray and is rejected rather than returned.
1008    /// - Propagates any error from [`StructField::column_default`] -- a malformed `CURRENT_DEFAULT`
1009    ///   (non-string metadata, or a non-NULL default on a non-primitive type).
1010    #[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    /// Validates that the table's logical schema supports data writes.
1033    ///
1034    /// Called at the top of [`partitioned_write_context`](Self::partitioned_write_context) and
1035    /// [`unpartitioned_write_context`](Self::unpartitioned_write_context), before any Parquet is
1036    /// written, so connectors fail fast when the schema contains void placements that cannot
1037    /// produce valid files (void inside Array/Map, all-void structs, all-void tables).
1038    /// The commit-time check in [`commit`](Self::commit) remains as defense-in-depth for callers
1039    /// that reach [`add_files`](Self::add_files) without going through a write context.
1040    fn validate_for_data_write(&self) -> DeltaResult<()> {
1041        validate_schema_for_write(&self.effective_table_config.logical_schema())
1042    }
1043
1044    /// Builds the [`SharedWriteState`] for a write context.
1045    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    /// Creates a write context for writing data to a specific partition.
1061    ///
1062    /// Performs the following validations and transformations:
1063    ///
1064    /// - **Key completeness**: ensures all partition columns are present and no extra keys exist.
1065    ///   For example, if the table has partition columns `["year", "region"]` and you pass
1066    ///   `{"year": Scalar::Integer(2024)}`, this returns an error for missing "region".
1067    ///
1068    /// - **Case normalization**: matches keys case-insensitively against the schema and normalizes
1069    ///   to schema case. For example, passing `"YEAR"` for a column named `"year"` is accepted and
1070    ///   normalized.
1071    ///
1072    /// - **Type checking**: rejects non-primitive partition column types (struct, array, map) and
1073    ///   validates that each non-null `Scalar`'s type matches the partition column's schema type.
1074    ///   For example, passing `Scalar::String("2024")` for an `INTEGER` column returns an error.
1075    ///   Null-equivalent scalars (null scalars, empty strings, and empty binary) all of which
1076    ///   collapse to JSON null in `partitionValues`) skip the value type check, but they are only
1077    ///   legal when the partition column is nullable; passing any of these for a `nullable: false`
1078    ///   partition column returns an error.
1079    ///
1080    /// - **Value serialization**: serializes each `Scalar` to a protocol-compliant string per the
1081    ///   Delta protocol's "Partition Value Serialization" rules. `Scalar::Null(...)` becomes `None`
1082    ///   in `add.partitionValues` (JSON null). `Scalar::String("")` also becomes `None` (empty
1083    ///   string equals null for all types). `Scalar::Date(19723)` becomes `Some("2024-01-01")`.
1084    ///
1085    /// - **Key translation**: translates logical column names to physical names using the table's
1086    ///   column mapping mode. For example, under `ColumnMappingMode::Name`, logical `"year"` might
1087    ///   become physical `"col-abc-123"` in the `partitionValues` map.
1088    ///
1089    /// - **Partition column materialization**: the returned [`WriteContext`]'s
1090    ///   [`logical_to_physical`] expression injects partition columns when the table requires
1091    ///   materializing partition columns (e.g. `materializePartitionColumns` or `icebergCompatV3`).
1092    ///   The input data fed to that expression must not contain partition columns.
1093    ///
1094    /// The returned [`WriteContext`] also provides a [`write_dir`] that returns the correct
1095    /// target directory (Hive-style paths when column mapping is off, random prefix when on).
1096    ///
1097    /// Returns an error if the table is not partitioned (use
1098    /// [`unpartitioned_write_context`](Self::unpartitioned_write_context) instead).
1099    ///
1100    /// [`write_dir`]: WriteContext::write_dir
1101    /// [`logical_to_physical`]: WriteContext::logical_to_physical
1102    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        // Validate keys (completeness, case normalization) and value types, then return
1114        // the map re-keyed to schema case.
1115        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        // Serialize values and translate keys from logical to physical names.
1123        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    /// Creates a write context for writing data to an unpartitioned table.
1152    ///
1153    /// Returns an error if the table has partition columns (use
1154    /// [`partitioned_write_context`](Self::partitioned_write_context) instead).
1155    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    /// Add files to include in this transaction. This API generally enables the engine to
1172    /// add/append/insert data (files) to the table. Note that this API can be called multiple times
1173    /// to add multiple batches.
1174    ///
1175    /// The expected schema for `add_metadata` is given by [`Transaction::add_files_schema`].
1176    pub fn add_files(&mut self, add_metadata: Box<dyn EngineData>) {
1177        self.add_files_metadata.push(add_metadata);
1178    }
1179}
1180
1181// =============================================================================
1182// Internal methods available on ALL transaction types (used by commit path)
1183// =============================================================================
1184impl<S> Transaction<S> {
1185    /// Validate that add files carry the per-file statistics required by the table's protocol.
1186    ///
1187    /// Currently checks two protocol requirements:
1188    /// - `stats.numRecords` must be present when [`requires_stats_num_records`] returns true.
1189    /// - Per-file min/max/nullCount must be present for clustering columns when the
1190    ///   `ClusteredTable` feature is enabled.
1191    ///
1192    /// Other stat columns (e.g. the conventional "first 32 columns") are not validated here
1193    /// because they are not protocol-required.
1194    ///
1195    /// Only add files are validated(remove files do not carry statistics).
1196    ///
1197    /// [`requires_stats_num_records`]: crate::table_configuration::TableConfiguration::requires_stats_num_records
1198    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            // TODO: Likely it's better to merge this with the clustering column validation below,
1204            // benchmark it and see if it's faster. If so, refactor this to do both validations in
1205            // one pass.
1206            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    /// Generates add actions and row tracking domain metadata for a commit.
1234    #[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        // Note: this does not require delta.enableRowTracking=true. "supported" is sufficient
1244        // for writers to assign row IDs.
1245        let row_tracking_supported = self.effective_table_config.should_write_row_tracking();
1246
1247        if self.add_files_metadata.is_empty() {
1248            // No files to add. For an empty CREATE TABLE with row tracking, emit the initial
1249            // high water mark domain metadata (rowIdHighWaterMark = -1) so subsequent writes
1250            // have a valid starting point. For all other empty commits (metadata-only, etc.),
1251            // nothing row-tracking-related needs to be written.
1252            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    /// Generates add actions with row tracking columns and the row ID high water mark
1274    /// domain metadata.
1275    ///
1276    /// Visits all add file batches once to read `numRecords` per file, assigning a unique
1277    /// non-overlapping `baseRowId` range to each file and computing the final high water mark
1278    /// for the domain metadata action. The initial high water mark is read from the snapshot
1279    /// for existing tables, or defaults to -1 for create-table (no prior log to read from).
1280    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        // Create a row tracking visitor and visit all files to collect row tracking information
1295        let mut row_tracking_visitor =
1296            RowTrackingVisitor::new(row_id_high_water_mark, Some(self.add_files_metadata.len()));
1297
1298        // We visit all files with the row visitor before creating the add action iterator because
1299        // we need to know the final row ID high water mark to create the domain metadata action.
1300        for add_files_batch in &self.add_files_metadata {
1301            row_tracking_visitor.visit_rows_of(add_files_batch.deref())?;
1302        }
1303
1304        // Destructure the visitor to move base_row_id_batches into the add-files iterator
1305        // while also extracting the final high water mark for the domain metadata action.
1306        let RowTrackingVisitor {
1307            base_row_id_batches,
1308            row_id_high_water_mark,
1309        } = row_tracking_visitor;
1310
1311        // Create extended add files with row tracking columns
1312        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        // Generate add actions including row tracking metadata
1330        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        // Generate a row tracking domain metadata based on the final high water mark
1338        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                // Existing table path: use the read snapshot to compute post-commit state.
1356                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                // CREATE TABLE path: build a fresh Snapshot at version 0.
1368                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    /// Build a [`CrcDelta`] from the transaction's commit state and a precomputed
1397    /// [`FileStatsDelta`].
1398    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        // TODO: drop these conversions by migrating the upstream chain
1405        //       (`CommitMetadata.domain_metadata_changes`, `Transaction.set_transactions`)
1406        //       to `HashMap<String, _>`, lifting protocol-mandated uniqueness from runtime
1407        //       checks into the type system.
1408        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        // Although `remove.size` is optional per the Delta protocol, the kernel write path
1418        // enforces presence: `try_compute_for_txn` above errors with `MissingData` if any
1419        // add or remove row lacks `size` (see `FileStatsVisitor::visit` in
1420        // `kernel/src/crc/file_stats.rs`). So at this point every size is known to be
1421        // present, and only operation classification can flip `is_incremental_safe`.
1422        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    /// Generates Remove actions from scan file metadata.
1456    ///
1457    /// This internal method transforms scan row metadata into Remove actions for the Delta log.
1458    /// It's called during commit to process files staged via [`remove_files`] or files being
1459    /// updated with new deletion vectors via [`update_deletion_vectors`].
1460    ///
1461    /// # Parameters
1462    ///
1463    /// - `engine`: The engine used for expression evaluation
1464    /// - `remove_files_metadata`: Iterator over scan file metadata to transform into Remove actions
1465    /// - `columns_to_drop`: Column names to drop from the scan metadata before transformation. This
1466    ///   is used to remove temporary columns like the intermediate deletion vector column added
1467    ///   during DV updates.
1468    ///
1469    /// # Returns
1470    ///
1471    /// An iterator of FilteredEngineData containing Remove actions in the log schema format.
1472    ///
1473    /// [`remove_files`]: Transaction::remove_files
1474    /// [`update_deletion_vectors`]: Transaction::update_deletion_vectors
1475    #[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        // Create-table transactions should not have any remove actions.
1483        // Only error if there are actually files queued for removal.
1484        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        // Build two evaluators: one for the common case where scan files do not include a
1510        // stats_parsed column, and one for predicate-based scans that include stats_parsed.
1511        // The stats_parsed evaluator coalesces stats with ToJson(stats_parsed) to handle the
1512        // case where stats is null (e.g., on V2 checkpoints with writeStatsAsJson=false) and
1513        // then drops the stats_parsed column.
1514        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
1534/// Builds the struct patch for converting scan row metadata into a Remove action.
1535///
1536/// Handles two "parsed" columns that predicate-based scans add to scan metadata:
1537///
1538/// - `stats_parsed`: when `coalesce_stats_with_parsed` is true, the `stats` field is replaced with
1539///   `COALESCE(stats, TO_JSON(stats_parsed))` and `stats_parsed` is dropped. The coalesce handles
1540///   cases where `stats` is null (e.g., V2 checkpoints with `writeStatsAsJson=false`) by
1541///   reconstructing the JSON from the parsed representation.
1542/// - `partitionValues_parsed`: dropped if present. Unlike stats, no reconstruction is needed: the
1543///   Remove action's `partitionValues` is sourced from `fileConstantValues.partitionValues`, which
1544///   scans always populate from `add.partitionValues`.
1545fn 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        // deletionTimestamp
1553        .insert_after("path", lit(commit_timestamp))
1554        // dataChange
1555        .insert_after("path", lit(data_change))
1556        // extended_file_metadata
1557        .insert_after("path", lit(true))
1558        .insert_after("path", col!(FILE_CONSTANT_VALUES_NAME, "partitionValues"));
1559
1560    if coalesce_stats_with_parsed {
1561        // Replace stats with COALESCE(stats, TO_JSON(stats_parsed)) and drop stats_parsed.
1562        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        // Added to scan output when the predicate touches a partition column.
1584        .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/// Kernel exposes information about the state of the table that engines might want to use to
1594/// trigger actions like checkpointing or log compaction. This struct holds that information.
1595#[derive(Debug)]
1596pub struct PostCommitStats {
1597    /// The number of commits since this table has been checkpointed. Note that commit 0 is
1598    /// considered a checkpoint for the purposes of this computation.
1599    pub commits_since_checkpoint: u64,
1600    /// The number of commits since the log has been compacted on this table. Note that a
1601    /// checkpoint is considered a compaction for the purposes of this computation. Thus this
1602    /// is really the number of commits since a compaction OR a checkpoint.
1603    pub commits_since_log_compaction: u64,
1604}
1605
1606/// The result of attempting to commit this transaction. If the commit was
1607/// successful/conflicted/retryable, the result is Ok(CommitResult), otherwise, if a nonrecoverable
1608/// error occurred, the result is Err(Error).
1609///
1610/// The commit result can be one of the following:
1611/// - [CommittedTransaction]: the transaction was successfully committed. [PostCommitStats] and in
1612///   the future a post-commit snapshot can be obtained from the committed transaction.
1613/// - [ConflictedTransaction]: the transaction conflicted with an existing version. This transcation
1614///   must be rebased before retrying. (currently no rebase APIs exist, caller must create new txn)
1615/// - [RetryableTransaction]: an IO (retryable) error occurred during the commit. This transaction
1616///   can be retried without rebasing.
1617#[derive(Debug)]
1618#[must_use]
1619pub enum CommitResult<S = ExistingTable> {
1620    /// The transaction was successfully committed.
1621    CommittedTransaction(CommittedTransaction),
1622    /// This transaction conflicted with an existing version (see
1623    /// [ConflictedTransaction::conflict_version]). The transaction
1624    /// is returned so the caller can resolve the conflict (along with the version which
1625    /// conflicted).
1626    // TODO(zach): in order to make the returning of a transaction useful, we need to add APIs to
1627    // update the transaction to a new version etc.
1628    ConflictedTransaction(ConflictedTransaction<S>),
1629    /// An IO (retryable) error occurred during the commit.
1630    RetryableTransaction(RetryableTransaction<S>),
1631}
1632
1633impl<S> CommitResult<S> {
1634    /// Returns true if the commit was successful.
1635    pub fn is_committed(&self) -> bool {
1636        matches!(self, CommitResult::CommittedTransaction(_))
1637    }
1638}
1639
1640impl<S: std::fmt::Debug> CommitResult<S> {
1641    /// Unwraps the [`CommittedTransaction`], panicking if the commit was not successful.
1642    #[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    /// Unwraps the post-commit snapshot of the [`CommittedTransaction`], panicking if the
1652    /// commit was not successful or the post-commit snapshot is missing.
1653    /// TODO(#2494): Refactor existing tests to use this.
1654    #[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/// This is the result of a successfully committed [Transaction]. One can retrieve the
1665/// [post_commit_stats], [commit version], and optionally the [post-commit snapshot] from this
1666/// struct.
1667///
1668/// [post_commit_stats]: Self::post_commit_stats
1669/// [commit version]: Self::commit_version
1670/// [post-commit snapshot]: Self::post_commit_snapshot
1671#[derive(Debug)]
1672pub struct CommittedTransaction {
1673    /// The version of the table that was just committed.
1674    commit_version: Version,
1675    /// The [`PostCommitStats`] for this transaction.
1676    post_commit_stats: PostCommitStats,
1677    /// The [`SnapshotRef`] of the table after this transaction was committed.
1678    ///
1679    /// This is optional to allow incremental development of new features (e.g., table creation,
1680    /// transaction retries) without blocking on implementing post-commit snapshot support.
1681    post_commit_snapshot: Option<SnapshotRef>,
1682}
1683
1684impl CommittedTransaction {
1685    /// The version of the table that was just sucessfully committed
1686    pub fn commit_version(&self) -> Version {
1687        self.commit_version
1688    }
1689
1690    /// The [`PostCommitStats`] for this transaction
1691    pub fn post_commit_stats(&self) -> &PostCommitStats {
1692        &self.post_commit_stats
1693    }
1694
1695    /// The [`SnapshotRef`] of the table after this transaction was committed.
1696    pub fn post_commit_snapshot(&self) -> Option<&SnapshotRef> {
1697        self.post_commit_snapshot.as_ref()
1698    }
1699}
1700
1701/// This is the result of a conflicted [Transaction]. One can retrieve the [conflict version] from
1702/// this struct. In the future a rebase API will be provided (issue #1389).
1703///
1704/// [conflict version]: Self::conflict_version
1705#[derive(Debug)]
1706pub struct ConflictedTransaction<S = ExistingTable> {
1707    // TODO: remove after rebase APIs
1708    #[allow(dead_code)]
1709    transaction: Transaction<S>,
1710    conflict_version: Version,
1711}
1712
1713impl<S> ConflictedTransaction<S> {
1714    /// The version attempted commit that yielded a conflict
1715    pub fn conflict_version(&self) -> Version {
1716        self.conflict_version
1717    }
1718}
1719
1720/// A transaction that failed to commit due to a retryable error (e.g. IO error). The transaction
1721/// can be recovered with `RetryableTransaction::transaction` and retried without rebasing. The
1722/// associated error can be inspected via `RetryableTransaction::error`.
1723#[derive(Debug)]
1724pub struct RetryableTransaction<S = ExistingTable> {
1725    /// The transaction that failed to commit due to a retryable error.
1726    pub transaction: Transaction<S>,
1727    /// Transient error that caused the commit to fail.
1728    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        /// Set clustering columns for testing purposes without needing a table
1774        /// with the ClusteredTable feature enabled.
1775        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    /// A mock committer that always returns an IOError, used to test the retryable error path.
1782    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    /// A mock committer that always returns a non-retryable (non-IO) error, used to test the
1806    /// terminal error path.
1807    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    /// A mock catalog committer, used to test catalog committer validation.
1831    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            // This won't be reached in tests — the validation error fires before commit.
1841            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    /// Sets up a snapshot for a table with deletion vector support at version 1
1856    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    /// Creates a test deletion vector descriptor with default values (the DV might not exist on
1915    /// disk)
1916    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    // TODO: create a finer-grained unit tests for transactions (issue#1091)
1940    #[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        // Test with empty prefix
2035        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        // Test with non-empty prefix
2040        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        // Test that two paths with same prefix are different (unique UUIDs)
2047        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        // Regression coverage for stale SharedWriteState caching: keep the first context alive
2064        // while the transaction's effective table config changes.
2065        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        // NB: `test_utils::schema_with_column_defaults` cannot be used here. In `--lib` unit tests
2109        // the crate under test and the `delta_kernel` that `test_utils` links are two distinct
2110        // crate instances, so kernel schema types don't unify across the `test_utils` boundary.
2111
2112        /// Builds a transaction whose effective logical schema is `schema`, with the
2113        /// `allowColumnDefaults` writer feature enabled so any declared defaults are honored.
2114        fn txn_with_schema(schema: StructType) -> Transaction {
2115            txn_with_schema_and_writer_features(schema, [TableFeature::AllowColumnDefaults])
2116        }
2117
2118        /// Like [`txn_with_schema`] but with an explicit writer-feature list, so a test can
2119        /// exercise a table that does *not* enable `allowColumnDefaults`. The schema and a
2120        /// synthetic protocol are swapped onto a real snapshot's table configuration so
2121        /// column-default discovery can be exercised without going through `create_table`.
2122        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        /// A nullable field carrying `raw_sql` as its `CURRENT_DEFAULT`.
2150        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        // Both schemas exclude partition columns.
2245        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        // Both should contain the non-partition columns
2255        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    /// Loads a snapshot from `table_path` and builds a partitioned write context for the given
2269    /// partition values. The table must be partitioned.
2270    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    /// Helper: evaluates the logical-to-physical transform on the given batch and returns the
2286    /// output RecordBatch.
2287    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        // The input data must exclude the partition column "letter".
2327        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        // Logical order: [p1, p2, d1, v(void), p3, p4, d2]; partition cols = p1, p2, p3, p4.
2370        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        // Input excludes partition columns but keeps the void column, in logical schema
2395        // order: [d1, v, d2].
2396        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        // With void stripped and partition literals inserted, the output names/order must match
2412        // the physical schema exactly.
2413        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        // Verify the transformed data.
2427        assert_eq!(get_column!(rb, names[0], StringArray).value(0), "aa"); // p1 (prepended)
2428        assert_eq!(get_column!(rb, names[1], Int32Array).value(0), 7); // p2 (prepended)
2429        assert_eq!(get_column!(rb, names[2], Int32Array).value(0), 10); // d1
2430        assert_eq!(get_column!(rb, names[3], StringArray).value(0), "cc"); // p3 (after d1, void skipped)
2431        assert_eq!(get_column!(rb, names[4], Int32Array).value(0), 9); // p4 (after d1)
2432        assert_eq!(get_column!(rb, names[5], Int32Array).value(0), 20); // d2
2433        Ok(())
2434    }
2435
2436    /// Physical schema should include partition columns when materializePartitionColumns is on.
2437    #[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    /// Using the wrong write context method for the table's partitioning returns an error.
2458    #[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    /// Tests that update_deletion_vectors validates table protocol requirements.
2493    /// Validates that attempting DV updates on unsupported tables returns protocol error.
2494    #[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    /// Tests that update_deletion_vectors validates DV descriptors match scan files.
2534    /// Validates detection of mismatch between provided DV descriptors and actual files.
2535    #[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    /// Tests that a mismatch after scanning some files does not leave staged DV updates behind.
2558    #[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    /// Tests that update_deletion_vectors handles empty DV updates correctly as a no-op.
2651    /// This edge case occurs when a DELETE operation matches no rows.
2652    #[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    // ============================================================================
2669    // validate_blind_append tests
2670    // ============================================================================
2671    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        // CreateTableTransaction does not expose with_blind_append() (compile-time
2752        // prevention per #1768). Directly set the field to test the runtime check.
2753        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        // No files added — commit should fail with blind append validation
2775        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        // Blind append with add files should pass validation and proceed to commit.
2792        // The commit itself may fail due to schema mismatch with the dummy data,
2793        // but we verify validation (line 415) passes on the Ok path.
2794        let result = txn.commit(engine.as_ref());
2795        // If it fails, it should NOT be an InvalidTransactionState error
2796        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    // Note: Additional test coverage for partial file matching (where some files in a scan
2806    // have DV updates but others don't) is provided by the end-to-end integration test
2807    // kernel/tests/features/dv.rs and kernel/tests/write/remove_dv.rs, which exercise
2808    // the full deletion vector write workflow including the DvMatchVisitor logic.
2809
2810    #[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        // Existing-table transactions should include the snapshot version number
2835        assert!(
2836            debug_str.contains("Transaction") && debug_str.contains("read_snapshot version"),
2837            "Debug output should contain Transaction info: {debug_str}"
2838        );
2839        // Should NOT contain "create_table"
2840        assert!(
2841            !debug_str.contains("create_table"),
2842            "Existing table debug should not contain create_table: {debug_str}"
2843        );
2844        Ok(())
2845    }
2846
2847    // Input schemas have no CM metadata; create_table automatically assigns IDs and
2848    // physical names when mode is Name or Id.
2849    #[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    /// Builds two-row [`EngineData`] with logical field names matching [`test_schema_nested`].
2877    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    /// Validates that [`WriteContext::logical_to_physical`] correctly renames fields at all nesting
2909    /// levels. Builds a RecordBatch with logical names, evaluates the transform, and checks
2910    /// that the output uses physical names from the physical schema — including nested struct
2911    /// children.
2912    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        // Evaluate the logical_to_physical expression
2930        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        // Verify: all field names, types, and metadata match the physical schema
2942        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        // Verify: data is preserved (id values)
2946        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    // =========================================================================
2965    // Stats validation tests for clustering columns
2966    // =========================================================================
2967
2968    /// Per-file stats configuration for test add file helpers.
2969    enum TestFileStats {
2970        /// No stats (null stats struct)
2971        None,
2972        /// Normal stats with non-null min/max
2973        Present,
2974        /// All-null column: nullCount == numRecords, null min/max
2975        AllNull,
2976    }
2977
2978    /// Creates test add file metadata with configurable stats for the "value" column.
2979    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            // Enable clustering columns for this test
3089            .with_clustering_columns_for_test(vec![ColumnName::new(["value"])]);
3090
3091        // Add files WITHOUT stats
3092        let add_files = create_test_add_files(vec!["file1.parquet"], vec![TestFileStats::None]);
3093
3094        // Directly test the validation method instead of committing
3095        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            // Enable clustering columns for this test
3117            .with_clustering_columns_for_test(vec![ColumnName::new(["value"])]);
3118
3119        // Add files WITH stats
3120        let add_files = create_test_add_files(vec!["file1.parquet"], vec![TestFileStats::Present]);
3121
3122        // Directly test the validation method
3123        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        // No clustering columns set (default)
3139
3140        // Add files WITHOUT stats
3141        let add_files = create_test_add_files(vec!["file1.parquet"], vec![TestFileStats::None]);
3142
3143        // Directly test the validation method - should pass because no clustering
3144        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        // Create a non-catalog-managed table (no catalogManaged feature)
3159        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        // Try to commit with a catalog committer to a non-catalog-managed table
3173        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        // Create a non-catalog-managed table using a catalog committer
3191        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        // Set up a table with ICT enabled and a very high previous ICT so that the
3249        // monotonicity rule (max(wall_time, prev_ict + 1)) produces a value strictly
3250        // greater than the current wall time. This lets us verify the computed ICT is
3251        // passed to CommitMetadata (not the wall-clock timestamp).
3252        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; // far-future timestamp in ms
3257        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        // The ICT in CommitMetadata must be prev_ict + 1 (monotonicity), NOT the wall time.
3313        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    // ===== Commit failure-metric tests =====
3327
3328    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}