Fix core Tokio runtime panic and add network logging

- Fix tokio runtime panic by spawning a dedicated Tokio runtime thread in CoreController.
- Add central log_msg utility in src/lib.rs for debugging.
- Add instrumentation/logs to join, leave, and gossip events in src/network/gossip.rs.
- Add test_net.rs bin for testing gossip loopback sync.
- Use std::sync::Mutex in IrohGossipState to resolve Tokio block-in-async panics.
This commit is contained in:
2026-05-27 05:48:04 -04:00
parent 1220d94e91
commit 5ddb792f0f
7 changed files with 279 additions and 72 deletions
+12
View File
@@ -3,6 +3,18 @@ name = "peerspeak"
version = "0.1.0"
edition = "2024"
[lib]
name = "peerspeak"
path = "src/lib.rs"
[[bin]]
name = "peerspeak"
path = "src/main.rs"
[[bin]]
name = "test_net"
path = "src/bin/test_net.rs"
[dependencies]
anyhow = "1.0.102"
async-trait = "0.1.89"
+1
View File
@@ -21,6 +21,7 @@ Before answering highly complex questions, writing macros, or optimizing code, y
- **Rule 1 (Absolute Ground Truth):** Never guess or hallucinate syntax rules, compiler behavior, or API surfaces. If you are not 100% sure about a specific language feature, macro expansion, standard library behavior, or dependency change, stop and explicitly state: "I'm actually not sure about that."
- **Rule 2 (No "C in Rust"):** Do not write C-style logic wrapped in Rust syntax. Prioritize idiomatic Rust patterns (e.g., using algebraic data types, proper trait bounds, combinators like `.map()` or `.and_then()`, and precise error handling with `Result` and `Option`).
- **Rule 3 (Safe by Default):** Always default to safe, idiomatic Rust code. Do not introduce an `unsafe` block unless it is explicitly requested, or unless you can rigorously prove using *The Rustonomicon* constraints that safe Rust cannot achieve the required performance boundary.
- **Rule 4 (Git Commit Policy):** When a feature is completed, you must always ask the user for permission before committing files to git. Never commit files automatically.
### 4. Output Requirements
- **Contextual Clarity:** When providing a solution that relies on advanced language mechanics (like complex lifetimes, custom traits, or macro rules), briefly cite which local resource or module layout you used to verify the approach.
+103
View File
@@ -0,0 +1,103 @@
use peerspeak::network::{
gossip::IrohGossipState,
RoomState, PeerState, RoomEvent,
};
use iroh::{Endpoint, endpoint::presets};
use iroh_gossip::net::Gossip;
use std::sync::Arc;
use tokio::time::{self, Duration};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("Starting network loopback test...");
// 1. Node A (Host) Setup
let lookup_a = iroh::address_lookup::memory::MemoryLookup::new();
let secret_a = iroh::SecretKey::generate();
let endpoint_a = Endpoint::builder(presets::N0)
.secret_key(secret_a)
.address_lookup(lookup_a.clone())
.bind()
.await?;
endpoint_a.online().await;
println!("Node A online. ID: {}", endpoint_a.id());
let gossip_a = Gossip::builder().spawn(endpoint_a.clone());
let _router_a = iroh::protocol::Router::builder(endpoint_a.clone())
.accept(iroh_gossip::net::GOSSIP_ALPN, gossip_a.clone())
.spawn();
let room_a = IrohGossipState::new(endpoint_a.clone(), gossip_a.clone(), lookup_a.clone());
// 2. Node B (Client) Setup
let lookup_b = iroh::address_lookup::memory::MemoryLookup::new();
let secret_b = iroh::SecretKey::generate();
let endpoint_b = Endpoint::builder(presets::N0)
.secret_key(secret_b)
.address_lookup(lookup_b.clone())
.bind()
.await?;
endpoint_b.online().await;
println!("Node B online. ID: {}", endpoint_b.id());
let gossip_b = Gossip::builder().spawn(endpoint_b.clone());
let _router_b = iroh::protocol::Router::builder(endpoint_b.clone())
.accept(iroh_gossip::net::GOSSIP_ALPN, gossip_b.clone())
.spawn();
let room_b = IrohGossipState::new(endpoint_b.clone(), gossip_b.clone(), lookup_b.clone());
// 3. Create room on Node A
let topic_id = rand::random();
let ticket = peerspeak::network::PeerSpeakTicket {
host_addr: endpoint_a.addr(),
topic_id,
};
let ticket_str = ticket.to_string();
println!("Ticket generated: {}", ticket_str);
let state_a = PeerState {
name: "Alice".to_string(),
is_muted: false,
addr: endpoint_a.addr(),
};
room_a.join(&ticket_str, state_a).await?;
println!("Node A joined topic.");
// Subscribe to events on Node A
let mut rx_a = room_a.subscribe_events().await?;
tokio::spawn(async move {
while let Some(event) = rx_a.recv().await {
println!("Node A Event: {:?}", event);
}
});
// 4. Join room on Node B
let state_b = PeerState {
name: "Bob".to_string(),
is_muted: false,
addr: endpoint_b.addr(),
};
room_b.join(&ticket_str, state_b).await?;
println!("Node B joined topic.");
// Subscribe to events on Node B
let mut rx_b = room_b.subscribe_events().await?;
tokio::spawn(async move {
while let Some(event) = rx_b.recv().await {
println!("Node B Event: {:?}", event);
}
});
// Wait and check connection
println!("Waiting 10 seconds for Gossip sync...");
time::sleep(Duration::from_secs(10)).await;
println!("Alice's peers: {:?}", room_a.active_peers());
println!("Bob's peers: {:?}", room_b.active_peers());
println!("Test finished.");
Ok(())
}
+22 -8
View File
@@ -25,11 +25,15 @@ impl CoreController {
pub fn new(ui_tx: mpsc::Sender<UiEvent>) -> Self {
let (cmd_tx, cmd_rx) = mpsc::channel(100);
tokio::spawn(async move {
std::thread::spawn(move || {
let rt = tokio::runtime::Runtime::new().expect("Failed to create Tokio runtime");
rt.block_on(async move {
crate::log_msg("Starting core network loop in dedicated Tokio runtime");
if let Err(e) = run_core_loop(cmd_rx, ui_tx).await {
eprintln!("App core loop failed: {:?}", e);
crate::log_msg(&format!("App core loop failed: {:?}", e));
}
});
});
Self { cmd_tx }
}
@@ -103,6 +107,7 @@ async fn run_core_loop(
// Clean up any existing session
if let Some(session) = active_session.take() {
crate::log_msg("Shutting down existing active session");
session.shutdown(audio_backend.clone()).await;
}
@@ -110,10 +115,13 @@ async fn run_core_loop(
let ticket_str = if ticket.trim().is_empty() || ticket == "create" {
let topic_id: [u8; 32] = rand::random();
let host_addr = endpoint.addr();
crate::log_msg(&format!("Creating room. host_addr={:?}, topic_id={:?}", host_addr, topic_id));
let ticket = PeerSpeakTicket { host_addr, topic_id };
ticket.to_string()
} else {
ticket.trim().to_string()
let ticket_str = ticket.trim().to_string();
crate::log_msg(&format!("Joining room with existing ticket={}", ticket_str));
ticket_str
};
// Initialize Gossip and Transport
@@ -139,11 +147,14 @@ async fn run_core_loop(
addr: endpoint.addr(),
};
crate::log_msg(&format!("Attempting room_state.join with self_state={:?}", self_state));
if let Err(e) = room_state.join(&ticket_str, self_state.clone()).await {
crate::log_msg(&format!("Error room_state.join failed: {:?}", e));
let _ = ui_tx.send(UiEvent::Error(format!("Failed to join room: {}", e))).await;
let _ = router.shutdown().await;
continue;
}
crate::log_msg("Joined room successfully via room_state");
// Setup raw audio channels
let (capture_tx, capture_rx) = std::sync::mpsc::channel();
@@ -177,7 +188,7 @@ async fn run_core_loop(
let mut encoder = match OpusEncoder::new(48000, Channels::Mono, Application::Voip) {
Ok(enc) => enc,
Err(e) => {
eprintln!("Capture thread error: {:?}", e);
crate::log_msg(&format!("Capture thread error: {:?}", e));
return;
}
};
@@ -193,7 +204,9 @@ async fn run_core_loop(
let transport = transport_clone.clone();
let bytes = bytes.clone();
tokio_handle.spawn(async move {
let _ = transport.send_datagram(peer_id, bytes).await;
if let Err(e) = transport.send_datagram(peer_id, bytes).await {
crate::log_msg(&format!("Failed to send datagram to peer {:?}: {:?}", peer_id, e));
}
});
}
}
@@ -208,20 +221,21 @@ async fn run_core_loop(
let mut datagram_rx = match transport_recv.receive_datagrams().await {
Ok(rx) => rx,
Err(e) => {
eprintln!("Receiver task error: {:?}", e);
crate::log_msg(&format!("Receiver task error: {:?}", e));
return;
}
};
let mut decoders: HashMap<EndpointId, OpusDecoder> = HashMap::new();
while let Some((from_peer, bytes)) = datagram_rx.recv().await {
crate::log_msg(&format!("Received datagram from peer={:?}, len={}", from_peer, bytes.len()));
let decoder = match decoders.entry(from_peer) {
std::collections::hash_map::Entry::Occupied(entry) => entry.into_mut(),
std::collections::hash_map::Entry::Vacant(entry) => {
match OpusDecoder::new(48000, Channels::Mono) {
Ok(dec) => entry.insert(dec),
Err(e) => {
eprintln!("Failed to initialize decoder for {:?}: {:?}", from_peer, e);
crate::log_msg(&format!("Failed to initialize decoder for {:?}: {:?}", from_peer, e));
continue;
}
}
@@ -235,7 +249,7 @@ async fn run_core_loop(
queue.extend(pcm);
}
Err(e) => {
eprintln!("Failed to decode packet from {:?}: {:?}", from_peer, e);
crate::log_msg(&format!("Failed to decode packet from {:?}: {:?}", from_peer, e));
}
}
}
+21
View File
@@ -0,0 +1,21 @@
pub mod audio;
pub mod codec;
pub mod network;
pub mod core;
pub mod app;
pub fn log_msg(msg: &str) {
if let Ok(mut file) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open("/home/mollusk/peerspeak.log")
{
use std::io::Write;
if let Ok(time) = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) {
let _ = writeln!(file, "[{}.{:03}] {}", time.as_secs(), time.subsec_millis(), msg);
} else {
let _ = writeln!(file, "{}", msg);
}
}
}
+2 -9
View File
@@ -1,12 +1,5 @@
pub mod audio;
pub mod codec;
pub mod network;
pub mod core;
pub mod app;
#[tokio::main]
async fn main() {
if let Err(e) = app::run_gui() {
fn main() {
if let Err(e) = peerspeak::app::run_gui() {
eprintln!("Error running GUI: {:?}", e);
}
}
+98 -35
View File
@@ -2,9 +2,9 @@ use crate::network::{RoomState, NetError, PeerState, RoomEvent, PeerSpeakTicket}
use iroh::{Endpoint, EndpointId};
use iroh_gossip::net::Gossip;
use iroh_gossip::proto::TopicId;
use tokio::sync::{mpsc, Mutex};
use tokio::sync::mpsc;
use tokio::sync::mpsc::Receiver;
use std::sync::Arc;
use std::sync::{Arc, Mutex};
use std::collections::HashMap;
use async_trait::async_trait;
use tokio_stream::StreamExt;
@@ -60,9 +60,12 @@ impl IrohGossipState {
#[async_trait]
impl RoomState for IrohGossipState {
async fn join(&self, ticket_str: &str, self_state: PeerState) -> Result<(), NetError> {
crate::log_msg(&format!("RoomState::join: self_id={:?}, self_name={:?}, ticket={}", self_state.addr.id, self_state.name, ticket_str));
let ticket = ticket_str.parse::<PeerSpeakTicket>()?;
let topic_id = TopicId::from_bytes(ticket.topic_id);
crate::log_msg(&format!("Parsed ticket. host_id={:?}, host_addrs={:?}, topic={:?}", ticket.host_addr.id, ticket.host_addr.addrs, topic_id));
// Stop any currently running topic
let _ = self.leave().await;
@@ -72,99 +75,154 @@ impl RoomState for IrohGossipState {
// Join the gossip topic. If we are the host, bootstrap list will be empty
// or contain ourselves (which is fine), but let's bootstrap to the ticket host.
let bootstrap_peers = if ticket.host_addr.id == self_state.addr.id {
crate::log_msg("We are the host. Bootstrap peers list is empty.");
vec![]
} else {
crate::log_msg(&format!("We are a client. Bootstrapping to host ID={:?}", ticket.host_addr.id));
vec![ticket.host_addr.id]
};
let gossip_topic = self.gossip.subscribe(topic_id, bootstrap_peers).await
.map_err(|e| NetError::Gossip(format!("Failed to join gossip topic: {}", e)))?;
.map_err(|e| {
let err = format!("Failed to join gossip topic: {}", e);
crate::log_msg(&err);
NetError::Gossip(err)
})?;
let (gossip_sender, mut gossip_receiver) = gossip_topic.split();
*self.self_state.lock().await = Some(self_state.clone());
*self.active_topic_id.lock().await = Some(topic_id);
*self.active_sender.lock().await = Some(gossip_sender.clone());
*self.self_state.lock().unwrap() = Some(self_state.clone());
*self.active_topic_id.lock().unwrap() = Some(topic_id);
*self.active_sender.lock().unwrap() = Some(gossip_sender.clone());
let event_tx = self.event_tx.clone();
let peers = self.peers.clone();
let address_lookup = self.address_lookup.clone();
let self_state_clone = self.self_state.clone();
let gossip_sender_clone = gossip_sender.clone();
let self_id = self_state.addr.id;
let handle = tokio::spawn(async move {
crate::log_msg(&format!("Spawned gossip topic loop for self_id={:?}", self_id));
// Broadcast initial state
let payload = GossipPayload {
author: self_state_clone.lock().await.as_ref().unwrap().addr.id,
msg: GossipMessage::Announce(self_state_clone.lock().await.clone().unwrap()),
let initial_payload = {
let guard = self_state_clone.lock().unwrap();
guard.as_ref().map(|s| GossipPayload {
author: s.addr.id,
msg: GossipMessage::Announce(s.clone()),
})
};
if let Some(payload) = initial_payload {
if let Ok(bytes) = serde_json::to_vec(&payload) {
crate::log_msg(&format!("Broadcasting initial state from self_id={:?}", self_id));
let _ = gossip_sender_clone.broadcast(bytes.into()).await;
}
}
// Stream topic messages
while let Some(res) = gossip_receiver.next().await {
match res {
Ok(iroh_gossip::api::Event::Received(msg)) => {
if let Ok(payload) = serde_json::from_slice::<GossipPayload>(&msg.content) {
crate::log_msg(&format!("Gossip received Event::Received from delivery={:?}", msg.delivered_from));
match serde_json::from_slice::<GossipPayload>(&msg.content) {
Ok(payload) => {
let our_id = {
self_state_clone.lock().unwrap()
.as_ref()
.map(|s| s.addr.id)
};
if Some(payload.author) == our_id {
crate::log_msg("Gossip Event::Received from ourselves; ignoring");
continue;
}
crate::log_msg(&format!("Gossip Event::Received from author={:?}, payload={:?}", payload.author, payload.msg));
match payload.msg {
GossipMessage::Announce(state) => {
if payload.author == self_state_clone.lock().await.as_ref().unwrap().addr.id {
continue; // Ignore our own announcements
}
let mut peer_map = peers.lock().await;
let (is_new, state_changed) = {
let mut peer_map = peers.lock().unwrap();
let is_new = !peer_map.contains_key(&payload.author);
let state_changed = peer_map.get(&payload.author) != Some(&state);
if is_new || state_changed {
peer_map.insert(payload.author, state.clone());
}
(is_new, state_changed)
};
if is_new {
crate::log_msg(&format!("Gossip new peer joined: {:?}, state: {:?}", payload.author, state));
address_lookup.add_endpoint_info(state.addr.clone());
peer_map.insert(payload.author, state.clone());
let _ = event_tx.send(RoomEvent::PeerJoined(payload.author, state)).await;
} else if state_changed {
peer_map.insert(payload.author, state.clone());
crate::log_msg(&format!("Gossip peer state updated: {:?}, state: {:?}", payload.author, state));
let _ = event_tx.send(RoomEvent::PeerUpdated(payload.author, state)).await;
}
}
GossipMessage::Leave => {
let mut peer_map = peers.lock().await;
if peer_map.remove(&payload.author).is_some() {
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 _ = event_tx.send(RoomEvent::PeerLeft(payload.author)).await;
}
}
}
}
Err(e) => {
crate::log_msg(&format!("Gossip failed to deserialize payload: {:?}", e));
}
Ok(iroh_gossip::api::Event::NeighborUp(_peer_id)) => {
}
}
Ok(iroh_gossip::api::Event::NeighborUp(peer_id)) => {
crate::log_msg(&format!("Gossip event: NeighborUp={:?}", peer_id));
// Resend state on new neighbor connection to guarantee synchronization
if let Some(state) = self_state_clone.lock().await.as_ref() {
let payload = GossipPayload {
let payload_opt = {
let guard = self_state_clone.lock().unwrap();
guard.as_ref().map(|state| GossipPayload {
author: state.addr.id,
msg: GossipMessage::Announce(state.clone()),
})
};
if let Some(payload) = payload_opt {
if let Ok(bytes) = serde_json::to_vec(&payload) {
crate::log_msg(&format!("Broadcasting state to new neighbor={:?}", peer_id));
let _ = gossip_sender_clone.broadcast(bytes.into()).await;
}
}
}
_ => {}
Ok(iroh_gossip::api::Event::NeighborDown(peer_id)) => {
crate::log_msg(&format!("Gossip event: NeighborDown={:?}", peer_id));
}
Ok(other) => {
crate::log_msg(&format!("Gossip other event: {:?}", other));
}
Err(e) => {
crate::log_msg(&format!("Gossip error event: {:?}", e));
}
}
}
crate::log_msg("Gossip topic loop terminated");
});
*self.active_topic.lock().await = Some(handle);
*self.active_topic.lock().unwrap() = Some(handle);
Ok(())
}
async fn update_self_state(&self, self_state: PeerState) -> Result<(), NetError> {
let mut self_guard = self.self_state.lock().await;
*self_guard = Some(self_state.clone());
crate::log_msg(&format!("RoomState::update_self_state: state={:?}", self_state));
*self.self_state.lock().unwrap() = Some(self_state.clone());
if let Some(sender) = self.active_sender.lock().await.as_ref() {
let sender_opt = self.active_sender.lock().unwrap().clone();
if let Some(sender) = sender_opt {
let payload = GossipPayload {
author: self_state.addr.id,
msg: GossipMessage::Announce(self_state),
};
if let Ok(bytes) = serde_json::to_vec(&payload) {
crate::log_msg("Broadcasting updated self state to gossip");
sender.broadcast(bytes.into()).await
.map_err(|e| NetError::Gossip(e.to_string()))?;
}
@@ -173,38 +231,43 @@ impl RoomState for IrohGossipState {
}
async fn leave(&self) -> Result<(), NetError> {
let mut handle_guard = self.active_topic.lock().await;
crate::log_msg("RoomState::leave called");
{
let mut handle_guard = self.active_topic.lock().unwrap();
if let Some(handle) = handle_guard.take() {
crate::log_msg("Aborting active gossip topic background task");
handle.abort();
}
}
let mut topic_id_guard = self.active_topic_id.lock().await;
let _ = topic_id_guard.take();
*self.active_topic_id.lock().unwrap() = None;
let mut sender_guard = self.active_sender.lock().await;
if let Some(sender) = sender_guard.take() {
if let Some(self_state) = self.self_state.lock().await.as_ref() {
let sender_opt = self.active_sender.lock().unwrap().take();
if let Some(sender) = sender_opt {
let self_state_opt = self.self_state.lock().unwrap().clone();
if let Some(self_state) = self_state_opt {
let payload = GossipPayload {
author: self_state.addr.id,
msg: GossipMessage::Leave,
};
if let Ok(bytes) = serde_json::to_vec(&payload) {
crate::log_msg("Broadcasting Leave message to gossip");
let _ = sender.broadcast(bytes.into()).await;
}
}
}
self.peers.lock().await.clear();
self.peers.lock().unwrap().clear();
Ok(())
}
fn active_peers(&self) -> Vec<(EndpointId, PeerState)> {
let guard = self.peers.blocking_lock();
let guard = self.peers.lock().unwrap();
guard.iter().map(|(k, v)| (*k, v.clone())).collect()
}
async fn subscribe_events(&self) -> Result<Receiver<RoomEvent>, NetError> {
let mut rx_guard = self.event_rx.lock().await;
let mut rx_guard = self.event_rx.lock().unwrap();
if let Some(rx) = rx_guard.take() {
Ok(rx)
} else {