Compare commits

...
Author SHA1 Message Date
mollusk 2c93c1c24f Add post-grace peer recovery coordinator 2026-06-20 15:37:33 -04:00
mollusk 5564af02f9 Spike targeted gossip rebootstrap 2026-06-20 04:28:53 -04:00
molluskandClaude Opus 4.8 ae29d1fea2 Merge windows-port-phase2: native Windows cpal/WASAPI audio port (b0fdd4e)
windows-build / windows-build (push) Has been cancelled
cargo-deny / cargo-deny (push) Has been cancelled
Brings the native Windows audio backend to main after a live host<->VM smoke
test: cpal/WASAPI capture+playback, device remap/resampling (W4/B5), cpal
RT-audit closed (B1-B5 + P3), Windows notification chimes, and Wine startup fix.

Verified on real Win11 (libvirt VM) this session: 2-way audio (host<->VM both
directions), audible join/leave/reconnect chimes, GUI renders, echo-cancel
correctly gated off. Linux unchanged (all changes cfg(windows); cargo test --lib
326/0, clippy clean). Windows build is GNU cross-compiled (b0fdd4e tester zip).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-20 02:55:10 -04:00
molluskandClaude Opus 4.8 7ff7766ede packaging: add test-pack split PKGBUILD (peerspeak + pixelpass)
One `makepkg -si` from packaging/test-pack/ builds and installs both
peerspeak and pixelpass from the public gitbutter repos over https, so a
tester can clone the repo and get a working voice+screenshare pair in one
command. pixelpass installs to /usr/bin so peerspeak's screen-share button
finds it.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-19 22:56:51 -04:00
10 changed files with 997 additions and 31 deletions
+112
View File
@@ -0,0 +1,112 @@
# Maintainer: mollusk <jitty+lc1iz0dc@protonmail.com>
#
# Test-pack split package: ONE `makepkg -si` builds + installs BOTH peerspeak
# (voice chat) and pixelpass (screen sharing) from the public gitbutter repos
# over https. pixelpass lands on /usr/bin so peerspeak's screen-share button
# finds it. Shared version string is derived from peerspeak's git.
#
# Clone this repo and build from here:
# git clone https://gitbutter.xyz/mollusk/peerspeak.git
# cd peerspeak/packaging/test-pack
# makepkg -si
pkgbase=peerspeak-git
pkgname=('peerspeak-git' 'pixelpass')
pkgver=0.1.0
pkgrel=1
arch=('x86_64')
url="https://gitbutter.xyz/mollusk/peerspeak"
license=('custom' 'MIT' 'Apache-2.0' 'OFL-1.1')
makedepends=('git' 'cargo' 'pkgconf')
options=('!lto' '!debug')
source=("peerspeak::git+https://gitbutter.xyz/mollusk/peerspeak.git"
"pixelpass::git+https://gitbutter.xyz/mollusk/pixelpass.git#branch=main")
sha256sums=('SKIP'
'SKIP')
pkgver() {
cd "$srcdir/peerspeak"
# Shared across both split packages. 0.1.0.r<commits>.g<short-sha>.
printf '%s.r%s.g%s' \
"$(awk -F'\"' '/^version =/{print $2; exit}' Cargo.toml)" \
"$(git rev-list --count HEAD)" \
"$(git rev-parse --short HEAD)"
}
prepare() {
# Vendor deps up front so build() can run --frozen (no surprise network).
export CARGO_HOME="$srcdir/cargo-home"
local host; host="$(rustc -vV | sed -n 's/host: //p')"
cd "$srcdir/peerspeak"; cargo fetch --locked --target "$host"
cd "$srcdir/pixelpass"; cargo fetch --locked --target "$host"
}
build() {
export CARGO_HOME="$srcdir/cargo-home"
export RUSTUP_TOOLCHAIN=stable
export CARGO_TARGET_DIR=target
cd "$srcdir/peerspeak"
cargo build --frozen --release --bin peerspeak
cd "$srcdir/pixelpass"
# --features gui so the .desktop launcher (pixelpass --gui) works.
cargo build --frozen --release --features gui
}
check() {
export CARGO_HOME="$srcdir/cargo-home"
export RUSTUP_TOOLCHAIN=stable
# peerspeak library unit tests only — its integration suites bind real
# iroh/QUIC endpoints and fail in a sandboxed/offline build environment.
cd "$srcdir/peerspeak"
cargo test --frozen --release --lib
}
package_peerspeak-git() {
pkgdesc="Decentralized peer-to-peer voice chat (Rust/iroh/PipeWire/Opus/iced)"
depends=('pipewire' 'opus')
optdepends=('pixelpass: screen sharing inside a room'
'mpv: screen-share viewer (vlc is used as a fallback)')
provides=('peerspeak')
conflicts=('peerspeak')
license=('custom')
cd "$srcdir/peerspeak"
install -Dm755 "target/release/peerspeak" "$pkgdir/usr/bin/peerspeak"
install -Dm644 "packaging/peerspeak.desktop" \
"$pkgdir/usr/share/applications/peerspeak.desktop"
# Hicolor icon theme (scalable SVG + the rendered raster sizes).
install -Dm644 "assets/icons/peerspeak.svg" \
"$pkgdir/usr/share/icons/hicolor/scalable/apps/peerspeak.svg"
local s
for s in 16 24 32 48 64 128 256 512; do
install -Dm644 "assets/icons/peerspeak-$s.png" \
"$pkgdir/usr/share/icons/hicolor/${s}x${s}/apps/peerspeak.png"
done
}
package_pixelpass() {
pkgdesc='P2P screen sharing over iroh — no port forwarding, no signup'
depends=('gstreamer' 'gst-plugins-base' 'gst-plugins-good' 'gst-plugins-bad'
'gst-libav' 'gst-plugin-va' 'libpulse' 'hicolor-icon-theme'
'libglvnd' 'libxkbcommon' 'wayland')
optdepends=('mpv: recommended stream viewer (the GUI launches mpv)'
'vlc: alternative stream viewer'
'gst-plugins-ugly: software x264 encoding for `pixelpass --no-hwencode`'
'gst-plugin-pipewire: screen capture on Wayland sessions'
'xorg-xwininfo: share a single window on X11 (`pixelpass --window`)')
license=('MIT' 'Apache-2.0' 'OFL-1.1')
cd "$srcdir/pixelpass"
install -Dm0755 "target/release/pixelpass" "$pkgdir/usr/bin/pixelpass"
install -Dm0644 assets/pixelpass.desktop \
"$pkgdir/usr/share/applications/pixelpass.desktop"
install -Dm0644 assets/pixelpass.svg \
"$pkgdir/usr/share/icons/hicolor/scalable/apps/pixelpass.svg"
install -Dm0644 README.md "$pkgdir/usr/share/doc/pixelpass/README.md"
install -Dm0644 LICENSE-MIT "$pkgdir/usr/share/licenses/pixelpass/LICENSE-MIT"
install -Dm0644 LICENSE-APACHE "$pkgdir/usr/share/licenses/pixelpass/LICENSE-APACHE"
install -Dm0644 assets/NotoSans-OFL.txt \
"$pkgdir/usr/share/licenses/pixelpass/NotoSans-OFL.txt"
}
+52
View File
@@ -0,0 +1,52 @@
# PeerSpeak + PixelPass — CachyOS/Arch test pack
A single **split PKGBUILD** that builds the latest code from the public gitbutter
repos and installs **both** programs at once:
- `peerspeak` — decentralized P2P voice chat
- `pixelpass` — P2P screen sharing (peerspeak launches it for the screen-share button)
## Build & install (one command)
```sh
git clone https://gitbutter.xyz/mollusk/peerspeak.git
cd peerspeak/packaging/test-pack
makepkg -si
```
`makepkg -si` auto-installs every dependency via pacman before building —
including the Rust toolchain itself (the `cargo` makedepend is provided by the
`rust` package), `git`, `pkgconf`, pipewire + opus for peerspeak, and the
gstreamer/VA-API stack for pixelpass. The only prerequisite is the `base-devel`
group (which provides `makepkg`). If you already use `rustup`, that satisfies the
`cargo` makedepend and the `rust` package won't be pulled in — no conflict.
When it finishes you'll have `peerspeak` and `pixelpass` on your PATH at
`/usr/bin`. To rebuild later with fresh upstream code, re-run `makepkg -si`; the
git sources re-pull `main` and the version bumps automatically.
> Skip the test step with `makepkg -si --nocheck` for a faster build.
## Running the cross-internet test
1. Launch `peerspeak` on both machines.
2. One person **creates** a room and shares the room code/ticket with the other.
3. The other **joins** with that code.
4. iroh does NAT hole-punching automatically; if a direct path can't be made it
falls back to a public n0 relay — **no port forwarding required**.
5. Allow the app through any local firewall if prompted (outbound UDP / QUIC;
nothing needs to be opened inbound for relay mode).
### What we're smoke-testing
- Two real humans, two networks, over the internet.
- Mic capture + remote playback both directions, no crackle/dropouts.
- Mute / deafen, push-to-talk.
- Text chat in-room.
- Avatars (presets + custom upload) show up on the other side.
- Screen share: click the screen-share control → it launches `pixelpass`; the
viewer opens in `mpv` on the receiving side.
- Notification chimes (join/leave/etc.).
- Leave / rejoin cleanly.
If anything misbehaves, grab the log path peerspeak prints on startup and the
exact repro steps.
+4
View File
@@ -888,6 +888,10 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
state.ever_connected.remove(&id);
notify::play(Sound::PeerLeave, state.config.custom_sound_peer_leave.as_deref());
}
// Core-only recovery phase: presentation for this state lands in
// the separate UI follow-up. In particular, do not play the
// terminal ReconnectFailed chime here.
UiEvent::PeerRecoveryStarted { .. } => {}
UiEvent::PeerConnectionFailed { id } => {
state.peers.remove(&id);
state.audio_levels.remove(&id);
+3
View File
@@ -85,6 +85,9 @@ pub enum UiEvent {
RoomLeft,
PeerJoined { id: EndpointId, state: PeerState },
PeerLeft { id: EndpointId },
/// The fixed reconnect grace expired and bounded background gossip recovery
/// has started. This is non-terminal and must not play the failure chime.
PeerRecoveryStarted { id: EndpointId },
PeerConnectionFailed { id: EndpointId },
PeerUpdated { id: EndpointId, state: PeerState },
/// Audio link to a peer is being (re)established — show a connecting state.
+144 -26
View File
@@ -1,5 +1,6 @@
pub mod messages;
pub mod jitter;
mod recovery;
use crate::audio::{AudioBackend, PlatformAudioBackend};
use crate::audio::eq::{Eq, EqSettings};
@@ -11,6 +12,7 @@ use crate::network::{
gossip::IrohGossipState,
};
use crate::core::messages::{CoreCommand, UiEvent};
use crate::core::recovery::RecoveryCoordinator;
use crate::config::{NetworkMode, RecordingMode};
use crate::presence::PresenceMode;
@@ -102,6 +104,39 @@ type GraceTimers = Arc<std::sync::Mutex<HashMap<EndpointId, tokio::task::JoinHan
/// Scrubbed whenever a peer is evicted or leaves so a later rejoin starts clean.
type SeenConnected = Arc<std::sync::Mutex<HashSet<EndpointId>>>;
type KnownPeers =
Arc<std::sync::Mutex<HashMap<String, HashMap<EndpointId, EndpointAddr>>>>;
#[derive(Clone)]
struct RecoveryContext {
coordinator: RecoveryCoordinator,
room_state: Arc<IrohGossipState>,
known_peers: KnownPeers,
ticket: String,
}
impl RecoveryContext {
fn retained_addr(&self, peer_id: &EndpointId) -> Option<EndpointAddr> {
self.known_peers
.lock()
.unwrap()
.get(&self.ticket)
.and_then(|peers| peers.get(peer_id))
.cloned()
}
fn cancel(&self, peer_id: EndpointId) {
self.coordinator.cancel(peer_id);
}
fn forget(&self, peer_id: EndpointId) {
if let Some(peers) = self.known_peers.lock().unwrap().get_mut(&self.ticket) {
peers.remove(&peer_id);
}
self.coordinator.cancel(peer_id);
}
}
/// Cancel and forget a peer's pending grace timer, if any. No-op if none is armed.
fn cancel_grace_timer(timers: &GraceTimers, peer_id: &EndpointId) {
if let Some(handle) = timers.lock().unwrap().remove(peer_id) {
@@ -116,12 +151,17 @@ fn cancel_grace_timer(timers: &GraceTimers, peer_id: &EndpointId) {
/// link repeatedly resetting the clock and dodging eviction forever. On firing it
/// also scrubs the peer from `seen_connected` so a later rejoin isn't treated as a
/// reconnect on its initial dial.
struct GraceExpiry<'a> {
transport: &'a Arc<IrohTransport>,
jitter: &'a Arc<Mutex<HashMap<EndpointId, JitterBuffer>>>,
ui_tx: &'a mpsc::Sender<UiEvent>,
recovery: Option<&'a RecoveryContext>,
}
fn arm_grace_timer(
timers: &GraceTimers,
seen_connected: &SeenConnected,
transport: &Arc<IrohTransport>,
jitter: &Arc<Mutex<HashMap<EndpointId, JitterBuffer>>>,
ui_tx: &mpsc::Sender<UiEvent>,
expiry: GraceExpiry<'_>,
grace: Duration,
peer_id: EndpointId,
) {
@@ -129,15 +169,29 @@ fn arm_grace_timer(
if timers_guard.contains_key(&peer_id) {
return;
}
let transport_evict = transport.clone();
let jitter_evict = jitter.clone();
let ui_evict = ui_tx.clone();
let transport_evict = expiry.transport.clone();
let jitter_evict = expiry.jitter.clone();
let ui_evict = expiry.ui_tx.clone();
let timers_evict = timers.clone();
let seen_evict = seen_connected.clone();
let recovery_evict = expiry.recovery.cloned();
let handle = tokio::spawn(async move {
tokio::time::sleep(grace).await;
crate::log_msg(&format!("Reconnect grace expired; evicting peer {:?}", peer_id));
crate::log_msg(&format!("Reconnect grace expired for peer {:?}", peer_id));
if let Some(recovery) = &recovery_evict
&& !recovery.coordinator.begin(peer_id)
{
return;
}
transport_evict.remove_audio_sender(peer_id);
if let Some(recovery) = &recovery_evict {
// Revoke roster authority before the first await in teardown. A
// verified Announce racing after this point is then a PeerJoined and
// cancels recovery instead of being erased after it was accepted.
recovery.room_state.mark_peer_disconnected(peer_id);
}
transport_evict.disconnect_peer(peer_id).await;
jitter_evict.lock().await.remove(&peer_id);
// Scrub our internal state *before* announcing the eviction, so anything
@@ -146,7 +200,38 @@ fn arm_grace_timer(
// reconnect.
timers_evict.lock().unwrap().remove(&peer_id);
seen_evict.lock().unwrap().remove(&peer_id);
let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await;
let Some(recovery) = recovery_evict else {
let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await;
return;
};
if !recovery.coordinator.is_active(&peer_id) {
return;
}
let Some(addr) = recovery.retained_addr(&peer_id) else {
crate::log_msg(&format!(
"Cannot recover peer {:?}: no retained authenticated address",
peer_id
));
recovery.cancel(peer_id);
let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await;
return;
};
match recovery.coordinator.activate(peer_id, addr) {
Ok(true) => {
let _ = ui_evict.send(UiEvent::PeerRecoveryStarted { id: peer_id }).await;
}
Ok(false) => {}
Err(()) => {
crate::log_msg(&format!(
"Cannot recover peer {:?}: recovery coordinator unavailable",
peer_id
));
let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await;
}
}
});
timers_guard.insert(peer_id, handle);
}
@@ -309,6 +394,7 @@ pub struct ConnEventHandler {
seen_connected: SeenConnected,
transport: Arc<IrohTransport>,
jitter: Arc<Mutex<HashMap<EndpointId, JitterBuffer>>>,
recovery: Option<RecoveryContext>,
grace: Duration,
}
@@ -326,6 +412,7 @@ impl ConnEventHandler {
seen_connected,
transport,
jitter,
recovery: None,
grace: RECONNECT_GRACE,
}
}
@@ -336,6 +423,11 @@ impl ConnEventHandler {
self
}
fn with_recovery(mut self, recovery: RecoveryContext) -> Self {
self.recovery = Some(recovery);
self
}
pub async fn handle(&self, event: ConnEvent) {
match event {
ConnEvent::Connecting(id) => {
@@ -349,9 +441,12 @@ impl ConnEventHandler {
arm_grace_timer(
&self.grace_timers,
&self.seen_connected,
&self.transport,
&self.jitter,
&self.ui_tx,
GraceExpiry {
transport: &self.transport,
jitter: &self.jitter,
ui_tx: &self.ui_tx,
recovery: self.recovery.as_ref(),
},
self.grace,
id,
);
@@ -359,6 +454,15 @@ impl ConnEventHandler {
let _ = self.ui_tx.send(UiEvent::PeerConnecting { id }).await;
}
ConnEvent::Connected(id) => {
// A transport event cannot readmit a grace-expired peer. Ignore a
// stale/racing link until authenticated gossip emits PeerJoined.
if self
.recovery
.as_ref()
.is_some_and(|recovery| recovery.coordinator.is_active(&id))
{
return;
}
// The audio link came back — the peer recovered within the grace
// window, so cancel its eviction.
cancel_grace_timer(&self.grace_timers, &id);
@@ -371,6 +475,9 @@ impl ConnEventHandler {
// until the grace timer or the slow gossip Leave.
cancel_grace_timer(&self.grace_timers, &id);
self.seen_connected.lock().unwrap().remove(&id);
if let Some(recovery) = &self.recovery {
recovery.forget(id);
}
self.transport.remove_audio_sender(id);
self.transport.disconnect_peer(id).await;
self.jitter.lock().await.remove(&id);
@@ -387,6 +494,7 @@ struct ActiveSession {
mixer_task: tokio::task::JoinHandle<()>,
event_task: tokio::task::JoinHandle<()>,
conn_event_task: tokio::task::JoinHandle<()>,
recovery_task: tokio::task::JoinHandle<()>,
grace_timers: GraceTimers,
transport: Arc<IrohTransport>,
/// Loaded PipeWire echo-cancel module (if enabled); unloads on drop.
@@ -421,6 +529,7 @@ impl ActiveSession {
for (_, handle) in self.grace_timers.lock().unwrap().drain() {
handle.abort();
}
self.recovery_task.abort();
crate::log_msg("Aborted tasks");
let audio_backend_clone = audio_backend.clone();
@@ -730,8 +839,7 @@ async fn run_core_loop(
// first room's peers — the old single-set version cleared them on any ticket
// change, so an A→B→A bounce stranded the rejoiner with an empty bootstrap.
// Inner map keyed by peer id so updates refresh the address.
let known_peers: Arc<std::sync::Mutex<HashMap<String, HashMap<EndpointId, EndpointAddr>>>> =
Arc::new(std::sync::Mutex::new(HashMap::new()));
let known_peers: KnownPeers = Arc::new(std::sync::Mutex::new(HashMap::new()));
let audio_backend = Arc::new(PlatformAudioBackend::new());
@@ -1474,6 +1582,15 @@ async fn run_core_loop(
// The ticket of the room this event loop serves, so peer add/remove
// updates the right per-ticket bucket in `known_peers` (A8 archive).
let ticket_events = ticket_str.clone();
let (recovery_coordinator, recovery_task) =
RecoveryCoordinator::spawn(room_state.clone());
let recovery_context = RecoveryContext {
coordinator: recovery_coordinator,
room_state: room_state.clone(),
known_peers: known_peers.clone(),
ticket: ticket_str.clone(),
};
let recovery_events = recovery_context.clone();
// Friends store + ui sender, so a connected peer who is a friend has
// their saved address auto-healed (W7) — populates `last_addr` so the
// presence scheduler can reach them later.
@@ -1486,6 +1603,7 @@ async fn run_core_loop(
// A (re)join means the peer is back — cancel any
// pending reconnect grace timer before re-adding it.
cancel_grace_timer(&grace_timers_events, &peer_id);
recovery_events.cancel(peer_id);
transport_events.admit_audio_sender(peer_id);
// Establish the audio connection as soon as the peer
// is known (the transport dedupes the full-mesh race).
@@ -1529,15 +1647,9 @@ async fn run_core_loop(
// Graceful leave — evict immediately.
cancel_grace_timer(&grace_timers_events, &peer_id);
seen_connected_events.lock().unwrap().remove(&peer_id);
// Graceful leave: drop them as a rejoin dial target
// for this room (a transient PeerConnectionLost
// deliberately does NOT, so we can still re-dial a
// peer who's still up).
if let Some(peers) =
known_peers_events.lock().unwrap().get_mut(&ticket_events)
{
peers.remove(&peer_id);
}
// A signed Leave cancels background recovery and
// drops the retained target. Transient loss keeps it.
recovery_events.forget(peer_id);
transport_events.remove_audio_sender(peer_id);
transport_events.disconnect_peer(peer_id).await;
jitter_events.lock().await.remove(&peer_id);
@@ -1551,6 +1663,7 @@ async fn run_core_loop(
// it. Idempotent: an ordinary mute/unmute update just
// re-records the same address.
cancel_grace_timer(&grace_timers_events, &peer_id);
recovery_events.cancel(peer_id);
transport_events.admit_audio_sender(peer_id);
transport_events.connect_peer(state.addr.clone()).await;
// Auto-heal a friend's saved address (W7) on the
@@ -1598,9 +1711,12 @@ async fn run_core_loop(
arm_grace_timer(
&grace_timers_events,
&seen_connected_events,
&transport_events,
&jitter_events,
&ui_tx_events,
GraceExpiry {
transport: &transport_events,
jitter: &jitter_events,
ui_tx: &ui_tx_events,
recovery: Some(&recovery_events),
},
RECONNECT_GRACE,
peer_id,
);
@@ -1624,7 +1740,8 @@ async fn run_core_loop(
seen_connected.clone(),
transport.clone(),
jitter.clone(),
);
)
.with_recovery(recovery_context);
let conn_event_task = tokio::spawn(async move {
while let Some(event) = conn_events.recv().await {
conn_handler.handle(event).await;
@@ -1638,6 +1755,7 @@ async fn run_core_loop(
mixer_task,
event_task,
conn_event_task,
recovery_task,
grace_timers,
transport: transport.clone(),
#[cfg(target_os = "linux")]
+264
View File
@@ -0,0 +1,264 @@
use crate::network::{RoomState, gossip::IrohGossipState};
use iroh::{EndpointAddr, EndpointId};
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use tokio::time::Instant;
const RECOVERY_COMMAND_CAPACITY: usize = 64;
const RECOVERY_DELAYS: [Duration; 7] = [
Duration::from_secs(1),
Duration::from_secs(2),
Duration::from_secs(4),
Duration::from_secs(8),
Duration::from_secs(15),
Duration::from_secs(30),
Duration::from_secs(60),
];
fn recovery_delay(attempt: usize) -> Duration {
RECOVERY_DELAYS[attempt.min(RECOVERY_DELAYS.len() - 1)]
}
enum RecoveryCommand {
Start {
peer_id: EndpointId,
addr: EndpointAddr,
},
Cancel(EndpointId),
}
struct RecoveryEntry {
addr: EndpointAddr,
attempt: usize,
next_attempt: Instant,
}
#[async_trait::async_trait]
trait RecoveryRoom: Send + Sync {
async fn rebootstrap_peers(&self, peers: Vec<EndpointAddr>) -> Result<(), String>;
}
#[async_trait::async_trait]
impl RecoveryRoom for IrohGossipState {
async fn rebootstrap_peers(&self, peers: Vec<EndpointAddr>) -> Result<(), String> {
RoomState::rebootstrap_peers(self, peers)
.await
.map_err(|error| error.to_string())
}
}
/// Cloneable command side of the single per-session recovery coordinator.
/// `active` is shared with transport/event handlers so cancellation is visible
/// immediately even while the coordinator is awaiting an in-flight gossip call.
#[derive(Clone)]
pub(super) struct RecoveryCoordinator {
tx: mpsc::Sender<RecoveryCommand>,
active: Arc<Mutex<HashSet<EndpointId>>>,
}
impl RecoveryCoordinator {
pub(super) fn spawn(room_state: Arc<IrohGossipState>) -> (Self, JoinHandle<()>) {
Self::spawn_inner(room_state)
}
fn spawn_inner(room_state: Arc<dyn RecoveryRoom>) -> (Self, JoinHandle<()>) {
let (tx, rx) = mpsc::channel(RECOVERY_COMMAND_CAPACITY);
let active = Arc::new(Mutex::new(HashSet::new()));
let handle = Self {
tx,
active: active.clone(),
};
let task = tokio::spawn(run_coordinator(room_state, active, rx));
(handle, task)
}
/// Reserve one recovery slot before grace-expiry teardown begins. Returns
/// false when the peer is already recovering, preventing duplicate work.
pub(super) fn begin(&self, peer_id: EndpointId) -> bool {
self.active.lock().unwrap().insert(peer_id)
}
/// Activate the reserved slot with its retained authenticated address.
/// Uses a bounded non-blocking send while holding the active-set lock so a
/// concurrent cancellation is ordered before or after this command.
pub(super) fn activate(&self, peer_id: EndpointId, addr: EndpointAddr) -> Result<bool, ()> {
let mut active = self.active.lock().unwrap();
if !active.contains(&peer_id) {
return Ok(false);
}
if self
.tx
.try_send(RecoveryCommand::Start { peer_id, addr })
.is_err()
{
active.remove(&peer_id);
return Err(());
}
Ok(true)
}
pub(super) fn cancel(&self, peer_id: EndpointId) {
self.active.lock().unwrap().remove(&peer_id);
// Cancellation is governed by the shared active set, so it remains
// immediate even if the bounded command queue is temporarily full.
let _ = self.tx.try_send(RecoveryCommand::Cancel(peer_id));
}
pub(super) fn is_active(&self, peer_id: &EndpointId) -> bool {
self.active.lock().unwrap().contains(peer_id)
}
}
async fn run_coordinator(
room_state: Arc<dyn RecoveryRoom>,
active: Arc<Mutex<HashSet<EndpointId>>>,
mut rx: mpsc::Receiver<RecoveryCommand>,
) {
let mut entries: HashMap<EndpointId, RecoveryEntry> = HashMap::new();
loop {
// The shared active set is the authoritative cancellation gate. Prune
// here as well as on Cancel commands so a saturated command queue cannot
// leave an inactive, past-due entry spinning the timer loop.
let active_snapshot = active.lock().unwrap().clone();
entries.retain(|peer_id, _| active_snapshot.contains(peer_id));
let next_deadline = entries.values().map(|entry| entry.next_attempt).min();
let command = match next_deadline {
Some(deadline) => {
tokio::select! {
command = rx.recv() => command,
_ = tokio::time::sleep_until(deadline) => {
let now = Instant::now();
let active_snapshot = active.lock().unwrap().clone();
let due: Vec<(EndpointId, EndpointAddr)> = entries
.iter()
.filter(|(id, entry)| {
entry.next_attempt <= now && active_snapshot.contains(*id)
})
.map(|(id, entry)| (*id, entry.addr.clone()))
.collect();
if !due.is_empty() {
let addrs = due.iter().map(|(_, addr)| addr.clone()).collect();
if let Err(error) = room_state.rebootstrap_peers(addrs).await {
crate::log_msg(&format!(
"Background peer recovery attempt failed: {error}"
));
}
let scheduled_at = Instant::now();
for (peer_id, _) in due {
if !active.lock().unwrap().contains(&peer_id) {
entries.remove(&peer_id);
continue;
}
if let Some(entry) = entries.get_mut(&peer_id) {
entry.next_attempt = scheduled_at + recovery_delay(entry.attempt);
entry.attempt = entry.attempt.saturating_add(1);
}
}
}
continue;
}
}
}
None => rx.recv().await,
};
match command {
Some(RecoveryCommand::Start { peer_id, addr }) => {
if active.lock().unwrap().contains(&peer_id) {
entries.entry(peer_id).or_insert(RecoveryEntry {
addr,
attempt: 0,
next_attempt: Instant::now(),
});
}
}
Some(RecoveryCommand::Cancel(peer_id)) => {
entries.remove(&peer_id);
}
None => break,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use iroh::SecretKey;
struct RecordingRoom {
attempts: mpsc::UnboundedSender<Vec<EndpointAddr>>,
}
#[async_trait::async_trait]
impl RecoveryRoom for RecordingRoom {
async fn rebootstrap_peers(&self, peers: Vec<EndpointAddr>) -> Result<(), String> {
self.attempts.send(peers).map_err(|error| error.to_string())
}
}
#[test]
fn retry_backoff_reaches_and_stays_at_sixty_seconds() {
let actual: Vec<u64> = (0..10)
.map(|attempt| recovery_delay(attempt).as_secs())
.collect();
assert_eq!(actual, vec![1, 2, 4, 8, 15, 30, 60, 60, 60, 60]);
}
#[test]
fn recovery_slots_are_deduplicated_and_cancel_immediately() {
let (tx, mut rx) = mpsc::channel(4);
let coordinator = RecoveryCoordinator {
tx,
active: Arc::new(Mutex::new(HashSet::new())),
};
let peer_id = SecretKey::generate().public();
assert!(coordinator.begin(peer_id));
assert!(
!coordinator.begin(peer_id),
"a peer gets only one recovery slot"
);
assert_eq!(
coordinator.activate(peer_id, EndpointAddr::from(peer_id)),
Ok(true)
);
assert!(matches!(
rx.try_recv(),
Ok(RecoveryCommand::Start { peer_id: id, .. }) if id == peer_id
));
coordinator.cancel(peer_id);
assert!(!coordinator.is_active(&peer_id));
assert!(matches!(
rx.try_recv(),
Ok(RecoveryCommand::Cancel(id)) if id == peer_id
));
}
#[tokio::test]
async fn coordinator_attempts_rebootstrap_immediately() {
let (attempts_tx, mut attempts_rx) = mpsc::unbounded_channel();
let (coordinator, task) = RecoveryCoordinator::spawn_inner(Arc::new(RecordingRoom {
attempts: attempts_tx,
}));
let peer_id = SecretKey::generate().public();
let addr = EndpointAddr::from(peer_id);
assert!(coordinator.begin(peer_id));
assert_eq!(coordinator.activate(peer_id, addr.clone()), Ok(true));
let attempted = tokio::time::timeout(Duration::from_secs(1), attempts_rx.recv())
.await
.expect("first recovery attempt should be immediate")
.expect("recording room remains subscribed");
assert_eq!(attempted, vec![addr]);
coordinator.cancel(peer_id);
task.abort();
}
}
+52 -2
View File
@@ -5,7 +5,7 @@ use iroh_gossip::proto::TopicId;
use tokio::sync::mpsc;
use tokio::sync::mpsc::Receiver;
use std::sync::{Arc, Mutex};
use std::collections::HashMap;
use std::collections::{HashMap, HashSet};
use async_trait::async_trait;
use tokio_stream::StreamExt;
use serde::{Serialize, Deserialize};
@@ -198,6 +198,10 @@ pub struct IrohGossipState {
secret_key: SecretKey,
self_state: Arc<Mutex<Option<PeerState>>>,
peers: Arc<Mutex<HashMap<EndpointId, PeerState>>>,
/// Previously verified peers whose live roster entry was removed by a
/// transient disconnect. Retained only so a later authenticated `Leave`
/// still reaches core and cancels background recovery.
disconnected_peers: Arc<Mutex<HashSet<EndpointId>>>,
event_tx: mpsc::Sender<RoomEvent>,
event_rx: Mutex<Option<mpsc::Receiver<RoomEvent>>>,
active_topic: Mutex<Option<tokio::task::JoinHandle<()>>>,
@@ -223,6 +227,7 @@ impl IrohGossipState {
secret_key,
self_state: Arc::new(Mutex::new(None)),
peers: Arc::new(Mutex::new(HashMap::new())),
disconnected_peers: Arc::new(Mutex::new(HashSet::new())),
event_tx,
event_rx: Mutex::new(Some(event_rx)),
active_topic: Mutex::new(None),
@@ -295,6 +300,7 @@ impl RoomState for IrohGossipState {
let event_tx = self.event_tx.clone();
let peers = self.peers.clone();
let disconnected_peers = self.disconnected_peers.clone();
let address_lookup = self.address_lookup.clone();
let self_state_clone = self.self_state.clone();
let gossip_sender_clone = gossip_sender.clone();
@@ -393,6 +399,7 @@ impl RoomState for IrohGossipState {
// peer-supplied: cap/validate once at ingest
// so invalid offers never render a Watch button.
state.sharing = state.sharing.and_then(crate::screenshare::sanitize_ticket);
disconnected_peers.lock().unwrap().remove(&payload.author);
let (is_new, state_changed) = {
let mut peer_map = peers.lock().unwrap();
let is_new = !peer_map.contains_key(&payload.author);
@@ -423,7 +430,11 @@ impl RoomState for IrohGossipState {
GossipMessage::Leave => {
crate::log_msg(&format!("Gossip peer leave request from author={:?}", payload.author));
let removed = peers.lock().unwrap().remove(&payload.author).is_some();
if removed {
let was_disconnected = disconnected_peers
.lock()
.unwrap()
.remove(&payload.author);
if removed || was_disconnected {
let _ = event_tx.send(RoomEvent::PeerLeft(payload.author)).await;
}
}
@@ -472,6 +483,7 @@ impl RoomState for IrohGossipState {
// cached presence entry; a rejoin re-announces as new.
let removed = peers.lock().unwrap().remove(&peer_id).is_some();
if removed {
disconnected_peers.lock().unwrap().insert(peer_id);
crate::log_msg(&format!("Peer connection lost (NeighborDown): {:?}", peer_id));
let _ = event_tx.send(RoomEvent::PeerConnectionLost(peer_id)).await;
}
@@ -516,6 +528,43 @@ impl RoomState for IrohGossipState {
Ok(())
}
async fn rebootstrap_peers(&self, peers: Vec<EndpointAddr>) -> Result<(), NetError> {
let self_id = self._endpoint.id();
let mut peer_ids = Vec::new();
for addr in peers {
if addr.id == self_id || peer_ids.contains(&addr.id) {
continue;
}
self.address_lookup.add_endpoint_info(addr.clone());
peer_ids.push(addr.id);
}
if peer_ids.is_empty() {
return Ok(());
}
// Clone the sender before awaiting: active_sender is a standard mutex and
// must never be held across an async gossip operation.
let sender = self
.active_sender
.lock()
.unwrap()
.clone()
.ok_or_else(|| NetError::Other("Not in a room".to_string()))?;
crate::log_msg(&format!("Rebootstrapping gossip peers: {:?}", peer_ids));
sender
.join_peers(peer_ids)
.await
.map_err(|e| NetError::Gossip(e.to_string()))
}
fn mark_peer_disconnected(&self, peer_id: EndpointId) {
if self.peers.lock().unwrap().remove(&peer_id).is_some() {
self.disconnected_peers.lock().unwrap().insert(peer_id);
}
}
async fn send_chat(&self, text: String) -> Result<(), NetError> {
let name = {
let guard = self.self_state.lock().unwrap();
@@ -571,6 +620,7 @@ impl RoomState for IrohGossipState {
}
self.peers.lock().unwrap().clear();
self.disconnected_peers.lock().unwrap().clear();
Ok(())
}
+11 -1
View File
@@ -194,6 +194,17 @@ pub trait RoomState: Send + Sync {
/// Updates our local state (e.g. when user mutes/unmutes) and broadcasts it.
async fn update_self_state(&self, self_state: PeerState) -> Result<(), NetError>;
/// Ask the active gossip topic to connect to retained peer addresses without
/// leaving or replacing the subscription. This is a recovery primitive only:
/// it does not add peers to the authenticated room roster. A peer becomes
/// active only after its normal signed `Announce` is received and verified.
async fn rebootstrap_peers(&self, peers: Vec<EndpointAddr>) -> Result<(), NetError>;
/// Remove a peer from the authenticated live roster before background
/// recovery. This only revokes membership; a fresh verified `Announce` is
/// required to add the peer again.
fn mark_peer_disconnected(&self, peer_id: EndpointId);
/// Broadcasts a room text-chat message authored by us (our display name is
/// taken from the current self-state).
async fn send_chat(&self, text: String) -> Result<(), NetError>;
@@ -342,4 +353,3 @@ mod tests {
assert_eq!(original, deserialized);
}
}
+354
View File
@@ -0,0 +1,354 @@
//! Phase-0 spike for post-grace gossip recovery.
//!
//! These tests prove that `GossipSender::join_peers` can restore an existing
//! topic subscription after the other peer drops and rejoins without its own
//! bootstrap target. The second case disables relays, clears the surviving
//! node's lookup, and moves the peer to a fresh endpoint address so only the
//! retained full address passed to `rebootstrap_peers` can drive recovery.
use std::sync::Arc;
use std::time::Duration;
use iroh::address_lookup::memory::MemoryLookup;
use iroh::endpoint::presets;
use iroh::protocol::Router;
use iroh::{Endpoint, EndpointAddr, EndpointId, RelayMode, SecretKey};
use iroh_gossip::net::Gossip;
use tokio::sync::mpsc;
use peerspeak::network::gossip::IrohGossipState;
use peerspeak::network::{PeerSpeakTicket, PeerState, RoomEvent, RoomState};
const EVENT_TIMEOUT: Duration = Duration::from_secs(10);
struct GossipNode {
endpoint: Endpoint,
lookup: MemoryLookup,
room: Arc<IrohGossipState>,
_router: Router,
}
async fn spawn_node(secret: SecretKey) -> GossipNode {
let lookup = MemoryLookup::new();
let endpoint = Endpoint::builder(presets::Minimal)
.secret_key(secret.clone())
.relay_mode(RelayMode::Disabled)
.address_lookup(lookup.clone())
.bind()
.await
.expect("bind gossip endpoint");
let gossip = Gossip::builder().spawn(endpoint.clone());
let router = Router::builder(endpoint.clone())
.accept(iroh_gossip::net::GOSSIP_ALPN, gossip.clone())
.spawn();
let room = Arc::new(IrohGossipState::new(
endpoint.clone(),
gossip,
lookup.clone(),
secret,
));
GossipNode {
endpoint,
lookup,
room,
_router: router,
}
}
fn state(name: &str, addr: EndpointAddr) -> PeerState {
PeerState {
name: name.to_string(),
is_muted: false,
addr,
sharing: None,
avatar: Default::default(),
}
}
async fn await_joined(rx: &mut mpsc::Receiver<RoomEvent>, peer_id: EndpointId) -> PeerState {
let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT;
loop {
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(RoomEvent::PeerJoined(id, peer_state))) if id == peer_id => return peer_state,
Ok(Some(_)) => continue,
Ok(None) => panic!("room event channel closed while waiting for PeerJoined"),
Err(_) => panic!("timed out waiting for PeerJoined({peer_id:?})"),
}
}
}
async fn await_joined_all(rx: &mut mpsc::Receiver<RoomEvent>, peer_ids: &[EndpointId]) {
let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT;
let mut remaining = peer_ids.to_vec();
while !remaining.is_empty() {
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(RoomEvent::PeerJoined(id, _))) => remaining.retain(|wanted| *wanted != id),
Ok(Some(_)) => {}
Ok(None) => panic!("room event channel closed while waiting for PeerJoined set"),
Err(_) => panic!("timed out waiting for PeerJoined set: {remaining:?}"),
}
}
}
async fn await_left(rx: &mut mpsc::Receiver<RoomEvent>, peer_id: EndpointId) {
let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT;
loop {
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(RoomEvent::PeerLeft(id))) if id == peer_id => return,
Ok(Some(_)) => continue,
Ok(None) => panic!("room event channel closed while waiting for PeerLeft"),
Err(_) => panic!("timed out waiting for PeerLeft({peer_id:?})"),
}
}
}
async fn await_absent(room: &IrohGossipState, peer_id: EndpointId) {
let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT;
loop {
if !room.active_peers().iter().any(|(id, _)| *id == peer_id) {
return;
}
assert!(
tokio::time::Instant::now() < deadline,
"timed out waiting for peer to leave the roster"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
fn ticket(host_addr: EndpointAddr) -> String {
PeerSpeakTicket {
host_addr,
topic_id: rand::random(),
name: "rebootstrap-spike".to_string(),
}
.to_string()
}
async fn establish_room(
a: &GossipNode,
b: &GossipNode,
ticket: &str,
events_a: &mut mpsc::Receiver<RoomEvent>,
) {
// B is the ticket host. Its own bootstrap set is empty; A is the only side
// that initially dials, which is also how the recovery setup is controlled.
b.room
.join(ticket, state("Bob", b.endpoint.addr()), vec![])
.await
.expect("host joins topic");
a.room
.join(ticket, state("Alice", a.endpoint.addr()), vec![])
.await
.expect("client joins topic");
let joined = await_joined(events_a, b.endpoint.id()).await;
assert_eq!(joined.name, "Bob");
}
#[tokio::test]
async fn rebootstrap_restores_roster_on_existing_subscription() {
let a = spawn_node(SecretKey::generate()).await;
let b = spawn_node(SecretKey::generate()).await;
let ticket = ticket(b.endpoint.addr());
let mut events_a = a.room.subscribe_events().await.expect("subscribe A events");
establish_room(&a, &b, &ticket, &mut events_a).await;
// Drop only B's topic subscription. A stays subscribed. B then rejoins as
// the ticket host, so compute_bootstrap removes self and B has nobody to dial.
b.room.leave().await.expect("B leaves topic");
await_absent(&a.room, b.endpoint.id()).await;
b.room
.join(&ticket, state("Bob recovered", b.endpoint.addr()), vec![])
.await
.expect("B rejoins without bootstrap peers");
tokio::time::sleep(Duration::from_millis(250)).await;
assert!(
!a.room
.active_peers()
.iter()
.any(|(id, _)| *id == b.endpoint.id()),
"B must not recover before A explicitly re-bootstraps it"
);
a.room
.rebootstrap_peers(vec![b.endpoint.addr()])
.await
.expect("targeted gossip re-bootstrap");
let recovered = await_joined(&mut events_a, b.endpoint.id()).await;
assert_eq!(recovered.name, "Bob recovered");
}
#[tokio::test]
async fn rebootstrap_uses_retained_full_address_with_empty_lookup() {
let a = spawn_node(SecretKey::generate()).await;
let b_secret = SecretKey::generate();
let b = spawn_node(b_secret.clone()).await;
let b_id = b.endpoint.id();
let old_b_addr = b.endpoint.addr();
let ticket = ticket(old_b_addr.clone());
let mut events_a = a.room.subscribe_events().await.expect("subscribe A events");
establish_room(&a, &b, &ticket, &mut events_a).await;
b.room.leave().await.expect("old B leaves topic");
await_absent(&a.room, b_id).await;
// Move the same authenticated identity to a newly-bound direct-only endpoint.
// The old cached path is now dead; the new full address is the only valid one.
b.endpoint.close().await;
drop(b);
tokio::time::sleep(Duration::from_millis(250)).await;
let b_rebound = spawn_node(b_secret).await;
let new_b_addr = b_rebound.endpoint.addr();
assert_eq!(new_b_addr.id, b_id, "identity must survive the rebind");
assert_ne!(
new_b_addr, old_b_addr,
"rebound peer must have a fresh address"
);
b_rebound
.room
.join(&ticket, state("Bob rebound", new_b_addr.clone()), vec![])
.await
.expect("rebound host joins without bootstrap peers");
// Remove the stale lookup entry. `rebootstrap_peers` must seed the retained
// new full address before asking gossip to join the peer by id.
a.lookup.remove_endpoint_info(b_id);
assert!(
a.lookup.get_endpoint_info(b_id).is_none(),
"A lookup starts empty for B"
);
tokio::time::sleep(Duration::from_millis(250)).await;
assert!(
!a.room.active_peers().iter().any(|(id, _)| *id == b_id),
"rebound B must not be rediscovered without the retained address"
);
a.room
.rebootstrap_peers(vec![new_b_addr])
.await
.expect("retained-address gossip re-bootstrap");
assert!(
a.lookup.get_endpoint_info(b_id).is_some(),
"re-bootstrap must restore B's address to the lookup"
);
let recovered = await_joined(&mut events_a, b_id).await;
assert_eq!(recovered.name, "Bob rebound");
}
#[tokio::test]
async fn demoted_peer_requires_a_fresh_signed_announce_to_rejoin() {
let a = spawn_node(SecretKey::generate()).await;
let b = spawn_node(SecretKey::generate()).await;
let ticket = ticket(b.endpoint.addr());
let mut events_a = a.room.subscribe_events().await.expect("subscribe A events");
establish_room(&a, &b, &ticket, &mut events_a).await;
a.room.mark_peer_disconnected(b.endpoint.id());
assert!(
!a.room
.active_peers()
.iter()
.any(|(id, _)| *id == b.endpoint.id()),
"demotion must revoke live roster membership"
);
b.room
.update_self_state(state("Bob authenticated again", b.endpoint.addr()))
.await
.expect("broadcast fresh signed announce");
let recovered = await_joined(&mut events_a, b.endpoint.id()).await;
assert_eq!(recovered.name, "Bob authenticated again");
}
#[tokio::test]
async fn signed_leave_after_demotion_still_emits_peer_left() {
let a = spawn_node(SecretKey::generate()).await;
let b = spawn_node(SecretKey::generate()).await;
let ticket = ticket(b.endpoint.addr());
let mut events_a = a.room.subscribe_events().await.expect("subscribe A events");
establish_room(&a, &b, &ticket, &mut events_a).await;
a.room.mark_peer_disconnected(b.endpoint.id());
b.room
.leave()
.await
.expect("broadcast signed Leave after demotion");
await_left(&mut events_a, b.endpoint.id()).await;
}
#[tokio::test]
async fn targeted_rebootstrap_preserves_healthy_peer_in_three_peer_room() {
let a = spawn_node(SecretKey::generate()).await;
let b = spawn_node(SecretKey::generate()).await;
let c_secret = SecretKey::generate();
let c = spawn_node(c_secret.clone()).await;
let c_id = c.endpoint.id();
let ticket = ticket(c.endpoint.addr());
let mut events_a = a.room.subscribe_events().await.expect("subscribe A events");
let mut events_b = b.room.subscribe_events().await.expect("subscribe B events");
c.room
.join(&ticket, state("Carol", c.endpoint.addr()), vec![])
.await
.expect("C hosts topic");
b.room
.join(&ticket, state("Bob", b.endpoint.addr()), vec![])
.await
.expect("B joins C");
await_joined(&mut events_b, c_id).await;
a.room
.join(&ticket, state("Alice", a.endpoint.addr()), vec![])
.await
.expect("A joins C");
await_joined_all(&mut events_a, &[b.endpoint.id(), c_id]).await;
await_joined(&mut events_b, a.endpoint.id()).await;
// Remove only C. A and B keep their existing topic subscriptions and remain
// mutually present while C is rebound to a fresh address.
c.endpoint.close().await;
drop(c);
a.room.mark_peer_disconnected(c_id);
b.room.mark_peer_disconnected(c_id);
let c_rebound = spawn_node(c_secret).await;
let rebound_addr = c_rebound.endpoint.addr();
c_rebound
.room
.join(
&ticket,
state("Carol recovered", rebound_addr.clone()),
vec![],
)
.await
.expect("rebound C rejoins as host without bootstrap");
a.room
.rebootstrap_peers(vec![rebound_addr])
.await
.expect("A targets only C for recovery");
let recovered = await_joined(&mut events_a, c_id).await;
assert_eq!(recovered.name, "Carol recovered");
await_joined(&mut events_b, c_id).await;
assert!(
a.room
.active_peers()
.iter()
.any(|(id, _)| *id == b.endpoint.id()),
"healthy B must remain present at A throughout C recovery"
);
assert!(
b.room
.active_peers()
.iter()
.any(|(id, _)| *id == a.endpoint.id()),
"healthy A must remain present at B throughout C recovery"
);
}
+1 -2
View File
@@ -25,8 +25,7 @@ use peerspeak::codec::opus_impl::OpusEncoder;
use peerspeak::core::jitter::{FRAME_SAMPLES, JitterBuffer};
use peerspeak::network::{ConnEvent, NetworkTransport};
use peerspeak::network::iroh_impl::{AudioRouter, IrohTransport};
const AUDIO_ALPN: &[u8] = b"peerspeak-audio";
use peerspeak::protocol::AUDIO_ALPN;
struct Node {
endpoint: Endpoint,