Skip to main content

sui_core/authority/
authority_store_pruner.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use super::authority_store_tables::AuthorityPerpetualTables;
5use crate::checkpoints::{CheckpointStore, CheckpointWatermark};
6use crate::jsonrpc_index::IndexStore;
7use anyhow::anyhow;
8use mysten_metrics::monitored_scope;
9#[cfg(not(tidehunter))]
10use mysten_metrics::spawn_monitored_task;
11#[cfg(not(tidehunter))]
12use once_cell::sync::Lazy;
13use prometheus::{
14    IntCounter, IntGauge, Registry, register_int_counter_with_registry,
15    register_int_gauge_with_registry,
16};
17#[cfg(tidehunter)]
18use serde::de::DeserializeOwned;
19#[cfg(not(tidehunter))]
20use std::cmp::max;
21use std::cmp::min;
22#[cfg(not(tidehunter))]
23use std::collections::{BTreeSet, HashMap};
24#[cfg(not(tidehunter))]
25use std::sync::Mutex;
26use std::sync::atomic::AtomicU64;
27#[cfg(not(tidehunter))]
28use std::time::{SystemTime, UNIX_EPOCH};
29use std::{sync::Arc, time::Duration};
30use sui_config::node::AuthorityStorePruningConfig;
31use sui_rpc_store::Store as RpcStore;
32#[cfg(not(tidehunter))]
33use sui_types::base_types::VersionNumber;
34use sui_types::committee::EpochId;
35use sui_types::effects::TransactionEffects;
36use sui_types::effects::TransactionEffectsAPI;
37use sui_types::message_envelope::Message;
38use sui_types::messages_checkpoint::{
39    CheckpointContents, CheckpointDigest, CheckpointSequenceNumber,
40};
41use sui_types::{
42    base_types::{ObjectID, SequenceNumber, TransactionDigest},
43    storage::ObjectKey,
44};
45use tokio::sync::oneshot::{self, Sender};
46use tokio::time::Instant;
47use tracing::{debug, error, info, warn};
48#[cfg(not(tidehunter))]
49use typed_store::rocksdb::LiveFile;
50use typed_store::{Map, TypedStoreError};
51
52#[cfg(not(tidehunter))]
53static PERIODIC_PRUNING_TABLES: Lazy<BTreeSet<String>> = Lazy::new(|| {
54    [
55        "objects",
56        "effects",
57        "transactions",
58        "events",
59        "executed_effects",
60        "executed_transactions_to_checkpoint",
61    ]
62    .into_iter()
63    .map(|cf| cf.to_string())
64    .collect()
65});
66pub const EPOCH_DURATION_MS_FOR_TESTING: u64 = 24 * 60 * 60 * 1000;
67pub struct AuthorityStorePruner {
68    _objects_pruner_cancel_handle: oneshot::Sender<()>,
69}
70
71#[derive(Default)]
72pub struct PrunerWatermarks {
73    pub epoch_id: Arc<AtomicU64>,
74    pub checkpoint_id: Arc<AtomicU64>,
75}
76
77static MIN_PRUNING_TICK_DURATION_MS: u64 = 10 * 1000;
78
79pub struct AuthorityStorePruningMetrics {
80    pub last_pruned_checkpoint: IntGauge,
81    pub num_pruned_objects: IntCounter,
82    pub num_pruned_tombstones: IntCounter,
83    pub last_pruned_effects_checkpoint: IntGauge,
84    pub last_pruned_indexes_transaction: IntGauge,
85    pub num_epochs_to_retain_for_objects: IntGauge,
86    pub num_epochs_to_retain_for_checkpoints: IntGauge,
87}
88
89impl AuthorityStorePruningMetrics {
90    pub fn new(registry: &Registry) -> Arc<Self> {
91        let this = Self {
92            last_pruned_checkpoint: register_int_gauge_with_registry!(
93                "last_pruned_checkpoint",
94                "Last pruned checkpoint",
95                registry
96            )
97            .unwrap(),
98            num_pruned_objects: register_int_counter_with_registry!(
99                "num_pruned_objects",
100                "Number of pruned objects",
101                registry
102            )
103            .unwrap(),
104            num_pruned_tombstones: register_int_counter_with_registry!(
105                "num_pruned_tombstones",
106                "Number of pruned tombstones",
107                registry
108            )
109            .unwrap(),
110            last_pruned_effects_checkpoint: register_int_gauge_with_registry!(
111                "last_pruned_effects_checkpoint",
112                "Last pruned effects checkpoint",
113                registry
114            )
115            .unwrap(),
116            last_pruned_indexes_transaction: register_int_gauge_with_registry!(
117                "last_pruned_indexes_transaction",
118                "Last pruned indexes transaction",
119                registry
120            )
121            .unwrap(),
122            num_epochs_to_retain_for_objects: register_int_gauge_with_registry!(
123                "num_epochs_to_retain_for_objects",
124                "Number of epochs to retain for objects",
125                registry
126            )
127            .unwrap(),
128            num_epochs_to_retain_for_checkpoints: register_int_gauge_with_registry!(
129                "num_epochs_to_retain_for_checkpoints",
130                "Number of epochs to retain for checkpoints",
131                registry
132            )
133            .unwrap(),
134        };
135        Arc::new(this)
136    }
137
138    pub fn new_for_test() -> Arc<Self> {
139        Self::new(&Registry::new())
140    }
141}
142
143#[derive(Debug, Clone, Copy, PartialEq)]
144pub enum PruningMode {
145    Objects,
146    Checkpoints,
147}
148
149impl AuthorityStorePruner {
150    /// prunes old versions of objects based on transaction effects
151    #[cfg(not(tidehunter))]
152    async fn prune_objects_and_indexes(
153        transaction_effects: Vec<(CheckpointSequenceNumber, TransactionEffects)>,
154        perpetual_db: &Arc<AuthorityPerpetualTables>,
155        checkpoint_number: CheckpointSequenceNumber,
156        metrics: Arc<AuthorityStorePruningMetrics>,
157        pruned_tx_seq_exclusive: u64,
158        rpc_store: Option<&RpcStore>,
159        enable_pruning_tombstones: bool,
160    ) -> anyhow::Result<()> {
161        let _scope = monitored_scope("ObjectsLivePruner");
162        let mut wb = perpetual_db.objects.batch();
163
164        // Collect objects keys that need to be deleted from `transaction_effects`.
165        let mut live_object_keys_to_prune = vec![];
166        let mut object_tombstones_to_prune = vec![];
167        for (_checkpoint, effects) in &transaction_effects {
168            for (object_id, seq_number) in effects.modified_at_versions() {
169                live_object_keys_to_prune.push(ObjectKey(object_id, seq_number));
170            }
171
172            if enable_pruning_tombstones {
173                for deleted_object_key in effects.all_tombstones() {
174                    object_tombstones_to_prune
175                        .push(ObjectKey(deleted_object_key.0, deleted_object_key.1));
176                }
177            }
178        }
179
180        metrics
181            .num_pruned_objects
182            .inc_by(live_object_keys_to_prune.len() as u64);
183        metrics
184            .num_pruned_tombstones
185            .inc_by(object_tombstones_to_prune.len() as u64);
186
187        let mut updates: HashMap<ObjectID, (VersionNumber, VersionNumber)> = HashMap::new();
188        for ObjectKey(object_id, seq_number) in live_object_keys_to_prune {
189            updates
190                .entry(object_id)
191                .and_modify(|range| *range = (min(range.0, seq_number), max(range.1, seq_number)))
192                .or_insert((seq_number, seq_number));
193        }
194
195        for (object_id, (min_version, max_version)) in updates {
196            debug!(
197                "Pruning object {:?} versions {:?} - {:?}",
198                object_id, min_version, max_version
199            );
200            let start_range = ObjectKey(object_id, min_version);
201            let end_range = ObjectKey(object_id, (max_version.value() + 1).into());
202            wb.schedule_delete_range(&perpetual_db.objects, &start_range, &end_range)?;
203        }
204
205        // When enable_pruning_tombstones is enabled, instead of using range deletes, we need to do a scan of all the keys
206        // for the deleted objects and then do point deletes to delete all the existing keys. This is because to improve read
207        // performance, we set `ignore_range_deletions` on all read options, and using range delete to delete tombstones
208        // may leak object (imagine a tombstone is compacted away, but earlier version is still not). Using point deletes
209        // guarantees that all earlier versions are deleted in the database.
210        if !object_tombstones_to_prune.is_empty() {
211            let mut object_keys_to_delete = vec![];
212            for ObjectKey(object_id, seq_number) in object_tombstones_to_prune {
213                for result in perpetual_db.objects.safe_iter_with_bounds(
214                    Some(ObjectKey(object_id, VersionNumber::MIN)),
215                    Some(ObjectKey(object_id, seq_number.next())),
216                ) {
217                    let (object_key, _) = result?;
218                    assert_eq!(object_key.0, object_id);
219                    object_keys_to_delete.push(object_key);
220                }
221            }
222
223            wb.delete_batch(&perpetual_db.objects, object_keys_to_delete)?;
224        }
225
226        perpetual_db.set_highest_pruned_checkpoint(&mut wb, checkpoint_number)?;
227        metrics.last_pruned_checkpoint.set(checkpoint_number as i64);
228
229        // Prune the embedded rpc-store's history cohort to the same floor
230        // BEFORE committing the perpetual batch. The rpc-store's
231        // `object_version_by_checkpoint` retraction is driven by the
232        // effects passed here; if the perpetual floor committed first, a
233        // crash between the two commits would resume the pruner past these
234        // checkpoints and their effects would never be replayed, leaking
235        // the rpc-store rows they retract. In this order a crash re-reads
236        // the same checkpoints on restart (the perpetual floor has not
237        // moved) and `prune_history_cohort` re-runs as an idempotent
238        // no-op.
239        if let Some(rpc_store) = rpc_store {
240            sui_rpc_store::prune_history_cohort(
241                rpc_store.db(),
242                rpc_store.schema(),
243                checkpoint_number,
244                pruned_tx_seq_exclusive,
245                &transaction_effects,
246            )?;
247        }
248
249        wb.write()?;
250
251        Ok(())
252    }
253
254    #[cfg(tidehunter)]
255    async fn prune_objects_and_indexes(
256        transaction_effects: Vec<(CheckpointSequenceNumber, TransactionEffects)>,
257        perpetual_db: &Arc<AuthorityPerpetualTables>,
258        checkpoint_number: CheckpointSequenceNumber,
259        metrics: Arc<AuthorityStorePruningMetrics>,
260        pruned_tx_seq_exclusive: u64,
261        rpc_store: Option<&RpcStore>,
262        _: bool,
263    ) -> anyhow::Result<()> {
264        let _scope = monitored_scope("ObjectsLivePruner");
265        let mut wb = perpetual_db.objects.batch();
266        let mut objects_to_prune = vec![];
267
268        for (_checkpoint, effects) in &transaction_effects {
269            for (object_id, version) in effects
270                .modified_at_versions()
271                .into_iter()
272                .chain(effects.all_tombstones())
273            {
274                debug!("Pruning object {:?} version {:?}", object_id, version);
275                objects_to_prune.push(ObjectKey(object_id, version));
276            }
277        }
278        metrics
279            .num_pruned_objects
280            .inc_by(objects_to_prune.len() as u64);
281        wb.delete_batch(&perpetual_db.objects, &objects_to_prune)?;
282
283        perpetual_db.set_highest_pruned_checkpoint(&mut wb, checkpoint_number)?;
284        metrics.last_pruned_checkpoint.set(checkpoint_number as i64);
285
286        // Prune the embedded rpc-store's history cohort BEFORE committing
287        // the perpetual batch; see the rocksdb variant above for the
288        // crash-ordering rationale.
289        if let Some(rpc_store) = rpc_store {
290            sui_rpc_store::prune_history_cohort(
291                rpc_store.db(),
292                rpc_store.schema(),
293                checkpoint_number,
294                pruned_tx_seq_exclusive,
295                &transaction_effects,
296            )?;
297        }
298
299        wb.write()?;
300
301        Ok(())
302    }
303
304    fn prune_checkpoints(
305        perpetual_db: &Arc<AuthorityPerpetualTables>,
306        checkpoint_db: &Arc<CheckpointStore>,
307        checkpoint_number: CheckpointSequenceNumber,
308        checkpoints_to_prune: Vec<CheckpointDigest>,
309        checkpoint_content_to_prune: Vec<CheckpointContents>,
310        effects_to_prune: &Vec<(CheckpointSequenceNumber, TransactionEffects)>,
311        metrics: Arc<AuthorityStorePruningMetrics>,
312    ) -> anyhow::Result<()> {
313        let _scope = monitored_scope("EffectsLivePruner");
314
315        let mut perpetual_batch = perpetual_db.objects.batch();
316        let transactions: Vec<_> = checkpoint_content_to_prune
317            .iter()
318            .flat_map(|content| content.iter().map(|tx| tx.transaction))
319            .collect();
320
321        perpetual_batch.delete_batch(&perpetual_db.transactions, transactions.iter())?;
322        perpetual_batch.delete_batch(&perpetual_db.executed_effects, transactions.iter())?;
323        perpetual_batch.delete_batch(
324            &perpetual_db.executed_transactions_to_checkpoint,
325            transactions.iter(),
326        )?;
327
328        let mut effect_digests = vec![];
329        for (_checkpoint, effects) in effects_to_prune {
330            let effects_digest = effects.digest();
331            debug!("Pruning effects {:?}", effects_digest);
332            effect_digests.push(effects_digest);
333
334            if effects.events_digest().is_some() {
335                perpetual_batch
336                    .delete_batch(&perpetual_db.events_2, [effects.transaction_digest()])?;
337            }
338        }
339        perpetual_batch.delete_batch(
340            &perpetual_db.unchanged_loaded_runtime_objects,
341            transactions.iter(),
342        )?;
343        perpetual_batch.delete_batch(&perpetual_db.effects, effect_digests)?;
344
345        let mut checkpoints_batch = checkpoint_db.tables.certified_checkpoints.batch();
346
347        let checkpoint_content_digests =
348            checkpoint_content_to_prune.iter().map(|ckpt| ckpt.digest());
349        checkpoints_batch.delete_batch(
350            &checkpoint_db.tables.checkpoint_content,
351            checkpoint_content_digests.clone(),
352        )?;
353        checkpoints_batch.delete_batch(
354            &checkpoint_db.tables.checkpoint_sequence_by_contents_digest,
355            checkpoint_content_digests,
356        )?;
357
358        checkpoints_batch.delete_batch(
359            &checkpoint_db.tables.checkpoint_by_digest,
360            checkpoints_to_prune,
361        )?;
362
363        checkpoints_batch.insert_batch(
364            &checkpoint_db.tables.watermarks,
365            [(
366                &CheckpointWatermark::HighestPruned,
367                &(checkpoint_number, CheckpointDigest::random()),
368            )],
369        )?;
370
371        perpetual_batch.write()?;
372        checkpoints_batch.write()?;
373        metrics
374            .last_pruned_effects_checkpoint
375            .set(checkpoint_number as i64);
376
377        Ok(())
378    }
379
380    /// The exclusive upper bound the embedded rpc-store imposes on the
381    /// pruner's eligible range: one past the highest checkpoint every
382    /// embedded-cohort pipeline has committed
383    /// ([`sui_rpc_store::embedded_prunable_checkpoint`]). The indexer's
384    /// backfill and gap-fill assemble full checkpoints from the
385    /// perpetual and checkpoint stores, so pruning a checkpoint's data
386    /// (contents, effects, or the object versions its transactions
387    /// read) before every pipeline has committed it would stall the
388    /// indexer permanently on a checkpoint that can no longer be
389    /// served. `u64::MAX` (no bound) when no embedded store is
390    /// configured; `0` (nothing eligible) while any cohort pipeline
391    /// has no watermark yet.
392    fn rpc_store_max_eligible_checkpoint(
393        rpc_store: Option<&RpcStore>,
394    ) -> anyhow::Result<CheckpointSequenceNumber> {
395        let Some(rpc_store) = rpc_store else {
396            return Ok(u64::MAX);
397        };
398        let indexed = sui_rpc_store::embedded_prunable_checkpoint(rpc_store.db())?;
399        // The eligible bound is exclusive, so `indexed + 1` still
400        // allows pruning the committed checkpoint itself.
401        Ok(indexed.map_or(0, |c| c.saturating_add(1)))
402    }
403
404    /// Prunes old data based on effects from all checkpoints from epochs eligible for pruning
405    pub async fn prune_objects_for_eligible_epochs(
406        perpetual_db: &Arc<AuthorityPerpetualTables>,
407        checkpoint_store: &Arc<CheckpointStore>,
408        rpc_store: Option<&RpcStore>,
409        config: AuthorityStorePruningConfig,
410        metrics: Arc<AuthorityStorePruningMetrics>,
411        epoch_duration_ms: u64,
412    ) -> anyhow::Result<()> {
413        let _scope = monitored_scope("PruneObjectsForEligibleEpochs");
414        let (mut max_eligible_checkpoint_number, epoch_id) = checkpoint_store
415            .get_highest_executed_checkpoint()?
416            .map(|c| (*c.sequence_number(), c.epoch))
417            .unwrap_or_default();
418        let pruned_checkpoint_number = perpetual_db
419            .get_highest_pruned_checkpoint()?
420            .unwrap_or_default();
421        if config.smooth && config.num_epochs_to_retain > 0 {
422            max_eligible_checkpoint_number = Self::smoothed_max_eligible_checkpoint_number(
423                checkpoint_store,
424                max_eligible_checkpoint_number,
425                pruned_checkpoint_number,
426                epoch_id,
427                epoch_duration_ms,
428                config.num_epochs_to_retain,
429            )?;
430        }
431        let rpc_store_bound = Self::rpc_store_max_eligible_checkpoint(rpc_store)?;
432        if rpc_store_bound < max_eligible_checkpoint_number {
433            info!(
434                "objects pruning gated by the embedded rpc-store indexer: \
435                 max eligible checkpoint {} -> {}",
436                max_eligible_checkpoint_number, rpc_store_bound,
437            );
438            max_eligible_checkpoint_number = rpc_store_bound;
439        }
440        Self::prune_for_eligible_epochs(
441            perpetual_db,
442            checkpoint_store,
443            rpc_store,
444            PruningMode::Objects,
445            config.num_epochs_to_retain,
446            pruned_checkpoint_number,
447            max_eligible_checkpoint_number,
448            config,
449            metrics.clone(),
450        )
451        .await
452    }
453
454    pub async fn prune_checkpoints_for_eligible_epochs(
455        perpetual_db: &Arc<AuthorityPerpetualTables>,
456        checkpoint_store: &Arc<CheckpointStore>,
457        rpc_store: Option<&RpcStore>,
458        config: AuthorityStorePruningConfig,
459        metrics: Arc<AuthorityStorePruningMetrics>,
460        epoch_duration_ms: u64,
461        pruner_watermarks: &Arc<PrunerWatermarks>,
462    ) -> anyhow::Result<()> {
463        let _scope = monitored_scope("PruneCheckpointsForEligibleEpochs");
464        let pruned_checkpoint_number = checkpoint_store
465            .get_highest_pruned_checkpoint_seq_number()?
466            .unwrap_or(0);
467        let (mut max_eligible_checkpoint, epoch_id) = checkpoint_store
468            .get_highest_executed_checkpoint()?
469            .map(|c| (*c.sequence_number(), c.epoch))
470            .unwrap_or_default();
471        if config.num_epochs_to_retain != u64::MAX {
472            max_eligible_checkpoint = min(
473                max_eligible_checkpoint,
474                perpetual_db
475                    .get_highest_pruned_checkpoint()?
476                    .unwrap_or_default(),
477            );
478        }
479        if config.smooth
480            && let Some(num_epochs_to_retain) = config.num_epochs_to_retain_for_checkpoints
481        {
482            max_eligible_checkpoint = Self::smoothed_max_eligible_checkpoint_number(
483                checkpoint_store,
484                max_eligible_checkpoint,
485                pruned_checkpoint_number,
486                epoch_id,
487                epoch_duration_ms,
488                num_epochs_to_retain,
489            )?;
490        }
491        // With `num_epochs_to_retain == u64::MAX` the objects floor
492        // never advances, so the clamp above does not apply and the
493        // rpc-store bound is the only thing keeping checkpoint contents
494        // around for the indexer; apply it in both cases regardless.
495        let rpc_store_bound = Self::rpc_store_max_eligible_checkpoint(rpc_store)?;
496        if rpc_store_bound < max_eligible_checkpoint {
497            info!(
498                "checkpoint pruning gated by the embedded rpc-store indexer: \
499                 max eligible checkpoint {} -> {}",
500                max_eligible_checkpoint, rpc_store_bound,
501            );
502            max_eligible_checkpoint = rpc_store_bound;
503        }
504        debug!("Max eligible checkpoint {}", max_eligible_checkpoint);
505        Self::prune_for_eligible_epochs(
506            perpetual_db,
507            checkpoint_store,
508            rpc_store,
509            PruningMode::Checkpoints,
510            config
511                .num_epochs_to_retain_for_checkpoints()
512                .ok_or_else(|| anyhow!("config value not set"))?,
513            pruned_checkpoint_number,
514            max_eligible_checkpoint,
515            config.clone(),
516            metrics.clone(),
517        )
518        .await?;
519
520        if let Some(num_epochs_to_retain) = config.num_epochs_to_retain_for_checkpoints() {
521            Self::update_pruning_watermarks(
522                perpetual_db,
523                checkpoint_store,
524                num_epochs_to_retain,
525                pruner_watermarks,
526                false,
527            )?;
528        }
529        Ok(())
530    }
531
532    /// Prunes old object versions based on effects from all checkpoints from epochs eligible for pruning
533    pub async fn prune_for_eligible_epochs(
534        perpetual_db: &Arc<AuthorityPerpetualTables>,
535        checkpoint_store: &Arc<CheckpointStore>,
536        rpc_store: Option<&RpcStore>,
537        mode: PruningMode,
538        num_epochs_to_retain: u64,
539        starting_checkpoint_number: CheckpointSequenceNumber,
540        max_eligible_checkpoint: CheckpointSequenceNumber,
541        config: AuthorityStorePruningConfig,
542        metrics: Arc<AuthorityStorePruningMetrics>,
543    ) -> anyhow::Result<()> {
544        let _scope = monitored_scope("PruneForEligibleEpochs");
545
546        let mut checkpoint_number = starting_checkpoint_number;
547        let current_epoch = checkpoint_store
548            .get_highest_executed_checkpoint()?
549            .map(|c| c.epoch())
550            .unwrap_or_default();
551
552        let mut checkpoints_to_prune = vec![];
553        let mut checkpoint_content_to_prune = vec![];
554        // Each effect is tagged with the checkpoint it is pruned from so the
555        // embedded rpc-store's `object_version_by_checkpoint` retraction can
556        // keep each object's anchor at its true supersession checkpoint.
557        let mut effects_to_prune: Vec<(CheckpointSequenceNumber, TransactionEffects)> = vec![];
558        // Absolute tx-seq floor (exclusive) after pruning the current
559        // batch — the last-pruned checkpoint's `network_total_transactions`.
560        // The embedded rpc-store's history-cohort prune consumes this
561        // directly instead of summing checkpoint content sizes.
562        let mut pruned_tx_seq_exclusive = 0u64;
563
564        while let Some(ckpt) = checkpoint_store
565            .tables
566            .certified_checkpoints
567            .get(&(checkpoint_number + 1))?
568        {
569            let checkpoint = ckpt.into_inner();
570            // Skipping because  checkpoint's epoch or checkpoint number is too new.
571            // We have to respect the highest executed checkpoint watermark (including the watermark itself)
572            // because there might be parts of the system that still require access to old object versions
573            // (i.e. state accumulator).
574            if (current_epoch < checkpoint.epoch() + num_epochs_to_retain)
575                || (*checkpoint.sequence_number() >= max_eligible_checkpoint)
576            {
577                break;
578            }
579            checkpoint_number = *checkpoint.sequence_number();
580            pruned_tx_seq_exclusive = checkpoint.network_total_transactions;
581
582            let content = checkpoint_store
583                .get_checkpoint_contents(&checkpoint.content_digest)?
584                .ok_or_else(|| {
585                    anyhow::anyhow!(
586                        "checkpoint content data is missing: {}",
587                        checkpoint.sequence_number
588                    )
589                })?;
590            let effects = perpetual_db
591                .effects
592                .multi_get(content.iter().map(|tx| tx.effects))?;
593
594            info!("scheduling pruning for checkpoint {:?}", checkpoint_number);
595            checkpoints_to_prune.push(*checkpoint.digest());
596            checkpoint_content_to_prune.push(content);
597            effects_to_prune.extend(
598                effects
599                    .into_iter()
600                    .flatten()
601                    .map(|effects| (checkpoint_number, effects)),
602            );
603
604            if effects_to_prune.len() >= config.max_transactions_in_batch
605                || checkpoints_to_prune.len() >= config.max_checkpoints_in_batch
606            {
607                match mode {
608                    PruningMode::Objects => {
609                        Self::prune_objects_and_indexes(
610                            effects_to_prune,
611                            perpetual_db,
612                            checkpoint_number,
613                            metrics.clone(),
614                            pruned_tx_seq_exclusive,
615                            rpc_store,
616                            !config.killswitch_tombstone_pruning,
617                        )
618                        .await?
619                    }
620                    PruningMode::Checkpoints => Self::prune_checkpoints(
621                        perpetual_db,
622                        checkpoint_store,
623                        checkpoint_number,
624                        checkpoints_to_prune,
625                        checkpoint_content_to_prune,
626                        &effects_to_prune,
627                        metrics.clone(),
628                    )?,
629                };
630                checkpoints_to_prune = vec![];
631                checkpoint_content_to_prune = vec![];
632                effects_to_prune = vec![];
633                // yield back to the tokio runtime. Prevent potential halt of other tasks
634                tokio::task::yield_now().await;
635            }
636        }
637
638        if !checkpoints_to_prune.is_empty() {
639            match mode {
640                PruningMode::Objects => {
641                    Self::prune_objects_and_indexes(
642                        effects_to_prune,
643                        perpetual_db,
644                        checkpoint_number,
645                        metrics.clone(),
646                        pruned_tx_seq_exclusive,
647                        rpc_store,
648                        !config.killswitch_tombstone_pruning,
649                    )
650                    .await?
651                }
652                PruningMode::Checkpoints => Self::prune_checkpoints(
653                    perpetual_db,
654                    checkpoint_store,
655                    checkpoint_number,
656                    checkpoints_to_prune,
657                    checkpoint_content_to_prune,
658                    &effects_to_prune,
659                    metrics.clone(),
660                )?,
661            };
662        }
663        Ok(())
664    }
665
666    #[cfg(not(tidehunter))]
667    fn prune_indexes(
668        indexes: Option<&IndexStore>,
669        config: &AuthorityStorePruningConfig,
670        epoch_duration_ms: u64,
671        metrics: &AuthorityStorePruningMetrics,
672    ) -> anyhow::Result<()> {
673        if let (Some(mut epochs_to_retain), Some(indexes)) =
674            (config.num_epochs_to_retain_for_indexes, indexes)
675        {
676            if epochs_to_retain < 7 {
677                warn!("num_epochs_to_retain_for_indexes is too low. Reseting it to 7");
678                epochs_to_retain = 7;
679            }
680            let now = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis();
681            if let Some(cut_time_ms) =
682                u64::try_from(now)?.checked_sub(epochs_to_retain * epoch_duration_ms)
683            {
684                let transaction_id = indexes.prune(cut_time_ms)?;
685                metrics
686                    .last_pruned_indexes_transaction
687                    .set(transaction_id as i64);
688            }
689        }
690        Ok(())
691    }
692
693    #[cfg(not(tidehunter))]
694    async fn prune_executed_tx_digests(
695        perpetual_db: &Arc<AuthorityPerpetualTables>,
696        checkpoint_store: &Arc<CheckpointStore>,
697    ) -> anyhow::Result<()> {
698        let current_epoch = checkpoint_store
699            .get_highest_executed_checkpoint()?
700            .map(|c| c.epoch)
701            .unwrap_or_default();
702
703        if current_epoch < 2 {
704            return Ok(());
705        }
706
707        let target_epoch = current_epoch - 1;
708
709        let start_key = (0u64, TransactionDigest::ZERO);
710        let end_key = (target_epoch, TransactionDigest::ZERO);
711
712        info!(
713            "Pruning executed_transaction_digests for epochs < {} (current epoch: {})",
714            target_epoch, current_epoch
715        );
716
717        let mut batch = perpetual_db.executed_transaction_digests.batch();
718        batch.schedule_delete_range(
719            &perpetual_db.executed_transaction_digests,
720            &start_key,
721            // `to` is non-inclusive so target_epoch and all later epochs are preserved
722            &end_key,
723        )?;
724        batch.write()?;
725        Ok(())
726    }
727
728    #[cfg(tidehunter)]
729    fn prune_executed_tx_digests_th(
730        perpetual_db: &Arc<AuthorityPerpetualTables>,
731        checkpoint_store: &Arc<CheckpointStore>,
732    ) -> anyhow::Result<()> {
733        let current_epoch = checkpoint_store
734            .get_highest_executed_checkpoint()?
735            .map(|c| c.epoch)
736            .unwrap_or_default();
737
738        if current_epoch < 2 {
739            return Ok(());
740        }
741
742        let last_epoch_to_delete = current_epoch - 2;
743        let from_key = (0u64, TransactionDigest::ZERO);
744        let to_key = (last_epoch_to_delete, TransactionDigest::new([0xff; 32]));
745        info!(
746            "Pruning executed_transaction_digests for epochs 0 to {} (current epoch: {})",
747            last_epoch_to_delete, current_epoch
748        );
749        perpetual_db
750            .executed_transaction_digests
751            .drop_cells_in_range(&from_key, &to_key)?;
752        Ok(())
753    }
754
755    fn update_pruning_watermarks(
756        perpetual_db: &Arc<AuthorityPerpetualTables>,
757        checkpoint_store: &Arc<CheckpointStore>,
758        num_epochs_to_retain: u64,
759        pruning_watermark: &Arc<PrunerWatermarks>,
760        objects_compactor_active: bool,
761    ) -> anyhow::Result<bool> {
762        use std::sync::atomic::Ordering;
763        let objects_pruning_checkpoint_id = perpetual_db
764            .get_highest_pruned_checkpoint()?
765            .unwrap_or_default();
766        let objects_pruning_epoch_id = checkpoint_store
767            .get_checkpoint_by_sequence_number(objects_pruning_checkpoint_id)?
768            .map(|chk| chk.epoch)
769            .unwrap_or_default();
770
771        let current_watermark = pruning_watermark.epoch_id.load(Ordering::Relaxed);
772        let current_epoch_id = checkpoint_store
773            .get_highest_executed_checkpoint()?
774            .map(|c| c.epoch)
775            .unwrap_or_default();
776        if current_epoch_id < num_epochs_to_retain {
777            return Ok(false);
778        }
779        let target_epoch_id = current_epoch_id - num_epochs_to_retain;
780        let checkpoint_id =
781            checkpoint_store.get_epoch_last_checkpoint_seq_number(target_epoch_id)?;
782
783        // The objects compactor handles object retention continuously without advancing
784        // `highest_pruned_checkpoint`, so capping on it would freeze the watermark at 0.
785        let new_watermark = if objects_compactor_active {
786            target_epoch_id + 1
787        } else {
788            min(target_epoch_id + 1, objects_pruning_epoch_id)
789        };
790        if current_watermark == new_watermark {
791            return Ok(false);
792        }
793        info!("relocation: setting epoch watermark to {}", new_watermark);
794        pruning_watermark
795            .epoch_id
796            .store(new_watermark, Ordering::Relaxed);
797        if let Some(checkpoint_id) = checkpoint_id {
798            let watermark = if objects_compactor_active {
799                checkpoint_id
800            } else {
801                min(checkpoint_id, objects_pruning_checkpoint_id)
802            };
803            info!("relocation: setting checkpoint watermark to {}", watermark);
804            pruning_watermark
805                .checkpoint_id
806                .store(watermark, Ordering::Relaxed);
807        }
808        Ok(true)
809    }
810
811    #[cfg(tidehunter)]
812    fn prune_th(
813        perpetual_db: &Arc<AuthorityPerpetualTables>,
814        checkpoint_store: &Arc<CheckpointStore>,
815        num_epochs_to_retain: u64,
816        pruning_watermark: Arc<PrunerWatermarks>,
817        objects_compactor_active: bool,
818    ) -> anyhow::Result<()> {
819        let watermark_updated = Self::update_pruning_watermarks(
820            perpetual_db,
821            checkpoint_store,
822            num_epochs_to_retain,
823            &pruning_watermark,
824            objects_compactor_active,
825        )?;
826        if !watermark_updated {
827            info!("skip relocation. Watermark hasn't changed");
828            return Ok(());
829        }
830        perpetual_db.objects.db.start_relocation()?;
831        checkpoint_store.tables.watermarks.db.start_relocation()?;
832        Self::prune_executed_tx_digests_th(perpetual_db, checkpoint_store)?;
833        Ok(())
834    }
835
836    #[cfg(not(tidehunter))]
837    fn compact_next_sst_file(
838        perpetual_db: Arc<AuthorityPerpetualTables>,
839        delay_days: usize,
840        last_processed: Arc<Mutex<HashMap<String, SystemTime>>>,
841    ) -> anyhow::Result<Option<LiveFile>> {
842        let db_path = perpetual_db.objects.db.path_for_pruning();
843        let mut state = last_processed
844            .lock()
845            .expect("failed to obtain a lock for last processed SST files");
846        let mut sst_file_for_compaction: Option<LiveFile> = None;
847        let time_threshold =
848            SystemTime::now() - Duration::from_secs(delay_days as u64 * 24 * 60 * 60);
849        for sst_file in perpetual_db.objects.db.live_files()? {
850            let file_path = db_path.join(sst_file.name.clone().trim_matches('/'));
851            let last_modified = std::fs::metadata(file_path)?.modified()?;
852            if !PERIODIC_PRUNING_TABLES.contains(&sst_file.column_family_name)
853                || sst_file.level < 1
854                || sst_file.start_key.is_none()
855                || sst_file.end_key.is_none()
856                || last_modified > time_threshold
857                || state.get(&sst_file.name).unwrap_or(&UNIX_EPOCH) > &time_threshold
858            {
859                continue;
860            }
861            if let Some(candidate) = &sst_file_for_compaction
862                && candidate.size > sst_file.size
863            {
864                continue;
865            }
866            sst_file_for_compaction = Some(sst_file);
867        }
868        let Some(sst_file) = sst_file_for_compaction else {
869            return Ok(None);
870        };
871        info!(
872            "Manual compaction of sst file {:?}. Size: {:?}, level: {:?}",
873            sst_file.name, sst_file.size, sst_file.level
874        );
875        perpetual_db.objects.compact_range_raw(
876            &sst_file.column_family_name,
877            sst_file.start_key.clone().unwrap(),
878            sst_file.end_key.clone().unwrap(),
879        )?;
880        state.insert(sst_file.name.clone(), SystemTime::now());
881        Ok(Some(sst_file))
882    }
883
884    fn pruning_tick_duration_ms(epoch_duration_ms: u64) -> u64 {
885        min(epoch_duration_ms / 2, MIN_PRUNING_TICK_DURATION_MS)
886    }
887
888    fn smoothed_max_eligible_checkpoint_number(
889        checkpoint_store: &Arc<CheckpointStore>,
890        mut max_eligible_checkpoint: CheckpointSequenceNumber,
891        pruned_checkpoint: CheckpointSequenceNumber,
892        epoch_id: EpochId,
893        epoch_duration_ms: u64,
894        num_epochs_to_retain: u64,
895    ) -> anyhow::Result<CheckpointSequenceNumber> {
896        if epoch_id < num_epochs_to_retain {
897            return Ok(0);
898        }
899        let last_checkpoint_in_epoch = checkpoint_store
900            .get_epoch_last_checkpoint(epoch_id - num_epochs_to_retain)?
901            .map(|checkpoint| checkpoint.sequence_number)
902            .unwrap_or_default();
903        max_eligible_checkpoint = max_eligible_checkpoint.min(last_checkpoint_in_epoch);
904        if max_eligible_checkpoint == 0 {
905            return Ok(max_eligible_checkpoint);
906        }
907        let num_intervals = epoch_duration_ms
908            .checked_div(Self::pruning_tick_duration_ms(epoch_duration_ms))
909            .unwrap_or(1);
910        let delta = max_eligible_checkpoint
911            .saturating_sub(pruned_checkpoint)
912            .checked_div(num_intervals)
913            .unwrap_or(1);
914        Ok(pruned_checkpoint + delta)
915    }
916
917    fn setup_pruning(
918        config: AuthorityStorePruningConfig,
919        epoch_duration_ms: u64,
920        perpetual_db: Arc<AuthorityPerpetualTables>,
921        checkpoint_store: Arc<CheckpointStore>,
922        rpc_store: Option<RpcStore>,
923        jsonrpc_index: Option<Arc<IndexStore>>,
924        metrics: Arc<AuthorityStorePruningMetrics>,
925        pruner_watermarks: Arc<PrunerWatermarks>,
926    ) -> Sender<()> {
927        let (sender, mut recv) = tokio::sync::oneshot::channel();
928        debug!(
929            "Starting object pruning service with num_epochs_to_retain={}",
930            config.num_epochs_to_retain
931        );
932
933        let tick_duration =
934            Duration::from_millis(Self::pruning_tick_duration_ms(epoch_duration_ms));
935        let pruning_initial_delay = if cfg!(msim) {
936            Duration::from_millis(1)
937        } else {
938            Duration::from_secs(config.pruning_run_delay_seconds.unwrap_or(60 * 60))
939        };
940        let mut objects_prune_interval =
941            tokio::time::interval_at(Instant::now() + pruning_initial_delay, tick_duration);
942        let mut checkpoints_prune_interval =
943            tokio::time::interval_at(Instant::now() + pruning_initial_delay, tick_duration);
944
945        metrics
946            .num_epochs_to_retain_for_objects
947            .set(config.num_epochs_to_retain as i64);
948        metrics.num_epochs_to_retain_for_checkpoints.set(
949            config
950                .num_epochs_to_retain_for_checkpoints
951                .unwrap_or_default() as i64,
952        );
953
954        #[cfg(tidehunter)]
955        {
956            // Index pruning is only implemented for the rocksdb backend.
957            let _ = jsonrpc_index;
958            if let Some(num_epochs_to_retain) = config.num_epochs_to_retain_for_checkpoints() {
959                let prune_objects = config.num_epochs_to_retain != u64::MAX;
960                tokio::task::spawn(async move {
961                    loop {
962                        tokio::select! {
963                            _ = objects_prune_interval.tick(), if prune_objects => {
964                                if let Err(err) = Self::prune_objects_for_eligible_epochs(&perpetual_db, &checkpoint_store, rpc_store.as_ref(), config.clone(), metrics.clone(), epoch_duration_ms).await {
965                                    error!("Failed to prune objects: {:?}", err);
966                                }
967                            },
968                            _ = checkpoints_prune_interval.tick() => {
969                                if let Err(err) = Self::prune_th(&perpetual_db, &checkpoint_store, num_epochs_to_retain, pruner_watermarks.clone(), !prune_objects) {
970                                    error!("Failed to prune checkpoints: {:?}", err);
971                                }
972                            },
973                            _ = &mut recv => break,
974                        }
975                    }
976                });
977            }
978        }
979        #[cfg(not(tidehunter))]
980        {
981            let mut indexes_prune_interval =
982                tokio::time::interval_at(Instant::now() + pruning_initial_delay, tick_duration);
983
984            let perpetual_db_for_compaction = perpetual_db.clone();
985            if let Some(delay_days) = config.periodic_compaction_threshold_days {
986                spawn_monitored_task!(async move {
987                    let last_processed = Arc::new(Mutex::new(HashMap::new()));
988                    loop {
989                        let db = perpetual_db_for_compaction.clone();
990                        let state = Arc::clone(&last_processed);
991                        let result = tokio::task::spawn_blocking(move || {
992                            Self::compact_next_sst_file(db, delay_days, state)
993                        })
994                        .await;
995                        let mut sleep_interval_secs = 1;
996                        match result {
997                            Err(err) => error!("Failed to compact sst file: {:?}", err),
998                            Ok(Err(err)) => error!("Failed to compact sst file: {:?}", err),
999                            Ok(Ok(None)) => {
1000                                sleep_interval_secs = 3600;
1001                            }
1002                            _ => {}
1003                        }
1004                        tokio::time::sleep(Duration::from_secs(sleep_interval_secs)).await;
1005                    }
1006                });
1007            }
1008            tokio::task::spawn(async move {
1009                loop {
1010                    tokio::select! {
1011                        _ = objects_prune_interval.tick(), if config.num_epochs_to_retain != u64::MAX => {
1012                            if let Err(err) = Self::prune_objects_for_eligible_epochs(&perpetual_db, &checkpoint_store, rpc_store.as_ref(), config.clone(), metrics.clone(), epoch_duration_ms).await {
1013                                error!("Failed to prune objects: {:?}", err);
1014                            }
1015                            if let Err(err) = Self::prune_executed_tx_digests(&perpetual_db, &checkpoint_store).await {
1016                                error!("Failed to prune executed_tx_digests: {:?}", err);
1017                            }
1018                        },
1019                        _ = checkpoints_prune_interval.tick(), if !matches!(config.num_epochs_to_retain_for_checkpoints(), None | Some(u64::MAX) | Some(0)) => {
1020                            if let Err(err) = Self::prune_checkpoints_for_eligible_epochs(&perpetual_db, &checkpoint_store, rpc_store.as_ref(), config.clone(), metrics.clone(), epoch_duration_ms, &pruner_watermarks).await {
1021                                error!("Failed to prune checkpoints: {:?}", err);
1022                            }
1023                        },
1024                        _ = indexes_prune_interval.tick(), if config.num_epochs_to_retain_for_indexes.is_some() => {
1025                            if let Err(err) = Self::prune_indexes(jsonrpc_index.as_deref(), &config, epoch_duration_ms, &metrics) {
1026                                error!("Failed to prune indexes: {:?}", err);
1027                            }
1028                        }
1029                        _ = &mut recv => break,
1030                    }
1031                }
1032            });
1033        }
1034        sender
1035    }
1036
1037    pub fn new(
1038        perpetual_db: Arc<AuthorityPerpetualTables>,
1039        checkpoint_store: Arc<CheckpointStore>,
1040        rpc_store: Option<RpcStore>,
1041        jsonrpc_index: Option<Arc<IndexStore>>,
1042        mut pruning_config: AuthorityStorePruningConfig,
1043        is_validator: bool,
1044        epoch_duration_ms: u64,
1045        registry: &Registry,
1046        pruner_watermarks: Arc<PrunerWatermarks>, // used by tidehunter relocation filters
1047    ) -> Self {
1048        // On tidehunter, the per-keyspace `objects_compactor`
1049        // (see `AuthorityPerpetualTables::open`) already retains only the latest
1050        // version per ObjectID during compaction, so running the object pruner
1051        // would duplicate that work. The compactor is enabled for validators
1052        // and for any node configured with `num_epochs_to_retain = 0` — keep
1053        // this predicate in sync with the one in `sui-node/src/lib.rs`.
1054        // Force-disable the pruner whenever the compactor is active.
1055        #[cfg(tidehunter)]
1056        {
1057            let objects_compactor_enabled =
1058                is_validator || pruning_config.num_epochs_to_retain == 0;
1059            if objects_compactor_enabled && pruning_config.num_epochs_to_retain != u64::MAX {
1060                info!(
1061                    "Tidehunter: disabling object pruner (was num_epochs_to_retain={}). The objects compactor performs equivalent compaction.",
1062                    pruning_config.num_epochs_to_retain
1063                );
1064                pruning_config.num_epochs_to_retain = u64::MAX;
1065            }
1066        }
1067
1068        if pruning_config.num_epochs_to_retain > 0 && pruning_config.num_epochs_to_retain < u64::MAX
1069        {
1070            warn!(
1071                "Using objects pruner with num_epochs_to_retain = {} can lead to performance issues",
1072                pruning_config.num_epochs_to_retain
1073            );
1074            if is_validator {
1075                warn!("Resetting to aggressive pruner.");
1076                pruning_config.num_epochs_to_retain = 0;
1077            } else {
1078                warn!("Consider using an aggressive pruner (num_epochs_to_retain = 0)");
1079            }
1080        }
1081        AuthorityStorePruner {
1082            _objects_pruner_cancel_handle: Self::setup_pruning(
1083                pruning_config,
1084                epoch_duration_ms,
1085                perpetual_db,
1086                checkpoint_store,
1087                rpc_store,
1088                jsonrpc_index,
1089                AuthorityStorePruningMetrics::new(registry),
1090                pruner_watermarks,
1091            ),
1092        }
1093    }
1094
1095    pub fn compact(perpetual_db: &Arc<AuthorityPerpetualTables>) -> Result<(), TypedStoreError> {
1096        perpetual_db.objects.compact_range(
1097            &ObjectKey(ObjectID::ZERO, SequenceNumber::MIN),
1098            &ObjectKey(ObjectID::MAX, SequenceNumber::MAX),
1099        )
1100    }
1101}
1102
1103#[cfg(tidehunter)]
1104pub(crate) fn apply_relocation_filter<T: DeserializeOwned>(
1105    config: typed_store::tidehunter_util::KeySpaceConfig,
1106    pruner_watermark: Arc<AtomicU64>,
1107    extractor: impl Fn(T) -> u64 + Send + Sync + 'static,
1108    by_key: bool,
1109) -> typed_store::tidehunter_util::KeySpaceConfig {
1110    use bincode::Options;
1111    use std::sync::atomic::Ordering;
1112    use typed_store::tidehunter_util::Decision;
1113    config.with_relocation_filter(move |key, value| {
1114        let data = if by_key {
1115            bincode::DefaultOptions::new()
1116                .with_big_endian()
1117                .with_fixint_encoding()
1118                .deserialize(key)
1119                .expect("relocation filter deserialization error")
1120        } else {
1121            bcs::from_bytes(value).expect("relocation filter deserialization error")
1122        };
1123        if extractor(data) < pruner_watermark.load(Ordering::Relaxed) {
1124            Decision::Remove
1125        } else {
1126            Decision::StopRelocation
1127        }
1128    })
1129}
1130
1131#[cfg(test)]
1132mod tests {
1133    use more_asserts as ma;
1134    #[cfg(not(tidehunter))]
1135    use std::collections::HashSet;
1136    use std::path::Path;
1137    use std::sync::Arc;
1138    #[cfg(not(tidehunter))]
1139    use std::time::Duration;
1140    use tracing::log::info;
1141
1142    use crate::authority::authority_store_pruner::AuthorityStorePruningMetrics;
1143    use crate::authority::authority_store_tables::AuthorityPerpetualTables;
1144    use crate::authority::authority_store_types::get_store_object;
1145    #[cfg(not(tidehunter))]
1146    use crate::authority::authority_store_types::{StoreObject, StoreObjectWrapper};
1147    use prometheus::Registry;
1148    use sui_types::base_types::ObjectDigest;
1149    use sui_types::effects::TransactionEffects;
1150    use sui_types::effects::TransactionEffectsAPI;
1151    use sui_types::{
1152        base_types::{ObjectID, SequenceNumber},
1153        object::Object,
1154        storage::ObjectKey,
1155    };
1156    use typed_store::Map;
1157    #[cfg(not(tidehunter))]
1158    use typed_store::rocks::{DBMap, MetricConf, ReadWriteOptions, default_db_options};
1159
1160    use super::AuthorityStorePruner;
1161
1162    /// The embedded rpc-store gate: no bound without a store, nothing
1163    /// eligible while any cohort pipeline is unwatermarked, and one
1164    /// past the slowest pipeline's watermark otherwise.
1165    #[test]
1166    fn rpc_store_gate_bounds_eligible_checkpoints() {
1167        use sui_consistent_store::Db;
1168        use sui_consistent_store::DbOptions;
1169        use sui_consistent_store::FrameworkSchema;
1170        use sui_consistent_store::PipelineTaskKey;
1171        use sui_consistent_store::Watermark;
1172        use sui_rpc_store::HISTORY_COHORT;
1173        use sui_rpc_store::LIVE_COHORT;
1174        use sui_rpc_store::RpcStoreSchema;
1175
1176        // No embedded store configured: pruning is unbounded.
1177        assert_eq!(
1178            AuthorityStorePruner::rpc_store_max_eligible_checkpoint(None).unwrap(),
1179            u64::MAX,
1180        );
1181
1182        let dir = tempfile::tempdir().unwrap();
1183        let (db, schema) = Db::open::<RpcStoreSchema>(dir.path(), DbOptions::default()).unwrap();
1184        let store = sui_rpc_store::Store::new(db.clone(), Arc::new(schema));
1185
1186        // Cohort pipelines with no watermarks yet (a from-genesis
1187        // build before its first commit): nothing is eligible.
1188        assert_eq!(
1189            AuthorityStorePruner::rpc_store_max_eligible_checkpoint(Some(&store)).unwrap(),
1190            0,
1191        );
1192
1193        // Every cohort pipeline committed through 41 except one
1194        // history pipeline lagging at 7: pruning may cover the
1195        // straggler's committed range [0, 7] and no further.
1196        let framework = FrameworkSchema::new(db.clone());
1197        let mut batch = db.batch();
1198        for name in LIVE_COHORT.iter().chain(HISTORY_COHORT) {
1199            batch
1200                .put(
1201                    &framework.watermarks,
1202                    &PipelineTaskKey::new(*name),
1203                    &Watermark::for_checkpoint(41),
1204                )
1205                .unwrap();
1206        }
1207        batch
1208            .put(
1209                &framework.watermarks,
1210                &PipelineTaskKey::new(HISTORY_COHORT[0]),
1211                &Watermark::for_checkpoint(7),
1212            )
1213            .unwrap();
1214        batch.commit().unwrap();
1215        assert_eq!(
1216            AuthorityStorePruner::rpc_store_max_eligible_checkpoint(Some(&store)).unwrap(),
1217            8,
1218        );
1219    }
1220
1221    #[cfg(not(tidehunter))]
1222    fn get_keys_after_pruning(path: &Path) -> anyhow::Result<HashSet<ObjectKey>> {
1223        let perpetual_db_path = path.join(Path::new("perpetual"));
1224        let cf_names = AuthorityPerpetualTables::describe_tables();
1225        let cfs: Vec<_> = cf_names
1226            .keys()
1227            .map(|x| (x.as_str(), default_db_options().options))
1228            .collect();
1229        let perpetual_db = typed_store::rocks::open_cf_opts(
1230            perpetual_db_path,
1231            None,
1232            MetricConf::new("perpetual_pruning"),
1233            &cfs,
1234        );
1235
1236        let mut after_pruning = HashSet::new();
1237        let objects = DBMap::<ObjectKey, StoreObjectWrapper>::reopen(
1238            &perpetual_db?,
1239            Some("objects"),
1240            // open the db to bypass default db options which ignores range tombstones
1241            // so we can read the accurate number of retained versions
1242            &ReadWriteOptions::default(),
1243            false,
1244        )?;
1245        let iter = objects.safe_iter();
1246        for item in iter {
1247            after_pruning.insert(item?.0);
1248        }
1249        Ok(after_pruning)
1250    }
1251
1252    #[cfg(not(tidehunter))]
1253    type GenerateTestDataResult = (Vec<ObjectKey>, Vec<ObjectKey>, Vec<ObjectKey>);
1254
1255    #[cfg(not(tidehunter))]
1256    fn generate_test_data(
1257        db: Arc<AuthorityPerpetualTables>,
1258        num_versions_per_object: u64,
1259        num_object_versions_to_retain: u64,
1260        total_unique_object_ids: u32,
1261    ) -> Result<GenerateTestDataResult, anyhow::Error> {
1262        assert!(num_versions_per_object >= num_object_versions_to_retain);
1263
1264        let (mut to_keep, mut to_delete, mut tombstones) = (vec![], vec![], vec![]);
1265        let mut batch = db.objects.batch();
1266
1267        let ids = ObjectID::in_range(ObjectID::ZERO, total_unique_object_ids.into())?;
1268        for id in ids {
1269            for (counter, seq) in (0..num_versions_per_object).rev().enumerate() {
1270                let object_key = ObjectKey(id, SequenceNumber::from_u64(seq));
1271                if counter < num_object_versions_to_retain.try_into().unwrap() {
1272                    // latest `num_object_versions_to_retain` should not have been pruned
1273                    to_keep.push(object_key);
1274                } else {
1275                    to_delete.push(object_key);
1276                }
1277                let obj = get_store_object(Object::immutable_with_id_for_testing(id));
1278                batch.insert_batch(
1279                    &db.objects,
1280                    [(ObjectKey(id, SequenceNumber::from(seq)), obj.clone())],
1281                )?;
1282            }
1283
1284            // Adding a tombstone for deleted object.
1285            if num_object_versions_to_retain == 0 {
1286                let tombstone_key = ObjectKey(id, SequenceNumber::from(num_versions_per_object));
1287                println!("Adding tombstone object {:?}", tombstone_key);
1288                batch.insert_batch(
1289                    &db.objects,
1290                    [(tombstone_key, StoreObjectWrapper::V1(StoreObject::Deleted))],
1291                )?;
1292                tombstones.push(tombstone_key);
1293            }
1294        }
1295        batch.write().unwrap();
1296        assert_eq!(
1297            to_keep.len() as u64,
1298            std::cmp::min(num_object_versions_to_retain, num_versions_per_object)
1299                * total_unique_object_ids as u64
1300        );
1301        assert_eq!(
1302            tombstones.len() as u64,
1303            if num_object_versions_to_retain == 0 {
1304                total_unique_object_ids as u64
1305            } else {
1306                0
1307            }
1308        );
1309        Ok((to_keep, to_delete, tombstones))
1310    }
1311
1312    #[cfg(not(tidehunter))]
1313    async fn run_pruner(
1314        path: &Path,
1315        num_versions_per_object: u64,
1316        num_object_versions_to_retain: u64,
1317        total_unique_object_ids: u32,
1318    ) -> Vec<ObjectKey> {
1319        let registry = Registry::default();
1320        let metrics = AuthorityStorePruningMetrics::new(&registry);
1321        let to_keep = {
1322            let db = Arc::new(AuthorityPerpetualTables::open(path, None, None));
1323            let (to_keep, to_delete, tombstones) = generate_test_data(
1324                db.clone(),
1325                num_versions_per_object,
1326                num_object_versions_to_retain,
1327                total_unique_object_ids,
1328            )
1329            .unwrap();
1330            let mut effects = TransactionEffects::default();
1331            for object in to_delete {
1332                effects.unsafe_add_deleted_live_object_for_testing((
1333                    object.0,
1334                    object.1,
1335                    ObjectDigest::MIN,
1336                ));
1337            }
1338            for object in tombstones {
1339                effects.unsafe_add_object_tombstone_for_testing((
1340                    object.0,
1341                    object.1,
1342                    ObjectDigest::MIN,
1343                ));
1344            }
1345            AuthorityStorePruner::prune_objects_and_indexes(
1346                vec![(0, effects)],
1347                &db,
1348                0,
1349                metrics,
1350                0,
1351                None,
1352                true,
1353            )
1354            .await
1355            .unwrap();
1356            to_keep
1357        };
1358        tokio::time::sleep(Duration::from_secs(3)).await;
1359        to_keep
1360    }
1361
1362    // Tests pruning old version of live objects.
1363    #[cfg(not(tidehunter))]
1364    #[tokio::test]
1365    async fn test_pruning_objects() {
1366        let path = tempfile::tempdir().unwrap().keep();
1367        let to_keep = run_pruner(&path, 3, 2, 1000).await;
1368        assert_eq!(
1369            HashSet::from_iter(to_keep),
1370            get_keys_after_pruning(&path).unwrap()
1371        );
1372        run_pruner(&tempfile::tempdir().unwrap().keep(), 3, 2, 1000).await;
1373    }
1374
1375    // Tests pruning deleted objects (object tombstones).
1376    #[cfg(not(tidehunter))]
1377    #[tokio::test]
1378    async fn test_pruning_tombstones() {
1379        let path = tempfile::tempdir().unwrap().keep();
1380        let to_keep = run_pruner(&path, 0, 0, 1000).await;
1381        assert_eq!(to_keep.len(), 0);
1382        assert_eq!(get_keys_after_pruning(&path).unwrap().len(), 0);
1383
1384        let path = tempfile::tempdir().unwrap().keep();
1385        let to_keep = run_pruner(&path, 3, 0, 1000).await;
1386        assert_eq!(to_keep.len(), 0);
1387        assert_eq!(get_keys_after_pruning(&path).unwrap().len(), 0);
1388    }
1389
1390    #[cfg(not(target_env = "msvc"))]
1391    #[tokio::test]
1392    async fn test_db_size_after_compaction() -> Result<(), anyhow::Error> {
1393        let primary_path = tempfile::tempdir()?.keep();
1394        let perpetual_db = Arc::new(AuthorityPerpetualTables::open(&primary_path, None, None));
1395        let total_unique_object_ids = 10_000;
1396        let num_versions_per_object = 10;
1397        let ids = ObjectID::in_range(ObjectID::ZERO, total_unique_object_ids)?;
1398        let mut to_delete = vec![];
1399        for id in ids {
1400            for i in (0..num_versions_per_object).rev() {
1401                if i < num_versions_per_object - 2 {
1402                    to_delete.push((id, SequenceNumber::from(i)));
1403                }
1404                let obj = get_store_object(Object::immutable_with_id_for_testing(id));
1405                perpetual_db
1406                    .objects
1407                    .insert(&ObjectKey(id, SequenceNumber::from(i)), &obj)?;
1408            }
1409        }
1410
1411        fn get_sst_size(path: &Path) -> u64 {
1412            let mut size = 0;
1413            for entry in std::fs::read_dir(path).unwrap() {
1414                let entry = entry.unwrap();
1415                let path = entry.path();
1416                if let Some(ext) = path.extension() {
1417                    if ext != "sst" {
1418                        continue;
1419                    }
1420                    size += std::fs::metadata(path).unwrap().len();
1421                }
1422            }
1423            size
1424        }
1425
1426        let db_path = primary_path.clone().join("perpetual");
1427        let start = ObjectKey(ObjectID::ZERO, SequenceNumber::MIN);
1428        let end = ObjectKey(ObjectID::MAX, SequenceNumber::MAX);
1429
1430        perpetual_db.objects.compact_range(&start, &end)?;
1431        let before_compaction_size = get_sst_size(&db_path);
1432
1433        let mut effects = TransactionEffects::default();
1434        for object in to_delete {
1435            effects.unsafe_add_deleted_live_object_for_testing((
1436                object.0,
1437                object.1,
1438                ObjectDigest::MIN,
1439            ));
1440        }
1441        let registry = Registry::default();
1442        let metrics = AuthorityStorePruningMetrics::new(&registry);
1443        let total_pruned = AuthorityStorePruner::prune_objects_and_indexes(
1444            vec![(0, effects)],
1445            &perpetual_db,
1446            0,
1447            metrics,
1448            0,
1449            None,
1450            true,
1451        )
1452        .await;
1453        info!("Total pruned keys = {:?}", total_pruned);
1454
1455        perpetual_db.objects.compact_range(&start, &end)?;
1456        let after_compaction_size = get_sst_size(&db_path);
1457
1458        info!(
1459            "Before compaction disk size = {:?}, after compaction disk size = {:?}",
1460            before_compaction_size, after_compaction_size
1461        );
1462        ma::assert_le!(after_compaction_size, before_compaction_size);
1463        Ok(())
1464    }
1465}