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
This commit is contained in:
Zhengchao An
2026-09-06 12:51:44 +08:00
committed by GitHub
parent 71f1dcf859
commit a8aaadb886
2 changed files with 716 additions and 74 deletions
+713 -73
View File
@@ -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<String> {
enum RemoteTargetWrite {
Create { replace_unreadable: bool },
Update { expected: Vec<u8> },
}
async fn persist_remote_target_repair(bucket: &str, target: BucketTarget, incarnation: uuid::Uuid) -> S3Result<String> {
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<String> {
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<std::sync::Mutex<Option<RequestPause>>>,
}
struct RequestPause {
reached: oneshot::Sender<()>,
release: oneshot::Receiver<()>,
}
static REMOVAL_PAUSE: std::sync::Mutex<Option<(String, RequestPause)>> = 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::<RequestPause>));
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<u8>) -> S3Request<Body> {
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"<ReplicationConfiguration>".to_vec()
} else {
format!("<ReplicationConfiguration><Role>{arn}</Role><Rule><ID>active</ID><Status>Enabled</Status><Filter><Prefix></Prefix></Filter><Priority>1</Priority><DeleteMarkerReplication><Status>Disabled</Status></DeleteMarkerReplication><Destination><Bucket>{arn}</Bucket></Destination></Rule></ReplicationConfiguration>").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;
}
}
+3 -1
View File
@@ -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;