Skip to main content

reifydb_core/key/
series.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// The tag byte is what keeps a series' metadata out of a view's; a bare id would collide.
78		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		// A series-backed view must keep its metadata under its own id, never a backing object's.
87		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		// Both narrow from the same numeric id, so identical bytes would silently share one row.
96		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		// Without the flag byte a missing tag shifts every later field by one and the key reads back wrong.
423		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		// A tagged key must report the exact tag it was written with, never a byte borrowed from the key.
440		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		// Widening off SeriesId is pointless unless a non-series storage keeps its own tag through decode.
457		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		// classify_range routes a series scan to its own physical entry; an undecodable range falls to Multi.
470		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		// A flag byte in the range bounds would pin the scan to flag=0 and silently drop every tagged row.
483		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		// The range must contain a key inside the window and exclude one outside, or eviction skips live rows.
507		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		// A shared kind byte made a series scan classify as a plain row scan of the same object id.
539		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		// Reads walk newest first; an ascending key encoding would hand back the oldest rows instead.
546		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		// Two rows at the same key must still come back newest sequence first.
567		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		// A range bounded on one side only must still pin its tag class: the untagged flag encodes 0xFF and
588		// the tagged flag 0xFE, so omitting the flag from the start bound lets every tagged row of the series
589		// sort into the window regardless of its key.
590		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		// Without the flag byte an untagged key shifts every later field by one and reads back wrong.
953		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		// The partition sits between the object id and the tag, so a mis-sized partition eats the tag.
967		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		// A partitioned series materialised into a view must keep its view tag, never narrow to a series.
981		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		// Range classification reads the storage without decoding the whole key; a wrong offset misroutes it.
995		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		// Reads walk newest first; an ascending key encoding would hand back the oldest rows instead.
1002		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		// A flag byte in the range bounds would pin the scan to flag=0 and silently drop every tagged row.
1012		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		// The range must contain a key inside the window and exclude ones outside, or eviction skips live rows.
1024		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		// Bounding only the key span would let a neighbouring partition's rows be evicted with this one.
1041		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		// The prefix range is the eviction unit for a partition, so it must not depend on the tag.
1053		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		// An object-wide read pages through every partition at once, so the cursor must only cut the
1067		// rows already returned and never fence the scan into the cursor's own partition.
1068		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		// The unpaged form is the object-wide read the planner falls back to when no partition was
1084		// pruned; missing a partition there silently halves the answer.
1085		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		// A resumed page must exclude the cursor itself, otherwise the last row of a page repeats forever.
1102		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		// One kind byte for two layouts is what let a series key answer to a plain partitioned row read.
1116		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		// Same tag-class pinning as the unpartitioned range: the flag byte must appear in the start bound
1128		// whenever either key bound is set, or tagged rows (flag 0xFE) sort inside an untagged window (0xFF).
1129		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		// the storage key is the cache key for the same rows the encoded key orders on disk, so a disagreement
1180		// here silently hands back a neighbouring row on any ordered lookup
1181		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		// none is written as a zero presence flag, which inverts to 0xff and therefore lands last; rust's
1199		// natural Option order puts none first, so Desc must be what reconciles the two
1200		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		// the odometer must carry left, otherwise an exclusive upper end skips every row that sorts
1216		// between a key and the value successor hands back
1217		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		// the partition is the outermost column, so it may only step once every inner column is exhausted
1252		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}