1use std::collections::Bound;
5
6use reifydb_codec::key::{
7 deserializer::KeyDeserializer,
8 encoded::{EncodedKey, EncodedKeyRange},
9};
10use reifydb_macro::KeyCodec;
11use reifydb_value::value::partition::Partition;
12
13use super::KeyTag;
14use crate::{
15 interface::catalog::{id::SeriesId, object::ObjectId, storage::StorageId},
16 key::{
17 any::{Field, KeyFields, TaggedKey, Width},
18 bound::{OwnedField, TaggedKeyBound, TaggedKeyBoundRange, object_fields},
19 catalog::{KeyDeserializerCatalogExt, KeySerializerCatalogExt},
20 typed::{BoundedKey, DenseKey, direction::Desc},
21 },
22 metrics::heap::HeapSize,
23};
24
25#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
26#[key(tag = Series)]
27pub struct SeriesKey {
28 pub series: SeriesId,
29}
30
31impl SeriesKey {
32 pub fn new(series: SeriesId) -> Self {
33 Self {
34 series,
35 }
36 }
37
38 pub fn encoded(series: impl Into<SeriesId>) -> EncodedKey {
39 Self::new(series.into()).encode()
40 }
41
42 pub fn full_scan() -> TaggedKeyBoundRange {
43 TaggedKeyBoundRange::kind(Self::TAG)
44 }
45}
46
47#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
48#[key(tag = SeriesMetadata)]
49pub struct SeriesMetadataKey {
50 pub storage: StorageId,
51}
52
53impl SeriesMetadataKey {
54 pub fn new(storage: impl Into<StorageId>) -> Self {
55 Self {
56 storage: storage.into(),
57 }
58 }
59
60 pub fn encoded(storage: impl Into<StorageId>) -> EncodedKey {
61 Self::new(storage).encode()
62 }
63}
64
65#[cfg(test)]
66mod series_metadata_key_tests {
67 use reifydb_codec::key::serializer::KeySerializer;
68
69 use super::{KeyTag, SeriesKey, SeriesMetadataKey};
70 use crate::interface::catalog::{
71 id::{SeriesId, ViewId},
72 storage::StorageId,
73 };
74
75 #[test]
76 fn test_metadata_key_roundtrip_series() {
77 let key = SeriesMetadataKey {
79 storage: StorageId::Series(SeriesId(7)),
80 };
81 assert_eq!(SeriesMetadataKey::decode(&key.encode()).unwrap(), key);
82 }
83
84 #[test]
85 fn test_metadata_key_roundtrip_view() {
86 let key = SeriesMetadataKey {
88 storage: StorageId::View(ViewId(7)),
89 };
90 assert_eq!(SeriesMetadataKey::decode(&key.encode()).unwrap(), key);
91 }
92
93 #[test]
94 fn test_metadata_key_separates_a_series_from_a_view_with_the_same_id() {
95 assert_ne!(SeriesMetadataKey::encoded(SeriesId(7)), SeriesMetadataKey::encoded(ViewId(7)));
97 }
98
99 #[test]
100 fn test_series_key_matches_legacy_byte_layout() {
101 for id in [SeriesId(0), SeriesId(1), SeriesId(u64::MAX)] {
102 let mut legacy = KeySerializer::with_capacity(9);
103 legacy.extend_u8(KeyTag::Series as u8).extend_u64(id);
104 assert_eq!(legacy.to_encoded_key().as_slice(), SeriesKey::encoded(id).as_slice());
105 }
106 }
107}
108
109#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
110#[key(tag = SeriesRow)]
111pub struct SeriesRowKey {
112 pub storage: StorageId,
113 pub variant_tag: Option<u8>,
114 pub key: u64,
115 pub sequence: u64,
116}
117
118#[derive(Debug, Clone)]
119pub struct SeriesRowKeyRange {
120 pub storage: StorageId,
121 pub variant_tag: Option<u8>,
122 pub key_start: Option<u64>,
123 pub key_end: Option<u64>,
124}
125
126impl SeriesRowKeyRange {
127 pub fn storage_start(storage: StorageId) -> EncodedKey {
128 TaggedKeyBound::prefix(SeriesRowKey::TAG, object_fields(ObjectId::from(storage))).encode()
129 }
130
131 pub fn storage_end(storage: StorageId) -> EncodedKey {
132 TaggedKeyBound::prefix(SeriesRowKey::TAG, object_fields(ObjectId::from(storage).prev())).encode()
133 }
134
135 pub fn full_scan(storage: StorageId, variant_tag: Option<u8>) -> TaggedKeyBoundRange {
136 SeriesRowKeyRange {
137 storage,
138 variant_tag,
139 key_start: None,
140 key_end: None,
141 }
142 .bounds()
143 }
144
145 pub fn scan_range(
146 storage: StorageId,
147 variant_tag: Option<u8>,
148 key_start: Option<u64>,
149 key_end: Option<u64>,
150 last: Option<&TaggedKey>,
151 ) -> TaggedKeyBoundRange {
152 if matches!(key_end, Some(0)) {
153 return TaggedKeyBoundRange::empty(SeriesRowKey::TAG);
154 }
155
156 SeriesRowKeyRange {
157 storage,
158 variant_tag,
159 key_start,
160 key_end,
161 }
162 .bounds()
163 .resume_after(last)
164 }
165
166 pub fn decode_storage(key: &EncodedKey) -> Option<StorageId> {
167 let mut de = KeyDeserializer::from_bytes(key.as_slice());
168
169 let kind: KeyTag = de.read_u8().ok()?.try_into().ok()?;
170 if kind != SeriesRowKey::TAG {
171 return None;
172 }
173
174 StorageId::from_object(de.read_object_id().ok()?)
175 }
176
177 pub fn decode(range: &EncodedKeyRange) -> (Option<StorageId>, Option<StorageId>) {
178 let start = match &range.start {
179 Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
180 Bound::Unbounded => None,
181 };
182
183 let end = match &range.end {
184 Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
185 Bound::Unbounded => None,
186 };
187
188 (start, end)
189 }
190
191 fn head_fields(&self, tagged: bool) -> Vec<OwnedField> {
192 let mut fields = object_fields(ObjectId::from(self.storage)).to_vec();
193 match self.variant_tag {
194 Some(tag) => {
195 fields.push(Field::UDesc(Width::U8, 1));
196 fields.push(Field::UDesc(Width::U8, tag as u128));
197 }
198 None if tagged || self.key_start.is_some() || self.key_end.is_some() => {
199 fields.push(Field::UDesc(Width::U8, 0));
200 fields.push(Field::UDesc(Width::U8, 0));
201 }
202 None => {}
203 }
204 fields
205 }
206
207 fn bounds(&self) -> TaggedKeyBoundRange {
208 let mut start = self.head_fields(false);
209 if let Some(key_val) = self.key_end {
210 start.push(Field::UDesc(Width::U64, (key_val - 1) as u128));
211 }
212 let start = Bound::Included(TaggedKeyBound::prefix(SeriesRowKey::TAG, start));
213
214 let end = match self.key_start {
215 Some(key_val) => {
216 let mut end = self.head_fields(true);
217 end.push(Field::UDesc(Width::U64, key_val as u128));
218 end.push(Field::UDesc(Width::U64, 0));
219 Bound::Included(TaggedKeyBound::prefix(SeriesRowKey::TAG, end))
220 }
221 None => Bound::Excluded(TaggedKeyBound::prefix_end(
222 SeriesRowKey::TAG,
223 object_fields(ObjectId::from(self.storage)),
224 )),
225 };
226
227 TaggedKeyBoundRange {
228 start,
229 end,
230 }
231 }
232}
233
234#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
235pub struct StorageSeriesKey {
236 pub variant_tag: Desc<Option<u8>>,
237 pub key: Desc<u64>,
238 pub sequence: Desc<u64>,
239}
240
241impl StorageSeriesKey {
242 pub fn new(variant_tag: Option<u8>, key: u64, sequence: u64) -> Self {
243 Self {
244 variant_tag: Desc(variant_tag),
245 key: Desc(key),
246 sequence: Desc(sequence),
247 }
248 }
249
250 pub fn variant_tag(self) -> Option<u8> {
251 self.variant_tag.0
252 }
253
254 pub fn key(self) -> u64 {
255 self.key.0
256 }
257
258 pub fn sequence(self) -> u64 {
259 self.sequence.0
260 }
261
262 pub fn to_sql_columns(self) -> SeriesKeyColumns {
263 SeriesKeyColumns {
264 variant_tag: series_variant_tag_to_sql(self.variant_tag()),
265 key: series_key_to_sql(self.key()),
266 sequence: series_sequence_to_sql(self.sequence()),
267 }
268 }
269
270 pub fn from_sql_columns(columns: SeriesKeyColumns) -> Self {
271 Self::new(
272 series_variant_tag_from_sql(columns.variant_tag),
273 series_key_from_sql(columns.key),
274 series_sequence_from_sql(columns.sequence),
275 )
276 }
277
278 pub fn with_storage(self, storage: StorageId) -> SeriesRowKey {
279 SeriesRowKey {
280 storage,
281 variant_tag: self.variant_tag(),
282 key: self.key(),
283 sequence: self.sequence(),
284 }
285 }
286}
287
288impl From<SeriesRowKey> for StorageSeriesKey {
289 fn from(key: SeriesRowKey) -> Self {
290 StorageSeriesKey::new(key.variant_tag, key.key, key.sequence)
291 }
292}
293
294impl HeapSize for StorageSeriesKey {
295 fn heap_size(&self) -> usize {
296 0
297 }
298}
299
300impl BoundedKey for StorageSeriesKey {
301 fn low() -> Self {
302 Self {
303 variant_tag: lowest_variant_tag(),
304 key: <Desc<u64> as BoundedKey>::low(),
305 sequence: <Desc<u64> as BoundedKey>::low(),
306 }
307 }
308}
309
310impl DenseKey for StorageSeriesKey {
311 fn successor(&self) -> Option<Self> {
312 if let Some(sequence) = self.sequence.successor() {
313 return Some(Self {
314 variant_tag: self.variant_tag,
315 key: self.key,
316 sequence,
317 });
318 }
319 if let Some(key) = self.key.successor() {
320 return Some(Self {
321 variant_tag: self.variant_tag,
322 key,
323 sequence: <Desc<u64> as BoundedKey>::low(),
324 });
325 }
326 Some(Self {
327 variant_tag: next_variant_tag(self.variant_tag)?,
328 key: <Desc<u64> as BoundedKey>::low(),
329 sequence: <Desc<u64> as BoundedKey>::low(),
330 })
331 }
332}
333
334fn lowest_variant_tag() -> Desc<Option<u8>> {
335 Desc(Some(u8::MAX))
336}
337
338fn next_variant_tag(tag: Desc<Option<u8>>) -> Option<Desc<Option<u8>>> {
339 match tag.0 {
340 Some(0) => Some(Desc(None)),
341 Some(value) => Some(Desc(Some(value - 1))),
342 None => None,
343 }
344}
345
346#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
347pub struct SeriesKeyColumns {
348 pub variant_tag: i64,
349 pub key: i64,
350 pub sequence: i64,
351}
352
353#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
354pub struct PartitionedSeriesKeyColumns {
355 pub partition_hi: i64,
356 pub partition_lo: i64,
357 pub variant_tag: i64,
358 pub key: i64,
359 pub sequence: i64,
360}
361
362const SERIES_VARIANT_TAG_NONE_SQL: i64 = (((!0u8) as i64) << 8) | ((!0u8) as i64);
363
364pub fn series_variant_tag_to_sql(variant_tag: Option<u8>) -> i64 {
365 match variant_tag {
366 Some(tag) => (((!1u8) as i64) << 8) | ((!tag) as i64),
367 None => SERIES_VARIANT_TAG_NONE_SQL,
368 }
369}
370
371pub fn series_variant_tag_from_sql(value: i64) -> Option<u8> {
372 if value == SERIES_VARIANT_TAG_NONE_SQL {
373 return None;
374 }
375 Some(!(value as u8))
376}
377
378pub fn series_key_to_sql(key: u64) -> i64 {
379 desc_u64_to_sql(key)
380}
381
382pub fn series_key_from_sql(value: i64) -> u64 {
383 desc_u64_from_sql(value)
384}
385
386pub fn series_sequence_to_sql(sequence: u64) -> i64 {
387 desc_u64_to_sql(sequence)
388}
389
390pub fn series_sequence_from_sql(value: i64) -> u64 {
391 desc_u64_from_sql(value)
392}
393
394pub fn series_partition_half_to_sql(half: u64) -> i64 {
395 desc_u64_to_sql(half)
396}
397
398pub fn series_partition_half_from_sql(value: i64) -> u64 {
399 desc_u64_from_sql(value)
400}
401
402fn desc_u64_to_sql(value: u64) -> i64 {
403 ((!value) ^ (1u64 << 63)) as i64
404}
405
406fn desc_u64_from_sql(value: i64) -> u64 {
407 !((value as u64) ^ (1u64 << 63))
408}
409
410#[cfg(test)]
411mod row_key_range_tests {
412 use std::ops::RangeBounds;
413
414 use super::*;
415 use crate::{
416 interface::catalog::id::{SeriesId, ViewId},
417 key::{KeyRangeCodec, row::RowKeyRange},
418 };
419
420 #[test]
421 fn test_encode_decode_without_tag() {
422 let key = SeriesRowKey {
424 storage: StorageId::Series(SeriesId(42)),
425 variant_tag: None,
426 key: 1706745600000,
427 sequence: 1,
428 };
429 let encoded = key.encode();
430 let decoded = SeriesRowKey::decode(&encoded).unwrap();
431 assert_eq!(decoded.storage, StorageId::Series(SeriesId(42)));
432 assert_eq!(decoded.variant_tag, None);
433 assert_eq!(decoded.key, 1706745600000);
434 assert_eq!(decoded.sequence, 1);
435 }
436
437 #[test]
438 fn test_encode_decode_with_tag() {
439 let key = SeriesRowKey {
441 storage: StorageId::Series(SeriesId(42)),
442 variant_tag: Some(3),
443 key: 1706745600000,
444 sequence: 5,
445 };
446 let encoded = key.encode();
447 let decoded = SeriesRowKey::decode(&encoded).unwrap();
448 assert_eq!(decoded.storage, StorageId::Series(SeriesId(42)));
449 assert_eq!(decoded.variant_tag, Some(3));
450 assert_eq!(decoded.key, 1706745600000);
451 assert_eq!(decoded.sequence, 5);
452 }
453
454 #[test]
455 fn test_a_view_storage_round_trips() {
456 let key = SeriesRowKey {
458 storage: StorageId::View(ViewId(42)),
459 variant_tag: Some(7),
460 key: 900,
461 sequence: 2,
462 };
463 let decoded = SeriesRowKey::decode(&key.encode()).unwrap();
464 assert_eq!(decoded, key);
465 }
466
467 #[test]
468 fn test_full_scan_range_still_names_its_storage() {
469 let range = SeriesRowKeyRange::full_scan(StorageId::Series(SeriesId(42)), None).encode();
471 let (start, end) = SeriesRowKeyRange::decode(&range);
472 assert_eq!(start, Some(StorageId::Series(SeriesId(42))));
473 assert_eq!(
474 end,
475 Some(StorageId::Series(SeriesId(41))),
476 "the exclusive end brackets the next series id down"
477 );
478 }
479
480 #[test]
481 fn test_untagged_full_scan_covers_tagged_rows() {
482 let storage = StorageId::Series(SeriesId(9));
484 let range = SeriesRowKeyRange::full_scan(storage, None).encode();
485 let untagged = SeriesRowKey {
486 storage,
487 variant_tag: None,
488 key: 500,
489 sequence: 0,
490 }
491 .encode();
492 let tagged = SeriesRowKey {
493 storage,
494 variant_tag: Some(4),
495 key: 500,
496 sequence: 0,
497 }
498 .encode();
499
500 assert!(range.contains(&untagged));
501 assert!(range.contains(&tagged), "an untagged full scan must still see tagged rows");
502 }
503
504 #[test]
505 fn test_scan_range_brackets_the_rows_it_selects() {
506 let storage = StorageId::Series(SeriesId(1));
508 let range = SeriesRowKeyRange::scan_range(storage, None, Some(100), Some(200), None).encode();
509 let inside = SeriesRowKey {
510 storage,
511 variant_tag: None,
512 key: 150,
513 sequence: 1,
514 }
515 .encode();
516 let below = SeriesRowKey {
517 storage,
518 variant_tag: None,
519 key: 99,
520 sequence: 1,
521 }
522 .encode();
523 let above = SeriesRowKey {
524 storage,
525 variant_tag: None,
526 key: 201,
527 sequence: 1,
528 }
529 .encode();
530
531 assert!(range.contains(&inside));
532 assert!(!range.contains(&below));
533 assert!(!range.contains(&above));
534 }
535
536 #[test]
537 fn test_row_key_range_never_claims_a_series_range() {
538 let range = SeriesRowKeyRange::full_scan(StorageId::Series(SeriesId(42)), None).encode();
540 assert_eq!(RowKeyRange::decode(&range), (None, None));
541 }
542
543 #[test]
544 fn test_ordering_by_key() {
545 let key1 = SeriesRowKey {
547 storage: StorageId::Series(SeriesId(1)),
548 variant_tag: None,
549 key: 100,
550 sequence: 0,
551 };
552 let key2 = SeriesRowKey {
553 storage: StorageId::Series(SeriesId(1)),
554 variant_tag: None,
555 key: 200,
556 sequence: 0,
557 };
558 let e1 = key1.encode();
559 let e2 = key2.encode();
560
561 assert!(e1 > e2, "key descending ordering not preserved");
562 }
563
564 #[test]
565 fn test_ordering_by_sequence() {
566 let key1 = SeriesRowKey {
568 storage: StorageId::Series(SeriesId(1)),
569 variant_tag: None,
570 key: 100,
571 sequence: 1,
572 };
573 let key2 = SeriesRowKey {
574 storage: StorageId::Series(SeriesId(1)),
575 variant_tag: None,
576 key: 100,
577 sequence: 2,
578 };
579 let e1 = key1.encode();
580 let e2 = key2.encode();
581
582 assert!(e1 > e2, "sequence descending ordering not preserved");
583 }
584
585 #[test]
586 fn test_half_bounded_untagged_range_excludes_tagged_rows() {
587 let range = SeriesRowKeyRange::scan_range(StorageId::Series(SeriesId(7)), None, Some(100), None, None)
591 .encode();
592
593 let untagged = SeriesRowKey {
594 storage: StorageId::Series(SeriesId(7)),
595 variant_tag: None,
596 key: 500,
597 sequence: 1,
598 }
599 .encode();
600 let tagged = SeriesRowKey {
601 storage: StorageId::Series(SeriesId(7)),
602 variant_tag: Some(3),
603 key: 500,
604 sequence: 1,
605 }
606 .encode();
607
608 assert!(range.contains(&untagged), "an untagged row above the lower bound must stay in range");
609 assert!(!range.contains(&tagged), "a tagged row must never leak into an untagged key-bounded range");
610 }
611}
612
613#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
614#[key(tag = PartitionedSeriesRow)]
615pub struct PartitionedSeriesRowKey {
616 pub storage: StorageId,
617 pub partition: Partition,
618 pub variant_tag: Option<u8>,
619 pub key: u64,
620 pub sequence: u64,
621}
622
623impl PartitionedSeriesRowKey {
624 pub fn new(
625 storage: impl Into<StorageId>,
626 partition: Partition,
627 variant_tag: Option<u8>,
628 key: u64,
629 sequence: u64,
630 ) -> Self {
631 Self {
632 storage: storage.into(),
633 partition,
634 variant_tag,
635 key,
636 sequence,
637 }
638 }
639
640 pub fn encoded(
641 storage: impl Into<StorageId>,
642 partition: Partition,
643 variant_tag: Option<u8>,
644 key: u64,
645 sequence: u64,
646 ) -> EncodedKey {
647 Self::new(storage, partition, variant_tag, key, sequence).encode()
648 }
649
650 pub fn storage_of(key: &EncodedKey) -> Option<StorageId> {
651 let mut de = KeyDeserializer::from_bytes(key.as_slice());
652 let kind: KeyTag = de.read_u8().ok()?.try_into().ok()?;
653 if kind != Self::TAG {
654 return None;
655 }
656 StorageId::from_object(de.read_object_id().ok()?)
657 }
658}
659
660#[derive(Debug, Clone)]
661pub struct PartitionedSeriesRowKeyRange {
662 pub storage: StorageId,
663 pub partition: Partition,
664 pub variant_tag: Option<u8>,
665 pub key_start: Option<u64>,
666 pub key_end: Option<u64>,
667}
668
669impl PartitionedSeriesRowKeyRange {
670 pub fn storage_start(storage: impl Into<StorageId>) -> EncodedKey {
671 TaggedKeyBound::prefix(PartitionedSeriesRowKey::TAG, object_fields(ObjectId::from(storage.into())))
672 .encode()
673 }
674
675 pub fn storage_end(storage: impl Into<StorageId>) -> EncodedKey {
676 TaggedKeyBound::prefix(
677 PartitionedSeriesRowKey::TAG,
678 object_fields(ObjectId::from(storage.into()).prev()),
679 )
680 .encode()
681 }
682
683 pub fn full_scan(storage: impl Into<StorageId>) -> TaggedKeyBoundRange {
684 TaggedKeyBoundRange::prefix(PartitionedSeriesRowKey::TAG, object_fields(ObjectId::from(storage.into())))
685 }
686
687 pub fn full_scan_range(storage: impl Into<StorageId>, last: Option<&TaggedKey>) -> TaggedKeyBoundRange {
688 Self::full_scan(storage).resume_after(last)
689 }
690
691 pub fn partition_range(storage: impl Into<StorageId>, partition: Partition) -> TaggedKeyBoundRange {
692 TaggedKeyBoundRange::prefix(
693 PartitionedSeriesRowKey::TAG,
694 Self::partition_fields(storage.into(), partition),
695 )
696 }
697
698 pub fn partition_scan_range(
699 storage: impl Into<StorageId>,
700 partition: Partition,
701 last: Option<&TaggedKey>,
702 ) -> TaggedKeyBoundRange {
703 Self::partition_range(storage, partition).resume_after(last)
704 }
705
706 pub fn scan_range(
707 storage: impl Into<StorageId>,
708 partition: Partition,
709 variant_tag: Option<u8>,
710 key_start: Option<u64>,
711 key_end: Option<u64>,
712 last: Option<&TaggedKey>,
713 ) -> TaggedKeyBoundRange {
714 if matches!(key_end, Some(0)) {
715 return TaggedKeyBoundRange::empty(PartitionedSeriesRowKey::TAG);
716 }
717
718 PartitionedSeriesRowKeyRange {
719 storage: storage.into(),
720 partition,
721 variant_tag,
722 key_start,
723 key_end,
724 }
725 .bounds()
726 .resume_after(last)
727 }
728
729 pub fn decode_storage(key: &EncodedKey) -> Option<StorageId> {
730 PartitionedSeriesRowKey::storage_of(key)
731 }
732
733 pub fn decode(range: &EncodedKeyRange) -> (Option<StorageId>, Option<StorageId>) {
734 let start = match &range.start {
735 Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
736 Bound::Unbounded => None,
737 };
738
739 let end = match &range.end {
740 Bound::Included(key) | Bound::Excluded(key) => Self::decode_storage(key),
741 Bound::Unbounded => None,
742 };
743
744 (start, end)
745 }
746
747 fn partition_fields(storage: StorageId, partition: Partition) -> Vec<OwnedField> {
748 let mut fields = object_fields(ObjectId::from(storage)).to_vec();
749 fields.push(Field::UDesc(Width::U128, partition.0));
750 fields
751 }
752
753 fn head_fields(&self, tagged: bool) -> Vec<OwnedField> {
754 let mut fields = Self::partition_fields(self.storage, self.partition);
755 match self.variant_tag {
756 Some(tag) => {
757 fields.push(Field::UDesc(Width::U8, 1));
758 fields.push(Field::UDesc(Width::U8, tag as u128));
759 }
760 None if tagged || self.key_start.is_some() || self.key_end.is_some() => {
761 fields.push(Field::UDesc(Width::U8, 0));
762 fields.push(Field::UDesc(Width::U8, 0));
763 }
764 None => {}
765 }
766 fields
767 }
768
769 fn bounds(&self) -> TaggedKeyBoundRange {
770 let mut start = self.head_fields(false);
771 if let Some(key_val) = self.key_end {
772 start.push(Field::UDesc(Width::U64, (key_val - 1) as u128));
773 }
774 let start = Bound::Included(TaggedKeyBound::prefix(PartitionedSeriesRowKey::TAG, start));
775
776 let end = match self.key_start {
777 Some(key_val) => {
778 let mut end = self.head_fields(true);
779 end.push(Field::UDesc(Width::U64, key_val as u128));
780 end.push(Field::UDesc(Width::U64, 0));
781 Bound::Included(TaggedKeyBound::prefix(PartitionedSeriesRowKey::TAG, end))
782 }
783 None => Bound::Excluded(TaggedKeyBound::prefix_end(
784 PartitionedSeriesRowKey::TAG,
785 Self::partition_fields(self.storage, self.partition),
786 )),
787 };
788
789 TaggedKeyBoundRange {
790 start,
791 end,
792 }
793 }
794}
795
796#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
797pub struct StoragePartitionedSeriesKey {
798 pub partition: Desc<Partition>,
799 pub variant_tag: Desc<Option<u8>>,
800 pub key: Desc<u64>,
801 pub sequence: Desc<u64>,
802}
803
804impl StoragePartitionedSeriesKey {
805 pub fn new(partition: Partition, variant_tag: Option<u8>, key: u64, sequence: u64) -> Self {
806 Self {
807 partition: Desc(partition),
808 variant_tag: Desc(variant_tag),
809 key: Desc(key),
810 sequence: Desc(sequence),
811 }
812 }
813
814 pub fn partition(self) -> Partition {
815 self.partition.0
816 }
817
818 pub fn variant_tag(self) -> Option<u8> {
819 self.variant_tag.0
820 }
821
822 pub fn key(self) -> u64 {
823 self.key.0
824 }
825
826 pub fn sequence(self) -> u64 {
827 self.sequence.0
828 }
829
830 pub fn partition_hi(self) -> u64 {
831 (self.partition().0 >> 64) as u64
832 }
833
834 pub fn partition_lo(self) -> u64 {
835 self.partition().0 as u64
836 }
837
838 pub fn from_halves(
839 partition_hi: u64,
840 partition_lo: u64,
841 variant_tag: Option<u8>,
842 key: u64,
843 sequence: u64,
844 ) -> Self {
845 Self::new(Partition(((partition_hi as u128) << 64) | partition_lo as u128), variant_tag, key, sequence)
846 }
847
848 pub fn to_sql_columns(self) -> PartitionedSeriesKeyColumns {
849 PartitionedSeriesKeyColumns {
850 partition_hi: series_partition_half_to_sql(self.partition_hi()),
851 partition_lo: series_partition_half_to_sql(self.partition_lo()),
852 variant_tag: series_variant_tag_to_sql(self.variant_tag()),
853 key: series_key_to_sql(self.key()),
854 sequence: series_sequence_to_sql(self.sequence()),
855 }
856 }
857
858 pub fn from_sql_columns(columns: PartitionedSeriesKeyColumns) -> Self {
859 Self::from_halves(
860 series_partition_half_from_sql(columns.partition_hi),
861 series_partition_half_from_sql(columns.partition_lo),
862 series_variant_tag_from_sql(columns.variant_tag),
863 series_key_from_sql(columns.key),
864 series_sequence_from_sql(columns.sequence),
865 )
866 }
867
868 pub fn with_storage(self, storage: StorageId) -> PartitionedSeriesRowKey {
869 PartitionedSeriesRowKey {
870 storage,
871 partition: self.partition(),
872 variant_tag: self.variant_tag(),
873 key: self.key(),
874 sequence: self.sequence(),
875 }
876 }
877}
878
879impl From<PartitionedSeriesRowKey> for StoragePartitionedSeriesKey {
880 fn from(key: PartitionedSeriesRowKey) -> Self {
881 StoragePartitionedSeriesKey::new(key.partition, key.variant_tag, key.key, key.sequence)
882 }
883}
884
885impl HeapSize for StoragePartitionedSeriesKey {
886 fn heap_size(&self) -> usize {
887 0
888 }
889}
890
891impl BoundedKey for StoragePartitionedSeriesKey {
892 fn low() -> Self {
893 Self {
894 partition: <Desc<Partition> as BoundedKey>::low(),
895 variant_tag: lowest_variant_tag(),
896 key: <Desc<u64> as BoundedKey>::low(),
897 sequence: <Desc<u64> as BoundedKey>::low(),
898 }
899 }
900}
901
902impl DenseKey for StoragePartitionedSeriesKey {
903 fn successor(&self) -> Option<Self> {
904 if let Some(sequence) = self.sequence.successor() {
905 return Some(Self {
906 sequence,
907 ..*self
908 });
909 }
910 if let Some(key) = self.key.successor() {
911 return Some(Self {
912 key,
913 sequence: <Desc<u64> as BoundedKey>::low(),
914 ..*self
915 });
916 }
917 if let Some(variant_tag) = next_variant_tag(self.variant_tag) {
918 return Some(Self {
919 partition: self.partition,
920 variant_tag,
921 key: <Desc<u64> as BoundedKey>::low(),
922 sequence: <Desc<u64> as BoundedKey>::low(),
923 });
924 }
925 Some(Self {
926 partition: self.partition.successor()?,
927 variant_tag: lowest_variant_tag(),
928 key: <Desc<u64> as BoundedKey>::low(),
929 sequence: <Desc<u64> as BoundedKey>::low(),
930 })
931 }
932}
933
934#[cfg(test)]
935mod partitioned_row_key_tests {
936 use std::ops::RangeBounds;
937
938 use reifydb_value::value::{Value, partition::Partition, row_number::RowNumber};
939
940 use super::*;
941 use crate::{
942 interface::catalog::id::{SeriesId, TableId, ViewId},
943 key::row::PartitionedRowKey,
944 };
945
946 fn part(v: &str) -> Partition {
947 Partition::of(&[Value::Utf8(v.to_string())])
948 }
949
950 #[test]
951 fn test_round_trip_without_tag() {
952 let key = PartitionedSeriesRowKey {
954 storage: StorageId::Series(SeriesId(3)),
955 partition: part("btc"),
956 variant_tag: None,
957 key: 1_700_000_000,
958 sequence: 9,
959 };
960 let decoded = PartitionedSeriesRowKey::decode(&key.encode()).unwrap();
961 assert_eq!(decoded, key);
962 }
963
964 #[test]
965 fn test_round_trip_with_tag() {
966 let key = PartitionedSeriesRowKey {
968 storage: StorageId::Series(SeriesId(3)),
969 partition: part("eth"),
970 variant_tag: Some(5),
971 key: 42,
972 sequence: 0,
973 };
974 let decoded = PartitionedSeriesRowKey::decode(&key.encode()).unwrap();
975 assert_eq!(decoded, key);
976 }
977
978 #[test]
979 fn test_a_view_storage_round_trips() {
980 let key = PartitionedSeriesRowKey {
982 storage: StorageId::View(ViewId(11)),
983 partition: part("us"),
984 variant_tag: Some(1),
985 key: 77,
986 sequence: 4,
987 };
988 let decoded = PartitionedSeriesRowKey::decode(&key.encode()).unwrap();
989 assert_eq!(decoded, key);
990 }
991
992 #[test]
993 fn test_storage_of() {
994 let key = PartitionedSeriesRowKey::encoded(StorageId::Series(SeriesId(42)), part("us"), None, 1, 0);
996 assert_eq!(PartitionedSeriesRowKey::storage_of(&key), Some(StorageId::Series(SeriesId(42))));
997 }
998
999 #[test]
1000 fn test_ordering_by_key_is_descending() {
1001 let storage = StorageId::Series(SeriesId(1));
1003 let low = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 100, 0);
1004 let high = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 200, 0);
1005
1006 assert!(low > high, "key descending ordering not preserved");
1007 }
1008
1009 #[test]
1010 fn test_untagged_full_scan_covers_tagged_rows() {
1011 let storage = StorageId::Series(SeriesId(9));
1013 let range = PartitionedSeriesRowKeyRange::full_scan(storage).encode();
1014 let untagged = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 500, 0);
1015 let tagged = PartitionedSeriesRowKey::encoded(storage, part("us"), Some(4), 500, 0);
1016
1017 assert!(range.contains(&untagged));
1018 assert!(range.contains(&tagged), "an untagged full scan must still see tagged rows");
1019 }
1020
1021 #[test]
1022 fn test_scan_range_brackets_the_rows_it_selects() {
1023 let storage = StorageId::Series(SeriesId(1));
1025 let partition = part("us");
1026 let range =
1027 PartitionedSeriesRowKeyRange::scan_range(storage, partition, None, Some(100), Some(200), None)
1028 .encode();
1029 let inside = PartitionedSeriesRowKey::encoded(storage, partition, None, 150, 1);
1030 let below = PartitionedSeriesRowKey::encoded(storage, partition, None, 99, 1);
1031 let above = PartitionedSeriesRowKey::encoded(storage, partition, None, 201, 1);
1032
1033 assert!(range.contains(&inside));
1034 assert!(!range.contains(&below));
1035 assert!(!range.contains(&above));
1036 }
1037
1038 #[test]
1039 fn test_scan_range_never_crosses_into_another_partition() {
1040 let storage = StorageId::Series(SeriesId(1));
1042 let range =
1043 PartitionedSeriesRowKeyRange::scan_range(storage, part("us"), None, Some(100), Some(200), None)
1044 .encode();
1045 let other = PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 150, 1);
1046
1047 assert!(!range.contains(&other), "an in-bounds key of another partition must stay outside");
1048 }
1049
1050 #[test]
1051 fn test_partition_range_covers_every_key_of_its_partition() {
1052 let storage = StorageId::Series(SeriesId(1));
1054 let range = PartitionedSeriesRowKeyRange::partition_range(storage, part("us")).encode();
1055 let untagged = PartitionedSeriesRowKey::encoded(storage, part("us"), None, 1, 0);
1056 let tagged = PartitionedSeriesRowKey::encoded(storage, part("us"), Some(9), u64::MAX, u64::MAX);
1057 let other = PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 1, 0);
1058
1059 assert!(range.contains(&untagged));
1060 assert!(range.contains(&tagged));
1061 assert!(!range.contains(&other));
1062 }
1063
1064 #[test]
1065 fn test_full_scan_range_resumes_across_partitions() {
1066 let storage = StorageId::Series(SeriesId(1));
1069 let cursor = TaggedKey::from(PartitionedSeriesRowKey::new(storage, part("us"), None, 200, 0));
1070 let range = PartitionedSeriesRowKeyRange::full_scan_range(storage, Some(&cursor)).encode();
1071
1072 assert!(!range.contains(&cursor.encode()));
1073 assert!(range.contains(&PartitionedSeriesRowKey::encoded(storage, part("us"), None, 100, 0)));
1074 assert!(
1075 range.contains(&PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 100, 0))
1076 || range.contains(&PartitionedSeriesRowKey::encoded(storage, part("eu"), None, 300, 0)),
1077 "a resumed object-wide scan must still be able to reach another partition"
1078 );
1079 }
1080
1081 #[test]
1082 fn test_full_scan_range_without_a_cursor_covers_every_partition() {
1083 let storage = StorageId::Series(SeriesId(1));
1086 let range = PartitionedSeriesRowKeyRange::full_scan_range(storage, None).encode();
1087
1088 assert!(range.contains(&PartitionedSeriesRowKey::encoded(storage, part("us"), None, 1, 0)));
1089 assert!(range.contains(&PartitionedSeriesRowKey::encoded(storage, part("eu"), Some(3), 9, 9)));
1090 assert!(!range.contains(&PartitionedSeriesRowKey::encoded(
1091 StorageId::Series(SeriesId(2)),
1092 part("us"),
1093 None,
1094 1,
1095 0
1096 )));
1097 }
1098
1099 #[test]
1100 fn test_partition_scan_range_resumes_after_the_cursor() {
1101 let storage = StorageId::Series(SeriesId(1));
1103 let partition = part("us");
1104 let cursor = TaggedKey::from(PartitionedSeriesRowKey::new(storage, partition, None, 200, 0));
1105 let range =
1106 PartitionedSeriesRowKeyRange::partition_scan_range(storage, partition, Some(&cursor)).encode();
1107 let next = PartitionedSeriesRowKey::encoded(storage, partition, None, 100, 0);
1108
1109 assert!(!range.contains(&cursor.encode()));
1110 assert!(range.contains(&next));
1111 }
1112
1113 #[test]
1114 fn test_the_two_partitioned_kinds_do_not_share_a_keyspace() {
1115 let series =
1117 PartitionedSeriesRowKey::encoded(StorageId::Series(SeriesId(1)), part("us"), Some(2), 7, 9);
1118 let row = PartitionedRowKey::encoded(StorageId::Table(TableId(1)), part("us"), RowNumber(9));
1119
1120 assert_ne!(series.as_slice()[0], row.as_slice()[0]);
1121 assert!(PartitionedRowKey::decode(&series).is_none());
1122 assert!(PartitionedSeriesRowKey::decode(&row).is_none());
1123 }
1124
1125 #[test]
1126 fn test_half_bounded_untagged_range_excludes_tagged_rows() {
1127 let range = PartitionedSeriesRowKeyRange::scan_range(
1130 StorageId::Series(SeriesId(7)),
1131 part("us"),
1132 None,
1133 Some(100),
1134 None,
1135 None,
1136 )
1137 .encode();
1138
1139 let untagged = PartitionedSeriesRowKey {
1140 storage: StorageId::Series(SeriesId(7)),
1141 partition: part("us"),
1142 variant_tag: None,
1143 key: 500,
1144 sequence: 1,
1145 }
1146 .encode();
1147 let tagged = PartitionedSeriesRowKey {
1148 storage: StorageId::Series(SeriesId(7)),
1149 partition: part("us"),
1150 variant_tag: Some(3),
1151 key: 500,
1152 sequence: 1,
1153 }
1154 .encode();
1155
1156 assert!(range.contains(&untagged), "an untagged row above the lower bound must stay in range");
1157 assert!(!range.contains(&tagged), "a tagged row must never leak into an untagged key-bounded range");
1158 }
1159}
1160
1161#[cfg(test)]
1162mod storage_series_key_tests {
1163 use reifydb_value::value::partition::Partition;
1164
1165 use super::{PartitionedSeriesRowKey, SeriesRowKey, StorageId, StoragePartitionedSeriesKey, StorageSeriesKey};
1166 use crate::{
1167 interface::catalog::id::SeriesId,
1168 key::typed::{BoundedKey, DenseKey},
1169 };
1170
1171 const STORAGE: StorageId = StorageId::Series(SeriesId(7));
1172
1173 fn series(variant_tag: Option<u8>, key: u64, sequence: u64) -> StorageSeriesKey {
1174 StorageSeriesKey::new(variant_tag, key, sequence)
1175 }
1176
1177 #[test]
1178 fn storage_key_order_matches_the_encoded_byte_order() {
1179 let mut keys = vec![
1182 series(None, 5, 1),
1183 series(Some(0), 5, 1),
1184 series(Some(9), 5, 1),
1185 series(Some(9), 5, 2),
1186 series(Some(9), 6, 1),
1187 ];
1188 let mut encoded: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
1189 keys.sort();
1190 encoded.sort();
1191
1192 let reordered: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
1193 assert_eq!(reordered, encoded);
1194 }
1195
1196 #[test]
1197 fn a_tagged_row_sorts_before_an_untagged_one() {
1198 assert!(series(Some(0), 1, 1) < series(None, 1, 1));
1201 assert!(series(Some(255), 1, 1) < series(Some(0), 1, 1));
1202 }
1203
1204 #[test]
1205 fn low_is_the_first_key_the_encoder_can_produce() {
1206 let low = <StorageSeriesKey as BoundedKey>::low();
1207 assert_eq!(low, series(Some(u8::MAX), u64::MAX, u64::MAX));
1208 for other in [series(None, 0, 0), series(Some(0), 0, 0), series(Some(200), 7, 3)] {
1209 assert!(low < other, "low must not sort above a real key");
1210 }
1211 }
1212
1213 #[test]
1214 fn successor_walks_the_sequence_then_the_key_then_the_tag() {
1215 assert_eq!(series(Some(9), 5, 2).successor(), Some(series(Some(9), 5, 1)));
1218 assert_eq!(series(Some(9), 5, 0).successor(), Some(series(Some(9), 4, u64::MAX)));
1219 assert_eq!(series(Some(9), 0, 0).successor(), Some(series(Some(8), u64::MAX, u64::MAX)));
1220 assert_eq!(series(Some(0), 0, 0).successor(), Some(series(None, u64::MAX, u64::MAX)));
1221 assert_eq!(series(None, 0, 0).successor(), None);
1222 }
1223
1224 #[test]
1225 fn every_successor_is_the_immediate_next_storage_key() {
1226 let mut walk = vec![series(Some(1), 1, 1)];
1227 for _ in 0..4 {
1228 walk.push(walk.last().unwrap().successor().unwrap());
1229 }
1230 assert!(walk.windows(2).all(|pair| pair[0] < pair[1]), "the walk must ascend: {walk:?}");
1231 }
1232
1233 #[test]
1234 fn partitioned_storage_key_order_matches_the_encoded_byte_order() {
1235 let mut keys = vec![
1236 StoragePartitionedSeriesKey::new(Partition(2), None, 5, 1),
1237 StoragePartitionedSeriesKey::new(Partition(2), Some(3), 5, 1),
1238 StoragePartitionedSeriesKey::new(Partition(1), Some(3), 5, 1),
1239 StoragePartitionedSeriesKey::new(Partition(2), Some(3), 5, 2),
1240 ];
1241 let mut encoded: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
1242 keys.sort();
1243 encoded.sort();
1244
1245 let reordered: Vec<_> = keys.iter().map(|it| it.with_storage(STORAGE).encode().to_vec()).collect();
1246 assert_eq!(reordered, encoded);
1247 }
1248
1249 #[test]
1250 fn partitioned_successor_carries_into_the_partition() {
1251 let last_of_partition = StoragePartitionedSeriesKey::new(Partition(5), None, 0, 0);
1253 assert_eq!(
1254 last_of_partition.successor(),
1255 Some(StoragePartitionedSeriesKey::new(Partition(4), Some(u8::MAX), u64::MAX, u64::MAX))
1256 );
1257 assert_eq!(StoragePartitionedSeriesKey::new(Partition(0), None, 0, 0).successor(), None);
1258 }
1259
1260 #[test]
1261 fn a_storage_key_round_trips_through_its_row_key() {
1262 let key = SeriesRowKey {
1263 storage: STORAGE,
1264 variant_tag: Some(4),
1265 key: 900,
1266 sequence: 12,
1267 };
1268 assert_eq!(StorageSeriesKey::from(key.clone()).with_storage(STORAGE), key);
1269
1270 let partitioned = PartitionedSeriesRowKey::new(STORAGE, Partition(3), None, 900, 12);
1271 assert_eq!(StoragePartitionedSeriesKey::from(partitioned.clone()).with_storage(STORAGE), partitioned);
1272 }
1273}