ckb_network/peer_store/
peer_store_impl.rs

1use crate::{
2    Flags, PeerId, SessionType,
3    errors::{PeerStoreError, Result},
4    extract_peer_id, multiaddr_to_socketaddr,
5    network_group::Group,
6    peer_store::{
7        ADDR_COUNT_LIMIT, ADDR_TIMEOUT_MS, ADDR_TRY_TIMEOUT_MS, Behaviour, DIAL_INTERVAL,
8        Multiaddr, PeerScoreConfig, ReportResult, Status,
9        addr_manager::AddrManager,
10        ban_list::BanList,
11        base_addr,
12        types::{AddrInfo, BannedAddr, PeerInfo, ip_to_network},
13    },
14};
15use ipnetwork::IpNetwork;
16use rand::prelude::IteratorRandom;
17use std::collections::{HashMap, hash_map::Entry};
18
19/// Peer store
20///
21/// | -- choose to identify --| --- choose to feeler --- | --      delete     -- |
22/// | 1      | 2     | 3      | 4    | 5    | 6   | 7    | More than seven days  |
23#[derive(Default)]
24pub struct PeerStore {
25    addr_manager: AddrManager,
26    ban_list: BanList,
27    connected_peers: HashMap<PeerId, PeerInfo>,
28    score_config: PeerScoreConfig,
29}
30
31impl PeerStore {
32    /// New with address list and ban list
33    pub fn new(addr_manager: AddrManager, ban_list: BanList) -> Self {
34        PeerStore {
35            addr_manager,
36            ban_list,
37            connected_peers: Default::default(),
38            score_config: Default::default(),
39        }
40    }
41
42    /// this method will assume peer is connected, which implies address is "verified".
43    pub fn add_connected_peer(&mut self, addr: Multiaddr, session_type: SessionType) {
44        let now_ms = ckb_systemtime::unix_time_as_millis();
45        match self
46            .connected_peers
47            .entry(extract_peer_id(&addr).expect("connected addr should have peer id"))
48        {
49            Entry::Occupied(mut entry) => {
50                let peer = entry.get_mut();
51                peer.connected_addr = addr;
52                peer.last_connected_at_ms = now_ms;
53                peer.session_type = session_type;
54            }
55            Entry::Vacant(entry) => {
56                let peer = PeerInfo::new(addr, session_type, now_ms);
57                entry.insert(peer);
58            }
59        }
60    }
61
62    /// Add discovered peer address
63    /// this method will assume peer and addr is untrust since we have not connected to it.
64    pub fn add_addr(&mut self, addr: Multiaddr, flags: Flags) -> Result<()> {
65        if self.ban_list.is_addr_banned(&addr) {
66            return Ok(());
67        }
68        self.check_purge()?;
69        let score = self.score_config.default_score;
70        self.addr_manager
71            .add(AddrInfo::new(addr, 0, score, flags.bits()));
72        Ok(())
73    }
74
75    #[cfg(feature = "fuzz")]
76    pub fn add_addr_fuzz(
77        &mut self,
78        addr: Multiaddr,
79        flags: Flags,
80        last_connected_at_ms: u64,
81        attempts_count: u32,
82    ) -> Result<()> {
83        if self.ban_list.is_addr_banned(&addr) {
84            return Ok(());
85        }
86        self.check_purge()?;
87        let score = self.score_config.default_score;
88        let mut addr_info = AddrInfo::new(addr, last_connected_at_ms, score, flags.bits());
89        addr_info.attempts_count = attempts_count;
90
91        self.addr_manager.add(addr_info);
92        Ok(())
93    }
94
95    /// Add outbound peer address
96    pub fn add_outbound_addr(&mut self, addr: Multiaddr, flags: Flags) {
97        if self.ban_list.is_addr_banned(&addr) {
98            return;
99        }
100        let score = self.score_config.default_score;
101        self.addr_manager.add(AddrInfo::new(
102            addr,
103            ckb_systemtime::unix_time_as_millis(),
104            score,
105            flags.bits(),
106        ));
107    }
108
109    /// Update outbound peer last connected ms
110    pub fn update_outbound_addr_last_connected_ms(&mut self, addr: Multiaddr) {
111        if self.ban_list.is_addr_banned(&addr) {
112            return;
113        }
114        let base_addr = base_addr(&addr);
115        if let Some(info) = self.addr_manager.get_mut(&base_addr) {
116            info.last_connected_at_ms = ckb_systemtime::unix_time_as_millis()
117        }
118    }
119
120    /// Get address manager
121    pub fn addr_manager(&self) -> &AddrManager {
122        &self.addr_manager
123    }
124
125    /// Get mut address manager
126    pub fn mut_addr_manager(&mut self) -> &mut AddrManager {
127        &mut self.addr_manager
128    }
129
130    /// Report peer behaviours
131    pub fn report(&mut self, addr: &Multiaddr, behaviour: Behaviour) -> ReportResult {
132        if let Some(peer_addr) = self.addr_manager.get_mut(addr) {
133            let score = peer_addr.score.saturating_add(behaviour.score());
134            peer_addr.score = score;
135            if score < self.score_config.ban_score {
136                self.ban_addr(
137                    addr,
138                    self.score_config.ban_timeout_ms,
139                    format!("report behaviour {behaviour:?}"),
140                );
141                return ReportResult::Banned;
142            }
143        }
144        ReportResult::Ok
145    }
146
147    /// Remove peer id
148    pub fn remove_disconnected_peer(&mut self, addr: &Multiaddr) -> Option<PeerInfo> {
149        extract_peer_id(addr).and_then(|peer_id| self.connected_peers.remove(&peer_id))
150    }
151
152    /// Get peer status
153    pub fn peer_status(&self, peer_id: &PeerId) -> Status {
154        if self.connected_peers.contains_key(peer_id) {
155            Status::Connected
156        } else {
157            Status::Disconnected
158        }
159    }
160
161    /// Get peers for outbound connection, this method randomly return recently connected peer addrs
162    pub fn fetch_addrs_to_attempt<F>(
163        &mut self,
164        count: usize,
165        required_flags: Flags,
166        filter: F,
167    ) -> Vec<AddrInfo>
168    where
169        F: Fn(&AddrInfo) -> bool,
170    {
171        // Get info:
172        // 1. Not already connected
173        // 2. Connected within 3 days
174
175        let now_ms = ckb_systemtime::unix_time_as_millis();
176        let peers = &self.connected_peers;
177        let addr_expired_ms = now_ms.saturating_sub(ADDR_TRY_TIMEOUT_MS);
178
179        let filter = |peer_addr: &AddrInfo| {
180            filter(peer_addr)
181                && extract_peer_id(&peer_addr.addr)
182                    .map(|peer_id| !peers.contains_key(&peer_id))
183                    .unwrap_or_default()
184                && peer_addr
185                    .connected(|t| t > addr_expired_ms && t <= now_ms.saturating_sub(DIAL_INTERVAL))
186                && required_flags_filter(required_flags, Flags::from_bits_truncate(peer_addr.flags))
187        };
188
189        // get addrs that can attempt.
190        self.addr_manager.fetch_random(count, filter)
191    }
192
193    /// Get peers for feeler connection, this method randomly return peer addrs that we never
194    /// connected to.
195    pub fn fetch_addrs_to_feeler<F>(&mut self, count: usize, filter: F) -> Vec<AddrInfo>
196    where
197        F: Fn(&AddrInfo) -> bool,
198    {
199        // Get info:
200        // 1. Not already connected
201        // 2. Not already tried in a minute
202        // 3. Not connected within 3 days
203
204        let now_ms = ckb_systemtime::unix_time_as_millis();
205        let addr_expired_ms = now_ms.saturating_sub(ADDR_TRY_TIMEOUT_MS);
206        let peers = &self.connected_peers;
207
208        let filter = |peer_addr: &AddrInfo| {
209            filter(peer_addr)
210                && extract_peer_id(&peer_addr.addr)
211                    .map(|peer_id| !peers.contains_key(&peer_id))
212                    .unwrap_or_default()
213                && !peer_addr.tried_in_last_minute(now_ms)
214                && !peer_addr.connected(|t| t > addr_expired_ms)
215        };
216
217        self.addr_manager.fetch_random(count, filter)
218    }
219
220    /// Return address that we never connected to, used for hole punching.
221    pub fn fetch_nat_addrs(&mut self, count: usize, required_flags: Flags) -> Vec<AddrInfo> {
222        // Get info:
223        // 1. Never connected
224        // 2. Not already connected
225        // 3. Ip4 / Ip6 address only
226
227        let peers = &self.connected_peers;
228
229        let filter = |peer_addr: &AddrInfo| {
230            required_flags_filter(required_flags, Flags::from_bits_truncate(peer_addr.flags))
231                && extract_peer_id(&peer_addr.addr)
232                    .map(|peer_id| !peers.contains_key(&peer_id))
233                    .unwrap_or_default()
234                && peer_addr.addr.iter().any(|p| {
235                    matches!(
236                        p,
237                        p2p::multiaddr::Protocol::Ip4(_) | p2p::multiaddr::Protocol::Ip6(_)
238                    )
239                })
240                && peer_addr.last_connected_at_ms == 0
241        };
242
243        self.addr_manager.fetch_random(count, filter)
244    }
245
246    /// Return valid addrs that success connected, used for discovery.
247    pub fn fetch_random_addrs(&mut self, count: usize, required_flags: Flags) -> Vec<AddrInfo> {
248        // Get info:
249        // 1. Connected within 7 days
250
251        let now_ms = ckb_systemtime::unix_time_as_millis();
252        let addr_expired_ms = now_ms.saturating_sub(ADDR_TIMEOUT_MS);
253
254        let filter = |peer_addr: &AddrInfo| {
255            required_flags_filter(required_flags, Flags::from_bits_truncate(peer_addr.flags))
256                && peer_addr.connected(|t| t > addr_expired_ms)
257        };
258
259        // get success connected addrs.
260        self.addr_manager.fetch_random(count, filter)
261    }
262
263    /// Ban an addr
264    pub(crate) fn ban_addr(&mut self, addr: &Multiaddr, timeout_ms: u64, ban_reason: String) {
265        if let Some(addr) = multiaddr_to_socketaddr(addr) {
266            let network = ip_to_network(addr.ip());
267            self.ban_network(network, timeout_ms, ban_reason)
268        }
269        self.addr_manager.remove(addr);
270    }
271
272    pub(crate) fn ban_network(&mut self, network: IpNetwork, timeout_ms: u64, ban_reason: String) {
273        let now_ms = ckb_systemtime::unix_time_as_millis();
274        let ban_addr = BannedAddr {
275            address: network,
276            ban_until: now_ms + timeout_ms,
277            created_at: now_ms,
278            ban_reason,
279        };
280        self.mut_ban_list().ban(ban_addr);
281    }
282
283    /// Whether the address is banned
284    pub fn is_addr_banned(&self, addr: &Multiaddr) -> bool {
285        self.ban_list().is_addr_banned(addr)
286    }
287
288    /// Get ban list
289    pub fn ban_list(&self) -> &BanList {
290        &self.ban_list
291    }
292
293    /// Get mut ban list
294    pub fn mut_ban_list(&mut self) -> &mut BanList {
295        &mut self.ban_list
296    }
297
298    /// Clear ban list
299    pub fn clear_ban_list(&mut self) {
300        std::mem::take(&mut self.ban_list);
301    }
302
303    /// Check and try delete addrs if reach limit
304    /// return Err if peer_store is full and can't be purge
305    fn check_purge(&mut self) -> Result<()> {
306        if self.addr_manager.count() < ADDR_COUNT_LIMIT {
307            return Ok(());
308        }
309
310        // Evicting invalid data in the peer store is a relatively rare operation
311        // There are certain cleanup strategies here:
312        // 1. First evict the nodes that have reached the eviction condition
313        // 2. If the first step is unsuccessful, enter the network segment grouping mode
314        //  2.1. Group current data according to network segment
315        //  2.2. Sort according to the amount of data in the same network segment
316        //  2.3. In the network segment with more than 4 peer, randomly evict 2 peer
317
318        let now_ms = ckb_systemtime::unix_time_as_millis();
319        let candidate_peers: Vec<_> = self
320            .addr_manager
321            .addrs_iter()
322            .filter_map(|addr| {
323                if !addr.is_connectable(now_ms) {
324                    Some(addr.addr.clone())
325                } else {
326                    None
327                }
328            })
329            .collect();
330
331        for key in candidate_peers.iter() {
332            self.addr_manager.remove(key);
333        }
334
335        if candidate_peers.is_empty() {
336            let candidate_peers: Vec<_> = {
337                let mut peers_by_network_group: HashMap<Group, Vec<_>> = HashMap::default();
338                for addr in self.addr_manager.addrs_iter() {
339                    peers_by_network_group
340                        .entry((&addr.addr).into())
341                        .or_default()
342                        .push(addr);
343                }
344                let len = peers_by_network_group.len();
345                let mut peers = peers_by_network_group
346                    .drain()
347                    .map(|(_, v)| v)
348                    .collect::<Vec<Vec<_>>>();
349
350                peers.sort_unstable_by_key(|k| std::cmp::Reverse(k.len()));
351
352                peers
353                    .into_iter()
354                    .take(len / 2)
355                    .flat_map(move |addrs| {
356                        if addrs.len() > 4 {
357                            Some(
358                                addrs
359                                    .iter()
360                                    .choose_multiple(&mut rand::thread_rng(), 2)
361                                    .into_iter()
362                                    .map(|addr| addr.addr.clone())
363                                    .collect::<Vec<Multiaddr>>(),
364                            )
365                        } else {
366                            None
367                        }
368                    })
369                    .flatten()
370                    .collect()
371            };
372
373            for key in candidate_peers.iter() {
374                self.addr_manager.remove(key);
375            }
376
377            if candidate_peers.is_empty() {
378                return Err(PeerStoreError::EvictionFailed.into());
379            }
380        }
381        Ok(())
382    }
383}
384
385pub(crate) fn required_flags_filter(required: Flags, t: Flags) -> bool {
386    if required == Flags::RELAY | Flags::DISCOVERY | Flags::SYNC {
387        t.contains(required) || t.contains(Flags::COMPATIBILITY)
388    } else {
389        t.contains(required)
390    }
391}