Skip to main content

reifydb_store_multi/gc/operator/
actor.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{collections::HashMap, mem::take};
5
6use reifydb_core::{
7	actors::operator_ttl::OperatorTtlMessage as Message,
8	common::CommitVersion,
9	event::row::OperatorRowsExpiredEvent,
10	interface::{
11		catalog::{config::ConfigKey, flow::FlowNodeId},
12		store::EntryKind,
13	},
14	key::flow_node_state::FlowNodeStateKey,
15	row::{Ttl, TtlCleanupMode},
16};
17use reifydb_runtime::{
18	actor::{
19		context::Context,
20		mailbox::ActorRef,
21		system::{ActorConfig, ActorSpawner},
22		timers::TimerHandle,
23		traits::{Actor as ActorTrait, Directive},
24	},
25	version_epoch::VersionEpoch,
26};
27use reifydb_value::{reifydb_assertions, value::datetime::DateTime};
28use tracing::{debug, trace, warn};
29
30use super::{ListOperatorSettings, OperatorScanStats, scanner};
31use crate::{
32	gc::row::scanner::ScanResult,
33	store::StandardMultiStore,
34	tier::{RangeCursor, commit::buffer::MultiCommitBufferTier, persistent::MultiPersistentTier},
35};
36
37#[derive(Default)]
38pub struct ScannerState {
39	cursors: HashMap<FlowNodeId, RangeCursor>,
40}
41
42pub struct ActorState {
43	_timer_handle: Option<TimerHandle>,
44	scanning: bool,
45	scanner: ScannerState,
46}
47
48pub struct Actor<P: ListOperatorSettings> {
49	store: StandardMultiStore,
50	provider: P,
51	epoch: VersionEpoch,
52}
53
54impl<P: ListOperatorSettings> Actor<P> {
55	pub fn new(store: StandardMultiStore, provider: P, epoch: VersionEpoch) -> Self {
56		Self {
57			store,
58			provider,
59			epoch,
60		}
61	}
62
63	pub fn spawn(
64		spawner: &ActorSpawner,
65		store: StandardMultiStore,
66		provider: P,
67		epoch: VersionEpoch,
68	) -> ActorRef<Message> {
69		let actor = Self::new(store, provider, epoch);
70		spawner.spawn_coordination("operator-row", actor).actor_ref().clone()
71	}
72
73	fn run_scan(&self, state: &mut ActorState, now: DateTime) {
74		if state.scanning {
75			debug!("Operator TTL scan already in progress, skipping tick");
76			return;
77		}
78
79		let buffer = self.store.commit();
80		let persistent = self.store.persistent();
81		if buffer.is_none() && persistent.is_none() {
82			warn!("Operator TTL scan skipped: no storage tier is configured");
83			return;
84		}
85
86		state.scanning = true;
87
88		let now_nanos = now.to_nanos();
89		trace!(now_nanos, "Starting operator TTL scan");
90
91		let (mut stats, persistent_rows_deleted) =
92			self.scan_all_operators(&mut state.scanner, buffer, persistent, now_nanos);
93
94		self.run_maintenance(buffer, persistent, &stats);
95		self.report_scan(&stats, persistent_rows_deleted);
96		self.emit_expired_event(&mut stats);
97
98		state.scanning = false;
99	}
100
101	#[inline]
102	fn scan_all_operators(
103		&self,
104		scan_state: &mut ScannerState,
105		buffer: Option<&MultiCommitBufferTier>,
106		persistent: Option<&MultiPersistentTier>,
107		now_nanos: u64,
108	) -> (OperatorScanStats, u64) {
109		let entries = self.provider.list_operator_settings();
110		let config = self.provider.config();
111		let mut stats = OperatorScanStats::default();
112		let mut persistent_rows_deleted: u64 = 0;
113
114		let batch_size = config.get_config_uint8(ConfigKey::OperatorTtlScanBatchSize) as usize;
115
116		for (node_id, settings) in &entries {
117			if let Some(join) = settings.join.as_ref() {
118				let left = join.left.as_ref();
119				let right = join.right.as_ref();
120				if left.is_none() && right.is_none() {
121					continue;
122				}
123
124				self.scan_join_entry(
125					scan_state,
126					buffer,
127					persistent,
128					*node_id,
129					left,
130					right,
131					now_nanos,
132					batch_size,
133					&mut stats,
134					&mut persistent_rows_deleted,
135				);
136				continue;
137			}
138
139			let Some(ttl) = settings.ttl.as_ref() else {
140				continue;
141			};
142			trace!(?node_id, ?ttl, "Evaluating TTL config for operator");
143			if ttl.cleanup_mode == TtlCleanupMode::Delete {
144				debug!(?node_id, "Skipping operator with TtlCleanupMode::Delete (not supported in V1)");
145				stats.operators_skipped += 1;
146				continue;
147			}
148
149			self.scan_ttl_entry(
150				scan_state,
151				buffer,
152				persistent,
153				*node_id,
154				ttl,
155				now_nanos,
156				batch_size,
157				&mut stats,
158				&mut persistent_rows_deleted,
159			);
160		}
161
162		(stats, persistent_rows_deleted)
163	}
164
165	#[allow(clippy::too_many_arguments)]
166	fn scan_join_entry(
167		&self,
168		scan_state: &mut ScannerState,
169		buffer: Option<&MultiCommitBufferTier>,
170		persistent: Option<&MultiPersistentTier>,
171		node_id: FlowNodeId,
172		left: Option<&Ttl>,
173		right: Option<&Ttl>,
174		now_nanos: u64,
175		batch_size: usize,
176		stats: &mut OperatorScanStats,
177		persistent_rows_deleted: &mut u64,
178	) {
179		reifydb_assertions! {
180			let both_none = left.is_none() && right.is_none();
181			assert!(
182				!both_none,
183				"scan_join_entry was called for node {node_id:?} with neither join side configured; \
184				 the caller's left/right guard let an idle join through, which wastes a buffer range \
185				 scan and a persistent delete_below_version on a node that can never expire rows"
186			);
187		}
188
189		let left_cutoff = left
190			.and_then(|ttl| DateTime::from_nanos(now_nanos).checked_sub(ttl.duration))
191			.and_then(|cutoff| self.epoch.floor_version_at(cutoff.to_nanos()))
192			.map(CommitVersion);
193		let right_cutoff = right
194			.and_then(|ttl| DateTime::from_nanos(now_nanos).checked_sub(ttl.duration))
195			.and_then(|cutoff| self.epoch.floor_version_at(cutoff.to_nanos()))
196			.map(CommitVersion);
197
198		if let Some(buffer) = buffer {
199			let mut cursor = scan_state.cursors.remove(&node_id).unwrap_or_default();
200			match scanner::scan_operator_join(
201				buffer,
202				node_id,
203				left_cutoff,
204				right_cutoff,
205				batch_size,
206				&mut cursor,
207			) {
208				Ok((expired, result)) => {
209					stats.operators_scanned += 1;
210					if !expired.is_empty() {
211						stats.rows_expired += expired.len() as u64;
212						for row in &expired {
213							*stats.bytes_discovered.entry(row.node_id).or_insert(0) +=
214								row.scanned_bytes;
215							self.store.remove_dropped_read_key(&row.key);
216						}
217						if let Err(e) =
218							scanner::drop_expired_operator_keys(buffer, &expired, stats)
219						{
220							warn!(?node_id, error = %e, "Failed to drop expired join-state keys");
221						}
222					}
223					if let ScanResult::Yielded = result {
224						scan_state.cursors.insert(node_id, cursor);
225					}
226				}
227				Err(e) => {
228					warn!(?node_id, error = %e, "Failed to scan join operator state for expired rows");
229				}
230			}
231		}
232
233		if let Some(persistent) = persistent {
234			for (side_cutoff, side_prefix) in
235				[(left_cutoff, scanner::JOIN_LEFT_PREFIX), (right_cutoff, scanner::JOIN_RIGHT_PREFIX)]
236			{
237				let Some(cutoff) = side_cutoff else {
238					continue;
239				};
240				let prefix = FlowNodeStateKey::encoded(node_id, vec![side_prefix]);
241				match persistent.delete_below_version(
242					EntryKind::Operator(node_id),
243					cutoff,
244					Some(prefix.as_ref()),
245				) {
246					Ok(keys) => {
247						*persistent_rows_deleted += keys.len() as u64;
248						for key in &keys {
249							self.store.remove_dropped_read_key(key);
250						}
251					}
252					Err(e) => {
253						warn!(?node_id, error = %e, "Failed to evict expired persistent join rows");
254					}
255				}
256			}
257		}
258	}
259
260	#[allow(clippy::too_many_arguments)]
261	fn scan_ttl_entry(
262		&self,
263		scan_state: &mut ScannerState,
264		buffer: Option<&MultiCommitBufferTier>,
265		persistent: Option<&MultiPersistentTier>,
266		node_id: FlowNodeId,
267		ttl: &Ttl,
268		now_nanos: u64,
269		batch_size: usize,
270		stats: &mut OperatorScanStats,
271		persistent_rows_deleted: &mut u64,
272	) {
273		reifydb_assertions! {
274			let is_delete = ttl.cleanup_mode == TtlCleanupMode::Delete;
275			assert!(
276				!is_delete,
277				"scan_ttl_entry was called for node {node_id:?} with TtlCleanupMode::Delete, which \
278				 the caller is supposed to skip and count as operators_skipped; dropping such rows here \
279				 would silently apply unsupported Delete semantics instead of the intended Drop"
280			);
281		}
282
283		let Some(cutoff) = DateTime::from_nanos(now_nanos).checked_sub(ttl.duration) else {
284			return;
285		};
286		let cutoff_version = self.epoch.floor_version_at(cutoff.to_nanos()).map(CommitVersion);
287
288		if let (Some(buffer), Some(cutoff_version)) = (buffer, cutoff_version) {
289			let mut cursor = scan_state.cursors.remove(&node_id).unwrap_or_default();
290
291			let scan_result = scanner::scan_operator_expired(
292				buffer,
293				node_id,
294				cutoff_version,
295				batch_size,
296				&mut cursor,
297			);
298
299			match scan_result {
300				Ok((expired, result)) => {
301					stats.operators_scanned += 1;
302
303					if !expired.is_empty() {
304						stats.rows_expired += expired.len() as u64;
305						for row in &expired {
306							*stats.bytes_discovered.entry(row.node_id).or_insert(0) +=
307								row.scanned_bytes;
308							self.store.remove_dropped_read_key(&row.key);
309						}
310
311						if let Err(e) =
312							scanner::drop_expired_operator_keys(buffer, &expired, stats)
313						{
314							warn!(?node_id, error = %e, "Failed to drop expired operator-state keys");
315						}
316					}
317
318					match result {
319						ScanResult::Yielded => {
320							scan_state.cursors.insert(node_id, cursor);
321						}
322						ScanResult::Exhausted => {}
323					}
324				}
325				Err(e) => {
326					warn!(?node_id, error = %e, "Failed to scan operator state for expired rows");
327				}
328			}
329		}
330
331		if let (Some(persistent), Some(cutoff_version)) = (persistent, cutoff_version) {
332			match persistent.delete_below_version(EntryKind::Operator(node_id), cutoff_version, None) {
333				Ok(keys) => {
334					*persistent_rows_deleted += keys.len() as u64;
335					if !keys.is_empty() {
336						for key in &keys {
337							self.store.remove_dropped_read_key(key);
338						}
339						debug!(
340							?node_id,
341							deleted = keys.len(),
342							"Evicted expired operator rows from persistent tier"
343						);
344					}
345				}
346				Err(e) => {
347					warn!(?node_id, error = %e, "Failed to evict expired persistent operator rows");
348				}
349			}
350		}
351	}
352
353	#[inline]
354	fn run_maintenance(
355		&self,
356		buffer: Option<&MultiCommitBufferTier>,
357		persistent: Option<&MultiPersistentTier>,
358		stats: &OperatorScanStats,
359	) {
360		if let Some(buffer) = buffer
361			&& stats.rows_expired > 0
362		{
363			buffer.maintenance();
364		}
365
366		if buffer.is_none()
367			&& let Some(persistent) = persistent
368			&& let Err(e) = persistent.maybe_checkpoint()
369		{
370			warn!(error = %e, "persistent WAL checkpoint failed");
371		}
372	}
373
374	#[inline]
375	fn report_scan(&self, stats: &OperatorScanStats, persistent_rows_deleted: u64) {
376		if stats.rows_expired > 0 || persistent_rows_deleted > 0 {
377			debug!(
378				operators_scanned = stats.operators_scanned,
379				operators_skipped = stats.operators_skipped,
380				rows_expired = stats.rows_expired,
381				versions_dropped = stats.versions_dropped,
382				persistent_rows_deleted,
383				"Operator TTL scan completed"
384			);
385		} else {
386			debug!(
387				operators_scanned = stats.operators_scanned,
388				operators_skipped = stats.operators_skipped,
389				"Operator TTL scan completed (no expired rows)"
390			);
391		}
392	}
393
394	#[inline]
395	fn emit_expired_event(&self, stats: &mut OperatorScanStats) {
396		self.store.event_bus.emit(OperatorRowsExpiredEvent::new(
397			stats.operators_scanned,
398			stats.operators_skipped,
399			stats.rows_expired,
400			stats.versions_dropped,
401			take(&mut stats.bytes_discovered),
402			take(&mut stats.bytes_reclaimed),
403		));
404	}
405}
406
407impl<P: ListOperatorSettings> ActorTrait for Actor<P> {
408	type State = ActorState;
409	type Message = Message;
410
411	fn init(&self, ctx: &Context<Message>) -> ActorState {
412		debug!("Operator TTL actor started");
413		let config = self.provider.config();
414		let scan_interval = config.get_config_duration(ConfigKey::OperatorTtlScanInterval);
415
416		let timer_handle = ctx.schedule_tick(scan_interval, |nanos| Message::Tick(DateTime::from_nanos(nanos)));
417		ActorState {
418			_timer_handle: Some(timer_handle),
419			scanning: false,
420			scanner: ScannerState::default(),
421		}
422	}
423
424	fn handle(&self, state: &mut ActorState, msg: Message, ctx: &Context<Message>) -> Directive {
425		if ctx.is_cancelled() {
426			return Directive::Stop;
427		}
428
429		match msg {
430			Message::Tick(now) => {
431				self.run_scan(state, now);
432			}
433			Message::Shutdown => {
434				debug!("Operator TTL actor shutting down");
435				return Directive::Stop;
436			}
437		}
438
439		Directive::Continue
440	}
441
442	fn post_stop(&self) {
443		debug!("Operator TTL actor stopped");
444	}
445
446	fn config(&self) -> ActorConfig {
447		ActorConfig::new().mailbox_capacity(64)
448	}
449}
450
451pub fn spawn_operator_settings_actor<P: ListOperatorSettings>(
452	store: StandardMultiStore,
453	spawner: ActorSpawner,
454	provider: P,
455	epoch: VersionEpoch,
456) -> ActorRef<Message> {
457	Actor::spawn(&spawner, store, provider, epoch)
458}
459
460#[cfg(all(test, feature = "sqlite", not(target_arch = "wasm32")))]
461mod tests {
462	use std::sync::Arc;
463
464	use reifydb_codec::encoded::row::{EncodedRow, SHAPE_HEADER_SIZE};
465	use reifydb_core::{
466		common::CommitVersion,
467		delta::Delta,
468		interface::{catalog::config::GetConfig, store::MultiVersionCommit},
469		row::OperatorSettings,
470	};
471	use reifydb_value::{
472		util::cowvec::CowVec,
473		value::{Value, duration::Duration},
474	};
475
476	use super::*;
477	use crate::tier::VersionedGetResult;
478
479	#[derive(Clone)]
480	struct TestProvider {
481		node: FlowNodeId,
482		ttl: Ttl,
483	}
484
485	impl ListOperatorSettings for TestProvider {
486		fn list_operator_settings(&self) -> Vec<(FlowNodeId, OperatorSettings)> {
487			vec![(
488				self.node,
489				OperatorSettings {
490					ttl: Some(self.ttl.clone()),
491					join: None,
492				},
493			)]
494		}
495
496		fn config(&self) -> Arc<dyn GetConfig> {
497			Arc::new(TestConfig)
498		}
499	}
500
501	struct TestConfig;
502
503	impl GetConfig for TestConfig {
504		fn get_config(&self, key: ConfigKey) -> Value {
505			key.default_value()
506		}
507
508		fn get_config_at(&self, key: ConfigKey, _version: CommitVersion) -> Value {
509			key.default_value()
510		}
511	}
512
513	fn row_with_created(payload: &[u8], created_at: u64) -> CowVec<u8> {
514		let mut buf = vec![0u8; SHAPE_HEADER_SIZE + payload.len()];
515		buf[8..16].copy_from_slice(&created_at.to_le_bytes());
516		buf[16..24].copy_from_slice(&created_at.to_le_bytes());
517		buf[SHAPE_HEADER_SIZE..].copy_from_slice(payload);
518		CowVec::new(buf)
519	}
520
521	#[test]
522	fn operator_ttl_gc_invalidates_read_cache_for_dropped_keys() {
523		let (store, _g) = StandardMultiStore::testing_memory_with_persistent_sqlite();
524		let read = store.read.clone().expect("read tier configured");
525
526		let node = FlowNodeId(1);
527		let opkey = FlowNodeStateKey::encoded(node, vec![1u8]);
528
529		MultiVersionCommit::commit(
530			&store,
531			CowVec::new(vec![Delta::Set {
532				key: opkey.clone(),
533				row: EncodedRow(row_with_created(b"state", 1)),
534			}]),
535			CommitVersion(1),
536		)
537		.unwrap();
538
539		assert!(
540			matches!(read.get(&opkey, CommitVersion(1)), VersionedGetResult::Value { .. }),
541			"write-through must have cached the operator state before GC, otherwise this test cannot \
542			 prove the GC clears a stale entry"
543		);
544
545		let ttl = Ttl {
546			duration: Duration::from_nanoseconds(100).unwrap(),
547			cleanup_mode: TtlCleanupMode::Drop,
548		};
549
550		let epoch = VersionEpoch::new();
551		epoch.record(1, 1);
552		let actor = Actor::new(
553			store.clone(),
554			TestProvider {
555				node,
556				ttl,
557			},
558			epoch,
559		);
560		let mut state = ActorState {
561			_timer_handle: None,
562			scanning: false,
563			scanner: ScannerState::default(),
564		};
565		actor.run_scan(&mut state, DateTime::from_nanos(1_000));
566
567		assert!(
568			matches!(read.get(&opkey, CommitVersion(1)), VersionedGetResult::NotFound),
569			"operator TTL GC must invalidate the read cache for reclaimed keys"
570		);
571	}
572}