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