zebra_network/
address_book_updater.rs1use std::{
4 cmp::max,
5 net::{IpAddr, SocketAddr},
6 sync::Arc,
7 time::Instant,
8};
9
10use indexmap::IndexMap;
11use thiserror::Error;
12use tokio::{
13 sync::{mpsc, watch},
14 task::JoinHandle,
15};
16use tracing::{Level, Span};
17
18use crate::{
19 address_book::AddressMetrics, meta_addr::MetaAddrChange, AddressBook, BoxError, Config,
20};
21
22pub const MIN_CHANNEL_SIZE: usize = 10;
24
25#[derive(Debug, Eq, PartialEq)]
30pub struct AddressBookUpdater;
31
32#[derive(Copy, Clone, Debug, Error, Eq, PartialEq, Hash)]
33#[error("all address book updater senders are closed")]
34pub struct AllAddressBookUpdaterSendersClosed;
35
36impl AddressBookUpdater {
37 pub fn spawn(
48 config: &Config,
49 local_listener: SocketAddr,
50 ) -> (
51 Arc<std::sync::Mutex<AddressBook>>,
52 watch::Receiver<Arc<IndexMap<IpAddr, Instant>>>,
53 mpsc::Sender<MetaAddrChange>,
54 watch::Receiver<AddressMetrics>,
55 JoinHandle<Result<(), BoxError>>,
56 ) {
57 let (worker_tx, mut worker_rx) = mpsc::channel::<MetaAddrChange>(max(
60 config.peerset_total_connection_limit(),
61 MIN_CHANNEL_SIZE,
62 ));
63
64 let address_book = AddressBook::new(
65 local_listener,
66 &config.network,
67 config.max_connections_per_ip,
68 span!(Level::TRACE, "address book"),
69 );
70 let address_metrics = address_book.address_metrics_watcher();
71 let address_book = Arc::new(std::sync::Mutex::new(address_book));
72
73 #[cfg(feature = "progress-bar")]
74 let (mut address_info, address_bar, never_bar, failed_bar) = {
75 let address_bar = howudoin::new_root().label("Known Peers");
76 let never_bar =
77 howudoin::new_with_parent(address_bar.id()).label("Never Attempted Peers");
78 let failed_bar = howudoin::new_with_parent(never_bar.id()).label("Failed Peers");
79
80 (address_metrics.clone(), address_bar, never_bar, failed_bar)
81 };
82
83 let worker_address_book = address_book.clone();
84 let (bans_sender, bans_receiver) = tokio::sync::watch::channel(
85 worker_address_book
86 .lock()
87 .expect("mutex should be unpoisoned")
88 .bans(),
89 );
90
91 let worker = move || {
92 info!("starting the address book updater");
93
94 while let Some(event) = worker_rx.blocking_recv() {
95 trace!(?event, "got address book change");
96
97 let event_ip = event.addr().ip();
102 let updated = worker_address_book
103 .lock()
104 .expect("mutex should be unpoisoned")
105 .update(event);
106
107 if updated.is_none() {
110 let bans = worker_address_book
111 .lock()
112 .expect("mutex should be unpoisoned")
113 .bans();
114
115 if bans.contains_key(&event_ip) {
116 let _ = bans_sender.send(bans);
117 }
118 }
119
120 #[cfg(feature = "progress-bar")]
121 if matches!(howudoin::cancelled(), Some(true)) {
122 address_bar.close();
123 never_bar.close();
124 failed_bar.close();
125 } else if address_info.has_changed()? {
126 let address_info = *address_info.borrow_and_update();
132
133 address_bar
134 .set_pos(u64::try_from(address_info.num_addresses).expect("fits in u64"));
135 never_bar.set_pos(
138 u64::try_from(address_info.never_attempted_gossiped).expect("fits in u64"),
139 );
140 failed_bar.set_pos(u64::try_from(address_info.failed).expect("fits in u64"));
143 }
145 }
146
147 #[cfg(feature = "progress-bar")]
148 {
149 address_bar.close();
150 never_bar.close();
151 failed_bar.close();
152 }
153
154 let error = Err(AllAddressBookUpdaterSendersClosed.into());
155 info!(?error, "stopping address book updater");
156 error
157 };
158
159 let span = Span::current();
162 let address_book_updater_task_handle =
163 tokio::task::spawn_blocking(move || span.in_scope(worker));
164
165 (
166 address_book,
167 bans_receiver,
168 worker_tx,
169 address_metrics,
170 address_book_updater_task_handle,
171 )
172 }
173}