diff --git a/.config/e2e-full-selection.txt b/.config/e2e-full-selection.txt index 395443805..ef0608575 100644 --- a/.config/e2e-full-selection.txt +++ b/.config/e2e-full-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=cca6d0bc1487f472dc354bbffcc2b5c7410edfccb499c6bce6638113fbe6bbec -sha256-linux=3b71936f6f4ea0cca3b5db6c2f387c990dddda7e96c982315eb42e3b2e6c92b8 +sha256-darwin=f0d15f2d1183be319d9d977d20b48bb05a29a424c266eed910c67aaa6f1ff955 +sha256-linux=340aa702576ebed5266b7c47e591f267a11178fa04bf3828d64f53fe92eb0907 diff --git a/.config/e2e-nightly-selection.txt b/.config/e2e-nightly-selection.txt index bec86f799..1fdba3204 100644 --- a/.config/e2e-nightly-selection.txt +++ b/.config/e2e-nightly-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=83a7dcaffd5a789517ae9f02a224f66a9713937885cff96fca2ad7e216f197ae -sha256-linux=626c10f8c964507ff987b6c86069e9019dc6d2ae7fb02db9be5df5aa8cc5145b +sha256-darwin=12d30fff5ed48fe95bbfb310dd507048f81782b2954e00eb72435708b1133f9c +sha256-linux=e917b2fdb303d01e6008ac4f1836c7b268698afe8383bd2aa4e697061df85cf6 diff --git a/crates/e2e_test/src/get_codec_streaming_compat_test.rs b/crates/e2e_test/src/get_codec_streaming_compat_test.rs index 4fdc8326f..3659af6d9 100644 --- a/crates/e2e_test/src/get_codec_streaming_compat_test.rs +++ b/crates/e2e_test/src/get_codec_streaming_compat_test.rs @@ -401,11 +401,11 @@ mod tests { let mp_view = get_full(&client, multipart_key).await?; assert_eq!(mp_view.sha256, sha256_hex(&multipart_body), "degraded baseline multipart body mismatch"); baseline_degraded.insert(multipart_key.to_string(), mp_view); - // Restore the disk so Phase B restarts from a clean, complete disk set. + // Stop disk writers before restoring the complete layout reused by Phase B. + harness.kill_server(); harness.bring_disk_online(0)?; // ---- Phase B: codec streaming (gates opened) ---- - harness.kill_server(); for (k, v) in codec_env() { harness.set_env(k, v); } @@ -497,6 +497,8 @@ mod tests { let mp_view = get_full(&client, multipart_key).await?; assert_eq!(mp_view.sha256, sha256_hex(&multipart_body), "degraded codec multipart body mismatch"); codec_degraded.insert(multipart_key.to_string(), mp_view); + // All server reads are complete; stop disk writers before restoring disk0. + harness.kill_server(); harness.bring_disk_online(0)?; // A/B under parity reconstruction: codec == legacy, byte-for-byte and diff --git a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs index c2f63155f..450bd6bd1 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -385,6 +385,79 @@ mod tests { ) } + async fn wait_for_admin_cluster_start_log( + log_path: &Path, + client_token: &str, + deadline: Instant, + ) -> Result<(), Box> { + loop { + let coordinator_log = std::fs::read_to_string(log_path)?; + if coordinator_log + .lines() + .filter_map(|line| serde_json::from_str::(line).ok()) + .any(|event| { + event["event"] == "heal_task_state" + && event["task_id"] == client_token + && event["heal_type"] == "cluster" + && event["state"] == "started" + }) + { + return Ok(()); + } + if Instant::now() >= deadline { + return Err( + format!("node 0 must have started the exact admin task before interruption: task_id={client_token}").into(), + ); + } + sleep(Duration::from_millis(10)).await; + } + } + + #[tokio::test] + async fn test_admin_cluster_start_log_waits_for_exact_delayed_event() -> Result<(), Box> { + use std::io::Write; + + let mut log = tempfile::NamedTempFile::new()?; + for (task_id, heal_type, state) in [ + ("other-task", "cluster", "started"), + ("admin-task", "object", "started"), + ("admin-task", "cluster", "completed"), + ] { + writeln!( + log, + "{}", + serde_json::json!({"event": "heal_task_state", "task_id": task_id, "heal_type": heal_type, "state": state}) + )?; + } + let log_path = log.path().to_path_buf(); + let started = wait_for_admin_cluster_start_log(&log_path, "admin-task", Instant::now() + Duration::from_secs(1)); + tokio::pin!(started); + // Poll the reader before publishing the start event, without depending + // on scheduling or a fixed writer delay to reproduce log visibility. + tokio::select! { + biased; + result = &mut started => panic!("unrelated events must leave the exact start pending: {result:?}"), + _ = std::future::ready(()) => {} + } + writeln!( + log, + "{}", + serde_json::json!({"event": "heal_task_state", "task_id": "admin-task", "heal_type": "cluster", "state": "started"}) + )?; + started.await?; + Ok(()) + } + + #[tokio::test] + async fn test_admin_cluster_start_log_respects_existing_deadline() -> Result<(), Box> { + let log = tempfile::NamedTempFile::new()?; + let error = wait_for_admin_cluster_start_log(log.path(), "admin-task", Instant::now()) + .await + .expect_err("missing exact start must fail at the supplied deadline"); + assert!(error.to_string().contains("task_id=admin-task"), "{error}"); + Ok(()) + } + fn cluster_heal_is_idle(status: &serde_json::Value) -> bool { let operations = &status["healOperations"]; status["clusterStatusComplete"] == serde_json::Value::Bool(true) @@ -1328,6 +1401,9 @@ mod tests { } sleep(Duration::from_millis(50)).await; } + // Task execution and its non-blocking log writer advance independently. + // Observe the exact start before taking the partial-rebuild snapshot. + wait_for_admin_cluster_start_log(&log_dir.join("node0.log"), client_token, partial_deadline).await?; let (partial_count, partial_manifest) = loop { // Hash one committed shard to prove progress without letting a // full-corpus hash pass consume the interruption window. @@ -1362,19 +1438,6 @@ mod tests { let pre_interrupt_status: serde_json::Value = serde_json::from_str(&pre_interrupt_status_body) .map_err(|err| format!("pre-interrupt background heal status is not JSON ({err}): {pre_interrupt_status_body}"))?; let pre_interrupt_replacement = replacement_recovery_status(&cluster).await?; - let coordinator_log = std::fs::read_to_string(log_dir.join("node0.log"))?; - assert!( - coordinator_log - .lines() - .filter_map(|line| serde_json::from_str::(line).ok()) - .any(|event| { - event["event"] == "heal_task_state" - && event["task_id"] == client_token - && event["heal_type"] == "cluster" - && event["state"] == "started" - }), - "node 0 must have started the exact admin task before interruption" - ); let pre_interrupt_operations = &pre_interrupt_status["healOperations"]; assert_eq!( pre_interrupt_operations["activeBySource"]["admin"].as_u64(),