| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314 |
- /* This file is part of DarkFi (https://dark.fi)
- *
- * Copyright (C) 2020-2024 Dyne.org foundation
- *
- * This program is free software: you can redistribute it and/or modify
- * it under the terms of the GNU Affero General Public License as
- * published by the Free Software Foundation, either version 3 of the
- * License, or (at your option) any later version.
- *
- * This program is distributed in the hope that it will be useful,
- * but WITHOUT ANY WARRANTY; without even the implied warranty of
- * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
- * GNU Affero General Public License for more details.
- *
- * You should have received a copy of the GNU Affero General Public License
- * along with this program. If not, see <https://www.gnu.org/licenses/>.
- */
- use std::{
- collections::HashMap,
- fmt, fs,
- fs::File,
- sync::Arc,
- time::{Instant, UNIX_EPOCH},
- };
- use log::{debug, error, info, trace, warn};
- use rand::{prelude::IteratorRandom, rngs::OsRng, Rng};
- use smol::lock::RwLock;
- use url::Url;
- use super::super::{settings::SettingsPtr, ChannelPtr};
- use crate::{
- system::{Subscriber, SubscriberPtr, Subscription},
- util::{
- file::{load_file, save_file},
- path::expand_path,
- },
- Error, Result,
- };
- // An array containing all possible local host strings
- // TODO: This could perhaps be more exhaustive?
- pub const LOCAL_HOST_STRS: [&str; 2] = ["localhost", "localhost.localdomain"];
- const WHITELIST_MAX_LEN: usize = 5000;
- const GREYLIST_MAX_LEN: usize = 2000;
- /// Atomic pointer to hosts object
- pub type HostsPtr = Arc<Hosts>;
- /// Keeps track of hosts and their current state. Prevents race conditions
- /// where multiple threads are simultaenously trying to change the state of
- /// a given host.
- pub type HostRegistry = RwLock<HashMap<Url, HostState>>;
- /// HostState is a set of mutually exclusive states that can be Pending,
- /// Connected, Disconnected or Refining. The state is `None` when the
- /// corresponding host has been removed from the HostRegistry.
- ///
- /// +----------+
- /// +-- | refining | --+
- /// | +----------+ |
- /// | |
- /// v v
- /// +---------+ +-----------+ +------+
- /// | pending | -> | connected | -> | None |
- /// +---------+ +-----------+ +------+
- /// | ^
- /// | |
- /// | +-------------+ |
- /// +-----> | downgrading | ------+
- /// +-------------+
- ///
- #[derive(Clone, Debug)]
- pub enum HostState {
- /// Hosts that are being connected to in Outbound and Manual Session.
- Pending,
- /// Hosts that have been successfully connected to.
- Connected(ChannelPtr),
- /// Hosts that we have repeatedly failed to connect to, and that are being
- /// removed from the anchorlist and whitelist and added to the greylist.
- Downgrading,
- /// Hosts that are migrating from the greylist to the whitelist or being
- /// removed from the greylist, as defined in `refinery.rs`.
- Refining,
- }
- impl HostState {
- // Try to change state to Downgrading. Only possible if this
- // connection is pending i.e. if we are trying to connect to this
- // host.
- fn try_downgrade(&self) -> Result<Self> {
- match self {
- HostState::Pending => Ok(HostState::Downgrading),
- HostState::Connected(_) => Err(Error::StateBlocked(self.to_string())),
- HostState::Downgrading => Err(Error::StateBlocked(self.to_string())),
- HostState::Refining => Err(Error::StateBlocked(self.to_string())),
- }
- }
- // Try to change state to Refining. Only possible if we are not yet
- // tracking this host in the HostRegistry.
- fn try_refine(&self) -> Result<Self> {
- match self {
- HostState::Pending => Err(Error::StateBlocked(self.to_string())),
- HostState::Connected(_) => Err(Error::StateBlocked(self.to_string())),
- HostState::Downgrading => Err(Error::StateBlocked(self.to_string())),
- HostState::Refining => Err(Error::StateBlocked(self.to_string())),
- }
- }
- // Try to change state to Connected. Possible if this peer is
- // currently Pending or being Refined. The latter is necessary since
- // the refinery process requires us to establish a connection to
- // a peer.
- fn try_connect(&self, channel: ChannelPtr) -> Result<Self> {
- match self {
- HostState::Pending => Ok(HostState::Connected(channel)),
- HostState::Connected(_) => Err(Error::StateBlocked(self.to_string())),
- HostState::Downgrading => Err(Error::StateBlocked(self.to_string())),
- HostState::Refining => Ok(HostState::Connected(channel)),
- }
- }
- // Try to change state to Pending. Only possible if we are not yet
- // tracking this host in the HostRegistry.
- fn try_pending(&self) -> Result<Self> {
- match self {
- HostState::Pending => Err(Error::StateBlocked(self.to_string())),
- HostState::Connected(_) => Err(Error::StateBlocked(self.to_string())),
- HostState::Downgrading => Err(Error::StateBlocked(self.to_string())),
- HostState::Refining => Err(Error::StateBlocked(self.to_string())),
- }
- }
- }
- impl fmt::Display for HostState {
- fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
- fmt::Debug::fmt(self, f)
- }
- }
- #[repr(u8)]
- pub enum HostColor {
- /// Intermediary nodes that are periodically probed and updated to White.
- Grey = 0,
- /// Recently seen hosts. Shared with other nodes.
- White = 1,
- /// Nodes to which we have already been able to establish a connection.
- Gold = 2,
- /// Hostile peers that can neither be connected to nor establish
- /// connections to us for the duration of the program.
- Black = 3,
- }
- /// A Container for managing Grey, White, Gold and Black
- /// hostlists. Exposes a common interface for writing to and querying
- /// hostlists.
- // TODO: Currently hosts (aside from hosts on the Black list) are on
- // multiple lists at once. This needs to be reconsidered.
- // Rethink upgrade/ downgrade methods and consider a single method move() which
- // removes from one hostlist and places on another.
- // TODO: Verify the performance overhead of using vectors for hostlists.
- // TODO: Check whether anchorlist (Gold) has a max size in Monero.
- pub struct HostContainer {
- pub hostlists: [RwLock<Vec<(Url, u64)>>; 4],
- }
- impl HostContainer {
- fn new() -> Self {
- let hostlists: [RwLock<Vec<(Url, u64)>>; 4] = [
- RwLock::new(Vec::new()),
- RwLock::new(Vec::new()),
- RwLock::new(Vec::new()),
- RwLock::new(Vec::new()),
- ];
- Self { hostlists }
- }
- /// Append host to a hostlist.
- pub async fn store(&self, color: usize, addr: Url, last_seen: u64) {
- trace!(target: "net::hosts::store()", "[START]");
- let mut list = self.hostlists[color].write().await;
- list.push((addr, last_seen));
- if color == 0 {
- if list.len() == GREYLIST_MAX_LEN {
- let last_entry = list.pop().unwrap();
- debug!(target: "net::hosts::store()",
- "Greylist reached max size. Removed {:?}", last_entry);
- }
- }
- if color == 1 {
- if list.len() == WHITELIST_MAX_LEN {
- let last_entry = list.pop().unwrap();
- debug!(target: "net::hosts::store()",
- "Whitelist reached max size. Removed {:?}", last_entry);
- }
- }
- // Sort the list by last_seen.
- list.sort_by_key(|entry| entry.1);
- list.reverse();
- trace!(target: "net::hosts::store()", "[END]");
- }
- /// Stores an address on a hostlist or updates its last_seen field if we already
- /// have the address.
- pub async fn store_or_update(&self, color: HostColor, addrs: &[(Url, u64)]) {
- trace!(target: "net::hosts::store_or_update()", "[START]");
- let parent_index = color as usize;
- for (addr, last_seen) in addrs {
- if !self.contains(parent_index, &addr).await {
- debug!(target: "net::hosts::store_or_update()",
- "We do not have this entry in the hostlist. Adding to store...");
- self.store(parent_index, addr.clone(), *last_seen).await;
- } else {
- debug!(target: "net::hosts::store_or_update()",
- "We have this entry in the hostlist. Updating last seen...");
- let child_index = self
- .get_index_at_addr(parent_index, addr.clone())
- .await
- .expect("Expected entry to exist");
- debug!(target: "net::hosts::store_or_update()",
- "Selected index, updating last seen...");
- self.update_last_seen(parent_index, &addr, *last_seen, child_index).await;
- }
- }
- }
- /// Update the last_seen field of a peer on a hostlist.
- pub async fn update_last_seen(&self, color: usize, addr: &Url, last_seen: u64, index: usize) {
- trace!(target: "net::hosts::update_last_seen()", "[START]");
- let mut list = self.hostlists[color].write().await;
- list[index] = (addr.clone(), last_seen);
- list.sort_by_key(|entry| entry.1);
- list.reverse();
- trace!(target: "net::hosts::update_last_seen()", "[END]");
- }
- /// Return all known hosts on a hostlist.
- pub async fn fetch_all(&self, color: HostColor) -> Vec<(Url, u64)> {
- self.hostlists[color as usize].read().await.iter().cloned().collect()
- }
- /// Get the oldest entry from a hostlist.
- pub async fn fetch_last(&self, color: HostColor) -> ((Url, u64), usize) {
- let list = self.hostlists[color as usize].read().await;
- let position = list.len() - 1;
- let entry = &list[position];
- (entry.clone(), position)
- }
- /// TODO: documentation
- pub async fn fetch_address(
- &self,
- color: HostColor,
- transports: &[String],
- transport_mixing: bool,
- ) -> Vec<(Url, u64)> {
- trace!(target: "net::hosts::fetch_address()", "[START]");
- let mut hosts = vec![];
- let index = color as usize;
- // If transport mixing is enabled, then for example we're allowed to
- // use tor:// to connect to tcp:// and tor+tls:// to connect to tcp+tls://.
- // However, **do not** mix tor:// and tcp+tls://, nor tor+tls:// and tcp://.
- macro_rules! mix_transport {
- ($a:expr, $b:expr) => {
- if transports.contains(&$a.to_string()) && transport_mixing {
- let mut a_to_b = self.fetch_with_schemes(index, &[$b.to_string()], None).await;
- for (addr, last_seen) in a_to_b.iter_mut() {
- addr.set_scheme($a).unwrap();
- hosts.push((addr.clone(), last_seen.clone()));
- }
- }
- };
- }
- mix_transport!("tor", "tcp");
- mix_transport!("tor+tls", "tcp+tls");
- mix_transport!("nym", "tcp");
- mix_transport!("nym+tls", "tcp+tls");
- // And now the actual requested transports
- for (addr, last_seen) in self.fetch_with_schemes(index, transports, None).await {
- hosts.push((addr, last_seen));
- }
- trace!(target: "net::hosts::fetch_address()", "Grabbed hosts, length: {}", hosts.len());
- hosts
- }
- /// Get up to limit peers that match the given transport schemes from a hostlist.
- /// If limit was not provided, return all matching peers.
- async fn fetch_with_schemes(
- &self,
- color: usize,
- schemes: &[String],
- limit: Option<usize>,
- ) -> Vec<(Url, u64)> {
- trace!(target: "net::hosts::fetch_with_schemes()", "[START]");
- let list = self.hostlists[color].read().await;
- let mut limit = match limit {
- Some(l) => l.min(list.len()),
- None => list.len(),
- };
- let mut ret = vec![];
- if limit == 0 {
- return ret
- }
- for (addr, last_seen) in list.iter() {
- if schemes.contains(&addr.scheme().to_string()) {
- ret.push((addr.clone(), *last_seen));
- limit -= 1;
- if limit == 0 {
- debug!(target: "net::hosts::fetch_with_schemes()",
- "Found matching scheme, returning {} addresses",
- ret.len());
- return ret
- }
- }
- }
- if ret.is_empty() {
- debug!(target: "net::hosts::fetch_with_schemes()",
- "No such schemes found!")
- }
- ret
- }
- /// Get up to limit peers that don't match the given transport schemes from a hostlist.
- /// If limit was not provided, return all matching peers.
- pub async fn fetch_excluding_schemes(
- &self,
- color: usize,
- schemes: &[String],
- limit: Option<usize>,
- ) -> Vec<(Url, u64)> {
- let list = self.hostlists[color].read().await;
- let mut limit = match limit {
- Some(l) => l.min(list.len()),
- None => list.len(),
- };
- let mut ret = vec![];
- if limit == 0 {
- return ret
- }
- for (addr, last_seen) in list.iter() {
- if !schemes.contains(&addr.scheme().to_string()) {
- ret.push((addr.clone(), *last_seen));
- limit -= 1;
- if limit == 0 {
- return ret
- }
- }
- }
- if ret.is_empty() {
- debug!(target: "net::hosts::fetch_excluding_schemes()",
- "No such schemes found!")
- }
- ret
- }
- /// Get a random peer from a hostlist.
- pub async fn fetch_random(&self, color: HostColor) -> ((Url, u64), usize) {
- let list = self.hostlists[color as usize].read().await;
- let position = rand::thread_rng().gen_range(0..list.len());
- let entry = &list[position];
- (entry.clone(), position)
- }
- /// Get a random peer from a hostlist that matches the given transport schemes.
- pub async fn fetch_random_with_schemes(
- &self,
- color: HostColor,
- schemes: &[String],
- ) -> Option<((Url, u64), usize)> {
- // Retrieve all peers corresponding to that transport schemes
- trace!(target: "net::hosts::fetch_random_with_schemes()", "[START]");
- let list = self.fetch_with_schemes(color as usize, schemes, None).await;
- if list.is_empty() {
- return None
- }
- let position = rand::thread_rng().gen_range(0..list.len());
- let entry = &list[position];
- Some((entry.clone(), position))
- }
- /// Get up to n random peers. Schemes are not taken into account.
- pub async fn fetch_n_random(&self, color: HostColor, n: u32) -> Vec<(Url, u64)> {
- trace!(target: "net::hosts::fetch_n_random()", "[START]");
- let n = n as usize;
- if n == 0 {
- return vec![]
- }
- let mut hosts = vec![];
- let list = self.hostlists[color as usize].read().await;
- for (addr, last_seen) in list.iter() {
- hosts.push((addr.clone(), *last_seen));
- }
- if hosts.is_empty() {
- debug!(target: "net::hosts::fetch_n_random()",
- "No entries found!");
- return hosts
- }
- // Grab random ones
- let urls = hosts.iter().choose_multiple(&mut OsRng, n.min(hosts.len()));
- urls.iter().map(|&url| url.clone()).collect()
- }
- /// Get up to n random peers that match the given transport schemes.
- pub async fn fetch_n_random_with_schemes(
- &self,
- color: HostColor,
- schemes: &[String],
- n: u32,
- ) -> Vec<(Url, u64)> {
- trace!(target: "net::hosts::fetch_n_random_with_schemes()", "[START]");
- let index = color as usize;
- let n = n as usize;
- if n == 0 {
- return vec![]
- }
- // Retrieve all peers corresponding to that transport schemes
- let hosts = self.fetch_with_schemes(index, schemes, None).await;
- if hosts.is_empty() {
- debug!(target: "net::hosts::fetch_n_random_with_schemes()",
- "No such schemes found!");
- return hosts
- }
- // Grab random ones
- let urls = hosts.iter().choose_multiple(&mut OsRng, n.min(hosts.len()));
- urls.iter().map(|&url| url.clone()).collect()
- }
- /// Get up to n random peers that don't match the given transport schemes from
- /// a hostlist.
- pub async fn fetch_n_random_excluding_schemes(
- &self,
- color: HostColor,
- schemes: &[String],
- n: u32,
- ) -> Vec<(Url, u64)> {
- trace!(target: "net::hosts::fetch_excluding_schemes()", "[START]");
- let index = color as usize;
- let n = n as usize;
- if n == 0 {
- return vec![]
- }
- // Retrieve all peers not corresponding to that transport schemes
- let hosts = self.fetch_excluding_schemes(index, schemes, None).await;
- if hosts.is_empty() {
- debug!(target: "net::hosts::fetch_n_random_excluding_schemes()",
- "No such schemes found!");
- return hosts
- }
- // Grab random ones
- let urls = hosts.iter().choose_multiple(&mut OsRng, n.min(hosts.len()));
- urls.iter().map(|&url| url.clone()).collect()
- }
- /// Upgrade a connection to the anchorlist. Called after a connection has been successfully
- /// established in Outbound and Manual sessions.
- pub async fn upgrade_host(&self, addr: &Url) {
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- self.store_or_update(HostColor::Gold, &[(addr.clone(), last_seen)]).await;
- }
- /// Remove an entry from a hostlist.
- pub async fn remove(&self, color: HostColor, addr: &Url, index: usize) {
- debug!(target: "net::hosts::remove()", "Removing peer {} from hostlist", addr);
- let mut list = self.hostlists[color as usize].write().await;
- list.remove(index);
- }
- /// Check if a hostlist is empty.
- pub async fn is_empty(&self, color: HostColor) -> bool {
- self.hostlists[color as usize].read().await.is_empty()
- }
- /// Check if host is in a hostlist
- pub async fn contains(&self, color: usize, addr: &Url) -> bool {
- self.hostlists[color].read().await.iter().any(|(u, _t)| u == addr)
- }
- /// Get the index for a given addr on a hostlist.
- pub async fn get_index_at_addr(&self, color: usize, addr: Url) -> Option<usize> {
- self.hostlists[color].read().await.iter().position(|a| a.0 == addr)
- }
- /// Get the entry for a given addr on the hostlist.
- pub async fn get_entry_at_addr(&self, color: usize, addr: &Url) -> Option<(Url, u64)> {
- self.hostlists[color]
- .read()
- .await
- .iter()
- .find(|(url, _)| url == addr)
- .map(|(url, time)| (url.clone(), *time))
- }
- /// Load the hostlists from a file.
- pub async fn load_all(&self, path: &String) -> Result<()> {
- let path = expand_path(path)?;
- if !path.exists() {
- if let Some(parent) = path.parent() {
- fs::create_dir_all(parent)?;
- }
- File::create(path.clone())?;
- }
- let contents = load_file(&path);
- if let Err(e) = contents {
- warn!(target: "net::hosts::load_hosts()", "Failed retrieving saved hosts: {}", e);
- return Ok(())
- }
- for line in contents.unwrap().lines() {
- let data: Vec<&str> = line.split('\t').collect();
- let url = match Url::parse(data[1]) {
- Ok(u) => u,
- Err(e) => {
- debug!(target: "net::hosts::load_hosts()", "Skipping malformed URL {}", e);
- continue
- }
- };
- let last_seen = match data[2].parse::<u64>() {
- Ok(t) => t,
- Err(e) => {
- debug!(target: "net::hosts::load_hosts()", "Skipping malformed last seen {}", e);
- continue
- }
- };
- match data[0] {
- "greylist" => {
- self.store(HostColor::Grey as usize, url, last_seen).await;
- }
- "whitelist" => {
- self.store(HostColor::White as usize, url, last_seen).await;
- }
- "anchorlist" => {
- self.store(HostColor::Gold as usize, url, last_seen).await;
- }
- _ => {
- debug!(target: "net::hosts::load_hosts()", "Malformed list name...");
- }
- }
- }
- Ok(())
- }
- /// Save the hostlist to a file. Whitelist gets written to the greylist to force
- /// whitelist entries through the refinery on start.
- pub async fn save_all(&self, path: &String) -> Result<()> {
- let path = expand_path(path)?;
- let mut tsv = String::new();
- let mut white = vec![];
- let mut greygold: HashMap<String, Vec<(Url, u64)>> = HashMap::new();
- // First gather all the whitelist entries we don't have in greylist.
- for (url, last_seen) in self.fetch_all(HostColor::White).await {
- if !self.contains(HostColor::Grey as usize, &url).await {
- white.push((url, last_seen))
- }
- }
- // Then gather the greylist and anchorlist entries.
- greygold.insert("anchorlist".to_string(), self.fetch_all(HostColor::Gold).await);
- greygold.insert("greylist".to_string(), self.fetch_all(HostColor::Grey).await);
- // We write whitelist entries to the greylist on p2p.stop() to force
- // them through the refinery on start().
- for (name, mut list) in greygold {
- if name == *"greylist".to_string() {
- list.append(&mut white)
- }
- for (url, last_seen) in list {
- tsv.push_str(&format!("{}\t{}\t{}\n", name, url, last_seen));
- }
- }
- if !tsv.eq("") {
- info!(target: "net::hosts::save_hosts()", "Saving hosts to: {:?}",
- path);
- if let Err(e) = save_file(&path, &tsv) {
- error!(target: "net::hosts::save_hosts()", "Failed saving hosts: {}", e);
- }
- }
- Ok(())
- }
- }
- /// TODO: documentation
- pub struct Hosts {
- /// Subscriber for notifications of new channels
- channel_subscriber: SubscriberPtr<Result<ChannelPtr>>,
- /// Set of stored addresses that are quarantined.
- /// We quarantine peers we've been unable to connect to, but we keep them
- /// around so we can potentially try them again, up to n tries. This should
- /// be helpful in order to self-heal the p2p connections in case we have an
- /// Internet interrupt (goblins unplugging cables)
- quarantine: RwLock<HashMap<Url, usize>>,
- /// A registry that tracks hosts and their current state.
- registry: HostRegistry,
- /// Subscriber listening for store updates
- store_subscriber: SubscriberPtr<usize>,
- /// Pointer to configured P2P settings
- settings: SettingsPtr,
- pub container: HostContainer,
- }
- impl Hosts {
- /// Create a new hosts list>
- pub fn new(settings: SettingsPtr) -> HostsPtr {
- Arc::new(Self {
- channel_subscriber: Subscriber::new(),
- quarantine: RwLock::new(HashMap::new()),
- registry: RwLock::new(HashMap::new()),
- store_subscriber: Subscriber::new(),
- settings,
- container: HostContainer::new(),
- })
- }
- /// Safely insert into the HostContainer. Filters the addresses first before storing and
- /// notifies the subscriber. Must be called when first receiving greylist addresses.
- pub async fn insert(&self, color: HostColor, addrs: &[(Url, u64)]) {
- trace!(target: "net::hosts:insert()", "[START]");
- let filtered_addrs = self.filter_addresses(self.settings.clone(), addrs).await;
- let filtered_addrs_len = filtered_addrs.len();
- if filtered_addrs.is_empty() {
- debug!(target: "net::hosts::insert()", "Filtered out all addresses");
- }
- self.container.store_or_update(color, &filtered_addrs).await;
- self.store_subscriber.notify(filtered_addrs_len).await;
- }
- /// Try to update the registry. If the host already exists, try to update its state.
- /// Otherwise add the host to the registry along with its state.
- pub async fn try_register(&self, addr: Url, new_state: HostState) -> Result<HostState> {
- let mut registry = self.registry.write().await;
- if registry.contains_key(&addr) {
- let current_state = registry.get(&addr).unwrap().clone();
- debug!(target: "net::hosts::try_update_registry()",
- "Attempting to update addr={} current_state={}, new_state={}",
- addr, current_state, new_state.to_string());
- let result: Result<HostState> = match new_state {
- HostState::Pending => current_state.try_pending(),
- HostState::Connected(c) => current_state.try_connect(c),
- HostState::Downgrading => current_state.try_downgrade(),
- HostState::Refining => current_state.try_refine(),
- };
- if let Ok(state) = &result {
- registry.insert(addr.clone(), state.clone());
- }
- result
- } else {
- // We don't know this peer. We can safely update the state.
- registry.insert(addr.clone(), new_state.clone());
- Ok(new_state)
- }
- }
- pub async fn check_address(&self, hosts: Vec<(Url, u64)>) -> Option<(Url, u64)> {
- // Try to find an unused host in the set.
- for (host, last_seen) in hosts {
- debug!(target: "net::hosts::check_address()", "Starting checks");
- if let Err(_) = self.try_register(host.clone(), HostState::Pending).await {
- continue
- }
- debug!(
- target: "net::hosts::check_address()",
- "Found valid host {}",
- host
- );
- return Some((host.clone(), last_seen))
- }
- None
- }
- /// Remove a host from the HostRegistry. Must be called after downgrade(), when the refinery
- /// process fails, or when a channel stops. Prevents hosts from getting trapped in the
- /// HostState logical machinery.
- pub async fn unregister(&self, addr: &Url) {
- debug!(target: "net::hosts::unregister()", "Removing {} from HostRegistry", addr);
- self.registry.write().await.remove(addr);
- }
- /// Returns the list of connected channels.
- pub async fn channels(&self) -> Vec<ChannelPtr> {
- let registry = self.registry.read().await;
- let mut channels = Vec::new();
- for (_, value) in registry.iter() {
- if let HostState::Connected(c) = value {
- channels.push(c.clone());
- }
- }
- channels
- }
- /// Retrieve a random connected channel
- pub async fn random_channel(&self) -> ChannelPtr {
- let channels = self.channels().await;
- let position = rand::thread_rng().gen_range(0..channels.len());
- channels[position].clone()
- }
- /// Add a channel to the set of connected channels
- pub async fn register_channel(&self, channel: ChannelPtr) -> Result<()> {
- let address = channel.address().clone();
- if let Err(e) =
- self.try_register(address.clone(), HostState::Connected(channel.clone())).await
- {
- return Err(e)
- }
- self.channel_subscriber.notify(Ok(channel)).await;
- Ok(())
- }
- pub async fn subscribe_store(&self) -> Result<Subscription<usize>> {
- let sub = self.store_subscriber.clone().subscribe().await;
- Ok(sub)
- }
- // Verify whether a URL is local.
- // NOTE: This function is stateless and not specific to
- // `Hosts`. For this reason, it might make more sense
- // to move this function to a more appropriate location
- // in the codebase.
- /// Check whether a URL is local host
- pub async fn is_local_host(&self, url: Url) -> bool {
- // Reject Urls without host strings.
- if url.host_str().is_none() {
- return false
- }
- // We do this hack in order to parse IPs properly.
- // https://github.com/whatwg/url/issues/749
- let addr = Url::parse(&url.as_str().replace(url.scheme(), "http")).unwrap();
- // Filter private IP ranges
- match addr.host().unwrap() {
- url::Host::Ipv4(ip) => {
- if !ip.is_global() {
- return true
- }
- }
- url::Host::Ipv6(ip) => {
- if !ip.is_global() {
- return true
- }
- }
- url::Host::Domain(d) => {
- if LOCAL_HOST_STRS.contains(&d) {
- return true
- }
- }
- }
- false
- }
- /// Filter given addresses based on certain rulesets and validity.
- async fn filter_addresses(
- &self,
- settings: SettingsPtr,
- addrs: &[(Url, u64)],
- ) -> Vec<(Url, u64)> {
- trace!(target: "net::hosts::filter_addresses()", "Filtering addrs: {:?}", addrs);
- let mut ret = vec![];
- let localnet = self.settings.localnet;
- 'addr_loop: for (addr_, last_seen) in addrs {
- // Validate that the format is `scheme://host_str:port`
- if addr_.host_str().is_none() ||
- addr_.port().is_none() ||
- addr_.cannot_be_a_base() ||
- addr_.path_segments().is_some()
- {
- continue
- }
- if self.container.contains(HostColor::Black as usize, addr_).await {
- warn!(target: "net::hosts::filter_addresses()",
- "Peer {} is blacklisted", addr_);
- continue
- }
- let host_str = addr_.host_str().unwrap();
- if !localnet {
- // Our own external addresses should never enter the hosts set.
- for ext in &settings.external_addrs {
- if host_str == ext.host_str().unwrap() {
- continue 'addr_loop
- }
- }
- }
- // On localnet, make sure ours ports don't enter the host set.
- for ext in &settings.external_addrs {
- if addr_.port() == ext.port() {
- continue 'addr_loop
- }
- }
- // We do this hack in order to parse IPs properly.
- // https://github.com/whatwg/url/issues/749
- let addr = Url::parse(&addr_.as_str().replace(addr_.scheme(), "http")).unwrap();
- // Filter non-global ranges if we're not allowing localnet.
- // Should never be allowed in production, so we don't really care
- // about some of them (e.g. 0.0.0.0, or broadcast, etc.).
- if !localnet && self.is_local_host(addr).await {
- continue
- }
- match addr_.scheme() {
- // Validate that the address is an actual onion.
- #[cfg(feature = "p2p-tor")]
- "tor" | "tor+tls" => {
- use std::str::FromStr;
- if tor_hscrypto::pk::HsId::from_str(host_str).is_err() {
- continue
- }
- trace!(target: "net::hosts::filter_addresses()",
- "[Tor] Valid: {}", host_str);
- }
- #[cfg(feature = "p2p-nym")]
- "nym" | "nym+tls" => continue, // <-- Temp skip
- #[cfg(feature = "p2p-tcp")]
- "tcp" | "tcp+tls" => {
- trace!(target: "net::hosts::filter_addresses()",
- "[TCP] Valid: {}", host_str);
- }
- _ => continue,
- }
- ret.push((addr_.clone(), *last_seen));
- }
- ret
- }
- /// Downgrade a host to greylist. If the host is on the anchorlist or whitelist, remove it.
- /// If it's already on the greylist we can't do anything here.
- pub async fn downgrade_host(&self, addr: &Url, last_seen: u64) {
- if let Err(_) = self.try_register(addr.clone(), HostState::Downgrading).await {
- return
- }
- debug!(target: "net::hosts::downgrade_host()", "Downgrading host {}", addr);
- if self.container.contains(HostColor::Grey as usize, addr).await {
- warn!(target: "net::hosts::downgrade_host()",
- "Cannot downgrade a host that is already on the greylist! {}", addr);
- }
- if self.container.contains(HostColor::Gold as usize, addr).await {
- debug!(target: "net::hosts::downgrade_host()", "Removing from anchorlist {}", addr);
- let index = self
- .container
- .get_index_at_addr(HostColor::Gold as usize, addr.clone())
- .await
- .expect("Expected anchorlist index to exist");
- self.container.remove(HostColor::Gold, addr, index).await;
- self.container.store_or_update(HostColor::Grey, &[(addr.clone(), last_seen)]).await;
- }
- if self.container.contains(HostColor::White as usize, addr).await {
- debug!(target: "net::hosts::downgrade_host()", "Removing from whitelist {}", addr);
- let index = self
- .container
- .get_index_at_addr(HostColor::White as usize, addr.clone())
- .await
- .expect("Expected whitelist index to exist");
- self.container.remove(HostColor::White, addr, index).await;
- self.container.store_or_update(HostColor::Grey, &[(addr.clone(), last_seen)]).await;
- }
- // Remove this entry from HostRegistry to avoid this host getting
- // stuck in the Downgrading state.
- self.unregister(&addr).await;
- }
- /// Quarantine a peer.
- /// If they've been quarantined for more than a configured limit, downgrade to greylist.
- pub async fn quarantine(&self, addr: &Url, last_seen: u64) {
- debug!(target: "net::hosts::quarantine()", "Quarantining peer {}", addr);
- let timer = Instant::now();
- let mut q = self.quarantine.write().await;
- if let Some(retries) = q.get_mut(addr) {
- *retries += 1;
- debug!(target: "net::hosts::quarantine()",
- "Peer {} quarantined {} times", addr, retries);
- if *retries == self.settings.hosts_quarantine_limit {
- debug!(target: "net::hosts::quarantine()",
- "Reached quarantine limited after {:?}", timer.elapsed());
- debug!(target: "net::hosts::quarantine()",
- "Removing from hostlist {}", addr);
- drop(q);
- self.downgrade_host(addr, last_seen).await;
- }
- } else {
- debug!(target: "net::hosts::quarantine()", "Added peer {} to quarantine", addr);
- q.insert(addr.clone(), 0);
- }
- }
- /// Mark a peer as blacklist.
- pub async fn blacklist(&self, peer: &Url) {
- // We ignore UNIX sockets here so we will just work
- // with stuff that has host_str().
- if let Some(_) = peer.host_str() {
- // Localhost connections should never enter the blacklist
- // This however allows any Tor and Nym connections.
- if self.is_local_host(peer.clone()).await {
- return
- }
- // Insert into the blacklist. We set last_seen to 0 (we don't care about this
- // field).
- self.container.hostlists[HostColor::Black as usize]
- .write()
- .await
- .push((peer.clone(), 0));
- }
- }
- }
- #[cfg(test)]
- mod tests {
- use super::{
- super::super::{settings::Settings, P2p},
- *,
- };
- use crate::{net::hosts::refinery::ping_node, system::sleep};
- use smol::Executor;
- #[test]
- fn test_ping_node() {
- smol::block_on(async {
- let settings = Settings {
- localnet: false,
- external_addrs: vec![
- Url::parse("tcp://foo.bar:123").unwrap(),
- Url::parse("tcp://lol.cat:321").unwrap(),
- ],
- ..Default::default()
- };
- let ex = Arc::new(Executor::new());
- let p2p = P2p::new(settings, ex.clone()).await;
- let url = Url::parse("tcp://xeno.systems.wtf").unwrap();
- println!("Pinging node...");
- let task = ex.spawn(ping_node(url.clone(), p2p));
- ex.run(task).await;
- println!("Ping node complete!");
- });
- }
- #[test]
- fn test_is_local_host() {
- smol::block_on(async {
- let settings = Settings {
- localnet: false,
- external_addrs: vec![
- Url::parse("tcp://foo.bar:123").unwrap(),
- Url::parse("tcp://lol.cat:321").unwrap(),
- ],
- ..Default::default()
- };
- let hosts = Hosts::new(Arc::new(settings.clone()));
- let local_hosts: Vec<Url> = vec![
- Url::parse("tcp://localhost").unwrap(),
- Url::parse("tcp://127.0.0.1").unwrap(),
- Url::parse("tcp+tls://[::1]").unwrap(),
- Url::parse("tcp://localhost.localdomain").unwrap(),
- Url::parse("tcp://192.168.10.65").unwrap(),
- ];
- for host in local_hosts {
- eprintln!("{}", host);
- assert!(hosts.is_local_host(host).await);
- }
- let remote_hosts: Vec<Url> = vec![
- Url::parse("https://dyne.org").unwrap(),
- Url::parse("tcp://77.168.10.65:2222").unwrap(),
- Url::parse("tcp://[2345:0425:2CA1:0000:0000:0567:5673:23b5]").unwrap(),
- Url::parse("http://eweiibe6tdjsdprb4px6rqrzzcsi22m4koia44kc5pcjr7nec2rlxyad.onion")
- .unwrap(),
- ];
- for host in remote_hosts {
- assert!(!hosts.is_local_host(host).await)
- }
- });
- }
- #[test]
- fn test_store() {
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- smol::block_on(async {
- let settings = Settings { ..Default::default() };
- let hosts = Hosts::new(Arc::new(settings.clone()));
- let grey_hosts = vec![
- Url::parse("tcp://localhost:3921").unwrap(),
- Url::parse("tor://[::1]:21481").unwrap(),
- Url::parse("tcp://192.168.10.65:311").unwrap(),
- Url::parse("tcp+tls://0.0.0.0:2312").unwrap(),
- Url::parse("tcp://255.255.255.255:2131").unwrap(),
- ];
- for addr in &grey_hosts {
- hosts.container.store(HostColor::Grey as usize, addr.clone(), last_seen).await;
- }
- assert!(!hosts.container.is_empty(HostColor::Grey).await);
- let white_hosts = vec![
- Url::parse("tcp://localhost:3921").unwrap(),
- Url::parse("tor://[::1]:21481").unwrap(),
- Url::parse("tcp://192.168.10.65:311").unwrap(),
- Url::parse("tcp+tls://0.0.0.0:2312").unwrap(),
- Url::parse("tcp://255.255.255.255:2131").unwrap(),
- ];
- for host in &white_hosts {
- hosts.container.store(HostColor::White as usize, host.clone(), last_seen).await;
- }
- assert!(!hosts.container.is_empty(HostColor::White).await);
- let gold_hosts = vec![
- Url::parse("tcp://dark.fi:80").unwrap(),
- Url::parse("tcp://http.cat:401").unwrap(),
- Url::parse("tcp://foo.bar:111").unwrap(),
- ];
- for host in &gold_hosts {
- hosts.container.store(HostColor::Gold as usize, host.clone(), last_seen).await;
- }
- assert!(hosts.container.contains(HostColor::Grey as usize, &grey_hosts[0]).await);
- assert!(hosts.container.contains(HostColor::White as usize, &white_hosts[1]).await);
- assert!(hosts.container.contains(HostColor::Gold as usize, &gold_hosts[2]).await);
- });
- }
- #[test]
- fn test_get_last() {
- smol::block_on(async {
- let settings = Settings { ..Default::default() };
- let hosts = Hosts::new(Arc::new(settings.clone()));
- // Build up a hostlist
- for i in 0..10 {
- sleep(1).await;
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- let url = Url::parse(&format!("tcp://whitelist{}:123", i)).unwrap();
- hosts.container.store(HostColor::White as usize, url.clone(), last_seen).await;
- }
- for (url, last_seen) in
- hosts.container.hostlists[HostColor::White as usize].read().await.iter()
- {
- println!("{} {}", url, last_seen);
- }
- let (entry, _position) = hosts.container.fetch_last(HostColor::White).await;
- println!("last entry: {} {}", entry.0, entry.1);
- });
- }
- #[test]
- fn test_get_entry() {
- smol::block_on(async {
- let settings = Settings { ..Default::default() };
- let hosts = Hosts::new(Arc::new(settings.clone()));
- let url = Url::parse("tcp://dark.renaissance:333").unwrap();
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- hosts.container.store(HostColor::White as usize, url.clone(), last_seen).await;
- hosts.container.store(HostColor::Gold as usize, url.clone(), last_seen).await;
- assert!(hosts
- .container
- .get_entry_at_addr(HostColor::White as usize, &url)
- .await
- .is_some());
- assert!(hosts
- .container
- .get_entry_at_addr(HostColor::Gold as usize, &url)
- .await
- .is_some());
- });
- }
- #[test]
- fn test_remove() {
- smol::block_on(async {
- let settings = Settings { ..Default::default() };
- let hosts = Hosts::new(Arc::new(settings.clone()));
- let url = Url::parse("tcp://dark.renaissance:333").unwrap();
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- hosts.container.store(HostColor::White as usize, url.clone(), last_seen).await;
- sleep(1).await;
- let url = Url::parse("tcp://milady:333").unwrap();
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- hosts.container.store(HostColor::White as usize, url.clone(), last_seen).await;
- sleep(1).await;
- let url = Url::parse("tcp://king-ted:333").unwrap();
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- hosts.container.store(HostColor::White as usize, url.clone(), last_seen).await;
- for (url, last_seen) in
- hosts.container.hostlists[HostColor::White as usize].read().await.iter()
- {
- println!("{}, {}", url, last_seen);
- }
- let position = hosts
- .container
- .get_index_at_addr(HostColor::White as usize, url.clone())
- .await
- .unwrap();
- hosts.container.remove(HostColor::White, &url, position).await;
- for (url, last_seen) in
- hosts.container.hostlists[HostColor::White as usize].read().await.iter()
- {
- println!("{}, {}", url, last_seen);
- }
- });
- }
- #[test]
- fn test_fetch_address() {
- smol::block_on(async {
- let mut hostlist = vec![];
- let mut grey_urls = vec![];
- let mut white_urls = vec![];
- let mut anchor_urls = vec![];
- let ex = Arc::new(Executor::new());
- let settings = Settings { ..Default::default() };
- let p2p = P2p::new(settings, ex.clone()).await;
- let hosts = &p2p.hosts().container;
- // Build up a hostlist
- for i in 0..5 {
- let last_seen = UNIX_EPOCH.elapsed().unwrap().as_secs();
- hosts
- .store(
- HostColor::Grey as usize,
- Url::parse(&format!("tcp://greylist{}:123", i)).unwrap(),
- last_seen,
- )
- .await;
- hosts
- .store(
- HostColor::White as usize,
- Url::parse(&format!("tcp://whitelist{}:123", i)).unwrap(),
- last_seen,
- )
- .await;
- hosts
- .store(
- HostColor::Gold as usize,
- Url::parse(&format!("tcp://anchorlist{}:123", i)).unwrap(),
- last_seen,
- )
- .await;
- grey_urls
- .push((Url::parse(&format!("tcp://greylist{}:123", i)).unwrap(), last_seen));
- white_urls
- .push((Url::parse(&format!("tcp://whitelist{}:123", i)).unwrap(), last_seen));
- anchor_urls
- .push((Url::parse(&format!("tcp://anchorlist{}:123", i)).unwrap(), last_seen));
- }
- assert!(!hosts.is_empty(HostColor::Grey).await);
- assert!(!hosts.is_empty(HostColor::White).await);
- assert!(!hosts.is_empty(HostColor::Gold).await);
- let transports = &vec!["tcp".to_string()];
- let white_count =
- p2p.settings().outbound_connections * p2p.settings().white_connection_percent / 100;
- let localnet = true;
- // Simulate the address selection logic found in outbound_session::fetch_address()
- for i in 0..8 {
- if i < p2p.settings().anchor_connection_count {
- if !hosts.fetch_address(HostColor::Gold, transports, localnet).await.is_empty()
- {
- let addrs =
- hosts.fetch_address(HostColor::Gold, transports, localnet).await;
- hostlist.push(addrs);
- }
- if !hosts.fetch_address(HostColor::White, transports, localnet).await.is_empty()
- {
- let addrs =
- hosts.fetch_address(HostColor::White, transports, localnet).await;
- hostlist.push(addrs);
- }
- if !hosts.fetch_address(HostColor::Grey, transports, localnet).await.is_empty()
- {
- let addrs =
- hosts.fetch_address(HostColor::Grey, transports, localnet).await;
- hostlist.push(addrs);
- }
- } else if i < white_count {
- if !hosts.fetch_address(HostColor::White, transports, localnet).await.is_empty()
- {
- let addrs =
- hosts.fetch_address(HostColor::White, transports, localnet).await;
- hostlist.push(addrs);
- }
- if !hosts.fetch_address(HostColor::Grey, transports, localnet).await.is_empty()
- {
- let addrs =
- hosts.fetch_address(HostColor::Grey, transports, localnet).await;
- hostlist.push(addrs);
- }
- } else if !hosts
- .fetch_address(HostColor::Grey, transports, localnet)
- .await
- .is_empty()
- {
- let addrs = hosts.fetch_address(HostColor::Grey, transports, localnet).await;
- hostlist.push(addrs);
- }
- }
- // Check we're returning the correct addresses.
- anchor_urls.sort();
- white_urls.sort();
- grey_urls.sort();
- hostlist[0].sort();
- hostlist[4].sort();
- hostlist[7].sort();
- assert!(anchor_urls == hostlist[0]);
- assert!(white_urls == hostlist[4]);
- assert!(grey_urls == hostlist[7]);
- })
- }
- }
|