#![cfg(feature = "net")]
mod common;
use common::*;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use net::adapter::net::behavior::capability::CapabilitySet;
use net::adapter::net::behavior::sensing::{
encode_attestation, encode_interest_frame, sign_attestation, AttestedStatus,
AudienceScopeCommitment, CanonicalConstraints, CapabilityId, DisclosureClass,
EvaluationRequest, Incarnation, InterestSpec, ProviderInterestKey, ProviderSelector,
ReadinessEvaluation, ReadinessEvaluator, ResultMode, SensingCounters, SensingInterestFrame,
StatusReason, UnsignedAttestation, WorkLatencyEnvelope, SUBPROTOCOL_READINESS_ATTESTATION,
SUBPROTOCOL_SENSING_INTEREST,
};
use net::adapter::net::{EntityKeypair, MeshNode, MeshNodeConfig, SocketBufferConfig};
use net::adapter::Adapter;
const D: Duration = Duration::from_millis(200);
const TTL: Duration = Duration::from_millis(1500);
const LONG_TTL: Duration = Duration::from_secs(10);
const REFRESH: Duration = Duration::from_millis(200);
fn base_config() -> MeshNodeConfig {
let addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
let mut cfg = MeshNodeConfig::new(addr, CHAOS_PSK)
.with_heartbeat_interval(Duration::from_millis(100))
.with_session_timeout(Duration::from_secs(10))
.with_handshake(3, Duration::from_secs(2));
cfg.socket_buffers = SocketBufferConfig {
send_buffer_size: CHAOS_BUFFER_SIZE,
recv_buffer_size: CHAOS_BUFFER_SIZE,
};
cfg
}
fn shared_spec(fleet: AudienceScopeCommitment) -> InterestSpec {
InterestSpec {
capability_id: CapabilityId::new("print.document"),
constraints: CanonicalConstraints::from_entries([("media", "a4")]).unwrap(),
work_latency: WorkLatencyEnvelope::start_within(Duration::from_secs(5)),
providers: ProviderSelector::AnyAuthorized,
result_mode: ResultMode::Any,
disclosure_class: DisclosureClass::Owner,
audience: fleet,
}
}
struct FlagEvaluator {
ready: Arc<AtomicBool>,
}
impl ReadinessEvaluator for FlagEvaluator {
fn evaluate(&self, _request: &EvaluationRequest<'_>) -> ReadinessEvaluation {
if self.ready.load(Ordering::Relaxed) {
ReadinessEvaluation::Ready {
estimated_start: Some(Duration::from_millis(3)),
}
} else {
ReadinessEvaluation::NotReady { reason: 7 }
}
}
}
async fn consumer_provider_pair(
provider_incarnation: Option<Incarnation>,
) -> (Arc<MeshNode>, Arc<MeshNode>, AudienceScopeCommitment) {
let fleet_kp = EntityKeypair::generate();
let fleet = AudienceScopeCommitment::owner_root(fleet_kp.entity_id());
let c = Arc::new(
MeshNode::new(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(fleet),
)
.await
.expect("MeshNode::new C"),
);
let mut p_cfg = base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(fleet);
if let Some(incarnation) = provider_incarnation {
p_cfg = p_cfg.with_sensing_incarnation(incarnation);
}
let p = Arc::new(
MeshNode::new(EntityKeypair::generate(), p_cfg)
.await
.expect("MeshNode::new P"),
);
connect_pair(&c, &p).await;
c.start();
p.start();
for node in [&c, &p] {
node.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
}
let c_id = c.node_id();
let p_id = p.node_id();
await_condition(Duration::from_secs(5), "entity pins established", || {
c.peer_entity_id(p_id).is_some() && p.peer_entity_id(c_id).is_some()
})
.await;
(c, p, fleet)
}
#[tokio::test]
async fn origin_streams_signed_readiness_then_edges_then_drains() {
let (c, p, fleet) = consumer_provider_pair(Some(Incarnation::new(3))).await;
let p_id = p.node_id();
let ready = Arc::new(AtomicBool::new(true));
p.register_readiness_evaluator(
CapabilityId::new("print.document"),
Arc::new(FlagEvaluator {
ready: ready.clone(),
}),
);
assert!(
p.sensing_origin_active(),
"incarnation supplied → origin role active"
);
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
let refresher = {
let c = c.clone();
let spec = spec.clone();
tokio::spawn(async move {
loop {
let _ = c.register_sensing_interest(&spec, p_id, D, TTL);
tokio::time::sleep(REFRESH).await;
}
})
};
await_condition(Duration::from_secs(5), "first admitted attestation", || {
c.sensing_latest_attestation(&branch).is_some()
})
.await;
let first = c.sensing_latest_attestation(&branch).expect("present");
assert_eq!(first.origin, p_id);
assert_eq!(first.origin_incarnation, Incarnation::new(3));
assert_eq!(first.status, AttestedStatus::Ready);
assert_eq!(first.status_reason, StatusReason::None);
assert_eq!(first.estimated_start, Some(Duration::from_millis(3)));
assert_eq!(
first.promised_cadence,
Duration::from_millis(100),
"promised cadence = max(D/2, floor)",
);
assert_eq!(first.audience_scope, fleet);
assert_eq!(p.sensing_live_streams(), 1);
let s0 = c.sensing_latest_attestation(&branch).expect("present").seq;
tokio::time::sleep(Duration::from_millis(450)).await;
let s1 = c.sensing_latest_attestation(&branch).expect("present").seq;
assert!(s1 > s0, "the stream advances ({s0} → {s1})");
assert!(
s1 - s0 <= 10,
"cadence-shaped emission, not a flood ({s0} → {s1})",
);
ready.store(false, Ordering::Relaxed);
p.notify_sensing_state_changed(&CapabilityId::new("print.document"));
await_condition(Duration::from_secs(2), "edge attestation lands", || {
c.sensing_latest_attestation(&branch)
.is_some_and(|a| a.status == AttestedStatus::NotReady)
})
.await;
let edge = c.sensing_latest_attestation(&branch).expect("present");
assert_eq!(edge.status_reason, StatusReason::Provider(7));
refresher.abort();
await_condition(Duration::from_secs(10), "P retires the stream", || {
p.sensing_live_streams() == 0
})
.await;
await_condition(Duration::from_secs(10), "P's table empties", || {
p.sensing_table_is_empty()
})
.await;
await_condition(Duration::from_secs(10), "C reclaims observations", || {
c.sensing_observation_count() == 0
})
.await;
tokio::time::sleep(Duration::from_millis(400)).await;
assert_eq!(p.sensing_live_streams(), 0, "the stream stays retired");
assert_eq!(
c.sensing_observation_count(),
0,
"nothing repopulates a drained hop",
);
for node in [&c, &p] {
let counters = node.sensing_counters();
assert_eq!(SensingCounters::get(&counters.protocol_invalid), 0);
assert_eq!(SensingCounters::get(&counters.scope_refusals), 0);
}
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn below_floor_interest_refused_with_signed_beat() {
let (c, p, fleet) = consumer_provider_pair(Some(Incarnation::new(1))).await;
let p_id = p.node_id();
p.register_readiness_evaluator(
CapabilityId::new("print.document"),
Arc::new(FlagEvaluator {
ready: Arc::new(AtomicBool::new(true)),
}),
);
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
let bad_d = Duration::from_millis(10);
await_condition(Duration::from_secs(10), "refusal beat lands at C", || {
if c.sensing_latest_refusal(&branch).is_none() {
let _ = c.register_sensing_interest(&spec, p_id, bad_d, LONG_TTL);
false
} else {
true
}
})
.await;
let refusal = c.sensing_latest_refusal(&branch).expect("present");
assert_eq!(refusal.status, AttestedStatus::ProviderUnknown);
assert_eq!(
refusal.status_reason,
StatusReason::SamplingIntervalUnsupported,
);
assert_eq!(
refusal.promised_cadence,
Duration::from_millis(50),
"the provider floor M rides promised_cadence",
);
assert_eq!(refusal.estimated_start, None);
assert!(
c.sensing_latest_attestation(&branch).is_none(),
"refusals never enter the warm-start observation store",
);
assert_eq!(p.sensing_live_streams(), 0);
let p_counters = p.sensing_counters();
assert!(SensingCounters::get(&p_counters.cadence_refusals) >= 1);
await_condition(Duration::from_secs(5), "P's table empties", || {
p.sensing_table_is_empty()
})
.await;
await_condition(
Duration::from_secs(5),
"C's sub-floor row partitioned",
|| c.sensing_downstreams(&branch).is_empty(),
)
.await;
let seq = refusal.seq;
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(
c.sensing_latest_refusal(&branch).expect("present").seq,
seq,
"no refusal stream — one signed beat per refused attempt",
);
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn tampered_and_equivocating_attestations_refused() {
let (c, p, fleet) = consumer_provider_pair(None).await;
let c_addr = c.local_addr();
let p_id = p.node_id();
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
let digest = spec.interest_digest();
assert!(!p.sensing_origin_active());
await_condition(Duration::from_secs(10), "row registers at P", || {
if p.sensing_interest_count() == 0 {
let _ = c.register_sensing_interest(&spec, p_id, D, LONG_TTL);
false
} else {
true
}
})
.await;
assert_eq!(
p.sensing_live_streams(),
0,
"fail-closed: no stream without incarnation"
);
let unsigned = UnsignedAttestation {
interest_digest: digest,
origin: p_id,
origin_incarnation: Incarnation::new(9),
capability_id: CapabilityId::new("print.document"),
capability_generation: 1,
status: AttestedStatus::Ready,
status_reason: StatusReason::None,
estimated_start: Some(Duration::from_millis(3)),
seq: 0,
promised_cadence: Duration::from_millis(100),
audience_scope: fleet,
};
let c_counters = c.sensing_counters();
let valid = sign_attestation(p.entity_keypair(), unsigned.clone()).expect("sign");
let mut tampered = valid.clone();
tampered.seq = 6;
let tampered_bytes = encode_attestation(&tampered).expect("encode");
await_condition(Duration::from_secs(10), "tampered frame refused", || {
if SensingCounters::get(&c_counters.protocol_invalid) == 0 {
let p = p.clone();
let bytes = tampered_bytes.clone();
tokio::spawn(async move {
let _ = p
.send_subprotocol(c_addr, SUBPROTOCOL_READINESS_ATTESTATION, &bytes)
.await;
});
false
} else {
true
}
})
.await;
assert!(
c.sensing_latest_attestation(&branch).is_none(),
"a tampered attestation never reaches the observation store",
);
let valid_bytes = encode_attestation(&valid).expect("encode");
await_condition(
Duration::from_secs(10),
"valid attestation admitted",
|| {
if c.sensing_latest_attestation(&branch).is_none() {
let p = p.clone();
let bytes = valid_bytes.clone();
tokio::spawn(async move {
let _ = p
.send_subprotocol(c_addr, SUBPROTOCOL_READINESS_ATTESTATION, &bytes)
.await;
});
false
} else {
true
}
},
)
.await;
let admitted = c.sensing_latest_attestation(&branch).expect("present");
assert_eq!(admitted.seq, 0);
assert_eq!(admitted.estimated_start, Some(Duration::from_millis(3)));
let mut twin = unsigned;
twin.estimated_start = Some(Duration::from_millis(4));
let twin = sign_attestation(p.entity_keypair(), twin).expect("sign twin");
let twin_bytes = encode_attestation(&twin).expect("encode twin");
await_condition(Duration::from_secs(10), "equivocation poisons", || {
if c.sensing_observer_poisoned(p_id, digest).is_none() {
let p = p.clone();
let bytes = twin_bytes.clone();
tokio::spawn(async move {
let _ = p
.send_subprotocol(c_addr, SUBPROTOCOL_READINESS_ATTESTATION, &bytes)
.await;
});
false
} else {
true
}
})
.await;
assert_eq!(
c.sensing_observer_poisoned(p_id, digest),
Some(Incarnation::new(9)),
);
assert_eq!(
c.sensing_latest_attestation(&branch)
.expect("present")
.estimated_start,
Some(Duration::from_millis(3)),
"the equivocating twin never displaces the admitted observation",
);
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn reentrant_evaluator_cannot_deadlock_the_emitter() {
struct ReentrantEvaluator {
node: parking_lot::Mutex<Option<Arc<MeshNode>>>,
}
impl ReadinessEvaluator for ReentrantEvaluator {
fn evaluate(&self, _request: &EvaluationRequest<'_>) -> ReadinessEvaluation {
if let Some(node) = self.node.lock().as_ref() {
node.notify_sensing_state_changed(&CapabilityId::new("print.document"));
let _ = node.sensing_live_streams();
}
ReadinessEvaluation::Ready {
estimated_start: None,
}
}
}
let (c, p, fleet) = consumer_provider_pair(Some(Incarnation::new(1))).await;
let p_id = p.node_id();
let evaluator = Arc::new(ReentrantEvaluator {
node: parking_lot::Mutex::new(None),
});
*evaluator.node.lock() = Some(p.clone());
p.register_readiness_evaluator(CapabilityId::new("print.document"), evaluator);
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
let refresher = {
let c = c.clone();
let spec = spec.clone();
tokio::spawn(async move {
loop {
let _ = c.register_sensing_interest(&spec, p_id, D, TTL);
tokio::time::sleep(REFRESH).await;
}
})
};
await_condition(Duration::from_secs(5), "stream starts", || {
c.sensing_latest_attestation(&branch).is_some()
})
.await;
let s0 = c.sensing_latest_attestation(&branch).expect("present").seq;
tokio::time::sleep(Duration::from_millis(450)).await;
let s1 = c.sensing_latest_attestation(&branch).expect("present").seq;
assert!(
s1 > s0,
"the stream advances despite reentrancy ({s0} → {s1})"
);
assert!(
s1 - s0 <= 15,
"poke-per-beat stays min-gapped at the floor ({s0} → {s1})",
);
refresher.abort();
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn out_of_bounds_intervals_refused_at_intake() {
let (c, p, fleet) = consumer_provider_pair(Some(Incarnation::new(1))).await;
let p_id = p.node_id();
let p_addr = p.local_addr();
p.register_readiness_evaluator(
CapabilityId::new("print.document"),
Arc::new(FlagEvaluator {
ready: Arc::new(AtomicBool::new(true)),
}),
);
let spec = shared_spec(fleet);
for bad in [Duration::ZERO, Duration::from_secs(3600)] {
let refused = c.register_sensing_interest(&spec, p_id, bad, TTL);
assert!(
matches!(
refused,
Err(net::adapter::net::SensingRegistrationError::Interval { .. })
),
"local registration with D={bad:?} must refuse, got {refused:?}",
);
}
assert!(
matches!(
c.register_sensing_interest(&spec, p_id, D, Duration::ZERO),
Err(net::adapter::net::SensingRegistrationError::ZeroTtl)
),
"a zero ttl is dead on arrival",
);
for (bad_d, bad_ttl) in [
(Duration::ZERO, TTL),
(Duration::from_secs(3600), TTL),
(D, Duration::ZERO),
] {
let frame = SensingInterestFrame::provider_registration(&spec, p_id, bad_d, bad_ttl);
let bytes = encode_interest_frame(&frame).expect("encode");
for _ in 0..5 {
let _ = c
.send_subprotocol(p_addr, SUBPROTOCOL_SENSING_INTEREST, &bytes)
.await;
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
assert_eq!(p.sensing_interest_count(), 0, "no row from bad intervals");
assert_eq!(p.sensing_live_streams(), 0, "no stream from bad intervals");
await_condition(Duration::from_secs(10), "legal D registers", || {
if p.sensing_interest_count() == 0 {
let _ = c.register_sensing_interest(&spec, p_id, D, TTL);
false
} else {
true
}
})
.await;
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn cached_floor_invalidates_on_origin_epoch_change() {
let (c, p, fleet) = consumer_provider_pair(None).await;
let c_addr = c.local_addr();
let p_id = p.node_id();
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
let digest = spec.interest_digest();
let beat = |incarnation: u64, generation: u64, seq: u64, refusal: bool| {
let unsigned = UnsignedAttestation {
interest_digest: digest,
origin: p_id,
origin_incarnation: Incarnation::new(incarnation),
capability_id: CapabilityId::new("print.document"),
capability_generation: generation,
status: if refusal {
AttestedStatus::ProviderUnknown
} else {
AttestedStatus::Ready
},
status_reason: if refusal {
StatusReason::SamplingIntervalUnsupported
} else {
StatusReason::None
},
estimated_start: None,
seq,
promised_cadence: Duration::from_millis(50),
audience_scope: fleet,
};
encode_attestation(&sign_attestation(p.entity_keypair(), unsigned).expect("sign"))
.expect("encode")
};
let send_until = |bytes: Vec<u8>, done: Box<dyn Fn() -> bool + Send>, what: &'static str| {
let p = p.clone();
async move {
await_condition(Duration::from_secs(10), what, || {
if done() {
true
} else {
let p = p.clone();
let bytes = bytes.clone();
tokio::spawn(async move {
let _ = p
.send_subprotocol(c_addr, SUBPROTOCOL_READINESS_ATTESTATION, &bytes)
.await;
});
false
}
})
.await;
}
};
await_condition(Duration::from_secs(10), "survivor row at C", || {
if c.sensing_downstreams(&branch).is_empty() {
let _ = c.register_sensing_interest(&spec, p_id, D, LONG_TTL);
false
} else {
true
}
})
.await;
{
let c = c.clone();
let branch = branch.clone();
send_until(
beat(5, 1, 0, true),
Box::new(move || c.sensing_latest_refusal(&branch).is_some()),
"refusal (inc 5) lands",
)
.await;
}
assert!(
matches!(
c.register_sensing_interest(&spec, p_id, Duration::from_millis(10), LONG_TTL),
Ok(net::adapter::net::behavior::sensing::RegisterOutcome::RefusedByCachedFloor { .. })
),
"the cached floor refuses a sub-floor joiner locally",
);
{
let c = c.clone();
let branch = branch.clone();
send_until(
beat(6, 1, 0, false),
Box::new(move || c.sensing_latest_attestation(&branch).is_some()),
"Ready beat (inc 6) lands",
)
.await;
}
assert!(
matches!(
c.register_sensing_interest(&spec, p_id, Duration::from_millis(10), LONG_TTL),
Ok(net::adapter::net::behavior::sensing::RegisterOutcome::Registered(_))
),
"a new incarnation invalidated the cached floor — the sub-floor request goes through again",
);
assert!(c
.register_sensing_interest(&spec, p_id, D, LONG_TTL)
.is_ok());
{
let c = c.clone();
let branch = branch.clone();
send_until(
beat(6, 1, 1, true),
Box::new(move || {
c.sensing_latest_refusal(&branch)
.is_some_and(|r| r.origin_incarnation == Incarnation::new(6))
}),
"refusal (inc 6, gen 1) lands",
)
.await;
}
assert!(
matches!(
c.register_sensing_interest(&spec, p_id, Duration::from_millis(10), LONG_TTL),
Ok(net::adapter::net::behavior::sensing::RegisterOutcome::RefusedByCachedFloor { .. })
),
"the re-cached floor refuses locally again",
);
{
let c = c.clone();
let branch = branch.clone();
send_until(
beat(6, 2, 2, false),
Box::new(move || {
c.sensing_latest_attestation(&branch)
.is_some_and(|a| a.capability_generation == 2)
}),
"Ready beat (gen 2) lands",
)
.await;
}
assert!(
matches!(
c.register_sensing_interest(&spec, p_id, Duration::from_millis(10), LONG_TTL),
Ok(net::adapter::net::behavior::sensing::RegisterOutcome::Registered(_))
),
"a new capability generation invalidated the cached floor too",
);
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn mixed_cadence_refusal_recovers_the_survivor_through_the_relay() {
let fleet_kp = EntityKeypair::generate();
let fleet = AudienceScopeCommitment::owner_root(fleet_kp.entity_id());
let mk = |incarnation: Option<Incarnation>| async move {
let mut cfg = base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(fleet);
if let Some(incarnation) = incarnation {
cfg = cfg.with_sensing_incarnation(incarnation);
}
Arc::new(
MeshNode::new(EntityKeypair::generate(), cfg)
.await
.expect("MeshNode::new"),
)
};
let a = mk(None).await;
let c = mk(None).await;
let r = mk(None).await;
let p = mk(Some(Incarnation::new(1))).await;
p.register_readiness_evaluator(
CapabilityId::new("print.document"),
Arc::new(FlagEvaluator {
ready: Arc::new(AtomicBool::new(true)),
}),
);
connect_pair(&a, &r).await;
connect_pair(&c, &r).await;
connect_pair(&r, &p).await;
a.start();
c.start();
r.start();
p.start();
for node in [&a, &c, &r, &p] {
node.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
}
let a_id = a.node_id();
let c_id = c.node_id();
let r_id = r.node_id();
let p_id = p.node_id();
await_condition(Duration::from_secs(5), "entity pins established", || {
r.peer_entity_id(a_id).is_some()
&& r.peer_entity_id(c_id).is_some()
&& r.peer_entity_id(p_id).is_some()
&& p.peer_entity_id(r_id).is_some()
})
.await;
a.router().add_route(p_id, r.local_addr());
c.router().add_route(p_id, r.local_addr());
let p_entity = p.entity_keypair().entity_id().clone();
a.test_pin_peer_entity(p_id, p_entity.clone());
c.test_pin_peer_entity(p_id, p_entity);
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
use net::adapter::net::behavior::sensing::DownstreamId;
let refresher_a = {
let a = a.clone();
let spec = spec.clone();
tokio::spawn(async move {
loop {
let _ = a.register_sensing_interest(&spec, p_id, Duration::from_millis(100), TTL);
tokio::time::sleep(REFRESH).await;
}
})
};
await_condition(Duration::from_secs(10), "P streams to R", || {
p.sensing_live_streams() == 1 && r.sensing_latest_attestation(&branch).is_some()
})
.await;
assert_eq!(
r.sensing_downstreams(&branch),
vec![DownstreamId::Peer(a_id)],
);
await_condition(Duration::from_secs(10), "refusal lands at R", || {
if r.sensing_latest_refusal(&branch).is_none() {
let _ = c.register_sensing_interest(&spec, p_id, Duration::from_millis(10), TTL);
false
} else {
true
}
})
.await;
let refusal = r.sensing_latest_refusal(&branch).expect("present");
assert_eq!(
refusal.status_reason,
StatusReason::SamplingIntervalUnsupported
);
await_condition(Duration::from_secs(5), "C partitioned out at R", || {
r.sensing_downstreams(&branch) == vec![DownstreamId::Peer(a_id)]
})
.await;
await_condition(Duration::from_secs(5), "forwarded refusal at C", || {
c.sensing_latest_refusal(&branch).is_some()
})
.await;
await_condition(
Duration::from_secs(10),
"P re-registers the survivor",
|| {
p.sensing_live_streams() == 1
&& p.sensing_downstream_entry(&branch, DownstreamId::Peer(r_id))
.is_some_and(|row| row.requested_sample_interval == Duration::from_millis(100))
},
)
.await;
await_condition(Duration::from_secs(10), "fresh beats at R", || {
r.sensing_latest_attestation(&branch)
.is_some_and(|beat| beat.status == AttestedStatus::Ready && beat.seq > refusal.seq)
})
.await;
for node in [&a, &c, &r, &p] {
let counters = node.sensing_counters();
assert_eq!(SensingCounters::get(&counters.protocol_invalid), 0);
}
refresher_a.abort();
a.shutdown().await.expect("shutdown A");
c.shutdown().await.expect("shutdown C");
r.shutdown().await.expect("shutdown R");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn malformed_refusals_and_stale_epochs_touch_nothing() {
let (c, p, fleet) = consumer_provider_pair(None).await;
let c_addr = c.local_addr();
let p_id = p.node_id();
let spec_a = shared_spec(fleet);
let branch_a = ProviderInterestKey::new(spec_a.key(), p_id);
let mut spec_b = shared_spec(fleet);
spec_b.constraints = CanonicalConstraints::from_entries([("media", "letter")]).unwrap();
let branch_b = ProviderInterestKey::new(spec_b.key(), p_id);
let beat = |spec: &InterestSpec,
incarnation: u64,
seq: u64,
refusal: bool,
floor_ms: u64,
start_ms: u64| {
let unsigned = UnsignedAttestation {
interest_digest: spec.interest_digest(),
origin: p_id,
origin_incarnation: Incarnation::new(incarnation),
capability_id: CapabilityId::new("print.document"),
capability_generation: 1,
status: if refusal {
AttestedStatus::ProviderUnknown
} else {
AttestedStatus::Ready
},
status_reason: if refusal {
StatusReason::SamplingIntervalUnsupported
} else {
StatusReason::None
},
estimated_start: (!refusal).then(|| Duration::from_millis(start_ms)),
seq,
promised_cadence: Duration::from_millis(floor_ms),
audience_scope: fleet,
};
encode_attestation(&sign_attestation(p.entity_keypair(), unsigned).expect("sign"))
.expect("encode")
};
macro_rules! send_until {
($bytes:expr, $what:literal, $done:expr) => {
let bytes = $bytes;
await_condition(Duration::from_secs(10), $what, || {
if $done {
true
} else {
let p = p.clone();
let bytes = bytes.clone();
tokio::spawn(async move {
let _ = p
.send_subprotocol(c_addr, SUBPROTOCOL_READINESS_ATTESTATION, &bytes)
.await;
});
false
}
})
.await;
};
}
let floor_cached = || {
matches!(
c.register_sensing_interest(&spec_a, p_id, Duration::from_millis(10), LONG_TTL),
Ok(net::adapter::net::behavior::sensing::RegisterOutcome::RefusedByCachedFloor { .. })
)
};
await_condition(Duration::from_secs(10), "survivor row at C", || {
if c.sensing_downstreams(&branch_a).is_empty() {
let _ = c.register_sensing_interest(&spec_a, p_id, D, LONG_TTL);
false
} else {
true
}
})
.await;
send_until!(
beat(&spec_a, 6, 0, false, 100, 3),
"Ready (inc 6) admits",
c.sensing_latest_attestation(&branch_a).is_some()
);
send_until!(beat(&spec_a, 6, 1, true, 50, 0), "refusal caches floor", {
c.sensing_latest_refusal(&branch_a).is_some()
});
assert!(floor_cached(), "the floor refuses a sub-floor joiner");
let invalid_before = SensingCounters::get(&c.sensing_counters().protocol_invalid);
send_until!(
beat(&spec_a, 7, 2, true, 0, 0),
"malformed refusal counted",
SensingCounters::get(&c.sensing_counters().protocol_invalid) > invalid_before
);
assert!(
floor_cached(),
"a malformed inc-7 refusal must not move the epoch or flush the floor",
);
send_until!(
beat(&spec_a, 6, 2, false, 100, 9),
"inc-6 seq-2 still admits",
{
c.sensing_latest_attestation(&branch_a)
.is_some_and(|a| a.estimated_start == Some(Duration::from_millis(9)))
}
);
await_condition(Duration::from_secs(10), "row on branch B", || {
if c.sensing_downstreams(&branch_b).is_empty() {
let _ = c.register_sensing_interest(&spec_b, p_id, D, LONG_TTL);
false
} else {
true
}
})
.await;
let superseded_before = SensingCounters::get(&c.sensing_counters().attestations_superseded);
send_until!(
beat(&spec_b, 5, 0, false, 100, 3),
"stale-epoch beat dropped at the provider-wide epoch gate",
SensingCounters::get(&c.sensing_counters().attestations_superseded) > superseded_before
);
assert!(
c.sensing_latest_attestation(&branch_b).is_none(),
"a globally stale epoch never becomes an observation (SI-5R P0)",
);
assert!(
floor_cached(),
"a cross-digest stale epoch must not flush valid floors",
);
send_until!(beat(&spec_b, 7, 1, false, 100, 3), "inc 7 advances", {
c.sensing_latest_attestation(&branch_b)
.is_some_and(|a| a.origin_incarnation == Incarnation::new(7))
});
await_condition(Duration::from_secs(5), "floor invalidated by inc 7", || {
matches!(
c.register_sensing_interest(&spec_a, p_id, Duration::from_millis(10), LONG_TTL),
Ok(net::adapter::net::behavior::sensing::RegisterOutcome::Registered(_))
)
})
.await;
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn short_ttl_rows_survive_the_damper() {
let (c, p, fleet) = consumer_provider_pair(Some(Incarnation::new(1))).await;
let c_id = c.node_id();
let p_id = p.node_id();
p.register_readiness_evaluator(
CapabilityId::new("print.document"),
Arc::new(FlagEvaluator {
ready: Arc::new(AtomicBool::new(true)),
}),
);
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
let short_ttl = Duration::from_millis(100);
let refresher = {
let c = c.clone();
let spec = spec.clone();
tokio::spawn(async move {
loop {
let _ = c.register_sensing_interest(&spec, p_id, short_ttl, short_ttl);
tokio::time::sleep(short_ttl / 2).await;
}
})
};
use net::adapter::net::behavior::sensing::DownstreamId;
await_condition(Duration::from_secs(10), "upstream row established", || {
p.sensing_downstream_entry(&branch, DownstreamId::Peer(c_id))
.is_some()
})
.await;
for _ in 0..20 {
tokio::time::sleep(Duration::from_millis(25)).await;
assert!(
p.sensing_downstream_entry(&branch, DownstreamId::Peer(c_id))
.is_some(),
"the upstream row must never lapse while refreshes flow at ttl/2",
);
}
assert_eq!(p.sensing_live_streams(), 1, "the stream rode along");
refresher.abort();
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn tombstone_ageout_drains_provider_epochs_too() {
let fleet_kp = EntityKeypair::generate();
let fleet = AudienceScopeCommitment::owner_root(fleet_kp.entity_id());
let c = Arc::new(
MeshNode::new(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(fleet)
.with_sensing_interest_ttl(Duration::from_secs(1)),
)
.await
.expect("MeshNode::new C"),
);
let p = Arc::new(
MeshNode::new(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(fleet),
)
.await
.expect("MeshNode::new P"),
);
connect_pair(&c, &p).await;
c.start();
p.start();
for node in [&c, &p] {
node.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
}
let p_id = p.node_id();
let c_id = c.node_id();
let c_addr = c.local_addr();
await_condition(Duration::from_secs(5), "entity pins established", || {
c.peer_entity_id(p_id).is_some() && p.peer_entity_id(c_id).is_some()
})
.await;
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
let unsigned = UnsignedAttestation {
interest_digest: spec.interest_digest(),
origin: p_id,
origin_incarnation: Incarnation::new(3),
capability_id: CapabilityId::new("print.document"),
capability_generation: 1,
status: AttestedStatus::ProviderUnknown,
status_reason: StatusReason::SamplingIntervalUnsupported,
estimated_start: None,
seq: 0,
promised_cadence: Duration::from_millis(50),
audience_scope: fleet,
};
let refusal_bytes =
encode_attestation(&sign_attestation(p.entity_keypair(), unsigned).expect("sign"))
.expect("encode");
await_condition(Duration::from_secs(10), "refusal tombstone lands", || {
if c.sensing_latest_refusal(&branch).is_some() {
return true;
}
let _ = c.register_sensing_interest(
&spec,
p_id,
Duration::from_millis(10),
Duration::from_secs(1),
);
let p = p.clone();
let bytes = refusal_bytes.clone();
tokio::spawn(async move {
let _ = p
.send_subprotocol(c_addr, SUBPROTOCOL_READINESS_ATTESTATION, &bytes)
.await;
});
false
})
.await;
assert_eq!(c.sensing_provider_epoch_count(), 1, "epoch recorded");
await_condition(Duration::from_secs(10), "tombstone and epoch drain", || {
c.sensing_latest_refusal(&branch).is_none() && c.sensing_provider_epoch_count() == 0
})
.await;
assert_eq!(c.sensing_observation_count(), 0);
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn at_capacity_registration_surfaced_and_rolled_back() {
use net::adapter::net::behavior::sensing::MAX_LIVE_SENSING_STREAMS;
use net::adapter::net::SensingRegistrationError;
let node = Arc::new(
MeshNode::new(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_incarnation(Incarnation::new(1))
.with_max_interests_per_peer(2 * MAX_LIVE_SENSING_STREAMS),
)
.await
.expect("MeshNode::new"),
);
let self_id = node.node_id();
let own_root = AudienceScopeCommitment::owner_root(node.entity_keypair().entity_id());
let spec_n = |i: usize| {
let mut spec = shared_spec(own_root);
spec.constraints =
CanonicalConstraints::from_entries([("slot", format!("{i}").as_str())]).unwrap();
spec
};
for i in 0..MAX_LIVE_SENSING_STREAMS {
node.register_sensing_interest(&spec_n(i), self_id, D, LONG_TTL)
.expect("registers under the cap");
}
assert_eq!(node.sensing_live_streams(), MAX_LIVE_SENSING_STREAMS);
let overflow = spec_n(MAX_LIVE_SENSING_STREAMS);
let refused = node.register_sensing_interest(&overflow, self_id, D, LONG_TTL);
assert!(
matches!(refused, Err(SensingRegistrationError::AtCapacity)),
"expected AtCapacity, got {refused:?}",
);
let overflow_branch = ProviderInterestKey::new(overflow.key(), self_id);
assert!(
node.sensing_downstreams(&overflow_branch).is_empty(),
"the refused registration's row is rolled back, never dark",
);
assert!(
node.sensing_consumer_cell_interval_for_test(&overflow_branch)
.is_none(),
"the eagerly-created consumer cell is rolled back with the row",
);
assert_eq!(node.sensing_live_streams(), MAX_LIVE_SENSING_STREAMS);
node.register_sensing_interest(&spec_n(0), self_id, D, LONG_TTL)
.expect("live refresh at capacity");
node.shutdown().await.expect("shutdown");
}
#[tokio::test]
async fn interval_changes_re_anchor_windows_with_no_intervening_beat() {
let (c, p, fleet) = consumer_provider_pair(Some(Incarnation::new(9))).await;
let p_id = p.node_id();
let ready = Arc::new(AtomicBool::new(true));
p.register_readiness_evaluator(
CapabilityId::new("print.document"),
Arc::new(FlagEvaluator {
ready: ready.clone(),
}),
);
let spec_tight = shared_spec(fleet);
let spec_loose = {
let mut spec = shared_spec(fleet);
spec.constraints = CanonicalConstraints::from_entries([("media", "letter")]).unwrap();
spec
};
let branch_tight = ProviderInterestKey::new(spec_tight.key(), p_id);
let branch_loose = ProviderInterestKey::new(spec_loose.key(), p_id);
let d0 = Duration::from_millis(400);
c.register_sensing_interest(&spec_tight, p_id, d0, LONG_TTL)
.expect("register tight");
c.register_sensing_interest(&spec_loose, p_id, d0, LONG_TTL)
.expect("register loose");
await_condition(Duration::from_secs(5), "both branches established", || {
use net::adapter::net::behavior::sensing::Continuity;
c.sensing_upstream_continuity(&branch_tight) == Some(Continuity::Established)
&& c.sensing_upstream_continuity(&branch_loose) == Some(Continuity::Established)
})
.await;
chaos_partition(&c, &p);
c.register_sensing_interest(&spec_tight, p_id, Duration::from_millis(100), LONG_TTL)
.expect("tighten");
c.register_sensing_interest(&spec_loose, p_id, Duration::from_millis(1200), LONG_TTL)
.expect("loosen");
let changed_at = std::time::Instant::now();
await_condition(
Duration::from_millis(700),
"tightened window expires early",
|| {
use net::adapter::net::behavior::sensing::Continuity;
c.sensing_upstream_continuity(&branch_tight) == Some(Continuity::Expired)
},
)
.await;
tokio::time::sleep(Duration::from_millis(1500).saturating_sub(changed_at.elapsed())).await;
{
use net::adapter::net::behavior::sensing::Continuity;
assert_eq!(
c.sensing_upstream_continuity(&branch_loose),
Some(Continuity::Established),
"loosened window must move the deadline outward immediately",
);
}
await_condition(
Duration::from_secs(4),
"loosened window expires eventually",
|| {
use net::adapter::net::behavior::sensing::Continuity;
c.sensing_upstream_continuity(&branch_loose) == Some(Continuity::Expired)
},
)
.await;
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}