1use api::v1::region::compact_request;
16use common_telemetry::{debug, error, info};
17use store_api::logstore::LogStore;
18use store_api::region_request::RegionCompactRequest;
19use store_api::storage::RegionId;
20
21use crate::compaction::{CompactionPickFinished, CompactionTransition};
22use crate::config::IndexBuildMode;
23use crate::error::{RegionNotFoundSnafu, StaleCompactionExecutionSnafu};
24use crate::metrics::COMPACTION_REQUEST_COUNT;
25use crate::region::MitoRegionRef;
26use crate::request::{
27 BuildIndexRequest, CompactionCancelled, CompactionFailed, CompactionFinished, OnFailure,
28 OptionOutputTx,
29};
30use crate::sst::index::IndexBuildType;
31use crate::worker::RegionWorkerLoop;
32
33fn made_progress(files_to_add: usize, files_to_remove: usize) -> bool {
34 files_to_add > 0 || files_to_remove > files_to_add
35}
36
37impl<S> RegionWorkerLoop<S> {
38 pub(crate) async fn handle_compaction_pick_finished(
39 &mut self,
40 region_id: RegionId,
41 request: CompactionPickFinished,
42 ) where
43 S: LogStore,
44 {
45 let Some(region) = self.regions.get_region(region_id) else {
46 return;
47 };
48 let transition = self
53 .compaction_scheduler
54 .handle_compaction_pick_finished(
55 request,
56 ®ion.manifest_ctx,
57 self.schema_metadata_manager.clone(),
58 )
59 .await;
60 match transition {
61 CompactionTransition::AutomaticFollowupScheduled => {
62 region.update_schedule_compaction_millis();
65 }
66 CompactionTransition::NoAction => {}
67 CompactionTransition::DdlReady(mut pending_ddls) => {
68 if !pending_ddls.is_empty() {
69 self.listener.on_compaction_result_notified(region_id).await;
72 self.handle_ddl_requests(&mut pending_ddls).await;
73 }
74 }
75 }
76 }
77
78 pub(crate) async fn handle_compaction_request(
80 &mut self,
81 region_id: RegionId,
82 req: RegionCompactRequest,
83 mut sender: OptionOutputTx,
84 ) {
85 let Some(region) = self.regions.writable_region_or(region_id, &mut sender) else {
86 return;
87 };
88 COMPACTION_REQUEST_COUNT.inc();
89 let parallelism = req.parallelism.unwrap_or(1) as usize;
90 match self.compaction_scheduler.schedule_manual_compaction(
91 req.options,
92 ®ion.version_control,
93 ®ion.access_layer,
94 sender,
95 ®ion.manifest_ctx,
96 self.schema_metadata_manager.clone(),
97 parallelism,
98 req.time_range,
99 ) {
100 Ok(true) => info!(
104 "Successfully scheduled compaction task for region: {}",
105 region_id
106 ),
107 Ok(false) => {}
108 Err(e) => {
109 error!(e; "Failed to schedule compaction task for region: {}", region_id);
110 }
111 }
112 }
113
114 pub(crate) async fn handle_compaction_finished(
116 &mut self,
117 region_id: RegionId,
118 mut request: CompactionFinished,
119 ) where
120 S: LogStore,
121 {
122 let region = match self.regions.get_region(region_id) {
123 Some(region) => region,
124 None => {
125 request.on_failure(RegionNotFoundSnafu { region_id }.build());
126 return;
127 }
128 };
129 if !self
131 .compaction_scheduler
132 .is_current_execution(region_id, &request.execution)
133 {
134 request.on_failure(StaleCompactionExecutionSnafu { region_id }.build());
135 return;
136 }
137 let execution = request.execution.clone();
138 let made_progress = made_progress(
139 request.edit.files_to_add.len(),
140 request.edit.files_to_remove.len(),
141 );
142
143 region.version_control.apply_edit(
144 Some(request.edit.clone()),
145 &[],
146 region.file_purger.clone(),
147 );
148
149 let index_build_file_metas = std::mem::take(&mut request.edit.files_to_add);
150
151 request.on_success();
153 self.listener.on_compaction_result_notified(region_id).await;
154
155 if self.config.index.build_mode == IndexBuildMode::Async
157 && !index_build_file_metas.is_empty()
158 {
159 self.handle_rebuild_index(
160 BuildIndexRequest {
161 region_id,
162 build_type: IndexBuildType::Compact,
163 file_metas: index_build_file_metas,
164 },
165 OptionOutputTx::new(None),
166 )
167 .await;
168 }
169
170 let transition = self
172 .compaction_scheduler
173 .on_execution_finished(
174 region_id,
175 &execution,
176 ®ion.manifest_ctx,
177 self.schema_metadata_manager.clone(),
178 made_progress,
179 )
180 .await;
181 match transition {
182 CompactionTransition::AutomaticFollowupScheduled => {
183 region.update_schedule_compaction_millis();
184 }
185 CompactionTransition::NoAction => {}
186 CompactionTransition::DdlReady(mut pending_ddls) => {
187 self.handle_ddl_requests(&mut pending_ddls).await;
188 }
189 }
190 }
191
192 pub(crate) async fn handle_compaction_cancelled(
193 &mut self,
194 region_id: RegionId,
195 request: CompactionCancelled,
196 ) where
197 S: LogStore,
198 {
199 let execution = request.execution.clone();
200 let is_current = self.regions.get_region(region_id).is_some_and(|_| {
201 self.compaction_scheduler
202 .is_current_execution(region_id, &execution)
203 });
204 request.on_success();
205
206 if !is_current {
207 return;
208 }
209
210 let mut pending_ddls = self
212 .compaction_scheduler
213 .on_execution_cancelled(region_id, &execution)
214 .await;
215 if !pending_ddls.is_empty() {
216 self.listener.on_compaction_result_notified(region_id).await;
217 }
218
219 self.handle_ddl_requests(&mut pending_ddls).await;
220 }
221
222 pub(crate) async fn handle_compaction_failure(&mut self, req: CompactionFailed) {
224 if self.regions.get_region(req.region_id).is_none() {
225 return;
226 }
227 if !self
228 .compaction_scheduler
229 .is_current_execution(req.region_id, &req.execution)
230 {
231 debug!(
232 "Ignores stale compaction failure for region {}: {:?}",
233 req.region_id, req.err
234 );
235 return;
236 }
237
238 error!(req.err; "Failed to compact region: {}", req.region_id);
239 self.compaction_scheduler
240 .on_execution_failed(req.region_id, &req.execution, req.err);
241 }
242
243 pub(crate) async fn schedule_compaction(&mut self, region: &MitoRegionRef) {
245 if region.is_staging() || region.is_enter_staging() {
246 info!(
247 "Region {} is staging or entering staging, skip compaction",
248 region.region_id
249 );
250 return;
251 }
252 let now = self.time_provider.current_time_millis();
253 if now - region.last_schedule_compaction_millis()
254 >= self.config.min_compaction_interval.as_millis() as i64
255 {
256 debug!(
257 "minimal compaction interval time {:?} has passed, scheduling next compaction",
258 self.config.min_compaction_interval
259 );
260 match self.compaction_scheduler.schedule_automatic_compaction(
261 compact_request::Options::Regular(Default::default()),
262 ®ion.version_control,
263 ®ion.access_layer,
264 ®ion.manifest_ctx,
265 self.schema_metadata_manager.clone(),
266 ) {
267 Ok(true) => region.update_schedule_compaction_millis(),
268 Ok(false) => {}
269 Err(e) => {
270 error!(e; "Failed to schedule compaction for region: {}", region.region_id)
271 }
272 }
273 }
274 }
275}
276
277#[cfg(test)]
278mod tests {
279 use super::made_progress;
280
281 #[test]
282 fn test_nonempty_output_or_file_reduction_is_progress() {
283 for (files_to_add, files_to_remove, expected) in [
284 (3, 3, true), (1, 3, true), (0, 3, true), (0, 0, false), (3, 2, true), ] {
290 assert_eq!(
291 expected,
292 made_progress(files_to_add, files_to_remove),
293 "files_to_remove: {files_to_remove}, files_to_add: {files_to_add}"
294 );
295 }
296 }
297}