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) source: IndexBuildSource,
711 pub reason: IndexBuildType,
712 pub access_layer: AccessLayerRef,
713 pub(crate) listener: WorkerListener,
714 pub(crate) manifest_ctx: ManifestContextRef,
715 pub write_cache: Option<WriteCacheRef>,
716 pub cache_manager: Option<CacheManagerRef>,
717 pub file_purger: FilePurgerRef,
718 pub indexer_builder: Arc<dyn IndexerBuilder + Send + Sync>,
721 pub(crate) request_sender: Sender<WorkerRequestWithTime>,
723 pub(crate) result_sender: ResultMpscSender,
725}
726
727impl std::fmt::Debug for IndexBuildTask {
728 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
729 f.debug_struct("IndexBuildTask")
730 .field("region_id", &self.region_id)
731 .field("origin_region_id", &self.source.file_meta.region_id)
732 .field("file_id", &self.source.file_meta.file_id)
733 .field("schema_version", &self.source.schema_version)
734 .field("reason", &self.reason)
735 .finish()
736 }
737}
738
739impl IndexBuildTask {
740 pub async fn on_success(&self, outcome: IndexBuildOutcome) {
742 let _ = self.result_sender.send(Ok(outcome)).await;
743 }
744
745 pub async fn on_failure(&self, err: Arc<Error>) {
747 let _ = self
748 .result_sender
749 .send(Err(err.clone()).context(BuildIndexAsyncSnafu {
750 region_id: self.region_id,
751 }))
752 .await;
753 }
754
755 fn into_index_build_job(mut self, version_control: VersionControlRef) -> Job {
756 Box::pin(async move {
757 self.do_index_build(version_control).await;
758 })
759 }
760
761 async fn do_index_build(&mut self, version_control: VersionControlRef) {
762 self.listener
763 .on_index_build_begin(RegionFileId::new(
764 self.source.file_meta.region_id,
765 self.source.file_meta.file_id,
766 ))
767 .await;
768 match self.index_build(version_control).await {
769 Ok(outcome) => self.on_success(outcome).await,
770 Err(e) => {
771 warn!(
772 e; "Index build task failed, region: {}, file_id: {}",
773 self.region_id, self.source.file_meta.file_id,
774 );
775 self.on_failure(e.into()).await
776 }
777 }
778 let worker_request = WorkerRequest::Background {
779 region_id: self.region_id,
780 notify: BackgroundNotify::IndexBuildStopped(IndexBuildStopped {
781 file_id: self.source.file_meta.file_id,
782 }),
783 };
784 let _ = self
785 .request_sender
786 .send(WorkerRequestWithTime::new(worker_request))
787 .await;
788 }
789
790 async fn check_sst_file_exists(&self, version_control: &VersionControlRef) -> bool {
792 let file_id = self.source.file_meta.file_id;
793 let level = self.source.file_meta.level;
794 let version = version_control.current().version;
796
797 let Some(level_files) = version.ssts.levels().get(level as usize) else {
798 warn!(
799 "File id {} not found in level {} for index build, region: {}",
800 file_id, level, self.region_id
801 );
802 return false;
803 };
804
805 match level_files.files.get(&file_id) {
806 Some(handle) if !handle.is_deleted() && !handle.compacting() => {
807 true
811 }
812 _ => {
813 warn!(
814 "File id {} not found in region version for index build, region: {}",
815 file_id, self.region_id
816 );
817 false
818 }
819 }
820 }
821
822 async fn index_build(
823 &mut self,
824 version_control: VersionControlRef,
825 ) -> Result<IndexBuildOutcome> {
826 let new_index_version = self
827 .source
828 .file_meta
829 .index_version()
830 .map_or(0, |version| version + 1);
831
832 if !self.check_sst_file_exists(&version_control).await {
834 self.listener
835 .on_index_build_abort(RegionFileId::new(
836 self.source.file_meta.region_id,
837 self.source.file_meta.file_id,
838 ))
839 .await;
840 return Ok(IndexBuildOutcome::Aborted(format!(
841 "SST file not found during index build, region: {}, file_id: {}",
842 self.region_id, self.source.file_meta.file_id
843 )));
844 }
845
846 let mut parquet_reader = self
847 .access_layer
848 .read_sst(self.file.clone()) .build()
850 .await?;
851
852 let row_group_size = parquet_reader.as_ref().and_then(|reader| {
853 reader
854 .parquet_metadata()
855 .row_groups()
856 .first()
857 .map(|row_group| row_group.num_rows() as usize)
858 .filter(|size| *size > 0)
859 });
860
861 let region_file_id = RegionFileId::new(
863 self.source.file_meta.region_id,
864 self.source.file_meta.file_id,
865 );
866 let index_file_id = region_file_id.file_id();
867 let mut indexer = self
868 .indexer_builder
869 .build(region_file_id, new_index_version, row_group_size)
870 .await;
871
872 if let Some(mut parquet_reader) = parquet_reader.take() {
873 loop {
875 match parquet_reader.next_record_batch().await {
876 Ok(Some(batch)) => {
877 indexer.update_flat(&batch).await;
878 }
879 Ok(None) => break,
880 Err(e) => {
881 indexer.abort().await;
882 return Err(e);
883 }
884 }
885 }
886 }
887 let index_output = indexer.finish().await;
888
889 if index_output.file_size > 0 {
890 if !self.check_sst_file_exists(&version_control).await {
892 indexer.abort().await;
894 self.listener
895 .on_index_build_abort(RegionFileId::new(
896 self.source.file_meta.region_id,
897 self.source.file_meta.file_id,
898 ))
899 .await;
900 return Ok(IndexBuildOutcome::Aborted(format!(
901 "SST file not found during index build, region: {}, file_id: {}",
902 self.region_id, self.source.file_meta.file_id
903 )));
904 }
905
906 self.maybe_upload_index_file(index_output.clone(), index_file_id, new_index_version)
908 .await?;
909
910 self.listener
911 .on_index_build_before_manifest_commit(region_file_id)
912 .await;
913
914 let worker_request = match self.update_manifest(index_output, new_index_version).await {
915 Ok(IndexPublication::Committed {
916 manifest_version,
917 file_meta,
918 }) => {
919 self.listener
920 .on_index_build_manifest_committed(region_file_id)
921 .await;
922 let index_build_finished = IndexBuildFinished {
923 manifest_version,
924 file_meta,
925 };
926 WorkerRequest::Background {
927 region_id: self.region_id,
928 notify: BackgroundNotify::IndexBuildFinished(index_build_finished),
929 }
930 }
931 Ok(IndexPublication::Stale(stale)) => {
932 INDEX_PUBLICATION_STALE_TOTAL
933 .with_label_values(&["manifest_commit"])
934 .inc();
935 let index_id = RegionIndexId::new(region_file_id, new_index_version);
936 cleanup_stale_index_caches(
940 index_id,
941 &self.access_layer,
942 self.cache_manager.as_ref(),
943 self.write_cache.as_ref(),
944 )
945 .await;
946 self.listener.on_index_build_abort(region_file_id).await;
947 if stale == IndexPublicationStale::SchemaChanged {
948 let retry = BuildIndexRequest {
953 region_id: self.region_id,
954 build_type: IndexBuildType::SchemaChange,
955 file_metas: vec![self.source.file_meta.clone()],
956 };
957 let worker_request = WorkerRequest::Background {
958 region_id: self.region_id,
959 notify: BackgroundNotify::IndexBuildRetry(retry),
960 };
961 let _ = self
962 .request_sender
963 .send(WorkerRequestWithTime::new(worker_request))
964 .await;
965 }
966 return Ok(IndexBuildOutcome::Aborted(format!(
967 "Index build source changed before publication, region: {}, file_id: {}",
968 self.region_id, self.source.file_meta.file_id
969 )));
970 }
971 Err(e) => {
972 let err = Arc::new(e);
973 WorkerRequest::Background {
974 region_id: self.region_id,
975 notify: BackgroundNotify::IndexBuildFailed(IndexBuildFailed { err }),
976 }
977 }
978 };
979
980 let _ = self
981 .request_sender
982 .send(WorkerRequestWithTime::new(worker_request))
983 .await;
984 }
985 Ok(IndexBuildOutcome::Finished)
986 }
987
988 async fn maybe_upload_index_file(
989 &self,
990 output: IndexOutput,
991 index_file_id: FileId,
992 index_version: u64,
993 ) -> Result<()> {
994 if let Some(write_cache) = &self.write_cache {
995 let file_id = self.source.file_meta.file_id;
996 let region_id = self.source.file_meta.region_id;
997 let remote_store = self.access_layer.object_store();
998 let mut upload_tracker = UploadTracker::new(region_id);
999 let mut err = None;
1000 let puffin_key =
1001 IndexKey::new(region_id, index_file_id, FileType::Puffin(output.version));
1002 let index_id = RegionIndexId::new(RegionFileId::new(region_id, file_id), index_version);
1003 let puffin_path = RegionFilePathFactory::new(
1004 self.access_layer.table_dir().to_string(),
1005 self.access_layer.path_type(),
1006 )
1007 .build_index_file_path_with_version(index_id);
1008 if let Err(e) = write_cache
1009 .upload(puffin_key, &puffin_path, remote_store)
1010 .await
1011 {
1012 err = Some(e);
1013 }
1014 upload_tracker.push_uploaded_file(puffin_path);
1015 if let Some(err) = err {
1016 upload_tracker
1018 .clean(
1019 &smallvec![SstInfo {
1020 file_id,
1021 index_metadata: output,
1022 ..Default::default()
1023 }],
1024 &write_cache.file_cache(),
1025 remote_store,
1026 )
1027 .await;
1028 return Err(err);
1029 }
1030 } else {
1031 debug!("write cache is not available, skip uploading index file");
1032 }
1033 Ok(())
1034 }
1035
1036 async fn update_manifest(
1037 &self,
1038 output: IndexOutput,
1039 new_index_version: u64,
1040 ) -> Result<IndexPublication> {
1041 let mut updated = self.source.file_meta.clone();
1042 updated.available_indexes = output.build_available_indexes();
1043 updated.indexes = output.build_indexes();
1044 updated.index_file_size = output.file_size;
1045 updated.index_version = new_index_version;
1046 let publication = self
1047 .manifest_ctx
1048 .update_manifest_for_index(&self.source, updated)
1049 .await?;
1050 if let IndexPublication::Committed {
1051 manifest_version, ..
1052 } = &publication
1053 {
1054 info!(
1055 "Successfully update manifest version to {}, region: {}, reason: {}",
1056 manifest_version,
1057 self.region_id,
1058 self.reason.as_str()
1059 );
1060 }
1061 Ok(publication)
1062 }
1063}
1064
1065impl PartialEq for IndexBuildTask {
1066 fn eq(&self, other: &Self) -> bool {
1067 self.reason.priority() == other.reason.priority()
1068 }
1069}
1070
1071impl Eq for IndexBuildTask {}
1072
1073impl PartialOrd for IndexBuildTask {
1074 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1075 Some(self.cmp(other))
1076 }
1077}
1078
1079impl Ord for IndexBuildTask {
1080 fn cmp(&self, other: &Self) -> Ordering {
1081 self.reason.priority().cmp(&other.reason.priority())
1082 }
1083}
1084
1085#[derive(Clone)]
1086struct PendingIndexBuild {
1087 task: IndexBuildTask,
1088 version_control: VersionControlRef,
1089}
1090
1091impl PartialEq for PendingIndexBuild {
1092 fn eq(&self, other: &Self) -> bool {
1093 self.task == other.task
1094 }
1095}
1096
1097impl Eq for PendingIndexBuild {}
1098
1099impl PartialOrd for PendingIndexBuild {
1100 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1101 Some(self.cmp(other))
1102 }
1103}
1104
1105impl Ord for PendingIndexBuild {
1106 fn cmp(&self, other: &Self) -> Ordering {
1107 self.task.cmp(&other.task)
1108 }
1109}
1110
1111impl PendingIndexBuild {
1112 fn supersedes(&self, other: &Self) -> bool {
1115 match self
1116 .task
1117 .source
1118 .schema_version
1119 .cmp(&other.task.source.schema_version)
1120 {
1121 Ordering::Greater => true,
1122 Ordering::Less => false,
1123 Ordering::Equal => match self
1124 .task
1125 .source
1126 .file_meta
1127 .index_version()
1128 .cmp(&other.task.source.file_meta.index_version())
1129 {
1130 Ordering::Greater => true,
1131 Ordering::Less => false,
1132 Ordering::Equal => self.task.reason.priority() > other.task.reason.priority(),
1133 },
1134 }
1135 }
1136}
1137
1138#[derive(Default)]
1139struct PendingIndexBuilds {
1140 tasks: HashMap<FileId, PendingIndexBuild>,
1141}
1142
1143impl PendingIndexBuilds {
1144 fn insert(&mut self, pending: PendingIndexBuild) -> Option<IndexBuildTask> {
1145 let file_id = pending.task.source.file_meta.file_id;
1146 match self.tasks.entry(file_id) {
1147 std::collections::hash_map::Entry::Vacant(entry) => {
1148 entry.insert(pending);
1149 None
1150 }
1151 std::collections::hash_map::Entry::Occupied(mut entry)
1152 if pending.supersedes(entry.get()) =>
1153 {
1154 Some(entry.insert(pending).task)
1155 }
1156 std::collections::hash_map::Entry::Occupied(_) => Some(pending.task),
1157 }
1158 }
1159
1160 fn highest_ready(&self, building_files: &HashSet<FileId>) -> Option<PendingIndexBuild> {
1161 self.tasks
1162 .values()
1163 .filter(|pending| !building_files.contains(&pending.task.source.file_meta.file_id))
1164 .max()
1165 .cloned()
1166 }
1167
1168 fn remove(&mut self, file_id: FileId) {
1169 self.tasks.remove(&file_id);
1170 }
1171
1172 fn drain(&mut self) -> impl Iterator<Item = PendingIndexBuild> {
1173 std::mem::take(&mut self.tasks).into_values()
1174 }
1175
1176 #[cfg(test)]
1177 fn len(&self) -> usize {
1178 self.tasks.len()
1179 }
1180
1181 fn is_empty(&self) -> bool {
1182 self.tasks.is_empty()
1183 }
1184}
1185
1186struct IndexBuildStatus {
1188 building_files: HashSet<FileId>,
1189 pending_tasks: PendingIndexBuilds,
1191 retiring: bool,
1194}
1195
1196impl IndexBuildStatus {
1197 fn new() -> Self {
1198 IndexBuildStatus {
1199 building_files: HashSet::new(),
1200 pending_tasks: PendingIndexBuilds::default(),
1201 retiring: false,
1202 }
1203 }
1204
1205 async fn fail_pending(&mut self, err: Arc<Error>) {
1206 for pending in self.pending_tasks.drain() {
1207 pending.task.on_failure(err.clone()).await;
1208 }
1209 }
1210
1211 async fn retire(&mut self, err: Arc<Error>) {
1212 self.retiring = !self.building_files.is_empty();
1213 self.fail_pending(err).await;
1214 }
1215}
1216
1217pub struct IndexBuildScheduler {
1218 scheduler: SchedulerRef,
1220 region_status: HashMap<RegionId, IndexBuildStatus>,
1222 files_limit: usize,
1224}
1225
1226impl IndexBuildScheduler {
1228 pub fn new(scheduler: SchedulerRef, files_limit: usize) -> Self {
1229 IndexBuildScheduler {
1230 scheduler,
1231 region_status: HashMap::new(),
1232 files_limit,
1233 }
1234 }
1235
1236 pub(crate) async fn schedule_build(
1237 &mut self,
1238 version_control: &VersionControlRef,
1239 task: IndexBuildTask,
1240 ) -> Result<()> {
1241 let status = self
1242 .region_status
1243 .entry(task.region_id)
1244 .or_insert_with(IndexBuildStatus::new);
1245
1246 let file_id = task.source.file_meta.file_id;
1247 let region_file_id = RegionFileId::new(
1248 task.source.file_meta.region_id,
1249 task.source.file_meta.file_id,
1250 );
1251 let can_wait_for_active = status.retiring || task.reason == IndexBuildType::SchemaChange;
1252 let (rejected, coalesced) =
1253 if status.building_files.contains(&file_id) && !can_wait_for_active {
1254 (Some(task), None)
1255 } else {
1256 let pending = PendingIndexBuild {
1257 task,
1258 version_control: version_control.clone(),
1259 };
1260 (None, status.pending_tasks.insert(pending))
1261 };
1262 let should_schedule = !status.retiring;
1263
1264 if let Some(rejected) = rejected {
1265 debug!(
1266 "Rejecting index build because region file {:?} is already being built",
1267 region_file_id
1268 );
1269 rejected
1270 .on_success(IndexBuildOutcome::Aborted(format!(
1271 "Index is already being built for region file {:?}",
1272 region_file_id
1273 )))
1274 .await;
1275 rejected.listener.on_index_build_abort(region_file_id).await;
1276 }
1277
1278 if let Some(coalesced) = coalesced {
1279 debug!(
1280 "Coalescing redundant index build task for region file {:?}",
1281 region_file_id
1282 );
1283 coalesced
1284 .on_success(IndexBuildOutcome::Aborted(format!(
1285 "Index build was coalesced for region file {:?}",
1286 region_file_id
1287 )))
1288 .await;
1289 coalesced
1290 .listener
1291 .on_index_build_abort(region_file_id)
1292 .await;
1293 }
1294
1295 if should_schedule {
1296 self.schedule_next_build_batch();
1297 }
1298 Ok(())
1299 }
1300
1301 fn schedule_next_build_batch(&mut self) {
1303 let mut building_count = self
1304 .region_status
1305 .values()
1306 .map(|status| status.building_files.len())
1307 .sum::<usize>();
1308
1309 while building_count < self.files_limit {
1310 let Some(pending) = self.find_next_task() else {
1311 break;
1312 };
1313
1314 let task = pending.task;
1315 let region_id = task.region_id;
1316 let file_id = task.source.file_meta.file_id;
1317 let job = task.clone().into_index_build_job(pending.version_control);
1318 match self.scheduler.schedule(job) {
1319 Ok(()) => {
1320 if let Some(status) = self.region_status.get_mut(®ion_id) {
1321 status.pending_tasks.remove(file_id);
1322 status.building_files.insert(file_id);
1323 building_count += 1;
1324 } else {
1325 error!(
1326 "Region status not found when scheduling index build task, region: {}",
1327 region_id
1328 );
1329 }
1330 }
1331 Err(err) => {
1332 error!(
1333 err;
1334 "Failed to schedule index build job, region: {}, file_id: {}",
1335 region_id,
1336 file_id
1337 );
1338 if let Some(status) = self.region_status.get_mut(®ion_id) {
1339 status.pending_tasks.remove(file_id);
1340 }
1341 common_runtime::spawn_global(async move {
1342 task.on_failure(Arc::new(err)).await;
1343 });
1344 }
1345 }
1346 }
1347
1348 self.region_status.retain(|_, status| {
1349 !status.building_files.is_empty() || !status.pending_tasks.is_empty()
1350 });
1351 }
1352
1353 fn find_next_task(&self) -> Option<PendingIndexBuild> {
1355 self.region_status
1356 .values()
1357 .filter(|status| !status.retiring)
1358 .filter_map(|status| status.pending_tasks.highest_ready(&status.building_files))
1359 .max()
1360 }
1361
1362 pub(crate) fn on_task_stopped(&mut self, region_id: RegionId, file_id: FileId) {
1363 if let Some(status) = self.region_status.get_mut(®ion_id) {
1364 if !status.building_files.remove(&file_id) {
1365 debug!(
1366 "Index build task is not tracked as building, region: {}, file: {}",
1367 region_id, file_id
1368 );
1369 return;
1370 }
1371 if status.building_files.is_empty() && status.pending_tasks.is_empty() {
1372 self.region_status.remove(®ion_id);
1374 } else if status.building_files.is_empty() {
1375 status.retiring = false;
1377 }
1378 }
1379
1380 self.schedule_next_build_batch();
1381 }
1382
1383 pub(crate) async fn on_failure(&mut self, region_id: RegionId, err: Arc<Error>) {
1384 let Some(status) = self.region_status.get_mut(®ion_id) else {
1385 error!(err; "Index build failed after scheduler state was removed, region: {}", region_id);
1386 return;
1387 };
1388 if status.retiring {
1389 error!(err; "Index build from a retiring region incarnation failed, region: {}", region_id);
1390 return;
1391 }
1392 error!(
1393 err; "Index build scheduler encountered failure for region {}, failing pending tasks.",
1394 region_id
1395 );
1396 status.fail_pending(err).await;
1397 if status.building_files.is_empty() {
1398 self.region_status.remove(®ion_id);
1399 }
1400 }
1401
1402 pub(crate) async fn on_region_dropped(&mut self, region_id: RegionId) {
1404 self.retire_region(
1405 region_id,
1406 Arc::new(RegionDroppedSnafu { region_id }.build()),
1407 )
1408 .await;
1409 }
1410
1411 pub(crate) async fn on_region_closed(&mut self, region_id: RegionId) {
1413 self.retire_region(region_id, Arc::new(RegionClosedSnafu { region_id }.build()))
1414 .await;
1415 }
1416
1417 pub(crate) async fn on_region_truncated(&mut self, region_id: RegionId) {
1419 self.retire_region(
1420 region_id,
1421 Arc::new(RegionTruncatedSnafu { region_id }.build()),
1422 )
1423 .await;
1424 }
1425
1426 async fn retire_region(&mut self, region_id: RegionId, err: Arc<Error>) {
1427 let Some(status) = self.region_status.get_mut(®ion_id) else {
1428 return;
1429 };
1430 status.retire(err).await;
1431 if status.building_files.is_empty() {
1432 self.region_status.remove(®ion_id);
1433 }
1434 }
1435}
1436
1437pub(crate) fn decode_primary_keys_with_counts(
1440 batch: &RecordBatch,
1441 codec: &IndexValuesCodec,
1442) -> Result<Vec<(CompositeValues, usize)>> {
1443 let primary_key_index = primary_key_column_index(batch.num_columns());
1444 let pk_dict_array = batch
1445 .column(primary_key_index)
1446 .as_any()
1447 .downcast_ref::<PrimaryKeyArray>()
1448 .context(InvalidRecordBatchSnafu {
1449 reason: "Primary key column is not a dictionary array",
1450 })?;
1451 let pk_values_array = pk_dict_array
1452 .values()
1453 .as_any()
1454 .downcast_ref::<BinaryArray>()
1455 .context(InvalidRecordBatchSnafu {
1456 reason: "Primary key values are not binary array",
1457 })?;
1458 let keys = pk_dict_array.keys();
1459
1460 let mut result: Vec<(CompositeValues, usize)> = Vec::new();
1462 let mut prev_key: Option<u32> = None;
1463
1464 let pk_indices = keys.values();
1465 for ¤t_key in pk_indices.iter().take(keys.len()) {
1466 if let Some(prev) = prev_key
1468 && prev == current_key
1469 {
1470 result.last_mut().unwrap().1 += 1;
1472 continue;
1473 }
1474
1475 let pk_bytes = pk_values_array.value(current_key as usize);
1477 let decoded_value = codec.decoder().decode(pk_bytes).context(DecodeSnafu)?;
1478
1479 result.push((decoded_value, 1));
1480 prev_key = Some(current_key);
1481 }
1482
1483 Ok(result)
1484}
1485
1486#[cfg(test)]
1487mod tests {
1488 use std::sync::Arc;
1489
1490 use api::v1::SemanticType;
1491 use common_base::readable_size::ReadableSize;
1492 use datafusion_common::HashMap;
1493 use datatypes::data_type::ConcreteDataType;
1494 use datatypes::schema::{
1495 ColumnSchema, FulltextOptions, SkippingIndexOptions, SkippingIndexType,
1496 };
1497 use object_store::ObjectStore;
1498 use object_store::services::Memory;
1499 use puffin_manager::PuffinManagerFactory;
1500 use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder};
1501 use tokio::sync::mpsc;
1502
1503 use super::*;
1504 use crate::access_layer::{FilePathProvider, Metrics, SstWriteRequest, WriteType};
1505 use crate::cache::write_cache::WriteCache;
1506 use crate::config::{FulltextIndexConfig, IndexBuildMode, MitoConfig, Mode};
1507 use crate::manifest::action::{RegionEdit, RegionMetaAction, RegionMetaActionList};
1508 use crate::memtable::time_partition::TimePartitions;
1509 use crate::region::RegionLeaderState;
1510 use crate::region::version::{VersionBuilder, VersionControl};
1511 use crate::schedule::scheduler::{LocalScheduler, Scheduler};
1512 use crate::sst::file::{FileMeta, RegionFileId};
1513 use crate::sst::file_purger::NoopFilePurger;
1514 use crate::sst::location;
1515 use crate::sst::parquet::WriteOptions;
1516 use crate::test_util::memtable_util::EmptyMemtableBuilder;
1517 use crate::test_util::scheduler_util::{SchedulerEnv, VecScheduler};
1518 use crate::test_util::sst_util::{
1519 new_flat_source_from_record_batches, new_record_batch_by_range, sst_region_metadata,
1520 };
1521
1522 struct MetaConfig {
1523 with_inverted: bool,
1524 with_fulltext: bool,
1525 with_skipping_bloom: bool,
1526 #[cfg(feature = "vector_index")]
1527 with_vector: bool,
1528 }
1529
1530 async fn seed_manifest_file(manifest_ctx: &ManifestContextRef, file_meta: &FileMeta) {
1531 manifest_ctx
1532 .update_manifest(
1533 RegionLeaderState::Writable,
1534 RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
1535 files_to_add: vec![file_meta.clone()],
1536 files_to_remove: Vec::new(),
1537 timestamp_ms: None,
1538 flushed_sequence: None,
1539 flushed_entry_id: None,
1540 committed_sequence: None,
1541 compaction_time_window: None,
1542 })),
1543 false,
1544 )
1545 .await
1546 .unwrap();
1547 }
1548
1549 fn mock_region_metadata(
1550 MetaConfig {
1551 with_inverted,
1552 with_fulltext,
1553 with_skipping_bloom,
1554 #[cfg(feature = "vector_index")]
1555 with_vector,
1556 }: MetaConfig,
1557 ) -> RegionMetadataRef {
1558 let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 2));
1559 let mut column_schema = ColumnSchema::new("a", ConcreteDataType::int64_datatype(), false);
1560 if with_inverted {
1561 column_schema = column_schema.with_inverted_index(true);
1562 }
1563 builder
1564 .push_column_metadata(ColumnMetadata {
1565 column_schema,
1566 semantic_type: SemanticType::Field,
1567 column_id: 1,
1568 })
1569 .push_column_metadata(ColumnMetadata {
1570 column_schema: ColumnSchema::new("b", ConcreteDataType::float64_datatype(), false),
1571 semantic_type: SemanticType::Field,
1572 column_id: 2,
1573 })
1574 .push_column_metadata(ColumnMetadata {
1575 column_schema: ColumnSchema::new(
1576 "c",
1577 ConcreteDataType::timestamp_millisecond_datatype(),
1578 false,
1579 ),
1580 semantic_type: SemanticType::Timestamp,
1581 column_id: 3,
1582 });
1583
1584 if with_fulltext {
1585 let column_schema =
1586 ColumnSchema::new("text", ConcreteDataType::string_datatype(), true)
1587 .with_fulltext_options(FulltextOptions {
1588 enable: true,
1589 ..Default::default()
1590 })
1591 .unwrap();
1592
1593 let column = ColumnMetadata {
1594 column_schema,
1595 semantic_type: SemanticType::Field,
1596 column_id: 4,
1597 };
1598
1599 builder.push_column_metadata(column);
1600 }
1601
1602 if with_skipping_bloom {
1603 let column_schema =
1604 ColumnSchema::new("bloom", ConcreteDataType::string_datatype(), false)
1605 .with_skipping_options(SkippingIndexOptions::new_unchecked(
1606 42,
1607 0.01,
1608 SkippingIndexType::BloomFilter,
1609 ))
1610 .unwrap();
1611
1612 let column = ColumnMetadata {
1613 column_schema,
1614 semantic_type: SemanticType::Field,
1615 column_id: 5,
1616 };
1617
1618 builder.push_column_metadata(column);
1619 }
1620
1621 #[cfg(feature = "vector_index")]
1622 if with_vector {
1623 use index::vector::VectorIndexOptions;
1624
1625 let options = VectorIndexOptions::default();
1626 let column_schema =
1627 ColumnSchema::new("vec", ConcreteDataType::vector_datatype(4), true)
1628 .with_vector_index_options(&options)
1629 .unwrap();
1630 let column = ColumnMetadata {
1631 column_schema,
1632 semantic_type: SemanticType::Field,
1633 column_id: 6,
1634 };
1635
1636 builder.push_column_metadata(column);
1637 }
1638
1639 Arc::new(builder.build().unwrap())
1640 }
1641
1642 fn mock_object_store() -> ObjectStore {
1643 ObjectStore::new(Memory::default()).unwrap().finish()
1644 }
1645
1646 async fn mock_intm_mgr(path: impl AsRef<str>) -> IntermediateManager {
1647 IntermediateManager::init_fs(path).await.unwrap()
1648 }
1649
1650 fn random_region_file_id() -> RegionFileId {
1651 RegionFileId::new(RegionId::new(1, 2), FileId::random())
1652 }
1653
1654 struct NoopPathProvider;
1655
1656 impl FilePathProvider for NoopPathProvider {
1657 fn build_index_file_path(&self, _file_id: RegionFileId) -> String {
1658 unreachable!()
1659 }
1660
1661 fn build_index_file_path_with_version(&self, _index_id: RegionIndexId) -> String {
1662 unreachable!()
1663 }
1664
1665 fn build_sst_file_path(&self, _file_id: RegionFileId) -> String {
1666 unreachable!()
1667 }
1668 }
1669
1670 async fn mock_sst_file(
1671 metadata: RegionMetadataRef,
1672 env: &SchedulerEnv,
1673 build_mode: IndexBuildMode,
1674 ) -> SstInfo {
1675 let source = new_flat_source_from_record_batches(vec![
1676 new_record_batch_by_range(&["a", "d"], 0, 60),
1677 new_record_batch_by_range(&["b", "f"], 0, 40),
1678 new_record_batch_by_range(&["b", "h"], 100, 200),
1679 ]);
1680 let mut index_config = MitoConfig::default().index;
1681 index_config.build_mode = build_mode;
1682 let write_request = SstWriteRequest {
1683 op_type: OperationType::Flush,
1684 metadata: metadata.clone(),
1685 source,
1686 storage: None,
1687 max_sequence: None,
1688 sst_write_format: Default::default(),
1689 cache_manager: Default::default(),
1690 index_options: IndexOptions::default(),
1691 index_config,
1692 inverted_index_config: Default::default(),
1693 fulltext_index_config: Default::default(),
1694 bloom_filter_index_config: Default::default(),
1695 #[cfg(feature = "vector_index")]
1696 vector_index_config: Default::default(),
1697 };
1698 let mut metrics = Metrics::new(WriteType::Flush);
1699 env.access_layer
1700 .write_sst(write_request, &WriteOptions::default(), &mut metrics)
1701 .await
1702 .unwrap()
1703 .remove(0)
1704 }
1705
1706 async fn mock_version_control(
1707 metadata: RegionMetadataRef,
1708 file_purger: FilePurgerRef,
1709 files: HashMap<FileId, FileMeta>,
1710 ) -> VersionControlRef {
1711 let mutable = Arc::new(TimePartitions::new(
1712 metadata.clone(),
1713 Arc::new(EmptyMemtableBuilder::default()),
1714 0,
1715 None,
1716 ));
1717 let version_builder = VersionBuilder::new(metadata, mutable)
1718 .add_files(file_purger, files.values().cloned())
1719 .build();
1720 Arc::new(VersionControl::new(version_builder))
1721 }
1722
1723 async fn mock_indexer_builder(
1724 metadata: RegionMetadataRef,
1725 env: &SchedulerEnv,
1726 ) -> Arc<dyn IndexerBuilder + Send + Sync> {
1727 let (dir, factory) = PuffinManagerFactory::new_for_test_async("mock_indexer_builder").await;
1728 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1729 let puffin_manager = factory.build(
1730 env.access_layer.object_store().clone(),
1731 RegionFilePathFactory::new(
1732 env.access_layer.table_dir().to_string(),
1733 env.access_layer.path_type(),
1734 ),
1735 );
1736 Arc::new(IndexerBuilderImpl {
1737 build_type: IndexBuildType::Flush,
1738 metadata,
1739 puffin_manager,
1740 write_cache_enabled: false,
1741 intermediate_manager: intm_manager,
1742 index_options: IndexOptions::default(),
1743 inverted_index_config: InvertedIndexConfig::default(),
1744 fulltext_index_config: FulltextIndexConfig::default(),
1745 bloom_filter_index_config: BloomFilterConfig::default(),
1746 #[cfg(feature = "vector_index")]
1747 vector_index_config: Default::default(),
1748 })
1749 }
1750
1751 #[tokio::test]
1752 async fn test_build_indexer_basic() {
1753 let (dir, factory) =
1754 PuffinManagerFactory::new_for_test_async("test_build_indexer_basic_").await;
1755 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1756
1757 let metadata = mock_region_metadata(MetaConfig {
1758 with_inverted: true,
1759 with_fulltext: true,
1760 with_skipping_bloom: true,
1761 #[cfg(feature = "vector_index")]
1762 with_vector: false,
1763 });
1764 let indexer = IndexerBuilderImpl {
1765 build_type: IndexBuildType::Flush,
1766 metadata,
1767 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1768 write_cache_enabled: false,
1769 intermediate_manager: intm_manager,
1770 index_options: IndexOptions::default(),
1771 inverted_index_config: InvertedIndexConfig::default(),
1772 fulltext_index_config: FulltextIndexConfig::default(),
1773 bloom_filter_index_config: BloomFilterConfig::default(),
1774 #[cfg(feature = "vector_index")]
1775 vector_index_config: Default::default(),
1776 }
1777 .build(random_region_file_id(), 0, Some(1024))
1778 .await;
1779
1780 assert!(indexer.inverted_indexer.is_some());
1781 assert!(indexer.fulltext_indexer.is_some());
1782 assert!(indexer.bloom_filter_indexer.is_some());
1783 }
1784
1785 #[tokio::test]
1786 async fn test_build_indexer_disable_create() {
1787 let (dir, factory) =
1788 PuffinManagerFactory::new_for_test_async("test_build_indexer_disable_create_").await;
1789 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1790
1791 let metadata = mock_region_metadata(MetaConfig {
1792 with_inverted: true,
1793 with_fulltext: true,
1794 with_skipping_bloom: true,
1795 #[cfg(feature = "vector_index")]
1796 with_vector: false,
1797 });
1798 let indexer = IndexerBuilderImpl {
1799 build_type: IndexBuildType::Flush,
1800 metadata: metadata.clone(),
1801 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1802 write_cache_enabled: false,
1803 intermediate_manager: intm_manager.clone(),
1804 index_options: IndexOptions::default(),
1805 inverted_index_config: InvertedIndexConfig {
1806 create_on_flush: Mode::Disable,
1807 ..Default::default()
1808 },
1809 fulltext_index_config: FulltextIndexConfig::default(),
1810 bloom_filter_index_config: BloomFilterConfig::default(),
1811 #[cfg(feature = "vector_index")]
1812 vector_index_config: Default::default(),
1813 }
1814 .build(random_region_file_id(), 0, Some(1024))
1815 .await;
1816
1817 assert!(indexer.inverted_indexer.is_none());
1818 assert!(indexer.fulltext_indexer.is_some());
1819 assert!(indexer.bloom_filter_indexer.is_some());
1820
1821 let indexer = IndexerBuilderImpl {
1822 build_type: IndexBuildType::Compact,
1823 metadata: metadata.clone(),
1824 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1825 write_cache_enabled: false,
1826 intermediate_manager: intm_manager.clone(),
1827 index_options: IndexOptions::default(),
1828 inverted_index_config: InvertedIndexConfig::default(),
1829 fulltext_index_config: FulltextIndexConfig {
1830 create_on_compaction: Mode::Disable,
1831 ..Default::default()
1832 },
1833 bloom_filter_index_config: BloomFilterConfig::default(),
1834 #[cfg(feature = "vector_index")]
1835 vector_index_config: Default::default(),
1836 }
1837 .build(random_region_file_id(), 0, Some(1024))
1838 .await;
1839
1840 assert!(indexer.inverted_indexer.is_some());
1841 assert!(indexer.fulltext_indexer.is_none());
1842 assert!(indexer.bloom_filter_indexer.is_some());
1843
1844 let indexer = IndexerBuilderImpl {
1845 build_type: IndexBuildType::Compact,
1846 metadata,
1847 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1848 write_cache_enabled: false,
1849 intermediate_manager: intm_manager,
1850 index_options: IndexOptions::default(),
1851 inverted_index_config: InvertedIndexConfig::default(),
1852 fulltext_index_config: FulltextIndexConfig::default(),
1853 bloom_filter_index_config: BloomFilterConfig {
1854 create_on_compaction: Mode::Disable,
1855 ..Default::default()
1856 },
1857 #[cfg(feature = "vector_index")]
1858 vector_index_config: Default::default(),
1859 }
1860 .build(random_region_file_id(), 0, Some(1024))
1861 .await;
1862
1863 assert!(indexer.inverted_indexer.is_some());
1864 assert!(indexer.fulltext_indexer.is_some());
1865 assert!(indexer.bloom_filter_indexer.is_none());
1866 }
1867
1868 #[tokio::test]
1869 async fn test_build_indexer_no_required() {
1870 let (dir, factory) =
1871 PuffinManagerFactory::new_for_test_async("test_build_indexer_no_required_").await;
1872 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1873
1874 let metadata = mock_region_metadata(MetaConfig {
1875 with_inverted: false,
1876 with_fulltext: true,
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_none());
1898 assert!(indexer.fulltext_indexer.is_some());
1899 assert!(indexer.bloom_filter_indexer.is_some());
1900
1901 let metadata = mock_region_metadata(MetaConfig {
1902 with_inverted: true,
1903 with_fulltext: false,
1904 with_skipping_bloom: true,
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.clone(),
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_none());
1926 assert!(indexer.bloom_filter_indexer.is_some());
1927
1928 let metadata = mock_region_metadata(MetaConfig {
1929 with_inverted: true,
1930 with_fulltext: true,
1931 with_skipping_bloom: false,
1932 #[cfg(feature = "vector_index")]
1933 with_vector: false,
1934 });
1935 let indexer = IndexerBuilderImpl {
1936 build_type: IndexBuildType::Flush,
1937 metadata: metadata.clone(),
1938 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1939 write_cache_enabled: false,
1940 intermediate_manager: intm_manager,
1941 index_options: IndexOptions::default(),
1942 inverted_index_config: InvertedIndexConfig::default(),
1943 fulltext_index_config: FulltextIndexConfig::default(),
1944 bloom_filter_index_config: BloomFilterConfig::default(),
1945 #[cfg(feature = "vector_index")]
1946 vector_index_config: Default::default(),
1947 }
1948 .build(random_region_file_id(), 0, Some(1024))
1949 .await;
1950
1951 assert!(indexer.inverted_indexer.is_some());
1952 assert!(indexer.fulltext_indexer.is_some());
1953 assert!(indexer.bloom_filter_indexer.is_none());
1954 }
1955
1956 #[tokio::test]
1957 async fn test_build_indexer_zero_row_group_hint() {
1958 let (dir, factory) =
1959 PuffinManagerFactory::new_for_test_async("test_build_indexer_zero_row_group_").await;
1960 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
1961
1962 let metadata = mock_region_metadata(MetaConfig {
1963 with_inverted: true,
1964 with_fulltext: true,
1965 with_skipping_bloom: true,
1966 #[cfg(feature = "vector_index")]
1967 with_vector: false,
1968 });
1969 let indexer = IndexerBuilderImpl {
1970 build_type: IndexBuildType::Flush,
1971 metadata,
1972 puffin_manager: factory.build(mock_object_store(), NoopPathProvider),
1973 write_cache_enabled: false,
1974 intermediate_manager: intm_manager,
1975 index_options: IndexOptions::default(),
1976 inverted_index_config: InvertedIndexConfig::default(),
1977 fulltext_index_config: FulltextIndexConfig::default(),
1978 bloom_filter_index_config: BloomFilterConfig::default(),
1979 #[cfg(feature = "vector_index")]
1980 vector_index_config: Default::default(),
1981 }
1982 .build(random_region_file_id(), 0, Some(0))
1983 .await;
1984
1985 assert!(indexer.inverted_indexer.is_some());
1986 }
1987
1988 #[cfg(feature = "vector_index")]
1989 #[tokio::test]
1990 async fn test_update_flat_builds_vector_index() {
1991 use datatypes::arrow::array::BinaryBuilder;
1992 use datatypes::arrow::datatypes::{DataType, Field, Schema};
1993
1994 struct TestPathProvider;
1995
1996 impl FilePathProvider for TestPathProvider {
1997 fn build_index_file_path(&self, file_id: RegionFileId) -> String {
1998 format!("index/{}.puffin", file_id)
1999 }
2000
2001 fn build_index_file_path_with_version(&self, index_id: RegionIndexId) -> String {
2002 format!("index/{}.puffin", index_id)
2003 }
2004
2005 fn build_sst_file_path(&self, file_id: RegionFileId) -> String {
2006 format!("sst/{}.parquet", file_id)
2007 }
2008 }
2009
2010 fn f32s_to_bytes(values: &[f32]) -> Vec<u8> {
2011 let mut bytes = Vec::with_capacity(values.len() * 4);
2012 for v in values {
2013 bytes.extend_from_slice(&v.to_le_bytes());
2014 }
2015 bytes
2016 }
2017
2018 let (dir, factory) =
2019 PuffinManagerFactory::new_for_test_async("test_update_flat_builds_vector_index_").await;
2020 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
2021
2022 let metadata = mock_region_metadata(MetaConfig {
2023 with_inverted: false,
2024 with_fulltext: false,
2025 with_skipping_bloom: false,
2026 with_vector: true,
2027 });
2028
2029 let mut indexer = IndexerBuilderImpl {
2030 build_type: IndexBuildType::Flush,
2031 metadata,
2032 puffin_manager: factory.build(mock_object_store(), TestPathProvider),
2033 write_cache_enabled: false,
2034 intermediate_manager: intm_manager,
2035 index_options: IndexOptions::default(),
2036 inverted_index_config: InvertedIndexConfig::default(),
2037 fulltext_index_config: FulltextIndexConfig::default(),
2038 bloom_filter_index_config: BloomFilterConfig::default(),
2039 vector_index_config: Default::default(),
2040 }
2041 .build(random_region_file_id(), 0, Some(1024))
2042 .await;
2043
2044 assert!(indexer.vector_indexer.is_some());
2045
2046 let vec1 = f32s_to_bytes(&[1.0, 0.0, 0.0, 0.0]);
2047 let vec2 = f32s_to_bytes(&[0.0, 1.0, 0.0, 0.0]);
2048
2049 let mut builder = BinaryBuilder::with_capacity(2, vec1.len() + vec2.len());
2050 builder.append_value(&vec1);
2051 builder.append_value(&vec2);
2052
2053 let schema = Arc::new(Schema::new(vec![Field::new("vec", DataType::Binary, true)]));
2054 let batch = RecordBatch::try_new(schema, vec![Arc::new(builder.finish())]).unwrap();
2055
2056 indexer.update_flat(&batch).await;
2057 let output = indexer.finish().await;
2058
2059 assert!(output.vector_index.is_available());
2060 assert!(output.vector_index.columns.contains(&6));
2061 }
2062
2063 #[tokio::test]
2064 async fn test_index_build_task_sst_not_exist() {
2065 let env = SchedulerEnv::new().await;
2066 let (tx, _rx) = mpsc::channel(4);
2067 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2068 let mut scheduler = env.mock_index_build_scheduler(4);
2069 let metadata = Arc::new(sst_region_metadata());
2070 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2071 let file_purger = Arc::new(NoopFilePurger {});
2072 let files = HashMap::new();
2073 let version_control =
2074 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2075 let region_id = metadata.region_id;
2076 let indexer_builder = mock_indexer_builder(metadata, &env).await;
2077
2078 let file_meta = FileMeta {
2079 region_id,
2080 file_id: FileId::random(),
2081 file_size: 100,
2082 ..Default::default()
2083 };
2084
2085 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2086
2087 let task = IndexBuildTask {
2089 region_id,
2090 file,
2091 source: IndexBuildSource::new(
2092 file_meta,
2093 version_control.current().version.metadata.schema_version,
2094 ),
2095 reason: IndexBuildType::Flush,
2096 access_layer: env.access_layer.clone(),
2097 listener: WorkerListener::default(),
2098 manifest_ctx,
2099 write_cache: None,
2100 cache_manager: None,
2101 file_purger,
2102 indexer_builder,
2103 request_sender: tx,
2104 result_sender: result_tx,
2105 };
2106
2107 scheduler
2109 .schedule_build(&version_control, task)
2110 .await
2111 .unwrap();
2112 match result_rx.recv().await.unwrap() {
2113 Ok(outcome) => {
2114 if outcome == IndexBuildOutcome::Finished {
2115 panic!("Expect aborted result due to missing SST file")
2116 }
2117 }
2118 _ => panic!("Expect aborted result due to missing SST file"),
2119 }
2120 }
2121
2122 #[tokio::test]
2123 async fn test_index_build_task_increments_legacy_index_version() {
2124 let env = SchedulerEnv::new().await;
2125 let mut scheduler = env.mock_index_build_scheduler(4);
2126 let metadata = Arc::new(sst_region_metadata());
2127 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2128 let region_id = metadata.region_id;
2129 let file_purger = Arc::new(NoopFilePurger {});
2130 let sst_info = mock_sst_file(metadata.clone(), &env, IndexBuildMode::Async).await;
2131 let file_meta = FileMeta {
2132 region_id,
2133 file_id: sst_info.file_id,
2134 file_size: sst_info.file_size,
2135 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
2136 available_indexes: smallvec![IndexType::InvertedIndex],
2137 index_file_size: 0,
2139 index_version: 0,
2140 num_rows: sst_info.num_rows as u64,
2141 num_row_groups: sst_info.num_row_groups,
2142 ..Default::default()
2143 };
2144 seed_manifest_file(&manifest_ctx, &file_meta).await;
2145 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2146 let version_control =
2147 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2148 let indexer_builder = mock_indexer_builder(metadata.clone(), &env).await;
2149
2150 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2151
2152 let (tx, mut rx) = mpsc::channel(4);
2154 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2155 let task = IndexBuildTask {
2156 region_id,
2157 file,
2158 source: IndexBuildSource::new(
2159 file_meta.clone(),
2160 version_control.current().version.metadata.schema_version,
2161 ),
2162 reason: IndexBuildType::Flush,
2163 access_layer: env.access_layer.clone(),
2164 listener: WorkerListener::default(),
2165 manifest_ctx,
2166 write_cache: None,
2167 cache_manager: None,
2168 file_purger,
2169 indexer_builder,
2170 request_sender: tx,
2171 result_sender: result_tx,
2172 };
2173
2174 scheduler
2175 .schedule_build(&version_control, task)
2176 .await
2177 .unwrap();
2178
2179 match result_rx.recv().await.unwrap() {
2181 Ok(outcome) => {
2182 assert_eq!(outcome, IndexBuildOutcome::Finished);
2183 }
2184 _ => panic!("Expect finished result"),
2185 }
2186
2187 let worker_req = rx.recv().await.unwrap().request;
2189 match worker_req {
2190 WorkerRequest::Background {
2191 region_id: req_region_id,
2192 notify: BackgroundNotify::IndexBuildFinished(finished),
2193 } => {
2194 assert_eq!(req_region_id, region_id);
2195 let updated_meta = &finished.file_meta;
2196
2197 assert!(!updated_meta.available_indexes.is_empty());
2199 assert!(updated_meta.index_file_size > 0);
2200 assert_eq!(updated_meta.file_id, file_meta.file_id);
2201 assert_eq!(updated_meta.index_version, 1);
2202 }
2203 _ => panic!("Unexpected worker request: {:?}", worker_req),
2204 }
2205 }
2206
2207 async fn schedule_index_build_task_with_mode(build_mode: IndexBuildMode) {
2208 let env = SchedulerEnv::new().await;
2209 let mut scheduler = env.mock_index_build_scheduler(4);
2210 let metadata = Arc::new(sst_region_metadata());
2211 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2212 let file_purger = Arc::new(NoopFilePurger {});
2213 let region_id = metadata.region_id;
2214 let sst_info = mock_sst_file(metadata.clone(), &env, build_mode.clone()).await;
2215 let file_meta = FileMeta {
2216 region_id,
2217 file_id: sst_info.file_id,
2218 file_size: sst_info.file_size,
2219 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
2220 index_file_size: sst_info.index_metadata.file_size,
2221 num_rows: sst_info.num_rows as u64,
2222 num_row_groups: sst_info.num_row_groups,
2223 ..Default::default()
2224 };
2225 seed_manifest_file(&manifest_ctx, &file_meta).await;
2226 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2227 let version_control =
2228 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2229 let indexer_builder = mock_indexer_builder(metadata.clone(), &env).await;
2230
2231 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2232
2233 let (tx, _rx) = mpsc::channel(4);
2235 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2236 let task = IndexBuildTask {
2237 region_id,
2238 file,
2239 source: IndexBuildSource::new(
2240 file_meta.clone(),
2241 version_control.current().version.metadata.schema_version,
2242 ),
2243 reason: IndexBuildType::Flush,
2244 access_layer: env.access_layer.clone(),
2245 listener: WorkerListener::default(),
2246 manifest_ctx,
2247 write_cache: None,
2248 cache_manager: None,
2249 file_purger,
2250 indexer_builder,
2251 request_sender: tx,
2252 result_sender: result_tx,
2253 };
2254
2255 scheduler
2256 .schedule_build(&version_control, task)
2257 .await
2258 .unwrap();
2259
2260 let puffin_path = location::index_file_path(
2261 env.access_layer.table_dir(),
2262 RegionIndexId::new(RegionFileId::new(region_id, file_meta.file_id), 0),
2263 env.access_layer.path_type(),
2264 );
2265
2266 if build_mode == IndexBuildMode::Async {
2267 assert!(
2269 !env.access_layer
2270 .object_store()
2271 .exists(&puffin_path)
2272 .await
2273 .unwrap()
2274 );
2275 } else {
2276 assert!(
2278 env.access_layer
2279 .object_store()
2280 .exists(&puffin_path)
2281 .await
2282 .unwrap()
2283 );
2284 }
2285
2286 match result_rx.recv().await.unwrap() {
2288 Ok(outcome) => {
2289 assert_eq!(outcome, IndexBuildOutcome::Finished);
2290 }
2291 _ => panic!("Expect finished result"),
2292 }
2293
2294 assert!(
2296 env.access_layer
2297 .object_store()
2298 .exists(&puffin_path)
2299 .await
2300 .unwrap()
2301 );
2302 }
2303
2304 #[tokio::test]
2305 async fn test_index_build_task_build_mode() {
2306 schedule_index_build_task_with_mode(IndexBuildMode::Async).await;
2307 schedule_index_build_task_with_mode(IndexBuildMode::Sync).await;
2308 }
2309
2310 #[tokio::test]
2311 async fn test_index_build_task_no_index() {
2312 let env = SchedulerEnv::new().await;
2313 let mut scheduler = env.mock_index_build_scheduler(4);
2314 let mut metadata = sst_region_metadata();
2315 metadata.column_metadatas.iter_mut().for_each(|col| {
2317 col.column_schema.set_inverted_index(false);
2318 let _ = col.column_schema.unset_skipping_options();
2319 });
2320 let region_id = metadata.region_id;
2321 let metadata = Arc::new(metadata);
2322 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2323 let file_purger = Arc::new(NoopFilePurger {});
2324 let sst_info = mock_sst_file(metadata.clone(), &env, IndexBuildMode::Async).await;
2325 let file_meta = FileMeta {
2326 region_id,
2327 file_id: sst_info.file_id,
2328 file_size: sst_info.file_size,
2329 max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
2330 index_file_size: sst_info.index_metadata.file_size,
2331 num_rows: sst_info.num_rows as u64,
2332 num_row_groups: sst_info.num_row_groups,
2333 ..Default::default()
2334 };
2335 seed_manifest_file(&manifest_ctx, &file_meta).await;
2336 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2337 let version_control =
2338 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2339 let indexer_builder = mock_indexer_builder(metadata.clone(), &env).await;
2340
2341 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2342
2343 let (tx, mut rx) = mpsc::channel(4);
2345 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2346 let task = IndexBuildTask {
2347 region_id,
2348 file,
2349 source: IndexBuildSource::new(
2350 file_meta.clone(),
2351 version_control.current().version.metadata.schema_version,
2352 ),
2353 reason: IndexBuildType::Flush,
2354 access_layer: env.access_layer.clone(),
2355 listener: WorkerListener::default(),
2356 manifest_ctx,
2357 write_cache: None,
2358 cache_manager: None,
2359 file_purger,
2360 indexer_builder,
2361 request_sender: tx,
2362 result_sender: result_tx,
2363 };
2364
2365 scheduler
2366 .schedule_build(&version_control, task)
2367 .await
2368 .unwrap();
2369
2370 match result_rx.recv().await.unwrap() {
2372 Ok(outcome) => {
2373 assert_eq!(outcome, IndexBuildOutcome::Finished);
2374 }
2375 _ => panic!("Expect finished result"),
2376 }
2377
2378 let _ = rx.recv().await.is_none();
2380 }
2381
2382 #[tokio::test]
2383 async fn test_index_build_task_with_write_cache() {
2384 let env = SchedulerEnv::new().await;
2385 let mut scheduler = env.mock_index_build_scheduler(4);
2386 let metadata = Arc::new(sst_region_metadata());
2387 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2388 let file_purger = Arc::new(NoopFilePurger {});
2389 let region_id = metadata.region_id;
2390
2391 let (dir, factory) = PuffinManagerFactory::new_for_test_async("test_write_cache").await;
2392 let intm_manager = mock_intm_mgr(dir.path().to_string_lossy()).await;
2393
2394 let write_cache = Arc::new(
2396 WriteCache::new_fs(
2397 dir.path().to_str().unwrap(),
2398 ReadableSize::mb(10),
2399 None,
2400 None,
2401 true, factory,
2403 intm_manager,
2404 ReadableSize::mb(10),
2405 )
2406 .await
2407 .unwrap(),
2408 );
2409 let indexer_builder = Arc::new(IndexerBuilderImpl {
2411 build_type: IndexBuildType::Flush,
2412 metadata: metadata.clone(),
2413 puffin_manager: write_cache.build_puffin_manager().clone(),
2414 write_cache_enabled: true,
2415 intermediate_manager: write_cache.intermediate_manager().clone(),
2416 index_options: IndexOptions::default(),
2417 inverted_index_config: InvertedIndexConfig::default(),
2418 fulltext_index_config: FulltextIndexConfig::default(),
2419 bloom_filter_index_config: BloomFilterConfig::default(),
2420 #[cfg(feature = "vector_index")]
2421 vector_index_config: Default::default(),
2422 });
2423
2424 let sst_info = mock_sst_file(metadata.clone(), &env, IndexBuildMode::Async).await;
2425 let file_meta = FileMeta {
2426 region_id,
2427 file_id: sst_info.file_id,
2428 file_size: sst_info.file_size,
2429 index_file_size: sst_info.index_metadata.file_size,
2430 num_rows: sst_info.num_rows as u64,
2431 num_row_groups: sst_info.num_row_groups,
2432 ..Default::default()
2433 };
2434 seed_manifest_file(&manifest_ctx, &file_meta).await;
2435 let files = HashMap::from([(file_meta.file_id, file_meta.clone())]);
2436 let version_control =
2437 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2438
2439 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2440
2441 let (tx, mut _rx) = mpsc::channel(4);
2443 let (result_tx, mut result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2444 let task = IndexBuildTask {
2445 region_id,
2446 file,
2447 source: IndexBuildSource::new(
2448 file_meta.clone(),
2449 version_control.current().version.metadata.schema_version,
2450 ),
2451 reason: IndexBuildType::Flush,
2452 access_layer: env.access_layer.clone(),
2453 listener: WorkerListener::default(),
2454 manifest_ctx,
2455 write_cache: Some(write_cache.clone()),
2456 cache_manager: None,
2457 file_purger,
2458 indexer_builder,
2459 request_sender: tx,
2460 result_sender: result_tx,
2461 };
2462
2463 scheduler
2464 .schedule_build(&version_control, task)
2465 .await
2466 .unwrap();
2467
2468 match result_rx.recv().await.unwrap() {
2470 Ok(outcome) => {
2471 assert_eq!(outcome, IndexBuildOutcome::Finished);
2472 }
2473 _ => panic!("Expect finished result"),
2474 }
2475
2476 let index_key = IndexKey::new(
2478 region_id,
2479 file_meta.file_id,
2480 FileType::Puffin(sst_info.index_metadata.version),
2481 );
2482 assert!(write_cache.file_cache().contains_key(&index_key));
2483 }
2484
2485 async fn create_mock_task_for_schedule(
2486 env: &SchedulerEnv,
2487 file_id: FileId,
2488 region_id: RegionId,
2489 reason: IndexBuildType,
2490 ) -> IndexBuildTask {
2491 create_mock_task_for_schedule_with_result(env, file_id, region_id, reason)
2492 .await
2493 .0
2494 }
2495
2496 async fn create_mock_task_for_schedule_with_result(
2499 env: &SchedulerEnv,
2500 file_id: FileId,
2501 region_id: RegionId,
2502 reason: IndexBuildType,
2503 ) -> (IndexBuildTask, mpsc::Receiver<Result<IndexBuildOutcome>>) {
2504 let metadata = Arc::new(sst_region_metadata());
2505 let schema_version = metadata.schema_version;
2506 let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
2507 let file_purger = Arc::new(NoopFilePurger {});
2508 let indexer_builder = mock_indexer_builder(metadata, env).await;
2509 let (tx, _rx) = mpsc::channel(4);
2510 let (result_tx, result_rx) = mpsc::channel::<Result<IndexBuildOutcome>>(4);
2511
2512 let file_meta = FileMeta {
2513 region_id,
2514 file_id,
2515 file_size: 100,
2516 ..Default::default()
2517 };
2518
2519 let file = FileHandle::new(file_meta.clone(), file_purger.clone());
2520
2521 let task = IndexBuildTask {
2522 region_id,
2523 file,
2524 source: IndexBuildSource::new(file_meta, schema_version),
2525 reason,
2526 access_layer: env.access_layer.clone(),
2527 listener: WorkerListener::default(),
2528 manifest_ctx,
2529 write_cache: None,
2530 cache_manager: None,
2531 file_purger,
2532 indexer_builder,
2533 request_sender: tx,
2534 result_sender: result_tx,
2535 };
2536 (task, result_rx)
2537 }
2538
2539 #[tokio::test]
2540 async fn test_scheduler_coalesces_latest_schema_generation_per_sst() {
2541 let job_scheduler = Arc::new(VecScheduler::default());
2542 let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2543 let mut scheduler = env.mock_index_build_scheduler(2);
2544 let metadata = Arc::new(sst_region_metadata());
2545 let region_id = metadata.region_id;
2546 let file_id = FileId::random();
2547 let file_purger = Arc::new(NoopFilePurger {});
2548 let files = HashMap::from([(
2549 file_id,
2550 FileMeta {
2551 region_id,
2552 file_id,
2553 file_size: 100,
2554 ..Default::default()
2555 },
2556 )]);
2557 let version_control = mock_version_control(metadata, file_purger, files).await;
2558
2559 let (active, _active_rx) = create_mock_task_for_schedule_with_result(
2560 &env,
2561 file_id,
2562 region_id,
2563 IndexBuildType::Manual,
2564 )
2565 .await;
2566 scheduler
2567 .schedule_build(&version_control, active)
2568 .await
2569 .unwrap();
2570 assert_eq!(job_scheduler.num_jobs(), 1);
2571
2572 let (mut older, mut older_rx) = create_mock_task_for_schedule_with_result(
2573 &env,
2574 file_id,
2575 region_id,
2576 IndexBuildType::SchemaChange,
2577 )
2578 .await;
2579 older.source.schema_version = 1;
2580 scheduler
2581 .schedule_build(&version_control, older)
2582 .await
2583 .unwrap();
2584
2585 let (mut latest, mut latest_rx) = create_mock_task_for_schedule_with_result(
2586 &env,
2587 file_id,
2588 region_id,
2589 IndexBuildType::SchemaChange,
2590 )
2591 .await;
2592 latest.source.schema_version = 2;
2593 scheduler
2594 .schedule_build(&version_control, latest)
2595 .await
2596 .unwrap();
2597
2598 let replaced = tokio::time::timeout(std::time::Duration::from_secs(5), older_rx.recv())
2599 .await
2600 .expect("replaced pending task result sender was not completed")
2601 .expect("replaced pending task result channel closed");
2602 assert!(matches!(
2603 replaced,
2604 Ok(IndexBuildOutcome::Aborted(reason)) if reason.contains("coalesced")
2605 ));
2606 assert!(matches!(
2607 latest_rx.try_recv(),
2608 Err(mpsc::error::TryRecvError::Empty)
2609 ));
2610
2611 let status = &scheduler.region_status[®ion_id];
2612 assert_eq!(status.building_files.len(), 1);
2613 assert_eq!(status.pending_tasks.len(), 1);
2614 assert_eq!(job_scheduler.num_jobs(), 1);
2615
2616 scheduler.on_task_stopped(region_id, file_id);
2617 let status = &scheduler.region_status[®ion_id];
2618 assert_eq!(status.building_files.len(), 1);
2619 assert!(status.pending_tasks.is_empty());
2620 assert_eq!(job_scheduler.num_jobs(), 2);
2621 }
2622
2623 #[tokio::test]
2624 async fn test_scheduler_completes_sender_when_job_is_rejected() {
2625 let job_scheduler = Arc::new(LocalScheduler::new(1));
2626 job_scheduler.stop(false).await.unwrap();
2627 let env = SchedulerEnv::new().await.scheduler(job_scheduler);
2628 let mut scheduler = env.mock_index_build_scheduler(1);
2629 let metadata = Arc::new(sst_region_metadata());
2630 let region_id = metadata.region_id;
2631 let file_id = FileId::random();
2632 let file_purger = Arc::new(NoopFilePurger {});
2633 let files = HashMap::from([(
2634 file_id,
2635 FileMeta {
2636 region_id,
2637 file_id,
2638 file_size: 100,
2639 ..Default::default()
2640 },
2641 )]);
2642 let version_control = mock_version_control(metadata, file_purger, files).await;
2643 let (task, mut result_rx) = create_mock_task_for_schedule_with_result(
2644 &env,
2645 file_id,
2646 region_id,
2647 IndexBuildType::Flush,
2648 )
2649 .await;
2650
2651 scheduler
2652 .schedule_build(&version_control, task)
2653 .await
2654 .unwrap();
2655
2656 let result = tokio::time::timeout(std::time::Duration::from_secs(5), result_rx.recv())
2657 .await
2658 .expect("scheduler rejection did not complete the result sender")
2659 .expect("result channel closed without a result");
2660 assert!(result.is_err());
2661 assert!(!scheduler.region_status.contains_key(®ion_id));
2662 }
2663
2664 #[tokio::test]
2665 async fn test_scheduler_comprehensive() {
2666 let env = SchedulerEnv::new().await;
2667 let mut scheduler = env.mock_index_build_scheduler(2);
2668 let metadata = Arc::new(sst_region_metadata());
2669 let region_id = metadata.region_id;
2670 let file_purger = Arc::new(NoopFilePurger {});
2671
2672 let file_id1 = FileId::random();
2674 let file_id2 = FileId::random();
2675 let file_id3 = FileId::random();
2676 let file_id4 = FileId::random();
2677 let file_id5 = FileId::random();
2678
2679 let mut files = HashMap::new();
2680 for file_id in [file_id1, file_id2, file_id3, file_id4, file_id5] {
2681 files.insert(
2682 file_id,
2683 FileMeta {
2684 region_id,
2685 file_id,
2686 file_size: 100,
2687 ..Default::default()
2688 },
2689 );
2690 }
2691
2692 let version_control = mock_version_control(metadata, file_purger, files).await;
2693
2694 let task1 =
2696 create_mock_task_for_schedule(&env, file_id1, region_id, IndexBuildType::Flush).await;
2697 assert!(
2698 scheduler
2699 .schedule_build(&version_control, task1)
2700 .await
2701 .is_ok()
2702 );
2703 assert!(scheduler.region_status.contains_key(®ion_id));
2704 let status = scheduler.region_status.get(®ion_id).unwrap();
2705 assert_eq!(status.building_files.len(), 1);
2706 assert!(status.building_files.contains(&file_id1));
2707
2708 let task1_dup =
2710 create_mock_task_for_schedule(&env, file_id1, region_id, IndexBuildType::Flush).await;
2711 scheduler
2712 .schedule_build(&version_control, task1_dup)
2713 .await
2714 .unwrap();
2715 let status = scheduler.region_status.get(®ion_id).unwrap();
2716 assert_eq!(status.building_files.len(), 1); let task2 =
2720 create_mock_task_for_schedule(&env, file_id2, region_id, IndexBuildType::Flush).await;
2721 scheduler
2722 .schedule_build(&version_control, task2)
2723 .await
2724 .unwrap();
2725 let status = scheduler.region_status.get(®ion_id).unwrap();
2726 assert_eq!(status.building_files.len(), 2); assert_eq!(status.pending_tasks.len(), 0);
2728
2729 let task3 =
2732 create_mock_task_for_schedule(&env, file_id3, region_id, IndexBuildType::Compact).await;
2733 let task4 =
2734 create_mock_task_for_schedule(&env, file_id4, region_id, IndexBuildType::SchemaChange)
2735 .await;
2736 let task5 =
2737 create_mock_task_for_schedule(&env, file_id5, region_id, IndexBuildType::Manual).await;
2738
2739 scheduler
2740 .schedule_build(&version_control, task3)
2741 .await
2742 .unwrap();
2743 scheduler
2744 .schedule_build(&version_control, task4)
2745 .await
2746 .unwrap();
2747 scheduler
2748 .schedule_build(&version_control, task5)
2749 .await
2750 .unwrap();
2751
2752 let status = scheduler.region_status.get(®ion_id).unwrap();
2753 assert_eq!(status.building_files.len(), 2); assert_eq!(status.pending_tasks.len(), 3); scheduler.on_task_stopped(region_id, file_id1);
2758 let status = scheduler.region_status.get(®ion_id).unwrap();
2759 assert!(!status.building_files.contains(&file_id1));
2760 assert_eq!(status.building_files.len(), 2); assert_eq!(status.pending_tasks.len(), 2); assert!(status.building_files.contains(&file_id5));
2764
2765 scheduler.on_task_stopped(region_id, file_id2);
2767 let status = scheduler.region_status.get(®ion_id).unwrap();
2768 assert_eq!(status.building_files.len(), 2);
2769 assert_eq!(status.pending_tasks.len(), 1); assert!(status.building_files.contains(&file_id4)); scheduler.on_task_stopped(region_id, file_id5);
2774 scheduler.on_task_stopped(region_id, file_id4);
2775
2776 let status = scheduler.region_status.get(®ion_id).unwrap();
2777 assert_eq!(status.building_files.len(), 1); assert_eq!(status.pending_tasks.len(), 0);
2779 assert!(status.building_files.contains(&file_id3));
2780
2781 scheduler.on_task_stopped(region_id, file_id3);
2782
2783 assert!(!scheduler.region_status.contains_key(®ion_id));
2785
2786 let task6 =
2788 create_mock_task_for_schedule(&env, file_id1, region_id, IndexBuildType::Flush).await;
2789 let task7 =
2790 create_mock_task_for_schedule(&env, file_id2, region_id, IndexBuildType::Flush).await;
2791 let task8 =
2792 create_mock_task_for_schedule(&env, file_id3, region_id, IndexBuildType::Manual).await;
2793
2794 scheduler
2795 .schedule_build(&version_control, task6)
2796 .await
2797 .unwrap();
2798 scheduler
2799 .schedule_build(&version_control, task7)
2800 .await
2801 .unwrap();
2802 scheduler
2803 .schedule_build(&version_control, task8)
2804 .await
2805 .unwrap();
2806
2807 assert!(scheduler.region_status.contains_key(®ion_id));
2808 let status = scheduler.region_status.get(®ion_id).unwrap();
2809 assert_eq!(status.building_files.len(), 2);
2810 assert_eq!(status.pending_tasks.len(), 1);
2811
2812 scheduler
2813 .on_failure(
2814 region_id,
2815 Arc::new(
2816 crate::error::UnexpectedSnafu {
2817 reason: "index build failed".to_string(),
2818 }
2819 .build(),
2820 ),
2821 )
2822 .await;
2823 let status = scheduler.region_status.get(®ion_id).unwrap();
2824 assert!(!status.retiring);
2825 assert_eq!(status.building_files.len(), 2);
2826 assert!(status.pending_tasks.is_empty());
2827
2828 scheduler.on_task_stopped(region_id, file_id1);
2829 assert!(scheduler.region_status.contains_key(®ion_id));
2830 scheduler.on_task_stopped(region_id, file_id2);
2831 assert!(!scheduler.region_status.contains_key(®ion_id));
2832 }
2833
2834 async fn setup_scheduler_with_pending_tasks(
2838 env: &SchedulerEnv,
2839 ) -> (
2840 IndexBuildScheduler,
2841 mpsc::Receiver<Result<IndexBuildOutcome>>,
2842 mpsc::Receiver<Result<IndexBuildOutcome>>,
2843 VersionControlRef,
2844 RegionId,
2845 FileId, ) {
2847 let metadata = Arc::new(sst_region_metadata());
2848 let region_id = metadata.region_id;
2849 let file_purger = Arc::new(NoopFilePurger {});
2850
2851 let file_id1 = FileId::random();
2852 let file_id2 = FileId::random();
2853 let file_id3 = FileId::random();
2854
2855 let files = HashMap::from([
2856 (
2857 file_id1,
2858 FileMeta {
2859 region_id,
2860 file_id: file_id1,
2861 file_size: 100,
2862 ..Default::default()
2863 },
2864 ),
2865 (
2866 file_id2,
2867 FileMeta {
2868 region_id,
2869 file_id: file_id2,
2870 file_size: 100,
2871 ..Default::default()
2872 },
2873 ),
2874 (
2875 file_id3,
2876 FileMeta {
2877 region_id,
2878 file_id: file_id3,
2879 file_size: 100,
2880 ..Default::default()
2881 },
2882 ),
2883 ]);
2884 let version_control =
2885 mock_version_control(metadata.clone(), file_purger.clone(), files).await;
2886
2887 let mut scheduler = env.mock_index_build_scheduler(1);
2888
2889 let (task1, _rx1) = create_mock_task_for_schedule_with_result(
2893 env,
2894 file_id1,
2895 region_id,
2896 IndexBuildType::Flush,
2897 )
2898 .await;
2899 let (task2, rx2) = create_mock_task_for_schedule_with_result(
2900 env,
2901 file_id2,
2902 region_id,
2903 IndexBuildType::Flush,
2904 )
2905 .await;
2906 let (task3, rx3) = create_mock_task_for_schedule_with_result(
2907 env,
2908 file_id3,
2909 region_id,
2910 IndexBuildType::Flush,
2911 )
2912 .await;
2913
2914 scheduler
2915 .schedule_build(&version_control, task1)
2916 .await
2917 .unwrap();
2918 scheduler
2919 .schedule_build(&version_control, task2)
2920 .await
2921 .unwrap();
2922 scheduler
2923 .schedule_build(&version_control, task3)
2924 .await
2925 .unwrap();
2926
2927 assert!(scheduler.region_status.contains_key(®ion_id));
2929 let status = scheduler.region_status.get(®ion_id).unwrap();
2930 assert_eq!(status.building_files.len(), 1);
2931 assert_eq!(status.pending_tasks.len(), 2);
2932
2933 (scheduler, rx2, rx3, version_control, region_id, file_id1)
2934 }
2935
2936 fn assert_lifecycle_error(
2941 err: crate::error::Error,
2942 expected_source: &str,
2943 lifecycle_name: &str,
2944 ) {
2945 let crate::error::Error::BuildIndexAsync { source, .. } = err else {
2946 panic!("[{lifecycle_name}] Expected BuildIndexAsync outer error, got: {err:?}");
2947 };
2948 let actual = match &*source {
2949 crate::error::Error::RegionDropped { .. } => "dropped",
2950 crate::error::Error::RegionClosed { .. } => "closed",
2951 crate::error::Error::RegionTruncated { .. } => "truncated",
2952 other => panic!(
2953 "[{lifecycle_name}] Expected lifecycle source variant (RegionDropped/RegionClosed/RegionTruncated), got: {other:?}"
2954 ),
2955 };
2956 assert_eq!(
2957 actual, expected_source,
2958 "[{lifecycle_name}] Source variant mismatch"
2959 );
2960 }
2961
2962 async fn recv_lifecycle_error(
2965 rx: &mut mpsc::Receiver<Result<IndexBuildOutcome>>,
2966 expected_source: &str,
2967 lifecycle_name: &str,
2968 ) {
2969 let result = tokio::time::timeout(std::time::Duration::from_secs(5), rx.recv())
2970 .await
2971 .unwrap_or_else(|_| {
2972 panic!(
2973 "[{lifecycle_name}] Timeout (5s) waiting for lifecycle error from pending task"
2974 )
2975 });
2976 let result = result.unwrap_or_else(|| {
2977 panic!("[{lifecycle_name}] Channel closed without receiving lifecycle error")
2978 });
2979 let err = result.unwrap_err();
2980 assert_lifecycle_error(err, expected_source, lifecycle_name);
2981 }
2982
2983 #[tokio::test]
2984 async fn test_scheduler_lifecycle_cleanup() {
2985 let env = SchedulerEnv::new().await;
2986
2987 {
2989 let (mut scheduler, mut rx2, mut rx3, _version_control, region_id, building_file_id) =
2990 setup_scheduler_with_pending_tasks(&env).await;
2991
2992 scheduler.on_region_dropped(region_id).await;
2993
2994 let status = scheduler.region_status.get(®ion_id).unwrap();
2995 assert!(status.retiring);
2996 assert_eq!(status.building_files.len(), 1);
2997 assert!(status.pending_tasks.is_empty());
2998
2999 recv_lifecycle_error(&mut rx2, "dropped", "on_region_dropped").await;
3001 recv_lifecycle_error(&mut rx3, "dropped", "on_region_dropped").await;
3002
3003 scheduler.on_task_stopped(region_id, building_file_id);
3005 assert!(!scheduler.region_status.contains_key(®ion_id));
3006 }
3007
3008 {
3010 let (mut scheduler, mut rx2, mut rx3, version_control, region_id, building_file_id) =
3011 setup_scheduler_with_pending_tasks(&env).await;
3012
3013 scheduler.on_region_closed(region_id).await;
3014
3015 let status = scheduler.region_status.get(®ion_id).unwrap();
3016 assert!(status.retiring);
3017 assert_eq!(status.building_files.len(), 1);
3018 assert!(status.pending_tasks.is_empty());
3019
3020 recv_lifecycle_error(&mut rx2, "closed", "on_region_closed").await;
3021 recv_lifecycle_error(&mut rx3, "closed", "on_region_closed").await;
3022
3023 let (reopened_task, _reopened_rx) = create_mock_task_for_schedule_with_result(
3026 &env,
3027 building_file_id,
3028 region_id,
3029 IndexBuildType::Manual,
3030 )
3031 .await;
3032 scheduler
3033 .schedule_build(&version_control, reopened_task)
3034 .await
3035 .unwrap();
3036 let status = scheduler.region_status.get(®ion_id).unwrap();
3037 assert!(status.retiring);
3038 assert_eq!(status.building_files.len(), 1);
3039 assert_eq!(status.pending_tasks.len(), 1);
3040
3041 scheduler
3044 .on_failure(region_id, Arc::new(RegionClosedSnafu { region_id }.build()))
3045 .await;
3046 assert_eq!(scheduler.region_status[®ion_id].pending_tasks.len(), 1);
3047
3048 scheduler.on_task_stopped(region_id, building_file_id);
3049 let status = scheduler.region_status.get(®ion_id).unwrap();
3050 assert!(!status.retiring);
3051 assert_eq!(status.building_files.len(), 1);
3052 assert!(status.pending_tasks.is_empty());
3053
3054 scheduler.on_task_stopped(region_id, building_file_id);
3055 assert!(!scheduler.region_status.contains_key(®ion_id));
3056 }
3057
3058 {
3060 let (mut scheduler, mut rx2, mut rx3, _version_control, region_id, building_file_id) =
3061 setup_scheduler_with_pending_tasks(&env).await;
3062
3063 scheduler.on_region_truncated(region_id).await;
3064
3065 let status = scheduler.region_status.get(®ion_id).unwrap();
3066 assert!(status.retiring);
3067 assert_eq!(status.building_files.len(), 1);
3068 assert!(status.pending_tasks.is_empty());
3069
3070 recv_lifecycle_error(&mut rx2, "truncated", "on_region_truncated").await;
3071 recv_lifecycle_error(&mut rx3, "truncated", "on_region_truncated").await;
3072
3073 scheduler.on_task_stopped(region_id, building_file_id);
3074 assert!(!scheduler.region_status.contains_key(®ion_id));
3075 }
3076 }
3077}