hopr_transport_p2p/
lib.rs1pub mod constants;
26
27pub mod errors;
29
30pub(crate) mod liveness;
32
33pub mod peer_store;
35
36pub mod swarm;
38
39mod behavior;
41
42use std::{collections::HashSet, sync::Arc};
43
44use dashmap::DashSet;
45use futures::{AsyncRead, AsyncWrite, StreamExt};
46pub use hopr_api::network::Health;
47use hopr_api::network::{NetworkView, traits::NetworkStreamControl};
48use libp2p::{Multiaddr, PeerId};
49
50use crate::liveness::{LivenessRegistry, LivenessStream};
51
52mod utils;
53
54pub use crate::{
55 behavior::{HoprNetworkBehavior, HoprNetworkBehaviorEvent},
56 swarm::HoprLibp2pNetworkBuilder,
57};
58
59#[derive(Debug, Clone)]
61pub enum PeerDiscovery {
62 Announce(PeerId, Vec<Multiaddr>),
63}
64
65#[derive(Clone)]
66pub struct HoprNetwork {
67 tracker: Arc<DashSet<PeerId>>,
68 store: Arc<crate::peer_store::NetworkPeerStore>,
69 control: libp2p_stream::Control,
70 protocol: libp2p::StreamProtocol,
71 event_rx: async_broadcast::InactiveReceiver<hopr_api::network::NetworkEvent>,
72 liveness: LivenessRegistry,
76}
77
78impl std::fmt::Debug for HoprNetwork {
79 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
80 f.debug_struct("HoprNetwork")
81 .field("tracker", &self.tracker)
82 .field("store", &self.store)
83 .field("protocol", &self.protocol)
84 .finish_non_exhaustive()
85 }
86}
87
88impl NetworkView for HoprNetwork {
89 fn listening_as(&self) -> HashSet<Multiaddr> {
90 self.store.get(self.store.me()).unwrap_or_else(|| {
91 tracing::error!("failed to get own peer info from the peer store");
92 std::collections::HashSet::new()
93 })
94 }
95
96 #[inline]
97 fn discovered_peers(&self) -> HashSet<PeerId> {
98 self.store.iter_keys().collect()
99 }
100
101 #[inline]
102 fn connected_peers(&self) -> HashSet<PeerId> {
103 self.tracker.iter().map(|r| *r).collect()
104 }
105
106 fn is_connected(&self, peer: &PeerId) -> bool {
107 self.tracker.contains(peer)
108 }
109
110 #[inline]
111 fn multiaddress_of(&self, peer: &PeerId) -> Option<HashSet<Multiaddr>> {
112 self.store.get(peer)
113 }
114
115 fn health(&self) -> Health {
116 match self.tracker.len() {
117 0 => Health::Red,
118 1 => Health::Orange,
119 2..4 => Health::Yellow,
120 _ => Health::Green,
121 }
122 }
123
124 fn subscribe_network_events(
125 &self,
126 ) -> impl futures::Stream<Item = hopr_api::network::NetworkEvent> + Send + 'static {
127 self.event_rx.clone().activate()
128 }
129}
130
131#[async_trait::async_trait]
132impl NetworkStreamControl for HoprNetwork {
133 fn accept(
134 mut self,
135 ) -> Result<impl futures::Stream<Item = (PeerId, impl AsyncRead + AsyncWrite + Send)> + Send, impl std::error::Error>
136 {
137 let liveness = self.liveness.clone();
138 self.control.accept(self.protocol).map(|stream| {
139 stream.map(move |(peer, inner)| {
140 let flag = liveness.get_or_create_connected(&peer);
141 (peer, LivenessStream::new(inner, flag))
142 })
143 })
144 }
145
146 async fn open(mut self, peer: PeerId) -> Result<impl AsyncRead + AsyncWrite + Send, impl std::error::Error> {
147 let flag = self.liveness.get_or_create_connected(&peer);
148 self.control
149 .open_stream(peer, self.protocol)
150 .await
151 .map(|inner| LivenessStream::new(inner, flag))
152 }
153}