Skip to main content

iceberg/transaction/
append.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use std::collections::{HashMap, HashSet};
19use std::sync::Arc;
20
21use async_trait::async_trait;
22use uuid::Uuid;
23
24use crate::error::Result;
25use crate::spec::{DataFile, ManifestEntry, ManifestFile, Operation};
26use crate::table::Table;
27use crate::transaction::snapshot::{
28    DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
29};
30use crate::transaction::{ActionCommit, TransactionAction};
31
32/// FastAppendAction is a transaction action for fast append data files to the table.
33pub struct FastAppendAction {
34    check_duplicate: bool,
35    // below are properties used to create SnapshotProducer when commit
36    commit_uuid: Option<Uuid>,
37    snapshot_properties: HashMap<String, String>,
38    added_data_files: Vec<DataFile>,
39}
40
41impl FastAppendAction {
42    pub(crate) fn new() -> Self {
43        Self {
44            check_duplicate: true,
45            commit_uuid: None,
46            snapshot_properties: HashMap::default(),
47            added_data_files: vec![],
48        }
49    }
50
51    /// Set whether to check duplicate files
52    pub fn with_check_duplicate(mut self, v: bool) -> Self {
53        self.check_duplicate = v;
54        self
55    }
56
57    /// Add data files to the snapshot.
58    pub fn add_data_files(mut self, data_files: impl IntoIterator<Item = DataFile>) -> Self {
59        self.added_data_files.extend(data_files);
60        self
61    }
62
63    /// Set commit UUID for the snapshot.
64    pub fn set_commit_uuid(mut self, commit_uuid: Uuid) -> Self {
65        self.commit_uuid = Some(commit_uuid);
66        self
67    }
68
69    /// Set snapshot summary properties.
70    pub fn set_snapshot_properties(mut self, snapshot_properties: HashMap<String, String>) -> Self {
71        self.snapshot_properties = snapshot_properties;
72        self
73    }
74
75    /// Collapse files sharing a path to their first occurrence, so a single
76    /// manifest never references the same file twice. Always runs (unlike the
77    /// `check_duplicate`-gated cross-snapshot check) since it is in-memory only.
78    fn dedupe_added_files(&self) -> Vec<DataFile> {
79        let mut seen = HashSet::with_capacity(self.added_data_files.len());
80        self.added_data_files
81            .iter()
82            .filter(|data_file| seen.insert(data_file.file_path.as_str()))
83            .cloned()
84            .collect()
85    }
86}
87
88#[async_trait]
89impl TransactionAction for FastAppendAction {
90    async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
91        let snapshot_producer = SnapshotProducer::new(
92            table,
93            self.commit_uuid.unwrap_or_else(Uuid::now_v7),
94            self.snapshot_properties.clone(),
95            self.dedupe_added_files(),
96        );
97
98        // validate added files
99        snapshot_producer.validate_added_data_files()?;
100
101        // Checks duplicate files
102        if self.check_duplicate {
103            snapshot_producer.validate_duplicate_files().await?;
104        }
105
106        snapshot_producer
107            .commit(FastAppendOperation, DefaultManifestProcess)
108            .await
109    }
110}
111
112struct FastAppendOperation;
113
114impl SnapshotProduceOperation for FastAppendOperation {
115    fn operation(&self) -> Operation {
116        Operation::Append
117    }
118
119    async fn delete_entries(
120        &self,
121        _snapshot_produce: &SnapshotProducer<'_>,
122    ) -> Result<Vec<ManifestEntry>> {
123        Ok(vec![])
124    }
125
126    async fn existing_manifest(
127        &self,
128        snapshot_produce: &SnapshotProducer<'_>,
129    ) -> Result<Vec<ManifestFile>> {
130        let Some(snapshot) = snapshot_produce.table.metadata().current_snapshot() else {
131            return Ok(vec![]);
132        };
133
134        let manifest_list = snapshot_produce
135            .table
136            .manifest_list_reader(snapshot)
137            .load()
138            .await?;
139
140        Ok(manifest_list
141            .entries()
142            .iter()
143            .filter(|entry| {
144                // Keep delete-only manifests too: they record which files were removed and
145                // must persist across snapshots until `expire_snapshots` cleans them up.
146                // Dropping them lets the removed files reappear as live data (see #2148).
147                entry.has_added_files() || entry.has_existing_files() || entry.has_deleted_files()
148            })
149            .cloned()
150            .collect())
151    }
152}
153
154#[cfg(test)]
155mod tests {
156    use std::collections::HashMap;
157    use std::fs;
158    use std::sync::Arc;
159
160    use minijinja::{AutoEscape, Environment, Value, context};
161    use tempfile::TempDir;
162    use uuid::Uuid;
163
164    use crate::encryption::kms::MemoryKeyManagementClient;
165    use crate::encryption::{SensitiveBytes, StandardKeyMetadata};
166    use crate::io::FileIO;
167    use crate::spec::{
168        DataContentType, DataFile, DataFileBuilder, DataFileFormat, Literal, MAIN_BRANCH,
169        ManifestEntry, ManifestListWriter, ManifestStatus, ManifestWriterBuilder, SnapshotRef,
170        Struct, TableMetadata,
171    };
172    use crate::table::Table;
173    use crate::test_utils::{make_encrypted_table, test_runtime};
174    use crate::transaction::tests::make_v2_minimal_table;
175    use crate::transaction::{Transaction, TransactionAction};
176    use crate::{TableIdent, TableRequirement, TableUpdate};
177
178    fn render_template(template: &str, ctx: Value) -> String {
179        let mut env = Environment::new();
180        env.set_auto_escape_callback(|_| AutoEscape::None);
181        env.render_str(template, ctx).unwrap()
182    }
183
184    /// Builds a table whose current snapshot's manifest list contains a data manifest
185    /// followed by a delete-only manifest (one entry with `ManifestStatus::Deleted`,
186    /// so `deleted_files_count > 0` while `added_files_count == existing_files_count == 0`).
187    ///
188    /// Returns the table plus the `manifest_path` of the delete-only manifest so callers
189    /// can assert whether a subsequent append carries it forward.
190    async fn make_table_with_delete_only_manifest() -> (Table, TempDir, String) {
191        let tmp_dir = TempDir::new().unwrap();
192        let table_location = tmp_dir.path().join("table1");
193        let manifest_list_location = table_location.join("metadata/manifests_list_1.avro");
194        let table_metadata_location = table_location.join("metadata/v1.json");
195
196        let file_io = FileIO::new_with_fs();
197
198        let template = fs::read_to_string(format!(
199            "{}/testdata/example_table_metadata_v2.json",
200            env!("CARGO_MANIFEST_DIR")
201        ))
202        .unwrap();
203        // The template has two snapshots; point the current one at our manifest list.
204        let metadata_json = render_template(&template, context! {
205            table_location => &table_location,
206            manifest_list_1_location => &manifest_list_location,
207            manifest_list_2_location => &manifest_list_location,
208            table_metadata_1_location => &table_metadata_location,
209        });
210        let table_metadata = serde_json::from_str::<TableMetadata>(&metadata_json).unwrap();
211
212        let table = Table::builder()
213            .metadata(table_metadata)
214            .identifier(TableIdent::from_strs(["db", "table1"]).unwrap())
215            .file_io(file_io)
216            .metadata_location(table_metadata_location.to_str().unwrap())
217            .runtime(test_runtime())
218            .build()
219            .unwrap();
220
221        let current_snapshot = table.metadata().current_snapshot().unwrap();
222        let schema = current_snapshot.schema(table.metadata()).unwrap();
223        let partition_spec = table.metadata().default_partition_spec();
224
225        let next_manifest_file = |location: &str| {
226            table
227                .file_io()
228                .new_output(format!(
229                    "{}/metadata/manifest_{}.avro",
230                    location,
231                    Uuid::new_v4()
232                ))
233                .unwrap()
234        };
235        let table_location_str = table_location.to_str().unwrap().to_string();
236
237        // Data manifest: one Added data file.
238        let mut data_writer = ManifestWriterBuilder::new(
239            next_manifest_file(&table_location_str),
240            Some(current_snapshot.snapshot_id()),
241            schema.clone(),
242            partition_spec.as_ref().clone(),
243        )
244        .build_v2_data();
245        data_writer
246            .add_entry(
247                ManifestEntry::builder()
248                    .status(ManifestStatus::Added)
249                    .data_file(
250                        DataFileBuilder::default()
251                            .partition_spec_id(0)
252                            .content(DataContentType::Data)
253                            .file_path(format!("{table_location_str}/data.parquet"))
254                            .file_format(DataFileFormat::Parquet)
255                            .file_size_in_bytes(100)
256                            .record_count(1)
257                            .partition(Struct::from_iter([Some(Literal::long(100))]))
258                            .build()
259                            .unwrap(),
260                    )
261                    .build(),
262            )
263            .unwrap();
264        let data_manifest = data_writer.write_manifest_file().await.unwrap();
265
266        // Delete-only manifest: a single Deleted entry, nothing added or existing.
267        let mut delete_writer = ManifestWriterBuilder::new(
268            next_manifest_file(&table_location_str),
269            Some(current_snapshot.snapshot_id()),
270            schema.clone(),
271            partition_spec.as_ref().clone(),
272        )
273        .build_v2_data();
274        delete_writer
275            .add_delete_entry(
276                ManifestEntry::builder()
277                    .status(ManifestStatus::Deleted)
278                    .sequence_number(0)
279                    .file_sequence_number(0)
280                    .data_file(
281                        DataFileBuilder::default()
282                            .partition_spec_id(0)
283                            .content(DataContentType::Data)
284                            .file_path(format!("{table_location_str}/removed.parquet"))
285                            .file_format(DataFileFormat::Parquet)
286                            .file_size_in_bytes(100)
287                            .record_count(1)
288                            .partition(Struct::from_iter([Some(Literal::long(100))]))
289                            .build()
290                            .unwrap(),
291                    )
292                    .build(),
293            )
294            .unwrap();
295        let delete_manifest = delete_writer.write_manifest_file().await.unwrap();
296        let delete_manifest_path = delete_manifest.manifest_path.clone();
297
298        // Sanity: the delete manifest really is delete-only.
299        assert!(delete_manifest.has_deleted_files());
300        assert!(!delete_manifest.has_added_files());
301        assert!(!delete_manifest.has_existing_files());
302
303        let mut manifest_list_writer = ManifestListWriter::v2(
304            table
305                .file_io()
306                .new_output(current_snapshot.manifest_list())
307                .unwrap()
308                .writer()
309                .await
310                .unwrap(),
311            current_snapshot.snapshot_id(),
312            current_snapshot.parent_snapshot_id(),
313            current_snapshot.sequence_number(),
314        );
315        manifest_list_writer
316            .add_manifests(vec![data_manifest, delete_manifest].into_iter())
317            .unwrap();
318        manifest_list_writer.close().await.unwrap();
319
320        (table, tmp_dir, delete_manifest_path)
321    }
322
323    /// Regression test for #2148: a `fast_append` must carry delete-only manifests
324    /// forward into the new snapshot. Dropping them lets the files they mark as
325    /// removed reappear as live data on the next append.
326    #[tokio::test]
327    async fn test_fast_append_preserves_delete_only_manifest() {
328        let (table, _tmp_dir, delete_manifest_path) = make_table_with_delete_only_manifest().await;
329
330        // Append a new data file via the public transaction API.
331        let new_file = DataFileBuilder::default()
332            .content(DataContentType::Data)
333            .file_path(format!("{}/appended.parquet", table.metadata().location()))
334            .file_format(DataFileFormat::Parquet)
335            .file_size_in_bytes(100)
336            .record_count(1)
337            .partition_spec_id(table.metadata().default_partition_spec_id())
338            .partition(Struct::from_iter([Some(Literal::long(100))]))
339            .build()
340            .unwrap();
341
342        let tx = Transaction::new(&table);
343        let action = tx.fast_append().add_data_files(vec![new_file]);
344        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
345        let updates = action_commit.take_updates();
346
347        let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
348            SnapshotRef::new(snapshot.clone())
349        } else {
350            unreachable!("first update of a fast append should be AddSnapshot")
351        };
352
353        let manifest_list = table
354            .manifest_list_reader(&new_snapshot)
355            .load()
356            .await
357            .unwrap();
358
359        assert!(
360            manifest_list
361                .entries()
362                .iter()
363                .any(|m| m.manifest_path == delete_manifest_path),
364            "delete-only manifest {delete_manifest_path} was dropped from the new snapshot's \
365             manifest list; the files it removed would reappear as live data"
366        );
367    }
368
369    /// Load the data files written by a single-manifest fast-append commit.
370    async fn committed_data_files(table: &Table, updates: &[TableUpdate]) -> Vec<DataFile> {
371        let TableUpdate::AddSnapshot { snapshot } = &updates[0] else {
372            unreachable!("first update is always AddSnapshot")
373        };
374        let manifest_list = table
375            .manifest_list_reader(&SnapshotRef::new(snapshot.clone()))
376            .load()
377            .await
378            .unwrap();
379        assert_eq!(1, manifest_list.entries().len());
380        table
381            .manifest_reader()
382            .read(&manifest_list.entries()[0])
383            .await
384            .unwrap()
385            .entries()
386            .iter()
387            .map(|entry| entry.data_file().clone())
388            .collect()
389    }
390
391    #[tokio::test]
392    async fn test_fast_append_writes_encrypted_manifest() {
393        let table = make_encrypted_table().await;
394        assert!(
395            table.encryption_manager().is_some(),
396            "fixture table should have an EncryptionManager"
397        );
398
399        let new_file = DataFileBuilder::default()
400            .content(DataContentType::Data)
401            .file_path("memory:///table/data/00000.parquet".to_string())
402            .file_format(DataFileFormat::Parquet)
403            .partition(Struct::empty())
404            .record_count(100)
405            .file_size_in_bytes(4096)
406            .partition_spec_id(table.metadata().default_partition_spec_id())
407            .build()
408            .unwrap();
409
410        let tx = Transaction::new(&table);
411        let action = tx.fast_append().add_data_files(vec![new_file]);
412        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
413        let updates = action_commit.take_updates();
414
415        let new_snapshot: SnapshotRef = updates
416            .iter()
417            .find_map(|u| match u {
418                TableUpdate::AddSnapshot { snapshot } => Some(SnapshotRef::new(snapshot.clone())),
419                _ => None,
420            })
421            .expect("a fast append should emit an AddSnapshot update");
422
423        let manifest_list_key_metadata = table
424            .encryption_manager()
425            .unwrap()
426            .decrypt_manifest_list_key_metadata(new_snapshot.encryption_key_id().unwrap())
427            .await
428            .unwrap();
429        let manifest_list_size = table
430            .file_io()
431            .new_input(new_snapshot.manifest_list())
432            .unwrap()
433            .metadata()
434            .await
435            .unwrap()
436            .size;
437        assert_eq!(
438            manifest_list_key_metadata.file_length(),
439            Some(manifest_list_size)
440        );
441
442        let manifest_list = table
443            .manifest_list_reader(&new_snapshot)
444            .load()
445            .await
446            .unwrap();
447        let manifest_file = manifest_list
448            .entries()
449            .iter()
450            .find(|m| m.added_files_count.unwrap_or(0) > 0)
451            .expect("new snapshot should carry the appended data manifest");
452
453        // The manifest list entry must carry decodable key metadata.
454        let key_metadata_bytes = manifest_file
455            .key_metadata
456            .as_ref()
457            .expect("encrypted manifest must record key metadata");
458        let key_metadata = StandardKeyMetadata::decode(key_metadata_bytes)
459            .expect("recorded key metadata must decode as StandardKeyMetadata");
460        let manifest_size = table
461            .file_io()
462            .new_input(&manifest_file.manifest_path)
463            .unwrap()
464            .metadata()
465            .await
466            .unwrap()
467            .size;
468        assert_eq!(key_metadata.file_length(), Some(manifest_size));
469        assert_eq!(manifest_file.manifest_length, manifest_size as i64);
470
471        // The reader self-decrypts using the recorded key metadata and must
472        // recover the entry we appended. Because the read goes through the
473        // decryption path, this succeeding also proves the bytes
474        // on disk were genuinely encrypted (not silently written as plaintext).
475        let manifest = table.manifest_reader().read(manifest_file).await.unwrap();
476        assert_eq!(manifest.entries().len(), 1);
477        assert_eq!(
478            manifest.entries()[0].data_file().file_path(),
479            "memory:///table/data/00000.parquet"
480        );
481    }
482
483    #[tokio::test]
484    async fn test_empty_data_append_action() {
485        let table = make_v2_minimal_table();
486        let tx = Transaction::new(&table);
487        let action = tx.fast_append().add_data_files(vec![]);
488        assert!(Arc::new(action).commit(&table).await.is_err());
489    }
490
491    /// A `fast_append` must write the manifest list and the manifest
492    /// files under the `write.metadata.path` prefix when configured,
493    /// rather than the default `<location>/metadata` directory.
494    #[tokio::test]
495    async fn test_fast_append_honors_write_metadata_path() {
496        let base = make_v2_minimal_table();
497        let metadata_root = format!("{}/custom-meta", base.metadata().location());
498        let metadata = base
499            .metadata()
500            .clone()
501            .into_builder(None)
502            .set_properties(HashMap::from([(
503                "write.metadata.path".to_string(),
504                metadata_root.clone(),
505            )]))
506            .unwrap()
507            .build()
508            .unwrap()
509            .metadata;
510        let table = base.with_metadata(Arc::new(metadata));
511
512        let data_file = DataFileBuilder::default()
513            .content(DataContentType::Data)
514            .file_path(format!("{}/data/1.parquet", table.metadata().location()))
515            .file_format(DataFileFormat::Parquet)
516            .file_size_in_bytes(100)
517            .record_count(1)
518            .partition_spec_id(table.metadata().default_partition_spec_id())
519            .partition(Struct::from_iter([Some(Literal::long(300))]))
520            .build()
521            .unwrap();
522
523        let tx = Transaction::new(&table);
524        let action = tx.fast_append().add_data_files(vec![data_file]);
525        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
526        let updates = action_commit.take_updates();
527
528        let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
529            SnapshotRef::new(snapshot.clone())
530        } else {
531            unreachable!("first update of a fast append should be AddSnapshot")
532        };
533
534        let prefix = format!("{metadata_root}/");
535
536        // Manifest list
537        assert!(
538            new_snapshot.manifest_list().starts_with(prefix.as_str()),
539            "manifest list {} not under configured write.metadata.path {metadata_root}",
540            new_snapshot.manifest_list()
541        );
542
543        // Manifest files
544        let manifest_list = table
545            .manifest_list_reader(&new_snapshot)
546            .load()
547            .await
548            .unwrap();
549        assert!(
550            !manifest_list.entries().is_empty(),
551            "expected at least one manifest entry"
552        );
553        for entry in manifest_list.entries() {
554            assert!(
555                entry.manifest_path.starts_with(prefix.as_str()),
556                "manifest {} not under configured write.metadata.path {metadata_root}",
557                entry.manifest_path
558            );
559        }
560    }
561
562    #[tokio::test]
563    async fn test_set_snapshot_properties() {
564        let table = make_v2_minimal_table();
565        let tx = Transaction::new(&table);
566
567        let mut snapshot_properties = HashMap::new();
568        snapshot_properties.insert("key".to_string(), "val".to_string());
569
570        let data_file = DataFileBuilder::default()
571            .content(DataContentType::Data)
572            .file_path("test/1.parquet".to_string())
573            .file_format(DataFileFormat::Parquet)
574            .file_size_in_bytes(100)
575            .record_count(1)
576            .partition_spec_id(table.metadata().default_partition_spec_id())
577            .partition(Struct::from_iter([Some(Literal::long(300))]))
578            .build()
579            .unwrap();
580
581        let action = tx
582            .fast_append()
583            .set_snapshot_properties(snapshot_properties)
584            .add_data_files(vec![data_file]);
585        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
586        let updates = action_commit.take_updates();
587
588        // Check customized properties is contained in snapshot summary properties.
589        let new_snapshot = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
590            snapshot
591        } else {
592            unreachable!()
593        };
594        assert_eq!(
595            new_snapshot
596                .summary()
597                .additional_properties
598                .get("key")
599                .unwrap(),
600            "val"
601        );
602    }
603
604    /// See `testdata/manifests_lists/README.md`.
605    const FIXTURE_MASTER_KEY_ID: &str = "master-1";
606    const FIXTURE_MASTER_KEY_BYTES: [u8; 16] = [
607        0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, 0x0d, 0x0e,
608        0x0f,
609    ];
610
611    async fn make_v3_encrypted_table() -> Table {
612        let json = fs::read_to_string(format!(
613            "{}/testdata/table_metadata/TableMetadataV3ValidEncryption.json",
614            env!("CARGO_MANIFEST_DIR")
615        ))
616        .unwrap();
617        let metadata = serde_json::from_str::<TableMetadata>(&json).unwrap();
618
619        let kms = MemoryKeyManagementClient::new();
620        kms.add_master_key_bytes(
621            FIXTURE_MASTER_KEY_ID,
622            SensitiveBytes::new(FIXTURE_MASTER_KEY_BYTES),
623        )
624        .unwrap();
625
626        let file_io = FileIO::new_with_memory();
627
628        let manifest_list_bytes = fs::read(format!(
629            "{}/testdata/manifests_lists/manifest-list-v3-encrypted.avro",
630            env!("CARGO_MANIFEST_DIR")
631        ))
632        .unwrap();
633        let parent_manifest_list = metadata.current_snapshot().unwrap().manifest_list();
634        file_io
635            .new_output(parent_manifest_list)
636            .unwrap()
637            .write(manifest_list_bytes.into())
638            .await
639            .unwrap();
640
641        Table::builder()
642            .metadata(metadata)
643            .metadata_location("memory:///table/metadata/v1.json")
644            .identifier(TableIdent::from_strs(["ns1", "enc"]).unwrap())
645            .file_io(file_io)
646            .kms_client(Arc::new(kms))
647            .runtime(test_runtime())
648            .build()
649            .unwrap()
650    }
651
652    #[tokio::test]
653    async fn test_commit_with_encryption_adds_keys_and_records_snapshot_key_id() {
654        let table = make_v3_encrypted_table().await;
655
656        let data_file = DataFileBuilder::default()
657            .content(DataContentType::Data)
658            .file_path("test/1.parquet".to_string())
659            .file_format(DataFileFormat::Parquet)
660            .file_size_in_bytes(100)
661            .record_count(1)
662            .partition_spec_id(table.metadata().default_partition_spec_id())
663            .partition(Struct::empty())
664            .build()
665            .unwrap();
666
667        let tx = Transaction::new(&table);
668        let action = tx.fast_append().add_data_files(vec![data_file]);
669        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
670        let updates = action_commit.take_updates();
671
672        let added_key_ids: Vec<String> = updates
673            .iter()
674            .filter_map(|u| match u {
675                TableUpdate::AddEncryptionKey { encryption_key } => {
676                    Some(encryption_key.key_id().to_string())
677                }
678                _ => None,
679            })
680            .collect();
681        assert!(!added_key_ids.is_empty(), "got {updates:?}");
682
683        // Encryption keys are added before the snapshot, so it isn't updates[0] here.
684        let new_snapshot = updates
685            .iter()
686            .find_map(|u| match u {
687                TableUpdate::AddSnapshot { snapshot } => Some(snapshot),
688                _ => None,
689            })
690            .expect("commit should add a snapshot");
691
692        let snapshot_key_id = new_snapshot
693            .encryption_key_id()
694            .expect("encrypted snapshot should record its manifest-list key id");
695        assert!(
696            added_key_ids.iter().any(|id| id == snapshot_key_id),
697            "snapshot key id {snapshot_key_id} not in added keys {added_key_ids:?}"
698        );
699
700        let new_snapshot_ref: SnapshotRef = Arc::new(new_snapshot.clone());
701        let manifest_list = table
702            .manifest_list_reader(&new_snapshot_ref)
703            .load()
704            .await
705            .expect("newly written encrypted manifest list should decrypt and parse");
706        assert_eq!(
707            manifest_list.entries().len(),
708            1,
709            "append should record exactly the one new data manifest"
710        );
711    }
712
713    #[tokio::test]
714    async fn test_snapshot_properties_cannot_override_computed_metrics() {
715        // A user-supplied snapshot property must not shadow a computed metric key
716        // such as `added-data-files`. Matching iceberg-java, the computed value
717        // wins, so the summary reflects the real count and a bad value can neither
718        // corrupt the summary nor panic total computation (see #2184-adjacent fix).
719        let table = make_v2_minimal_table();
720        let tx = Transaction::new(&table);
721
722        let mut snapshot_properties = HashMap::new();
723        // Both a benign-but-wrong value and a non-integer value collide with
724        // computed metric keys; neither should reach the final summary.
725        snapshot_properties.insert("added-data-files".to_string(), "9999".to_string());
726        snapshot_properties.insert("added-records".to_string(), "not-a-number".to_string());
727
728        let data_file = DataFileBuilder::default()
729            .content(DataContentType::Data)
730            .file_path("test/1.parquet".to_string())
731            .file_format(DataFileFormat::Parquet)
732            .file_size_in_bytes(100)
733            .record_count(1)
734            .partition_spec_id(table.metadata().default_partition_spec_id())
735            .partition(Struct::from_iter([Some(Literal::long(300))]))
736            .build()
737            .unwrap();
738
739        let action = tx
740            .fast_append()
741            .set_snapshot_properties(snapshot_properties)
742            .add_data_files(vec![data_file]);
743        // Must not panic during total computation.
744        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
745        let updates = action_commit.take_updates();
746
747        let new_snapshot = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
748            snapshot
749        } else {
750            unreachable!()
751        };
752        let props = &new_snapshot.summary().additional_properties;
753
754        // Computed metric wins over the user's colliding values.
755        assert_eq!(
756            props.get("added-data-files").unwrap(),
757            "1",
758            "computed added-data-files must override the user-supplied value"
759        );
760        assert_eq!(
761            props.get("added-records").unwrap(),
762            "1",
763            "computed added-records must override the user-supplied non-integer value"
764        );
765    }
766
767    #[tokio::test]
768    async fn test_append_snapshot_properties() {
769        let table = make_v2_minimal_table();
770        let tx = Transaction::new(&table);
771
772        let mut snapshot_properties = HashMap::new();
773        snapshot_properties.insert("key".to_string(), "val".to_string());
774
775        let action = tx
776            .fast_append()
777            .set_snapshot_properties(snapshot_properties);
778        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
779        let updates = action_commit.take_updates();
780
781        // Check customized properties is contained in snapshot summary properties.
782        let new_snapshot = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
783            snapshot
784        } else {
785            unreachable!()
786        };
787        assert_eq!(
788            new_snapshot
789                .summary()
790                .additional_properties
791                .get("key")
792                .unwrap(),
793            "val"
794        );
795    }
796
797    #[tokio::test]
798    async fn test_fast_append_file_with_incompatible_partition_value() {
799        let table = make_v2_minimal_table();
800        let tx = Transaction::new(&table);
801        let action = tx.fast_append();
802
803        // check add data file with incompatible partition value
804        let data_file = DataFileBuilder::default()
805            .content(DataContentType::Data)
806            .file_path("test/3.parquet".to_string())
807            .file_format(DataFileFormat::Parquet)
808            .file_size_in_bytes(100)
809            .record_count(1)
810            .partition_spec_id(table.metadata().default_partition_spec_id())
811            .partition(Struct::from_iter([Some(Literal::string("test"))]))
812            .build()
813            .unwrap();
814
815        let action = action.add_data_files(vec![data_file.clone()]);
816
817        assert!(Arc::new(action).commit(&table).await.is_err());
818    }
819
820    #[tokio::test]
821    async fn test_fast_append_dedupes_intra_batch_duplicate_paths() {
822        let table = make_v2_minimal_table();
823        let tx = Transaction::new(&table);
824
825        let make_file = |size: u64, records: u64| {
826            DataFileBuilder::default()
827                .content(DataContentType::Data)
828                .file_path("test/dup.parquet".to_string())
829                .file_format(DataFileFormat::Parquet)
830                .file_size_in_bytes(size)
831                .record_count(records)
832                .partition_spec_id(table.metadata().default_partition_spec_id())
833                .partition(Struct::from_iter([Some(Literal::long(1))]))
834                .build()
835                .unwrap()
836        };
837
838        // Same path three times: the manifest keeps a single entry, the first one.
839        let action = tx.fast_append().add_data_files(vec![
840            make_file(100, 10),
841            make_file(200, 20),
842            make_file(300, 30),
843        ]);
844        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
845        let files = committed_data_files(&table, &action_commit.take_updates()).await;
846        assert_eq!(1, files.len());
847        assert_eq!(100, files[0].file_size_in_bytes());
848    }
849
850    #[tokio::test]
851    async fn test_fast_append_dedupes_regardless_of_check_duplicate_flag() {
852        let table = make_v2_minimal_table();
853        let tx = Transaction::new(&table);
854
855        let make_file = || {
856            DataFileBuilder::default()
857                .content(DataContentType::Data)
858                .file_path("test/dup.parquet".to_string())
859                .file_format(DataFileFormat::Parquet)
860                .file_size_in_bytes(100)
861                .record_count(10)
862                .partition_spec_id(table.metadata().default_partition_spec_id())
863                .partition(Struct::from_iter([Some(Literal::long(1))]))
864                .build()
865                .unwrap()
866        };
867
868        // `check_duplicate` only gates the cross-snapshot check; intra-batch dedupe runs regardless.
869        let action = tx
870            .fast_append()
871            .with_check_duplicate(false)
872            .add_data_files(vec![make_file(), make_file()]);
873        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
874        let files = committed_data_files(&table, &action_commit.take_updates()).await;
875        assert_eq!(1, files.len());
876    }
877
878    #[tokio::test]
879    async fn test_fast_append() {
880        let table = make_v2_minimal_table();
881        let tx = Transaction::new(&table);
882        let action = tx.fast_append();
883
884        let data_file = DataFileBuilder::default()
885            .content(DataContentType::Data)
886            .file_path("test/3.parquet".to_string())
887            .file_format(DataFileFormat::Parquet)
888            .file_size_in_bytes(100)
889            .record_count(1)
890            .partition_spec_id(table.metadata().default_partition_spec_id())
891            .partition(Struct::from_iter([Some(Literal::long(300))]))
892            .build()
893            .unwrap();
894
895        let action = action.add_data_files(vec![data_file.clone()]);
896        let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
897        let updates = action_commit.take_updates();
898        let requirements = action_commit.take_requirements();
899
900        // check updates and requirements
901        assert!(
902            matches!((&updates[0],&updates[1]), (TableUpdate::AddSnapshot { snapshot },TableUpdate::SetSnapshotRef { reference,ref_name }) if snapshot.snapshot_id() == reference.snapshot_id && ref_name == MAIN_BRANCH)
903        );
904        assert_eq!(
905            vec![
906                TableRequirement::UuidMatch {
907                    uuid: table.metadata().uuid()
908                },
909                TableRequirement::RefSnapshotIdMatch {
910                    r#ref: MAIN_BRANCH.to_string(),
911                    snapshot_id: table.metadata().current_snapshot_id
912                }
913            ],
914            requirements
915        );
916
917        // check manifest list
918        let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
919            SnapshotRef::new(snapshot.clone())
920        } else {
921            unreachable!()
922        };
923        let manifest_list = table
924            .manifest_list_reader(&new_snapshot)
925            .load()
926            .await
927            .unwrap();
928        assert_eq!(1, manifest_list.entries().len());
929        assert_eq!(
930            manifest_list.entries()[0].sequence_number,
931            new_snapshot.sequence_number()
932        );
933
934        // check manifest
935        let manifest = table
936            .manifest_reader()
937            .read(&manifest_list.entries()[0])
938            .await
939            .unwrap();
940        assert_eq!(1, manifest.entries().len());
941        assert_eq!(
942            new_snapshot.sequence_number(),
943            manifest.entries()[0]
944                .sequence_number()
945                .expect("Inherit sequence number by load manifest")
946        );
947
948        assert_eq!(
949            new_snapshot.snapshot_id(),
950            manifest.entries()[0].snapshot_id().unwrap()
951        );
952        assert_eq!(data_file, *manifest.entries()[0].data_file());
953    }
954}