#![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::{
encode_attestation, sign_attestation, AttestedStatus, AudienceScopeCommitment,
CanonicalConstraints, CapabilityId, Continuity, DisclosureClass, EvaluationRequest,
Incarnation, InterestSpec, ProjectedReadiness, ProviderInterestKey, ProviderSelector,
ReadinessEvaluation, ReadinessEvaluator, ResultMode, StatusReason, UnsignedAttestation,
WorkLatencyEnvelope, SUBPROTOCOL_READINESS_ATTESTATION,
};
use net::adapter::net::{EntityKeypair, MeshNode, MeshNodeConfig, SocketBufferConfig};
use net::adapter::Adapter;
const LONG_TTL: Duration = Duration::from_secs(10);
const D: Duration = Duration::from_secs(2);
const FD_TIMEOUT: Duration = Duration::from_millis(300);
const NO_FD: Duration = Duration::from_secs(10);
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_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 AlwaysReady;
impl ReadinessEvaluator for AlwaysReady {
fn evaluate(&self, _request: &EvaluationRequest<'_>) -> ReadinessEvaluation {
ReadinessEvaluation::Ready {
estimated_start: Some(Duration::from_millis(3)),
}
}
}
async fn sensing_node(
kp: EntityKeypair,
fleet: AudienceScopeCommitment,
incarnation: Option<Incarnation>,
session_timeout: Duration,
) -> Arc<MeshNode> {
let mut cfg = base_config()
.with_session_timeout(session_timeout)
.with_sensing_coalescing(true)
.with_sensing_owner_root(fleet);
if let Some(incarnation) = incarnation {
cfg = cfg.with_sensing_incarnation(incarnation);
}
Arc::new(MeshNode::new(kp, cfg).await.expect("MeshNode::new"))
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn provider_failure_expires_observations_and_recovery_re_establishes() {
let fleet_kp = EntityKeypair::generate();
let fleet = AudienceScopeCommitment::owner_root(fleet_kp.entity_id());
let c = sensing_node(EntityKeypair::generate(), fleet, None, FD_TIMEOUT).await;
let p = sensing_node(
EntityKeypair::generate(),
fleet,
Some(Incarnation::new(1)),
FD_TIMEOUT,
)
.await;
p.register_readiness_evaluator(CapabilityId::new("print.document"), Arc::new(AlwaysReady));
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, p_id) = (c.node_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;
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
c.register_sensing_interest(&spec, p_id, D, LONG_TTL)
.expect("register");
await_condition(Duration::from_secs(10), "established + Ready", || {
c.sensing_upstream_continuity(&branch) == Some(Continuity::Established)
&& c.sensing_projected(&branch) == ProjectedReadiness::Ready
})
.await;
let overlay = c.subscribe_sensing_overlay_changes();
let generation_before = *overlay.borrow();
let partitioned_at = std::time::Instant::now();
chaos_partition(&c, &p);
await_condition(
Duration::from_secs(4),
"observations expire on the failure edge",
|| {
c.sensing_upstream_continuity(&branch) == Some(Continuity::Expired)
&& c.sensing_projected(&branch) == ProjectedReadiness::Unknown
},
)
.await;
assert!(
partitioned_at.elapsed() < Duration::from_secs(5),
"disruption must be event-driven",
);
assert!(
*overlay.borrow() > generation_before,
"a disappearing projection fires the overlay signal",
);
await_condition(
Duration::from_secs(4),
"origin retires the dead consumer's stream",
|| p.sensing_live_streams() == 0,
)
.await;
chaos_heal(&c, &p);
await_peer_recovered(&c, p_id, Duration::from_secs(10)).await;
await_peer_recovered(&p, c_id, Duration::from_secs(10)).await;
await_condition(Duration::from_secs(10), "pins re-established", || {
if c.peer_entity_id(p_id).is_some() && p.peer_entity_id(c_id).is_some() {
return true;
}
let (c, p) = (c.clone(), p.clone());
tokio::spawn(async move {
let _ = c.announce_capabilities(CapabilitySet::new()).await;
let _ = p.announce_capabilities(CapabilitySet::new()).await;
});
false
})
.await;
await_condition(Duration::from_secs(15), "recovery re-establishes", || {
let _ = c.register_sensing_interest(&spec, p_id, D, LONG_TTL);
c.sensing_projected(&branch) == ProjectedReadiness::Ready
})
.await;
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn next_hop_failure_disrupts_multi_hop_provider_branches() {
let fleet_kp = EntityKeypair::generate();
let fleet = AudienceScopeCommitment::owner_root(fleet_kp.entity_id());
let c = sensing_node(EntityKeypair::generate(), fleet, None, FD_TIMEOUT).await;
let r = sensing_node(EntityKeypair::generate(), fleet, None, FD_TIMEOUT).await;
let p = sensing_node(
EntityKeypair::generate(),
fleet,
Some(Incarnation::new(1)),
FD_TIMEOUT,
)
.await;
p.register_readiness_evaluator(CapabilityId::new("print.document"), Arc::new(AlwaysReady));
connect_pair(&c, &r).await;
connect_pair(&r, &p).await;
for node in [&c, &r, &p] {
node.start();
}
for node in [&c, &r, &p] {
node.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
}
let (c_id, r_id, p_id) = (c.node_id(), r.node_id(), p.node_id());
await_condition(Duration::from_secs(5), "entity pins established", || {
c.peer_entity_id(r_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;
let p_entity = p.entity_keypair().entity_id().clone();
c.test_pin_peer_entity(p_id, p_entity);
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
await_condition(Duration::from_secs(5), "C routes to P via R", || {
c.register_sensing_interest(&spec, p_id, D, LONG_TTL)
.is_ok()
&& c.sensing_upstream_continuity(&branch) == Some(Continuity::Established)
&& c.sensing_projected(&branch) == ProjectedReadiness::Ready
})
.await;
let partitioned_at = std::time::Instant::now();
chaos_partition(&c, &r);
await_condition(
Duration::from_secs(4),
"next_hop failure expires the provider's observations",
|| {
c.sensing_upstream_continuity(&branch) == Some(Continuity::Expired)
&& c.sensing_projected(&branch) == ProjectedReadiness::Unknown
},
)
.await;
assert!(
partitioned_at.elapsed() < Duration::from_secs(5),
"disruption must be event-driven, not the 6 s window",
);
c.shutdown().await.expect("shutdown C");
r.shutdown().await.expect("shutdown R");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn route_withdrawal_disrupts_provider_branches_at_remote_consumers() {
let fleet_kp = EntityKeypair::generate();
let fleet = AudienceScopeCommitment::owner_root(fleet_kp.entity_id());
let c = sensing_node(EntityKeypair::generate(), fleet, None, FD_TIMEOUT).await;
let x = sensing_node(
EntityKeypair::generate(),
fleet,
None,
Duration::from_millis(700),
)
.await;
let p = sensing_node(
EntityKeypair::generate(),
fleet,
Some(Incarnation::new(1)),
FD_TIMEOUT,
)
.await;
p.register_readiness_evaluator(CapabilityId::new("print.document"), Arc::new(AlwaysReady));
connect_pair(&c, &x).await;
connect_pair(&x, &p).await;
for node in [&c, &x, &p] {
node.start();
}
for node in [&c, &x, &p] {
node.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
}
let (c_id, x_id, p_id) = (c.node_id(), x.node_id(), p.node_id());
await_condition(Duration::from_secs(5), "entity pins established", || {
c.peer_entity_id(x_id).is_some()
&& x.peer_entity_id(c_id).is_some()
&& x.peer_entity_id(p_id).is_some()
&& p.peer_entity_id(x_id).is_some()
})
.await;
let p_entity = p.entity_keypair().entity_id().clone();
c.test_pin_peer_entity(p_id, p_entity);
let x_bind = x.local_addr();
let p_pub = *p.public_key();
c.connect_via(x_bind, &p_pub, p_id)
.await
.expect("relay-routed connect_via");
assert_eq!(
c.peer_addr(p_id),
Some(x_bind),
"precondition: C's session to P rides the relay",
);
let spec = shared_spec(fleet);
let branch = ProviderInterestKey::new(spec.key(), p_id);
await_condition(Duration::from_secs(5), "C senses P via X", || {
c.register_sensing_interest(&spec, p_id, D, LONG_TTL)
.is_ok()
&& c.sensing_upstream_continuity(&branch) == Some(Continuity::Established)
&& c.sensing_projected(&branch) == ProjectedReadiness::Ready
})
.await;
let partitioned_at = std::time::Instant::now();
chaos_partition(&x, &p);
await_condition(
Duration::from_secs(5),
"the received withdrawal expires the provider's observations",
|| {
c.sensing_upstream_continuity(&branch) == Some(Continuity::Expired)
&& c.sensing_projected(&branch) == ProjectedReadiness::Unknown
},
)
.await;
assert!(
partitioned_at.elapsed() < Duration::from_secs(5),
"disruption must ride the withdrawal, not the 6 s window",
);
c.shutdown().await.expect("shutdown C");
x.shutdown().await.expect("shutdown X");
p.shutdown().await.expect("shutdown P");
}
#[tokio::test]
async fn epoch_supersession_disrupts_sibling_branches() {
let fleet_kp = EntityKeypair::generate();
let fleet = AudienceScopeCommitment::owner_root(fleet_kp.entity_id());
let c = sensing_node(EntityKeypair::generate(), fleet, None, NO_FD).await;
let p = sensing_node(EntityKeypair::generate(), fleet, None, NO_FD).await;
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, p_id) = (c.node_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;
let spec_a = shared_spec(fleet);
let spec_b = {
let mut spec = shared_spec(fleet);
spec.constraints = CanonicalConstraints::from_entries([("media", "letter")]).unwrap();
spec
};
let branch_a = ProviderInterestKey::new(spec_a.key(), p_id);
let branch_b = ProviderInterestKey::new(spec_b.key(), p_id);
c.register_sensing_interest(&spec_a, p_id, D, LONG_TTL)
.expect("register A");
c.register_sensing_interest(&spec_b, p_id, D, LONG_TTL)
.expect("register B");
let c_addr = c.local_addr();
let craft = |digest: &InterestSpec, incarnation: u64, generation: u64, seq: u64| {
let unsigned = UnsignedAttestation {
interest_digest: digest.interest_digest(),
origin: p_id,
origin_incarnation: Incarnation::new(incarnation),
capability_id: CapabilityId::new("print.document"),
capability_generation: generation,
status: AttestedStatus::Ready,
status_reason: StatusReason::None,
estimated_start: None,
seq,
promised_cadence: Duration::from_millis(1000),
audience_scope: fleet,
};
let signed = sign_attestation(p.entity_keypair(), unsigned).expect("sign");
encode_attestation(&signed).expect("encode")
};
let send = |bytes: Vec<u8>| {
let p = p.clone();
async move {
p.send_subprotocol(c_addr, SUBPROTOCOL_READINESS_ATTESTATION, &bytes)
.await
.expect("send");
}
};
for seq in 1..=2u64 {
send(craft(&spec_a, 5, 1, seq)).await;
send(craft(&spec_b, 5, 1, seq)).await;
tokio::time::sleep(Duration::from_millis(100)).await;
}
await_condition(Duration::from_secs(5), "both branches established", || {
c.sensing_upstream_continuity(&branch_a) == Some(Continuity::Established)
&& c.sensing_upstream_continuity(&branch_b) == Some(Continuity::Established)
&& c.sensing_projected(&branch_b) == ProjectedReadiness::Ready
})
.await;
send(craft(&spec_a, 6, 1, 1)).await;
await_condition(
Duration::from_secs(2),
"sibling branch expires on the incarnation move",
|| {
c.sensing_upstream_continuity(&branch_b) == Some(Continuity::Expired)
&& c.sensing_projected(&branch_b) == ProjectedReadiness::Unknown
},
)
.await;
assert_eq!(
c.sensing_upstream_continuity(&branch_a),
Some(Continuity::Established),
"the epoch-advancing beat's own branch re-establishes",
);
send(craft(&spec_b, 5, 1, 10)).await;
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(
c.sensing_upstream_continuity(&branch_b),
Some(Continuity::Expired),
"a globally stale incarnation must not revive a sibling branch",
);
assert_eq!(c.sensing_projected(&branch_b), ProjectedReadiness::Unknown);
for seq in 2..=3u64 {
send(craft(&spec_b, 6, 1, seq)).await;
tokio::time::sleep(Duration::from_millis(100)).await;
}
await_condition(Duration::from_secs(5), "B re-establishes", || {
c.sensing_upstream_continuity(&branch_b) == Some(Continuity::Established)
})
.await;
send(craft(&spec_a, 6, 2, 2)).await;
await_condition(
Duration::from_secs(2),
"sibling branch expires on the generation move",
|| c.sensing_upstream_continuity(&branch_b) == Some(Continuity::Expired),
)
.await;
send(craft(&spec_b, 6, 1, 4)).await;
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(
c.sensing_upstream_continuity(&branch_b),
Some(Continuity::Expired),
"a globally stale generation must not revive a sibling branch",
);
send(craft(&spec_b, 6, 2, 5)).await;
await_condition(
Duration::from_secs(5),
"B resumes under the new epoch",
|| c.sensing_upstream_continuity(&branch_b) == Some(Continuity::Established),
)
.await;
c.shutdown().await.expect("shutdown C");
p.shutdown().await.expect("shutdown P");
}
#[cfg(feature = "redex")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn consumer_failure_drains_leader_demand() {
let owner_kp = EntityKeypair::generate();
let owner_entity = owner_kp.entity_id().clone();
let owner = AudienceScopeCommitment::owner_root(&owner_entity);
let a = sensing_node(EntityKeypair::generate(), owner, None, FD_TIMEOUT).await;
let r = sensing_node(EntityKeypair::generate(), owner, None, FD_TIMEOUT).await;
let p = sensing_node(owner_kp, owner, Some(Incarnation::new(1)), FD_TIMEOUT).await;
p.register_readiness_evaluator(CapabilityId::new("print.document"), Arc::new(AlwaysReady));
connect_pair(&a, &r).await;
connect_pair(&r, &p).await;
for node in [&a, &r, &p] {
node.start();
}
for node in [&a, &r] {
node.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
}
p.announce_capabilities(CapabilitySet::new().add_tag("print.document"))
.await
.expect("announce P");
let (a_id, r_id, p_id) = (a.node_id(), r.node_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(p_id).is_some()
&& p.peer_entity_id(r_id).is_some()
})
.await;
a.test_pin_peer_entity(p_id, owner_entity.clone());
assert!(r.assume_sensing_leader(), "leader role installs at R");
let cap = CapabilityId::new("print.document");
await_condition(Duration::from_secs(5), "R's snapshot authorizes P", || {
r.sensing_candidate_snapshot(&cap)
.iter()
.any(|candidate| candidate.node_id == p_id && candidate.authorized)
})
.await;
let spec = shared_spec(owner);
let branch = ProviderInterestKey::new(spec.key(), p_id);
await_condition(Duration::from_secs(10), "provider-free path live", || {
let _ = a.register_capability_interest(&spec, r_id, D, LONG_TTL);
r.sensing_leader_interest_count() == Some(1)
&& p.sensing_live_streams() == 1
&& a.sensing_projected(&branch) == ProjectedReadiness::Ready
})
.await;
let partitioned_at = std::time::Instant::now();
chaos_partition(&a, &r);
await_condition(
Duration::from_secs(5),
"leader demand drains on the failure edge",
|| r.sensing_leader_interest_count() == Some(0),
)
.await;
await_condition(
Duration::from_secs(5),
"provider retires the drained stream",
|| p.sensing_live_streams() == 0,
)
.await;
assert!(
partitioned_at.elapsed() < Duration::from_secs(7),
"the drain must be event-driven, never the 10 s row ttl",
);
a.shutdown().await.expect("shutdown A");
r.shutdown().await.expect("shutdown R");
p.shutdown().await.expect("shutdown P");
}