use std::collections::Bound;
use reifydb_codec::key::{
deserializer::KeyDeserializer,
encoded::{EncodedKey, EncodedKeyRange},
};
use reifydb_macro::KeyCodec;
use reifydb_value::value::partition::Partition;
use super::KeyTag;
use crate::{
interface::catalog::{id::SeriesId, object::ObjectId, storage::StorageId},
key::{
any::{Field, KeyFields, TaggedKey, Width},
bound::{OwnedField, TaggedKeyBound, TaggedKeyBoundRange, object_fields},
catalog::{KeyDeserializerCatalogExt, KeySerializerCatalogExt},
typed::{BoundedKey, DenseKey, direction::Desc},
},
metrics::heap::HeapSize,
};
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = Series)]
pub struct SeriesKey {
pub series: SeriesId,
}
impl SeriesKey {
pub fn new(series: SeriesId) -> Self {
Self {
series,
}
}
pub fn encoded(series: impl Into<SeriesId>) -> EncodedKey {
Self::new(series.into()).encode()
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = SeriesMetadata)]
pub struct SeriesMetadataKey {
pub storage: StorageId,
}
impl SeriesMetadataKey {
pub fn new(storage: impl Into<StorageId>) -> Self {
Self {
storage: storage.into(),
}
}
pub fn encoded(storage: impl Into<StorageId>) -> EncodedKey {
Self::new(storage).encode()
}
}
#[cfg(test)]
mod series_metadata_key_tests {
use reifydb_codec::key::serializer::KeySerializer;
use super::{KeyTag, SeriesKey, SeriesMetadataKey};
use crate::interface::catalog::{
id::{SeriesId, ViewId},
storage::StorageId,
};
#[test]
fn test_metadata_key_roundtrip_series() {
let key = SeriesMetadataKey {
storage: StorageId::Series(SeriesId(7)),
};
assert_eq!(SeriesMetadataKey::decode(&key.encode()).unwrap(), key);
}
#[test]
fn test_metadata_key_roundtrip_view() {
let key = SeriesMetadataKey {
storage: StorageId::View(ViewId(7)),
};
assert_eq!(SeriesMetadataKey::decode(&key.encode()).unwrap(), key);
}
#[test]
fn test_metadata_key_separates_a_series_from_a_view_with_the_same_id() {
assert_ne!(SeriesMetadataKey::encoded(SeriesId(7)), SeriesMetadataKey::encoded(ViewId(7)));
}
#[test]
fn test_series_key_matches_legacy_byte_layout() {
for id in [SeriesId(0), SeriesId(1), SeriesId(u64::MAX)] {
let mut legacy = KeySerializer::with_capacity(9);
legacy.extend_u8(KeyTag::Series as u8).extend_u64(id);
assert_eq!(legacy.to_encoded_key().as_slice(), SeriesKey::encoded(id).as_slice());
}
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = SeriesRow)]
pub struct SeriesRowKey {
pub storage: StorageId,
pub variant_tag: Option<u8>,
pub key: u64,
pub sequence: u64,
}
#[derive(Debug, Clone)]
pub struct SeriesRowKeyRange {
pub storage: StorageId,
pub variant_tag: Option<u8>,
pub key_start: Option<u64>,
pub key_end: Option<u64>,
}
impl SeriesRowKeyRange {
pub fn storage_start(storage: StorageId) -> EncodedKey {
TaggedKeyBound::prefix(SeriesRowKey::TAG, object_fields(ObjectId::from(storage))).encode()
}
pub fn storage_end(storage: StorageId) -> EncodedKey {
TaggedKeyBound::prefix(SeriesRowKey::TAG, object_fields(ObjectId::from(storage).prev())).encode()
}
pub fn full_scan(storage: StorageId, variant_tag: Option<u8>) -> TaggedKeyBoundRange {
SeriesRowKeyRange {
storage,
variant_tag,
key_start: None,
key_end: None,
}
.bounds()
}
pub fn scan_range(
storage: StorageId,
variant_tag: Option<u8>,
key_start: Option<u64>,
key_end: Option<u64>,
last: Option<&TaggedKey>,
) -> TaggedKeyBoundRange {
if matches!(key_end, Some(0)) {
return TaggedKeyBoundRange::empty(SeriesRowKey::TAG);
}
SeriesRowKeyRange {
storage,
variant_tag,
key_start,
key_end,
}
.bounds()
.resume_after(last)
}
pub fn decode_storage(key: &EncodedKey) -> Option<StorageId> {
let mut de = KeyDeserializer::from_bytes(key.as_slice());
let kind: KeyTag = de.read_u8().ok()?.try_into().ok()?;
if kind != SeriesRowKey::TAG {
return None;
}
StorageId::from_object(de.read_object_id().ok()?)
}
pub fn decode(range: &EncodedKeyRange) -> (Option<StorageId>, Option<StorageId>) {
let start = match &range.start {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
Bound::Unbounded => None,
};
let end = match &range.end {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
Bound::Unbounded => None,
};
(start, end)
}
fn head_fields(&self, tagged: bool) -> Vec<OwnedField> {
let mut fields = object_fields(ObjectId::from(self.storage)).to_vec();
match self.variant_tag {
Some(tag) => {
fields.push(Field::UDesc(Width::U8, 1));
fields.push(Field::UDesc(Width::U8, tag as u128));
}
None if tagged || self.key_start.is_some() || self.key_end.is_some() => {
fields.push(Field::UDesc(Width::U8, 0));
fields.push(Field::UDesc(Width::U8, 0));
}
None => {}
}
fields
}
fn bounds(&self) -> TaggedKeyBoundRange {
let mut start = self.head_fields(false);
if let Some(key_val) = self.key_end {
start.push(Field::UDesc(Width::U64, (key_val - 1) as u128));
}
let start = Bound::Included(TaggedKeyBound::prefix(SeriesRowKey::TAG, start));
let end = match self.key_start {
Some(key_val) => {
let mut end = self.head_fields(true);
end.push(Field::UDesc(Width::U64, key_val as u128));
end.push(Field::UDesc(Width::U64, 0));
Bound::Included(TaggedKeyBound::prefix(SeriesRowKey::TAG, end))
}
None => Bound::Excluded(TaggedKeyBound::prefix_end(
SeriesRowKey::TAG,
object_fields(ObjectId::from(self.storage)),
)),
};
TaggedKeyBoundRange {
start,
end,
}
}
}
#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct StorageSeriesKey {
pub variant_tag: Desc<Option<u8>>,
pub key: Desc<u64>,
pub sequence: Desc<u64>,
}
impl StorageSeriesKey {
pub fn new(variant_tag: Option<u8>, key: u64, sequence: u64) -> Self {
Self {
variant_tag: Desc(variant_tag),
key: Desc(key),
sequence: Desc(sequence),
}
}
pub fn variant_tag(self) -> Option<u8> {
self.variant_tag.0
}
pub fn key(self) -> u64 {
self.key.0
}
pub fn sequence(self) -> u64 {
self.sequence.0
}
pub fn to_sql_columns(self) -> SeriesKeyColumns {
SeriesKeyColumns {
variant_tag: series_variant_tag_to_sql(self.variant_tag()),
key: series_key_to_sql(self.key()),
sequence: series_sequence_to_sql(self.sequence()),
}
}
pub fn from_sql_columns(columns: SeriesKeyColumns) -> Self {
Self::new(
series_variant_tag_from_sql(columns.variant_tag),
series_key_from_sql(columns.key),
series_sequence_from_sql(columns.sequence),
)
}
pub fn with_storage(self, storage: StorageId) -> SeriesRowKey {
SeriesRowKey {
storage,
variant_tag: self.variant_tag(),
key: self.key(),
sequence: self.sequence(),
}
}
}
impl From<SeriesRowKey> for StorageSeriesKey {
fn from(key: SeriesRowKey) -> Self {
StorageSeriesKey::new(key.variant_tag, key.key, key.sequence)
}
}
impl HeapSize for StorageSeriesKey {
fn heap_size(&self) -> usize {
0
}
}
impl BoundedKey for StorageSeriesKey {
fn low() -> Self {
Self {
variant_tag: lowest_variant_tag(),
key: <Desc<u64> as BoundedKey>::low(),
sequence: <Desc<u64> as BoundedKey>::low(),
}
}
}
impl DenseKey for StorageSeriesKey {
fn successor(&self) -> Option<Self> {
if let Some(sequence) = self.sequence.successor() {
return Some(Self {
variant_tag: self.variant_tag,
key: self.key,
sequence,
});
}
if let Some(key) = self.key.successor() {
return Some(Self {
variant_tag: self.variant_tag,
key,
sequence: <Desc<u64> as BoundedKey>::low(),
});
}
Some(Self {
variant_tag: next_variant_tag(self.variant_tag)?,
key: <Desc<u64> as BoundedKey>::low(),
sequence: <Desc<u64> as BoundedKey>::low(),
})
}
}
fn lowest_variant_tag() -> Desc<Option<u8>> {
Desc(Some(u8::MAX))
}
fn next_variant_tag(tag: Desc<Option<u8>>) -> Option<Desc<Option<u8>>> {
match tag.0 {
Some(0) => Some(Desc(None)),
Some(value) => Some(Desc(Some(value - 1))),
None => None,
}
}
#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct SeriesKeyColumns {
pub variant_tag: i64,
pub key: i64,
pub sequence: i64,
}
#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct PartitionedSeriesKeyColumns {
pub partition_hi: i64,
pub partition_lo: i64,
pub variant_tag: i64,
pub key: i64,
pub sequence: i64,
}
const SERIES_VARIANT_TAG_NONE_SQL: i64 = (((!0u8) as i64) << 8) | ((!0u8) as i64);
pub fn series_variant_tag_to_sql(variant_tag: Option<u8>) -> i64 {
match variant_tag {
Some(tag) => (((!1u8) as i64) << 8) | ((!tag) as i64),
None => SERIES_VARIANT_TAG_NONE_SQL,
}
}
pub fn series_variant_tag_from_sql(value: i64) -> Option<u8> {
if value == SERIES_VARIANT_TAG_NONE_SQL {
return None;
}
Some(!(value as u8))
}
pub fn series_key_to_sql(key: u64) -> i64 {
desc_u64_to_sql(key)
}
pub fn series_key_from_sql(value: i64) -> u64 {
desc_u64_from_sql(value)
}
pub fn series_sequence_to_sql(sequence: u64) -> i64 {
desc_u64_to_sql(sequence)
}
pub fn series_sequence_from_sql(value: i64) -> u64 {
desc_u64_from_sql(value)
}
pub fn series_partition_half_to_sql(half: u64) -> i64 {
desc_u64_to_sql(half)
}
pub fn series_partition_half_from_sql(value: i64) -> u64 {
desc_u64_from_sql(value)
}
fn desc_u64_to_sql(value: u64) -> i64 {
((!value) ^ (1u64 << 63)) as i64
}
fn desc_u64_from_sql(value: i64) -> u64 {
!((value as u64) ^ (1u64 << 63))
}
#[cfg(test)]
mod row_key_range_tests {
use std::ops::RangeBounds;
use super::*;
use crate::{
interface::catalog::id::{SeriesId, ViewId},
key::{KeyRangeCodec, row::RowKeyRange},
};
#[test]
fn test_encode_decode_without_tag() {
let key = SeriesRowKey {
storage: StorageId::Series(SeriesId(42)),
variant_tag: None,
key: 1706745600000,
sequence: 1,
};
let encoded = key.encode();
let decoded = SeriesRowKey::decode(&encoded).unwrap();
assert_eq!(decoded.storage, StorageId::Series(SeriesId(42)));
assert_eq!(decoded.variant_tag, None);
assert_eq!(decoded.key, 1706745600000);
assert_eq!(decoded.sequence, 1);
}
#[test]
fn test_encode_decode_with_tag() {
let key = SeriesRowKey {
storage: StorageId::Series(SeriesId(42)),
variant_tag: Some(3),
key: 1706745600000,
sequence: 5,
};
let encoded = key.encode();
let decoded = SeriesRowKey::decode(&encoded).unwrap();
assert_eq!(decoded.storage, StorageId::Series(SeriesId(42)));
assert_eq!(decoded.variant_tag, Some(3));
assert_eq!(decoded.key, 1706745600000);
assert_eq!(decoded.sequence, 5);
}
#[test]
fn test_a_view_storage_round_trips() {
let key = SeriesRowKey {
storage: StorageId::View(ViewId(42)),
variant_tag: Some(7),
key: 900,
sequence: 2,
};
let decoded = SeriesRowKey::decode(&key.encode()).unwrap();
assert_eq!(decoded, key);
}
#[test]
fn test_full_scan_range_still_names_its_storage() {
let range = SeriesRowKeyRange::full_scan(StorageId::Series(SeriesId(42)), None).encode();
let (start, end) = SeriesRowKeyRange::decode(&range);
assert_eq!(start, Some(StorageId::Series(SeriesId(42))));
assert_eq!(
end,
Some(StorageId::Series(SeriesId(41))),
"the exclusive end brackets the next series id down"
);
}
#[test]
fn test_untagged_full_scan_covers_tagged_rows() {
let storage = StorageId::Series(SeriesId(9));
let range = SeriesRowKeyRange::full_scan(storage, None).encode();
let untagged = SeriesRowKey {
storage,
variant_tag: None,
key: 500,
sequence: 0,
}
.encode();
let tagged = SeriesRowKey {
storage,
variant_tag: Some(4),
key: 500,
sequence: 0,
}
.encode();
assert!(range.contains(&untagged));
assert!(range.contains(&tagged), "an untagged full scan must still see tagged rows");
}
#[test]
fn test_scan_range_brackets_the_rows_it_selects() {
let storage = StorageId::Series(SeriesId(1));
let range = SeriesRowKeyRange::scan_range(storage, None, Some(100), Some(200), None).encode();
let inside = SeriesRowKey {
storage,
variant_tag: None,
key: 150,
sequence: 1,
}
.encode();
let below = SeriesRowKey {
storage,
variant_tag: None,
key: 99,
sequence: 1,
}
.encode();
let above = SeriesRowKey {
storage,
variant_tag: None,
key: 201,
sequence: 1,
}
.encode();
assert!(range.contains(&inside));
assert!(!range.contains(&below));
assert!(!range.contains(&above));
}
#[test]
fn test_row_key_range_never_claims_a_series_range() {
let range = SeriesRowKeyRange::full_scan(StorageId::Series(SeriesId(42)), None).encode();
assert_eq!(RowKeyRange::decode(&range), (None, None));
}
#[test]
fn test_ordering_by_key() {
let key1 = SeriesRowKey {
storage: StorageId::Series(SeriesId(1)),
variant_tag: None,
key: 100,
sequence: 0,
};
let key2 = SeriesRowKey {
storage: StorageId::Series(SeriesId(1)),
variant_tag: None,
key: 200,
sequence: 0,
};
let e1 = key1.encode();
let e2 = key2.encode();
assert!(e1 > e2, "key descending ordering not preserved");
}
#[test]
fn test_ordering_by_sequence() {
let key1 = SeriesRowKey {
storage: StorageId::Series(SeriesId(1)),
variant_tag: None,
key: 100,
sequence: 1,
};
let key2 = SeriesRowKey {
storage: StorageId::Series(SeriesId(1)),
variant_tag: None,
key: 100,
sequence: 2,
};
let e1 = key1.encode();
let e2 = key2.encode();
assert!(e1 > e2, "sequence descending ordering not preserved");
}
#[test]
fn test_half_bounded_untagged_range_excludes_tagged_rows() {
let range = SeriesRowKeyRange::scan_range(StorageId::Series(SeriesId(7)), None, Some(100), None, None)
.encode();
let untagged = SeriesRowKey {
storage: StorageId::Series(SeriesId(7)),
variant_tag: None,
key: 500,
sequence: 1,
}
.encode();
let tagged = SeriesRowKey {
storage: StorageId::Series(SeriesId(7)),
variant_tag: Some(3),
key: 500,
sequence: 1,
}
.encode();
assert!(range.contains(&untagged), "an untagged row above the lower bound must stay in range");
assert!(!range.contains(&tagged), "a tagged row must never leak into an untagged key-bounded range");
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = PartitionedSeriesRow)]
pub struct PartitionedSeriesRowKey {
pub storage: StorageId,
pub partition: Partition,
pub variant_tag: Option<u8>,
pub key: u64,
pub sequence: u64,
}
impl PartitionedSeriesRowKey {
pub fn new(
storage: impl Into<StorageId>,
partition: Partition,
variant_tag: Option<u8>,
key: u64,
sequence: u64,
) -> Self {
Self {
storage: storage.into(),
partition,
variant_tag,
key,
sequence,
}
}
pub fn encoded(
storage: impl Into<StorageId>,
partition: Partition,
variant_tag: Option<u8>,
key: u64,
sequence: u64,
) -> EncodedKey {
Self::new(storage, partition, variant_tag, key, sequence).encode()
}
pub fn storage_of(key: &EncodedKey) -> Option<StorageId> {
let mut de = KeyDeserializer::from_bytes(key.as_slice());
let kind: KeyTag = de.read_u8().ok()?.try_into().ok()?;
if kind != Self::TAG {
return None;
}
StorageId::from_object(de.read_object_id().ok()?)
}
}
#[derive(Debug, Clone)]
pub struct PartitionedSeriesRowKeyRange {
pub storage: StorageId,
pub partition: Partition,
pub variant_tag: Option<u8>,
pub key_start: Option<u64>,
pub key_end: Option<u64>,
}
impl PartitionedSeriesRowKeyRange {
pub fn storage_start(storage: impl Into<StorageId>) -> EncodedKey {
TaggedKeyBound::prefix(PartitionedSeriesRowKey::TAG, object_fields(ObjectId::from(storage.into())))
.encode()
}
pub fn storage_end(storage: impl Into<StorageId>) -> EncodedKey {
TaggedKeyBound::prefix(
PartitionedSeriesRowKey::TAG,
object_fields(ObjectId::from(storage.into()).prev()),
)
.encode()
}
pub fn full_scan(storage: impl Into<StorageId>) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(PartitionedSeriesRowKey::TAG, object_fields(ObjectId::from(storage.into())))
}
pub fn full_scan_range(storage: impl Into<StorageId>, last: Option<&TaggedKey>) -> TaggedKeyBoundRange {
Self::full_scan(storage).resume_after(last)
}
pub fn partition_range(storage: impl Into<StorageId>, partition: Partition) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
PartitionedSeriesRowKey::TAG,
Self::partition_fields(storage.into(), partition),
)
}
pub fn partition_scan_range(
storage: impl Into<StorageId>,
partition: Partition,
last: Option<&TaggedKey>,
) -> TaggedKeyBoundRange {
Self::partition_range(storage, partition).resume_after(last)
}
pub fn scan_range(
storage: impl Into<StorageId>,
partition: Partition,
variant_tag: Option<u8>,
key_start: Option<u64>,
key_end: Option<u64>,
last: Option<&TaggedKey>,
) -> TaggedKeyBoundRange {
if matches!(key_end, Some(0)) {
return TaggedKeyBoundRange::empty(PartitionedSeriesRowKey::TAG);
}
PartitionedSeriesRowKeyRange {
storage: storage.into(),
partition,
variant_tag,
key_start,
key_end,
}
.bounds()
.resume_after(last)
}
pub fn decode_storage(key: &EncodedKey) -> Option<StorageId> {
PartitionedSeriesRowKey::storage_of(key)
}
pub fn decode(range: &EncodedKeyRange) -> (Option<StorageId>, Option<StorageId>) {
let start = match &range.start {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
Bound::Unbounded => None,
};
let end = match &range.end {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
Bound::Unbounded => None,
};
(start, end)
}
fn partition_fields(storage: StorageId, partition: Partition) -> Vec<OwnedField> {
let mut fields = object_fields(ObjectId::from(storage)).to_vec();
fields.push(Field::UDesc(Width::U128, partition.0));
fields
}
fn head_fields(&self, tagged: bool) -> Vec<OwnedField> {
let mut fields = Self::partition_fields(self.storage, self.partition);
match self.variant_tag {
Some(tag) => {
fields.push(Field::UDesc(Width::U8, 1));
fields.push(Field::UDesc(Width::U8, tag as u128));
}
None if tagged || self.key_start.is_some() || self.key_end.is_some() => {
fields.push(Field::UDesc(Width::U8, 0));
fields.push(Field::UDesc(Width::U8, 0));
}
None => {}
}
fields
}
fn bounds(&self) -> TaggedKeyBoundRange {
let mut start = self.head_fields(false);
if let Some(key_val) = self.key_end {
start.push(Field::UDesc(Width::U64, (key_val - 1) as u128));
}
let start = Bound::Included(TaggedKeyBound::prefix(PartitionedSeriesRowKey::TAG, start));
let end = match self.key_start {
Some(key_val) => {
let mut end = self.head_fields(true);
end.push(Field::UDesc(Width::U64, key_val as u128));
end.push(Field::UDesc(Width::U64, 0));
Bound::Included(TaggedKeyBound::prefix(PartitionedSeriesRowKey::TAG, end))
}
None => Bound::Excluded(TaggedKeyBound::prefix_end(
PartitionedSeriesRowKey::TAG,
Self::partition_fields(self.storage, self.partition),
)),
};
TaggedKeyBoundRange {
start,
end,
}
}
}
#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct StoragePartitionedSeriesKey {
pub partition: Desc<Partition>,
pub variant_tag: Desc<Option<u8>>,
pub key: Desc<u64>,
pub sequence: Desc<u64>,
}
impl StoragePartitionedSeriesKey {
pub fn new(partition: Partition, variant_tag: Option<u8>, key: u64, sequence: u64) -> Self {
Self {
partition: Desc(partition),
variant_tag: Desc(variant_tag),
key: Desc(key),
sequence: Desc(sequence),
}
}
pub fn partition(self) -> Partition {
self.partition.0
}
pub fn variant_tag(self) -> Option<u8> {
self.variant_tag.0
}
pub fn key(self) -> u64 {
self.key.0
}
pub fn sequence(self) -> u64 {
self.sequence.0
}
pub fn partition_hi(self) -> u64 {
(self.partition().0 >> 64) as u64
}
pub fn partition_lo(self) -> u64 {
self.partition().0 as u64
}
pub fn from_halves(
partition_hi: u64,
partition_lo: u64,
variant_tag: Option<u8>,
key: u64,
sequence: u64,
) -> Self {
Self::new(Partition(((partition_hi as u128) << 64) | partition_lo as u128), variant_tag, key, sequence)
}
pub fn to_sql_columns(self) -> PartitionedSeriesKeyColumns {
PartitionedSeriesKeyColumns {
partition_hi: series_partition_half_to_sql(self.partition_hi()),
partition_lo: series_partition_half_to_sql(self.partition_lo()),
variant_tag: series_variant_tag_to_sql(self.variant_tag()),
key: series_key_to_sql(self.key()),
sequence: series_sequence_to_sql(self.sequence()),
}
}
pub fn from_sql_columns(columns: PartitionedSeriesKeyColumns) -> Self {
Self::from_halves(
series_partition_half_from_sql(columns.partition_hi),
series_partition_half_from_sql(columns.partition_lo),
series_variant_tag_from_sql(columns.variant_tag),
series_key_from_sql(columns.key),
series_sequence_from_sql(columns.sequence),
)
}
pub fn with_storage(self, storage: StorageId) -> PartitionedSeriesRowKey {
PartitionedSeriesRowKey {
storage,
partition: self.partition(),
variant_tag: self.variant_tag(),
key: self.key(),
sequence: self.sequence(),
}
}
}
impl From<PartitionedSeriesRowKey> for StoragePartitionedSeriesKey {
fn from(key: PartitionedSeriesRowKey) -> Self {
StoragePartitionedSeriesKey::new(key.partition, key.variant_tag, key.key, key.sequence)
}
}
impl HeapSize for StoragePartitionedSeriesKey {
fn heap_size(&self) -> usize {
0
}
}
impl BoundedKey for StoragePartitionedSeriesKey {
fn low() -> Self {
Self {
partition: <Desc<Partition> as BoundedKey>::low(),
variant_tag: lowest_variant_tag(),
key: <Desc<u64> as BoundedKey>::low(),
sequence: <Desc<u64> as BoundedKey>::low(),
}
}
}
impl DenseKey for StoragePartitionedSeriesKey {
fn successor(&self) -> Option<Self> {
if let Some(sequence) = self.sequence.successor() {
return Some(Self {
sequence,
..*self
});
}
if let Some(key) = self.key.successor() {
return Some(Self {
key,
sequence: <Desc<u64> as BoundedKey>::low(),
..*self
});
}
if let Some(variant_tag) = next_variant_tag(self.variant_tag) {
return Some(Self {
partition: self.partition,
variant_tag,
key: <Desc<u64> as BoundedKey>::low(),
sequence: <Desc<u64> as BoundedKey>::low(),
});
}
Some(Self {
partition: self.partition.successor()?,
variant_tag: lowest_variant_tag(),
key: <Desc<u64> as BoundedKey>::low(),
sequence: <Desc<u64> as BoundedKey>::low(),
})
}
}
#[cfg(test)]
mod partitioned_row_key_tests {
use std::ops::RangeBounds;
use reifydb_value::value::{Value, partition::Partition, row_number::RowNumber};
use super::*;
use crate::{
interface::catalog::id::{SeriesId, TableId, ViewId},
key::row::PartitionedRowKey,
};
fn part(v: &str) -> Partition {
Partition::of(&[Value::Utf8(v.to_string())])
}
#[test]
fn test_round_trip_without_tag() {
let key = PartitionedSeriesRowKey {
storage: StorageId::Series(SeriesId(3)),
partition: part("btc"),
variant_tag: None,
key: 1_700_000_000,
sequence: 9,
};
let decoded = PartitionedSeriesRowKey::decode(&key.encode()).unwrap();
assert_eq!(decoded, key);
}
#[test]
fn test_round_trip_with_tag() {
let key = PartitionedSeriesRowKey {
storage: StorageId::Series(SeriesId(3)),
partition: part("eth"),
variant_tag: Some(5),
key: 42,
sequence: 0,
};
let decoded = PartitionedSeriesRowKey::decode(&key.encode()).unwrap();
assert_eq!(decoded, key);
}
#[test]
fn test_a_view_storage_round_trips() {
let key = PartitionedSeriesRowKey {
storage: StorageId::View(ViewId(11)),
partition: part("us"),
variant_tag: Some(1),
key: 77,
sequence: 4,
};
let decoded = PartitionedSeriesRowKey::decode(&key.encode()).unwrap();
assert_eq!(decoded, key);
}
#[test]
fn test_storage_of() {
let key = PartitionedSeriesRowKey::encoded(StorageId::Series(SeriesId(42)), part("us"), None, 1, 0);
assert_eq!(PartitionedSeriesRowKey::storage_of(&key), Some(StorageId::Series(SeriesId(42))));
}
#[test]
fn test_ordering_by_key_is_descending() {
let storage = StorageId::Series(SeriesId(1));
let low = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 100, 0);
let high = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 200, 0);
assert!(low > high, "key descending ordering not preserved");
}
#[test]
fn test_untagged_full_scan_covers_tagged_rows() {
let storage = StorageId::Series(SeriesId(9));
let range = PartitionedSeriesRowKeyRange::full_scan(storage).encode();
let untagged = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 500, 0);
let tagged = PartitionedSeriesRowKey::encoded(storage, part("us"), Some(4), 500, 0);
assert!(range.contains(&untagged));
assert!(range.contains(&tagged), "an untagged full scan must still see tagged rows");
}
#[test]
fn test_scan_range_brackets_the_rows_it_selects() {
let storage = StorageId::Series(SeriesId(1));
let partition = part("us");
let range =
PartitionedSeriesRowKeyRange::scan_range(storage, partition, None, Some(100), Some(200), None)
.encode();
let inside = PartitionedSeriesRowKey::encoded(storage, partition, None, 150, 1);
let below = PartitionedSeriesRowKey::encoded(storage, partition, None, 99, 1);
let above = PartitionedSeriesRowKey::encoded(storage, partition, None, 201, 1);
assert!(range.contains(&inside));
assert!(!range.contains(&below));
assert!(!range.contains(&above));
}
#[test]
fn test_scan_range_never_crosses_into_another_partition() {
let storage = StorageId::Series(SeriesId(1));
let range =
PartitionedSeriesRowKeyRange::scan_range(storage, part("us"), None, Some(100), Some(200), None)
.encode();
let other = PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 150, 1);
assert!(!range.contains(&other), "an in-bounds key of another partition must stay outside");
}
#[test]
fn test_partition_range_covers_every_key_of_its_partition() {
let storage = StorageId::Series(SeriesId(1));
let range = PartitionedSeriesRowKeyRange::partition_range(storage, part("us")).encode();
let untagged = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 1, 0);
let tagged = PartitionedSeriesRowKey::encoded(storage, part("us"), Some(9), u64::MAX, u64::MAX);
let other = PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 1, 0);
assert!(range.contains(&untagged));
assert!(range.contains(&tagged));
assert!(!range.contains(&other));
}
#[test]
fn test_full_scan_range_resumes_across_partitions() {
let storage = StorageId::Series(SeriesId(1));
let cursor = TaggedKey::from(PartitionedSeriesRowKey::new(storage, part("us"), None, 200, 0));
let range = PartitionedSeriesRowKeyRange::full_scan_range(storage, Some(&cursor)).encode();
assert!(!range.contains(&cursor.encode()));
assert!(range.contains(&PartitionedSeriesRowKey::encoded(storage, part("us"), None, 100, 0)));
assert!(
range.contains(&PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 100, 0))
|| range.contains(&PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 300, 0)),
"a resumed object-wide scan must still be able to reach another partition"
);
}
#[test]
fn test_full_scan_range_without_a_cursor_covers_every_partition() {
let storage = StorageId::Series(SeriesId(1));
let range = PartitionedSeriesRowKeyRange::full_scan_range(storage, None).encode();
assert!(range.contains(&PartitionedSeriesRowKey::encoded(storage, part("us"), None, 1, 0)));
assert!(range.contains(&PartitionedSeriesRowKey::encoded(storage, part("eu"), Some(3), 9, 9)));
assert!(!range.contains(&PartitionedSeriesRowKey::encoded(
StorageId::Series(SeriesId(2)),
part("us"),
None,
1,
0
)));
}
#[test]
fn test_partition_scan_range_resumes_after_the_cursor() {
let storage = StorageId::Series(SeriesId(1));
let partition = part("us");
let cursor = TaggedKey::from(PartitionedSeriesRowKey::new(storage, partition, None, 200, 0));
let range =
PartitionedSeriesRowKeyRange::partition_scan_range(storage, partition, Some(&cursor)).encode();
let next = PartitionedSeriesRowKey::encoded(storage, partition, None, 100, 0);
assert!(!range.contains(&cursor.encode()));
assert!(range.contains(&next));
}
#[test]
fn test_the_two_partitioned_kinds_do_not_share_a_keyspace() {
let series =
PartitionedSeriesRowKey::encoded(StorageId::Series(SeriesId(1)), part("us"), Some(2), 7, 9);
let row = PartitionedRowKey::encoded(StorageId::Table(TableId(1)), part("us"), RowNumber(9));
assert_ne!(series.as_slice()[0], row.as_slice()[0]);
assert!(PartitionedRowKey::decode(&series).is_none());
assert!(PartitionedSeriesRowKey::decode(&row).is_none());
}
#[test]
fn test_half_bounded_untagged_range_excludes_tagged_rows() {
let range = PartitionedSeriesRowKeyRange::scan_range(
StorageId::Series(SeriesId(7)),
part("us"),
None,
Some(100),
None,
None,
)
.encode();
let untagged = PartitionedSeriesRowKey {
storage: StorageId::Series(SeriesId(7)),
partition: part("us"),
variant_tag: None,
key: 500,
sequence: 1,
}
.encode();
let tagged = PartitionedSeriesRowKey {
storage: StorageId::Series(SeriesId(7)),
partition: part("us"),
variant_tag: Some(3),
key: 500,
sequence: 1,
}
.encode();
assert!(range.contains(&untagged), "an untagged row above the lower bound must stay in range");
assert!(!range.contains(&tagged), "a tagged row must never leak into an untagged key-bounded range");
}
}
#[cfg(test)]
mod storage_series_key_tests {
use reifydb_value::value::partition::Partition;
use super::{PartitionedSeriesRowKey, SeriesRowKey, StorageId, StoragePartitionedSeriesKey, StorageSeriesKey};
use crate::{
interface::catalog::id::SeriesId,
key::typed::{BoundedKey, DenseKey},
};
const STORAGE: StorageId = StorageId::Series(SeriesId(7));
fn series(variant_tag: Option<u8>, key: u64, sequence: u64) -> StorageSeriesKey {
StorageSeriesKey::new(variant_tag, key, sequence)
}
#[test]
fn storage_key_order_matches_the_encoded_byte_order() {
let mut keys = vec![
series(None, 5, 1),
series(Some(0), 5, 1),
series(Some(9), 5, 1),
series(Some(9), 5, 2),
series(Some(9), 6, 1),
];
let mut encoded: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
keys.sort();
encoded.sort();
let reordered: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
assert_eq!(reordered, encoded);
}
#[test]
fn a_tagged_row_sorts_before_an_untagged_one() {
assert!(series(Some(0), 1, 1) < series(None, 1, 1));
assert!(series(Some(255), 1, 1) < series(Some(0), 1, 1));
}
#[test]
fn low_is_the_first_key_the_encoder_can_produce() {
let low = <StorageSeriesKey as BoundedKey>::low();
assert_eq!(low, series(Some(u8::MAX), u64::MAX, u64::MAX));
for other in [series(None, 0, 0), series(Some(0), 0, 0), series(Some(200), 7, 3)] {
assert!(low < other, "low must not sort above a real key");
}
}
#[test]
fn successor_walks_the_sequence_then_the_key_then_the_tag() {
assert_eq!(series(Some(9), 5, 2).successor(), Some(series(Some(9), 5, 1)));
assert_eq!(series(Some(9), 5, 0).successor(), Some(series(Some(9), 4, u64::MAX)));
assert_eq!(series(Some(9), 0, 0).successor(), Some(series(Some(8), u64::MAX, u64::MAX)));
assert_eq!(series(Some(0), 0, 0).successor(), Some(series(None, u64::MAX, u64::MAX)));
assert_eq!(series(None, 0, 0).successor(), None);
}
#[test]
fn every_successor_is_the_immediate_next_storage_key() {
let mut walk = vec![series(Some(1), 1, 1)];
for _ in 0..4 {
walk.push(walk.last().unwrap().successor().unwrap());
}
assert!(walk.windows(2).all(|pair| pair[0] < pair[1]), "the walk must ascend: {walk:?}");
}
#[test]
fn partitioned_storage_key_order_matches_the_encoded_byte_order() {
let mut keys = vec![
StoragePartitionedSeriesKey::new(Partition(2), None, 5, 1),
StoragePartitionedSeriesKey::new(Partition(2), Some(3), 5, 1),
StoragePartitionedSeriesKey::new(Partition(1), Some(3), 5, 1),
StoragePartitionedSeriesKey::new(Partition(2), Some(3), 5, 2),
];
let mut encoded: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
keys.sort();
encoded.sort();
let reordered: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
assert_eq!(reordered, encoded);
}
#[test]
fn partitioned_successor_carries_into_the_partition() {
let last_of_partition = StoragePartitionedSeriesKey::new(Partition(5), None, 0, 0);
assert_eq!(
last_of_partition.successor(),
Some(StoragePartitionedSeriesKey::new(Partition(4), Some(u8::MAX), u64::MAX, u64::MAX))
);
assert_eq!(StoragePartitionedSeriesKey::new(Partition(0), None, 0, 0).successor(), None);
}
#[test]
fn a_storage_key_round_trips_through_its_row_key() {
let key = SeriesRowKey {
storage: STORAGE,
variant_tag: Some(4),
key: 900,
sequence: 12,
};
assert_eq!(StorageSeriesKey::from(key.clone()).with_storage(STORAGE), key);
let partitioned = PartitionedSeriesRowKey::new(STORAGE, Partition(3), None, 900, 12);
assert_eq!(StoragePartitionedSeriesKey::from(partitioned.clone()).with_storage(STORAGE), partitioned);
}
}