fix(replication): prevent target state loss across buckets (#2704)

This commit is contained in:
weisd
2026-04-27 17:27:33 +08:00
committed by GitHub
parent cfbd094bc4
commit 334184b005
4 changed files with 271 additions and 73 deletions
@@ -21,6 +21,7 @@ use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
use aws_sdk_s3::{Client, Config};
use http::header::{CONTENT_TYPE, HOST};
use reqwest::StatusCode;
use rustfs_ecstore::bucket::bucket_target_sys::BucketTargetSys;
use rustfs_madmin::{
AddServiceAccountReq, ListServiceAccountsResp, PeerInfo, PeerSite, ReplicateAddStatus, ReplicateEditStatus,
ReplicateRemoveStatus, SRRemoveReq, SRResyncOpStatus, SRStatusInfo, SiteReplicationInfo, SyncStatus,
@@ -1648,6 +1649,90 @@ async fn test_single_bucket_replication_fans_out_to_multiple_targets() -> Result
Ok(())
}
#[tokio::test]
#[serial]
async fn test_sequential_bucket_replication_succeeds_for_multiple_buckets() -> Result<(), Box<dyn Error + Send + Sync>> {
init_logging();
let mut source_env = RustFSTestEnvironment::new().await?;
source_env.start_rustfs_server(vec![]).await?;
let mut target_env = RustFSTestEnvironment::new().await?;
target_env.start_rustfs_server_without_cleanup(vec![]).await?;
let source_client = source_env.create_s3_client();
let target_client = target_env.create_s3_client();
for idx in 1..=5 {
let source_bucket = format!("replication-multi-src-{idx}");
let target_bucket = format!("replication-multi-dst-{idx}");
let object_key = format!("probe-{idx}.txt");
let body = format!("payload-{idx}");
source_client.create_bucket().bucket(&source_bucket).send().await?;
target_client.create_bucket().bucket(&target_bucket).send().await?;
enable_bucket_versioning(&source_env, &source_bucket).await?;
enable_bucket_versioning(&target_env, &target_bucket).await?;
let target_arn = set_replication_target(&source_env, &source_bucket, &target_env, &target_bucket).await?;
put_bucket_replication(&source_env, &source_bucket, &target_arn).await?;
source_client
.put_object()
.bucket(&source_bucket)
.key(&object_key)
.body(ByteStream::from(body.clone().into_bytes()))
.send()
.await?;
wait_for_replicated_object(&target_client, &target_bucket, &object_key, &body).await?;
}
Ok(())
}
#[tokio::test]
#[serial]
async fn test_replication_recovers_after_runtime_target_cache_is_cleared() -> Result<(), Box<dyn Error + Send + Sync>> {
init_logging();
let mut source_env = RustFSTestEnvironment::new().await?;
source_env.start_rustfs_server(vec![]).await?;
let mut target_env = RustFSTestEnvironment::new().await?;
target_env.start_rustfs_server_without_cleanup(vec![]).await?;
let source_bucket = "replication-refresh-src";
let target_bucket = "replication-refresh-dst";
let object_key = "probe-refresh.txt";
let body = "payload-refresh";
let source_client = source_env.create_s3_client();
let target_client = target_env.create_s3_client();
source_client.create_bucket().bucket(source_bucket).send().await?;
target_client.create_bucket().bucket(target_bucket).send().await?;
enable_bucket_versioning(&source_env, source_bucket).await?;
enable_bucket_versioning(&target_env, target_bucket).await?;
let target_arn = set_replication_target(&source_env, source_bucket, &target_env, target_bucket).await?;
put_bucket_replication(&source_env, source_bucket, &target_arn).await?;
BucketTargetSys::get().delete(source_bucket).await;
source_client
.put_object()
.bucket(source_bucket)
.key(object_key)
.body(ByteStream::from(body.as_bytes().to_vec()))
.send()
.await?;
wait_for_replicated_object(&target_client, target_bucket, object_key, body).await?;
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_resync_start_cancel_restart_real_dual_node() -> Result<(), Box<dyn Error + Send + Sync>> {
+63 -59
View File
@@ -395,8 +395,27 @@ impl BucketTargetSys {
}
}
pub async fn set_target(&self, bucket: &str, target: &BucketTarget, update: bool) -> Result<(), BucketTargetError> {
if !target.target_type.is_valid() && !update {
pub async fn set_target(
&self,
bucket: &str,
target: &BucketTarget,
update: bool,
) -> Result<BucketTargets, BucketTargetError> {
self.validate_target(bucket, target).await?;
let mut bucket_targets = match self.list_bucket_targets(bucket).await {
Ok(targets) => targets,
Err(BucketTargetError::BucketRemoteTargetNotFound { .. }) => BucketTargets::default(),
Err(err) => return Err(err),
};
Self::upsert_target_entry(&mut bucket_targets.targets, target, update)?;
Ok(bucket_targets)
}
pub async fn validate_target(&self, bucket: &str, target: &BucketTarget) -> Result<(), BucketTargetError> {
if !target.target_type.is_valid() {
return Err(BucketTargetError::BucketRemoteArnTypeInvalid {
bucket: bucket.to_string(),
});
@@ -450,52 +469,44 @@ impl BucketTargetSys {
}
}
{
let mut targets_map = self.targets_map.write().await;
let bucket_targets = targets_map.entry(bucket.to_string()).or_insert_with(Vec::new);
let mut found = false;
Ok(())
}
for (idx, existing_target) in bucket_targets.iter().enumerate() {
if existing_target.target_type.to_string() == target.target_type.to_string() {
if existing_target.arn == target.arn {
if !update {
return Err(BucketTargetError::BucketRemoteAlreadyExists {
bucket: existing_target.target_bucket.clone(),
});
}
bucket_targets[idx] = target.clone();
found = true;
break;
}
if existing_target.endpoint == target.endpoint {
fn upsert_target_entry(
bucket_targets: &mut Vec<BucketTarget>,
target: &BucketTarget,
update: bool,
) -> Result<(), BucketTargetError> {
let mut found = false;
for (idx, existing_target) in bucket_targets.iter().enumerate() {
if existing_target.target_type.to_string() == target.target_type.to_string() {
if existing_target.arn == target.arn {
if !update {
return Err(BucketTargetError::BucketRemoteAlreadyExists {
bucket: existing_target.target_bucket.clone(),
});
}
bucket_targets[idx] = target.clone();
found = true;
break;
}
if existing_target.endpoint == target.endpoint {
return Err(BucketTargetError::BucketRemoteAlreadyExists {
bucket: existing_target.target_bucket.clone(),
});
}
}
if !found && !update {
bucket_targets.push(target.clone());
}
}
{
let mut arn_remotes_map = self.arn_remotes_map.write().await;
arn_remotes_map.insert(
target.arn.clone(),
ArnTarget {
client: Some(Arc::new(target_client)),
last_refresh: OffsetDateTime::now_utc(),
},
);
if !found && !update {
bucket_targets.push(target.clone());
}
self.update_bandwidth_limit(bucket, &target.arn, target.bandwidth_limit);
Ok(())
}
pub async fn remove_target(&self, bucket: &str, arn_str: &str) -> Result<(), BucketTargetError> {
pub async fn remove_target(&self, bucket: &str, arn_str: &str) -> Result<BucketTargets, BucketTargetError> {
if arn_str.is_empty() {
return Err(BucketTargetError::BucketRemoteArnInvalid {
bucket: bucket.to_string(),
@@ -524,33 +535,16 @@ impl BucketTargetSys {
}
}
{
let mut targets_map = self.targets_map.write().await;
let targets = self.list_bucket_targets(bucket).await?;
let new_targets: Vec<BucketTarget> = targets.targets.iter().filter(|t| t.arn != arn_str).cloned().collect();
let Some(targets) = targets_map.get(bucket) else {
return Err(BucketTargetError::BucketRemoteTargetNotFound {
bucket: bucket.to_string(),
});
};
let new_targets: Vec<BucketTarget> = targets.iter().filter(|t| t.arn != arn_str).cloned().collect();
if new_targets.len() == targets.len() {
return Err(BucketTargetError::BucketRemoteTargetNotFound {
bucket: bucket.to_string(),
});
}
targets_map.insert(bucket.to_string(), new_targets);
if new_targets.len() == targets.targets.len() {
return Err(BucketTargetError::BucketRemoteTargetNotFound {
bucket: bucket.to_string(),
});
}
{
self.arn_remotes_map.write().await.remove(arn_str);
}
self.update_bandwidth_limit(bucket, arn_str, 0);
Ok(())
Ok(BucketTargets { targets: new_targets })
}
pub async fn mark_refresh_in_progress(&self, bucket: &str, arn: &str) {
@@ -603,7 +597,7 @@ impl BucketTargetSys {
if let Some(last_refresh) = last_refresh {
let now = OffsetDateTime::now_utc();
if now - last_refresh > Duration::from_secs(60 * 5) {
if now - last_refresh < Duration::from_secs(60 * 5) {
return None;
}
}
@@ -619,6 +613,16 @@ impl BucketTargetSys {
}
};
let cli = self
.arn_remotes_map
.read()
.await
.get(arn)
.and_then(|target| target.client.clone());
if cli.is_some() {
return cli;
}
self.inc_arn_errs(bucket, arn).await;
None
}