From eb87bb1fafcaa08f82a053b65964b6fd8a989ebc Mon Sep 17 00:00:00 2001 From: cxymds Date: Tue, 4 Aug 2026 21:40:50 +0800 Subject: [PATCH] fix(replication): harden resync and MRF recovery (#5694) * fix(replication): harden resync and MRF recovery * fix(replication): correct MRF validation regressions * fix(replication): address CI validation failures * fix(heal): initialize decode error in merge test --- .../replication/replication_config_store.rs | 7 + .../replication_metadata_boundary.rs | 5 + .../bucket/replication/replication_pool.rs | 1146 ++++++++--------- .../replication_resync_boundary.rs | 1 + crates/ecstore/src/config/com.rs | 37 +- crates/ecstore/src/set_disk/ops/heal_walk.rs | 1 + crates/replication/src/filemeta.rs | 2 +- crates/replication/src/lib.rs | 6 +- crates/replication/src/resync.rs | 484 +++++-- rustfs/src/admin/handlers/replication.rs | 4 +- rustfs/src/admin/route_registration_test.rs | 26 + 11 files changed, 1004 insertions(+), 715 deletions(-) diff --git a/crates/ecstore/src/bucket/replication/replication_config_store.rs b/crates/ecstore/src/bucket/replication/replication_config_store.rs index 25402c394..de0e14c61 100644 --- a/crates/ecstore/src/bucket/replication/replication_config_store.rs +++ b/crates/ecstore/src/bucket/replication/replication_config_store.rs @@ -30,6 +30,13 @@ impl ReplicationConfigStore { com::read_config(api, file).await } + pub(crate) async fn read_limited(api: Arc, file: &str, max_bytes: usize) -> Result> + where + S: ReplicationObjectIO, + { + com::read_config_limited(api, file, max_bytes).await + } + pub(crate) async fn read_no_lock(api: Arc, file: &str) -> Result> where S: ReplicationObjectIO, diff --git a/crates/ecstore/src/bucket/replication/replication_metadata_boundary.rs b/crates/ecstore/src/bucket/replication/replication_metadata_boundary.rs index c7d7ba37c..de74f9b0a 100644 --- a/crates/ecstore/src/bucket/replication/replication_metadata_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_metadata_boundary.rs @@ -32,6 +32,7 @@ pub(crate) struct ReplicationMetadataStore; impl ReplicationMetadataStore { pub(crate) const MRF_REPLICATION_FILE: &'static str = "config/replication/mrf.bin"; + pub(crate) const MRF_REPLICATION_RECOVERY_LOCK: &'static str = "config/replication/mrf.bin.recovery"; pub(crate) const FORCE_DELETE_REPLICATION_FILE: &'static str = "config/replication/force-delete.bin"; pub(crate) const FORCE_DELETE_REPLICATION_TRANSACTION_LOCK: &'static str = "config/replication/force-delete.bin.transaction"; @@ -111,6 +112,10 @@ mod tests { "buckets/bucket-a/.replication/resync.bin" ); assert_eq!(ReplicationMetadataStore::MRF_REPLICATION_FILE, "config/replication/mrf.bin"); + assert_eq!( + ReplicationMetadataStore::MRF_REPLICATION_RECOVERY_LOCK, + "config/replication/mrf.bin.recovery" + ); assert_eq!( ReplicationMetadataStore::FORCE_DELETE_REPLICATION_FILE, "config/replication/force-delete.bin" diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index dba231b95..9edc4fe2d 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -34,8 +34,8 @@ use super::replication_queue_boundary::{ }; use super::replication_resync_boundary::ResyncStatusType; use super::replication_resync_boundary::{ - BucketReplicationResyncStatus, ResyncOpts, TargetReplicationResyncStatus, decode_mrf_file, decode_resync_file, - encode_mrf_file, should_auto_resume_resync, + BucketReplicationResyncStatus, RESYNC_FILE_MAX_BYTES, ResyncOpts, TargetReplicationResyncStatus, decode_mrf_file, + decode_resync_file, encode_mrf_file, should_auto_resume_resync, }; use super::replication_resyncer::{ ReplicationResyncer, get_heal_replicate_object_info, replicate_delete, replicate_delete_with_outcome, replicate_object, @@ -64,7 +64,6 @@ use std::time::Instant; use time::OffsetDateTime; use time::format_description::well_known::Rfc3339; use tokio::sync::Mutex; -use tokio::sync::Notify; use tokio::sync::RwLock; use tokio::sync::mpsc; use tokio::sync::mpsc::Receiver; @@ -229,20 +228,18 @@ fn durable_mrf_backlog_tracker_from_entries(entries: &[MrfReplicateEntry]) -> Du tracker } -fn move_staged_mrf_entries(pending: &mut Vec, staged: &mut Vec) -> Vec { - let pending_capacity = MRF_PENDING_CAP.saturating_sub(pending.len()); - let mut staged_entries = std::mem::take(staged); - let capped_batch = staged_entries.split_off(pending_capacity.min(staged_entries.len())); - pending.append(&mut staged_entries); - capped_batch -} - #[derive(Debug)] struct PendingMrfAppend { digest: [u8; 32], entry_count: usize, } +#[derive(Debug)] +struct MrfAppendResult { + duration_millis: u64, + backlog: DurableMrfBacklogSnapshot, +} + fn mrf_payload_digest(data: &[u8]) -> [u8; 32] { let encoded = HashAlgorithm::SHA256.hash_encode(data); let mut digest = [0; 32]; @@ -250,13 +247,6 @@ fn mrf_payload_digest(data: &[u8]) -> [u8; 32] { digest } -fn add_durable_mrf_suffix(tracker: &mut DurableMrfBacklogTracker, entries: &[MrfReplicateEntry], count: usize) { - let start = entries.len().saturating_sub(count); - for entry in &entries[start..] { - tracker.add_entry(entry); - } -} - #[derive(Debug, Clone, Default)] struct MrfBacklogObservabilityTracker { buckets: HashMap, @@ -607,64 +597,72 @@ fn mrf_journal_preconditions(etag: Option<&str>, exists: bool) -> Option( +fn mrf_prefix_matches(current: &[MrfReplicateEntry], prefix: &[MrfReplicateEntry]) -> bool { + current.starts_with(prefix) +} + +async fn read_mrf_entries(storage: Arc) -> Result, EcstoreError> { + match ReplicationConfigStore::read(storage, ReplicationMetadataStore::MRF_REPLICATION_FILE).await { + Ok(data) if data.is_empty() => Ok(Vec::new()), + Ok(data) => decode_mrf_file(&data), + Err(EcstoreError::ConfigNotFound) => Ok(Vec::new()), + Err(error) => Err(error), + } +} + +/// Acknowledge only the generation read by the recovery leader. Entries appended +/// after that generation are retained as a suffix and replayed on the next startup. +/// Lock order: recovery leader lock -> MRF journal object lock. +async fn acknowledge_mrf_recovery( storage: Arc, - desired: &[MrfReplicateEntry], -) -> Result<(), EcstoreError> { + recovery_guard: &rustfs_lock::NamespaceLockGuard, + replayed_prefix: &[MrfReplicateEntry], + retry_entries: &[MrfReplicateEntry], +) -> Result, EcstoreError> { let file = ReplicationMetadataStore::MRF_REPLICATION_FILE; - let mut merged = desired.to_vec(); - let mut saw_conflict = false; for _attempt in 0..=FORCE_DELETE_INTENT_CAS_RETRIES { let lock = storage .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_preserve_empty(storage.clone(), file).await; - let (etag, exists) = match current { - Ok((data, object_info)) => { - if saw_conflict { - let current = decode_mrf_file(&data)?; - for entry in current { - if !merged.iter().any(|existing| mrf_entries_same(existing, &entry)) { - merged.push(entry); - } - } - } - (object_info.etag, true) - } - Err(EcstoreError::ConfigNotFound) => (None, false), - Err(err) => return Err(err), + let (current_data, current_etag, current_exists) = match current { + Ok((data, object_info)) => (data, object_info.etag, true), + Err(EcstoreError::ConfigNotFound) => (Vec::new(), None, false), + Err(error) => return Err(error), }; - if guard.is_lock_lost() { - return Err(EcstoreError::other("MRF journal namespace lock was lost before commit")); - } - 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() { + let current_entries = if current_data.is_empty() { Vec::new() } else { - encode_mrf_file(&merged)? + decode_mrf_file(¤t_data)? }; - match ReplicationConfigStore::save_conditional(storage.clone(), file, data, preconditions).await { - Ok(()) => return Ok(()), - Err(EcstoreError::PreconditionFailed) => saw_conflict = true, - Err(err) => return Err(err), + if !mrf_prefix_matches(¤t_entries, replayed_prefix) { + return Err(EcstoreError::other("MRF recovery prefix changed before acknowledgement")); + } + + let mut retained = Vec::with_capacity(retry_entries.len() + current_entries.len().saturating_sub(replayed_prefix.len())); + retained.extend_from_slice(retry_entries); + retained.extend_from_slice(¤t_entries[replayed_prefix.len()..]); + if recovery_guard.is_lock_lost() || guard.is_lock_lost() { + return Err(EcstoreError::other("MRF recovery lock lost before acknowledgement")); + } + let Some(preconditions) = mrf_journal_preconditions(current_etag.as_deref(), current_exists) else { + return Err(EcstoreError::other("MRF journal has no ETag for recovery acknowledgement")); + }; + let data = if retained.is_empty() { + Vec::new() + } else { + encode_mrf_file(&retained)? + }; + match ReplicationConfigStore::save_conditional_no_lock(storage.clone(), file, data, preconditions).await { + Ok(()) => return Ok(retained), + Err(EcstoreError::PreconditionFailed) => continue, + Err(error) => return Err(error), } } Err(EcstoreError::PreconditionFailed) } -fn mrf_entries_same(left: &MrfReplicateEntry, right: &MrfReplicateEntry) -> bool { - left.bucket == right.bucket - && left.object == right.object - && left.version_id == right.version_id - && left.op == right.op - && left.target_arns == right.target_arns - && left.force_delete_id == right.force_delete_id - && left.delete_marker_version_id == right.delete_marker_version_id -} - #[derive(Debug, thiserror::Error)] #[error("replication resync {active_resync_id} is already active for {bucket}/{arn}")] struct ResyncActiveConflictError { @@ -711,8 +709,6 @@ pub struct ReplicationPool { mrf_replica_rx: Arc>>, mrf_save_tx: Sender, mrf_save_rx: Mutex>>, - mrf_recovery_complete: Arc, - mrf_recovery_result: Arc>>>, // Control channels mrf_worker_kill_tx: Sender<()>, @@ -756,8 +752,6 @@ impl ReplicationPool { mrf_replica_rx: Arc::new(Mutex::new(mrf_replica_rx)), mrf_save_tx, mrf_save_rx: Mutex::new(Some(mrf_save_rx)), - mrf_recovery_complete: Arc::new(Notify::new()), - mrf_recovery_result: Arc::new(Mutex::new(None)), mrf_worker_kill_tx, mrf_stop_tx, mrf_worker_size: AtomicI32::new(0), @@ -1197,10 +1191,41 @@ impl ReplicationPool { /// rewrites any entries that could not be admitted for a later startup retry. async fn start_mrf_processor(&self) { let storage = self.storage.clone(); - let recovery_complete = self.mrf_recovery_complete.clone(); - let recovery_result = self.mrf_recovery_result.clone(); let handle = tokio::spawn(async move { + let recovery_lock = match storage + .new_ns_lock( + ReplicationMetadataStore::rustfs_meta_bucket(), + ReplicationMetadataStore::MRF_REPLICATION_RECOVERY_LOCK, + ) + .await + { + Ok(lock) => lock, + Err(error) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %error, + "Failed to create the MRF recovery leader lock" + ); + return; + } + }; + let recovery_guard = match recovery_lock + .get_write_lock_quiet(ReplicationLockTiming::acquire_timeout()) + .await + { + Ok(guard) => guard, + Err(_) => { + debug!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + "Another node is already processing the MRF recovery backlog" + ); + return; + } + }; + let data = match ReplicationConfigStore::read(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await { Ok(d) => d, Err(EcstoreError::ConfigNotFound) => { @@ -1208,8 +1233,6 @@ impl ReplicationPool { available: true, buckets: Vec::new(), }); - *recovery_result.lock().await = Some(Vec::new()); - recovery_complete.notify_one(); return; } Err(e) => { @@ -1219,7 +1242,6 @@ impl ReplicationPool { error = %e, "Failed to load MRF recovery file" ); - recovery_complete.notify_one(); return; } }; @@ -1234,7 +1256,6 @@ impl ReplicationPool { "Failed to decode MRF recovery file — preserving corrupt data" ); quarantine_mrf_file(&storage, &data).await; - recovery_complete.notify_one(); return; } }; @@ -1435,9 +1456,31 @@ impl ReplicationPool { } } - let retained_count = retry_entries.len(); - *recovery_result.lock().await = Some(retry_entries); - recovery_complete.notify_one(); + let retained = match acknowledge_mrf_recovery(storage.clone(), &recovery_guard, &entries, &retry_entries).await { + Ok(retained) => retained, + Err(error) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %error, + "Failed to acknowledge the MRF recovery prefix; preserving it for the next startup" + ); + match read_mrf_entries(storage.clone()).await { + Ok(current) => current, + Err(read_error) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %read_error, + "Failed to refresh the MRF backlog after acknowledgement failure" + ); + entries.clone() + } + } + } + }; + let retained_count = retained.len(); + set_durable_mrf_backlog_snapshot(durable_mrf_backlog_summary_from_entries(&retained)); if queued_count > 0 { info!( @@ -1519,96 +1562,20 @@ impl ReplicationPool { /// Drains `mrf_save_rx` (entries that overflowed the normal worker channels) and /// writes them to the on-disk MRF file every flush interval (default 10s, /// overridable via `RUSTFS_REPL_MRF_FLUSH_INTERVAL_MS`) or when 1 000 new - /// entries accumulate. Each flush rewrites the whole cumulative backlog so no - /// previously-persisted (and not-yet-replayed) entry is lost; the file is only - /// consumed and cleared at startup. + /// entries accumulate. Runtime writers append under the journal lock; only the + /// startup recovery leader may remove an acknowledged prefix. async fn start_mrf_persister(&self) { let Some(mut rx) = self.mrf_save_rx.lock().await.take() else { return; }; let storage = self.storage.clone(); let stats = self.stats.clone(); - let recovery_complete = self.mrf_recovery_complete.clone(); - let recovery_result = self.mrf_recovery_result.clone(); let handle = tokio::spawn(async move { - let mut staged = Vec::new(); - let mut staging_closed = false; - let retry_timer = tokio::time::sleep(Duration::ZERO); - tokio::pin!(retry_timer); - let mut pending = loop { - tokio::select! { - entry = rx.recv(), if staged.len() < MRF_PENDING_CAP && !staging_closed => match entry { - Some(entry) => { - observe_mrf_pending(&entry); - staged.push(entry); - }, - None => staging_closed = true, - }, - _ = &mut retry_timer => { - match ReplicationConfigStore::read(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await { - Ok(data) => match decode_mrf_file(&data) { - Ok(entries) => break entries, - Err(error) => { - warn!( - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION, - error = %error, - "Failed to seed MRF persister from the startup recovery file; retrying without overwriting it" - ); - } - }, - Err(EcstoreError::ConfigNotFound) => break Vec::new(), - Err(error) => { - warn!( - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION, - error = %error, - "Failed to read the startup MRF backlog for persister seeding; retrying without overwriting it" - ); - } - } - retry_timer.as_mut().reset(tokio::time::Instant::now() + MRF_RETRY_INITIAL_DELAY); - } - } - }; - let initial_pending_len = pending.len(); - let mut capped_batch = move_staged_mrf_entries(&mut pending, &mut staged); - let pending_staged_len = pending.len().saturating_sub(initial_pending_len); - // The on-disk MRF file is a restart-recovery backstop: entries are - // only replayed (and the file cleared) at startup, never during the - // run. So the file must hold the *cumulative* set of overflow entries - // written this run. `pending` is therefore kept cumulative and the - // whole set is rewritten on each flush — clearing it after a flush - // let the next flush overwrite the file and drop everything written - // earlier (backlog#859 / #799 B10). Bounded by `MRF_PENDING_CAP` so a - // sustained failure storm can't grow it without limit. - let mut durable_tracker = DurableMrfBacklogTracker { - available: true, - ..Default::default() - }; - for entry in &pending[..initial_pending_len] { - durable_tracker.add_entry(entry); - } - if initial_pending_len > 0 { - set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - } - let mut new_entries_pending_stats = pending_staged_len; - let mut new_entries_pending_observability = pending_staged_len; - let mut dirty = pending_staged_len > 0; - let mut capped = initial_pending_len >= MRF_PENDING_CAP; - let mut recovery_applied = false; + let mut pending = Vec::new(); + let mut pending_payload = None; let mut channel_closed = false; - let mut capped_payload = None; - if capped { - warn!( - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION, - cap = MRF_PENDING_CAP, - pending = initial_pending_len, - "MRF pending backlog is at capacity — applying backpressure" - ); - } + let mut capped = false; // Flush interval: `RUSTFS_REPL_MRF_FLUSH_INTERVAL_MS` (default 10000ms, // clamped to >=10ms), read once when the persister task starts. let mut interval = tokio::time::interval(super::replication_timing::mrf_flush_interval()); @@ -1618,134 +1585,52 @@ impl ReplicationPool { if !channel_closed && rx.is_closed() && rx.is_empty() { channel_closed = true; } - if recovery_applied && channel_closed && !dirty && capped_batch.is_empty() { + if channel_closed && pending.is_empty() { break; } - if recovery_applied && (pending.len() >= MRF_PENDING_CAP || !capped_batch.is_empty()) { - if dirty { - if let Some(duration_millis) = flush_mrf_to_disk(&pending, &storage).await { - add_durable_mrf_suffix(&mut durable_tracker, &pending, new_entries_pending_stats); - set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - let observe_start = pending.len().saturating_sub(new_entries_pending_observability); - observe_mrf_pending_flushed(&pending[observe_start..], duration_millis); - new_entries_pending_observability = 0; - if new_entries_pending_stats > 0 { - let stats_start = pending.len().saturating_sub(new_entries_pending_stats); - dec_mrf_entries(stats.as_ref(), &pending[stats_start..]); - new_entries_pending_stats = 0; + let flush_requested = if channel_closed || pending.len() >= MRF_PENDING_CAP || pending_payload.is_some() { + true + } else { + tokio::select! { + entry = rx.recv() => match entry { + Some(entry) => { + observe_mrf_pending(&entry); + pending.push(entry); + pending.len() >= 1000 } - dirty = false; - } else { - // Keep the channel bounded while the current backlog - // cannot be persisted; draining it here would turn a - // failed flush into unbounded in-memory growth. - interval.tick().await; - continue; - } + None => { + channel_closed = true; + true + } + }, + _ = interval.tick() => !pending.is_empty(), } - if !capped { - capped = true; - warn!( - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION, - cap = MRF_PENDING_CAP, - "MRF pending backlog reached capacity — applying backpressure" - ); - } - if capped_batch.is_empty() { - while let Ok(entry) = rx.try_recv() { - capped_batch.push(entry); - } - for entry in &capped_batch { - observe_mrf_pending(entry); - } - } - if !capped_batch.is_empty() - && let Some(duration_millis) = - append_mrf_entries_to_disk(&capped_batch, &storage, &mut capped_payload, &pending).await - { - for entry in &capped_batch { - durable_tracker.add_entry(entry); - } - set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - observe_mrf_pending_flushed(&capped_batch, duration_millis); - dec_mrf_entries(stats.as_ref(), &capped_batch); - capped_batch.clear(); - capped_payload = None; - } - if channel_closed && capped_batch.is_empty() && !dirty { - break; - } - interval.tick().await; + }; + + if !flush_requested || pending.is_empty() { continue; } - tokio::select! { - entry = rx.recv(), if (!channel_closed || !rx.is_empty()) && pending.len() < MRF_PENDING_CAP => match entry { - Some(e) => { - observe_mrf_pending(&e); - pending.push(e); - new_entries_pending_stats += 1; - new_entries_pending_observability += 1; - dirty = true; - // Flush eagerly once enough new entries have accumulated - // since the last write (measured against the flushed - // set, not the absolute length, so a large backlog is - // not rewritten on every single add). - if recovery_applied - && new_entries_pending_stats >= 1000 - && let Some(duration_millis) = flush_mrf_to_disk(&pending, &storage).await - { - add_durable_mrf_suffix(&mut durable_tracker, &pending, new_entries_pending_stats); - set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - let observe_start = pending.len().saturating_sub(new_entries_pending_observability); - observe_mrf_pending_flushed(&pending[observe_start..], duration_millis); - new_entries_pending_observability = 0; - let stats_start = pending.len().saturating_sub(new_entries_pending_stats); - dec_mrf_entries(stats.as_ref(), &pending[stats_start..]); - new_entries_pending_stats = 0; - dirty = false; - } - } - None => { - channel_closed = true; - } - }, - _ = recovery_complete.notified(), if !recovery_applied => { - recovery_applied = true; - if let Some(retry_entries) = recovery_result.lock().await.take() { - let new_entries = pending.split_off(initial_pending_len.min(pending.len())); - pending = retry_entries; - pending.extend(new_entries); - let pending_capacity = MRF_PENDING_CAP.saturating_sub(pending.len()); - let moved_count = pending_capacity.min(capped_batch.len()); - if moved_count > 0 { - pending.extend(capped_batch.drain(..moved_count)); - new_entries_pending_stats += moved_count; - new_entries_pending_observability += moved_count; - } - durable_tracker = durable_mrf_backlog_tracker_from_entries( - &pending[..pending.len().saturating_sub(new_entries_pending_stats)], - ); - dirty = true; - } - }, - _ = interval.tick() => { - if recovery_applied - && dirty - && let Some(duration_millis) = flush_mrf_to_disk(&pending, &storage).await - { - add_durable_mrf_suffix(&mut durable_tracker, &pending, new_entries_pending_stats); - set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - let observe_start = pending.len().saturating_sub(new_entries_pending_observability); - observe_mrf_pending_flushed(&pending[observe_start..], duration_millis); - new_entries_pending_observability = 0; - if new_entries_pending_stats > 0 { - let stats_start = pending.len().saturating_sub(new_entries_pending_stats); - dec_mrf_entries(stats.as_ref(), &pending[stats_start..]); - new_entries_pending_stats = 0; - } - dirty = false; - } + if pending.len() >= MRF_PENDING_CAP && !capped { + capped = true; + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + cap = MRF_PENDING_CAP, + "MRF pending backlog reached capacity — applying backpressure" + ); + } + + match flush_mrf_to_disk(&pending, &storage, &mut pending_payload).await { + Some(result) => { + set_durable_mrf_backlog_snapshot(result.backlog); + observe_mrf_pending_flushed(&pending, result.duration_millis); + dec_mrf_entries(stats.as_ref(), &pending); + pending.clear(); + pending_payload = None; + capped = false; + } + None => { + interval.tick().await; } } } @@ -2356,27 +2241,17 @@ fn dec_mrf_entries(stats: &ReplicationStats, entries: &[MrfReplicateEntry]) { } } -/// Encodes `entries` and overwrites the MRF persistence file. -/// Returns the flush duration on success; on failure logs the error and returns `None`. +/// Appends `entries` to the MRF persistence file. +/// Returns the committed backlog and flush duration on success; on failure logs +/// the error and returns `None`. /// Callers must NOT clear their in-memory buffer on `None` 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) -> Option { - let started = Instant::now(); - match write_mrf_journal_snapshot(storage.clone(), entries).await { - Ok(()) => Some(duration_millis_u64(started.elapsed())), - Err(e) => { - 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 = %e, - "Failed to flush MRF entries to disk" - ); - None - } - } +async fn flush_mrf_to_disk( + entries: &[MrfReplicateEntry], + storage: &Arc, + pending_payload: &mut Option, +) -> Option { + append_mrf_entries_to_disk(entries, storage, pending_payload, &[]).await } async fn recover_corrupt_mrf_generation( @@ -2387,7 +2262,7 @@ async fn recover_corrupt_mrf_generation( guard: &rustfs_lock::NamespaceLockGuard, pending_payload: &mut Option, started: Instant, -) -> Option { +) -> Option { let quarantine_file = format!("{MRF_CORRUPT_FILE_PREFIX}.{}.bin", OffsetDateTime::now_utc().unix_timestamp_nanos()); if let Err(error) = ReplicationConfigStore::save_no_lock(storage.clone(), &quarantine_file, corrupt_generation.to_vec()).await { @@ -2458,7 +2333,10 @@ async fn recover_corrupt_mrf_generation( ); return None; } - Some(duration_millis_u64(started.elapsed())) + Some(MrfAppendResult { + duration_millis: duration_millis_u64(started.elapsed()), + backlog: durable_mrf_backlog_summary_from_entries(entries), + }) } async fn append_mrf_entries_to_disk( @@ -2466,9 +2344,9 @@ async fn append_mrf_entries_to_disk( storage: &Arc, pending_payload: &mut Option, known_pending: &[MrfReplicateEntry], -) -> Option { +) -> Option { if entries_to_append.is_empty() { - return Some(0); + return None; } let started = Instant::now(); let lock = match storage @@ -2484,7 +2362,7 @@ async fn append_mrf_entries_to_disk( component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_REPLICATION, error = %error, - "Failed to acquire the MRF lock before appending a capped entry" + "Failed to acquire the MRF lock before appending entries" ); return None; } @@ -2496,7 +2374,7 @@ async fn append_mrf_entries_to_disk( component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_REPLICATION, error = %error, - "Failed to acquire the MRF write lock before appending a capped entry" + "Failed to acquire the MRF write lock before appending entries" ); return None; } @@ -2515,19 +2393,12 @@ async fn append_mrf_entries_to_disk( component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_REPLICATION, error = %error, - "Failed to read MRF backlog before appending a capped entry" + "Failed to read MRF backlog before appending entries" ); return None; } }; - if pending_payload - .as_ref() - .is_some_and(|pending| pending.digest == mrf_payload_digest(¤t)) - { - return Some(duration_millis_u64(started.elapsed())); - } - let mut entries = match decode_mrf_file(¤t) { Ok(entries) => entries, Err(_) if current.is_empty() => Vec::new(), @@ -2536,7 +2407,7 @@ async fn append_mrf_entries_to_disk( component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_REPLICATION, error = %error, - "Failed to decode MRF backlog before appending a capped entry" + "Failed to decode MRF backlog before appending entries" ); let mut recovery_entries = Vec::with_capacity(known_pending.len().saturating_add(entries_to_append.len())); recovery_entries.extend_from_slice(known_pending); @@ -2557,7 +2428,12 @@ async fn append_mrf_entries_to_disk( && entries.len() >= pending.entry_count { match encode_mrf_file(&entries[..pending.entry_count]) { - Ok(prefix) if mrf_payload_digest(&prefix) == pending.digest => return Some(duration_millis_u64(started.elapsed())), + Ok(prefix) if mrf_payload_digest(&prefix) == pending.digest => { + return Some(MrfAppendResult { + duration_millis: duration_millis_u64(started.elapsed()), + backlog: durable_mrf_backlog_summary_from_entries(&entries), + }); + } Ok(_) => {} Err(error) => { observe_mrf_flush_failure(0); @@ -2565,7 +2441,7 @@ async fn append_mrf_entries_to_disk( component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_REPLICATION, error = %error, - "Failed to verify a capped MRF append after an ambiguous save error" + "Failed to verify an MRF append after an ambiguous save error" ); return None; } @@ -2581,7 +2457,7 @@ async fn append_mrf_entries_to_disk( subsystem = LOG_SUBSYSTEM_REPLICATION, count = entries.len(), error = %error, - "Failed to encode capped MRF entries for disk append" + "Failed to encode MRF entries for disk append" ); return None; } @@ -2597,7 +2473,7 @@ async fn append_mrf_entries_to_disk( component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_REPLICATION, count = entries.len(), - "Failed to append capped MRF entries because the current generation has no ETag" + "Failed to append MRF entries because the current generation has no ETag" ); return None; }; @@ -2609,7 +2485,7 @@ async fn append_mrf_entries_to_disk( subsystem = LOG_SUBSYSTEM_REPLICATION, count = entries.len(), error = %error, - "Failed to append capped MRF entries after losing the namespace lock" + "Failed to append MRF entries after losing the namespace lock" ); return None; } @@ -2628,11 +2504,14 @@ async fn append_mrf_entries_to_disk( subsystem = LOG_SUBSYSTEM_REPLICATION, count = entries.len(), error = %error, - "Failed to append capped MRF entries to disk" + "Failed to append MRF entries to disk" ); return None; } - Some(duration_millis_u64(started.elapsed())) + Some(MrfAppendResult { + duration_millis: duration_millis_u64(started.elapsed()), + backlog: durable_mrf_backlog_summary_from_entries(&entries), + }) } fn duration_millis_u64(duration: std::time::Duration) -> u64 { @@ -2648,7 +2527,7 @@ async fn load_bucket_resync_metadata( let resync_file_path = ReplicationMetadataStore::bucket_resync_file_path(bucket); - let data = match ReplicationConfigStore::read(obj_api, &resync_file_path).await { + let data = match ReplicationConfigStore::read_limited(obj_api, &resync_file_path, RESYNC_FILE_MAX_BYTES).await { Ok(data) => data, Err(EcstoreError::ConfigNotFound) => return Ok(brs), Err(err) => return Err(err), @@ -3059,10 +2938,12 @@ mod tests { use super::*; use std::collections::{HashMap, VecDeque}; use std::fmt::{Debug, Formatter}; - use std::io::Cursor; + use std::io::{self, Cursor}; + use std::pin::Pin; use std::sync::Mutex as StdMutex; use std::sync::atomic::{AtomicBool, AtomicUsize}; - use tokio::io::AsyncReadExt; + use std::task::{Context, Poll}; + use tokio::io::{AsyncRead, AsyncReadExt, ReadBuf}; use tokio::sync::Notify; use uuid::Uuid; @@ -3085,6 +2966,8 @@ mod tests { hold_first_read: AtomicBool, allow_first_read: Notify, read_count: AtomicUsize, + reported_size: StdMutex>, + stream_read_bytes: Arc, write_count: AtomicUsize, fail_next_write: AtomicBool, fail_after_write: AtomicBool, @@ -3098,6 +2981,25 @@ mod tests { shared: Arc, } + struct CountingReader { + inner: Cursor>, + bytes_read: Arc, + } + + impl AsyncRead for CountingReader { + fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + let filled_before = buf.filled().len(); + match Pin::new(&mut self.inner).poll_read(cx, buf) { + Poll::Ready(Ok(())) => { + self.bytes_read + .fetch_add(buf.filled().len().saturating_sub(filled_before), Ordering::SeqCst); + Poll::Ready(Ok(())) + } + other => other, + } + } + } + impl LoadResyncNodeStore { fn new(owner: &str, shared: Arc) -> Self { Self { @@ -3157,12 +3059,21 @@ mod tests { 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"); + let actual_size = i64::try_from(data.len()).expect("test metadata length should fit i64"); + let size = self + .shared + .reported_size + .lock() + .expect("test reported size lock should not be poisoned") + .unwrap_or(actual_size); Ok(Self::GetObjectReader { - stream: Box::new(Cursor::new(data)), + stream: Box::new(CountingReader { + inner: Cursor::new(data), + bytes_read: self.shared.stream_read_bytes.clone(), + }), object_info: ObjectInfo { size, - actual_size: size, + actual_size, etag: (!self.shared.omit_etag.load(Ordering::SeqCst)) .then(|| format!("mrf-{}", self.shared.etag_revision.load(Ordering::SeqCst))), ..Default::default() @@ -3457,8 +3368,6 @@ mod tests { mrf_replica_rx: Arc::new(Mutex::new(mrf_replica_rx)), mrf_save_tx, mrf_save_rx: Mutex::new(Some(mrf_save_rx)), - mrf_recovery_complete: Arc::new(Notify::new()), - mrf_recovery_result: Arc::new(Mutex::new(None)), mrf_worker_kill_tx, mrf_stop_tx, mrf_worker_size: AtomicI32::new(0), @@ -3709,6 +3618,8 @@ mod tests { hold_first_read: AtomicBool::new(false), allow_first_read: Notify::new(), read_count: AtomicUsize::new(0), + reported_size: StdMutex::new(None), + stream_read_bytes: Arc::new(AtomicUsize::new(0)), write_count: AtomicUsize::new(0), fail_next_write: AtomicBool::new(false), fail_after_write: AtomicBool::new(false), @@ -4230,10 +4141,11 @@ mod tests { ..Default::default() }; observe_mrf_pending(&entry); + let mut pending_payload = None; - let result = flush_mrf_to_disk(std::slice::from_ref(&entry), &storage).await; + let result = flush_mrf_to_disk(std::slice::from_ref(&entry), &storage, &mut pending_payload).await; - assert_eq!(result, None); + assert!(result.is_none()); let snapshot = mrf_backlog_observability_snapshot(); let bucket = snapshot .buckets @@ -4347,6 +4259,55 @@ mod tests { assert!(!should_auto_resume_resync(ResyncStatusType::ResyncFailed)); } + #[tokio::test] + async fn bounded_replication_config_read_accepts_exact_limit_and_caps_underreported_stream() { + const TEST_LIMIT: usize = 32; + + let shared = empty_resync_shared_state(); + *shared.data.lock().expect("test data lock should not be poisoned") = vec![0xaa; TEST_LIMIT]; + let storage = Arc::new(LoadResyncNodeStore::new("node-a", shared.clone())); + let file = ReplicationMetadataStore::bucket_resync_file_path("bounded-read"); + + let data = ReplicationConfigStore::read_limited(storage.clone(), &file, TEST_LIMIT) + .await + .expect("payload ending at the read limit should succeed"); + assert_eq!(data.len(), TEST_LIMIT); + assert_eq!(shared.stream_read_bytes.load(Ordering::SeqCst), TEST_LIMIT); + + *shared.data.lock().expect("test data lock should not be poisoned") = vec![0xbb; TEST_LIMIT * 2]; + *shared + .reported_size + .lock() + .expect("test reported size lock should not be poisoned") = + Some(i64::try_from(TEST_LIMIT).expect("test limit should fit i64")); + shared.stream_read_bytes.store(0, Ordering::SeqCst); + + let error = ReplicationConfigStore::read_limited(storage, &file, TEST_LIMIT) + .await + .expect_err("an underreported oversized payload should fail"); + assert!(matches!(error, EcstoreError::CorruptedFormat)); + assert_eq!(shared.stream_read_bytes.load(Ordering::SeqCst), TEST_LIMIT + 1); + } + + #[tokio::test] + async fn load_bucket_resync_metadata_rejects_declared_oversize_before_body_read() { + let shared = empty_resync_shared_state(); + *shared.data.lock().expect("test data lock should not be poisoned") = load_resync_test_metadata(); + *shared + .reported_size + .lock() + .expect("test reported size lock should not be poisoned") = + Some(i64::try_from(RESYNC_FILE_MAX_BYTES + 1).expect("resync limit should fit i64")); + let storage = Arc::new(LoadResyncNodeStore::new("node-a", shared.clone())); + + let error = load_bucket_resync_metadata("bounded-read", storage) + .await + .expect_err("declared oversized resync metadata should fail"); + + assert!(matches!(error, EcstoreError::CorruptedFormat)); + assert_eq!(shared.stream_read_bytes.load(Ordering::SeqCst), 0); + } + #[tokio::test] async fn load_resync_leader_lock_allows_only_one_startup_recovery() { temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { @@ -4365,6 +4326,8 @@ mod tests { hold_first_read: AtomicBool::new(false), allow_first_read: Notify::new(), read_count: AtomicUsize::new(0), + reported_size: StdMutex::new(None), + stream_read_bytes: Arc::new(AtomicUsize::new(0)), write_count: AtomicUsize::new(0), fail_next_write: AtomicBool::new(false), fail_after_write: AtomicBool::new(false), @@ -4591,29 +4554,29 @@ mod tests { tokio::time::timeout(Duration::from_secs(2), async { loop { - let quarantine_complete = { + let quarantined = { let writes = shared.writes.lock().expect("test writes lock should not be poisoned"); - let quarantined = writes + writes .iter() - .any(|(file, data)| file.starts_with(MRF_CORRUPT_FILE_PREFIX) && data == &[0xde, 0xad, 0xbe, 0xef]); - let cleared = writes - .iter() - .any(|(file, data)| file == ReplicationMetadataStore::MRF_REPLICATION_FILE && data.is_empty()); - quarantined && cleared + .any(|(file, data)| file.starts_with(MRF_CORRUPT_FILE_PREFIX) && data == &[0xde, 0xad, 0xbe, 0xef]) }; - if quarantine_complete { + if quarantined { break; } tokio::task::yield_now().await; } }) .await - .expect("quarantine should retry after the injected write failure and clear the active path"); + .expect("quarantine should retry after the injected write failure"); tokio::time::timeout(Duration::from_secs(2), 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.iter().any(|entry| entry.object == "new-failure")) { + if decode_mrf_file(&data).is_ok_and(|entries| { + ["new-failure", "second-failure"] + .iter() + .all(|object| entries.iter().any(|entry| entry.object == *object)) + }) { break; } tokio::task::yield_now().await; @@ -4634,62 +4597,43 @@ mod tests { } #[tokio::test] - async fn mrf_startup_staging_does_not_publish_entries_before_flush() { + async fn mrf_persister_appends_without_waiting_for_recovery() { temp_env::async_with_vars([("RUSTFS_REPL_MRF_FLUSH_INTERVAL_MS", Some("10"))], async { let shared = empty_resync_shared_state(); - *shared.data.lock().expect("test data lock should not be poisoned") = encode_mrf_file(&[MrfReplicateEntry { - bucket: "mrf-durable-seed".to_string(), + let seed = MrfReplicateEntry { + bucket: "mrf-append-only".to_string(), object: "seed".to_string(), - size: 7, op: MrfOpKind::Object, ..Default::default() - }]) - .expect("seed 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-durable-staged", shared.clone()))).await; - let read_started = shared.first_read_started.notified(); - let write_started = shared.write_started.notified(); + }; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(std::slice::from_ref(&seed)).expect("seed MRF backlog should encode"); + let pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-append-only", shared.clone()))).await; pool.start_mrf_persister().await; - tokio::time::timeout(Duration::from_secs(2), read_started) - .await - .expect("startup MRF read should be delayed"); pool.mrf_save_tx .send(MrfReplicateEntry { - bucket: "mrf-not-yet-durable".to_string(), - object: "staged".to_string(), - size: 11, + bucket: seed.bucket.clone(), + object: "new-failure".to_string(), op: MrfOpKind::Object, ..Default::default() }) .await - .expect("staged MRF entry should be accepted"); - tokio::time::sleep(Duration::from_millis(100)).await; - assert_eq!( - shared.write_count.load(Ordering::SeqCst), - 0, - "staged entries must not flush before startup recovery is applied" - ); - *pool.mrf_recovery_result.lock().await = Some(Vec::new()); - pool.mrf_recovery_complete.notify_one(); - tokio::time::timeout(Duration::from_secs(3), write_started) - .await - .expect("first flush should block after startup recovery is applied"); - - let snapshot = durable_mrf_backlog_summary_snapshot(); - assert_eq!( - snapshot.buckets.iter().find(|bucket| bucket.bucket == "mrf-durable-seed"), - Some(&DurableMrfBucketBacklog { - bucket: "mrf-durable-seed".to_string(), - count: 1, - bytes: 7, - }) - ); - assert!( - snapshot.buckets.iter().all(|bucket| bucket.bucket != "mrf-not-yet-durable"), - "staged entries must not appear in durable metrics before the first successful flush" - ); + .expect("new MRF failure should be accepted"); + tokio::time::timeout(Duration::from_secs(3), 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() == 2 && entries[0].object == seed.object && entries[1].object == "new-failure" + }) { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("persister should append without a recovery handoff"); + assert!(shared.last_put_no_lock.load(Ordering::SeqCst)); let handle = pool .task_handles @@ -4698,80 +4642,57 @@ mod tests { .pop() .expect("MRF persister task should be registered"); handle.abort(); - let _ = handle.await; - set_durable_mrf_backlog_snapshot(DurableMrfBacklogSnapshot::default()); }) .await; } #[tokio::test] - async fn mrf_persister_does_not_eager_flush_before_recovery_snapshot() { - temp_env::async_with_vars([("RUSTFS_REPL_MRF_FLUSH_INTERVAL_MS", Some("10"))], async { - let shared = empty_resync_shared_state(); - *shared.data.lock().expect("test data lock should not be poisoned") = - encode_mrf_file(&[MrfReplicateEntry::default()]).expect("seed MRF backlog should encode"); - shared.delay_first_read.store(true, Ordering::SeqCst); - shared.hold_first_read.store(true, Ordering::SeqCst); - let pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-recovery-gate", shared.clone()))).await; - - pool.start_mrf_processor().await; - tokio::time::timeout(Duration::from_secs(2), shared.first_read_started.notified()) - .await - .expect("processor startup read should be delayed"); - pool.start_mrf_persister().await; - tokio::time::timeout(Duration::from_secs(2), async { - loop { - if shared.read_count.load(Ordering::SeqCst) >= 2 { - break; - } - tokio::task::yield_now().await; - } - }) + async fn mrf_recovery_acknowledgement_preserves_concurrent_suffix() { + let shared = empty_resync_shared_state(); + let completed = MrfReplicateEntry { + bucket: "mrf-recovery-prefix".to_string(), + object: "completed".to_string(), + op: MrfOpKind::Delete, + force_delete_id: Some(Uuid::new_v4()), + ..Default::default() + }; + let retry = MrfReplicateEntry { + object: "retry".to_string(), + ..completed.clone() + }; + let suffix = MrfReplicateEntry { + object: "concurrent-suffix".to_string(), + ..completed.clone() + }; + let prefix = vec![completed, retry.clone()]; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(&prefix).expect("MRF recovery prefix should encode"); + let storage = Arc::new(LoadResyncNodeStore::new("mrf-recovery-prefix", shared.clone())); + let recovery_lock = storage + .new_ns_lock( + ReplicationMetadataStore::rustfs_meta_bucket(), + ReplicationMetadataStore::MRF_REPLICATION_RECOVERY_LOCK, + ) .await - .expect("persister should finish its startup snapshot while processor is delayed"); - - for index in 0..1000 { - pool.mrf_save_tx - .send(MrfReplicateEntry { - bucket: "mrf-recovery-gate".to_string(), - object: format!("staged-{index}"), - op: MrfOpKind::Object, - ..Default::default() - }) - .await - .expect("failure should be accepted before recovery completes"); - } - assert_eq!( - shared.write_count.load(Ordering::SeqCst), - 0, - "the eager 1,000-entry threshold must not flush before recovery snapshot completion" - ); - - shared.allow_first_read.notify_one(); - let processor_handle = pool.task_handles.lock().await.remove(0); - processor_handle.await.expect("processor recovery should complete"); - - tokio::time::timeout(Duration::from_secs(3), 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() == 1001) { - break; - } - tokio::task::yield_now().await; - } - }) + .expect("recovery leader lock should be created"); + let recovery_guard = recovery_lock + .get_write_lock(Duration::from_secs(1)) .await - .expect("staged failures should flush after recovery snapshot completion"); - - let persister_handle = pool - .task_handles - .lock() + .expect("recovery leader lock should be acquired"); + let mut pending_payload = None; + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&suffix), &storage, &mut pending_payload, &[]) .await - .pop() - .expect("MRF persister task should be registered"); - persister_handle.abort(); - }) - .await; + .is_some() + ); + + let retained = acknowledge_mrf_recovery(storage, &recovery_guard, &prefix, std::slice::from_ref(&retry)) + .await + .expect("recovery acknowledgement should preserve the suffix"); + + assert_eq!(retained.len(), 2); + assert_eq!(retained[0].object, retry.object); + assert_eq!(retained[1].object, suffix.object); } #[tokio::test] @@ -4942,9 +4863,10 @@ mod tests { let storage = Arc::new(LoadResyncNodeStore::new("mrf-append-idempotency", shared.clone())); let mut pending_payload = None; - assert_eq!( - append_mrf_entries_to_disk(std::slice::from_ref(&appended), &storage, &mut pending_payload, &[]).await, - None, + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&appended), &storage, &mut pending_payload, &[]) + .await + .is_none(), "the injected post-commit error should leave the capped batch pending" ); let concurrent = MrfReplicateEntry { @@ -5148,196 +5070,192 @@ mod tests { } #[tokio::test] - async fn mrf_persister_appends_startup_staging_overflow_after_recovery() { + async fn mrf_recovery_leader_lock_allows_only_one_processor() { assert!( runtime_sources::replication_pool().is_none(), "test requires the runtime replication pool to be unavailable" ); - temp_env::async_with_vars([("RUSTFS_REPL_MRF_FLUSH_INTERVAL_MS", Some("10"))], async { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { let shared = empty_resync_shared_state(); - let retained = vec![MrfReplicateEntry::default(); MRF_PENDING_CAP]; + let entry = MrfReplicateEntry { + bucket: "mrf-recovery-leader".to_string(), + object: "pending".to_string(), + op: MrfOpKind::Object, + ..Default::default() + }; *shared.data.lock().expect("test data lock should not be poisoned") = - encode_mrf_file(&retained).expect("full startup MRF backlog should encode"); + encode_mrf_file(std::slice::from_ref(&entry)).expect("MRF entry should encode"); shared.delay_first_read.store(true, Ordering::SeqCst); - let pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-staged-overflow", shared.clone()))).await; - let read_started = shared.first_read_started.notified(); + let leader_pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-leader-a", shared.clone()))).await; + let skipped_pool = + new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-leader-b", shared.clone()))).await; - pool.start_mrf_persister().await; - tokio::time::timeout(Duration::from_secs(2), read_started) + leader_pool.start_mrf_processor().await; + tokio::time::timeout(Duration::from_secs(2), shared.first_read_started.notified()) .await - .expect("persister startup read should be delayed while staging the overflow entry"); - pool.mrf_save_tx - .send(MrfReplicateEntry { - bucket: "mrf-staged-overflow".to_string(), - object: "staged-overflow".to_string(), - op: MrfOpKind::Object, - ..Default::default() - }) - .await - .expect("overflow entry should be staged before startup recovery completes"); - *pool.mrf_recovery_result.lock().await = Some(Vec::new()); - pool.mrf_recovery_complete.notify_one(); + .expect("the leader should start reading the MRF backlog"); + skipped_pool.start_mrf_processor().await; - 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() == 1 && entries[0].object == "staged-overflow") { - break; - } - tokio::task::yield_now().await; - } - }) - .await - .expect("persister should append staged overflow through its capped recovery path"); - - let handle = pool + let skipped_handle = skipped_pool .task_handles .lock() .await .pop() - .expect("MRF persister task should be registered"); - handle.abort(); - }) - .await; - } + .expect("MRF processor task should be registered for skipped node"); + skipped_handle.await.expect("skipped MRF processor should not panic"); - #[tokio::test] - async fn mrf_recovery_shrink_refills_pending_and_flushes_before_shutdown() { - assert!( - runtime_sources::replication_pool().is_none(), - "test requires the runtime replication pool to be unavailable" - ); - temp_env::async_with_vars([("RUSTFS_REPL_MRF_FLUSH_INTERVAL_MS", Some("60000"))], async { - let shared = empty_resync_shared_state(); - let retained = vec![MrfReplicateEntry::default(); MRF_PENDING_CAP]; - *shared.data.lock().expect("test data lock should not be poisoned") = - 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_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), 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 { - bucket: "mrf-recovery-shrink".to_string(), - object: object.to_string(), - op: MrfOpKind::Object, - ..Default::default() - }) - .await - .expect("overflow entry should be staged before recovery completes"); - } - *pool.mrf_recovery_result.lock().await = Some(vec![MrfReplicateEntry::default(); MRF_PENDING_CAP - 2]); - pool.mrf_recovery_complete.notify_one(); - - tokio::time::timeout(Duration::from_secs(5), write_started) + let leader_handle = leader_pool + .task_handles + .lock() .await - .expect("recovery shrink should trigger the normal pending flush before shutdown"); - assert!( - decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) - .expect("the blocked write should leave the old backlog readable") - .iter() - .all(|entry| !entry.object.starts_with("staged-overflow-")), - "the staged overflow must remain pending until the flush completes" + .pop() + .expect("MRF processor task should be registered for leader"); + leader_handle.await.expect("leader MRF processor should not panic"); + + assert_eq!( + shared.read_count.load(Ordering::SeqCst), + 2, + "only the leader should read the backlog: once for replay and once for acknowledgement" ); - shared.allow_write.notify_one(); - - tokio::time::timeout(Duration::from_secs(30), async { - loop { - if shared.write_count.load(Ordering::SeqCst) >= 2 { - break; - } - tokio::task::yield_now().await; - } - }) - .await - .expect("the recovered pending prefix should persist before the task is stopped"); - - let persisted = shared - .writes - .lock() - .expect("test writes lock should not be poisoned") - .iter() - .rev() - .find(|(file, _)| file == ReplicationMetadataStore::MRF_REPLICATION_FILE) - .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 - 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 - .task_handles - .lock() - .await - .pop() - .expect("MRF persister task should be registered"); - handle.abort(); }) .await; } - #[test] - fn mrf_durable_tracker_flushes_only_the_new_suffix() { - let retained = MrfReplicateEntry { - bucket: "mrf-tracker-suffix".to_string(), - size: 3, + #[tokio::test] + async fn mrf_recovery_acknowledgement_rejects_changed_prefix() { + let shared = empty_resync_shared_state(); + let completed = MrfReplicateEntry { + bucket: "mrf-prefix-changed".to_string(), + object: "completed".to_string(), + retry_count: 1, + op: MrfOpKind::Object, ..Default::default() }; - let appended = MrfReplicateEntry { - bucket: retained.bucket.clone(), - size: 5, - ..Default::default() + let retry = MrfReplicateEntry { + object: "retry".to_string(), + ..completed.clone() }; - let pending = [retained.clone(), appended]; - let mut tracker = durable_mrf_backlog_tracker_from_entries(std::slice::from_ref(&retained)); + let mut changed = completed.clone(); + changed.retry_count = 2; + let prefix = vec![completed, retry.clone()]; + let current = vec![changed, retry]; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(¤t).expect("changed MRF prefix should encode"); + let storage = Arc::new(LoadResyncNodeStore::new("mrf-prefix-changed", shared.clone())); + let recovery_lock = storage + .new_ns_lock( + ReplicationMetadataStore::rustfs_meta_bucket(), + ReplicationMetadataStore::MRF_REPLICATION_RECOVERY_LOCK, + ) + .await + .expect("recovery leader lock should be created"); + let recovery_guard = recovery_lock + .get_write_lock(Duration::from_secs(1)) + .await + .expect("recovery leader lock should be acquired"); - add_durable_mrf_suffix(&mut tracker, &pending, 1); + let error = acknowledge_mrf_recovery(storage.clone(), &recovery_guard, &prefix, &[]) + .await + .expect_err("acknowledgement must reject a changed recovery prefix"); - let snapshot = tracker.into_snapshot(); - assert_eq!(snapshot.summary.buckets.len(), 1); - assert_eq!(snapshot.summary.buckets[0].count, 2); - assert_eq!(snapshot.summary.buckets[0].bytes, 8); - } - - #[test] - fn move_staged_mrf_entries_releases_staging_capacity() { - let mut pending = vec![MrfReplicateEntry::default(); MRF_PENDING_CAP - 1]; - let mut staged = Vec::with_capacity(2); - staged.push(MrfReplicateEntry { - object: "pending-entry".to_string(), - ..Default::default() - }); - staged.push(MrfReplicateEntry { - object: "capped-entry".to_string(), - ..Default::default() - }); - - let capped_batch = move_staged_mrf_entries(&mut pending, &mut staged); - - assert_eq!(pending.len(), MRF_PENDING_CAP); - assert_eq!(pending.last().expect("pending prefix should be filled").object, "pending-entry"); - assert_eq!(capped_batch.len(), 1); - assert_eq!(capped_batch[0].object, "capped-entry"); - assert_eq!(staged.capacity(), 0, "the staging allocation should be released after the move"); + assert!(error.to_string().contains("prefix changed")); + let persisted = decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) + .expect("changed MRF prefix should remain readable"); + assert_eq!(persisted, current); } #[tokio::test] - async fn mrf_delete_replay_result_is_retained_when_runtime_pool_is_unavailable() { + async fn mrf_recovery_acknowledgement_write_failure_preserves_prefix() { + let shared = empty_resync_shared_state(); + let entry = MrfReplicateEntry { + bucket: "mrf-ack-write-failure".to_string(), + object: "completed".to_string(), + op: MrfOpKind::Object, + ..Default::default() + }; + let prefix = vec![entry.clone()]; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(&prefix).expect("MRF prefix should encode"); + shared.fail_next_write.store(true, Ordering::SeqCst); + let storage = Arc::new(LoadResyncNodeStore::new("mrf-ack-write-failure", shared.clone())); + let recovery_lock = storage + .new_ns_lock( + ReplicationMetadataStore::rustfs_meta_bucket(), + ReplicationMetadataStore::MRF_REPLICATION_RECOVERY_LOCK, + ) + .await + .expect("recovery leader lock should be created"); + let recovery_guard = recovery_lock + .get_write_lock(Duration::from_secs(1)) + .await + .expect("recovery leader lock should be acquired"); + + assert!( + acknowledge_mrf_recovery(storage, &recovery_guard, &prefix, &[]) + .await + .is_err(), + "acknowledgement must fail closed when the conditional save fails" + ); + let persisted = decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) + .expect("the original MRF prefix should remain readable"); + assert_eq!(persisted, prefix); + } + + #[tokio::test] + async fn mrf_appenders_accumulate_without_overwriting_each_other() { + let shared = empty_resync_shared_state(); + let first = MrfReplicateEntry { + bucket: "mrf-multi-writer".to_string(), + object: "first".to_string(), + op: MrfOpKind::Object, + ..Default::default() + }; + let second = MrfReplicateEntry { + object: "second".to_string(), + ..first.clone() + }; + let first_storage = Arc::new(LoadResyncNodeStore::new("mrf-writer-a", shared.clone())); + let second_storage = Arc::new(LoadResyncNodeStore::new("mrf-writer-b", shared.clone())); + let mut first_payload = None; + let mut second_payload = None; + + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&first), &first_storage, &mut first_payload, &[]) + .await + .is_some() + ); + assert!( + append_mrf_entries_to_disk(std::slice::from_ref(&second), &second_storage, &mut second_payload, &[]) + .await + .is_some() + ); + + let persisted = decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) + .expect("combined MRF backlog should decode"); + assert_eq!(persisted, vec![first, second]); + } + + #[test] + fn mrf_recovery_prefix_matching_checks_all_persisted_fields() { + let original = MrfReplicateEntry { + bucket: "mrf-prefix".to_string(), + object: "object".to_string(), + retry_count: 1, + size: 10, + op: MrfOpKind::Object, + ..Default::default() + }; + let mut changed = original.clone(); + changed.retry_count = 2; + assert!(!mrf_prefix_matches(&[changed], std::slice::from_ref(&original))); + + let mut suffix = original.clone(); + suffix.object = "suffix".to_string(); + assert!(mrf_prefix_matches(&[original.clone(), suffix], &[original])); + } + + #[tokio::test] + async fn mrf_delete_replay_retry_is_retained_on_disk_when_runtime_pool_is_unavailable() { assert!( runtime_sources::replication_pool().is_none(), "test requires the runtime replication pool to be unavailable" @@ -5362,14 +5280,10 @@ mod tests { .pop() .expect("MRF processor task should be registered"); handle.await.expect("MRF processor should not panic"); - let retained = pool - .mrf_recovery_result - .lock() - .await - .take() - .expect("processor should publish retry entries") + let retained = decode_mrf_file(&shared.data.lock().expect("test data lock should not be poisoned")) + .expect("processor should keep retry entries readable") .pop() - .expect("the unavailable runtime pool should retain the entry"); + .expect("the unavailable runtime pool should retain the entry on disk"); assert_eq!(retained.bucket, entry.bucket); assert_eq!(retained.object, entry.object); assert_eq!(retained.version_id, entry.version_id); diff --git a/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs b/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs index 3c4da2413..c41c6ab79 100644 --- a/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_resync_boundary.rs @@ -23,6 +23,7 @@ pub(crate) use rustfs_replication::{ pub(crate) const RESYNC_META_FORMAT: u16 = rustfs_replication::resync::RESYNC_META_FORMAT; pub(crate) const RESYNC_META_VERSION: u16 = rustfs_replication::resync::RESYNC_META_VERSION; +pub(crate) const RESYNC_FILE_MAX_BYTES: usize = rustfs_replication::RESYNC_FILE_MAX_BYTES; pub(crate) const WIRE_ZERO_TIME_UNIX: i64 = rustfs_replication::resync::WIRE_ZERO_TIME_UNIX; pub(crate) const MRF_META_FORMAT: u16 = rustfs_replication::mrf::MRF_META_FORMAT; pub(crate) const MRF_META_VERSION: u16 = rustfs_replication::mrf::MRF_META_VERSION; diff --git a/crates/ecstore/src/config/com.rs b/crates/ecstore/src/config/com.rs index f74379069..5498e570d 100644 --- a/crates/ecstore/src/config/com.rs +++ b/crates/ecstore/src/config/com.rs @@ -51,6 +51,7 @@ use serde_json::{Map, Value}; use std::collections::{HashMap, HashSet}; use std::sync::LazyLock; use std::sync::{Arc, RwLock}; +use tokio::io::AsyncReadExt; use tokio::sync::{OwnedRwLockWriteGuard, RwLock as AsyncRwLock}; use tracing::{debug, error, info, instrument, warn}; use uuid::Uuid; @@ -400,6 +401,14 @@ where Ok(data) } +pub(crate) async fn read_config_limited(api: Arc, file: &str, max_bytes: usize) -> Result> +where + S: EcstoreObjectIO, +{ + let (data, _obj) = read_config_with_metadata_inner(api, file, &ObjectOptions::default(), false, Some(max_bytes)).await?; + Ok(data) +} + /// Read an existing config object without treating an empty payload as absent. /// Callers that validate their own payload format need to distinguish corruption /// from `ConfigNotFound`. @@ -407,7 +416,7 @@ pub(crate) async fn read_config_preserve_empty(api: Arc, file: &str) -> Re where S: EcstoreObjectIO, { - let (data, _obj) = read_config_with_metadata_inner(api, file, &ObjectOptions::default(), true).await?; + let (data, _obj) = read_config_with_metadata_inner(api, file, &ObjectOptions::default(), true, None).await?; Ok(data) } @@ -447,6 +456,7 @@ where ..Default::default() }, true, + None, ) .await } @@ -463,7 +473,7 @@ where PutObjectReader = PutObjReader, >, { - read_config_with_metadata_inner(api, file, opts, false).await + read_config_with_metadata_inner(api, file, opts, false, None).await } async fn read_config_with_metadata_inner( @@ -471,6 +481,7 @@ async fn read_config_with_metadata_inner( file: &str, opts: &ObjectOptions, preserve_empty: bool, + max_bytes: Option, ) -> Result<(Vec, ObjectInfo)> where S: ObjectIO< @@ -496,7 +507,25 @@ where } })?; - let data = rd.read_all().await?; + let data = if let Some(max_bytes) = max_bytes { + let object_size = usize::try_from(rd.object_info.size).map_err(|_| Error::CorruptedFormat)?; + if object_size > max_bytes { + return Err(Error::CorruptedFormat); + } + + let read_limit = max_bytes.checked_add(1).ok_or(Error::CorruptedFormat)?; + let mut data = Vec::with_capacity(read_limit.min(64 * 1024)); + (&mut rd) + .take(u64::try_from(read_limit).map_err(|_| Error::CorruptedFormat)?) + .read_to_end(&mut data) + .await?; + if data.len() > max_bytes { + return Err(Error::CorruptedFormat); + } + data + } else { + rd.read_all().await? + }; if data.is_empty() && !preserve_empty { return Err(Error::ConfigNotFound); @@ -2354,7 +2383,7 @@ where let lock = api.new_ns_lock(RUSTFS_META_BUCKET, &transaction_lock).await?; let guard = lock.get_write_lock(get_lock_acquire_timeout()).await?; let read_options = ObjectOptions::default(); - match read_config_with_metadata_inner(api, &config_file, &read_options, true).await { + match read_config_with_metadata_inner(api, &config_file, &read_options, true, None).await { Ok((raw, object_info)) => { let (config, seed) = decode_persisted_server_config_with_seed(&raw)?; Ok(ServerConfigSnapshot { diff --git a/crates/ecstore/src/set_disk/ops/heal_walk.rs b/crates/ecstore/src/set_disk/ops/heal_walk.rs index 9bcc0a8b1..e7d920491 100644 --- a/crates/ecstore/src/set_disk/ops/heal_walk.rs +++ b/crates/ecstore/src/set_disk/ops/heal_walk.rs @@ -492,6 +492,7 @@ mod tests { batch_objects: 1000, version_budget: 10_000, objects: Mutex::new(Vec::new()), + decode_error: Mutex::new(None), version_total: AtomicUsize::new(0), truncated: AtomicBool::new(false), cancel: CancellationToken::new(), diff --git a/crates/replication/src/filemeta.rs b/crates/replication/src/filemeta.rs index 7d3a8593f..e6a57cac6 100644 --- a/crates/replication/src/filemeta.rs +++ b/crates/replication/src/filemeta.rs @@ -583,7 +583,7 @@ impl MrfOpKind { } } -#[derive(Serialize, Deserialize, Debug, Clone, Default)] +#[derive(Serialize, Deserialize, Debug, Clone, Default, PartialEq, Eq)] pub struct MrfReplicateEntry { #[serde(rename = "bucket")] pub bucket: String, diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index 330425303..6c4ec5620 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -73,9 +73,9 @@ pub use queue::{ replication_heal_queue_action, worker_queue_for_replication_type, }; pub use resync::{ - BucketReplicationResyncStatus, Error, Result, ResyncOpts, ResyncStatusType, TargetReplicationResyncStatus, - decode_resync_file, encode_resync_file, is_version_id_mismatch, resync_state_accepts_update, sanitize_resync_error_detail, - should_auto_resume_resync, should_count_head_proxy_failure, + BucketReplicationResyncStatus, Error, RESYNC_FILE_MAX_BYTES, Result, ResyncOpts, ResyncStatusType, + TargetReplicationResyncStatus, decode_resync_file, encode_resync_file, is_version_id_mismatch, resync_state_accepts_update, + sanitize_resync_error_detail, should_auto_resume_resync, should_count_head_proxy_failure, }; pub use rule::ReplicationRuleExt; pub use runtime::{ diff --git a/crates/replication/src/resync.rs b/crates/replication/src/resync.rs index 04d08f891..2380d9e18 100644 --- a/crates/replication/src/resync.rs +++ b/crates/replication/src/resync.rs @@ -24,6 +24,15 @@ use time::OffsetDateTime; pub const RESYNC_META_FORMAT: u16 = 1; pub const RESYNC_META_VERSION: u16 = 1; +// A resync snapshot contains small fixed-shape records for configured targets. These +// limits leave room for thousands of targets while bounding corrupt MessagePack work. +const RESYNC_FILE_HEADER_LEN: usize = 4; +pub const RESYNC_FILE_MAX_BYTES: usize = 16 * 1024 * 1024; +const RESYNC_MSGP_MAX_BYTES: usize = RESYNC_FILE_MAX_BYTES - RESYNC_FILE_HEADER_LEN; +const RESYNC_MSGP_MAX_ELEMENT_BYTES: usize = 1024 * 1024; +const RESYNC_MSGP_MAX_COLLECTION_ITEMS: usize = 4096; +const RESYNC_MSGP_MAX_VALUES: usize = 131_072; +const RESYNC_MSGP_MAX_DEPTH: usize = 32; const MSGP_TIME_EXT_TYPE: i8 = 5; const MSGP_TIME_LEN: u8 = 12; pub const WIRE_ZERO_TIME_UNIX: i64 = -62_135_596_800; @@ -326,7 +335,8 @@ impl BucketReplicationResyncStatus { rmp::encode::write_str(&mut wr, "v")?; rmp::encode::write_i32(&mut wr, i32::from(self.version))?; rmp::encode::write_str(&mut wr, "brs")?; - rmp::encode::write_map_len(&mut wr, self.targets_map.len() as u32)?; + let target_count = u32::try_from(self.targets_map.len()).map_err(|_| Error::CorruptedFormat)?; + rmp::encode::write_map_len(&mut wr, target_count)?; for (arn, status) in &self.targets_map { rmp::encode::write_str(&mut wr, arn)?; status.marshal_wire_msg(&mut wr)?; @@ -335,10 +345,12 @@ impl BucketReplicationResyncStatus { rmp::encode::write_i32(&mut wr, self.id)?; rmp::encode::write_str(&mut wr, "lu")?; write_msgp_time(&mut wr, wire_time_or_default(self.last_update))?; + validate_msgp_payload(&wr)?; Ok(wr) } pub fn unmarshal_msg(data: &[u8]) -> Result { + validate_msgp_payload(data)?; let mut rd = Cursor::new(data); let mut out = Self::new(); let mut fields = rmp::decode::read_map_len(&mut rd)?; @@ -353,7 +365,8 @@ impl BucketReplicationResyncStatus { } "brs" => { let map_len = rmp::decode::read_map_len(&mut rd)?; - let mut targets = HashMap::with_capacity(map_len as usize); + let target_count = usize::try_from(map_len).map_err(|_| Error::CorruptedFormat)?; + let mut targets = HashMap::with_capacity(target_count); for _ in 0..map_len { let arn = read_msgp_str(&mut rd)?; let status = TargetReplicationResyncStatus::unmarshal_wire_msg(&mut rd)?; @@ -374,6 +387,7 @@ impl BucketReplicationResyncStatus { } pub fn unmarshal_legacy_msg(data: &[u8]) -> Result { + validate_msgp_payload(data)?; let mut status: Self = rmp_serde::from_slice(data)?; for target in status.targets_map.values_mut() { target.error = target.error.as_deref().and_then(sanitize_resync_error_detail); @@ -384,7 +398,7 @@ impl BucketReplicationResyncStatus { pub fn encode_resync_file(status: &BucketReplicationResyncStatus) -> Result> { let payload = status.marshal_msg()?; - let mut data = Vec::with_capacity(4 + payload.len()); + let mut data = Vec::with_capacity(RESYNC_FILE_HEADER_LEN + payload.len()); let mut major = [0u8; 2]; LittleEndian::write_u16(&mut major, RESYNC_META_FORMAT); data.extend_from_slice(&major); @@ -396,7 +410,7 @@ pub fn encode_resync_file(status: &BucketReplicationResyncStatus) -> Result Result { - if data.len() <= 4 { + if data.len() <= RESYNC_FILE_HEADER_LEN || data.len() > RESYNC_FILE_MAX_BYTES { return Err(Error::CorruptedFormat); } @@ -437,6 +451,113 @@ fn wire_zero_time() -> OffsetDateTime { OffsetDateTime::from_unix_timestamp(WIRE_ZERO_TIME_UNIX).unwrap_or(OffsetDateTime::UNIX_EPOCH) } +fn validate_msgp_payload(data: &[u8]) -> Result<()> { + if data.len() > RESYNC_MSGP_MAX_BYTES { + return Err(Error::CorruptedFormat); + } + + let mut rd = Cursor::new(data); + let mut values = 0usize; + validate_msgp_value(&mut rd, 0, &mut values)?; + if usize::try_from(rd.position()).ok() != Some(data.len()) { + return Err(Error::CorruptedFormat); + } + Ok(()) +} + +fn validate_msgp_value(rd: &mut R, depth: usize, values: &mut usize) -> Result<()> { + if depth > RESYNC_MSGP_MAX_DEPTH { + return Err(Error::CorruptedFormat); + } + *values = values.checked_add(1).ok_or(Error::CorruptedFormat)?; + if *values > RESYNC_MSGP_MAX_VALUES { + return Err(Error::CorruptedFormat); + } + + let marker = rmp::decode::read_marker(rd).map_err(|e| Error::other(format!("{e:?}")))?; + let skip_len = match marker { + Marker::Null | Marker::False | Marker::True | Marker::FixPos(_) | Marker::FixNeg(_) => 0, + Marker::U8 | Marker::I8 => 1, + Marker::U16 | Marker::I16 => 2, + Marker::U32 | Marker::I32 | Marker::F32 => 4, + Marker::U64 | Marker::I64 | Marker::F64 => 8, + Marker::FixStr(len) => usize::from(len), + Marker::Str8 | Marker::Bin8 => read_skip_len(rd, 1)?, + Marker::Str16 | Marker::Bin16 => read_skip_len(rd, 2)?, + Marker::Str32 | Marker::Bin32 => read_skip_len(rd, 4)?, + Marker::FixArray(len) => { + validate_msgp_collection(rd, usize::from(len), depth, values, false)?; + return Ok(()); + } + Marker::Array16 => { + let len = read_skip_len(rd, 2)?; + validate_msgp_collection(rd, len, depth, values, false)?; + return Ok(()); + } + Marker::Array32 => { + let len = read_skip_len(rd, 4)?; + validate_msgp_collection(rd, len, depth, values, false)?; + return Ok(()); + } + Marker::FixMap(len) => { + validate_msgp_collection(rd, usize::from(len), depth, values, true)?; + return Ok(()); + } + Marker::Map16 => { + let len = read_skip_len(rd, 2)?; + validate_msgp_collection(rd, len, depth, values, true)?; + return Ok(()); + } + Marker::Map32 => { + let len = read_skip_len(rd, 4)?; + validate_msgp_collection(rd, len, depth, values, true)?; + return Ok(()); + } + Marker::FixExt1 => 2, + Marker::FixExt2 => 3, + Marker::FixExt4 => 5, + Marker::FixExt8 => 9, + Marker::FixExt16 => 17, + Marker::Ext8 => { + let len = validate_msgp_element_len(read_skip_len(rd, 1)?)?; + skip_exact(rd, 1)?; + return skip_exact(rd, len); + } + Marker::Ext16 => { + let len = validate_msgp_element_len(read_skip_len(rd, 2)?)?; + skip_exact(rd, 1)?; + return skip_exact(rd, len); + } + Marker::Ext32 => { + let len = validate_msgp_element_len(read_skip_len(rd, 4)?)?; + skip_exact(rd, 1)?; + return skip_exact(rd, len); + } + Marker::Reserved => return Err(Error::CorruptedFormat), + }; + let skip_len = validate_msgp_element_len(skip_len)?; + skip_exact(rd, skip_len) +} + +fn validate_msgp_collection(rd: &mut R, len: usize, depth: usize, values: &mut usize, is_map: bool) -> Result<()> { + if len > RESYNC_MSGP_MAX_COLLECTION_ITEMS { + return Err(Error::CorruptedFormat); + } + let values_per_item = if is_map { 2 } else { 1 }; + let child_count = len.checked_mul(values_per_item).ok_or(Error::CorruptedFormat)?; + for _ in 0..child_count { + validate_msgp_value(rd, depth + 1, values)?; + } + Ok(()) +} + +fn validate_msgp_element_len(len: usize) -> Result { + if len > RESYNC_MSGP_MAX_ELEMENT_BYTES { + return Err(Error::CorruptedFormat); + } + Ok(len) +} + fn read_msgp_str(rd: &mut R) -> Result { let len = rmp::decode::read_str_len(rd)? as usize; let mut buf = vec![0u8; len]; @@ -485,81 +606,8 @@ fn write_msgp_time(wr: &mut W, time: OffsetDateTime) -> Result<()> { } fn skip_msgp_value(rd: &mut R) -> Result<()> { - let marker = rmp::decode::read_marker(rd).map_err(|e| Error::other(format!("{e:?}")))?; - let skip_len: usize = match marker { - Marker::Null | Marker::False | Marker::True => 0, - Marker::FixPos(_) | Marker::FixNeg(_) => 0, - Marker::U8 => 1, - Marker::U16 => 2, - Marker::U32 => 4, - Marker::U64 => 8, - Marker::I8 => 1, - Marker::I16 => 2, - Marker::I32 => 4, - Marker::I64 => 8, - Marker::F32 => 4, - Marker::F64 => 8, - Marker::FixStr(n) => n as usize, - Marker::Str8 | Marker::Bin8 => read_skip_len(rd, 1)?, - Marker::Str16 | Marker::Bin16 => read_skip_len(rd, 2)?, - Marker::Str32 | Marker::Bin32 => read_skip_len(rd, 4)?, - Marker::FixArray(n) => { - for _ in 0..n { - skip_msgp_value(rd)?; - } - return Ok(()); - } - Marker::Array16 => { - let n = read_skip_len(rd, 2)?; - for _ in 0..n { - skip_msgp_value(rd)?; - } - return Ok(()); - } - Marker::Array32 => { - let n = read_skip_len(rd, 4)?; - for _ in 0..n { - skip_msgp_value(rd)?; - } - return Ok(()); - } - Marker::FixMap(n) => { - for _ in 0..n { - skip_msgp_value(rd)?; - skip_msgp_value(rd)?; - } - return Ok(()); - } - Marker::Map16 => { - let n = read_skip_len(rd, 2)?; - for _ in 0..n { - skip_msgp_value(rd)?; - skip_msgp_value(rd)?; - } - return Ok(()); - } - Marker::Map32 => { - let n = read_skip_len(rd, 4)?; - for _ in 0..n { - skip_msgp_value(rd)?; - skip_msgp_value(rd)?; - } - return Ok(()); - } - Marker::FixExt1 => 1 + 1, - Marker::FixExt2 => 1 + 2, - Marker::FixExt4 => 1 + 4, - Marker::FixExt8 => 1 + 8, - Marker::FixExt16 => 1 + 16, - Marker::Ext8 => 1 + read_skip_len(rd, 1)?, - Marker::Ext16 => 1 + read_skip_len(rd, 2)?, - Marker::Ext32 => 1 + read_skip_len(rd, 4)?, - Marker::Reserved => 0, - }; - if skip_len > 0 { - skip_exact(rd, skip_len)?; - } - Ok(()) + let mut values = 0usize; + validate_msgp_value(rd, 0, &mut values) } fn skip_exact(rd: &mut R, mut len: usize) -> Result<()> { @@ -573,15 +621,24 @@ fn skip_exact(rd: &mut R, mut len: usize) -> Result<()> { } fn read_skip_len(rd: &mut R, bytes: usize) -> Result { - let mut buf = vec![0u8; bytes]; - rd.read_exact(&mut buf)?; - let len = match bytes { - 1 => buf[0] as usize, - 2 => u16::from_be_bytes([buf[0], buf[1]]) as usize, - 4 => u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]) as usize, - _ => return Err(Error::other("invalid MessagePack length width")), - }; - Ok(len) + match bytes { + 1 => { + let mut buf = [0u8; 1]; + rd.read_exact(&mut buf)?; + Ok(usize::from(buf[0])) + } + 2 => { + let mut buf = [0u8; 2]; + rd.read_exact(&mut buf)?; + Ok(usize::from(u16::from_be_bytes(buf))) + } + 4 => { + let mut buf = [0u8; 4]; + rd.read_exact(&mut buf)?; + usize::try_from(u32::from_be_bytes(buf)).map_err(|_| Error::CorruptedFormat) + } + _ => Err(Error::other("invalid MessagePack length width")), + } } fn resync_status_to_i32(status: ResyncStatusType) -> i32 { @@ -611,6 +668,33 @@ fn resync_status_from_i32(code: i32) -> Result { mod tests { use super::*; + fn wrap_resync_payload(payload: &[u8]) -> Vec { + let mut data = Vec::with_capacity(RESYNC_FILE_HEADER_LEN + payload.len()); + data.extend_from_slice(&RESYNC_META_FORMAT.to_le_bytes()); + data.extend_from_slice(&RESYNC_META_VERSION.to_le_bytes()); + data.extend_from_slice(payload); + data + } + + fn resync_file_with_unknown_value(value: &[u8]) -> Vec { + let mut payload = Vec::with_capacity(32 + value.len()); + rmp::encode::write_map_len(&mut payload, 2).expect("test payload map length should encode"); + rmp::encode::write_str(&mut payload, "v").expect("test version key should encode"); + rmp::encode::write_i32(&mut payload, i32::from(RESYNC_META_VERSION)).expect("test version value should encode"); + rmp::encode::write_str(&mut payload, "future").expect("test unknown key should encode"); + payload.extend_from_slice(value); + wrap_resync_payload(&payload) + } + + fn msgp_str32(len: usize) -> Vec { + let wire_len = u32::try_from(len).expect("test string length should fit u32"); + let mut value = Vec::with_capacity(5 + len); + value.push(0xdb); + value.extend_from_slice(&wire_len.to_be_bytes()); + value.resize(5 + len, b'x'); + value + } + #[test] fn resync_status_display_matches_admin_contract() { assert_eq!(ResyncStatusType::ResyncStarted.to_string(), "Ongoing"); @@ -716,6 +800,212 @@ mod tests { assert_eq!(got.targets_map["arn:replication:a"].error.as_deref(), Some("durable failure")); } + #[test] + fn resync_file_rejects_trailing_messagepack_value() { + let status = BucketReplicationResyncStatus::new(); + let mut data = encode_resync_file(&status).expect("resync status should encode"); + data.push(0xc0); + + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_rejects_reserved_messagepack_marker() { + let data = resync_file_with_unknown_value(&[0xc1]); + + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_rejects_oversized_messagepack_element() { + let accepted = decode_resync_file(&resync_file_with_unknown_value(&msgp_str32(RESYNC_MSGP_MAX_ELEMENT_BYTES))) + .expect("element ending at the byte limit should decode"); + assert_eq!(accepted.version, RESYNC_META_VERSION); + + let value = msgp_str32(RESYNC_MSGP_MAX_ELEMENT_BYTES + 1); + let data = resync_file_with_unknown_value(&value); + + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_rejects_oversized_messagepack_collection() { + let accepted_wire_len = u32::try_from(RESYNC_MSGP_MAX_COLLECTION_ITEMS).expect("test collection length should fit u32"); + let mut accepted = Vec::with_capacity(5 + RESYNC_MSGP_MAX_COLLECTION_ITEMS); + accepted.push(0xdd); + accepted.extend_from_slice(&accepted_wire_len.to_be_bytes()); + accepted.resize(5 + RESYNC_MSGP_MAX_COLLECTION_ITEMS, 0xc0); + let mut accepted_map = Vec::with_capacity(5 + RESYNC_MSGP_MAX_COLLECTION_ITEMS * 2); + accepted_map.push(0xdf); + accepted_map.extend_from_slice(&accepted_wire_len.to_be_bytes()); + accepted_map.resize(5 + RESYNC_MSGP_MAX_COLLECTION_ITEMS * 2, 0xc0); + for value in [accepted, accepted_map] { + let accepted = decode_resync_file(&resync_file_with_unknown_value(&value)) + .expect("collection ending at the item limit should decode"); + assert_eq!(accepted.version, RESYNC_META_VERSION); + } + + let collection_len = RESYNC_MSGP_MAX_COLLECTION_ITEMS + 1; + let wire_len = u32::try_from(collection_len).expect("test collection length should fit u32"); + let mut array = Vec::with_capacity(5 + collection_len); + array.push(0xdd); + array.extend_from_slice(&wire_len.to_be_bytes()); + array.resize(5 + collection_len, 0xc0); + let mut map = Vec::with_capacity(5 + collection_len * 2); + map.push(0xdf); + map.extend_from_slice(&wire_len.to_be_bytes()); + map.resize(5 + collection_len * 2, 0xc0); + + for value in [array, map] { + let data = resync_file_with_unknown_value(&value); + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + } + + #[test] + fn resync_file_rejects_oversized_messagepack_extension() { + let extension = |len: usize| { + let wire_len = u32::try_from(len).expect("test extension length should fit u32"); + let mut extension = Vec::with_capacity(6 + len); + extension.push(0xc9); + extension.extend_from_slice(&wire_len.to_be_bytes()); + extension.push(u8::try_from(MSGP_TIME_EXT_TYPE).expect("test extension type should fit u8")); + extension.resize(6 + len, 0xaa); + extension + }; + + let accepted = decode_resync_file(&resync_file_with_unknown_value(&extension(RESYNC_MSGP_MAX_ELEMENT_BYTES))) + .expect("extension ending at the byte limit should decode"); + assert_eq!(accepted.version, RESYNC_META_VERSION); + + let data = resync_file_with_unknown_value(&extension(RESYNC_MSGP_MAX_ELEMENT_BYTES + 1)); + + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_enforces_messagepack_depth_limit() { + let mut accepted = vec![0x91; RESYNC_MSGP_MAX_DEPTH - 1]; + accepted.push(0xc0); + let accepted = decode_resync_file(&resync_file_with_unknown_value(&accepted)) + .expect("payload ending at the depth limit should decode"); + assert_eq!(accepted.version, RESYNC_META_VERSION); + + let mut rejected = vec![0x91; RESYNC_MSGP_MAX_DEPTH]; + rejected.push(0xc0); + let data = resync_file_with_unknown_value(&rejected); + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_rejects_excessive_messagepack_values() { + const ENVELOPE_VALUES: usize = 4; + const INNER_ARRAYS: usize = 32; + + let nested_arrays = |leaf_count: usize| { + let base_len = leaf_count / INNER_ARRAYS; + let longer_arrays = leaf_count % INNER_ARRAYS; + assert!(base_len + usize::from(longer_arrays > 0) <= RESYNC_MSGP_MAX_COLLECTION_ITEMS); + + let mut value = Vec::new(); + rmp::encode::write_array_len( + &mut value, + u32::try_from(INNER_ARRAYS).expect("test outer array length should fit u32"), + ) + .expect("test outer array should encode"); + for index in 0..INNER_ARRAYS { + let inner_len = base_len + usize::from(index < longer_arrays); + rmp::encode::write_array_len( + &mut value, + u32::try_from(inner_len).expect("test inner array length should fit u32"), + ) + .expect("test inner array should encode"); + value.resize(value.len() + inner_len, 0xc0); + } + value + }; + + let accepted_leaf_count = RESYNC_MSGP_MAX_VALUES - ENVELOPE_VALUES - 1 - INNER_ARRAYS; + let accepted = decode_resync_file(&resync_file_with_unknown_value(&nested_arrays(accepted_leaf_count))) + .expect("payload ending at the value limit should decode"); + assert_eq!(accepted.version, RESYNC_META_VERSION); + + let data = resync_file_with_unknown_value(&nested_arrays(accepted_leaf_count + 1)); + + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_rejects_payload_over_file_limit() { + const FULL_CHUNKS: usize = 15; + + let field_count = u32::try_from(FULL_CHUNKS + 2).expect("test field count should fit u32"); + let chunk = msgp_str32(RESYNC_MSGP_MAX_ELEMENT_BYTES); + let mut payload = Vec::with_capacity(RESYNC_MSGP_MAX_BYTES); + rmp::encode::write_map_len(&mut payload, field_count).expect("test payload map should encode"); + rmp::encode::write_str(&mut payload, "v").expect("test version key should encode"); + rmp::encode::write_i32(&mut payload, i32::from(RESYNC_META_VERSION)).expect("test version value should encode"); + for _ in 0..FULL_CHUNKS { + rmp::encode::write_str(&mut payload, "future").expect("test unknown key should encode"); + payload.extend_from_slice(&chunk); + } + rmp::encode::write_str(&mut payload, "future").expect("test final unknown key should encode"); + let remaining = RESYNC_MSGP_MAX_BYTES + .checked_sub(payload.len() + 5) + .expect("test payload should leave room for the final string"); + assert!(remaining <= RESYNC_MSGP_MAX_ELEMENT_BYTES); + payload.extend_from_slice(&msgp_str32(remaining)); + + let mut data = wrap_resync_payload(&payload); + assert_eq!(data.len(), RESYNC_FILE_MAX_BYTES); + let accepted = decode_resync_file(&data).expect("payload ending at the file limit should decode"); + assert_eq!(accepted.version, RESYNC_META_VERSION); + + data.push(0xc0); + assert_eq!(data.len(), RESYNC_FILE_MAX_BYTES + 1); + assert!(matches!(decode_resync_file(&data), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_encoder_rejects_payload_over_file_limit() { + let mut status = BucketReplicationResyncStatus::new(); + status.targets_map.insert( + "arn:replication:0".to_string(), + TargetReplicationResyncStatus { + object: "x".repeat(RESYNC_MSGP_MAX_ELEMENT_BYTES), + ..Default::default() + }, + ); + assert!(encode_resync_file(&status).is_ok(), "a single maximum-sized element should encode"); + + let target_count = RESYNC_FILE_MAX_BYTES / RESYNC_MSGP_MAX_ELEMENT_BYTES + 1; + for index in 1..target_count { + status.targets_map.insert( + format!("arn:replication:{index}"), + TargetReplicationResyncStatus { + object: "x".repeat(RESYNC_MSGP_MAX_ELEMENT_BYTES), + ..Default::default() + }, + ); + } + + assert!(matches!(encode_resync_file(&status), Err(Error::CorruptedFormat))); + } + + #[test] + fn resync_file_encoder_rejects_oversized_messagepack_element() { + let mut status = BucketReplicationResyncStatus::new(); + status.targets_map.insert( + "arn:replication:oversized-element".to_string(), + TargetReplicationResyncStatus { + object: "x".repeat(RESYNC_MSGP_MAX_ELEMENT_BYTES + 1), + ..Default::default() + }, + ); + + assert!(matches!(encode_resync_file(&status), Err(Error::CorruptedFormat))); + } + #[test] fn resync_error_detail_is_bounded_and_unicode_safe() { let detail = format!("{}尾", "x".repeat(RESYNC_ERROR_DETAIL_MAX_CHARS + 32)); @@ -770,6 +1060,24 @@ mod tests { assert_eq!(decoded.targets_map["arn:replication:a"].error.as_deref(), Some(RESYNC_ERROR_REDACTED)); } + #[test] + fn legacy_resync_payload_rejects_oversized_messagepack_element() { + let mut status = BucketReplicationResyncStatus::new(); + status.targets_map.insert( + "arn:replication:legacy-oversized".to_string(), + TargetReplicationResyncStatus { + object: "x".repeat(RESYNC_MSGP_MAX_ELEMENT_BYTES + 1), + ..Default::default() + }, + ); + let payload = rmp_serde::to_vec(&status).expect("legacy test status should encode"); + + assert!(matches!( + BucketReplicationResyncStatus::unmarshal_legacy_msg(&payload), + Err(Error::CorruptedFormat) + )); + } + #[test] fn resync_file_retains_error_for_restartable_and_failed_states() { let mut status = BucketReplicationResyncStatus::new(); diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index 098778c58..9fefb7adf 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -777,9 +777,7 @@ impl Operation for ReplicationDiffHandler { // Optional prefix can be supplied either as a query parameter (MinIO // clients) or, for RustFS clients, as a small JSON body. let mut prefix = queries.get("prefix").cloned().unwrap_or_default(); - let body = read_compatible_admin_body(req.input, MAX_ADMIN_REQUEST_BODY_SIZE, req.uri.path(), &cred.secret_key) - .await - .unwrap_or_default(); + let body = read_compatible_admin_body(req.input, MAX_ADMIN_REQUEST_BODY_SIZE, req.uri.path(), &cred.secret_key).await?; if prefix.is_empty() && !body.trim_ascii().is_empty() { match serde_json::from_slice::(&body) { Ok(parsed) => prefix = parsed.prefix, diff --git a/rustfs/src/admin/route_registration_test.rs b/rustfs/src/admin/route_registration_test.rs index a03be7560..acc1cbaab 100644 --- a/rustfs/src/admin/route_registration_test.rs +++ b/rustfs/src/admin/route_registration_test.rs @@ -1447,3 +1447,29 @@ fn test_replication_set_remote_target_compat_contract() { "set-remote-target should not re-encrypt ARN success responses" ); } + +#[test] +fn test_replication_diff_body_read_errors_are_not_ignored() { + let replication_src = include_str!("handlers/replication.rs"); + let handler_block = extract_block_between_markers( + replication_src, + "impl Operation for ReplicationDiffHandler", + "/// Failed-replication totals for one remote target", + ); + + assert!( + handler_block.contains("read_compatible_admin_body("), + "replication diff must read the optional body through the compatible admin payload reader" + ); + + let body_read_block = + extract_block_between_markers(handler_block, "let body = read_compatible_admin_body(", "if prefix.is_empty()"); + assert!( + body_read_block.contains(".await?;"), + "replication diff must fail closed when body read or MinIO-compatible decryption fails" + ); + assert!( + !body_read_block.contains(".unwrap_or_default()"), + "replication diff must not silently treat body read/decryption failures as an empty request" + ); +}