From a8aaadb886ccb1a63a159b050679b46440f71e3a Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sun, 6 Sep 2026 12:51:44 +0800 Subject: [PATCH] fix(replication): preserve concurrent remote target changes (#7262) * test(replication): cover concurrent remote target writes * fix(replication): merge remote target writes under transactions * test(replication): preserve concurrent repairs during target updates * fix(replication): retain removal guards after request cancellation * test(replication): enable loopback in ordinary target fixtures * test(replication): isolate remote target mutation scenarios * refactor(admin): keep target writes within existing boundaries * test(replication): assert removed target cache state * test(replication): inspect target cache before client refresh --- rustfs/src/admin/handlers/replication.rs | 786 ++++++++++++++++++++--- rustfs/src/admin/storage_api.rs | 4 +- 2 files changed, 716 insertions(+), 74 deletions(-) diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index 40205b703..9c0e485c0 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -19,7 +19,6 @@ use crate::admin::runtime_sources::{ AppContext, app_context_from_req, current_notification_system_for_context, current_replication_pool_handle, current_replication_stats_handle_for_context, current_runtime_port, object_store_from_req, }; -use crate::admin::storage_api::AdminVersioningConfigExt as _; use crate::admin::storage_api::bucket::metadata::BUCKET_TARGETS_FILE; use crate::admin::storage_api::bucket::metadata_sys; use crate::admin::storage_api::bucket::metadata_sys::get_replication_config; @@ -28,9 +27,10 @@ use crate::admin::storage_api::bucket::replication::{BucketStats, ReplicationSta #[cfg(test)] use crate::admin::storage_api::bucket::replication::{REMOTE_TARGET_READ_ONLY_HISTORICAL_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS}; use crate::admin::storage_api::bucket::target::{ - BucketTarget, BucketTargetType, Credentials as TargetCredentials, LatencyStat, duration_from_secs_or_nanos, + ARN, BucketTarget, BucketTargetType, Credentials as TargetCredentials, LatencyStat, duration_from_secs_or_nanos, }; use crate::admin::storage_api::bucket::target_sys::{BucketTargetError, BucketTargetSys}; +use crate::admin::storage_api::bucket::{AdminReplicationConfigExt as _, AdminVersioningConfigExt as _}; use crate::admin::storage_api::contract::bucket::{BucketOperations, BucketOptions}; use crate::admin::storage_api::contract::list::ListOperations as _; use crate::admin::storage_api::error::StorageError; @@ -50,6 +50,7 @@ use s3s::header::CONTENT_TYPE; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use serde::{Deserialize, Serialize}; use std::collections::{HashMap, HashSet}; +use std::str::FromStr as _; use std::sync::Arc; use std::time::Duration; use time::OffsetDateTime; @@ -103,11 +104,34 @@ fn parse_remote_target_write_modes(uri: &http::Uri) -> S3Result<(bool, bool)> { Ok((update, replace_unreadable)) } -/// Repair decisions use the disk snapshot protected by the metadata transaction, -/// never the stale target cache retained after an unreadable configuration load. -async fn persist_remote_target_repair(bucket: &str, mut target: BucketTarget, incarnation: uuid::Uuid) -> S3Result { +enum RemoteTargetWrite { + Create { replace_unreadable: bool }, + Update { expected: Vec }, +} + +async fn persist_remote_target_repair(bucket: &str, target: BucketTarget, incarnation: uuid::Uuid) -> S3Result { + persist_remote_target_write( + bucket, + target, + incarnation, + RemoteTargetWrite::Create { + replace_unreadable: true, + }, + ) + .await +} + +/// Merge only into the disk snapshot protected by the metadata transaction. +/// An update may reuse remote validation only while that target is unchanged. +async fn persist_remote_target_write( + bucket: &str, + mut target: BucketTarget, + incarnation: uuid::Uuid, + mode: RemoteTargetWrite, +) -> S3Result { let mut discarded_unreadable = false; let mut target_error = None; + let mut conflict = false; let updated = metadata_sys::update_config_with(bucket, BUCKET_TARGETS_FILE, |metadata| { if metadata.bucket_incarnation_id != incarnation { return Err(StorageError::BucketNotFound(bucket.to_string())); @@ -118,12 +142,44 @@ async fn persist_remote_target_repair(bucket: &str, mut target: BucketTarget, in target_error = Some(BucketTargetError::BucketReplicationSourceNotVersioned { bucket: bucket.to_string(), }); - return Err(StorageError::other("source bucket versioning changed before target repair")); + return Err(StorageError::other("source bucket versioning changed before target write")); } discarded_unreadable = metadata.bucket_targets_unreadable(); + if discarded_unreadable + && !matches!( + &mode, + RemoteTargetWrite::Create { + replace_unreadable: true + } + ) + { + target_error = Some(BucketTargetError::BucketRemoteTargetsUnreadable { + bucket: bucket.to_string(), + }); + return Err(StorageError::other("persisted remote targets cannot be decoded")); + } let mut targets = metadata.bucket_target_config.clone().unwrap_or_default(); - let (arn, exists) = BucketTargetSys::remote_arn_for_targets(&targets.targets, &target, &target.deployment_id); - target.arn = arn; + let (update, exists) = match &mode { + RemoteTargetWrite::Create { .. } => { + let (arn, exists) = BucketTargetSys::remote_arn_for_targets(&targets.targets, &target, &target.deployment_id); + target.arn = arn; + (false, exists) + } + RemoteTargetWrite::Update { expected } => { + let current = targets.targets.iter().find(|current| current.arn == target.arn); + if current + .map(serde_json::to_vec) + .transpose() + .map_err(StorageError::other)? + .as_ref() + != Some(expected) + { + conflict = true; + return Err(StorageError::other("remote target changed during validation")); + } + (true, false) + } + }; if target.arn.is_empty() { target_error = Some(BucketTargetError::BucketRemoteArnInvalid { bucket: bucket.to_string(), @@ -131,7 +187,7 @@ async fn persist_remote_target_repair(bucket: &str, mut target: BucketTarget, in return Err(StorageError::other("remote target ARN is empty")); } if !exists { - BucketTargetSys::upsert_target_entry(&mut targets.targets, &target, false).map_err(|error| { + BucketTargetSys::upsert_target_entry(&mut targets.targets, &target, update).map_err(|error| { target_error = Some(error); StorageError::other("remote target merge failed") })?; @@ -139,6 +195,12 @@ async fn persist_remote_target_repair(bucket: &str, mut target: BucketTarget, in serde_json::to_vec(&targets).map_err(StorageError::other) }) .await; + if conflict { + return Err(S3Error::with_message( + S3ErrorCode::OperationAborted, + "remote target changed during validation; retry the request", + )); + } if let Some(error) = target_error { return Err(map_bucket_target_error(error)); } @@ -741,32 +803,60 @@ impl Operation for SetRemoteTargetHandler { return Ok(S3Response::new((StatusCode::OK, Body::from(arn_str)))); } - if !update { - let (arn, exist) = bucket_target_sys - .get_remote_arn(bucket, Some(&remote_target), remote_target.deployment_id.as_str()) - .await; - remote_target.arn = arn.clone(); - if exist && !arn.is_empty() { - let arn_str = serde_json::to_string(&arn).unwrap_or_default(); - - warn!("return exists, arn: {}", arn_str); - // MinIO-compatible clients encrypt the request payload for this endpoint, - // but they parse the success response directly as plain JSON string ARN. - return Ok(S3Response::new((StatusCode::OK, Body::from(arn_str)))); - } - } - - 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 - .get_remote_bucket_target_by_arn(bucket, &remote_target.arn) + let incarnation = metadata_sys::capture_bucket_metadata_incarnation(bucket) + .await + .map_err(ApiError::from)?; + // Preserve create idempotency even when an existing destination is + // offline, but decide it from disk under the same transaction as writers. + let metadata = { + let transaction_guard = metadata_sys::acquire_bucket_metadata_transaction_lock_for_incarnation(bucket, incarnation) .await + .map_err(ApiError::from)?; + let metadata = metadata_sys::get_config_from_disk(bucket).await.map_err(ApiError::from)?; + transaction_guard.checked_bucket_incarnation().map_err(ApiError::from)?; + if metadata.bucket_incarnation_id != incarnation { + return Err(ApiError::from(StorageError::BucketNotFound(bucket.to_string())).into()); + } + if metadata.bucket_targets_unreadable() { + return Err(map_bucket_target_error(BucketTargetError::BucketRemoteTargetsUnreadable { + bucket: bucket.to_string(), + })); + } + if !update { + let targets = metadata + .bucket_target_config + .as_ref() + .map(|targets| targets.targets.as_slice()) + .unwrap_or_default(); + let (arn, exists) = + BucketTargetSys::remote_arn_for_targets(targets, &remote_target, &remote_target.deployment_id); + if exists && !arn.is_empty() { + let body = serde_json::to_string(&arn) + .map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize target ARN"))?; + return Ok(S3Response::new((StatusCode::OK, Body::from(body)))); + } + } + metadata + }; + let mut mode = RemoteTargetWrite::Create { + replace_unreadable: false, + }; + if update { + if remote_target.arn.is_empty() { + return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, "ARN is empty")); + } + let Some(mut target) = metadata + .bucket_target_config + .unwrap_or_default() + .targets + .into_iter() + .find(|target| target.arn == remote_target.arn) else { - return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, "Target not found".to_string())); + return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, "Target not found")); + }; + mode = RemoteTargetWrite::Update { + expected: serde_json::to_vec(&target) + .map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize target"))?, }; // Overlay only the requested field groups onto the stored target @@ -828,26 +918,16 @@ impl Operation for SetRemoteTargetHandler { remote_target = target; } - let arn = remote_target.arn.clone(); - - let targets = bucket_target_sys - .set_target(bucket, &remote_target, update) + // Neither the process-local target lock nor the cluster metadata + // transaction may be held while the remote endpoint is responding. + bucket_target_sys + .validate_target(bucket, &remote_target) .await .map_err(map_bucket_target_error)?; - let json_targets = serde_json::to_vec(&targets).map_err(|e| { - error!("Serialization error: {}", e); - S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize targets".to_string()) - })?; - - metadata_sys::update(bucket, BUCKET_TARGETS_FILE, json_targets) - .await - .map_err(|e| { - error!("Failed to update bucket targets: {}", e); - S3Error::with_message(S3ErrorCode::InternalError, format!("Failed to update bucket targets: {e}")) - })?; - bucket_target_sys.update_all_targets(bucket, Some(&targets)).await; - - let arn_str = serde_json::to_string(&arn).unwrap_or_default(); + let _targets_guard = lock_bucket_targets_metadata(bucket).await; + let arn = persist_remote_target_write(bucket, remote_target, incarnation, mode).await?; + let arn_str = serde_json::to_string(&arn) + .map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize target ARN"))?; // MinIO-compatible clients encrypt the request payload for this endpoint, // but they parse the success response directly as plain JSON string ARN. @@ -952,25 +1032,74 @@ impl Operation for RemoveRemoteTargetHandler { .await .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)?; - - cancel_active_resync_intent(bucket, arn_str).await?; - - let json_targets = serde_json::to_vec(&targets).map_err(|e| { - error!("Serialization error: {}", e); - S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize targets".to_string()) + let arn = ARN::from_str(arn_str).map_err(|_| { + map_bucket_target_error(BucketTargetError::BucketRemoteArnInvalid { + bucket: bucket.to_string(), + }) })?; - - metadata_sys::update(bucket, BUCKET_TARGETS_FILE, json_targets) + let incarnation = metadata_sys::capture_bucket_metadata_incarnation(bucket) .await - .map_err(|e| { - error!("Failed to update bucket targets: {}", e); - S3Error::with_message(S3ErrorCode::InternalError, format!("Failed to update bucket targets: {e}")) - })?; - sys.update_all_targets(bucket, Some(&targets)).await; + .map_err(ApiError::from)?; + let targets_guard = lock_bucket_targets_metadata(bucket).await; + // Match resync start: local targets -> lifecycle -> metadata transaction + // -> resync admission -> resync status CAS. Keep the transaction through + // cancellation and target persistence so another start cannot interleave. + let transaction_guard = metadata_sys::acquire_bucket_metadata_transaction_lock_for_incarnation(bucket, incarnation) + .await + .map_err(ApiError::from)?; + let metadata = metadata_sys::get_config_from_disk(bucket).await.map_err(ApiError::from)?; + if metadata.bucket_targets_unreadable() { + return Err(map_bucket_target_error(BucketTargetError::BucketRemoteTargetsUnreadable { + bucket: bucket.to_string(), + })); + } + if arn.arn_type == BucketTargetType::ReplicationService { + if !metadata.replication_config_xml.is_empty() && metadata.replication_config.is_none() { + return Err(S3Error::with_message( + S3ErrorCode::InternalError, + "persisted replication rules cannot be decoded", + )); + } + if metadata.replication_config.as_ref().is_some_and(|config| { + config + .filter_all_replication_target_arns() + .iter() + .any(|target| target == arn_str) + }) { + return Err(map_bucket_target_error(BucketTargetError::BucketRemoteRemoveDisallowed { + bucket: bucket.to_string(), + })); + } + } + let mut targets = metadata.bucket_target_config.unwrap_or_default(); + let previous_len = targets.targets.len(); + targets.targets.retain(|target| target.arn != *arn_str); + if targets.targets.len() == previous_len { + return Err(map_bucket_target_error(BucketTargetError::BucketRemoteTargetNotFound { + bucket: bucket.to_string(), + })); + } + let json_targets = serde_json::to_vec(&targets) + .map_err(|_| S3Error::with_message(S3ErrorCode::InternalError, "Failed to serialize targets"))?; + let bucket = bucket.clone(); + let arn = arn_str.clone(); + // The pool cancellation owns a detached task. Both outer guards must + // outlive it, even when the HTTP future is dropped while awaiting it. + tokio::spawn(async move { + let _targets_guard = targets_guard; + #[cfg(test)] + target_repair_tests::pause_target_removal(&bucket).await; + transaction_guard.checked_bucket_incarnation().map_err(ApiError::from)?; + cancel_active_resync_intent(&bucket, &arn).await?; + metadata_sys::update_bucket_targets_under_transaction_lock(&transaction_guard, &bucket, json_targets) + .await + .map_err(ApiError::from)?; + Ok::<(), S3Error>(()) + }) + .await + .map_err(|error| { + S3Error::with_message(S3ErrorCode::InternalError, format!("remote target removal task failed: {error}")) + })??; Ok(S3Response::new((StatusCode::NO_CONTENT, Body::from("".to_string())))) } @@ -2804,6 +2933,7 @@ mod target_repair_tests { use rustfs_iam::store::{Store as _, object::IAM_CONFIG_PREFIX}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpListener; + use tokio::sync::oneshot; use tracing::instrument::WithSubscriber as _; const ACCESS_KEY: &str = "TARGETREPAIRROOT"; @@ -2819,6 +2949,29 @@ mod target_repair_tests { struct RemoteTargetServer { endpoint: String, task: tokio::task::JoinHandle<()>, + next_request: Arc>>, + } + + struct RequestPause { + reached: oneshot::Sender<()>, + release: oneshot::Receiver<()>, + } + + static REMOVAL_PAUSE: std::sync::Mutex> = std::sync::Mutex::new(None); + + pub(super) async fn pause_target_removal(bucket: &str) { + let pause = { + let mut pending = REMOVAL_PAUSE.lock().expect("removal pause lock"); + if pending.as_ref().is_some_and(|(expected, _)| expected == bucket) { + pending.take().map(|(_, pause)| pause) + } else { + None + } + }; + if let Some(pause) = pause { + pause.reached.send(()).expect("observe removal before cancellation"); + pause.release.await.expect("release removal cancellation"); + } } impl Drop for RemoteTargetServer { @@ -2831,6 +2984,8 @@ mod target_repair_tests { async fn start() -> Self { let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind remote target"); let endpoint = listener.local_addr().expect("remote target address").to_string(); + let next_request = Arc::new(std::sync::Mutex::new(None::)); + let request_pause = Arc::clone(&next_request); let task = tokio::spawn(async move { loop { let (mut socket, _) = listener.accept().await.expect("accept remote target request"); @@ -2845,6 +3000,11 @@ mod target_repair_tests { } let head = String::from_utf8_lossy(&request); let first_line = head.lines().next().unwrap_or_default(); + let pause = request_pause.lock().expect("request pause lock").take(); + if let Some(pause) = pause { + pause.reached.send(()).expect("notify request arrival"); + pause.release.await.expect("release remote response"); + } let body = if first_line.starts_with("HEAD /") { "" } else if first_line.starts_with("GET /") && first_line.contains("versioning") { @@ -2862,7 +3022,27 @@ mod target_repair_tests { .expect("reply to remote target request"); } }); - Self { endpoint, task } + Self { + endpoint, + task, + next_request, + } + } + + fn pause_next_request(&self) -> (oneshot::Receiver<()>, oneshot::Sender<()>) { + let (reached, observed) = oneshot::channel(); + let (release, resume) = oneshot::channel(); + assert!( + self.next_request + .lock() + .expect("request pause lock") + .replace(RequestPause { + reached, + release: resume, + }) + .is_none() + ); + (observed, release) } fn target(&self) -> BucketTarget { @@ -2914,10 +3094,10 @@ mod target_repair_tests { } fn request(method: Method, query: &str, body: Vec) -> S3Request { - let operation = if method == Method::GET { - "list-remote-targets" - } else { - "set-remote-target" + let operation = match method { + Method::GET => "list-remote-targets", + Method::DELETE => "remove-remote-target", + _ => "set-remote-target", }; S3Request { input: Body::from(body), @@ -3272,4 +3452,464 @@ mod target_repair_tests { }) .await; } + + async fn remove(arn: &str) -> S3Result<()> { + let query = url::form_urlencoded::Serializer::new(String::new()) + .append_pair("arn", arn) + .finish(); + let response = RemoveRemoteTargetHandler {} + .call(request(Method::DELETE, &query, Vec::new()), Params::new()) + .await?; + assert_eq!(response.output.0, StatusCode::NO_CONTENT); + Ok(()) + } + + async fn persisted_targets() -> BucketTargets { + metadata_sys::get_config_from_disk(BUCKET) + .await + .expect("read persisted targets") + .bucket_target_config + .expect("decode persisted targets") + } + + async fn assert_published_targets(targets: &BucketTargets) { + let cached = match BucketTargetSys::get().list_bucket_targets(BUCKET).await { + Ok(cached) => cached, + Err(BucketTargetError::BucketRemoteTargetNotFound { bucket }) if bucket == BUCKET && targets.targets.is_empty() => { + BucketTargets::default() + } + Err(error) => panic!("read published targets: {error}"), + }; + assert_eq!( + serde_json::to_value(cached).expect("encode cache"), + serde_json::to_value(targets).expect("encode disk") + ); + } + + const ORDINARY_TARGET_ENV: [(&str, Option<&str>); 3] = [ + ("NO_PROXY", Some("*")), + ("no_proxy", Some("*")), + ("RUSTFS_REPLICATION_ALLOW_LOOPBACK_TARGET", Some("true")), + ]; + + async fn assert_target_write_preserves_uncached_repair(operation: &str) { + temp_env::async_with_vars( + ORDINARY_TARGET_ENV, + Box::pin(async move { + let (_temp, _env) = test_env().await; + let first = RemoteTargetServer::start().await; + let second = RemoteTargetServer::start().await; + let third = RemoteTargetServer::start().await; + let first_arn = repair(&first.target(), "").await.expect("create initial target"); + let stale = persisted_targets().await; + let second_arn = repair(&second.target(), "replace-unreadable=true") + .await + .expect("commit peer repair"); + let repaired = persisted_targets() + .await + .targets + .into_iter() + .find(|target| target.arn == second_arn) + .expect("persisted peer repair"); + // Model a node whose target cache predates the committed repair. + BucketTargetSys::get().update_all_targets(BUCKET, Some(&stale)).await; + let created = match operation { + "create" => Some(repair(&third.target(), "").await.expect("merge target create")), + "update" => { + let mut changed = stale.targets[0].clone(); + changed.replication_sync = true; + assert_eq!(repair(&changed, "update=true&sync=true").await.expect("merge target update"), first_arn); + None + } + "remove" => { + remove(&first_arn).await.expect("merge target removal"); + None + } + _ => unreachable!(), + }; + let targets = persisted_targets().await; + assert_eq!( + targets.targets.len(), + match operation { + "create" => 3, + "update" => 2, + _ => 1, + }, + "{operation}" + ); + let kept = targets + .targets + .iter() + .find(|target| target.arn == second_arn) + .expect("unrelated repair must survive"); + assert_eq!( + serde_json::to_value(kept).expect("encode kept target"), + serde_json::to_value(repaired).expect("encode repair") + ); + if let Some(created) = created { + assert!( + targets + .targets + .iter() + .any(|target| target.arn == created && target.endpoint == third.endpoint) + ); + } + if operation == "update" { + let updated = targets + .targets + .iter() + .find(|target| target.arn == first_arn) + .expect("updated target"); + assert!(updated.replication_sync); + assert_eq!(updated.credentials.as_ref().expect("retained credentials").secret_key, "remote-secret"); + } + assert_published_targets(&targets).await; + }), + ) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_create_preserves_repairs_missing_from_the_cache() { + assert_target_write_preserves_uncached_repair("create").await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_update_preserves_repairs_missing_from_the_cache() { + assert_target_write_preserves_uncached_repair("update").await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_remove_preserves_repairs_missing_from_the_cache() { + assert_target_write_preserves_uncached_repair("remove").await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_writes_refuse_unreadable_disk_even_with_a_readable_cache() { + temp_env::async_with_vars(ORDINARY_TARGET_ENV, async { + let (_temp, env) = test_env().await; + let server = RemoteTargetServer::start().await; + let arn = repair(&server.target(), "").await.expect("create initial target"); + let stale = persisted_targets().await; + seed_unreadable(&env).await; + BucketTargetSys::get().update_all_targets(BUCKET, Some(&stale)).await; + let file = metadata_sys::get_config_from_disk(BUCKET) + .await + .expect("read metadata") + .save_file_path(); + let before = crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file) + .await + .expect("read original bytes"); + for operation in ["create", "update", "remove"] { + let error = match operation { + "create" => repair(&server.target(), "") + .await + .expect_err("cached idempotency cannot bypass unreadable disk"), + "update" => repair(&stale.targets[0], "update=true&sync=true") + .await + .expect_err("update cannot replace unreadable targets"), + "remove" => remove(&arn).await.expect_err("remove cannot replace unreadable targets"), + _ => unreachable!(), + }; + assert_eq!(error.code(), &S3ErrorCode::InternalError, "{operation}"); + assert_eq!( + crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file) + .await + .expect("read rejected write"), + before + ); + assert_published_targets(&stale).await; + } + }) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_create_is_idempotent_from_disk_when_the_remote_is_offline() { + temp_env::async_with_vars(ORDINARY_TARGET_ENV, async { + let (_temp, env) = test_env().await; + let mut server = RemoteTargetServer::start().await; + let target = server.target(); + let arn = repair(&target, "").await.expect("create initial target"); + let file = metadata_sys::get_config_from_disk(BUCKET) + .await + .expect("read metadata") + .save_file_path(); + let before = crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file) + .await + .expect("read original bytes"); + BucketTargetSys::get() + .update_all_targets(BUCKET, Some(&BucketTargets::default())) + .await; + server.task.abort(); + assert!((&mut server.task).await.expect_err("remote server must stop").is_cancelled()); + assert_eq!( + tokio::time::timeout(Duration::from_secs(10), repair(&target, "")) + .await + .expect("idempotent create must not wait for the remote") + .expect("recognize persisted target"), + arn + ); + assert_eq!( + crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file) + .await + .expect("read idempotent create bytes"), + before + ); + let targets = persisted_targets().await; + assert_eq!(targets.targets.len(), 1); + assert_eq!(targets.targets[0].arn, arn); + }) + .await; + } + + async fn assert_target_update_rejects_concurrent_change(deleted: bool) { + temp_env::async_with_vars( + ORDINARY_TARGET_ENV, + Box::pin(async move { + let (_temp, _env) = test_env().await; + let server = RemoteTargetServer::start().await; + let arn = repair(&server.target(), "").await.expect("create initial target"); + let mut original = persisted_targets().await; + let mut requested = original.targets[0].clone(); + requested.replication_sync = true; + if deleted { + original.targets.clear(); + } else { + original.targets[0].bandwidth_limit = 1024; + } + let (observed, release) = server.pause_next_request(); + let update = repair(&requested, "update=true&sync=true"); + let peer_write = async { + tokio::time::timeout(Duration::from_secs(10), observed) + .await + .expect("validation must reach HTTP source") + .expect("observe validation"); + metadata_sys::update(BUCKET, BUCKET_TARGETS_FILE, serde_json::to_vec(&original).expect("encode peer change")) + .await + .expect("commit peer change during validation"); + release.send(()).expect("release validation response"); + }; + let (result, ()) = tokio::join!(update, peer_write); + assert_eq!( + result.expect_err("stale target update must conflict").code(), + &S3ErrorCode::OperationAborted + ); + let targets = persisted_targets().await; + assert_eq!( + serde_json::to_value(&targets).expect("encode current targets"), + serde_json::to_value(&original).expect("encode peer targets") + ); + assert_published_targets(&targets).await; + if !deleted { + assert_eq!(targets.targets[0].arn, arn); + assert!(!targets.targets[0].replication_sync); + } else { + assert!(BucketTargetSys::get().get_remote_target_client(BUCKET, &arn).await.is_none()); + } + }), + ) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_update_rejects_same_target_changes_during_remote_validation() { + assert_target_update_rejects_concurrent_change(false).await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_update_rejects_target_deletion_during_remote_validation() { + assert_target_update_rejects_concurrent_change(true).await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_update_preserves_a_repair_committed_during_validation() { + temp_env::async_with_vars(ORDINARY_TARGET_ENV, async { + let (_temp, _env) = test_env().await; + let first = RemoteTargetServer::start().await; + let second = RemoteTargetServer::start().await; + let first_arn = repair(&first.target(), "").await.expect("create initial target"); + let mut requested = persisted_targets().await.targets.remove(0); + requested.replication_sync = true; + let (observed, release) = first.pause_next_request(); + let update = repair(&requested, "update=true&sync=true"); + let concurrent_repair = async { + tokio::time::timeout(Duration::from_secs(10), observed) + .await + .expect("update validation must reach HTTP source") + .expect("observe update validation"); + let repaired = + tokio::time::timeout(Duration::from_secs(10), repair(&second.target(), "replace-unreadable=true")).await; + let second_arn = repaired + .expect("repair must complete during validation") + .expect("commit concurrent repair"); + let repaired = persisted_targets() + .await + .targets + .into_iter() + .find(|target| target.arn == second_arn) + .expect("read committed repair"); + release.send(()).expect("release update validation response"); + repaired + }; + let (updated_arn, repaired) = tokio::join!(update, concurrent_repair); + assert_eq!(updated_arn.expect("an unrelated target change must not conflict"), first_arn); + let targets = persisted_targets().await; + assert_eq!(targets.targets.len(), 2); + let updated = targets + .targets + .iter() + .find(|target| target.arn == first_arn) + .expect("updated target"); + assert_eq!( + serde_json::to_value(updated).expect("encode updated target"), + serde_json::to_value(requested).expect("encode requested target") + ); + let kept = targets + .targets + .iter() + .find(|target| target.arn == repaired.arn) + .expect("retained repair"); + assert_eq!( + serde_json::to_value(kept).expect("encode retained repair"), + serde_json::to_value(repaired).expect("encode committed repair") + ); + assert_published_targets(&targets).await; + }) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_validation_does_not_lock_out_a_concurrent_repair() { + temp_env::async_with_vars(ORDINARY_TARGET_ENV, async { + let (_temp, _env) = test_env().await; + let first = RemoteTargetServer::start().await; + let second = RemoteTargetServer::start().await; + let target = first.target(); + BucketTargetSys::get() + .validate_target(BUCKET, &target) + .await + .expect("remote validation must succeed before pausing the create request"); + let (observed, release) = first.pause_next_request(); + let create = repair(&target, ""); + let concurrent_repair = async { + tokio::time::timeout(Duration::from_secs(10), observed) + .await + .expect("validation must reach HTTP source") + .expect("observe validation"); + let repaired = + tokio::time::timeout(Duration::from_secs(10), repair(&second.target(), "replace-unreadable=true")).await; + release.send(()).expect("release validation response"); + repaired + .expect("repair must complete while another remote validation is paused") + .expect("concurrent repair") + }; + let (first_arn, second_arn) = tokio::join!(create, concurrent_repair); + let first_arn = first_arn.expect("merge validated create"); + let targets = persisted_targets().await; + assert_eq!(targets.targets.len(), 2); + assert!(targets.targets.iter().any(|target| target.arn == first_arn)); + assert!(targets.targets.iter().any(|target| target.arn == second_arn)); + assert_published_targets(&targets).await; + }) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_remove_keeps_both_guards_after_the_http_future_is_dropped() { + temp_env::async_with_vars(ORDINARY_TARGET_ENV, async { + let (_temp, _env) = test_env().await; + let server = RemoteTargetServer::start().await; + let arn = repair(&server.target(), "").await.expect("create initial target"); + let mut replacement = persisted_targets().await.targets.remove(0); + replacement.arn = "arn:rustfs:replication:us-east-1:replacement:remote".to_string(); + let expected = BucketTargets { targets: vec![replacement] }; + let encoded = serde_json::to_vec(&expected).expect("encode next target writer"); + let (reached, observed) = oneshot::channel(); + let (release, resume) = oneshot::channel(); + assert!(REMOVAL_PAUSE.lock().expect("removal pause lock").replace(( + BUCKET.to_string(), RequestPause { reached, release: resume } + )).is_none()); + let removing = tokio::spawn(async move { remove(&arn).await }); + tokio::time::timeout(Duration::from_secs(10), observed).await + .expect("removal must reach cancellation boundary").expect("observe removal"); + removing.abort(); + assert!(removing.await.expect_err("HTTP future must be dropped").is_cancelled()); + + let mut local_writer = Box::pin(lock_bucket_targets_metadata(BUCKET)); + assert!(futures::poll!(local_writer.as_mut()).is_pending(), "HTTP cancellation must not release the local target guard"); + drop(local_writer); + let probe = metadata_sys::ConfigWriteLockProbe::install(BUCKET); + let (entered, mut entered_rx) = oneshot::channel(); + let mut peer_writer = Box::pin(metadata_sys::update_config_with(BUCKET, BUCKET_TARGETS_FILE, move |metadata| { + entered.send(()).expect("observe peer transaction"); + assert!(metadata.bucket_target_config.as_ref().expect("read removed targets").targets.is_empty(), + "a competing writer must observe the completed removal before entering"); + Ok(encoded) + })); + tokio::select! { + result = peer_writer.as_mut() => panic!("HTTP cancellation released the metadata guard before commit: {result:?}"), + () = probe.wait_until_attempted() => {} + } + assert!(matches!(entered_rx.try_recv(), Err(oneshot::error::TryRecvError::Empty)), + "the peer writer must remain outside the metadata transaction"); + release.send(()).expect("allow cancellation and target commit"); + tokio::time::timeout(Duration::from_secs(10), peer_writer).await + .expect("peer writer must proceed after removal commits").expect("persist peer target"); + entered_rx.await.expect("peer transaction entered"); + let targets = persisted_targets().await; + assert_eq!(serde_json::to_value(&targets).expect("encode disk targets"), serde_json::to_value(expected).expect("encode expected targets")); + assert_published_targets(&targets).await; + }).await; + } + + async fn assert_target_remove_checks_persisted_rules(malformed: bool) { + temp_env::async_with_vars(ORDINARY_TARGET_ENV, Box::pin(async move { + let (_temp, env) = test_env().await; + let server = RemoteTargetServer::start().await; + let arn = repair(&server.target(), "").await.expect("create initial target"); + let mut metadata = metadata_sys::get_config_from_disk(BUCKET).await.expect("read initial metadata"); + let file = metadata.save_file_path(); + metadata.replication_config_xml = if malformed { + b"".to_vec() + } else { + format!("{arn}activeEnabled1Disabled{arn}").into_bytes() + }; + // Another node has persisted rules before this node reloads them. + metadata.save_with_store(Arc::clone(&env.ecstore)).await.expect("persist peer rules without refreshing cache"); + assert!(metadata_sys::get_replication_config(BUCKET).await.is_err(), "precondition: cached rules are absent"); + assert_eq!(metadata_sys::get_config_from_disk(BUCKET).await.expect("read peer metadata").replication_config.is_none(), malformed); + let before = crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file).await.expect("read peer bytes"); + let error = remove(&arn).await.expect_err("rules must prevent unsafe target removal"); + assert_eq!(error.code(), if malformed { &S3ErrorCode::InternalError } else { &S3ErrorCode::InvalidRequest }); + assert_eq!(crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file).await.expect("read rejected removal bytes"), before); + let targets = persisted_targets().await; + assert_eq!(targets.targets.len(), 1); + assert_eq!(targets.targets[0].arn, arn); + assert_published_targets(&targets).await; + })) + .await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_remove_checks_current_persisted_replication_rules() { + assert_target_remove_checks_persisted_rules(false).await; + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_target_remove_rejects_malformed_persisted_replication_rules() { + assert_target_remove_checks_persisted_rules(true).await; + } } diff --git a/rustfs/src/admin/storage_api.rs b/rustfs/src/admin/storage_api.rs index ecf18cd6e..56f634c42 100644 --- a/rustfs/src/admin/storage_api.rs +++ b/rustfs/src/admin/storage_api.rs @@ -295,6 +295,8 @@ pub(crate) mod remote_s3_client { } pub(crate) mod metadata_sys { + #[cfg(test)] + pub(crate) use super::ecstore_bucket::metadata_sys::ConfigWriteLockProbe; use std::sync::Arc; use rustfs_policy::policy::BucketPolicy; @@ -667,7 +669,7 @@ pub(crate) mod replication { } pub(crate) mod target { - pub(crate) use super::ecstore_bucket::target::duration_from_secs_or_nanos; + pub(crate) use super::ecstore_bucket::target::{ARN, duration_from_secs_or_nanos}; pub(crate) type BucketTarget = super::ecstore_bucket::target::BucketTarget; pub(crate) type BucketTargetType = super::ecstore_bucket::target::BucketTargetType; pub(crate) type BucketTargets = super::ecstore_bucket::target::BucketTargets;