From d9b58edbb0c78c9ddc9200fa49eee72a2fb56bb1 Mon Sep 17 00:00:00 2001 From: Ning Sun Date: Wed, 22 Jul 2026 11:27:04 +0800 Subject: [PATCH] feat: add a region hook for gc cleanup (#8547) * feat: add a region hook for gc cleanup * chore: more tests and review fix * fix: address review comments * fix: address review comments --- .../src/heartbeat/handler/gc_worker.rs | 1 + src/meta-srv/src/gc/procedure.rs | 15 +- src/mito2/src/engine.rs | 15 + src/mito2/src/engine/region_hook.rs | 61 ++- src/mito2/src/gc.rs | 96 +++- src/mito2/src/gc/worker_test.rs | 486 ++++++++++++++++++ 6 files changed, 657 insertions(+), 17 deletions(-) diff --git a/src/datanode/src/heartbeat/handler/gc_worker.rs b/src/datanode/src/heartbeat/handler/gc_worker.rs index 9489f1461d..b32a1c29ad 100644 --- a/src/datanode/src/heartbeat/handler/gc_worker.rs +++ b/src/datanode/src/heartbeat/handler/gc_worker.rs @@ -169,6 +169,7 @@ impl GcRegionsHandler { file_ref_manifest.clone(), &mito_engine.gc_limiter(), full_file_listing, + mito_engine.region_hook(), ) .await .context(GcMitoEngineSnafu { diff --git a/src/meta-srv/src/gc/procedure.rs b/src/meta-srv/src/gc/procedure.rs index 05ad098749..c1841d2c77 100644 --- a/src/meta-srv/src/gc/procedure.rs +++ b/src/meta-srv/src/gc/procedure.rs @@ -382,6 +382,15 @@ impl BatchGcProcedure { let repart_mgr = self.table_metadata_manager.table_repart_manager(); + // Regions whose extension sidecar cleanup failed and need retry. Keep + // their tombstone so the next GC cycle can replay the cleanup. + let need_retry: HashSet = self + .data + .gc_report + .as_ref() + .map(|r| r.need_retry_regions.clone()) + .unwrap_or_default(); + let mut table_ids: HashSet = cross_refs_grouped .keys() .copied() @@ -432,9 +441,9 @@ impl BatchGcProcedure { let mut set = BTreeSet::new(); set.extend(dst_regions.iter().copied()); new_value.src_to_dst.insert(src_region, set); - } else if has_tmp_ref { - // Keep a tombstone entry with an empty set so dropped regions that still - // have tmp refs are preserved; removing it would lose the repartition trace. + } else if has_tmp_ref || need_retry.contains(&src_region) { + // Keep the tombstone: tmp refs or pending extension cleanup + // still need a future GC pass. new_value.src_to_dst.insert(src_region, BTreeSet::new()); } else { new_value.src_to_dst.remove(&src_region); diff --git a/src/mito2/src/engine.rs b/src/mito2/src/engine.rs index 94911814b5..a128531b53 100644 --- a/src/mito2/src/engine.rs +++ b/src/mito2/src/engine.rs @@ -105,6 +105,7 @@ use common_wal::options::WalOptions; use futures::future::{join_all, try_join_all}; use futures::stream::{self, Stream, StreamExt}; use object_store::manager::ObjectStoreManagerRef; +use region_hook::RegionHookRef; use snafu::{OptionExt, ResultExt, ensure}; use store_api::ManifestVersion; use store_api::codec::PrimaryKeyEncoding; @@ -219,6 +220,9 @@ impl<'a, S: LogStore> MitoEngineBuilder<'a, S> { self.config.sanitize(self.data_home)?; let config = Arc::new(self.config); + // Extract the region hook before `plugins` is moved into the WorkerGroup, + // so the engine (and thus the GC worker) can fire `on_region_gc`. + let region_hook = self.plugins.get::(); let workers = WorkerGroup::start( config.clone(), self.log_store.clone(), @@ -250,6 +254,7 @@ impl<'a, S: LogStore> MitoEngineBuilder<'a, S> { config, wal_raw_entry_reader, scan_memory_tracker, + region_hook, #[cfg(feature = "enterprise")] extension_range_provider_factory: None, }; @@ -328,6 +333,12 @@ impl MitoEngine { self.inner.workers.schema_metadata_manager() } + /// Returns the registered region hook (if any), for the GC worker to fire + /// [`RegionHook::on_region_gc`]. + pub fn region_hook(&self) -> Option { + self.inner.region_hook.clone() + } + /// Get all tmp ref files for given region ids, excluding files that's already in manifest. pub async fn get_snapshot_of_file_refs( &self, @@ -728,6 +739,9 @@ struct EngineInner { wal_raw_entry_reader: Arc, /// Memory tracker for table scans. scan_memory_tracker: QueryMemoryTracker, + /// The region hook (if any) registered via plugins; exposed for the GC worker + /// to fire [`RegionHook::on_region_gc`]. + region_hook: Option, #[cfg(feature = "enterprise")] extension_range_provider_factory: Option, } @@ -1458,6 +1472,7 @@ impl MitoEngine { config, wal_raw_entry_reader, scan_memory_tracker, + region_hook: None, #[cfg(feature = "enterprise")] extension_range_provider_factory: None, }), diff --git a/src/mito2/src/engine/region_hook.rs b/src/mito2/src/engine/region_hook.rs index 9c5a7dfc95..4d816463b7 100644 --- a/src/mito2/src/engine/region_hook.rs +++ b/src/mito2/src/engine/region_hook.rs @@ -89,6 +89,7 @@ //! | Close | [`on_region_closed`] | A close request (or a close-after-flush) removes the region from the active set. Data files, manifest and WAL state are **preserved**; the region may be reopened. | //! | Logical drop | [`on_region_dropped`] | A drop request has been handled: the region leaves the active set and its WAL entries are marked obsolete. Data files are **not yet deleted**. | //! | Physical file removal | [`on_region_files_removed`] | The drop GC worker has deleted the region directory. Terminal file-lifecycle event. | +//! | Global GC pass | [`on_region_gc`] | The datanode's global GC worker finished a GC pass for a region — both periodic GC for live regions and the global reclamation of dropped/repartitioned regions (`is_region_dropped`). | //! //! Notes: //! - `on_region_closed` / `on_region_dropped` run **inline in the region worker loop**, @@ -128,8 +129,9 @@ //! //! `on_region_files_removed` currently covers only the **drop** GC worker's physical //! directory removal. A broader per-file `on_files_removed` hook covering compaction -//! removal, truncate, and the global GC reclamation path is not yet implemented -//! (though logical file removal is already observable via `on_manifest_updated`). +//! removal and truncate is not yet implemented +//! (though logical file removal is already observable via `on_manifest_updated`, +//! and the global GC reclamation path is covered by `on_region_gc`). //! Role/leadership transitions (`on_region_role_changed`) are also not hooked. //! //! [`on_sst_files_written`]: RegionHook::on_sst_files_written @@ -138,6 +140,7 @@ //! [`on_region_closed`]: RegionHook::on_region_closed //! [`on_region_dropped`]: RegionHook::on_region_dropped //! [`on_region_files_removed`]: RegionHook::on_region_files_removed +//! [`on_region_gc`]: RegionHook::on_region_gc //! [`RegionManifestManager::update`]: crate::manifest::manager::RegionManifestManager::update //! [`ManifestContext::update_locked`]: crate::region::ManifestContext::update_locked //! [`ManifestContext::update_manifest`]: crate::region::ManifestContext::update_manifest @@ -150,7 +153,9 @@ use store_api::ManifestVersion; use store_api::metadata::RegionMetadataRef; use store_api::storage::RegionId; -use crate::manifest::action::RegionMetaActionList; +use crate::access_layer::AccessLayerRef; +use crate::error::Result; +use crate::manifest::action::{RegionMetaActionList, RemovedFile}; use crate::sst::file::FileMeta; use crate::sst::parquet::SstInfo; @@ -382,6 +387,56 @@ pub trait RegionHook: Send + Sync + Debug { ) { let _ = (region_id, region_metadata); } + + /// Called after the datanode's global GC worker (`LocalGcWorker`) finishes a + /// GC pass for a region — live regions (periodic GC) or dropped/repartitioned + /// regions (the global reclamation path). Lets extensions with sidecar files + /// outside mito2's region dir clean up residual files. + /// + /// Always scoped to [`RegionGcInfo::removed_files`]: clean only sidecar + /// artifacts for the files mito deleted this pass. On a + /// [`RegionGcInfo::full_file_listing`] pass you may also reconcile + /// sidecar-only orphans (e.g. detect a fully-reaped region by cross-checking + /// the region dir against your own manifest). This callback never authorizes + /// blind whole-directory removal — derive "fully reaped" from your own state + /// so a stale mito snapshot can't cause accidental deletion. + /// [`RegionGcInfo::is_region_dropped`] is context, not authorization. + /// + /// # Retry + /// + /// Returning `Err` keeps the region un-acknowledged (`need_retry_regions`) + /// for a future replay. **Dropped/repartitioned**: guaranteed — metasrv keeps + /// the `table_repart` tombstone until `Ok`, so full-listing passes continue + /// until cleanup finishes. **Live**: best-effort — not expedited through + /// candidate selection, and a later full-listing pass may not reconstruct the + /// same `removed_files`; treat live cleanup as opportunistic. + /// + /// `region_metadata` is `None` for dropped regions; use `region_id` + + /// `access_layer`. Idempotent; runs on the background GC task. + async fn on_region_gc( + &self, + region_id: RegionId, + region_metadata: Option<&RegionMetadataRef>, + access_layer: &AccessLayerRef, + info: &RegionGcInfo<'_>, + ) -> Result<()> { + let _ = (region_id, region_metadata, access_layer, info); + Ok(()) + } +} + +/// What mito2's GC pass deleted for a region, handed to +/// [`RegionHook::on_region_gc`]. +pub struct RegionGcInfo<'a> { + /// Files mito2 physically deleted this pass. Cleanup must be scoped to these + /// (plus sidecar-only orphans you can identify on a [`Self::full_file_listing`] + /// pass). + pub removed_files: &'a [RemovedFile], + /// `true` when the region is dropped/absent (e.g. a repartitioned source). + /// Context only — not authorization. `region_metadata` is `None` when `true`. + pub is_region_dropped: bool, + /// Whether this pass did a full object-store listing. + pub full_file_listing: bool, } pub type RegionHookRef = Arc; diff --git a/src/mito2/src/gc.rs b/src/mito2/src/gc.rs index 61ba362a7b..8c0e72fc08 100644 --- a/src/mito2/src/gc.rs +++ b/src/mito2/src/gc.rs @@ -41,6 +41,7 @@ use crate::access_layer::AccessLayerRef; use crate::cache::CacheManagerRef; use crate::cache::file_cache::FileType; use crate::config::MitoConfig; +use crate::engine::region_hook::{RegionGcInfo, RegionHookRef}; use crate::error::{ DurationOutOfRangeSnafu, InvalidRequestSnafu, JoinSnafu, OpenDalSnafu, Result, TooManyGcJobsSnafu, UnexpectedSnafu, @@ -205,6 +206,9 @@ pub struct LocalGcWorker { /// Set to false for regular GC operations to optimize performance. /// Set to true periodically or when you need to clean up orphan files. pub full_file_listing: bool, + /// The region hook (if any), fired via `on_region_gc` after each GC pass so + /// extensions with sidecar files outside the mito2 region dir can clean up. + pub(crate) hook: Option, } pub struct ManifestOpenConfig { @@ -240,6 +244,7 @@ impl LocalGcWorker { file_ref_manifest: FileRefsManifest, limiter: &GcLimiterRef, full_file_listing: bool, + hook: Option, ) -> Result { if let Some(first_region_id) = regions_to_gc.keys().next() { let table_id = first_region_id.table_id(); @@ -267,6 +272,7 @@ impl LocalGcWorker { file_ref_manifest, _permit: permit, full_file_listing, + hook, }) } @@ -303,6 +309,7 @@ impl LocalGcWorker { let mut deleted_files = HashMap::new(); let mut deleted_indexes = HashMap::new(); let mut processed_regions = HashSet::new(); + let mut need_retry_regions = HashSet::new(); let tmp_ref_files = self.read_tmp_ref_files().await?; for (region_id, region) in &self.regions { let per_region_time = std::time::Instant::now(); @@ -321,23 +328,33 @@ impl LocalGcWorker { .get(region_id) .cloned() .unwrap_or_else(HashSet::new); - let files = self + let outcome = self .do_region_gc(*region_id, region.clone(), &tmp_ref_files) .await?; - let index_files = files + let RegionGcOutcome { + removed_files, + extension_cleanup_failed, + } = outcome; + let index_files = removed_files .iter() .filter_map(|f| f.index_version().map(|v| (f.file_id(), v))) .collect_vec(); - let data_files = files + let data_files = removed_files .into_iter() .filter_map(|f| match f { RemovedFile::File(file_id, _) => Some(file_id), RemovedFile::Index(_, _) => None, }) .collect(); - deleted_files.insert(*region_id, data_files); - deleted_indexes.insert(*region_id, index_files); - processed_regions.insert(*region_id); + // Don't acknowledge the region as processed until extension cleanup + // succeeds; retry it next pass instead. + if extension_cleanup_failed { + need_retry_regions.insert(*region_id); + } else { + deleted_files.insert(*region_id, data_files); + deleted_indexes.insert(*region_id, index_files); + processed_regions.insert(*region_id); + } debug!( "GC for region {} took {} secs.", region_id, @@ -351,13 +368,26 @@ impl LocalGcWorker { let report = GcReport { deleted_files, deleted_indexes, - need_retry_regions: HashSet::new(), + need_retry_regions, processed_regions, }; Ok(report) } } +/// Per-region outcome of [`LocalGcWorker::do_region_gc`]. +/// +/// `extension_cleanup_failed` records whether a registered extension's +/// [`RegionHook::on_region_gc`] could not finish; the caller must then keep the +/// region un-acknowledged (retry set) so the next GC pass replays the callback. +pub(crate) struct RegionGcOutcome { + /// Files physically deleted this pass. + pub removed_files: Vec, + /// `true` if a registered extension's `on_region_gc` returned `Err`; the + /// caller retries the region next pass. + pub extension_cleanup_failed: bool, +} + impl LocalGcWorker { /// concurrency of listing files per region. /// This is used to limit the number of concurrent listing operations and speed up listing @@ -384,7 +414,7 @@ impl LocalGcWorker { region_id: RegionId, region: Option, tmp_ref_files: &HashSet, - ) -> Result> { + ) -> Result { let mode = if self.full_file_listing { "full_listing" } else { @@ -428,7 +458,10 @@ impl LocalGcWorker { GC_ERRORS_TOTAL .with_label_values(&["manifest_mismatch"]) .inc(); - return Ok(vec![]); + return Ok(RegionGcOutcome { + removed_files: vec![], + extension_cleanup_failed: false, + }); } Some(manifest) } else { @@ -531,7 +564,45 @@ impl LocalGcWorker { "Successfully deleted {} unused files for region {}", unused_file_cnt, region_id ); - if let Some(region) = ®ion { + + // Notify extensions so they can clean up sidecar files. Fire when there + // are removed files, or on a full-listing pass (lets extensions reconcile + // orphans even when mito deleted nothing). Cleanup is always scoped to + // `removed_files`; see `RegionGcInfo`. + let extension_cleanup_failed = if let Some(hook) = &self.hook + && (!deletable_files.is_empty() || self.full_file_listing) + { + let region_metadata = region.as_ref().map(|r| r.metadata()); + let result = hook + .on_region_gc( + region_id, + region_metadata.as_ref(), + &self.access_layer, + &RegionGcInfo { + removed_files: &deletable_files, + is_region_dropped, + full_file_listing: self.full_file_listing, + }, + ) + .await; + if let Err(err) = result { + warn!( + err; + "Region hook on_region_gc failed for region {}, will retry on the next GC pass", + region_id + ); + true + } else { + false + } + } else { + false + }; + + // Defer clearing the manifest tracking until the extension succeeds, so + // a failed live-region cleanup is replayed next pass. (Dropped regions + // have no manifest.) + if !extension_cleanup_failed && let Some(region) = ®ion { let _update_timer = GC_DURATION_SECONDS .with_label_values(&["update_manifest"]) .start_timer(); @@ -539,7 +610,10 @@ impl LocalGcWorker { .await?; } - Ok(deletable_files) + Ok(RegionGcOutcome { + removed_files: deletable_files, + extension_cleanup_failed, + }) } #[common_telemetry::tracing::instrument( diff --git a/src/mito2/src/gc/worker_test.rs b/src/mito2/src/gc/worker_test.rs index 794b5d61c0..1efa70e6e9 100644 --- a/src/mito2/src/gc/worker_test.rs +++ b/src/mito2/src/gc/worker_test.rs @@ -14,18 +14,25 @@ use std::collections::{BTreeMap, HashMap, HashSet}; use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use api::v1::Rows; +use async_trait::async_trait; +use common_base::Plugins; use common_telemetry::init_default_ut_logging; use futures::TryStreamExt; use object_store::{Entry, ObjectStore, services}; +use store_api::metadata::RegionMetadataRef; use store_api::region_engine::RegionEngine as _; use store_api::region_request::{RegionCompactRequest, RegionRequest}; use store_api::storage::{FileRef, FileRefsManifest, RegionId}; +use crate::access_layer::AccessLayerRef; use crate::config::MitoConfig; use crate::engine::MitoEngine; use crate::engine::compaction_test::{delete_and_flush, put_and_flush}; +use crate::engine::region_hook::{RegionGcInfo, RegionHook, RegionHookRef}; +use crate::error::UnexpectedSnafu; use crate::gc::{GcConfig, LocalGcWorker, should_delete_file}; use crate::manifest::action::RemovedFile; use crate::region::MitoRegionRef; @@ -58,6 +65,7 @@ async fn create_gc_worker( file_ref_manifest.clone(), &mito_engine.gc_limiter(), full_file_listing, + mito_engine.region_hook(), ) .await .unwrap() @@ -738,3 +746,481 @@ async fn test_file_in_tmp_ref_old_mtime_kept() { "File in tmp_ref should NOT be deleted even with old last-modified time" ); } + +/// A region hook that records `on_region_gc` calls, for testing the GC hook. +/// `fail_remaining` simulates a transient extension-cleanup failure: while it +/// is greater than zero, `on_region_gc` returns `Err` (after decrementing), so +/// the GC retry path can be exercised. +#[derive(Debug, Default)] +struct GcCountingHook { + gc_calls: AtomicUsize, + last_is_region_dropped: AtomicBool, + fail_remaining: AtomicUsize, +} + +#[async_trait] +impl RegionHook for GcCountingHook { + async fn on_region_gc( + &self, + _region_id: RegionId, + _region_metadata: Option<&RegionMetadataRef>, + _access_layer: &AccessLayerRef, + info: &RegionGcInfo<'_>, + ) -> Result<(), crate::error::Error> { + self.gc_calls.fetch_add(1, Ordering::Relaxed); + self.last_is_region_dropped + .store(info.is_region_dropped, Ordering::Relaxed); + // Atomically decrement-and-return the previous value; fail while > 0. + let prev = self.fail_remaining.fetch_sub(1, Ordering::Relaxed); + if prev > 0 { + return UnexpectedSnafu { + reason: "simulated extension cleanup failure".to_string(), + } + .fail(); + } + Ok(()) + } +} + +/// `on_region_gc` is **skipped** for a live region in **fast mode** (no full +/// listing) when the GC pass deleted no files — the common case for a periodic +/// pass — to avoid unnecessary hook overhead/I/O. A freshly-flushed region has +/// no deletable files. (A full-listing pass instead always fires the hook; see +/// `test_on_region_gc_fires_on_full_listing_without_deleted_files`.) +#[tokio::test] +async fn test_on_region_gc_skipped_when_no_files_deleted() { + init_default_ut_logging(); + + let mut env = TestEnv::new().await; + + let hook = Arc::new(GcCountingHook::default()); + let plugins = Plugins::new(); + plugins.insert(hook.clone() as RegionHookRef); + + let engine = env + .create_engine_with_plugins( + MitoConfig { + gc: GcConfig { + enable: true, + lingering_time: None, + ..Default::default() + }, + ..Default::default() + }, + plugins, + ) + .await; + + let region_id = RegionId::new(1, 1); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build(); + engine + .handle_request(region_id, RegionRequest::Create(request.clone())) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(0, 3), + }, + ) + .await; + flush_region(&engine, region_id, None).await; + + let region = engine.get_region(region_id).unwrap(); + let version = region.manifest_ctx.manifest().await.manifest_version; + let regions = BTreeMap::from([(region_id, Some(region.clone()))]); + let file_ref_manifest = FileRefsManifest { + file_refs: Default::default(), + manifest_version: [(region_id, version)].into(), + cross_region_refs: HashMap::new(), + }; + + let gc_worker = create_gc_worker(&engine, regions, &file_ref_manifest, false).await; + gc_worker.run().await.unwrap(); + + assert_eq!( + hook.gc_calls.load(Ordering::Relaxed), + 0, + "on_region_gc must be skipped when a live region's fast-mode GC pass deleted no files" + ); +} + +/// `on_region_gc` **fires** (with `is_region_dropped = false`) for a live region +/// when the GC pass actually deleted files. Truncating makes the flushed SST +/// deletable, so the pass reclaims it. +#[tokio::test] +async fn test_on_region_gc_fires_when_live_region_has_deleted_files() { + init_default_ut_logging(); + + let mut env = TestEnv::new().await; + + let hook = Arc::new(GcCountingHook::default()); + let plugins = Plugins::new(); + plugins.insert(hook.clone() as RegionHookRef); + + let engine = env + .create_engine_with_plugins( + MitoConfig { + gc: GcConfig { + enable: true, + lingering_time: None, + ..Default::default() + }, + ..Default::default() + }, + plugins, + ) + .await; + + let region_id = RegionId::new(1, 2); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build(); + engine + .handle_request(region_id, RegionRequest::Create(request.clone())) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(0, 3), + }, + ) + .await; + flush_region(&engine, region_id, None).await; + + // Truncate so the flushed SST becomes a removed file the next GC pass reclaims + // (`lingering_time` is None -> immediately deletable). + engine + .handle_request( + region_id, + RegionRequest::Truncate(store_api::region_request::RegionTruncateRequest::All), + ) + .await + .unwrap(); + + let region = engine.get_region(region_id).unwrap(); + let version = region.manifest_ctx.manifest().await.manifest_version; + let regions = BTreeMap::from([(region_id, Some(region.clone()))]); + let file_ref_manifest = FileRefsManifest { + file_refs: Default::default(), + manifest_version: [(region_id, version)].into(), + cross_region_refs: HashMap::new(), + }; + + let gc_worker = create_gc_worker(&engine, regions, &file_ref_manifest, true).await; + let report = gc_worker.run().await.unwrap(); + assert!( + report + .deleted_files + .get(®ion_id) + .map(|v| !v.is_empty()) + .unwrap_or(false), + "precondition: the GC pass should have deleted the truncated SST" + ); + + assert!( + hook.gc_calls.load(Ordering::Relaxed) >= 1, + "on_region_gc should fire when a live region's GC pass deleted files" + ); + assert!( + !hook.last_is_region_dropped.load(Ordering::Relaxed), + "a live region's GC pass must report is_region_dropped=false" + ); +} + +/// `on_region_gc` fires with `is_region_dropped = true` on the global-GC +/// reclamation path for a dropped/absent region — e.g. a repartitioned-away +/// source region whose mito2 files global GC reclaims. Mirrors +/// [`test_on_region_gc_fires_on_gc_pass`] but drives GC with the region absent +/// (`region: None`), the state `do_region_gc` reports as dropped. This is the +/// path extensions rely on to clean up a deallocated region's sidecar files. +#[tokio::test] +async fn test_on_region_gc_fires_for_dropped_region() { + init_default_ut_logging(); + + let mut env = TestEnv::new().await; + + let hook = Arc::new(GcCountingHook::default()); + let plugins = Plugins::new(); + plugins.insert(hook.clone() as RegionHookRef); + + let engine = env + .create_engine_with_plugins( + MitoConfig { + gc: GcConfig { + enable: true, + lingering_time: None, + ..Default::default() + }, + ..Default::default() + }, + plugins, + ) + .await; + + let region_id = RegionId::new(2, 1); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build(); + engine + .handle_request(region_id, RegionRequest::Create(request.clone())) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(0, 3), + }, + ) + .await; + flush_region(&engine, region_id, None).await; + + // Grab the access layer + manifest version while the region is live, then + // treat the region as dropped (absent) for GC — the global reclamation path + // where `do_region_gc` sees `region: None` and sets `is_region_dropped`. + let region = engine.get_region(region_id).unwrap(); + let access_layer = region.access_layer.clone(); + let version = region.manifest_ctx.manifest().await.manifest_version; + let regions = BTreeMap::from([(region_id, None)]); // dropped/absent + let file_ref_manifest = FileRefsManifest { + file_refs: Default::default(), + manifest_version: [(region_id, version)].into(), + cross_region_refs: HashMap::new(), + }; + + let gc_worker = LocalGcWorker::try_new( + access_layer, + Some(engine.cache_manager()), + regions, + engine.mito_config().gc.clone(), + file_ref_manifest, + &engine.gc_limiter(), + true, // full_file_listing is required when the region is absent + engine.region_hook(), + ) + .await + .unwrap(); + gc_worker.run().await.unwrap(); + + assert!( + hook.gc_calls.load(Ordering::Relaxed) >= 1, + "on_region_gc should fire after a GC pass for a dropped region" + ); + assert!( + hook.last_is_region_dropped.load(Ordering::Relaxed), + "a dropped region's GC pass must report is_region_dropped=true" + ); +} + +/// `on_region_gc` **fires** on a full-listing pass even when mito itself deleted +/// no files, so an extension can reconcile sidecar-only orphans (failed writes, +/// leftovers from a previous cleanup). Contrast with the fast-mode skip in +/// `test_on_region_gc_skipped_when_no_files_deleted`. +#[tokio::test] +async fn test_on_region_gc_fires_on_full_listing_without_deleted_files() { + init_default_ut_logging(); + + let mut env = TestEnv::new().await; + + let hook = Arc::new(GcCountingHook::default()); + let plugins = Plugins::new(); + plugins.insert(hook.clone() as RegionHookRef); + + let engine = env + .create_engine_with_plugins( + MitoConfig { + gc: GcConfig { + enable: true, + lingering_time: None, + ..Default::default() + }, + ..Default::default() + }, + plugins, + ) + .await; + + let region_id = RegionId::new(1, 3); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build(); + engine + .handle_request(region_id, RegionRequest::Create(request.clone())) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(0, 3), + }, + ) + .await; + flush_region(&engine, region_id, None).await; + + // No truncate: there are no removed files to reclaim this pass. + let region = engine.get_region(region_id).unwrap(); + let version = region.manifest_ctx.manifest().await.manifest_version; + let regions = BTreeMap::from([(region_id, Some(region.clone()))]); + let file_ref_manifest = FileRefsManifest { + file_refs: Default::default(), + manifest_version: [(region_id, version)].into(), + cross_region_refs: HashMap::new(), + }; + + let gc_worker = create_gc_worker(&engine, regions, &file_ref_manifest, true).await; + gc_worker.run().await.unwrap(); + + assert!( + hook.gc_calls.load(Ordering::Relaxed) >= 1, + "on_region_gc must fire on a full-listing pass even when no files were deleted" + ); + assert!( + !hook.last_is_region_dropped.load(Ordering::Relaxed), + "a live region's GC pass must report is_region_dropped=false" + ); +} + +/// A failing `on_region_gc` must keep the region un-acknowledged: it lands in +/// `need_retry_regions` and is excluded from `processed_regions` / `deleted_files`, +/// so metasrv does not treat it as fully cleaned and the next GC pass replays +/// the callback. +#[tokio::test] +async fn test_on_region_gc_failure_routes_region_to_retry() { + init_default_ut_logging(); + + let mut env = TestEnv::new().await; + + // Fail the first (and only, in this test) extension cleanup. + let hook = Arc::new(GcCountingHook { + fail_remaining: AtomicUsize::new(1), + ..Default::default() + }); + let plugins = Plugins::new(); + plugins.insert(hook.clone() as RegionHookRef); + + let engine = env + .create_engine_with_plugins( + MitoConfig { + gc: GcConfig { + enable: true, + lingering_time: None, + ..Default::default() + }, + ..Default::default() + }, + plugins, + ) + .await; + + let region_id = RegionId::new(1, 4); + env.get_schema_metadata_manager() + .register_region_table_info( + region_id.table_id(), + "test_table", + "test_catalog", + "test_schema", + None, + env.get_kv_backend(), + ) + .await; + + let request = CreateRequestBuilder::new().build(); + engine + .handle_request(region_id, RegionRequest::Create(request.clone())) + .await + .unwrap(); + put_rows( + &engine, + region_id, + Rows { + schema: rows_schema(&request), + rows: build_rows(0, 3), + }, + ) + .await; + flush_region(&engine, region_id, None).await; + // Truncate so the flushed SST becomes deletable. + engine + .handle_request( + region_id, + RegionRequest::Truncate(store_api::region_request::RegionTruncateRequest::All), + ) + .await + .unwrap(); + + let region = engine.get_region(region_id).unwrap(); + let version = region.manifest_ctx.manifest().await.manifest_version; + let regions = BTreeMap::from([(region_id, Some(region.clone()))]); + let file_ref_manifest = FileRefsManifest { + file_refs: Default::default(), + manifest_version: [(region_id, version)].into(), + cross_region_refs: HashMap::new(), + }; + + let gc_worker = create_gc_worker(&engine, regions, &file_ref_manifest, true).await; + let report = gc_worker.run().await.unwrap(); + + assert!( + hook.gc_calls.load(Ordering::Relaxed) >= 1, + "on_region_gc should fire when the GC pass deleted files" + ); + assert!( + report.need_retry_regions.contains(®ion_id), + "a failed extension cleanup must keep the region in need_retry_regions" + ); + assert!( + !report.processed_regions.contains(®ion_id), + "a failed extension cleanup must not acknowledge the region as processed" + ); + assert!( + !report.deleted_files.contains_key(®ion_id), + "a failed extension cleanup must not report deleted_files for the region" + ); +}