1use 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#[derive(Clone, Debug)]
57pub struct SparsePrimaryKeyCodec {
58 inner: Arc<SparsePrimaryKeyCodecInner>,
59}
60
61#[derive(Debug)]
62struct SparsePrimaryKeyCodecInner {
63 table_id_field: SortField,
65 tsid_field: SortField,
67 label_field: SortField,
69 columns: Option<HashSet<ColumnId>>,
73}
74
75#[derive(Debug, Clone, PartialEq, Eq, Default)]
80pub struct SparseValues {
81 values: Vec<(ColumnId, Value)>,
82}
83
84impl SparseValues {
85 pub fn new() -> Self {
87 Self { values: Vec::new() }
88 }
89
90 pub fn with_capacity(cap: usize) -> Self {
92 Self {
93 values: Vec::with_capacity(cap),
94 }
95 }
96
97 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 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 pub fn insert(&mut self, column_id: ColumnId, value: Value) {
121 self.values.push((column_id, value));
122 }
123
124 pub fn iter(&self) -> impl Iterator<Item = (&ColumnId, &Value)> {
126 self.values.iter().map(|(id, value)| (id, value))
127 }
128}
129
130pub const RESERVED_COLUMN_ID_TSID: ColumnId = ReservedColumnId::tsid();
132pub const RESERVED_COLUMN_ID_TABLE_ID: ColumnId = ReservedColumnId::table_id();
134pub const COLUMN_ID_ENCODE_SIZE: usize = 4;
136
137const TABLE_ID_VALUE_OFFSET: usize = COLUMN_ID_ENCODE_SIZE;
141const TSID_VALUE_OFFSET: usize = COLUMN_ID_ENCODE_SIZE + 5 + COLUMN_ID_ENCODE_SIZE;
143const TAGS_START_OFFSET: usize = COLUMN_ID_ENCODE_SIZE + 5 + COLUMN_ID_ENCODE_SIZE + 9;
145
146const SPARSE_OFFSETS_INLINE_CAP: usize = 32;
152
153#[derive(Debug, Clone)]
155pub struct SparseOffsetsCache {
156 inline: Vec<(ColumnId, usize)>,
159 overflow: HashMap<ColumnId, usize>,
161 cursor: usize,
163 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 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 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 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 pub fn new(region_metadata: &RegionMetadataRef) -> Self {
235 Self::from_columns(region_metadata.primary_key_columns().map(|c| c.column_id))
236 }
237
238 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 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 fn get_field(&self, column_id: ColumnId) -> Option<&SortField> {
266 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 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 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 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 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 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 pub fn decode_ids(&self, bytes: &[u8]) -> Result<(u32, u64)> {
354 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 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 fn decode_leftmost(&self, bytes: &[u8]) -> Result<Option<Value>> {
426 let mut deserializer = Deserializer::new(bytes);
427 deserializer.advance(COLUMN_ID_ENCODE_SIZE);
429 let value = self.inner.table_id_field.deserialize(&mut deserializer)?;
430 Ok(Some(value))
431 }
432
433 pub fn has_column(
443 &self,
444 pk: &[u8],
445 cache: &mut SparseOffsetsCache,
446 column_id: ColumnId,
447 ) -> Option<usize> {
448 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 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 let field = self.get_field(column_id).unwrap();
501 field.deserialize(&mut deserializer)
502 }
503
504 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
581pub struct FieldWithId {
583 pub field: SortField,
584 pub column_id: ColumnId,
585}
586
587pub 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 (1, ValueRef::String("greptime-frontend-6989d9899-22222")),
752 (2, ValueRef::String("greptime-cluster")),
754 (3, ValueRef::String("greptime-frontend-6989d9899-22222")),
756 (4, ValueRef::String("greptime-frontend-6989d9899-22222")),
758 (5, ValueRef::String("10.10.10.10")),
760 (6, ValueRef::Float64(OrderedFloat(1.0))),
762 (
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(®ion_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(®ion_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(®ion_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(®ion_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(®ion_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(®ion_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 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 assert!(codec.has_column(&buffer, &mut cache, 5).is_some());
922 assert!(cache.contains(5));
923
924 assert!(codec.has_column(&buffer, &mut cache, 2).is_some());
926
927 assert!(codec.has_column(&buffer, &mut cache, 999).is_none());
929 assert!(cache.finished);
930 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(®ion_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(®ion_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}