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..5240806aea --- /dev/null +++ b/mm2src/database/stats_nodes.rs @@ -0,0 +1,84 @@ +/// 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), + 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, 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); + 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, + 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_native_dex.rs b/mm2src/lp_native_dex.rs index 5f070a0cf2..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::{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,36 +60,17 @@ 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() - .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() } } -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 { @@ -329,31 +310,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 c80dc00a0a..67e60b3432 100644 --- a/mm2src/lp_network.rs +++ b/mm2src/lp_network.rs @@ -19,20 +19,45 @@ use common::executor::spawn; use common::log; use common::mm_ctx::{MmArc, MmWeak}; +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; +use std::net::ToSocketAddrs; use std::sync::Arc; -use crate::mm2::{lp_ordermatch, lp_swap}; +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), + NetworkInfo(lp_stats::NetworkInfoRequest), } pub struct P2PContext { @@ -137,10 +162,11 @@ 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)?; 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 { @@ -151,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(()) } @@ -182,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)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -191,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)?; Ok(Some((response, from_peer))) }, None => Ok(None), @@ -211,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)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -220,8 +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); + 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)) } @@ -229,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)?; let (response_tx, response_rx) = oneshot::channel(); let p2p_ctx = P2PContext::fetch_from_mm_arc(&ctx); @@ -239,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)) } @@ -248,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)), } } @@ -297,3 +350,82 @@ pub fn propagate_message(ctx: &MmArc, message_id: MessageId, propagation_source: }; }); } + +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::AddReservedPeer { peer, addresses }; + if let Err(e) = p2p_ctx.cmd_tx.lock().await.try_send(cmd) { + 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_ordermatch.rs b/mm2src/lp_ordermatch.rs index 19aaf22bf0..fc7bf66849 100644 --- a/mm2src/lp_ordermatch.rs +++ b/mm2src/lp_ordermatch.rs @@ -3808,7 +3808,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 new file mode 100644 index 0000000000..c365bfd1ff --- /dev/null +++ b/mm2src/lp_stats.rs @@ -0,0 +1,288 @@ +/// The module is responsible for mm2 network stats collection +/// +use common::executor::{spawn, Timer}; +use common::mm_ctx::MmArc; +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, PeerId}; +use serde_json::{self as json, Value as Json}; +use std::collections::{HashMap, HashSet}; + +use crate::mm2::lp_network::{add_reserved_peer_addresses, lp_ports, request_peers, NetIdError, P2PRequest, + ParseAddressError, 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 { + pub name: String, + pub address: String, + pub peer_id: String, +} + +#[derive(Serialize, Deserialize)] +pub struct NodeVersionStat { + pub name: String, + pub version: Option, + pub timestamp: u64, + pub error: Option, +} + +#[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| e.to_string()) +} + +#[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| e.to_string()) +} + +#[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| e.to_string()) +} + +#[cfg(target_arch = "wasm32")] +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> { + crate::mm2::database::stats_nodes::select_peers_addresses(ctx).map_err(|e| e.to_string()) +} + +#[cfg(target_arch = "wasm32")] +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) -> 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)?; + let node_info_with_ipv4_addr = NodeInfo { + name: node_info.name, + address: ipv4_addr, + peer_id: node_info.peer_id, + }; + + insert_node_info_to_db(&ctx, &node_info_with_ipv4_addr).map_to_mm(NodeVersionError::DatabaseError)?; + + Ok("success".into()) +} + +#[cfg(target_arch = "wasm32")] +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) -> NodeVersionResult { + 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()) +} + +#[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) -> 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) -> NodeVersionResult { + 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)); + + 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_names; + + loop { + if ctx.is_stopping() { + break; + }; + { + 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 peers: Vec = peers_names.keys().cloned().collect(); + + let timestamp = now_ms() / 1000; + let get_versions_res = match request_peers::( + ctx.clone(), + P2PRequest::NetworkInfo(NetworkInfoRequest::GetMm2Version), + peers, + ) + .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(v) => { + let node_version_stat = NodeVersionStat { + 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 node {} version {} into db: {}", name, v, e); + }; + }, + PeerDecodedResponse::Err(e) => { + 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.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); + }; + }, + } + } + } + Timer::sleep(interval).await; + } +} diff --git a/mm2src/mm2.rs b/mm2src/mm2.rs index 001256b2a9..3f743d3d71 100644 --- a/mm2src/mm2.rs +++ b/mm2src/mm2.rs @@ -49,6 +49,7 @@ pub mod database; #[path = "lp_network.rs"] pub mod lp_network; #[path = "lp_ordermatch.rs"] pub mod lp_ordermatch; +#[path = "lp_stats.rs"] pub mod lp_stats; #[path = "lp_swap.rs"] pub mod lp_swap; #[path = "rpc.rs"] pub mod rpc; diff --git a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs index b0da3b71ea..e0dab61876 100644 --- a/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs +++ b/mm2src/mm2_libp2p/src/atomicdex_behaviour.rs @@ -139,6 +139,11 @@ pub enum AdexBehaviourCmd { GetRelayMesh { result_tx: oneshot::Sender>, }, + /// Add a reserved peer to the peer exchange. + AddReservedPeer { + peer: PeerId, + addresses: PeerAddresses, + }, PropagateMessage { message_id: MessageId, propagation_source: PeerId, @@ -368,6 +373,10 @@ impl AtomicDexBehaviour { error!("Result rx is dropped"); } }, + AdexBehaviourCmd::AddReservedPeer { peer, addresses } => { + self.peers_exchange + .add_peer_addresses_to_reserved_peers(&peer, addresses); + }, AdexBehaviourCmd::PropagateMessage { message_id, propagation_source, @@ -406,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); } } } @@ -501,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); + } } } } @@ -721,7 +733,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); } } @@ -836,7 +848,7 @@ fn generate_ed25519_keypair(rng: &mut R, force_key: Option<[u8; 32]>) -> /// `addr` is expected to be either `/dns//tcp/` or `/ipv4//tcp/` or an IPv4 address. #[cfg(target_arch = "wasm32")] -fn parse_relay_address(addr: String, port: u16) -> Multiaddr { +pub fn parse_relay_address(addr: String, port: u16) -> Multiaddr { let dns = addr.starts_with("/dns") && addr.contains("/tcp/") && addr.ends_with("/ws"); let ip4 = addr.starts_with("/ip4/") && addr.contains("/tcp/") && addr.ends_with("/ws"); if dns || ip4 { @@ -850,7 +862,7 @@ fn parse_relay_address(addr: String, port: u16) -> Multiaddr { /// If the `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 the `addr` and `port` values. #[cfg(all(test, not(target_arch = "wasm32")))] -fn parse_relay_address(addr: String, port: u16) -> Multiaddr { +pub fn parse_relay_address(addr: String, port: u16) -> Multiaddr { if addr.starts_with("/ip4/") && addr.contains("/tcp/") { return addr.parse().unwrap(); } @@ -860,7 +872,7 @@ fn parse_relay_address(addr: String, port: u16) -> Multiaddr { /// The `addr` is expected to be an IP of the relay. #[cfg(all(not(test), not(target_arch = "wasm32")))] -fn parse_relay_address(ipv4_addr: String, port: u16) -> Multiaddr { +pub fn parse_relay_address(ipv4_addr: String, port: u16) -> Multiaddr { format!("/ip4/{}/tcp/{}", ipv4_addr, port).parse().unwrap() } 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/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 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 }), } }