Skip to main content

reifydb_core/key/
cdc_consumer.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}