use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use tokio::time::Instant;
use either::Either;
use freenet_stdlib::client_api::{ContractResponse, ErrorKind, HostResponse};
use freenet_stdlib::prelude::*;
use crate::client_events::HostResult;
use crate::config::{GlobalExecutor, GlobalRng};
use crate::contract::{ContractHandlerEvent, StoreResponse};
use crate::message::{DeltaOrFullState, InterestMessage, NetMessage, NodeEvent, Transaction};
use crate::node::OpManager;
use crate::operations::OpError;
use crate::operations::bootstrap::bootstrap_gateway_target;
use crate::ring::interest::RESYNC_REQUEST_MIN_INTERVAL;
use crate::ring::{PeerKeyLocation, RingError};
use crate::tracing::{NetEventLog, OperationFailure, state_hash_full};
use super::{
AutoFetchReason, BroadcastStreamingPayload, UpdateExecution, UpdateMsg, UpdateStreamingPayload,
};
use crate::transport::peer_connection::StreamId;
#[cfg(any(test, feature = "testing"))]
pub static RELAY_UPDATE_DRIVER_CALL_COUNT: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub static RELAY_UPDATE_INFLIGHT: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub static RELAY_UPDATE_SPAWNED_TOTAL: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub static RELAY_UPDATE_COMPLETED_TOTAL: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub static RELAY_UPDATE_DEDUP_REJECTS: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
#[cfg(any(test, feature = "testing"))]
pub static RELAY_UPDATE_STREAMING_DRIVER_CALL_COUNT: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub static RELAY_UPDATE_STREAMING_INFLIGHT: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub static RELAY_UPDATE_STREAMING_SPAWNED_TOTAL: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub static RELAY_UPDATE_STREAMING_COMPLETED_TOTAL: std::sync::atomic::AtomicUsize =
std::sync::atomic::AtomicUsize::new(0);
pub(crate) async fn start_client_update(
op_manager: Arc<OpManager>,
client_tx: Transaction,
key: ContractKey,
update_data: UpdateData<'static>,
related_contracts: RelatedContracts<'static>,
) -> Result<Transaction, OpError> {
crate::operations::reject_if_contract_banned(&op_manager, key.id())?;
tracing::debug!(
tx = %client_tx,
contract = %key,
"update: spawning client-initiated task"
);
let inflight_guard = match op_manager.admit_client_op() {
Some(g) => g,
None => return Err(OpError::NodeShuttingDown),
};
GlobalExecutor::spawn(async move {
let _inflight_guard = inflight_guard;
run_client_update(op_manager, client_tx, key, update_data, related_contracts).await;
});
Ok(client_tx)
}
async fn run_client_update(
op_manager: Arc<OpManager>,
client_tx: Transaction,
key: ContractKey,
update_data: UpdateData<'static>,
related_contracts: RelatedContracts<'static>,
) {
let outcome = drive_client_update(&op_manager, client_tx, key, update_data, related_contracts)
.await
.unwrap_or_else(DriverOutcome::InfrastructureError);
deliver_outcome(&op_manager, client_tx, outcome);
}
#[derive(Debug)]
enum DriverOutcome {
Publish(HostResult),
InfrastructureError(OpError),
}
async fn drive_client_update(
op_manager: &OpManager,
client_tx: Transaction,
key: ContractKey,
update_data: UpdateData<'static>,
related_contracts: RelatedContracts<'static>,
) -> Result<DriverOutcome, OpError> {
let sender_addr = op_manager.ring.connection_manager.peer_addr()?;
let is_delta_update = matches!(
&update_data,
UpdateData::Delta(_) | UpdateData::RelatedDelta { .. }
);
let proximity_neighbors: Vec<_> = op_manager.neighbor_hosting.neighbors_with_contract(&key);
let mut target_from_proximity: Option<PeerKeyLocation> = None;
for pub_key in &proximity_neighbors {
match op_manager
.ring
.connection_manager
.get_peer_by_pub_key(pub_key)
{
Some(peer) => {
if peer
.socket_addr()
.map(|a| a == sender_addr)
.unwrap_or(false)
{
continue;
}
target_from_proximity = Some(peer);
break;
}
None => {
tracing::debug!(
%key,
peer = %pub_key,
"update: proximity cache neighbor not connected, trying next"
);
}
}
}
let target = if let Some(proximity_neighbor) = target_from_proximity {
tracing::debug!(
%key,
target = ?proximity_neighbor.socket_addr(),
proximity_neighbors_found = proximity_neighbors.len(),
"update: using proximity cache neighbor as target"
);
Some(proximity_neighbor)
} else {
op_manager
.ring
.closest_potentially_hosting(&key, [sender_addr].as_slice())
.or_else(|| {
bootstrap_gateway_target(op_manager, |addr| addr == sender_addr).map(
|(gateway, gateway_addr)| {
tracing::info!(
%key,
gateway = %gateway_addr,
"update: ring empty — target falls back to configured gateway"
);
gateway
},
)
})
};
match target {
None => {
tracing::debug!(
tx = %client_tx,
%key,
"update: no remote peers, handling locally"
);
let is_hosting = op_manager.ring.is_hosting_contract(&key);
if !is_hosting {
tracing::error!(
contract = %key,
phase = "error",
"update: cannot update contract on isolated node — not hosted"
);
return Err(OpError::RingError(RingError::NoHostingPeers(*key.id())));
}
let UpdateExecution {
value: _,
summary,
changed,
..
} = match super::update_contract(
op_manager,
key,
update_data,
related_contracts,
crate::contract::Priority::ClientLocal,
)
.await
{
Ok(execution) => {
if is_delta_update {
op_manager
.ring
.merge_backoff
.record_success_local(key.id(), execution.changed);
} else if execution.changed {
op_manager
.ring
.merge_backoff
.invalidate_payload_memo(key.id());
}
execution
}
Err(err) if err.is_missing_contract_parameters() => {
tracing::error!(
tx = %client_tx,
contract = %key,
error = %err,
phase = "ring_state_store_inconsistency",
"update: is_hosting_contract reports true \
but state_store has no params for this contract; cannot \
auto-fetch (no remote target) — the ring/state-store \
divergence must be repaired by a re-PUT or a manual \
seeding step (#4066)"
);
return Ok(DriverOutcome::Publish(Err(ErrorKind::OperationError {
cause: format!(
"originator hosts {key} per ring but local state_store \
is missing contract code/params; no remote peer \
available to auto-fetch from"
)
.into(),
}
.into())));
}
Err(err) => return Err(err),
};
if !changed {
tracing::debug!(
tx = %client_tx,
%key,
"update: local update resulted in no change"
);
} else {
tracing::debug!(
tx = %client_tx,
%key,
"update: local-only update complete"
);
}
let host_result: HostResult = Ok(HostResponse::ContractResponse(
ContractResponse::UpdateResponse {
key,
summary: summary.clone(),
},
));
Ok(DriverOutcome::Publish(host_result))
}
Some(target) => {
let target_addr = match target.socket_addr() {
Some(addr) => addr,
None => {
tracing::error!(
tx = %client_tx,
%key,
target_pub_key = %target.pub_key(),
"update: target peer has no socket address"
);
return Err(OpError::RingError(RingError::NoHostingPeers(*key.id())));
}
};
tracing::debug!(
tx = %client_tx,
%key,
target_peer = %target_addr,
"update: applying locally before forwarding"
);
let local_apply = super::update_contract(
op_manager,
key,
update_data.clone(),
related_contracts.clone(),
crate::contract::Priority::ClientLocal,
)
.await;
let UpdateExecution {
value: updated_value,
summary,
changed: _,
..
} = match local_apply {
Ok(execution) => {
if is_delta_update {
op_manager
.ring
.merge_backoff
.record_success_local(key.id(), execution.changed);
} else if execution.changed {
op_manager
.ring
.merge_backoff
.invalidate_payload_memo(key.id());
}
execution
}
Err(err) => {
if err.is_missing_contract_parameters() {
tracing::warn!(
tx = %client_tx,
contract = %key,
error = %err,
target = %target_addr,
phase = "auto_fetch_originator",
"update: originator has no contract \
code/params for this contract; triggering \
auto-fetch from target and asking client to retry"
);
op_manager.try_auto_fetch_contract(
&key,
target_addr,
AutoFetchReason::Originator,
);
return Ok(DriverOutcome::Publish(Err(ErrorKind::OperationError {
cause: format!(
"originator missing contract code/params for {key}; \
auto-fetch triggered from {target_addr}, \
client should retry after the local store \
is primed (typically <5s in healthy \
networks). Note: `try_auto_fetch_contract` \
is rate-limited 5min/contract, so retries \
within that window will surface the same \
error until the in-flight GET completes."
)
.into(),
}
.into())));
}
if err.is_contract_queue_full() {
tracing::debug!(
tx = %client_tx,
contract = %key,
error = %err,
event = "queue_full",
"update: per-contract queue saturated before forwarding"
);
} else {
tracing::error!(
tx = %client_tx,
contract = %key,
error = %err,
phase = "error",
"update: failed to apply update locally before forwarding"
);
}
return Err(err);
}
};
if let Some(event) =
NetEventLog::update_request(&client_tx, &op_manager.ring, key, target.clone())
{
op_manager.ring.register_events(Either::Left(event)).await;
}
let msg = NetMessage::from(UpdateMsg::RequestUpdate {
id: client_tx,
key,
related_contracts,
value: updated_value,
});
let mut ctx = op_manager.op_ctx(client_tx);
ctx.send_fire_and_forget(target_addr, msg).await?;
tracing::debug!(
tx = %client_tx,
%key,
target_peer = %target_addr,
"update: forwarded to target, operation complete"
);
let host_result: HostResult = Ok(HostResponse::ContractResponse(
ContractResponse::UpdateResponse {
key,
summary: summary.clone(),
},
));
Ok(DriverOutcome::Publish(host_result))
}
}
}
fn classify_update_outcome_for_op_stats(outcome: &DriverOutcome) -> bool {
matches!(outcome, DriverOutcome::Publish(Ok(_)))
}
fn deliver_outcome(op_manager: &OpManager, client_tx: Transaction, outcome: DriverOutcome) {
if !client_tx.is_sub_operation() {
let success = classify_update_outcome_for_op_stats(&outcome);
crate::node::network_status::record_op_result(
crate::node::network_status::OpType::Update,
success,
);
}
match outcome {
DriverOutcome::Publish(result) => {
op_manager.send_client_result(client_tx, result);
}
DriverOutcome::InfrastructureError(err) => {
tracing::warn!(
tx = %client_tx,
error = %err,
"update: infrastructure error; publishing synthesized client error"
);
let synthesized: HostResult = Err(ErrorKind::OperationError {
cause: format!("UPDATE failed: {err}").into(),
}
.into());
op_manager.send_client_result(client_tx, synthesized);
}
}
op_manager.completed(client_tx);
}
pub(crate) async fn start_relay_request_update(
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
related_contracts: RelatedContracts<'static>,
value: WrappedState,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
#[cfg(any(test, feature = "testing"))]
RELAY_UPDATE_DRIVER_CALL_COUNT.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if !op_manager.active_relay_update_txs.insert(incoming_tx) {
RELAY_UPDATE_DEDUP_REJECTS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
phase = "relay_update_dedup_reject",
"UPDATE relay: duplicate RequestUpdate for in-flight tx, dropping"
);
return Ok(());
}
crate::node::network_status::record_relayed_update();
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
phase = "relay_update_request_start",
"UPDATE relay: spawning RequestUpdate driver"
);
RELAY_UPDATE_INFLIGHT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
RELAY_UPDATE_SPAWNED_TOTAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let guard = RelayUpdateInflightGuard {
op_manager: op_manager.clone(),
incoming_tx,
};
GlobalExecutor::spawn(run_relay_request_update(
guard,
op_manager,
incoming_tx,
key,
related_contracts,
value,
sender_addr,
));
Ok(())
}
pub(crate) async fn start_relay_broadcast_to(
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
payload: DeltaOrFullState,
sender_summary_bytes: Vec<u8>,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
#[cfg(any(test, feature = "testing"))]
RELAY_UPDATE_DRIVER_CALL_COUNT.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if !op_manager.active_relay_update_txs.insert(incoming_tx) {
RELAY_UPDATE_DEDUP_REJECTS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
phase = "relay_update_dedup_reject",
"UPDATE relay: duplicate BroadcastTo for in-flight tx, dropping"
);
return Ok(());
}
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
phase = "relay_update_broadcast_start",
"UPDATE relay: spawning BroadcastTo driver"
);
RELAY_UPDATE_INFLIGHT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
RELAY_UPDATE_SPAWNED_TOTAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let guard = RelayUpdateInflightGuard {
op_manager: op_manager.clone(),
incoming_tx,
};
GlobalExecutor::spawn(run_relay_broadcast_to(
guard,
op_manager,
incoming_tx,
key,
payload,
sender_summary_bytes,
sender_addr,
));
Ok(())
}
struct RelayUpdateInflightGuard {
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
}
impl Drop for RelayUpdateInflightGuard {
fn drop(&mut self) {
self.op_manager
.active_relay_update_txs
.remove(&self.incoming_tx);
RELAY_UPDATE_INFLIGHT.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
RELAY_UPDATE_COMPLETED_TOTAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}
async fn run_relay_request_update(
guard: RelayUpdateInflightGuard,
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
related_contracts: RelatedContracts<'static>,
value: WrappedState,
sender_addr: SocketAddr,
) {
let _guard = guard;
if let Err(err) = drive_relay_request_update(
&op_manager,
incoming_tx,
key,
related_contracts,
value,
sender_addr,
)
.await
{
if err.is_contract_queue_full() {
tracing::debug!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_request_error",
event = "queue_full",
"UPDATE relay: RequestUpdate driver returned error"
);
} else {
tracing::warn!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_request_error",
"UPDATE relay: RequestUpdate driver returned error"
);
}
}
}
async fn run_relay_broadcast_to(
guard: RelayUpdateInflightGuard,
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
payload: DeltaOrFullState,
sender_summary_bytes: Vec<u8>,
sender_addr: SocketAddr,
) {
let _guard = guard;
if let Err(err) = drive_relay_broadcast_to(
&op_manager,
incoming_tx,
key,
payload,
sender_summary_bytes,
sender_addr,
)
.await
{
if err.is_contract_queue_full() {
tracing::debug!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_broadcast_error",
event = "queue_full",
"UPDATE relay: BroadcastTo driver returned error"
);
} else {
tracing::warn!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_broadcast_error",
"UPDATE relay: BroadcastTo driver returned error"
);
}
}
}
async fn drive_relay_request_update(
op_manager: &OpManager,
incoming_tx: Transaction,
key: ContractKey,
related_contracts: RelatedContracts<'static>,
value: WrappedState,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
let executing_addr = op_manager
.ring
.connection_manager
.own_location()
.socket_addr();
tracing::debug!(
tx = %incoming_tx,
%key,
executing_peer = ?executing_addr,
request_sender = %sender_addr,
"UPDATE relay: processing RequestUpdate"
);
let state_before = match op_manager
.notify_contract_handler(ContractHandlerEvent::GetQuery {
instance_id: *key.id(),
return_contract_code: false,
})
.await
{
Ok(ContractHandlerEvent::GetResponse {
response: Ok(StoreResponse { state: Some(s), .. }),
..
}) => Some(s),
_ => None,
};
if state_before.is_some() {
let hash_before = state_before.as_ref().map(state_hash_full);
let UpdateExecution {
value: updated_value,
summary: _,
changed,
..
} = super::update_contract(
op_manager,
key,
UpdateData::State(State::from(value.clone())),
related_contracts.clone(),
crate::contract::Priority::NetworkRelay,
)
.await?;
let hash_after = Some(state_hash_full(&updated_value));
if let Some(requester_pkl) = op_manager
.ring
.connection_manager
.get_peer_by_addr(sender_addr)
{
if let Some(event) = NetEventLog::update_success(
&incoming_tx,
&op_manager.ring,
key,
requester_pkl,
hash_before,
hash_after,
Some(updated_value.len()),
) {
op_manager.ring.register_events(Either::Left(event)).await;
}
}
if !changed {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay: yielded no state change, skipping broadcast"
);
} else {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay: RequestUpdate succeeded, state changed"
);
}
return Ok(());
}
let self_addr = op_manager.ring.connection_manager.peer_addr()?;
let skip_list = vec![self_addr, sender_addr];
let next_target = op_manager
.ring
.closest_potentially_hosting(&key, skip_list.as_slice());
let forward_target = match next_target {
Some(t) => t,
None => {
let candidates = op_manager
.ring
.k_closest_potentially_hosting(&key, skip_list.as_slice(), 5)
.into_iter()
.filter_map(|loc| loc.socket_addr())
.map(|addr| format!("{:.8}", addr))
.collect::<Vec<_>>();
let connection_count = op_manager.ring.connection_manager.num_connections();
tracing::error!(
tx = %incoming_tx,
contract = %key,
?candidates,
connection_count,
peer_addr = %sender_addr,
phase = "error",
"UPDATE relay: contract not local and no peers to forward to"
);
if let Some(event) = NetEventLog::update_failure(
&incoming_tx,
&op_manager.ring,
key,
OperationFailure::NoPeersAvailable,
) {
op_manager.ring.register_events(Either::Left(event)).await;
}
return Err(OpError::RingError(RingError::NoHostingPeers(*key.id())));
}
};
let forward_addr = match forward_target.socket_addr() {
Some(addr) => addr,
None => {
tracing::error!(
tx = %incoming_tx,
%key,
target_pub_key = %forward_target.pub_key(),
"UPDATE relay: forward target has no socket address"
);
return Err(OpError::RingError(RingError::NoHostingPeers(*key.id())));
}
};
tracing::debug!(
tx = %incoming_tx,
%key,
next_peer = %forward_addr,
"UPDATE relay: forwarding to peer that might have contract"
);
if let Some(event) =
NetEventLog::update_request(&incoming_tx, &op_manager.ring, key, forward_target.clone())
{
op_manager.ring.register_events(Either::Left(event)).await;
}
let request = NetMessage::from(UpdateMsg::RequestUpdate {
id: incoming_tx,
key,
related_contracts,
value,
});
let mut ctx = op_manager.op_ctx(incoming_tx);
if let Err(err) = ctx.send_fire_and_forget(forward_addr, request).await {
tracing::warn!(
tx = %incoming_tx,
%key,
target = %forward_addr,
error = %err,
"UPDATE relay: forward send failed"
);
return Err(err);
}
Ok(())
}
async fn send_queue_full_resync_request(
op_manager: &OpManager,
key: ContractKey,
sender_addr: SocketAddr,
incoming_tx: Transaction,
) {
let Some(reservation_deadline) = op_manager
.interest_manager
.begin_resync_request(&key, sender_addr)
else {
tracing::debug!(
tx = %incoming_tx,
contract = %key,
sender = %sender_addr,
event = "queue_full_resync_throttled",
"UPDATE relay: queue-full ResyncRequest throttled (rate limit) to avoid amplification"
);
return;
};
if !op_manager
.ring
.resync_emit_limiter
.check_and_record(*key.id())
{
op_manager
.interest_manager
.cancel_resync_request(&key, sender_addr);
crate::config::GlobalTestMetrics::record_resync_request_suppressed();
tracing::debug!(
tx = %incoming_tx,
contract = %key,
sender = %sender_addr,
event = "queue_full_resync_suppressed_global",
"UPDATE relay: queue-full ResyncRequest suppressed by global per-contract cap"
);
return;
}
tracing::info!(
tx = %incoming_tx,
contract = %key,
target = %sender_addr,
event = "resync_request_sent",
reason = "queue_full",
"UPDATE relay: sending rate-limited ResyncRequest after queue-full drop (#4857)"
);
if let Some(slot) = op_manager.interest_manager.try_reserve_resync_retry_slot() {
let op_mgr = op_manager.clone();
GlobalExecutor::spawn(async move {
let _slot = slot; resend_queue_full_resync_request(
&op_mgr,
key,
sender_addr,
incoming_tx,
reservation_deadline,
)
.await;
});
}
op_manager
.ring
.outstanding_resync_requests
.record(*key.id(), sender_addr);
if let Err(e) = op_manager
.notify_node_event(NodeEvent::SendInterestMessage {
target: sender_addr,
message: InterestMessage::ResyncRequest { key },
})
.await
{
tracing::warn!(
tx = %incoming_tx,
error = %e,
"UPDATE relay: failed to send queue-full ResyncRequest"
);
}
}
const QUEUE_FULL_RESYNC_MAX_RETRIES: u32 = 2;
const QUEUE_FULL_RESYNC_RETRY_BASE_DELAY: Duration = Duration::from_secs(2);
const QUEUE_FULL_RESYNC_RETRY_POLL_INTERVAL: Duration = Duration::from_millis(100);
fn jittered_resync_retry_delay(attempt: u32) -> Duration {
let factor: f64 = GlobalRng::random_range(0.8..1.2);
QUEUE_FULL_RESYNC_RETRY_BASE_DELAY.mul_f64(attempt as f64 * factor)
}
async fn resend_queue_full_resync_request(
op_manager: &OpManager,
key: ContractKey,
sender_addr: SocketAddr,
incoming_tx: Transaction,
reservation_deadline: Instant,
) {
let started = tokio::time::Instant::now();
for attempt in 1..=QUEUE_FULL_RESYNC_MAX_RETRIES {
let now = op_manager.interest_manager.now();
let remaining = reservation_deadline.saturating_duration_since(now);
let delay = jittered_resync_retry_delay(attempt).min(remaining / 2);
let target = now + delay;
loop {
let now = op_manager.interest_manager.now();
if now >= reservation_deadline {
tracing::debug!(
tx = %incoming_tx,
contract = %key,
attempt,
event = "queue_full_resync_retry_window_elapsed",
"UPDATE relay: reservation window elapsed, stopping queue-full \
ResyncRequest retries (#4857 P2)"
);
return;
}
if now >= target {
break;
}
if started.elapsed() >= RESYNC_REQUEST_MIN_INTERVAL {
tracing::debug!(
tx = %incoming_tx,
contract = %key,
attempt,
event = "queue_full_resync_retry_backstop",
"UPDATE relay: tokio liveness backstop elapsed (injected clock \
stalled), stopping queue-full ResyncRequest retries (#4857 P2)"
);
return;
}
tokio::time::sleep(QUEUE_FULL_RESYNC_RETRY_POLL_INTERVAL).await;
}
match op_manager.try_notify_node_event(NodeEvent::SendInterestMessage {
target: sender_addr,
message: InterestMessage::ResyncRequest { key },
}) {
Ok(()) => tracing::debug!(
tx = %incoming_tx,
contract = %key,
target = %sender_addr,
attempt,
event = "queue_full_resync_retry_sent",
"UPDATE relay: re-dispatched queue-full ResyncRequest (#4857 P2)"
),
Err(e) => tracing::debug!(
tx = %incoming_tx,
contract = %key,
error = %e,
attempt,
event = "queue_full_resync_retry_dropped",
"UPDATE relay: queue-full ResyncRequest retry dropped (heartbeat backstops)"
),
}
}
}
async fn drive_relay_broadcast_to(
op_manager: &OpManager,
incoming_tx: Transaction,
key: ContractKey,
payload: DeltaOrFullState,
sender_summary_bytes: Vec<u8>,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
let self_location = op_manager.ring.connection_manager.own_location();
let sender_summary = StateSummary::from(sender_summary_bytes.clone());
if let Some(sender_pkl) = op_manager
.ring
.connection_manager
.get_peer_by_addr(sender_addr)
{
let sender_key = crate::ring::PeerKey::from(sender_pkl.pub_key().clone());
op_manager
.interest_manager
.update_peer_summary(&key, &sender_key, Some(sender_summary));
}
let (update_data, payload_bytes) = match &payload {
DeltaOrFullState::Delta(bytes) => {
tracing::debug!(
contract = %key,
delta_size = bytes.len(),
"UPDATE relay: received delta broadcast"
);
(
UpdateData::Delta(StateDelta::from(bytes.clone())),
bytes.clone(),
)
}
DeltaOrFullState::FullState(bytes) => {
tracing::debug!(
contract = %key,
state_size = bytes.len(),
"UPDATE relay: received full state broadcast"
);
(UpdateData::State(State::from(bytes.clone())), bytes.clone())
}
};
let is_delta = matches!(payload, DeltaOrFullState::Delta(_));
let state_for_telemetry = WrappedState::from(payload_bytes.clone());
if let Some(requester_pkl) = op_manager
.ring
.connection_manager
.get_peer_by_addr(sender_addr)
{
if let Some(event) = NetEventLog::update_broadcast_received(
&incoming_tx,
&op_manager.ring,
key,
requester_pkl,
state_for_telemetry.clone(),
) {
op_manager.ring.register_events(Either::Left(event)).await;
}
}
if op_manager.broadcast_dedup_cache.check_and_insert(
&key,
&payload_bytes,
is_delta,
op_manager.interest_manager.now(),
) {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay: BroadcastTo skipped — duplicate payload (dedup hit)"
);
return Ok(());
}
let payload_hash = crate::ring::merge_backoff::merge_payload_hash(is_delta, &payload_bytes);
match op_manager
.ring
.merge_backoff
.check(key.id(), sender_addr, payload_hash)
{
crate::ring::merge_backoff::MergeDecision::Allow => {}
decision @ (crate::ring::merge_backoff::MergeDecision::InBackoff
| crate::ring::merge_backoff::MergeDecision::KnownFailedPayload) => {
crate::config::GlobalTestMetrics::record_merge_suppressed_by_backoff();
tracing::debug!(
tx = %incoming_tx,
%key,
?decision,
"UPDATE relay: BroadcastTo merge skipped — contract in merge-failure backoff"
);
return Ok(());
}
}
let update_result = super::update_contract(
op_manager,
key,
update_data,
RelatedContracts::default(),
crate::contract::Priority::NetworkRelay,
)
.await;
let UpdateExecution {
value: updated_value,
summary: update_summary,
changed,
..
} = match update_result {
Ok(result) => {
if is_delta {
op_manager.ring.merge_backoff.record_success_from_sender(
key.id(),
sender_addr,
result.changed,
);
op_manager
.ring
.delta_incompat
.record_delta_success(key.id());
} else if result.changed {
op_manager
.ring
.merge_backoff
.invalidate_payload_memo(key.id());
}
result
}
Err(err) => {
if err.is_contract_exec_rejection() && !err.is_scheduler_timeout() {
let class = if err.is_wasm_timeout() {
crate::ring::merge_backoff::MergeFailureClass::Timeout
} else {
crate::ring::merge_backoff::MergeFailureClass::Invalid
};
op_manager.ring.merge_backoff.record_failure(
key.id(),
sender_addr,
class,
payload_hash,
);
if is_delta && class == crate::ring::merge_backoff::MergeFailureClass::Invalid {
op_manager
.ring
.delta_incompat
.note_delta_apply_failed(*key.id());
}
}
if err.is_invalid_update_rejection() {
let op_mgr = op_manager.clone();
let contract_key = key;
let sender_summary = sender_summary_bytes.clone();
tokio::spawn(async move {
super::send_summary_back_on_rejection(
&op_mgr,
&contract_key,
sender_addr,
sender_summary,
)
.await;
});
}
let queue_full = err.is_contract_queue_full() || err.is_scheduler_timeout();
if is_delta && !queue_full {
if op_manager
.ring
.resync_emit_limiter
.check_and_record(*key.id())
{
tracing::warn!(
tx = %incoming_tx,
contract = %key,
sender = %sender_addr,
error = %err,
event = "delta_apply_failed",
"UPDATE relay: delta apply failed, sending ResyncRequest"
);
if let Some(sender_pkl) = op_manager
.ring
.connection_manager
.get_peer_by_addr(sender_addr)
{
let sender_key = crate::ring::PeerKey::from(sender_pkl.pub_key().clone());
op_manager
.interest_manager
.update_peer_summary(&key, &sender_key, None);
}
tracing::info!(
tx = %incoming_tx,
contract = %key,
target = %sender_addr,
event = "resync_request_sent",
"UPDATE relay: sending ResyncRequest after delta failure"
);
op_manager
.ring
.outstanding_resync_requests
.record(*key.id(), sender_addr);
if let Err(e) = op_manager
.notify_node_event(NodeEvent::SendInterestMessage {
target: sender_addr,
message: InterestMessage::ResyncRequest { key },
})
.await
{
tracing::warn!(
tx = %incoming_tx,
error = %e,
"UPDATE relay: failed to send ResyncRequest"
);
}
} else {
crate::config::GlobalTestMetrics::record_resync_request_suppressed();
tracing::debug!(
tx = %incoming_tx,
contract = %key,
target = %sender_addr,
event = "resync_request_suppressed",
"UPDATE relay: ResyncRequest suppressed by per-contract rate limit"
);
}
} else if !is_delta && !err.is_contract_exec_rejection() && !queue_full {
op_manager.try_auto_fetch_contract(
&key,
sender_addr,
AutoFetchReason::InboundRelay,
);
} else if queue_full {
send_queue_full_resync_request(op_manager, key, sender_addr, incoming_tx).await;
}
return Err(err);
}
};
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay: BroadcastTo applied"
);
if let Some(event) = NetEventLog::update_broadcast_applied(
&incoming_tx,
&op_manager.ring,
key,
&state_for_telemetry,
&updated_value,
changed,
) {
op_manager.ring.register_events(Either::Left(event)).await;
}
if !changed {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay: BroadcastTo produced no change, ending propagation"
);
return Ok(());
}
crate::node::network_status::record_update_received();
tracing::debug!(
"UPDATE relay: contract {} @ {:?} updated via BroadcastTo",
key,
self_location.location()
);
let op_mgr = op_manager.clone();
let summary = update_summary.clone();
GlobalExecutor::spawn(async move {
super::send_proactive_summary_notification(&op_mgr, &key, sender_addr, summary).await;
});
Ok(())
}
pub(crate) async fn start_relay_request_update_streaming(
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
stream_id: StreamId,
total_size: u64,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
#[cfg(any(test, feature = "testing"))]
RELAY_UPDATE_STREAMING_DRIVER_CALL_COUNT.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if !op_manager.active_relay_update_txs.insert(incoming_tx) {
RELAY_UPDATE_DEDUP_REJECTS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
phase = "relay_update_streaming_dedup_reject",
"UPDATE relay (driver streaming): duplicate RequestUpdateStreaming for in-flight tx, dropping"
);
return Ok(());
}
crate::node::network_status::record_relayed_update();
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
%stream_id,
total_size,
phase = "relay_update_streaming_request_start",
"UPDATE relay (driver streaming): spawning RequestUpdateStreaming driver"
);
RELAY_UPDATE_STREAMING_INFLIGHT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
RELAY_UPDATE_STREAMING_SPAWNED_TOTAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let guard = RelayUpdateStreamingInflightGuard {
op_manager: op_manager.clone(),
incoming_tx,
};
GlobalExecutor::spawn(run_relay_request_update_streaming(
guard,
op_manager,
incoming_tx,
key,
stream_id,
total_size,
sender_addr,
));
Ok(())
}
pub(crate) async fn start_relay_broadcast_to_streaming(
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
stream_id: StreamId,
total_size: u64,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
#[cfg(any(test, feature = "testing"))]
RELAY_UPDATE_STREAMING_DRIVER_CALL_COUNT.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if !op_manager.active_relay_update_txs.insert(incoming_tx) {
RELAY_UPDATE_DEDUP_REJECTS.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
phase = "relay_update_streaming_dedup_reject",
"UPDATE relay (driver streaming): duplicate BroadcastToStreaming for in-flight tx, dropping"
);
return Ok(());
}
tracing::debug!(
tx = %incoming_tx,
%key,
%sender_addr,
%stream_id,
total_size,
phase = "relay_update_streaming_broadcast_start",
"UPDATE relay (driver streaming): spawning BroadcastToStreaming driver"
);
RELAY_UPDATE_STREAMING_INFLIGHT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
RELAY_UPDATE_STREAMING_SPAWNED_TOTAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let guard = RelayUpdateStreamingInflightGuard {
op_manager: op_manager.clone(),
incoming_tx,
};
GlobalExecutor::spawn(run_relay_broadcast_to_streaming(
guard,
op_manager,
incoming_tx,
key,
stream_id,
total_size,
sender_addr,
));
Ok(())
}
struct RelayUpdateStreamingInflightGuard {
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
}
impl Drop for RelayUpdateStreamingInflightGuard {
fn drop(&mut self) {
self.op_manager
.active_relay_update_txs
.remove(&self.incoming_tx);
RELAY_UPDATE_STREAMING_INFLIGHT.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
RELAY_UPDATE_STREAMING_COMPLETED_TOTAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}
async fn run_relay_request_update_streaming(
guard: RelayUpdateStreamingInflightGuard,
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
stream_id: StreamId,
total_size: u64,
sender_addr: SocketAddr,
) {
let _guard = guard;
if let Err(err) = drive_relay_request_update_streaming(
&op_manager,
incoming_tx,
key,
stream_id,
total_size,
sender_addr,
)
.await
{
if err.is_contract_queue_full() {
tracing::debug!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_streaming_request_error",
event = "queue_full",
"UPDATE relay (driver streaming): RequestUpdateStreaming driver returned error"
);
} else {
tracing::warn!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_streaming_request_error",
"UPDATE relay (driver streaming): RequestUpdateStreaming driver returned error"
);
}
}
}
async fn run_relay_broadcast_to_streaming(
guard: RelayUpdateStreamingInflightGuard,
op_manager: Arc<OpManager>,
incoming_tx: Transaction,
key: ContractKey,
stream_id: StreamId,
total_size: u64,
sender_addr: SocketAddr,
) {
let _guard = guard;
if let Err(err) = drive_relay_broadcast_to_streaming(
&op_manager,
incoming_tx,
key,
stream_id,
total_size,
sender_addr,
)
.await
{
if err.is_contract_queue_full() {
tracing::debug!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_streaming_broadcast_error",
event = "queue_full",
"UPDATE relay (driver streaming): BroadcastToStreaming driver returned error"
);
} else {
tracing::warn!(
tx = %incoming_tx,
%key,
error = %err,
phase = "relay_update_streaming_broadcast_error",
"UPDATE relay (driver streaming): BroadcastToStreaming driver returned error"
);
}
}
}
async fn drive_relay_request_update_streaming(
op_manager: &OpManager,
incoming_tx: Transaction,
key: ContractKey,
stream_id: StreamId,
total_size: u64,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
use crate::operations::orphan_streams::{OrphanStreamError, STREAM_CLAIM_TIMEOUT};
tracing::info!(
tx = %incoming_tx,
contract = %key,
%stream_id,
total_size,
"UPDATE relay (driver streaming): processing RequestUpdateStreaming"
);
let stream_handle = match op_manager
.orphan_stream_registry()
.claim_or_wait(sender_addr, stream_id, STREAM_CLAIM_TIMEOUT)
.await
{
Ok(handle) => handle,
Err(OrphanStreamError::AlreadyClaimed) => {
tracing::debug!(
tx = %incoming_tx,
%stream_id,
"UPDATE relay (driver streaming): RequestUpdateStreaming skipped — stream already claimed (dedup)"
);
return Ok(());
}
Err(e) => {
tracing::error!(
tx = %incoming_tx,
%stream_id,
error = %e,
"UPDATE relay (driver streaming): failed to claim stream from orphan registry"
);
return Err(OpError::OrphanStreamClaimFailed);
}
};
let stream_data = match stream_handle.assemble().await {
Ok(data) => {
tracing::debug!(
tx = %incoming_tx,
%stream_id,
assembled_size = data.len(),
expected_size = total_size,
"UPDATE relay (driver streaming): stream assembled"
);
data
}
Err(e) => {
tracing::error!(
tx = %incoming_tx,
%stream_id,
error = %e,
"UPDATE relay (driver streaming): failed to assemble stream"
);
return Err(OpError::StreamCancelled);
}
};
let payload: UpdateStreamingPayload = match bincode::deserialize(&stream_data) {
Ok(p) => p,
Err(e) => {
tracing::error!(
tx = %incoming_tx,
error = %e,
"UPDATE relay (driver streaming): failed to deserialize UpdateStreamingPayload"
);
return Err(OpError::invalid_transition(incoming_tx));
}
};
let UpdateStreamingPayload {
related_contracts,
value,
} = payload;
let UpdateExecution {
value: updated_value,
summary: _,
changed,
..
} = super::update_contract(
op_manager,
key,
UpdateData::State(State::from(value.clone())),
related_contracts,
crate::contract::Priority::NetworkRelay,
)
.await?;
let hash_after = Some(state_hash_full(&updated_value));
if let Some(requester_pkl) = op_manager
.ring
.connection_manager
.get_peer_by_addr(sender_addr)
{
if let Some(event) = NetEventLog::update_success(
&incoming_tx,
&op_manager.ring,
key,
requester_pkl,
None, hash_after,
Some(updated_value.len()),
) {
op_manager.ring.register_events(Either::Left(event)).await;
}
}
if changed {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay (driver streaming): RequestUpdateStreaming succeeded, state changed"
);
} else {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay (driver streaming): RequestUpdateStreaming yielded no state change"
);
}
Ok(())
}
async fn drive_relay_broadcast_to_streaming(
op_manager: &OpManager,
incoming_tx: Transaction,
key: ContractKey,
stream_id: StreamId,
total_size: u64,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
use crate::operations::orphan_streams::{OrphanStreamError, STREAM_CLAIM_TIMEOUT};
tracing::info!(
tx = %incoming_tx,
contract = %key,
%stream_id,
total_size,
sender = %sender_addr,
"UPDATE relay (driver streaming): processing BroadcastToStreaming"
);
let stream_handle = match op_manager
.orphan_stream_registry()
.claim_or_wait(sender_addr, stream_id, STREAM_CLAIM_TIMEOUT)
.await
{
Ok(handle) => handle,
Err(OrphanStreamError::AlreadyClaimed) => {
tracing::debug!(
tx = %incoming_tx,
%stream_id,
"UPDATE relay (driver streaming): BroadcastToStreaming skipped — stream already claimed (dedup)"
);
return Ok(());
}
Err(e) => {
tracing::error!(
tx = %incoming_tx,
%stream_id,
error = %e,
"UPDATE relay (driver streaming): failed to claim stream from orphan registry (broadcast)"
);
return Err(OpError::OrphanStreamClaimFailed);
}
};
let stream_data = match stream_handle.assemble().await {
Ok(data) => {
tracing::debug!(
tx = %incoming_tx,
%stream_id,
assembled_size = data.len(),
expected_size = total_size,
"UPDATE relay (driver streaming): stream assembled (broadcast)"
);
data
}
Err(e) => {
tracing::error!(
tx = %incoming_tx,
%stream_id,
error = %e,
"UPDATE relay (driver streaming): failed to assemble stream (broadcast)"
);
return Err(OpError::StreamCancelled);
}
};
let payload: BroadcastStreamingPayload = match bincode::deserialize(&stream_data) {
Ok(p) => p,
Err(e) => {
tracing::error!(
tx = %incoming_tx,
error = %e,
"UPDATE relay (driver streaming): failed to deserialize BroadcastStreamingPayload"
);
return Err(OpError::invalid_transition(incoming_tx));
}
};
let BroadcastStreamingPayload {
state_bytes,
sender_summary_bytes,
} = payload;
apply_streaming_broadcast(
op_manager,
incoming_tx,
key,
state_bytes,
sender_summary_bytes,
sender_addr,
)
.await
}
async fn apply_streaming_broadcast(
op_manager: &OpManager,
incoming_tx: Transaction,
key: ContractKey,
state_bytes: Vec<u8>,
sender_summary_bytes: Vec<u8>,
sender_addr: SocketAddr,
) -> Result<(), OpError> {
let sender_summary = StateSummary::from(sender_summary_bytes.clone());
if let Some(sender_pkl) = op_manager
.ring
.connection_manager
.get_peer_by_addr(sender_addr)
{
let sender_key = crate::ring::PeerKey::from(sender_pkl.pub_key().clone());
op_manager.interest_manager.update_peer_summary(
&key,
&sender_key,
Some(sender_summary.clone()),
);
}
if op_manager.broadcast_dedup_cache.check_and_insert(
&key,
&state_bytes,
false,
op_manager.interest_manager.now(),
) {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay (driver streaming): BroadcastToStreaming skipped — duplicate payload (dedup hit)"
);
return Ok(());
}
let stream_payload_hash = crate::ring::merge_backoff::merge_payload_hash(false, &state_bytes);
match op_manager
.ring
.merge_backoff
.check(key.id(), sender_addr, stream_payload_hash)
{
crate::ring::merge_backoff::MergeDecision::Allow => {}
decision @ (crate::ring::merge_backoff::MergeDecision::InBackoff
| crate::ring::merge_backoff::MergeDecision::KnownFailedPayload) => {
crate::config::GlobalTestMetrics::record_merge_suppressed_by_backoff();
tracing::debug!(
tx = %incoming_tx,
%key,
?decision,
"UPDATE relay (driver streaming): merge skipped — contract in merge-failure backoff"
);
return Ok(());
}
}
let state_for_telemetry = WrappedState::from(state_bytes.clone());
if let Some(requester_pkl) = op_manager
.ring
.connection_manager
.get_peer_by_addr(sender_addr)
{
if let Some(event) = NetEventLog::update_broadcast_received(
&incoming_tx,
&op_manager.ring,
key,
requester_pkl,
state_for_telemetry.clone(),
) {
op_manager.ring.register_events(Either::Left(event)).await;
}
}
let update_result = super::update_contract(
op_manager,
key,
UpdateData::State(State::from(state_bytes.clone())),
RelatedContracts::default(),
crate::contract::Priority::NetworkRelay,
)
.await;
let UpdateExecution {
value: updated_value,
summary: streaming_update_summary,
changed,
..
} = match update_result {
Ok(exec) => {
if exec.changed {
op_manager
.ring
.merge_backoff
.invalidate_payload_memo(key.id());
}
exec
}
Err(err) => {
if err.is_contract_exec_rejection() && !err.is_scheduler_timeout() {
let class = if err.is_wasm_timeout() {
crate::ring::merge_backoff::MergeFailureClass::Timeout
} else {
crate::ring::merge_backoff::MergeFailureClass::Invalid
};
op_manager.ring.merge_backoff.record_failure(
key.id(),
sender_addr,
class,
stream_payload_hash,
);
}
let queue_full = err.is_contract_queue_full() || err.is_scheduler_timeout();
if super::log_broadcast_to_streaming_failure(&incoming_tx, &key, &err) {
op_manager.try_auto_fetch_contract(
&key,
sender_addr,
AutoFetchReason::InboundRelay,
);
} else if err.is_invalid_update_rejection() {
let op_mgr = op_manager.clone();
let contract_key = key;
let sender_summary_bytes = sender_summary_bytes.clone();
GlobalExecutor::spawn(async move {
super::send_summary_back_on_rejection(
&op_mgr,
&contract_key,
sender_addr,
sender_summary_bytes,
)
.await;
});
} else if queue_full {
send_queue_full_resync_request(op_manager, key, sender_addr, incoming_tx).await;
}
return Err(err);
}
};
if let Some(event) = NetEventLog::update_broadcast_applied(
&incoming_tx,
&op_manager.ring,
key,
&state_for_telemetry,
&updated_value,
changed,
) {
op_manager.ring.register_events(Either::Left(event)).await;
}
if !changed {
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay (driver streaming): BroadcastToStreaming produced no change"
);
return Ok(());
}
crate::node::network_status::record_update_received();
tracing::debug!(
tx = %incoming_tx,
%key,
"UPDATE relay (driver streaming): BroadcastToStreaming applied (state changed)"
);
let op_mgr = op_manager.clone();
let summary = streaming_update_summary.clone();
GlobalExecutor::spawn(async move {
super::send_proactive_summary_notification(&op_mgr, &key, sender_addr, summary).await;
});
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn start_client_update_acquires_inflight_guard_before_spawn() {
let src = include_str!("op_ctx_task.rs");
let entry = src
.find("pub(crate) async fn start_client_update(")
.expect("start_client_update must exist");
let after_spawn = src[entry..]
.find("GlobalExecutor::spawn(")
.expect("start_client_update must spawn a driver task");
let before_spawn = &src[entry..entry + after_spawn];
assert!(
before_spawn.contains("op_manager.admit_client_op()"),
"start_client_update must call op_manager.admit_client_op() \
before GlobalExecutor::spawn (atomic admission gate + \
counter bump; closes the Codex r2 TOCTOU)."
);
assert!(
before_spawn.contains("OpError::NodeShuttingDown"),
"start_client_update must early-return OpError::NodeShuttingDown \
on a refused admission."
);
let spawned = &src[entry + after_spawn..];
let block_end = spawned
.find("\n Ok(client_tx)")
.expect("start_client_update must return Ok(client_tx)");
let spawn_block = &spawned[..block_end];
assert!(
spawn_block.contains("let _inflight_guard = inflight_guard;"),
"the ClientOpGuard must be moved into the spawned future."
);
}
#[test]
fn start_client_update_gates_banned_contracts_before_spawn() {
let src = include_str!("op_ctx_task.rs");
let entry = src
.find("pub(crate) async fn start_client_update(")
.expect("start_client_update must exist");
let after_spawn = src[entry..]
.find("GlobalExecutor::spawn(")
.expect("start_client_update must spawn a driver task");
let before_spawn = &src[entry..entry + after_spawn];
assert!(
before_spawn.contains("reject_if_contract_banned"),
"start_client_update must call reject_if_contract_banned() \
before GlobalExecutor::spawn so a banned contract's UPDATE \
is rejected with a typed error instead of being driven to \
peers (#4300 egress self-block)."
);
}
#[test]
fn report_op_init_error_handles_contract_banned() {
let src = include_str!("../../client_events.rs");
assert!(
src.contains("OpError::ContractBanned"),
"report_op_init_error in client_events.rs must explicitly \
handle OpError::ContractBanned so the egress self-block \
(#4300) surfaces a typed error to the client instead of a \
silent proceed-then-timeout."
);
}
#[test]
fn client_events_calls_start_client_update() {
let src = include_str!("../../client_events.rs");
assert!(
src.contains("op_ctx_task::start_client_update"),
"client_events.rs must call update::op_ctx_task::start_client_update \
for client-initiated UPDATEs (not the legacy request_update path)"
);
let update_section = src
.split("ContractRequest::Update")
.nth(1)
.expect("client_events.rs must contain a ContractRequest::Update handler");
let update_section = update_section
.split("ContractRequest::")
.next()
.unwrap_or(update_section);
assert!(
!update_section.contains("request_update("),
"client_events.rs UPDATE handler must NOT call request_update() directly \
— it should use op_ctx_task::start_client_update instead"
);
}
#[test]
fn driver_uses_send_client_result_not_raw_try_send() {
let src = include_str!("op_ctx_task.rs");
assert!(
src.contains("send_client_result"),
"driver must use op_manager.send_client_result() for result delivery"
);
let deliver_fn = src
.find("fn deliver_outcome(")
.expect("deliver_outcome must exist");
let deliver_body = &src[deliver_fn..];
let end = deliver_body
.find("#[cfg(test)]")
.unwrap_or(deliver_body.len());
let deliver_body = &deliver_body[..end];
assert!(
deliver_body.contains("send_client_result("),
"deliver_outcome must call send_client_result"
);
}
#[test]
fn driver_calls_completed() {
let src = include_str!("op_ctx_task.rs");
assert!(
src.contains("op_manager.completed("),
"driver must call op_manager.completed() for cleanup"
);
}
#[test]
fn driver_calls_update_contract() {
let src = include_str!("op_ctx_task.rs");
assert!(
src.contains("update_contract("),
"driver must call update_contract() to preserve WASM merge + \
BroadcastStateChange + persistence side effects"
);
}
#[test]
fn driver_sends_updated_value_in_request() {
let src = include_str!("op_ctx_task.rs");
assert!(
src.contains("value: updated_value"),
"RequestUpdate must use the post-merge updated_value, not the original \
client update_data (mirrors legacy update.rs:2055)"
);
}
#[test]
fn drive_client_update_uses_bootstrap_gateway_fallback() {
const SOURCE: &str = include_str!("op_ctx_task.rs");
let prod = production_source(SOURCE);
let body = extract_fn_body(prod, "async fn drive_client_update(");
assert!(
body.contains("bootstrap_gateway_target("),
"drive_client_update must fall back to bootstrap_gateway_target \
on an empty ring instead of applying locally to no audience or \
failing with NoHostingPeers (#4361 / #4365)"
);
assert!(
body.contains("bootstrap_gateway_target(op_manager, |addr| addr == sender_addr)"),
"drive_client_update must exclude its own address (sender_addr) \
when selecting the bootstrap gateway (#4361 / #4365)"
);
}
#[test]
fn relay_drivers_never_call_send_and_await() {
let src = include_str!("op_ctx_task.rs");
let relay_start = src
.find("pub(crate) async fn start_relay_request_update(")
.expect("start_relay_request_update not found");
let relay_end = src
.find("#[cfg(test)]")
.expect("test module marker not found");
let relay_src = &src[relay_start..relay_end];
assert!(
!relay_src.contains(".send_and_await("),
"Relay UPDATE drivers MUST NOT call send_and_await. UPDATE is \
fire-and-forget; introducing a wait reintroduces phase-5 GET's \
amplifier surface."
);
}
#[test]
fn relay_request_update_uses_fire_and_forget() {
let src = include_str!("op_ctx_task.rs");
let driver_start = src
.find("async fn drive_relay_request_update(")
.expect("drive_relay_request_update not found");
let driver_end_marker = src[driver_start..]
.find("\n}")
.expect("driver body end not found");
let driver_src = &src[driver_start..driver_start + driver_end_marker];
assert!(
driver_src.contains("send_fire_and_forget"),
"drive_relay_request_update must use send_fire_and_forget for the \
forward path"
);
}
#[test]
fn broadcast_to_dedup_runs_before_merge() {
let src = include_str!("op_ctx_task.rs");
let driver_start = src
.find("async fn drive_relay_broadcast_to(")
.expect("drive_relay_broadcast_to not found");
let driver_end_marker = src[driver_start..]
.find("\n}")
.expect("driver body end not found");
let driver_src = &src[driver_start..driver_start + driver_end_marker];
let dedup_pos = driver_src
.find("broadcast_dedup_cache.check_and_insert")
.expect("dedup cache check missing in BroadcastTo driver");
let merge_pos = driver_src
.find("super::update_contract(")
.expect("update_contract call missing in BroadcastTo driver");
assert!(
dedup_pos < merge_pos,
"broadcast_dedup_cache.check_and_insert MUST appear before \
update_contract() in drive_relay_broadcast_to. If the order \
flips, duplicate broadcasts re-run WASM merge and amplify cost."
);
}
#[test]
fn broadcast_to_sends_resync_on_delta_failure() {
let src = include_str!("op_ctx_task.rs");
let driver_start = src
.find("async fn drive_relay_broadcast_to(")
.expect("drive_relay_broadcast_to not found");
let driver_src = &src[driver_start..];
assert!(
driver_src.contains("InterestMessage::ResyncRequest"),
"drive_relay_broadcast_to must emit InterestMessage::ResyncRequest \
when a delta payload fails to apply (mirrors legacy update.rs:745-758)."
);
}
#[test]
fn broadcast_to_triggers_auto_fetch_on_full_state_failure() {
let src = include_str!("op_ctx_task.rs");
let driver_start = src
.find("async fn drive_relay_broadcast_to(")
.expect("drive_relay_broadcast_to not found");
let driver_src = &src[driver_start..];
assert!(
driver_src.contains("try_auto_fetch_contract"),
"drive_relay_broadcast_to must call try_auto_fetch_contract on \
non-rejection full-state failures (mirrors legacy update.rs:764)."
);
}
#[test]
fn broadcast_to_spawns_proactive_summary() {
let src = include_str!("op_ctx_task.rs");
let driver_start = src
.find("async fn drive_relay_broadcast_to(")
.expect("drive_relay_broadcast_to not found");
let driver_src = &src[driver_start..];
assert!(
driver_src.contains("send_proactive_summary_notification"),
"drive_relay_broadcast_to must spawn send_proactive_summary_notification \
after a successful state change (mirrors legacy update.rs:806-819)."
);
}
#[test]
fn broadcast_to_sends_summary_back_on_rejection() {
let src = include_str!("op_ctx_task.rs");
let driver_start = src
.find("async fn drive_relay_broadcast_to(")
.expect("drive_relay_broadcast_to not found");
let after_start = &src[driver_start + 1..];
let end_offset = after_start
.find("\nasync fn ")
.or_else(|| after_start.find("\n#[cfg(test)]"))
.unwrap_or(after_start.len());
let driver_src = &src[driver_start..driver_start + 1 + end_offset];
assert!(
driver_src.contains("err.is_invalid_update_rejection()"),
"drive_relay_broadcast_to must gate summary-back on \
is_invalid_update_rejection (not is_contract_exec_rejection): \
the broader predicate matches OOG/traps which are \
attacker-inducible and must not amplify into summary-back"
);
assert!(
driver_src.contains("send_summary_back_on_rejection"),
"drive_relay_broadcast_to must spawn send_summary_back_on_rejection \
when the WASM merge rejects — otherwise the sender's cached view \
of us stays wrong and it keeps full-state-broadcasting identical \
content"
);
}
#[test]
fn broadcast_to_feeds_delta_incompat_memo() {
let driver_src = broadcast_to_driver_src();
let arm_pos = driver_src
.find(".note_delta_apply_failed(")
.expect("drive_relay_broadcast_to must arm the delta-incompat memo on delta failure");
assert!(
driver_src.contains(
"if is_delta && class == crate::ring::merge_backoff::MergeFailureClass::Invalid"
),
"the delta-incompat arm must be gated on is_delta AND the Invalid \
failure class — arming on Timeout (load) or full-state failures \
would flip healthy contracts to full-state fan-out"
);
let clear_pos = driver_src
.find(".record_delta_success(")
.expect("drive_relay_broadcast_to must clear the delta-incompat memo on delta success");
let is_delta_gate_pos = driver_src
.find("if is_delta {")
.expect("is_delta success gate not found");
assert!(
is_delta_gate_pos < clear_pos && clear_pos < arm_pos,
"record_delta_success must sit inside the `if is_delta` success arm \
(delta-only signal), before the failure classification \
(order: is_delta gate {is_delta_gate_pos} < clear {clear_pos} < arm {arm_pos})"
);
}
#[cfg(test)]
fn broadcast_to_driver_src() -> &'static str {
let src = include_str!("op_ctx_task.rs");
let start = src
.find("async fn drive_relay_broadcast_to(")
.expect("drive_relay_broadcast_to not found");
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
&src[start..start + 1 + end]
}
#[test]
fn broadcast_to_backoff_gate_runs_before_merge() {
let driver_src = broadcast_to_driver_src();
let dedup_pos = driver_src
.find("broadcast_dedup_cache.check_and_insert")
.expect("dedup check missing");
let backoff_pos = driver_src
.find(".check(key.id(), sender_addr, payload_hash)")
.expect("merge_backoff.check gate missing in BroadcastTo driver");
let merge_pos = driver_src
.find("super::update_contract(")
.expect("update_contract call missing");
assert!(
dedup_pos < backoff_pos && backoff_pos < merge_pos,
"merge_backoff.check MUST appear after the dedup check and before \
update_contract() (order: dedup {dedup_pos} < backoff {backoff_pos} \
< merge {merge_pos})"
);
}
#[test]
fn broadcast_to_backoff_records_success_and_gated_failure() {
let driver_src = broadcast_to_driver_src();
let success_pos = driver_src
.find(".record_success_from_sender(")
.expect("drive_relay_broadcast_to must clear the backoff on a successful merge");
let is_delta_gate_pos = driver_src
.find("if is_delta {")
.expect("backoff reset must be gated on `if is_delta` (#4861)");
assert!(
is_delta_gate_pos < success_pos,
"merge_backoff.record_success_from_sender MUST be inside an `if is_delta` \
gate so a full-state apply does not reset the backoff (#4861 fork oscillation)"
);
let record_pos = driver_src
.find(".record_failure(")
.expect("merge_backoff.record_failure missing in BroadcastTo driver");
let guard_pos = driver_src
.find("if err.is_contract_exec_rejection() && !err.is_scheduler_timeout() {")
.expect(
"record_failure must be gated on \
is_contract_exec_rejection() && !err.is_scheduler_timeout()",
);
assert!(
guard_pos < record_pos,
"merge_backoff.record_failure MUST be gated on \
is_contract_exec_rejection() so queue-full and missing-contract \
failures do NOT create a backoff entry (#4861)"
);
assert!(
driver_src.contains("!err.is_scheduler_timeout()"),
"the record gate MUST exclude scheduler timeouts via \
!err.is_scheduler_timeout() (#4864 round-6): a queued-never-ran \
merge must not quarantine the contract"
);
assert!(
driver_src.contains("is_wasm_timeout()"),
"the failure class must distinguish a WASM timeout (longer cooldown) \
via is_wasm_timeout()"
);
}
#[test]
fn broadcast_to_resync_emit_is_rate_limited() {
let driver_src = broadcast_to_driver_src();
let gate_pos = driver_src
.find("resync_emit_limiter")
.expect("ResyncRequest emit must be gated by resync_emit_limiter");
let summary_clear_pos = driver_src
.find("update_peer_summary(&key, &sender_key, None)")
.expect("sender summary-clear missing");
let resync_pos = driver_src
.find("InterestMessage::ResyncRequest")
.expect("ResyncRequest emission missing");
assert!(
gate_pos < summary_clear_pos && summary_clear_pos < resync_pos,
"emit rate-limit gate ({gate_pos}) must wrap BOTH the summary-clear \
({summary_clear_pos}) and the ResyncRequest emission ({resync_pos})"
);
assert!(
driver_src.contains("record_resync_request_suppressed()"),
"the suppressed branch must record the resync-request-suppressed metric"
);
}
#[test]
fn broadcast_to_streaming_has_backoff_wiring() {
let src = include_str!("op_ctx_task.rs");
let start = src
.find("async fn apply_streaming_broadcast(")
.expect("apply_streaming_broadcast (streaming merge path) not found");
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let driver_src = &src[start..start + 1 + end];
let backoff_pos = driver_src
.find(".check(key.id(), sender_addr, stream_payload_hash)")
.expect("streaming merge path missing merge_backoff.check gate");
let merge_pos = driver_src
.find("super::update_contract(")
.expect("streaming merge path missing update_contract");
assert!(
backoff_pos < merge_pos,
"streaming merge_backoff.check MUST run before update_contract"
);
let record_success = concat!("record_", "success");
assert!(
!driver_src.contains(record_success),
"streaming path must NOT reset the backoff — a full-state merge is \
not convergence evidence (only a clean DELTA apply is). See #4861."
);
assert!(
driver_src.contains(".record_failure(")
&& driver_src.contains("is_contract_exec_rejection()"),
"streaming path must record gated failures like the non-streaming driver"
);
assert!(
driver_src.contains("!err.is_scheduler_timeout()"),
"streaming record gate MUST exclude scheduler timeouts \
(!err.is_scheduler_timeout(), #4864 round-6/7): a queued-never-ran \
failure is transient load, not a contract fault — mirror of the relay pin"
);
}
#[test]
fn client_local_delta_success_resets_backoff() {
let src = include_str!("op_ctx_task.rs");
let start = src
.find("async fn drive_client_update(")
.expect("drive_client_update not found");
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let body = &src[start..start + 1 + end];
let success = concat!(".record_", "success_local(");
assert!(
body.contains(success),
"drive_client_update must reset the backoff on a successful merge via \
record_success_local (contract-wide only, not per-sender) (#4864 P2)"
);
assert!(
body.contains("is_delta_update"),
"the client-local backoff reset must be gated delta-only (is_delta_update)"
);
}
#[test]
fn relay_drivers_use_per_node_dedup_gate() {
let src = include_str!("op_ctx_task.rs");
let request_start = src
.find("pub(crate) async fn start_relay_request_update(")
.expect("start_relay_request_update not found");
let broadcast_start = src
.find("pub(crate) async fn start_relay_broadcast_to(")
.expect("start_relay_broadcast_to not found");
let req_body = &src[request_start..request_start + 1500];
let bc_body = &src[broadcast_start..broadcast_start + 1500];
assert!(
req_body.contains("active_relay_update_txs.insert"),
"start_relay_request_update must check active_relay_update_txs"
);
assert!(
bc_body.contains("active_relay_update_txs.insert"),
"start_relay_broadcast_to must check active_relay_update_txs"
);
}
#[test]
fn dedup_rejection_increments_counter() {
let src = include_str!("op_ctx_task.rs");
let relay_start = src
.find("pub(crate) async fn start_relay_request_update(")
.expect("start_relay_request_update not found");
let relay_end = src
.find("#[cfg(test)]")
.expect("test module marker not found");
let relay_src = &src[relay_start..relay_end];
let sites: Vec<&str> = relay_src
.split("active_relay_update_txs.insert(incoming_tx)")
.skip(1)
.collect();
assert_eq!(
sites.len(),
4,
"expected four dedup sites (RequestUpdate + BroadcastTo + \
RequestUpdateStreaming + BroadcastToStreaming)"
);
for (idx, section) in sites.iter().enumerate() {
let window = §ion[..500.min(section.len())];
assert!(
window.contains("RELAY_UPDATE_DEDUP_REJECTS.fetch_add"),
"dedup gate site #{idx} does not increment RELAY_UPDATE_DEDUP_REJECTS"
);
}
}
#[test]
fn raii_guard_clears_dedup_set_on_drop() {
let src = include_str!("op_ctx_task.rs");
let drop_start = src
.find("impl Drop for RelayUpdateInflightGuard")
.expect("RelayUpdateInflightGuard Drop impl not found");
let drop_body = &src[drop_start..drop_start + 600];
assert!(
drop_body.contains("active_relay_update_txs"),
"RelayUpdateInflightGuard::drop must remove from active_relay_update_txs"
);
assert!(
drop_body.contains("RELAY_UPDATE_INFLIGHT.fetch_sub"),
"RelayUpdateInflightGuard::drop must decrement RELAY_UPDATE_INFLIGHT"
);
assert!(
drop_body.contains("RELAY_UPDATE_COMPLETED_TOTAL.fetch_add"),
"RelayUpdateInflightGuard::drop must increment RELAY_UPDATE_COMPLETED_TOTAL"
);
}
#[test]
fn slice_a_drivers_do_not_touch_streaming_variants() {
let src = include_str!("op_ctx_task.rs");
for driver_name in ["drive_relay_request_update(", "drive_relay_broadcast_to("] {
let start = src
.find(&format!("async fn {driver_name}"))
.unwrap_or_else(|| panic!("{driver_name} not found"));
let after = &src[start..];
let end = after.find("\n}\n").expect("fn body end not found");
let driver_src = &after[..end];
assert!(
!driver_src.contains("RequestUpdateStreaming"),
"{driver_name} (slice A) must not reference RequestUpdateStreaming"
);
assert!(
!driver_src.contains("BroadcastToStreaming"),
"{driver_name} (slice A) must not reference BroadcastToStreaming"
);
}
}
#[test]
fn slice_c_streaming_drivers_exist() {
let src = include_str!("op_ctx_task.rs");
assert!(
src.contains("pub(crate) async fn start_relay_request_update_streaming("),
"start_relay_request_update_streaming must exist (slice C entry point)"
);
assert!(
src.contains("pub(crate) async fn start_relay_broadcast_to_streaming("),
"start_relay_broadcast_to_streaming must exist (slice C entry point)"
);
}
#[test]
fn slice_c_drivers_use_tx_dedup_gate() {
let src = include_str!("op_ctx_task.rs");
for entry in [
"pub(crate) async fn start_relay_request_update_streaming(",
"pub(crate) async fn start_relay_broadcast_to_streaming(",
] {
let start = src.find(entry).unwrap_or_else(|| panic!("{entry} missing"));
let body = &src[start..start + 1500];
assert!(
body.contains("active_relay_update_txs.insert"),
"{entry} must insert into active_relay_update_txs for per-tx dedup"
);
assert!(
body.contains("RELAY_UPDATE_DEDUP_REJECTS.fetch_add"),
"{entry} must increment RELAY_UPDATE_DEDUP_REJECTS on dedup reject"
);
}
}
#[test]
fn slice_c_drivers_claim_orphan_stream() {
let src = include_str!("op_ctx_task.rs");
for driver in [
"drive_relay_request_update_streaming(",
"drive_relay_broadcast_to_streaming(",
] {
let start = src
.find(&format!("async fn {driver}"))
.unwrap_or_else(|| panic!("{driver} not found"));
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let driver_src = &src[start..start + 1 + end];
assert!(
driver_src.contains("orphan_stream_registry()"),
"{driver} must claim via orphan_stream_registry() for atomic \
stream dedup"
);
assert!(
driver_src.contains("claim_or_wait("),
"{driver} must use claim_or_wait on the orphan stream registry"
);
}
}
#[test]
fn slice_c_drivers_are_fire_and_forget() {
let src = include_str!("op_ctx_task.rs");
for driver in [
"drive_relay_request_update_streaming(",
"drive_relay_broadcast_to_streaming(",
] {
let start = src
.find(&format!("async fn {driver}"))
.unwrap_or_else(|| panic!("{driver} not found"));
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let driver_src = &src[start..start + 1 + end];
assert!(
!driver_src.contains(".send_and_await("),
"{driver} must not call send_and_await — UPDATE streaming is \
fire-and-forget"
);
assert!(
!driver_src.contains(".pipe_stream("),
"{driver} must not call pipe_stream — streaming UPDATE relay \
propagation is via BroadcastStateChange after local apply, \
not explicit downstream piping"
);
}
}
#[test]
fn broadcast_to_streaming_dedup_runs_before_merge() {
let src = include_str!("op_ctx_task.rs");
let start = src
.find("async fn apply_streaming_broadcast(")
.expect("apply_streaming_broadcast not found");
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let driver_src = &src[start..start + 1 + end];
let dedup_pos = driver_src
.find("broadcast_dedup_cache.check_and_insert")
.expect("dedup check missing in BroadcastToStreaming driver");
let merge_pos = driver_src
.find("super::update_contract(")
.expect("update_contract call missing in BroadcastToStreaming driver");
assert!(
dedup_pos < merge_pos,
"broadcast_dedup_cache.check_and_insert MUST appear before \
update_contract() in apply_streaming_broadcast"
);
}
#[test]
fn broadcast_to_sends_throttled_resync_on_queue_full() {
let src = include_str!("op_ctx_task.rs");
let fn_body = |name: &str| -> &str {
let start = src.find(name).unwrap_or_else(|| panic!("{name} not found"));
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
&src[start..start + 1 + end]
};
let queue_full_arm = |body: &'static str, marker: &str| -> &'static str {
let arm_start = body
.find(marker)
.unwrap_or_else(|| panic!("queue-full arm `{marker}` missing"));
let arm_end = body[arm_start..]
.find("return Err(err);")
.expect("return after queue-full arm missing");
&body[arm_start..arm_start + arm_end]
};
let ns = fn_body("async fn drive_relay_broadcast_to(");
assert!(
ns.contains(
"let queue_full = err.is_contract_queue_full() || err.is_scheduler_timeout();"
),
"drive_relay_broadcast_to must still classify queue-full and fold in \
scheduler timeouts (#4251, #4864 round-6)"
);
let ns_arm = queue_full_arm(ns, "} else if queue_full {");
assert!(
ns_arm.contains("send_queue_full_resync_request("),
"drive_relay_broadcast_to queue-full arm MUST delegate to \
send_queue_full_resync_request (heals #4857)"
);
assert!(
!ns_arm.contains("try_auto_fetch_contract"),
"drive_relay_broadcast_to queue-full arm must NOT auto-fetch (#4251)"
);
assert!(
fn_body("async fn drive_relay_broadcast_to_streaming(")
.contains("apply_streaming_broadcast("),
"drive_relay_broadcast_to_streaming must delegate to \
apply_streaming_broadcast — otherwise the streaming apply/heal logic \
is unreachable"
);
let st = fn_body("async fn apply_streaming_broadcast(");
assert!(
st.contains(
"let queue_full = err.is_contract_queue_full() || err.is_scheduler_timeout();"
),
"apply_streaming_broadcast must classify queue-full and fold in \
scheduler timeouts (#4251, #4864 round-6)"
);
let st_arm = queue_full_arm(st, "} else if queue_full {");
assert!(
st_arm.contains("send_queue_full_resync_request("),
"apply_streaming_broadcast queue-full arm MUST delegate to \
send_queue_full_resync_request — otherwise #4857 stays OPEN on the \
streaming full-state path"
);
assert!(
!st_arm.contains("try_auto_fetch_contract"),
"streaming queue-full arm must NOT auto-fetch (#4251)"
);
let helper = fn_body("async fn send_queue_full_resync_request(");
let per_sender = helper.find(".begin_resync_request(").expect(
"helper must RESERVE the per-sender throttle atomically via \
begin_resync_request (#4864 round-6 item 2)",
);
let global_cap = helper.find(".check_and_record(*key.id())").expect(
"helper must ALSO gate on the global per-contract emit cap \
(resync_emit_limiter.check_and_record, #4864 round-4)",
);
let emit = helper
.find("InterestMessage::ResyncRequest")
.expect("helper must emit InterestMessage::ResyncRequest (heal, #4857)");
let cancel = helper.find(".cancel_resync_request(").expect(
"helper must CANCEL the reservation when the global cap rejects, so the \
window is released (#4864 round-6 item 2)",
);
assert!(
per_sender < emit && global_cap < emit,
"the ResyncRequest emit ({emit}) MUST come AFTER both the per-sender \
throttle reservation ({per_sender}) and the global emit cap ({global_cap}), \
so an unthrottled emission cannot re-open the #4251 storm"
);
assert!(
per_sender < global_cap && global_cap < cancel && cancel < emit,
"the cancel ({cancel}) MUST be in the global-cap reject branch — after the \
reservation ({per_sender}) and the global cap ({global_cap}), before the \
emit ({emit}) — so a global-cap suppression releases the window (#4864 round-6)"
);
}
fn drain_resync_targets(
rx: &mut crate::node::EventLoopNotificationsReceiver,
key: ContractKey,
) -> Vec<SocketAddr> {
let mut targets = Vec::new();
while let Ok(ev) = rx.notifications_receiver.try_recv() {
if let Either::Right(NodeEvent::SendInterestMessage {
target,
message: InterestMessage::ResyncRequest { key: k },
}) = ev
{
if k == key {
targets.push(target);
}
}
}
targets
}
async fn build_queue_full_test_node(
id: &str,
) -> (
Arc<OpManager>,
crate::node::EventLoopNotificationsReceiver,
Box<dyn std::any::Any>,
) {
build_queue_full_test_node_with_clock(id, None).await
}
async fn build_queue_full_test_node_with_clock(
id: &str,
override_clock: Option<crate::util::time_source::DynTimeSource>,
) -> (
Arc<OpManager>,
crate::node::EventLoopNotificationsReceiver,
Box<dyn std::any::Any>,
) {
let config_args = crate::config::ConfigArgs {
id: Some(id.to_string()),
mode: Some(crate::contract::OperationMode::Local),
..Default::default()
};
let mut node_config =
crate::node::NodeConfig::new(config_args.build().await.expect("build Config"))
.await
.expect("build NodeConfig");
node_config.hosting_time_source_override = override_clock;
let (notification_rx, notification_tx) = crate::node::event_loop_notification_channel();
let (ops_ch_channel, mut ch_channel, wait_for_event) =
crate::contract::contract_handler_channel();
let connection_manager = crate::ring::ConnectionManager::new(&node_config);
let (result_router_tx, result_router_rx) = tokio::sync::mpsc::channel(100);
let task_monitor = crate::node::background_task_monitor::BackgroundTaskMonitor::new();
let op_manager = Arc::new(
OpManager::new(
notification_tx,
ops_ch_channel,
&node_config,
crate::tracing::DynamicRegister::new(vec![]),
connection_manager,
result_router_tx,
&task_monitor,
)
.expect("build OpManager"),
);
op_manager.ring.attach_op_manager(&op_manager);
let self_addr: SocketAddr = "127.0.0.1:12000".parse().unwrap();
op_manager.ring.connection_manager.set_own_addr(self_addr);
let handler = tokio::spawn(async move {
while let Ok((id, ev, _priority)) = ch_channel.recv_from_sender().await {
let response = match ev {
ContractHandlerEvent::GetQuery { .. } => ContractHandlerEvent::GetResponse {
key: None,
response: Ok(StoreResponse {
state: None,
contract: None,
}),
},
ContractHandlerEvent::UpdateQuery { .. } => {
ContractHandlerEvent::UpdateResponse {
new_value: Err(crate::contract::ExecutorError::other(
crate::contract::ContractQueueFull,
)),
state_changed: false,
}
}
other => panic!("unexpected handler event in stand-in: {other:?}"),
};
if ch_channel.send_to_sender(id, response).await.is_err() {
break;
}
}
});
let guard: Box<dyn std::any::Any> =
Box::new((handler, result_router_rx, task_monitor, wait_for_event));
(op_manager, notification_rx, guard)
}
#[tokio::test]
async fn queue_full_broadcast_emits_throttled_resync_request() {
let (op_manager, mut notification_rx, _guard) =
build_queue_full_test_node("queue-full-resync-4857-delta").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([7u8; 32]),
CodeHash::new([8u8; 32]),
);
let sender_addr: SocketAddr = "127.0.0.1:12100".parse().unwrap();
let r1 = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![1, 2, 3]),
Vec::new(),
sender_addr,
)
.await;
assert!(
r1.as_ref()
.err()
.is_some_and(|e| e.is_contract_queue_full()),
"expected a queue-full error from the stand-in handler, got {r1:?}"
);
assert_eq!(
drain_resync_targets(&mut notification_rx, key),
vec![sender_addr],
"a queue-full drop MUST emit exactly one ResyncRequest to the sender \
so it clears its poisoned summary and re-sends (heals #4857)"
);
let r2 = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![4, 5, 6]),
Vec::new(),
sender_addr,
)
.await;
assert!(
r2.as_ref()
.err()
.is_some_and(|e| e.is_contract_queue_full()),
"second broadcast should still hit queue-full, got {r2:?}"
);
assert!(
drain_resync_targets(&mut notification_rx, key).is_empty(),
"a second queue-full drop within the throttle window MUST NOT emit \
another ResyncRequest (bounds #4251 amplification)"
);
}
#[tokio::test]
async fn queue_full_streaming_broadcast_emits_throttled_resync_request() {
let (op_manager, mut notification_rx, _guard) =
build_queue_full_test_node("queue-full-resync-4857-streaming").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([9u8; 32]),
CodeHash::new([10u8; 32]),
);
let sender_addr: SocketAddr = "127.0.0.1:12200".parse().unwrap();
let r1 = apply_streaming_broadcast(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
vec![1, 2, 3],
Vec::new(),
sender_addr,
)
.await;
assert!(
r1.as_ref()
.err()
.is_some_and(|e| e.is_contract_queue_full()),
"expected a queue-full error from the stand-in handler, got {r1:?}"
);
assert_eq!(
drain_resync_targets(&mut notification_rx, key),
vec![sender_addr],
"a queue-full streaming full-state drop MUST emit exactly one \
ResyncRequest (heals #4857 on the streaming path)"
);
let r2 = apply_streaming_broadcast(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
vec![4, 5, 6],
Vec::new(),
sender_addr,
)
.await;
assert!(
r2.as_ref()
.err()
.is_some_and(|e| e.is_contract_queue_full()),
"second streaming broadcast should still hit queue-full, got {r2:?}"
);
assert!(
drain_resync_targets(&mut notification_rx, key).is_empty(),
"a second streaming queue-full drop within the throttle window MUST \
NOT emit another ResyncRequest (bounds #4251 amplification)"
);
}
#[cfg(test)]
fn exec_reject_err(key: ContractKey) -> crate::contract::ExecutorError {
crate::contract::ExecutorError::from(
freenet_stdlib::client_api::RequestError::ContractError(
freenet_stdlib::client_api::ContractError::update_exec_error(
key,
"contract merge failed (test poison)",
),
),
)
}
fn queue_full_err() -> crate::contract::ExecutorError {
crate::contract::ExecutorError::other(crate::contract::ContractQueueFull)
}
#[derive(Clone, Copy)]
enum DeltaOutcome {
ExecReject,
QueueFull,
Success,
SuccessNoChange,
}
async fn build_broadcast_test_node(
id: &str,
delta_outcome: DeltaOutcome,
) -> (
Arc<OpManager>,
crate::node::EventLoopNotificationsReceiver,
std::sync::Arc<std::sync::atomic::AtomicUsize>,
Box<dyn std::any::Any>,
) {
use std::sync::atomic::Ordering;
let config_args = crate::config::ConfigArgs {
id: Some(id.to_string()),
mode: Some(crate::contract::OperationMode::Local),
..Default::default()
};
let node_config =
crate::node::NodeConfig::new(config_args.build().await.expect("build Config"))
.await
.expect("build NodeConfig");
let (notification_rx, notification_tx) = crate::node::event_loop_notification_channel();
let (ops_ch_channel, mut ch_channel, wait_for_event) =
crate::contract::contract_handler_channel();
let connection_manager = crate::ring::ConnectionManager::new(&node_config);
let (result_router_tx, result_router_rx) = tokio::sync::mpsc::channel(100);
let task_monitor = crate::node::background_task_monitor::BackgroundTaskMonitor::new();
let op_manager = Arc::new(
OpManager::new(
notification_tx,
ops_ch_channel,
&node_config,
crate::tracing::DynamicRegister::new(vec![]),
connection_manager,
result_router_tx,
&task_monitor,
)
.expect("build OpManager"),
);
op_manager.ring.attach_op_manager(&op_manager);
let self_addr: SocketAddr = "127.0.0.1:12000".parse().unwrap();
op_manager.ring.connection_manager.set_own_addr(self_addr);
let update_query_count = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let counter = update_query_count.clone();
let handler = tokio::spawn(async move {
while let Ok((id, ev, _priority)) = ch_channel.recv_from_sender().await {
let response = match ev {
ContractHandlerEvent::GetQuery { .. } => ContractHandlerEvent::GetResponse {
key: None,
response: Ok(StoreResponse {
state: None,
contract: None,
}),
},
ContractHandlerEvent::GetSummaryQuery { key } => {
ContractHandlerEvent::GetSummaryResponse {
key,
summary: Ok(StateSummary::from(Vec::new())),
}
}
ContractHandlerEvent::UpdateQuery { key, data, .. } => {
counter.fetch_add(1, Ordering::Relaxed);
let is_delta = matches!(
&data,
UpdateData::Delta(_) | UpdateData::RelatedDelta { .. }
);
if is_delta {
match delta_outcome {
DeltaOutcome::ExecReject => ContractHandlerEvent::UpdateResponse {
new_value: Err(exec_reject_err(key)),
state_changed: false,
},
DeltaOutcome::QueueFull => ContractHandlerEvent::UpdateResponse {
new_value: Err(queue_full_err()),
state_changed: false,
},
DeltaOutcome::Success => ContractHandlerEvent::UpdateResponse {
new_value: Ok(WrappedState::from(vec![1u8, 2, 3])),
state_changed: true,
},
DeltaOutcome::SuccessNoChange => {
ContractHandlerEvent::UpdateResponse {
new_value: Ok(WrappedState::from(vec![1u8, 2, 3])),
state_changed: false,
}
}
}
} else {
ContractHandlerEvent::UpdateResponse {
new_value: Ok(WrappedState::from(vec![1u8, 2, 3])),
state_changed: true,
}
}
}
other => {
panic!("unexpected handler event in exec-reject stand-in: {other:?}")
}
};
if ch_channel.send_to_sender(id, response).await.is_err() {
break;
}
}
});
let guard: Box<dyn std::any::Any> =
Box::new((handler, result_router_rx, task_monitor, wait_for_event));
(op_manager, notification_rx, update_query_count, guard)
}
async fn build_exec_reject_test_node(
id: &str,
) -> (
Arc<OpManager>,
crate::node::EventLoopNotificationsReceiver,
std::sync::Arc<std::sync::atomic::AtomicUsize>,
Box<dyn std::any::Any>,
) {
build_broadcast_test_node(id, DeltaOutcome::ExecReject).await
}
#[tokio::test]
async fn backoff_gate_skips_merge_after_invalid_threshold_trips() {
use std::sync::atomic::Ordering;
crate::config::GlobalTestMetrics::reset();
let (op_manager, mut notification_rx, update_queries, _guard) =
build_exec_reject_test_node("backoff-gate-4864").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([21u8; 32]),
CodeHash::new([22u8; 32]),
);
let sender: SocketAddr = "127.0.0.1:12200".parse().unwrap();
for bytes in [vec![1u8], vec![2u8], vec![3u8]] {
let r = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(bytes),
Vec::new(),
sender,
)
.await;
assert!(r.is_err(), "each poison delta must fail the merge");
let _ = drain_resync_targets(&mut notification_rx, key); }
assert_eq!(
update_queries.load(Ordering::Relaxed),
3,
"3 failing merges reached the executor"
);
let suppressed_before = crate::config::GlobalTestMetrics::merges_suppressed_by_backoff();
let r4 = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![4u8]),
Vec::new(),
sender,
)
.await;
assert!(
r4.is_ok(),
"a suppressed broadcast completes Ok (like a dedup skip)"
);
assert_eq!(
update_queries.load(Ordering::Relaxed),
3,
"the suppressed merge MUST NOT reach the executor (no 4th UpdateQuery)"
);
assert!(
crate::config::GlobalTestMetrics::merges_suppressed_by_backoff() > suppressed_before,
"the suppressed merge MUST bump merges_suppressed_by_backoff"
);
assert!(
drain_resync_targets(&mut notification_rx, key).is_empty(),
"no ResyncRequest may be emitted while the merge is suppressed"
);
}
#[tokio::test]
async fn full_state_success_does_not_reset_backoff() {
use std::sync::atomic::Ordering;
crate::config::GlobalTestMetrics::reset();
let (op_manager, mut notification_rx, update_queries, _guard) =
build_exec_reject_test_node("backoff-fullstate-noreset-4864").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([31u8; 32]),
CodeHash::new([32u8; 32]),
);
let sender: SocketAddr = "127.0.0.1:12300".parse().unwrap();
for bytes in [vec![1u8], vec![2u8]] {
let _ = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(bytes),
Vec::new(),
sender,
)
.await;
let _ = drain_resync_targets(&mut notification_rx, key);
}
let ok = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::FullState(vec![9, 9, 9]),
Vec::new(),
sender,
)
.await;
assert!(
ok.is_ok(),
"full-state merge should succeed via the stand-in"
);
let _ = drain_resync_targets(&mut notification_rx, key);
let _ = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![3u8]),
Vec::new(),
sender,
)
.await;
let _ = drain_resync_targets(&mut notification_rx, key);
let queries_after_trip = update_queries.load(Ordering::Relaxed);
let r = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![4u8]),
Vec::new(),
sender,
)
.await;
assert!(r.is_ok());
assert_eq!(
update_queries.load(Ordering::Relaxed),
queries_after_trip,
"the backoff tripped across the full-state success (no reset) — the \
following delta merge is skipped"
);
}
#[tokio::test]
async fn backoff_gate_is_scoped_per_sender_in_driver() {
use std::sync::atomic::Ordering;
crate::config::GlobalTestMetrics::reset();
let (op_manager, mut notification_rx, update_queries, _guard) =
build_exec_reject_test_node("backoff-per-sender-4864").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([41u8; 32]),
CodeHash::new([42u8; 32]),
);
let sender_a: SocketAddr = "127.0.0.1:12400".parse().unwrap();
let sender_b: SocketAddr = "127.0.0.1:12500".parse().unwrap();
for bytes in [vec![1u8], vec![2u8], vec![3u8]] {
let _ = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(bytes),
Vec::new(),
sender_a,
)
.await;
let _ = drain_resync_targets(&mut notification_rx, key);
}
assert_eq!(
update_queries.load(Ordering::Relaxed),
3,
"A's 3 failing merges reached the executor"
);
let ra = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![4u8]),
Vec::new(),
sender_a,
)
.await;
assert!(ra.is_ok(), "A's suppressed broadcast completes Ok");
assert_eq!(
update_queries.load(Ordering::Relaxed),
3,
"A's suppressed delta MUST NOT reach the executor"
);
let rb = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![5u8]),
Vec::new(),
sender_b,
)
.await;
let _ = drain_resync_targets(&mut notification_rx, key);
assert!(
rb.is_err(),
"B's delta reaches the merge (fails via stand-in)"
);
assert_eq!(
update_queries.load(Ordering::Relaxed),
4,
"B's delta MUST reach the executor despite A's per-sender backoff"
);
}
#[tokio::test]
async fn queue_full_failures_do_not_trip_backoff() {
use std::sync::atomic::Ordering;
crate::config::GlobalTestMetrics::reset();
let (op_manager, mut notification_rx, update_queries, _guard) =
build_broadcast_test_node("backoff-queuefull-4864", DeltaOutcome::QueueFull).await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([51u8; 32]),
CodeHash::new([52u8; 32]),
);
let sender: SocketAddr = "127.0.0.1:12600".parse().unwrap();
for bytes in [vec![1u8], vec![2u8], vec![3u8], vec![4u8], vec![5u8]] {
let _ = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(bytes),
Vec::new(),
sender,
)
.await;
let _ = drain_resync_targets(&mut notification_rx, key);
}
assert_eq!(
update_queries.load(Ordering::Relaxed),
5,
"every queue-full delta must still reach the executor — queue-full \
must NOT trip the backoff (it is transient load, not poison)"
);
assert_eq!(
crate::config::GlobalTestMetrics::merges_suppressed_by_backoff(),
0,
"no merge may be suppressed — queue-full never records a backoff entry"
);
assert_eq!(
op_manager.ring.merge_backoff.contracts_in_backoff(),
0,
"no contract may be in backoff after only queue-full failures"
);
}
#[tokio::test]
async fn queue_full_resync_flood_is_bounded_by_global_emit_cap() {
crate::config::GlobalTestMetrics::reset();
let (op_manager, mut notification_rx, _update_queries, _guard) =
build_broadcast_test_node("queuefull-emit-cap-4864", DeltaOutcome::QueueFull).await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([71u8; 32]),
CodeHash::new([72u8; 32]),
);
let sender_count = 6u16;
for i in 0..sender_count {
let sender: SocketAddr = format!("127.0.0.1:{}", 13000 + i).parse().unwrap();
let _ = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![i as u8]),
Vec::new(),
sender,
)
.await;
}
let emissions = drain_resync_targets(&mut notification_rx, key);
let burst = crate::ring::resync_rate_limit::EMIT_BURST as usize;
assert_eq!(
emissions.len(),
burst,
"queue-full resync emissions ({}) must be bounded by the global \
per-contract cap (EMIT_BURST = {}), not one per sender ({sender_count})",
emissions.len(),
burst
);
}
#[test]
fn queue_full_resync_helper_goes_through_global_emit_cap() {
let src = include_str!("op_ctx_task.rs");
let start = src
.find("async fn send_queue_full_resync_request(")
.expect("send_queue_full_resync_request not found");
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let body = &src[start..start + 1 + end];
assert!(
body.contains("resync_emit_limiter"),
"the queue-full resync helper must route through the GLOBAL \
per-contract emit cap (resync_emit_limiter) (#4864 round-4)"
);
assert!(
body.contains("begin_resync_request"),
"and keep the per-(contract, sender) throttle"
);
}
#[test]
fn every_production_resync_request_emit_records_outstanding() {
let src = include_str!("op_ctx_task.rs");
let prod_end = src.find("\nmod tests {").unwrap_or(src.len());
let prod = &src[..prod_end];
let retry_start = prod
.find("async fn resend_queue_full_resync_request(")
.expect("resend_queue_full_resync_request (the #4862 retry) not found");
let retry_end = retry_start
+ 1
+ prod[retry_start + 1..]
.find("\nasync fn ")
.unwrap_or(prod.len() - retry_start - 1);
let mut emits = 0usize;
let mut cursor = 0usize;
while let Some(rel) = prod[cursor..].find("message: InterestMessage::ResyncRequest") {
let pos = cursor + rel;
cursor = pos + 1;
if (retry_start..retry_end).contains(&pos) {
continue;
}
emits += 1;
let window_start = pos.saturating_sub(600);
assert!(
prod[window_start..pos].contains("outstanding_resync_requests"),
"a production ResyncRequest emit at byte {pos} is NOT preceded by an \
outstanding_resync_requests.record(...) within 600 bytes — the \
receive arm would drop its response as unsolicited (#4864 round-8)"
);
}
assert!(
emits >= 2,
"expected at least the queue-full-helper and delta-failure ResyncRequest \
emit sites in production ({emits} found) — has the emit path moved?"
);
}
#[tokio::test]
async fn memoized_payload_replay_from_second_sender_is_hard_skipped() {
use std::sync::atomic::Ordering;
crate::config::GlobalTestMetrics::reset();
let (op_manager, mut notification_rx, update_queries, _guard) =
build_exec_reject_test_node("memo-replay-4864").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([83u8; 32]),
CodeHash::new([84u8; 32]),
);
let sender_a: SocketAddr = "127.0.0.1:13300".parse().unwrap();
let sender_b: SocketAddr = "127.0.0.1:13400".parse().unwrap();
let poisons = [vec![91u8], vec![92u8], vec![93u8]];
for p in &poisons {
let h = crate::ring::merge_backoff::merge_payload_hash(true, p);
op_manager.ring.merge_backoff.record_failure(
key.id(),
sender_a,
crate::ring::merge_backoff::MergeFailureClass::Invalid,
h,
);
}
for p in &poisons {
let r = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(p.clone()),
Vec::new(),
sender_b,
)
.await;
assert!(
r.is_ok(),
"a memoized-payload skip completes Ok (like a dedup skip)"
);
}
assert_eq!(
update_queries.load(Ordering::Relaxed),
0,
"B's replays of memoized payloads MUST NOT reach the executor"
);
assert_eq!(
crate::config::GlobalTestMetrics::merges_suppressed_by_backoff(),
poisons.len() as u64,
"each replay must be a memo skip (bumps merges_suppressed_by_backoff), \
not a dedup skip"
);
assert!(
drain_resync_targets(&mut notification_rx, key).is_empty(),
"a memo skip must NOT emit a ResyncRequest"
);
let novel = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![94u8]),
Vec::new(),
sender_b,
)
.await;
let _ = drain_resync_targets(&mut notification_rx, key);
assert!(
novel.is_err(),
"B's novel delta reaches the merge (fails via stand-in)"
);
assert_eq!(
update_queries.load(Ordering::Relaxed),
1,
"B's novel delta MUST reach the executor — the memo skips never advanced \
B's channel"
);
}
#[tokio::test]
async fn client_local_delta_success_clears_contract_side_not_other_senders() {
let (op_manager, _notification_rx, _update_queries, _guard) =
build_broadcast_test_node("backoff-client-local-4864", DeltaOutcome::Success).await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([61u8; 32]),
CodeHash::new([62u8; 32]),
);
op_manager
.ring
.host_contract(key, 100, crate::ring::AccessType::Put);
let remote_sender: SocketAddr = "127.0.0.1:12700".parse().unwrap();
let probe_sender: SocketAddr = "127.0.0.1:12800".parse().unwrap();
op_manager.ring.merge_backoff.record_failure(
key.id(),
remote_sender,
crate::ring::merge_backoff::MergeFailureClass::Timeout,
0x01,
);
for h in [0x11u64, 0x12, 0x13] {
op_manager.ring.merge_backoff.record_failure(
key.id(),
remote_sender,
crate::ring::merge_backoff::MergeFailureClass::Invalid,
h,
);
}
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), probe_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"seeded contract-wide Timeout should suppress before the client success"
);
let outcome = drive_client_update(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
UpdateData::Delta(StateDelta::from(vec![7u8, 7, 7])),
RelatedContracts::default(),
)
.await;
assert!(
outcome.is_ok(),
"client-local delta merge should succeed via the stand-in: {outcome:?}"
);
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), probe_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::Allow,
"record_success_local must clear the contract-wide Timeout + memo"
);
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), remote_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"a local client's success must NOT clear a remote sender's Invalid \
channel (it says nothing about that sender's fork)"
);
}
#[tokio::test]
async fn client_local_full_state_success_does_not_reset_backoff() {
let (op_manager, _notification_rx, _update_queries, _guard) =
build_broadcast_test_node("backoff-client-fullstate-4864", DeltaOutcome::Success).await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([63u8; 32]),
CodeHash::new([64u8; 32]),
);
op_manager
.ring
.host_contract(key, 100, crate::ring::AccessType::Put);
let remote_sender: SocketAddr = "127.0.0.1:15100".parse().unwrap();
let probe_sender: SocketAddr = "127.0.0.1:15200".parse().unwrap();
op_manager.ring.merge_backoff.record_failure(
key.id(),
remote_sender,
crate::ring::merge_backoff::MergeFailureClass::Timeout,
0x01,
);
for h in [0x11u64, 0x12, 0x13] {
op_manager.ring.merge_backoff.record_failure(
key.id(),
remote_sender,
crate::ring::merge_backoff::MergeFailureClass::Invalid,
h,
);
}
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), probe_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"seeded contract-wide Timeout should suppress before the client update"
);
let outcome = drive_client_update(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
UpdateData::State(State::from(vec![9u8, 9, 9])),
RelatedContracts::default(),
)
.await;
assert!(
outcome.is_ok(),
"client-local full-state update should succeed via the stand-in: {outcome:?}"
);
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), probe_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"a client FULL-STATE success must NOT clear the contract-wide Timeout \
(record_success_local is gated is_delta_update)"
);
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), remote_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"a client FULL-STATE success must NOT clear a remote sender's Invalid channel"
);
}
#[tokio::test]
async fn client_local_no_op_delta_success_does_not_clear_contract_side() {
let (op_manager, _notification_rx, _update_queries, _guard) = build_broadcast_test_node(
"backoff-client-nochange-4864",
DeltaOutcome::SuccessNoChange,
)
.await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([65u8; 32]),
CodeHash::new([66u8; 32]),
);
op_manager
.ring
.host_contract(key, 100, crate::ring::AccessType::Put);
let remote_sender: SocketAddr = "127.0.0.1:15300".parse().unwrap();
let probe_sender: SocketAddr = "127.0.0.1:15400".parse().unwrap();
op_manager.ring.merge_backoff.record_failure(
key.id(),
remote_sender,
crate::ring::merge_backoff::MergeFailureClass::Timeout,
0x01,
);
for h in [0x11u64, 0x12, 0x13] {
op_manager.ring.merge_backoff.record_failure(
key.id(),
remote_sender,
crate::ring::merge_backoff::MergeFailureClass::Invalid,
h,
);
}
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), probe_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"seeded contract-wide Timeout should suppress before the client update"
);
let outcome = drive_client_update(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
UpdateData::Delta(StateDelta::from(vec![7u8, 7, 7])),
RelatedContracts::default(),
)
.await;
assert!(
outcome.is_ok(),
"no-op client-local delta should succeed via the stand-in: {outcome:?}"
);
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), probe_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"a no-op (changed=false) client delta success must NOT clear the \
contract-wide Timeout (record_success_local is gated on changed)"
);
assert_eq!(
op_manager
.ring
.merge_backoff
.check(key.id(), remote_sender, 0x99),
crate::ring::merge_backoff::MergeDecision::InBackoff,
"a no-op client delta success must NOT clear a remote sender's Invalid channel"
);
}
#[test]
fn success_sites_pass_real_changed_flag_not_a_literal() {
let src = include_str!("op_ctx_task.rs");
let client = extract_fn_body(src, "async fn drive_client_update(");
let client_calls = client
.matches("record_success_local(key.id(), execution.changed)")
.count();
assert!(
client_calls >= 2,
"both drive_client_update success arms must pass the REAL \
execution.changed to record_success_local (found {client_calls}, \
expected >= 2) — not a literal `true` (#4864 round-5 item 4)"
);
assert!(
!client.contains("record_success_local(key.id(), true)"),
"drive_client_update must NOT hardcode `true` for the changed flag"
);
let relay = broadcast_to_driver_src();
let call = relay.find(".record_success_from_sender(").expect(
"drive_relay_broadcast_to must clear the backoff via record_success_from_sender",
);
let call_args = &relay[call..(call + 160).min(relay.len())];
assert!(
call_args.contains("result.changed"),
"record_success_from_sender must be passed the REAL result.changed \
(found in its arg list), not a literal `true` (#4864 round-5 item 4)"
);
}
#[test]
fn streaming_changed_apply_invalidates_payload_memo() {
let src = include_str!("op_ctx_task.rs");
let st = extract_fn_body(src, "async fn apply_streaming_broadcast(");
let changed_gate = st
.find("if exec.changed")
.expect("streaming Ok arm must gate on `if exec.changed`");
let open_brace = changed_gate
+ st[changed_gate..]
.find('{')
.expect("if exec.changed block must have a body");
let mut depth = 0usize;
let mut close_brace = None;
for (i, ch) in st[open_brace..].char_indices() {
match ch {
'{' => depth += 1,
'}' => {
depth -= 1;
if depth == 0 {
close_brace = Some(open_brace + i);
break;
}
}
_ => {}
}
}
let close_brace = close_brace.expect("if exec.changed block must be balanced");
let invalidate = concat!("invalidate_", "payload_memo(");
let inv_pos = st.find(invalidate).expect(
"the CHANGED streaming full-state apply must call invalidate_payload_memo \
(#4864 round-5 item 6)",
);
assert!(
open_brace < inv_pos && inv_pos < close_brace,
"invalidate_payload_memo ({inv_pos}) must be INSIDE the `if exec.changed \
{{ .. }}` block ({open_brace}..{close_brace}) so a no-op streaming apply \
does not invalidate the memo"
);
}
#[test]
fn nonstreaming_changed_fullstate_success_invalidates_payload_memo() {
let src = include_str!("op_ctx_task.rs");
let invalidate = concat!("invalidate_", "payload_memo(");
for (fn_needle, changed_expr) in [
(
"async fn drive_relay_broadcast_to(",
"else if result.changed",
),
("async fn drive_client_update(", "else if execution.changed"),
] {
let body = extract_fn_body(src, fn_needle);
let gate = body.find(changed_expr).unwrap_or_else(|| {
panic!(
"`{fn_needle}` must have an `{changed_expr}` full-state success arm \
(#4864 round-9 item 3)"
)
});
let open_brace = gate
+ body[gate..]
.find('{')
.expect("else-if block must have a body");
let mut depth = 0usize;
let mut close = None;
for (i, ch) in body[open_brace..].char_indices() {
match ch {
'{' => depth += 1,
'}' => {
depth -= 1;
if depth == 0 {
close = Some(open_brace + i);
break;
}
}
_ => {}
}
}
let close = close.expect("else-if block must be balanced");
let inv = body
.find(invalidate)
.unwrap_or_else(|| panic!("`{fn_needle}` must call invalidate_payload_memo"));
assert!(
open_brace < inv && inv < close,
"`{fn_needle}`: invalidate_payload_memo must be INSIDE the \
`{changed_expr}` full-state block (memo-only, gated on changed)"
);
}
}
#[tokio::test(start_paused = true)]
async fn queue_full_resync_retry_redispatches_within_reservation() {
let (op_manager, mut notification_rx, _guard) =
build_queue_full_test_node("queue-full-resync-4857-retry").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([9u8; 32]),
CodeHash::new([10u8; 32]),
);
let sender_addr: SocketAddr = "127.0.0.1:12300".parse().unwrap();
let r1 = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![1, 2, 3]),
Vec::new(),
sender_addr,
)
.await;
assert!(
r1.as_ref()
.err()
.is_some_and(|e| e.is_contract_queue_full()),
"expected a queue-full error from the stand-in handler, got {r1:?}"
);
assert_eq!(
drain_resync_targets(&mut notification_rx, key),
vec![sender_addr],
"the immediate (pre-retry) ResyncRequest must fire exactly once"
);
let mut retries = Vec::new();
for _ in 0..40 {
if retries.len() >= QUEUE_FULL_RESYNC_MAX_RETRIES as usize {
break;
}
tokio::time::advance(Duration::from_millis(500)).await;
tokio::task::yield_now().await;
retries.extend(drain_resync_targets(&mut notification_rx, key));
}
assert_eq!(
retries.len(),
QUEUE_FULL_RESYNC_MAX_RETRIES as usize,
"the trailing retry must re-dispatch the ResyncRequest exactly \
QUEUE_FULL_RESYNC_MAX_RETRIES times within the single reservation, \
so a bridge-dropped resync heals before the ~5-min heartbeat — and \
must NOT exceed that bound (storm bound #4251)"
);
assert!(
retries.iter().all(|t| *t == sender_addr),
"every retry ResyncRequest must target the original sender, got {retries:?}"
);
}
#[tokio::test(start_paused = true)]
async fn queue_full_resync_retry_stops_after_reservation_deadline() {
let (op_manager, mut notification_rx, _guard) =
build_queue_full_test_node("queue-full-resync-4857-deadline").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([11u8; 32]),
CodeHash::new([12u8; 32]),
);
let sender_addr: SocketAddr = "127.0.0.1:12400".parse().unwrap();
let r1 = drive_relay_broadcast_to(
&op_manager,
Transaction::new::<UpdateMsg>(),
key,
DeltaOrFullState::Delta(vec![1, 2, 3]),
Vec::new(),
sender_addr,
)
.await;
assert!(
r1.as_ref()
.err()
.is_some_and(|e| e.is_contract_queue_full()),
"expected a queue-full error from the stand-in handler, got {r1:?}"
);
assert_eq!(
drain_resync_targets(&mut notification_rx, key),
vec![sender_addr],
"the immediate (pre-retry) ResyncRequest must fire exactly once"
);
tokio::task::yield_now().await;
tokio::time::advance(Duration::from_secs(35)).await;
let mut retries = Vec::new();
for _ in 0..5 {
tokio::task::yield_now().await;
retries.extend(drain_resync_targets(&mut notification_rx, key));
}
assert!(
retries.is_empty(),
"once the injected clock passes the reservation deadline, the retry \
MUST stop and emit no further ResyncRequests (P2-A), got {retries:?}"
);
}
#[tokio::test(start_paused = true)]
async fn queue_full_resync_retry_terminates_when_injected_clock_frozen() {
let (op_manager, mut notification_rx, _guard) = build_queue_full_test_node_with_clock(
"queue-full-resync-4857-frozen-clock",
Some(
Arc::new(crate::util::time_source::SharedMockTimeSource::new())
as crate::util::time_source::DynTimeSource,
),
)
.await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([13u8; 32]),
CodeHash::new([14u8; 32]),
);
let sender_addr: SocketAddr = "127.0.0.1:12500".parse().unwrap();
let reservation_deadline = op_manager.interest_manager.now() + Duration::from_secs(30);
let op_mgr = op_manager.clone();
let handle = tokio::spawn(async move {
resend_queue_full_resync_request(
&op_mgr,
key,
sender_addr,
Transaction::new::<UpdateMsg>(),
reservation_deadline,
)
.await;
});
let terminated = tokio::time::timeout(Duration::from_secs(300), handle).await;
assert!(
terminated.is_ok(),
"the retry task MUST terminate via the tokio liveness backstop when \
the injected clock is frozen — a spin would hang here (DST safety, \
#4857 P2)"
);
terminated.unwrap().expect("retry task must not panic");
assert!(
drain_resync_targets(&mut notification_rx, key).is_empty(),
"a frozen injected clock never crosses a retry target, so no retry \
ResyncRequest should dispatch"
);
}
#[tokio::test]
async fn queue_full_resync_retry_slot_cap_bounds_outstanding_tasks() {
let (op_manager, _rx, _guard) =
build_queue_full_test_node("queue-full-resync-4857-cap").await;
let im = &op_manager.interest_manager;
let cap = crate::ring::interest::MAX_OUTSTANDING_QUEUE_FULL_RESYNC_RETRIES;
let mut slots = Vec::new();
for i in 0..cap {
let slot = im
.try_reserve_resync_retry_slot()
.unwrap_or_else(|| panic!("reservation {i} within the cap must succeed"));
slots.push(slot);
}
assert_eq!(
im.outstanding_resync_retries(),
cap,
"outstanding count must reach exactly the cap"
);
assert!(
im.try_reserve_resync_retry_slot().is_none(),
"at cap, a further retry-slot reservation MUST be refused (#4862 P1)"
);
assert_eq!(
im.outstanding_resync_retries(),
cap,
"a refused reservation must not increment the count (hard cap)"
);
slots.pop();
assert_eq!(
im.outstanding_resync_retries(),
cap - 1,
"dropping a slot guard must free it"
);
let readmit = im.try_reserve_resync_retry_slot();
assert!(
readmit.is_some(),
"a freed slot must re-admit exactly one reservation"
);
assert_eq!(
im.outstanding_resync_retries(),
cap,
"re-admission returns the count to the cap"
);
drop(slots);
drop(readmit);
assert_eq!(
im.outstanding_resync_retries(),
0,
"every slot guard frees its slot on drop"
);
}
#[tokio::test(start_paused = true)]
async fn queue_full_resync_retry_fires_within_a_short_remaining_window() {
let (op_manager, mut notification_rx, _guard) =
build_queue_full_test_node("queue-full-resync-4857-shortwindow").await;
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([15u8; 32]),
CodeHash::new([16u8; 32]),
);
let sender_addr: SocketAddr = "127.0.0.1:12600".parse().unwrap();
let reservation_deadline = op_manager.interest_manager.now() + Duration::from_millis(500);
let op_mgr = op_manager.clone();
let handle = tokio::spawn(async move {
resend_queue_full_resync_request(
&op_mgr,
key,
sender_addr,
Transaction::new::<UpdateMsg>(),
reservation_deadline,
)
.await;
});
let mut fired = Vec::new();
for _ in 0..20 {
tokio::time::advance(Duration::from_millis(50)).await;
tokio::task::yield_now().await;
fired.extend(drain_resync_targets(&mut notification_rx, key));
if !fired.is_empty() {
break;
}
}
let _ = tokio::time::timeout(Duration::from_secs(5), handle).await;
assert!(
!fired.is_empty(),
"with a short remaining window the clamp must still fire ≥1 retry \
within [now, reservation_deadline) (#4862 P2), got none"
);
assert!(
fired.iter().all(|t| *t == sender_addr),
"retry must target the original sender, got {fired:?}"
);
}
#[test]
fn queue_full_resync_retry_is_bounded_and_reservation_scoped() {
let src = include_str!("op_ctx_task.rs");
let fn_body = |name: &str| -> &str {
let start = src.find(name).unwrap_or_else(|| panic!("{name} not found"));
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
&src[start..start + 1 + end]
};
let helper = fn_body("async fn send_queue_full_resync_request(");
let gate = helper
.find(".begin_resync_request(")
.expect("helper must gate on begin_resync_request (per-sender throttle, #4251)");
let global_cap = helper
.find("resync_emit_limiter")
.expect("helper must gate on the global per-contract emit cap (#4864)");
let spawn = helper
.find("resend_queue_full_resync_request(")
.expect("helper must schedule the trailing retry (heal #4857 P2)");
assert!(
gate < spawn && global_cap < spawn,
"the trailing-retry spawn MUST come AFTER both the begin_resync_request \
reservation and the global emit-cap gate, so retries only run within a \
granted, globally-authorized reservation (never on a throttled or \
suppressed drop)"
);
assert!(
helper.contains("GlobalExecutor::spawn("),
"the trailing retry must be a spawned background task"
);
let reserve = helper.find("try_reserve_resync_retry_slot(").expect(
"retry spawn must be gated on interest_manager.try_reserve_resync_retry_slot() (#4862 P1)",
);
assert!(
reserve < spawn,
"the slot reservation must gate the spawn (P1): a spawn without a slot \
would defeat the node-wide retry-task cap"
);
assert!(
helper.contains("let _slot = slot;"),
"the retry slot guard must be moved into the spawned task so it frees \
on task completion (#4862 P1)"
);
let immediate_send = helper
.find(".notify_node_event(")
.expect("helper must perform the immediate (blocking) ResyncRequest enqueue");
assert!(
spawn < immediate_send,
"the retry spawn MUST come BEFORE the blocking .notify_node_event() \
immediate send (#4862 P2), so the retry timer runs during a \
backpressured enqueue"
);
assert!(
helper.contains("reservation_deadline"),
"the helper must thread the reservation deadline from \
begin_resync_request into the retry (anchor to the window, P2-A)"
);
let retry = fn_body("async fn resend_queue_full_resync_request(");
assert!(
!retry.contains("begin_resync_request")
&& !retry.contains("resync_emit_limiter")
&& !retry.contains("outstanding_resync_requests"),
"resend_queue_full_resync_request must NOT re-consult the throttle, the \
global emit cap, or re-record the outstanding request — its sends \
belong to the caller's single granted reservation (storm bound \
#4251 / #4864 cap)"
);
assert!(
retry.contains("in 1..=QUEUE_FULL_RESYNC_MAX_RETRIES"),
"the retry loop must be bounded by QUEUE_FULL_RESYNC_MAX_RETRIES"
);
assert!(
retry.contains("try_notify_node_event"),
"retries must use the NON-blocking try_notify_node_event \
(bounded-channel rule; blocking would pin the task up to 30s)"
);
assert!(
retry.contains("interest_manager.now()"),
"the retry must pace + gate on interest_manager.now() (the throttle's \
injected clock), not tokio's independent clock (P2-B)"
);
assert!(
retry.contains("reservation_deadline"),
"the retry must stop once interest_manager.now() reaches the \
reservation deadline (P2-A)"
);
assert!(
retry.contains("started.elapsed()") && retry.contains("RESYNC_REQUEST_MIN_INTERVAL"),
"the retry must have a tokio-clock liveness backstop \
(started.elapsed() >= RESYNC_REQUEST_MIN_INTERVAL) so a frozen \
injected clock cannot make it spin — while keeping the injected \
reservation deadline as the primary stop"
);
}
#[test]
fn resync_retry_span_stays_within_throttle_window() {
let max = QUEUE_FULL_RESYNC_MAX_RETRIES as u64;
let triangular = max * (max + 1) / 2; let worst_case = QUEUE_FULL_RESYNC_RETRY_BASE_DELAY.mul_f64(triangular as f64 * 1.2);
assert!(
worst_case < Duration::from_secs(30),
"worst-case retry span {worst_case:?} must stay under the 30s \
RESYNC_REQUEST_MIN_INTERVAL reservation window (interest.rs)"
);
}
#[test]
fn jittered_resync_retry_delay_stays_within_bounds() {
for attempt in 1..=QUEUE_FULL_RESYNC_MAX_RETRIES {
let lo = QUEUE_FULL_RESYNC_RETRY_BASE_DELAY.mul_f64(attempt as f64 * 0.8);
let hi = QUEUE_FULL_RESYNC_RETRY_BASE_DELAY.mul_f64(attempt as f64 * 1.2);
for seed in 0..128u64 {
let _guard = GlobalRng::seed_guard(seed);
for _ in 0..16 {
let d = jittered_resync_retry_delay(attempt);
assert!(
d >= lo,
"delay {d:?} below 0.8× for attempt {attempt} (seed {seed})"
);
assert!(
d < hi,
"delay {d:?} at/above 1.2× for attempt {attempt} (seed {seed})"
);
}
}
}
}
#[test]
fn run_relay_wrappers_gate_queue_full_log_severity() {
let src = include_str!("op_ctx_task.rs");
for wrapper in [
"async fn run_relay_request_update(",
"async fn run_relay_broadcast_to(",
"async fn run_relay_request_update_streaming(",
"async fn run_relay_broadcast_to_streaming(",
] {
let start = src
.find(wrapper)
.unwrap_or_else(|| panic!("{wrapper} not found"));
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let body = &src[start..start + 1 + end];
assert!(
body.contains("is_contract_queue_full()"),
"{wrapper} must gate its WARN log on \
err.is_contract_queue_full() — see issue #4251 and PR #4253"
);
assert!(
body.contains("event = \"queue_full\""),
"{wrapper} must tag the DEBUG branch with \
event = \"queue_full\" so log filtering / telemetry can \
distinguish queue-full backpressure from real failures"
);
assert!(
body.contains("tracing::debug!") && body.contains("tracing::warn!"),
"{wrapper} must keep BOTH a debug! (queue_full) and a warn! \
(real failures) call — an inversion that maps queue_full to \
warn would re-open the spam"
);
}
}
#[test]
fn broadcast_to_streaming_classifies_failures() {
let src = include_str!("op_ctx_task.rs");
let start = src
.find("async fn apply_streaming_broadcast(")
.expect("apply_streaming_broadcast not found");
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let driver_src = &src[start..start + 1 + end];
assert!(
driver_src.contains("log_broadcast_to_streaming_failure"),
"apply_streaming_broadcast must call \
log_broadcast_to_streaming_failure for failure classification"
);
assert!(
driver_src.contains("try_auto_fetch_contract"),
"apply_streaming_broadcast must call try_auto_fetch_contract \
on real (non-benign) failures for self-heal"
);
assert!(
driver_src.contains("send_summary_back_on_rejection"),
"apply_streaming_broadcast must spawn \
send_summary_back_on_rejection on is_invalid_update_rejection"
);
}
#[test]
fn broadcast_to_streaming_spawns_proactive_summary() {
let src = include_str!("op_ctx_task.rs");
let start = src
.find("async fn apply_streaming_broadcast(")
.expect("apply_streaming_broadcast not found");
let after = &src[start + 1..];
let end = after
.find("\nasync fn ")
.or_else(|| after.find("\n#[cfg(test)]"))
.unwrap_or(after.len());
let driver_src = &src[start..start + 1 + end];
assert!(
driver_src.contains("send_proactive_summary_notification"),
"apply_streaming_broadcast must spawn \
send_proactive_summary_notification on successful state change"
);
}
#[test]
fn streaming_raii_guard_clears_dedup_and_streaming_counters() {
let src = include_str!("op_ctx_task.rs");
let drop_start = src
.find("impl Drop for RelayUpdateStreamingInflightGuard")
.expect("RelayUpdateStreamingInflightGuard Drop impl not found");
let drop_body = &src[drop_start..drop_start + 700];
assert!(
drop_body.contains("active_relay_update_txs"),
"streaming guard Drop must remove from active_relay_update_txs"
);
assert!(
drop_body.contains("RELAY_UPDATE_STREAMING_INFLIGHT.fetch_sub"),
"streaming guard Drop must decrement RELAY_UPDATE_STREAMING_INFLIGHT"
);
assert!(
drop_body.contains("RELAY_UPDATE_STREAMING_COMPLETED_TOTAL.fetch_add"),
"streaming guard Drop must increment \
RELAY_UPDATE_STREAMING_COMPLETED_TOTAL"
);
}
#[test]
fn log_broadcast_to_streaming_failure_is_pub_crate() {
let src = include_str!("../update.rs");
assert!(
src.contains("pub(crate) fn log_broadcast_to_streaming_failure("),
"log_broadcast_to_streaming_failure must remain pub(crate) — \
reused by apply_streaming_broadcast"
);
}
#[test]
fn deprecated_broadcasting_variant_must_not_return() {
let src = include_str!("../update.rs");
assert!(
!src.contains("Broadcasting {"),
"UpdateMsg::Broadcasting must remain deleted; appending new variants \
at the end of UpdateMsg preserves bincode discriminant tags."
);
}
#[test]
fn deliver_outcome_records_update_op_result() {
const SOURCE: &str = include_str!("op_ctx_task.rs");
let prod = production_source(SOURCE);
let body = extract_fn_body(prod, "fn deliver_outcome(");
assert!(
body.contains("record_op_result"),
"deliver_outcome must call record_op_result so the dashboard \
UPDATE counter advances on driver terminal replies. \
Issue #4010."
);
assert!(
body.contains("OpType::Update"),
"record_op_result inside deliver_outcome must be passed \
OpType::Update (not Get/Put/Subscribe)."
);
assert!(
body.contains("!client_tx.is_sub_operation()"),
"deliver_outcome must gate record_op_result on \
`!client_tx.is_sub_operation()` (note the leading `!`) for \
parity with SUBSCRIBE; defensive against a future change \
that routes a sub-op tx through run_client_update and \
silently inflates the user-facing dashboard counter. \
Issue #4010."
);
}
#[test]
fn classify_update_outcome_covers_all_variants() {
use freenet_stdlib::client_api::{ContractResponse, ErrorKind, HostResponse};
use freenet_stdlib::prelude::{CodeHash, ContractInstanceId, ContractKey, StateSummary};
let key = ContractKey::from_id_and_code(
ContractInstanceId::new([7u8; 32]),
CodeHash::new([8u8; 32]),
);
let publish_ok = DriverOutcome::Publish(Ok(HostResponse::ContractResponse(
ContractResponse::UpdateResponse {
key,
summary: StateSummary::from(Vec::<u8>::new()),
},
)));
let publish_err = DriverOutcome::Publish(Err(ErrorKind::OperationError {
cause: "synthetic".into(),
}
.into()));
let infra = DriverOutcome::InfrastructureError(OpError::UnexpectedOpState);
assert!(classify_update_outcome_for_op_stats(&publish_ok));
assert!(!classify_update_outcome_for_op_stats(&publish_err));
assert!(!classify_update_outcome_for_op_stats(&infra));
}
#[test]
fn drive_client_update_remote_target_auto_fetches_on_missing_params() {
const SOURCE: &str = include_str!("op_ctx_task.rs");
let prod = production_source(SOURCE);
let body = extract_fn_body(prod, "async fn drive_client_update(");
let arm = body
.find("Some(target) =>")
.expect("drive_client_update must have a Some(target) remote-target arm");
let arm_body = &body[arm..];
assert!(
arm_body.contains("is_missing_contract_parameters"),
"remote-target arm must distinguish missing-code/params errors \
via the NARROW `is_missing_contract_parameters` predicate, \
NOT `!is_contract_exec_rejection` — see #4066. Auto-fetching \
on the broader negation triggers unnecessary GETs for \
malformed-input and storage failures."
);
assert!(
!arm_body.contains("!err.is_contract_exec_rejection()"),
"remote-target arm must NOT use the broad \
`!is_contract_exec_rejection` discriminator — see the \
skeptical-review callout on PR #4072 #4. Use the narrow \
`is_missing_contract_parameters` predicate instead."
);
assert!(
arm_body.contains("try_auto_fetch_contract"),
"remote-target arm must call `op_manager.try_auto_fetch_contract` \
when the local apply fails with the missing-params error so the \
contract code/params get pulled from the chosen target peer. \
See #4066."
);
}
#[test]
fn drive_client_update_local_branch_surfaces_ring_state_store_divergence() {
const SOURCE: &str = include_str!("op_ctx_task.rs");
let prod = production_source(SOURCE);
let body = extract_fn_body(prod, "async fn drive_client_update(");
let arm = body
.find("None =>")
.expect("drive_client_update must have a None local-only arm");
let arm_body = &body[arm..];
let arm_end = arm_body[1..]
.find("Some(target) =>")
.map(|p| p + 1)
.unwrap_or(arm_body.len());
let arm_slice = &arm_body[..arm_end];
assert!(
arm_slice.contains("is_missing_contract_parameters"),
"local-only arm must explicitly handle the \
`is_missing_contract_parameters` failure — see #4066. \
Without this, an `is_hosting_contract`/`state_store` \
divergence surfaces as an opaque OpError rather than a \
structured error operators can act on."
);
assert!(
arm_slice.contains("ring_state_store_inconsistency"),
"local-only arm must log the divergence with phase = \
\"ring_state_store_inconsistency\" so the inconsistency \
is searchable in operator dashboards (#4066)"
);
}
fn production_source(full: &str) -> &str {
let cutoff = full
.find("#[cfg(test)]")
.expect("file must have a #[cfg(test)] section");
&full[..cutoff]
}
fn extract_fn_body<'a>(source: &'a str, signature_prefix: &str) -> &'a str {
let start = source
.find(signature_prefix)
.unwrap_or_else(|| panic!("could not find {signature_prefix}"));
let brace = source[start..]
.find('{')
.expect("fn signature must have a body");
let body_start = start + brace + 1;
let bytes = source.as_bytes();
let mut depth: i32 = 1;
let mut i = body_start;
while i < bytes.len() {
match bytes[i] {
b'{' => depth += 1,
b'}' => {
depth -= 1;
if depth == 0 {
return &source[body_start..i];
}
}
_ => {}
}
i += 1;
}
panic!("unterminated fn body for {signature_prefix}");
}
#[test]
fn both_broadcast_drivers_record_update_received() {
let full = include_str!("op_ctx_task.rs");
let src = production_source(full);
for sig in [
"async fn drive_relay_broadcast_to(",
"async fn apply_streaming_broadcast(",
] {
let body = extract_fn_body(src, sig);
let record_calls = body.matches("record_update_received()").count();
assert_eq!(
record_calls, 1,
"{sig} must call record_update_received() EXACTLY once — \
`updates_received` is the only UPDATE number the dashboard \
shows for subscriber nodes, and broadcast_queue.rs picks the \
NON-streaming variant for every payload under \
streaming_threshold (64 KB), i.e. the normal small-delta \
case. Found {record_calls}. Issue #4828."
);
let changed_guard = body
.find("if !changed {")
.unwrap_or_else(|| panic!("{sig} must keep its `if !changed` early-return"));
let guard_open = changed_guard
+ body[changed_guard..]
.find('{')
.expect("the `!changed` guard must have an opening brace");
let guard_end = {
let bytes = body.as_bytes();
let mut depth = 0i32;
let mut end = None;
for (i, &b) in bytes.iter().enumerate().skip(guard_open) {
match b {
b'{' => depth += 1,
b'}' => {
depth -= 1;
if depth == 0 {
end = Some(i);
break;
}
}
_ => {}
}
}
end.expect("the `!changed` guard block must have a matching closing brace")
};
let record_pos = body
.find("record_update_received()")
.expect("checked above");
assert!(
record_pos > guard_end,
"{sig} must call record_update_received() AFTER the entire \
`if !changed` guard block (past its closing brace), so no \
in-block early-return path can record on a no-op broadcast. \
Issue #4828."
);
}
let streaming_driver = extract_fn_body(src, "async fn drive_relay_broadcast_to_streaming(");
let delegations = streaming_driver
.matches("apply_streaming_broadcast(")
.count();
assert_eq!(
delegations, 1,
"drive_relay_broadcast_to_streaming must delegate to \
apply_streaming_broadcast EXACTLY once (found {delegations}); \
that helper owns the streaming record_update_received() call \
(#4857). Zero means the streaming path never records; two would \
double-count. Issue #4828."
);
assert!(
!streaming_driver.contains("record_update_received()"),
"drive_relay_broadcast_to_streaming must NOT call \
record_update_received() directly — the record lives in \
apply_streaming_broadcast (#4857). A direct call here would \
double-count the streaming path. Issue #4828."
);
}
#[test]
fn relayed_update_counter_scope_requests_only() {
let full = include_str!("op_ctx_task.rs");
let src = production_source(full);
for sig in [
"pub(crate) async fn start_relay_request_update(",
"pub(crate) async fn start_relay_request_update_streaming(",
] {
assert!(
extract_fn_body(src, sig).contains("record_relayed_update()"),
"{sig} must record the relayed UPDATE request counter"
);
}
for sig in [
"pub(crate) async fn start_relay_broadcast_to(",
"pub(crate) async fn start_relay_broadcast_to_streaming(",
] {
assert!(
!extract_fn_body(src, sig).contains("record_relayed_update()"),
"{sig} must NOT record relayed_updates_total (fan-out, not a \
relayed request)"
);
}
}
}