test(e2e): synchronize disk restore and heal log assertions

This commit is contained in:
overtrue
2026-09-09 13:27:18 +08:00
parent 6920abfe29
commit 81038fc4e2
4 changed files with 84 additions and 19 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=cca6d0bc1487f472dc354bbffcc2b5c7410edfccb499c6bce6638113fbe6bbec
sha256-linux=3b71936f6f4ea0cca3b5db6c2f387c990dddda7e96c982315eb42e3b2e6c92b8
sha256-darwin=f0d15f2d1183be319d9d977d20b48bb05a29a424c266eed910c67aaa6f1ff955
sha256-linux=340aa702576ebed5266b7c47e591f267a11178fa04bf3828d64f53fe92eb0907
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=83a7dcaffd5a789517ae9f02a224f66a9713937885cff96fca2ad7e216f197ae
sha256-linux=626c10f8c964507ff987b6c86069e9019dc6d2ae7fb02db9be5df5aa8cc5145b
sha256-darwin=12d30fff5ed48fe95bbfb310dd507048f81782b2954e00eb72435708b1133f9c
sha256-linux=e917b2fdb303d01e6008ac4f1836c7b268698afe8383bd2aa4e697061df85cf6
@@ -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
@@ -385,6 +385,79 @@ mod tests {
)
}
async fn wait_for_admin_cluster_start_log(
log_path: &Path,
client_token: &str,
deadline: Instant,
) -> Result<(), Box<dyn Error + Send + Sync>> {
loop {
let coordinator_log = std::fs::read_to_string(log_path)?;
if coordinator_log
.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(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<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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::<serde_json::Value>(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(),