1use std::any::Any;
16use std::collections::{BTreeSet, HashMap, HashSet};
17use std::sync::Arc;
18use std::time::Duration;
19
20use api::v1::meta::MailboxMessage;
21use common_error::ext::PlainError;
22use common_error::status_code::StatusCode;
23use common_meta::instruction::{
24 self, GetPackedFileRefs, GetPackedFileRefsReply, InstructionReply, PackedFileRefsManifest,
25 PackedGcRegions, PackedRegionFileRefs,
26};
27use common_meta::key::TableMetadataManagerRef;
28use common_meta::key::runtime_switch::RuntimeSwitchManagerRef;
29use common_meta::key::table_repart::TableRepartValue;
30use common_meta::key::table_route::PhysicalTableRouteValue;
31use common_meta::lock_key::{RegionLock, TableLock};
32use common_meta::peer::Peer;
33use common_meta::rpc::ddl::TriggerReason;
34use common_procedure::error::ToJsonSnafu;
35use common_procedure::{
36 Context as ProcedureContext, Error as ProcedureError, EventContext, EventTrigger, LockKey,
37 Procedure, ProcedureState, Result as ProcedureResult, Status,
38};
39use common_telemetry::tracing::Instrument as _;
40use common_telemetry::tracing_context::TracingContext;
41use common_telemetry::{debug, error, info, warn};
42use futures::future::join_all;
43use futures::stream::{FuturesUnordered, StreamExt};
44use itertools::Itertools as _;
45use serde::{Deserialize, Serialize};
46use snafu::ResultExt as _;
47use store_api::storage::{FileRefsManifest, GcReport, RegionId};
48use table::metadata::TableId;
49
50use crate::error::{self, KvBackendSnafu, Result, SerializeToJsonSnafu, TableMetadataManagerSnafu};
51use crate::event::gc::{BATCH_GC_EVENT_TYPE, BatchGcEvent};
52use crate::gc::util::table_route_to_region;
53use crate::gc::{Peer2Regions, Region2Peers};
54use crate::handler::HeartbeatMailbox;
55use crate::metrics::{METRIC_META_GC_DATANODE_CALLS_TOTAL, METRIC_META_GC_FAILED_REGIONS_TOTAL};
56use crate::procedure::utils::{instruction_error_result, instruction_to_error};
57use crate::service::mailbox::{Channel, MailboxReceiver, MailboxRef};
58
59async fn send_get_packed_file_refs_inner(
60 mailbox: &MailboxRef,
61 server_addr: &str,
62 peer: &Peer,
63 instruction: GetPackedFileRefs,
64 timeout: Duration,
65) -> Result<MailboxReceiver> {
66 let instruction = instruction::Instruction::GetPackedFileRefs(instruction);
67 let tracing_ctx = TracingContext::from_current_span();
68 let msg = MailboxMessage::json_message(
69 &format!("Get file references: {}", instruction),
70 &format!("Metasrv@{}", server_addr),
71 &format!("Datanode-{}@{}", peer.id, peer.addr),
72 common_time::util::current_time_millis(),
73 &instruction,
74 Some(tracing_ctx.to_w3c()),
75 )
76 .with_context(|_| SerializeToJsonSnafu {
77 input: instruction.to_string(),
78 })?;
79
80 mailbox
81 .send(&Channel::Datanode(peer.id), msg, timeout)
82 .await
83}
84
85async fn recv_get_packed_file_refs_reply(
86 peer: &Peer,
87 mailbox_rx: MailboxReceiver,
88) -> Result<GetPackedFileRefsReply> {
89 let reply = match mailbox_rx.await {
90 Ok(reply_msg) => HeartbeatMailbox::json_reply(&reply_msg)?,
91 Err(e) => {
92 error!(
93 e; "Failed to receive reply from datanode {} for GetPackedFileRefs instruction",
94 peer,
95 );
96 return Err(e);
97 }
98 };
99
100 let InstructionReply::GetPackedFileRefs(reply) = reply else {
101 return error::UnexpectedInstructionReplySnafu {
102 mailbox_message: "unexpected instruction reply for GetPackedFileRefs".to_string(),
103 reason: "Unexpected reply of the GetPackedFileRefs instruction",
104 }
105 .fail();
106 };
107
108 Ok(reply)
109}
110
111async fn send_gc_regions_inner(
112 mailbox: &MailboxRef,
113 peer: &Peer,
114 instruction: instruction::Instruction,
115 server_addr: &str,
116 timeout: Duration,
117 description: &str,
118) -> Result<MailboxReceiver> {
119 let tracing_ctx = TracingContext::from_current_span();
120 let msg = MailboxMessage::json_message(
121 &format!("{}: {}", description, instruction),
122 &format!("Metasrv@{}", server_addr),
123 &format!("Datanode-{}@{}", peer.id, peer.addr),
124 common_time::util::current_time_millis(),
125 &instruction,
126 Some(tracing_ctx.to_w3c()),
127 )
128 .with_context(|_| SerializeToJsonSnafu {
129 input: instruction.to_string(),
130 })?;
131
132 mailbox
133 .send(&Channel::Datanode(peer.id), msg, timeout)
134 .await
135}
136
137fn scoped_gc_instruction(
138 regions: Vec<RegionId>,
139 full_manifest: &FileRefsManifest,
140 full_file_listing: bool,
141) -> instruction::Instruction {
142 let packed_file_refs_manifest = PackedFileRefsManifest {
143 file_refs: regions
144 .iter()
145 .filter_map(|region| {
146 full_manifest
147 .file_refs
148 .get(region)
149 .map(|refs| (*region, PackedRegionFileRefs::from_refs(refs)))
150 })
151 .collect(),
152 manifest_version: regions
153 .iter()
154 .filter_map(|region| {
155 full_manifest
156 .manifest_version
157 .get(region)
158 .map(|version| (*region, *version))
159 })
160 .collect(),
161 cross_region_refs: HashMap::new(),
162 };
163
164 instruction::Instruction::PackedGcRegions(PackedGcRegions {
165 regions,
166 packed_file_refs_manifest,
167 full_file_listing,
168 })
169}
170
171async fn recv_gc_regions_reply(
172 peer: &Peer,
173 regions: &[RegionId],
174 description: &str,
175 mailbox_rx: MailboxReceiver,
176) -> Result<GcReport> {
177 let reply = match mailbox_rx.await {
178 Ok(reply_msg) => HeartbeatMailbox::json_reply(&reply_msg)?,
179 Err(e) => {
180 error!(
181 e; "Failed to receive reply from datanode {} for {}",
182 peer, description
183 );
184 return Err(e);
185 }
186 };
187
188 let InstructionReply::GcRegions(reply) = reply else {
189 return error::UnexpectedInstructionReplySnafu {
190 mailbox_message: "unexpected instruction reply for GcRegions".to_string(),
191 reason: "Unexpected reply of the GcRegions instruction",
192 }
193 .fail();
194 };
195
196 let res = reply.result;
197 match res {
198 Ok(report) => Ok(report),
199 Err(e) => {
200 error!(
201 e; "Datanode {} reported error during GC for {} regions",
202 peer, regions.len()
203 );
204 instruction_error_result(
205 &e,
206 format!(
207 "Datanode {} reported error during GC for {} regions: {}",
208 peer,
209 regions.len(),
210 e
211 ),
212 )
213 }
214 }
215}
216
217pub struct BatchGcProcedure {
220 mailbox: MailboxRef,
221 table_metadata_manager: TableMetadataManagerRef,
222 runtime_switch_manager: RuntimeSwitchManagerRef,
223 data: BatchGcData,
224}
225
226#[derive(Serialize, Deserialize)]
227pub struct BatchGcData {
228 state: State,
229 server_addr: String,
231 regions: Vec<RegionId>,
233 full_file_listing: bool,
234 region_routes: Region2Peers,
235 #[serde(default)]
237 region_routes_override: Region2Peers,
238 related_regions: HashMap<RegionId, HashSet<RegionId>>,
241 file_refs: FileRefsManifest,
243 timeout: Duration,
245 gc_report: Option<GcReport>,
246}
247
248#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
249pub enum State {
250 Start,
252 Acquiring,
254 Gcing,
256 UpdateRepartition,
258}
259
260impl BatchGcProcedure {
261 pub const TYPE_NAME: &'static str = "metasrv-procedure::BatchGcProcedure";
262
263 #[allow(clippy::too_many_arguments)]
264 pub fn new(
265 mailbox: MailboxRef,
266 table_metadata_manager: TableMetadataManagerRef,
267 runtime_switch_manager: RuntimeSwitchManagerRef,
268 server_addr: String,
269 regions: Vec<RegionId>,
270 full_file_listing: bool,
271 timeout: Duration,
272 region_routes_override: Region2Peers,
273 ) -> Self {
274 Self {
275 mailbox,
276 table_metadata_manager,
277 runtime_switch_manager,
278 data: BatchGcData {
279 state: State::Start,
280 server_addr,
281 regions,
282 full_file_listing,
283 timeout,
284 region_routes: HashMap::new(),
285 region_routes_override,
286 related_regions: HashMap::new(),
287 file_refs: FileRefsManifest::default(),
288 gc_report: None,
289 },
290 }
291 }
292
293 #[cfg(feature = "mock")]
297 pub fn new_update_repartition_for_test(
298 mailbox: MailboxRef,
299 table_metadata_manager: TableMetadataManagerRef,
300 runtime_switch_manager: RuntimeSwitchManagerRef,
301 server_addr: String,
302 regions: Vec<RegionId>,
303 file_refs: FileRefsManifest,
304 timeout: Duration,
305 ) -> Self {
306 Self {
307 mailbox,
308 table_metadata_manager,
309 runtime_switch_manager,
310 data: BatchGcData {
311 state: State::UpdateRepartition,
312 server_addr,
313 regions,
314 full_file_listing: false,
315 timeout,
316 region_routes: HashMap::new(),
317 region_routes_override: HashMap::new(),
318 related_regions: HashMap::new(),
319 file_refs,
320 gc_report: Some(GcReport::default()),
321 },
322 }
323 }
324
325 pub fn cast_result(res: Arc<dyn Any>) -> Result<GcReport> {
326 res.downcast_ref::<GcReport>().cloned().ok_or_else(|| {
327 error::UnexpectedSnafu {
328 violated: format!(
329 "Failed to downcast procedure result to GcReport, got {:?}",
330 std::any::type_name_of_val(&res.as_ref())
331 ),
332 }
333 .build()
334 })
335 }
336
337 async fn check_maintenance_mode(&self) -> ProcedureResult<()> {
338 let enabled = self
339 .runtime_switch_manager
340 .maintenance_mode()
341 .await
342 .context(error::RuntimeSwitchManagerSnafu)
343 .map_err(ProcedureError::retry_later)?;
344 if enabled {
345 return Err(ProcedureError::external(PlainError::new(
346 "maintenance mode is enabled".to_string(),
347 StatusCode::IllegalState,
348 )));
349 }
350 Ok(())
351 }
352
353 fn merge_gc_report(&mut self, report: GcReport) {
354 let accumulated = self.data.gc_report.get_or_insert_default();
355 let affected_regions: HashSet<_> = report
358 .processed_regions
359 .iter()
360 .chain(&report.need_retry_regions)
361 .copied()
362 .collect();
363
364 let mut processed_regions = std::mem::take(&mut accumulated.processed_regions);
365 processed_regions.retain(|region| !affected_regions.contains(region));
366 processed_regions.extend(report.processed_regions.iter().copied());
367
368 let mut need_retry_regions = std::mem::take(&mut accumulated.need_retry_regions);
369 need_retry_regions.retain(|region| !affected_regions.contains(region));
370 need_retry_regions.extend(report.need_retry_regions.iter().copied());
371
372 accumulated.merge(report);
373 accumulated.processed_regions = processed_regions;
374 accumulated.need_retry_regions = need_retry_regions;
375 }
376
377 fn done_with_gc_report(&self) -> ProcedureResult<Status> {
378 let Some(report) = self.data.gc_report.clone() else {
379 return common_procedure::error::UnexpectedSnafu {
380 err_msg: "GC report should be present after GC completion".to_string(),
381 }
382 .fail();
383 };
384
385 Ok(Status::done_with_output(report))
386 }
387
388 #[cfg(test)]
389 pub(crate) fn set_gc_report_for_test(&mut self, report: GcReport) {
390 self.data.gc_report = Some(report);
391 }
392
393 async fn get_table_route(
394 &self,
395 table_id: TableId,
396 ) -> Result<(TableId, PhysicalTableRouteValue)> {
397 self.table_metadata_manager
398 .table_route_manager()
399 .get_physical_table_route(table_id)
400 .await
401 .context(TableMetadataManagerSnafu)
402 }
403
404 async fn find_related_regions(
408 &self,
409 regions: &[RegionId],
410 ) -> Result<HashMap<RegionId, HashSet<RegionId>>> {
411 let table_ids: HashSet<TableId> = regions.iter().map(|r| r.table_id()).collect();
412 let table_ids = table_ids.into_iter().collect::<Vec<_>>();
413
414 let table_routes = self
415 .table_metadata_manager
416 .table_route_manager()
417 .batch_get_physical_table_routes(&table_ids)
418 .await
419 .context(TableMetadataManagerSnafu)?;
420
421 if table_routes.len() != table_ids.len() {
422 for table_id in &table_ids {
424 if !table_routes.contains_key(table_id) {
425 return error::InvalidArgumentsSnafu {
427 err_msg: format!(
428 "Unexpected logical table route: table {} resolved to physical table regions",
429 table_id
430 ),
431 }
432 .fail();
433 }
434 }
435 }
436
437 let mut table_all_regions: HashMap<TableId, HashSet<RegionId>> = HashMap::new();
438 for (table_id, table_route) in table_routes {
439 let all_regions: HashSet<RegionId> = table_route
440 .region_routes
441 .iter()
442 .map(|r| r.region.id)
443 .collect();
444
445 table_all_regions.insert(table_id, all_regions);
446 }
447
448 let mut related_regions: HashMap<RegionId, HashSet<RegionId>> = HashMap::new();
449 for region_id in regions {
450 let table_id = region_id.table_id();
451 if let Some(all_regions) = table_all_regions.get(&table_id) {
452 let mut related: HashSet<RegionId> = all_regions.clone();
453 related.remove(region_id);
454 related_regions.insert(*region_id, related);
455 } else {
456 related_regions.insert(*region_id, Default::default());
457 }
458 }
459
460 Ok(related_regions)
461 }
462
463 async fn cleanup_region_repartition(&self, procedure_ctx: &ProcedureContext) -> Result<()> {
466 let mut cross_refs_grouped: HashMap<TableId, HashMap<RegionId, HashSet<RegionId>>> =
467 HashMap::new();
468 for (src_region, dst_regions) in &self.data.file_refs.cross_region_refs {
469 cross_refs_grouped
470 .entry(src_region.table_id())
471 .or_default()
472 .entry(*src_region)
473 .or_default()
474 .extend(dst_regions.iter().copied());
475 }
476
477 let mut tmp_refs_grouped: HashMap<TableId, HashSet<RegionId>> = HashMap::new();
478 for (src_region, refs) in &self.data.file_refs.file_refs {
479 if refs.is_empty() {
480 continue;
481 }
482
483 tmp_refs_grouped
484 .entry(src_region.table_id())
485 .or_default()
486 .insert(*src_region);
487 }
488
489 let repart_mgr = self.table_metadata_manager.table_repart_manager();
490
491 let need_retry: HashSet<RegionId> = self
494 .data
495 .gc_report
496 .as_ref()
497 .map(|r| r.need_retry_regions.clone())
498 .unwrap_or_default();
499
500 let mut table_ids: HashSet<TableId> = cross_refs_grouped
501 .keys()
502 .copied()
503 .chain(tmp_refs_grouped.keys().copied())
504 .collect();
505 table_ids.extend(self.data.regions.iter().map(|r| r.table_id()));
506
507 for table_id in table_ids {
508 let table_lock = TableLock::Write(table_id).into();
509 let _guard = procedure_ctx.provider.acquire_lock(&table_lock).await;
510
511 let cross_refs = cross_refs_grouped
512 .get(&table_id)
513 .cloned()
514 .unwrap_or_default();
515 let tmp_refs = tmp_refs_grouped.get(&table_id).cloned().unwrap_or_default();
516
517 let current = repart_mgr
518 .get_with_raw_bytes(table_id)
519 .await
520 .context(KvBackendSnafu)?;
521
522 let mut new_value = current
523 .as_ref()
524 .map(|v| (**v).clone())
525 .unwrap_or_else(TableRepartValue::new);
526
527 let batch_src_regions: HashSet<RegionId> = self
530 .data
531 .regions
532 .iter()
533 .copied()
534 .filter(|r| r.table_id() == table_id)
535 .collect();
536
537 let all_src_regions: HashSet<RegionId> = batch_src_regions;
541
542 for src_region in all_src_regions {
543 let cross_dst = cross_refs.get(&src_region);
544 let has_tmp_ref = tmp_refs.contains(&src_region);
545
546 if let Some(dst_regions) = cross_dst {
547 let mut set = BTreeSet::new();
548 set.extend(dst_regions.iter().copied());
549 new_value.src_to_dst.insert(src_region, set);
550 } else if has_tmp_ref || need_retry.contains(&src_region) {
551 new_value.src_to_dst.insert(src_region, BTreeSet::new());
554 } else {
555 new_value.src_to_dst.remove(&src_region);
556 }
557 }
558
559 if new_value.src_to_dst.is_empty() && current.is_none() {
561 continue;
562 }
563
564 repart_mgr
565 .upsert_value(table_id, current, &new_value)
566 .await
567 .context(KvBackendSnafu)?;
568 }
569
570 Ok(())
571 }
572
573 async fn discover_route_for_regions(
575 &self,
576 regions: &[RegionId],
577 ) -> Result<(Region2Peers, Peer2Regions)> {
578 let mut region_to_peer = HashMap::new();
579 let mut peer_to_regions = HashMap::new();
580
581 let mut table_to_regions: HashMap<TableId, Vec<RegionId>> = HashMap::new();
583 for region_id in regions {
584 let table_id = region_id.table_id();
585 table_to_regions
586 .entry(table_id)
587 .or_default()
588 .push(*region_id);
589 }
590
591 for (table_id, table_regions) in table_to_regions {
593 match self.get_table_route(table_id).await {
594 Ok((_phy_table_id, table_route)) => {
595 table_route_to_region(
596 &table_route,
597 &table_regions,
598 &mut region_to_peer,
599 &mut peer_to_regions,
600 );
601 }
602 Err(e) => {
603 warn!(
606 "Failed to get table route for table {}: {}, skipping its regions",
607 table_id, e
608 );
609 continue;
610 }
611 }
612 }
613
614 Ok((region_to_peer, peer_to_regions))
615 }
616
617 async fn set_routes_and_related_regions(&mut self) -> Result<()> {
619 let related_regions = self.find_related_regions(&self.data.regions).await?;
620
621 self.data.related_regions = related_regions.clone();
622
623 let mut regions_set: HashSet<RegionId> = self.data.regions.iter().cloned().collect();
626
627 regions_set.extend(related_regions.keys().cloned());
628 regions_set.extend(related_regions.values().flat_map(|v| v.iter()).cloned());
629
630 let regions_to_discover = regions_set.into_iter().collect_vec();
631
632 let (mut region_to_peer, _) = self
633 .discover_route_for_regions(®ions_to_discover)
634 .await?;
635
636 for (region_id, route) in &self.data.region_routes_override {
637 region_to_peer
638 .entry(*region_id)
639 .or_insert_with(|| route.clone());
640 }
641
642 self.data.region_routes = region_to_peer;
643
644 Ok(())
645 }
646
647 async fn get_file_references(&mut self) -> Result<FileRefsManifest> {
649 let region_count = self.data.regions.len();
650 self.set_routes_and_related_regions()
651 .instrument(common_telemetry::tracing::info_span!(
652 "meta_gc_procedure_prepare_routes",
653 region_count = region_count
654 ))
655 .await?;
656
657 let query_regions = &self.data.regions;
658 let related_regions = &self.data.related_regions;
659 let region_routes = &self.data.region_routes;
660 let timeout = self.data.timeout;
661 let dropped_regions = self
662 .data
663 .region_routes_override
664 .keys()
665 .collect::<HashSet<_>>();
666
667 let mut datanode2query_regions: HashMap<Peer, Vec<RegionId>> = HashMap::new();
669
670 for region_id in query_regions {
671 if dropped_regions.contains(region_id) {
672 continue;
673 }
674 if let Some((leader, followers)) = region_routes.get(region_id) {
675 datanode2query_regions
676 .entry(leader.clone())
677 .or_default()
678 .push(*region_id);
679 for follower in followers {
681 datanode2query_regions
682 .entry(follower.clone())
683 .or_default()
684 .push(*region_id);
685 }
686 } else {
687 return error::UnexpectedSnafu {
688 violated: format!(
689 "region_routes: {region_routes:?} does not contain region_id: {region_id}",
690 ),
691 }
692 .fail();
693 }
694 }
695
696 let mut datanode2related_regions: HashMap<Peer, HashMap<RegionId, HashSet<RegionId>>> =
697 HashMap::new();
698 for (src_region, dst_regions) in related_regions {
699 for dst_region in dst_regions {
700 if let Some((leader, _followers)) = region_routes.get(dst_region) {
701 datanode2related_regions
702 .entry(leader.clone())
703 .or_default()
704 .entry(*src_region)
705 .or_default()
706 .insert(*dst_region);
707 } }
709 }
710
711 let mut all_file_refs: HashMap<RegionId, HashSet<_>> = HashMap::new();
713 let mut all_manifest_versions = HashMap::new();
714 let mut all_cross_region_refs: HashMap<RegionId, HashSet<RegionId>> = HashMap::new();
715
716 let mut peers = HashSet::new();
717 peers.extend(datanode2query_regions.keys().cloned());
718 peers.extend(datanode2related_regions.keys().cloned());
719
720 let mailbox = &self.mailbox;
721 let server_addr = &self.data.server_addr;
722 let mut tasks = FuturesUnordered::new();
725 for peer in peers {
726 let regions = datanode2query_regions.remove(&peer).unwrap_or_default();
727 let related_regions_for_peer =
728 datanode2related_regions.remove(&peer).unwrap_or_default();
729 if regions.is_empty() && related_regions_for_peer.is_empty() {
730 continue;
731 }
732 tasks.push(async move {
733 let instruction = GetPackedFileRefs {
734 query_regions: regions,
735 related_regions: related_regions_for_peer,
736 };
737 let rx = send_get_packed_file_refs_inner(
738 mailbox,
739 server_addr,
740 &peer,
741 instruction,
742 timeout,
743 )
744 .await?;
745 let reply = recv_get_packed_file_refs_reply(&peer, rx).await?;
746 Ok::<_, crate::error::Error>((peer, reply))
747 });
748 }
749
750 let mut first_error = None;
751 let mut record_get_file_refs_error = |e| {
752 METRIC_META_GC_DATANODE_CALLS_TOTAL
753 .with_label_values(&["get_file_refs", "error"])
754 .inc();
755 if first_error.is_none() {
756 first_error = Some(e);
757 }
758 };
759 while let Some(result) = tasks.next().await {
760 let (peer, reply) = match result {
761 Ok(reply) => reply,
762 Err(e) => {
763 record_get_file_refs_error(e);
764 continue;
765 }
766 };
767 if !reply.success {
768 let err = if let Some(error) = &reply.error {
769 instruction_to_error(
770 error,
771 format!("Failed to get file references from datanode {peer}"),
772 )
773 } else {
774 error::UnexpectedSnafu {
775 violated:
776 "Datanode returned an unsuccessful GetPackedFileRefs reply without an error"
777 .to_string(),
778 }
779 .build()
780 };
781 record_get_file_refs_error(err);
782 continue;
783 }
784 let manifest = match reply.packed_file_refs_manifest.into_manifest() {
785 Ok(manifest) => manifest,
786 Err(err) => {
787 record_get_file_refs_error(error::Error::Other {
788 source: common_error::ext::BoxedError::new(err),
789 location: snafu::Location::new(file!(), line!(), 0),
790 });
791 continue;
792 }
793 };
794 METRIC_META_GC_DATANODE_CALLS_TOTAL
795 .with_label_values(&["get_file_refs", "success"])
796 .inc();
797 for (region_id, refs) in manifest.file_refs {
798 all_file_refs.entry(region_id).or_default().extend(refs);
799 }
800 for (region_id, version) in manifest.manifest_version {
801 let entry = all_manifest_versions.entry(region_id).or_insert(version);
802 *entry = (*entry).min(version);
803 }
804 for (region_id, related) in manifest.cross_region_refs {
805 all_cross_region_refs
806 .entry(region_id)
807 .or_default()
808 .extend(related);
809 }
810 }
811
812 if let Some(e) = first_error {
813 return Err(e);
814 }
815
816 Ok(FileRefsManifest {
817 file_refs: all_file_refs,
818 manifest_version: all_manifest_versions,
819 cross_region_refs: all_cross_region_refs,
820 })
821 }
822
823 async fn send_gc_instructions(&mut self) -> Result<()> {
825 let regions = &self.data.regions;
826 let region_routes = &self.data.region_routes;
827 let file_refs = &self.data.file_refs;
828 let timeout = self.data.timeout;
829
830 let mut datanode2regions: HashMap<Peer, Vec<RegionId>> = HashMap::new();
832 let mut all_report = GcReport::default();
833
834 for region_id in regions {
835 if let Some((leader, _followers)) = region_routes.get(region_id) {
836 datanode2regions
837 .entry(leader.clone())
838 .or_default()
839 .push(*region_id);
840 } else {
841 return error::UnexpectedSnafu {
842 violated: format!(
843 "region_routes: {region_routes:?} does not contain region_id: {region_id}",
844 ),
845 }
846 .fail();
847 }
848 }
849
850 let mut all_need_retry = HashSet::new();
851 let mailbox = &self.mailbox;
852 let server_addr = self.data.server_addr.as_str();
853 let full_file_listing = self.data.full_file_listing;
854 let tasks = datanode2regions
855 .into_iter()
856 .map(|(peer, regions_for_peer)| {
857 let region_count = regions_for_peer.len() as u64;
858 let regions = regions_for_peer.clone();
859 let instruction =
860 scoped_gc_instruction(regions_for_peer, file_refs, full_file_listing);
861
862 async move {
863 let report = send_gc_regions_inner(
864 mailbox,
865 &peer,
866 instruction,
867 server_addr,
868 timeout,
869 "Batch GC",
870 )
871 .await;
872
873 (peer, regions, region_count, report)
874 }
875 });
876
877 let mut recv_tasks = Vec::new();
878 let mut first_error = None;
879 let mut record_gc_error = |e, region_count| {
880 METRIC_META_GC_DATANODE_CALLS_TOTAL
881 .with_label_values(&["gc_regions", "error"])
882 .inc();
883 if region_count > 0 {
884 METRIC_META_GC_FAILED_REGIONS_TOTAL.inc_by(region_count);
885 }
886 if first_error.is_none() {
887 first_error = Some(e);
888 }
889 };
890 for (peer, regions, region_count, report) in join_all(tasks).await {
891 match report {
892 Ok(mailbox_rx) => {
893 recv_tasks.push(async move {
894 let report =
895 recv_gc_regions_reply(&peer, ®ions, "Batch GC", mailbox_rx).await;
896 (peer, region_count, report)
897 });
898 }
899 Err(e) => record_gc_error(e, region_count),
900 }
901 }
902
903 for (peer, region_count, report) in join_all(recv_tasks).await {
904 let report = match report {
905 Ok(report) => {
906 METRIC_META_GC_DATANODE_CALLS_TOTAL
907 .with_label_values(&["gc_regions", "success"])
908 .inc();
909 let need_retry_count = report.need_retry_regions.len() as u64;
910 if need_retry_count > 0 {
911 METRIC_META_GC_FAILED_REGIONS_TOTAL.inc_by(need_retry_count);
912 }
913 report
914 }
915 Err(e) => {
916 record_gc_error(e, region_count);
917 continue;
918 }
919 };
920
921 let success = report.deleted_files.keys().collect_vec();
922 let need_retry = report.need_retry_regions.iter().cloned().collect_vec();
923
924 if need_retry.is_empty() {
925 info!(
926 "GC report from datanode {}: successfully deleted files for region IDs {:?}",
927 peer, success
928 );
929 } else {
930 warn!(
931 "GC report from datanode {}: successfully deleted files for region IDs {:?}, need retry for region IDs {:?}",
932 peer, success, need_retry
933 );
934 }
935 all_need_retry.extend(report.need_retry_regions.clone());
936 all_report.merge(report);
937 }
938
939 self.merge_gc_report(all_report);
940
941 if let Some(e) = first_error {
942 return Err(e);
943 }
944
945 if !all_need_retry.is_empty() {
946 warn!("Regions need retry after batch GC: {:?}", all_need_retry);
947 }
948
949 Ok(())
950 }
951}
952
953#[async_trait::async_trait]
954impl Procedure for BatchGcProcedure {
955 fn type_name(&self) -> &str {
956 Self::TYPE_NAME
957 }
958
959 async fn execute(&mut self, ctx: &ProcedureContext) -> ProcedureResult<Status> {
960 self.check_maintenance_mode().await?;
961
962 match self.data.state {
963 State::Start => {
964 let _regions_span = common_telemetry::tracing::debug_span!(
965 "meta_gc_procedure_regions",
966 state = "start",
967 regions = ?self.data.regions
968 )
969 .entered();
970 info!(
971 "Batch GC procedure transitioning from Start to Acquiring for {} regions",
972 self.data.regions.len()
973 );
974 self.data.state = State::Acquiring;
976 Ok(Status::executing(false))
977 }
978 State::Acquiring => {
979 let region_count = self.data.regions.len();
980 let full_file_listing = self.data.full_file_listing;
981 let regions = self.data.regions.clone();
982 info!(
983 "Batch GC procedure acquiring file references for {} regions",
984 region_count
985 );
986 match self
988 .get_file_references()
989 .instrument(common_telemetry::tracing::debug_span!(
990 "meta_gc_procedure_regions",
991 state = "acquiring",
992 regions = ?regions
993 ))
994 .instrument(common_telemetry::tracing::info_span!(
995 "meta_gc_procedure_get_file_references",
996 region_count = region_count,
997 full_file_listing = full_file_listing
998 ))
999 .await
1000 {
1001 Ok(file_refs) => {
1002 info!(
1003 "Batch GC procedure acquired file references for {} regions",
1004 file_refs.file_refs.len()
1005 );
1006 self.data.file_refs = file_refs;
1007 self.data.state = State::Gcing;
1008 Ok(Status::executing(false))
1009 }
1010 Err(e) => {
1011 error!(e; "Failed to get file references");
1012 Err(ProcedureError::external(e))
1013 }
1014 }
1015 }
1016 State::Gcing => {
1017 info!(
1018 "Batch GC procedure sending GC instructions for {} regions",
1019 self.data.regions.len()
1020 );
1021 let debug_span = common_telemetry::tracing::debug_span!(
1024 "meta_gc_procedure_regions",
1025 state = "gcing",
1026 regions = ?self.data.regions
1027 );
1028 let info_span = common_telemetry::tracing::info_span!(
1029 "meta_gc_procedure_send_gc_instructions",
1030 region_count = self.data.regions.len(),
1031 full_file_listing = self.data.full_file_listing
1032 );
1033 match self
1034 .send_gc_instructions()
1035 .instrument(debug_span)
1036 .instrument(info_span)
1037 .await
1038 {
1039 Ok(()) => {
1040 info!(
1041 "Batch GC procedure received GC report, retry region count: {}",
1042 self.data
1043 .gc_report
1044 .as_ref()
1045 .map_or(0, |report| report.need_retry_regions.len())
1046 );
1047 self.data.state = State::UpdateRepartition;
1048 Ok(Status::executing(false))
1049 }
1050 Err(e) => {
1051 error!(e; "Failed to send GC instructions");
1052 Err(ProcedureError::external(e))
1053 }
1054 }
1055 }
1056 State::UpdateRepartition => match self
1057 .cleanup_region_repartition(ctx)
1058 .instrument(common_telemetry::tracing::debug_span!(
1059 "meta_gc_procedure_regions",
1060 state = "update_repartition",
1061 regions = ?self.data.regions
1062 ))
1063 .instrument(common_telemetry::tracing::info_span!(
1064 "meta_gc_procedure_update_repartition",
1065 region_count = self.data.regions.len()
1066 ))
1067 .await
1068 {
1069 Ok(()) => {
1070 debug!(
1071 "Cleanup region repartition info completed successfully for regions {:?}",
1072 self.data.regions
1073 );
1074 info!(
1075 "Batch GC completed successfully for regions {:?}",
1076 self.data.regions
1077 );
1078 info!("GC report: {:?}", self.data.gc_report);
1079 self.done_with_gc_report()
1080 }
1081 Err(e) => {
1082 error!(e; "Failed to cleanup region repartition info");
1083 Err(ProcedureError::external(e))
1084 }
1085 },
1086 }
1087 }
1088
1089 fn dump(&self) -> ProcedureResult<String> {
1090 serde_json::to_string(&self.data).context(ToJsonSnafu)
1091 }
1092
1093 fn lock_key(&self) -> LockKey {
1096 let lock_key: Vec<_> = self
1097 .data
1098 .regions
1099 .iter()
1100 .sorted() .map(|id| RegionLock::Read(*id).into())
1102 .collect();
1103
1104 LockKey::new(lock_key)
1105 }
1106
1107 fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
1108 if !ctx.event_type_filter.allows(BATCH_GC_EVENT_TYPE) {
1109 return None;
1110 }
1111
1112 let event = match &ctx.trigger {
1113 EventTrigger::Submitted => ctx
1115 .event_context
1116 .is_some_and(|context| context.reason == TriggerReason::Manual)
1117 .then(|| {
1118 BatchGcEvent::with_config(
1119 &self.data.regions,
1120 self.data.full_file_listing,
1121 self.data.timeout,
1122 )
1123 })?,
1124 EventTrigger::Recovered | EventTrigger::ChildSubmitted { .. } => return None,
1125 EventTrigger::Succeeded => {
1126 let ProcedureState::Done {
1127 output: Some(output),
1128 } = ctx.lifecycle_state
1129 else {
1130 return None;
1131 };
1132 let report = output.downcast_ref::<GcReport>()?;
1133 BatchGcEvent::with_report(report)?
1134 }
1135 EventTrigger::Retrying { .. } | EventTrigger::RollingBack => BatchGcEvent::with_config(
1136 &self.data.regions,
1137 self.data.full_file_listing,
1138 self.data.timeout,
1139 ),
1140 EventTrigger::Failed | EventTrigger::Poisoned => self
1141 .data
1142 .gc_report
1143 .as_ref()
1144 .and_then(BatchGcEvent::with_report)
1145 .unwrap_or_else(|| {
1146 BatchGcEvent::with_config(
1147 &self.data.regions,
1148 self.data.full_file_listing,
1149 self.data.timeout,
1150 )
1151 }),
1152 };
1153 Some(Box::new(event))
1154 }
1155}
1156
1157#[cfg(test)]
1158mod tests {
1159 use std::collections::HashMap;
1160
1161 use api::v1::meta::MailboxMessage;
1162 use api::v1::meta::mailbox_message::Payload;
1163 use common_meta::instruction::{GcRegionsReply, Instruction, InstructionReply};
1164 use common_meta::key::TableMetadataManager;
1165 use common_meta::key::runtime_switch::RuntimeSwitchManager;
1166 use common_meta::kv_backend::memory::MemoryKvBackend;
1167 use common_meta::kv_backend::test_util::MockKvBackendBuilder;
1168 use common_meta::peer::Peer;
1169 use common_meta::sequence::SequenceBuilder;
1170 use common_procedure::Context as ProcedureContext;
1171 use common_procedure_test::MockContextProvider;
1172 use common_time::util::current_time_millis;
1173 use store_api::storage::FileId;
1174 use tokio::sync::mpsc;
1175
1176 use super::*;
1177 use crate::procedure::test_util::{MailboxContext, send_mock_reply};
1178 use crate::service::mailbox::Channel;
1179
1180 #[test]
1181 fn test_scoped_gc_instruction_selects_and_scopes_manifest() {
1182 let region = RegionId::new(7, 3);
1183 let other = RegionId::new(7, 4);
1184 let mut manifest = FileRefsManifest::default();
1185 manifest.file_refs.insert(region, HashSet::new());
1186 manifest.file_refs.insert(other, HashSet::new());
1187 manifest.manifest_version.insert(region, 42);
1188
1189 let packed = scoped_gc_instruction(vec![region], &manifest, true);
1190 let Instruction::PackedGcRegions(packed) = packed else {
1191 panic!("expected packed instruction");
1192 };
1193 assert_eq!(packed.regions, vec![region]);
1194 assert_eq!(
1195 packed.packed_file_refs_manifest.manifest_version[®ion],
1196 42
1197 );
1198 assert!(
1199 !packed
1200 .packed_file_refs_manifest
1201 .file_refs
1202 .contains_key(&other)
1203 );
1204 assert!(
1205 packed
1206 .packed_file_refs_manifest
1207 .cross_region_refs
1208 .is_empty()
1209 );
1210 }
1211
1212 #[test]
1213 fn test_done_with_gc_report_keeps_report() {
1214 let region_id = RegionId::new(1024, 1);
1215 let file_id = FileId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
1216 let mut procedure = batch_gc_procedure();
1217 procedure.data.gc_report = Some(GcReport {
1218 deleted_files: HashMap::from([(region_id, vec![file_id])]),
1219 ..Default::default()
1220 });
1221
1222 for _ in 0..2 {
1223 let status = procedure.done_with_gc_report().unwrap();
1224 assert_eq!(
1225 status.downcast_output_ref::<GcReport>(),
1226 procedure.data.gc_report.as_ref()
1227 );
1228 }
1229 }
1230
1231 #[test]
1232 fn test_merge_gc_report_preserves_partial_outcomes() {
1233 let first_region = RegionId::new(1024, 1);
1234 let second_region = RegionId::new(1024, 2);
1235 let first_file = FileId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
1236 let second_file = FileId::parse_str("00000000-0000-0000-0000-000000000002").unwrap();
1237 let mut procedure = batch_gc_procedure();
1238
1239 procedure.merge_gc_report(GcReport {
1240 deleted_files: HashMap::from([(first_region, vec![first_file])]),
1241 ..Default::default()
1242 });
1243 procedure.merge_gc_report(GcReport {
1244 deleted_files: HashMap::from([
1245 (first_region, vec![first_file]),
1246 (second_region, vec![second_file]),
1247 ]),
1248 ..Default::default()
1249 });
1250
1251 let report = procedure.data.gc_report.unwrap();
1252 assert_eq!(report.deleted_files.len(), 2);
1253 assert_eq!(report.deleted_files[&first_region], vec![first_file]);
1254 assert_eq!(report.deleted_files[&second_region], vec![second_file]);
1255 }
1256
1257 #[test]
1258 fn test_merge_gc_report_uses_latest_region_outcome() {
1259 let region_id = RegionId::new(1024, 1);
1260 let mut procedure = batch_gc_procedure();
1261
1262 procedure.merge_gc_report(GcReport {
1263 deleted_files: HashMap::from([(region_id, vec![])]),
1264 processed_regions: HashSet::from([region_id]),
1265 ..Default::default()
1266 });
1267 procedure.merge_gc_report(GcReport {
1268 need_retry_regions: HashSet::from([region_id]),
1269 ..Default::default()
1270 });
1271
1272 let report = procedure.data.gc_report.unwrap();
1273 assert_eq!(report.deleted_files[®ion_id], Vec::<FileId>::new());
1274 assert!(!report.processed_regions.contains(®ion_id));
1275 assert!(report.need_retry_regions.contains(®ion_id));
1276 }
1277
1278 #[tokio::test]
1279 async fn test_send_gc_instructions_preserves_partial_report() {
1280 let first_region = RegionId::new(1024, 1);
1281 let second_region = RegionId::new(1024, 2);
1282 let first_peer = Peer::new(1, "first");
1283 let second_peer = Peer::new(2, "second");
1284 let file_id = FileId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
1285 let report = GcReport {
1286 deleted_files: HashMap::from([(first_region, vec![file_id])]),
1287 ..Default::default()
1288 };
1289
1290 let kv_backend = Arc::new(MemoryKvBackend::new());
1291 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1292 let runtime_switch_manager = Arc::new(RuntimeSwitchManager::new(kv_backend.clone()));
1293 let mailbox_sequence =
1294 SequenceBuilder::new("test_batch_gc_partial_report", kv_backend).build();
1295 let mut mailbox = MailboxContext::new(mailbox_sequence);
1296 let (tx, rx) = mpsc::channel(1);
1297 mailbox
1298 .insert_heartbeat_response_receiver(Channel::Datanode(first_peer.id), tx)
1299 .await;
1300 send_mock_reply(mailbox.mailbox().clone(), rx, {
1301 let report = report.clone();
1302 move |id| gc_reply(id, report.clone())
1303 });
1304
1305 let mut procedure = BatchGcProcedure::new(
1306 mailbox.mailbox().clone(),
1307 table_metadata_manager,
1308 runtime_switch_manager,
1309 "localhost".to_string(),
1310 vec![first_region, second_region],
1311 true,
1312 Duration::from_secs(10),
1313 HashMap::new(),
1314 );
1315 procedure.data.region_routes = HashMap::from([
1316 (first_region, (first_peer, vec![])),
1317 (second_region, (second_peer, vec![])),
1318 ]);
1319
1320 assert!(procedure.send_gc_instructions().await.is_err());
1321 assert_eq!(procedure.data.gc_report.as_ref(), Some(&report));
1322 }
1323
1324 fn gc_reply(id: u64, report: GcReport) -> Result<MailboxMessage> {
1325 Ok(MailboxMessage {
1326 id,
1327 subject: "mock".to_string(),
1328 from: "datanode".to_string(),
1329 to: "meta".to_string(),
1330 timestamp_millis: current_time_millis(),
1331 payload: Some(Payload::Json(
1332 serde_json::to_string(&InstructionReply::GcRegions(GcRegionsReply {
1333 result: Ok(report),
1334 }))
1335 .unwrap(),
1336 )),
1337 header: None,
1338 })
1339 }
1340
1341 fn batch_gc_procedure() -> BatchGcProcedure {
1342 let kv_backend = Arc::new(MemoryKvBackend::new());
1343 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1344 let runtime_switch_manager = Arc::new(RuntimeSwitchManager::new(kv_backend));
1345 let mailbox_sequence =
1346 SequenceBuilder::new("test_batch_gc_procedure", Arc::new(MemoryKvBackend::new()))
1347 .build();
1348 let mailbox = MailboxContext::new(mailbox_sequence);
1349 BatchGcProcedure::new(
1350 mailbox.mailbox().clone(),
1351 table_metadata_manager,
1352 runtime_switch_manager,
1353 "localhost".to_string(),
1354 vec![RegionId::new(1024, 1)],
1355 true,
1356 Duration::from_secs(10),
1357 HashMap::new(),
1358 )
1359 }
1360
1361 #[tokio::test]
1362 async fn test_maintenance_mode_gates_every_gc_state() {
1363 let states = [
1364 State::Start,
1365 State::Acquiring,
1366 State::Gcing,
1367 State::UpdateRepartition,
1368 ];
1369
1370 for state in states {
1371 let kv_backend = Arc::new(MemoryKvBackend::new());
1372 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1373 let runtime_switch_manager = Arc::new(RuntimeSwitchManager::new(kv_backend.clone()));
1374 runtime_switch_manager.set_maintenance_mode().await.unwrap();
1375 let mailbox_sequence =
1376 SequenceBuilder::new("test_batch_gc_maintenance_gate", kv_backend).build();
1377 let mut mailbox = MailboxContext::new(mailbox_sequence);
1378 let (tx, mut rx) = mpsc::channel(1);
1379 mailbox
1380 .insert_heartbeat_response_receiver(Channel::Datanode(1), tx)
1381 .await;
1382 let mut procedure = BatchGcProcedure::new(
1383 mailbox.mailbox().clone(),
1384 table_metadata_manager,
1385 runtime_switch_manager,
1386 "localhost".to_string(),
1387 vec![RegionId::new(1024, 1)],
1388 true,
1389 Duration::from_secs(10),
1390 HashMap::new(),
1391 );
1392 procedure.data.state = state.clone();
1393 let dump_before = procedure.dump().unwrap();
1394
1395 let ctx = ProcedureContext {
1396 procedure_id: common_procedure::ProcedureId::random(),
1397 provider: Arc::new(MockContextProvider::default()),
1398 event_context: None,
1399 };
1400 let err = procedure.execute(&ctx).await.unwrap_err();
1401
1402 assert!(!err.is_retry_later());
1403 assert_eq!(
1404 common_error::ext::ErrorExt::retry_hint(&err),
1405 common_error::ext::RetryHint::NonRetryable
1406 );
1407 assert_eq!(procedure.data.state, state);
1408 assert_eq!(procedure.dump().unwrap(), dump_before);
1409 if matches!(state, State::Acquiring | State::Gcing) {
1410 assert!(rx.try_recv().is_err());
1411 }
1412 }
1413 }
1414
1415 #[tokio::test]
1416 async fn test_maintenance_mode_read_error_retries_without_advancing_state() {
1417 let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1418 let failing_once_calls = calls.clone();
1419 let kv_backend = Arc::new(
1420 MockKvBackendBuilder::default()
1421 .range_fn(Arc::new(move |_| {
1422 if failing_once_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) == 0 {
1423 common_meta::error::UnexpectedSnafu {
1424 err_msg: "maintenance read failed",
1425 }
1426 .fail()
1427 } else {
1428 Ok(common_meta::rpc::store::RangeResponse {
1429 kvs: vec![],
1430 more: false,
1431 })
1432 }
1433 }))
1434 .build()
1435 .unwrap(),
1436 );
1437 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1438 let runtime_switch_manager = Arc::new(RuntimeSwitchManager::new(kv_backend.clone()));
1439 let mailbox_sequence =
1440 SequenceBuilder::new("test_batch_gc_maintenance_read_error", kv_backend).build();
1441 let mailbox = MailboxContext::new(mailbox_sequence);
1442 let mut procedure = BatchGcProcedure::new(
1443 mailbox.mailbox().clone(),
1444 table_metadata_manager,
1445 runtime_switch_manager,
1446 "localhost".to_string(),
1447 vec![RegionId::new(1024, 1)],
1448 true,
1449 Duration::from_secs(10),
1450 HashMap::new(),
1451 );
1452 let ctx = ProcedureContext {
1453 procedure_id: common_procedure::ProcedureId::random(),
1454 provider: Arc::new(MockContextProvider::default()),
1455 event_context: None,
1456 };
1457 let dump_before = procedure.dump().unwrap();
1458
1459 assert!(procedure.execute(&ctx).await.unwrap_err().is_retry_later());
1460 assert_eq!(procedure.data.state, State::Start);
1461 assert_eq!(procedure.dump().unwrap(), dump_before);
1462
1463 assert!(matches!(
1464 procedure.execute(&ctx).await.unwrap(),
1465 Status::Executing { .. }
1466 ));
1467 assert_eq!(procedure.data.state, State::Acquiring);
1468 }
1469}