use std::io;
use std::net::SocketAddr;
use std::sync::Arc;
use bytes::BytesMut;
use parking_lot::Mutex;
use pea2pea::{Node as P2pNode, Pea2Pea};
use tokio::sync::broadcast;
use tokio::task::JoinHandle;
use tracing::debug;
use crate::config::NodeConfig;
use crate::error::Error;
use crate::gossip::GossipState;
#[derive(Clone)]
pub struct Node {
config: NodeConfig,
peashape: peashape::Node,
gossip: Arc<GossipState>,
incoming: broadcast::Sender<BytesMut>,
forwarding_handle: Arc<Mutex<Option<JoinHandle<()>>>>,
}
impl Node {
pub fn new(config: NodeConfig) -> Self {
assert!(config.dedup_capacity > 0, "dedup_capacity must be non-zero",);
let peashape = peashape::Node::new(config.to_shape_config());
let gossip = Arc::new(GossipState::new(config.dedup_capacity));
let (incoming, _) = broadcast::channel(peashape::SUBSCRIBER_CAPACITY);
Self {
config,
peashape,
gossip,
incoming,
forwarding_handle: Arc::new(Mutex::new(None)),
}
}
pub async fn spawn(&self) -> io::Result<Option<SocketAddr>> {
let addr = self.peashape.spawn().await?;
let mut rx = self.peashape.subscribe();
let shaper = self.peashape.shaper();
let gossip = self.gossip.clone();
let incoming = self.incoming.clone();
let handle = tokio::spawn(async move {
while let Ok(frame) = rx.recv().await {
if frame.len() < peashape::ID_SIZE {
continue;
}
let mut id = [0u8; peashape::ID_SIZE];
id.copy_from_slice(&frame[..peashape::ID_SIZE]);
if gossip.check_and_record(&id) {
if let Err(e) = shaper.enqueue_raw(
peashape::Lane::Low,
peashape::Target::Broadcast,
frame.clone(),
) {
debug!("could not enqueue a relay frame: {e}");
}
let _ = incoming.send(frame);
}
}
});
*self.forwarding_handle.lock() = Some(handle);
Ok(addr)
}
pub fn publish(&self, payload: &[u8]) -> Result<[u8; peashape::ID_SIZE], Error> {
self.peashape.broadcast_shaped(payload).map_err(Error::from)
}
pub fn subscribe(&self) -> broadcast::Receiver<BytesMut> {
self.incoming.subscribe()
}
pub async fn connect(&self, addr: SocketAddr) -> io::Result<()> {
self.peashape.connect(addr).await
}
pub async fn disconnect(&self, addr: SocketAddr) -> bool {
self.peashape.disconnect(addr).await
}
pub fn connected_peers(&self) -> Vec<SocketAddr> {
self.peashape.connected_peers()
}
pub async fn local_addr(&self) -> io::Result<SocketAddr> {
self.peashape.local_addr().await
}
pub fn peashape(&self) -> &peashape::Node {
&self.peashape
}
pub fn p2p(&self) -> &P2pNode {
self.peashape.p2p()
}
pub fn config(&self) -> &NodeConfig {
&self.config
}
pub async fn shutdown(&self) {
let handle = self.forwarding_handle.lock().take();
if let Some(handle) = handle {
handle.abort();
}
self.peashape.shutdown().await;
}
}
impl Pea2Pea for Node {
fn node(&self) -> &P2pNode {
self.peashape.node()
}
}