diff --git a/Cargo.toml b/Cargo.toml index 3c7fa44..67e8bb0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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" diff --git a/antigravity.toml b/antigravity.toml index 7fd338c..07454fe 100644 --- a/antigravity.toml +++ b/antigravity.toml @@ -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. diff --git a/src/bin/test_net.rs b/src/bin/test_net.rs new file mode 100644 index 0000000..1cd9866 --- /dev/null +++ b/src/bin/test_net.rs @@ -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> { + 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(()) +} diff --git a/src/core/mod.rs b/src/core/mod.rs index 173907b..9ecdce9 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -25,10 +25,14 @@ impl CoreController { pub fn new(ui_tx: mpsc::Sender) -> Self { let (cmd_tx, cmd_rx) = mpsc::channel(100); - tokio::spawn(async move { - if let Err(e) = run_core_loop(cmd_rx, ui_tx).await { - eprintln!("App core loop failed: {:?}", e); - } + 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 { + 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 = 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)); } } } diff --git a/src/lib.rs b/src/lib.rs new file mode 100644 index 0000000..3b98cf9 --- /dev/null +++ b/src/lib.rs @@ -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); + } + } +} + diff --git a/src/main.rs b/src/main.rs index a3d65fb..af0e8fb 100644 --- a/src/main.rs +++ b/src/main.rs @@ -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); } } diff --git a/src/network/gossip.rs b/src/network/gossip.rs index 11f94df..2fd5e76 100644 --- a/src/network/gossip.rs +++ b/src/network/gossip.rs @@ -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::()?; 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 Ok(bytes) = serde_json::to_vec(&payload) { - let _ = gossip_sender_clone.broadcast(bytes.into()).await; + + 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::(&msg.content) { - 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 = !peer_map.contains_key(&payload.author); - let state_changed = peer_map.get(&payload.author) != Some(&state); - - if is_new { - 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()); - let _ = event_tx.send(RoomEvent::PeerUpdated(payload.author, state)).await; - } + crate::log_msg(&format!("Gossip received Event::Received from delivery={:?}", msg.delivered_from)); + match serde_json::from_slice::(&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; } - GossipMessage::Leave => { - let mut peer_map = peers.lock().await; - if peer_map.remove(&payload.author).is_some() { - let _ = event_tx.send(RoomEvent::PeerLeft(payload.author)).await; + + crate::log_msg(&format!("Gossip Event::Received from author={:?}, payload={:?}", payload.author, payload.msg)); + + match payload.msg { + GossipMessage::Announce(state) => { + 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()); + let _ = event_tx.send(RoomEvent::PeerJoined(payload.author, state)).await; + } else if state_changed { + crate::log_msg(&format!("Gossip peer state updated: {:?}, state: {:?}", payload.author, state)); + let _ = event_tx.send(RoomEvent::PeerUpdated(payload.author, state)).await; + } + } + 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 _ = 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; - if let Some(handle) = handle_guard.take() { - handle.abort(); + 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, 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 {