use super::identity::CommunitasIdentity;
use crate::gossip::GossipContext;
use anyhow::Result;
use async_trait::async_trait;
use blake3;
use saorsa_gossip_types::TopicId;
use saorsa_webrtc::signaling::{SignalingMessage, SignalingTransport};
use std::fmt;
use std::net::SocketAddr;
use std::str::FromStr;
use std::sync::Arc;
use tokio::sync::RwLock;
use tracing::{debug, info, warn};
const WEBRTC_TOPIC_PREFIX: &str = "webrtc.signaling";
#[derive(Debug)]
pub struct GossipSignalingError(anyhow::Error);
impl fmt::Display for GossipSignalingError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}", self.0)
}
}
impl std::error::Error for GossipSignalingError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
self.0.source()
}
}
impl From<anyhow::Error> for GossipSignalingError {
fn from(err: anyhow::Error) -> Self {
GossipSignalingError(err)
}
}
pub struct GossipSignalingTransport {
gossip: Arc<GossipContext>,
local_identity: CommunitasIdentity,
message_queue: Arc<RwLock<Vec<(CommunitasIdentity, SignalingMessage)>>>,
}
impl GossipSignalingTransport {
pub fn new(gossip: Arc<GossipContext>) -> Result<Self> {
let four_words = gossip.four_words.clone();
let local_identity = CommunitasIdentity::new(four_words)?;
Ok(Self {
gossip,
local_identity,
message_queue: Arc::new(RwLock::new(Vec::new())),
})
}
fn peer_topic(&self, peer: &CommunitasIdentity) -> TopicId {
let topic_str = format!("{}.{}", WEBRTC_TOPIC_PREFIX, peer.four_words());
let hash = blake3::hash(topic_str.as_bytes());
TopicId::new(*hash.as_bytes())
}
pub async fn subscribe_to_signaling(&self) -> Result<()> {
let topic = self.peer_topic(&self.local_identity);
info!(
"Subscribing to WebRTC signaling topic for {}",
self.local_identity
);
let pubsub = self.gossip.pubsub.read().await;
let _rx = pubsub.subscribe(topic);
debug!("Subscribed to topic: {:?}", topic);
Ok(())
}
pub async fn process_incoming_messages(&self) -> Result<()> {
Ok(())
}
}
#[async_trait]
impl SignalingTransport for GossipSignalingTransport {
type PeerId = CommunitasIdentity;
type Error = GossipSignalingError;
async fn send_message(
&self,
peer: &Self::PeerId,
message: SignalingMessage,
) -> Result<(), Self::Error> {
info!("Sending signaling message to {}: {:?}", peer, message);
let topic = self.peer_topic(peer);
let payload = (self.local_identity.clone(), message);
let message_bytes = serde_json::to_vec(&payload)
.map_err(|e| anyhow::anyhow!("Failed to serialize signaling message: {}", e))?;
let pubsub = self.gossip.pubsub.write().await;
pubsub
.publish(topic, message_bytes.into())
.await
.map_err(|e| anyhow::anyhow!("Failed to publish signaling message: {}", e))?;
debug!("Signaling message sent to {}", peer);
Ok(())
}
async fn receive_message(&self) -> Result<(Self::PeerId, SignalingMessage), Self::Error> {
self.process_incoming_messages().await?;
let mut queue = self.message_queue.write().await;
if let Some((sender, message)) = queue.pop() {
debug!("Dequeued signaling message from {}", sender);
return Ok((sender, message));
}
Err(anyhow::anyhow!("No signaling messages available").into())
}
async fn discover_peer_endpoint(
&self,
peer: &Self::PeerId,
) -> Result<Option<SocketAddr>, Self::Error> {
info!("Discovering endpoint for peer: {}", peer);
let target_hash = blake3::hash(peer.four_words().as_bytes());
let target_id: [u8; 32] = *target_hash.as_bytes();
let rendezvous = self.gossip.rendezvous.clone();
if let Err(e) = rendezvous.subscribe_to_shard(&target_id).await {
warn!(
"Failed to subscribe to rendezvous shard for {}: {}",
peer, e
);
return Ok(None);
}
let providers = rendezvous.get_providers_for_target(&target_id).await;
if providers.is_empty() {
warn!("Peer {} not found in rendezvous", peer);
return Ok(None);
}
debug!(
"Peer {} is discoverable (found {} providers)",
peer,
providers.len()
);
Ok(None)
}
}
impl FromStr for CommunitasIdentity {
type Err = anyhow::Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
CommunitasIdentity::new(s.to_string())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_peer_topic_generation() {
let identity =
CommunitasIdentity::new("ocean-forest-moon-star".to_string()).expect("valid identity");
let topic_str = format!("{}.{}", WEBRTC_TOPIC_PREFIX, identity.four_words());
let hash = blake3::hash(topic_str.as_bytes());
let topic = TopicId::new(*hash.as_bytes());
assert_eq!(topic_str, "webrtc.signaling.ocean-forest-moon-star");
let hash2 = blake3::hash(topic_str.as_bytes());
let topic2 = TopicId::new(*hash2.as_bytes());
assert_eq!(topic, topic2);
}
#[test]
fn test_identity_from_str() {
let identity =
CommunitasIdentity::from_str("ocean-forest-moon-star").expect("valid identity");
assert_eq!(identity.four_words(), "ocean-forest-moon-star");
}
}