reifydb_core/key/
flow_node_internal_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 FlowNodeInternalStateKey {
17 pub node: FlowNodeId,
18 pub key: Vec<u8>,
19}
20
21impl EncodableKey for FlowNodeInternalStateKey {
22 const KIND: KeyKind = KeyKind::FlowNodeInternalState;
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 FlowNodeInternalStateKey {
49 pub const ROW_NUMBER_COUNTER_TAG: u8 = b'C';
50
51 pub const ROW_NUMBER_MAPPING_TAG: u8 = b'M';
52
53 pub const WINDOW_META_TAG: u8 = b'W';
54
55 pub const WINDOW_EXPIRY_TAG: u8 = b'X';
56
57 pub const WINDOW_RUNNING_TAG: u8 = b'R';
58
59 pub const WINDOW_COORD_TAG: u8 = b'S';
60
61 pub const WINDOW_ROW_STATE_TAG: u8 = b'A';
62
63 pub const WINDOW_EMIT_TAG: u8 = b'E';
64
65 pub const GATE_VISIBILITY_TAG: u8 = b'G';
66
67 pub fn is_row_number_counter(&self) -> bool {
68 self.key.as_slice() == [Self::ROW_NUMBER_COUNTER_TAG]
69 }
70
71 pub fn is_row_number_mapping(&self) -> bool {
72 self.key.first() == Some(&Self::ROW_NUMBER_MAPPING_TAG)
73 }
74
75 pub fn is_window_meta(&self) -> bool {
76 self.key.first() == Some(&Self::WINDOW_META_TAG)
77 }
78
79 pub fn is_window_expiry(&self) -> bool {
80 self.key.first() == Some(&Self::WINDOW_EXPIRY_TAG)
81 }
82
83 pub fn is_gate_visibility(&self) -> bool {
84 self.key.first() == Some(&Self::GATE_VISIBILITY_TAG)
85 }
86
87 pub fn new(node: FlowNodeId, key: Vec<u8>) -> Self {
88 Self {
89 node,
90 key,
91 }
92 }
93
94 pub fn new_empty(node: FlowNodeId) -> Self {
95 Self {
96 node,
97 key: Vec::new(),
98 }
99 }
100
101 pub fn encoded(node: impl Into<FlowNodeId>, key: impl Into<Vec<u8>>) -> EncodedKey {
102 Self::new(node.into(), key.into()).encode()
103 }
104
105 pub fn node_range(node: FlowNodeId) -> EncodedKeyRange {
106 let range = FlowNodeInternalStateKeyRange::new(node);
107 EncodedKeyRange::start_end(range.start(), range.end())
108 }
109}
110
111#[derive(Debug, Clone, PartialEq, Eq)]
112pub struct FlowNodeInternalStateKeyRange {
113 pub node: FlowNodeId,
114}
115
116impl FlowNodeInternalStateKeyRange {
117 pub fn new(node: FlowNodeId) -> Self {
118 Self {
119 node,
120 }
121 }
122
123 fn decode_key(key: &EncodedKey) -> Option<Self> {
124 let mut de = KeyDeserializer::from_bytes(key.as_slice());
125
126 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
127 if kind != FlowNodeInternalStateKey::KIND {
128 return None;
129 }
130
131 let node_id = de.read_u64().ok()?;
132
133 Some(Self {
134 node: FlowNodeId(node_id),
135 })
136 }
137}
138
139impl EncodableKeyRange for FlowNodeInternalStateKeyRange {
140 const KIND: KeyKind = KeyKind::FlowNodeInternalState;
141
142 fn start(&self) -> Option<EncodedKey> {
143 let mut serializer = KeySerializer::with_capacity(9);
144 serializer.extend_u8(Self::KIND as u8).extend_u64(self.node.0);
145 Some(serializer.to_encoded_key())
146 }
147
148 fn end(&self) -> Option<EncodedKey> {
149 let mut serializer = KeySerializer::with_capacity(9);
150 serializer.extend_u8(Self::KIND as u8).extend_u64(self.node.0.wrapping_sub(1));
151 Some(serializer.to_encoded_key())
152 }
153
154 fn decode(range: &EncodedKeyRange) -> (Option<Self>, Option<Self>)
155 where
156 Self: Sized,
157 {
158 let start_key = match &range.start {
159 Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
160 Bound::Unbounded => None,
161 };
162
163 let end_key = match &range.end {
164 Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
165 Bound::Unbounded => None,
166 };
167
168 (start_key, end_key)
169 }
170}
171
172#[cfg(test)]
173pub mod tests {
174 use reifydb_codec::key::encoded::{EncodedKey, EncodedKeyRange};
175
176 use super::{EncodableKey, EncodableKeyRange, FlowNodeInternalStateKey, FlowNodeInternalStateKeyRange};
177 use crate::interface::catalog::flow::FlowNodeId;
178
179 #[test]
180 fn test_encode_decode() {
181 let key = FlowNodeInternalStateKey {
182 node: FlowNodeId(0xDEADBEEF),
183 key: vec![1, 2, 3, 4],
184 };
185 let encoded = key.encode();
186
187 assert_eq!(encoded[0], 0xE0);
188
189 let decoded = FlowNodeInternalStateKey::decode(&encoded).unwrap();
190 assert_eq!(decoded.node.0, 0xDEADBEEF);
191 assert_eq!(decoded.key, vec![1, 2, 3, 4]);
192 }
193
194 #[test]
195 fn test_encode_decode_empty_key() {
196 let key = FlowNodeInternalStateKey {
197 node: FlowNodeId(0xDEADBEEF),
198 key: vec![],
199 };
200 let encoded = key.encode();
201
202 let decoded = FlowNodeInternalStateKey::decode(&encoded).unwrap();
203 assert_eq!(decoded.node.0, 0xDEADBEEF);
204 assert_eq!(decoded.key, Vec::<u8>::new());
205 }
206
207 #[test]
208 fn test_new() {
209 let key = FlowNodeInternalStateKey::new(FlowNodeId(42), vec![5, 6, 7]);
210 assert_eq!(key.node.0, 42);
211 assert_eq!(key.key, vec![5, 6, 7]);
212 }
213
214 #[test]
215 fn test_new_empty() {
216 let key = FlowNodeInternalStateKey::new_empty(FlowNodeId(42));
217 assert_eq!(key.node.0, 42);
218 assert_eq!(key.key, Vec::<u8>::new());
219 }
220
221 #[test]
222 fn test_roundtrip() {
223 let original = FlowNodeInternalStateKey {
224 node: FlowNodeId(999_999_999),
225 key: vec![10, 20, 30, 40, 50],
226 };
227 let encoded = original.encode();
228 let decoded = FlowNodeInternalStateKey::decode(&encoded).unwrap();
229 assert_eq!(original, decoded);
230 }
231
232 #[test]
233 fn test_decode_invalid_version() {
234 let mut encoded = Vec::new();
235 encoded.push(0xFF);
236 encoded.push(0xE5);
237 encoded.extend(&999u64.to_be_bytes());
238 let key = EncodedKey::new(encoded);
239 assert!(FlowNodeInternalStateKey::decode(&key).is_none());
240 }
241
242 #[test]
243 fn test_decode_invalid_kind() {
244 let mut encoded = Vec::new();
245 encoded.push(0xFE);
246 encoded.push(0xFF);
247 encoded.extend(&999u64.to_be_bytes());
248 let key = EncodedKey::new(encoded);
249 assert!(FlowNodeInternalStateKey::decode(&key).is_none());
250 }
251
252 #[test]
253 fn test_decode_too_short() {
254 let mut encoded = Vec::new();
255 encoded.push(0xFE);
256 encoded.push(0xE5);
257 encoded.extend(&999u32.to_be_bytes());
258 let key = EncodedKey::new(encoded);
259 assert!(FlowNodeInternalStateKey::decode(&key).is_none());
260 }
261
262 #[test]
263 fn test_flow_node_internal_state_key_range() {
264 let node = FlowNodeId(42);
265 let range = FlowNodeInternalStateKeyRange::new(node);
266
267 let start = range.start().unwrap();
268 let decoded_start = FlowNodeInternalStateKey::decode(&start).unwrap();
269 assert_eq!(decoded_start.node, node);
270 assert_eq!(decoded_start.key, Vec::<u8>::new());
271
272 let end = range.end().unwrap();
273 let decoded_end = FlowNodeInternalStateKey::decode(&end).unwrap();
274 assert_eq!(decoded_end.node.0, 41);
275 assert_eq!(decoded_end.key, Vec::<u8>::new());
276 }
277
278 #[test]
279 fn test_flow_node_internal_state_key_range_decode() {
280 let node = FlowNodeId(100);
281 let range = FlowNodeInternalStateKeyRange::new(node);
282
283 let encoded_range = EncodedKeyRange::start_end(range.start(), range.end());
284
285 let (start_decoded, end_decoded) = FlowNodeInternalStateKeyRange::decode(&encoded_range);
286
287 assert!(start_decoded.is_some());
288 assert_eq!(start_decoded.unwrap().node, node);
289
290 assert!(end_decoded.is_some());
291 assert_eq!(end_decoded.unwrap().node.0, 99);
292 }
293}