reifydb_core/key/
namespace_queue.rs1use reifydb_codec::key::{
5 deserializer::KeyDeserializer,
6 encoded::{EncodedKey, EncodedKeyRange},
7 serializer::KeySerializer,
8};
9
10use super::{EncodableKey, KeyKind};
11use crate::interface::catalog::id::{NamespaceId, QueueId};
12
13#[derive(Debug, Clone, PartialEq)]
14pub struct NamespaceQueueKey {
15 pub namespace: NamespaceId,
16 pub queue: QueueId,
17}
18
19impl NamespaceQueueKey {
20 pub fn new(namespace: NamespaceId, queue: QueueId) -> Self {
21 Self {
22 namespace,
23 queue,
24 }
25 }
26
27 pub fn encoded(namespace: impl Into<NamespaceId>, queue: impl Into<QueueId>) -> EncodedKey {
28 Self::new(namespace.into(), queue.into()).encode()
29 }
30
31 pub fn full_scan(namespace: NamespaceId) -> EncodedKeyRange {
32 EncodedKeyRange::start_end(Some(Self::link_start(namespace)), Some(Self::link_end(namespace)))
33 }
34
35 fn link_start(namespace: NamespaceId) -> EncodedKey {
36 let mut serializer = KeySerializer::with_capacity(9);
37 serializer.extend_u8(Self::KIND as u8).extend_u64(namespace);
38 serializer.to_encoded_key()
39 }
40
41 fn link_end(namespace: NamespaceId) -> EncodedKey {
42 let mut serializer = KeySerializer::with_capacity(9);
43 serializer.extend_u8(Self::KIND as u8).extend_u64(*namespace - 1);
44 serializer.to_encoded_key()
45 }
46}
47
48impl EncodableKey for NamespaceQueueKey {
49 const KIND: KeyKind = KeyKind::NamespaceQueue;
50
51 fn encode(&self) -> EncodedKey {
52 let mut serializer = KeySerializer::with_capacity(17);
53 serializer.extend_u8(Self::KIND as u8).extend_u64(self.namespace).extend_u64(self.queue);
54 serializer.to_encoded_key()
55 }
56
57 fn decode(key: &EncodedKey) -> Option<Self> {
58 let mut de = KeyDeserializer::from_bytes(key.as_slice());
59
60 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
61 if kind != Self::KIND {
62 return None;
63 }
64
65 let namespace = de.read_u64().ok()?;
66 let queue = de.read_u64().ok()?;
67
68 Some(Self {
69 namespace: NamespaceId(namespace),
70 queue: QueueId(queue),
71 })
72 }
73}
74
75#[cfg(test)]
76mod tests {
77 use std::ops::Bound;
78
79 use super::*;
80
81 #[test]
82 fn test_encode_decode_roundtrip() {
83 let encoded = NamespaceQueueKey::encoded(NamespaceId(3), QueueId(42));
86 let decoded = NamespaceQueueKey::decode(&encoded).unwrap();
87 assert_eq!(decoded.namespace, NamespaceId(3));
88 assert_eq!(decoded.queue, QueueId(42));
89 }
90
91 #[test]
92 fn test_full_scan_contains_only_the_target_namespace() {
93 let range = NamespaceQueueKey::full_scan(NamespaceId(3));
96 let Bound::Included(start) = &range.start else {
97 panic!("expected an included start bound")
98 };
99 let Bound::Included(end) = &range.end else {
100 panic!("expected an included end bound")
101 };
102
103 assert!(start.as_slice() < end.as_slice(), "the range must be non-empty under byte order");
104
105 for queue in [QueueId(1), QueueId(u64::MAX)] {
106 let inside = NamespaceQueueKey::encoded(NamespaceId(3), queue);
107 assert!(
108 inside.as_slice() >= start.as_slice() && inside.as_slice() <= end.as_slice(),
109 "queue {queue:?} in namespace 3 must fall inside the scan range"
110 );
111 }
112
113 for namespace in [NamespaceId(2), NamespaceId(4)] {
114 let neighbour = NamespaceQueueKey::encoded(namespace, QueueId(1));
115 assert!(
116 neighbour.as_slice() < start.as_slice() || neighbour.as_slice() > end.as_slice(),
117 "namespace {namespace:?} must fall outside namespace 3's scan range"
118 );
119 }
120 }
121}