Skip to main content

reifydb_cdc/compact/
actor.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}