mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-05 04:47:43 +00:00
feat(replication): purge delete markers by the target's own version id (#5676)
* feat(replication): purge delete markers by the target's own version id When a delete marker is replicated, the target assigns it a version id. The purge that follows derived one from the *source* uuid instead, which is only correct when the target mirrors source version ids. A generic S3 target does not: the derived id addresses a version that does not exist there, so the purge is a no-op and the replica keeps a marker the source has already removed. Same failure class as #4401. Record the id the target reports and address it directly on purge. Data path, all of it driven by the object's internal metadata rather than the `ReplicationState` wire form, which encodes positionally and cannot carry a map: - `rustfs-utils`: the `replication-delete-marker-version-<arn>` key family, plus `strip_internal_prefix_preserving_case` — ARNs are case-sensitive and the existing `strip_internal_prefix` lowercases. - `ReplicationState` gains the map and a `..._corrupt` flag, both `#[serde(skip)]`; `ReplicatedTargetInfo` carries the per-target id. - `persist_target_delete_marker_versions` is merge-only. A delete arriving over internode RPC has an empty map, so treating it as authoritative would let a remote disk erase an id the local disk still holds. - `delete_object_version` copies the map into `fi.metadata` before dispatch, so the durable carrier crosses the wire even though the field does not. - The keys are folded into the quorum hash through their normalized form: the dual internal prefixes carrying one mapping share an identity, while a genuine disagreement between disks still shows up as a quorum difference. - `corrupt` (the prefixes disagreed) fails closed: skip the purge and warn rather than guess an id and risk destroying a live version on the target. Ported from the rc.1 branch, which cannot merge as a whole: its MRF replay rewrite collides with #5659/#5671/#5672/#5673 and regressed `MRF_PENDING_CAP`. main's MRF machinery is kept; only this capability moves across. It touches no MRF code. Two things did not survive the port, deliberately. The branch's `missing_is_complete` purge regression does not exist here — it came from its own HEAD-precheck rewrite, and main's simpler path never had it. And the branch's `MrfReplicateEntry` ordering fields are MRF-redesign scope, left behind. Verification: cargo fmt --all --check, git diff --check, cargo check --workspace --all-targets, and the suites for the four touched crates — 4070 tests, 2 pre-existing failures unrelated to this change (`system_resolver_negative_result_reaches_the_dns_allowlist`, `test_resolve_domain_preserves_system_resolver_error_provenance`; both are the sandbox DNS interception, they fail on a clean checkout too). * fix(replication): keep the layer guard happy scripts/check_architecture_migration_rules.sh matches on text, so the doc comments naming `rustfs_filemeta::` read as a cross-layer dependency even though nothing imports it. Reword them; the guard passes. * fix(replication): make the target-version cap deterministic Two defects in this PR, both found in review. The cap was applied while iterating a `HashMap`, so *which* 1000 entries survived depended on iteration order. Two disks decoding the same oversized metadata could keep different subsets, hash differently, and lose quorum — instead of both reporting the same corruption. Collect first, then truncate in `BTreeMap` order, which is total and identical everywhere. And `persist_target_delete_marker_versions` discarded the `corrupt` flag from the RPC carrier, committing a delete-marker update that looked clean while the exact remote marker identity was unknown. It now declines to merge a corrupt carrier. Because the helper only ever inserts, declining leaves the durable keys already on the object untouched, which is strictly safer than writing a mapping we cannot trust. Residual, stated rather than papered over: corruption confined to the RPC carrier is not persisted as a sentinel, so a later reader of an object that carried no durable keys still sees "legacy, no mapping" rather than "corrupt". Persisting that would need a wire-format addition; the consumer already fails closed on any corruption it can observe. New test: `target_delete_marker_versions_cap_is_deterministic_across_decodes` decodes the same 1050-entry map twice and asserts both the corrupt flag and the retained subset agree. * fix(replication): preserve multipart source mtime (#5669) * fix(kms): repair unopenable ciphertext and cover the Vault backends (#5668) * Add black-box behavior tests for KMS resilience and serialization * fix(kms): repair unopenable ciphertext across backends Black-box testing of the KMS crate surfaced several defects that make encrypted data permanently unreadable. Symmetric envelopes. The Local and Vault Transit backends returned raw cipher output from `encrypt` while `decrypt` parsed a JSON envelope, so anything sealed through the master-key path could never be opened again. Local also discarded the AES-GCM nonce. Both now emit the same envelope `decrypt` consumes, matching the Static backend. Deterministic AAD. The object layer derived AEAD additional data by serializing a `HashMap` directly. Iteration order differs per instance, so a context rebuilt from storage produced different AAD bytes than the one used to seal and the object stopped opening. Ordering by key removes that dependency, matching the Static backend's existing `context_aad`. Objects written with the default single-key context are unaffected, since a one-entry map has only one serialization. Cipher in the header projection. `metadata_to_headers` recorded the SSE mode (`AES256` / `aws:kms`), which cannot represent ChaCha20-Poly1305, so a ChaCha-sealed object came back claiming `aws:kms` and was opened with the wrong cipher. The cipher now travels in `x-rustfs-encryption-algorithm` — the header the storage layer already reads but nothing ever wrote. Objects without it fall back as before. Also: the Static backend ignored `key_spec` and always issued 256-bit data keys; Local `list_keys` hardcoded `truncated: false`, ignored `marker`, and paginated over unordered `read_dir`, so a paginating client silently saw a partial key list; and Local and Vault KV2 reported `key_id: "unknown"` from `decrypt` despite the envelope naming the master key. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * test(kms): cover both Vault backends and key rotation The behavior suite ran only against Local and Static, and its own harness documented the gap: the Vault backends had no business-capability coverage at all. Setting `RUSTFS_KMS_VAULT_TOKEN` now adds Vault KV2 and Vault Transit to every `for_each_backend` spec against a live server. That lane is what surfaced the Transit envelope defect fixed in the previous commit. `rotate` and `versioning` are advertised only by the Vault backends, so until now every capability-gated branch for them took the `UnsupportedCapability` side and the working half was never asserted — a rotation that dropped prior key versions would have gone green. The new `behavior_rotation.rs` pins that half: material sealed before a rotation still opens after it, repeated rotations accumulate versions rather than overwriting a single spare, and the history survives a restart. Two test defects fixed. `objects_round_trip_across_sizes_and_algorithms` asserted a 1-byte object differs from its own ciphertext, which collides once every 256 runs; the assertion now applies only where a collision is not realistic, and small objects stay covered by the tag check and the decrypt round-trip. `test_from_env_selects_token_file` depended on `RUSTFS_KMS_VAULT_TOKEN` being absent from the caller's environment and now clears it explicitly. The snapshots directory was also removed from `.gitignore`: insta snapshots are the assertions themselves, so leaving them untracked gives CI nothing to compare against. Only `.snap.new` scratch files are ignored now. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * test(kms): adapt behavior suite to current key APIs Rebasing onto main brought four API changes the suite predates. `DeleteKeyRequest` gained `confirm_key_id`, and immediate deletion is now gated on the server's `allow_immediate_deletion`. Scheduled deletions pass `None`; the four specs that destroy a key outright echo the key id back and opt the harness config in, which is what the gate asks of a real caller. `LocalBackupExportRequest` gained `sanitized_config`. These specs cover the key-material path, so they seal no configuration and pass `None`. `KmsCacheStats` became a named struct with real hit, miss, and eviction counters. `cache_stats_returns_an_entry_count_and_no_hit_or_miss_data` existed to pin the old placeholder behavior — that the second tuple element was always zero — which main has since fixed, so it is now `cache_stats_reports_hits_and_misses_separately` and asserts the counters actually move. Starting the service provisions the reserved probe key, so it shows up in listings and backup bundles. Exact-set assertions filter it through a new `without_probe_key` helper rather than naming it, keeping those specs about the keys they seeded. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix(kms): bind the AAD to the stored context bytes Review caught that canonicalizing the AAD on decrypt breaks objects sealed before canonicalization existed, and it was right. The AAD is the *serialization* of the encryption context, and `x-rustfs-encryption-context` stores that exact byte sequence: `encrypt_object` fed one `HashMap` to the AEAD and then moved the same map into the metadata the header is written from, so the stored string is byte-identical to the AAD the object was sealed under. Those objects are therefore recoverable — but only while nothing round-trips the value through a `HashMap` and re-serializes it. Recomputing sorted AAD on decrypt would have turned a readable object into a permanently unreadable one. The previous behavior was worse than the first analysis credited: it did not merely fail intermittently, it made the failure deterministic. `EncryptionMetadata` now carries `context_aad`, the bytes the object was actually sealed with. Encryption records what it fed the AEAD, the header projection stores those bytes verbatim (and preserves a legacy ordering across a re-projection rather than rewriting it into sorted form), and `headers_to_metadata` carries the stored string through untouched. Both decrypt paths, SSE-KMS and SSE-C, prefer it and fall back to canonical serialization only when no stored serialization exists. Canonicalization still applies to everything newly sealed, so the original ordering bug cannot recur. Two tests pin this: a legacy record whose sealed bytes are non-canonical must survive a full header round trip unchanged, and a context header rewritten to an equivalent-but-reordered serialization must fail authentication rather than silently re-deriving a working AAD. Both were mutation-checked against the reinstated bug on each side. Also from review: the lifecycle churn test asserted only that every request was accounted for, which holds whether the state gate exists or not, so both branches are now pinned deterministically after the churn (asserting `refused > 0` on the concurrent phase would only trade the hole for a scheduling flake). And the Local and Vault KV2 envelopes compare `encryption_context` without authenticating it — `DekCrypto` seals only the plaintext — which is now documented at both sites; closing it needs a versioned envelope, since existing ciphertext was sealed without AAD. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: ccccpj <ccccpj@outlook.com> Co-authored-by: 唐小鸭 <tangtang1251@qq.com> Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -2000,7 +2000,7 @@ impl TargetClient {
|
||||
object: &str,
|
||||
version_id: Option<String>,
|
||||
opts: RemoveObjectOptions,
|
||||
) -> Result<(), S3ClientError> {
|
||||
) -> Result<Option<String>, 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();
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1634,6 +1634,29 @@ async fn source_delete_marker_missing<S: EcstoreObjectOperations>(
|
||||
}
|
||||
}
|
||||
|
||||
/// 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<Option<String>> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<String, String>) {
|
||||
@@ -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<String, String>) {
|
||||
let (versions, corrupt) = http::target_delete_marker_versions(metadata);
|
||||
hasher.update([u8::from(corrupt)]);
|
||||
let mut versions = versions.iter().collect::<Vec<_>>();
|
||||
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"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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-<arn>` 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<String, Vec<u8>>,
|
||||
versions: &HashMap<String, String>,
|
||||
transport_metadata: &HashMap<String, String>,
|
||||
) {
|
||||
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<String, Vec<u8>>, reset_statuses_map: &HashMap<String, String>) {
|
||||
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";
|
||||
|
||||
@@ -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<String, String>) -> 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
|
||||
|
||||
@@ -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<OffsetDateTime>,
|
||||
@@ -239,6 +246,16 @@ pub struct ReplicationState {
|
||||
pub targets: HashMap<String, ReplicationStatusType>,
|
||||
pub purge_targets: HashMap<String, VersionPurgeStatusType>,
|
||||
pub reset_statuses_map: HashMap<String, String>,
|
||||
/// 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<String, String>,
|
||||
/// 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<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub target_delete_marker_version_id: Option<String>,
|
||||
}
|
||||
|
||||
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()
|
||||
}
|
||||
|
||||
@@ -239,6 +239,13 @@ pub struct ReplicationState {
|
||||
pub targets: HashMap<String, ReplicationStatusType>,
|
||||
pub purge_targets: HashMap<String, VersionPurgeStatusType>,
|
||||
pub reset_statuses_map: HashMap<String, String>,
|
||||
/// 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<String, String>,
|
||||
#[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<String>,
|
||||
/// 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<String>,
|
||||
}
|
||||
|
||||
impl ReplicatedTargetInfo {
|
||||
@@ -1022,6 +1033,11 @@ fn version_purge_statuses_string(targets: &HashMap<String, VersionPurgeStatusTyp
|
||||
if result.is_empty() { None } else { Some(result) }
|
||||
}
|
||||
|
||||
/// Kept in step with the bounds used by the filemeta crate's copy.
|
||||
const MAX_REPLICATION_TARGET_VERSION_ENTRIES: usize = 1_000;
|
||||
const MAX_REPLICATION_TARGET_ARN_LEN: usize = 1_024;
|
||||
const MAX_REPLICATION_TARGET_VERSION_ID_LEN: usize = 1_024;
|
||||
|
||||
pub fn get_replication_state(rinfos: &ReplicatedInfos, prev_state: &ReplicationState, _vid: Option<String>) -> 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()
|
||||
}
|
||||
|
||||
@@ -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<String, Vec<u8>>, 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<String, String>) -> (HashMap<String, String>, 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::<String, Option<String>>::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::<Vec<_>>();
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user