1use std::time::Instant;
16
17use common_telemetry::{error, info, warn};
18use store_api::logstore::LogStore;
19use store_api::region_request::{EnterStagingRequest, StagingPartitionDirective};
20use store_api::storage::RegionId;
21
22use crate::error::{RegionNotFoundSnafu, Result, StagingPartitionExprMismatchSnafu};
23use crate::flush::FlushReason;
24use crate::manifest::action::{RegionMetaAction, RegionMetaActionList, RegionPartitionExprChange};
25use crate::region::{MitoRegionRef, RegionLeaderState, StagingPartitionInfo};
26use crate::request::{
27 BackgroundNotify, DdlRequest, EnterStagingResult, OptionOutputTx, SenderDdlRequest,
28 WorkerRequest, WorkerRequestWithTime,
29};
30use crate::worker::RegionWorkerLoop;
31
32impl<S: LogStore> RegionWorkerLoop<S> {
33 pub(crate) async fn handle_enter_staging_request(
34 &mut self,
35 region_id: RegionId,
36 partition_directive: StagingPartitionDirective,
37 mut sender: OptionOutputTx,
38 ) {
39 let Some(region) = self.regions.writable_region_or(region_id, &mut sender) else {
40 return;
41 };
42
43 if region.is_staging() {
45 let staging_partition_info = region.manifest_ctx.staging_partition_info();
46 if staging_partition_info
48 .as_ref()
49 .map(|info| &info.partition_directive)
50 != Some(&partition_directive)
51 {
52 sender.send(Err(StagingPartitionExprMismatchSnafu {
53 manifest_expr: staging_partition_info
54 .as_ref()
55 .and_then(|info| info.partition_expr().map(ToString::to_string)),
56 request_expr: format!("{:?}", partition_directive),
57 }
58 .build()));
59 return;
60 }
61
62 sender.send(Ok(0));
64 return;
65 }
66
67 let version = region.version();
68 if !version.memtables.is_empty() {
69 info!("Flush region: {} before entering staging", region_id);
72 debug_assert!(!region.is_staging());
73 let task = self.new_flush_task(
74 ®ion,
75 FlushReason::EnterStaging,
76 None,
77 self.config.clone(),
78 );
79 if let Err(e) =
80 self.flush_scheduler
81 .schedule_flush(region.region_id, ®ion.version_control, task)
82 {
83 sender.send(Err(e));
85 return;
86 }
87
88 self.flush_scheduler
90 .add_ddl_request_to_pending(SenderDdlRequest {
91 region_id,
92 sender,
93 request: DdlRequest::EnterStaging(EnterStagingRequest {
94 partition_directive: partition_directive.clone(),
95 }),
96 });
97
98 return;
99 }
100
101 let (sender, partition_directive) = match self.compaction_scheduler.try_cancel_and_add_ddl(
102 region_id,
103 sender,
104 partition_directive,
105 |partition_directive| {
106 DdlRequest::EnterStaging(EnterStagingRequest {
107 partition_directive,
108 })
109 },
110 ) {
111 Ok(()) => {
112 self.listener.on_compaction_cancel_requested(region_id);
113 return;
114 }
115 Err(request) => request,
116 };
117
118 self.handle_enter_staging(region, partition_directive, sender);
119 }
120
121 async fn enter_staging(
122 region: &MitoRegionRef,
123 partition_directive: &StagingPartitionDirective,
124 ) -> Result<()> {
125 let now = Instant::now();
126 {
128 let mut manager = region.manifest_ctx.manifest_manager.write().await;
129 manager
130 .clear_staging_manifest_and_dir()
131 .await
132 .inspect_err(|e| {
133 error!(
134 e;
135 "Failed to clear staging manifest files for region {}",
136 region.region_id
137 );
138 })?;
139
140 info!(
141 "Cleared all staging manifest files for region {}, elapsed: {:?}",
142 region.region_id,
143 now.elapsed(),
144 );
145 }
146
147 let partition_expr = match partition_directive {
148 StagingPartitionDirective::UpdatePartitionExpr(partition_expr) => {
149 partition_expr.clone()
150 }
151 StagingPartitionDirective::RejectAllWrites => {
152 info!(
153 "Enter staging with reject all writes, region_id: {}",
154 region.region_id
155 );
156 return Ok(());
158 }
159 };
160
161 let change = RegionPartitionExprChange {
163 partition_expr: Some(partition_expr.clone()),
164 };
165 let action_list =
166 RegionMetaActionList::with_action(RegionMetaAction::PartitionExprChange(change));
167 region
168 .manifest_ctx
169 .update_manifest(RegionLeaderState::EnteringStaging, action_list, true)
170 .await?;
171
172 Ok(())
173 }
174
175 fn handle_enter_staging(
176 &self,
177 region: MitoRegionRef,
178 partition_directive: StagingPartitionDirective,
179 sender: OptionOutputTx,
180 ) {
181 if let Err(e) = region.set_entering_staging() {
182 sender.send(Err(e));
183 return;
184 }
185
186 let listener = self.listener.clone();
187 let request_sender = self.sender.clone();
188 common_runtime::spawn_global(async move {
189 let now = Instant::now();
190 let result = Self::enter_staging(®ion, &partition_directive).await;
191 match result {
192 Ok(_) => {
193 info!(
194 "Created staging manifest for region {}, elapsed: {:?}",
195 region.region_id,
196 now.elapsed(),
197 );
198 }
199 Err(ref e) => {
200 region
202 .manifest_ctx
203 .manifest_manager
204 .write()
205 .await
206 .unset_staging_manifest();
207 error!(
208 "Failed to create staging manifest for region {}: {:?}, elapsed: {:?}",
209 region.region_id,
210 e,
211 now.elapsed(),
212 );
213 }
214 }
215
216 let notify = WorkerRequest::Background {
217 region_id: region.region_id,
218 notify: BackgroundNotify::EnterStaging(EnterStagingResult {
219 region_id: region.region_id,
220 sender,
221 result,
222 partition_directive,
223 }),
224 };
225 listener
226 .on_enter_staging_result_begin(region.region_id)
227 .await;
228
229 if let Err(res) = request_sender
230 .send(WorkerRequestWithTime::new(notify))
231 .await
232 {
233 warn!(
234 "Failed to send enter staging result back to the worker, region_id: {}, res: {:?}",
235 region.region_id, res
236 );
237 }
238 });
239 }
240
241 pub(crate) async fn handle_enter_staging_result(
243 &mut self,
244 enter_staging_result: EnterStagingResult,
245 ) {
246 let region = match self.regions.get_region(enter_staging_result.region_id) {
247 Some(region) => region,
248 None => {
249 self.reject_region_stalled_requests(&enter_staging_result.region_id);
250 enter_staging_result.sender.send(
251 RegionNotFoundSnafu {
252 region_id: enter_staging_result.region_id,
253 }
254 .fail(),
255 );
256 return;
257 }
258 };
259
260 if enter_staging_result.result.is_ok() {
261 info!(
262 "Updating region {} staging partition directive to {:?}",
263 region.region_id, enter_staging_result.partition_directive
264 );
265 Self::update_region_staging_partition_info(
266 ®ion,
267 enter_staging_result.partition_directive,
268 );
269 region.switch_state_to_staging(RegionLeaderState::EnteringStaging);
270 } else {
271 region.switch_state_to_writable(RegionLeaderState::EnteringStaging);
272 }
273 enter_staging_result
274 .sender
275 .send(enter_staging_result.result.map(|_| 0));
276 self.handle_region_stalled_requests(&enter_staging_result.region_id, true)
278 .await;
279 }
280
281 fn update_region_staging_partition_info(
282 region: &MitoRegionRef,
283 partition_directive: StagingPartitionDirective,
284 ) {
285 region.manifest_ctx.set_staging_partition_info(
286 StagingPartitionInfo::from_partition_directive(partition_directive),
287 );
288 }
289}