From f6a990949e658f553e3880a70efe7e7956979c1c Mon Sep 17 00:00:00 2001 From: shamardy Date: Fri, 30 Jul 2021 15:51:48 +0200 Subject: [PATCH 1/6] version p2p request --- build.rs | 72 ++--- mm2src/database.rs | 9 + mm2src/database/stats_nodes.rs | 82 ++++++ mm2src/lp_native_dex.rs | 17 -- mm2src/lp_network.rs | 23 +- mm2src/lp_stats.rs | 260 +++++++++++++++++++ mm2src/mm2.rs | 5 +- mm2src/mm2_libp2p/src/atomicdex_behaviour.rs | 42 ++- mm2src/rpc/dispatcher/dispatcher_legacy.rs | 4 + 9 files changed, 447 insertions(+), 67 deletions(-) create mode 100644 mm2src/database/stats_nodes.rs create mode 100644 mm2src/lp_stats.rs diff --git a/build.rs b/build.rs index c6ee66cf10..d2facb83a9 100644 --- a/build.rs +++ b/build.rs @@ -2,14 +2,11 @@ use chrono::DateTime; use gstuff::slurp; use regex::Regex; use std::fs; -use std::io::{Read, Write}; +use std::io::Write; use std::path::{Path, PathBuf}; use std::process::Command; use std::str::from_utf8; -/// Absolute path taken from SuperNET's root + `path`. -fn rabs(rrel: &str) -> PathBuf { root().join(rrel) } - fn path2s(path: PathBuf) -> String { path.to_str() .unwrap_or_else(|| panic!("Non-stringy path {:?}", path)) @@ -40,7 +37,6 @@ fn root() -> PathBuf { /// /// The build script will usually help us by putting the MarketMaker version into the “MM_VERSION” file /// and the corresponding ISO 8601 time into the “MM_DATETIME” file -/// (environment variable isn't as useful because we can't `rerun-if-changed` on it). /// /// For the nightly builds the version contains the short commit hash. /// @@ -50,43 +46,32 @@ fn root() -> PathBuf { /// we might skip synchronizing the Git repository there), /// but if it is, then we're going to check if the “MM_DATETIME” and the Git data match. fn mm_version() -> String { - // Try to load the variable from the file. - let mm_version_p = root().join("MM_VERSION"); - let mut buf; - let version = if let Ok(mut mm_version_f) = fs::File::open(&mm_version_p) { - buf = String::new(); - mm_version_f - .read_to_string(&mut buf) - .expect("Can't read from MM_VERSION"); - buf.trim().to_string() - } else { - // If the “MM_VERSION” file is absent then we should create it - // in order for the Cargo dependency management to see it, - // because Cargo will keep rebuilding the `common` crate otherwise. - // - // We should probably fetch the actual git version here, - // with something like `git log '--pretty=format:%h' -n 1` for the nightlies, - // and a release tag when building from some kind of a stable branch, - // though we should keep the ability for the tooling to provide the “MM_VERSION” - // externally, because moving the entire ".git" around is not always practical. - - let mut version = "UNKNOWN".to_string(); - let mut command = Command::new("git"); - command.arg("log").arg("--pretty=format:%h").arg("-n1"); - if let Ok(go) = command.output() { - if go.status.success() { - version = from_utf8(&go.stdout).unwrap().trim().to_string(); - if !Regex::new(r"^\w+$").unwrap().is_match(&version) { - panic!("{}", version) - } + // We fetch the actual git version here, + // with `git log '--pretty=format:%h' -n 1` for the nightlies, + // and a release tag when building from some kind of a stable branch, + // though we should keep the ability for the tooling to provide the “MM_VERSION” + // externally, because moving the entire ".git" around is not always practical. + let mut version = "UNKNOWN".to_string(); + let mut command = Command::new("git"); + command.arg("log").arg("--pretty=format:%h").arg("-n1"); + if let Ok(go) = command.output() { + if go.status.success() { + version = from_utf8(&go.stdout).unwrap().trim().to_string(); + if !Regex::new(r"^\w+$").unwrap().is_match(&version) { + panic!("{}", version) } } + } + + let mm_version_p = root().join("MM_VERSION"); + let v_file = String::from_utf8(slurp(&mm_version_p)).unwrap(); + let v_file = v_file.trim().to_string(); + if version[..] != v_file[..] { + // Create or update the MM_VERSION file in order to appease the Cargo dependency management. + let mut mm_version_f = fs::File::create(&mm_version_p).unwrap(); + mm_version_f.write_all(version.as_bytes()).unwrap(); + } - if let Ok(mut mm_version_f) = fs::File::create(&mm_version_p) { - mm_version_f.write_all(version.as_bytes()).unwrap(); - } - version - }; println!("cargo:rustc-env=MM_VERSION={}", version); let mut dt_git = None; @@ -102,13 +87,12 @@ fn mm_version() -> String { let mm_datetime_p = root().join("MM_DATETIME"); let dt_file = String::from_utf8(slurp(&mm_datetime_p)).unwrap(); - let mut dt_file = dt_file.trim().to_string(); + let dt_file = dt_file.trim().to_string(); if let Some(ref dt_git) = dt_git { if dt_git[..] != dt_file[..] { // Create or update the “MM_DATETIME” file in order to appease the Cargo dependency management. let mut mm_datetime_f = fs::File::create(&mm_datetime_p).unwrap(); mm_datetime_f.write_all(dt_git.as_bytes()).unwrap(); - dt_file = dt_git.clone(); } } @@ -117,8 +101,4 @@ fn mm_version() -> String { version } -fn main() { - println!("cargo:rerun-if-changed={}", path2s(rabs("MM_VERSION"))); - println!("cargo:rerun-if-changed={}", path2s(rabs("MM_DATETIME"))); - mm_version(); -} +fn main() { mm_version(); } diff --git a/mm2src/database.rs b/mm2src/database.rs index 749ac92e35..a1ce6e308a 100644 --- a/mm2src/database.rs +++ b/mm2src/database.rs @@ -4,6 +4,7 @@ pub mod database_common; #[path = "database/my_orders.rs"] pub mod my_orders; #[path = "database/my_swaps.rs"] pub mod my_swaps; +#[path = "database/stats_nodes.rs"] pub mod stats_nodes; #[path = "database/stats_swaps.rs"] pub mod stats_swaps; use crate::CREATE_MY_SWAPS_TABLE; @@ -71,6 +72,13 @@ fn migration_4() -> Vec<(&'static str, Vec)> { stats_swaps::add_and_spli fn migration_5() -> Vec<(&'static str, Vec)> { vec![(my_orders::CREATE_MY_ORDERS_TABLE, vec![])] } +fn migration_6() -> Vec<(&'static str, Vec)> { + vec![ + (stats_nodes::CREATE_NODES_TABLE, vec![]), + (stats_nodes::CREATE_STATS_NODES_TABLE, vec![]), + ] +} + fn statements_for_migration(ctx: &MmArc, current_migration: i64) -> Option)>> { match current_migration { 1 => Some(migration_1(ctx)), @@ -78,6 +86,7 @@ fn statements_for_migration(ctx: &MmArc, current_migration: i64) -> Option Some(migration_3()), 4 => Some(migration_4()), 5 => Some(migration_5()), + 6 => Some(migration_6()), _ => None, } } diff --git a/mm2src/database/stats_nodes.rs b/mm2src/database/stats_nodes.rs new file mode 100644 index 0000000000..4b57e38217 --- /dev/null +++ b/mm2src/database/stats_nodes.rs @@ -0,0 +1,82 @@ +/// This module contains code to work with nodes table for stats collection in MM2 SQLite DB +use crate::mm2::lp_stats::{NodeInfo, NodeVersionStat}; +use common::log::debug; +use common::mm_ctx::MmArc; +use common::rusqlite::{Error as SqlError, Result as SqlResult, NO_PARAMS}; +use std::collections::hash_map::HashMap; + +pub const CREATE_NODES_TABLE: &str = "CREATE TABLE IF NOT EXISTS nodes ( + id INTEGER NOT NULL PRIMARY KEY, + name VARCHAR(255) NOT NULL UNIQUE, + address VARCHAR(255) NOT NULL, + peer_id VARCHAR(255) NOT NULL UNIQUE +);"; + +pub const CREATE_STATS_NODES_TABLE: &str = "CREATE TABLE IF NOT EXISTS stats_nodes ( + id INTEGER NOT NULL PRIMARY KEY, + name VARCHAR(255) NOT NULL, + version VARCHAR(255) NOT NULL, + timestamp INTEGER NOT NULL +);"; + +const INSERT_NODE: &str = "INSERT INTO nodes (name, address, peer_id) VALUES (?1, ?2, ?3)"; + +const DELETE_NODE: &str = "DELETE FROM nodes WHERE name = ?1"; + +const SELECT_PEERS_ADDRESSES: &str = "SELECT peer_id, address FROM nodes"; + +const SELECT_PEERS_NAMES: &str = "SELECT peer_id, name FROM nodes"; + +const INSERT_STAT: &str = "INSERT INTO stats_nodes (name, version, timestamp) VALUES (?1, ?2, ?3)"; + +pub fn insert_node_info(ctx: &MmArc, node_info: &NodeInfo) -> SqlResult<()> { + debug!("Inserting info about node {} to the SQLite database", node_info.name); + let params = vec![ + node_info.name.clone(), + node_info.address.clone(), + node_info.peer_id.clone(), + ]; + let conn = ctx.sqlite_connection(); + conn.execute(INSERT_NODE, ¶ms).map(|_| ()) +} + +pub fn delete_node_info(ctx: &MmArc, name: String) -> SqlResult<()> { + debug!("Deleting info about node {} from the SQLite database", name); + let params = vec![name]; + let conn = ctx.sqlite_connection(); + conn.execute(DELETE_NODE, ¶ms).map(|_| ()) +} + +pub fn select_peers_addresses(ctx: &MmArc) -> SqlResult, SqlError> { + let conn = ctx.sqlite_connection(); + let mut stmt = conn.prepare(SELECT_PEERS_ADDRESSES)?; + let peers_addresses = stmt + .query_map(NO_PARAMS, |row| Ok((row.get(0)?, row.get(1)?)))? + .collect::>>()?; + + Ok(peers_addresses) +} + +pub fn select_peers_names(ctx: &MmArc) -> SqlResult, SqlError> { + let conn = ctx.sqlite_connection(); + let mut stmt = conn.prepare(SELECT_PEERS_NAMES)?; + let peers_names = stmt + .query_map(NO_PARAMS, |row| Ok((row.get(0)?, row.get(1)?)))? + .collect::>>(); + + peers_names +} + +pub fn insert_node_version_stat(ctx: &MmArc, node_version_stat: &NodeVersionStat) -> SqlResult<()> { + debug!( + "Inserting new version stat for node {} to the SQLite database", + node_version_stat.name + ); + let params = vec![ + node_version_stat.name.clone(), + node_version_stat.version.clone(), + node_version_stat.timestamp.to_string(), + ]; + let conn = ctx.sqlite_connection(); + conn.execute(INSERT_STAT, ¶ms).map(|_| ()) +} diff --git a/mm2src/lp_native_dex.rs b/mm2src/lp_native_dex.rs index 4d810be8c4..2318bcedb3 100644 --- a/mm2src/lp_native_dex.rs +++ b/mm2src/lp_native_dex.rs @@ -50,23 +50,6 @@ const NETID_7777_SEEDNODES: [&str; 3] = [ "seed3.defimania.live:0", ]; -pub fn lp_ports(netid: u16) -> Result<(u16, u16, u16), String> { - const LP_RPCPORT: u16 = 7783; - let max_netid = (65535 - 40 - LP_RPCPORT) / 4; - if netid > max_netid { - return ERR!("Netid {} is larger than max {}", netid, max_netid); - } - - let other_ports = if netid != 0 { - let net_mod = netid % 10; - let net_div = netid / 10; - (net_div * 40) + LP_RPCPORT + net_mod - } else { - LP_RPCPORT - }; - Ok((other_ports + 10, other_ports + 20, other_ports + 30)) -} - /// Invokes `OS_ensure_directory`, /// then prints an error and returns `false` if the directory is not writable. fn ensure_dir_is_writable(dir_path: &Path) -> bool { diff --git a/mm2src/lp_network.rs b/mm2src/lp_network.rs index a933ad91b8..00990d65bf 100644 --- a/mm2src/lp_network.rs +++ b/mm2src/lp_network.rs @@ -28,11 +28,12 @@ use mm2_libp2p::{decode_message, encode_message, GossipsubMessage, MessageId, Pe use serde::de; use std::sync::Arc; -use crate::mm2::{lp_ordermatch, lp_swap}; +use crate::mm2::{lp_ordermatch, lp_stats, lp_swap}; #[derive(Eq, Debug, Deserialize, PartialEq, Serialize)] pub enum P2PRequest { Ordermatch(lp_ordermatch::OrdermatchRequest), + NetworkInfo(lp_stats::NetworkInfoRequest), } pub struct P2PContext { @@ -142,6 +143,7 @@ async fn process_p2p_request( let request = try_s!(decode_message::(&request)); let result = match request { P2PRequest::Ordermatch(req) => lp_ordermatch::process_peer_request(ctx.clone(), req).await, + P2PRequest::NetworkInfo(req) => lp_stats::process_info_request(ctx.clone(), req).await, }; let res = match result { @@ -226,6 +228,25 @@ pub async fn request_relays( Ok(parse_peers_responses(responses)) } +pub async fn request_addresses( + ctx: MmArc, + req: P2PRequest, + peers_addresses: Vec<(String, String)>, +) -> Result)>, String> { + let encoded = try_s!(encode_message(&req)); + + let (response_tx, response_rx) = oneshot::channel(); + let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); + let cmd = AdexBehaviourCmd::RequestAddresses { + req: encoded, + peers_addresses, + response_tx, + }; + try_s!(p2p_ctx.cmd_tx.lock().await.try_send(cmd)); + let responses = try_s!(response_rx.await); + Ok(parse_peers_responses(responses)) +} + pub async fn request_peers( ctx: MmArc, req: P2PRequest, diff --git a/mm2src/lp_stats.rs b/mm2src/lp_stats.rs new file mode 100644 index 0000000000..cf5098bbf7 --- /dev/null +++ b/mm2src/lp_stats.rs @@ -0,0 +1,260 @@ +/// The module is responsible for mm2 network stats collection +/// +use common::executor::{spawn, Timer}; +use common::mm_ctx::MmArc; +use common::{log, now_ms}; +use http::Response; +use mm2_libp2p::atomicdex_behaviour::parse_relay_address; +use mm2_libp2p::encode_message; +use serde_json::{self as json, Value as Json}; +use std::collections::HashMap; +use std::net::ToSocketAddrs; + +use crate::mm2::lp_network::{request_addresses, P2PRequest, PeerDecodedResponse}; + +#[derive(Serialize, Deserialize)] +pub struct NodeInfo { + pub name: String, + pub address: String, + pub peer_id: String, +} + +#[derive(Serialize, Deserialize)] +pub struct NodeVersionStat { + pub name: String, + pub version: String, + pub timestamp: u64, +} + +#[cfg(target_arch = "wasm32")] +fn insert_node_info_to_db(_ctx: &MmArc, _node_info: &NodeInfo) -> Result<(), String> { Ok(()) } + +#[cfg(not(target_arch = "wasm32"))] +fn insert_node_info_to_db(ctx: &MmArc, node_info: &NodeInfo) -> Result<(), String> { + crate::mm2::database::stats_nodes::insert_node_info(ctx, node_info).map_err(|e| ERRL!("{}", e)) +} + +#[cfg(target_arch = "wasm32")] +fn insert_node_version_stat_to_db(_ctx: &MmArc, _node_version_stat: &NodeVersionStat) -> Result<(), String> { Ok(()) } + +#[cfg(not(target_arch = "wasm32"))] +fn insert_node_version_stat_to_db(ctx: &MmArc, node_version_stat: &NodeVersionStat) -> Result<(), String> { + crate::mm2::database::stats_nodes::insert_node_version_stat(ctx, node_version_stat).map_err(|e| ERRL!("{}", e)) +} + +#[cfg(not(target_arch = "wasm32"))] +fn delete_node_info_from_db(ctx: &MmArc, name: String) -> Result<(), String> { + crate::mm2::database::stats_nodes::delete_node_info(ctx, name).map_err(|e| ERRL!("{}", e)) +} + +#[cfg(target_arch = "wasm32")] +fn delete_node_info_from_db(_ctx: &MmArc, _name: String) -> Result<(), String> { Ok(()) } + +#[cfg(target_arch = "wasm32")] +pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> Result>, String> {} + +/// Adds node info. to db to be used later for stats collection +#[cfg(not(target_arch = "wasm32"))] +pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> Result>, String> { + let node_info: NodeInfo = try_s!(json::from_value(req)); + let netid = ctx.conf["netid"].as_u64().unwrap_or(0) as u16; + let (_, pubport, _) = try_s!(lp_ports(netid)); + let addr = try_s!(addr_to_ipv4_string(&node_info.address)); + let relay_address = parse_relay_address(addr, pubport); + + let node_info_with_formated_addr = NodeInfo { + name: node_info.name, + address: relay_address.to_string(), + peer_id: node_info.peer_id, + }; + + if let Err(e) = insert_node_info_to_db(&ctx, &node_info_with_formated_addr) { + return ERR!("Error {} on node insertion", e); + } + let res = json!({ + "result": "success" + }); + + return Response::builder() + .body(json::to_vec(&res).expect("Serialization failed")) + .map_err(|e| ERRL!("{}", e)); +} + +#[cfg(target_arch = "wasm32")] +pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> Result>, String> {} + +/// Removes node info. from db to skip collecting stats for this node +#[cfg(not(target_arch = "wasm32"))] +pub async fn remove_node_from_version_stat(ctx: MmArc, req: Json) -> Result>, String> { + let node_name: String = try_s!(json::from_value(req["name"].clone())); + if let Err(e) = delete_node_info_from_db(&ctx, node_name) { + return ERR!("Error {} on node deletion", e); + } + let res = json!({ + "result": "success" + }); + + return Response::builder() + .body(json::to_vec(&res).expect("Serialization failed")) + .map_err(|e| ERRL!("{}", e)); +} + +#[derive(Debug, Deserialize, Serialize)] +struct Mm2VersionRes { + nodes: HashMap, +} + +#[derive(Debug, Deserialize, Eq, PartialEq, Serialize)] +pub enum NetworkInfoRequest { + /// Get MM2 version of nodes added to stats collection + GetMm2Version, +} + +async fn process_get_version_request(ctx: MmArc) -> Result>, String> { + let response = ctx.mm_version().to_string(); + let encoded = try_s!(encode_message(&response)); + Ok(Some(encoded)) +} + +pub async fn process_info_request(ctx: MmArc, request: NetworkInfoRequest) -> Result>, String> { + log::debug!("Got stats request {:?}", request); + match request { + NetworkInfoRequest::GetMm2Version => process_get_version_request(ctx).await, + } +} + +#[cfg(target_arch = "wasm32")] +pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> Result>, String> {} + +#[cfg(not(target_arch = "wasm32"))] +pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> Result>, String> { + let interval: f64 = try_s!(json::from_value(req["interval"].clone())); + + spawn(stat_collection_loop(ctx, interval)); + + let res = json!({ + "result": "success" + }); + Response::builder() + .body(json::to_vec(&res).expect("Serialization failed")) + .map_err(|e| ERRL!("{}", e)) +} + +async fn stat_collection_loop(ctx: MmArc, interval: f64) { + use crate::mm2::database::stats_nodes::{select_peers_addresses, select_peers_names}; + + loop { + if ctx.is_stopping() { + break; + }; + { + let peers_addresses = match select_peers_addresses(&ctx) { + Ok(p) => p, + Err(e) => { + log::error!("Error selecting peers addresses from db: {}", e); + Timer::sleep(10.).await; + continue; + }, + }; + + let peers_names = match select_peers_names(&ctx) { + Ok(n) => n, + Err(e) => { + log::error!("Error selecting peers names from db: {}", e); + Timer::sleep(10.).await; + continue; + }, + }; + + let timestamp = now_ms() / 1000; + let get_versions_res = match request_addresses::( + ctx.clone(), + P2PRequest::NetworkInfo(NetworkInfoRequest::GetMm2Version), + peers_addresses, + ) + .await + { + Ok(res) => res, + Err(e) => { + log::error!("Error getting nodes versions from peers: {}", e); + Timer::sleep(10.).await; + continue; + }, + }; + + for (peer_id, response) in get_versions_res { + let name = match peers_names.get(&peer_id.to_string()) { + Some(n) => n.clone(), + None => continue, + }; + + match response { + PeerDecodedResponse::Ok(version) => { + let node_version_stat = NodeVersionStat { + name, + version, + timestamp, + }; + if let Err(e) = insert_node_version_stat_to_db(&ctx, &node_version_stat) { + log::error!("Error inserting nodes versions into db: {}", e); + continue; + }; + }, + // If a node returns an error or no response it will not be added to the stats table + // A simple count for every node in db will return a count for the number of responses recieved + PeerDecodedResponse::Err(e) => { + log::error!("Node {} responded to version request with error: {}", name, e); + continue; + }, + PeerDecodedResponse::None => { + log::debug!("Node {} did not respond to version request", name); + continue; + }, + } + } + } + Timer::sleep(interval).await; + } +} + +fn addr_to_ipv4_string(address: &str) -> Result { + let address_with_port = if address.contains(':') { + address.to_string() + } else { + format!("{}{}", address, ":0") + }; + match address_with_port.as_str().to_socket_addrs() { + Ok(mut iter) => match iter.next() { + Some(addr) => { + if addr.is_ipv4() { + Ok(addr.ip().to_string()) + } else { + ERR!("Address {} resolved to IPv6 {} which is not supported", address, addr) + } + }, + None => { + ERR!("Address {} to_socket_addrs empty iter", address) + }, + }, + Err(e) => { + ERR!("Couldn't resolve '{}' Address: {}", address, e) + }, + } +} + +pub fn lp_ports(netid: u16) -> Result<(u16, u16, u16), String> { + const LP_RPCPORT: u16 = 7783; + let max_netid = (65535 - 40 - LP_RPCPORT) / 4; + if netid > max_netid { + return ERR!("Netid {} is larger than max {}", netid, max_netid); + } + + let other_ports = if netid != 0 { + let net_mod = netid % 10; + let net_div = netid / 10; + (net_div * 40) + LP_RPCPORT + net_mod + } else { + LP_RPCPORT + }; + Ok((other_ports + 10, other_ports + 20, other_ports + 30)) +} diff --git a/mm2src/mm2.rs b/mm2src/mm2.rs index cb6575dfe4..6342379778 100644 --- a/mm2src/mm2.rs +++ b/mm2src/mm2.rs @@ -39,7 +39,7 @@ use std::ptr::null; use std::str; #[path = "lp_native_dex.rs"] mod lp_native_dex; -use self::lp_native_dex::{lp_init, lp_ports}; +use self::lp_native_dex::lp_init; use coins::update_coins_config; #[cfg(not(target_arch = "wasm32"))] @@ -52,6 +52,9 @@ pub mod database; #[path = "lp_swap.rs"] pub mod lp_swap; #[path = "rpc.rs"] pub mod rpc; +#[path = "lp_stats.rs"] pub mod lp_stats; +use self::lp_stats::lp_ports; + #[cfg(any(test, target_arch = "wasm32"))] #[path = "mm2_tests.rs"] pub mod mm2_tests; diff --git a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs index 0d92dc4650..2e56c747b7 100644 --- a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs +++ b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs @@ -117,6 +117,12 @@ pub enum AdexBehaviourCmd { req: Vec, response_tx: oneshot::Sender>, }, + /// Add addresses to peer exchange then request peers and collect all their responses. + RequestAddresses { + req: Vec, + peers_addresses: Vec<(String, String)>, + response_tx: oneshot::Sender>, + }, /// Send a response using a `response_channel`. SendResponse { /// Response to a request. @@ -286,6 +292,36 @@ impl AtomicDexBehaviour { let future = request_peers(relays, req, self.request_response.sender(), response_tx); self.spawn(future); }, + AdexBehaviourCmd::RequestAddresses { + req, + peers_addresses, + response_tx, + } => { + let peers = peers_addresses + .into_iter() + .filter_map(|(peer, address)| match peer.parse() { + Ok(p) => { + let multi_address = match Multiaddr::from_str(&address) { + Ok(addr) => addr, + Err(e) => { + error!("Error on parse address {}: {:?}", address, e); + return None; + }, + }; + let mut addresses = HashSet::new(); + addresses.insert(multi_address); + self.peers_exchange.add_peer_addresses(&p, addresses); + Some(p) + }, + Err(e) => { + error!("Error on parse peer id {:?}: {:?}", peer, e); + None + }, + }) + .collect(); + let future = request_peers(peers, req, self.request_response.sender(), response_tx); + self.spawn(future); + }, AdexBehaviourCmd::SendResponse { res, response_channel } => { self.request_response.send_response(response_channel.into(), res.into()); }, @@ -770,7 +806,7 @@ pub fn start_gossipsub( /// If te `addr` is in the "/ip4/{addr}/tcp/{port}" format then parse the `addr` immediately to the `Multiaddr`, /// else construct the "/ip4/{addr}/tcp/{port}" from `addr` and `port` values. #[cfg(test)] -fn parse_relay_address(addr: String, port: u16) -> Multiaddr { +pub fn parse_relay_address(addr: String, port: u16) -> Multiaddr { if addr.contains("/ip4/") && addr.contains("/tcp/") { addr.parse().unwrap() } else { @@ -780,7 +816,9 @@ fn parse_relay_address(addr: String, port: u16) -> Multiaddr { /// The addr is expected to be an IP of the relay #[cfg(not(test))] -fn parse_relay_address(addr: String, port: u16) -> Multiaddr { format!("/ip4/{}/tcp/{}", addr, port).parse().unwrap() } +pub fn parse_relay_address(addr: String, port: u16) -> Multiaddr { + format!("/ip4/{}/tcp/{}", addr, port).parse().unwrap() +} /// Request the peers sequential until a `PeerResponse::Ok()` will not be received. async fn request_any_peer( diff --git a/mm2src/rpc/dispatcher/dispatcher_legacy.rs b/mm2src/rpc/dispatcher/dispatcher_legacy.rs index 33f6348fda..2ff5f9f7e4 100644 --- a/mm2src/rpc/dispatcher/dispatcher_legacy.rs +++ b/mm2src/rpc/dispatcher/dispatcher_legacy.rs @@ -12,6 +12,7 @@ use super::lp_commands::*; use crate::mm2::lp_ordermatch::{best_orders_rpc, buy, cancel_all_orders, cancel_order, my_orders, order_status, orderbook_depth_rpc, orderbook_rpc, orders_history_by_filter, sell, set_price, update_maker_order}; +use crate::mm2::lp_stats::{add_node_to_version_stat, remove_node_from_version_stat, start_version_stat_collection}; use crate::mm2::lp_swap::{active_swaps_rpc, all_swaps_uuids_by_filter, ban_pubkey_rpc, coins_needed_for_kick_start, import_swaps, list_banned_pubkeys_rpc, max_taker_vol, my_recent_swaps, my_swap_status, recover_funds_of_swap, stats_swap_status, unban_pubkeys_rpc}; @@ -60,6 +61,7 @@ pub fn dispatcher(req: Json, ctx: MmArc) -> DispatcherRes { // Sorted alphanumerically (on the first latter) for readability. // "autoprice" => lp_autoprice (ctx, req), "active_swaps" => hyres(active_swaps_rpc(ctx, req)), + "add_node_to_version_stat" => hyres(add_node_to_version_stat(ctx, req)), "all_swaps_uuids_by_filter" => all_swaps_uuids_by_filter(ctx, req), "ban_pubkey" => hyres(ban_pubkey_rpc(ctx, req)), "best_orders" => hyres(best_orders_rpc(ctx, req)), @@ -118,12 +120,14 @@ pub fn dispatcher(req: Json, ctx: MmArc) -> DispatcherRes { return DispatcherRes::NoMatch(req); } }, + "remove_node_from_version_stat" => hyres(remove_node_from_version_stat(ctx, req)), "sell" => hyres(sell(ctx, req)), "show_priv_key" => hyres(show_priv_key(ctx, req)), "send_raw_transaction" => hyres(send_raw_transaction(ctx, req)), "set_required_confirmations" => hyres(set_required_confirmations(ctx, req)), "set_requires_notarization" => hyres(set_requires_notarization(ctx, req)), "setprice" => hyres(set_price(ctx, req)), + "start_version_stat_collection" => hyres(start_version_stat_collection(ctx, req)), "stats_swap_status" => stats_swap_status(ctx, req), "stop" => stop(ctx), "trade_preimage" => hyres(into_legacy::trade_preimage(ctx, req)), From 3316dc4b903e55cc29e71f46c17af05e86fed384 Mon Sep 17 00:00:00 2001 From: shamardy Date: Fri, 30 Jul 2021 17:58:30 +0200 Subject: [PATCH 2/6] fix wasm errors --- mm2src/lp_ordermatch.rs | 2 +- mm2src/lp_stats.rs | 37 +++++++++++++++++++++++++++++++------ 2 files changed, 32 insertions(+), 7 deletions(-) diff --git a/mm2src/lp_ordermatch.rs b/mm2src/lp_ordermatch.rs index 254f6bf4cd..772f256c15 100644 --- a/mm2src/lp_ordermatch.rs +++ b/mm2src/lp_ordermatch.rs @@ -3804,7 +3804,7 @@ struct OrderForRpcWithCancellationReason<'a> { } #[cfg(target_arch = "wasm32")] -pub async fn order_status(ctx: MmArc, req: Json) -> Result>, String> { +pub async fn order_status(_ctx: MmArc, _req: Json) -> Result>, String> { let res = json!({ "error": format!("'order_status' is only supported in native mode"), }); diff --git a/mm2src/lp_stats.rs b/mm2src/lp_stats.rs index cf5098bbf7..f54c0fb0ad 100644 --- a/mm2src/lp_stats.rs +++ b/mm2src/lp_stats.rs @@ -42,16 +42,24 @@ fn insert_node_version_stat_to_db(ctx: &MmArc, node_version_stat: &NodeVersionSt crate::mm2::database::stats_nodes::insert_node_version_stat(ctx, node_version_stat).map_err(|e| ERRL!("{}", e)) } +#[cfg(target_arch = "wasm32")] +fn delete_node_info_from_db(_ctx: &MmArc, _name: String) -> Result<(), String> { Ok(()) } + #[cfg(not(target_arch = "wasm32"))] fn delete_node_info_from_db(ctx: &MmArc, name: String) -> Result<(), String> { crate::mm2::database::stats_nodes::delete_node_info(ctx, name).map_err(|e| ERRL!("{}", e)) } #[cfg(target_arch = "wasm32")] -fn delete_node_info_from_db(_ctx: &MmArc, _name: String) -> Result<(), String> { Ok(()) } - -#[cfg(target_arch = "wasm32")] -pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> Result>, String> {} +pub async fn add_node_to_version_stat(_ctx: MmArc, _req: Json) -> Result>, String> { + let res = json!({ + "error": format!("'add_node_to_version_stat' is only supported in native mode"), + }); + Response::builder() + .status(404) + .body(json::to_vec(&res).expect("Serialization failed")) + .map_err(|e| ERRL!("{}", e)) +} /// Adds node info. to db to be used later for stats collection #[cfg(not(target_arch = "wasm32"))] @@ -81,7 +89,15 @@ pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> Result Result>, String> {} +pub async fn remove_node_from_version_stat(_ctx: MmArc, _req: Json) -> Result>, String> { + let res = json!({ + "error": format!("'remove_node_from_version_stat' is only supported in native mode"), + }); + Response::builder() + .status(404) + .body(json::to_vec(&res).expect("Serialization failed")) + .map_err(|e| ERRL!("{}", e)) +} /// Removes node info. from db to skip collecting stats for this node #[cfg(not(target_arch = "wasm32"))] @@ -124,7 +140,15 @@ pub async fn process_info_request(ctx: MmArc, request: NetworkInfoRequest) -> Re } #[cfg(target_arch = "wasm32")] -pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> Result>, String> {} +pub async fn start_version_stat_collection(_ctx: MmArc, _req: Json) -> Result>, String> { + let res = json!({ + "error": format!("'start_version_stat_collection' is only supported in native mode"), + }); + Response::builder() + .status(404) + .body(json::to_vec(&res).expect("Serialization failed")) + .map_err(|e| ERRL!("{}", e)) +} #[cfg(not(target_arch = "wasm32"))] pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> Result>, String> { @@ -140,6 +164,7 @@ pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> Result Date: Tue, 3 Aug 2021 13:17:05 +0200 Subject: [PATCH 3/6] first review fixes --- mm2src/database/stats_nodes.rs | 26 +- mm2src/lp_network.rs | 120 +++++++--- mm2src/lp_stats.rs | 237 +++++++++++-------- mm2src/mm2_libp2p/src/atomicdex_behaviour.rs | 47 +--- mm2src/mm2_libp2p/src/lib.rs | 3 +- mm2src/rpc/dispatcher/dispatcher_legacy.rs | 4 - mm2src/rpc/dispatcher/dispatcher_v2.rs | 6 +- 7 files changed, 248 insertions(+), 195 deletions(-) diff --git a/mm2src/database/stats_nodes.rs b/mm2src/database/stats_nodes.rs index 4b57e38217..0f0b28fc52 100644 --- a/mm2src/database/stats_nodes.rs +++ b/mm2src/database/stats_nodes.rs @@ -15,19 +15,18 @@ pub const CREATE_NODES_TABLE: &str = "CREATE TABLE IF NOT EXISTS nodes ( pub const CREATE_STATS_NODES_TABLE: &str = "CREATE TABLE IF NOT EXISTS stats_nodes ( id INTEGER NOT NULL PRIMARY KEY, name VARCHAR(255) NOT NULL, - version VARCHAR(255) NOT NULL, - timestamp INTEGER NOT NULL + version VARCHAR(255), + timestamp INTEGER NOT NULL, + error VARCHAR(255) );"; const INSERT_NODE: &str = "INSERT INTO nodes (name, address, peer_id) VALUES (?1, ?2, ?3)"; const DELETE_NODE: &str = "DELETE FROM nodes WHERE name = ?1"; -const SELECT_PEERS_ADDRESSES: &str = "SELECT peer_id, address FROM nodes"; - const SELECT_PEERS_NAMES: &str = "SELECT peer_id, name FROM nodes"; -const INSERT_STAT: &str = "INSERT INTO stats_nodes (name, version, timestamp) VALUES (?1, ?2, ?3)"; +const INSERT_STAT: &str = "INSERT INTO stats_nodes (name, version, timestamp, error) VALUES (?1, ?2, ?3, ?4)"; pub fn insert_node_info(ctx: &MmArc, node_info: &NodeInfo) -> SqlResult<()> { debug!("Inserting info about node {} to the SQLite database", node_info.name); @@ -47,16 +46,6 @@ pub fn delete_node_info(ctx: &MmArc, name: String) -> SqlResult<()> { conn.execute(DELETE_NODE, ¶ms).map(|_| ()) } -pub fn select_peers_addresses(ctx: &MmArc) -> SqlResult, SqlError> { - let conn = ctx.sqlite_connection(); - let mut stmt = conn.prepare(SELECT_PEERS_ADDRESSES)?; - let peers_addresses = stmt - .query_map(NO_PARAMS, |row| Ok((row.get(0)?, row.get(1)?)))? - .collect::>>()?; - - Ok(peers_addresses) -} - pub fn select_peers_names(ctx: &MmArc) -> SqlResult, SqlError> { let conn = ctx.sqlite_connection(); let mut stmt = conn.prepare(SELECT_PEERS_NAMES)?; @@ -67,15 +56,16 @@ pub fn select_peers_names(ctx: &MmArc) -> SqlResult, Sql peers_names } -pub fn insert_node_version_stat(ctx: &MmArc, node_version_stat: &NodeVersionStat) -> SqlResult<()> { +pub fn insert_node_version_stat(ctx: &MmArc, node_version_stat: NodeVersionStat) -> SqlResult<()> { debug!( "Inserting new version stat for node {} to the SQLite database", node_version_stat.name ); let params = vec![ - node_version_stat.name.clone(), - node_version_stat.version.clone(), + node_version_stat.name, + node_version_stat.version.unwrap_or_default(), node_version_stat.timestamp.to_string(), + node_version_stat.error.unwrap_or_default(), ]; let conn = ctx.sqlite_connection(); conn.execute(INSERT_STAT, ¶ms).map(|_| ()) diff --git a/mm2src/lp_network.rs b/mm2src/lp_network.rs index 00990d65bf..03314c3abf 100644 --- a/mm2src/lp_network.rs +++ b/mm2src/lp_network.rs @@ -19,10 +19,13 @@ use common::executor::spawn; use common::log; use common::mm_ctx::MmArc; +use common::mm_error::prelude::*; use common::mm_metrics::{ClockOps, MetricsOps}; +use derive_more::Display; use futures::{channel::oneshot, lock::Mutex as AsyncMutex, StreamExt}; use mm2_libp2p::atomicdex_behaviour::{AdexBehaviourCmd, AdexBehaviourEvent, AdexCmdTx, AdexEventRx, AdexResponse, AdexResponseChannel}; +use mm2_libp2p::peers_exchange::PeerAddresses; use mm2_libp2p::{decode_message, encode_message, GossipsubMessage, MessageId, PeerId, TOPIC_SEPARATOR}; #[cfg(test)] use mocktopus::macros::*; use serde::de; @@ -30,6 +33,26 @@ use std::sync::Arc; use crate::mm2::{lp_ordermatch, lp_stats, lp_swap}; +pub type P2PRequestResult = Result>; + +#[derive(Debug, Display)] +pub enum P2PRequestError { + EncodeError(String), + DecodeError(String), + SendError(String), + ResponseError(String), + #[display(fmt = "Expected 1 response, found {}", _0)] + ExpectedSingleResponseError(usize), +} + +impl From for P2PRequestError { + fn from(e: rmp_serde::encode::Error) -> Self { P2PRequestError::EncodeError(e.to_string()) } +} + +impl From for P2PRequestError { + fn from(e: rmp_serde::decode::Error) -> Self { P2PRequestError::DecodeError(e.to_string()) } +} + #[derive(Eq, Debug, Deserialize, PartialEq, Serialize)] pub enum P2PRequest { Ordermatch(lp_ordermatch::OrdermatchRequest), @@ -139,8 +162,8 @@ async fn process_p2p_request( _peer_id: PeerId, request: Vec, response_channel: AdexResponseChannel, -) -> Result<(), String> { - let request = try_s!(decode_message::(&request)); +) -> P2PRequestResult<()> { + let request = decode_message::(&request).map_to_mm(P2PRequestError::from)?; let result = match request { P2PRequest::Ordermatch(req) => lp_ordermatch::process_peer_request(ctx.clone(), req).await, P2PRequest::NetworkInfo(req) => lp_stats::process_info_request(ctx.clone(), req).await, @@ -154,7 +177,12 @@ async fn process_p2p_request( let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); let cmd = AdexBehaviourCmd::SendResponse { res, response_channel }; - try_s!(p2p_ctx.cmd_tx.lock().await.try_send(cmd)); + p2p_ctx + .cmd_tx + .lock() + .await + .try_send(cmd) + .map_to_mm(|e| P2PRequestError::SendError(e.to_string()))?; Ok(()) } @@ -185,8 +213,8 @@ pub async fn subscribe_to_topic(ctx: &MmArc, topic: String) { pub async fn request_any_relay( ctx: MmArc, req: P2PRequest, -) -> Result, String> { - let encoded = try_s!(encode_message(&req)); +) -> P2PRequestResult> { + let encoded = encode_message(&req).map_to_mm(P2PRequestError::from)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -194,10 +222,18 @@ pub async fn request_any_relay( req: encoded, response_tx, }; - try_s!(p2p_ctx.cmd_tx.lock().await.try_send(cmd)); - match try_s!(response_rx.await) { + p2p_ctx + .cmd_tx + .lock() + .await + .try_send(cmd) + .map_to_mm(|e| P2PRequestError::SendError(e.to_string()))?; + match response_rx + .await + .map_to_mm(|e| P2PRequestError::ResponseError(e.to_string()))? + { Some((from_peer, response)) => { - let response = try_s!(decode_message::(&response)); + let response = decode_message::(&response).map_to_mm(P2PRequestError::from)?; Ok(Some((response, from_peer))) }, None => Ok(None), @@ -214,8 +250,8 @@ pub enum PeerDecodedResponse { pub async fn request_relays( ctx: MmArc, req: P2PRequest, -) -> Result)>, String> { - let encoded = try_s!(encode_message(&req)); +) -> P2PRequestResult)>> { + let encoded = encode_message(&req).map_to_mm(P2PRequestError::from)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -223,27 +259,15 @@ pub async fn request_relays( req: encoded, response_tx, }; - try_s!(p2p_ctx.cmd_tx.lock().await.try_send(cmd)); - let responses = try_s!(response_rx.await); - Ok(parse_peers_responses(responses)) -} - -pub async fn request_addresses( - ctx: MmArc, - req: P2PRequest, - peers_addresses: Vec<(String, String)>, -) -> Result)>, String> { - let encoded = try_s!(encode_message(&req)); - - let (response_tx, response_rx) = oneshot::channel(); - let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); - let cmd = AdexBehaviourCmd::RequestAddresses { - req: encoded, - peers_addresses, - response_tx, - }; - try_s!(p2p_ctx.cmd_tx.lock().await.try_send(cmd)); - let responses = try_s!(response_rx.await); + p2p_ctx + .cmd_tx + .lock() + .await + .try_send(cmd) + .map_to_mm(|e| P2PRequestError::SendError(e.to_string()))?; + let responses = response_rx + .await + .map_to_mm(|e| P2PRequestError::ResponseError(e.to_string()))?; Ok(parse_peers_responses(responses)) } @@ -251,8 +275,8 @@ pub async fn request_peers( ctx: MmArc, req: P2PRequest, peers: Vec, -) -> Result)>, String> { - let encoded = try_s!(encode_message(&req)); +) -> P2PRequestResult)>> { + let encoded = encode_message(&req).map_to_mm(P2PRequestError::from)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -261,8 +285,15 @@ pub async fn request_peers( peers, response_tx, }; - try_s!(p2p_ctx.cmd_tx.lock().await.try_send(cmd)); - let responses = try_s!(response_rx.await); + p2p_ctx + .cmd_tx + .lock() + .await + .try_send(cmd) + .map_to_mm(|e| P2PRequestError::SendError(e.to_string()))?; + let responses = response_rx + .await + .map_to_mm(|e| P2PRequestError::ResponseError(e.to_string()))?; Ok(parse_peers_responses(responses)) } @@ -270,20 +301,20 @@ pub async fn request_one_peer( ctx: MmArc, req: P2PRequest, peer: String, -) -> Result, String> { +) -> P2PRequestResult> { let clock = ctx.metrics.clock().expect("Metrics clock is not available"); let start = clock.now(); - let mut responses = try_s!(request_peers::(ctx.clone(), req, vec![peer.clone()]).await); + let mut responses = request_peers::(ctx.clone(), req, vec![peer.clone()]).await?; let end = clock.now(); mm_timing!(ctx.metrics, "peer.outgoing_request.timing", start, end, "peer" => peer); if responses.len() != 1 { - return ERR!("Expected 1 response, found {}", responses.len()); + return MmError::err(P2PRequestError::ExpectedSingleResponseError(responses.len())); } let (_, response) = responses.remove(0); match response { PeerDecodedResponse::Ok(response) => Ok(Some(response)), PeerDecodedResponse::None => Ok(None), - PeerDecodedResponse::Err(e) => ERR!("{}", e), + PeerDecodedResponse::Err(e) => MmError::err(P2PRequestError::ResponseError(e)), } } @@ -319,3 +350,14 @@ pub fn propagate_message(ctx: &MmArc, message_id: MessageId, propagation_source: }; }); } + +pub fn add_peer_addresses(ctx: &MmArc, peer: PeerId, addresses: PeerAddresses) { + let ctx = ctx.clone(); + spawn(async move { + let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); + let cmd = AdexBehaviourCmd::AddPeerAddresses { peer, addresses }; + if let Err(e) = p2p_ctx.cmd_tx.lock().await.try_send(cmd) { + log::error!("add_peer_addresses cmd_tx.send error {:?}", e); + }; + }); +} diff --git a/mm2src/lp_stats.rs b/mm2src/lp_stats.rs index f54c0fb0ad..01d83dffa9 100644 --- a/mm2src/lp_stats.rs +++ b/mm2src/lp_stats.rs @@ -2,15 +2,58 @@ /// use common::executor::{spawn, Timer}; use common::mm_ctx::MmArc; -use common::{log, now_ms}; -use http::Response; +use common::mm_error::prelude::*; +use common::{log, now_ms, HttpStatusCode}; +use derive_more::Display; +use http::StatusCode; use mm2_libp2p::atomicdex_behaviour::parse_relay_address; -use mm2_libp2p::encode_message; +use mm2_libp2p::{encode_message, PeerId}; use serde_json::{self as json, Value as Json}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::net::ToSocketAddrs; -use crate::mm2::lp_network::{request_addresses, P2PRequest, PeerDecodedResponse}; +use crate::mm2::lp_network::{add_peer_addresses, request_peers, P2PRequest, PeerDecodedResponse}; + +pub type NodeVersionResult = Result>; + +#[derive(Debug, Deserialize, Display, Serialize, SerializeErrorType)] +#[serde(tag = "error_type", content = "error_data")] +pub enum NodeVersionError { + #[display(fmt = "Invalid request: {}", _0)] + InvalidRequest(String), + #[display(fmt = "Database error: {}", _0)] + DatabaseError(String), + #[display(fmt = "Invalid address: {}", _0)] + InvalidAddress(String), + #[display(fmt = "Error on parse peer id {}", _0)] + PeerIdParseError(String), + #[display(fmt = "{} is only supported in native mode", _0)] + UnsupportedMode(String), +} + +impl HttpStatusCode for NodeVersionError { + fn status_code(&self) -> StatusCode { + match self { + NodeVersionError::InvalidRequest(_) + | NodeVersionError::InvalidAddress(_) + | NodeVersionError::PeerIdParseError(_) => StatusCode::BAD_REQUEST, + NodeVersionError::UnsupportedMode(_) => StatusCode::METHOD_NOT_ALLOWED, + NodeVersionError::DatabaseError(_) => StatusCode::INTERNAL_SERVER_ERROR, + } + } +} + +impl From for NodeVersionError { + fn from(e: serde_json::Error) -> Self { NodeVersionError::InvalidRequest(e.to_string()) } +} + +impl From for NodeVersionError { + fn from(e: NetIdError) -> Self { NodeVersionError::InvalidAddress(e.to_string()) } +} + +impl From for NodeVersionError { + fn from(e: ParseAddressError) -> Self { NodeVersionError::InvalidAddress(e.to_string()) } +} #[derive(Serialize, Deserialize)] pub struct NodeInfo { @@ -22,8 +65,9 @@ pub struct NodeInfo { #[derive(Serialize, Deserialize)] pub struct NodeVersionStat { pub name: String, - pub version: String, + pub version: Option, pub timestamp: u64, + pub error: Option, } #[cfg(target_arch = "wasm32")] @@ -35,10 +79,10 @@ fn insert_node_info_to_db(ctx: &MmArc, node_info: &NodeInfo) -> Result<(), Strin } #[cfg(target_arch = "wasm32")] -fn insert_node_version_stat_to_db(_ctx: &MmArc, _node_version_stat: &NodeVersionStat) -> Result<(), String> { Ok(()) } +fn insert_node_version_stat_to_db(_ctx: &MmArc, _node_version_stat: NodeVersionStat) -> Result<(), String> { Ok(()) } #[cfg(not(target_arch = "wasm32"))] -fn insert_node_version_stat_to_db(ctx: &MmArc, node_version_stat: &NodeVersionStat) -> Result<(), String> { +fn insert_node_version_stat_to_db(ctx: &MmArc, node_version_stat: NodeVersionStat) -> Result<(), String> { crate::mm2::database::stats_nodes::insert_node_version_stat(ctx, node_version_stat).map_err(|e| ERRL!("{}", e)) } @@ -51,25 +95,29 @@ fn delete_node_info_from_db(ctx: &MmArc, name: String) -> Result<(), String> { } #[cfg(target_arch = "wasm32")] -pub async fn add_node_to_version_stat(_ctx: MmArc, _req: Json) -> Result>, String> { - let res = json!({ - "error": format!("'add_node_to_version_stat' is only supported in native mode"), - }); - Response::builder() - .status(404) - .body(json::to_vec(&res).expect("Serialization failed")) - .map_err(|e| ERRL!("{}", e)) +pub async fn add_node_to_version_stat(_ctx: MmArc, _req: Json) -> NodeVersionResult { + MmError::err(NodeVersionError::UnsupportedMode("'add_node_to_version_stat'".into())) } /// Adds node info. to db to be used later for stats collection #[cfg(not(target_arch = "wasm32"))] -pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> Result>, String> { - let node_info: NodeInfo = try_s!(json::from_value(req)); +pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> NodeVersionResult { + let node_info: NodeInfo = json::from_value(req).map_to_mm(NodeVersionError::from)?; let netid = ctx.conf["netid"].as_u64().unwrap_or(0) as u16; - let (_, pubport, _) = try_s!(lp_ports(netid)); - let addr = try_s!(addr_to_ipv4_string(&node_info.address)); + let (_, pubport, _) = lp_ports(netid)?; + let addr = addr_to_ipv4_string(&node_info.address)?; let relay_address = parse_relay_address(addr, pubport); + let mut addresses = HashSet::new(); + addresses.insert(relay_address.clone()); + + let peer_id: PeerId = match node_info.peer_id.parse() { + Ok(p) => p, + Err(e) => return MmError::err(NodeVersionError::PeerIdParseError(e.to_string())), + }; + + add_peer_addresses(&ctx, peer_id, addresses); + let node_info_with_formated_addr = NodeInfo { name: node_info.name, address: relay_address.to_string(), @@ -77,42 +125,28 @@ pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> Result Result>, String> { - let res = json!({ - "error": format!("'remove_node_from_version_stat' is only supported in native mode"), - }); - Response::builder() - .status(404) - .body(json::to_vec(&res).expect("Serialization failed")) - .map_err(|e| ERRL!("{}", e)) +pub async fn remove_node_from_version_stat(_ctx: MmArc, _req: Json) -> NodeVersionResult { + MmError::err(NodeVersionError::UnsupportedMode( + "'remove_node_from_version_stat'".into(), + )) } /// Removes node info. from db to skip collecting stats for this node #[cfg(not(target_arch = "wasm32"))] -pub async fn remove_node_from_version_stat(ctx: MmArc, req: Json) -> Result>, String> { - let node_name: String = try_s!(json::from_value(req["name"].clone())); +pub async fn remove_node_from_version_stat(ctx: MmArc, req: Json) -> NodeVersionResult { + let node_name: String = json::from_value(req["name"].clone()).map_to_mm(NodeVersionError::from)?; if let Err(e) = delete_node_info_from_db(&ctx, node_name) { - return ERR!("Error {} on node deletion", e); + return MmError::err(NodeVersionError::DatabaseError(e)); } - let res = json!({ - "result": "success" - }); - return Response::builder() - .body(json::to_vec(&res).expect("Serialization failed")) - .map_err(|e| ERRL!("{}", e)); + Ok("success".into()) } #[derive(Debug, Deserialize, Serialize)] @@ -140,48 +174,30 @@ pub async fn process_info_request(ctx: MmArc, request: NetworkInfoRequest) -> Re } #[cfg(target_arch = "wasm32")] -pub async fn start_version_stat_collection(_ctx: MmArc, _req: Json) -> Result>, String> { - let res = json!({ - "error": format!("'start_version_stat_collection' is only supported in native mode"), - }); - Response::builder() - .status(404) - .body(json::to_vec(&res).expect("Serialization failed")) - .map_err(|e| ERRL!("{}", e)) +pub async fn start_version_stat_collection(_ctx: MmArc, _req: Json) -> NodeVersionResult { + MmError::err(NodeVersionError::UnsupportedMode( + "'start_version_stat_collection'".into(), + )) } #[cfg(not(target_arch = "wasm32"))] -pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> Result>, String> { - let interval: f64 = try_s!(json::from_value(req["interval"].clone())); +pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> NodeVersionResult { + let interval: f64 = json::from_value(req["interval"].clone()).map_to_mm(NodeVersionError::from)?; spawn(stat_collection_loop(ctx, interval)); - let res = json!({ - "result": "success" - }); - Response::builder() - .body(json::to_vec(&res).expect("Serialization failed")) - .map_err(|e| ERRL!("{}", e)) + Ok("success".into()) } #[cfg(not(target_arch = "wasm32"))] async fn stat_collection_loop(ctx: MmArc, interval: f64) { - use crate::mm2::database::stats_nodes::{select_peers_addresses, select_peers_names}; + use crate::mm2::database::stats_nodes::select_peers_names; loop { if ctx.is_stopping() { break; }; { - let peers_addresses = match select_peers_addresses(&ctx) { - Ok(p) => p, - Err(e) => { - log::error!("Error selecting peers addresses from db: {}", e); - Timer::sleep(10.).await; - continue; - }, - }; - let peers_names = match select_peers_names(&ctx) { Ok(n) => n, Err(e) => { @@ -191,11 +207,13 @@ async fn stat_collection_loop(ctx: MmArc, interval: f64) { }, }; + let peers: Vec = peers_names.keys().cloned().collect(); + let timestamp = now_ms() / 1000; - let get_versions_res = match request_addresses::( + let get_versions_res = match request_peers::( ctx.clone(), P2PRequest::NetworkInfo(NetworkInfoRequest::GetMm2Version), - peers_addresses, + peers, ) .await { @@ -214,26 +232,44 @@ async fn stat_collection_loop(ctx: MmArc, interval: f64) { }; match response { - PeerDecodedResponse::Ok(version) => { + PeerDecodedResponse::Ok(v) => { let node_version_stat = NodeVersionStat { - name, - version, + name: name.clone(), + version: Some(v.clone()), timestamp, + error: None, }; - if let Err(e) = insert_node_version_stat_to_db(&ctx, &node_version_stat) { - log::error!("Error inserting nodes versions into db: {}", e); - continue; + if let Err(e) = insert_node_version_stat_to_db(&ctx, node_version_stat) { + log::error!("Error inserting node {} version {} into db: {}", name, v, e); }; }, - // If a node returns an error or no response it will not be added to the stats table - // A simple count for every node in db will return a count for the number of responses recieved PeerDecodedResponse::Err(e) => { - log::error!("Node {} responded to version request with error: {}", name, e); - continue; + log::error!( + "Node {} responded to version request with error: {}", + name.clone(), + e.clone() + ); + let node_version_stat = NodeVersionStat { + name: name.clone(), + version: None, + timestamp, + error: Some(e.clone()), + }; + if let Err(e) = insert_node_version_stat_to_db(&ctx, node_version_stat) { + log::error!("Error inserting node {} error into db: {}", name, e); + }; }, PeerDecodedResponse::None => { - log::debug!("Node {} did not respond to version request", name); - continue; + log::debug!("Node {} did not respond to version request", name.clone()); + let node_version_stat = NodeVersionStat { + name: name.clone(), + version: None, + timestamp, + error: None, + }; + if let Err(e) = insert_node_version_stat_to_db(&ctx, node_version_stat) { + log::error!("Error inserting no response for node {} into db: {}", name, e); + }; }, } } @@ -242,7 +278,18 @@ async fn stat_collection_loop(ctx: MmArc, interval: f64) { } } -fn addr_to_ipv4_string(address: &str) -> Result { +#[derive(Debug, Display)] +pub enum ParseAddressError { + #[display(fmt = "Address {} resolved to IPv6 which is not supported", _0)] + UnsupportedIPv6Address(String), + #[display(fmt = "Address {} to_socket_addrs empty iter", _0)] + EmptyIterator(String), + // error return for second string + #[display(fmt = "Couldn't resolve '{}' Address: {}", _0, _1)] + UnresolvedAddress(String, String), +} + +fn addr_to_ipv4_string(address: &str) -> Result> { let address_with_port = if address.contains(':') { address.to_string() } else { @@ -254,24 +301,26 @@ fn addr_to_ipv4_string(address: &str) -> Result { if addr.is_ipv4() { Ok(addr.ip().to_string()) } else { - ERR!("Address {} resolved to IPv6 {} which is not supported", address, addr) + MmError::err(ParseAddressError::UnsupportedIPv6Address(address.into())) } }, - None => { - ERR!("Address {} to_socket_addrs empty iter", address) - }, - }, - Err(e) => { - ERR!("Couldn't resolve '{}' Address: {}", address, e) + None => MmError::err(ParseAddressError::EmptyIterator(address.into())), }, + Err(e) => MmError::err(ParseAddressError::UnresolvedAddress(address.into(), e.to_string())), } } -pub fn lp_ports(netid: u16) -> Result<(u16, u16, u16), String> { +#[derive(Debug, Display)] +pub enum NetIdError { + #[display(fmt = "Netid {} is larger than max {}", netid, max_netid)] + LargerThanMax { netid: u16, max_netid: u16 }, +} + +pub fn lp_ports(netid: u16) -> Result<(u16, u16, u16), MmError> { const LP_RPCPORT: u16 = 7783; let max_netid = (65535 - 40 - LP_RPCPORT) / 4; if netid > max_netid { - return ERR!("Netid {} is larger than max {}", netid, max_netid); + return MmError::err(NetIdError::LargerThanMax { netid, max_netid }); } let other_ports = if netid != 0 { diff --git a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs index d332e6912a..5dfe0cb8f3 100644 --- a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs +++ b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs @@ -21,8 +21,7 @@ use libp2p_floodsub::{Floodsub, FloodsubEvent, Topic as FloodsubTopic}; use log::{debug, error, info}; use rand::seq::SliceRandom; use rand::Rng; -use std::{collections::{hash_map::{DefaultHasher, HashMap}, - HashSet}, +use std::{collections::hash_map::{DefaultHasher, HashMap}, hash::{Hash, Hasher}, iter, net::IpAddr, @@ -118,12 +117,6 @@ pub enum AdexBehaviourCmd { req: Vec, response_tx: oneshot::Sender>, }, - /// Add addresses to peer exchange then request peers and collect all their responses. - RequestAddresses { - req: Vec, - peers_addresses: Vec<(String, String)>, - response_tx: oneshot::Sender>, - }, /// Send a response using a `response_channel`. SendResponse { /// Response to a request. @@ -146,6 +139,11 @@ pub enum AdexBehaviourCmd { GetRelayMesh { result_tx: oneshot::Sender>, }, + /// Add addresses for a peer to peer exchange. + AddPeerAddresses { + peer: PeerId, + addresses: PeerAddresses, + }, PropagateMessage { message_id: MessageId, propagation_source: PeerId, @@ -293,36 +291,6 @@ impl AtomicDexBehaviour { let future = request_peers(relays, req, self.request_response.sender(), response_tx); self.spawn(future); }, - AdexBehaviourCmd::RequestAddresses { - req, - peers_addresses, - response_tx, - } => { - let peers = peers_addresses - .into_iter() - .filter_map(|(peer, address)| match peer.parse() { - Ok(p) => { - let multi_address = match Multiaddr::from_str(&address) { - Ok(addr) => addr, - Err(e) => { - error!("Error on parse address {}: {:?}", address, e); - return None; - }, - }; - let mut addresses = HashSet::new(); - addresses.insert(multi_address); - self.peers_exchange.add_peer_addresses(&p, addresses); - Some(p) - }, - Err(e) => { - error!("Error on parse peer id {:?}: {:?}", peer, e); - None - }, - }) - .collect(); - let future = request_peers(peers, req, self.request_response.sender(), response_tx); - self.spawn(future); - }, AdexBehaviourCmd::SendResponse { res, response_channel } => { if let Err(response) = self.request_response.send_response(response_channel.into(), res.into()) { error!("Error sending response: {:?}", response); @@ -405,6 +373,9 @@ impl AtomicDexBehaviour { error!("Result rx is dropped"); } }, + AdexBehaviourCmd::AddPeerAddresses { peer, addresses } => { + self.peers_exchange.add_peer_addresses(&peer, addresses); + }, AdexBehaviourCmd::PropagateMessage { message_id, propagation_source, diff --git a/mm2src/mm2_libp2p/src/lib.rs b/mm2src/mm2_libp2p/src/lib.rs index dd312ea865..674ddde210 100644 --- a/mm2src/mm2_libp2p/src/lib.rs +++ b/mm2src/mm2_libp2p/src/lib.rs @@ -2,7 +2,7 @@ mod adex_ping; pub mod atomicdex_behaviour; -mod peers_exchange; +pub mod peers_exchange; pub mod request_response; mod runtime; @@ -15,6 +15,7 @@ use sha2::{Digest, Sha256}; pub use atomicdex_behaviour::{spawn_gossipsub, NodeType}; pub use atomicdex_gossipsub::{GossipsubEvent, GossipsubMessage, MessageId}; pub use libp2p::PeerId; +pub use peers_exchange::PeerAddresses; lazy_static! { static ref SECP_VERIFY: Secp256k1 = Secp256k1::verification_only(); diff --git a/mm2src/rpc/dispatcher/dispatcher_legacy.rs b/mm2src/rpc/dispatcher/dispatcher_legacy.rs index 2ff5f9f7e4..33f6348fda 100644 --- a/mm2src/rpc/dispatcher/dispatcher_legacy.rs +++ b/mm2src/rpc/dispatcher/dispatcher_legacy.rs @@ -12,7 +12,6 @@ use super::lp_commands::*; use crate::mm2::lp_ordermatch::{best_orders_rpc, buy, cancel_all_orders, cancel_order, my_orders, order_status, orderbook_depth_rpc, orderbook_rpc, orders_history_by_filter, sell, set_price, update_maker_order}; -use crate::mm2::lp_stats::{add_node_to_version_stat, remove_node_from_version_stat, start_version_stat_collection}; use crate::mm2::lp_swap::{active_swaps_rpc, all_swaps_uuids_by_filter, ban_pubkey_rpc, coins_needed_for_kick_start, import_swaps, list_banned_pubkeys_rpc, max_taker_vol, my_recent_swaps, my_swap_status, recover_funds_of_swap, stats_swap_status, unban_pubkeys_rpc}; @@ -61,7 +60,6 @@ pub fn dispatcher(req: Json, ctx: MmArc) -> DispatcherRes { // Sorted alphanumerically (on the first latter) for readability. // "autoprice" => lp_autoprice (ctx, req), "active_swaps" => hyres(active_swaps_rpc(ctx, req)), - "add_node_to_version_stat" => hyres(add_node_to_version_stat(ctx, req)), "all_swaps_uuids_by_filter" => all_swaps_uuids_by_filter(ctx, req), "ban_pubkey" => hyres(ban_pubkey_rpc(ctx, req)), "best_orders" => hyres(best_orders_rpc(ctx, req)), @@ -120,14 +118,12 @@ pub fn dispatcher(req: Json, ctx: MmArc) -> DispatcherRes { return DispatcherRes::NoMatch(req); } }, - "remove_node_from_version_stat" => hyres(remove_node_from_version_stat(ctx, req)), "sell" => hyres(sell(ctx, req)), "show_priv_key" => hyres(show_priv_key(ctx, req)), "send_raw_transaction" => hyres(send_raw_transaction(ctx, req)), "set_required_confirmations" => hyres(set_required_confirmations(ctx, req)), "set_requires_notarization" => hyres(set_requires_notarization(ctx, req)), "setprice" => hyres(set_price(ctx, req)), - "start_version_stat_collection" => hyres(start_version_stat_collection(ctx, req)), "stats_swap_status" => stats_swap_status(ctx, req), "stop" => stop(ctx), "trade_preimage" => hyres(into_legacy::trade_preimage(ctx, req)), diff --git a/mm2src/rpc/dispatcher/dispatcher_v2.rs b/mm2src/rpc/dispatcher/dispatcher_v2.rs index 67dff2d783..ffbc665f1a 100644 --- a/mm2src/rpc/dispatcher/dispatcher_v2.rs +++ b/mm2src/rpc/dispatcher/dispatcher_v2.rs @@ -1,5 +1,6 @@ use super::lp_protocol::{MmRpcBuilder, MmRpcRequest}; use super::{DispatcherError, DispatcherResult, PUBLIC_METHODS}; +use crate::mm2::lp_stats::{add_node_to_version_stat, remove_node_from_version_stat, start_version_stat_collection}; use crate::mm2::lp_swap::trade_preimage_rpc; use coins::withdraw; use common::log::{error, warn}; @@ -83,8 +84,11 @@ fn auth(request: &MmRpcRequest, ctx: &MmArc) -> DispatcherResult<()> { async fn dispatcher(request: MmRpcRequest, ctx: MmArc) -> DispatcherResult>> { match request.method.as_str() { - "withdraw" => handle_mmrpc(ctx, request, withdraw).await, + "add_node_to_version_stat" => handle_mmrpc(ctx, request, add_node_to_version_stat).await, + "remove_node_from_version_stat" => handle_mmrpc(ctx, request, remove_node_from_version_stat).await, + "start_version_stat_collection" => handle_mmrpc(ctx, request, start_version_stat_collection).await, "trade_preimage" => handle_mmrpc(ctx, request, trade_preimage_rpc).await, + "withdraw" => handle_mmrpc(ctx, request, withdraw).await, _ => MmError::err(DispatcherError::NoSuchMethod { method: request.method }), } } From de83a69e0cf491ddadf3cac331d895b37f016140 Mon Sep 17 00:00:00 2001 From: shamardy Date: Wed, 4 Aug 2021 22:42:03 +0200 Subject: [PATCH 4/6] second review fixes --- mm2src/database/stats_nodes.rs | 12 ++ mm2src/lp_native_dex.rs | 33 +---- mm2src/lp_network.rs | 85 +++++++++++-- mm2src/lp_stats.rs | 122 ++++++------------- mm2src/mm2_libp2p/src/atomicdex_behaviour.rs | 14 ++- mm2src/mm2_libp2p/src/peers_exchange.rs | 35 +++++- 6 files changed, 169 insertions(+), 132 deletions(-) diff --git a/mm2src/database/stats_nodes.rs b/mm2src/database/stats_nodes.rs index 0f0b28fc52..5240806aea 100644 --- a/mm2src/database/stats_nodes.rs +++ b/mm2src/database/stats_nodes.rs @@ -24,6 +24,8 @@ const INSERT_NODE: &str = "INSERT INTO nodes (name, address, peer_id) VALUES (?1 const DELETE_NODE: &str = "DELETE FROM nodes WHERE name = ?1"; +const SELECT_PEERS_ADDRESSES: &str = "SELECT peer_id, address FROM nodes"; + const SELECT_PEERS_NAMES: &str = "SELECT peer_id, name FROM nodes"; const INSERT_STAT: &str = "INSERT INTO stats_nodes (name, version, timestamp, error) VALUES (?1, ?2, ?3, ?4)"; @@ -46,6 +48,16 @@ pub fn delete_node_info(ctx: &MmArc, name: String) -> SqlResult<()> { conn.execute(DELETE_NODE, ¶ms).map(|_| ()) } +pub fn select_peers_addresses(ctx: &MmArc) -> SqlResult, SqlError> { + let conn = ctx.sqlite_connection(); + let mut stmt = conn.prepare(SELECT_PEERS_ADDRESSES)?; + let peers_addresses = stmt + .query_map(NO_PARAMS, |row| Ok((row.get(0)?, row.get(1)?)))? + .collect::>>()?; + + Ok(peers_addresses) +} + pub fn select_peers_names(ctx: &MmArc) -> SqlResult, SqlError> { let conn = ctx.sqlite_connection(); let mut stmt = conn.prepare(SELECT_PEERS_NAMES)?; diff --git a/mm2src/lp_native_dex.rs b/mm2src/lp_native_dex.rs index b6f77f6b89..86cdae1666 100644 --- a/mm2src/lp_native_dex.rs +++ b/mm2src/lp_native_dex.rs @@ -30,10 +30,9 @@ use std::str; #[cfg(not(target_arch = "wasm32"))] use crate::mm2::database::init_and_migrate_db; -use crate::mm2::lp_network::{p2p_event_process_loop, P2PContext}; +use crate::mm2::lp_network::{addr_to_ipv4_string, lp_ports, p2p_event_process_loop, P2PContext}; use crate::mm2::lp_ordermatch::{broadcast_maker_orders_keep_alive_loop, lp_ordermatch_loop, orders_kick_start, BalanceUpdateOrdermatchHandler}; -use crate::mm2::lp_stats::lp_ports; use crate::mm2::lp_swap::{running_swaps_num, swap_kick_starts}; use crate::mm2::rpc::spawn_rpc; use crate::mm2::{MM_DATETIME, MM_VERSION}; @@ -64,10 +63,7 @@ fn default_seednodes(netid: u16) -> Vec { if netid == 7777 { NETID_7777_SEEDNODES .iter() - .filter_map(|seed| { - let seed_url = format!("{}:0", *seed); - seed_to_ipv4_string(&seed_url) - }) + .filter_map(|seed| addr_to_ipv4_string(*seed).ok()) .collect() } else { Vec::new() @@ -313,31 +309,6 @@ fn test_ip(ctx: &MmArc, ip: IpAddr) -> Result<(), String> { } } -#[cfg(not(target_arch = "wasm32"))] -fn seed_to_ipv4_string(seed: &str) -> Option { - use std::net::ToSocketAddrs; - match seed.to_socket_addrs() { - Ok(mut iter) => match iter.next() { - Some(addr) => { - if addr.is_ipv4() { - Some(addr.ip().to_string()) - } else { - warn!("Seed {} resolved to IPv6 {} which is not supported", seed, addr); - None - } - }, - None => { - warn!("Seed {} to_socket_addrs empty iter", seed); - None - }, - }, - Err(e) => { - error!("Couldn't resolve '{}' seed: {}", seed, e); - None - }, - } -} - #[cfg_attr(target_arch = "wasm32", allow(unused_variables))] /// * `ctx_cb` - callback used to share the `MmCtx` ID with the call site. pub async fn lp_init(ctx: MmArc) -> Result<(), String> { diff --git a/mm2src/lp_network.rs b/mm2src/lp_network.rs index f1f474d3c6..67e60b3432 100644 --- a/mm2src/lp_network.rs +++ b/mm2src/lp_network.rs @@ -29,6 +29,7 @@ use mm2_libp2p::peers_exchange::PeerAddresses; use mm2_libp2p::{decode_message, encode_message, GossipsubMessage, MessageId, PeerId, TOPIC_SEPARATOR}; #[cfg(test)] use mocktopus::macros::*; use serde::de; +use std::net::ToSocketAddrs; use std::sync::Arc; use crate::mm2::{lp_ordermatch, lp_stats, lp_swap}; @@ -162,7 +163,7 @@ async fn process_p2p_request( request: Vec, response_channel: AdexResponseChannel, ) -> P2PRequestResult<()> { - let request = decode_message::(&request).map_to_mm(P2PRequestError::from)?; + let request = decode_message::(&request)?; let result = match request { P2PRequest::Ordermatch(req) => lp_ordermatch::process_peer_request(ctx.clone(), req).await, P2PRequest::NetworkInfo(req) => lp_stats::process_info_request(ctx.clone(), req).await, @@ -213,7 +214,7 @@ pub async fn request_any_relay( ctx: MmArc, req: P2PRequest, ) -> P2PRequestResult> { - let encoded = encode_message(&req).map_to_mm(P2PRequestError::from)?; + let encoded = encode_message(&req)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -232,7 +233,7 @@ pub async fn request_any_relay( .map_to_mm(|e| P2PRequestError::ResponseError(e.to_string()))? { Some((from_peer, response)) => { - let response = decode_message::(&response).map_to_mm(P2PRequestError::from)?; + let response = decode_message::(&response)?; Ok(Some((response, from_peer))) }, None => Ok(None), @@ -250,7 +251,7 @@ pub async fn request_relays( ctx: MmArc, req: P2PRequest, ) -> P2PRequestResult)>> { - let encoded = encode_message(&req).map_to_mm(P2PRequestError::from)?; + let encoded = encode_message(&req)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -275,7 +276,7 @@ pub async fn request_peers( req: P2PRequest, peers: Vec, ) -> P2PRequestResult)>> { - let encoded = encode_message(&req).map_to_mm(P2PRequestError::from)?; + let encoded = encode_message(&req)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -350,13 +351,81 @@ pub fn propagate_message(ctx: &MmArc, message_id: MessageId, propagation_source: }); } -pub fn add_peer_addresses(ctx: &MmArc, peer: PeerId, addresses: PeerAddresses) { +pub fn add_reserved_peer_addresses(ctx: &MmArc, peer: PeerId, addresses: PeerAddresses) { let ctx = ctx.clone(); spawn(async move { let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); - let cmd = AdexBehaviourCmd::AddPeerAddresses { peer, addresses }; + let cmd = AdexBehaviourCmd::AddReservedPeer { peer, addresses }; if let Err(e) = p2p_ctx.cmd_tx.lock().await.try_send(cmd) { - log::error!("add_peer_addresses cmd_tx.send error {:?}", e); + log::error!("add_reserved_peer_addresses cmd_tx.send error {:?}", e); }; }); } + +#[derive(Debug, Display)] +pub enum ParseAddressError { + #[display(fmt = "Address/Seed {} resolved to IPv6 which is not supported", _0)] + UnsupportedIPv6Address(String), + #[display(fmt = "Address/Seed {} to_socket_addrs empty iter", _0)] + EmptyIterator(String), + #[display(fmt = "Couldn't resolve '{}' Address/Seed: {}", _0, _1)] + UnresolvedAddress(String, String), +} + +#[cfg(not(target_arch = "wasm32"))] +pub fn addr_to_ipv4_string(address: &str) -> Result> { + // Remove "https:// or http://" etc.. from address str + let formated_address = address.split("://").last().unwrap_or(address); + let address_with_port = if formated_address.contains(':') { + formated_address.to_string() + } else { + format!("{}:0", formated_address) + }; + match address_with_port.as_str().to_socket_addrs() { + Ok(mut iter) => match iter.next() { + Some(addr) => { + if addr.is_ipv4() { + Ok(addr.ip().to_string()) + } else { + log::warn!( + "Address/Seed {} resolved to IPv6 {} which is not supported", + address, + addr + ); + MmError::err(ParseAddressError::UnsupportedIPv6Address(address.into())) + } + }, + None => { + log::warn!("Address/Seed {} to_socket_addrs empty iter", address); + MmError::err(ParseAddressError::EmptyIterator(address.into())) + }, + }, + Err(e) => { + log::error!("Couldn't resolve '{}' seed: {}", address, e); + MmError::err(ParseAddressError::UnresolvedAddress(address.into(), e.to_string())) + }, + } +} + +#[derive(Debug, Display)] +pub enum NetIdError { + #[display(fmt = "Netid {} is larger than max {}", netid, max_netid)] + LargerThanMax { netid: u16, max_netid: u16 }, +} + +pub fn lp_ports(netid: u16) -> Result<(u16, u16, u16), MmError> { + const LP_RPCPORT: u16 = 7783; + let max_netid = (65535 - 40 - LP_RPCPORT) / 4; + if netid > max_netid { + return MmError::err(NetIdError::LargerThanMax { netid, max_netid }); + } + + let other_ports = if netid != 0 { + let net_mod = netid % 10; + let net_div = netid / 10; + (net_div * 40) + LP_RPCPORT + net_mod + } else { + LP_RPCPORT + }; + Ok((other_ports + 10, other_ports + 20, other_ports + 30)) +} diff --git a/mm2src/lp_stats.rs b/mm2src/lp_stats.rs index 01d83dffa9..bff8760b06 100644 --- a/mm2src/lp_stats.rs +++ b/mm2src/lp_stats.rs @@ -10,9 +10,9 @@ use mm2_libp2p::atomicdex_behaviour::parse_relay_address; use mm2_libp2p::{encode_message, PeerId}; use serde_json::{self as json, Value as Json}; use std::collections::{HashMap, HashSet}; -use std::net::ToSocketAddrs; -use crate::mm2::lp_network::{add_peer_addresses, request_peers, P2PRequest, PeerDecodedResponse}; +use crate::mm2::lp_network::{add_reserved_peer_addresses, addr_to_ipv4_string, lp_ports, request_peers, NetIdError, + P2PRequest, ParseAddressError, PeerDecodedResponse}; pub type NodeVersionResult = Result>; @@ -75,7 +75,7 @@ fn insert_node_info_to_db(_ctx: &MmArc, _node_info: &NodeInfo) -> Result<(), Str #[cfg(not(target_arch = "wasm32"))] fn insert_node_info_to_db(ctx: &MmArc, node_info: &NodeInfo) -> Result<(), String> { - crate::mm2::database::stats_nodes::insert_node_info(ctx, node_info).map_err(|e| ERRL!("{}", e)) + crate::mm2::database::stats_nodes::insert_node_info(ctx, node_info).map_err(|e| e.to_string()) } #[cfg(target_arch = "wasm32")] @@ -83,7 +83,7 @@ fn insert_node_version_stat_to_db(_ctx: &MmArc, _node_version_stat: NodeVersionS #[cfg(not(target_arch = "wasm32"))] fn insert_node_version_stat_to_db(ctx: &MmArc, node_version_stat: NodeVersionStat) -> Result<(), String> { - crate::mm2::database::stats_nodes::insert_node_version_stat(ctx, node_version_stat).map_err(|e| ERRL!("{}", e)) + crate::mm2::database::stats_nodes::insert_node_version_stat(ctx, node_version_stat).map_err(|e| e.to_string()) } #[cfg(target_arch = "wasm32")] @@ -91,7 +91,15 @@ fn delete_node_info_from_db(_ctx: &MmArc, _name: String) -> Result<(), String> { #[cfg(not(target_arch = "wasm32"))] fn delete_node_info_from_db(ctx: &MmArc, name: String) -> Result<(), String> { - crate::mm2::database::stats_nodes::delete_node_info(ctx, name).map_err(|e| ERRL!("{}", e)) + crate::mm2::database::stats_nodes::delete_node_info(ctx, name).map_err(|e| e.to_string()) +} + +#[cfg(target_arch = "wasm32")] +fn select_peers_addresses_from_db(_ctx: &MmArc) -> Result, String> { Ok(()) } + +#[cfg(not(target_arch = "wasm32"))] +fn select_peers_addresses_from_db(ctx: &MmArc) -> Result, String> { + crate::mm2::database::stats_nodes::select_peers_addresses(ctx).map_err(|e| e.to_string()) } #[cfg(target_arch = "wasm32")] @@ -102,31 +110,16 @@ pub async fn add_node_to_version_stat(_ctx: MmArc, _req: Json) -> NodeVersionRes /// Adds node info. to db to be used later for stats collection #[cfg(not(target_arch = "wasm32"))] pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> NodeVersionResult { - let node_info: NodeInfo = json::from_value(req).map_to_mm(NodeVersionError::from)?; - let netid = ctx.conf["netid"].as_u64().unwrap_or(0) as u16; - let (_, pubport, _) = lp_ports(netid)?; - let addr = addr_to_ipv4_string(&node_info.address)?; - let relay_address = parse_relay_address(addr, pubport); + let node_info: NodeInfo = json::from_value(req)?; - let mut addresses = HashSet::new(); - addresses.insert(relay_address.clone()); - - let peer_id: PeerId = match node_info.peer_id.parse() { - Ok(p) => p, - Err(e) => return MmError::err(NodeVersionError::PeerIdParseError(e.to_string())), - }; - - add_peer_addresses(&ctx, peer_id, addresses); - - let node_info_with_formated_addr = NodeInfo { + let ipv4_addr = addr_to_ipv4_string(&node_info.address)?; + let node_info_with_ipv4_addr = NodeInfo { name: node_info.name, - address: relay_address.to_string(), + address: ipv4_addr, peer_id: node_info.peer_id, }; - if let Err(e) = insert_node_info_to_db(&ctx, &node_info_with_formated_addr) { - return MmError::err(NodeVersionError::DatabaseError(e)); - } + insert_node_info_to_db(&ctx, &node_info_with_ipv4_addr).map_to_mm(NodeVersionError::DatabaseError)?; Ok("success".into()) } @@ -141,10 +134,9 @@ pub async fn remove_node_from_version_stat(_ctx: MmArc, _req: Json) -> NodeVersi /// Removes node info. from db to skip collecting stats for this node #[cfg(not(target_arch = "wasm32"))] pub async fn remove_node_from_version_stat(ctx: MmArc, req: Json) -> NodeVersionResult { - let node_name: String = json::from_value(req["name"].clone()).map_to_mm(NodeVersionError::from)?; - if let Err(e) = delete_node_info_from_db(&ctx, node_name) { - return MmError::err(NodeVersionError::DatabaseError(e)); - } + let node_name: String = json::from_value(req["name"].clone())?; + + delete_node_info_from_db(&ctx, node_name).map_to_mm(NodeVersionError::DatabaseError)?; Ok("success".into()) } @@ -182,7 +174,22 @@ pub async fn start_version_stat_collection(_ctx: MmArc, _req: Json) -> NodeVersi #[cfg(not(target_arch = "wasm32"))] pub async fn start_version_stat_collection(ctx: MmArc, req: Json) -> NodeVersionResult { - let interval: f64 = json::from_value(req["interval"].clone()).map_to_mm(NodeVersionError::from)?; + let interval: f64 = json::from_value(req["interval"].clone())?; + + let peers_addresses = select_peers_addresses_from_db(&ctx).map_to_mm(NodeVersionError::DatabaseError)?; + + let netid = ctx.conf["netid"].as_u64().unwrap_or(0) as u16; + let (_, pubport, _) = lp_ports(netid)?; + + for (peer_id, address) in peers_addresses { + let peer_id = peer_id + .parse::() + .map_to_mm(|e| NodeVersionError::PeerIdParseError(e.to_string()))?; + let mut addresses = HashSet::new(); + let multi_address = parse_relay_address(address, pubport); + addresses.insert(multi_address); + add_reserved_peer_addresses(&ctx, peer_id, addresses); + } spawn(stat_collection_loop(ctx, interval)); @@ -277,58 +284,3 @@ async fn stat_collection_loop(ctx: MmArc, interval: f64) { Timer::sleep(interval).await; } } - -#[derive(Debug, Display)] -pub enum ParseAddressError { - #[display(fmt = "Address {} resolved to IPv6 which is not supported", _0)] - UnsupportedIPv6Address(String), - #[display(fmt = "Address {} to_socket_addrs empty iter", _0)] - EmptyIterator(String), - // error return for second string - #[display(fmt = "Couldn't resolve '{}' Address: {}", _0, _1)] - UnresolvedAddress(String, String), -} - -fn addr_to_ipv4_string(address: &str) -> Result> { - let address_with_port = if address.contains(':') { - address.to_string() - } else { - format!("{}{}", address, ":0") - }; - match address_with_port.as_str().to_socket_addrs() { - Ok(mut iter) => match iter.next() { - Some(addr) => { - if addr.is_ipv4() { - Ok(addr.ip().to_string()) - } else { - MmError::err(ParseAddressError::UnsupportedIPv6Address(address.into())) - } - }, - None => MmError::err(ParseAddressError::EmptyIterator(address.into())), - }, - Err(e) => MmError::err(ParseAddressError::UnresolvedAddress(address.into(), e.to_string())), - } -} - -#[derive(Debug, Display)] -pub enum NetIdError { - #[display(fmt = "Netid {} is larger than max {}", netid, max_netid)] - LargerThanMax { netid: u16, max_netid: u16 }, -} - -pub fn lp_ports(netid: u16) -> Result<(u16, u16, u16), MmError> { - const LP_RPCPORT: u16 = 7783; - let max_netid = (65535 - 40 - LP_RPCPORT) / 4; - if netid > max_netid { - return MmError::err(NetIdError::LargerThanMax { netid, max_netid }); - } - - let other_ports = if netid != 0 { - let net_mod = netid % 10; - let net_div = netid / 10; - (net_div * 40) + LP_RPCPORT + net_mod - } else { - LP_RPCPORT - }; - Ok((other_ports + 10, other_ports + 20, other_ports + 30)) -} diff --git a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs index 5dfe0cb8f3..be33e1ff45 100644 --- a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs +++ b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs @@ -139,8 +139,8 @@ pub enum AdexBehaviourCmd { GetRelayMesh { result_tx: oneshot::Sender>, }, - /// Add addresses for a peer to peer exchange. - AddPeerAddresses { + /// Add a reserved peer to the peer exchange. + AddReservedPeer { peer: PeerId, addresses: PeerAddresses, }, @@ -373,8 +373,9 @@ impl AtomicDexBehaviour { error!("Result rx is dropped"); } }, - AdexBehaviourCmd::AddPeerAddresses { peer, addresses } => { - self.peers_exchange.add_peer_addresses(&peer, addresses); + AdexBehaviourCmd::AddReservedPeer { peer, addresses } => { + self.peers_exchange + .add_peer_addresses_to_reserved_peers(&peer, addresses); }, AdexBehaviourCmd::PropagateMessage { message_id, @@ -414,7 +415,8 @@ impl NetworkBehaviourEventProcess for AtomicDexBehaviour { Ok(a) => a, Err(_) => return, }; - self.peers_exchange.add_peer_addresses(&message.source, addresses); + self.peers_exchange + .add_peer_addresses_to_known_peers(&message.source, addresses); } } } @@ -729,7 +731,7 @@ fn start_gossipsub( for (peer_id, address) in ALL_NETID_7777_SEEDNODES { let peer_id = PeerId::from_str(peer_id).expect("valid peer id"); let multiaddr = parse_relay_address((*address).to_owned(), network_port); - peers_exchange.add_peer_addresses(&peer_id, iter::once(multiaddr).collect()); + peers_exchange.add_peer_addresses_to_known_peers(&peer_id, iter::once(multiaddr).collect()); gossipsub.add_explicit_relay(peer_id); } } diff --git a/mm2src/mm2_libp2p/src/peers_exchange.rs b/mm2src/mm2_libp2p/src/peers_exchange.rs index 4c398a065e..e86c58cf05 100644 --- a/mm2src/mm2_libp2p/src/peers_exchange.rs +++ b/mm2src/mm2_libp2p/src/peers_exchange.rs @@ -77,6 +77,8 @@ pub struct PeersExchange { #[behaviour(ignore)] known_peers: Vec, #[behaviour(ignore)] + reserved_peers: Vec, + #[behaviour(ignore)] events: VecDeque, ()>>, #[behaviour(ignore)] maintain_peers_interval: Interval, @@ -94,6 +96,7 @@ impl PeersExchange { PeersExchange { request_response, known_peers: Vec::new(), + reserved_peers: Vec::new(), events: VecDeque::new(), maintain_peers_interval: Interval::new_at( Instant::now() + Duration::from_secs(REQUEST_PEERS_INITIAL_DELAY), @@ -125,7 +128,7 @@ impl PeersExchange { } } - pub fn add_peer_addresses(&mut self, peer: &PeerId, addresses: PeerAddresses) { + pub fn add_peer_addresses_to_known_peers(&mut self, peer: &PeerId, addresses: PeerAddresses) { if addresses.len() > 1 { return; } @@ -146,6 +149,29 @@ impl PeersExchange { } } + pub fn add_peer_addresses_to_reserved_peers(&mut self, peer: &PeerId, addresses: PeerAddresses) { + if addresses.len() > 1 { + return; + } + + for address in addresses.iter() { + if !self.validate_global_multiaddr(address) { + return; + } + } + + if !self.reserved_peers.contains(&peer) && !addresses.is_empty() { + self.reserved_peers.push(*peer); + } + + let already_reserved = self.request_response.addresses_of_peer(peer); + for address in addresses { + if !already_reserved.contains(&address) { + self.request_response.add_address(&peer, address); + } + } + } + fn maintain_known_peers(&mut self) { if self.known_peers.len() > MAX_PEERS { let mut rng = rand::thread_rng(); @@ -185,6 +211,8 @@ impl PeersExchange { pub fn is_known_peer(&self, peer: &PeerId) -> bool { self.known_peers.contains(peer) } + pub fn is_reserved_peer(&self, peer: &PeerId) -> bool { self.reserved_peers.contains(peer) } + pub fn add_known_peer(&mut self, peer: PeerId) { if !self.is_known_peer(&peer) { self.known_peers.push(peer) @@ -282,7 +310,10 @@ impl NetworkBehaviourEventProcess Date: Wed, 4 Aug 2021 23:10:31 +0200 Subject: [PATCH 5/6] fix wasm errors --- mm2src/lp_native_dex.rs | 3 ++- mm2src/lp_stats.rs | 8 +++++--- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/mm2src/lp_native_dex.rs b/mm2src/lp_native_dex.rs index 86cdae1666..50ff199421 100644 --- a/mm2src/lp_native_dex.rs +++ b/mm2src/lp_native_dex.rs @@ -30,7 +30,7 @@ use std::str; #[cfg(not(target_arch = "wasm32"))] use crate::mm2::database::init_and_migrate_db; -use crate::mm2::lp_network::{addr_to_ipv4_string, lp_ports, p2p_event_process_loop, P2PContext}; +use crate::mm2::lp_network::{lp_ports, p2p_event_process_loop, P2PContext}; use crate::mm2::lp_ordermatch::{broadcast_maker_orders_keep_alive_loop, lp_ordermatch_loop, orders_kick_start, BalanceUpdateOrdermatchHandler}; use crate::mm2::lp_swap::{running_swaps_num, swap_kick_starts}; @@ -60,6 +60,7 @@ fn default_seednodes(netid: u16) -> Vec { #[cfg(not(target_arch = "wasm32"))] fn default_seednodes(netid: u16) -> Vec { + use crate::mm2::lp_network::addr_to_ipv4_string; if netid == 7777 { NETID_7777_SEEDNODES .iter() diff --git a/mm2src/lp_stats.rs b/mm2src/lp_stats.rs index bff8760b06..c365bfd1ff 100644 --- a/mm2src/lp_stats.rs +++ b/mm2src/lp_stats.rs @@ -11,8 +11,8 @@ use mm2_libp2p::{encode_message, PeerId}; use serde_json::{self as json, Value as Json}; use std::collections::{HashMap, HashSet}; -use crate::mm2::lp_network::{add_reserved_peer_addresses, addr_to_ipv4_string, lp_ports, request_peers, NetIdError, - P2PRequest, ParseAddressError, PeerDecodedResponse}; +use crate::mm2::lp_network::{add_reserved_peer_addresses, lp_ports, request_peers, NetIdError, P2PRequest, + ParseAddressError, PeerDecodedResponse}; pub type NodeVersionResult = Result>; @@ -95,7 +95,7 @@ fn delete_node_info_from_db(ctx: &MmArc, name: String) -> Result<(), String> { } #[cfg(target_arch = "wasm32")] -fn select_peers_addresses_from_db(_ctx: &MmArc) -> Result, String> { Ok(()) } +fn select_peers_addresses_from_db(_ctx: &MmArc) -> Result, String> { Ok(Vec::new()) } #[cfg(not(target_arch = "wasm32"))] fn select_peers_addresses_from_db(ctx: &MmArc) -> Result, String> { @@ -110,6 +110,8 @@ pub async fn add_node_to_version_stat(_ctx: MmArc, _req: Json) -> NodeVersionRes /// Adds node info. to db to be used later for stats collection #[cfg(not(target_arch = "wasm32"))] pub async fn add_node_to_version_stat(ctx: MmArc, req: Json) -> NodeVersionResult { + use crate::mm2::lp_network::addr_to_ipv4_string; + let node_info: NodeInfo = json::from_value(req)?; let ipv4_addr = addr_to_ipv4_string(&node_info.address)?; From 86114c310397d2feb6af32767f91978a5c8e51e1 Mon Sep 17 00:00:00 2001 From: shamardy Date: Thu, 5 Aug 2021 13:19:46 +0200 Subject: [PATCH 6/6] avoid reserved peers disconnect in maintain relays --- mm2src/mm2_libp2p/src/atomicdex_behaviour.rs | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs index be33e1ff45..e0dab61876 100644 --- a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs +++ b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs @@ -511,9 +511,11 @@ fn maintain_connection_to_relays(swarm: &mut AtomicDexSwarm, bootstrap_addresses .filter(|peer| !relays_mesh.contains(peer)) .collect(); for peer in not_in_mesh.choose_multiple(&mut rng, to_disconnect_num) { - info!("Disconnecting peer {}", peer); - if Swarm::disconnect_peer_id(swarm, **peer).is_err() { - error!("Peer {} disconnect error", peer); + if !swarm.behaviour().peers_exchange.is_reserved_peer(*peer) { + info!("Disconnecting peer {}", peer); + if Swarm::disconnect_peer_id(swarm, **peer).is_err() { + error!("Peer {} disconnect error", peer); + } } } }