Skip to main content

meta_srv/gc/
ctx.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::collections::HashMap;
16#[cfg(feature = "enterprise")]
17use std::collections::HashSet;
18#[cfg(feature = "enterprise")]
19use std::sync::atomic::{AtomicUsize, Ordering};
20#[cfg(feature = "enterprise")]
21use std::sync::{Arc, Mutex};
22use std::time::Duration;
23
24use common_meta::datanode::RegionStat;
25#[cfg(feature = "enterprise")]
26use common_meta::ddl_manager::DdlManagerRef;
27#[cfg(feature = "enterprise")]
28use common_meta::key::DroppedTableName;
29use common_meta::key::TableMetadataManagerRef;
30use common_meta::key::runtime_switch::RuntimeSwitchManagerRef;
31use common_meta::key::table_repart::TableRepartValue;
32use common_meta::key::table_route::PhysicalTableRouteValue;
33#[cfg(feature = "enterprise")]
34use common_meta::rpc::ddl::PurgeDroppedTableTask;
35use common_procedure::{ProcedureContext, ProcedureManagerRef, ProcedureWithId, watcher};
36use common_telemetry::debug;
37use snafu::{OptionExt as _, ResultExt as _};
38use store_api::storage::{GcReport, RegionId};
39use table::metadata::TableId;
40
41use crate::cluster::MetaPeerClientRef;
42use crate::error::{self, Result, TableMetadataManagerSnafu};
43use crate::gc::Region2Peers;
44use crate::gc::procedure::BatchGcProcedure;
45#[cfg(feature = "enterprise")]
46use crate::metrics::METRIC_META_GC_SOFT_DROP_PURGES_TOTAL;
47use crate::service::mailbox::MailboxRef;
48
49#[async_trait::async_trait]
50pub(crate) trait SchedulerCtx: Send + Sync {
51    async fn get_table_to_region_stats(&self) -> Result<HashMap<TableId, Vec<RegionStat>>>;
52
53    async fn get_table_reparts(&self) -> Result<Vec<(TableId, TableRepartValue)>>;
54
55    async fn get_table_route(
56        &self,
57        table_id: TableId,
58    ) -> Result<(TableId, PhysicalTableRouteValue)>;
59
60    async fn batch_get_table_route(
61        &self,
62        table_ids: &[TableId],
63    ) -> Result<HashMap<TableId, PhysicalTableRouteValue>>;
64
65    async fn gc_regions(
66        &self,
67        region_ids: &[RegionId],
68        full_file_listing: bool,
69        timeout: Duration,
70        region_routes_override: Region2Peers,
71        procedure_context: ProcedureContext,
72    ) -> Result<GcReport>;
73
74    #[cfg(feature = "enterprise")]
75    async fn list_dropped_tables(&self) -> Result<Vec<DroppedTableName>>;
76
77    #[cfg(feature = "enterprise")]
78    async fn purge_dropped_table(&self, table_id: TableId) -> Result<()>;
79
80    #[cfg(feature = "enterprise")]
81    fn try_reserve_purge(
82        &self,
83        table_id: TableId,
84        max_in_flight: usize,
85    ) -> Option<PurgeReservation>;
86
87    #[cfg(feature = "enterprise")]
88    fn next_purge_scan_start(&self, _table_count: usize) -> usize {
89        0
90    }
91}
92
93#[cfg(feature = "enterprise")]
94pub(crate) enum PurgeOutcome {
95    Succeeded,
96    Failed,
97}
98
99#[cfg(feature = "enterprise")]
100pub(crate) struct PurgeReservation {
101    table_id: TableId,
102    in_flight: Arc<Mutex<HashSet<TableId>>>,
103    outcome_recorded: bool,
104}
105
106#[cfg(feature = "enterprise")]
107impl PurgeReservation {
108    pub(crate) fn try_new(
109        in_flight: Arc<Mutex<HashSet<TableId>>>,
110        table_id: TableId,
111        max_in_flight: usize,
112    ) -> Option<Self> {
113        let mut tables = in_flight
114            .lock()
115            .unwrap_or_else(|poisoned| poisoned.into_inner());
116        if tables.len() >= max_in_flight || !tables.insert(table_id) {
117            return None;
118        }
119        drop(tables);
120        Some(Self {
121            table_id,
122            in_flight,
123            outcome_recorded: false,
124        })
125    }
126
127    pub(crate) fn record_outcome(mut self, outcome: PurgeOutcome) {
128        let status = match outcome {
129            PurgeOutcome::Succeeded => "succeeded",
130            PurgeOutcome::Failed => "failed",
131        };
132        METRIC_META_GC_SOFT_DROP_PURGES_TOTAL
133            .with_label_values(&[status])
134            .inc();
135        self.outcome_recorded = true;
136    }
137}
138
139#[cfg(feature = "enterprise")]
140impl Drop for PurgeReservation {
141    fn drop(&mut self) {
142        if !self.outcome_recorded {
143            METRIC_META_GC_SOFT_DROP_PURGES_TOTAL
144                .with_label_values(&["cancelled"])
145                .inc();
146        }
147        self.in_flight
148            .lock()
149            .unwrap_or_else(|poisoned| poisoned.into_inner())
150            .remove(&self.table_id);
151    }
152}
153
154pub(crate) struct DefaultGcSchedulerCtx {
155    /// The metadata manager.
156    pub(crate) table_metadata_manager: TableMetadataManagerRef,
157    /// Procedure manager.
158    pub(crate) procedure_manager: ProcedureManagerRef,
159    /// Runtime switch manager used by recovered and newly submitted GC procedures.
160    pub(crate) runtime_switch_manager: RuntimeSwitchManagerRef,
161    /// DDL manager used to submit the existing purge procedure.
162    #[cfg(feature = "enterprise")]
163    pub(crate) ddl_manager: DdlManagerRef,
164    /// Process-local reservations for purge procedures submitted by this scheduler.
165    /// Procedure recovery after a metasrv restart may outlive this set.
166    #[cfg(feature = "enterprise")]
167    in_flight_purges: Arc<Mutex<HashSet<TableId>>>,
168    #[cfg(feature = "enterprise")]
169    purge_scan_cursor: AtomicUsize,
170    /// For getting `RegionStats`.
171    pub(crate) meta_peer_client: MetaPeerClientRef,
172    /// The mailbox to send messages.
173    pub(crate) mailbox: MailboxRef,
174    /// The server address.
175    pub(crate) server_addr: String,
176}
177
178impl DefaultGcSchedulerCtx {
179    pub fn try_new(
180        table_metadata_manager: TableMetadataManagerRef,
181        procedure_manager: ProcedureManagerRef,
182        runtime_switch_manager: RuntimeSwitchManagerRef,
183        #[cfg(feature = "enterprise")] ddl_manager: DdlManagerRef,
184        meta_peer_client: MetaPeerClientRef,
185        mailbox: MailboxRef,
186        server_addr: String,
187    ) -> Result<Self> {
188        Ok(Self {
189            table_metadata_manager,
190            procedure_manager,
191            runtime_switch_manager,
192            #[cfg(feature = "enterprise")]
193            ddl_manager,
194            #[cfg(feature = "enterprise")]
195            in_flight_purges: Arc::new(Mutex::new(HashSet::new())),
196            #[cfg(feature = "enterprise")]
197            purge_scan_cursor: AtomicUsize::new(0),
198            meta_peer_client,
199            mailbox,
200            server_addr,
201        })
202    }
203}
204
205#[async_trait::async_trait]
206impl SchedulerCtx for DefaultGcSchedulerCtx {
207    async fn get_table_to_region_stats(&self) -> Result<HashMap<TableId, Vec<RegionStat>>> {
208        let dn_stats = self.meta_peer_client.get_all_dn_stat_kvs().await?;
209        let mut table_to_region_stats: HashMap<TableId, Vec<RegionStat>> = HashMap::new();
210        for (_dn_id, stats) in dn_stats {
211            let stats = stats.stats;
212
213            let Some(latest_stat) = stats.iter().max_by_key(|s| s.timestamp_millis).cloned() else {
214                continue;
215            };
216
217            for region_stat in latest_stat.region_stats {
218                table_to_region_stats
219                    .entry(region_stat.id.table_id())
220                    .or_default()
221                    .push(region_stat);
222            }
223        }
224        Ok(table_to_region_stats)
225    }
226
227    async fn get_table_reparts(&self) -> Result<Vec<(TableId, TableRepartValue)>> {
228        self.table_metadata_manager
229            .table_repart_manager()
230            .table_reparts()
231            .await
232            .context(TableMetadataManagerSnafu)
233    }
234
235    async fn get_table_route(
236        &self,
237        table_id: TableId,
238    ) -> Result<(TableId, PhysicalTableRouteValue)> {
239        self.table_metadata_manager
240            .table_route_manager()
241            .get_physical_table_route(table_id)
242            .await
243            .context(TableMetadataManagerSnafu)
244    }
245
246    async fn batch_get_table_route(
247        &self,
248        table_ids: &[TableId],
249    ) -> Result<HashMap<TableId, PhysicalTableRouteValue>> {
250        self.table_metadata_manager
251            .table_route_manager()
252            .batch_get_physical_table_routes(table_ids)
253            .await
254            .context(TableMetadataManagerSnafu)
255    }
256
257    async fn gc_regions(
258        &self,
259        region_ids: &[RegionId],
260        full_file_listing: bool,
261        timeout: Duration,
262        region_routes_override: Region2Peers,
263        procedure_context: ProcedureContext,
264    ) -> Result<GcReport> {
265        self.gc_regions_inner(
266            region_ids,
267            full_file_listing,
268            timeout,
269            region_routes_override,
270            procedure_context,
271        )
272        .await
273    }
274
275    #[cfg(feature = "enterprise")]
276    async fn list_dropped_tables(&self) -> Result<Vec<DroppedTableName>> {
277        self.table_metadata_manager
278            .list_dropped_tables()
279            .await
280            .context(TableMetadataManagerSnafu)
281    }
282
283    #[cfg(feature = "enterprise")]
284    async fn purge_dropped_table(&self, table_id: TableId) -> Result<()> {
285        self.ddl_manager
286            .submit_expired_purge_dropped_table_task(PurgeDroppedTableTask { table_id })
287            .await
288            .context(error::SubmitDdlTaskSnafu)?;
289        Ok(())
290    }
291
292    #[cfg(feature = "enterprise")]
293    fn try_reserve_purge(
294        &self,
295        table_id: TableId,
296        max_in_flight: usize,
297    ) -> Option<PurgeReservation> {
298        PurgeReservation::try_new(self.in_flight_purges.clone(), table_id, max_in_flight)
299    }
300
301    #[cfg(feature = "enterprise")]
302    fn next_purge_scan_start(&self, table_count: usize) -> usize {
303        if table_count == 0 {
304            return 0;
305        }
306        self.purge_scan_cursor.fetch_add(1, Ordering::Relaxed) % table_count
307    }
308}
309
310impl DefaultGcSchedulerCtx {
311    async fn gc_regions_inner(
312        &self,
313        region_ids: &[RegionId],
314        full_file_listing: bool,
315        timeout: Duration,
316        region_routes_override: Region2Peers,
317        procedure_context: ProcedureContext,
318    ) -> Result<GcReport> {
319        debug!(
320            "Sending GC instruction for {} regions (full_file_listing: {})",
321            region_ids.len(),
322            full_file_listing
323        );
324
325        let procedure = BatchGcProcedure::new(
326            self.mailbox.clone(),
327            self.table_metadata_manager.clone(),
328            self.runtime_switch_manager.clone(),
329            self.server_addr.clone(),
330            region_ids.to_vec(),
331            full_file_listing,
332            timeout,
333            region_routes_override,
334        );
335        let procedure_with_id =
336            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
337
338        let id = procedure_with_id.id;
339
340        let mut watcher = self
341            .procedure_manager
342            .submit(procedure_with_id)
343            .await
344            .context(error::SubmitProcedureSnafu)?;
345        let res = watcher::wait(&mut watcher)
346            .await
347            .context(error::WaitProcedureSnafu)?
348            .with_context(|| error::UnexpectedSnafu {
349                violated: format!(
350                    "GC procedure {id} successfully completed but no result returned"
351                ),
352            })?;
353
354        let gc_report = BatchGcProcedure::cast_result(res)?;
355
356        Ok(gc_report)
357    }
358}
359
360#[cfg(all(test, feature = "enterprise"))]
361mod tests {
362    use std::panic::{AssertUnwindSafe, catch_unwind};
363
364    use super::*;
365    use crate::metrics::METRIC_META_GC_SOFT_DROP_PURGES_TOTAL;
366
367    #[test]
368    fn test_purge_reservation_releases_slot_on_panic() {
369        let in_flight = Arc::new(std::sync::Mutex::new(HashSet::new()));
370        let cancelled = METRIC_META_GC_SOFT_DROP_PURGES_TOTAL.with_label_values(&["cancelled"]);
371        let before = cancelled.get();
372        let reservation = PurgeReservation::try_new(in_flight.clone(), 1, 1).unwrap();
373
374        let result = catch_unwind(AssertUnwindSafe(|| {
375            let _reservation = reservation;
376            panic!("mock purge panic");
377        }));
378
379        assert!(result.is_err());
380        assert!(PurgeReservation::try_new(in_flight, 2, 1).is_some());
381        assert!(cancelled.get() > before);
382    }
383
384    #[tokio::test]
385    async fn test_purge_reservation_releases_slot_on_task_abort() {
386        let in_flight = Arc::new(std::sync::Mutex::new(HashSet::new()));
387        let cancelled = METRIC_META_GC_SOFT_DROP_PURGES_TOTAL.with_label_values(&["cancelled"]);
388        let before = cancelled.get();
389        let reservation = PurgeReservation::try_new(in_flight.clone(), 1, 1).unwrap();
390        let handle = tokio::spawn(async move {
391            let _reservation = reservation;
392            std::future::pending::<()>().await;
393        });
394        tokio::task::yield_now().await;
395
396        handle.abort();
397        let error = handle.await.unwrap_err();
398
399        assert!(error.is_cancelled());
400        assert!(PurgeReservation::try_new(in_flight, 2, 1).is_some());
401        assert!(cancelled.get() > before);
402    }
403}