test(heal): cover degraded MRF replay after deletion (#8249)

* test(heal): cover deleted versions during MRF replay

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* test(heal): exercise degraded MRF replay after delete

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* test(heal): route quorum errors through test storage API

Co-Authored-By: heihutu <heihutu@gmail.com>
Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
Hauser
2026-09-30 09:14:02 +08:00
committed by GitHub
parent 56468fd553
commit 296854c3d2
2 changed files with 194 additions and 53 deletions
+193 -53
View File
@@ -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;
}
+1
View File
@@ -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;