Skip to main content

meta_srv/gc/
procedure.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
217/// Procedure to perform get file refs then batch GC for multiple regions,
218/// it holds locks for all regions during the whole procedure.
219pub 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    /// Meta server address
230    server_addr: String,
231    /// The regions to be GC-ed
232    regions: Vec<RegionId>,
233    full_file_listing: bool,
234    region_routes: Region2Peers,
235    /// Routes assigned by the scheduler for regions missing from table routes.
236    #[serde(default)]
237    region_routes_override: Region2Peers,
238    /// Related regions (e.g., for shared files after repartition).
239    /// The source regions (where those files originally came from) are used as the key, and the destination regions (where files are currently stored) are used as the value.
240    related_regions: HashMap<RegionId, HashSet<RegionId>>,
241    /// Acquired file references (Populated in Acquiring state)
242    file_refs: FileRefsManifest,
243    /// mailbox timeout duration
244    timeout: Duration,
245    gc_report: Option<GcReport>,
246}
247
248#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
249pub enum State {
250    /// Initial state
251    Start,
252    /// Fetching file references from datanodes
253    Acquiring,
254    /// Sending GC instruction to the target datanode
255    Gcing,
256    /// Updating region repartition info in kvbackend after GC based on the GC result
257    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    /// Test-only constructor to jump directly into the repartition update state.
294    /// Intended for integration tests that validate `cleanup_region_repartition` without
295    /// running the full batch GC state machine.
296    #[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        // Deleted objects are cumulative, while these sets describe the latest outcome
356        // for each region covered by this report.
357        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    /// Return related regions for the given regions.
405    /// The returned map uses the input region as key, and all other regions
406    /// from the same table as values (excluding the input region itself).
407    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            // batch_get_physical_table_routes returns a subset on misses; treat that as error
423            for table_id in &table_ids {
424                if !table_routes.contains_key(table_id) {
425                    // indicate is a logical table id
426                    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    /// Clean up region repartition info in kvbackend after GC
464    /// according to cross reference in `FileRefsManifest`.
465    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        // Regions whose extension sidecar cleanup failed and need retry. Keep
492        // their tombstone so the next GC cycle can replay the cleanup.
493        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            // We only touch regions involved in this GC batch for the current table to avoid
528            // clobbering unrelated repart entries. Start from the batch regions of this table.
529            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            // Merge targets: only the batch regions of this table. This avoids touching unrelated
538            // repart entries; we just reconcile mappings for regions involved in the current GC
539            // cycle for this table.
540            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                    // Keep the tombstone: tmp refs or pending extension cleanup
552                    // still need a future GC pass.
553                    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 there is no repartition info to persist, skip creating/updating the key
560            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    /// Discover region routes for the given regions.
574    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        // Group regions by table ID for batch processing
582        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        // Process each table's regions together for efficiency
592        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                    // Continue with other tables instead of failing completely
604                    // TODO(discord9): consider failing here instead
605                    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    /// Set region routes and related regions for GC procedure
618    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        // Discover routes for all regions involved in GC, including both the
624        // primary GC regions and their related regions.
625        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(&regions_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    /// Get file references from all datanodes that host the regions
648    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        // Group regions by datanode to minimize RPC calls
668        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                // also need to send for follower regions for file refs in case query is running on follower
680                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                } // since read from manifest, no need to send to followers
708            }
709        }
710
711        // Send packed GetFileRefs instructions to each datanode
712        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        // Each future owns both send and receive, so a completed reply is merged and
723        // dropped immediately rather than retained in a second join_all buffer.
724        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    /// Sends GC instructions to all datanodes that host the regions.
824    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        // Group regions by datanode
831        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, &regions, "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                // Transition to Acquiring state
975                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                // Get file references from all datanodes
987                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                // Send GC instructions to all datanodes
1022                // TODO(discord9): handle need-retry regions
1023                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    /// Read lock all regions involved in this GC procedure.
1094    /// So i.e. region migration won't happen during GC and cause race conditions.
1095    fn lock_key(&self) -> LockKey {
1096        let lock_key: Vec<_> = self
1097            .data
1098            .regions
1099            .iter()
1100            .sorted() // sort to have a deterministic lock order
1101            .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            // Keep scheduled GC low-noise; record submitted manual requests for auditability.
1114            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[&region],
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[&region_id], Vec::<FileId>::new());
1274        assert!(!report.processed_regions.contains(&region_id));
1275        assert!(report.need_retry_regions.contains(&region_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}