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_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
175/// Procedure to perform get file refs then batch GC for multiple regions,
176/// it holds locks for all regions during the whole procedure.
177pub 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    /// Meta server address
187    server_addr: String,
188    /// The regions to be GC-ed
189    regions: Vec<RegionId>,
190    full_file_listing: bool,
191    region_routes: Region2Peers,
192    /// Routes assigned by the scheduler for regions missing from table routes.
193    #[serde(default)]
194    region_routes_override: Region2Peers,
195    /// Related regions (e.g., for shared files after repartition).
196    /// 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.
197    related_regions: HashMap<RegionId, HashSet<RegionId>>,
198    /// Acquired file references (Populated in Acquiring state)
199    file_refs: FileRefsManifest,
200    /// mailbox timeout duration
201    timeout: Duration,
202    gc_report: Option<GcReport>,
203}
204
205#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
206pub enum State {
207    /// Initial state
208    Start,
209    /// Fetching file references from datanodes
210    Acquiring,
211    /// Sending GC instruction to the target datanode
212    Gcing,
213    /// Updating region repartition info in kvbackend after GC based on the GC result
214    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    /// Test-only constructor to jump directly into the repartition update state.
248    /// Intended for integration tests that validate `cleanup_region_repartition` without
249    /// running the full batch GC state machine.
250    #[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        // Deleted objects are cumulative, while these sets describe the latest outcome
292        // for each region covered by this report.
293        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    /// Return related regions for the given regions.
341    /// The returned map uses the input region as key, and all other regions
342    /// from the same table as values (excluding the input region itself).
343    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            // batch_get_physical_table_routes returns a subset on misses; treat that as error
359            for table_id in &table_ids {
360                if !table_routes.contains_key(table_id) {
361                    // indicate is a logical table id
362                    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    /// Clean up region repartition info in kvbackend after GC
400    /// according to cross reference in `FileRefsManifest`.
401    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        // Regions whose extension sidecar cleanup failed and need retry. Keep
428        // their tombstone so the next GC cycle can replay the cleanup.
429        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            // We only touch regions involved in this GC batch for the current table to avoid
464            // clobbering unrelated repart entries. Start from the batch regions of this table.
465            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            // Merge targets: only the batch regions of this table. This avoids touching unrelated
474            // repart entries; we just reconcile mappings for regions involved in the current GC
475            // cycle for this table.
476            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                    // Keep the tombstone: tmp refs or pending extension cleanup
488                    // still need a future GC pass.
489                    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 there is no repartition info to persist, skip creating/updating the key
496            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    /// Discover region routes for the given regions.
510    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        // Group regions by table ID for batch processing
518        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        // Process each table's regions together for efficiency
528        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                    // Continue with other tables instead of failing completely
540                    // TODO(discord9): consider failing here instead
541                    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    /// Set region routes and related regions for GC procedure
554    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        // Discover routes for all regions involved in GC, including both the
560        // primary GC regions and their related regions.
561        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(&regions_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    /// Get file references from all datanodes that host the regions
584    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        // Group regions by datanode to minimize RPC calls
604        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                // also need to send for follower regions for file refs in case query is running on follower
616                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                } // since read from manifest, no need to send to followers
644            }
645        }
646
647        // Send GetFileRefs instructions to each datanode
648        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        // store error to make sure metrics doesn't ignore other peers
685        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            // Merge the file references from this datanode
750            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            // region manifest version should be the smallest one among all peers, so outdated region can be detected
758            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    /// Sends GC instructions to all datanodes that host the regions.
783    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        // Group regions by datanode
790        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 could be somewhere large. But still intentionally clone per datanode here:
819                    // this path is admin-triggered or scheduler-triggered, peer count is expected to be bounded, and
820                    // and abnormal manifest growth should be addressed at the source
821                    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                // Transition to Acquiring state
937                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                // Get file references from all datanodes
949                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                // Send GC instructions to all datanodes
984                // TODO(discord9): handle need-retry regions
985                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    /// Read lock all regions involved in this GC procedure.
1056    /// So i.e. region migration won't happen during GC and cause race conditions.
1057    fn lock_key(&self) -> LockKey {
1058        let lock_key: Vec<_> = self
1059            .data
1060            .regions
1061            .iter()
1062            .sorted() // sort to have a deterministic lock order
1063            .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            // Keep scheduled GC low-noise; record submitted manual requests for auditability.
1076            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[&region_id], Vec::<FileId>::new());
1200        assert!(!report.processed_regions.contains(&region_id));
1201        assert!(report.need_retry_regions.contains(&region_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}