mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-08-18 12:08:22 +00:00
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
This commit is contained in:
@@ -169,6 +169,7 @@ impl GcRegionsHandler {
|
||||
file_ref_manifest.clone(),
|
||||
&mito_engine.gc_limiter(),
|
||||
full_file_listing,
|
||||
mito_engine.region_hook(),
|
||||
)
|
||||
.await
|
||||
.context(GcMitoEngineSnafu {
|
||||
|
||||
@@ -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<RegionId> = self
|
||||
.data
|
||||
.gc_report
|
||||
.as_ref()
|
||||
.map(|r| r.need_retry_regions.clone())
|
||||
.unwrap_or_default();
|
||||
|
||||
let mut table_ids: HashSet<TableId> = 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);
|
||||
|
||||
@@ -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::<RegionHookRef>();
|
||||
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<RegionHookRef> {
|
||||
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<dyn RawEntryReader>,
|
||||
/// 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<RegionHookRef>,
|
||||
#[cfg(feature = "enterprise")]
|
||||
extension_range_provider_factory: Option<BoxedExtensionRangeProviderFactory>,
|
||||
}
|
||||
@@ -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,
|
||||
}),
|
||||
|
||||
@@ -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<dyn RegionHook>;
|
||||
|
||||
+85
-11
@@ -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<RegionHookRef>,
|
||||
}
|
||||
|
||||
pub struct ManifestOpenConfig {
|
||||
@@ -240,6 +244,7 @@ impl LocalGcWorker {
|
||||
file_ref_manifest: FileRefsManifest,
|
||||
limiter: &GcLimiterRef,
|
||||
full_file_listing: bool,
|
||||
hook: Option<RegionHookRef>,
|
||||
) -> Result<Self> {
|
||||
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<RemovedFile>,
|
||||
/// `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<MitoRegionRef>,
|
||||
tmp_ref_files: &HashSet<FileRef>,
|
||||
) -> Result<Vec<RemovedFile>> {
|
||||
) -> Result<RegionGcOutcome> {
|
||||
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(
|
||||
|
||||
@@ -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"
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user