laskoviymishka commented on code in PR #3253:
URL: https://github.com/apache/iceberg-rust/pull/3253#discussion_r4153177113


##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -82,9 +83,24 @@ pub(crate) trait SnapshotProduceOperation: Send + Sync {
     /// - **Overwrite operations**: May exclude manifests for partitions being 
overwritten
     /// - **Delete operations**: May exclude manifests for partitions being 
deleted
     fn existing_manifest(
-        &self,
+        &mut self,
         snapshot_produce: &SnapshotProducer<'_>,
     ) -> impl Future<Output = Result<Vec<ManifestFile>>> + Send;
+
+    /// Returns the data files this operation actually removed, each with the 
schema and
+    /// partition spec of the manifest that recorded it.
+    ///
+    /// Only populated once [`Self::existing_manifest`] has run, so the 
snapshot summary counts
+    /// what was really removed rather than what the caller asked to remove.
+    fn removed_data_files(&self) -> &[(DataFile, SchemaRef, PartitionSpecRef)] 
{
+        &[]
+    }
+
+    /// Returns whether this operation replaces the whole table, dropping the 
previous totals
+    /// from the snapshot summary. A partial overwrite must return `false`.
+    fn truncate_full_table(&self) -> bool {

Review Comment:
   Flipping this from `operation() == Operation::Overwrite` to an opt-in flag 
is the right move for partial overwrites — this PR's own 
`RemoveExistingFileOperation` depends on the `false` default to report actual 
removed files instead of truncating. My worry is what it does to the ops that 
come after it: a full-table overwrite (replace-partitions, overwrite-by-filter 
with no predicate) returns `Overwrite` too, and if it forgets to override this 
it silently accumulates totals instead of resetting them — no error, and no 
repair path once the baseline is wrong.
   
   Since this is step 1 and the follow-ons are the exact callers at risk, I'd 
bind the invariant to something harder than a doc comment: at minimum a 
`debug_assert!`/warn in `summary()` when `operation() == Overwrite && 
!truncate_full_table()`, or better, tie truncation to a dedicated 
operation/path so the type carries it. I'd also add one e2e test that drives 
`truncate_full_table() == true` through `commit()` — right now only the 
internal `truncate_table_summary` fn is covered, so the wiring here is untested.



##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -82,9 +83,24 @@ pub(crate) trait SnapshotProduceOperation: Send + Sync {
     /// - **Overwrite operations**: May exclude manifests for partitions being 
overwritten
     /// - **Delete operations**: May exclude manifests for partitions being 
deleted
     fn existing_manifest(
-        &self,
+        &mut self,
         snapshot_produce: &SnapshotProducer<'_>,
     ) -> impl Future<Output = Result<Vec<ManifestFile>>> + Send;
+
+    /// Returns the data files this operation actually removed, each with the 
schema and
+    /// partition spec of the manifest that recorded it.
+    ///
+    /// Only populated once [`Self::existing_manifest`] has run, so the 
snapshot summary counts
+    /// what was really removed rather than what the caller asked to remove.
+    fn removed_data_files(&self) -> &[(DataFile, SchemaRef, PartitionSpecRef)] 
{

Review Comment:
   `&[(DataFile, SchemaRef, PartitionSpecRef)]` is a positional three-tuple 
whose roles you have to infer, and the COW series is going to stack callers on 
top of it. A small named struct (`RemovedFile { data_file, schema, 
partition_spec }`) costs nothing now and saves a coordinated touch at every 
call site the first time this needs a fourth field.



##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -403,6 +433,10 @@ impl<'a> SnapshotProducer<'a> {
             );
         }
 
+        for (data_file, schema, partition_spec) in 
snapshot_produce_operation.removed_data_files() {

Review Comment:
   This loop builds the summary's `deleted-*` counts purely from 
`removed_data_files()`, while the manifest's `Deleted` entries come from 
whatever `existing_manifest()` passed to `add_delete_entry()`. Nothing binds 
the two, so an operation that marks entries deleted in the manifest without 
reporting them here (or the reverse) commits manifests and a summary that 
disagree — and `total-data-files` drifts from there on every subsequent commit.
   
   `test_summary_counts_files_the_operation_reports_as_removed` is the case in 
point: `RemoveOneFileOperation` reports a removal but writes no `Deleted` 
entry, so the summary claims `deleted-data-files=1` over manifests with zero 
deleted entries, and it passes. I'd either sum `deleted_files_count` across the 
produced data manifests and `debug_assert!` it matches 
`removed_data_files().len()` in `produce_manifests`, or have 
`existing_manifest` return the removed set so the two can't diverge by 
construction.



##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -308,7 +328,9 @@ impl<'a> SnapshotProducer<'a> {
 
     // Write manifest file for added data files and return the ManifestFile 
for ManifestList.
     async fn write_added_manifest(&mut self) -> Result<ManifestFile> {
-        let added_data_files = std::mem::take(&mut self.added_data_files);
+        // Cloned rather than taken: the summary is built after the manifests, 
so the added
+        // files must still be here.
+        let added_data_files = self.added_data_files.clone();

Review Comment:
   The ordering fix is right, but cloning the whole `added_data_files` vec here 
is a real cost on the commit hot path — each `DataFile` carries several 
`HashMap`s, so a large append allocates a full duplicate where `mem::take` was 
free. Since the only reason to keep them is so `summary()` can read them 
afterward, I'd pass `&[DataFile]` into `write_added_manifest` (or iterate 
borrowing) instead of cloning the vec.



##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -234,22 +250,26 @@ impl<'a> SnapshotProducer<'a> {
         snapshot_id
     }
 
-    fn new_manifest_writer(&mut self, content: ManifestContentType) -> 
Result<ManifestWriter> {
-        let new_manifest_path = format!(
+    /// Returns the path for the next manifest file of this commit.
+    fn new_manifest_path(&self) -> Result<String> {
+        Ok(format!(
             "{}/{}-m{}.{}",
             self.table.metadata().metadata_location()?,
             self.commit_uuid,
-            self.manifest_counter.next().unwrap(),
+            self.manifest_counter.fetch_add(1, Ordering::Relaxed),
             DataFileFormat::Avro
-        );
-        let output_file = self.table.file_io().new_output(new_manifest_path)?;
-        let partition_spec = self
-            .table
-            .metadata()
-            .default_partition_spec()
-            .as_ref()
-            .clone();
-        let schema = self.table.metadata().current_schema().clone();
+        ))
+    }
+
+    /// Returns a writer for the next manifest file of this commit, routed 
through the table's
+    /// encryption manager when one is configured.
+    pub(crate) fn new_manifest_writer(
+        &self,
+        content: ManifestContentType,
+        schema: SchemaRef,
+        partition_spec: PartitionSpec,

Review Comment:
   `schema` comes in as `SchemaRef` but `partition_spec` is taken by value, so 
every caller holding an `Arc<PartitionSpec>` has to `.as_ref().clone()` a full 
spec to hand it over — `write_added_manifest` does exactly that just below, and 
`removed_data_files()` already standardizes on `PartitionSpecRef`. I'd take 
`partition_spec: PartitionSpecRef` here and clone once internally if the 
builder needs an owned value, so the call sites stop deep-cloning and the Arc 
convention stays consistent.



##########
crates/iceberg/src/spec/snapshot_summary.rs:
##########
@@ -532,7 +532,7 @@ fn update_totals(
         return;
     };
 
-    let new_total = previous_total + added - removed;
+    let new_total = (previous_total + added).saturating_sub(removed);

Review Comment:
   The `saturating_sub` fixes the subtraction underflow, but `previous_total + 
added` on the left is still plain `u64` addition — a corrupt or adversarial 
previous summary with `total-data-files` near `u64::MAX` panics in debug and 
wraps to a small number in release. `saturating_add` closes it: `let new_total 
= previous_total.saturating_add(added).saturating_sub(removed);`.
   
   Separately, saturating to `0` diverges from Java: 
`SnapshotProducer.updateTotal` *omits* the `total-*` property entirely when the 
value would go negative, and skips recomputation when the key is absent — so 
writing `0` makes a Java reader treat it as a real basis and carry the 
undercount forward. I'd match Java and omit the property on underflow, and flip 
`test_update_totals_saturates_...` to assert the key is absent rather than 
`"0"`.



##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -116,7 +132,7 @@ pub(crate) struct SnapshotProducer<'a> {
     // A counter used to generate unique manifest file names.
     // It starts from 0 and increments for each new manifest file.
     // Note: This counter is limited to the range of (0..u64::MAX).
-    manifest_counter: RangeFrom<u64>,
+    manifest_counter: AtomicU64,

Review Comment:
   The carried-over comment just above says the counter is "limited to the 
range of (0..u64::MAX)", but `fetch_add` wraps at `u64::MAX` rather than 
stopping — so the note now reads as a guarantee the type doesn't give (a wrap 
hands out a duplicate manifest path). While we're here, a word on why it's 
`AtomicU64` would help: it's interior mutability so `new_manifest_writer` can 
take `&self`, not cross-thread sharing — which is why `Relaxed` is fine.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to