use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use freenet_stdlib::prelude::{ContractKey, WrappedState};
use crate::node::OpManager;
use crate::ring::PeerKeyLocation;
use crate::transport::BroadcastDeliveryOutcome;
use super::broadcast_payload_mix::PayloadArm;
use super::p2p_protoc::P2pBridge;
const STREAM_COMPLETION_TIMEOUT: Duration = Duration::from_secs(120);
const BROADCAST_QUEUE_PAYLOAD_SIZE_THRESHOLD: usize = 64 * 1024;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[cfg_attr(feature = "simulation_tests", allow(dead_code))]
pub(super) enum QueuedPayloadClass {
Small,
Large,
}
pub(crate) static BROADCAST_STREAM_METRICS: BroadcastStreamMetrics = BroadcastStreamMetrics::new();
pub(crate) struct BroadcastStreamMetrics {
streaming_attempts_total: AtomicU64,
streaming_failures_total: AtomicU64,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct BroadcastStreamMetricsSnapshot {
pub streaming_attempts_total: u64,
pub streaming_failures_total: u64,
}
impl BroadcastStreamMetrics {
const fn new() -> Self {
Self {
streaming_attempts_total: AtomicU64::new(0),
streaming_failures_total: AtomicU64::new(0),
}
}
fn record_attempt(&self, delivered: bool) {
self.streaming_attempts_total
.fetch_add(1, Ordering::Relaxed);
if !delivered {
self.streaming_failures_total
.fetch_add(1, Ordering::Relaxed);
}
}
pub(crate) fn snapshot(&self) -> BroadcastStreamMetricsSnapshot {
BroadcastStreamMetricsSnapshot {
streaming_attempts_total: self.streaming_attempts_total.load(Ordering::Relaxed),
streaming_failures_total: self.streaming_failures_total.load(Ordering::Relaxed),
}
}
}
pub(super) fn should_broadcast_contract(op_manager: &Arc<OpManager>, key: &ContractKey) -> bool {
op_manager.ring.should_summarize_or_broadcast(key)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum FanoutSendPlan {
Skip,
Send,
Probe,
}
#[derive(Clone, Copy)]
pub(super) struct SummaryPair<'a> {
pub ours: &'a freenet_stdlib::prelude::StateSummary<'static>,
pub theirs: &'a freenet_stdlib::prelude::StateSummary<'static>,
}
pub(super) fn plan_fanout_send<T: crate::util::time_source::TimeSource + Sync>(
interest_manager: &crate::ring::interest::InterestManager<T>,
key: &ContractKey,
summaries: SummaryPair<'_>,
probes_used: usize,
) -> FanoutSendPlan {
use crate::node::{StalenessProbeAction, plan_staleness_probe};
let SummaryPair { ours, theirs } = summaries;
if ours.as_ref() == theirs.as_ref() {
return FanoutSendPlan::Skip;
}
let cached = interest_manager.cached_staleness_verdict(key, theirs.as_ref(), ours.as_ref());
match plan_staleness_probe(cached, probes_used) {
StalenessProbeAction::UseCached(true) => FanoutSendPlan::Send,
StalenessProbeAction::UseCached(false) => FanoutSendPlan::Skip,
StalenessProbeAction::RunProbe => FanoutSendPlan::Probe,
StalenessProbeAction::BudgetExhaustedFallBack => FanoutSendPlan::Send,
}
}
pub(super) async fn fanout_send_needed(
op_manager: &OpManager,
key: &ContractKey,
summaries: SummaryPair<'_>,
probes_used: &mut usize,
) -> bool {
match plan_fanout_send(&op_manager.interest_manager, key, summaries, *probes_used) {
FanoutSendPlan::Send => true,
FanoutSendPlan::Skip => false,
FanoutSendPlan::Probe => {
*probes_used += 1;
let SummaryPair { ours, theirs } = summaries;
let verdict = op_manager
.interest_manager
.peer_summary_has_pending_state(op_manager, key, theirs, ours)
.await;
crate::ring::interest::summary_indicates_stale_peer(ours, theirs, verdict)
}
}
}
#[cfg(not(feature = "simulation_tests"))]
mod queue {
use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Instant;
use freenet_stdlib::prelude::{ContractKey, WrappedState};
use tokio::sync::{Mutex, Notify, Semaphore};
use crate::node::OpManager;
use crate::ring::PeerKeyLocation;
use super::super::p2p_protoc::P2pBridge;
use super::{QueuedPayloadClass, broadcast_to_single_peer};
const DEFAULT_SMALL_PAYLOAD_CONCURRENCY: usize = 12;
const DEFAULT_LARGE_PAYLOAD_CONCURRENCY: usize = 2;
const DEFAULT_MAX_QUEUE_DEPTH: usize = 256;
type DedupeKey = (ContractKey, PeerKeyLocation);
struct BroadcastEntry {
key: ContractKey,
target: PeerKeyLocation,
new_state: WrappedState,
payload_size: usize,
}
struct QueueState {
order: VecDeque<DedupeKey>,
entries: HashMap<DedupeKey, BroadcastEntry>,
active: HashMap<DedupeKey, u64>,
small_queued: usize,
hol_block: Option<HolBlockHandle>,
}
#[derive(Clone)]
struct HolBlockHandle(Arc<HolBlockState>);
struct HolBlockState {
finished: AtomicBool,
observation: std::sync::Mutex<HolBlockObservation>,
}
struct HolBlockObservation {
started: Instant,
last_updated: Instant,
small_queued: usize,
small_entry_millis: u128,
observed_small: bool,
}
impl HolBlockHandle {
fn is_finished(&self) -> bool {
self.0.finished.load(Ordering::Acquire)
}
fn set_small_queued(&self, next: usize, now: Instant) {
if self.is_finished() {
return;
}
let mut block = self.0.observation.lock().unwrap();
if self.is_finished() {
return;
}
block.advance(now);
block.small_queued = next;
block.observed_small |= next > 0;
}
fn finish_now(&self) -> Option<(u64, u64)> {
let mut block = self.0.observation.lock().unwrap();
if self.is_finished() {
return None;
}
let now = Instant::now();
self.finish_locked(&mut block, now)
}
#[cfg(test)]
fn finish_at(&self, now: Instant) -> Option<(u64, u64)> {
let mut block = self.0.observation.lock().unwrap();
if self.is_finished() {
return None;
}
self.finish_locked(&mut block, now)
}
fn finish_locked(
&self,
block: &mut HolBlockObservation,
now: Instant,
) -> Option<(u64, u64)> {
block.advance(now);
let result = block.observed_small.then(|| {
let blocked_millis = now
.saturating_duration_since(block.started)
.as_millis()
.min(u128::from(u64::MAX)) as u64;
let small_entry_millis = block.small_entry_millis.min(u128::from(u64::MAX)) as u64;
(blocked_millis, small_entry_millis)
});
self.0.finished.store(true, Ordering::Release);
result
}
}
impl HolBlockObservation {
fn advance(&mut self, now: Instant) {
if now <= self.last_updated {
return;
}
let millis = now.saturating_duration_since(self.last_updated).as_millis();
self.small_entry_millis = self
.small_entry_millis
.saturating_add(millis.saturating_mul(self.small_queued as u128));
self.last_updated = now;
}
}
impl QueueState {
fn new() -> Self {
Self {
order: VecDeque::new(),
entries: HashMap::new(),
active: HashMap::new(),
small_queued: 0,
hol_block: None,
}
}
fn len(&self) -> usize {
self.entries.len()
}
fn pop_front(&mut self, now: Instant) -> Option<BroadcastEntry> {
while let Some(key) = self.order.pop_front() {
if let Some(entry) = self.entries.remove(&key) {
if entry.payload_size < super::BROADCAST_QUEUE_PAYLOAD_SIZE_THRESHOLD {
self.set_small_queued(self.small_queued.saturating_sub(1), now);
}
return Some(entry);
}
}
None
}
fn set_small_queued(&mut self, next: usize, now: Instant) {
self.small_queued = next;
if self
.hol_block
.as_ref()
.is_some_and(HolBlockHandle::is_finished)
{
self.hol_block = None;
}
if let Some(block) = &self.hol_block {
block.set_small_queued(next, now);
}
}
fn start_hol(&mut self, now: Instant) -> HolBlockHandle {
let block = HolBlockHandle(Arc::new(HolBlockState {
finished: AtomicBool::new(false),
observation: std::sync::Mutex::new(HolBlockObservation {
started: now,
last_updated: now,
small_queued: self.small_queued,
small_entry_millis: 0,
observed_small: self.small_queued > 0,
}),
}));
self.hol_block = Some(block.clone());
block
}
fn track_active(&mut self, key: DedupeKey) -> bool {
if let Some(count) = self.active.get_mut(&key) {
*count = count.saturating_add(1);
return true;
}
if self.active.len() >= DEFAULT_MAX_QUEUE_DEPTH {
return false;
}
self.active.insert(key, 1);
true
}
fn untrack_active(&mut self, key: &DedupeKey) {
if let Some(count) = self.active.get_mut(key) {
if *count > 1 {
*count -= 1;
} else {
self.active.remove(key);
}
}
}
}
struct ActiveSendGuard {
queue: Arc<Mutex<QueueState>>,
key: DedupeKey,
tracked: bool,
}
impl ActiveSendGuard {
async fn finish(mut self) {
if self.tracked {
self.queue.lock().await.untrack_active(&self.key);
self.tracked = false;
}
}
}
impl Drop for ActiveSendGuard {
fn drop(&mut self) {
if !self.tracked {
return;
}
if let Ok(mut queue) = self.queue.try_lock() {
queue.untrack_active(&self.key);
return;
}
let queue = self.queue.clone();
let key = self.key.clone();
if let Ok(runtime) = tokio::runtime::Handle::try_current() {
runtime.spawn(async move {
queue.lock().await.untrack_active(&key);
});
}
}
}
#[derive(Clone)]
pub(crate) struct BroadcastQueue {
queue: Arc<Mutex<QueueState>>,
notify: Arc<Notify>,
small_payload_concurrency: usize,
large_payload_concurrency: usize,
max_queue_depth: usize,
}
impl BroadcastQueue {
pub(crate) fn new() -> Self {
Self {
queue: Arc::new(Mutex::new(QueueState::new())),
notify: Arc::new(Notify::new()),
small_payload_concurrency: DEFAULT_SMALL_PAYLOAD_CONCURRENCY,
large_payload_concurrency: DEFAULT_LARGE_PAYLOAD_CONCURRENCY,
max_queue_depth: DEFAULT_MAX_QUEUE_DEPTH,
}
}
pub(crate) async fn enqueue(
&self,
key: ContractKey,
target: PeerKeyLocation,
new_state: WrappedState,
) {
let dedup_key = (key, target.clone());
let mut queue = self.queue.lock().await;
if queue.active.contains_key(&dedup_key) {
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS.record_enqueue_while_active();
}
if let Some(existing) = queue.entries.get_mut(&dedup_key) {
existing.new_state = new_state;
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS.record_dedup_replacement();
tracing::trace!(
contract = %dedup_key.0,
peer = ?target.socket_addr(),
"Broadcast queue: replaced stale entry with newer state"
);
} else {
while queue.len() >= self.max_queue_depth {
if let Some(entry) = queue.pop_front(Instant::now()) {
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS.record_capacity_eviction();
tracing::warn!(
contract = %entry.key,
peer = ?entry.target.socket_addr(),
queue_depth = self.max_queue_depth,
"Broadcast queue full, evicted oldest entry"
);
} else {
break;
}
}
let payload_size = new_state.size();
if payload_size < super::BROADCAST_QUEUE_PAYLOAD_SIZE_THRESHOLD {
let next = queue.small_queued.saturating_add(1);
queue.set_small_queued(next, Instant::now());
}
queue.entries.insert(
dedup_key.clone(),
BroadcastEntry {
key,
target,
new_state,
payload_size,
},
);
queue.order.push_back(dedup_key);
}
crate::transport::shadow_demand::record_broadcast_queue_depth(queue.len());
drop(queue);
self.notify.notify_one();
}
pub(crate) fn start_worker(
&self,
bridge: P2pBridge,
op_manager: Arc<OpManager>,
) -> tokio::task::JoinHandle<()> {
let queue = self.queue.clone();
let notify = self.notify.clone();
let small_semaphore = Arc::new(Semaphore::new(self.small_payload_concurrency));
let large_semaphore = Arc::new(Semaphore::new(self.large_payload_concurrency));
tokio::spawn(async move {
loop {
let notified = notify.notified();
let mut drained_any = false;
loop {
let (entry, active_tracked) = {
let mut q = queue.lock().await;
let entry = q.pop_front(Instant::now());
let active_tracked = entry.as_ref().is_none_or(|entry| {
let tracked = q.track_active((entry.key, entry.target.clone()));
if !tracked {
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS
.record_active_tracking_overflow();
}
tracked
});
crate::transport::shadow_demand::record_broadcast_queue_depth(q.len());
(entry, active_tracked)
};
let Some(entry) = entry else {
break; };
drained_any = true;
let nominally_large =
entry.payload_size >= super::BROADCAST_QUEUE_PAYLOAD_SIZE_THRESHOLD;
let queued_class = if nominally_large {
QueuedPayloadClass::Large
} else {
QueuedPayloadClass::Small
};
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS
.record_scheduled(nominally_large, entry.payload_size);
let sem = if !nominally_large {
small_semaphore.clone()
} else {
large_semaphore.clone()
};
let blocked = nominally_large && sem.available_permits() == 0;
let hol_block = if blocked {
Some(queue.lock().await.start_hol(Instant::now()))
} else {
None
};
let active_guard = ActiveSendGuard {
queue: queue.clone(),
key: (entry.key, entry.target.clone()),
tracked: active_tracked,
};
let permit = sem.acquire_owned().await;
let hol_metrics = hol_block.as_ref().and_then(HolBlockHandle::finish_now);
let Ok(permit) = permit else {
if let Some((blocked_millis, small_entry_millis)) = hol_metrics {
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS
.record_large_head_block(blocked_millis, small_entry_millis);
}
active_guard.finish().await;
tracing::error!("Broadcast queue semaphore closed unexpectedly");
return;
};
let bridge = bridge.clone();
let op_manager = op_manager.clone();
tokio::spawn(async move {
broadcast_to_single_peer(
&bridge,
&op_manager,
entry.key,
entry.new_state,
entry.target,
Some(queued_class),
)
.await;
drop(permit);
active_guard.finish().await;
});
if let Some((blocked_millis, small_entry_millis)) = hol_metrics {
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS
.record_large_head_block(blocked_millis, small_entry_millis);
}
}
if !drained_any {
notified.await;
}
}
})
}
}
#[cfg(test)]
mod observation_tests {
use super::*;
use std::time::Duration;
#[test]
fn hol_integral_includes_small_entries_arriving_during_wait() {
let start = Instant::now();
let mut queue = QueueState::new();
let hol = queue.start_hol(start);
queue.set_small_queued(1, start + Duration::from_millis(10));
queue.set_small_queued(2, start + Duration::from_millis(20));
assert_eq!(
hol.finish_at(start + Duration::from_millis(30)),
Some((30, 30)),
"integral is 0×10ms + 1×10ms + 2×10ms"
);
queue.set_small_queued(9, start + Duration::from_millis(40));
assert!(hol.is_finished());
assert!(
queue.hol_block.is_none(),
"first later mutation clears the handle"
);
assert_eq!(hol.finish_at(start + Duration::from_millis(40)), None);
}
#[test]
fn active_refcount_survives_one_of_two_overlapping_completions() {
let mut queue = QueueState::new();
let code = freenet_stdlib::prelude::ContractCode::from(vec![7]);
let params = freenet_stdlib::prelude::Parameters::from(vec![9]);
let key = (
ContractKey::from_params_and_code(¶ms, &code),
PeerKeyLocation::random(),
);
assert!(queue.track_active(key.clone()));
assert!(queue.track_active(key.clone()));
queue.untrack_active(&key);
assert_eq!(queue.active.get(&key), Some(&1));
queue.untrack_active(&key);
assert!(!queue.active.contains_key(&key));
}
#[test]
fn worker_wires_hol_boundaries_and_direct_normal_cleanup() {
let src = include_str!("broadcast_queue.rs");
let start = src.find("pub(crate) fn start_worker(").unwrap();
let end = src[start..].find("} // end `mod queue`").unwrap() + start;
let worker = &src[start..end];
let hol_start = worker.find("start_hol(Instant::now())").unwrap();
let acquire = worker.find("sem.acquire_owned().await").unwrap();
assert!(
hol_start < acquire,
"HOL observation must start before waiting"
);
let freeze = worker.find("and_then(HolBlockHandle::finish_now)").unwrap();
let permit_match = worker.find("let Ok(permit) = permit").unwrap();
assert!(
acquire < freeze && freeze < permit_match,
"linearize the integral immediately after permit acquisition"
);
assert!(
!worker[acquire + "sem.acquire_owned().await".len()..freeze].contains(".await"),
"no further await may separate permit acquisition from HOL linearization"
);
let success = &worker[worker.find("let Ok(permit) = permit").unwrap()..];
let spawn = success.find("tokio::spawn(async move").unwrap();
let task = &success[spawn..];
assert!(
task.find("drop(permit);").unwrap()
< task.find("active_guard.finish().await").unwrap(),
"normal tracking cleanup must run directly after releasing capacity"
);
}
}
}
#[cfg(not(feature = "simulation_tests"))]
pub(crate) use queue::BroadcastQueue;
fn streaming_completion_delivered(completion: StreamCompletionResult) -> bool {
matches!(completion, Ok(Ok(BroadcastDeliveryOutcome::Delivered)))
}
type StreamCompletionResult = Result<
Result<BroadcastDeliveryOutcome, tokio::sync::oneshot::error::RecvError>,
tokio::time::error::Elapsed,
>;
#[allow(clippy::too_many_arguments)]
fn record_streaming_delivery<T: crate::util::time_source::TimeSource + Sync>(
interest_manager: &crate::ring::interest::InterestManager<T>,
completion: StreamCompletionResult,
sent_delta: bool,
key: &ContractKey,
peer_key: &crate::ring::PeerKey,
our_summary: Option<&freenet_stdlib::prelude::StateSummary<'static>>,
state_size: usize,
payload_size: usize,
) -> bool {
let delivered = streaming_completion_delivered(completion);
if delivered {
record_delivery_to_interest(
interest_manager,
sent_delta,
key,
peer_key,
our_summary,
state_size,
payload_size,
);
}
delivered
}
fn record_delivery_to_interest<T: crate::util::time_source::TimeSource + Sync>(
interest_manager: &crate::ring::interest::InterestManager<T>,
sent_delta: bool,
key: &ContractKey,
peer_key: &crate::ring::PeerKey,
our_summary: Option<&freenet_stdlib::prelude::StateSummary<'static>>,
state_size: usize,
payload_size: usize,
) {
if sent_delta {
interest_manager.record_delta_send(state_size, payload_size);
crate::config::GlobalTestMetrics::record_delta_send();
} else {
interest_manager.record_full_state_send();
crate::config::GlobalTestMetrics::record_full_state_send();
}
interest_manager.refresh_peer_interest(key, peer_key);
if let Some(summary) = our_summary {
interest_manager.upsert_peer_summary_from(
key,
peer_key,
summary.clone(),
crate::ring::interest::SummaryPopulationSource::Delivery,
);
}
}
pub(super) async fn broadcast_to_single_peer(
bridge: &P2pBridge,
op_manager: &Arc<OpManager>,
key: ContractKey,
new_state: WrappedState,
target: PeerKeyLocation,
queued_class: Option<QueuedPayloadClass>,
) {
use crate::message::{DeltaOrFullState, NetMessage};
use crate::node::network_bridge::NetworkBridge;
use crate::operations::update::{BroadcastStreamingPayload, UpdateMsg};
use crate::ring::PeerKey;
use crate::transport::peer_connection::StreamId;
let Some(peer_addr) = target.socket_addr() else {
return;
};
if !should_broadcast_contract(op_manager, &key) {
tracing::trace!(
contract = %key,
peer = %peer_addr,
"Skipping broadcast - contract not hosted or in use"
);
return;
}
let peer_key = PeerKey::from(target.pub_key().clone());
let cost_clock = op_manager.ring.time_source.clone();
let send_wasm_started = cost_clock.now();
let report_send_cpu = |op_manager: &Arc<OpManager>| {
use crate::topology::meter::ResourceType;
let elapsed_us = cost_clock
.now()
.saturating_duration_since(send_wasm_started)
.as_micros() as f64;
op_manager.ring.report_contract_resource_usage(
*key.id(),
ResourceType::ExecCpuMicros,
elapsed_us,
);
};
let our_summary = op_manager
.interest_manager
.get_contract_summary(op_manager, &key)
.await;
let (their_summary, tracked_missing_reason, missing_attempt) = if our_summary.is_some() {
match op_manager
.interest_manager
.begin_peer_summary_broadcast(&key, &peer_key)
{
crate::ring::interest::PeerSummaryForBroadcast::Known(summary) => {
(Some(summary), None, None)
}
crate::ring::interest::PeerSummaryForBroadcast::Missing { reason, attempt } => {
(None, reason, attempt)
}
}
} else {
(None, None, None)
};
let mut missing_attempt_guard = missing_attempt.map(|attempt| {
op_manager
.interest_manager
.missing_summary_attempt_guard(attempt)
});
if let (Some(ours), Some(theirs)) = (&our_summary, &their_summary) {
let mut staleness_probes_used = 0usize;
if !fanout_send_needed(
op_manager,
&key,
SummaryPair { ours, theirs },
&mut staleness_probes_used,
)
.await
{
tracing::trace!(
contract = %key,
peer = %peer_addr,
"Skipping broadcast - peer already has our state (byte-equal \
or logically converged summaries)"
);
op_manager
.interest_manager
.refresh_peer_interest(&key, &peer_key);
report_send_cpu(op_manager);
return;
}
}
let deltas_suppressed = op_manager.ring.delta_incompat.suppress_deltas(key.id());
if deltas_suppressed {
tracing::debug!(
contract = %key,
peer = %peer_addr,
event = "delta_suppressed_incompat",
"Contract is in delta-incompat backoff — sending full state instead of a delta"
);
}
let mut not_efficient_gate_inputs: Option<(usize, usize)> = None;
let (payload, sent_delta, payload_arm) = match (&our_summary, &their_summary) {
(Some(_), Some(_)) if deltas_suppressed => (
DeltaOrFullState::FullState(new_state.as_ref().to_vec()),
false,
PayloadArm::FullDeltaSuppressed,
),
(Some(ours), Some(theirs)) => {
match op_manager
.interest_manager
.compute_delta(op_manager, &key, theirs, ours, new_state.size())
.await
{
Ok(Some(delta)) => (
DeltaOrFullState::Delta(delta.as_ref().to_vec()),
true,
PayloadArm::Delta,
),
Ok(None) => {
tracing::trace!(
contract = %key,
peer = %peer_addr,
"Skipping broadcast - contract reported empty delta \
(peer converged)"
);
op_manager
.interest_manager
.refresh_peer_interest(&key, &peer_key);
report_send_cpu(op_manager);
return;
}
Err(err) => {
tracing::debug!(
contract = %key,
error = %err,
"Delta computation failed, falling back to full state"
);
let arm = match err {
crate::ring::interest::DeltaUnavailable::NotEfficient {
summary_size,
state_size,
} => {
not_efficient_gate_inputs = Some((summary_size, state_size));
PayloadArm::FullNotEfficient
}
crate::ring::interest::DeltaUnavailable::ComputeFailed(_) => {
PayloadArm::FullComputeFailed
}
};
(
DeltaOrFullState::FullState(new_state.as_ref().to_vec()),
false,
arm,
)
}
}
}
(None, _) => (
DeltaOrFullState::FullState(new_state.as_ref().to_vec()),
false,
PayloadArm::FullNoOurSummary,
),
(Some(_), None) => {
let arm = if tracked_missing_reason.is_some() {
PayloadArm::FullNoTheirSummaryTracked
} else {
PayloadArm::FullNoTheirSummaryUntracked
};
(
DeltaOrFullState::FullState(new_state.as_ref().to_vec()),
false,
arm,
)
}
};
let payload_size = payload.size();
match (
queued_class,
payload_size < BROADCAST_QUEUE_PAYLOAD_SIZE_THRESHOLD,
) {
(Some(QueuedPayloadClass::Large), true) => {
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS
.record_queued_large_actual_small(payload_size);
}
(Some(QueuedPayloadClass::Small), false) => {
crate::node::BROADCAST_QUEUE_EFFICIENCY_METRICS
.record_queued_small_actual_large(payload_size);
}
_ => {}
}
report_send_cpu(op_manager);
let update_tx = crate::message::Transaction::new::<crate::operations::update::UpdateMsg>();
let use_streaming = matches!(&payload, DeltaOrFullState::FullState(_))
&& crate::operations::should_use_streaming(op_manager.streaming_threshold, payload_size);
let send_result = if use_streaming {
let sender_summary_bytes = our_summary
.as_ref()
.map(|s| s.as_ref().to_vec())
.unwrap_or_default();
let state_bytes = match payload {
DeltaOrFullState::FullState(data) => data,
_ => unreachable!("checked above"),
};
let streaming_payload = BroadcastStreamingPayload {
state_bytes,
sender_summary_bytes,
};
let payload_bytes = match bincode::serialize(&streaming_payload) {
Ok(b) => b,
Err(e) => {
tracing::warn!(
tx = %update_tx,
error = %e,
"Failed to serialize BroadcastStreamingPayload, skipping"
);
return;
}
};
let sid = StreamId::next_operations();
tracing::debug!(
tx = %update_tx,
contract = %key,
peer = %peer_addr,
stream_id = %sid,
payload_size,
"Using streaming for BroadcastTo (via queue)"
);
let msg = UpdateMsg::BroadcastToStreaming {
id: update_tx,
stream_id: sid,
key,
total_size: payload_bytes.len() as u64,
};
let net_msg: NetMessage = msg.into();
let metadata = match bincode::serialize(&net_msg) {
Ok(bytes) => Some(bytes::Bytes::from(bytes)),
Err(e) => {
tracing::warn!(
?peer_addr,
error = %e,
"Failed to serialize BroadcastTo metadata for embedding"
);
None
}
};
let send_res = bridge.send(peer_addr, net_msg).await;
if send_res.is_err() {
BROADCAST_STREAM_METRICS.record_attempt(false);
} else {
let (completion_tx, completion_rx) = tokio::sync::oneshot::channel();
if let Err(err) = bridge
.send_stream_with_completion(
peer_addr,
sid,
bytes::Bytes::from(payload_bytes),
metadata,
Some(completion_tx),
None,
)
.await
{
BROADCAST_STREAM_METRICS.record_attempt(false);
tracing::warn!(
tx = %update_tx,
peer = %peer_addr,
error = %err,
"Failed to send broadcast stream data"
);
} else {
let completion =
tokio::time::timeout(STREAM_COMPLETION_TIMEOUT, completion_rx).await;
let delivered = record_streaming_delivery(
&op_manager.interest_manager,
completion,
sent_delta,
&key,
&peer_key,
our_summary.as_ref(),
new_state.size(),
payload_size,
);
BROADCAST_STREAM_METRICS.record_attempt(delivered);
if delivered {
if let Some(guard) = missing_attempt_guard.as_mut() {
guard.mark_delivered(payload_size);
}
op_manager.ring.report_contract_resource_usage(
*key.id(),
crate::topology::meter::ResourceType::BroadcastFanoutCost,
payload_size as f64,
);
op_manager.payload_mix.record_delivered(
payload_arm,
key.id(),
payload_size,
not_efficient_gate_inputs,
tracked_missing_reason,
);
tracing::debug!(
tx = %update_tx,
peer = %peer_addr,
"Broadcast stream completed successfully"
);
} else {
tracing::debug!(
tx = %update_tx,
peer = %peer_addr,
timeout_secs = STREAM_COMPLETION_TIMEOUT.as_secs(),
"Broadcast stream dropped or timed out before delivery \
(permit released, interest NOT refreshed)"
);
}
}
}
send_res
} else {
let msg = UpdateMsg::BroadcastTo {
id: update_tx,
key,
payload,
sender_summary_bytes: our_summary
.as_ref()
.map(|s| s.as_ref().to_vec())
.unwrap_or_default(),
};
let res = bridge.send(peer_addr, msg.into()).await;
if res.is_ok() {
if let Some(guard) = missing_attempt_guard.as_mut() {
guard.mark_delivered(payload_size);
}
op_manager.ring.report_contract_resource_usage(
*key.id(),
crate::topology::meter::ResourceType::BroadcastFanoutCost,
payload_size as f64,
);
op_manager.payload_mix.record_delivered(
payload_arm,
key.id(),
payload_size,
not_efficient_gate_inputs,
tracked_missing_reason,
);
if sent_delta {
op_manager
.ring
.delta_incompat
.record_delta_sent(*key.id(), peer_addr);
}
record_delivery_to_interest(
&op_manager.interest_manager,
sent_delta,
&key,
&peer_key,
our_summary.as_ref(),
new_state.size(),
payload_size,
);
}
res
};
if let Err(err) = &send_result {
tracing::warn!(
tx = %update_tx,
peer = %peer_addr,
error = %err,
"Failed to send state change broadcast (queued)"
);
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use freenet_stdlib::prelude::{
CodeHash, ContractInstanceId, ContractKey, StateDelta, StateSummary,
};
use crate::ring::PeerKey;
use crate::ring::interest::InterestManager;
use crate::transport::{BroadcastDeliveryOutcome, TransportKeypair};
use crate::util::time_source::SharedMockTimeSource;
use super::{
BroadcastStreamMetrics, FanoutSendPlan, SummaryPair, plan_fanout_send,
record_streaming_delivery, streaming_completion_delivered,
};
#[test]
fn broadcast_stream_metrics_counts_attempts_and_failures() {
let m = BroadcastStreamMetrics::new();
let s = m.snapshot();
assert_eq!(s.streaming_attempts_total, 0, "starts at zero");
assert_eq!(s.streaming_failures_total, 0, "starts at zero");
m.record_attempt(true);
let s = m.snapshot();
assert_eq!(s.streaming_attempts_total, 1);
assert_eq!(s.streaming_failures_total, 0, "delivered is not a failure");
m.record_attempt(false);
let s = m.snapshot();
assert_eq!(s.streaming_attempts_total, 2, "every attempt counts");
assert_eq!(s.streaming_failures_total, 1, "the drop is counted");
m.record_attempt(false);
m.record_attempt(true);
let s = m.snapshot();
assert_eq!(s.streaming_attempts_total, 4);
assert_eq!(s.streaming_failures_total, 2);
}
fn make_contract_key(seed: u8) -> ContractKey {
ContractKey::from_id_and_code(
ContractInstanceId::new([seed; 32]),
CodeHash::new([seed.wrapping_add(1); 32]),
)
}
fn make_peer_key() -> PeerKey {
PeerKey(TransportKeypair::new().public().clone())
}
async fn dropped_oneshot()
-> Result<BroadcastDeliveryOutcome, tokio::sync::oneshot::error::RecvError> {
let (tx, rx) = tokio::sync::oneshot::channel::<BroadcastDeliveryOutcome>();
drop(tx);
rx.await.map(|_| unreachable!("sender was dropped"))
}
async fn elapsed_timeout() -> tokio::time::error::Elapsed {
let (tx, rx) = tokio::sync::oneshot::channel::<BroadcastDeliveryOutcome>();
let res = tokio::time::timeout(Duration::from_millis(1), rx).await;
drop(tx);
res.expect_err("never-resolving recv must time out")
}
#[tokio::test]
async fn streaming_completion_delivered_only_on_explicit_delivery() {
assert!(
streaming_completion_delivered(Ok(Ok(BroadcastDeliveryOutcome::Delivered))),
"an explicit Delivered outcome must be treated as a delivery"
);
assert!(
!streaming_completion_delivered(Ok(Ok(BroadcastDeliveryOutcome::Dropped))),
"an explicit Dropped outcome must NOT be treated as a delivery (#4235)"
);
assert!(
!streaming_completion_delivered(Ok(dropped_oneshot().await)),
"a dropped completion oneshot must NOT be treated as a delivery (#4235)"
);
assert!(
!streaming_completion_delivered(Err(elapsed_timeout().await)),
"a completion-wait timeout must NOT be treated as a delivery (#4235)"
);
}
#[tokio::test]
async fn drop_outcome_does_not_refresh_interest_or_cache_summary() {
let our_summary = StateSummary::from(vec![9, 9, 9, 9]);
let dropped = dropped_oneshot().await;
let timed_out = elapsed_timeout().await;
let cases: Vec<(&str, super::StreamCompletionResult, bool)> = vec![
(
"delivered",
Ok(Ok(BroadcastDeliveryOutcome::Delivered)),
true,
),
(
"explicit-drop",
Ok(Ok(BroadcastDeliveryOutcome::Dropped)),
false,
),
("dropped-oneshot", Ok(dropped), false),
("timeout", Err(timed_out), false),
];
for (name, completion, expect_delivered) in cases {
let time_source = SharedMockTimeSource::new();
let manager = InterestManager::new(time_source.clone());
let contract = make_contract_key(7);
let peer = make_peer_key();
manager.register_peer_interest(&contract, peer.clone(), None, false);
let baseline = manager
.get_peer_interest(&contract, &peer)
.expect("peer interest registered")
.last_refreshed;
time_source.advance_time(Duration::from_secs(5));
let delivered = record_streaming_delivery(
&manager,
completion,
true,
&contract,
&peer,
Some(&our_summary),
1024,
64,
);
assert_eq!(
delivered, expect_delivered,
"[{name}] classification mismatch"
);
let interest = manager
.get_peer_interest(&contract, &peer)
.expect("peer interest still registered");
if expect_delivered {
assert!(
interest.last_refreshed > baseline,
"[{name}] a real delivery MUST refresh the peer interest TTL"
);
assert_eq!(
manager.get_peer_summary(&contract, &peer),
Some(our_summary.clone()),
"[{name}] a real delivery MUST cache the peer summary"
);
} else {
assert_eq!(
interest.last_refreshed, baseline,
"[{name}] a dropped/timed-out broadcast MUST NOT refresh the \
peer interest TTL (#4235)"
);
assert_eq!(
manager.get_peer_summary(&contract, &peer),
None,
"[{name}] a dropped/timed-out broadcast MUST NOT cache the peer \
summary, or the next summary-mismatch resend is suppressed (#4235)"
);
}
}
}
#[tokio::test]
async fn full_state_delivery_caches_summary_so_next_broadcast_is_delta() {
let our_summary = StateSummary::from(vec![1, 2, 3, 4]);
let time_source = SharedMockTimeSource::new();
let manager = InterestManager::new(time_source.clone());
let contract = make_contract_key(42);
let peer = make_peer_key();
manager.register_peer_interest(&contract, peer.clone(), None, false);
assert_eq!(
manager.get_peer_summary(&contract, &peer),
None,
"precondition: a brand-new subscriber has no cached summary, so the \
first broadcast must be full state"
);
let delivered = record_streaming_delivery(
&manager,
Ok(Ok(BroadcastDeliveryOutcome::Delivered)),
false,
&contract,
&peer,
Some(&our_summary),
4096,
4096,
);
assert!(delivered, "a Delivered outcome must classify as delivered");
assert_eq!(
manager.get_peer_summary(&contract, &peer),
Some(our_summary.clone()),
"#4145: a delivered FULL-STATE broadcast must cache the peer summary, \
so the next broadcast can be a delta — otherwise the peer is trapped \
sending full state forever (the #4233 storm)"
);
let their_summary = manager.get_peer_summary(&contract, &peer);
assert!(
their_summary.is_some(),
"#4145: with a cached peer summary the next broadcast takes the delta \
path (compute_delta), not another full state"
);
}
#[tokio::test]
async fn untracked_peer_delivery_seeds_interest_and_summary() {
let our_summary = StateSummary::from(vec![5, 6, 7, 8]);
let time_source = SharedMockTimeSource::new();
let manager = InterestManager::new(time_source.clone());
let contract = make_contract_key(11);
let peer = make_peer_key();
assert!(
manager.get_peer_interest(&contract, &peer).is_none(),
"precondition: the peer must be untracked"
);
let delivered = record_streaming_delivery(
&manager,
Ok(Ok(BroadcastDeliveryOutcome::Delivered)),
false,
&contract,
&peer,
Some(&our_summary),
4096,
4096,
);
assert!(delivered);
assert_eq!(
manager.get_peer_summary(&contract, &peer),
Some(our_summary),
"#4952: a delivered full-state broadcast to an UNTRACKED peer must \
seed the interest entry + summary, so the next broadcast is a \
delta — otherwise the pair is a full-state fixed point"
);
}
#[tokio::test]
async fn untracked_peer_drop_outcome_does_not_fabricate_interest() {
let our_summary = StateSummary::from(vec![3, 3, 3]);
let dropped = dropped_oneshot().await;
let timed_out = elapsed_timeout().await;
let cases: Vec<(&str, super::StreamCompletionResult)> = vec![
("explicit-drop", Ok(Ok(BroadcastDeliveryOutcome::Dropped))),
("dropped-oneshot", Ok(dropped)),
("timeout", Err(timed_out)),
];
for (name, completion) in cases {
let time_source = SharedMockTimeSource::new();
let manager = InterestManager::new(time_source.clone());
let contract = make_contract_key(12);
let peer = make_peer_key();
let delivered = record_streaming_delivery(
&manager,
completion,
false,
&contract,
&peer,
Some(&our_summary),
2048,
2048,
);
assert!(!delivered, "[{name}] must not classify as delivered");
assert!(
manager.get_peer_interest(&contract, &peer).is_none(),
"[{name}] a non-delivered send must NOT fabricate an interest \
entry for an untracked peer — a summary the peer never \
received would suppress the mismatch resend (#4235)"
);
}
}
#[tokio::test]
async fn dropped_full_state_stream_does_not_cache_summary() {
let our_summary = StateSummary::from(vec![5, 6, 7, 8]);
let dropped = dropped_oneshot().await;
let timed_out = elapsed_timeout().await;
let cases: Vec<(&str, super::StreamCompletionResult)> = vec![
("explicit-drop", Ok(Ok(BroadcastDeliveryOutcome::Dropped))),
("dropped-oneshot", Ok(dropped)),
("timeout", Err(timed_out)),
];
for (name, completion) in cases {
let time_source = SharedMockTimeSource::new();
let manager = InterestManager::new(time_source.clone());
let contract = make_contract_key(43);
let peer = make_peer_key();
manager.register_peer_interest(&contract, peer.clone(), None, false);
let delivered = record_streaming_delivery(
&manager,
completion,
false,
&contract,
&peer,
Some(&our_summary),
4096,
4096,
);
assert!(
!delivered,
"[{name}] a dropped/timed-out full-state stream must NOT classify \
as delivered"
);
assert_eq!(
manager.get_peer_summary(&contract, &peer),
None,
"[{name}] #4145 must not weaken the #2763/#4235 guard: a DROPPED \
full-state stream must NOT cache the summary (the peer never got \
the state), or the next summary-mismatch resend is suppressed"
);
}
}
#[test]
fn broadcast_single_peer_gates_summarize_on_hosted_or_in_use_pin() {
let src = include_str!("broadcast_queue.rs");
let helper_start = src
.find("pub(super) fn should_broadcast_contract(")
.expect("should_broadcast_contract helper not found");
let helper_end = helper_start
+ src[helper_start..]
.find("\n}\n")
.expect("should_broadcast_contract body end not found");
let helper_src = &src[helper_start..helper_end];
assert!(
helper_src.contains("should_summarize_or_broadcast"),
"should_broadcast_contract must delegate to the composed \
should_summarize_or_broadcast predicate (single source of truth, \
#4610), not re-inline a partial (is_hosting || in_use) gate that \
would re-admit phantom stateless contracts"
);
let fn_start = src
.find("pub(super) async fn broadcast_to_single_peer(")
.expect("broadcast_to_single_peer not found");
let fn_src = &src[fn_start..];
let gate_off = fn_src.find("should_broadcast_contract(op_manager").expect(
"broadcast_to_single_peer must call should_broadcast_contract — a bare \
get_contract_summary here reintroduces the #4473 storm",
);
let summarize_off = fn_src
.find("get_contract_summary(")
.expect("broadcast_to_single_peer get_contract_summary call not found");
assert!(
gate_off < summarize_off,
"broadcast_to_single_peer must gate on should_broadcast_contract BEFORE \
calling get_contract_summary (#4473) — otherwise the summarize storm \
fires for every phantom contract before the gate can skip it"
);
}
#[test]
fn record_delivery_routes_summary_cache_through_upsert() {
let src = include_str!("broadcast_queue.rs");
let fn_start = src
.find("fn record_delivery_to_interest<")
.expect("record_delivery_to_interest not found");
let after = &src[fn_start..];
let fn_end = after
.find("\npub(super) async fn broadcast_to_single_peer(")
.expect("end of record_delivery_to_interest not found");
let body: String = after[..fn_end].split_whitespace().collect();
assert!(
body.contains(
"interest_manager.upsert_peer_summary_from(key,peer_key,summary.clone(),"
),
"post-delivery cache must upsert (create-if-absent) the peer summary"
);
assert!(
!body.contains("interest_manager.update_peer_summary("),
"update_peer_summary silently no-ops for untracked peers — the \
#4952 fixed point. Use upsert_peer_summary here."
);
}
#[test]
fn broadcast_to_single_peer_gates_deltas_on_incompat_memo() {
let src = include_str!("broadcast_queue.rs");
let fn_start = src
.find("pub(super) async fn broadcast_to_single_peer(")
.expect("broadcast_to_single_peer not found");
let after = &src[fn_start..];
let fn_end = after
.find("\nmod tests {")
.or_else(|| after.find("\n#[cfg(test)]"))
.expect("end of broadcast_to_single_peer not found");
let body = &after[..fn_end];
let gate_pos = body
.find(".suppress_deltas(")
.expect("broadcast_to_single_peer must consult the delta-incompat memo");
let delta_pos = body
.find(".compute_delta(")
.expect("compute_delta call not found");
assert!(
gate_pos < delta_pos,
"the delta-incompat gate must be consulted BEFORE compute_delta \
(gate {gate_pos} < compute_delta {delta_pos}) — otherwise the \
doomed delta is still computed and sent"
);
assert!(
body.contains("if deltas_suppressed => ("),
"suppression must short-circuit the payload match to FullState"
);
assert!(
body.contains("(Some(_), Some(_)) if deltas_suppressed => ("),
"the suppression guard must be scoped to the both-summaries-present \
case — a wildcard guard mis-attributes missing-summary full states \
to FullDeltaSuppressed (#3335 payload-mix accuracy)"
);
let guard_arm = body
.find("if deltas_suppressed => (")
.expect("guard arm not found");
let compute_arm = body
.find("(Some(ours), Some(theirs)) => {")
.expect("compute_delta arm `(Some(ours), Some(theirs))` not found");
assert!(
guard_arm < compute_arm,
"the `_ if deltas_suppressed` guard arm must come BEFORE the \
`(Some(ours), Some(theirs))` compute_delta arm (guard {guard_arm} \
< compute {compute_arm}) — a suppressed delta-incapable contract \
must never reach compute_delta"
);
let record_pos = body
.find(".record_delta_sent(")
.expect("broadcast_to_single_peer must record delivered delta sends");
let sent_delta_gate = body
.find("if sent_delta {")
.expect("record_delta_sent must be gated on sent_delta");
assert!(
sent_delta_gate < record_pos,
"record_delta_sent must sit inside the `if sent_delta` gate \
(gate {sent_delta_gate} < record {record_pos})"
);
}
#[test]
fn broadcast_to_single_peer_records_attempt_on_every_streaming_exit_pin() {
let src = include_str!("broadcast_queue.rs");
let fn_start = src
.find("pub(super) async fn broadcast_to_single_peer(")
.expect("broadcast_to_single_peer not found");
let after = &src[fn_start..];
let fn_end = after
.find("\nmod tests {")
.or_else(|| after.find("\n#[cfg(test)]"))
.expect("end of broadcast_to_single_peer (start of tests module) not found");
let body = &after[..fn_end];
let record_calls = body.matches(".record_attempt(").count();
assert_eq!(
record_calls, 3,
"broadcast_to_single_peer's streaming branch must call record_attempt \
on all three exits (initial-send Err, dispatch Err, post-dispatch \
outcome) — got {record_calls}. A dropped early-exit record silently \
biases the v0.2.73 incident gauge LOW under congestion."
);
let failure_calls = body.matches(".record_attempt(false)").count();
assert_eq!(
failure_calls, 2,
"exactly the two early-failure exits must record record_attempt(false) \
(got {failure_calls}); the third exit records record_attempt(delivered)"
);
}
fn make_manager() -> InterestManager<SharedMockTimeSource> {
InterestManager::new(SharedMockTimeSource::new())
}
#[test]
fn nondeterministic_converged_summaries_skip_fanout_resend() {
let manager = make_manager();
let contract = make_contract_key(50);
let ours = StateSummary::from(vec![1u8, 2, 3]);
let theirs = StateSummary::from(vec![3u8, 2, 1]);
assert_ne!(
ours.as_ref(),
theirs.as_ref(),
"precondition: summaries differ byte-wise (the pre-fix byte-compare \
would NOT skip, and the delta path fell back to full state)"
);
manager.cache_delta(
&contract,
theirs.as_ref(),
ours.as_ref(),
StateDelta::from(Vec::<u8>::new()),
);
assert_eq!(
plan_fanout_send(
&manager,
&contract,
SummaryPair {
ours: &ours,
theirs: &theirs
},
0
),
FanoutSendPlan::Skip,
"a converged-but-byte-differing pair must be skipped by the fan-out \
(pre-fix: full state was re-sent on every fan-out — the heal storm)"
);
}
#[test]
fn genuinely_diverged_summaries_still_send() {
let manager = make_manager();
let contract = make_contract_key(51);
let ours = StateSummary::from(vec![9u8, 9, 9]);
let theirs = StateSummary::from(vec![1u8]);
manager.cache_delta(
&contract,
theirs.as_ref(),
ours.as_ref(),
StateDelta::from(vec![42u8]),
);
assert_eq!(
plan_fanout_send(
&manager,
&contract,
SummaryPair {
ours: &ours,
theirs: &theirs
},
0
),
FanoutSendPlan::Send,
"a genuine divergence (non-empty delta) must still be sent"
);
}
#[test]
fn byte_equal_summaries_skip_before_cache_lookup() {
let manager = make_manager();
let contract = make_contract_key(52);
let ours = StateSummary::from(vec![7u8, 7, 7]);
let theirs = StateSummary::from(vec![7u8, 7, 7]);
manager.cache_delta(
&contract,
theirs.as_ref(),
ours.as_ref(),
StateDelta::from(vec![1u8]),
);
assert_eq!(
plan_fanout_send(
&manager,
&contract,
SummaryPair {
ours: &ours,
theirs: &theirs
},
0
),
FanoutSendPlan::Skip,
"byte-identical summaries are trivially converged; the byte-equal \
short-circuit must precede any delta-cache verdict"
);
}
#[test]
fn probe_budget_gates_wasm_probe_and_falls_back_to_send() {
use crate::node::MAX_STALENESS_PROBES_PER_SUMMARIES;
let manager = make_manager();
let contract = make_contract_key(53);
let ours = StateSummary::from(vec![1u8, 2, 3]);
let theirs = StateSummary::from(vec![3u8, 2, 1]);
assert_eq!(
plan_fanout_send(
&manager,
&contract,
SummaryPair {
ours: &ours,
theirs: &theirs
},
0
),
FanoutSendPlan::Probe,
"a cache miss within budget must run the bounded WASM probe"
);
assert_eq!(
plan_fanout_send(
&manager,
&contract,
SummaryPair {
ours: &ours,
theirs: &theirs
},
MAX_STALENESS_PROBES_PER_SUMMARIES - 1
),
FanoutSendPlan::Probe,
"the last budget slot is still spendable"
);
assert_eq!(
plan_fanout_send(
&manager,
&contract,
SummaryPair {
ours: &ours,
theirs: &theirs
},
MAX_STALENESS_PROBES_PER_SUMMARIES
),
FanoutSendPlan::Send,
"an exhausted probe budget must fall back to the conservative \
byte-differ ⇒ send behavior, never a silent skip"
);
manager.cache_delta(
&contract,
theirs.as_ref(),
ours.as_ref(),
StateDelta::from(Vec::<u8>::new()),
);
assert_eq!(
plan_fanout_send(
&manager,
&contract,
SummaryPair {
ours: &ours,
theirs: &theirs
},
MAX_STALENESS_PROBES_PER_SUMMARIES * 10
),
FanoutSendPlan::Skip,
"cache hits never consume budget and still answer (converged ⇒ skip)"
);
}
#[test]
fn fanout_path_uses_semantic_delta_skip_pin() {
let src = include_str!("broadcast_queue.rs");
let fn_start = src
.find("pub(super) async fn broadcast_to_single_peer(")
.expect("broadcast_to_single_peer not found");
let after = &src[fn_start..];
let fn_end = after
.find("\nmod tests {")
.or_else(|| after.find("\n#[cfg(test)]"))
.expect("end of broadcast_to_single_peer (start of tests module) not found");
let body = &after[..fn_end];
assert!(
body.contains("fanout_send_needed("),
"broadcast_to_single_peer must route the per-peer skip decision \
through fanout_send_needed — a bare summary byte comparison \
re-opens the nondeterministic-summary heal storm"
);
let ok_none_off = body
.find("Ok(None) =>")
.expect("compute_delta Ok(None) arm not found in broadcast_to_single_peer");
let err_off = body[ok_none_off..]
.find("Err(err) =>")
.expect("compute_delta Err arm not found after Ok(None) arm");
let ok_none_arm = &body[ok_none_off..ok_none_off + err_off];
assert!(
!ok_none_arm.contains("FullState"),
"the Ok(None) (empty delta = converged) arm must NOT fall back to \
sending full state — that re-flood on every fan-out IS the heal \
storm. Arm body:\n{ok_none_arm}"
);
assert!(
ok_none_arm.contains("return;"),
"the Ok(None) (empty delta = converged) arm must skip the send \
entirely (return). Arm body:\n{ok_none_arm}"
);
let helpers_start = src
.find("pub(super) fn plan_fanout_send")
.expect("plan_fanout_send not found");
let helpers_end = src
.find("// The `BroadcastQueue` struct (constants, types, impl)")
.expect("queue module comment anchor not found");
assert!(
helpers_start < helpers_end,
"plan_fanout_send / fanout_send_needed must be defined before the \
queue module"
);
let helpers = &src[helpers_start..helpers_end];
assert!(
helpers.contains("plan_staleness_probe"),
"plan_fanout_send must ration WASM probes through \
plan_staleness_probe (the MAX_STALENESS_PROBES_PER_SUMMARIES cap)"
);
assert!(
helpers.contains("cached_staleness_verdict"),
"plan_fanout_send must consult the shared delta cache \
(cached_staleness_verdict) before trusting summary bytes"
);
assert!(
helpers.contains("peer_summary_has_pending_state"),
"fanout_send_needed must resolve cache misses via the bounded \
contract delta probe (peer_summary_has_pending_state)"
);
assert!(
helpers.contains("summary_indicates_stale_peer"),
"fanout_send_needed must decide from the probe verdict via \
summary_indicates_stale_peer (semantic policy), not inline byte \
inequality"
);
}
#[test]
fn broadcast_to_single_peer_reports_send_cost_pin() {
let src = include_str!("broadcast_queue.rs");
let fn_start = src
.find("pub(super) async fn broadcast_to_single_peer(")
.expect("broadcast_to_single_peer not found");
let after = &src[fn_start..];
let fn_end = after
.find("\nmod tests {")
.or_else(|| after.find("\n#[cfg(test)]"))
.expect("end of broadcast_to_single_peer (start of tests module) not found");
let body = &after[..fn_end];
let report_invocations = body.matches("report_send_cpu(op_manager)").count();
assert_eq!(
report_invocations, 3,
"broadcast_to_single_peer must invoke report_send_cpu at all three \
exits (summaries-equal skip + empty-delta converged skip + send \
attempt) — got {report_invocations}. A dropped report blinds \
cost-pressure eviction (#4861) to per-send CPU."
);
let cpu_needle = concat!("ResourceType::", "Exec", "CpuMicros");
let bytes_needle = concat!("ResourceType::", "Broadcast", "FanoutCost");
assert!(
body.contains(cpu_needle),
"report_send_cpu must attribute on the ExecCpuMicros axis"
);
assert_eq!(
body.matches(bytes_needle).count(),
2,
"fan-out bytes must be charged on the BroadcastFanoutCost axis at \
exactly the two real-delivery sites (streaming Delivered + inline \
send Ok) — never up-front (review round-3 Fix 4)"
);
assert_eq!(
body.matches("payload_size as f64").count(),
2,
"each delivery-gated bytes report must charge the selected \
payload_size (delta or full state), not the pre-delta full-state size"
);
}
#[test]
fn broadcast_to_single_peer_refreshes_interest_on_every_skip_pin() {
let src = include_str!("broadcast_queue.rs");
let fn_start = src
.find("pub(super) async fn broadcast_to_single_peer(")
.expect("broadcast_to_single_peer not found");
let after = &src[fn_start..];
let fn_end = after
.find("\nmod tests {")
.or_else(|| after.find("\n#[cfg(test)]"))
.expect("end of broadcast_to_single_peer (start of tests module) not found");
let body = &after[..fn_end];
let refreshes = body
.matches("refresh_peer_interest(&key,&peer_key)")
.count()
+ body
.matches("refresh_peer_interest(&key, &peer_key)")
.count();
assert_eq!(
refreshes, 2,
"broadcast_to_single_peer must refresh the interest TTL at BOTH \
converged-skip exits (summaries-equal and empty-delta) — got \
{refreshes}. Skipping the payload must not expire the interest."
);
let equal_skip = body
.find("Skipping broadcast - peer already has our state")
.expect("summaries-equal skip arm not found");
let empty_delta_skip = body
.find("Skipping broadcast - contract reported empty delta")
.expect("empty-delta skip arm not found");
assert!(
equal_skip < empty_delta_skip,
"unexpected arm order; the offsets below assume summaries-equal \
precedes empty-delta"
);
assert!(
body[equal_skip..empty_delta_skip].contains("refresh_peer_interest("),
"the summaries-equal skip must refresh the interest TTL before it \
returns"
);
assert!(
body[empty_delta_skip..].contains("refresh_peer_interest("),
"the empty-delta (converged) skip must refresh the interest TTL \
before it returns — same outcome as the skip above, same obligation"
);
assert!(
!body[equal_skip..empty_delta_skip].contains("record_delivery_to_interest("),
"the summaries-equal skip must refresh ONLY — recording a delivery \
that did not happen corrupts the delta/full-state telemetry"
);
}
#[test]
fn broadcast_to_single_peer_tags_every_payload_arm_pin() {
let src = include_str!("broadcast_queue.rs");
let fn_start = src
.find("pub(super) async fn broadcast_to_single_peer(")
.expect("broadcast_to_single_peer not found");
let after = &src[fn_start..];
let fn_end = after
.find("\nmod tests {")
.or_else(|| after.find("\n#[cfg(test)]"))
.expect("end of broadcast_to_single_peer (start of tests module) not found");
let body = &after[..fn_end];
for arm in [
"PayloadArm::Delta",
"PayloadArm::FullDeltaSuppressed",
"PayloadArm::FullNotEfficient",
"PayloadArm::FullComputeFailed",
"PayloadArm::FullNoOurSummary",
"PayloadArm::FullNoTheirSummaryUntracked",
"PayloadArm::FullNoTheirSummaryTracked",
] {
assert!(
body.contains(arm),
"broadcast_to_single_peer must tag the {arm} arm — an untagged \
fallback makes the #3335 payload-mix measurement attribute \
bytes to the wrong cause"
);
}
assert_eq!(
body.matches("begin_peer_summary_broadcast(&key, &peer_key)")
.count(),
1,
"broadcast payload selection must use the atomic missing-summary \
observation/classification operation exactly once"
);
let collapsed: String = body.chars().filter(|c| !c.is_whitespace()).collect();
let record_needle = concat!(
".record_",
"delivered(payload_arm,key.id(),payload_size,not_efficient_gate_inputs,tracked_missing_reason,)"
);
assert_eq!(
collapsed.matches(record_needle).count(),
2,
"payload mix must be recorded at exactly the two real-delivery \
sites (streaming Delivered + inline send Ok) — recording up-front \
would count dropped/failed sends as bytes on the wire"
);
assert_eq!(
collapsed
.matches(concat!("op_manager.payload_mix.record_", "delivered("))
.count(),
2,
"the payload mix must be recorded on op_manager.payload_mix (the \
per-node accumulator), not a process-global static"
);
assert!(
body.contains("DeltaUnavailable::NotEfficient")
&& body.contains("DeltaUnavailable::ComputeFailed"),
"the delta-failure arm must keep the typed NotEfficient vs \
ComputeFailed split — collapsing them re-blinds the measurement"
);
assert!(
collapsed.contains("(None,_)=>") && collapsed.contains("(Some(_),None)=>"),
"the no-summary arms must branch on WHICH side of the pair is \
missing — a catch-all `_` arm re-blinds the split"
);
assert!(
collapsed.contains(
"PeerSummaryForBroadcast::Missing{reason,attempt}=>{(None,reason,attempt)}"
),
"the atomic peer-summary result must carry the reason through so \
tracked and untracked missing-summary sends remain distinct"
);
}
}