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