Skip to main content

reifydb_core/key/
flow.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_codec::key::encoded::EncodedKey;
5use reifydb_macro::KeyCodec;
6
7use super::KeyTag;
8use crate::{
9	interface::catalog::flow::{FlowEdgeId, FlowId},
10	key::{
11		any::{Field, KeyFields, Width},
12		bound::TaggedKeyBoundRange,
13	},
14};
15
16#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
17#[key(tag = Flow)]
18pub struct FlowKey {
19	pub flow: FlowId,
20}
21
22impl FlowKey {
23	pub fn new(flow: impl Into<FlowId>) -> Self {
24		Self {
25			flow: flow.into(),
26		}
27	}
28
29	pub fn encoded(flow: impl Into<FlowId>) -> EncodedKey {
30		Self::new(flow).encode()
31	}
32
33	pub fn full_scan() -> TaggedKeyBoundRange {
34		TaggedKeyBoundRange::kind(Self::TAG)
35	}
36}
37
38#[cfg(test)]
39pub mod flow_key_tests {
40	use super::FlowKey;
41	use crate::interface::catalog::flow::FlowId;
42
43	#[test]
44	fn test_encode_decode() {
45		let key = FlowKey {
46			flow: FlowId(0x1234),
47		};
48		let encoded = key.encode();
49		let decoded = FlowKey::decode(&encoded).unwrap();
50		assert_eq!(decoded.flow, FlowId(0x1234));
51		assert_eq!(key, decoded);
52	}
53
54	#[test]
55	fn test_order_preserving() {
56		let key1 = FlowKey {
57			flow: FlowId(1),
58		};
59		let key2 = FlowKey {
60			flow: FlowId(2),
61		};
62
63		let encoded1 = key1.encode();
64		let encoded2 = key2.encode();
65
66		assert!(encoded2 < encoded1, "ordering not preserved");
67	}
68}
69
70#[cfg(test)]
71mod verify_byte_identical_flow_key {
72	use reifydb_codec::key::serializer::KeySerializer;
73
74	use super::FlowKey;
75	use crate::interface::catalog::flow::FlowId;
76
77	fn legacy_encode(key: &FlowKey) -> Vec<u8> {
78		let mut serializer = KeySerializer::with_capacity(9);
79		serializer.extend_u8(FlowKey::TAG as u8).extend_u64(key.flow);
80		serializer.to_encoded_key().as_slice().to_vec()
81	}
82
83	#[test]
84	fn matches_legacy_byte_layout() {
85		for flow in [0u64, 1, 42, 0x1234, u64::MAX] {
86			let key = FlowKey {
87				flow: FlowId(flow),
88			};
89			assert_eq!(legacy_encode(&key), key.encode().as_slice().to_vec(), "flow={flow:#x}");
90		}
91	}
92}
93
94#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
95#[key(tag = FlowEdge)]
96pub struct FlowEdgeKey {
97	pub edge: FlowEdgeId,
98}
99
100impl FlowEdgeKey {
101	pub fn new(edge: impl Into<FlowEdgeId>) -> Self {
102		Self {
103			edge: edge.into(),
104		}
105	}
106
107	pub fn encoded(edge: impl Into<FlowEdgeId>) -> EncodedKey {
108		Self::new(edge).encode()
109	}
110
111	pub fn full_scan() -> TaggedKeyBoundRange {
112		TaggedKeyBoundRange::kind(Self::TAG)
113	}
114}
115
116#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
117#[key(tag = FlowEdgeByFlow)]
118pub struct FlowEdgeByFlowKey {
119	pub flow: FlowId,
120	pub edge: FlowEdgeId,
121}
122
123impl FlowEdgeByFlowKey {
124	pub fn new(flow: impl Into<FlowId>, edge: impl Into<FlowEdgeId>) -> Self {
125		Self {
126			flow: flow.into(),
127			edge: edge.into(),
128		}
129	}
130
131	pub fn encoded(flow: impl Into<FlowId>, edge: impl Into<FlowEdgeId>) -> EncodedKey {
132		Self::new(flow, edge).encode()
133	}
134
135	pub fn full_scan(flow: FlowId) -> TaggedKeyBoundRange {
136		TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, flow.0 as u128)])
137	}
138}
139
140#[cfg(test)]
141pub mod flow_edge_by_flow_key_tests {
142	use super::{FlowEdgeByFlowKey, FlowEdgeKey};
143	use crate::interface::catalog::flow::{FlowEdgeId, FlowId};
144
145	#[test]
146	fn test_flow_edge_key_encode_decode() {
147		let key = FlowEdgeKey {
148			edge: FlowEdgeId(0x1234),
149		};
150		let encoded = key.encode();
151		let decoded = FlowEdgeKey::decode(&encoded).unwrap();
152		assert_eq!(decoded.edge, FlowEdgeId(0x1234));
153		assert_eq!(key, decoded);
154	}
155
156	#[test]
157	fn test_flow_edge_key_order_preserving() {
158		let key1 = FlowEdgeKey {
159			edge: FlowEdgeId(1),
160		};
161		let key2 = FlowEdgeKey {
162			edge: FlowEdgeId(2),
163		};
164
165		let encoded1 = key1.encode();
166		let encoded2 = key2.encode();
167
168		assert!(encoded2 < encoded1, "ordering not preserved");
169	}
170
171	#[test]
172	fn test_flow_edge_by_flow_key_encode_decode() {
173		let key = FlowEdgeByFlowKey {
174			flow: FlowId(0x42),
175			edge: FlowEdgeId(0x1234),
176		};
177		let encoded = key.encode();
178		let decoded = FlowEdgeByFlowKey::decode(&encoded).unwrap();
179		assert_eq!(decoded.flow, FlowId(0x42));
180		assert_eq!(decoded.edge, FlowEdgeId(0x1234));
181		assert_eq!(key, decoded);
182	}
183
184	#[test]
185	fn test_flow_edge_by_flow_key_order_preserving() {
186		let key1 = FlowEdgeByFlowKey {
187			flow: FlowId(1),
188			edge: FlowEdgeId(100),
189		};
190		let key2 = FlowEdgeByFlowKey {
191			flow: FlowId(1),
192			edge: FlowEdgeId(200),
193		};
194		let key3 = FlowEdgeByFlowKey {
195			flow: FlowId(2),
196			edge: FlowEdgeId(1),
197		};
198
199		let encoded1 = key1.encode();
200		let encoded2 = key2.encode();
201		let encoded3 = key3.encode();
202
203		assert!(encoded2 < encoded1, "edge ordering not preserved within same flow");
204		assert!(encoded3 < encoded2, "flow ordering not preserved");
205	}
206}
207
208#[cfg(test)]
209mod verify_byte_identical_flow_edge_by_flow_key {
210	use reifydb_codec::key::serializer::KeySerializer;
211
212	use super::{FlowEdgeByFlowKey, FlowEdgeKey};
213	use crate::interface::catalog::flow::{FlowEdgeId, FlowId};
214
215	fn legacy_encode_edge(key: &FlowEdgeKey) -> Vec<u8> {
216		let mut serializer = KeySerializer::with_capacity(9);
217		serializer.extend_u8(FlowEdgeKey::TAG as u8).extend_u64(key.edge);
218		serializer.to_encoded_key().as_slice().to_vec()
219	}
220
221	fn legacy_encode_by_flow(key: &FlowEdgeByFlowKey) -> Vec<u8> {
222		let mut serializer = KeySerializer::with_capacity(17);
223		serializer.extend_u8(FlowEdgeByFlowKey::TAG as u8).extend_u64(key.flow).extend_u64(key.edge);
224		serializer.to_encoded_key().as_slice().to_vec()
225	}
226
227	#[test]
228	fn flow_edge_key_matches_legacy_byte_layout() {
229		for edge in [0u64, 1, 42, 0x1234, u64::MAX] {
230			let key = FlowEdgeKey {
231				edge: FlowEdgeId(edge),
232			};
233			assert_eq!(legacy_encode_edge(&key), key.encode().as_slice().to_vec(), "edge={edge:#x}");
234		}
235	}
236
237	#[test]
238	fn flow_edge_by_flow_key_matches_legacy_byte_layout() {
239		for (flow, edge) in [(0u64, 0u64), (1, 2), (0x42, 0x1234), (u64::MAX, u64::MAX)] {
240			let key = FlowEdgeByFlowKey {
241				flow: FlowId(flow),
242				edge: FlowEdgeId(edge),
243			};
244			assert_eq!(
245				legacy_encode_by_flow(&key),
246				key.encode().as_slice().to_vec(),
247				"flow={flow:#x} edge={edge:#x}"
248			);
249		}
250	}
251}
252
253#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
254#[key(tag = FlowVersion)]
255pub struct FlowVersionKey {
256	pub flow: FlowId,
257}
258
259impl FlowVersionKey {
260	pub fn new(flow: impl Into<FlowId>) -> Self {
261		Self {
262			flow: flow.into(),
263		}
264	}
265
266	pub fn encoded(flow: impl Into<FlowId>) -> EncodedKey {
267		Self::new(flow).encode()
268	}
269}
270
271#[cfg(test)]
272pub mod flow_version_key_tests {
273	use super::FlowVersionKey;
274	use crate::interface::catalog::flow::FlowId;
275
276	#[test]
277	fn test_encode_decode() {
278		let key = FlowVersionKey {
279			flow: FlowId(0x1234),
280		};
281		let encoded = key.encode();
282		let decoded = FlowVersionKey::decode(&encoded).unwrap();
283		assert_eq!(decoded.flow, FlowId(0x1234));
284		assert_eq!(key, decoded);
285	}
286
287	#[test]
288	fn test_new_and_encoded() {
289		let key = FlowVersionKey::new(42u64);
290		assert_eq!(key.flow, FlowId(42));
291
292		let encoded = FlowVersionKey::encoded(42u64);
293		let decoded = FlowVersionKey::decode(&encoded).unwrap();
294		assert_eq!(decoded.flow, FlowId(42));
295	}
296
297	#[test]
298	fn test_order_preserving() {
299		let key1 = FlowVersionKey {
300			flow: FlowId(1),
301		};
302		let key2 = FlowVersionKey {
303			flow: FlowId(2),
304		};
305
306		let encoded1 = key1.encode();
307		let encoded2 = key2.encode();
308
309		assert!(encoded2 < encoded1, "ordering not preserved");
310	}
311}
312
313#[cfg(test)]
314mod verify_byte_identical_flow_version_key {
315	use reifydb_codec::key::serializer::KeySerializer;
316
317	use super::FlowVersionKey;
318	use crate::interface::catalog::flow::FlowId;
319
320	fn legacy_encode(key: &FlowVersionKey) -> Vec<u8> {
321		let mut serializer = KeySerializer::with_capacity(9);
322		serializer.extend_u8(FlowVersionKey::TAG as u8).extend_u64(key.flow);
323		serializer.to_encoded_key().as_slice().to_vec()
324	}
325
326	#[test]
327	fn matches_legacy_byte_layout() {
328		for flow in [0u64, 1, 42, 0x1234, u64::MAX] {
329			let key = FlowVersionKey {
330				flow: FlowId(flow),
331			};
332			assert_eq!(legacy_encode(&key), key.encode().as_slice().to_vec(), "flow={flow:#x}");
333		}
334	}
335}