diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 301c022a1..abdb9f446 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -779,7 +779,7 @@ fn free_version_physical_topology_generation(api: &ECStore) -> String { rustfs_utils::crypto::hex(hasher.finalize().as_slice()) } -fn free_version_remote_tuple_matches(candidate: &ObjectInfo, expected: &ObjectInfo) -> std::io::Result { +pub(crate) fn free_version_remote_tuple_matches(candidate: &ObjectInfo, expected: &ObjectInfo) -> std::io::Result { if candidate.transitioned_object.tier != expected.transitioned_object.tier || candidate.transitioned_object.name != expected.transitioned_object.name { diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 8635b57ac..1706433c6 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -7470,6 +7470,20 @@ impl PoolMeta { .is_some_and(is_decommission_suspended) } + pub(crate) fn has_active_decommission_capacity_reservation(&self, idx: usize) -> bool { + self.pools + .get(idx) + .and_then(|pool| pool.decommission.as_ref()) + .is_some_and(|info| { + info.has_decommission_state() + && is_decommission_active(info.complete, info.failed, info.canceled) + && info + .capacity_reservation + .as_ref() + .is_some_and(DecommissionCapacityReservation::active) + }) + } + pub(crate) fn scanner_pause_backlog_pool_writable(&self, idx: usize) -> bool { self.pools.get(idx).is_some_and(|pool| { !pool diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 719c07319..38e442ab2 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -20,6 +20,7 @@ fn to_filemeta_err(err: Error) -> rustfs_filemeta::Error { err.narrow_to_filemeta().unwrap_or_else(rustfs_filemeta::Error::other) } +use crate::bucket::lifecycle::bucket_lifecycle_ops::free_version_remote_tuple_matches; use crate::bucket::metadata_sys::{ get_versioning_config, has_authoritative_never_versioned_state, has_authoritative_never_versioned_state_in, }; @@ -5118,6 +5119,30 @@ fn merge_object_entry_versions(first: &mut MetaCacheEntry, others: impl Iterator } } } + let mut live_remote_references = HashMap::>>::new(); + for (_, info) in versions.values() { + if !info.transitioned_object.free_version && info.transitioned_object.status == rustfs_filemeta::TRANSITION_COMPLETE { + live_remote_references + .entry(info.transitioned_object.tier.clone()) + .or_default() + .entry(info.transitioned_object.name.clone()) + .or_default() + .push(info.clone()); + } + } + // Keep cleanup durable in its source xl.meta, but do not expose it to a + // merged recovery walk while another physical pool still owns the tuple. + versions.retain(|_, (_, info)| { + !info.transitioned_object.free_version + || !live_remote_references + .get(info.transitioned_object.tier.as_str()) + .and_then(|by_name| by_name.get(info.transitioned_object.name.as_str())) + .is_some_and(|candidates| { + candidates + .iter() + .any(|live| free_version_remote_tuple_matches(info, live).unwrap_or(false)) + }) + }); let mut merged = FileMeta::new(); merged.versions = versions.into_values().map(|(version, _)| version).collect(); merged.versions.sort_by(|a, b| { @@ -7435,6 +7460,47 @@ mod test { } } + fn test_transitioned_meta_entry(name: &str, remote_object: &str, delete_source: bool) -> MetaCacheEntry { + let mut source = FileInfo::new(name, 2, 2); + source.volume = "bucket".to_string(); + source.name = name.to_string(); + source.version_id = Some(Uuid::from_u128(1)); + source.versioned = true; + source.size = 1; + source.mod_time = Some(time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")); + source.transition_status = rustfs_filemeta::TRANSITION_COMPLETE.to_string(); + source.transition_tier = "WARM".to_string(); + source.transitioned_objname = remote_object.to_string(); + source.transition_version = Some("remote-version".to_string()); + source.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact; + rustfs_utils::http::metadata_compat::insert_str( + &mut source.metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + "00".repeat(32), + ); + + let mut meta = FileMeta::new(); + meta.add_version(source.clone()) + .expect("test metadata should accept transitioned source"); + if delete_source { + let mut delete = FileInfo { + name: name.to_string(), + version_id: source.version_id, + ..Default::default() + }; + delete.set_tier_free_version_id(&Uuid::from_u128(2).to_string()); + meta.delete_version(&delete) + .expect("transitioned delete should create a free-version owner"); + } + let metadata = meta.marshal_msg().expect("test transitioned metadata should marshal"); + MetaCacheEntry { + name: name.to_string(), + metadata, + cached: Some(meta), + reusable: false, + } + } + fn test_object_with_delete_marker_meta_entry( name: &str, object_mod_time: time::OffsetDateTime, @@ -10638,6 +10704,32 @@ mod test { assert_eq!(versions.versions[0].metadata.get("etag").map(String::as_str), Some("same-etag")); } + #[tokio::test] + async fn merge_entry_channels_defers_free_version_while_same_remote_source_is_live() { + let live = test_transitioned_meta_entry("key", "remote/shared", false); + let free = test_transitioned_meta_entry("key", "remote/shared", true); + for inputs in [vec![live.clone(), free.clone()], vec![free.clone(), live.clone()]] { + let merged = merge_test_object_entries(inputs) + .await + .expect("same remote source and cleanup owner should merge"); + let versions = merged + .file_info_versions_with_free_versions("bucket") + .expect("merged transition history should decode"); + assert_eq!(versions.versions.len(), 1); + assert!(versions.free_versions.is_empty(), "a live remote reference must defer cleanup discovery"); + } + + let unrelated = test_transitioned_meta_entry("key", "remote/other", false); + let merged = merge_test_object_entries(vec![free, unrelated]) + .await + .expect("unrelated remote references should merge"); + let versions = merged + .file_info_versions_with_free_versions("bucket") + .expect("merged transition history should decode"); + assert_eq!(versions.versions.len(), 1); + assert_eq!(versions.free_versions.len(), 1, "an unrelated source must not suppress cleanup"); + } + #[tokio::test] async fn merge_entry_channels_rejects_conflicting_version_identity_and_metadata() { let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"); diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index 58355f48b..22bb521da 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -4965,6 +4965,17 @@ impl ECStore { } }; + if creates_latest_marker && self.is_suspended(pinfo.index).await { + let has_active_reservation = self + .pool_meta + .read() + .await + .has_active_decommission_capacity_reservation(pinfo.index); + if has_active_reservation { + pinfo.index = self.get_pool_idx_no_lock(bucket, object, 0).await?; + } + } + if pinfo.object_info.delete_marker && opts.version_id.is_none() && !creates_latest_marker { pinfo.object_info.name = decode_dir_object(object); return Ok(pinfo.object_info); diff --git a/crates/ecstore/src/store/rebalance.rs b/crates/ecstore/src/store/rebalance.rs index ac313f2c0..e55cc5fdc 100644 --- a/crates/ecstore/src/store/rebalance.rs +++ b/crates/ecstore/src/store/rebalance.rs @@ -677,6 +677,34 @@ impl ECStore { } } + if require_all_pool_reads { + let suspended_pools = { + let pool_meta = self.pool_meta.read().await; + (0..self.pools.len()) + .map(|idx| pool_meta.is_suspended(idx)) + .collect::>() + }; + let candidates = ress + .iter() + .map(|pinfo| LatestObjectInfoCandidate { + info: pinfo.err.is_none().then(|| pinfo.object_info.clone()), + idx: pinfo.index, + err: pinfo.err.clone(), + }) + .collect(); + let (object_info, index) = + resolve_latest_object_info_candidates_with_pool_state(candidates, &suspended_pools, bucket, object, opts)?; + let pools_with_object = self.pools_with_object(&ress, opts).await; + return Ok(( + PoolObjInfo { + index, + object_info, + err: None, + }, + pools_with_object, + )); + } + ress.sort_by(|a, b| { let at = a.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH); let bt = b.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH);