risingwave_storage/hummock/compactor/iceberg_compaction/
report.rs1use 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}