diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index ab5294d42..33a35adaf 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -1022,6 +1022,7 @@ impl DefaultBucketUsecase { &self, req: S3Request, ) -> S3Result> { + let request_context = req.extensions.get::().cloned(); let DeleteBucketPolicyInput { bucket, .. } = req.input; let Some(store) = new_object_layer_fn() else { @@ -1037,6 +1038,8 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; + notify_bucket_metadata_reload(bucket.clone(), "delete bucket policy", request_context); + let item = sr_bucket_meta_item(bucket.clone(), "policy"); if let Err(err) = site_replication_bucket_meta_hook(item).await { warn!(bucket = %bucket, error = ?err, "site replication bucket policy delete hook failed"); @@ -1108,6 +1111,7 @@ impl DefaultBucketUsecase { &self, req: S3Request, ) -> S3Result> { + let request_context = req.extensions.get::().cloned(); let DeletePublicAccessBlockInput { bucket, .. } = req.input; let Some(store) = new_object_layer_fn() else { @@ -1123,6 +1127,8 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; + notify_bucket_metadata_reload(bucket.clone(), "delete public access block", request_context); + Ok(S3Response::with_status(DeletePublicAccessBlockOutput::default(), StatusCode::NO_CONTENT)) } @@ -1678,6 +1684,7 @@ impl DefaultBucketUsecase { &self, req: S3Request, ) -> S3Result> { + let request_context = req.extensions.get::().cloned(); let PutBucketPolicyInput { bucket, policy, .. } = req.input; let Some(store) = new_object_layer_fn() else { @@ -1724,6 +1731,8 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; + notify_bucket_metadata_reload(bucket.clone(), "put bucket policy", request_context); + let mut item = sr_bucket_meta_item(bucket.clone(), "policy"); item.policy = Some(serde_json::from_str(&policy).map_err(|e| s3_error!(InvalidArgument, "parse policy failed {:?}", e))?); if let Err(err) = site_replication_bucket_meta_hook(item).await { @@ -1807,6 +1816,7 @@ impl DefaultBucketUsecase { &self, req: S3Request, ) -> S3Result> { + let request_context = req.extensions.get::().cloned(); let PutPublicAccessBlockInput { bucket, public_access_block_configuration, @@ -1827,6 +1837,8 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; + notify_bucket_metadata_reload(bucket.clone(), "put public access block", request_context); + Ok(S3Response::new(PutPublicAccessBlockOutput::default())) } @@ -2103,6 +2115,35 @@ mod tests { req } + fn usecase_method_source<'a>(source: &'a str, method: &str) -> &'a str { + let start_marker = format!("pub async fn {method}"); + let start = source.find(&start_marker).expect("method should exist"); + let rest = &source[start + start_marker.len()..]; + let end = rest.find("\n pub async fn ").unwrap_or(rest.len()); + &rest[..end] + } + + #[test] + fn bucket_policy_and_public_access_block_changes_notify_peer_metadata_reload() { + let source = include_str!("bucket_usecase.rs"); + for (method, operation) in [ + ("execute_delete_bucket_policy", "delete bucket policy"), + ("execute_put_bucket_policy", "put bucket policy"), + ("execute_delete_public_access_block", "delete public access block"), + ("execute_put_public_access_block", "put public access block"), + ] { + let body = usecase_method_source(source, method); + assert!( + body.contains("notify_bucket_metadata_reload("), + "{method} should notify peers to reload cached bucket metadata" + ); + assert!( + body.contains(operation), + "{method} should identify the bucket metadata operation in reload logs" + ); + } + } + fn replication_rule_for_target(arn: &str) -> ReplicationRule { ReplicationRule { delete_marker_replication: None,