reifydb_store_multi/gc/operator/
actor.rs1use 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}