fix(heal): finalize durable MRF legacy responsibility lifecycle (#8254)

## Related Issues

Resolves the corrected lifecycle fixture findings in PR #8254.

## Summary of Changes

Preserve separate incarnation-bound durable MRF responsibilities and their lifecycle audit records, and align replay fixtures with their checkpoint identities.

## Verification

Two independent source reviews and changed-delta reviews are complete. The final review approved 607a0b9bb0 with no remaining supported findings. Local runtime claims were not independently reproduced for this pull request.

## Impact

Keeps storage generation fences and retained responsibility semantics. Test-capacity reservations preserve existing deadlines and assertions.

## Additional Notes

Squash merge of the currently approved fix under the authorized CI-bypass exception. Main CI and release acceptance remain required.
This commit is contained in:
Hauser
2026-09-30 14:02:32 +08:00
committed by GitHub
parent 6bc2e10bc1
commit d60dfbb826
21 changed files with 3226 additions and 234 deletions
+18
View File
@@ -242,6 +242,13 @@ test-group = 'e2e-cluster-nightly'
filter = 'package(e2e_test) & (test(/^kms::kms_vault_test::/) | test(/^kms::kms_rekey_sweep_test::/) | test(/^kms::configured_roundtrip_test::test_configured_vault_kms_admin_and_versioned_cleanup$/))'
test-group = 'e2e-vault'
# These server-heavy smoke cases perform 1k+ sequential S3 writes or run the
# scanner and a remote tier server together. Reserve the full test budget so
# unrelated E2E processes cannot starve their storage and polling deadlines.
[[profile.default.overrides]]
filter = 'package(e2e_test) & (test(/^list_objects_v2_pagination_test::tests::(test_list_objects_v2_delimiter_small_page_traverses_all|test_list_objects_v2_max_keys_above_limit_returns_token|test_list_objects_v2_maxkeys_above_limit_with_delimiter)$/) | test(=reliant::tiering::test_hermetic_transition_restore_failure_expiry_and_retry))'
threads-required = "num-test-threads"
# This four-disk, 65-member rollback probe already drives up to 32 concurrent
# durable deletions. Reserve this nextest run's capacity for its progress oracle.
[[profile.default.overrides]]
@@ -525,6 +532,11 @@ path = "junit.xml"
[[profile.e2e-smoke.overrides]]
filter = 'package(e2e_test) & test(/^list_objects_v2_pagination_test::tests::(test_list_objects_v2_delimiter_small_page_traverses_all|test_list_objects_v2_max_keys_above_limit_returns_token|test_list_objects_v2_maxkeys_above_limit_with_delimiter)$/)'
slow-timeout = { period = "60s", terminate-after = 2, grace-period = "10s" }
threads-required = "num-test-threads"
[[profile.e2e-smoke.overrides]]
filter = 'package(e2e_test) & test(=reliant::tiering::test_hermetic_transition_restore_failure_expiry_and_retry)'
threads-required = "num-test-threads"
# ---------------------------------------------------------------------------
# e2e-repl-nightly profile — scheduled full replication e2e lane (repl-1)
@@ -763,3 +775,9 @@ test-group = 'e2e-cluster-nightly'
[[profile.e2e-full.overrides]]
filter = 'package(e2e_test) & (test(/^kms::kms_vault_test::/) | test(/^kms::kms_rekey_sweep_test::/) | test(/^kms::configured_roundtrip_test::test_configured_vault_kms_admin_and_versioned_cleanup$/))'
test-group = 'e2e-vault'
# Keep the storage-heavy smoke cases isolated in the full suite too. They
# start multiple server fixtures and issue many sequential S3 operations.
[[profile.e2e-full.overrides]]
filter = 'package(e2e_test) & (test(/^list_objects_v2_pagination_test::tests::(test_list_objects_v2_delimiter_small_page_traverses_all|test_list_objects_v2_max_keys_above_limit_returns_token|test_list_objects_v2_maxkeys_above_limit_with_delimiter)$/) | test(=reliant::tiering::test_hermetic_transition_restore_failure_expiry_and_retry))'
threads-required = "num-test-threads"
+33 -6
View File
@@ -333,6 +333,9 @@ pub enum MrfDurableAdmissionError {
/// cancel repair.
pub struct MrfDurableSubmission {
pub intent: MrfIntent,
/// Bucket incarnation observed by the committed-write caller, when
/// available. Old journal records and legacy callers remain unbound.
pub source_bucket_incarnation_id: Option<Uuid>,
pub response: oneshot::Sender<Result<(), MrfDurableAdmissionError>>,
}
@@ -353,6 +356,19 @@ pub async fn persist_partial_write_intent(
version_id: Option<Uuid>,
scope: MrfScope,
) -> Result<(), MrfDurableAdmissionError> {
persist_partial_write_intent_with_incarnation(bucket, object, version_id, scope, None).await
}
pub async fn persist_partial_write_intent_with_incarnation(
bucket: &str,
object: &str,
version_id: Option<Uuid>,
scope: MrfScope,
source_bucket_incarnation_id: Option<Uuid>,
) -> Result<(), MrfDurableAdmissionError> {
if source_bucket_incarnation_id.is_some_and(|incarnation| incarnation.is_nil()) {
return Err(MrfDurableAdmissionError::InvalidIdentity);
}
let intent = MrfIntent {
bucket: Arc::from(bucket),
object: Arc::from(object),
@@ -364,7 +380,7 @@ pub async fn persist_partial_write_intent(
enqueued_at_ms: unix_now_ms(),
attempts: 0,
};
persist_durable_intent(intent).await
persist_durable_intent(intent, source_bucket_incarnation_id).await
}
pub async fn persist_delete_marker_purge_intent(
@@ -377,6 +393,7 @@ pub async fn persist_delete_marker_purge_intent(
if version_id.is_nil() {
return Err(MrfDurableAdmissionError::InvalidIdentity);
}
let source_bucket_incarnation_id = delete_marker_purge.bucket_incarnation_id;
let intent = MrfIntent {
bucket: Arc::from(bucket),
object: Arc::from(object),
@@ -388,17 +405,23 @@ pub async fn persist_delete_marker_purge_intent(
enqueued_at_ms: unix_now_ms(),
attempts: 0,
};
persist_durable_intent_unconditionally(intent).await
persist_durable_intent_unconditionally(intent, Some(source_bucket_incarnation_id)).await
}
async fn persist_durable_intent(intent: MrfIntent) -> Result<(), MrfDurableAdmissionError> {
async fn persist_durable_intent(
intent: MrfIntent,
source_bucket_incarnation_id: Option<Uuid>,
) -> Result<(), MrfDurableAdmissionError> {
if !mrf_delivery_enabled() {
return Err(MrfDurableAdmissionError::Disabled);
}
persist_durable_intent_unconditionally(intent).await
persist_durable_intent_unconditionally(intent, source_bucket_incarnation_id).await
}
async fn persist_durable_intent_unconditionally(mut intent: MrfIntent) -> Result<(), MrfDurableAdmissionError> {
async fn persist_durable_intent_unconditionally(
mut intent: MrfIntent,
source_bucket_incarnation_id: Option<Uuid>,
) -> Result<(), MrfDurableAdmissionError> {
if intent.bucket.is_empty()
|| intent.object.is_empty()
|| intent.bucket.len() > MRF_MAX_IDENTITY_COMPONENT
@@ -413,7 +436,11 @@ async fn persist_durable_intent_unconditionally(mut intent: MrfIntent) -> Result
}
let (response, receipt) = oneshot::channel();
sender
.try_send(MrfDurableSubmission { intent, response })
.try_send(MrfDurableSubmission {
intent,
source_bucket_incarnation_id,
response,
})
.map_err(|err| match err {
mpsc::error::TrySendError::Full(_) => MrfDurableAdmissionError::Full,
mpsc::error::TrySendError::Closed(_) => MrfDurableAdmissionError::Unavailable,
+22 -5
View File
@@ -4339,9 +4339,15 @@ impl SetDisks {
}
}
pub(in crate::set_disk) async fn persist_partial_write(&self, bucket: &str, object: &str, version_id: Option<&str>) -> bool {
pub(in crate::set_disk) async fn persist_partial_write(
&self,
bucket: &str,
object: &str,
version_id: Option<&str>,
source_bucket_incarnation_id: Option<Uuid>,
) -> bool {
use rustfs_common::mrf_channel::{
MrfDurableAdmissionError, MrfScope, mrf_delivery_enabled, persist_partial_write_intent,
MrfDurableAdmissionError, MrfScope, mrf_delivery_enabled, persist_partial_write_intent_with_incarnation,
};
if !mrf_delivery_enabled() {
@@ -4360,7 +4366,9 @@ impl SetDisks {
Ok::<_, MrfDurableAdmissionError>((version, scope))
})();
let result = match identity {
Ok((version, scope)) => persist_partial_write_intent(bucket, object, version, scope).await,
Ok((version, scope)) => {
persist_partial_write_intent_with_incarnation(bucket, object, version, scope, source_bucket_incarnation_id).await
}
Err(err) => Err(err),
};
match result {
@@ -4396,15 +4404,24 @@ impl SetDisks {
pub(in crate::set_disk) async fn submit_rename_tail_heal(
&self,
request: rustfs_heal_contracts::heal_channel::HealChannelRequest,
mut request: rustfs_heal_contracts::heal_channel::HealChannelRequest,
) {
if let Some(object) = request.object_prefix.as_deref()
&& self
.persist_partial_write(&request.bucket, object, request.object_version_id.as_deref())
.persist_partial_write(
&request.bucket,
object,
request.object_version_id.as_deref(),
request.expected_bucket_incarnation_id,
)
.await
{
return;
}
if request.expected_bucket_incarnation_id.is_none() {
return;
}
request.source = rustfs_heal_contracts::heal_channel::HealRequestSource::Mrf;
#[cfg(test)]
{
let capture = self
+124 -52
View File
@@ -3500,6 +3500,10 @@ impl SetDisks {
mut publication_fence: Option<RemoteTuplePublicationFence>,
) -> Result<(ObjectInfo, Option<OldCurrentSize>)> {
let protect_write = opts.shard_integrity_write_enabled();
let source_bucket_incarnation_id = match opts.expected_bucket_incarnation_id {
Some(incarnation_id) => Some(incarnation_id),
None => self.bucket_incarnation_id_from_disk(bucket).await.ok(),
};
if publication_fence.is_none()
&& opts.data_movement
&& rustfs_utils::http::metadata_compat::contains_key_str(
@@ -4031,7 +4035,7 @@ impl SetDisks {
{
#[cfg(any(test, feature = "test-util"))]
pause_put_object_commit(bucket, object, PutObjectCommitPause::BeforeNamespace).await;
if let Some(expected_incarnation_id) = opts.expected_bucket_incarnation_id
if let Some(expected_incarnation_id) = source_bucket_incarnation_id
&& opts.bucket_lifecycle_lock_fence.is_none()
{
bucket_lifecycle_guard = Some(
@@ -4642,6 +4646,7 @@ impl SetDisks {
request.object_version_id = committed_version_id
.or_else(|| commit_version_suspended.then(Uuid::nil))
.map(|version_id| version_id.to_string());
request.expected_bucket_incarnation_id = source_bucket_incarnation_id;
let object_lock_guard = _object_lock_guard.take();
let publication_guard = _publication_guard.take();
let bucket_lifecycle_guard = _bucket_lifecycle_guard.take();
@@ -4790,6 +4795,7 @@ impl SetDisks {
request.object_version_id = committed_version_id
.or_else(|| commit_version_suspended.then(Uuid::nil))
.map(|version_id| version_id.to_string());
request.expected_bucket_incarnation_id = source_bucket_incarnation_id;
commit_set.submit_rename_tail_heal(request).await;
}
@@ -7492,8 +7498,15 @@ impl SetDisks {
ensure_delete_commit_locks_held(None, bucket, &encoded_object, opts)?;
authorization.authorized_journal_name(&candidate)?;
begin_scanner_publication_delete_mutation(opts.scanner_publication_commit_scope.as_ref())?;
self.delete_object_version(bucket, &encoded_object, &delete_request, false)
.await?;
self.delete_object_version_with_purge(
bucket,
&encoded_object,
&delete_request,
false,
None,
opts.expected_bucket_incarnation_id,
)
.await?;
if let Some((_, deleted_object)) = replication_delete {
ReplicationLifecycleBridge::schedule_delete(bucket.to_string(), deleted_object).await;
}
@@ -7573,6 +7586,7 @@ impl SetDisks {
fi: &FileInfo,
force_del_marker: bool,
delete_marker_purge: Option<MrfDeleteMarkerPurge>,
source_bucket_incarnation_id: Option<Uuid>,
) -> Result<()> {
let transported = delete_file_info_with_replication_transport_metadata(fi);
let fi = &transported;
@@ -7727,13 +7741,55 @@ impl SetDisks {
{
let version_id = fi.version_id.map(|version| version.to_string());
let _ = self
.add_partial(bucket, object, version_id.as_deref().unwrap_or_default())
.add_partial_with_source_incarnation(
bucket,
object,
version_id.as_deref().unwrap_or_default(),
source_bucket_incarnation_id,
)
.await;
}
quorum_result
}
}
impl SetDisks {
async fn add_partial_with_source_incarnation(
&self,
bucket: &str,
object: &str,
version_id: &str,
source_bucket_incarnation_id: Option<Uuid>,
) -> Result<()> {
if self
.persist_partial_write(bucket, object, Some(version_id), source_bucket_incarnation_id)
.await
{
return Ok(());
}
let Some(source_bucket_incarnation_id) = source_bucket_incarnation_id else {
// The fallback is best-effort. Without the source generation it
// cannot safely target the object name after bucket recreation.
return Ok(());
};
let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options(
bucket.to_string(),
Some(object.to_string()),
false,
Some(HealChannelPriority::Normal),
Some(self.pool_index),
Some(self.set_index),
);
request.object_version_id = (!version_id.is_empty()).then(|| version_id.to_string());
request.source = rustfs_heal_contracts::heal_channel::HealRequestSource::Mrf;
request.expected_bucket_incarnation_id = Some(source_bucket_incarnation_id);
if let Err(error) = rustfs_heal_contracts::heal_channel::send_heal_request(request).await {
warn!(bucket, object, version_id, error = %error, "Failed to enqueue heal request for partial object");
}
Ok(())
}
}
#[async_trait::async_trait]
impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
type Error = Error;
@@ -8021,7 +8077,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}
#[tracing::instrument(skip(self))]
async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, force_del_marker: bool) -> Result<()> {
self.delete_object_version_with_purge(bucket, object, fi, force_del_marker, None)
let source_bucket_incarnation_id = self.bucket_incarnation_id_from_disk(bucket).await.ok();
self.delete_object_version_with_purge(bucket, object, fi, force_del_marker, None, source_bucket_incarnation_id)
.await
}
@@ -8843,7 +8900,15 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
delete_request.set_skip_tier_free_version();
}
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
self.delete_object_version(bucket, object, &delete_request, false).await?;
self.delete_object_version_with_purge(
bucket,
object,
&delete_request,
false,
None,
opts.expected_bucket_incarnation_id,
)
.await?;
if let Some((_, deleted_object)) = replication_delete {
ReplicationLifecycleBridge::schedule_delete(bucket.to_string(), deleted_object).await;
}
@@ -8858,7 +8923,15 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
};
delete_request.set_tier_free_version_id(&Uuid::new_v4().to_string());
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
self.delete_object_version(bucket, object, &delete_request, false).await?;
self.delete_object_version_with_purge(
bucket,
object,
&delete_request,
false,
None,
opts.expected_bucket_incarnation_id,
)
.await?;
}
for version in &versions.free_versions {
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
@@ -8870,7 +8943,15 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
};
delete_request.set_tier_free_version();
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
self.delete_object_version(bucket, object, &delete_request, false).await?;
self.delete_object_version_with_purge(
bucket,
object,
&delete_request,
false,
None,
opts.expected_bucket_incarnation_id,
)
.await?;
}
}
}
@@ -8978,7 +9059,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
};
ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?;
begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?;
self.delete_object_version(bucket, object, &dfi, false)
self.delete_object_version_with_purge(bucket, object, &dfi, false, None, opts.expected_bucket_incarnation_id)
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
self.invalidate_get_object_metadata_cache(bucket, object).await;
@@ -9072,6 +9153,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
&fi,
should_force_delete_marker_for_missing_version(&opts),
delete_marker_purge.clone(),
opts.expected_bucket_incarnation_id,
)
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
@@ -9124,9 +9206,16 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
if opts.skip_free_version {
dfi.set_skip_tier_free_version();
}
self.delete_object_version_with_purge(bucket, object, &dfi, opts.delete_marker, delete_marker_purge)
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
self.delete_object_version_with_purge(
bucket,
object,
&dfi,
opts.delete_marker,
delete_marker_purge,
opts.expected_bucket_incarnation_id,
)
.await
.map_err(|e| to_object_err(e, vec![bucket, object]))?;
#[cfg(test)]
pause_delete_object_commit_after_publish(bucket, object).await;
@@ -9188,28 +9277,9 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
#[tracing::instrument(skip(self))]
async fn add_partial(&self, bucket: &str, object: &str, version_id: &str) -> Result<()> {
if self.persist_partial_write(bucket, object, Some(version_id)).await {
return Ok(());
}
let mut request = rustfs_heal_contracts::heal_channel::create_heal_request_with_options(
bucket.to_string(),
Some(object.to_string()),
false,
Some(HealChannelPriority::Normal),
Some(self.pool_index),
Some(self.set_index),
);
request.object_version_id = (!version_id.is_empty()).then(|| version_id.to_string());
if let Err(e) = rustfs_heal_contracts::heal_channel::send_heal_request(request).await {
warn!(
bucket,
object,
version_id,
error = %e,
"Failed to enqueue heal request for partial object"
);
}
Ok(())
let source_bucket_incarnation_id = self.bucket_incarnation_id_from_disk(bucket).await.ok();
self.add_partial_with_source_incarnation(bucket, object, version_id, source_bucket_incarnation_id)
.await
}
#[tracing::instrument(skip(self))]
@@ -9768,7 +9838,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
#[cfg(all(test, feature = "test-util"))]
pause_transition_transaction_at(bucket, object, TransitionTransactionKillPoint::CommitFenceBeforeLocalCommit).await;
upload_cleanup.disarm();
if let Err(err) = self.delete_object_version(bucket, object, &fi, false).await {
if let Err(err) = self
.delete_object_version_with_purge(bucket, object, &fi, false, None, opts.expected_bucket_incarnation_id)
.await
{
warn!(
bucket = bucket,
object = object,
@@ -9839,7 +9912,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
continue;
}
let _ = self
.add_partial(bucket, object, opts.version_id.as_deref().unwrap_or_default())
.add_partial_with_source_incarnation(
bucket,
object,
opts.version_id.as_deref().unwrap_or_default(),
opts.expected_bucket_incarnation_id,
)
.await;
break;
}
@@ -19558,15 +19636,11 @@ mod put_object_tmp_cleanup_tests {
expected.is_subset(&after_tail),
"the detached tail must re-mark capacity after the first scope was drained"
);
let request = tokio::time::timeout(Duration::from_secs(30), heal_requests.recv())
.await
.expect("a failed tail should submit heal")
.expect("the per-set heal capture should stay connected");
assert_eq!(request.bucket, bucket);
assert_eq!(request.object_prefix.as_deref(), Some(object));
assert_eq!(request.object_version_id.as_deref(), Some(expected_version_id.as_str()));
assert_eq!(request.pool_index, Some(set_disks.pool_index));
assert_eq!(request.set_index, Some(set_disks.set_index));
assert!(
heal_requests.try_recv().is_err(),
"a legacy fixture without a bucket-generation record must not enqueue an unfenced fallback heal"
);
assert_eq!(expected_version_id, Uuid::nil().to_string());
})
.await;
}
@@ -19687,7 +19761,7 @@ mod put_object_tmp_cleanup_tests {
#[tokio::test]
#[serial_test::serial(capacity_dirty_scope)]
async fn tail_drained_put_preserves_quorum_success_and_heals_failed_tail() {
async fn tail_drained_put_preserves_quorum_success_without_unfenced_fallback_heal() {
let (_dirs, disks, set) = hermetic_set_disks(4).await;
let bucket = "put-full-tail-heal";
let object = "full-tail-heal-object";
@@ -19723,12 +19797,10 @@ mod put_object_tmp_cleanup_tests {
.expect("PUT task should join")
.expect("a minority tail error must not negate committed quorum");
assert_eq!(tasks.running(), 0);
let heal = tokio::time::timeout(Duration::from_secs(30), heals.recv())
.await
.expect("failed tail must schedule heal")
.expect("heal capture must remain connected");
assert_eq!(heal.bucket, bucket);
assert_eq!(heal.object_prefix.as_deref(), Some(object));
assert!(
heals.try_recv().is_err(),
"a legacy fixture without a bucket-generation record must not enqueue an unfenced fallback heal"
);
let info = set
.get_object_info(bucket, object, &ObjectOptions::default())
.await
@@ -361,6 +361,9 @@ pub struct HealChannelRequest {
pub object_prefix: Option<String>,
/// Object version ID (optional)
pub object_version_id: Option<String>,
/// Bucket incarnation observed by the object write that produced a local
/// durable partial-write responsibility.
pub expected_bucket_incarnation_id: Option<Uuid>,
/// Force start heal
pub force_start: bool,
/// Priority
@@ -598,6 +601,7 @@ pub fn create_heal_request(
bucket,
object_prefix,
object_version_id: None,
expected_bucket_incarnation_id: None,
force_start,
priority: priority.unwrap_or_default(),
pool_index: None,
@@ -630,6 +634,7 @@ pub fn create_heal_request_with_options(
bucket,
object_prefix,
object_version_id: None,
expected_bucket_incarnation_id: None,
force_start,
priority: priority.unwrap_or_default(),
pool_index,
@@ -661,6 +666,7 @@ fn create_auto_heal_disk_request(set_disk_id: String, priority: Option<HealChann
disk: Some(set_disk_id),
heal_endpoints: Vec::new(),
object_version_id: None,
expected_bucket_incarnation_id: None,
force_start: false,
priority: priority.unwrap_or(HealChannelPriority::Low),
pool_index: None,
+38
View File
@@ -766,6 +766,9 @@ impl HealChannelProcessor {
let mut heal_request = HealRequest::new(heal_type, options, priority);
heal_request.id = request.id;
heal_request.source = request.source;
if request.source == HealRequestSource::Mrf {
heal_request.expected_mrf_bucket_incarnation_id = request.expected_bucket_incarnation_id;
}
heal_request.heal_endpoints = request.heal_endpoints;
// force_start controls admission/queue semantics only. Do not reinterpret it as
// destructive heal options: admin clients commonly pass forceStart=true together
@@ -1143,6 +1146,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Normal,
@@ -1176,6 +1180,7 @@ mod tests {
bucket: String::new(),
object_prefix: None,
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::High,
@@ -1209,6 +1214,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::High,
@@ -1236,6 +1242,27 @@ mod tests {
assert!(heal_request.options.no_lock);
}
#[tokio::test]
async fn test_convert_mrf_fallback_preserves_source_bucket_incarnation() {
let incarnation = uuid::Uuid::new_v4();
let processor = HealChannelProcessor::new(create_test_heal_manager());
let request = HealChannelRequest {
id: "mrf-fallback".to_string(),
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
expected_bucket_incarnation_id: Some(incarnation),
source: HealRequestSource::Mrf,
..Default::default()
};
let heal_request = processor
.convert_to_heal_request(request)
.expect("MRF fallback request converts");
assert_eq!(heal_request.source, HealRequestSource::Mrf);
assert_eq!(heal_request.expected_mrf_bucket_incarnation_id, Some(incarnation));
}
#[tokio::test]
async fn test_convert_to_heal_request_admin_cannot_bypass_object_lock() {
let heal_manager = create_test_heal_manager();
@@ -1263,6 +1290,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Low,
@@ -1301,6 +1329,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Normal,
@@ -1341,6 +1370,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: Some("test-object".to_string()),
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Normal,
@@ -1374,6 +1404,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: Some("logs/".to_string()),
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::High,
@@ -1410,6 +1441,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: Some("pool_0_set_1".to_string()),
heal_endpoints: vec!["http://node0:9000/drive1".to_string()],
priority: HealChannelPriority::Critical,
@@ -1445,6 +1477,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: Some("invalid-disk-id".to_string()),
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Normal,
@@ -1484,6 +1517,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: channel_priority,
@@ -1516,6 +1550,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: None,
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Normal,
@@ -1550,6 +1585,7 @@ mod tests {
bucket: "test-bucket".to_string(),
object_prefix: Some("".to_string()), // Empty prefix should be treated as bucket heal
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Normal,
@@ -1588,6 +1624,7 @@ mod tests {
bucket: "bucket".to_string(),
object_prefix: Some("object".to_string()),
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: None,
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Low,
@@ -2006,6 +2043,7 @@ mod tests {
bucket: "bucket".to_string(),
object_prefix: None,
object_version_id: None,
expected_bucket_incarnation_id: None,
disk: Some("invalid".to_string()),
heal_endpoints: Vec::new(),
priority: HealChannelPriority::Normal,
+2
View File
@@ -563,6 +563,7 @@ fn active_heal_for_dedup_key(active_heals: &HashMap<String, Arc<HealTask>>, key:
fn request_matches_task(request: &HealRequest, task: &HealTask) -> bool {
request.heal_type == task.heal_type
&& request.bucket_incarnation_id == task.bucket_incarnation_id
&& request.expected_mrf_bucket_incarnation_id == task.expected_mrf_bucket_incarnation_id
&& request.options == task.options
&& request.priority == task.priority
&& request.source == task.source
@@ -573,6 +574,7 @@ fn request_matches_task(request: &HealRequest, task: &HealTask) -> bool {
fn request_matches_request(request: &HealRequest, existing: &HealRequest) -> bool {
request.heal_type == existing.heal_type
&& request.bucket_incarnation_id == existing.bucket_incarnation_id
&& request.expected_mrf_bucket_incarnation_id == existing.expected_mrf_bucket_incarnation_id
&& request.options == existing.options
&& request.priority == existing.priority
&& request.source == existing.source
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
+174
View File
@@ -365,6 +365,19 @@ pub async fn publish_committed_snapshot(
sequence: u64,
payload: &[u8],
limit: usize,
) -> Result<SnapshotPublication, SnapshotError> {
publish_committed_snapshot_with_companion(disks, owner, sequence, payload, limit, None).await
}
/// Publish a companion lifecycle record on the same disk before the commit
/// manifest. A replica is committed only after both payloads have been stored.
pub async fn publish_committed_snapshot_with_companion(
disks: &[EcstoreDiskStore],
owner: Uuid,
sequence: u64,
payload: &[u8],
limit: usize,
companion: Option<(&[&str; 2], &[u8], usize)>,
) -> Result<SnapshotPublication, SnapshotError> {
if disks.is_empty() {
return Err(SnapshotError::NoWritableReplica);
@@ -413,6 +426,18 @@ pub async fn publish_committed_snapshot(
continue;
}
}
if let Some((companion_paths, companion_bytes, companion_limit)) = companion {
match cas_replace(disk, companion_paths[slot], companion_bytes, companion_limit).await {
Ok(EcstoreConditionalFileUpdate::Updated) => {}
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => continue,
Err(error) => {
if first_error.is_none() {
first_error = Some(error);
}
continue;
}
}
}
match cas_replace_expected(disk, MANIFEST_PATHS[slot], expected_manifest.map(EcstoreDiskBytes::from), &manifest).await {
Ok(EcstoreConditionalFileUpdate::Updated) => manifest_replicas += 1,
Ok(EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch) => {}
@@ -436,6 +461,80 @@ pub async fn publish_committed_snapshot(
})
}
/// Read a companion only from replicas that also contain its matching
/// committed checkpoint. Differing companions for the same checkpoint fail
/// closed rather than selecting one replica arbitrarily.
pub async fn read_committed_companion(
disks: &[EcstoreDiskStore],
owner: Uuid,
sequence: u64,
companion_paths: &[&str; 2],
checkpoint_limit: usize,
companion_limit: usize,
) -> Result<Option<Vec<u8>>, SnapshotError> {
let mut selected: Option<Vec<u8>> = None;
let mut first_error = None;
for disk in disks {
let mut checkpoint_slot = None;
for (slot, (manifest_path, payload_path)) in MANIFEST_PATHS.into_iter().zip(PAYLOAD_PATHS).enumerate() {
let manifest_bytes = match read_bounded(disk, manifest_path, MANIFEST_LEN).await {
Ok(Some(bytes)) => bytes,
Ok(None) => continue,
Err(error @ SnapshotError::Unsupported) => return Err(error),
Err(error) => {
if first_error.is_none() {
first_error = Some(error);
}
continue;
}
};
let manifest = match Manifest::decode(&manifest_bytes, checkpoint_limit) {
Ok(manifest) if manifest.owner == owner && manifest.sequence == sequence => manifest,
Ok(_) | Err(SnapshotError::Corrupt) | Err(SnapshotError::TooLarge) => continue,
Err(error) => return Err(error),
};
let payload = match read_bounded(disk, payload_path, manifest.payload_len).await {
Ok(Some(payload)) => payload,
Ok(None) => continue,
Err(error) => {
if first_error.is_none() {
first_error = Some(error);
}
continue;
}
};
if CommittedSnapshot::decode(slot, &manifest_bytes, payload, checkpoint_limit).is_ok() {
checkpoint_slot = Some(slot);
break;
}
}
let Some(checkpoint_slot) = checkpoint_slot else {
continue;
};
let companion = match read_bounded(disk, companion_paths[checkpoint_slot], companion_limit).await {
Ok(Some(companion)) => companion,
Ok(None) => continue,
Err(error) => {
if first_error.is_none() {
first_error = Some(error);
}
continue;
}
};
if selected.as_ref().is_some_and(|existing| existing != &companion) {
return Err(SnapshotError::Conflict);
}
selected = Some(companion);
}
if selected.is_some() {
Ok(selected)
} else if let Some(error) = first_error {
Err(error)
} else {
Ok(None)
}
}
/// Reclaim the slot superseded by an already committed checkpoint.
///
/// This is a narrow cleanup primitive: it first reads back the current
@@ -842,6 +941,81 @@ mod tests {
}
}
#[tokio::test]
async fn lifecycle_companion_is_read_only_with_its_matching_checkpoint_replica() {
let root = TempDir::new().expect("test directory");
let checkpoint_disk = disk(&root, "checkpoint").await;
let sidecar_disk = disk(&root, "sidecar").await;
let owner = Uuid::new_v4();
let checkpoint_payload = payload("responsibility");
let sidecar = b"lifecycle state";
commit(&checkpoint_disk, 0, owner, 7, &checkpoint_payload).await;
install(&sidecar_disk, ".heal-mrf-lifecycle.0.bin", sidecar).await;
let absent = read_committed_companion(
&[checkpoint_disk.clone(), sidecar_disk.clone()],
owner,
7,
&[".heal-mrf-lifecycle.0.bin", ".heal-mrf-lifecycle.1.bin"],
4096,
4096,
)
.await
.expect("unpaired lifecycle sidecar is ignored");
assert!(absent.is_none());
publish_committed_snapshot_with_companion(
&[checkpoint_disk.clone(), sidecar_disk.clone()],
owner,
8,
&checkpoint_payload,
4096,
Some((&[".heal-mrf-lifecycle.0.bin", ".heal-mrf-lifecycle.1.bin"], sidecar, 4096)),
)
.await
.expect("paired checkpoint and lifecycle state publish");
let paired = read_committed_companion(
&[checkpoint_disk, sidecar_disk],
owner,
8,
&[".heal-mrf-lifecycle.0.bin", ".heal-mrf-lifecycle.1.bin"],
4096,
4096,
)
.await
.expect("read paired lifecycle sidecar")
.expect("paired sidecar exists");
assert_eq!(paired, sidecar);
}
#[tokio::test]
async fn lifecycle_companion_recovers_from_a_healthy_replica_after_peer_read_error() {
let root = TempDir::new().expect("test directory");
let failing_disk = disk(&root, "failing").await;
let healthy_disk = disk(&root, "healthy").await;
let owner = Uuid::new_v4();
let checkpoint_payload = payload("responsibility");
let sidecar = b"operator audit and source incarnation";
commit(&healthy_disk, 0, owner, 9, &checkpoint_payload).await;
install(&healthy_disk, ".heal-mrf-lifecycle.0.bin", sidecar).await;
std::fs::create_dir(root.path().join("failing").join(RUSTFS_META_BUCKET).join(MANIFEST_PATHS[0]))
.expect("simulate one replica read error");
let recovered = read_committed_companion(
&[failing_disk, healthy_disk],
owner,
9,
&[".heal-mrf-lifecycle.0.bin", ".heal-mrf-lifecycle.1.bin"],
4096,
4096,
)
.await
.expect("one unreadable replica must not hide a healthy paired sidecar")
.expect("healthy paired sidecar must be recovered");
assert_eq!(recovered, sidecar);
}
#[tokio::test]
async fn divergent_commits_at_same_sequence_fail_closed() {
let root = TempDir::new().expect("test directory");
+8
View File
@@ -323,6 +323,10 @@ pub struct HealRequest {
pub heal_type: HealType,
/// Admission identity for an explicit administrator bucket heal. Never rebound on replay.
pub bucket_incarnation_id: Option<Uuid>,
/// Source bucket generation captured for a durable MRF object repair.
/// Unlike the admin admission fence above, this must survive queueing and
/// retries so execution cannot target a later bucket incarnation.
pub expected_mrf_bucket_incarnation_id: Option<Uuid>,
/// Heal options
pub options: HealOptions,
/// Priority
@@ -351,6 +355,7 @@ impl HealRequest {
id: Uuid::new_v4().to_string(),
heal_type,
bucket_incarnation_id: None,
expected_mrf_bucket_incarnation_id: None,
options,
priority,
source: HealRequestSource::Internal,
@@ -417,6 +422,7 @@ pub struct HealTask {
/// Heal type
pub heal_type: HealType,
pub bucket_incarnation_id: Option<Uuid>,
pub expected_mrf_bucket_incarnation_id: Option<Uuid>,
/// Heal options
pub options: HealOptions,
/// Priority inherited from the request
@@ -494,6 +500,7 @@ impl HealTask {
id: request.id,
heal_type: request.heal_type,
bucket_incarnation_id: request.bucket_incarnation_id,
expected_mrf_bucket_incarnation_id: request.expected_mrf_bucket_incarnation_id,
options: request.options,
priority: request.priority,
source: request.source,
@@ -530,6 +537,7 @@ impl HealTask {
id: self.id.clone(),
heal_type: self.heal_type.clone(),
bucket_incarnation_id: self.bucket_incarnation_id,
expected_mrf_bucket_incarnation_id: self.expected_mrf_bucket_incarnation_id,
options: self.options.clone(),
priority: self.priority,
source: self.source,
+11 -1
View File
@@ -551,12 +551,22 @@ impl HealTask {
set: self.options.set_index,
};
let mut expected = self.outcome_identity(bucket, object, version_id, self.options.pool_index, self.options.set_index);
let bucket_incarnation_id = self
let current_bucket_incarnation_id = self
.outcome_bucket_incarnation_id(bucket, self.options.dry_run)
.await?
.ok_or_else(|| Error::TaskExecutionFailed {
message: format!("Missing bucket incarnation for durable MRF repair {bucket}/{object}"),
})?;
let bucket_incarnation_id = self
.expected_mrf_bucket_incarnation_id
.ok_or_else(|| Error::TaskExecutionFailed {
message: format!("Missing source bucket incarnation for durable MRF repair {bucket}/{object}"),
})?;
if current_bucket_incarnation_id != bucket_incarnation_id {
return Err(Error::TaskExecutionFailed {
message: format!("Bucket incarnation changed before durable MRF repair {bucket}/{object}"),
});
}
expected.bucket_incarnation_id = Some(bucket_incarnation_id);
let storage_result = self
+20
View File
@@ -2678,6 +2678,25 @@ async fn read_repair_object_heal_sets_read_repair_option() {
assert!(!opts[0].no_lock);
}
#[tokio::test]
async fn durable_mrf_heal_rejects_a_recreated_bucket_before_storage_heal() {
let original_incarnation = Uuid::new_v4();
let recreated_incarnation = Uuid::new_v4();
let storage = Arc::new(MockStorage {
bucket_incarnation_id: Mutex::new(Some(recreated_incarnation)),
..Default::default()
});
let mut request = HealRequest::object("bucket".to_string(), "object".to_string(), None);
request.source = HealRequestSource::Mrf;
request.expected_mrf_bucket_incarnation_id = Some(original_incarnation);
let task = HealTask::from_request(request, storage.clone());
let result = task.heal_object("bucket", "object", None).await;
assert!(matches!(result, Err(Error::TaskExecutionFailed { .. })));
assert!(storage.heal_object_calls.lock().unwrap().is_empty());
}
#[tokio::test(start_paused = true)]
async fn read_repair_object_heal_is_not_failed_by_flat_task_timeout() {
let storage = Arc::new(MockStorage {
@@ -4002,6 +4021,7 @@ async fn mrf_recreate_missing_object_records_exact_absence_receipt_with_scope()
HealPriority::Normal,
);
request.source = HealRequestSource::Mrf;
request.expected_mrf_bucket_incarnation_id = Some(incarnation);
let task = HealTask::from_request(request, storage.clone());
task.execute()
+349 -3
View File
@@ -91,6 +91,201 @@ async fn partial_write_persistence_failure_is_reported_and_retained_for_retry()
);
}
#[test]
fn legacy_unbound_generation_is_parked_and_requires_explicit_risk_acceptance() {
const STACK_SIZE: usize = 8 * 1024 * 1024;
std::thread::Builder::new()
.name("mrf-unbound-lifecycle".to_owned())
.stack_size(STACK_SIZE)
.spawn(|| {
let runtime = tokio::runtime::Builder::new_current_thread()
.thread_stack_size(STACK_SIZE)
.enable_all()
.build()
.expect("unbound lifecycle runtime should build");
runtime.block_on(legacy_unbound_generation_is_parked_and_requires_explicit_risk_acceptance_inner());
})
.expect("unbound lifecycle test thread should spawn")
.join()
.expect("unbound lifecycle test thread should finish");
}
async fn legacy_unbound_generation_is_parked_and_requires_explicit_risk_acceptance_inner() {
use rustfs_common::mrf_channel::{MrfScope, persist_partial_write_intent};
temp_env::async_with_vars([("RUSTFS_HEAL_MRF_ENABLE", Some("true"))], async {
let root = tempfile::tempdir().expect("unbound lifecycle fixture directory");
let env = TestECStoreEnv::builder().base_dir(root.path()).build().await;
env.make_bucket("unbound-lifecycle", false).await;
let manager = manager(&env);
manager.start().await.expect("manager should start");
mrf_queue::spawn_mrf_consumer(manager.clone());
persist_partial_write_intent(
"unbound-lifecycle",
"legacy.bin",
None,
MrfScope {
pool_index: 0,
set_index: 0,
},
)
.await
.expect("old-format intent should still be durably admitted");
assert!(
wait_until(|| async {
mrf_queue::list_legacy_responsibilities(None, 8).await.is_ok_and(|snapshot| {
snapshot
.responsibilities
.iter()
.any(|item| item.bucket == "unbound-lifecycle" && item.object == "legacy.bin")
})
})
.await,
"unbound old journal intent must become visible"
);
let snapshot = mrf_queue::list_legacy_responsibilities(None, 8)
.await
.expect("listing should be available");
let entry = snapshot
.responsibilities
.iter()
.find(|item| item.bucket == "unbound-lifecycle" && item.object == "legacy.bin")
.expect("visible generation-unknown entry");
assert_eq!(entry.status, "legacy_generation_unknown");
assert!(entry.source_bucket_incarnation_id.is_none());
assert!(
snapshot_contains("legacy.bin").await,
"generation uncertainty must preserve journal responsibility"
);
let operations = manager.operations_snapshot().await;
assert_eq!(
operations.queue_length, 0,
"unbound old intent must not be sent to the current bucket generation"
);
assert!(matches!(
mrf_queue::recheck_legacy_responsibility(entry.responsibility_id, entry.bucket_incarnation_id).await,
Err(rustfs_heal::heal::mrf_queue::MrfLifecycleControlError::InvalidAction(_))
));
mrf_queue::accept_unverified_legacy_risk(mrf_queue::MrfLegacyRiskAcceptanceRequest {
responsibility_id: entry.responsibility_id,
expected_bucket_incarnation_id: entry.bucket_incarnation_id,
acknowledge_unknown_source_incarnation: true,
acknowledge_incarnation_mismatch: false,
actor: "integration-test-operator".to_string(),
reason: "The old journal has no source bucket-generation binding".to_string(),
reference: "TEST-ISSUE-2682-UNBOUND".to_string(),
request_id: uuid::Uuid::new_v4(),
})
.await
.expect("unbound risk disposition requires explicit acknowledgment and must be durable");
assert!(
snapshot_contains("legacy.bin").await,
"risk acknowledgment must not delete the intent record"
);
let accepted = mrf_queue::list_legacy_responsibilities(None, 8)
.await
.expect("accepted status should be readable");
let accepted = accepted
.responsibilities
.iter()
.find(|item| item.responsibility_id == entry.responsibility_id)
.expect("accepted unbound generation remains visible");
assert_eq!(accepted.status, "operator_accepted_unverified");
assert!(
accepted
.accepted
.as_ref()
.is_some_and(|audit| audit.acknowledged_unknown_source_incarnation)
);
manager.stop().await.expect("manager should stop");
})
.await;
}
#[test]
fn unversioned_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())
.stack_size(STACK_SIZE)
.spawn(|| {
let runtime = tokio::runtime::Builder::new_current_thread()
.thread_stack_size(STACK_SIZE)
.enable_all()
.build()
.expect("partial-write absence runtime should build");
runtime.block_on(unversioned_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_with_incarnation};
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 source_bucket_incarnation_id = env.ecstore.pools[0]
.get_disks(0)
.bucket_incarnation_id_from_disk("partial-absence")
.await
.expect("fixture bucket source incarnation");
for (object, version_id) in [
("deleted-unversioned.bin", None),
("deleted-versioned.bin", Some(uuid::Uuid::new_v4())),
] {
persist_partial_write_intent_with_incarnation(
"partial-absence",
object,
version_id,
MrfScope {
pool_index: 0,
set_index: 0,
},
Some(source_bucket_incarnation_id),
)
.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"] {
assert!(
wait_until(|| async { !snapshot_contains(object).await }).await,
"complete absence proof must discharge the durable 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");
})
.await;
}
#[test]
fn degraded_deleted_partial_write_is_discharged_by_an_absence_proof() {
const STACK_SIZE: usize = 8 * 1024 * 1024;
@@ -703,9 +898,12 @@ async fn partial_write_sigkill_replay_scenario(protected: bool) {
mrf_queue::spawn_mrf_consumer(manager.clone());
assert!(snapshot_contains("crash.bin").await, "restart must find durable responsibility");
*set.disks.write().await = all.iter().cloned().map(Some).collect();
let healed = wait_until(|| async { replicas(&all, "partial-crash", "crash.bin", None, false).await == 4 }).await;
assert!(
wait_until(|| async { replicas(&all, "partial-crash", "crash.bin", None, false).await == 4 }).await,
"replayed responsibility must heal the returning member"
healed,
"replayed responsibility must heal the returning member; manager={:?}; responsibility={:?}",
manager.operations_snapshot().await,
mrf_queue::list_legacy_responsibilities(None, 16).await
);
assert_payload(&env, "partial-crash", "crash.bin", None, b"durable partial write across SIGKILL").await;
if protected {
@@ -734,10 +932,158 @@ async fn partial_write_sigkill_replay_scenario(protected: bool) {
while tokio::time::Instant::now() < retry_window {
let snapshot = manager.operations_snapshot().await;
assert_eq!(snapshot.queue_length, 0, "an unverified legacy result must not refill the manager queue");
assert_eq!(snapshot.active_tasks, 0, "a held intent must not stay active");
assert_eq!(
snapshot.active_tasks,
0,
"a held intent must not stay active; lifecycle state: {:?}",
mrf_queue::list_legacy_responsibilities(None, 32)
.await
.expect("lifecycle state should be inspectable")
.responsibilities
);
tokio::time::sleep(Duration::from_millis(100)).await;
}
assert!(snapshot_contains("crash.bin").await, "unverified legacy responsibility must remain");
let before = mrf_queue::list_legacy_responsibilities(None, 32)
.await
.expect("held lifecycle listing should be available");
let held = before
.responsibilities
.iter()
.find(|entry| entry.bucket == "partial-crash" && entry.object == "crash.bin")
.expect("legacy durable intent must be visible with its exact identity");
assert_eq!(held.status, "held_unverified_legacy");
assert_eq!(
held.source_bucket_incarnation_id,
Some(
env.ecstore.pools[0]
.get_disks(0)
.bucket_incarnation_id_from_disk("partial-crash")
.await
.expect("source bucket incarnation should remain stable")
),
"the producer-bound bucket generation must survive journal replay"
);
let responsibility_id = held.responsibility_id;
let bucket_incarnation_id = held.bucket_incarnation_id;
let request_id = uuid::Uuid::new_v4();
mrf_queue::accept_unverified_legacy_risk(mrf_queue::MrfLegacyRiskAcceptanceRequest {
responsibility_id,
expected_bucket_incarnation_id: bucket_incarnation_id,
acknowledge_unknown_source_incarnation: true,
acknowledge_incarnation_mismatch: false,
actor: "integration-test-operator".to_string(),
reason: "The operator has accepted that legacy object identity cannot be proven automatically".to_string(),
reference: "TEST-ISSUE-2682".to_string(),
request_id,
})
.await
.expect("explicit risk acceptance must persist before success");
assert!(snapshot_contains("crash.bin").await, "risk acceptance must retain the MRF responsibility");
let accepted = mrf_queue::list_legacy_responsibilities(None, 32)
.await
.expect("accepted lifecycle state should remain queryable");
let accepted = accepted
.responsibilities
.iter()
.find(|entry| entry.responsibility_id == responsibility_id)
.expect("accepted responsibility must remain visible");
assert_eq!(accepted.status, "operator_accepted_unverified");
assert_eq!(accepted.accepted.as_ref().map(|audit| audit.request_id), Some(request_id));
assert_eq!(
accepted.accepted.as_ref().map(|audit| audit.actor.as_str()),
Some("integration-test-operator")
);
manager
.stop()
.await
.expect("first manager should stop before restart verification");
let log = std::fs::File::create(root.path().join("lifecycle-restart.log")).expect("restart child log");
let mut child = Command::new(std::env::current_exe().expect("integration test executable"))
.args(["--exact", "mrf_legacy_lifecycle_restore_fixture", "--nocapture"])
.env("RUSTFS_TEST_MRF_LIFECYCLE_ROOT", root.path())
.env("RUSTFS_TEST_MRF_LIFECYCLE_ID", responsibility_id.to_string())
.stdout(Stdio::from(log.try_clone().expect("clone child log")))
.stderr(Stdio::from(log))
.spawn()
.expect("lifecycle restart fixture should start");
let restored = wait_until(|| async { root.path().join("lifecycle-restored").exists() }).await;
let status = child.wait().expect("lifecycle restart fixture should exit");
assert!(
restored && status.success(),
"operator-accepted state should survive process restart: {}",
std::fs::read_to_string(root.path().join("lifecycle-restart.log")).expect("read child evidence")
);
return;
}
manager.stop().await.expect("restarted manager should stop");
}
#[test]
fn mrf_legacy_lifecycle_restore_fixture() {
const STACK_SIZE: usize = 8 * 1024 * 1024;
std::thread::Builder::new()
.name("mrf-lifecycle-restore".to_owned())
.stack_size(STACK_SIZE)
.spawn(|| {
let runtime = tokio::runtime::Builder::new_current_thread()
.thread_stack_size(STACK_SIZE)
.enable_all()
.build()
.expect("lifecycle restore runtime should build");
runtime.block_on(mrf_legacy_lifecycle_restore_fixture_inner());
})
.expect("lifecycle restore thread should spawn")
.join()
.expect("lifecycle restore thread should finish");
}
async fn mrf_legacy_lifecycle_restore_fixture_inner() {
let Ok(root) = std::env::var("RUSTFS_TEST_MRF_LIFECYCLE_ROOT") else {
return;
};
let expected_id = std::env::var("RUSTFS_TEST_MRF_LIFECYCLE_ID")
.expect("lifecycle child responsibility ID")
.parse::<uuid::Uuid>()
.expect("valid lifecycle child responsibility ID");
let root = std::path::PathBuf::from(root);
let env = TestECStoreEnv::builder().base_dir(&root).build().await;
let manager = manager(&env);
manager.start().await.expect("restarted heal manager should start");
mrf_queue::spawn_mrf_consumer(manager.clone());
assert!(
wait_until(|| async {
mrf_queue::list_legacy_responsibilities(None, 32).await.is_ok_and(|snapshot| {
snapshot
.responsibilities
.iter()
.any(|entry| entry.responsibility_id == expected_id)
})
})
.await,
"durable operator-accepted lifecycle state must replay"
);
let restored = mrf_queue::list_legacy_responsibilities(None, 32)
.await
.expect("restored lifecycle listing should be available");
let entry = restored
.responsibilities
.iter()
.find(|entry| entry.responsibility_id == expected_id)
.expect("restart listing matched the stable generation");
assert_eq!(entry.status, "operator_accepted_unverified");
assert_eq!(
entry.accepted.as_ref().map(|audit| audit.actor.as_str()),
Some("integration-test-operator")
);
assert!(
snapshot_contains("crash.bin").await,
"risk-accepted responsibility remains in the durable journal"
);
tokio::fs::write(root.join("lifecycle-restored"), b"restored")
.await
.expect("signal lifecycle restore");
manager.stop().await.expect("restarted manager should stop");
}
+1
View File
@@ -153,6 +153,7 @@ impl StartCommand {
bucket: self.bucket,
object_prefix: self.object_prefix,
object_version_id: self.object_version_id,
expected_bucket_incarnation_id: None,
force_start: self.force_start,
priority: self.priority.into(),
pool_index: self
@@ -18,6 +18,7 @@
- `backlog-2519` retained admin heal reports: keep the schema-1 terminal as the commit and replay fence, with a bounded versioned report in a separate namespace on the same disk. Older rollback readers can still query the terminal and suppress replay, but cannot expose its outcome; newer readers mark missing outcomes as unavailable. Remove the legacy marker and missing-report adapter only after all supported direct-upgrade and rollback readers understand the report format and retained schema-1-only receipts have expired.
- `odm-list-bare-envelope` historical ODM continuation tokens: preserve complete bare v1/v2 envelopes. Framed issuance defaults on for the deployed framed-only generation; upgrades from older bare-only readers must explicitly disable it before starting new nodes and keep it off until reader convergence. Remove the legacy classifier and framing issuance override only after every supported reader accepts framing and outstanding bare listings have drained or clients explicitly restarted them; tokens have no automatic expiry. Exact full-envelope object keys remain intrinsically ambiguous during this compatibility period.
- `backlog-2263` legacy heal MRF inspection: retained per-record journals remain readable while committed-snapshot ownership and writer activation are staged. Remove legacy import only after all supported direct-upgrade and rollback readers understand committed snapshots and migration tooling confirms that no retained or restorable legacy journal requires it. This does not enable a new writer or change the automatic legacy consumer.
- `backlog-2682` MRF lifecycle downgrade fence: the lifecycle sidecar binds durable responsibilities to a source bucket incarnation and records operator disposition, but pre-lifecycle readers ignore it and may replay the compatibility journal against a same-name recreated bucket. Downgrade to those readers is unsupported while any durable MRF responsibility remains. Remove the restriction only after every supported rollback reader enforces the source-incarnation fence and operators have verified that no old-format or lifecycle responsibility remains.
- `backlog-1337` legacy restore orphan recovery: releases that predate the restore worker-lock marker can leave a valid operation-id and `ongoing-request="true"` after cancellation or process failure, with no durable liveness proof. New servers allow an exact, non-nil legacy generation to be superseded only when its consistently parsed request date is at least 24 hours old. Remove the clock-based legacy fallback after the minimum supported direct-upgrade release writes the v1 worker-lock marker on every restore and operators have resolved every retained pre-v1 ongoing generation.
- `backlog-2133-tier-delete-chunk-parent` bounded tier-delete dispatch compatibility: prefixes at or below the legacy manifest limit keep the byte-compatible v1 single-manifest protocol, while larger prefixes place a chunk-parent sentinel at the original deterministic root path and use operation-scoped child manifests. Older binaries reject the sentinel and child paths, preserving the v6 sole-owner downgrade fence instead of starting a competing local delete. Remove the v1 reader and fail-closed mixed-version sentinel only after every supported rollback release validates the parent/child protocol and migration tooling confirms that no retained v1 dispatch manifest remains.
- `backlog-2102` rc.2/rc.3 empty scanner usage floor recovery: old DeleteBucket cleanup could synthesize an empty incomplete v2 usage primary/backup before leadership added an epoch, while newer scanners require a durable authoritative baseline identity. New scanners recognize only that exact serialized empty-fence shape, preserve its epoch through a CAS-protected recovery marker, and rebuild namespace coverage without treating zero usage as authoritative. Remove this recovery path and marker after rc.2 and rc.3 are no longer supported direct-upgrade sources.
+27 -3
View File
@@ -288,11 +288,35 @@ Durable partial-write responsibilities also have node-local MRF metrics:
| Metric | Meaning |
|---|---|
| `rustfs_heal_mrf_queue_depth` | In-memory MRF work plus retained durable responsibilities, including held legacy entries. |
| `rustfs_heal_mrf_queue_depth` | In-memory MRF work plus retained durable responsibilities, including held and operator-accepted unverified entries. |
| `rustfs_heal_mrf_unverified_legacy` | Durable `PartialWrite` entries whose completed deep check found healthy legacy data without an independent payload identity proof. The journal entry remains intact; the current process pauses automatic re-dispatch for that entry. |
| `rustfs_heal_mrf_unverified_legacy_oldest_age_seconds` | Age of the oldest such retained responsibility on this node. |
| `rustfs_heal_mrf_legacy_generation_unknown` | Durable entries replayed without a source bucket-incarnation binding; they are not dispatched until an operator chooses an explicit disposition. |
| `rustfs_heal_mrf_bucket_incarnation_changed` | Durable entries whose source bucket incarnation differs from the current bucket; the old intent is parked rather than sent to the recreated bucket. |
| `rustfs_heal_mrf_operator_accepted_unverified` | Responsibilities for which an administrator explicitly accepted the unresolved data risk. They remain durable and are never reported as repaired or verified. |
| `rustfs_heal_mrf_operator_accepted_unverified_oldest_age_seconds` | Age of the oldest durable operator-accepted, still-unverified responsibility on this node. |
These gauges are emitted per node through the configured OTLP metrics exporter (`RUSTFS_OBS_ENDPOINT` or `RUSTFS_OBS_METRIC_ENDPOINT`). The background-heal status endpoint reports execution tasks; it does not include durable MRF responsibilities. A process restart replays unchanged journal records and checks them again; it does not delete or certify an entry, and it replays all other durable intents on that node as well. For an unversioned object, restore the expected content from a trusted canonical source with protected shard integrity enabled, then restart the node that owns the journal so its intent can obtain an exact receipt. For a versioned object, a protected rewrite creates a new version and does not resolve a responsibility for the old version; preserve the source and only retire that exact version when the intended data is backed up and an authoritative absence proof is appropriate. A successful receipt may discharge only the matching responsibility; an unverified result remains retained and becomes held again. If no trusted copy or expected checksum is available, preserve the held responsibility and investigate its source of truth before retrying. The protected-copy migration in the [shard-integrity audit workflow](shard-integrity-audit.md) writes a new key and preserves its source; completing that migration alone does not prove or discharge an intent for the original key. Never remove MRF journal files manually.
These gauges are emitted per node through the configured OTLP metrics exporter (`RUSTFS_OBS_ENDPOINT` or `RUSTFS_OBS_METRIC_ENDPOINT`). The background-heal status endpoint reports execution tasks; it does not include durable MRF responsibilities. Use the authenticated node-local `GET /rustfs/admin/v4/heal/mrf/responsibilities?limit=100&cursor=<nextCursor>` endpoint to list held, generation-unknown, incarnation-mismatched, or operator-accepted entries, including the stable responsibility ID, exact object/version/scope, source/current bucket incarnation, checkpoint owner/sequence, and any risk-acceptance record. `limit` is bounded to 1–256 and `nextCursor` is absent on the last page. If the listed bucket incarnation may have changed since the responsibility was parked, call the action endpoint with `action: "refresh"`, the listed `responsibilityId`, and its current `expectedBucketIncarnationId`. Refresh reads the live bucket identity, persists a changed-generation classification, and never dispatches a heal. Then review the fresh listing before taking another action. Use `action: "recheck"` to resume one exact source-bound held generation after restoring a trusted source or preparing an exact-version deletion. A recheck is persisted before the task is dispatched.
An operator may instead choose `action: "acceptUnverifiedRisk"` only with `acknowledgeUnverifiedDataRisk: true`, explicit booleans for both source-incarnation acknowledgments, a non-empty reason, an audit reference, and a non-nil request ID. When `sourceBucketIncarnationId` is absent, set `acknowledgeUnknownSourceIncarnation: true`; when it differs from the current `bucketIncarnationId`, set `acknowledgeBucketIncarnationMismatch: true`. A generation with a known source/current mismatch is parked before dispatch, so its old intent cannot be applied to a recreated bucket. Legacy journal records without a source incarnation are listed as `legacy_generation_unknown` and are not automatically dispatched; refresh can record the live generation but does not bind the old responsibility to it. This does not create a payload proof, delete the object, or remove the durable intent. It persists an `operator_accepted_unverified` disposition that stops automatic retries while leaving the responsibility visible. The actor, reason, reference, time, request ID, target identity, source/current bucket incarnations, and acknowledgment flags are kept in a checksummed lifecycle record in the same disk replica and alternating slot as its committed MRF checkpoint. If that lifecycle record is missing, corrupt, or belongs to another checkpoint, unbound records return to `legacy_generation_unknown`; they are not treated as proof or dispatched against an unbound bucket. Pre-lifecycle binaries ignore this state and can heal a compatibility-journal object into a recreated bucket generation; downgrading while any durable MRF responsibility remains is unsupported. See `backlog-2682` in the compatibility cleanup register.
Example risk-acceptance body (substitute values from the node-local listing):
```json
{
"action": "acceptUnverifiedRisk",
"responsibilityId": "<responsibility-id>",
"expectedBucketIncarnationId": "<bucket-incarnation-id>",
"acknowledgeUnverifiedDataRisk": true,
"acknowledgeUnknownSourceIncarnation": false,
"acknowledgeBucketIncarnationMismatch": false,
"reason": "The source was reviewed and the remaining identity risk is accepted",
"reference": "INC-1234",
"requestId": "<new-request-uuid>"
}
```
A process restart restores lifecycle state only when the checksummed sidecar matches the committed checkpoint. Otherwise the durable intent remains visible and generation-unknown; no record is deleted or certified, and other durable intents on that node replay normally. For source-bound unversioned objects, restore expected content from a trusted canonical source with protected shard integrity enabled, then request a targeted recheck so the exact intent can obtain a receipt. For a legacy record without a source generation, the operator must first establish which bucket generation the responsibility belongs to; if that cannot be established, keep it generation-unknown or explicitly accept the risk. For a versioned object, a protected rewrite creates a new version and does not resolve a responsibility for the old version; preserve the source and only retire that exact version when the intended data is backed up and an authoritative absence proof is appropriate. A successful receipt may discharge only the matching responsibility; an unverified result remains retained and becomes held again. If no trusted copy or expected checksum is available, keep the responsibility held or explicitly accept the risk with the node-local admin operation. The protected-copy migration in the [shard-integrity audit workflow](shard-integrity-audit.md) writes a new key and preserves its source; completing that migration alone does not prove or discharge an intent for the original key. Never remove MRF journal files manually.
## Heal runtime controls
@@ -382,7 +406,7 @@ Rediscovery and admission observations update the recorded result but do not pos
With MRF enabled, a partial commit waits up to ten seconds for its checkpoint receipt. Fully converged writes do not enter this path or add MRF journal writes. The consumer batches available submissions and uses the existing count/byte budgets; an offline member or exhausted hint retry count does not evict an admitted partial-write obligation. New responsibility for the same unversioned object invalidates the previous lease, and a task that already started cannot accept that new responsibility as a merged proof target.
Disabled/unavailable delivery, full queues, invalid identities, persistence failures and receipt timeouts are not durable admission. An already committed object is not rolled back: the producer falls back to the existing in-memory Heal channel. When MRF is enabled, a failed admission attempt also records `mrf_durable_admission_failed`. An admission timeout does not cancel a record already retained by the consumer. These failure cases therefore do not promise that every successful S3 response has a durable repair obligation. Early-ACK PUTs can also return before their rename tail settles; the tail submits responsibility when its final outcome requires repair. The MRF switch is independent of automatic disk scanning, so returning members can be repaired with scanner and auto-heal disabled.
Disabled/unavailable delivery, full queues, invalid identities, persistence failures and receipt timeouts are not durable admission. An already committed object is not rolled back: when its source bucket incarnation is known, the producer falls back to the in-memory Heal channel with that same generation fence. If the source incarnation is unavailable, it does not dispatch an unfenced heal that could mutate a recreated bucket. When MRF is enabled, a failed admission attempt also records `mrf_durable_admission_failed`. An admission timeout does not cancel a record already retained by the consumer. These failure cases therefore do not promise that every successful S3 response has a durable repair obligation. Early-ACK PUTs can also return before their rename tail settles; the tail submits responsibility when its final outcome requires repair. The MRF switch is independent of automatic disk scanning, so returning members can be repaired with scanner and auto-heal disabled.
Legacy notices carry only bucket/object/version, not a verified storage disposition, incarnation, scope, or durable responsibility generation. They are drained without clearing hints. Terminal callbacks release only their exact node-local ingress lease so rediscovery remains possible; lease generations are not durable successor receipts. Pending migration staging is not activated, and this change does not enable durable tombstones or garbage collection. Positive cleanup requires a storage-owner receipt with the complete responsibility identity and validated commit/fence evidence; neither task status nor the bounded diagnostic outcome window supplies it.
+359 -9
View File
@@ -57,7 +57,11 @@ const EVENT_ADMIN_REQUEST_FAILED: &str = "admin_request_failed";
const EVENT_ADMIN_RESPONSE_EMITTED: &str = "admin_response_emitted";
const LEGACY_ROOT_HEAL_RESPONSE_ID: &str = ".";
const PEER_HEAL_STATUS_TIMEOUT: Duration = Duration::from_secs(5);
const MRF_RESPONSIBILITY_LIST_DEFAULT_LIMIT: usize = 100;
const MRF_RESPONSIBILITY_LIST_MAX_LIMIT: usize = 256;
pub(crate) const REPLACEMENT_RECOVERY_STATUS_ROUTE_SUFFIX: &str = "/v4/heal/replacement-recovery";
pub(crate) const MRF_LEGACY_RESPONSIBILITIES_ROUTE_SUFFIX: &str = "/v4/heal/mrf/responsibilities";
pub(crate) const MRF_LEGACY_RESPONSIBILITIES_ACTIONS_ROUTE_SUFFIX: &str = "/v4/heal/mrf/responsibilities/actions";
const REPLACEMENT_RECOVERY_STATUS_CONTRACT_VERSION: u32 = 2;
#[derive(Debug, Default, Serialize, Deserialize)]
@@ -221,6 +225,18 @@ pub fn register_heal_route(r: &mut S3Router<AdminOperation>) -> std::io::Result<
AdminOperation(&BackgroundHealStatusHandler {}),
)?;
r.insert(
Method::GET,
format!("{}{}", ADMIN_PREFIX, MRF_LEGACY_RESPONSIBILITIES_ROUTE_SUFFIX).as_str(),
AdminOperation(&MrfLegacyResponsibilitiesHandler {}),
)?;
r.insert(
Method::POST,
format!("{}{}", ADMIN_PREFIX, MRF_LEGACY_RESPONSIBILITIES_ACTIONS_ROUTE_SUFFIX).as_str(),
AdminOperation(&MrfLegacyResponsibilitiesActionHandler {}),
)?;
r.insert(
Method::GET,
format!("{}{}", ADMIN_PREFIX, REPLACEMENT_RECOVERY_STATUS_ROUTE_SUFFIX).as_str(),
@@ -1414,7 +1430,7 @@ fn encode_replacement_recovery_status(response: &ReplacementRecoveryStatusRespon
})
}
async fn validate_heal_admin_request(req: &S3Request<Body>) -> S3Result<()> {
async fn authenticate_heal_admin_request(req: &S3Request<Body>) -> S3Result<String> {
let Some(input_cred) = req.credentials.as_ref() else {
return Err(s3_error!(InvalidRequest, "authentication required"));
};
@@ -1429,7 +1445,12 @@ async fn validate_heal_admin_request(req: &S3Request<Body>) -> S3Result<()> {
vec![Action::AdminAction(AdminAction::HealAdminAction)],
req.extensions.get::<Option<RemoteAddr>>().and_then(|opt| opt.map(|a| a.0)),
)
.await
.await?;
Ok(cred.access_key)
}
async fn validate_heal_admin_request(req: &S3Request<Body>) -> S3Result<()> {
authenticate_heal_admin_request(req).await.map(|_| ())
}
pub struct HealHandler {}
@@ -1576,6 +1597,248 @@ impl Operation for HealHandler {
pub struct BackgroundHealStatusHandler {}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct MrfLegacyResponsibilityActionRequest {
action: MrfLegacyResponsibilityAction,
responsibility_id: uuid::Uuid,
expected_bucket_incarnation_id: uuid::Uuid,
#[serde(default)]
acknowledge_unverified_data_risk: Option<bool>,
#[serde(default)]
acknowledge_unknown_source_incarnation: Option<bool>,
#[serde(default)]
acknowledge_bucket_incarnation_mismatch: Option<bool>,
#[serde(default)]
reason: Option<String>,
#[serde(default)]
reference: Option<String>,
#[serde(default)]
request_id: Option<uuid::Uuid>,
}
#[derive(Clone, Copy, Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
enum MrfLegacyResponsibilityAction {
Refresh,
Recheck,
AcceptUnverifiedRisk,
}
fn validate_mrf_legacy_responsibility_action(action: &MrfLegacyResponsibilityActionRequest) -> Result<(), &'static str> {
if action.responsibility_id.is_nil() || action.expected_bucket_incarnation_id.is_nil() {
return Err("responsibility and bucket incarnation IDs must be non-nil");
}
match action.action {
MrfLegacyResponsibilityAction::Refresh | MrfLegacyResponsibilityAction::Recheck => {
if action.acknowledge_unverified_data_risk.is_some()
|| action.acknowledge_unknown_source_incarnation.is_some()
|| action.acknowledge_bucket_incarnation_mismatch.is_some()
|| action.reason.is_some()
|| action.reference.is_some()
{
return Err("recheck does not accept risk-acknowledgment fields");
}
}
MrfLegacyResponsibilityAction::AcceptUnverifiedRisk => {
if action.acknowledge_unverified_data_risk != Some(true) {
return Err("acknowledgeUnverifiedDataRisk must be true");
}
if action.acknowledge_unknown_source_incarnation.is_none() || action.acknowledge_bucket_incarnation_mismatch.is_none()
{
return Err("both source-incarnation acknowledgment fields are required");
}
if action
.reason
.as_deref()
.is_none_or(|value| value.trim().is_empty() || value.len() > 1024)
{
return Err("reason is required and limited to 1024 UTF-8 bytes");
}
if action
.reference
.as_deref()
.is_none_or(|value| value.trim().is_empty() || value.len() > 256)
{
return Err("reference is required and limited to 256 UTF-8 bytes");
}
if action.request_id.is_none_or(|request_id| request_id.is_nil()) {
return Err("requestId must be non-nil");
}
}
}
Ok(())
}
fn parse_mrf_legacy_responsibility_query(uri: &Uri) -> S3Result<(Option<uuid::Uuid>, usize)> {
let mut cursor = None;
let mut limit = MRF_RESPONSIBILITY_LIST_DEFAULT_LIMIT;
let mut seen = HashSet::with_capacity(2);
if let Some(query) = uri.query() {
for (key, value) in url::form_urlencoded::parse(query.as_bytes()) {
match key.as_ref() {
"cursor" if seen.insert("cursor") => {
cursor = Some(
value
.parse()
.map_err(|_| admin_error(S3ErrorCode::InvalidArgument, "cursor must be a UUID"))?,
);
}
"limit" if seen.insert("limit") => {
limit = value
.parse::<usize>()
.map_err(|_| admin_error(S3ErrorCode::InvalidArgument, "limit must be an integer between 1 and 256"))?;
if limit == 0 || limit > MRF_RESPONSIBILITY_LIST_MAX_LIMIT {
return Err(admin_error(S3ErrorCode::InvalidArgument, "limit must be between 1 and 256"));
}
}
"cursor" | "limit" => {
return Err(admin_error(S3ErrorCode::InvalidArgument, "duplicate MRF responsibility query parameter"));
}
_ => return Err(admin_error(S3ErrorCode::InvalidArgument, "unknown MRF responsibility query parameter")),
}
}
}
Ok((cursor, limit))
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct MrfLegacyResponsibilityActionResponse {
responsibility_id: uuid::Uuid,
state: &'static str,
data_verified: bool,
durable_responsibility_retained: bool,
}
fn map_mrf_lifecycle_control_error(error: rustfs_heal::heal::mrf_queue::MrfLifecycleControlError) -> s3s::S3Error {
use rustfs_heal::heal::mrf_queue::MrfLifecycleControlError;
match error {
MrfLifecycleControlError::Unavailable | MrfLifecycleControlError::Persistence => {
admin_error(S3ErrorCode::ServiceUnavailable, "MRF lifecycle operation could not be completed")
}
MrfLifecycleControlError::StaleGeneration => {
admin_error(S3ErrorCode::InvalidRequest, "MRF responsibility generation changed")
}
MrfLifecycleControlError::NotHeld => {
admin_error(S3ErrorCode::InvalidRequest, "MRF responsibility is not held as unverified legacy")
}
MrfLifecycleControlError::IncarnationChanged => admin_error(
S3ErrorCode::InvalidRequest,
"bucket incarnation changed; refresh the MRF responsibility listing",
),
MrfLifecycleControlError::InvalidAction(reason) => admin_error(S3ErrorCode::InvalidRequest, reason),
}
}
pub struct MrfLegacyResponsibilitiesHandler {}
#[async_trait::async_trait]
impl Operation for MrfLegacyResponsibilitiesHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
validate_heal_admin_request(&req).await?;
let (cursor, limit) = parse_mrf_legacy_responsibility_query(&req.uri)?;
let snapshot = timeout(
Duration::from_secs(5),
rustfs_heal::heal::mrf_queue::list_legacy_responsibilities(cursor, limit),
)
.await
.map_err(|_| admin_error(S3ErrorCode::ServiceUnavailable, "MRF responsibility listing timed out"))?
.map_err(map_mrf_lifecycle_control_error)?;
let body = serde_json::to_vec(&snapshot)
.map_err(|_| admin_error(S3ErrorCode::InternalError, "failed to encode MRF responsibility listing"))?;
Ok(json_response(StatusCode::OK, body))
}
}
pub struct MrfLegacyResponsibilitiesActionHandler {}
#[async_trait::async_trait]
impl Operation for MrfLegacyResponsibilitiesActionHandler {
async fn call(&self, mut req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let actor = authenticate_heal_admin_request(&req).await?;
let bytes = req
.input
.store_all_limited(rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE)
.await
.map_err(|_| admin_error(S3ErrorCode::InvalidRequest, "MRF responsibility action body is too large or unreadable"))?;
let action: MrfLegacyResponsibilityActionRequest = serde_json::from_slice(&bytes)
.map_err(|_| admin_error(S3ErrorCode::InvalidRequest, "invalid MRF responsibility action body"))?;
validate_mrf_legacy_responsibility_action(&action).map_err(|reason| admin_error(S3ErrorCode::InvalidRequest, reason))?;
let (state, data_verified, durable_responsibility_retained) = match action.action {
MrfLegacyResponsibilityAction::Refresh => {
timeout(
Duration::from_secs(30),
rustfs_heal::heal::mrf_queue::refresh_legacy_responsibility(
action.responsibility_id,
action.expected_bucket_incarnation_id,
),
)
.await
.map_err(|_| admin_error(S3ErrorCode::ServiceUnavailable, "MRF responsibility refresh timed out"))?
.map_err(map_mrf_lifecycle_control_error)?;
("refreshed", false, true)
}
MrfLegacyResponsibilityAction::Recheck => {
timeout(
Duration::from_secs(30),
rustfs_heal::heal::mrf_queue::recheck_legacy_responsibility(
action.responsibility_id,
action.expected_bucket_incarnation_id,
),
)
.await
.map_err(|_| admin_error(S3ErrorCode::ServiceUnavailable, "MRF responsibility recheck timed out"))?
.map_err(map_mrf_lifecycle_control_error)?;
("active", false, true)
}
MrfLegacyResponsibilityAction::AcceptUnverifiedRisk => {
let reason = action
.reason
.ok_or_else(|| admin_error(S3ErrorCode::InvalidRequest, "reason is required"))?;
let reference = action
.reference
.ok_or_else(|| admin_error(S3ErrorCode::InvalidRequest, "reference is required"))?;
let request_id = action
.request_id
.ok_or_else(|| admin_error(S3ErrorCode::InvalidRequest, "requestId is required"))?;
timeout(
Duration::from_secs(30),
rustfs_heal::heal::mrf_queue::accept_unverified_legacy_risk(
rustfs_heal::heal::mrf_queue::MrfLegacyRiskAcceptanceRequest {
responsibility_id: action.responsibility_id,
expected_bucket_incarnation_id: action.expected_bucket_incarnation_id,
acknowledge_unknown_source_incarnation: action.acknowledge_unknown_source_incarnation == Some(true),
acknowledge_incarnation_mismatch: action.acknowledge_bucket_incarnation_mismatch == Some(true),
actor,
reason,
reference,
request_id,
},
),
)
.await
.map_err(|_| {
admin_error(
S3ErrorCode::ServiceUnavailable,
"MRF risk disposition timed out; read status before retrying",
)
})?
.map_err(map_mrf_lifecycle_control_error)?;
("operatorAcceptedUnverified", false, true)
}
};
let body = serde_json::to_vec(&MrfLegacyResponsibilityActionResponse {
responsibility_id: action.responsibility_id,
state,
data_verified,
durable_responsibility_retained,
})
.map_err(|_| admin_error(S3ErrorCode::InternalError, "failed to encode MRF responsibility action response"))?;
Ok(json_response(StatusCode::OK, body))
}
}
#[async_trait::async_trait]
impl Operation for BackgroundHealStatusHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
@@ -1667,13 +1930,14 @@ mod tests {
use super::extract_heal_init_params;
use super::{
BackgroundHealCoverage, BackgroundHealCoverageReason, BackgroundHealProgress, HealInitParams, HealResp, HealRuntimeState,
aggregate_cluster_heal_status, aggregate_replacement_recovery_cluster_status, background_heal_runtime_state,
build_heal_channel_request, build_replacement_recovery_status_response, encode_background_heal_status,
encode_heal_control_path, encode_heal_start_success, encode_heal_task_status, execute_after_heal_start_preflight,
heal_channel_response_items, heal_channel_response_progress, heal_channel_response_summary, heal_control_response_id,
json_response, map_heal_response, merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status,
query_peer_replacement_recovery_status, read_cluster_heal_status, reject_heal_admission, validate_heal_request_mode,
validate_heal_target,
MrfLegacyResponsibilityAction, MrfLegacyResponsibilityActionRequest, aggregate_cluster_heal_status,
aggregate_replacement_recovery_cluster_status, background_heal_runtime_state, build_heal_channel_request,
build_replacement_recovery_status_response, encode_background_heal_status, encode_heal_control_path,
encode_heal_start_success, encode_heal_task_status, execute_after_heal_start_preflight, heal_channel_response_items,
heal_channel_response_progress, heal_channel_response_summary, heal_control_response_id, json_response,
map_heal_response, merge_peer_heal_statuses, parse_mrf_legacy_responsibility_query, peer_topology_complete,
query_peer_heal_status, query_peer_replacement_recovery_status, read_cluster_heal_status, reject_heal_admission,
validate_heal_request_mode, validate_heal_target, validate_mrf_legacy_responsibility_action,
};
use crate::storage::rpc::node_service::heal::{
NodeHealProgress, NodeHealStatusSnapshot, NodeReplacementRecoveryStatusSnapshot, encode_node_replacement_recovery_status,
@@ -1692,6 +1956,92 @@ mod tests {
};
use serde_json::json;
use std::sync::atomic::{AtomicBool, Ordering};
use uuid::Uuid;
#[test]
fn mrf_legacy_risk_action_requires_explicit_acknowledgment_and_audit_fields() {
let base = json!({
"action": "acceptUnverifiedRisk",
"responsibilityId": Uuid::new_v4(),
"expectedBucketIncarnationId": Uuid::new_v4(),
"acknowledgeUnverifiedDataRisk": true,
"acknowledgeUnknownSourceIncarnation": false,
"acknowledgeBucketIncarnationMismatch": false,
"reason": "Verified against the retained canonical backup",
"reference": "INC-1234",
"requestId": Uuid::new_v4(),
});
let action: MrfLegacyResponsibilityActionRequest = serde_json::from_value(base.clone()).expect("valid risk action");
assert!(validate_mrf_legacy_responsibility_action(&action).is_ok());
let mut missing_ack = base.clone();
missing_ack["acknowledgeUnverifiedDataRisk"] = json!(false);
let missing_ack: MrfLegacyResponsibilityActionRequest = serde_json::from_value(missing_ack).expect("well-formed request");
assert!(validate_mrf_legacy_responsibility_action(&missing_ack).is_err());
let mut missing_reference = base.clone();
missing_reference.as_object_mut().expect("request object").remove("reference");
let missing_reference: MrfLegacyResponsibilityActionRequest =
serde_json::from_value(missing_reference).expect("well-formed request");
assert!(validate_mrf_legacy_responsibility_action(&missing_reference).is_err());
let mut missing_source_ack = base.clone();
missing_source_ack
.as_object_mut()
.expect("request object")
.remove("acknowledgeUnknownSourceIncarnation");
let missing_source_ack: MrfLegacyResponsibilityActionRequest =
serde_json::from_value(missing_source_ack).expect("well-formed request");
assert!(validate_mrf_legacy_responsibility_action(&missing_source_ack).is_err());
let mut oversized_reason = base.clone();
oversized_reason["reason"] = json!("x".repeat(1025));
let oversized_reason: MrfLegacyResponsibilityActionRequest =
serde_json::from_value(oversized_reason).expect("well-formed request");
assert!(validate_mrf_legacy_responsibility_action(&oversized_reason).is_err());
let mut unknown_field = base;
unknown_field["clearJournal"] = json!(true);
assert!(serde_json::from_value::<MrfLegacyResponsibilityActionRequest>(unknown_field).is_err());
let recheck = MrfLegacyResponsibilityActionRequest {
action: MrfLegacyResponsibilityAction::Recheck,
responsibility_id: Uuid::new_v4(),
expected_bucket_incarnation_id: Uuid::new_v4(),
acknowledge_unverified_data_risk: None,
acknowledge_unknown_source_incarnation: None,
acknowledge_bucket_incarnation_mismatch: None,
reason: None,
reference: None,
request_id: None,
};
assert!(validate_mrf_legacy_responsibility_action(&recheck).is_ok());
let refresh = MrfLegacyResponsibilityActionRequest {
action: MrfLegacyResponsibilityAction::Refresh,
..recheck
};
assert!(validate_mrf_legacy_responsibility_action(&refresh).is_ok());
}
#[test]
fn mrf_legacy_listing_query_is_bounded_and_rejects_ambiguous_parameters() {
let cursor = Uuid::new_v4();
let uri = Uri::try_from(format!("/rustfs/admin/v4/heal/mrf/responsibilities?cursor={cursor}&limit=256"))
.expect("valid responsibility URI");
assert_eq!(parse_mrf_legacy_responsibility_query(&uri).expect("valid query"), (Some(cursor), 256));
for query in [
"limit=0",
"limit=257",
"limit=-1",
"cursor=bad",
"limit=1&limit=2",
"object=secret",
] {
let uri = Uri::try_from(format!("/rustfs/admin/v4/heal/mrf/responsibilities?{query}")).expect("well-formed URI");
assert!(parse_mrf_legacy_responsibility_query(&uri).is_err(), "accepted invalid query {query}");
}
}
use time::{OffsetDateTime, format_description::well_known::Rfc3339};
use tokio::sync::mpsc;
use tokio::time::Duration;
+12
View File
@@ -354,6 +354,18 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
admin(HttpMethod::Post, "/rustfs/admin/v3/heal/{bucket}", HEAL, RouteRiskLevel::High),
admin(HttpMethod::Post, "/rustfs/admin/v3/heal/{bucket}/{*prefix}", HEAL, RouteRiskLevel::High),
admin(HttpMethod::Post, "/rustfs/admin/v3/background-heal/status", HEAL, RouteRiskLevel::High),
admin(
HttpMethod::Get,
"/rustfs/admin/v4/heal/mrf/responsibilities",
HEAL,
RouteRiskLevel::Sensitive,
),
admin(
HttpMethod::Post,
"/rustfs/admin/v4/heal/mrf/responsibilities/actions",
HEAL,
RouteRiskLevel::High,
),
admin(
HttpMethod::Get,
"/rustfs/admin/v4/heal/replacement-recovery",
@@ -206,6 +206,8 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
admin_route_sample(Method::POST, "/v3/heal/{bucket}", "/v3/heal/test-bucket"),
admin_route_sample(Method::POST, "/v3/heal/{bucket}/{*prefix}", "/v3/heal/test-bucket/prefix"),
admin_route(Method::POST, "/v3/background-heal/status"),
admin_route(Method::GET, "/v4/heal/mrf/responsibilities"),
admin_route(Method::POST, "/v4/heal/mrf/responsibilities/actions"),
admin_route(Method::GET, "/v4/heal/replacement-recovery"),
admin_route(Method::GET, "/v3/tier"),
admin_route(Method::GET, "/v3/tier-stats"),
@@ -1356,6 +1358,8 @@ fn test_register_routes_cover_representative_admin_paths() {
assert_route(&router, Method::POST, &admin_path("/v3/heal/test-bucket"));
assert_route(&router, Method::POST, &admin_path("/v3/heal/test-bucket/prefix"));
assert_route(&router, Method::POST, &admin_path("/v3/background-heal/status"));
assert_route(&router, Method::GET, &admin_path("/v4/heal/mrf/responsibilities"));
assert_route(&router, Method::POST, &admin_path("/v4/heal/mrf/responsibilities/actions"));
assert_route(&router, Method::GET, &admin_path("/v4/heal/replacement-recovery"));
assert_route(&router, Method::GET, &admin_path("/v3/tier"));
@@ -1464,6 +1468,8 @@ fn test_admin_alias_paths_match_existing_admin_routes() {
(Method::POST, compat_admin_alias_path("/v3/heal/test-bucket")),
(Method::POST, compat_admin_alias_path("/v3/heal/test-bucket/prefix")),
(Method::POST, compat_admin_alias_path("/v3/background-heal/status")),
(Method::GET, compat_admin_alias_path("/v4/heal/mrf/responsibilities")),
(Method::POST, compat_admin_alias_path("/v4/heal/mrf/responsibilities/actions")),
(Method::GET, compat_admin_alias_path("/v3/tier/HOT")),
(Method::GET, compat_admin_alias_path("/v3/export-bucket-metadata")),
(Method::PUT, compat_admin_alias_path("/v3/import-bucket-metadata")),
@@ -517,36 +517,38 @@ mod tests {
assert!(!local.delete_marker);
}
#[tokio::test]
#[test]
#[serial_test::serial]
async fn write_back_rejects_unsupported_topology_before_any_mutation() {
let (_dir, _paths, store) = crate::app::gating_test_env::isolated_multi_pool_ecstore().await;
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
let bucket = "odm-unsupported";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket");
let write_back = OnDemandMigrationWriteBack::new();
let req = request(
bucket,
store.bucket_incarnation_id(bucket).await.expect("bucket incarnation"),
"key",
source_head(b"source"),
);
assert!(matches!(
write_back.put_object(&req, body_stream(b"source")).await,
Err(WriteBackError::Unsupported(_))
));
assert!(matches!(
write_back.create_multipart_upload(&req).await,
Err(WriteBackError::Unsupported(_))
));
assert!(matches!(
write_back.complete_multipart_upload(&req, "no-session", Vec::new()).await,
Err(WriteBackError::Unsupported(_))
));
assert_nothing_left(&store, bucket, "key").await;
fn write_back_rejects_unsupported_topology_before_any_mutation() {
crate::app::gating_test_env::run_large_stack_test("odm-unsupported-topology", || async {
let (_dir, _paths, store) = crate::app::gating_test_env::isolated_multi_pool_ecstore().await;
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
let bucket = "odm-unsupported";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket");
let write_back = OnDemandMigrationWriteBack::new();
let req = request(
bucket,
store.bucket_incarnation_id(bucket).await.expect("bucket incarnation"),
"key",
source_head(b"source"),
);
assert!(matches!(
write_back.put_object(&req, body_stream(b"source")).await,
Err(WriteBackError::Unsupported(_))
));
assert!(matches!(
write_back.create_multipart_upload(&req).await,
Err(WriteBackError::Unsupported(_))
));
assert!(matches!(
write_back.complete_multipart_upload(&req, "no-session", Vec::new()).await,
Err(WriteBackError::Unsupported(_))
));
assert_nothing_left(&store, bucket, "key").await;
});
}
#[tokio::test]