reifydb_core/key/
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::QueueId;
12
13#[derive(Debug, Clone, PartialEq)]
14pub struct QueueKey {
15 pub queue: QueueId,
16}
17
18impl QueueKey {
19 pub fn new(queue: QueueId) -> Self {
20 Self {
21 queue,
22 }
23 }
24
25 pub fn encoded(queue: impl Into<QueueId>) -> EncodedKey {
26 Self::new(queue.into()).encode()
27 }
28
29 pub fn full_scan() -> EncodedKeyRange {
30 EncodedKeyRange::start_end(Some(Self::queue_start()), Some(Self::queue_end()))
31 }
32
33 fn queue_start() -> EncodedKey {
34 let mut serializer = KeySerializer::with_capacity(1);
35 serializer.extend_u8(Self::KIND as u8);
36 serializer.to_encoded_key()
37 }
38
39 fn queue_end() -> EncodedKey {
40 let mut serializer = KeySerializer::with_capacity(1);
41 serializer.extend_u8(Self::KIND as u8 - 1);
42 serializer.to_encoded_key()
43 }
44}
45
46impl EncodableKey for QueueKey {
47 const KIND: KeyKind = KeyKind::Queue;
48
49 fn encode(&self) -> EncodedKey {
50 let mut serializer = KeySerializer::with_capacity(9);
51 serializer.extend_u8(Self::KIND as u8).extend_u64(self.queue);
52 serializer.to_encoded_key()
53 }
54
55 fn decode(key: &EncodedKey) -> Option<Self> {
56 let mut de = KeyDeserializer::from_bytes(key.as_slice());
57
58 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
59 if kind != Self::KIND {
60 return None;
61 }
62
63 let queue = de.read_u64().ok()?;
64
65 Some(Self {
66 queue: QueueId(queue),
67 })
68 }
69}
70
71#[cfg(test)]
72mod tests {
73 use std::ops::Bound;
74
75 use super::*;
76
77 #[test]
78 fn test_encode_decode_roundtrip() {
79 let encoded = QueueKey::encoded(QueueId(42));
81 let decoded = QueueKey::decode(&encoded).unwrap();
82 assert_eq!(decoded.queue, QueueId(42));
83 }
84
85 #[test]
86 fn test_decode_rejects_foreign_kind() {
87 let mut serializer = KeySerializer::with_capacity(9);
90 serializer.extend_u8(KeyKind::NamespaceQueue as u8).extend_u64(7u64);
91 assert!(QueueKey::decode(&serializer.to_encoded_key()).is_none());
92 }
93
94 #[test]
95 fn test_full_scan_brackets_every_queue_key() {
96 let range = QueueKey::full_scan();
99
100 let Bound::Included(start) = &range.start else {
101 panic!("expected an included start bound")
102 };
103 let Bound::Included(end) = &range.end else {
104 panic!("expected an included end bound")
105 };
106
107 assert_eq!(start.as_slice(), &[!(KeyKind::Queue as u8)]);
108 assert_eq!(end.as_slice(), &[!(KeyKind::Queue as u8 - 1)]);
109 assert!(start.as_slice() < end.as_slice(), "the range must be non-empty under byte order");
110
111 for id in [QueueId(1), QueueId(u64::MAX)] {
112 let key = QueueKey::encoded(id);
113 assert!(
114 key.as_slice() >= start.as_slice() && key.as_slice() <= end.as_slice(),
115 "queue {id:?} must fall inside the scan range"
116 );
117 }
118 }
119
120 #[test]
121 fn test_full_scan_excludes_the_neighbouring_kind() {
122 let range = QueueKey::full_scan();
125 let Bound::Included(start) = &range.start else {
126 panic!("expected an included start bound")
127 };
128 let Bound::Included(end) = &range.end else {
129 panic!("expected an included end bound")
130 };
131
132 let mut serializer = KeySerializer::with_capacity(9);
133 serializer.extend_u8(KeyKind::NamespaceQueue as u8).extend_u64(1u64);
134 let foreign = serializer.to_encoded_key();
135
136 assert!(
137 foreign.as_slice() < start.as_slice() || foreign.as_slice() > end.as_slice(),
138 "a NamespaceQueue key must fall outside the QueueKey scan range"
139 );
140 }
141}