Skip to main content

iceberg/transaction/
expire_snapshots.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::HashSet;
19use std::sync::Arc;
20
21use async_trait::async_trait;
22use chrono::Utc;
23
24use crate::error::invalid_data;
25use crate::spec::{
26    MAIN_BRANCH, SnapshotReference, SnapshotRetention, TableMetadata, TableProperties,
27};
28use crate::table::Table;
29use crate::transaction::action::{ActionCommit, TransactionAction};
30use crate::{Error, Result, TableRequirement, TableUpdate};
31
32/// A transaction action that removes snapshots from table metadata.
33///
34/// This only rewrites metadata; the now-unreferenced data and metadata files are left untouched.
35/// Physical file cleanup is the responsibility of a higher-level maintenance operation built on
36/// top of this action.
37///
38/// Selection follows Java `RemoveSnapshots`:
39/// - Explicit ids ([`expire_snapshot_ids`](Self::expire_snapshot_ids)) and age-based expiry are
40///   combined: a snapshot is expired if it is named explicitly *or* selected by age.
41/// - Age-based expiry always runs. The cutoff is [`expire_older_than_ms`](Self::expire_older_than_ms)
42///   when set, a per-branch `max_snapshot_age_ms` for that branch, otherwise
43///   `now - history.expire.max-snapshot-age-ms` (default 5 days), matching Java's constructor default.
44/// - Expiry is computed per branch along each branch's ancestry: each branch keeps its most recent
45///   [`retain_last`](Self::retain_last) snapshots — defaulting to `history.expire.min-snapshots-to-keep`,
46///   with a per-ref `min_snapshots_to_keep` overriding both — plus any ancestor newer than the cutoff,
47///   so a shared ancestor reachable from a retained branch is never expired.
48/// - Refs are aged out first: a non-`main` branch or tag whose head is older than its
49///   `max_ref_age_ms` (defaulting to `history.expire.max-ref-age-ms`) is removed, and snapshots only
50///   that ref retained then become expirable.
51/// - Heads of retained refs (including the current snapshot) are never expired, and naming one
52///   explicitly is an error, since
53///   [`remove_snapshots`](crate::spec::TableMetadataBuilder::remove_snapshots) would otherwise
54///   drop the ref silently.
55pub struct ExpireSnapshotsAction {
56    explicit_ids_to_remove: Vec<i64>,
57    older_than_ms: Option<i64>,
58    retain_last: Option<usize>,
59}
60
61impl ExpireSnapshotsAction {
62    pub(crate) fn new() -> Self {
63        Self {
64            explicit_ids_to_remove: vec![],
65            older_than_ms: None,
66            retain_last: None,
67        }
68    }
69
70    /// Expire these snapshot ids in addition to any age-based selection.
71    ///
72    /// Age-based expiry runs by default (see the type-level docs), so a call that only names ids
73    /// still expires snapshots older than `history.expire.max-snapshot-age-ms`. Pin
74    /// [`expire_older_than_ms`](Self::expire_older_than_ms) to a very old timestamp to expire by id
75    /// alone.
76    ///
77    /// Ids accumulate across calls (like [`add_data_files`](crate::transaction::Transaction::fast_append)).
78    /// An id that is still referenced by a branch or tag cannot be expired and causes commits to fail.
79    pub fn expire_snapshot_ids(mut self, snapshot_ids: impl IntoIterator<Item = i64>) -> Self {
80        self.explicit_ids_to_remove.extend(snapshot_ids);
81        self
82    }
83
84    /// Expire snapshots whose timestamp is strictly older than `older_than_ms`.
85    pub fn expire_older_than_ms(mut self, older_than_ms: i64) -> Self {
86        self.older_than_ms = Some(older_than_ms);
87        self
88    }
89
90    /// Keep at least the `retain_last` most recent snapshots of each branch when expiring by age
91    /// (defaults to the table's `history.expire.min-snapshots-to-keep`, must be at least 1).
92    ///
93    /// This only bounds the age cutoff; it does not protect snapshots named via
94    /// [`expire_snapshot_ids`](Self::expire_snapshot_ids). Setting it to 0 makes commit fail.
95    pub fn retain_last(mut self, retain_last: usize) -> Self {
96        self.retain_last = Some(retain_last);
97        self
98    }
99
100    /// Resolves the snapshots and refs to remove, following Java `RemoveSnapshots.internalApply`.
101    fn plan(&self, table: &Table, properties: &TableProperties<'_>) -> Result<ExpirePlan> {
102        // Matches Java `RemoveSnapshots.retainLast`, which requires at least one snapshot.
103        if self.retain_last == Some(0) {
104            return Err(invalid_data!(
105                "Number of snapshots to retain must be at least 1"
106            ));
107        }
108
109        let metadata = table.metadata();
110        let now = Utc::now().timestamp_millis();
111        // When a knob is not set explicitly, fall back to the table's `history.expire.*` properties,
112        // matching Java `RemoveSnapshots`' constructor. With the default `max-snapshot-age-ms` (5
113        // days) the age path always runs, so even an explicit-id-only call applies the default cutoff.
114        let default_cutoff = match self.older_than_ms {
115            Some(older_than_ms) => older_than_ms,
116            None => now.saturating_sub(properties.max_snapshot_age_ms()?),
117        };
118        let default_min_to_keep = match self.retain_last {
119            Some(retain_last) => retain_last,
120            None => properties.min_snapshots_to_keep()?,
121        };
122
123        // Ref aging: `main` is always kept; any other ref whose head is older than its
124        // `max_ref_age_ms` (defaulting to `history.expire.max-ref-age-ms`) is dropped, like Java's
125        // `computeRetainedRefs`.
126        let default_max_ref_age_ms = properties.max_ref_age_ms()?;
127        let mut removed_ref_names: Vec<String> = vec![];
128        let mut retained_refs: Vec<&SnapshotReference> = vec![];
129        for (ref_name, snapshot_ref) in &metadata.refs {
130            if ref_name == MAIN_BRANCH
131                || !Self::ref_aged_out(metadata, snapshot_ref, now, default_max_ref_age_ms)
132            {
133                retained_refs.push(snapshot_ref);
134            } else {
135                removed_ref_names.push(ref_name.clone());
136            }
137        }
138
139        // Heads of retained refs (plus the current snapshot) are never expired; naming one
140        // explicitly is an error, since `remove_snapshots` would otherwise drop the ref silently.
141        let mut ref_head_ids: HashSet<i64> = retained_refs.iter().map(|r| r.snapshot_id).collect();
142        if let Some(current_id) = metadata.current_snapshot_id() {
143            ref_head_ids.insert(current_id);
144        }
145
146        let existing_ids: HashSet<i64> = metadata.snapshots().map(|s| s.snapshot_id()).collect();
147        let mut expiring_ids: HashSet<i64> = HashSet::new();
148        for id in &self.explicit_ids_to_remove {
149            if ref_head_ids.contains(id) {
150                return Err(Self::reference_error(metadata, *id));
151            }
152            if existing_ids.contains(id) {
153                expiring_ids.insert(*id);
154            }
155        }
156
157        // Per-branch retention: keep each branch's most recent `min_to_keep` ancestors plus any
158        // newer than the branch cutoff. The current snapshot is treated as a branch (default policy)
159        // so its lineage is protected even when there is no explicit `main` ref.
160        let mut retained_ids = ref_head_ids.clone();
161        let mut referenced_ids = ref_head_ids.clone();
162        let mut branches: Vec<(i64, usize, i64)> = vec![];
163        for snapshot_ref in &retained_refs {
164            match &snapshot_ref.retention {
165                SnapshotRetention::Branch {
166                    min_snapshots_to_keep,
167                    max_snapshot_age_ms,
168                    ..
169                } => {
170                    let min_to_keep =
171                        min_snapshots_to_keep.map_or(default_min_to_keep, |m| m as usize);
172                    let cutoff =
173                        max_snapshot_age_ms.map_or(default_cutoff, |age| now.saturating_sub(age));
174                    branches.push((snapshot_ref.snapshot_id, min_to_keep, cutoff));
175                }
176                SnapshotRetention::Tag { .. } => {
177                    referenced_ids.insert(snapshot_ref.snapshot_id);
178                }
179            }
180        }
181        if let Some(current_id) = metadata.current_snapshot_id()
182            && !branches
183                .iter()
184                .any(|(head_id, _, _)| *head_id == current_id)
185        {
186            branches.push((current_id, default_min_to_keep, default_cutoff));
187        }
188        for (head_id, min_to_keep, cutoff) in branches {
189            Self::retain_branch(
190                metadata,
191                head_id,
192                min_to_keep,
193                cutoff,
194                &mut retained_ids,
195                &mut referenced_ids,
196            );
197        }
198
199        // Unreferenced snapshots newer than the default cutoff are kept (Java's
200        // `unreferencedSnapshotsToRetain`); everything else not retained is expired.
201        for snapshot in metadata.snapshots() {
202            let id = snapshot.snapshot_id();
203            if !referenced_ids.contains(&id) && snapshot.timestamp_ms() >= default_cutoff {
204                retained_ids.insert(id);
205            }
206        }
207        for snapshot in metadata.snapshots() {
208            if !retained_ids.contains(&snapshot.snapshot_id()) {
209                expiring_ids.insert(snapshot.snapshot_id());
210            }
211        }
212
213        let mut ids_to_remove: Vec<i64> = expiring_ids.into_iter().collect();
214        ids_to_remove.sort_unstable();
215        removed_ref_names.sort();
216        Ok(ExpirePlan {
217            ids_to_remove,
218            refs_to_remove: removed_ref_names,
219        })
220    }
221
222    /// Whether a non-main ref should be dropped because its head is older than its `max_ref_age_ms`,
223    /// defaulting to `default_max_ref_age_ms` (`history.expire.max-ref-age-ms`) when the ref sets no
224    /// window of its own. The default `i64::MAX` effectively never ages a ref out.
225    fn ref_aged_out(
226        metadata: &TableMetadata,
227        snapshot_ref: &SnapshotReference,
228        now: i64,
229        default_max_ref_age_ms: i64,
230    ) -> bool {
231        let max_ref_age_ms = match snapshot_ref.retention {
232            SnapshotRetention::Branch { max_ref_age_ms, .. }
233            | SnapshotRetention::Tag { max_ref_age_ms } => max_ref_age_ms,
234        }
235        .unwrap_or(default_max_ref_age_ms);
236        match metadata.snapshot_by_id(snapshot_ref.snapshot_id) {
237            Some(snapshot) => now.saturating_sub(snapshot.timestamp_ms()) > max_ref_age_ms,
238            None => false,
239        }
240    }
241
242    /// Walks a branch's ancestry (Java's `computeBranchSnapshotsToRetain`), retaining each ancestor
243    /// while fewer than `min` are kept or it is newer than `cutoff`, and recording every ancestor as
244    /// referenced. Ancestry timestamps decrease monotonically, so this needs no early break.
245    fn retain_branch(
246        metadata: &TableMetadata,
247        head_id: i64,
248        min_to_keep: usize,
249        cutoff: i64,
250        retained_ids: &mut HashSet<i64>,
251        referenced_ids: &mut HashSet<i64>,
252    ) {
253        let mut kept_count = 0usize;
254        for ancestor_id in Self::ancestors(metadata, head_id) {
255            referenced_ids.insert(ancestor_id);
256            let timestamp = metadata
257                .snapshot_by_id(ancestor_id)
258                .map_or(i64::MIN, |snapshot| snapshot.timestamp_ms());
259            if kept_count < min_to_keep || timestamp >= cutoff {
260                retained_ids.insert(ancestor_id);
261                kept_count += 1;
262            }
263        }
264    }
265
266    /// Iterates a snapshot and its ancestors, newest first, following `parent_snapshot_id`.
267    fn ancestors(metadata: &TableMetadata, head_id: i64) -> impl Iterator<Item = i64> + '_ {
268        let mut next_id = Some(head_id);
269        std::iter::from_fn(move || {
270            let id = next_id?;
271            next_id = metadata
272                .snapshot_by_id(id)
273                .and_then(|snapshot| snapshot.parent_snapshot_id());
274            Some(id)
275        })
276    }
277
278    fn reference_error(metadata: &TableMetadata, snapshot_id: i64) -> Error {
279        if metadata.current_snapshot_id() == Some(snapshot_id) {
280            return invalid_data!("Cannot expire the current snapshot");
281        }
282        let ref_names: Vec<&str> = metadata
283            .refs
284            .iter()
285            .filter(|(_, snapshot_ref)| snapshot_ref.snapshot_id == snapshot_id)
286            .map(|(ref_name, _)| ref_name.as_str())
287            .collect();
288        invalid_data!("Cannot expire snapshot {snapshot_id}: still referenced by {ref_names:?}")
289    }
290}
291
292/// Snapshots and refs an [`ExpireSnapshotsAction`] resolves to remove.
293struct ExpirePlan {
294    ids_to_remove: Vec<i64>,
295    refs_to_remove: Vec<String>,
296}
297
298#[async_trait]
299impl TransactionAction for ExpireSnapshotsAction {
300    async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
301        let metadata = table.metadata();
302        let properties = metadata.table_properties();
303
304        // Expiring metadata defeats a user's explicit decision to disable GC (Java refuses too).
305        if !properties.gc_enabled()? {
306            return Err(invalid_data!(
307                "Cannot expire snapshots: gc.enabled is false"
308            ));
309        }
310
311        let plan = self.plan(table, &properties)?;
312
313        if plan.ids_to_remove.is_empty() && plan.refs_to_remove.is_empty() {
314            return Ok(ActionCommit::new(vec![], vec![]));
315        }
316
317        // Drop aged-out refs first, then the snapshots no ref retains anymore.
318        let mut updates: Vec<TableUpdate> = plan
319            .refs_to_remove
320            .into_iter()
321            .map(|ref_name| TableUpdate::RemoveSnapshotRef { ref_name })
322            .collect();
323
324        // Drop statistics metadata for expired snapshots.
325        // This only updates metadata; puffin files are cleaned up separately.
326        let mut stats_updates: Vec<TableUpdate> = vec![];
327        for &snapshot_id in &plan.ids_to_remove {
328            stats_updates.extend(
329                metadata
330                    .statistics_for_snapshot(snapshot_id)
331                    .is_some()
332                    .then_some(TableUpdate::RemoveStatistics { snapshot_id }),
333            );
334            stats_updates.extend(
335                metadata
336                    .partition_statistics_for_snapshot(snapshot_id)
337                    .is_some()
338                    .then_some(TableUpdate::RemovePartitionStatistics { snapshot_id }),
339            );
340        }
341
342        if !plan.ids_to_remove.is_empty() {
343            updates.push(TableUpdate::RemoveSnapshots {
344                snapshot_ids: plan.ids_to_remove,
345            });
346        }
347        updates.extend(stats_updates);
348
349        // The ref assertion closes the race where a concurrent writer advances `main` between
350        // selection and commit, which could orphan a snapshot whose parent we are about to remove.
351        Ok(ActionCommit::new(updates, vec![
352            TableRequirement::UuidMatch {
353                uuid: metadata.uuid(),
354            },
355            TableRequirement::RefSnapshotIdMatch {
356                r#ref: MAIN_BRANCH.to_string(),
357                snapshot_id: metadata.current_snapshot_id(),
358            },
359        ]))
360    }
361}
362
363#[cfg(test)]
364mod tests {
365    use std::collections::HashMap;
366    use std::sync::Arc;
367
368    use chrono::Utc;
369
370    use crate::spec::{
371        MAIN_BRANCH, Operation, PartitionStatisticsFile, Snapshot, SnapshotReference,
372        SnapshotRetention, StatisticsFile, Summary,
373    };
374    use crate::table::Table;
375    use crate::transaction::Transaction;
376    use crate::transaction::action::{ApplyTransactionAction, TransactionAction};
377    use crate::transaction::expire_snapshots::ExpireSnapshotsAction;
378    use crate::transaction::tests::{make_v2_minimal_table, make_v2_table};
379    use crate::{TableRequirement, TableUpdate};
380
381    // `make_v2_table` carries an older snapshot (ts 1515100955770) and a current
382    // snapshot (ts 1555100955770).
383    const OLD_SNAPSHOT: i64 = 3051729675574597004;
384    const CURRENT_SNAPSHOT: i64 = 3055729675574597004;
385    // Well after the minimal table's last-updated-ms, so synthetic snapshots pass timestamp checks.
386    const TS: i64 = 1_700_000_000_000;
387    // A cutoff at the epoch makes age-based expiry a no-op (every snapshot is newer), isolating the
388    // explicit-id behavior under test now that the default cutoff (`now - 5 days`) always runs.
389    const NO_AGE_EXPIRY: i64 = 0;
390
391    fn action() -> ExpireSnapshotsAction {
392        ExpireSnapshotsAction::new()
393    }
394
395    async fn removed_ids(action: ExpireSnapshotsAction) -> Vec<i64> {
396        expired(&make_v2_table(), action).await
397    }
398
399    async fn updates_of(table: &Table, action: ExpireSnapshotsAction) -> Vec<TableUpdate> {
400        Arc::new(action).commit(table).await.unwrap().take_updates()
401    }
402
403    async fn expired(table: &Table, action: ExpireSnapshotsAction) -> Vec<i64> {
404        updates_of(table, action)
405            .await
406            .into_iter()
407            .find_map(|update| match update {
408                TableUpdate::RemoveSnapshots { snapshot_ids } => Some(snapshot_ids),
409                _ => None,
410            })
411            .unwrap_or_default()
412    }
413
414    fn removed_refs(updates: &[TableUpdate]) -> Vec<String> {
415        let mut refs: Vec<String> = updates
416            .iter()
417            .filter_map(|update| match update {
418                TableUpdate::RemoveSnapshotRef { ref_name } => Some(ref_name.clone()),
419                _ => None,
420            })
421            .collect();
422        refs.sort();
423        refs
424    }
425
426    fn snapshot(id: i64, parent: Option<i64>, sequence_number: i64, timestamp_ms: i64) -> Snapshot {
427        Snapshot::builder()
428            .with_snapshot_id(id)
429            .with_parent_snapshot_id(parent)
430            .with_sequence_number(sequence_number)
431            .with_timestamp_ms(timestamp_ms)
432            .with_schema_id(0)
433            .with_manifest_list(format!("/snap-{id}.avro"))
434            .with_summary(Summary {
435                operation: Operation::Append,
436                additional_properties: HashMap::new(),
437            })
438            .build()
439    }
440
441    fn branch(snapshot_id: i64, min_snapshots_to_keep: Option<i32>) -> SnapshotReference {
442        branch_with(snapshot_id, min_snapshots_to_keep, None, None)
443    }
444
445    fn branch_with(
446        snapshot_id: i64,
447        min_snapshots_to_keep: Option<i32>,
448        max_snapshot_age_ms: Option<i64>,
449        max_ref_age_ms: Option<i64>,
450    ) -> SnapshotReference {
451        SnapshotReference {
452            snapshot_id,
453            retention: SnapshotRetention::Branch {
454                min_snapshots_to_keep,
455                max_snapshot_age_ms,
456                max_ref_age_ms,
457            },
458        }
459    }
460
461    fn tag(snapshot_id: i64, max_ref_age_ms: Option<i64>) -> SnapshotReference {
462        SnapshotReference {
463            snapshot_id,
464            retention: SnapshotRetention::Tag { max_ref_age_ms },
465        }
466    }
467
468    /// Builds a table from synthetic snapshots and refs on top of an empty base.
469    fn table_with(snapshots: Vec<Snapshot>, refs: Vec<(&str, SnapshotReference)>) -> Table {
470        table_with_props(snapshots, refs, HashMap::new())
471    }
472
473    /// Like [`table_with`], but also seeds table properties (e.g. `history.expire.*` defaults).
474    fn table_with_props(
475        snapshots: Vec<Snapshot>,
476        refs: Vec<(&str, SnapshotReference)>,
477        properties: HashMap<String, String>,
478    ) -> Table {
479        let base = make_v2_minimal_table();
480        let mut builder = base
481            .metadata()
482            .clone()
483            .into_builder(None)
484            .set_properties(properties)
485            .unwrap();
486        for snapshot in snapshots {
487            builder = builder.add_snapshot(snapshot).unwrap();
488        }
489        for (name, reference) in refs {
490            builder = builder.set_ref(name, reference).unwrap();
491        }
492        base.with_metadata(Arc::new(builder.build().unwrap().metadata))
493    }
494
495    /// Like [`table_with`], but also attaches statistics and partition-statistics files. Reuses
496    /// [`table_with`] for the snapshot/ref wiring so the two can't drift, then layers stats on top.
497    fn table_with_stats(
498        snapshots: Vec<Snapshot>,
499        refs: Vec<(&str, SnapshotReference)>,
500        statistics: Vec<StatisticsFile>,
501        partition_statistics: Vec<PartitionStatisticsFile>,
502    ) -> Table {
503        let table = table_with(snapshots, refs);
504        let mut builder = table.metadata().clone().into_builder(None);
505        for stats in statistics {
506            builder = builder.set_statistics(stats);
507        }
508        for stats in partition_statistics {
509            builder = builder.set_partition_statistics(stats);
510        }
511        table.with_metadata(Arc::new(builder.build().unwrap().metadata))
512    }
513
514    fn stats_file(snapshot_id: i64) -> StatisticsFile {
515        StatisticsFile {
516            snapshot_id,
517            statistics_path: format!("/stats-{snapshot_id}.puffin"),
518            file_size_in_bytes: 1,
519            file_footer_size_in_bytes: 1,
520            key_metadata: None,
521            blob_metadata: vec![],
522        }
523    }
524
525    fn partition_stats_file(snapshot_id: i64) -> PartitionStatisticsFile {
526        PartitionStatisticsFile {
527            snapshot_id,
528            statistics_path: format!("/partition-stats-{snapshot_id}.puffin"),
529            file_size_in_bytes: 1,
530        }
531    }
532
533    fn removed_statistics(updates: &[TableUpdate]) -> Vec<i64> {
534        updates
535            .iter()
536            .filter_map(|update| match update {
537                TableUpdate::RemoveStatistics { snapshot_id } => Some(*snapshot_id),
538                _ => None,
539            })
540            .collect()
541    }
542
543    fn removed_partition_statistics(updates: &[TableUpdate]) -> Vec<i64> {
544        updates
545            .iter()
546            .filter_map(|update| match update {
547                TableUpdate::RemovePartitionStatistics { snapshot_id } => Some(*snapshot_id),
548                _ => None,
549            })
550            .collect()
551    }
552
553    #[tokio::test]
554    async fn test_expire_explicit_snapshot_id() {
555        assert_eq!(
556            removed_ids(
557                action()
558                    .expire_snapshot_ids(vec![OLD_SNAPSHOT])
559                    .expire_older_than_ms(NO_AGE_EXPIRY)
560            )
561            .await,
562            vec![OLD_SNAPSHOT]
563        );
564    }
565
566    #[tokio::test]
567    async fn test_explicit_unknown_id_is_ignored() {
568        assert!(
569            removed_ids(
570                action()
571                    .expire_snapshot_ids(vec![42])
572                    .expire_older_than_ms(NO_AGE_EXPIRY)
573            )
574            .await
575            .is_empty()
576        );
577    }
578
579    #[tokio::test]
580    async fn test_cannot_expire_current_snapshot() {
581        let table = make_v2_table();
582        let action = action().expire_snapshot_ids(vec![CURRENT_SNAPSHOT]);
583        assert!(Arc::new(action).commit(&table).await.is_err());
584    }
585
586    /// `make_v2_table` with a tag pointing at the older snapshot.
587    fn table_with_tag_on_old() -> Table {
588        let table = make_v2_table();
589        let metadata = table
590            .metadata()
591            .clone()
592            .into_builder(None)
593            .set_ref("history-tag", SnapshotReference {
594                snapshot_id: OLD_SNAPSHOT,
595                retention: SnapshotRetention::Tag {
596                    max_ref_age_ms: None,
597                },
598            })
599            .unwrap()
600            .build()
601            .unwrap()
602            .metadata;
603        table.with_metadata(Arc::new(metadata))
604    }
605
606    #[tokio::test]
607    async fn test_cannot_expire_tagged_snapshot_explicitly() {
608        let table = table_with_tag_on_old();
609        let action = action().expire_snapshot_ids(vec![OLD_SNAPSHOT]);
610        assert!(Arc::new(action).commit(&table).await.is_err());
611    }
612
613    #[tokio::test]
614    async fn test_age_expiry_skips_tagged_snapshot() {
615        let table = table_with_tag_on_old();
616        let mut commit = Arc::new(action().expire_older_than_ms(i64::MAX))
617            .commit(&table)
618            .await
619            .unwrap();
620        // Both snapshots are referenced (current + tag), so nothing is expired.
621        assert!(commit.take_updates().is_empty());
622    }
623
624    #[tokio::test]
625    async fn test_retain_last_default_expires_older_non_current() {
626        assert_eq!(
627            removed_ids(action().expire_older_than_ms(i64::MAX)).await,
628            vec![OLD_SNAPSHOT]
629        );
630    }
631
632    #[tokio::test]
633    async fn test_retain_last_noop_when_enough_retained() {
634        assert!(removed_ids(action().retain_last(5)).await.is_empty());
635    }
636
637    #[tokio::test]
638    async fn test_older_than_excludes_newer_snapshots() {
639        // Threshold older than every snapshot -> nothing qualifies.
640        assert!(
641            removed_ids(action().expire_older_than_ms(1))
642                .await
643                .is_empty()
644        );
645    }
646
647    #[tokio::test]
648    async fn test_apply_registers_action() {
649        let table = make_v2_table();
650        let tx = Transaction::new(&table);
651        let tx = tx
652            .expire_snapshots()
653            .expire_snapshot_ids(vec![OLD_SNAPSHOT])
654            .apply(tx)
655            .unwrap();
656        assert_eq!(tx.actions.len(), 1);
657    }
658
659    #[tokio::test]
660    async fn test_per_branch_retention_protects_shared_ancestor() {
661        // main: 1 -> 2 -> 3 ; branch `b`: 1 -> 2 -> 4. Snapshot 2 is a shared ancestor of both
662        // branches but is not a ref head.
663        let table = table_with(
664            vec![
665                snapshot(1, None, 35, TS + 1),
666                snapshot(2, Some(1), 36, TS + 2),
667                snapshot(3, Some(2), 37, TS + 3),
668                snapshot(4, Some(2), 38, TS + 4),
669            ],
670            vec![(MAIN_BRANCH, branch(3, None)), ("b", branch(4, None))],
671        );
672
673        // A global "newest 2" would expire 2 and orphan branch `b`; per-branch retention keeps it.
674        let removed = expired(
675            &table,
676            action().retain_last(2).expire_older_than_ms(i64::MAX),
677        )
678        .await;
679        assert_eq!(removed, vec![1]);
680    }
681
682    #[tokio::test]
683    async fn test_per_ref_min_snapshots_to_keep_overrides_retain_last() {
684        let table = table_with(
685            vec![
686                snapshot(1, None, 35, TS + 1),
687                snapshot(2, Some(1), 36, TS + 2),
688                snapshot(3, Some(2), 37, TS + 3),
689            ],
690            vec![(MAIN_BRANCH, branch(3, Some(3)))],
691        );
692
693        // The branch's own min_snapshots_to_keep=3 wins over the action's retain_last(1).
694        let removed = expired(
695            &table,
696            action().retain_last(1).expire_older_than_ms(i64::MAX),
697        )
698        .await;
699        assert!(removed.is_empty());
700    }
701
702    #[tokio::test]
703    async fn test_explicit_and_age_combine() {
704        let table = table_with(
705            vec![
706                snapshot(1, None, 35, TS + 1),
707                snapshot(2, Some(1), 36, TS + 2),
708                snapshot(3, Some(2), 37, TS + 3),
709                snapshot(4, Some(3), 38, TS + 4),
710            ],
711            vec![(MAIN_BRANCH, branch(4, None))],
712        );
713
714        // Age expires 1 and 2 (older than the cutoff, beyond retain_last). 3 is newer than the
715        // cutoff so age keeps it, but it is named explicitly, so all three are expired.
716        let removed = expired(
717            &table,
718            action()
719                .retain_last(1)
720                .expire_older_than_ms(TS + 3)
721                .expire_snapshot_ids(vec![3]),
722        )
723        .await;
724        assert_eq!(removed, vec![1, 2, 3]);
725    }
726
727    #[tokio::test]
728    async fn test_expire_snapshot_ids_accumulates() {
729        let table = table_with(
730            vec![
731                snapshot(1, None, 35, TS + 1),
732                snapshot(2, Some(1), 36, TS + 2),
733                snapshot(3, Some(2), 37, TS + 3),
734            ],
735            vec![(MAIN_BRANCH, branch(3, None))],
736        );
737
738        // Two separate calls both take effect (age expiry pinned off to isolate accumulation).
739        let removed = expired(
740            &table,
741            action()
742                .expire_snapshot_ids(vec![1])
743                .expire_snapshot_ids(vec![2])
744                .expire_older_than_ms(NO_AGE_EXPIRY),
745        )
746        .await;
747        assert_eq!(removed, vec![1, 2]);
748    }
749
750    #[tokio::test]
751    async fn test_gc_disabled_errors() {
752        let table = make_v2_table();
753        let metadata = table
754            .metadata()
755            .clone()
756            .into_builder(None)
757            .set_properties(HashMap::from([(
758                "gc.enabled".to_string(),
759                "false".to_string(),
760            )]))
761            .unwrap()
762            .build()
763            .unwrap()
764            .metadata;
765        let table = table.with_metadata(Arc::new(metadata));
766
767        let action = action().expire_snapshot_ids(vec![OLD_SNAPSHOT]);
768        assert!(Arc::new(action).commit(&table).await.is_err());
769    }
770
771    #[tokio::test]
772    async fn test_commit_asserts_main_ref() {
773        let table = make_v2_table();
774        let mut commit = Arc::new(action().expire_snapshot_ids(vec![OLD_SNAPSHOT]))
775            .commit(&table)
776            .await
777            .unwrap();
778        assert!(
779            commit
780                .take_requirements()
781                .iter()
782                .any(|requirement| matches!(
783                    requirement,
784                    TableRequirement::RefSnapshotIdMatch { r#ref, snapshot_id }
785                        if r#ref == MAIN_BRANCH && *snapshot_id == Some(CURRENT_SNAPSHOT)
786                ))
787        );
788    }
789
790    #[tokio::test]
791    async fn test_ref_aging_drops_old_tag_and_expires_its_snapshot() {
792        let now = Utc::now().timestamp_millis();
793        let day_ms = 24 * 60 * 60 * 1000;
794        // main: 2 (recent). Isolated snapshot 1 (old) is only kept alive by `old-tag`, whose own
795        // age (10 days) exceeds its max-ref-age of 1 day.
796        let table = table_with(
797            vec![
798                snapshot(1, None, 35, now - 10 * day_ms),
799                snapshot(2, None, 36, now - 1000),
800            ],
801            // Set the tag before main so the builder's last-updated bookkeeping stays monotonic.
802            vec![
803                ("old-tag", tag(1, Some(day_ms))),
804                (MAIN_BRANCH, branch(2, None)),
805            ],
806        );
807
808        // An explicit cutoff lets the freed snapshot 1 actually expire once the tag is gone.
809        let updates = updates_of(&table, action().expire_older_than_ms(now - 5 * day_ms)).await;
810        // The tag is dropped, and snapshot 1 (now unreferenced and old) is expired with it.
811        assert_eq!(removed_refs(&updates), vec!["old-tag".to_string()]);
812        assert!(updates.iter().any(
813            |u| matches!(u, TableUpdate::RemoveSnapshots { snapshot_ids } if snapshot_ids == &[1])
814        ));
815    }
816
817    #[tokio::test]
818    async fn test_ref_aging_keeps_recent_tag() {
819        let now = Utc::now().timestamp_millis();
820        let day_ms = 24 * 60 * 60 * 1000;
821        let table = table_with(
822            vec![
823                snapshot(1, None, 35, now - 1000),
824                snapshot(2, None, 36, now - 500),
825            ],
826            vec![
827                ("fresh-tag", tag(1, Some(day_ms))),
828                (MAIN_BRANCH, branch(2, None)),
829            ],
830        );
831
832        let updates = updates_of(&table, action()).await;
833        // The tag is younger than its max-ref-age, so neither it nor its snapshot is removed.
834        assert!(removed_refs(&updates).is_empty());
835        assert_eq!(expired(&table, action()).await, Vec::<i64>::new());
836    }
837
838    #[tokio::test]
839    async fn test_per_ref_max_snapshot_age_overrides_default() {
840        let now = Utc::now().timestamp_millis();
841        let day_ms = 24 * 60 * 60 * 1000;
842        // main: 1 (3 days old) -> 2 (1 day old) -> 3 (recent), with a 2-day per-ref window.
843        let table = table_with(
844            vec![
845                snapshot(1, None, 35, now - 3 * day_ms),
846                snapshot(2, Some(1), 36, now - day_ms),
847                snapshot(3, Some(2), 37, now - 1000),
848            ],
849            vec![(MAIN_BRANCH, branch_with(3, Some(1), Some(2 * day_ms), None))],
850        );
851
852        // With no explicit cutoff nothing would expire, but main's 2-day window expires snapshot 1.
853        let removed = expired(&table, action()).await;
854        assert_eq!(removed, vec![1]);
855    }
856
857    #[tokio::test]
858    async fn test_ref_aging_drops_old_branch() {
859        let now = Utc::now().timestamp_millis();
860        let day_ms = 24 * 60 * 60 * 1000;
861        // A stale non-main branch (head 1, 10 days old) past its 1-day max-ref-age.
862        let table = table_with(
863            vec![
864                snapshot(1, None, 35, now - 10 * day_ms),
865                snapshot(2, None, 36, now - 1000),
866            ],
867            vec![
868                ("stale", branch_with(1, None, None, Some(day_ms))),
869                (MAIN_BRANCH, branch(2, None)),
870            ],
871        );
872
873        let updates = updates_of(&table, action().expire_older_than_ms(now - 5 * day_ms)).await;
874        // The stale branch is dropped and snapshot 1 (now unreferenced and old) is expired.
875        assert_eq!(removed_refs(&updates), vec!["stale".to_string()]);
876        assert!(updates.iter().any(
877            |u| matches!(u, TableUpdate::RemoveSnapshots { snapshot_ids } if snapshot_ids == &[1])
878        ));
879    }
880
881    #[tokio::test]
882    async fn test_unreferenced_snapshots_retained_only_while_young() {
883        let now = Utc::now().timestamp_millis();
884        let day_ms = 24 * 60 * 60 * 1000;
885        // Snapshot 1 is the main head; 2 and 3 are orphans (no ref, not ancestors). Add oldest
886        // first so the builder's timestamp bookkeeping stays monotonic.
887        let table = table_with(
888            vec![
889                snapshot(3, None, 35, now - 10 * day_ms), // old orphan
890                snapshot(1, None, 36, now - 1000),        // main head
891                snapshot(2, None, 37, now - 1000),        // young orphan
892            ],
893            vec![(MAIN_BRANCH, branch(1, None))],
894        );
895
896        // The young orphan (2) is kept; only the old orphan (3) is expired.
897        let removed = expired(&table, action().expire_older_than_ms(now - 5 * day_ms)).await;
898        assert_eq!(removed, vec![3]);
899    }
900
901    #[tokio::test]
902    async fn test_retain_last_zero_errors() {
903        let table = make_v2_table();
904        assert!(
905            Arc::new(action().retain_last(0))
906                .commit(&table)
907                .await
908                .is_err()
909        );
910    }
911
912    #[tokio::test]
913    async fn test_cannot_expire_branch_head_explicitly() {
914        let table = table_with(
915            vec![snapshot(1, None, 35, TS + 1), snapshot(2, None, 36, TS + 2)],
916            vec![(MAIN_BRANCH, branch(1, None)), ("branch", branch(2, None))],
917        );
918        // 2 is the head of a non-main branch, so it cannot be expired explicitly.
919        let action = action().expire_snapshot_ids(vec![2]);
920        assert!(Arc::new(action).commit(&table).await.is_err());
921    }
922
923    #[tokio::test]
924    async fn test_tag_does_not_protect_its_ancestry() {
925        let now = Utc::now().timestamp_millis();
926        let day_ms = 24 * 60 * 60 * 1000;
927        // Chain 1 -> 2 -> 3 with main rewound to 1 and a tag on 3. Snapshot 2 (the tag target's
928        // parent) is reachable from no ref's retained set, so it is expired.
929        let table = table_with(
930            vec![
931                snapshot(1, None, 35, now - 10 * day_ms),
932                snapshot(2, Some(1), 36, now - 8 * day_ms),
933                snapshot(3, Some(2), 37, now - 1000),
934            ],
935            vec![(MAIN_BRANCH, branch(1, None)), ("tag", tag(3, None))],
936        );
937
938        let removed = expired(&table, action().expire_older_than_ms(now - 5 * day_ms)).await;
939        assert_eq!(removed, vec![2]);
940    }
941
942    #[tokio::test]
943    async fn test_branch_protects_its_ancestry() {
944        let now = Utc::now().timestamp_millis();
945        let day_ms = 24 * 60 * 60 * 1000;
946        // Same topology as the tag case, but a branch (keeping its whole history) replaces the tag.
947        // The branch's parent snapshot 2 is now reachable history and is not expired.
948        let table = table_with(
949            vec![
950                snapshot(1, None, 35, now - 10 * day_ms),
951                snapshot(2, Some(1), 36, now - 8 * day_ms),
952                snapshot(3, Some(2), 37, now - 1000),
953            ],
954            vec![
955                (MAIN_BRANCH, branch(1, None)),
956                ("branch", branch_with(3, None, Some(i64::MAX), None)),
957            ],
958        );
959
960        let removed = expired(&table, action().expire_older_than_ms(now - 5 * day_ms)).await;
961        assert!(removed.is_empty());
962    }
963
964    #[tokio::test]
965    async fn test_per_branch_max_snapshot_age_differs_across_branches() {
966        let now = Utc::now().timestamp_millis();
967        let day_ms = 24 * 60 * 60 * 1000;
968        // main (1 -> 2) keeps 5 days; branch `keep` (3 -> 4) keeps 60 days. Snapshots 1 and 3 are
969        // both 30 days old, but only main's short window expires its ancestor.
970        let table = table_with(
971            vec![
972                snapshot(1, None, 35, now - 30 * day_ms),
973                snapshot(3, None, 36, now - 30 * day_ms),
974                snapshot(2, Some(1), 37, now - 2 * day_ms),
975                snapshot(4, Some(3), 38, now - 2 * day_ms),
976            ],
977            vec![
978                (MAIN_BRANCH, branch_with(2, None, Some(5 * day_ms), None)),
979                ("keep", branch_with(4, None, Some(60 * day_ms), None)),
980            ],
981        );
982
983        // main's 5-day window expires its old ancestor 1; keep's 60-day window retains its old
984        // ancestor 3.
985        let removed = expired(&table, action()).await;
986        assert_eq!(removed, vec![1]);
987    }
988
989    #[tokio::test]
990    async fn test_default_cutoff_expires_snapshots_older_than_max_age() {
991        let now = Utc::now().timestamp_millis();
992        let day_ms = 24 * 60 * 60 * 1000;
993        let table = table_with(
994            vec![
995                snapshot(1, None, 35, now - 10 * day_ms), // older than the default 5-day cutoff
996                snapshot(2, Some(1), 36, now - 1000),     // recent
997            ],
998            vec![(MAIN_BRANCH, branch(2, None))],
999        );
1000
1001        // No explicit cutoff: defaults to now - history.expire.max-snapshot-age-ms (5 days).
1002        let removed = expired(&table, action()).await;
1003        assert_eq!(removed, vec![1]);
1004    }
1005
1006    #[tokio::test]
1007    async fn test_min_snapshots_to_keep_property_is_the_default_floor() {
1008        let table = table_with_props(
1009            vec![
1010                snapshot(1, None, 35, TS + 1),
1011                snapshot(2, Some(1), 36, TS + 2),
1012                snapshot(3, Some(2), 37, TS + 3),
1013            ],
1014            vec![(MAIN_BRANCH, branch(3, None))],
1015            HashMap::from([(
1016                "history.expire.min-snapshots-to-keep".to_string(),
1017                "3".to_string(),
1018            )]),
1019        );
1020
1021        // The snapshots predate the default cutoff, but the table's min-snapshots-to-keep=3 keeps
1022        // the whole chain.
1023        let removed = expired(&table, action()).await;
1024        assert!(removed.is_empty());
1025    }
1026
1027    #[tokio::test]
1028    async fn test_max_snapshot_age_ms_property_sets_the_cutoff() {
1029        let now = Utc::now().timestamp_millis();
1030        let day_ms = 24 * 60 * 60 * 1000;
1031        // Snapshot 1 is 2 days old: the built-in 5-day default would keep it, but the table's
1032        // 1-day history.expire.max-snapshot-age-ms expires it. This pins the cutoff to the property
1033        // value rather than the hardcoded default.
1034        let table = table_with_props(
1035            vec![
1036                snapshot(1, None, 35, now - 2 * day_ms),
1037                snapshot(2, Some(1), 36, now - 1000),
1038            ],
1039            vec![(MAIN_BRANCH, branch(2, None))],
1040            HashMap::from([(
1041                "history.expire.max-snapshot-age-ms".to_string(),
1042                day_ms.to_string(),
1043            )]),
1044        );
1045
1046        let removed = expired(&table, action()).await;
1047        assert_eq!(removed, vec![1]);
1048    }
1049
1050    #[tokio::test]
1051    async fn test_max_ref_age_ms_property_ages_out_ref_without_its_own_window() {
1052        let now = Utc::now().timestamp_millis();
1053        let day_ms = 24 * 60 * 60 * 1000;
1054        // `old-tag` sets no max_ref_age_ms of its own, so it ages against the table's
1055        // history.expire.max-ref-age-ms (1 day); its head is 10 days old, so the tag is dropped.
1056        // Set the tag before main so the builder's last-updated bookkeeping stays monotonic.
1057        let table = table_with_props(
1058            vec![
1059                snapshot(1, None, 35, now - 10 * day_ms),
1060                snapshot(2, None, 36, now - 1000),
1061            ],
1062            vec![("old-tag", tag(1, None)), (MAIN_BRANCH, branch(2, None))],
1063            HashMap::from([(
1064                "history.expire.max-ref-age-ms".to_string(),
1065                day_ms.to_string(),
1066            )]),
1067        );
1068
1069        let updates = updates_of(&table, action()).await;
1070        assert_eq!(removed_refs(&updates), vec!["old-tag".to_string()]);
1071        // Once the tag is gone, snapshot 1 (unreferenced and well past the default cutoff) expires.
1072        assert!(updates.iter().any(
1073            |u| matches!(u, TableUpdate::RemoveSnapshots { snapshot_ids } if snapshot_ids == &[1])
1074        ));
1075    }
1076
1077    #[tokio::test]
1078    async fn test_expiring_snapshot_drops_its_statistics() {
1079        // main: 1 -> 2, both carrying statistics. retain_last(1) keeps only the head (2).
1080        let table = table_with_stats(
1081            vec![
1082                snapshot(1, None, 35, TS + 1),
1083                snapshot(2, Some(1), 36, TS + 2),
1084            ],
1085            vec![(MAIN_BRANCH, branch(2, None))],
1086            vec![stats_file(1), stats_file(2)],
1087            vec![partition_stats_file(1), partition_stats_file(2)],
1088        );
1089
1090        let updates = updates_of(
1091            &table,
1092            action().retain_last(1).expire_older_than_ms(i64::MAX),
1093        )
1094        .await;
1095
1096        // Snapshot 1 expires, so its stats entries are dropped; the retained head 2 keeps its stats.
1097        assert_eq!(removed_statistics(&updates), vec![1]);
1098        assert_eq!(removed_partition_statistics(&updates), vec![1]);
1099    }
1100
1101    #[tokio::test]
1102    async fn test_expiring_snapshot_without_statistics_emits_no_removal() {
1103        // Same expiry, but no statistics attached to any snapshot.
1104        let table = table_with(
1105            vec![
1106                snapshot(1, None, 35, TS + 1),
1107                snapshot(2, Some(1), 36, TS + 2),
1108            ],
1109            vec![(MAIN_BRANCH, branch(2, None))],
1110        );
1111
1112        let updates = updates_of(
1113            &table,
1114            action().retain_last(1).expire_older_than_ms(i64::MAX),
1115        )
1116        .await;
1117
1118        assert!(removed_statistics(&updates).is_empty());
1119        assert!(removed_partition_statistics(&updates).is_empty());
1120    }
1121
1122    #[tokio::test]
1123    async fn test_only_present_statistics_variant_is_removed() {
1124        // Snapshot 1 has statistics but no partition statistics.
1125        let table = table_with_stats(
1126            vec![
1127                snapshot(1, None, 35, TS + 1),
1128                snapshot(2, Some(1), 36, TS + 2),
1129            ],
1130            vec![(MAIN_BRANCH, branch(2, None))],
1131            vec![stats_file(1)],
1132            vec![],
1133        );
1134
1135        let updates = updates_of(
1136            &table,
1137            action().retain_last(1).expire_older_than_ms(i64::MAX),
1138        )
1139        .await;
1140
1141        assert_eq!(removed_statistics(&updates), vec![1]);
1142        assert!(removed_partition_statistics(&updates).is_empty());
1143    }
1144
1145    #[tokio::test]
1146    async fn test_ref_aging_expiry_drops_statistics() {
1147        let now = Utc::now().timestamp_millis();
1148        let day_ms = 24 * 60 * 60 * 1000;
1149        // Snapshot 1 is kept alive only by `old-tag`; once the tag ages out (1-day max-ref-age vs a
1150        // 10-day-old head) snapshot 1 becomes expirable and its statistics go with it. This reaches
1151        // the stats loop through the ref-aging branch of `plan()`, distinct from the retain/age path.
1152        let table = table_with_stats(
1153            vec![
1154                snapshot(1, None, 35, now - 10 * day_ms),
1155                snapshot(2, None, 36, now - 1000),
1156            ],
1157            vec![
1158                ("old-tag", tag(1, Some(day_ms))),
1159                (MAIN_BRANCH, branch(2, None)),
1160            ],
1161            vec![stats_file(1)],
1162            vec![partition_stats_file(1)],
1163        );
1164
1165        let updates = updates_of(&table, action().expire_older_than_ms(now - 5 * day_ms)).await;
1166        assert_eq!(removed_refs(&updates), vec!["old-tag".to_string()]);
1167        assert_eq!(removed_statistics(&updates), vec![1]);
1168        assert_eq!(removed_partition_statistics(&updates), vec![1]);
1169    }
1170
1171    #[tokio::test]
1172    async fn test_multiple_expired_snapshots_drop_their_statistics() {
1173        // Chain 1 -> 2 -> 3 on main; retain_last(1) keeps only head 3, expiring 1 and 2, both
1174        // carrying statistics.
1175        let table = table_with_stats(
1176            vec![
1177                snapshot(1, None, 35, TS + 1),
1178                snapshot(2, Some(1), 36, TS + 2),
1179                snapshot(3, Some(2), 37, TS + 3),
1180            ],
1181            vec![(MAIN_BRANCH, branch(3, None))],
1182            vec![stats_file(1), stats_file(2)],
1183            vec![partition_stats_file(1), partition_stats_file(2)],
1184        );
1185
1186        let updates = updates_of(
1187            &table,
1188            action().retain_last(1).expire_older_than_ms(i64::MAX),
1189        )
1190        .await;
1191
1192        assert_eq!(removed_statistics(&updates), vec![1, 2]);
1193        assert_eq!(removed_partition_statistics(&updates), vec![1, 2]);
1194    }
1195}