diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs
index 0aafe9ded..0d3caa90c 100644
--- a/crates/ecstore/src/api/mod.rs
+++ b/crates/ecstore/src/api/mod.rs
@@ -127,8 +127,8 @@ pub mod bucket {
get_global_bucket_metadata_sys, get_lifecycle_config, get_logging_config, get_notification_config,
get_object_lock_config, get_public_access_block_config, get_quota_config, get_replication_config,
get_request_payment_config, get_sse_config, get_tagging_config, get_versioning_config, get_website_config,
- init_bucket_metadata_sys, list_bucket_targets, remove_bucket_metadata, set_bucket_metadata, update,
- update_bucket_targets_under_transaction_lock, update_config_with,
+ init_bucket_metadata_sys, list_bucket_targets, reload_bucket_metadata, remove_bucket_metadata, set_bucket_metadata,
+ update, update_bucket_targets_under_transaction_lock, update_config_with,
};
}
diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs
index 2a30a17be..70a84851f 100644
--- a/crates/ecstore/src/bucket/metadata_sys.rs
+++ b/crates/ecstore/src/bucket/metadata_sys.rs
@@ -91,6 +91,20 @@ pub async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) -> Result<(
Ok(())
}
+/// Peer LoadBucketMetadata entry point; see
+/// [`BucketMetadataSys::reload_from_store`] for the caching contract.
+///
+/// The outer write guard spans the disk load, mirroring [`update`]: every
+/// other cache installer holds this lock (read or write), so the snapshot
+/// read here can never land after — and roll back — a newer concurrent
+/// install, and the install-plus-registry-sync sequence stays atomic
+/// against concurrent removes and reloads.
+pub async fn reload_bucket_metadata(bucket: &str) -> Result<()> {
+ let sys = get_bucket_metadata_sys()?;
+ let lock = sys.write().await;
+ lock.reload_from_store(bucket).await
+}
+
/// Drop a bucket's cached metadata from the in-memory map.
///
/// This is the counterpart to [`set_bucket_metadata`] and is invoked when a
@@ -659,6 +673,43 @@ impl BucketMetadataSys {
}
}
+ /// Reload `bucket`'s metadata from this system's own store and cache it,
+ /// refusing to treat a load miss as authoritative (the peer
+ /// LoadBucketMetadata notification path, [`reload_bucket_metadata`]).
+ ///
+ /// Only metadata actually read from persisted storage reaches the cache.
+ /// On a miss the fabricated default is discarded and an error is
+ /// returned: installing it would let a transient ConfigNotFound during
+ /// the notification overwrite a lock-enabled bucket's cached metadata
+ /// with an authoritative "no Object Lock" default, disabling the
+ /// batch-delete retention gate (`object_lock_delete_check_required`) on
+ /// this node until the next refresh. A miss is also not treated as
+ /// deletion: bucket deletion propagates through the dedicated
+ /// DeleteBucketMetadata notification ([`remove_bucket_metadata`]), which
+ /// is best-effort — a reload racing it can still re-install a just
+ /// deleted bucket's entry (pre-existing, bounded by the next delete or
+ /// restart) — but a reload miss removing entries would turn every
+ /// transient quorum dip into dropped metadata and spurious
+ /// target/durability teardown.
+ ///
+ /// The peer-visible error text is deliberately fixed: the notifying peer
+ /// matches error strings against network-failure needles
+ /// (`is_network_like_error`), so interpolating a caller-controlled
+ /// bucket name here could mark a healthy peer offline.
+ ///
+ /// Lock order: the caller holds the outer metadata-sys guard, and the
+ /// load acquires the namespace lock on the bucket's metadata config
+ /// object — the same `outer guard → meta-config namespace lock` order
+ /// `update`'s load takes; no path acquires these in reverse.
+ pub(crate) async fn reload_from_store(&self, bucket: &str) -> Result<()> {
+ let (bm, persisted) = load_bucket_metadata_parse_with_presence(self.api.clone(), bucket, true).await?;
+ if !persisted {
+ return Err(Error::other("no persisted bucket metadata readable; peer cache left unchanged"));
+ }
+ self.set(bucket.to_string(), Arc::new(bm)).await;
+ Ok(())
+ }
+
/// Remove a bucket's cached metadata from the in-memory map.
///
/// Returns `true` if an entry was present. Reserved meta buckets are ignored.
@@ -1431,6 +1482,71 @@ mod tests {
assert_eq!(tags.tag_set.len(), WRITERS, "every concurrent rewrite must survive: {tags:?}");
}
+ /// Pins the peer reload-notification contract (`reload_from_store`, the
+ /// LoadBucketMetadata RPC path): only metadata actually read from
+ /// persisted storage enters the cache. A load miss errors out and leaves
+ /// the cache untouched — it must neither install a fabricated default
+ /// for an unknown bucket nor replace an existing entry, since a
+ /// transient ConfigNotFound during the notification would otherwise
+ /// downgrade a lock-enabled bucket to an authoritative "no Object Lock"
+ /// default and disable the batch-delete retention gate on this peer.
+ #[tokio::test]
+ async fn peer_reload_never_caches_fabricated_defaults_as_authoritative() {
+ let (_dirs, ecstore) = isolated_store_over_temp_disks().await;
+ let sys = BucketMetadataSys::new(ecstore.clone());
+
+ // (a) Miss with no cached entry: the reload fails and installs nothing.
+ let err = sys
+ .reload_from_store("reload-bucket")
+ .await
+ .expect_err("a reload miss must be reported to the notifying peer");
+ assert!(
+ err.to_string().contains("no persisted bucket metadata readable"),
+ "the miss must surface through the dedicated non-persisted branch, got: {err}"
+ );
+ assert!(
+ sys.get("reload-bucket").await.is_err(),
+ "a reload miss must not install a fabricated default"
+ );
+
+ // (b) Miss with an existing entry: the reload fails and the entry
+ // (standing in for a lock-enabled bucket's metadata) survives intact.
+ let mut kept = BucketMetadata::new("reload-bucket");
+ kept.object_lock_config_xml = b"".to_vec();
+ sys.set("reload-bucket".to_string(), Arc::new(kept)).await;
+ assert!(sys.reload_from_store("reload-bucket").await.is_err());
+ let cached = sys
+ .get("reload-bucket")
+ .await
+ .expect("existing entry must survive a reload miss");
+ assert_eq!(
+ cached.object_lock_config_xml,
+ b"".to_vec(),
+ "a reload miss must not replace the cached entry with a fabricated default"
+ );
+
+ // (c) Persisted metadata reloads over a stale cached entry: the
+ // reload converges the cache to disk truth.
+ let mut persisted = BucketMetadata::new("reload-bucket");
+ persisted.policy_config_json = b"persisted-marker".to_vec();
+ sys.persist_and_set(persisted).await.expect("metadata should persist");
+ let mut stale = BucketMetadata::new("reload-bucket");
+ stale.policy_config_json = b"stale-cache-marker".to_vec();
+ sys.set("reload-bucket".to_string(), Arc::new(stale)).await;
+ sys.reload_from_store("reload-bucket")
+ .await
+ .expect("persisted metadata should reload");
+ let cached = sys
+ .get("reload-bucket")
+ .await
+ .expect("reloaded persisted metadata must be cached");
+ assert_eq!(
+ cached.policy_config_json,
+ b"persisted-marker".to_vec(),
+ "a reload must converge the cache to the persisted disk state"
+ );
+ }
+
fn target(bucket: &str, id: &str) -> BucketTarget {
BucketTarget {
source_bucket: bucket.to_string(),
diff --git a/rustfs/src/storage/mod.rs b/rustfs/src/storage/mod.rs
index 8b68f3aeb..9033400c4 100644
--- a/rustfs/src/storage/mod.rs
+++ b/rustfs/src/storage/mod.rs
@@ -69,9 +69,9 @@ pub(crate) use storage_api::{
get_lock_acquire_timeout, get_public_access_block_config, head_prefix_consumer, helper_consumer, init_background_replication,
init_bucket_metadata_sys, init_ecstore_config, init_local_disks_with_instance_ctx, init_lock_clients,
is_all_buckets_not_found, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found, is_valid_storage_class,
- load_bucket_metadata, options_consumer, prewarm_local_disk_id_map_with_instance_ctx, read_config, record_replication_proxy,
- rpc_consumer, runtime_sources_consumer, s3_api_consumer, serialize, set_bucket_metadata, table_catalog_path_hash,
- to_s3s_etag, topology_snapshot_from_endpoint_pools_with_capabilities, try_migrate_bucket_metadata, try_migrate_iam_config,
+ options_consumer, prewarm_local_disk_id_map_with_instance_ctx, read_config, record_replication_proxy, rpc_consumer,
+ runtime_sources_consumer, s3_api_consumer, serialize, table_catalog_path_hash, to_s3s_etag,
+ topology_snapshot_from_endpoint_pools_with_capabilities, try_migrate_bucket_metadata, try_migrate_iam_config,
try_migrate_server_config, update_bucket_metadata_config, verify_rpc_signature, wrap_reader,
};
diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs
index bf2f81f71..4b8b94a8d 100644
--- a/rustfs/src/storage/rpc/node_service.rs
+++ b/rustfs/src/storage/rpc/node_service.rs
@@ -4439,6 +4439,31 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn test_load_bucket_metadata_failure_skips_scanner_maintenance() {
+ let service = create_test_node_service();
+ let maintenance_generation = rustfs_scanner::scanner_maintenance_generation();
+
+ let request = Request::new(LoadBucketMetadataRequest {
+ bucket: "reload-miss-scanner-guard-bucket".to_string(),
+ scanner_maintenance_change: true,
+ });
+
+ let response = service.load_bucket_metadata(request).await.expect("rpc should reply");
+ let load_response = response.into_inner();
+
+ // Whether the reload fails on missing server state or on the absent
+ // persisted metadata, a failed reload must report failure and must
+ // not tell the scanner a maintenance change landed.
+ assert!(!load_response.success);
+ assert!(load_response.error_info.is_some());
+ assert_eq!(
+ rustfs_scanner::scanner_maintenance_generation(),
+ maintenance_generation,
+ "a failed metadata reload must not advance scanner maintenance activity"
+ );
+ }
+
#[tokio::test]
#[ignore = "requires isolated global object layer state"]
async fn test_load_bucket_metadata_no_object_layer() {
diff --git a/rustfs/src/storage/rpc/node_service/bucket.rs b/rustfs/src/storage/rpc/node_service/bucket.rs
index 2d383f65b..2f9c490c7 100644
--- a/rustfs/src/storage/rpc/node_service/bucket.rs
+++ b/rustfs/src/storage/rpc/node_service/bucket.rs
@@ -17,7 +17,7 @@ use crate::storage::storage_api::rpc_consumer::node_service::contract::bucket::{
BucketOptions, DeleteBucketOptions, MakeBucketOptions,
};
use crate::storage::storage_api::rpc_consumer::node_service::{
- DiskError, StoragePeerS3ClientExt as _, load_bucket_metadata, remove_bucket_metadata, set_bucket_metadata,
+ DiskError, StoragePeerS3ClientExt as _, reload_bucket_metadata, remove_bucket_metadata,
};
use rustfs_common::heal_channel::HealOpts;
use rustfs_protos::proto_gen::node_service::*;
@@ -66,21 +66,15 @@ impl NodeService {
}));
}
- let Some(store) = self.resolve_object_store() else {
+ let Some(_store) = self.resolve_object_store() else {
return Ok(Response::new(LoadBucketMetadataResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}));
};
- match load_bucket_metadata(store, &bucket).await {
- Ok(meta) => {
- if let Err(err) = set_bucket_metadata(bucket.clone(), meta).await {
- return Ok(Response::new(LoadBucketMetadataResponse {
- success: false,
- error_info: Some(err.to_string()),
- }));
- };
+ match reload_bucket_metadata(&bucket).await {
+ Ok(()) => {
if scanner_maintenance_change {
rustfs_scanner::record_scanner_maintenance_change(&bucket);
}
diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs
index 859c9e384..1b9b27056 100644
--- a/rustfs/src/storage/storage_api.rs
+++ b/rustfs/src/storage/storage_api.rs
@@ -236,8 +236,8 @@ pub(crate) mod rpc_consumer {
ECStore, Error, FileInfoVersions, LocalPeerS3Client, MetricType, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS,
ReadMultipleReq, ReadMultipleResp, ReadOptions, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC,
StorageDiskRpcExt, StoragePeerS3ClientExt, UpdateMetadataOpts, all_local_disk_path, collect_local_metrics,
- find_local_disk_by_ref, get_local_server_property, load_bucket_metadata, reload_transition_tier_config,
- remove_bucket_metadata, set_bucket_metadata, validate_batch_read_version_item_count,
+ find_local_disk_by_ref, get_local_server_property, reload_bucket_metadata, reload_transition_tier_config,
+ remove_bucket_metadata, validate_batch_read_version_item_count,
};
pub(crate) type StorageResult = super::super::Result;
@@ -1365,10 +1365,6 @@ impl StoragePeerS3ClientExt for LocalPeerS3Client {
}
}
-pub(crate) async fn load_bucket_metadata(api: Arc, bucket: &str) -> Result {
- ecstore_bucket::metadata::load_bucket_metadata(api, bucket).await
-}
-
#[cfg(test)]
pub(crate) fn bucket_metadata_sys_initialized() -> bool {
ecstore_bucket::metadata_sys::get_global_bucket_metadata_sys().is_some()
@@ -1441,10 +1437,15 @@ pub(crate) async fn get_bucket_website_config(bucket: &str) -> Result<(s3s::dto:
ecstore_bucket::metadata_sys::get_website_config(bucket).await
}
+#[cfg(test)]
pub(crate) async fn set_bucket_metadata(bucket: String, bm: BucketMetadata) -> Result<()> {
ecstore_bucket::metadata_sys::set_bucket_metadata(bucket, bm).await
}
+pub(crate) async fn reload_bucket_metadata(bucket: &str) -> Result<()> {
+ ecstore_bucket::metadata_sys::reload_bucket_metadata(bucket).await
+}
+
pub(crate) async fn remove_bucket_metadata(bucket: &str) -> Result {
ecstore_bucket::metadata_sys::remove_bucket_metadata(bucket).await
}