reifydb_core/interface/catalog/
series.rs1use 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}