reifydb_core/key/
cdc_consumer.rs1use reifydb_codec::key::{
5 deserializer::KeyDeserializer,
6 encoded::{EncodedKey, EncodedKeyRange},
7 serializer::KeySerializer,
8};
9
10use super::{EncodableKey, KeyKind};
11use crate::interface::{catalog::flow::FlowId, cdc::CdcConsumerId};
12
13pub trait ToConsumerKey {
14 fn to_consumer_key(&self) -> EncodedKey;
15}
16
17impl ToConsumerKey for EncodedKey {
18 fn to_consumer_key(&self) -> EncodedKey {
19 self.clone()
20 }
21}
22
23impl ToConsumerKey for CdcConsumerId {
24 fn to_consumer_key(&self) -> EncodedKey {
25 CdcConsumerKey {
26 consumer: self.clone(),
27 }
28 .encode()
29 }
30}
31
32impl ToConsumerKey for FlowId {
33 fn to_consumer_key(&self) -> EncodedKey {
34 CdcConsumerKey::encoded(CdcConsumerId::new(format!("flow:{}", self.0)))
35 }
36}
37
38#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
39pub struct CdcConsumerKey {
40 pub consumer: CdcConsumerId,
41}
42
43impl CdcConsumerKey {
44 pub fn encoded(consumer: impl Into<CdcConsumerId>) -> EncodedKey {
45 Self {
46 consumer: consumer.into(),
47 }
48 .encode()
49 }
50}
51
52impl EncodableKey for CdcConsumerKey {
53 const KIND: KeyKind = KeyKind::CdcConsumer;
54
55 fn encode(&self) -> EncodedKey {
56 let mut serializer = KeySerializer::new();
57 serializer.extend_u8(Self::KIND as u8).extend_str(&self.consumer);
58 serializer.to_encoded_key()
59 }
60
61 fn decode(key: &EncodedKey) -> Option<Self>
62 where
63 Self: Sized,
64 {
65 let mut de = KeyDeserializer::from_bytes(key.as_slice());
66
67 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
68 if kind != Self::KIND {
69 return None;
70 }
71
72 let consumer_id = de.read_str().ok()?;
73
74 Some(Self {
75 consumer: CdcConsumerId(consumer_id),
76 })
77 }
78}
79
80#[derive(Debug, Clone, PartialEq, Eq)]
81pub struct CdcConsumerKeyRange;
82
83impl CdcConsumerKeyRange {
84 pub fn full_scan() -> EncodedKeyRange {
85 EncodedKeyRange::start_end(Some(Self::start()), Some(Self::end()))
86 }
87
88 fn start() -> EncodedKey {
89 let mut serializer = KeySerializer::with_capacity(1);
90 serializer.extend_u8(CdcConsumerKey::KIND as u8);
91 serializer.to_encoded_key()
92 }
93
94 fn end() -> EncodedKey {
95 let mut serializer = KeySerializer::with_capacity(1);
96 serializer.extend_u8((CdcConsumerKey::KIND as u8).wrapping_sub(1));
97 serializer.to_encoded_key()
98 }
99}
100
101#[cfg(test)]
102pub mod tests {
103 use std::ops::RangeBounds;
104
105 use super::{CdcConsumerKey, CdcConsumerKeyRange, EncodableKey, ToConsumerKey};
106 use crate::interface::{catalog::flow::FlowId, cdc::CdcConsumerId};
107
108 #[test]
109 fn test_encode_decode_cdc_consumer() {
110 let key = CdcConsumerKey {
111 consumer: CdcConsumerId::new("test-consumer"),
112 };
113
114 let encoded = key.encode();
115 let decoded = CdcConsumerKey::decode(&encoded).expect("Failed to decode key");
116
117 assert_eq!(decoded.consumer, CdcConsumerId::new("test-consumer"));
118 }
119
120 #[test]
121 fn test_cdc_consumer_keys_within_range() {
122 let key1 = CdcConsumerKey {
123 consumer: CdcConsumerId::new("consumer-a"),
124 }
125 .encode();
126
127 let key2 = CdcConsumerKey {
128 consumer: CdcConsumerId::new("consumer-b"),
129 }
130 .encode();
131
132 let key3 = CdcConsumerKey {
133 consumer: CdcConsumerId::new("consumer-z"),
134 }
135 .encode();
136
137 let range = CdcConsumerKeyRange::full_scan();
138
139 assert!(range.contains(&key1), "consumer-a key should be in range");
140 assert!(range.contains(&key2), "consumer-b key should be in range");
141 assert!(range.contains(&key3), "consumer-z key should be in range");
142 }
143
144 #[test]
145 fn test_flow_id_to_consumer_key() {
146 let flow_id = FlowId(42);
147 let encoded = flow_id.to_consumer_key();
148
149 let decoded = CdcConsumerKey::decode(&encoded).expect("Failed to decode key");
150 assert_eq!(decoded.consumer, CdcConsumerId::new("flow:42"));
151 }
152
153 #[test]
154 fn test_flow_id_keys_within_range() {
155 let flow1 = FlowId(1).to_consumer_key();
156 let flow2 = FlowId(100).to_consumer_key();
157 let flow3 = FlowId(999).to_consumer_key();
158
159 let range = CdcConsumerKeyRange::full_scan();
160
161 assert!(range.contains(&flow1), "flow:1 key should be in range");
162 assert!(range.contains(&flow2), "flow:100 key should be in range");
163 assert!(range.contains(&flow3), "flow:999 key should be in range");
164 }
165}