feat(admin): complete site replication support (#2346)

This commit is contained in:
weisd
2026-03-31 13:01:38 +08:00
committed by GitHub
parent 6bf0c542a1
commit 15995aae14
22 changed files with 5429 additions and 101 deletions
@@ -13,14 +13,38 @@
// limitations under the License.
use crate::common::{RustFSTestEnvironment, init_logging, local_http_client};
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
use http::header::{CONTENT_TYPE, HOST};
use reqwest::StatusCode;
use rustfs_madmin::{
PeerInfo, PeerSite, ReplicateAddStatus, ReplicateEditStatus, ReplicateRemoveStatus, SRRemoveReq, SRResyncOpStatus,
SRStatusInfo, SiteReplicationInfo, SyncStatus,
};
use rustfs_signer::constants::UNSIGNED_PAYLOAD;
use rustfs_signer::sign_v4;
use s3s::Body;
use serial_test::serial;
use std::collections::BTreeMap;
use std::error::Error;
use time::Duration as TimeDuration;
use tokio::time::{Duration, sleep};
#[derive(Debug, Clone, serde::Deserialize)]
struct ReplicationResetStatusResponse {
#[serde(rename = "Targets", default)]
targets: Vec<ReplicationResetStatusTarget>,
}
#[derive(Debug, Clone, serde::Deserialize)]
struct ReplicationResetStatusTarget {
#[serde(rename = "Arn", default)]
arn: String,
#[serde(rename = "ResetID", default)]
reset_id: String,
#[serde(rename = "Status", default)]
status: String,
}
async fn signed_request(
method: http::Method,
@@ -239,6 +263,265 @@ async fn list_replication_targets_request(
signed_request(http::Method::GET, &url, &env.access_key, &env.secret_key, None, None).await
}
async fn site_replication_add(
env: &RustFSTestEnvironment,
sites: &[PeerSite],
) -> Result<ReplicateAddStatus, Box<dyn Error + Send + Sync>> {
let url = format!("{}/rustfs/admin/v3/site-replication/add?replicateILMExpiry=false", env.url);
let response = signed_request(
http::Method::PUT,
&url,
&env.access_key,
&env.secret_key,
Some(serde_json::to_vec(sites)?),
Some("application/json"),
)
.await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("site replication add failed: {status} {body}").into());
}
Ok(serde_json::from_slice(&response.bytes().await?)?)
}
async fn site_replication_info(env: &RustFSTestEnvironment) -> Result<SiteReplicationInfo, Box<dyn Error + Send + Sync>> {
let url = format!("{}/rustfs/admin/v3/site-replication/info", env.url);
let response = signed_request(http::Method::GET, &url, &env.access_key, &env.secret_key, None, None).await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("site replication info failed: {status} {body}").into());
}
Ok(serde_json::from_slice(&response.bytes().await?)?)
}
async fn site_replication_resync_op(
env: &RustFSTestEnvironment,
operation: &str,
peer: &PeerInfo,
) -> Result<SRResyncOpStatus, Box<dyn Error + Send + Sync>> {
let url = format!("{}/rustfs/admin/v3/site-replication/resync/op?operation={operation}", env.url);
let response = signed_request(
http::Method::PUT,
&url,
&env.access_key,
&env.secret_key,
Some(serde_json::to_vec(peer)?),
Some("application/json"),
)
.await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("site replication resync {operation} failed: {status} {body}").into());
}
Ok(serde_json::from_slice(&response.bytes().await?)?)
}
async fn site_replication_edit(
env: &RustFSTestEnvironment,
query: &str,
peer: &PeerInfo,
) -> Result<ReplicateEditStatus, Box<dyn Error + Send + Sync>> {
let url = if query.is_empty() {
format!("{}/rustfs/admin/v3/site-replication/edit", env.url)
} else {
format!("{}/rustfs/admin/v3/site-replication/edit?{query}", env.url)
};
let response = signed_request(
http::Method::PUT,
&url,
&env.access_key,
&env.secret_key,
Some(serde_json::to_vec(peer)?),
Some("application/json"),
)
.await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("site replication edit failed: {status} {body}").into());
}
Ok(serde_json::from_slice(&response.bytes().await?)?)
}
async fn site_replication_status(env: &RustFSTestEnvironment, query: &str) -> Result<SRStatusInfo, Box<dyn Error + Send + Sync>> {
let url = if query.is_empty() {
format!("{}/rustfs/admin/v3/site-replication/status", env.url)
} else {
format!("{}/rustfs/admin/v3/site-replication/status?{query}", env.url)
};
let response = signed_request(http::Method::GET, &url, &env.access_key, &env.secret_key, None, None).await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("site replication status failed: {status} {body}").into());
}
Ok(serde_json::from_slice(&response.bytes().await?)?)
}
async fn site_replication_remove(
env: &RustFSTestEnvironment,
req: &SRRemoveReq,
) -> Result<ReplicateRemoveStatus, Box<dyn Error + Send + Sync>> {
let url = format!("{}/rustfs/admin/v3/site-replication/remove", env.url);
let response = signed_request(
http::Method::PUT,
&url,
&env.access_key,
&env.secret_key,
Some(serde_json::to_vec(req)?),
Some("application/json"),
)
.await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("site replication remove failed: {status} {body}").into());
}
Ok(serde_json::from_slice(&response.bytes().await?)?)
}
async fn site_replication_state_edit(
env: &RustFSTestEnvironment,
body: &rustfs_madmin::SRStateEditReq,
) -> Result<(), Box<dyn Error + Send + Sync>> {
let url = format!("{}/rustfs/admin/v3/site-replication/state/edit", env.url);
let response = signed_request(
http::Method::PUT,
&url,
&env.access_key,
&env.secret_key,
Some(serde_json::to_vec(body)?),
Some("application/json"),
)
.await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("site replication state edit failed: {status} {body}").into());
}
Ok(())
}
async fn get_replication_reset_status(
env: &RustFSTestEnvironment,
bucket: &str,
arn: &str,
) -> Result<ReplicationResetStatusResponse, Box<dyn Error + Send + Sync>> {
let url = format!("{}/{bucket}?replication-reset-status&arn={}", env.url, urlencoding::encode(arn));
let response = signed_request(http::Method::GET, &url, &env.access_key, &env.secret_key, None, None).await?;
if response.status() != StatusCode::OK {
let status = response.status();
let body = response.text().await.unwrap_or_default();
return Err(format!("replication reset status failed: {status} {body}").into());
}
Ok(serde_json::from_slice(&response.bytes().await?)?)
}
async fn wait_for_site_replication_enabled(
env: &RustFSTestEnvironment,
expected_sites: usize,
) -> Result<SiteReplicationInfo, Box<dyn Error + Send + Sync>> {
for _ in 0..40 {
let info = site_replication_info(env).await?;
if info.enabled && info.sites.len() == expected_sites {
return Ok(info);
}
sleep(Duration::from_millis(250)).await;
}
Err(format!("site replication did not reach {expected_sites} sites on {}", env.address).into())
}
async fn wait_for_site_replication_disabled(
env: &RustFSTestEnvironment,
) -> Result<SiteReplicationInfo, Box<dyn Error + Send + Sync>> {
wait_for_site_replication_info(env, |info| !info.enabled && info.sites.is_empty()).await
}
async fn wait_for_site_replication_info<F>(
env: &RustFSTestEnvironment,
predicate: F,
) -> Result<SiteReplicationInfo, Box<dyn Error + Send + Sync>>
where
F: Fn(&SiteReplicationInfo) -> bool,
{
for _ in 0..40 {
let info = site_replication_info(env).await?;
if predicate(&info) {
return Ok(info);
}
sleep(Duration::from_millis(250)).await;
}
Err(format!("site replication info did not reach expected state on {}", env.address).into())
}
async fn wait_for_site_replication_status<F>(
env: &RustFSTestEnvironment,
query: &str,
predicate: F,
) -> Result<SRStatusInfo, Box<dyn Error + Send + Sync>>
where
F: Fn(&SRStatusInfo) -> bool,
{
for _ in 0..40 {
let status = site_replication_status(env, query).await?;
if predicate(&status) {
return Ok(status);
}
sleep(Duration::from_millis(250)).await;
}
Err(format!("site replication status did not reach expected state on {}", env.address).into())
}
async fn wait_for_replication_reset_target<F>(
env: &RustFSTestEnvironment,
bucket: &str,
arn: &str,
predicate: F,
) -> Result<ReplicationResetStatusTarget, Box<dyn Error + Send + Sync>>
where
F: Fn(&ReplicationResetStatusTarget) -> bool,
{
let mut last_seen = None;
for _ in 0..40 {
let status = get_replication_reset_status(env, bucket, arn).await?;
if let Some(target) = status.targets.into_iter().find(|target| target.arn == arn) {
if predicate(&target) {
return Ok(target);
}
last_seen = Some(target);
}
sleep(Duration::from_millis(250)).await;
}
Err(format!(
"replication reset target {arn} for bucket {bucket} did not reach expected state; last seen: {:?}",
last_seen
)
.into())
}
async fn build_replication_pair(
enable_target_versioning: bool,
) -> Result<(RustFSTestEnvironment, RustFSTestEnvironment, String), Box<dyn Error + Send + Sync>> {
@@ -800,3 +1083,410 @@ async fn test_remove_remote_target_rejects_target_used_by_replication() -> Resul
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_resync_start_cancel_restart_real_dual_node() -> 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 = "site-repl-resync-src";
let target_bucket = "site-repl-resync-dst";
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 add_status = site_replication_add(
&source_env,
&[
PeerSite {
name: "source-site".to_string(),
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_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 remote_peer = source_info
.sites
.into_iter()
.find(|peer| peer.endpoint == target_env.url)
.ok_or("target peer missing from source site replication info")?;
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?;
for idx in 0..32 {
source_client
.put_object()
.bucket(source_bucket)
.key(format!("resync-object-{idx:02}"))
.body(ByteStream::from(vec![b'x'; 256 * 1024]))
.send()
.await?;
}
let started = site_replication_resync_op(&source_env, "start", &remote_peer).await?;
assert_eq!(started.status, "success", "unexpected start result: {:?}", started);
assert!(
started
.buckets
.iter()
.any(|bucket| bucket.bucket == source_bucket && matches!(bucket.status.as_str(), "started" | "success")),
"source bucket start status missing: {:?}",
started
);
let started_target =
wait_for_replication_reset_target(&source_env, source_bucket, &target_arn, |target| !target.reset_id.is_empty()).await?;
let started_reset_id = started_target.reset_id.clone();
assert!(
matches!(started_target.status.as_str(), "Pending" | "Started" | "InProgress" | "Completed"),
"unexpected start status: {:?}",
started_target
);
let canceled = site_replication_resync_op(&source_env, "cancel", &remote_peer).await?;
assert_eq!(canceled.status, "success", "unexpected cancel result: {:?}", canceled);
assert!(
canceled
.buckets
.iter()
.any(|bucket| bucket.bucket == source_bucket && matches!(bucket.status.as_str(), "canceled" | "success")),
"source bucket cancel status missing: {:?}",
canceled
);
let canceled_target =
wait_for_replication_reset_target(&source_env, source_bucket, &target_arn, |target| target.status == "Canceled").await?;
assert_eq!(canceled_target.status, "Canceled");
assert_eq!(canceled_target.reset_id, started_reset_id);
let restarted = site_replication_resync_op(&source_env, "start", &remote_peer).await?;
assert_eq!(restarted.status, "success", "unexpected restart result: {:?}", restarted);
assert!(
restarted
.buckets
.iter()
.any(|bucket| bucket.bucket == source_bucket && matches!(bucket.status.as_str(), "started" | "success")),
"source bucket restart status missing: {:?}",
restarted
);
let restart_snapshot = get_replication_reset_status(&source_env, source_bucket, &target_arn).await?;
let restarted_target = wait_for_replication_reset_target(&source_env, source_bucket, &target_arn, |target| {
!target.reset_id.is_empty() && target.reset_id != started_reset_id
})
.await
.map_err(|err| {
format!(
"restart ids: start={} restart={} snapshot={:?}; {err}",
started_reset_id, restarted.resync_id, restart_snapshot.targets
)
})?;
assert_ne!(restarted_target.reset_id, started_reset_id);
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_edit_and_status_peer_state_real_dual_node() -> 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 add_status = site_replication_add(
&source_env,
&[
PeerSite {
name: "source-site".to_string(),
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_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 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")?;
remote_peer.sync_state = SyncStatus::Enable;
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)
})
.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)
})
.await?;
assert!(
source_after_sync
.sites
.iter()
.any(|peer| peer.endpoint == target_env.url && peer.sync_state == SyncStatus::Enable)
);
assert!(
target_after_sync
.sites
.iter()
.any(|peer| peer.endpoint == target_env.url && peer.sync_state == SyncStatus::Enable)
);
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)
})
.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)
})
.await?;
assert!(source_after_ilm.sites.iter().all(|peer| peer.replicate_ilm_expiry));
assert!(target_after_ilm.sites.iter().all(|peer| peer.replicate_ilm_expiry));
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
.values()
.all(|state| state.peers.len() == 2 && 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
.values()
.all(|state| state.peers.len() == 2 && 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!(
source_status
.peer_states
.values()
.all(|state| state.peers.values().all(|peer| peer.replicate_ilm_expiry))
);
assert!(
target_status
.peer_states
.values()
.all(|state| state.peers.values().all(|peer| peer.replicate_ilm_expiry))
);
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_remove_all_real_dual_node() -> 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 add_status = site_replication_add(
&source_env,
&[
PeerSite {
name: "source-site".to_string(),
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_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 remove_status = site_replication_remove(
&source_env,
&SRRemoveReq {
remove_all: true,
..Default::default()
},
)
.await?;
assert!(
!remove_status.status.is_empty() && remove_status.err_detail.is_empty(),
"unexpected site remove result: {:?}",
remove_status
);
let source_after_remove = wait_for_site_replication_disabled(&source_env).await?;
let target_after_remove = wait_for_site_replication_disabled(&target_env).await?;
assert!(!source_after_remove.enabled);
assert!(source_after_remove.sites.is_empty());
assert!(!target_after_remove.enabled);
assert!(target_after_remove.sites.is_empty());
Ok(())
}
#[tokio::test]
#[serial]
async fn test_site_replication_state_edit_fresh_and_stale_real_dual_node() -> 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 add_status = site_replication_add(
&source_env,
&[
PeerSite {
name: "source-site".to_string(),
endpoint: source_env.url.clone(),
access_key: source_env.access_key.clone(),
secret_key: source_env.secret_key.clone(),
},
PeerSite {
name: "target-site".to_string(),
endpoint: target_env.url.clone(),
access_key: target_env.access_key.clone(),
secret_key: target_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?;
assert!(source_info.sites.iter().all(|peer| !peer.replicate_ilm_expiry));
assert!(target_info.sites.iter().all(|peer| !peer.replicate_ilm_expiry));
let target_status =
wait_for_site_replication_status(&target_env, "peer-state=true", |status| status.peer_states.len() == 2).await?;
let current_updated_at = target_status
.peer_states
.values()
.find_map(|state| state.updated_at)
.ok_or("missing target site replication updated_at")?;
let mut stale_peers = BTreeMap::new();
for peer in target_info.sites {
let mut peer = peer;
peer.replicate_ilm_expiry = true;
stale_peers.insert(peer.deployment_id.clone(), peer);
}
site_replication_state_edit(
&target_env,
&rustfs_madmin::SRStateEditReq {
peers: stale_peers,
updated_at: Some(current_updated_at - TimeDuration::seconds(1)),
},
)
.await?;
let target_after_stale = site_replication_info(&target_env).await?;
let source_after_stale = site_replication_info(&source_env).await?;
assert!(target_after_stale.sites.iter().all(|peer| !peer.replicate_ilm_expiry));
assert!(source_after_stale.sites.iter().all(|peer| !peer.replicate_ilm_expiry));
let mut fresh_peers = BTreeMap::new();
for peer in target_after_stale.sites {
let mut peer = peer;
peer.replicate_ilm_expiry = true;
fresh_peers.insert(peer.deployment_id.clone(), peer);
}
let fresh_updated_at = current_updated_at + TimeDuration::seconds(1);
site_replication_state_edit(
&target_env,
&rustfs_madmin::SRStateEditReq {
peers: fresh_peers,
updated_at: Some(fresh_updated_at),
},
)
.await?;
let target_after_fresh = wait_for_site_replication_info(&target_env, |info| {
info.sites.len() == 2 && info.sites.iter().all(|peer| peer.replicate_ilm_expiry)
})
.await?;
assert!(target_after_fresh.sites.iter().all(|peer| peer.replicate_ilm_expiry));
let target_status_after_fresh = wait_for_site_replication_status(&target_env, "peer-state=true", |status| {
status.peer_states.len() == 2
&& status.peer_states.values().all(|state| {
state.updated_at == Some(fresh_updated_at) && state.peers.values().all(|peer| peer.replicate_ilm_expiry)
})
})
.await?;
assert!(target_status_after_fresh.peer_states.values().all(|state| {
state.updated_at == Some(fresh_updated_at) && state.peers.values().all(|peer| peer.replicate_ilm_expiry)
}));
let source_after_fresh = site_replication_info(&source_env).await?;
assert!(source_after_fresh.sites.iter().all(|peer| !peer.replicate_ilm_expiry));
Ok(())
}
@@ -777,6 +777,14 @@ impl<S: StorageAPI> ReplicationPool<S> {
Ok(status)
}
pub async fn cancel_bucket_resync(&self, opts: ResyncOpts) -> Result<(), EcstoreError> {
self.resyncer.cancel(&opts).await;
self.resyncer
.mark_status(ResyncStatusType::ResyncCanceled, opts, self.storage.clone())
.await?;
Ok(())
}
pub async fn start_bucket_resync(self: Arc<Self>, opts: ResyncOpts) -> Result<(), EcstoreError> {
let now = OffsetDateTime::now_utc();
let bucket_status = {
@@ -813,8 +821,14 @@ impl<S: StorageAPI> ReplicationPool<S> {
let resyncer = self.resyncer.clone();
let storage = self.storage.clone();
let cancel_token = CancellationToken::new();
resyncer.register_cancel_token(&opts, cancel_token.clone()).await;
tokio::spawn(async move {
resyncer.resync_bucket(CancellationToken::new(), storage, false, opts).await;
resyncer
.clone()
.resync_bucket(cancel_token, storage, false, opts.clone())
.await;
resyncer.clear_cancel_token(&opts).await;
});
Ok(())
@@ -852,8 +866,12 @@ impl<S: StorageAPI> ReplicationPool<S> {
}
/// Load bucket replication resync statuses into memory
#[instrument(skip(cancellation_token))]
async fn load_resync(self: Arc<Self>, buckets: &[String], cancellation_token: CancellationToken) -> Result<(), EcstoreError> {
#[instrument(skip(_cancellation_token))]
async fn load_resync(
self: Arc<Self>,
buckets: &[String],
_cancellation_token: CancellationToken,
) -> Result<(), EcstoreError> {
// TODO: add leader_lock
// Make sure only one node running resync on the cluster
// Note: Leader lock implementation would be needed here
@@ -884,24 +902,20 @@ impl<S: StorageAPI> ReplicationPool<S> {
// Note: This would spawn a resync task in a real implementation
// For now, we just log the resync request
let ctx = cancellation_token.clone();
let ctx = CancellationToken::new();
let bucket_clone = bucket.clone();
let resync = self.resyncer.clone();
let storage = self.storage.clone();
let opts = ResyncOpts {
bucket: bucket_clone,
arn,
resync_id: stats.resync_id,
resync_before: stats.resync_before_date,
};
tokio::spawn(async move {
resync
.resync_bucket(
ctx,
storage,
true,
ResyncOpts {
bucket: bucket_clone,
arn,
resync_id: stats.resync_id,
resync_before: stats.resync_before_date,
},
)
.await;
resync.register_cancel_token(&opts, ctx.clone()).await;
resync.clone().resync_bucket(ctx, storage, true, opts.clone()).await;
resync.clear_cancel_token(&opts).await;
});
}
_ => {}
@@ -949,6 +963,7 @@ pub trait ReplicationPoolTrait: std::fmt::Debug {
async fn queue_replica_delete_task(&self, ri: DeletedObjectReplicationInfo);
async fn resize(&self, priority: ReplicationPriority, max_workers: usize, max_l_workers: usize);
async fn get_bucket_resync_status(&self, bucket: &str) -> Result<BucketReplicationResyncStatus, EcstoreError>;
async fn cancel_bucket_resync(&self, opts: ResyncOpts) -> Result<(), EcstoreError>;
async fn start_bucket_resync(self: Arc<Self>, opts: ResyncOpts) -> Result<(), EcstoreError>;
async fn init_resync(
self: Arc<Self>,
@@ -976,6 +991,10 @@ impl<S: StorageAPI> ReplicationPoolTrait for ReplicationPool<S> {
self.get_bucket_resync_status(bucket).await
}
async fn cancel_bucket_resync(&self, opts: ResyncOpts) -> Result<(), EcstoreError> {
self.cancel_bucket_resync(opts).await
}
async fn start_bucket_resync(self: Arc<Self>, opts: ResyncOpts) -> Result<(), EcstoreError> {
self.start_bucket_resync(opts).await
}
@@ -110,6 +110,10 @@ fn normalize_wire_time(value: Option<OffsetDateTime>) -> Option<OffsetDateTime>
}
}
fn resync_state_accepts_update(state: &TargetReplicationResyncStatus, opts: &ResyncOpts) -> bool {
state.resync_id.is_empty() || opts.resync_id.is_empty() || state.resync_id == opts.resync_id
}
#[derive(Debug, Clone, Default)]
pub struct ResyncOpts {
pub bucket: String,
@@ -360,16 +364,13 @@ static RESYNC_WORKER_COUNT: usize = 10;
pub struct ReplicationResyncer {
pub status_map: Arc<RwLock<HashMap<String, BucketReplicationResyncStatus>>>,
pub worker_size: usize,
pub resync_cancel_tx: CancellationToken,
pub resync_cancel_rx: CancellationToken,
pub cancel_tokens: Arc<RwLock<HashMap<String, CancellationToken>>>,
pub worker_tx: tokio::sync::broadcast::Sender<()>,
pub worker_rx: tokio::sync::broadcast::Receiver<()>,
}
impl ReplicationResyncer {
pub async fn new() -> Self {
let resync_cancel_tx = CancellationToken::new();
let resync_cancel_rx = resync_cancel_tx.clone();
let (worker_tx, worker_rx) = tokio::sync::broadcast::channel(RESYNC_WORKER_COUNT);
for _ in 0..RESYNC_WORKER_COUNT {
@@ -381,16 +382,34 @@ impl ReplicationResyncer {
Self {
status_map: Arc::new(RwLock::new(HashMap::new())),
worker_size: RESYNC_WORKER_COUNT,
resync_cancel_tx,
resync_cancel_rx,
cancel_tokens: Arc::new(RwLock::new(HashMap::new())),
worker_tx,
worker_rx,
}
}
fn cancel_key(opts: &ResyncOpts) -> String {
format!("{}:{}", opts.bucket, opts.arn)
}
pub async fn register_cancel_token(&self, opts: &ResyncOpts, token: CancellationToken) {
self.cancel_tokens.write().await.insert(Self::cancel_key(opts), token);
}
pub async fn clear_cancel_token(&self, opts: &ResyncOpts) {
self.cancel_tokens.write().await.remove(&Self::cancel_key(opts));
}
pub async fn cancel(&self, opts: &ResyncOpts) {
if let Some(token) = self.cancel_tokens.write().await.remove(&Self::cancel_key(opts)) {
token.cancel();
}
}
pub async fn mark_status<S: StorageAPI>(&self, status: ResyncStatusType, opts: ResyncOpts, obj_layer: Arc<S>) -> Result<()> {
let bucket_status = {
let mut status_map = self.status_map.write().await;
let now = OffsetDateTime::now_utc();
let bucket_status = if let Some(bucket_status) = status_map.get_mut(&opts.bucket) {
bucket_status
@@ -409,10 +428,33 @@ impl ReplicationResyncer {
bucket_status.targets_map.get_mut(&opts.arn).unwrap()
};
state.resync_status = status;
state.last_update = Some(OffsetDateTime::now_utc());
if !resync_state_accepts_update(state, &opts) {
warn!(
bucket = %opts.bucket,
arn = %opts.arn,
incoming_resync_id = %opts.resync_id,
current_resync_id = %state.resync_id,
"ignoring stale resync status update"
);
return Ok(());
}
bucket_status.last_update = Some(OffsetDateTime::now_utc());
if state.resync_id.is_empty() {
state.resync_id = opts.resync_id.clone();
}
if state.resync_before_date.is_none() {
state.resync_before_date = opts.resync_before;
}
if state.bucket.is_empty() {
state.bucket = opts.bucket.clone();
}
if status == ResyncStatusType::ResyncStarted && state.start_time.is_none() {
state.start_time = Some(now);
}
state.resync_status = status;
state.last_update = Some(now);
bucket_status.last_update = Some(now);
bucket_status.clone()
};
@@ -424,6 +466,7 @@ impl ReplicationResyncer {
pub async fn inc_stats(&self, status: &TargetReplicationResyncStatus, opts: ResyncOpts) {
let mut status_map = self.status_map.write().await;
let now = OffsetDateTime::now_utc();
let bucket_status = if let Some(bucket_status) = status_map.get_mut(&opts.bucket) {
bucket_status
@@ -442,13 +485,30 @@ impl ReplicationResyncer {
bucket_status.targets_map.get_mut(&opts.arn).unwrap()
};
if !resync_state_accepts_update(state, &opts) {
warn!(
bucket = %opts.bucket,
arn = %opts.arn,
incoming_resync_id = %opts.resync_id,
current_resync_id = %state.resync_id,
"ignoring stale resync stats update"
);
return;
}
if state.resync_id.is_empty() {
state.resync_id = opts.resync_id.clone();
}
if state.bucket.is_empty() {
state.bucket = opts.bucket.clone();
}
state.object = status.object.clone();
state.replicated_count += status.replicated_count;
state.replicated_size += status.replicated_size;
state.failed_count += status.failed_count;
state.failed_size += status.failed_size;
state.last_update = Some(OffsetDateTime::now_utc());
bucket_status.last_update = Some(OffsetDateTime::now_utc());
state.last_update = Some(now);
bucket_status.last_update = Some(now);
}
pub async fn persist_to_disk<S: StorageAPI>(&self, cancel_token: CancellationToken, api: Arc<S>) {
@@ -640,7 +700,6 @@ impl ReplicationResyncer {
let cancel_token = cancellation_token.clone();
let target_client = target_client.clone();
let resync_cancel_rx = self.resync_cancel_rx.clone();
let storage = storage.clone();
let results_tx = results_tx.clone();
let bucket_name = opts.bucket.clone();
@@ -714,10 +773,6 @@ impl ReplicationResyncer {
err,
);
if resync_cancel_rx.is_cancelled() {
return;
}
if cancel_token.is_cancelled() {
return;
}
@@ -731,8 +786,6 @@ impl ReplicationResyncer {
futures.push(f);
}
let resync_cancel_rx = self.resync_cancel_rx.clone();
while let Some(res) = rx.recv().await {
if let Some(err) = res.err {
error!("Failed to get object info: {}", err);
@@ -741,14 +794,8 @@ impl ReplicationResyncer {
return;
}
if resync_cancel_rx.is_cancelled() {
self.resync_bucket_mark_status(ResyncStatusType::ResyncCanceled, opts.clone(), storage.clone())
.await;
return;
}
if cancellation_token.is_cancelled() {
self.resync_bucket_mark_status(ResyncStatusType::ResyncFailed, opts.clone(), storage.clone())
self.resync_bucket_mark_status(ResyncStatusType::ResyncCanceled, opts.clone(), storage.clone())
.await;
return;
}
@@ -770,14 +817,8 @@ impl ReplicationResyncer {
continue;
}
if resync_cancel_rx.is_cancelled() {
self.resync_bucket_mark_status(ResyncStatusType::ResyncCanceled, opts.clone(), storage.clone())
.await;
return;
}
if cancellation_token.is_cancelled() {
self.resync_bucket_mark_status(ResyncStatusType::ResyncFailed, opts.clone(), storage.clone())
self.resync_bucket_mark_status(ResyncStatusType::ResyncCanceled, opts.clone(), storage.clone())
.await;
return;
}
@@ -3430,4 +3471,54 @@ mod tests {
"With no replication config, dsc may be empty; with config, replicate_any() would be true and queueing would occur"
);
}
#[tokio::test]
async fn test_cancel_marks_only_matching_bucket_target_token() {
let resyncer = ReplicationResyncer::new().await;
let opts_a = ResyncOpts {
bucket: "bucket-a".to_string(),
arn: "arn:replication::a".to_string(),
resync_id: "rid-a".to_string(),
resync_before: None,
};
let opts_b = ResyncOpts {
bucket: "bucket-b".to_string(),
arn: "arn:replication::b".to_string(),
resync_id: "rid-b".to_string(),
resync_before: None,
};
let token_a = CancellationToken::new();
let token_b = CancellationToken::new();
resyncer.register_cancel_token(&opts_a, token_a.clone()).await;
resyncer.register_cancel_token(&opts_b, token_b.clone()).await;
resyncer.cancel(&opts_a).await;
assert!(token_a.is_cancelled());
assert!(!token_b.is_cancelled());
}
#[test]
fn test_resync_state_accepts_update_only_for_matching_run() {
let current = TargetReplicationResyncStatus {
resync_id: "run-new".to_string(),
..Default::default()
};
let matching = ResyncOpts {
bucket: "bucket".to_string(),
arn: "arn:replication::dest".to_string(),
resync_id: "run-new".to_string(),
resync_before: None,
};
let stale = ResyncOpts {
bucket: "bucket".to_string(),
arn: "arn:replication::dest".to_string(),
resync_id: "run-old".to_string(),
resync_before: None,
};
assert!(resync_state_accepts_update(&TargetReplicationResyncStatus::default(), &matching));
assert!(resync_state_accepts_update(&current, &matching));
assert!(!resync_state_accepts_update(&current, &stale));
}
}
+2 -2
View File
@@ -17,7 +17,7 @@ use serde::Deserializer;
use serde::Serialize;
use time::OffsetDateTime;
#[derive(Debug, Serialize, Default, PartialEq, Eq)]
#[derive(Debug, Clone, Serialize, Default, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum GroupStatus {
#[default]
@@ -39,7 +39,7 @@ impl<'de> Deserialize<'de> for GroupStatus {
}
}
#[derive(Debug, Serialize, Deserialize, Default)]
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct GroupAddRemove {
pub group: String,
pub members: Vec<String>,
+2
View File
@@ -20,6 +20,7 @@ pub mod metrics;
pub mod net;
pub mod policy;
pub mod service_commands;
pub mod site_replication;
pub mod trace;
pub mod user;
pub mod utils;
@@ -27,4 +28,5 @@ pub mod utils;
pub use group::*;
pub use info_commands::*;
pub use policy::*;
pub use site_replication::*;
pub use user::*;
File diff suppressed because it is too large Load Diff
+2 -2
View File
@@ -21,7 +21,7 @@ use time::format_description::well_known::Rfc3339;
use crate::BackendInfo;
#[derive(Debug, Serialize, Deserialize, Default, PartialEq, Eq)]
#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
pub enum AccountStatus {
#[serde(rename = "enabled")]
Enabled,
@@ -94,7 +94,7 @@ pub struct UserInfo {
pub updated_at: Option<OffsetDateTime>,
}
#[derive(Debug, Serialize, Deserialize)]
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AddOrUpdateUserReq {
#[serde(rename = "secretKey")]
pub secret_key: String,