Skip to main content

reifydb_core/interface/
cdc.rs

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