diff --git a/crates/e2e_test/src/cluster_concurrency_test.rs b/crates/e2e_test/src/cluster_concurrency_test.rs index cc2f2cddc..62c9e694a 100644 --- a/crates/e2e_test/src/cluster_concurrency_test.rs +++ b/crates/e2e_test/src/cluster_concurrency_test.rs @@ -15,12 +15,14 @@ use crate::common::RustFSTestClusterEnvironment; use aws_sdk_s3::Client; use aws_sdk_s3::error::SdkError; +use aws_sdk_s3::types::{CorsConfiguration, CorsRule}; use bytes::Bytes; use std::sync::Arc; use tokio::sync::Barrier; use tracing::{info, warn}; const BUCKET: &str = "conditional-put-race-bucket"; +const BUCKET_METADATA_RELOAD_BUCKET: &str = "bucket-metadata-reload-barrier"; async fn cleanup_object(client: &Client, key: &str) { if let Err(e) = client.delete_object().bucket(BUCKET).key(key).send().await { @@ -28,6 +30,16 @@ async fn cleanup_object(client: &Client, key: &str) { } } +async fn assert_bucket_cors_missing(client: &Client) { + let result = client.get_bucket_cors().bucket(BUCKET_METADATA_RELOAD_BUCKET).send().await; + match result { + Err(SdkError::ServiceError(error)) => { + assert_eq!(error.err().meta().code(), Some("NoSuchCORSConfiguration")); + } + result => panic!("expected the peer to report a missing CORS configuration: {result:?}"), + } +} + async fn conditional_put( client: &Client, key: &str, @@ -236,3 +248,48 @@ async fn test_conditional_put_basic_cluster() -> Result<(), Box Result<(), Box> { + crate::common::init_logging(); + + let mut cluster = RustFSTestClusterEnvironment::new(2).await?; + cluster.start().await?; + cluster.create_test_bucket(BUCKET_METADATA_RELOAD_BUCKET).await?; + + let writer = cluster.create_s3_client(0)?; + let reader = cluster.create_s3_client(1)?; + assert_bucket_cors_missing(&reader).await; + + let rule = CorsRule::builder() + .allowed_methods("GET") + .allowed_origins("https://example.com") + .build()?; + let configuration = CorsConfiguration::builder().cors_rules(rule).build()?; + + writer + .put_bucket_cors() + .bucket(BUCKET_METADATA_RELOAD_BUCKET) + .cors_configuration(configuration) + .send() + .await?; + + let response = reader.get_bucket_cors().bucket(BUCKET_METADATA_RELOAD_BUCKET).send().await?; + let rules = response.cors_rules(); + assert_eq!( + rules.len(), + 1, + "peer should observe the committed CORS rule before the write response returns" + ); + assert_eq!(rules[0].allowed_methods(), ["GET"]); + assert_eq!(rules[0].allowed_origins(), ["https://example.com"]); + + writer + .delete_bucket_cors() + .bucket(BUCKET_METADATA_RELOAD_BUCKET) + .send() + .await?; + assert_bucket_cors_missing(&reader).await; + writer.delete_bucket().bucket(BUCKET_METADATA_RELOAD_BUCKET).send().await?; + Ok(()) +} diff --git a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs index 3bac43266..edfec184f 100644 --- a/crates/ecstore/src/cluster/rpc/peer_rest_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_rest_client.rs @@ -85,6 +85,7 @@ const PEER_REST_RECOVERY_MAX_ATTEMPTS: u32 = 60; const PEER_REST_RECOVERY_MAX_BACKOFF: Duration = Duration::from_secs(30); const SCANNER_ACTIVITY_MAX_MESSAGE_SIZE: usize = 1024; const REPLICATION_STATS_MAX_MESSAGE_SIZE: usize = 8 * 1024 * 1024; +const BUCKET_METADATA_RELOAD_TIMEOUT: Duration = Duration::from_secs(5); /// Error for a peer that reported `success = false` without an `error_info` payload. /// @@ -1328,27 +1329,38 @@ impl PeerRestClient { } pub async fn load_bucket_metadata(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> { - self.finalize_result( - async { - let mut client = self.get_client().await?; - let mut request = Request::new(LoadBucketMetadataRequest { - bucket: bucket.to_string(), - scanner_maintenance_change, - }); - set_tonic_mutation_body_digest(&mut request)?; - - let response = client.load_bucket_metadata(request).await?.into_inner(); - if !response.success { - if let Some(msg) = response.error_info { - return Err(Error::other(msg)); - } - return Err(peer_failure_without_details("load_bucket_metadata", Some(bucket))); - } - Ok(()) + let result = tokio::time::timeout(BUCKET_METADATA_RELOAD_TIMEOUT, async { + let result = self.load_bucket_metadata_once(bucket, scanner_maintenance_change).await; + if let Err(err) = &result + && Self::is_network_like_error(err) + { + self.prepare_retry().await; + return self.load_bucket_metadata_once(bucket, scanner_maintenance_change).await; } - .await, - ) + result + }) .await + .unwrap_or_else(|_| Err(Error::other(format!("load_bucket_metadata({bucket}) timed out")))); + self.finalize_result(result).await + } + + async fn load_bucket_metadata_once(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> { + let mut client = self.get_client().await?; + let mut request = Request::new(LoadBucketMetadataRequest { + bucket: bucket.to_string(), + scanner_maintenance_change, + }); + set_tonic_mutation_body_digest(&mut request)?; + request.set_timeout(BUCKET_METADATA_RELOAD_TIMEOUT); + + let response = client.load_bucket_metadata(request).await?.into_inner(); + if !response.success { + if let Some(msg) = response.error_info { + return Err(Error::other(msg)); + } + return Err(peer_failure_without_details("load_bucket_metadata", Some(bucket))); + } + Ok(()) } pub async fn delete_bucket_metadata(&self, bucket: &str) -> Result<()> { diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index 7d55035e9..25e37179e 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -513,13 +513,15 @@ fn sr_bucket_meta_item(bucket: String, item_type: &str) -> SRBucketMeta { } } -fn notify_bucket_metadata_reload( +async fn notify_bucket_metadata_reload( bucket: String, operation: &'static str, request_context: Option, scanner_maintenance_change: bool, ) { record_local_scanner_maintenance_reload(&bucket, scanner_maintenance_change); + // Keep reload detached across request cancellation, but wait before a healthy peer can serve the previous config. + let (completed_tx, completed_rx) = tokio::sync::oneshot::channel(); spawn_background_with_context(request_context, async move { if let Some(notification_sys) = current_notification_system() { let result = if scanner_maintenance_change { @@ -531,7 +533,9 @@ fn notify_bucket_metadata_reload( warn!(bucket = %bucket, error = %err, "failed to notify peers after {operation}"); } } + let _ = completed_tx.send(()); }); + let _ = completed_rx.await; } fn record_local_scanner_maintenance_reload(bucket: &str, scanner_maintenance_change: bool) { @@ -1476,7 +1480,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "delete bucket encryption", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "delete bucket encryption", request_context, false).await; let item = sr_bucket_meta_item(bucket.clone(), "sse-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { @@ -1508,7 +1512,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "delete bucket cors", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "delete bucket cors", request_context, false).await; let item = sr_bucket_meta_item(bucket.clone(), "cors-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { @@ -1540,7 +1544,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "delete bucket lifecycle", request_context, true); + notify_bucket_metadata_reload(bucket.clone(), "delete bucket lifecycle", request_context, true).await; let item = sr_bucket_meta_item(bucket.clone(), "lc-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { @@ -1572,7 +1576,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "delete bucket policy", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "delete bucket policy", request_context, false).await; let item = sr_bucket_meta_item(bucket.clone(), "policy"); if let Err(err) = site_replication_bucket_meta_hook(item).await { @@ -1630,7 +1634,7 @@ impl DefaultBucketUsecase { } drop(targets_guard); - notify_bucket_metadata_reload(bucket.clone(), "delete bucket replication", request_context, true); + notify_bucket_metadata_reload(bucket.clone(), "delete bucket replication", request_context, true).await; let item = sr_bucket_meta_item(bucket.clone(), "replication-config"); if let Err(err) = site_replication_bucket_meta_hook(item).await { @@ -1655,7 +1659,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "delete bucket tagging", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "delete bucket tagging", request_context, false).await; let item = sr_bucket_meta_item(bucket.clone(), "tags"); if let Err(err) = site_replication_bucket_meta_hook(item).await { @@ -1688,7 +1692,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "delete public access block", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "delete public access block", request_context, false).await; Ok(S3Response::with_status(DeletePublicAccessBlockOutput::default(), StatusCode::NO_CONTENT)) } @@ -2143,7 +2147,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put bucket encryption", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "put bucket encryption", request_context, false).await; let mut item = sr_bucket_meta_item(bucket.clone(), "sse-config"); item.sse_config = Some( @@ -2222,7 +2226,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put bucket lifecycle", request_context, true); + notify_bucket_metadata_reload(bucket.clone(), "put bucket lifecycle", request_context, true).await; let mut item = sr_bucket_meta_item(bucket.clone(), "lc-config"); item.expiry_lc_config = @@ -2307,7 +2311,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put bucket notification", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "put bucket notification", request_context, false).await; let region = resolve_notification_region(self.global_region(), request_region); let notify = current_notify_interface_for_context(self.context.as_deref()); @@ -2412,7 +2416,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put bucket policy", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "put bucket policy", request_context, false).await; 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))?); @@ -2447,7 +2451,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put bucket cors", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "put bucket cors", request_context, false).await; let mut item = sr_bucket_meta_item(bucket.clone(), "cors-config"); item.cors = @@ -2491,7 +2495,7 @@ impl DefaultBucketUsecase { .map_err(ApiError::from)?; drop(targets_guard); - notify_bucket_metadata_reload(bucket.clone(), "put bucket replication", request_context, true); + notify_bucket_metadata_reload(bucket.clone(), "put bucket replication", request_context, true).await; let mut item = sr_bucket_meta_item(bucket.clone(), "replication-config"); item.replication_config = Some( @@ -2531,7 +2535,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put public access block", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "put public access block", request_context, false).await; Ok(S3Response::new(PutPublicAccessBlockOutput::default())) } @@ -2560,7 +2564,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put bucket tagging", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "put bucket tagging", request_context, false).await; let mut item = sr_bucket_meta_item(bucket.clone(), "tags"); item.tags = Some(serialize_config(&tagging).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?); @@ -2593,7 +2597,7 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; - notify_bucket_metadata_reload(bucket.clone(), "put bucket versioning", request_context, false); + notify_bucket_metadata_reload(bucket.clone(), "put bucket versioning", request_context, false).await; let mut item = sr_bucket_meta_item(bucket.clone(), "version-config"); item.versioning = Some( @@ -3044,7 +3048,7 @@ mod tests { "{method} should identify the bucket metadata operation in reload logs" ); let expected_reload = format!( - "notify_bucket_metadata_reload(bucket.clone(), \"{operation}\", request_context, {scanner_maintenance_change});" + "notify_bucket_metadata_reload(bucket.clone(), \"{operation}\", request_context, {scanner_maintenance_change}).await;" ); assert!( body.contains(&expected_reload),