reifydb_cdc/compact/
actor.rs1use std::sync::Arc;
5
6use reifydb_core::{
7 common::CommitVersion,
8 interface::catalog::config::{ConfigKey, GetConfig},
9};
10use reifydb_runtime::actor::{
11 context::Context,
12 system::ActorConfig,
13 traits::{Actor, Directive},
14};
15use reifydb_value::value::duration::Duration;
16use tracing::{debug, error, info, trace};
17
18use crate::{produce::watermark::CdcProducerWatermark, storage::sqlite::storage::SqliteCdcStorage};
19
20pub enum CompactMessage {
21 Tick,
22
23 CompactAll,
24}
25
26pub struct CompactActor {
27 config: Arc<dyn GetConfig>,
28 store: SqliteCdcStorage,
29 watermark: CdcProducerWatermark,
30}
31
32impl CompactActor {
33 pub fn new(config: Arc<dyn GetConfig>, store: SqliteCdcStorage, watermark: CdcProducerWatermark) -> Self {
34 Self {
35 config,
36 store,
37 watermark,
38 }
39 }
40
41 fn read_block_size(&self) -> usize {
42 self.config.get_config_uint8(ConfigKey::CdcCompactBlockSize) as usize
43 }
44
45 fn read_safety_lag(&self) -> u64 {
46 self.config.get_config_uint8(ConfigKey::CdcCompactSafetyLag)
47 }
48
49 fn read_max_blocks_per_tick(&self) -> usize {
50 self.config.get_config_uint8(ConfigKey::CdcCompactMaxBlocksPerTick) as usize
51 }
52
53 fn read_interval(&self) -> Duration {
54 self.config.get_config_duration(ConfigKey::CdcCompactInterval)
55 }
56
57 fn read_zstd_level(&self) -> u8 {
58 self.config.get_config_uint1(ConfigKey::CdcCompactZstdLevel)
59 }
60}
61
62impl Actor for CompactActor {
63 type State = ();
64 type Message = CompactMessage;
65
66 fn init(&self, ctx: &Context<Self::Message>) -> Self::State {
67 let interval = self.read_interval();
68 info!("[CdcCompact] started: interval={:?}", interval);
69 ctx.schedule_once(interval, || CompactMessage::Tick);
70 }
71
72 fn handle(&self, _state: &mut Self::State, msg: Self::Message, ctx: &Context<Self::Message>) -> Directive {
73 if ctx.is_cancelled() {
74 info!("[CdcCompact] stopped");
75 return Directive::Stop;
76 }
77 match msg {
78 CompactMessage::Tick => self.on_tick(ctx),
79 CompactMessage::CompactAll => self.on_compact_all(),
80 }
81 Directive::Continue
82 }
83
84 fn config(&self) -> ActorConfig {
85 ActorConfig::new().mailbox_capacity(8)
86 }
87}
88
89impl CompactActor {
90 #[inline]
91 fn on_tick(&self, ctx: &Context<CompactMessage>) {
92 let block_size = self.read_block_size();
93 let safety_lag = self.read_safety_lag();
94 let max_blocks = self.read_max_blocks_per_tick();
95 let zstd_level = self.read_zstd_level();
96 let watermark = self.watermark.get();
97
98 let produced = self.run_tick_loop(block_size, safety_lag, zstd_level, watermark, max_blocks);
99 if produced > 0 {
100 debug!("[CdcCompact] produced {produced} block(s) this tick");
101 }
102
103 ctx.schedule_once(self.read_interval(), || CompactMessage::Tick);
104 }
105
106 #[inline]
107 fn run_tick_loop(
108 &self,
109 block_size: usize,
110 safety_lag: u64,
111 zstd_level: u8,
112 watermark: CommitVersion,
113 max_blocks: usize,
114 ) -> usize {
115 let mut produced = 0usize;
116 while produced < max_blocks {
117 match self.store.compact_oldest(block_size, safety_lag, zstd_level, watermark) {
118 Ok(Some(s)) => {
119 trace!(
120 "[CdcCompact] block: [{}..{}] entries={} bytes={}",
121 s.min_version.0, s.max_version.0, s.num_entries, s.compressed_bytes,
122 );
123 produced += 1;
124 }
125 Ok(None) => break,
126 Err(e) => {
127 error!("[CdcCompact] {e}");
128 break;
129 }
130 }
131 }
132 produced
133 }
134
135 #[inline]
136 fn on_compact_all(&self) {
137 let block_size = self.read_block_size();
138 let zstd_level = self.read_zstd_level();
139 let watermark = self.watermark.get();
140 match self.store.compact_all(block_size, zstd_level, watermark) {
141 Ok(s) => debug!("[CdcCompact] CompactAll produced {} block(s)", s.len()),
142 Err(e) => error!("[CdcCompact] CompactAll error: {e}"),
143 }
144 }
145}