use anyhow::{anyhow, bail, ensure, Context, Result};
use base64::Engine as _;
use ed25519_dalek::{Signature, Verifier, VerifyingKey};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
pub const SMALL_MESH_LIMIT: usize = 8;
pub const ACTIVE_NEIGHBOR_LIMIT: usize = 4;
pub const MAX_REVIEWED_MEMBERS: usize = 99;
pub const MAX_OVERLAY_PAYLOAD_BYTES: usize = 64 * 1024;
pub const MAX_OVERLAY_HOPS: u8 = 32;
pub const MAX_SEEN_MESSAGES: usize = 4_096;
pub const OVERLAY_MESSAGE_LIFETIME_MS: u64 = 30_000;
pub const MAX_OVERLAY_FORWARD_PEERS: usize = SMALL_MESH_LIMIT - 1;
const MAX_OVERLAY_ENVELOPE_BYTES: usize = (MAX_OVERLAY_PAYLOAD_BYTES * 4 / 3) + 4_096;
const OVERLAY_BUDGET_WINDOW_MS: u64 = 1_000;
const MAX_INGRESS_ACTIONS_PER_WINDOW: u64 = 128;
const MAX_INGRESS_BYTES_PER_WINDOW: u64 = (MAX_OVERLAY_ENVELOPE_BYTES as u64) * 8;
const MAX_FORWARD_ACTIONS_PER_WINDOW: u64 = 64;
pub const MAX_FORWARD_QUEUE_ACTIONS: usize = 64;
const MAX_FORWARDED_BYTES_PER_WINDOW: u64 =
(MAX_OVERLAY_ENVELOPE_BYTES as u64) * (ACTIVE_NEIGHBOR_LIMIT as u64) * 8;
const AUTHOR_REPLAY_WINDOW_BITS: u64 = 128;
const PROTOCOL_VERSION: &str = "openrtc-sparse-fanout/1";
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SparseFanoutProjection {
pub topology_revision: u64,
pub member_count: usize,
pub active_neighbor_count: usize,
pub sparse: bool,
}
#[derive(Debug, Clone)]
struct Member {
device_id: String,
node_id: Option<String>,
device_key_x: String,
value: Value,
}
#[derive(Debug, Clone)]
struct MemberIdentity {
device_id: String,
device_key_x: String,
}
#[derive(Debug, Clone)]
struct MembershipRevision {
revision: u64,
local_device_id: String,
local_node_id: String,
local_device_key_x: String,
by_node: HashMap<String, MemberIdentity>,
by_device: HashMap<String, MemberIdentity>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct UnsignedEnvelope {
protocol: String,
capability: String,
topology_revision: u64,
author_device_id: String,
author_peer_id: String,
sender_epoch_ms: u64,
sequence: u64,
message_id: String,
issued_at_ms: u64,
expires_at_ms: u64,
max_hops: u8,
payload_hash: String,
payload: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct SignedEnvelope {
#[serde(rename = "__openrtcFanout")]
marker: u8,
#[serde(flatten)]
unsigned: UnsignedEnvelope,
hop: u8,
signature: String,
}
#[derive(Debug, Clone)]
struct PendingEnvelope {
capability: String,
unsigned: UnsignedEnvelope,
signing_input: Vec<u8>,
payload_bytes: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SparseFanoutSigningRequest {
pub request_id: String,
pub signing_input: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SparseFanoutDecision {
pub payload: Vec<u8>,
pub author_device_id: String,
pub author_peer_id: String,
pub topology_revision: u64,
pub hop: u8,
pub forward_peer_ids: Vec<String>,
pub forward_peer_limit: usize,
pub forward_queue_capacity: usize,
#[serde(skip_serializing_if = "Option::is_none")]
pub forward: Option<Vec<u8>>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SparseFanoutDiagnostics {
pub configured_revisions: u64,
pub encoded: u64,
pub accepted: u64,
pub forwarded: u64,
pub duplicate_drops: u64,
pub rejected: u64,
pub budget_drops: u64,
pub replay_window_drops: u64,
pub forward_queue_drops: u64,
pub ingress_bytes: u64,
pub payload_bytes: u64,
pub forwarded_bytes: u64,
pub active_neighbors: u64,
pub overlap_neighbors: u64,
pub max_forward_peers: u64,
pub roster_members: u64,
pub topology_revision: u64,
}
#[derive(Debug, Default)]
pub struct SparseFanoutState {
memberships: HashMap<String, VecDeque<MembershipRevision>>,
sequences: HashMap<String, u64>,
sender_epochs: HashMap<String, u64>,
author_replay: HashMap<String, AuthorReplayWindow>,
traffic: HashMap<String, TrafficWindow>,
seen: HashSet<String>,
seen_order: VecDeque<String>,
pending: HashMap<String, PendingEnvelope>,
pending_order: VecDeque<String>,
diagnostics: HashMap<String, SparseFanoutDiagnostics>,
}
const MAX_PENDING_SIGNATURES: usize = 64;
#[derive(Debug, Default)]
struct AuthorReplayWindow {
highest_sequence: u64,
seen_bits: u128,
}
impl AuthorReplayWindow {
fn accept(&mut self, sequence: u64) -> bool {
if sequence == 0 {
return false;
}
if sequence > self.highest_sequence {
let shift = sequence - self.highest_sequence;
self.seen_bits = if shift >= AUTHOR_REPLAY_WINDOW_BITS {
1
} else {
(self.seen_bits << shift) | 1
};
self.highest_sequence = sequence;
return true;
}
let distance = self.highest_sequence - sequence;
if distance >= AUTHOR_REPLAY_WINDOW_BITS {
return false;
}
let bit = 1_u128 << distance;
if self.seen_bits & bit != 0 {
return false;
}
self.seen_bits |= bit;
true
}
}
#[derive(Debug, Default)]
struct TrafficWindow {
started_at_ms: u64,
ingress_actions: u64,
ingress_bytes: u64,
forward_actions: u64,
forwarded_bytes: u64,
}
impl TrafficWindow {
fn refresh(&mut self, now_ms: u64) {
if self.started_at_ms == 0
|| now_ms.saturating_sub(self.started_at_ms) >= OVERLAY_BUDGET_WINDOW_MS
{
*self = Self {
started_at_ms: now_ms,
..Self::default()
};
}
}
}
fn canonical_message_id(
author_device_id: &str,
author_peer_id: &str,
topology_revision: u64,
sender_epoch_ms: u64,
sequence: u64,
) -> String {
blake3::hash(
format!(
"{PROTOCOL_VERSION}\0{author_device_id}\0{author_peer_id}\0{topology_revision}\0{sender_epoch_ms}\0{sequence}"
)
.as_bytes(),
)
.to_hex()
.to_string()
}
fn required_part(value: &str, label: &str, max: usize) -> Result<String> {
let value = value.trim();
ensure!(!value.is_empty(), "sparse fanout {label} is empty");
ensure!(value.len() <= max, "sparse fanout {label} is too large");
ensure!(
value
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || b"_.:@-".contains(&byte)),
"sparse fanout {label} contains unsupported characters"
);
Ok(value.to_string())
}
fn public_key_x(value: &Value) -> Option<String> {
let jwk = value.as_object()?;
if jwk.get("kty").and_then(Value::as_str) != Some("OKP")
|| jwk.get("crv").and_then(Value::as_str) != Some("Ed25519")
|| jwk.contains_key("d")
{
return None;
}
let x = jwk.get("x")?.as_str()?.trim();
(x.len() == 43
&& x.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-')))
.then(|| x.to_string())
}
fn normalized_members(peers: Vec<Value>, local_device_id: &str) -> Result<Vec<Member>> {
let mut by_device = BTreeMap::<String, Member>::new();
for value in peers {
let Some(object) = value.as_object() else {
continue;
};
let device_id = object
.get("deviceId")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty());
let node_id = object
.get("nodeId")
.and_then(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty());
let online = object
.get("online")
.and_then(Value::as_bool)
.unwrap_or(true);
let Some(device_id) = device_id else {
continue;
};
if !online || device_id == local_device_id {
continue;
}
let device_key_x = object
.get("devicePublicKeyJwk")
.and_then(public_key_x)
.ok_or_else(|| {
anyhow::anyhow!(
"sparse fanout member {device_id} is missing a valid device signing key"
)
})?;
by_device.insert(
device_id.to_string(),
Member {
device_id: device_id.to_string(),
node_id: node_id.map(ToString::to_string),
device_key_x,
value,
},
);
}
Ok(by_device.into_values().collect())
}
fn selected_device_ids(members: &[Member], local_device_id: &str) -> HashSet<String> {
if members.len() + 1 <= SMALL_MESH_LIMIT {
return members
.iter()
.map(|member| member.device_id.clone())
.collect();
}
let mut ids = members
.iter()
.map(|member| member.device_id.clone())
.chain(std::iter::once(local_device_id.to_string()))
.collect::<Vec<_>>();
ids.sort();
ids.dedup();
let Some(local_index) = ids.iter().position(|id| id == local_device_id) else {
return HashSet::new();
};
let mut selected = HashSet::new();
for distance in 1..ids.len() {
selected.insert(ids[(local_index + distance) % ids.len()].clone());
selected.insert(ids[(local_index + ids.len() - distance) % ids.len()].clone());
selected.remove(local_device_id);
if selected.len() >= ACTIVE_NEIGHBOR_LIMIT {
break;
}
}
selected
}
fn directed_forward_peer_ids(
membership: &MembershipRevision,
author_device_id: &str,
) -> Vec<String> {
let mut device_ids = membership.by_device.keys().cloned().collect::<Vec<_>>();
device_ids.sort();
let Some(author_index) = device_ids.iter().position(|id| id == author_device_id) else {
return Vec::new();
};
let Some(local_index) = device_ids
.iter()
.position(|id| id == &membership.local_device_id)
else {
return Vec::new();
};
if local_index == author_index || device_ids.len() < 2 {
return Vec::new();
}
let clockwise = (local_index + device_ids.len() - author_index) % device_ids.len();
let counter_clockwise = (author_index + device_ids.len() - local_index) % device_ids.len();
if clockwise == counter_clockwise {
return Vec::new();
}
let target_index = if clockwise < counter_clockwise {
(local_index + 1) % device_ids.len()
} else {
(local_index + device_ids.len() - 1) % device_ids.len()
};
let target_device_id = &device_ids[target_index];
membership
.by_node
.iter()
.find_map(|(node_id, identity)| {
(identity.device_id == *target_device_id).then(|| node_id.clone())
})
.into_iter()
.collect()
}
impl SparseFanoutState {
fn charge_ingress(&mut self, capability: &str, bytes: usize, now_ms: u64) -> Result<()> {
let budget = self.traffic.entry(capability.to_string()).or_default();
budget.refresh(now_ms);
budget.ingress_actions = budget.ingress_actions.saturating_add(1);
budget.ingress_bytes = budget.ingress_bytes.saturating_add(bytes as u64);
let diagnostics = self.diagnostics.entry(capability.to_string()).or_default();
diagnostics.ingress_bytes = diagnostics.ingress_bytes.saturating_add(bytes as u64);
if budget.ingress_actions > MAX_INGRESS_ACTIONS_PER_WINDOW
|| budget.ingress_bytes > MAX_INGRESS_BYTES_PER_WINDOW
{
diagnostics.budget_drops = diagnostics.budget_drops.saturating_add(1);
bail!("sparse fanout ingress budget exceeded");
}
Ok(())
}
fn authorize_forward(
&mut self,
capability: &str,
encoded_bytes: usize,
peer_limit: usize,
now_ms: u64,
) -> Result<()> {
ensure!(
peer_limit <= MAX_OVERLAY_FORWARD_PEERS,
"sparse fanout per-message peer limit exceeded"
);
let projected_bytes = (encoded_bytes as u64).saturating_mul(peer_limit as u64);
let budget = self.traffic.entry(capability.to_string()).or_default();
budget.refresh(now_ms);
budget.forward_actions = budget.forward_actions.saturating_add(1);
budget.forwarded_bytes = budget.forwarded_bytes.saturating_add(projected_bytes);
if budget.forward_actions > MAX_FORWARD_ACTIONS_PER_WINDOW
|| budget.forwarded_bytes > MAX_FORWARDED_BYTES_PER_WINDOW
{
self.diagnostics
.entry(capability.to_string())
.or_default()
.budget_drops += 1;
bail!("sparse fanout forward action budget exceeded");
}
Ok(())
}
pub fn configure(
&mut self,
capability: &str,
revision: u64,
local_device_id: &str,
local_node_id: &str,
local_device_key_x: &str,
peers: Vec<Value>,
sparse: bool,
) -> Result<(Vec<Value>, SparseFanoutProjection)> {
let capability = required_part(capability, "capability", 192)?;
let local_device_id = required_part(local_device_id, "local device id", 256)?;
let local_node_id = required_part(local_node_id, "local peer id", 256)?;
let local_device_key_x = required_part(local_device_key_x, "local device key", 43)?;
ensure!(
local_device_key_x.len() == 43,
"sparse fanout local device key is invalid"
);
ensure!(revision > 0, "sparse fanout revision must be positive");
let history = self.memberships.entry(capability.clone()).or_default();
if history
.back()
.is_some_and(|current| revision <= current.revision)
{
bail!("sparse fanout revision is stale");
}
let members = normalized_members(peers, &local_device_id)?;
ensure!(
members.len() + 1 <= MAX_REVIEWED_MEMBERS,
"sparse fanout roster exceeds {MAX_REVIEWED_MEMBERS} members"
);
let by_node = members
.iter()
.filter_map(|member| {
let node_id = member.node_id.as_ref()?;
(
node_id.clone(),
MemberIdentity {
device_id: member.device_id.clone(),
device_key_x: member.device_key_x.clone(),
},
)
.into()
})
.chain(std::iter::once((
local_node_id.clone(),
MemberIdentity {
device_id: local_device_id.clone(),
device_key_x: local_device_key_x.clone(),
},
)))
.collect::<HashMap<_, _>>();
let by_device = members
.iter()
.map(|member| {
(
member.device_id.clone(),
MemberIdentity {
device_id: member.device_id.clone(),
device_key_x: member.device_key_x.clone(),
},
)
})
.chain(std::iter::once((
local_device_id.clone(),
MemberIdentity {
device_id: local_device_id.clone(),
device_key_x: local_device_key_x.clone(),
},
)))
.collect::<HashMap<_, _>>();
history.push_back(MembershipRevision {
revision,
local_device_id: local_device_id.clone(),
local_node_id,
local_device_key_x,
by_node,
by_device,
});
while history.len() > 2 {
history.pop_front();
}
let sparse = sparse && members.len() + 1 > SMALL_MESH_LIMIT;
let gateway_lease = sparse
&& members.iter().any(|member| {
member.value.get("topologyRevision").is_some()
|| member.value.get("topologyRole").is_some()
});
if gateway_lease {
let routes = members
.iter()
.filter(|member| member.node_id.is_some())
.collect::<Vec<_>>();
let revisions = routes
.iter()
.filter_map(|member| member.value.get("topologyRevision").and_then(Value::as_u64))
.collect::<HashSet<_>>();
let active = routes
.iter()
.filter(|member| {
member.value.get("topologyRole").and_then(Value::as_str) == Some("active")
})
.count();
let backup = routes
.iter()
.filter(|member| {
member.value.get("topologyRole").and_then(Value::as_str) == Some("backup")
})
.count();
ensure!(
revisions.len() == 1
&& revisions.iter().all(|revision| *revision > 0)
&& active > 0
&& active <= ACTIVE_NEIGHBOR_LIMIT
&& backup <= 4
&& active + backup == routes.len(),
"sparse fanout gateway topology lease is invalid"
);
}
let selected =
(sparse && !gateway_lease).then(|| selected_device_ids(&members, &local_device_id));
let projected = members
.into_iter()
.filter(|member| {
member.node_id.is_some()
&& selected
.as_ref()
.is_none_or(|selected| selected.contains(&member.device_id))
})
.map(|member| member.value)
.collect::<Vec<_>>();
let overlap_neighbors = projected
.iter()
.filter(|peer| peer.get("topologyRole").and_then(Value::as_str) == Some("backup"))
.count();
let active_neighbors = if gateway_lease {
projected
.iter()
.filter(|peer| peer.get("topologyRole").and_then(Value::as_str) == Some("active"))
.count()
} else {
projected.len()
};
let diagnostics = self.diagnostics.entry(capability).or_default();
diagnostics.configured_revisions += 1;
diagnostics.active_neighbors = active_neighbors as u64;
diagnostics.overlap_neighbors = overlap_neighbors as u64;
diagnostics.max_forward_peers = active_neighbors.min(MAX_OVERLAY_FORWARD_PEERS) as u64;
diagnostics.roster_members = history
.back()
.map(|membership| membership.by_device.len() as u64)
.unwrap_or(1);
diagnostics.topology_revision = revision;
Ok((
projected.clone(),
SparseFanoutProjection {
topology_revision: revision,
member_count: history
.back()
.map(|membership| membership.by_device.len())
.unwrap_or(1),
active_neighbor_count: active_neighbors,
sparse,
},
))
}
pub fn prepare(
&mut self,
capability: &str,
payload: &[u8],
now_ms: u64,
) -> Result<SparseFanoutSigningRequest> {
ensure!(
payload.len() <= MAX_OVERLAY_PAYLOAD_BYTES,
"sparse fanout payload exceeds {MAX_OVERLAY_PAYLOAD_BYTES} bytes"
);
let capability = required_part(capability, "capability", 192)?;
let membership = self
.memberships
.get(&capability)
.and_then(|history| history.back())
.ok_or_else(|| anyhow!("sparse fanout capability is not configured"))?;
let sequence = self.sequences.entry(capability.clone()).or_default();
*sequence = sequence.saturating_add(1);
let sender_epoch_ms = *self
.sender_epochs
.entry(capability.clone())
.or_insert(now_ms);
let author_peer_id = membership.local_node_id.clone();
let message_id = canonical_message_id(
&membership.local_device_id,
&author_peer_id,
membership.revision,
sender_epoch_ms,
*sequence,
);
let unsigned = UnsignedEnvelope {
protocol: PROTOCOL_VERSION.to_string(),
capability: capability.clone(),
topology_revision: membership.revision,
author_device_id: membership.local_device_id.clone(),
author_peer_id: author_peer_id.clone(),
sender_epoch_ms,
sequence: *sequence,
message_id,
issued_at_ms: now_ms,
expires_at_ms: now_ms.saturating_add(OVERLAY_MESSAGE_LIFETIME_MS),
max_hops: MAX_OVERLAY_HOPS,
payload_hash: blake3::hash(payload).to_hex().to_string(),
payload: base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(payload),
};
let signing_input = serde_json::to_vec(&unsigned).context("encode fanout signing input")?;
let request_id = format!("{}:{}", unsigned.message_id, unsigned.topology_revision);
self.pending.insert(
request_id.clone(),
PendingEnvelope {
capability,
unsigned,
signing_input: signing_input.clone(),
payload_bytes: payload.len(),
},
);
self.pending_order.push_back(request_id.clone());
while self.pending_order.len() > MAX_PENDING_SIGNATURES {
if let Some(retired) = self.pending_order.pop_front() {
self.pending.remove(&retired);
}
}
Ok(SparseFanoutSigningRequest {
request_id,
signing_input,
})
}
pub fn finalize(&mut self, request_id: &str, signature: &str, now_ms: u64) -> Result<Vec<u8>> {
let pending = self.pending.remove(request_id).ok_or_else(|| {
anyhow!("sparse fanout signing request is missing or already consumed")
})?;
ensure!(
now_ms <= pending.unsigned.expires_at_ms,
"sparse fanout signing request expired"
);
let membership = self
.memberships
.get(&pending.capability)
.and_then(|history| history.back())
.ok_or_else(|| anyhow!("sparse fanout capability is not configured"))?;
ensure!(
membership.revision == pending.unsigned.topology_revision,
"sparse fanout topology changed while signing"
);
verify_signature(
&membership.local_device_key_x,
&pending.signing_input,
signature,
)?;
let encoded = serde_json::to_vec(&SignedEnvelope {
marker: 1,
unsigned: pending.unsigned,
hop: 0,
signature: signature.to_string(),
})
.context("encode sparse fanout envelope")?;
let diagnostics = self.diagnostics.entry(pending.capability).or_default();
diagnostics.encoded += 1;
diagnostics.payload_bytes = diagnostics
.payload_bytes
.saturating_add(pending.payload_bytes as u64);
Ok(encoded)
}
pub fn accept(
&mut self,
capability: &str,
source_peer_id: &str,
encoded: &[u8],
now_ms: u64,
) -> Result<SparseFanoutDecision> {
let capability = required_part(capability, "capability", 192)?;
let result = self.accept_inner(&capability, source_peer_id, encoded, now_ms);
if result.is_err() {
self.diagnostics.entry(capability).or_default().rejected += 1;
}
result
}
fn accept_inner(
&mut self,
capability: &str,
source_peer_id: &str,
encoded: &[u8],
now_ms: u64,
) -> Result<SparseFanoutDecision> {
ensure!(
encoded.len() <= MAX_OVERLAY_ENVELOPE_BYTES,
"sparse fanout envelope is too large"
);
self.charge_ingress(capability, encoded.len(), now_ms)?;
let mut envelope: SignedEnvelope =
serde_json::from_slice(encoded).context("decode sparse fanout envelope")?;
ensure!(
envelope.marker == 1,
"unsupported sparse fanout envelope marker"
);
ensure!(
envelope.unsigned.protocol == PROTOCOL_VERSION,
"unsupported sparse fanout protocol"
);
ensure!(
envelope.unsigned.capability == capability,
"sparse fanout capability mismatch"
);
ensure!(
envelope.unsigned.issued_at_ms <= now_ms.saturating_add(5_000)
&& now_ms <= envelope.unsigned.expires_at_ms,
"sparse fanout envelope is outside its lifetime"
);
ensure!(
envelope
.unsigned
.expires_at_ms
.saturating_sub(envelope.unsigned.issued_at_ms)
<= OVERLAY_MESSAGE_LIFETIME_MS,
"sparse fanout envelope lifetime is too large"
);
ensure!(
envelope.unsigned.max_hops <= MAX_OVERLAY_HOPS
&& envelope.hop <= envelope.unsigned.max_hops,
"sparse fanout hop bound exceeded"
);
let source_peer_id = required_part(source_peer_id, "source peer id", 256)?;
let history = self
.memberships
.get(capability)
.ok_or_else(|| anyhow!("sparse fanout capability is not configured"))?;
let membership = history
.iter()
.find(|membership| membership.revision == envelope.unsigned.topology_revision)
.ok_or_else(|| anyhow!("sparse fanout topology revision is not current"))?;
let current = history
.back()
.ok_or_else(|| anyhow!("sparse fanout capability is not configured"))?;
if membership.revision != current.revision {
ensure!(
current.by_node.contains_key(&source_peer_id)
&& current
.by_device
.contains_key(&envelope.unsigned.author_device_id),
"sparse fanout stale generation cannot revive a retired edge"
);
let previous_author = membership
.by_device
.get(&envelope.unsigned.author_device_id)
.ok_or_else(|| anyhow!("sparse fanout author is not an avenue member"))?;
let current_author = current
.by_device
.get(&envelope.unsigned.author_device_id)
.ok_or_else(|| anyhow!("sparse fanout author is not an avenue member"))?;
ensure!(
previous_author.device_key_x == current_author.device_key_x,
"sparse fanout stale author key was rotated"
);
let current_author_node = current
.by_node
.get(&envelope.unsigned.author_peer_id)
.ok_or_else(|| anyhow!("sparse fanout stale author node binding changed"))?;
ensure!(
current_author_node.device_id == envelope.unsigned.author_device_id
&& current_author_node.device_key_x == current_author.device_key_x,
"sparse fanout stale author node binding changed"
);
}
ensure!(
membership.by_node.contains_key(&source_peer_id),
"sparse fanout source peer is not an authorized avenue member"
);
ensure!(
membership
.by_device
.contains_key(&envelope.unsigned.author_device_id),
"sparse fanout author is not an avenue member"
);
let author = membership
.by_device
.get(&envelope.unsigned.author_device_id)
.ok_or_else(|| anyhow!("sparse fanout author is not an avenue member"))?;
let author_peer_id = match membership.by_node.get(&envelope.unsigned.author_peer_id) {
Some(bound) => {
ensure!(
bound.device_id == envelope.unsigned.author_device_id,
"sparse fanout author peer is bound to another device"
);
envelope.unsigned.author_peer_id.clone()
}
None => String::new(),
};
ensure!(
envelope.unsigned.message_id
== canonical_message_id(
&envelope.unsigned.author_device_id,
&envelope.unsigned.author_peer_id,
envelope.unsigned.topology_revision,
envelope.unsigned.sender_epoch_ms,
envelope.unsigned.sequence,
),
"sparse fanout message id is not canonical"
);
let signing_input =
serde_json::to_vec(&envelope.unsigned).context("encode fanout verification input")?;
verify_signature(&author.device_key_x, &signing_input, &envelope.signature)?;
let payload = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(&envelope.unsigned.payload)
.context("decode sparse fanout payload")?;
ensure!(
payload.len() <= MAX_OVERLAY_PAYLOAD_BYTES,
"sparse fanout payload exceeds {MAX_OVERLAY_PAYLOAD_BYTES} bytes"
);
ensure!(
blake3::hash(&payload).to_hex().as_str() == envelope.unsigned.payload_hash,
"sparse fanout payload hash mismatch"
);
let seen_key = format!("{capability}:{}", envelope.unsigned.message_id);
if !self.seen.insert(seen_key.clone()) {
self.diagnostics
.entry(capability.to_string())
.or_default()
.duplicate_drops += 1;
bail!("sparse fanout duplicate message");
}
self.seen_order.push_back(seen_key);
while self.seen_order.len() > MAX_SEEN_MESSAGES {
if let Some(retired) = self.seen_order.pop_front() {
self.seen.remove(&retired);
}
}
let replay_key = format!(
"{capability}:{}:{}:{}",
envelope.unsigned.topology_revision,
envelope.unsigned.author_device_id,
envelope.unsigned.sender_epoch_ms,
);
if !self
.author_replay
.entry(replay_key)
.or_default()
.accept(envelope.unsigned.sequence)
{
self.diagnostics
.entry(capability.to_string())
.or_default()
.replay_window_drops += 1;
bail!("sparse fanout author replay window rejected the sequence");
}
let forward = if envelope.hop < envelope.unsigned.max_hops {
envelope.hop += 1;
Some(serde_json::to_vec(&envelope).context("encode forwarded fanout envelope")?)
} else {
None
};
let forward_peer_ids = forward
.as_ref()
.map(|_| directed_forward_peer_ids(current, &envelope.unsigned.author_device_id))
.unwrap_or_default();
let forward_peer_limit = forward_peer_ids.len().min(MAX_OVERLAY_FORWARD_PEERS);
if let Some(forward) = &forward {
self.authorize_forward(capability, forward.len(), forward_peer_limit, now_ms)?;
}
let diagnostics = self.diagnostics.entry(capability.to_string()).or_default();
diagnostics.accepted += 1;
if let Some(forward) = &forward {
diagnostics.forwarded += 1;
diagnostics.forwarded_bytes = diagnostics
.forwarded_bytes
.saturating_add((forward.len() as u64).saturating_mul(forward_peer_limit as u64));
}
Ok(SparseFanoutDecision {
payload,
author_device_id: envelope.unsigned.author_device_id,
author_peer_id,
topology_revision: envelope.unsigned.topology_revision,
hop: envelope.hop.saturating_sub(1),
forward_peer_ids,
forward_peer_limit,
forward_queue_capacity: MAX_FORWARD_QUEUE_ACTIONS,
forward,
})
}
pub fn diagnostics(&self, capability: &str) -> SparseFanoutDiagnostics {
self.diagnostics
.get(capability)
.cloned()
.unwrap_or_default()
}
pub fn record_forward_queue_drop(&mut self, capability: &str, count: u64) -> Result<()> {
let capability = required_part(capability, "capability", 192)?;
let diagnostics = self.diagnostics.entry(capability).or_default();
diagnostics.forward_queue_drops = diagnostics.forward_queue_drops.saturating_add(count);
Ok(())
}
}
fn verify_signature(public_key_x: &str, message: &[u8], signature: &str) -> Result<()> {
let public_key = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(public_key_x)
.context("decode sparse fanout device public key")?;
let public_key: [u8; 32] = public_key
.try_into()
.map_err(|_| anyhow!("sparse fanout device public key has an invalid length"))?;
let public_key = VerifyingKey::from_bytes(&public_key)
.map_err(|_| anyhow!("sparse fanout device public key is invalid"))?;
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(signature)
.context("decode sparse fanout signature")?;
let signature: [u8; 64] = signature
.try_into()
.map_err(|_| anyhow!("sparse fanout signature has an invalid length"))?;
public_key
.verify(message, &Signature::from_bytes(&signature))
.map_err(|_| anyhow!("sparse fanout author signature is invalid"))
}
#[cfg(test)]
mod tests {
use super::*;
use ed25519_dalek::{Signer, SigningKey};
fn device_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 device_key_x(key: &SigningKey) -> String {
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(key.verifying_key().as_bytes())
}
fn peer(index: usize, secret: &iroh::SecretKey, key: &SigningKey) -> Value {
serde_json::json!({
"deviceId": format!("device-{index:03}"),
"nodeId": secret.public().to_string(),
"devicePublicKeyJwk": {
"kty": "OKP",
"crv": "Ed25519",
"x": device_key_x(key),
},
"online": true,
"ticket": format!("ticket-{index}"),
})
}
#[test]
fn sparse_ring_is_symmetric_bounded_and_locally_stable() {
let secrets = (0..MAX_REVIEWED_MEMBERS)
.map(|_| iroh::SecretKey::generate())
.collect::<Vec<_>>();
let device_keys = (0..MAX_REVIEWED_MEMBERS)
.map(device_key)
.collect::<Vec<_>>();
let roster = secrets
.iter()
.zip(&device_keys)
.enumerate()
.map(|(index, (secret, key))| peer(index, secret, key))
.collect::<Vec<_>>();
let mut adjacency = HashMap::<String, HashSet<String>>::new();
for (index, secret) in secrets.iter().enumerate() {
let local_device_id = format!("device-{index:03}");
let peers = roster
.iter()
.filter(|value| value["deviceId"] != local_device_id)
.cloned()
.collect();
let mut state = SparseFanoutState::default();
let (selected, projection) = state
.configure(
"room:test",
1,
&local_device_id,
&secret.public().to_string(),
&device_key_x(&device_keys[index]),
peers,
true,
)
.unwrap();
assert!(projection.sparse);
assert_eq!(selected.len(), ACTIVE_NEIGHBOR_LIMIT);
adjacency.insert(
local_device_id,
selected
.iter()
.filter_map(|value| value["deviceId"].as_str().map(ToString::to_string))
.collect(),
);
}
let edge_count = adjacency.values().map(HashSet::len).sum::<usize>() / 2;
assert_eq!(edge_count, 198);
for (device, neighbors) in &adjacency {
for neighbor in neighbors {
assert!(adjacency
.get(neighbor)
.is_some_and(|set| set.contains(device)));
}
}
}
#[test]
fn fifty_member_broadcast_reaches_every_member_with_bounded_amplification() {
const MEMBERS: usize = 50;
let secrets = (0..MEMBERS)
.map(|_| iroh::SecretKey::generate())
.collect::<Vec<_>>();
let device_keys = (0..MEMBERS).map(device_key).collect::<Vec<_>>();
let roster = secrets
.iter()
.zip(&device_keys)
.enumerate()
.map(|(index, (secret, key))| peer(index, secret, key))
.collect::<Vec<_>>();
let mut routes = Vec::<HashSet<String>>::with_capacity(MEMBERS);
let mut states = Vec::<SparseFanoutState>::with_capacity(MEMBERS);
for index in 0..MEMBERS {
let mut state = SparseFanoutState::default();
let projected = state
.configure(
"room:fifty",
1,
&format!("device-{index:03}"),
&secrets[index].public().to_string(),
&device_key_x(&device_keys[index]),
roster
.iter()
.enumerate()
.filter(|(peer_index, _)| *peer_index != index)
.map(|(_, value)| value.clone())
.collect(),
true,
)
.unwrap()
.0
.into_iter()
.filter_map(|value| value["nodeId"].as_str().map(ToString::to_string))
.collect();
routes.push(projected);
states.push(state);
}
let request = states[0]
.prepare("room:fifty", b"latest-state", 1_000)
.unwrap();
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(device_keys[0].sign(&request.signing_input).to_bytes());
let encoded = states[0]
.finalize(&request.request_id, &signature, 1_001)
.unwrap();
let by_node = secrets
.iter()
.enumerate()
.map(|(index, secret)| (secret.public().to_string(), index))
.collect::<HashMap<_, _>>();
let mut deliveries = HashSet::new();
let mut queue = routes[0]
.iter()
.map(|target| {
(
secrets[0].public().to_string(),
target.clone(),
encoded.clone(),
)
})
.collect::<VecDeque<_>>();
let mut transmissions = queue.len();
while let Some((source_node_id, target_node_id, envelope)) = queue.pop_front() {
let target_index = by_node[&target_node_id];
let decision = match states[target_index].accept(
"room:fifty",
&source_node_id,
&envelope,
1_001,
) {
Ok(decision) => decision,
Err(error) if error.to_string().contains("duplicate") => continue,
Err(error) => panic!("broadcast failed at member {target_index}: {error:#}"),
};
assert_eq!(decision.payload, b"latest-state");
assert!(deliveries.insert(target_index));
if let Some(forward) = decision.forward {
for peer_id in decision.forward_peer_ids {
assert!(
routes[target_index].contains(&peer_id),
"Rust selected a peer outside the admitted route set",
);
queue.push_back((target_node_id.clone(), peer_id, forward.clone()));
transmissions += 1;
}
}
}
assert_eq!(deliveries.len(), MEMBERS - 1);
assert!(
transmissions * 4 <= (MEMBERS - 1) * 5,
"{transmissions} transmissions exceed the 1.25x amplification ceiling",
);
}
#[test]
fn gateway_active_and_backup_routes_are_one_generation_fenced_desired_set() {
let local = iroh::SecretKey::generate();
let local_key = device_key(0);
let secrets = (1..MAX_REVIEWED_MEMBERS)
.map(|_| iroh::SecretKey::generate())
.collect::<Vec<_>>();
let keys = (1..MAX_REVIEWED_MEMBERS)
.map(device_key)
.collect::<Vec<_>>();
let mut roster = secrets
.iter()
.zip(&keys)
.enumerate()
.map(|(offset, (secret, key))| {
let mut value = peer(offset + 1, secret, key);
if offset < ACTIVE_NEIGHBOR_LIMIT + 4 {
value["topologyRevision"] = Value::from(1);
value["topologyRole"] = Value::String(
if offset < ACTIVE_NEIGHBOR_LIMIT {
"active"
} else {
"backup"
}
.to_string(),
);
} else {
value.as_object_mut().unwrap().remove("nodeId");
value.as_object_mut().unwrap().remove("ticket");
}
value
})
.collect::<Vec<_>>();
let mut state = SparseFanoutState::default();
let (projected, projection) = state
.configure(
"room:test",
1,
"device-000",
&local.public().to_string(),
&device_key_x(&local_key),
roster.clone(),
true,
)
.unwrap();
assert!(projection.sparse);
assert_eq!(projection.member_count, MAX_REVIEWED_MEMBERS);
assert_eq!(projected.len(), ACTIVE_NEIGHBOR_LIMIT + 4);
roster.swap(0, 8);
roster[0]["nodeId"] = Value::String(secrets[8].public().to_string());
roster[0]["ticket"] = Value::String("ticket-next".to_string());
roster[0]["topologyRevision"] = Value::from(2);
roster[0]["topologyRole"] = Value::String("active".to_string());
roster[8].as_object_mut().unwrap().remove("nodeId");
roster[8].as_object_mut().unwrap().remove("ticket");
for value in roster
.iter_mut()
.filter(|value| value.get("nodeId").is_some())
{
value["topologyRevision"] = Value::from(2);
}
let (next, _) = state
.configure(
"room:test",
2,
"device-000",
&local.public().to_string(),
&device_key_x(&local_key),
roster.clone(),
true,
)
.unwrap();
assert_eq!(next.len(), ACTIVE_NEIGHBOR_LIMIT + 4);
assert!(state
.configure(
"room:test",
1,
"device-000",
&local.public().to_string(),
&device_key_x(&local_key),
roster,
true,
)
.unwrap_err()
.to_string()
.contains("revision is stale"));
}
#[test]
fn signed_envelope_rejects_tamper_replay_and_stale_membership() {
let local = iroh::SecretKey::generate();
let remote = iroh::SecretKey::generate();
let local_key = device_key(0);
let remote_key = device_key(1);
let mut sender = SparseFanoutState::default();
let mut receiver = SparseFanoutState::default();
let local_peer = peer(0, &local, &local_key);
let remote_peer = peer(1, &remote, &remote_key);
sender
.configure(
"room:test",
7,
"device-000",
&local.public().to_string(),
&device_key_x(&local_key),
vec![remote_peer.clone()],
true,
)
.unwrap();
receiver
.configure(
"room:test",
7,
"device-001",
&remote.public().to_string(),
&device_key_x(&remote_key),
vec![local_peer],
true,
)
.unwrap();
let request = sender
.prepare("room:test", br#"{"kind":"channel"}"#, 1_000)
.unwrap();
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(local_key.sign(&request.signing_input).to_bytes());
let encoded = sender
.finalize(&request.request_id, &signature, 1_001)
.unwrap();
let accepted = receiver
.accept("room:test", &local.public().to_string(), &encoded, 1_001)
.unwrap();
assert_eq!(accepted.payload, br#"{"kind":"channel"}"#);
assert!(accepted.forward.is_some());
assert!(receiver
.accept("room:test", &local.public().to_string(), &encoded, 1_002)
.unwrap_err()
.to_string()
.contains("duplicate"));
let mut tampered: Value = serde_json::from_slice(&encoded).unwrap();
tampered["payloadHash"] = Value::String("00".repeat(32));
let tampered = serde_json::to_vec(&tampered).unwrap();
let mut fresh_receiver = SparseFanoutState::default();
fresh_receiver
.configure(
"room:test",
7,
"device-001",
&remote.public().to_string(),
&device_key_x(&remote_key),
vec![peer(0, &local, &local_key)],
true,
)
.unwrap();
assert!(fresh_receiver
.accept("room:test", &local.public().to_string(), &tampered, 1_001)
.is_err());
assert!(fresh_receiver
.accept("room:test", &local.public().to_string(), &encoded, 40_001)
.is_err());
receiver
.configure(
"room:test",
8,
"device-001",
&remote.public().to_string(),
&device_key_x(&remote_key),
vec![],
true,
)
.unwrap();
assert!(receiver
.accept("room:test", &local.public().to_string(), &encoded, 1_003)
.unwrap_err()
.to_string()
.contains("stale generation cannot revive"));
}
#[test]
fn canonical_message_ids_survive_restart_and_reject_dedupe_poisoning() {
let local = iroh::SecretKey::generate();
let remote = iroh::SecretKey::generate();
let local_key = device_key(0);
let remote_key = device_key(1);
let local_peer = peer(0, &local, &local_key);
let remote_peer = peer(1, &remote, &remote_key);
let configure_sender = |state: &mut SparseFanoutState| {
state
.configure(
"room:test",
7,
"device-000",
&local.public().to_string(),
&device_key_x(&local_key),
vec![remote_peer.clone()],
true,
)
.unwrap();
};
let mut receiver = SparseFanoutState::default();
receiver
.configure(
"room:test",
7,
"device-001",
&remote.public().to_string(),
&device_key_x(&remote_key),
vec![local_peer],
true,
)
.unwrap();
let mut first_sender = SparseFanoutState::default();
configure_sender(&mut first_sender);
let first_request = first_sender.prepare("room:test", b"first", 1_000).unwrap();
let first_signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(local_key.sign(&first_request.signing_input).to_bytes());
let first = first_sender
.finalize(&first_request.request_id, &first_signature, 1_001)
.unwrap();
let mut restarted_sender = SparseFanoutState::default();
configure_sender(&mut restarted_sender);
let restarted_request = restarted_sender
.prepare("room:test", b"after-restart", 2_000)
.unwrap();
let restarted_signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(local_key.sign(&restarted_request.signing_input).to_bytes());
let restarted = restarted_sender
.finalize(&restarted_request.request_id, &restarted_signature, 2_001)
.unwrap();
let first_envelope: SignedEnvelope = serde_json::from_slice(&first).unwrap();
let restarted_envelope: SignedEnvelope = serde_json::from_slice(&restarted).unwrap();
assert_eq!(first_envelope.unsigned.sequence, 1);
assert_eq!(restarted_envelope.unsigned.sequence, 1);
assert_ne!(
first_envelope.unsigned.message_id,
restarted_envelope.unsigned.message_id
);
receiver
.accept("room:test", &local.public().to_string(), &first, 1_001)
.unwrap();
receiver
.accept("room:test", &local.public().to_string(), &restarted, 2_001)
.unwrap();
let mut poisoned = restarted_envelope;
poisoned.unsigned.message_id = first_envelope.unsigned.message_id;
poisoned.signature = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(
local_key
.sign(&serde_json::to_vec(&poisoned.unsigned).unwrap())
.to_bytes(),
);
let mut fresh_receiver = SparseFanoutState::default();
fresh_receiver
.configure(
"room:test",
7,
"device-001",
&remote.public().to_string(),
&device_key_x(&remote_key),
vec![peer(0, &local, &local_key)],
true,
)
.unwrap();
assert!(fresh_receiver
.accept(
"room:test",
&local.public().to_string(),
&serde_json::to_vec(&poisoned).unwrap(),
2_001,
)
.unwrap_err()
.to_string()
.contains("message id is not canonical"));
let mut impersonated: SignedEnvelope = serde_json::from_slice(&restarted).unwrap();
impersonated.unsigned.author_peer_id = remote.public().to_string();
impersonated.unsigned.message_id = canonical_message_id(
&impersonated.unsigned.author_device_id,
&impersonated.unsigned.author_peer_id,
impersonated.unsigned.topology_revision,
impersonated.unsigned.sender_epoch_ms,
impersonated.unsigned.sequence,
);
impersonated.signature = base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(
local_key
.sign(&serde_json::to_vec(&impersonated.unsigned).unwrap())
.to_bytes(),
);
assert!(fresh_receiver
.accept(
"room:test",
&local.public().to_string(),
&serde_json::to_vec(&impersonated).unwrap(),
2_001,
)
.unwrap_err()
.to_string()
.contains("author peer is bound to another device"));
}
#[test]
fn stale_topology_rejects_a_rotated_author_key() {
let local = iroh::SecretKey::generate();
let remote = iroh::SecretKey::generate();
let old_key = device_key(0);
let rotated_key = device_key(2);
let remote_key = device_key(1);
let mut sender = SparseFanoutState::default();
sender
.configure(
"room:test",
7,
"device-000",
&local.public().to_string(),
&device_key_x(&old_key),
vec![peer(1, &remote, &remote_key)],
true,
)
.unwrap();
let request = sender.prepare("room:test", b"old-key", 1_000).unwrap();
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(old_key.sign(&request.signing_input).to_bytes());
let encoded = sender
.finalize(&request.request_id, &signature, 1_001)
.unwrap();
let mut receiver = SparseFanoutState::default();
receiver
.configure(
"room:test",
7,
"device-001",
&remote.public().to_string(),
&device_key_x(&remote_key),
vec![peer(0, &local, &old_key)],
true,
)
.unwrap();
receiver
.configure(
"room:test",
8,
"device-001",
&remote.public().to_string(),
&device_key_x(&remote_key),
vec![peer(0, &local, &rotated_key)],
true,
)
.unwrap();
assert!(receiver
.accept("room:test", &local.public().to_string(), &encoded, 1_002)
.unwrap_err()
.to_string()
.contains("author key was rotated"));
}
#[test]
fn replay_and_traffic_windows_are_bounded_before_forwarding() {
let mut replay = AuthorReplayWindow::default();
assert!(replay.accept(10));
assert!(
replay.accept(9),
"bounded out-of-order delivery remains valid"
);
assert!(!replay.accept(9));
assert!(!replay.accept(10));
assert!(replay.accept(200));
assert!(
!replay.accept(1),
"retired sequence windows cannot be revived"
);
let mut state = SparseFanoutState::default();
for _ in 0..8 {
state
.charge_ingress("room:test", MAX_OVERLAY_ENVELOPE_BYTES, 1_000)
.unwrap();
}
assert!(state
.charge_ingress("room:test", 1, 1_000)
.unwrap_err()
.to_string()
.contains("ingress budget"));
state.record_forward_queue_drop("room:test", 3).unwrap();
let diagnostics = state.diagnostics("room:test");
assert_eq!(diagnostics.budget_drops, 1);
assert_eq!(diagnostics.forward_queue_drops, 3);
}
#[test]
fn forward_actions_and_peer_count_are_rust_bounded() {
let secrets = (0..9)
.map(|_| iroh::SecretKey::generate())
.collect::<Vec<_>>();
let keys = (0..9).map(device_key).collect::<Vec<_>>();
let roster = secrets
.iter()
.zip(&keys)
.enumerate()
.map(|(index, (secret, key))| peer(index, secret, key))
.collect::<Vec<_>>();
let mut sender = SparseFanoutState::default();
let mut receiver = SparseFanoutState::default();
sender
.configure(
"room:test",
7,
"device-000",
&secrets[0].public().to_string(),
&device_key_x(&keys[0]),
roster.iter().skip(1).cloned().collect(),
true,
)
.unwrap();
receiver
.configure(
"room:test",
7,
"device-001",
&secrets[1].public().to_string(),
&device_key_x(&keys[1]),
roster
.iter()
.enumerate()
.filter(|(index, _)| *index != 1)
.map(|(_, value)| value.clone())
.collect(),
true,
)
.unwrap();
for index in 0..MAX_FORWARD_ACTIONS_PER_WINDOW {
let request = sender
.prepare("room:test", format!("message-{index}").as_bytes(), 1_000)
.unwrap();
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(keys[0].sign(&request.signing_input).to_bytes());
let encoded = sender
.finalize(&request.request_id, &signature, 1_001)
.unwrap();
let decision = receiver
.accept(
"room:test",
&secrets[0].public().to_string(),
&encoded,
1_001,
)
.unwrap();
assert_eq!(decision.forward_peer_limit, 1);
assert_eq!(
decision.forward_peer_ids,
vec![secrets[2].public().to_string()]
);
assert_eq!(decision.forward_queue_capacity, MAX_FORWARD_QUEUE_ACTIONS);
}
let request = sender.prepare("room:test", b"over-budget", 1_000).unwrap();
let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(keys[0].sign(&request.signing_input).to_bytes());
let encoded = sender
.finalize(&request.request_id, &signature, 1_001)
.unwrap();
assert!(receiver
.accept(
"room:test",
&secrets[0].public().to_string(),
&encoded,
1_001,
)
.unwrap_err()
.to_string()
.contains("forward action budget"));
assert_eq!(receiver.diagnostics("room:test").budget_drops, 1);
assert!(receiver
.accept(
"room:test",
&secrets[0].public().to_string(),
&vec![0; MAX_OVERLAY_ENVELOPE_BYTES + 1],
2_001,
)
.unwrap_err()
.to_string()
.contains("envelope is too large"));
}
}