1use 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
32pub 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 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 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 pub fn retain_last(mut self, retain_last: usize) -> Self {
96 self.retain_last = Some(retain_last);
97 self
98 }
99
100 fn plan(&self, table: &Table, properties: &TableProperties<'_>) -> Result<ExpirePlan> {
102 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 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 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 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 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 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 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 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 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
292struct 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 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 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 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 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 const OLD_SNAPSHOT: i64 = 3051729675574597004;
384 const CURRENT_SNAPSHOT: i64 = 3055729675574597004;
385 const TS: i64 = 1_700_000_000_000;
387 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 fn table_with(snapshots: Vec<Snapshot>, refs: Vec<(&str, SnapshotReference)>) -> Table {
470 table_with_props(snapshots, refs, HashMap::new())
471 }
472
473 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 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 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 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 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 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 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 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 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 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 let table = table_with(
797 vec![
798 snapshot(1, None, 35, now - 10 * day_ms),
799 snapshot(2, None, 36, now - 1000),
800 ],
801 vec![
803 ("old-tag", tag(1, Some(day_ms))),
804 (MAIN_BRANCH, branch(2, None)),
805 ],
806 );
807
808 let updates = updates_of(&table, action().expire_older_than_ms(now - 5 * day_ms)).await;
810 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 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 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 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 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 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 let table = table_with(
888 vec![
889 snapshot(3, None, 35, now - 10 * day_ms), snapshot(1, None, 36, now - 1000), snapshot(2, None, 37, now - 1000), ],
893 vec![(MAIN_BRANCH, branch(1, None))],
894 );
895
896 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 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 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 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 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 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), snapshot(2, Some(1), 36, now - 1000), ],
998 vec![(MAIN_BRANCH, branch(2, None))],
999 );
1000
1001 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 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 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 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 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 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 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 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 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 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 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}