diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index 2ac994976..a0b32f596 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=5de6852b2e606eca4ee0d86dbb3b725a3eb8923143fd00df96c549acc36040bd +sha256=7e233cd838efd2688a9cf11ba23ad60297172d124067e91f59502d10c9633464 diff --git a/.config/nextest.toml b/.config/nextest.toml index 70efc267d..b75177486 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -142,11 +142,11 @@ test-group = 'ecstore-serial-flaky' filter = 'package(rustfs) & (binary(/^embedded.*_test$/) | binary(admin_diagnostic_capability_e2e))' test-group = 'embedded-test-ports' -# Inventory delivery and real drive/object probes include durable filesystem IO -# in short deadlines. Reserve capacity so unrelated storage fixtures cannot -# exhaust those budgets. Keep the deadlines and assertions unchanged. +# Registration, inventory delivery, and real drive/object probes include +# durable filesystem IO in short deadlines. Reserve capacity so sibling test +# processes cannot exhaust those budgets. Keep deadlines and assertions intact. [[profile.default.overrides]] -filter = 'package(rustfs) & (binary(connect_inventory) | binary(connect_perf_drive) | (binary(connect_perf_object) & test(=real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put)))' +filter = 'package(rustfs) & (binary(/^connect_registration$/) | binary(connect_inventory) | binary(connect_perf_drive) | (binary(connect_perf_object) & test(=real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put)))' threads-required = "num-test-threads" # These concurrent state writers keep the production 5s lock-acquisition limit @@ -363,7 +363,7 @@ filter = 'package(rustfs) & (binary(/^embedded.*_test$/) | binary(admin_diagnost test-group = 'embedded-test-ports' [[profile.ci.overrides]] -filter = 'package(rustfs) & (binary(connect_inventory) | binary(connect_perf_drive) | (binary(connect_perf_object) & test(=real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put)))' +filter = 'package(rustfs) & (binary(/^connect_registration$/) | binary(connect_inventory) | binary(connect_perf_drive) | (binary(connect_perf_object) & test(=real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put)))' threads-required = "num-test-threads" # Keep the same state-writer capacity reservation in CI without changing the @@ -509,6 +509,7 @@ default-filter = """ | test(/^replication_extension_test::(test_replication_check_succeeds_with_remote_target|test_replication_check_rejects_target_without_object_lock|test_set_remote_target_rejects_unversioned_source_bucket|test_replication_check_rejects_unversioned_source_bucket|test_replication_check_rejects_missing_replication_config|test_replication_check_rejects_invalid_bucket|test_set_remote_target_rejects_same_bucket_on_same_deployment|test_set_remote_target_rejects_unversioned_target_bucket|test_set_remote_target_update_requires_arn|test_set_remote_target_update_rejects_missing_target|test_set_remote_target_rejects_invalid_target_url|test_set_remote_target_rejects_self_signed_https_target_without_skip_tls_verify|test_set_remote_target_rejects_private_ca_https_target_without_ca_cert_pem|test_list_remote_targets_rejects_empty_bucket|test_list_remote_targets_rejects_invalid_bucket|test_remove_remote_target_rejects_missing_target|test_remove_remote_target_rejects_missing_arn|test_remove_remote_target_rejects_invalid_bucket|test_remove_remote_target_rejects_target_used_by_replication|test_delete_bucket_replication_removes_remote_target)$/) | test(/^reliant::lifecycle::/) | test(/^reliant::tiering::/) + | test(/^upgrade_compatibility_test::upgrade_write_readiness_tests::/) | test(/^on_demand_migration::(get_basic_test::(get_miss_pulls_inline_and_serves_locally_afterwards|head_miss_answers_from_the_source_without_persisting)|interaction_test::test_odm_admin_config_is_redacted_and_status_counts_match_the_source)$/) ) """ @@ -733,6 +734,7 @@ default-filter = """ & !test(/^replication_extension_test::/) & !test(/^replication_target_matrix_test::/) & !test(/^on_demand_migration::(concurrency_test|fault_test|interop_test|real_source_test)::/) + & !test(/^upgrade_compatibility_test::upgrade_write_readiness_tests::/) """ fail-fast = false diff --git a/crates/e2e_test/src/upgrade_compatibility_test.rs b/crates/e2e_test/src/upgrade_compatibility_test.rs index 6b89f8148..c7b607282 100644 --- a/crates/e2e_test/src/upgrade_compatibility_test.rs +++ b/crates/e2e_test/src/upgrade_compatibility_test.rs @@ -395,6 +395,7 @@ async fn exercise_mixed_cluster( previous_node: usize, ) -> TestResult { let clients = cluster.create_all_clients()?; + wait_for_upgrade_write_readiness(&clients, phase, LISTING_CONVERGENCE_TIMEOUT).await?; let current_client = &clients[current_node]; let previous_client = &clients[previous_node]; @@ -441,6 +442,190 @@ async fn exercise_mixed_cluster( Ok(()) } +async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budget: Duration) -> TestResult { + // ListBuckets can succeed before peers recover a restarted disk. Cluster + // health also accepts Returning disks whose write health is still FAULTY. + // Probe every writer outside the asserted phase prefix; compatibility + // writes still execute once and retain their original assertions. + let deadline = Instant::now() + budget; + for (node, client) in clients.iter().enumerate() { + let client = Client::from_conf( + client + .config() + .to_builder() + .retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1)) + .build(), + ); + let key = format!(".upgrade-readiness/{phase}/node-{node}"); + let mut last_response = "no response".to_string(); + loop { + if Instant::now() >= deadline { + return Err( + format!("{phase}: node {node} write readiness deadline exceeded; last response: {last_response}").into(), + ); + } + let response = tokio::time::timeout_at( + deadline, + client + .put_object() + .bucket(MIXED_BUCKET) + .key(&key) + .body(ByteStream::from_static(b"upgrade write readiness")) + .send(), + ) + .await; + match response { + Ok(Ok(_)) => break, + Ok(Err(error)) => { + if error.raw_response().map(|response| response.status().as_u16()) != Some(503) + || error.as_service_error().and_then(ProvideErrorMetadata::code) != Some("ServiceUnavailable") + { + return Err(format!("{phase}: node {node} write readiness failed: {error:?}").into()); + } + last_response = format!("{error:?}"); + } + Err(_) => { + return Err(format!( + "{phase}: node {node} write readiness deadline exceeded during PutObject; last response: {last_response}" + ) + .into()); + } + } + tokio::time::sleep_until(deadline.min(Instant::now() + Duration::from_millis(500))).await; + } + } + Ok(()) +} + +#[cfg(test)] +mod upgrade_write_readiness_tests { + use super::*; + use crate::fake_s3_target::FaultAction; + + #[tokio::test] + async fn waits_for_each_writer_after_metadata_is_ready() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket(MIXED_BUCKET); + let client = fake_source_client(&target); + client.head_bucket().bucket(MIXED_BUCKET).send().await?; + + let phase = "one-previous-node"; + let first_key = format!(".upgrade-readiness/{phase}/node-0"); + let second_key = format!(".upgrade-readiness/{phase}/node-1"); + target.inject_for_key( + FakeTargetOperation::PutObject, + &first_key, + FaultAction::Status(StatusCode::SERVICE_UNAVAILABLE), + 2, + ); + target.inject_for_key( + FakeTargetOperation::PutObject, + &second_key, + FaultAction::Status(StatusCode::SERVICE_UNAVAILABLE), + 1, + ); + // Metadata readiness does not prove that a data write can succeed. + let premature = client + .put_object() + .bucket(MIXED_BUCKET) + .key(&first_key) + .body(ByteStream::from_static(b"upgrade write readiness")) + .send() + .await + .expect_err("metadata readiness does not prove write readiness"); + assert_eq!(premature.raw_response().map(|response| response.status().as_u16()), Some(503)); + + wait_for_upgrade_write_readiness(&[client.clone(), client.clone()], phase, Duration::from_secs(5)).await?; + assert_eq!(target.count_requests(FakeTargetOperation::PutObject, &first_key), 3); + assert_eq!(target.count_requests(FakeTargetOperation::PutObject, &second_key), 2); + assert!(target.has_object(MIXED_BUCKET, &first_key)); + assert!(target.has_object(MIXED_BUCKET, &second_key)); + assert!( + client + .list_objects_v2() + .bucket(MIXED_BUCKET) + .prefix(format!("{phase}/")) + .send() + .await? + .contents() + .is_empty() + ); + target.shutdown().await; + Ok(()) + } + + #[tokio::test] + async fn rejects_permanent_errors_without_sdk_retries() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket(MIXED_BUCKET); + let client = Client::from_conf( + fake_source_client(&target) + .config() + .to_builder() + .retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(3)) + .build(), + ); + for status in [ + StatusCode::INTERNAL_SERVER_ERROR, + StatusCode::FORBIDDEN, + StatusCode::NOT_FOUND, + ] { + let phase = format!("permanent-{}", status.as_u16()); + let key = format!(".upgrade-readiness/{phase}/node-0"); + target.inject_for_key(FakeTargetOperation::PutObject, &key, FaultAction::Status(status), 1); + let error = wait_for_upgrade_write_readiness(std::slice::from_ref(&client), &phase, Duration::from_secs(5)) + .await + .expect_err("a permanent error must not be retried into success"); + assert!(error.to_string().contains("node 0 write readiness failed"), "{error}"); + assert_eq!(target.count_requests(FakeTargetOperation::PutObject, &key), 1); + assert!(!target.has_object(MIXED_BUCKET, &key)); + } + target.shutdown().await; + Ok(()) + } + + #[tokio::test] + async fn transient_errors_stop_at_the_deadline() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket(MIXED_BUCKET); + let client = fake_source_client(&target); + client.head_bucket().bucket(MIXED_BUCKET).send().await?; + let phase = "deadline"; + let key = format!(".upgrade-readiness/{phase}/node-0"); + target.inject_for_key( + FakeTargetOperation::PutObject, + &key, + FaultAction::Status(StatusCode::SERVICE_UNAVAILABLE), + 10, + ); + let error = wait_for_upgrade_write_readiness(&[client], phase, Duration::from_secs(1)) + .await + .expect_err("persistent unavailability must exhaust the shared deadline"); + // Transport scheduling consumes the same budget; do not require a + // response to reach the fake target before the deadline on a busy host. + assert!(error.to_string().contains("deadline exceeded"), "{error}"); + target.shutdown().await; + Ok(()) + } + + #[tokio::test] + async fn in_flight_requests_are_bounded_by_the_deadline() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket(MIXED_BUCKET); + let client = fake_source_client(&target); + client.head_bucket().bucket(MIXED_BUCKET).send().await?; + let phase = "stalled"; + let key = format!(".upgrade-readiness/{phase}/node-0"); + target.inject_for_key(FakeTargetOperation::PutObject, &key, FaultAction::Stall(Duration::from_secs(30)), 1); + let error = wait_for_upgrade_write_readiness(&[client], phase, Duration::from_secs(1)) + .await + .expect_err("a stalled request must not outlive the readiness deadline"); + assert!(error.to_string().contains("deadline exceeded during PutObject"), "{error}"); + target.shutdown().await; + Ok(()) + } +} + /// Pins the published old writer's limitation and the supported recovery /// procedure. This is not a promise that mixed-version ODM is supported. /// Replace the loss assertion when ODM gains independent persistence; diff --git a/rustfs/tests/connect_registration.rs b/rustfs/tests/connect_registration.rs index bf957e3b2..2385dd960 100644 --- a/rustfs/tests/connect_registration.rs +++ b/rustfs/tests/connect_registration.rs @@ -881,11 +881,13 @@ async fn wait_for_heartbeat_status( if predicate(¤t) { return current; } - status.changed().await.expect("heartbeat status channel"); + status.changed().await.unwrap_or_else(|error| { + panic!("heartbeat status channel: {error}; last status: {:?}", *status.borrow()); + }); } }) .await - .expect("heartbeat status") + .unwrap_or_else(|error| panic!("heartbeat status: {error}; last status: {:?}", *status.borrow())) } fn rotation_response(pki: &TestPki, identity: &rustfs::connect::DeviceIdentity, serial: u8) -> (Value, Value) { @@ -1647,11 +1649,13 @@ async fn inventory_first_recovers_a_saved_reenrollment_before_telemetry() { if matches!(current, InventoryStatus::Online { .. }) { break current; } - inventory_status.changed().await.expect("inventory status channel"); + inventory_status.changed().await.unwrap_or_else(|error| { + panic!("inventory status channel: {error}; last status: {:?}", *inventory_status.borrow()); + }); } }) .await - .expect("inventory online status"), + .unwrap_or_else(|error| panic!("inventory online status: {error}; last status: {:?}", *inventory_status.borrow())), InventoryStatus::Online { .. } ));