diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 95d791218..51eb7cf8e 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -625,6 +625,7 @@ impl SetDisks { } } + #[hotpath::measure(impl_type = "SetDisks")] pub(super) async fn get_object_info_fileinfo_and_quorum( &self, bucket: &str, diff --git a/crates/ecstore/tests/hotpath_cpu_attribution_contract_test.rs b/crates/ecstore/tests/hotpath_cpu_attribution_contract_test.rs index 38c82329b..00e91653e 100644 --- a/crates/ecstore/tests/hotpath_cpu_attribution_contract_test.rs +++ b/crates/ecstore/tests/hotpath_cpu_attribution_contract_test.rs @@ -73,7 +73,7 @@ fn inherent_hotpath_measurements_have_cpu_attribution_types() { for function in [ "read_version_optimized", "get_object_fileinfo", - "get_object_info_and_quorum", + "get_object_info_fileinfo_and_quorum", "get_object_with_fileinfo", "get_object_decode_reader_with_fileinfo", "build_codec_streaming_part_reader", diff --git a/crates/heal/tests/mrf_partial_write_test.rs b/crates/heal/tests/mrf_partial_write_test.rs index 75f2c9fde..26035987d 100644 --- a/crates/heal/tests/mrf_partial_write_test.rs +++ b/crates/heal/tests/mrf_partial_write_test.rs @@ -30,7 +30,7 @@ use tokio::io::AsyncReadExt; mod storage_api; use storage_api::endpoint_index::{EndpointServerPools, Endpoints, init_local_disks}; use storage_api::integration::{ - DiskAPI, DiskStore, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, RUSTFS_META_BUCKET, ReadOptions, + DiskAPI, DiskError, DiskStore, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, RUSTFS_META_BUCKET, ReadOptions, }; const SNAPSHOT_LIMIT: usize = 64 * 1024 * 1024; @@ -348,15 +348,12 @@ async fn partial_write_ec12_4_ack_rejoin_repairs_versions_and_delete_marker_inne ); assert!(snapshot_contains("new.bin").await); assert!(snapshot_contains("versioned.bin").await); - manager.stop().await.expect("test manager should stop"); let before_purge = inspect_local_committed_snapshot(SNAPSHOT_LIMIT) .await .expect("legacy checkpoint must validate") - .expect("unverified legacy responsibility must remain") - .payload() - .to_vec(); + .expect("unverified legacy responsibility must remain"); // Deleting an existing marker by VersionId removes a version; it - // must not create a new partial-write repair responsibility. + // must persist a purge so the returning members cannot resurrect it. for path in &env.disk_paths[12..] { tokio::fs::rename(path, path.with_extension("offline")) .await @@ -380,14 +377,14 @@ async fn partial_write_ec12_4_ack_rejoin_repairs_versions_and_delete_marker_inne ) .await .expect("explicit marker purge should retain quorum"); - assert_eq!( + assert!( inspect_local_committed_snapshot(SNAPSHOT_LIMIT) .await .expect("purge checkpoint must validate") - .expect("existing legacy checkpoint must remain") - .payload(), - before_purge, - "a physical marker purge must not add a marker creation repair" + .expect("purge responsibility must be committed") + .sequence() + > before_purge.sequence(), + "a degraded marker purge must commit a successor checkpoint" ); for (path, disk) in env.disk_paths[12..].iter().zip(&all[12..]) { tokio::fs::remove_file(path).await.expect("remove purge outage sentinel"); @@ -397,6 +394,23 @@ async fn partial_write_ec12_4_ack_rejoin_repairs_versions_and_delete_marker_inne disk.reset_health_for_store_init_retry(); } *set.disks.write().await = all.iter().cloned().map(Some).collect(); + assert!( + wait_until(|| async { + let results = futures::future::join_all(all.iter().map(|disk| async { + disk.read_version("", "partial-versions", "versioned.bin", &marker, &ReadOptions::default()) + .await + })) + .await; + results + .iter() + .all(|result| matches!(result, Err(DiskError::FileVersionNotFound))) + }) + .await, + "MRF must complete the marker purge on the returning members" + ); + assert_payload(&env, "partial-versions", "versioned.bin", Some(&first), &payload1).await; + assert_payload(&env, "partial-versions", "versioned.bin", Some(&second), &payload2).await; + manager.stop().await.expect("test manager should stop"); }, ) .await; diff --git a/rustfs/tests/connect_perf_network.rs b/rustfs/tests/connect_perf_network.rs index bab33e923..53606369b 100644 --- a/rustfs/tests/connect_perf_network.rs +++ b/rustfs/tests/connect_perf_network.rs @@ -252,17 +252,18 @@ async fn controlled_peer_reports_exact_bytes_duration_latency_and_attributed_fai assert!(aggregate.get("peers").is_none(), "frozen aggregate schema has no per-peer field"); } -#[tokio::test] +#[tokio::test(start_paused = true)] async fn slow_peer_is_attributed_and_stops_at_the_wall_clock_limit() { let _guard = TEST_HARNESS_LOCK.lock().await; let mut request = request(1); request.duration = Duration::from_millis(20); request.traffic_bytes_per_peer = 1_000; - let started = Instant::now(); + let started = tokio::time::Instant::now(); let measurement = measure_network_with_harness(&request, &BlockingHarness, &CancellationToken::new()) .await .expect("typed timeout result"); assert!(started.elapsed() < Duration::from_millis(100)); + assert!(started.elapsed() >= request.duration); assert_eq!(measurement.result.outcome(), NetworkOutcome::Failed); assert_eq!(measurement.result.reason_code(), NetworkReasonCode::CollectionFailed); assert!(measurement.result.data().is_none());