Skip to main content

reifydb_store_multi/gc/operator/
actor.rs

1// SPDX-License-Identifier: AGPL-3.0-or-later
2// Copyright (c) 2026 ReifyDB
3
4use std::collections::HashMap;
5
6use reifydb_core::{
7	actors::operator_ttl::OperatorTtlMessage as Message,
8	event::row::OperatorRowsExpiredEvent,
9	interface::{
10		catalog::{config::ConfigKey, flow::FlowNodeId},
11		store::EntryKind,
12	},
13	key::flow_node_state::FlowNodeStateKey,
14	row::{TtlAnchor, TtlCleanupMode},
15};
16use reifydb_runtime::actor::{
17	context::Context,
18	mailbox::ActorRef,
19	system::{ActorConfig, ActorSystem},
20	timers::TimerHandle,
21	traits::{Actor as ActorTrait, Directive},
22};
23use reifydb_value::value::datetime::DateTime;
24use tracing::{debug, info, trace, warn};
25
26use super::{ListOperatorSettings, OperatorScanStats, scanner};
27use crate::{gc::row::scanner::ScanResult, store::StandardMultiStore, tier::RangeCursor};
28
29#[derive(Default)]
30pub struct ScannerState {
31	cursors: HashMap<FlowNodeId, RangeCursor>,
32}
33
34pub struct ActorState {
35	_timer_handle: Option<TimerHandle>,
36	scanning: bool,
37	scanner: ScannerState,
38}
39
40pub struct Actor<P: ListOperatorSettings> {
41	store: StandardMultiStore,
42	provider: P,
43}
44
45impl<P: ListOperatorSettings> Actor<P> {
46	pub fn new(store: StandardMultiStore, provider: P) -> Self {
47		Self {
48			store,
49			provider,
50		}
51	}
52
53	pub fn spawn(system: &ActorSystem, store: StandardMultiStore, provider: P) -> ActorRef<Message> {
54		let actor = Self::new(store, provider);
55		system.spawn_background("operator-row", actor).actor_ref().clone()
56	}
57
58	fn run_scan(&self, state: &mut ActorState, now: DateTime) {
59		if state.scanning {
60			debug!("Operator TTL scan already in progress, skipping tick");
61			return;
62		}
63
64		let buffer = self.store.commit();
65		let persistent = self.store.persistent();
66		if buffer.is_none() && persistent.is_none() {
67			warn!("Operator TTL scan skipped: no storage tier is configured");
68			return;
69		}
70
71		state.scanning = true;
72
73		let now_nanos = now.to_nanos();
74		trace!(now_nanos, "Starting operator TTL scan");
75
76		let entries = self.provider.list_operator_settings();
77		let config = self.provider.config();
78		let mut stats = OperatorScanStats::default();
79		let mut persistent_rows_deleted: u64 = 0;
80
81		let batch_size = config.get_config_uint8(ConfigKey::OperatorTtlScanBatchSize) as usize;
82
83		for (node_id, settings) in &entries {
84			if let Some(join) = settings.join.as_ref() {
85				let left = join.left.as_ref();
86				let right = join.right.as_ref();
87				if left.is_none() && right.is_none() {
88					continue;
89				}
90
91				if let Some(buffer) = buffer {
92					let mut cursor = state.scanner.cursors.remove(node_id).unwrap_or_default();
93					match scanner::scan_operator_join(
94						buffer,
95						*node_id,
96						left,
97						right,
98						now_nanos,
99						batch_size,
100						&mut cursor,
101					) {
102						Ok((expired, result)) => {
103							stats.operators_scanned += 1;
104							if !expired.is_empty() {
105								stats.rows_expired += expired.len() as u64;
106								for row in &expired {
107									*stats.bytes_discovered
108										.entry(row.node_id)
109										.or_insert(0) += row.scanned_bytes;
110								}
111								if let Err(e) = scanner::drop_expired_operator_keys(
112									buffer, &expired, &mut stats,
113								) {
114									warn!(?node_id, error = %e, "Failed to drop expired join-state keys");
115								}
116							}
117							if let ScanResult::Yielded = result {
118								state.scanner.cursors.insert(*node_id, cursor);
119							}
120						}
121						Err(e) => {
122							warn!(?node_id, error = %e, "Failed to scan join operator state for expired rows");
123						}
124					}
125				}
126
127				if let Some(persistent) = persistent {
128					for (side_ttl, side_prefix) in
129						[(left, scanner::JOIN_LEFT_PREFIX), (right, scanner::JOIN_RIGHT_PREFIX)]
130					{
131						let Some(ttl) = side_ttl else {
132							continue;
133						};
134						let cutoff = now_nanos.saturating_sub(ttl.duration_nanos);
135						let prefix = FlowNodeStateKey::encoded(*node_id, vec![side_prefix]);
136						match persistent.delete_expired(
137							EntryKind::Operator(*node_id),
138							ttl.anchor,
139							cutoff,
140							Some(prefix.as_ref()),
141						) {
142							Ok(deleted) => persistent_rows_deleted += deleted,
143							Err(e) => {
144								warn!(?node_id, error = %e, "Failed to evict expired persistent join rows");
145							}
146						}
147					}
148				}
149
150				continue;
151			}
152
153			let Some(ttl) = settings.ttl.as_ref() else {
154				continue;
155			};
156			trace!(?node_id, ?ttl, "Evaluating TTL config for operator");
157			if ttl.cleanup_mode == TtlCleanupMode::Delete {
158				debug!(?node_id, "Skipping operator with TtlCleanupMode::Delete (not supported in V1)");
159				stats.operators_skipped += 1;
160				continue;
161			}
162
163			if let Some(buffer) = buffer {
164				let mut cursor = state.scanner.cursors.remove(node_id).unwrap_or_default();
165
166				let scan_result = match ttl.anchor {
167					TtlAnchor::Created => scanner::scan_operator_by_created_at(
168						buffer,
169						*node_id,
170						ttl,
171						now_nanos,
172						batch_size,
173						&mut cursor,
174					),
175					TtlAnchor::Updated => scanner::scan_operator_by_updated_at(
176						buffer,
177						*node_id,
178						ttl,
179						now_nanos,
180						batch_size,
181						&mut cursor,
182					),
183				};
184
185				match scan_result {
186					Ok((expired, result)) => {
187						stats.operators_scanned += 1;
188
189						if !expired.is_empty() {
190							stats.rows_expired += expired.len() as u64;
191							for row in &expired {
192								*stats.bytes_discovered
193									.entry(row.node_id)
194									.or_insert(0) += row.scanned_bytes;
195							}
196
197							if let Err(e) = scanner::drop_expired_operator_keys(
198								buffer, &expired, &mut stats,
199							) {
200								warn!(?node_id, error = %e, "Failed to drop expired operator-state keys");
201							}
202						}
203
204						match result {
205							ScanResult::Yielded => {
206								state.scanner.cursors.insert(*node_id, cursor);
207							}
208							ScanResult::Exhausted => {}
209						}
210					}
211					Err(e) => {
212						warn!(?node_id, error = %e, "Failed to scan operator state for expired rows");
213					}
214				}
215			}
216
217			if let Some(persistent) = persistent {
218				let cutoff = now_nanos.saturating_sub(ttl.duration_nanos);
219				match persistent.delete_expired(EntryKind::Operator(*node_id), ttl.anchor, cutoff, None)
220				{
221					Ok(deleted) => {
222						persistent_rows_deleted += deleted;
223						if deleted > 0 {
224							debug!(
225								?node_id,
226								deleted,
227								"Evicted expired operator rows from persistent tier"
228							);
229						}
230					}
231					Err(e) => {
232						warn!(?node_id, error = %e, "Failed to evict expired persistent operator rows");
233					}
234				}
235			}
236		}
237
238		if let Some(buffer) = buffer
239			&& stats.rows_expired > 0
240		{
241			buffer.maintenance();
242		}
243
244		if buffer.is_none()
245			&& let Some(persistent) = persistent
246			&& let Err(e) = persistent.maybe_checkpoint()
247		{
248			warn!(error = %e, "persistent WAL checkpoint failed");
249		}
250
251		if stats.rows_expired > 0 || persistent_rows_deleted > 0 {
252			info!(
253				operators_scanned = stats.operators_scanned,
254				operators_skipped = stats.operators_skipped,
255				rows_expired = stats.rows_expired,
256				versions_dropped = stats.versions_dropped,
257				persistent_rows_deleted,
258				"Operator TTL scan completed"
259			);
260		} else {
261			debug!(
262				operators_scanned = stats.operators_scanned,
263				operators_skipped = stats.operators_skipped,
264				"Operator TTL scan completed (no expired rows)"
265			);
266		}
267
268		self.store.event_bus.emit(OperatorRowsExpiredEvent::new(
269			stats.operators_scanned,
270			stats.operators_skipped,
271			stats.rows_expired,
272			stats.versions_dropped,
273			stats.bytes_discovered,
274			stats.bytes_reclaimed,
275		));
276
277		state.scanning = false;
278	}
279}
280
281impl<P: ListOperatorSettings> ActorTrait for Actor<P> {
282	type State = ActorState;
283	type Message = Message;
284
285	fn init(&self, ctx: &Context<Message>) -> ActorState {
286		debug!("Operator TTL actor started");
287		let config = self.provider.config();
288		let scan_interval = config.get_config_duration(ConfigKey::OperatorTtlScanInterval);
289
290		let timer_handle = ctx.schedule_tick(scan_interval, |nanos| Message::Tick(DateTime::from_nanos(nanos)));
291		ActorState {
292			_timer_handle: Some(timer_handle),
293			scanning: false,
294			scanner: ScannerState::default(),
295		}
296	}
297
298	fn handle(&self, state: &mut ActorState, msg: Message, ctx: &Context<Message>) -> Directive {
299		if ctx.is_cancelled() {
300			return Directive::Stop;
301		}
302
303		match msg {
304			Message::Tick(now) => {
305				self.run_scan(state, now);
306			}
307			Message::Shutdown => {
308				debug!("Operator TTL actor shutting down");
309				return Directive::Stop;
310			}
311		}
312
313		Directive::Continue
314	}
315
316	fn post_stop(&self) {
317		debug!("Operator TTL actor stopped");
318	}
319
320	fn config(&self) -> ActorConfig {
321		ActorConfig::new().mailbox_capacity(64)
322	}
323}
324
325pub fn spawn_operator_settings_actor<P: ListOperatorSettings>(
326	store: StandardMultiStore,
327	system: ActorSystem,
328	provider: P,
329) -> ActorRef<Message> {
330	Actor::spawn(&system, store, provider)
331}