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