Skip to main content

reifydb_core/interface/
cdc.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_value::value::datetime::DateTime;
5use serde::{Deserialize, Serialize};
6
7use crate::{
8	common::CommitVersion,
9	encoded::{key::EncodedKey, row::EncodedRow},
10	interface::change::Change,
11};
12
13#[repr(transparent)]
14#[derive(Debug, Clone, PartialOrd, PartialEq, Ord, Eq, Hash)]
15pub struct CdcConsumerId(pub(crate) String);
16
17impl CdcConsumerId {
18	pub fn new(id: impl Into<String>) -> Self {
19		let id = id.into();
20		assert_ne!(id, "__FLOW_COORDINATOR");
21		assert_ne!(id, "__SUBSCRIPTION_CONSUMER");
22		Self(id)
23	}
24
25	pub fn flow_consumer() -> Self {
26		Self("__FLOW_COORDINATOR".to_string())
27	}
28
29	pub fn subscription_consumer() -> Self {
30		Self("__SUBSCRIPTION_CONSUMER".to_string())
31	}
32}
33
34impl AsRef<str> for CdcConsumerId {
35	fn as_ref(&self) -> &str {
36		&self.0
37	}
38}
39
40#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
41pub enum SystemChange {
42	Insert {
43		key: EncodedKey,
44		post: EncodedRow,
45	},
46	Update {
47		key: EncodedKey,
48		pre: EncodedRow,
49		post: EncodedRow,
50	},
51	Delete {
52		key: EncodedKey,
53		pre: Option<EncodedRow>,
54	},
55}
56
57impl SystemChange {
58	pub fn key(&self) -> &EncodedKey {
59		match self {
60			SystemChange::Insert {
61				key,
62				..
63			} => key,
64			SystemChange::Update {
65				key,
66				..
67			} => key,
68			SystemChange::Delete {
69				key,
70				..
71			} => key,
72		}
73	}
74
75	pub fn value_bytes(&self) -> usize {
76		match self {
77			SystemChange::Insert {
78				post,
79				..
80			} => post.len(),
81			SystemChange::Update {
82				pre,
83				post,
84				..
85			} => pre.len() + post.len(),
86			SystemChange::Delete {
87				pre,
88				..
89			} => pre.as_ref().map(|p| p.len()).unwrap_or(0),
90		}
91	}
92}
93
94#[derive(Debug, Clone, Serialize, Deserialize)]
95pub struct Cdc {
96	pub version: CommitVersion,
97	pub timestamp: DateTime,
98
99	pub changes: Vec<Change>,
100
101	pub system_changes: Vec<SystemChange>,
102}
103
104impl Cdc {
105	pub fn new(
106		version: CommitVersion,
107		timestamp: DateTime,
108		changes: Vec<Change>,
109		system_changes: Vec<SystemChange>,
110	) -> Self {
111		Self {
112			version,
113			timestamp,
114			changes,
115			system_changes,
116		}
117	}
118}
119
120#[derive(Debug, Clone, PartialEq, Eq)]
121pub struct ConsumerState {
122	pub consumer_id: CdcConsumerId,
123	pub checkpoint: CommitVersion,
124}
125
126#[derive(Debug, Clone)]
127pub struct CdcBatch {
128	pub items: Vec<Cdc>,
129
130	pub has_more: bool,
131}
132
133impl CdcBatch {
134	pub fn empty() -> Self {
135		Self {
136			items: Vec::new(),
137			has_more: false,
138		}
139	}
140
141	pub fn is_empty(&self) -> bool {
142		self.items.is_empty()
143	}
144}