diff --git a/crates/ecstore/src/bucket/replication/replication_pool.rs b/crates/ecstore/src/bucket/replication/replication_pool.rs index 3159452cf..ee27c422b 100644 --- a/crates/ecstore/src/bucket/replication/replication_pool.rs +++ b/crates/ecstore/src/bucket/replication/replication_pool.rs @@ -13,7 +13,7 @@ // limitations under the License. use super::replication_config_store::ReplicationConfigStore; -use super::replication_error_boundary::Error as EcstoreError; +use super::replication_error_boundary::{Error as EcstoreError, is_err_object_not_found, is_err_version_not_found}; use super::replication_filemeta_boundary::{ MrfOpKind, MrfReplicateEntry, REPLICATE_HEAL_DELETE, ReplicateDecision, ReplicateObjectInfo, ReplicatedTargetInfo, ReplicationStatusType, ReplicationType, ReplicationWorkerOperation, ResyncDecision, replicate_decision_for_admitted_targets, @@ -22,7 +22,7 @@ use super::replication_filemeta_boundary::{ use super::replication_lock_boundary::ReplicationLockTiming; use super::replication_logging::{EVENT_REPLICATION_CONFIG_LOOKUP_SKIPPED, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REPLICATION}; use super::replication_metadata_boundary::ReplicationMetadataStore; -use super::replication_object_config::{ReplicationConfig, check_replicate_delete, must_replicate}; +use super::replication_object_config::{ReplicationConfig, check_replicate_delete_strict, must_replicate}; use super::replication_object_decision_boundary::MustReplicateOptions; use super::replication_queue_boundary::{ DeletedObjectReplicationInfo, LARGE_WORKER_COUNT, ReplicationBackpressureRecommendation, ReplicationBackpressureState, @@ -61,6 +61,7 @@ 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; @@ -79,6 +80,7 @@ const EVENT_REPLICATION_MRF_QUEUE_UNAVAILABLE: &str = "replication_mrf_queue_una const DELETE_BATCH_ADMISSION_CONCURRENCY: usize = 16; const METRIC_DELETE_BATCH_ITEMS_TOTAL: &str = "rustfs_replication_delete_batch_items_total"; const METRIC_DELETE_BATCH_SIZE: &str = "rustfs_replication_delete_batch_size"; +const MRF_CORRUPT_FILE_PREFIX: &str = "config/replication/mrf.corrupt"; #[derive(Debug, Default)] pub struct DurableMrfBacklog { @@ -247,6 +249,7 @@ impl MrfBacklogObservabilityTracker { } } + #[cfg(test)] fn record_drop(&mut self, entry: &MrfReplicateEntry) { let bucket = self.bucket_mut(&entry.bucket); bucket.dropped_count = bucket.dropped_count.saturating_add(1); @@ -366,10 +369,6 @@ fn observe_mrf_pending_flushed(entries: &[MrfReplicateEntry], duration_millis: u update_mrf_backlog_observability(|tracker| tracker.flush_pending_entries(entries, duration_millis)); } -fn observe_mrf_drop(entry: &MrfReplicateEntry) { - update_mrf_backlog_observability(|tracker| tracker.record_drop(entry)); -} - fn observe_mrf_missed(bucket: &str) { update_mrf_backlog_observability(|tracker| tracker.record_missed(bucket)); } @@ -521,6 +520,8 @@ 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<()>, @@ -564,6 +565,8 @@ 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), @@ -577,9 +580,9 @@ impl ReplicationPool { pool.resize_failed_workers(worker_counts.mrf_workers_i32()).await; // Start background tasks + pool.start_mrf_persister().await; pool.start_mrf_processor().await; pool.start_force_delete_processor().await; - pool.start_mrf_persister().await; pool } @@ -999,12 +1002,12 @@ impl ReplicationPool { /// 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. + /// Reads the on-disk MRF file, re-injects admitted entries as Heal operations, and + /// 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 data = match ReplicationConfigStore::read(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await { @@ -1014,6 +1017,8 @@ impl ReplicationPool { available: true, buckets: Vec::new(), }); + *recovery_result.lock().await = Some(Vec::new()); + recovery_complete.notify_one(); return; } Err(e) => { @@ -1023,6 +1028,7 @@ impl ReplicationPool { error = %e, "Failed to load MRF recovery file" ); + recovery_complete.notify_one(); return; } }; @@ -1034,15 +1040,10 @@ impl ReplicationPool { component = LOG_COMPONENT_ECSTORE, subsystem = LOG_SUBSYSTEM_REPLICATION, error = %e, - "Failed to decode MRF recovery file — discarding corrupt data" + "Failed to decode MRF recovery file — preserving corrupt data" ); - // Overwrite the corrupt file so we don't fail again on next restart. - let _ = ReplicationConfigStore::save( - storage, - ReplicationMetadataStore::MRF_REPLICATION_FILE, - encode_mrf_file(&[]).unwrap_or_default(), - ) - .await; + quarantine_mrf_file(&storage, &data).await; + recovery_complete.notify_one(); return; } }; @@ -1050,9 +1051,10 @@ impl ReplicationPool { let total = entries.len(); let mut queued_count = 0usize; + let mut retry_entries = Vec::new(); for entry in entries.iter() { - match entry.op { + let admission = match entry.op { MrfOpKind::Delete => { if should_replay_force_delete_intent(entry) { let Some(operation_id) = entry.force_delete_id else { @@ -1072,81 +1074,86 @@ impl ReplicationPool { event_type: REPLICATE_HEAL_DELETE.to_string(), ..Default::default() }) - .await; - queued_count += 1; - continue; - } - if entry.force_delete_id.is_some() { - continue; - } + .await + } else if entry.force_delete_id.is_some() { + ReplicationQueueAdmission::Skipped + } else { + // 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. + // + // The MRF entry does not persist the replication decision and the + // source object is gone, so re-derive the decision from the live + // bucket config (mirroring get_heal_replicate_object_info) and set + // it on the reconstructed delete. Without this the decision string + // is empty and the delete replicates to zero targets — a silent + // no-op that leaves replicas diverged (backlog#858 / #799 B9). + let versioned = ReplicationVersioningStore::prefix_enabled(&entry.bucket, &entry.object).await; + let oi = ObjectInfo { + bucket: entry.bucket.clone(), + name: entry.object.clone(), + version_id: entry.version_id, + delete_marker: entry.delete_marker, + ..Default::default() + }; + let dsc = if entry.target_arns.is_empty() { + match ReplicationMetadataStore::optional_replication_config(&entry.bucket).await { + Ok(None) => continue, + Err(_) => { + retry_entries.push(entry.clone()); + continue; + } + Ok(Some(_)) => match check_replicate_delete_strict( + &entry.bucket, + &ObjectToDelete { + object_name: entry.object.clone(), + version_id: entry.version_id, + ..Default::default() + }, + &oi, + &ObjectOptions { + versioned, + ..Default::default() + }, + None, + ) + .await + { + Ok(dsc) => dsc, + Err(_) => { + retry_entries.push(entry.clone()); + continue; + } + }, + } + } else { + replicate_decision_for_admitted_targets(&entry.target_arns) + }; + let mut rstate = oi.replication_state(); + rstate.replicate_decision_str = dsc.to_string(); - // 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. - // - // The MRF entry does not persist the replication decision and the - // source object is gone, so re-derive the decision from the live - // bucket config (mirroring get_heal_replicate_object_info) and set - // it on the reconstructed delete. Without this the decision string - // is empty and the delete replicates to zero targets — a silent - // no-op that leaves replicas diverged (backlog#858 / #799 B9). - let versioned = ReplicationVersioningStore::prefix_enabled(&entry.bucket, &entry.object).await; - let oi = ObjectInfo { - bucket: entry.bucket.clone(), - name: entry.object.clone(), - version_id: entry.version_id, - delete_marker: entry.delete_marker, - ..Default::default() - }; - let dsc = if entry.target_arns.is_empty() { - check_replicate_delete( - &entry.bucket, - &ObjectToDelete { + let delete_marker_mtime = entry + .delete_marker_mtime + .and_then(|nanos| OffsetDateTime::from_unix_timestamp_nanos(i128::from(nanos)).ok()); + + let dv = DeletedObjectReplicationInfo { + delete_object: ReplicationDeletedObject { object_name: entry.object.clone(), version_id: entry.version_id, + delete_marker_version_id: entry.delete_marker_version_id, + delete_marker: entry.delete_marker, + delete_marker_mtime, + force_delete: entry.force_delete, + replication_state: Some(rstate), ..Default::default() }, - &oi, - &ObjectOptions { - versioned, - ..Default::default() - }, - None, - ) - .await - } else { - replicate_decision_for_admitted_targets(&entry.target_arns) - }; - let mut rstate = oi.replication_state(); - rstate.replicate_decision_str = dsc.to_string(); - - // Restore the original delete-marker mtime persisted with the entry so - // the replica keeps the source timestamp. Old MRF files lack this field - // (delete_marker_mtime = None) — fall back to None so the replica is - // stamped with the current time, preserving pre-#867 behaviour - // (backlog#867). - let delete_marker_mtime = entry - .delete_marker_mtime - .and_then(|nanos| OffsetDateTime::from_unix_timestamp_nanos(nanos as i128).ok()); - - let dv = DeletedObjectReplicationInfo { - delete_object: ReplicationDeletedObject { - object_name: entry.object.clone(), - version_id: entry.version_id, - delete_marker_version_id: entry.delete_marker_version_id, - delete_marker: entry.delete_marker, - delete_marker_mtime, - force_delete: entry.force_delete, - replication_state: Some(rstate), + bucket: entry.bucket.clone(), + op_type: ReplicationType::Heal, + event_type: REPLICATE_HEAL_DELETE.to_string(), ..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; + }; + schedule_replication_delete(dv).await + } } MrfOpKind::Object | MrfOpKind::Heal | MrfOpKind::ExistingObject => { let opts = ObjectOptions { @@ -1162,22 +1169,26 @@ impl ReplicationPool { bucket = %entry.bucket, object = %entry.object, error = %e, - "MRF recovery: object not found, skipping" + "MRF recovery: source object lookup failed" ); + if should_retry_mrf_source_lookup(&e) { + retry_entries.push(entry.clone()); + } continue; } }; if entry.target_arns.is_empty() { // Legacy entries predate target admission persistence. They cannot // be safely attributed, so retain the old live-config fallback. - queue_replication_heal(&entry.bucket, oi, entry.retry_count as u32).await; + queue_replication_heal(&entry.bucket, oi, entry.retry_count.max(0) as u32).await } else if let Some(pool) = runtime_sources::replication_pool() { let dsc = replicate_decision_for_admitted_targets(&entry.target_arns); let mut roi = replicate_object_info_from_object_info(oi, dsc, entry.op.replication_type()); roi.retry_count = entry.retry_count.max(0) as u32; - let _ = pool.queue_replica_task(roi).await; + pool.queue_replica_task(roi).await + } else { + ReplicationQueueAdmission::Missed } - queued_count += 1; } MrfOpKind::Metadata => { let opts = ObjectOptions { @@ -1193,38 +1204,28 @@ impl ReplicationPool { bucket = %entry.bucket, object = %entry.object, error = %e, - "MRF metadata recovery: object not found, skipping" + "MRF metadata recovery: source object lookup failed" ); + if should_retry_mrf_source_lookup(&e) { + retry_entries.push(entry.clone()); + } continue; } }; - queue_replication_metadata(&entry.bucket, oi, entry.retry_count as u32).await; - queued_count += 1; + queue_replication_metadata(&entry.bucket, oi, entry.retry_count.max(0) as u32).await } + }; + + if admission == ReplicationQueueAdmission::Missed { + retry_entries.push(entry.clone()); + } else if admission == ReplicationQueueAdmission::Queued { + 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) = ReplicationConfigStore::save( - storage, - ReplicationMetadataStore::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" - ); - } else { - set_durable_mrf_backlog_summary(DurableMrfBacklogSummary { - available: true, - buckets: Vec::new(), - }); - } + let retained_count = retry_entries.len(); + *recovery_result.lock().await = Some(retry_entries); + recovery_complete.notify_one(); if queued_count > 0 { info!( @@ -1232,7 +1233,8 @@ impl ReplicationPool { subsystem = LOG_SUBSYSTEM_REPLICATION, recovered = queued_count, total, - "Recovered MRF entries from disk and queued for retry" + retained = retained_count, + "Replayed MRF entries admitted for retry" ); } }); @@ -1314,8 +1316,37 @@ impl ReplicationPool { }; 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 pending = loop { + 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" + ); + tokio::time::sleep(Duration::from_millis(100)).await; + } + }, + 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" + ); + tokio::time::sleep(Duration::from_millis(100)).await; + } + } + }; + let initial_pending_len = 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 @@ -1325,52 +1356,115 @@ impl ReplicationPool { // earlier (backlog#859 / #799 B10). Bounded by `MRF_PENDING_CAP` so a // sustained failure storm can't grow it without limit. const MRF_PENDING_CAP: usize = 200_000; - let mut pending: Vec = Vec::new(); let mut durable_tracker = DurableMrfBacklogTracker { available: true, ..Default::default() }; - let mut flushed_len = 0usize; + for entry in &pending { + 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 = 0usize; + let mut new_entries_pending_observability = 0usize; let mut dirty = false; - let mut capped = false; + let mut capped = initial_pending_len >= MRF_PENDING_CAP; + let mut recovery_applied = false; + 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" + ); + } // 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()); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { - tokio::select! { - entry = rx.recv() => match entry { - Some(e) => { - if pending.len() >= MRF_PENDING_CAP { - observe_mrf_drop(&e); - dec_mrf_entries(stats.as_ref(), std::slice::from_ref(&e)); - if !capped { - capped = true; - warn!( - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REPLICATION, - cap = MRF_PENDING_CAP, - "MRF pending backlog hit cap — dropping further recovery entries for this run" - ); - } - continue; + if pending.len() >= MRF_PENDING_CAP && recovery_applied { + if dirty { + if let Some(duration_millis) = flush_mrf_to_disk(&pending, &storage).await { + 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; + } 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; + } + } + if !capped { + capped = true; + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + cap = MRF_PENDING_CAP, + "MRF pending backlog reached capacity — applying backpressure" + ); + } + let mut batch = Vec::new(); + while let Ok(entry) = rx.try_recv() { + batch.push(entry); + } + if !batch.is_empty() { + for entry in &batch { + durable_tracker.add_entry(entry); + observe_mrf_pending(entry); + } + if !dirty && let Some(duration_millis) = append_mrf_entries_to_disk(&batch, &storage).await { + set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); + observe_mrf_pending_flushed(&batch, duration_millis); + dec_mrf_entries(stats.as_ref(), &batch); + } else { + new_entries_pending_stats += batch.len(); + new_entries_pending_observability += batch.len(); + pending.extend(batch); + dirty = true; + } + } + if rx.is_closed() && rx.is_empty() && !dirty { + break; + } + interval.tick().await; + continue; + } + tokio::select! { + entry = rx.recv(), if pending.len() < MRF_PENDING_CAP => match entry { + Some(e) => { durable_tracker.add_entry(&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 pending.len() - flushed_len >= 1000 + if new_entries_pending_stats >= 1000 && let Some(duration_millis) = flush_mrf_to_disk(&pending, &storage).await { set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - observe_mrf_pending_flushed(&pending[flushed_len..], duration_millis); - dec_mrf_entries(stats.as_ref(), &pending[flushed_len..]); - flushed_len = pending.len(); + 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; } } @@ -1378,18 +1472,43 @@ impl ReplicationPool { // Channel closed (pool shutting down) — final flush. if dirty && let Some(duration_millis) = flush_mrf_to_disk(&pending, &storage).await { set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - observe_mrf_pending_flushed(&pending[flushed_len..], duration_millis); - dec_mrf_entries(stats.as_ref(), &pending[flushed_len..]); + let observe_start = pending.len().saturating_sub(new_entries_pending_observability); + observe_mrf_pending_flushed(&pending[observe_start..], duration_millis); + 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..]); + } } break; } }, + _ = 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); + durable_tracker = DurableMrfBacklogTracker { + available: true, + ..Default::default() + }; + for entry in &pending { + durable_tracker.add_entry(entry); + } + dirty = true; + } + }, _ = interval.tick() => { if dirty && let Some(duration_millis) = flush_mrf_to_disk(&pending, &storage).await { set_durable_mrf_backlog_snapshot(durable_tracker.clone().into_snapshot()); - observe_mrf_pending_flushed(&pending[flushed_len..], duration_millis); - dec_mrf_entries(stats.as_ref(), &pending[flushed_len..]); - flushed_len = pending.len(); + 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; } } @@ -1863,6 +1982,44 @@ async fn queue_mrf_save_entry( ReplicationQueueAdmission::Missed } +async fn quarantine_mrf_file(storage: &Arc, data: &[u8]) { + let quarantine_file = format!("{MRF_CORRUPT_FILE_PREFIX}.{}.bin", OffsetDateTime::now_utc().unix_timestamp_nanos()); + match ReplicationConfigStore::save(storage.clone(), &quarantine_file, data.to_vec()).await { + Ok(()) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + file = %quarantine_file, + "Quarantined corrupt MRF recovery file" + ); + // Clear the active path only after the quarantine copy succeeds. An empty + // config is treated as absent, preventing the same corrupt bytes from being + // quarantined again on every restart while preserving the forensic copy. + if let Err(error) = + ReplicationConfigStore::save(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE, Vec::new()).await + { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %error, + "Failed to clear the corrupt MRF recovery path after quarantine" + ); + } + } + Err(error) => warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + file = %quarantine_file, + error = %error, + "Failed to quarantine corrupt MRF recovery file; original was preserved" + ), + } +} + +fn should_retry_mrf_source_lookup(error: &EcstoreError) -> bool { + !is_err_object_not_found(error) && !is_err_version_not_found(error) +} + fn dec_mrf_entries(stats: &ReplicationStats, entries: &[MrfReplicateEntry]) { for entry in entries { stats.dec_q(&entry.bucket, entry.size, matches!(entry.op, MrfOpKind::Delete), ReplicationType::Heal); @@ -1908,6 +2065,41 @@ async fn flush_mrf_to_disk(entries: &[MrfReplicateEntry] } } +async fn append_mrf_entries_to_disk( + entries_to_append: &[MrfReplicateEntry], + storage: &Arc, +) -> Option { + if entries_to_append.is_empty() { + return Some(0); + } + let mut entries = match ReplicationConfigStore::read(storage.clone(), ReplicationMetadataStore::MRF_REPLICATION_FILE).await { + Ok(data) => match decode_mrf_file(&data) { + Ok(entries) => entries, + Err(error) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %error, + "Failed to decode MRF backlog before appending a capped entry" + ); + return None; + } + }, + Err(EcstoreError::ConfigNotFound) => Vec::new(), + Err(error) => { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + error = %error, + "Failed to read MRF backlog before appending a capped entry" + ); + return None; + } + }; + entries.extend_from_slice(entries_to_append); + flush_mrf_to_disk(&entries, storage).await +} + fn duration_millis_u64(duration: std::time::Duration) -> u64 { u64::try_from(duration.as_millis()).unwrap_or(u64::MAX) } @@ -2126,10 +2318,12 @@ fn replicate_object_info_from_object_info( } } -pub(crate) async fn schedule_replication_delete(dv: DeletedObjectReplicationInfo) { - if let Some(pool) = runtime_sources::replication_pool() { - let _ = pool.queue_replica_delete_task(dv.clone()).await; - } +pub(crate) async fn schedule_replication_delete(dv: DeletedObjectReplicationInfo) -> ReplicationQueueAdmission { + let admission = if let Some(pool) = runtime_sources::replication_pool() { + pool.queue_replica_delete_task(dv.clone()).await + } else { + ReplicationQueueAdmission::Missed + }; if let Some(stats) = runtime_sources::replication_stats() { let target_arns = dv.admitted_target_arns(); @@ -2151,17 +2345,20 @@ pub(crate) async fn schedule_replication_delete(dv: DeletedObjectReplicationInfo } } } + + admission } /// QueueReplicationHeal is a wrapper for queue_replication_heal_internal -pub async fn queue_replication_heal(bucket: &str, oi: ObjectInfo, retry_count: u32) { +pub async fn queue_replication_heal(bucket: &str, oi: ObjectInfo, retry_count: u32) -> ReplicationQueueAdmission { // ignore modtime zero objects if oi.mod_time.is_none() || oi.mod_time == Some(OffsetDateTime::UNIX_EPOCH) { - return; + return ReplicationQueueAdmission::Skipped; } - let rcfg = match ReplicationMetadataStore::replication_config(bucket).await { - Ok((config, _)) => config, + let rcfg = match ReplicationMetadataStore::optional_replication_config(bucket).await { + Ok(Some(config)) => config, + Ok(None) => return ReplicationQueueAdmission::Skipped, Err(err) => { debug!( event = EVENT_REPLICATION_CONFIG_LOOKUP_SKIPPED, @@ -2173,7 +2370,7 @@ pub async fn queue_replication_heal(bucket: &str, oi: ObjectInfo, retry_count: u "Skipped replication heal queue due to missing replication config" ); - return; + return ReplicationQueueAdmission::Missed; } }; @@ -2194,10 +2391,12 @@ pub async fn queue_replication_heal(bucket: &str, oi: ObjectInfo, retry_count: u }; let rcfg_wrapper = ReplicationConfig::new(Some(rcfg), tgts); - queue_replication_heal_internal(bucket, oi, rcfg_wrapper, retry_count).await; + queue_replication_heal_internal(bucket, oi, rcfg_wrapper, retry_count) + .await + .admission } -pub async fn queue_replication_metadata(bucket: &str, oi: ObjectInfo, retry_count: u32) { +pub async fn queue_replication_metadata(bucket: &str, oi: ObjectInfo, retry_count: u32) -> ReplicationQueueAdmission { let dsc = must_replicate( bucket, &oi.name, @@ -2207,13 +2406,15 @@ pub async fn queue_replication_metadata(bucket: &str, oi: ObjectInfo, retry_coun .await; if !dsc.replicate_any() { - return; + return ReplicationQueueAdmission::Skipped; } let mut roi = replicate_object_info_from_object_info(oi, dsc, ReplicationType::Metadata); roi.retry_count = retry_count; if let Some(pool) = runtime_sources::replication_pool() { - let _ = pool.queue_replica_task(roi).await; + pool.queue_replica_task(roi).await + } else { + ReplicationQueueAdmission::Missed } } @@ -2336,6 +2537,7 @@ mod tests { struct LoadResyncSharedState { data: StdMutex>, + writes: StdMutex)>>, lock_manager: Arc, first_read_started: Notify, delay_first_read: AtomicBool, @@ -2385,7 +2587,10 @@ mod tests { _h: Self::HeaderMap, _opts: &Self::ObjectOptions, ) -> Result { - if !object.ends_with("/.replication/resync.bin") && !object.ends_with("config/replication/force-delete.bin") { + if !object.ends_with("/.replication/resync.bin") + && !object.ends_with("config/replication/mrf.bin") + && !object.ends_with("config/replication/force-delete.bin") + { return Err(EcstoreError::FileNotFound); } @@ -2420,7 +2625,7 @@ mod tests { async fn put_object( &self, _bucket: &str, - _object: &str, + object: &str, data: &mut Self::PutObjectReader, _opts: &Self::ObjectOptions, ) -> Result { @@ -2433,6 +2638,11 @@ mod tests { } let mut encoded = Vec::new(); data.stream.read_to_end(&mut encoded).await.map_err(EcstoreError::from)?; + self.shared + .writes + .lock() + .expect("test writes lock should not be poisoned") + .push((object.to_string(), encoded.clone())); *self.shared.data.lock().expect("test data lock should not be poisoned") = encoded; self.shared.write_count.fetch_add(1, Ordering::SeqCst); Ok(ObjectInfo::default()) @@ -2647,6 +2857,8 @@ 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), @@ -2884,6 +3096,7 @@ mod tests { fn empty_resync_shared_state() -> Arc { Arc::new(LoadResyncSharedState { data: StdMutex::new(Vec::new()), + writes: StdMutex::new(Vec::new()), lock_manager: Arc::new(rustfs_lock::GlobalLockManager::new()), first_read_started: Notify::new(), delay_first_read: AtomicBool::new(false), @@ -3530,6 +3743,7 @@ mod tests { temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { let shared = Arc::new(LoadResyncSharedState { data: StdMutex::new(load_resync_test_metadata()), + writes: StdMutex::new(Vec::new()), lock_manager: Arc::new(rustfs_lock::GlobalLockManager::new()), first_read_started: Notify::new(), delay_first_read: AtomicBool::new(true), @@ -3617,6 +3831,164 @@ mod tests { assert!(!got.delete_marker); } + #[test] + fn mrf_object_replay_source_lookup_discards_missing_objects_and_retries_transient_errors() { + assert!(!should_retry_mrf_source_lookup(&EcstoreError::FileNotFound)); + assert!(!should_retry_mrf_source_lookup(&EcstoreError::FileVersionNotFound)); + assert!(!should_retry_mrf_source_lookup(&EcstoreError::VersionNotFound( + "bucket".to_string(), + "object".to_string(), + "version".to_string(), + ))); + assert!(should_retry_mrf_source_lookup(&EcstoreError::Unexpected)); + } + + #[test] + fn mrf_metadata_replay_source_lookup_discards_missing_objects_and_retries_transient_errors() { + for error in [EcstoreError::FileNotFound, EcstoreError::FileVersionNotFound] { + assert!(!should_retry_mrf_source_lookup(&error)); + } + assert!(should_retry_mrf_source_lookup(&EcstoreError::Unexpected)); + } + + #[tokio::test] + async fn corrupt_mrf_file_is_quarantined_without_overwriting_recovery_data() { + let shared = empty_resync_shared_state(); + let corrupt = vec![0xde, 0xad, 0xbe, 0xef]; + *shared.data.lock().expect("test data lock should not be poisoned") = corrupt.clone(); + let pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-corrupt", shared.clone()))).await; + + pool.start_mrf_processor().await; + let handle = pool + .task_handles + .lock() + .await + .pop() + .expect("MRF processor task should be registered"); + handle.await.expect("MRF processor should not panic"); + + let writes = shared.writes.lock().expect("test writes lock should not be poisoned"); + let (file, data) = writes.first().expect("corrupt MRF data should be quarantined"); + assert!(file.starts_with(MRF_CORRUPT_FILE_PREFIX)); + assert_eq!(data, &corrupt); + let marker = writes + .iter() + .find(|(file, _)| file == ReplicationMetadataStore::MRF_REPLICATION_FILE) + .expect("active MRF path should be cleared after quarantine"); + assert!(marker.1.is_empty(), "the active MRF path should be marked absent"); + } + + #[tokio::test] + async fn mrf_persister_seeds_retained_startup_entries() { + 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 { + let shared = empty_resync_shared_state(); + let retained = MrfReplicateEntry { + bucket: "mrf-replay-seed".to_string(), + object: "retained-delete".to_string(), + op: MrfOpKind::Delete, + target_arns: vec!["arn:rustfs:replication:target-a".to_string()], + ..Default::default() + }; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(std::slice::from_ref(&retained)).expect("MRF entry should encode"); + let pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-seed", shared.clone()))).await; + + pool.start_mrf_persister().await; + pool.start_mrf_processor().await; + let processor_handle = pool + .task_handles + .lock() + .await + .pop() + .expect("MRF processor task should be registered"); + processor_handle.await.expect("MRF processor should not panic"); + pool.mrf_save_tx + .send(MrfReplicateEntry { + bucket: "mrf-replay-seed".to_string(), + object: "new-failure".to_string(), + op: MrfOpKind::Object, + ..Default::default() + }) + .await + .expect("new MRF failure should be accepted"); + + tokio::time::timeout(Duration::from_secs(2), async { + loop { + let persisted = { + let writes = shared.writes.lock().expect("test writes lock should not be poisoned"); + writes + .iter() + .rev() + .find(|(file, _)| file == ReplicationMetadataStore::MRF_REPLICATION_FILE) + .map(|(_, data)| decode_mrf_file(data).expect("persisted MRF data should decode")) + }; + if let Some(entries) = persisted + && entries.iter().any(|entry| entry.object == retained.object) + && entries.iter().any(|entry| entry.object == "new-failure") + { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("persister flush should retain startup entries"); + + let persister_handle = pool + .task_handles + .lock() + .await + .pop() + .expect("MRF persister task should be registered"); + persister_handle.abort(); + }) + .await; + } + + #[tokio::test] + async fn mrf_delete_replay_result_is_retained_when_runtime_pool_is_unavailable() { + assert!( + runtime_sources::replication_pool().is_none(), + "test requires the runtime replication pool to be unavailable" + ); + let shared = empty_resync_shared_state(); + let entry = MrfReplicateEntry { + bucket: "mrf-replay-retry".to_string(), + object: "destructive-delete".to_string(), + op: MrfOpKind::Delete, + target_arns: vec!["arn:rustfs:replication:target-a".to_string()], + ..Default::default() + }; + *shared.data.lock().expect("test data lock should not be poisoned") = + encode_mrf_file(std::slice::from_ref(&entry)).expect("MRF entry should encode"); + let pool = new_test_replication_pool(Arc::new(LoadResyncNodeStore::new("mrf-retry", shared.clone()))).await; + + pool.start_mrf_processor().await; + let handle = pool + .task_handles + .lock() + .await + .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") + .pop() + .expect("the unavailable runtime pool should retain the entry"); + assert_eq!(retained.bucket, entry.bucket); + assert_eq!(retained.object, entry.object); + assert_eq!(retained.version_id, entry.version_id); + assert_eq!(retained.target_arns, entry.target_arns); + } + #[test] fn mrf_entry_delete_marker_roundtrip() { let dm_vid = Uuid::new_v4(); diff --git a/crates/replication/src/lib.rs b/crates/replication/src/lib.rs index 525325db8..c140b4fb9 100644 --- a/crates/replication/src/lib.rs +++ b/crates/replication/src/lib.rs @@ -47,7 +47,10 @@ pub use filemeta::{ VersionPurgeStatusType, get_replication_state, parse_replicate_decision, replicate_decision_for_admitted_targets, replication_statuses_map, target_reset_header, version_purge_statuses_map, }; -pub use mrf::{MrfOpKind, MrfReplicateEntry, decode_mrf_file, encode_mrf_file}; +pub use mrf::{ + MrfCapabilities, MrfCapability, MrfEnvelope, MrfEnvelopeError, MrfOpKind, MrfProtocolCapabilities, MrfReplicateEntry, + decode_mrf_file, encode_mrf_file, +}; pub use multipart::{ ReplicationMultipartPartInput, ReplicationMultipartPartPlan, ReplicationMultipartPlanError, ReplicationMultipartRange, replication_multipart_complete_actual_size, replication_multipart_part_plan, diff --git a/crates/replication/src/mrf.rs b/crates/replication/src/mrf.rs index 0c30972b4..add3b0c11 100644 --- a/crates/replication/src/mrf.rs +++ b/crates/replication/src/mrf.rs @@ -13,6 +13,7 @@ // limitations under the License. use byteorder::{ByteOrder, LittleEndian}; +use std::fmt; use crate::{Error, Result}; @@ -21,6 +22,324 @@ pub use crate::filemeta::{MrfOpKind, MrfReplicateEntry}; pub const MRF_META_FORMAT: u16 = 1; pub const MRF_META_VERSION: u16 = 1; +const MRF_ENVELOPE_MAGIC: [u8; 4] = *b"MRFE"; +const MRF_ENVELOPE_HEADER_LEN: usize = 24; +pub const MRF_ENVELOPE_FORMAT: u16 = 1; +pub const MRF_ENVELOPE_VERSION: u16 = 1; + +const CAPABILITY_OPERATION_KIND: u64 = 1 << 0; +const CAPABILITY_TARGET_ARNS: u64 = 1 << 1; +const CAPABILITY_FORCE_DELETE: u64 = 1 << 2; +const CAPABILITY_DELETE_MARKER_MTIME: u64 = 1 << 3; +const MRF_KNOWN_CAPABILITIES: u64 = + CAPABILITY_OPERATION_KIND | CAPABILITY_TARGET_ARNS | CAPABILITY_FORCE_DELETE | CAPABILITY_DELETE_MARKER_MTIME; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum MrfCapability { + OperationKind, + TargetArns, + ForceDelete, + DeleteMarkerMtime, +} + +impl MrfCapability { + const fn bit(self) -> u64 { + match self { + Self::OperationKind => CAPABILITY_OPERATION_KIND, + Self::TargetArns => CAPABILITY_TARGET_ARNS, + Self::ForceDelete => CAPABILITY_FORCE_DELETE, + Self::DeleteMarkerMtime => CAPABILITY_DELETE_MARKER_MTIME, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub struct MrfCapabilities(u64); + +impl MrfCapabilities { + pub const fn current() -> Self { + Self(MRF_KNOWN_CAPABILITIES) + } + + pub const fn empty() -> Self { + Self(0) + } + + pub const fn with(capability: MrfCapability) -> Self { + Self(capability.bit()) + } + + pub const fn bits(self) -> u64 { + self.0 + } + + pub const fn contains(self, capability: MrfCapability) -> bool { + self.0 & capability.bit() != 0 + } + + pub const fn supports(self, required: Self) -> bool { + self.0 & required.0 == required.0 + } + + pub fn from_bits(bits: u64) -> std::result::Result { + let unknown = bits & !MRF_KNOWN_CAPABILITIES; + if unknown != 0 { + return Err(MrfEnvelopeError::UnknownCapabilities { bits: unknown }); + } + Ok(Self(bits)) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct MrfProtocolCapabilities { + version: u16, + min_reader_version: u16, + capabilities: MrfCapabilities, +} + +impl MrfProtocolCapabilities { + pub const fn new(version: u16, min_reader_version: u16, capabilities: MrfCapabilities) -> Self { + Self { + version, + min_reader_version, + capabilities, + } + } + + pub const fn current() -> Self { + Self { + version: MRF_ENVELOPE_VERSION, + min_reader_version: MRF_ENVELOPE_VERSION, + capabilities: MrfCapabilities::current(), + } + } + + pub const fn version(self) -> u16 { + self.version + } + + pub const fn min_reader_version(self) -> u16 { + self.min_reader_version + } + + pub const fn capabilities(self) -> MrfCapabilities { + self.capabilities + } + + pub fn negotiate(self, peer: Self) -> std::result::Result { + MrfCapabilities::from_bits(self.capabilities.bits())?; + MrfCapabilities::from_bits(peer.capabilities.bits())?; + if self.min_reader_version > self.version { + return Err(MrfEnvelopeError::InvalidVersionRange { + version: self.version, + min_reader_version: self.min_reader_version, + }); + } + if peer.min_reader_version > peer.version { + return Err(MrfEnvelopeError::InvalidVersionRange { + version: peer.version, + min_reader_version: peer.min_reader_version, + }); + } + if self.version < MRF_ENVELOPE_VERSION { + return Err(MrfEnvelopeError::UnsupportedVersion { version: self.version }); + } + if peer.version < MRF_ENVELOPE_VERSION { + return Err(MrfEnvelopeError::UnsupportedVersion { version: peer.version }); + } + if peer.min_reader_version > self.version { + return Err(MrfEnvelopeError::RollbackFenced { + min_reader_version: peer.min_reader_version, + supported_version: self.version, + }); + } + if self.min_reader_version > peer.version { + return Err(MrfEnvelopeError::RollbackFenced { + min_reader_version: self.min_reader_version, + supported_version: peer.version, + }); + } + Ok(Self { + version: self.version.min(peer.version), + min_reader_version: self.min_reader_version.max(peer.min_reader_version), + capabilities: MrfCapabilities(self.capabilities.bits() & peer.capabilities.bits()), + }) + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MrfEnvelope { + protocol: MrfProtocolCapabilities, + payload: Vec, +} + +impl MrfEnvelope { + pub fn new(protocol: MrfProtocolCapabilities, payload: Vec) -> std::result::Result { + if protocol.version != MRF_ENVELOPE_VERSION { + return Err(MrfEnvelopeError::UnsupportedVersion { + version: protocol.version, + }); + } + if protocol.min_reader_version > protocol.version { + return Err(MrfEnvelopeError::InvalidVersionRange { + version: protocol.version, + min_reader_version: protocol.min_reader_version, + }); + } + MrfCapabilities::from_bits(protocol.capabilities.bits())?; + Ok(Self { protocol, payload }) + } + + pub const fn protocol(&self) -> MrfProtocolCapabilities { + self.protocol + } + + pub fn payload(&self) -> &[u8] { + &self.payload + } + + pub fn encode(&self) -> std::result::Result, MrfEnvelopeError> { + let payload_len: u32 = self.payload.len().try_into().map_err(|_| MrfEnvelopeError::PayloadTooLarge)?; + let mut data = Vec::with_capacity(MRF_ENVELOPE_HEADER_LEN + self.payload.len()); + data.extend_from_slice(&MRF_ENVELOPE_MAGIC); + data.extend_from_slice(&MRF_ENVELOPE_FORMAT.to_le_bytes()); + data.extend_from_slice(&self.protocol.version.to_le_bytes()); + data.extend_from_slice(&self.protocol.min_reader_version.to_le_bytes()); + data.extend_from_slice(&0u16.to_le_bytes()); + data.extend_from_slice(&self.protocol.capabilities.bits().to_le_bytes()); + data.extend_from_slice(&payload_len.to_le_bytes()); + data.extend_from_slice(&self.payload); + Ok(data) + } + + pub fn decode(data: &[u8], supported: MrfProtocolCapabilities) -> std::result::Result { + if data.len() < MRF_ENVELOPE_HEADER_LEN { + return Err(MrfEnvelopeError::Truncated); + } + if data[..4] != MRF_ENVELOPE_MAGIC { + return Err(MrfEnvelopeError::InvalidMagic); + } + let format = LittleEndian::read_u16(&data[4..6]); + if format != MRF_ENVELOPE_FORMAT { + return Err(MrfEnvelopeError::UnsupportedFormat { format }); + } + let version = LittleEndian::read_u16(&data[6..8]); + if version < MRF_ENVELOPE_VERSION { + return Err(MrfEnvelopeError::UnsupportedVersion { version }); + } + let min_reader_version = LittleEndian::read_u16(&data[8..10]); + if min_reader_version > supported.version { + return Err(MrfEnvelopeError::RollbackFenced { + min_reader_version, + supported_version: supported.version, + }); + } + if min_reader_version > version { + return Err(MrfEnvelopeError::InvalidVersionRange { + version, + min_reader_version, + }); + } + let reserved = LittleEndian::read_u16(&data[10..12]); + if reserved != 0 { + return Err(MrfEnvelopeError::ReservedHeaderBits { bits: reserved }); + } + let capabilities = MrfCapabilities::from_bits(LittleEndian::read_u64(&data[12..20]))?; + if !supported.capabilities.supports(capabilities) { + return Err(MrfEnvelopeError::MissingCapabilities { + required: capabilities.bits(), + available: supported.capabilities.bits(), + }); + } + let payload_len = LittleEndian::read_u32(&data[20..24]); + let actual_len = data.len() - MRF_ENVELOPE_HEADER_LEN; + if usize::try_from(payload_len).map_err(|_| MrfEnvelopeError::PayloadTooLarge)? != actual_len { + return Err(MrfEnvelopeError::PayloadLengthMismatch { + declared: payload_len, + actual: actual_len, + }); + } + Ok(Self { + protocol: MrfProtocolCapabilities { + version, + min_reader_version, + capabilities, + }, + payload: data[MRF_ENVELOPE_HEADER_LEN..].to_vec(), + }) + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum MrfEnvelopeError { + Truncated, + InvalidMagic, + UnsupportedFormat { + format: u16, + }, + UnsupportedVersion { + version: u16, + }, + InvalidVersionRange { + version: u16, + min_reader_version: u16, + }, + RollbackFenced { + min_reader_version: u16, + supported_version: u16, + }, + ReservedHeaderBits { + bits: u16, + }, + UnknownCapabilities { + bits: u64, + }, + MissingCapabilities { + required: u64, + available: u64, + }, + PayloadLengthMismatch { + declared: u32, + actual: usize, + }, + PayloadTooLarge, +} + +impl fmt::Display for MrfEnvelopeError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Truncated => write!(f, "truncated MRF envelope"), + Self::InvalidMagic => write!(f, "invalid MRF envelope magic"), + Self::UnsupportedFormat { format } => write!(f, "unsupported MRF envelope format {format}"), + Self::UnsupportedVersion { version } => write!(f, "unsupported MRF envelope version {version}"), + Self::InvalidVersionRange { + version, + min_reader_version, + } => { + write!(f, "invalid MRF version range: version {version}, minimum reader {min_reader_version}") + } + Self::RollbackFenced { + min_reader_version, + supported_version, + } => write!( + f, + "MRF rollback fenced: reader version {supported_version} is below required {min_reader_version}" + ), + Self::ReservedHeaderBits { bits } => write!(f, "reserved MRF envelope header bits are set: 0x{bits:04x}"), + Self::UnknownCapabilities { bits } => write!(f, "unknown MRF capability bits 0x{bits:016x}"), + Self::MissingCapabilities { required, available } => { + write!(f, "MRF capabilities 0x{available:016x} do not satisfy required 0x{required:016x}") + } + Self::PayloadLengthMismatch { declared, actual } => { + write!(f, "MRF payload length is {actual}, expected {declared}") + } + Self::PayloadTooLarge => write!(f, "MRF payload exceeds the envelope length limit"), + } + } +} + +impl std::error::Error for MrfEnvelopeError {} + pub fn encode_mrf_file(entries: &[MrfReplicateEntry]) -> 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()); @@ -56,6 +375,10 @@ mod tests { use super::*; use uuid::Uuid; + const ENVELOPE_FIXTURE: &[u8] = &[ + b'M', b'R', b'F', b'E', 1, 0, 1, 0, 1, 0, 0, 0, 15, 0, 0, 0, 0, 0, 0, 0, 3, 0, 0, 0, 1, 2, 3, + ]; + #[test] fn mrf_file_round_trips_object_metadata_and_delete_entries() { let obj_vid = Uuid::new_v4(); @@ -206,4 +529,128 @@ mod tests { assert!(matches!(decode_mrf_file(&data), Err(Error::CorruptedFormat))); } + + #[test] + fn envelope_fixture_is_stable_and_round_trips() { + let envelope = + MrfEnvelope::new(MrfProtocolCapabilities::current(), vec![1, 2, 3]).expect("current MRF envelope should be valid"); + assert_eq!(envelope.encode().expect("envelope should encode"), ENVELOPE_FIXTURE); + let decoded = MrfEnvelope::decode(ENVELOPE_FIXTURE, MrfProtocolCapabilities::current()).expect("fixture should decode"); + assert_eq!(decoded.protocol(), MrfProtocolCapabilities::current()); + assert_eq!(decoded.payload(), &[1, 2, 3]); + } + + #[test] + fn envelope_accepts_a_forward_compatible_writer_version() { + let mut version = ENVELOPE_FIXTURE.to_vec(); + version[6..8].copy_from_slice(&2u16.to_le_bytes()); + let decoded = + MrfEnvelope::decode(&version, MrfProtocolCapabilities::current()).expect("compatible v2 envelope should decode"); + assert_eq!(decoded.protocol().version(), 2); + assert_eq!(decoded.protocol().min_reader_version(), 1); + + let mut fenced = version; + fenced[8..10].copy_from_slice(&2u16.to_le_bytes()); + assert_eq!( + MrfEnvelope::decode(&fenced, MrfProtocolCapabilities::current()), + Err(MrfEnvelopeError::RollbackFenced { + min_reader_version: 2, + supported_version: 1, + }) + ); + } + + #[test] + fn envelope_rejects_unknown_capability_bits() { + let mut capabilities = ENVELOPE_FIXTURE.to_vec(); + capabilities[12..20].copy_from_slice(&(1u64 << 63).to_le_bytes()); + + assert_eq!( + MrfEnvelope::decode(&capabilities, MrfProtocolCapabilities::current()), + Err(MrfEnvelopeError::UnknownCapabilities { bits: 1u64 << 63 }) + ); + + let mut reserved = ENVELOPE_FIXTURE.to_vec(); + reserved[10..12].copy_from_slice(&1u16.to_le_bytes()); + assert_eq!( + MrfEnvelope::decode(&reserved, MrfProtocolCapabilities::current()), + Err(MrfEnvelopeError::ReservedHeaderBits { bits: 1 }) + ); + + let mut legacy = ENVELOPE_FIXTURE.to_vec(); + legacy[6..8].copy_from_slice(&0u16.to_le_bytes()); + assert_eq!( + MrfEnvelope::decode(&legacy, MrfProtocolCapabilities::current()), + Err(MrfEnvelopeError::UnsupportedVersion { version: 0 }) + ); + } + + #[test] + fn envelope_rejects_rollback_and_missing_capabilities() { + let mut rollback = ENVELOPE_FIXTURE.to_vec(); + rollback[8..10].copy_from_slice(&2u16.to_le_bytes()); + assert_eq!( + MrfEnvelope::decode(&rollback, MrfProtocolCapabilities::current()), + Err(MrfEnvelopeError::RollbackFenced { + min_reader_version: 2, + supported_version: 1, + }) + ); + + let required = MrfCapabilities::with(MrfCapability::TargetArns); + let envelope = + MrfEnvelope::new(MrfProtocolCapabilities::new(1, 1, required), Vec::new()).expect("known capability should be valid"); + let encoded = envelope.encode().expect("envelope should encode"); + assert_eq!( + MrfEnvelope::decode(&encoded, MrfProtocolCapabilities::new(1, 1, MrfCapabilities::empty())), + Err(MrfEnvelopeError::MissingCapabilities { + required: required.bits(), + available: 0, + }) + ); + } + + #[test] + fn protocol_negotiation_fences_rollback() { + let current = MrfProtocolCapabilities::current(); + let rollback = MrfProtocolCapabilities::new(1, 2, MrfCapabilities::current()); + assert_eq!( + current.negotiate(rollback), + Err(MrfEnvelopeError::InvalidVersionRange { + version: 1, + min_reader_version: 2, + }) + ); + } + + #[test] + fn protocol_negotiation_rejects_invalid_local_version_range() { + let invalid = MrfProtocolCapabilities::new(1, 2, MrfCapabilities::current()); + assert_eq!( + invalid.negotiate(MrfProtocolCapabilities::current()), + Err(MrfEnvelopeError::InvalidVersionRange { + version: 1, + min_reader_version: 2, + }) + ); + } + + #[test] + fn protocol_negotiation_intersects_capabilities() { + let local = MrfProtocolCapabilities::new(1, 1, MrfCapabilities::with(MrfCapability::TargetArns)); + let peer = MrfProtocolCapabilities::new(1, 1, MrfCapabilities::with(MrfCapability::ForceDelete)); + let negotiated = local.negotiate(peer).expect("same-version peers should negotiate"); + assert_eq!(negotiated.capabilities(), MrfCapabilities::empty()); + } + + #[test] + fn protocol_negotiation_accepts_a_forward_compatible_peer() { + let reader = MrfProtocolCapabilities::current(); + let writer = MrfProtocolCapabilities::new(2, 1, MrfCapabilities::current()); + let negotiated = reader + .negotiate(writer) + .expect("v1 reader should negotiate with a compatible v2 writer"); + assert_eq!(negotiated.version(), 1); + assert_eq!(negotiated.min_reader_version(), 1); + } }