From 6e305911454736a47efd61aba46bbadb41787fdb Mon Sep 17 00:00:00 2001 From: overtrue Date: Sun, 6 Sep 2026 11:27:58 +0800 Subject: [PATCH] fix(iam): publish group changes after concurrent cache reloads --- crates/iam/src/manager.rs | 9 +++-- crates/iam/src/sys.rs | 73 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 77 insertions(+), 5 deletions(-) diff --git a/crates/iam/src/manager.rs b/crates/iam/src/manager.rs index f7ad9e6a7..a859fe59a 100644 --- a/crates/iam/src/manager.rs +++ b/crates/iam/src/manager.rs @@ -1812,12 +1812,12 @@ where gi } }; - let now = OffsetDateTime::now_utc(); drop(cache); self.api.save_group_info(group, gi.clone()).await?; self.cache.with_write_lock(|cache| { + let now = OffsetDateTime::now_utc(); cache.add_or_update_group(group, &gi, now); let user_group_memberships = Arc::clone(&cache.state().user_group_memberships); @@ -1960,15 +1960,14 @@ where let d: HashSet<&String> = HashSet::from_iter(members.iter()); gi.members = s.difference(&d).map(|v| v.to_string()).collect::>(); gi.update_at = Some(updated_at); - // Cache publication time is the local clock, not the record stamp - // (see `add_users_to_group_at`). - let now = OffsetDateTime::now_utc(); - if !update_cache_only { self.api.save_group_info(name, gi.clone()).await?; } self.cache.with_write_lock(|cache| { + // Sample after storage completes so a concurrent reload cannot + // make this publication older than the cache it must update. + let now = OffsetDateTime::now_utc(); cache.add_or_update_group(name, &gi, now); let user_group_memberships = Arc::clone(&cache.state().user_group_memberships); diff --git a/crates/iam/src/sys.rs b/crates/iam/src/sys.rs index 16bc8de0a..5b0a58c92 100644 --- a/crates/iam/src/sys.rs +++ b/crates/iam/src/sys.rs @@ -2210,6 +2210,9 @@ mod tests { block_delete: Arc, delete_started: Arc, release_delete: Arc, + block_group_save: Arc, + group_save_started: Arc, + group_save_release: Arc, } impl StsTestMockStore { @@ -2223,6 +2226,9 @@ mod tests { block_delete: Arc::new(std::sync::atomic::AtomicBool::new(false)), delete_started: Arc::new(tokio::sync::Notify::new()), release_delete: Arc::new(tokio::sync::Notify::new()), + block_group_save: Arc::new(std::sync::atomic::AtomicBool::new(false)), + group_save_started: Arc::new(tokio::sync::Notify::new()), + group_save_release: Arc::new(tokio::sync::Notify::new()), } } @@ -2326,6 +2332,10 @@ mod tests { } async fn save_group_info(&self, _name: &str, _item: GroupInfo) -> Result<()> { + if self.block_group_save.load(std::sync::atomic::Ordering::SeqCst) { + self.group_save_started.notify_one(); + self.group_save_release.notified().await; + } Ok(()) } @@ -2507,6 +2517,69 @@ mod tests { IamSys::new(cache) } + async fn assert_group_write_during_reload_is_published(remove: bool) { + let iam_sys = Arc::new(temp_env::async_with_vars([("RUSTFS_SKIP_BACKGROUND_TASK", Some("1"))], test_iam_sys()).await); + let member = "sts-fallback-test-parent"; + let group = if remove { "testgroup" } else { "new-published-group" }; + let source_time = OffsetDateTime::now_utc() - time::Duration::hours(1); + iam_sys + .store + .api + .block_group_save + .store(true, std::sync::atomic::Ordering::SeqCst); + let before = iam_sys.store.cache.snapshot(); + let writer_iam = iam_sys.clone(); + let writer = tokio::spawn(async move { + if remove { + writer_iam + .remove_users_from_group_at(group, vec![member.to_string()], source_time) + .await + } else { + writer_iam + .add_users_to_group_at(group, vec![member.to_string()], source_time) + .await + } + }); + tokio::time::timeout(std::time::Duration::from_secs(5), iam_sys.store.api.group_save_started.notified()) + .await + .expect("group save should reach the barrier"); + // The pending store write has not changed the cache, so the production + // full-reload snapshot guard permits this replacement. + assert!(iam_sys.store.cache.with_write_lock(|cache| cache.matches_snapshot(&before))); + iam_sys + .store + .api + .load_all(&iam_sys.store.cache) + .await + .expect("reload while group save is pending"); + iam_sys.store.api.group_save_release.notify_one(); + assert_eq!(writer.await.expect("join group writer").expect("group write should succeed"), source_time); + let info = iam_sys + .get_group_info(group) + .await + .expect("successful group write must remain readable after reload"); + assert_eq!(info.update_at, Some(source_time), "source timestamp must remain on the record"); + assert_eq!(info.members, if remove { Vec::new() } else { vec![member.to_string()] }); + let groups = iam_sys.store.cache.snapshot().user_group_memberships.get(member).cloned(); + assert_eq!( + groups.is_some_and(|groups| groups.contains(group)), + !remove, + "membership index must reflect the write" + ); + } + + #[tokio::test] + #[serial] + async fn add_group_write_during_reload_publishes_after_store_save() { + assert_group_write_during_reload_is_published(false).await; + } + + #[tokio::test] + #[serial] + async fn remove_group_write_during_reload_publishes_after_store_save() { + assert_group_write_during_reload_is_published(true).await; + } + /// Review finding on rustfs#7195: a replicated group edit carries a source /// stamp that may predate this node's cache load time. The stamp belongs on /// the record only; publishing the cache with it makes `LockedCache::exec`