use std::collections::{BTreeMap, HashMap, HashSet};
use std::sync::Arc;
use std::time::{Duration, Instant};
use anyhow::Result;
use tokio::sync::{broadcast, watch, Mutex};
use tokio::task::JoinSet;
#[cfg(test)]
use super::WIRE_PREFIX;
use super::{
decode_wire, encode_wire, mesh_error, mesh_owner_error, AckKind, TicketInvite,
TicketMeshErrorCode, TicketMeshMember, TicketMeshOptions, TicketMeshRoster, Wire,
ADMISSION_LEASE, ADMISSION_REFRESH, GUEST_GRACE, HELLO_RETRY, INVITE_REFRESH,
INVITE_REFRESH_RETRY, MAX_RETRY, ROSTER_RESEND,
};
use crate::client::{NativePeerDataEvent, PeerSessionSnapshot};
use crate::session_token::now_unix_ms;
use crate::Client;
const CONTROL_SEND_TIMEOUT: Duration = Duration::from_secs(2);
const CONTROL_ACK_TIMEOUT: Duration = Duration::from_secs(2);
const CONTROL_NOTICE_RETRY: Duration = Duration::from_millis(250);
enum Role {
Issuer {
members: BTreeMap<String, TicketMeshMember>,
unseen_since: HashMap<String, Instant>,
revision: u64,
last_roster_sent: Option<Instant>,
},
Guest {
roster: TicketMeshRoster,
local_ticket: String,
last_hello_sent: Option<Instant>,
local_ticket_issued_at: Instant,
last_ticket_refresh_attempt: Option<Instant>,
},
}
struct Retry {
next: Instant,
delay: Duration,
}
pub struct TicketMesh {
client: Arc<Client>,
id: String,
invite: watch::Receiver<String>,
peer_updates: watch::Sender<Vec<PeerSessionSnapshot>>,
peers: watch::Receiver<Vec<PeerSessionSnapshot>>,
issuer: bool,
stop: watch::Sender<bool>,
task: Mutex<Option<tokio::task::JoinHandle<()>>>,
}
impl Drop for TicketMesh {
fn drop(&mut self) {
if *self.stop.borrow() {
return;
}
self.peer_updates.send_replace(Vec::new());
self.stop.send_replace(true);
if let Ok(mut task) = self.task.try_lock() {
if let Some(task) = task.take() {
task.abort();
}
}
let client = self.client.clone();
let scope = format!("ticket:{}", self.id);
if let Ok(runtime) = tokio::runtime::Handle::try_current() {
runtime.spawn(async move {
client.revoke_tokens_by_scope(&scope).await;
client
.session_token_registry
.release_ticket_peer_admission(&scope);
client
.session_token_registry
.release_ticket_peer_limit(&scope);
});
}
}
}
impl TicketMesh {
pub async fn issue(client: Arc<Client>, id: &str, options: TicketMeshOptions) -> Result<Self> {
let invite = client.issue_ticket_invite(id, options).await?;
Ok(Self::start(
client,
invite,
Role::Issuer {
members: BTreeMap::new(),
unseen_since: HashMap::new(),
revision: 1,
last_roster_sent: None,
},
true,
))
}
pub async fn join(client: Arc<Client>, invitation: &str) -> Result<Self> {
let invite = TicketInvite::parse(invitation)?;
if !invite.is_mesh() {
return Err(mesh_error(
TicketMeshErrorCode::InvalidInvitation,
"ticket invitation is not a mesh",
));
}
let local_node = client.current_node_id().await.ok_or_else(|| {
mesh_error(
TicketMeshErrorCode::Unavailable,
"ticket mesh requires an initialized endpoint",
)
})?;
let roster = TicketMeshRoster::new(&invite, &local_node)?;
let scope = format!("ticket:{}", invite.id());
client
.session_token_registry
.set_ticket_peer_limit(&scope, invite.max_peers())
.map_err(mesh_owner_error)?;
let local_ticket = match async {
client
.session_token_registry
.require_scope_peer_admission(&scope)
.map_err(anyhow::Error::msg)?;
client
.session_token_registry
.update_scope_peer_admission(
&scope,
1,
now_unix_ms() + ADMISSION_LEASE.as_millis() as u64,
roster.allowed_nodes(),
)
.map_err(anyhow::Error::msg)?;
let local_ticket = client.endpoint_ticket_with_token(&scope, 0).await?;
TicketInvite::new(
invite.id(),
&local_ticket,
TicketMeshOptions {
max_peers: invite.max_peers(),
},
)?;
Ok::<_, anyhow::Error>(local_ticket)
}
.await
{
Ok(ticket) => ticket,
Err(error) => {
client
.session_token_registry
.release_ticket_peer_admission(&scope);
client
.session_token_registry
.release_ticket_peer_limit(&scope);
return Err(error);
}
};
Ok(Self::start(
client,
invite,
Role::Guest {
roster,
local_ticket,
last_hello_sent: None,
local_ticket_issued_at: Instant::now(),
last_ticket_refresh_attempt: None,
},
false,
))
}
fn start(client: Arc<Client>, invite: TicketInvite, role: Role, issuer: bool) -> Self {
let (peers_tx, peers) = watch::channel(Vec::new());
let (invite_tx, current_invite) = watch::channel(invite.encode());
let (stop, stop_rx) = watch::channel(false);
let actor = Actor {
client: client.clone(),
invite: invite.clone(),
invite_tx,
invite_issued_at: Instant::now(),
last_invite_refresh_attempt: None,
role,
peers: peers_tx.clone(),
stop: stop_rx,
events: client.subscribe_native_peer_data(),
retries: HashMap::new(),
in_flight: HashSet::new(),
dials: JoinSet::new(),
control_sends: JoinSet::new(),
admission_revision: 1,
last_admission_update: Instant::now(),
last_allowed_nodes: if issuer {
Vec::new()
} else {
vec![invite.issuer_node().expect("validated invitation")]
},
last_peer_projection: Vec::new(),
terminal: false,
};
let task = tokio::spawn(async move { actor.run().await });
Self {
client,
id: invite.id().to_string(),
invite: current_invite,
peer_updates: peers_tx.clone(),
peers,
issuer,
stop,
task: Mutex::new(Some(task)),
}
}
pub fn invite(&self) -> String {
self.invite.borrow().clone()
}
pub fn id(&self) -> &str {
&self.id
}
pub fn issuer_node(&self) -> Result<String> {
TicketInvite::parse(self.invite.borrow().as_str())?.issuer_node()
}
pub fn watch_invite(&self) -> watch::Receiver<String> {
self.invite.clone()
}
pub fn watch_peers(&self) -> watch::Receiver<Vec<PeerSessionSnapshot>> {
self.peers.clone()
}
pub async fn close(&self) {
let mut task = self.task.lock().await;
if task.is_none() {
return;
}
if !*self.stop.borrow() {
let (frame, nodes, expected_ack) = if self.issuer {
(
encode_wire(&Wire::Closed {
id: self.id.clone(),
}),
self.peers
.borrow()
.iter()
.filter_map(|peer| peer.node_id.clone())
.collect::<Vec<_>>(),
AckKind::Close,
)
} else {
(
encode_wire(&Wire::Leave {
id: self.id.clone(),
}),
self.issuer_node().into_iter().collect(),
AckKind::Leave,
)
};
let mut events = self.client.subscribe_native_peer_data();
let mut pending: HashSet<String> = nodes.iter().cloned().collect();
let scope = format!("ticket:{}", self.id);
let notify_until_acked = async {
while !pending.is_empty() {
let mut notices = JoinSet::new();
for node in pending.iter().cloned() {
let client = self.client.clone();
let frame = frame.clone();
notices.spawn(async move {
let _ = tokio::time::timeout(
CONTROL_SEND_TIMEOUT,
client.send_peer(&node, &frame),
)
.await;
});
}
while notices.join_next().await.is_some() {}
let retry_at = Instant::now() + CONTROL_NOTICE_RETRY;
while !pending.is_empty() {
let remaining = retry_at.saturating_duration_since(Instant::now());
if remaining.is_zero() {
break;
}
let event = match tokio::time::timeout(remaining, events.recv()).await {
Ok(Ok(event)) => event,
Ok(Err(broadcast::error::RecvError::Lagged(_))) => continue,
Ok(Err(broadcast::error::RecvError::Closed)) | Err(_) => break,
};
let Some(Wire::Ack { id, kind }) = decode_wire(&event.payload) else {
continue;
};
if id != self.id || kind != expected_ack {
continue;
}
if let Some(source) =
authenticated_ticket_source(&self.client, &event, &scope).await
{
pending.remove(&source);
}
}
}
};
let _ = tokio::time::timeout(CONTROL_ACK_TIMEOUT, notify_until_acked).await;
}
self.peer_updates.send_replace(Vec::new());
self.stop.send_replace(true);
if let Some(actor) = task.take() {
actor.abort();
let _ = actor.await;
}
let scope = format!("ticket:{}", self.id);
self.client.revoke_tokens_by_scope(&scope).await;
self.client
.session_token_registry
.release_ticket_peer_admission(&scope);
self.client
.session_token_registry
.release_ticket_peer_limit(&scope);
}
}
struct Actor {
client: Arc<Client>,
invite: TicketInvite,
invite_tx: watch::Sender<String>,
invite_issued_at: Instant,
last_invite_refresh_attempt: Option<Instant>,
role: Role,
peers: watch::Sender<Vec<PeerSessionSnapshot>>,
stop: watch::Receiver<bool>,
events: broadcast::Receiver<NativePeerDataEvent>,
retries: HashMap<String, Retry>,
in_flight: HashSet<String>,
dials: JoinSet<String>,
control_sends: JoinSet<()>,
admission_revision: u64,
last_admission_update: Instant,
last_allowed_nodes: Vec<String>,
last_peer_projection: Vec<u8>,
terminal: bool,
}
impl Actor {
async fn run(mut self) {
let mut tick = tokio::time::interval(Duration::from_secs(1));
loop {
tokio::select! {
_ = tick.tick() => self.settle().await,
event = self.events.recv() => {
if let Ok(event) = event {
self.handle_event(event).await;
if self.terminal { break; }
}
}
changed = self.stop.changed() => {
if changed.is_err() || *self.stop.borrow() { break; }
}
}
if *self.stop.borrow() {
break;
}
if self.terminal {
break;
}
}
if self.terminal {
let scope = format!("ticket:{}", self.invite.id());
self.client.revoke_tokens_by_scope(&scope).await;
self.client
.session_token_registry
.release_ticket_peer_admission(&scope);
self.client
.session_token_registry
.release_ticket_peer_limit(&scope);
}
}
async fn ticket_peers(&self) -> Vec<PeerSessionSnapshot> {
let scope = format!("ticket:{}", self.invite.id());
self.client
.peer_sessions()
.await
.into_iter()
.filter(|peer| {
peer.settled_ready
&& !peer.logical_session_terminal
&& peer.scopes.iter().any(|candidate| candidate == &scope)
&& peer.node_id.is_some()
})
.collect()
}
async fn handle_event(&mut self, event: NativePeerDataEvent) {
let Some(message) = decode_wire(&event.payload) else {
return;
};
let scope = format!("ticket:{}", self.invite.id());
let Some(source) = authenticated_ticket_source(&self.client, &event, &scope).await else {
return;
};
let mut reply = None;
match (&mut self.role, message) {
(
Role::Issuer {
members,
revision,
last_roster_sent,
..
},
Wire::Hello { id, ticket },
) if id == self.invite.id() => {
let member = TicketMeshMember {
node_id: source.clone(),
ticket,
};
let valid = TicketInvite::new(
self.invite.id(),
&member.ticket,
TicketMeshOptions {
max_peers: self.invite.max_peers(),
},
)
.and_then(|grant| grant.issuer_node())
.is_ok_and(|node| node == source);
if !valid
|| (!members.contains_key(&source)
&& members.len() + 1 >= self.invite.max_peers() as usize)
{
return;
}
if members
.get(&source)
.is_none_or(|old| old.ticket != member.ticket)
{
members.insert(source.clone(), member);
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
}
(
Role::Issuer {
members,
unseen_since,
revision,
last_roster_sent,
},
Wire::Leave { id },
) if id == self.invite.id() => {
if members.remove(&source).is_some() {
unseen_since.remove(&source);
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
reply = Some(AckKind::Leave);
}
(
Role::Guest { roster, .. },
Wire::Roster {
id,
revision,
members,
issuer_invite,
},
) if id == self.invite.id() => {
let renewed = if issuer_invite == self.invite.encode() {
None
} else {
match self.invite.renewed(&issuer_invite) {
Ok(invite) => Some(invite),
Err(_) => return,
}
};
if roster.accept(&source, revision, members).unwrap_or(false) {
if let Some(invite) = renewed {
self.invite_tx.send_replace(invite.encode());
self.invite = invite;
}
}
}
(Role::Guest { .. }, Wire::Closed { id })
if id == self.invite.id()
&& self.invite.issuer_node().ok().as_deref() == Some(source.as_str()) =>
{
reply = Some(AckKind::Close);
self.terminal = true;
self.peers.send_replace(Vec::new());
}
(_, Wire::Ack { .. }) => {}
_ => {}
}
if let Some(kind) = reply {
let frame = encode_wire(&Wire::Ack {
id: self.invite.id().to_string(),
kind,
});
let _ =
tokio::time::timeout(CONTROL_SEND_TIMEOUT, self.client.send_peer(&source, &frame))
.await;
}
}
async fn settle(&mut self) {
while let Some(done) = self.dials.try_join_next() {
if let Ok(node) = done {
self.in_flight.remove(&node);
}
}
while self.control_sends.try_join_next().is_some() {}
let mut peers = self.ticket_peers().await;
peers.sort_by(|left, right| left.node_id.cmp(&right.node_id));
let nodes: std::collections::HashSet<String> = peers
.iter()
.filter_map(|peer| peer.node_id.clone())
.collect();
let projection = serde_json::to_vec(&peers).unwrap_or_default();
if projection != self.last_peer_projection {
self.last_peer_projection = projection;
self.peers.send_replace(peers);
}
let now = Instant::now();
if matches!(self.role, Role::Issuer { .. })
&& now.duration_since(self.invite_issued_at) >= INVITE_REFRESH
&& self
.last_invite_refresh_attempt
.is_none_or(|attempt| now.duration_since(attempt) >= INVITE_REFRESH_RETRY)
{
self.last_invite_refresh_attempt = Some(now);
let scope = format!("ticket:{}", self.invite.id());
if let Ok(ticket) = self.client.endpoint_ticket_with_token(&scope, 0).await {
if let Ok(next) = TicketInvite::new(
self.invite.id(),
&ticket,
TicketMeshOptions {
max_peers: self.invite.max_peers(),
},
)
.and_then(|next| self.invite.renewed(&next.encode()))
{
self.invite_tx.send_replace(next.encode());
self.invite = next;
self.invite_issued_at = now;
if let Role::Issuer {
revision,
last_roster_sent,
..
} = &mut self.role
{
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
}
}
}
let refresh_guest_ticket = matches!(&self.role,
Role::Guest { local_ticket_issued_at, last_ticket_refresh_attempt, .. }
if now.duration_since(*local_ticket_issued_at) >= INVITE_REFRESH
&& last_ticket_refresh_attempt.is_none_or(|attempt| now.duration_since(attempt) >= INVITE_REFRESH_RETRY)
);
if refresh_guest_ticket {
if let Role::Guest {
last_ticket_refresh_attempt,
..
} = &mut self.role
{
*last_ticket_refresh_attempt = Some(now);
}
let scope = format!("ticket:{}", self.invite.id());
if let Ok(ticket) = self.client.endpoint_ticket_with_token(&scope, 0).await {
let local_node = self.client.current_node_id().await;
let valid = TicketInvite::new(
self.invite.id(),
&ticket,
TicketMeshOptions {
max_peers: self.invite.max_peers(),
},
)
.and_then(|invite| invite.issuer_node())
.is_ok_and(|node| local_node.as_deref() == Some(node.as_str()));
if valid {
if let Role::Guest {
local_ticket,
last_hello_sent,
local_ticket_issued_at,
..
} = &mut self.role
{
*local_ticket = ticket;
*last_hello_sent = None;
*local_ticket_issued_at = now;
}
}
}
}
let mut wanted = Vec::new();
let mut control_messages = Vec::new();
match &mut self.role {
Role::Issuer {
members,
unseen_since,
revision,
last_roster_sent,
} => {
let stale: Vec<_> = members
.keys()
.filter_map(|node| {
if nodes.contains(node) {
unseen_since.remove(node);
None
} else if now
.duration_since(*unseen_since.entry(node.clone()).or_insert(now))
>= GUEST_GRACE
{
Some(node.clone())
} else {
None
}
})
.collect();
for node in stale {
members.remove(&node);
unseen_since.remove(&node);
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
if last_roster_sent.is_none_or(|sent| now.duration_since(sent) >= ROSTER_RESEND) {
let frame = encode_wire(&Wire::Roster {
id: self.invite.id().to_string(),
revision: *revision,
members: members.values().cloned().collect(),
issuer_invite: self.invite.encode(),
});
for node in members.keys() {
control_messages.push((node.clone(), frame.clone()));
}
*last_roster_sent = Some(now);
}
}
Role::Guest {
roster,
local_ticket,
last_hello_sent,
..
} => {
let allowed = roster.allowed_nodes();
if allowed != self.last_allowed_nodes
|| now.duration_since(self.last_admission_update) >= ADMISSION_REFRESH
{
self.admission_revision = self.admission_revision.saturating_add(1);
let scope = format!("ticket:{}", self.invite.id());
if self
.client
.session_token_registry
.update_scope_peer_admission(
&scope,
self.admission_revision,
now_unix_ms() + ADMISSION_LEASE.as_millis() as u64,
allowed.clone(),
)
.is_err()
{
return;
}
self.last_allowed_nodes = allowed;
self.last_admission_update = now;
}
let issuer = match self.invite.issuer_node() {
Ok(node) => node,
Err(_) => return,
};
wanted.push((issuer.clone(), self.invite.endpoint_ticket().to_string()));
wanted.extend(
roster
.desired()
.into_iter()
.map(|member| (member.node_id, member.ticket)),
);
if nodes.contains(&issuer)
&& !roster.contains_local()
&& last_hello_sent.is_none_or(|sent| now.duration_since(sent) >= HELLO_RETRY)
{
let frame = encode_wire(&Wire::Hello {
id: self.invite.id().to_string(),
ticket: local_ticket.clone(),
});
control_messages.push((issuer, frame));
*last_hello_sent = Some(now);
}
}
}
for (node, frame) in control_messages {
let client = self.client.clone();
self.control_sends.spawn(async move {
let _ = tokio::time::timeout(CONTROL_SEND_TIMEOUT, client.send_peer(&node, &frame))
.await;
});
}
self.retries
.retain(|node, _| wanted.iter().any(|(desired, _)| desired == node));
for (node, ticket) in wanted {
if nodes.contains(&node) {
self.retries.remove(&node);
continue;
}
if self
.retries
.get(&node)
.is_some_and(|retry| now < retry.next)
|| self.in_flight.contains(&node)
{
continue;
}
let delay = self
.retries
.get(&node)
.map_or(Duration::from_millis(500), |retry| {
(retry.delay * 2).min(MAX_RETRY)
});
self.retries.insert(
node.clone(),
Retry {
next: now + delay,
delay,
},
);
self.in_flight.insert(node.clone());
let client = self.client.clone();
self.dials.spawn(async move {
let _ = client.connect_device(None, &ticket).await;
node
});
}
}
}
async fn authenticated_ticket_source(
client: &Client,
event: &NativePeerDataEvent,
scope: &str,
) -> Option<String> {
let peer = client.peer_session(&event.connection_id).await?;
match peer.node_id.as_deref() {
Some(source)
if peer.active_connection_id.as_deref() == Some(event.connection_id.as_str())
&& peer.settled_ready
&& !peer.logical_session_terminal
&& event.remote_node_id.as_deref() == Some(source)
&& event.transport_generation == peer.transport_generation
&& event.route_generation == peer.route_generation
&& peer.scopes.iter().any(|candidate| candidate == scope) =>
{
Some(source.to_string())
}
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::session_token::SessionTokenRegistry;
#[test]
fn mesh_wire_is_namespaced_and_bounded() {
let hello = Wire::Hello {
id: "share-1".into(),
ticket: "secret".into(),
};
assert!(
matches!(decode_wire(&encode_wire(&hello)), Some(Wire::Hello { id, ticket }) if id == "share-1" && ticket == "secret")
);
assert!(matches!(
decode_wire(&encode_wire(&Wire::Leave { id: "share-1".into() })),
Some(Wire::Leave { id }) if id == "share-1"
));
assert!(matches!(
decode_wire(&encode_wire(&Wire::Ack {
id: "share-1".into(),
kind: AckKind::Leave,
})),
Some(Wire::Ack { id, kind: AckKind::Leave }) if id == "share-1"
));
assert!(decode_wire(br#"{"type":"hello","id":"share-1"}"#).is_none());
let mut oversized = WIRE_PREFIX.to_vec();
oversized.extend(vec![b' '; 128 * 1024 + 1]);
assert!(decode_wire(&oversized).is_none());
}
#[test]
fn guest_scope_can_be_closed_and_reissued_without_a_second_owner() {
let registry = SessionTokenRegistry::new();
let scope = "ticket:share-1";
let issuer = iroh::SecretKey::generate().public().to_string();
registry.set_ticket_peer_limit(scope, 3).unwrap();
registry.require_scope_peer_admission(scope).unwrap();
registry
.update_scope_peer_admission(scope, 1, now_unix_ms() + 90_000, vec![issuer.clone()])
.unwrap();
assert!(registry.set_ticket_peer_limit(scope, 3).is_err());
registry.revoke_by_scope(scope);
registry.release_ticket_peer_admission(scope);
registry.set_ticket_peer_limit(scope, 3).unwrap();
registry.require_scope_peer_admission(scope).unwrap();
registry
.update_scope_peer_admission(scope, 1, now_unix_ms() + 90_000, vec![issuer])
.unwrap();
}
#[tokio::test]
async fn three_native_members_converge_on_one_ticket_mesh() {
use iroh::endpoint::{presets, RelayMode};
use iroh::Endpoint;
async fn client() -> Arc<Client> {
let client = Arc::new(Client::new_provider_neutral(
"ticket-mesh-test".into(),
Box::new(|| None),
));
let endpoint = Endpoint::builder(presets::Minimal)
.alpns(vec![crate::native_node::PlutoniumProtocol::ALPN.to_vec()])
.relay_mode(RelayMode::Disabled)
.bind_addr("127.0.0.1:0".parse::<std::net::SocketAddr>().unwrap())
.unwrap()
.bind()
.await
.unwrap();
client.adopt_endpoint(endpoint).await;
client
}
let host_client = client().await;
let guest_a_client = client().await;
let guest_b_client = client().await;
let host = TicketMesh::issue(
host_client,
"three-native-members",
TicketMeshOptions { max_peers: 3 },
)
.await
.unwrap();
let invite = host.invite();
let guest_a = TicketMesh::join(guest_a_client.clone(), &invite)
.await
.unwrap();
let guest_b = TicketMesh::join(guest_b_client.clone(), &invite)
.await
.unwrap();
let result = tokio::time::timeout(Duration::from_secs(25), async {
loop {
if host.watch_peers().borrow().len() == 2
&& guest_a.watch_peers().borrow().len() == 2
&& guest_b.watch_peers().borrow().len() == 2
{
break;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
})
.await;
let counts = (
host.watch_peers().borrow().len(),
guest_a.watch_peers().borrow().len(),
guest_b.watch_peers().borrow().len(),
);
assert!(result.is_ok(), "ticket mesh did not converge: {counts:?}");
let guest_b_node = guest_b_client.current_node_id().await.unwrap();
let original = guest_a_client.peer_session(&guest_b_node).await.unwrap();
guest_a_client
.disconnect_with_reason(
guest_b_node.parse().unwrap(),
crate::lifecycle_reason::REASON_NETWORK_CHANGE_RECONNECT,
)
.await
.unwrap();
let recovered = tokio::time::timeout(Duration::from_secs(15), async {
loop {
if let Some(peer) = guest_a_client.peer_session(&guest_b_node).await {
if peer.settled_ready
&& peer.active_transport_stable_id != original.active_transport_stable_id
{
break;
}
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
})
.await;
assert!(
recovered.is_ok(),
"ticket guest route did not recover after transport loss"
);
guest_b.close().await;
let left = tokio::time::timeout(Duration::from_secs(5), async {
loop {
if host.watch_peers().borrow().len() == 1
&& guest_a.watch_peers().borrow().len() == 1
{
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
})
.await;
assert!(
left.is_ok(),
"explicit guest leave did not prune its mesh route"
);
host.close().await;
let closed = tokio::time::timeout(Duration::from_secs(5), async {
loop {
if guest_a.watch_peers().borrow().is_empty()
&& guest_b.watch_peers().borrow().is_empty()
{
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
})
.await;
assert!(
closed.is_ok(),
"issuer close did not retire guest mesh projections"
);
guest_a.close().await;
}
}