use std::{
borrow::Cow,
fmt::{Display, Formatter, Result as FmtResult},
ops::Bound,
};
use reifydb_codec::key::{
deserializer::KeyDeserializer,
encode_u8,
encoded::{EncodedKey, EncodedKeyRange},
serializer::KeySerializer,
};
use reifydb_value::util::hash::{Hash128, xxh3_128};
use serde::{Deserialize, Serialize};
use smallvec::{SmallVec, smallvec};
use super::super::KeyTag;
use crate::{
interface::catalog::flow::OperatorId,
key::{
any::{ByteEncoding, Field, KeyFields, RawEncoding, Width},
bound::{TaggedKeyBound, TaggedKeyBoundRange},
operator::{
keyspace::{
KeyspaceVisitor, REGISTERED, dispatch,
root::{CustomNotCachedSuffix, NodeCounter, NodeCounterKey, NodeCounterKind},
suffix_width_of,
},
traits::Keyspace,
},
typed::{BoundedKey, layout::KeyLayout},
},
metrics::heap::HeapSize,
state::typed::{SuffixBytes, typed_key},
};
#[repr(transparent)]
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
pub struct GroupId([u8; GroupId::WIDTH]);
impl GroupId {
pub const WIDTH: usize = 24;
const HASH_OFFSET: usize = size_of::<u64>();
const NOT_A_WINDOW: u64 = u64::MAX;
pub const ROOT: Self = Self([0; Self::WIDTH]);
pub const FIRST_NON_ROOT: Self = {
let mut bytes = [0u8; Self::WIDTH];
bytes[Self::WIDTH - 1] = 1;
Self(bytes)
};
pub const MIN: Self = Self([u8::MIN; Self::WIDTH]);
pub const MAX: Self = Self([u8::MAX; Self::WIDTH]);
pub const fn from_bytes(bytes: [u8; Self::WIDTH]) -> Self {
Self(bytes)
}
pub const fn as_bytes(&self) -> &[u8; Self::WIDTH] {
&self.0
}
pub fn of(key: &EncodedKey) -> Self {
Self::hashed(xxh3_128(key.as_slice()))
}
pub fn hashed(hash: Hash128) -> Self {
match hash.0 {
0 => Self::FIRST_NON_ROOT,
carried => Self::at(Self::NOT_A_WINDOW, carried),
}
}
pub fn window(partition: Hash128, window_id: u64) -> Self {
assert!(
window_id != Self::NOT_A_WINDOW,
"window id {window_id} is the sentinel that marks a group as having no window; a window \
minted at it would collide with the join and distinct groups of the same operator"
);
Self::at(window_id, partition.0)
}
pub fn window_span(window_id: u64) -> (Self, Self) {
(Self::at(window_id, u128::MIN), Self::at(window_id, u128::MAX))
}
pub fn window_id(&self) -> Option<u64> {
let mut leading = [0u8; size_of::<u64>()];
leading.copy_from_slice(&self.0[..Self::HASH_OFFSET]);
match !u64::from_be_bytes(leading) {
Self::NOT_A_WINDOW => None,
window_id => Some(window_id),
}
}
fn at(window_id: u64, hash: u128) -> Self {
let mut bytes = [0u8; Self::WIDTH];
bytes[..Self::HASH_OFFSET].copy_from_slice(&(!window_id).to_be_bytes());
bytes[Self::HASH_OFFSET..].copy_from_slice(&hash.to_be_bytes());
Self(bytes)
}
pub fn is_root(&self) -> bool {
*self == Self::ROOT
}
pub fn successor(&self) -> Option<Self> {
let mut bytes = self.0;
for byte in bytes.iter_mut().rev() {
let (stepped, carried) = byte.overflowing_add(1);
*byte = stepped;
if !carried {
return Some(Self(bytes));
}
}
None
}
pub fn predecessor(&self) -> Option<Self> {
let mut bytes = self.0;
for byte in bytes.iter_mut().rev() {
let (stepped, borrowed) = byte.overflowing_sub(1);
*byte = stepped;
if !borrowed {
return Some(Self(bytes));
}
}
None
}
}
impl Display for GroupId {
fn fmt(&self, f: &mut Formatter<'_>) -> FmtResult {
for byte in self.0 {
write!(f, "{byte:02x}")?;
}
Ok(())
}
}
impl HeapSize for GroupId {
fn heap_size(&self) -> usize {
0
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct GroupSet(Vec<GroupId>);
impl GroupSet {
pub fn new(groups: impl IntoIterator<Item = GroupId>) -> Self {
let mut groups: Vec<GroupId> = groups.into_iter().filter(|g| !g.is_root()).collect();
groups.sort_unstable();
groups.dedup();
Self(groups)
}
pub fn contains(&self, group: GroupId) -> bool {
self.0.binary_search(&group).is_ok()
}
pub fn as_slice(&self) -> &[GroupId] {
&self.0
}
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
}
pub fn group_data_of_inner(inner: &[u8]) -> Option<GroupId> {
let mut de = KeyDeserializer::from_bytes(inner);
let group = GroupId::from_bytes(de.read_fixed().ok()?);
let keyspace = KeyspaceId(de.read_u8().ok()?);
if !keyspace.is_data() {
return None;
}
inner.starts_with(&group_inner_prefix(group)).then_some(group)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct KeyspaceId(pub u8);
impl KeyspaceId {
pub const HIGHEST_DATA: u8 = 0x24;
pub const NODE_COUNTER: Self = Self(0xFF);
pub const SOURCE_WATERMARK: Self = Self(0xFE);
pub const TIMER_WHEEL: Self = Self(0xFD);
pub const TIMER_INDEX: Self = Self(0xFC);
pub const JOIN_ROW_MAPPING: Self = Self(0xFB);
pub const GROUP_ROW_MAPPING: Self = Self(0xFA);
pub const GUEST_ROW_MAPPING: Self = Self(0xF9);
pub const ACCUMULATOR: Self = Self(0x00);
pub const BUFFER: Self = Self(0x01);
pub const RUNNING: Self = Self(0x02);
pub const EMIT: Self = Self(0x03);
pub const ROLLING_EXPIRY: Self = Self(0x04);
pub const COUNT: Self = Self(0x05);
pub const ROW_INDEX: Self = Self(0x06);
pub const SESSION: Self = Self(0x07);
pub const ROLLING_META: Self = Self(0x08);
pub const ENGINE_META: Self = Self(0x09);
pub const DISTINCT_ENTRY: Self = Self(0x0A);
pub const WINDOW_META: Self = Self(0x0B);
pub const JOIN_LEFT: Self = Self(0x0C);
pub const JOIN_RIGHT: Self = Self(0x0D);
pub const JOIN_SCHEMA: Self = Self(0x0E);
pub const RINGBUFFER_FORWARD: Self = Self(0x0F);
pub const RINGBUFFER_ENTRY: Self = Self(0x10);
pub const GATE_VISIBILITY: Self = Self(0x11);
pub const DISTINCT_LAYOUT: Self = Self(0x12);
pub const RINGBUFFER_EXPIRY: Self = Self(0x13);
pub const RINGBUFFER_TTL_ARM: Self = Self(0x14);
pub const SEAL_LEDGER: Self = Self(0x15);
pub const JOIN_PUBLISHED: Self = Self(0x16);
pub const JOIN_PIN: Self = Self(0x17);
pub const RINGBUFFER_META: Self = Self(0x18);
pub const REAP_QUEUE: Self = Self(0x19);
pub const JOIN_ROW_EXPIRY: Self = Self(0x1A);
pub const GUEST_ACCUMULATOR: Self = Self(0x1B);
pub const GUEST_BUFFER: Self = Self(0x1C);
pub const GUEST_RUNNING: Self = Self(0x1D);
pub const TUMBLING_EXPIRY: Self = Self(0x1E);
pub const PARTITIONED_RINGBUFFER_ENTRY: Self = Self(0x1F);
pub const PARTITIONED_RINGBUFFER_EXPIRY: Self = Self(0x20);
pub const PARTITIONED_RINGBUFFER_TTL_ARM: Self = Self(0x21);
pub const PARTITIONED_RINGBUFFER_META: Self = Self(0x22);
pub const CUSTOM_NOT_CACHED: Self = Self(0x23);
pub const JOIN_EXPIRY_DUE: Self = Self(0x24);
pub fn name(&self) -> Cow<'static, str> {
match *self {
Self::NODE_COUNTER => "NODE_COUNTER",
Self::SOURCE_WATERMARK => "SOURCE_WATERMARK",
Self::TIMER_WHEEL => "TIMER_WHEEL",
Self::TIMER_INDEX => "TIMER_INDEX",
Self::JOIN_ROW_MAPPING => "JOIN_ROW_MAPPING",
Self::GROUP_ROW_MAPPING => "GROUP_ROW_MAPPING",
Self::GUEST_ROW_MAPPING => "GUEST_ROW_MAPPING",
Self::ACCUMULATOR => "ACCUMULATOR",
Self::BUFFER => "BUFFER",
Self::RUNNING => "RUNNING",
Self::EMIT => "EMIT",
Self::ROLLING_EXPIRY => "ROLLING_EXPIRY",
Self::COUNT => "COUNT",
Self::ROW_INDEX => "ROW_INDEX",
Self::SESSION => "SESSION",
Self::ROLLING_META => "ROLLING_META",
Self::ENGINE_META => "ENGINE_META",
Self::DISTINCT_ENTRY => "DISTINCT_ENTRY",
Self::WINDOW_META => "WINDOW_META",
Self::JOIN_LEFT => "JOIN_LEFT",
Self::JOIN_RIGHT => "JOIN_RIGHT",
Self::JOIN_SCHEMA => "JOIN_SCHEMA",
Self::RINGBUFFER_FORWARD => "RINGBUFFER_FORWARD",
Self::RINGBUFFER_ENTRY => "RINGBUFFER_ENTRY",
Self::GATE_VISIBILITY => "GATE_VISIBILITY",
Self::DISTINCT_LAYOUT => "DISTINCT_LAYOUT",
Self::RINGBUFFER_EXPIRY => "RINGBUFFER_EXPIRY",
Self::RINGBUFFER_TTL_ARM => "RINGBUFFER_TTL_ARM",
Self::SEAL_LEDGER => "SEAL_LEDGER",
Self::JOIN_PUBLISHED => "JOIN_PUBLISHED",
Self::JOIN_PIN => "JOIN_PIN",
Self::RINGBUFFER_META => "RINGBUFFER_META",
Self::REAP_QUEUE => "REAP_QUEUE",
Self::JOIN_ROW_EXPIRY => "JOIN_ROW_EXPIRY",
Self::JOIN_EXPIRY_DUE => "JOIN_EXPIRY_DUE",
Self::GUEST_ACCUMULATOR => "GUEST_ACCUMULATOR",
Self::GUEST_BUFFER => "GUEST_BUFFER",
Self::GUEST_RUNNING => "GUEST_RUNNING",
Self::TUMBLING_EXPIRY => "TUMBLING_EXPIRY",
Self::PARTITIONED_RINGBUFFER_ENTRY => "PARTITIONED_RINGBUFFER_ENTRY",
Self::PARTITIONED_RINGBUFFER_EXPIRY => "PARTITIONED_RINGBUFFER_EXPIRY",
Self::PARTITIONED_RINGBUFFER_TTL_ARM => "PARTITIONED_RINGBUFFER_TTL_ARM",
Self::PARTITIONED_RINGBUFFER_META => "PARTITIONED_RINGBUFFER_META",
Self::CUSTOM_NOT_CACHED => "CUSTOM_NOT_CACHED",
_ => return Cow::Owned(format!("{:#04x}", self.0)),
}
.into()
}
pub fn is_data(&self) -> bool {
self.0 <= Self::HIGHEST_DATA
}
pub fn is_identity(&self) -> bool {
!self.is_data()
}
pub fn caches_ranges(&self) -> bool {
*self != Self::CUSTOM_NOT_CACHED
}
pub fn is_guest_owned(&self) -> bool {
matches!(*self, Self::CUSTOM_NOT_CACHED)
}
pub const fn is_known(&self) -> bool {
REGISTERED[(self.0 >> 6) as usize] & (1u64 << (self.0 & 63)) != 0
}
}
pub fn is_framed_inner(inner: &[u8]) -> bool {
inner.is_empty() || OperatorStateKey::decode_inner(inner).is_some_and(|(_, keyspace, _)| keyspace.is_known())
}
pub fn is_guest_framed_inner(inner: &[u8]) -> bool {
OperatorStateKey::decode_inner(inner).is_some_and(|(_, keyspace, suffix)| {
keyspace.is_guest_owned() && suffix_width_of(keyspace) == Some(suffix.len())
})
}
pub fn is_identity_framed_inner(inner: &[u8]) -> bool {
OperatorStateKey::decode_inner(inner)
.is_some_and(|(_, keyspace, _)| keyspace.is_identity() && keyspace.is_known())
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct OperatorStateKey {
pub operator: OperatorId,
pub group: GroupId,
pub keyspace: KeyspaceId,
pub suffix: Vec<u8>,
}
impl OperatorStateKey {
pub fn new(operator: OperatorId, group: GroupId, keyspace: KeyspaceId, suffix: impl Into<Vec<u8>>) -> Self {
Self {
operator,
group,
keyspace,
suffix: suffix.into(),
}
}
pub fn root(operator: OperatorId, keyspace: KeyspaceId, suffix: impl Into<Vec<u8>>) -> Self {
Self::new(operator, GroupId::ROOT, keyspace, suffix)
}
pub fn encoded(
operator: OperatorId,
group: GroupId,
keyspace: KeyspaceId,
suffix: impl AsRef<[u8]>,
) -> EncodedKey {
let suffix = suffix.as_ref();
let mut serializer = KeySerializer::with_capacity(NODE_GROUP_PREFIX_LEN + 1 + suffix.len());
serializer
.extend_u8(KeyTag::OperatorState as u8)
.extend_u64(operator.0)
.extend_fixed(*group.as_bytes())
.extend_u8(keyspace.0)
.extend_raw(suffix);
serializer.to_encoded_key()
}
pub fn inner(&self) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(KEYSPACE_INNER_PREFIX_LEN + self.suffix.len());
serializer.extend_fixed(*self.group.as_bytes()).extend_u8(self.keyspace.0).extend_raw(&self.suffix);
serializer.to_encoded_key()
}
pub const KEYSPACE_INNER_OFFSET: u32 = GroupId::WIDTH as u32;
pub fn decode_keyspace(stored: u8) -> KeyspaceId {
KeyspaceId(KeyDeserializer::from_bytes(&[stored]).read_u8().expect("a single byte decodes as u8"))
}
pub fn inner_encoded(group: GroupId, keyspace: KeyspaceId, suffix: impl AsRef<[u8]>) -> GroupStateKey {
let suffix = suffix.as_ref();
let mut serializer = KeySerializer::with_capacity(KEYSPACE_INNER_PREFIX_LEN + suffix.len());
serializer.extend_fixed(*group.as_bytes()).extend_u8(keyspace.0).extend_raw(suffix);
GroupStateKey(serializer.to_encoded_key())
}
pub fn decode_inner(inner: &[u8]) -> Option<(GroupId, KeyspaceId, &[u8])> {
let mut de = KeyDeserializer::from_bytes(inner);
let group = de.read_fixed().ok()?;
let keyspace = de.read_u8().ok()?;
let suffix = de.read_raw(de.remaining()).ok()?;
Some((GroupId::from_bytes(group), KeyspaceId(keyspace), suffix))
}
pub fn node_range(operator: OperatorId) -> TaggedKeyBoundRange {
node_range(operator)
}
pub fn decode_operator(key: &EncodedKey) -> Option<(OperatorId, EncodedKey)> {
let mut de = KeyDeserializer::from_bytes(key.as_slice());
let kind: KeyTag = de.read_u8().ok()?.try_into().ok()?;
if kind != KeyTag::OperatorState {
return None;
}
let operator = de.read_u64().ok()?;
let inner = de.read_raw(de.remaining()).ok()?.to_vec();
Some((OperatorId(operator), EncodedKey::new(inner)))
}
}
impl OperatorStateKey {
pub const TAG: KeyTag = KeyTag::OperatorState;
pub fn encode(&self) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(NODE_GROUP_PREFIX_LEN + 1 + self.suffix.len());
serializer
.extend_u8(KeyTag::OperatorState as u8)
.extend_u64(self.operator.0)
.extend_fixed(*self.group.as_bytes())
.extend_u8(self.keyspace.0)
.extend_raw(&self.suffix);
serializer.to_encoded_key()
}
pub fn decode(key: &EncodedKey) -> Option<Self> {
let mut de = KeyDeserializer::from_bytes(key.as_slice());
let kind: KeyTag = de.read_u8().ok()?.try_into().ok()?;
if kind != KeyTag::OperatorState {
return None;
}
let operator = de.read_u64().ok()?;
let group = de.read_fixed().ok()?;
let keyspace = de.read_u8().ok()?;
let suffix = de.read_raw(de.remaining()).ok()?.to_vec();
Some(Self {
operator: OperatorId(operator),
group: GroupId::from_bytes(group),
keyspace: KeyspaceId(keyspace),
suffix,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct GroupStateKey(EncodedKey);
impl GroupStateKey {
pub fn new(group: GroupId, keyspace: KeyspaceId, suffix: impl AsRef<[u8]>) -> Self {
OperatorStateKey::inner_encoded(group, keyspace, suffix)
}
pub fn root(keyspace: KeyspaceId, suffix: impl AsRef<[u8]>) -> Self {
Self::new(GroupId::ROOT, keyspace, suffix)
}
pub fn from_framed(key: EncodedKey) -> Option<Self> {
is_framed_inner(key.as_slice()).then_some(Self(key))
}
pub fn from_guest_framed(key: EncodedKey) -> Option<Self> {
is_guest_framed_inner(key.as_slice()).then_some(Self(key))
}
pub fn from_identity_framed(key: EncodedKey) -> Option<Self> {
is_identity_framed_inner(key.as_slice()).then_some(Self(key))
}
pub fn bound_unchecked(key: EncodedKey) -> Self {
Self(key)
}
pub fn as_encoded(&self) -> &EncodedKey {
&self.0
}
pub fn into_encoded(self) -> EncodedKey {
self.0
}
pub fn as_slice(&self) -> &[u8] {
self.0.as_slice()
}
pub fn as_bytes(&self) -> &[u8] {
self.0.as_bytes()
}
pub fn group(&self) -> Option<GroupId> {
OperatorStateKey::decode_inner(self.0.as_slice()).map(|(group, _, _)| group)
}
pub fn keyspace(&self) -> Option<KeyspaceId> {
let bytes = self.0.as_slice();
let offset = OperatorStateKey::KEYSPACE_INNER_OFFSET as usize;
(bytes.len() > offset).then(|| KeyspaceId(encode_u8(bytes[offset])))
}
}
impl AsRef<[u8]> for GroupStateKey {
fn as_ref(&self) -> &[u8] {
self.0.as_slice()
}
}
impl AsRef<EncodedKey> for GroupStateKey {
fn as_ref(&self) -> &EncodedKey {
&self.0
}
}
pub trait IntoGroupStateKey {
fn into_group_state_key(self) -> GroupStateKey;
}
impl IntoGroupStateKey for GroupStateKey {
fn into_group_state_key(self) -> GroupStateKey {
self
}
}
fn group_inner_prefix(group: GroupId) -> Vec<u8> {
let mut serializer = KeySerializer::with_capacity(GroupId::WIDTH);
serializer.extend_fixed(*group.as_bytes());
serializer.finish().as_ref().to_vec()
}
fn keyspace_inner_prefix(group: GroupId, keyspace: KeyspaceId) -> Vec<u8> {
let mut prefix = group_inner_prefix(group);
prefix.push(encode_u8(keyspace.0));
prefix
}
pub fn group_inner_range(group: GroupId) -> EncodedKeyRange {
EncodedKeyRange::prefix(&group_inner_prefix(group))
}
pub fn keyspace_inner_range(group: GroupId, keyspace: KeyspaceId) -> EncodedKeyRange {
EncodedKeyRange::prefix(&keyspace_inner_prefix(group, keyspace))
}
enum SuffixEdge {
Low,
High,
}
fn suffix_at_edge(keyspace: KeyspaceId, suffix: &[u8], edge: SuffixEdge) -> Vec<u8> {
struct Pad<'a> {
suffix: &'a [u8],
edge: SuffixEdge,
}
impl KeyspaceVisitor for Pad<'_> {
type Output = Vec<u8>;
fn visit<K: Keyspace>(self) -> Self::Output {
let template = match self.edge {
SuffixEdge::Low => <K::Suffix as BoundedKey>::low().to_suffix_bytes(),
SuffixEdge::High => <K::Suffix as KeyLayout>::high().to_suffix_bytes(),
};
let mut bytes = self.suffix.to_vec();
bytes.truncate(template.len());
bytes.extend_from_slice(&template[bytes.len()..]);
bytes
}
}
dispatch(
keyspace,
Pad {
suffix,
edge,
},
)
.unwrap_or_else(|| suffix.to_vec())
}
pub fn keyspace_inner_range_in(
group: GroupId,
keyspace: KeyspaceId,
start: Bound<&[u8]>,
end: Bound<&[u8]>,
) -> EncodedKeyRange {
let prefix = keyspace_inner_prefix(group, keyspace);
let whole = EncodedKeyRange::prefix(&prefix);
let at = |suffix: &[u8], edge: SuffixEdge| {
let mut key = prefix.clone();
key.extend_from_slice(&suffix_at_edge(keyspace, suffix, edge));
EncodedKey::new(key)
};
let lower = match start {
Bound::Unbounded => whole.start.clone(),
Bound::Included(suffix) => Bound::Included(at(suffix, SuffixEdge::Low)),
Bound::Excluded(suffix) => Bound::Excluded(at(suffix, SuffixEdge::High)),
};
let upper = match end {
Bound::Unbounded => whole.end.clone(),
Bound::Included(suffix) => Bound::Included(at(suffix, SuffixEdge::High)),
Bound::Excluded(suffix) => Bound::Excluded(at(suffix, SuffixEdge::Low)),
};
EncodedKeyRange::new(lower, upper)
}
pub type KeyspaceInnerRangeSplit = (GroupId, KeyspaceId, Bound<Vec<u8>>, Bound<Vec<u8>>);
pub fn keyspace_inner_range_split(range: &EncodedKeyRange) -> Option<KeyspaceInnerRangeSplit> {
let (group, keyspace, start) = match &range.start {
Bound::Included(key) => {
let (group, keyspace, suffix) = OperatorStateKey::decode_inner(key.as_slice())?;
(group, keyspace, Bound::Included(suffix.to_vec()))
}
Bound::Excluded(key) => {
let (group, keyspace, suffix) = OperatorStateKey::decode_inner(key.as_slice())?;
(group, keyspace, Bound::Excluded(suffix.to_vec()))
}
Bound::Unbounded => return None,
};
let whole = EncodedKeyRange::prefix(&keyspace_inner_prefix(group, keyspace));
let end = if range.end == whole.end {
Bound::Unbounded
} else {
match &range.end {
Bound::Included(key) => match OperatorStateKey::decode_inner(key.as_slice())? {
(g, k, suffix) if g == group && k == keyspace => Bound::Included(suffix.to_vec()),
_ => return None,
},
Bound::Excluded(key) => match OperatorStateKey::decode_inner(key.as_slice())? {
(g, k, suffix) if g == group && k == keyspace => Bound::Excluded(suffix.to_vec()),
_ => return None,
},
Bound::Unbounded => return None,
}
};
Some((group, keyspace, start, end))
}
pub fn group_inner_range_split(range: &EncodedKeyRange) -> Option<GroupId> {
let key = match &range.start {
Bound::Included(key) | Bound::Excluded(key) => key,
Bound::Unbounded => return None,
};
let group = GroupId::from_bytes(KeyDeserializer::from_bytes(key.as_slice()).read_fixed().ok()?);
for candidate in [group_inner_range(group), group_data_inner_range(group)] {
if range.start == candidate.start && range.end == candidate.end {
return Some(group);
}
}
None
}
pub fn custom_not_cached_key_in(group: GroupId, id: &[u8]) -> Option<GroupStateKey> {
CustomNotCachedSuffix::of(id)
.map(|key| OperatorStateKey::inner_encoded(group, KeyspaceId::CUSTOM_NOT_CACHED, key.to_suffix_bytes()))
}
pub fn custom_not_cached_key(id: &[u8]) -> Option<GroupStateKey> {
custom_not_cached_key_in(GroupId::ROOT, id)
}
pub fn node_counter_key(kind: NodeCounterKind) -> GroupStateKey {
typed_key::<NodeCounter>(GroupId::ROOT, &NodeCounterKey::of(kind))
}
pub fn row_number_counter_key() -> GroupStateKey {
node_counter_key(NodeCounterKind::RowNumber)
}
pub fn keyspace_inner_range_upto(group: GroupId, keyspace: KeyspaceId, suffix: &[u8]) -> EncodedKeyRange {
let mut bound = keyspace_inner_prefix(group, keyspace);
bound.extend_from_slice(suffix);
EncodedKeyRange::new(keyspace_inner_range(group, keyspace).start, EncodedKeyRange::prefix(&bound).end)
}
pub fn group_data_inner_range(group: GroupId) -> EncodedKeyRange {
let prefix = group_inner_prefix(group);
let mut start = prefix.clone();
start.push(encode_u8(KeyspaceId::HIGHEST_DATA));
EncodedKeyRange::new(Bound::Included(EncodedKey::new(start)), EncodedKeyRange::prefix(&prefix).end)
}
pub fn group_identity_inner_range(group: GroupId) -> EncodedKeyRange {
let prefix = group_inner_prefix(group);
let mut end = prefix.clone();
end.push(encode_u8(KeyspaceId::HIGHEST_DATA));
EncodedKeyRange::new(Bound::Included(EncodedKey::new(prefix)), Bound::Excluded(EncodedKey::new(end)))
}
pub const NODE_PREFIX_LEN: usize = 9;
pub const KEYSPACE_INNER_PREFIX_LEN: usize = GroupId::WIDTH + size_of::<u8>();
pub const NODE_GROUP_PREFIX_LEN: usize = NODE_PREFIX_LEN + GroupId::WIDTH;
pub fn extend_node_prefix(serializer: &mut KeySerializer, operator: OperatorId) {
serializer.extend_u8(KeyTag::OperatorState as u8).extend_u64(operator.0);
}
pub fn node_prefix(operator: OperatorId) -> Vec<u8> {
let mut serializer = KeySerializer::with_capacity(NODE_PREFIX_LEN);
extend_node_prefix(&mut serializer, operator);
serializer.finish().as_ref().to_vec()
}
pub fn node_range(operator: OperatorId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(KeyTag::OperatorState, [Field::UDesc(Width::U64, operator.0 as u128)])
}
pub fn group_range(operator: OperatorId, group: GroupId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
KeyTag::OperatorState,
[
Field::UDesc(Width::U64, operator.0 as u128),
Field::BytesDesc(ByteEncoding::Fixed, Cow::Owned(group.as_bytes().to_vec())),
],
)
}
pub fn keyspace_range(operator: OperatorId, group: GroupId, keyspace: KeyspaceId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
KeyTag::OperatorState,
[
Field::UDesc(Width::U64, operator.0 as u128),
Field::BytesDesc(ByteEncoding::Fixed, Cow::Owned(group.as_bytes().to_vec())),
Field::UDesc(Width::U8, keyspace.0 as u128),
],
)
}
pub fn group_data_range(operator: OperatorId, group: GroupId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange {
start: Bound::Included(TaggedKeyBound::prefix(
KeyTag::OperatorState,
[
Field::UDesc(Width::U64, operator.0 as u128),
Field::BytesDesc(ByteEncoding::Fixed, Cow::Owned(group.as_bytes().to_vec())),
Field::UDesc(Width::U8, KeyspaceId::HIGHEST_DATA as u128),
],
)),
end: Bound::Excluded(TaggedKeyBound::prefix_end(
KeyTag::OperatorState,
[
Field::UDesc(Width::U64, operator.0 as u128),
Field::BytesDesc(ByteEncoding::Fixed, Cow::Owned(group.as_bytes().to_vec())),
],
)),
}
}
pub fn group_identity_range(operator: OperatorId, group: GroupId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange {
start: Bound::Included(TaggedKeyBound::prefix(
KeyTag::OperatorState,
[
Field::UDesc(Width::U64, operator.0 as u128),
Field::BytesDesc(ByteEncoding::Fixed, Cow::Owned(group.as_bytes().to_vec())),
],
)),
end: Bound::Excluded(TaggedKeyBound::prefix(
KeyTag::OperatorState,
[
Field::UDesc(Width::U64, operator.0 as u128),
Field::BytesDesc(ByteEncoding::Fixed, Cow::Owned(group.as_bytes().to_vec())),
Field::UDesc(Width::U8, KeyspaceId::HIGHEST_DATA as u128),
],
)),
}
}
#[cfg(test)]
mod tests {
use std::ops::Bound;
use reifydb_value::util::hash::Hash128;
use super::{
EncodedKey, EncodedKeyRange, GroupId, GroupSet, GroupStateKey, KeySerializer, KeyspaceId,
OperatorStateKey, custom_not_cached_key_in, group_data_inner_range, group_data_of_inner,
group_data_range, group_identity_inner_range, group_identity_range, group_inner_prefix,
group_inner_range, group_range, is_framed_inner, is_guest_framed_inner, keyspace_range, node_prefix,
node_range,
};
use crate::interface::catalog::flow::OperatorId;
const NODES: [u64; 4] = [1, 17, 300, 70_000];
const GROUPS: [u128; 8] = [1, 2, 127, 128, 1000, 100_000, 1 << 30, u128::MAX];
const DATA_KEYSPACES: [KeyspaceId; 4] =
[KeyspaceId::ACCUMULATOR, KeyspaceId::BUFFER, KeyspaceId::RUNNING, KeyspaceId::CUSTOM_NOT_CACHED];
const IDENTITY_KEYSPACES: [KeyspaceId; 1] = [KeyspaceId::GUEST_ROW_MAPPING];
#[derive(Clone, Copy, PartialEq, Debug)]
enum Phase {
Data,
Identity,
}
const CENSUS: [(&str, KeyspaceId, Phase, bool); 44] = [
("NODE_COUNTER", KeyspaceId::NODE_COUNTER, Phase::Identity, true),
("SOURCE_WATERMARK", KeyspaceId::SOURCE_WATERMARK, Phase::Identity, true),
("TIMER_WHEEL", KeyspaceId::TIMER_WHEEL, Phase::Identity, true),
("TIMER_INDEX", KeyspaceId::TIMER_INDEX, Phase::Identity, true),
("JOIN_ROW_MAPPING", KeyspaceId::JOIN_ROW_MAPPING, Phase::Identity, true),
("GROUP_ROW_MAPPING", KeyspaceId::GROUP_ROW_MAPPING, Phase::Identity, true),
("GUEST_ROW_MAPPING", KeyspaceId::GUEST_ROW_MAPPING, Phase::Identity, true),
("ACCUMULATOR", KeyspaceId::ACCUMULATOR, Phase::Data, true),
("BUFFER", KeyspaceId::BUFFER, Phase::Data, true),
("RUNNING", KeyspaceId::RUNNING, Phase::Data, true),
("EMIT", KeyspaceId::EMIT, Phase::Data, true),
("ROLLING_EXPIRY", KeyspaceId::ROLLING_EXPIRY, Phase::Data, true),
("COUNT", KeyspaceId::COUNT, Phase::Data, true),
("ROW_INDEX", KeyspaceId::ROW_INDEX, Phase::Data, true),
("SESSION", KeyspaceId::SESSION, Phase::Data, true),
("ROLLING_META", KeyspaceId::ROLLING_META, Phase::Data, true),
("ENGINE_META", KeyspaceId::ENGINE_META, Phase::Data, true),
("DISTINCT_ENTRY", KeyspaceId::DISTINCT_ENTRY, Phase::Data, true),
("WINDOW_META", KeyspaceId::WINDOW_META, Phase::Data, true),
("JOIN_LEFT", KeyspaceId::JOIN_LEFT, Phase::Data, true),
("JOIN_RIGHT", KeyspaceId::JOIN_RIGHT, Phase::Data, true),
("JOIN_SCHEMA", KeyspaceId::JOIN_SCHEMA, Phase::Data, true),
("RINGBUFFER_FORWARD", KeyspaceId::RINGBUFFER_FORWARD, Phase::Data, true),
("RINGBUFFER_ENTRY", KeyspaceId::RINGBUFFER_ENTRY, Phase::Data, true),
("GATE_VISIBILITY", KeyspaceId::GATE_VISIBILITY, Phase::Data, true),
("DISTINCT_LAYOUT", KeyspaceId::DISTINCT_LAYOUT, Phase::Data, true),
("RINGBUFFER_EXPIRY", KeyspaceId::RINGBUFFER_EXPIRY, Phase::Data, true),
("RINGBUFFER_TTL_ARM", KeyspaceId::RINGBUFFER_TTL_ARM, Phase::Data, true),
("SEAL_LEDGER", KeyspaceId::SEAL_LEDGER, Phase::Data, true),
("JOIN_PUBLISHED", KeyspaceId::JOIN_PUBLISHED, Phase::Data, true),
("JOIN_PIN", KeyspaceId::JOIN_PIN, Phase::Data, true),
("RINGBUFFER_META", KeyspaceId::RINGBUFFER_META, Phase::Data, true),
("REAP_QUEUE", KeyspaceId::REAP_QUEUE, Phase::Data, true),
("JOIN_ROW_EXPIRY", KeyspaceId::JOIN_ROW_EXPIRY, Phase::Data, true),
("JOIN_EXPIRY_DUE", KeyspaceId::JOIN_EXPIRY_DUE, Phase::Data, true),
("GUEST_ACCUMULATOR", KeyspaceId::GUEST_ACCUMULATOR, Phase::Data, true),
("GUEST_BUFFER", KeyspaceId::GUEST_BUFFER, Phase::Data, true),
("GUEST_RUNNING", KeyspaceId::GUEST_RUNNING, Phase::Data, true),
("TUMBLING_EXPIRY", KeyspaceId::TUMBLING_EXPIRY, Phase::Data, true),
("PARTITIONED_RINGBUFFER_ENTRY", KeyspaceId::PARTITIONED_RINGBUFFER_ENTRY, Phase::Data, true),
("PARTITIONED_RINGBUFFER_EXPIRY", KeyspaceId::PARTITIONED_RINGBUFFER_EXPIRY, Phase::Data, true),
("PARTITIONED_RINGBUFFER_TTL_ARM", KeyspaceId::PARTITIONED_RINGBUFFER_TTL_ARM, Phase::Data, true),
("PARTITIONED_RINGBUFFER_META", KeyspaceId::PARTITIONED_RINGBUFFER_META, Phase::Data, true),
("CUSTOM_NOT_CACHED", KeyspaceId::CUSTOM_NOT_CACHED, Phase::Data, false),
];
fn declared_keyspaces() -> usize {
let source = include_str!("state.rs");
let body = source
.split("impl KeyspaceId {")
.nth(1)
.expect("the KeyspaceId impl block is where the constants are declared");
let body = body.split("\n}\n").next().expect("the impl block is closed");
body.lines()
.filter(|line| {
let line = line.trim_start();
line.starts_with("pub const") && line.contains("Self(")
})
.count()
}
#[test]
fn a_bare_row_number_key_is_too_short_to_be_read_as_a_framed_key() {
let mut bare = KeySerializer::with_capacity(4);
bare.extend_u64(7u64);
let bare = bare.finish().as_ref().to_vec();
assert!(
bare.len() < group_inner_prefix(GroupId::hashed(Hash128(7))).len(),
"a bare row number cannot span a group"
);
assert!(OperatorStateKey::decode_inner(&bare).is_none());
assert!(!is_framed_inner(&bare));
let framed = OperatorStateKey::inner_encoded(
GroupId::ROOT,
KeyspaceId::CUSTOM_NOT_CACHED,
7u64.to_be_bytes(),
);
assert!(is_framed_inner(framed.as_slice()));
assert!(
!contains(&group_identity_inner_range(GroupId::hashed(Hash128(7))), framed.as_slice()),
"the framed form must sit outside every other group's range"
);
}
#[test]
fn the_empty_key_is_framing_because_it_sorts_below_every_group() {
let empty: &[u8] = &[];
assert!(is_framed_inner(empty));
for group in GROUPS {
let range = group_inner_range(GroupId::hashed(Hash128(group)));
assert!(
!contains(&range, empty),
"the empty key must sit outside group {group}'s range, not merely be unattributed"
);
}
}
#[test]
fn the_empty_key_is_not_guest_framing_even_though_it_is_host_framing() {
let empty: &[u8] = &[];
assert!(is_framed_inner(empty));
assert!(!is_guest_framed_inner(empty));
assert!(GroupStateKey::from_guest_framed(EncodedKey::new(Vec::new())).is_none());
assert!(is_guest_framed_inner(
custom_not_cached_key_in(GroupId::hashed(Hash128(3)), &[])
.expect("an empty id fits the keyspace")
.as_slice()
));
assert!(
!is_guest_framed_inner(
OperatorStateKey::inner_encoded(
GroupId::hashed(Hash128(3)),
KeyspaceId::CUSTOM_NOT_CACHED,
[]
)
.as_slice()
),
"a suffix narrower than its keyspace declares must be refused at the wall, or it reaches the typed bucket and panics there instead"
);
}
#[test]
fn a_keyspace_this_substrate_never_defines_is_not_framing() {
let mut stray = KeySerializer::with_capacity(4);
stray.extend_u64(3u64).extend_u8(0x90u8);
assert!(!is_framed_inner(stray.finish().as_ref()));
for keyspace in DATA_KEYSPACES.iter().chain(IDENTITY_KEYSPACES.iter()) {
assert!(
is_framed_inner(
OperatorStateKey::inner_encoded(GroupId::hashed(Hash128(3)), *keyspace, [])
.as_slice()
),
"keyspace {keyspace:?} is one the substrate writes and must pass"
);
}
}
fn contains(range: &EncodedKeyRange, key: &[u8]) -> bool {
let after_start = match &range.start {
Bound::Included(start) => key >= start.as_slice(),
Bound::Excluded(start) => key > start.as_slice(),
Bound::Unbounded => true,
};
let before_end = match &range.end {
Bound::Included(end) => key <= end.as_slice(),
Bound::Excluded(end) => key < end.as_slice(),
Bound::Unbounded => true,
};
after_start && before_end
}
fn population() -> Vec<OperatorStateKey> {
let mut keys = Vec::new();
for operator in NODES {
for group in GROUPS {
for keyspace in DATA_KEYSPACES.iter().chain(IDENTITY_KEYSPACES.iter()) {
for coord in [0u64, 1, 999, u64::MAX] {
keys.push(OperatorStateKey::new(
OperatorId(operator),
GroupId::hashed(Hash128(group)),
*keyspace,
coord.to_be_bytes().to_vec(),
));
}
}
}
keys.push(OperatorStateKey::root(
OperatorId(operator),
KeyspaceId::NODE_COUNTER,
b"7xKXtg2CW87d97TXJSDpbD5jBkheTqA83TZRuJosgAsU".to_vec(),
));
}
keys
}
#[test]
fn a_group_range_contains_exactly_that_groups_keys() {
let population = population();
for operator in NODES {
for group in GROUPS {
let group_id = GroupId::hashed(Hash128(group));
let range = group_range(OperatorId(operator), group_id).encode();
for key in &population {
let encoded = key.encode();
let expected = key.operator.0 == operator && key.group == group_id;
assert_eq!(
contains(&range, encoded.as_slice()),
expected,
"operator {operator} group {group} range disagreed about a key of operator {} \
group {}",
key.operator.0,
key.group
);
}
}
}
}
#[test]
fn variable_length_group_ids_cannot_prefix_one_another() {
let encodings: Vec<Vec<u8>> = GROUPS
.iter()
.map(|group| {
OperatorStateKey::new(
OperatorId(1),
GroupId::hashed(Hash128(*group)),
KeyspaceId::ACCUMULATOR,
vec![],
)
.encode()
.as_slice()
.to_vec()
})
.collect();
for (i, a) in encodings.iter().enumerate() {
for (j, b) in encodings.iter().enumerate() {
if i != j {
assert!(
!b.starts_with(a.as_slice()),
"group {} encodes as a prefix of group {}",
GROUPS[i],
GROUPS[j]
);
}
}
}
}
#[test]
fn the_data_and_identity_ranges_partition_the_group() {
for operator in NODES {
for group in GROUPS {
let group = GroupId::hashed(Hash128(group));
let data = group_data_range(OperatorId(operator), group).encode();
let identity = group_identity_range(OperatorId(operator), group).encode();
for keyspace in DATA_KEYSPACES {
let key = OperatorStateKey::new(
OperatorId(operator),
group,
keyspace,
vec![7, 7],
)
.encode();
assert!(
contains(&data, key.as_slice()),
"data keyspace {keyspace:?} must fall in the phase-1 range"
);
assert!(
!contains(&identity, key.as_slice()),
"data keyspace {keyspace:?} must not fall in the phase-2 range"
);
}
for keyspace in IDENTITY_KEYSPACES {
let key = OperatorStateKey::new(
OperatorId(operator),
group,
keyspace,
vec![7, 7],
)
.encode();
assert!(
contains(&identity, key.as_slice()),
"identity keyspace {keyspace:?} must fall in the phase-2 range"
);
assert!(
!contains(&data, key.as_slice()),
"identity keyspace {keyspace:?} must survive phase 1"
);
}
}
}
}
#[test]
fn every_declared_keyspace_names_itself_for_offline_attribution() {
for (name, keyspace, _, _) in CENSUS {
assert_eq!(
keyspace.name(),
name,
"{name} ({:#04x}) does not name itself, so an offline census reports it as CUSTOM",
keyspace.0
);
}
assert_eq!(
KeyspaceId::CUSTOM_NOT_CACHED.name(),
"CUSTOM_NOT_CACHED",
"a custom keyspace names the admission side it sits on; there is no unnamed fallback to absorb it"
);
}
#[test]
fn every_declared_keyspace_states_whether_it_may_be_range_cached() {
for (name, keyspace, _, policy) in CENSUS {
assert_eq!(
keyspace.caches_ranges(),
policy,
"{name} ({:#04x}) is cached on a different side than the census records",
keyspace.0
);
}
let uncached: Vec<&str> = CENSUS.iter().filter(|(_, _, _, p)| !*p).map(|(n, ..)| *n).collect();
assert_eq!(
uncached,
["CUSTOM_NOT_CACHED"],
"widening the set the tier refuses turns that tier into an off switch and only shows up as a \
throughput loss in a replay, so every move in or out is a measured decision"
);
assert!(
KeyspaceId(0x43).caches_ranges(),
"an undeclared keyspace must default to cacheable, or a custom operator silently loses the \
range tier"
);
}
#[test]
fn every_declared_keyspace_is_distinct_framing_and_swept_by_exactly_one_phase() {
assert_eq!(
CENSUS.len(),
declared_keyspaces(),
"a keyspace was added to KeyspaceId without being added to the census, so nothing below \
ever looks at its byte"
);
let mut seen: Vec<(&str, u8)> = Vec::new();
for (name, keyspace, phase, _) in CENSUS {
if let Some((other, _)) = seen.iter().find(|(_, byte)| *byte == keyspace.0) {
panic!("{name} and {other} both claim keyspace byte {:#04x}", keyspace.0);
}
seen.push((name, keyspace.0));
assert!(
keyspace.is_known(),
"{name} is declared but not framing, so the sweep panics on the first row it holds"
);
let group = GroupId::hashed(Hash128(4));
let key = OperatorStateKey::new(OperatorId(9), group, keyspace, vec![7, 7]).encode();
let data = contains(&group_data_range(OperatorId(9), group).encode(), key.as_slice());
let identity = contains(&group_identity_range(OperatorId(9), group).encode(), key.as_slice());
assert!(data != identity, "{name} must fall in exactly one phase, not {data} and {identity}");
assert_eq!(
data,
phase == Phase::Data,
"{name} is declared {phase:?} but the phase-1 range says data={data}"
);
assert_eq!(
keyspace.is_data(),
phase == Phase::Data,
"{name} is declared {phase:?} but is_data says {}",
keyspace.is_data()
);
}
}
#[test]
fn root_entries_sit_outside_every_group_range() {
for operator in NODES {
let counter = OperatorStateKey::root(
OperatorId(operator),
KeyspaceId::NODE_COUNTER,
b"mint".to_vec(),
)
.encode();
for group in GROUPS {
let range = group_range(OperatorId(operator), GroupId::hashed(Hash128(group))).encode();
assert!(
!contains(&range, counter.as_slice()),
"group {group} range must not contain the root group's counter"
);
}
}
}
#[test]
fn a_node_range_contains_exactly_that_nodes_keys() {
let population = population();
for operator in NODES {
let range = node_range(OperatorId(operator)).encode();
for key in &population {
let encoded = key.encode();
assert_eq!(
contains(&range, encoded.as_slice()),
key.operator.0 == operator,
"operator {operator} range disagreed about a key of operator {}",
key.operator.0
);
}
}
}
#[test]
fn a_keyspace_range_isolates_one_keyspace_of_one_group() {
let operator = OperatorId(17);
let group = GroupId::hashed(Hash128(42));
let range = keyspace_range(operator, group, KeyspaceId::BUFFER).encode();
let inside = OperatorStateKey::new(operator, group, KeyspaceId::BUFFER, vec![1]).encode();
assert!(contains(&range, inside.as_slice()));
for other in [KeyspaceId::ACCUMULATOR, KeyspaceId::RUNNING, KeyspaceId::GUEST_ROW_MAPPING] {
let key = OperatorStateKey::new(operator, group, other, vec![1]).encode();
assert!(!contains(&range, key.as_slice()), "keyspace {other:?} leaked into the buffer range");
}
let other_group =
OperatorStateKey::new(operator, GroupId::hashed(Hash128(43)), KeyspaceId::BUFFER, vec![1])
.encode();
assert!(!contains(&range, other_group.as_slice()), "another group's buffer leaked into the range");
}
#[test]
fn encode_decode_round_trips_every_component() {
let key = OperatorStateKey::new(
OperatorId(0xDEAD_BEEF),
GroupId::hashed(Hash128(123_456)),
KeyspaceId::CUSTOM_NOT_CACHED,
vec![1, 2, 3, 4],
);
assert_eq!(OperatorStateKey::decode(&key.encode()), Some(key));
}
#[test]
fn keys_still_decode_as_operator_state_of_their_node() {
let key = OperatorStateKey::new(
OperatorId(9),
GroupId::hashed(Hash128(4)),
KeyspaceId::ACCUMULATOR,
vec![1],
)
.encode();
let decoded = OperatorStateKey::decode(&key).expect("must remain decodable as its key kind");
assert_eq!(decoded.operator, OperatorId(9));
}
#[test]
fn an_inner_key_composed_with_its_node_prefix_reproduces_the_full_key() {
let key = OperatorStateKey::new(
OperatorId(17),
GroupId::hashed(Hash128(42)),
KeyspaceId::BUFFER,
vec![9, 9],
);
let mut composed = node_prefix(OperatorId(17));
composed.extend_from_slice(key.inner().as_slice());
assert_eq!(composed, key.encode().as_slice(), "inner key plus operator prefix must equal the full key");
}
#[test]
fn the_root_group_range_stays_inside_its_node() {
let range = group_inner_range(GroupId::ROOT).with_prefix(EncodedKey::new(node_prefix(OperatorId(17))));
let own = OperatorStateKey::root(OperatorId(17), KeyspaceId::NODE_COUNTER, vec![1]).encode();
assert!(contains(&range, own.as_slice()), "the operator's own counter must be in range");
for operator in NODES {
if operator == 17 {
continue;
}
for keyspace in [KeyspaceId::NODE_COUNTER, KeyspaceId::ACCUMULATOR] {
let foreign =
OperatorStateKey::new(OperatorId(operator), GroupId::ROOT, keyspace, vec![1])
.encode();
assert!(
!contains(&range, foreign.as_slice()),
"operator {operator} leaked into operator 17's root-group range"
);
}
}
}
#[test]
fn inner_ranges_partition_the_group_like_their_full_key_counterparts() {
let operator = OperatorId(17);
let prefix = EncodedKey::new(node_prefix(operator));
for group in GROUPS {
let group = GroupId::hashed(Hash128(group));
let data = group_data_inner_range(group).with_prefix(prefix.clone());
let identity = group_identity_inner_range(group).with_prefix(prefix.clone());
for keyspace in DATA_KEYSPACES {
let key = OperatorStateKey::new(operator, group, keyspace, vec![7]).encode();
assert!(contains(&data, key.as_slice()));
assert!(!contains(&identity, key.as_slice()));
}
for keyspace in IDENTITY_KEYSPACES {
let key = OperatorStateKey::new(operator, group, keyspace, vec![7]).encode();
assert!(contains(&identity, key.as_slice()));
assert!(!contains(&data, key.as_slice()));
}
}
}
#[test]
fn decode_inner_round_trips_the_tail() {
let key = OperatorStateKey::new(
OperatorId(3),
GroupId::hashed(Hash128(77)),
KeyspaceId::EMIT,
vec![4, 5, 6],
);
let inner = key.inner();
let (group, keyspace, suffix) =
OperatorStateKey::decode_inner(inner.as_slice()).expect("inner must decode");
assert_eq!(group, GroupId::hashed(Hash128(77)));
assert_eq!(keyspace, KeyspaceId::EMIT);
assert_eq!(suffix, [4, 5, 6]);
}
#[test]
fn a_state_key_stays_compact_however_long_the_group_key_is() {
let long = b"7xKXtg2CW87d97TXJSDpbD5jBkheTqA83TZRuJosgAsU\
So11111111111111111111111111111111111111112";
let hashed = OperatorStateKey::new(
OperatorId(17),
GroupId::of(&EncodedKey::new(long.to_vec())),
KeyspaceId::ACCUMULATOR,
vec![0; 8],
)
.encode();
let short = OperatorStateKey::new(
OperatorId(17),
GroupId::of(&EncodedKey::new(b"g".to_vec())),
KeyspaceId::ACCUMULATOR,
vec![0; 8],
)
.encode();
assert_eq!(hashed.as_slice().len(), short.as_slice().len(), "group key length must not reach the key");
assert!(hashed.as_slice().len() * 2 < long.len() + 8, "and must stay far below embedding the bytes");
}
#[test]
fn the_ram_predicate_and_the_disk_range_agree_on_every_key() {
for group in GROUPS.map(|group| GroupId::hashed(Hash128(group))) {
let range = group_data_inner_range(group);
for other in GROUPS.map(|group| GroupId::hashed(Hash128(group))) {
for keyspace in DATA_KEYSPACES.iter().chain(IDENTITY_KEYSPACES.iter()) {
let key = OperatorStateKey::inner_encoded(other, *keyspace, vec![7, 7]);
let in_range = contains(&range, key.as_slice());
let in_predicate = group_data_of_inner(key.as_slice()) == Some(group);
assert_eq!(
in_range, in_predicate,
"disk range and RAM predicate disagree for group {group:?} on a \
{keyspace:?} key of group {other:?}"
);
}
}
}
}
#[test]
fn the_ram_predicate_refuses_identity_keyspaces() {
for keyspace in IDENTITY_KEYSPACES {
let key = OperatorStateKey::inner_encoded(GroupId::hashed(Hash128(9)), keyspace, vec![1]);
assert_eq!(
group_data_of_inner(key.as_slice()),
None,
"{keyspace:?} must not be reported as reclaimable group data"
);
}
}
#[test]
fn a_key_too_short_to_carry_a_keyspace_is_refused() {
assert_eq!(group_data_of_inner(&[]), None);
assert_eq!(group_data_of_inner(&[0xAB]), None, "a group with no keyspace byte must not decode");
}
#[test]
fn the_predicate_agrees_with_the_disk_range_on_arbitrary_bytes() {
let mut seed = 0x2545F4914F6CDD1Du64;
let mut next = move || {
seed ^= seed << 13;
seed ^= seed >> 7;
seed ^= seed << 17;
seed
};
for _ in 0..2000 {
let len = (next() % 12) as usize;
let key: Vec<u8> = (0..len).map(|_| (next() % 256) as u8).collect();
let Some(group) = group_data_of_inner(&key) else {
continue;
};
assert!(
contains(&group_data_inner_range(group), &key),
"predicate attributed {key:?} to {group:?} but the disk range excludes it"
);
}
}
#[test]
fn a_group_set_is_sorted_deduped_and_never_admits_root() {
let two = GroupId::hashed(Hash128(2));
let five = GroupId::hashed(Hash128(5));
let nine = GroupId::hashed(Hash128(9));
let set = GroupSet::new([nine, two, nine, GroupId::ROOT, five]);
assert_eq!(set.as_slice(), &[two, five, nine]);
assert_eq!(set.len(), 3);
assert!(set.contains(five));
assert!(!set.contains(GroupId::hashed(Hash128(3))));
assert!(!set.contains(GroupId::ROOT), "the root group must be filtered out, not merely unsorted");
}
#[test]
fn an_empty_group_set_matches_nothing() {
let set = GroupSet::new([]);
assert!(set.is_empty());
assert!(!set.contains(GroupId::FIRST_NON_ROOT));
}
}
impl KeyFields for OperatorStateKey {
fn fields(&self) -> SmallVec<[Field<'_>; 6]> {
smallvec![
Field::UDesc(Width::U64, self.operator.0 as u128),
Field::BytesDesc(ByteEncoding::Fixed, Cow::Borrowed(self.group.as_bytes())),
Field::UDesc(Width::U8, self.keyspace.0 as u128),
Field::RawAsc(RawEncoding::Verbatim, Cow::Borrowed(&self.suffix)),
]
}
}