1use std::{
5 ops::Bound,
6 process,
7 sync::{
8 Arc,
9 atomic::{AtomicBool, Ordering},
10 },
11};
12
13use reifydb_codec::key::encoded::EncodedKey;
14use reifydb_core::{
15 actors::cdc::CdcPollMessage,
16 common::CommitVersion,
17 interface::{
18 catalog::config::{ConfigKey, GetConfig},
19 cdc::{Cdc, CdcConsumerId, SystemChange},
20 },
21 key::{EncodableKey, Key, cdc_consumer::CdcConsumerKey, kind::KeyKind},
22};
23use reifydb_runtime::actor::{
24 context::Context,
25 system::ActorConfig,
26 traits::{Actor, Directive},
27};
28use reifydb_transaction::transaction::Transaction;
29use reifydb_value::{Result, error::Error, reifydb_assertions, value::duration::Duration};
30use tracing::{debug, error};
31
32use super::{checkpoint::CdcCheckpoint, consumer::CdcConsume, host::CdcHost, watermark::CdcConsumerWatermark};
33use crate::storage::CdcStore;
34
35#[derive(Debug, Clone)]
36pub struct PollActorConfig {
37 pub consumer_id: CdcConsumerId,
38
39 pub poll_interval: Duration,
40
41 pub max_batch_size: Option<u64>,
42}
43
44pub struct PollActor<H: CdcHost, C: CdcConsume> {
45 config: PollActorConfig,
46 host: H,
47 consumer: Box<C>,
48 store: CdcStore,
49 consumer_key: EncodedKey,
50 consumer_watermark: Option<CdcConsumerWatermark>,
51 wake_armed: Arc<AtomicBool>,
52}
53
54impl<H: CdcHost, C: CdcConsume> PollActor<H, C> {
55 pub fn new(
56 config: PollActorConfig,
57 host: H,
58 consumer: C,
59 store: CdcStore,
60 consumer_watermark: Option<CdcConsumerWatermark>,
61 wake_armed: Arc<AtomicBool>,
62 ) -> Self {
63 let consumer_key = CdcConsumerKey {
64 consumer: config.consumer_id.clone(),
65 }
66 .encode();
67
68 Self {
69 config,
70 host,
71 consumer: Box::new(consumer),
72 store,
73 consumer_key,
74 consumer_watermark,
75 wake_armed,
76 }
77 }
78
79 #[inline]
80 fn publish_watermark(&self, version: CommitVersion) {
81 if let Some(wm) = &self.consumer_watermark {
82 wm.store(version);
83 }
84 }
85
86 #[inline]
87 fn watermark_wait_timeout(&self) -> Duration {
88 self.host.catalog().get_config_duration(ConfigKey::CdcWatermarkWaitTimeout)
89 }
90
91 #[inline]
92 fn consume_wait_timeout(&self) -> Duration {
93 self.host.catalog().get_config_duration(ConfigKey::CdcConsumeWaitTimeout)
94 }
95}
96
97pub enum Phase {
98 Ready,
99
100 WaitingForWatermark,
101
102 WaitingForConsume {
103 latest_version: CommitVersion,
104
105 count: usize,
106
107 generation: u64,
108 },
109}
110
111pub struct PollState {
112 phase: Phase,
113
114 cached_checkpoint: Option<CommitVersion>,
115
116 consume_generation: u64,
117}
118
119impl<H: CdcHost, C: CdcConsume + Send + Sync + 'static> Actor for PollActor<H, C> {
120 type State = PollState;
121 type Message = CdcPollMessage;
122
123 fn init(&self, ctx: &Context<Self::Message>) -> Self::State {
124 debug!(
125 "[Consumer {:?}] Started polling with interval {:?}",
126 self.config.consumer_id, self.config.poll_interval
127 );
128
129 let _ = ctx.self_ref().send(CdcPollMessage::Poll);
130
131 PollState {
132 phase: Phase::Ready,
133 cached_checkpoint: None,
134 consume_generation: 0,
135 }
136 }
137
138 fn handle(&self, state: &mut Self::State, msg: Self::Message, ctx: &Context<Self::Message>) -> Directive {
139 match msg {
140 CdcPollMessage::Poll => self.on_poll(state, ctx),
141 CdcPollMessage::CheckWatermark => self.on_check_watermark(state, ctx),
142 CdcPollMessage::ConsumeResponse {
143 generation,
144 result,
145 } => self.on_consume_response(state, ctx, generation, result),
146 CdcPollMessage::CheckConsume {
147 generation,
148 } => self.on_check_consume(state, ctx, generation),
149 CdcPollMessage::Shutdown => {
150 debug!("[Consumer {:?}] Shutdown", self.config.consumer_id);
151 Directive::Stop
152 }
153 }
154 }
155
156 fn config(&self) -> ActorConfig {
157 ActorConfig::new()
158 }
159}
160
161impl<H: CdcHost, C: CdcConsume> PollActor<H, C> {
162 #[inline]
163 fn on_poll(&self, state: &mut PollState, ctx: &Context<CdcPollMessage>) -> Directive {
164 if !matches!(state.phase, Phase::Ready) {
165 return Directive::Continue;
166 }
167 if ctx.is_cancelled() {
168 debug!("[Consumer {:?}] Stopped", self.config.consumer_id);
169 return Directive::Stop;
170 }
171 let current_version = match self.host.current_version() {
172 Ok(v) => v,
173 Err(e) => {
174 error!("[Consumer {:?}] Error getting current version: {}", self.config.consumer_id, e);
175 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
176 return Directive::Continue;
177 }
178 };
179 if self.host.done_until() >= current_version {
180 self.start_consume(state, ctx);
181 } else {
182 state.phase = Phase::WaitingForWatermark;
183 let self_ref = ctx.self_ref();
184 self.host.notify_on_mark(
185 current_version,
186 Box::new(move || {
187 let _ = self_ref.send(CdcPollMessage::CheckWatermark);
188 }),
189 );
190 ctx.schedule_once(self.watermark_wait_timeout(), || CdcPollMessage::CheckWatermark);
191 }
192 Directive::Continue
193 }
194
195 #[inline]
196 fn on_check_watermark(&self, state: &mut PollState, ctx: &Context<CdcPollMessage>) -> Directive {
197 if !matches!(state.phase, Phase::WaitingForWatermark) {
198 return Directive::Continue;
199 }
200 if ctx.is_cancelled() {
201 debug!("[Consumer {:?}] Stopped", self.config.consumer_id);
202 return Directive::Stop;
203 }
204 state.phase = Phase::Ready;
205 self.start_consume(state, ctx);
206 Directive::Continue
207 }
208
209 #[inline]
210 fn on_consume_response(
211 &self,
212 state: &mut PollState,
213 ctx: &Context<CdcPollMessage>,
214 generation: u64,
215 result: Result<()>,
216 ) -> Directive {
217 if let Phase::WaitingForConsume {
218 latest_version,
219 count,
220 generation: pending,
221 } = state.phase
222 {
223 if pending != generation {
224 return Directive::Continue;
225 }
226 state.phase = Phase::Ready;
227 self.finish_consume(state, ctx, latest_version, count, result);
228 }
229 Directive::Continue
230 }
231
232 #[inline]
233 fn on_check_consume(&self, state: &mut PollState, ctx: &Context<CdcPollMessage>, generation: u64) -> Directive {
234 let still_waiting = matches!(
235 state.phase,
236 Phase::WaitingForConsume {
237 generation: pending,
238 ..
239 } if pending == generation
240 );
241 if !still_waiting {
242 return Directive::Continue;
243 }
244 if ctx.is_cancelled() {
245 debug!("[Consumer {:?}] Stopped", self.config.consumer_id);
246 return Directive::Stop;
247 }
248 error!(
249 "[Consumer {:?}] consume reply not received within {:?}; re-dispatching batch",
250 self.config.consumer_id,
251 self.consume_wait_timeout()
252 );
253 state.phase = Phase::Ready;
254 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
255 Directive::Continue
256 }
257
258 fn start_consume(&self, state: &mut PollState, ctx: &Context<CdcPollMessage>) {
259 state.phase = Phase::Ready;
260 self.wake_armed.store(false, Ordering::Release);
261 let safe_version = self.host.cdc_producer_watermark();
262 if safe_version > self.host.done_until() {
263 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
264 return;
265 }
266
267 let Some(checkpoint) = self.resolve_checkpoint(state, ctx) else {
268 return;
269 };
270 if safe_version <= checkpoint {
271 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
272 return;
273 }
274
275 let Some(transactions) = self.fetch_or_reschedule(checkpoint, safe_version, ctx) else {
276 return;
277 };
278 if transactions.is_empty() {
279 self.advance_checkpoint_skip_ahead(state, ctx, safe_version);
280 return;
281 }
282
283 let (count, latest_version) = summarize_batch(checkpoint, &transactions);
284 let relevant_cdcs: Vec<Cdc> = transactions.into_iter().filter(is_relevant_cdc).collect();
285
286 if relevant_cdcs.is_empty() {
287 self.advance_checkpoint_skip_ahead(state, ctx, latest_version);
288 return;
289 }
290
291 state.consume_generation = state.consume_generation.wrapping_add(1);
292 let generation = state.consume_generation;
293 state.phase = Phase::WaitingForConsume {
294 latest_version,
295 count,
296 generation,
297 };
298 self.dispatch_to_consumer(relevant_cdcs, generation, ctx);
299 ctx.schedule_once(self.consume_wait_timeout(), move || CdcPollMessage::CheckConsume {
300 generation,
301 });
302 }
303
304 #[inline]
305 fn advance_checkpoint_skip_ahead(
306 &self,
307 state: &mut PollState,
308 ctx: &Context<CdcPollMessage>,
309 latest_version: CommitVersion,
310 ) {
311 reifydb_assertions! {
312 if let Some(prev) = state.cached_checkpoint {
313 assert!(
314 latest_version >= prev,
315 "the consumer checkpoint moved backwards, so CDC that was already consumed would be \
316 re-delivered (cached checkpoint prev={}, new latest={})",
317 prev.0,
318 latest_version.0
319 );
320 }
321 }
322 state.cached_checkpoint = Some(latest_version);
323 self.publish_watermark(latest_version);
324 let _ = ctx.self_ref().send(CdcPollMessage::Poll);
325 }
326
327 #[inline]
328 fn resolve_checkpoint(&self, state: &mut PollState, ctx: &Context<CdcPollMessage>) -> Option<CommitVersion> {
329 if let Some(v) = state.cached_checkpoint {
330 return Some(v);
331 }
332 let v = self.seed_checkpoint_from_durable(ctx)?;
333 state.cached_checkpoint = Some(v);
334 self.publish_watermark(v);
335 Some(v)
336 }
337
338 #[inline]
339 fn seed_checkpoint_from_durable(&self, ctx: &Context<CdcPollMessage>) -> Option<CommitVersion> {
340 let mut query = match self.host.begin_query() {
341 Ok(q) => q,
342 Err(e) => {
343 error!("[Consumer {:?}] Error beginning query: {}", self.config.consumer_id, e);
344 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
345 return None;
346 }
347 };
348 let v = match CdcCheckpoint::fetch(&mut Transaction::Query(&mut query), &self.consumer_key) {
349 Ok(c) => c,
350 Err(e) => {
351 error!("[Consumer {:?}] Error fetching checkpoint: {}", self.config.consumer_id, e);
352 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
353 return None;
354 }
355 };
356 drop(query);
357 Some(v)
358 }
359
360 #[inline]
361 fn fetch_or_reschedule(
362 &self,
363 checkpoint: CommitVersion,
364 safe_version: CommitVersion,
365 ctx: &Context<CdcPollMessage>,
366 ) -> Option<Vec<Cdc>> {
367 match self.fetch_cdcs_until(checkpoint, safe_version) {
368 Ok(t) => Some(t),
369 Err(e) => {
370 error!("[Consumer {:?}] Error fetching CDCs: {}", self.config.consumer_id, e);
371 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
372 None
373 }
374 }
375 }
376
377 #[inline]
378 fn dispatch_to_consumer(&self, cdcs: Vec<Cdc>, generation: u64, ctx: &Context<CdcPollMessage>) {
379 let self_ref = ctx.self_ref().clone();
380 let reply: Box<dyn FnOnce(Result<()>) + Send> = Box::new(move |result| {
381 let _ = self_ref.send(CdcPollMessage::ConsumeResponse {
382 generation,
383 result,
384 });
385 });
386 self.consumer.consume(cdcs, reply);
387 }
388
389 fn finish_consume(
390 &self,
391 state: &mut PollState,
392 ctx: &Context<CdcPollMessage>,
393 latest_version: CommitVersion,
394 count: usize,
395 result: Result<()>,
396 ) {
397 state.phase = Phase::Ready;
398 match result {
399 Ok(()) => self.advance_after_success(state, ctx, latest_version, count),
400 Err(e) => self.abort_on_error(e),
401 }
402 }
403
404 #[inline]
405 fn advance_after_success(
406 &self,
407 state: &mut PollState,
408 ctx: &Context<CdcPollMessage>,
409 latest_version: CommitVersion,
410 count: usize,
411 ) {
412 reifydb_assertions! {
413 if let Some(prev) = state.cached_checkpoint {
414 assert!(
415 latest_version >= prev,
416 "the consumer checkpoint moved backwards, so CDC that was already consumed would be \
417 re-delivered (cached checkpoint prev={}, new latest={})",
418 prev.0,
419 latest_version.0
420 );
421 }
422 }
423 state.cached_checkpoint = Some(latest_version);
424 self.publish_watermark(latest_version);
425 if count > 0 {
426 let _ = ctx.self_ref().send(CdcPollMessage::Poll);
427 } else {
428 ctx.schedule_once(self.config.poll_interval, || CdcPollMessage::Poll);
429 }
430 }
431
432 #[inline]
433 fn abort_on_error(&self, err: Error) -> ! {
434 error!(
435 "[Consumer {:?}] fatal error consuming events, aborting application: {}",
436 self.config.consumer_id, err
437 );
438 process::abort();
439 }
440
441 fn fetch_cdcs_until(&self, since_version: CommitVersion, until_version: CommitVersion) -> Result<Vec<Cdc>> {
442 let batch_size = self.config.max_batch_size.unwrap_or(1024);
443 let batch = self.store.read_range(
444 Bound::Excluded(since_version),
445 Bound::Included(until_version),
446 batch_size,
447 )?;
448 Ok(batch.items)
449 }
450}
451
452#[inline]
453fn summarize_batch(checkpoint: CommitVersion, transactions: &[Cdc]) -> (usize, CommitVersion) {
454 let count = transactions.len();
455 let latest_version = transactions.iter().map(|tx| tx.version).max().unwrap_or(checkpoint);
456 (count, latest_version)
457}
458
459fn is_relevant_cdc(cdc: &Cdc) -> bool {
460 !cdc.changes.is_empty() || cdc.system_changes.iter().any(is_relevant_system_change)
461}
462
463fn is_relevant_system_change(change: &SystemChange) -> bool {
464 let key = match change {
465 SystemChange::Insert {
466 key,
467 ..
468 }
469 | SystemChange::Update {
470 key,
471 ..
472 }
473 | SystemChange::Delete {
474 key,
475 ..
476 } => key,
477 };
478 Key::kind(key)
479 .map(|kind| {
480 matches!(
481 kind,
482 KeyKind::Row
483 | KeyKind::Flow | KeyKind::FlowNode | KeyKind::FlowNodeByFlow
484 | KeyKind::FlowEdge | KeyKind::FlowEdgeByFlow
485 | KeyKind::NamespaceFlow
486 )
487 })
488 .unwrap_or(false)
489}