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;
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 #[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 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 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 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 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 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 Ok(indexed.map_or(0, |c| c.saturating_add(1)))
406 }
407
408 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 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 &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 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 let mut effects_to_prune: Vec<(CheckpointSequenceNumber, TransactionEffects)> = vec![];
568 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 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 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 &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 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 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>, ) -> Self {
1100 #[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 #[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 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 assert_eq!(
1241 AuthorityStorePruner::rpc_store_max_eligible_checkpoint(Some(&store)).unwrap(),
1242 0,
1243 );
1244
1245 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 &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 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 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(®istry);
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 #[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 #[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(®istry);
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}