1use std::ops::Bound;
5
6use reifydb_codec::key::{
7 deserializer::KeyDeserializer,
8 encode_u8,
9 encoded::{EncodedKey, EncodedKeyRange},
10 serializer::KeySerializer,
11};
12
13use super::{EncodableKey, KeyKind};
14use crate::interface::catalog::flow::OperatorId;
15
16#[repr(transparent)]
17#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
18pub struct GroupId(pub u64);
19
20impl GroupId {
21 pub const ROOT: Self = Self(0);
22
23 pub const FIRST: Self = Self(1);
24
25 pub fn is_root(&self) -> bool {
26 *self == Self::ROOT
27 }
28}
29
30#[derive(Debug, Clone, Default, PartialEq, Eq)]
31pub struct GroupSet(Vec<GroupId>);
32
33impl GroupSet {
34 pub fn new(groups: impl IntoIterator<Item = GroupId>) -> Self {
35 let mut groups: Vec<GroupId> = groups.into_iter().filter(|g| !g.is_root()).collect();
36 groups.sort_unstable();
37 groups.dedup();
38 Self(groups)
39 }
40
41 pub fn contains(&self, group: GroupId) -> bool {
42 self.0.binary_search(&group).is_ok()
43 }
44
45 pub fn as_slice(&self) -> &[GroupId] {
46 &self.0
47 }
48
49 pub fn len(&self) -> usize {
50 self.0.len()
51 }
52
53 pub fn is_empty(&self) -> bool {
54 self.0.is_empty()
55 }
56
57 pub fn as_raw_parts(&self) -> (*const u64, usize) {
58 (self.0.as_ptr() as *const u64, self.0.len())
59 }
60}
61
62pub fn group_data_of_inner(inner: &[u8]) -> Option<GroupId> {
63 let mut de = KeyDeserializer::from_bytes(inner);
64 let group = GroupId(de.read_u64().ok()?);
65 let keyspace = Keyspace(de.read_u8().ok()?);
66 if !keyspace.is_data() {
67 return None;
68 }
69 inner.starts_with(&group_inner_prefix(group)).then_some(group)
70}
71
72#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
73pub struct Keyspace(pub u8);
74
75impl Keyspace {
76 pub const HIGHEST_DATA: u8 = 0x7F;
77
78 pub const ROW_NUMBER_MAPPING: Self = Self(0xFE);
79
80 pub const GROUP_DICTIONARY: Self = Self(0xFD);
81
82 pub const NODE_COUNTER: Self = Self(0xFC);
83
84 pub const GROUP_RECORD: Self = Self(0xFB);
85
86 pub const SOURCE_WATERMARK: Self = Self(0xFA);
87
88 pub const TIMER_WHEEL: Self = Self(0xF9);
89
90 pub const TIMER_INDEX: Self = Self(0xF8);
91
92 pub const ACCUMULATOR: Self = Self(0x10);
93
94 pub const BUFFER: Self = Self(0x11);
95
96 pub const RUNNING: Self = Self(0x12);
97
98 pub const EMIT: Self = Self(0x13);
99
100 pub const EXPIRY: Self = Self(0x14);
101
102 pub const COUNT: Self = Self(0x16);
103
104 pub const ROW_INDEX: Self = Self(0x17);
105
106 pub const SESSION: Self = Self(0x18);
107
108 pub const ROLLING_META: Self = Self(0x19);
109
110 pub const ENGINE_META: Self = Self(0x1A);
111
112 pub const DISTINCT_ENTRY: Self = Self(0x1B);
113
114 pub const WINDOW_META: Self = Self(0x1C);
115
116 pub const JOIN_LEFT: Self = Self(0x1D);
117
118 pub const JOIN_RIGHT: Self = Self(0x1E);
119
120 pub const JOIN_SCHEMA: Self = Self(0x1F);
121
122 pub const RINGBUFFER_FORWARD: Self = Self(0x20);
123
124 pub const RINGBUFFER_ENTRY: Self = Self(0x21);
125
126 pub const GATE_VISIBILITY: Self = Self(0x22);
127
128 pub const DISTINCT_LAYOUT: Self = Self(0x23);
129
130 pub const RINGBUFFER_EXPIRY: Self = Self(0x24);
131
132 pub const RINGBUFFER_TTL_ARM: Self = Self(0x25);
133
134 pub const SEAL_LEDGER: Self = Self(0x26);
135
136 pub const JOIN_PUBLISHED: Self = Self(0x27);
137
138 pub const JOIN_PIN: Self = Self(0x28);
139
140 pub const RINGBUFFER_META: Self = Self(0x29);
141
142 pub const REAP_QUEUE: Self = Self(0x2A);
143
144 pub const SEAL_ANCHOR: Self = Self(0x2B);
145
146 pub const CUSTOM: Self = Self(0x40);
147
148 pub fn name(&self) -> &'static str {
149 match *self {
150 Self::ROW_NUMBER_MAPPING => "ROW_NUMBER_MAPPING",
151 Self::GROUP_DICTIONARY => "GROUP_DICTIONARY",
152 Self::NODE_COUNTER => "NODE_COUNTER",
153 Self::GROUP_RECORD => "GROUP_RECORD",
154 Self::SOURCE_WATERMARK => "SOURCE_WATERMARK",
155 Self::TIMER_WHEEL => "TIMER_WHEEL",
156 Self::TIMER_INDEX => "TIMER_INDEX",
157 Self::ACCUMULATOR => "ACCUMULATOR",
158 Self::BUFFER => "BUFFER",
159 Self::RUNNING => "RUNNING",
160 Self::EMIT => "EMIT",
161 Self::EXPIRY => "EXPIRY",
162 Self::COUNT => "COUNT",
163 Self::ROW_INDEX => "ROW_INDEX",
164 Self::SESSION => "SESSION",
165 Self::ROLLING_META => "ROLLING_META",
166 Self::ENGINE_META => "ENGINE_META",
167 Self::DISTINCT_ENTRY => "DISTINCT_ENTRY",
168 Self::WINDOW_META => "WINDOW_META",
169 Self::JOIN_LEFT => "JOIN_LEFT",
170 Self::JOIN_RIGHT => "JOIN_RIGHT",
171 Self::JOIN_SCHEMA => "JOIN_SCHEMA",
172 Self::RINGBUFFER_FORWARD => "RINGBUFFER_FORWARD",
173 Self::RINGBUFFER_ENTRY => "RINGBUFFER_ENTRY",
174 Self::GATE_VISIBILITY => "GATE_VISIBILITY",
175 Self::DISTINCT_LAYOUT => "DISTINCT_LAYOUT",
176 Self::RINGBUFFER_EXPIRY => "RINGBUFFER_EXPIRY",
177 Self::RINGBUFFER_TTL_ARM => "RINGBUFFER_TTL_ARM",
178 Self::SEAL_LEDGER => "SEAL_LEDGER",
179 Self::JOIN_PUBLISHED => "JOIN_PUBLISHED",
180 Self::JOIN_PIN => "JOIN_PIN",
181 Self::RINGBUFFER_META => "RINGBUFFER_META",
182 Self::REAP_QUEUE => "REAP_QUEUE",
183 Self::SEAL_ANCHOR => "SEAL_ANCHOR",
184 Self::CUSTOM => "CUSTOM",
185 _ => "CUSTOM",
186 }
187 }
188
189 pub fn is_data(&self) -> bool {
190 self.0 <= Self::HIGHEST_DATA
191 }
192
193 pub fn is_identity(&self) -> bool {
194 !self.is_data()
195 }
196
197 pub fn is_known(&self) -> bool {
198 self.is_data()
199 || matches!(
200 *self,
201 Self::ROW_NUMBER_MAPPING
202 | Self::GROUP_DICTIONARY | Self::NODE_COUNTER
203 | Self::GROUP_RECORD | Self::SOURCE_WATERMARK
204 | Self::TIMER_WHEEL | Self::TIMER_INDEX
205 )
206 }
207}
208
209pub fn is_framed_inner(inner: &[u8]) -> bool {
210 inner.is_empty() || OperatorStateKey::decode_inner(inner).is_some_and(|(_, keyspace, _)| keyspace.is_known())
211}
212
213#[derive(Debug, Clone, PartialEq, Eq)]
214pub struct OperatorStateKey {
215 pub operator: OperatorId,
216 pub group: GroupId,
217 pub keyspace: Keyspace,
218 pub suffix: Vec<u8>,
219}
220
221impl OperatorStateKey {
222 pub fn new(operator: OperatorId, group: GroupId, keyspace: Keyspace, suffix: impl Into<Vec<u8>>) -> Self {
223 Self {
224 operator,
225 group,
226 keyspace,
227 suffix: suffix.into(),
228 }
229 }
230
231 pub fn root(operator: OperatorId, keyspace: Keyspace, suffix: impl Into<Vec<u8>>) -> Self {
232 Self::new(operator, GroupId::ROOT, keyspace, suffix)
233 }
234
235 pub fn encoded(
236 operator: OperatorId,
237 group: GroupId,
238 keyspace: Keyspace,
239 suffix: impl AsRef<[u8]>,
240 ) -> EncodedKey {
241 let suffix = suffix.as_ref();
242 let mut serializer = KeySerializer::with_capacity(20 + suffix.len());
243 serializer
244 .extend_u8(KeyKind::OperatorState as u8)
245 .extend_u64(operator.0)
246 .extend_u64(group.0)
247 .extend_u8(keyspace.0)
248 .extend_raw(suffix);
249 serializer.to_encoded_key()
250 }
251
252 pub fn inner(&self) -> EncodedKey {
253 let mut serializer = KeySerializer::with_capacity(12 + self.suffix.len());
254 serializer.extend_u64(self.group.0).extend_u8(self.keyspace.0).extend_raw(&self.suffix);
255 serializer.to_encoded_key()
256 }
257
258 pub const KEYSPACE_INNER_OFFSET: u32 = size_of::<u64>() as u32;
259
260 pub fn decode_keyspace(stored: u8) -> Keyspace {
261 Keyspace(KeyDeserializer::from_bytes(&[stored]).read_u8().expect("a single byte decodes as u8"))
262 }
263
264 pub fn inner_encoded(group: GroupId, keyspace: Keyspace, suffix: impl AsRef<[u8]>) -> GroupStateKey {
265 let suffix = suffix.as_ref();
266 let mut serializer = KeySerializer::with_capacity(12 + suffix.len());
267 serializer.extend_u64(group.0).extend_u8(keyspace.0).extend_raw(suffix);
268 GroupStateKey(serializer.to_encoded_key())
269 }
270
271 pub fn decode_inner(inner: &[u8]) -> Option<(GroupId, Keyspace, Vec<u8>)> {
272 let mut de = KeyDeserializer::from_bytes(inner);
273 let group = de.read_u64().ok()?;
274 let keyspace = de.read_u8().ok()?;
275 let suffix = de.read_raw(de.remaining()).ok()?.to_vec();
276 Some((GroupId(group), Keyspace(keyspace), suffix))
277 }
278
279 pub fn node_range(operator: OperatorId) -> EncodedKeyRange {
280 node_range(operator)
281 }
282
283 pub fn decode_operator(key: &EncodedKey) -> Option<(OperatorId, EncodedKey)> {
284 let mut de = KeyDeserializer::from_bytes(key.as_slice());
285 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
286 if kind != KeyKind::OperatorState {
287 return None;
288 }
289 let operator = de.read_u64().ok()?;
290 let inner = de.read_raw(de.remaining()).ok()?.to_vec();
291 Some((OperatorId(operator), EncodedKey::new(inner)))
292 }
293}
294
295impl EncodableKey for OperatorStateKey {
296 const KIND: KeyKind = KeyKind::OperatorState;
297
298 fn encode(&self) -> EncodedKey {
299 let mut serializer = KeySerializer::with_capacity(20 + self.suffix.len());
300 serializer
301 .extend_u8(KeyKind::OperatorState as u8)
302 .extend_u64(self.operator.0)
303 .extend_u64(self.group.0)
304 .extend_u8(self.keyspace.0)
305 .extend_raw(&self.suffix);
306 serializer.to_encoded_key()
307 }
308
309 fn decode(key: &EncodedKey) -> Option<Self> {
310 let mut de = KeyDeserializer::from_bytes(key.as_slice());
311
312 let kind: KeyKind = de.read_u8().ok()?.try_into().ok()?;
313 if kind != KeyKind::OperatorState {
314 return None;
315 }
316
317 let operator = de.read_u64().ok()?;
318 let group = de.read_u64().ok()?;
319 let keyspace = de.read_u8().ok()?;
320 let suffix = de.read_raw(de.remaining()).ok()?.to_vec();
321
322 Some(Self {
323 operator: OperatorId(operator),
324 group: GroupId(group),
325 keyspace: Keyspace(keyspace),
326 suffix,
327 })
328 }
329}
330
331#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
332pub struct GroupStateKey(EncodedKey);
333
334impl GroupStateKey {
335 pub fn new(group: GroupId, keyspace: Keyspace, suffix: impl AsRef<[u8]>) -> Self {
336 OperatorStateKey::inner_encoded(group, keyspace, suffix)
337 }
338
339 pub fn root(keyspace: Keyspace, suffix: impl AsRef<[u8]>) -> Self {
340 Self::new(GroupId::ROOT, keyspace, suffix)
341 }
342
343 pub fn from_framed(key: EncodedKey) -> Option<Self> {
344 is_framed_inner(key.as_slice()).then_some(Self(key))
345 }
346
347 pub fn bound_unchecked(key: EncodedKey) -> Self {
348 Self(key)
349 }
350
351 pub fn as_encoded(&self) -> &EncodedKey {
352 &self.0
353 }
354
355 pub fn into_encoded(self) -> EncodedKey {
356 self.0
357 }
358
359 pub fn as_slice(&self) -> &[u8] {
360 self.0.as_slice()
361 }
362
363 pub fn as_bytes(&self) -> &[u8] {
364 self.0.as_bytes()
365 }
366
367 pub fn group(&self) -> Option<GroupId> {
368 OperatorStateKey::decode_inner(self.0.as_slice()).map(|(group, _, _)| group)
369 }
370
371 pub fn keyspace(&self) -> Option<Keyspace> {
372 let bytes = self.0.as_slice();
373 let offset = OperatorStateKey::KEYSPACE_INNER_OFFSET as usize;
374 (bytes.len() > offset).then(|| Keyspace(encode_u8(bytes[offset])))
375 }
376}
377
378impl AsRef<[u8]> for GroupStateKey {
379 fn as_ref(&self) -> &[u8] {
380 self.0.as_slice()
381 }
382}
383
384impl AsRef<EncodedKey> for GroupStateKey {
385 fn as_ref(&self) -> &EncodedKey {
386 &self.0
387 }
388}
389
390pub trait IntoGroupStateKey {
391 fn into_group_state_key(self) -> GroupStateKey;
392}
393
394impl IntoGroupStateKey for GroupStateKey {
395 fn into_group_state_key(self) -> GroupStateKey {
396 self
397 }
398}
399
400fn group_inner_prefix(group: GroupId) -> Vec<u8> {
401 let mut serializer = KeySerializer::with_capacity(12);
402 serializer.extend_u64(group.0);
403 serializer.finish().as_ref().to_vec()
404}
405
406fn keyspace_inner_prefix(group: GroupId, keyspace: Keyspace) -> Vec<u8> {
407 let mut prefix = group_inner_prefix(group);
408 prefix.push(encode_u8(keyspace.0));
409 prefix
410}
411
412pub fn group_inner_range(group: GroupId) -> EncodedKeyRange {
413 EncodedKeyRange::prefix(&group_inner_prefix(group))
414}
415
416pub fn keyspace_inner_range(group: GroupId, keyspace: Keyspace) -> EncodedKeyRange {
417 EncodedKeyRange::prefix(&keyspace_inner_prefix(group, keyspace))
418}
419
420pub fn keyspace_inner_range_upto(group: GroupId, keyspace: Keyspace, suffix: &[u8]) -> EncodedKeyRange {
421 let mut bound = keyspace_inner_prefix(group, keyspace);
422 bound.extend_from_slice(suffix);
423 EncodedKeyRange::new(keyspace_inner_range(group, keyspace).start, EncodedKeyRange::prefix(&bound).end)
424}
425
426pub fn group_data_inner_range(group: GroupId) -> EncodedKeyRange {
427 let prefix = group_inner_prefix(group);
428 let mut start = prefix.clone();
429 start.push(encode_u8(Keyspace::HIGHEST_DATA));
430 EncodedKeyRange::new(Bound::Included(EncodedKey::new(start)), EncodedKeyRange::prefix(&prefix).end)
431}
432
433pub fn group_identity_inner_range(group: GroupId) -> EncodedKeyRange {
434 let prefix = group_inner_prefix(group);
435 let mut end = prefix.clone();
436 end.push(encode_u8(Keyspace::HIGHEST_DATA));
437 EncodedKeyRange::new(Bound::Included(EncodedKey::new(prefix)), Bound::Excluded(EncodedKey::new(end)))
438}
439
440pub fn node_prefix(operator: OperatorId) -> Vec<u8> {
441 let mut serializer = KeySerializer::with_capacity(12);
442 serializer.extend_u8(KeyKind::OperatorState as u8).extend_u64(operator.0);
443 serializer.finish().as_ref().to_vec()
444}
445
446fn group_prefix(operator: OperatorId, group: GroupId) -> Vec<u8> {
447 let mut serializer = KeySerializer::with_capacity(20);
448 serializer.extend_u8(KeyKind::OperatorState as u8).extend_u64(operator.0).extend_u64(group.0);
449 serializer.finish().as_ref().to_vec()
450}
451
452fn keyspace_prefix(operator: OperatorId, group: GroupId, keyspace: Keyspace) -> Vec<u8> {
453 let mut prefix = group_prefix(operator, group);
454 prefix.push(encode_u8(keyspace.0));
455 prefix
456}
457
458pub fn node_range(operator: OperatorId) -> EncodedKeyRange {
459 EncodedKeyRange::prefix(&node_prefix(operator))
460}
461
462pub fn group_range(operator: OperatorId, group: GroupId) -> EncodedKeyRange {
463 EncodedKeyRange::prefix(&group_prefix(operator, group))
464}
465
466pub fn keyspace_range(operator: OperatorId, group: GroupId, keyspace: Keyspace) -> EncodedKeyRange {
467 EncodedKeyRange::prefix(&keyspace_prefix(operator, group, keyspace))
468}
469
470pub fn group_data_range(operator: OperatorId, group: GroupId) -> EncodedKeyRange {
471 let prefix = group_prefix(operator, group);
472 let mut start = prefix.clone();
473 start.push(encode_u8(Keyspace::HIGHEST_DATA));
474 EncodedKeyRange::new(Bound::Included(EncodedKey::new(start)), EncodedKeyRange::prefix(&prefix).end)
475}
476
477pub fn group_identity_range(operator: OperatorId, group: GroupId) -> EncodedKeyRange {
478 let prefix = group_prefix(operator, group);
479 let mut end = prefix.clone();
480 end.push(encode_u8(Keyspace::HIGHEST_DATA));
481 EncodedKeyRange::new(Bound::Included(EncodedKey::new(prefix)), Bound::Excluded(EncodedKey::new(end)))
482}
483
484#[cfg(test)]
485mod tests {
486 use std::{ops::Bound, slice};
487
488 use super::{
489 EncodedKey, EncodedKeyRange, GroupId, GroupSet, KeySerializer, Keyspace, OperatorStateKey,
490 group_data_inner_range, group_data_of_inner, group_data_range, group_identity_inner_range,
491 group_identity_range, group_inner_prefix, group_inner_range, group_range, is_framed_inner,
492 keyspace_range, node_prefix, node_range,
493 };
494 use crate::{interface::catalog::flow::OperatorId, key::EncodableKey};
495
496 const NODES: [u64; 4] = [1, 17, 300, 70_000];
497 const GROUPS: [u64; 8] = [1, 2, 127, 128, 1000, 100_000, 1 << 30, u64::MAX];
498 const DATA_KEYSPACES: [Keyspace; 4] =
499 [Keyspace::ACCUMULATOR, Keyspace::BUFFER, Keyspace::RUNNING, Keyspace::CUSTOM];
500 const IDENTITY_KEYSPACES: [Keyspace; 2] = [Keyspace::GROUP_RECORD, Keyspace::ROW_NUMBER_MAPPING];
501
502 #[derive(Clone, Copy, PartialEq, Debug)]
503 enum Phase {
504 Data,
505 Identity,
506 }
507
508 const CENSUS: [(&str, Keyspace, Phase); 35] = [
511 ("ROW_NUMBER_MAPPING", Keyspace::ROW_NUMBER_MAPPING, Phase::Identity),
512 ("GROUP_DICTIONARY", Keyspace::GROUP_DICTIONARY, Phase::Identity),
513 ("NODE_COUNTER", Keyspace::NODE_COUNTER, Phase::Identity),
514 ("GROUP_RECORD", Keyspace::GROUP_RECORD, Phase::Identity),
515 ("SOURCE_WATERMARK", Keyspace::SOURCE_WATERMARK, Phase::Identity),
516 ("TIMER_WHEEL", Keyspace::TIMER_WHEEL, Phase::Identity),
517 ("TIMER_INDEX", Keyspace::TIMER_INDEX, Phase::Identity),
518 ("ACCUMULATOR", Keyspace::ACCUMULATOR, Phase::Data),
519 ("BUFFER", Keyspace::BUFFER, Phase::Data),
520 ("RUNNING", Keyspace::RUNNING, Phase::Data),
521 ("EMIT", Keyspace::EMIT, Phase::Data),
522 ("EXPIRY", Keyspace::EXPIRY, Phase::Data),
523 ("COUNT", Keyspace::COUNT, Phase::Data),
524 ("ROW_INDEX", Keyspace::ROW_INDEX, Phase::Data),
525 ("SESSION", Keyspace::SESSION, Phase::Data),
526 ("ROLLING_META", Keyspace::ROLLING_META, Phase::Data),
527 ("ENGINE_META", Keyspace::ENGINE_META, Phase::Data),
528 ("DISTINCT_ENTRY", Keyspace::DISTINCT_ENTRY, Phase::Data),
529 ("WINDOW_META", Keyspace::WINDOW_META, Phase::Data),
530 ("JOIN_LEFT", Keyspace::JOIN_LEFT, Phase::Data),
531 ("JOIN_RIGHT", Keyspace::JOIN_RIGHT, Phase::Data),
532 ("JOIN_SCHEMA", Keyspace::JOIN_SCHEMA, Phase::Data),
533 ("RINGBUFFER_FORWARD", Keyspace::RINGBUFFER_FORWARD, Phase::Data),
534 ("RINGBUFFER_ENTRY", Keyspace::RINGBUFFER_ENTRY, Phase::Data),
535 ("GATE_VISIBILITY", Keyspace::GATE_VISIBILITY, Phase::Data),
536 ("DISTINCT_LAYOUT", Keyspace::DISTINCT_LAYOUT, Phase::Data),
537 ("RINGBUFFER_EXPIRY", Keyspace::RINGBUFFER_EXPIRY, Phase::Data),
538 ("RINGBUFFER_TTL_ARM", Keyspace::RINGBUFFER_TTL_ARM, Phase::Data),
539 ("SEAL_LEDGER", Keyspace::SEAL_LEDGER, Phase::Data),
540 ("JOIN_PUBLISHED", Keyspace::JOIN_PUBLISHED, Phase::Data),
541 ("JOIN_PIN", Keyspace::JOIN_PIN, Phase::Data),
542 ("RINGBUFFER_META", Keyspace::RINGBUFFER_META, Phase::Data),
543 ("REAP_QUEUE", Keyspace::REAP_QUEUE, Phase::Data),
544 ("SEAL_ANCHOR", Keyspace::SEAL_ANCHOR, Phase::Data),
545 ("CUSTOM", Keyspace::CUSTOM, Phase::Data),
546 ];
547
548 fn declared_keyspaces() -> usize {
551 let source = include_str!("operator_state.rs");
552 let body = source
553 .split("impl Keyspace {")
554 .nth(1)
555 .expect("the Keyspace impl block is where the constants are declared");
556 let body = body.split("\n}\n").next().expect("the impl block is closed");
557 body.lines()
558 .filter(|line| {
559 let line = line.trim_start();
560 line.starts_with("pub const") && line.contains("Self(")
561 })
562 .count()
563 }
564
565 #[test]
566 fn a_bare_row_number_key_is_indistinguishable_from_another_groups_prefix() {
567 let mut bare = KeySerializer::with_capacity(4);
569 bare.extend_u64(7u64);
570 let bare = bare.finish().as_ref().to_vec();
571
572 assert_eq!(bare, group_inner_prefix(GroupId(7)), "a bare row number encodes as a group prefix");
573 assert!(
574 contains(&group_identity_inner_range(GroupId(7)), &bare),
575 "so reclaiming group 7 erases it with the row-number mappings"
576 );
577 assert!(!is_framed_inner(&bare));
578
579 let framed = OperatorStateKey::inner_encoded(GroupId::ROOT, Keyspace::CUSTOM, 7u64.to_be_bytes());
580 assert!(is_framed_inner(framed.as_slice()));
581 assert!(
582 !contains(&group_identity_inner_range(GroupId(7)), framed.as_slice()),
583 "the framed form must sit outside every other group's range"
584 );
585 }
586
587 #[test]
588 fn the_empty_key_is_framing_because_it_sorts_below_every_group() {
589 let empty: &[u8] = &[];
591 assert!(is_framed_inner(empty));
592
593 for group in GROUPS {
594 let range = group_inner_range(GroupId(group));
595 assert!(
596 !contains(&range, empty),
597 "the empty key must sit outside group {group}'s range, not merely be unattributed"
598 );
599 }
600 }
601
602 #[test]
603 fn a_keyspace_this_substrate_never_defines_is_not_framing() {
604 let mut stray = KeySerializer::with_capacity(4);
606 stray.extend_u64(3u64).extend_u8(0x90u8);
607 assert!(!is_framed_inner(stray.finish().as_ref()));
608
609 for keyspace in DATA_KEYSPACES.iter().chain(IDENTITY_KEYSPACES.iter()) {
610 assert!(
611 is_framed_inner(OperatorStateKey::inner_encoded(GroupId(3), *keyspace, []).as_slice()),
612 "keyspace {keyspace:?} is one the substrate writes and must pass"
613 );
614 }
615 }
616
617 fn contains(range: &EncodedKeyRange, key: &[u8]) -> bool {
618 let after_start = match &range.start {
619 Bound::Included(start) => key >= start.as_slice(),
620 Bound::Excluded(start) => key > start.as_slice(),
621 Bound::Unbounded => true,
622 };
623 let before_end = match &range.end {
624 Bound::Included(end) => key <= end.as_slice(),
625 Bound::Excluded(end) => key < end.as_slice(),
626 Bound::Unbounded => true,
627 };
628 after_start && before_end
629 }
630
631 fn population() -> Vec<OperatorStateKey> {
632 let mut keys = Vec::new();
633 for operator in NODES {
634 for group in GROUPS {
635 for keyspace in DATA_KEYSPACES.iter().chain(IDENTITY_KEYSPACES.iter()) {
636 for coord in [0u64, 1, 999, u64::MAX] {
637 keys.push(OperatorStateKey::new(
638 OperatorId(operator),
639 GroupId(group),
640 *keyspace,
641 coord.to_be_bytes().to_vec(),
642 ));
643 }
644 }
645 }
646 keys.push(OperatorStateKey::root(
647 OperatorId(operator),
648 Keyspace::GROUP_DICTIONARY,
649 b"7xKXtg2CW87d97TXJSDpbD5jBkheTqA83TZRuJosgAsU".to_vec(),
650 ));
651 }
652 keys
653 }
654
655 #[test]
656 fn a_group_range_contains_exactly_that_groups_keys() {
657 let population = population();
659 for operator in NODES {
660 for group in GROUPS {
661 let range = group_range(OperatorId(operator), GroupId(group));
662 for key in &population {
663 let encoded = key.encode();
664 let expected = key.operator.0 == operator && key.group.0 == group;
665 assert_eq!(
666 contains(&range, encoded.as_slice()),
667 expected,
668 "operator {operator} group {group} range disagreed about a key of operator {} \
669 group {}",
670 key.operator.0,
671 key.group.0
672 );
673 }
674 }
675 }
676 }
677
678 #[test]
679 fn variable_length_group_ids_cannot_prefix_one_another() {
680 let encodings: Vec<Vec<u8>> = GROUPS
683 .iter()
684 .map(|group| {
685 OperatorStateKey::new(OperatorId(1), GroupId(*group), Keyspace::ACCUMULATOR, vec![])
686 .encode()
687 .as_slice()
688 .to_vec()
689 })
690 .collect();
691
692 for (i, a) in encodings.iter().enumerate() {
693 for (j, b) in encodings.iter().enumerate() {
694 if i != j {
695 assert!(
696 !b.starts_with(a.as_slice()),
697 "group {} encodes as a prefix of group {}",
698 GROUPS[i],
699 GROUPS[j]
700 );
701 }
702 }
703 }
704 }
705
706 #[test]
707 fn the_data_and_identity_ranges_partition_the_group() {
708 for operator in NODES {
710 for group in GROUPS {
711 let data = group_data_range(OperatorId(operator), GroupId(group));
712 let identity = group_identity_range(OperatorId(operator), GroupId(group));
713
714 for keyspace in DATA_KEYSPACES {
715 let key = OperatorStateKey::new(
716 OperatorId(operator),
717 GroupId(group),
718 keyspace,
719 vec![7, 7],
720 )
721 .encode();
722 assert!(
723 contains(&data, key.as_slice()),
724 "data keyspace {keyspace:?} must fall in the phase-1 range"
725 );
726 assert!(
727 !contains(&identity, key.as_slice()),
728 "data keyspace {keyspace:?} must not fall in the phase-2 range"
729 );
730 }
731
732 for keyspace in IDENTITY_KEYSPACES {
733 let key = OperatorStateKey::new(
734 OperatorId(operator),
735 GroupId(group),
736 keyspace,
737 vec![7, 7],
738 )
739 .encode();
740 assert!(
741 contains(&identity, key.as_slice()),
742 "identity keyspace {keyspace:?} must fall in the phase-2 range"
743 );
744 assert!(
745 !contains(&data, key.as_slice()),
746 "identity keyspace {keyspace:?} must survive phase 1"
747 );
748 }
749 }
750 }
751 }
752
753 #[test]
754 fn every_declared_keyspace_names_itself_for_offline_attribution() {
755 for (name, keyspace, _) in CENSUS {
757 assert_eq!(
758 keyspace.name(),
759 name,
760 "{name} ({:#04x}) does not name itself, so an offline census reports it as CUSTOM",
761 keyspace.0
762 );
763 }
764
765 assert_eq!(
766 Keyspace(0x41).name(),
767 "CUSTOM",
768 "a byte no constant claims must fall through rather than borrow a neighbour's name"
769 );
770 }
771
772 #[test]
773 fn every_declared_keyspace_is_distinct_framing_and_swept_by_exactly_one_phase() {
774 assert_eq!(
776 CENSUS.len(),
777 declared_keyspaces(),
778 "a keyspace was added to Keyspace without being added to the census, so nothing below \
779 ever looks at its byte"
780 );
781
782 let mut seen: Vec<(&str, u8)> = Vec::new();
783 for (name, keyspace, phase) in CENSUS {
784 if let Some((other, _)) = seen.iter().find(|(_, byte)| *byte == keyspace.0) {
785 panic!("{name} and {other} both claim keyspace byte {:#04x}", keyspace.0);
786 }
787 seen.push((name, keyspace.0));
788
789 assert!(
790 keyspace.is_known(),
791 "{name} is declared but not framing, so the sweep panics on the first row it holds"
792 );
793
794 let key = OperatorStateKey::new(OperatorId(9), GroupId(4), keyspace, vec![7, 7]).encode();
795 let data = contains(&group_data_range(OperatorId(9), GroupId(4)), key.as_slice());
796 let identity = contains(&group_identity_range(OperatorId(9), GroupId(4)), key.as_slice());
797
798 assert!(data != identity, "{name} must fall in exactly one phase, not {data} and {identity}");
799 assert_eq!(
800 data,
801 phase == Phase::Data,
802 "{name} is declared {phase:?} but the phase-1 range says data={data}"
803 );
804 assert_eq!(
805 keyspace.is_data(),
806 phase == Phase::Data,
807 "{name} is declared {phase:?} but is_data says {}",
808 keyspace.is_data()
809 );
810 }
811 }
812
813 #[test]
814 fn root_entries_sit_outside_every_group_range() {
815 for operator in NODES {
817 let dictionary = OperatorStateKey::root(
818 OperatorId(operator),
819 Keyspace::GROUP_DICTIONARY,
820 b"mint-pubkey".to_vec(),
821 )
822 .encode();
823 for group in GROUPS {
824 let range = group_range(OperatorId(operator), GroupId(group));
825 assert!(
826 !contains(&range, dictionary.as_slice()),
827 "group {group} range must not contain the root group's dictionary entry"
828 );
829 }
830 }
831 }
832
833 #[test]
834 fn a_node_range_contains_exactly_that_nodes_keys() {
835 let population = population();
837 for operator in NODES {
838 let range = node_range(OperatorId(operator));
839 for key in &population {
840 let encoded = key.encode();
841 assert_eq!(
842 contains(&range, encoded.as_slice()),
843 key.operator.0 == operator,
844 "operator {operator} range disagreed about a key of operator {}",
845 key.operator.0
846 );
847 }
848 }
849 }
850
851 #[test]
852 fn a_keyspace_range_isolates_one_keyspace_of_one_group() {
853 let operator = OperatorId(17);
855 let group = GroupId(42);
856 let range = keyspace_range(operator, group, Keyspace::BUFFER);
857
858 let inside = OperatorStateKey::new(operator, group, Keyspace::BUFFER, vec![1]).encode();
859 assert!(contains(&range, inside.as_slice()));
860
861 for other in [Keyspace::ACCUMULATOR, Keyspace::RUNNING, Keyspace::GROUP_RECORD] {
862 let key = OperatorStateKey::new(operator, group, other, vec![1]).encode();
863 assert!(!contains(&range, key.as_slice()), "keyspace {other:?} leaked into the buffer range");
864 }
865
866 let other_group = OperatorStateKey::new(operator, GroupId(43), Keyspace::BUFFER, vec![1]).encode();
867 assert!(!contains(&range, other_group.as_slice()), "another group's buffer leaked into the range");
868 }
869
870 #[test]
871 fn encode_decode_round_trips_every_component() {
872 let key = OperatorStateKey::new(
873 OperatorId(0xDEAD_BEEF),
874 GroupId(123_456),
875 Keyspace::CUSTOM,
876 vec![1, 2, 3, 4],
877 );
878 assert_eq!(OperatorStateKey::decode(&key.encode()), Some(key));
879 }
880
881 #[test]
882 fn keys_still_decode_as_operator_state_of_their_node() {
883 let key = OperatorStateKey::new(OperatorId(9), GroupId(4), Keyspace::ACCUMULATOR, vec![1]).encode();
885
886 let decoded = OperatorStateKey::decode(&key).expect("must remain decodable as its key kind");
887 assert_eq!(decoded.operator, OperatorId(9));
888 }
889
890 #[test]
891 fn an_inner_key_composed_with_its_node_prefix_reproduces_the_full_key() {
892 let key = OperatorStateKey::new(OperatorId(17), GroupId(42), Keyspace::BUFFER, vec![9, 9]);
895
896 let mut composed = node_prefix(OperatorId(17));
897 composed.extend_from_slice(key.inner().as_slice());
898
899 assert_eq!(composed, key.encode().as_slice(), "inner key plus operator prefix must equal the full key");
900 }
901
902 #[test]
903 fn the_root_group_range_stays_inside_its_node() {
904 let range = group_inner_range(GroupId::ROOT).with_prefix(EncodedKey::new(node_prefix(OperatorId(17))));
906
907 let own = OperatorStateKey::root(OperatorId(17), Keyspace::GROUP_DICTIONARY, vec![1]).encode();
908 assert!(contains(&range, own.as_slice()), "the operator's own dictionary entry must be in range");
909
910 for operator in NODES {
911 if operator == 17 {
912 continue;
913 }
914 for keyspace in [Keyspace::GROUP_DICTIONARY, Keyspace::ACCUMULATOR] {
915 let foreign =
916 OperatorStateKey::new(OperatorId(operator), GroupId::ROOT, keyspace, vec![1])
917 .encode();
918 assert!(
919 !contains(&range, foreign.as_slice()),
920 "operator {operator} leaked into operator 17's root-group range"
921 );
922 }
923 }
924 }
925
926 #[test]
927 fn inner_ranges_partition_the_group_like_their_full_key_counterparts() {
928 let operator = OperatorId(17);
931 let prefix = EncodedKey::new(node_prefix(operator));
932 for group in GROUPS {
933 let data = group_data_inner_range(GroupId(group)).with_prefix(prefix.clone());
934 let identity = group_identity_inner_range(GroupId(group)).with_prefix(prefix.clone());
935
936 for keyspace in DATA_KEYSPACES {
937 let key = OperatorStateKey::new(operator, GroupId(group), keyspace, vec![7]).encode();
938 assert!(contains(&data, key.as_slice()));
939 assert!(!contains(&identity, key.as_slice()));
940 }
941 for keyspace in IDENTITY_KEYSPACES {
942 let key = OperatorStateKey::new(operator, GroupId(group), keyspace, vec![7]).encode();
943 assert!(contains(&identity, key.as_slice()));
944 assert!(!contains(&data, key.as_slice()));
945 }
946 }
947 }
948
949 #[test]
950 fn decode_inner_round_trips_the_tail() {
951 let key = OperatorStateKey::new(OperatorId(3), GroupId(77), Keyspace::EMIT, vec![4, 5, 6]);
952 let (group, keyspace, suffix) =
953 OperatorStateKey::decode_inner(key.inner().as_slice()).expect("inner must decode");
954
955 assert_eq!(group, GroupId(77));
956 assert_eq!(keyspace, Keyspace::EMIT);
957 assert_eq!(suffix, vec![4, 5, 6]);
958 }
959
960 #[test]
961 fn interned_group_keys_stay_compact() {
962 let interned =
964 OperatorStateKey::new(OperatorId(17), GroupId(123_456), Keyspace::ACCUMULATOR, vec![0; 8])
965 .encode();
966 let raw_group_bytes = b"7xKXtg2CW87d97TXJSDpbD5jBkheTqA83TZRuJosgAsU\
967 So11111111111111111111111111111111111111112"
968 .len() + 8;
969
970 assert!(
971 interned.as_slice().len() * 3 < raw_group_bytes,
972 "an interned state key ({} bytes) must stay far below a raw-group key ({} bytes)",
973 interned.as_slice().len(),
974 raw_group_bytes
975 );
976 }
977
978 #[test]
979 fn the_ram_predicate_and_the_disk_range_agree_on_every_key() {
980 for group in GROUPS.map(GroupId) {
983 let range = group_data_inner_range(group);
984 for other in GROUPS.map(GroupId) {
985 for keyspace in DATA_KEYSPACES.iter().chain(IDENTITY_KEYSPACES.iter()) {
986 let key = OperatorStateKey::inner_encoded(other, *keyspace, vec![7, 7]);
987 let in_range = contains(&range, key.as_slice());
988 let in_predicate = group_data_of_inner(key.as_slice()) == Some(group);
989 assert_eq!(
990 in_range, in_predicate,
991 "disk range and RAM predicate disagree for group {group:?} on a \
992 {keyspace:?} key of group {other:?}"
993 );
994 }
995 }
996 }
997 }
998
999 #[test]
1000 fn the_ram_predicate_refuses_identity_keyspaces() {
1001 for keyspace in IDENTITY_KEYSPACES {
1003 let key = OperatorStateKey::inner_encoded(GroupId(9), keyspace, vec![1]);
1004 assert_eq!(
1005 group_data_of_inner(key.as_slice()),
1006 None,
1007 "{keyspace:?} must not be reported as reclaimable group data"
1008 );
1009 }
1010 }
1011
1012 #[test]
1013 fn a_key_too_short_to_carry_a_keyspace_is_refused() {
1014 assert_eq!(group_data_of_inner(&[]), None);
1016 assert_eq!(group_data_of_inner(&[0xAB]), None, "a group with no keyspace byte must not decode");
1017 }
1018
1019 #[test]
1020 fn the_predicate_agrees_with_the_disk_range_on_arbitrary_bytes() {
1021 let mut seed = 0x2545F4914F6CDD1Du64;
1023 let mut next = move || {
1024 seed ^= seed << 13;
1025 seed ^= seed >> 7;
1026 seed ^= seed << 17;
1027 seed
1028 };
1029
1030 for _ in 0..2000 {
1031 let len = (next() % 12) as usize;
1032 let key: Vec<u8> = (0..len).map(|_| (next() % 256) as u8).collect();
1033 let Some(group) = group_data_of_inner(&key) else {
1034 continue;
1035 };
1036 assert!(
1037 contains(&group_data_inner_range(group), &key),
1038 "predicate attributed {key:?} to {group:?} but the disk range excludes it"
1039 );
1040 }
1041 }
1042
1043 #[test]
1044 fn a_group_set_is_sorted_deduped_and_never_admits_root() {
1045 let set = GroupSet::new([GroupId(9), GroupId(2), GroupId(9), GroupId::ROOT, GroupId(5)]);
1047
1048 assert_eq!(set.as_slice(), &[GroupId(2), GroupId(5), GroupId(9)]);
1049 assert_eq!(set.len(), 3);
1050 assert!(set.contains(GroupId(5)));
1051 assert!(!set.contains(GroupId(3)));
1052 assert!(!set.contains(GroupId::ROOT), "the root group must be filtered out, not merely unsorted");
1053 }
1054
1055 #[test]
1056 fn an_empty_group_set_matches_nothing() {
1057 let set = GroupSet::new([]);
1058
1059 assert!(set.is_empty());
1060 assert!(!set.contains(GroupId::FIRST));
1061 }
1062
1063 #[test]
1064 fn a_group_set_hands_the_ffi_boundary_a_plain_u64_array() {
1065 let set = GroupSet::new([GroupId(3), GroupId(1), GroupId(2)]);
1067 let (ptr, len) = set.as_raw_parts();
1068
1069 assert_eq!(len, 3);
1070 let raw = unsafe { slice::from_raw_parts(ptr, len) };
1072 assert_eq!(raw, &[1u64, 2, 3]);
1073 }
1074}