Skip to main content

reifydb_core/key/
flow_node_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 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}