use std::collections::VecDeque;
use std::net::SocketAddr;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::time::Duration;
use dashmap::DashMap;
use freenet_stdlib::prelude::ContractInstanceId;
use tokio::time::Instant;
use crate::config::GlobalRng;
use crate::util::time_source::TimeSource;
const INVALID_BASE: Duration = Duration::from_secs(30);
const INVALID_CAP: Duration = Duration::from_secs(30 * 60);
const TIMEOUT_BASE: Duration = Duration::from_secs(120);
const TIMEOUT_CAP: Duration = Duration::from_secs(2 * 60 * 60);
const INVALID_TRIP_THRESHOLD: u32 = 3;
const TIMEOUT_TRIP_THRESHOLD: u32 = 1;
const FAILED_PAYLOAD_TTL: Duration = Duration::from_secs(10 * 60);
const MAX_FAILED_PAYLOADS_PER_CONTRACT: usize = 32;
pub(crate) const MAX_TRACKED_CONTRACTS: usize = 16_384;
const INVALID_CLEANUP_AGE: Duration = INVALID_CAP;
const CONTRACT_CLEANUP_AGE: Duration = TIMEOUT_CAP;
pub(crate) fn merge_payload_hash(is_delta: bool, payload_bytes: &[u8]) -> u64 {
use ahash::AHasher;
use std::hash::Hasher;
let mut hasher = AHasher::default();
hasher.write_u8(u8::from(is_delta));
hasher.write(payload_bytes);
hasher.finish()
}
fn cooldown_for(class: MergeFailureClass, consecutive_failures: u32) -> Duration {
let base = class.base();
let cap = class.cap();
let exponent = consecutive_failures
.saturating_sub(class.trip_threshold())
.min(20);
let raw = base.saturating_mul(1u32 << exponent.min(30));
let capped = raw.min(cap);
let jitter: f64 = GlobalRng::random_range(0.8_f64..=1.2_f64);
capped.mul_f64(jitter)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum MergeFailureClass {
Invalid,
Timeout,
}
impl MergeFailureClass {
fn base(self) -> Duration {
match self {
MergeFailureClass::Invalid => INVALID_BASE,
MergeFailureClass::Timeout => TIMEOUT_BASE,
}
}
fn cap(self) -> Duration {
match self {
MergeFailureClass::Invalid => INVALID_CAP,
MergeFailureClass::Timeout => TIMEOUT_CAP,
}
}
fn trip_threshold(self) -> u32 {
match self {
MergeFailureClass::Invalid => INVALID_TRIP_THRESHOLD,
MergeFailureClass::Timeout => TIMEOUT_TRIP_THRESHOLD,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum MergeDecision {
Allow,
InBackoff,
KnownFailedPayload,
}
impl MergeDecision {
#[cfg_attr(not(test), allow(dead_code))]
pub fn is_allowed(self) -> bool {
matches!(self, MergeDecision::Allow)
}
}
struct Cooldown {
consecutive_failures: u32,
next_allowed: Instant,
last_failure: Instant,
}
impl Cooldown {
fn inert(now: Instant) -> Self {
Self {
consecutive_failures: 0,
next_allowed: now,
last_failure: now,
}
}
fn record(&mut self, class: MergeFailureClass, now: Instant) {
self.consecutive_failures = self.consecutive_failures.saturating_add(1);
self.last_failure = now;
self.next_allowed = now + cooldown_for(class, self.consecutive_failures);
}
fn tripped(&self, class: MergeFailureClass) -> bool {
self.consecutive_failures >= class.trip_threshold()
}
fn suppressing(&self, class: MergeFailureClass, now: Instant) -> bool {
self.tripped(class) && now < self.next_allowed
}
fn past_cooldown(&self, now: Instant) -> bool {
now >= self.next_allowed
}
fn idle(&self, now: Instant, grace: Duration) -> bool {
now.saturating_duration_since(self.last_failure) > grace
}
}
struct ContractState {
timeout: Cooldown,
failed_payloads: VecDeque<(u64, Instant)>,
last_touch: Instant,
}
impl ContractState {
fn new(now: Instant) -> Self {
Self {
timeout: Cooldown::inert(now),
failed_payloads: VecDeque::new(),
last_touch: now,
}
}
fn prune_payloads(&mut self, now: Instant) {
while let Some((_, at)) = self.failed_payloads.front() {
if now.saturating_duration_since(*at) > FAILED_PAYLOAD_TTL {
self.failed_payloads.pop_front();
} else {
break;
}
}
}
fn memoize(&mut self, payload_hash: u64, now: Instant) {
self.prune_payloads(now);
if self.failed_payloads.iter().any(|(h, _)| *h == payload_hash) {
return;
}
if self.failed_payloads.len() >= MAX_FAILED_PAYLOADS_PER_CONTRACT {
self.failed_payloads.pop_front();
}
self.failed_payloads.push_back((payload_hash, now));
}
}
pub(crate) struct MergeBackoff {
invalid_by_sender: DashMap<(ContractInstanceId, SocketAddr), Cooldown>,
invalid_size: AtomicUsize,
contract: DashMap<ContractInstanceId, ContractState>,
contract_size: AtomicUsize,
max_tracked: usize,
time_source: Arc<dyn TimeSource + Send + Sync>,
suppressed_total: AtomicU64,
}
impl MergeBackoff {
pub fn new(time_source: Arc<dyn TimeSource + Send + Sync>) -> Self {
Self::with_max(time_source, MAX_TRACKED_CONTRACTS)
}
pub fn with_max(time_source: Arc<dyn TimeSource + Send + Sync>, max_tracked: usize) -> Self {
Self {
invalid_by_sender: DashMap::new(),
invalid_size: AtomicUsize::new(0),
contract: DashMap::new(),
contract_size: AtomicUsize::new(0),
max_tracked,
time_source,
suppressed_total: AtomicU64::new(0),
}
}
pub fn check(
&self,
contract: &ContractInstanceId,
sender: SocketAddr,
payload_hash: u64,
) -> MergeDecision {
let now = self.time_source.now();
let sender_suppressing = self
.invalid_by_sender
.get(&(*contract, sender))
.map(|cd| cd.suppressing(MergeFailureClass::Invalid, now))
.unwrap_or(false);
if let Some(mut cs) = self.contract.get_mut(contract) {
cs.prune_payloads(now);
if cs.failed_payloads.iter().any(|(h, _)| *h == payload_hash) {
self.suppressed_total.fetch_add(1, Ordering::Relaxed);
return MergeDecision::KnownFailedPayload;
}
if cs.timeout.suppressing(MergeFailureClass::Timeout, now) {
self.suppressed_total.fetch_add(1, Ordering::Relaxed);
return MergeDecision::InBackoff;
}
}
if sender_suppressing {
self.suppressed_total.fetch_add(1, Ordering::Relaxed);
return MergeDecision::InBackoff;
}
MergeDecision::Allow
}
pub fn record_failure(
&self,
contract: &ContractInstanceId,
sender: SocketAddr,
class: MergeFailureClass,
payload_hash: u64,
) {
let now = self.time_source.now();
if class == MergeFailureClass::Invalid {
self.record_invalid_sender(*contract, sender, now);
}
self.with_contract_entry(*contract, now, |cs| {
if class == MergeFailureClass::Timeout {
cs.timeout.record(MergeFailureClass::Timeout, now);
}
cs.last_touch = now;
cs.memoize(payload_hash, now);
});
}
fn record_invalid_sender(
&self,
contract: ContractInstanceId,
sender: SocketAddr,
now: Instant,
) {
use dashmap::mapref::entry::Entry as DEntry;
match self.invalid_by_sender.entry((contract, sender)) {
DEntry::Occupied(mut occ) => occ.get_mut().record(MergeFailureClass::Invalid, now),
DEntry::Vacant(vac) => {
let prev = self.invalid_size.fetch_add(1, Ordering::Relaxed);
if prev >= self.max_tracked {
self.invalid_size.fetch_sub(1, Ordering::Relaxed);
return;
}
let mut cd = Cooldown::inert(now);
cd.record(MergeFailureClass::Invalid, now);
vac.insert(cd);
}
}
}
fn with_contract_entry(
&self,
contract: ContractInstanceId,
now: Instant,
f: impl FnOnce(&mut ContractState),
) {
use dashmap::mapref::entry::Entry as DEntry;
match self.contract.entry(contract) {
DEntry::Occupied(mut occ) => f(occ.get_mut()),
DEntry::Vacant(vac) => {
let prev = self.contract_size.fetch_add(1, Ordering::Relaxed);
if prev >= self.max_tracked {
self.contract_size.fetch_sub(1, Ordering::Relaxed);
return;
}
let mut cs = ContractState::new(now);
f(&mut cs);
vac.insert(cs);
}
}
}
pub fn record_success_from_sender(
&self,
contract: &ContractInstanceId,
sender: SocketAddr,
changed: bool,
) {
if self
.invalid_by_sender
.remove(&(*contract, sender))
.is_some()
{
self.invalid_size.fetch_sub(1, Ordering::Relaxed);
}
if changed {
self.clear_contract_side(contract);
}
}
pub fn record_success_local(&self, contract: &ContractInstanceId, changed: bool) {
if changed {
self.clear_contract_side(contract);
}
}
pub fn invalidate_payload_memo(&self, contract: &ContractInstanceId) {
if let Some(mut cs) = self.contract.get_mut(contract) {
cs.failed_payloads.clear();
}
}
fn clear_contract_side(&self, contract: &ContractInstanceId) {
if self.contract.remove(contract).is_some() {
self.contract_size.fetch_sub(1, Ordering::Relaxed);
}
}
pub fn cleanup_expired(&self) {
let now = self.time_source.now();
let mut removed_invalid = 0usize;
self.invalid_by_sender.retain(|_, cd| {
if cd.past_cooldown(now) && cd.idle(now, INVALID_CLEANUP_AGE) {
removed_invalid += 1;
return false;
}
true
});
if removed_invalid > 0 {
self.invalid_size
.fetch_sub(removed_invalid, Ordering::Relaxed);
}
let mut removed_contract = 0usize;
self.contract.retain(|_, cs| {
cs.prune_payloads(now);
let past_cooldown = cs.timeout.past_cooldown(now);
let idle = now.saturating_duration_since(cs.last_touch) > CONTRACT_CLEANUP_AGE;
if past_cooldown && cs.failed_payloads.is_empty() && idle {
removed_contract += 1;
return false;
}
true
});
if removed_contract > 0 {
self.contract_size
.fetch_sub(removed_contract, Ordering::Relaxed);
}
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn suppressed_total(&self) -> u64 {
self.suppressed_total.load(Ordering::Relaxed)
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn contracts_in_backoff(&self) -> usize {
let now = self.time_source.now();
let mut suppressing = std::collections::HashSet::new();
for e in self.contract.iter() {
if e.value()
.timeout
.suppressing(MergeFailureClass::Timeout, now)
{
suppressing.insert(*e.key());
}
}
for e in self.invalid_by_sender.iter() {
if e.value().suppressing(MergeFailureClass::Invalid, now) {
suppressing.insert(e.key().0);
}
}
suppressing.len()
}
#[cfg(test)]
fn invalid_len(&self) -> usize {
self.invalid_by_sender.len()
}
#[cfg(test)]
fn contract_len(&self) -> usize {
self.contract.len()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::util::time_source::SharedMockTimeSource;
fn mk_contract(byte: u8) -> ContractInstanceId {
ContractInstanceId::new([byte; 32])
}
fn mk_addr(port: u16) -> SocketAddr {
SocketAddr::from(([127, 0, 0, 1], port))
}
fn mk() -> (MergeBackoff, SharedMockTimeSource) {
let ts = SharedMockTimeSource::new();
(MergeBackoff::new(Arc::new(ts.clone())), ts)
}
#[test]
fn untracked_contract_is_allowed() {
let (b, _ts) = mk();
assert_eq!(
b.check(&mk_contract(1), mk_addr(1), 42),
MergeDecision::Allow
);
assert_eq!(b.invalid_len(), 0, "check must not create an entry");
assert_eq!(b.contract_len(), 0, "check must not create an entry");
}
#[test]
fn invalid_trips_only_after_three_consecutive_failures() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Invalid, 1);
assert_eq!(b.check(&c, s, 2), MergeDecision::Allow, "failure 1");
b.record_failure(&c, s, MergeFailureClass::Invalid, 3);
assert_eq!(b.check(&c, s, 4), MergeDecision::Allow, "failure 2");
b.record_failure(&c, s, MergeFailureClass::Invalid, 5);
assert_eq!(
b.check(&c, s, 6),
MergeDecision::InBackoff,
"failure 3 trips"
);
assert_eq!(b.suppressed_total(), 1);
}
#[test]
fn invalid_backoff_is_scoped_per_sender() {
let (b, _ts) = mk();
let c = mk_contract(1);
let poison = mk_addr(10);
let healthy = mk_addr(20);
for h in [1u64, 2, 3] {
b.record_failure(&c, poison, MergeFailureClass::Invalid, h);
}
assert_eq!(
b.check(&c, poison, 99),
MergeDecision::InBackoff,
"the poison sender is suppressed on its own channel"
);
assert_eq!(
b.check(&c, healthy, 99),
MergeDecision::Allow,
"a healthy sender must still reach the merge (per-sender scoping)"
);
}
#[test]
fn invalid_two_failures_then_success_never_suppressed() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Invalid, 1);
assert!(b.check(&c, s, 2).is_allowed());
b.record_failure(&c, s, MergeFailureClass::Invalid, 3);
assert!(b.check(&c, s, 4).is_allowed());
b.record_success_from_sender(&c, s, true);
assert_eq!(b.invalid_len(), 0, "success clears the sender's channel");
assert_eq!(b.suppressed_total(), 0, "nothing was ever suppressed");
b.record_failure(&c, s, MergeFailureClass::Invalid, 5);
assert_eq!(
b.check(&c, s, 6),
MergeDecision::Allow,
"post-success failure 1 must not suppress"
);
}
#[test]
fn success_from_sender_clears_that_sender_and_contract_side() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for h in [1u64, 3, 5] {
b.record_failure(&c, s, MergeFailureClass::Invalid, h);
}
assert!(!b.check(&c, s, 2).is_allowed());
b.record_success_from_sender(&c, s, true);
assert_eq!(b.check(&c, s, 2), MergeDecision::Allow);
assert_eq!(
b.check(&c, s, 1),
MergeDecision::Allow,
"memoized payload cleared too"
);
assert_eq!(b.invalid_len(), 0);
assert_eq!(b.contract_len(), 0);
}
#[test]
fn success_from_sender_leaves_other_senders_intact() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s1 = mk_addr(10);
let s2 = mk_addr(20);
for h in [1u64, 2, 3] {
b.record_failure(&c, s1, MergeFailureClass::Invalid, h);
}
for h in [4u64, 5, 6] {
b.record_failure(&c, s2, MergeFailureClass::Invalid, h);
}
b.record_success_from_sender(&c, s1, true);
assert_eq!(b.check(&c, s1, 100), MergeDecision::Allow, "s1 cleared");
assert_eq!(
b.check(&c, s2, 101),
MergeDecision::InBackoff,
"s2's independent channel is untouched"
);
}
#[test]
fn record_success_local_clears_contract_side_leaves_sender_channels() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Timeout, 1);
for h in [2u64, 3, 4] {
b.record_failure(&c, s, MergeFailureClass::Invalid, h);
}
b.record_success_local(&c, true);
assert_eq!(b.check(&c, mk_addr(99), 100), MergeDecision::Allow);
assert_eq!(b.check(&c, s, 100), MergeDecision::InBackoff);
}
#[test]
fn contract_side_clear_requires_changed_true() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Timeout, 0xBAD);
b.record_success_from_sender(&c, s, false);
assert_eq!(
b.check(&c, mk_addr(99), 0xBAD),
MergeDecision::KnownFailedPayload,
"memo must survive a no-change success (state did not advance)"
);
assert_eq!(
b.check(&c, mk_addr(99), 0x11),
MergeDecision::InBackoff,
"Timeout cooldown must survive a no-change success"
);
b.record_success_from_sender(&c, s, true);
assert_eq!(
b.check(&c, mk_addr(99), 0xBAD),
MergeDecision::Allow,
"memo cleared on a state-advancing success"
);
assert_eq!(b.check(&c, mk_addr(99), 0x11), MergeDecision::Allow);
}
#[test]
fn sender_invalid_channel_clears_on_clean_delta_regardless_of_changed() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for h in [1u64, 2, 3] {
b.record_failure(&c, s, MergeFailureClass::Invalid, h);
}
assert_eq!(b.invalid_len(), 1);
b.record_success_from_sender(&c, s, false);
assert_eq!(
b.invalid_len(),
0,
"the sender's Invalid channel clears even on a no-change apply"
);
}
#[test]
fn invalidate_payload_memo_readmits_payload_but_keeps_cooldown_channel() {
let (b, ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for _ in 0..3 {
b.record_failure(&c, s, MergeFailureClass::Invalid, 0xBAD);
}
ts.advance_time(Duration::from_secs(60));
assert_eq!(b.check(&c, s, 0xBAD), MergeDecision::KnownFailedPayload);
assert_eq!(b.check(&c, s, 0x11), MergeDecision::Allow);
b.invalidate_payload_memo(&c);
assert_eq!(
b.check(&c, s, 0xBAD),
MergeDecision::Allow,
"memo cleared → the payload is re-admitted after the state advanced"
);
assert_eq!(
b.invalid_len(),
1,
"the cooldown channel is NOT cleared by memo invalidation"
);
}
#[test]
fn known_failed_payload_skipped_even_after_cooldown() {
let (b, ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for _ in 0..3 {
b.record_failure(&c, s, MergeFailureClass::Invalid, 99);
}
ts.advance_time(Duration::from_secs(60));
assert_eq!(b.check(&c, s, 99), MergeDecision::KnownFailedPayload);
assert_eq!(b.check(&c, s, 7), MergeDecision::Allow);
}
#[test]
fn known_failed_payload_expires_after_ttl() {
let (b, ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for _ in 0..3 {
b.record_failure(&c, s, MergeFailureClass::Invalid, 99);
}
ts.advance_time(FAILED_PAYLOAD_TTL + Duration::from_secs(1));
assert_eq!(b.check(&c, s, 99), MergeDecision::Allow);
}
#[test]
fn timeout_trips_at_first_failure_and_is_contract_wide() {
let (b, _ts) = mk();
let c = mk_contract(1);
b.record_failure(&c, mk_addr(10), MergeFailureClass::Timeout, 1);
assert_eq!(
b.check(&c, mk_addr(10), 2),
MergeDecision::InBackoff,
"a single timeout (full ~5s CPU burn) must trip immediately"
);
assert_eq!(
b.check(&c, mk_addr(20), 3),
MergeDecision::InBackoff,
"Timeout suppression is contract-wide, not per-sender"
);
}
#[test]
fn timeout_class_gets_longer_cooldown_than_invalid() {
let (b, ts) = mk();
let invalid = mk_contract(1);
let timeout = mk_contract(2);
let s = mk_addr(10);
for h in [1u64, 2, 3] {
b.record_failure(&invalid, s, MergeFailureClass::Invalid, h);
}
b.record_failure(&timeout, s, MergeFailureClass::Timeout, 1);
ts.advance_time(Duration::from_secs(60));
assert_eq!(
b.check(&invalid, s, 555),
MergeDecision::Allow,
"invalid-class cooldown (30s) should have elapsed by 60s"
);
assert_eq!(
b.check(&timeout, s, 555),
MergeDecision::InBackoff,
"timeout-class cooldown (120s) should still be active at 60s"
);
}
#[test]
fn timeout_and_invalid_tracked_independently() {
let (b, ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Invalid, 1);
b.record_failure(&c, s, MergeFailureClass::Invalid, 2);
b.record_failure(&c, s, MergeFailureClass::Timeout, 3);
let cs = b.contract.get(&c).unwrap();
let cooldown = cs
.timeout
.next_allowed
.saturating_duration_since(cs.timeout.last_failure);
drop(cs);
assert!(
cooldown >= TIMEOUT_BASE.mul_f64(0.8) && cooldown <= TIMEOUT_BASE.mul_f64(1.2),
"first Timeout cooldown {cooldown:?} must be TIMEOUT_BASE ±20%, not \
inflated by the two preceding Invalid failures"
);
ts.advance_time(Duration::from_secs(60));
b.record_failure(&c, s, MergeFailureClass::Invalid, 4);
assert_eq!(
b.check(&c, mk_addr(99), 999),
MergeDecision::InBackoff,
"the contract-wide Timeout keeps suppressing (never downgrades)"
);
}
#[test]
fn cooldown_escalates_after_trip() {
let (b, ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for h in [1u64, 2, 3] {
b.record_failure(&c, s, MergeFailureClass::Invalid, h);
}
ts.advance_time(Duration::from_secs(40));
assert_eq!(
b.check(&c, s, 10),
MergeDecision::Allow,
"first (base) cooldown should have elapsed by 40s"
);
b.record_failure(&c, s, MergeFailureClass::Invalid, 4);
ts.advance_time(Duration::from_secs(40));
assert_eq!(
b.check(&c, s, 11),
MergeDecision::InBackoff,
"post-trip failure must escalate the cooldown past 40s"
);
}
#[test]
fn cooldown_jitter_stays_within_bounds() {
for seed_byte in 0..32u8 {
let (b, _ts) = mk();
let c = mk_contract(seed_byte);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Invalid, 1);
let cd = b.invalid_by_sender.get(&(c, s)).unwrap();
let cooldown = cd.next_allowed.saturating_duration_since(cd.last_failure);
assert!(
cooldown >= Duration::from_secs(24) && cooldown <= Duration::from_secs(36),
"first-failure cooldown {cooldown:?} must be 30s ±20%"
);
}
}
#[test]
fn cooldown_capped_at_class_max() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for i in 0..40u64 {
b.record_failure(&c, s, MergeFailureClass::Invalid, i);
}
let cd = b.invalid_by_sender.get(&(c, s)).unwrap();
let cooldown = cd.next_allowed.saturating_duration_since(cd.last_failure);
assert!(
cooldown <= INVALID_CAP.mul_f64(1.2),
"escalating cooldown {cooldown:?} must stay capped near INVALID_CAP"
);
assert!(cooldown >= INVALID_CAP.mul_f64(0.8));
}
#[test]
fn tracked_senders_are_capped() {
let ts = SharedMockTimeSource::new();
let b = MergeBackoff::with_max(Arc::new(ts.clone()), 4);
let c = mk_contract(1);
for port in 0..4u16 {
b.record_failure(&c, mk_addr(port), MergeFailureClass::Invalid, 1);
}
assert_eq!(b.invalid_len(), 4);
b.record_failure(&c, mk_addr(99), MergeFailureClass::Invalid, 1);
assert_eq!(b.invalid_len(), 4, "cap must bound the per-sender channels");
b.record_failure(&c, mk_addr(0), MergeFailureClass::Invalid, 2);
assert_eq!(b.invalid_len(), 4);
}
#[test]
fn tracked_contracts_are_capped() {
let ts = SharedMockTimeSource::new();
let b = MergeBackoff::with_max(Arc::new(ts.clone()), 4);
let s = mk_addr(10);
for i in 0..4u8 {
b.record_failure(&mk_contract(i), s, MergeFailureClass::Timeout, 1);
}
assert_eq!(b.contract_len(), 4);
b.record_failure(&mk_contract(99), s, MergeFailureClass::Timeout, 1);
assert_eq!(b.contract_len(), 4, "cap must bound the contract entries");
}
#[test]
fn failed_payloads_are_bounded_per_contract() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
for i in 0..(MAX_FAILED_PAYLOADS_PER_CONTRACT as u64 + 10) {
b.record_failure(&c, s, MergeFailureClass::Invalid, i);
}
let cs = b.contract.get(&c).unwrap();
assert!(
cs.failed_payloads.len() <= MAX_FAILED_PAYLOADS_PER_CONTRACT,
"per-contract failed-payload history must stay bounded"
);
}
#[test]
fn cleanup_uses_per_class_idle_grace() {
let (b, ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Invalid, 1);
assert_eq!(b.invalid_len(), 1);
assert_eq!(b.contract_len(), 1);
ts.advance_time(INVALID_CLEANUP_AGE + Duration::from_secs(1));
b.cleanup_expired();
assert_eq!(
b.invalid_len(),
0,
"idle Invalid channel must be swept after INVALID_CLEANUP_AGE"
);
assert_eq!(
b.contract_len(),
1,
"contract entry keeps the longer CONTRACT_CLEANUP_AGE grace"
);
ts.advance_time(CONTRACT_CLEANUP_AGE);
b.cleanup_expired();
assert_eq!(
b.contract_len(),
0,
"idle contract entry must be swept after CONTRACT_CLEANUP_AGE"
);
assert_eq!(b.invalid_size.load(Ordering::Relaxed), 0);
assert_eq!(b.contract_size.load(Ordering::Relaxed), 0);
}
#[test]
fn cleanup_preserves_active_cooldown() {
let (b, ts) = mk();
let c = mk_contract(1);
b.record_failure(&c, mk_addr(10), MergeFailureClass::Timeout, 1);
ts.advance_time(Duration::from_secs(60));
b.cleanup_expired();
assert_eq!(
b.contract_len(),
1,
"an entry still in cooldown must not be swept"
);
}
#[test]
fn fork_oscillation_escalates_into_backoff_without_reset() {
let (b, ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Invalid, 0xAAAA);
ts.advance_time(Duration::from_secs(40));
b.record_failure(&c, s, MergeFailureClass::Invalid, 0xBBBB);
ts.advance_time(Duration::from_secs(40));
b.record_failure(&c, s, MergeFailureClass::Invalid, 0xCCCC);
assert_eq!(
b.check(&c, s, 0xDDDD),
MergeDecision::InBackoff,
"alternating fork deltas must trip the backoff (no resync reset)"
);
}
#[test]
fn contracts_in_backoff_gauge() {
let (b, ts) = mk();
let s = mk_addr(10);
for h in [1u64, 2, 3] {
b.record_failure(&mk_contract(1), s, MergeFailureClass::Invalid, h);
}
b.record_failure(&mk_contract(2), s, MergeFailureClass::Timeout, 1);
assert_eq!(b.contracts_in_backoff(), 2);
ts.advance_time(Duration::from_secs(60));
assert_eq!(
b.contracts_in_backoff(),
1,
"only the still-cooling Timeout contract should count"
);
}
#[test]
fn contracts_in_backoff_gauge_excludes_pre_trip_entries() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s = mk_addr(10);
b.record_failure(&c, s, MergeFailureClass::Invalid, 1);
b.record_failure(&c, s, MergeFailureClass::Invalid, 2);
assert_eq!(b.contracts_in_backoff(), 0);
assert!(b.check(&c, s, 3).is_allowed());
b.record_failure(&c, s, MergeFailureClass::Invalid, 3);
assert_eq!(b.contracts_in_backoff(), 1);
}
#[test]
fn memoized_payload_is_hard_skipped_for_any_sender_pre_trip() {
let (b, _ts) = mk();
let c = mk_contract(1);
let s1 = mk_addr(10);
let s2 = mk_addr(20);
b.record_failure(&c, s1, MergeFailureClass::Invalid, 0x11);
b.record_failure(&c, s1, MergeFailureClass::Invalid, 0x22);
assert_eq!(
b.check(&c, s1, 0x11),
MergeDecision::KnownFailedPayload,
"a memoized payload is hard-skipped immediately, even pre-trip"
);
assert_eq!(
b.check(&c, s2, 0x22),
MergeDecision::KnownFailedPayload,
"the memo is content-addressed: any sender replaying it is hard-skipped"
);
assert_eq!(
b.check(&c, s2, 0x33),
MergeDecision::Allow,
"a novel payload must still run — the memo never blocks new content"
);
}
}