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                // A scheduling snapshot may still contain a group deleted before we took the lock.
422                continue;
423            }
424
425            // When the last table of a compaction group is deleted, the compaction group (and its
426            // config) is destroyed as well. Then a compaction task for this group may come later and
427            // cannot find its config.
428            let group_config = {
429                let config_manager = self
430                    .compaction_group_manager
431                    .read_with_process_name("get_compact_tasks_impl")
432                    .await;
433
434                match config_manager.try_get_compaction_group_config(compaction_group_id) {
435                    Some(config) => config,
436                    None => continue,
437                }
438            };
439
440            // Use prefetched task id if available; when exhausted, refill in chunks first.
441            let task_id = self
442                .next_compaction_task_id_with_prefetch(
443                    self.env.opts.compaction_task_id_refill_capacity,
444                )
445                .await?;
446
447            if !compaction_statuses.contains_key(&compaction_group_id) {
448                // lazy initialize.
449                compaction_statuses.insert(
450                    compaction_group_id,
451                    CompactStatus::new(
452                        compaction_group_id,
453                        group_config.compaction_config.max_level,
454                    ),
455                );
456            }
457            let mut compact_status = compaction_statuses.get_mut(compaction_group_id).unwrap();
458
459            let mut stats = LocalSelectorStatistic::default();
460            let member_table_ids: Vec<_> = version
461                .latest_version()
462                .state_table_info
463                .compaction_group_member_table_ids(compaction_group_id)
464                .iter()
465                .copied()
466                .collect();
467
468            let mut table_id_to_option: HashMap<TableId, _> = HashMap::default();
469
470            {
471                let guard = self.table_id_to_table_option.read();
472                for table_id in &member_table_ids {
473                    if let Some(opts) = guard.get(table_id) {
474                        table_id_to_option.insert(*table_id, *opts);
475                    }
476                }
477            }
478
479            let in_progress_compactions = InProgressCompactionView::for_group(
480                compact_task_assignment.tree_ref().values(),
481                compaction_group_id,
482            );
483
484            while let Some(picked_task) = compact_status.get_compact_task(
485                version
486                    .latest_version()
487                    .get_compaction_group_levels(compaction_group_id),
488                version
489                    .latest_version()
490                    .state_table_info
491                    .compaction_group_member_table_ids(compaction_group_id),
492                task_id as HummockCompactionTaskId,
493                &group_config,
494                &mut stats,
495                selector,
496                &table_id_to_option,
497                developer_config.clone(),
498                &version.latest_version().table_watermarks,
499                &version.latest_version().state_table_info,
500                &in_progress_compactions,
501            ) {
502                let compaction_group_levels = version
503                    .latest_version()
504                    .get_compaction_group_levels(compaction_group_id);
505                let target_level_id = picked_task.input.target_level as u32;
506                let is_target_level_last = compaction_group_levels.is_last_level(target_level_id);
507                let table_options = table_id_to_option
508                    .iter()
509                    .map(|(table_id, table_option)| (*table_id, TableOption::from(table_option)))
510                    .collect();
511                let built_compact_task = self.build_ready_compact_task(
512                    picked_task,
513                    CompactTaskBuildContext {
514                        task_id,
515                        compaction_group_id: group_config.group_id,
516                        compaction_group_version_id: compaction_group_levels
517                            .compaction_group_version_id,
518                        existing_table_ids: member_table_ids.clone(),
519                        table_options,
520                        is_target_level_last,
521                        compaction_config: group_config.compaction_config.clone(),
522                        current_epoch_time: Epoch::now().0,
523                    },
524                    &version.latest_version().table_watermarks,
525                    &all_versioned_table_schemas,
526                );
527
528                match built_compact_task {
529                    BuiltCompactTask::MetaFinished(compact_task) => {
530                        let label = compact_task.task_label();
531                        tracing::debug!(
532                            "{} for compaction group {}: input: {:?}, cost time: {:?}",
533                            label,
534                            compact_task.compaction_group_id,
535                            compact_task.input_ssts,
536                            start_time.elapsed()
537                        );
538                        compact_status.report_compact_task(&compact_task);
539                        update_table_stats_for_vnode_watermark_trivial_reclaim(
540                            &mut version_stats.table_stats,
541                            &compact_task,
542                        );
543                        self.metrics
544                            .compact_frequency
545                            .with_label_values(&[
546                                label,
547                                &compact_task.compaction_group_id.to_string(),
548                                selector.task_type().as_str_name(),
549                                "SUCCESS",
550                            ])
551                            .inc();
552
553                        version.apply_compact_task(&compact_task);
554                        trivial_tasks.push(compact_task);
555                        if trivial_tasks.len() >= self.env.opts.max_trivial_move_task_count_per_loop
556                        {
557                            break 'outside;
558                        }
559                    }
560                    BuiltCompactTask::PendingAssignment(compact_task) => {
561                        compact_task_assignment.insert(
562                            compact_task.task_id,
563                            CompactTaskAssignment {
564                                compact_task: compact_task.clone(),
565                                context_id: META_NODE_ID, // deprecated
566                            },
567                        );
568
569                        pick_tasks.push(compact_task);
570                        break;
571                    }
572                }
573
574                stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
575                stats = LocalSelectorStatistic::default();
576            }
577            if pick_tasks
578                .last()
579                .map(|task| task.compaction_group_id != compaction_group_id)
580                .unwrap_or(true)
581            {
582                unschedule_groups.push(compaction_group_id);
583            }
584            stats.report_to_metrics(compaction_group_id, self.metrics.as_ref());
585        }
586
587        if !trivial_tasks.is_empty() {
588            commit_multi_var!(
589                self.meta_store_ref(),
590                compaction_statuses,
591                compact_task_assignment,
592                version,
593                version_stats
594            )?;
595            self.metrics
596                .compact_task_batch_count
597                .with_label_values(&["batch_trivial_move"])
598                .observe(trivial_tasks.len() as f64);
599
600            for trivial_task in &trivial_tasks {
601                self.metrics
602                    .compact_task_trivial_move_sst_count
603                    .with_label_values(&[&trivial_task.compaction_group_id.to_string()])
604                    .observe(trivial_task.input_ssts[0].table_infos.len() as _);
605            }
606
607            drop(versioning_guard);
608        } else {
609            // We are using a single transaction to ensure that each task has progress when it is
610            // created.
611            drop(versioning_guard);
612            commit_multi_var!(
613                self.meta_store_ref(),
614                compaction_statuses,
615                compact_task_assignment
616            )?;
617        }
618        drop(compaction_guard);
619        if !pick_tasks.is_empty() {
620            self.metrics
621                .compact_task_batch_count
622                .with_label_values(&["batch_get_compact_task"])
623                .observe(pick_tasks.len() as f64);
624        }
625
626        for compact_task in &mut pick_tasks {
627            let compaction_group_id = compact_task.compaction_group_id;
628
629            // Initiate heartbeat for the task to track its progress.
630            self.compactor_manager
631                .initiate_task_heartbeat(compact_task.clone());
632
633            // this task has been finished.
634            compact_task.task_status = TaskStatus::Pending;
635            let compact_task_statistics = statistics_compact_task(compact_task);
636
637            let level_type_label = build_compact_task_level_type_metrics_label(
638                compact_task.input_ssts[0].level_idx as usize,
639                compact_task.input_ssts.last().unwrap().level_idx as usize,
640            );
641
642            let level_count = compact_task.input_ssts.len();
643            if compact_task.input_ssts[0].level_idx == 0 {
644                self.metrics
645                    .l0_compact_level_count
646                    .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
647                    .observe(level_count as _);
648            }
649
650            self.metrics
651                .compact_task_size
652                .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
653                .observe(compact_task_statistics.total_file_size as _);
654
655            self.metrics
656                .compact_task_size
657                .with_label_values(&[
658                    &compaction_group_id.to_string(),
659                    &format!("{} uncompressed", level_type_label),
660                ])
661                .observe(compact_task_statistics.total_uncompressed_file_size as _);
662
663            self.metrics
664                .compact_task_file_count
665                .with_label_values(&[&compaction_group_id.to_string(), &level_type_label])
666                .observe(compact_task_statistics.total_file_count as _);
667
668            tracing::trace!(
669                "For compaction group {}: pick up {} {} sub_level in level {} to compact to target {}. cost time: {:?} compact_task_statistics {:?}",
670                compaction_group_id,
671                level_count,
672                compact_task.input_ssts[0].level_type.as_str_name(),
673                compact_task.input_ssts[0].level_idx,
674                compact_task.target_level,
675                start_time.elapsed(),
676                compact_task_statistics
677            );
678        }
679
680        #[cfg(test)]
681        {
682            self.check_state_consistency().await;
683        }
684        pick_tasks.extend(trivial_tasks);
685        Ok((pick_tasks, unschedule_groups))
686    }
687
688    /// Cancels a compaction task no matter it's assigned or unassigned.
689    pub async fn cancel_compact_task(&self, task_id: u64, task_status: TaskStatus) -> Result<bool> {
690        fail_point!("fp_cancel_compact_task", |_| Err(Error::MetaStore(
691            anyhow::anyhow!("failpoint metastore err")
692        )));
693        let ret = self
694            .cancel_compact_task_impl(vec![task_id], task_status)
695            .await?;
696        Ok(ret[0])
697    }
698
699    pub async fn cancel_compact_tasks(
700        &self,
701        tasks: Vec<u64>,
702        task_status: TaskStatus,
703    ) -> Result<Vec<bool>> {
704        self.cancel_compact_task_impl(tasks, task_status).await
705    }
706
707    async fn cancel_compact_task_impl(
708        &self,
709        task_ids: Vec<u64>,
710        task_status: TaskStatus,
711    ) -> Result<Vec<bool>> {
712        assert!(CANCEL_STATUS_SET.contains(&task_status));
713        let tasks = task_ids
714            .into_iter()
715            .map(|task_id| ReportTask {
716                task_id,
717                task_status,
718                sorted_output_ssts: vec![],
719                table_stats_change: HashMap::default(),
720                object_timestamps: HashMap::default(),
721            })
722            .collect_vec();
723        let rets = self.report_compact_tasks(tasks).await?;
724        #[cfg(test)]
725        {
726            self.check_state_consistency().await;
727        }
728        Ok(rets)
729    }
730
731    async fn get_compact_tasks(
732        &self,
733        mut compaction_groups: Vec<CompactionGroupId>,
734        max_select_count: usize,
735        selector: &mut dyn CompactionSelector,
736    ) -> Result<(Vec<CompactTask>, Vec<CompactionGroupId>)> {
737        fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
738            anyhow::anyhow!("failpoint metastore error")
739        )));
740        compaction_groups.shuffle(&mut thread_rng());
741        let (mut tasks, groups) = self
742            .get_compact_tasks_impl(compaction_groups, max_select_count, selector)
743            .await?;
744        tasks.retain(|task| {
745            if task.task_status == TaskStatus::Success {
746                debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
747                false
748            } else {
749                true
750            }
751        });
752        Ok((tasks, groups))
753    }
754
755    pub async fn get_compact_task(
756        &self,
757        compaction_group_id: CompactionGroupId,
758        selector: &mut dyn CompactionSelector,
759    ) -> Result<Option<CompactTask>> {
760        fail_point!("fp_get_compact_task", |_| Err(Error::MetaStore(
761            anyhow::anyhow!("failpoint metastore error")
762        )));
763
764        let (normal_tasks, _) = self
765            .get_compact_tasks_impl(vec![compaction_group_id], 1, selector)
766            .await?;
767        for task in normal_tasks {
768            if task.task_status != TaskStatus::Success {
769                return Ok(Some(task));
770            }
771            debug_assert!(task.is_trivial_reclaim() || task.is_trivial_move_task());
772        }
773        Ok(None)
774    }
775
776    pub async fn manual_get_compact_task(
777        &self,
778        compaction_group_id: CompactionGroupId,
779        manual_compaction_option: ManualCompactionOption,
780    ) -> Result<Option<CompactTask>> {
781        let (task, _) = self
782            .manual_get_compact_task_with_info(compaction_group_id, manual_compaction_option)
783            .await?;
784        Ok(task)
785    }
786
787    pub async fn manual_get_compact_task_with_info(
788        &self,
789        compaction_group_id: CompactionGroupId,
790        manual_compaction_option: ManualCompactionOption,
791    ) -> Result<(Option<CompactTask>, bool)> {
792        let mut selector = ManualCompactionSelector::new(manual_compaction_option);
793        let task = self
794            .get_compact_task(compaction_group_id, &mut selector)
795            .await?;
796        if let Some(err) = selector.validation_error() {
797            return Err(Error::InvalidManualCompactionOption(err.to_owned()));
798        }
799        Ok((task, selector.blocked_by_pending()))
800    }
801
802    pub async fn report_compact_task(
803        &self,
804        task_id: u64,
805        task_status: TaskStatus,
806        sorted_output_ssts: Vec<SstableInfo>,
807        table_stats_change: Option<PbTableStatsMap>,
808        object_timestamps: HashMap<HummockSstableObjectId, u64>,
809    ) -> Result<bool> {
810        let rets = self
811            .report_compact_tasks(vec![ReportTask {
812                task_id,
813                task_status,
814                sorted_output_ssts,
815                table_stats_change: table_stats_change.unwrap_or_default(),
816                object_timestamps,
817            }])
818            .await?;
819        Ok(rets[0])
820    }
821
822    pub async fn report_compact_tasks(&self, report_tasks: Vec<ReportTask>) -> Result<Vec<bool>> {
823        let compaction_guard = self
824            .compaction
825            .write_with_process_name("report_compact_tasks")
826            .await;
827        let versioning_guard = self
828            .versioning
829            .write_with_process_name("report_compact_tasks")
830            .await;
831
832        self.report_compact_tasks_impl(report_tasks, compaction_guard, versioning_guard)
833            .await
834    }
835
836    /// Finishes or cancels a compaction task, according to `task_status`.
837    ///
838    /// If `context_id` is not None, its validity will be checked when writing meta store.
839    /// Its ownership of the task is checked as well.
840    ///
841    /// Return Ok(false) indicates either the task is not found,
842    /// or the task is not owned by `context_id` when `context_id` is not None.
843    pub async fn report_compact_tasks_impl(
844        &self,
845        report_tasks: Vec<ReportTask>,
846        mut compaction_guard: impl DerefMut<Target = Compaction>,
847        mut versioning_guard: impl DerefMut<Target = Versioning>,
848    ) -> Result<Vec<bool>> {
849        let deterministic_mode = self.env.opts.compaction_deterministic_test;
850        let compaction: &mut Compaction = &mut compaction_guard;
851        let start_time = Instant::now();
852        let original_keys = compaction.compaction_statuses.keys().cloned().collect_vec();
853        let mut compact_statuses = BTreeMapTransaction::new(&mut compaction.compaction_statuses);
854        let mut rets = vec![false; report_tasks.len()];
855        let mut compact_task_assignment =
856            BTreeMapTransaction::new(&mut compaction.compact_task_assignment);
857        // The compaction task is finished.
858        let versioning: &mut Versioning = &mut versioning_guard;
859
860        // purge stale compact_status
861        for group_id in original_keys {
862            if !versioning.current_version.levels.contains_key(&group_id) {
863                compact_statuses.remove(group_id);
864            }
865        }
866        let mut tasks = vec![];
867
868        let mut version = HummockVersionTransaction::new(
869            &mut versioning.current_version,
870            &mut versioning.hummock_version_deltas,
871            &mut versioning.table_change_log,
872            self.env.notification_manager(),
873            None,
874            &self.metrics,
875            &self.env.opts,
876            &self.version_stat_tx,
877        );
878
879        if deterministic_mode {
880            version.disable_apply_to_txn();
881        }
882
883        let mut version_stats = HummockVersionStatsTransaction::new(
884            &mut versioning.version_stats,
885            self.env.notification_manager(),
886        );
887        let mut success_count = 0;
888        let mut report_results = Vec::with_capacity(rets.len());
889        for (idx, task) in report_tasks.into_iter().enumerate() {
890            rets[idx] = true;
891            let task_id = task.task_id;
892            let mut task_status = task.task_status;
893            let mut compact_task = match compact_task_assignment.remove(task.task_id) {
894                Some(compact_task_assignment) => compact_task_assignment.compact_task,
895                None => {
896                    tracing::warn!("{}", format!("compact task {} not found", task.task_id));
897                    rets[idx] = false;
898                    report_results.push(CompactionTaskReportResult {
899                        task_id,
900                        task_status,
901                        reported: false,
902                    });
903                    continue;
904                }
905            };
906
907            {
908                // apply result
909                compact_task.task_status = task.task_status;
910                compact_task.sorted_output_ssts = task.sorted_output_ssts;
911            }
912
913            match compact_statuses.get_mut(compact_task.compaction_group_id) {
914                Some(mut compact_status) => {
915                    compact_status.report_compact_task(&compact_task);
916                }
917                None => {
918                    // When the group_id is not found in the compaction_statuses, it means the group has been removed.
919                    // The task is invalid and should be canceled.
920                    // e.g.
921                    // 1. The group is removed by the user unregistering the tables
922                    // 2. The group is removed by the group scheduling algorithm
923                    compact_task.task_status = TaskStatus::InvalidGroupCanceled;
924                }
925            }
926
927            let is_success = if let TaskStatus::Success = compact_task.task_status {
928                match self
929                    .report_compaction_sanity_check(&task.object_timestamps)
930                    .await
931                {
932                    Err(e) => {
933                        warn!(
934                            "failed to commit compaction task {} {}",
935                            compact_task.task_id,
936                            e.as_report()
937                        );
938                        compact_task.task_status = TaskStatus::RetentionTimeRejected;
939                        false
940                    }
941                    _ => {
942                        let group = version
943                            .latest_version()
944                            .levels
945                            .get(&compact_task.compaction_group_id)
946                            .unwrap();
947                        let is_expired = compact_task.is_expired(group.compaction_group_version_id);
948                        if is_expired {
949                            compact_task.task_status = TaskStatus::InputOutdatedCanceled;
950                            warn!(
951                                "The task may be expired because of group split, task:\n {:?}",
952                                compact_task_to_string(&compact_task)
953                            );
954                        }
955                        !is_expired
956                    }
957                }
958            } else {
959                false
960            };
961            if is_success {
962                success_count += 1;
963                version.apply_compact_task(&compact_task);
964                if purge_prost_table_stats(
965                    &mut version_stats.table_stats,
966                    version.latest_version(),
967                    &HashSet::default(),
968                ) {
969                    self.metrics.version_stats.reset();
970                    versioning.local_metrics.clear();
971                }
972                add_prost_table_stats_map(&mut version_stats.table_stats, &task.table_stats_change);
973                trigger_local_table_stat(
974                    &self.metrics,
975                    &mut versioning.local_metrics,
976                    &version_stats,
977                    &task.table_stats_change,
978                );
979            }
980            task_status = compact_task.task_status;
981            report_results.push(CompactionTaskReportResult {
982                task_id,
983                task_status,
984                reported: rets[idx],
985            });
986            tasks.push(compact_task);
987        }
988        if success_count > 0 {
989            commit_multi_var!(
990                self.meta_store_ref(),
991                compact_statuses,
992                compact_task_assignment,
993                version,
994                version_stats
995            )?;
996
997            self.metrics
998                .compact_task_batch_count
999                .with_label_values(&["batch_report_task"])
1000                .observe(success_count as f64);
1001        } else {
1002            // The compaction task is cancelled or failed.
1003            commit_multi_var!(
1004                self.meta_store_ref(),
1005                compact_statuses,
1006                compact_task_assignment
1007            )?;
1008        }
1009
1010        self.notify_compaction_task_report_waiters(report_results);
1011
1012        let mut success_groups = vec![];
1013        for compact_task in &tasks {
1014            self.compactor_manager
1015                .remove_task_heartbeat(compact_task.task_id);
1016            tracing::trace!(
1017                "Reported compaction task. {}. cost time: {:?}",
1018                compact_task_to_string(compact_task),
1019                start_time.elapsed(),
1020            );
1021
1022            if !deterministic_mode
1023                && versioning_guard
1024                    .current_version
1025                    .levels
1026                    .contains_key(&compact_task.compaction_group_id)
1027                && (matches!(compact_task.task_type, compact_task::TaskType::Dynamic)
1028                    || matches!(compact_task.task_type, compact_task::TaskType::Emergency))
1029            {
1030                // only try send Dynamic compaction
1031                self.try_send_compaction_request(
1032                    compact_task.compaction_group_id,
1033                    compact_task::TaskType::Dynamic,
1034                );
1035            }
1036
1037            if compact_task.task_status == TaskStatus::Success {
1038                success_groups.push(compact_task.compaction_group_id);
1039            }
1040        }
1041
1042        trigger_compact_tasks_stat(
1043            &self.metrics,
1044            &tasks,
1045            &compaction.compaction_statuses,
1046            &versioning_guard.current_version,
1047        );
1048        drop(versioning_guard);
1049        if !success_groups.is_empty() {
1050            self.try_update_write_limits(&success_groups).await;
1051        }
1052        Ok(rets)
1053    }
1054
1055    /// Triggers compacitons to specified compaction groups.
1056    /// Don't wait for compaction finish
1057    pub async fn trigger_compaction_deterministic(
1058        &self,
1059        _base_version_id: HummockVersionId,
1060        compaction_groups: Vec<CompactionGroupId>,
1061    ) -> Result<()> {
1062        self.on_current_version(|old_version| {
1063            tracing::info!(
1064                "Trigger compaction for version {}, groups {:?}",
1065                old_version.id,
1066                compaction_groups
1067            );
1068            for &compaction_group in &compaction_groups {
1069                if old_version.levels.contains_key(&compaction_group) {
1070                    self.try_send_compaction_request(
1071                        compaction_group,
1072                        compact_task::TaskType::Dynamic,
1073                    );
1074                }
1075            }
1076        })
1077        .await;
1078
1079        Ok(())
1080    }
1081
1082    pub async fn trigger_manual_compaction(
1083        &self,
1084        compaction_group: CompactionGroupId,
1085        manual_compaction_option: ManualCompactionOption,
1086    ) -> Result<ManualCompactionTriggerResult> {
1087        let start_time = Instant::now();
1088        let exclusive = manual_compaction_option.exclusive;
1089
1090        // 1. Get idle compactor.
1091        let compactor = match self.compactor_manager.next_compactor() {
1092            Some(compactor) => compactor,
1093            None => {
1094                tracing::warn!("trigger_manual_compaction No compactor is available.");
1095                return Err(anyhow::anyhow!(
1096                    "trigger_manual_compaction No compactor is available. compaction_group {}",
1097                    compaction_group
1098                )
1099                .into());
1100            }
1101        };
1102
1103        // 2. Get manual compaction task.
1104        let compact_task = self
1105            .manual_get_compact_task_with_info(compaction_group, manual_compaction_option)
1106            .await;
1107        let (compact_task, blocked_by_pending) = match compact_task {
1108            Ok((compact_task, blocked_by_pending)) => (compact_task, blocked_by_pending),
1109            Err(err) => {
1110                tracing::warn!(error = %err.as_report(), "Failed to get compaction task");
1111                if matches!(err, Error::InvalidManualCompactionOption(_)) {
1112                    return Err(err);
1113                }
1114
1115                return Err(anyhow::anyhow!(err)
1116                    .context(format!(
1117                        "Failed to get compaction task for compaction_group {}",
1118                        compaction_group,
1119                    ))
1120                    .into());
1121            }
1122        };
1123        let compact_task = match compact_task {
1124            Some(compact_task) => compact_task,
1125            None => {
1126                if exclusive && blocked_by_pending {
1127                    return Ok(ManualCompactionTriggerResult::Retry);
1128                }
1129                // No compaction task available.
1130                return Err(anyhow::anyhow!(
1131                    "trigger_manual_compaction No compaction_task is available. compaction_group {}",
1132                    compaction_group
1133                )
1134                .into());
1135            }
1136        };
1137
1138        // 3. send task to compactor
1139        let task_id = compact_task.task_id;
1140        let compact_task_string = compact_task_to_string(&compact_task);
1141        tracing::info!(
1142            compact_task_string,
1143            duration = ?start_time.elapsed(),
1144            "Triggered manual compaction task."
1145        );
1146
1147        let report_rx = self.register_compaction_task_report_waiter(task_id);
1148        if let Err(err) = compactor
1149            .send_event(ResponseEvent::CompactTask(compact_task.into()))
1150            .with_context(|| {
1151                format!(
1152                    "Failed to trigger compaction task for compaction_group {}",
1153                    compaction_group,
1154                )
1155            })
1156        {
1157            self.remove_compaction_task_report_waiter(task_id);
1158            return Err(err.into());
1159        }
1160
1161        let report_result = match report_rx.await {
1162            Ok(result) => result,
1163            Err(_) => {
1164                self.remove_compaction_task_report_waiter(task_id);
1165                return Err(anyhow::anyhow!(
1166                    "trigger_manual_compaction wait report failed. compaction_group {}",
1167                    compaction_group
1168                )
1169                .into());
1170            }
1171        };
1172        if !report_result.reported {
1173            return Err(anyhow::anyhow!(
1174                "trigger_manual_compaction report not accepted. task_id {}",
1175                report_result.task_id
1176            )
1177            .into());
1178        }
1179
1180        if report_result.task_status == TaskStatus::NoAvailCpuResourceCanceled
1181            || report_result.task_status == TaskStatus::NoAvailMemoryResourceCanceled
1182        {
1183            return Ok(ManualCompactionTriggerResult::Retry);
1184        }
1185
1186        tracing::info!(
1187            ?report_result,
1188            duration = ?start_time.elapsed(),
1189            "Completed manual compaction task."
1190        );
1191
1192        Ok(ManualCompactionTriggerResult::Submitted)
1193    }
1194
1195    /// Sends a compaction request for new data or a topology change (clears cooldown).
1196    /// The caller must hold the version lock and ensure the group exists, so deletion cannot
1197    /// race with a late request. Lock order is `versioning` -> `compaction_state`.
1198    pub(crate) fn try_send_compaction_request(
1199        &self,
1200        compaction_group: CompactionGroupId,
1201        task_type: compact_task::TaskType,
1202    ) -> bool {
1203        self.compaction_state.try_sched_compaction(
1204            compaction_group,
1205            task_type,
1206            ScheduleTrigger::NewData,
1207        )
1208    }
1209
1210    /// Schedules all current groups, respecting the cooldown policy of the trigger.
1211    pub async fn trigger_compaction_for_all_groups(
1212        &self,
1213        task_type: TaskType,
1214        trigger: ScheduleTrigger,
1215    ) {
1216        // Keep membership stable until publication, so deletion cannot be followed by a late enqueue.
1217        let versioning = self
1218            .versioning
1219            .read_with_process_name("on_handle_trigger_multi_group")
1220            .await;
1221        #[cfg(test)]
1222        compaction_state_tests::before_trigger().await;
1223        for &group_id in versioning.current_version.levels.keys() {
1224            self.compaction_state
1225                .try_sched_compaction(group_id, task_type, trigger);
1226        }
1227    }
1228
1229    /// Apply `split_weight_by_vnode` based partition strategy.
1230    /// This handles dynamic partitioning based on table size and write throughput.
1231    fn apply_split_weight_by_vnode_partition(
1232        &self,
1233        compact_task: &mut CompactTask,
1234        compaction_config: &CompactionConfig,
1235        compact_table_ids: &[TableId],
1236    ) {
1237        if compaction_config.split_weight_by_vnode > 0 {
1238            for table_id in compact_table_ids {
1239                compact_task
1240                    .table_vnode_partition
1241                    .insert(*table_id, compact_task.split_weight_by_vnode);
1242            }
1243
1244            return;
1245        }
1246
1247        // Calculate per-table size from normalized input SSTs.
1248        let mut table_size_info: HashMap<TableId, u64> = HashMap::default();
1249        for input_ssts in &compact_task.input_ssts {
1250            for sst in &input_ssts.table_infos {
1251                for table_id in &sst.table_ids {
1252                    *table_size_info.entry(*table_id).or_default() +=
1253                        sst.sst_size / (sst.table_ids.len() as u64);
1254                }
1255            }
1256        }
1257
1258        let hybrid_vnode_count = self.env.opts.hybrid_partition_node_count;
1259        let default_partition_count = self.env.opts.partition_vnode_count;
1260        let compact_task_table_size_partition_threshold_low = self
1261            .env
1262            .opts
1263            .compact_task_table_size_partition_threshold_low;
1264        let compact_task_table_size_partition_threshold_high = self
1265            .env
1266            .opts
1267            .compact_task_table_size_partition_threshold_high;
1268
1269        // Check latest write throughput
1270        let table_write_throughput_statistic_manager =
1271            self.table_write_throughput_statistic_manager.read();
1272
1273        for (table_id, compact_table_size) in table_size_info {
1274            let write_throughput = table_write_throughput_statistic_manager
1275                .latest_table_throughput(table_id)
1276                .unwrap_or(0);
1277
1278            if compact_table_size > compact_task_table_size_partition_threshold_high
1279                && default_partition_count > 0
1280            {
1281                compact_task
1282                    .table_vnode_partition
1283                    .insert(table_id, default_partition_count);
1284            } else if (compact_table_size > compact_task_table_size_partition_threshold_low
1285                || (write_throughput > self.env.opts.table_high_write_throughput_threshold
1286                    && compact_table_size > compaction_config.target_file_size_base))
1287                && hybrid_vnode_count > 0
1288            {
1289                compact_task
1290                    .table_vnode_partition
1291                    .insert(table_id, hybrid_vnode_count);
1292            } else if compact_table_size > compaction_config.target_file_size_base {
1293                compact_task.table_vnode_partition.insert(table_id, 1);
1294            }
1295        }
1296
1297        compact_task
1298            .table_vnode_partition
1299            .retain(|table_id, _| compact_table_ids.contains(table_id));
1300    }
1301
1302    pub(crate) fn calculate_vnode_partition(
1303        &self,
1304        compact_task: &mut CompactTask,
1305        compaction_config: &CompactionConfig,
1306        compact_table_ids: &[TableId],
1307    ) {
1308        // Do not split sst by vnode partition when target_level > base_level
1309        // The purpose of data alignment is mainly to improve the parallelism of base level compaction
1310        // and reduce write amplification. However, at high level, the size of the sst file is often
1311        // larger and only contains the data of a single table_id, so there is no need to cut it.
1312        if compact_task.target_level > compact_task.base_level {
1313            return;
1314        }
1315
1316        // Apply split_weight_by_vnode based partition strategy
1317        self.apply_split_weight_by_vnode_partition(
1318            compact_task,
1319            compaction_config,
1320            compact_table_ids,
1321        );
1322    }
1323
1324    fn build_ready_compact_task(
1325        &self,
1326        picked_task: PickedCompactionTask,
1327        context: CompactTaskBuildContext,
1328        table_watermarks: &HashMap<TableId, Arc<TableWatermarks>>,
1329        all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1330    ) -> BuiltCompactTask {
1331        let compaction_config = context.compaction_config.clone();
1332        let (mut compact_task, compact_table_ids) = build_base_compact_task(picked_task, context);
1333
1334        if compact_task.is_trivial_reclaim() {
1335            compact_task.task_status = TaskStatus::Success;
1336            compact_task.sorted_output_ssts.clear();
1337            return BuiltCompactTask::MetaFinished(compact_task);
1338        }
1339
1340        if compact_task.is_trivial_move_task() {
1341            compact_task.task_status = TaskStatus::Success;
1342            compact_task.sorted_output_ssts = compact_task.input_ssts[0]
1343                .read_sstable_infos()
1344                .cloned()
1345                .collect();
1346            return BuiltCompactTask::MetaFinished(compact_task);
1347        }
1348
1349        self.prepare_compact_task_for_assignment(
1350            &mut compact_task,
1351            compaction_config.as_ref(),
1352            &compact_table_ids,
1353            safe_epoch_table_watermarks_impl(table_watermarks, &compact_table_ids),
1354            all_versioned_table_schemas,
1355        );
1356
1357        BuiltCompactTask::PendingAssignment(compact_task)
1358    }
1359
1360    fn prepare_compact_task_for_assignment(
1361        &self,
1362        compact_task: &mut CompactTask,
1363        compaction_config: &CompactionConfig,
1364        compact_table_ids: &[TableId],
1365        table_watermarks: BTreeMap<TableId, TableWatermarks>,
1366        all_versioned_table_schemas: &HashMap<TableId, Vec<i32>>,
1367    ) {
1368        self.calculate_vnode_partition(compact_task, compaction_config, compact_table_ids);
1369        attach_compact_task_table_metadata(
1370            compact_task,
1371            compact_table_ids,
1372            table_watermarks,
1373            all_versioned_table_schemas,
1374        );
1375    }
1376
1377    pub fn compactor_manager_ref(&self) -> crate::hummock::CompactorManagerRef {
1378        self.compactor_manager.clone()
1379    }
1380
1381    fn register_compaction_task_report_waiter(
1382        &self,
1383        task_id: HummockCompactionTaskId,
1384    ) -> Receiver<CompactionTaskReportResult> {
1385        let (tx, rx) = tokio::sync::oneshot::channel();
1386        self.compaction_task_report_notifiers
1387            .lock()
1388            .register(task_id, tx);
1389        rx
1390    }
1391
1392    fn remove_compaction_task_report_waiter(&self, task_id: HummockCompactionTaskId) {
1393        self.compaction_task_report_notifiers.lock().remove(task_id);
1394    }
1395
1396    fn notify_compaction_task_report_waiters(&self, results: Vec<CompactionTaskReportResult>) {
1397        let mut guard = self.compaction_task_report_notifiers.lock();
1398        for result in results {
1399            guard.notify(result);
1400        }
1401    }
1402}
1403
1404#[cfg(any(test, feature = "test"))]
1405impl HummockManager {
1406    pub async fn compaction_task_from_assignment_for_test(
1407        &self,
1408        task_id: u64,
1409    ) -> Option<CompactTaskAssignment> {
1410        let compaction_guard = self
1411            .compaction
1412            .read_with_process_name("compaction_task_from_assignment_for_test")
1413            .await;
1414        let assignment_ref = &compaction_guard.compact_task_assignment;
1415        assignment_ref.get(&task_id).cloned()
1416    }
1417
1418    pub async fn report_compact_task_for_test(
1419        &self,
1420        task_id: u64,
1421        compact_task: Option<CompactTask>,
1422        task_status: TaskStatus,
1423        sorted_output_ssts: Vec<SstableInfo>,
1424        table_stats_change: Option<PbTableStatsMap>,
1425    ) -> Result<()> {
1426        if let Some(task) = compact_task {
1427            let mut guard = self
1428                .compaction
1429                .write_with_process_name("report_compact_task_for_test")
1430                .await;
1431            guard.compact_task_assignment.insert(
1432                task_id,
1433                CompactTaskAssignment {
1434                    compact_task: task,
1435                    context_id: 0.into(),
1436                },
1437            );
1438        }
1439
1440        // In the test, the contents of the compact task may have been modified directly, while the contents of compact_task_assignment were not modified.
1441        // So we pass the modified compact_task directly into the `report_compact_task_impl`
1442        self.report_compact_tasks(vec![ReportTask {
1443            task_id,
1444            task_status,
1445            sorted_output_ssts,
1446            table_stats_change: table_stats_change.unwrap_or_default(),
1447            object_timestamps: HashMap::default(),
1448        }])
1449        .await?;
1450        Ok(())
1451    }
1452}
1453
1454/// What triggered the compaction schedule request.
1455#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1456pub enum ScheduleTrigger {
1457    /// New data arrived or group topology changed. Clears cooldown.
1458    NewData,
1459    /// Periodic timer. Respects cooldown for Dynamic type.
1460    Periodic,
1461}
1462
1463/// A point-in-time snapshot of the compaction schedule state.
1464///
1465/// `generation` is used by `unschedule()` to detect whether new data arrived
1466/// after the snapshot was taken, preserving newer requests and preventing incorrect cooldown.
1467pub struct CompactionScheduleSnapshot {
1468    scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1469    generation: u64,
1470}
1471
1472impl CompactionScheduleSnapshot {
1473    /// Task type priority order for scheduling (checked first = higher priority).
1474    const TASK_TYPE_PRIORITY: &[TaskType] = &[
1475        TaskType::Dynamic,
1476        TaskType::SpaceReclaim,
1477        TaskType::Ttl,
1478        TaskType::Tombstone,
1479        TaskType::VnodeWatermark,
1480    ];
1481
1482    pub fn generation(&self) -> u64 {
1483        self.generation
1484    }
1485
1486    /// Pick compaction groups and task type from this snapshot.
1487    ///
1488    /// Returns groups in shuffled order. Non-Dynamic types have higher priority
1489    /// and return a single group; Dynamic groups are batched together.
1490    pub fn pick_compaction_groups_and_type(&self) -> Option<(Vec<CompactionGroupId>, TaskType)> {
1491        let group_ids = self.group_ids_shuffled();
1492        let mut normal_groups = vec![];
1493        for cg_id in group_ids {
1494            if let Some(pick_type) = self.pick_type(cg_id) {
1495                if pick_type == TaskType::Dynamic {
1496                    normal_groups.push(cg_id);
1497                } else if normal_groups.is_empty() {
1498                    return Some((vec![cg_id], pick_type));
1499                }
1500            }
1501        }
1502        if normal_groups.is_empty() {
1503            None
1504        } else {
1505            Some((normal_groups, TaskType::Dynamic))
1506        }
1507    }
1508
1509    fn group_ids_shuffled(&self) -> Vec<CompactionGroupId> {
1510        let mut group_ids: Vec<_> = self.scheduled.iter().map(|(g, _)| *g).unique().collect();
1511        group_ids.shuffle(&mut thread_rng());
1512        group_ids
1513    }
1514
1515    fn pick_type(&self, group: CompactionGroupId) -> Option<TaskType> {
1516        Self::TASK_TYPE_PRIORITY
1517            .iter()
1518            .find(|t| self.scheduled.contains(&(group, **t)))
1519            .copied()
1520    }
1521}
1522
1523/// Tracks which (`compaction_group`, `task_type`) pairs are scheduled for compaction.
1524///
1525/// For `Dynamic` type, includes a cooldown mechanism: groups with no compaction work
1526/// are skipped by periodic triggers until new data arrives via `commit_epoch`.
1527#[derive(Debug, Default)]
1528pub struct CompactionState {
1529    inner: Mutex<CompactionStateInner>,
1530}
1531
1532#[derive(Debug, Default)]
1533struct CompactionStateInner {
1534    generation: u64,
1535    scheduled: HashSet<(CompactionGroupId, compact_task::TaskType)>,
1536    /// Groups skipped by periodic Dynamic trigger until new data arrives.
1537    dynamic_cooldown: HashSet<CompactionGroupId>,
1538    /// Tracks new-data generations so old picker results cannot remove newer requests.
1539    last_new_data_generation: HashMap<CompactionGroupId, u64>,
1540}
1541
1542impl CompactionState {
1543    pub fn new() -> Self {
1544        Self {
1545            inner: Default::default(),
1546        }
1547    }
1548
1549    /// Enqueues a compaction request. Returns `true` if newly scheduled.
1550    ///
1551    /// `trigger` only affects `Dynamic` type — see [`ScheduleTrigger`].
1552    pub fn try_sched_compaction(
1553        &self,
1554        compaction_group: CompactionGroupId,
1555        task_type: TaskType,
1556        trigger: ScheduleTrigger,
1557    ) -> bool {
1558        let mut guard = self.inner.lock();
1559        if task_type == TaskType::Dynamic {
1560            match trigger {
1561                ScheduleTrigger::NewData => {
1562                    guard.dynamic_cooldown.remove(&compaction_group);
1563                    // A coalesced request still invalidates an older picker's no-task result.
1564                    guard.generation += 1;
1565                    let generation = guard.generation;
1566                    guard
1567                        .last_new_data_generation
1568                        .insert(compaction_group, generation);
1569                }
1570                ScheduleTrigger::Periodic => {
1571                    if guard.dynamic_cooldown.contains(&compaction_group) {
1572                        return false;
1573                    }
1574                }
1575            }
1576        }
1577        guard.scheduled.insert((compaction_group, task_type))
1578    }
1579
1580    /// Removes a scheduled entry. For Dynamic type, adds to cooldown unless
1581    /// a newer request exists than the snapshot generation.
1582    pub fn unschedule(
1583        &self,
1584        compaction_group: CompactionGroupId,
1585        task_type: compact_task::TaskType,
1586        generation: u64,
1587    ) {
1588        let mut guard = self.inner.lock();
1589        if task_type == TaskType::Dynamic
1590            && guard
1591                .last_new_data_generation
1592                .get(&compaction_group)
1593                .is_some_and(|g| *g > generation)
1594        {
1595            return;
1596        }
1597        // Removal can race with an old picker result. Do not recreate state for an absent entry.
1598        if guard.scheduled.remove(&(compaction_group, task_type)) && task_type == TaskType::Dynamic
1599        {
1600            guard.dynamic_cooldown.insert(compaction_group);
1601        }
1602    }
1603
1604    /// Takes a snapshot of the current schedule state.
1605    pub fn snapshot(&self) -> CompactionScheduleSnapshot {
1606        let guard = self.inner.lock();
1607        let generation = guard.generation;
1608        CompactionScheduleSnapshot {
1609            scheduled: guard.scheduled.clone(),
1610            generation,
1611        }
1612    }
1613
1614    /// Removes all schedule state for a deleted or merged group.
1615    pub fn remove_compaction_group(&self, compaction_group: CompactionGroupId) {
1616        let mut guard = self.inner.lock();
1617        guard
1618            .scheduled
1619            .retain(|(group, _)| *group != compaction_group);
1620        guard.dynamic_cooldown.remove(&compaction_group);
1621        guard.last_new_data_generation.remove(&compaction_group);
1622    }
1623}
1624
1625impl Compaction {
1626    pub fn get_compact_task_assignments_by_group_id(
1627        &self,
1628        compaction_group_id: CompactionGroupId,
1629    ) -> Vec<CompactTaskAssignment> {
1630        self.compact_task_assignment
1631            .values()
1632            .filter_map(|assignment| {
1633                if assignment.compact_task.compaction_group_id == compaction_group_id {
1634                    Some(assignment.clone())
1635                } else {
1636                    None
1637                }
1638            })
1639            .collect()
1640    }
1641}
1642
1643#[derive(Clone, Default)]
1644pub struct CompactionGroupStatistic {
1645    pub group_id: CompactionGroupId,
1646    pub group_size: u64,
1647    pub table_statistic: BTreeMap<StateTableId, u64>,
1648    pub compaction_group_config: CompactionGroup,
1649}
1650
1651/// Updates table stats caused by vnode watermark trivial reclaim compaction.
1652fn update_table_stats_for_vnode_watermark_trivial_reclaim(
1653    table_stats: &mut PbTableStatsMap,
1654    task: &CompactTask,
1655) {
1656    if task.task_type != TaskType::VnodeWatermark {
1657        return;
1658    }
1659    let mut deleted_table_keys: HashMap<TableId, u64> = HashMap::default();
1660    for s in task.input_ssts.iter().flat_map(|l| l.table_infos.iter()) {
1661        assert_eq!(s.table_ids.len(), 1);
1662        let e = deleted_table_keys.entry(s.table_ids[0]).or_insert(0);
1663        *e += s.total_key_count;
1664    }
1665    for (table_id, delete_count) in deleted_table_keys {
1666        let Some(stats) = table_stats.get_mut(&table_id) else {
1667            continue;
1668        };
1669        if stats.total_key_count == 0 {
1670            continue;
1671        }
1672        let new_total_key_count = stats.total_key_count.saturating_sub(delete_count as i64);
1673        let ratio = new_total_key_count as f64 / stats.total_key_count as f64;
1674        // total_key_count is updated accurately.
1675        stats.total_key_count = new_total_key_count;
1676        // others are updated approximately.
1677        stats.total_key_size = (stats.total_key_size as f64 * ratio).ceil() as i64;
1678        stats.total_value_size = (stats.total_value_size as f64 * ratio).ceil() as i64;
1679    }
1680}
1681
1682#[derive(Debug, Clone)]
1683pub enum GroupState {
1684    /// The compaction group is not in emergency state.
1685    Normal,
1686
1687    /// The compaction group is in emergency state.
1688    Emergency(String), // reason
1689
1690    /// The compaction group is in write stop state.
1691    WriteStop(String), // reason
1692}
1693
1694impl GroupState {
1695    pub fn is_write_stop(&self) -> bool {
1696        matches!(self, Self::WriteStop(_))
1697    }
1698
1699    pub fn is_emergency(&self) -> bool {
1700        matches!(self, Self::Emergency(_))
1701    }
1702
1703    pub fn reason(&self) -> Option<&str> {
1704        match self {
1705            Self::Emergency(reason) | Self::WriteStop(reason) => Some(reason),
1706            _ => None,
1707        }
1708    }
1709}
1710
1711#[derive(Clone, Default)]
1712pub struct GroupStateValidator;
1713
1714impl GroupStateValidator {
1715    pub fn write_stop_sub_level_count(
1716        level_count: usize,
1717        compaction_config: &CompactionConfig,
1718    ) -> bool {
1719        let threshold = compaction_config.level0_stop_write_threshold_sub_level_number as usize;
1720        level_count > threshold
1721    }
1722
1723    pub fn write_stop_l0_size(l0_size: u64, compaction_config: &CompactionConfig) -> bool {
1724        l0_size
1725            > compaction_config
1726                .level0_stop_write_threshold_max_size
1727                .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1728    }
1729
1730    pub fn write_stop_l0_file_count(
1731        l0_file_count: usize,
1732        compaction_config: &CompactionConfig,
1733    ) -> bool {
1734        l0_file_count
1735            > compaction_config
1736                .level0_stop_write_threshold_max_sst_count
1737                .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1738                as usize
1739    }
1740
1741    pub fn emergency_l0_file_count(
1742        l0_file_count: usize,
1743        compaction_config: &CompactionConfig,
1744    ) -> bool {
1745        l0_file_count
1746            > compaction_config
1747                .emergency_level0_sst_file_count
1748                .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1749                as usize
1750    }
1751
1752    pub fn emergency_l0_partition_count(
1753        last_l0_sub_level_partition_count: usize,
1754        compaction_config: &CompactionConfig,
1755    ) -> bool {
1756        last_l0_sub_level_partition_count
1757            > compaction_config
1758                .emergency_level0_sub_level_partition
1759                .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1760                as usize
1761    }
1762
1763    pub fn check_single_group_write_stop(
1764        levels: &Levels,
1765        compaction_config: &CompactionConfig,
1766    ) -> GroupState {
1767        if Self::write_stop_sub_level_count(levels.l0.sub_levels.len(), compaction_config) {
1768            return GroupState::WriteStop(format!(
1769                "WriteStop(l0_level_count: {}, threshold: {}) too many L0 sub levels",
1770                levels.l0.sub_levels.len(),
1771                compaction_config.level0_stop_write_threshold_sub_level_number
1772            ));
1773        }
1774
1775        if Self::write_stop_l0_file_count(
1776            levels
1777                .l0
1778                .sub_levels
1779                .iter()
1780                .map(|l| l.table_infos.len())
1781                .sum(),
1782            compaction_config,
1783        ) {
1784            return GroupState::WriteStop(format!(
1785                "WriteStop(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1786                levels
1787                    .l0
1788                    .sub_levels
1789                    .iter()
1790                    .map(|l| l.table_infos.len())
1791                    .sum::<usize>(),
1792                compaction_config
1793                    .level0_stop_write_threshold_max_sst_count
1794                    .unwrap_or(compaction_config::level0_stop_write_threshold_max_sst_count())
1795            ));
1796        }
1797
1798        if Self::write_stop_l0_size(levels.l0.total_file_size, compaction_config) {
1799            return GroupState::WriteStop(format!(
1800                "WriteStop(l0_size: {}, threshold: {}) too large L0 size",
1801                levels.l0.total_file_size,
1802                compaction_config
1803                    .level0_stop_write_threshold_max_size
1804                    .unwrap_or(compaction_config::level0_stop_write_threshold_max_size())
1805            ));
1806        }
1807
1808        GroupState::Normal
1809    }
1810
1811    pub fn check_single_group_emergency(
1812        levels: &Levels,
1813        compaction_config: &CompactionConfig,
1814    ) -> GroupState {
1815        if Self::emergency_l0_file_count(
1816            levels
1817                .l0
1818                .sub_levels
1819                .iter()
1820                .map(|l| l.table_infos.len())
1821                .sum(),
1822            compaction_config,
1823        ) {
1824            return GroupState::Emergency(format!(
1825                "Emergency(l0_sst_count: {}, threshold: {}) too many L0 sst files",
1826                levels
1827                    .l0
1828                    .sub_levels
1829                    .iter()
1830                    .map(|l| l.table_infos.len())
1831                    .sum::<usize>(),
1832                compaction_config
1833                    .emergency_level0_sst_file_count
1834                    .unwrap_or(compaction_config::emergency_level0_sst_file_count())
1835            ));
1836        }
1837
1838        if Self::emergency_l0_partition_count(
1839            levels
1840                .l0
1841                .sub_levels
1842                .first()
1843                .map(|l| l.table_infos.len())
1844                .unwrap_or(0),
1845            compaction_config,
1846        ) {
1847            return GroupState::Emergency(format!(
1848                "Emergency(l0_partition_count: {}, threshold: {}) too many L0 partitions",
1849                levels
1850                    .l0
1851                    .sub_levels
1852                    .first()
1853                    .map(|l| l.table_infos.len())
1854                    .unwrap_or(0),
1855                compaction_config
1856                    .emergency_level0_sub_level_partition
1857                    .unwrap_or(compaction_config::emergency_level0_sub_level_partition())
1858            ));
1859        }
1860
1861        GroupState::Normal
1862    }
1863
1864    pub fn group_state(levels: &Levels, compaction_config: &CompactionConfig) -> GroupState {
1865        let state = Self::check_single_group_write_stop(levels, compaction_config);
1866        if state.is_write_stop() {
1867            return state;
1868        }
1869
1870        Self::check_single_group_emergency(levels, compaction_config)
1871    }
1872}
1873
1874#[cfg(test)]
1875mod compaction_state_tests {
1876    use risingwave_pb::hummock::compact_task::TaskType;
1877
1878    use super::*;
1879
1880    tokio::task_local! {
1881        static BEFORE_TRIGGER: Arc<tokio::sync::Barrier>;
1882    }
1883
1884    pub(super) async fn before_trigger() {
1885        if let Ok(barrier) = BEFORE_TRIGGER.try_with(Arc::clone) {
1886            barrier.wait().await;
1887            barrier.wait().await;
1888        }
1889    }
1890
1891    #[tokio::test]
1892    #[cfg(not(madsim))]
1893    async fn test_trigger_all_groups_serializes_with_deletion() {
1894        use std::time::Duration;
1895
1896        use crate::hummock::test_utils::setup_compute_env;
1897
1898        for trigger in [ScheduleTrigger::NewData, ScheduleTrigger::Periodic] {
1899            let (_, manager, _, _) = setup_compute_env(80).await;
1900            manager
1901                .register_table_ids_for_test(&[(100, 2.into()), (101, 2.into())])
1902                .await
1903                .unwrap();
1904            let (group, _) = manager
1905                .move_state_tables_to_dedicated_compaction_group(2.into(), &[101.into()], None)
1906                .await
1907                .unwrap();
1908            let snapshot = manager.compaction_state.snapshot();
1909            manager
1910                .compaction_state
1911                .unschedule(2.into(), TaskType::Dynamic, snapshot.generation());
1912            let barrier = Arc::new(tokio::sync::Barrier::new(2));
1913            let delete = async {
1914                barrier.wait().await;
1915                {
1916                    let write = tokio::task::unconstrained(manager.versioning.write());
1917                    tokio::pin!(write);
1918                    assert!(
1919                        futures::poll!(write.as_mut()).is_pending(),
1920                        "group deletion must wait until candidate publication finishes"
1921                    );
1922                }
1923                barrier.wait().await;
1924                manager.unregister_table_ids([101.into()]).await.unwrap();
1925            };
1926            tokio::time::timeout(Duration::from_secs(10), async {
1927                tokio::join!(
1928                    BEFORE_TRIGGER.scope(
1929                        barrier.clone(),
1930                        manager.trigger_compaction_for_all_groups(TaskType::Dynamic, trigger),
1931                    ),
1932                    delete,
1933                )
1934            })
1935            .await
1936            .unwrap();
1937            assert!(!manager.compaction_group_ids().await.contains(&group));
1938            assert_eq!(
1939                manager
1940                    .compaction_state
1941                    .snapshot()
1942                    .scheduled
1943                    .contains(&(2.into(), TaskType::Dynamic)),
1944                trigger == ScheduleTrigger::NewData,
1945                "registration wakes cooled groups; periodic scans must respect cooldown"
1946            );
1947            assert!(
1948                !manager
1949                    .compaction_state
1950                    .snapshot()
1951                    .scheduled
1952                    .iter()
1953                    .any(|(id, _)| *id == group)
1954            );
1955        }
1956    }
1957
1958    #[test]
1959    fn test_basic_schedule_and_unschedule() {
1960        let state = CompactionState::new();
1961        let group_id: CompactionGroupId = 1.into();
1962
1963        // First schedule should succeed
1964        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1965        // Duplicate schedule should fail
1966        assert!(!state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1967        // Different task type should succeed
1968        assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
1969
1970        // Snapshot should contain both
1971        let snapshot = state.snapshot();
1972        assert!(snapshot.scheduled.contains(&(group_id, TaskType::Dynamic)));
1973        assert!(snapshot.scheduled.contains(&(group_id, TaskType::Ttl)));
1974
1975        // Unschedule removes from scheduled set
1976        state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
1977        let snapshot2 = state.snapshot();
1978        assert!(!snapshot2.scheduled.contains(&(group_id, TaskType::Dynamic)));
1979        assert!(snapshot2.scheduled.contains(&(group_id, TaskType::Ttl)));
1980    }
1981
1982    #[test]
1983    fn test_cooldown_blocks_periodic_trigger() {
1984        let state = CompactionState::new();
1985        let group_id: CompactionGroupId = 1.into();
1986
1987        // Schedule then unschedule - should add to cooldown
1988        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
1989        let snapshot = state.snapshot();
1990        state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
1991
1992        // Verify in cooldown
1993        assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
1994
1995        // Periodic trigger should be blocked
1996        assert!(!state.try_sched_compaction(
1997            group_id,
1998            TaskType::Dynamic,
1999            ScheduleTrigger::Periodic
2000        ));
2001    }
2002
2003    #[test]
2004    fn test_new_data_clears_cooldown() {
2005        let state = CompactionState::new();
2006        let group_id: CompactionGroupId = 1.into();
2007
2008        // Put group in cooldown
2009        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2010        let snapshot = state.snapshot();
2011        state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2012        assert!(state.inner.lock().dynamic_cooldown.contains(&group_id));
2013
2014        // NewData trigger should clear cooldown and schedule
2015        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2016        assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
2017    }
2018
2019    #[test]
2020    fn test_cooldown_only_affects_dynamic_type() {
2021        let state = CompactionState::new();
2022        let group_id: CompactionGroupId = 1.into();
2023
2024        // Put group in cooldown for Dynamic
2025        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2026        let snapshot = state.snapshot();
2027        state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2028
2029        // Ttl unschedule should NOT add to cooldown
2030        let group_id_2: CompactionGroupId = 2.into();
2031        assert!(state.try_sched_compaction(group_id_2, TaskType::Ttl, ScheduleTrigger::Periodic));
2032        let snapshot2 = state.snapshot();
2033        state.unschedule(group_id_2, TaskType::Ttl, snapshot2.generation());
2034        assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id_2));
2035
2036        // Other task types should work regardless of cooldown
2037        assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
2038        assert!(state.try_sched_compaction(
2039            group_id,
2040            TaskType::SpaceReclaim,
2041            ScheduleTrigger::Periodic
2042        ));
2043    }
2044
2045    #[test]
2046    fn test_race_condition_new_data_after_snapshot() {
2047        let state = CompactionState::new();
2048        let group_id: CompactionGroupId = 1.into();
2049
2050        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2051        let snapshot = state.snapshot();
2052
2053        // A duplicate request changes the generation even though set membership stays the same.
2054        assert!(!state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2055        let latest = state.snapshot();
2056        state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2057        assert!(
2058            state
2059                .snapshot()
2060                .scheduled
2061                .contains(&(group_id, TaskType::Dynamic))
2062        );
2063        assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
2064
2065        // Another group's newer request must not prevent this group's fresh result from cooling it.
2066        state.try_sched_compaction(2.into(), TaskType::Dynamic, ScheduleTrigger::NewData);
2067        state.unschedule(group_id, TaskType::Dynamic, latest.generation());
2068        assert!(
2069            !state
2070                .snapshot()
2071                .scheduled
2072                .contains(&(group_id, TaskType::Dynamic))
2073        );
2074        assert!(!state.try_sched_compaction(
2075            group_id,
2076            TaskType::Dynamic,
2077            ScheduleTrigger::Periodic
2078        ));
2079
2080        // Defensive API-level ABA: remove -> requeue -> delayed old result. This does not
2081        // assert that production has concurrent pickers or reuses deleted group IDs.
2082        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2083        state.unschedule(group_id, TaskType::Dynamic, latest.generation());
2084        assert!(
2085            state
2086                .snapshot()
2087                .scheduled
2088                .contains(&(group_id, TaskType::Dynamic))
2089        );
2090        assert!(!state.inner.lock().dynamic_cooldown.contains(&group_id));
2091    }
2092
2093    #[test]
2094    fn test_remove_compaction_group_cleans_all_state() {
2095        let state = CompactionState::new();
2096        let group_id: CompactionGroupId = 1.into();
2097
2098        // Set up state
2099        assert!(state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData));
2100        assert!(state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic));
2101        state.inner.lock().dynamic_cooldown.insert(group_id);
2102
2103        // Remove group
2104        let snapshot = state.snapshot();
2105        state.remove_compaction_group(group_id);
2106
2107        // Verify all state cleaned up
2108        let guard = state.inner.lock();
2109        assert!(!guard.scheduled.contains(&(group_id, TaskType::Dynamic)));
2110        assert!(!guard.scheduled.contains(&(group_id, TaskType::Ttl)));
2111        assert!(!guard.dynamic_cooldown.contains(&group_id));
2112        assert!(!guard.last_new_data_generation.contains_key(&group_id));
2113        drop(guard);
2114
2115        // A delayed picker result must not recreate any state for the removed group.
2116        state.unschedule(group_id, TaskType::Dynamic, snapshot.generation());
2117        let guard = state.inner.lock();
2118        assert!(guard.scheduled.is_empty());
2119        assert!(guard.dynamic_cooldown.is_empty());
2120        assert!(guard.last_new_data_generation.is_empty());
2121    }
2122
2123    #[test]
2124    fn test_snapshot_pick_type_priority() {
2125        let state = CompactionState::new();
2126        let group_id: CompactionGroupId = 1.into();
2127
2128        // Empty group returns None
2129        assert_eq!(state.snapshot().pick_type(group_id), None);
2130
2131        // Priority order: Dynamic > SpaceReclaim > Ttl > Tombstone > VnodeWatermark
2132        state.try_sched_compaction(
2133            group_id,
2134            TaskType::VnodeWatermark,
2135            ScheduleTrigger::Periodic,
2136        );
2137        assert_eq!(
2138            state.snapshot().pick_type(group_id),
2139            Some(TaskType::VnodeWatermark)
2140        );
2141
2142        state.try_sched_compaction(group_id, TaskType::Tombstone, ScheduleTrigger::Periodic);
2143        assert_eq!(
2144            state.snapshot().pick_type(group_id),
2145            Some(TaskType::Tombstone)
2146        );
2147
2148        state.try_sched_compaction(group_id, TaskType::Ttl, ScheduleTrigger::Periodic);
2149        assert_eq!(state.snapshot().pick_type(group_id), Some(TaskType::Ttl));
2150
2151        state.try_sched_compaction(group_id, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2152        assert_eq!(
2153            state.snapshot().pick_type(group_id),
2154            Some(TaskType::SpaceReclaim)
2155        );
2156
2157        state.try_sched_compaction(group_id, TaskType::Dynamic, ScheduleTrigger::NewData);
2158        assert_eq!(
2159            state.snapshot().pick_type(group_id),
2160            Some(TaskType::Dynamic)
2161        );
2162    }
2163
2164    #[test]
2165    fn test_multiple_groups_independent_cooldown() {
2166        let state = CompactionState::new();
2167        let g1: CompactionGroupId = 1.into();
2168        let g2: CompactionGroupId = 2.into();
2169
2170        state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2171        state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2172        let snapshot = state.snapshot();
2173
2174        // Only unschedule g1
2175        state.unschedule(g1, TaskType::Dynamic, snapshot.generation());
2176
2177        let guard = state.inner.lock();
2178        assert!(guard.dynamic_cooldown.contains(&g1));
2179        assert!(!guard.dynamic_cooldown.contains(&g2));
2180    }
2181
2182    #[test]
2183    fn test_pick_compaction_groups_empty() {
2184        let state = CompactionState::new();
2185        let snapshot = state.snapshot();
2186        // No scheduled groups → returns None
2187        assert!(snapshot.pick_compaction_groups_and_type().is_none());
2188    }
2189
2190    #[test]
2191    fn test_pick_compaction_groups_mixed_types() {
2192        let state = CompactionState::new();
2193        let g1: CompactionGroupId = 1.into();
2194        let g2: CompactionGroupId = 2.into();
2195        let g3: CompactionGroupId = 3.into();
2196
2197        // g1: Dynamic, g2: Ttl, g3: Dynamic
2198        state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2199        state.try_sched_compaction(g2, TaskType::Ttl, ScheduleTrigger::Periodic);
2200        state.try_sched_compaction(g3, TaskType::Dynamic, ScheduleTrigger::NewData);
2201
2202        let snapshot = state.snapshot();
2203        let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2204
2205        // Due to shuffle, either:
2206        // - Ttl group is encountered first → returns (vec![g2], Ttl)
2207        // - Dynamic group is encountered first → collects all Dynamic, skips Ttl
2208        //   → returns ([g1, g3] in some order, Dynamic)
2209        if task_type == TaskType::Dynamic {
2210            assert!(groups.contains(&g1));
2211            assert!(groups.contains(&g3));
2212            assert!(!groups.contains(&g2)); // Ttl group excluded from Dynamic result
2213        } else {
2214            assert_eq!(task_type, TaskType::Ttl);
2215            assert_eq!(groups, vec![g2]);
2216        }
2217    }
2218
2219    #[test]
2220    fn test_pick_compaction_groups_all_dynamic() {
2221        let state = CompactionState::new();
2222        let g1: CompactionGroupId = 1.into();
2223        let g2: CompactionGroupId = 2.into();
2224
2225        state.try_sched_compaction(g1, TaskType::Dynamic, ScheduleTrigger::NewData);
2226        state.try_sched_compaction(g2, TaskType::Dynamic, ScheduleTrigger::NewData);
2227
2228        let snapshot = state.snapshot();
2229        let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2230        assert_eq!(task_type, TaskType::Dynamic);
2231        assert!(groups.contains(&g1));
2232        assert!(groups.contains(&g2));
2233    }
2234
2235    #[test]
2236    fn test_pick_compaction_groups_single_non_dynamic() {
2237        let state = CompactionState::new();
2238        let g1: CompactionGroupId = 1.into();
2239
2240        state.try_sched_compaction(g1, TaskType::SpaceReclaim, ScheduleTrigger::Periodic);
2241
2242        let snapshot = state.snapshot();
2243        let (groups, task_type) = snapshot.pick_compaction_groups_and_type().unwrap();
2244        assert_eq!(task_type, TaskType::SpaceReclaim);
2245        assert_eq!(groups, vec![g1]);
2246    }
2247}