diff --git a/crates/ecstore/src/bucket/bucket_target_sys.rs b/crates/ecstore/src/bucket/bucket_target_sys.rs index 628db8f2e..e401c0539 100644 --- a/crates/ecstore/src/bucket/bucket_target_sys.rs +++ b/crates/ecstore/src/bucket/bucket_target_sys.rs @@ -2000,7 +2000,7 @@ impl TargetClient { object: &str, version_id: Option, opts: RemoveObjectOptions, - ) -> Result<(), S3ClientError> { + ) -> Result, S3ClientError> { let headers = build_remove_object_headers(version_id.as_deref(), &opts); let api_version_id = resolve_delete_api_version_id(version_id, &opts); @@ -2023,7 +2023,11 @@ impl TargetClient { .send() .await { - Ok(_res) => Ok(()), + // A DELETE without a version id on a versioned target creates a delete + // marker and reports the version it assigned. That id is the only + // reliable handle for purging the marker later: a generic S3 target + // does not mirror source version ids. + Ok(res) => Ok(res.version_id().map(ToOwned::to_owned)), Err(e) => match e { SdkError::ServiceError(service_err) => { let err = service_err.into_err(); diff --git a/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs b/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs index 86171b5fd..ae56fec17 100644 --- a/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_filemeta_boundary.rs @@ -52,6 +52,8 @@ pub(crate) fn replication_state_from_filemeta(state: &rustfs_filemeta::Replicati .map(|(arn, status)| (arn.clone(), version_purge_status_from_filemeta(status.clone()))) .collect(), reset_statuses_map: state.reset_statuses_map.clone(), + target_delete_marker_version_ids: state.target_delete_marker_version_ids.clone(), + target_delete_marker_version_ids_corrupt: state.target_delete_marker_version_ids_corrupt, } } @@ -83,5 +85,7 @@ pub fn replication_state_to_filemeta(state: &ReplicationState) -> rustfs_filemet .map(|(arn, status)| (arn.clone(), version_purge_status_to_filemeta(status.clone()))) .collect(), reset_statuses_map: state.reset_statuses_map.clone(), + target_delete_marker_version_ids: state.target_delete_marker_version_ids.clone(), + target_delete_marker_version_ids_corrupt: state.target_delete_marker_version_ids_corrupt, } } diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index 929b9b331..a351b0aa7 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -1634,6 +1634,29 @@ async fn source_delete_marker_missing( } } +/// Which version a delete-marker purge should address on one target. +/// +/// `None` means do not purge at all: the recorded mapping disagreed across the +/// dual internal prefixes, and guessing an id could destroy a live version on +/// the target. `Some(id)` is the exact version the target reported when it +/// accepted the marker; falling back to a source-derived id is only correct +/// when the target mirrors source version ids, which a generic S3 target does +/// not. +fn delete_marker_purge_version_id( + state: Option<&ReplicationState>, + arn: &str, + delete_marker_version_id: Uuid, +) -> Option> { + if state.is_some_and(|state| state.target_delete_marker_version_ids_corrupt) { + return None; + } + let recorded = state.and_then(|state| state.target_delete_marker_version_ids.get(arn).cloned()); + Some(match recorded { + Some(version_id) => Some(version_id), + None => target_delete_version_id(delete_marker_version_id, true), + }) +} + async fn replicate_delete_marker_purge_to_targets(bucket: &str, dobj: &DeletedObjectReplicationInfo, dsc: &ReplicateDecision) { let Some(delete_marker_version_id) = dobj.delete_object.delete_marker_version_id else { return; @@ -1651,11 +1674,27 @@ async fn replicate_delete_marker_purge_to_targets(bucket: &str, dobj: &DeletedOb continue; }; + let Some(purge_version_id) = delete_marker_purge_version_id( + dobj.delete_object.replication_state.as_ref(), + &tgt_entry.arn, + delete_marker_version_id, + ) else { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + bucket, + object = dobj.delete_object.object_name, + arn = tgt_entry.arn, + "Skipping delete-marker purge: recorded target version metadata is inconsistent" + ); + continue; + }; + let _ = tgt_client .remove_object( &tgt_client.bucket, &dobj.delete_object.object_name, - target_delete_version_id(delete_marker_version_id, true), + purge_version_id, replication_delete_marker_purge_remove_options(dobj.delete_object.delete_marker_mtime), ) .await; @@ -2007,16 +2046,24 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli ) .await { - Ok(_) => { + Ok(assigned_version_id) => { debug!( bucket = tgt_client.bucket, object = dobj.delete_object.object_name, version_id = ?version_id, + assigned_version_id = ?assigned_version_id, delete_marker = dobj.delete_object.delete_marker, is_version_purge, "replicate_delete_to_target succeeded" ); if !is_version_purge { + // Record the version the target actually assigned to the marker it + // just created. A later purge addresses that id directly instead of + // deriving one from the source uuid, which only holds when the + // target mirrors source version ids. + if dobj.delete_object.delete_marker { + rinfo.target_delete_marker_version_id = assigned_version_id.filter(|version_id| !version_id.is_empty()); + } rinfo.replication_status = ReplicationStatusType::Completed; } else { rinfo.version_purge_status = VersionPurgeStatusType::Complete; @@ -4009,4 +4056,35 @@ mod tests { assert_eq!(target_delete_version_id(Uuid::nil(), true).as_deref(), Some(NULL_VERSION_ID)); assert_eq!(target_delete_version_id(Uuid::nil(), false), None); } + + #[test] + fn delete_marker_purge_prefers_the_recorded_target_version() { + let source = Uuid::new_v4(); + let arn = "arn:rustfs:replication::target:bucket"; + + // No recorded mapping: fall back to deriving from the source uuid. + assert_eq!(delete_marker_purge_version_id(None, arn, source), Some(Some(source.to_string()))); + + // Recorded mapping wins — a generic S3 target assigns its own id, so the + // derived one would purge the wrong version or nothing at all. + let mut state = ReplicationState::default(); + state + .target_delete_marker_version_ids + .insert(arn.to_string(), "target-assigned-id".to_string()); + assert_eq!( + delete_marker_purge_version_id(Some(&state), arn, source), + Some(Some("target-assigned-id".to_string())) + ); + + // A mapping recorded for a different ARN must not be reused. + assert_eq!( + delete_marker_purge_version_id(Some(&state), "arn:rustfs:replication::other:bucket", source), + Some(Some(source.to_string())) + ); + + // Inconsistent persisted metadata: refuse to purge rather than guess. + let mut corrupt = state.clone(); + corrupt.target_delete_marker_version_ids_corrupt = true; + assert_eq!(delete_marker_purge_version_id(Some(&corrupt), arn, source), None); + } } diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index b2fcc0469..738a09e23 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -762,6 +762,10 @@ impl ObjectInfo { } pub fn replication_state(&self) -> ReplicationState { + // Derived from the durable internal keys, not from the wire form: the + // state's positional encoding skips this map. + let (target_delete_marker_version_ids, target_delete_marker_version_ids_corrupt) = + rustfs_utils::http::target_delete_marker_versions(&self.user_defined); ReplicationState { replication_status_internal: self.replication_status_internal.clone(), version_purge_status_internal: self.version_purge_status_internal.clone(), @@ -779,6 +783,8 @@ impl ObjectInfo { .map(|arn| (arn, v.clone())) }) .collect(), + target_delete_marker_version_ids, + target_delete_marker_version_ids_corrupt, ..Default::default() } } diff --git a/crates/ecstore/src/set_disk/metadata.rs b/crates/ecstore/src/set_disk/metadata.rs index b19d502a0..4d2d9d7e2 100644 --- a/crates/ecstore/src/set_disk/metadata.rs +++ b/crates/ecstore/src/set_disk/metadata.rs @@ -538,6 +538,8 @@ impl SetDisks { || suffix.eq_ignore_ascii_case(http::SUFFIX_REPLICATION_TIMESTAMP) || suffix.eq_ignore_ascii_case(http::SUFFIX_PURGESTATUS) || Self::starts_with_ignore_ascii_case(suffix, http::SUFFIX_REPLICATION_RESET_ARN_PREFIX) + // Raw compatibility keys are normalized and hashed separately below. + || Self::starts_with_ignore_ascii_case(suffix, http::SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX) } fn update_hash_quorum_metadata_map(hasher: &mut Sha256, entries: &HashMap) { @@ -562,6 +564,22 @@ impl SetDisks { key } + /// Hash the per-target delete-marker versions through their normalized form + /// so the dual internal prefixes carrying the same mapping share one + /// identity, while a genuine disagreement between disks still changes the + /// hash and surfaces as a quorum difference. + fn update_hash_target_delete_marker_versions(hasher: &mut Sha256, metadata: &HashMap) { + let (versions, corrupt) = http::target_delete_marker_versions(metadata); + hasher.update([u8::from(corrupt)]); + let mut versions = versions.iter().collect::>(); + versions.sort_by(|left, right| left.0.cmp(right.0)); + hasher.update(versions.len().to_le_bytes()); + for (arn, version_id) in versions { + Self::update_hash_str(hasher, arn); + Self::update_hash_str(hasher, version_id); + } + } + fn update_file_info_quorum_hash(hasher: &mut Sha256, meta: &FileInfo) { hasher.update(meta.size.to_le_bytes()); hasher.update([u8::from(meta.deleted), u8::from(meta.mark_deleted)]); @@ -592,6 +610,7 @@ impl SetDisks { Self::update_hash_optional_bytes(hasher, meta.checksum.as_ref()); Self::update_hash_quorum_metadata_map(hasher, &meta.metadata); + Self::update_hash_target_delete_marker_versions(hasher, &meta.metadata); hasher.update(meta.parts.len().to_le_bytes()); for part in meta.parts.iter() { @@ -1414,4 +1433,39 @@ mod tests { assert!(fallback_disks.iter().any(Option::is_none)); assert!(fallback_parts.iter().any(|part| !part.is_valid())); } + + #[test] + fn target_delete_marker_version_metadata_is_included_in_quorum_hash() { + let suffix = "replication-delete-marker-version-arn:rustfs:replication::target:bucket"; + assert!(SetDisks::is_replication_quorum_metadata_key(&format!( + "{}{}", + http::RUSTFS_INTERNAL_PREFIX, + suffix + ))); + assert!(SetDisks::is_replication_quorum_metadata_key(&format!( + "{}{}", + http::MINIO_INTERNAL_PREFIX, + suffix + ))); + assert!(!SetDisks::is_replication_quorum_metadata_key("x-rustfs-internal-unrelated")); + + let mut left = metadata_quorum_test_fileinfo(OffsetDateTime::now_utc(), 1); + let mut right = left.clone(); + left.metadata + .insert(format!("{}{}", http::RUSTFS_INTERNAL_PREFIX, suffix), "target-version-a".to_string()); + right + .metadata + .insert(format!("{}{}", http::MINIO_INTERNAL_PREFIX, suffix), "target-version-b".to_string()); + assert_ne!(SetDisks::file_info_quorum_hash(&left), SetDisks::file_info_quorum_hash(&right)); + + let mut dual_prefixed = left.clone(); + dual_prefixed + .metadata + .insert(format!("{}{}", http::MINIO_INTERNAL_PREFIX, suffix), "target-version-a".to_string()); + assert_eq!( + SetDisks::file_info_quorum_hash(&left), + SetDisks::file_info_quorum_hash(&dual_prefixed), + "compatible prefixes carrying the same mapping must share one identity" + ); + } } diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index be4a27264..ad3d68797 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -868,6 +868,31 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { } } +/// `ReplicationState::target_delete_marker_version_ids` is skipped by the +/// positional `FileInfo` wire form, so a remote disk would otherwise receive a +/// delete with an empty map and lose the exact per-target version. Copy it into +/// the object's internal metadata — the durable carrier both sides already +/// agree on — before the delete is dispatched. Bounds mirror +/// `persist_target_delete_marker_versions`; anything outside them is dropped +/// rather than forwarded. +fn delete_file_info_with_replication_transport_metadata(fi: &FileInfo) -> FileInfo { + let mut transported = fi.clone(); + let Some(state) = transported.replication_state_internal.as_ref() else { + return transported; + }; + if state.target_delete_marker_version_ids.len() > 1_000 { + return transported; + } + for (arn, version_id) in &state.target_delete_marker_version_ids { + if !arn.starts_with("arn:") || arn.len() > 1_024 || version_id.is_empty() || version_id.len() > 1_024 { + continue; + } + let suffix = format!("{}{}", rustfs_utils::http::SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX, arn); + rustfs_utils::http::insert_str(&mut transported.metadata, &suffix, version_id.clone()); + } + transported +} + impl SetDisks { /// `put_object` plus the destination key's previous current-version size, /// quorum-reduced from the dst `xl.meta` copies `rename_data` reads while @@ -3001,6 +3026,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { } #[tracing::instrument(skip(self))] async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, force_del_marker: bool) -> Result<()> { + let transported = delete_file_info_with_replication_transport_metadata(fi); + let fi = &transported; let disks = self.disk_inventory().await; let write_quorum = disks.len() / 2 + 1; let rollback_dir = Uuid::new_v4(); diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index a49063eec..da3c336c8 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -12,6 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::replication::{ + MAX_REPLICATION_TARGET_ARN_LEN, MAX_REPLICATION_TARGET_VERSION_ENTRIES, MAX_REPLICATION_TARGET_VERSION_ID_LEN, +}; use crate::{ ErasureAlgo, ErasureInfo, Error, FileInfo, FileInfoVersions, InlineData, NULL_VERSION_ID, ObjectPartInfo, RawFileInfo, ReplicationState, ReplicationStatusType, Result, VersionPurgeStatusType, is_restored_object_on_disk, @@ -25,13 +28,14 @@ use rustfs_utils::http::headers::{ }; use rustfs_utils::http::{ AMZ_BUCKET_REPLICATION_STATUS, MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX, SUFFIX_CRC, SUFFIX_DATA_MOV, SUFFIX_HEALING, - SUFFIX_PURGESTATUS, SUFFIX_REPLICA_STATUS, SUFFIX_REPLICA_TIMESTAMP, SUFFIX_REPLICATION_RESET, SUFFIX_REPLICATION_STATUS, - SUFFIX_REPLICATION_TIMESTAMP, SUFFIX_RESTORE_OPERATION_ID, contains_key_str, has_internal_suffix, insert_bytes, - is_internal_key, remove_bytes, + SUFFIX_PURGESTATUS, SUFFIX_REPLICA_STATUS, SUFFIX_REPLICA_TIMESTAMP, SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX, + SUFFIX_REPLICATION_RESET, SUFFIX_REPLICATION_STATUS, SUFFIX_REPLICATION_TIMESTAMP, SUFFIX_RESTORE_OPERATION_ID, + contains_key_str, has_internal_suffix, insert_bytes, is_internal_key, remove_bytes, }; use s3s::header::X_AMZ_RESTORE; use serde::{Deserialize, Serialize}; use std::cmp::Ordering; +use std::collections::BTreeMap; use std::convert::TryFrom; use std::hash::Hasher; use std::io::{Read, Write}; @@ -139,6 +143,58 @@ fn cmp_shallow_versions_for_order(a: &FileMetaShallowVersion, b: &FileMetaShallo /// the reset state, and a rustfs-only key was invisible to MinIO-compatible /// readers (backlog#799 B16). Normalize every entry to the canonical /// `replication-reset-` suffix and write both prefixes. +fn valid_target_delete_marker_version(arn: &str, version_id: &str) -> bool { + arn.starts_with("arn:") + && arn.len() <= MAX_REPLICATION_TARGET_ARN_LEN + && !version_id.is_empty() + && version_id.len() <= MAX_REPLICATION_TARGET_VERSION_ID_LEN +} + +/// Merge-only, never destructive. +/// +/// `ReplicationState::target_delete_marker_version_ids` is skipped by the +/// positional `FileInfo` wire form, so a delete that arrives over internode RPC +/// carries an empty map. Treating that as authoritative would let a remote disk +/// erase an exact target version that the local disk still holds. The key is +/// included in the quorum hash, so such a divergence does surface — but as a +/// quorum failure on an otherwise healthy object, which is not a state worth +/// reaching. Merge the RPC metadata carrier instead, and only ever insert. +fn persist_target_delete_marker_versions( + meta_sys: &mut HashMap>, + versions: &HashMap, + transport_metadata: &HashMap, +) { + let mut bounded = BTreeMap::new(); + // A corrupt carrier means the dual internal prefixes disagreed. Do not merge + // anything derived from it: this helper only ever inserts, so declining to + // merge leaves whatever durable keys the object already carries untouched, + // which is strictly safer than committing a mapping we cannot trust. + let (transport_versions, transport_corrupt) = rustfs_utils::http::target_delete_marker_versions(transport_metadata); + if transport_corrupt { + warn!("delete-marker target version transport metadata is inconsistent; leaving the persisted mapping unchanged"); + return; + } + for (arn, version_id) in transport_versions.iter().chain(versions.iter()) { + if !valid_target_delete_marker_version(arn, version_id) + || (bounded.len() >= MAX_REPLICATION_TARGET_VERSION_ENTRIES + && bounded.last_key_value().is_some_and(|(largest, _)| arn >= *largest)) + { + continue; + } + bounded.insert(arn, version_id); + if bounded.len() > MAX_REPLICATION_TARGET_VERSION_ENTRIES { + bounded.pop_last(); + } + } + for (arn, version_id) in bounded { + insert_bytes( + meta_sys, + &format!("{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}{arn}"), + version_id.as_bytes().to_vec(), + ); + } +} + fn persist_reset_statuses(meta_sys: &mut HashMap>, reset_statuses_map: &HashMap) { for (k, v) in reset_statuses_map { let suffix = k @@ -521,6 +577,11 @@ impl FileMeta { && let Some(state) = fi.replication_state_internal.as_ref() { persist_reset_statuses(&mut delete_marker.meta_sys, &state.reset_statuses_map); + persist_target_delete_marker_versions( + &mut delete_marker.meta_sys, + &state.target_delete_marker_version_ids, + &fi.metadata, + ); } } @@ -597,6 +658,11 @@ impl FileMeta { if let Some(state) = fi.replication_state_internal.as_ref() { persist_reset_statuses(&mut delete_marker.meta_sys, &state.reset_statuses_map); + persist_target_delete_marker_versions( + &mut delete_marker.meta_sys, + &state.target_delete_marker_version_ids, + &fi.metadata, + ); } } @@ -1278,6 +1344,90 @@ mod test { /// persisted under both internal prefixes, never as a bare ARN. A bare-ARN /// key (produced by `ObjectInfo::replication_state`) has no internal prefix, /// so read-back — which only recognizes prefixed keys — silently dropped it. + #[test] + fn persist_target_delete_marker_versions_uses_bounded_dual_prefixed_keys() { + let arn = "arn:rustfs:replication::target:bucket"; + let version_id = "opaque-target-version"; + let suffix = format!("{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}{arn}"); + let versions = HashMap::from([ + (arn.to_string(), version_id.to_string()), + ("not-an-arn".to_string(), "ignored".to_string()), + ("arn:too-long".to_string(), "x".repeat(MAX_REPLICATION_TARGET_VERSION_ID_LEN + 1)), + ]); + let mut meta_sys = HashMap::new(); + + persist_target_delete_marker_versions(&mut meta_sys, &versions, &HashMap::new()); + + assert_eq!( + meta_sys.get(&format!("{RUSTFS_INTERNAL_PREFIX}{suffix}")).map(Vec::as_slice), + Some(version_id.as_bytes()) + ); + assert_eq!( + meta_sys.get(&format!("{MINIO_INTERNAL_PREFIX}{suffix}")).map(Vec::as_slice), + Some(version_id.as_bytes()) + ); + assert_eq!(meta_sys.len(), 2, "invalid or oversized mappings must not expand xl.meta"); + } + + #[test] + fn persist_target_delete_marker_versions_caps_target_count() { + let versions = (0..=MAX_REPLICATION_TARGET_VERSION_ENTRIES) + .map(|index| (format!("arn:rustfs:replication::target:{index:04}"), format!("version-{index}"))) + .collect(); + let mut meta_sys = HashMap::new(); + + persist_target_delete_marker_versions(&mut meta_sys, &versions, &HashMap::new()); + + assert_eq!(meta_sys.len(), MAX_REPLICATION_TARGET_VERSION_ENTRIES * 2); + let excluded_suffix = format!( + "{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}arn:rustfs:replication::target:{MAX_REPLICATION_TARGET_VERSION_ENTRIES:04}" + ); + assert!(!meta_sys.contains_key(&format!("{RUSTFS_INTERNAL_PREFIX}{excluded_suffix}"))); + assert!(!meta_sys.contains_key(&format!("{MINIO_INTERNAL_PREFIX}{excluded_suffix}"))); + } + + #[test] + fn persist_target_delete_marker_versions_preserves_existing_targets_on_empty_update() { + let stale_suffix = format!("{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}arn:rustfs:replication::target:stale"); + let mut meta_sys = HashMap::from([ + (format!("{RUSTFS_INTERNAL_PREFIX}{stale_suffix}"), b"stale-rustfs".to_vec()), + (format!("{MINIO_INTERNAL_PREFIX}{stale_suffix}"), b"stale-minio".to_vec()), + ("unrelated".to_string(), b"kept".to_vec()), + ]); + + persist_target_delete_marker_versions(&mut meta_sys, &HashMap::new(), &HashMap::new()); + + assert_eq!( + meta_sys, + HashMap::from([ + (format!("{RUSTFS_INTERNAL_PREFIX}{stale_suffix}"), b"stale-rustfs".to_vec()), + (format!("{MINIO_INTERNAL_PREFIX}{stale_suffix}"), b"stale-minio".to_vec()), + ("unrelated".to_string(), b"kept".to_vec()), + ]) + ); + } + + #[test] + fn persist_target_delete_marker_versions_reads_rpc_transport_metadata() { + let arn = "arn:rustfs:replication::target:remote"; + let version_id = "opaque-remote-version"; + let suffix = format!("{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}{arn}"); + let mut transport_metadata = HashMap::new(); + rustfs_utils::http::insert_str(&mut transport_metadata, &suffix, version_id.to_string()); + let mut meta_sys = HashMap::new(); + + persist_target_delete_marker_versions(&mut meta_sys, &HashMap::new(), &transport_metadata); + + assert_eq!( + meta_sys.get(&format!("{RUSTFS_INTERNAL_PREFIX}{suffix}")).map(Vec::as_slice), + Some(version_id.as_bytes()) + ); + assert_eq!( + meta_sys.get(&format!("{MINIO_INTERNAL_PREFIX}{suffix}")).map(Vec::as_slice), + Some(version_id.as_bytes()) + ); + } + #[test] fn persist_reset_statuses_normalizes_to_dual_prefixed_keys() { let arn = "arn:rustfs:replication::target:bucket"; diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index 091bfb527..28973452f 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -33,7 +33,7 @@ use rustfs_utils::http::{ SUFFIX_TIER_FV_MARKER, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID, SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, contains_key_bytes, get_bytes, get_consistent_bytes, get_str, has_internal_suffix, insert_bytes, is_internal_key, remove_bytes, - strip_internal_prefix, + strip_internal_prefix, target_delete_marker_versions, }; const MSGPACK_EXT8: u8 = 0xc7; @@ -2725,6 +2725,15 @@ fn get_internal_replication_state(metadata: &HashMap) -> Option< } } + // Re-derive the per-target delete-marker versions from the durable keys. + // `ReplicationState` skips this map on the wire, so the metadata is the only + // authority. `corrupt` means the dual internal prefixes disagreed: surface it + // rather than guessing, so callers fail closed instead of purging the wrong + // target version. + (rs.target_delete_marker_version_ids, rs.target_delete_marker_version_ids_corrupt) = target_delete_marker_versions(metadata); + has |= !rs.target_delete_marker_version_ids.is_empty(); + has |= rs.target_delete_marker_version_ids_corrupt; + if has { Some(rs) } else { None } } @@ -4725,6 +4734,47 @@ mod tests { assert!(dm.mod_time.is_some()); } + #[test] + fn target_delete_marker_version_metadata_is_forward_and_backward_compatible() { + let arn = "arn:rustfs:replication:us-east-1:target:bucket"; + let suffix = format!("{}{arn}", rustfs_utils::http::SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX); + let mut metadata = HashMap::from([(format!("{RUSTFS_INTERNAL_PREFIX}replication-status"), format!("{arn}=COMPLETED;"))]); + + let old = get_internal_replication_state(&metadata).expect("legacy replication metadata should parse"); + assert!( + old.target_delete_marker_version_ids.is_empty(), + "a new reader must treat the missing legacy field as empty" + ); + + metadata.insert(format!("{}{suffix}", rustfs_utils::http::MINIO_INTERNAL_PREFIX), "target-id".to_string()); + metadata.insert(format!("{RUSTFS_INTERNAL_PREFIX}{suffix}"), "target-id".to_string()); + let current = get_internal_replication_state(&metadata).expect("new replication metadata should parse"); + assert_eq!( + current.target_delete_marker_version_ids.get(arn).map(String::as_str), + Some("target-id"), + "matching dual-prefix values must remain readable" + ); + assert_eq!( + current.targets.get(arn), + Some(&ReplicationStatusType::Completed), + "the added metadata key must not disturb fields understood by old nodes" + ); + + metadata.insert( + format!("{}{suffix}", rustfs_utils::http::MINIO_INTERNAL_PREFIX), + "conflicting-id".to_string(), + ); + let conflicted = get_internal_replication_state(&metadata).expect("replication status should still parse"); + assert!( + conflicted.target_delete_marker_version_ids.is_empty(), + "a destructive version ID must fail closed when the dual prefixes disagree" + ); + assert!( + conflicted.target_delete_marker_version_ids_corrupt, + "a dual-prefix conflict must remain distinguishable from missing legacy metadata" + ); + } + #[test] fn get_internal_replication_state_keeps_canonical_reset_key() { // The reset status is stored on disk under the full internal key; parsing diff --git a/crates/filemeta/src/replication.rs b/crates/filemeta/src/replication.rs index 8c822df40..3b008bc4b 100644 --- a/crates/filemeta/src/replication.rs +++ b/crates/filemeta/src/replication.rs @@ -227,6 +227,13 @@ impl From<&str> for ReplicationType { } /// ReplicationState represents internal replication state +/// Bounds on the per-target delete-marker version map. The map is rebuilt from +/// attacker-influenced object metadata, so cap the entry count and both string +/// lengths rather than trusting what was persisted. +pub(crate) const MAX_REPLICATION_TARGET_VERSION_ENTRIES: usize = 1_000; +pub(crate) const MAX_REPLICATION_TARGET_ARN_LEN: usize = 1_024; +pub(crate) const MAX_REPLICATION_TARGET_VERSION_ID_LEN: usize = 1_024; + #[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)] pub struct ReplicationState { pub replica_timestamp: Option, @@ -239,6 +246,16 @@ pub struct ReplicationState { pub targets: HashMap, pub purge_targets: HashMap, pub reset_statuses_map: HashMap, + /// Exact version id the delete marker got on each replication target, keyed + /// by target ARN. Skipped by serde: `ReplicationState` has a positional wire + /// form, so this travels in the object's internal metadata instead and is + /// re-derived on read. See `persist_target_delete_marker_versions`. + #[serde(skip)] + pub target_delete_marker_version_ids: HashMap, + /// Set when the persisted keys disagreed across the dual internal prefixes, + /// so callers fail closed instead of purging the wrong target version. + #[serde(skip)] + pub target_delete_marker_version_ids_corrupt: bool, } impl ReplicationState { @@ -311,6 +328,7 @@ impl ReplicationState { arn: arn.to_string(), prev_replication_status: self.targets.get(arn).cloned().unwrap_or_default(), version_purge_status: self.purge_targets.get(arn).cloned().unwrap_or_default(), + target_delete_marker_version_id: self.target_delete_marker_version_ids.get(arn).cloned(), resync_timestamp, ..Default::default() } @@ -414,6 +432,8 @@ pub struct ReplicatedTargetInfo { pub endpoint: String, pub secure: bool, pub error: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub target_delete_marker_version_id: Option, } impl ReplicatedTargetInfo { @@ -926,6 +946,36 @@ pub fn get_replication_state(rinfos: &ReplicatedInfos, prev_state: &ReplicationS reset_statuses_map.insert(key, value); } + // Carry the previously recorded per-target delete-marker versions forward, + // dropping anything that no longer satisfies the bounds, then fold in the + // versions this round's targets reported. A map that has already grown past + // the cap is discarded rather than trusted. + let mut target_delete_marker_version_ids = prev_state.target_delete_marker_version_ids.clone(); + target_delete_marker_version_ids.retain(|arn, version_id| { + !arn.is_empty() + && arn.len() <= MAX_REPLICATION_TARGET_ARN_LEN + && !version_id.is_empty() + && version_id.len() <= MAX_REPLICATION_TARGET_VERSION_ID_LEN + }); + if target_delete_marker_version_ids.len() > MAX_REPLICATION_TARGET_VERSION_ENTRIES { + target_delete_marker_version_ids.clear(); + } + for target in &rinfos.targets { + let Some(version_id) = target.target_delete_marker_version_id.as_ref() else { + continue; + }; + if (!target_delete_marker_version_ids.contains_key(&target.arn) + && target_delete_marker_version_ids.len() >= MAX_REPLICATION_TARGET_VERSION_ENTRIES) + || target.arn.is_empty() + || target.arn.len() > MAX_REPLICATION_TARGET_ARN_LEN + || version_id.is_empty() + || version_id.len() > MAX_REPLICATION_TARGET_VERSION_ID_LEN + { + continue; + } + target_delete_marker_version_ids.insert(target.arn.clone(), version_id.clone()); + } + ReplicationState { replicate_decision_str: prev_state.replicate_decision_str.clone(), reset_statuses_map, @@ -936,6 +986,8 @@ pub fn get_replication_state(rinfos: &ReplicatedInfos, prev_state: &ReplicationS replication_timestamp: rinfos.replication_timestamp, purge_targets, version_purge_status_internal: vpurge_statuses, + target_delete_marker_version_ids, + target_delete_marker_version_ids_corrupt: prev_state.target_delete_marker_version_ids_corrupt, ..Default::default() } diff --git a/crates/replication/src/filemeta.rs b/crates/replication/src/filemeta.rs index 2605353c9..7d3a8593f 100644 --- a/crates/replication/src/filemeta.rs +++ b/crates/replication/src/filemeta.rs @@ -239,6 +239,13 @@ pub struct ReplicationState { pub targets: HashMap, pub purge_targets: HashMap, pub reset_statuses_map: HashMap, + /// Skipped by serde: this state has a positional wire form, so the map + /// travels in the object's internal metadata and is re-derived on read. + /// Kept in step with the filemeta crate's copy of the same state. + #[serde(skip)] + pub target_delete_marker_version_ids: HashMap, + #[serde(skip)] + pub target_delete_marker_version_ids_corrupt: bool, } impl ReplicationState { @@ -311,6 +318,7 @@ impl ReplicationState { arn: arn.to_string(), prev_replication_status: self.targets.get(arn).cloned().unwrap_or_default(), version_purge_status: self.purge_targets.get(arn).cloned().unwrap_or_default(), + target_delete_marker_version_id: self.target_delete_marker_version_ids.get(arn).cloned(), resync_timestamp, ..Default::default() } @@ -414,6 +422,9 @@ pub struct ReplicatedTargetInfo { pub endpoint: String, pub secure: bool, pub error: Option, + /// Version the target assigned to the delete marker it just created. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub target_delete_marker_version_id: Option, } impl ReplicatedTargetInfo { @@ -1022,6 +1033,11 @@ fn version_purge_statuses_string(targets: &HashMap) -> ReplicationState { let reset_status_map: Vec<(String, String)> = rinfos .targets @@ -1048,6 +1064,35 @@ pub fn get_replication_state(rinfos: &ReplicatedInfos, prev_state: &ReplicationS reset_statuses_map.insert(key, value); } + // Carry forward the recorded per-target delete-marker versions, dropping + // anything outside the bounds, then fold in what this round's targets + // reported. A map already past the cap is discarded rather than trusted. + let mut target_delete_marker_version_ids = prev_state.target_delete_marker_version_ids.clone(); + target_delete_marker_version_ids.retain(|arn, version_id| { + !arn.is_empty() + && arn.len() <= MAX_REPLICATION_TARGET_ARN_LEN + && !version_id.is_empty() + && version_id.len() <= MAX_REPLICATION_TARGET_VERSION_ID_LEN + }); + if target_delete_marker_version_ids.len() > MAX_REPLICATION_TARGET_VERSION_ENTRIES { + target_delete_marker_version_ids.clear(); + } + for target in &rinfos.targets { + let Some(version_id) = target.target_delete_marker_version_id.as_ref() else { + continue; + }; + if (!target_delete_marker_version_ids.contains_key(&target.arn) + && target_delete_marker_version_ids.len() >= MAX_REPLICATION_TARGET_VERSION_ENTRIES) + || target.arn.is_empty() + || target.arn.len() > MAX_REPLICATION_TARGET_ARN_LEN + || version_id.is_empty() + || version_id.len() > MAX_REPLICATION_TARGET_VERSION_ID_LEN + { + continue; + } + target_delete_marker_version_ids.insert(target.arn.clone(), version_id.clone()); + } + ReplicationState { replicate_decision_str: prev_state.replicate_decision_str.clone(), reset_statuses_map, @@ -1058,6 +1103,8 @@ pub fn get_replication_state(rinfos: &ReplicatedInfos, prev_state: &ReplicationS replication_timestamp: rinfos.replication_timestamp, purge_targets, version_purge_status_internal: vpurge_statuses, + target_delete_marker_version_ids, + target_delete_marker_version_ids_corrupt: prev_state.target_delete_marker_version_ids_corrupt, ..Default::default() } diff --git a/crates/utils/src/http/metadata_compat.rs b/crates/utils/src/http/metadata_compat.rs index ab680c953..c99d8ed00 100644 --- a/crates/utils/src/http/metadata_compat.rs +++ b/crates/utils/src/http/metadata_compat.rs @@ -15,7 +15,7 @@ //! System metadata compatibility: write both x-rustfs-internal-* and x-minio-internal-* //! for MinIO interoperability. Read prefers RustFS, fallback to MinIO. -use std::collections::HashMap; +use std::collections::{BTreeMap, HashMap}; pub const RUSTFS_INTERNAL_PREFIX: &str = "x-rustfs-internal-"; pub const MINIO_INTERNAL_PREFIX: &str = "x-minio-internal-"; @@ -58,6 +58,9 @@ pub const SUFFIX_TIER_FV_ID: &str = "tier-free-versionID"; pub const SUFFIX_TIER_FV_MARKER: &str = "tier-free-marker"; pub const SUFFIX_TIER_SKIP_FV_ID: &str = "tier-skip-fvid"; +/// Per-target delete-marker version ids are stored one key per target ARN. +pub const SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX: &str = "replication-delete-marker-version-"; + /// Case-insensitive (ASCII) check that `s` begins with `prefix`. Equivalent to /// `s.to_lowercase().starts_with(prefix)` when `prefix` is ASCII (as both internal prefixes are), /// but without allocating. @@ -247,6 +250,74 @@ pub fn remove_bytes(map: &mut HashMap>, suffix: &str) { with_internal_key(MINIO_INTERNAL_PREFIX, suffix, |k2| map.remove(k2)); } +/// Strips an internal metadata prefix while preserving the suffix casing. +pub fn strip_internal_prefix_preserving_case(key: &str) -> Option<&str> { + if starts_with_ignore_ascii_case(key, RUSTFS_INTERNAL_PREFIX) { + key.get(RUSTFS_INTERNAL_PREFIX.len()..) + } else if starts_with_ignore_ascii_case(key, MINIO_INTERNAL_PREFIX) { + key.get(MINIO_INTERNAL_PREFIX.len()..) + } else { + None + } +} + +/// Reads the bounded per-target delete-marker version map in one metadata scan. +/// The boolean is set when matching metadata is malformed or compatibility keys disagree. +pub fn target_delete_marker_versions(map: &HashMap) -> (HashMap, bool) { + const MAX_ENTRIES: usize = 1_000; + const MAX_ARN_LEN: usize = 1_024; + const MAX_VERSION_ID_LEN: usize = 1_024; + + let mut versions = BTreeMap::>::new(); + let mut corrupt = false; + for (key, value) in map { + let Some(suffix) = strip_internal_prefix_preserving_case(key) else { + continue; + }; + let Some(prefix) = suffix.get(..SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX.len()) else { + continue; + }; + if !prefix.eq_ignore_ascii_case(SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX) { + continue; + } + let arn = &suffix[SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX.len()..]; + if !arn.starts_with("arn:") || arn.len() > MAX_ARN_LEN || value.is_empty() || value.len() > MAX_VERSION_ID_LEN { + corrupt = true; + continue; + } + match versions.entry(arn.to_string()) { + std::collections::btree_map::Entry::Vacant(entry) => { + entry.insert(Some(value.clone())); + } + std::collections::btree_map::Entry::Occupied(mut entry) => { + if entry.get().as_deref() != Some(value.as_str()) { + entry.insert(None); + corrupt = true; + } + } + } + } + + // Apply the cap after collecting, never during. Capping mid-iteration made + // the surviving subset depend on `HashMap` order, so two disks decoding the + // same metadata could keep different entries and hash differently — turning + // an over-cap object into a quorum failure instead of a reported corruption. + // `BTreeMap` order is total, so truncating here is identical everywhere. + if versions.len() > MAX_ENTRIES { + corrupt = true; + let retained = versions.keys().take(MAX_ENTRIES).cloned().collect::>(); + versions.retain(|arn, _| retained.binary_search(arn).is_ok()); + } + + ( + versions + .into_iter() + .filter_map(|(arn, version_id)| version_id.map(|version_id| (arn, version_id))) + .collect(), + corrupt, + ) +} + #[cfg(test)] mod tests { use super::*; @@ -478,4 +549,59 @@ mod tests { remove_bytes(&mut meta_sys, &long_suffix); assert!(!contains_key_bytes(&meta_sys, &long_suffix)); } + + #[test] + fn target_delete_marker_versions_preserve_arn_case_and_report_conflicts() { + let arn = "arn:rustfs:replication::Target:Bucket"; + let suffix = format!("{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}{arn}"); + let mut metadata = HashMap::new(); + insert_str(&mut metadata, &suffix, "target-version".to_string()); + + let (versions, corrupt) = target_delete_marker_versions(&metadata); + assert_eq!(versions.get(arn).map(String::as_str), Some("target-version")); + assert!(!corrupt); + + metadata.insert(format!("{MINIO_INTERNAL_PREFIX}{suffix}"), "other-version".to_string()); + let (versions, corrupt) = target_delete_marker_versions(&metadata); + assert!(versions.is_empty()); + assert!(corrupt); + } + + #[test] + fn target_delete_marker_versions_bound_distinct_entries_during_scan() { + let metadata = (0..=1_000) + .map(|index| { + ( + format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}arn:target:{index:04}"), + format!("version-{index}"), + ) + }) + .collect(); + + let (versions, corrupt) = target_delete_marker_versions(&metadata); + + assert_eq!(versions.len(), 1_000); + assert!(corrupt); + } + #[test] + fn target_delete_marker_versions_cap_is_deterministic_across_decodes() { + // Two decodes of the same oversized metadata must agree, or the two disks + // holding it hash differently and the object loses quorum instead of + // reporting corruption. + let mut metadata = HashMap::new(); + for index in 0..1_050 { + metadata.insert( + format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX}arn:target:{index:05}"), + format!("version-{index}"), + ); + } + + let (first, first_corrupt) = target_delete_marker_versions(&metadata); + let (second, second_corrupt) = target_delete_marker_versions(&metadata); + + assert!(first_corrupt, "exceeding the cap must be reported as corrupt"); + assert_eq!(first_corrupt, second_corrupt); + assert_eq!(first, second, "the retained subset must not depend on map iteration order"); + assert_eq!(first.len(), 1_000); + } }