mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-05 12:57:42 +00:00
fix(replication): harden MRF replay durability (#5659)
* feat(replication): add MRF envelope capabilities * fix(replication): retain failed MRF replay entries * fix(replication): retain transient MRF source failures * fix(replication): address MRF durability review feedback * fix(replication): preserve MRF recovery handoff * fix(replication): harden MRF recovery handoff
This commit is contained in:
@@ -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<S: ReplicationStorage> {
|
||||
mrf_replica_rx: Arc<Mutex<Receiver<ReplicationOperation>>>,
|
||||
mrf_save_tx: Sender<MrfReplicateEntry>,
|
||||
mrf_save_rx: Mutex<Option<Receiver<MrfReplicateEntry>>>,
|
||||
mrf_recovery_complete: Arc<Notify>,
|
||||
mrf_recovery_result: Arc<Mutex<Option<Vec<MrfReplicateEntry>>>>,
|
||||
|
||||
// Control channels
|
||||
mrf_worker_kill_tx: Sender<()>,
|
||||
@@ -564,6 +565,8 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
|
||||
/// 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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
available: true,
|
||||
buckets: Vec::new(),
|
||||
});
|
||||
*recovery_result.lock().await = Some(Vec::new());
|
||||
recovery_complete.notify_one();
|
||||
return;
|
||||
}
|
||||
Err(e) => {
|
||||
@@ -1023,6 +1028,7 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
|
||||
error = %e,
|
||||
"Failed to load MRF recovery file"
|
||||
);
|
||||
recovery_complete.notify_one();
|
||||
return;
|
||||
}
|
||||
};
|
||||
@@ -1034,15 +1040,10 @@ impl<S: ReplicationStorage> ReplicationPool<S> {
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
};
|
||||
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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
// 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<MrfReplicateEntry> = 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<S: ReplicationStorage> ReplicationPool<S> {
|
||||
// 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<S: ReplicationStorage>(storage: &Arc<S>, 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<S: ReplicationObjectIO>(entries: &[MrfReplicateEntry]
|
||||
}
|
||||
}
|
||||
|
||||
async fn append_mrf_entries_to_disk<S: ReplicationObjectIO>(
|
||||
entries_to_append: &[MrfReplicateEntry],
|
||||
storage: &Arc<S>,
|
||||
) -> Option<u64> {
|
||||
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<Vec<u8>>,
|
||||
writes: StdMutex<Vec<(String, Vec<u8>)>>,
|
||||
lock_manager: Arc<rustfs_lock::GlobalLockManager>,
|
||||
first_read_started: Notify,
|
||||
delay_first_read: AtomicBool,
|
||||
@@ -2385,7 +2587,10 @@ mod tests {
|
||||
_h: Self::HeaderMap,
|
||||
_opts: &Self::ObjectOptions,
|
||||
) -> Result<Self::GetObjectReader, Self::Error> {
|
||||
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<Self::ObjectInfo, Self::Error> {
|
||||
@@ -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<LoadResyncSharedState> {
|
||||
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();
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<Self, MrfEnvelopeError> {
|
||||
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<Self, MrfEnvelopeError> {
|
||||
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<u8>,
|
||||
}
|
||||
|
||||
impl MrfEnvelope {
|
||||
pub fn new(protocol: MrfProtocolCapabilities, payload: Vec<u8>) -> std::result::Result<Self, MrfEnvelopeError> {
|
||||
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<Vec<u8>, 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<Self, MrfEnvelopeError> {
|
||||
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<Vec<u8>> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user