Skip to main content

risingwave_storage/hummock/compactor/iceberg_compaction/
mod.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
15pub use self::iceberg_compactor_runner::create_task_execution;
16#[cfg(madsim)]
17pub use self::report::set_simulated_pk_index_compaction_result;
18#[cfg(madsim)]
19pub(crate) use self::report::simulated_pk_index_compaction_result;
20pub(crate) use self::report::{
21    IcebergPlanCompletion, IcebergTaskReport, IcebergTaskTracker, PkIndexCompactionResult,
22    ReportSendResult, build_drained_iceberg_task_report, build_iceberg_task_report,
23    build_pk_index_compaction_result, flush_pending_iceberg_task_reports,
24    send_or_buffer_iceberg_task_report,
25};
26use crate::hummock::compactor::iceberg_compaction::iceberg_compactor_runner::IcebergCompactionPlanRunner;
27
28pub(crate) mod iceberg_compactor_runner;
29pub(crate) mod memory;
30pub(crate) mod report;
31
32use std::collections::{HashMap, VecDeque};
33use std::sync::Arc;
34
35use risingwave_pb::id::IcebergCompactionTaskId;
36use tokio::sync::Notify;
37
38use crate::monitor::CompactorMetrics;
39
40/// Unique key combining `(task_id, plan_index)` since one task can have multiple plans.
41pub(crate) type TaskKey = (IcebergCompactionTaskId, usize);
42
43/// Task metadata for queue operations.
44#[derive(Debug, Clone)]
45pub struct IcebergTaskMeta {
46    pub task_id: IcebergCompactionTaskId,
47    pub plan_index: usize,
48    /// Must be in range `1..=max_parallelism`
49    pub required_parallelism: u32,
50    /// Estimated heap peak used for admission control.
51    /// Must be in range `1..=total_memory_budget_bytes`.
52    pub memory_reservation_bytes: usize,
53}
54
55#[derive(Debug)]
56pub struct PoppedIcebergTask {
57    pub meta: IcebergTaskMeta,
58    pub runner: Option<IcebergCompactionPlanRunner>,
59}
60
61impl IcebergTaskMeta {
62    fn key(&self) -> TaskKey {
63        (self.task_id, self.plan_index)
64    }
65}
66
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
68struct IcebergTaskResources {
69    required_parallelism: u32,
70    memory_reservation_bytes: usize,
71}
72
73/// Internal storage for the task queue.
74struct IcebergTaskQueueInner {
75    /// FIFO queue of waiting task metadata
76    deque: VecDeque<IcebergTaskMeta>,
77    /// Maps `(task_id, plan_index)` to its reserved resources for tracking
78    resource_map: HashMap<TaskKey, IcebergTaskResources>,
79    /// Sum of `required_parallelism` for all waiting tasks
80    waiting_parallelism_sum: u32,
81    /// Sum of `required_parallelism` for all running tasks
82    running_parallelism_sum: u32,
83    /// Sum of estimated heap peak reservations for all running tasks
84    running_memory_reservation_bytes: usize,
85    /// Optional runner payloads indexed by `(task_id, plan_index)`
86    runners: HashMap<TaskKey, IcebergCompactionPlanRunner>,
87}
88
89/// FIFO task queue with parallelism- and memory-based scheduling for Iceberg compaction.
90///
91/// Tasks execute in submission order when sufficient parallelism is available.
92/// The queue tracks waiting and running tasks to prevent over-commitment of running resources.
93///
94/// Constraints:
95/// - Each task requires `1..=max_parallelism` units
96/// - Each task reserves `1..=total_memory_budget_bytes` bytes while running
97/// - Total waiting parallelism cannot exceed `pending_parallelism_budget`
98/// - Total running parallelism cannot exceed `max_parallelism`
99/// - Total running memory cannot exceed `total_memory_budget_bytes`
100/// - Tasks block until enough parallelism and memory are available
101///
102/// Note: The queue rejects duplicate task keys but does not reorder or merge tasks.
103/// Higher-level task management remains Meta's responsibility.
104pub struct IcebergTaskQueue {
105    inner: IcebergTaskQueueInner,
106    /// Maximum concurrent parallelism for running tasks
107    max_parallelism: u32,
108    /// Maximum total parallelism for waiting tasks (backpressure limit)
109    pending_parallelism_budget: u32,
110    /// Maximum total memory reserved by running tasks
111    total_memory_budget_bytes: usize,
112    /// Notification for event-driven scheduling
113    schedule_notify: Arc<Notify>,
114    metrics: Arc<CompactorMetrics>,
115}
116
117#[derive(Debug, PartialEq, Eq)]
118pub enum PushResult {
119    Added,
120    /// Would exceed `pending_parallelism_budget`
121    RejectedCapacity,
122    /// The task's parallelism or estimated memory exceeds the worker capacity.
123    RejectedTooLarge,
124    /// `required_parallelism` == 0
125    RejectedInvalidParallelism,
126    /// Task with same `(task_id, plan_index)` already exists
127    RejectedDuplicate,
128}
129
130impl IcebergTaskQueue {
131    #[cfg(test)]
132    pub fn new(
133        max_parallelism: u32,
134        pending_parallelism_budget: u32,
135        total_memory_budget_bytes: usize,
136    ) -> Self {
137        Self::new_with_metrics(
138            max_parallelism,
139            pending_parallelism_budget,
140            total_memory_budget_bytes,
141            Arc::new(CompactorMetrics::unused()),
142        )
143    }
144
145    pub(crate) fn new_with_metrics(
146        max_parallelism: u32,
147        pending_parallelism_budget: u32,
148        total_memory_budget_bytes: usize,
149        metrics: Arc<CompactorMetrics>,
150    ) -> Self {
151        assert!(max_parallelism > 0, "max_parallelism must be > 0");
152        assert!(
153            pending_parallelism_budget >= max_parallelism,
154            "pending budget should allow at least one task"
155        );
156        assert!(
157            total_memory_budget_bytes > 0,
158            "total memory budget must be > 0"
159        );
160        metrics
161            .iceberg_compaction_memory_budget_bytes
162            .set(i64::try_from(total_memory_budget_bytes).unwrap_or(i64::MAX));
163        metrics
164            .iceberg_compaction_running_memory_reservation_bytes
165            .set(0);
166        Self {
167            inner: IcebergTaskQueueInner {
168                deque: VecDeque::new(),
169                resource_map: HashMap::new(),
170                waiting_parallelism_sum: 0,
171                running_parallelism_sum: 0,
172                running_memory_reservation_bytes: 0,
173                runners: HashMap::new(),
174            },
175            max_parallelism,
176            pending_parallelism_budget,
177            total_memory_budget_bytes,
178            schedule_notify: Arc::new(Notify::new()),
179            metrics,
180        }
181    }
182
183    /// Waits until there are tasks that can be scheduled.
184    ///
185    /// Returns `true` if there are schedulable tasks, `false` otherwise.
186    /// Use this in a `tokio::select!` to wake up when tasks become schedulable.
187    pub async fn wait_schedulable(&self) -> bool {
188        // Check if we have tasks that can be scheduled right now
189        if self.has_schedulable_tasks() {
190            return true;
191        }
192        // Otherwise wait for notification
193        self.schedule_notify.notified().await;
194        self.has_schedulable_tasks()
195    }
196
197    fn has_schedulable_tasks(&self) -> bool {
198        if let Some(front_task) = self.inner.deque.front() {
199            let available_parallelism = self
200                .max_parallelism
201                .saturating_sub(self.inner.running_parallelism_sum);
202            let available_memory = self
203                .total_memory_budget_bytes
204                .saturating_sub(self.inner.running_memory_reservation_bytes);
205            available_parallelism >= front_task.required_parallelism
206                && available_memory >= front_task.memory_reservation_bytes
207        } else {
208            false
209        }
210    }
211
212    fn notify_schedulable(&self) {
213        if self.has_schedulable_tasks() {
214            self.schedule_notify.notify_one();
215        }
216    }
217
218    pub fn running_parallelism_sum(&self) -> u32 {
219        self.inner.running_parallelism_sum
220    }
221
222    pub fn waiting_parallelism_sum(&self) -> u32 {
223        self.inner.waiting_parallelism_sum
224    }
225
226    fn available_parallelism(&self) -> u32 {
227        self.max_parallelism
228            .saturating_sub(self.inner.running_parallelism_sum)
229    }
230
231    fn available_memory_reservation_bytes(&self) -> usize {
232        self.total_memory_budget_bytes
233            .saturating_sub(self.inner.running_memory_reservation_bytes)
234    }
235
236    fn update_running_memory_metric(&self) {
237        self.metrics
238            .iceberg_compaction_running_memory_reservation_bytes
239            .set(i64::try_from(self.inner.running_memory_reservation_bytes).unwrap_or(i64::MAX));
240    }
241
242    /// Push a task into the queue.
243    ///
244    /// The task is validated and added to the end of the FIFO queue if constraints are met.
245    pub fn push(
246        &mut self,
247        meta: IcebergTaskMeta,
248        runner: Option<IcebergCompactionPlanRunner>,
249    ) -> PushResult {
250        if meta.required_parallelism == 0 {
251            return PushResult::RejectedInvalidParallelism;
252        }
253        if meta.required_parallelism > self.max_parallelism
254            || meta.memory_reservation_bytes > self.total_memory_budget_bytes
255        {
256            return PushResult::RejectedTooLarge;
257        }
258        let key = meta.key();
259
260        // Reject duplicate keys to prevent inconsistent state between resource_map and deque.
261        if self.inner.resource_map.contains_key(&key) {
262            return PushResult::RejectedDuplicate;
263        }
264
265        let Some(new_parallelism_total) = self
266            .inner
267            .waiting_parallelism_sum
268            .checked_add(meta.required_parallelism)
269        else {
270            return PushResult::RejectedCapacity;
271        };
272        if new_parallelism_total > self.pending_parallelism_budget {
273            return PushResult::RejectedCapacity;
274        }
275        self.inner.resource_map.insert(
276            key,
277            IcebergTaskResources {
278                required_parallelism: meta.required_parallelism,
279                memory_reservation_bytes: meta.memory_reservation_bytes,
280            },
281        );
282        self.inner.waiting_parallelism_sum = new_parallelism_total;
283
284        if let Some(r) = runner {
285            self.inner.runners.insert(key, r);
286        }
287
288        self.inner.deque.push_back(meta);
289
290        self.notify_schedulable();
291        PushResult::Added
292    }
293
294    /// Pop the next task if sufficient parallelism and memory are available.
295    ///
296    /// Returns `None` if the queue is empty or the front task cannot fit
297    /// within the available running resource budgets.
298    pub fn pop(&mut self) -> Option<PoppedIcebergTask> {
299        let front = self.inner.deque.front()?;
300        if front.required_parallelism > self.available_parallelism() {
301            return None;
302        }
303        if front.memory_reservation_bytes > self.available_memory_reservation_bytes() {
304            return None;
305        }
306
307        let meta = self.inner.deque.pop_front()?;
308        self.inner.waiting_parallelism_sum = self
309            .inner
310            .waiting_parallelism_sum
311            .saturating_sub(meta.required_parallelism);
312        self.inner.running_parallelism_sum = self
313            .inner
314            .running_parallelism_sum
315            .saturating_add(meta.required_parallelism);
316        self.inner.running_memory_reservation_bytes = self
317            .inner
318            .running_memory_reservation_bytes
319            .checked_add(meta.memory_reservation_bytes)
320            .expect("running memory reservation sum overflowed");
321
322        let key = meta.key();
323        self.update_running_memory_metric();
324        let runner = self.inner.runners.remove(&key);
325        Some(PoppedIcebergTask { meta, runner })
326    }
327
328    /// Mark a task as finished, freeing its parallelism and memory for other tasks.
329    ///
330    /// Returns `true` if the task was found and removed, `false` otherwise.
331    pub fn finish_running(&mut self, task_key: TaskKey) -> bool {
332        let Some(resources) = self.inner.resource_map.remove(&task_key) else {
333            tracing::warn!(
334                task_id = %task_key.0,
335                plan_index = task_key.1,
336                "finish_running called for unknown task key, possible bug: double-finish or invalid key"
337            );
338            return false;
339        };
340        self.inner.running_parallelism_sum = self
341            .inner
342            .running_parallelism_sum
343            .saturating_sub(resources.required_parallelism);
344        self.inner.running_memory_reservation_bytes = self
345            .inner
346            .running_memory_reservation_bytes
347            .checked_sub(resources.memory_reservation_bytes)
348            .expect("running memory reservation bookkeeping underflowed");
349        self.inner.runners.remove(&task_key);
350        self.update_running_memory_metric();
351        self.notify_schedulable();
352        true
353    }
354
355    /// Cancel all waiting plans belonging to the given task.
356    ///
357    /// Returns the number of waiting plans removed from the queue.
358    pub fn cancel_waiting_task(&mut self, task_id: IcebergCompactionTaskId) -> usize {
359        let mut retained = VecDeque::with_capacity(self.inner.deque.len());
360        let mut cancelled_parallelism = 0;
361        let mut cancelled_count = 0;
362
363        while let Some(meta) = self.inner.deque.pop_front() {
364            if meta.task_id == task_id {
365                cancelled_parallelism += meta.required_parallelism;
366                cancelled_count += 1;
367                self.inner.resource_map.remove(&meta.key());
368                self.inner.runners.remove(&meta.key());
369            } else {
370                retained.push_back(meta);
371            }
372        }
373
374        self.inner.deque = retained;
375        self.inner.waiting_parallelism_sum = self
376            .inner
377            .waiting_parallelism_sum
378            .saturating_sub(cancelled_parallelism);
379
380        if cancelled_count > 0 {
381            self.notify_schedulable();
382        }
383
384        cancelled_count
385    }
386}
387
388#[cfg(test)]
389mod tests {
390    use super::*;
391
392    fn mk_meta(id: u64, plan_index: usize, p: u32) -> IcebergTaskMeta {
393        IcebergTaskMeta {
394            task_id: id.into(),
395            plan_index,
396            required_parallelism: p,
397            memory_reservation_bytes: 1,
398        }
399    }
400
401    #[test]
402    fn test_basic_push_pop() {
403        let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
404        assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
405        assert_eq!(q.waiting_parallelism_sum(), 4);
406
407        let popped = q.pop().expect("should pop");
408        assert_eq!(popped.meta.task_id.as_raw_id(), 1);
409        assert_eq!(q.waiting_parallelism_sum(), 0);
410        assert_eq!(q.running_parallelism_sum(), 4);
411
412        assert!(q.finish_running((1.into(), 0)));
413        assert_eq!(q.running_parallelism_sum(), 0);
414    }
415
416    #[test]
417    fn test_fifo_ordering() {
418        let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
419        assert_eq!(q.push(mk_meta(1, 0, 2), None), PushResult::Added);
420        assert_eq!(q.push(mk_meta(2, 0, 2), None), PushResult::Added);
421        assert_eq!(q.push(mk_meta(3, 0, 2), None), PushResult::Added);
422
423        assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 1);
424        assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 2);
425        assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 3);
426    }
427
428    #[test]
429    fn test_capacity_reject() {
430        let mut q = IcebergTaskQueue::new(4, 6, usize::MAX);
431        assert_eq!(q.push(mk_meta(1, 0, 3), None), PushResult::Added);
432        assert_eq!(q.push(mk_meta(2, 0, 3), None), PushResult::Added); // sum=6
433        assert_eq!(q.push(mk_meta(3, 0, 1), None), PushResult::RejectedCapacity); // would exceed
434    }
435
436    #[test]
437    fn test_invalid_parallelism() {
438        let mut q = IcebergTaskQueue::new(4, 10, usize::MAX);
439        assert_eq!(
440            q.push(mk_meta(1, 0, 0), None),
441            PushResult::RejectedInvalidParallelism
442        );
443        assert_eq!(q.push(mk_meta(2, 0, 5), None), PushResult::RejectedTooLarge); // > max
444    }
445
446    #[test]
447    fn test_memory_estimate_exceeding_budget_is_rejected() {
448        let mut q = IcebergTaskQueue::new(4, 10, 100);
449        let mut meta = mk_meta(1, 0, 1);
450        meta.memory_reservation_bytes = 101;
451
452        assert_eq!(q.push(meta, None), PushResult::RejectedTooLarge);
453    }
454
455    #[test]
456    fn test_duplicate_key_rejected() {
457        let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
458        assert_eq!(q.push(mk_meta(1, 0, 3), None), PushResult::Added);
459        // Same (task_id, plan_index) should be rejected
460        assert_eq!(
461            q.push(mk_meta(1, 0, 5), None),
462            PushResult::RejectedDuplicate
463        );
464        // Parallelism sum should not have changed
465        assert_eq!(q.waiting_parallelism_sum(), 3);
466
467        // Different plan_index is allowed
468        assert_eq!(q.push(mk_meta(1, 1, 2), None), PushResult::Added);
469        assert_eq!(q.waiting_parallelism_sum(), 5);
470
471        // After pop and finish, the key can be reused
472        let p = q.pop().unwrap();
473        assert_eq!(p.meta.task_id.as_raw_id(), 1);
474        assert_eq!(p.meta.plan_index, 0);
475        q.finish_running((1.into(), 0));
476
477        // Now the same key can be pushed again
478        assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
479    }
480
481    #[test]
482    fn test_pop_insufficient_parallelism() {
483        let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
484        assert_eq!(q.push(mk_meta(1, 0, 6), None), PushResult::Added);
485        assert_eq!(q.push(mk_meta(2, 0, 4), None), PushResult::Added);
486
487        let p1 = q.pop().unwrap();
488        assert_eq!(p1.meta.task_id.as_raw_id(), 1);
489        // Not enough remaining parallelism (only 2 left)
490        assert!(q.pop().is_none());
491
492        // Finish first, then second becomes schedulable
493        assert!(q.finish_running((1.into(), 0)));
494        let p2 = q.pop().unwrap();
495        assert_eq!(p2.meta.task_id.as_raw_id(), 2);
496    }
497
498    #[test]
499    fn test_finish_nonexistent_task() {
500        let mut q = IcebergTaskQueue::new(4, 16, usize::MAX);
501        assert!(!q.finish_running((999.into(), 0)));
502        assert_eq!(q.running_parallelism_sum(), 0);
503    }
504
505    #[test]
506    fn test_double_finish() {
507        let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
508        assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
509        q.pop().unwrap();
510        assert_eq!(q.running_parallelism_sum(), 4);
511
512        // First finish succeeds
513        assert!(q.finish_running((1.into(), 0)));
514        assert_eq!(q.running_parallelism_sum(), 0);
515
516        // Second finish on same key returns false (triggers warn log)
517        assert!(!q.finish_running((1.into(), 0)));
518        assert_eq!(q.running_parallelism_sum(), 0);
519    }
520
521    #[test]
522    fn test_max_parallelism_boundary() {
523        // pending_budget == max_parallelism: minimal valid configuration
524        let mut q = IcebergTaskQueue::new(4, 4, usize::MAX);
525
526        // Task with required_parallelism == max_parallelism should be accepted
527        assert_eq!(q.push(mk_meta(1, 0, 4), None), PushResult::Added);
528        // Budget exhausted
529        assert_eq!(q.push(mk_meta(2, 0, 1), None), PushResult::RejectedCapacity);
530
531        // Can pop and run at full parallelism
532        let p = q.pop().unwrap();
533        assert_eq!(p.meta.required_parallelism, 4);
534        assert_eq!(q.running_parallelism_sum(), 4);
535
536        // No room for any new running task
537        assert_eq!(q.push(mk_meta(3, 0, 1), None), PushResult::Added);
538        assert!(q.pop().is_none());
539    }
540
541    #[test]
542    fn test_same_task_id_multiple_plans() {
543        let mut q = IcebergTaskQueue::new(10, 30, usize::MAX);
544        let task_id = 1u64;
545
546        // Same task_id with different plan_index are independent
547        assert_eq!(q.push(mk_meta(task_id, 0, 3), None), PushResult::Added);
548        assert_eq!(q.push(mk_meta(task_id, 1, 4), None), PushResult::Added);
549        assert_eq!(q.push(mk_meta(task_id, 2, 2), None), PushResult::Added);
550        assert_eq!(q.waiting_parallelism_sum(), 9);
551
552        // Pop all
553        for i in 0..3 {
554            let p = q.pop().unwrap();
555            assert_eq!(p.meta.task_id.as_raw_id(), task_id);
556            assert_eq!(p.meta.plan_index, i);
557        }
558        assert_eq!(q.running_parallelism_sum(), 9);
559
560        // Finish out of order
561        assert!(q.finish_running((task_id.into(), 1)));
562        assert_eq!(q.running_parallelism_sum(), 5);
563        assert!(q.finish_running((task_id.into(), 0)));
564        assert!(q.finish_running((task_id.into(), 2)));
565        assert_eq!(q.running_parallelism_sum(), 0);
566    }
567
568    #[test]
569    fn test_cancel_waiting_task_only_removes_waiting_plans() {
570        let mut q = IcebergTaskQueue::new(10, 30, usize::MAX);
571        let task_id = 1u64;
572
573        assert_eq!(q.push(mk_meta(task_id, 0, 3), None), PushResult::Added);
574        assert_eq!(q.push(mk_meta(task_id, 1, 4), None), PushResult::Added);
575        assert_eq!(q.push(mk_meta(2, 0, 2), None), PushResult::Added);
576
577        let popped = q.pop().unwrap();
578        assert_eq!(popped.meta.task_id.as_raw_id(), task_id);
579        assert_eq!(popped.meta.plan_index, 0);
580        assert_eq!(q.running_parallelism_sum(), 3);
581        assert_eq!(q.waiting_parallelism_sum(), 6);
582
583        assert_eq!(q.cancel_waiting_task(task_id.into()), 1);
584        assert_eq!(q.running_parallelism_sum(), 3);
585        assert_eq!(q.waiting_parallelism_sum(), 2);
586
587        assert!(q.finish_running((task_id.into(), 0)));
588        let next = q.pop().unwrap();
589        assert_eq!(next.meta.task_id.as_raw_id(), 2);
590        assert_eq!(next.meta.plan_index, 0);
591    }
592
593    #[test]
594    fn test_empty_queue_behavior() {
595        let mut q = IcebergTaskQueue::new(8, 32, usize::MAX);
596        assert!(q.pop().is_none());
597        assert!(!q.finish_running((1.into(), 0)));
598        assert_eq!(q.waiting_parallelism_sum(), 0);
599        assert_eq!(q.running_parallelism_sum(), 0);
600    }
601    #[test]
602    fn waiting_memory_does_not_count_against_running_budget() {
603        let mut q = IcebergTaskQueue::new(8, 32, 100);
604        let meta = |id, memory_reservation_bytes| IcebergTaskMeta {
605            task_id: IcebergCompactionTaskId::new(id),
606            plan_index: 0,
607            required_parallelism: 1,
608            memory_reservation_bytes,
609        };
610        assert_eq!(q.push(meta(1, 80), None), PushResult::Added);
611        q.pop().unwrap();
612
613        assert_eq!(q.push(meta(2, 60), None), PushResult::Added);
614        assert!(q.pop().is_none());
615
616        assert!(q.finish_running((1.into(), 0)));
617        assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 2);
618    }
619
620    #[test]
621    fn memory_head_of_line_blocking_preserves_fifo() {
622        let mut q = IcebergTaskQueue::new(8, 32, 100);
623        let meta = |id, memory_reservation_bytes| IcebergTaskMeta {
624            task_id: IcebergCompactionTaskId::new(id),
625            plan_index: 0,
626            required_parallelism: 1,
627            memory_reservation_bytes,
628        };
629        assert_eq!(q.push(meta(1, 60), None), PushResult::Added);
630        q.pop().unwrap();
631        assert_eq!(q.push(meta(2, 50), None), PushResult::Added);
632        assert_eq!(q.push(meta(3, 40), None), PushResult::Added);
633
634        assert!(q.pop().is_none());
635        assert!(q.finish_running((1.into(), 0)));
636        assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 2);
637        assert_eq!(q.pop().unwrap().meta.task_id.as_raw_id(), 3);
638    }
639}