Skip to main content

iceberg/spec/
snapshot_summary.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;
19
20use itertools::Itertools;
21
22use super::{DataContentType, DataFile, PartitionSpecRef};
23use crate::error::invalid_data;
24use crate::spec::{ManifestContentType, ManifestFile, Operation, SchemaRef, Summary};
25use crate::{Error, ErrorKind, Result};
26
27const ADDED_DATA_FILES: &str = "added-data-files";
28const ADDED_DELETE_FILES: &str = "added-delete-files";
29const ADDED_EQUALITY_DELETES: &str = "added-equality-deletes";
30const ADDED_FILE_SIZE: &str = "added-files-size";
31const ADDED_POSITION_DELETES: &str = "added-position-deletes";
32const ADDED_POSITION_DELETE_FILES: &str = "added-position-delete-files";
33const ADDED_RECORDS: &str = "added-records";
34const DELETED_DATA_FILES: &str = "deleted-data-files";
35const DELETED_RECORDS: &str = "deleted-records";
36const ADDED_EQUALITY_DELETE_FILES: &str = "added-equality-delete-files";
37const REMOVED_DELETE_FILES: &str = "removed-delete-files";
38const REMOVED_EQUALITY_DELETES: &str = "removed-equality-deletes";
39const REMOVED_EQUALITY_DELETE_FILES: &str = "removed-equality-delete-files";
40const REMOVED_FILE_SIZE: &str = "removed-files-size";
41const REMOVED_POSITION_DELETES: &str = "removed-position-deletes";
42const REMOVED_POSITION_DELETE_FILES: &str = "removed-position-delete-files";
43const TOTAL_EQUALITY_DELETES: &str = "total-equality-deletes";
44const TOTAL_POSITION_DELETES: &str = "total-position-deletes";
45const TOTAL_DATA_FILES: &str = "total-data-files";
46const TOTAL_DELETE_FILES: &str = "total-delete-files";
47const TOTAL_RECORDS: &str = "total-records";
48const TOTAL_FILE_SIZE: &str = "total-files-size";
49const CHANGED_PARTITION_COUNT_PROP: &str = "changed-partition-count";
50const CHANGED_PARTITION_PREFIX: &str = "partitions.";
51
52/// `SnapshotSummaryCollector` collects and aggregates snapshot update metrics.
53/// It gathers metrics about added or removed data files and manifests, and tracks
54/// partition-specific updates.
55#[derive(Default)]
56pub struct SnapshotSummaryCollector {
57    metrics: UpdateMetrics,
58    partition_metrics: HashMap<String, UpdateMetrics>,
59    max_changed_partitions_for_summaries: u64,
60    properties: HashMap<String, String>,
61    trust_partition_metrics: bool,
62}
63
64impl SnapshotSummaryCollector {
65    /// Set properties for snapshot summary
66    pub fn set(&mut self, key: &str, value: &str) {
67        self.properties.insert(key.to_string(), value.to_string());
68    }
69
70    /// Sets the limit for including partition summaries. Summaries are not
71    /// included if the number of partitions is exceeded.
72    pub fn set_partition_summary_limit(&mut self, limit: u64) {
73        self.max_changed_partitions_for_summaries = limit;
74    }
75
76    /// Adds a data file to the summary collector
77    pub fn add_file(
78        &mut self,
79        data_file: &DataFile,
80        schema: SchemaRef,
81        partition_spec: PartitionSpecRef,
82    ) {
83        self.metrics.add_file(data_file);
84        if !data_file.partition.fields().is_empty() {
85            self.update_partition_metrics(schema, partition_spec, data_file, true);
86        }
87    }
88
89    /// Removes a data file from the summary collector
90    pub fn remove_file(
91        &mut self,
92        data_file: &DataFile,
93        schema: SchemaRef,
94        partition_spec: PartitionSpecRef,
95    ) {
96        self.metrics.remove_file(data_file);
97        if !data_file.partition.fields().is_empty() {
98            self.update_partition_metrics(schema, partition_spec, data_file, false);
99        }
100    }
101
102    /// Adds a manifest to the summary collector
103    pub fn add_manifest(&mut self, manifest: &ManifestFile) {
104        self.trust_partition_metrics = false;
105        self.partition_metrics.clear();
106        self.metrics.add_manifest(manifest);
107    }
108
109    /// Updates partition-specific metrics for a data file.
110    pub fn update_partition_metrics(
111        &mut self,
112        schema: SchemaRef,
113        partition_spec: PartitionSpecRef,
114        data_file: &DataFile,
115        is_add_file: bool,
116    ) {
117        let partition_path = partition_spec.partition_to_path(&data_file.partition, schema);
118        let metrics = self.partition_metrics.entry(partition_path).or_default();
119
120        if is_add_file {
121            metrics.add_file(data_file);
122        } else {
123            metrics.remove_file(data_file);
124        }
125    }
126
127    /// Merges another `SnapshotSummaryCollector` into the current one
128    pub fn merge(&mut self, summary: SnapshotSummaryCollector) {
129        self.metrics.merge(&summary.metrics);
130        self.properties.extend(summary.properties);
131
132        if self.trust_partition_metrics && summary.trust_partition_metrics {
133            for (partition, partition_metric) in summary.partition_metrics.iter() {
134                self.partition_metrics
135                    .entry(partition.to_string())
136                    .or_default()
137                    .merge(partition_metric);
138            }
139        } else {
140            self.partition_metrics.clear();
141            self.trust_partition_metrics = false;
142        }
143    }
144
145    /// Builds final map of summaries
146    pub fn build(&self) -> HashMap<String, String> {
147        let mut properties = self.metrics.to_map();
148        let changed_partitions_count = self.partition_metrics.len() as u64;
149        set_if_positive(
150            &mut properties,
151            changed_partitions_count,
152            CHANGED_PARTITION_COUNT_PROP,
153        );
154
155        if changed_partitions_count <= self.max_changed_partitions_for_summaries {
156            for (partition_path, update_metrics_partition) in &self.partition_metrics {
157                let property_key = format!("{CHANGED_PARTITION_PREFIX}{partition_path}");
158                let partition_summary = update_metrics_partition
159                    .to_map()
160                    .into_iter()
161                    .map(|(property, value)| format!("{property}={value}"))
162                    .join(",");
163
164                if !partition_summary.is_empty() {
165                    properties.insert(property_key, partition_summary);
166                }
167            }
168        }
169        properties
170    }
171}
172
173#[derive(Debug, Default)]
174struct UpdateMetrics {
175    added_file_size: u64,
176    removed_file_size: u64,
177    added_data_files: u32,
178    removed_data_files: u32,
179    added_eq_delete_files: u64,
180    removed_eq_delete_files: u64,
181    added_pos_delete_files: u64,
182    removed_pos_delete_files: u64,
183    added_delete_files: u32,
184    removed_delete_files: u32,
185    added_records: u64,
186    deleted_records: u64,
187    added_pos_deletes: u64,
188    removed_pos_deletes: u64,
189    added_eq_deletes: u64,
190    removed_eq_deletes: u64,
191}
192
193impl UpdateMetrics {
194    fn add_file(&mut self, data_file: &DataFile) {
195        self.added_file_size += data_file.file_size_in_bytes;
196        match data_file.content_type() {
197            DataContentType::Data => {
198                self.added_data_files += 1;
199                self.added_records += data_file.record_count;
200            }
201            DataContentType::PositionDeletes => {
202                self.added_delete_files += 1;
203                self.added_pos_delete_files += 1;
204                self.added_pos_deletes += data_file.record_count;
205            }
206            DataContentType::EqualityDeletes => {
207                self.added_delete_files += 1;
208                self.added_eq_delete_files += 1;
209                self.added_eq_deletes += data_file.record_count;
210            }
211        }
212    }
213
214    fn remove_file(&mut self, data_file: &DataFile) {
215        self.removed_file_size += data_file.file_size_in_bytes;
216        match data_file.content_type() {
217            DataContentType::Data => {
218                self.removed_data_files += 1;
219                self.deleted_records += data_file.record_count;
220            }
221            DataContentType::PositionDeletes => {
222                self.removed_delete_files += 1;
223                self.removed_pos_delete_files += 1;
224                self.removed_pos_deletes += data_file.record_count;
225            }
226            DataContentType::EqualityDeletes => {
227                self.removed_delete_files += 1;
228                self.removed_eq_delete_files += 1;
229                self.removed_eq_deletes += data_file.record_count;
230            }
231        }
232    }
233
234    fn add_manifest(&mut self, manifest: &ManifestFile) {
235        match manifest.content {
236            ManifestContentType::Data => {
237                self.added_data_files += manifest.added_files_count.unwrap_or(0);
238                self.added_records += manifest.added_rows_count.unwrap_or(0);
239                self.removed_data_files += manifest.deleted_files_count.unwrap_or(0);
240                self.deleted_records += manifest.deleted_rows_count.unwrap_or(0);
241            }
242            ManifestContentType::Deletes => {
243                self.added_delete_files += manifest.added_files_count.unwrap_or(0);
244                self.removed_delete_files += manifest.deleted_files_count.unwrap_or(0);
245            }
246        }
247    }
248
249    fn to_map(&self) -> HashMap<String, String> {
250        let mut properties = HashMap::new();
251        set_if_positive(&mut properties, self.added_file_size, ADDED_FILE_SIZE);
252        set_if_positive(&mut properties, self.removed_file_size, REMOVED_FILE_SIZE);
253        set_if_positive(&mut properties, self.added_data_files, ADDED_DATA_FILES);
254        set_if_positive(&mut properties, self.removed_data_files, DELETED_DATA_FILES);
255        set_if_positive(
256            &mut properties,
257            self.added_eq_delete_files,
258            ADDED_EQUALITY_DELETE_FILES,
259        );
260        set_if_positive(
261            &mut properties,
262            self.removed_eq_delete_files,
263            REMOVED_EQUALITY_DELETE_FILES,
264        );
265        set_if_positive(
266            &mut properties,
267            self.added_pos_delete_files,
268            ADDED_POSITION_DELETE_FILES,
269        );
270        set_if_positive(
271            &mut properties,
272            self.removed_pos_delete_files,
273            REMOVED_POSITION_DELETE_FILES,
274        );
275        set_if_positive(&mut properties, self.added_delete_files, ADDED_DELETE_FILES);
276        set_if_positive(
277            &mut properties,
278            self.removed_delete_files,
279            REMOVED_DELETE_FILES,
280        );
281        set_if_positive(&mut properties, self.added_records, ADDED_RECORDS);
282        set_if_positive(&mut properties, self.deleted_records, DELETED_RECORDS);
283        set_if_positive(
284            &mut properties,
285            self.added_pos_deletes,
286            ADDED_POSITION_DELETES,
287        );
288        set_if_positive(
289            &mut properties,
290            self.removed_pos_deletes,
291            REMOVED_POSITION_DELETES,
292        );
293        set_if_positive(
294            &mut properties,
295            self.added_eq_deletes,
296            ADDED_EQUALITY_DELETES,
297        );
298        set_if_positive(
299            &mut properties,
300            self.removed_eq_deletes,
301            REMOVED_EQUALITY_DELETES,
302        );
303        properties
304    }
305
306    fn merge(&mut self, other: &UpdateMetrics) {
307        self.added_file_size += other.added_file_size;
308        self.removed_file_size += other.removed_file_size;
309        self.added_data_files += other.added_data_files;
310        self.removed_data_files += other.removed_data_files;
311        self.added_eq_delete_files += other.added_eq_delete_files;
312        self.removed_eq_delete_files += other.removed_eq_delete_files;
313        self.added_pos_delete_files += other.added_pos_delete_files;
314        self.removed_pos_delete_files += other.removed_pos_delete_files;
315        self.added_delete_files += other.added_delete_files;
316        self.removed_delete_files += other.removed_delete_files;
317        self.added_records += other.added_records;
318        self.deleted_records += other.deleted_records;
319        self.added_pos_deletes += other.added_pos_deletes;
320        self.removed_pos_deletes += other.removed_pos_deletes;
321        self.added_eq_deletes += other.added_eq_deletes;
322        self.removed_eq_deletes += other.removed_eq_deletes;
323    }
324}
325
326fn set_if_positive<T>(properties: &mut HashMap<String, String>, value: T, property_name: &str)
327where T: PartialOrd + Default + ToString {
328    if value > T::default() {
329        properties.insert(property_name.to_string(), value.to_string());
330    }
331}
332
333pub(crate) fn update_snapshot_summaries(
334    summary: Summary,
335    previous_summary: Option<&Summary>,
336    truncate_full_table: bool,
337) -> Result<Summary> {
338    // Validate that the operation is supported
339    if summary.operation != Operation::Append
340        && summary.operation != Operation::Overwrite
341        && summary.operation != Operation::Delete
342    {
343        return Err(invalid_data!("Operation is not supported."));
344    }
345
346    let mut summary = match previous_summary {
347        Some(prev_summary) if truncate_full_table && summary.operation == Operation::Overwrite => {
348            truncate_table_summary(summary, prev_summary).map_err(|err| {
349                Error::new(ErrorKind::Unexpected, "Failed to truncate table summary.")
350                    .with_source(err)
351            })?
352        }
353        _ => summary,
354    };
355
356    update_totals(
357        &mut summary,
358        previous_summary,
359        TOTAL_DATA_FILES,
360        ADDED_DATA_FILES,
361        DELETED_DATA_FILES,
362    );
363
364    update_totals(
365        &mut summary,
366        previous_summary,
367        TOTAL_DELETE_FILES,
368        ADDED_DELETE_FILES,
369        REMOVED_DELETE_FILES,
370    );
371
372    update_totals(
373        &mut summary,
374        previous_summary,
375        TOTAL_RECORDS,
376        ADDED_RECORDS,
377        DELETED_RECORDS,
378    );
379
380    update_totals(
381        &mut summary,
382        previous_summary,
383        TOTAL_FILE_SIZE,
384        ADDED_FILE_SIZE,
385        REMOVED_FILE_SIZE,
386    );
387
388    update_totals(
389        &mut summary,
390        previous_summary,
391        TOTAL_POSITION_DELETES,
392        ADDED_POSITION_DELETES,
393        REMOVED_POSITION_DELETES,
394    );
395
396    update_totals(
397        &mut summary,
398        previous_summary,
399        TOTAL_EQUALITY_DELETES,
400        ADDED_EQUALITY_DELETES,
401        REMOVED_EQUALITY_DELETES,
402    );
403    Ok(summary)
404}
405
406fn get_prop(previous_summary: &Summary, prop: &str) -> Result<u64> {
407    let value_str = previous_summary
408        .additional_properties
409        .get(prop)
410        .map(String::as_str)
411        .unwrap_or("0");
412    value_str.parse::<u64>().map_err(|err| {
413        Error::new(
414            ErrorKind::Unexpected,
415            format!("Failed to parse summary property '{prop}' value '{value_str}' as u64."),
416        )
417        .with_source(err)
418    })
419}
420
421fn truncate_table_summary(mut summary: Summary, previous_summary: &Summary) -> Result<Summary> {
422    for prop in [
423        TOTAL_DATA_FILES,
424        TOTAL_DELETE_FILES,
425        TOTAL_RECORDS,
426        TOTAL_FILE_SIZE,
427        TOTAL_POSITION_DELETES,
428        TOTAL_EQUALITY_DELETES,
429    ] {
430        summary
431            .additional_properties
432            .insert(prop.to_string(), "0".to_string());
433    }
434
435    let value = get_prop(previous_summary, TOTAL_DATA_FILES)?;
436    if value != 0 {
437        summary
438            .additional_properties
439            .insert(DELETED_DATA_FILES.to_string(), value.to_string());
440    }
441    let value = get_prop(previous_summary, TOTAL_DELETE_FILES)?;
442    if value != 0 {
443        summary
444            .additional_properties
445            .insert(REMOVED_DELETE_FILES.to_string(), value.to_string());
446    }
447    let value = get_prop(previous_summary, TOTAL_RECORDS)?;
448    if value != 0 {
449        summary
450            .additional_properties
451            .insert(DELETED_RECORDS.to_string(), value.to_string());
452    }
453    let value = get_prop(previous_summary, TOTAL_FILE_SIZE)?;
454    if value != 0 {
455        summary
456            .additional_properties
457            .insert(REMOVED_FILE_SIZE.to_string(), value.to_string());
458    }
459
460    let value = get_prop(previous_summary, TOTAL_POSITION_DELETES)?;
461    if value != 0 {
462        summary
463            .additional_properties
464            .insert(REMOVED_POSITION_DELETES.to_string(), value.to_string());
465    }
466
467    let value = get_prop(previous_summary, TOTAL_EQUALITY_DELETES)?;
468    if value != 0 {
469        summary
470            .additional_properties
471            .insert(REMOVED_EQUALITY_DELETES.to_string(), value.to_string());
472    }
473
474    Ok(summary)
475}
476
477fn update_totals(
478    summary: &mut Summary,
479    previous_summary: Option<&Summary>,
480    total_property: &str,
481    added_property: &str,
482    removed_property: &str,
483) {
484    let previous_total = match previous_summary {
485        None => 0,
486        Some(prev_summary) => match prev_summary.additional_properties.get(total_property) {
487            Some(value_str) => match value_str.parse::<u64>() {
488                Ok(v) => v,
489                Err(parse_err) => {
490                    tracing::warn!(
491                        "Property '{total_property}' could not be parsed in the previous snapshot summary: {parse_err}. \
492                         Skipping total computation.",
493                    );
494                    return;
495                }
496            },
497            None => {
498                tracing::debug!(
499                    "Property '{total_property}' was not set in the previous snapshot summary. \
500                     Skipping total computation."
501                );
502                return;
503            }
504        },
505    };
506
507    // Parse the added/removed deltas, tolerating an unparsable value by skipping
508    // the total entirely rather than panicking. Computed metrics always overwrite
509    // user-supplied summary properties (see `SnapshotProducer::summary`), so a bad
510    // value should only ever come from a previous snapshot's summary; matching
511    // iceberg-java's `updateTotal`, we ignore it instead of failing the commit.
512    let parse_delta = |property: &str| -> Option<u64> {
513        match summary.additional_properties.get(property) {
514            None => Some(0),
515            Some(value) => match value.parse::<u64>() {
516                Ok(v) => Some(v),
517                Err(parse_err) => {
518                    tracing::warn!(
519                        "Property '{property}' could not be parsed when computing '{total_property}': {parse_err}. \
520                         Skipping total computation.",
521                    );
522                    None
523                }
524            },
525        }
526    };
527
528    let (Some(added), Some(removed)) = (parse_delta(added_property), parse_delta(removed_property))
529    else {
530        return;
531    };
532
533    let new_total = previous_total + added - removed;
534    summary
535        .additional_properties
536        .insert(total_property.to_string(), new_total.to_string());
537}
538
539#[cfg(test)]
540mod tests {
541    use std::collections::HashMap;
542    use std::sync::Arc;
543
544    use super::*;
545    use crate::spec::{
546        DataFileFormat, Datum, Literal, NestedField, PartitionSpec, PrimitiveType, Schema, Struct,
547        Transform, Type, UnboundPartitionField,
548    };
549
550    #[test]
551    fn test_update_snapshot_summaries_append() {
552        let prev_props: HashMap<String, String> = [
553            (TOTAL_DATA_FILES.to_string(), "10".to_string()),
554            (TOTAL_DELETE_FILES.to_string(), "5".to_string()),
555            (TOTAL_RECORDS.to_string(), "100".to_string()),
556            (TOTAL_FILE_SIZE.to_string(), "1000".to_string()),
557            (TOTAL_POSITION_DELETES.to_string(), "3".to_string()),
558            (TOTAL_EQUALITY_DELETES.to_string(), "2".to_string()),
559        ]
560        .into_iter()
561        .collect();
562
563        let previous_summary = Summary {
564            operation: Operation::Append,
565            additional_properties: prev_props,
566        };
567
568        let new_props: HashMap<String, String> = [
569            (ADDED_DATA_FILES.to_string(), "4".to_string()),
570            (DELETED_DATA_FILES.to_string(), "1".to_string()),
571            (ADDED_DELETE_FILES.to_string(), "2".to_string()),
572            (REMOVED_DELETE_FILES.to_string(), "1".to_string()),
573            (ADDED_RECORDS.to_string(), "40".to_string()),
574            (DELETED_RECORDS.to_string(), "10".to_string()),
575            (ADDED_FILE_SIZE.to_string(), "400".to_string()),
576            (REMOVED_FILE_SIZE.to_string(), "100".to_string()),
577            (ADDED_POSITION_DELETES.to_string(), "5".to_string()),
578            (REMOVED_POSITION_DELETES.to_string(), "2".to_string()),
579            (ADDED_EQUALITY_DELETES.to_string(), "3".to_string()),
580            (REMOVED_EQUALITY_DELETES.to_string(), "1".to_string()),
581        ]
582        .into_iter()
583        .collect();
584
585        let summary = Summary {
586            operation: Operation::Append,
587            additional_properties: new_props,
588        };
589
590        let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
591
592        assert_eq!(
593            updated.additional_properties.get(TOTAL_DATA_FILES).unwrap(),
594            "13"
595        );
596        assert_eq!(
597            updated
598                .additional_properties
599                .get(TOTAL_DELETE_FILES)
600                .unwrap(),
601            "6"
602        );
603        assert_eq!(
604            updated.additional_properties.get(TOTAL_RECORDS).unwrap(),
605            "130"
606        );
607        assert_eq!(
608            updated.additional_properties.get(TOTAL_FILE_SIZE).unwrap(),
609            "1300"
610        );
611        assert_eq!(
612            updated
613                .additional_properties
614                .get(TOTAL_POSITION_DELETES)
615                .unwrap(),
616            "6"
617        );
618        assert_eq!(
619            updated
620                .additional_properties
621                .get(TOTAL_EQUALITY_DELETES)
622                .unwrap(),
623            "4"
624        );
625    }
626
627    #[test]
628    fn test_truncate_table_summary() {
629        let prev_props: HashMap<String, String> = [
630            (TOTAL_DATA_FILES.to_string(), "10".to_string()),
631            (TOTAL_DELETE_FILES.to_string(), "5".to_string()),
632            (TOTAL_RECORDS.to_string(), "100".to_string()),
633            (TOTAL_FILE_SIZE.to_string(), "1000".to_string()),
634            (TOTAL_POSITION_DELETES.to_string(), "3".to_string()),
635            (TOTAL_EQUALITY_DELETES.to_string(), "2".to_string()),
636        ]
637        .into_iter()
638        .collect();
639
640        let previous_summary = Summary {
641            operation: Operation::Overwrite,
642            additional_properties: prev_props,
643        };
644
645        let mut new_props = HashMap::new();
646        new_props.insert("dummy".to_string(), "value".to_string());
647        let summary = Summary {
648            operation: Operation::Overwrite,
649            additional_properties: new_props,
650        };
651
652        let truncated = truncate_table_summary(summary, &previous_summary).unwrap();
653
654        assert_eq!(
655            truncated
656                .additional_properties
657                .get(TOTAL_DATA_FILES)
658                .unwrap(),
659            "0"
660        );
661        assert_eq!(
662            truncated
663                .additional_properties
664                .get(TOTAL_DELETE_FILES)
665                .unwrap(),
666            "0"
667        );
668        assert_eq!(
669            truncated.additional_properties.get(TOTAL_RECORDS).unwrap(),
670            "0"
671        );
672        assert_eq!(
673            truncated
674                .additional_properties
675                .get(TOTAL_FILE_SIZE)
676                .unwrap(),
677            "0"
678        );
679        assert_eq!(
680            truncated
681                .additional_properties
682                .get(TOTAL_POSITION_DELETES)
683                .unwrap(),
684            "0"
685        );
686        assert_eq!(
687            truncated
688                .additional_properties
689                .get(TOTAL_EQUALITY_DELETES)
690                .unwrap(),
691            "0"
692        );
693
694        assert_eq!(
695            truncated
696                .additional_properties
697                .get(DELETED_DATA_FILES)
698                .unwrap(),
699            "10"
700        );
701        assert_eq!(
702            truncated
703                .additional_properties
704                .get(REMOVED_DELETE_FILES)
705                .unwrap(),
706            "5"
707        );
708        assert_eq!(
709            truncated
710                .additional_properties
711                .get(DELETED_RECORDS)
712                .unwrap(),
713            "100"
714        );
715        assert_eq!(
716            truncated
717                .additional_properties
718                .get(REMOVED_FILE_SIZE)
719                .unwrap(),
720            "1000"
721        );
722        assert_eq!(
723            truncated
724                .additional_properties
725                .get(REMOVED_POSITION_DELETES)
726                .unwrap(),
727            "3"
728        );
729        assert_eq!(
730            truncated
731                .additional_properties
732                .get(REMOVED_EQUALITY_DELETES)
733                .unwrap(),
734            "2"
735        );
736    }
737
738    #[test]
739    fn test_update_snapshot_summaries_overwrite_truncate_handles_totals_above_i32_max() {
740        // A table can legitimately accumulate more than i32::MAX rows or files
741        // over its lifetime. Truncating such a table on overwrite must succeed
742        // and surface the previous totals into the deleted-* counters.
743        let big = (i32::MAX as u64 + 1).to_string(); // 2_147_483_648
744        let prev_props: HashMap<String, String> = [
745            (TOTAL_DATA_FILES.to_string(), big.clone()),
746            (TOTAL_DELETE_FILES.to_string(), "0".to_string()),
747            (TOTAL_RECORDS.to_string(), big.clone()),
748            (TOTAL_FILE_SIZE.to_string(), "0".to_string()),
749            (TOTAL_POSITION_DELETES.to_string(), "0".to_string()),
750            (TOTAL_EQUALITY_DELETES.to_string(), "0".to_string()),
751        ]
752        .into_iter()
753        .collect();
754
755        let previous_summary = Summary {
756            operation: Operation::Overwrite,
757            additional_properties: prev_props,
758        };
759
760        let summary = Summary {
761            operation: Operation::Overwrite,
762            additional_properties: HashMap::new(),
763        };
764
765        let updated = update_snapshot_summaries(summary, Some(&previous_summary), true)
766            .expect("overwrite truncation should accept totals above i32::MAX");
767        assert_eq!(
768            updated
769                .additional_properties
770                .get(DELETED_DATA_FILES)
771                .unwrap(),
772            &big
773        );
774        assert_eq!(
775            updated.additional_properties.get(DELETED_RECORDS).unwrap(),
776            &big
777        );
778    }
779
780    #[test]
781    fn test_update_snapshot_summaries_overwrite_truncate_returns_err_on_malformed_total() {
782        // Non-numeric values in the previous summary (corruption, manual edits,
783        // a foreign implementation) must surface as a recoverable Err - not
784        // crash the process.
785        let prev_props: HashMap<String, String> = [
786            (TOTAL_DATA_FILES.to_string(), "not_a_number".to_string()),
787            (TOTAL_DELETE_FILES.to_string(), "0".to_string()),
788            (TOTAL_RECORDS.to_string(), "0".to_string()),
789            (TOTAL_FILE_SIZE.to_string(), "0".to_string()),
790            (TOTAL_POSITION_DELETES.to_string(), "0".to_string()),
791            (TOTAL_EQUALITY_DELETES.to_string(), "0".to_string()),
792        ]
793        .into_iter()
794        .collect();
795
796        let previous_summary = Summary {
797            operation: Operation::Overwrite,
798            additional_properties: prev_props,
799        };
800
801        let summary = Summary {
802            operation: Operation::Overwrite,
803            additional_properties: HashMap::new(),
804        };
805
806        let err = update_snapshot_summaries(summary, Some(&previous_summary), true)
807            .expect_err("malformed previous summary must produce an Err, not a panic");
808        assert!(
809            err.message().contains("truncate table summary"),
810            "expected wrapped 'Failed to truncate table summary' context, got: {}",
811            err.message()
812        );
813    }
814
815    #[test]
816    fn test_snapshot_summary_collector_build() {
817        let schema = Arc::new(
818            Schema::builder()
819                .with_fields(vec![
820                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
821                    NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
822                ])
823                .build()
824                .unwrap(),
825        );
826
827        let partition_spec = Arc::new(
828            PartitionSpec::builder(schema.clone())
829                .add_unbound_fields(vec![
830                    UnboundPartitionField::builder()
831                        .source_ids(vec![2])
832                        .name("year".to_string())
833                        .transform(Transform::Identity)
834                        .build()
835                        .unwrap(),
836                ])
837                .unwrap()
838                .with_spec_id(1)
839                .build()
840                .unwrap(),
841        );
842
843        let mut collector = SnapshotSummaryCollector::default();
844        collector.set_partition_summary_limit(10);
845
846        let file1 = DataFile {
847            content: DataContentType::Data,
848            file_path: "s3://testbucket/path/to/file1.parquet".to_string(),
849            file_format: DataFileFormat::Parquet,
850            partition: Struct::from_iter(vec![]),
851            record_count: 10,
852            file_size_in_bytes: 100,
853            column_sizes: HashMap::from([(1, 46), (2, 48), (3, 48)]),
854            value_counts: HashMap::from([(1, 10), (2, 10), (3, 10)]),
855            null_value_counts: HashMap::from([(1, 0), (2, 0), (3, 0)]),
856            nan_value_counts: HashMap::new(),
857            lower_bounds: HashMap::from([
858                (1, Datum::long(1)),
859                (2, Datum::string("a")),
860                (3, Datum::string("x")),
861            ]),
862            upper_bounds: HashMap::from([
863                (1, Datum::long(1)),
864                (2, Datum::string("a")),
865                (3, Datum::string("x")),
866            ]),
867            key_metadata: None,
868            split_offsets: Some(vec![4]),
869            equality_ids: None,
870            sort_order_id: Some(0),
871            partition_spec_id: 0,
872            first_row_id: None,
873            referenced_data_file: None,
874            content_offset: None,
875            content_size_in_bytes: None,
876        };
877
878        let file2 = DataFile {
879            content: DataContentType::Data,
880            file_path: "s3://testbucket/path/to/file2.parquet".to_string(),
881            file_format: DataFileFormat::Parquet,
882            partition: Struct::from_iter(vec![Some(Literal::string("2025"))]),
883            record_count: 20,
884            file_size_in_bytes: 200,
885            column_sizes: HashMap::from([(1, 46), (2, 48), (3, 48)]),
886            value_counts: HashMap::from([(1, 20), (2, 20), (3, 20)]),
887            null_value_counts: HashMap::from([(1, 0), (2, 0), (3, 0)]),
888            nan_value_counts: HashMap::new(),
889            lower_bounds: HashMap::from([
890                (1, Datum::long(1)),
891                (2, Datum::string("a")),
892                (3, Datum::string("x")),
893            ]),
894            upper_bounds: HashMap::from([
895                (1, Datum::long(1)),
896                (2, Datum::string("a")),
897                (3, Datum::string("x")),
898            ]),
899            key_metadata: None,
900            split_offsets: Some(vec![4]),
901            equality_ids: None,
902            sort_order_id: Some(0),
903            partition_spec_id: 0,
904            first_row_id: None,
905            referenced_data_file: None,
906            content_offset: None,
907            content_size_in_bytes: None,
908        };
909
910        collector.add_file(&file1, schema.clone(), partition_spec.clone());
911        collector.add_file(&file2, schema.clone(), partition_spec.clone());
912
913        collector.remove_file(&file1, schema.clone(), partition_spec.clone());
914
915        let props = collector.build();
916
917        assert_eq!(props.get(ADDED_FILE_SIZE).unwrap(), "300");
918        assert_eq!(props.get(REMOVED_FILE_SIZE).unwrap(), "100");
919
920        let partition_key = format!("{}{}", CHANGED_PARTITION_PREFIX, "year=2025");
921
922        assert!(props.contains_key(&partition_key));
923
924        let partition_summary = props.get(&partition_key).unwrap();
925        assert!(partition_summary.contains(&format!("{ADDED_FILE_SIZE}=200")));
926        assert!(partition_summary.contains(&format!("{ADDED_DATA_FILES}=1")));
927        assert!(partition_summary.contains(&format!("{ADDED_RECORDS}=20")));
928    }
929
930    #[test]
931    fn test_snapshot_summary_collector_add_manifest() {
932        let mut collector = SnapshotSummaryCollector::default();
933        collector.set_partition_summary_limit(10);
934
935        let manifest = ManifestFile {
936            manifest_path: "file://dummy.manifest".to_string(),
937            manifest_length: 0,
938            partition_spec_id: 0,
939            content: ManifestContentType::Data,
940            sequence_number: 0,
941            min_sequence_number: 0,
942            added_snapshot_id: 0,
943            added_files_count: Some(3),
944            existing_files_count: Some(0),
945            deleted_files_count: Some(1),
946            added_rows_count: Some(100),
947            existing_rows_count: Some(0),
948            deleted_rows_count: Some(50),
949            partitions: Some(Vec::new()),
950            key_metadata: None,
951            first_row_id: None,
952        };
953
954        collector
955            .partition_metrics
956            .insert("dummy".to_string(), UpdateMetrics::default());
957        collector.add_manifest(&manifest);
958
959        let props = collector.build();
960        assert_eq!(props.get(ADDED_DATA_FILES).unwrap(), "3");
961        assert_eq!(props.get(DELETED_DATA_FILES).unwrap(), "1");
962        assert_eq!(props.get(ADDED_RECORDS).unwrap(), "100");
963        assert_eq!(props.get(DELETED_RECORDS).unwrap(), "50");
964    }
965
966    #[test]
967    fn test_snapshot_summary_collector_merge() {
968        let schema = Arc::new(
969            Schema::builder()
970                .with_fields(vec![
971                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
972                    NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
973                ])
974                .build()
975                .unwrap(),
976        );
977
978        let partition_spec = Arc::new(
979            PartitionSpec::builder(schema.clone())
980                .add_unbound_fields(vec![
981                    UnboundPartitionField::builder()
982                        .source_ids(vec![2])
983                        .name("year".to_string())
984                        .transform(Transform::Identity)
985                        .build()
986                        .unwrap(),
987                ])
988                .unwrap()
989                .with_spec_id(1)
990                .build()
991                .unwrap(),
992        );
993
994        let mut summary_one = SnapshotSummaryCollector::default();
995        let mut summary_two = SnapshotSummaryCollector::default();
996
997        summary_one.add_file(
998            &DataFile {
999                content: DataContentType::Data,
1000                file_path: "test.parquet".into(),
1001                file_format: DataFileFormat::Parquet,
1002                partition: Struct::from_iter(vec![]),
1003                record_count: 10,
1004                file_size_in_bytes: 100,
1005                column_sizes: HashMap::new(),
1006                value_counts: HashMap::new(),
1007                null_value_counts: HashMap::new(),
1008                nan_value_counts: HashMap::new(),
1009                lower_bounds: HashMap::new(),
1010                upper_bounds: HashMap::new(),
1011                key_metadata: None,
1012                split_offsets: None,
1013                equality_ids: None,
1014                sort_order_id: None,
1015                partition_spec_id: 0,
1016                first_row_id: None,
1017                referenced_data_file: None,
1018                content_offset: None,
1019                content_size_in_bytes: None,
1020            },
1021            schema.clone(),
1022            partition_spec.clone(),
1023        );
1024
1025        summary_two.add_file(
1026            &DataFile {
1027                content: DataContentType::Data,
1028                file_path: "test.parquet".into(),
1029                file_format: DataFileFormat::Parquet,
1030                partition: Struct::from_iter(vec![]),
1031                record_count: 20,
1032                file_size_in_bytes: 200,
1033                column_sizes: HashMap::new(),
1034                value_counts: HashMap::new(),
1035                null_value_counts: HashMap::new(),
1036                nan_value_counts: HashMap::new(),
1037                lower_bounds: HashMap::new(),
1038                upper_bounds: HashMap::new(),
1039                key_metadata: None,
1040                split_offsets: None,
1041                equality_ids: None,
1042                sort_order_id: None,
1043                partition_spec_id: 0,
1044                first_row_id: None,
1045                referenced_data_file: None,
1046                content_offset: None,
1047                content_size_in_bytes: None,
1048            },
1049            schema.clone(),
1050            partition_spec.clone(),
1051        );
1052
1053        summary_one.merge(summary_two);
1054        let props = summary_one.build();
1055        assert_eq!(props.get(ADDED_DATA_FILES).unwrap(), "2");
1056        assert_eq!(props.get(ADDED_RECORDS).unwrap(), "30");
1057
1058        let mut summary_three = SnapshotSummaryCollector::default();
1059        let mut summary_four = SnapshotSummaryCollector::default();
1060
1061        summary_three.add_manifest(&ManifestFile {
1062            manifest_path: "test.manifest".to_string(),
1063            manifest_length: 0,
1064            partition_spec_id: 0,
1065            content: ManifestContentType::Data,
1066            sequence_number: 0,
1067            min_sequence_number: 0,
1068            added_snapshot_id: 0,
1069            added_files_count: Some(1),
1070            existing_files_count: Some(0),
1071            deleted_files_count: Some(0),
1072            added_rows_count: Some(5),
1073            existing_rows_count: Some(0),
1074            deleted_rows_count: Some(0),
1075            partitions: Some(Vec::new()),
1076            key_metadata: None,
1077            first_row_id: None,
1078        });
1079
1080        summary_four.add_file(
1081            &DataFile {
1082                content: DataContentType::Data,
1083                file_path: "test.parquet".into(),
1084                file_format: DataFileFormat::Parquet,
1085                partition: Struct::from_iter(vec![]),
1086                record_count: 1,
1087                file_size_in_bytes: 10,
1088                column_sizes: HashMap::new(),
1089                value_counts: HashMap::new(),
1090                null_value_counts: HashMap::new(),
1091                nan_value_counts: HashMap::new(),
1092                lower_bounds: HashMap::new(),
1093                upper_bounds: HashMap::new(),
1094                key_metadata: None,
1095                split_offsets: None,
1096                equality_ids: None,
1097                sort_order_id: None,
1098                partition_spec_id: 0,
1099                first_row_id: None,
1100                referenced_data_file: None,
1101                content_offset: None,
1102                content_size_in_bytes: None,
1103            },
1104            schema.clone(),
1105            partition_spec.clone(),
1106        );
1107
1108        summary_three.merge(summary_four);
1109        let props = summary_three.build();
1110
1111        assert_eq!(props.get(ADDED_DATA_FILES).unwrap(), "2");
1112        assert_eq!(props.get(ADDED_RECORDS).unwrap(), "6");
1113        assert!(
1114            props
1115                .iter()
1116                .all(|(k, _)| !k.starts_with(CHANGED_PARTITION_PREFIX))
1117        );
1118    }
1119
1120    #[test]
1121    fn test_update_totals_skipped_when_previous_summary_missing_totals() {
1122        let prev_props: HashMap<String, String> = [(TOTAL_DATA_FILES, "8")]
1123            .into_iter()
1124            .map(|(k, v)| (k.to_string(), v.to_string()))
1125            .collect();
1126
1127        let previous_summary = Summary {
1128            operation: Operation::Overwrite,
1129            additional_properties: prev_props,
1130        };
1131
1132        let new_props: HashMap<String, String> = [
1133            (ADDED_DATA_FILES, "4"),
1134            (ADDED_DELETE_FILES, "2"),
1135            (ADDED_RECORDS, "40"),
1136            (ADDED_FILE_SIZE, "400"),
1137            (ADDED_POSITION_DELETES, "5"),
1138            (ADDED_EQUALITY_DELETES, "3"),
1139        ]
1140        .into_iter()
1141        .map(|(k, v)| (k.to_string(), v.to_string()))
1142        .collect();
1143
1144        let summary = Summary {
1145            operation: Operation::Append,
1146            additional_properties: new_props,
1147        };
1148
1149        let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
1150        let props = &updated.additional_properties;
1151
1152        assert_eq!(props.get(TOTAL_DATA_FILES).unwrap(), "12");
1153
1154        for total_field in [
1155            TOTAL_DELETE_FILES,
1156            TOTAL_RECORDS,
1157            TOTAL_FILE_SIZE,
1158            TOTAL_POSITION_DELETES,
1159            TOTAL_EQUALITY_DELETES,
1160        ] {
1161            assert!(
1162                !props.contains_key(total_field),
1163                "{total_field} should not be set when previous summary lacks it",
1164            );
1165        }
1166    }
1167
1168    #[test]
1169    fn test_update_totals_tolerates_unparsable_added_value() {
1170        // A non-integer added value (which can survive in a previous snapshot's
1171        // summary) must not panic the commit. Matching iceberg-java's `updateTotal`
1172        // try/catch, the affected total is skipped while other totals still compute.
1173        let prev_props: HashMap<String, String> = [(TOTAL_DATA_FILES, "8"), (TOTAL_RECORDS, "80")]
1174            .into_iter()
1175            .map(|(k, v)| (k.to_string(), v.to_string()))
1176            .collect();
1177
1178        let previous_summary = Summary {
1179            operation: Operation::Append,
1180            additional_properties: prev_props,
1181        };
1182
1183        let new_props: HashMap<String, String> =
1184            [(ADDED_DATA_FILES, "not-a-number"), (ADDED_RECORDS, "40")]
1185                .into_iter()
1186                .map(|(k, v)| (k.to_string(), v.to_string()))
1187                .collect();
1188
1189        let summary = Summary {
1190            operation: Operation::Append,
1191            additional_properties: new_props,
1192        };
1193
1194        // Must not panic.
1195        let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
1196        let props = &updated.additional_properties;
1197
1198        // The total whose added delta was unparsable is skipped...
1199        assert!(
1200            !props.contains_key(TOTAL_DATA_FILES),
1201            "TOTAL_DATA_FILES should be skipped when its added value is unparsable",
1202        );
1203        // ...while a sibling total with valid deltas still computes.
1204        assert_eq!(props.get(TOTAL_RECORDS).unwrap(), "120");
1205    }
1206
1207    #[test]
1208    fn test_update_totals_computed_when_no_previous_summary() {
1209        let new_props: HashMap<String, String> = [
1210            (ADDED_DATA_FILES, "4"),
1211            (ADDED_RECORDS, "40"),
1212            (ADDED_FILE_SIZE, "400"),
1213        ]
1214        .into_iter()
1215        .map(|(k, v)| (k.to_string(), v.to_string()))
1216        .collect();
1217
1218        let summary = Summary {
1219            operation: Operation::Append,
1220            additional_properties: new_props,
1221        };
1222
1223        let updated = update_snapshot_summaries(summary, None, false).unwrap();
1224        let props = &updated.additional_properties;
1225
1226        assert_eq!(props.get(TOTAL_DATA_FILES).unwrap(), "4");
1227        assert_eq!(props.get(TOTAL_RECORDS).unwrap(), "40");
1228        assert_eq!(props.get(TOTAL_FILE_SIZE).unwrap(), "400");
1229    }
1230
1231    #[test]
1232    fn test_update_totals_with_removes_only() {
1233        let prev_props: HashMap<String, String> = [
1234            (TOTAL_DATA_FILES, "10"),
1235            (TOTAL_DELETE_FILES, "5"),
1236            (TOTAL_RECORDS, "100"),
1237            (TOTAL_FILE_SIZE, "1000"),
1238            (TOTAL_POSITION_DELETES, "3"),
1239            (TOTAL_EQUALITY_DELETES, "2"),
1240        ]
1241        .into_iter()
1242        .map(|(k, v)| (k.to_string(), v.to_string()))
1243        .collect();
1244
1245        let previous_summary = Summary {
1246            operation: Operation::Overwrite,
1247            additional_properties: prev_props,
1248        };
1249
1250        let new_props: HashMap<String, String> = [
1251            (DELETED_DATA_FILES, "2"),
1252            (REMOVED_DELETE_FILES, "1"),
1253            (DELETED_RECORDS, "20"),
1254            (REMOVED_FILE_SIZE, "200"),
1255            (REMOVED_POSITION_DELETES, "1"),
1256            (REMOVED_EQUALITY_DELETES, "1"),
1257        ]
1258        .into_iter()
1259        .map(|(k, v)| (k.to_string(), v.to_string()))
1260        .collect();
1261
1262        let summary = Summary {
1263            operation: Operation::Delete,
1264            additional_properties: new_props,
1265        };
1266
1267        let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
1268        let props = &updated.additional_properties;
1269
1270        assert_eq!(props.get(TOTAL_DATA_FILES).unwrap(), "8");
1271        assert_eq!(props.get(TOTAL_DELETE_FILES).unwrap(), "4");
1272        assert_eq!(props.get(TOTAL_RECORDS).unwrap(), "80");
1273        assert_eq!(props.get(TOTAL_FILE_SIZE).unwrap(), "800");
1274        assert_eq!(props.get(TOTAL_POSITION_DELETES).unwrap(), "2");
1275        assert_eq!(props.get(TOTAL_EQUALITY_DELETES).unwrap(), "1");
1276    }
1277}