test: stabilize heal timing and improve heartbeat replay coverage (#8150)

fix: preserve heartbeat replay and stabilize timing tests
This commit is contained in:
Chris
2026-09-28 22:07:15 +08:00
committed by GitHub
parent 4d782141d8
commit 5ebc1e480b
5 changed files with 79 additions and 36 deletions
+12
View File
@@ -139,6 +139,12 @@ test-group = 'embedded-test-ports'
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)))'
threads-required = "num-test-threads"
# Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129).
# Their existing deadlines, assertions, and retry policy remain in force.
[[profile.default.overrides]]
filter = 'package(rustfs) & binary(connect_heartbeat) & test(/^(restart_replays_pending_request_then_advances_sequence|restart_replays_exact_legacy_capabilities_with_and_without_jobs)$/)'
threads-required = "num-test-threads"
# Serialize the durable manual-transition checkpoint test across nextest's
# process boundary; it mutates bucket lifecycle metadata and is not quarantined.
[[profile.default.overrides]]
@@ -338,6 +344,12 @@ test-group = 'embedded-test-ports'
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)))'
threads-required = "num-test-threads"
# Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129).
# Their existing deadlines, assertions, and retry policy remain in force.
[[profile.ci.overrides]]
filter = 'package(rustfs) & binary(connect_heartbeat) & test(/^(restart_replays_pending_request_then_advances_sequence|restart_replays_exact_legacy_capabilities_with_and_without_jobs)$/)'
threads-required = "num-test-threads"
# Serialize the durable manual-transition checkpoint test under the ci profile
# too. No retries: failures stay visible.
[[profile.ci.overrides]]
+33 -28
View File
@@ -2678,7 +2678,7 @@ async fn read_repair_object_heal_sets_read_repair_option() {
assert!(!opts[0].no_lock);
}
#[tokio::test]
#[tokio::test(start_paused = true)]
async fn read_repair_object_heal_is_not_failed_by_flat_task_timeout() {
let storage = Arc::new(MockStorage {
block_heal_object: Mutex::new(true),
@@ -2687,31 +2687,23 @@ async fn read_repair_object_heal_is_not_failed_by_flat_task_timeout() {
let mut request = HealRequest::object("bucket".to_string(), "object".to_string(), None);
request.source = HealRequestSource::ReadRepair;
request.options.timeout = Some(Duration::from_millis(1));
let task = Arc::new(HealTask::from_request(request, storage.clone()));
let execution = tokio::spawn({
let task = task.clone();
async move { task.execute().await }
});
let task = HealTask::from_request(request, storage.clone());
// Isolate the storage boundary from preflight's separately enforced budget.
let execution = task.heal_object("bucket", "object", None);
tokio::pin!(execution);
assert!(futures::poll!(&mut execution).is_pending());
assert!(storage.object_heal_opts.lock().expect("heal options")[0].read_repair);
tokio::time::timeout(Duration::from_secs(1), async {
loop {
if !storage.object_heal_opts.lock().unwrap().is_empty() {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("read-repair object heal should start");
tokio::time::sleep(Duration::from_millis(20)).await;
assert!(!execution.is_finished(), "read repair must not be failed by the flat task timeout");
execution.abort();
assert!(execution.await.is_err(), "aborted mock execution should not join successfully");
assert!(storage.object_heal_opts.lock().unwrap()[0].read_repair);
*task.task_start_instant.write().await = Some(Instant::now() - Duration::from_millis(2));
assert!(matches!(task.remaining_timeout().await, Err(Error::TaskTimeout)));
tokio::time::advance(Duration::from_millis(20)).await;
assert!(
futures::poll!(&mut execution).is_pending(),
"read repair must remain in storage after the flat task budget expires"
);
}
#[tokio::test]
#[tokio::test(start_paused = true)]
async fn non_read_repair_object_heal_still_uses_flat_timeout() {
let storage = Arc::new(MockStorage {
block_heal_object: Mutex::new(true),
@@ -2719,13 +2711,26 @@ async fn non_read_repair_object_heal_still_uses_flat_timeout() {
});
let mut request = HealRequest::object("bucket".to_string(), "object".to_string(), None);
request.options.timeout = Some(Duration::from_millis(1));
let task = HealTask::from_request(request, storage);
let task = HealTask::from_request(request, storage.clone());
let execution = task.heal_object("bucket", "object", None);
tokio::pin!(execution);
assert!(futures::poll!(&mut execution).is_pending());
assert!(!storage.object_heal_opts.lock().expect("heal options")[0].read_repair);
tokio::time::advance(Duration::from_millis(20)).await;
assert!(matches!(futures::poll!(&mut execution), std::task::Poll::Ready(Err(Error::TaskTimeout))));
}
let result = tokio::time::timeout(Duration::from_secs(1), task.execute())
.await
.expect("flat timeout should finish the task");
#[tokio::test]
async fn read_repair_preflight_still_rejects_an_expired_task_budget() {
let storage = Arc::new(MockStorage::default());
let mut request = HealRequest::object("bucket".to_string(), "object".to_string(), None);
request.source = HealRequestSource::ReadRepair;
request.options.timeout = Some(Duration::from_millis(1));
let task = HealTask::from_request(request, storage.clone());
*task.task_start_instant.write().await = Some(Instant::now() - Duration::from_millis(2));
assert!(matches!(result, Err(Error::TaskTimeout)));
assert!(matches!(task.heal_object("bucket", "object", None).await, Err(Error::TaskTimeout)));
assert!(storage.object_heal_opts.lock().expect("heal options").is_empty());
}
async fn make_resume_disk(temp: &TempDir) -> DiskStore {
@@ -1479,7 +1479,10 @@ async fn cancelled_transition_waiting_for_prepared_reader_cleans_remote() {
.transition_object(&transition_bucket, object, &transition_opts)
.await
});
put_barrier.wait_until_paused().await;
tokio::select! {
result = &mut transition => panic!("transition ended before reaching the tier PUT barrier: {result:?}"),
() = put_barrier.wait_until_paused() => {}
}
let prepared_reader = ecstore
.prepare_get_object_reader(bucket.as_str(), object, None, HeaderMap::new(), &ObjectOptions::default())
+25 -5
View File
@@ -323,11 +323,14 @@ async fn wait_for(
if predicate(&current) {
return current;
}
status.changed().await.expect("status channel");
status
.changed()
.await
.unwrap_or_else(|error| panic!("heartbeat status channel closed: {error}; last status: {:?}", *status.borrow()));
}
})
.await
.expect("heartbeat status timeout")
.unwrap_or_else(|error| panic!("heartbeat status timeout: {error}; last status: {:?}", *status.borrow()))
}
async fn assert_credential_failure(config: HeartbeatConfig, server: &TestServer, expected: &str) {
@@ -634,7 +637,13 @@ async fn restart_replays_pending_request_then_advances_sequence() {
}
})
.await
.expect("two heartbeats");
.unwrap_or_else(|error| {
panic!(
"two heartbeats: {error}; last status: {:?}; received: {}",
*runtime.status().borrow(),
second_server.seen.lock().expect("seen lock").len()
)
});
runtime.shutdown().await;
let seen = second_server.seen.lock().expect("seen lock");
@@ -694,14 +703,25 @@ fn pending_heartbeat_with_capabilities(capabilities: &[&str]) -> Value {
#[tokio::test]
async fn restart_replays_exact_legacy_capabilities_with_and_without_jobs() {
for (job_capable, service_memory) in [(false, false), (true, false), (false, true), (true, true)] {
for (job_capable, service_memory, health_service) in [
(false, false, false),
(true, false, false),
(false, true, false),
(true, true, false),
(false, false, true),
(true, false, true),
] {
let pki = TestPki::new();
let server = server(&pki, vec![Reply::ok("2026-08-22T01:02:03Z")]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let config = config(&temp, &pki, &server);
fs::create_dir_all(config.state_path.parent().expect("state directory")).expect("create state directory");
let pending = pending_heartbeat_with_capabilities(&legacy_capabilities(job_capable, service_memory));
let mut capabilities = legacy_capabilities(job_capable, service_memory);
if health_service {
capabilities.insert(if job_capable { 4 } else { 3 }, "health.check.service@1");
}
let pending = pending_heartbeat_with_capabilities(&capabilities);
let state = json!({"nextSequence": 0, "pending": pending});
fs::write(&config.state_path, serde_json::to_vec(&state).expect("heartbeat state JSON")).expect("write heartbeat state");
private_mode(&config.state_path);
+5 -2
View File
@@ -308,11 +308,14 @@ async fn wait_for(
if predicate(&current) {
return current;
}
status.changed().await.expect("status channel");
status
.changed()
.await
.unwrap_or_else(|error| panic!("inventory status channel closed: {error}; last status: {:?}", *status.borrow()));
}
})
.await
.expect("inventory status timeout")
.unwrap_or_else(|error| panic!("inventory status timeout: {error}; last status: {:?}", *status.borrow()))
}
#[test]