reifydb_core/interface/
cdc.rs1use 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}