diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index 9c0e485c0..e4ca6ec9d 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -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; }