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