1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
use crate::{
application::{
storage::{LockingHashMap, PeerMetadataStorage},
types::{PeerInfo, PeerState},
},
error::NetworkError,
protocols::network::{ApplicationNetworkSender, Message, RpcError},
};
use aptos_config::network_id::{NetworkId, PeerNetworkId};
use aptos_types::PeerId;
use async_trait::async_trait;
use itertools::Itertools;
use std::{collections::HashMap, fmt::Debug, hash::Hash, marker::PhantomData, time::Duration};
#[async_trait]
pub trait NetworkInterface<TMessage: Message + Send, NetworkSender> {
type AppDataKey: Clone + Debug + Eq + Hash;
type AppData: Clone + Debug;
fn peer_metadata_storage(&self) -> &PeerMetadataStorage;
fn sender(&self) -> NetworkSender;
fn connected_peers(&self, network_id: NetworkId) -> HashMap<PeerNetworkId, PeerInfo> {
self.filtered_peers(network_id, |(_, peer_info)| {
peer_info.status == PeerState::Connected
})
}
fn filtered_peers<F: FnMut(&(&PeerId, &PeerInfo)) -> bool>(
&self,
network_id: NetworkId,
filter: F,
) -> HashMap<PeerNetworkId, PeerInfo> {
self.peer_metadata_storage()
.read_filtered(network_id, filter)
}
fn peers(&self, network_id: NetworkId) -> HashMap<PeerNetworkId, PeerInfo> {
self.peer_metadata_storage().read_all(network_id)
}
fn app_data(&self) -> &LockingHashMap<Self::AppDataKey, Self::AppData>;
}
#[derive(Clone, Debug)]
pub struct MultiNetworkSender<
TMessage: Message + Send,
Sender: ApplicationNetworkSender<TMessage> + Send,
> {
senders: HashMap<NetworkId, Sender>,
_phantom: PhantomData<TMessage>,
}
impl<TMessage: Clone + Message + Send, Sender: ApplicationNetworkSender<TMessage> + Send>
MultiNetworkSender<TMessage, Sender>
{
pub fn new(senders: HashMap<NetworkId, Sender>) -> Self {
MultiNetworkSender {
senders,
_phantom: Default::default(),
}
}
fn sender(&self, network_id: &NetworkId) -> &Sender {
self.senders.get(network_id).expect("Unknown NetworkId")
}
pub fn send_to(&self, recipient: PeerNetworkId, message: TMessage) -> Result<(), NetworkError> {
self.sender(&recipient.network_id())
.send_to(recipient.peer_id(), message)
}
pub fn send_to_many(
&self,
recipients: impl Iterator<Item = PeerNetworkId>,
message: TMessage,
) -> Result<(), NetworkError> {
for (network_id, recipients) in
&recipients.group_by(|peer_network_id| peer_network_id.network_id())
{
let sender = self.sender(&network_id);
let peer_ids = recipients.map(|peer_network_id| peer_network_id.peer_id());
sender.send_to_many(peer_ids, message.clone())?;
}
Ok(())
}
pub async fn send_rpc(
&self,
recipient: PeerNetworkId,
req_msg: TMessage,
timeout: Duration,
) -> Result<TMessage, RpcError> {
self.sender(&recipient.network_id())
.send_rpc(recipient.peer_id(), req_msg, timeout)
.await
}
}