reifydb_core/key/
flow_node_state.rs1use std::ops::Bound;
5
6use reifydb_codec::key::{
7 deserializer::KeyDeserializer,
8 encoded::{EncodedKey, EncodedKeyRange},
9 serializer::KeySerializer,
10};
11
12use super::{EncodableKey, EncodableKeyRange, KeyKind};
13use crate::interface::catalog::flow::FlowNodeId;
14
15#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct FlowNodeStateKey {
17 pub node: FlowNodeId,
18 pub key: Vec<u8>,
19}
20
21impl EncodableKey for FlowNodeStateKey {
22 const KIND: KeyKind = KeyKind::FlowNodeState;
23
24 fn encode(&self) -> EncodedKey {
25 let mut serializer = KeySerializer::with_capacity(10 + self.key.len());
26 serializer.extend_u8(Self::KIND as u8).extend_u64(self.node.0).extend_raw(&self.key);
27 serializer.to_encoded_key()
28 }
29
30 fn decode(key: &EncodedKey) -> Option<Self> {
31 let mut de = KeyDeserializer::from_bytes(key.as_slice());
32
33 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
34 if kind != Self::KIND {
35 return None;
36 }
37
38 let node_id = de.read_u64().ok()?;
39 let key_bytes = de.read_raw(de.remaining()).ok()?.to_vec();
40
41 Some(Self {
42 node: FlowNodeId(node_id),
43 key: key_bytes,
44 })
45 }
46}
47
48impl FlowNodeStateKey {
49 pub fn new(node: FlowNodeId, key: Vec<u8>) -> Self {
50 Self {
51 node,
52 key,
53 }
54 }
55
56 pub fn new_empty(node: FlowNodeId) -> Self {
57 Self {
58 node,
59 key: Vec::new(),
60 }
61 }
62
63 pub fn encoded(node: impl Into<FlowNodeId>, key: impl Into<Vec<u8>>) -> EncodedKey {
64 Self::new(node.into(), key.into()).encode()
65 }
66
67 pub fn node_range(node: FlowNodeId) -> EncodedKeyRange {
68 let range = FlowNodeStateKeyRange::new(node);
69 EncodedKeyRange::start_end(range.start(), range.end())
70 }
71}
72
73#[derive(Debug, Clone, PartialEq, Eq)]
74pub struct FlowNodeStateKeyRange {
75 pub node: FlowNodeId,
76}
77
78impl FlowNodeStateKeyRange {
79 pub fn new(node: FlowNodeId) -> Self {
80 Self {
81 node,
82 }
83 }
84
85 fn decode_key(key: &EncodedKey) -> Option<Self> {
86 let mut de = KeyDeserializer::from_bytes(key.as_slice());
87
88 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
89 if kind != FlowNodeStateKey::KIND {
90 return None;
91 }
92
93 let node_id = de.read_u64().ok()?;
94
95 Some(Self {
96 node: FlowNodeId(node_id),
97 })
98 }
99}
100
101impl EncodableKeyRange for FlowNodeStateKeyRange {
102 const KIND: KeyKind = KeyKind::FlowNodeState;
103
104 fn start(&self) -> Option<EncodedKey> {
105 let mut serializer = KeySerializer::with_capacity(9);
106 serializer.extend_u8(Self::KIND as u8).extend_u64(self.node.0);
107 Some(serializer.to_encoded_key())
108 }
109
110 fn end(&self) -> Option<EncodedKey> {
111 let mut serializer = KeySerializer::with_capacity(9);
112 serializer.extend_u8(Self::KIND as u8).extend_u64(self.node.0.wrapping_sub(1));
113 Some(serializer.to_encoded_key())
114 }
115
116 fn decode(range: &EncodedKeyRange) -> (Option<Self>, Option<Self>)
117 where
118 Self: Sized,
119 {
120 let start_key = match &range.start {
121 Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
122 Bound::Unbounded => None,
123 };
124
125 let end_key = match &range.end {
126 Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
127 Bound::Unbounded => None,
128 };
129
130 (start_key, end_key)
131 }
132}
133
134#[cfg(test)]
135pub mod tests {
136 use reifydb_codec::key::encoded::{EncodedKey, EncodedKeyRange};
137
138 use super::{EncodableKey, EncodableKeyRange, FlowNodeStateKey, FlowNodeStateKeyRange};
139 use crate::interface::catalog::flow::FlowNodeId;
140
141 #[test]
142 fn test_encode_decode() {
143 let key = FlowNodeStateKey {
144 node: FlowNodeId(0xDEADBEEF),
145 key: vec![1, 2, 3, 4],
146 };
147 let encoded = key.encode();
148
149 assert_eq!(encoded[0], 0xEC);
150
151 let decoded = FlowNodeStateKey::decode(&encoded).unwrap();
152 assert_eq!(decoded.node.0, 0xDEADBEEF);
153 assert_eq!(decoded.key, vec![1, 2, 3, 4]);
154 }
155
156 #[test]
157 fn test_encode_decode_empty_key() {
158 let key = FlowNodeStateKey {
159 node: FlowNodeId(0xDEADBEEF),
160 key: vec![],
161 };
162 let encoded = key.encode();
163
164 let decoded = FlowNodeStateKey::decode(&encoded).unwrap();
165 assert_eq!(decoded.node.0, 0xDEADBEEF);
166 assert_eq!(decoded.key, Vec::<u8>::new());
167 }
168
169 #[test]
170 fn test_new() {
171 let key = FlowNodeStateKey::new(FlowNodeId(42), vec![5, 6, 7]);
172 assert_eq!(key.node.0, 42);
173 assert_eq!(key.key, vec![5, 6, 7]);
174 }
175
176 #[test]
177 fn test_new_empty() {
178 let key = FlowNodeStateKey::new_empty(FlowNodeId(42));
179 assert_eq!(key.node.0, 42);
180 assert_eq!(key.key, Vec::<u8>::new());
181 }
182
183 #[test]
184 fn test_roundtrip() {
185 let original = FlowNodeStateKey {
186 node: FlowNodeId(999_999_999),
187 key: vec![10, 20, 30, 40, 50],
188 };
189 let encoded = original.encode();
190 let decoded = FlowNodeStateKey::decode(&encoded).unwrap();
191 assert_eq!(original, decoded);
192 }
193
194 #[test]
195 fn test_decode_invalid_version() {
196 let mut encoded = Vec::new();
197 encoded.push(0xFF);
198 encoded.push(0xEC);
199 encoded.extend(&999u64.to_be_bytes());
200 let key = EncodedKey::new(encoded);
201 assert!(FlowNodeStateKey::decode(&key).is_none());
202 }
203
204 #[test]
205 fn test_decode_invalid_kind() {
206 let mut encoded = Vec::new();
207 encoded.push(0xFE);
208 encoded.push(0xFF);
209 encoded.extend(&999u64.to_be_bytes());
210 let key = EncodedKey::new(encoded);
211 assert!(FlowNodeStateKey::decode(&key).is_none());
212 }
213
214 #[test]
215 fn test_decode_too_short() {
216 let mut encoded = Vec::new();
217 encoded.push(0xFE);
218 encoded.push(0xEC);
219 encoded.extend(&999u32.to_be_bytes());
220 let key = EncodedKey::new(encoded);
221 assert!(FlowNodeStateKey::decode(&key).is_none());
222 }
223
224 #[test]
225 fn test_flow_node_state_key_range() {
226 let node = FlowNodeId(42);
227 let range = FlowNodeStateKeyRange::new(node);
228
229 let start = range.start().unwrap();
230 let decoded_start = FlowNodeStateKey::decode(&start).unwrap();
231 assert_eq!(decoded_start.node, node);
232 assert_eq!(decoded_start.key, Vec::<u8>::new());
233
234 let end = range.end().unwrap();
235 let decoded_end = FlowNodeStateKey::decode(&end).unwrap();
236 assert_eq!(decoded_end.node.0, 41);
237 assert_eq!(decoded_end.key, Vec::<u8>::new());
238 }
239
240 #[test]
241 fn test_flow_node_state_key_range_decode() {
242 let node = FlowNodeId(100);
243 let range = FlowNodeStateKeyRange::new(node);
244
245 let encoded_range = EncodedKeyRange::start_end(range.start(), range.end());
246
247 let (start_decoded, end_decoded) = FlowNodeStateKeyRange::decode(&encoded_range);
248
249 assert!(start_decoded.is_some());
250 assert_eq!(start_decoded.unwrap().node, node);
251
252 assert!(end_decoded.is_some());
253 assert_eq!(end_decoded.unwrap().node.0, 99);
254 }
255
256 #[test]
257 fn test_node_range_method() {
258 let node = FlowNodeId(555);
259 let range = FlowNodeStateKey::node_range(node);
260
261 let (start_range, end_range) = FlowNodeStateKeyRange::decode(&range);
262
263 assert!(start_range.is_some());
264 assert_eq!(start_range.unwrap().node, node);
265
266 assert!(end_range.is_some());
267 assert_eq!(end_range.unwrap().node.0, 554);
268 }
269}