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, warn};
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;
42use crate::sst::primary_key::PrimaryKeyRangeMapper;
43
44fn serialize_bytes_option<S>(bytes: &Option<Bytes>, serializer: S) -> Result<S::Ok, S::Error>
46where
47 S: serde::Serializer,
48{
49 match bytes {
50 None => serializer.serialize_none(),
51 Some(b) => serializer.serialize_some(&BASE64_STANDARD.encode(b)),
52 }
53}
54
55fn deserialize_bytes_option<'de, D>(deserializer: D) -> Result<Option<Bytes>, D::Error>
56where
57 D: serde::Deserializer<'de>,
58{
59 let opt: Option<String> = Option::deserialize(deserializer)?;
60 match opt {
61 None => Ok(None),
62 Some(s) => {
63 let decoded = BASE64_STANDARD
64 .decode(&s)
65 .map_err(serde::de::Error::custom)?;
66 Ok(Some(Bytes::from(decoded)))
67 }
68 }
69}
70
71fn serialize_partition_expr<S>(
73 partition_expr: &Option<PartitionExpr>,
74 serializer: S,
75) -> Result<S::Ok, S::Error>
76where
77 S: serde::Serializer,
78{
79 use serde::ser::Error;
80
81 match partition_expr {
82 None => serializer.serialize_none(),
83 Some(expr) => {
84 let json_str = expr.as_json_str().map_err(S::Error::custom)?;
85 serializer.serialize_some(&json_str)
86 }
87 }
88}
89
90fn deserialize_partition_expr<'de, D>(deserializer: D) -> Result<Option<PartitionExpr>, D::Error>
91where
92 D: serde::Deserializer<'de>,
93{
94 use serde::de::Error;
95
96 let opt_json_str: Option<String> = Option::deserialize(deserializer)?;
97 match opt_json_str {
98 None => Ok(None),
99 Some(json_str) => {
100 if json_str.is_empty() {
101 Ok(None)
103 } else {
104 PartitionExpr::from_json_str(&json_str).map_err(D::Error::custom)
106 }
107 }
108 }
109}
110
111pub type Level = u8;
113pub const MAX_LEVEL: Level = 2;
115pub type IndexTypes = SmallVec<[IndexType; 4]>;
117
118#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
122pub struct RegionFileId {
123 region_id: RegionId,
125 file_id: FileId,
127}
128
129impl RegionFileId {
130 pub fn new(region_id: RegionId, file_id: FileId) -> Self {
132 Self { region_id, file_id }
133 }
134
135 pub fn region_id(&self) -> RegionId {
137 self.region_id
138 }
139
140 pub fn file_id(&self) -> FileId {
142 self.file_id
143 }
144}
145
146impl fmt::Display for RegionFileId {
147 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
148 write!(f, "{}/{}", self.region_id, self.file_id)
149 }
150}
151
152#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
154pub struct RegionIndexId {
155 pub file_id: RegionFileId,
156 pub version: IndexVersion,
157}
158
159impl RegionIndexId {
160 pub fn new(file_id: RegionFileId, version: IndexVersion) -> Self {
161 Self { file_id, version }
162 }
163
164 pub fn region_id(&self) -> RegionId {
165 self.file_id.region_id
166 }
167
168 pub fn file_id(&self) -> FileId {
169 self.file_id.file_id
170 }
171}
172
173impl fmt::Display for RegionIndexId {
174 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
175 if self.version == 0 {
176 write!(f, "{}/{}", self.file_id.region_id, self.file_id.file_id)
177 } else {
178 write!(
179 f,
180 "{}/{}.{}",
181 self.file_id.region_id, self.file_id.file_id, self.version
182 )
183 }
184 }
185}
186
187pub type FileTimeRange = (Timestamp, Timestamp);
190
191pub(crate) fn overlaps(l: &FileTimeRange, r: &FileTimeRange) -> bool {
193 let (l, r) = if l.0 <= r.0 { (l, r) } else { (r, l) };
194 let (_, l_end) = l;
195 let (r_start, _) = r;
196
197 r_start <= l_end
198}
199
200#[derive(Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
202#[serde(default)]
203pub struct FileMeta {
204 pub region_id: RegionId,
206 pub file_id: FileId,
208 pub time_range: FileTimeRange,
211 pub level: Level,
213 pub file_size: u64,
215 pub max_row_group_uncompressed_size: u64,
217 pub available_indexes: IndexTypes,
219 pub indexes: Vec<ColumnIndexMetadata>,
229 pub index_file_size: u64,
231 pub index_version: u64,
235 pub num_rows: u64,
241 pub num_row_groups: u64,
247 pub sequence: Option<NonZeroU64>,
255 #[serde(
263 serialize_with = "serialize_partition_expr",
264 deserialize_with = "deserialize_partition_expr"
265 )]
266 pub partition_expr: Option<PartitionExpr>,
267 pub num_series: u64,
271 #[serde(
274 default,
275 skip_serializing_if = "Option::is_none",
276 serialize_with = "serialize_bytes_option",
277 deserialize_with = "deserialize_bytes_option"
278 )]
279 pub primary_key_min: Option<Bytes>,
280 #[serde(
283 default,
284 skip_serializing_if = "Option::is_none",
285 serialize_with = "serialize_bytes_option",
286 deserialize_with = "deserialize_bytes_option"
287 )]
288 pub primary_key_max: Option<Bytes>,
289 #[serde(default, skip_serializing_if = "is_false")]
294 pub preserve_row_sequence: bool,
295}
296
297fn is_false(value: &bool) -> bool {
298 !*value
299}
300
301struct DebugFmt<F: Fn(&mut Formatter<'_>) -> fmt::Result>(F);
304
305impl<F: Fn(&mut Formatter<'_>) -> fmt::Result> std::fmt::Debug for DebugFmt<F> {
306 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
307 (self.0)(f)
308 }
309}
310
311impl Debug for FileMeta {
312 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
313 let mut debug_struct = f.debug_struct("FileMeta");
314 debug_struct
315 .field("region_id", &self.region_id)
316 .field("file_id", &DebugFmt(|f| write!(f, "{} ", self.file_id)))
317 .field(
318 "time_range",
319 &DebugFmt(|f| {
320 write!(
321 f,
322 "({}, {}) ",
323 self.time_range.0.to_iso8601_string(),
324 self.time_range.1.to_iso8601_string()
325 )
326 }),
327 )
328 .field("level", &self.level)
329 .field("file_size", &ReadableSize(self.file_size))
330 .field(
331 "max_row_group_uncompressed_size",
332 &ReadableSize(self.max_row_group_uncompressed_size),
333 );
334 if !self.available_indexes.is_empty() {
335 debug_struct
336 .field("available_indexes", &self.available_indexes)
337 .field("indexes", &self.indexes)
338 .field("index_file_size", &ReadableSize(self.index_file_size));
339 }
340 debug_struct
341 .field("num_rows", &self.num_rows)
342 .field("num_row_groups", &self.num_row_groups)
343 .field(
344 "sequence",
345 &DebugFmt(|f| match self.sequence {
346 None => {
347 write!(f, "None")
348 }
349 Some(seq) => {
350 write!(f, "{}", seq)
351 }
352 }),
353 )
354 .field("partition_expr", &self.partition_expr)
355 .field("num_series", &self.num_series);
356 if self.primary_key_min.is_some() || self.primary_key_max.is_some() {
357 debug_struct
358 .field(
359 "primary_key_min",
360 &self.primary_key_min.as_ref().map(|b| b.len()),
361 )
362 .field(
363 "primary_key_max",
364 &self.primary_key_max.as_ref().map(|b| b.len()),
365 );
366 }
367 debug_struct.finish()
368 }
369}
370
371#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
373pub enum IndexType {
374 InvertedIndex,
376 FulltextIndex,
378 BloomFilterIndex,
380 VectorIndex,
382}
383
384#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize, Default)]
394#[serde(default)]
395pub struct ColumnIndexMetadata {
396 pub column_id: ColumnId,
398 pub created_indexes: IndexTypes,
400}
401
402impl FileMeta {
403 pub fn primary_key_range(&self) -> Option<(Bytes, Bytes)> {
405 match (&self.primary_key_min, &self.primary_key_max) {
406 (Some(min), Some(max)) => Some((min.clone(), max.clone())),
407 _ => None,
408 }
409 }
410
411 pub fn exists_index(&self) -> bool {
412 !self.available_indexes.is_empty()
413 }
414
415 pub fn index_version(&self) -> Option<IndexVersion> {
416 if self.exists_index() {
417 Some(self.index_version)
418 } else {
419 None
420 }
421 }
422
423 pub fn is_index_up_to_date(&self, other: &FileMeta) -> bool {
425 self.exists_index() && other.exists_index() && self.index_version >= other.index_version
426 }
427
428 pub fn inverted_index_available(&self) -> bool {
430 self.available_indexes.contains(&IndexType::InvertedIndex)
431 }
432
433 pub fn fulltext_index_available(&self) -> bool {
435 self.available_indexes.contains(&IndexType::FulltextIndex)
436 }
437
438 pub fn bloom_filter_index_available(&self) -> bool {
440 self.available_indexes
441 .contains(&IndexType::BloomFilterIndex)
442 }
443
444 pub fn index_file_size(&self) -> u64 {
445 self.index_file_size
446 }
447
448 pub fn is_index_consistent_with_region(&self, metadata: &[ColumnMetadata]) -> bool {
450 let id_to_indexes = self
451 .indexes
452 .iter()
453 .map(|index| (index.column_id, index.created_indexes.clone()))
454 .collect::<std::collections::HashMap<_, _>>();
455 for column in metadata {
456 if !column.column_schema.is_indexed() {
457 continue;
458 }
459 if let Some(indexes) = id_to_indexes.get(&column.column_id) {
460 if column.column_schema.is_inverted_indexed()
461 && !indexes.contains(&IndexType::InvertedIndex)
462 {
463 return false;
464 }
465 if column.column_schema.is_fulltext_indexed()
466 && !indexes.contains(&IndexType::FulltextIndex)
467 {
468 return false;
469 }
470 if column.column_schema.is_skipping_indexed()
471 && !indexes.contains(&IndexType::BloomFilterIndex)
472 {
473 return false;
474 }
475 } else {
476 return false;
477 }
478 }
479 true
480 }
481
482 pub fn file_id(&self) -> RegionFileId {
484 RegionFileId::new(self.region_id, self.file_id)
485 }
486
487 pub fn index_id(&self) -> RegionIndexId {
489 RegionIndexId::new(self.file_id(), self.index_version)
490 }
491}
492
493#[derive(Clone)]
495pub struct FileHandle {
496 inner: Arc<FileHandleInner>,
497}
498
499impl fmt::Debug for FileHandle {
500 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
501 f.debug_struct("FileHandle")
502 .field("meta", self.meta_ref())
503 .field("compacting", &self.compacting())
504 .field("deleted", &self.inner.deleted.load(Ordering::Relaxed))
505 .finish()
506 }
507}
508
509impl FileHandle {
510 pub fn new(meta: FileMeta, file_purger: FilePurgerRef) -> FileHandle {
512 let pk_range = meta.primary_key_range();
513 FileHandle {
514 inner: Arc::new(FileHandleInner::new(meta, file_purger, pk_range)),
515 }
516 }
517
518 #[cfg(test)]
519 pub fn new_with_primary_key_range(
520 meta: FileMeta,
521 file_purger: FilePurgerRef,
522 primary_key_range: Option<(Bytes, Bytes)>,
523 ) -> FileHandle {
524 FileHandle {
525 inner: Arc::new(FileHandleInner::new(meta, file_purger, primary_key_range)),
526 }
527 }
528
529 pub fn region_id(&self) -> RegionId {
531 self.inner.meta.region_id
532 }
533
534 pub(crate) fn is_effective_target_sequence_trusted(&self, target_region_id: RegionId) -> bool {
539 if self.region_id() != target_region_id {
540 self.meta_ref().sequence.is_some()
541 } else {
542 self.meta_ref().preserve_row_sequence
543 }
544 }
545
546 pub fn file_id(&self) -> RegionFileId {
548 RegionFileId::new(self.inner.meta.region_id, self.inner.meta.file_id)
549 }
550
551 pub fn index_id(&self) -> RegionIndexId {
553 RegionIndexId::new(self.file_id(), self.inner.meta.index_version)
554 }
555
556 pub fn file_path(&self, table_dir: &str, path_type: PathType) -> String {
558 location::sst_file_path(table_dir, self.file_id(), path_type)
559 }
560
561 pub fn time_range(&self) -> FileTimeRange {
563 self.inner.meta.time_range
564 }
565
566 pub fn mark_deleted(&self) {
568 self.inner.deleted.store(true, Ordering::Relaxed);
569 }
570
571 pub fn compacting(&self) -> bool {
572 self.inner.compacting.load(Ordering::Relaxed)
573 }
574
575 pub fn set_compacting(&self, compacting: bool) {
576 self.inner.compacting.store(compacting, Ordering::Relaxed);
577 }
578
579 pub fn try_set_compacting(&self) -> bool {
581 self.inner
582 .compacting
583 .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
584 .is_ok()
585 }
586
587 pub fn index_outdated(&self) -> bool {
588 self.inner.index_outdated.load(Ordering::Relaxed)
589 }
590
591 pub fn set_index_outdated(&self, index_outdated: bool) {
592 self.inner
593 .index_outdated
594 .store(index_outdated, Ordering::Relaxed);
595 }
596
597 pub fn meta_ref(&self) -> &FileMeta {
599 &self.inner.meta
600 }
601
602 pub fn file_purger(&self) -> FilePurgerRef {
603 self.inner.file_purger.clone()
604 }
605
606 pub fn size(&self) -> u64 {
607 self.inner.meta.file_size
608 }
609
610 pub fn index_size(&self) -> u64 {
611 self.inner.meta.index_file_size
612 }
613
614 pub fn num_rows(&self) -> usize {
615 self.inner.meta.num_rows as usize
616 }
617
618 pub fn level(&self) -> Level {
619 self.inner.meta.level
620 }
621
622 pub fn is_deleted(&self) -> bool {
623 self.inner.deleted.load(Ordering::Relaxed)
624 }
625
626 pub(crate) fn primary_key_range(
629 &self,
630 mapper: &PrimaryKeyRangeMapper,
631 ) -> Option<(Bytes, Bytes)> {
632 debug_assert_eq!(self.region_id().table_id(), mapper.region_id().table_id());
633 if let Some(range) = self
634 .inner
635 .primary_key_range
636 .read()
637 .unwrap()
638 .aligned(mapper.schema_version())
639 {
640 return range.clone();
641 }
642 let aligned = self.inner.primary_key_range.write().unwrap().align(mapper);
644 match aligned {
645 Ok(range) => range,
646 Err(err) => {
647 warn!(err; "Invalid SST primary key range; using unknown bounds, region: {}, file: {}, schema version: {}",
648 self.region_id(), self.file_id(), mapper.schema_version());
649 None
650 }
651 }
652 }
653
654 pub fn raw_primary_key_range(&self) -> Option<(Bytes, Bytes)> {
656 self.inner.primary_key_range.read().unwrap().raw().cloned()
657 }
658
659 pub(crate) fn set_primary_key_range(&self, primary_key_range: (Bytes, Bytes)) {
660 let mut range = self.inner.primary_key_range.write().unwrap();
663 if matches!(*range, PrimaryKeyRange::Missing) {
664 *range = PrimaryKeyRange::Raw(primary_key_range);
665 }
666 }
667}
668
669type PrimaryKeyBounds = (Bytes, Bytes);
670
671enum PrimaryKeyRange {
674 Missing,
675 Raw(PrimaryKeyBounds),
676 Aligned {
677 raw: PrimaryKeyBounds,
678 schema_version: u64,
679 bounds: Option<PrimaryKeyBounds>,
680 },
681}
682
683impl PrimaryKeyRange {
684 fn raw(&self) -> Option<&PrimaryKeyBounds> {
685 match self {
686 Self::Missing => None,
687 Self::Raw(raw) | Self::Aligned { raw, .. } => Some(raw),
688 }
689 }
690
691 fn aligned(&self, target_version: u64) -> Option<&Option<PrimaryKeyBounds>> {
692 match self {
693 Self::Aligned {
694 schema_version,
695 bounds,
696 ..
697 } if *schema_version == target_version => Some(bounds),
698 _ => None,
699 }
700 }
701
702 fn align(
703 &mut self,
704 mapper: &PrimaryKeyRangeMapper,
705 ) -> crate::error::Result<Option<PrimaryKeyBounds>> {
706 if let Some(bounds) = self.aligned(mapper.schema_version()) {
707 return Ok(bounds.clone());
708 }
709 let Some(raw) = self.raw().cloned() else {
710 return Ok(None);
711 };
712 let aligned = mapper.map(raw.clone());
713 *self = Self::Aligned {
715 raw,
716 schema_version: mapper.schema_version(),
717 bounds: aligned.as_ref().ok().cloned().flatten(),
718 };
719 aligned
720 }
721}
722
723struct FileHandleInner {
727 meta: FileMeta,
728 compacting: AtomicBool,
729 deleted: AtomicBool,
730 index_outdated: AtomicBool,
731 primary_key_range: RwLock<PrimaryKeyRange>,
732 file_purger: FilePurgerRef,
733}
734
735impl Drop for FileHandleInner {
736 fn drop(&mut self) {
737 self.file_purger.remove_file(
738 self.meta.clone(),
739 self.deleted.load(Ordering::Acquire),
740 self.index_outdated.load(Ordering::Acquire),
741 );
742 }
743}
744
745impl FileHandleInner {
746 fn new(
748 meta: FileMeta,
749 file_purger: FilePurgerRef,
750 primary_key_range: Option<(Bytes, Bytes)>,
751 ) -> FileHandleInner {
752 file_purger.new_file(&meta);
753 FileHandleInner {
754 meta,
755 compacting: AtomicBool::new(false),
756 deleted: AtomicBool::new(false),
757 index_outdated: AtomicBool::new(false),
758 primary_key_range: RwLock::new(
759 primary_key_range.map_or(PrimaryKeyRange::Missing, PrimaryKeyRange::Raw),
760 ),
761 file_purger,
762 }
763 }
764}
765
766pub async fn delete_files(
773 region_id: RegionId,
774 file_ids: &[(FileId, u64)],
775 delete_index: bool,
776 access_layer: &AccessLayerRef,
777 cache_manager: &Option<CacheManagerRef>,
778) -> crate::error::Result<()> {
779 if let Some(cache) = &cache_manager {
781 for (file_id, _) in file_ids {
782 cache.remove_parquet_meta_data(RegionFileId::new(region_id, *file_id));
783 }
784 }
785 let mut attempted_files = Vec::with_capacity(file_ids.len());
786 let mut index_ids = Vec::new();
787
788 for (file_id, index_version) in file_ids {
789 let region_file_id = RegionFileId::new(region_id, *file_id);
790 attempted_files.push(*file_id);
791 index_ids.extend(
792 (0..=*index_version).map(|version| RegionIndexId::new(region_file_id, version)),
793 );
794 }
795
796 access_layer
797 .delete_ssts(region_id, &attempted_files)
798 .await?;
799 access_layer.delete_indexes(&index_ids).await?;
800
801 debug!(
802 "Attempted to delete {} files for region {}: {:?}",
803 attempted_files.len(),
804 region_id,
805 attempted_files
806 );
807
808 for (file_id, index_version) in file_ids {
809 purge_index_cache_stager(
810 region_id,
811 delete_index,
812 access_layer,
813 cache_manager,
814 *file_id,
815 *index_version,
816 )
817 .await;
818 }
819 Ok(())
820}
821
822#[derive(Clone)]
827pub(crate) struct UncommittedSsts {
828 region_id: RegionId,
829 files: Arc<Mutex<HashMap<FileId, (u64, bool)>>>,
830 access_layer: AccessLayerRef,
831 cache_manager: Option<CacheManagerRef>,
832}
833
834impl UncommittedSsts {
835 pub(crate) fn new(
836 region_id: RegionId,
837 access_layer: AccessLayerRef,
838 cache_manager: Option<CacheManagerRef>,
839 ) -> Self {
840 Self {
841 region_id,
842 files: Arc::new(Mutex::new(HashMap::new())),
843 access_layer,
844 cache_manager,
845 }
846 }
847
848 pub(crate) fn track(&self, ssts: &[SstInfo]) {
850 let mut files = self.files.lock().unwrap();
851 for sst in ssts {
852 files.insert(
853 sst.file_id,
854 (sst.index_metadata.version, sst.index_metadata.file_size > 0),
855 );
856 }
857 }
858
859 pub(crate) fn disarm_cleanup(&self) {
864 self.files.lock().unwrap().clear();
865 }
866
867 #[cfg(test)]
868 pub(crate) fn num_tracked_files(&self) -> usize {
869 self.files.lock().unwrap().len()
870 }
871
872 pub(crate) async fn cleanup(&self) {
874 if let Err(err) = self.try_cleanup().await {
875 error!(err; "Failed to clean uncommitted SSTs for region {}", self.region_id);
876 }
877 }
878
879 async fn try_cleanup(&self) -> crate::error::Result<()> {
880 let files = std::mem::take(&mut *self.files.lock().unwrap());
881 if files.is_empty() {
882 return Ok(());
883 }
884
885 let delete_index = files.values().any(|(_, exists_index)| *exists_index);
886 let file_ids = files
887 .into_iter()
888 .map(|(file_id, (index_version, _))| (file_id, index_version))
889 .collect::<Vec<_>>();
890 delete_files(
891 self.region_id,
892 &file_ids,
893 delete_index,
894 &self.access_layer,
895 &self.cache_manager,
896 )
897 .await
898 }
899
900 #[cfg(test)]
901 pub(crate) async fn cleanup_for_test(&self) -> crate::error::Result<()> {
902 self.try_cleanup().await
903 }
904}
905
906pub async fn delete_index(
907 region_index_id: RegionIndexId,
908 access_layer: &AccessLayerRef,
909 cache_manager: &Option<CacheManagerRef>,
910) -> crate::error::Result<()> {
911 delete_index_and_purge(region_index_id, access_layer, cache_manager).await?;
912
913 Ok(())
914}
915
916pub async fn delete_indexes(
917 index_ids: &[RegionIndexId],
918 access_layer: &AccessLayerRef,
919 cache_manager: &Option<CacheManagerRef>,
920) -> crate::error::Result<()> {
921 if index_ids.is_empty() {
922 return Ok(());
923 }
924
925 if let Err(e) = access_layer.delete_indexes(index_ids).await {
926 error!(e; "Failed to batch delete index files");
927
928 for index_id in index_ids {
929 delete_index_and_purge(*index_id, access_layer, cache_manager).await?;
930 }
931
932 return Ok(());
933 }
934
935 purge_indexes(index_ids, access_layer, cache_manager).await;
936
937 Ok(())
938}
939
940async fn delete_index_and_purge(
941 index_id: RegionIndexId,
942 access_layer: &AccessLayerRef,
943 cache_manager: &Option<CacheManagerRef>,
944) -> crate::error::Result<()> {
945 access_layer.delete_index(index_id).await?;
946 purge_index_cache_stager(
947 index_id.region_id(),
948 true,
949 access_layer,
950 cache_manager,
951 index_id.file_id(),
952 index_id.version,
953 )
954 .await;
955 Ok(())
956}
957
958async fn purge_indexes(
959 index_ids: &[RegionIndexId],
960 access_layer: &AccessLayerRef,
961 cache_manager: &Option<CacheManagerRef>,
962) {
963 for index_id in index_ids {
964 purge_index_cache_stager(
965 index_id.region_id(),
966 true,
967 access_layer,
968 cache_manager,
969 index_id.file_id(),
970 index_id.version,
971 )
972 .await;
973 }
974}
975
976async fn purge_index_cache_stager(
977 region_id: RegionId,
978 delete_index: bool,
979 access_layer: &AccessLayerRef,
980 cache_manager: &Option<CacheManagerRef>,
981 file_id: FileId,
982 index_version: u64,
983) {
984 if let Some(write_cache) = cache_manager.as_ref().and_then(|cache| cache.write_cache()) {
985 if delete_index {
987 write_cache
988 .remove(IndexKey::new(
989 region_id,
990 file_id,
991 FileType::Puffin(index_version),
992 ))
993 .await;
994 }
995
996 write_cache
998 .remove(IndexKey::new(region_id, file_id, FileType::Parquet))
999 .await;
1000 }
1001
1002 if let Err(e) = access_layer
1004 .puffin_manager_factory()
1005 .purge_stager(RegionIndexId::new(
1006 RegionFileId::new(region_id, file_id),
1007 index_version,
1008 ))
1009 .await
1010 {
1011 error!(e; "Failed to purge stager with index file, file_id: {}, index_version: {}, region: {}",
1012 file_id, index_version, region_id);
1013 }
1014}
1015
1016#[cfg(test)]
1017mod tests {
1018 use std::str::FromStr;
1019
1020 use datatypes::prelude::ConcreteDataType;
1021 use datatypes::schema::{
1022 ColumnSchema, FulltextAnalyzer, FulltextBackend, FulltextOptions, SkippingIndexOptions,
1023 };
1024 use datatypes::value::Value;
1025 use partition::expr::{PartitionExpr, col};
1026
1027 use super::*;
1028
1029 fn create_file_meta(file_id: FileId, level: Level) -> FileMeta {
1030 FileMeta {
1031 region_id: 0.into(),
1032 file_id,
1033 time_range: FileTimeRange::default(),
1034 level,
1035 file_size: 0,
1036 max_row_group_uncompressed_size: 0,
1037 available_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1038 indexes: vec![ColumnIndexMetadata {
1039 column_id: 0,
1040 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1041 }],
1042 index_file_size: 0,
1043 index_version: 0,
1044 num_rows: 0,
1045 num_row_groups: 0,
1046 sequence: None,
1047 partition_expr: None,
1048 num_series: 0,
1049 ..Default::default()
1050 }
1051 }
1052
1053 #[test]
1054 fn test_deserialize_file_meta() {
1055 let file_meta = create_file_meta(FileId::random(), 0);
1056 let serialized_file_meta = serde_json::to_string(&file_meta).unwrap();
1057 let deserialized_file_meta = serde_json::from_str(&serialized_file_meta);
1058 assert_eq!(file_meta, deserialized_file_meta.unwrap());
1059 }
1060
1061 #[test]
1062 fn test_deserialize_from_string() {
1063 let json_file_meta = "{\"region_id\":0,\"file_id\":\"bc5896ec-e4d8-4017-a80d-f2de73188d55\",\
1064 \"time_range\":[{\"value\":0,\"unit\":\"Millisecond\"},{\"value\":0,\"unit\":\"Millisecond\"}],\
1065 \"available_indexes\":[\"InvertedIndex\"],\"indexes\":[{\"column_id\": 0, \"created_indexes\": [\"InvertedIndex\"]}],\"level\":0}";
1066 let file_meta = create_file_meta(
1067 FileId::from_str("bc5896ec-e4d8-4017-a80d-f2de73188d55").unwrap(),
1068 0,
1069 );
1070 let deserialized_file_meta: FileMeta = serde_json::from_str(json_file_meta).unwrap();
1071 assert_eq!(file_meta, deserialized_file_meta);
1072 }
1073
1074 #[test]
1075 fn test_deserialize_legacy_vector_index() {
1076 let json = r#"{
1077 "region_id": 0,
1078 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1079 "time_range": [{"value":0,"unit":"Millisecond"},{"value":0,"unit":"Millisecond"}],
1080 "available_indexes": ["VectorIndex", "InvertedIndex"],
1081 "indexes": [{"column_id": 1, "created_indexes": ["VectorIndex"]}],
1082 "level": 0
1083 }"#;
1084 let meta: FileMeta = serde_json::from_str(json).unwrap();
1085 assert!(meta.inverted_index_available());
1086 assert!(!meta.fulltext_index_available());
1087 assert!(!meta.bloom_filter_index_available());
1088 let encoded = serde_json::to_value(&meta).unwrap();
1089 assert_eq!(
1090 encoded["available_indexes"],
1091 serde_json::json!(["VectorIndex", "InvertedIndex"])
1092 );
1093 assert_eq!(
1094 encoded["indexes"][0]["created_indexes"],
1095 serde_json::json!(["VectorIndex"])
1096 );
1097 }
1098
1099 #[test]
1100 fn test_file_meta_with_partition_expr() {
1101 let file_id = FileId::random();
1102 let partition_expr = PartitionExpr::new(
1103 col("a"),
1104 partition::expr::RestrictedOp::GtEq,
1105 Value::UInt32(10).into(),
1106 );
1107
1108 let file_meta_with_partition = FileMeta {
1109 region_id: 0.into(),
1110 file_id,
1111 time_range: FileTimeRange::default(),
1112 level: 0,
1113 file_size: 0,
1114 max_row_group_uncompressed_size: 0,
1115 available_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1116 indexes: vec![ColumnIndexMetadata {
1117 column_id: 0,
1118 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1119 }],
1120 index_file_size: 0,
1121 index_version: 0,
1122 num_rows: 0,
1123 num_row_groups: 0,
1124 sequence: None,
1125 partition_expr: Some(partition_expr.clone()),
1126 num_series: 0,
1127 ..Default::default()
1128 };
1129
1130 let serialized = serde_json::to_string(&file_meta_with_partition).unwrap();
1132 let deserialized: FileMeta = serde_json::from_str(&serialized).unwrap();
1133 assert_eq!(file_meta_with_partition, deserialized);
1134
1135 let serialized_value: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1137 assert!(serialized_value["partition_expr"].as_str().is_some());
1138 let partition_expr_json = serialized_value["partition_expr"].as_str().unwrap();
1139 assert!(partition_expr_json.contains("\"Column\":\"a\""));
1140 assert!(partition_expr_json.contains("\"op\":\"GtEq\""));
1141
1142 let file_meta_none = FileMeta {
1144 partition_expr: None,
1145 ..file_meta_with_partition.clone()
1146 };
1147 let serialized_none = serde_json::to_string(&file_meta_none).unwrap();
1148 let deserialized_none: FileMeta = serde_json::from_str(&serialized_none).unwrap();
1149 assert_eq!(file_meta_none, deserialized_none);
1150 }
1151
1152 #[test]
1153 fn test_file_meta_partition_expr_backward_compatibility() {
1154 let json_with_partition_expr = r#"{
1156 "region_id": 0,
1157 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1158 "time_range": [
1159 {"value": 0, "unit": "Millisecond"},
1160 {"value": 0, "unit": "Millisecond"}
1161 ],
1162 "level": 0,
1163 "file_size": 0,
1164 "available_indexes": ["InvertedIndex"],
1165 "index_file_size": 0,
1166 "num_rows": 0,
1167 "num_row_groups": 0,
1168 "sequence": null,
1169 "partition_expr": "{\"Expr\":{\"lhs\":{\"Column\":\"a\"},\"op\":\"GtEq\",\"rhs\":{\"Value\":{\"UInt32\":10}}}}"
1170 }"#;
1171
1172 let file_meta: FileMeta = serde_json::from_str(json_with_partition_expr).unwrap();
1173 assert!(file_meta.partition_expr.is_some());
1174 let expr = file_meta.partition_expr.unwrap();
1175 assert_eq!(format!("{}", expr), "a >= 10");
1176
1177 let json_with_empty_expr = r#"{
1179 "region_id": 0,
1180 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1181 "time_range": [
1182 {"value": 0, "unit": "Millisecond"},
1183 {"value": 0, "unit": "Millisecond"}
1184 ],
1185 "level": 0,
1186 "file_size": 0,
1187 "available_indexes": [],
1188 "index_file_size": 0,
1189 "num_rows": 0,
1190 "num_row_groups": 0,
1191 "sequence": null,
1192 "partition_expr": ""
1193 }"#;
1194
1195 let file_meta_empty: FileMeta = serde_json::from_str(json_with_empty_expr).unwrap();
1196 assert!(file_meta_empty.partition_expr.is_none());
1197
1198 let json_with_null_expr = r#"{
1200 "region_id": 0,
1201 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1202 "time_range": [
1203 {"value": 0, "unit": "Millisecond"},
1204 {"value": 0, "unit": "Millisecond"}
1205 ],
1206 "level": 0,
1207 "file_size": 0,
1208 "available_indexes": [],
1209 "index_file_size": 0,
1210 "num_rows": 0,
1211 "num_row_groups": 0,
1212 "sequence": null,
1213 "partition_expr": null
1214 }"#;
1215
1216 let file_meta_null: FileMeta = serde_json::from_str(json_with_null_expr).unwrap();
1217 assert!(file_meta_null.partition_expr.is_none());
1218
1219 let json_with_empty_expr = r#"{
1221 "region_id": 0,
1222 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1223 "time_range": [
1224 {"value": 0, "unit": "Millisecond"},
1225 {"value": 0, "unit": "Millisecond"}
1226 ],
1227 "level": 0,
1228 "file_size": 0,
1229 "available_indexes": [],
1230 "index_file_size": 0,
1231 "num_rows": 0,
1232 "num_row_groups": 0,
1233 "sequence": null
1234 }"#;
1235
1236 let file_meta_empty: FileMeta = serde_json::from_str(json_with_empty_expr).unwrap();
1237 assert!(file_meta_empty.partition_expr.is_none());
1238 }
1239
1240 #[test]
1241 fn test_file_meta_indexes_backward_compatibility() {
1242 let json_old_file_meta = r#"{
1244 "region_id": 0,
1245 "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1246 "time_range": [
1247 {"value": 0, "unit": "Millisecond"},
1248 {"value": 0, "unit": "Millisecond"}
1249 ],
1250 "available_indexes": ["InvertedIndex"],
1251 "level": 0,
1252 "file_size": 0,
1253 "index_file_size": 0,
1254 "num_rows": 0,
1255 "num_row_groups": 0
1256 }"#;
1257
1258 let deserialized_file_meta: FileMeta = serde_json::from_str(json_old_file_meta).unwrap();
1259
1260 assert_eq!(deserialized_file_meta.indexes, vec![]);
1262
1263 let expected_indexes: IndexTypes = SmallVec::from_iter([IndexType::InvertedIndex]);
1264 assert_eq!(deserialized_file_meta.available_indexes, expected_indexes);
1265
1266 assert_eq!(
1267 deserialized_file_meta.file_id,
1268 FileId::from_str("bc5896ec-e4d8-4017-a80d-f2de73188d55").unwrap()
1269 );
1270 assert!(!deserialized_file_meta.preserve_row_sequence);
1271 }
1272
1273 #[test]
1274 fn test_file_meta_preserve_row_sequence_serde() {
1275 let file_meta = FileMeta {
1276 preserve_row_sequence: true,
1277 ..Default::default()
1278 };
1279
1280 let serialized = serde_json::to_string(&file_meta).unwrap();
1281 let value: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1282 assert_eq!(value["preserve_row_sequence"], true);
1283
1284 let deserialized: FileMeta = serde_json::from_str(&serialized).unwrap();
1285 assert_eq!(file_meta, deserialized);
1286
1287 let file_meta_false = FileMeta {
1288 preserve_row_sequence: false,
1289 ..file_meta.clone()
1290 };
1291 let serialized_false = serde_json::to_string(&file_meta_false).unwrap();
1292 let value_false: serde_json::Value = serde_json::from_str(&serialized_false).unwrap();
1293 assert!(value_false.get("preserve_row_sequence").is_none());
1294 }
1295 #[test]
1296 fn test_is_index_consistent_with_region() {
1297 fn new_column_meta(
1298 id: ColumnId,
1299 name: &str,
1300 inverted: bool,
1301 fulltext: bool,
1302 skipping: bool,
1303 ) -> ColumnMetadata {
1304 let mut column_schema =
1305 ColumnSchema::new(name, ConcreteDataType::string_datatype(), true);
1306 if inverted {
1307 column_schema = column_schema.with_inverted_index(true);
1308 }
1309 if fulltext {
1310 column_schema = column_schema
1311 .with_fulltext_options(FulltextOptions::new_unchecked(
1312 true,
1313 FulltextAnalyzer::English,
1314 false,
1315 FulltextBackend::Bloom,
1316 1000,
1317 0.01,
1318 ))
1319 .unwrap();
1320 }
1321 if skipping {
1322 column_schema = column_schema
1323 .with_skipping_options(SkippingIndexOptions::new_unchecked(
1324 1024,
1325 0.01,
1326 datatypes::schema::SkippingIndexType::BloomFilter,
1327 ))
1328 .unwrap();
1329 }
1330
1331 ColumnMetadata {
1332 column_schema,
1333 semantic_type: api::v1::SemanticType::Tag,
1334 column_id: id,
1335 }
1336 }
1337
1338 let mut file_meta = FileMeta {
1340 indexes: vec![ColumnIndexMetadata {
1341 column_id: 1,
1342 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1343 }],
1344 ..Default::default()
1345 };
1346 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1347 assert!(file_meta.is_index_consistent_with_region(®ion_meta));
1348
1349 file_meta.indexes = vec![ColumnIndexMetadata {
1351 column_id: 1,
1352 created_indexes: SmallVec::from_iter([
1353 IndexType::InvertedIndex,
1354 IndexType::BloomFilterIndex,
1355 ]),
1356 }];
1357 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1358 assert!(file_meta.is_index_consistent_with_region(®ion_meta));
1359
1360 file_meta.indexes = vec![ColumnIndexMetadata {
1362 column_id: 1,
1363 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1364 }];
1365 let region_meta = vec![new_column_meta(1, "tag1", true, true, false)]; assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1367
1368 file_meta.indexes = vec![ColumnIndexMetadata {
1370 column_id: 2, created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1372 }];
1373 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)]; assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1375
1376 file_meta.indexes = vec![ColumnIndexMetadata {
1378 column_id: 1,
1379 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1380 }];
1381 let region_meta = vec![new_column_meta(1, "tag1", false, false, false)]; assert!(file_meta.is_index_consistent_with_region(®ion_meta));
1383
1384 file_meta.indexes = vec![];
1386 let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1387 assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1388
1389 file_meta.indexes = vec![
1391 ColumnIndexMetadata {
1392 column_id: 1,
1393 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1394 },
1395 ColumnIndexMetadata {
1396 column_id: 2, created_indexes: SmallVec::from_iter([IndexType::FulltextIndex]),
1398 },
1399 ];
1400 let region_meta = vec![
1401 new_column_meta(1, "tag1", true, false, false),
1402 new_column_meta(2, "tag2", false, true, true), ];
1404 assert!(!file_meta.is_index_consistent_with_region(®ion_meta));
1405 }
1406}