use std::collections::HashSet;
use std::net::SocketAddr;
use std::time::{Duration, Instant};
use dashmap::DashMap;
use freenet_stdlib::prelude::ContractInstanceId;
use serde::{Deserialize, Serialize};
use crate::message::Transaction;
use crate::ring::PeerKey;
use crate::transport::TransportPublicKey;
pub(crate) type PeerHash = [u8; 8];
pub(crate) const MAX_COVERED_PEERS: usize = 64;
const COVERAGE_TTL: Duration = Duration::from_secs(10);
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct CoveredPeers {
hashes: Vec<PeerHash>,
}
impl CoveredPeers {
pub(crate) fn empty() -> Self {
Self::default()
}
pub(crate) fn from_targets<'a>(
tx: &Transaction,
targets: impl IntoIterator<Item = &'a TransportPublicKey>,
) -> Self {
let mut hashes: Vec<PeerHash> = targets
.into_iter()
.map(|pub_key| peer_hash(tx, pub_key))
.collect();
hashes.sort_unstable();
hashes.dedup();
hashes.truncate(MAX_COVERED_PEERS);
Self { hashes }
}
pub(crate) fn len(&self) -> usize {
self.hashes.len()
}
pub(crate) fn is_empty(&self) -> bool {
self.hashes.is_empty()
}
pub(crate) fn resolve<'a>(
&self,
tx: &Transaction,
candidates: impl IntoIterator<Item = &'a TransportPublicKey>,
) -> HashSet<PeerKey> {
if self.hashes.is_empty() {
return HashSet::new();
}
let named: HashSet<PeerHash> = self
.hashes
.iter()
.take(MAX_COVERED_PEERS)
.copied()
.collect();
candidates
.into_iter()
.filter(|pub_key| named.contains(&peer_hash(tx, pub_key)))
.map(|pub_key| PeerKey::from(pub_key.clone()))
.collect()
}
}
fn peer_hash(tx: &Transaction, pub_key: &TransportPublicKey) -> PeerHash {
let mut hasher = blake3::Hasher::new();
hasher.update(&tx.id_bytes());
hasher.update(pub_key.as_bytes());
let digest = hasher.finalize();
let mut out = [0u8; 8];
out.copy_from_slice(&digest.as_bytes()[..8]);
out
}
pub(crate) struct BroadcastCoverageStore {
entries: DashMap<ContractInstanceId, CoverageEntry>,
}
const SWEEP_THRESHOLD: usize = 512;
struct CoverageEntry {
origin: BroadcastOrigin,
expires_at: Instant,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub(crate) struct BroadcastOrigin {
sender: Option<SocketAddr>,
covered: HashSet<PeerKey>,
}
impl BroadcastOrigin {
pub(crate) fn local() -> Self {
Self::default()
}
pub(crate) fn relayed(sender: SocketAddr, covered: HashSet<PeerKey>) -> Self {
Self {
sender: Some(sender),
covered,
}
}
pub(crate) fn relayed_list_only(covered: HashSet<PeerKey>) -> Self {
Self {
sender: None,
covered,
}
}
pub(crate) fn sender(&self) -> Option<&SocketAddr> {
self.sender.as_ref()
}
pub(crate) fn covers(&self, pub_key: &TransportPublicKey) -> bool {
!self.covered.is_empty() && self.covered.contains(&PeerKey::from(pub_key.clone()))
}
pub(crate) fn covered_len(&self) -> usize {
self.covered.len()
}
fn narrow(&mut self, other: &BroadcastOrigin) {
if self.sender != other.sender {
self.sender = None;
}
self.covered.retain(|peer| other.covered.contains(peer));
}
}
impl BroadcastCoverageStore {
pub(crate) fn new() -> Self {
Self {
entries: DashMap::new(),
}
}
pub(crate) fn register(&self, key: &ContractInstanceId, origin: BroadcastOrigin) {
let now = Instant::now();
let expires_at = now + COVERAGE_TTL;
let mut inserted_new_contract = false;
match self.entries.entry(*key) {
dashmap::mapref::entry::Entry::Occupied(mut occupied) => {
let existing = occupied.get_mut();
if existing.expires_at <= now {
existing.origin = origin;
existing.expires_at = expires_at;
} else {
existing.origin.narrow(&origin);
existing.expires_at = existing.expires_at.min(expires_at);
}
}
dashmap::mapref::entry::Entry::Vacant(vacant) => {
vacant.insert(CoverageEntry { origin, expires_at });
inserted_new_contract = true;
}
}
if inserted_new_contract {
self.sweep_expired(now);
}
}
fn sweep_expired(&self, now: Instant) {
if self.entries.len() <= SWEEP_THRESHOLD {
return;
}
self.entries.retain(|_, entry| entry.expires_at > now);
}
pub(crate) fn take(&self, key: &ContractInstanceId) -> BroadcastOrigin {
match self.entries.remove(key) {
Some((_, entry)) if entry.expires_at > Instant::now() => entry.origin,
_ => BroadcastOrigin::local(),
}
}
pub(crate) fn discard(&self, key: &ContractInstanceId) {
self.entries.remove(key);
}
#[cfg(test)]
pub(crate) fn expire_all_for_test(&self) {
let past = Instant::now() - COVERAGE_TTL - Duration::from_secs(1);
for mut entry in self.entries.iter_mut() {
entry.expires_at = past;
}
}
#[cfg(test)]
pub(crate) fn live_entries(&self) -> usize {
let now = Instant::now();
self.entries
.iter()
.filter(|entry| entry.expires_at > now)
.count()
}
#[cfg(test)]
pub(crate) fn resident_entries(&self) -> usize {
self.entries.len()
}
#[cfg(test)]
pub(crate) fn insert_with_deadline(
&self,
key: &ContractInstanceId,
origin: BroadcastOrigin,
expires_at: Instant,
) {
self.entries
.insert(*key, CoverageEntry { origin, expires_at });
}
}
impl Default for BroadcastCoverageStore {
fn default() -> Self {
Self::new()
}
}
pub(crate) struct CoverageRegistration<'a> {
store: &'a BroadcastCoverageStore,
key: ContractInstanceId,
keep: bool,
}
impl<'a> CoverageRegistration<'a> {
pub(crate) fn new(
store: &'a BroadcastCoverageStore,
key: ContractInstanceId,
origin: BroadcastOrigin,
) -> Self {
store.register(&key, origin);
Self {
store,
key,
keep: false,
}
}
pub(crate) fn keep(mut self) {
self.keep = true;
}
}
impl Drop for CoverageRegistration<'_> {
fn drop(&mut self) {
if !self.keep {
self.store.discard(&self.key);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::transport::TransportKeypair;
fn pub_key() -> TransportPublicKey {
TransportKeypair::new().public().clone()
}
fn tx() -> Transaction {
Transaction::new::<crate::operations::update::UpdateMsg>()
}
fn relayed(covered: HashSet<PeerKey>) -> BroadcastOrigin {
BroadcastOrigin::relayed("127.0.0.1:1".parse().unwrap(), covered)
}
fn instance_id() -> ContractInstanceId {
*crate::operations::test_utils::make_contract_key(1).id()
}
#[test]
fn replacing_a_stale_entry_yields_a_live_claim() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let stale_peer = pub_key();
let fresh_peer = pub_key();
store.insert_with_deadline(
&key,
relayed(HashSet::from([PeerKey::from(stale_peer.clone())])),
Instant::now() - Duration::from_secs(1),
);
store.register(
&key,
relayed(HashSet::from([PeerKey::from(fresh_peer.clone())])),
);
let taken = store.take(&key);
assert!(
taken.covers(&fresh_peer),
"the replacing apply's coverage must survive — if this is empty the \
entry was born expired and `take` handed back a local() claim, so \
the stale-replacement branch is inert"
);
assert!(
!taken.covers(&stale_peer),
"the orphaned claim must be REPLACED, not merged: no live apply \
stands behind it"
);
}
#[test]
fn an_expired_entry_suppresses_nothing() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let peer = pub_key();
store.insert_with_deadline(
&key,
relayed(HashSet::from([PeerKey::from(peer.clone())])),
Instant::now() - Duration::from_millis(1),
);
let taken = store.take(&key);
assert!(
!taken.covers(&peer),
"a claim past COVERAGE_TTL must be ignored; mutation: delete the \
`expires_at > Instant::now()` guard in `take`"
);
}
#[test]
fn expired_entries_are_reclaimed_once_the_map_grows() {
let store = BroadcastCoverageStore::new();
let past = Instant::now() - Duration::from_secs(1);
for i in 0..(SWEEP_THRESHOLD + 2) {
let mut bytes = [0u8; 32];
bytes[0..8].copy_from_slice(&(i as u64).to_le_bytes());
store.insert_with_deadline(
&ContractInstanceId::new(bytes),
relayed(HashSet::from([PeerKey::from(pub_key())])),
past,
);
}
assert!(
store.resident_entries() > SWEEP_THRESHOLD,
"premise: the map must actually be over the sweep threshold"
);
assert_eq!(
store.live_entries(),
0,
"premise: every seeded entry is expired"
);
let fresh_key = instance_id();
store.register(
&fresh_key,
relayed(HashSet::from([PeerKey::from(pub_key())])),
);
assert_eq!(
store.resident_entries(),
1,
"expired entries must be reclaimed, leaving only the live one; \
found {} resident",
store.resident_entries()
);
assert!(
store.take(&fresh_key).covered_len() > 0,
"the sweep must not evict the live entry it ran alongside"
);
}
#[test]
fn a_named_peer_resolves_and_an_unnamed_one_does_not() {
let tx = tx();
let named = pub_key();
let unnamed = pub_key();
let covered = CoveredPeers::from_targets(&tx, [&named]);
let resolved = covered.resolve(&tx, [&named, &unnamed]);
assert!(resolved.contains(&PeerKey::from(named)));
assert!(!resolved.contains(&PeerKey::from(unnamed)));
assert_eq!(resolved.len(), 1);
}
#[test]
fn the_same_peer_hashes_differently_under_a_different_transaction() {
let peer = pub_key();
let first = CoveredPeers::from_targets(&tx(), [&peer]);
let second = CoveredPeers::from_targets(&tx(), [&peer]);
assert_ne!(first, second);
}
#[test]
fn a_list_does_not_resolve_under_a_different_transaction() {
let peer = pub_key();
let covered = CoveredPeers::from_targets(&tx(), [&peer]);
assert!(covered.resolve(&tx(), [&peer]).is_empty());
}
#[test]
fn resolve_bounds_an_oversized_inbound_list() {
let tx = tx();
let peers: Vec<TransportPublicKey> =
(0..(MAX_COVERED_PEERS + 12)).map(|_| pub_key()).collect();
let mut hashes: Vec<PeerHash> = peers.iter().map(|pk| peer_hash(&tx, pk)).collect();
hashes.sort_unstable();
hashes.dedup();
let oversized = CoveredPeers { hashes };
assert!(
oversized.len() > MAX_COVERED_PEERS,
"premise: the fixture must actually exceed the cap, or this test \
passes against any implementation"
);
let resolved = oversized.resolve(&tx, peers.iter());
assert!(
resolved.len() <= MAX_COVERED_PEERS,
"an inbound list of {} names resolved {} peers; the cap must hold \
for a list we did not build",
oversized.len(),
resolved.len()
);
}
#[test]
fn covered_peers_over_cap_truncates_rather_than_omitting() {
let tx = tx();
let peers: Vec<TransportPublicKey> =
(0..MAX_COVERED_PEERS * 2).map(|_| pub_key()).collect();
let covered = CoveredPeers::from_targets(&tx, peers.iter());
assert_eq!(
covered.len(),
MAX_COVERED_PEERS,
"over-cap must truncate to the cap"
);
assert!(
!covered.is_empty(),
"over-cap must NEVER degrade to an empty/omitted list — that is \
indistinguishable from a peer that does not implement this"
);
let resolved = covered.resolve(&tx, peers.iter());
assert_eq!(resolved.len(), MAX_COVERED_PEERS);
let strangers: Vec<TransportPublicKey> = (0..8).map(|_| pub_key()).collect();
let resolved_strangers = covered.resolve(&tx, strangers.iter());
assert!(
resolved_strangers.is_empty(),
"a truncated list resolved against peers it never named suppressed \
{} of them; truncation may only LOSE suppression, never invent it",
resolved_strangers.len()
);
}
#[test]
fn encoding_is_order_independent() {
let tx = tx();
let peers: Vec<TransportPublicKey> = (0..8).map(|_| pub_key()).collect();
let forward = CoveredPeers::from_targets(&tx, peers.iter());
let backward = CoveredPeers::from_targets(&tx, peers.iter().rev());
assert_eq!(forward, backward);
}
#[test]
fn a_registered_entry_is_taken_by_the_fan_out() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let peer = PeerKey::from(pub_key());
store.register(&key, relayed(HashSet::from([peer.clone()])));
assert_eq!(store.take(&key), relayed(HashSet::from([peer])));
assert_eq!(
store.take(&key),
BroadcastOrigin::local(),
"coverage is consumed once"
);
}
#[test]
fn concurrent_registrations_intersect_and_never_over_suppress() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let shared = PeerKey::from(pub_key());
let only_first = PeerKey::from(pub_key());
let only_second = PeerKey::from(pub_key());
store.register(
&key,
relayed(HashSet::from([shared.clone(), only_first.clone()])),
);
store.register(
&key,
relayed(HashSet::from([shared.clone(), only_second.clone()])),
);
let taken = store.take(&key);
assert_eq!(taken, relayed(HashSet::from([shared])));
assert!(!taken.covers(&only_first.0));
assert!(!taken.covers(&only_second.0));
}
#[test]
fn an_empty_registration_collapses_concurrent_coverage() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
store.register(&key, relayed(HashSet::from([PeerKey::from(pub_key())])));
store.register(&key, BroadcastOrigin::local());
assert_eq!(store.take(&key).covered_len(), 0);
}
#[test]
fn a_no_change_apply_discards_its_registration() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let peer = PeerKey::from(pub_key());
{
let _registration =
CoverageRegistration::new(&store, key, relayed(HashSet::from([peer])));
assert_eq!(store.live_entries(), 1);
}
assert_eq!(
store.live_entries(),
0,
"an apply that does not change state must not leave coverage \
behind for an unrelated fan-out to consume"
);
assert_eq!(store.take(&key).covered_len(), 0);
}
#[test]
fn a_changed_apply_keeps_its_registration_for_the_fan_out() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let peer = PeerKey::from(pub_key());
{
let registration =
CoverageRegistration::new(&store, key, relayed(HashSet::from([peer.clone()])));
registration.keep();
}
assert_eq!(store.take(&key), relayed(HashSet::from([peer])));
}
#[test]
fn replacing_an_expired_entry_leaves_the_new_claim_usable() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let orphaned = PeerKey::from(pub_key());
let live = PeerKey::from(pub_key());
store.register(&key, relayed(HashSet::from([orphaned.clone()])));
store.expire_all_for_test();
store.register(&key, relayed(HashSet::from([live.clone()])));
let taken = store.take(&key);
assert!(
taken.covers(&live.0),
"the replacement claim must be usable; it was written after the \
stale one expired and belongs to a live apply"
);
assert!(
!taken.covers(&orphaned.0),
"the orphaned claim must not survive its replacement"
);
}
#[test]
fn narrow_drops_a_disagreeing_sender() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let shared = PeerKey::from(pub_key());
let first: SocketAddr = "127.0.0.1:1".parse().unwrap();
let second: SocketAddr = "127.0.0.1:2".parse().unwrap();
store.register(
&key,
BroadcastOrigin::relayed(first, HashSet::from([shared.clone()])),
);
store.register(
&key,
BroadcastOrigin::relayed(second, HashSet::from([shared.clone()])),
);
let taken = store.take(&key);
assert_eq!(
taken.sender(),
None,
"with two candidate senders we cannot tell which delivery this \
fan-out belongs to, so neither may be excluded"
);
assert!(
taken.covers(&shared.0),
"the agreed coverage survives; only the sender is dropped"
);
}
#[test]
fn narrow_keeps_an_agreeing_sender() {
let store = BroadcastCoverageStore::new();
let key = instance_id();
let sender: SocketAddr = "127.0.0.1:1".parse().unwrap();
store.register(&key, BroadcastOrigin::relayed(sender, HashSet::new()));
store.register(&key, BroadcastOrigin::relayed(sender, HashSet::new()));
assert_eq!(store.take(&key).sender(), Some(&sender));
}
#[test]
fn covered_peers_round_trips_over_the_wire_with_content() {
let tx = tx();
let peers: Vec<TransportPublicKey> = (0..17).map(|_| pub_key()).collect();
let covered = CoveredPeers::from_targets(&tx, peers.iter());
let bytes = bincode::serialize(&covered).expect("serialize");
let decoded: CoveredPeers = bincode::deserialize(&bytes).expect("deserialize");
assert_eq!(decoded, covered);
assert_eq!(decoded.resolve(&tx, peers.iter()).len(), peers.len());
assert_eq!(
bytes.len(),
8 + 17 * 8,
"17 eight-byte hashes plus bincode's u64 length prefix; if this \
grows, the cost argument against the retired Bloom design moves \
with it"
);
}
#[test]
fn coverage_for_one_contract_does_not_leak_into_another() {
let store = BroadcastCoverageStore::new();
let first = instance_id();
let second = *crate::operations::test_utils::make_contract_key(2).id();
assert_ne!(first, second);
store.register(&first, relayed(HashSet::from([PeerKey::from(pub_key())])));
assert_eq!(store.take(&second), BroadcastOrigin::local());
}
}