This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new 2fa48384 fix(table): reject duplicate ADD and DELETE file entries 
(#904)
2fa48384 is described below

commit 2fa483845c58dc28bd36a909d3c596a4ab541969
Author: jerry <[email protected]>
AuthorDate: Mon Sep 21 18:38:28 2026 +0800

    fix(table): reject duplicate ADD and DELETE file entries (#904)
---
 crates/paimon/src/table/table_commit.rs | 164 +++++++++++++++++++++++++++++---
 crates/paimon/src/table/table_scan.rs   |   2 +-
 2 files changed, 151 insertions(+), 15 deletions(-)

diff --git a/crates/paimon/src/table/table_commit.rs 
b/crates/paimon/src/table/table_commit.rs
index 6d3f7738..d25a6547 100644
--- a/crates/paimon/src/table/table_commit.rs
+++ b/crates/paimon/src/table/table_commit.rs
@@ -52,6 +52,63 @@ type PartitionBucketKey = (Vec<u8>, i32);
 type RowIdRange = (i64, i64);
 type ExistingRowIdRanges = HashMap<PartitionBucketKey, Vec<RowIdRange>>;
 
+fn validate_file_entries<'a>(entries: impl IntoIterator<Item = &'a 
ManifestEntry>) -> Result<()> {
+    // Mirror Java FileEntry.mergeEntries while also rejecting repeated entries
+    // of the same kind instead of letting two DELETEs cancel each other.
+    #[derive(Default)]
+    struct State {
+        added: bool,
+        deleted: bool,
+        present: bool,
+    }
+
+    let mut files = HashMap::new();
+    for entry in entries {
+        let state = files
+            .entry(entry.identifier())
+            .or_insert_with(State::default);
+        let duplicate = match entry.kind() {
+            FileKind::Add => {
+                if state.present && !state.added {
+                    return Err(crate::Error::DataInvalid {
+                        message: format!(
+                            "File conflict: trying to add file '{}' in bucket 
{} which is already present.",
+                            entry.file().file_name,
+                            entry.bucket(),
+                        ),
+                        source: None,
+                    });
+                }
+                let duplicate = state.added;
+                state.added = true;
+                state.present = true;
+                duplicate
+            }
+            FileKind::Delete => {
+                let duplicate = state.deleted;
+                state.deleted = true;
+                state.present = !state.present;
+                duplicate
+            }
+        };
+        if duplicate {
+            return Err(crate::Error::DataInvalid {
+                message: format!(
+                    "Duplicate {} entry for file '{}' in bucket {}.",
+                    match entry.kind() {
+                        FileKind::Add => "ADD",
+                        FileKind::Delete => "DELETE",
+                    },
+                    entry.file().file_name,
+                    entry.bucket(),
+                ),
+                source: None,
+            });
+        }
+    }
+    Ok(())
+}
+
 fn validate_bucket_ownership(messages: &[CommitMessage]) -> Result<()> {
     let mut owners = HashSet::new();
     for message in messages {
@@ -1289,6 +1346,8 @@ impl TableCommit {
                 new_index_entries,
                 check_from_snapshot,
             } => {
+                validate_file_entries(entries.iter())?;
+
                 // Auto-promote to OVERWRITE when CoW rewrites produce Delete 
entries.
                 // This ensures the snapshot correctly reflects file 
replacements.
                 let has_delete = entries.iter().any(|e| *e.kind() == 
FileKind::Delete);
@@ -1365,6 +1424,8 @@ impl TableCommit {
                 let entries = self
                     .provide_overwrite_entries(plan, latest_snapshot)
                     .await?;
+                validate_file_entries(entries.iter())?;
+
                 let (partition_filter, new_index_entries, check_from_snapshot) 
= match plan {
                     CommitEntriesPlan::Overwrite {
                         partition_filter,
@@ -1925,6 +1986,7 @@ impl TableCommit {
         check_from_snapshot: Option<i64>,
     ) -> Result<()> {
         self.check_delete_entries_against_base(base_entries, delta_entries)?;
+        validate_file_entries(base_entries.iter().chain(delta_entries))?;
 
         // Validate delta entries before duplicate files are merged.
         self.check_total_bucket_conflicts(delta_entries)?;
@@ -2123,23 +2185,27 @@ impl TableCommit {
         base_entries: &[ManifestEntry],
         delta_entries: &[ManifestEntry],
     ) -> Result<()> {
-        let base_identifiers = base_entries
+        let mut active_identifiers = base_entries
             .iter()
             .map(ManifestEntry::identifier)
             .collect::<HashSet<_>>();
-        for entry in delta_entries
-            .iter()
-            .filter(|entry| *entry.kind() == FileKind::Delete)
-        {
-            if !base_identifiers.contains(&entry.identifier()) {
-                return Err(crate::Error::DataInvalid {
-                    message: format!(
-                        "Delete conflict: file '{}' in bucket {} does not 
exist in the current snapshot.",
-                        entry.file().file_name,
-                        entry.bucket(),
-                    ),
-                    source: None,
-                });
+        for entry in delta_entries {
+            let identifier = entry.identifier();
+            match entry.kind() {
+                FileKind::Add => {
+                    active_identifiers.insert(identifier);
+                }
+                FileKind::Delete if !active_identifiers.remove(&identifier) => 
{
+                    return Err(crate::Error::DataInvalid {
+                        message: format!(
+                            "Delete conflict: file '{}' in bucket {} does not 
exist in the current snapshot.",
+                            entry.file().file_name,
+                            entry.bucket(),
+                        ),
+                        source: None,
+                    });
+                }
+                FileKind::Delete => {}
             }
         }
         Ok(())
@@ -3459,6 +3525,76 @@ mod tests {
         assert_eq!(entries[0].file().file_name, "data-0.parquet");
     }
 
+    #[tokio::test]
+    async fn test_commit_rejects_conflicting_data_file_entries() {
+        let file_io = test_file_io();
+        let table_path = "memory:/test_duplicate_data_file_entries";
+        setup_dirs(&file_io, table_path).await;
+        let commit = setup_commit(&file_io, table_path);
+        let partition = EMPTY_SERIALIZED_ROW.clone();
+        let file = test_data_file("data-0.parquet", 100);
+
+        let transient = test_data_file("transient.parquet", 1);
+        let mut net_zero = CommitMessage::new(partition.clone(), 0, 
vec![transient.clone()]);
+        net_zero.deleted_files = vec![transient];
+        commit
+            .commit(vec![net_zero])
+            .await
+            .expect("ADD followed by DELETE should cancel out");
+
+        commit
+            .commit(vec![CommitMessage::new(
+                partition.clone(),
+                0,
+                vec![file.clone()],
+            )])
+            .await
+            .unwrap();
+
+        let mut duplicate_delete = CommitMessage::new(partition.clone(), 0, 
vec![]);
+        duplicate_delete.deleted_files = vec![file.clone(), file.clone()];
+        let error = commit
+            .commit(vec![duplicate_delete])
+            .await
+            .expect_err("duplicate DELETE entries must be rejected");
+        assert!(
+            error.to_string().contains("Duplicate DELETE"),
+            "unexpected error: {error}"
+        );
+
+        let duplicate_add = test_data_file("data-1.parquet", 200);
+        let error = commit
+            .commit(vec![CommitMessage::new(
+                partition.clone(),
+                0,
+                vec![duplicate_add.clone(), duplicate_add],
+            )])
+            .await
+            .expect_err("duplicate ADD entries must be rejected");
+        assert!(
+            error.to_string().contains("Duplicate ADD"),
+            "unexpected error: {error}"
+        );
+
+        let mut add_then_delete = CommitMessage::new(partition.clone(), 0, 
vec![file.clone()]);
+        add_then_delete.deleted_files = vec![file.clone()];
+        let error = commit
+            .commit(vec![add_then_delete])
+            .await
+            .expect_err("a commit must reject adding an existing file before 
deleting it");
+        assert!(matches!(error, crate::Error::DataInvalid { .. }));
+
+        let error = commit
+            .overwrite(vec![CommitMessage::new(partition, 0, vec![file])], 
None)
+            .await
+            .expect_err("overwrite must reject deleting and adding the same 
file");
+        assert!(matches!(error, crate::Error::DataInvalid { .. }));
+
+        let snapshot = latest_snapshot(&file_io, table_path).await.unwrap();
+        assert_eq!(snapshot.id(), 2);
+        assert_eq!(snapshot.total_record_count(), Some(100));
+    }
+
     #[tokio::test]
     async fn test_commit_writes_and_publishes_manifest_sidecar() {
         let file_io = test_file_io();
diff --git a/crates/paimon/src/table/table_scan.rs 
b/crates/paimon/src/table/table_scan.rs
index 3ccb3562..0decc94e 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -4073,7 +4073,7 @@ mod tests {
             .unwrap();
 
         let anchor = make_evo_file_with_cols("anchor.parquet", 1_000, 2, 0, 
&["id"]);
-        let mut replacement = CommitMessage::new(partition, 0, vec![anchor, 
deleted.clone()]);
+        let mut replacement = CommitMessage::new(partition, 0, vec![anchor]);
         replacement.deleted_files = vec![deleted];
         TableCommit::new(table.clone(), 
"delta-netting-replacement".to_string())
             .commit(vec![replacement])

Reply via email to