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