diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index f8f245d1a..de8fe6565 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -87,7 +87,7 @@ use std::sync::{LazyLock, Mutex as StdMutex}; use std::time::Duration; use time::OffsetDateTime; use tokio::sync::Mutex; -use tracing::{info, warn}; +use tracing::{debug, info, warn}; use url::Url; use url::form_urlencoded; use uuid::Uuid; @@ -3497,6 +3497,7 @@ fn remove_sites(mut state: SiteReplicationState, req: SRRemoveReq) -> SiteReplic state.resync_status.clear(); state.retry_queue.clear(); state.iam_deletion_replays.clear(); + state.iam_deletion_marks.clear(); state.pending_endpoint_refresh = None; state.updated_at = Some(OffsetDateTime::now_utc()); return state; @@ -3508,6 +3509,7 @@ fn remove_sites(mut state: SiteReplicationState, req: SRRemoveReq) -> SiteReplic state.resync_status.clear(); state.retry_queue.clear(); state.iam_deletion_replays.clear(); + state.iam_deletion_marks.clear(); state.pending_endpoint_refresh = None; state.updated_at = Some(OffsetDateTime::now_utc()); return state; @@ -5758,23 +5760,65 @@ async fn apply_iam_item(item: SRIAMItem) -> S3Result<()> { return Err(s3_error!(InvalidRequest, "iam not init")); }; let incoming_updated_at = item.updated_at; + let deletion_mark_entities = iam_item_deletion_mark_entities(&item); - match item.r#type.as_str() { - "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, + let verdict = match item.r#type.as_str() { + "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 // upgrades; it is a compatibility layer, not temporary code. - SR_IAM_ITEM_STS_ACC | SR_IAM_ITEM_STS_ACC_LEGACY => apply_iam_sts_account_item(&iam_sys, item.sts_credential).await, - "iam-user" => apply_iam_user_item(&iam_sys, item.iam_user, incoming_updated_at).await, - "service-account" => apply_iam_service_account_item(&iam_sys, item.svc_acc_change, incoming_updated_at).await, - _ => Err(s3_error!( - NotImplemented, - "site replication IAM item type `{}` is not supported", - item.r#type - )), + SR_IAM_ITEM_STS_ACC | SR_IAM_ITEM_STS_ACC_LEGACY => { + apply_iam_sts_account_item(&iam_sys, item.sts_credential).await?; + IamItemVerdict::Apply + } + "iam-user" => apply_iam_user_item(&iam_sys, item.iam_user, incoming_updated_at).await?, + "service-account" => { + apply_iam_service_account_item(&iam_sys, item.svc_acc_change, incoming_updated_at).await?; + IamItemVerdict::Apply + } + _ => { + return Err(s3_error!( + NotImplemented, + "site replication IAM item type `{}` is not supported", + item.r#type + )); + } + }; + // A committed deletion leaves no record for the gate to judge later items + // against, so its source timestamp is kept as a mark (backlog#2291). The + // mark is part of applying the deletion: failing here makes the sender + // retry the (idempotent) deletion rather than leave a revoke that a stale + // grant could still undo. + if verdict == IamItemVerdict::Apply + && let Some(deleted_at) = incoming_updated_at.filter(|_| !deletion_mark_entities.is_empty()) + { + commit_iam_deletion_marks(deletion_mark_entities, deleted_at).await?; + } + Ok(()) +} + +/// The deletion mark consulted by the staleness gate when the targeted record +/// is absent: the newest recorded deletion of any of `entities`. An +/// unreadable state falls back to today's behaviour (no mark, the item is +/// applied) — the gate must not turn a state-object outage into rejected +/// IAM replication. +async fn local_iam_deletion_mark(entities: &[String]) -> Option { + match load_site_replication_state().await { + Ok(state) => iam_deletion_mark(&state, entities), + Err(err) => { + debug!( + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + result = "iam_deletion_mark_unavailable", + error = ?err, + "site replication state unreadable; applying IAM item without a deletion mark" + ); + None + } } } @@ -5783,16 +5827,19 @@ async fn apply_iam_policy_item( name: &str, policy: Option, incoming_updated_at: Option, -) -> S3Result<()> { +) -> 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). + // older body (or delete) cannot overwrite a newer edit; once the document + // is deleted, its deletion mark stands in for it (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) if rustfs_iam::error::is_err_no_such_policy(&err) => { + local_iam_deletion_mark(&[iam_policy_deletion_mark_entity(name)]).await + } Err(err) => return Err(ApiError::from(err).into()), }; if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale { - return Ok(()); + return Ok(IamItemVerdict::SkipStale); } if let Some(policy) = policy { let policy: Policy = @@ -5808,41 +5855,51 @@ async fn apply_iam_policy_item( Err(err) => return Err(ApiError::from(err).into()), } } - Ok(()) + Ok(IamItemVerdict::Apply) } async fn apply_iam_policy_mapping_item( iam_sys: &IamSys, policy_mapping: Option, incoming_updated_at: Option, -) -> S3Result<()> { +) -> 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 + // (backlog#2291). A detach removes the mapping outright, so once it is + // gone the detach's deletion mark stands in for the record. + let local_updated_at = match iam_sys .get_mapped_policy_record(&mapping.user_or_group, user_type, mapping.is_group) .await - .map(|record| record.update_at); + { + Some(record) => Some(record.update_at), + None => { + local_iam_deletion_mark(&[iam_policy_mapping_deletion_mark_entity( + &mapping.user_or_group, + mapping.user_type, + mapping.is_group, + )]) + .await + } + }; if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale { - return Ok(()); + return Ok(IamItemVerdict::SkipStale); } iam_sys .policy_db_set(&mapping.user_or_group, user_type, mapping.is_group, &mapping.policy) .await .map_err(ApiError::from)?; - Ok(()) + Ok(IamItemVerdict::Apply) } async fn apply_iam_group_info_item( iam_sys: &IamSys, group_info: Option, incoming_updated_at: Option, -) -> S3Result<()> { +) -> S3Result { let Some(group_info) = group_info else { return Err(s3_error!(InvalidRequest, "groupInfo is required")); }; @@ -5850,13 +5907,25 @@ async fn apply_iam_group_info_item( // 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)); + // delete) cannot undo a newer add (backlog#2291). Once the group is gone + // the marks of its deletion and of its members' removals stand in for it, + // so a stale add cannot re-create it or re-add a removed member. + let local_updated_at = match iam_sys.get_group_info(&update.group).await { + Some(group) => Some(group.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH)), + None => { + let entities: Vec = std::iter::once(iam_group_deletion_mark_entity(&update.group)) + .chain( + update + .members + .iter() + .map(|member| iam_group_member_deletion_mark_entity(&update.group, member)), + ) + .collect(); + local_iam_deletion_mark(&entities).await + } + }; if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale { - return Ok(()); + return Ok(IamItemVerdict::SkipStale); } if !group_info_requires_upsert(&update) { // Idempotent removal: a replayed deletion may find the group or a @@ -5872,14 +5941,14 @@ async fn apply_iam_group_info_item( } } if members.is_empty() && !update.members.is_empty() { - return Ok(()); + return Ok(IamItemVerdict::Apply); } match iam_sys.remove_users_from_group(&update.group, members).await { Ok(_) => {} Err(err) if rustfs_iam::error::is_err_no_such_group(&err) => {} Err(err) => return Err(ApiError::from(err).into()), } - return Ok(()); + return Ok(IamItemVerdict::Apply); } iam_sys @@ -5890,7 +5959,7 @@ async fn apply_iam_group_info_item( .set_group_status(&update.group, matches!(update.status, GroupStatus::Enabled)) .await .map_err(ApiError::from)?; - Ok(()) + Ok(IamItemVerdict::Apply) } async fn apply_iam_sts_account_item(iam_sys: &IamSys, sts_credential: Option) -> S3Result<()> { @@ -5932,14 +6001,18 @@ async fn apply_iam_user_item( iam_sys: &IamSys, iam_user: Option, incoming_updated_at: Option, -) -> S3Result<()> { +) -> S3Result { let Some(user) = iam_user else { return Err(s3_error!(InvalidRequest, "iamUser is required")); }; - if let Some(local) = iam_sys.get_user(&user.access_key).await - && is_stale_update(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH), incoming_updated_at) - { - return Ok(()); + // Once the identity is deleted, its deletion mark stands in for the + // record so a stale re-create cannot resurrect it (backlog#2291). + let local_updated_at = match iam_sys.get_user(&user.access_key).await { + Some(local) => Some(local.update_at.unwrap_or(OffsetDateTime::UNIX_EPOCH)), + None => local_iam_deletion_mark(&[iam_user_deletion_mark_entity(&user.access_key)]).await, + }; + if judge_iam_item_staleness(local_updated_at, incoming_updated_at) == IamItemVerdict::SkipStale { + return Ok(IamItemVerdict::SkipStale); } if user.is_delete_req { iam_sys.delete_user(&user.access_key, true).await.map_err(ApiError::from)?; @@ -5960,7 +6033,7 @@ async fn apply_iam_user_item( .map_err(ApiError::from)?; } } - Ok(()) + Ok(IamItemVerdict::Apply) } async fn apply_iam_service_account_item( @@ -6176,6 +6249,7 @@ fn adopt_add_commit_state(state: &mut SiteReplicationState, next_state: SiteRepl sync_state_initialized, edit_generation: _, applied_edit_generations: _, + iam_deletion_marks: _, } = next_state; state.name = name; state.service_account_access_key = service_account_access_key; @@ -12365,10 +12439,151 @@ mod tests { ); } + /// The apply paths once the record is gone: the local timestamp the gate + /// sees is the deletion mark (or `None` when no deletion was recorded), + /// and a committed deletion records its source timestamp as the mark — + /// the same sequence `apply_iam_item` runs. + fn apply_iam_item_to_deleted_record_model( + marks: &mut SiteReplicationState, + entity: &str, + incoming_is_delete: bool, + incoming_updated_at: Option, + ) -> IamItemVerdict { + let entities = vec![entity.to_string()]; + let verdict = judge_iam_item_staleness(iam_deletion_mark(marks, &entities), incoming_updated_at); + if verdict == IamItemVerdict::Apply + && incoming_is_delete + && let Some(deleted_at) = incoming_updated_at + { + record_iam_deletion_marks(marks, &entities, deleted_at); + } + verdict + } + + /// backlog#2291 (real-VM case R6.3a of backlog#2080): a detach deletes the + /// mapping outright, so the older grant that arrives afterwards finds no + /// record — the deletion mark must stand in for it and reject the grant. + /// The same holds for a deleted policy document, user or group. + #[test] + fn test_iam_item_stale_grant_after_record_deletion_is_not_applied() { + let mut marks = SiteReplicationState::default(); + let entity = "policy-mapping:alice:0:false"; + + // The revoke (detach) is applied first: the record is gone, the mark stays. + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(20))), + IamItemVerdict::Apply + ); + assert_eq!(marks.iam_deletion_marks.get(entity), Some(&at(20))); + + // The older grant is delivered after the revoke. + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, false, Some(at(10))), + IamItemVerdict::SkipStale, + "a grant older than the recorded deletion must not re-create the record" + ); + // A replayed copy of the same revoke stays a no-op and keeps the mark. + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(20))), + IamItemVerdict::Apply + ); + assert_eq!(marks.iam_deletion_marks.get(entity), Some(&at(20))); + // An older replayed revoke is stale against the newer one. + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(5))), + IamItemVerdict::SkipStale + ); + assert_eq!( + marks.iam_deletion_marks.get(entity), + Some(&at(20)), + "an older deletion never lowers the mark" + ); + } + + /// backlog#2291: a mark only fences items older than the deletion. A grant + /// newer than (or as new as) the recorded deletion re-creates the record, + /// an unmarked entity and an item without a source timestamp keep today's + /// behaviour. + #[test] + fn test_iam_item_newer_than_deletion_mark_is_applied() { + let mut marks = SiteReplicationState::default(); + let entity = "policy:readonly"; + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, true, Some(at(20))), + IamItemVerdict::Apply + ); + + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, false, Some(at(20))), + IamItemVerdict::Apply, + "a grant as new as the deletion is not stale" + ); + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, false, Some(at(30))), + IamItemVerdict::Apply + ); + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, entity, false, None), + IamItemVerdict::Apply, + "an item from a peer without timestamps keeps last-writer-wins" + ); + assert_eq!( + apply_iam_item_to_deleted_record_model(&mut marks, "policy:other", false, Some(at(1))), + IamItemVerdict::Apply, + "no mark, no record: nothing to be stale against" + ); + assert_eq!( + judge_iam_item_staleness(iam_deletion_mark(&marks, &[]), Some(at(1))), + IamItemVerdict::Apply + ); + } + + /// backlog#2291: a group's removal marks are per member (plus the group + /// itself for a group delete), so with the group gone a stale add is + /// judged against the newest mark among the group and the members it + /// would add. + #[test] + fn test_iam_group_item_after_deletion_is_judged_against_member_marks() { + let mut marks = SiteReplicationState::default(); + let bob = iam_group_member_deletion_mark_entity("devs", "bob"); + let group = iam_group_deletion_mark_entity("devs"); + record_iam_deletion_marks(&mut marks, std::slice::from_ref(&bob), at(20)); + record_iam_deletion_marks(&mut marks, std::slice::from_ref(&group), at(30)); + + // The gate for an add of `bob` to the (deleted) group. + let add_bob = [group.clone(), bob.clone()]; + assert_eq!(iam_deletion_mark(&marks, &add_bob), Some(at(30))); + assert_eq!( + judge_iam_item_staleness(iam_deletion_mark(&marks, &add_bob), Some(at(25))), + IamItemVerdict::SkipStale + ); + assert_eq!( + judge_iam_item_staleness(iam_deletion_mark(&marks, &add_bob), Some(at(30))), + IamItemVerdict::Apply + ); + + // An add of `carol` to a group that was only ever partially emptied + // (no group delete) is judged against carol's own mark only. + marks.iam_deletion_marks.remove(&group); + let add_carol = [group.clone(), iam_group_member_deletion_mark_entity("devs", "carol")]; + assert_eq!(iam_deletion_mark(&marks, &add_carol), None); + assert_eq!( + judge_iam_item_staleness(iam_deletion_mark(&marks, &add_carol), Some(at(1))), + IamItemVerdict::Apply + ); + let add_bob = [group, bob]; + assert_eq!( + judge_iam_item_staleness(iam_deletion_mark(&marks, &add_bob), Some(at(10))), + IamItemVerdict::SkipStale + ); + } + /// 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 + /// `policy-mapping`, `group-info` and `iam-user` apply paths must route + /// through the shared staleness verdict before it writes or deletes + /// anything, and must fall back to the deletion mark when the record is + /// absent; `apply_iam_item` must record the mark of a committed deletion. + /// 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() { @@ -12377,6 +12592,7 @@ mod tests { ("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("), + ("async fn apply_iam_user_item(", "async fn apply_iam_service_account_item("), ] { let body = source .split(start) @@ -12387,7 +12603,20 @@ mod tests { body.contains("judge_iam_item_staleness(local_updated_at, incoming_updated_at)"), "{start} must judge the item against the local record before applying it" ); + assert!( + body.contains("local_iam_deletion_mark("), + "{start} must fall back to the deletion mark when the record is absent" + ); } + let dispatch = source + .split("async fn apply_iam_item(") + .nth(1) + .and_then(|rest| rest.split("async fn local_iam_deletion_mark(").next()) + .expect("apply_iam_item"); + assert!( + dispatch.contains("commit_iam_deletion_marks(deletion_mark_entities, deleted_at)"), + "apply_iam_item must record the mark of a deletion it committed" + ); } #[test] diff --git a/rustfs/src/site_replication/hooks.rs b/rustfs/src/site_replication/hooks.rs index 514864831..2af6aa000 100644 --- a/rustfs/src/site_replication/hooks.rs +++ b/rustfs/src/site_replication/hooks.rs @@ -1028,6 +1028,21 @@ pub async fn site_replication_iam_change_hook(item: SRIAMItem) -> S3Result<()> { let Some(runtime) = runtime_site_replication_targets().await? else { return Ok(()); }; + // A local revoke must out-rank a stale grant a peer delivers later, so its + // mark is committed before the broadcast (backlog#2291). The broadcast + // still goes out when the mark cannot be persisted: the peers' own records + // remain the primary gate, the mark only covers the deleted case. + if let Err(err) = record_iam_deletion_marks_for_item(&item).await { + warn!( + component = LOG_COMPONENT_ADMIN, + subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, + event = EVENT_ADMIN_SITE_REPLICATION_STATE, + item_type = %item.r#type, + result = "iam_deletion_mark_not_recorded", + error = ?err, + "failed to record local IAM deletion mark before broadcast" + ); + } let mut first_error: Option = None; for peer in runtime.state.peers.values() { if peer.deployment_id == runtime.local_peer.deployment_id diff --git a/rustfs/src/site_replication/retry.rs b/rustfs/src/site_replication/retry.rs index 2227b3763..a627e0b48 100644 --- a/rustfs/src/site_replication/retry.rs +++ b/rustfs/src/site_replication/retry.rs @@ -397,12 +397,12 @@ pub(crate) fn iam_deletion_replay_matches(record: &SiteReplicationIamDeletionRep /// newer revision of one another. pub(crate) fn iam_item_deletion_entity(item: &SRIAMItem) -> Option { match item.r#type.as_str() { - "policy" if item.policy.is_none() => Some(format!("policy:{}", item.name)), + "policy" if item.policy.is_none() => Some(iam_policy_deletion_mark_entity(&item.name)), "iam-user" => item .iam_user .as_ref() .filter(|user| user.is_delete_req) - .map(|user| format!("iam-user:{}", user.access_key)), + .map(|user| iam_user_deletion_mark_entity(&user.access_key)), "group-info" => item .group_info .as_ref() @@ -416,7 +416,7 @@ pub(crate) fn iam_item_deletion_entity(item: &SRIAMItem) -> Option { .policy_mapping .as_ref() .filter(|mapping| mapping.policy.is_empty()) - .map(|mapping| format!("policy-mapping:{}:{}:{}", mapping.user_or_group, mapping.user_type, mapping.is_group)), + .map(|mapping| iam_policy_mapping_deletion_mark_entity(&mapping.user_or_group, mapping.user_type, mapping.is_group)), "service-account" => item .svc_acc_change .as_ref() @@ -426,6 +426,82 @@ pub(crate) fn iam_item_deletion_entity(item: &SRIAMItem) -> Option { } } +/// The entities whose deletion a deletion-shaped IAM item commits, keyed the +/// way the receive-side staleness gate looks them up once the local record is +/// gone (backlog#2291); empty for creates and updates. Group member removal +/// yields one entity per removed member so a stale re-add of that member can +/// be judged, and a group delete (no members) yields the group itself. +pub(crate) fn iam_item_deletion_mark_entities(item: &SRIAMItem) -> Vec { + if item.r#type == "group-info" { + let Some(update) = item + .group_info + .as_ref() + .map(|group| &group.update_req) + .filter(|update| update.is_remove) + else { + return Vec::new(); + }; + if update.members.is_empty() { + return vec![iam_group_deletion_mark_entity(&update.group)]; + } + return update + .members + .iter() + .map(|member| iam_group_member_deletion_mark_entity(&update.group, member)) + .collect(); + } + iam_item_deletion_entity(item).into_iter().collect() +} + +pub(crate) fn iam_policy_deletion_mark_entity(name: &str) -> String { + format!("policy:{name}") +} + +pub(crate) fn iam_user_deletion_mark_entity(access_key: &str) -> String { + format!("iam-user:{access_key}") +} + +/// `user_type` is the SR wire integer, as carried by the item on both sides. +pub(crate) fn iam_policy_mapping_deletion_mark_entity(user_or_group: &str, user_type: i64, is_group: bool) -> String { + format!("policy-mapping:{user_or_group}:{user_type}:{is_group}") +} + +pub(crate) fn iam_group_deletion_mark_entity(group: &str) -> String { + format!("group:{group}") +} + +pub(crate) fn iam_group_member_deletion_mark_entity(group: &str, member: &str) -> String { + format!("group-member:{group}:{member}") +} + +/// Persist the deletion marks of `item` (its source `updated_at` per entity +/// of [`iam_item_deletion_mark_entities`]) through the state transaction. +/// No-op for creates/updates and for items without a source timestamp +/// (older peers): a mark without a source clock could not be ordered against +/// later items. Called before a local deletion is broadcast and after a +/// replicated deletion is applied, so both sides out-rank a stale grant that +/// arrives later. +pub(crate) async fn record_iam_deletion_marks_for_item(item: &SRIAMItem) -> S3Result<()> { + let entities = iam_item_deletion_mark_entities(item); + let Some(deleted_at) = item.updated_at.filter(|_| !entities.is_empty()) else { + return Ok(()); + }; + commit_iam_deletion_marks(entities, deleted_at).await +} + +/// [`record_iam_deletion_marks`] under the state transaction; the write is +/// skipped when no mark moves. +pub(crate) async fn commit_iam_deletion_marks(entities: Vec, deleted_at: OffsetDateTime) -> S3Result<()> { + update_site_replication_state_when_changed(move |state| { + Ok(if record_iam_deletion_marks(state, &entities, deleted_at) { + StateCommit::Changed(()) + } else { + StateCommit::Unchanged(()) + }) + }) + .await +} + /// Failure bookkeeping for one IAM item delivery: upsert the collapsed retry /// event and, when the item is a deletion, record its body for replay. Both /// live in the same state so the caller commits them in one transaction — a diff --git a/rustfs/src/site_replication/state.rs b/rustfs/src/site_replication/state.rs index 04619a9ac..0d9bf5f6e 100644 --- a/rustfs/src/site_replication/state.rs +++ b/rustfs/src/site_replication/state.rs @@ -64,6 +64,86 @@ pub(crate) struct SiteReplicationState { /// newer edit that already landed. #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] pub(crate) applied_edit_generations: BTreeMap, + /// Source timestamp of the newest IAM deletion committed on this site, + /// keyed by the deleted entity (`iam_item_deletion_mark_entities`). A + /// deletion leaves no local record to judge a later item against, so this + /// is what lets the receive-side staleness gate reject a grant that is + /// older than the revoke it would otherwise undo (backlog#2291). Bounded + /// by [`SITE_REPLICATION_IAM_DELETION_MARK_LIMIT`]; the oldest mark is + /// evicted first. + #[serde(default, with = "rfc3339_map", skip_serializing_if = "BTreeMap::is_empty")] + pub(crate) iam_deletion_marks: BTreeMap, +} + +/// Upper bound on [`SiteReplicationState::iam_deletion_marks`]. +pub(crate) const SITE_REPLICATION_IAM_DELETION_MARK_LIMIT: usize = 1024; + +/// Record that deletions of `entities` with source timestamp `deleted_at` +/// were committed here. Newest wins per entity: an older deletion never +/// lowers a mark. Returns whether the state changed. +pub(crate) fn record_iam_deletion_marks( + state: &mut SiteReplicationState, + entities: &[String], + deleted_at: OffsetDateTime, +) -> bool { + let mut changed = false; + for entity in entities { + if state + .iam_deletion_marks + .get(entity) + .is_some_and(|existing| *existing >= deleted_at) + { + continue; + } + state.iam_deletion_marks.insert(entity.clone(), deleted_at); + changed = true; + } + while state.iam_deletion_marks.len() > SITE_REPLICATION_IAM_DELETION_MARK_LIMIT { + let Some(oldest) = state + .iam_deletion_marks + .iter() + .min_by_key(|(_, deleted_at)| **deleted_at) + .map(|(entity, _)| entity.clone()) + else { + break; + }; + state.iam_deletion_marks.remove(&oldest); + } + changed +} + +/// Newest deletion mark among `entities`, or `None` when no deletion of any +/// of them was recorded here. The receive-side staleness gate feeds this in +/// as the local timestamp when the targeted record is absent. +pub(crate) fn iam_deletion_mark(state: &SiteReplicationState, entities: &[String]) -> Option { + entities + .iter() + .filter_map(|entity| state.iam_deletion_marks.get(entity).copied()) + .max() +} + +/// RFC 3339 map values, matching the other timestamps in the state object +/// (`time::serde::rfc3339` only applies to a single field). +mod rfc3339_map { + use serde::{Deserialize, Deserializer, Serialize, Serializer}; + use std::collections::BTreeMap; + use time::OffsetDateTime; + + #[derive(Serialize, Deserialize)] + #[serde(transparent)] + struct Stamp(#[serde(with = "time::serde::rfc3339")] OffsetDateTime); + + pub(super) fn serialize(map: &BTreeMap, serializer: S) -> Result { + serializer.collect_map(map.iter().map(|(entity, deleted_at)| (entity, Stamp(*deleted_at)))) + } + + pub(super) fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result, D::Error> { + let map = BTreeMap::::deserialize(deserializer)?; + Ok(map + .into_iter() + .map(|(entity, Stamp(deleted_at))| (entity, deleted_at)) + .collect()) + } } #[derive(Debug, Clone, Serialize, Deserialize, Default)] diff --git a/rustfs/src/site_replication/tests.rs b/rustfs/src/site_replication/tests.rs index 9960bbdec..2ac08ed2e 100644 --- a/rustfs/src/site_replication/tests.rs +++ b/rustfs/src/site_replication/tests.rs @@ -554,6 +554,121 @@ fn test_iam_item_deletion_entity_shapes() { assert!(iam_item_deletion_entity(&policy_set).is_none()); } +/// Deletion marks (backlog#2291) key on the same entities as the replay +/// records, except that a group member removal is marked per member (so a +/// stale re-add of one member can be judged) and a group delete marks the +/// group itself. Creates and updates leave no mark. +#[test] +fn test_iam_item_deletion_mark_entities_shapes() { + assert_eq!( + iam_item_deletion_mark_entities(&user_delete_item("alice")), + vec!["iam-user:alice".to_string()] + ); + assert_eq!( + iam_item_deletion_mark_entities(&policy_delete_item("readonly")), + vec!["policy:readonly".to_string()] + ); + + let mut group_remove = SRIAMItem { + r#type: "group-info".to_string(), + group_info: Some(SRGroupInfo { + update_req: GroupAddRemove { + group: "devs".to_string(), + members: vec!["bob".to_string(), "alice".to_string()], + status: GroupStatus::Enabled, + is_remove: true, + }, + api_version: None, + }), + ..Default::default() + }; + assert_eq!( + iam_item_deletion_mark_entities(&group_remove), + vec!["group-member:devs:bob".to_string(), "group-member:devs:alice".to_string()] + ); + group_remove + .group_info + .as_mut() + .expect("group info") + .update_req + .members + .clear(); + assert_eq!( + iam_item_deletion_mark_entities(&group_remove), + vec!["group:devs".to_string()], + "a removal without members deletes the group" + ); + group_remove.group_info.as_mut().expect("group info").update_req.is_remove = false; + assert!(iam_item_deletion_mark_entities(&group_remove).is_empty()); + + let mapping_clear = SRIAMItem { + r#type: "policy-mapping".to_string(), + policy_mapping: Some(SRPolicyMapping { + user_or_group: "alice".to_string(), + user_type: 0, + is_group: false, + policy: String::new(), + ..Default::default() + }), + ..Default::default() + }; + assert_eq!( + iam_item_deletion_mark_entities(&mapping_clear), + vec!["policy-mapping:alice:0:false".to_string()] + ); + + let mut user_create = user_delete_item("alice"); + user_create.iam_user.as_mut().expect("iam user").is_delete_req = false; + assert!(iam_item_deletion_mark_entities(&user_create).is_empty()); +} + +/// Newest wins per entity, the map stays bounded by evicting the oldest +/// mark, and the timestamps survive the state object as RFC 3339. +#[test] +fn test_record_iam_deletion_marks_newest_wins_and_stays_bounded() { + let at = |seconds: i64| OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(seconds); + let mut state = SiteReplicationState::default(); + let alice = vec!["iam-user:alice".to_string()]; + + assert!(record_iam_deletion_marks(&mut state, &alice, at(20))); + assert!( + !record_iam_deletion_marks(&mut state, &alice, at(10)), + "an older deletion does not move the mark" + ); + assert!( + !record_iam_deletion_marks(&mut state, &alice, at(20)), + "a replayed deletion is not a change" + ); + assert_eq!(iam_deletion_mark(&state, &alice), Some(at(20))); + assert!(record_iam_deletion_marks(&mut state, &alice, at(30))); + assert_eq!(iam_deletion_mark(&state, &alice), Some(at(30))); + assert_eq!(iam_deletion_mark(&state, &["iam-user:bob".to_string()]), None); + assert!(!record_iam_deletion_marks(&mut state, &[], at(40))); + + // Fill past the bound with marks older than alice's; the oldest go first. + let members: Vec = (0..SITE_REPLICATION_IAM_DELETION_MARK_LIMIT) + .map(|index| format!("group-member:devs:user-{index:04}")) + .collect(); + for (index, member) in members.iter().enumerate() { + record_iam_deletion_marks(&mut state, std::slice::from_ref(member), at(index as i64 - 2000)); + } + assert_eq!(state.iam_deletion_marks.len(), SITE_REPLICATION_IAM_DELETION_MARK_LIMIT); + assert_eq!(iam_deletion_mark(&state, &alice), Some(at(30)), "the newest mark survives eviction"); + assert_eq!(iam_deletion_mark(&state, &members[..1]), None, "the oldest mark is evicted first"); + assert_eq!(iam_deletion_mark(&state, &members[1..2]), Some(at(-1999))); + + let json = serde_json::to_value(&state).expect("serialize state"); + assert_eq!(json["iam_deletion_marks"]["iam-user:alice"], serde_json::json!("1970-01-01T00:00:30Z")); + let reloaded = parse_site_replication_state(&serde_json::to_vec(&state).expect("serialize state")).expect("parse state"); + assert_eq!(reloaded.iam_deletion_marks, state.iam_deletion_marks); + assert!( + parse_site_replication_state(br#"{"name":"a","service_account_access_key":"","service_account_parent":"","peers":{},"updated_at":null,"resync_status":{}}"#) + .expect("state without marks") + .iam_deletion_marks + .is_empty() + ); +} + /// A failed deletion delivery persists a replay record next to the collapsed /// retry entry; a fresh entry is stamped `deletions_recorded` so a later /// replay can settle it, and a repeated deletion of the same entity keeps the