Skip to main content

mito2/compaction/
compactor.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
15use std::num::NonZero;
16use std::sync::Arc;
17use std::time::Duration;
18
19use common_base::Plugins;
20use common_base::cancellation::{CancellableFuture, CancellationHandle};
21use common_meta::key::SchemaMetadataManagerRef;
22use common_telemetry::{debug, info, warn};
23use common_time::TimeToLive;
24use either::Either;
25use itertools::Itertools;
26use object_store::manager::ObjectStoreManagerRef;
27use partition::expr::PartitionExpr;
28use serde::{Deserialize, Serialize};
29use snafu::{OptionExt, ResultExt};
30use store_api::ManifestVersion;
31use store_api::metadata::RegionMetadataRef;
32use store_api::region_request::PathType;
33use store_api::storage::RegionId;
34
35use crate::access_layer::{
36    AccessLayer, AccessLayerRef, Metrics, OperationType, SstWriteRequest, WriteType,
37};
38use crate::cache::{CacheManager, CacheManagerRef};
39use crate::compaction::picker::PickerOutput;
40use crate::compaction::reader::CompactionSstReaderBuilder;
41use crate::compaction::{CompactionOutput, find_dynamic_options};
42use crate::config::MitoConfig;
43use crate::engine::region_hook::{RegionHookRef, SstFileInfo};
44use crate::error;
45use crate::error::{
46    EmptyRegionDirSnafu, InvalidPartitionExprSnafu, ObjectStoreNotFoundSnafu, Result,
47};
48use crate::manifest::action::{RegionEdit, RegionMetaAction, RegionMetaActionList};
49use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
50use crate::region::options::RegionOptions;
51use crate::region::version::VersionRef;
52use crate::region::{ManifestContext, RegionLeaderState, RegionRoleState};
53use crate::schedule::scheduler::LocalScheduler;
54use crate::sst::FormatType;
55use crate::sst::file::{FileMeta, UncommittedSsts};
56use crate::sst::file_purger::LocalFilePurger;
57use crate::sst::index::intermediate::IntermediateManager;
58use crate::sst::index::puffin_manager::PuffinManagerFactory;
59use crate::sst::location::region_dir_from_table_dir;
60use crate::sst::parquet::metadata::extract_primary_key_range;
61use crate::sst::parquet::{SstInfo, WriteOptions};
62use crate::sst::version::{SstVersion, SstVersionRef};
63
64/// Region version for compaction that does not hold memtables.
65#[derive(Clone)]
66pub struct CompactionVersion {
67    /// Metadata of the region.
68    ///
69    /// Altering metadata isn't frequent, storing metadata in Arc to allow sharing
70    /// metadata and reuse metadata when creating a new `Version`.
71    pub(crate) metadata: RegionMetadataRef,
72    /// Options of the region.
73    pub(crate) options: RegionOptions,
74    /// SSTs of the region.
75    pub(crate) ssts: SstVersionRef,
76    /// Inferred compaction time window.
77    pub(crate) compaction_time_window: Option<Duration>,
78}
79
80impl From<VersionRef> for CompactionVersion {
81    fn from(value: VersionRef) -> Self {
82        Self {
83            metadata: value.metadata.clone(),
84            options: value.options.clone(),
85            ssts: value.ssts.clone(),
86            compaction_time_window: value.compaction_time_window,
87        }
88    }
89}
90
91/// CompactionRegion represents a region that needs to be compacted.
92/// It's the subset of MitoRegion.
93#[derive(Clone)]
94pub struct CompactionRegion {
95    pub region_id: RegionId,
96    pub region_options: RegionOptions,
97
98    pub(crate) engine_config: Arc<MitoConfig>,
99    pub(crate) region_metadata: RegionMetadataRef,
100    pub(crate) cache_manager: CacheManagerRef,
101    /// Access layer to get the table path and path type.
102    pub access_layer: AccessLayerRef,
103    pub(crate) manifest_ctx: Arc<ManifestContext>,
104    pub(crate) current_version: CompactionVersion,
105    pub(crate) file_purger: Option<Arc<LocalFilePurger>>,
106    pub(crate) ttl: Option<TimeToLive>,
107
108    /// Controls the parallelism of this compaction task. Default is 1.
109    ///
110    /// The parallel is inside this compaction task, not across different compaction tasks.
111    /// It can be different windows of the same compaction task or something like this.
112    pub max_parallelism: usize,
113
114    pub(crate) plugins: Plugins,
115}
116
117/// OpenCompactionRegionRequest represents the request to open a compaction region.
118#[derive(Clone)]
119pub struct OpenCompactionRegionRequest {
120    pub region_id: RegionId,
121    pub table_dir: String,
122    pub path_type: PathType,
123    pub region_options: RegionOptions,
124    pub max_parallelism: usize,
125    /// Plugins for the compaction region, used to look up the [`RegionHook`](crate::engine::region_hook::RegionHook).
126    pub plugins: Plugins,
127}
128
129/// Open a compaction region from a compaction request.
130/// It's simple version of RegionOpener::open().
131pub async fn open_compaction_region(
132    req: &OpenCompactionRegionRequest,
133    mito_config: &MitoConfig,
134    object_store_manager: ObjectStoreManagerRef,
135    ttl_provider: Either<TimeToLive, SchemaMetadataManagerRef>,
136) -> Result<CompactionRegion> {
137    let object_store = {
138        let name = &req.region_options.storage;
139        if let Some(name) = name {
140            object_store_manager
141                .find(name)
142                .with_context(|| ObjectStoreNotFoundSnafu {
143                    object_store: name.clone(),
144                })?
145        } else {
146            object_store_manager.default_object_store()
147        }
148    };
149
150    let access_layer = {
151        let puffin_manager_factory = PuffinManagerFactory::new(
152            &mito_config.index.aux_path,
153            mito_config.index.staging_size.as_bytes(),
154            Some(mito_config.index.write_buffer_size.as_bytes() as _),
155            mito_config.index.staging_ttl,
156        )
157        .await?;
158        let intermediate_manager =
159            IntermediateManager::init_fs(mito_config.index.aux_path.clone()).await?;
160
161        Arc::new(AccessLayer::new(
162            &req.table_dir,
163            req.path_type,
164            object_store.clone(),
165            puffin_manager_factory,
166            intermediate_manager,
167        ))
168    };
169
170    let manifest_manager = {
171        let region_dir = region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type);
172        let region_manifest_options =
173            RegionManifestOptions::new(mito_config, &region_dir, object_store);
174
175        RegionManifestManager::open(region_manifest_options, &Default::default())
176            .await?
177            .with_context(|| EmptyRegionDirSnafu {
178                region_id: req.region_id,
179                region_dir: region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type),
180            })?
181    };
182
183    let manifest = manifest_manager.manifest();
184    let region_metadata = manifest.metadata.clone();
185    let hook: Option<RegionHookRef> = req.plugins.get();
186    let manifest_ctx = Arc::new(ManifestContext::new(
187        manifest_manager,
188        RegionRoleState::Leader(RegionLeaderState::Writable),
189        hook,
190    ));
191
192    let file_purger = {
193        let purge_scheduler = Arc::new(LocalScheduler::new(mito_config.max_background_purges));
194        Arc::new(LocalFilePurger::new(
195            purge_scheduler.clone(),
196            access_layer.clone(),
197            None,
198        ))
199    };
200
201    let current_version = {
202        let mut ssts = SstVersion::new();
203        ssts.add_files(file_purger.clone(), manifest.files.values().cloned());
204        CompactionVersion {
205            metadata: region_metadata.clone(),
206            options: req.region_options.clone(),
207            ssts: Arc::new(ssts),
208            compaction_time_window: manifest.compaction_time_window,
209        }
210    };
211
212    let ttl = match ttl_provider {
213        // Use the specified ttl.
214        Either::Left(ttl) => ttl,
215        // Get the ttl from the schema metadata manager.
216        Either::Right(schema_metadata_manager) => {
217            let (_, ttl) =
218                find_dynamic_options(req.region_id, &req.region_options, &schema_metadata_manager)
219                    .await
220                    .unwrap_or_else(|e| {
221                        warn!(e; "Failed to get ttl for region: {}", region_metadata.region_id);
222                        (
223                            crate::region::options::CompactionOptions::default(),
224                            TimeToLive::default(),
225                        )
226                    });
227            ttl
228        }
229    };
230
231    Ok(CompactionRegion {
232        region_id: req.region_id,
233        region_options: req.region_options.clone(),
234        engine_config: Arc::new(mito_config.clone()),
235        region_metadata: region_metadata.clone(),
236        cache_manager: Arc::new(CacheManager::default()),
237        access_layer,
238        manifest_ctx,
239        current_version,
240        file_purger: Some(file_purger),
241        ttl: Some(ttl),
242        max_parallelism: req.max_parallelism,
243        plugins: req.plugins.clone(),
244    })
245}
246
247impl CompactionRegion {
248    /// Get the file purger of the compaction region.
249    pub fn file_purger(&self) -> Option<Arc<LocalFilePurger>> {
250        self.file_purger.clone()
251    }
252
253    /// Stop the file purger scheduler of the compaction region.
254    pub async fn stop_purger_scheduler(&self) -> Result<()> {
255        if let Some(file_purger) = &self.file_purger {
256            file_purger.stop_scheduler().await
257        } else {
258            Ok(())
259        }
260    }
261
262    /// Fires [`RegionHook::on_sst_files_written`] for the freshly-merged SST
263    /// files in `merge_output`. Shared by both compaction paths, it must run
264    /// before [`Compactor::update_manifest`], whose `on_manifest_updated`
265    /// drains the per-region state this hook populates.
266    pub async fn invoke_sst_hook(&self, merge_output: &MergeOutput) {
267        let Some(hook) = self.plugins.get::<RegionHookRef>() else {
268            return;
269        };
270
271        // Remote compaction deserializes `MergeOutput` over gRPC, where
272        // `sst_infos` is `#[serde(skip)]`. Rebuild each `SstInfo` from the
273        // matching `FileMeta` so the hook observes real ids/rows/sizes rather
274        // than zeros — see [`sst_info_from_file_meta`] for what stays default.
275        let synthesized: Vec<SstInfo>;
276        let infos: &[SstInfo] = if merge_output.sst_infos.is_empty() {
277            synthesized = merge_output
278                .files_to_add
279                .iter()
280                .map(sst_info_from_file_meta)
281                .collect();
282            &synthesized
283        } else {
284            // `sst_infos` and `files_to_add` are documented as 1:1. If they
285            // ever diverge, `zip` below would silently truncate; warn so the
286            // mismatch is observable rather than dropping files from the hook.
287            if merge_output.sst_infos.len() != merge_output.files_to_add.len() {
288                warn!(
289                    "sst_infos length ({}) does not match files_to_add length ({}) for region {}",
290                    merge_output.sst_infos.len(),
291                    merge_output.files_to_add.len(),
292                    self.region_id
293                );
294            }
295            &merge_output.sst_infos
296        };
297
298        let files: Vec<SstFileInfo<'_>> = merge_output
299            .files_to_add
300            .iter()
301            .zip(infos)
302            .map(|(meta, info)| SstFileInfo {
303                sst_info_ref: info,
304                file_meta: meta,
305            })
306            .collect();
307        hook.on_sst_files_written(self.region_id, &self.region_metadata, &files)
308            .await;
309    }
310}
311
312/// Builds an [`SstInfo`] for the remote-compaction path by copying the scalar
313/// fields that [`FileMeta`] also carries.
314///
315/// Remote compaction runs off-process and ships `MergeOutput` over the wire;
316/// [`MergeOutput::sst_infos`] is `#[serde(skip)]` because the parquet footer
317/// (`ParquetMetaData`) is not serde-serializable. To avoid feeding the hook
318/// default (zero) values for real SSTs, this rebuilds the seven fields `FileMeta`
319/// and `SstInfo` share: `file_id`, `time_range`, `file_size`,
320/// `max_row_group_uncompressed_size`, `num_rows`, `num_row_groups`, `num_series`.
321///
322/// Two fields stay default and **cannot be recovered on the datanode**:
323/// - `file_metadata` — the parquet footer, whose column min/max/null-count
324///   statistics are required by hooks building richer artifacts (e.g. an Iceberg
325///   manifest). That data is produced by the compactor's writer and exists only
326///   in the freshly-written SST on object storage.
327/// - `index_metadata` — `FileMeta` tracks indexes via `available_indexes` /
328///   `indexes`, not as an [`IndexOutput`].
329///
330/// A hook that needs column stats must fetch the footer from object storage
331/// (the [`CompactionRegion`]'s `access_layer` reaches the store), or the footer
332/// must be shipped over the wire (revisit the `#[serde(skip)]` on `sst_infos`).
333/// `num_rows` may read `0` for legacy `FileMeta`s where the count is unknown.
334fn sst_info_from_file_meta(meta: &FileMeta) -> SstInfo {
335    SstInfo {
336        file_id: meta.file_id,
337        time_range: meta.time_range,
338        file_size: meta.file_size,
339        max_row_group_uncompressed_size: meta.max_row_group_uncompressed_size,
340        num_rows: meta.num_rows as usize,
341        num_row_groups: meta.num_row_groups,
342        num_series: meta.num_series,
343        ..Default::default()
344    }
345}
346
347/// `[MergeOutput]` represents the output of merging SST files.
348#[derive(Default, Debug, Serialize, Deserialize)]
349pub struct MergeOutput {
350    pub files_to_add: Vec<FileMeta>,
351    pub files_to_remove: Vec<FileMeta>,
352    pub compaction_time_window: Option<i64>,
353    #[serde(skip)]
354    pub sst_infos: Vec<SstInfo>,
355}
356
357impl MergeOutput {
358    pub fn is_empty(&self) -> bool {
359        self.files_to_add.is_empty() && self.files_to_remove.is_empty()
360    }
361
362    pub fn input_file_size(&self) -> u64 {
363        self.files_to_remove.iter().map(|f| f.file_size).sum()
364    }
365
366    pub fn output_file_size(&self) -> u64 {
367        self.files_to_add.iter().map(|f| f.file_size).sum()
368    }
369}
370
371/// Compactor is the trait that defines the compaction logic.
372#[async_trait::async_trait]
373pub trait Compactor: Send + Sync + 'static {
374    /// Merge SST files for a region.
375    async fn merge_ssts(
376        &self,
377        compaction_region: &CompactionRegion,
378        picker_output: PickerOutput,
379    ) -> Result<MergeOutput>;
380
381    /// Update the manifest after merging SST files.
382    async fn update_manifest(
383        &self,
384        compaction_region: &CompactionRegion,
385        merge_output: MergeOutput,
386    ) -> Result<(RegionEdit, ManifestVersion)>;
387}
388
389/// Trait for merging a single compaction output into SST files.
390///
391/// This is extracted from `DefaultCompactor` to allow injecting mock
392/// implementations in tests.
393#[async_trait::async_trait]
394pub trait SstMerger: Send + Sync + 'static {
395    async fn merge_single_output(
396        &self,
397        compaction_region: CompactionRegion,
398        output: CompactionOutput,
399        write_opts: WriteOptions,
400    ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)>;
401}
402
403/// The production [`SstMerger`] that reads, merges, and writes SST files.
404#[derive(Clone)]
405pub struct DefaultSstMerger;
406
407#[async_trait::async_trait]
408impl SstMerger for DefaultSstMerger {
409    async fn merge_single_output(
410        &self,
411        compaction_region: CompactionRegion,
412        output: CompactionOutput,
413        write_opts: WriteOptions,
414    ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
415        let region_id = compaction_region.region_id;
416        let storage = compaction_region.region_options.storage.clone();
417        let index_options = compaction_region
418            .current_version
419            .options
420            .index_options
421            .clone();
422        let append_mode = compaction_region.current_version.options.append_mode;
423        let merge_mode = compaction_region.current_version.options.merge_mode();
424        let flat_format = compaction_region
425            .region_options
426            .sst_format
427            .map(|format| format == FormatType::Flat)
428            .unwrap_or(compaction_region.engine_config.default_flat_format);
429
430        let index_config = compaction_region.engine_config.index.clone();
431        let inverted_index_config = compaction_region.engine_config.inverted_index.clone();
432        let fulltext_index_config = compaction_region.engine_config.fulltext_index.clone();
433        let bloom_filter_index_config = compaction_region.engine_config.bloom_filter_index.clone();
434        #[cfg(feature = "vector_index")]
435        let vector_index_config = compaction_region.engine_config.vector_index.clone();
436
437        let input_file_names = output
438            .inputs
439            .iter()
440            .map(|f| f.file_id().to_string())
441            .join(",");
442        let max_sequence = output
443            .inputs
444            .iter()
445            .map(|f| f.meta_ref().sequence)
446            .max()
447            .flatten();
448        let builder = CompactionSstReaderBuilder {
449            metadata: compaction_region.region_metadata.clone(),
450            sst_layer: compaction_region.access_layer.clone(),
451            cache: compaction_region.cache_manager.clone(),
452            inputs: &output.inputs,
453            append_mode,
454            filter_deleted: output.filter_deleted,
455            time_range: output.output_time_range,
456            merge_mode,
457        };
458        let source = builder.build_flat_sst_reader().await?;
459
460        let mut metrics = Metrics::new(WriteType::Compaction);
461        let region_metadata = compaction_region.region_metadata.clone();
462        let sst_infos = compaction_region
463            .access_layer
464            .write_sst(
465                SstWriteRequest {
466                    op_type: OperationType::Compact,
467                    metadata: region_metadata.clone(),
468                    source,
469                    cache_manager: compaction_region.cache_manager.clone(),
470                    storage,
471                    max_sequence: max_sequence.map(NonZero::get),
472                    sst_write_format: if flat_format {
473                        FormatType::Flat
474                    } else {
475                        FormatType::PrimaryKey
476                    },
477                    index_options,
478                    index_config,
479                    inverted_index_config,
480                    fulltext_index_config,
481                    bloom_filter_index_config,
482                    #[cfg(feature = "vector_index")]
483                    vector_index_config,
484                },
485                &write_opts,
486                &mut metrics,
487            )
488            .await?;
489        // Convert partition expression once outside the map
490        let partition_expr = match &region_metadata.partition_expr {
491            None => None,
492            Some(json_str) if json_str.is_empty() => None,
493            Some(json_str) => PartitionExpr::from_json_str(json_str).with_context(|_| {
494                InvalidPartitionExprSnafu {
495                    expr: json_str.clone(),
496                }
497            })?,
498        };
499
500        let output_files = sst_infos
501            .iter()
502            .map(|sst_info| {
503                let pk_range = sst_info
504                    .file_metadata
505                    .as_ref()
506                    .and_then(|meta| extract_primary_key_range(meta, &region_metadata));
507                let (primary_key_min, primary_key_max) = match pk_range {
508                    Some((min, max)) => (Some(min), Some(max)),
509                    None => (None, None),
510                };
511
512                FileMeta {
513                    region_id,
514                    file_id: sst_info.file_id,
515                    time_range: sst_info.time_range,
516                    level: output.output_level,
517                    file_size: sst_info.file_size,
518                    max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
519                    available_indexes: sst_info.index_metadata.build_available_indexes(),
520                    indexes: sst_info.index_metadata.build_indexes(),
521                    index_file_size: sst_info.index_metadata.file_size,
522                    index_version: 0,
523                    num_rows: sst_info.num_rows as u64,
524                    num_row_groups: sst_info.num_row_groups,
525                    sequence: max_sequence,
526                    partition_expr: partition_expr.clone(),
527                    num_series: sst_info.num_series,
528                    primary_key_min,
529                    primary_key_max,
530                }
531            })
532            .collect::<Vec<_>>();
533        let output_file_names = output_files.iter().map(|f| f.file_id.to_string()).join(",");
534        info!(
535            "Region {} compaction inputs: [{}], outputs: [{}], flat_format: {}, metrics: {:?}",
536            region_id, input_file_names, output_file_names, flat_format, metrics
537        );
538        metrics.observe();
539        Ok((output_files, sst_infos.into_iter().collect()))
540    }
541}
542
543/// DefaultCompactor is the default implementation of Compactor.
544///
545/// It is parameterized by an [`SstMerger`] to allow injecting mock
546/// implementations in tests.
547pub struct DefaultCompactor<M = DefaultSstMerger> {
548    merger: M,
549    cancel_handle: Arc<CancellationHandle>,
550    uncommitted: Option<UncommittedSsts>,
551}
552
553#[cfg(test)]
554impl<M: SstMerger> DefaultCompactor<M> {
555    pub fn with_merger(merger: M) -> Self {
556        Self {
557            merger,
558            cancel_handle: Arc::new(CancellationHandle::default()),
559            uncommitted: None,
560        }
561    }
562}
563
564impl DefaultCompactor {
565    /// Creates a new `DefaultCompactor` with the given cancel handle and no
566    /// uncommitted SSTs tracking.
567    ///
568    /// This is the public entry point for external crates that want a
569    /// cancellable compactor without opting into uncommitted-SST tracking.
570    pub fn new_with_cancel_handle(cancel_handle: Arc<CancellationHandle>) -> Self {
571        Self {
572            merger: DefaultSstMerger,
573            cancel_handle,
574            uncommitted: None,
575        }
576    }
577
578    pub(crate) fn with_cancel_handle(
579        cancel_handle: Arc<CancellationHandle>,
580        uncommitted: UncommittedSsts,
581    ) -> Self {
582        Self {
583            merger: DefaultSstMerger,
584            cancel_handle,
585            uncommitted: Some(uncommitted),
586        }
587    }
588}
589
590#[async_trait::async_trait]
591impl<M: SstMerger> Compactor for DefaultCompactor<M>
592where
593    M: Clone,
594{
595    async fn merge_ssts(
596        &self,
597        compaction_region: &CompactionRegion,
598        mut picker_output: PickerOutput,
599    ) -> Result<MergeOutput> {
600        let internal_parallelism = compaction_region.max_parallelism.max(1);
601        let compaction_time_window = picker_output.time_window_size;
602        let region_id = compaction_region.region_id;
603
604        // Build tasks along with their input file metas so we can track which
605        // inputs correspond to each task.
606        let mut tasks: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(picker_output.outputs.len());
607
608        for output in picker_output.outputs.drain(..) {
609            let inputs_to_remove: Vec<_> =
610                output.inputs.iter().map(|f| f.meta_ref().clone()).collect();
611            let write_opts = WriteOptions {
612                write_buffer_size: compaction_region.engine_config.sst_write_buffer_size,
613                max_file_size: picker_output.max_file_size,
614                row_group_size: compaction_region.region_options.row_group_size(),
615            };
616            let merger = self.merger.clone();
617            let compaction_region = compaction_region.clone();
618            let uncommitted = self.uncommitted.clone();
619            let fut = async move {
620                let result = merger
621                    .merge_single_output(compaction_region, output, write_opts)
622                    .await;
623                if let (Some(uncommitted), Ok((_, infos))) = (&uncommitted, &result) {
624                    uncommitted.track(infos);
625                }
626                result
627            };
628            tasks.push((inputs_to_remove, fut));
629        }
630
631        let hook: Option<RegionHookRef> = compaction_region.plugins.get();
632        let mut output_files = Vec::with_capacity(tasks.len());
633        let mut all_sst_infos: Vec<SstInfo> = Vec::new();
634        let mut compacted_inputs = Vec::with_capacity(
635            tasks.iter().map(|(inputs, _)| inputs.len()).sum::<usize>()
636                + picker_output.expired_ssts.len(),
637        );
638
639        while !tasks.is_empty() {
640            let mut chunk: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(internal_parallelism);
641            for _ in 0..internal_parallelism {
642                if let Some(task) = tasks.pop() {
643                    chunk.push(task);
644                }
645            }
646            let mut spawned: Vec<_> = chunk
647                .into_iter()
648                .map(|(inputs, fut)| {
649                    let handle = common_runtime::spawn_compact(fut);
650                    (inputs, handle)
651                })
652                .collect();
653
654            while let Some((inputs, mut handle)) = spawned.pop() {
655                let abort_handle = handle.abort_handle();
656                match CancellableFuture::new(&mut handle, self.cancel_handle.clone()).await {
657                    Ok(Ok(Ok((files, infos)))) => {
658                        output_files.extend(files);
659                        if hook.is_some() {
660                            all_sst_infos.extend(infos);
661                        }
662                        compacted_inputs.extend(inputs);
663                    }
664                    Ok(Ok(Err(e))) => {
665                        warn!(
666                            e; "Failed to merge compaction output for region: {}, inputs: [{}]",
667                            region_id,
668                            inputs.iter().map(|f| f.file_id.to_string()).join(",")
669                        );
670                    }
671                    Ok(Err(e)) => {
672                        warn!(
673                            "Region {} compaction task join error for inputs: [{}], skipping: {}",
674                            region_id,
675                            inputs.iter().map(|f| f.file_id.to_string()).join(","),
676                            e
677                        );
678                        // If the cancel handle is cancelled,
679                        // cancel the remaining tasks before returns the error.
680                        if self.cancel_handle.is_cancelled() {
681                            abort_handle.abort();
682                            for (_, handle) in &spawned {
683                                handle.abort();
684                            }
685                            for (_, handle) in spawned {
686                                let _ = handle.await;
687                            }
688                        }
689                        return Err(e).context(error::JoinSnafu);
690                    }
691                    Err(_) => {
692                        debug!(
693                            "Compaction merge cancelled for region: {}, aborting remaining {} spawned tasks",
694                            region_id,
695                            spawned.len(),
696                        );
697                        abort_handle.abort();
698                        for (_, handle) in &spawned {
699                            handle.abort();
700                        }
701                        let _ = handle.await;
702                        for (_, handle) in spawned {
703                            let _ = handle.await;
704                        }
705                        break;
706                    }
707                }
708            }
709
710            if self.cancel_handle.is_cancelled() {
711                info!("Compaction merge cancelled for region: {}", region_id);
712                break;
713            }
714        }
715
716        // Include expired SSTs in removals — these don't depend on merge success.
717        compacted_inputs.extend(
718            picker_output
719                .expired_ssts
720                .iter()
721                .map(|f| f.meta_ref().clone()),
722        );
723
724        Ok(MergeOutput {
725            files_to_add: output_files,
726            files_to_remove: compacted_inputs,
727            compaction_time_window: Some(compaction_time_window),
728            sst_infos: all_sst_infos,
729        })
730    }
731
732    async fn update_manifest(
733        &self,
734        compaction_region: &CompactionRegion,
735        merge_output: MergeOutput,
736    ) -> Result<(RegionEdit, ManifestVersion)> {
737        // Write region edit to manifest.
738        let edit = RegionEdit {
739            files_to_add: merge_output.files_to_add,
740            files_to_remove: merge_output.files_to_remove,
741            // Use current timestamp as the edit timestamp.
742            timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
743            compaction_time_window: merge_output
744                .compaction_time_window
745                .map(|seconds| Duration::from_secs(seconds as u64)),
746            flushed_entry_id: None,
747            flushed_sequence: None,
748            committed_sequence: None,
749        };
750
751        let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone()));
752        // TODO: We might leak files if we fail to update manifest. We can add a cleanup task to remove them later.
753        let manifest_version = compaction_region
754            .manifest_ctx
755            .update_manifest_for_compaction(action_list)
756            .await?;
757
758        Ok((edit, manifest_version))
759    }
760}
761
762#[cfg(test)]
763mod tests {
764    use std::sync::atomic::{AtomicUsize, Ordering};
765    use std::sync::{Arc, Mutex};
766    use std::time::Duration;
767
768    use store_api::storage::{FileId, RegionId};
769    use tokio::time::sleep;
770
771    use super::{DefaultCompactor, *};
772    use crate::cache::CacheManager;
773    use crate::compaction::picker::PickerOutput;
774    use crate::error::Result;
775    use crate::sst::file::FileHandle;
776    use crate::sst::file_purger::NoopFilePurger;
777    use crate::sst::version::SstVersion;
778    use crate::test_util::memtable_util::metadata_for_test;
779    use crate::test_util::scheduler_util::SchedulerEnv;
780
781    fn dummy_file_meta() -> FileMeta {
782        FileMeta {
783            region_id: RegionId::new(1, 1),
784            file_id: FileId::random(),
785            file_size: 100,
786            ..Default::default()
787        }
788    }
789
790    fn new_file_handle(meta: FileMeta) -> FileHandle {
791        FileHandle::new(meta, Arc::new(NoopFilePurger))
792    }
793
794    /// Build a minimal [`CompactionRegion`] suitable for tests where the
795    /// [`SstMerger`] is mocked and never touches the access layer.
796    async fn new_test_compaction_region() -> CompactionRegion {
797        let env = SchedulerEnv::new().await;
798        let metadata = metadata_for_test();
799        let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
800        CompactionRegion {
801            region_id: RegionId::new(1, 1),
802            region_options: RegionOptions::default(),
803            engine_config: Arc::new(MitoConfig::default()),
804            region_metadata: metadata.clone(),
805            cache_manager: Arc::new(CacheManager::default()),
806            access_layer: env.access_layer.clone(),
807            manifest_ctx,
808            current_version: CompactionVersion {
809                metadata,
810                options: RegionOptions::default(),
811                ssts: Arc::new(SstVersion::new()),
812                compaction_time_window: None,
813            },
814            file_purger: None,
815            ttl: None,
816            max_parallelism: 1,
817            plugins: Plugins::new(),
818        }
819    }
820
821    /// An [`SstMerger`] that returns pre-configured results per call index.
822    ///
823    /// Call 0 gets `results[0]`, call 1 gets `results[1]`, etc.
824    #[derive(Clone)]
825    struct MockMerger {
826        results: Arc<Mutex<Vec<Result<Vec<FileMeta>>>>>,
827        call_idx: Arc<AtomicUsize>,
828    }
829
830    impl MockMerger {
831        fn new(results: Vec<Result<Vec<FileMeta>>>) -> Self {
832            Self {
833                results: Arc::new(Mutex::new(results)),
834                call_idx: Arc::new(AtomicUsize::new(0)),
835            }
836        }
837    }
838
839    #[async_trait::async_trait]
840    impl SstMerger for MockMerger {
841        async fn merge_single_output(
842            &self,
843            _compaction_region: CompactionRegion,
844            _output: CompactionOutput,
845            _write_opts: WriteOptions,
846        ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
847            let idx = self.call_idx.fetch_add(1, Ordering::SeqCst);
848            match self.results.lock().unwrap().get(idx) {
849                Some(Ok(files)) => Ok((files.clone(), Vec::new())),
850                Some(Err(_)) => error::InvalidMetaSnafu {
851                    reason: format!("simulated failure at index {idx}"),
852                }
853                .fail(),
854                None => panic!("MockMerger: no result configured for call index {idx}"),
855            }
856        }
857    }
858
859    #[tokio::test]
860    async fn test_partial_merge_failure_collects_only_successful_outputs() {
861        common_telemetry::init_default_ut_logging();
862
863        let compaction_region = new_test_compaction_region().await;
864
865        // Prepare 3 compaction outputs: output 0 and 2 succeed, output 1 fails.
866        let input_meta_0 = dummy_file_meta();
867        let input_meta_1 = dummy_file_meta();
868        let input_meta_2 = dummy_file_meta();
869
870        let output_meta_0 = vec![dummy_file_meta()];
871        let output_meta_2 = vec![dummy_file_meta(), dummy_file_meta()];
872
873        let merger = MockMerger::new(vec![
874            Ok(output_meta_0.clone()),
875            Err(error::InvalidMetaSnafu {
876                reason: "boom".to_string(),
877            }
878            .build()),
879            Ok(output_meta_2.clone()),
880        ]);
881        let compactor = DefaultCompactor::with_merger(merger);
882
883        let picker_output = PickerOutput {
884            outputs: vec![
885                CompactionOutput {
886                    output_level: 1,
887                    inputs: vec![new_file_handle(input_meta_0.clone())],
888                    filter_deleted: false,
889                    output_time_range: None,
890                },
891                CompactionOutput {
892                    output_level: 1,
893                    inputs: vec![new_file_handle(input_meta_1.clone())],
894                    filter_deleted: false,
895                    output_time_range: None,
896                },
897                CompactionOutput {
898                    output_level: 1,
899                    inputs: vec![new_file_handle(input_meta_2.clone())],
900                    filter_deleted: false,
901                    output_time_range: None,
902                },
903            ],
904            expired_ssts: vec![],
905            time_window_size: 3600,
906            max_file_size: None,
907        };
908
909        let merge_output = compactor
910            .merge_ssts(&compaction_region, picker_output)
911            .await
912            .unwrap();
913
914        // Outputs 0 and 2 succeeded (1 + 2 = 3 files added).
915        assert_eq!(merge_output.files_to_add.len(), 3);
916        // Only inputs from successful merges should be removed.
917        assert_eq!(merge_output.files_to_remove.len(), 2);
918
919        let removed_ids: Vec<_> = merge_output
920            .files_to_remove
921            .iter()
922            .map(|f| f.file_id)
923            .collect();
924        assert!(removed_ids.contains(&input_meta_0.file_id));
925        assert!(removed_ids.contains(&input_meta_2.file_id));
926        // The failed output's input must NOT be removed.
927        assert!(!removed_ids.contains(&input_meta_1.file_id));
928    }
929
930    #[tokio::test]
931    async fn test_all_outputs_succeed() {
932        common_telemetry::init_default_ut_logging();
933
934        let compaction_region = new_test_compaction_region().await;
935        let input_meta = dummy_file_meta();
936        let output_meta = vec![dummy_file_meta()];
937
938        let merger = MockMerger::new(vec![Ok(output_meta.clone())]);
939        let compactor = DefaultCompactor::with_merger(merger);
940
941        let picker_output = PickerOutput {
942            outputs: vec![CompactionOutput {
943                output_level: 1,
944                inputs: vec![new_file_handle(input_meta.clone())],
945                filter_deleted: false,
946                output_time_range: None,
947            }],
948            expired_ssts: vec![],
949            time_window_size: 3600,
950            max_file_size: None,
951        };
952
953        let merge_output = compactor
954            .merge_ssts(&compaction_region, picker_output)
955            .await
956            .unwrap();
957
958        assert_eq!(merge_output.files_to_add.len(), 1);
959        assert_eq!(merge_output.files_to_add[0].file_id, output_meta[0].file_id);
960        assert_eq!(merge_output.files_to_remove.len(), 1);
961        assert_eq!(merge_output.files_to_remove[0].file_id, input_meta.file_id);
962    }
963
964    #[tokio::test]
965    async fn test_expired_ssts_always_removed() {
966        common_telemetry::init_default_ut_logging();
967
968        let compaction_region = new_test_compaction_region().await;
969        let input_meta = dummy_file_meta();
970        let expired_meta = dummy_file_meta();
971
972        // The single merge output fails, but expired SSTs should still be removed.
973        let merger = MockMerger::new(vec![Err(error::InvalidMetaSnafu {
974            reason: "fail".to_string(),
975        }
976        .build())]);
977        let compactor = DefaultCompactor::with_merger(merger);
978
979        let picker_output = PickerOutput {
980            outputs: vec![CompactionOutput {
981                output_level: 1,
982                inputs: vec![new_file_handle(input_meta.clone())],
983                filter_deleted: false,
984                output_time_range: None,
985            }],
986            expired_ssts: vec![new_file_handle(expired_meta.clone())],
987            time_window_size: 3600,
988            max_file_size: None,
989        };
990
991        let merge_output = compactor
992            .merge_ssts(&compaction_region, picker_output)
993            .await
994            .unwrap();
995
996        // No files added (merge failed).
997        assert!(merge_output.files_to_add.is_empty());
998        // Only the expired SST should be in files_to_remove (not the failed merge's input).
999        assert_eq!(merge_output.files_to_remove.len(), 1);
1000        assert_eq!(
1001            merge_output.files_to_remove[0].file_id,
1002            expired_meta.file_id
1003        );
1004    }
1005
1006    #[derive(Clone)]
1007    struct BlockingMerger {
1008        call_idx: Arc<AtomicUsize>,
1009    }
1010
1011    #[async_trait::async_trait]
1012    impl SstMerger for BlockingMerger {
1013        async fn merge_single_output(
1014            &self,
1015            _compaction_region: CompactionRegion,
1016            _output: CompactionOutput,
1017            _write_opts: WriteOptions,
1018        ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
1019            self.call_idx.fetch_add(1, Ordering::SeqCst);
1020            std::future::pending().await
1021        }
1022    }
1023
1024    #[tokio::test(flavor = "multi_thread")]
1025    async fn test_merge_ssts_cancels_spawned_tasks() {
1026        common_telemetry::init_default_ut_logging();
1027
1028        let mut compaction_region = new_test_compaction_region().await;
1029        compaction_region.max_parallelism = 2;
1030
1031        let cancel_handle = Arc::new(CancellationHandle::default());
1032        let call_idx = Arc::new(AtomicUsize::new(0));
1033        let compactor = DefaultCompactor {
1034            merger: BlockingMerger {
1035                call_idx: call_idx.clone(),
1036            },
1037            cancel_handle: cancel_handle.clone(),
1038            uncommitted: None,
1039        };
1040
1041        let picker_output = PickerOutput {
1042            outputs: vec![
1043                CompactionOutput {
1044                    output_level: 1,
1045                    inputs: vec![new_file_handle(dummy_file_meta())],
1046                    filter_deleted: false,
1047                    output_time_range: None,
1048                },
1049                CompactionOutput {
1050                    output_level: 1,
1051                    inputs: vec![new_file_handle(dummy_file_meta())],
1052                    filter_deleted: false,
1053                    output_time_range: None,
1054                },
1055                CompactionOutput {
1056                    output_level: 1,
1057                    inputs: vec![new_file_handle(dummy_file_meta())],
1058                    filter_deleted: false,
1059                    output_time_range: None,
1060                },
1061            ],
1062            expired_ssts: vec![],
1063            time_window_size: 3600,
1064            max_file_size: None,
1065        };
1066
1067        let task = tokio::spawn(async move {
1068            compactor
1069                .merge_ssts(&compaction_region, picker_output)
1070                .await
1071        });
1072
1073        sleep(Duration::from_millis(100)).await;
1074        cancel_handle.cancel();
1075
1076        let merge_output = task
1077            .await
1078            .expect("merge_ssts should stop after cancellation")
1079            .unwrap();
1080
1081        let started = call_idx.load(Ordering::SeqCst);
1082
1083        assert!(merge_output.files_to_add.is_empty());
1084        assert!(merge_output.files_to_remove.is_empty());
1085        assert_eq!(started, 2);
1086    }
1087}