1use 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 pub(crate) table_metadata_manager: TableMetadataManagerRef,
157 pub(crate) procedure_manager: ProcedureManagerRef,
159 pub(crate) runtime_switch_manager: RuntimeSwitchManagerRef,
161 #[cfg(feature = "enterprise")]
163 pub(crate) ddl_manager: DdlManagerRef,
164 #[cfg(feature = "enterprise")]
167 in_flight_purges: Arc<Mutex<HashSet<TableId>>>,
168 #[cfg(feature = "enterprise")]
169 purge_scan_cursor: AtomicUsize,
170 pub(crate) meta_peer_client: MetaPeerClientRef,
172 pub(crate) mailbox: MailboxRef,
174 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}