|
|
@@ -16,28 +16,26 @@
|
|
|
* along with this program. If not, see <https://www.gnu.org/licenses/>.
|
|
|
*/
|
|
|
|
|
|
-use std::{error, fs::File, io::stdin};
|
|
|
-
|
|
|
-// ANCHOR: daemon_deps
|
|
|
-use darkfi::system::StoppableTask;
|
|
|
use easy_parallel::Parallel;
|
|
|
-use smol::{lock::Mutex, Executor};
|
|
|
-use std::{collections::HashSet, sync::Arc};
|
|
|
-// ANCHOR_END: daemon_deps
|
|
|
-
|
|
|
-use log::{debug, error};
|
|
|
-use simplelog::WriteLogger;
|
|
|
+use log::{debug, error, info, warn};
|
|
|
+use smol::{lock::Mutex, stream::StreamExt, Executor};
|
|
|
+use std::{collections::HashSet, error, fs::File, io::stdin, sync::Arc};
|
|
|
use url::Url;
|
|
|
|
|
|
use darkfi::{
|
|
|
- net,
|
|
|
- net::Settings,
|
|
|
+ async_daemonize, cli_desc, net,
|
|
|
+ net::{settings::SettingsOpt, Settings},
|
|
|
rpc::{
|
|
|
jsonrpc::JsonSubscriber,
|
|
|
server::{listen_and_serve, RequestHandler},
|
|
|
},
|
|
|
+ system::StoppableTask,
|
|
|
+ util::path::get_config_path,
|
|
|
+ Error, Result,
|
|
|
};
|
|
|
|
|
|
+use structopt_toml::{serde::Deserialize, structopt::StructOpt, StructOptToml};
|
|
|
+
|
|
|
use crate::{
|
|
|
dchat_error::ErrorMissingSpecifier,
|
|
|
dchatmsg::{DchatMsg, DchatMsgsBuffer},
|
|
|
@@ -50,11 +48,33 @@ pub mod dchatmsg;
|
|
|
pub mod protocol_dchat;
|
|
|
pub mod rpc;
|
|
|
|
|
|
-// ANCHOR: error
|
|
|
-pub type DarkfiError = darkfi::Error;
|
|
|
-pub type Error = Box<dyn error::Error>;
|
|
|
-pub type Result<T> = std::result::Result<T, Error>;
|
|
|
-// ANCHOR_END: error
|
|
|
+const CONFIG_FILE: &str = "dchat_config.toml";
|
|
|
+const CONFIG_FILE_CONTENTS: &str = include_str!("../dchat_config.toml");
|
|
|
+
|
|
|
+#[derive(Clone, Debug, Deserialize, StructOpt, StructOptToml)]
|
|
|
+#[serde(default)]
|
|
|
+#[structopt(name = "dchat", about = cli_desc!())]
|
|
|
+struct Args {
|
|
|
+ #[structopt(long, default_value = "tcp://127.0.0.1:55054")]
|
|
|
+ /// RPC server listen address
|
|
|
+ rpc_listen: Url,
|
|
|
+
|
|
|
+ #[structopt(short, long)]
|
|
|
+ /// Configuration file to use
|
|
|
+ config: Option<String>,
|
|
|
+
|
|
|
+ #[structopt(short, long)]
|
|
|
+ /// Set log file to ouput into
|
|
|
+ log: Option<String>,
|
|
|
+
|
|
|
+ #[structopt(short, parse(from_occurrences))]
|
|
|
+ /// Increase verbosity (-vvv supported)
|
|
|
+ verbose: u8,
|
|
|
+
|
|
|
+ /// P2P network settings
|
|
|
+ #[structopt(flatten)]
|
|
|
+ net: SettingsOpt,
|
|
|
+}
|
|
|
|
|
|
// ANCHOR: dchat
|
|
|
struct Dchat {
|
|
|
@@ -68,87 +88,6 @@ impl Dchat {
|
|
|
Self { p2p, recv_msgs }
|
|
|
}
|
|
|
|
|
|
- // ANCHOR: menu
|
|
|
- async fn menu(&self) -> Result<()> {
|
|
|
- let mut buffer = String::new();
|
|
|
- let stdin = stdin();
|
|
|
- loop {
|
|
|
- println!(
|
|
|
- "Welcome to dchat.
|
|
|
- s: send message
|
|
|
- i: inbox
|
|
|
- q: quit "
|
|
|
- );
|
|
|
- stdin.read_line(&mut buffer)?;
|
|
|
- // Remove trailing \n
|
|
|
- buffer.pop();
|
|
|
- match buffer.as_str() {
|
|
|
- "q" => return Ok(()),
|
|
|
- "s" => {
|
|
|
- // Remove trailing s
|
|
|
- buffer.pop();
|
|
|
- stdin.read_line(&mut buffer)?;
|
|
|
- match self.send(buffer.clone()).await {
|
|
|
- Ok(_) => {
|
|
|
- println!("you sent: {}", buffer);
|
|
|
- }
|
|
|
- Err(e) => {
|
|
|
- println!("send failed for reason: {}", e);
|
|
|
- }
|
|
|
- }
|
|
|
- buffer.clear();
|
|
|
- }
|
|
|
- "i" => {
|
|
|
- let msgs = self.recv_msgs.lock().await;
|
|
|
- if msgs.is_empty() {
|
|
|
- println!("inbox is empty")
|
|
|
- } else {
|
|
|
- println!("received:");
|
|
|
- for i in msgs.iter() {
|
|
|
- if !i.msg.is_empty() {
|
|
|
- println!("{}", i.msg);
|
|
|
- }
|
|
|
- }
|
|
|
- }
|
|
|
- buffer.clear();
|
|
|
- }
|
|
|
- _ => {}
|
|
|
- }
|
|
|
- }
|
|
|
- }
|
|
|
- // ANCHOR_END: menu
|
|
|
-
|
|
|
- // ANCHOR: register_protocol
|
|
|
- async fn register_protocol(&self, msgs: DchatMsgsBuffer) -> Result<()> {
|
|
|
- debug!(target: "dchat", "Dchat::register_protocol() [START]");
|
|
|
- let registry = self.p2p.protocol_registry();
|
|
|
- registry
|
|
|
- .register(!net::session::SESSION_SEED, move |channel, _p2p| {
|
|
|
- let msgs2 = msgs.clone();
|
|
|
- async move { ProtocolDchat::init(channel, msgs2).await }
|
|
|
- })
|
|
|
- .await;
|
|
|
- debug!(target: "dchat", "Dchat::register_protocol() [STOP]");
|
|
|
- Ok(())
|
|
|
- }
|
|
|
- // ANCHOR_END: register_protocol
|
|
|
-
|
|
|
- // ANCHOR: start
|
|
|
- async fn start(&mut self) -> Result<()> {
|
|
|
- debug!(target: "dchat", "Dchat::start() [START]");
|
|
|
-
|
|
|
- self.register_protocol(self.recv_msgs.clone()).await?;
|
|
|
- self.p2p.clone().start().await?;
|
|
|
-
|
|
|
- self.menu().await?;
|
|
|
-
|
|
|
- self.p2p.stop().await;
|
|
|
-
|
|
|
- debug!(target: "dchat", "Dchat::start() [STOP]");
|
|
|
- Ok(())
|
|
|
- }
|
|
|
- // ANCHOR_END: start
|
|
|
-
|
|
|
// ANCHOR: send
|
|
|
async fn send(&self, msg: String) -> Result<()> {
|
|
|
let dchatmsg = DchatMsg { msg };
|
|
|
@@ -170,145 +109,89 @@ impl AppSettings {
|
|
|
Self { accept_addr, net }
|
|
|
}
|
|
|
}
|
|
|
-// ANCHOR_END: app_settings
|
|
|
-
|
|
|
-// ANCHOR: alice
|
|
|
-fn alice() -> Result<AppSettings> {
|
|
|
- let log_level = simplelog::LevelFilter::Debug;
|
|
|
- let log_config = simplelog::Config::default();
|
|
|
-
|
|
|
- let log_path = "/tmp/alice.log";
|
|
|
- let file = File::create(log_path).unwrap();
|
|
|
- WriteLogger::init(log_level, log_config, file)?;
|
|
|
-
|
|
|
- let seed = Url::parse("tcp://127.0.0.1:50515").unwrap();
|
|
|
- let inbound = Url::parse("tcp://127.0.0.1:51554").unwrap();
|
|
|
- let ext_addr = Url::parse("tcp://127.0.0.1:51554").unwrap();
|
|
|
-
|
|
|
- let net = Settings {
|
|
|
- inbound_addrs: vec![inbound],
|
|
|
- external_addrs: vec![ext_addr],
|
|
|
- seeds: vec![seed],
|
|
|
- localnet: true,
|
|
|
- ..Default::default()
|
|
|
- };
|
|
|
-
|
|
|
- let accept_addr = Url::parse("tcp://127.0.0.1:55054").unwrap();
|
|
|
- let settings = AppSettings::new(accept_addr, net);
|
|
|
-
|
|
|
- Ok(settings)
|
|
|
-}
|
|
|
-// ANCHOR_END: alice
|
|
|
-
|
|
|
-// ANCHOR: bob
|
|
|
-fn bob() -> Result<AppSettings> {
|
|
|
- let log_level = simplelog::LevelFilter::Debug;
|
|
|
- let log_config = simplelog::Config::default();
|
|
|
-
|
|
|
- let log_path = "/tmp/bob.log";
|
|
|
- let file = File::create(log_path).unwrap();
|
|
|
- WriteLogger::init(log_level, log_config, file)?;
|
|
|
-
|
|
|
- let seed = Url::parse("tcp://127.0.0.1:50515").unwrap();
|
|
|
-
|
|
|
- let net = Settings {
|
|
|
- inbound_addrs: vec![],
|
|
|
- outbound_connections: 5,
|
|
|
- seeds: vec![seed],
|
|
|
- localnet: true,
|
|
|
- ..Default::default()
|
|
|
- };
|
|
|
-
|
|
|
- let accept_addr = Url::parse("tcp://127.0.0.1:51054").unwrap();
|
|
|
- let settings = AppSettings::new(accept_addr, net);
|
|
|
-
|
|
|
- Ok(settings)
|
|
|
-}
|
|
|
-// ANCHOR_END: bob
|
|
|
-
|
|
|
-// ANCHOR: main
|
|
|
-fn main() -> Result<()> {
|
|
|
- smol::block_on(async {
|
|
|
- let settings: Result<AppSettings> = match std::env::args().nth(1) {
|
|
|
- Some(id) => match id.as_str() {
|
|
|
- "a" => alice(),
|
|
|
- "b" => bob(),
|
|
|
- _ => Err(ErrorMissingSpecifier.into()),
|
|
|
- },
|
|
|
- None => Err(ErrorMissingSpecifier.into()),
|
|
|
- };
|
|
|
+async_daemonize!(realmain);
|
|
|
+async fn realmain(args: Args, ex: Arc<smol::Executor<'static>>) -> Result<()> {
|
|
|
+ let cfg_path = get_config_path(args.config, CONFIG_FILE)?;
|
|
|
+ let toml_contents = std::fs::read_to_string(cfg_path)?;
|
|
|
|
|
|
- let settings = settings?.clone();
|
|
|
- let ex = Arc::new(Executor::new());
|
|
|
- let p2p = net::P2p::new(settings.net, ex.clone()).await;
|
|
|
- let msgs: DchatMsgsBuffer = Arc::new(Mutex::new(vec![DchatMsg { msg: String::new() }]));
|
|
|
- let mut dchat = Dchat::new(p2p.clone(), msgs);
|
|
|
+ // ANCHOR: dchat
|
|
|
+ let p2p = net::P2p::new(args.net.into(), ex.clone()).await;
|
|
|
+ let msgs: DchatMsgsBuffer = Arc::new(Mutex::new(vec![DchatMsg { msg: String::new() }]));
|
|
|
+ let dchat = Dchat::new(p2p.clone(), msgs.clone());
|
|
|
|
|
|
- // ANCHOR: dnet_sub
|
|
|
- let dnet_sub = JsonSubscriber::new("dnet.subscribe_events");
|
|
|
- let dnet_sub_ = dnet_sub.clone();
|
|
|
- let p2p_ = p2p.clone();
|
|
|
- let dnet_task = StoppableTask::new();
|
|
|
- dnet_task.clone().start(
|
|
|
- async move {
|
|
|
- let dnet_sub = p2p_.dnet_subscribe().await;
|
|
|
- loop {
|
|
|
- let event = dnet_sub.receive().await;
|
|
|
- //debug!("Got dnet event: {:?}", event);
|
|
|
- dnet_sub_.notify(vec![event.into()].into()).await;
|
|
|
- }
|
|
|
- },
|
|
|
- |res| async {
|
|
|
- match res {
|
|
|
- Ok(()) | Err(DarkfiError::DetachedTaskStopped) => { /* Do nothing */ }
|
|
|
- Err(e) => panic!("{}", e),
|
|
|
- }
|
|
|
- },
|
|
|
- DarkfiError::DetachedTaskStopped,
|
|
|
- ex.clone(),
|
|
|
- );
|
|
|
- // ANCHOR_END: dnet_sub
|
|
|
+ // ANCHOR: register_protocol
|
|
|
+ info!("Registering Dchat protocol");
|
|
|
+ let registry = p2p.protocol_registry();
|
|
|
+ registry
|
|
|
+ .register(!net::session::SESSION_SEED, move |channel, _p2p| {
|
|
|
+ let msgs_ = msgs.clone();
|
|
|
+ async move { ProtocolDchat::init(channel, msgs_).await }
|
|
|
+ })
|
|
|
+ .await;
|
|
|
+ // ANCHOR_END: register_protocol
|
|
|
|
|
|
- // ANCHOR: json_init
|
|
|
- let accept_addr = settings.accept_addr.clone();
|
|
|
+ // ANCHOR: dnet
|
|
|
+ info!("Starting dnet subs task");
|
|
|
+ let dnet_sub = JsonSubscriber::new("dnet.subscribe_events");
|
|
|
+ let dnet_sub_ = dnet_sub.clone();
|
|
|
+ let p2p_ = p2p.clone();
|
|
|
+ let dnet_task = StoppableTask::new();
|
|
|
+ dnet_task.clone().start(
|
|
|
+ async move {
|
|
|
+ let dnet_sub = p2p_.dnet_subscribe().await;
|
|
|
+ loop {
|
|
|
+ let event = dnet_sub.receive().await;
|
|
|
+ debug!("Got dnet event: {:?}", event);
|
|
|
+ dnet_sub_.notify(vec![event.into()].into()).await;
|
|
|
+ }
|
|
|
+ },
|
|
|
+ |res| async {
|
|
|
+ match res {
|
|
|
+ Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
|
|
|
+ Err(e) => panic!("{}", e),
|
|
|
+ }
|
|
|
+ },
|
|
|
+ Error::DetachedTaskStopped,
|
|
|
+ ex.clone(),
|
|
|
+ );
|
|
|
+ // ANCHOR_end: dnet
|
|
|
+
|
|
|
+ // ANCHOR: rpc
|
|
|
+ //info!("Starting JSON-RPC server");
|
|
|
+ let rpc_connections = Mutex::new(HashSet::new());
|
|
|
+ let rpc = Arc::new(JsonRpcInterface { p2p: p2p.clone(), rpc_connections, dnet_sub });
|
|
|
+ let _ex = ex.clone();
|
|
|
+
|
|
|
+ let rpc_task = StoppableTask::new();
|
|
|
+ rpc_task.clone().start(
|
|
|
+ listen_and_serve(args.rpc_listen, rpc.clone(), None, ex.clone()),
|
|
|
+ |res| async move {
|
|
|
+ match res {
|
|
|
+ Ok(()) | Err(Error::RpcServerStopped) => rpc.stop_connections().await,
|
|
|
+ Err(e) => error!("Failed stopping JSON-RPC server: {}", e),
|
|
|
+ }
|
|
|
+ },
|
|
|
+ Error::RpcServerStopped,
|
|
|
+ ex.clone(),
|
|
|
+ );
|
|
|
+ // ANCHOR_end: rpc
|
|
|
|
|
|
- let rpc_connections = Mutex::new(HashSet::new());
|
|
|
- let rpc = Arc::new(JsonRpcInterface {
|
|
|
- addr: accept_addr.clone(),
|
|
|
- p2p,
|
|
|
- rpc_connections,
|
|
|
- dnet_sub,
|
|
|
- });
|
|
|
- let _ex = ex.clone();
|
|
|
+ info!("Starting P2P network");
|
|
|
+ p2p.clone().start().await?;
|
|
|
|
|
|
- let rpc_task = StoppableTask::new();
|
|
|
- rpc_task.clone().start(
|
|
|
- listen_and_serve(accept_addr.clone(), rpc.clone(), None, ex.clone()),
|
|
|
- |res| async move {
|
|
|
- match res {
|
|
|
- Ok(()) | Err(DarkfiError::RpcServerStopped) => rpc.stop_connections().await,
|
|
|
- Err(e) => error!("Failed stopping JSON-RPC server: {}", e),
|
|
|
- }
|
|
|
- },
|
|
|
- DarkfiError::RpcServerStopped,
|
|
|
- ex.clone(),
|
|
|
- );
|
|
|
- // ANCHOR_END: json_init
|
|
|
+ // Signal handling for graceful termination.
|
|
|
+ let (signals_handler, signals_task) = SignalHandler::new(ex)?;
|
|
|
+ signals_handler.wait_termination(signals_task).await?;
|
|
|
+ info!("Caught termination signal, cleaning up and exiting...");
|
|
|
|
|
|
- let nthreads = std::thread::available_parallelism().unwrap().get();
|
|
|
- let (signal, shutdown) = smol::channel::unbounded::<()>();
|
|
|
+ info!("Stopping P2P network");
|
|
|
+ p2p.stop().await;
|
|
|
|
|
|
- let (_, result) = Parallel::new()
|
|
|
- .each(0..nthreads, |_| smol::future::block_on(ex.run(shutdown.recv())))
|
|
|
- .finish(|| {
|
|
|
- smol::future::block_on(async move {
|
|
|
- dchat.start().await?;
|
|
|
- drop(signal);
|
|
|
- Ok(())
|
|
|
- })
|
|
|
- });
|
|
|
+ info!("Stopping JSON-RPC server");
|
|
|
+ rpc_task.stop().await;
|
|
|
+ dnet_task.stop().await;
|
|
|
|
|
|
- result
|
|
|
- })
|
|
|
+ info!("Shut down successfully");
|
|
|
+ Ok(())
|
|
|
}
|
|
|
-// ANCHOR_END: main
|
|
|
+//// ANCHOR: main
|