mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-31 09:18:28 +00:00
fix(site-replication): merge incoming ILM expiry documents instead of overwriting (#6130)
* test(site-replication): pin ILM expiry merge contract for incoming lc-config Red-light evidence for backlog#1675 P1-1: the lc-config receiver overwrites the whole local lifecycle config with whatever the peer sends (and deletes it wholesale on peer delete), so an expiry-only document erases the receiver's local tier/transition rules, and peer transition rules get installed across sites. The new tests pin the MinIO mergeWithCurrentLCConfig semantics plus RustFS hardening: - incoming expiry documents merge with (never replace) local rules - local transition sides are authoritative for same-id rules - incoming transition fields are discarded at the trust boundary - dropped expiry rules strip the expiry side but keep transitions; pure-expiry rules are removed - delete merges with the empty set instead of dropping the config - disabled rules survive; abort-mpu-only rules stay site-local - deterministic order (idempotent re-delivery) and expiry_updated_at stamping for the staleness axis All fail against the current overwrite implementation (identity extraction of merge_incoming_lifecycle_config). * 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. * fix(site-replication): close ILM merge review findings Adversarial review of the lc-config merge surfaced four real defects, all fixed here: - Deletion tombstone regression: with the staleness axis moved to the in-config expiry_updated_at, a deleted lifecycle config fell back to UNIX_EPOCH and any delayed stale broadcast could resurrect deleted expiry rules. The axis now falls back to the whole-config write time (which survives deletion in bucket metadata as the deletion's lower bound), also covering legacy configs that predate the axis field. - MinIO zero-rule documents: MinIO's delete tombstone / transition-only state marshals a lifecycle document with no <Rule>, which the strict s3s deserializer rejects — the receiver now recognizes it as the 'no expiry rules here' statement (delete semantics) instead of erroring on every MinIO heal pass. - Inflated expiry axis at the sender: PutBucketLifecycle stamped expiry_updated_at unconditionally, so a transition-only edit advanced the axis and let this site's stale expiry subset shadow and roll back newer peer expiry edits fleet-wide. The stamp is now conditional (expiry subset present before or after the edit, MinIO parity), the hook item travels with the config's expiry axis (UNIX_EPOCH when the site has none), and the SRInfo bucket entry feeds bootstrap/repair the same axis instead of the whole-config write time. - Del-marker parity: MinIO's CloneNonTransition never emits del-marker or abort-mpu fields, so treating del_marker_expiration as traveling expiry let a MinIO broadcast delete this site's del-marker-only rules. Both fields are now site-local on every edge: stripped from outbound subsets and inbound rules, restored from the local side on same-id merges, and never a deletion criterion. Receiver-side validation of merged configs (object-lock / tier constraints, MinIO runs finalLcCfg.Validate) remains a follow-up. * fix(site-replication): close the second ILM review round - Missed-delete repair: a deleted expiry state now travels through bootstrap/repair as an explicit timestamped lc-config delete item (lifecycle_expiry_statement distinguishes deletion — whole-config write time advanced past the created backfill — from never-configured buckets and from transition-only configs without an expiry axis, which say nothing). A peer that missed the live delete converges on repair; the receiver's staleness guard protects newer peer state. - Strict tombstone recognition: only a well-delimited zero-rule <LifecycleConfiguration> document maps to delete semantics; truncated or foreign payloads that fail the strict deserializer are rejected instead of being treated as a delete that erases local expiry rules. - Staleness fallback axis narrowed: the whole-config write time is used only for deleted or legacy-with-expiry state. A present transition-only config without an expiry axis compares at epoch — its whole-config time moves on transition edits and must not shadow or block independent peer expiry updates and same-timestamp repairs. * fix(site-replication): validate tombstone children structurally Second review round: a well-delimited root could still smuggle malformed content — e.g. <LifecycleConfiguration><ExpiryUpdatedAt> </LifecycleConfiguration> passed the no-<Rule check and was applied as a delete. The tombstone body must now be a sequence of well-formed simple children (matching open/close or self-closing, no nested markup, no stray text, none named Rule); anything else surfaces InvalidRequest. Malformed-child cases pinned in the recognition test. * fix(site-replication): serialize lifecycle merges --------- Co-authored-by: overtrue <anzhengchao@gmail.com>
This commit is contained in:
@@ -135,7 +135,8 @@ pub mod bucket {
|
|||||||
pub use crate::bucket::metadata_sys::ConfigWriteLockProbe;
|
pub use crate::bucket::metadata_sys::ConfigWriteLockProbe;
|
||||||
pub use crate::bucket::metadata_sys::{
|
pub use crate::bucket::metadata_sys::{
|
||||||
BucketMetadataMutationGuard, BucketMetadataSys, ObjectLockConfigState, acquire_bucket_metadata_transaction_lock,
|
BucketMetadataMutationGuard, BucketMetadataSys, ObjectLockConfigState, acquire_bucket_metadata_transaction_lock,
|
||||||
capture_bucket_metadata_incarnation, delete, delete_if_incarnation, get, get_accelerate_config, get_bucket_policy,
|
acquire_bucket_metadata_transaction_lock_for_incarnation, capture_bucket_metadata_incarnation, delete,
|
||||||
|
delete_if_incarnation, delete_under_transaction_lock, get, get_accelerate_config, get_bucket_policy,
|
||||||
get_bucket_policy_raw, get_bucket_targets_config, get_config_from_disk, get_cors_config, get_durability_config,
|
get_bucket_policy_raw, get_bucket_targets_config, get_config_from_disk, get_cors_config, get_durability_config,
|
||||||
get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config, get_notification_config,
|
get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config, get_notification_config,
|
||||||
get_object_lock_config, get_object_lock_config_state, get_public_access_block_config, get_quota_config,
|
get_object_lock_config, get_object_lock_config_state, get_public_access_block_config, get_quota_config,
|
||||||
|
|||||||
@@ -656,6 +656,16 @@ pub async fn update_under_transaction_lock(
|
|||||||
update_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file, data).await
|
update_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file, data).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Clear one config file while the caller holds this bucket's transaction lock.
|
||||||
|
pub async fn delete_under_transaction_lock(
|
||||||
|
guard: &BucketMetadataMutationGuard,
|
||||||
|
bucket: &str,
|
||||||
|
config_file: &str,
|
||||||
|
) -> Result<OffsetDateTime> {
|
||||||
|
guard.ensure_valid(bucket)?;
|
||||||
|
delete_under_config_write_guard(get_bucket_metadata_sys()?, guard, config_file).await
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn update_quota_if_incarnation(
|
pub async fn update_quota_if_incarnation(
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
data: Vec<u8>,
|
data: Vec<u8>,
|
||||||
@@ -795,6 +805,14 @@ pub async fn acquire_bucket_metadata_transaction_lock(bucket: &str) -> Result<Bu
|
|||||||
acquire_config_write_guard(get_bucket_metadata_sys()?, bucket).await
|
acquire_config_write_guard(get_bucket_metadata_sys()?, bucket).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Acquire the bucket transaction lock only if its incarnation still matches.
|
||||||
|
pub async fn acquire_bucket_metadata_transaction_lock_for_incarnation(
|
||||||
|
bucket: &str,
|
||||||
|
expected_incarnation_id: Uuid,
|
||||||
|
) -> Result<BucketMetadataMutationGuard> {
|
||||||
|
acquire_config_write_guard_for_incarnation(get_bucket_metadata_sys()?, bucket, Some(expected_incarnation_id)).await
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn acquire_bucket_metadata_transaction_lock_in(
|
pub(crate) async fn acquire_bucket_metadata_transaction_lock_in(
|
||||||
ctx: &crate::runtime::instance::InstanceContext,
|
ctx: &crate::runtime::instance::InstanceContext,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
|
|||||||
@@ -2323,6 +2323,7 @@ fn append_bootstrap_bucket_items(
|
|||||||
},
|
},
|
||||||
)?;
|
)?;
|
||||||
if replicate_ilm_expiry {
|
if replicate_ilm_expiry {
|
||||||
|
if bucket.expiry_lc_config.is_some() {
|
||||||
append_bootstrap_bucket_item(
|
append_bootstrap_bucket_item(
|
||||||
&mut plan.bucket_items,
|
&mut plan.bucket_items,
|
||||||
bucket,
|
bucket,
|
||||||
@@ -2331,10 +2332,22 @@ fn append_bootstrap_bucket_items(
|
|||||||
bucket.expiry_lc_config_updated_at,
|
bucket.expiry_lc_config_updated_at,
|
||||||
|item, value| {
|
|item, value| {
|
||||||
item.expiry_lc_config = Some(value);
|
item.expiry_lc_config = Some(value);
|
||||||
|
// `updated_at` here is the entry's expiry axis (see the
|
||||||
|
// SRBucketInfo construction), not the wall clock.
|
||||||
item.expiry_updated_at = item.updated_at;
|
item.expiry_updated_at = item.updated_at;
|
||||||
Ok(())
|
Ok(())
|
||||||
},
|
},
|
||||||
)?;
|
)?;
|
||||||
|
} else if bucket.expiry_lc_config_updated_at.is_some() {
|
||||||
|
// Expiry rules were removed at this axis (lifecycle_expiry_statement):
|
||||||
|
// an explicit timestamped delete item, so a peer that missed the
|
||||||
|
// live delete converges on bootstrap/repair instead of keeping
|
||||||
|
// stale expiry rules. The receiver's staleness guard protects a
|
||||||
|
// peer whose expiry state is newer.
|
||||||
|
let mut item = bootstrap_bucket_meta_item(bucket, "lc-config", bucket.expiry_lc_config_updated_at);
|
||||||
|
item.expiry_updated_at = item.updated_at;
|
||||||
|
plan.bucket_items.push(item);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
append_bootstrap_bucket_item(
|
append_bootstrap_bucket_item(
|
||||||
&mut plan.bucket_items,
|
&mut plan.bucket_items,
|
||||||
@@ -4220,13 +4233,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 +4342,14 @@ 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). A deleted expiry state is a `None`
|
||||||
|
// value with the deletion's axis so repair can converge peers
|
||||||
|
// that missed the live delete.
|
||||||
|
let expiry_statement = lifecycle_expiry_statement(&metadata);
|
||||||
|
entry.expiry_lc_config = expiry_statement.as_ref().and_then(|(subset, _)| subset.clone());
|
||||||
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);
|
||||||
@@ -4328,7 +4358,11 @@ async fn build_sr_info(state: &SiteReplicationState, local_peer: &PeerInfo) -> S
|
|||||||
entry.versioning_config_updated_at = maybe_time(metadata.versioning_config_updated_at);
|
entry.versioning_config_updated_at = maybe_time(metadata.versioning_config_updated_at);
|
||||||
entry.replication_config_updated_at = maybe_time(metadata.replication_config_updated_at);
|
entry.replication_config_updated_at = maybe_time(metadata.replication_config_updated_at);
|
||||||
entry.quota_config_updated_at = maybe_time(metadata.quota_config_updated_at);
|
entry.quota_config_updated_at = maybe_time(metadata.quota_config_updated_at);
|
||||||
entry.expiry_lc_config_updated_at = maybe_time(metadata.lifecycle_config_updated_at);
|
// The expiry axis, not the whole-config write time: local
|
||||||
|
// transition-only edits inflate the latter, and a repair item
|
||||||
|
// stamped with it could out-rank a newer real expiry edit on a
|
||||||
|
// third site.
|
||||||
|
entry.expiry_lc_config_updated_at = expiry_statement.map(|(_, axis)| axis);
|
||||||
entry.cors_config_updated_at = maybe_time(metadata.cors_config_updated_at);
|
entry.cors_config_updated_at = maybe_time(metadata.cors_config_updated_at);
|
||||||
entry.replication_targets_online =
|
entry.replication_targets_online =
|
||||||
Some(site_replication_targets_online(&bucket.name, &metadata.replication_config_xml).await);
|
Some(site_replication_targets_online(&bucket.name, &metadata.replication_config_xml).await);
|
||||||
@@ -6847,6 +6881,300 @@ fn merge_incoming_replication_config(
|
|||||||
Some(ReplicationConfiguration { role, rules })
|
Some(ReplicationConfiguration { role, rules })
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Merge a peer's ILM expiry document into the local lifecycle config.
|
||||||
|
///
|
||||||
|
/// Mirrors MinIO's `mergeWithCurrentLCConfig` with one hardening: incoming
|
||||||
|
/// site-local fields (transitions, abort-multipart, del-marker expiration —
|
||||||
|
/// exactly what MinIO's `CloneNonTransition` sender never emits) are
|
||||||
|
/// discarded outright at the trust boundary, whatever the peer sends. Local
|
||||||
|
/// site-local fields always survive; a delete (`incoming == None`) therefore
|
||||||
|
/// merges with the empty set instead of dropping the whole config.
|
||||||
|
fn merge_incoming_lifecycle_config(
|
||||||
|
incoming: Option<s3s::dto::BucketLifecycleConfiguration>,
|
||||||
|
local: Option<s3s::dto::BucketLifecycleConfiguration>,
|
||||||
|
updated_at: Option<OffsetDateTime>,
|
||||||
|
) -> Option<s3s::dto::BucketLifecycleConfiguration> {
|
||||||
|
// Incoming rules reduced to their traveling expiry side. Rules with no
|
||||||
|
// expiry semantics after the strip 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) {
|
||||||
|
strip_site_local_lifecycle_fields(&mut rule);
|
||||||
|
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 site-local side is
|
||||||
|
// authoritative (MinIO CloneNonTransition + restore).
|
||||||
|
incoming_rule.transitions = rule.transitions.take();
|
||||||
|
incoming_rule.noncurrent_version_transitions = rule.noncurrent_version_transitions.take();
|
||||||
|
incoming_rule.abort_incomplete_multipart_upload = rule.abort_incomplete_multipart_upload.take();
|
||||||
|
incoming_rule.del_marker_expiration = rule.del_marker_expiration.take();
|
||||||
|
rules.push(incoming_rule);
|
||||||
|
} else if lifecycle_rule_has_expiry(&rule) {
|
||||||
|
// Expiry rule dropped upstream: strip only the traveling expiry
|
||||||
|
// side; the rule survives while any site-local action remains.
|
||||||
|
rule.expiration = None;
|
||||||
|
rule.noncurrent_version_expiration = None;
|
||||||
|
if lifecycle_rule_has_transition(&rule)
|
||||||
|
|| rule.abort_incomplete_multipart_upload.is_some()
|
||||||
|
|| rule.del_marker_expiration.is_some()
|
||||||
|
{
|
||||||
|
rules.push(rule);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// No traveling expiry semantics (transition-only / abort-mpu-only
|
||||||
|
// / del-marker-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,
|
||||||
|
// Record the expiry axis the staleness guard compares on. (The PUT
|
||||||
|
// path stamps `expiry_updated_at` only when the expiry subset
|
||||||
|
// changes, so this axis is not inflated by transition-only edits.)
|
||||||
|
expiry_updated_at: updated_at.map(s3s::dto::Timestamp::from).or(local_expiry_updated_at),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// True when the rule carries the expiry semantics that `replicateILMExpiry`
|
||||||
|
/// propagates. Del-marker expiration and abort-multipart are deliberately
|
||||||
|
/// excluded: MinIO's sender never emits them (`CloneNonTransition` drops
|
||||||
|
/// both), so treating them as traveling state would let a MinIO peer's
|
||||||
|
/// broadcast delete this site's del-marker-only rules.
|
||||||
|
fn lifecycle_rule_has_expiry(rule: &s3s::dto::LifecycleRule) -> bool {
|
||||||
|
rule.expiration.is_some() || rule.noncurrent_version_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())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Remove the fields that never travel between sites (MinIO
|
||||||
|
/// `CloneNonTransition` parity).
|
||||||
|
fn strip_site_local_lifecycle_fields(rule: &mut s3s::dto::LifecycleRule) {
|
||||||
|
rule.transitions = None;
|
||||||
|
rule.noncurrent_version_transitions = None;
|
||||||
|
rule.abort_incomplete_multipart_upload = None;
|
||||||
|
rule.del_marker_expiration = None;
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 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| {
|
||||||
|
strip_site_local_lifecycle_fields(&mut rule);
|
||||||
|
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())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The expiry replication axis persisted in a lifecycle XML document, if any.
|
||||||
|
/// Used for the SRInfo bucket entry so bootstrap/repair items carry the
|
||||||
|
/// expiry axis instead of the whole-config write time (which local
|
||||||
|
/// transition-only edits inflate).
|
||||||
|
fn lifecycle_expiry_updated_at(raw: &[u8]) -> Option<OffsetDateTime> {
|
||||||
|
if raw.is_empty() {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
deserialize::<s3s::dto::BucketLifecycleConfiguration>(raw)
|
||||||
|
.ok()
|
||||||
|
.and_then(|config| config.expiry_updated_at)
|
||||||
|
.map(OffsetDateTime::from)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The timestamp an incoming lc-config item must beat to be applied.
|
||||||
|
///
|
||||||
|
/// - Present config with the expiry axis: the axis itself.
|
||||||
|
/// - Present legacy config that has expiry rules but predates the axis
|
||||||
|
/// field: the whole-config write time bounds its last expiry edit.
|
||||||
|
/// - Present transition-only config without the axis: `UNIX_EPOCH` — there
|
||||||
|
/// is no local expiry state to protect, and the whole-config time moves on
|
||||||
|
/// transition edits, which must not shadow independent peer expiry updates.
|
||||||
|
/// - Absent config: the whole-config write time — it survives deletion in
|
||||||
|
/// bucket metadata as the deletion's lower bound, so a delayed stale
|
||||||
|
/// broadcast cannot resurrect deleted expiry rules.
|
||||||
|
fn local_lifecycle_staleness_axis(
|
||||||
|
local: Option<&s3s::dto::BucketLifecycleConfiguration>,
|
||||||
|
whole_config_axis: OffsetDateTime,
|
||||||
|
) -> OffsetDateTime {
|
||||||
|
match local {
|
||||||
|
Some(config) => match config.expiry_updated_at.clone() {
|
||||||
|
Some(axis) => OffsetDateTime::from(axis),
|
||||||
|
None if config.rules.iter().any(lifecycle_rule_has_expiry) => whole_config_axis,
|
||||||
|
None => OffsetDateTime::UNIX_EPOCH,
|
||||||
|
},
|
||||||
|
None => whole_config_axis,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Recognize MinIO's zero-rule lifecycle tombstone (its delete /
|
||||||
|
/// transition-only state marshals `<LifecycleConfiguration>` with no `<Rule>`
|
||||||
|
/// child, which the strict s3s deserializer rejects). Only a well-delimited
|
||||||
|
/// document qualifies as the "no expiry rules here" statement; truncated or
|
||||||
|
/// otherwise malformed payloads are rejected rather than treated as a delete
|
||||||
|
/// that would erase local expiry rules.
|
||||||
|
fn is_zero_rule_lifecycle_tombstone(raw: &[u8]) -> bool {
|
||||||
|
#[derive(Deserialize)]
|
||||||
|
#[serde(deny_unknown_fields)]
|
||||||
|
struct Tombstone {
|
||||||
|
#[serde(rename = "@xmlns")]
|
||||||
|
_xmlns: Option<String>,
|
||||||
|
#[serde(rename = "ExpiryUpdatedAt")]
|
||||||
|
_expiry_updated_at: Option<s3s::dto::Timestamp>,
|
||||||
|
}
|
||||||
|
|
||||||
|
let mut reader = quick_xml::Reader::from_reader(raw);
|
||||||
|
let mut depth = 0usize;
|
||||||
|
let mut seen_root = false;
|
||||||
|
let mut closed_root = false;
|
||||||
|
let mut seen_declaration = false;
|
||||||
|
let well_formed_document = loop {
|
||||||
|
match reader.read_event() {
|
||||||
|
Ok(quick_xml::events::Event::Start(element)) => {
|
||||||
|
if depth == 0 {
|
||||||
|
if seen_root || closed_root || element.name().as_ref() != b"LifecycleConfiguration" {
|
||||||
|
break false;
|
||||||
|
}
|
||||||
|
seen_root = true;
|
||||||
|
}
|
||||||
|
depth += 1;
|
||||||
|
}
|
||||||
|
Ok(quick_xml::events::Event::Empty(element)) => {
|
||||||
|
if depth == 0 {
|
||||||
|
if seen_root || closed_root || element.name().as_ref() != b"LifecycleConfiguration" {
|
||||||
|
break false;
|
||||||
|
}
|
||||||
|
seen_root = true;
|
||||||
|
closed_root = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(quick_xml::events::Event::End(_)) => {
|
||||||
|
if depth == 0 {
|
||||||
|
break false;
|
||||||
|
}
|
||||||
|
depth -= 1;
|
||||||
|
if depth == 0 {
|
||||||
|
closed_root = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(quick_xml::events::Event::Decl(_)) => {
|
||||||
|
if seen_declaration || seen_root || depth != 0 {
|
||||||
|
break false;
|
||||||
|
}
|
||||||
|
seen_declaration = true;
|
||||||
|
}
|
||||||
|
Ok(quick_xml::events::Event::DocType(_)) => break false,
|
||||||
|
Ok(quick_xml::events::Event::Text(text)) if depth == 0 && !text.iter().all(u8::is_ascii_whitespace) => {
|
||||||
|
break false;
|
||||||
|
}
|
||||||
|
Ok(quick_xml::events::Event::Text(_)) => {}
|
||||||
|
Ok(quick_xml::events::Event::CData(_)) if depth == 0 => break false,
|
||||||
|
Ok(quick_xml::events::Event::Comment(_) | quick_xml::events::Event::PI(_)) => {}
|
||||||
|
Ok(quick_xml::events::Event::Eof) => break seen_root && closed_root && depth == 0,
|
||||||
|
Ok(_) if depth == 0 => break false,
|
||||||
|
Ok(_) => {}
|
||||||
|
Err(_) => break false,
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
well_formed_document && quick_xml::de::from_reader::<_, Tombstone>(raw).is_ok()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The ILM expiry statement this site contributes to its SRInfo bucket entry
|
||||||
|
/// (feeding bootstrap/repair and consistency views), if any.
|
||||||
|
/// `Some((subset_b64, axis))` — a `None` subset means "expiry rules were
|
||||||
|
/// removed at `axis`" and travels as an explicit timestamped delete item, so
|
||||||
|
/// a peer that missed the live delete still converges on repair.
|
||||||
|
fn lifecycle_expiry_statement(
|
||||||
|
metadata: &crate::admin::storage_api::bucket::metadata::BucketMetadata,
|
||||||
|
) -> Option<(Option<String>, OffsetDateTime)> {
|
||||||
|
if metadata.lifecycle_config_xml.is_empty() {
|
||||||
|
// Deleted vs never configured: the whole-config write time survives
|
||||||
|
// deletion in bucket metadata and strictly exceeds the created-time
|
||||||
|
// backfill only after a real write.
|
||||||
|
return (metadata.lifecycle_config_updated_at > metadata.created).then_some((None, metadata.lifecycle_config_updated_at));
|
||||||
|
}
|
||||||
|
let axis = lifecycle_expiry_updated_at(&metadata.lifecycle_config_xml);
|
||||||
|
match lifecycle_expiry_subset_xml(&metadata.lifecycle_config_xml) {
|
||||||
|
Some(subset) => {
|
||||||
|
// Legacy documents predate the axis field; their whole-config
|
||||||
|
// write time bounds the last expiry edit.
|
||||||
|
let axis = axis.unwrap_or(metadata.lifecycle_config_updated_at);
|
||||||
|
Some((raw_config_to_base64(&subset), axis))
|
||||||
|
}
|
||||||
|
// Transition-only config: with an expiry axis the site once had
|
||||||
|
// expiry rules and properly removed them — the delete travels at
|
||||||
|
// that axis. Without one there is nothing to say (a delete stamped
|
||||||
|
// off the whole-config time would let a local transition edit erase
|
||||||
|
// newer peer expiry state).
|
||||||
|
None => axis.map(|axis| (None, axis)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn replication_rule_deployment_id(rule: &ReplicationRule) -> Option<String> {
|
fn replication_rule_deployment_id(rule: &ReplicationRule) -> Option<String> {
|
||||||
if let Some(rule_id) = rule.id.as_deref() {
|
if let Some(rule_id) = rule.id.as_deref() {
|
||||||
if let Some(deployment_id) = rule_id.strip_prefix("site-repl-")
|
if let Some(deployment_id) = rule_id.strip_prefix("site-repl-")
|
||||||
@@ -7928,7 +8256,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(());
|
||||||
@@ -7970,6 +8303,72 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
|
|||||||
None
|
None
|
||||||
};
|
};
|
||||||
|
|
||||||
|
let (merged_lifecycle_config, lifecycle_guard) = 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. Note
|
||||||
|
// the gate acks with 200 — the sender treats the item as delivered
|
||||||
|
// and will not retry; items skipped inside the enable-flag
|
||||||
|
// propagation window are healed by repair, not by retry.
|
||||||
|
if let Ok(state) = load_site_replication_state().await
|
||||||
|
&& !site_replication_state_replicates_ilm_expiry(&state)
|
||||||
|
{
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
|
||||||
|
let incoming = match item.expiry_lc_config.as_ref() {
|
||||||
|
Some(raw) => {
|
||||||
|
let data = decode_bucket_meta_wire_value(raw);
|
||||||
|
match deserialize::<s3s::dto::BucketLifecycleConfiguration>(&data) {
|
||||||
|
Ok(config) => Some(config),
|
||||||
|
// MinIO's delete tombstone / transition-only state is a
|
||||||
|
// zero-rule document the strict deserializer rejects; it
|
||||||
|
// means "no expiry rules here" (delete semantics). Any
|
||||||
|
// other malformed payload is rejected — treating it as a
|
||||||
|
// delete would let a bad payload erase local expiry rules.
|
||||||
|
Err(_) if is_zero_rule_lifecycle_tombstone(&data) => None,
|
||||||
|
Err(e) => return Err(s3_error!(InvalidRequest, "invalid lifecycle config: {e}")),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
None => None,
|
||||||
|
};
|
||||||
|
let lifecycle_guard =
|
||||||
|
metadata_sys::acquire_bucket_metadata_transaction_lock_for_incarnation(&item.bucket, expected_incarnation_id)
|
||||||
|
.await
|
||||||
|
.map_err(ApiError::from)?;
|
||||||
|
let local_metadata = metadata_sys::get_config_from_disk(&item.bucket)
|
||||||
|
.await
|
||||||
|
.map_err(ApiError::from)?;
|
||||||
|
let local = if local_metadata.lifecycle_config_xml.is_empty() {
|
||||||
|
None
|
||||||
|
} else {
|
||||||
|
Some(
|
||||||
|
deserialize::<s3s::dto::BucketLifecycleConfiguration>(&local_metadata.lifecycle_config_xml).map_err(|e| {
|
||||||
|
S3Error::with_message(S3ErrorCode::InternalError, format!("invalid local lifecycle config: {e}"))
|
||||||
|
})?,
|
||||||
|
)
|
||||||
|
};
|
||||||
|
let whole_config_axis = local_metadata.lifecycle_config_updated_at;
|
||||||
|
if is_stale_update(local_lifecycle_staleness_axis(local.as_ref(), whole_config_axis), incoming_updated_at) {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
let local_absent = local.is_none();
|
||||||
|
let merged = 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
|
||||||
|
}
|
||||||
|
};
|
||||||
|
(merged, Some(lifecycle_guard))
|
||||||
|
} else {
|
||||||
|
(None, None)
|
||||||
|
};
|
||||||
|
|
||||||
let data = match item.r#type.as_str() {
|
let data = match item.r#type.as_str() {
|
||||||
"policy" => item
|
"policy" => item
|
||||||
.policy
|
.policy
|
||||||
@@ -7986,7 +8385,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!(),
|
||||||
};
|
};
|
||||||
@@ -8017,17 +8416,30 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
|
|||||||
.await
|
.await
|
||||||
.map_err(ApiError::from)?;
|
.map_err(ApiError::from)?;
|
||||||
}
|
}
|
||||||
|
} else {
|
||||||
|
if let Some(guard) = lifecycle_guard.as_ref() {
|
||||||
|
metadata_sys::update_under_transaction_lock(guard, &item.bucket, config_file, data)
|
||||||
|
.await
|
||||||
|
.map_err(ApiError::from)?;
|
||||||
} else {
|
} else {
|
||||||
metadata_sys::update_if_incarnation(&item.bucket, config_file, data, expected_incarnation_id)
|
metadata_sys::update_if_incarnation(&item.bucket, config_file, data, expected_incarnation_id)
|
||||||
.await
|
.await
|
||||||
.map_err(ApiError::from)?;
|
.map_err(ApiError::from)?;
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
if let Some(guard) = lifecycle_guard.as_ref() {
|
||||||
|
metadata_sys::delete_under_transaction_lock(guard, &item.bucket, config_file)
|
||||||
|
.await
|
||||||
|
.map_err(ApiError::from)?;
|
||||||
} else {
|
} else {
|
||||||
metadata_sys::delete_if_incarnation(&item.bucket, config_file, expected_incarnation_id)
|
metadata_sys::delete_if_incarnation(&item.bucket, config_file, expected_incarnation_id)
|
||||||
.await
|
.await
|
||||||
.map_err(ApiError::from)?;
|
.map_err(ApiError::from)?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
drop(lifecycle_guard);
|
||||||
drop(targets_guard);
|
drop(targets_guard);
|
||||||
|
|
||||||
if item.r#type == "replication-config" {
|
if item.r#type == "replication-config" {
|
||||||
@@ -12645,6 +13057,87 @@ mod tests {
|
|||||||
assert!(!plan.bucket_items.iter().any(|item| item.r#type == "lc-config"));
|
assert!(!plan.bucket_items.iter().any(|item| item.r#type == "lc-config"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A deleted expiry state (entry value None, axis set) must travel as an
|
||||||
|
/// explicit timestamped delete item — a peer that missed the live delete
|
||||||
|
/// otherwise keeps stale expiry rules through every repair (review
|
||||||
|
/// finding).
|
||||||
|
#[test]
|
||||||
|
fn test_site_replication_bootstrap_plan_emits_timestamped_lifecycle_delete() {
|
||||||
|
let deleted_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
|
||||||
|
let mut info = SRInfo::default();
|
||||||
|
info.state.peers.insert(
|
||||||
|
"remote-dep".to_string(),
|
||||||
|
PeerInfo {
|
||||||
|
replicate_ilm_expiry: true,
|
||||||
|
..peer("remote", "https://remote.example.com")
|
||||||
|
},
|
||||||
|
);
|
||||||
|
info.buckets.insert(
|
||||||
|
"photos".to_string(),
|
||||||
|
SRBucketInfo {
|
||||||
|
bucket: "photos".to_string(),
|
||||||
|
expiry_lc_config: None,
|
||||||
|
expiry_lc_config_updated_at: Some(deleted_at),
|
||||||
|
api_version: Some(SITE_REPL_API_VERSION.to_string()),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
let plan = site_replication_bootstrap_plan(&info).expect("bootstrap plan should build");
|
||||||
|
|
||||||
|
let item = plan
|
||||||
|
.bucket_items
|
||||||
|
.iter()
|
||||||
|
.find(|item| item.r#type == "lc-config")
|
||||||
|
.expect("a deleted expiry state must produce an lc-config delete item");
|
||||||
|
assert!(item.expiry_lc_config.is_none(), "delete items carry no config body");
|
||||||
|
assert_eq!(item.expiry_updated_at, Some(deleted_at));
|
||||||
|
assert_eq!(item.updated_at, Some(deleted_at));
|
||||||
|
}
|
||||||
|
|
||||||
|
/// What each local lifecycle state contributes to the SRInfo entry:
|
||||||
|
/// deletions are timestamped statements, never-configured buckets and
|
||||||
|
/// transition-only configs without an expiry axis say nothing.
|
||||||
|
#[test]
|
||||||
|
fn test_lifecycle_expiry_statement_matrix() {
|
||||||
|
let created = OffsetDateTime::from_unix_timestamp(1_600_000_000).expect("timestamp");
|
||||||
|
let mut meta = crate::admin::storage_api::bucket::metadata::BucketMetadata::new("photos");
|
||||||
|
meta.created = created;
|
||||||
|
// Never configured: load backfills the write time to `created`.
|
||||||
|
meta.lifecycle_config_updated_at = created;
|
||||||
|
assert!(lifecycle_expiry_statement(&meta).is_none());
|
||||||
|
|
||||||
|
// Deleted: the write time survives deletion and exceeds creation.
|
||||||
|
let deleted_at = created + time::Duration::seconds(100);
|
||||||
|
meta.lifecycle_config_updated_at = deleted_at;
|
||||||
|
let (subset, axis) = lifecycle_expiry_statement(&meta).expect("deletion is a statement");
|
||||||
|
assert!(subset.is_none());
|
||||||
|
assert_eq!(axis, deleted_at);
|
||||||
|
|
||||||
|
// Present with expiry rules and the axis: subset + axis travel.
|
||||||
|
let expiry_axis = created + time::Duration::seconds(50);
|
||||||
|
let mut config = lc_config(vec![lc_rule("e1", Some(7), None)]);
|
||||||
|
config.expiry_updated_at = Some(s3s::dto::Timestamp::from(expiry_axis));
|
||||||
|
meta.lifecycle_config_xml = serialize(&config).expect("serialize config");
|
||||||
|
let (subset, axis) = lifecycle_expiry_statement(&meta).expect("expiry config is a statement");
|
||||||
|
assert!(subset.is_some());
|
||||||
|
assert_eq!(axis.unix_timestamp(), expiry_axis.unix_timestamp());
|
||||||
|
|
||||||
|
// Transition-only without an axis: nothing to say (a delete stamped
|
||||||
|
// off the whole-config time would erase newer peer expiry state).
|
||||||
|
meta.lifecycle_config_xml = serialize(&lc_config(vec![lc_rule("t1", None, Some(30))])).expect("serialize config");
|
||||||
|
assert!(lifecycle_expiry_statement(&meta).is_none());
|
||||||
|
|
||||||
|
// Transition-only WITH an axis: expiry rules were properly removed —
|
||||||
|
// the delete travels at that axis.
|
||||||
|
let mut transition_only = lc_config(vec![lc_rule("t1", None, Some(30))]);
|
||||||
|
transition_only.expiry_updated_at = Some(s3s::dto::Timestamp::from(expiry_axis));
|
||||||
|
meta.lifecycle_config_xml = serialize(&transition_only).expect("serialize config");
|
||||||
|
let (subset, axis) = lifecycle_expiry_statement(&meta).expect("removed expiry state is a statement");
|
||||||
|
assert!(subset.is_none());
|
||||||
|
assert_eq!(axis.unix_timestamp(), expiry_axis.unix_timestamp());
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_site_replication_repair_request_is_strict_and_requires_explicit_mode() {
|
fn test_site_replication_repair_request_is_strict_and_requires_explicit_mode() {
|
||||||
assert!(serde_json::from_str::<SiteReplicationRepairRequest>(r#"{"mode":"dry-run"}"#).is_ok());
|
assert!(serde_json::from_str::<SiteReplicationRepairRequest>(r#"{"mode":"dry-run"}"#).is_ok());
|
||||||
@@ -14299,6 +14792,405 @@ mod tests {
|
|||||||
assert!(merge_incoming_replication_config(Some(site_repl_config("home")), None).is_none());
|
assert!(merge_incoming_replication_config(Some(site_repl_config("home")), None).is_none());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn lc_rule(id: &str, expiry_days: Option<i32>, transition_days: Option<i32>) -> s3s::dto::LifecycleRule {
|
||||||
|
s3s::dto::LifecycleRule {
|
||||||
|
id: Some(id.to_string()),
|
||||||
|
status: s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::ENABLED),
|
||||||
|
prefix: Some(String::new()),
|
||||||
|
expiration: expiry_days.map(|days| s3s::dto::LifecycleExpiration {
|
||||||
|
days: Some(days),
|
||||||
|
..Default::default()
|
||||||
|
}),
|
||||||
|
transitions: transition_days.map(|days| {
|
||||||
|
vec![s3s::dto::Transition {
|
||||||
|
days: Some(days),
|
||||||
|
storage_class: Some(s3s::dto::TransitionStorageClass::from_static(s3s::dto::TransitionStorageClass::GLACIER)),
|
||||||
|
date: None,
|
||||||
|
}]
|
||||||
|
}),
|
||||||
|
abort_incomplete_multipart_upload: None,
|
||||||
|
del_marker_expiration: None,
|
||||||
|
filter: None,
|
||||||
|
noncurrent_version_expiration: None,
|
||||||
|
noncurrent_version_transitions: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn lc_config(rules: Vec<s3s::dto::LifecycleRule>) -> s3s::dto::BucketLifecycleConfiguration {
|
||||||
|
s3s::dto::BucketLifecycleConfiguration {
|
||||||
|
rules,
|
||||||
|
expiry_updated_at: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn rule_ids(config: &s3s::dto::BucketLifecycleConfiguration) -> Vec<&str> {
|
||||||
|
config.rules.iter().filter_map(|rule| rule.id.as_deref()).collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// P1-1 red-light: an incoming expiry-only document must not erase the
|
||||||
|
/// receiver's local transition/tiering rules (today the receiver
|
||||||
|
/// overwrites the whole lifecycle config).
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_preserves_local_transition_rule() {
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
Some(lc_config(vec![lc_rule("e1", Some(7), None)])),
|
||||||
|
Some(lc_config(vec![lc_rule("t1", None, Some(30))])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
|
||||||
|
let mut ids = rule_ids(&merged);
|
||||||
|
ids.sort_unstable();
|
||||||
|
assert_eq!(ids, vec!["e1", "t1"]);
|
||||||
|
let t1 = merged.rules.iter().find(|rule| rule.id.as_deref() == Some("t1")).unwrap();
|
||||||
|
assert!(t1.transitions.as_ref().is_some_and(|t| !t.is_empty()), "local transition must survive");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Same-id incoming rule updates the expiry side but the local transition
|
||||||
|
/// side is authoritative (MinIO `CloneNonTransition` + restore).
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_same_id_keeps_local_transition() {
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
Some(lc_config(vec![lc_rule("r1", Some(7), None)])),
|
||||||
|
Some(lc_config(vec![lc_rule("r1", Some(1), Some(30))])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
|
||||||
|
assert_eq!(merged.rules.len(), 1);
|
||||||
|
let r1 = &merged.rules[0];
|
||||||
|
assert_eq!(r1.expiration.as_ref().and_then(|e| e.days), Some(7), "incoming expiry wins");
|
||||||
|
assert!(
|
||||||
|
r1.transitions.as_ref().is_some_and(|t| !t.is_empty()),
|
||||||
|
"local transition is authoritative"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Trust boundary: whatever the peer sends, its transition fields never
|
||||||
|
/// land here — a new incoming rule is stripped to its expiry parts.
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_strips_incoming_transitions() {
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
Some(lc_config(vec![lc_rule("r1", Some(7), Some(1))])),
|
||||||
|
Some(lc_config(vec![lc_rule("t1", None, Some(30))])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
|
||||||
|
let r1 = merged.rules.iter().find(|rule| rule.id.as_deref() == Some("r1")).unwrap();
|
||||||
|
assert!(
|
||||||
|
r1.transitions.as_ref().is_none_or(|t| t.is_empty()),
|
||||||
|
"incoming transition fields must be discarded"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A local rule whose expiry part was dropped upstream loses only the
|
||||||
|
/// expiry fields; a pure-expiry rule disappears entirely.
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_dropped_rule_strips_expiry_keeps_transition() {
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
Some(lc_config(vec![lc_rule("other", Some(3), None)])),
|
||||||
|
Some(lc_config(vec![
|
||||||
|
lc_rule("mixed", Some(1), Some(30)),
|
||||||
|
lc_rule("pure-expiry", Some(2), None),
|
||||||
|
])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
|
||||||
|
let mut ids = rule_ids(&merged);
|
||||||
|
ids.sort_unstable();
|
||||||
|
assert_eq!(ids, vec!["mixed", "other"], "pure-expiry rule not in the incoming set is removed");
|
||||||
|
let mixed = merged.rules.iter().find(|rule| rule.id.as_deref() == Some("mixed")).unwrap();
|
||||||
|
assert!(mixed.expiration.is_none(), "expiry side cleared");
|
||||||
|
assert!(mixed.transitions.as_ref().is_some_and(|t| !t.is_empty()), "transition side kept");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Peer lifecycle delete merges with the empty set: local transition rules
|
||||||
|
/// survive with their expiry parts cleared; only when nothing remains does
|
||||||
|
/// the whole config disappear.
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_delete_merges_with_empty() {
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
None,
|
||||||
|
Some(lc_config(vec![
|
||||||
|
lc_rule("mixed", Some(1), Some(30)),
|
||||||
|
lc_rule("pure-expiry", Some(2), None),
|
||||||
|
])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("transition rules must survive a peer lifecycle delete");
|
||||||
|
assert_eq!(rule_ids(&merged), vec!["mixed"]);
|
||||||
|
assert!(merged.rules[0].expiration.is_none());
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
merge_incoming_lifecycle_config(None, Some(lc_config(vec![lc_rule("pure-expiry", Some(2), None)])), None).is_none(),
|
||||||
|
"an all-expiry config deletes cleanly"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Disabled rules must survive the merge like enabled ones — the merge
|
||||||
|
/// must not reuse ENABLED-filtered helpers.
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_keeps_disabled_transition_rule() {
|
||||||
|
let mut disabled = lc_rule("t-disabled", None, Some(30));
|
||||||
|
disabled.status = s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::DISABLED);
|
||||||
|
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
Some(lc_config(vec![lc_rule("e1", Some(7), None)])),
|
||||||
|
Some(lc_config(vec![disabled])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
|
||||||
|
let mut ids = rule_ids(&merged);
|
||||||
|
ids.sort_unstable();
|
||||||
|
assert_eq!(ids, vec!["e1", "t-disabled"]);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Abort-multipart-only rules carry no expiry semantics: local ones stay
|
||||||
|
/// untouched, incoming ones are not installed (they are site-local, like
|
||||||
|
/// MinIO's sender-side filter).
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_abort_mpu_rules_stay_local() {
|
||||||
|
let abort_only = |id: &str| s3s::dto::LifecycleRule {
|
||||||
|
id: Some(id.to_string()),
|
||||||
|
status: s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::ENABLED),
|
||||||
|
prefix: Some(String::new()),
|
||||||
|
abort_incomplete_multipart_upload: Some(s3s::dto::AbortIncompleteMultipartUpload {
|
||||||
|
days_after_initiation: Some(3),
|
||||||
|
}),
|
||||||
|
del_marker_expiration: None,
|
||||||
|
expiration: None,
|
||||||
|
filter: None,
|
||||||
|
noncurrent_version_expiration: None,
|
||||||
|
noncurrent_version_transitions: None,
|
||||||
|
transitions: None,
|
||||||
|
};
|
||||||
|
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
Some(lc_config(vec![abort_only("incoming-abort"), lc_rule("e1", Some(7), None)])),
|
||||||
|
Some(lc_config(vec![abort_only("local-abort")])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
|
||||||
|
let mut ids = rule_ids(&merged);
|
||||||
|
ids.sort_unstable();
|
||||||
|
assert_eq!(
|
||||||
|
ids,
|
||||||
|
vec!["e1", "local-abort"],
|
||||||
|
"incoming abort-mpu rule is not installed; local one survives"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Repeated delivery of the same document must be byte-stable (rule order
|
||||||
|
/// deterministic), or every broadcast rewrites bucket metadata.
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_is_idempotent() {
|
||||||
|
let incoming = || Some(lc_config(vec![lc_rule("e1", Some(7), None), lc_rule("e2", Some(9), None)]));
|
||||||
|
let local = Some(lc_config(vec![lc_rule("t1", None, Some(30))]));
|
||||||
|
|
||||||
|
let once = merge_incoming_lifecycle_config(incoming(), local, None).expect("first merge");
|
||||||
|
let twice = merge_incoming_lifecycle_config(incoming(), Some(once.clone()), None).expect("second merge");
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
serialize(&once).expect("serialize once"),
|
||||||
|
serialize(&twice).expect("serialize twice"),
|
||||||
|
"merge must be idempotent for identical input"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The merged config records the expiry axis timestamp so the staleness
|
||||||
|
/// guard compares expiry updates against expiry updates (a local
|
||||||
|
/// transition-only edit must not shadow newer peer expiry updates).
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_stamps_expiry_updated_at() {
|
||||||
|
let updated_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
|
||||||
|
let merged = merge_incoming_lifecycle_config(Some(lc_config(vec![lc_rule("e1", Some(7), None)])), None, Some(updated_at))
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
|
||||||
|
let stamped = merged.expiry_updated_at.expect("expiry_updated_at must be stamped");
|
||||||
|
assert_eq!(OffsetDateTime::from(stamped).unix_timestamp(), updated_at.unix_timestamp());
|
||||||
|
}
|
||||||
|
|
||||||
|
/// MinIO's sender never emits del-marker-expiration rules
|
||||||
|
/// (CloneNonTransition drops them), so a MinIO expiry broadcast must not
|
||||||
|
/// delete this site's del-marker-only rules, and the local del-marker
|
||||||
|
/// side of a same-id rule is authoritative.
|
||||||
|
#[test]
|
||||||
|
fn test_merge_incoming_lifecycle_del_marker_rules_stay_local() {
|
||||||
|
let del_marker_only = |id: &str| {
|
||||||
|
let mut rule = lc_rule(id, None, None);
|
||||||
|
rule.del_marker_expiration = Some(s3s::dto::DelMarkerExpiration { days: Some(3) });
|
||||||
|
rule
|
||||||
|
};
|
||||||
|
|
||||||
|
// A local del-marker-only rule survives an incoming expiry document
|
||||||
|
// that does not mention it.
|
||||||
|
let merged = merge_incoming_lifecycle_config(
|
||||||
|
Some(lc_config(vec![lc_rule("e1", Some(7), None)])),
|
||||||
|
Some(lc_config(vec![del_marker_only("dm-local")])),
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
let mut ids = rule_ids(&merged);
|
||||||
|
ids.sort_unstable();
|
||||||
|
assert_eq!(ids, vec!["dm-local", "e1"]);
|
||||||
|
|
||||||
|
// Same-id: the incoming expiry side wins, the local del-marker /
|
||||||
|
// abort-mpu side is authoritative and an incoming del-marker field is
|
||||||
|
// discarded at the trust boundary.
|
||||||
|
let mut local_mixed = lc_rule("r1", Some(1), None);
|
||||||
|
local_mixed.del_marker_expiration = Some(s3s::dto::DelMarkerExpiration { days: Some(3) });
|
||||||
|
local_mixed.abort_incomplete_multipart_upload = Some(s3s::dto::AbortIncompleteMultipartUpload {
|
||||||
|
days_after_initiation: Some(5),
|
||||||
|
});
|
||||||
|
let mut incoming_mixed = lc_rule("r1", Some(7), None);
|
||||||
|
incoming_mixed.del_marker_expiration = Some(s3s::dto::DelMarkerExpiration { days: Some(9) });
|
||||||
|
|
||||||
|
let merged =
|
||||||
|
merge_incoming_lifecycle_config(Some(lc_config(vec![incoming_mixed])), Some(lc_config(vec![local_mixed])), None)
|
||||||
|
.expect("merge should keep rules");
|
||||||
|
let r1 = &merged.rules[0];
|
||||||
|
assert_eq!(r1.expiration.as_ref().and_then(|e| e.days), Some(7));
|
||||||
|
assert_eq!(r1.del_marker_expiration.as_ref().and_then(|d| d.days), Some(3), "local del-marker wins");
|
||||||
|
assert_eq!(
|
||||||
|
r1.abort_incomplete_multipart_upload
|
||||||
|
.as_ref()
|
||||||
|
.and_then(|a| a.days_after_initiation),
|
||||||
|
Some(5),
|
||||||
|
"local abort-mpu wins"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Only a well-delimited zero-rule `<LifecycleConfiguration>` document is
|
||||||
|
/// the delete statement; truncated or foreign payloads must be rejected,
|
||||||
|
/// not treated as a delete that erases local expiry rules.
|
||||||
|
#[test]
|
||||||
|
fn test_zero_rule_lifecycle_tombstone_recognition() {
|
||||||
|
assert!(is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration><ExpiryUpdatedAt>2026-01-01T00:00:00Z</ExpiryUpdatedAt></LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
assert!(is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<?xml version=\"1.0\" encoding=\"UTF-8\"?>\n<LifecycleConfiguration xmlns=\"http://s3.amazonaws.com/doc/2006-03-01/\"></LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
assert!(is_zero_rule_lifecycle_tombstone(b"<LifecycleConfiguration/>"));
|
||||||
|
|
||||||
|
// Documents with rules are not tombstones (they must parse strictly).
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration><Rule><ID>x</ID></Rule></LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
// Truncated / malformed / foreign payloads are rejected.
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(b"<LifecycleConfiguration><ExpiryUpdatedAt>"));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(b"<LifecycleConfiguration><Rule></Broken>"));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(b"garbage"));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(b"<SomethingElse></SomethingElse>"));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(b""));
|
||||||
|
// Malformed children inside a well-delimited root are still rejected
|
||||||
|
// (second review round): a dangling open tag, stray text, an
|
||||||
|
// unclosed child, or nested markup is not a tombstone.
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration><ExpiryUpdatedAt></LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration>stray text</LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration><A><Rule/></A></LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration><Marker/></LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration><ExpiryUpdatedAt>&bogus;</ExpiryUpdatedAt></LifecycleConfiguration>"
|
||||||
|
));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(b"<evil:LifecycleConfiguration/>"));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(
|
||||||
|
b"<LifecycleConfiguration><ExpiryUpdatedAt>2026-01-01T00:00:00Z</ExpiryUpdatedAt></LifecycleConfiguration><Marker/>"
|
||||||
|
));
|
||||||
|
assert!(!is_zero_rule_lifecycle_tombstone(b"<LifecycleConfiguration/>&bogus;"));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_lifecycle_merge_holds_metadata_transaction_across_read_and_write() {
|
||||||
|
let source = include_str!("site_replication.rs");
|
||||||
|
let apply = source
|
||||||
|
.split("async fn apply_bucket_meta_item")
|
||||||
|
.nth(1)
|
||||||
|
.and_then(|rest| rest.split("fn group_info_requires_upsert").next())
|
||||||
|
.expect("apply_bucket_meta_item source");
|
||||||
|
let acquire = apply
|
||||||
|
.find("acquire_bucket_metadata_transaction_lock_for_incarnation")
|
||||||
|
.expect("lifecycle merge transaction acquisition");
|
||||||
|
let read = apply.find("get_config_from_disk").expect("fresh lifecycle config read");
|
||||||
|
let write = apply
|
||||||
|
.find("update_under_transaction_lock")
|
||||||
|
.expect("lifecycle config write under transaction");
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
acquire < read && read < write,
|
||||||
|
"the transaction must span the lifecycle read, merge, and write"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The staleness axis an incoming lc-config item must beat: the expiry
|
||||||
|
/// axis when present; the whole-config write time only for deleted or
|
||||||
|
/// legacy-with-expiry state; epoch for a transition-only config (its
|
||||||
|
/// whole-config time moves on transition edits and must not shadow
|
||||||
|
/// independent peer expiry updates — review finding).
|
||||||
|
#[test]
|
||||||
|
fn test_local_lifecycle_staleness_axis_selection() {
|
||||||
|
let whole = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
|
||||||
|
let axis_ts = OffsetDateTime::from_unix_timestamp(1_600_000_000).expect("timestamp");
|
||||||
|
|
||||||
|
let mut with_axis = lc_config(vec![lc_rule("e1", Some(7), None)]);
|
||||||
|
with_axis.expiry_updated_at = Some(s3s::dto::Timestamp::from(axis_ts));
|
||||||
|
assert_eq!(local_lifecycle_staleness_axis(Some(&with_axis), whole), axis_ts);
|
||||||
|
|
||||||
|
let legacy_with_expiry = lc_config(vec![lc_rule("e1", Some(7), None)]);
|
||||||
|
assert_eq!(local_lifecycle_staleness_axis(Some(&legacy_with_expiry), whole), whole);
|
||||||
|
|
||||||
|
let transition_only = lc_config(vec![lc_rule("t1", None, Some(30))]);
|
||||||
|
assert_eq!(
|
||||||
|
local_lifecycle_staleness_axis(Some(&transition_only), whole),
|
||||||
|
OffsetDateTime::UNIX_EPOCH,
|
||||||
|
"a transition-only config has no expiry state to protect"
|
||||||
|
);
|
||||||
|
|
||||||
|
assert_eq!(local_lifecycle_staleness_axis(None, whole), whole, "deletion lower bound");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 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.
|
||||||
|
|||||||
@@ -321,6 +321,34 @@ pub(crate) mod metadata_sys {
|
|||||||
crate::storage::storage_api::acquire_bucket_metadata_transaction_lock(bucket).await
|
crate::storage::storage_api::acquire_bucket_metadata_transaction_lock(bucket).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn acquire_bucket_metadata_transaction_lock_for_incarnation(
|
||||||
|
bucket: &str,
|
||||||
|
expected_incarnation_id: uuid::Uuid,
|
||||||
|
) -> Result<super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard> {
|
||||||
|
super::ecstore_bucket::metadata_sys::acquire_bucket_metadata_transaction_lock_for_incarnation(
|
||||||
|
bucket,
|
||||||
|
expected_incarnation_id,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn update_under_transaction_lock(
|
||||||
|
guard: &super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard,
|
||||||
|
bucket: &str,
|
||||||
|
config_file: &str,
|
||||||
|
data: Vec<u8>,
|
||||||
|
) -> Result<OffsetDateTime> {
|
||||||
|
super::ecstore_bucket::metadata_sys::update_under_transaction_lock(guard, bucket, config_file, data).await
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn delete_under_transaction_lock(
|
||||||
|
guard: &super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard,
|
||||||
|
bucket: &str,
|
||||||
|
config_file: &str,
|
||||||
|
) -> Result<OffsetDateTime> {
|
||||||
|
super::ecstore_bucket::metadata_sys::delete_under_transaction_lock(guard, bucket, config_file).await
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) async fn update_bucket_targets_under_transaction_lock(
|
pub(crate) async fn update_bucket_targets_under_transaction_lock(
|
||||||
guard: &super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard,
|
guard: &super::ecstore_bucket::metadata_sys::BucketMetadataMutationGuard,
|
||||||
bucket: &str,
|
bucket: &str,
|
||||||
|
|||||||
@@ -1158,6 +1158,19 @@ fn lifecycle_has_expiry_rules(config: &BucketLifecycleConfiguration) -> bool {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Status-independent presence of the expiry subset that site replication
|
||||||
|
/// propagates (`replicateILMExpiry`): expiration / noncurrent-version
|
||||||
|
/// expiration only. Distinct from [`lifecycle_has_expiry_rules`], which
|
||||||
|
/// filters on ENABLED for scanner scheduling — editing a Disabled expiry rule
|
||||||
|
/// must still advance the replication axis. Del-marker expiration and
|
||||||
|
/// abort-multipart are site-local and never travel.
|
||||||
|
fn lifecycle_rules_have_expiry(config: &BucketLifecycleConfiguration) -> bool {
|
||||||
|
config
|
||||||
|
.rules
|
||||||
|
.iter()
|
||||||
|
.any(|rule| rule.expiration.is_some() || rule.noncurrent_version_expiration.is_some())
|
||||||
|
}
|
||||||
|
|
||||||
fn lifecycle_has_abort_multipart_rules(config: &BucketLifecycleConfiguration) -> bool {
|
fn lifecycle_has_abort_multipart_rules(config: &BucketLifecycleConfiguration) -> bool {
|
||||||
config.rules.iter().any(|rule| {
|
config.rules.iter().any(|rule| {
|
||||||
rule.status == ExpirationStatus::from_static(ExpirationStatus::ENABLED)
|
rule.status == ExpirationStatus::from_static(ExpirationStatus::ENABLED)
|
||||||
@@ -2186,7 +2199,24 @@ impl DefaultBucketUsecase {
|
|||||||
return Err(s3_error!(InvalidArgument, "{err}"));
|
return Err(s3_error!(InvalidArgument, "{err}"));
|
||||||
}
|
}
|
||||||
|
|
||||||
input_cfg.expiry_updated_at = Some(Timestamp::from(time::OffsetDateTime::now_utc()));
|
// Stamp the expiry axis only when the expiry subset can have changed
|
||||||
|
// (MinIO: HasExpiry() || expiryRuleRemoved). Site-replication peers
|
||||||
|
// judge lc-config staleness on this axis; a transition-only edit that
|
||||||
|
// advanced it would let this site's stale expiry subset shadow — and
|
||||||
|
// roll back — a newer peer expiry edit fleet-wide.
|
||||||
|
let previous_expiry_updated_at = match metadata_sys::get_lifecycle_config(&bucket).await {
|
||||||
|
Ok((previous, _)) => {
|
||||||
|
if lifecycle_rules_have_expiry(&input_cfg) || lifecycle_rules_have_expiry(&previous) {
|
||||||
|
Some(Timestamp::from(time::OffsetDateTime::now_utc()))
|
||||||
|
} else {
|
||||||
|
previous.expiry_updated_at
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// No previous config (or unreadable): stamping is the
|
||||||
|
// conservative pre-existing behavior.
|
||||||
|
Err(_) => lifecycle_rules_have_expiry(&input_cfg).then(|| Timestamp::from(time::OffsetDateTime::now_utc())),
|
||||||
|
};
|
||||||
|
input_cfg.expiry_updated_at = previous_expiry_updated_at;
|
||||||
let data = serialize_config(&input_cfg)?;
|
let data = serialize_config(&input_cfg)?;
|
||||||
update_bucket_config_for_incarnation(&bucket, BUCKET_LIFECYCLE_CONFIG, data, expected_incarnation_id)
|
update_bucket_config_for_incarnation(&bucket, BUCKET_LIFECYCLE_CONFIG, data, expected_incarnation_id)
|
||||||
.await
|
.await
|
||||||
@@ -2197,7 +2227,14 @@ impl DefaultBucketUsecase {
|
|||||||
let mut item = sr_bucket_meta_item(bucket.clone(), "lc-config");
|
let mut item = sr_bucket_meta_item(bucket.clone(), "lc-config");
|
||||||
item.expiry_lc_config =
|
item.expiry_lc_config =
|
||||||
Some(serialize_config(&input_cfg).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?);
|
Some(serialize_config(&input_cfg).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?);
|
||||||
item.expiry_updated_at = item.updated_at;
|
// The item travels with the expiry axis, not the wall clock: a site
|
||||||
|
// whose expiry knowledge is old (or absent — UNIX_EPOCH) must not
|
||||||
|
// out-rank newer peer expiry state at the receivers.
|
||||||
|
item.expiry_updated_at = input_cfg
|
||||||
|
.expiry_updated_at
|
||||||
|
.clone()
|
||||||
|
.map(time::OffsetDateTime::from)
|
||||||
|
.or(Some(time::OffsetDateTime::UNIX_EPOCH));
|
||||||
if let Err(err) = site_replication_bucket_meta_hook(item).await {
|
if let Err(err) = site_replication_bucket_meta_hook(item).await {
|
||||||
warn!(bucket = %bucket, error = ?err, "site replication bucket lifecycle hook failed");
|
warn!(bucket = %bucket, error = ?err, "site replication bucket lifecycle hook failed");
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user