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