use super::Client;
use base64::Engine as _;
use ed25519_dalek::{Signer, SigningKey};
use futures::{stream::FuturesUnordered, StreamExt};
use iroh::{Endpoint, EndpointAddr, EndpointId, RelayMode};
use serde::{Deserialize, Serialize};
use std::{
collections::{BTreeSet, HashMap, HashSet},
sync::{
atomic::{AtomicU64, AtomicUsize, Ordering},
Arc,
},
time::{Duration, Instant, SystemTime, UNIX_EPOCH},
};
use tokio::{sync::mpsc, task::JoinHandle};
const MEMBERS: usize = 50;
const MAX_IN_FLIGHT_SENDS: usize = MEMBERS * crate::sparse_fanout::ACTIVE_NEIGHBOR_LIMIT;
const CAPABILITY: &str = "room:bounded-fifty-iroh";
const CROSS_CAPABILITY: &str = "room:bounded-fifty-cross-avenue";
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct LatestStatePayload {
avenue: String,
sequence: u64,
issued_at_ms: u64,
}
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct LaneWireFrame {
test_id: String,
capability: String,
envelope_base64: String,
expectation: WireExpectation,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
enum WireExpectation {
Deliver,
Reject,
Isolate,
}
#[derive(Debug)]
enum LaneEvent {
Delivery {
member: usize,
sequence: u64,
latency_ms: u64,
payload: Vec<u8>,
},
ExpectedRejection(String),
ExpectedIsolation {
test_id: String,
payload: Vec<u8>,
},
Error(String),
}
#[derive(Default)]
struct LaneMetrics {
transmissions: AtomicU64,
in_flight: AtomicUsize,
max_in_flight: AtomicUsize,
app_visible_duplicates: AtomicU64,
corrupt_deliveries: AtomicU64,
cross_avenue_deliveries: AtomicU64,
stale_generation_retirements: AtomicU64,
}
struct InFlight<'a>(&'a AtomicUsize);
impl Drop for InFlight<'_> {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::SeqCst);
}
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.min(u128::from(u64::MAX)) as u64
}
fn signing_key(index: usize) -> SigningKey {
let mut seed = [0_u8; 32];
seed[..8].copy_from_slice(&(index as u64 + 1).to_le_bytes());
SigningKey::from_bytes(&seed)
}
fn cross_avenue_signing_key(index: usize) -> SigningKey {
let mut seed = [0_u8; 32];
seed[..8].copy_from_slice(&(10_001_u64 + index as u64).to_le_bytes());
SigningKey::from_bytes(&seed)
}
fn key_x(key: &SigningKey) -> String {
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(key.verifying_key().as_bytes())
}
fn percentile_ms(values: &[u64], percentile: f64) -> u64 {
let mut values = values.to_vec();
values.sort_unstable();
let index = ((values.len() as f64 * percentile).ceil() as usize)
.saturating_sub(1)
.min(values.len().saturating_sub(1));
values[index]
}
async fn wait_for_endpoint_addr(endpoint: &Endpoint) -> anyhow::Result<EndpointAddr> {
tokio::time::timeout(Duration::from_secs(5), async {
loop {
let address = endpoint.addr();
if !address.is_empty() {
break address;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.map_err(|_| anyhow::anyhow!("endpoint {} did not publish a local address", endpoint.id()))
}
async fn send_wire(
client: &Client,
endpoint_id: EndpointId,
frame: &[u8],
metrics: &LaneMetrics,
) -> anyhow::Result<()> {
let in_flight = metrics.in_flight.fetch_add(1, Ordering::SeqCst) + 1;
metrics.max_in_flight.fetch_max(in_flight, Ordering::SeqCst);
let _guard = InFlight(&metrics.in_flight);
let mut send = client.open_uni_internal(endpoint_id).await?;
send.write_all(frame).await?;
send.finish()?;
metrics.transmissions.fetch_add(1, Ordering::SeqCst);
Ok(())
}
fn spawn_receiver(
member: usize,
client: Arc<Client>,
incoming: async_channel::Receiver<crate::native_node::IncomingStream>,
endpoints_by_node: Arc<HashMap<String, EndpointId>>,
events: mpsc::UnboundedSender<LaneEvent>,
metrics: Arc<LaneMetrics>,
) -> JoinHandle<()> {
tokio::spawn(async move {
while let Ok(incoming) = incoming.recv().await {
let crate::native_node::IncomingStreamType::Uni(mut recv) = incoming.stream else {
let _ = events.send(LaneEvent::Error(format!(
"member {member} received an unexpected bidirectional stream"
)));
continue;
};
let mut bytes = Vec::new();
if let Err(error) = tokio::io::AsyncReadExt::read_to_end(&mut recv, &mut bytes).await {
let _ = events.send(LaneEvent::Error(format!(
"member {member} failed reading Iroh stream: {error}"
)));
continue;
}
let frame = match serde_json::from_slice::<LaneWireFrame>(&bytes) {
Ok(frame) => frame,
Err(error) => {
let _ = events.send(LaneEvent::Error(format!(
"member {member} received corrupt lane framing: {error}"
)));
continue;
}
};
let encoded = match base64::engine::general_purpose::STANDARD
.decode(frame.envelope_base64.as_bytes())
{
Ok(encoded) => encoded,
Err(error) => {
let _ = events.send(LaneEvent::Error(format!(
"member {member} received invalid lane base64: {error}"
)));
continue;
}
};
let source = incoming.endpoint_id.to_string();
let decision = client
.accept_sparse_fanout_message(&frame.capability, &source, &encoded)
.await;
let decision = match decision {
Ok(decision) if frame.expectation == WireExpectation::Deliver => decision,
Ok(decision) if frame.expectation == WireExpectation::Isolate => {
let payload =
match serde_json::from_slice::<LatestStatePayload>(&decision.payload) {
Ok(payload)
if frame.capability == CROSS_CAPABILITY
&& payload.avenue == CROSS_CAPABILITY =>
{
decision.payload
}
Ok(payload) => {
metrics
.cross_avenue_deliveries
.fetch_add(1, Ordering::SeqCst);
let _ = events.send(LaneEvent::Error(format!(
"member {member} isolated probe escaped avenue {}",
payload.avenue
)));
continue;
}
Err(error) => {
metrics.corrupt_deliveries.fetch_add(1, Ordering::SeqCst);
let _ = events.send(LaneEvent::Error(format!(
"member {member} received corrupt isolated payload: {error}"
)));
continue;
}
};
let _ = events.send(LaneEvent::ExpectedIsolation {
test_id: frame.test_id,
payload,
});
continue;
}
Ok(_) => {
let _ = events.send(LaneEvent::Error(format!(
"member {member} accepted rejected probe {}",
frame.test_id
)));
continue;
}
Err(_) if frame.expectation == WireExpectation::Reject => {
let _ = events.send(LaneEvent::ExpectedRejection(frame.test_id));
continue;
}
Err(error) if error.to_string().contains("duplicate message") => continue,
Err(error) => {
let _ = events.send(LaneEvent::Error(format!(
"member {member} rejected {}: {error:#}",
frame.test_id
)));
continue;
}
};
let payload = match serde_json::from_slice::<LatestStatePayload>(&decision.payload) {
Ok(payload) => payload,
Err(error) => {
metrics.corrupt_deliveries.fetch_add(1, Ordering::SeqCst);
let _ = events.send(LaneEvent::Error(format!(
"member {member} received corrupt application payload: {error}"
)));
continue;
}
};
if payload.avenue != CAPABILITY {
metrics
.cross_avenue_deliveries
.fetch_add(1, Ordering::SeqCst);
let _ = events.send(LaneEvent::Error(format!(
"member {member} received cross-avenue payload {}",
payload.avenue
)));
continue;
}
let _ = events.send(LaneEvent::Delivery {
member,
sequence: payload.sequence,
latency_ms: now_ms().saturating_sub(payload.issued_at_ms),
payload: decision.payload,
});
let Some(forward) = decision.forward else {
continue;
};
let forward_frame = match serde_json::to_vec(&LaneWireFrame {
test_id: frame.test_id,
capability: frame.capability,
envelope_base64: base64::engine::general_purpose::STANDARD.encode(forward),
expectation: WireExpectation::Deliver,
}) {
Ok(frame) => Arc::new(frame),
Err(error) => {
let _ = events.send(LaneEvent::Error(format!(
"member {member} failed encoding forward frame: {error}"
)));
continue;
}
};
let mut sends = FuturesUnordered::new();
for peer_id in decision.forward_peer_ids {
let Some(endpoint_id) = endpoints_by_node.get(&peer_id).copied() else {
let _ = events.send(LaneEvent::Error(format!(
"member {member} selected unknown forward peer {peer_id}"
)));
continue;
};
let client = client.clone();
let frame = forward_frame.clone();
let metrics = metrics.clone();
sends.push(async move { send_wire(&client, endpoint_id, &frame, &metrics).await });
}
while let Some(result) = sends.next().await {
if let Err(error) = result {
let _ = events.send(LaneEvent::Error(format!(
"member {member} failed forwarding over Iroh: {error:#}"
)));
}
}
}
})
}
async fn signed_envelope(
client: &Client,
key: &SigningKey,
capability: &str,
sequence: u64,
) -> anyhow::Result<(Vec<u8>, Vec<u8>)> {
let payload = serde_json::to_vec(&LatestStatePayload {
avenue: capability.to_string(),
sequence,
issued_at_ms: now_ms(),
})?;
let request = client
.prepare_sparse_fanout_message(capability, &payload)
.await?;
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(key.sign(&request.signing_input).to_bytes());
let envelope = client
.finalize_sparse_fanout_message(&request.request_id, &signature)
.await?;
Ok((payload, envelope))
}
async fn send_broadcast(
author: usize,
sequence: u64,
clients: &[Arc<Client>],
routes: &[Vec<EndpointId>],
keys: &[SigningKey],
events: &mut mpsc::UnboundedReceiver<LaneEvent>,
metrics: &Arc<LaneMetrics>,
timeout: Duration,
) -> anyhow::Result<(usize, usize)> {
let (expected_payload, envelope) =
signed_envelope(&clients[author], &keys[author], CAPABILITY, sequence).await?;
let frame = Arc::new(serde_json::to_vec(&LaneWireFrame {
test_id: format!("latest-state-{sequence}"),
capability: CAPABILITY.to_string(),
envelope_base64: base64::engine::general_purpose::STANDARD.encode(envelope),
expectation: WireExpectation::Deliver,
})?);
let mut sends = FuturesUnordered::new();
for endpoint_id in routes[author].iter().copied() {
let client = clients[author].clone();
let frame = frame.clone();
let metrics = metrics.clone();
sends.push(async move { send_wire(&client, endpoint_id, &frame, &metrics).await });
}
while let Some(result) = sends.next().await {
result?;
}
let mut delivered = HashSet::new();
let mut fresh = 0;
tokio::time::timeout(timeout, async {
while delivered.len() < MEMBERS - 1 {
match events.recv().await {
Some(LaneEvent::Delivery {
member,
sequence: observed_sequence,
latency_ms,
payload,
}) if observed_sequence == sequence => {
anyhow::ensure!(member != author, "author received its own sparse broadcast");
anyhow::ensure!(
if payload == expected_payload {
true
} else {
metrics.corrupt_deliveries.fetch_add(1, Ordering::SeqCst);
false
},
"member {member} received corrupt latest-state payload"
);
anyhow::ensure!(
if delivered.insert(member) {
true
} else {
metrics
.app_visible_duplicates
.fetch_add(1, Ordering::SeqCst);
false
},
"member {member} received an app-visible duplicate"
);
fresh += usize::from(latency_ms <= 250);
}
Some(LaneEvent::Delivery {
sequence: other, ..
}) => {
anyhow::bail!("delivery from sequence {other} leaked into sequence {sequence}")
}
Some(LaneEvent::ExpectedRejection(id)) => {
anyhow::bail!("unexpected rejection event {id} during broadcast")
}
Some(LaneEvent::ExpectedIsolation { test_id, .. }) => {
anyhow::bail!("unexpected isolation event {test_id} during broadcast")
}
Some(LaneEvent::Error(error)) => anyhow::bail!(error),
None => anyhow::bail!("Iroh delivery event channel closed"),
}
}
anyhow::Ok(())
})
.await
.map_err(|_| anyhow::anyhow!("sequence {sequence} did not reach all 49 peers in time"))??;
Ok((delivered.len(), fresh))
}
async fn expect_rejected_probe(
test_id: &str,
source: usize,
target: EndpointId,
capability: &str,
envelope: Vec<u8>,
clients: &[Arc<Client>],
events: &mut mpsc::UnboundedReceiver<LaneEvent>,
metrics: &Arc<LaneMetrics>,
) -> anyhow::Result<u64> {
let frame = serde_json::to_vec(&LaneWireFrame {
test_id: test_id.to_string(),
capability: capability.to_string(),
envelope_base64: base64::engine::general_purpose::STANDARD.encode(envelope),
expectation: WireExpectation::Reject,
})?;
send_wire(&clients[source], target, &frame, metrics).await?;
tokio::time::timeout(Duration::from_secs(2), async {
loop {
match events.recv().await {
Some(LaneEvent::ExpectedRejection(id)) if id == test_id => return Ok(1),
Some(LaneEvent::ExpectedRejection(_)) => continue,
Some(LaneEvent::ExpectedIsolation { test_id, .. }) => {
anyhow::bail!("unexpected isolation event {test_id}")
}
Some(LaneEvent::Error(error)) => anyhow::bail!(error),
Some(LaneEvent::Delivery { .. }) => {
anyhow::bail!("rejected probe became app-visible")
}
None => anyhow::bail!("Iroh delivery event channel closed"),
}
}
})
.await
.map_err(|_| anyhow::anyhow!("rejected probe {test_id} timed out"))?
}
async fn expect_isolated_probe(
test_id: &str,
source: usize,
target: EndpointId,
envelope: Vec<u8>,
expected_payload: &[u8],
clients: &[Arc<Client>],
events: &mut mpsc::UnboundedReceiver<LaneEvent>,
metrics: &Arc<LaneMetrics>,
) -> anyhow::Result<u64> {
let frame = serde_json::to_vec(&LaneWireFrame {
test_id: test_id.to_string(),
capability: CROSS_CAPABILITY.to_string(),
envelope_base64: base64::engine::general_purpose::STANDARD.encode(envelope),
expectation: WireExpectation::Isolate,
})?;
send_wire(&clients[source], target, &frame, metrics).await?;
tokio::time::timeout(Duration::from_secs(2), async {
loop {
match events.recv().await {
Some(LaneEvent::ExpectedIsolation {
test_id: id,
payload,
}) if id == test_id => {
anyhow::ensure!(
payload == expected_payload,
"isolated avenue changed its application payload"
);
return Ok(1);
}
Some(LaneEvent::ExpectedIsolation { .. }) => continue,
Some(LaneEvent::ExpectedRejection(id)) => {
anyhow::bail!("valid isolated avenue was rejected as {id}")
}
Some(LaneEvent::Error(error)) => anyhow::bail!(error),
Some(LaneEvent::Delivery { .. }) => {
anyhow::bail!("isolated avenue became visible to the primary avenue")
}
None => anyhow::bail!("Iroh delivery event channel closed"),
}
}
})
.await
.map_err(|_| anyhow::anyhow!("isolated probe {test_id} timed out"))?
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
#[ignore = "required local lane runs 300 seconds; use the package command"]
async fn fifty_member_sparse_iroh_latest_state_lane() -> anyhow::Result<()> {
rustls::crypto::ring::default_provider()
.install_default()
.ok();
let sanity = std::env::var("OPENRTC_BOUNDED_50_SPARSE_IROH_SANITY").as_deref() == Ok("1");
let active_duration = if sanity {
Duration::from_secs(8)
} else {
Duration::from_secs(240)
};
let quiet_duration = if sanity {
Duration::from_secs(2)
} else {
Duration::from_secs(60)
};
let broadcast_interval = if sanity {
Duration::from_millis(500)
} else {
Duration::from_millis(500)
};
let churn_interval = if sanity {
Duration::from_secs(3)
} else {
Duration::from_secs(10)
};
let keys = (0..MEMBERS).map(signing_key).collect::<Vec<_>>();
let cross_keys = (0..MEMBERS)
.map(cross_avenue_signing_key)
.collect::<Vec<_>>();
let mut clients = Vec::with_capacity(MEMBERS);
let mut endpoint_ids = Vec::with_capacity(MEMBERS);
let mut endpoint_addrs = Vec::<EndpointAddr>::with_capacity(MEMBERS);
let mut receivers = Vec::with_capacity(MEMBERS);
for member in 0..MEMBERS {
let client = Arc::new(Client::new_with_app_tag(
crate::test_constants::TEST_PROJECT_ID.to_string(),
format!("bounded-50-iroh-{member:02}"),
Box::new(|| None),
));
let endpoint = Endpoint::builder(iroh::endpoint::presets::N0)
.alpns(vec![crate::native_node::PlutoniumProtocol::ALPN.to_vec()])
.relay_mode(RelayMode::Disabled)
.bind_addr("127.0.0.1:0")?
.bind()
.await?;
client.adopt_endpoint(endpoint).await;
let endpoint = client.get_endpoint().await?;
endpoint_ids.push(endpoint.id());
endpoint_addrs.push(wait_for_endpoint_addr(&endpoint).await?);
receivers.push(client.incoming_streams().await?);
clients.push(client);
}
let roster = endpoint_ids
.iter()
.enumerate()
.map(|(member, endpoint_id)| {
serde_json::json!({
"deviceId": format!("device-{member:02}"),
"nodeId": endpoint_id.to_string(),
"devicePublicKeyJwk": {
"kty": "OKP",
"crv": "Ed25519",
"x": key_x(&keys[member]),
},
"online": true,
"ticket": format!("local-only-{member:02}"),
})
})
.collect::<Vec<_>>();
let cross_roster = endpoint_ids
.iter()
.enumerate()
.map(|(member, endpoint_id)| {
serde_json::json!({
"deviceId": format!("cross-device-{member:02}"),
"nodeId": endpoint_id.to_string(),
"devicePublicKeyJwk": {
"kty": "OKP",
"crv": "Ed25519",
"x": key_x(&cross_keys[member]),
},
"online": true,
"ticket": format!("cross-local-only-{member:02}"),
})
})
.collect::<Vec<_>>();
let by_node = endpoint_ids
.iter()
.copied()
.map(|endpoint_id| (endpoint_id.to_string(), endpoint_id))
.collect::<HashMap<_, _>>();
let mut routes = Vec::with_capacity(MEMBERS);
let mut edges = BTreeSet::new();
for member in 0..MEMBERS {
let peers = roster
.iter()
.enumerate()
.filter(|(index, _)| *index != member)
.map(|(_, peer)| peer.clone())
.collect::<Vec<_>>();
let (selected, projection) = clients[member]
.configure_sparse_fanout(
CAPABILITY,
1,
&format!("device-{member:02}"),
&key_x(&keys[member]),
&serde_json::to_string(&peers)?,
true,
)
.await?;
anyhow::ensure!(projection.sparse, "member {member} did not select sparse");
anyhow::ensure!(projection.member_count == MEMBERS, "member count drift");
anyhow::ensure!(
projection.active_neighbor_count == crate::sparse_fanout::ACTIVE_NEIGHBOR_LIMIT,
"member {member} did not receive four sparse neighbors"
);
let cross_peers = cross_roster
.iter()
.enumerate()
.filter(|(index, _)| *index != member)
.map(|(_, peer)| peer.clone())
.collect::<Vec<_>>();
let (_, cross_projection) = clients[member]
.configure_sparse_fanout(
CROSS_CAPABILITY,
1,
&format!("cross-device-{member:02}"),
&key_x(&cross_keys[member]),
&serde_json::to_string(&cross_peers)?,
true,
)
.await?;
anyhow::ensure!(
cross_projection.sparse && cross_projection.member_count == MEMBERS,
"member {member} did not configure the isolated cross avenue"
);
let selected = serde_json::from_str::<Vec<serde_json::Value>>(&selected)?
.into_iter()
.map(|peer| {
peer["nodeId"]
.as_str()
.and_then(|node_id| by_node.get(node_id).copied())
.ok_or_else(|| anyhow::anyhow!("member {member} received an invalid route"))
})
.collect::<anyhow::Result<Vec<_>>>()?;
for target in &selected {
let target_index = endpoint_ids
.iter()
.position(|endpoint_id| endpoint_id == target)
.ok_or_else(|| anyhow::anyhow!("route target is outside roster"))?;
edges.insert((member.min(target_index), member.max(target_index)));
}
routes.push(selected);
}
anyhow::ensure!(edges.len() == 100, "expected 100 undirected sparse edges");
for &(left, right) in &edges {
clients[left]
.ensure_connected_addr(endpoint_ids[right], endpoint_addrs[right].clone())
.await?;
}
tokio::time::timeout(Duration::from_secs(10), async {
loop {
let mut ready = true;
for member in 0..MEMBERS {
for endpoint_id in &routes[member] {
if clients[member].get_connection(*endpoint_id).await.is_none() {
ready = false;
break;
}
}
}
if ready {
break;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.map_err(|_| anyhow::anyhow!("50/50 sparse Iroh setup did not settle"))?;
let metrics = Arc::new(LaneMetrics::default());
let (event_tx, mut event_rx) = mpsc::unbounded_channel();
let endpoints_by_node = Arc::new(by_node);
let receiver_tasks = receivers
.iter()
.cloned()
.enumerate()
.map(|(member, receiver)| {
spawn_receiver(
member,
clients[member].clone(),
receiver,
endpoints_by_node.clone(),
event_tx.clone(),
metrics.clone(),
)
})
.collect::<Vec<_>>();
let (cross_payload, cross_envelope) =
signed_envelope(&clients[0], &cross_keys[0], CROSS_CAPABILITY, 1).await?;
let isolated_cross_avenue_messages = expect_isolated_probe(
"valid-cross-avenue",
0,
routes[0][0],
cross_envelope.clone(),
&cross_payload,
&clients,
&mut event_rx,
&metrics,
)
.await?;
let rejected_cross_avenue_replays = expect_rejected_probe(
"cross-avenue-as-primary",
0,
routes[0][0],
CAPABILITY,
cross_envelope,
&clients,
&mut event_rx,
&metrics,
)
.await?;
let (_, mut corrupt_envelope) = signed_envelope(&clients[0], &keys[0], CAPABILITY, 2).await?;
if let Some(last) = corrupt_envelope.last_mut() {
*last ^= 1;
}
let rejected_corrupt_envelopes = expect_rejected_probe(
"corrupt-envelope",
0,
routes[0][0],
CAPABILITY,
corrupt_envelope,
&clients,
&mut event_rx,
&metrics,
)
.await?;
let probe_transmissions = metrics.transmissions.load(Ordering::SeqCst);
let started = Instant::now();
let total_duration = active_duration + quiet_duration;
let mut next_broadcast = started;
let mut next_churn = started + churn_interval;
let edge_list = edges.iter().copied().collect::<Vec<_>>();
let mut edge_cursor = 0;
let mut sequence = 10_u64;
let mut broadcasts = 0_u64;
let mut delivery_samples = 0_u64;
let mut fresh_samples = 0_u64;
let mut reconnect_ms = Vec::new();
let mut quiet_queue_start = None;
while started.elapsed() < total_duration {
tokio::time::sleep_until(next_broadcast.into()).await;
next_broadcast += broadcast_interval;
let elapsed = started.elapsed();
if elapsed >= active_duration && quiet_queue_start.is_none() {
quiet_queue_start = Some(
receivers
.iter()
.map(|receiver| receiver.len())
.sum::<usize>(),
);
}
let mut churn = None;
if elapsed < active_duration && Instant::now() >= next_churn {
let (left, right) = edge_list[edge_cursor % edge_list.len()];
edge_cursor += 1;
next_churn += churn_interval;
let remote = endpoint_ids[right];
let connection = clients[left]
.get_connection(remote)
.await
.ok_or_else(|| anyhow::anyhow!("churn edge {left}->{right} is not connected"))?;
let old_stable_id = crate::transport_generation::for_connection(&connection);
let connection_id = Client::deterministic_connection_id(
&endpoint_ids[left].to_string(),
&endpoint_ids[right].to_string(),
);
let fault_started = Instant::now();
clients[left]
.disconnect_with_reason(
remote,
crate::lifecycle_reason::REASON_NETWORK_CHANGE_RECONNECT,
)
.await?;
churn = Some((
left,
right,
remote,
connection_id,
old_stable_id,
fault_started,
));
}
let author = churn
.as_ref()
.map(|churn| churn.0)
.unwrap_or(sequence as usize % MEMBERS);
let (delivered, fresh) = send_broadcast(
author,
sequence,
&clients,
&routes,
&keys,
&mut event_rx,
&metrics,
if churn.is_some() {
Duration::from_secs(10)
} else {
Duration::from_secs(2)
},
)
.await?;
delivery_samples += delivered as u64;
fresh_samples += fresh as u64;
broadcasts += 1;
sequence += 1;
if let Some((left, right, remote, connection_id, old_stable_id, fault_started)) = churn {
let (new_stable_id, current_record) =
tokio::time::timeout(Duration::from_secs(10), async {
loop {
let physical =
clients[left]
.get_connection(remote)
.await
.map(|connection| {
crate::transport_generation::for_connection(&connection)
});
let record = clients[left]
.connection_manager
.get_by_connection_id(&connection_id)
.await;
if let (Some(new_stable_id), Some(record)) = (physical, record) {
if new_stable_id != old_stable_id
&& record.transport_stable_id == Some(new_stable_id)
{
break (new_stable_id, record);
}
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.map_err(|_| {
anyhow::anyhow!("edge {left}->{right} did not install a replacement")
})?;
reconnect_ms.push(fault_started.elapsed().as_millis() as u64);
let stale = clients[left]
.connection_manager
.set_closed_if_current(
&connection_id,
old_stable_id,
Some("delayed-old-physical-generation".to_string()),
)
.await;
if stale.is_some() {
metrics
.stale_generation_retirements
.fetch_add(1, Ordering::SeqCst);
}
anyhow::ensure!(
stale.is_none(),
"stale physical generation retired replacement"
);
let after_stale = clients[left]
.connection_manager
.get_by_connection_id(&connection_id)
.await
.ok_or_else(|| anyhow::anyhow!("replacement record disappeared"))?;
anyhow::ensure!(
after_stale.transport_stable_id == Some(new_stable_id)
&& after_stale.transport_generation == current_record.transport_generation,
"stale close changed replacement generation"
);
}
}
tokio::time::sleep(Duration::from_millis(500)).await;
let measured_ms = started.elapsed().as_millis() as u64;
let queue_end = receivers
.iter()
.map(|receiver| receiver.len())
.sum::<usize>();
let queue_start = quiet_queue_start.unwrap_or(queue_end);
let transmissions = metrics
.transmissions
.load(Ordering::SeqCst)
.saturating_sub(probe_transmissions);
let expected_deliveries = broadcasts * (MEMBERS as u64 - 1);
let p95 = percentile_ms(&reconnect_ms, 0.95);
let p99 = percentile_ms(&reconnect_ms, 0.99);
anyhow::ensure!(
delivery_samples == expected_deliveries,
"delivery sample count drift"
);
anyhow::ensure!(
fresh_samples.saturating_mul(1_000) >= delivery_samples.saturating_mul(999),
"latest-state freshness fell below 99.9% within 250ms"
);
anyhow::ensure!(
transmissions.saturating_mul(4) <= expected_deliveries.saturating_mul(5),
"actual Iroh transmission amplification exceeded 1.25x"
);
anyhow::ensure!(
!reconnect_ms.is_empty(),
"lane did not exercise reconnection"
);
anyhow::ensure!(p95 <= 5_000, "reconnect p95 exceeded five seconds");
anyhow::ensure!(p99 <= 10_000, "reconnect p99 exceeded ten seconds");
anyhow::ensure!(
queue_end <= queue_start,
"queues grew through final quiet window"
);
anyhow::ensure!(queue_end == 0, "application stream queues did not drain");
anyhow::ensure!(
metrics.in_flight.load(Ordering::SeqCst) == 0,
"Iroh send tasks did not drain"
);
anyhow::ensure!(
metrics.max_in_flight.load(Ordering::SeqCst) <= MAX_IN_FLIGHT_SENDS,
"Iroh send concurrency exceeded its fixed local bound"
);
for (member, client) in clients.iter().enumerate() {
let diagnostics = client.sparse_fanout_diagnostics(CAPABILITY).await;
anyhow::ensure!(
diagnostics.forward_queue_drops == 0,
"member {member} dropped forward work"
);
anyhow::ensure!(
diagnostics.budget_drops == 0,
"member {member} exceeded sparse budget"
);
anyhow::ensure!(
diagnostics.roster_members == MEMBERS as u64,
"member {member} roster drifted"
);
}
if !sanity {
anyhow::ensure!(
measured_ms >= 300_000,
"release lane ran for less than 300 seconds"
);
}
let summary = serde_json::json!({
"schemaVersion": 1,
"mode": if sanity { "sanity-only" } else { "release" },
"status": "passed",
"members": MEMBERS,
"setupMembers": MEMBERS,
"sparseEdges": edges.len(),
"measurementDurationMs": measured_ms,
"activeChurnDurationMs": active_duration.as_millis() as u64,
"quietDurationMs": quiet_duration.as_millis() as u64,
"broadcasts": broadcasts,
"deliverySamples": delivery_samples,
"freshWithin250Ms": fresh_samples,
"transmissions": transmissions,
"deliveryAmplification": transmissions as f64 / expected_deliveries as f64,
"reconnections": reconnect_ms.len(),
"reconnectP95Ms": p95,
"reconnectP99Ms": p99,
"maxInFlightSends": metrics.max_in_flight.load(Ordering::SeqCst),
"maxInFlightSendLimit": MAX_IN_FLIGHT_SENDS,
"quietQueueStart": queue_start,
"quietQueueEnd": queue_end,
"appVisibleDuplicates": metrics.app_visible_duplicates.load(Ordering::SeqCst),
"corruptDeliveries": metrics.corrupt_deliveries.load(Ordering::SeqCst),
"crossAvenueDeliveries": metrics.cross_avenue_deliveries.load(Ordering::SeqCst),
"staleGenerationRetirements": metrics.stale_generation_retirements.load(Ordering::SeqCst),
"isolatedCrossAvenueMessages": isolated_cross_avenue_messages,
"rejectedCrossAvenueReplays": rejected_cross_avenue_replays,
"rejectedCorruptEnvelopes": rejected_corrupt_envelopes,
});
println!("[openrtc-bounded-50-sparse-iroh] {summary}");
for task in receiver_tasks {
task.abort();
}
Ok(())
}