Skip to main content

mito_codec/row_converter/
sparse.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{HashMap, HashSet};
16use std::sync::Arc;
17
18use bytes::BufMut;
19use common_recordbatch::filter::SimpleFilterEvaluator;
20use datatypes::prelude::ConcreteDataType;
21use datatypes::value::{Value, ValueRef};
22use memcomparable::{Deserializer, Serializer};
23use serde::{Deserialize, Serialize};
24use snafu::{ResultExt, ensure};
25use store_api::codec::PrimaryKeyEncoding;
26use store_api::metadata::RegionMetadataRef;
27use store_api::storage::ColumnId;
28use store_api::storage::consts::ReservedColumnId;
29
30use crate::error::{
31    DeserializeFieldSnafu, InvalidSparsePrimaryKeySnafu, Result, SerializeFieldSnafu,
32    UnsupportedOperationSnafu,
33};
34use crate::key_values::KeyValue;
35use crate::primary_key_filter::SparsePrimaryKeyFilter;
36use crate::row_converter::dense::SortField;
37use crate::row_converter::{CompositeValues, PrimaryKeyCodec, PrimaryKeyFilter};
38
39/// A codec for sparse key of metrics.
40///
41/// ## Encoding format
42/// Each primary key is encoded as a sequence of `(column_id, value)` pairs:
43/// - The first two fields are always the reserved `table_id` (uint32) and `tsid` (uint64).
44/// - User-defined labels follow, sorted by **column name** in lexicographical order.
45/// - Null values are omitted (not encoded).
46///
47/// The `column_id` is encoded as a 4-byte big-endian integer, and the value is encoded
48/// using memcomparable serialization.
49///
50/// `decode_leftmost` always decodes the first value from the encoded bytes (i.e., the
51/// `table_id` field).
52///
53/// ## Requirements
54/// It requires the input primary key columns are sorted by the column name in lexicographical order.
55/// It encodes the column id of the physical region.
56#[derive(Clone, Debug)]
57pub struct SparsePrimaryKeyCodec {
58    inner: Arc<SparsePrimaryKeyCodecInner>,
59}
60
61#[derive(Debug)]
62struct SparsePrimaryKeyCodecInner {
63    // Internal fields
64    table_id_field: SortField,
65    // Internal fields
66    tsid_field: SortField,
67    // User defined label field
68    label_field: SortField,
69    // Columns in primary key
70    //
71    // None means all unknown columns is primary key(`Self::label_field`).
72    columns: Option<HashSet<ColumnId>>,
73}
74
75/// Sparse values representation.
76///
77/// Callers must not insert a column id that is already present; otherwise
78/// the existing entry will shadow the newly inserted value on lookup.
79#[derive(Debug, Clone, PartialEq, Eq, Default)]
80pub struct SparseValues {
81    values: Vec<(ColumnId, Value)>,
82}
83
84impl SparseValues {
85    /// Creates an empty [`SparseValues`].
86    pub fn new() -> Self {
87        Self { values: Vec::new() }
88    }
89
90    /// Creates an empty [`SparseValues`] with space reserved for `cap` entries.
91    pub fn with_capacity(cap: usize) -> Self {
92        Self {
93            values: Vec::with_capacity(cap),
94        }
95    }
96
97    /// Returns the value of the given column, or [`Value::Null`] if the column is not present.
98    pub fn get_or_null(&self, column_id: ColumnId) -> &Value {
99        for (id, value) in &self.values {
100            if *id == column_id {
101                return value;
102            }
103        }
104        &Value::Null
105    }
106
107    /// Returns the value of the given column, or [`None`] if the column is not present.
108    pub fn get(&self, column_id: &ColumnId) -> Option<&Value> {
109        for (id, value) in &self.values {
110            if id == column_id {
111                return Some(value);
112            }
113        }
114        None
115    }
116
117    /// Appends a new `(column_id, value)` pair.
118    ///
119    /// Append-only: the caller must ensure `column_id` is not already present.
120    pub fn insert(&mut self, column_id: ColumnId, value: Value) {
121        self.values.push((column_id, value));
122    }
123
124    /// Returns an iterator over all stored column id/value pairs.
125    pub fn iter(&self) -> impl Iterator<Item = (&ColumnId, &Value)> {
126        self.values.iter().map(|(id, value)| (id, value))
127    }
128}
129
130/// The column id of the tsid.
131pub const RESERVED_COLUMN_ID_TSID: ColumnId = ReservedColumnId::tsid();
132/// The column id of the table id.
133pub const RESERVED_COLUMN_ID_TABLE_ID: ColumnId = ReservedColumnId::table_id();
134/// The size of the column id in the encoded sparse row.
135pub const COLUMN_ID_ENCODE_SIZE: usize = 4;
136
137// Fixed byte offsets for reserved columns in the sparse encoding.
138// Layout: [table_id_col_id: 4B][marker: 1B][table_id: 4B][tsid_col_id: 4B][marker: 1B][tsid: 8B]
139/// Byte offset to the table_id value (after its 4-byte column id).
140const TABLE_ID_VALUE_OFFSET: usize = COLUMN_ID_ENCODE_SIZE;
141/// Byte offset to the tsid value (after 9-byte table_id entry + 4-byte tsid column id).
142const TSID_VALUE_OFFSET: usize = COLUMN_ID_ENCODE_SIZE + 5 + COLUMN_ID_ENCODE_SIZE;
143/// Byte offset where tag columns start (after 9-byte table_id + 13-byte tsid entries).
144const TAGS_START_OFFSET: usize = COLUMN_ID_ENCODE_SIZE + 5 + COLUMN_ID_ENCODE_SIZE + 9;
145
146/// Inline capacity for the small-vec fast path of [`SparseOffsetsCache`].
147///
148/// Most sparse primary keys carry only a handful of tags, so a linear scan
149/// over a short `Vec` beats a `HashMap` lookup. Tags beyond this capacity
150/// spill into the overflow `HashMap`.
151const SPARSE_OFFSETS_INLINE_CAP: usize = 32;
152
153/// A lazily populated cache of tag column offsets inside a sparse primary key.
154#[derive(Debug, Clone)]
155pub struct SparseOffsetsCache {
156    /// Small-vec fast path. Reserves [`SPARSE_OFFSETS_INLINE_CAP`] slots on
157    /// the first insert.
158    inline: Vec<(ColumnId, usize)>,
159    /// Overflow for columns beyond the inline capacity. Lazily allocated.
160    overflow: HashMap<ColumnId, usize>,
161    /// Next byte position in the pk to resume parsing from.
162    cursor: usize,
163    /// True once the decoder has walked past the last tag column (or stopped
164    /// on an unknown column id); no further offsets can be discovered.
165    finished: bool,
166}
167
168impl Default for SparseOffsetsCache {
169    fn default() -> Self {
170        Self::new()
171    }
172}
173
174impl SparseOffsetsCache {
175    pub fn new() -> Self {
176        Self {
177            inline: Vec::new(),
178            overflow: HashMap::new(),
179            cursor: TAGS_START_OFFSET,
180            finished: false,
181        }
182    }
183
184    pub fn clear(&mut self) {
185        self.inline.clear();
186        self.overflow.clear();
187        self.cursor = TAGS_START_OFFSET;
188        self.finished = false;
189    }
190
191    /// Returns the cached offset for `column_id`, if any.
192    fn get(&self, column_id: ColumnId) -> Option<usize> {
193        for entry in &self.inline {
194            if entry.0 == column_id {
195                return Some(entry.1);
196            }
197        }
198        self.overflow.get(&column_id).copied()
199    }
200
201    /// Records a new `(column_id, offset)` entry.
202    fn insert(&mut self, column_id: ColumnId, offset: usize) {
203        if self.inline.len() < SPARSE_OFFSETS_INLINE_CAP {
204            if self.inline.capacity() == 0 {
205                self.inline.reserve_exact(SPARSE_OFFSETS_INLINE_CAP);
206            }
207            self.inline.push((column_id, offset));
208        } else {
209            self.overflow.insert(column_id, offset);
210        }
211    }
212
213    #[cfg(test)]
214    fn contains(&self, column_id: ColumnId) -> bool {
215        self.get(column_id).is_some()
216    }
217}
218
219impl SparsePrimaryKeyCodec {
220    /// Creates a new [`SparsePrimaryKeyCodec`] instance.
221    pub fn from_columns(columns_ids: impl Iterator<Item = ColumnId>) -> Self {
222        let columns = columns_ids.collect();
223        Self {
224            inner: Arc::new(SparsePrimaryKeyCodecInner {
225                table_id_field: SortField::new(ConcreteDataType::uint32_datatype()),
226                tsid_field: SortField::new(ConcreteDataType::uint64_datatype()),
227                label_field: SortField::new(ConcreteDataType::string_datatype()),
228                columns: Some(columns),
229            }),
230        }
231    }
232
233    /// Creates a new [`SparsePrimaryKeyCodec`] instance.
234    pub fn new(region_metadata: &RegionMetadataRef) -> Self {
235        Self::from_columns(region_metadata.primary_key_columns().map(|c| c.column_id))
236    }
237
238    /// Returns a new [`SparsePrimaryKeyCodec`] instance.
239    ///
240    /// It treats all unknown columns as primary key(label field).
241    pub fn schemaless() -> Self {
242        Self {
243            inner: Arc::new(SparsePrimaryKeyCodecInner {
244                table_id_field: SortField::new(ConcreteDataType::uint32_datatype()),
245                tsid_field: SortField::new(ConcreteDataType::uint64_datatype()),
246                label_field: SortField::new(ConcreteDataType::string_datatype()),
247                columns: None,
248            }),
249        }
250    }
251
252    /// Creates a new [`SparsePrimaryKeyCodec`] instance with additional label `fields`.
253    pub fn with_fields(fields: Vec<(ColumnId, SortField)>) -> Self {
254        Self {
255            inner: Arc::new(SparsePrimaryKeyCodecInner {
256                columns: Some(fields.iter().map(|f| f.0).collect()),
257                table_id_field: SortField::new(ConcreteDataType::uint32_datatype()),
258                tsid_field: SortField::new(ConcreteDataType::uint64_datatype()),
259                label_field: SortField::new(ConcreteDataType::string_datatype()),
260            }),
261        }
262    }
263
264    /// Returns the field of the given column id.
265    fn get_field(&self, column_id: ColumnId) -> Option<&SortField> {
266        // if the `columns` is not specified, all unknown columns is primary key(label field).
267        if let Some(columns) = &self.inner.columns
268            && !columns.contains(&column_id)
269        {
270            return None;
271        }
272
273        match column_id {
274            RESERVED_COLUMN_ID_TABLE_ID => Some(&self.inner.table_id_field),
275            RESERVED_COLUMN_ID_TSID => Some(&self.inner.tsid_field),
276            _ => Some(&self.inner.label_field),
277        }
278    }
279
280    /// Encodes the given bytes into a [`SparseValues`].
281    pub fn encode_to_vec<'a, I>(&self, row: I, buffer: &mut Vec<u8>) -> Result<()>
282    where
283        I: Iterator<Item = (ColumnId, ValueRef<'a>)>,
284    {
285        let mut serializer = Serializer::new(buffer);
286        for (column_id, value) in row {
287            if value.is_null() {
288                continue;
289            }
290
291            if let Some(field) = self.get_field(column_id) {
292                column_id
293                    .serialize(&mut serializer)
294                    .context(SerializeFieldSnafu)?;
295                field.serialize(&mut serializer, &value)?;
296            } else {
297                // TODO(weny): handle the error.
298                common_telemetry::warn!("Column {} is not in primary key, skipping", column_id);
299            }
300        }
301        Ok(())
302    }
303
304    pub fn encode_raw_tag_value<'a, I>(&self, row: I, buffer: &mut Vec<u8>) -> Result<()>
305    where
306        I: Iterator<Item = (ColumnId, &'a [u8])>,
307    {
308        for (tag_column_id, tag_value) in row {
309            let value_len = tag_value.len();
310            buffer.reserve(6 + value_len / 8 * 9);
311            buffer.put_u32(tag_column_id);
312            buffer.put_u8(1);
313            buffer.put_u8(!tag_value.is_empty() as u8);
314
315            // Manual implementation of memcomparable::ser::Serializer::serialize_bytes
316            // to avoid byte-by-byte put.
317            let mut len = 0;
318            let num_chucks = value_len / 8;
319            let remainder = value_len % 8;
320
321            for idx in 0..num_chucks {
322                buffer.extend_from_slice(&tag_value[idx * 8..idx * 8 + 8]);
323                len += 8;
324                // append an extra byte that signals the number of significant bytes in this chunk
325                // 1-8: many bytes were significant and this group is the last group
326                // 9: all 8 bytes were significant and there is more data to come
327                let extra = if len == value_len { 8 } else { 9 };
328                buffer.put_u8(extra);
329            }
330
331            if remainder != 0 {
332                buffer.extend_from_slice(&tag_value[len..value_len]);
333                buffer.put_bytes(0, 8 - remainder);
334                buffer.put_u8(remainder as u8);
335            }
336        }
337        Ok(())
338    }
339
340    /// Encodes the given bytes into a [`SparseValues`].
341    pub fn encode_internal(&self, table_id: u32, tsid: u64, buffer: &mut Vec<u8>) -> Result<()> {
342        buffer.reserve_exact(22);
343        buffer.put_u32(RESERVED_COLUMN_ID_TABLE_ID);
344        buffer.put_u8(1);
345        buffer.put_u32(table_id);
346        buffer.put_u32(RESERVED_COLUMN_ID_TSID);
347        buffer.put_u8(1);
348        buffer.put_u64(tsid);
349        Ok(())
350    }
351
352    /// Decodes the reserved `(table_id, tsid)` prefix from a sparse primary key.
353    pub fn decode_ids(&self, bytes: &[u8]) -> Result<(u32, u64)> {
354        // Two column IDs, two non-null markers, a u32 table ID, and a u64 TSID.
355        const INTERNAL_PREFIX_LEN: usize = 4 + 1 + 4 + 4 + 1 + 8;
356        ensure!(
357            bytes.len() >= INTERNAL_PREFIX_LEN,
358            InvalidSparsePrimaryKeySnafu {
359                reason: format!(
360                    "internal prefix requires at least {INTERNAL_PREFIX_LEN} bytes, got {}",
361                    bytes.len()
362                ),
363            }
364        );
365
366        let mut deserializer = Deserializer::new(bytes);
367
368        let table_id_column = u32::deserialize(&mut deserializer).context(DeserializeFieldSnafu)?;
369        ensure!(
370            table_id_column == RESERVED_COLUMN_ID_TABLE_ID,
371            InvalidSparsePrimaryKeySnafu {
372                reason: format!(
373                    "expected table id column {}, got {}",
374                    RESERVED_COLUMN_ID_TABLE_ID, table_id_column
375                ),
376            }
377        );
378        let table_id = self.inner.table_id_field.deserialize(&mut deserializer)?;
379
380        let tsid_column = u32::deserialize(&mut deserializer).context(DeserializeFieldSnafu)?;
381        ensure!(
382            tsid_column == RESERVED_COLUMN_ID_TSID,
383            InvalidSparsePrimaryKeySnafu {
384                reason: format!(
385                    "expected tsid column {}, got {}",
386                    RESERVED_COLUMN_ID_TSID, tsid_column
387                ),
388            }
389        );
390        let tsid = self.inner.tsid_field.deserialize(&mut deserializer)?;
391
392        match (table_id, tsid) {
393            (Value::UInt32(table_id), Value::UInt64(tsid)) => Ok((table_id, tsid)),
394            (table_id, tsid) => InvalidSparsePrimaryKeySnafu {
395                reason: format!(
396                    "expected UInt32 table id and UInt64 tsid, got {table_id:?} and {tsid:?}"
397                ),
398            }
399            .fail(),
400        }
401    }
402
403    /// Decodes the given bytes into a [`SparseValues`].
404    fn decode_sparse(&self, bytes: &[u8]) -> Result<SparseValues> {
405        let mut deserializer = Deserializer::new(bytes);
406        let mut values = SparseValues::with_capacity(16);
407
408        let column_id = u32::deserialize(&mut deserializer).context(DeserializeFieldSnafu)?;
409        let value = self.inner.table_id_field.deserialize(&mut deserializer)?;
410        values.insert(column_id, value);
411
412        let column_id = u32::deserialize(&mut deserializer).context(DeserializeFieldSnafu)?;
413        let value = self.inner.tsid_field.deserialize(&mut deserializer)?;
414        values.insert(column_id, value);
415        while deserializer.has_remaining() {
416            let column_id = u32::deserialize(&mut deserializer).context(DeserializeFieldSnafu)?;
417            let value = self.inner.label_field.deserialize(&mut deserializer)?;
418            values.insert(column_id, value);
419        }
420
421        Ok(values)
422    }
423
424    /// Decodes the given bytes into a [`Value`].
425    fn decode_leftmost(&self, bytes: &[u8]) -> Result<Option<Value>> {
426        let mut deserializer = Deserializer::new(bytes);
427        // Skip the column id.
428        deserializer.advance(COLUMN_ID_ENCODE_SIZE);
429        let value = self.inner.table_id_field.deserialize(&mut deserializer)?;
430        Ok(Some(value))
431    }
432
433    /// Returns the offset of the given column id in the given primary key.
434    ///
435    /// The pk must start with the table_id + tsid prefix written by
436    /// `encode_internal`.
437    ///
438    /// # Panics
439    ///
440    /// Panics if `pk` is not a well-formed sparse primary key produced by
441    /// this codec (e.g. truncated or otherwise malformed bytes).
442    pub fn has_column(
443        &self,
444        pk: &[u8],
445        cache: &mut SparseOffsetsCache,
446        column_id: ColumnId,
447    ) -> Option<usize> {
448        // Decoding is lazy: on each call we only advance the cache's cursor as
449        // far as needed to answer the query. A column that has already been
450        // seen returns immediately; a column we haven't reached yet causes the
451        // parser to resume from `cache.cursor` and stop as soon as the column
452        // is located. Once the cursor walks off the end (or hits an unknown
453        // column id) the cache is marked finished, so subsequent misses are
454        // O(1).
455        // table_id and tsid are at fixed offsets.
456        match column_id {
457            RESERVED_COLUMN_ID_TABLE_ID => return Some(TABLE_ID_VALUE_OFFSET),
458            RESERVED_COLUMN_ID_TSID => return Some(TSID_VALUE_OFFSET),
459            _ => {}
460        }
461
462        if let Some(offset) = cache.get(column_id) {
463            return Some(offset);
464        }
465        if cache.finished {
466            return None;
467        }
468
469        let mut deserializer = Deserializer::new(pk);
470        deserializer.advance(cache.cursor);
471        let mut offset = cache.cursor;
472        while deserializer.has_remaining() {
473            let col = u32::deserialize(&mut deserializer).unwrap();
474            offset += COLUMN_ID_ENCODE_SIZE;
475            let value_offset = offset;
476            cache.insert(col, value_offset);
477            let Some(field) = self.get_field(col) else {
478                cache.finished = true;
479                cache.cursor = offset;
480                return None;
481            };
482
483            let skip = field.skip_deserialize(pk, &mut deserializer).unwrap();
484            offset += skip;
485            cache.cursor = offset;
486            if col == column_id {
487                return Some(value_offset);
488            }
489        }
490
491        cache.finished = true;
492        None
493    }
494
495    /// Decode value at `offset` in `pk`.
496    pub fn decode_value_at(&self, pk: &[u8], offset: usize, column_id: ColumnId) -> Result<Value> {
497        let mut deserializer = Deserializer::new(pk);
498        deserializer.advance(offset);
499        // Safety: checked by `has_column`
500        let field = self.get_field(column_id).unwrap();
501        field.deserialize(&mut deserializer)
502    }
503
504    /// Returns the encoded bytes of the given `column_id` in `pk`.
505    ///
506    /// Returns `Ok(None)` if the `column_id` is missing in `pk`.
507    pub fn encoded_value_for_column<'a>(
508        &self,
509        pk: &'a [u8],
510        cache: &mut SparseOffsetsCache,
511        column_id: ColumnId,
512    ) -> Result<Option<&'a [u8]>> {
513        let Some(offset) = self.has_column(pk, cache, column_id) else {
514            return Ok(None);
515        };
516
517        let Some(field) = self.get_field(column_id) else {
518            return Ok(None);
519        };
520
521        let mut deserializer = Deserializer::new(pk);
522        deserializer.advance(offset);
523        let len = field.skip_deserialize(pk, &mut deserializer)?;
524        Ok(Some(&pk[offset..offset + len]))
525    }
526}
527
528impl PrimaryKeyCodec for SparsePrimaryKeyCodec {
529    fn encode_key_value(&self, _key_value: &KeyValue, _buffer: &mut Vec<u8>) -> Result<()> {
530        UnsupportedOperationSnafu {
531            err_msg: "The encode_key_value method is not supported in SparsePrimaryKeyCodec.",
532        }
533        .fail()
534    }
535
536    fn encode_values(&self, values: &[(ColumnId, Value)], buffer: &mut Vec<u8>) -> Result<()> {
537        self.encode_to_vec(values.iter().map(|v| (v.0, v.1.as_value_ref())), buffer)
538    }
539
540    fn encode_value_refs(
541        &self,
542        values: &[(ColumnId, ValueRef)],
543        buffer: &mut Vec<u8>,
544    ) -> Result<()> {
545        self.encode_to_vec(values.iter().map(|v| (v.0, v.1.clone())), buffer)
546    }
547
548    fn estimated_size(&self) -> Option<usize> {
549        None
550    }
551
552    fn num_fields(&self) -> Option<usize> {
553        None
554    }
555
556    fn encoding(&self) -> PrimaryKeyEncoding {
557        PrimaryKeyEncoding::Sparse
558    }
559
560    fn primary_key_filter(
561        &self,
562        metadata: &RegionMetadataRef,
563        filters: Arc<Vec<SimpleFilterEvaluator>>,
564    ) -> Box<dyn PrimaryKeyFilter> {
565        Box::new(SparsePrimaryKeyFilter::new(
566            metadata.clone(),
567            filters,
568            self.clone(),
569        ))
570    }
571
572    fn decode(&self, bytes: &[u8]) -> Result<CompositeValues> {
573        Ok(CompositeValues::Sparse(self.decode_sparse(bytes)?))
574    }
575
576    fn decode_leftmost(&self, bytes: &[u8]) -> Result<Option<Value>> {
577        self.decode_leftmost(bytes)
578    }
579}
580
581/// Field with column id.
582pub struct FieldWithId {
583    pub field: SortField,
584    pub column_id: ColumnId,
585}
586
587/// A special encoder for memtable.
588pub struct SparseEncoder {
589    fields: Vec<FieldWithId>,
590}
591
592impl SparseEncoder {
593    pub fn new(fields: Vec<FieldWithId>) -> Self {
594        Self { fields }
595    }
596
597    pub fn encode_to_vec<'a, I>(&self, row: I, buffer: &mut Vec<u8>) -> Result<()>
598    where
599        I: Iterator<Item = ValueRef<'a>>,
600    {
601        let mut serializer = Serializer::new(buffer);
602        for (value, field) in row.zip(self.fields.iter()) {
603            if !value.is_null() {
604                field
605                    .column_id
606                    .serialize(&mut serializer)
607                    .context(SerializeFieldSnafu)?;
608                field.field.serialize(&mut serializer, &value)?;
609            }
610        }
611        Ok(())
612    }
613}
614
615#[cfg(test)]
616mod tests {
617    use std::sync::Arc;
618
619    use api::v1::SemanticType;
620    use common_query::prelude::{greptime_timestamp, greptime_value};
621    use common_time::Timestamp;
622    use common_time::timestamp::TimeUnit;
623    use datatypes::schema::ColumnSchema;
624    use datatypes::value::{OrderedFloat, Value};
625    use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
626    use store_api::metric_engine_consts::{
627        DATA_SCHEMA_TABLE_ID_COLUMN_NAME, DATA_SCHEMA_TSID_COLUMN_NAME,
628    };
629    use store_api::storage::{ColumnId, RegionId};
630
631    use super::*;
632
633    fn test_region_metadata() -> RegionMetadataRef {
634        let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1));
635        builder
636            .push_column_metadata(ColumnMetadata {
637                column_schema: ColumnSchema::new(
638                    DATA_SCHEMA_TABLE_ID_COLUMN_NAME,
639                    ConcreteDataType::uint32_datatype(),
640                    false,
641                ),
642                semantic_type: SemanticType::Tag,
643                column_id: ReservedColumnId::table_id(),
644            })
645            .push_column_metadata(ColumnMetadata {
646                column_schema: ColumnSchema::new(
647                    DATA_SCHEMA_TSID_COLUMN_NAME,
648                    ConcreteDataType::uint64_datatype(),
649                    false,
650                ),
651                semantic_type: SemanticType::Tag,
652                column_id: ReservedColumnId::tsid(),
653            })
654            .push_column_metadata(ColumnMetadata {
655                column_schema: ColumnSchema::new("pod", ConcreteDataType::string_datatype(), true),
656                semantic_type: SemanticType::Tag,
657                column_id: 1,
658            })
659            .push_column_metadata(ColumnMetadata {
660                column_schema: ColumnSchema::new(
661                    "namespace",
662                    ConcreteDataType::string_datatype(),
663                    true,
664                ),
665                semantic_type: SemanticType::Tag,
666                column_id: 2,
667            })
668            .push_column_metadata(ColumnMetadata {
669                column_schema: ColumnSchema::new(
670                    "container",
671                    ConcreteDataType::string_datatype(),
672                    true,
673                ),
674                semantic_type: SemanticType::Tag,
675                column_id: 3,
676            })
677            .push_column_metadata(ColumnMetadata {
678                column_schema: ColumnSchema::new(
679                    "pod_name",
680                    ConcreteDataType::string_datatype(),
681                    true,
682                ),
683                semantic_type: SemanticType::Tag,
684                column_id: 4,
685            })
686            .push_column_metadata(ColumnMetadata {
687                column_schema: ColumnSchema::new(
688                    "pod_ip",
689                    ConcreteDataType::string_datatype(),
690                    true,
691                ),
692                semantic_type: SemanticType::Tag,
693                column_id: 5,
694            })
695            .push_column_metadata(ColumnMetadata {
696                column_schema: ColumnSchema::new(
697                    greptime_value(),
698                    ConcreteDataType::float64_datatype(),
699                    false,
700                ),
701                semantic_type: SemanticType::Field,
702                column_id: 6,
703            })
704            .push_column_metadata(ColumnMetadata {
705                column_schema: ColumnSchema::new(
706                    greptime_timestamp(),
707                    ConcreteDataType::timestamp_nanosecond_datatype(),
708                    false,
709                ),
710                semantic_type: SemanticType::Timestamp,
711                column_id: 7,
712            })
713            .primary_key(vec![
714                ReservedColumnId::table_id(),
715                ReservedColumnId::tsid(),
716                1,
717                2,
718                3,
719                4,
720                5,
721            ]);
722        let metadata = builder.build().unwrap();
723        Arc::new(metadata)
724    }
725
726    #[test]
727    fn test_sparse_value_new_and_get_or_null() {
728        let mut sparse_value = SparseValues::new();
729        sparse_value.insert(1, Value::Int32(42));
730
731        assert_eq!(sparse_value.get_or_null(1), &Value::Int32(42));
732        assert_eq!(sparse_value.get_or_null(2), &Value::Null);
733    }
734
735    #[test]
736    fn test_sparse_value_insert() {
737        let mut sparse_value = SparseValues::new();
738        sparse_value.insert(1, Value::Int32(42));
739
740        assert_eq!(sparse_value.get_or_null(1), &Value::Int32(42));
741    }
742
743    fn test_row() -> Vec<(ColumnId, ValueRef<'static>)> {
744        vec![
745            (RESERVED_COLUMN_ID_TABLE_ID, ValueRef::UInt32(42)),
746            (
747                RESERVED_COLUMN_ID_TSID,
748                ValueRef::UInt64(123843349035232323),
749            ),
750            // label: pod
751            (1, ValueRef::String("greptime-frontend-6989d9899-22222")),
752            // label: namespace
753            (2, ValueRef::String("greptime-cluster")),
754            // label: container
755            (3, ValueRef::String("greptime-frontend-6989d9899-22222")),
756            // label: pod_name
757            (4, ValueRef::String("greptime-frontend-6989d9899-22222")),
758            // label: pod_ip
759            (5, ValueRef::String("10.10.10.10")),
760            // field: greptime_value
761            (6, ValueRef::Float64(OrderedFloat(1.0))),
762            // field: greptime_timestamp
763            (
764                7,
765                ValueRef::Timestamp(Timestamp::new(1618876800000000000, TimeUnit::Nanosecond)),
766            ),
767        ]
768    }
769
770    #[test]
771    fn test_encode_by_short_cuts() {
772        let region_metadata = test_region_metadata();
773        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
774        let mut buffer = Vec::new();
775        let internal_columns = [
776            (RESERVED_COLUMN_ID_TABLE_ID, ValueRef::UInt32(1024)),
777            (RESERVED_COLUMN_ID_TSID, ValueRef::UInt64(42)),
778        ];
779        let tags = [
780            (1, "greptime-frontend-6989d9899-22222"),
781            (2, "greptime-cluster"),
782            (3, "greptime-frontend-6989d9899-22222"),
783            (4, "greptime-frontend-6989d9899-22222"),
784            (5, "10.10.10.10"),
785        ];
786        codec
787            .encode_to_vec(internal_columns.into_iter(), &mut buffer)
788            .unwrap();
789        codec
790            .encode_to_vec(
791                tags.iter()
792                    .map(|(col_id, tag_value)| (*col_id, ValueRef::String(tag_value))),
793                &mut buffer,
794            )
795            .unwrap();
796
797        let mut buffer_by_raw_encoding = Vec::new();
798        codec
799            .encode_internal(1024, 42, &mut buffer_by_raw_encoding)
800            .unwrap();
801        let tags: Vec<_> = tags
802            .into_iter()
803            .map(|(col_id, tag_value)| (col_id, tag_value.as_bytes()))
804            .collect();
805        codec
806            .encode_raw_tag_value(
807                tags.iter().map(|(c, b)| (*c, *b)),
808                &mut buffer_by_raw_encoding,
809            )
810            .unwrap();
811        assert_eq!(buffer, buffer_by_raw_encoding);
812    }
813
814    #[test]
815    fn test_encode_to_vec() {
816        let region_metadata = test_region_metadata();
817        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
818        let mut buffer = Vec::new();
819
820        let row = test_row();
821        codec.encode_to_vec(row.into_iter(), &mut buffer).unwrap();
822        assert!(!buffer.is_empty());
823        let sparse_value = codec.decode_sparse(&buffer).unwrap();
824        assert_eq!(
825            sparse_value.get_or_null(RESERVED_COLUMN_ID_TABLE_ID),
826            &Value::UInt32(42)
827        );
828        assert_eq!(
829            sparse_value.get_or_null(1),
830            &Value::String("greptime-frontend-6989d9899-22222".into())
831        );
832        assert_eq!(
833            sparse_value.get_or_null(2),
834            &Value::String("greptime-cluster".into())
835        );
836        assert_eq!(
837            sparse_value.get_or_null(3),
838            &Value::String("greptime-frontend-6989d9899-22222".into())
839        );
840        assert_eq!(
841            sparse_value.get_or_null(4),
842            &Value::String("greptime-frontend-6989d9899-22222".into())
843        );
844        assert_eq!(
845            sparse_value.get_or_null(5),
846            &Value::String("10.10.10.10".into())
847        );
848    }
849
850    #[test]
851    fn test_decode_leftmost() {
852        let region_metadata = test_region_metadata();
853        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
854        let mut buffer = Vec::new();
855        let row = test_row();
856        codec.encode_to_vec(row.into_iter(), &mut buffer).unwrap();
857        assert!(!buffer.is_empty());
858        let result = codec.decode_leftmost(&buffer).unwrap().unwrap();
859        assert_eq!(result, Value::UInt32(42));
860    }
861
862    #[test]
863    fn test_decode_ids() {
864        let region_metadata = test_region_metadata();
865        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
866        let mut buffer = Vec::new();
867        codec.encode_internal(42, 100, &mut buffer).unwrap();
868
869        assert_eq!((42, 100), codec.decode_ids(&buffer).unwrap());
870
871        let mut invalid = buffer.clone();
872        invalid[0..4].copy_from_slice(&1_u32.to_be_bytes());
873        assert!(codec.decode_ids(&invalid).is_err());
874        assert!(codec.decode_ids(&buffer[..buffer.len() - 1]).is_err());
875    }
876
877    #[test]
878    fn test_has_column() {
879        let region_metadata = test_region_metadata();
880        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
881        let mut buffer = Vec::new();
882        let row = test_row();
883        codec.encode_to_vec(row.into_iter(), &mut buffer).unwrap();
884        assert!(!buffer.is_empty());
885
886        let mut offsets_map = SparseOffsetsCache::new();
887        for column_id in [
888            RESERVED_COLUMN_ID_TABLE_ID,
889            RESERVED_COLUMN_ID_TSID,
890            1,
891            2,
892            3,
893            4,
894            5,
895        ] {
896            let offset = codec.has_column(&buffer, &mut offsets_map, column_id);
897            assert!(offset.is_some());
898        }
899
900        let offset = codec.has_column(&buffer, &mut offsets_map, 6);
901        assert!(offset.is_none());
902    }
903
904    #[test]
905    fn test_has_column_lazy_resume() {
906        let region_metadata = test_region_metadata();
907        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
908        let mut buffer = Vec::new();
909        codec
910            .encode_to_vec(test_row().into_iter(), &mut buffer)
911            .unwrap();
912
913        let mut cache = SparseOffsetsCache::new();
914        // Look up an early column: only a prefix of tags is decoded.
915        assert!(codec.has_column(&buffer, &mut cache, 1).is_some());
916        assert!(!cache.finished);
917        assert!(cache.contains(1));
918        assert!(!cache.contains(5));
919
920        // A later column resumes from the cursor.
921        assert!(codec.has_column(&buffer, &mut cache, 5).is_some());
922        assert!(cache.contains(5));
923
924        // An earlier column that was already cached still resolves.
925        assert!(codec.has_column(&buffer, &mut cache, 2).is_some());
926
927        // A non-existent column walks off the end and marks the cache finished.
928        assert!(codec.has_column(&buffer, &mut cache, 999).is_none());
929        assert!(cache.finished);
930        // Further misses are O(1).
931        assert!(codec.has_column(&buffer, &mut cache, 998).is_none());
932    }
933
934    #[test]
935    fn test_decode_value_at() {
936        let region_metadata = test_region_metadata();
937        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
938        let mut buffer = Vec::new();
939        let row = test_row();
940        codec.encode_to_vec(row.into_iter(), &mut buffer).unwrap();
941        assert!(!buffer.is_empty());
942
943        let row = test_row();
944        let mut offsets_map = SparseOffsetsCache::new();
945        for column_id in [
946            RESERVED_COLUMN_ID_TABLE_ID,
947            RESERVED_COLUMN_ID_TSID,
948            1,
949            2,
950            3,
951            4,
952            5,
953        ] {
954            let offset = codec
955                .has_column(&buffer, &mut offsets_map, column_id)
956                .unwrap();
957            let value = codec.decode_value_at(&buffer, offset, column_id).unwrap();
958            let expected_value = row
959                .iter()
960                .find(|(id, _)| *id == column_id)
961                .unwrap()
962                .1
963                .clone();
964            assert_eq!(value.as_value_ref(), expected_value);
965        }
966    }
967
968    #[test]
969    fn test_encoded_value_for_column() {
970        let region_metadata = test_region_metadata();
971        let codec = SparsePrimaryKeyCodec::new(&region_metadata);
972        let mut buffer = Vec::new();
973        let row = test_row();
974        codec
975            .encode_to_vec(row.clone().into_iter(), &mut buffer)
976            .unwrap();
977        assert!(!buffer.is_empty());
978
979        let mut offsets_map = SparseOffsetsCache::new();
980        for column_id in [
981            RESERVED_COLUMN_ID_TABLE_ID,
982            RESERVED_COLUMN_ID_TSID,
983            1,
984            2,
985            3,
986            4,
987            5,
988        ] {
989            let encoded_value = codec
990                .encoded_value_for_column(&buffer, &mut offsets_map, column_id)
991                .unwrap()
992                .unwrap();
993            let expected_value = row
994                .iter()
995                .find(|(id, _)| *id == column_id)
996                .unwrap()
997                .1
998                .clone();
999            let data_type = match column_id {
1000                RESERVED_COLUMN_ID_TABLE_ID => ConcreteDataType::uint32_datatype(),
1001                RESERVED_COLUMN_ID_TSID => ConcreteDataType::uint64_datatype(),
1002                _ => ConcreteDataType::string_datatype(),
1003            };
1004            let field = SortField::new(data_type);
1005            let mut expected_encoded = Vec::new();
1006            let mut serializer = Serializer::new(&mut expected_encoded);
1007            field.serialize(&mut serializer, &expected_value).unwrap();
1008            assert_eq!(encoded_value, expected_encoded.as_slice());
1009        }
1010
1011        for column_id in [6_u32, 7_u32, 999_u32] {
1012            let encoded_value = codec
1013                .encoded_value_for_column(&buffer, &mut offsets_map, column_id)
1014                .unwrap();
1015            assert!(encoded_value.is_none());
1016        }
1017    }
1018}