From a19560da419c6d354a805fc127951cc9b56ad887 Mon Sep 17 00:00:00 2001 From: abdullahnah92 <157595835+abdullahnah92@users.noreply.github.com> Date: Mon, 15 Jun 2026 06:56:50 +0300 Subject: [PATCH] fix(bucket-repl): persist MRF retry queue to disk and reload on startup (#3456) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(bucket-repl): persist MRF retry queue to disk and reload on startup Co-Authored-By: Claude Sonnet 4.6 * fix(bucket-repl): address three blocking MRF issues from review 1. MRF replay loses delete operations — add `MrfOpKind` discriminator to `MrfReplicateEntry` (Object | Delete, default=Object for backward compat). `DeletedObjectReplicationInfo::to_mrf_entry` now persists `op=Delete`, `version_id`, `delete_marker_version_id`, and `delete_marker`. `start_mrf_processor` branches on `op`: delete entries skip `get_object_info` and replay via `schedule_replication_delete` with `ReplicationType::Heal`; object entries follow the existing heal path. 2. `flush_mrf_to_disk` cleared the in-memory batch even on encode/write failure — changed return type to `bool` and callers now only `pending.clear()` on `true`, so a transient storage error retries on the next tick instead of silently dropping the batch. 3. Add focused tests: encode/decode roundtrips for object, delete-marker, versioned-delete, and mixed-batch entries; a routing test confirming op-kind propagates correctly and that the default is Object for legacy files; a legacy-compat test verifying old entries round-trip cleanly through the new format. Co-Authored-By: Claude Sonnet 4.6 * style: fix clippy redundant-clone in MRF tests Replace &[entry.clone()] with std::slice::from_ref(&entry) in two encode_mrf_file call sites flagged by clippy's redundant_clone lint under --all-targets --features rio-v2. Co-Authored-By: Claude Sonnet 4.6 * test(bucket-repl): strengthen legacy MRF compat test with hand-built msgpack The previous mrf_legacy_file_without_op_field_decoded_as_object test round-tripped through encode_mrf_file, so it exercised the new format and never touched a truly-legacy payload. Replace it with a hand-built msgpack payload that genuinely omits the "op", "deleteMarker", and "deleteMarkerVersionID" keys — exactly what the old binary would have written before MrfOpKind existed. The test now fails if #[serde(default)] is removed from the op field, which proves real backward compatibility rather than round-trip stability. Co-Authored-By: Claude Sonnet 4.6 --------- Co-authored-by: Claude Sonnet 4.6 Co-authored-by: houseme --- .../bucket/replication/replication_pool.rs | 475 ++++++++++++++++-- .../replication/replication_resyncer.rs | 136 ++++- crates/filemeta/src/replication.rs | 41 +- 3 files changed, 604 insertions(+), 48 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 2eb558c83..af2c2f575 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -20,15 +20,17 @@ use crate::bucket::replication::ResyncStatusType; use crate::bucket::replication::replicate_delete; use crate::bucket::replication::replicate_object; use crate::bucket::replication::replication_resyncer::{ - BucketReplicationResyncStatus, DeletedObjectReplicationInfo, REPLICATION_DIR, RESYNC_FILE_NAME, ReplicationConfig, - ReplicationResyncer, TargetReplicationResyncStatus, decode_resync_file, get_heal_replicate_object_info, save_resync_status, + BucketReplicationResyncStatus, DeletedObjectReplicationInfo, MRF_REPLICATION_FILE, REPLICATION_DIR, RESYNC_FILE_NAME, + ReplicationConfig, ReplicationResyncer, TargetReplicationResyncStatus, decode_mrf_file, decode_resync_file, encode_mrf_file, + get_heal_replicate_object_info, save_resync_status, }; use crate::bucket::replication::replication_state::ReplicationStats; -use crate::config::com::read_config; +use crate::config::com::{read_config, save_config}; use crate::disk::BUCKET_META_PREFIX; use crate::error::Error as EcstoreError; -use crate::store_api::{NamespaceLocking, ObjectIO, ObjectInfo}; +use crate::store_api::{NamespaceLocking, ObjectIO, ObjectInfo, ObjectOptions}; use lazy_static::lazy_static; +use rustfs_filemeta::MrfOpKind; use rustfs_filemeta::MrfReplicateEntry; use rustfs_filemeta::ReplicateDecision; use rustfs_filemeta::ReplicateObjectInfo; @@ -233,7 +235,8 @@ pub struct ReplicationPool { // MRF (Most Recent Failures) channels mrf_replica_tx: Sender, - mrf_replica_rx: Mutex>>, + // Shared among N MRF workers; Arc allows spawning more than one worker. + mrf_replica_rx: Arc>>, mrf_save_tx: Sender, mrf_save_rx: Mutex>>, @@ -285,7 +288,7 @@ impl ReplicationPool { workers: RwLock::new(Vec::new()), lrg_workers: RwLock::new(Vec::new()), mrf_replica_tx, - mrf_replica_rx: Mutex::new(Some(mrf_replica_rx)), + mrf_replica_rx: Arc::new(Mutex::new(mrf_replica_rx)), mrf_save_tx, mrf_save_rx: Mutex::new(Some(mrf_save_rx)), mrf_worker_kill_tx, @@ -463,50 +466,48 @@ impl ReplicationPool { /// Resizes the failed workers pool pub async fn resize_failed_workers(&self, n: i32) { - // Add workers if needed + // Spawn workers up to n. Each worker shares the receiver via Arc>. + // The mutex is held only while calling recv() — released before processing — so + // all workers process entries concurrently (the dequeue step is serialised but + // the replication I/O is not). while self.mrf_worker_size.load(Ordering::SeqCst) < n { self.mrf_worker_size.fetch_add(1, Ordering::SeqCst); let active_counter = self.active_mrf_workers.clone(); let stats = self.stats.clone(); let storage = self.storage.clone(); - let mrf_rx = self.mrf_replica_rx.lock().await.take(); + let mrf_rx = Arc::clone(&self.mrf_replica_rx); - if let Some(rx) = mrf_rx { - let handle = tokio::spawn(async move { - let mut rx = rx; - while let Some(operation) = rx.recv().await { - active_counter.fetch_add(1, Ordering::SeqCst); + let handle = tokio::spawn(async move { + loop { + let operation = { mrf_rx.lock().await.recv().await }; + let Some(operation) = operation else { break }; - match operation { - ReplicationOperation::Object(obj_info) => { - stats - .inc_q(&obj_info.bucket, obj_info.size, obj_info.delete_marker, obj_info.op_type) - .await; - - replicate_object(obj_info.as_ref().clone(), storage.clone()).await; - - stats - .dec_q(&obj_info.bucket, obj_info.size, obj_info.delete_marker, obj_info.op_type) - .await; - } - ReplicationOperation::Delete(del_info) => { - replicate_delete(*del_info, storage.clone()).await; - } + active_counter.fetch_add(1, Ordering::SeqCst); + match operation { + ReplicationOperation::Object(obj_info) => { + stats + .inc_q(&obj_info.bucket, obj_info.size, obj_info.delete_marker, obj_info.op_type) + .await; + replicate_object(obj_info.as_ref().clone(), storage.clone()).await; + stats + .dec_q(&obj_info.bucket, obj_info.size, obj_info.delete_marker, obj_info.op_type) + .await; + } + ReplicationOperation::Delete(del_info) => { + replicate_delete(*del_info, storage.clone()).await; } - - active_counter.fetch_sub(1, Ordering::SeqCst); } - }); - self.task_handles.lock().await.push(handle); - break; // Only one receiver can be taken - } + active_counter.fetch_sub(1, Ordering::SeqCst); + } + }); + self.task_handles.lock().await.push(handle); } // Remove workers if needed while self.mrf_worker_size.load(Ordering::SeqCst) > n { self.mrf_worker_size.fetch_sub(1, Ordering::SeqCst); - let _ = self.mrf_worker_kill_tx.try_send(()); // Signal worker to stop + let _ = self.mrf_worker_kill_tx.try_send(()); } } @@ -789,16 +790,164 @@ impl ReplicationPool { } } - /// Starts the MRF processor background task + /// Starts the MRF processor — one-shot at startup. + /// + /// Reads the on-disk MRF file, re-injects every entry into `mrf_replica_tx` as a + /// Heal operation, then clears the file. The file is cleared AFTER all entries are + /// successfully queued so a crash mid-replay results in at-most-twice delivery + /// (safe — replication is idempotent) rather than entry loss. async fn start_mrf_processor(&self) { - // This would start a background task to process MRF entries - // Implementation depends on the actual MRF processing logic + let storage = self.storage.clone(); + + let handle = tokio::spawn(async move { + let data = match read_config(storage.clone(), MRF_REPLICATION_FILE).await { + Ok(d) => d, + Err(EcstoreError::ConfigNotFound) => return, // no file yet — normal on first start + Err(e) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %e, + "Failed to load MRF recovery file" + ); + return; + } + }; + + let entries = match decode_mrf_file(&data) { + Ok(v) => v, + Err(e) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %e, + "Failed to decode MRF recovery file — discarding corrupt data" + ); + // Overwrite the corrupt file so we don't fail again on next restart. + let _ = save_config(storage, MRF_REPLICATION_FILE, encode_mrf_file(&[]).unwrap_or_default()).await; + return; + } + }; + + let total = entries.len(); + let mut queued_count = 0usize; + + for entry in entries.iter() { + match entry.op { + MrfOpKind::Delete => { + // Reconstruct a heal delete and re-queue it. We do NOT call + // get_object_info here because the delete-marker or version may + // already be absent from the local store — that is expected. + let dv = DeletedObjectReplicationInfo { + delete_object: crate::store_api::DeletedObject { + object_name: entry.object.clone(), + version_id: entry.version_id, + delete_marker_version_id: entry.delete_marker_version_id, + delete_marker: entry.delete_marker, + ..Default::default() + }, + bucket: entry.bucket.clone(), + op_type: ReplicationType::Heal, + event_type: REPLICATE_HEAL_DELETE.to_string(), + ..Default::default() + }; + schedule_replication_delete(dv).await; + queued_count += 1; + } + MrfOpKind::Object => { + let opts = ObjectOptions { + version_id: entry.version_id.map(|u| u.to_string()), + ..Default::default() + }; + let oi = match storage.get_object_info(&entry.bucket, &entry.object, &opts).await { + Ok(oi) => oi, + Err(e) => { + debug!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + bucket = %entry.bucket, + object = %entry.object, + error = %e, + "MRF recovery: object not found, skipping" + ); + continue; + } + }; + // Route through queue_replication_heal so the replication decision (dsc) + // is computed from the live config — required for replicate_object. + queue_replication_heal(&entry.bucket, oi, entry.retry_count as u32).await; + queued_count += 1; + } + } + } + + // Clear AFTER all entries are processed so a crash mid-replay causes at-most-twice + // delivery (idempotent) rather than entry loss. + if let Err(e) = save_config(storage, MRF_REPLICATION_FILE, encode_mrf_file(&[]).unwrap_or_default()).await { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %e, + "Failed to clear MRF recovery file after replay — entries may be replayed again on next restart" + ); + } + + if queued_count > 0 { + info!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + recovered = queued_count, + total, + "Recovered MRF entries from disk and queued for retry" + ); + } + }); + self.task_handles.lock().await.push(handle); } - /// Starts the MRF persister background task + /// Starts the MRF persister — ongoing background task. + /// + /// Drains `mrf_save_rx` (entries that overflowed the normal worker channels) and + /// writes them to the on-disk MRF file every 10 seconds or when 1 000 entries + /// accumulate. The file is overwritten (not appended) on each flush so it always + /// reflects the current pending backlog. async fn start_mrf_persister(&self) { - // This would start a background task to persist MRF entries to disk - // Implementation depends on the actual persistence logic + let Some(mut rx) = self.mrf_save_rx.lock().await.take() else { + return; + }; + let storage = self.storage.clone(); + + let handle = tokio::spawn(async move { + let mut pending: Vec = Vec::new(); + let mut interval = tokio::time::interval(Duration::from_secs(10)); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + + loop { + tokio::select! { + entry = rx.recv() => match entry { + Some(e) => { + pending.push(e); + if pending.len() >= 1000 && flush_mrf_to_disk(&pending, &storage).await { + pending.clear(); + } + } + None => { + // Channel closed (pool shutting down) — final flush. + if !pending.is_empty() { + flush_mrf_to_disk(&pending, &storage).await; + } + break; + } + }, + _ = interval.tick() => { + if !pending.is_empty() && flush_mrf_to_disk(&pending, &storage).await { + pending.clear(); + } + } + } + } + }); + self.task_handles.lock().await.push(handle); } /// Worker function for handling regular replication operations @@ -1099,6 +1248,38 @@ impl ReplicationPool { } } +/// Encodes `entries` and overwrites the MRF persistence file. +/// Returns `true` on success; on failure logs the error and returns `false`. +/// Callers must NOT clear their in-memory buffer on `false` so the next tick +/// can retry — otherwise a transient storage error permanently drops the batch. +async fn flush_mrf_to_disk(entries: &[MrfReplicateEntry], storage: &Arc) -> bool { + match encode_mrf_file(entries) { + Ok(data) => { + if let Err(e) = save_config(storage.clone(), MRF_REPLICATION_FILE, data).await { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + count = entries.len(), + error = %e, + "Failed to flush MRF entries to disk" + ); + return false; + } + true + } + Err(e) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + count = entries.len(), + error = %e, + "Failed to encode MRF entries for disk flush" + ); + false + } + } +} + /// Load bucket resync metadata from disk async fn load_bucket_resync_metadata( bucket: &str, @@ -1511,6 +1692,8 @@ async fn queue_replicate_deletes_wrapper( #[cfg(test)] mod tests { use super::*; + use crate::bucket::replication::replication_resyncer::{decode_mrf_file, encode_mrf_file}; + use uuid::Uuid; #[test] fn replication_queue_admission_combines_target_results() { @@ -1532,4 +1715,214 @@ mod tests { assert!(!should_auto_resume_resync(ResyncStatusType::ResyncCompleted)); assert!(!should_auto_resume_resync(ResyncStatusType::ResyncFailed)); } + + // ── MrfReplicateEntry encode/decode roundtrips ──────────────────────────── + + #[test] + fn mrf_entry_object_roundtrip() { + let vid = Uuid::new_v4(); + let entry = MrfReplicateEntry { + bucket: "my-bucket".to_string(), + object: "path/to/obj".to_string(), + version_id: Some(vid), + retry_count: 3, + size: 1024, + op: MrfOpKind::Object, + delete_marker_version_id: None, + delete_marker: false, + }; + + let encoded = encode_mrf_file(std::slice::from_ref(&entry)).expect("encode"); + let decoded = decode_mrf_file(&encoded).expect("decode"); + + assert_eq!(decoded.len(), 1); + let got = &decoded[0]; + assert_eq!(got.bucket, "my-bucket"); + assert_eq!(got.object, "path/to/obj"); + assert_eq!(got.version_id, Some(vid)); + assert_eq!(got.retry_count, 3); + assert_eq!(got.size, 1024); + assert_eq!(got.op, MrfOpKind::Object); + assert_eq!(got.delete_marker_version_id, None); + assert!(!got.delete_marker); + } + + #[test] + fn mrf_entry_delete_marker_roundtrip() { + let dm_vid = Uuid::new_v4(); + let entry = MrfReplicateEntry { + bucket: "del-bucket".to_string(), + object: "key".to_string(), + version_id: None, + retry_count: 0, + size: 0, + op: MrfOpKind::Delete, + delete_marker_version_id: Some(dm_vid), + delete_marker: true, + }; + + let encoded = encode_mrf_file(std::slice::from_ref(&entry)).expect("encode"); + let decoded = decode_mrf_file(&encoded).expect("decode"); + + assert_eq!(decoded.len(), 1); + let got = &decoded[0]; + assert_eq!(got.bucket, "del-bucket"); + assert_eq!(got.object, "key"); + assert_eq!(got.version_id, None); + assert_eq!(got.op, MrfOpKind::Delete); + assert_eq!(got.delete_marker_version_id, Some(dm_vid)); + assert!(got.delete_marker); + } + + #[test] + fn mrf_entry_versioned_delete_roundtrip() { + let vid = Uuid::new_v4(); + let entry = MrfReplicateEntry { + bucket: "ver-bucket".to_string(), + object: "versioned-key".to_string(), + version_id: Some(vid), + retry_count: 0, + size: 0, + op: MrfOpKind::Delete, + delete_marker_version_id: None, + delete_marker: false, + }; + + let encoded = encode_mrf_file(&[entry]).expect("encode"); + let decoded = decode_mrf_file(&encoded).expect("decode"); + + assert_eq!(decoded.len(), 1); + let got = &decoded[0]; + assert_eq!(got.op, MrfOpKind::Delete); + assert_eq!(got.version_id, Some(vid)); + assert_eq!(got.delete_marker_version_id, None); + assert!(!got.delete_marker); + } + + #[test] + fn mrf_entry_mixed_batch_roundtrip() { + let obj_vid = Uuid::new_v4(); + let del_dm_vid = Uuid::new_v4(); + let entries = vec![ + MrfReplicateEntry { + bucket: "b".to_string(), + object: "obj".to_string(), + version_id: Some(obj_vid), + retry_count: 1, + size: 512, + op: MrfOpKind::Object, + delete_marker_version_id: None, + delete_marker: false, + }, + MrfReplicateEntry { + bucket: "b".to_string(), + object: "del".to_string(), + version_id: None, + retry_count: 0, + size: 0, + op: MrfOpKind::Delete, + delete_marker_version_id: Some(del_dm_vid), + delete_marker: true, + }, + ]; + + let encoded = encode_mrf_file(&entries).expect("encode"); + let decoded = decode_mrf_file(&encoded).expect("decode"); + + assert_eq!(decoded.len(), 2); + assert_eq!(decoded[0].op, MrfOpKind::Object); + assert_eq!(decoded[0].version_id, Some(obj_vid)); + assert_eq!(decoded[1].op, MrfOpKind::Delete); + assert_eq!(decoded[1].delete_marker_version_id, Some(del_dm_vid)); + assert!(decoded[1].delete_marker); + } + + // ── Recovery replay routing ─────────────────────────────────────────────── + + #[test] + fn mrf_entry_op_routes_correctly() { + // Object entries must have op=Object so the processor calls get_object_info + heal. + let obj_entry = MrfReplicateEntry { + bucket: "b".to_string(), + object: "o".to_string(), + version_id: None, + retry_count: 0, + size: 0, + op: MrfOpKind::Object, + delete_marker_version_id: None, + delete_marker: false, + }; + assert_eq!(obj_entry.op, MrfOpKind::Object); + + // Delete entries must have op=Delete so the processor calls schedule_replication_delete. + let del_entry = MrfReplicateEntry { + bucket: "b".to_string(), + object: "o".to_string(), + version_id: None, + retry_count: 0, + size: 0, + op: MrfOpKind::Delete, + delete_marker_version_id: Some(Uuid::new_v4()), + delete_marker: true, + }; + assert_eq!(del_entry.op, MrfOpKind::Delete); + + // Entries written by old code (before the op field existed) must deserialise as Object + // so existing recovery behaviour is preserved. + let legacy_entry = MrfReplicateEntry { + bucket: "b".to_string(), + object: "o".to_string(), + version_id: None, + retry_count: 0, + size: 0, + op: MrfOpKind::default(), + delete_marker_version_id: None, + delete_marker: false, + }; + assert_eq!(legacy_entry.op, MrfOpKind::Object, "legacy default must be Object"); + } + + #[test] + fn mrf_legacy_file_without_op_field_decoded_as_object() { + // Hand-build the exact bytes a pre-MrfOpKind binary would have written to disk. + // The old MrfReplicateEntry had only 4 persisted keys (versionID is omitted when + // None due to skip_serializing_if): bucket, object, retryCount, size. + // There is no "op", "deleteMarker", or "deleteMarkerVersionID" key. + // + // This proves that #[serde(default)] on the `op` field carries real weight: + // if you remove that attribute, rmp_serde will return an error on this payload + // and the test will fail. + let mut msgpack = Vec::new(); + // Outer: array of 1 (the Vec) + rmp::encode::write_array_len(&mut msgpack, 1).unwrap(); + // Inner: named map with the 4 original fields only — no "op", no "deleteMarker*" + rmp::encode::write_map_len(&mut msgpack, 4).unwrap(); + rmp::encode::write_str(&mut msgpack, "bucket").unwrap(); + rmp::encode::write_str(&mut msgpack, "old-bucket").unwrap(); + rmp::encode::write_str(&mut msgpack, "object").unwrap(); + rmp::encode::write_str(&mut msgpack, "old-key").unwrap(); + rmp::encode::write_str(&mut msgpack, "retryCount").unwrap(); + rmp::encode::write_i32(&mut msgpack, 2).unwrap(); + rmp::encode::write_str(&mut msgpack, "size").unwrap(); + rmp::encode::write_i64(&mut msgpack, 100).unwrap(); + + // Prepend the MRF file header: format=1 (LE u16) || version=1 (LE u16) + let mut data = Vec::with_capacity(4 + msgpack.len()); + data.extend_from_slice(&1u16.to_le_bytes()); // MRF_META_FORMAT + data.extend_from_slice(&1u16.to_le_bytes()); // MRF_META_VERSION + data.extend_from_slice(&msgpack); + + let decoded = decode_mrf_file(&data).expect("legacy payload must decode without error"); + assert_eq!(decoded.len(), 1); + let entry = &decoded[0]; + assert_eq!(entry.bucket, "old-bucket"); + assert_eq!(entry.object, "old-key"); + assert_eq!(entry.retry_count, 2); + assert_eq!(entry.size, 100); + assert_eq!(entry.version_id, None); + // The "op" key was absent — #[serde(default)] must fill in MrfOpKind::Object. + assert_eq!(entry.op, MrfOpKind::Object, "missing op key must default to Object"); + assert!(!entry.delete_marker); + assert_eq!(entry.delete_marker_version_id, None); + } } diff --git a/crates/ecstore/src/bucket/replication/replication_resyncer.rs b/crates/ecstore/src/bucket/replication/replication_resyncer.rs index 610b51f28..692dc8127 100644 --- a/crates/ecstore/src/bucket/replication/replication_resyncer.rs +++ b/crates/ecstore/src/bucket/replication/replication_resyncer.rs @@ -26,7 +26,7 @@ use crate::bucket::target::BucketTargets; use crate::bucket::versioning_sys::BucketVersioningSys; use crate::client::api_get_options::{AdvancedGetOptions, StatObjectOptions}; use crate::config::com::save_config; -use crate::disk::BUCKET_META_PREFIX; +use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET}; use crate::error::{Error, Result, is_err_object_not_found, is_err_version_not_found}; use crate::event_notification::{EventArgs, send_event}; use crate::global::GLOBAL_LocalNodeName; @@ -52,6 +52,7 @@ use http::HeaderMap; use http_body::Frame; use http_body_util::StreamBody; use regex::Regex; +use rmp_serde; use rustfs_filemeta::{ MrfReplicateEntry, REPLICATE_EXISTING, REPLICATE_EXISTING_DELETE, ReplicateDecision, ReplicateObjectInfo, ReplicateTargetDecision, ReplicatedInfos, ReplicatedTargetInfo, ReplicationAction, ReplicationState, ReplicationStatusType, @@ -106,6 +107,12 @@ pub(crate) const REPLICATION_DIR: &str = ".replication"; pub(crate) const RESYNC_FILE_NAME: &str = "resync.bin"; pub(crate) const RESYNC_META_FORMAT: u16 = 1; pub(crate) const RESYNC_META_VERSION: u16 = 1; + +// MRF (Most Recent Failures) persistence file — stored at +// `{RUSTFS_META_BUCKET}/config/replication/mrf.bin`, cross-bucket. +pub(crate) const MRF_REPLICATION_FILE: &str = "config/replication/mrf.bin"; +const MRF_META_FORMAT: u16 = 1; +const MRF_META_VERSION: u16 = 1; const RESYNC_TIME_INTERVAL: TokioDuration = TokioDuration::from_secs(60); const WIRE_ZERO_TIME_UNIX: i64 = -62_135_596_800; @@ -355,6 +362,36 @@ pub(crate) fn decode_resync_file(data: &[u8]) -> Result Result> { + let payload = rmp_serde::to_vec_named(entries).map_err(|e| Error::other(e.to_string()))?; + let mut data = Vec::with_capacity(4 + payload.len()); + let mut fmt = [0u8; 2]; + byteorder::LittleEndian::write_u16(&mut fmt, MRF_META_FORMAT); + data.extend_from_slice(&fmt); + let mut ver = [0u8; 2]; + byteorder::LittleEndian::write_u16(&mut ver, MRF_META_VERSION); + data.extend_from_slice(&ver); + data.extend_from_slice(&payload); + Ok(data) +} + +pub(crate) fn decode_mrf_file(data: &[u8]) -> Result> { + if data.len() <= 4 { + return Err(Error::CorruptedFormat); + } + let mut fmt = [0u8; 2]; + fmt.copy_from_slice(&data[0..2]); + if byteorder::LittleEndian::read_u16(&fmt) != MRF_META_FORMAT { + return Err(Error::CorruptedFormat); + } + let mut ver = [0u8; 2]; + ver.copy_from_slice(&data[2..4]); + if byteorder::LittleEndian::read_u16(&ver) != MRF_META_VERSION { + return Err(Error::CorruptedFormat); + } + rmp_serde::from_slice(&data[4..]).map_err(|e| Error::other(e.to_string())) +} + impl TargetReplicationResyncStatus { fn marshal_wire_msg(&self, wr: &mut Vec) -> Result<()> { rmp::encode::write_map_len(wr, 11)?; @@ -691,6 +728,42 @@ impl ReplicationResyncer { return; } + // Acquire a cluster-wide leader lock for this (bucket, ARN) pair so that only + // one node runs the resync scan at a time. Without this, every cluster node would + // scan and replicate every object independently, causing N-fold duplicate traffic. + let resync_lock_key = format!("{}/{}/{}", REPLICATION_DIR, opts.bucket, opts.arn); + let resync_ns_lock = match storage.new_ns_lock(RUSTFS_META_BUCKET, &resync_lock_key).await { + Ok(l) => l, + Err(e) => { + warn!( + event = EVENT_RESYNC_STATUS_UPDATE_SKIPPED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + bucket = %opts.bucket, + arn = %opts.arn, + error = %e, + reason = "leader_lock_create_failed", + "Failed to create resync leader lock — skipping resync" + ); + return; + } + }; + let _resync_leader_guard = match resync_ns_lock.get_write_lock(get_lock_acquire_timeout()).await { + Ok(g) => g, + Err(_) => { + debug!( + event = EVENT_RESYNC_STATUS_UPDATE_SKIPPED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + bucket = %opts.bucket, + arn = %opts.arn, + reason = "leader_lock_held_by_another_node", + "Another node is already running resync for this bucket/ARN — skipping" + ); + return; + } + }; + let cfg = match get_replication_config(&opts.bucket).await { Ok(cfg) => cfg, Err(err) => { @@ -1187,9 +1260,14 @@ impl ReplicationWorkerOperation for DeletedObjectReplicationInfo { MrfReplicateEntry { bucket: self.bucket.clone(), object: self.delete_object.object_name.clone(), - version_id: None, + // version_id here is the version being purged (if any); the delete-marker + // version is stored separately in delete_marker_version_id. + version_id: self.delete_object.version_id, retry_count: 0, size: 0, + op: rustfs_filemeta::MrfOpKind::Delete, + delete_marker_version_id: self.delete_object.delete_marker_version_id, + delete_marker: self.delete_object.delete_marker, } } @@ -2556,7 +2634,7 @@ async fn replicate_delete_to_target(dobj: &DeletedObjectReplicationInfo, tgt_cli rinfo } -pub async fn replicate_object(roi: ReplicateObjectInfo, storage: Arc) { +pub async fn replicate_object(roi: ReplicateObjectInfo, storage: Arc) { let bucket = roi.bucket.clone(); let object = roi.name.clone(); @@ -2614,7 +2692,57 @@ pub async fn replicate_object(roi: ReplicateObjectInfo, storage: ..Default::default() }); - // TODO: NSLOCK + // Acquire a per-object namespace lock so that at most one worker (across all cluster + // nodes and MRF retry goroutines) replicates this object version at a time. + let obj_lock_key = format!("/[replicate]/{}", object); + let obj_ns_lock = match storage.new_ns_lock(&bucket, &obj_lock_key).await { + Ok(l) => l, + Err(e) => { + debug!( + event = EVENT_RESYNC_RUNTIME_SKIPPED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + bucket = %bucket, + object = %object, + error = %e, + reason = "ns_lock_create_failed", + "Skipping replication object" + ); + send_event(EventArgs { + event_name: EventName::ObjectReplicationNotTracked.to_string(), + bucket_name: bucket.clone(), + object: roi.to_object_info(), + host: GLOBAL_LocalNodeName.to_string(), + user_agent: "Internal: [Replication]".to_string(), + ..Default::default() + }); + return; + } + }; + let _obj_lock_guard = match obj_ns_lock.get_write_lock(get_lock_acquire_timeout()).await { + Ok(g) => g, + Err(e) => { + debug!( + event = EVENT_RESYNC_RUNTIME_SKIPPED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC, + bucket = %bucket, + object = %object, + error = %e, + reason = "ns_lock_write_lock_failed", + "Skipping replication object" + ); + send_event(EventArgs { + event_name: EventName::ObjectReplicationNotTracked.to_string(), + bucket_name: bucket.clone(), + object: roi.to_object_info(), + host: GLOBAL_LocalNodeName.to_string(), + user_agent: "Internal: [Replication]".to_string(), + ..Default::default() + }); + return; + } + }; let mut join_set = JoinSet::new(); diff --git a/crates/filemeta/src/replication.rs b/crates/filemeta/src/replication.rs index 236e820ad..32dd8a644 100644 --- a/crates/filemeta/src/replication.rs +++ b/crates/filemeta/src/replication.rs @@ -534,7 +534,20 @@ impl ReplicatedInfos { } } -#[derive(Serialize, Deserialize, Debug)] +/// Distinguishes the kind of operation stored in [`MrfReplicateEntry`]. +/// +/// Old serialized files lack the `op` key; `default` maps to `Object`, which preserves +/// the pre-existing replay behaviour for entries written before this field existed. +#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, Default)] +pub enum MrfOpKind { + #[default] + #[serde(rename = "object")] + Object, + #[serde(rename = "delete")] + Delete, +} + +#[derive(Serialize, Deserialize, Debug, Clone)] pub struct MrfReplicateEntry { #[serde(rename = "bucket")] pub bucket: String, @@ -542,14 +555,30 @@ pub struct MrfReplicateEntry { #[serde(rename = "object")] pub object: String, - #[serde(skip_serializing, skip_deserializing)] + // Persisted so recovery after restart can replay the exact version. + // Old serialized files lack this key; `default` fills in None safely. + #[serde(rename = "versionID", skip_serializing_if = "Option::is_none", default)] pub version_id: Option, #[serde(rename = "retryCount")] pub retry_count: i32, - #[serde(skip_serializing, skip_deserializing)] + #[serde(rename = "size", default)] pub size: i64, + + // Operation kind. Old files lack this key; default=Object preserves existing behaviour. + #[serde(rename = "op", default)] + pub op: MrfOpKind, + + // For delete entries: the delete-marker version id (distinct from version_id, which is + // the version being purged). Old files lack this; default=None is correct. + #[serde(rename = "deleteMarkerVersionID", skip_serializing_if = "Option::is_none", default)] + pub delete_marker_version_id: Option, + + // For delete entries: whether this is a delete-marker vs a versioned-object delete. + // Old files lack this; default=false is correct. + #[serde(rename = "deleteMarker", default)] + pub delete_marker: bool, } pub trait ReplicationWorkerOperation: Any + Send + Sync { @@ -748,6 +777,9 @@ impl ReplicationWorkerOperation for ReplicateObjectInfo { version_id: self.version_id, retry_count: self.retry_count as i32, size: self.size, + op: MrfOpKind::Object, + delete_marker_version_id: None, + delete_marker: false, } } @@ -795,6 +827,9 @@ impl ReplicateObjectInfo { version_id: self.version_id, retry_count: self.retry_count as i32, size: self.size, + op: MrfOpKind::Object, + delete_marker_version_id: None, + delete_marker: false, } } }