diff --git a/crates/e2e_test/src/kms/bucket_default_encryption_test.rs b/crates/e2e_test/src/kms/bucket_default_encryption_test.rs index fecba2b89..c0f0f0181 100644 --- a/crates/e2e_test/src/kms/bucket_default_encryption_test.rs +++ b/crates/e2e_test/src/kms/bucket_default_encryption_test.rs @@ -37,7 +37,7 @@ async fn test_bucket_default_sse_s3_put_object() -> Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box> { - let total_deadline = Duration::from_secs(5); + wait_for_kms_ready_with_timeout(base_url, access_key, secret_key, Duration::from_secs(5)).await +} + +async fn wait_for_kms_ready_with_timeout( + base_url: &str, + access_key: &str, + secret_key: &str, + total_deadline: Duration, +) -> Result<(), Box> { let start = tokio::time::Instant::now(); + let deadline = start + total_deadline; let mut backoff = Duration::from_millis(200); let max_backoff = Duration::from_secs(1); - let mut first_attempt = true; loop { - if !first_attempt { - if start.elapsed() >= total_deadline { - return Err("KMS failed to become ready within 5 seconds".into()); - } - sleep(backoff).await; - backoff = (backoff * 2).min(max_backoff); - } - first_attempt = false; - - match get_kms_status(base_url, access_key, secret_key).await { - Ok(status) => { - info!("KMS is ready (status: {})", status); - return Ok(()); - } - Err(e) => { - if start.elapsed() >= total_deadline { - return Err(format!("KMS did not become ready within 5 s: last error: {e}").into()); + match tokio::time::timeout_at(deadline, get_kms_status(base_url, access_key, secret_key)).await { + Ok(Ok(status)) => { + let backend_status = serde_json::from_str::(&status) + .ok() + .and_then(|value| value.get("backend_status")?.as_str().map(str::to_owned)); + if backend_status.as_deref() == Some("healthy") { + info!("KMS is ready (status: {})", status); + return Ok(()); } - warn!(error = %e, elapsed_ms = start.elapsed().as_millis() as u64, "KMS not ready yet, retrying…"); + warn!( + backend_status = backend_status.as_deref().unwrap_or("missing"), + elapsed_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX), + "KMS not ready yet, retrying…" + ); } + Ok(Err(e)) => { + let elapsed_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX); + warn!(error = %e, elapsed_ms, "KMS not ready yet, retrying…"); + } + Err(_) => return Err(format!("KMS failed to become ready within {} ms", total_deadline.as_millis()).into()), } + + let now = tokio::time::Instant::now(); + if now >= deadline { + return Err(format!("KMS failed to become ready within {} ms", total_deadline.as_millis()).into()); + } + sleep((now + backoff).min(deadline) - now).await; + backoff = (backoff * 2).min(max_backoff); + } +} + +#[cfg(test)] +mod readiness_tests { + use super::{wait_for_kms_ready, wait_for_kms_ready_with_timeout}; + use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }; + use std::time::Duration; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + #[tokio::test] + async fn kms_readiness_retries_http_success_until_backend_is_healthy() { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind readiness test server"); + let address = listener.local_addr().expect("read readiness test server address"); + let requests = Arc::new(AtomicUsize::new(0)); + let server_requests = Arc::clone(&requests); + let server = tokio::spawn(async move { + for backend_status in ["error", "healthy"] { + let (mut socket, _) = listener.accept().await.expect("accept readiness request"); + let mut request = Vec::new(); + let mut chunk = [0_u8; 1024]; + while !request.windows(4).any(|window| window == b"\r\n\r\n") { + let read = socket.read(&mut chunk).await.expect("read readiness request"); + if read == 0 { + break; + } + request.extend_from_slice(&chunk[..read]); + } + server_requests.fetch_add(1, Ordering::SeqCst); + + let body = format!(r#"{{"backend_status":"{backend_status}"}}"#); + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", + body.len() + ); + socket.write_all(response.as_bytes()).await.expect("write readiness response"); + } + }); + + wait_for_kms_ready(&format!("http://{address}"), "access-key", "secret-key") + .await + .expect("KMS should become ready after the healthy response"); + + let observed_requests = requests.load(Ordering::SeqCst); + server.abort(); + assert_eq!(observed_requests, 2, "an HTTP 200 unhealthy status must be retried"); + } + + #[tokio::test] + async fn kms_readiness_deadline_covers_a_stalled_status_request() { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind readiness test server"); + let address = listener.local_addr().expect("read readiness test server address"); + let server = tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept readiness request"); + let mut request = [0_u8; 1024]; + let _ = socket.read(&mut request).await.expect("read readiness request"); + std::future::pending::<()>().await; + }); + + let result = tokio::time::timeout( + Duration::from_secs(1), + wait_for_kms_ready_with_timeout(&format!("http://{address}"), "access-key", "secret-key", Duration::from_millis(50)), + ) + .await + .expect("readiness helper must enforce its own deadline"); + + server.abort(); + assert!(result.is_err(), "a stalled status request must not outlive the readiness deadline"); } } diff --git a/crates/e2e_test/src/kms/copy_object_self_copy_sse_test.rs b/crates/e2e_test/src/kms/copy_object_self_copy_sse_test.rs index 1cf19a565..85a013ab3 100644 --- a/crates/e2e_test/src/kms/copy_object_self_copy_sse_test.rs +++ b/crates/e2e_test/src/kms/copy_object_self_copy_sse_test.rs @@ -61,7 +61,7 @@ async fn test_metadata_replace_self_copy_of_sse_object_stays_decryptable() { ) .await .expect("failed to start RustFS with local KMS"); - tokio::time::sleep(tokio::time::Duration::from_secs(3)).await; + kms_env.wait_for_kms_ready().await.expect("KMS ready"); let client = kms_env.base_env.create_s3_client(); // Deliberately an UNVERSIONED bucket: that is the branch where the store layer can service @@ -160,7 +160,7 @@ async fn test_metadata_replace_self_copy_dropping_sse_rewrites_plaintext() { ) .await .expect("failed to start RustFS with local KMS"); - tokio::time::sleep(tokio::time::Duration::from_secs(3)).await; + kms_env.wait_for_kms_ready().await.expect("KMS ready"); let client = kms_env.base_env.create_s3_client(); // Unversioned, and deliberately WITHOUT a bucket default-encryption rule, so the copy below @@ -256,7 +256,7 @@ async fn test_metadata_replace_self_copy_under_bucket_default_sse_stays_decrypta ) .await .expect("failed to start RustFS with local KMS"); - tokio::time::sleep(tokio::time::Duration::from_secs(3)).await; + kms_env.wait_for_kms_ready().await.expect("KMS ready"); let client = kms_env.base_env.create_s3_client(); let bucket = "copy-object-self-copy-bucket-default-sse-test"; diff --git a/crates/e2e_test/src/kms/copy_object_version_restore_sse_test.rs b/crates/e2e_test/src/kms/copy_object_version_restore_sse_test.rs index 3241a217d..c7a572e93 100644 --- a/crates/e2e_test/src/kms/copy_object_version_restore_sse_test.rs +++ b/crates/e2e_test/src/kms/copy_object_version_restore_sse_test.rs @@ -56,7 +56,7 @@ async fn test_self_copy_of_historical_sse_s3_version_is_readable() { ) .await .expect("failed to start RustFS with local KMS"); - tokio::time::sleep(tokio::time::Duration::from_secs(3)).await; + kms_env.wait_for_kms_ready().await.expect("KMS ready"); let client = kms_env.base_env.create_s3_client(); let bucket = "copy-object-version-restore-sse-test"; diff --git a/crates/e2e_test/src/kms/encryption_metadata_test.rs b/crates/e2e_test/src/kms/encryption_metadata_test.rs index a316668f6..59501430b 100644 --- a/crates/e2e_test/src/kms/encryption_metadata_test.rs +++ b/crates/e2e_test/src/kms/encryption_metadata_test.rs @@ -87,7 +87,7 @@ async fn test_head_reports_managed_metadata_for_sse_s3() -> Result<(), Box Result<(), let mut kms_env = LocalKMSTestEnvironment::new().await?; let default_key_id = kms_env.start_rustfs_for_local_kms().await?; - tokio::time::sleep(tokio::time::Duration::from_secs(3)).await; + kms_env.wait_for_kms_ready().await?; let s3_client = kms_env.base_env.create_s3_client(); kms_env.base_env.create_test_bucket(TEST_BUCKET).await?; @@ -250,7 +250,7 @@ async fn test_multipart_upload_writes_encrypted_data() -> Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Box Result<(), Bo let mut kms_env = LocalKMSTestEnvironment::new().await?; let _default_key_id = kms_env.start_rustfs_for_local_kms().await?; - tokio::time::sleep(tokio::time::Duration::from_secs(3)).await; + kms_env.wait_for_kms_ready().await?; let s3_client = kms_env.base_env.create_s3_client(); kms_env.base_env.create_test_bucket(TEST_BUCKET).await?; @@ -187,7 +187,7 @@ async fn test_step3_multipart_upload_with_sse_s3() -> Result<(), Box Result<(), Box Result<(), Box