Skip to main content

reifydb_core/key/
operator_state.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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	/// Every keyspace the substrate declares, with the phase allowed to erase it. The phase is written
509	/// down rather than read back from `is_data`, or a keyspace changing sides would pass unremarked.
510	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	/// Counts `Keyspace` constants from the source text. There is no reflection over associated
549	/// constants, so this is the only way the census can notice a keyspace nobody listed.
550	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		// a bare row number equals a group prefix, so it is erased with that group on reclaim, never errors
568		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		// an empty inner key must sort below every group's prefix, or a reclaim phase could reach it
590		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		// a two-byte group+keyspace pair must not be framing unless the keyspace is one the substrate declares
605		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		// a group range must contain exactly that operator+group's keys, or reclaim destroys or leaks state
658		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		// no group id's varint encoding may prefix another's, or reclaiming the shorter erases the longer's
681		// keys
682		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		// data and identity keyspaces must fall in exactly one reclamation phase, never both or neither
709		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		// every declared keyspace must name itself, or an offline census misattributes it as CUSTOM
756		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		// every declared keyspace must have a unique byte and belong to exactly one reclamation phase
775		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		// the root-scoped dictionary entry must sit outside every group range, or reclaiming a group erases it
816		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		// a node range must contain exactly its own operator's keys, since drop_operator deletes by range
836		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		// a keyspace range must isolate exactly one keyspace of one group, or scans mix incompatible payloads
854		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		// every key of this kind must still decode to its own operator, or state is misrouted into the CDC log
884		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		// inner key plus node prefix must reproduce the full key, or state written through the API is
893		// unreachable
894		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		// group 0's inner range has no byte-wise successor, so it must stay bounded by the operator prefix
905		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		// inner data/identity ranges must partition a group like their full-key counterparts, since reclamation
929		// uses them
930		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		// interning must keep state keys far smaller than embedding raw group bytes in every key
963		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		// the RAM predicate and the disk range must agree on every key, or a phase-1 delete ghosts or strands a
981		// row
982		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		// the RAM predicate must never report an identity keyspace as reclaimable group data
1002		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		// a key without both a group and a keyspace byte must not decode as group data
1015		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		// the predicate and the disk range must agree even on arbitrary bytes no encoder produced
1022		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		// a group set must stay sorted and deduped for binary_search, and must never admit the root group
1046		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		// GroupId must stay repr(transparent) over u64, or the FFI slice cast reads the wrong bytes
1066		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		// SAFETY: the slice is alive for the whole assertion and GroupId is repr(transparent) over u64.
1071		let raw = unsafe { slice::from_raw_parts(ptr, len) };
1072		assert_eq!(raw, &[1u64, 2, 3]);
1073	}
1074}