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