1use 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 #[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 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 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 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 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 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 Ok(indexed.map_or(0, |c| c.saturating_add(1)))
402 }
403
404 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 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 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 let mut effects_to_prune: Vec<(CheckpointSequenceNumber, TransactionEffects)> = vec![];
558 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 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 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 &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 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 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>, ) -> Self {
1048 #[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 #[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 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 assert_eq!(
1189 AuthorityStorePruner::rpc_store_max_eligible_checkpoint(Some(&store)).unwrap(),
1190 0,
1191 );
1192
1193 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 &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 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 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(®istry);
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 #[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 #[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(®istry);
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}