Skip to main content

reifydb_core/interface/catalog/
series.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_codec::row::pod::EncodedPodRow;
5use reifydb_value::{
6	Result,
7	value::{Value, datetime::DateTime, sumtype::SumTypeId, value_type::ValueType},
8};
9use serde::{Deserialize, Serialize};
10
11use crate::{
12	common::TimeSource,
13	interface::catalog::{
14		column::Column,
15		id::{NamespaceId, SeriesId},
16		key::PrimaryKey,
17	},
18	return_internal_error,
19	value::column::{buffer::ColumnBuffer, columns::Columns},
20};
21
22#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
23#[serde(rename_all = "lowercase")]
24#[derive(Default)]
25pub enum TimestampPrecision {
26	#[default]
27	Millisecond = 0,
28	Microsecond = 1,
29	Nanosecond = 2,
30	Second = 3,
31}
32
33#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
34pub enum SeriesKey {
35	DateTime {
36		column: String,
37		precision: TimestampPrecision,
38	},
39	Integer {
40		column: String,
41	},
42}
43
44impl SeriesKey {
45	pub fn column(&self) -> &str {
46		match self {
47			SeriesKey::DateTime {
48				column,
49				..
50			} => column,
51			SeriesKey::Integer {
52				column,
53			} => column,
54		}
55	}
56
57	pub fn extract_key(&self, columns: &Columns, row_idx: usize) -> Option<u64> {
58		let key_column = self.column();
59		columns.iter()
60			.find(|col| col.name().text() == key_column)
61			.and_then(|col| self.key_to_u64(col.data().get_value(row_idx)))
62	}
63
64	pub fn key_to_u64(&self, value: Value) -> Option<u64> {
65		match value {
66			Value::Int1(v) => u64::try_from(v).ok(),
67			Value::Int2(v) => u64::try_from(v).ok(),
68			Value::Int4(v) => u64::try_from(v).ok(),
69			Value::Int8(v) => u64::try_from(v).ok(),
70			Value::Int16(v) => u64::try_from(v).ok(),
71			Value::Uint1(v) => Some(v as u64),
72			Value::Uint2(v) => Some(v as u64),
73			Value::Uint4(v) => Some(v as u64),
74			Value::Uint8(v) => Some(v),
75			Value::Uint16(v) => u64::try_from(v).ok(),
76			Value::DateTime(dt) => {
77				let nanos = dt.to_nanos();
78				match self {
79					SeriesKey::DateTime {
80						precision,
81						..
82					} => Some(match precision {
83						TimestampPrecision::Second => nanos / 1_000_000_000,
84						TimestampPrecision::Millisecond => nanos / 1_000_000,
85						TimestampPrecision::Microsecond => nanos / 1_000,
86						TimestampPrecision::Nanosecond => nanos,
87					}),
88					_ => Some(nanos),
89				}
90			}
91			_ => None,
92		}
93	}
94
95	pub fn key_from_u64(&self, v: u64, key_type: Option<ValueType>) -> Value {
96		match key_type.as_ref() {
97			Some(ValueType::Int1) => Value::Int1(v as i8),
98			Some(ValueType::Int2) => Value::Int2(v as i16),
99			Some(ValueType::Int4) => Value::Int4(v as i32),
100			Some(ValueType::Int8) => Value::Int8(v as i64),
101			Some(ValueType::Uint1) => Value::Uint1(v as u8),
102			Some(ValueType::Uint2) => Value::Uint2(v as u16),
103			Some(ValueType::Uint4) => Value::Uint4(v as u32),
104			Some(ValueType::Uint8) => Value::Uint8(v),
105			Some(ValueType::Uint16) => Value::Uint16(v as u128),
106			Some(ValueType::Int16) => Value::Int16(v as i128),
107			Some(ValueType::DateTime) => {
108				let nanos: u64 = match self {
109					SeriesKey::DateTime {
110						precision,
111						..
112					} => match precision {
113						TimestampPrecision::Second => v * 1_000_000_000,
114						TimestampPrecision::Millisecond => v * 1_000_000,
115						TimestampPrecision::Microsecond => v * 1_000,
116						TimestampPrecision::Nanosecond => v,
117					},
118					_ => v,
119				};
120				Value::DateTime(DateTime::from_nanos(nanos))
121			}
122			_ => Value::Uint8(v),
123		}
124	}
125
126	pub fn decode(key_kind: u8, precision_raw: u8, column: String) -> Self {
127		match key_kind {
128			1 => SeriesKey::Integer {
129				column,
130			},
131			_ => {
132				let precision = match precision_raw {
133					1 => TimestampPrecision::Microsecond,
134					2 => TimestampPrecision::Nanosecond,
135					3 => TimestampPrecision::Second,
136					_ => TimestampPrecision::Millisecond,
137				};
138				SeriesKey::DateTime {
139					column,
140					precision,
141				}
142			}
143		}
144	}
145}
146
147#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
148pub struct Series {
149	pub id: SeriesId,
150	pub namespace: NamespaceId,
151	pub name: String,
152	pub columns: Vec<Column>,
153	pub tag: Option<SumTypeId>,
154	pub key: SeriesKey,
155	pub primary_key: Option<PrimaryKey>,
156	pub partition_by: Vec<String>,
157	pub time: TimeSource,
158}
159
160impl Series {
161	pub fn name(&self) -> &str {
162		&self.name
163	}
164
165	pub fn key_column_type(&self) -> Option<ValueType> {
166		let key_col_name = self.key.column();
167		self.columns.iter().find(|c| c.name == key_col_name).map(|c| c.constraint.get_type())
168	}
169
170	pub fn key_to_u64(&self, value: Value) -> Option<u64> {
171		self.key.key_to_u64(value)
172	}
173
174	pub fn key_from_u64(&self, v: u64) -> Value {
175		self.key.key_from_u64(v, self.key_column_type())
176	}
177
178	pub fn key_column_data(&self, keys: Vec<u64>) -> ColumnBuffer {
179		let key_type = self.key_column_type();
180		match &key_type {
181			Some(ty) => {
182				let mut data = ColumnBuffer::with_capacity(ty.clone(), keys.len());
183				for k in keys {
184					data.push_value(self.key_from_u64(k));
185				}
186				data
187			}
188			None => ColumnBuffer::uint8(keys),
189		}
190	}
191
192	pub fn data_columns(&self) -> impl Iterator<Item = &Column> {
193		let key_column = self.key.column().to_string();
194		self.columns.iter().filter(move |c| c.name != key_column)
195	}
196}
197
198#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
199pub struct SeriesMetadata {
200	pub row_count: u64,
201	pub oldest_key: u64,
202	pub newest_key: u64,
203	pub sequence_counter: u64,
204}
205
206impl SeriesMetadata {
207	pub fn new() -> Self {
208		Self {
209			row_count: 0,
210			oldest_key: 0,
211			newest_key: 0,
212			sequence_counter: 0,
213		}
214	}
215}
216
217impl Default for SeriesMetadata {
218	fn default() -> Self {
219		Self::new()
220	}
221}
222
223const SERIES_METADATA_WIDTH: usize = 32;
224
225pub fn encode_series_metadata(metadata: &SeriesMetadata) -> EncodedPodRow {
226	let mut bytes = Vec::with_capacity(SERIES_METADATA_WIDTH);
227	bytes.extend_from_slice(&metadata.row_count.to_be_bytes());
228	bytes.extend_from_slice(&metadata.oldest_key.to_be_bytes());
229	bytes.extend_from_slice(&metadata.newest_key.to_be_bytes());
230	bytes.extend_from_slice(&metadata.sequence_counter.to_be_bytes());
231	EncodedPodRow::new(&bytes)
232}
233
234pub fn decode_series_metadata(row: &EncodedPodRow) -> Result<SeriesMetadata> {
235	let bytes = row.body();
236	if bytes.len() != SERIES_METADATA_WIDTH {
237		return_internal_error!(
238			"Series metadata is {} bytes wide, expected {}. This indicates a corrupt metadata row.",
239			bytes.len(),
240			SERIES_METADATA_WIDTH
241		)
242	}
243	Ok(SeriesMetadata {
244		row_count: u64::from_be_bytes(bytes[0..8].try_into().unwrap()),
245		oldest_key: u64::from_be_bytes(bytes[8..16].try_into().unwrap()),
246		newest_key: u64::from_be_bytes(bytes[16..24].try_into().unwrap()),
247		sequence_counter: u64::from_be_bytes(bytes[24..32].try_into().unwrap()),
248	})
249}
250
251#[cfg(test)]
252mod series_metadata_tests {
253	use super::*;
254
255	#[test]
256	fn every_field_survives_a_round_trip_at_the_declared_width() {
257		let metadata = SeriesMetadata {
258			row_count: 42,
259			oldest_key: 100,
260			newest_key: 900,
261			sequence_counter: 7,
262		};
263
264		let row = encode_series_metadata(&metadata);
265
266		assert_eq!(row.len(), SERIES_METADATA_WIDTH);
267		assert_eq!(decode_series_metadata(&row).unwrap(), metadata);
268	}
269
270	#[test]
271	fn the_key_bounds_do_not_swap_because_they_select_which_buckets_materialise() {
272		let metadata = SeriesMetadata {
273			row_count: 1,
274			oldest_key: 1,
275			newest_key: u64::MAX,
276			sequence_counter: 0,
277		};
278
279		let decoded = decode_series_metadata(&encode_series_metadata(&metadata)).unwrap();
280
281		assert_eq!(decoded.oldest_key, 1);
282		assert_eq!(decoded.newest_key, u64::MAX);
283	}
284
285	#[test]
286	fn a_row_of_the_wrong_width_is_rejected_rather_than_rewinding_the_sequence_counter() {
287		assert!(decode_series_metadata(&EncodedPodRow::new(&[0u8; 31])).is_err());
288		assert!(decode_series_metadata(&EncodedPodRow::new(&[0u8; 33])).is_err());
289		assert!(decode_series_metadata(&EncodedPodRow::new(&[0u8; 40])).is_err());
290	}
291}