use crate::{common::ProtocolError, pb::KaspadMessage, ConnectionInitializer, Peer, Router};
use kaspa_core::{debug, info};
use parking_lot::RwLock;
use std::{collections::HashMap, sync::Arc};
use tokio::sync::mpsc::Receiver as MpscReceiver;
use uuid::Uuid;
#[derive(Debug)]
pub(crate) enum HubEvent {
NewPeer(Arc<Router>),
PeerClosing(Uuid),
}
#[derive(Debug, Clone)]
pub struct Hub {
pub(crate) peers: Arc<RwLock<HashMap<Uuid, Arc<Router>>>>,
}
impl Hub {
pub fn new() -> Self {
Self { peers: Arc::new(RwLock::new(HashMap::new())) }
}
pub(crate) fn start_event_loop(self, mut hub_receiver: MpscReceiver<HubEvent>, initializer: Arc<dyn ConnectionInitializer>) {
tokio::spawn(async move {
while let Some(new_event) = hub_receiver.recv().await {
match new_event {
HubEvent::NewPeer(new_router) => {
match initializer.initialize_connection(new_router.clone()).await {
Ok(_) => {
info!("P2P Connected to {}", new_router);
self.peers.write().insert(new_router.identity(), new_router);
}
Err(err) => {
new_router.close().await;
debug!("P2P, flow initialization for router-id {:?} failed: {}", new_router.identity(), err);
}
}
}
HubEvent::PeerClosing(peer_id) => {
if let Some(router) = self.peers.write().remove(&peer_id) {
debug!("P2P, Hub event loop, removing peer, router-id: {}", router.identity());
}
}
}
}
debug!("P2P, Hub event loop exiting");
});
}
pub async fn send(&self, peer_id: Uuid, msg: KaspadMessage) -> Result<bool, ProtocolError> {
let op = self.peers.read().get(&peer_id).cloned();
if let Some(router) = op {
router.enqueue(msg).await?;
Ok(true)
} else {
Ok(false)
}
}
pub async fn broadcast(&self, msg: KaspadMessage) {
let peers = self.peers.read().values().cloned().collect::<Vec<_>>();
for router in peers {
let _ = router.enqueue(msg.clone()).await;
}
}
pub async fn terminate(&self, peer_id: Uuid) {
let op = self.peers.read().get(&peer_id).cloned();
if let Some(router) = op {
router.close().await;
}
}
pub async fn terminate_all_peers(&self) {
let peers = self.peers.write().drain().map(|(_, r)| r).collect::<Vec<_>>();
for router in peers {
router.close().await;
}
}
pub fn active_peers(&self) -> Vec<Peer> {
self.peers.read().values().map(|r| r.as_ref().into()).collect()
}
pub fn has_peers(&self) -> bool {
!self.peers.read().is_empty()
}
}
impl Default for Hub {
fn default() -> Self {
Self::new()
}
}