fix(replication): fence MRF journal updates (#5686)

This commit is contained in:
cxymds
2026-08-04 11:49:16 +08:00
committed by GitHub
parent d6e11cf018
commit cebc28f678
2 changed files with 352 additions and 58 deletions
@@ -52,6 +52,13 @@ impl ReplicationConfigStore {
.await
}
pub(crate) async fn read_no_lock_with_metadata_preserve_empty<S>(api: Arc<S>, file: &str) -> Result<(Vec<u8>, ObjectInfo)>
where
S: ReplicationObjectIO,
{
com::read_config_no_lock_preserve_empty_with_metadata(api, file).await
}
pub(crate) async fn save<S>(api: Arc<S>, file: &str, data: Vec<u8>) -> Result<()>
where
S: ReplicationObjectIO,
@@ -87,4 +94,27 @@ impl ReplicationConfigStore {
)
.await
}
pub(crate) async fn save_conditional_no_lock<S>(
api: Arc<S>,
file: &str,
data: Vec<u8>,
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
}
}
@@ -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<HTTPPreconditions> {
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<S: ReplicationStorage>(
storage: Arc<S>,
desired: &[MrfReplicateEntry],
@@ -598,8 +619,8 @@ async fn write_mrf_journal_snapshot<S: ReplicationStorage>(
.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<S: ReplicationStorage>(
}
}
}
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<S: ReplicationStorage>(storage: &Arc<S>, 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<S: ReplicationStorage>(storage: &Arc<S>, 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<S: ReplicationStorage>(storage: &Arc<S>, 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<S: ReplicationStorage>(entries: &[MrfReplicateEntry],
async fn recover_corrupt_mrf_generation<S: ReplicationStorage>(
corrupt_generation: &[u8],
entries_to_append: &[MrfReplicateEntry],
known_pending: &[MrfReplicateEntry],
entries: &[MrfReplicateEntry],
preconditions: Option<HTTPPreconditions>,
storage: &Arc<S>,
guard: &rustfs_lock::NamespaceLockGuard,
pending_payload: &mut Option<PendingMrfAppend>,
started: Instant,
) -> Option<u64> {
@@ -2357,10 +2401,7 @@ async fn recover_corrupt_mrf_generation<S: ReplicationStorage>(
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<S: ReplicationStorage>(
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<S: ReplicationStorage>(
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<S: ReplicationStorage>(
}
};
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<S: ReplicationStorage>(
error = %error,
"Failed to decode MRF backlog before appending a capped entry"
);
return recover_corrupt_mrf_generation(&current, 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(
&current,
&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<S: ReplicationStorage>(
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<Vec<u8>>,
empty_object_exists: AtomicBool,
etag_revision: AtomicUsize,
last_put_preconditions: StdMutex<Option<HTTPPreconditions>>,
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<LoadResyncNodeStore>) -> Arc<ReplicationPool<LoadResyncNodeStore>> {
new_test_replication_pool_with_mrf_capacity(storage, 1).await
}
async fn new_test_replication_pool_with_mrf_capacity(
storage: Arc<LoadResyncNodeStore>,
mrf_save_capacity: usize,
) -> Arc<ReplicationPool<LoadResyncNodeStore>> {
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<LoadResyncSharedState> {
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();