From cebc28f678d12e951cb48e68bc1e0c3d61a70eed Mon Sep 17 00:00:00 2001 From: cxymds Date: Tue, 4 Aug 2026 11:49:16 +0800 Subject: [PATCH] fix(replication): fence MRF journal updates (#5686) --- .../replication/replication_config_store.rs | 30 ++ .../bucket/replication/replication_pool.rs | 380 +++++++++++++++--- 2 files changed, 352 insertions(+), 58 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_config_store.rs b/crates/ecstore/src/bucket/replication/replication_config_store.rs index f0db89f2e..25402c394 100644 --- a/crates/ecstore/src/bucket/replication/replication_config_store.rs +++ b/crates/ecstore/src/bucket/replication/replication_config_store.rs @@ -52,6 +52,13 @@ impl ReplicationConfigStore { .await } + pub(crate) async fn read_no_lock_with_metadata_preserve_empty(api: Arc, file: &str) -> Result<(Vec, ObjectInfo)> + where + S: ReplicationObjectIO, + { + com::read_config_no_lock_preserve_empty_with_metadata(api, file).await + } + pub(crate) async fn save(api: Arc, file: &str, data: Vec) -> Result<()> where S: ReplicationObjectIO, @@ -87,4 +94,27 @@ impl ReplicationConfigStore { ) .await } + + pub(crate) async fn save_conditional_no_lock( + api: Arc, + file: &str, + data: Vec, + http_preconditions: HTTPPreconditions, + ) -> Result<()> + where + S: ReplicationObjectIO, + { + com::save_config_with_opts_quiet( + api, + file, + data, + &ObjectOptions { + max_parity: true, + no_lock: true, + http_preconditions: Some(http_preconditions), + ..Default::default() + }, + ) + .await + } } diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 78b945670..dba231b95 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -586,6 +586,27 @@ fn ensure_force_delete_journal_lock_held(lock_lost: bool) -> Result<(), EcstoreE Ok(()) } +fn ensure_mrf_journal_lock_held(lock_lost: bool) -> Result<(), EcstoreError> { + if lock_lost { + return Err(EcstoreError::other("MRF journal lock lost before conditional update")); + } + Ok(()) +} + +fn mrf_journal_preconditions(etag: Option<&str>, exists: bool) -> Option { + if exists { + etag.filter(|value| !value.trim().is_empty()).map(|etag| HTTPPreconditions { + if_match: Some(etag.to_string()), + ..Default::default() + }) + } else { + Some(HTTPPreconditions { + if_none_match: Some("*".to_string()), + ..Default::default() + }) + } +} + async fn write_mrf_journal_snapshot( storage: Arc, desired: &[MrfReplicateEntry], @@ -598,8 +619,8 @@ async fn write_mrf_journal_snapshot( .new_ns_lock(ReplicationMetadataStore::rustfs_meta_bucket(), file) .await?; let guard = lock.get_write_lock(ReplicationLockTiming::acquire_timeout()).await?; - let current = ReplicationConfigStore::read_no_lock_with_metadata(storage.clone(), file).await; - let etag = match current { + let current = ReplicationConfigStore::read_no_lock_with_metadata_preserve_empty(storage.clone(), file).await; + let (etag, exists) = match current { Ok((data, object_info)) => { if saw_conflict { let current = decode_mrf_file(&data)?; @@ -609,23 +630,16 @@ async fn write_mrf_journal_snapshot( } } } - object_info.etag + (object_info.etag, true) } - Err(EcstoreError::ConfigNotFound) => None, + Err(EcstoreError::ConfigNotFound) => (None, false), Err(err) => return Err(err), }; if guard.is_lock_lost() { return Err(EcstoreError::other("MRF journal namespace lock was lost before commit")); } - let preconditions = match etag.filter(|value| !value.trim().is_empty()) { - Some(etag) => HTTPPreconditions { - if_match: Some(etag), - ..Default::default() - }, - None => HTTPPreconditions { - if_none_match: Some("*".to_string()), - ..Default::default() - }, + let Some(preconditions) = mrf_journal_preconditions(etag.as_deref(), exists) else { + return Err(EcstoreError::other("MRF journal has no ETag for conditional update")); }; let data = if merged.is_empty() { Vec::new() @@ -2255,7 +2269,7 @@ async fn quarantine_mrf_file(storage: &Arc, data: &[u8 continue; } }; - let _guard = match lock.get_write_lock(ReplicationLockTiming::acquire_timeout()).await { + let guard = match lock.get_write_lock(ReplicationLockTiming::acquire_timeout()).await { Ok(guard) => guard, Err(error) => { warn!( @@ -2269,14 +2283,43 @@ async fn quarantine_mrf_file(storage: &Arc, data: &[u8 continue; } }; - match ReplicationConfigStore::read_no_lock(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await { + match ReplicationConfigStore::read_no_lock_with_metadata_preserve_empty( + storage.clone(), + ReplicationMetadataStore::MRF_REPLICATION_FILE, + ) + .await + { Err(EcstoreError::ConfigNotFound) => return, - Ok(current) if current != data => return, - Ok(_) => { - match ReplicationConfigStore::save_no_lock( + Ok((current, _)) if current != data => return, + Ok((_, object_info)) => { + let Some(preconditions) = mrf_journal_preconditions(object_info.etag.as_deref(), true) else { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + "Cannot clear the corrupt MRF recovery path without an ETag; retrying" + ); + drop(guard); + tokio::time::sleep(retry_delay).await; + retry_delay = retry_delay.saturating_mul(2).min(MRF_RETRY_MAX_DELAY); + continue; + }; + if let Err(error) = ensure_mrf_journal_lock_held(guard.is_lock_lost()) { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %error, + "MRF journal lock was lost before clearing the corrupt recovery path; retrying" + ); + drop(guard); + tokio::time::sleep(retry_delay).await; + retry_delay = retry_delay.saturating_mul(2).min(MRF_RETRY_MAX_DELAY); + continue; + } + match ReplicationConfigStore::save_conditional_no_lock( storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE, Vec::new(), + preconditions, ) .await { @@ -2296,7 +2339,7 @@ async fn quarantine_mrf_file(storage: &Arc, data: &[u8 "Failed to verify the corrupt MRF recovery path before clearing; retrying" ), } - drop(_guard); + drop(guard); tokio::time::sleep(retry_delay).await; retry_delay = retry_delay.saturating_mul(2).min(MRF_RETRY_MAX_DELAY); } @@ -2338,9 +2381,10 @@ async fn flush_mrf_to_disk(entries: &[MrfReplicateEntry], async fn recover_corrupt_mrf_generation( corrupt_generation: &[u8], - entries_to_append: &[MrfReplicateEntry], - known_pending: &[MrfReplicateEntry], + entries: &[MrfReplicateEntry], + preconditions: Option, storage: &Arc, + guard: &rustfs_lock::NamespaceLockGuard, pending_payload: &mut Option, started: Instant, ) -> Option { @@ -2357,10 +2401,7 @@ async fn recover_corrupt_mrf_generation( return None; } - let mut entries = Vec::with_capacity(known_pending.len().saturating_add(entries_to_append.len())); - entries.extend_from_slice(known_pending); - entries.extend_from_slice(entries_to_append); - let data = match encode_mrf_file(&entries) { + let data = match encode_mrf_file(entries) { Ok(data) => data, Err(error) => { observe_mrf_flush_failure(0); @@ -2378,7 +2419,35 @@ async fn recover_corrupt_mrf_generation( digest: mrf_payload_digest(&data), entry_count: entries.len(), }); - if let Err(error) = write_mrf_journal_snapshot(storage.clone(), &entries).await { + let Some(preconditions) = preconditions else { + observe_mrf_flush_failure(duration_millis_u64(started.elapsed())); + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + count = entries.len(), + "Failed to rebuild the active MRF generation because its ETag is unavailable" + ); + return None; + }; + if let Err(error) = ensure_mrf_journal_lock_held(guard.is_lock_lost()) { + observe_mrf_flush_failure(duration_millis_u64(started.elapsed())); + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + count = entries.len(), + error = %error, + "Failed to replace the active MRF generation after losing its namespace lock" + ); + return None; + } + if let Err(error) = ReplicationConfigStore::save_conditional_no_lock( + storage.clone(), + ReplicationMetadataStore::MRF_REPLICATION_FILE, + data, + preconditions, + ) + .await + { observe_mrf_flush_failure(duration_millis_u64(started.elapsed())); warn!( component = LOG_COMPONENT_ECSTORE, @@ -2420,7 +2489,7 @@ async fn append_mrf_entries_to_disk( return None; } }; - let _guard = match lock.get_write_lock(ReplicationLockTiming::acquire_timeout()).await { + let guard = match lock.get_write_lock(ReplicationLockTiming::acquire_timeout()).await { Ok(guard) => guard, Err(error) => { warn!( @@ -2433,20 +2502,24 @@ async fn append_mrf_entries_to_disk( } }; - let current = - match ReplicationConfigStore::read_no_lock(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await { - Ok(data) => data, - Err(EcstoreError::ConfigNotFound) => Vec::new(), - Err(error) => { - warn!( - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION, - error = %error, - "Failed to read MRF backlog before appending a capped entry" - ); - return None; - } - }; + let (current, current_etag, current_exists) = match ReplicationConfigStore::read_no_lock_with_metadata_preserve_empty( + storage.clone(), + ReplicationMetadataStore::MRF_REPLICATION_FILE, + ) + .await + { + Ok((data, object_info)) => (data, object_info.etag, true), + Err(EcstoreError::ConfigNotFound) => (Vec::new(), None, false), + Err(error) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %error, + "Failed to read MRF backlog before appending a capped entry" + ); + return None; + } + }; if pending_payload .as_ref() @@ -2465,8 +2538,19 @@ async fn append_mrf_entries_to_disk( error = %error, "Failed to decode MRF backlog before appending a capped entry" ); - return recover_corrupt_mrf_generation(¤t, entries_to_append, known_pending, storage, pending_payload, started) - .await; + let mut recovery_entries = Vec::with_capacity(known_pending.len().saturating_add(entries_to_append.len())); + recovery_entries.extend_from_slice(known_pending); + recovery_entries.extend_from_slice(entries_to_append); + return recover_corrupt_mrf_generation( + ¤t, + &recovery_entries, + mrf_journal_preconditions(current_etag.as_deref(), current_exists), + storage, + &guard, + pending_payload, + started, + ) + .await; } }; if let Some(pending) = pending_payload.as_ref() @@ -2506,8 +2590,36 @@ async fn append_mrf_entries_to_disk( digest: mrf_payload_digest(&data), entry_count: entries.len(), }); - if let Err(error) = - ReplicationConfigStore::save_no_lock(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE, data).await + let Some(preconditions) = mrf_journal_preconditions(current_etag.as_deref(), current_exists) else { + let duration_millis = duration_millis_u64(started.elapsed()); + observe_mrf_flush_failure(duration_millis); + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + count = entries.len(), + "Failed to append capped MRF entries because the current generation has no ETag" + ); + return None; + }; + if let Err(error) = ensure_mrf_journal_lock_held(guard.is_lock_lost()) { + let duration_millis = duration_millis_u64(started.elapsed()); + observe_mrf_flush_failure(duration_millis); + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + count = entries.len(), + error = %error, + "Failed to append capped MRF entries after losing the namespace lock" + ); + return None; + } + if let Err(error) = ReplicationConfigStore::save_conditional_no_lock( + storage.clone(), + ReplicationMetadataStore::MRF_REPLICATION_FILE, + data, + preconditions, + ) + .await { let duration_millis = duration_millis_u64(started.elapsed()); observe_mrf_flush_failure(duration_millis); @@ -2960,6 +3072,7 @@ mod tests { struct LoadResyncSharedState { data: StdMutex>, + empty_object_exists: AtomicBool, etag_revision: AtomicUsize, last_put_preconditions: StdMutex>, last_put_no_lock: AtomicBool, @@ -3041,7 +3154,7 @@ mod tests { .lock() .expect("test data lock should not be poisoned") .clone(); - if data.is_empty() { + if data.is_empty() && !self.shared.empty_object_exists.load(Ordering::SeqCst) { return Err(EcstoreError::FileNotFound); } let size = i64::try_from(data.len()).expect("test metadata length should fit i64"); @@ -3083,6 +3196,7 @@ mod tests { .lock() .expect("test data lock should not be poisoned") .is_empty() + && !self.shared.empty_object_exists.load(Ordering::SeqCst) { None } else { @@ -3316,8 +3430,15 @@ mod tests { } async fn new_test_replication_pool(storage: Arc) -> Arc> { + new_test_replication_pool_with_mrf_capacity(storage, 1).await + } + + async fn new_test_replication_pool_with_mrf_capacity( + storage: Arc, + mrf_save_capacity: usize, + ) -> Arc> { let (mrf_replica_tx, mrf_replica_rx) = mpsc::channel(1); - let (mrf_save_tx, mrf_save_rx) = mpsc::channel(1); + let (mrf_save_tx, mrf_save_rx) = mpsc::channel(mrf_save_capacity); let (mrf_worker_kill_tx, _) = mpsc::channel(1); let (mrf_stop_tx, _) = mpsc::channel(1); @@ -3575,6 +3696,7 @@ mod tests { fn empty_resync_shared_state() -> Arc { Arc::new(LoadResyncSharedState { data: StdMutex::new(Vec::new()), + empty_object_exists: AtomicBool::new(false), etag_revision: AtomicUsize::new(0), last_put_preconditions: StdMutex::new(None), last_put_no_lock: AtomicBool::new(false), @@ -4230,6 +4352,7 @@ mod tests { temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { let shared = Arc::new(LoadResyncSharedState { data: StdMutex::new(load_resync_test_metadata()), + empty_object_exists: AtomicBool::new(false), etag_revision: AtomicUsize::new(1), last_put_preconditions: StdMutex::new(None), last_put_no_lock: AtomicBool::new(false), @@ -4371,6 +4494,15 @@ mod tests { .find(|(file, _)| file == ReplicationMetadataStore::MRF_REPLICATION_FILE) .expect("active MRF path should be cleared after quarantine"); assert!(marker.1.is_empty(), "the active MRF path should be marked absent"); + let preconditions = shared + .last_put_preconditions + .lock() + .expect("test preconditions lock should not be poisoned") + .clone() + .expect("active MRF cleanup should be conditional"); + assert_eq!(preconditions.if_match_value(), Some("mrf-0")); + assert_eq!(preconditions.if_none_match_value(), None); + assert!(shared.last_put_no_lock.load(Ordering::SeqCst)); } #[tokio::test] @@ -4849,6 +4981,127 @@ mod tests { ); } + #[tokio::test] + async fn mrf_capped_append_retries_after_a_conditional_generation_conflict() { + let shared = empty_resync_shared_state(); + let initial = MrfReplicateEntry { + bucket: "mrf-append-cas".to_string(), + object: "retained".to_string(), + op: MrfOpKind::Object, + ..Default::default() + }; + let concurrent = MrfReplicateEntry { + object: "concurrent".to_string(), + ..initial.clone() + }; + let appended = MrfReplicateEntry { + object: "new-failure".to_string(), + ..initial.clone() + }; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(std::slice::from_ref(&initial)).expect("initial MRF backlog should encode"); + shared + .conditional_write_replacements + .lock() + .expect("test replacement lock should not be poisoned") + .push_back(encode_mrf_file(&[initial.clone(), concurrent.clone()]).expect("replacement should encode")); + let storage = Arc::new(LoadResyncNodeStore::new("mrf-append-cas", shared.clone())); + let mut pending_payload = None; + + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&appended), &storage, &mut pending_payload, &[]) + .await + .is_none(), + "a stale capped append must not overwrite a concurrent generation" + ); + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&appended), &storage, &mut pending_payload, &[]) + .await + .is_some(), + "the capped append should retry against the concurrent generation" + ); + + let entries = decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) + .expect("persisted MRF backlog should decode"); + assert_eq!(entries.len(), 3); + assert_eq!(entries[0].object, initial.object); + assert_eq!(entries[1].object, concurrent.object); + assert_eq!(entries[2].object, appended.object); + } + + #[tokio::test] + async fn mrf_capped_append_rejects_an_existing_generation_without_an_etag() { + let shared = empty_resync_shared_state(); + let initial = MrfReplicateEntry { + bucket: "mrf-append-no-etag".to_string(), + object: "retained".to_string(), + op: MrfOpKind::Object, + ..Default::default() + }; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(std::slice::from_ref(&initial)).expect("initial MRF backlog should encode"); + shared.omit_etag.store(true, Ordering::SeqCst); + let storage = Arc::new(LoadResyncNodeStore::new("mrf-append-no-etag", shared.clone())); + let appended = MrfReplicateEntry { + object: "new-failure".to_string(), + ..initial.clone() + }; + let mut pending_payload = None; + + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&appended), &storage, &mut pending_payload, &[]) + .await + .is_none(), + "an existing MRF generation without an ETag must fail closed" + ); + assert!( + shared + .writes + .lock() + .expect("test writes lock should not be poisoned") + .is_empty(), + "the un-fenced capped append must not write" + ); + let entries = decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) + .expect("the original MRF backlog should remain readable"); + assert_eq!(entries.len(), 1); + assert_eq!(entries[0].object, initial.object); + } + + #[tokio::test] + async fn mrf_capped_append_replaces_an_existing_empty_generation_conditionally() { + let shared = empty_resync_shared_state(); + shared.empty_object_exists.store(true, Ordering::SeqCst); + let storage = Arc::new(LoadResyncNodeStore::new("mrf-append-empty", shared.clone())); + let appended = MrfReplicateEntry { + bucket: "mrf-append-empty".to_string(), + object: "new-failure".to_string(), + op: MrfOpKind::Object, + ..Default::default() + }; + let mut pending_payload = None; + + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&appended), &storage, &mut pending_payload, &[]) + .await + .is_some(), + "an existing empty MRF generation should be replaced" + ); + let preconditions = shared + .last_put_preconditions + .lock() + .expect("test preconditions lock should not be poisoned") + .clone() + .expect("the empty generation replacement should be conditional"); + assert_eq!(preconditions.if_match_value(), Some("mrf-0")); + assert_eq!(preconditions.if_none_match_value(), None); + assert!(shared.last_put_no_lock.load(Ordering::SeqCst)); + let entries = decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) + .expect("the replaced MRF backlog should decode"); + assert_eq!(entries.len(), 1); + assert_eq!(entries[0].object, appended.object); + } + #[tokio::test] async fn mrf_capped_append_recovers_a_late_corrupt_generation() { let shared = empty_resync_shared_state(); @@ -4928,10 +5181,7 @@ mod tests { tokio::time::timeout(Duration::from_secs(30), async { loop { let data = shared.data.lock().expect("test data lock should not be poisoned").clone(); - if decode_mrf_file(&data).is_ok_and(|entries| { - entries.len() == MRF_PENDING_CAP + 1 - && entries.last().is_some_and(|entry| entry.object == "staged-overflow") - }) { + if decode_mrf_file(&data).is_ok_and(|entries| entries.len() == 1 && entries[0].object == "staged-overflow") { break; } tokio::task::yield_now().await; @@ -4964,14 +5214,21 @@ mod tests { encode_mrf_file(&retained).expect("full startup MRF backlog should encode"); shared.delay_first_read.store(true, Ordering::SeqCst); shared.block_next_write.store(true, Ordering::SeqCst); - let pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-recovery-shrink", shared.clone()))).await; - let read_started = shared.first_read_started.notified(); + let pool = new_test_replication_pool_with_mrf_capacity( + Arc::new(LoadResyncNodeStore::new("mrf-recovery-shrink", shared.clone())), + 3, + ) + .await; let write_started = shared.write_started.notified(); pool.start_mrf_persister().await; - tokio::time::timeout(Duration::from_secs(2), read_started) - .await - .expect("startup MRF read should be delayed"); + tokio::time::timeout(Duration::from_secs(2), async { + while shared.read_count.load(Ordering::SeqCst) == 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("startup MRF read should be delayed"); for object in ["staged-overflow-1", "staged-overflow-2", "staged-overflow-3"] { pool.mrf_save_tx .send(MrfReplicateEntry { @@ -5019,8 +5276,8 @@ mod tests { .map(|(_, data)| decode_mrf_file(data).expect("persisted MRF data should decode")) .expect("the capped suffix should be persisted after the pending flush"); assert_eq!(persisted.len(), MRF_PENDING_CAP + 1); - assert_eq!(persisted[MRF_PENDING_CAP - 1].object, "staged-overflow-1"); - assert_eq!(persisted[MRF_PENDING_CAP].object, "staged-overflow-2"); + assert_eq!(persisted[MRF_PENDING_CAP - 2].object, "staged-overflow-1"); + assert_eq!(persisted[MRF_PENDING_CAP - 1].object, "staged-overflow-2"); assert_eq!(persisted.last().expect("capped suffix should be present").object, "staged-overflow-3"); let handle = pool @@ -5711,6 +5968,13 @@ mod tests { assert!(err.to_string().contains("lock lost")); } + #[test] + fn mrf_journal_rejects_a_lost_transaction_lease() { + let err = ensure_mrf_journal_lock_held(true).expect_err("lost MRF lease must fence the journal write"); + + assert!(err.to_string().contains("lock lost")); + } + #[tokio::test] async fn force_delete_journal_rejects_a_stale_conditional_write() { let shared = empty_resync_shared_state();