fix(ecstore): migrate tier free versions during decommission

This commit is contained in:
overtrue
2026-08-23 06:45:57 +08:00
parent a206895fad
commit e1b5c665e8
6 changed files with 394 additions and 118 deletions
+160 -56
View File
@@ -1130,21 +1130,29 @@ fn should_cleanup_decommission_source_entry(decommissioned: usize, total_version
decommissioned.saturating_add(expired) == total_versions
}
/// Disposition reason logged for tier free-version records that decommission
/// skips instead of migrating.
const DECOMMISSION_FREE_VERSION_SKIP_REASON: &str = "tier_free_version_not_migrated";
const DECOMMISSION_FREE_VERSION_MIGRATED_REASON: &str = "tier_free_version_migrated";
const DECOMMISSION_FREE_VERSION_RETAINED_REASON: &str = "tier_free_version_migration_failed";
const DECOMMISSION_FREE_VERSION_SWEEP_REASON: &str = "tier_free_version_unresolved_after_decommission";
const DECOMMISSION_FREE_VERSION_DISPOSITION_REASON: &str = "tier_free_version_disposition_recorded";
/// Counts the tier free-version records present in a decommission entry
/// inventory. The exact loader (`load_file_info_versions_exact`) keeps these
/// records inline in `versions` instead of separating them into
/// `free_versions`, and the migration loop then routes them through the
/// generic delete-marker path: the free-version flag and its remote-tier
/// identity are never carried to the target pool, and a lone record is skipped
/// by the empty-delete-marker rule. Accounting for them here keeps the final
/// sweep from silently omitting records whose free-version disposition was
/// dropped (see docs/architecture/decommission-compatibility.md).
fn decommission_free_versions_skipped(fivs: &FileInfoVersions) -> usize {
fivs.versions.iter().filter(|version| version.tier_free_version()).count()
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
struct DecommissionFreeVersionDisposition {
migrated: usize,
retained: usize,
}
impl DecommissionFreeVersionDisposition {
fn record_migrated(&mut self) {
self.migrated += 1;
}
fn record_retained(&mut self) {
self.retained += 1;
}
fn total(self) -> usize {
self.migrated.saturating_add(self.retained)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -2478,6 +2486,7 @@ fn decommission_remote_tiered_opts(
user_defined: version.metadata.clone(),
src_pool_idx,
data_movement: true,
incl_free_versions: version.tier_free_version(),
include_part_checksums: true,
http_preconditions: Some(crate::data_movement::data_movement_target_precondition()),
expected_bucket_incarnation_id,
@@ -3114,24 +3123,9 @@ impl ECStore {
fivs.versions
.sort_by_key(|v| (v.mod_time.is_none(), std::cmp::Reverse(v.mod_time)));
let skipped_free_versions = decommission_free_versions_skipped(&fivs);
if skipped_free_versions > 0 {
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %entry.name,
skipped_free_versions,
reason = DECOMMISSION_FREE_VERSION_SKIP_REASON,
state = "free_versions_skipped",
"Decommission skipped free-version migration"
);
}
let mut decommissioned: usize = 0;
let mut expired: usize = 0;
let mut free_version_disposition = DecommissionFreeVersionDisposition::default();
let mut cleanup_preflight_allowed_missing = Vec::new();
for version in fivs.versions.iter() {
@@ -3140,6 +3134,80 @@ impl ECStore {
}
decommission_cancel_signal_result(rx.is_cancelled())?;
if version.tier_free_version() {
let version_id = version.version_id.map(|v| v.to_string());
let mut migration_error = None;
let mut migrated = false;
for _ in 0..3 {
match self
.decommission_tiered_object(
bucket.as_str(),
&version.name,
version,
&decommission_remote_tiered_opts(version, version_id.clone(), idx, expected_bucket_incarnation_id),
)
.await
{
Ok(()) => {
migrated = true;
migration_error = None;
break;
}
Err(err) if is_decommission_target_capacity_error(&err) => {
return Err(with_decommission_entry_context(
"decommission_tier_free_version",
bucket.as_str(),
version.name.as_str(),
err,
));
}
Err(err) => migration_error = Some(err),
}
}
{
let mut pool_meta = self.pool_meta.write().await;
if let Err(err) = count_decommission_item(&mut pool_meta, idx, 0, !migrated) {
return Err(with_decommission_entry_context(
"count_decommission_item",
bucket.as_str(),
entry.name.as_str(),
err,
));
}
}
if migrated {
decommissioned += 1;
free_version_disposition.record_migrated();
} else {
free_version_disposition.record_retained();
}
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %version.name,
version_id = ?version_id,
result = ?migration_error,
reason = if migrated {
DECOMMISSION_FREE_VERSION_MIGRATED_REASON
} else {
DECOMMISSION_FREE_VERSION_RETAINED_REASON
},
state = if migrated { "free_version_migrated" } else { "free_version_retained" },
"Decommission free-version disposition recorded"
);
if !migrated {
break;
}
continue;
}
if should_skip_lifecycle_for_data_movement(
self.clone(),
&bucket,
@@ -3427,6 +3495,23 @@ impl ECStore {
}
}
if free_version_disposition.total() > 0 {
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket,
object = %entry.name,
free_versions_migrated = free_version_disposition.migrated,
free_versions_retained = free_version_disposition.retained,
free_versions_total = free_version_disposition.total(),
reason = DECOMMISSION_FREE_VERSION_DISPOSITION_REASON,
state = "free_version_disposition",
"Decommission free-version disposition summary"
);
}
if should_cleanup_decommission_source_entry(decommissioned, fivs.versions.len(), expired) {
if bucket_incarnation_fence.as_ref().is_some_and(|guard| guard.is_lock_lost()) {
return Err(Error::other("decommission bucket incarnation fence was lost before source cleanup"));
@@ -4393,6 +4478,7 @@ impl ECStore {
let lifecycle_config_cb = lifecycle_config.clone();
let object_lock_config_cb = object_lock_config.clone();
let store = Arc::clone(self);
let set_cb = Arc::clone(set);
let callback_rx_cb = callback_rx.clone();
let callback: ListCallback = Arc::new(move |entry: MetaCacheEntry| {
@@ -4402,6 +4488,7 @@ impl ECStore {
let lifecycle_config = lifecycle_config_cb.clone();
let object_lock_config = object_lock_config_cb.clone();
let store = Arc::clone(&store);
let set = Arc::clone(&set_cb);
let callback_rx = callback_rx_cb.clone();
Box::pin(async move {
if callback_rx.is_cancelled() {
@@ -4416,11 +4503,14 @@ impl ECStore {
return;
}
let fivs = match load_decommission_entry_versions(
let fivs = match load_decommission_entry_exact_versions(
&set,
&entry,
&bucket_name,
"check_after_decommission.file_info_versions",
) {
)
.await
{
Ok(fivs) => fivs,
Err(err) => {
let mut first_err = entry_error.lock().await;
@@ -4433,7 +4523,23 @@ impl ECStore {
};
let mut remaining = 0;
for version in &fivs.versions {
for version in fivs.versions.iter().chain(fivs.free_versions.iter()) {
if version.tier_free_version() {
remaining += 1;
debug!(
event = EVENT_DECOMMISSION_ENTRY,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_POOLS,
pool_index = idx,
bucket = %bucket_name,
object = %entry.name,
version_id = ?version.version_id,
reason = DECOMMISSION_FREE_VERSION_SWEEP_REASON,
state = "free_version_retained",
"Decommission final sweep retained a free version"
);
continue;
}
if version.deleted {
continue;
}
@@ -4787,6 +4893,12 @@ mod tests {
assert!(opts.include_part_checksums);
assert!(opts.http_preconditions.is_some());
assert_eq!(opts.expected_bucket_incarnation_id, Some(incarnation));
assert!(!opts.incl_free_versions);
let mut free_version = version;
free_version.set_tier_free_version();
let free_opts = decommission_remote_tiered_opts(&free_version, Some("free-version-id".to_string()), 9, Some(incarnation));
assert!(free_opts.incl_free_versions);
}
#[test]
@@ -5491,13 +5603,13 @@ pub(crate) fn fallback_free_capacity_dedup(disks: &[rustfs_madmin::Disk]) -> usi
#[cfg(test)]
mod pools_tests {
use super::{
DECOMMISSION_FREE_VERSION_SKIP_REASON, DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD,
DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF, DecomBucketInfo, DecommissionStartPoolState, DecommissionTerminalState,
ListCallback, PoolDecommissionInfo, PoolMeta, PoolSpaceInfo, PoolStatus, apply_decommission_status_space_info,
DECOMMISSION_FREE_VERSION_MIGRATED_REASON, DECOMMISSION_FREE_VERSION_RETAINED_REASON,
DECOMMISSION_PROGRESS_SAVE_INTERVAL, DECOMMISSION_PROGRESS_SAVE_ITEM_THRESHOLD, DECOMMISSION_PROGRESS_SAVE_RETRY_BACKOFF,
DecomBucketInfo, DecommissionFreeVersionDisposition, 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, count_decommission_item, decommission_cancel_signal_result,
decommission_free_versions_skipped, decommission_item_size, decommission_meta_bucket_options,
decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result, decommission_item_size,
decommission_meta_bucket_options, decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency,
ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available,
ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool,
ensure_decommission_start_local_leader, ensure_decommission_start_pool_states,
@@ -6797,24 +6909,16 @@ mod pools_tests {
}
#[test]
fn decommission_free_version_accounting_reports_skipped_records() {
let mut fivs = FileInfoVersions::default();
assert_eq!(decommission_free_versions_skipped(&fivs), 0);
fn decommission_free_version_accounting_records_migration_and_retention() {
let mut disposition = DecommissionFreeVersionDisposition::default();
disposition.record_migrated();
disposition.record_retained();
fivs.versions.push(FileInfo {
name: "object.txt".to_string(),
..Default::default()
});
let mut free_one = FileInfo::default();
free_one.set_tier_free_version();
fivs.versions.push(free_one);
let mut free_two = FileInfo::default();
free_two.set_tier_free_version();
free_two.transition_tier = "WARM".to_string();
fivs.versions.push(free_two);
assert_eq!(decommission_free_versions_skipped(&fivs), 2);
assert_eq!(DECOMMISSION_FREE_VERSION_SKIP_REASON, "tier_free_version_not_migrated");
assert_eq!(disposition.migrated, 1);
assert_eq!(disposition.retained, 1);
assert_eq!(disposition.total(), 2);
assert_eq!(DECOMMISSION_FREE_VERSION_MIGRATED_REASON, "tier_free_version_migrated");
assert_eq!(DECOMMISSION_FREE_VERSION_RETAINED_REASON, "tier_free_version_migration_failed");
}
#[test]
+108
View File
@@ -4642,6 +4642,63 @@ fn resolve_delete_version_state(opts: &ObjectOptions, goi: &ObjectInfo, version_
}
impl SetDisks {
/// Publish an internal tier free-version record without changing its
/// delete-marker shape or remote-tier identity. The caller holds the
/// source and target object locks; a write quorum is required before the
/// source cleanup may remove the original record.
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tier_free_version(
&self,
bucket: &str,
object: &str,
fi: &FileInfo,
opts: &ObjectOptions,
) -> Result<()> {
if !fi.deleted || !fi.tier_free_version() {
return Err(Error::other("decommission tier free-version write requires a free version record"));
}
if opts
.namespace_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
|| opts
.bucket_lifecycle_lock_fence
.as_ref()
.is_some_and(NamespaceLockFence::is_lock_lost)
{
return Err(StorageError::NamespaceLockQuorumUnavailable {
mode: "decommission_tier_free_version_commit",
bucket: bucket.to_string(),
object: object.to_string(),
required: 1,
achieved: 0,
});
}
let disks = self.disks.read().await.clone();
let write_quorum = self.default_write_quorum();
let futures = disks.into_iter().map(|disk| {
let file_info = fi.clone();
async move {
if let Some(disk) = disk {
disk.write_metadata("", bucket, object, file_info).await
} else {
Err(DiskError::DiskNotFound)
}
}
});
let mut errs = Vec::new();
for result in join_all(futures).await {
match result {
Ok(_) => errs.push(None),
Err(err) => errs.push(Some(err)),
}
}
resolve_tiered_decommission_write_quorum_result(&errs, write_quorum, bucket, object)
}
#[tracing::instrument(skip(self, fi, opts))]
pub(crate) async fn decommission_tiered_object(
&self,
@@ -9864,6 +9921,57 @@ mod tests {
assert_ne!(updated.erasure.distribution, original.erasure.distribution);
}
#[tokio::test]
async fn decommission_tier_free_version_preserves_remote_identity() {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "free-version-decommission";
let object = "object.txt";
let version_id = Uuid::new_v4();
let mut free_version = FileInfo {
name: object.to_string(),
volume: bucket.to_string(),
version_id: Some(version_id),
mod_time: Some(time::OffsetDateTime::now_utc()),
deleted: true,
transition_tier: "WARM-TIER".to_string(),
transitioned_objname: "remote/object".to_string(),
..Default::default()
};
free_version.set_tier_free_version();
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("free-version metadata should reach the target quorum");
set_disks
.decommission_tier_free_version(bucket, object, &free_version, &ObjectOptions::default())
.await
.expect("replaying the same free-version metadata should be idempotent");
let versions = set_disks
.load_file_info_versions_exact(bucket, object)
.await
.expect("migrated free-version metadata should decode")
.expect("migrated free-version metadata should exist");
let migrated = versions
.versions
.iter()
.find(|version| version.version_id == Some(version_id))
.expect("free version should be present on the target");
assert_eq!(
versions
.versions
.iter()
.filter(|version| version.version_id == Some(version_id))
.count(),
1
);
assert!(migrated.tier_free_version());
assert_eq!(migrated.transition_tier, "WARM-TIER");
assert_eq!(migrated.transitioned_objname, "remote/object");
}
#[test]
fn test_resolve_tiered_decommission_write_quorum_result_allows_successful_quorum() {
let errs = vec![None, None, Some(DiskError::DiskNotFound), None];
+80 -14
View File
@@ -1152,6 +1152,15 @@ fn tiered_data_movement_source_matches(
&& expected_backend == current_backend)
}
fn decommission_free_version_overwrite_error(bucket: &str, object: &str, version_id: Option<Uuid>) -> Error {
StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
version_id.map(|id| id.to_string()).unwrap_or_default(),
)
.into()
}
fn should_check_data_movement_resume_target(src_pool_idx: usize, target_pool_idx: usize) -> bool {
target_pool_idx != src_pool_idx
}
@@ -1776,6 +1785,43 @@ impl ECStore {
)
}
async fn has_equivalent_data_movement_tier_free_version(
&self,
bucket: &str,
object: &str,
source: &rustfs_filemeta::FileInfo,
target_pool_idx: usize,
) -> Result<bool> {
let pool = self
.pools
.get(target_pool_idx)
.ok_or_else(|| Error::other(format!("invalid tiered data movement target pool {target_pool_idx}")))?;
let logical_object = decode_dir_object(object);
let Some(versions) = pool
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, &logical_object)
.await?
else {
return Ok(false);
};
let Some(target) = versions
.versions
.iter()
.find(|version| version.version_id == source.version_id)
else {
return Ok(false);
};
if !target.tier_free_version() {
return Err(decommission_free_version_overwrite_error(bucket, object, source.version_id));
}
if tiered_data_movement_source_matches(source, target)? {
Ok(true)
} else {
Err(decommission_free_version_overwrite_error(bucket, object, source.version_id))
}
}
fn resolve_decommission_target_pool_idx_result(result: Result<usize>, bucket: &str, object: &str) -> Result<usize> {
result.map_err(|err| Error::other(format!("failed to select decommission target pool for {bucket}/{object}: {err}")))
}
@@ -1795,6 +1841,10 @@ impl ECStore {
check_put_object_args(bucket, object)?;
let mut opts = opts.clone();
let is_free_version = fi.tier_free_version();
if is_free_version {
opts.incl_free_versions = true;
}
let bucket_incarnation_fence = if is_meta_bucketname(bucket) {
None
} else {
@@ -1849,7 +1899,7 @@ impl ECStore {
versions
.versions
.iter()
.find(|current| current.version_id == fi.version_id && !current.tier_free_version())
.find(|current| current.version_id == fi.version_id && current.tier_free_version() == is_free_version)
})
.ok_or_else(|| to_object_err(StorageError::FileNotFound, vec![bucket, object.as_str()]))?;
if !tiered_data_movement_source_matches(fi, current_source)? {
@@ -1864,24 +1914,40 @@ impl ECStore {
.get_available_pool_idx_excluding(bucket, &object, fi.size, opts.src_pool_idx)
.await;
let target_pool_idx = resolve_data_movement_resume_target_pool(idx, resume_target_pool_idx, opts.src_pool_idx);
if is_free_version && target_pool_idx == opts.src_pool_idx {
return Err(Error::DiskFull);
}
let equivalent = if is_free_version {
self.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, target_pool_idx)
.await?
} else {
self.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.await?
};
if equivalent {
return Ok(());
}
return Err(decommission_free_version_overwrite_error(bucket, &object, fi.version_id));
}
let result = if is_free_version {
if self
.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, target_pool_idx)
.has_equivalent_data_movement_tier_free_version(bucket, &object, &fi, idx)
.await?
{
return Ok(());
}
return Err(StorageError::DataMovementOverwriteErr(
bucket.to_owned(),
object.to_owned(),
opts.version_id.clone().unwrap_or_default(),
));
}
let result = self.pools[idx]
.get_disks_by_key(&object)
.decommission_tiered_object(bucket, &object, &fi, &opts)
.await;
self.pools[idx]
.get_disks_by_key(&object)
.decommission_tier_free_version(bucket, &object, &fi, &opts)
.await
} else {
self.pools[idx]
.get_disks_by_key(&object)
.decommission_tiered_object(bucket, &object, &fi, &opts)
.await
};
if matches!(result, Err(Error::PreconditionFailed)) {
if self
.has_equivalent_data_movement_tiered_object(bucket, &object, &fi, &opts, idx)
+3 -4
View File
@@ -102,10 +102,9 @@ pub const TRANSITION_PENDING: &str = "pending";
/// delete paths the same obligation is also carried by a committed tier-journal
/// entry; deletes without such an entry (for example a removed version whose
/// transition state decodes as unknown) rely on this record alone until the
/// worker removes it after a successful remote delete. Decommission does not
/// preserve these semantics: its exact inventory keeps the records inline in
/// `versions` and the migration loop treats them as ordinary delete markers —
/// see docs/architecture/decommission-compatibility.md.
/// worker removes it after a successful remote delete. Decommission preserves
/// the record and its remote identity on the target pool before source cleanup
/// — see docs/architecture/decommission-compatibility.md.
pub const FREE_VERSION: &str = "free-version";
pub const TRANSITION_STATUS: &str = "transition-status";
+4 -4
View File
@@ -2730,10 +2730,10 @@ impl MetaObject {
/// identity so the lifecycle worker can issue the idempotent remote delete
/// and only then remove the record; until then the recovery scan and the
/// usage scanner keep re-enqueueing it. S3 and lifecycle deletes also
/// persist a committed tier-journal entry for the same remote delete, so a
/// record destroyed without its remote delete (as decommission does when it
/// treats these records as ordinary delete markers) strands only the
/// journal-less cases — see docs/architecture/decommission-compatibility.md.
/// persist a committed tier-journal entry for the same remote delete. The
/// decommission path copies this record unchanged before source cleanup,
/// including when the transition state is unknown — see
/// docs/architecture/decommission-compatibility.md.
pub fn init_free_version(&self, fi: &FileInfo) -> Result<(FileMetaVersion, bool)> {
if fi.skip_tier_free_version() {
return Ok((FileMetaVersion::default(), false));