use std::{borrow::Cow, cmp::Ordering, collections::Bound};
use reifydb_codec::{
key::{
ByteSink,
deserializer::KeyDeserializer,
encoded::{EncodedKey, EncodedKeyRange},
serializer::KeySerializer,
},
row::shape::fingerprint::RowShapeFingerprint,
};
use reifydb_macro::KeyCodec;
use reifydb_value::value::{partition::Partition, row_number::RowNumber};
use serde::{Deserialize, Serialize};
use smallvec::{SmallVec, smallvec};
use super::{KeyRangeCodec, KeyTag};
use crate::{
interface::catalog::{object::ObjectId, storage::StorageId},
key::{
any::{Field, KeyFields, RawEncoding, TaggedKey, Width},
bound::{TaggedKeyBoundRange, object_fields},
catalog::{KeyDeserializerCatalogExt, KeySerializerCatalogExt},
sort_run::SortRun,
typed::{
BoundedKey, DenseKey,
direction::{Asc, Desc},
},
},
metrics::heap::HeapSize,
};
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = Row)]
pub struct RowKey {
pub storage: StorageId,
pub row: RowNumber,
}
#[derive(Debug, Clone, PartialEq)]
pub struct RowKeyRange {
pub storage: StorageId,
}
impl RowKeyRange {
fn decode_key(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 != Self::TAG {
return None;
}
let storage = StorageId::from_object(de.read_object_id().ok()?)?;
Some(RowKeyRange {
storage,
})
}
pub fn storage_scan(storage: StorageId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, object_fields(ObjectId::from(storage)))
}
pub fn scan_range(storage: StorageId, last: Option<&TaggedKey>) -> TaggedKeyBoundRange {
Self::storage_scan(storage).resume_after(last)
}
pub fn scan_range_rev(storage: StorageId, last: Option<&TaggedKey>) -> TaggedKeyBoundRange {
Self::storage_scan(storage).resume_before(last)
}
}
impl KeyRangeCodec for RowKeyRange {
const TAG: KeyTag = KeyTag::Row;
fn start(&self) -> Option<EncodedKey> {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(Self::TAG as u8).extend_object_id(self.storage);
Some(serializer.to_encoded_key())
}
fn end(&self) -> Option<EncodedKey> {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(Self::TAG as u8).extend_object_id(ObjectId::from(self.storage).prev());
Some(serializer.to_encoded_key())
}
fn decode(range: &EncodedKeyRange) -> (Option<Self>, Option<Self>)
where
Self: Sized,
{
let start_key = match &range.start {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
Bound::Unbounded => None,
};
let end_key = match &range.end {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
Bound::Unbounded => None,
};
(start_key, end_key)
}
}
impl RowKey {
pub fn new(storage: impl Into<StorageId>, row: impl Into<RowNumber>) -> Self {
Self {
storage: storage.into(),
row: row.into(),
}
}
pub fn encoded(storage: impl Into<StorageId>, row: impl Into<RowNumber>) -> EncodedKey {
Self {
storage: storage.into(),
row: row.into(),
}
.encode()
}
pub fn full_scan(storage: impl Into<StorageId>) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, object_fields(ObjectId::from(storage.into())))
}
pub fn storage_start(storage: impl Into<StorageId>) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(RowKey::TAG as u8).extend_object_id(storage.into());
serializer.to_encoded_key()
}
pub fn storage_end(storage: impl Into<StorageId>) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(RowKey::TAG as u8).extend_object_id(ObjectId::from(storage.into()).prev());
serializer.to_encoded_key()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct SortedViewRowKey {
pub storage: StorageId,
pub run: SortRun,
pub row: Asc<RowNumber>,
}
impl SortedViewRowKey {
pub const TAG: KeyTag = KeyTag::SortedViewRow;
pub fn encode(&self) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(20 + self.run.len());
serializer.extend_u8(Self::TAG as u8).extend_object_id(self.storage);
extend_sort_run(&mut serializer, &self.run);
serializer.extend_raw(&self.row.0.0.to_be_bytes());
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 != Self::TAG {
return None;
}
let storage = StorageId::from_object(de.read_object_id().ok()?)?;
let run = read_sort_run(&mut de)?;
let row = read_row_tail(&mut de)?;
if !de.is_empty() {
return None;
}
Some(Self {
storage,
run,
row,
})
}
}
impl Ord for SortedViewRowKey {
fn cmp(&self, other: &Self) -> Ordering {
storage_order(self.storage)
.cmp(&storage_order(other.storage))
.then_with(|| self.run.cmp(&other.run))
.then_with(|| self.row.cmp(&other.row))
}
}
impl PartialOrd for SortedViewRowKey {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl SortedViewRowKey {
pub fn new(storage: impl Into<StorageId>, run: SortRun, row: RowNumber) -> Self {
Self {
storage: storage.into(),
run,
row: Asc(row),
}
}
pub fn encoded(storage: impl Into<StorageId>, run: SortRun, row: RowNumber) -> EncodedKey {
Self::new(storage, run, row).encode()
}
pub fn storage_start(storage: impl Into<StorageId>) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(Self::TAG as u8).extend_object_id(storage.into());
serializer.to_encoded_key()
}
pub fn storage_scan(storage: impl Into<StorageId>) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, object_fields(ObjectId::from(storage.into())))
}
pub fn scan_range(storage: impl Into<StorageId>, last: Option<&TaggedKey>) -> TaggedKeyBoundRange {
Self::storage_scan(storage).resume_after(last)
}
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()?)
}
pub fn row_of(key: &EncodedKey) -> Option<RowNumber> {
let bytes = key.as_slice();
if KeyTag::of(bytes)? != Self::TAG {
return None;
}
let tail = bytes.len().checked_sub(8)?;
Some(RowNumber(u64::from_be_bytes(bytes[tail..].try_into().ok()?)))
}
pub fn range_storage_of(range: &EncodedKeyRange) -> Option<StorageId> {
let start = bound_storage_of(&range.start, Self::TAG)?;
bound_storage_of(&range.end, Self::TAG)?;
Some(start)
}
}
const SORT_RUN_MARKER: u8 = 0x00;
const SORT_RUN_ZERO: u8 = 0xff;
const SORT_RUN_END: u8 = 0x00;
pub(crate) fn encode_sort_run<B: ByteSink>(run: &[u8], out: &mut B) {
for &byte in run {
if byte == SORT_RUN_MARKER {
out.extend_from_slice(&[SORT_RUN_MARKER, SORT_RUN_ZERO]);
} else {
out.push(byte);
}
}
out.extend_from_slice(&[SORT_RUN_MARKER, SORT_RUN_END]);
}
fn extend_sort_run(serializer: &mut KeySerializer, run: &SortRun) {
encode_sort_run(run.as_slice(), serializer);
}
fn read_sort_run(de: &mut KeyDeserializer) -> Option<SortRun> {
let mut run: Vec<u8> = Vec::new();
loop {
let byte = de.read_raw(1).ok()?[0];
if byte != SORT_RUN_MARKER {
run.push(byte);
continue;
}
match de.read_raw(1).ok()?[0] {
SORT_RUN_END => return Some(SortRun::new(run)),
SORT_RUN_ZERO => run.push(0x00),
_ => return None,
}
}
}
fn read_row_tail(de: &mut KeyDeserializer) -> Option<Asc<RowNumber>> {
let bytes: [u8; 8] = de.read_raw(8).ok()?.try_into().ok()?;
Some(Asc(RowNumber(u64::from_be_bytes(bytes))))
}
fn storage_order(storage: StorageId) -> (u8, Desc<u64>) {
(ObjectId::from(storage).type_tag(), Desc(storage.as_u64()))
}
fn bound_storage_of(bound: &Bound<EncodedKey>, kind: KeyTag) -> Option<StorageId> {
let key = match bound {
Bound::Included(key) | Bound::Excluded(key) => key,
Bound::Unbounded => return None,
};
let mut de = KeyDeserializer::from_bytes(key.as_slice());
if KeyTag::try_from(de.read_u8().ok()?).ok()? != kind {
return None;
}
StorageId::from_object(de.read_object_id().ok()?)
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct PartitionedSortedViewRowKey {
pub storage: StorageId,
pub partition: Partition,
pub run: SortRun,
pub row: Asc<RowNumber>,
}
impl PartitionedSortedViewRowKey {
pub const TAG: KeyTag = KeyTag::PartitionedSortedViewRow;
pub fn encode(&self) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(36 + self.run.len());
serializer.extend_u8(Self::TAG as u8).extend_object_id(self.storage).extend_u128(self.partition.0);
extend_sort_run(&mut serializer, &self.run);
serializer.extend_raw(&self.row.0.0.to_be_bytes());
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 != Self::TAG {
return None;
}
let storage = StorageId::from_object(de.read_object_id().ok()?)?;
let partition = Partition(de.read_u128().ok()?);
let run = read_sort_run(&mut de)?;
let row = read_row_tail(&mut de)?;
if !de.is_empty() {
return None;
}
Some(Self {
storage,
partition,
run,
row,
})
}
}
impl Ord for PartitionedSortedViewRowKey {
fn cmp(&self, other: &Self) -> Ordering {
storage_order(self.storage)
.cmp(&storage_order(other.storage))
.then_with(|| Desc(self.partition).cmp(&Desc(other.partition)))
.then_with(|| self.run.cmp(&other.run))
.then_with(|| self.row.cmp(&other.row))
}
}
impl PartialOrd for PartitionedSortedViewRowKey {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl PartitionedSortedViewRowKey {
pub fn new(storage: impl Into<StorageId>, partition: Partition, run: SortRun, row: RowNumber) -> Self {
Self {
storage: storage.into(),
partition,
run,
row: Asc(row),
}
}
pub fn encoded(
storage: impl Into<StorageId>,
partition: Partition,
run: SortRun,
row: RowNumber,
) -> EncodedKey {
Self::new(storage, partition, run, row).encode()
}
pub fn storage_start(storage: impl Into<StorageId>) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(Self::TAG as u8).extend_object_id(storage.into());
serializer.to_encoded_key()
}
pub fn storage_scan(storage: impl Into<StorageId>) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, object_fields(ObjectId::from(storage.into())))
}
pub fn scan_range(storage: impl Into<StorageId>, last: Option<&TaggedKey>) -> TaggedKeyBoundRange {
Self::storage_scan(storage).resume_after(last)
}
pub fn partition_range(storage: impl Into<StorageId>, partition: Partition) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
Self::TAG,
object_fields(ObjectId::from(storage.into()))
.into_iter()
.chain([Field::UDesc(Width::U128, partition.0)])
.collect::<Vec<_>>(),
)
}
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 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()?)
}
pub fn row_of(key: &EncodedKey) -> Option<RowNumber> {
let bytes = key.as_slice();
if KeyTag::of(bytes)? != Self::TAG {
return None;
}
let tail = bytes.len().checked_sub(8)?;
Some(RowNumber(u64::from_be_bytes(bytes[tail..].try_into().ok()?)))
}
pub fn range_storage_of(range: &EncodedKeyRange) -> Option<StorageId> {
let start = bound_storage_of(&range.start, Self::TAG)?;
bound_storage_of(&range.end, Self::TAG)?;
Some(start)
}
}
#[cfg(test)]
mod sorted_view_row_key_tests {
use std::ops::RangeBounds;
use reifydb_codec::key::encoded::EncodedKey;
use reifydb_value::value::{Value, partition::Partition, row_number::RowNumber};
use super::{PartitionedSortedViewRowKey, RowKey, SortedViewRowKey, TaggedKey};
use crate::{interface::catalog::storage::StorageId, key::sort_run::SortRun};
fn part(v: &str) -> Partition {
Partition::of(&[Value::Utf8(v.to_string())])
}
fn sorted_view(storage: StorageId, sort: &[u8], row: RowNumber) -> EncodedKey {
SortedViewRowKey::encoded(storage, SortRun::new(sort), row)
}
fn sorted_view_cursor(storage: StorageId, sort: &[u8], row: RowNumber) -> TaggedKey {
TaggedKey::from(SortedViewRowKey::new(storage, SortRun::new(sort), row))
}
fn partitioned(storage: StorageId, partition: Partition, sort: &[u8], row: RowNumber) -> EncodedKey {
PartitionedSortedViewRowKey::encoded(storage, partition, SortRun::new(sort), row)
}
#[test]
fn test_row_comes_from_the_tail_not_the_sort_prefix() {
let storage = StorageId::view(3);
let a = sorted_view(storage, &[0xAA; 8], RowNumber(7));
let b = sorted_view(storage, &[0xBB; 24], RowNumber(9));
assert_eq!(SortedViewRowKey::row_of(&a), Some(RowNumber(7)));
assert_eq!(SortedViewRowKey::row_of(&b), Some(RowNumber(9)));
assert_eq!(SortedViewRowKey::storage_of(&a), Some(storage));
}
#[test]
fn test_row_of_rejects_a_plain_row_key() {
let plain = RowKey::encoded(StorageId::view(3), RowNumber(7));
assert_eq!(SortedViewRowKey::row_of(&plain), None);
assert_eq!(PartitionedSortedViewRowKey::row_of(&plain), None);
}
#[test]
fn test_scan_range_covers_its_storage_and_nothing_else() {
let storage = StorageId::view(3);
let range = SortedViewRowKey::scan_range(storage, None).encode();
assert!(range.contains(&sorted_view(storage, &[0x00; 8], RowNumber(1))));
assert!(range.contains(&sorted_view(storage, &[0xFF; 8], RowNumber(u64::MAX))));
assert!(!range.contains(&sorted_view(StorageId::view(4), &[0x00; 8], RowNumber(1))));
}
#[test]
fn test_scan_range_resumes_strictly_after_the_last_key() {
let storage = StorageId::view(3);
let last = sorted_view_cursor(storage, &[0x40; 8], RowNumber(5));
let range = SortedViewRowKey::scan_range(storage, Some(&last)).encode();
assert!(!range.contains(&last.encode()));
assert!(range.contains(&sorted_view(storage, &[0x41; 8], RowNumber(1))));
}
#[test]
fn test_sort_prefix_orders_the_keyspace() {
let storage = StorageId::view(3);
let early_sort_late_row = sorted_view(storage, &[0x10; 8], RowNumber(999));
let late_sort_early_row = sorted_view(storage, &[0x20; 8], RowNumber(1));
assert!(early_sort_late_row < late_sort_early_row);
}
#[test]
fn test_partition_range_contains_only_its_partition() {
let storage = StorageId::view(3);
let range = PartitionedSortedViewRowKey::partition_range(storage, part("us")).encode();
assert!(range.contains(&partitioned(storage, part("us"), &[0x10; 8], RowNumber(1))));
assert!(!range.contains(&partitioned(storage, part("eu"), &[0x10; 8], RowNumber(1))));
}
#[test]
fn test_partitioned_row_comes_from_the_tail() {
let storage = StorageId::view(3);
let key = partitioned(storage, part("us"), &[0xAA; 8], RowNumber(42));
assert_eq!(PartitionedSortedViewRowKey::row_of(&key), Some(RowNumber(42)));
assert_eq!(PartitionedSortedViewRowKey::storage_of(&key), Some(storage));
assert_eq!(SortedViewRowKey::row_of(&key), None);
}
}
#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct StorageRowKey(pub Desc<RowNumber>);
impl From<RowKey> for StorageRowKey {
fn from(key: RowKey) -> Self {
StorageRowKey::new(key.row)
}
}
impl StorageRowKey {
pub fn new(row: RowNumber) -> Self {
StorageRowKey(Desc(row))
}
pub fn row(self) -> RowNumber {
self.0.0
}
pub fn with_storage(self, storage: StorageId) -> RowKey {
RowKey {
storage,
row: self.row(),
}
}
}
impl HeapSize for StorageRowKey {
fn heap_size(&self) -> usize {
0
}
}
impl BoundedKey for StorageRowKey {
fn low() -> Self {
StorageRowKey(<Desc<RowNumber> as BoundedKey>::low())
}
}
impl DenseKey for StorageRowKey {
fn successor(&self) -> Option<Self> {
self.0.successor().map(StorageRowKey)
}
}
#[cfg(test)]
pub mod row_key_tests {
use reifydb_value::value::row_number::RowNumber;
use super::{RowKey, StorageRowKey};
use crate::{
interface::catalog::storage::StorageId,
key::typed::{BoundedKey, DenseKey},
};
#[test]
fn test_encode_decode() {
let key = RowKey {
storage: StorageId::table(0xABCD),
row: RowNumber(0x123456789ABCDEF0),
};
let encoded = key.encode();
let expected: Vec<u8> = vec![
0xFC, 0x01, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x54, 0x32, 0xED, 0xCB, 0xA9, 0x87, 0x65, 0x43,
0x21, 0x0F,
];
assert_eq!(encoded.as_slice(), expected);
let key = RowKey::decode(&encoded).unwrap();
assert_eq!(key.storage, StorageId::table(0xABCD));
assert_eq!(key.row, 0x123456789ABCDEF0);
}
#[test]
fn test_encode_decode_view() {
let key = RowKey {
storage: StorageId::view(0xABCD),
row: RowNumber(0x123456789ABCDEF0),
};
let encoded = key.encode();
let expected: Vec<u8> = vec![
0xFC, 0x02, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x54, 0x32, 0xED, 0xCB, 0xA9, 0x87, 0x65, 0x43,
0x21, 0x0F,
];
assert_eq!(encoded.as_slice(), expected);
let key = RowKey::decode(&encoded).unwrap();
assert_eq!(key.storage, StorageId::view(0xABCD));
assert_eq!(key.row, 0x123456789ABCDEF0);
}
#[test]
fn test_order_preserving() {
let key1 = RowKey {
storage: StorageId::table(1),
row: RowNumber(100),
};
let key2 = RowKey {
storage: StorageId::table(1),
row: RowNumber(200),
};
let key3 = RowKey {
storage: StorageId::table(2),
row: RowNumber(1),
};
let encoded1 = key1.encode();
let encoded2 = key2.encode();
let encoded3 = key3.encode();
assert!(encoded3 < encoded2, "ordering not preserved");
assert!(encoded2 < encoded1, "ordering not preserved");
}
#[test]
fn test_row_ident_roundtrip() {
let key = RowKey {
storage: StorageId::table(7),
row: RowNumber(42),
};
let ident: StorageRowKey = key.clone().into();
let restored = ident.with_storage(key.storage);
assert_eq!(restored, key);
}
#[test]
fn test_row_ident_ordering_matches_the_encoded_key() {
let storage = StorageId::table(1);
let one = StorageRowKey::new(RowNumber(1));
let two = StorageRowKey::new(RowNumber(2));
assert!(two < one);
assert_eq!(
two < one,
RowKey::encoded(storage, RowNumber(2)).as_slice()
< RowKey::encoded(storage, RowNumber(1)).as_slice()
);
}
#[test]
fn test_row_ident_low_is_the_greatest_row() {
assert_eq!(<StorageRowKey as BoundedKey>::low(), StorageRowKey::new(RowNumber(u64::MAX)));
}
#[test]
fn test_row_ident_successor_is_the_next_key_in_scan_order() {
let ident = StorageRowKey::new(RowNumber(5));
let next = ident.successor().unwrap();
assert_eq!(next, StorageRowKey::new(RowNumber(4)));
assert!(next > ident);
assert!(StorageRowKey::new(RowNumber(3)) > next);
}
#[test]
fn test_row_ident_successor_runs_out_at_row_zero() {
assert_eq!(StorageRowKey::new(RowNumber(0)).successor(), None);
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = RowSequence)]
pub struct RowSequenceKey {
pub storage: StorageId,
}
impl RowSequenceKey {
pub fn new(storage: StorageId) -> Self {
Self {
storage,
}
}
pub fn encoded(storage: impl Into<StorageId>) -> EncodedKey {
Self::new(storage.into()).encode()
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[cfg(test)]
pub mod row_sequence_key_tests {
use super::RowSequenceKey;
use crate::interface::catalog::storage::StorageId;
#[test]
fn test_encode_decode() {
let key = RowSequenceKey {
storage: StorageId::table(0xABCD),
};
let encoded = key.encode();
let expected = vec![0xF7, 0x01, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x54, 0x32];
assert_eq!(encoded.as_slice(), expected);
let key = RowSequenceKey::decode(&encoded).unwrap();
assert_eq!(key.storage, StorageId::table(0xABCD));
}
#[test]
fn test_encode_decode_view() {
let key = RowSequenceKey {
storage: StorageId::view(0xABCD),
};
let encoded = key.encode();
let expected = vec![0xF7, 0x02, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x54, 0x32];
assert_eq!(encoded.as_slice(), expected);
let key = RowSequenceKey::decode(&encoded).unwrap();
assert_eq!(key.storage, StorageId::view(0xABCD));
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, KeyCodec, Hash)]
#[key(tag = RowSettings)]
pub struct RowSettingsKey {
pub storage: StorageId,
}
impl RowSettingsKey {
pub fn new(storage: StorageId) -> Self {
Self {
storage,
}
}
pub fn encoded(storage: StorageId) -> EncodedKey {
Self::new(storage).encode()
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[cfg(test)]
pub mod row_settings_key_tests {
use super::*;
use crate::interface::catalog::id::{RingBufferId, SeriesId, TableId, ViewId};
#[test]
fn test_row_settings_key_encoding() {
let key = RowSettingsKey {
storage: StorageId::Table(TableId(42)),
};
let encoded = key.encode();
let decoded = RowSettingsKey::decode(&encoded).unwrap();
assert_eq!(key, decoded);
}
#[test]
fn test_row_settings_key_roundtrip_view() {
let key = RowSettingsKey {
storage: StorageId::View(ViewId(13)),
};
let encoded = key.encode();
let decoded = RowSettingsKey::decode(&encoded).unwrap();
assert_eq!(key, decoded);
}
#[test]
fn test_row_settings_key_roundtrip_ringbuffer() {
let key = RowSettingsKey {
storage: StorageId::RingBuffer(RingBufferId(99)),
};
let encoded = key.encode();
let decoded = RowSettingsKey::decode(&encoded).unwrap();
assert_eq!(key, decoded);
}
#[test]
fn test_row_settings_key_roundtrip_series() {
let key = RowSettingsKey {
storage: StorageId::Series(SeriesId(7)),
};
let encoded = key.encode();
let decoded = RowSettingsKey::decode(&encoded).unwrap();
assert_eq!(key, decoded);
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = RowShape)]
pub struct RowShapeKey {
pub fingerprint: RowShapeFingerprint,
}
impl RowShapeKey {
pub fn new(fingerprint: RowShapeFingerprint) -> Self {
Self {
fingerprint,
}
}
pub fn encoded(fingerprint: RowShapeFingerprint) -> EncodedKey {
Self {
fingerprint,
}
.encode()
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = RowShapeField)]
pub struct RowShapeFieldKey {
pub shape_fingerprint: RowShapeFingerprint,
pub field_index: u16,
}
impl RowShapeFieldKey {
pub fn new(shape_fingerprint: RowShapeFingerprint, field_index: u16) -> Self {
Self {
shape_fingerprint,
field_index,
}
}
pub fn encoded(shape_fingerprint: RowShapeFingerprint, field_index: u16) -> EncodedKey {
Self::new(shape_fingerprint, field_index).encode()
}
pub fn scan_for_shape(fingerprint: RowShapeFingerprint) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, fingerprint.as_u64() as u128)])
}
}
#[cfg(test)]
mod row_shape_key_tests {
use super::*;
#[test]
fn test_shape_key_encode_decode() {
let key = RowShapeKey {
fingerprint: RowShapeFingerprint::new(0xDEADBEEFCAFEBABE),
};
let encoded = key.encode();
let decoded = RowShapeKey::decode(&encoded).unwrap();
assert_eq!(decoded.fingerprint, RowShapeFingerprint::new(0xDEADBEEFCAFEBABE));
}
#[test]
fn test_shape_field_key_encode_decode() {
let key = RowShapeFieldKey {
shape_fingerprint: RowShapeFingerprint::new(0x1234567890ABCDEF),
field_index: 42,
};
let encoded = key.encode();
let decoded = RowShapeFieldKey::decode(&encoded).unwrap();
assert_eq!(decoded.shape_fingerprint, RowShapeFingerprint::new(0x1234567890ABCDEF));
assert_eq!(decoded.field_index, 42);
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = PartitionedRow)]
pub struct PartitionedRowKey {
pub storage: StorageId,
pub partition: Partition,
pub row: RowNumber,
}
impl PartitionedRowKey {
pub fn new(storage: impl Into<StorageId>, partition: Partition, row: RowNumber) -> Self {
Self {
storage: storage.into(),
partition,
row,
}
}
pub fn encoded(storage: impl Into<StorageId>, partition: Partition, row: RowNumber) -> EncodedKey {
Self::new(storage, partition, row).encode()
}
pub fn storage_start(storage: impl Into<StorageId>) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(PartitionedRowKey::TAG as u8).extend_object_id(storage.into());
serializer.to_encoded_key()
}
pub fn storage_end(storage: impl Into<StorageId>) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(10);
serializer
.extend_u8(PartitionedRowKey::TAG as u8)
.extend_object_id(ObjectId::from(storage.into()).prev());
serializer.to_encoded_key()
}
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()?)
}
pub fn full_scan(storage: impl Into<StorageId>) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, object_fields(ObjectId::from(storage.into())))
}
pub fn 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(
Self::TAG,
object_fields(ObjectId::from(storage.into()))
.into_iter()
.chain([Field::UDesc(Width::U128, partition.0)])
.collect::<Vec<_>>(),
)
}
pub fn partition_scan_range(
storage: impl Into<StorageId>,
partition: Partition,
last: Option<&TaggedKey>,
) -> TaggedKeyBoundRange {
Self::partition_range(storage, partition).resume_after(last)
}
}
#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct StoragePartitionedRowKey {
pub partition: Desc<Partition>,
pub row: Desc<RowNumber>,
}
impl StoragePartitionedRowKey {
pub fn new(partition: Partition, row: RowNumber) -> Self {
Self {
partition: Desc(partition),
row: Desc(row),
}
}
pub fn partition(self) -> Partition {
self.partition.0
}
pub fn row(self) -> RowNumber {
self.row.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, row: RowNumber) -> Self {
Self::new(Partition(((partition_hi as u128) << 64) | partition_lo as u128), row)
}
pub fn with_storage(self, storage: StorageId) -> PartitionedRowKey {
PartitionedRowKey {
storage,
partition: self.partition(),
row: self.row(),
}
}
}
impl From<PartitionedRowKey> for StoragePartitionedRowKey {
fn from(key: PartitionedRowKey) -> Self {
StoragePartitionedRowKey::new(key.partition, key.row)
}
}
impl HeapSize for StoragePartitionedRowKey {
fn heap_size(&self) -> usize {
0
}
}
impl BoundedKey for StoragePartitionedRowKey {
fn low() -> Self {
Self {
partition: <Desc<Partition> as BoundedKey>::low(),
row: <Desc<RowNumber> as BoundedKey>::low(),
}
}
}
impl DenseKey for StoragePartitionedRowKey {
fn successor(&self) -> Option<Self> {
if let Some(row) = self.row.successor() {
return Some(Self {
partition: self.partition,
row,
});
}
Some(Self {
partition: self.partition.successor()?,
row: <Desc<RowNumber> as BoundedKey>::low(),
})
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct PartitionedRowKeyRange {
pub storage: StorageId,
}
impl PartitionedRowKeyRange {
fn decode_key(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 != Self::TAG {
return None;
}
let storage = StorageId::from_object(de.read_object_id().ok()?)?;
Some(PartitionedRowKeyRange {
storage,
})
}
}
impl KeyRangeCodec for PartitionedRowKeyRange {
const TAG: KeyTag = KeyTag::PartitionedRow;
fn start(&self) -> Option<EncodedKey> {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(Self::TAG as u8).extend_object_id(self.storage);
Some(serializer.to_encoded_key())
}
fn end(&self) -> Option<EncodedKey> {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(Self::TAG as u8).extend_object_id(ObjectId::from(self.storage).prev());
Some(serializer.to_encoded_key())
}
fn decode(range: &EncodedKeyRange) -> (Option<Self>, Option<Self>)
where
Self: Sized,
{
let start_key = match &range.start {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
Bound::Unbounded => None,
};
let end_key = match &range.end {
Bound::Included(key) | Bound::Excluded(key) => Self::decode_key(key),
Bound::Unbounded => None,
};
(start_key, end_key)
}
}
#[cfg(test)]
mod partitioned_row_key_tests {
use std::ops::RangeBounds;
use reifydb_codec::key::{encoded::EncodedKey, serializer::KeySerializer};
use reifydb_value::value::{Value, partition::Partition, row_number::RowNumber};
use super::{PartitionedRowKey, RowKey, StoragePartitionedRowKey};
use crate::{
interface::catalog::{
id::{TableId, ViewId},
object::ObjectId,
storage::StorageId,
},
key::{
catalog::KeySerializerCatalogExt,
typed::{BoundedKey, DenseKey},
},
};
fn part(v: &str) -> Partition {
Partition::of(&[Value::Utf8(v.to_string())])
}
#[test]
fn test_table_roundtrip() {
let key = PartitionedRowKey {
storage: StorageId::Table(TableId(7)),
partition: part("us"),
row: RowNumber(42),
};
let decoded = PartitionedRowKey::decode(&key.encode()).unwrap();
assert_eq!(decoded, key);
}
#[test]
fn test_view_roundtrip() {
let key = PartitionedRowKey {
storage: StorageId::View(ViewId(11)),
partition: part("us"),
row: RowNumber(42),
};
let decoded = PartitionedRowKey::decode(&key.encode()).unwrap();
assert_eq!(decoded, key);
}
#[test]
fn test_storage_of() {
let key = PartitionedRowKey::encoded(StorageId::Table(TableId(42)), part("us"), RowNumber(1));
assert_eq!(PartitionedRowKey::storage_of(&key), Some(StorageId::Table(TableId(42))));
}
#[test]
fn test_storage_of_rejects_a_rowless_object() {
let mut serializer = KeySerializer::with_capacity(10);
serializer.extend_u8(PartitionedRowKey::TAG as u8).extend_object_id(ObjectId::vtable(42));
assert_eq!(PartitionedRowKey::storage_of(&serializer.to_encoded_key()), None);
}
#[test]
fn test_partition_rows_cluster_together() {
let storage = StorageId::Table(TableId(1));
let us_a = PartitionedRowKey::encoded(storage, part("us"), RowNumber(1));
let us_b = PartitionedRowKey::encoded(storage, part("us"), RowNumber(2));
let eu = PartitionedRowKey::encoded(storage, part("eu"), RowNumber(1));
let mut keys = [us_a.clone(), us_b.clone(), eu.clone()];
keys.sort();
let us_positions: Vec<usize> =
keys.iter().enumerate().filter(|(_, k)| **k == us_a || **k == us_b).map(|(i, _)| i).collect();
assert_eq!(us_positions[1] - us_positions[0], 1, "us partition rows must be contiguous");
}
#[test]
fn test_partition_range_contains_only_its_partition() {
let storage = StorageId::Table(TableId(1));
let range = PartitionedRowKey::partition_range(storage, part("us")).encode();
let us = PartitionedRowKey::encoded(storage, part("us"), RowNumber(500));
let eu = PartitionedRowKey::encoded(storage, part("eu"), RowNumber(1));
assert!(range.contains(&us), "us row must be inside the us partition range");
assert!(!range.contains(&eu), "eu row must be outside the us partition range");
}
#[test]
fn test_partitioned_row_ident_roundtrip() {
let key = PartitionedRowKey {
storage: StorageId::Table(TableId(7)),
partition: part("us"),
row: RowNumber(42),
};
let ident: StoragePartitionedRowKey = key.clone().into();
let restored = ident.with_storage(key.storage);
assert_eq!(restored, key);
}
#[test]
fn test_partitioned_row_ident_halves_split_correctly() {
let partition = Partition(0x1122334455667788_99AABBCCDDEEFF00);
let ident = StoragePartitionedRowKey::new(partition, RowNumber(1));
assert_eq!(ident.partition_hi(), 0x1122334455667788);
assert_eq!(ident.partition_lo(), 0x99AABBCCDDEEFF00);
assert_eq!(ident.partition(), partition);
assert_eq!(
StoragePartitionedRowKey::from_halves(ident.partition_hi(), ident.partition_lo(), RowNumber(1)),
ident
);
}
#[test]
fn test_partitioned_row_ident_ordering_matches_field_order() {
let lower_partition = StoragePartitionedRowKey::new(Partition(1), RowNumber(999));
let higher_partition = StoragePartitionedRowKey::new(Partition(2), RowNumber(1));
let same_partition_lower_row = StoragePartitionedRowKey::new(Partition(2), RowNumber(1));
let same_partition_higher_row = StoragePartitionedRowKey::new(Partition(2), RowNumber(2));
assert!(higher_partition < lower_partition);
assert!(same_partition_higher_row < same_partition_lower_row);
}
#[test]
fn test_partitioned_row_ident_successor_carries_into_the_partition() {
let last_row = StoragePartitionedRowKey::new(Partition(5), RowNumber(0));
let next = last_row.successor().unwrap();
assert_eq!(next, StoragePartitionedRowKey::new(Partition(4), RowNumber(u64::MAX)));
assert!(next > last_row);
}
#[test]
fn test_partitioned_row_ident_low_is_the_greatest_partition_and_row() {
assert_eq!(
<StoragePartitionedRowKey as BoundedKey>::low(),
StoragePartitionedRowKey::new(Partition(u128::MAX), RowNumber(u64::MAX))
);
}
#[test]
fn test_decode_rejects_trailing_bytes() {
let exact = RowKey::encoded(StorageId::table(7), RowNumber(42));
assert_eq!(exact.as_slice().len(), 18);
assert_eq!(
RowKey::decode(&exact),
Some(RowKey {
storage: StorageId::table(7),
row: RowNumber(42)
})
);
let mut longer = exact.as_slice().to_vec();
longer.push(0x00);
assert_eq!(RowKey::decode(&EncodedKey::new(longer)), None);
}
#[test]
fn test_decode_rejects_a_longer_key_that_shares_the_kind_byte() {
let mut sorted_view = RowKey::encoded(StorageId::view(3), RowNumber(1)).as_slice().to_vec();
sorted_view.extend_from_slice(&[0xAA; 8]);
sorted_view.extend_from_slice(&99u64.to_be_bytes());
assert_eq!(RowKey::decode(&EncodedKey::new(sorted_view)), None);
}
}
#[cfg(test)]
mod sorted_view_run_tests {
use reifydb_codec::key::encoded::EncodedKey;
use reifydb_value::value::{partition::Partition, row_number::RowNumber};
use super::{PartitionedSortedViewRowKey, SortRun, SortedViewRowKey};
use crate::interface::catalog::storage::StorageId;
fn run(bytes: &[u8]) -> SortRun {
SortRun::new(bytes)
}
#[test]
fn test_a_run_round_trips_through_the_terminator() {
let key = SortedViewRowKey::new(StorageId::view(3), run(&[0x10, 0x00, 0xff, 0x00]), RowNumber(42));
let decoded = SortedViewRowKey::decode(&key.encode()).unwrap();
assert_eq!(decoded, key);
assert_eq!(decoded.run.as_slice(), &[0x10, 0x00, 0xff, 0x00]);
}
#[test]
fn test_an_empty_run_round_trips() {
let key = SortedViewRowKey::new(StorageId::view(3), run(&[]), RowNumber(1));
assert_eq!(SortedViewRowKey::decode(&key.encode()), Some(key));
}
#[test]
fn test_a_partitioned_run_round_trips() {
let key = PartitionedSortedViewRowKey::new(
StorageId::view(3),
Partition(0x1122334455667788_99AABBCCDDEEFF00),
run(&[0x00, 0x00, 0x01]),
RowNumber(7),
);
assert_eq!(PartitionedSortedViewRowKey::decode(&key.encode()), Some(key));
}
#[test]
fn test_decode_refuses_the_other_kind_and_trailing_bytes() {
let partitioned = PartitionedSortedViewRowKey::encoded(
StorageId::view(3),
Partition(1),
run(&[0x10]),
RowNumber(1),
);
assert_eq!(SortedViewRowKey::decode(&partitioned), None);
let mut longer = SortedViewRowKey::encoded(StorageId::view(3), run(&[0x10]), RowNumber(1)).to_vec();
longer.push(0x00);
assert_eq!(SortedViewRowKey::decode(&EncodedKey::new(longer)), None);
}
#[test]
fn test_ord_matches_the_encoded_byte_order() {
let mut keys = vec![
SortedViewRowKey::new(StorageId::table(3), run(&[0x10]), RowNumber(1)),
SortedViewRowKey::new(StorageId::view(3), run(&[0x10]), RowNumber(1)),
SortedViewRowKey::new(StorageId::view(4), run(&[0x10]), RowNumber(1)),
SortedViewRowKey::new(StorageId::view(3), run(&[0x10]), RowNumber(2)),
SortedViewRowKey::new(StorageId::view(3), run(&[0x10, 0x00]), RowNumber(0)),
SortedViewRowKey::new(StorageId::view(3), run(&[0x10, 0x01]), RowNumber(0)),
SortedViewRowKey::new(StorageId::view(3), run(&[0x00]), RowNumber(0)),
SortedViewRowKey::new(StorageId::view(3), run(&[]), RowNumber(u64::MAX)),
];
keys.sort();
let encoded: Vec<EncodedKey> = keys.iter().map(|key| key.encode()).collect();
let mut sorted_bytes = encoded.clone();
sorted_bytes.sort();
assert_eq!(encoded, sorted_bytes);
}
#[test]
fn test_partitioned_ord_matches_the_encoded_byte_order() {
let mut keys = vec![
PartitionedSortedViewRowKey::new(StorageId::view(3), Partition(1), run(&[0x10]), RowNumber(1)),
PartitionedSortedViewRowKey::new(StorageId::view(3), Partition(2), run(&[0x10]), RowNumber(1)),
PartitionedSortedViewRowKey::new(StorageId::view(3), Partition(2), run(&[0x00]), RowNumber(1)),
PartitionedSortedViewRowKey::new(StorageId::view(3), Partition(2), run(&[0x10]), RowNumber(0)),
PartitionedSortedViewRowKey::new(StorageId::view(4), Partition(2), run(&[0x10]), RowNumber(0)),
];
keys.sort();
let encoded: Vec<EncodedKey> = keys.iter().map(|key| key.encode()).collect();
let mut sorted_bytes = encoded.clone();
sorted_bytes.sort();
assert_eq!(encoded, sorted_bytes);
}
#[test]
fn test_a_shorter_run_sorts_before_the_run_that_extends_it() {
let short = SortedViewRowKey::new(StorageId::view(3), run(&[0x10]), RowNumber(u64::MAX));
let long = SortedViewRowKey::new(StorageId::view(3), run(&[0x10, 0x00]), RowNumber(0));
assert!(short < long);
assert!(short.encode() < long.encode());
}
#[test]
fn test_the_helpers_still_read_the_new_layout() {
let storage = StorageId::view(3);
let key = SortedViewRowKey::encoded(storage, run(&[0x00, 0xAA]), RowNumber(9));
assert_eq!(SortedViewRowKey::storage_of(&key), Some(storage));
assert_eq!(SortedViewRowKey::row_of(&key), Some(RowNumber(9)));
let partitioned =
PartitionedSortedViewRowKey::encoded(storage, Partition(5), run(&[0x00, 0xAA]), RowNumber(9));
assert_eq!(PartitionedSortedViewRowKey::storage_of(&partitioned), Some(storage));
assert_eq!(PartitionedSortedViewRowKey::row_of(&partitioned), Some(RowNumber(9)));
assert_eq!(SortedViewRowKey::row_of(&partitioned), None);
}
}
impl KeyFields for SortedViewRowKey {
fn fields(&self) -> SmallVec<[Field<'_>; 6]> {
smallvec![
Field::UAsc(Width::U8, ObjectId::from(self.storage).type_tag() as u128),
Field::UDesc(Width::U64, ObjectId::from(self.storage).as_u64() as u128),
Field::RawAsc(RawEncoding::SortRun, Cow::Borrowed(self.run.as_slice())),
Field::UAsc(Width::U64, self.row.0.0 as u128),
]
}
}
impl KeyFields for PartitionedSortedViewRowKey {
fn fields(&self) -> SmallVec<[Field<'_>; 6]> {
smallvec![
Field::UAsc(Width::U8, ObjectId::from(self.storage).type_tag() as u128),
Field::UDesc(Width::U64, ObjectId::from(self.storage).as_u64() as u128),
Field::UDesc(Width::U128, self.partition.0),
Field::RawAsc(RawEncoding::SortRun, Cow::Borrowed(self.run.as_slice())),
Field::UAsc(Width::U64, self.row.0.0 as u128),
]
}
}