Skip to main content

risingwave_storage/hummock/compactor/iceberg_compaction/
report.rs

1// Copyright 2026 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::VecDeque;
16use std::time::SystemTime;
17
18use risingwave_pb::iceberg_compaction::{
19    SubscribeIcebergCompactionEventRequest, subscribe_iceberg_compaction_event_request,
20};
21use risingwave_pb::id::IcebergCompactionTaskId;
22use thiserror_ext::AsReport;
23use tokio::sync::mpsc;
24
25use super::TaskKey;
26
27#[derive(Debug)]
28pub(crate) struct IcebergPlanCompletion {
29    pub(crate) task_key: TaskKey,
30    pub(crate) error_message: Option<String>,
31}
32
33pub(crate) type IcebergTaskReport = subscribe_iceberg_compaction_event_request::ReportTask;
34
35pub(crate) enum ReportSendResult {
36    Sent,
37    RestartStream,
38}
39
40pub(crate) struct IcebergTaskTracker {
41    sink_id: u32,
42    total_plans: usize,
43    remaining_plans: usize,
44    successful_plans: usize,
45    failed_plans: usize,
46    first_error: Option<String>,
47}
48
49impl IcebergTaskTracker {
50    pub(crate) fn new(sink_id: u32, remaining_plans: usize) -> Self {
51        Self {
52            sink_id,
53            total_plans: remaining_plans,
54            remaining_plans,
55            successful_plans: 0,
56            failed_plans: 0,
57            first_error: None,
58        }
59    }
60
61    pub(crate) fn record_completion(&mut self, error_message: Option<String>) {
62        debug_assert!(self.remaining_plans > 0);
63        self.remaining_plans -= 1;
64        if let Some(error_message) = error_message {
65            self.failed_plans += 1;
66            if self.first_error.is_none() {
67                self.first_error = Some(error_message);
68            }
69        } else {
70            self.successful_plans += 1;
71        }
72    }
73
74    pub(crate) fn is_finished(&self) -> bool {
75        self.remaining_plans == 0
76    }
77
78    pub(crate) fn sink_id(&self) -> u32 {
79        self.sink_id
80    }
81
82    pub(crate) fn total_plans(&self) -> usize {
83        self.total_plans
84    }
85
86    pub(crate) fn successful_plans(&self) -> usize {
87        self.successful_plans
88    }
89
90    pub(crate) fn failed_plans(&self) -> usize {
91        self.failed_plans
92    }
93
94    pub(crate) fn into_report(self, task_id: IcebergCompactionTaskId) -> IcebergTaskReport {
95        let error_message = if self.successful_plans > 0 {
96            None
97        } else {
98            Some(
99                self.first_error
100                    .unwrap_or_else(|| "All admitted iceberg compaction plans failed".to_owned()),
101            )
102        };
103        build_iceberg_task_report(task_id, self.sink_id, error_message)
104    }
105}
106
107pub(crate) fn build_iceberg_task_report(
108    task_id: IcebergCompactionTaskId,
109    sink_id: u32,
110    error_message: Option<String>,
111) -> IcebergTaskReport {
112    subscribe_iceberg_compaction_event_request::ReportTask {
113        task_id,
114        sink_id,
115        status: if error_message.is_some() {
116            subscribe_iceberg_compaction_event_request::report_task::Status::Failed as i32
117        } else {
118            subscribe_iceberg_compaction_event_request::report_task::Status::Success as i32
119        },
120        error_message,
121    }
122}
123
124pub(crate) fn send_iceberg_task_report(
125    request_sender: &mpsc::UnboundedSender<SubscribeIcebergCompactionEventRequest>,
126    report_event: IcebergTaskReport,
127) -> Result<(), IcebergTaskReport> {
128    if let Err(e) = request_sender.send(SubscribeIcebergCompactionEventRequest {
129        event: Some(
130            subscribe_iceberg_compaction_event_request::Event::ReportTask(report_event.clone()),
131        ),
132        create_at: SystemTime::now()
133            .duration_since(std::time::UNIX_EPOCH)
134            .expect("Clock may have gone backwards")
135            .as_millis() as u64,
136    }) {
137        tracing::warn!(
138            iceberg_component = "compaction_worker",
139            iceberg_operation = "report_task",
140            error = %e.as_report(),
141            task_id = %report_event.task_id,
142            sink_id = report_event.sink_id,
143            "iceberg_compaction_task_report_send_failed",
144        );
145        return Err(report_event);
146    }
147
148    Ok(())
149}
150
151pub(crate) fn send_or_buffer_iceberg_task_report(
152    request_sender: &mpsc::UnboundedSender<SubscribeIcebergCompactionEventRequest>,
153    pending_task_reports: &mut VecDeque<IcebergTaskReport>,
154    report: IcebergTaskReport,
155) -> ReportSendResult {
156    if let Err(report) = send_iceberg_task_report(request_sender, report) {
157        pending_task_reports.push_back(report);
158        return ReportSendResult::RestartStream;
159    }
160    ReportSendResult::Sent
161}
162
163pub(crate) fn flush_pending_iceberg_task_reports(
164    request_sender: &mpsc::UnboundedSender<SubscribeIcebergCompactionEventRequest>,
165    pending_task_reports: &mut VecDeque<IcebergTaskReport>,
166) -> ReportSendResult {
167    while let Some(report_event) = pending_task_reports.pop_front() {
168        if let Err(report_event) = send_iceberg_task_report(request_sender, report_event) {
169            pending_task_reports.push_front(report_event);
170            return ReportSendResult::RestartStream;
171        }
172    }
173    ReportSendResult::Sent
174}
175
176#[cfg(test)]
177mod tests {
178    use risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_request;
179
180    use super::*;
181
182    #[test]
183    fn test_send_iceberg_task_report_returns_payload_on_send_failure() {
184        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
185        drop(rx);
186
187        let report = build_iceberg_task_report(7.into(), 9, Some("send failure".to_owned()));
188        let failed_report = send_iceberg_task_report(&tx, report.clone()).unwrap_err();
189
190        assert_eq!(failed_report.task_id, report.task_id);
191        assert_eq!(failed_report.sink_id, report.sink_id);
192        assert_eq!(failed_report.error_message, report.error_message);
193    }
194
195    #[test]
196    fn test_build_iceberg_task_result_partial_enqueue_is_success_if_admitted_plan_succeeds() {
197        let mut tracker = IcebergTaskTracker::new(9, 1);
198        tracker.record_completion(None);
199
200        let report = tracker.into_report(7.into());
201
202        assert_eq!(
203            report.status,
204            subscribe_iceberg_compaction_event_request::report_task::Status::Success as i32
205        );
206        assert!(report.error_message.is_none());
207    }
208
209    #[test]
210    fn test_build_iceberg_task_result_fails_if_all_admitted_plans_fail() {
211        let mut tracker = IcebergTaskTracker::new(9, 2);
212        tracker.record_completion(Some("first failure".to_owned()));
213        tracker.record_completion(Some("second failure".to_owned()));
214
215        let report = tracker.into_report(7.into());
216
217        assert_eq!(
218            report.status,
219            subscribe_iceberg_compaction_event_request::report_task::Status::Failed as i32
220        );
221        assert_eq!(report.error_message.as_deref(), Some("first failure"));
222    }
223}