1pub(crate) mod bloom_filter;
16pub(crate) mod fulltext_index;
17mod indexer;
18pub mod intermediate;
19pub(crate) mod inverted_index;
20pub mod puffin_manager;
21mod statistics;
22pub(crate) mod store;
23#[cfg(feature = "vector_index")]
24pub(crate) mod vector_index;
25
26use std::cmp::Ordering;
27use std::collections::{HashMap, HashSet};
28use std::num::NonZeroUsize;
29use std::sync::Arc;
30
31use bloom_filter::creator::BloomFilterIndexer;
32use common_telemetry::{debug, error, info, warn};
33use datatypes::arrow::array::BinaryArray;
34use datatypes::arrow::record_batch::RecordBatch;
35use mito_codec::index::IndexValuesCodec;
36use mito_codec::row_converter::CompositeValues;
37use object_store::ObjectStore;
38use puffin_manager::SstPuffinManager;
39use smallvec::{SmallVec, smallvec};
40use snafu::{OptionExt, ResultExt};
41use statistics::{ByteCount, RowCount};
42use store_api::metadata::RegionMetadataRef;
43use store_api::storage::{ColumnId, FileId, RegionId};
44use strum::IntoStaticStr;
45use tokio::sync::mpsc::Sender;
46#[cfg(feature = "vector_index")]
47use vector_index::creator::VectorIndexer;
48
49use crate::access_layer::{AccessLayerRef, FilePathProvider, OperationType, RegionFilePathFactory};
50use crate::cache::file_cache::{FileCacheRef, FileType, IndexKey};
51use crate::cache::write_cache::{UploadTracker, WriteCacheRef};
52use crate::cache::{CacheManagerRef, CacheStrategy};
53#[cfg(feature = "vector_index")]
54use crate::config::VectorIndexConfig;
55use crate::config::{BloomFilterConfig, FulltextIndexConfig, InvertedIndexConfig};
56use crate::error::{
57 BuildIndexAsyncSnafu, DecodeSnafu, Error, InvalidRecordBatchSnafu, RegionClosedSnafu,
58 RegionDroppedSnafu, RegionTruncatedSnafu, Result,
59};
60use crate::metrics::{
61 INDEX_ARTIFACT_CLEANUP_FAILURE_TOTAL, INDEX_CREATE_MEMORY_USAGE, INDEX_PUBLICATION_STALE_TOTAL,
62};
63use crate::read::Batch;
64use crate::region::options::IndexOptions;
65use crate::region::version::VersionControlRef;
66use crate::region::{
67 IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContextRef,
68};
69use crate::request::{
70 BackgroundNotify, BuildIndexRequest, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped,
71 WorkerRequest, WorkerRequestWithTime,
72};
73use crate::schedule::scheduler::{Job, SchedulerRef};
74use crate::sst::file::{
75 ColumnIndexMetadata, FileHandle, IndexType, IndexTypes, RegionFileId, RegionIndexId,
76};
77use crate::sst::file_purger::FilePurgerRef;
78use crate::sst::index::fulltext_index::creator::FulltextIndexer;
79use crate::sst::index::intermediate::IntermediateManager;
80use crate::sst::index::inverted_index::creator::InvertedIndexer;
81use crate::sst::parquet::SstInfo;
82use crate::sst::parquet::flat_format::primary_key_column_index;
83use crate::sst::parquet::format::PrimaryKeyArray;
84use crate::worker::WorkerListener;
85
86pub(crate) const TYPE_INVERTED_INDEX: &str = "inverted_index";
87pub(crate) const TYPE_FULLTEXT_INDEX: &str = "fulltext_index";
88pub(crate) const TYPE_BLOOM_FILTER_INDEX: &str = "bloom_filter_index";
89#[cfg(feature = "vector_index")]
90pub(crate) const TYPE_VECTOR_INDEX: &str = "vector_index";
91
92pub(crate) async fn cleanup_stale_index_caches(
98 index_id: RegionIndexId,
99 access_layer: &AccessLayerRef,
100 cache_manager: Option<&CacheManagerRef>,
101 write_cache: Option<&WriteCacheRef>,
102) {
103 if let Some(cache_manager) = cache_manager {
104 CacheStrategy::EnableAll(cache_manager.clone())
105 .evict_puffin_cache(index_id)
106 .await;
107 }
108 if let Some(write_cache) = write_cache {
109 write_cache
110 .remove(IndexKey::new(
111 index_id.region_id(),
112 index_id.file_id(),
113 FileType::Puffin(index_id.version),
114 ))
115 .await;
116 }
117 if let Err(e) = access_layer
118 .puffin_manager_factory()
119 .purge_stager(index_id)
120 .await
121 {
122 INDEX_ARTIFACT_CLEANUP_FAILURE_TOTAL.inc();
123 warn!(
124 e;
125 "Failed to purge stale index artifact from stager, region: {}, file_id: {}, index_version: {}",
126 index_id.region_id(),
127 index_id.file_id(),
128 index_id.version
129 );
130 }
131}
132
133pub(crate) fn trigger_index_background_download(
135 file_cache: Option<&FileCacheRef>,
136 file_id: &RegionIndexId,
137 file_size_hint: Option<u64>,
138 path_factory: &RegionFilePathFactory,
139 object_store: &ObjectStore,
140) {
141 if let (Some(file_cache), Some(file_size)) = (file_cache, file_size_hint) {
142 let index_key = IndexKey::new(
143 file_id.region_id(),
144 file_id.file_id(),
145 FileType::Puffin(file_id.version),
146 );
147 let remote_path = path_factory.build_index_file_path(file_id.file_id);
148 file_cache.maybe_download_background(
149 index_key,
150 remote_path,
151 object_store.clone(),
152 file_size,
153 );
154 }
155}
156
157#[derive(Debug, Clone, Default)]
159pub struct IndexOutput {
160 pub file_size: u64,
162 pub version: u64,
164 pub inverted_index: InvertedIndexOutput,
166 pub fulltext_index: FulltextIndexOutput,
168 pub bloom_filter: BloomFilterOutput,
170 #[cfg(feature = "vector_index")]
172 pub vector_index: VectorIndexOutput,
173}
174
175impl IndexOutput {
176 pub fn build_available_indexes(&self) -> SmallVec<[IndexType; 4]> {
177 let mut indexes = SmallVec::new();
178 if self.inverted_index.is_available() {
179 indexes.push(IndexType::InvertedIndex);
180 }
181 if self.fulltext_index.is_available() {
182 indexes.push(IndexType::FulltextIndex);
183 }
184 if self.bloom_filter.is_available() {
185 indexes.push(IndexType::BloomFilterIndex);
186 }
187 #[cfg(feature = "vector_index")]
188 if self.vector_index.is_available() {
189 indexes.push(IndexType::VectorIndex);
190 }
191 indexes
192 }
193
194 pub fn build_indexes(&self) -> Vec<ColumnIndexMetadata> {
195 let mut map: HashMap<ColumnId, IndexTypes> = HashMap::new();
196
197 if self.inverted_index.is_available() {
198 for &col in &self.inverted_index.columns {
199 map.entry(col).or_default().push(IndexType::InvertedIndex);
200 }
201 }
202 if self.fulltext_index.is_available() {
203 for &col in &self.fulltext_index.columns {
204 map.entry(col).or_default().push(IndexType::FulltextIndex);
205 }
206 }
207 if self.bloom_filter.is_available() {
208 for &col in &self.bloom_filter.columns {
209 map.entry(col)
210 .or_default()
211 .push(IndexType::BloomFilterIndex);
212 }
213 }
214 #[cfg(feature = "vector_index")]
215 if self.vector_index.is_available() {
216 for &col in &self.vector_index.columns {
217 map.entry(col).or_default().push(IndexType::VectorIndex);
218 }
219 }
220
221 map.into_iter()
222 .map(|(column_id, created_indexes)| ColumnIndexMetadata {
223 column_id,
224 created_indexes,
225 })
226 .collect::<Vec<_>>()
227 }
228}
229
230#[derive(Debug, Clone, Default)]
232pub struct IndexBaseOutput {
233 pub index_size: ByteCount,
235 pub row_count: RowCount,
237 pub columns: Vec<ColumnId>,
239}
240
241impl IndexBaseOutput {
242 pub fn is_available(&self) -> bool {
243 self.index_size > 0
244 }
245}
246
247pub type InvertedIndexOutput = IndexBaseOutput;
249pub type FulltextIndexOutput = IndexBaseOutput;
251pub type BloomFilterOutput = IndexBaseOutput;
253#[cfg(feature = "vector_index")]
255pub type VectorIndexOutput = IndexBaseOutput;
256
257#[derive(Default)]
259pub struct Indexer {
260 file_id: FileId,
261 region_id: RegionId,
263 physical_region_id: RegionId,
265 index_version: u64,
266 puffin_manager: Option<SstPuffinManager>,
267 write_cache_enabled: bool,
268 inverted_indexer: Option<InvertedIndexer>,
269 last_mem_inverted_index: usize,
270 fulltext_indexer: Option<FulltextIndexer>,
271 last_mem_fulltext_index: usize,
272 bloom_filter_indexer: Option<BloomFilterIndexer>,
273 last_mem_bloom_filter: usize,
274 #[cfg(feature = "vector_index")]
275 vector_indexer: Option<VectorIndexer>,
276 #[cfg(feature = "vector_index")]
277 last_mem_vector_index: usize,
278 intermediate_manager: Option<IntermediateManager>,
279}
280
281impl Indexer {
282 pub async fn update(&mut self, batch: &mut Batch) {
284 self.do_update(batch).await;
285
286 self.flush_mem_metrics();
287 }
288
289 pub async fn update_flat(&mut self, batch: &RecordBatch) {
291 self.do_update_flat(batch).await;
292
293 self.flush_mem_metrics();
294 }
295
296 pub async fn finish(&mut self) -> IndexOutput {
298 let output = self.do_finish().await;
299
300 self.flush_mem_metrics();
301 output
302 }
303
304 pub async fn abort(&mut self) {
306 self.do_abort().await;
307
308 self.flush_mem_metrics();
309 }
310
311 fn flush_mem_metrics(&mut self) {
312 let inverted_mem = self
313 .inverted_indexer
314 .as_ref()
315 .map_or(0, |creator| creator.memory_usage());
316 INDEX_CREATE_MEMORY_USAGE
317 .with_label_values(&[TYPE_INVERTED_INDEX])
318 .add(inverted_mem as i64 - self.last_mem_inverted_index as i64);
319 self.last_mem_inverted_index = inverted_mem;
320
321 let fulltext_mem = self
322 .fulltext_indexer
323 .as_ref()
324 .map_or(0, |creator| creator.memory_usage());
325 INDEX_CREATE_MEMORY_USAGE
326 .with_label_values(&[TYPE_FULLTEXT_INDEX])
327 .add(fulltext_mem as i64 - self.last_mem_fulltext_index as i64);
328 self.last_mem_fulltext_index = fulltext_mem;
329
330 let bloom_filter_mem = self
331 .bloom_filter_indexer
332 .as_ref()
333 .map_or(0, |creator| creator.memory_usage());
334 INDEX_CREATE_MEMORY_USAGE
335 .with_label_values(&[TYPE_BLOOM_FILTER_INDEX])
336 .add(bloom_filter_mem as i64 - self.last_mem_bloom_filter as i64);
337 self.last_mem_bloom_filter = bloom_filter_mem;
338
339 #[cfg(feature = "vector_index")]
340 {
341 let vector_mem = self
342 .vector_indexer
343 .as_ref()
344 .map_or(0, |creator| creator.memory_usage());
345 INDEX_CREATE_MEMORY_USAGE
346 .with_label_values(&[TYPE_VECTOR_INDEX])
347 .add(vector_mem as i64 - self.last_mem_vector_index as i64);
348 self.last_mem_vector_index = vector_mem;
349 }
350 }
351}
352
353#[async_trait::async_trait]
354pub trait IndexerBuilder {
355 async fn build(
357 &self,
358 region_file_id: RegionFileId,
359 index_version: u64,
360 row_group_size: Option<usize>,
361 ) -> Indexer;
362}
363#[derive(Clone)]
364pub(crate) struct IndexerBuilderImpl {
365 pub(crate) build_type: IndexBuildType,
366 pub(crate) metadata: RegionMetadataRef,
367 pub(crate) puffin_manager: SstPuffinManager,
368 pub(crate) write_cache_enabled: bool,
369 pub(crate) intermediate_manager: IntermediateManager,
370 pub(crate) index_options: IndexOptions,
371 pub(crate) inverted_index_config: InvertedIndexConfig,
372 pub(crate) fulltext_index_config: FulltextIndexConfig,
373 pub(crate) bloom_filter_index_config: BloomFilterConfig,
374 #[cfg(feature = "vector_index")]
375 pub(crate) vector_index_config: VectorIndexConfig,
376}
377
378#[async_trait::async_trait]
379impl IndexerBuilder for IndexerBuilderImpl {
380 async fn build(
382 &self,
383 region_file_id: RegionFileId,
384 index_version: u64,
385 row_group_size: Option<usize>,
386 ) -> Indexer {
387 let mut indexer = Indexer {
388 file_id: region_file_id.file_id(),
389 region_id: self.metadata.region_id,
390 physical_region_id: region_file_id.region_id(),
391 index_version,
392 write_cache_enabled: self.write_cache_enabled,
393 ..Default::default()
394 };
395
396 indexer.inverted_indexer =
397 self.build_inverted_indexer(region_file_id.file_id(), row_group_size);
398 indexer.fulltext_indexer = self.build_fulltext_indexer(region_file_id.file_id()).await;
399 indexer.bloom_filter_indexer = self.build_bloom_filter_indexer(region_file_id.file_id());
400 #[cfg(feature = "vector_index")]
401 {
402 indexer.vector_indexer = self.build_vector_indexer(region_file_id.file_id());
403 }
404 indexer.intermediate_manager = Some(self.intermediate_manager.clone());
405
406 #[cfg(feature = "vector_index")]
407 let has_any_indexer = indexer.inverted_indexer.is_some()
408 || indexer.fulltext_indexer.is_some()
409 || indexer.bloom_filter_indexer.is_some()
410 || indexer.vector_indexer.is_some();
411 #[cfg(not(feature = "vector_index"))]
412 let has_any_indexer = indexer.inverted_indexer.is_some()
413 || indexer.fulltext_indexer.is_some()
414 || indexer.bloom_filter_indexer.is_some();
415
416 if !has_any_indexer {
417 indexer.abort().await;
418 return Indexer::default();
419 }
420
421 indexer.puffin_manager = Some(self.puffin_manager.clone());
422 indexer
423 }
424}
425
426impl IndexerBuilderImpl {
427 fn build_inverted_indexer(
428 &self,
429 file_id: FileId,
430 row_group_size: Option<usize>,
431 ) -> Option<InvertedIndexer> {
432 let create = match self.build_type {
433 IndexBuildType::Flush => self.inverted_index_config.create_on_flush.auto(),
434 IndexBuildType::Compact => self.inverted_index_config.create_on_compaction.auto(),
435 _ => true,
436 };
437
438 if !create {
439 debug!(
440 "Skip creating inverted index due to config, region_id: {}, file_id: {}",
441 self.metadata.region_id, file_id,
442 );
443 return None;
444 }
445
446 let indexed_column_ids = self.metadata.inverted_indexed_column_ids(
447 self.index_options.inverted_index.ignore_column_ids.iter(),
448 );
449 if indexed_column_ids.is_empty() {
450 debug!(
451 "No columns to be indexed, skip creating inverted index, region_id: {}, file_id: {}",
452 self.metadata.region_id, file_id,
453 );
454 return None;
455 }
456
457 let Some(mut segment_row_count) =
458 NonZeroUsize::new(self.index_options.inverted_index.segment_row_count)
459 else {
460 warn!(
461 "Segment row count is 0, skip creating index, region_id: {}, file_id: {}",
462 self.metadata.region_id, file_id,
463 );
464 return None;
465 };
466
467 if let Some(row_group_size) = row_group_size.and_then(NonZeroUsize::new)
469 && row_group_size.get() % segment_row_count.get() != 0
470 {
471 segment_row_count = row_group_size;
472 }
473
474 let indexer = InvertedIndexer::new(
475 file_id,
476 &self.metadata,
477 self.intermediate_manager.clone(),
478 self.inverted_index_config.mem_threshold_on_create(),
479 segment_row_count,
480 indexed_column_ids,
481 );
482
483 Some(indexer)
484 }
485
486 async fn build_fulltext_indexer(&self, file_id: FileId) -> Option<FulltextIndexer> {
487 let create = match self.build_type {
488 IndexBuildType::Flush => self.fulltext_index_config.create_on_flush.auto(),
489 IndexBuildType::Compact => self.fulltext_index_config.create_on_compaction.auto(),
490 _ => true,
491 };
492
493 if !create {
494 debug!(
495 "Skip creating full-text index due to config, region_id: {}, file_id: {}",
496 self.metadata.region_id, file_id,
497 );
498 return None;
499 }
500
501 let mem_limit = self.fulltext_index_config.mem_threshold_on_create();
502 let creator = FulltextIndexer::new(
503 &self.metadata.region_id,
504 &file_id,
505 &self.intermediate_manager,
506 &self.metadata,
507 self.fulltext_index_config.compress,
508 mem_limit,
509 )
510 .await;
511
512 let err = match creator {
513 Ok(creator) => {
514 if creator.is_none() {
515 debug!(
516 "Skip creating full-text index due to no columns require indexing, region_id: {}, file_id: {}",
517 self.metadata.region_id, file_id,
518 );
519 }
520 return creator;
521 }
522 Err(err) => err,
523 };
524
525 if cfg!(any(test, feature = "test")) {
526 panic!(
527 "Failed to create full-text indexer, region_id: {}, file_id: {}, err: {:?}",
528 self.metadata.region_id, file_id, err
529 );
530 } else {
531 warn!(
532 err; "Failed to create full-text indexer, region_id: {}, file_id: {}",
533 self.metadata.region_id, file_id,
534 );
535 }
536
537 None
538 }
539
540 fn build_bloom_filter_indexer(&self, file_id: FileId) -> Option<BloomFilterIndexer> {
541 let create = match self.build_type {
542 IndexBuildType::Flush => self.bloom_filter_index_config.create_on_flush.auto(),
543 IndexBuildType::Compact => self.bloom_filter_index_config.create_on_compaction.auto(),
544 _ => true,
545 };
546
547 if !create {
548 debug!(
549 "Skip creating bloom filter due to config, region_id: {}, file_id: {}",
550 self.metadata.region_id, file_id,
551 );
552 return None;
553 }
554
555 let mem_limit = self.bloom_filter_index_config.mem_threshold_on_create();
556 let indexer = BloomFilterIndexer::new(
557 file_id,
558 &self.metadata,
559 self.intermediate_manager.clone(),
560 mem_limit,
561 );
562
563 let err = match indexer {
564 Ok(indexer) => {
565 if indexer.is_none() {
566 debug!(
567 "Skip creating bloom filter due to no columns require indexing, region_id: {}, file_id: {}",
568 self.metadata.region_id, file_id,
569 );
570 }
571 return indexer;
572 }
573 Err(err) => err,
574 };
575
576 if cfg!(any(test, feature = "test")) {
577 panic!(
578 "Failed to create bloom filter, region_id: {}, file_id: {}, err: {:?}",
579 self.metadata.region_id, file_id, err
580 );
581 } else {
582 warn!(
583 err; "Failed to create bloom filter, region_id: {}, file_id: {}",
584 self.metadata.region_id, file_id,
585 );
586 }
587
588 None
589 }
590
591 #[cfg(feature = "vector_index")]
592 fn build_vector_indexer(&self, file_id: FileId) -> Option<VectorIndexer> {
593 let create = match self.build_type {
594 IndexBuildType::Flush => self.vector_index_config.create_on_flush.auto(),
595 IndexBuildType::Compact => self.vector_index_config.create_on_compaction.auto(),
596 _ => true,
597 };
598
599 if !create {
600 debug!(
601 "Skip creating vector index due to config, region_id: {}, file_id: {}",
602 self.metadata.region_id, file_id,
603 );
604 return None;
605 }
606
607 let vector_index_options = self.metadata.vector_indexed_column_ids();
609 if vector_index_options.is_empty() {
610 debug!(
611 "No vector columns to index, skip creating vector index, region_id: {}, file_id: {}",
612 self.metadata.region_id, file_id,
613 );
614 return None;
615 }
616
617 let mem_limit = self.vector_index_config.mem_threshold_on_create();
618 let indexer = VectorIndexer::new(
619 file_id,
620 &self.metadata,
621 self.intermediate_manager.clone(),
622 mem_limit,
623 &vector_index_options,
624 );
625
626 let err = match indexer {
627 Ok(indexer) => {
628 if indexer.is_none() {
629 debug!(
630 "Skip creating vector index due to no columns require indexing, region_id: {}, file_id: {}",
631 self.metadata.region_id, file_id,
632 );
633 }
634 return indexer;
635 }
636 Err(err) => err,
637 };
638
639 if cfg!(any(test, feature = "test")) {
640 panic!(
641 "Failed to create vector index, region_id: {}, file_id: {}, err: {:?}",
642 self.metadata.region_id, file_id, err
643 );
644 } else {
645 warn!(
646 err; "Failed to create vector index, region_id: {}, file_id: {}",
647 self.metadata.region_id, file_id,
648 );
649 }
650
651 None
652 }
653}
654
655#[derive(Debug, Clone, IntoStaticStr, PartialEq)]
657pub enum IndexBuildType {
658 SchemaChange,
660 Flush,
662 Compact,
664 Manual,
666}
667
668impl IndexBuildType {
669 fn as_str(&self) -> &'static str {
670 self.into()
671 }
672
673 fn priority(&self) -> u8 {
675 match self {
676 IndexBuildType::Manual => 3,
677 IndexBuildType::SchemaChange => 2,
678 IndexBuildType::Flush => 1,
679 IndexBuildType::Compact => 0,
680 }
681 }
682}
683
684impl From<OperationType> for IndexBuildType {
685 fn from(op_type: OperationType) -> Self {
686 match op_type {
687 OperationType::Flush => IndexBuildType::Flush,
688 OperationType::Compact => IndexBuildType::Compact,
689 }
690 }
691}
692
693#[derive(Debug, Clone, PartialEq, Eq, Hash)]
695pub enum IndexBuildOutcome {
696 Finished,
697 Aborted(String),
698}
699
700pub type ResultMpscSender = Sender<Result<IndexBuildOutcome>>;
702
703#[derive(Clone)]
704pub struct IndexBuildTask {
705 pub region_id: RegionId,
707 pub file: FileHandle,
709 pub(crate) target_region_metadata: RegionMetadataRef,
717 pub(crate) source: IndexBuildSource,
719 pub reason: IndexBuildType,
720 pub access_layer: AccessLayerRef,
721 pub(crate) listener: WorkerListener,
722 pub(crate) manifest_ctx: ManifestContextRef,
723 pub write_cache: Option<WriteCacheRef>,
724 pub cache_manager: Option<CacheManagerRef>,
725 pub file_purger: FilePurgerRef,
726 pub indexer_builder: Arc<dyn IndexerBuilder + Send + Sync>,
729 pub(crate) request_sender: Sender<WorkerRequestWithTime>,
731 pub(crate) result_sender: ResultMpscSender,
733}
734
735impl std::fmt::Debug for IndexBuildTask {
736 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
737 f.debug_struct("IndexBuildTask")
738 .field("region_id", &self.region_id)
739 .field("origin_region_id", &self.source.file_meta.region_id)
740 .field("file_id", &self.source.file_meta.file_id)
741 .field("schema_version", &self.source.schema_version)
742 .field("reason", &self.reason)
743 .finish()
744 }
745}
746
747impl IndexBuildTask {
748 pub async fn on_success(&self, outcome: IndexBuildOutcome) {
750 let _ = self.result_sender.send(Ok(outcome)).await;
751 }
752
753 pub async fn on_failure(&self, err: Arc<Error>) {
755 let _ = self
756 .result_sender
757 .send(Err(err.clone()).context(BuildIndexAsyncSnafu {
758 region_id: self.region_id,
759 }))
760 .await;
761 }
762
763 fn into_index_build_job(mut self, version_control: VersionControlRef) -> Job {
764 Box::pin(async move {
765 self.do_index_build(version_control).await;
766 })
767 }
768
769 async fn do_index_build(&mut self, version_control: VersionControlRef) {
770 self.listener
771 .on_index_build_begin(RegionFileId::new(
772 self.source.file_meta.region_id,
773 self.source.file_meta.file_id,
774 ))
775 .await;
776 match self.index_build(version_control).await {
777 Ok(outcome) => self.on_success(outcome).await,
778 Err(e) => {
779 warn!(
780 e; "Index build task failed, region: {}, file_id: {}",
781 self.region_id, self.source.file_meta.file_id,
782 );
783 self.on_failure(e.into()).await
784 }
785 }
786 let worker_request = WorkerRequest::Background {
787 region_id: self.region_id,
788 notify: BackgroundNotify::IndexBuildStopped(IndexBuildStopped {
789 file_id: self.source.file_meta.file_id,
790 }),
791 };
792 let _ = self
793 .request_sender
794 .send(WorkerRequestWithTime::new(worker_request))
795 .await;
796 }
797
798 async fn check_sst_file_exists(&self, version_control: &VersionControlRef) -> bool {
800 let file_id = self.source.file_meta.file_id;
801 let level = self.source.file_meta.level;
802 let version = version_control.current().version;
804
805 let Some(level_files) = version.ssts.levels().get(level as usize) else {
806 warn!(
807 "File id {} not found in level {} for index build, region: {}",
808 file_id, level, self.region_id
809 );
810 return false;
811 };
812
813 match level_files.files.get(&file_id) {
814 Some(handle) if !handle.is_deleted() && !handle.compacting() => {
815 true
819 }
820 _ => {
821 warn!(
822 "File id {} not found in region version for index build, region: {}",
823 file_id, self.region_id
824 );
825 false
826 }
827 }
828 }
829
830 async fn index_build(
831 &mut self,
832 version_control: VersionControlRef,
833 ) -> Result<IndexBuildOutcome> {
834 let new_index_version = self
835 .source
836 .file_meta
837 .index_version()
838 .map_or(0, |version| version + 1);
839
840 if !self.check_sst_file_exists(&version_control).await {
842 self.listener
843 .on_index_build_abort(RegionFileId::new(
844 self.source.file_meta.region_id,
845 self.source.file_meta.file_id,
846 ))
847 .await;
848 return Ok(IndexBuildOutcome::Aborted(format!(
849 "SST file not found during index build, region: {}, file_id: {}",
850 self.region_id, self.source.file_meta.file_id
851 )));
852 }
853
854 let mut parquet_reader = self
855 .access_layer
856 .read_sst(self.file.clone()) .expected_metadata(Some(self.target_region_metadata.clone()))
858 .build()
859 .await?;
860
861 let row_group_size = parquet_reader.as_ref().and_then(|reader| {
862 reader
863 .parquet_metadata()
864 .row_groups()
865 .first()
866 .map(|row_group| row_group.num_rows() as usize)
867 .filter(|size| *size > 0)
868 });
869
870 let region_file_id = RegionFileId::new(
872 self.source.file_meta.region_id,
873 self.source.file_meta.file_id,
874 );
875 let index_file_id = region_file_id.file_id();
876 let mut indexer = self
877 .indexer_builder
878 .build(region_file_id, new_index_version, row_group_size)
879 .await;
880
881 if let Some(mut parquet_reader) = parquet_reader.take() {
882 loop {
884 match parquet_reader.next_record_batch().await {
885 Ok(Some(batch)) => {
886 indexer.update_flat(&batch).await;
887 }
888 Ok(None) => break,
889 Err(e) => {
890 indexer.abort().await;
891 return Err(e);
892 }
893 }
894 }
895 }
896 let index_output = indexer.finish().await;
897
898 if index_output.file_size > 0 {
899 if !self.check_sst_file_exists(&version_control).await {
901 indexer.abort().await;
903 self.listener
904 .on_index_build_abort(RegionFileId::new(
905 self.source.file_meta.region_id,
906 self.source.file_meta.file_id,
907 ))
908 .await;
909 return Ok(IndexBuildOutcome::Aborted(format!(
910 "SST file not found during index build, region: {}, file_id: {}",
911 self.region_id, self.source.file_meta.file_id
912 )));
913 }
914
915 self.maybe_upload_index_file(index_output.clone(), index_file_id, new_index_version)
917 .await?;
918
919 self.listener
920 .on_index_build_before_manifest_commit(region_file_id)
921 .await;
922
923 let worker_request = match self.update_manifest(index_output, new_index_version).await {
924 Ok(IndexPublication::Committed {
925 manifest_version,
926 file_meta,
927 }) => {
928 self.listener
929 .on_index_build_manifest_committed(region_file_id)
930 .await;
931 let index_build_finished = IndexBuildFinished {
932 manifest_version,
933 file_meta,
934 };
935 WorkerRequest::Background {
936 region_id: self.region_id,
937 notify: BackgroundNotify::IndexBuildFinished(index_build_finished),
938 }
939 }
940 Ok(IndexPublication::Stale(stale)) => {
941 INDEX_PUBLICATION_STALE_TOTAL
942 .with_label_values(&["manifest_commit"])
943 .inc();
944 let index_id = RegionIndexId::new(region_file_id, new_index_version);
945 cleanup_stale_index_caches(
949 index_id,
950 &self.access_layer,
951 self.cache_manager.as_ref(),
952 self.write_cache.as_ref(),
953 )
954 .await;
955 self.listener.on_index_build_abort(region_file_id).await;
956 if stale == IndexPublicationStale::SchemaChanged {
957 let retry = BuildIndexRequest {
962 region_id: self.region_id,
963 build_type: IndexBuildType::SchemaChange,
964 file_metas: vec![self.source.file_meta.clone()],
965 };
966 let worker_request = WorkerRequest::Background {
967 region_id: self.region_id,
968 notify: BackgroundNotify::IndexBuildRetry(retry),
969 };
970 let _ = self
971 .request_sender
972 .send(WorkerRequestWithTime::new(worker_request))
973 .await;
974 }
975 return Ok(IndexBuildOutcome::Aborted(format!(
976 "Index build source changed before publication, region: {}, file_id: {}",
977 self.region_id, self.source.file_meta.file_id
978 )));
979 }
980 Err(e) => {
981 let err = Arc::new(e);
982 WorkerRequest::Background {
983 region_id: self.region_id,
984 notify: BackgroundNotify::IndexBuildFailed(IndexBuildFailed { err }),
985 }
986 }
987 };
988
989 let _ = self
990 .request_sender
991 .send(WorkerRequestWithTime::new(worker_request))
992 .await;
993 }
994 Ok(IndexBuildOutcome::Finished)
995 }
996
997 async fn maybe_upload_index_file(
998 &self,
999 output: IndexOutput,
1000 index_file_id: FileId,
1001 index_version: u64,
1002 ) -> Result<()> {
1003 if let Some(write_cache) = &self.write_cache {
1004 let file_id = self.source.file_meta.file_id;
1005 let region_id = self.source.file_meta.region_id;
1006 let remote_store = self.access_layer.object_store();
1007 let mut upload_tracker = UploadTracker::new(region_id);
1008 let mut err = None;
1009 let puffin_key =
1010 IndexKey::new(region_id, index_file_id, FileType::Puffin(output.version));
1011 let index_id = RegionIndexId::new(RegionFileId::new(region_id, file_id), index_version);
1012 let puffin_path = RegionFilePathFactory::new(
1013 self.access_layer.table_dir().to_string(),
1014 self.access_layer.path_type(),
1015 )
1016 .build_index_file_path_with_version(index_id);
1017 if let Err(e) = write_cache
1018 .upload(
1021 puffin_key,
1022 &puffin_path,
1023 remote_store,
1024 OperationType::Compact,
1025 )
1026 .await
1027 {
1028 err = Some(e);
1029 }
1030 upload_tracker.push_uploaded_file(puffin_path);
1031 if let Some(err) = err {
1032 upload_tracker
1034 .clean(
1035 &smallvec![SstInfo {
1036 file_id,
1037 index_metadata: output,
1038 ..Default::default()
1039 }],
1040 &write_cache.file_cache(),
1041 remote_store,
1042 )
1043 .await;
1044 return Err(err);
1045 }
1046 } else {
1047 debug!("write cache is not available, skip uploading index file");
1048 }
1049 Ok(())
1050 }
1051
1052 async fn update_manifest(
1053 &self,
1054 output: IndexOutput,
1055 new_index_version: u64,
1056 ) -> Result<IndexPublication> {
1057 let mut updated = self.source.file_meta.clone();
1058 updated.available_indexes = output.build_available_indexes();
1059 updated.indexes = output.build_indexes();
1060 updated.index_file_size = output.file_size;
1061 updated.index_version = new_index_version;
1062 let publication = self
1063 .manifest_ctx
1064 .update_manifest_for_index(&self.source, updated)
1065 .await?;
1066 if let IndexPublication::Committed {
1067 manifest_version, ..
1068 } = &publication
1069 {
1070 info!(
1071 "Successfully update manifest version to {}, region: {}, reason: {}",
1072 manifest_version,
1073 self.region_id,
1074 self.reason.as_str()
1075 );
1076 }
1077 Ok(publication)
1078 }
1079}
1080
1081impl PartialEq for IndexBuildTask {
1082 fn eq(&self, other: &Self) -> bool {
1083 self.reason.priority() == other.reason.priority()
1084 }
1085}
1086
1087impl Eq for IndexBuildTask {}
1088
1089impl PartialOrd for IndexBuildTask {
1090 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1091 Some(self.cmp(other))
1092 }
1093}
1094
1095impl Ord for IndexBuildTask {
1096 fn cmp(&self, other: &Self) -> Ordering {
1097 self.reason.priority().cmp(&other.reason.priority())
1098 }
1099}
1100
1101#[derive(Clone)]
1102struct PendingIndexBuild {
1103 task: IndexBuildTask,
1104 version_control: VersionControlRef,
1105}
1106
1107impl PartialEq for PendingIndexBuild {
1108 fn eq(&self, other: &Self) -> bool {
1109 self.task == other.task
1110 }
1111}
1112
1113impl Eq for PendingIndexBuild {}
1114
1115impl PartialOrd for PendingIndexBuild {
1116 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1117 Some(self.cmp(other))
1118 }
1119}
1120
1121impl Ord for PendingIndexBuild {
1122 fn cmp(&self, other: &Self) -> Ordering {
1123 self.task.cmp(&other.task)
1124 }
1125}
1126
1127impl PendingIndexBuild {
1128 fn supersedes(&self, other: &Self) -> bool {
1131 match self
1132 .task
1133 .source
1134 .schema_version
1135 .cmp(&other.task.source.schema_version)
1136 {
1137 Ordering::Greater => true,
1138 Ordering::Less => false,
1139 Ordering::Equal => match self
1140 .task
1141 .source
1142 .file_meta
1143 .index_version()
1144 .cmp(&other.task.source.file_meta.index_version())
1145 {
1146 Ordering::Greater => true,
1147 Ordering::Less => false,
1148 Ordering::Equal => self.task.reason.priority() > other.task.reason.priority(),
1149 },
1150 }
1151 }
1152}
1153
1154#[derive(Default)]
1155struct PendingIndexBuilds {
1156 tasks: HashMap<FileId, PendingIndexBuild>,
1157}
1158
1159impl PendingIndexBuilds {
1160 fn insert(&mut self, pending: PendingIndexBuild) -> Option<IndexBuildTask> {
1161 let file_id = pending.task.source.file_meta.file_id;
1162 match self.tasks.entry(file_id) {
1163 std::collections::hash_map::Entry::Vacant(entry) => {
1164 entry.insert(pending);
1165 None
1166 }
1167 std::collections::hash_map::Entry::Occupied(mut entry)
1168 if pending.supersedes(entry.get()) =>
1169 {
1170 Some(entry.insert(pending).task)
1171 }
1172 std::collections::hash_map::Entry::Occupied(_) => Some(pending.task),
1173 }
1174 }
1175
1176 fn highest_ready(&self, building_files: &HashSet<FileId>) -> Option<PendingIndexBuild> {
1177 self.tasks
1178 .values()
1179 .filter(|pending| !building_files.contains(&pending.task.source.file_meta.file_id))
1180 .max()
1181 .cloned()
1182 }
1183
1184 fn remove(&mut self, file_id: FileId) {
1185 self.tasks.remove(&file_id);
1186 }
1187
1188 fn drain(&mut self) -> impl Iterator<Item = PendingIndexBuild> {
1189 std::mem::take(&mut self.tasks).into_values()
1190 }
1191
1192 #[cfg(test)]
1193 fn len(&self) -> usize {
1194 self.tasks.len()
1195 }
1196
1197 fn is_empty(&self) -> bool {
1198 self.tasks.is_empty()
1199 }
1200}
1201
1202struct IndexBuildStatus {
1204 building_files: HashSet<FileId>,
1205 pending_tasks: PendingIndexBuilds,
1207 retiring: bool,
1210}
1211
1212impl IndexBuildStatus {
1213 fn new() -> Self {
1214 IndexBuildStatus {
1215 building_files: HashSet::new(),
1216 pending_tasks: PendingIndexBuilds::default(),
1217 retiring: false,
1218 }
1219 }
1220
1221 async fn fail_pending(&mut self, err: Arc<Error>) {
1222 for pending in self.pending_tasks.drain() {
1223 pending.task.on_failure(err.clone()).await;
1224 }
1225 }
1226
1227 async fn retire(&mut self, err: Arc<Error>) {
1228 self.retiring = !self.building_files.is_empty();
1229 self.fail_pending(err).await;
1230 }
1231}
1232
1233pub struct IndexBuildScheduler {
1234 scheduler: SchedulerRef,
1236 region_status: HashMap<RegionId, IndexBuildStatus>,
1238 files_limit: usize,
1240}
1241
1242impl IndexBuildScheduler {
1244 pub fn new(scheduler: SchedulerRef, files_limit: usize) -> Self {
1245 IndexBuildScheduler {
1246 scheduler,
1247 region_status: HashMap::new(),
1248 files_limit,
1249 }
1250 }
1251
1252 pub(crate) async fn schedule_build(
1253 &mut self,
1254 version_control: &VersionControlRef,
1255 task: IndexBuildTask,
1256 ) -> Result<()> {
1257 let status = self
1258 .region_status
1259 .entry(task.region_id)
1260 .or_insert_with(IndexBuildStatus::new);
1261
1262 let file_id = task.source.file_meta.file_id;
1263 let region_file_id = RegionFileId::new(
1264 task.source.file_meta.region_id,
1265 task.source.file_meta.file_id,
1266 );
1267 let can_wait_for_active = status.retiring || task.reason == IndexBuildType::SchemaChange;
1268 let (rejected, coalesced) =
1269 if status.building_files.contains(&file_id) && !can_wait_for_active {
1270 (Some(task), None)
1271 } else {
1272 let pending = PendingIndexBuild {
1273 task,
1274 version_control: version_control.clone(),
1275 };
1276 (None, status.pending_tasks.insert(pending))
1277 };
1278 let should_schedule = !status.retiring;
1279
1280 if let Some(rejected) = rejected {
1281 debug!(
1282 "Rejecting index build because region file {:?} is already being built",
1283 region_file_id
1284 );
1285 rejected
1286 .on_success(IndexBuildOutcome::Aborted(format!(
1287 "Index is already being built for region file {:?}",
1288 region_file_id
1289 )))
1290 .await;
1291 rejected.listener.on_index_build_abort(region_file_id).await;
1292 }
1293
1294 if let Some(coalesced) = coalesced {
1295 debug!(
1296 "Coalescing redundant index build task for region file {:?}",
1297 region_file_id
1298 );
1299 coalesced
1300 .on_success(IndexBuildOutcome::Aborted(format!(
1301 "Index build was coalesced for region file {:?}",
1302 region_file_id
1303 )))
1304 .await;
1305 coalesced
1306 .listener
1307 .on_index_build_abort(region_file_id)
1308 .await;
1309 }
1310
1311 if should_schedule {
1312 self.schedule_next_build_batch();
1313 }
1314 Ok(())
1315 }
1316
1317 fn schedule_next_build_batch(&mut self) {
1319 let mut building_count = self
1320 .region_status
1321 .values()
1322 .map(|status| status.building_files.len())
1323 .sum::<usize>();
1324
1325 while building_count < self.files_limit {
1326 let Some(pending) = self.find_next_task() else {
1327 break;
1328 };
1329
1330 let task = pending.task;
1331 let region_id = task.region_id;
1332 let file_id = task.source.file_meta.file_id;
1333 let job = task.clone().into_index_build_job(pending.version_control);
1334 match self.scheduler.schedule(job) {
1335 Ok(()) => {
1336 if let Some(status) = self.region_status.get_mut(®ion_id) {
1337 status.pending_tasks.remove(file_id);
1338 status.building_files.insert(file_id);
1339 building_count += 1;
1340 } else {
1341 error!(
1342 "Region status not found when scheduling index build task, region: {}",
1343 region_id
1344 );
1345 }
1346 }
1347 Err(err) => {
1348 error!(
1349 err;
1350 "Failed to schedule index build job, region: {}, file_id: {}",
1351 region_id,
1352 file_id
1353 );
1354 if let Some(status) = self.region_status.get_mut(®ion_id) {
1355 status.pending_tasks.remove(file_id);
1356 }
1357 common_runtime::spawn_global(async move {
1358 task.on_failure(Arc::new(err)).await;
1359 });
1360 }
1361 }
1362 }
1363
1364 self.region_status.retain(|_, status| {
1365 !status.building_files.is_empty() || !status.pending_tasks.is_empty()
1366 });
1367 }
1368
1369 fn find_next_task(&self) -> Option<PendingIndexBuild> {
1371 self.region_status
1372 .values()
1373 .filter(|status| !status.retiring)
1374 .filter_map(|status| status.pending_tasks.highest_ready(&status.building_files))
1375 .max()
1376 }
1377
1378 pub(crate) fn on_task_stopped(&mut self, region_id: RegionId, file_id: FileId) {
1379 if let Some(status) = self.region_status.get_mut(®ion_id) {
1380 if !status.building_files.remove(&file_id) {
1381 debug!(
1382 "Index build task is not tracked as building, region: {}, file: {}",
1383 region_id, file_id
1384 );
1385 return;
1386 }
1387 if status.building_files.is_empty() && status.pending_tasks.is_empty() {
1388 self.region_status.remove(®ion_id);
1390 } else if status.building_files.is_empty() {
1391 status.retiring = false;
1393 }
1394 }
1395
1396 self.schedule_next_build_batch();
1397 }
1398
1399 pub(crate) async fn on_failure(&mut self, region_id: RegionId, err: Arc<Error>) {
1400 let Some(status) = self.region_status.get_mut(®ion_id) else {
1401 error!(err; "Index build failed after scheduler state was removed, region: {}", region_id);
1402 return;
1403 };
1404 if status.retiring {
1405 error!(err; "Index build from a retiring region incarnation failed, region: {}", region_id);
1406 return;
1407 }
1408 error!(
1409 err; "Index build scheduler encountered failure for region {}, failing pending tasks.",
1410 region_id
1411 );
1412 status.fail_pending(err).await;
1413 if status.building_files.is_empty() {
1414 self.region_status.remove(®ion_id);
1415 }
1416 }
1417
1418 pub(crate) async fn on_region_dropped(&mut self, region_id: RegionId) {
1420 self.retire_region(
1421 region_id,
1422 Arc::new(RegionDroppedSnafu { region_id }.build()),
1423 )
1424 .await;
1425 }
1426
1427 pub(crate) async fn on_region_closed(&mut self, region_id: RegionId) {
1429 self.retire_region(region_id, Arc::new(RegionClosedSnafu { region_id }.build()))
1430 .await;
1431 }
1432
1433 pub(crate) async fn on_region_truncated(&mut self, region_id: RegionId) {
1435 self.retire_region(
1436 region_id,
1437 Arc::new(RegionTruncatedSnafu { region_id }.build()),
1438 )
1439 .await;
1440 }
1441
1442 async fn retire_region(&mut self, region_id: RegionId, err: Arc<Error>) {
1443 let Some(status) = self.region_status.get_mut(®ion_id) else {
1444 return;
1445 };
1446 status.retire(err).await;
1447 if status.building_files.is_empty() {
1448 self.region_status.remove(®ion_id);
1449 }
1450 }
1451}
1452
1453pub(crate) fn decode_primary_keys_with_counts(
1456 batch: &RecordBatch,
1457 codec: &IndexValuesCodec,
1458) -> Result<Vec<(CompositeValues, usize)>> {
1459 let primary_key_index = primary_key_column_index(batch.num_columns());
1460 let pk_dict_array = batch
1461 .column(primary_key_index)
1462 .as_any()
1463 .downcast_ref::<PrimaryKeyArray>()
1464 .context(InvalidRecordBatchSnafu {
1465 reason: "Primary key column is not a dictionary array",
1466 })?;
1467 let pk_values_array = pk_dict_array
1468 .values()
1469 .as_any()
1470 .downcast_ref::<BinaryArray>()
1471 .context(InvalidRecordBatchSnafu {
1472 reason: "Primary key values are not binary array",
1473 })?;
1474 let keys = pk_dict_array.keys();
1475
1476 let mut result: Vec<(CompositeValues, usize)> = Vec::new();
1478 let mut prev_key: Option<u32> = None;
1479
1480 let pk_indices = keys.values();
1481 for ¤t_key in pk_indices.iter().take(keys.len()) {
1482 if let Some(prev) = prev_key
1484 && prev == current_key
1485 {
1486 result.last_mut().unwrap().1 += 1;
1488 continue;
1489 }
1490
1491 let pk_bytes = pk_values_array.value(current_key as usize);
1493 let decoded_value = codec.decoder().decode(pk_bytes).context(DecodeSnafu)?;
1494
1495 result.push((decoded_value, 1));
1496 prev_key = Some(current_key);
1497 }
1498
1499 Ok(result)
1500}
1501
1502#[cfg(test)]
1503mod tests {
1504 use std::sync::Arc;
1505
1506 use api::v1::SemanticType;
1507 use common_base::readable_size::ReadableSize;
1508 use datafusion_common::HashMap;
1509 use datatypes::data_type::ConcreteDataType;
1510 use datatypes::schema::{
1511 ColumnSchema, FulltextOptions, SkippingIndexOptions, SkippingIndexType,
1512 };
1513 use datatypes::value::Value;
1514 use index::inverted_index::format::reader::InvertedIndexReader;
1515 use object_store::ObjectStore;
1516 use object_store::services::Memory;
1517 use partition::expr::col;
1518 use puffin::puffin_manager::{PuffinManager, PuffinReader};
1519 use puffin_manager::PuffinManagerFactory;
1520 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
1521 use tokio::sync::mpsc;
1522
1523 use super::*;
1524 use crate::access_layer::{FilePathProvider, Metrics, SstWriteRequest, WriteType};
1525 use crate::cache::write_cache::WriteCache;
1526 use crate::config::{FulltextIndexConfig, IndexBuildMode, MitoConfig, Mode};
1527 use crate::manifest::action::{RegionEdit, RegionMetaAction, RegionMetaActionList};
1528 use crate::memtable::time_partition::TimePartitions;
1529 use crate::region::RegionLeaderState;
1530 use crate::region::version::{VersionBuilder, VersionControl};
1531 use crate::schedule::scheduler::{LocalScheduler, Scheduler};
1532 use crate::sst::file::{FileMeta, RegionFileId};
1533 use crate::sst::file_purger::NoopFilePurger;
1534 use crate::sst::location;
1535 use crate::sst::parquet::WriteOptions;
1536 use crate::test_util::memtable_util::EmptyMemtableBuilder;
1537 use crate::test_util::scheduler_util::{SchedulerEnv, VecScheduler};
1538 use crate::test_util::sst_util::{
1539 new_flat_source_from_record_batches, new_record_batch_by_range, sst_region_metadata,
1540 };
1541
1542 struct MetaConfig {
1543 with_inverted: bool,
1544 with_fulltext: bool,
1545 with_skipping_bloom: bool,
1546 #[cfg(feature = "vector_index")]
1547 with_vector: bool,
1548 }
1549
1550 async fn seed_manifest_file(manifest_ctx: &ManifestContextRef, file_meta: &FileMeta) {
1551 manifest_ctx
1552 .update_manifest(
1553 RegionLeaderState::Writable,
1554 RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
1555 files_to_add: vec![file_meta.clone()],
1556 files_to_remove: Vec::new(),
1557 timestamp_ms: None,
1558 flushed_sequence: None,
1559 flushed_entry_id: None,
1560 committed_sequence: None,
1561 compaction_time_window: None,
1562 })),
1563 false,
1564 )
1565 .await
1566 .unwrap();
1567 }
1568
1569 fn mock_region_metadata(
1570 MetaConfig {
1571 with_inverted,
1572 with_fulltext,
1573 with_skipping_bloom,
1574 #[cfg(feature = "vector_index")]
1575 with_vector,
1576 }: MetaConfig,
1577 ) -> RegionMetadataRef {
1578 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 2));
1579 let mut column_schema = ColumnSchema::new("a", ConcreteDataType::int64_datatype(), false);
1580 if with_inverted {
1581 column_schema = column_schema.with_inverted_index(true);
1582 }
1583 builder
1584 .push_column_metadata(ColumnMetadata {
1585 column_schema,
1586 semantic_type: SemanticType::Field,
1587 column_id: 1,
1588 })
1589 .push_column_metadata(ColumnMetadata {
1590 column_schema: ColumnSchema::new("b", ConcreteDataType::float64_datatype(), false),
1591 semantic_type: SemanticType::Field,
1592 column_id: 2,
1593 })
1594 .push_column_metadata(ColumnMetadata {
1595 column_schema: ColumnSchema::new(
1596 "c",
1597 ConcreteDataType::timestamp_millisecond_datatype(),
1598 false,
1599 ),
1600 semantic_type: SemanticType::Timestamp,
1601 column_id: 3,
1602 });
1603
1604 if with_fulltext {
1605 let column_schema =
1606 ColumnSchema::new("text", ConcreteDataType::string_datatype(), true)
1607 .with_fulltext_options(FulltextOptions {
1608 enable: true,
1609 ..Default::default()
1610 })
1611 .unwrap();
1612
1613 let column = ColumnMetadata {
1614 column_schema,
1615 semantic_type: SemanticType::Field,
1616 column_id: 4,
1617 };
1618
1619 builder.push_column_metadata(column);
1620 }
1621
1622 if with_skipping_bloom {
1623 let column_schema =
1624 ColumnSchema::new("bloom", ConcreteDataType::string_datatype(), false)
1625 .with_skipping_options(SkippingIndexOptions::new_unchecked(
1626 42,
1627 0.01,
1628 SkippingIndexType::BloomFilter,
1629 ))
1630 .unwrap();
1631
1632 let column = ColumnMetadata {
1633 column_schema,
1634 semantic_type: SemanticType::Field,
1635 column_id: 5,
1636 };
1637
1638 builder.push_column_metadata(column);
1639 }
1640
1641 #[cfg(feature = "vector_index")]
1642 if with_vector {
1643 use index::vector::VectorIndexOptions;
1644
1645 let options = VectorIndexOptions::default();
1646 let column_schema =
1647 ColumnSchema::new("vec", ConcreteDataType::vector_datatype(4), true)
1648 .with_vector_index_options(&options)
1649 .unwrap();
1650 let column = ColumnMetadata {
1651 column_schema,
1652 semantic_type: SemanticType::Field,
1653 column_id: 6,
1654 };
1655
1656 builder.push_column_metadata(column);
1657 }
1658
1659 Arc::new(builder.build().unwrap())
1660 }
1661
1662 fn mock_object_store() -> ObjectStore {
1663 ObjectStore::new(Memory::default()).unwrap()
1664 }
1665
1666 async fn mock_intm_mgr(path: impl AsRef<str>) -> IntermediateManager {
1667 IntermediateManager::init_fs(path).await.unwrap()
1668 }
1669
1670 fn random_region_file_id() -> RegionFileId {
1671 RegionFileId::new(RegionId::new(1, 2), FileId::random())
1672 }
1673
1674 struct NoopPathProvider;
1675
1676 impl FilePathProvider for NoopPathProvider {
1677 fn build_index_file_path(&self, _file_id: RegionFileId) -> String {
1678 unreachable!()
1679 }
1680
1681 fn build_index_file_path_with_version(&self, _index_id: RegionIndexId) -> String {
1682 unreachable!()
1683 }
1684
1685 fn build_sst_file_path(&self, _file_id: RegionFileId) -> String {
1686 unreachable!()
1687 }
1688 }
1689
1690 async fn mock_sst_file(
1691 metadata: RegionMetadataRef,
1692 env: &SchedulerEnv,
1693 build_mode: IndexBuildMode,
1694 ) -> SstInfo {
1695 let source = new_flat_source_from_record_batches(vec![
1696 new_record_batch_by_range(&["a", "d"], 0, 60),
1697 new_record_batch_by_range(&["b", "f"], 0, 40),
1698 new_record_batch_by_range(&["b", "h"], 100, 200),
1699 ]);
1700 let mut index_config = MitoConfig::default().index;
1701 index_config.build_mode = build_mode;
1702 let write_request = SstWriteRequest {
1703 op_type: OperationType::Flush,
1704 metadata: metadata.clone(),
1705 source,
1706 storage: None,
1707 max_sequence: None,
1708 sst_write_format: Default::default(),
1709 cache_manager: Default::default(),
1710 preserve_row_sequence: false,
1711 index_options: IndexOptions::default(),
1712 index_config,
1713 inverted_index_config: Default::default(),
1714 fulltext_index_config: Default::default(),
1715 bloom_filter_index_config: Default::default(),
1716 #[cfg(feature = "vector_index")]
1717 vector_index_config: Default::default(),
1718 };
1719 let mut metrics = Metrics::new(WriteType::Flush);
1720 env.access_layer
1721 .write_sst(write_request, &WriteOptions::default(), &mut metrics)
1722 .await
1723 .unwrap()
1724 .remove(0)
1725 }
1726
1727 async fn mock_version_control(
1728 metadata: RegionMetadataRef,
1729 file_purger: FilePurgerRef,
1730 files: HashMap<FileId, FileMeta>,
1731 ) -> VersionControlRef {
1732 let mutable = Arc::new(TimePartitions::new(
1733 metadata.clone(),
1734 Arc::new(EmptyMemtableBuilder::default()),
1735 0,
1736 None,
1737 ));
1738 let version_builder = VersionBuilder::new(metadata, mutable)
1739 .add_files(file_purger, files.values().cloned())
1740 .build();
1741 Arc::new(VersionControl::new(version_builder))
1742 }
1743
1744 async fn mock_indexer_builder(
1745 metadata: RegionMetadataRef,
1746 env: &SchedulerEnv,
1747 ) -> Arc<dyn IndexerBuilder + Send + Sync> {
1748 let (dir, factory) = PuffinManagerFactory::new_for_test_async("mock_indexer_builder").await;
1749 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1750 let puffin_manager = factory.build(
1751 env.access_layer.object_store().clone(),
1752 RegionFilePathFactory::new(
1753 env.access_layer.table_dir().to_string(),
1754 env.access_layer.path_type(),
1755 ),
1756 );
1757 Arc::new(IndexerBuilderImpl {
1758 build_type: IndexBuildType::Flush,
1759 metadata,
1760 puffin_manager,
1761 write_cache_enabled: false,
1762 intermediate_manager: intm_manager,
1763 index_options: IndexOptions::default(),
1764 inverted_index_config: InvertedIndexConfig::default(),
1765 fulltext_index_config: FulltextIndexConfig::default(),
1766 bloom_filter_index_config: BloomFilterConfig::default(),
1767 #[cfg(feature = "vector_index")]
1768 vector_index_config: Default::default(),
1769 })
1770 }
1771
1772 #[tokio::test]
1773 async fn test_build_indexer_basic() {
1774 let (dir, factory) =
1775 PuffinManagerFactory::new_for_test_async("test_build_indexer_basic_").await;
1776 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1777
1778 let metadata = mock_region_metadata(MetaConfig {
1779 with_inverted: true,
1780 with_fulltext: true,
1781 with_skipping_bloom: true,
1782 #[cfg(feature = "vector_index")]
1783 with_vector: false,
1784 });
1785 let indexer = IndexerBuilderImpl {
1786 build_type: IndexBuildType::Flush,
1787 metadata,
1788 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1789 write_cache_enabled: false,
1790 intermediate_manager: intm_manager,
1791 index_options: IndexOptions::default(),
1792 inverted_index_config: InvertedIndexConfig::default(),
1793 fulltext_index_config: FulltextIndexConfig::default(),
1794 bloom_filter_index_config: BloomFilterConfig::default(),
1795 #[cfg(feature = "vector_index")]
1796 vector_index_config: Default::default(),
1797 }
1798 .build(random_region_file_id(), 0, Some(1024))
1799 .await;
1800
1801 assert!(indexer.inverted_indexer.is_some());
1802 assert!(indexer.fulltext_indexer.is_some());
1803 assert!(indexer.bloom_filter_indexer.is_some());
1804 }
1805
1806 #[tokio::test]
1807 async fn test_build_indexer_disable_create() {
1808 let (dir, factory) =
1809 PuffinManagerFactory::new_for_test_async("test_build_indexer_disable_create_").await;
1810 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1811
1812 let metadata = mock_region_metadata(MetaConfig {
1813 with_inverted: true,
1814 with_fulltext: true,
1815 with_skipping_bloom: true,
1816 #[cfg(feature = "vector_index")]
1817 with_vector: false,
1818 });
1819 let indexer = IndexerBuilderImpl {
1820 build_type: IndexBuildType::Flush,
1821 metadata: metadata.clone(),
1822 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1823 write_cache_enabled: false,
1824 intermediate_manager: intm_manager.clone(),
1825 index_options: IndexOptions::default(),
1826 inverted_index_config: InvertedIndexConfig {
1827 create_on_flush: Mode::Disable,
1828 ..Default::default()
1829 },
1830 fulltext_index_config: FulltextIndexConfig::default(),
1831 bloom_filter_index_config: BloomFilterConfig::default(),
1832 #[cfg(feature = "vector_index")]
1833 vector_index_config: Default::default(),
1834 }
1835 .build(random_region_file_id(), 0, Some(1024))
1836 .await;
1837
1838 assert!(indexer.inverted_indexer.is_none());
1839 assert!(indexer.fulltext_indexer.is_some());
1840 assert!(indexer.bloom_filter_indexer.is_some());
1841
1842 let indexer = IndexerBuilderImpl {
1843 build_type: IndexBuildType::Compact,
1844 metadata: metadata.clone(),
1845 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1846 write_cache_enabled: false,
1847 intermediate_manager: intm_manager.clone(),
1848 index_options: IndexOptions::default(),
1849 inverted_index_config: InvertedIndexConfig::default(),
1850 fulltext_index_config: FulltextIndexConfig {
1851 create_on_compaction: Mode::Disable,
1852 ..Default::default()
1853 },
1854 bloom_filter_index_config: BloomFilterConfig::default(),
1855 #[cfg(feature = "vector_index")]
1856 vector_index_config: Default::default(),
1857 }
1858 .build(random_region_file_id(), 0, Some(1024))
1859 .await;
1860
1861 assert!(indexer.inverted_indexer.is_some());
1862 assert!(indexer.fulltext_indexer.is_none());
1863 assert!(indexer.bloom_filter_indexer.is_some());
1864
1865 let indexer = IndexerBuilderImpl {
1866 build_type: IndexBuildType::Compact,
1867 metadata,
1868 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1869 write_cache_enabled: false,
1870 intermediate_manager: intm_manager,
1871 index_options: IndexOptions::default(),
1872 inverted_index_config: InvertedIndexConfig::default(),
1873 fulltext_index_config: FulltextIndexConfig::default(),
1874 bloom_filter_index_config: BloomFilterConfig {
1875 create_on_compaction: Mode::Disable,
1876 ..Default::default()
1877 },
1878 #[cfg(feature = "vector_index")]
1879 vector_index_config: Default::default(),
1880 }
1881 .build(random_region_file_id(), 0, Some(1024))
1882 .await;
1883
1884 assert!(indexer.inverted_indexer.is_some());
1885 assert!(indexer.fulltext_indexer.is_some());
1886 assert!(indexer.bloom_filter_indexer.is_none());
1887 }
1888
1889 #[tokio::test]
1890 async fn test_build_indexer_no_required() {
1891 let (dir, factory) =
1892 PuffinManagerFactory::new_for_test_async("test_build_indexer_no_required_").await;
1893 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1894
1895 let metadata = mock_region_metadata(MetaConfig {
1896 with_inverted: false,
1897 with_fulltext: true,
1898 with_skipping_bloom: true,
1899 #[cfg(feature = "vector_index")]
1900 with_vector: false,
1901 });
1902 let indexer = IndexerBuilderImpl {
1903 build_type: IndexBuildType::Flush,
1904 metadata: metadata.clone(),
1905 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1906 write_cache_enabled: false,
1907 intermediate_manager: intm_manager.clone(),
1908 index_options: IndexOptions::default(),
1909 inverted_index_config: InvertedIndexConfig::default(),
1910 fulltext_index_config: FulltextIndexConfig::default(),
1911 bloom_filter_index_config: BloomFilterConfig::default(),
1912 #[cfg(feature = "vector_index")]
1913 vector_index_config: Default::default(),
1914 }
1915 .build(random_region_file_id(), 0, Some(1024))
1916 .await;
1917
1918 assert!(indexer.inverted_indexer.is_none());
1919 assert!(indexer.fulltext_indexer.is_some());
1920 assert!(indexer.bloom_filter_indexer.is_some());
1921
1922 let metadata = mock_region_metadata(MetaConfig {
1923 with_inverted: true,
1924 with_fulltext: false,
1925 with_skipping_bloom: true,
1926 #[cfg(feature = "vector_index")]
1927 with_vector: false,
1928 });
1929 let indexer = IndexerBuilderImpl {
1930 build_type: IndexBuildType::Flush,
1931 metadata: metadata.clone(),
1932 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1933 write_cache_enabled: false,
1934 intermediate_manager: intm_manager.clone(),
1935 index_options: IndexOptions::default(),
1936 inverted_index_config: InvertedIndexConfig::default(),
1937 fulltext_index_config: FulltextIndexConfig::default(),
1938 bloom_filter_index_config: BloomFilterConfig::default(),
1939 #[cfg(feature = "vector_index")]
1940 vector_index_config: Default::default(),
1941 }
1942 .build(random_region_file_id(), 0, Some(1024))
1943 .await;
1944
1945 assert!(indexer.inverted_indexer.is_some());
1946 assert!(indexer.fulltext_indexer.is_none());
1947 assert!(indexer.bloom_filter_indexer.is_some());
1948
1949 let metadata = mock_region_metadata(MetaConfig {
1950 with_inverted: true,
1951 with_fulltext: true,
1952 with_skipping_bloom: false,
1953 #[cfg(feature = "vector_index")]
1954 with_vector: false,
1955 });
1956 let indexer = IndexerBuilderImpl {
1957 build_type: IndexBuildType::Flush,
1958 metadata: metadata.clone(),
1959 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1960 write_cache_enabled: false,
1961 intermediate_manager: intm_manager,
1962 index_options: IndexOptions::default(),
1963 inverted_index_config: InvertedIndexConfig::default(),
1964 fulltext_index_config: FulltextIndexConfig::default(),
1965 bloom_filter_index_config: BloomFilterConfig::default(),
1966 #[cfg(feature = "vector_index")]
1967 vector_index_config: Default::default(),
1968 }
1969 .build(random_region_file_id(), 0, Some(1024))
1970 .await;
1971
1972 assert!(indexer.inverted_indexer.is_some());
1973 assert!(indexer.fulltext_indexer.is_some());
1974 assert!(indexer.bloom_filter_indexer.is_none());
1975 }
1976
1977 #[tokio::test]
1978 async fn test_build_indexer_zero_row_group_hint() {
1979 let (dir, factory) =
1980 PuffinManagerFactory::new_for_test_async("test_build_indexer_zero_row_group_").await;
1981 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1982
1983 let metadata = mock_region_metadata(MetaConfig {
1984 with_inverted: true,
1985 with_fulltext: true,
1986 with_skipping_bloom: true,
1987 #[cfg(feature = "vector_index")]
1988 with_vector: false,
1989 });
1990 let indexer = IndexerBuilderImpl {
1991 build_type: IndexBuildType::Flush,
1992 metadata,
1993 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1994 write_cache_enabled: false,
1995 intermediate_manager: intm_manager,
1996 index_options: IndexOptions::default(),
1997 inverted_index_config: InvertedIndexConfig::default(),
1998 fulltext_index_config: FulltextIndexConfig::default(),
1999 bloom_filter_index_config: BloomFilterConfig::default(),
2000 #[cfg(feature = "vector_index")]
2001 vector_index_config: Default::default(),
2002 }
2003 .build(random_region_file_id(), 0, Some(0))
2004 .await;
2005
2006 assert!(indexer.inverted_indexer.is_some());
2007 }
2008
2009 #[cfg(feature = "vector_index")]
2010 #[tokio::test]
2011 async fn test_update_flat_builds_vector_index() {
2012 use datatypes::arrow::array::BinaryBuilder;
2013 use datatypes::arrow::datatypes::{DataType, Field, Schema};
2014
2015 struct TestPathProvider;
2016
2017 impl FilePathProvider for TestPathProvider {
2018 fn build_index_file_path(&self, file_id: RegionFileId) -> String {
2019 format!("index/{}.puffin", file_id)
2020 }
2021
2022 fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String {
2023 format!("index/{}.puffin", index_id)
2024 }
2025
2026 fn build_sst_file_path(&self, file_id: RegionFileId) -> String {
2027 format!("sst/{}.parquet", file_id)
2028 }
2029 }
2030
2031 fn f32s_to_bytes(values: &[f32]) -> Vec<u8> {
2032 let mut bytes = Vec::with_capacity(values.len() * 4);
2033 for v in values {
2034 bytes.extend_from_slice(&v.to_le_bytes());
2035 }
2036 bytes
2037 }
2038
2039 let (dir, factory) =
2040 PuffinManagerFactory::new_for_test_async("test_update_flat_builds_vector_index_").await;
2041 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
2042
2043 let metadata = mock_region_metadata(MetaConfig {
2044 with_inverted: false,
2045 with_fulltext: false,
2046 with_skipping_bloom: false,
2047 with_vector: true,
2048 });
2049
2050 let mut indexer = IndexerBuilderImpl {
2051 build_type: IndexBuildType::Flush,
2052 metadata,
2053 puffin_manager: factory.build(mock_object_store(), TestPathProvider),
2054 write_cache_enabled: false,
2055 intermediate_manager: intm_manager,
2056 index_options: IndexOptions::default(),
2057 inverted_index_config: InvertedIndexConfig::default(),
2058 fulltext_index_config: FulltextIndexConfig::default(),
2059 bloom_filter_index_config: BloomFilterConfig::default(),
2060 vector_index_config: Default::default(),
2061 }
2062 .build(random_region_file_id(), 0, Some(1024))
2063 .await;
2064
2065 assert!(indexer.vector_indexer.is_some());
2066
2067 let vec1 = f32s_to_bytes(&[1.0, 0.0, 0.0, 0.0]);
2068 let vec2 = f32s_to_bytes(&[0.0, 1.0, 0.0, 0.0]);
2069
2070 let mut builder = BinaryBuilder::with_capacity(2, vec1.len() + vec2.len());
2071 builder.append_value(&vec1);
2072 builder.append_value(&vec2);
2073
2074 let schema = Arc::new(Schema::new(vec![Field::new("vec", DataType::Binary, true)]));
2075 let batch = RecordBatch::try_new(schema, vec![Arc::new(builder.finish())]).unwrap();
2076
2077 indexer.update_flat(&batch).await;
2078 let output = indexer.finish().await;
2079
2080 assert!(output.vector_index.is_available());
2081 assert!(output.vector_index.columns.contains(&6));
2082 }
2083
2084 #[tokio::test]
2085 async fn test_index_build_task_sst_not_exist() {
2086 let env = SchedulerEnv::new().await;
2087 let (tx, _rx) = mpsc::channel(4);
2088 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2089 let mut scheduler = env.mock_index_build_scheduler(4);
2090 let metadata = Arc::new(sst_region_metadata());
2091 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2092 let file_purger = Arc::new(NoopFilePurger {});
2093 let files = HashMap::new();
2094 let version_control =
2095 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2096 let region_id = metadata.region_id;
2097 let indexer_builder = mock_indexer_builder(metadata, &env).await;
2098
2099 let file_meta = FileMeta {
2100 region_id,
2101 file_id: FileId::random(),
2102 file_size: 100,
2103 ..Default::default()
2104 };
2105
2106 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2107
2108 let task = IndexBuildTask {
2110 region_id,
2111 file,
2112 target_region_metadata: version_control.current().version.metadata.clone(),
2113 source: IndexBuildSource::new(
2114 file_meta,
2115 version_control.current().version.metadata.schema_version,
2116 ),
2117 reason: IndexBuildType::Flush,
2118 access_layer: env.access_layer.clone(),
2119 listener: WorkerListener::default(),
2120 manifest_ctx,
2121 write_cache: None,
2122 cache_manager: None,
2123 file_purger,
2124 indexer_builder,
2125 request_sender: tx,
2126 result_sender: result_tx,
2127 };
2128
2129 scheduler
2131 .schedule_build(&version_control, task)
2132 .await
2133 .unwrap();
2134 match result_rx.recv().await.unwrap() {
2135 Ok(outcome) => {
2136 if outcome == IndexBuildOutcome::Finished {
2137 panic!("Expect aborted result due to missing SST file")
2138 }
2139 }
2140 _ => panic!("Expect aborted result due to missing SST file"),
2141 }
2142 }
2143
2144 #[tokio::test]
2145 async fn test_index_build_task_foreign_file_uses_target_metadata() {
2146 let env = SchedulerEnv::new().await;
2147 let mut scheduler = env.mock_index_build_scheduler(4);
2148 let source_metadata = Arc::new(sst_region_metadata());
2149 let mut target_metadata = (*source_metadata).clone();
2150 target_metadata.region_id = RegionId::new(1, 3);
2151 let mut target_builder = RegionMetadataBuilder::new(target_metadata.region_id);
2152 for mut column_metadata in target_metadata.column_metadatas.clone() {
2153 if column_metadata.column_id == 2 {
2154 column_metadata.column_schema =
2155 column_metadata.column_schema.with_inverted_index(true);
2156 }
2157 target_builder.push_column_metadata(column_metadata);
2158 }
2159 let partition_expr = col("field_0")
2160 .gt_eq(Value::UInt64(100))
2161 .and(col("field_0").lt(Value::UInt64(200)));
2162 target_builder
2163 .primary_key(target_metadata.primary_key.clone())
2164 .partition_expr_json(Some(partition_expr.as_json_str().unwrap()))
2165 .bump_version();
2166 let target_metadata = Arc::new(target_builder.build().unwrap());
2167 let manifest_ctx = env.mock_manifest_context(target_metadata.clone()).await;
2168 let region_id = target_metadata.region_id;
2169 let file_purger = Arc::new(NoopFilePurger {});
2170 let sst_info = mock_sst_file(source_metadata.clone(), &env, IndexBuildMode::Async).await;
2171 let file_meta = FileMeta {
2172 region_id: source_metadata.region_id,
2173 file_id: sst_info.file_id,
2174 file_size: sst_info.file_size,
2175 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
2176 available_indexes: smallvec![IndexType::InvertedIndex],
2177 index_file_size: 0,
2179 index_version: 0,
2180 num_rows: sst_info.num_rows as u64,
2181 num_row_groups: sst_info.num_row_groups,
2182 ..Default::default()
2183 };
2184 seed_manifest_file(&manifest_ctx, &file_meta).await;
2185 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2186 let version_control =
2187 mock_version_control(target_metadata.clone(), file_purger.clone(), files).await;
2188 let indexer_builder = mock_indexer_builder(target_metadata.clone(), &env).await;
2189
2190 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2191
2192 let (tx, mut rx) = mpsc::channel(4);
2194 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2195 let task = IndexBuildTask {
2196 region_id,
2197 file,
2198 target_region_metadata: version_control.current().version.metadata.clone(),
2199 source: IndexBuildSource::new(
2200 file_meta.clone(),
2201 version_control.current().version.metadata.schema_version,
2202 ),
2203 reason: IndexBuildType::Flush,
2204 access_layer: env.access_layer.clone(),
2205 listener: WorkerListener::default(),
2206 manifest_ctx,
2207 write_cache: None,
2208 cache_manager: None,
2209 file_purger,
2210 indexer_builder,
2211 request_sender: tx,
2212 result_sender: result_tx,
2213 };
2214
2215 scheduler
2216 .schedule_build(&version_control, task)
2217 .await
2218 .unwrap();
2219
2220 match result_rx.recv().await.unwrap() {
2222 Ok(outcome) => {
2223 assert_eq!(outcome, IndexBuildOutcome::Finished);
2224 }
2225 _ => panic!("Expect finished result"),
2226 }
2227
2228 let worker_req = rx.recv().await.unwrap().request;
2230 match worker_req {
2231 WorkerRequest::Background {
2232 region_id: req_region_id,
2233 notify: BackgroundNotify::IndexBuildFinished(finished),
2234 } => {
2235 assert_eq!(req_region_id, region_id);
2236 let updated_meta = &finished.file_meta;
2237
2238 assert!(!updated_meta.available_indexes.is_empty());
2240 assert!(updated_meta.index_file_size > 0);
2241 assert_eq!(updated_meta.file_id, file_meta.file_id);
2242 assert_eq!(updated_meta.index_version, 1);
2243 let field_0_index = updated_meta
2244 .indexes
2245 .iter()
2246 .find(|index| index.column_id == 2)
2247 .expect("field_0 should have an inverted index");
2248 assert_eq!(
2249 field_0_index.created_indexes.as_slice(),
2250 [IndexType::InvertedIndex]
2251 );
2252 }
2253 _ => panic!("Unexpected worker request: {:?}", worker_req),
2254 }
2255
2256 let puffin_reader = env
2257 .access_layer
2258 .build_puffin_manager()
2259 .reader(&RegionIndexId::new(
2260 RegionFileId::new(source_metadata.region_id, file_meta.file_id),
2261 1,
2262 ))
2263 .await
2264 .unwrap();
2265 let blob = puffin_reader
2266 .blob(inverted_index::INDEX_BLOB_TYPE)
2267 .await
2268 .unwrap();
2269 let blob_reader = blob.reader().await.unwrap();
2270 let index_metadata =
2271 index::inverted_index::format::reader::InvertedIndexBlobReader::new(blob_reader)
2272 .metadata(None)
2273 .await
2274 .unwrap();
2275 assert!(index_metadata.metas.contains_key("2"));
2276 assert_eq!(index_metadata.total_row_count, 100);
2277 }
2278
2279 async fn schedule_index_build_task_with_mode(build_mode: IndexBuildMode) {
2280 let env = SchedulerEnv::new().await;
2281 let mut scheduler = env.mock_index_build_scheduler(4);
2282 let metadata = Arc::new(sst_region_metadata());
2283 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2284 let file_purger = Arc::new(NoopFilePurger {});
2285 let region_id = metadata.region_id;
2286 let sst_info = mock_sst_file(metadata.clone(), &env, build_mode.clone()).await;
2287 let file_meta = FileMeta {
2288 region_id,
2289 file_id: sst_info.file_id,
2290 file_size: sst_info.file_size,
2291 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
2292 index_file_size: sst_info.index_metadata.file_size,
2293 num_rows: sst_info.num_rows as u64,
2294 num_row_groups: sst_info.num_row_groups,
2295 ..Default::default()
2296 };
2297 seed_manifest_file(&manifest_ctx, &file_meta).await;
2298 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2299 let version_control =
2300 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2301 let indexer_builder = mock_indexer_builder(metadata.clone(), &env).await;
2302
2303 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2304
2305 let (tx, _rx) = mpsc::channel(4);
2307 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2308 let task = IndexBuildTask {
2309 region_id,
2310 file,
2311 target_region_metadata: version_control.current().version.metadata.clone(),
2312 source: IndexBuildSource::new(
2313 file_meta.clone(),
2314 version_control.current().version.metadata.schema_version,
2315 ),
2316 reason: IndexBuildType::Flush,
2317 access_layer: env.access_layer.clone(),
2318 listener: WorkerListener::default(),
2319 manifest_ctx,
2320 write_cache: None,
2321 cache_manager: None,
2322 file_purger,
2323 indexer_builder,
2324 request_sender: tx,
2325 result_sender: result_tx,
2326 };
2327
2328 scheduler
2329 .schedule_build(&version_control, task)
2330 .await
2331 .unwrap();
2332
2333 let puffin_path = location::index_file_path(
2334 env.access_layer.table_dir(),
2335 RegionIndexId::new(RegionFileId::new(region_id, file_meta.file_id), 0),
2336 env.access_layer.path_type(),
2337 );
2338
2339 if build_mode == IndexBuildMode::Async {
2340 assert!(
2342 !env.access_layer
2343 .object_store()
2344 .exists(&puffin_path)
2345 .await
2346 .unwrap()
2347 );
2348 } else {
2349 assert!(
2351 env.access_layer
2352 .object_store()
2353 .exists(&puffin_path)
2354 .await
2355 .unwrap()
2356 );
2357 }
2358
2359 match result_rx.recv().await.unwrap() {
2361 Ok(outcome) => {
2362 assert_eq!(outcome, IndexBuildOutcome::Finished);
2363 }
2364 _ => panic!("Expect finished result"),
2365 }
2366
2367 assert!(
2369 env.access_layer
2370 .object_store()
2371 .exists(&puffin_path)
2372 .await
2373 .unwrap()
2374 );
2375 }
2376
2377 #[tokio::test]
2378 async fn test_index_build_task_build_mode() {
2379 schedule_index_build_task_with_mode(IndexBuildMode::Async).await;
2380 schedule_index_build_task_with_mode(IndexBuildMode::Sync).await;
2381 }
2382
2383 #[tokio::test]
2384 async fn test_index_build_task_no_index() {
2385 let env = SchedulerEnv::new().await;
2386 let mut scheduler = env.mock_index_build_scheduler(4);
2387 let mut metadata = sst_region_metadata();
2388 metadata.column_metadatas.iter_mut().for_each(|col| {
2390 col.column_schema.set_inverted_index(false);
2391 let _ = col.column_schema.unset_skipping_options();
2392 });
2393 let region_id = metadata.region_id;
2394 let metadata = Arc::new(metadata);
2395 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2396 let file_purger = Arc::new(NoopFilePurger {});
2397 let sst_info = mock_sst_file(metadata.clone(), &env, IndexBuildMode::Async).await;
2398 let file_meta = FileMeta {
2399 region_id,
2400 file_id: sst_info.file_id,
2401 file_size: sst_info.file_size,
2402 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
2403 index_file_size: sst_info.index_metadata.file_size,
2404 num_rows: sst_info.num_rows as u64,
2405 num_row_groups: sst_info.num_row_groups,
2406 ..Default::default()
2407 };
2408 seed_manifest_file(&manifest_ctx, &file_meta).await;
2409 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2410 let version_control =
2411 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2412 let indexer_builder = mock_indexer_builder(metadata.clone(), &env).await;
2413
2414 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2415
2416 let (tx, mut rx) = mpsc::channel(4);
2418 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2419 let task = IndexBuildTask {
2420 region_id,
2421 file,
2422 target_region_metadata: version_control.current().version.metadata.clone(),
2423 source: IndexBuildSource::new(
2424 file_meta.clone(),
2425 version_control.current().version.metadata.schema_version,
2426 ),
2427 reason: IndexBuildType::Flush,
2428 access_layer: env.access_layer.clone(),
2429 listener: WorkerListener::default(),
2430 manifest_ctx,
2431 write_cache: None,
2432 cache_manager: None,
2433 file_purger,
2434 indexer_builder,
2435 request_sender: tx,
2436 result_sender: result_tx,
2437 };
2438
2439 scheduler
2440 .schedule_build(&version_control, task)
2441 .await
2442 .unwrap();
2443
2444 match result_rx.recv().await.unwrap() {
2446 Ok(outcome) => {
2447 assert_eq!(outcome, IndexBuildOutcome::Finished);
2448 }
2449 _ => panic!("Expect finished result"),
2450 }
2451
2452 let _ = rx.recv().await.is_none();
2454 }
2455
2456 #[tokio::test]
2457 async fn test_index_build_task_with_write_cache() {
2458 let env = SchedulerEnv::new().await;
2459 let mut scheduler = env.mock_index_build_scheduler(4);
2460 let metadata = Arc::new(sst_region_metadata());
2461 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2462 let file_purger = Arc::new(NoopFilePurger {});
2463 let region_id = metadata.region_id;
2464
2465 let (dir, factory) = PuffinManagerFactory::new_for_test_async("test_write_cache").await;
2466 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
2467
2468 let write_cache = Arc::new(
2470 WriteCache::new_fs(
2471 dir.path().to_str().unwrap(),
2472 ReadableSize::mb(10),
2473 None,
2474 None,
2475 true, factory,
2477 intm_manager,
2478 ReadableSize::mb(10),
2479 )
2480 .await
2481 .unwrap(),
2482 );
2483 let indexer_builder = Arc::new(IndexerBuilderImpl {
2485 build_type: IndexBuildType::Flush,
2486 metadata: metadata.clone(),
2487 puffin_manager: write_cache.build_puffin_manager().clone(),
2488 write_cache_enabled: true,
2489 intermediate_manager: write_cache.intermediate_manager().clone(),
2490 index_options: IndexOptions::default(),
2491 inverted_index_config: InvertedIndexConfig::default(),
2492 fulltext_index_config: FulltextIndexConfig::default(),
2493 bloom_filter_index_config: BloomFilterConfig::default(),
2494 #[cfg(feature = "vector_index")]
2495 vector_index_config: Default::default(),
2496 });
2497
2498 let sst_info = mock_sst_file(metadata.clone(), &env, IndexBuildMode::Async).await;
2499 let file_meta = FileMeta {
2500 region_id,
2501 file_id: sst_info.file_id,
2502 file_size: sst_info.file_size,
2503 index_file_size: sst_info.index_metadata.file_size,
2504 num_rows: sst_info.num_rows as u64,
2505 num_row_groups: sst_info.num_row_groups,
2506 ..Default::default()
2507 };
2508 seed_manifest_file(&manifest_ctx, &file_meta).await;
2509 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2510 let version_control =
2511 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2512
2513 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2514
2515 let (tx, mut _rx) = mpsc::channel(4);
2517 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2518 let task = IndexBuildTask {
2519 region_id,
2520 file,
2521 target_region_metadata: version_control.current().version.metadata.clone(),
2522 source: IndexBuildSource::new(
2523 file_meta.clone(),
2524 version_control.current().version.metadata.schema_version,
2525 ),
2526 reason: IndexBuildType::Flush,
2527 access_layer: env.access_layer.clone(),
2528 listener: WorkerListener::default(),
2529 manifest_ctx,
2530 write_cache: Some(write_cache.clone()),
2531 cache_manager: None,
2532 file_purger,
2533 indexer_builder,
2534 request_sender: tx,
2535 result_sender: result_tx,
2536 };
2537
2538 scheduler
2539 .schedule_build(&version_control, task)
2540 .await
2541 .unwrap();
2542
2543 match result_rx.recv().await.unwrap() {
2545 Ok(outcome) => {
2546 assert_eq!(outcome, IndexBuildOutcome::Finished);
2547 }
2548 _ => panic!("Expect finished result"),
2549 }
2550
2551 let index_key = IndexKey::new(
2553 region_id,
2554 file_meta.file_id,
2555 FileType::Puffin(sst_info.index_metadata.version),
2556 );
2557 assert!(write_cache.file_cache().contains_key(&index_key));
2558 }
2559
2560 async fn create_mock_task_for_schedule(
2561 env: &SchedulerEnv,
2562 file_id: FileId,
2563 region_id: RegionId,
2564 reason: IndexBuildType,
2565 ) -> IndexBuildTask {
2566 create_mock_task_for_schedule_with_result(env, file_id, region_id, reason)
2567 .await
2568 .0
2569 }
2570
2571 async fn create_mock_task_for_schedule_with_result(
2574 env: &SchedulerEnv,
2575 file_id: FileId,
2576 region_id: RegionId,
2577 reason: IndexBuildType,
2578 ) -> (IndexBuildTask, mpsc::Receiver<Result<IndexBuildOutcome>>) {
2579 let metadata = Arc::new(sst_region_metadata());
2580 let schema_version = metadata.schema_version;
2581 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2582 let file_purger = Arc::new(NoopFilePurger {});
2583 let indexer_builder = mock_indexer_builder(metadata.clone(), env).await;
2584 let (tx, _rx) = mpsc::channel(4);
2585 let (result_tx, result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2586
2587 let file_meta = FileMeta {
2588 region_id,
2589 file_id,
2590 file_size: 100,
2591 ..Default::default()
2592 };
2593
2594 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2595
2596 let task = IndexBuildTask {
2597 region_id,
2598 file,
2599 target_region_metadata: metadata,
2600 source: IndexBuildSource::new(file_meta, schema_version),
2601 reason,
2602 access_layer: env.access_layer.clone(),
2603 listener: WorkerListener::default(),
2604 manifest_ctx,
2605 write_cache: None,
2606 cache_manager: None,
2607 file_purger,
2608 indexer_builder,
2609 request_sender: tx,
2610 result_sender: result_tx,
2611 };
2612 (task, result_rx)
2613 }
2614
2615 #[tokio::test]
2616 async fn test_scheduler_coalesces_latest_schema_generation_per_sst() {
2617 let job_scheduler = Arc::new(VecScheduler::default());
2618 let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2619 let mut scheduler = env.mock_index_build_scheduler(2);
2620 let metadata = Arc::new(sst_region_metadata());
2621 let region_id = metadata.region_id;
2622 let file_id = FileId::random();
2623 let file_purger = Arc::new(NoopFilePurger {});
2624 let files = HashMap::from([(
2625 file_id,
2626 FileMeta {
2627 region_id,
2628 file_id,
2629 file_size: 100,
2630 ..Default::default()
2631 },
2632 )]);
2633 let version_control = mock_version_control(metadata, file_purger, files).await;
2634
2635 let (active, _active_rx) = create_mock_task_for_schedule_with_result(
2636 &env,
2637 file_id,
2638 region_id,
2639 IndexBuildType::Manual,
2640 )
2641 .await;
2642 scheduler
2643 .schedule_build(&version_control, active)
2644 .await
2645 .unwrap();
2646 assert_eq!(job_scheduler.num_jobs(), 1);
2647
2648 let (mut older, mut older_rx) = create_mock_task_for_schedule_with_result(
2649 &env,
2650 file_id,
2651 region_id,
2652 IndexBuildType::SchemaChange,
2653 )
2654 .await;
2655 older.source.schema_version = 1;
2656 scheduler
2657 .schedule_build(&version_control, older)
2658 .await
2659 .unwrap();
2660
2661 let (mut latest, mut latest_rx) = create_mock_task_for_schedule_with_result(
2662 &env,
2663 file_id,
2664 region_id,
2665 IndexBuildType::SchemaChange,
2666 )
2667 .await;
2668 latest.source.schema_version = 2;
2669 scheduler
2670 .schedule_build(&version_control, latest)
2671 .await
2672 .unwrap();
2673
2674 let replaced = tokio::time::timeout(std::time::Duration::from_secs(5), older_rx.recv())
2675 .await
2676 .expect("replaced pending task result sender was not completed")
2677 .expect("replaced pending task result channel closed");
2678 assert!(matches!(
2679 replaced,
2680 Ok(IndexBuildOutcome::Aborted(reason)) if reason.contains("coalesced")
2681 ));
2682 assert!(matches!(
2683 latest_rx.try_recv(),
2684 Err(mpsc::error::TryRecvError::Empty)
2685 ));
2686
2687 let status = &scheduler.region_status[®ion_id];
2688 assert_eq!(status.building_files.len(), 1);
2689 assert_eq!(status.pending_tasks.len(), 1);
2690 assert_eq!(job_scheduler.num_jobs(), 1);
2691
2692 scheduler.on_task_stopped(region_id, file_id);
2693 let status = &scheduler.region_status[®ion_id];
2694 assert_eq!(status.building_files.len(), 1);
2695 assert!(status.pending_tasks.is_empty());
2696 assert_eq!(job_scheduler.num_jobs(), 2);
2697 }
2698
2699 #[tokio::test]
2700 async fn test_scheduler_completes_sender_when_job_is_rejected() {
2701 let job_scheduler = Arc::new(LocalScheduler::new(1));
2702 job_scheduler.stop(false).await.unwrap();
2703 let env = SchedulerEnv::new().await.scheduler(job_scheduler);
2704 let mut scheduler = env.mock_index_build_scheduler(1);
2705 let metadata = Arc::new(sst_region_metadata());
2706 let region_id = metadata.region_id;
2707 let file_id = FileId::random();
2708 let file_purger = Arc::new(NoopFilePurger {});
2709 let files = HashMap::from([(
2710 file_id,
2711 FileMeta {
2712 region_id,
2713 file_id,
2714 file_size: 100,
2715 ..Default::default()
2716 },
2717 )]);
2718 let version_control = mock_version_control(metadata, file_purger, files).await;
2719 let (task, mut result_rx) = create_mock_task_for_schedule_with_result(
2720 &env,
2721 file_id,
2722 region_id,
2723 IndexBuildType::Flush,
2724 )
2725 .await;
2726
2727 scheduler
2728 .schedule_build(&version_control, task)
2729 .await
2730 .unwrap();
2731
2732 let result = tokio::time::timeout(std::time::Duration::from_secs(5), result_rx.recv())
2733 .await
2734 .expect("scheduler rejection did not complete the result sender")
2735 .expect("result channel closed without a result");
2736 assert!(result.is_err());
2737 assert!(!scheduler.region_status.contains_key(®ion_id));
2738 }
2739
2740 #[tokio::test]
2741 async fn test_scheduler_comprehensive() {
2742 let env = SchedulerEnv::new().await;
2743 let mut scheduler = env.mock_index_build_scheduler(2);
2744 let metadata = Arc::new(sst_region_metadata());
2745 let region_id = metadata.region_id;
2746 let file_purger = Arc::new(NoopFilePurger {});
2747
2748 let file_id1 = FileId::random();
2750 let file_id2 = FileId::random();
2751 let file_id3 = FileId::random();
2752 let file_id4 = FileId::random();
2753 let file_id5 = FileId::random();
2754
2755 let mut files = HashMap::new();
2756 for file_id in [file_id1, file_id2, file_id3, file_id4, file_id5] {
2757 files.insert(
2758 file_id,
2759 FileMeta {
2760 region_id,
2761 file_id,
2762 file_size: 100,
2763 ..Default::default()
2764 },
2765 );
2766 }
2767
2768 let version_control = mock_version_control(metadata, file_purger, files).await;
2769
2770 let task1 =
2772 create_mock_task_for_schedule(&env, file_id1, region_id, IndexBuildType::Flush).await;
2773 assert!(
2774 scheduler
2775 .schedule_build(&version_control, task1)
2776 .await
2777 .is_ok()
2778 );
2779 assert!(scheduler.region_status.contains_key(®ion_id));
2780 let status = scheduler.region_status.get(®ion_id).unwrap();
2781 assert_eq!(status.building_files.len(), 1);
2782 assert!(status.building_files.contains(&file_id1));
2783
2784 let task1_dup =
2786 create_mock_task_for_schedule(&env, file_id1, region_id, IndexBuildType::Flush).await;
2787 scheduler
2788 .schedule_build(&version_control, task1_dup)
2789 .await
2790 .unwrap();
2791 let status = scheduler.region_status.get(®ion_id).unwrap();
2792 assert_eq!(status.building_files.len(), 1); let task2 =
2796 create_mock_task_for_schedule(&env, file_id2, region_id, IndexBuildType::Flush).await;
2797 scheduler
2798 .schedule_build(&version_control, task2)
2799 .await
2800 .unwrap();
2801 let status = scheduler.region_status.get(®ion_id).unwrap();
2802 assert_eq!(status.building_files.len(), 2); assert_eq!(status.pending_tasks.len(), 0);
2804
2805 let task3 =
2808 create_mock_task_for_schedule(&env, file_id3, region_id, IndexBuildType::Compact).await;
2809 let task4 =
2810 create_mock_task_for_schedule(&env, file_id4, region_id, IndexBuildType::SchemaChange)
2811 .await;
2812 let task5 =
2813 create_mock_task_for_schedule(&env, file_id5, region_id, IndexBuildType::Manual).await;
2814
2815 scheduler
2816 .schedule_build(&version_control, task3)
2817 .await
2818 .unwrap();
2819 scheduler
2820 .schedule_build(&version_control, task4)
2821 .await
2822 .unwrap();
2823 scheduler
2824 .schedule_build(&version_control, task5)
2825 .await
2826 .unwrap();
2827
2828 let status = scheduler.region_status.get(®ion_id).unwrap();
2829 assert_eq!(status.building_files.len(), 2); assert_eq!(status.pending_tasks.len(), 3); scheduler.on_task_stopped(region_id, file_id1);
2834 let status = scheduler.region_status.get(®ion_id).unwrap();
2835 assert!(!status.building_files.contains(&file_id1));
2836 assert_eq!(status.building_files.len(), 2); assert_eq!(status.pending_tasks.len(), 2); assert!(status.building_files.contains(&file_id5));
2840
2841 scheduler.on_task_stopped(region_id, file_id2);
2843 let status = scheduler.region_status.get(®ion_id).unwrap();
2844 assert_eq!(status.building_files.len(), 2);
2845 assert_eq!(status.pending_tasks.len(), 1); assert!(status.building_files.contains(&file_id4)); scheduler.on_task_stopped(region_id, file_id5);
2850 scheduler.on_task_stopped(region_id, file_id4);
2851
2852 let status = scheduler.region_status.get(®ion_id).unwrap();
2853 assert_eq!(status.building_files.len(), 1); assert_eq!(status.pending_tasks.len(), 0);
2855 assert!(status.building_files.contains(&file_id3));
2856
2857 scheduler.on_task_stopped(region_id, file_id3);
2858
2859 assert!(!scheduler.region_status.contains_key(®ion_id));
2861
2862 let task6 =
2864 create_mock_task_for_schedule(&env, file_id1, region_id, IndexBuildType::Flush).await;
2865 let task7 =
2866 create_mock_task_for_schedule(&env, file_id2, region_id, IndexBuildType::Flush).await;
2867 let task8 =
2868 create_mock_task_for_schedule(&env, file_id3, region_id, IndexBuildType::Manual).await;
2869
2870 scheduler
2871 .schedule_build(&version_control, task6)
2872 .await
2873 .unwrap();
2874 scheduler
2875 .schedule_build(&version_control, task7)
2876 .await
2877 .unwrap();
2878 scheduler
2879 .schedule_build(&version_control, task8)
2880 .await
2881 .unwrap();
2882
2883 assert!(scheduler.region_status.contains_key(®ion_id));
2884 let status = scheduler.region_status.get(®ion_id).unwrap();
2885 assert_eq!(status.building_files.len(), 2);
2886 assert_eq!(status.pending_tasks.len(), 1);
2887
2888 scheduler
2889 .on_failure(
2890 region_id,
2891 Arc::new(
2892 crate::error::UnexpectedSnafu {
2893 reason: "index build failed".to_string(),
2894 }
2895 .build(),
2896 ),
2897 )
2898 .await;
2899 let status = scheduler.region_status.get(®ion_id).unwrap();
2900 assert!(!status.retiring);
2901 assert_eq!(status.building_files.len(), 2);
2902 assert!(status.pending_tasks.is_empty());
2903
2904 scheduler.on_task_stopped(region_id, file_id1);
2905 assert!(scheduler.region_status.contains_key(®ion_id));
2906 scheduler.on_task_stopped(region_id, file_id2);
2907 assert!(!scheduler.region_status.contains_key(®ion_id));
2908 }
2909
2910 async fn setup_scheduler_with_pending_tasks(
2914 env: &SchedulerEnv,
2915 ) -> (
2916 IndexBuildScheduler,
2917 mpsc::Receiver<Result<IndexBuildOutcome>>,
2918 mpsc::Receiver<Result<IndexBuildOutcome>>,
2919 VersionControlRef,
2920 RegionId,
2921 FileId, ) {
2923 let metadata = Arc::new(sst_region_metadata());
2924 let region_id = metadata.region_id;
2925 let file_purger = Arc::new(NoopFilePurger {});
2926
2927 let file_id1 = FileId::random();
2928 let file_id2 = FileId::random();
2929 let file_id3 = FileId::random();
2930
2931 let files = HashMap::from([
2932 (
2933 file_id1,
2934 FileMeta {
2935 region_id,
2936 file_id: file_id1,
2937 file_size: 100,
2938 ..Default::default()
2939 },
2940 ),
2941 (
2942 file_id2,
2943 FileMeta {
2944 region_id,
2945 file_id: file_id2,
2946 file_size: 100,
2947 ..Default::default()
2948 },
2949 ),
2950 (
2951 file_id3,
2952 FileMeta {
2953 region_id,
2954 file_id: file_id3,
2955 file_size: 100,
2956 ..Default::default()
2957 },
2958 ),
2959 ]);
2960 let version_control =
2961 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2962
2963 let mut scheduler = env.mock_index_build_scheduler(1);
2964
2965 let (task1, _rx1) = create_mock_task_for_schedule_with_result(
2969 env,
2970 file_id1,
2971 region_id,
2972 IndexBuildType::Flush,
2973 )
2974 .await;
2975 let (task2, rx2) = create_mock_task_for_schedule_with_result(
2976 env,
2977 file_id2,
2978 region_id,
2979 IndexBuildType::Flush,
2980 )
2981 .await;
2982 let (task3, rx3) = create_mock_task_for_schedule_with_result(
2983 env,
2984 file_id3,
2985 region_id,
2986 IndexBuildType::Flush,
2987 )
2988 .await;
2989
2990 scheduler
2991 .schedule_build(&version_control, task1)
2992 .await
2993 .unwrap();
2994 scheduler
2995 .schedule_build(&version_control, task2)
2996 .await
2997 .unwrap();
2998 scheduler
2999 .schedule_build(&version_control, task3)
3000 .await
3001 .unwrap();
3002
3003 assert!(scheduler.region_status.contains_key(®ion_id));
3005 let status = scheduler.region_status.get(®ion_id).unwrap();
3006 assert_eq!(status.building_files.len(), 1);
3007 assert_eq!(status.pending_tasks.len(), 2);
3008
3009 (scheduler, rx2, rx3, version_control, region_id, file_id1)
3010 }
3011
3012 fn assert_lifecycle_error(
3017 err: crate::error::Error,
3018 expected_source: &str,
3019 lifecycle_name: &str,
3020 ) {
3021 let crate::error::Error::BuildIndexAsync { source, .. } = err else {
3022 panic!("[{lifecycle_name}] Expected BuildIndexAsync outer error, got: {err:?}");
3023 };
3024 let actual = match &*source {
3025 crate::error::Error::RegionDropped { .. } => "dropped",
3026 crate::error::Error::RegionClosed { .. } => "closed",
3027 crate::error::Error::RegionTruncated { .. } => "truncated",
3028 other => panic!(
3029 "[{lifecycle_name}] Expected lifecycle source variant (RegionDropped/RegionClosed/RegionTruncated), got: {other:?}"
3030 ),
3031 };
3032 assert_eq!(
3033 actual, expected_source,
3034 "[{lifecycle_name}] Source variant mismatch"
3035 );
3036 }
3037
3038 async fn recv_lifecycle_error(
3041 rx: &mut mpsc::Receiver<Result<IndexBuildOutcome>>,
3042 expected_source: &str,
3043 lifecycle_name: &str,
3044 ) {
3045 let result = tokio::time::timeout(std::time::Duration::from_secs(5), rx.recv())
3046 .await
3047 .unwrap_or_else(|_| {
3048 panic!(
3049 "[{lifecycle_name}] Timeout (5s) waiting for lifecycle error from pending task"
3050 )
3051 });
3052 let result = result.unwrap_or_else(|| {
3053 panic!("[{lifecycle_name}] Channel closed without receiving lifecycle error")
3054 });
3055 let err = result.unwrap_err();
3056 assert_lifecycle_error(err, expected_source, lifecycle_name);
3057 }
3058
3059 #[tokio::test]
3060 async fn test_scheduler_lifecycle_cleanup() {
3061 let env = SchedulerEnv::new().await;
3062
3063 {
3065 let (mut scheduler, mut rx2, mut rx3, _version_control, region_id, building_file_id) =
3066 setup_scheduler_with_pending_tasks(&env).await;
3067
3068 scheduler.on_region_dropped(region_id).await;
3069
3070 let status = scheduler.region_status.get(®ion_id).unwrap();
3071 assert!(status.retiring);
3072 assert_eq!(status.building_files.len(), 1);
3073 assert!(status.pending_tasks.is_empty());
3074
3075 recv_lifecycle_error(&mut rx2, "dropped", "on_region_dropped").await;
3077 recv_lifecycle_error(&mut rx3, "dropped", "on_region_dropped").await;
3078
3079 scheduler.on_task_stopped(region_id, building_file_id);
3081 assert!(!scheduler.region_status.contains_key(®ion_id));
3082 }
3083
3084 {
3086 let (mut scheduler, mut rx2, mut rx3, version_control, region_id, building_file_id) =
3087 setup_scheduler_with_pending_tasks(&env).await;
3088
3089 scheduler.on_region_closed(region_id).await;
3090
3091 let status = scheduler.region_status.get(®ion_id).unwrap();
3092 assert!(status.retiring);
3093 assert_eq!(status.building_files.len(), 1);
3094 assert!(status.pending_tasks.is_empty());
3095
3096 recv_lifecycle_error(&mut rx2, "closed", "on_region_closed").await;
3097 recv_lifecycle_error(&mut rx3, "closed", "on_region_closed").await;
3098
3099 let (reopened_task, _reopened_rx) = create_mock_task_for_schedule_with_result(
3102 &env,
3103 building_file_id,
3104 region_id,
3105 IndexBuildType::Manual,
3106 )
3107 .await;
3108 scheduler
3109 .schedule_build(&version_control, reopened_task)
3110 .await
3111 .unwrap();
3112 let status = scheduler.region_status.get(®ion_id).unwrap();
3113 assert!(status.retiring);
3114 assert_eq!(status.building_files.len(), 1);
3115 assert_eq!(status.pending_tasks.len(), 1);
3116
3117 scheduler
3120 .on_failure(region_id, Arc::new(RegionClosedSnafu { region_id }.build()))
3121 .await;
3122 assert_eq!(scheduler.region_status[®ion_id].pending_tasks.len(), 1);
3123
3124 scheduler.on_task_stopped(region_id, building_file_id);
3125 let status = scheduler.region_status.get(®ion_id).unwrap();
3126 assert!(!status.retiring);
3127 assert_eq!(status.building_files.len(), 1);
3128 assert!(status.pending_tasks.is_empty());
3129
3130 scheduler.on_task_stopped(region_id, building_file_id);
3131 assert!(!scheduler.region_status.contains_key(®ion_id));
3132 }
3133
3134 {
3136 let (mut scheduler, mut rx2, mut rx3, _version_control, region_id, building_file_id) =
3137 setup_scheduler_with_pending_tasks(&env).await;
3138
3139 scheduler.on_region_truncated(region_id).await;
3140
3141 let status = scheduler.region_status.get(®ion_id).unwrap();
3142 assert!(status.retiring);
3143 assert_eq!(status.building_files.len(), 1);
3144 assert!(status.pending_tasks.is_empty());
3145
3146 recv_lifecycle_error(&mut rx2, "truncated", "on_region_truncated").await;
3147 recv_lifecycle_error(&mut rx3, "truncated", "on_region_truncated").await;
3148
3149 scheduler.on_task_stopped(region_id, building_file_id);
3150 assert!(!scheduler.region_status.contains_key(®ion_id));
3151 }
3152 }
3153}