Skip to main content

reifydb_core/key/
flow_node_internal_state.rs

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