Skip to main content

reifydb_core/key/
namespace_queue.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// The link row is what makes a queue findable by name; losing either component makes DROP
84		// NAMESPACE miss its queues.
85		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		// Keys are stored bitwise-inverted, so a bound derived with the wrong sign would make DROP
94		// NAMESPACE either miss its queues or reach into a sibling.
95		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}