Skip to main content

reifydb_cdc/consume/
actor.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}