diff --git a/crates/heal/tests/mrf_partial_write_test.rs b/crates/heal/tests/mrf_partial_write_test.rs index 6c8bbc430..440fb4bf1 100644 --- a/crates/heal/tests/mrf_partial_write_test.rs +++ b/crates/heal/tests/mrf_partial_write_test.rs @@ -30,7 +30,8 @@ use tokio::io::AsyncReadExt; mod storage_api; use storage_api::endpoint_index::{EndpointServerPools, Endpoints, init_local_disks}; use storage_api::integration::{ - DiskAPI, DiskError, DiskStore, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, RUSTFS_META_BUCKET, ReadOptions, + DiskAPI, DiskError, DiskStore, EcstoreStorageError, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, + RUSTFS_META_BUCKET, ReadOptions, }; const SNAPSHOT_LIMIT: usize = 64 * 1024 * 1024; @@ -91,7 +92,7 @@ async fn partial_write_persistence_failure_is_reported_and_retained_for_retry() } #[test] -fn unversioned_deleted_partial_write_is_discharged_by_an_absence_proof() { +fn degraded_deleted_partial_write_is_discharged_by_an_absence_proof() { const STACK_SIZE: usize = 8 * 1024 * 1024; std::thread::Builder::new() .name("mrf-partial-write-absence".to_owned()) @@ -102,67 +103,206 @@ fn unversioned_deleted_partial_write_is_discharged_by_an_absence_proof() { .enable_all() .build() .expect("partial-write absence runtime should build"); - runtime.block_on(unversioned_deleted_partial_write_is_discharged_by_an_absence_proof_inner()); + runtime.block_on(degraded_deleted_partial_write_is_discharged_by_an_absence_proof_inner()); }) .expect("partial-write absence test thread should spawn") .join() .expect("partial-write absence test thread should finish"); } -async fn unversioned_deleted_partial_write_is_discharged_by_an_absence_proof_inner() { - use rustfs_common::mrf_channel::{MrfScope, persist_partial_write_intent}; +async fn degraded_deleted_partial_write_is_discharged_by_an_absence_proof_inner() { + temp_env::async_with_vars( + [ + ("RUSTFS_HEAL_MRF_ENABLE", Some("true")), + ("RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE", Some("false")), + ("RUSTFS_SHARD_INTEGRITY_WRITE", Some("false")), + ("RUSTFS_SHARD_INTEGRITY_FLEET_CONFIRMED", Some("false")), + ], + async { + let root = tempfile::tempdir().expect("partial-write absence fixture directory"); + let env = TestECStoreEnv::builder().disk_count(16).base_dir(root.path()).build().await; + let set = env.ecstore.pools[0].get_disks(0); + assert_eq!(env.ecstore.pools[0].parity_count, 4, "fixture must be EC12+4"); + env.make_bucket("partial-absence", false).await; + env.make_bucket("partial-absence-versioned", true).await; + let mut coordinator_pool = env.endpoint_pools.as_ref()[0].clone(); + let mut endpoints = coordinator_pool.endpoints.as_ref().to_vec(); + for endpoint in endpoints.iter_mut().skip(4) { + endpoint.is_local = false; + } + coordinator_pool.endpoints = Endpoints::from(endpoints); + init_local_disks(EndpointServerPools::from(vec![coordinator_pool])) + .await + .expect("coordinator journal disks"); - temp_env::async_with_vars([("RUSTFS_HEAL_MRF_ENABLE", Some("true"))], async { - let root = tempfile::tempdir().expect("partial-write absence fixture directory"); - let env = TestECStoreEnv::builder().disk_count(16).base_dir(root.path()).build().await; - env.make_bucket("partial-absence", false).await; - let mut coordinator_pool = env.endpoint_pools.as_ref()[0].clone(); - let mut endpoints = coordinator_pool.endpoints.as_ref().to_vec(); - for endpoint in endpoints.iter_mut().skip(4) { - endpoint.is_local = false; - } - coordinator_pool.endpoints = Endpoints::from(endpoints); - init_local_disks(EndpointServerPools::from(vec![coordinator_pool])) - .await - .expect("coordinator journal disks"); + let manager = manager(&env); + mrf_queue::spawn_mrf_consumer(manager.clone()); + let all: Vec<_> = set + .disks + .read() + .await + .iter() + .map(|disk| disk.clone().expect("all sixteen members start online")) + .collect(); + for path in &env.disk_paths[12..] { + tokio::fs::rename(path, path.with_extension("offline")) + .await + .expect("detach the unavailable test node"); + tokio::fs::write(path, b"offline member") + .await + .expect("prevent the endpoint monitor from reopening the node"); + } + for disk in &mut set.disks.write().await[12..] { + *disk = None; + } - let manager = manager(&env); - mrf_queue::spawn_mrf_consumer(manager.clone()); - for (object, version_id) in [ - ("deleted-unversioned.bin", None), - ("deleted-versioned.bin", Some(uuid::Uuid::new_v4())), - ] { - persist_partial_write_intent( - "partial-absence", - object, - version_id, - MrfScope { - pool_index: 0, - set_index: 0, - }, - ) - .await - .expect("durable partial-write responsibility must commit before scheduling"); - assert!(snapshot_contains(object).await, "committed responsibility must exist before repair runs"); - } - - manager.start().await.expect("MRF scheduler should start"); - for object in ["deleted-unversioned.bin", "deleted-versioned.bin"] { + let deleted_object = "flink-key/.incomplete/upload-unversioned/part-1"; + put(&env, "partial-absence", deleted_object, b"object to delete", false).await; assert!( - wait_until(|| async { !snapshot_contains(object).await }).await, - "complete absence proof must discharge the durable responsibility" + snapshot_contains(deleted_object).await, + "a degraded PUT must durably admit its partial-write responsibility" ); - } - assert!( - wait_until(|| async { - let snapshot = manager.operations_snapshot().await; - snapshot.queue_length == 0 && snapshot.active_tasks == 0 - }) - .await, - "discharged absence repair must leave no queued work" - ); - manager.stop().await.expect("absence manager should stop"); - }) + let versioned_bucket = "partial-absence-versioned"; + let versioned_object = "flink-key/.incomplete/upload-versioned/part-1"; + let version_id = put(&env, versioned_bucket, versioned_object, b"version to delete", true) + .await + .expect("degraded versioned PUT must return a version ID"); + assert!( + snapshot_contains(versioned_object).await, + "a degraded versioned PUT must durably admit its partial-write responsibility" + ); + assert_eq!(replicas(&all, "partial-absence", deleted_object, None, false).await, 12); + assert_eq!(replicas(&all, versioned_bucket, versioned_object, Some(&version_id), false).await, 12); + + for path in &env.disk_paths[8..12] { + tokio::fs::rename(path, path.with_extension("offline")) + .await + .expect("detach a second unavailable test node"); + tokio::fs::write(path, b"offline member") + .await + .expect("prevent the endpoint monitor from reopening the second node"); + } + for disk in &mut set.disks.write().await[8..12] { + *disk = None; + } + let failed_object = "flink-key/.incomplete/upload-under-quorum/part-1"; + let mut failed_reader = PutObjReader::from_vec(b"write below quorum".to_vec()); + let failed_put = env + .ecstore + .put_object("partial-absence", failed_object, &mut failed_reader, &Default::default()) + .await; + let failed_put = failed_put.expect_err("two unavailable nodes must reject this write"); + assert!( + matches!( + &failed_put, + EcstoreStorageError::ErasureWriteQuorum | EcstoreStorageError::InsufficientWriteQuorum(_, _) + ), + "the rejected write must report a write-quorum failure: {failed_put:?}" + ); + assert!( + !snapshot_contains(failed_object).await, + "a subquorum PUT that was never committed must not create MRF responsibility" + ); + + for (path, disk) in env.disk_paths[12..].iter().zip(&all[12..]) { + tokio::fs::remove_file(path).await.expect("remove offline sentinel"); + tokio::fs::rename(path.with_extension("offline"), path) + .await + .expect("restore the same member data"); + disk.reset_health_for_store_init_retry(); + } + for (path, disk) in env.disk_paths[8..12].iter().zip(&all[8..12]) { + tokio::fs::remove_file(path) + .await + .expect("remove second-node offline sentinel"); + tokio::fs::rename(path.with_extension("offline"), path) + .await + .expect("restore the second node's original member data"); + disk.reset_health_for_store_init_retry(); + } + *set.disks.write().await = all.iter().cloned().map(Some).collect(); + assert!( + env.ecstore + .get_object_info("partial-absence", failed_object, &ObjectOptions::default()) + .await + .is_err(), + "the rejected subquorum PUT must not become a visible object after rejoin" + ); + + env.ecstore + .delete_object("partial-absence", deleted_object, ObjectOptions::default()) + .await + .expect("delete the exact object after its durable heal responsibility commits"); + assert!( + env.ecstore + .get_object_info("partial-absence", deleted_object, &ObjectOptions::default()) + .await + .is_err(), + "the deleted key must be absent before replay" + ); + env.ecstore + .delete_object( + versioned_bucket, + versioned_object, + ObjectOptions { + versioned: true, + version_id: Some(version_id.clone()), + ..Default::default() + }, + ) + .await + .expect("delete the exact version after its durable heal responsibility commits"); + let deleted_version_options = ObjectOptions { + versioned: true, + version_id: Some(version_id), + ..Default::default() + }; + assert!( + env.ecstore + .get_object_info(versioned_bucket, versioned_object, &deleted_version_options) + .await + .is_err(), + "the deleted version must be absent before replay" + ); + for object in [deleted_object, versioned_object] { + assert!( + snapshot_contains(object).await, + "deletion must not release the pending MRF responsibility" + ); + } + + manager.start().await.expect("MRF scheduler should start"); + for object in [deleted_object, versioned_object] { + assert!( + wait_until(|| async { !snapshot_contains(object).await }).await, + "complete absence proof must discharge the durable responsibility" + ); + } + assert!( + env.ecstore + .get_object_info("partial-absence", deleted_object, &ObjectOptions::default()) + .await + .is_err(), + "absence-proof replay must not recreate the deleted object" + ); + assert!( + env.ecstore + .get_object_info(versioned_bucket, versioned_object, &deleted_version_options) + .await + .is_err(), + "absence-proof replay must not recreate the deleted version" + ); + assert!( + wait_until(|| async { + let snapshot = manager.operations_snapshot().await; + snapshot.queue_length == 0 && snapshot.active_tasks == 0 + }) + .await, + "discharged absence repair must leave no queued work" + ); + manager.stop().await.expect("absence manager should stop"); + }, + ) .await; } diff --git a/crates/heal/tests/storage_api.rs b/crates/heal/tests/storage_api.rs index d35e48a4f..4edf01007 100644 --- a/crates/heal/tests/storage_api.rs +++ b/crates/heal/tests/storage_api.rs @@ -23,6 +23,7 @@ pub(crate) mod integration { pub(crate) use rustfs_ecstore::api::disk::{ DiskAPI, DiskError, DiskOption, DiskStore, Endpoint, RUSTFS_META_BUCKET, ReadOptions, new_disk, }; + pub(crate) use rustfs_ecstore::api::error::StorageError as EcstoreStorageError; pub(crate) use rustfs_ecstore::api::object::{ObjectOptions, PutObjReader, ShardIntegrityWriteMode, WriteCompletion}; pub(crate) use rustfs_ecstore::api::storage::ECStore; pub(crate) use rustfs_storage_api::BucketOperations;