diff --git a/.config/nextest.toml b/.config/nextest.toml index 9e7351fa3..d685fd1da 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -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]] diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 14d1b1188..b79555201 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -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 { diff --git a/rustfs/src/app/lifecycle_transition_api_test.rs b/rustfs/src/app/lifecycle_transition_api_test.rs index cf784282b..cc4b3a9c7 100644 --- a/rustfs/src/app/lifecycle_transition_api_test.rs +++ b/rustfs/src/app/lifecycle_transition_api_test.rs @@ -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()) diff --git a/rustfs/tests/connect_heartbeat.rs b/rustfs/tests/connect_heartbeat.rs index 814b8ad71..f02d2679d 100644 --- a/rustfs/tests/connect_heartbeat.rs +++ b/rustfs/tests/connect_heartbeat.rs @@ -323,11 +323,14 @@ async fn wait_for( if predicate(¤t) { 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); diff --git a/rustfs/tests/connect_inventory.rs b/rustfs/tests/connect_inventory.rs index c1766d38c..292f5f5eb 100644 --- a/rustfs/tests/connect_inventory.rs +++ b/rustfs/tests/connect_inventory.rs @@ -308,11 +308,14 @@ async fn wait_for( if predicate(¤t) { 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]