mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-18 02:33:15 +00:00
fix(site-replication): merge incoming ILM expiry documents instead of overwriting
The lc-config receiver replaced the whole local lifecycle config with the peer's document (and deleted it wholesale on peer delete), so an expiry-only update erased the receiver's local tier/transition rules, and a peer's transition rules were installed across sites (backlog#1675 P1-1). Receiver (apply_bucket_meta_item): - lc-config now merges via merge_incoming_lifecycle_config, mirroring MinIO's mergeWithCurrentLCConfig with a trust-boundary hardening: incoming transition fields are discarded outright; the local transition side of a same-id rule is authoritative. A peer delete merges with the empty set — pure-expiry rules go away, transition rules survive with their expiry side cleared, and only an empty result deletes the config file. - Staleness moves to the expiry axis (config.expiry_updated_at): lifecycle_config_updated_at also moves on local transition-only edits, which shadowed newer peer expiry updates. - Receiver-side replicateILMExpiry gate, symmetric with the sender hook (previously any peer could install expiry rules while the option was off). - Rule order is deterministic (local order, incoming-new appended), so re-delivering the same document is byte-stable and does not rewrite bucket metadata per broadcast. Sender: - Both admin choke points — the bucket-meta hook and the SRInfo bucket entry feeding bootstrap/repair and consistency views — now emit only the expiry subset (transition fields stripped, non-expiry rules dropped). MinIO receivers install incoming rules verbatim, so transition rules must never leave the site. An unparseable local config is forwarded unfiltered rather than degraded to a delete. Not covered here (follow-up): a two-site e2e with a real tier backend to exercise transition-rule preservation end to end; receiver-side validate_transition_tier for merged configs.
This commit is contained in:
@@ -4220,13 +4220,23 @@ pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: boo
|
|||||||
broadcast_site_replication_json(&path, &serde_json::json!({})).await
|
broadcast_site_replication_json(&path, &serde_json::json!({})).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn site_replication_bucket_meta_hook(item: SRBucketMeta) -> S3Result<()> {
|
pub async fn site_replication_bucket_meta_hook(mut item: SRBucketMeta) -> S3Result<()> {
|
||||||
let Some(runtime) = runtime_site_replication_targets().await? else {
|
let Some(runtime) = runtime_site_replication_targets().await? else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
};
|
};
|
||||||
if item.r#type == "lc-config" && !site_replication_state_replicates_ilm_expiry(&runtime.state) {
|
if item.r#type == "lc-config" && !site_replication_state_replicates_ilm_expiry(&runtime.state) {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
if item.r#type == "lc-config" {
|
||||||
|
// Only the expiry subset travels (MinIO peers install incoming rules
|
||||||
|
// verbatim, so transition rules must never leave this site). An empty
|
||||||
|
// subset becomes a delete, which the receiver merges with the empty
|
||||||
|
// set — local transition rules there survive.
|
||||||
|
item.expiry_lc_config = item
|
||||||
|
.expiry_lc_config
|
||||||
|
.and_then(|raw| lifecycle_expiry_subset_xml(raw.as_bytes()))
|
||||||
|
.map(|data| String::from_utf8_lossy(&data).into_owned());
|
||||||
|
}
|
||||||
broadcast_site_replication_json_with_runtime(
|
broadcast_site_replication_json_with_runtime(
|
||||||
&runtime,
|
&runtime,
|
||||||
"/rustfs/admin/v3/site-replication/peer/bucket-meta",
|
"/rustfs/admin/v3/site-replication/peer/bucket-meta",
|
||||||
@@ -4319,7 +4329,12 @@ async fn build_sr_info(state: &SiteReplicationState, local_peer: &PeerInfo) -> S
|
|||||||
entry.sse_config = raw_config_to_base64(&metadata.encryption_config_xml);
|
entry.sse_config = raw_config_to_base64(&metadata.encryption_config_xml);
|
||||||
entry.replication_config = raw_config_to_base64(&metadata.replication_config_xml);
|
entry.replication_config = raw_config_to_base64(&metadata.replication_config_xml);
|
||||||
entry.quota_config = raw_config_to_base64(&metadata.quota_config_json);
|
entry.quota_config = raw_config_to_base64(&metadata.quota_config_json);
|
||||||
entry.expiry_lc_config = raw_config_to_base64(&metadata.lifecycle_config_xml);
|
// Expiry subset only: this entry feeds both the bootstrap/repair
|
||||||
|
// plan (peers must not receive transition rules) and cross-site
|
||||||
|
// consistency views (transition rules are site-local and would
|
||||||
|
// read as false mismatches).
|
||||||
|
entry.expiry_lc_config =
|
||||||
|
lifecycle_expiry_subset_xml(&metadata.lifecycle_config_xml).and_then(|data| raw_config_to_base64(&data));
|
||||||
entry.cors_config = raw_config_to_base64(&metadata.cors_config_xml);
|
entry.cors_config = raw_config_to_base64(&metadata.cors_config_xml);
|
||||||
entry.policy_updated_at = maybe_time(metadata.policy_config_updated_at);
|
entry.policy_updated_at = maybe_time(metadata.policy_config_updated_at);
|
||||||
entry.tag_config_updated_at = maybe_time(metadata.tagging_config_updated_at);
|
entry.tag_config_updated_at = maybe_time(metadata.tagging_config_updated_at);
|
||||||
@@ -6859,8 +6874,128 @@ fn merge_incoming_lifecycle_config(
|
|||||||
local: Option<s3s::dto::BucketLifecycleConfiguration>,
|
local: Option<s3s::dto::BucketLifecycleConfiguration>,
|
||||||
updated_at: Option<OffsetDateTime>,
|
updated_at: Option<OffsetDateTime>,
|
||||||
) -> Option<s3s::dto::BucketLifecycleConfiguration> {
|
) -> Option<s3s::dto::BucketLifecycleConfiguration> {
|
||||||
let _ = (local, updated_at);
|
// Incoming rules reduced to their expiry side (trust boundary). Rules
|
||||||
incoming
|
// that carry no expiry semantics after the strip — transition-only or
|
||||||
|
// abort-multipart-only — are site-local and are not installed.
|
||||||
|
let mut incoming_by_id: HashMap<String, s3s::dto::LifecycleRule> = HashMap::new();
|
||||||
|
let mut incoming_order: Vec<String> = Vec::new();
|
||||||
|
for mut rule in incoming.into_iter().flat_map(|config| config.rules) {
|
||||||
|
rule.transitions = None;
|
||||||
|
rule.noncurrent_version_transitions = None;
|
||||||
|
if !lifecycle_rule_has_expiry(&rule) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
let id = rule.id.clone().unwrap_or_default();
|
||||||
|
if incoming_by_id.insert(id.clone(), rule).is_none() {
|
||||||
|
incoming_order.push(id);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Local order first, incoming-new appended: repeated delivery of the same
|
||||||
|
// document is byte-stable, so bucket metadata is written once, not on
|
||||||
|
// every broadcast.
|
||||||
|
let local_expiry_updated_at = local.as_ref().and_then(|config| config.expiry_updated_at.clone());
|
||||||
|
let mut rules: Vec<s3s::dto::LifecycleRule> = Vec::new();
|
||||||
|
for mut rule in local.into_iter().flat_map(|config| config.rules) {
|
||||||
|
let id = rule.id.clone().unwrap_or_default();
|
||||||
|
if let Some(mut incoming_rule) = incoming_by_id.remove(&id) {
|
||||||
|
incoming_order.retain(|pending| pending != &id);
|
||||||
|
// The incoming expiry side wins; the local transition side is
|
||||||
|
// authoritative (MinIO CloneNonTransition + restore).
|
||||||
|
incoming_rule.transitions = rule.transitions.take();
|
||||||
|
incoming_rule.noncurrent_version_transitions = rule.noncurrent_version_transitions.take();
|
||||||
|
rules.push(incoming_rule);
|
||||||
|
} else if lifecycle_rule_has_expiry(&rule) {
|
||||||
|
// Expiry rule dropped upstream: remove a pure-expiry rule, or
|
||||||
|
// strip only the expiry side when a transition side remains.
|
||||||
|
if lifecycle_rule_has_transition(&rule) {
|
||||||
|
rule.expiration = None;
|
||||||
|
rule.noncurrent_version_expiration = None;
|
||||||
|
rule.del_marker_expiration = None;
|
||||||
|
rules.push(rule);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// No expiry semantics (transition-only / abort-mpu-only): not
|
||||||
|
// managed by expiry replication, keep untouched.
|
||||||
|
rules.push(rule);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
for id in incoming_order {
|
||||||
|
if let Some(rule) = incoming_by_id.remove(&id) {
|
||||||
|
rules.push(rule);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if rules.is_empty() {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
|
||||||
|
Some(s3s::dto::BucketLifecycleConfiguration {
|
||||||
|
rules,
|
||||||
|
// Stamp the expiry axis so the staleness guard compares expiry
|
||||||
|
// updates against expiry updates; a local transition-only edit moves
|
||||||
|
// only the whole-config timestamp, not this one.
|
||||||
|
expiry_updated_at: updated_at.map(s3s::dto::Timestamp::from).or(local_expiry_updated_at),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// True when the rule carries expiry semantics that `replicateILMExpiry` is
|
||||||
|
/// allowed to propagate.
|
||||||
|
fn lifecycle_rule_has_expiry(rule: &s3s::dto::LifecycleRule) -> bool {
|
||||||
|
rule.expiration.is_some() || rule.noncurrent_version_expiration.is_some() || rule.del_marker_expiration.is_some()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn lifecycle_rule_has_transition(rule: &s3s::dto::LifecycleRule) -> bool {
|
||||||
|
rule.transitions.as_ref().is_some_and(|transitions| !transitions.is_empty())
|
||||||
|
|| rule
|
||||||
|
.noncurrent_version_transitions
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|transitions| !transitions.is_empty())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Reduce a lifecycle XML document to the expiry subset that is allowed to
|
||||||
|
/// travel between sites (what MinIO's sender emits): transition fields are
|
||||||
|
/// stripped and rules left with no expiry semantics are dropped. Returns
|
||||||
|
/// `None` when nothing remains — the receiver then merges with the empty set,
|
||||||
|
/// which is exactly the "no expiry rules here" statement. A document that
|
||||||
|
/// fails to parse is forwarded unfiltered (`Some(original)`): the receiver
|
||||||
|
/// merge strips it anyway, and turning a local parse error into a `None`
|
||||||
|
/// would delete the peers' replicated expiry rules.
|
||||||
|
fn lifecycle_expiry_subset_xml(raw: &[u8]) -> Option<Vec<u8>> {
|
||||||
|
if raw.is_empty() {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let config: s3s::dto::BucketLifecycleConfiguration = match deserialize(raw) {
|
||||||
|
Ok(config) => config,
|
||||||
|
Err(err) => {
|
||||||
|
warn!("failed to parse local lifecycle config for expiry replication; forwarding unfiltered: {err}");
|
||||||
|
return Some(raw.to_vec());
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let expiry_updated_at = config.expiry_updated_at.clone();
|
||||||
|
let rules: Vec<s3s::dto::LifecycleRule> = config
|
||||||
|
.rules
|
||||||
|
.into_iter()
|
||||||
|
.filter_map(|mut rule| {
|
||||||
|
rule.transitions = None;
|
||||||
|
rule.noncurrent_version_transitions = None;
|
||||||
|
lifecycle_rule_has_expiry(&rule).then_some(rule)
|
||||||
|
})
|
||||||
|
.collect();
|
||||||
|
if rules.is_empty() {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let subset = s3s::dto::BucketLifecycleConfiguration {
|
||||||
|
rules,
|
||||||
|
expiry_updated_at,
|
||||||
|
};
|
||||||
|
match serialize(&subset) {
|
||||||
|
Ok(data) => Some(data),
|
||||||
|
Err(err) => {
|
||||||
|
warn!("failed to serialize lifecycle expiry subset; forwarding unfiltered: {err}");
|
||||||
|
Some(raw.to_vec())
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn replication_rule_deployment_id(rule: &ReplicationRule) -> Option<String> {
|
fn replication_rule_deployment_id(rule: &ReplicationRule) -> Option<String> {
|
||||||
@@ -7944,7 +8079,12 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
|
|||||||
} else {
|
} else {
|
||||||
None
|
None
|
||||||
};
|
};
|
||||||
if let Ok(bucket_meta) = metadata_sys::get(&item.bucket).await {
|
// lc-config staleness is judged on the expiry axis inside its merge block
|
||||||
|
// below: `lifecycle_config_updated_at` moves on local transition-only
|
||||||
|
// edits too, which would shadow newer peer expiry updates.
|
||||||
|
if item.r#type != "lc-config"
|
||||||
|
&& let Ok(bucket_meta) = metadata_sys::get(&item.bucket).await
|
||||||
|
{
|
||||||
let local_updated_at = bucket_meta_local_updated_at(&bucket_meta, config_file);
|
let local_updated_at = bucket_meta_local_updated_at(&bucket_meta, config_file);
|
||||||
if is_stale_update(local_updated_at, incoming_updated_at) {
|
if is_stale_update(local_updated_at, incoming_updated_at) {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
@@ -7986,6 +8126,54 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
|
|||||||
None
|
None
|
||||||
};
|
};
|
||||||
|
|
||||||
|
let merged_lifecycle_config = if item.r#type == "lc-config" {
|
||||||
|
// Receiver-side gate, symmetric with the sender hook: a peer must not
|
||||||
|
// install expiry rules here while `replicateILMExpiry` is off. When
|
||||||
|
// the state cannot be read, fall through and apply (pre-gate
|
||||||
|
// behavior) rather than silently dropping a legitimate update.
|
||||||
|
if let Ok(state) = load_site_replication_state().await
|
||||||
|
&& !site_replication_state_replicates_ilm_expiry(&state)
|
||||||
|
{
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
let incoming = item
|
||||||
|
.expiry_lc_config
|
||||||
|
.as_ref()
|
||||||
|
.map(|raw| {
|
||||||
|
let data = decode_bucket_meta_wire_value(raw);
|
||||||
|
deserialize::<s3s::dto::BucketLifecycleConfiguration>(&data)
|
||||||
|
})
|
||||||
|
.transpose()
|
||||||
|
.map_err(|e| s3_error!(InvalidRequest, "invalid lifecycle config: {e}"))?;
|
||||||
|
let local = match metadata_sys::get_lifecycle_config(&item.bucket).await {
|
||||||
|
Ok((config, _)) => Some(config),
|
||||||
|
Err(StorageError::ConfigNotFound) => None,
|
||||||
|
Err(err) => return Err(ApiError::from(err).into()),
|
||||||
|
};
|
||||||
|
let local_expiry_updated_at = local
|
||||||
|
.as_ref()
|
||||||
|
.and_then(|config| config.expiry_updated_at.clone())
|
||||||
|
.map(OffsetDateTime::from)
|
||||||
|
.unwrap_or(OffsetDateTime::UNIX_EPOCH);
|
||||||
|
if is_stale_update(local_expiry_updated_at, incoming_updated_at) {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
let local_absent = local.is_none();
|
||||||
|
match merge_incoming_lifecycle_config(incoming, local, incoming_updated_at) {
|
||||||
|
Some(config) => Some(
|
||||||
|
serialize(&config)
|
||||||
|
.map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize lifecycle failed: {e}")))?,
|
||||||
|
),
|
||||||
|
None => {
|
||||||
|
skip_config_write = local_absent;
|
||||||
|
None
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
};
|
||||||
|
|
||||||
let data = match item.r#type.as_str() {
|
let data = match item.r#type.as_str() {
|
||||||
"policy" => item
|
"policy" => item
|
||||||
.policy
|
.policy
|
||||||
@@ -8002,7 +8190,7 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
|
|||||||
"object-lock-config" => decode_bucket_meta_wire_option(item.object_lock_config),
|
"object-lock-config" => decode_bucket_meta_wire_option(item.object_lock_config),
|
||||||
"sse-config" => decode_bucket_meta_wire_option(item.sse_config),
|
"sse-config" => decode_bucket_meta_wire_option(item.sse_config),
|
||||||
"replication-config" => merged_replication_config,
|
"replication-config" => merged_replication_config,
|
||||||
"lc-config" => decode_bucket_meta_wire_option(item.expiry_lc_config),
|
"lc-config" => merged_lifecycle_config,
|
||||||
"cors-config" => decode_bucket_meta_wire_option(item.cors),
|
"cors-config" => decode_bucket_meta_wire_option(item.cors),
|
||||||
_ => unreachable!(),
|
_ => unreachable!(),
|
||||||
};
|
};
|
||||||
@@ -14537,6 +14725,36 @@ mod tests {
|
|||||||
assert_eq!(OffsetDateTime::from(stamped).unix_timestamp(), updated_at.unix_timestamp());
|
assert_eq!(OffsetDateTime::from(stamped).unix_timestamp(), updated_at.unix_timestamp());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Sender-side filter: only the expiry subset leaves this site. MinIO
|
||||||
|
/// peers install incoming rules verbatim, so a full document would plant
|
||||||
|
/// this site's transition rules there.
|
||||||
|
#[test]
|
||||||
|
fn test_lifecycle_expiry_subset_xml_strips_transitions() {
|
||||||
|
let full = serialize(&lc_config(vec![lc_rule("mixed", Some(1), Some(30)), lc_rule("t-only", None, Some(7))]))
|
||||||
|
.expect("serialize full config");
|
||||||
|
|
||||||
|
let subset = lifecycle_expiry_subset_xml(&full).expect("expiry subset should remain");
|
||||||
|
let parsed: s3s::dto::BucketLifecycleConfiguration = deserialize(&subset).expect("subset should parse");
|
||||||
|
assert_eq!(rule_ids(&parsed), vec!["mixed"]);
|
||||||
|
assert!(parsed.rules[0].transitions.is_none(), "transition side must not travel");
|
||||||
|
|
||||||
|
let transition_only =
|
||||||
|
serialize(&lc_config(vec![lc_rule("t-only", None, Some(7))])).expect("serialize transition-only config");
|
||||||
|
assert!(
|
||||||
|
lifecycle_expiry_subset_xml(&transition_only).is_none(),
|
||||||
|
"a transition-only config states 'no expiry rules' (delete semantics)"
|
||||||
|
);
|
||||||
|
assert!(lifecycle_expiry_subset_xml(b"").is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A local parse failure must forward the document unfiltered — mapping
|
||||||
|
/// it to `None` would delete the peers' replicated expiry rules.
|
||||||
|
#[test]
|
||||||
|
fn test_lifecycle_expiry_subset_xml_forwards_unparseable_config() {
|
||||||
|
let garbage = b"<LifecycleConfiguration><Rule></Broken>";
|
||||||
|
assert_eq!(lifecycle_expiry_subset_xml(garbage).as_deref(), Some(garbage.as_slice()));
|
||||||
|
}
|
||||||
|
|
||||||
// `role` is part of the bucket's S3-visible configuration. Repairing a reverse rule must
|
// `role` is part of the bucket's S3-visible configuration. Repairing a reverse rule must
|
||||||
// drop only a sender-owned site-replication ARN, never an operator's own role — the same
|
// drop only a sender-owned site-replication ARN, never an operator's own role — the same
|
||||||
// rule the merge path applies, so both paths agree on what is ours to rewrite.
|
// rule the merge path applies, so both paths agree on what is ours to rewrite.
|
||||||
|
|||||||
Reference in New Issue
Block a user