test(admin): box remote target repair scenarios (#7277)

This commit is contained in:
Zhengchao An
2026-09-06 15:23:30 +08:00
committed by GitHub
parent 27d593fa35
commit c99efd9477
+148 -139
View File
@@ -3146,64 +3146,67 @@ mod target_repair_tests {
#[tokio::test]
#[serial_test::serial]
async fn repair_existing_cached_target_persists_readable_targets_and_lists() {
temp_env::async_with_vars(TARGET_REPAIR_ENV, async {
let (_temp, env) = test_env().await;
let server = RemoteTargetServer::start().await;
let mut target = server.target();
target.arn = "arn:rustfs:replication:us-east-1:cached:remote".to_string();
let targets = BucketTargets {
targets: vec![target.clone()],
};
metadata_sys::update(BUCKET, BUCKET_TARGETS_FILE, serde_json::to_vec(&targets).expect("encode cached target"))
.await
.expect("seed cached target");
seed_unreadable(&env).await;
assert_eq!(
BucketTargetSys::get().get_remote_arn(BUCKET, Some(&target), "").await,
(target.arn.clone(), true)
);
assert!(
BucketTargetSys::get()
.get_remote_target_client(BUCKET, &target.arn)
temp_env::async_with_vars(
TARGET_REPAIR_ENV,
Box::pin(async {
let (_temp, env) = test_env().await;
let server = RemoteTargetServer::start().await;
let mut target = server.target();
target.arn = "arn:rustfs:replication:us-east-1:cached:remote".to_string();
let targets = BucketTargets {
targets: vec![target.clone()],
};
metadata_sys::update(BUCKET, BUCKET_TARGETS_FILE, serde_json::to_vec(&targets).expect("encode cached target"))
.await
.is_some()
);
.expect("seed cached target");
seed_unreadable(&env).await;
assert_eq!(
BucketTargetSys::get().get_remote_arn(BUCKET, Some(&target), "").await,
(target.arn.clone(), true)
);
assert!(
BucketTargetSys::get()
.get_remote_target_client(BUCKET, &target.arn)
.await
.is_some()
);
let arn = repair(&target, "replace-unreadable=true").await.expect("repair must commit");
assert_ne!(arn, target.arn);
assert!(
BucketTargetSys::get()
.get_remote_target_client(BUCKET, &target.arn)
let arn = repair(&target, "replace-unreadable=true").await.expect("repair must commit");
assert_ne!(arn, target.arn);
assert!(
BucketTargetSys::get()
.get_remote_target_client(BUCKET, &target.arn)
.await
.is_none()
);
assert!(BucketTargetSys::get().get_remote_target_client(BUCKET, &arn).await.is_some());
let persisted = metadata_sys::get_config_from_disk(BUCKET)
.await
.is_none()
);
assert!(BucketTargetSys::get().get_remote_target_client(BUCKET, &arn).await.is_some());
let persisted = metadata_sys::get_config_from_disk(BUCKET)
.await
.expect("read repaired metadata");
assert!(!persisted.bucket_targets_unreadable());
let targets = persisted.bucket_target_config.expect("decode persisted repair");
assert_eq!(targets.targets.len(), 1);
assert_eq!(targets.targets[0].arn, arn);
assert_eq!(
targets.targets[0]
.credentials
.as_ref()
.expect("persist credentials")
.secret_key,
"remote-secret"
);
let list = ListRemoteTargetHandler {}
.call(request(Method::GET, "", Vec::new()), Params::new())
.await
.expect("list repaired targets");
assert_eq!(list.output.0, StatusCode::OK);
let listed: serde_json::Value =
serde_json::from_slice(&list.output.1.collect().await.expect("collect target list").to_bytes())
.expect("decode list");
assert_eq!(listed.as_array().expect("targets list").len(), 1);
assert_eq!(listed[0]["arn"], arn);
})
.expect("read repaired metadata");
assert!(!persisted.bucket_targets_unreadable());
let targets = persisted.bucket_target_config.expect("decode persisted repair");
assert_eq!(targets.targets.len(), 1);
assert_eq!(targets.targets[0].arn, arn);
assert_eq!(
targets.targets[0]
.credentials
.as_ref()
.expect("persist credentials")
.secret_key,
"remote-secret"
);
let list = ListRemoteTargetHandler {}
.call(request(Method::GET, "", Vec::new()), Params::new())
.await
.expect("list repaired targets");
assert_eq!(list.output.0, StatusCode::OK);
let listed: serde_json::Value =
serde_json::from_slice(&list.output.1.collect().await.expect("collect target list").to_bytes())
.expect("decode list");
assert_eq!(listed.as_array().expect("targets list").len(), 1);
assert_eq!(listed[0]["arn"], arn);
}),
)
.await;
}
@@ -3260,51 +3263,54 @@ mod target_repair_tests {
#[tokio::test]
#[serial_test::serial]
async fn repair_with_stale_unreadable_cache_preserves_another_committed_repair() {
temp_env::async_with_vars(TARGET_REPAIR_ENV, async {
let (_temp, env) = test_env().await;
let first = RemoteTargetServer::start().await;
let second = RemoteTargetServer::start().await;
seed_unreadable(&env).await;
let first_arn = repair(&first.target(), "replace-unreadable=true")
.await
.expect("commit first repair");
// Model another node which still retains the original unreadable
// verdict when it begins its repair after this commit.
BucketTargetSys::get().mark_targets_unreadable(BUCKET).await;
let second_arn = repair(&second.target(), "replace-unreadable=true")
.await
.expect("merge second repair");
let persisted = metadata_sys::get_config_from_disk(BUCKET).await.expect("read both repairs");
let targets = persisted.bucket_target_config.expect("decode both repairs");
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_eq!(
BucketTargetSys::get()
.list_bucket_targets(BUCKET)
temp_env::async_with_vars(
TARGET_REPAIR_ENV,
Box::pin(async {
let (_temp, env) = test_env().await;
let first = RemoteTargetServer::start().await;
let second = RemoteTargetServer::start().await;
seed_unreadable(&env).await;
let first_arn = repair(&first.target(), "replace-unreadable=true")
.await
.expect("published repair")
.targets
.len(),
2
);
assert_eq!(
repair(&first.target(), "replace-unreadable=true")
.expect("commit first repair");
// Model another node which still retains the original unreadable
// verdict when it begins its repair after this commit.
BucketTargetSys::get().mark_targets_unreadable(BUCKET).await;
let second_arn = repair(&second.target(), "replace-unreadable=true")
.await
.expect("idempotent repair"),
first_arn
);
assert_eq!(
metadata_sys::get_config_from_disk(BUCKET)
.await
.expect("read repeated repair")
.bucket_target_config
.expect("decode repeated repair")
.targets
.len(),
2
);
})
.expect("merge second repair");
let persisted = metadata_sys::get_config_from_disk(BUCKET).await.expect("read both repairs");
let targets = persisted.bucket_target_config.expect("decode both repairs");
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_eq!(
BucketTargetSys::get()
.list_bucket_targets(BUCKET)
.await
.expect("published repair")
.targets
.len(),
2
);
assert_eq!(
repair(&first.target(), "replace-unreadable=true")
.await
.expect("idempotent repair"),
first_arn
);
assert_eq!(
metadata_sys::get_config_from_disk(BUCKET)
.await
.expect("read repeated repair")
.bucket_target_config
.expect("decode repeated repair")
.targets
.len(),
2
);
}),
)
.await;
}
@@ -3408,48 +3414,51 @@ mod target_repair_tests {
#[tokio::test]
#[serial_test::serial]
async fn failed_repair_transaction_never_reports_a_successful_replacement() {
temp_env::async_with_vars(TARGET_REPAIR_ENV, async {
let (_temp, env) = test_env().await;
let server = RemoteTargetServer::start().await;
seed_unreadable(&env).await;
let target = server.target();
BucketTargetSys::get()
.validate_target(BUCKET, &target)
.await
.expect("remote validation must succeed before injecting the metadata failure");
let file = metadata_sys::get_config_from_disk(BUCKET)
.await
.expect("read source metadata")
.save_file_path();
// Keep the source versioning and unreadable-target caches intact,
// but make the transaction's fresh disk load fail.
let corrupt = b"invalid metadata envelope".to_vec();
env.put_object_bytes(".rustfs.sys", &file, corrupt.clone()).await;
assert!(metadata_sys::get_config_from_disk(BUCKET).await.is_err());
let log = tempfile::NamedTempFile::new().expect("create captured log");
let writer = log.reopen().expect("open captured log writer");
let subscriber = tracing_subscriber::fmt()
.with_ansi(false)
.without_time()
.with_writer(writer)
.finish();
let error = repair(&target, "replace-unreadable=true")
.with_subscriber(subscriber)
.await
.expect_err("repair must fail on the unreadable metadata envelope");
assert_eq!(error.code(), &S3ErrorCode::InternalError);
let lines = std::fs::read_to_string(log.path()).expect("read captured log");
assert!(
!lines.contains("unreadable_targets_replaced"),
"a failed transaction must not claim success: {lines}"
);
assert_eq!(
crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file)
temp_env::async_with_vars(
TARGET_REPAIR_ENV,
Box::pin(async {
let (_temp, env) = test_env().await;
let server = RemoteTargetServer::start().await;
seed_unreadable(&env).await;
let target = server.target();
BucketTargetSys::get()
.validate_target(BUCKET, &target)
.await
.expect("read failed repair bytes"),
corrupt
);
})
.expect("remote validation must succeed before injecting the metadata failure");
let file = metadata_sys::get_config_from_disk(BUCKET)
.await
.expect("read source metadata")
.save_file_path();
// Keep the source versioning and unreadable-target caches intact,
// but make the transaction's fresh disk load fail.
let corrupt = b"invalid metadata envelope".to_vec();
env.put_object_bytes(".rustfs.sys", &file, corrupt.clone()).await;
assert!(metadata_sys::get_config_from_disk(BUCKET).await.is_err());
let log = tempfile::NamedTempFile::new().expect("create captured log");
let writer = log.reopen().expect("open captured log writer");
let subscriber = tracing_subscriber::fmt()
.with_ansi(false)
.without_time()
.with_writer(writer)
.finish();
let error = repair(&target, "replace-unreadable=true")
.with_subscriber(subscriber)
.await
.expect_err("repair must fail on the unreadable metadata envelope");
assert_eq!(error.code(), &S3ErrorCode::InternalError);
let lines = std::fs::read_to_string(log.path()).expect("read captured log");
assert!(
!lines.contains("unreadable_targets_replaced"),
"a failed transaction must not claim success: {lines}"
);
assert_eq!(
crate::admin::storage_api::config::read_admin_config(Arc::clone(&env.ecstore), &file)
.await
.expect("read failed repair bytes"),
corrupt
);
}),
)
.await;
}