mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 12:26:37 +00:00
fix(ecstore): migrate ILM metadata during decommission
This commit is contained in:
@@ -31,3 +31,5 @@ pub mod tier_free_version_recovery;
|
||||
pub mod tier_last_day_stats;
|
||||
pub mod tier_sweeper;
|
||||
pub mod transition_transaction;
|
||||
|
||||
pub(crate) const ILM_META_PREFIX: &str = "ilm";
|
||||
|
||||
@@ -16,7 +16,7 @@ use crate::bucket::replication::replication_state_from_filemeta;
|
||||
use crate::bucket::versioning_sys::BucketVersioningSys;
|
||||
use crate::bucket::{
|
||||
lifecycle::{
|
||||
LifecycleExpiryConfigs,
|
||||
ILM_META_PREFIX, LifecycleExpiryConfigs,
|
||||
bucket_lifecycle_audit::LcEventSrc,
|
||||
bucket_lifecycle_ops::{
|
||||
LifecycleOps, apply_expiry_on_transitioned_object, apply_expiry_rule_in, eval_action_from_lifecycle,
|
||||
@@ -93,6 +93,7 @@ const DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD: usize = 1000;
|
||||
const DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF: Duration = Duration::seconds(1);
|
||||
const DECOMMISSION_BUCKET_CONCURRENCY_ENV: &str = "RUSTFS_DECOMMISSION_BUCKET_CONCURRENCY";
|
||||
const DECOMMISSION_BUCKET_CONCURRENCY_DEFAULT_CAP: usize = 4;
|
||||
const DECOMMISSION_META_PREFIXES: [&str; 3] = [CONFIG_PREFIX, BUCKET_META_PREFIX, ILM_META_PREFIX];
|
||||
const DECOMMISSION_TARGET_CAPACITY_OVERHEAD_PERCENT: usize = 30;
|
||||
const DECOMMISSION_LISTING_MAX_ATTEMPTS: usize = 3;
|
||||
const DECOMMISSION_LISTING_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(5);
|
||||
@@ -204,6 +205,19 @@ fn is_decommission_meta_bucket(bucket: &DecomBucketInfo) -> bool {
|
||||
bucket.name == RUSTFS_META_BUCKET
|
||||
}
|
||||
|
||||
fn decommission_meta_buckets() -> [DecomBucketInfo; DECOMMISSION_META_PREFIXES.len()] {
|
||||
DECOMMISSION_META_PREFIXES.map(|prefix| DecomBucketInfo {
|
||||
name: RUSTFS_META_BUCKET.to_owned(),
|
||||
prefix: prefix.to_owned(),
|
||||
})
|
||||
}
|
||||
|
||||
fn reconcile_decommission_meta_buckets(meta: &mut PoolMeta, idx: usize) -> bool {
|
||||
let before = meta.pending_buckets(idx).len();
|
||||
meta.queue_buckets(idx, decommission_meta_buckets().into());
|
||||
meta.pending_buckets(idx).len() != before
|
||||
}
|
||||
|
||||
fn split_decommission_buckets(buckets: Vec<DecomBucketInfo>) -> (Vec<DecomBucketInfo>, Vec<DecomBucketInfo>) {
|
||||
let mut regular = Vec::with_capacity(buckets.len());
|
||||
let mut meta = Vec::new();
|
||||
@@ -2879,13 +2893,20 @@ impl ECStore {
|
||||
}
|
||||
|
||||
async fn promote_queued_decommission(&self, idx: usize) -> Result<()> {
|
||||
let promoted = {
|
||||
let (changed, previous_pool_meta) = {
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
pool_meta.promote_queued_decommission(idx)
|
||||
let previous_pool_meta = pool_meta.clone();
|
||||
let reconciled = reconcile_decommission_meta_buckets(&mut pool_meta, idx);
|
||||
let promoted = pool_meta.promote_queued_decommission(idx);
|
||||
(reconciled || promoted, previous_pool_meta)
|
||||
};
|
||||
|
||||
if promoted {
|
||||
self.save_current_pool_meta().await?;
|
||||
if changed {
|
||||
if let Err(err) = self.save_current_pool_meta().await {
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
rollback_decommission_pool_meta(&mut pool_meta, previous_pool_meta);
|
||||
return Err(err);
|
||||
}
|
||||
if let Some(notification_sys) = runtime_sources::notification_sys() {
|
||||
let stage = format!("promote_queued_decommission for pool {idx}");
|
||||
resolve_decommission_pool_meta_reload_result(notification_sys.reload_pool_meta().await, stage.as_str())?;
|
||||
@@ -4206,17 +4227,16 @@ impl ECStore {
|
||||
|
||||
let decom_buckets = self.get_buckets_to_decommission().await?;
|
||||
|
||||
let mut healed_buckets = HashSet::with_capacity(decom_buckets.len());
|
||||
for bk in decom_buckets.iter() {
|
||||
resolve_decommission_preflight_heal_result(&bk.name, self.heal_bucket(&bk.name, &HealOpts::default()).await)?;
|
||||
if healed_buckets.insert(bk.name.as_str()) {
|
||||
resolve_decommission_preflight_heal_result(&bk.name, self.heal_bucket(&bk.name, &HealOpts::default()).await)?;
|
||||
}
|
||||
}
|
||||
|
||||
let meta_buckets = [
|
||||
path_join(&[PathBuf::from(RUSTFS_META_BUCKET), PathBuf::from(CONFIG_PREFIX)]),
|
||||
path_join(&[PathBuf::from(RUSTFS_META_BUCKET), PathBuf::from(BUCKET_META_PREFIX)]),
|
||||
];
|
||||
|
||||
let meta_bucket_opts = decommission_meta_bucket_options();
|
||||
for bk in meta_buckets.iter() {
|
||||
for prefix in DECOMMISSION_META_PREFIXES {
|
||||
let bk = path_join(&[PathBuf::from(RUSTFS_META_BUCKET), PathBuf::from(prefix)]);
|
||||
if let Err(err) = self
|
||||
.make_bucket(bk.to_string_lossy().to_string().as_str(), &meta_bucket_opts)
|
||||
.await
|
||||
@@ -4325,14 +4345,7 @@ impl ECStore {
|
||||
})
|
||||
.collect();
|
||||
|
||||
ret.push(DecomBucketInfo {
|
||||
name: RUSTFS_META_BUCKET.to_owned(),
|
||||
prefix: CONFIG_PREFIX.to_owned(),
|
||||
});
|
||||
ret.push(DecomBucketInfo {
|
||||
name: RUSTFS_META_BUCKET.to_owned(),
|
||||
prefix: BUCKET_META_PREFIX.to_owned(),
|
||||
});
|
||||
ret.extend(decommission_meta_buckets());
|
||||
|
||||
Ok(ret)
|
||||
}
|
||||
@@ -5458,7 +5471,7 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi
|
||||
#[cfg(test)]
|
||||
mod pools_tests {
|
||||
use super::{
|
||||
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF,
|
||||
DECOMMISSION_META_PREFIXES, DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
|
||||
DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState, ListCallback, PoolDecommissionInfo, PoolMeta,
|
||||
PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info, bind_decommission_cancelers,
|
||||
bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state,
|
||||
@@ -5473,18 +5486,17 @@ mod pools_tests {
|
||||
has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested,
|
||||
load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done,
|
||||
merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result,
|
||||
pool_meta_has_active_decommission, record_decommission_entry_error, require_decommission_store,
|
||||
pool_meta_has_active_decommission, reconcile_decommission_meta_buckets, require_decommission_store,
|
||||
resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state,
|
||||
resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result,
|
||||
resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error,
|
||||
resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result,
|
||||
resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result,
|
||||
resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result,
|
||||
resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result,
|
||||
resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result,
|
||||
resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result,
|
||||
resolve_start_decommission_pool_meta_reload_result, rollback_start_decommission_pool_meta,
|
||||
run_decommission_buckets_bounded, run_decommission_listing_with_retry, should_cleanup_decommission_source_entry,
|
||||
should_continue_decommission_queue, should_count_decommission_version_complete,
|
||||
resolve_decommission_pool_meta_reload_result, resolve_decommission_preflight_heal_result,
|
||||
resolve_decommission_progress_save_result, resolve_decommission_spawn_failure_result,
|
||||
resolve_decommission_terminal_mark_after_error_result, resolve_decommission_terminal_mark_result,
|
||||
resolve_decommission_update_after_result, resolve_start_decommission_pool_meta_reload_result,
|
||||
rollback_start_decommission_pool_meta, run_decommission_buckets_bounded, run_decommission_listing_with_retry,
|
||||
should_cleanup_decommission_source_entry, should_continue_decommission_queue, should_count_decommission_version_complete,
|
||||
should_preserve_decommission_canceled_state, should_reject_decommission_cancel_as_terminal,
|
||||
should_retry_decommission_cancel_reload, should_retry_decommission_listing, should_skip_canceled_decommission_routine,
|
||||
split_decommission_buckets, take_and_cancel_decommission_canceler, take_decommission_canceler,
|
||||
@@ -5745,6 +5757,10 @@ mod pools_tests {
|
||||
name: crate::disk::RUSTFS_META_BUCKET.to_string(),
|
||||
prefix: crate::disk::BUCKET_META_PREFIX.to_string(),
|
||||
},
|
||||
DecomBucketInfo {
|
||||
name: crate::disk::RUSTFS_META_BUCKET.to_string(),
|
||||
prefix: crate::bucket::lifecycle::ILM_META_PREFIX.to_string(),
|
||||
},
|
||||
]);
|
||||
|
||||
assert_eq!(
|
||||
@@ -5753,10 +5769,61 @@ mod pools_tests {
|
||||
);
|
||||
assert_eq!(
|
||||
meta.iter().map(|bucket| bucket.prefix.as_str()).collect::<Vec<_>>(),
|
||||
vec![crate::config::com::CONFIG_PREFIX, crate::disk::BUCKET_META_PREFIX,]
|
||||
vec![
|
||||
crate::config::com::CONFIG_PREFIX,
|
||||
crate::disk::BUCKET_META_PREFIX,
|
||||
crate::bucket::lifecycle::ILM_META_PREFIX,
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_decommission_meta_prefixes_cover_durable_ilm_records() {
|
||||
let ilm_prefix = format!("{}/", crate::bucket::lifecycle::ILM_META_PREFIX);
|
||||
let durable_prefixes = [
|
||||
crate::bucket::lifecycle::tier_delete_journal::TIER_DELETE_JOURNAL_PREFIX,
|
||||
crate::bucket::lifecycle::transition_transaction::TRANSITION_TRANSACTION_RECORD_PREFIX,
|
||||
crate::bucket::lifecycle::manual_transition_job::MANUAL_TRANSITION_JOB_RECORD_PREFIX,
|
||||
crate::bucket::lifecycle::manual_transition_job::MANUAL_TRANSITION_SCOPE_RECORD_PREFIX,
|
||||
crate::bucket::lifecycle::manual_transition_job::MANUAL_TRANSITION_TASK_PREFIX,
|
||||
crate::bucket::lifecycle::manual_transition_job::MANUAL_TRANSITION_WORKER_RESULT_PREFIX,
|
||||
];
|
||||
|
||||
assert!(DECOMMISSION_META_PREFIXES.contains(&crate::bucket::lifecycle::ILM_META_PREFIX));
|
||||
assert!(
|
||||
durable_prefixes.iter().all(|prefix| prefix.starts_with(&ilm_prefix)),
|
||||
"every durable ILM record must stay under the decommissioned ILM namespace"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_resume_reconciles_missing_decommission_meta_prefixes() {
|
||||
let mut meta = PoolMeta {
|
||||
pools: vec![decommission_test_pool_status(
|
||||
0,
|
||||
Some(PoolDecommissionInfo {
|
||||
queued_buckets: vec![
|
||||
format!("{}/{}", crate::disk::RUSTFS_META_BUCKET, crate::config::com::CONFIG_PREFIX),
|
||||
format!("{}/{}", crate::disk::RUSTFS_META_BUCKET, crate::disk::BUCKET_META_PREFIX),
|
||||
],
|
||||
..Default::default()
|
||||
}),
|
||||
)],
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
assert!(reconcile_decommission_meta_buckets(&mut meta, 0));
|
||||
assert_eq!(
|
||||
meta.pending_buckets(0)
|
||||
.iter()
|
||||
.filter(|bucket| bucket.name == crate::disk::RUSTFS_META_BUCKET)
|
||||
.map(|bucket| bucket.prefix.as_str())
|
||||
.collect::<Vec<_>>(),
|
||||
DECOMMISSION_META_PREFIXES
|
||||
);
|
||||
assert!(!reconcile_decommission_meta_buckets(&mut meta, 0));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_run_decommission_buckets_bounded_respects_limit() {
|
||||
let rx = CancellationToken::new();
|
||||
|
||||
Reference in New Issue
Block a user