1use std::collections::HashMap;
18use std::fmt;
19use std::fmt::{Debug, Formatter};
20use std::num::NonZeroU64;
21use std::sync::atomic::{AtomicBool, Ordering};
22use std::sync::{Arc, Mutex, RwLock};
23
24use base64::prelude::{BASE64_STANDARD, Engine};
25use bytes::Bytes;
26use common_base::readable_size::ReadableSize;
27use common_telemetry::{debug, error};
28use common_time::Timestamp;
29use partition::expr::PartitionExpr;
30use serde::{Deserialize, Serialize};
31use smallvec::SmallVec;
32use store_api::metadata::ColumnMetadata;
33use store_api::region_request::PathType;
34use store_api::storage::{ColumnId, FileId, IndexVersion, RegionId};
35
36use crate::access_layer::AccessLayerRef;
37use crate::cache::CacheManagerRef;
38use crate::cache::file_cache::{FileType, IndexKey};
39use crate::sst::file_purger::FilePurgerRef;
40use crate::sst::location;
41use crate::sst::parquet::SstInfo;
42
43fn serialize_bytes_option<S>(bytes: &Option<Bytes>, serializer: S) -> Result<S::Ok, S::Error>
45where
46 S: serde::Serializer,
47{
48 match bytes {
49 None => serializer.serialize_none(),
50 Some(b) => serializer.serialize_some(&BASE64_STANDARD.encode(b)),
51 }
52}
53
54fn deserialize_bytes_option<'de, D>(deserializer: D) -> Result<Option<Bytes>, D::Error>
55where
56 D: serde::Deserializer<'de>,
57{
58 let opt: Option<String> = Option::deserialize(deserializer)?;
59 match opt {
60 None => Ok(None),
61 Some(s) => {
62 let decoded = BASE64_STANDARD
63 .decode(&s)
64 .map_err(serde::de::Error::custom)?;
65 Ok(Some(Bytes::from(decoded)))
66 }
67 }
68}
69
70fn serialize_partition_expr<S>(
72 partition_expr: &Option<PartitionExpr>,
73 serializer: S,
74) -> Result<S::Ok, S::Error>
75where
76 S: serde::Serializer,
77{
78 use serde::ser::Error;
79
80 match partition_expr {
81 None => serializer.serialize_none(),
82 Some(expr) => {
83 let json_str = expr.as_json_str().map_err(S::Error::custom)?;
84 serializer.serialize_some(&json_str)
85 }
86 }
87}
88
89fn deserialize_partition_expr<'de, D>(deserializer: D) -> Result<Option<PartitionExpr>, D::Error>
90where
91 D: serde::Deserializer<'de>,
92{
93 use serde::de::Error;
94
95 let opt_json_str: Option<String> = Option::deserialize(deserializer)?;
96 match opt_json_str {
97 None => Ok(None),
98 Some(json_str) => {
99 if json_str.is_empty() {
100 Ok(None)
102 } else {
103 PartitionExpr::from_json_str(&json_str).map_err(D::Error::custom)
105 }
106 }
107 }
108}
109
110pub type Level = u8;
112pub const MAX_LEVEL: Level = 2;
114pub type IndexTypes = SmallVec<[IndexType; 4]>;
116
117#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
121pub struct RegionFileId {
122 region_id: RegionId,
124 file_id: FileId,
126}
127
128impl RegionFileId {
129 pub fn new(region_id: RegionId, file_id: FileId) -> Self {
131 Self { region_id, file_id }
132 }
133
134 pub fn region_id(&self) -> RegionId {
136 self.region_id
137 }
138
139 pub fn file_id(&self) -> FileId {
141 self.file_id
142 }
143}
144
145impl fmt::Display for RegionFileId {
146 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
147 write!(f, "{}/{}", self.region_id, self.file_id)
148 }
149}
150
151#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
153pub struct RegionIndexId {
154 pub file_id: RegionFileId,
155 pub version: IndexVersion,
156}
157
158impl RegionIndexId {
159 pub fn new(file_id: RegionFileId, version: IndexVersion) -> Self {
160 Self { file_id, version }
161 }
162
163 pub fn region_id(&self) -> RegionId {
164 self.file_id.region_id
165 }
166
167 pub fn file_id(&self) -> FileId {
168 self.file_id.file_id
169 }
170}
171
172impl fmt::Display for RegionIndexId {
173 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
174 if self.version == 0 {
175 write!(f, "{}/{}", self.file_id.region_id, self.file_id.file_id)
176 } else {
177 write!(
178 f,
179 "{}/{}.{}",
180 self.file_id.region_id, self.file_id.file_id, self.version
181 )
182 }
183 }
184}
185
186pub type FileTimeRange = (Timestamp, Timestamp);
189
190pub(crate) fn overlaps(l: &FileTimeRange, r: &FileTimeRange) -> bool {
192 let (l, r) = if l.0 <= r.0 { (l, r) } else { (r, l) };
193 let (_, l_end) = l;
194 let (r_start, _) = r;
195
196 r_start <= l_end
197}
198
199#[derive(Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
201#[serde(default)]
202pub struct FileMeta {
203 pub region_id: RegionId,
205 pub file_id: FileId,
207 pub time_range: FileTimeRange,
210 pub level: Level,
212 pub file_size: u64,
214 pub max_row_group_uncompressed_size: u64,
216 pub available_indexes: IndexTypes,
218 pub indexes: Vec<ColumnIndexMetadata>,
228 pub index_file_size: u64,
230 pub index_version: u64,
234 pub num_rows: u64,
240 pub num_row_groups: u64,
246 pub sequence: Option<NonZeroU64>,
251 #[serde(
259 serialize_with = "serialize_partition_expr",
260 deserialize_with = "deserialize_partition_expr"
261 )]
262 pub partition_expr: Option<PartitionExpr>,
263 pub num_series: u64,
267 #[serde(
270 default,
271 skip_serializing_if = "Option::is_none",
272 serialize_with = "serialize_bytes_option",
273 deserialize_with = "deserialize_bytes_option"
274 )]
275 pub primary_key_min: Option<Bytes>,
276 #[serde(
279 default,
280 skip_serializing_if = "Option::is_none",
281 serialize_with = "serialize_bytes_option",
282 deserialize_with = "deserialize_bytes_option"
283 )]
284 pub primary_key_max: Option<Bytes>,
285 #[serde(default, skip_serializing_if = "is_false")]
288 pub preserve_row_sequence: bool,
289}
290
291fn is_false(value: &bool) -> bool {
292 !*value
293}
294
295impl Debug for FileMeta {
296 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
297 let mut debug_struct = f.debug_struct("FileMeta");
298 debug_struct
299 .field("region_id", &self.region_id)
300 .field_with("file_id", |f| write!(f, "{} ", self.file_id))
301 .field_with("time_range", |f| {
302 write!(
303 f,
304 "({}, {}) ",
305 self.time_range.0.to_iso8601_string(),
306 self.time_range.1.to_iso8601_string()
307 )
308 })
309 .field("level", &self.level)
310 .field("file_size", &ReadableSize(self.file_size))
311 .field(
312 "max_row_group_uncompressed_size",
313 &ReadableSize(self.max_row_group_uncompressed_size),
314 );
315 if !self.available_indexes.is_empty() {
316 debug_struct
317 .field("available_indexes", &self.available_indexes)
318 .field("indexes", &self.indexes)
319 .field("index_file_size", &ReadableSize(self.index_file_size));
320 }
321 debug_struct
322 .field("num_rows", &self.num_rows)
323 .field("num_row_groups", &self.num_row_groups)
324 .field_with("sequence", |f| match self.sequence {
325 None => {
326 write!(f, "None")
327 }
328 Some(seq) => {
329 write!(f, "{}", seq)
330 }
331 })
332 .field("partition_expr", &self.partition_expr)
333 .field("num_series", &self.num_series);
334 if self.primary_key_min.is_some() || self.primary_key_max.is_some() {
335 debug_struct
336 .field(
337 "primary_key_min",
338 &self.primary_key_min.as_ref().map(|b| b.len()),
339 )
340 .field(
341 "primary_key_max",
342 &self.primary_key_max.as_ref().map(|b| b.len()),
343 );
344 }
345 debug_struct.finish()
346 }
347}
348
349#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
351pub enum IndexType {
352 InvertedIndex,
354 FulltextIndex,
356 BloomFilterIndex,
358 #[cfg(feature = "vector_index")]
360 VectorIndex,
361}
362
363#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize, Default)]
373#[serde(default)]
374pub struct ColumnIndexMetadata {
375 pub column_id: ColumnId,
377 pub created_indexes: IndexTypes,
379}
380
381impl FileMeta {
382 pub fn primary_key_range(&self) -> Option<(Bytes, Bytes)> {
384 match (&self.primary_key_min, &self.primary_key_max) {
385 (Some(min), Some(max)) => Some((min.clone(), max.clone())),
386 _ => None,
387 }
388 }
389
390 pub fn exists_index(&self) -> bool {
391 !self.available_indexes.is_empty()
392 }
393
394 pub fn index_version(&self) -> Option<IndexVersion> {
395 if self.exists_index() {
396 Some(self.index_version)
397 } else {
398 None
399 }
400 }
401
402 pub fn is_index_up_to_date(&self, other: &FileMeta) -> bool {
404 self.exists_index() && other.exists_index() && self.index_version >= other.index_version
405 }
406
407 pub fn inverted_index_available(&self) -> bool {
409 self.available_indexes.contains(&IndexType::InvertedIndex)
410 }
411
412 pub fn fulltext_index_available(&self) -> bool {
414 self.available_indexes.contains(&IndexType::FulltextIndex)
415 }
416
417 pub fn bloom_filter_index_available(&self) -> bool {
419 self.available_indexes
420 .contains(&IndexType::BloomFilterIndex)
421 }
422
423 #[cfg(feature = "vector_index")]
425 pub fn vector_index_available(&self) -> bool {
426 self.available_indexes.contains(&IndexType::VectorIndex)
427 }
428
429 pub fn index_file_size(&self) -> u64 {
430 self.index_file_size
431 }
432
433 pub fn is_index_consistent_with_region(&self, metadata: &[ColumnMetadata]) -> bool {
435 let id_to_indexes = self
436 .indexes
437 .iter()
438 .map(|index| (index.column_id, index.created_indexes.clone()))
439 .collect::<std::collections::HashMap<_, _>>();
440 for column in metadata {
441 if !column.column_schema.is_indexed() {
442 continue;
443 }
444 if let Some(indexes) = id_to_indexes.get(&column.column_id) {
445 if column.column_schema.is_inverted_indexed()
446 && !indexes.contains(&IndexType::InvertedIndex)
447 {
448 return false;
449 }
450 if column.column_schema.is_fulltext_indexed()
451 && !indexes.contains(&IndexType::FulltextIndex)
452 {
453 return false;
454 }
455 if column.column_schema.is_skipping_indexed()
456 && !indexes.contains(&IndexType::BloomFilterIndex)
457 {
458 return false;
459 }
460 } else {
461 return false;
462 }
463 }
464 true
465 }
466
467 pub fn file_id(&self) -> RegionFileId {
469 RegionFileId::new(self.region_id, self.file_id)
470 }
471
472 pub fn index_id(&self) -> RegionIndexId {
474 RegionIndexId::new(self.file_id(), self.index_version)
475 }
476}
477
478#[derive(Clone)]
480pub struct FileHandle {
481 inner: Arc<FileHandleInner>,
482}
483
484impl fmt::Debug for FileHandle {
485 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
486 f.debug_struct("FileHandle")
487 .field("meta", self.meta_ref())
488 .field("compacting", &self.compacting())
489 .field("deleted", &self.inner.deleted.load(Ordering::Relaxed))
490 .finish()
491 }
492}
493
494impl FileHandle {
495 pub fn new(meta: FileMeta, file_purger: FilePurgerRef) -> FileHandle {
496 let pk_range = meta.primary_key_range();
497 FileHandle {
498 inner: Arc::new(FileHandleInner::new(meta, file_purger, pk_range)),
499 }
500 }
501
502 #[cfg(test)]
503 pub fn new_with_primary_key_range(
504 meta: FileMeta,
505 file_purger: FilePurgerRef,
506 primary_key_range: Option<(Bytes, Bytes)>,
507 ) -> FileHandle {
508 FileHandle {
509 inner: Arc::new(FileHandleInner::new(meta, file_purger, primary_key_range)),
510 }
511 }
512
513 pub fn region_id(&self) -> RegionId {
515 self.inner.meta.region_id
516 }
517
518 pub(crate) fn is_effective_target_sequence_trusted(&self, target_region_id: RegionId) -> bool {
523 if self.region_id() != target_region_id {
524 self.meta_ref().sequence.is_some()
525 } else {
526 self.meta_ref().preserve_row_sequence
527 }
528 }
529
530 pub fn file_id(&self) -> RegionFileId {
532 RegionFileId::new(self.inner.meta.region_id, self.inner.meta.file_id)
533 }
534
535 pub fn index_id(&self) -> RegionIndexId {
537 RegionIndexId::new(self.file_id(), self.inner.meta.index_version)
538 }
539
540 pub fn file_path(&self, table_dir: &str, path_type: PathType) -> String {
542 location::sst_file_path(table_dir, self.file_id(), path_type)
543 }
544
545 pub fn time_range(&self) -> FileTimeRange {
547 self.inner.meta.time_range
548 }
549
550 pub fn mark_deleted(&self) {
552 self.inner.deleted.store(true, Ordering::Relaxed);
553 }
554
555 pub fn compacting(&self) -> bool {
556 self.inner.compacting.load(Ordering::Relaxed)
557 }
558
559 pub fn set_compacting(&self, compacting: bool) {
560 self.inner.compacting.store(compacting, Ordering::Relaxed);
561 }
562
563 pub fn try_set_compacting(&self) -> bool {
565 self.inner
566 .compacting
567 .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
568 .is_ok()
569 }
570
571 pub fn index_outdated(&self) -> bool {
572 self.inner.index_outdated.load(Ordering::Relaxed)
573 }
574
575 pub fn set_index_outdated(&self, index_outdated: bool) {
576 self.inner
577 .index_outdated
578 .store(index_outdated, Ordering::Relaxed);
579 }
580
581 pub fn meta_ref(&self) -> &FileMeta {
583 &self.inner.meta
584 }
585
586 pub fn file_purger(&self) -> FilePurgerRef {
587 self.inner.file_purger.clone()
588 }
589
590 pub fn size(&self) -> u64 {
591 self.inner.meta.file_size
592 }
593
594 pub fn index_size(&self) -> u64 {
595 self.inner.meta.index_file_size
596 }
597
598 pub fn num_rows(&self) -> usize {
599 self.inner.meta.num_rows as usize
600 }
601
602 pub fn level(&self) -> Level {
603 self.inner.meta.level
604 }
605
606 pub fn is_deleted(&self) -> bool {
607 self.inner.deleted.load(Ordering::Relaxed)
608 }
609
610 pub fn primary_key_range(&self) -> Option<(Bytes, Bytes)> {
611 self.inner.primary_key_range.read().unwrap().clone()
612 }
613
614 pub(crate) fn set_primary_key_range(&self, primary_key_range: (Bytes, Bytes)) {
615 *self.inner.primary_key_range.write().unwrap() = Some(primary_key_range);
616 }
617}
618
619struct FileHandleInner {
623 meta: FileMeta,
624 compacting: AtomicBool,
625 deleted: AtomicBool,
626 index_outdated: AtomicBool,
627 primary_key_range: RwLock<Option<(Bytes, Bytes)>>,
628 file_purger: FilePurgerRef,
629}
630
631impl Drop for FileHandleInner {
632 fn drop(&mut self) {
633 self.file_purger.remove_file(
634 self.meta.clone(),
635 self.deleted.load(Ordering::Acquire),
636 self.index_outdated.load(Ordering::Acquire),
637 );
638 }
639}
640
641impl FileHandleInner {
642 fn new(
644 meta: FileMeta,
645 file_purger: FilePurgerRef,
646 primary_key_range: Option<(Bytes, Bytes)>,
647 ) -> FileHandleInner {
648 file_purger.new_file(&meta);
649 FileHandleInner {
650 meta,
651 compacting: AtomicBool::new(false),
652 deleted: AtomicBool::new(false),
653 index_outdated: AtomicBool::new(false),
654 primary_key_range: RwLock::new(primary_key_range),
655 file_purger,
656 }
657 }
658}
659
660pub async fn delete_files(
667 region_id: RegionId,
668 file_ids: &[(FileId, u64)],
669 delete_index: bool,
670 access_layer: &AccessLayerRef,
671 cache_manager: &Option<CacheManagerRef>,
672) -> crate::error::Result<()> {
673 if let Some(cache) = &cache_manager {
675 for (file_id, _) in file_ids {
676 cache.remove_parquet_meta_data(RegionFileId::new(region_id, *file_id));
677 }
678 }
679 let mut attempted_files = Vec::with_capacity(file_ids.len());
680 let mut index_ids = Vec::new();
681
682 for (file_id, index_version) in file_ids {
683 let region_file_id = RegionFileId::new(region_id, *file_id);
684 attempted_files.push(*file_id);
685 index_ids.extend(
686 (0..=*index_version).map(|version| RegionIndexId::new(region_file_id, version)),
687 );
688 }
689
690 access_layer
691 .delete_ssts(region_id, &attempted_files)
692 .await?;
693 access_layer.delete_indexes(&index_ids).await?;
694
695 debug!(
696 "Attempted to delete {} files for region {}: {:?}",
697 attempted_files.len(),
698 region_id,
699 attempted_files
700 );
701
702 for (file_id, index_version) in file_ids {
703 purge_index_cache_stager(
704 region_id,
705 delete_index,
706 access_layer,
707 cache_manager,
708 *file_id,
709 *index_version,
710 )
711 .await;
712 }
713 Ok(())
714}
715
716#[derive(Clone)]
721pub(crate) struct UncommittedSsts {
722 region_id: RegionId,
723 files: Arc<Mutex<HashMap<FileId, (u64, bool)>>>,
724 access_layer: AccessLayerRef,
725 cache_manager: Option<CacheManagerRef>,
726}
727
728impl UncommittedSsts {
729 pub(crate) fn new(
730 region_id: RegionId,
731 access_layer: AccessLayerRef,
732 cache_manager: Option<CacheManagerRef>,
733 ) -> Self {
734 Self {
735 region_id,
736 files: Arc::new(Mutex::new(HashMap::new())),
737 access_layer,
738 cache_manager,
739 }
740 }
741
742 pub(crate) fn track(&self, ssts: &[SstInfo]) {
744 let mut files = self.files.lock().unwrap();
745 for sst in ssts {
746 files.insert(
747 sst.file_id,
748 (sst.index_metadata.version, sst.index_metadata.file_size > 0),
749 );
750 }
751 }
752
753 pub(crate) fn disarm_cleanup(&self) {
758 self.files.lock().unwrap().clear();
759 }
760
761 #[cfg(test)]
762 pub(crate) fn num_tracked_files(&self) -> usize {
763 self.files.lock().unwrap().len()
764 }
765
766 pub(crate) async fn cleanup(&self) {
768 if let Err(err) = self.try_cleanup().await {
769 error!(err; "Failed to clean uncommitted SSTs for region {}", self.region_id);
770 }
771 }
772
773 async fn try_cleanup(&self) -> crate::error::Result<()> {
774 let files = std::mem::take(&mut *self.files.lock().unwrap());
775 if files.is_empty() {
776 return Ok(());
777 }
778
779 let delete_index = files.values().any(|(_, exists_index)| *exists_index);
780 let file_ids = files
781 .into_iter()
782 .map(|(file_id, (index_version, _))| (file_id, index_version))
783 .collect::<Vec<_>>();
784 delete_files(
785 self.region_id,
786 &file_ids,
787 delete_index,
788 &self.access_layer,
789 &self.cache_manager,
790 )
791 .await
792 }
793
794 #[cfg(test)]
795 pub(crate) async fn cleanup_for_test(&self) -> crate::error::Result<()> {
796 self.try_cleanup().await
797 }
798}
799
800pub async fn delete_index(
801 region_index_id: RegionIndexId,
802 access_layer: &AccessLayerRef,
803 cache_manager: &Option<CacheManagerRef>,
804) -> crate::error::Result<()> {
805 delete_index_and_purge(region_index_id, access_layer, cache_manager).await?;
806
807 Ok(())
808}
809
810pub async fn delete_indexes(
811 index_ids: &[RegionIndexId],
812 access_layer: &AccessLayerRef,
813 cache_manager: &Option<CacheManagerRef>,
814) -> crate::error::Result<()> {
815 if index_ids.is_empty() {
816 return Ok(());
817 }
818
819 if let Err(e) = access_layer.delete_indexes(index_ids).await {
820 error!(e; "Failed to batch delete index files");
821
822 for index_id in index_ids {
823 delete_index_and_purge(*index_id, access_layer, cache_manager).await?;
824 }
825
826 return Ok(());
827 }
828
829 purge_indexes(index_ids, access_layer, cache_manager).await;
830
831 Ok(())
832}
833
834async fn delete_index_and_purge(
835 index_id: RegionIndexId,
836 access_layer: &AccessLayerRef,
837 cache_manager: &Option<CacheManagerRef>,
838) -> crate::error::Result<()> {
839 access_layer.delete_index(index_id).await?;
840 purge_index_cache_stager(
841 index_id.region_id(),
842 true,
843 access_layer,
844 cache_manager,
845 index_id.file_id(),
846 index_id.version,
847 )
848 .await;
849 Ok(())
850}
851
852async fn purge_indexes(
853 index_ids: &[RegionIndexId],
854 access_layer: &AccessLayerRef,
855 cache_manager: &Option<CacheManagerRef>,
856) {
857 for index_id in index_ids {
858 purge_index_cache_stager(
859 index_id.region_id(),
860 true,
861 access_layer,
862 cache_manager,
863 index_id.file_id(),
864 index_id.version,
865 )
866 .await;
867 }
868}
869
870async fn purge_index_cache_stager(
871 region_id: RegionId,
872 delete_index: bool,
873 access_layer: &AccessLayerRef,
874 cache_manager: &Option<CacheManagerRef>,
875 file_id: FileId,
876 index_version: u64,
877) {
878 if let Some(write_cache) = cache_manager.as_ref().and_then(|cache| cache.write_cache()) {
879 if delete_index {
881 write_cache
882 .remove(IndexKey::new(
883 region_id,
884 file_id,
885 FileType::Puffin(index_version),
886 ))
887 .await;
888 }
889
890 write_cache
892 .remove(IndexKey::new(region_id, file_id, FileType::Parquet))
893 .await;
894 }
895
896 if let Err(e) = access_layer
898 .puffin_manager_factory()
899 .purge_stager(RegionIndexId::new(
900 RegionFileId::new(region_id, file_id),
901 index_version,
902 ))
903 .await
904 {
905 error!(e; "Failed to purge stager with index file, file_id: {}, index_version: {}, region: {}",
906 file_id, index_version, region_id);
907 }
908}
909
910#[cfg(test)]
911mod tests {
912 use std::str::FromStr;
913
914 use datatypes::prelude::ConcreteDataType;
915 use datatypes::schema::{
916 ColumnSchema, FulltextAnalyzer, FulltextBackend, FulltextOptions, SkippingIndexOptions,
917 };
918 use datatypes::value::Value;
919 use partition::expr::{PartitionExpr, col};
920
921 use super::*;
922
923 fn create_file_meta(file_id: FileId, level: Level) -> FileMeta {
924 FileMeta {
925 region_id: 0.into(),
926 file_id,
927 time_range: FileTimeRange::default(),
928 level,
929 file_size: 0,
930 max_row_group_uncompressed_size: 0,
931 available_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
932 indexes: vec![ColumnIndexMetadata {
933 column_id: 0,
934 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
935 }],
936 index_file_size: 0,
937 index_version: 0,
938 num_rows: 0,
939 num_row_groups: 0,
940 sequence: None,
941 partition_expr: None,
942 num_series: 0,
943 ..Default::default()
944 }
945 }
946
947 #[test]
948 fn test_deserialize_file_meta() {
949 let file_meta = create_file_meta(FileId::random(), 0);
950 let serialized_file_meta = serde_json::to_string(&file_meta).unwrap();
951 let deserialized_file_meta = serde_json::from_str(&serialized_file_meta);
952 assert_eq!(file_meta, deserialized_file_meta.unwrap());
953 }
954
955 #[test]
956 fn test_deserialize_from_string() {
957 let json_file_meta = "{\"region_id\":0,\"file_id\":\"bc5896ec-e4d8-4017-a80d-f2de73188d55\",\
958 \"time_range\":[{\"value\":0,\"unit\":\"Millisecond\"},{\"value\":0,\"unit\":\"Millisecond\"}],\
959 \"available_indexes\":[\"InvertedIndex\"],\"indexes\":[{\"column_id\": 0, \"created_indexes\": [\"InvertedIndex\"]}],\"level\":0}";
960 let file_meta = create_file_meta(
961 FileId::from_str("bc5896ec-e4d8-4017-a80d-f2de73188d55").unwrap(),
962 0,
963 );
964 let deserialized_file_meta: FileMeta = serde_json::from_str(json_file_meta).unwrap();
965 assert_eq!(file_meta, deserialized_file_meta);
966 }
967
968 #[test]
969 fn test_file_meta_with_partition_expr() {
970 let file_id = FileId::random();
971 let partition_expr = PartitionExpr::new(
972 col("a"),
973 partition::expr::RestrictedOp::GtEq,
974 Value::UInt32(10).into(),
975 );
976
977 let file_meta_with_partition = FileMeta {
978 region_id: 0.into(),
979 file_id,
980 time_range: FileTimeRange::default(),
981 level: 0,
982 file_size: 0,
983 max_row_group_uncompressed_size: 0,
984 available_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
985 indexes: vec![ColumnIndexMetadata {
986 column_id: 0,
987 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
988 }],
989 index_file_size: 0,
990 index_version: 0,
991 num_rows: 0,
992 num_row_groups: 0,
993 sequence: None,
994 partition_expr: Some(partition_expr.clone()),
995 num_series: 0,
996 ..Default::default()
997 };
998
999 let serialized = serde_json::to_string(&file_meta_with_partition).unwrap();
1001 let deserialized: FileMeta = serde_json::from_str(&serialized).unwrap();
1002 assert_eq!(file_meta_with_partition, deserialized);
1003
1004 let serialized_value: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1006 assert!(serialized_value["partition_expr"].as_str().is_some());
1007 let partition_expr_json = serialized_value["partition_expr"].as_str().unwrap();
1008 assert!(partition_expr_json.contains("\"Column\":\"a\""));
1009 assert!(partition_expr_json.contains("\"op\":\"GtEq\""));
1010
1011 let file_meta_none = FileMeta {
1013 partition_expr: None,
1014 ..file_meta_with_partition.clone()
1015 };
1016 let serialized_none = serde_json::to_string(&file_meta_none).unwrap();
1017 let deserialized_none: FileMeta = serde_json::from_str(&serialized_none).unwrap();
1018 assert_eq!(file_meta_none, deserialized_none);
1019 }
1020
1021 #[test]
1022 fn test_file_meta_partition_expr_backward_compatibility() {
1023 let json_with_partition_expr = r#"{
1025 "region_id": 0,
1026 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1027 "time_range": [
1028 {"value": 0, "unit": "Millisecond"},
1029 {"value": 0, "unit": "Millisecond"}
1030 ],
1031 "level": 0,
1032 "file_size": 0,
1033 "available_indexes": ["InvertedIndex"],
1034 "index_file_size": 0,
1035 "num_rows": 0,
1036 "num_row_groups": 0,
1037 "sequence": null,
1038 "partition_expr": "{\"Expr\":{\"lhs\":{\"Column\":\"a\"},\"op\":\"GtEq\",\"rhs\":{\"Value\":{\"UInt32\":10}}}}"
1039 }"#;
1040
1041 let file_meta: FileMeta = serde_json::from_str(json_with_partition_expr).unwrap();
1042 assert!(file_meta.partition_expr.is_some());
1043 let expr = file_meta.partition_expr.unwrap();
1044 assert_eq!(format!("{}", expr), "a >= 10");
1045
1046 let json_with_empty_expr = r#"{
1048 "region_id": 0,
1049 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1050 "time_range": [
1051 {"value": 0, "unit": "Millisecond"},
1052 {"value": 0, "unit": "Millisecond"}
1053 ],
1054 "level": 0,
1055 "file_size": 0,
1056 "available_indexes": [],
1057 "index_file_size": 0,
1058 "num_rows": 0,
1059 "num_row_groups": 0,
1060 "sequence": null,
1061 "partition_expr": ""
1062 }"#;
1063
1064 let file_meta_empty: FileMeta = serde_json::from_str(json_with_empty_expr).unwrap();
1065 assert!(file_meta_empty.partition_expr.is_none());
1066
1067 let json_with_null_expr = r#"{
1069 "region_id": 0,
1070 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1071 "time_range": [
1072 {"value": 0, "unit": "Millisecond"},
1073 {"value": 0, "unit": "Millisecond"}
1074 ],
1075 "level": 0,
1076 "file_size": 0,
1077 "available_indexes": [],
1078 "index_file_size": 0,
1079 "num_rows": 0,
1080 "num_row_groups": 0,
1081 "sequence": null,
1082 "partition_expr": null
1083 }"#;
1084
1085 let file_meta_null: FileMeta = serde_json::from_str(json_with_null_expr).unwrap();
1086 assert!(file_meta_null.partition_expr.is_none());
1087
1088 let json_with_empty_expr = r#"{
1090 "region_id": 0,
1091 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1092 "time_range": [
1093 {"value": 0, "unit": "Millisecond"},
1094 {"value": 0, "unit": "Millisecond"}
1095 ],
1096 "level": 0,
1097 "file_size": 0,
1098 "available_indexes": [],
1099 "index_file_size": 0,
1100 "num_rows": 0,
1101 "num_row_groups": 0,
1102 "sequence": null
1103 }"#;
1104
1105 let file_meta_empty: FileMeta = serde_json::from_str(json_with_empty_expr).unwrap();
1106 assert!(file_meta_empty.partition_expr.is_none());
1107 }
1108
1109 #[test]
1110 fn test_file_meta_indexes_backward_compatibility() {
1111 let json_old_file_meta = r#"{
1113 "region_id": 0,
1114 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1115 "time_range": [
1116 {"value": 0, "unit": "Millisecond"},
1117 {"value": 0, "unit": "Millisecond"}
1118 ],
1119 "available_indexes": ["InvertedIndex"],
1120 "level": 0,
1121 "file_size": 0,
1122 "index_file_size": 0,
1123 "num_rows": 0,
1124 "num_row_groups": 0
1125 }"#;
1126
1127 let deserialized_file_meta: FileMeta = serde_json::from_str(json_old_file_meta).unwrap();
1128
1129 assert_eq!(deserialized_file_meta.indexes, vec![]);
1131
1132 let expected_indexes: IndexTypes = SmallVec::from_iter([IndexType::InvertedIndex]);
1133 assert_eq!(deserialized_file_meta.available_indexes, expected_indexes);
1134
1135 assert_eq!(
1136 deserialized_file_meta.file_id,
1137 FileId::from_str("bc5896ec-e4d8-4017-a80d-f2de73188d55").unwrap()
1138 );
1139 assert!(!deserialized_file_meta.preserve_row_sequence);
1140 }
1141
1142 #[test]
1143 fn test_file_meta_preserve_row_sequence_serde() {
1144 let file_meta = FileMeta {
1145 preserve_row_sequence: true,
1146 ..Default::default()
1147 };
1148
1149 let serialized = serde_json::to_string(&file_meta).unwrap();
1150 let value: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1151 assert_eq!(value["preserve_row_sequence"], true);
1152
1153 let deserialized: FileMeta = serde_json::from_str(&serialized).unwrap();
1154 assert_eq!(file_meta, deserialized);
1155
1156 let file_meta_false = FileMeta {
1157 preserve_row_sequence: false,
1158 ..file_meta.clone()
1159 };
1160 let serialized_false = serde_json::to_string(&file_meta_false).unwrap();
1161 let value_false: serde_json::Value = serde_json::from_str(&serialized_false).unwrap();
1162 assert!(value_false.get("preserve_row_sequence").is_none());
1163 }
1164 #[test]
1165 fn test_is_index_consistent_with_region() {
1166 fn new_column_meta(
1167 id: ColumnId,
1168 name: &str,
1169 inverted: bool,
1170 fulltext: bool,
1171 skipping: bool,
1172 ) -> ColumnMetadata {
1173 let mut column_schema =
1174 ColumnSchema::new(name, ConcreteDataType::string_datatype(), true);
1175 if inverted {
1176 column_schema = column_schema.with_inverted_index(true);
1177 }
1178 if fulltext {
1179 column_schema = column_schema
1180 .with_fulltext_options(FulltextOptions::new_unchecked(
1181 true,
1182 FulltextAnalyzer::English,
1183 false,
1184 FulltextBackend::Bloom,
1185 1000,
1186 0.01,
1187 ))
1188 .unwrap();
1189 }
1190 if skipping {
1191 column_schema = column_schema
1192 .with_skipping_options(SkippingIndexOptions::new_unchecked(
1193 1024,
1194 0.01,
1195 datatypes::schema::SkippingIndexType::BloomFilter,
1196 ))
1197 .unwrap();
1198 }
1199
1200 ColumnMetadata {
1201 column_schema,
1202 semantic_type: api::v1::SemanticType::Tag,
1203 column_id: id,
1204 }
1205 }
1206
1207 let mut file_meta = FileMeta {
1209 indexes: vec![ColumnIndexMetadata {
1210 column_id: 1,
1211 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1212 }],
1213 ..Default::default()
1214 };
1215 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1216 assert!(file_meta.is_index_consistent_with_region(®ion_meta));
1217
1218 file_meta.indexes = vec![ColumnIndexMetadata {
1220 column_id: 1,
1221 created_indexes: SmallVec::from_iter([
1222 IndexType::InvertedIndex,
1223 IndexType::BloomFilterIndex,
1224 ]),
1225 }];
1226 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1227 assert!(file_meta.is_index_consistent_with_region(®ion_meta));
1228
1229 file_meta.indexes = vec![ColumnIndexMetadata {
1231 column_id: 1,
1232 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1233 }];
1234 let region_meta = vec![new_column_meta(1, "tag1", true, true, false)]; assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1236
1237 file_meta.indexes = vec![ColumnIndexMetadata {
1239 column_id: 2, created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1241 }];
1242 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)]; assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1244
1245 file_meta.indexes = vec![ColumnIndexMetadata {
1247 column_id: 1,
1248 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1249 }];
1250 let region_meta = vec![new_column_meta(1, "tag1", false, false, false)]; assert!(file_meta.is_index_consistent_with_region(®ion_meta));
1252
1253 file_meta.indexes = vec![];
1255 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1256 assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1257
1258 file_meta.indexes = vec![
1260 ColumnIndexMetadata {
1261 column_id: 1,
1262 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1263 },
1264 ColumnIndexMetadata {
1265 column_id: 2, created_indexes: SmallVec::from_iter([IndexType::FulltextIndex]),
1267 },
1268 ];
1269 let region_meta = vec![
1270 new_column_meta(1, "tag1", true, false, false),
1271 new_column_meta(2, "tag2", false, true, true), ];
1273 assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1274 }
1275}