1use 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 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 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 self.try_dispatch_tasks(
276 &compactor,
277 pull_task_count,
278 compaction_selectors,
279 max_get_task_probe_times,
280 )
281 .await;
282
283 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 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 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 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 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 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 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 if let Ok(barrier) = BEFORE_UNSCHEDULE.try_with(Arc::clone) {
825 barrier.wait().await;
826 barrier.wait().await;
827 }
828 }
829
830 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 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 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 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 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 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 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 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}