Skip to main content

risingwave_meta/hummock/manager/compaction/
mod.rs

1// Copyright 2024 RisingWave Labs
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::collections::{BTreeMap, HashMap, HashSet};
16use std::ops::DerefMut;
17use std::sync::{Arc, LazyLock};
18use std::time::Instant;
19
20use anyhow::Context;
21use compaction_event_loop::{
22    HummockCompactionEventDispatcher, HummockCompactionEventHandler, HummockCompactionEventLoop,
23    HummockCompactorDedicatedEventLoop,
24};
25use fail::fail_point;
26use itertools::Itertools;
27use parking_lot::Mutex;
28use rand::rng as thread_rng;
29use rand::seq::SliceRandom;
30use risingwave_common::catalog::TableId;
31use risingwave_common::config::meta::default::compaction_config;
32use risingwave_common::util::epoch::Epoch;
33use risingwave_hummock_sdk::compact_task::{CompactTask, CompactTaskAssignment, ReportTask};
34use risingwave_hummock_sdk::compaction_group::StateTableId;
35use risingwave_hummock_sdk::compaction_group::hummock_version_ext::safe_epoch_table_watermarks_impl;
36use risingwave_hummock_sdk::level::Levels;
37use risingwave_hummock_sdk::sstable_info::SstableInfo;
38use risingwave_hummock_sdk::table_stats::{
39    PbTableStatsMap, add_prost_table_stats_map, purge_prost_table_stats,
40};
41use risingwave_hummock_sdk::table_watermark::TableWatermarks;
42use risingwave_hummock_sdk::version::{GroupDelta, IntraLevelDelta};
43use risingwave_hummock_sdk::{
44    CompactionGroupId, HummockCompactionTaskId, HummockContextId, HummockSstableId,
45    HummockSstableObjectId, HummockVersionId, compact_task_to_string, statistics_compact_task,
46};
47use risingwave_meta_model::hummock_sequence::COMPACTION_TASK_ID;
48use risingwave_pb::hummock::compact_task::{TaskStatus, TaskType};
49use risingwave_pb::hummock::subscribe_compaction_event_response::Event as ResponseEvent;
50use risingwave_pb::hummock::{
51    CompactTaskAssignment as PbCompactTaskAssignment, CompactionConfig, PbCompactStatus,
52    SubscribeCompactionEventRequest, TableOption, compact_task,
53};
54use thiserror_ext::AsReport;
55use tokio::sync::mpsc::UnboundedReceiver;
56use tokio::sync::oneshot::{Receiver, Sender};
57use tokio::task::JoinHandle;
58use tonic::Streaming;
59use tracing::warn;
60
61use crate::hummock::compaction::in_progress_compaction::InProgressCompactionView;
62use crate::hummock::compaction::selector::level_selector::PickerInfo;
63use crate::hummock::compaction::selector::{
64    DynamicLevelSelector, DynamicLevelSelectorCore, LocalSelectorStatistic, ManualCompactionOption,
65    ManualCompactionSelector, SpaceReclaimCompactionSelector, TombstoneCompactionSelector,
66    TtlCompactionSelector, VnodeWatermarkCompactionSelector,
67};
68use crate::hummock::compaction::{
69    CompactStatus, CompactionDeveloperConfig, CompactionSelector,
70    CompactionTask as PickedCompactionTask,
71};
72use crate::hummock::error::{Error, Result};
73use crate::hummock::manager::CompactionTaskReportResult;
74use crate::hummock::manager::compaction::compact_task_builder::{
75    CompactTaskBuildContext, attach_compact_task_table_metadata, build_base_compact_task,
76};
77use crate::hummock::manager::transaction::{
78    HummockVersionStatsTransaction, HummockVersionTransaction,
79};
80use crate::hummock::manager::versioning::Versioning;
81use crate::hummock::metrics_utils::{
82    build_compact_task_level_type_metrics_label, trigger_compact_tasks_stat,
83    trigger_local_table_stat,
84};
85use crate::hummock::model::CompactionGroup;
86use crate::hummock::{HummockManager, commit_multi_var};
87use crate::manager::META_NODE_ID;
88use crate::model::BTreeMapTransaction;
89
90#[derive(Debug, Clone, Copy, PartialEq, Eq)]
91pub enum ManualCompactionTriggerResult {
92    Submitted,
93    Retry,
94}
95
96mod compact_task_builder;
97pub mod compaction_event_loop;
98pub mod compaction_group_manager;
99pub mod compaction_group_schedule;
100
101static CANCEL_STATUS_SET: LazyLock<HashSet<TaskStatus>> = LazyLock::new(|| {
102    [
103        TaskStatus::ManualCanceled,
104        TaskStatus::SendFailCanceled,
105        TaskStatus::AssignFailCanceled,
106        TaskStatus::HeartbeatCanceled,
107        TaskStatus::InvalidGroupCanceled,
108        TaskStatus::NoAvailMemoryResourceCanceled,
109        TaskStatus::NoAvailCpuResourceCanceled,
110        TaskStatus::HeartbeatProgressCanceled,
111    ]
112    .into_iter()
113    .collect()
114});
115
116fn init_selectors() -> HashMap<compact_task::TaskType, Box<dyn CompactionSelector>> {
117    let mut compaction_selectors: HashMap<compact_task::TaskType, Box<dyn CompactionSelector>> =
118        HashMap::default();
119    compaction_selectors.insert(
120        compact_task::TaskType::Dynamic,
121        Box::<DynamicLevelSelector>::default(),
122    );
123    compaction_selectors.insert(
124        compact_task::TaskType::SpaceReclaim,
125        Box::<SpaceReclaimCompactionSelector>::default(),
126    );
127    compaction_selectors.insert(
128        compact_task::TaskType::Ttl,
129        Box::<TtlCompactionSelector>::default(),
130    );
131    compaction_selectors.insert(
132        compact_task::TaskType::Tombstone,
133        Box::<TombstoneCompactionSelector>::default(),
134    );
135    compaction_selectors.insert(
136        compact_task::TaskType::VnodeWatermark,
137        Box::<VnodeWatermarkCompactionSelector>::default(),
138    );
139    compaction_selectors
140}
141
142enum BuiltCompactTask {
143    MetaFinished(CompactTask),
144    PendingAssignment(CompactTask),
145}
146
147impl HummockVersionTransaction<'_> {
148    fn apply_compact_task(&mut self, compact_task: &CompactTask) {
149        let mut version_delta = self.new_delta();
150        let trivial_move = compact_task.is_trivial_move_task();
151        version_delta.trivial_move = trivial_move;
152
153        let group_deltas = &mut version_delta
154            .group_deltas
155            .entry(compact_task.compaction_group_id)
156            .or_default()
157            .group_deltas;
158        let mut removed_table_ids_map: BTreeMap<u32, HashSet<HummockSstableId>> =
159            BTreeMap::default();
160
161        for level in &compact_task.input_ssts {
162            let level_idx = level.level_idx;
163
164            removed_table_ids_map
165                .entry(level_idx)
166                .or_default()
167                .extend(level.table_infos.iter().map(|sst| sst.sst_id));
168        }
169
170        for (level_idx, removed_table_ids) in removed_table_ids_map {
171            let group_delta = GroupDelta::IntraLevel(IntraLevelDelta::new(
172                level_idx,
173                0, // default
174                removed_table_ids,
175                vec![], // default
176                0,      // default
177                compact_task.compaction_group_version_id,
178            ));
179
180            group_deltas.push(group_delta);
181        }
182
183        let group_delta = GroupDelta::IntraLevel(IntraLevelDelta::new(
184            compact_task.target_level,
185            compact_task.target_sub_level_id,
186            HashSet::new(), // default
187            compact_task.sorted_output_ssts.clone(),
188            compact_task.split_weight_by_vnode,
189            compact_task.compaction_group_version_id,
190        ));
191
192        group_deltas.push(group_delta);
193        version_delta.pre_apply();
194    }
195}
196
197#[derive(Default)]
198pub struct Compaction {
199    /// Compaction task that is already assigned to a compactor
200    pub compact_task_assignment: BTreeMap<HummockCompactionTaskId, CompactTaskAssignment>,
201    /// `CompactStatus` of each compaction group
202    pub compaction_statuses: BTreeMap<CompactionGroupId, CompactStatus>,
203
204    pub _deterministic_mode: bool,
205}
206
207impl HummockManager {
208    pub async fn get_assigned_compact_task_num(&self) -> u64 {
209        self.compaction
210            .read_with_process_name("get_assigned_compact_task_num")
211            .await
212            .compact_task_assignment
213            .len() as u64
214    }
215
216    pub async fn list_compaction_status(
217        &self,
218    ) -> (Vec<PbCompactStatus>, Vec<PbCompactTaskAssignment>) {
219        let (compaction_statuses, compact_task_assignments) = {
220            let compaction = self
221                .compaction
222                .read_with_process_name("list_compaction_status")
223                .await;
224            (
225                compaction
226                    .compaction_statuses
227                    .values()
228                    .map_into()
229                    .collect_vec(),
230                compaction
231                    .compact_task_assignment
232                    .values()
233                    .cloned()
234                    .collect_vec(),
235            )
236        };
237
238        (
239            compaction_statuses,
240            compact_task_assignments
241                .into_iter()
242                .map(PbCompactTaskAssignment::from)
243                .collect(),
244        )
245    }
246
247    pub async fn get_compaction_scores(
248        &self,
249        compaction_group_id: CompactionGroupId,
250    ) -> Vec<PickerInfo> {
251        let (status, levels, group) = {
252            let compaction = self
253                .compaction
254                .read_with_process_name("get_compaction_scores")
255                .await;
256            let versioning = self
257                .versioning
258                .read_with_process_name("get_compaction_scores")
259                .await;
260            let config_manager = self
261                .compaction_group_manager
262                .read_with_process_name("get_compaction_scores")
263                .await;
264            match (
265                compaction.compaction_statuses.get(&compaction_group_id),
266                versioning.current_version.levels.get(&compaction_group_id),
267                config_manager.try_get_compaction_group_config(compaction_group_id),
268            ) {
269                (Some(cs), Some(v), Some(cf)) => (cs.to_owned(), v.to_owned(), cf),
270                _ => {
271                    return vec![];
272                }
273            }
274        };
275        let dynamic_level_core = DynamicLevelSelectorCore::new(
276            group.compaction_config,
277            Arc::new(CompactionDeveloperConfig::default()),
278        );
279        let ctx = dynamic_level_core.get_priority_levels(&levels, &status.level_handlers);
280        ctx.score_levels
281    }
282}
283
284impl HummockManager {
285    pub fn compaction_event_loop(
286        hummock_manager: Arc<Self>,
287        compactor_streams_change_rx: UnboundedReceiver<(
288            HummockContextId,
289            Streaming<SubscribeCompactionEventRequest>,
290        )>,
291    ) -> Vec<(JoinHandle<()>, Sender<()>)> {
292        let mut join_handle_vec = Vec::default();
293
294        let hummock_compaction_event_handler =
295            HummockCompactionEventHandler::new(hummock_manager.clone());
296
297        let dedicated_event_loop = HummockCompactorDedicatedEventLoop::new(
298            hummock_manager.clone(),
299            hummock_compaction_event_handler.clone(),
300        );
301
302        let (dedicated_event_loop_join_handle, event_tx, shutdown_tx) = dedicated_event_loop.run();
303        join_handle_vec.push((dedicated_event_loop_join_handle, shutdown_tx));
304
305        let hummock_compaction_event_dispatcher = HummockCompactionEventDispatcher::new(
306            hummock_manager.env.opts.clone(),
307            hummock_compaction_event_handler,
308            Some(event_tx),
309        );
310
311        let event_loop = HummockCompactionEventLoop::new(
312            hummock_compaction_event_dispatcher,
313            hummock_manager.metrics.clone(),
314            compactor_streams_change_rx,
315        );
316
317        let (event_loop_join_handle, event_loop_shutdown_tx) = event_loop.run();
318        join_handle_vec.push((event_loop_join_handle, event_loop_shutdown_tx));
319
320        join_handle_vec
321    }
322
323    pub fn add_compactor_stream(
324        &self,
325        context_id: HummockContextId,
326        req_stream: Streaming<SubscribeCompactionEventRequest>,
327    ) {
328        self.compactor_streams_change_tx
329            .send((context_id, req_stream))
330            .unwrap();
331    }
332}
333
334impl HummockManager {
335    /// Gets one compaction task id with best-effort batching while ensuring concurrent callers
336    /// share the same refill instead of wasting an allocated range.
337    async fn next_compaction_task_id_with_prefetch(&self, refill_capacity: u32) -> Result<u64> {
338        self.prefetched_compaction_task_ids
339            .next(refill_capacity, |count| async move {
340                self.env
341                    .hummock_seq
342                    .next_interval(COMPACTION_TASK_ID, count)
343                    .await
344            })
345            .await
346    }
347
348    pub async fn get_compact_tasks_impl(
349        &self,
350        compaction_groups: Vec<CompactionGroupId>,
351        max_select_count: usize,
352        selector: &mut dyn CompactionSelector,
353    ) -> Result<(Vec<CompactTask>, Vec<CompactionGroupId>)> {
354        let deterministic_mode = self.env.opts.compaction_deterministic_test;
355
356        let mut compaction_guard = self
357            .compaction
358            .write_with_process_name("get_compact_tasks_impl")
359            .await;
360        let mut versioning_guard = self
361            .versioning
362            .write_with_process_name("get_compact_tasks_impl")
363            .await;
364        let compaction: &mut Compaction = &mut compaction_guard;
365        let versioning: &mut Versioning = &mut versioning_guard;
366
367        let start_time = Instant::now();
368        let mut compaction_statuses = BTreeMapTransaction::new(&mut compaction.compaction_statuses);
369
370        let mut compact_task_assignment =
371            BTreeMapTransaction::new(&mut compaction.compact_task_assignment);
372
373        let mut version = HummockVersionTransaction::new(
374            &mut versioning.current_version,
375            &mut versioning.hummock_version_deltas,
376            &mut versioning.table_change_log,
377            self.env.notification_manager(),
378            None,
379            &self.metrics,
380            &self.env.opts,
381            &self.version_stat_tx,
382        );
383        // Apply stats changes.
384        let mut version_stats = HummockVersionStatsTransaction::new(
385            &mut versioning.version_stats,
386            self.env.notification_manager(),
387        );
388
389        if deterministic_mode {
390            version.disable_apply_to_txn();
391        }
392        let all_versioned_table_schemas = if self.env.opts.enable_dropped_column_reclaim {
393            self.metadata_manager
394                .catalog_controller
395                .get_versioned_table_schemas()
396                .await
397                .map_err(|e| Error::Internal(e.into()))?
398        } else {
399            HashMap::default()
400        };
401        let mut unschedule_groups = vec![];
402        let mut trivial_tasks = vec![];
403        let mut pick_tasks = vec![];
404        let developer_config = Arc::new(CompactionDeveloperConfig::new_from_meta_opts(
405            &self.env.opts,
406        ));
407        // Reuse prefetched task ids from previous loops.
408        // Each group consumes at most one task_id (trivial tasks share the same id with normal
409        // task). When prefetched ids are exhausted, refill in fixed-size chunks to avoid
410        // per-group SQL transactions while keeping the in-memory cache small.
411        'outside: for compaction_group_id in compaction_groups {
412            if pick_tasks.len() >= max_select_count {
413                break;
414            }
415
416            if !version
417                .latest_version()
418                .levels
419                .contains_key(&compaction_group_id)
420            {
421                continue;
422            }
423
424            // When the last table of a compaction group is deleted, the compaction group (and its
425            // config) is destroyed as well. Then a compaction task for this group may come later and
426            // cannot find its config.
427            let group_config = {
428                let config_manager = self
429                    .compaction_group_manager
430                    .read_with_process_name("get_compact_tasks_impl")
431                    .await;
432
433                match config_manager.try_get_compaction_group_config(compaction_group_id) {
434                    Some(config) => config,
435                    None => continue,
436                }
437            };
438
439            // Use prefetched task id if available; when exhausted, refill in chunks first.
440            let task_id = self
441                .next_compaction_task_id_with_prefetch(
442                    self.env.opts.compaction_task_id_refill_capacity,
443                )
444                .await?;
445
446            if !compaction_statuses.contains_key(&compaction_group_id) {
447                // lazy initialize.
448                compaction_statuses.insert(
449                    compaction_group_id,
450                    CompactStatus::new(
451                        compaction_group_id,
452                        group_config.compaction_config.max_level,
453                    ),
454                );
455            }
456            let mut compact_status = compaction_statuses.get_mut(compaction_group_id).unwrap();
457
458            let mut stats = LocalSelectorStatistic::default();
459            let member_table_ids: Vec<_> = version
460                .latest_version()
461                .state_table_info
462                .compaction_group_member_table_ids(compaction_group_id)
463                .iter()
464                .copied()
465                .collect();
466
467            let mut table_id_to_option: HashMap<TableId, _> = HashMap::default();
468
469            {
470                let guard = self.table_id_to_table_option.read();
471                for table_id in &member_table_ids {
472                    if let Some(opts) = guard.get(table_id) {
473                        table_id_to_option.insert(*table_id, *opts);
474                    }
475                }
476            }
477
478            let in_progress_compactions = InProgressCompactionView::for_group(
479                compact_task_assignment.tree_ref().values(),
480                compaction_group_id,
481            );
482
483            while let Some(picked_task) = compact_status.get_compact_task(
484                version
485                    .latest_version()
486                    .get_compaction_group_levels(compaction_group_id),
487                version
488                    .latest_version()
489                    .state_table_info
490                    .compaction_group_member_table_ids(compaction_group_id),
491                task_id as HummockCompactionTaskId,
492                &group_config,
493                &mut stats,
494                selector,
495                &table_id_to_option,
496                developer_config.clone(),
497                &version.latest_version().table_watermarks,
498                &version.latest_version().state_table_info,
499                &in_progress_compactions,
500            ) {
501                let compaction_group_levels = version
502                    .latest_version()
503                    .get_compaction_group_levels(compaction_group_id);
504                let target_level_id = picked_task.input.target_level as u32;
505                let is_target_level_last = compaction_group_levels.is_last_level(target_level_id);
506                let table_options = table_id_to_option
507                    .iter()
508                    .map(|(table_id, table_option)| (*table_id, TableOption::from(table_option)))
509                    .collect();
510                let built_compact_task = self.build_ready_compact_task(
511                    picked_task,
512                    CompactTaskBuildContext {
513                        task_id,
514                        compaction_group_id: group_config.group_id,
515                        compaction_group_version_id: compaction_group_levels
516                            .compaction_group_version_id,
517                        existing_table_ids: member_table_ids.clone(),
518                        table_options,
519                        is_target_level_last,
520                        compaction_config: group_config.compaction_config.clone(),
521                        current_epoch_time: Epoch::now().0,
522                    },
523                    &version.latest_version().table_watermarks,
524                    &all_versioned_table_schemas,
525                );
526
527                match built_compact_task {
528                    BuiltCompactTask::MetaFinished(compact_task) => {
529                        let label = compact_task.task_label();
530                        tracing::debug!(
531                            "{} for compaction group {}: input: {:?}, cost time: {:?}",
532                            label,
533                            compact_task.compaction_group_id,
534                            compact_task.input_ssts,
535                            start_time.elapsed()
536                        );
537                        compact_status.report_compact_task(&compact_task);
538                        update_table_stats_for_vnode_watermark_trivial_reclaim(
539                            &mut version_stats.table_stats,
540                            &compact_task,
541                        );
542                        self.metrics
543                            .compact_frequency
544                            .with_label_values(&[
545                                label,
546                                &compact_task.compaction_group_id.to_string(),
547                                selector.task_type().as_str_name(),
548                                "SUCCESS",
549                            ])
550                            .inc();
551
552                        version.apply_compact_task(&compact_task);
553                        trivial_tasks.push(compact_task);
554                        if trivial_tasks.len() >= self.env.opts.max_trivial_move_task_count_per_loop
555                        {
556                            break 'outside;
557                        }
558                    }
559                    BuiltCompactTask::PendingAssignment(compact_task) => {
560                        compact_task_assignment.insert(
561                            compact_task.task_id,
562                            CompactTaskAssignment {
563                                compact_task: compact_task.clone(),
564                                context_id: META_NODE_ID, // deprecated
565                            },
566                        );
567
568                        pick_tasks.push(compact_task);
569                        break;
570                    }
571                }
572
573                stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
574                stats = LocalSelectorStatistic::default();
575            }
576            if pick_tasks
577                .last()
578                .map(|task| task.compaction_group_id != compaction_group_id)
579                .unwrap_or(true)
580            {
581                unschedule_groups.push(compaction_group_id);
582            }
583            stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
584        }
585
586        if !trivial_tasks.is_empty() {
587            commit_multi_var!(
588                self.meta_store_ref(),
589                compaction_statuses,
590                compact_task_assignment,
591                version,
592                version_stats
593            )?;
594            self.metrics
595                .compact_task_batch_count
596                .with_label_values(&["batch_trivial_move"])
597                .observe(trivial_tasks.len() as f64);
598
599            for trivial_task in &trivial_tasks {
600                self.metrics
601                    .compact_task_trivial_move_sst_count
602                    .with_label_values(&[&trivial_task.compaction_group_id.to_string()])
603                    .observe(trivial_task.input_ssts[0].table_infos.len() as _);
604            }
605
606            drop(versioning_guard);
607        } else {
608            // We are using a single transaction to ensure that each task has progress when it is
609            // created.
610            drop(versioning_guard);
611            commit_multi_var!(
612                self.meta_store_ref(),
613                compaction_statuses,
614                compact_task_assignment
615            )?;
616        }
617        drop(compaction_guard);
618        if !pick_tasks.is_empty() {
619            self.metrics
620                .compact_task_batch_count
621                .with_label_values(&["batch_get_compact_task"])
622                .observe(pick_tasks.len() as f64);
623        }
624
625        for compact_task in &mut pick_tasks {
626            let compaction_group_id = compact_task.compaction_group_id;
627
628            // Initiate heartbeat for the task to track its progress.
629            self.compactor_manager
630                .initiate_task_heartbeat(compact_task.clone());
631
632            // this task has been finished.
633            compact_task.task_status = TaskStatus::Pending;
634            let compact_task_statistics = statistics_compact_task(compact_task);
635
636            let level_type_label = build_compact_task_level_type_metrics_label(
637                compact_task.input_ssts[0].level_idx as usize,
638                compact_task.input_ssts.last().unwrap().level_idx as usize,
639            );
640
641            let level_count = compact_task.input_ssts.len();
642            if compact_task.input_ssts[0].level_idx == 0 {
643                self.metrics
644                    .l0_compact_level_count
645                    .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
646                    .observe(level_count as _);
647            }
648
649            self.metrics
650                .compact_task_size
651                .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
652                .observe(compact_task_statistics.total_file_size as _);
653
654            self.metrics
655                .compact_task_size
656                .with_label_values(&[
657                    &compaction_group_id.to_string(),
658                    &format!("{} uncompressed", level_type_label),
659                ])
660                .observe(compact_task_statistics.total_uncompressed_file_size as _);
661
662            self.metrics
663                .compact_task_file_count
664                .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
665                .observe(compact_task_statistics.total_file_count as _);
666
667            tracing::trace!(
668                "For compaction group {}: pick up {} {} sub_level in level {} to compact to target {}. cost time: {:?} compact_task_statistics {:?}",
669                compaction_group_id,
670                level_count,
671                compact_task.input_ssts[0].level_type.as_str_name(),
672                compact_task.input_ssts[0].level_idx,
673                compact_task.target_level,
674                start_time.elapsed(),
675                compact_task_statistics
676            );
677        }
678
679        #[cfg(test)]
680        {
681            self.check_state_consistency().await;
682        }
683        pick_tasks.extend(trivial_tasks);
684        Ok((pick_tasks, unschedule_groups))
685    }
686
687    /// Cancels a compaction task no matter it's assigned or unassigned.
688    pub async fn cancel_compact_task(&self, task_id: u64, task_status: TaskStatus) -> Result<bool> {
689        fail_point!("fp_cancel_compact_task", |_| Err(Error::MetaStore(
690            anyhow::anyhow!("failpoint metastore err")
691        )));
692        let ret = self
693            .cancel_compact_task_impl(vec![task_id], task_status)
694            .await?;
695        Ok(ret[0])
696    }
697
698    pub async fn cancel_compact_tasks(
699        &self,
700        tasks: Vec<u64>,
701        task_status: TaskStatus,
702    ) -> Result<Vec<bool>> {
703        self.cancel_compact_task_impl(tasks, task_status).await
704    }
705
706    async fn cancel_compact_task_impl(
707        &self,
708        task_ids: Vec<u64>,
709        task_status: TaskStatus,
710    ) -> Result<Vec<bool>> {
711        assert!(CANCEL_STATUS_SET.contains(&task_status));
712        let tasks = task_ids
713            .into_iter()
714            .map(|task_id| ReportTask {
715                task_id,
716                task_status,
717                sorted_output_ssts: vec![],
718                table_stats_change: HashMap::default(),
719                object_timestamps: HashMap::default(),
720            })
721            .collect_vec();
722        let rets = self.report_compact_tasks(tasks).await?;
723        #[cfg(test)]
724        {
725            self.check_state_consistency().await;
726        }
727        Ok(rets)
728    }
729
730    async fn get_compact_tasks(
731        &self,
732        mut compaction_groups: Vec<CompactionGroupId>,
733        max_select_count: usize,
734        selector: &mut dyn CompactionSelector,
735    ) -> Result<(Vec<CompactTask>, Vec<CompactionGroupId>)> {
736        fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
737            anyhow::anyhow!("failpoint metastore error")
738        )));
739        compaction_groups.shuffle(&mut thread_rng());
740        let (mut tasks, groups) = self
741            .get_compact_tasks_impl(compaction_groups, max_select_count, selector)
742            .await?;
743        tasks.retain(|task| {
744            if task.task_status == TaskStatus::Success {
745                debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
746                false
747            } else {
748                true
749            }
750        });
751        Ok((tasks, groups))
752    }
753
754    pub async fn get_compact_task(
755        &self,
756        compaction_group_id: CompactionGroupId,
757        selector: &mut dyn CompactionSelector,
758    ) -> Result<Option<CompactTask>> {
759        fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
760            anyhow::anyhow!("failpoint metastore error")
761        )));
762
763        let (normal_tasks, _) = self
764            .get_compact_tasks_impl(vec![compaction_group_id], 1, selector)
765            .await?;
766        for task in normal_tasks {
767            if task.task_status != TaskStatus::Success {
768                return Ok(Some(task));
769            }
770            debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
771        }
772        Ok(None)
773    }
774
775    pub async fn manual_get_compact_task(
776        &self,
777        compaction_group_id: CompactionGroupId,
778        manual_compaction_option: ManualCompactionOption,
779    ) -> Result<Option<CompactTask>> {
780        let (task, _) = self
781            .manual_get_compact_task_with_info(compaction_group_id, manual_compaction_option)
782            .await?;
783        Ok(task)
784    }
785
786    pub async fn manual_get_compact_task_with_info(
787        &self,
788        compaction_group_id: CompactionGroupId,
789        manual_compaction_option: ManualCompactionOption,
790    ) -> Result<(Option<CompactTask>, bool)> {
791        let mut selector = ManualCompactionSelector::new(manual_compaction_option);
792        let task = self
793            .get_compact_task(compaction_group_id, &mut selector)
794            .await?;
795        if let Some(err) = selector.validation_error() {
796            return Err(Error::InvalidManualCompactionOption(err.to_owned()));
797        }
798        Ok((task, selector.blocked_by_pending()))
799    }
800
801    pub async fn report_compact_task(
802        &self,
803        task_id: u64,
804        task_status: TaskStatus,
805        sorted_output_ssts: Vec<SstableInfo>,
806        table_stats_change: Option<PbTableStatsMap>,
807        object_timestamps: HashMap<HummockSstableObjectId, u64>,
808    ) -> Result<bool> {
809        let rets = self
810            .report_compact_tasks(vec![ReportTask {
811                task_id,
812                task_status,
813                sorted_output_ssts,
814                table_stats_change: table_stats_change.unwrap_or_default(),
815                object_timestamps,
816            }])
817            .await?;
818        Ok(rets[0])
819    }
820
821    pub async fn report_compact_tasks(&self, report_tasks: Vec<ReportTask>) -> Result<Vec<bool>> {
822        let compaction_guard = self
823            .compaction
824            .write_with_process_name("report_compact_tasks")
825            .await;
826        let versioning_guard = self
827            .versioning
828            .write_with_process_name("report_compact_tasks")
829            .await;
830
831        self.report_compact_tasks_impl(report_tasks, compaction_guard, versioning_guard)
832            .await
833    }
834
835    /// Finishes or cancels a compaction task, according to `task_status`.
836    ///
837    /// If `context_id` is not None, its validity will be checked when writing meta store.
838    /// Its ownership of the task is checked as well.
839    ///
840    /// Return Ok(false) indicates either the task is not found,
841    /// or the task is not owned by `context_id` when `context_id` is not None.
842    pub async fn report_compact_tasks_impl(
843        &self,
844        report_tasks: Vec<ReportTask>,
845        mut compaction_guard: impl DerefMut<Target = Compaction>,
846        mut versioning_guard: impl DerefMut<Target = Versioning>,
847    ) -> Result<Vec<bool>> {
848        let deterministic_mode = self.env.opts.compaction_deterministic_test;
849        let compaction: &mut Compaction = &mut compaction_guard;
850        let start_time = Instant::now();
851        let original_keys = compaction.compaction_statuses.keys().cloned().collect_vec();
852        let mut compact_statuses = BTreeMapTransaction::new(&mut compaction.compaction_statuses);
853        let mut rets = vec![false; report_tasks.len()];
854        let mut compact_task_assignment =
855            BTreeMapTransaction::new(&mut compaction.compact_task_assignment);
856        // The compaction task is finished.
857        let versioning: &mut Versioning = &mut versioning_guard;
858
859        // purge stale compact_status
860        for group_id in original_keys {
861            if !versioning.current_version.levels.contains_key(&group_id) {
862                compact_statuses.remove(group_id);
863            }
864        }
865        let mut tasks = vec![];
866
867        let mut version = HummockVersionTransaction::new(
868            &mut versioning.current_version,
869            &mut versioning.hummock_version_deltas,
870            &mut versioning.table_change_log,
871            self.env.notification_manager(),
872            None,
873            &self.metrics,
874            &self.env.opts,
875            &self.version_stat_tx,
876        );
877
878        if deterministic_mode {
879            version.disable_apply_to_txn();
880        }
881
882        let mut version_stats = HummockVersionStatsTransaction::new(
883            &mut versioning.version_stats,
884            self.env.notification_manager(),
885        );
886        let mut success_count = 0;
887        let mut report_results = Vec::with_capacity(rets.len());
888        for (idx, task) in report_tasks.into_iter().enumerate() {
889            rets[idx] = true;
890            let task_id = task.task_id;
891            let mut task_status = task.task_status;
892            let mut compact_task = match compact_task_assignment.remove(task.task_id) {
893                Some(compact_task_assignment) => compact_task_assignment.compact_task,
894                None => {
895                    tracing::warn!("{}", format!("compact task {} not found", task.task_id));
896                    rets[idx] = false;
897                    report_results.push(CompactionTaskReportResult {
898                        task_id,
899                        task_status,
900                        reported: false,
901                    });
902                    continue;
903                }
904            };
905
906            {
907                // apply result
908                compact_task.task_status = task.task_status;
909                compact_task.sorted_output_ssts = task.sorted_output_ssts;
910            }
911
912            match compact_statuses.get_mut(compact_task.compaction_group_id) {
913                Some(mut compact_status) => {
914                    compact_status.report_compact_task(&compact_task);
915                }
916                None => {
917                    // When the group_id is not found in the compaction_statuses, it means the group has been removed.
918                    // The task is invalid and should be canceled.
919                    // e.g.
920                    // 1. The group is removed by the user unregistering the tables
921                    // 2. The group is removed by the group scheduling algorithm
922                    compact_task.task_status = TaskStatus::InvalidGroupCanceled;
923                }
924            }
925
926            let is_success = if let TaskStatus::Success = compact_task.task_status {
927                match self
928                    .report_compaction_sanity_check(&task.object_timestamps)
929                    .await
930                {
931                    Err(e) => {
932                        warn!(
933                            "failed to commit compaction task {} {}",
934                            compact_task.task_id,
935                            e.as_report()
936                        );
937                        compact_task.task_status = TaskStatus::RetentionTimeRejected;
938                        false
939                    }
940                    _ => {
941                        let group = version
942                            .latest_version()
943                            .levels
944                            .get(&compact_task.compaction_group_id)
945                            .unwrap();
946                        let is_expired = compact_task.is_expired(group.compaction_group_version_id);
947                        if is_expired {
948                            compact_task.task_status = TaskStatus::InputOutdatedCanceled;
949                            warn!(
950                                "The task may be expired because of group split, task:\n {:?}",
951                                compact_task_to_string(&compact_task)
952                            );
953                        }
954                        !is_expired
955                    }
956                }
957            } else {
958                false
959            };
960            if is_success {
961                success_count += 1;
962                version.apply_compact_task(&compact_task);
963                if purge_prost_table_stats(
964                    &mut version_stats.table_stats,
965                    version.latest_version(),
966                    &HashSet::default(),
967                ) {
968                    self.metrics.version_stats.reset();
969                    versioning.local_metrics.clear();
970                }
971                add_prost_table_stats_map(&mut version_stats.table_stats, &task.table_stats_change);
972                trigger_local_table_stat(
973                    &self.metrics,
974                    &mut versioning.local_metrics,
975                    &version_stats,
976                    &task.table_stats_change,
977                );
978            }
979            task_status = compact_task.task_status;
980            report_results.push(CompactionTaskReportResult {
981                task_id,
982                task_status,
983                reported: rets[idx],
984            });
985            tasks.push(compact_task);
986        }
987        if success_count > 0 {
988            commit_multi_var!(
989                self.meta_store_ref(),
990                compact_statuses,
991                compact_task_assignment,
992                version,
993                version_stats
994            )?;
995
996            self.metrics
997                .compact_task_batch_count
998                .with_label_values(&["batch_report_task"])
999                .observe(success_count as f64);
1000        } else {
1001            // The compaction task is cancelled or failed.
1002            commit_multi_var!(
1003                self.meta_store_ref(),
1004                compact_statuses,
1005                compact_task_assignment
1006            )?;
1007        }
1008
1009        self.notify_compaction_task_report_waiters(report_results);
1010
1011        let mut success_groups = vec![];
1012        for compact_task in &tasks {
1013            self.compactor_manager
1014                .remove_task_heartbeat(compact_task.task_id);
1015            tracing::trace!(
1016                "Reported compaction task. {}. cost time: {:?}",
1017                compact_task_to_string(compact_task),
1018                start_time.elapsed(),
1019            );
1020
1021            if !deterministic_mode
1022                && (matches!(compact_task.task_type, compact_task::TaskType::Dynamic)
1023                    || matches!(compact_task.task_type, compact_task::TaskType::Emergency))
1024            {
1025                // only try send Dynamic compaction
1026                self.try_send_compaction_request(
1027                    compact_task.compaction_group_id,
1028                    compact_task::TaskType::Dynamic,
1029                );
1030            }
1031
1032            if compact_task.task_status == TaskStatus::Success {
1033                success_groups.push(compact_task.compaction_group_id);
1034            }
1035        }
1036
1037        trigger_compact_tasks_stat(
1038            &self.metrics,
1039            &tasks,
1040            &compaction.compaction_statuses,
1041            &versioning_guard.current_version,
1042        );
1043        drop(versioning_guard);
1044        if !success_groups.is_empty() {
1045            self.try_update_write_limits(&success_groups).await;
1046        }
1047        Ok(rets)
1048    }
1049
1050    /// Triggers compacitons to specified compaction groups.
1051    /// Don't wait for compaction finish
1052    pub async fn trigger_compaction_deterministic(
1053        &self,
1054        _base_version_id: HummockVersionId,
1055        compaction_groups: Vec<CompactionGroupId>,
1056    ) -> Result<()> {
1057        self.on_current_version(|old_version| {
1058            tracing::info!(
1059                "Trigger compaction for version {}, groups {:?}",
1060                old_version.id,
1061                compaction_groups
1062            );
1063        })
1064        .await;
1065
1066        if compaction_groups.is_empty() {
1067            return Ok(());
1068        }
1069        for compaction_group in compaction_groups {
1070            self.try_send_compaction_request(compaction_group, compact_task::TaskType::Dynamic);
1071        }
1072        Ok(())
1073    }
1074
1075    pub async fn trigger_manual_compaction(
1076        &self,
1077        compaction_group: CompactionGroupId,
1078        manual_compaction_option: ManualCompactionOption,
1079    ) -> Result<ManualCompactionTriggerResult> {
1080        let start_time = Instant::now();
1081        let exclusive = manual_compaction_option.exclusive;
1082
1083        // 1. Get idle compactor.
1084        let compactor = match self.compactor_manager.next_compactor() {
1085            Some(compactor) => compactor,
1086            None => {
1087                tracing::warn!("trigger_manual_compaction No compactor is available.");
1088                return Err(anyhow::anyhow!(
1089                    "trigger_manual_compaction No compactor is available. compaction_group {}",
1090                    compaction_group
1091                )
1092                .into());
1093            }
1094        };
1095
1096        // 2. Get manual compaction task.
1097        let compact_task = self
1098            .manual_get_compact_task_with_info(compaction_group, manual_compaction_option)
1099            .await;
1100        let (compact_task, blocked_by_pending) = match compact_task {
1101            Ok((compact_task, blocked_by_pending)) => (compact_task, blocked_by_pending),
1102            Err(err) => {
1103                tracing::warn!(error = %err.as_report(), "Failed to get compaction task");
1104                if matches!(err, Error::InvalidManualCompactionOption(_)) {
1105                    return Err(err);
1106                }
1107
1108                return Err(anyhow::anyhow!(err)
1109                    .context(format!(
1110                        "Failed to get compaction task for compaction_group {}",
1111                        compaction_group,
1112                    ))
1113                    .into());
1114            }
1115        };
1116        let compact_task = match compact_task {
1117            Some(compact_task) => compact_task,
1118            None => {
1119                if exclusive && blocked_by_pending {
1120                    return Ok(ManualCompactionTriggerResult::Retry);
1121                }
1122                // No compaction task available.
1123                return Err(anyhow::anyhow!(
1124                    "trigger_manual_compaction No compaction_task is available. compaction_group {}",
1125                    compaction_group
1126                )
1127                .into());
1128            }
1129        };
1130
1131        // 3. send task to compactor
1132        let task_id = compact_task.task_id;
1133        let compact_task_string = compact_task_to_string(&compact_task);
1134        tracing::info!(
1135            compact_task_string,
1136            duration = ?start_time.elapsed(),
1137            "Triggered manual compaction task."
1138        );
1139
1140        let report_rx = self.register_compaction_task_report_waiter(task_id);
1141        if let Err(err) = compactor
1142            .send_event(ResponseEvent::CompactTask(compact_task.into()))
1143            .with_context(|| {
1144                format!(
1145                    "Failed to trigger compaction task for compaction_group {}",
1146                    compaction_group,
1147                )
1148            })
1149        {
1150            self.remove_compaction_task_report_waiter(task_id);
1151            return Err(err.into());
1152        }
1153
1154        let report_result = match report_rx.await {
1155            Ok(result) => result,
1156            Err(_) => {
1157                self.remove_compaction_task_report_waiter(task_id);
1158                return Err(anyhow::anyhow!(
1159                    "trigger_manual_compaction wait report failed. compaction_group {}",
1160                    compaction_group
1161                )
1162                .into());
1163            }
1164        };
1165        if !report_result.reported {
1166            return Err(anyhow::anyhow!(
1167                "trigger_manual_compaction report not accepted. task_id {}",
1168                report_result.task_id
1169            )
1170            .into());
1171        }
1172
1173        if report_result.task_status == TaskStatus::NoAvailCpuResourceCanceled
1174            || report_result.task_status == TaskStatus::NoAvailMemoryResourceCanceled
1175        {
1176            return Ok(ManualCompactionTriggerResult::Retry);
1177        }
1178
1179        tracing::info!(
1180            ?report_result,
1181            duration = ?start_time.elapsed(),
1182            "Completed manual compaction task."
1183        );
1184
1185        Ok(ManualCompactionTriggerResult::Submitted)
1186    }
1187
1188    /// Sends a compaction request for new data (clears cooldown).
1189    pub fn try_send_compaction_request(
1190        &self,
1191        compaction_group: CompactionGroupId,
1192        task_type: compact_task::TaskType,
1193    ) -> bool {
1194        self.compaction_state.try_sched_compaction(
1195            compaction_group,
1196            task_type,
1197            ScheduleTrigger::NewData,
1198        )
1199    }
1200
1201    /// Apply `split_weight_by_vnode` based partition strategy.
1202    /// This handles dynamic partitioning based on table size and write throughput.
1203    fn apply_split_weight_by_vnode_partition(
1204        &self,
1205        compact_task: &mut CompactTask,
1206        compaction_config: &CompactionConfig,
1207        compact_table_ids: &[TableId],
1208    ) {
1209        if compaction_config.split_weight_by_vnode > 0 {
1210            for table_id in compact_table_ids {
1211                compact_task
1212                    .table_vnode_partition
1213                    .insert(*table_id, compact_task.split_weight_by_vnode);
1214            }
1215
1216            return;
1217        }
1218
1219        // Calculate per-table size from normalized input SSTs.
1220        let mut table_size_info: HashMap<TableId, u64> = HashMap::default();
1221        for input_ssts in &compact_task.input_ssts {
1222            for sst in &input_ssts.table_infos {
1223                for table_id in &sst.table_ids {
1224                    *table_size_info.entry(*table_id).or_default() +=
1225                        sst.sst_size / (sst.table_ids.len() as u64);
1226                }
1227            }
1228        }
1229
1230        let hybrid_vnode_count = self.env.opts.hybrid_partition_node_count;
1231        let default_partition_count = self.env.opts.partition_vnode_count;
1232        let compact_task_table_size_partition_threshold_low = self
1233            .env
1234            .opts
1235            .compact_task_table_size_partition_threshold_low;
1236        let compact_task_table_size_partition_threshold_high = self
1237            .env
1238            .opts
1239            .compact_task_table_size_partition_threshold_high;
1240
1241        // Check latest write throughput
1242        let table_write_throughput_statistic_manager =
1243            self.table_write_throughput_statistic_manager.read();
1244        let timestamp = chrono::Utc::now().timestamp();
1245
1246        for (table_id, compact_table_size) in table_size_info {
1247            let write_throughput = table_write_throughput_statistic_manager
1248                .get_table_throughput_descending(table_id, timestamp)
1249                .peekable()
1250                .peek()
1251                .map(|item| item.throughput)
1252                .unwrap_or(0);
1253
1254            if compact_table_size > compact_task_table_size_partition_threshold_high
1255                && default_partition_count > 0
1256            {
1257                compact_task
1258                    .table_vnode_partition
1259                    .insert(table_id, default_partition_count);
1260            } else if (compact_table_size > compact_task_table_size_partition_threshold_low
1261                || (write_throughput > self.env.opts.table_high_write_throughput_threshold
1262                    && compact_table_size > compaction_config.target_file_size_base))
1263                && hybrid_vnode_count > 0
1264            {
1265                compact_task
1266                    .table_vnode_partition
1267                    .insert(table_id, hybrid_vnode_count);
1268            } else if compact_table_size > compaction_config.target_file_size_base {
1269                compact_task.table_vnode_partition.insert(table_id, 1);
1270            }
1271        }
1272
1273        compact_task
1274            .table_vnode_partition
1275            .retain(|table_id, _| compact_table_ids.contains(table_id));
1276    }
1277
1278    pub(crate) fn calculate_vnode_partition(
1279        &self,
1280        compact_task: &mut CompactTask,
1281        compaction_config: &CompactionConfig,
1282        compact_table_ids: &[TableId],
1283    ) {
1284        // Do not split sst by vnode partition when target_level > base_level
1285        // The purpose of data alignment is mainly to improve the parallelism of base level compaction
1286        // and reduce write amplification. However, at high level, the size of the sst file is often
1287        // larger and only contains the data of a single table_id, so there is no need to cut it.
1288        if compact_task.target_level > compact_task.base_level {
1289            return;
1290        }
1291
1292        // Apply split_weight_by_vnode based partition strategy
1293        self.apply_split_weight_by_vnode_partition(
1294            compact_task,
1295            compaction_config,
1296            compact_table_ids,
1297        );
1298    }
1299
1300    fn build_ready_compact_task(
1301        &self,
1302        picked_task: PickedCompactionTask,
1303        context: CompactTaskBuildContext,
1304        table_watermarks: &HashMap<TableId, Arc<TableWatermarks>>,
1305        all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1306    ) -> BuiltCompactTask {
1307        let compaction_config = context.compaction_config.clone();
1308        let (mut compact_task, compact_table_ids) = build_base_compact_task(picked_task, context);
1309
1310        if compact_task.is_trivial_reclaim() {
1311            compact_task.task_status = TaskStatus::Success;
1312            compact_task.sorted_output_ssts.clear();
1313            return BuiltCompactTask::MetaFinished(compact_task);
1314        }
1315
1316        if compact_task.is_trivial_move_task() {
1317            compact_task.task_status = TaskStatus::Success;
1318            compact_task.sorted_output_ssts = compact_task.input_ssts[0]
1319                .read_sstable_infos()
1320                .cloned()
1321                .collect();
1322            return BuiltCompactTask::MetaFinished(compact_task);
1323        }
1324
1325        self.prepare_compact_task_for_assignment(
1326            &mut compact_task,
1327            compaction_config.as_ref(),
1328            &compact_table_ids,
1329            safe_epoch_table_watermarks_impl(table_watermarks, &compact_table_ids),
1330            all_versioned_table_schemas,
1331        );
1332
1333        BuiltCompactTask::PendingAssignment(compact_task)
1334    }
1335
1336    fn prepare_compact_task_for_assignment(
1337        &self,
1338        compact_task: &mut CompactTask,
1339        compaction_config: &CompactionConfig,
1340        compact_table_ids: &[TableId],
1341        table_watermarks: BTreeMap<TableId, TableWatermarks>,
1342        all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1343    ) {
1344        self.calculate_vnode_partition(compact_task, compaction_config, compact_table_ids);
1345        attach_compact_task_table_metadata(
1346            compact_task,
1347            compact_table_ids,
1348            table_watermarks,
1349            all_versioned_table_schemas,
1350        );
1351    }
1352
1353    pub fn compactor_manager_ref(&self) -> crate::hummock::CompactorManagerRef {
1354        self.compactor_manager.clone()
1355    }
1356
1357    fn register_compaction_task_report_waiter(
1358        &self,
1359        task_id: HummockCompactionTaskId,
1360    ) -> Receiver<CompactionTaskReportResult> {
1361        let (tx, rx) = tokio::sync::oneshot::channel();
1362        self.compaction_task_report_notifiers
1363            .lock()
1364            .register(task_id, tx);
1365        rx
1366    }
1367
1368    fn remove_compaction_task_report_waiter(&self, task_id: HummockCompactionTaskId) {
1369        self.compaction_task_report_notifiers.lock().remove(task_id);
1370    }
1371
1372    fn notify_compaction_task_report_waiters(&self, results: Vec<CompactionTaskReportResult>) {
1373        let mut guard = self.compaction_task_report_notifiers.lock();
1374        for result in results {
1375            guard.notify(result);
1376        }
1377    }
1378}
1379
1380#[cfg(any(test, feature = "test"))]
1381impl HummockManager {
1382    pub async fn compaction_task_from_assignment_for_test(
1383        &self,
1384        task_id: u64,
1385    ) -> Option<CompactTaskAssignment> {
1386        let compaction_guard = self
1387            .compaction
1388            .read_with_process_name("compaction_task_from_assignment_for_test")
1389            .await;
1390        let assignment_ref = &compaction_guard.compact_task_assignment;
1391        assignment_ref.get(&task_id).cloned()
1392    }
1393
1394    pub async fn report_compact_task_for_test(
1395        &self,
1396        task_id: u64,
1397        compact_task: Option<CompactTask>,
1398        task_status: TaskStatus,
1399        sorted_output_ssts: Vec<SstableInfo>,
1400        table_stats_change: Option<PbTableStatsMap>,
1401    ) -> Result<()> {
1402        if let Some(task) = compact_task {
1403            let mut guard = self
1404                .compaction
1405                .write_with_process_name("report_compact_task_for_test")
1406                .await;
1407            guard.compact_task_assignment.insert(
1408                task_id,
1409                CompactTaskAssignment {
1410                    compact_task: task,
1411                    context_id: 0.into(),
1412                },
1413            );
1414        }
1415
1416        // In the test, the contents of the compact task may have been modified directly, while the contents of compact_task_assignment were not modified.
1417        // So we pass the modified compact_task directly into the `report_compact_task_impl`
1418        self.report_compact_tasks(vec![ReportTask {
1419            task_id,
1420            task_status,
1421            sorted_output_ssts,
1422            table_stats_change: table_stats_change.unwrap_or_default(),
1423            object_timestamps: HashMap::default(),
1424        }])
1425        .await?;
1426        Ok(())
1427    }
1428}
1429
1430/// What triggered the compaction schedule request.
1431#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1432pub enum ScheduleTrigger {
1433    /// New data arrived (e.g., `commit_epoch`). Clears cooldown.
1434    NewData,
1435    /// Periodic timer. Respects cooldown for Dynamic type.
1436    Periodic,
1437}
1438
1439/// A point-in-time snapshot of the compaction schedule state.
1440///
1441/// `snapshot_time` is used by `unschedule()` to detect whether new data arrived
1442/// after the snapshot was taken, preventing incorrect cooldown.
1443pub struct CompactionScheduleSnapshot {
1444    scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1445    snapshot_time: Instant,
1446}
1447
1448impl CompactionScheduleSnapshot {
1449    /// Task type priority order for scheduling (checked first = higher priority).
1450    const TASK_TYPE_PRIORITY: &[TaskType] = &[
1451        TaskType::Dynamic,
1452        TaskType::SpaceReclaim,
1453        TaskType::Ttl,
1454        TaskType::Tombstone,
1455        TaskType::VnodeWatermark,
1456    ];
1457
1458    pub fn snapshot_time(&self) -> Instant {
1459        self.snapshot_time
1460    }
1461
1462    /// Pick compaction groups and task type from this snapshot.
1463    ///
1464    /// Returns groups in shuffled order. Non-Dynamic types have higher priority
1465    /// and return a single group; Dynamic groups are batched together.
1466    pub fn pick_compaction_groups_and_type(&self) -> Option<(Vec<CompactionGroupId>, TaskType)> {
1467        let group_ids = self.group_ids_shuffled();
1468        let mut normal_groups = vec![];
1469        for cg_id in group_ids {
1470            if let Some(pick_type) = self.pick_type(cg_id) {
1471                if pick_type == TaskType::Dynamic {
1472                    normal_groups.push(cg_id);
1473                } else if normal_groups.is_empty() {
1474                    return Some((vec![cg_id], pick_type));
1475                }
1476            }
1477        }
1478        if normal_groups.is_empty() {
1479            None
1480        } else {
1481            Some((normal_groups, TaskType::Dynamic))
1482        }
1483    }
1484
1485    fn group_ids_shuffled(&self) -> Vec<CompactionGroupId> {
1486        let mut group_ids: Vec<_> = self.scheduled.iter().map(|(g, _)| *g).unique().collect();
1487        group_ids.shuffle(&mut thread_rng());
1488        group_ids
1489    }
1490
1491    fn pick_type(&self, group: CompactionGroupId) -> Option<TaskType> {
1492        Self::TASK_TYPE_PRIORITY
1493            .iter()
1494            .find(|t| self.scheduled.contains(&(group, **t)))
1495            .copied()
1496    }
1497}
1498
1499/// Tracks which (`compaction_group`, `task_type`) pairs are scheduled for compaction.
1500///
1501/// For `Dynamic` type, includes a cooldown mechanism: groups with no compaction work
1502/// are skipped by periodic triggers until new data arrives via `commit_epoch`.
1503#[derive(Debug, Default)]
1504pub struct CompactionState {
1505    inner: Mutex<CompactionStateInner>,
1506}
1507
1508#[derive(Debug, Default)]
1509struct CompactionStateInner {
1510    scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1511    /// Groups skipped by periodic Dynamic trigger until new data arrives.
1512    dynamic_cooldown: HashSet<CompactionGroupId>,
1513    /// Tracks new-data arrival time per group for cooldown race detection.
1514    last_new_data_time: HashMap<CompactionGroupId, Instant>,
1515}
1516
1517impl CompactionState {
1518    pub fn new() -> Self {
1519        Self {
1520            inner: Default::default(),
1521        }
1522    }
1523
1524    /// Enqueues a compaction request. Returns `true` if newly scheduled.
1525    ///
1526    /// `trigger` only affects `Dynamic` type — see [`ScheduleTrigger`].
1527    pub fn try_sched_compaction(
1528        &self,
1529        compaction_group: CompactionGroupId,
1530        task_type: TaskType,
1531        trigger: ScheduleTrigger,
1532    ) -> bool {
1533        let mut guard = self.inner.lock();
1534        if task_type == TaskType::Dynamic {
1535            match trigger {
1536                ScheduleTrigger::NewData => {
1537                    guard.dynamic_cooldown.remove(&compaction_group);
1538                    guard
1539                        .last_new_data_time
1540                        .insert(compaction_group, Instant::now());
1541                }
1542                ScheduleTrigger::Periodic => {
1543                    if guard.dynamic_cooldown.contains(&compaction_group) {
1544                        return false;
1545                    }
1546                }
1547            }
1548        }
1549        guard.scheduled.insert((compaction_group, task_type))
1550    }
1551
1552    /// Removes a scheduled entry. For Dynamic type, adds to cooldown unless
1553    /// new data arrived after `snapshot_time`.
1554    pub fn unschedule(
1555        &self,
1556        compaction_group: CompactionGroupId,
1557        task_type: compact_task::TaskType,
1558        snapshot_time: Instant,
1559    ) {
1560        let mut guard = self.inner.lock();
1561        guard.scheduled.remove(&(compaction_group, task_type));
1562        if task_type == TaskType::Dynamic {
1563            let has_new_data = guard
1564                .last_new_data_time
1565                .get(&compaction_group)
1566                .is_some_and(|t| *t > snapshot_time);
1567            if !has_new_data {
1568                guard.dynamic_cooldown.insert(compaction_group);
1569            }
1570        }
1571    }
1572
1573    /// Takes a snapshot of the current schedule state.
1574    pub fn snapshot(&self) -> CompactionScheduleSnapshot {
1575        let guard = self.inner.lock();
1576        // Record time after lock to ensure accurate ordering vs. try_sched_compaction
1577        let snapshot_time = Instant::now();
1578        CompactionScheduleSnapshot {
1579            scheduled: guard.scheduled.clone(),
1580            snapshot_time,
1581        }
1582    }
1583
1584    /// Removes all schedule state for a deleted or merged group.
1585    pub fn remove_compaction_group(&self, compaction_group: CompactionGroupId) {
1586        let mut guard = self.inner.lock();
1587        guard
1588            .scheduled
1589            .retain(|(group, _)| *group != compaction_group);
1590        guard.dynamic_cooldown.remove(&compaction_group);
1591        guard.last_new_data_time.remove(&compaction_group);
1592    }
1593}
1594
1595impl Compaction {
1596    pub fn get_compact_task_assignments_by_group_id(
1597        &self,
1598        compaction_group_id: CompactionGroupId,
1599    ) -> Vec<CompactTaskAssignment> {
1600        self.compact_task_assignment
1601            .values()
1602            .filter_map(|assignment| {
1603                if assignment.compact_task.compaction_group_id == compaction_group_id {
1604                    Some(assignment.clone())
1605                } else {
1606                    None
1607                }
1608            })
1609            .collect()
1610    }
1611}
1612
1613#[derive(Clone, Default)]
1614pub struct CompactionGroupStatistic {
1615    pub group_id: CompactionGroupId,
1616    pub group_size: u64,
1617    pub table_statistic: BTreeMap<StateTableId, u64>,
1618    pub compaction_group_config: CompactionGroup,
1619}
1620
1621/// Updates table stats caused by vnode watermark trivial reclaim compaction.
1622fn update_table_stats_for_vnode_watermark_trivial_reclaim(
1623    table_stats: &mut PbTableStatsMap,
1624    task: &CompactTask,
1625) {
1626    if task.task_type != TaskType::VnodeWatermark {
1627        return;
1628    }
1629    let mut deleted_table_keys: HashMap<TableId, u64> = HashMap::default();
1630    for s in task.input_ssts.iter().flat_map(|l| l.table_infos.iter()) {
1631        assert_eq!(s.table_ids.len(), 1);
1632        let e = deleted_table_keys.entry(s.table_ids[0]).or_insert(0);
1633        *e += s.total_key_count;
1634    }
1635    for (table_id, delete_count) in deleted_table_keys {
1636        let Some(stats) = table_stats.get_mut(&table_id) else {
1637            continue;
1638        };
1639        if stats.total_key_count == 0 {
1640            continue;
1641        }
1642        let new_total_key_count = stats.total_key_count.saturating_sub(delete_count as i64);
1643        let ratio = new_total_key_count as f64 / stats.total_key_count as f64;
1644        // total_key_count is updated accurately.
1645        stats.total_key_count = new_total_key_count;
1646        // others are updated approximately.
1647        stats.total_key_size = (stats.total_key_size as f64 * ratio).ceil() as i64;
1648        stats.total_value_size = (stats.total_value_size as f64 * ratio).ceil() as i64;
1649    }
1650}
1651
1652#[derive(Debug, Clone)]
1653pub enum GroupState {
1654    /// The compaction group is not in emergency state.
1655    Normal,
1656
1657    /// The compaction group is in emergency state.
1658    Emergency(String), // reason
1659
1660    /// The compaction group is in write stop state.
1661    WriteStop(String), // reason
1662}
1663
1664impl GroupState {
1665    pub fn is_write_stop(&self) -> bool {
1666        matches!(self, Self::WriteStop(_))
1667    }
1668
1669    pub fn is_emergency(&self) -> bool {
1670        matches!(self, Self::Emergency(_))
1671    }
1672
1673    pub fn reason(&self) -> Option<&str> {
1674        match self {
1675            Self::Emergency(reason) | Self::WriteStop(reason) => Some(reason),
1676            _ => None,
1677        }
1678    }
1679}
1680
1681#[derive(Clone, Default)]
1682pub struct GroupStateValidator;
1683
1684impl GroupStateValidator {
1685    pub fn write_stop_sub_level_count(
1686        level_count: usize,
1687        compaction_config: &CompactionConfig,
1688    ) -> bool {
1689        let threshold = compaction_config.level0_stop_write_threshold_sub_level_number as usize;
1690        level_count > threshold
1691    }
1692
1693    pub fn write_stop_l0_size(l0_size: u64, compaction_config: &CompactionConfig) -> bool {
1694        l0_size
1695            > compaction_config
1696                .level0_stop_write_threshold_max_size
1697                .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1698    }
1699
1700    pub fn write_stop_l0_file_count(
1701        l0_file_count: usize,
1702        compaction_config: &CompactionConfig,
1703    ) -> bool {
1704        l0_file_count
1705            > compaction_config
1706                .level0_stop_write_threshold_max_sst_count
1707                .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1708                as usize
1709    }
1710
1711    pub fn emergency_l0_file_count(
1712        l0_file_count: usize,
1713        compaction_config: &CompactionConfig,
1714    ) -> bool {
1715        l0_file_count
1716            > compaction_config
1717                .emergency_level0_sst_file_count
1718                .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1719                as usize
1720    }
1721
1722    pub fn emergency_l0_partition_count(
1723        last_l0_sub_level_partition_count: usize,
1724        compaction_config: &CompactionConfig,
1725    ) -> bool {
1726        last_l0_sub_level_partition_count
1727            > compaction_config
1728                .emergency_level0_sub_level_partition
1729                .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1730                as usize
1731    }
1732
1733    pub fn check_single_group_write_stop(
1734        levels: &Levels,
1735        compaction_config: &CompactionConfig,
1736    ) -> GroupState {
1737        if Self::write_stop_sub_level_count(levels.l0.sub_levels.len(), compaction_config) {
1738            return GroupState::WriteStop(format!(
1739                "WriteStop(l0_level_count: {}, threshold: {}) too many L0 sub levels",
1740                levels.l0.sub_levels.len(),
1741                compaction_config.level0_stop_write_threshold_sub_level_number
1742            ));
1743        }
1744
1745        if Self::write_stop_l0_file_count(
1746            levels
1747                .l0
1748                .sub_levels
1749                .iter()
1750                .map(|l| l.table_infos.len())
1751                .sum(),
1752            compaction_config,
1753        ) {
1754            return GroupState::WriteStop(format!(
1755                "WriteStop(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1756                levels
1757                    .l0
1758                    .sub_levels
1759                    .iter()
1760                    .map(|l| l.table_infos.len())
1761                    .sum::<usize>(),
1762                compaction_config
1763                    .level0_stop_write_threshold_max_sst_count
1764                    .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1765            ));
1766        }
1767
1768        if Self::write_stop_l0_size(levels.l0.total_file_size, compaction_config) {
1769            return GroupState::WriteStop(format!(
1770                "WriteStop(l0_size: {}, threshold: {}) too large L0 size",
1771                levels.l0.total_file_size,
1772                compaction_config
1773                    .level0_stop_write_threshold_max_size
1774                    .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1775            ));
1776        }
1777
1778        GroupState::Normal
1779    }
1780
1781    pub fn check_single_group_emergency(
1782        levels: &Levels,
1783        compaction_config: &CompactionConfig,
1784    ) -> GroupState {
1785        if Self::emergency_l0_file_count(
1786            levels
1787                .l0
1788                .sub_levels
1789                .iter()
1790                .map(|l| l.table_infos.len())
1791                .sum(),
1792            compaction_config,
1793        ) {
1794            return GroupState::Emergency(format!(
1795                "Emergency(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1796                levels
1797                    .l0
1798                    .sub_levels
1799                    .iter()
1800                    .map(|l| l.table_infos.len())
1801                    .sum::<usize>(),
1802                compaction_config
1803                    .emergency_level0_sst_file_count
1804                    .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1805            ));
1806        }
1807
1808        if Self::emergency_l0_partition_count(
1809            levels
1810                .l0
1811                .sub_levels
1812                .first()
1813                .map(|l| l.table_infos.len())
1814                .unwrap_or(0),
1815            compaction_config,
1816        ) {
1817            return GroupState::Emergency(format!(
1818                "Emergency(l0_partition_count: {}, threshold: {}) too many L0 partitions",
1819                levels
1820                    .l0
1821                    .sub_levels
1822                    .first()
1823                    .map(|l| l.table_infos.len())
1824                    .unwrap_or(0),
1825                compaction_config
1826                    .emergency_level0_sub_level_partition
1827                    .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1828            ));
1829        }
1830
1831        GroupState::Normal
1832    }
1833
1834    pub fn group_state(levels: &Levels, compaction_config: &CompactionConfig) -> GroupState {
1835        let state = Self::check_single_group_write_stop(levels, compaction_config);
1836        if state.is_write_stop() {
1837            return state;
1838        }
1839
1840        Self::check_single_group_emergency(levels, compaction_config)
1841    }
1842}
1843
1844#[cfg(test)]
1845mod compaction_state_tests {
1846    use risingwave_pb::hummock::compact_task::TaskType;
1847
1848    use super::*;
1849
1850    #[test]
1851    fn test_basic_schedule_and_unschedule() {
1852        let state = CompactionState::new();
1853        let group_id: CompactionGroupId = 1.into();
1854
1855        // First schedule should succeed
1856        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1857        // Duplicate schedule should fail
1858        assert!(!state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1859        // Different task type should succeed
1860        assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1861
1862        // Snapshot should contain both
1863        let snapshot = state.snapshot();
1864        assert!(snapshot.scheduled.contains(&(group_id, TaskType::Dynamic)));
1865        assert!(snapshot.scheduled.contains(&(group_id, TaskType::Ttl)));
1866
1867        // Unschedule removes from scheduled set
1868        state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1869        let snapshot2 = state.snapshot();
1870        assert!(!snapshot2.scheduled.contains(&(group_id, TaskType::Dynamic)));
1871        assert!(snapshot2.scheduled.contains(&(group_id, TaskType::Ttl)));
1872    }
1873
1874    #[test]
1875    fn test_cooldown_blocks_periodic_trigger() {
1876        let state = CompactionState::new();
1877        let group_id: CompactionGroupId = 1.into();
1878
1879        // Schedule then unschedule - should add to cooldown
1880        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1881        let snapshot = state.snapshot();
1882        state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1883
1884        // Verify in cooldown
1885        assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
1886
1887        // Periodic trigger should be blocked
1888        assert!(!state.try_sched_compaction(
1889            group_id,
1890            TaskType::Dynamic,
1891            ScheduleTrigger::Periodic
1892        ));
1893    }
1894
1895    #[test]
1896    fn test_new_data_clears_cooldown() {
1897        let state = CompactionState::new();
1898        let group_id: CompactionGroupId = 1.into();
1899
1900        // Put group in cooldown
1901        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1902        let snapshot = state.snapshot();
1903        state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1904        assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
1905
1906        // NewData trigger should clear cooldown and schedule
1907        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1908        assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
1909    }
1910
1911    #[test]
1912    fn test_cooldown_only_affects_dynamic_type() {
1913        let state = CompactionState::new();
1914        let group_id: CompactionGroupId = 1.into();
1915
1916        // Put group in cooldown for Dynamic
1917        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1918        let snapshot = state.snapshot();
1919        state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1920
1921        // Ttl unschedule should NOT add to cooldown
1922        let group_id_2: CompactionGroupId = 2.into();
1923        assert!(state.try_sched_compaction(group_id_2, TaskType::Ttl, ScheduleTrigger::Periodic));
1924        let snapshot2 = state.snapshot();
1925        state.unschedule(group_id_2, TaskType::Ttl, snapshot2.snapshot_time());
1926        assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id_2));
1927
1928        // Other task types should work regardless of cooldown
1929        assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1930        assert!(state.try_sched_compaction(
1931            group_id,
1932            TaskType::SpaceReclaim,
1933            ScheduleTrigger::Periodic
1934        ));
1935    }
1936
1937    #[test]
1938    fn test_race_condition_new_data_after_snapshot() {
1939        let state = CompactionState::new();
1940        let group_id: CompactionGroupId = 1.into();
1941
1942        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1943        let snapshot = state.snapshot();
1944
1945        // Simulate new data arriving AFTER snapshot
1946        {
1947            let mut guard = state.inner.lock();
1948            guard.last_new_data_time.insert(group_id, Instant::now());
1949        }
1950
1951        // Unschedule should NOT add to cooldown (new data arrived after snapshot)
1952        state.unschedule(group_id, TaskType::Dynamic, snapshot.snapshot_time());
1953        assert!(
1954            !state.inner.lock().dynamic_cooldown.contains(&group_id),
1955            "Should skip cooldown when new data arrived after snapshot"
1956        );
1957    }
1958
1959    #[test]
1960    fn test_remove_compaction_group_cleans_all_state() {
1961        let state = CompactionState::new();
1962        let group_id: CompactionGroupId = 1.into();
1963
1964        // Set up state
1965        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1966        assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1967        state.inner.lock().dynamic_cooldown.insert(group_id);
1968
1969        // Remove group
1970        state.remove_compaction_group(group_id);
1971
1972        // Verify all state cleaned up
1973        let guard = state.inner.lock();
1974        assert!(!guard.scheduled.contains(&(group_id, TaskType::Dynamic)));
1975        assert!(!guard.scheduled.contains(&(group_id, TaskType::Ttl)));
1976        assert!(!guard.dynamic_cooldown.contains(&group_id));
1977        assert!(!guard.last_new_data_time.contains_key(&group_id));
1978    }
1979
1980    #[test]
1981    fn test_snapshot_pick_type_priority() {
1982        let state = CompactionState::new();
1983        let group_id: CompactionGroupId = 1.into();
1984
1985        // Empty group returns None
1986        assert_eq!(state.snapshot().pick_type(group_id), None);
1987
1988        // Priority order: Dynamic > SpaceReclaim > Ttl > Tombstone > VnodeWatermark
1989        state.try_sched_compaction(
1990            group_id,
1991            TaskType::VnodeWatermark,
1992            ScheduleTrigger::Periodic,
1993        );
1994        assert_eq!(
1995            state.snapshot().pick_type(group_id),
1996            Some(TaskType::VnodeWatermark)
1997        );
1998
1999        state.try_sched_compaction(group_id, TaskType::Tombstone, ScheduleTrigger::Periodic);
2000        assert_eq!(
2001            state.snapshot().pick_type(group_id),
2002            Some(TaskType::Tombstone)
2003        );
2004
2005        state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic);
2006        assert_eq!(state.snapshot().pick_type(group_id), Some(TaskType::Ttl));
2007
2008        state.try_sched_compaction(group_id, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2009        assert_eq!(
2010            state.snapshot().pick_type(group_id),
2011            Some(TaskType::SpaceReclaim)
2012        );
2013
2014        state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData);
2015        assert_eq!(
2016            state.snapshot().pick_type(group_id),
2017            Some(TaskType::Dynamic)
2018        );
2019    }
2020
2021    #[test]
2022    fn test_multiple_groups_independent_cooldown() {
2023        let state = CompactionState::new();
2024        let g1: CompactionGroupId = 1.into();
2025        let g2: CompactionGroupId = 2.into();
2026
2027        state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2028        state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2029        let snapshot = state.snapshot();
2030
2031        // Only unschedule g1
2032        state.unschedule(g1, TaskType::Dynamic, snapshot.snapshot_time());
2033
2034        let guard = state.inner.lock();
2035        assert!(guard.dynamic_cooldown.contains(&g1));
2036        assert!(!guard.dynamic_cooldown.contains(&g2));
2037    }
2038
2039    #[test]
2040    fn test_pick_compaction_groups_empty() {
2041        let state = CompactionState::new();
2042        let snapshot = state.snapshot();
2043        // No scheduled groups → returns None
2044        assert!(snapshot.pick_compaction_groups_and_type().is_none());
2045    }
2046
2047    #[test]
2048    fn test_pick_compaction_groups_mixed_types() {
2049        let state = CompactionState::new();
2050        let g1: CompactionGroupId = 1.into();
2051        let g2: CompactionGroupId = 2.into();
2052        let g3: CompactionGroupId = 3.into();
2053
2054        // g1: Dynamic, g2: Ttl, g3: Dynamic
2055        state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2056        state.try_sched_compaction(g2, TaskType::Ttl, ScheduleTrigger::Periodic);
2057        state.try_sched_compaction(g3, TaskType::Dynamic, ScheduleTrigger::NewData);
2058
2059        let snapshot = state.snapshot();
2060        let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2061
2062        // Due to shuffle, either:
2063        // - Ttl group is encountered first → returns (vec![g2], Ttl)
2064        // - Dynamic group is encountered first → collects all Dynamic, skips Ttl
2065        //   → returns ([g1, g3] in some order, Dynamic)
2066        if task_type == TaskType::Dynamic {
2067            assert!(groups.contains(&g1));
2068            assert!(groups.contains(&g3));
2069            assert!(!groups.contains(&g2)); // Ttl group excluded from Dynamic result
2070        } else {
2071            assert_eq!(task_type, TaskType::Ttl);
2072            assert_eq!(groups, vec![g2]);
2073        }
2074    }
2075
2076    #[test]
2077    fn test_pick_compaction_groups_all_dynamic() {
2078        let state = CompactionState::new();
2079        let g1: CompactionGroupId = 1.into();
2080        let g2: CompactionGroupId = 2.into();
2081
2082        state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2083        state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2084
2085        let snapshot = state.snapshot();
2086        let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2087        assert_eq!(task_type, TaskType::Dynamic);
2088        assert!(groups.contains(&g1));
2089        assert!(groups.contains(&g2));
2090    }
2091
2092    #[test]
2093    fn test_pick_compaction_groups_single_non_dynamic() {
2094        let state = CompactionState::new();
2095        let g1: CompactionGroupId = 1.into();
2096
2097        state.try_sched_compaction(g1, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2098
2099        let snapshot = state.snapshot();
2100        let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2101        assert_eq!(task_type, TaskType::SpaceReclaim);
2102        assert_eq!(groups, vec![g1]);
2103    }
2104}