1use 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}