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::NonZeroU64;
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, SequenceNumber};
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::{MergeMode, 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::{FileHandle, 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    /// Lower bound of pending memtable sequences from the same Version as `ssts`.
77    /// None means no barrier; Some(0) conservatively represents an unknown bound.
78    /// Keep only the value so background compaction cannot pin flushed memtables.
79    pub(crate) memtable_min_sequence: Option<SequenceNumber>,
80    /// Inferred compaction time window.
81    pub(crate) compaction_time_window: Option<Duration>,
82}
83
84impl From<VersionRef> for CompactionVersion {
85    fn from(value: VersionRef) -> Self {
86        Self {
87            metadata: value.metadata.clone(),
88            options: value.options.clone(),
89            ssts: value.ssts.clone(),
90            memtable_min_sequence: if value.options.merge_mode() == MergeMode::LastNonNull {
91                value.memtables.min_sequence()
92            } else {
93                None
94            },
95            compaction_time_window: value.compaction_time_window,
96        }
97    }
98}
99
100/// CompactionRegion represents a region that needs to be compacted.
101/// It's the subset of MitoRegion.
102#[derive(Clone)]
103pub struct CompactionRegion {
104    pub region_id: RegionId,
105    pub region_options: RegionOptions,
106
107    pub(crate) engine_config: Arc<MitoConfig>,
108    pub(crate) region_metadata: RegionMetadataRef,
109    pub(crate) cache_manager: CacheManagerRef,
110    /// Access layer to get the table path and path type.
111    pub access_layer: AccessLayerRef,
112    pub(crate) manifest_ctx: Arc<ManifestContext>,
113    pub(crate) current_version: CompactionVersion,
114    pub(crate) file_purger: Option<Arc<LocalFilePurger>>,
115    pub(crate) ttl: Option<TimeToLive>,
116
117    /// Controls the parallelism of this compaction task. Default is 1.
118    ///
119    /// The parallel is inside this compaction task, not across different compaction tasks.
120    /// It can be different windows of the same compaction task or something like this.
121    pub max_parallelism: usize,
122
123    pub(crate) plugins: Plugins,
124}
125
126/// OpenCompactionRegionRequest represents the request to open a compaction region.
127#[derive(Clone)]
128pub struct OpenCompactionRegionRequest {
129    pub region_id: RegionId,
130    pub table_dir: String,
131    pub path_type: PathType,
132    pub region_options: RegionOptions,
133    pub max_parallelism: usize,
134    /// Plugins for the compaction region, used to look up the [`RegionHook`](crate::engine::region_hook::RegionHook).
135    pub plugins: Plugins,
136}
137
138/// Open a compaction region from a compaction request.
139/// It's simple version of RegionOpener::open().
140pub async fn open_compaction_region(
141    req: &OpenCompactionRegionRequest,
142    mito_config: &MitoConfig,
143    object_store_manager: ObjectStoreManagerRef,
144    ttl_provider: Either<TimeToLive, SchemaMetadataManagerRef>,
145) -> Result<CompactionRegion> {
146    let object_store = {
147        let name = &req.region_options.storage;
148        if let Some(name) = name {
149            object_store_manager
150                .find(name)
151                .with_context(|| ObjectStoreNotFoundSnafu {
152                    object_store: name.clone(),
153                })?
154        } else {
155            object_store_manager.default_object_store()
156        }
157    };
158
159    let access_layer = {
160        let puffin_manager_factory = PuffinManagerFactory::new(
161            &mito_config.index.aux_path,
162            mito_config.index.staging_size.as_bytes(),
163            Some(mito_config.index.write_buffer_size.as_bytes() as _),
164            mito_config.index.staging_ttl,
165        )
166        .await?;
167        let intermediate_manager =
168            IntermediateManager::init_fs(mito_config.index.aux_path.clone()).await?;
169
170        Arc::new(AccessLayer::new(
171            &req.table_dir,
172            req.path_type,
173            object_store.clone(),
174            puffin_manager_factory,
175            intermediate_manager,
176        ))
177    };
178
179    let manifest_manager = {
180        let region_dir = region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type);
181        let region_manifest_options =
182            RegionManifestOptions::new(mito_config, &region_dir, object_store);
183
184        RegionManifestManager::open(region_manifest_options, &Default::default())
185            .await?
186            .with_context(|| EmptyRegionDirSnafu {
187                region_id: req.region_id,
188                region_dir: region_dir_from_table_dir(&req.table_dir, req.region_id, req.path_type),
189            })?
190    };
191
192    let manifest = manifest_manager.manifest();
193    let region_metadata = manifest.metadata.clone();
194    let hook: Option<RegionHookRef> = req.plugins.get();
195    let manifest_ctx = Arc::new(ManifestContext::new(
196        manifest_manager,
197        RegionRoleState::Leader(RegionLeaderState::Writable),
198        hook,
199    ));
200
201    let file_purger = {
202        let purge_scheduler = Arc::new(LocalScheduler::new(mito_config.max_background_purges));
203        Arc::new(LocalFilePurger::new(
204            purge_scheduler.clone(),
205            access_layer.clone(),
206            None,
207        ))
208    };
209
210    let current_version = {
211        let mut ssts = SstVersion::new(region_metadata.clone());
212        ssts.add_files(file_purger.clone(), manifest.files.values().cloned());
213        CompactionVersion {
214            metadata: region_metadata.clone(),
215            options: req.region_options.clone(),
216            ssts: Arc::new(ssts),
217            // Remote execution uses an already picked output. A manifest-only
218            // opener cannot prove that a new LastNonNull pick is memtable-safe.
219            memtable_min_sequence: Some(0),
220            compaction_time_window: manifest.compaction_time_window,
221        }
222    };
223
224    let ttl = match ttl_provider {
225        // Use the specified ttl.
226        Either::Left(ttl) => ttl,
227        // Get the ttl from the schema metadata manager.
228        Either::Right(schema_metadata_manager) => {
229            let (_, ttl) =
230                find_dynamic_options(req.region_id, &req.region_options, &schema_metadata_manager)
231                    .await
232                    .unwrap_or_else(|e| {
233                        warn!(e; "Failed to get ttl for region: {}", region_metadata.region_id);
234                        (
235                            crate::region::options::CompactionOptions::default(),
236                            TimeToLive::default(),
237                        )
238                    });
239            ttl
240        }
241    };
242
243    Ok(CompactionRegion {
244        region_id: req.region_id,
245        region_options: req.region_options.clone(),
246        engine_config: Arc::new(mito_config.clone()),
247        region_metadata: region_metadata.clone(),
248        cache_manager: Arc::new(CacheManager::default()),
249        access_layer,
250        manifest_ctx,
251        current_version,
252        file_purger: Some(file_purger),
253        ttl: Some(ttl),
254        max_parallelism: req.max_parallelism,
255        plugins: req.plugins.clone(),
256    })
257}
258
259impl CompactionRegion {
260    /// Get the file purger of the compaction region.
261    pub fn file_purger(&self) -> Option<Arc<LocalFilePurger>> {
262        self.file_purger.clone()
263    }
264
265    /// Stop the file purger scheduler of the compaction region.
266    pub async fn stop_purger_scheduler(&self) -> Result<()> {
267        if let Some(file_purger) = &self.file_purger {
268            file_purger.stop_scheduler().await
269        } else {
270            Ok(())
271        }
272    }
273
274    /// Fires [`RegionHook::on_sst_files_written`] for the freshly-merged SST
275    /// files in `merge_output`. Shared by both compaction paths, it must run
276    /// before [`Compactor::update_manifest`], whose `on_manifest_updated`
277    /// drains the per-region state this hook populates.
278    pub async fn invoke_sst_hook(&self, merge_output: &MergeOutput) {
279        let Some(hook) = self.plugins.get::<RegionHookRef>() else {
280            return;
281        };
282
283        // Remote compaction deserializes `MergeOutput` over gRPC, where
284        // `sst_infos` is `#[serde(skip)]`. Rebuild each `SstInfo` from the
285        // matching `FileMeta` so the hook observes real ids/rows/sizes rather
286        // than zeros — see [`sst_info_from_file_meta`] for what stays default.
287        let synthesized: Vec<SstInfo>;
288        let infos: &[SstInfo] = if merge_output.sst_infos.is_empty() {
289            synthesized = merge_output
290                .files_to_add
291                .iter()
292                .map(sst_info_from_file_meta)
293                .collect();
294            &synthesized
295        } else {
296            // `sst_infos` and `files_to_add` are documented as 1:1. If they
297            // ever diverge, `zip` below would silently truncate; warn so the
298            // mismatch is observable rather than dropping files from the hook.
299            if merge_output.sst_infos.len() != merge_output.files_to_add.len() {
300                warn!(
301                    "sst_infos length ({}) does not match files_to_add length ({}) for region {}",
302                    merge_output.sst_infos.len(),
303                    merge_output.files_to_add.len(),
304                    self.region_id
305                );
306            }
307            &merge_output.sst_infos
308        };
309
310        let files: Vec<SstFileInfo<'_>> = merge_output
311            .files_to_add
312            .iter()
313            .zip(infos)
314            .map(|(meta, info)| SstFileInfo {
315                sst_info_ref: info,
316                file_meta: meta,
317            })
318            .collect();
319        hook.on_sst_files_written(self.region_id, &self.region_metadata, &files)
320            .await;
321    }
322}
323
324/// Builds an [`SstInfo`] for the remote-compaction path by copying the scalar
325/// fields that [`FileMeta`] also carries.
326///
327/// Remote compaction runs off-process and ships `MergeOutput` over the wire;
328/// [`MergeOutput::sst_infos`] is `#[serde(skip)]` because the parquet footer
329/// (`ParquetMetaData`) is not serde-serializable. To avoid feeding the hook
330/// default (zero) values for real SSTs, this rebuilds the seven fields `FileMeta`
331/// and `SstInfo` share: `file_id`, `time_range`, `file_size`,
332/// `max_row_group_uncompressed_size`, `num_rows`, `num_row_groups`, `num_series`.
333///
334/// Two fields stay default and **cannot be recovered on the datanode**:
335/// - `file_metadata` — the parquet footer, whose column min/max/null-count
336///   statistics are required by hooks building richer artifacts (e.g. an Iceberg
337///   manifest). That data is produced by the compactor's writer and exists only
338///   in the freshly-written SST on object storage.
339/// - `index_metadata` — `FileMeta` tracks indexes via `available_indexes` /
340///   `indexes`, not as an [`IndexOutput`].
341///
342/// A hook that needs column stats must fetch the footer from object storage
343/// (the [`CompactionRegion`]'s `access_layer` reaches the store), or the footer
344/// must be shipped over the wire (revisit the `#[serde(skip)]` on `sst_infos`).
345/// `num_rows` may read `0` for legacy `FileMeta`s where the count is unknown.
346fn sst_info_from_file_meta(meta: &FileMeta) -> SstInfo {
347    SstInfo {
348        file_id: meta.file_id,
349        time_range: meta.time_range,
350        file_size: meta.file_size,
351        max_row_group_uncompressed_size: meta.max_row_group_uncompressed_size,
352        num_rows: meta.num_rows as usize,
353        num_row_groups: meta.num_row_groups,
354        num_series: meta.num_series,
355        ..Default::default()
356    }
357}
358
359/// `[MergeOutput]` represents the output of merging SST files.
360#[derive(Default, Debug, Serialize, Deserialize)]
361pub struct MergeOutput {
362    pub files_to_add: Vec<FileMeta>,
363    pub files_to_remove: Vec<FileMeta>,
364    pub compaction_time_window: Option<i64>,
365    #[serde(skip)]
366    pub sst_infos: Vec<SstInfo>,
367}
368
369impl MergeOutput {
370    pub fn is_empty(&self) -> bool {
371        self.files_to_add.is_empty() && self.files_to_remove.is_empty()
372    }
373
374    pub fn input_file_size(&self) -> u64 {
375        self.files_to_remove.iter().map(|f| f.file_size).sum()
376    }
377
378    pub fn output_file_size(&self) -> u64 {
379        self.files_to_add.iter().map(|f| f.file_size).sum()
380    }
381}
382
383/// Compactor is the trait that defines the compaction logic.
384#[async_trait::async_trait]
385pub trait Compactor: Send + Sync + 'static {
386    /// Merge SST files for a region.
387    async fn merge_ssts(
388        &self,
389        compaction_region: &CompactionRegion,
390        picker_output: PickerOutput,
391    ) -> Result<MergeOutput>;
392
393    /// Update the manifest after merging SST files.
394    async fn update_manifest(
395        &self,
396        compaction_region: &CompactionRegion,
397        merge_output: MergeOutput,
398    ) -> Result<(RegionEdit, ManifestVersion)>;
399}
400
401/// Trait for merging a single compaction output into SST files.
402///
403/// This is extracted from `DefaultCompactor` to allow injecting mock
404/// implementations in tests.
405#[async_trait::async_trait]
406pub trait SstMerger: Send + Sync + 'static {
407    async fn merge_single_output(
408        &self,
409        compaction_region: CompactionRegion,
410        output: CompactionOutput,
411        write_opts: WriteOptions,
412    ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)>;
413}
414
415/// The production [`SstMerger`] that reads, merges, and writes SST files.
416#[derive(Clone)]
417pub struct DefaultSstMerger;
418
419/// File-level sequence metadata is independent of physical row encoding.
420struct OutputSequenceMetadata {
421    file_sequence: Option<NonZeroU64>,
422    exact_sequence_trusted: bool,
423}
424
425/// Compaction does not admit new data, so it inherits the input sequence bounds
426/// without allocating a new admission marker. Retaining physical row sequences
427/// does not by itself restore exact-read capability for untrusted inputs.
428fn output_sequence_metadata(
429    region: &CompactionRegion,
430    inputs: &[FileHandle],
431) -> OutputSequenceMetadata {
432    let exact_sequence_trusted = region.region_options.preserve_row_sequence
433        && inputs
434            .iter()
435            .all(|f| f.is_effective_target_sequence_trusted(region.region_id));
436    OutputSequenceMetadata {
437        file_sequence: known_max_input_sequence(inputs),
438        exact_sequence_trusted,
439    }
440}
441
442/// Computes the maximum target-domain sequence bound of the output of merging
443/// `inputs`. A bound can be a row maximum or an admission marker; compaction
444/// must preserve both for sequence-based pruning. Unknown bounds must remain
445/// unknown rather than being replaced with a partial maximum.
446fn known_max_input_sequence(inputs: &[FileHandle]) -> Option<NonZeroU64> {
447    let mut max: Option<NonZeroU64> = None;
448    for input in inputs {
449        let sequence = input.meta_ref().sequence?;
450        max = Some(max.map_or(sequence, |current| current.max(sequence)));
451    }
452    max
453}
454
455#[async_trait::async_trait]
456impl SstMerger for DefaultSstMerger {
457    async fn merge_single_output(
458        &self,
459        compaction_region: CompactionRegion,
460        output: CompactionOutput,
461        write_opts: WriteOptions,
462    ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
463        let region_id = compaction_region.region_id;
464        let storage = compaction_region.region_options.storage.clone();
465        let index_options = compaction_region
466            .current_version
467            .options
468            .index_options
469            .clone();
470        let append_mode = compaction_region.current_version.options.append_mode;
471        let merge_mode = compaction_region.current_version.options.merge_mode();
472        let flat_format = compaction_region
473            .region_options
474            .sst_format
475            .map(|format| format == FormatType::Flat)
476            .unwrap_or(compaction_region.engine_config.default_flat_format);
477
478        let index_config = compaction_region.engine_config.index.clone();
479        let inverted_index_config = compaction_region.engine_config.inverted_index.clone();
480        let fulltext_index_config = compaction_region.engine_config.fulltext_index.clone();
481        let bloom_filter_index_config = compaction_region.engine_config.bloom_filter_index.clone();
482        #[cfg(feature = "vector_index")]
483        let vector_index_config = compaction_region.engine_config.vector_index.clone();
484
485        let input_file_names = output
486            .inputs
487            .iter()
488            .map(|f| f.file_id().to_string())
489            .join(",");
490        let sequence_metadata = output_sequence_metadata(&compaction_region, &output.inputs);
491        let builder = CompactionSstReaderBuilder {
492            metadata: compaction_region.region_metadata.clone(),
493            sst_layer: compaction_region.access_layer.clone(),
494            cache: compaction_region.cache_manager.clone(),
495            inputs: &output.inputs,
496            append_mode,
497            filter_deleted: output.filter_deleted,
498            time_range: output.output_time_range,
499            merge_mode,
500        };
501        let source = builder.build_flat_sst_reader().await?;
502
503        let mut metrics = Metrics::new(WriteType::Compaction);
504        let region_metadata = compaction_region.region_metadata.clone();
505        let sst_infos = compaction_region
506            .access_layer
507            .write_sst(
508                SstWriteRequest {
509                    op_type: OperationType::Compact,
510                    metadata: region_metadata.clone(),
511                    source,
512                    cache_manager: compaction_region.cache_manager.clone(),
513                    storage,
514                    // Readers resolve file overrides before merge/dedup. Replacing
515                    // their effective sequences here could promote old rows above
516                    // versions in SSTs that were not part of this merge.
517                    max_sequence: None,
518                    sst_write_format: if flat_format {
519                        FormatType::Flat
520                    } else {
521                        FormatType::PrimaryKey
522                    },
523                    preserve_row_sequence: sequence_metadata.exact_sequence_trusted,
524                    index_options,
525                    index_config,
526                    inverted_index_config,
527                    fulltext_index_config,
528                    bloom_filter_index_config,
529                    #[cfg(feature = "vector_index")]
530                    vector_index_config,
531                },
532                &write_opts,
533                &mut metrics,
534            )
535            .await?;
536        // Convert partition expression once outside the map
537        let partition_expr = match &region_metadata.partition_expr {
538            None => None,
539            Some(json_str) if json_str.is_empty() => None,
540            Some(json_str) => PartitionExpr::from_json_str(json_str).with_context(|_| {
541                InvalidPartitionExprSnafu {
542                    expr: json_str.clone(),
543                }
544            })?,
545        };
546
547        let output_files = sst_infos
548            .iter()
549            .map(|sst_info| {
550                let pk_range = sst_info
551                    .file_metadata
552                    .as_ref()
553                    .and_then(|meta| extract_primary_key_range(meta, &region_metadata));
554                let (primary_key_min, primary_key_max) = match pk_range {
555                    Some((min, max)) => (Some(min), Some(max)),
556                    None => (None, None),
557                };
558
559                FileMeta {
560                    region_id,
561                    file_id: sst_info.file_id,
562                    time_range: sst_info.time_range,
563                    level: output.output_level,
564                    file_size: sst_info.file_size,
565                    max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
566                    available_indexes: sst_info.index_metadata.build_available_indexes(),
567                    indexes: sst_info.index_metadata.build_indexes(),
568                    index_file_size: sst_info.index_metadata.file_size,
569                    index_version: 0,
570                    num_rows: sst_info.num_rows as u64,
571                    num_row_groups: sst_info.num_row_groups,
572                    sequence: sequence_metadata.file_sequence,
573                    partition_expr: partition_expr.clone(),
574                    num_series: sst_info.num_series,
575                    primary_key_min,
576                    primary_key_max,
577                    preserve_row_sequence: sequence_metadata.exact_sequence_trusted,
578                }
579            })
580            .collect::<Vec<_>>();
581        let output_file_names = output_files.iter().map(|f| f.file_id.to_string()).join(",");
582        info!(
583            "Region {} compaction inputs: [{}], outputs: [{}], flat_format: {}, metrics: {:?}",
584            region_id, input_file_names, output_file_names, flat_format, metrics
585        );
586        metrics.observe();
587        Ok((output_files, sst_infos.into_iter().collect()))
588    }
589}
590
591/// DefaultCompactor is the default implementation of Compactor.
592///
593/// It is parameterized by an [`SstMerger`] to allow injecting mock
594/// implementations in tests.
595pub struct DefaultCompactor<M = DefaultSstMerger> {
596    merger: M,
597    cancel_handle: Arc<CancellationHandle>,
598    uncommitted: Option<UncommittedSsts>,
599}
600
601#[cfg(test)]
602impl<M: SstMerger> DefaultCompactor<M> {
603    pub fn with_merger(merger: M) -> Self {
604        Self {
605            merger,
606            cancel_handle: Arc::new(CancellationHandle::default()),
607            uncommitted: None,
608        }
609    }
610}
611
612impl DefaultCompactor {
613    /// Creates a new `DefaultCompactor` with the given cancel handle and no
614    /// uncommitted SSTs tracking.
615    ///
616    /// This is the public entry point for external crates that want a
617    /// cancellable compactor without opting into uncommitted-SST tracking.
618    pub fn new_with_cancel_handle(cancel_handle: Arc<CancellationHandle>) -> Self {
619        Self {
620            merger: DefaultSstMerger,
621            cancel_handle,
622            uncommitted: None,
623        }
624    }
625
626    pub(crate) fn with_cancel_handle(
627        cancel_handle: Arc<CancellationHandle>,
628        uncommitted: UncommittedSsts,
629    ) -> Self {
630        Self {
631            merger: DefaultSstMerger,
632            cancel_handle,
633            uncommitted: Some(uncommitted),
634        }
635    }
636}
637
638#[async_trait::async_trait]
639impl<M: SstMerger> Compactor for DefaultCompactor<M>
640where
641    M: Clone,
642{
643    async fn merge_ssts(
644        &self,
645        compaction_region: &CompactionRegion,
646        mut picker_output: PickerOutput,
647    ) -> Result<MergeOutput> {
648        let internal_parallelism = compaction_region.max_parallelism.max(1);
649        let compaction_time_window = picker_output.time_window_size;
650        let region_id = compaction_region.region_id;
651
652        // Build tasks along with their input file metas so we can track which
653        // inputs correspond to each task.
654        let mut tasks: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(picker_output.outputs.len());
655
656        for output in picker_output.outputs.drain(..) {
657            let inputs_to_remove: Vec<_> =
658                output.inputs.iter().map(|f| f.meta_ref().clone()).collect();
659            let write_opts = WriteOptions {
660                write_buffer_size: compaction_region.engine_config.sst_write_buffer_size,
661                max_file_size: picker_output.max_file_size,
662                row_group_size: compaction_region.region_options.row_group_size(),
663                float_field_encoding: compaction_region.region_options.float_field_encoding,
664            };
665            let merger = self.merger.clone();
666            let compaction_region = compaction_region.clone();
667            let uncommitted = self.uncommitted.clone();
668            let fut = async move {
669                let result = merger
670                    .merge_single_output(compaction_region, output, write_opts)
671                    .await;
672                if let (Some(uncommitted), Ok((_, infos))) = (&uncommitted, &result) {
673                    uncommitted.track(infos);
674                }
675                result
676            };
677            tasks.push((inputs_to_remove, fut));
678        }
679
680        let hook: Option<RegionHookRef> = compaction_region.plugins.get();
681        let mut output_files = Vec::with_capacity(tasks.len());
682        let mut all_sst_infos: Vec<SstInfo> = Vec::new();
683        let mut compacted_inputs = Vec::with_capacity(
684            tasks.iter().map(|(inputs, _)| inputs.len()).sum::<usize>()
685                + picker_output.expired_ssts.len(),
686        );
687
688        while !tasks.is_empty() {
689            let mut chunk: Vec<(Vec<FileMeta>, _)> = Vec::with_capacity(internal_parallelism);
690            for _ in 0..internal_parallelism {
691                if let Some(task) = tasks.pop() {
692                    chunk.push(task);
693                }
694            }
695            let mut spawned: Vec<_> = chunk
696                .into_iter()
697                .map(|(inputs, fut)| {
698                    let handle = common_runtime::spawn_compact(fut);
699                    (inputs, handle)
700                })
701                .collect();
702
703            while let Some((inputs, mut handle)) = spawned.pop() {
704                let abort_handle = handle.abort_handle();
705                match CancellableFuture::new(&mut handle, self.cancel_handle.clone()).await {
706                    Ok(Ok(Ok((files, infos)))) => {
707                        output_files.extend(files);
708                        if hook.is_some() {
709                            all_sst_infos.extend(infos);
710                        }
711                        compacted_inputs.extend(inputs);
712                    }
713                    Ok(Ok(Err(e))) => {
714                        warn!(
715                            e; "Failed to merge compaction output for region: {}, inputs: [{}]",
716                            region_id,
717                            inputs.iter().map(|f| f.file_id.to_string()).join(",")
718                        );
719                    }
720                    Ok(Err(e)) => {
721                        warn!(
722                            "Region {} compaction task join error for inputs: [{}], skipping: {}",
723                            region_id,
724                            inputs.iter().map(|f| f.file_id.to_string()).join(","),
725                            e
726                        );
727                        // If the cancel handle is cancelled,
728                        // cancel the remaining tasks before returns the error.
729                        if self.cancel_handle.is_cancelled() {
730                            abort_handle.abort();
731                            for (_, handle) in &spawned {
732                                handle.abort();
733                            }
734                            for (_, handle) in spawned {
735                                let _ = handle.await;
736                            }
737                        }
738                        return Err(e).context(error::JoinSnafu);
739                    }
740                    Err(_) => {
741                        debug!(
742                            "Compaction merge cancelled for region: {}, aborting remaining {} spawned tasks",
743                            region_id,
744                            spawned.len(),
745                        );
746                        abort_handle.abort();
747                        for (_, handle) in &spawned {
748                            handle.abort();
749                        }
750                        let _ = handle.await;
751                        for (_, handle) in spawned {
752                            let _ = handle.await;
753                        }
754                        break;
755                    }
756                }
757            }
758
759            if self.cancel_handle.is_cancelled() {
760                info!("Compaction merge cancelled for region: {}", region_id);
761                break;
762            }
763        }
764
765        // Include expired SSTs in removals — these don't depend on merge success.
766        compacted_inputs.extend(
767            picker_output
768                .expired_ssts
769                .iter()
770                .map(|f| f.meta_ref().clone()),
771        );
772
773        Ok(MergeOutput {
774            files_to_add: output_files,
775            files_to_remove: compacted_inputs,
776            compaction_time_window: Some(compaction_time_window),
777            sst_infos: all_sst_infos,
778        })
779    }
780
781    async fn update_manifest(
782        &self,
783        compaction_region: &CompactionRegion,
784        merge_output: MergeOutput,
785    ) -> Result<(RegionEdit, ManifestVersion)> {
786        // Write region edit to manifest.
787        let edit = RegionEdit {
788            files_to_add: merge_output.files_to_add,
789            files_to_remove: merge_output.files_to_remove,
790            // Use current timestamp as the edit timestamp.
791            timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
792            compaction_time_window: merge_output
793                .compaction_time_window
794                .map(|seconds| Duration::from_secs(seconds as u64)),
795            flushed_entry_id: None,
796            flushed_sequence: None,
797            committed_sequence: None,
798        };
799
800        let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone()));
801        // TODO: We might leak files if we fail to update manifest. We can add a cleanup task to remove them later.
802        let manifest_version = compaction_region
803            .manifest_ctx
804            .update_manifest_for_compaction(action_list)
805            .await?;
806
807        Ok((edit, manifest_version))
808    }
809}
810
811#[cfg(test)]
812mod tests {
813    use std::sync::atomic::{AtomicUsize, Ordering};
814    use std::sync::{Arc, Mutex};
815    use std::time::Duration;
816
817    use store_api::storage::{FileId, RegionId};
818    use tokio::time::sleep;
819
820    use super::{DefaultCompactor, *};
821    use crate::cache::CacheManager;
822    use crate::compaction::picker::PickerOutput;
823    use crate::error::Result;
824    use crate::sst::file::FileHandle;
825    use crate::sst::file_purger::NoopFilePurger;
826    use crate::sst::version::SstVersion;
827    use crate::test_util::memtable_util::metadata_for_test;
828    use crate::test_util::scheduler_util::SchedulerEnv;
829
830    fn dummy_file_meta() -> FileMeta {
831        FileMeta {
832            region_id: RegionId::new(1, 1),
833            file_id: FileId::random(),
834            file_size: 100,
835            ..Default::default()
836        }
837    }
838
839    fn new_file_handle(meta: FileMeta) -> FileHandle {
840        FileHandle::new(meta, Arc::new(NoopFilePurger))
841    }
842
843    #[test]
844    fn test_effective_target_trust_and_max_sequence() {
845        let target = RegionId::new(1, 1);
846        let local_trusted = new_file_handle(FileMeta {
847            region_id: target,
848            sequence: NonZeroU64::new(3),
849            preserve_row_sequence: true,
850            ..dummy_file_meta()
851        });
852        let foreign_untrusted_marker = new_file_handle(FileMeta {
853            region_id: RegionId::new(1, 2),
854            sequence: NonZeroU64::new(9),
855            preserve_row_sequence: false,
856            ..dummy_file_meta()
857        });
858        let local_untrusted = new_file_handle(FileMeta {
859            region_id: target,
860            sequence: NonZeroU64::new(11),
861            preserve_row_sequence: false,
862            ..dummy_file_meta()
863        });
864        assert!(local_trusted.is_effective_target_sequence_trusted(target));
865        assert!(foreign_untrusted_marker.is_effective_target_sequence_trusted(target));
866        assert!(!local_untrusted.is_effective_target_sequence_trusted(target));
867        assert_eq!(
868            NonZeroU64::new(9),
869            known_max_input_sequence(&[local_trusted, foreign_untrusted_marker])
870        );
871    }
872
873    #[test]
874    fn test_known_max_input_sequence() {
875        let meta_with = |sequence| FileMeta {
876            sequence,
877            ..dummy_file_meta()
878        };
879
880        let inputs = vec![
881            new_file_handle(meta_with(NonZeroU64::new(3))),
882            new_file_handle(meta_with(NonZeroU64::new(9))),
883            new_file_handle(meta_with(NonZeroU64::new(5))),
884        ];
885        assert_eq!(NonZeroU64::new(9), known_max_input_sequence(&inputs));
886
887        let inputs = vec![
888            new_file_handle(meta_with(NonZeroU64::new(9))),
889            new_file_handle(meta_with(None)),
890        ];
891        assert_eq!(None, known_max_input_sequence(&inputs));
892
893        assert_eq!(None, known_max_input_sequence(&[]));
894    }
895
896    #[tokio::test]
897    async fn test_output_sequence_metadata_preserves_trust_and_unknown_bounds() {
898        let mut region = new_test_compaction_region().await;
899        let local = region.region_id;
900        let foreign = RegionId::new(1, 2);
901        let file = |region_id, sequence: Option<u64>, preserve_row_sequence| {
902            new_file_handle(FileMeta {
903                region_id,
904                sequence: sequence.and_then(NonZeroU64::new),
905                preserve_row_sequence,
906                ..dummy_file_meta()
907            })
908        };
909        // Bounds are independent of row trust: known, mixed trust, unknown in
910        // either sequence domain, empty inputs, and the largest possible bound.
911        let cases = [
912            (
913                vec![file(local, Some(3), true), file(foreign, Some(9), false)],
914                Some(9),
915                true,
916            ),
917            (
918                vec![file(local, Some(3), true), file(local, Some(6), false)],
919                Some(6),
920                false,
921            ),
922            (
923                vec![file(foreign, Some(9), false), file(foreign, Some(11), true)],
924                Some(11),
925                true,
926            ),
927            (
928                vec![file(local, Some(3), true), file(local, None, true)],
929                None,
930                true,
931            ),
932            (
933                vec![file(local, None, false), file(local, Some(6), true)],
934                None,
935                false,
936            ),
937            (
938                vec![file(local, Some(3), true), file(foreign, None, true)],
939                None,
940                false,
941            ),
942            (vec![], None, true),
943            (
944                vec![file(local, Some(u64::MAX), false)],
945                Some(u64::MAX),
946                false,
947            ),
948        ];
949        for committed_sequence in [100, 101] {
950            region
951                .manifest_ctx
952                .manifest_manager
953                .write()
954                .await
955                .update(
956                    RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
957                        committed_sequence: Some(committed_sequence),
958                        files_to_add: vec![],
959                        files_to_remove: vec![],
960                        timestamp_ms: None,
961                        compaction_time_window: None,
962                        flushed_entry_id: None,
963                        flushed_sequence: None,
964                    })),
965                    false,
966                )
967                .await
968                .unwrap();
969            for preserve in [false, true] {
970                region.region_options.preserve_row_sequence = preserve;
971                for (inputs, known_bound, inputs_trusted) in &cases {
972                    let metadata = output_sequence_metadata(&region, inputs);
973                    let trusted = preserve && *inputs_trusted;
974                    assert_eq!(trusted, metadata.exact_sequence_trusted);
975                    assert_eq!(
976                        known_bound.and_then(NonZeroU64::new),
977                        metadata.file_sequence
978                    );
979                }
980            }
981        }
982    }
983
984    /// Independent outputs must be safe in either visibility order, including when
985    /// only one output succeeds. Merely committing both outputs together is not a
986    /// substitute for preserving version order across their separate merge streams.
987    #[rstest::rstest]
988    #[case(true)]
989    #[case(false)]
990    #[tokio::test]
991    async fn test_independent_outputs_preserve_delete_in_every_visibility_state(
992        #[case] flat_format: bool,
993        #[values(false, true)] unknown_put_bound: bool,
994    ) {
995        use datatypes::arrow::array::AsArray;
996        use datatypes::arrow::datatypes::TimestampMillisecondType;
997        use store_api::region_engine::RegionEngine;
998        use store_api::region_request::RegionRequest;
999
1000        use crate::compaction::CompactionOutput;
1001        use crate::compaction::compactor::{CompactionRegion, Compactor, DefaultCompactor};
1002        use crate::compaction::picker::PickerOutput;
1003        use crate::compaction::reader::CompactionSstReaderBuilder;
1004        use crate::engine::compaction_test::{delete_and_flush, put_and_flush};
1005        use crate::region::options::MergeMode;
1006        use crate::sst::file::FileHandle;
1007        use crate::sst::file_purger::NoopFilePurger;
1008        use crate::test_util::{CreateRequestBuilder, TestEnv, rows_schema};
1009
1010        let mut env = TestEnv::new().await;
1011        let region_id = RegionId::new(1, 1);
1012        let config = MitoConfig {
1013            default_flat_format: flat_format,
1014            min_compaction_interval: Duration::from_secs(3600),
1015            ..Default::default()
1016        };
1017        let engine = env.create_engine(config).await;
1018        let request = CreateRequestBuilder::new().build();
1019        let columns = rows_schema(&request);
1020        engine
1021            .handle_request(region_id, RegionRequest::Create(request))
1022            .await
1023            .unwrap();
1024        put_and_flush(&engine, region_id, &columns, 0..1).await;
1025        delete_and_flush(&engine, region_id, &columns, 0..1).await;
1026        put_and_flush(&engine, region_id, &columns, 10..11).await;
1027        let region = engine.get_region(region_id).unwrap();
1028        let version = region.version();
1029        let mut files: Vec<_> = version
1030            .ssts
1031            .levels()
1032            .iter()
1033            .flat_map(|level| level.files())
1034            .cloned()
1035            .collect();
1036        assert_eq!(3, files.len());
1037        files.sort_unstable_by_key(|file| file.meta_ref().sequence);
1038        if unknown_put_bound {
1039            // Model a legacy SST whose physical rows are readable but whose
1040            // file-level boundary is unknown. Merging cannot invent that bound.
1041            let mut meta = files[0].meta_ref().clone();
1042            meta.sequence = None;
1043            files[0] = FileHandle::new(meta, Arc::new(NoopFilePurger));
1044        }
1045        let compaction_region = CompactionRegion {
1046            region_id,
1047            region_options: version.options.clone(),
1048            engine_config: Arc::new(MitoConfig {
1049                default_flat_format: flat_format,
1050                ..Default::default()
1051            }),
1052            region_metadata: version.metadata.clone(),
1053            cache_manager: engine.cache_manager(),
1054            access_layer: region.access_layer.clone(),
1055            manifest_ctx: region.manifest_ctx.clone(),
1056            current_version: version.into(),
1057            file_purger: None,
1058            ttl: None,
1059            max_parallelism: 2,
1060            plugins: common_base::Plugins::new(),
1061        };
1062        let puts = vec![files[0].clone(), files[2].clone()];
1063        let deletes = vec![files[1].clone()];
1064        let picker_output = PickerOutput {
1065            outputs: [puts.clone(), deletes.clone()]
1066                .into_iter()
1067                .map(|inputs| CompactionOutput {
1068                    output_level: 1,
1069                    inputs,
1070                    filter_deleted: false,
1071                    output_time_range: None,
1072                })
1073                .collect(),
1074            expired_ssts: vec![],
1075            time_window_size: 3600,
1076            max_file_size: None,
1077        };
1078        let merged = DefaultCompactor::with_merger(DefaultSstMerger)
1079            .merge_ssts(&compaction_region, picker_output)
1080            .await
1081            .unwrap();
1082        assert_eq!(2, merged.files_to_add.len());
1083        let outputs: Vec<_> = merged
1084            .files_to_add
1085            .into_iter()
1086            .map(|meta| FileHandle::new(meta, Arc::new(NoopFilePurger)))
1087            .collect();
1088        let put_output = outputs
1089            .iter()
1090            .find(|file| file.meta_ref().num_rows == 2)
1091            .unwrap();
1092        let delete_output = outputs
1093            .iter()
1094            .find(|file| file.meta_ref().num_rows == 1)
1095            .unwrap();
1096        assert_eq!(
1097            if unknown_put_bound {
1098                None
1099            } else {
1100                NonZeroU64::new(3)
1101            },
1102            put_output.meta_ref().sequence,
1103        );
1104        assert_eq!(NonZeroU64::new(2), delete_output.meta_ref().sequence);
1105        assert!(
1106            outputs
1107                .iter()
1108                .all(|file| !file.meta_ref().preserve_row_sequence)
1109        );
1110
1111        // Cover the original snapshot, each partial success/cancellation state,
1112        // and both completed outputs using real persisted SSTs and merge readers.
1113        for (put_done, delete_done) in [(false, false), (true, false), (false, true), (true, true)]
1114        {
1115            let mut visible = if put_done {
1116                vec![put_output.clone()]
1117            } else {
1118                puts.clone()
1119            };
1120            visible.extend(if delete_done {
1121                vec![delete_output.clone()]
1122            } else {
1123                deletes.clone()
1124            });
1125            let mut reader = CompactionSstReaderBuilder {
1126                metadata: compaction_region.region_metadata.clone(),
1127                sst_layer: region.access_layer.clone(),
1128                cache: engine.cache_manager(),
1129                inputs: &visible,
1130                append_mode: false,
1131                filter_deleted: true,
1132                time_range: None,
1133                merge_mode: MergeMode::LastRow,
1134            }
1135            .build_flat_sst_reader()
1136            .await
1137            .unwrap();
1138            let mut timestamps = Vec::new();
1139            while let Some(batch) = reader.next_batch().await.unwrap() {
1140                timestamps.extend(
1141                    batch
1142                        .column_by_name("ts")
1143                        .unwrap()
1144                        .as_primitive::<TimestampMillisecondType>()
1145                        .values()
1146                        .iter()
1147                        .copied(),
1148                );
1149            }
1150            assert_eq!(
1151                vec![10_000],
1152                timestamps,
1153                "put_done={put_done}, delete_done={delete_done}"
1154            );
1155        }
1156    }
1157
1158    /// Build a minimal [`CompactionRegion`] suitable for tests where the
1159    /// [`SstMerger`] is mocked and never touches the access layer.
1160    async fn new_test_compaction_region() -> CompactionRegion {
1161        let env = SchedulerEnv::new().await;
1162        let metadata = metadata_for_test();
1163        let manifest_ctx = env.mock_manifest_context(metadata.clone()).await;
1164        CompactionRegion {
1165            region_id: RegionId::new(1, 1),
1166            region_options: RegionOptions::default(),
1167            engine_config: Arc::new(MitoConfig::default()),
1168            region_metadata: metadata.clone(),
1169            cache_manager: Arc::new(CacheManager::default()),
1170            access_layer: env.access_layer.clone(),
1171            manifest_ctx,
1172            current_version: CompactionVersion {
1173                metadata: metadata.clone(),
1174                options: RegionOptions::default(),
1175                ssts: Arc::new(SstVersion::new(metadata)),
1176                memtable_min_sequence: None,
1177                compaction_time_window: None,
1178            },
1179            file_purger: None,
1180            ttl: None,
1181            max_parallelism: 1,
1182            plugins: Plugins::new(),
1183        }
1184    }
1185
1186    /// An [`SstMerger`] that returns pre-configured results per call index.
1187    ///
1188    /// Call 0 gets `results[0]`, call 1 gets `results[1]`, etc.
1189    #[derive(Clone)]
1190    struct MockMerger {
1191        results: Arc<Mutex<Vec<Result<Vec<FileMeta>>>>>,
1192        call_idx: Arc<AtomicUsize>,
1193    }
1194
1195    impl MockMerger {
1196        fn new(results: Vec<Result<Vec<FileMeta>>>) -> Self {
1197            Self {
1198                results: Arc::new(Mutex::new(results)),
1199                call_idx: Arc::new(AtomicUsize::new(0)),
1200            }
1201        }
1202    }
1203
1204    #[async_trait::async_trait]
1205    impl SstMerger for MockMerger {
1206        async fn merge_single_output(
1207            &self,
1208            _compaction_region: CompactionRegion,
1209            _output: CompactionOutput,
1210            _write_opts: WriteOptions,
1211        ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
1212            let idx = self.call_idx.fetch_add(1, Ordering::SeqCst);
1213            match self.results.lock().unwrap().get(idx) {
1214                Some(Ok(files)) => Ok((files.clone(), Vec::new())),
1215                Some(Err(_)) => error::InvalidMetaSnafu {
1216                    reason: format!("simulated failure at index {idx}"),
1217                }
1218                .fail(),
1219                None => panic!("MockMerger: no result configured for call index {idx}"),
1220            }
1221        }
1222    }
1223
1224    #[tokio::test]
1225    async fn test_partial_merge_failure_collects_only_successful_outputs() {
1226        common_telemetry::init_default_ut_logging();
1227
1228        let compaction_region = new_test_compaction_region().await;
1229
1230        // Prepare 3 compaction outputs: output 0 and 2 succeed, output 1 fails.
1231        let input_meta_0 = dummy_file_meta();
1232        let input_meta_1 = dummy_file_meta();
1233        let input_meta_2 = dummy_file_meta();
1234
1235        let output_meta_0 = vec![dummy_file_meta()];
1236        let output_meta_2 = vec![dummy_file_meta(), dummy_file_meta()];
1237
1238        let merger = MockMerger::new(vec![
1239            Ok(output_meta_0.clone()),
1240            Err(error::InvalidMetaSnafu {
1241                reason: "boom".to_string(),
1242            }
1243            .build()),
1244            Ok(output_meta_2.clone()),
1245        ]);
1246        let compactor = DefaultCompactor::with_merger(merger);
1247
1248        let picker_output = PickerOutput {
1249            outputs: vec![
1250                CompactionOutput {
1251                    output_level: 1,
1252                    inputs: vec![new_file_handle(input_meta_0.clone())],
1253                    filter_deleted: false,
1254                    output_time_range: None,
1255                },
1256                CompactionOutput {
1257                    output_level: 1,
1258                    inputs: vec![new_file_handle(input_meta_1.clone())],
1259                    filter_deleted: false,
1260                    output_time_range: None,
1261                },
1262                CompactionOutput {
1263                    output_level: 1,
1264                    inputs: vec![new_file_handle(input_meta_2.clone())],
1265                    filter_deleted: false,
1266                    output_time_range: None,
1267                },
1268            ],
1269            expired_ssts: vec![],
1270            time_window_size: 3600,
1271            max_file_size: None,
1272        };
1273
1274        let merge_output = compactor
1275            .merge_ssts(&compaction_region, picker_output)
1276            .await
1277            .unwrap();
1278
1279        // Outputs 0 and 2 succeeded (1 + 2 = 3 files added).
1280        assert_eq!(merge_output.files_to_add.len(), 3);
1281        // Only inputs from successful merges should be removed.
1282        assert_eq!(merge_output.files_to_remove.len(), 2);
1283
1284        let removed_ids: Vec<_> = merge_output
1285            .files_to_remove
1286            .iter()
1287            .map(|f| f.file_id)
1288            .collect();
1289        assert!(removed_ids.contains(&input_meta_0.file_id));
1290        assert!(removed_ids.contains(&input_meta_2.file_id));
1291        // The failed output's input must NOT be removed.
1292        assert!(!removed_ids.contains(&input_meta_1.file_id));
1293    }
1294
1295    #[tokio::test]
1296    async fn test_all_outputs_succeed() {
1297        common_telemetry::init_default_ut_logging();
1298
1299        let compaction_region = new_test_compaction_region().await;
1300        let input_meta = dummy_file_meta();
1301        let output_meta = vec![dummy_file_meta()];
1302
1303        let merger = MockMerger::new(vec![Ok(output_meta.clone())]);
1304        let compactor = DefaultCompactor::with_merger(merger);
1305
1306        let picker_output = PickerOutput {
1307            outputs: vec![CompactionOutput {
1308                output_level: 1,
1309                inputs: vec![new_file_handle(input_meta.clone())],
1310                filter_deleted: false,
1311                output_time_range: None,
1312            }],
1313            expired_ssts: vec![],
1314            time_window_size: 3600,
1315            max_file_size: None,
1316        };
1317
1318        let merge_output = compactor
1319            .merge_ssts(&compaction_region, picker_output)
1320            .await
1321            .unwrap();
1322
1323        assert_eq!(merge_output.files_to_add.len(), 1);
1324        assert_eq!(merge_output.files_to_add[0].file_id, output_meta[0].file_id);
1325        assert_eq!(merge_output.files_to_remove.len(), 1);
1326        assert_eq!(merge_output.files_to_remove[0].file_id, input_meta.file_id);
1327    }
1328
1329    #[tokio::test]
1330    async fn test_expired_ssts_always_removed() {
1331        common_telemetry::init_default_ut_logging();
1332
1333        let compaction_region = new_test_compaction_region().await;
1334        let input_meta = dummy_file_meta();
1335        let expired_meta = dummy_file_meta();
1336
1337        // The single merge output fails, but expired SSTs should still be removed.
1338        let merger = MockMerger::new(vec![Err(error::InvalidMetaSnafu {
1339            reason: "fail".to_string(),
1340        }
1341        .build())]);
1342        let compactor = DefaultCompactor::with_merger(merger);
1343
1344        let picker_output = PickerOutput {
1345            outputs: vec![CompactionOutput {
1346                output_level: 1,
1347                inputs: vec![new_file_handle(input_meta.clone())],
1348                filter_deleted: false,
1349                output_time_range: None,
1350            }],
1351            expired_ssts: vec![new_file_handle(expired_meta.clone())],
1352            time_window_size: 3600,
1353            max_file_size: None,
1354        };
1355
1356        let merge_output = compactor
1357            .merge_ssts(&compaction_region, picker_output)
1358            .await
1359            .unwrap();
1360
1361        // No files added (merge failed).
1362        assert!(merge_output.files_to_add.is_empty());
1363        // Only the expired SST should be in files_to_remove (not the failed merge's input).
1364        assert_eq!(merge_output.files_to_remove.len(), 1);
1365        assert_eq!(
1366            merge_output.files_to_remove[0].file_id,
1367            expired_meta.file_id
1368        );
1369    }
1370
1371    #[derive(Clone)]
1372    struct BlockingMerger {
1373        call_idx: Arc<AtomicUsize>,
1374    }
1375
1376    #[async_trait::async_trait]
1377    impl SstMerger for BlockingMerger {
1378        async fn merge_single_output(
1379            &self,
1380            _compaction_region: CompactionRegion,
1381            _output: CompactionOutput,
1382            _write_opts: WriteOptions,
1383        ) -> Result<(Vec<FileMeta>, Vec<SstInfo>)> {
1384            self.call_idx.fetch_add(1, Ordering::SeqCst);
1385            std::future::pending().await
1386        }
1387    }
1388
1389    #[tokio::test(flavor = "multi_thread")]
1390    async fn test_merge_ssts_cancels_spawned_tasks() {
1391        common_telemetry::init_default_ut_logging();
1392
1393        let mut compaction_region = new_test_compaction_region().await;
1394        compaction_region.max_parallelism = 2;
1395
1396        let cancel_handle = Arc::new(CancellationHandle::default());
1397        let call_idx = Arc::new(AtomicUsize::new(0));
1398        let compactor = DefaultCompactor {
1399            merger: BlockingMerger {
1400                call_idx: call_idx.clone(),
1401            },
1402            cancel_handle: cancel_handle.clone(),
1403            uncommitted: None,
1404        };
1405
1406        let picker_output = PickerOutput {
1407            outputs: vec![
1408                CompactionOutput {
1409                    output_level: 1,
1410                    inputs: vec![new_file_handle(dummy_file_meta())],
1411                    filter_deleted: false,
1412                    output_time_range: None,
1413                },
1414                CompactionOutput {
1415                    output_level: 1,
1416                    inputs: vec![new_file_handle(dummy_file_meta())],
1417                    filter_deleted: false,
1418                    output_time_range: None,
1419                },
1420                CompactionOutput {
1421                    output_level: 1,
1422                    inputs: vec![new_file_handle(dummy_file_meta())],
1423                    filter_deleted: false,
1424                    output_time_range: None,
1425                },
1426            ],
1427            expired_ssts: vec![],
1428            time_window_size: 3600,
1429            max_file_size: None,
1430        };
1431
1432        let task = tokio::spawn(async move {
1433            compactor
1434                .merge_ssts(&compaction_region, picker_output)
1435                .await
1436        });
1437
1438        sleep(Duration::from_millis(100)).await;
1439        cancel_handle.cancel();
1440
1441        let merge_output = task
1442            .await
1443            .expect("merge_ssts should stop after cancellation")
1444            .unwrap();
1445
1446        let started = call_idx.load(Ordering::SeqCst);
1447
1448        assert!(merge_output.files_to_add.is_empty());
1449        assert!(merge_output.files_to_remove.is_empty());
1450        assert_eq!(started, 2);
1451    }
1452}