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; use tokio::sync::mpsc::Receiver; use std::sync::{Arc, Mutex}; use std::collections::HashMap; use async_trait::async_trait; use tokio_stream::StreamExt; use serde::{Serialize, Deserialize}; #[derive(Serialize, Deserialize, Clone, Debug)] pub struct GossipPayload { pub author: EndpointId, pub msg: GossipMessage, } #[derive(Serialize, Deserialize, Clone, Debug)] pub enum GossipMessage { Announce(PeerState), Leave, /// A room text-chat message: the author's display name, the text, and a /// sender-stamped millisecond timestamp. Chat { name: String, text: String, ts: u64 }, } pub struct IrohGossipState { _endpoint: Endpoint, gossip: Gossip, address_lookup: iroh::address_lookup::memory::MemoryLookup, self_state: Arc>>, peers: Arc>>, event_tx: mpsc::Sender, event_rx: Mutex>>, active_topic: Mutex>>, active_topic_id: Mutex>, active_sender: Mutex>, } impl IrohGossipState { pub fn new( endpoint: Endpoint, gossip: Gossip, address_lookup: iroh::address_lookup::memory::MemoryLookup, ) -> Self { let (event_tx, event_rx) = mpsc::channel(100); Self { _endpoint: endpoint, gossip, address_lookup, self_state: Arc::new(Mutex::new(None)), peers: Arc::new(Mutex::new(HashMap::new())), event_tx, event_rx: Mutex::new(Some(event_rx)), active_topic: Mutex::new(None), active_topic_id: Mutex::new(None), active_sender: Mutex::new(None), } } } #[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; // Add the host to the address book self.address_lookup.add_endpoint_info(ticket.host_addr.clone()); // 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| { 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().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 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 && 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)) => { 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; } 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; } } GossipMessage::Chat { name, text, ts } => { crate::log_msg(&format!("Gossip chat from author={:?}", payload.author)); let _ = event_tx.send(RoomEvent::ChatMessage { from: payload.author, name, text, ts, }).await; } } } Err(e) => { crate::log_msg(&format!("Gossip failed to deserialize payload: {:?}", e)); } } } 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 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 && 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)); // A NeighborDown is a *transient* loss, not a graceful // leave: emit PeerConnectionLost so the core marks the peer // "reconnecting" and keeps its audio supervisor redialing, // rather than tearing everything down. (Treating this as a // PeerLeft is exactly what defeated reconnect in the field — // it aborted the supervisor ~30s in.) We still drop our // cached presence entry; a rejoin re-announces as new. let removed = peers.lock().unwrap().remove(&peer_id).is_some(); if removed { crate::log_msg(&format!("Peer connection lost (NeighborDown): {:?}", peer_id)); let _ = event_tx.send(RoomEvent::PeerConnectionLost(peer_id)).await; } } 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().unwrap() = Some(handle); Ok(()) } async fn update_self_state(&self, self_state: PeerState) -> Result<(), NetError> { crate::log_msg(&format!("RoomState::update_self_state: state={:?}", self_state)); *self.self_state.lock().unwrap() = Some(self_state.clone()); 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()))?; } } Ok(()) } async fn send_chat(&self, text: String) -> Result<(), NetError> { let (name, author) = { let guard = self.self_state.lock().unwrap(); match guard.as_ref() { Some(s) => (s.name.clone(), s.addr.id), None => return Err(NetError::Other("Not in a room".to_string())), } }; let ts = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_millis() as u64) .unwrap_or(0); let sender_opt = self.active_sender.lock().unwrap().clone(); if let Some(sender) = sender_opt { let payload = GossipPayload { author, msg: GossipMessage::Chat { name, text, ts }, }; if let Ok(bytes) = serde_json::to_vec(&payload) { sender.broadcast(bytes.into()).await .map_err(|e| NetError::Gossip(e.to_string()))?; } } Ok(()) } async fn leave(&self) -> Result<(), NetError> { 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(); } } *self.active_topic_id.lock().unwrap() = None; 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().unwrap().clear(); Ok(()) } fn active_peers(&self) -> Vec<(EndpointId, PeerState)> { 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().unwrap(); if let Some(rx) = rx_guard.take() { Ok(rx) } else { Err(NetError::Other("Events already subscribed".to_string())) } } } #[cfg(test)] mod tests { use super::*; use crate::network::PeerState; use iroh::SecretKey; fn sample_peer_state() -> PeerState { let secret = SecretKey::generate(); let public = secret.public(); let addr = iroh::EndpointAddr::from(public); PeerState { name: "TestPeerGossip".to_string(), is_muted: true, addr, sharing: None, } } #[test] fn test_gossip_message_leave_round_trip() { let original = GossipMessage::Leave; let serialized = serde_json::to_string(&original).unwrap(); let deserialized: GossipMessage = serde_json::from_str(&serialized).unwrap(); assert!(matches!(deserialized, GossipMessage::Leave)); } #[test] fn test_gossip_payload_announce_round_trip() { let peer_state = sample_peer_state(); let author = peer_state.addr.id; let payload = GossipPayload { author, msg: GossipMessage::Announce(peer_state.clone()), }; let serialized = serde_json::to_string(&payload).unwrap(); let deserialized: GossipPayload = serde_json::from_str(&serialized).unwrap(); assert_eq!(deserialized.author, author); match deserialized.msg { GossipMessage::Announce(state) => { assert_eq!(state, peer_state); } GossipMessage::Leave | GossipMessage::Chat { .. } => { panic!("Expected GossipMessage::Announce"); } } } #[test] fn test_gossip_message_chat_round_trip() { // Test normal chat message let original = GossipMessage::Chat { name: "Alice".to_string(), text: "Hello".to_string(), ts: 123456789, }; let serialized = serde_json::to_string(&original).unwrap(); let deserialized: GossipMessage = serde_json::from_str(&serialized).unwrap(); if let GossipMessage::Chat { name, text, ts } = deserialized { assert_eq!(name, "Alice"); assert_eq!(text, "Hello"); assert_eq!(ts, 123456789); } else { panic!("Expected GossipMessage::Chat"); } // Test empty strings and large timestamp let original_empty = GossipMessage::Chat { name: "".to_string(), text: "".to_string(), ts: u64::MAX, }; let serialized_empty = serde_json::to_string(&original_empty).unwrap(); let deserialized_empty: GossipMessage = serde_json::from_str(&serialized_empty).unwrap(); if let GossipMessage::Chat { name, text, ts } = deserialized_empty { assert_eq!(name, ""); assert_eq!(text, ""); assert_eq!(ts, u64::MAX); } else { panic!("Expected GossipMessage::Chat"); } } #[test] fn test_gossip_payload_chat_round_trip() { let peer_state = sample_peer_state(); let author = peer_state.addr.id; let payload = GossipPayload { author, msg: GossipMessage::Chat { name: "Bob".to_string(), text: "Hi there".to_string(), ts: 987654321, }, }; let serialized = serde_json::to_string(&payload).unwrap(); let deserialized: GossipPayload = serde_json::from_str(&serialized).unwrap(); assert_eq!(deserialized.author, author); if let GossipMessage::Chat { name, text, ts } = deserialized.msg { assert_eq!(name, "Bob"); assert_eq!(text, "Hi there"); assert_eq!(ts, 987654321); } else { panic!("Expected GossipMessage::Chat"); } } #[test] fn test_gossip_chat_unicode_round_trip() { let original = GossipMessage::Chat { name: "🎙 User".to_string(), text: "héllo 🎙 世界".to_string(), ts: 1717171717, }; let serialized = serde_json::to_string(&original).unwrap(); let deserialized: GossipMessage = serde_json::from_str(&serialized).unwrap(); if let GossipMessage::Chat { name, text, ts } = deserialized { assert_eq!(name, "🎙 User"); assert_eq!(text, "héllo 🎙 世界"); assert_eq!(ts, 1717171717); } else { panic!("Expected GossipMessage::Chat"); } } }