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::util::time_source::TimeSource;
pub(crate) const DELTA_INCOMPAT_TTL: Duration = Duration::from_secs(10 * 60);
pub(crate) const DELTA_ATTRIBUTION_WINDOW: Duration = Duration::from_secs(60);
pub(crate) const INCOMPAT_TRIP_THRESHOLD: u32 = 3;
pub(crate) const MAX_TRACKED_CONTRACTS: usize = 8_192;
pub(crate) const MAX_TRACKED_DELTA_SENDS: usize = 16_384;
const CONTRACT_CLEANUP_AGE: Duration = Duration::from_secs(20 * 60);
struct ContractEntry {
local_strikes: u32,
resync_strikes: u32,
first_resync_peer: Option<SocketAddr>,
multi_peer_resync: bool,
armed_until: Option<Instant>,
last_event: Instant,
}
impl ContractEntry {
fn new(now: Instant) -> Self {
Self {
local_strikes: 0,
resync_strikes: 0,
first_resync_peer: None,
multi_peer_resync: false,
armed_until: None,
last_event: now,
}
}
fn effective_strikes(&self) -> u32 {
let resync = if self.multi_peer_resync {
self.resync_strikes
} else {
0
};
self.local_strikes.saturating_add(resync)
}
}
pub(crate) struct DeltaIncompat {
contracts: DashMap<ContractInstanceId, ContractEntry>,
contracts_size: AtomicUsize,
recent_delta_sends: DashMap<(ContractInstanceId, SocketAddr), Instant>,
sends_size: AtomicUsize,
time_source: Arc<dyn TimeSource + Send + Sync>,
armed_total: AtomicU64,
suppressed_total: AtomicU64,
}
impl DeltaIncompat {
pub fn new(time_source: Arc<dyn TimeSource + Send + Sync>) -> Self {
Self {
contracts: DashMap::new(),
contracts_size: AtomicUsize::new(0),
recent_delta_sends: DashMap::new(),
sends_size: AtomicUsize::new(0),
time_source,
armed_total: AtomicU64::new(0),
suppressed_total: AtomicU64::new(0),
}
}
pub fn record_delta_sent(&self, contract: ContractInstanceId, peer: SocketAddr) {
let now = self.time_source.now();
match self.recent_delta_sends.entry((contract, peer)) {
dashmap::mapref::entry::Entry::Occupied(mut e) => {
*e.get_mut() = now;
}
dashmap::mapref::entry::Entry::Vacant(e) => {
if self.sends_size.load(Ordering::Relaxed) >= MAX_TRACKED_DELTA_SENDS {
return;
}
e.insert(now);
self.sends_size.fetch_add(1, Ordering::Relaxed);
}
}
}
pub fn note_resync_request(&self, contract: ContractInstanceId, peer: SocketAddr) -> bool {
let now = self.time_source.now();
let Some((_, sent_at)) = self.recent_delta_sends.remove(&(contract, peer)) else {
return false;
};
self.sends_size.fetch_sub(1, Ordering::Relaxed);
if now.saturating_duration_since(sent_at) > DELTA_ATTRIBUTION_WINDOW {
return false;
}
self.note_failure(contract, now, Some(peer));
true
}
pub fn note_delta_apply_failed(&self, contract: ContractInstanceId) {
let now = self.time_source.now();
self.note_failure(contract, now, None);
}
fn note_failure(
&self,
contract: ContractInstanceId,
now: Instant,
resync_peer: Option<SocketAddr>,
) {
let mut armed = false;
{
let mut entry = match self.contracts.entry(contract) {
dashmap::mapref::entry::Entry::Occupied(e) => e.into_ref(),
dashmap::mapref::entry::Entry::Vacant(e) => {
if self.contracts_size.load(Ordering::Relaxed) >= MAX_TRACKED_CONTRACTS {
return;
}
self.contracts_size.fetch_add(1, Ordering::Relaxed);
e.insert(ContractEntry::new(now))
}
};
match resync_peer {
Some(peer) => {
entry.resync_strikes = entry.resync_strikes.saturating_add(1);
match entry.first_resync_peer {
None => entry.first_resync_peer = Some(peer),
Some(first) if first != peer => entry.multi_peer_resync = true,
Some(_) => {}
}
}
None => {
entry.local_strikes = entry.local_strikes.saturating_add(1);
}
}
entry.last_event = now;
if entry.effective_strikes() >= INCOMPAT_TRIP_THRESHOLD {
let was_armed = matches!(entry.armed_until, Some(until) if until > now);
entry.armed_until = Some(now + DELTA_INCOMPAT_TTL);
armed = !was_armed;
}
}
if armed {
self.armed_total.fetch_add(1, Ordering::Relaxed);
tracing::info!(
contract = %contract,
ttl_secs = DELTA_INCOMPAT_TTL.as_secs(),
event = "delta_incompat_armed",
"Contract repeatedly rejects deltas — sending full state instead for the TTL"
);
}
}
pub fn suppress_deltas(&self, contract: &ContractInstanceId) -> bool {
let now = self.time_source.now();
let suppressed = self
.contracts
.get(contract)
.is_some_and(|e| matches!(e.armed_until, Some(until) if until > now));
if suppressed {
self.suppressed_total.fetch_add(1, Ordering::Relaxed);
}
suppressed
}
pub fn record_delta_success(&self, contract: &ContractInstanceId) {
if self.contracts.remove(contract).is_some() {
self.contracts_size.fetch_sub(1, Ordering::Relaxed);
}
}
pub fn cleanup(&self) {
let now = self.time_source.now();
let mut removed_sends = 0usize;
self.recent_delta_sends.retain(|_, sent_at| {
let keep = now.saturating_duration_since(*sent_at) <= DELTA_ATTRIBUTION_WINDOW;
if !keep {
removed_sends += 1;
}
keep
});
if removed_sends > 0 {
self.sends_size.fetch_sub(removed_sends, Ordering::Relaxed);
}
let mut removed_contracts = 0usize;
self.contracts.retain(|_, e| {
let armed = matches!(e.armed_until, Some(until) if until > now);
let keep = armed || now.saturating_duration_since(e.last_event) <= CONTRACT_CLEANUP_AGE;
if !keep {
removed_contracts += 1;
}
keep
});
if removed_contracts > 0 {
self.contracts_size
.fetch_sub(removed_contracts, Ordering::Relaxed);
}
}
#[allow(dead_code)]
pub fn armed_total(&self) -> u64 {
self.armed_total.load(Ordering::Relaxed)
}
#[allow(dead_code)]
pub fn suppressed_total(&self) -> u64 {
self.suppressed_total.load(Ordering::Relaxed)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::util::time_source::SharedMockTimeSource;
fn contract(byte: u8) -> ContractInstanceId {
ContractInstanceId::new([byte; 32])
}
fn peer(port: u16) -> SocketAddr {
format!("127.0.0.1:{port}").parse().unwrap()
}
fn setup() -> (SharedMockTimeSource, DeltaIncompat) {
let ts = SharedMockTimeSource::new();
let memo = DeltaIncompat::new(Arc::new(ts.clone()));
(ts, memo)
}
#[test]
fn resync_after_delta_arms_after_threshold() {
let (_ts, memo) = setup();
let c = contract(1);
for i in 0..INCOMPAT_TRIP_THRESHOLD {
assert!(
!memo.suppress_deltas(&c),
"must not suppress before the threshold trips (i={i})"
);
memo.record_delta_sent(c, peer(4000 + i as u16));
assert!(
memo.note_resync_request(c, peer(4000 + i as u16)),
"attributed resync must count as a failure signal"
);
}
assert!(
memo.suppress_deltas(&c),
"after {INCOMPAT_TRIP_THRESHOLD} attributed resyncs from distinct \
peers the sender must fall back to full-state sends"
);
assert_eq!(memo.armed_total(), 1);
}
#[test]
fn single_peer_resyncs_do_not_arm_but_two_distinct_peers_do() {
let (_ts, memo) = setup();
let c = contract(11);
let lone = peer(4300);
for i in 0..(INCOMPAT_TRIP_THRESHOLD * 3) {
memo.record_delta_sent(c, lone);
assert!(
memo.note_resync_request(c, lone),
"attributed resync must still be consumed/counted (i={i})"
);
assert!(
!memo.suppress_deltas(&c),
"resync strikes from a single peer must NOT arm (i={i}) — \
one peer bouncing deltas is load/malice, not contract evidence"
);
}
assert_eq!(memo.armed_total(), 0);
let second = peer(4301);
memo.record_delta_sent(c, second);
assert!(memo.note_resync_request(c, second));
assert!(
memo.suppress_deltas(&c),
"a second distinct peer must corroborate the resync strikes and arm"
);
assert_eq!(memo.armed_total(), 1);
}
#[test]
fn unattributed_resync_does_not_count() {
let (_ts, memo) = setup();
let c = contract(2);
for _ in 0..10 {
assert!(!memo.note_resync_request(c, peer(4100)));
}
assert!(!memo.suppress_deltas(&c));
assert_eq!(memo.armed_total(), 0);
}
#[test]
fn stale_attribution_does_not_count() {
let (ts, memo) = setup();
let c = contract(3);
memo.record_delta_sent(c, peer(4200));
ts.advance_time(DELTA_ATTRIBUTION_WINDOW + Duration::from_secs(1));
assert!(!memo.note_resync_request(c, peer(4200)));
assert!(!memo.suppress_deltas(&c));
}
#[test]
fn local_delta_apply_failures_arm() {
let (_ts, memo) = setup();
let c = contract(4);
for _ in 0..INCOMPAT_TRIP_THRESHOLD {
assert!(!memo.suppress_deltas(&c));
memo.note_delta_apply_failed(c);
}
assert!(memo.suppress_deltas(&c));
}
#[test]
fn delta_success_clears_memo_and_counter() {
let (_ts, memo) = setup();
let c = contract(5);
for _ in 0..INCOMPAT_TRIP_THRESHOLD {
memo.note_delta_apply_failed(c);
}
assert!(memo.suppress_deltas(&c));
memo.record_delta_success(&c);
assert!(
!memo.suppress_deltas(&c),
"a successful delta apply must clear the memo"
);
for _ in 0..(INCOMPAT_TRIP_THRESHOLD - 1) {
memo.note_delta_apply_failed(c);
}
assert!(
!memo.suppress_deltas(&c),
"the failure counter must restart from zero after a success"
);
}
#[test]
fn memo_expires_after_ttl() {
let (ts, memo) = setup();
let c = contract(6);
for _ in 0..INCOMPAT_TRIP_THRESHOLD {
memo.note_delta_apply_failed(c);
}
assert!(memo.suppress_deltas(&c));
ts.advance_time(DELTA_INCOMPAT_TTL + Duration::from_secs(1));
assert!(
!memo.suppress_deltas(&c),
"the memo must expire after DELTA_INCOMPAT_TTL"
);
}
#[test]
fn refailure_after_expiry_rearms() {
let (ts, memo) = setup();
let c = contract(7);
for _ in 0..INCOMPAT_TRIP_THRESHOLD {
memo.note_delta_apply_failed(c);
}
ts.advance_time(DELTA_INCOMPAT_TTL + Duration::from_secs(1));
assert!(!memo.suppress_deltas(&c));
memo.note_delta_apply_failed(c);
assert!(
memo.suppress_deltas(&c),
"a failure after expiry must re-arm without re-counting from zero"
);
assert_eq!(memo.armed_total(), 2);
}
#[test]
fn maps_are_capacity_capped() {
let (_ts, memo) = setup();
for i in 0..(MAX_TRACKED_DELTA_SENDS + 100) {
let mut id = [0u8; 32];
id[..8].copy_from_slice(&(i as u64).to_be_bytes());
memo.record_delta_sent(ContractInstanceId::new(id), peer(5000));
}
assert!(memo.recent_delta_sends.len() <= MAX_TRACKED_DELTA_SENDS);
for i in 0..(MAX_TRACKED_CONTRACTS + 100) {
let mut id = [0u8; 32];
id[..8].copy_from_slice(&(i as u64).to_be_bytes());
memo.note_delta_apply_failed(ContractInstanceId::new(id));
}
assert!(memo.contracts.len() <= MAX_TRACKED_CONTRACTS);
}
#[test]
fn cleanup_sweeps_stale_but_keeps_armed() {
let (ts, memo) = setup();
let armed = contract(8);
for _ in 0..INCOMPAT_TRIP_THRESHOLD {
memo.note_delta_apply_failed(armed);
}
let idle = contract(9);
memo.note_delta_apply_failed(idle);
memo.record_delta_sent(contract(10), peer(6000));
ts.advance_time(DELTA_ATTRIBUTION_WINDOW + Duration::from_secs(1));
memo.cleanup();
assert!(
memo.recent_delta_sends.is_empty(),
"stale delta-send attributions must be swept"
);
assert!(memo.suppress_deltas(&armed), "armed entry must survive");
ts.advance_time(CONTRACT_CLEANUP_AGE);
memo.cleanup();
assert!(
!memo.contracts.contains_key(&idle),
"idle unarmed entry must be swept"
);
}
}