From 68499b6549af2a34ef31c8ade89456e99cd69062 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Sat, 5 Sep 2026 16:00:56 +0800 Subject: [PATCH] fix(site-replication): gate policy, mapping and group items on source updated_at (backlog#2291) The `policy`, `policy-mapping` and `group-info` receive paths applied every incoming item unconditionally, so a delayed older grant (wide policy body, old mapping, old group add) overwrote a newer revoke on the peer. `iam-user` and `service-account` already compared the item's `updatedAt` with the local record. Route the three paths through one pure verdict helper: an item older than the local record is acknowledged without being applied; items without a source timestamp and items targeting an absent record keep today's behaviour (older peers, idempotent deletes from backlog#2071). Deletes are gated the same way so an older delete cannot remove a newer record. The group record's own timestamp now moves on every membership and status change instead of staying at creation, so the gate judges group items against the last change. Add the IamSys accessors the gate reads (`get_policy_doc`, `get_mapped_policy_record`, `get_group_info`). (cherry picked from commit 98c32093406cb47014b7eda2fe139f01079de337) --- crates/iam/src/manager.rs | 42 +++- crates/iam/src/sys.rs | 16 ++ rustfs/src/admin/handlers/site_replication.rs | 227 +++++++++++++++++- 3 files changed, 271 insertions(+), 14 deletions(-) diff --git a/crates/iam/src/manager.rs b/crates/iam/src/manager.rs index 384ed9a0a..105254456 100644 --- a/crates/iam/src/manager.rs +++ b/crates/iam/src/manager.rs @@ -429,6 +429,27 @@ where } } + /// The cached mapping record for one user or group, looked up in the same + /// cache partition `policy_db_set` writes it to (group / STS / regular+service + /// user). `None` when no mapping is stored. + pub async fn get_mapped_policy_record(&self, name: &str, user_type: UserType, is_group: bool) -> Option { + let cache = self.cache.snapshot(); + if is_group { + cache.group_policies.get(name).cloned() + } else if user_type == UserType::Sts { + cache.sts_policies.get(name).cloned() + } else { + cache.user_policies.get(name).cloned() + } + } + + /// The cached group record (members, status, own timestamp) without the + /// mapped-policy overlay `get_group_description` applies. `None` when the + /// group does not exist. + pub async fn get_group_info(&self, name: &str) -> Option { + self.cache.snapshot().groups.get(name).cloned() + } + pub async fn get_policy(&self, name: &str) -> Result { if name.is_empty() { return Err(Error::InvalidArgument); @@ -1693,6 +1714,10 @@ where } } + // The group's own timestamp moves with every membership or status + // change: site replication judges an incoming group item against it + // (backlog#2291), so it must reflect the last change, not creation. + let now = OffsetDateTime::now_utc(); let gi = match cache.groups.get(group) { Some(res) => { let mut gi = res.clone(); @@ -1701,6 +1726,7 @@ where uniq_set.extend(members.iter().cloned()); gi.members = uniq_set.into_iter().collect(); + gi.update_at = Some(now); gi } None => GroupInfo::new(members.clone()), @@ -1709,8 +1735,7 @@ where self.api.save_group_info(group, gi.clone()).await?; - let now = self.cache.with_write_lock(|cache| { - let now = OffsetDateTime::now_utc(); + self.cache.with_write_lock(|cache| { cache.add_or_update_group(group, &gi, now); let user_group_memberships = Arc::clone(&cache.state().user_group_memberships); @@ -1719,7 +1744,6 @@ where m.insert(group.to_string()); cache.add_or_update_user_group_membership(member, &m, now); }); - now }); Ok(now) @@ -1743,12 +1767,14 @@ where } else { gi.status = STATUS_DISABLED.to_owned(); } + let now = OffsetDateTime::now_utc(); + gi.update_at = Some(now); self.api.save_group_info(name, gi.clone()).await?; - self.cache.add_or_update_group(name, &gi, OffsetDateTime::now_utc()); + self.cache.add_or_update_group(name, &gi, now); - Ok(OffsetDateTime::now_utc()) + Ok(now) } pub async fn get_group_description(&self, name: &str) -> Result { @@ -1830,13 +1856,14 @@ where let s: HashSet<&String> = HashSet::from_iter(gi.members.iter()); let d: HashSet<&String> = HashSet::from_iter(members.iter()); gi.members = s.difference(&d).map(|v| v.to_string()).collect::>(); + let now = OffsetDateTime::now_utc(); + gi.update_at = Some(now); if !update_cache_only { self.api.save_group_info(name, gi.clone()).await?; } - let now = self.cache.with_write_lock(|cache| { - let now = OffsetDateTime::now_utc(); + self.cache.with_write_lock(|cache| { cache.add_or_update_group(name, &gi, now); let user_group_memberships = Arc::clone(&cache.state().user_group_memberships); @@ -1847,7 +1874,6 @@ where cache.add_or_update_user_group_membership(member, &m, now); } }); - now }); Ok(now) diff --git a/crates/iam/src/sys.rs b/crates/iam/src/sys.rs index cd38d5f4e..58ccbeb9a 100644 --- a/crates/iam/src/sys.rs +++ b/crates/iam/src/sys.rs @@ -1055,6 +1055,22 @@ impl IamSys { self.store.get_group_description(group).await } + /// The stored group record itself (see `IamCache::get_group_info`). + pub async fn get_group_info(&self, group: &str) -> Option { + self.store.get_group_info(group).await + } + + /// The stored policy document, `Error::NoSuchPolicy` when absent. + pub async fn get_policy_doc(&self, name: &str) -> Result { + self.store.get_policy_doc(name).await + } + + /// The stored mapping record for one user or group (see + /// `IamCache::get_mapped_policy_record`). + pub async fn get_mapped_policy_record(&self, name: &str, user_type: UserType, is_group: bool) -> Option { + self.store.get_mapped_policy_record(name, user_type, is_group).await + } + pub async fn list_groups_load(&self) -> Result> { self.store.update_groups().await } diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index eca73bb40..f8f245d1a 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -5385,6 +5385,38 @@ fn is_stale_update(local_updated_at: OffsetDateTime, incoming_updated_at: Option incoming_updated_at.is_some_and(|incoming_updated_at| incoming_updated_at < local_updated_at) } +/// Verdict for an incoming IAM item judged against the local record it would +/// overwrite or delete (backlog#2291). +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum IamItemVerdict { + /// Apply the item: there is no local record, the item carries no source + /// timestamp (older peer), or it is at least as new as the local record. + Apply, + /// The local record was written from a newer source change; acknowledge the + /// item without touching the record. Covers both directions: a delayed + /// grant must not undo a newer revoke, and a delayed revoke must not undo a + /// newer grant. + SkipStale, +} + +/// Ordering rule shared by the `policy`, `policy-mapping` and `group-info` +/// item paths (and matching `iam-user` / `service-account`). +/// +/// `local_record_updated_at` is `None` when the targeted record does not +/// exist locally: nothing can be stale relative to an absent record, so a +/// create is applied and a delete falls through to the idempotent no-op paths +/// (backlog#2071). A record that exists but predates timestamps passes +/// `Some(UNIX_EPOCH)` and therefore never rejects an item. +fn judge_iam_item_staleness( + local_record_updated_at: Option, + incoming_updated_at: Option, +) -> IamItemVerdict { + match local_record_updated_at { + Some(local_updated_at) if is_stale_update(local_updated_at, incoming_updated_at) => IamItemVerdict::SkipStale, + _ => IamItemVerdict::Apply, + } +} + fn bucket_meta_local_updated_at( bucket_meta: &crate::admin::storage_api::bucket::metadata::BucketMetadata, config_file: &str, @@ -5728,9 +5760,9 @@ async fn apply_iam_item(item: SRIAMItem) -> S3Result<()> { let incoming_updated_at = item.updated_at; match item.r#type.as_str() { - "policy" => apply_iam_policy_item(&iam_sys, &item.name, item.policy).await, - "policy-mapping" => apply_iam_policy_mapping_item(&iam_sys, item.policy_mapping).await, - "group-info" => apply_iam_group_info_item(&iam_sys, item.group_info).await, + "policy" => apply_iam_policy_item(&iam_sys, &item.name, item.policy, incoming_updated_at).await, + "policy-mapping" => apply_iam_policy_mapping_item(&iam_sys, item.policy_mapping, incoming_updated_at).await, + "group-info" => apply_iam_group_info_item(&iam_sys, item.group_info, incoming_updated_at).await, // MinIO madmin-go sends `SRIAMItemSTSAcc = "sts-account"`. The legacy alias // `sts-credential` (emitted by older RustFS releases) stays accepted permanently // so mixed-version RustFS sites keep replicating STS credentials during rolling @@ -5746,7 +5778,22 @@ async fn apply_iam_item(item: SRIAMItem) -> S3Result<()> { } } -async fn apply_iam_policy_item(iam_sys: &IamSys, name: &str, policy: Option) -> S3Result<()> { +async fn apply_iam_policy_item( + iam_sys: &IamSys, + name: &str, + policy: Option, + incoming_updated_at: Option, +) -> S3Result<()> { + // Judge the item against the local document's own timestamp so a delayed + // older body (or delete) cannot overwrite a newer edit (backlog#2291). + let local_updated_at = match iam_sys.get_policy_doc(name).await { + Ok(doc) => Some(doc.update_date.unwrap_or(OffsetDateTime::UNIX_EPOCH)), + Err(err) if rustfs_iam::error::is_err_no_such_policy(&err) => None, + Err(err) => return Err(ApiError::from(err).into()), + }; + if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale { + return Ok(()); + } if let Some(policy) = policy { let policy: Policy = serde_json::from_value(policy).map_err(|e| s3_error!(InvalidRequest, "invalid policy body: {}", e))?; @@ -5764,11 +5811,26 @@ async fn apply_iam_policy_item(iam_sys: &IamSys, name: &str, policy Ok(()) } -async fn apply_iam_policy_mapping_item(iam_sys: &IamSys, policy_mapping: Option) -> S3Result<()> { +async fn apply_iam_policy_mapping_item( + iam_sys: &IamSys, + policy_mapping: Option, + incoming_updated_at: Option, +) -> S3Result<()> { let Some(mapping) = policy_mapping else { return Err(s3_error!(InvalidRequest, "policyMapping is required")); }; let user_type = user_type_from_sr_wire(mapping.user_type).ok_or_else(|| s3_error!(InvalidRequest, "invalid userType"))?; + // Judge the item against the stored mapping's timestamp so a delayed older + // attach (or an older detach, `policy == ""`) cannot overwrite a newer one + // (backlog#2291). A detach that already removed the mapping leaves no + // record, and an item targeting an absent mapping is applied as before. + let local_updated_at = iam_sys + .get_mapped_policy_record(&mapping.user_or_group, user_type, mapping.is_group) + .await + .map(|record| record.update_at); + if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale { + return Ok(()); + } iam_sys .policy_db_set(&mapping.user_or_group, user_type, mapping.is_group, &mapping.policy) .await @@ -5776,11 +5838,26 @@ async fn apply_iam_policy_mapping_item(iam_sys: &IamSys, policy_map Ok(()) } -async fn apply_iam_group_info_item(iam_sys: &IamSys, group_info: Option) -> S3Result<()> { +async fn apply_iam_group_info_item( + iam_sys: &IamSys, + group_info: Option, + incoming_updated_at: Option, +) -> S3Result<()> { let Some(group_info) = group_info else { return Err(s3_error!(InvalidRequest, "groupInfo is required")); }; let update = group_info.update_req; + // The record is the group itself: its own timestamp moves on every + // membership or status change, so a delayed older add cannot re-add a + // member a newer removal took out, and a delayed older removal (or group + // delete) cannot undo a newer add (backlog#2291). + let local_updated_at = iam_sys + .get_group_info(&update.group) + .await + .map(|group| group.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH)); + if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale { + return Ok(()); + } if !group_info_requires_upsert(&update) { // Idempotent removal: a replayed deletion may find the group or a // member already gone (deleted here earlier, or the user tombstone @@ -12175,6 +12252,144 @@ mod tests { assert!(!is_stale_update(local, None)); } + /// Minimal model of one replicated IAM record (a policy document body, a + /// user/group mapping, or a group's member set) as the apply paths treat + /// it: `None` is "absent", `Some((content, stamp))` is the local record + /// with the timestamp of the change that last wrote it. Applying an item + /// goes through `judge_iam_item_staleness` exactly like the three apply + /// functions do; a delete (`incoming == None`) on an absent record is the + /// idempotent no-op of backlog#2071. + fn apply_iam_item_to_model( + record: &mut Option<(&'static str, OffsetDateTime)>, + incoming: Option<&'static str>, + incoming_updated_at: Option, + ) -> IamItemVerdict { + let verdict = judge_iam_item_staleness(record.map(|(_, stamp)| stamp), incoming_updated_at); + if verdict == IamItemVerdict::Apply { + *record = incoming.map(|content| (content, incoming_updated_at.unwrap_or(OffsetDateTime::UNIX_EPOCH))); + } + verdict + } + + fn at(seconds: i64) -> OffsetDateTime { + OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(seconds) + } + + /// backlog#2291: a revoke (narrowed policy body, detached mapping, member + /// removed from the group) followed by the delayed delivery of the older + /// grant must leave the revoke in place. + #[test] + fn test_iam_item_stale_grant_after_revoke_is_not_applied() { + let mut record = None; + assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(10))), IamItemVerdict::Apply); + assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(20))), IamItemVerdict::Apply); + + // The older grant is redelivered (retry drain, slow peer) after the revoke. + assert_eq!( + apply_iam_item_to_model(&mut record, Some("grant"), Some(at(10))), + IamItemVerdict::SkipStale, + "a grant older than the local revoke must be acknowledged without being applied" + ); + assert_eq!(record, Some(("revoke", at(20))), "the revoke must survive the stale grant"); + } + + /// backlog#2291: the mirror image — a grant followed by the delayed delivery + /// of an older revoke (older body, older detach, older member removal, or + /// an older delete) must leave the grant in place. + #[test] + fn test_iam_item_stale_revoke_after_grant_is_not_applied() { + let mut record = None; + assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(10))), IamItemVerdict::Apply); + assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(20))), IamItemVerdict::Apply); + + assert_eq!( + apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(10))), + IamItemVerdict::SkipStale, + "a revoke older than the local grant must not be applied" + ); + assert_eq!( + apply_iam_item_to_model(&mut record, None, Some(at(15))), + IamItemVerdict::SkipStale, + "a delete older than the local record must not remove it" + ); + assert_eq!(record, Some(("grant", at(20)))); + } + + /// backlog#2291: an item at least as new as the local record is applied, + /// including a newer delete; equal timestamps are not stale (same rule as + /// `iam-user`). + #[test] + fn test_iam_item_newer_than_local_record_is_applied() { + let mut record = Some(("grant", at(20))); + assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), Some(at(20))), IamItemVerdict::Apply); + assert_eq!(record, Some(("revoke", at(20)))); + + assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(30))), IamItemVerdict::Apply); + assert_eq!(record, Some(("grant", at(30)))); + + assert_eq!(apply_iam_item_to_model(&mut record, None, Some(at(40))), IamItemVerdict::Apply); + assert_eq!(record, None, "a newer delete removes the record"); + } + + /// backlog#2291: peers that predate item timestamps keep today's + /// last-writer-wins behaviour — an item without `updatedAt` is applied even + /// over a newer local record. + #[test] + fn test_iam_item_without_source_timestamp_is_applied() { + let mut record = Some(("grant", at(20))); + assert_eq!(apply_iam_item_to_model(&mut record, Some("revoke"), None), IamItemVerdict::Apply); + assert_eq!(record.map(|(content, _)| content), Some("revoke")); + + assert_eq!(judge_iam_item_staleness(Some(at(20)), None), IamItemVerdict::Apply); + assert_eq!(judge_iam_item_staleness(None, None), IamItemVerdict::Apply); + } + + /// backlog#2291: nothing is stale relative to an absent record. A create + /// with any timestamp is applied, and a delete falls through to the + /// idempotent no-op paths (backlog#2071) instead of being judged. + #[test] + fn test_iam_item_targeting_absent_record_is_applied() { + assert_eq!(judge_iam_item_staleness(None, Some(at(1))), IamItemVerdict::Apply); + + let mut record = None; + assert_eq!(apply_iam_item_to_model(&mut record, None, Some(at(1))), IamItemVerdict::Apply); + assert_eq!(record, None); + assert_eq!(apply_iam_item_to_model(&mut record, Some("grant"), Some(at(1))), IamItemVerdict::Apply); + assert_eq!(record, Some(("grant", at(1)))); + + // A record that predates timestamps is reported as UNIX_EPOCH by the + // apply paths and therefore never rejects an item. + assert_eq!( + judge_iam_item_staleness(Some(OffsetDateTime::UNIX_EPOCH), Some(at(1))), + IamItemVerdict::Apply + ); + } + + /// Cheap wiring guard for backlog#2291: every one of the `policy`, + /// `policy-mapping` and `group-info` apply paths must route through the + /// shared staleness verdict before it writes or deletes anything. The + /// ordering rule itself is covered by the `test_iam_item_*` behaviour + /// tests above; this only pins that no path bypasses it again. + #[test] + fn test_iam_policy_mapping_and_group_items_gate_on_incoming_updated_at() { + let source = include_str!("site_replication.rs"); + for (start, end) in [ + ("async fn apply_iam_policy_item(", "async fn apply_iam_policy_mapping_item("), + ("async fn apply_iam_policy_mapping_item(", "async fn apply_iam_group_info_item("), + ("async fn apply_iam_group_info_item(", "async fn apply_iam_sts_account_item("), + ] { + let body = source + .split(start) + .nth(1) + .and_then(|rest| rest.split(end).next()) + .expect(start); + assert!( + body.contains("judge_iam_item_staleness(local_updated_at, incoming_updated_at)"), + "{start} must judge the item against the local record before applying it" + ); + } + } + #[test] fn test_apply_state_edit_req_only_updates_ilm_expiry_flags() { let mut state = SiteReplicationState::default();