fix(iam): avoid stale cache replacement on walk errors (#3094)

* fix(iam): avoid stale cache replacement on walk errors

* fix(iam): guard full reload cache commits

---------

Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: loverustfs <hello@rustfs.com>
Co-authored-by: 季宏伟 <jihongwei@jihongweis-MacBook-Pro.local>
This commit is contained in:
GatewayJ
2026-05-28 17:54:06 +08:00
committed by GitHub
parent 28bac7fbd6
commit 8d20e89bf8
+56 -9
View File
@@ -362,6 +362,7 @@ impl ObjectStore {
let (tx, mut rx) = mpsc::channel::<ObjectInfoOrErr>(100);
let path = prefix.to_owned();
let sender_on_error = sender.clone();
tokio::spawn(async move {
if let Err(err) = store
.walk(ctx.clone(), Self::BUCKET_NAME, &path, tx, WalkOptions::default())
@@ -378,6 +379,12 @@ impl ObjectStore {
error = %err,
"system path walk failed"
);
let _ = sender_on_error
.send(StringOrErr {
item: None,
err: Some(err.into()),
})
.await;
}
});
@@ -1018,6 +1025,14 @@ impl Store for ObjectStore {
}
async fn load_all(&self, cache: &Cache) -> Result<()> {
let policy_docs_snapshot = cache.policy_docs.load();
let users_snapshot = cache.users.load();
let user_policies_snapshot = cache.user_policies.load();
let groups_snapshot = cache.groups.load();
let user_group_memberships_snapshot = cache.user_group_memberships.load();
let group_policies_snapshot = cache.group_policies.load();
let sts_accounts_snapshot = cache.sts_accounts.load();
let sts_policies_snapshot = cache.sts_policies.load();
let listed_config_items = self.list_all_iamconfig_items().await?;
let mut policy_docs_cache = CacheEntity::new(get_default_policyes());
@@ -1059,8 +1074,6 @@ impl Store for ObjectStore {
}
}
cache.policy_docs.store(Arc::new(policy_docs_cache.update_load_time()));
let mut user_items_cache = CacheEntity::default();
// users
@@ -1104,6 +1117,7 @@ impl Store for ObjectStore {
}
// groups
let mut groups_cache = None;
if let Some(item_name_list) = listed_config_items.get(GROUPS_LIST_KEY) {
let mut items_cache = CacheEntity::default();
@@ -1115,10 +1129,11 @@ impl Store for ObjectStore {
};
}
cache.groups.store(Arc::new(items_cache.update_load_time()));
groups_cache = Some(items_cache);
}
// user policies
let mut user_policies_cache = None;
if let Some(item_name_list) = listed_config_items.get(POLICY_DB_USERS_LIST_KEY) {
let mut item_name_list = item_name_list.clone();
@@ -1159,10 +1174,11 @@ impl Store for ObjectStore {
item_name_list = item_name_list.split_off(32);
}
cache.user_policies.store(Arc::new(items_cache.update_load_time()));
user_policies_cache = Some(items_cache);
}
// group policy
let mut group_policies_cache = None;
if let Some(item_name_list) = listed_config_items.get(POLICY_DB_GROUPS_LIST_KEY) {
let mut items_cache = CacheEntity::default();
@@ -1177,7 +1193,7 @@ impl Store for ObjectStore {
};
}
cache.group_policies.store(Arc::new(items_cache.update_load_time()));
group_policies_cache = Some(items_cache);
}
let mut sts_policies_cache = CacheEntity::default();
@@ -1216,7 +1232,6 @@ impl Store for ObjectStore {
// cache.users.store(Arc::new(items_cache.update_load_time()));
}
cache.build_user_group_memberships();
let mut sts_items_cache = CacheEntity::default();
// sts users
if let Some(item_name_list) = listed_config_items.get(STS_LIST_KEY) {
@@ -1245,9 +1260,41 @@ impl Store for ObjectStore {
}
}
cache.users.store(Arc::new(user_items_cache.update_load_time()));
cache.sts_accounts.store(Arc::new(sts_items_cache.update_load_time()));
cache.sts_policies.store(Arc::new(sts_policies_cache.update_load_time()));
let policy_docs_current = cache.policy_docs.load();
let users_current = cache.users.load();
let user_policies_current = cache.user_policies.load();
let groups_current = cache.groups.load();
let user_group_memberships_current = cache.user_group_memberships.load();
let group_policies_current = cache.group_policies.load();
let sts_accounts_current = cache.sts_accounts.load();
let sts_policies_current = cache.sts_policies.load();
if Cache::ptr_eq(&*policy_docs_snapshot, &*policy_docs_current)
&& Cache::ptr_eq(&*users_snapshot, &*users_current)
&& Cache::ptr_eq(&*user_policies_snapshot, &*user_policies_current)
&& Cache::ptr_eq(&*groups_snapshot, &*groups_current)
&& Cache::ptr_eq(&*user_group_memberships_snapshot, &*user_group_memberships_current)
&& Cache::ptr_eq(&*group_policies_snapshot, &*group_policies_current)
&& Cache::ptr_eq(&*sts_accounts_snapshot, &*sts_accounts_current)
&& Cache::ptr_eq(&*sts_policies_snapshot, &*sts_policies_current)
{
cache.policy_docs.store(Arc::new(policy_docs_cache.update_load_time()));
if let Some(groups_cache) = groups_cache {
cache.groups.store(Arc::new(groups_cache.update_load_time()));
}
if let Some(user_policies_cache) = user_policies_cache {
cache.user_policies.store(Arc::new(user_policies_cache.update_load_time()));
}
if let Some(group_policies_cache) = group_policies_cache {
cache.group_policies.store(Arc::new(group_policies_cache.update_load_time()));
}
cache.users.store(Arc::new(user_items_cache.update_load_time()));
cache.sts_accounts.store(Arc::new(sts_items_cache.update_load_time()));
cache.sts_policies.store(Arc::new(sts_policies_cache.update_load_time()));
cache.build_user_group_memberships();
} else {
warn!("skip IAM full reload cache commit because one or more IAM caches changed during reload");
}
Ok(())
}