Skip to main content

mito2/sst/
index.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15pub(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
93/// Evicts local state for a stale index artifact.
94///
95/// The remote artifact may already be referenced by another region manifest
96/// because repartitioned regions share physical SST and index paths. With GC
97/// enabled, remote deletion therefore stays in the reference-aware GC lifecycle.
98pub(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
134/// Triggers background download of an index file to the local cache.
135pub(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/// Output of the index creation.
159#[derive(Debug, Clone, Default)]
160pub struct IndexOutput {
161    /// Size of the file.
162    pub file_size: u64,
163    /// Index version.
164    pub version: u64,
165    /// Inverted index output.
166    pub inverted_index: InvertedIndexOutput,
167    /// Fulltext index output.
168    pub fulltext_index: FulltextIndexOutput,
169    /// Bloom filter output.
170    pub bloom_filter: BloomFilterOutput,
171    /// Vector index output.
172    #[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/// Base output of the index creation.
232#[derive(Debug, Clone, Default)]
233pub struct IndexBaseOutput {
234    /// Size of the index.
235    pub index_size: ByteCount,
236    /// Number of rows in the index.
237    pub row_count: RowCount,
238    /// Available columns in the index.
239    pub columns: Vec<ColumnId>,
240}
241
242impl IndexBaseOutput {
243    pub fn is_available(&self) -> bool {
244        self.index_size > 0
245    }
246}
247
248/// Output of the inverted index creation.
249pub type InvertedIndexOutput = IndexBaseOutput;
250/// Output of the fulltext index creation.
251pub type FulltextIndexOutput = IndexBaseOutput;
252/// Output of the bloom filter creation.
253pub type BloomFilterOutput = IndexBaseOutput;
254/// Output of the vector index creation.
255#[cfg(feature = "vector_index")]
256pub type VectorIndexOutput = IndexBaseOutput;
257
258/// The index creator that hides the error handling details.
259#[derive(Default)]
260pub struct Indexer {
261    file_id: FileId,
262    /// The logical region used for metadata and intermediate index files.
263    region_id: RegionId,
264    /// The region that owns the physical SST and Puffin files.
265    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    /// Updates the index with the given batch.
284    pub async fn update(&mut self, batch: &mut Batch) {
285        self.do_update(batch).await;
286
287        self.flush_mem_metrics();
288    }
289
290    /// Updates the index with the given flat format RecordBatch.
291    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    /// Finalizes the index creation.
298    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    /// Aborts the index creation.
306    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    /// Builds an indexer for the physical SST file.
357    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    /// Sanity check for arguments and create a new [Indexer] if arguments are valid.
382    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 segment row count not aligned with row group size, adjust it to be aligned.
469        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        // Get vector index column IDs and options from metadata
609        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/// Type of an index build task.
657#[derive(Debug, Clone, IntoStaticStr, PartialEq)]
658pub enum IndexBuildType {
659    /// Build index when schema change.
660    SchemaChange,
661    /// Create or update index after flush.
662    Flush,
663    /// Create or update index after compact.
664    Compact,
665    /// Manually build index.
666    Manual,
667}
668
669impl IndexBuildType {
670    fn as_str(&self) -> &'static str {
671        self.into()
672    }
673
674    // Higher value means higher priority.
675    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/// Outcome of an index build task.
695#[derive(Debug, Clone, PartialEq, Eq, Hash)]
696pub enum IndexBuildOutcome {
697    Finished,
698    Aborted(String),
699}
700
701/// Mpsc output result sender.
702pub type ResultMpscSender = Sender<Result<IndexBuildOutcome>>;
703
704#[derive(Clone)]
705pub struct IndexBuildTask {
706    /// The logical region whose manifest and in-memory version this task updates.
707    pub region_id: RegionId,
708    /// The SST file handle to build index for.
709    pub file: FileHandle,
710    /// The target region metadata used to decode rows from the SST.
711    ///
712    /// An SST may originate in another region while being visible in the target
713    /// manifest. This metadata defines the target schema and sequence domain;
714    /// applying the staging manifest only makes imported files visible. Index
715    /// rebuild happens later when a flush, compaction, schema change, or manual
716    /// index build request schedules it.
717    pub(crate) target_region_metadata: RegionMetadataRef,
718    /// The manifest state this build is based on.
719    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    /// When write cache is enabled, the indexer builder should be built from the write cache.
728    /// Otherwise, it should be built from the access layer.
729    pub indexer_builder: Arc<dyn IndexerBuilder + Send + Sync>,
730    /// Request sender to notify the region worker.
731    pub(crate) request_sender: Sender<WorkerRequestWithTime>,
732    /// Index build result sender.
733    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    /// Notify the caller the job is success.
750    pub async fn on_success(&self, outcome: IndexBuildOutcome) {
751        let _ = self.result_sender.send(Ok(outcome)).await;
752    }
753
754    /// Send index build error to waiter.
755    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    // Checks if the SST file still exists in object store and version to avoid conflict with compaction.
800    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        // We should check current version instead of the version when the job is created.
804        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                // If the file's metadata is present in the current version, the physical SST file
817                // is guaranteed to exist on object store. The file purger removes the physical
818                // file only after its metadata is removed from the version.
819                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        // Check SST file existence before building index to avoid failure of parquet reader.
842        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()) // use the latest file handle instead of creating a new one
858            .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        // Use the same file_id but with new version for index file.
872        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            // TODO(SNC123): optimize index batch
884            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            // Check SST file existence again after building index.
901            if !self.check_sst_file_exists(&version_control).await {
902                // Calls abort to clean up index files.
903                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            // Upload index file if write cache is enabled.
917            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                    // If no successor publishes the same version, GC collects the
947                    // remote artifact. Repartition requires GC, while local mode
948                    // accepts this narrow orphan window when GC is disabled.
949                    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                        // An index-unrelated schema change also invalidates the
959                        // generation fence. Retry to avoid leaving the SST
960                        // unindexed; per-SST coalescing limits this to one active
961                        // and one pending task, so repeated changes cannot storm.
962                        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                // Index rebuild is background maintenance work, so its uploads
1020                // are reported as compaction uploads.
1021                .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                // Cleans index files on failure.
1034                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    /// Returns whether this task should replace another pending task for the
1130    /// same SST.
1131    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
1203/// Tracks the index build status of a region scheduled by the [IndexBuildScheduler].
1204struct IndexBuildStatus {
1205    building_files: HashSet<FileId>,
1206    /// At most one coalesced pending task is kept for each SST.
1207    pending_tasks: PendingIndexBuilds,
1208    /// Whether the active builds belong to a region incarnation that has
1209    /// stopped accepting index publications.
1210    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    /// Background job scheduler.
1236    scheduler: SchedulerRef,
1237    /// Tracks regions need to build index.
1238    region_status: HashMap<RegionId, IndexBuildStatus>,
1239    /// Limit of files allowed to build index concurrently on this worker.
1240    files_limit: usize,
1241}
1242
1243/// Manager background index build tasks of a worker.
1244impl 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    /// Schedule tasks until reaching the files limit or no more tasks.
1319    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(&region_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(&region_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    /// Find the ready task with the highest priority.
1371    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(&region_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                // No more tasks for this region, remove it.
1390                self.region_status.remove(&region_id);
1391            } else if status.building_files.is_empty() {
1392                // All builds from the previous region incarnation have stopped.
1393                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(&region_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(&region_id);
1416        }
1417    }
1418
1419    /// Notifies the scheduler that the region is dropped.
1420    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    /// Notifies the scheduler that the region is closed.
1429    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    /// Notifies the scheduler that the region is truncated.
1435    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(&region_id) else {
1445            return;
1446        };
1447        status.retire(err).await;
1448        if status.building_files.is_empty() {
1449            self.region_status.remove(&region_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        // Create mock task.
2061        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        // Schedule the build task and check result.
2082        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            // Old manifests may publish an index without recording its size.
2130            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        // Create mock task.
2145        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        // The task should finish successfully.
2173        match result_rx.recv().await.unwrap() {
2174            Ok(outcome) => {
2175                assert_eq!(outcome, IndexBuildOutcome::Finished);
2176            }
2177            _ => panic!("Expect finished result"),
2178        }
2179
2180        // A notification should be sent to the worker to update the manifest.
2181        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                // The mock indexer builder creates all index types.
2191                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        // Create mock task.
2258        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            // The index file should not exist before the task finishes.
2293            assert!(
2294                !env.access_layer
2295                    .object_store()
2296                    .exists(&puffin_path)
2297                    .await
2298                    .unwrap()
2299            );
2300        } else {
2301            // The index file should exist before the task finishes.
2302            assert!(
2303                env.access_layer
2304                    .object_store()
2305                    .exists(&puffin_path)
2306                    .await
2307                    .unwrap()
2308            );
2309        }
2310
2311        // The task should finish successfully.
2312        match result_rx.recv().await.unwrap() {
2313            Ok(outcome) => {
2314                assert_eq!(outcome, IndexBuildOutcome::Finished);
2315            }
2316            _ => panic!("Expect finished result"),
2317        }
2318
2319        // The index file should exist after the task finishes.
2320        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        // Unset indexes in metadata to simulate no index scenario.
2341        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        // Create mock task.
2369        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        // The task should finish successfully.
2397        match result_rx.recv().await.unwrap() {
2398            Ok(outcome) => {
2399                assert_eq!(outcome, IndexBuildOutcome::Finished);
2400            }
2401            _ => panic!("Expect finished result"),
2402        }
2403
2404        // No index is built, so no notification should be sent to the worker.
2405        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        // Create mock write cache
2421        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, // enable_background_worker
2428                factory,
2429                intm_manager,
2430                ReadableSize::mb(10),
2431            )
2432            .await
2433            .unwrap(),
2434        );
2435        // Indexer builder built from write cache.
2436        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        // Create mock task.
2468        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        // The task should finish successfully.
2496        match result_rx.recv().await.unwrap() {
2497            Ok(outcome) => {
2498                assert_eq!(outcome, IndexBuildOutcome::Finished);
2499            }
2500            _ => panic!("Expect finished result"),
2501        }
2502
2503        // The write cache should contain the uploaded index file.
2504        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    /// Like [`create_mock_task_for_schedule`] but also returns the result receiver
2524    /// so tests can verify pending task cancellation errors.
2525    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[&region_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[&region_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(&region_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        // Prepare multiple files for testing
2701        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        // Test 1: Basic scheduling
2723        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(&region_id));
2732        let status = scheduler.region_status.get(&region_id).unwrap();
2733        assert_eq!(status.building_files.len(), 1);
2734        assert!(status.building_files.contains(&file_id1));
2735
2736        // Test 2: Duplicate file scheduling (should be skipped)
2737        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(&region_id).unwrap();
2744        assert_eq!(status.building_files.len(), 1); // Still only one
2745
2746        // Test 3: Fill up to limit (2 building tasks)
2747        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(&region_id).unwrap();
2754        assert_eq!(status.building_files.len(), 2); // Reached limit
2755        assert_eq!(status.pending_tasks.len(), 0);
2756
2757        // Test 4: Add tasks with different priorities to pending queue
2758        // Now all new tasks will be pending since we reached the limit
2759        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(&region_id).unwrap();
2781        assert_eq!(status.building_files.len(), 2); // Still at limit
2782        assert_eq!(status.pending_tasks.len(), 3); // Three pending
2783
2784        // Test 5: Task completion triggers scheduling next highest priority task (Manual)
2785        scheduler.on_task_stopped(region_id, file_id1);
2786        let status = scheduler.region_status.get(&region_id).unwrap();
2787        assert!(!status.building_files.contains(&file_id1));
2788        assert_eq!(status.building_files.len(), 2); // Should schedule next task
2789        assert_eq!(status.pending_tasks.len(), 2); // One less pending
2790        // The highest priority task (Manual) should now be building
2791        assert!(status.building_files.contains(&file_id5));
2792
2793        // Test 6: Complete another task, should schedule SchemaChange (second highest priority)
2794        scheduler.on_task_stopped(region_id, file_id2);
2795        let status = scheduler.region_status.get(&region_id).unwrap();
2796        assert_eq!(status.building_files.len(), 2);
2797        assert_eq!(status.pending_tasks.len(), 1); // One less pending
2798        assert!(status.building_files.contains(&file_id4)); // SchemaChange should be building
2799
2800        // Test 7: Complete remaining tasks and cleanup
2801        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(&region_id).unwrap();
2805        assert_eq!(status.building_files.len(), 1); // Last task (Compact) should be building
2806        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        // Region should be removed when all tasks complete
2812        assert!(!scheduler.region_status.contains_key(&region_id));
2813
2814        // Test 8: A build failure keeps active leases without retiring the region
2815        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(&region_id));
2836        let status = scheduler.region_status.get(&region_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(&region_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(&region_id));
2858        scheduler.on_task_stopped(region_id, file_id2);
2859        assert!(!scheduler.region_status.contains_key(&region_id));
2860    }
2861
2862    /// Helper to set up a scheduler with files_limit=1 and 3 scheduled tasks,
2863    /// returning the scheduler, the two pending-task result receivers, and the
2864    /// version control.
2865    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, // building file_id for no-op assertion
2874    ) {
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        // task1 becomes the "building" task (files_limit=1).
2918        // We intentionally drop its result receiver: the building task's late-stop
2919        // behavior is covered by the manual on_task_stopped no-op assertion below.
2920        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        // Verify: 1 building + 2 pending.
2956        assert!(scheduler.region_status.contains_key(&region_id));
2957        let status = scheduler.region_status.get(&region_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    /// Pattern‑matches a pending‑task cancellation error: outer **must** be
2965    /// [`crate::error::Error::BuildIndexAsync`] and the inner source **must** be
2966    /// the lifecycle variant named by `expected_source` (`"dropped"`, `"closed"`,
2967    /// or `"truncated"`).  Panics with a descriptive message on mismatch.
2968    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    /// Receives a pending‑task error from `rx` with a 5‑second timeout, then
2991    /// delegates to [`assert_lifecycle_error`].
2992    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        // --- on_region_dropped ---
3016        {
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(&region_id).unwrap();
3023            assert!(status.retiring);
3024            assert_eq!(status.building_files.len(), 1);
3025            assert!(status.pending_tasks.is_empty());
3026
3027            // Pending-task receivers get lifecycle errors (with timeout).
3028            recv_lifecycle_error(&mut rx2, "dropped", "on_region_dropped").await;
3029            recv_lifecycle_error(&mut rx3, "dropped", "on_region_dropped").await;
3030
3031            // The active build keeps its lease until it stops.
3032            scheduler.on_task_stopped(region_id, building_file_id);
3033            assert!(!scheduler.region_status.contains_key(&region_id));
3034        }
3035
3036        // --- on_region_closed ---
3037        {
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(&region_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            // A build scheduled by the reopened region waits behind the old
3052            // incarnation's active lease.
3053            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(&region_id).unwrap();
3065            assert!(status.retiring);
3066            assert_eq!(status.building_files.len(), 1);
3067            assert_eq!(status.pending_tasks.len(), 1);
3068
3069            // A manifest error from the old incarnation must not discard the
3070            // reopened task before the old build reports stopped.
3071            scheduler
3072                .on_failure(region_id, Arc::new(RegionClosedSnafu { region_id }.build()))
3073                .await;
3074            assert_eq!(scheduler.region_status[&region_id].pending_tasks.len(), 1);
3075
3076            scheduler.on_task_stopped(region_id, building_file_id);
3077            let status = scheduler.region_status.get(&region_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(&region_id));
3084        }
3085
3086        // --- on_region_truncated ---
3087        {
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(&region_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(&region_id));
3103        }
3104    }
3105}