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 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}