#![cfg(feature = "net")]
mod common;
use common::*;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use net::adapter::net::behavior::capability::CapabilitySet;
use net::adapter::net::behavior::sensing::{
AudienceScopeCommitment, CanonicalConstraints, CapabilityId, DisclosureClass, DownstreamId,
InterestSpec, ProviderInterestKey, ProviderSelector, RegisterOutcome, ResultMode,
SensingCounters, WorkLatencyEnvelope,
};
use net::adapter::net::{
EntityKeypair, MeshNode, MeshNodeConfig, SensingRegistrationError, SocketBufferConfig,
};
use net::adapter::Adapter;
const TTL: Duration = Duration::from_secs(30);
const D: Duration = Duration::from_millis(100);
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
}
async fn build_with_keypair(keypair: EntityKeypair, cfg: MeshNodeConfig) -> Arc<MeshNode> {
Arc::new(MeshNode::new(keypair, cfg).await.expect("MeshNode::new"))
}
async fn bring_up(a: &Arc<MeshNode>, b: &Arc<MeshNode>) {
connect_pair(a, b).await;
a.start();
b.start();
a.announce_capabilities(CapabilitySet::new())
.await
.expect("A announce");
b.announce_capabilities(CapabilitySet::new())
.await
.expect("B announce");
let a_id = a.node_id();
let b_id = b.node_id();
await_condition(Duration::from_secs(5), "entity pins established", || {
a.peer_entity_id(b_id).is_some() && b.peer_entity_id(a_id).is_some()
})
.await;
}
fn spec_for(owner: AudienceScopeCommitment, provider: u64, marker: &str) -> InterestSpec {
InterestSpec {
capability_id: CapabilityId::new("print.document"),
constraints: CanonicalConstraints::from_entries([("media", marker)]).unwrap(),
work_latency: WorkLatencyEnvelope::start_within(Duration::from_secs(5)),
providers: ProviderSelector::Node(provider),
result_mode: ResultMode::Any,
disclosure_class: DisclosureClass::Owner,
audience: owner,
}
}
fn counter_snapshot(counters: &SensingCounters) -> [u64; 4] {
[
SensingCounters::get(&counters.invalid_constraints),
SensingCounters::get(&counters.protocol_invalid),
SensingCounters::get(&counters.cadence_refusals),
SensingCounters::get(&counters.scope_refusals),
]
}
#[tokio::test]
async fn provider_registration_lands_a_peer_row_with_proven_root() {
let a_kp = EntityKeypair::generate();
let owner = AudienceScopeCommitment::owner_root(a_kp.entity_id());
let a = build_with_keypair(a_kp, base_config().with_sensing_coalescing(true)).await;
let b = build_with_keypair(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(owner),
)
.await;
bring_up(&a, &b).await;
let a_id = a.node_id();
let b_id = b.node_id();
let spec = spec_for(owner, b_id, "a4");
let key = ProviderInterestKey::new(spec.key(), b_id);
let mut landed = false;
for _ in 0..40 {
let outcome = a
.register_sensing_interest(&spec, b_id, D, TTL)
.expect("A's local registration is in-scope");
assert!(matches!(outcome, RegisterOutcome::Registered(_)));
if poll_until(Duration::from_millis(250), || {
!b.sensing_downstreams(&key).is_empty()
})
.await
{
landed = true;
break;
}
}
assert!(landed, "B never gained a row for A's interest");
assert_eq!(b.sensing_downstreams(&key), vec![DownstreamId::Peer(a_id)]);
let row = b
.sensing_downstream_entry(&key, DownstreamId::Peer(a_id))
.expect("row present");
assert_eq!(row.owner_root, owner, "the proven root is stored");
assert_eq!(row.requested_sample_interval, D);
assert_eq!(b.sensing_interest_count(), 1, "exactly one branch key");
assert_eq!(a.sensing_downstreams(&key), vec![DownstreamId::Local]);
let counters = b.sensing_counters();
assert_eq!(SensingCounters::get(&counters.protocol_invalid), 0);
assert_eq!(SensingCounters::get(&counters.scope_refusals), 0);
a.shutdown().await.expect("shutdown A");
b.shutdown().await.expect("shutdown B");
}
#[tokio::test]
async fn disabled_receiver_is_inert_and_moves_no_counters() {
let addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
assert!(
!MeshNodeConfig::new(addr, CHAOS_PSK).enable_sensing_coalescing,
"enable_sensing_coalescing must default to false",
);
let a_kp = EntityKeypair::generate();
let owner = AudienceScopeCommitment::owner_root(a_kp.entity_id());
let a = build_with_keypair(a_kp, base_config().with_sensing_coalescing(true)).await;
let b = build_with_keypair(EntityKeypair::generate(), base_config()).await;
bring_up(&a, &b).await;
let b_id = b.node_id();
let spec = spec_for(owner, b_id, "a4");
let dark = build_with_keypair(EntityKeypair::generate(), base_config()).await;
assert_eq!(
dark.register_sensing_interest(&spec, b_id, D, TTL),
Err(SensingRegistrationError::Disabled),
);
for _ in 0..10 {
a.register_sensing_interest(&spec, b_id, D, TTL)
.expect("A registers");
tokio::time::sleep(Duration::from_millis(100)).await;
}
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
b.sensing_table_is_empty(),
"a dark receiver must never gain sensing rows",
);
assert_eq!(
counter_snapshot(&b.sensing_counters()),
[0, 0, 0, 0],
"a dark receiver must move zero sensing counters",
);
assert_eq!(b.sensing_over_cap_refusals(), 0);
a.shutdown().await.expect("shutdown A");
b.shutdown().await.expect("shutdown B");
dark.shutdown().await.expect("shutdown dark");
}
#[tokio::test]
async fn unbacked_scope_claim_is_rejected_as_protocol_invalid() {
let a_kp = EntityKeypair::generate();
let a_session_root = AudienceScopeCommitment::owner_root(a_kp.entity_id());
let foreign = AudienceScopeCommitment::owner_root(EntityKeypair::generate().entity_id());
let a = build_with_keypair(
a_kp,
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(foreign),
)
.await;
let b = build_with_keypair(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(a_session_root),
)
.await;
bring_up(&a, &b).await;
let b_id = b.node_id();
let spec = spec_for(foreign, b_id, "a4");
let counters = b.sensing_counters();
let mut rejected = false;
for _ in 0..40 {
a.register_sensing_interest(&spec, b_id, D, TTL)
.expect("locally in-scope at A");
if poll_until(Duration::from_millis(250), || {
SensingCounters::get(&counters.protocol_invalid) >= 1
})
.await
{
rejected = true;
break;
}
}
assert!(rejected, "B never rejected the unbacked scope claim");
assert!(
SensingCounters::get(&counters.scope_refusals) >= 1,
"scope refusal counter moved",
);
assert!(
b.sensing_table_is_empty(),
"a rejected registration must record nothing",
);
a.shutdown().await.expect("shutdown A");
b.shutdown().await.expect("shutdown B");
}
#[tokio::test]
async fn per_peer_cap_bounds_inbound_interests() {
const CAP: usize = 4;
let a_kp = EntityKeypair::generate();
let owner = AudienceScopeCommitment::owner_root(a_kp.entity_id());
let a = build_with_keypair(a_kp, base_config().with_sensing_coalescing(true)).await;
let b = build_with_keypair(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(owner)
.with_max_interests_per_peer(CAP),
)
.await;
bring_up(&a, &b).await;
let b_id = b.node_id();
let specs: Vec<InterestSpec> = (0..=CAP)
.map(|i| spec_for(owner, b_id, &format!("m{i}")))
.collect();
let mut converged = false;
for _ in 0..40 {
for spec in &specs {
a.register_sensing_interest(spec, b_id, D, TTL)
.expect("in-scope at A");
}
if poll_until(Duration::from_millis(250), || {
b.sensing_interest_count() == CAP && b.sensing_over_cap_refusals() >= 1
})
.await
{
converged = true;
break;
}
}
assert!(
converged,
"B never converged to cap rows + a surfaced OverCap refusal \
(rows: {}, over-cap: {})",
b.sensing_interest_count(),
b.sensing_over_cap_refusals(),
);
for _ in 0..5 {
for spec in &specs {
a.register_sensing_interest(spec, b_id, D, TTL)
.expect("in-scope at A");
}
tokio::time::sleep(Duration::from_millis(150)).await;
assert_eq!(
b.sensing_interest_count(),
CAP,
"row count must never exceed the per-peer cap",
);
}
let present = specs
.iter()
.filter(|spec| {
!b.sensing_downstreams(&ProviderInterestKey::new(spec.key(), b_id))
.is_empty()
})
.count();
assert_eq!(present, CAP, "exactly cap keys admitted");
a.shutdown().await.expect("shutdown A");
b.shutdown().await.expect("shutdown B");
}
#[tokio::test]
async fn short_ttl_rows_are_swept_on_the_heartbeat_tick() {
const SHORT_TTL: Duration = Duration::from_millis(300);
let a_kp = EntityKeypair::generate();
let owner = AudienceScopeCommitment::owner_root(a_kp.entity_id());
let a = build_with_keypair(a_kp, base_config().with_sensing_coalescing(true)).await;
let b = build_with_keypair(
EntityKeypair::generate(),
base_config()
.with_sensing_coalescing(true)
.with_sensing_owner_root(owner),
)
.await;
bring_up(&a, &b).await;
let b_id = b.node_id();
let spec = spec_for(owner, b_id, "a4");
let key = ProviderInterestKey::new(spec.key(), b_id);
let mut landed = false;
for _ in 0..40 {
a.register_sensing_interest(&spec, b_id, D, SHORT_TTL)
.expect("in-scope at A");
if poll_until(Duration::from_millis(200), || {
!b.sensing_downstreams(&key).is_empty()
})
.await
{
landed = true;
break;
}
}
assert!(landed, "precondition: the row never landed at B");
assert!(
poll_until(Duration::from_secs(5), || b.sensing_table_is_empty()).await,
"B's sweep never expired the unrefreshed row",
);
assert!(
poll_until(Duration::from_secs(5), || a.sensing_table_is_empty()).await,
"A's sweep never expired its LOCAL row",
);
a.shutdown().await.expect("shutdown A");
b.shutdown().await.expect("shutdown B");
}