From 676f2276b4cfc87d7ef0206a5d52340c87c2d11b Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sun, 12 Jul 2026 05:03:17 +0800 Subject: [PATCH] fix(replication): refresh targets after site endpoint edits (#4756) * fix(replication): refresh targets after site endpoint edits * fix(replication): serialize site bucket lifecycle --- .config/nextest.toml | 6 - .../src/replication_extension_test.rs | 283 +++- rustfs/src/admin/handlers/bucket_meta.rs | 7 + rustfs/src/admin/handlers/replication.rs | 3 + rustfs/src/admin/handlers/site_replication.rs | 1412 ++++++++++++++--- rustfs/src/admin/route_policy.rs | 6 + rustfs/src/admin/route_registration_test.rs | 2 + rustfs/src/admin/router.rs | 3 + rustfs/src/app/bucket_usecase.rs | 5 + rustfs/src/storage/storage_api.rs | 58 +- 10 files changed, 1548 insertions(+), 237 deletions(-) diff --git a/.config/nextest.toml b/.config/nextest.toml index ac87015fb..6a967d759 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -101,12 +101,6 @@ retries = 2 filter = 'package(rustfs-ecstore) & test(walk_dir_does_not_charge_consumer_backpressure_to_the_stall_budget)' retries = 2 -# QUARANTINE: OPEN rustfs#4740 — the embedded S3 test intermittently observes -# BucketNotEmpty immediately after an acknowledged unversioned object delete. -[[profile.ci.overrides]] -filter = 'package(rustfs) & binary(embedded_test) & test(test_embedded_server_basic_s3_operations)' -retries = 2 - # Serialize the 4-disk reliability / degraded-read e2e tests under the ci # profile too (see the e2e-reliability test-group note near the top). Not a # quarantine: no retries, just single-threaded so several 4-disk servers never diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index 6a55150c3..057760a8a 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -405,6 +405,14 @@ async fn delete_bucket_replication( signed_request(http::Method::DELETE, &url, &env.access_key, &env.secret_key, None, None).await } +async fn get_bucket_replication( + env: &RustFSTestEnvironment, + bucket: &str, +) -> Result> { + let url = format!("{}/{bucket}?replication", env.url); + signed_request(http::Method::GET, &url, &env.access_key, &env.secret_key, None, None).await +} + async fn enable_bucket_versioning(env: &RustFSTestEnvironment, bucket: &str) -> Result<(), Box> { let client = env.create_s3_client(); client @@ -1315,6 +1323,29 @@ async fn wait_for_site_replication_disabled( wait_for_site_replication_info(env, |info| !info.enabled && info.sites.is_empty()).await } +async fn assert_site_replication_bucket_detached( + env: &RustFSTestEnvironment, + bucket: &str, +) -> Result<(), Box> { + let targets_response = list_replication_targets_request(env, Some(bucket)).await?; + if targets_response.status() != StatusCode::OK { + return Err(format!("list remote targets failed after site removal: {}", targets_response.status()).into()); + } + let targets: Vec = targets_response.json().await?; + if !targets.is_empty() { + return Err(format!("site removal left remote targets for {bucket}: {targets:?}").into()); + } + + let replication_response = get_bucket_replication(env, bucket).await?; + if replication_response.status() != StatusCode::NOT_FOUND { + let status = replication_response.status(); + let body = replication_response.text().await.unwrap_or_default(); + return Err(format!("site removal left replication config for {bucket}: {status} {body}").into()); + } + + Ok(()) +} + async fn wait_for_site_replication_info( env: &RustFSTestEnvironment, predicate: F, @@ -2632,10 +2663,10 @@ async fn test_site_replication_resync_start_cancel_restart_real_dual_node() -> R let source_client = source_env.create_s3_client(); let target_client = target_env.create_s3_client(); - // Site replication rejects initialization unless all but one site is empty, so only - // the source is seeded before joining; the target bucket is created afterwards. - source_client.create_bucket().bucket(source_bucket).send().await?; - enable_bucket_versioning(&source_env, source_bucket).await?; + // Seed only the joining site so its synchronous backfill must call the initiating + // site before the add handler has persisted its final enabled state. + target_client.create_bucket().bucket(source_bucket).send().await?; + enable_bucket_versioning(&target_env, source_bucket).await?; let add_status = site_replication_add( &source_env, @@ -2665,11 +2696,9 @@ async fn test_site_replication_resync_start_cancel_restart_real_dual_node() -> R .find(|peer| peer.endpoint == target_env.url) .ok_or("target peer missing from source site replication info")?; - // Wait for the joined target site to converge: site replication propagates the - // source bucket to the target and auto-configures its replication target. Drive the - // resync against that auto-created target instead of a redundant manual one (which - // now collides — "Remote target already exists"). - wait_for_bucket_on_target(&target_client, source_bucket).await?; + // Wait for the initiating site to accept the join callback and configure the + // backfilled bucket before driving resync in the normal source-to-target direction. + wait_for_bucket_on_target(&source_client, source_bucket).await?; let target_arn = wait_for_remote_target_arn(&source_env, source_bucket).await?; for idx in 0..32 { @@ -2739,7 +2768,7 @@ async fn test_site_replication_resync_start_cancel_restart_real_dual_node() -> R #[tokio::test] #[serial] -async fn test_site_replication_edit_and_status_peer_state_real_dual_node() -> Result<(), Box> { +async fn test_site_replication_edit_and_status_peer_state_real_three_node() -> Result<(), Box> { init_logging(); let mut source_env = RustFSTestEnvironment::new().await?; @@ -2748,7 +2777,25 @@ async fn test_site_replication_edit_and_status_peer_state_real_dual_node() -> Re .await?; let mut target_env = RustFSTestEnvironment::new().await?; - target_env.start_rustfs_server_without_cleanup(vec![]).await?; + target_env + .start_rustfs_server_with_env(vec![], LOOPBACK_REPLICATION_TARGET_ENV) + .await?; + + let mut relay_env = RustFSTestEnvironment::new().await?; + relay_env + .start_rustfs_server_with_env(vec![], LOOPBACK_REPLICATION_TARGET_ENV) + .await?; + + let source_client = source_env.create_s3_client(); + let target_client = target_env.create_s3_client(); + let relay_client = relay_env.create_s3_client(); + let bucket = "site-repl-edit-endpoint"; + let baseline_key = "before-edit.txt"; + let baseline_payload = b"site replication before endpoint edit".to_vec(); + let moved_key = "after-edit.txt"; + let moved_payload = b"site replication after endpoint edit".to_vec(); + let relayed_key = "after-edit-from-relay.txt"; + let relayed_payload = b"site replication after endpoint edit from relay".to_vec(); let add_status = site_replication_add( &source_env, @@ -2765,57 +2812,132 @@ async fn test_site_replication_edit_and_status_peer_state_real_dual_node() -> Re access_key: target_env.access_key.clone(), secret_key: target_env.secret_key.clone(), }, + PeerSite { + name: "relay-site".to_string(), + endpoint: relay_env.url.clone(), + access_key: relay_env.access_key.clone(), + secret_key: relay_env.secret_key.clone(), + }, ], ) .await?; assert!(add_status.success, "unexpected site add result: {:?}", add_status); - let source_info = wait_for_site_replication_enabled(&source_env, 2).await?; - let _target_info = wait_for_site_replication_enabled(&target_env, 2).await?; + let source_info = wait_for_site_replication_enabled(&source_env, 3).await?; + let _target_info = wait_for_site_replication_enabled(&target_env, 3).await?; + let _relay_info = wait_for_site_replication_enabled(&relay_env, 3).await?; let mut remote_peer = source_info .sites .into_iter() .find(|peer| peer.endpoint == target_env.url) .ok_or("target peer missing from source site replication info")?; + source_client.create_bucket().bucket(bucket).send().await?; + enable_bucket_versioning(&source_env, bucket).await?; + wait_for_bucket_on_target(&target_client, bucket).await?; + wait_for_bucket_on_target(&relay_client, bucket).await?; + source_client + .put_object() + .bucket(bucket) + .key(baseline_key) + .body(ByteStream::from(baseline_payload.clone())) + .send() + .await?; + let replicated_baseline = wait_for_object_on_target(&target_client, bucket, baseline_key).await?; + assert_eq!(replicated_baseline, baseline_payload); + + let old_target_address = target_env.address.clone(); + let new_target_port = RustFSTestEnvironment::find_available_port().await?; + let new_target_address = format!("127.0.0.1:{new_target_port}"); + let new_target_url = format!("http://{new_target_address}"); remote_peer.sync_state = SyncStatus::Enable; + remote_peer.endpoint = new_target_url.clone(); let edit_status = site_replication_edit(&source_env, "", &remote_peer).await?; assert!(edit_status.success, "unexpected site edit result: {:?}", edit_status); let source_after_sync = wait_for_site_replication_info(&source_env, |info| { info.sites .iter() - .any(|peer| peer.endpoint == target_env.url && peer.sync_state == SyncStatus::Enable) + .any(|peer| peer.endpoint == new_target_url && peer.sync_state == SyncStatus::Enable) }) .await?; let target_after_sync = wait_for_site_replication_info(&target_env, |info| { info.sites .iter() - .any(|peer| peer.endpoint == target_env.url && peer.sync_state == SyncStatus::Enable) + .any(|peer| peer.endpoint == new_target_url && peer.sync_state == SyncStatus::Enable) + }) + .await?; + let relay_after_sync = wait_for_site_replication_info(&relay_env, |info| { + info.sites + .iter() + .any(|peer| peer.endpoint == new_target_url && peer.sync_state == SyncStatus::Enable) }) .await?; assert!( source_after_sync .sites .iter() - .any(|peer| peer.endpoint == target_env.url && peer.sync_state == SyncStatus::Enable) + .any(|peer| peer.endpoint == new_target_url && peer.sync_state == SyncStatus::Enable) ); assert!( target_after_sync .sites .iter() - .any(|peer| peer.endpoint == target_env.url && peer.sync_state == SyncStatus::Enable) + .any(|peer| peer.endpoint == new_target_url && peer.sync_state == SyncStatus::Enable) ); + assert!( + relay_after_sync + .sites + .iter() + .any(|peer| peer.endpoint == new_target_url && peer.sync_state == SyncStatus::Enable) + ); + assert_eq!(relay_after_sync.sites.len(), 3); + + for (env, unchanged_endpoint) in [(&source_env, &relay_env.address), (&relay_env, &source_env.address)] { + let mut endpoints_replaced = false; + let mut last_targets = Vec::new(); + for _ in 0..40 { + let response = list_replication_targets_request(env, Some(bucket)).await?; + if response.status() == StatusCode::OK { + let targets: Vec = response.json().await?; + let mut endpoints = targets + .iter() + .filter_map(|target| target.get("endpoint").and_then(|endpoint| endpoint.as_str())) + .collect::>(); + endpoints.sort_unstable(); + let mut expected = vec![new_target_address.as_str(), unchanged_endpoint.as_str()]; + expected.sort_unstable(); + endpoints_replaced = endpoints == expected; + last_targets = targets; + if endpoints_replaced { + break; + } + } + sleep(Duration::from_millis(250)).await; + } + assert!( + endpoints_replaced, + "site edit did not replace the bucket target endpoint on {}: {last_targets:?}", + env.address + ); + } + + target_env.stop_server(); + let old_endpoint_listener = tokio::net::TcpListener::bind(&old_target_address).await?; + target_env.address = new_target_address; + target_env.url = new_target_url.clone(); + target_env.start_rustfs_server_without_cleanup(vec![]).await?; + let moved_target_client = target_env.create_s3_client(); let ilm_edit_status = site_replication_edit(&source_env, "enableILMExpiryReplication=true", &PeerInfo::default()).await?; assert!(ilm_edit_status.success, "unexpected ilm edit result: {:?}", ilm_edit_status); let source_after_ilm = wait_for_site_replication_info(&source_env, |info| { - info.sites.len() == 2 && info.sites.iter().all(|peer| peer.replicate_ilm_expiry) + info.sites.len() == 3 && info.sites.iter().all(|peer| peer.replicate_ilm_expiry) }) .await?; let target_after_ilm = wait_for_site_replication_info(&target_env, |info| { - info.sites.len() == 2 && info.sites.iter().all(|peer| peer.replicate_ilm_expiry) + info.sites.len() == 3 && info.sites.iter().all(|peer| peer.replicate_ilm_expiry) }) .await?; assert!(source_after_ilm.sites.iter().all(|peer| peer.replicate_ilm_expiry)); @@ -2823,26 +2945,26 @@ async fn test_site_replication_edit_and_status_peer_state_real_dual_node() -> Re let status_query = "peer-state=true"; let source_status = wait_for_site_replication_status(&source_env, status_query, |status| { - status.peer_states.len() == 2 + status.peer_states.len() == 3 && status .peer_states .values() - .all(|state| state.peers.len() == 2 && state.peers.values().all(|peer| peer.replicate_ilm_expiry)) + .all(|state| state.peers.len() == 3 && state.peers.values().all(|peer| peer.replicate_ilm_expiry)) }) .await?; let target_status = wait_for_site_replication_status(&target_env, status_query, |status| { - status.peer_states.len() == 2 + status.peer_states.len() == 3 && status .peer_states .values() - .all(|state| state.peers.len() == 2 && state.peers.values().all(|peer| peer.replicate_ilm_expiry)) + .all(|state| state.peers.len() == 3 && state.peers.values().all(|peer| peer.replicate_ilm_expiry)) }) .await?; - assert_eq!(source_status.peer_states.len(), 2); - assert_eq!(target_status.peer_states.len(), 2); - assert!(source_status.peer_states.values().all(|state| state.peers.len() == 2)); - assert!(target_status.peer_states.values().all(|state| state.peers.len() == 2)); + assert_eq!(source_status.peer_states.len(), 3); + assert_eq!(target_status.peer_states.len(), 3); + assert!(source_status.peer_states.values().all(|state| state.peers.len() == 3)); + assert!(target_status.peer_states.values().all(|state| state.peers.len() == 3)); assert!( source_status .peer_states @@ -2856,6 +2978,35 @@ async fn test_site_replication_edit_and_status_peer_state_real_dual_node() -> Re .all(|state| state.peers.values().all(|peer| peer.replicate_ilm_expiry)) ); + source_client + .put_object() + .bucket(bucket) + .key(moved_key) + .body(ByteStream::from(moved_payload.clone())) + .send() + .await?; + relay_client + .put_object() + .bucket(bucket) + .key(relayed_key) + .body(ByteStream::from(relayed_payload.clone())) + .send() + .await?; + let no_old_endpoint_connection = async { + match tokio::time::timeout(Duration::from_secs(3), old_endpoint_listener.accept()).await { + Err(_) => Ok::<(), Box>(()), + Ok(Ok((_, peer_address))) => Err(format!("source contacted the old site endpoint from {peer_address}").into()), + Ok(Err(err)) => Err(err.into()), + } + }; + let (replicated_after_edit, replicated_from_relay, ()) = tokio::try_join!( + wait_for_object_on_target(&moved_target_client, bucket, moved_key), + wait_for_object_on_target(&moved_target_client, bucket, relayed_key), + no_old_endpoint_connection, + )?; + assert_eq!(replicated_after_edit, moved_payload); + assert_eq!(replicated_from_relay, relayed_payload); + Ok(()) } @@ -2872,6 +3023,14 @@ async fn test_site_replication_remove_all_real_dual_node() -> Result<(), Box Result<(), Box Result<(), Box assert_site_replication_bucket_detached(&target_env, racing_bucket).await?, + Err(err) if matches!(err.code(), Some("NoSuchBucket" | "NotFound")) => {} + Err(err) => return Err(err.into()), + } + assert_site_replication_bucket_detached(&target_env, bucket).await?; + + source_client + .put_object() + .bucket(bucket) + .key(post_remove_key) + .body(ByteStream::from_static(b"site replication must stay stopped")) + .send() + .await?; + let absence_deadline = tokio::time::Instant::now() + Duration::from_secs(10); + loop { + match target_client.get_object().bucket(bucket).key(post_remove_key).send().await { + Ok(_) => return Err("object reached the removed site replication target".into()), + Err(err) if matches!(err.code(), Some("NoSuchKey" | "NotFound" | "NoSuchVersion")) => { + if tokio::time::Instant::now() >= absence_deadline { + break; + } + sleep(Duration::from_millis(250)).await; + } + Err(err) => return Err(err.into()), + } + } + Ok(()) } @@ -3068,7 +3274,12 @@ async fn test_site_replication_replicates_object_with_bucket_versioning_real_dua let _target_info = wait_for_site_replication_enabled(&target_env, 2).await?; source_client.create_bucket().bucket(bucket).send().await?; - enable_bucket_versioning(&source_env, bucket).await?; + let versioning = source_client.get_bucket_versioning().bucket(bucket).send().await?; + assert_eq!( + versioning.status(), + Some(&BucketVersioningStatus::Enabled), + "site replication did not enable source bucket versioning" + ); let replication_response = signed_request( http::Method::GET, &format!("{}/{bucket}?replication", source_env.url), diff --git a/rustfs/src/admin/handlers/bucket_meta.rs b/rustfs/src/admin/handlers/bucket_meta.rs index 2251fdfaf..1a3c2b383 100644 --- a/rustfs/src/admin/handlers/bucket_meta.rs +++ b/rustfs/src/admin/handlers/bucket_meta.rs @@ -27,6 +27,7 @@ use crate::admin::storage_api::bucket::{ }; use crate::admin::storage_api::contract::bucket::{BucketOperations, BucketOptions, MakeBucketOptions}; use crate::admin::storage_api::error::StorageError; +use crate::storage::storage_api::lock_bucket_targets_metadata; use crate::{ admin::{ auth::validate_admin_request, @@ -836,6 +837,11 @@ impl Operation for ImportBucketMetadata { for (bucket_name, metadata) in &bucket_metadatas { for (config_file, data) in imported_configs_to_persist(metadata) { let site_replication_item = imported_config_to_site_replication_item(bucket_name, metadata, config_file, &data)?; + let targets_guard = if matches!(config_file, BUCKET_REPLICATION_CONFIG | BUCKET_TARGETS_FILE) { + Some(lock_bucket_targets_metadata(bucket_name).await) + } else { + None + }; if let Err(e) = metadata_sys::update(bucket_name, config_file, data).await { warn!( event = EVENT_ADMIN_BUCKET_META_STATE, @@ -853,6 +859,7 @@ impl Operation for ImportBucketMetadata { "failed to persist imported bucket metadata for {bucket_name}/{config_file}: {e}" )); } + drop(targets_guard); if let Some(item) = site_replication_item && let Err(err) = site_replication_bucket_meta_hook(item).await { diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index 4d659ae62..74f437170 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -29,6 +29,7 @@ use crate::admin::utils::read_compatible_admin_body; use crate::auth::{check_key_valid, get_session_token}; use crate::error::ApiError; use crate::server::{ADMIN_PREFIX, RemoteAddr}; +use crate::storage::storage_api::lock_bucket_targets_metadata; use http::{HeaderMap, HeaderValue, Uri}; use hyper::{Method, StatusCode}; use matchit::Params; @@ -431,6 +432,7 @@ impl Operation for SetRemoteTargetHandler { if remote_target.arn.is_empty() { return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, "ARN is empty".to_string())); } + let _targets_guard = lock_bucket_targets_metadata(bucket).await; if update { let Some(mut target) = bucket_target_sys @@ -577,6 +579,7 @@ impl Operation for RemoveRemoteTargetHandler { .map_err(ApiError::from)?; let sys = BucketTargetSys::get(); + let _targets_guard = lock_bucket_targets_metadata(bucket).await; let targets = sys.remove_target(bucket, arn_str).await.map_err(map_bucket_target_error)?; diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index ce05978ac..568aece4d 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -44,6 +44,7 @@ use crate::auth::{check_key_valid, get_session_token}; use crate::config::get_config_snapshot; use crate::error::ApiError; use crate::server::{ADMIN_PREFIX, RemoteAddr}; +use crate::storage::storage_api::lock_bucket_targets_metadata; use base64::Engine; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use base64::engine::general_purpose::URL_SAFE_NO_PAD; @@ -90,10 +91,10 @@ use serde::de::DeserializeOwned; use serde_json::Value; use sha2::{Digest, Sha256}; use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet}; -use std::sync::LazyLock; +use std::sync::{LazyLock, Mutex as StdMutex}; use std::time::{Duration, Instant}; use time::OffsetDateTime; -use tokio::sync::Mutex; +use tokio::sync::{Mutex, RwLock}; use tracing::{info, warn}; use url::{Url, form_urlencoded}; use uuid::Uuid; @@ -123,6 +124,10 @@ const IDENTITY_LDAP_SUB_SYS: &str = "identity_ldap"; const LEGACY_LDAP_SUB_SYS: &str = "ldapserverconfig"; const SITE_REPLICATION_PEER_JOIN_PATH: &str = "/rustfs/admin/v3/site-replication/peer/join"; const SITE_REPLICATION_PEER_EDIT_PATH: &str = "/rustfs/admin/v3/site-replication/peer/edit"; +const SITE_REPLICATION_PEER_EDIT_CAPABILITY_PATH: &str = + "/rustfs/admin/v3/site-replication/peer/edit-capabilities?capability=endpoint-target-refresh"; +const SITE_REPLICATION_PEER_EDIT_REFRESH_PATH: &str = "/rustfs/admin/v3/site-replication/peer/edit?refresh-targets=true"; +const SITE_REPLICATION_ENDPOINT_REFRESH_RETRY_PATH: &str = "internal:endpoint-target-refresh"; const SITE_REPLICATION_PEER_REMOVE_PATH: &str = "/rustfs/admin/v3/site-replication/peer/remove"; const RUSTFS_ADMIN_V3_PREFIX: &str = "/rustfs/admin/v3"; const MINIO_ADMIN_V3_PREFIX: &str = "/minio/admin/v3"; @@ -196,7 +201,71 @@ struct SiteReplicationPeerClientCache { } static SITE_REPLICATION_PEER_CLIENT: LazyLock>> = LazyLock::new(|| Mutex::new(None)); +// Lock order: lifecycle -> bucket operation -> state -> per-bucket metadata. static SITE_REPLICATION_STATE_LOCK: LazyLock> = LazyLock::new(|| Mutex::new(())); +static SITE_REPLICATION_LIFECYCLE_LOCK: LazyLock> = LazyLock::new(|| Mutex::new(())); +static SITE_REPLICATION_BUCKET_OP_LOCK: LazyLock> = LazyLock::new(|| RwLock::new(())); +static SITE_REPLICATION_ADD_BOOTSTRAP: LazyLock>> = + LazyLock::new(|| StdMutex::new(None)); + +struct SiteReplicationAddBootstrap { + token: Uuid, + buckets: HashSet, +} + +struct SiteReplicationAddInProgressGuard { + token: Uuid, + _lifecycle: SiteReplicationLifecycleGuard, +} + +struct SiteReplicationLifecycleGuard { + _guard: tokio::sync::MutexGuard<'static, ()>, +} + +impl SiteReplicationLifecycleGuard { + async fn acquire() -> Self { + Self { + _guard: SITE_REPLICATION_LIFECYCLE_LOCK.lock().await, + } + } +} + +impl SiteReplicationAddInProgressGuard { + fn start(lifecycle: SiteReplicationLifecycleGuard, buckets: HashSet) -> S3Result { + let token = Uuid::new_v4(); + let mut pending = SITE_REPLICATION_ADD_BOOTSTRAP.lock().map_err(|_| { + S3Error::with_message(S3ErrorCode::InternalError, "site replication bootstrap lock poisoned".to_string()) + })?; + *pending = Some(SiteReplicationAddBootstrap { token, buckets }); + Ok(Self { + token, + _lifecycle: lifecycle, + }) + } +} + +impl Drop for SiteReplicationAddInProgressGuard { + fn drop(&mut self) { + if let Ok(mut pending) = SITE_REPLICATION_ADD_BOOTSTRAP.lock() + && pending.as_ref().is_some_and(|bootstrap| bootstrap.token == self.token) + { + *pending = None; + } + } +} + +fn bootstrap_peer_bucket_operation_allowed(bucket: &str, operation: &str, bootstrap_token: Option<&str>) -> bool { + if !matches!(operation, "make-with-versioning" | "configure-replication") { + return false; + } + let parsed_token = bootstrap_token.and_then(|value| Uuid::parse_str(value).ok()); + SITE_REPLICATION_ADD_BOOTSTRAP.lock().is_ok_and(|pending| { + pending.as_ref().is_some_and(|bootstrap| { + parsed_token.is_some_and(|token| token == bootstrap.token) + || (bootstrap_token.is_none() && bootstrap.buckets.contains(bucket)) + }) + }) +} fn site_replication_peer_client_cache_hit( cache: &Option, @@ -229,6 +298,8 @@ struct SiteReplicationState { pending_rotation: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pending_remove: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pending_endpoint_refresh: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] retry_queue: Vec, } @@ -246,6 +317,23 @@ struct SiteReplicationRetryEvent { updated_at: Option, } +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +struct PendingEndpointRefresh { + id: String, + peer: PeerInfo, + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] + remote_peers: BTreeMap, + #[serde(default, skip_serializing_if = "BTreeSet::is_empty")] + acked_deployment_ids: BTreeSet, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct EndpointRefreshRequest { + id: String, + peer: PeerInfo, +} + #[derive(Debug, Clone, Serialize, Deserialize, Default)] struct PendingRotation { id: String, @@ -290,6 +378,7 @@ struct SiteReplicationAddPreflightInfo { deployment_id: String, enabled: bool, bucket_count: usize, + bucket_names: HashSet, peer_deployment_ids: BTreeSet, idp_settings: serde_json::Value, } @@ -422,6 +511,11 @@ pub fn register_site_replication_route(r: &mut S3Router) -> std: AdminOperation(&SRPeerGetIDPSettingsHandler {}), ), (Method::PUT, "/v3/site-replication/edit", AdminOperation(&SiteReplicationEditHandler {})), + ( + Method::PUT, + "/v3/site-replication/peer/edit-capabilities", + AdminOperation(&SRPeerEditCapabilitiesHandler {}), + ), (Method::PUT, "/v3/site-replication/peer/edit", AdminOperation(&SRPeerEditHandler {})), (Method::PUT, "/v3/site-replication/peer/remove", AdminOperation(&SRPeerRemoveHandler {})), ( @@ -1118,12 +1212,14 @@ fn add_preflight_info_from_sr_info( info: SRInfo, idp_settings: IDPSettings, ) -> S3Result { + let bucket_names = info.buckets.keys().cloned().collect(); Ok(SiteReplicationAddPreflightInfo { name: if info.name.is_empty() { site.name.clone() } else { info.name }, endpoint: site.endpoint.clone(), deployment_id: info.deployment_id, enabled: info.enabled, bucket_count: info.buckets.len(), + bucket_names, peer_deployment_ids: info.state.peers.keys().cloned().collect(), idp_settings: idp_settings_value(&idp_settings)?, }) @@ -1248,6 +1344,18 @@ fn bootstrap_bucket_op_path(bucket: &str, operation: &str) -> String { ) } +fn with_site_replication_bootstrap_token(path: &str, token: &str) -> String { + let separator = if path.contains('?') { '&' } else { '?' }; + let query = form_urlencoded::Serializer::new(String::new()) + .append_pair("bootstrapToken", token) + .finish(); + format!("{path}{separator}{query}") +} + +fn site_replication_bootstrap_token(uri: &Uri) -> Option { + query_pairs(uri).get("bootstrapToken").cloned() +} + fn bootstrap_bucket_make_op_path(bucket: &SRBucketInfo) -> String { let mut query = form_urlencoded::Serializer::new(String::new()); query.append_pair("bucket", &bucket.bucket); @@ -1683,13 +1791,13 @@ fn site_replication_peer_payload(path: &str, secret_key: &str, payload: Vec) } } -async fn send_peer_admin_request( +async fn send_peer_admin_request_raw( endpoint: &str, path: &str, access_key: &str, secret_key: &str, body: &T, -) -> S3Result> { +) -> S3Result<(StatusCode, Vec)> { let path = site_replication_peer_wire_path(path); let base = endpoint.trim_end_matches('/'); let url = format!("{base}{path}"); @@ -1750,15 +1858,26 @@ async fn send_peer_admin_request( .await .map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("read peer response failed: {e}")))?; - if !status.is_success() { - let detail = String::from_utf8_lossy(&body).into_owned(); - return Err(S3Error::with_message( - S3ErrorCode::InternalError, - format!("peer request to {url} failed with {status}: {detail}"), - )); + Ok((status, body.to_vec())) +} + +async fn send_peer_admin_request( + endpoint: &str, + path: &str, + access_key: &str, + secret_key: &str, + body: &T, +) -> S3Result> { + let (status, body) = send_peer_admin_request_raw(endpoint, path, access_key, secret_key, body).await?; + if status.is_success() { + return Ok(body); } - Ok(body.to_vec()) + let detail = String::from_utf8_lossy(&body).into_owned(); + Err(S3Error::with_message( + S3ErrorCode::InternalError, + format!("peer request to {endpoint}{path} failed with {status}: {detail}"), + )) } async fn send_peer_admin_request_with_secret_candidates( @@ -2050,27 +2169,38 @@ async fn bootstrap_existing_metadata_after_add( } pub async fn site_replication_make_bucket_hook(bucket: &str, lock_enabled: bool) -> S3Result<()> { - let Some(runtime) = runtime_site_replication_targets().await? else { - return Ok(()); + let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await; + let runtime = { + let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; + let Some(runtime) = runtime_site_replication_targets().await? else { + return Ok(()); + }; + + ensure_site_replication_bucket_versioning(bucket).await?; + ensure_site_replication_bucket_setup_with_runtime(bucket, &runtime).await?; + runtime }; - ensure_site_replication_bucket_versioning(bucket).await?; - ensure_site_replication_bucket_targets( - bucket, - &runtime.state, - &runtime.local_peer, - None, - &runtime.service_account_secret_key, - ) - .await?; - ensure_site_replication_bucket_replication_config( - bucket, - &runtime.state, - &runtime.local_peer, - &runtime.service_account_secret_key, - ) - .await?; + broadcast_site_replication_make_bucket(bucket, lock_enabled, Some(&runtime), None).await +} +async fn broadcast_site_replication_json_using_runtime( + runtime: Option<&SiteReplicationRuntime>, + path: &str, + body: &T, +) -> S3Result<()> { + match runtime { + Some(runtime) => broadcast_site_replication_json_with_runtime(runtime, path, body).await, + None => broadcast_site_replication_json(path, body).await, + } +} + +async fn broadcast_site_replication_make_bucket( + bucket: &str, + lock_enabled: bool, + runtime: Option<&SiteReplicationRuntime>, + bootstrap_token: Option<&str>, +) -> S3Result<()> { let created_at = current_object_store_handle() .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))? .get_bucket_info(bucket, &BucketOptions::default()) @@ -2091,16 +2221,20 @@ pub async fn site_replication_make_bucket_hook(bucket: &str, lock_enabled: bool) } format!("/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", query.finish()) }; - broadcast_site_replication_json(&path, &serde_json::json!({})).await?; + let path = if let Some(token) = bootstrap_token { + with_site_replication_bootstrap_token(&path, token) + } else { + path + }; + broadcast_site_replication_json_using_runtime(runtime, &path, &serde_json::json!({})).await?; - let configure_path = format!( - "/rustfs/admin/v3/site-replication/peer/bucket-ops?{}", - form_urlencoded::Serializer::new(String::new()) - .append_pair("bucket", bucket) - .append_pair("operation", "configure-replication") - .finish() - ); - broadcast_site_replication_json(&configure_path, &serde_json::json!({})).await + let configure_path = bootstrap_bucket_op_path(bucket, "configure-replication"); + let configure_path = if let Some(token) = bootstrap_token { + with_site_replication_bootstrap_token(&configure_path, token) + } else { + configure_path + }; + broadcast_site_replication_json_using_runtime(runtime, &configure_path, &serde_json::json!({})).await } pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: bool) -> S3Result<()> { @@ -3093,11 +3227,281 @@ fn edit_state(mut state: SiteReplicationState, incoming: PeerInfo, ilm_expiry_ov state } +fn peer_endpoint_edit_requested(state: &SiteReplicationState, incoming: &PeerInfo) -> bool { + !incoming.deployment_id.is_empty() && !incoming.endpoint.is_empty() && state.peers.contains_key(&incoming.deployment_id) +} + +fn peer_endpoint_refresh_requested(state: &SiteReplicationState, incoming: &PeerInfo) -> bool { + if !peer_endpoint_edit_requested(state, incoming) { + return false; + } + if let Some(pending) = pending_endpoint_refresh(state) { + return pending.peer.deployment_id == incoming.deployment_id + && canonical_endpoint(&pending.peer.endpoint) == canonical_endpoint(&incoming.endpoint); + } + state + .peers + .get(&incoming.deployment_id) + .is_some_and(|peer| canonical_endpoint(&peer.endpoint) != canonical_endpoint(&incoming.endpoint)) +} + +fn pending_endpoint_refresh(state: &SiteReplicationState) -> Option { + state.pending_endpoint_refresh.clone().or_else(|| { + state + .retry_queue + .iter() + .find(|event| event.path == SITE_REPLICATION_ENDPOINT_REFRESH_RETRY_PATH) + .and_then(|event| serde_json::from_str(&event.last_error).ok()) + }) +} + +fn set_pending_endpoint_refresh(state: &mut SiteReplicationState, pending: PendingEndpointRefresh) -> S3Result<()> { + let encoded = serde_json::to_string(&pending) + .map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize endpoint refresh state: {err}")))?; + state + .retry_queue + .retain(|event| event.path != SITE_REPLICATION_ENDPOINT_REFRESH_RETRY_PATH); + state.retry_queue.push(SiteReplicationRetryEvent { + id: pending.id.clone(), + peer_deployment_id: pending.peer.deployment_id.clone(), + peer_endpoint: pending.peer.endpoint.clone(), + path: SITE_REPLICATION_ENDPOINT_REFRESH_RETRY_PATH.to_string(), + retry_count: 0, + failed: false, + last_error: encoded, + updated_at: Some(OffsetDateTime::now_utc()), + }); + state.pending_endpoint_refresh = Some(pending); + Ok(()) +} + +fn clear_pending_endpoint_refresh(state: &mut SiteReplicationState) { + state.pending_endpoint_refresh = None; + state + .retry_queue + .retain(|event| event.path != SITE_REPLICATION_ENDPOINT_REFRESH_RETRY_PATH); +} + +fn parse_endpoint_refresh_status(peer: &PeerInfo, body: &[u8]) -> S3Result<()> { + let status: ReplicateEditStatus = serde_json::from_slice(body).map_err(|_| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("peer {} does not support endpoint target refresh", peer.endpoint), + ) + })?; + if status.success { + Ok(()) + } else { + Err(S3Error::with_message( + S3ErrorCode::InternalError, + format!("peer {} failed endpoint target refresh: {}", peer.endpoint, status.err_detail), + )) + } +} + +fn endpoint_refresh_capability_supported(peer: &PeerInfo, status: StatusCode, body: &[u8]) -> S3Result { + if status.is_success() { + return Ok(parse_endpoint_refresh_status(peer, body).is_ok()); + } + if matches!(status, StatusCode::BAD_REQUEST | StatusCode::NOT_FOUND | StatusCode::METHOD_NOT_ALLOWED) { + return Ok(false); + } + + Err(S3Error::with_message( + S3ErrorCode::InternalError, + format!( + "probe endpoint target refresh on peer {} failed with {status}: {}", + peer.endpoint, + String::from_utf8_lossy(body) + ), + )) +} + +fn endpoint_refresh_route_endpoints(target: &PeerInfo, pending: &PendingEndpointRefresh) -> Vec { + let mut endpoints = vec![target.endpoint.clone()]; + if target.deployment_id == pending.peer.deployment_id + && canonical_endpoint(&target.endpoint) != canonical_endpoint(&pending.peer.endpoint) + { + endpoints.push(pending.peer.endpoint.clone()); + } + endpoints +} + +fn endpoint_refresh_remote_targets<'a>( + routing_peers: &'a BTreeMap, + pending: Option<&PendingEndpointRefresh>, + local_deployment_id: Option<&str>, +) -> Vec<&'a PeerInfo> { + routing_peers + .values() + .filter(|target| { + local_deployment_id.is_none_or(|deployment_id| deployment_id != target.deployment_id) + && pending.is_none_or(|pending| !pending.acked_deployment_ids.contains(&target.deployment_id)) + }) + .collect() +} + +async fn send_endpoint_refresh_admin_request( + target: &PeerInfo, + pending: &PendingEndpointRefresh, + path: &str, + access_key: &str, + secret_key: &str, + body: &T, +) -> S3Result> { + let (status, response) = send_endpoint_refresh_admin_request_raw(target, pending, path, access_key, secret_key, body).await?; + if status.is_success() { + return Ok(response); + } + + Err(S3Error::with_message( + S3ErrorCode::InternalError, + format!( + "peer {} endpoint target refresh failed with {status}: {}", + target.endpoint, + String::from_utf8_lossy(&response) + ), + )) +} + +async fn send_endpoint_refresh_admin_request_raw( + target: &PeerInfo, + pending: &PendingEndpointRefresh, + path: &str, + access_key: &str, + secret_key: &str, + body: &T, +) -> S3Result<(StatusCode, Vec)> { + let mut last_error = None; + let mut last_response = None; + for endpoint in endpoint_refresh_route_endpoints(target, pending) { + match send_peer_admin_request_raw(&endpoint, path, access_key, secret_key, body).await { + Ok((status, response)) + if matches!(status, StatusCode::NOT_FOUND | StatusCode::METHOD_NOT_ALLOWED | StatusCode::GONE) + || status.is_server_error() => + { + last_response = Some((status, response)); + } + Ok(response) => return Ok(response), + Err(err) => last_error = Some(err), + } + } + + if let Some(response) = last_response { + return Ok(response); + } + Err(last_error.unwrap_or_else(|| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("peer {} endpoint target refresh failed", target.endpoint), + ) + })) +} + +async fn legacy_peer_bucket_names( + target: &PeerInfo, + pending: &PendingEndpointRefresh, + access_key: &str, + secret_key: &str, +) -> S3Result> { + let mut last_error = None; + for endpoint in endpoint_refresh_route_endpoints(target, pending) { + match send_peer_admin_get_request( + &endpoint, + "/rustfs/admin/v3/site-replication/metainfo?buckets=true", + access_key, + secret_key, + ) + .await + { + Ok(body) => return peer_bucket_names_from_metainfo(&endpoint, &body), + Err(err) => last_error = Some(err), + } + } + + Err(last_error.unwrap_or_else(|| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("list site replication buckets on peer {} failed", target.endpoint), + ) + })) +} + +fn peer_bucket_names_from_metainfo(endpoint: &str, body: &[u8]) -> S3Result> { + let info: Value = serde_json::from_slice(body).map_err(|err| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("parse site replication metainfo from {endpoint} failed: {err}"), + ) + })?; + let Some(buckets) = info.get("buckets").or_else(|| info.get("Buckets")) else { + return Ok(Vec::new()); + }; + let buckets = buckets.as_object().ok_or_else(|| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("site replication metainfo from {endpoint} has invalid buckets"), + ) + })?; + Ok(buckets.keys().cloned().collect()) +} + +async fn refresh_legacy_peer_bucket_targets( + target: &PeerInfo, + pending: &PendingEndpointRefresh, + access_key: &str, + secret_key: &str, +) -> S3Result<()> { + send_endpoint_refresh_admin_request(target, pending, SITE_REPLICATION_PEER_EDIT_PATH, access_key, secret_key, &pending.peer) + .await?; + + let buckets = legacy_peer_bucket_names(target, pending, access_key, secret_key).await?; + let mut configure_operation = None; + for bucket in &buckets { + if let Some(operation) = configure_operation { + let path = bootstrap_bucket_op_path(bucket, operation); + send_endpoint_refresh_admin_request(target, pending, &path, access_key, secret_key, &serde_json::json!({})).await?; + continue; + } + + let minio_path = bootstrap_bucket_op_path(bucket, "ConfigureReplication"); + let (status, _) = + send_endpoint_refresh_admin_request_raw(target, pending, &minio_path, access_key, secret_key, &serde_json::json!({})) + .await?; + if status.is_success() { + configure_operation = Some("ConfigureReplication"); + continue; + } + + let rustfs_path = bootstrap_bucket_op_path(bucket, SITE_REPLICATION_BUCKET_OP_CONFIGURE_REPLICATION); + send_endpoint_refresh_admin_request(target, pending, &rustfs_path, access_key, secret_key, &serde_json::json!({})) + .await?; + configure_operation = Some(SITE_REPLICATION_BUCKET_OP_CONFIGURE_REPLICATION); + } + + Ok(()) +} + +fn align_peer_edit_deployment_id(state: &SiteReplicationState, incoming: &mut PeerInfo) { + if incoming.name.is_empty() || state.peers.contains_key(&incoming.deployment_id) { + return; + } + + let mut matches = state.peers.values().filter(|peer| peer.name == incoming.name); + let Some(peer) = matches.next() else { + return; + }; + if matches.next().is_none() { + incoming.deployment_id = peer.deployment_id.clone(); + } +} + fn remove_sites(mut state: SiteReplicationState, req: SRRemoveReq) -> SiteReplicationState { if req.remove_all { state.peers.clear(); state.resync_status.clear(); state.retry_queue.clear(); + state.pending_endpoint_refresh = None; state.updated_at = Some(OffsetDateTime::now_utc()); return state; } @@ -3107,6 +3511,7 @@ fn remove_sites(mut state: SiteReplicationState, req: SRRemoveReq) -> SiteReplic state.peers.clear(); state.resync_status.clear(); state.retry_queue.clear(); + state.pending_endpoint_refresh = None; state.updated_at = Some(OffsetDateTime::now_utc()); return state; } @@ -3129,6 +3534,13 @@ fn remove_sites(mut state: SiteReplicationState, req: SRRemoveReq) -> SiteReplic state .resync_status .retain(|deployment_id, _| state.peers.contains_key(deployment_id)); + if state + .pending_endpoint_refresh + .as_ref() + .is_some_and(|pending| !state.peers.contains_key(&pending.peer.deployment_id)) + { + clear_pending_endpoint_refresh(&mut state); + } state.updated_at = Some(OffsetDateTime::now_utc()); state } @@ -3577,11 +3989,8 @@ fn site_replication_target_arns_by_peer(config: Option<&s3s::dto::ReplicationCon } for arn in configured_arns { - if let Ok(parsed) = arn.parse::() - && parsed.arn_type == BucketTargetType::ReplicationService - && !parsed.id.is_empty() - { - arns_by_peer.entry(parsed.id).or_insert(arn); + if let Some(deployment_id) = replication_target_arn_deployment_id(&arn) { + arns_by_peer.entry(deployment_id).or_insert(arn); } } @@ -3711,24 +4120,14 @@ fn bucket_target_deployment_id(target: &BucketTarget) -> Option { } fn replication_target_arn_deployment_id(arn: &str) -> Option { - if let Ok(parsed) = arn.parse::() - && parsed.arn_type == BucketTargetType::ReplicationService - { - if !parsed.id.is_empty() { - return Some(parsed.id); - } - if !parsed.region.is_empty() { - return Some(parsed.region); - } - } - let parts: Vec<_> = arn.split(':').collect(); - if parts.len() == 6 && parts[0] == "arn" && parts[1] == "rustfs" && parts[2] == "replication" { - return if parts[3].is_empty() { - (!parts[4].is_empty()).then(|| parts[4].to_string()) - } else { - Some(parts[3].to_string()) - }; + if parts.len() == 6 + && parts[0] == "arn" + && matches!(parts[1], "rustfs" | "minio") + && parts[2] == "replication" + && !parts[4].is_empty() + { + return Some(parts[4].to_string()); } None @@ -3870,7 +4269,7 @@ fn build_site_replication_config( } } -async fn ensure_site_replication_bucket_targets( +async fn ensure_site_replication_bucket_targets_with_runtime( bucket: &str, state: &SiteReplicationState, local_peer: &PeerInfo, @@ -3899,7 +4298,31 @@ async fn ensure_site_replication_bucket_targets( Ok(()) } -async fn ensure_site_replication_bucket_replication_config( +async fn bucket_replication_config_for_target_refresh(bucket: &str) -> S3Result> { + match metadata_sys::get_replication_config(bucket).await { + Ok((config, _)) => Ok(Some(config)), + Err(StorageError::ConfigNotFound) => Ok(None), + Err(err) => Err(ApiError::from(err).into()), + } +} + +async fn ensure_site_replication_bucket_targets(bucket: &str) -> S3Result<()> { + let _targets_guard = lock_bucket_targets_metadata(bucket).await; + let Some(runtime) = runtime_site_replication_targets().await? else { + return Ok(()); + }; + let config = bucket_replication_config_for_target_refresh(bucket).await?; + ensure_site_replication_bucket_targets_with_runtime( + bucket, + &runtime.state, + &runtime.local_peer, + config.as_ref(), + &runtime.service_account_secret_key, + ) + .await +} + +async fn ensure_site_replication_bucket_replication_config_with_runtime( bucket: &str, state: &SiteReplicationState, local_peer: &PeerInfo, @@ -3961,7 +4384,37 @@ async fn ensure_site_replication_bucket_replication_config( Ok(()) } +async fn ensure_site_replication_bucket_setup(bucket: &str) -> S3Result { + let Some(runtime) = runtime_site_replication_targets().await? else { + return Ok(false); + }; + ensure_site_replication_bucket_setup_with_runtime(bucket, &runtime).await?; + Ok(true) +} + +async fn ensure_site_replication_bucket_setup_with_runtime(bucket: &str, runtime: &SiteReplicationRuntime) -> S3Result<()> { + let _targets_guard = lock_bucket_targets_metadata(bucket).await; + let config = bucket_replication_config_for_target_refresh(bucket).await?; + ensure_site_replication_bucket_targets_with_runtime( + bucket, + &runtime.state, + &runtime.local_peer, + config.as_ref(), + &runtime.service_account_secret_key, + ) + .await?; + ensure_site_replication_bucket_replication_config_with_runtime( + bucket, + &runtime.state, + &runtime.local_peer, + &runtime.service_account_secret_key, + ) + .await?; + Ok(()) +} + async fn cleanup_removed_site_replication_bucket(bucket: &str, removed_deployment_ids: &HashSet) -> S3Result { + let _targets_guard = lock_bucket_targets_metadata(bucket).await; let mut removed = 0usize; match metadata_sys::list_bucket_targets(bucket).await { @@ -4054,11 +4507,7 @@ pub async fn site_replication_peer_deployment_id_for_endpoint(endpoint: &str) -> /// that already exists locally, wire up versioning + targets + replication config for each, and /// kick a resync toward every remote peer so pre-existing objects back-fill. Errors are logged /// but never abort the caller — the admin can run a manual resync if needed. -async fn backfill_existing_buckets_after_add( - state: &SiteReplicationState, - local_peer: &PeerInfo, - service_account_secret_key: &str, -) { +async fn backfill_existing_buckets_after_add(state: &SiteReplicationState, local_peer: &PeerInfo, bootstrap_token: Option<&str>) { let Some(store) = current_object_store_handle() else { return; }; @@ -4093,27 +4542,13 @@ async fn backfill_existing_buckets_after_add( ); continue; } - if let Err(err) = ensure_site_replication_bucket_targets(name, state, local_peer, None, service_account_secret_key).await - { + if let Err(err) = ensure_site_replication_bucket_setup(name).await { warn!( event = EVENT_ADMIN_SITE_REPLICATION_STATE, component = LOG_COMPONENT_ADMIN, subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, bucket = %name, - result = "backfill_targets_setup_failed", - error = ?err, - "admin site replication state" - ); - } - if let Err(err) = - ensure_site_replication_bucket_replication_config(name, state, local_peer, service_account_secret_key).await - { - warn!( - event = EVENT_ADMIN_SITE_REPLICATION_STATE, - component = LOG_COMPONENT_ADMIN, - subsystem = LOG_SUBSYSTEM_SITE_REPLICATION, - bucket = %name, - result = "backfill_replication_config_setup_failed", + result = "backfill_bucket_setup_failed", error = ?err, "admin site replication state" ); @@ -4137,7 +4572,7 @@ async fn backfill_existing_buckets_after_add( false } }; - if let Err(err) = site_replication_make_bucket_hook(name, lock_enabled).await { + if let Err(err) = broadcast_site_replication_make_bucket(name, lock_enabled, None, bootstrap_token).await { warn!( event = EVENT_ADMIN_SITE_REPLICATION_STATE, component = LOG_COMPONENT_ADMIN, @@ -4170,11 +4605,7 @@ async fn backfill_existing_buckets_after_add( } } -async fn refresh_bucket_targets_after_service_account_rotation( - state: &SiteReplicationState, - local_peer: &PeerInfo, - service_account_secret_key: &str, -) { +async fn refresh_bucket_targets_after_service_account_rotation() { let Some(store) = current_object_store_handle() else { return; }; @@ -4194,19 +4625,7 @@ async fn refresh_bucket_targets_after_service_account_rotation( }; for bucket in buckets { - let replication_config = metadata_sys::get_replication_config(&bucket.name) - .await - .ok() - .map(|(config, _)| config); - if let Err(err) = ensure_site_replication_bucket_targets( - &bucket.name, - state, - local_peer, - replication_config.as_ref(), - service_account_secret_key, - ) - .await - { + if let Err(err) = ensure_site_replication_bucket_targets(&bucket.name).await { warn!( event = EVENT_ADMIN_SITE_REPLICATION_STATE, component = LOG_COMPONENT_ADMIN, @@ -4220,12 +4639,40 @@ async fn refresh_bucket_targets_after_service_account_rotation( } } +async fn refresh_bucket_targets_after_endpoint_edit(pending_id: &str, service_account_secret_key: &str) -> S3Result<()> { + let store = current_object_store_handle() + .ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "object store is not initialized".to_string()))?; + let buckets = store.list_bucket(&BucketOptions::default()).await.map_err(ApiError::from)?; + + for bucket in buckets { + let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; + let state = load_site_replication_state().await?; + if pending_endpoint_refresh(&state).is_none_or(|pending| pending.id != pending_id) { + return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); + } + let local_peer = current_local_runtime_peer(&state); + let _targets_guard = lock_bucket_targets_metadata(&bucket.name).await; + let replication_config = bucket_replication_config_for_target_refresh(&bucket.name).await?; + ensure_site_replication_bucket_targets_with_runtime( + &bucket.name, + &state, + &local_peer, + replication_config.as_ref(), + service_account_secret_key, + ) + .await?; + } + + Ok(()) +} + async fn start_site_bucket_resync(bucket: &str, peer: &PeerInfo, resync_id: &str) -> ResyncBucketStatus { let mut bucket_status = ResyncBucketStatus { bucket: bucket.to_string(), status: "started".to_string(), ..Default::default() }; + let targets_guard = lock_bucket_targets_metadata(bucket).await; let (config, _) = match metadata_sys::get_replication_config(bucket).await { Ok(config) => config, @@ -4282,6 +4729,7 @@ async fn start_site_bucket_resync(bucket: &str, peer: &PeerInfo, resync_id: &str return bucket_status; } BucketTargetSys::get().update_all_targets(bucket, Some(&targets)).await; + drop(targets_guard); let Some(pool) = current_replication_pool_handle() else { bucket_status.status = "failed".to_string(); @@ -4306,6 +4754,7 @@ async fn cancel_site_bucket_resync(bucket: &str, peer: &PeerInfo, resync_id: &st status: "canceled".to_string(), ..Default::default() }; + let targets_guard = lock_bucket_targets_metadata(bucket).await; let mut targets = match metadata_sys::list_bucket_targets(bucket).await { Ok(targets) => targets, @@ -4345,6 +4794,7 @@ async fn cancel_site_bucket_resync(bucket: &str, peer: &PeerInfo, resync_id: &st return bucket_status; } BucketTargetSys::get().update_all_targets(bucket, Some(&targets)).await; + drop(targets_guard); let Some(pool) = current_replication_pool_handle() else { bucket_status.status = "failed".to_string(); @@ -4464,6 +4914,11 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> { } else { item.updated_at }; + let targets_guard = if item.r#type == "replication-config" { + Some(lock_bucket_targets_metadata(&item.bucket).await) + } else { + None + }; if let Ok(bucket_meta) = metadata_sys::get(&item.bucket).await { let local_updated_at = bucket_meta_local_updated_at(&bucket_meta, config_file); if is_stale_update(local_updated_at, incoming_updated_at) { @@ -4471,7 +4926,7 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> { } } - let replication_config = if item.r#type == "replication-config" { + if item.r#type == "replication-config" { item.replication_config .as_ref() .map(|raw| { @@ -4479,10 +4934,8 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> { deserialize::(&data) }) .transpose() - .map_err(|e| s3_error!(InvalidRequest, "invalid replication config: {e}"))? - } else { - None - }; + .map_err(|e| s3_error!(InvalidRequest, "invalid replication config: {e}"))?; + } let data = match item.r#type.as_str() { "policy" => item @@ -4514,18 +4967,10 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> { .await .map_err(ApiError::from)?; } + drop(targets_guard); - if item.r#type == "replication-config" - && let Some(runtime) = runtime_site_replication_targets().await? - { - ensure_site_replication_bucket_targets( - &item.bucket, - &runtime.state, - &runtime.local_peer, - replication_config.as_ref(), - &runtime.service_account_secret_key, - ) - .await?; + if item.r#type == "replication-config" { + ensure_site_replication_bucket_targets(&item.bucket).await?; } if item.r#type == "version-config" @@ -4533,23 +4978,8 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> { .await .ok() .is_some_and(|(config, _)| config.enabled()) - && let Some(runtime) = runtime_site_replication_targets().await? { - ensure_site_replication_bucket_targets( - &item.bucket, - &runtime.state, - &runtime.local_peer, - replication_config.as_ref(), - &runtime.service_account_secret_key, - ) - .await?; - ensure_site_replication_bucket_replication_config( - &item.bucket, - &runtime.state, - &runtime.local_peer, - &runtime.service_account_secret_key, - ) - .await?; + ensure_site_replication_bucket_setup(&item.bucket).await?; } Ok(()) @@ -4808,8 +5238,12 @@ impl Operation for SiteReplicationAddHandler { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationAddAction).await?; reject_site_replicator_on_public_admin(&cred)?; let replicate_ilm_expiry = sr_add_replicate_ilm_expiry(&req.uri); + let lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; let current_state = load_site_replication_state().await?; + if pending_endpoint_refresh(¤t_state).is_some() { + return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); + } let local_peer = current_local_peer(&req, ¤t_state); let sites: Vec = read_site_replication_json(req, &cred.secret_key, true).await?; validate_add_sites(&sites, &local_peer)?; @@ -4824,6 +5258,12 @@ impl Operation for SiteReplicationAddHandler { validate_add_preflight_topology(&preflight_infos, &local_peer)?; let (service_account_access_key, service_account_secret_key) = ensure_site_replicator_service_account(&cred.access_key, false).await?; + let bootstrap_buckets = preflight_infos + .iter() + .filter(|info| !same_identity_endpoint(&info.endpoint, &local_peer.endpoint)) + .flat_map(|info| info.bucket_names.iter().cloned()) + .collect(); + let add_in_progress_guard = SiteReplicationAddInProgressGuard::start(lifecycle_guard, bootstrap_buckets)?; let mut state = merge_add_sites( current_state, local_peer.clone(), @@ -4839,6 +5279,8 @@ impl Operation for SiteReplicationAddHandler { peers: state.peers.clone(), updated_at: state.updated_at, }; + let peer_join_path = + with_site_replication_bootstrap_token(SITE_REPLICATION_PEER_JOIN_PATH, &add_in_progress_guard.token.to_string()); let mut joined_endpoints = HashSet::new(); for site in &sites { @@ -4850,14 +5292,9 @@ impl Operation for SiteReplicationAddHandler { let mut peer_join_req = join_req.clone(); peer_join_req.svc_acct_parent = site.access_key.clone(); - let body = send_peer_admin_request( - &site.endpoint, - SITE_REPLICATION_PEER_JOIN_PATH, - &site.access_key, - &site.secret_key, - &peer_join_req, - ) - .await?; + let body = + send_peer_admin_request(&site.endpoint, &peer_join_path, &site.access_key, &site.secret_key, &peer_join_req) + .await?; let join_response: SRPeerJoinResponse = serde_json::from_slice(&body).map_err(|e| { S3Error::with_message( @@ -4874,7 +5311,7 @@ impl Operation for SiteReplicationAddHandler { // Fix 1: back-fill pre-existing buckets so objects created before `replicate add` // are not silently left out of replication. Failures are logged but do not abort // the overall add operation — the admin can trigger a manual resync if needed. - backfill_existing_buckets_after_add(&state, &local_peer, &service_account_secret_key).await; + backfill_existing_buckets_after_add(&state, &local_peer, None).await; json_response(&ReplicateAddStatus { success: true, @@ -4893,9 +5330,17 @@ impl Operation for SiteReplicationRemoveHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationRemoveAction).await?; reject_site_replicator_on_public_admin(&cred)?; + let _lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; let (pending_remove, local_peer) = { + let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.write().await; let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; let current_state = load_site_replication_state().await?; + if pending_endpoint_refresh(¤t_state).is_some() { + return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); + } + if current_state.pending_rotation.is_some() { + return Err(s3_error!(InvalidRequest, "service account rotation is pending")); + } let local_peer = current_local_peer(&req, ¤t_state); let remove_req: SRRemoveReq = read_site_replication_json(req, "", false).await?; @@ -4969,6 +5414,7 @@ impl Operation for SiteReplicationRemoveHandler { let finalize_candidate = pending_remove_ready_to_finalize(&pending_remove.id, &local_peer).await?; let complete = if let Some(finalized_remove) = finalize_candidate { + let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.write().await; let removed_deployment_ids = removed_deployment_ids_for_pending_remove(&finalized_remove, &local_peer); match cleanup_removed_site_replication_buckets(&removed_deployment_ids).await { Ok(removed) => { @@ -5094,6 +5540,7 @@ pub struct SRPeerJoinHandler {} impl Operation for SRPeerJoinHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationAddAction).await?; + let bootstrap_token = site_replication_bootstrap_token(&req.uri); let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; let mut state = load_site_replication_state().await?; let local_peer = current_local_peer(&req, &state); @@ -5175,7 +5622,7 @@ impl Operation for SRPeerJoinHandler { persist_site_replication_state(&state).await?; // Fix 1 (receiving side): ensure the joining peer also sets up replication for any // buckets it already owns so the reverse direction works from the start. - backfill_existing_buckets_after_add(&state, &local_peer, &join_req.svc_acct_secret_key).await; + backfill_existing_buckets_after_add(&state, &local_peer, bootstrap_token.as_deref()).await; json_response(&SRPeerJoinResponse { peer: state.peers.get(&local_peer.deployment_id).cloned().unwrap_or(local_peer), }) @@ -5188,6 +5635,8 @@ pub struct SRPeerBucketOpsHandler {} impl Operation for SRPeerBucketOpsHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; + let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await; + let state = load_site_replication_state().await?; let queries = query_pairs(&req.uri); let bucket = queries .get("bucket") @@ -5199,6 +5648,16 @@ impl Operation for SRPeerBucketOpsHandler { .filter(|value| !value.is_empty()) .cloned() .ok_or_else(|| s3_error!(InvalidRequest, "operation is required"))?; + if state.pending_remove.is_some() + || (!state.enabled() + && !bootstrap_peer_bucket_operation_allowed( + &bucket, + &operation, + queries.get("bootstrapToken").map(String::as_str), + )) + { + return Err(s3_error!(InvalidRequest, "site replication is not enabled")); + } let Some(store) = object_store_from_req(&req) else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); @@ -5232,27 +5691,7 @@ impl Operation for SRPeerBucketOpsHandler { .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; - if let Some(runtime) = runtime_site_replication_targets().await? { - let replication_config = metadata_sys::get_replication_config(&bucket) - .await - .ok() - .map(|(config, _)| config); - ensure_site_replication_bucket_targets( - &bucket, - &runtime.state, - &runtime.local_peer, - replication_config.as_ref(), - &runtime.service_account_secret_key, - ) - .await?; - ensure_site_replication_bucket_replication_config( - &bucket, - &runtime.state, - &runtime.local_peer, - &runtime.service_account_secret_key, - ) - .await?; - } + ensure_site_replication_bucket_setup(&bucket).await?; } "delete-bucket" => { store @@ -5343,41 +5782,160 @@ impl Operation for SiteReplicationEditHandler { reject_site_replicator_on_public_admin(&cred)?; let ilm_expiry_override = sr_edit_ilm_expiry_override(&req.uri); let incoming: PeerInfo = read_site_replication_json(req, &cred.secret_key, true).await?; - let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; + let state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; let current_state = load_site_replication_state().await?; - let state = edit_state(current_state.clone(), incoming.clone(), ilm_expiry_override); + if current_state.pending_rotation.is_some() || current_state.pending_remove.is_some() { + return Err(s3_error!(InvalidRequest, "another site replication operation is pending")); + } + let persisted_pending = pending_endpoint_refresh(¤t_state); + let endpoint_refresh_requested = peer_endpoint_refresh_requested(¤t_state, &incoming); + if persisted_pending.is_some() && !endpoint_refresh_requested { + return Err(s3_error!(InvalidRequest, "an endpoint target refresh is already pending")); + } + let pending = endpoint_refresh_requested.then(|| { + persisted_pending.clone().unwrap_or_else(|| PendingEndpointRefresh { + id: Uuid::new_v4().to_string(), + peer: normalize_peer_info(incoming.clone()), + remote_peers: current_state.peers.clone(), + acked_deployment_ids: BTreeSet::new(), + }) + }); + let mut state = edit_state(current_state.clone(), incoming.clone(), ilm_expiry_override); - if !current_state.service_account_access_key.is_empty() { + if endpoint_refresh_requested && current_state.service_account_access_key.is_empty() { + return Err(s3_error!(InvalidRequest, "site replication service account is not configured")); + } + if current_state.service_account_access_key.is_empty() { + save_site_replication_state(&state).await?; + } else { let service_account_secret_key = site_replicator_service_account_secret(¤t_state.service_account_access_key).await?; let peers_to_send: Vec = if ilm_expiry_override.is_some() { state.peers.values().cloned().collect() + } else if let Some(pending) = pending.as_ref() { + vec![pending.peer.clone()] } else { vec![normalize_peer_info(incoming)] }; + let routing_peers = pending + .as_ref() + .map(|pending| &pending.remote_peers) + .unwrap_or(¤t_state.peers); + let local_deployment_id = current_deployment_id(); + let remote_targets = endpoint_refresh_remote_targets(routing_peers, pending.as_ref(), local_deployment_id.as_deref()); - for target in current_state.peers.values() { - let local_target = current_deployment_id() - .as_ref() - .is_some_and(|deployment_id| deployment_id == &target.deployment_id); - if local_target { - continue; - } - - for peer in &peers_to_send { - send_peer_admin_request_with_retry_event( + if endpoint_refresh_requested { + let pending = pending.clone().ok_or_else(|| { + S3Error::with_message(S3ErrorCode::InternalError, "endpoint refresh state is missing".to_string()) + })?; + let probes = futures::future::join_all(remote_targets.iter().map(|target| { + send_endpoint_refresh_admin_request_raw( target, - SITE_REPLICATION_PEER_EDIT_PATH, + &pending, + SITE_REPLICATION_PEER_EDIT_CAPABILITY_PATH, ¤t_state.service_account_access_key, &service_account_secret_key, - peer, + &(), ) - .await?; + })) + .await; + let mut legacy_deployment_ids = BTreeSet::new(); + for (target, probe) in remote_targets.iter().zip(probes) { + let (status, body) = probe.map_err(|err| { + S3Error::with_message( + S3ErrorCode::InternalError, + format!("probe endpoint target refresh on peer {} failed: {err}", target.endpoint), + ) + })?; + if endpoint_refresh_capability_supported(target, status, &body)? { + continue; + } else { + legacy_deployment_ids.insert(target.deployment_id.clone()); + } } + + let pending_id = pending.id.clone(); + let refresh_request = EndpointRefreshRequest { + id: pending.id.clone(), + peer: pending.peer.clone(), + }; + set_pending_endpoint_refresh(&mut state, pending.clone())?; + save_site_replication_state(&state).await?; + drop(state_guard); + let responses = futures::future::join_all(remote_targets.iter().map(|target| async { + if legacy_deployment_ids.contains(&target.deployment_id) { + refresh_legacy_peer_bucket_targets( + target, + &pending, + ¤t_state.service_account_access_key, + &service_account_secret_key, + ) + .await + } else { + let body = send_endpoint_refresh_admin_request( + target, + &pending, + SITE_REPLICATION_PEER_EDIT_REFRESH_PATH, + ¤t_state.service_account_access_key, + &service_account_secret_key, + &refresh_request, + ) + .await?; + parse_endpoint_refresh_status(target, &body) + } + })) + .await; + let mut acked_deployment_ids = BTreeSet::new(); + let mut refresh_error = None; + for (target, response) in remote_targets.iter().zip(responses) { + match response { + Ok(()) => { + acked_deployment_ids.insert(target.deployment_id.clone()); + } + Err(err) if refresh_error.is_none() => refresh_error = Some(err), + Err(_) => {} + } + } + + let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; + let mut state = load_site_replication_state().await?; + let Some(mut pending) = pending_endpoint_refresh(&state).filter(|pending| pending.id == pending_id) else { + return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); + }; + pending.acked_deployment_ids.extend(acked_deployment_ids); + set_pending_endpoint_refresh(&mut state, pending)?; + save_site_replication_state(&state).await?; + if let Some(err) = refresh_error { + return Err(err); + } + let service_account_secret_key = + site_replicator_service_account_secret(&state.service_account_access_key).await?; + drop(_state_guard); + refresh_bucket_targets_after_endpoint_edit(&pending_id, &service_account_secret_key).await?; + let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; + let mut state = load_site_replication_state().await?; + if pending_endpoint_refresh(&state).is_none_or(|pending| pending.id != pending_id) { + return Err(s3_error!(InvalidRequest, "endpoint target refresh state changed during update")); + } + clear_pending_endpoint_refresh(&mut state); + save_site_replication_state(&state).await?; + } else { + for target in remote_targets { + for peer in &peers_to_send { + send_peer_admin_request_with_retry_event( + target, + SITE_REPLICATION_PEER_EDIT_PATH, + ¤t_state.service_account_access_key, + &service_account_secret_key, + peer, + ) + .await?; + } + } + save_site_replication_state(&state).await?; } } - save_site_replication_state(&state).await?; json_response(&ReplicateEditStatus { success: true, status: SITE_REPL_EDIT_SUCCESS.to_string(), @@ -5387,26 +5945,123 @@ impl Operation for SiteReplicationEditHandler { } } +pub struct SRPeerEditCapabilitiesHandler {} + +#[async_trait::async_trait] +impl Operation for SRPeerEditCapabilitiesHandler { + async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { + validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; + json_response(&ReplicateEditStatus { + success: query_pairs(&req.uri) + .get("capability") + .is_some_and(|value| value == "endpoint-target-refresh"), + status: SITE_REPL_EDIT_SUCCESS.to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + ..Default::default() + }) + } +} + pub struct SRPeerEditHandler {} #[async_trait::async_trait] impl Operation for SRPeerEditHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { validate_site_replication_admin_request(&req, AdminAction::SiteReplicationOperationAction).await?; + let queries = query_pairs(&req.uri); let ilm_expiry_override = sr_edit_ilm_expiry_override(&req.uri); - let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; + let endpoint_refresh_requested = queries.get("refresh-targets").is_some_and(|value| value == "true"); + let state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; let state = load_site_replication_state().await?; + if endpoint_refresh_requested && (state.pending_rotation.is_some() || state.pending_remove.is_some()) { + return json_response(&ReplicateEditStatus { + success: false, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + err_detail: "another site replication operation is pending".to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }); + } let local_peer = current_local_peer(&req, &state); - let mut incoming: PeerInfo = read_site_replication_json(req, "", false).await?; + let (refresh_id, mut incoming) = if endpoint_refresh_requested { + let refresh: EndpointRefreshRequest = read_site_replication_json(req, "", false).await?; + (Some(refresh.id), refresh.peer) + } else { + (None, read_site_replication_json(req, "", false).await?) + }; if same_identity_endpoint(&incoming.endpoint, &local_peer.endpoint) { incoming.deployment_id = local_peer.deployment_id.clone(); if incoming.name.is_empty() { incoming.name = local_peer.name.clone(); } } - let state = + align_peer_edit_deployment_id(&state, &mut incoming); + if endpoint_refresh_requested + && pending_endpoint_refresh(&state).is_some_and(|pending| refresh_id.as_deref() != Some(&pending.id)) + { + return json_response(&ReplicateEditStatus { + success: false, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + err_detail: "another endpoint target refresh is pending".to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }); + } + if endpoint_refresh_requested + && (refresh_id.as_ref().is_none_or(String::is_empty) || !peer_endpoint_edit_requested(&state, &incoming)) + { + return json_response(&ReplicateEditStatus { + success: false, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + err_detail: "peer endpoint was not found".to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }); + } + + let mut state = sync_state_name_for_local_peer(update_peer(state, incoming.clone(), ilm_expiry_override), &local_peer, &incoming); + if endpoint_refresh_requested { + set_pending_endpoint_refresh( + &mut state, + PendingEndpointRefresh { + id: refresh_id.clone().unwrap_or_default(), + peer: incoming.clone(), + remote_peers: BTreeMap::new(), + acked_deployment_ids: BTreeSet::new(), + }, + )?; + } save_site_replication_state(&state).await?; + if endpoint_refresh_requested { + if state.service_account_access_key.is_empty() { + return json_response(&ReplicateEditStatus { + success: false, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + err_detail: "site replicator service account is not configured".to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }); + } + let service_account_secret_key = site_replicator_service_account_secret(&state.service_account_access_key).await?; + let pending_id = refresh_id.unwrap_or_default(); + drop(state_guard); + refresh_bucket_targets_after_endpoint_edit(&pending_id, &service_account_secret_key).await?; + let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; + let mut state = load_site_replication_state().await?; + if pending_endpoint_refresh(&state).is_none_or(|pending| pending.id != pending_id) { + return json_response(&ReplicateEditStatus { + success: false, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + err_detail: "endpoint target refresh state changed during update".to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + }); + } + clear_pending_endpoint_refresh(&mut state); + save_site_replication_state(&state).await?; + return json_response(&ReplicateEditStatus { + success: true, + status: SITE_REPL_EDIT_SUCCESS.to_string(), + api_version: Some(SITE_REPL_API_VERSION.to_string()), + ..Default::default() + }); + } Ok(empty_response(StatusCode::OK)) } } @@ -5418,8 +6073,16 @@ impl Operation for SRPeerRemoveHandler { async fn call(&self, req: S3Request, _params: Params<'_, '_>) -> S3Result> { validate_site_replication_admin_request(&req, AdminAction::SiteReplicationRemoveAction).await?; let remove_req: SRRemoveReq = read_site_replication_json(req, "", false).await?; + let _lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await; + let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.write().await; let _state_guard = SITE_REPLICATION_STATE_LOCK.lock().await; let current_state = load_site_replication_state().await?; + if pending_endpoint_refresh(¤t_state).is_some() { + return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); + } + if current_state.pending_rotation.is_some() { + return Err(s3_error!(InvalidRequest, "service account rotation is pending")); + } let removed_deployment_ids = removed_deployment_ids_for_remove_req(¤t_state, &remove_req); @@ -5620,6 +6283,12 @@ impl Operation for SRRotateServiceAccountHandler { if !state.enabled() { return Err(s3_error!(InvalidRequest, "site replication is not configured")); } + if pending_endpoint_refresh(&state).is_some() { + return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending")); + } + if state.pending_remove.is_some() { + return Err(s3_error!(InvalidRequest, "site replication remove is pending")); + } let local_peer = current_local_peer(&req, &state); let previous_access_key = state.service_account_access_key.clone(); @@ -5655,15 +6324,7 @@ impl Operation for SRRotateServiceAccountHandler { set_site_replicator_service_account_secret(&pending_rotation.parent, pending_rotation.new_secret_key.clone()).await?; - let rotation_state = SiteReplicationState { - service_account_access_key: pending_rotation.access_key.clone(), - service_account_parent: pending_rotation.parent.clone(), - peers: pending_rotation.peers.clone(), - updated_at: pending_rotation.updated_at, - ..Default::default() - }; - refresh_bucket_targets_after_service_account_rotation(&rotation_state, &local_peer, &pending_rotation.new_secret_key) - .await; + refresh_bucket_targets_after_service_account_rotation().await; let mut secret_candidates = pending_rotation.secret_candidates.clone(); if let Ok(current_secret) = site_replicator_service_account_secret(&pending_rotation.access_key).await { @@ -5736,10 +6397,16 @@ mod tests { use crate::admin::runtime_sources::{current_outbound_tls_generation, set_test_outbound_tls_generation}; use crate::admin::storage_api::runtime::Endpoint; use crate::admin::storage_api::runtime::{EndpointServerPools, Endpoints, PoolEndpoints}; + use axum::{Router, extract::State, routing::any}; use http::{HeaderMap, HeaderValue, Uri}; use rustfs_policy::policy::action::S3Action; use serial_test::serial; + use std::sync::{ + Arc, Mutex as StdMutex, + atomic::{AtomicBool, Ordering}, + }; use temp_env::with_var; + use tokio::net::TcpListener; fn peer(name: &str, endpoint: &str) -> PeerInfo { PeerInfo { @@ -5754,6 +6421,101 @@ mod tests { } } + #[derive(Clone, Default)] + struct LegacyPeerTestState { + requests: Arc>>, + minio_operation_supported: Arc, + } + + async fn legacy_peer_test_handler( + State(state): State, + method: Method, + uri: Uri, + ) -> (StatusCode, String) { + state + .requests + .lock() + .expect("legacy peer request log") + .push(format!("{method} {}", uri.path_and_query().map_or("/", |value| value.as_str()))); + match (method, uri.path()) { + (Method::GET, "/minio/admin/v3/site-replication/metainfo") => { + (StatusCode::OK, r#"{"Buckets":{"photos":{}}}"#.to_string()) + } + (Method::PUT, "/minio/admin/v3/site-replication/peer/edit") => (StatusCode::OK, String::new()), + (Method::PUT, "/minio/admin/v3/site-replication/peer/bucket-ops") + if !uri + .query() + .is_some_and(|query| query.contains("operation=ConfigureReplication")) + || state.minio_operation_supported.load(Ordering::Relaxed) => + { + (StatusCode::OK, String::new()) + } + _ => (StatusCode::NOT_FOUND, String::new()), + } + } + + #[tokio::test] + async fn legacy_endpoint_refresh_executes_peer_edit_and_bucket_repair() { + let listener = match TcpListener::bind("127.0.0.1:0").await { + Ok(listener) => listener, + Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("bind legacy peer test server: {err}"), + }; + let endpoint = format!("http://{}", listener.local_addr().expect("legacy peer test address")); + let state = LegacyPeerTestState::default(); + state.minio_operation_supported.store(true, Ordering::Relaxed); + let requests = state.requests.clone(); + let minio_operation_supported = state.minio_operation_supported.clone(); + let server = tokio::spawn(async move { + axum::serve(listener, Router::new().fallback(any(legacy_peer_test_handler)).with_state(state)) + .await + .expect("serve legacy peer test requests"); + }); + + let target = PeerInfo { + deployment_id: "remote".to_string(), + endpoint: endpoint.clone(), + ..Default::default() + }; + let pending = PendingEndpointRefresh { + id: "refresh-legacy".to_string(), + peer: PeerInfo { + deployment_id: "remote".to_string(), + endpoint, + ..Default::default() + }, + ..Default::default() + }; + + refresh_legacy_peer_bucket_targets(&target, &pending, "site-replicator-0", "test-secret") + .await + .expect("legacy peer endpoint refresh"); + assert_eq!( + *requests.lock().expect("legacy peer request log"), + vec![ + "PUT /minio/admin/v3/site-replication/peer/edit".to_string(), + "GET /minio/admin/v3/site-replication/metainfo?buckets=true".to_string(), + "PUT /minio/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=ConfigureReplication".to_string(), + ] + ); + + requests.lock().expect("legacy peer request log").clear(); + minio_operation_supported.store(false, Ordering::Relaxed); + refresh_legacy_peer_bucket_targets(&target, &pending, "site-replicator-0", "test-secret") + .await + .expect("legacy RustFS peer endpoint refresh"); + server.abort(); + assert_eq!( + *requests.lock().expect("legacy peer request log"), + vec![ + "PUT /minio/admin/v3/site-replication/peer/edit".to_string(), + "GET /minio/admin/v3/site-replication/metainfo?buckets=true".to_string(), + "PUT /minio/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=ConfigureReplication".to_string(), + "PUT /minio/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=configure-replication".to_string(), + ] + ); + } + #[tokio::test] async fn test_site_replicator_service_account_policy_allows_peer_and_object_replication() { let policy = site_replicator_service_account_policy().expect("site replicator policy should parse"); @@ -5955,6 +6717,97 @@ mod tests { assert!(!query_flag(&uri, "missing")); } + #[tokio::test] + #[serial] + async fn test_add_bootstrap_scope_only_allows_expected_bucket_setup_until_guard_drops() { + let token; + { + let lifecycle = SiteReplicationLifecycleGuard::acquire().await; + let guard = SiteReplicationAddInProgressGuard::start(lifecycle, HashSet::from(["legacy-bucket".to_string()])) + .expect("start site replication add guard"); + token = guard.token.to_string(); + assert!(bootstrap_peer_bucket_operation_allowed( + "new-bucket", + "make-with-versioning", + Some(&token) + )); + assert!(bootstrap_peer_bucket_operation_allowed( + "new-bucket", + "configure-replication", + Some(&token) + )); + assert!(bootstrap_peer_bucket_operation_allowed("legacy-bucket", "make-with-versioning", None)); + assert!(!bootstrap_peer_bucket_operation_allowed( + "unexpected-bucket", + "make-with-versioning", + None + )); + assert!(!bootstrap_peer_bucket_operation_allowed( + "legacy-bucket", + "force-delete-bucket", + Some(&token) + )); + assert!(!bootstrap_peer_bucket_operation_allowed( + "legacy-bucket", + "make-with-versioning", + Some(&Uuid::new_v4().to_string()) + )); + } + assert!(!bootstrap_peer_bucket_operation_allowed( + "new-bucket", + "make-with-versioning", + Some(&token) + )); + } + + #[test] + fn test_add_bootstrap_token_round_trips_from_join_to_bucket_operation() { + let token = Uuid::new_v4().to_string(); + let join_path = with_site_replication_bootstrap_token(SITE_REPLICATION_PEER_JOIN_PATH, &token); + let join_uri: Uri = join_path.parse().expect("parse peer join path"); + let received_token = site_replication_bootstrap_token(&join_uri).expect("peer join bootstrap token"); + let bucket_path = + with_site_replication_bootstrap_token(&bootstrap_bucket_op_path("photos", "configure-replication"), &received_token); + let bucket_uri: Uri = bucket_path.parse().expect("parse bucket operation path"); + let query = query_pairs(&bucket_uri); + + assert_eq!(query.get("bootstrapToken"), Some(&token)); + assert_eq!(query.get("operation").map(String::as_str), Some("configure-replication")); + } + + #[tokio::test] + #[serial] + async fn test_add_lifecycle_allows_callback_before_remove_writer() { + let lifecycle = SiteReplicationLifecycleGuard::acquire().await; + let add_guard = + SiteReplicationAddInProgressGuard::start(lifecycle, HashSet::new()).expect("start site replication add guard"); + let add_state = SITE_REPLICATION_STATE_LOCK.lock().await; + let (started_tx, started_rx) = tokio::sync::oneshot::channel(); + let (entered_tx, mut entered_rx) = tokio::sync::oneshot::channel(); + let remove = tokio::spawn(async move { + let _ = started_tx.send(()); + let _lifecycle = SiteReplicationLifecycleGuard::acquire().await; + let _bucket_op = SITE_REPLICATION_BUCKET_OP_LOCK.write().await; + let _state = SITE_REPLICATION_STATE_LOCK.lock().await; + let _ = entered_tx.send(()); + }); + started_rx.await.expect("remove task started"); + + let callback = tokio::time::timeout(Duration::from_millis(500), SITE_REPLICATION_BUCKET_OP_LOCK.read()) + .await + .expect("callback read lock should not wait behind remove"); + assert!(matches!(entered_rx.try_recv(), Err(tokio::sync::oneshot::error::TryRecvError::Empty))); + + drop(callback); + drop(add_state); + drop(add_guard); + tokio::time::timeout(Duration::from_millis(500), remove) + .await + .expect("remove should enter after add finishes") + .expect("remove task should finish"); + entered_rx.await.expect("remove entered lifecycle"); + } + #[test] fn test_merge_add_sites_propagates_replicate_ilm_expiry() { let state = merge_add_sites( @@ -6050,6 +6903,7 @@ mod tests { deployment_id: deployment_id.to_string(), enabled: false, bucket_count, + bucket_names: HashSet::new(), peer_deployment_ids: BTreeSet::new(), idp_settings: serde_json::json!({"provider": "same"}), } @@ -7114,6 +7968,176 @@ mod tests { .expect("site replication target should carry credentials"); assert_eq!(credentials.access_key, "site-replicator-0"); assert_eq!(credentials.secret_key, "runtime-iam-secret"); + + let regional_arn = "arn:rustfs:replication:eu-west-1:remote:photos"; + let config = ReplicationConfiguration { + role: String::new(), + rules: vec![build_site_replication_rule(regional_arn, 1, "site-repl-remote")], + }; + state.peers.get_mut("remote").expect("remote peer should exist").endpoint = "http://moved.example.com:9001".to_string(); + let targets = reconcile_site_replication_bucket_targets( + targets, + "photos", + &state, + &PeerInfo { + deployment_id: "local".to_string(), + ..peer("local", "https://local.example.com") + }, + Some(&config), + "runtime-iam-secret", + ) + .expect("reconcile moved peer target"); + + assert_eq!(targets.targets.len(), 1); + assert_eq!(targets.targets[0].endpoint, "moved.example.com:9001"); + assert_eq!(targets.targets[0].arn, regional_arn); + assert_eq!( + replication_target_arn_deployment_id("arn:minio:replication:eu-west-1:remote:photos").as_deref(), + Some("remote") + ); + + let retry = PeerInfo { + deployment_id: "remote".to_string(), + endpoint: "https://moved.example.com:9001".to_string(), + ..Default::default() + }; + assert!(peer_endpoint_edit_requested(&state, &retry)); + state.peers.get_mut("remote").expect("remote peer should exist").endpoint = retry.endpoint; + let targets = reconcile_site_replication_bucket_targets( + targets, + "photos", + &state, + &PeerInfo { + deployment_id: "local".to_string(), + ..peer("local", "https://local.example.com") + }, + Some(&config), + "runtime-iam-secret", + ) + .expect("reconcile secure peer target"); + + assert_eq!(targets.targets.len(), 1); + assert_eq!(targets.targets[0].endpoint, "moved.example.com:9001"); + assert_eq!(targets.targets[0].arn, regional_arn); + assert!(targets.targets[0].secure); + + let mut mismatched = PeerInfo { + deployment_id: "source-view-remote".to_string(), + name: "remote".to_string(), + endpoint: "https://moved.example.com:9001".to_string(), + ..Default::default() + }; + align_peer_edit_deployment_id(&state, &mut mismatched); + assert_eq!(mismatched.deployment_id, "remote"); + + let retry_peer = PeerInfo { + deployment_id: "remote".to_string(), + endpoint: "https://moved.example.com:9001".to_string(), + ..Default::default() + }; + assert!(!peer_endpoint_refresh_requested(&state, &retry_peer)); + + let mut ambiguous_state = state.clone(); + ambiguous_state.peers.insert( + "remote-duplicate".to_string(), + PeerInfo { + deployment_id: "remote-duplicate".to_string(), + name: "remote".to_string(), + endpoint: "https://duplicate.example.com:9001".to_string(), + ..Default::default() + }, + ); + let mut ambiguous = mismatched.clone(); + ambiguous.deployment_id = "source-view-remote".to_string(); + align_peer_edit_deployment_id(&ambiguous_state, &mut ambiguous); + assert_eq!(ambiguous.deployment_id, "source-view-remote"); + + let remote_peers = state.peers.clone(); + set_pending_endpoint_refresh( + &mut state, + PendingEndpointRefresh { + id: "refresh-1".to_string(), + peer: retry_peer.clone(), + remote_peers, + acked_deployment_ids: BTreeSet::new(), + }, + ) + .expect("set pending endpoint refresh"); + assert!(peer_endpoint_refresh_requested(&state, &retry_peer)); + state.pending_endpoint_refresh = None; + assert_eq!(pending_endpoint_refresh(&state).map(|pending| pending.id), Some("refresh-1".to_string())); + clear_pending_endpoint_refresh(&mut state); + assert!(pending_endpoint_refresh(&state).is_none()); + assert!(!peer_endpoint_refresh_requested(&state, &retry_peer)); + assert!(parse_endpoint_refresh_status(&mismatched, b"").is_err()); + assert!(parse_endpoint_refresh_status(&mismatched, br#"{"success":false,"errorDetail":"refresh failed"}"#).is_err()); + assert!(parse_endpoint_refresh_status(&mismatched, br#"{"success":true}"#).is_ok()); + assert!( + endpoint_refresh_capability_supported(&mismatched, StatusCode::OK, br#"{"success":true}"#) + .expect("current peer capability response") + ); + assert!(!endpoint_refresh_capability_supported(&mismatched, StatusCode::OK, b"").expect("legacy empty response")); + assert!( + !endpoint_refresh_capability_supported(&mismatched, StatusCode::BAD_REQUEST, b"unsupported") + .expect("legacy bad request response") + ); + assert!(endpoint_refresh_capability_supported(&mismatched, StatusCode::UNAUTHORIZED, b"denied").is_err()); + + let old_target = PeerInfo { + deployment_id: "remote".to_string(), + endpoint: "http://old.example.com:9000".to_string(), + ..Default::default() + }; + let pending = PendingEndpointRefresh { + id: "refresh-2".to_string(), + peer: PeerInfo { + deployment_id: "remote".to_string(), + endpoint: "https://new.example.com:9001".to_string(), + ..Default::default() + }, + ..Default::default() + }; + assert_eq!( + endpoint_refresh_route_endpoints(&old_target, &pending), + vec![ + "http://old.example.com:9000".to_string(), + "https://new.example.com:9001".to_string() + ] + ); + let routing_peers = BTreeMap::from([ + ( + "local".to_string(), + PeerInfo { + deployment_id: "local".to_string(), + ..Default::default() + }, + ), + ("remote".to_string(), old_target), + ]); + let mut acked_pending = pending.clone(); + acked_pending.acked_deployment_ids.insert("remote".to_string()); + assert!(endpoint_refresh_remote_targets(&routing_peers, Some(&acked_pending), Some("local")).is_empty()); + + let mut request = serde_json::to_value(EndpointRefreshRequest { + id: pending.id, + peer: pending.peer, + }) + .expect("serialize endpoint refresh request"); + request + .as_object_mut() + .expect("endpoint refresh request object") + .insert("unexpected".to_string(), Value::Bool(true)); + assert!(serde_json::from_value::(request).is_err()); + assert_eq!( + peer_bucket_names_from_metainfo("https://minio.example.com", br#"{"Buckets":{"archive":{},"photos":{}}}"#) + .expect("MinIO metainfo bucket inventory"), + vec!["archive".to_string(), "photos".to_string()] + ); + assert_eq!( + peer_bucket_names_from_metainfo("https://rustfs.example.com", br#"{"buckets":{"photos":{}}}"#) + .expect("RustFS metainfo bucket inventory"), + vec!["photos".to_string()] + ); } #[test] diff --git a/rustfs/src/admin/route_policy.rs b/rustfs/src/admin/route_policy.rs index 525bd8573..b413742b6 100644 --- a/rustfs/src/admin/route_policy.rs +++ b/rustfs/src/admin/route_policy.rs @@ -598,6 +598,12 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[ SITE_REPLICATION_ADD, RouteRiskLevel::High, ), + admin( + HttpMethod::Put, + "/rustfs/admin/v3/site-replication/peer/edit-capabilities", + SITE_REPLICATION_OPERATION, + RouteRiskLevel::High, + ), admin( HttpMethod::Put, "/rustfs/admin/v3/site-replication/peer/edit", diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index e6b08cbcf..18f3ffd36 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -272,6 +272,7 @@ fn expected_admin_route_matrix() -> Vec { admin_route(Method::PUT, "/v3/site-replication/peer/bucket-meta"), admin_route(Method::GET, "/v3/site-replication/peer/idp-settings"), admin_route(Method::PUT, "/v3/site-replication/edit"), + admin_route(Method::PUT, "/v3/site-replication/peer/edit-capabilities"), admin_route(Method::PUT, "/v3/site-replication/peer/edit"), admin_route(Method::PUT, "/v3/site-replication/peer/remove"), admin_route(Method::PUT, "/v3/site-replication/resync/op"), @@ -1192,6 +1193,7 @@ fn test_register_routes_cover_representative_admin_paths() { assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/peer/bucket-meta")); assert_route(&router, Method::GET, &admin_path("/v3/site-replication/peer/idp-settings")); assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/edit")); + assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/peer/edit-capabilities")); assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/peer/edit")); assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/peer/remove")); assert_route(&router, Method::PUT, &admin_path("/v3/site-replication/resync/op")); diff --git a/rustfs/src/admin/router.rs b/rustfs/src/admin/router.rs index 3d0d53d7e..aa137d677 100644 --- a/rustfs/src/admin/router.rs +++ b/rustfs/src/admin/router.rs @@ -40,6 +40,7 @@ use crate::license::license_check; use crate::server::{ ADMIN_PREFIX, HEALTH_PREFIX, HEALTH_READY_PATH, MINIO_ADMIN_PREFIX, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH, is_admin_path, }; +use crate::storage::storage_api::lock_bucket_targets_metadata; use aws_sdk_s3::primitives::ByteStream as AwsByteStream; use bytes::Bytes; use futures::{Stream, StreamExt}; @@ -2136,6 +2137,7 @@ async fn target_client_object_lock_enabled(bucket: &str, target: &BucketTarget) } async fn start_replication_resync(bucket: &str, reset: &ReplicationResetStartRequest) -> S3Result { + let targets_guard = lock_bucket_targets_metadata(bucket).await; let (config, _) = metadata_sys::get_replication_config(bucket).await.map_err(ApiError::from)?; let resolved_arn = resolve_replication_reset_target_arn(&config, &reset.arn)?; let mut resolved_reset = reset.clone(); @@ -2149,6 +2151,7 @@ async fn start_replication_resync(bucket: &str, reset: &ReplicationResetStartReq .await .map_err(ApiError::from)?; BucketTargetSys::get().update_all_targets(bucket, Some(&targets)).await; + drop(targets_guard); let Some(pool) = current_replication_pool_handle() else { return Err(s3_error!(InternalError, "replication pool is not initialized")); diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index bcc69b0da..4eabbb720 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -71,6 +71,7 @@ use crate::app::runtime_sources::{ use crate::auth::get_condition_values_with_client_info; use crate::error::ApiError; use crate::server::RemoteAddr; +use crate::storage::storage_api::lock_bucket_targets_metadata; use futures::StreamExt; use http::StatusCode; use metrics::counter; @@ -1455,6 +1456,7 @@ impl DefaultBucketUsecase { .get_bucket_info(&bucket, &BucketOptions::default()) .await .map_err(ApiError::from)?; + let targets_guard = lock_bucket_targets_metadata(&bucket).await; let replication_config = match metadata_sys::get_replication_config(&bucket).await { Ok((config, _)) => Some(config), Err(StorageError::ConfigNotFound) => None, @@ -1477,6 +1479,7 @@ impl DefaultBucketUsecase { } return Err(err); } + drop(targets_guard); notify_bucket_metadata_reload(bucket.clone(), "delete bucket replication", request_context); @@ -2296,11 +2299,13 @@ impl DefaultBucketUsecase { .await .map_err(ApiError::from)?; + let targets_guard = lock_bucket_targets_metadata(&bucket).await; validate_bucket_replication_update(&bucket, &replication_configuration).await?; let data = serialize_config(&replication_configuration)?; metadata_sys::update(&bucket, BUCKET_REPLICATION_CONFIG, data) .await .map_err(ApiError::from)?; + drop(targets_guard); notify_bucket_metadata_reload(bucket.clone(), "put bucket replication", request_context); diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index efeef4b8c..d573af5fd 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -14,11 +14,34 @@ //! Storage owner-local boundary for ECStore facade and storage contract symbols. -use std::sync::Arc; +use std::collections::hash_map::DefaultHasher; +use std::hash::{Hash, Hasher}; +use std::sync::{Arc, LazyLock}; use rustfs_storage_api as storage_contracts; +use tokio::sync::{Mutex, OwnedMutexGuard}; use tokio_util::sync::CancellationToken; +const BUCKET_TARGETS_METADATA_LOCK_SHARDS: usize = 256; +static BUCKET_TARGETS_METADATA_LOCKS: LazyLock>>> = LazyLock::new(|| { + (0..BUCKET_TARGETS_METADATA_LOCK_SHARDS) + .map(|_| Arc::new(Mutex::new(()))) + .collect() +}); + +pub(crate) async fn lock_bucket_targets_metadata(bucket: &str) -> OwnedMutexGuard<()> { + BUCKET_TARGETS_METADATA_LOCKS[bucket_targets_metadata_lock_shard(bucket)] + .clone() + .lock_owned() + .await +} + +fn bucket_targets_metadata_lock_shard(bucket: &str) -> usize { + let mut hasher = DefaultHasher::new(); + bucket.hash(&mut hasher); + hasher.finish() as usize % BUCKET_TARGETS_METADATA_LOCK_SHARDS +} + pub(crate) mod contract { pub(crate) mod admin { pub(crate) use super::super::storage_contracts::StorageAdminApi; @@ -1435,3 +1458,36 @@ pub(crate) async fn store_compression_total_in_backend() { pub(crate) async fn init_compression_total_memory_from_backend(store: Arc) { ecstore_data_usage::init_compression_total_memory_from_backend(store).await } + +#[cfg(test)] +mod tests { + use super::{bucket_targets_metadata_lock_shard, lock_bucket_targets_metadata}; + use std::time::Duration; + + #[tokio::test] + async fn bucket_target_metadata_locks_serialize_only_matching_shards() { + let bucket = "bucket-target-lock"; + let other = (0..1000) + .map(|index| format!("other-bucket-{index}")) + .find(|candidate| bucket_targets_metadata_lock_shard(candidate) != bucket_targets_metadata_lock_shard(bucket)) + .expect("find bucket on another lock shard"); + + let guard = lock_bucket_targets_metadata(bucket).await; + assert!( + tokio::time::timeout(Duration::from_millis(20), lock_bucket_targets_metadata(bucket)) + .await + .is_err() + ); + assert!( + tokio::time::timeout(Duration::from_secs(1), lock_bucket_targets_metadata(&other)) + .await + .is_ok() + ); + drop(guard); + assert!( + tokio::time::timeout(Duration::from_secs(1), lock_bucket_targets_metadata(bucket)) + .await + .is_ok() + ); + } +}