Skip to main content

risingwave_meta/hummock/manager/compaction/
compaction_event_loop.rs

1// Copyright 2025 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::{HashMap, HashSet};
16use std::sync::Arc;
17use std::time::{Instant, SystemTime};
18
19use anyhow::Context;
20use futures::stream::FuturesUnordered;
21use futures::{FutureExt, StreamExt};
22use risingwave_hummock_sdk::compact_task::ReportTask;
23use risingwave_hummock_sdk::{CompactionGroupId, HummockContextId};
24use risingwave_pb::hummock::compact_task::{TaskStatus, TaskType};
25use risingwave_pb::hummock::subscribe_compaction_event_request::{
26    Event as RequestEvent, HeartBeat, PullTask,
27};
28use risingwave_pb::hummock::subscribe_compaction_event_response::{
29    Event as ResponseEvent, PullTaskAck,
30};
31use risingwave_pb::hummock::{CompactTaskProgress, SubscribeCompactionEventRequest};
32use risingwave_pb::iceberg_compaction::SubscribeIcebergCompactionEventRequest;
33use risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_request::{
34    Event as IcebergRequestEvent, PullTask as IcebergPullTask, ReportTask as IcebergReportTask,
35};
36use risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_response::{
37    Event as IcebergResponseEvent, PullTaskAck as IcebergPullTaskAck,
38};
39use rw_futures_util::pending_on_none;
40use thiserror_ext::AsReport;
41use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel};
42use tokio::sync::oneshot::{Receiver as OneShotReceiver, Sender};
43use tokio::task::JoinHandle;
44use tonic::Streaming;
45use tracing::warn;
46
47use super::init_selectors;
48use crate::hummock::HummockManager;
49use crate::hummock::compaction::CompactionSelector;
50use crate::hummock::compactor_manager::Compactor;
51use crate::hummock::error::{Error, Result};
52use crate::hummock::sequence::next_compaction_task_id;
53use crate::manager::MetaOpts;
54use crate::manager::iceberg_compaction::IcebergCompactionManagerRef;
55use crate::rpc::metrics::MetaMetrics;
56
57const MAX_SKIP_TIMES: usize = 8;
58const MAX_REPORT_COUNT: usize = 16;
59
60#[async_trait::async_trait]
61pub trait CompactionEventDispatcher: Send + Sync + 'static {
62    type EventType: Send + Sync + 'static;
63
64    async fn on_event_locally(&self, context_id: HummockContextId, event: Self::EventType) -> bool;
65
66    async fn on_event_remotely(
67        &self,
68        context_id: HummockContextId,
69        event: Self::EventType,
70    ) -> Result<()>;
71
72    fn should_forward(&self, event: &Self::EventType) -> bool;
73
74    fn remove_compactor(&self, context_id: HummockContextId);
75}
76
77pub trait CompactorStreamEvent: Send + Sync + 'static {
78    type EventType: Send + Sync + 'static;
79    fn take_event(self) -> Self::EventType;
80    fn create_at(&self) -> u64;
81}
82
83pub struct CompactionEventLoop<
84    D: CompactionEventDispatcher<EventType = E::EventType>,
85    E: CompactorStreamEvent,
86> {
87    hummock_compactor_dispatcher: D,
88    metrics: Arc<MetaMetrics>,
89    compactor_streams_change_rx: UnboundedReceiver<(HummockContextId, Streaming<E>)>,
90}
91
92pub type HummockCompactionEventLoop =
93    CompactionEventLoop<HummockCompactionEventDispatcher, SubscribeCompactionEventRequest>;
94
95pub type IcebergCompactionEventLoop =
96    CompactionEventLoop<IcebergCompactionEventDispatcher, SubscribeIcebergCompactionEventRequest>;
97
98pub struct HummockCompactionEventDispatcher {
99    meta_opts: Arc<MetaOpts>,
100    hummock_compaction_event_handler: HummockCompactionEventHandler,
101    tx: Option<UnboundedSender<(HummockContextId, RequestEvent)>>,
102}
103
104#[async_trait::async_trait]
105impl CompactionEventDispatcher for HummockCompactionEventDispatcher {
106    type EventType = RequestEvent;
107
108    fn should_forward(&self, event: &Self::EventType) -> bool {
109        if self.tx.is_none() {
110            return false;
111        }
112
113        matches!(event, RequestEvent::PullTask(_)) || matches!(event, RequestEvent::ReportTask(_))
114    }
115
116    async fn on_event_locally(&self, context_id: HummockContextId, event: Self::EventType) -> bool {
117        let mut compactor_alive = true;
118        match event {
119            RequestEvent::HeartBeat(HeartBeat { progress }) => {
120                compactor_alive = self
121                    .hummock_compaction_event_handler
122                    .handle_heartbeat(context_id, progress)
123                    .await;
124            }
125
126            RequestEvent::Register(_event) => {
127                unreachable!()
128            }
129
130            RequestEvent::PullTask(pull_task) => {
131                compactor_alive = self
132                    .hummock_compaction_event_handler
133                    .handle_pull_task_event(
134                        context_id,
135                        pull_task.pull_task_count as usize,
136                        &mut init_selectors(),
137                        self.meta_opts.max_get_task_probe_times,
138                    )
139                    .await;
140            }
141
142            RequestEvent::ReportTask(report_event) => {
143                if let Err(e) = self
144                    .hummock_compaction_event_handler
145                    .handle_report_task_event(vec![report_event.into()])
146                    .await
147                {
148                    tracing::error!(error = %e.as_report(), "report compact_tack fail")
149                }
150            }
151        }
152
153        compactor_alive
154    }
155
156    async fn on_event_remotely(
157        &self,
158        context_id: HummockContextId,
159        event: Self::EventType,
160    ) -> Result<()> {
161        if let Some(tx) = &self.tx {
162            tx.send((context_id, event))
163                .with_context(|| format!("Failed to send event to compactor {context_id}"))?;
164        } else {
165            unreachable!();
166        }
167        Ok(())
168    }
169
170    fn remove_compactor(&self, context_id: HummockContextId) {
171        self.hummock_compaction_event_handler
172            .hummock_manager
173            .compactor_manager
174            .remove_compactor(context_id);
175    }
176}
177
178impl HummockCompactionEventDispatcher {
179    pub fn new(
180        meta_opts: Arc<MetaOpts>,
181        hummock_compaction_event_handler: HummockCompactionEventHandler,
182        tx: Option<UnboundedSender<(HummockContextId, RequestEvent)>>,
183    ) -> Self {
184        Self {
185            meta_opts,
186            hummock_compaction_event_handler,
187            tx,
188        }
189    }
190}
191
192#[derive(Clone)]
193pub struct HummockCompactionEventHandler {
194    pub hummock_manager: Arc<HummockManager>,
195}
196
197impl HummockCompactionEventHandler {
198    pub fn new(hummock_manager: Arc<HummockManager>) -> Self {
199        Self { hummock_manager }
200    }
201
202    async fn handle_heartbeat(
203        &self,
204        context_id: HummockContextId,
205        progress: Vec<CompactTaskProgress>,
206    ) -> bool {
207        let mut compactor_alive = true;
208        let compactor_manager = self.hummock_manager.compactor_manager.clone();
209        let cancel_tasks = compactor_manager
210            .update_task_heartbeats(&progress)
211            .into_iter()
212            .map(|task| task.task_id)
213            .collect::<Vec<_>>();
214        if !cancel_tasks.is_empty() {
215            tracing::info!(
216                ?cancel_tasks,
217                %context_id,
218                "Tasks cancel has expired due to lack of visible progress",
219            );
220
221            if let Err(e) = self
222                .hummock_manager
223                .cancel_compact_tasks(cancel_tasks.clone(), TaskStatus::HeartbeatProgressCanceled)
224                .await
225            {
226                tracing::error!(
227                    error = %e.as_report(),
228                    "Attempt to remove compaction task due to elapsed heartbeat failed. We will continue to track its heartbeat
229                    until we can successfully report its status."
230                );
231            }
232        }
233
234        match compactor_manager.get_compactor(context_id) {
235            Some(compactor) => {
236                // Forcefully cancel the task so that it terminates
237                // early on the compactor
238                // node.
239                if !cancel_tasks.is_empty() {
240                    let _ = compactor.cancel_tasks(&cancel_tasks);
241                    tracing::info!(
242                        ?cancel_tasks,
243                        %context_id,
244                        "CancelTask operation has been sent to compactor node",
245                    );
246                }
247            }
248            _ => {
249                // Determine the validity of the compactor streaming rpc. When the compactor no longer exists in the manager, the stream will be removed.
250                // Tip: Connectivity to the compactor will be determined through the `send_event` operation. When send fails, it will be removed from the manager
251                compactor_alive = false;
252            }
253        }
254
255        compactor_alive
256    }
257
258    async fn handle_pull_task_event(
259        &self,
260        context_id: HummockContextId,
261        pull_task_count: usize,
262        compaction_selectors: &mut HashMap<TaskType, Box<dyn CompactionSelector>>,
263        max_get_task_probe_times: usize,
264    ) -> bool {
265        assert_ne!(0, pull_task_count);
266        let Some(compactor) = self
267            .hummock_manager
268            .compactor_manager
269            .get_compactor(context_id)
270        else {
271            return false;
272        };
273
274        // Task selection + dispatch (any early return inside is safe — PullTaskAck below is unreachable only if we return here)
275        self.try_dispatch_tasks(
276            &compactor,
277            pull_task_count,
278            compaction_selectors,
279            max_get_task_probe_times,
280        )
281        .await;
282
283        // PullTaskAck: structurally guaranteed to execute after try_dispatch_tasks returns
284        if let Err(e) = compactor.send_event(ResponseEvent::PullTaskAck(PullTaskAck {})) {
285            tracing::warn!(
286                error = %e.as_report(),
287                "Failed to send ack to {}",
288                context_id,
289            );
290            return false;
291        }
292
293        true
294    }
295
296    /// Selects and dispatches compaction tasks to a compactor.
297    ///
298    /// Separated from `handle_pull_task_event` so that `PullTaskAck` cannot be
299    /// accidentally skipped by early returns.
300    async fn try_dispatch_tasks(
301        &self,
302        compactor: &Arc<Compactor>,
303        pull_task_count: usize,
304        compaction_selectors: &mut HashMap<TaskType, Box<dyn CompactionSelector>>,
305        max_get_task_probe_times: usize,
306    ) {
307        let snapshot = self.hummock_manager.compaction_state.snapshot();
308        let Some((groups, task_type)) = snapshot.pick_compaction_groups_and_type() else {
309            return;
310        };
311
312        if let TaskType::Ttl = task_type {
313            match self
314                .hummock_manager
315                .metadata_manager
316                .get_all_table_options()
317                .await
318                .map_err(|err| Error::MetaStore(err.into()))
319            {
320                Ok(table_options) => {
321                    self.hummock_manager
322                        .update_table_id_to_table_option(table_options);
323                }
324                Err(err) => {
325                    warn!(error = %err.as_report(), "Failed to get table options");
326                }
327            }
328        }
329
330        let selector: &mut Box<dyn CompactionSelector> =
331            compaction_selectors.get_mut(&task_type).unwrap();
332
333        let mut generated_task_count = 0;
334        let mut existed_groups = groups.clone();
335        let mut no_task_groups: HashSet<CompactionGroupId> = HashSet::default();
336        let mut failed_tasks = vec![];
337        let mut loop_times = 0;
338
339        while generated_task_count < pull_task_count
340            && failed_tasks.is_empty()
341            && loop_times < max_get_task_probe_times
342        {
343            loop_times += 1;
344            let compact_ret = self
345                .hummock_manager
346                .get_compact_tasks(
347                    existed_groups.clone(),
348                    pull_task_count - generated_task_count,
349                    &mut **selector,
350                )
351                .await;
352
353            match compact_ret {
354                Ok((compact_tasks, unschedule_groups)) => {
355                    no_task_groups.extend(unschedule_groups);
356                    if compact_tasks.is_empty() {
357                        break;
358                    }
359                    generated_task_count += compact_tasks.len();
360                    for task in compact_tasks {
361                        let task_id = task.task_id;
362                        if let Err(e) =
363                            compactor.send_event(ResponseEvent::CompactTask(task.into()))
364                        {
365                            tracing::warn!(
366                                error = %e.as_report(),
367                                "Failed to send task {} to {}",
368                                task_id,
369                                compactor.context_id(),
370                            );
371                            failed_tasks.push(task_id);
372                        }
373                    }
374                    if !failed_tasks.is_empty() {
375                        self.hummock_manager
376                            .compactor_manager
377                            .remove_compactor(compactor.context_id());
378                    }
379                    existed_groups.retain(|group_id| !no_task_groups.contains(group_id));
380                }
381                Err(err) => {
382                    tracing::warn!(error = %err.as_report(), "Failed to get compaction task");
383                    break;
384                }
385            };
386        }
387        #[cfg(test)]
388        scheduling_tests::before_unschedule().await;
389        for group in no_task_groups {
390            self.hummock_manager.compaction_state.unschedule(
391                group,
392                task_type,
393                snapshot.generation(),
394            );
395        }
396        if let Err(err) = self
397            .hummock_manager
398            .cancel_compact_tasks(failed_tasks, TaskStatus::SendFailCanceled)
399            .await
400        {
401            tracing::warn!(error = %err.as_report(), "Failed to cancel compaction task");
402        }
403    }
404
405    async fn handle_report_task_event(&self, report_events: Vec<ReportTask>) -> Result<()> {
406        if let Err(e) = self
407            .hummock_manager
408            .report_compact_tasks(report_events)
409            .await
410        {
411            tracing::error!(error = %e.as_report(), "report compact_tack fail")
412        }
413        Ok(())
414    }
415}
416
417impl<D: CompactionEventDispatcher<EventType = E::EventType>, E: CompactorStreamEvent>
418    CompactionEventLoop<D, E>
419{
420    pub fn new(
421        hummock_compactor_dispatcher: D,
422        metrics: Arc<MetaMetrics>,
423        compactor_streams_change_rx: UnboundedReceiver<(HummockContextId, Streaming<E>)>,
424    ) -> Self {
425        Self {
426            hummock_compactor_dispatcher,
427            metrics,
428            compactor_streams_change_rx,
429        }
430    }
431
432    pub fn run(mut self) -> (JoinHandle<()>, Sender<()>) {
433        let mut compactor_request_streams = FuturesUnordered::new();
434        let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
435        let shutdown_rx_shared = shutdown_rx.shared();
436
437        let join_handle = tokio::spawn(async move {
438            // A compactor reconnect reuses the same `context_id`, so an older stream may still
439            // resolve after a newer stream has already been registered. Track a local generation
440            // per context to fence off stale stream results from the previous session.
441            let mut stream_generations = HashMap::<HummockContextId, u64>::new();
442            let push_stream =
443                |context_id: HummockContextId,
444                 stream_generation: u64,
445                 stream: Streaming<E>,
446                 compactor_request_streams: &mut FuturesUnordered<_>| {
447                    let future = StreamExt::into_future(stream)
448                        .map(move |stream_future| (context_id, stream_generation, stream_future));
449
450                    compactor_request_streams.push(future);
451                };
452
453            let mut event_loop_iteration_now = Instant::now();
454
455            loop {
456                let shutdown_rx_shared = shutdown_rx_shared.clone();
457                self.metrics
458                    .compaction_event_loop_iteration_latency
459                    .observe(event_loop_iteration_now.elapsed().as_millis() as _);
460                event_loop_iteration_now = Instant::now();
461
462                tokio::select! {
463                    _ = shutdown_rx_shared => { return; },
464
465                    compactor_stream = self.compactor_streams_change_rx.recv() => {
466                        if let Some((context_id, stream)) = compactor_stream {
467                            tracing::info!("compactor {} enters the cluster", context_id);
468                            let stream_generation = stream_generations
469                                .entry(context_id)
470                                .and_modify(|generation| *generation += 1)
471                                .or_insert(1);
472                            push_stream(
473                                context_id,
474                                *stream_generation,
475                                stream,
476                                &mut compactor_request_streams,
477                            );
478                        }
479                    },
480
481                    result = pending_on_none(compactor_request_streams.next()) => {
482                        let (context_id, stream_generation, compactor_stream_req): (_, _, (std::option::Option<std::result::Result<E, _>>, _)) = result;
483                        let Some(current_generation) = stream_generations.get(&context_id).copied() else {
484                            continue;
485                        };
486                        if current_generation != stream_generation {
487                            // Ignore events, EOF and poll errors from superseded streams. Without
488                            // this fence, a late error from an old stream could remove the current
489                            // compactor session for the same `context_id`.
490                            continue;
491                        }
492                        let (event, create_at, stream) = match compactor_stream_req {
493                            (Some(Ok(req)), stream) => {
494                                let create_at = req.create_at();
495                                let event  = req.take_event();
496                                (event, create_at, stream)
497                            }
498
499                            (Some(Err(err)), _stream) => {
500                                tracing::warn!(error = %err.as_report(), %context_id, "compactor stream poll with err, recv stream may be destroyed");
501                                self.hummock_compactor_dispatcher.remove_compactor(context_id);
502                                continue
503                            }
504
505                            _ => {
506                                // remove compactor from compactor manager
507                                tracing::warn!(%context_id, "compactor stream poll err, recv stream may be destroyed");
508                                self.hummock_compactor_dispatcher.remove_compactor(context_id);
509                                continue
510                            },
511                        };
512
513                        {
514                            let consumed_latency_ms = SystemTime::now()
515                                .duration_since(std::time::UNIX_EPOCH)
516                                .expect("Clock may have gone backwards")
517                                .as_millis()
518                                as u64
519                            - create_at;
520                            self.metrics
521                                .compaction_event_consumed_latency
522                                .observe(consumed_latency_ms as _);
523                        }
524
525                        let mut compactor_alive = true;
526                        if self
527                            .hummock_compactor_dispatcher
528                            .should_forward(&event)
529                        {
530                            if let Err(e) = self
531                                .hummock_compactor_dispatcher
532                                .on_event_remotely(context_id, event)
533                                .await
534                            {
535                                tracing::warn!(error = %e.as_report(), "Failed to forward event");
536                            }
537                        } else {
538                            compactor_alive = self.hummock_compactor_dispatcher.on_event_locally(
539                                context_id,
540                                event,
541                            ).await;
542                        }
543
544                        if compactor_alive {
545                            push_stream(
546                                context_id,
547                                stream_generation,
548                                stream,
549                                &mut compactor_request_streams,
550                            );
551                        } else {
552                            tracing::warn!(%context_id, "compactor stream error, send stream may be destroyed");
553                            self
554                            .hummock_compactor_dispatcher
555                            .remove_compactor(context_id);
556                        }
557                    },
558                }
559            }
560        });
561
562        (join_handle, shutdown_tx)
563    }
564}
565
566impl CompactorStreamEvent for SubscribeCompactionEventRequest {
567    type EventType = RequestEvent;
568
569    fn take_event(self) -> Self::EventType {
570        self.event.unwrap()
571    }
572
573    fn create_at(&self) -> u64 {
574        self.create_at
575    }
576}
577
578pub struct HummockCompactorDedicatedEventLoop {
579    hummock_manager: Arc<HummockManager>,
580    hummock_compaction_event_handler: HummockCompactionEventHandler,
581}
582
583impl HummockCompactorDedicatedEventLoop {
584    pub fn new(
585        hummock_manager: Arc<HummockManager>,
586        hummock_compaction_event_handler: HummockCompactionEventHandler,
587    ) -> Self {
588        Self {
589            hummock_manager,
590            hummock_compaction_event_handler,
591        }
592    }
593
594    /// dedicated event runtime for CPU/IO bound event
595    async fn compact_task_dedicated_event_handler(
596        &self,
597        mut rx: UnboundedReceiver<(HummockContextId, RequestEvent)>,
598        shutdown_rx: OneShotReceiver<()>,
599    ) {
600        let mut compaction_selectors = init_selectors();
601
602        tokio::select! {
603            _ = shutdown_rx => {}
604
605            _ = async {
606                while let Some((context_id, event)) = rx.recv().await {
607                    let mut report_events = vec![];
608                    let mut skip_times = 0;
609                    match event {
610                        RequestEvent::PullTask(PullTask { pull_task_count }) => {
611                            self.hummock_compaction_event_handler.handle_pull_task_event(context_id, pull_task_count as usize, &mut compaction_selectors, self.hummock_manager.env.opts.max_get_task_probe_times).await;
612                        }
613
614                        RequestEvent::ReportTask(task) => {
615                           report_events.push(task.into());
616                        }
617
618                        _ => unreachable!(),
619                    }
620                    while let Ok((context_id, event)) = rx.try_recv() {
621                        match event {
622                            RequestEvent::PullTask(PullTask { pull_task_count }) => {
623                                self.hummock_compaction_event_handler.handle_pull_task_event(context_id, pull_task_count as usize, &mut compaction_selectors, self.hummock_manager.env.opts.max_get_task_probe_times).await;
624                                if !report_events.is_empty() {
625                                    if skip_times > MAX_SKIP_TIMES {
626                                        break;
627                                    }
628                                    skip_times += 1;
629                                }
630                            }
631
632                            RequestEvent::ReportTask(task) => {
633                                report_events.push(task.into());
634                                if report_events.len() >= MAX_REPORT_COUNT {
635                                    break;
636                                }
637                            }
638                        _ => unreachable!(),
639                        }
640                    }
641                    if !report_events.is_empty()
642                        && let Err(e) = self.hummock_compaction_event_handler.handle_report_task_event(report_events).await
643                        {
644                            tracing::error!(error = %e.as_report(), "report compact_tack fail")
645                        }
646                }
647            } => {}
648        }
649    }
650
651    pub fn run(
652        self,
653    ) -> (
654        JoinHandle<()>,
655        UnboundedSender<(HummockContextId, RequestEvent)>,
656        Sender<()>,
657    ) {
658        let (tx, rx) = unbounded_channel();
659        let (shutdon_tx, shutdown_rx) = tokio::sync::oneshot::channel();
660        let join_handler = tokio::spawn(async move {
661            self.compact_task_dedicated_event_handler(rx, shutdown_rx)
662                .await;
663        });
664        (join_handler, tx, shutdon_tx)
665    }
666}
667
668pub struct IcebergCompactionEventHandler {
669    compaction_manager: IcebergCompactionManagerRef,
670}
671
672impl IcebergCompactionEventHandler {
673    pub fn new(compaction_manager: IcebergCompactionManagerRef) -> Self {
674        Self { compaction_manager }
675    }
676
677    async fn handle_pull_task_event(
678        &self,
679        context_id: HummockContextId,
680        pull_task_count: usize,
681    ) -> bool {
682        assert_ne!(0, pull_task_count);
683        if let Some(compactor) = self
684            .compaction_manager
685            .iceberg_compactor_manager
686            .get_compactor(context_id)
687        {
688            let mut compactor_alive = true;
689
690            let iceberg_compaction_handles = self
691                .compaction_manager
692                .get_top_n_iceberg_commit_sink_ids(pull_task_count);
693
694            for handle in iceberg_compaction_handles {
695                let compactor = compactor.clone();
696                // send iceberg commit task to compactor
697                if let Err(e) = async move {
698                    handle
699                        .send_compact_task(
700                            compactor,
701                            next_compaction_task_id(&self.compaction_manager.env)
702                                .await?
703                                .into(),
704                        )
705                        .await
706                }
707                .await
708                {
709                    tracing::warn!(
710                        error = %e.as_report(),
711                        "Failed to send iceberg commit task to {}",
712                        context_id,
713                    );
714                    compactor_alive = false;
715                }
716            }
717
718            if let Err(e) =
719                compactor.send_event(IcebergResponseEvent::PullTaskAck(IcebergPullTaskAck {}))
720            {
721                tracing::warn!(
722                    error = %e.as_report(),
723                    "Failed to send ask to {}",
724                    context_id,
725                );
726                compactor_alive = false;
727            }
728
729            return compactor_alive;
730        }
731
732        false
733    }
734
735    fn apply_report_task_event(&self, report: IcebergReportTask) {
736        self.compaction_manager.handle_report_task(report);
737    }
738}
739
740pub struct IcebergCompactionEventDispatcher {
741    compaction_event_handler: IcebergCompactionEventHandler,
742}
743
744#[async_trait::async_trait]
745impl CompactionEventDispatcher for IcebergCompactionEventDispatcher {
746    type EventType = IcebergRequestEvent;
747
748    async fn on_event_locally(&self, context_id: HummockContextId, event: Self::EventType) -> bool {
749        match event {
750            IcebergRequestEvent::PullTask(IcebergPullTask { pull_task_count }) => {
751                return self
752                    .compaction_event_handler
753                    .handle_pull_task_event(context_id, pull_task_count as usize)
754                    .await;
755            }
756            IcebergRequestEvent::ReportTask(report) => {
757                self.compaction_event_handler
758                    .apply_report_task_event(report);
759                return true;
760            }
761            _ => unreachable!(),
762        }
763    }
764
765    async fn on_event_remotely(
766        &self,
767        _context_id: HummockContextId,
768        _event: Self::EventType,
769    ) -> Result<()> {
770        unreachable!()
771    }
772
773    fn should_forward(&self, _event: &Self::EventType) -> bool {
774        false
775    }
776
777    fn remove_compactor(&self, context_id: HummockContextId) {
778        self.compaction_event_handler
779            .compaction_manager
780            .iceberg_compactor_manager
781            .remove_compactor(context_id);
782    }
783}
784
785impl IcebergCompactionEventDispatcher {
786    pub fn new(compaction_event_handler: IcebergCompactionEventHandler) -> Self {
787        Self {
788            compaction_event_handler,
789        }
790    }
791}
792
793impl CompactorStreamEvent for SubscribeIcebergCompactionEventRequest {
794    type EventType = IcebergRequestEvent;
795
796    fn take_event(self) -> Self::EventType {
797        self.event.unwrap()
798    }
799
800    fn create_at(&self) -> u64 {
801        self.create_at
802    }
803}
804
805#[cfg(test)]
806mod scheduling_tests {
807    use std::time::Duration;
808
809    use risingwave_hummock_sdk::{LocalSstableInfo, SyncResult};
810    use risingwave_rpc_client::HummockMetaClient;
811
812    use super::*;
813    use crate::hummock::MockHummockMetaClient;
814    use crate::hummock::compaction::selector::default_compaction_selector;
815    use crate::hummock::manager::compaction::ScheduleTrigger;
816    use crate::hummock::manager::tests::{gen_sstable_info, setup_compute_env_with_meta_opts};
817
818    tokio::task_local! {
819        static BEFORE_UNSCHEDULE: Arc<tokio::sync::Barrier>;
820    }
821
822    pub(super) async fn before_unschedule() {
823        // Only the explicitly scoped dispatch pauses; parallel tests are unaffected.
824        if let Ok(barrier) = BEFORE_UNSCHEDULE.try_with(Arc::clone) {
825            barrier.wait().await;
826            barrier.wait().await;
827        }
828    }
829
830    // Exercise the actual candidate -> picker -> compactor stream path, without a running cluster.
831    async fn dispatch(manager: Arc<HummockManager>) -> Vec<risingwave_pb::hummock::CompactTask> {
832        let mut receiver = manager.compactor_manager.add_compactor(99.into());
833        let compactor = manager.compactor_manager.get_compactor(99.into()).unwrap();
834        let mut selectors = HashMap::from([(TaskType::Dynamic, default_compaction_selector())]);
835        HummockCompactionEventHandler::new(manager)
836            .try_dispatch_tasks(&compactor, 8, &mut selectors, 1)
837            .await;
838        let mut tasks = vec![];
839        while let Ok(event) = receiver.try_recv() {
840            if let Some(ResponseEvent::CompactTask(task)) = event.unwrap().event {
841                tasks.push(task);
842            }
843        }
844        tasks
845    }
846
847    async fn fixture(deterministic: bool) -> (Arc<HummockManager>, MockHummockMetaClient) {
848        let mut opts = MetaOpts::test(false);
849        opts.compaction_deterministic_test = deterministic;
850        let (_, manager, _, worker_id) = setup_compute_env_with_meta_opts(80, opts).await;
851        manager
852            .register_table_ids_for_test(&[(100, 2.into()), (101, 2.into())])
853            .await
854            .unwrap();
855        let client = MockHummockMetaClient::new(manager.clone(), worker_id as _);
856        commit_ssts(&client, 30, 1).await;
857        (manager, client)
858    }
859
860    async fn commit_ssts(client: &MockHummockMetaClient, epoch: u64, first_sst: u64) {
861        client
862            .commit_epoch(
863                risingwave_common::util::epoch::test_epoch(epoch),
864                SyncResult {
865                    uncommitted_ssts: (first_sst..first_sst + 4)
866                        .map(|id| LocalSstableInfo {
867                            sst_info: gen_sstable_info(
868                                id,
869                                vec![100, 101],
870                                risingwave_common::util::epoch::test_epoch(epoch - 10),
871                            ),
872                            table_stats: Default::default(),
873                            created_at: u64::MAX,
874                        })
875                        .collect(),
876                    ..Default::default()
877                },
878            )
879            .await
880            .unwrap();
881    }
882
883    #[tokio::test]
884    async fn test_commit_after_no_task_preserves_candidate() {
885        let (manager, client) = fixture(false).await;
886        // Keep all existing inputs assigned so the next real picker finds no work.
887        assert!(!dispatch(manager.clone()).await.is_empty());
888        let barrier = Arc::new(tokio::sync::Barrier::new(2));
889        let commit = async {
890            barrier.wait().await;
891            commit_ssts(&client, 40, 5).await;
892            barrier.wait().await;
893        };
894        let (old_tasks, ()) = tokio::time::timeout(Duration::from_secs(10), async {
895            tokio::join!(
896                BEFORE_UNSCHEDULE.scope(barrier.clone(), dispatch(manager.clone())),
897                commit,
898            )
899        })
900        .await
901        .expect("picker and commit must finish without holding each other's locks");
902        assert!(
903            old_tasks.is_empty(),
904            "the paused picker must have found no task"
905        );
906        assert!(
907            manager
908                .compaction_state
909                .snapshot()
910                .scheduled
911                .contains(&(2.into(), TaskType::Dynamic))
912        );
913        let tasks = dispatch(manager.clone()).await;
914        assert!(
915            tasks.iter().any(|task| task
916                .input_ssts
917                .iter()
918                .any(|level| { level.table_infos.iter().any(|sst| sst.sst_id >= 5.into()) })),
919            "the next dispatch must select the newly committed SSTs"
920        );
921        // Once the new inputs are assigned, a fresh no-task result must still cool down the group.
922        assert!(dispatch(manager.clone()).await.is_empty());
923        assert!(!manager.compaction_state.try_sched_compaction(
924            2.into(),
925            TaskType::Dynamic,
926            ScheduleTrigger::Periodic,
927        ));
928    }
929
930    #[tokio::test]
931    async fn test_split_schedules_children_without_commit_or_timer() {
932        let (manager, _) = fixture(false).await;
933        let snapshot = manager.compaction_state.snapshot();
934        manager
935            .compaction_state
936            .unschedule(2.into(), TaskType::Dynamic, snapshot.generation());
937        manager
938            .move_state_tables_to_dedicated_compaction_group(2.into(), &[100.into()], None)
939            .await
940            .unwrap();
941        let tasks = dispatch(manager.clone()).await;
942        let groups: HashSet<_> = tasks.iter().map(|task| task.compaction_group_id).collect();
943        let version = manager.get_current_version().await;
944        for table in [100, 101] {
945            let group = version.state_table_info.info()
946                [&risingwave_common::catalog::TableId::new(table)]
947                .compaction_group_id;
948            assert!(
949                groups.contains(&group),
950                "split group {group} did not reach the picker"
951            );
952        }
953    }
954
955    #[tokio::test]
956    async fn test_merge_wakes_cooled_survivor_without_commit() {
957        let (manager, _) = fixture(false).await;
958        let (left, mapping) = manager
959            .move_state_tables_to_dedicated_compaction_group(2.into(), &[100.into()], None)
960            .await
961            .unwrap();
962        let right = *mapping.keys().find(|&&id| id != left).unwrap();
963        // Obtain a real no-task result before merge by keeping the existing left SSTs assigned.
964        let pending = manager
965            .get_compact_task(left, &mut *default_compaction_selector())
966            .await
967            .unwrap()
968            .unwrap();
969        dispatch(manager.clone()).await;
970        assert!(
971            manager
972                .compaction_state
973                .inner
974                .lock()
975                .dynamic_cooldown
976                .contains(&left)
977        );
978        assert!(!manager.compaction_state.try_sched_compaction(
979            left,
980            TaskType::Dynamic,
981            ScheduleTrigger::Periodic
982        ));
983        manager
984            .merge_compaction_group_for_test(left, right, HashSet::from([100.into(), 101.into()]))
985            .await
986            .unwrap();
987        let tasks = dispatch(manager.clone()).await;
988        assert!(
989            tasks
990                .iter()
991                .any(|task| task.compaction_group_id == left && task.task_id != pending.task_id),
992            "merged SSTs must reach the picker without another commit"
993        );
994    }
995
996    #[tokio::test]
997    async fn test_merge_cancellation_does_not_reschedule_deleted_group() {
998        for task_type in [TaskType::Dynamic, TaskType::Emergency] {
999            let (manager, _) = fixture(false).await;
1000            let (left, mapping) = manager
1001                .move_state_tables_to_dedicated_compaction_group(2.into(), &[100.into()], None)
1002                .await
1003                .unwrap();
1004            let right = *mapping.keys().find(|&&id| id != left).unwrap();
1005            let task = manager
1006                .get_compact_task(right, &mut *default_compaction_selector())
1007                .await
1008                .unwrap()
1009                .unwrap();
1010            // The report path consumes the assigned task type. Keep a real assignment and inputs.
1011            manager
1012                .compaction
1013                .write()
1014                .await
1015                .compact_task_assignment
1016                .get_mut(&task.task_id)
1017                .unwrap()
1018                .compact_task
1019                .task_type = task_type;
1020            for kind in [
1021                TaskType::Dynamic,
1022                TaskType::Emergency,
1023                TaskType::Ttl,
1024                TaskType::SpaceReclaim,
1025                TaskType::Tombstone,
1026                TaskType::VnodeWatermark,
1027            ] {
1028                manager.compaction_state.try_sched_compaction(
1029                    right,
1030                    kind,
1031                    ScheduleTrigger::NewData,
1032                );
1033            }
1034            let snapshot = manager.compaction_state.snapshot();
1035            manager
1036                .merge_compaction_group_for_test(
1037                    left,
1038                    right,
1039                    HashSet::from([100.into(), 101.into()]),
1040                )
1041                .await
1042                .unwrap();
1043            assert!(
1044                !manager
1045                    .compaction_state
1046                    .snapshot()
1047                    .scheduled
1048                    .iter()
1049                    .any(|(id, _)| *id == right)
1050            );
1051            // An old picker result must not recreate cooldown for the deleted group.
1052            manager
1053                .compaction_state
1054                .unschedule(right, TaskType::Dynamic, snapshot.generation());
1055            assert!(
1056                !manager
1057                    .compaction_state
1058                    .inner
1059                    .lock()
1060                    .dynamic_cooldown
1061                    .contains(&right)
1062            );
1063            // A getter or explicit trigger using the deleted group must also leave it absent.
1064            let (tasks, _) = manager
1065                .get_compact_tasks(vec![right], 1, &mut *default_compaction_selector())
1066                .await
1067                .unwrap();
1068            assert!(tasks.is_empty());
1069            manager
1070                .trigger_compaction_deterministic(
1071                    manager.get_current_version().await.id,
1072                    vec![right],
1073                )
1074                .await
1075                .unwrap();
1076            assert!(
1077                !manager
1078                    .compaction_state
1079                    .snapshot()
1080                    .scheduled
1081                    .iter()
1082                    .any(|(id, _)| *id == right)
1083            );
1084        }
1085    }
1086
1087    #[tokio::test]
1088    async fn test_normalization_schedules_all_affected_groups() {
1089        for deterministic in [false, true] {
1090            let mut opts = MetaOpts::test(false);
1091            opts.compaction_deterministic_test = deterministic;
1092            let (_, manager, _, _) = setup_compute_env_with_meta_opts(80, opts).await;
1093            manager
1094                .register_table_ids_for_test(&[(100, 2.into()), (102, 2.into()), (101, 3.into())])
1095                .await
1096                .unwrap();
1097            let old = manager.compaction_state.snapshot();
1098            assert_eq!(
1099                manager
1100                    .normalize_overlapping_compaction_groups()
1101                    .await
1102                    .unwrap(),
1103                1
1104            );
1105            let version = manager.get_current_version().await;
1106            for table in [100, 102] {
1107                let group = version.state_table_info.info()
1108                    [&risingwave_common::catalog::TableId::new(table)]
1109                    .compaction_group_id;
1110                // An older no-task result cannot erase the topology notification.
1111                manager
1112                    .compaction_state
1113                    .unschedule(group, TaskType::Dynamic, old.generation());
1114                assert_eq!(
1115                    manager
1116                        .compaction_state
1117                        .snapshot()
1118                        .scheduled
1119                        .contains(&(group, TaskType::Dynamic)),
1120                    !deterministic
1121                );
1122            }
1123        }
1124    }
1125
1126    #[tokio::test]
1127    async fn test_topology_changes_respect_deterministic_mode() {
1128        let (manager, _) = fixture(true).await;
1129        let (left, mapping) = manager
1130            .move_state_tables_to_dedicated_compaction_group(2.into(), &[100.into()], None)
1131            .await
1132            .unwrap();
1133        assert!(manager.compaction_state.snapshot().scheduled.is_empty());
1134        let right = *mapping.keys().find(|&&id| id != left).unwrap();
1135        manager
1136            .merge_compaction_group_for_test(left, right, HashSet::from([100.into(), 101.into()]))
1137            .await
1138            .unwrap();
1139        assert!(manager.compaction_state.snapshot().scheduled.is_empty());
1140    }
1141}