Explorar o código

ircd: General code cleanup.

parazyd %!s(int64=4) %!d(string=hai) anos
pai
achega
a2e6917f3c

+ 16 - 16
Cargo.lock

@@ -849,7 +849,7 @@ dependencies = [
  "async-trait",
  "bdk",
  "bitcoin",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "futures",
@@ -974,9 +974,9 @@ dependencies = [
 
 [[package]]
 name = "clap"
-version = "3.1.6"
+version = "3.1.8"
 source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "d8c93436c21e4698bacadf42917db28b23017027a4deccb35dbe47a7e7840123"
+checksum = "71c47df61d9e16dc010b55dba1952a57d8c215dbb533fd13cdd13369aac73b1c"
 dependencies = [
  "atty",
  "bitflags",
@@ -991,9 +991,9 @@ dependencies = [
 
 [[package]]
 name = "clap_derive"
-version = "3.1.4"
+version = "3.1.7"
 source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "da95d038ede1a964ce99f49cbe27a7fb538d1da595e4b4f70b8c8f338d17bf16"
+checksum = "a3aab4734e083b809aaf5794e14e756d1c798d2c69c7f7de7a09a2f5214993c1"
 dependencies = [
  "heck 0.4.0",
  "proc-macro-error",
@@ -1516,7 +1516,7 @@ dependencies = [
  "async-executor",
  "async-std",
  "async-trait",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "futures",
@@ -1564,7 +1564,7 @@ dependencies = [
  "bs58",
  "bytes",
  "chrono",
- "clap 3.1.6",
+ "clap 3.1.8",
  "crypto_api_chachapoly",
  "darkfi-derive",
  "darkfi-derive-internal",
@@ -1642,7 +1642,7 @@ dependencies = [
  "async-executor",
  "async-std",
  "async-trait",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "fxhash",
@@ -1910,7 +1910,7 @@ version = "0.3.0"
 dependencies = [
  "async-channel",
  "async-std",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "fxhash",
@@ -1937,7 +1937,7 @@ name = "drk"
 version = "0.3.0"
 dependencies = [
  "async-std",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "log",
  "prettytable-rs",
@@ -2517,7 +2517,7 @@ dependencies = [
  "async-channel",
  "async-executor",
  "async-std",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "log",
@@ -2995,7 +2995,7 @@ dependencies = [
  "async-executor",
  "async-std",
  "async-trait",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "futures",
@@ -6057,7 +6057,7 @@ dependencies = [
  "async-std",
  "async-trait",
  "chrono",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "futures",
@@ -6080,7 +6080,7 @@ dependencies = [
  "async-std",
  "async-trait",
  "chrono",
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
  "easy-parallel",
  "futures",
@@ -6574,7 +6574,7 @@ name = "vanityaddr"
 version = "0.3.0"
 dependencies = [
  "bs58",
- "clap 3.1.6",
+ "clap 3.1.8",
  "ctrlc",
  "darkfi",
  "indicatif 0.17.0-rc.9",
@@ -7239,7 +7239,7 @@ dependencies = [
 name = "zkas"
 version = "0.3.0"
 dependencies = [
- "clap 3.1.6",
+ "clap 3.1.8",
  "darkfi",
 ]
 

+ 7 - 5
bin/ircd/Cargo.toml

@@ -1,13 +1,15 @@
 [package]
 name = "ircd"
 version = "0.3.0"
+homepage = "https://dark.fi"
+description = "P2P IRC daemon"
+authors = ["darkfi <dev@dark.fi>"]
+repository = "https://github.com/darkrenaissance/darkfi"
+license = "AGPL-3.0-only"
 edition = "2021"
 
-[dependencies.darkfi]
-path = "../../"
-features = ["net", "rpc"]
-
 [dependencies]
+darkfi = {path = "../../", features = ["net", "rpc"]}
 # Async
 smol = "1.2.5"
 futures = "0.3.21"
@@ -21,7 +23,7 @@ easy-parallel = "3.2.0"
 rand = "0.8.5"
 
 # Misc
-clap = "3.1.6"
+clap = {version = "3.1.8", features = ["derive"]}
 log = "0.4.16"
 simplelog = "0.12.0-alpha1"
 fxhash = "0.2.1"

+ 141 - 160
bin/ircd/src/main.rs

@@ -1,173 +1,184 @@
-use async_std::io::BufReader;
-use std::{
-    net::{SocketAddr, TcpListener, TcpStream},
-    sync::Arc,
-};
+use std::{net::SocketAddr, sync::Arc};
 
-extern crate clap;
+use async_channel::Receiver;
 use async_executor::Executor;
-use async_trait::async_trait;
+use async_std::net::{TcpListener, TcpStream};
+use clap::Parser;
 use easy_parallel::Parallel;
-use futures::{AsyncBufReadExt, AsyncReadExt, FutureExt};
+use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
 use log::{debug, error, info, warn};
-use serde_json::{json, Value};
 use simplelog::{ColorChoice, TermLogger, TerminalMode};
-use smol::Async;
 
 use darkfi::{
-    net,
-    rpc::{
-        jsonrpc::{error as jsonerr, response as jsonresp, ErrorCode::*, JsonRequest, JsonResult},
-        rpcserver::{listen_and_serve, RequestHandler, RpcServerConfig},
-    },
+    cli_desc, net,
+    rpc::rpcserver::{listen_and_serve, RpcServerConfig},
     util::cli::log_config,
     Error, Result,
 };
 
-mod irc_server;
-mod privmsg;
-mod program_options;
-mod protocol_privmsg;
+pub(crate) mod proto;
+pub(crate) mod rpc;
+pub(crate) mod server;
 
 use crate::{
-    irc_server::IrcServerConnection,
-    privmsg::{PrivMsg, SeenPrivMsgIds, SeenPrivMsgIdsPtr},
-    program_options::ProgramOptions,
-    protocol_privmsg::ProtocolPrivMsg,
+    proto::privmsg::{Privmsg, ProtocolPrivmsg, SeenPrivmsgIds, SeenPrivmsgIdsPtr},
+    rpc::JsonRpcInterface,
+    server::IrcServerConnection,
 };
 
-async fn process(
-    recvr: async_channel::Receiver<Arc<PrivMsg>>,
-    stream: Async<TcpStream>,
-    peer_addr: SocketAddr,
-    p2p: net::P2pPtr,
-    seen_privmsg_ids: SeenPrivMsgIdsPtr,
-    _executor: Arc<Executor<'_>>,
-) -> Result<()> {
-    let (reader, writer) = stream.split();
+#[derive(Parser)]
+#[clap(name = "ircd", about = cli_desc!(), version)]
+struct Args {
+    /// Accept address
+    #[clap(short, long)]
+    accept: Option<SocketAddr>,
 
-    let mut reader = BufReader::new(reader);
-    let mut connection = IrcServerConnection::new(writer, seen_privmsg_ids);
+    /// Seed node (repeatable)
+    #[clap(short, long)]
+    seed: Vec<SocketAddr>,
 
-    loop {
-        let mut line = String::new();
-        futures::select! {
-            privmsg = recvr.recv().fuse() => {
-                let privmsg = privmsg.expect("internal message queue error");
-                debug!("ABOUT TO SEND {:?}", privmsg);
-                let irc_msg = format!(
-                    ":{}!darkfi@127.0.0.1 PRIVMSG {} :{}\n",
-                    privmsg.nickname,
-                    privmsg.channel,
-                    privmsg.message
-                );
+    /// Manual connection (repeatable)
+    #[clap(short, long)]
+    connect: Vec<SocketAddr>,
 
-                connection.reply(&irc_msg).await?;
-            }
-            err = reader.read_line(&mut line).fuse() => {
-                if let Err(err) = err {
-                    warn!("Read line error. Closing stream for {}: {}", peer_addr, err);
-                    return Ok(())
-                }
-                process_user_input(line, peer_addr, &mut connection, p2p.clone()).await?;
-            }
-        };
-    }
+    /// Connection slots
+    #[clap(long, default_value_t = 0)]
+    slots: u32,
+
+    /// External address
+    #[clap(short, long)]
+    external: Option<SocketAddr>,
+
+    /// IRC listen address
+    #[clap(short = 'r', long, default_value = "127.0.0.1:6667")]
+    irc: SocketAddr,
+
+    /// RPC listen address
+    #[clap(long, default_value = "127.0.0.1:8000")]
+    rpc: SocketAddr,
+
+    /// Verbosity level
+    #[clap(short, parse(from_occurrences))]
+    verbose: u8,
 }
 
 async fn process_user_input(
     mut line: String,
     peer_addr: SocketAddr,
-    connection: &mut IrcServerConnection,
+    conn: &mut IrcServerConnection,
     p2p: net::P2pPtr,
 ) -> Result<()> {
     if line.is_empty() {
         warn!("Received empty line from {}. Closing connection.", peer_addr);
         return Err(Error::ChannelStopped)
     }
-    assert!(&line[(line.len() - 1)..] == "\n");
-    // Remove the \n character
+
+    assert!(&line[(line.len() - 2)..] == "\r\n");
+    // Remove CRLF
+    line.pop();
     line.pop();
 
     debug!("Received '{}' from {}", line, peer_addr);
 
-    if let Err(err) = connection.update(line, p2p.clone()).await {
-        warn!("Connection error: {} for {}", err, peer_addr);
+    if let Err(e) = conn.update(line, p2p.clone()).await {
+        warn!("Connection error: {} for {}", e, peer_addr);
         return Err(Error::ChannelStopped)
     }
 
     Ok(())
 }
 
-async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<()> {
-    let listener = match Async::<TcpListener>::bind(options.irc_accept_addr) {
-        Ok(listener) => listener,
-        Err(err) => {
-            error!("Bind listener failed: {}", err);
-            return Err(Error::OperationFailed)
-        }
-    };
-    let local_addr = match listener.get_ref().local_addr() {
-        Ok(addr) => addr,
-        Err(err) => {
-            error!("Failed to get local address: {}", err);
-            return Err(Error::OperationFailed)
-        }
-    };
+async fn process(
+    receiver: Receiver<Arc<Privmsg>>,
+    stream: TcpStream,
+    peer_addr: SocketAddr,
+    p2p: net::P2pPtr,
+    seen_privmsg_ids: SeenPrivmsgIdsPtr,
+) -> Result<()> {
+    let (reader, writer) = stream.split();
+
+    let mut reader = BufReader::new(reader);
+    let mut conn = IrcServerConnection::new(writer, seen_privmsg_ids);
+
+    loop {
+        let mut line = String::new();
+        futures::select! {
+            privmsg = receiver.recv().fuse() => {
+                let msg = privmsg.expect("internal message queue error");
+                debug!("ABOUT TO SEND: {:?}", msg);
+                let irc_msg = format!(":{}!anon@dark.fi PRIVMSG {} :{}\r\n",
+                    msg.nickname,
+                    msg.channel,
+                    msg.message,
+                );
+
+                conn.reply(&irc_msg).await?;
+            }
+
+            err = reader.read_line(&mut line).fuse() => {
+                if let Err(e) = err {
+                    warn!("Read line error. Closing stream for {}: {}", peer_addr, e);
+                    return Ok(())
+                }
+
+                process_user_input(line, peer_addr, &mut conn, p2p.clone()).await?;
+            }
+        };
+    }
+}
+
+async fn start(executor: Arc<Executor<'_>>, args: Args, net_settings: net::Settings) -> Result<()> {
+    let listener = TcpListener::bind(args.irc).await?;
+    let local_addr = listener.local_addr()?;
     info!("Listening on {}", local_addr);
 
-    let server_config = RpcServerConfig {
-        socket_addr: options.rpc_listen_addr,
+    let rpc_config = RpcServerConfig {
+        socket_addr: args.rpc,
+        // TODO: Use net/transport:
         use_tls: false,
-        // this is all random filler that is meaningless bc tls is disabled
         identity_path: Default::default(),
         identity_pass: Default::default(),
     };
 
-    let seen_privmsg_ids = SeenPrivMsgIds::new();
-
     //
-    // PrivMsg protocol
+    // Privmsg protocol
     //
-    let p2p = net::P2p::new(options.network_settings).await;
-    let registry = p2p.protocol_registry();
+    let seen_privmsg_ids = SeenPrivmsgIds::new();
+    let seen_privmsg_ids_clone = seen_privmsg_ids.clone();
+
+    let (sender, receiver) = async_channel::unbounded();
+    let sender_clone = sender.clone();
 
-    let (sender, recvr) = async_channel::unbounded();
-    let seen_privmsg_ids2 = seen_privmsg_ids.clone();
-    let sender2 = sender.clone();
+    let p2p = net::P2p::new(net_settings).await;
+    let registry = p2p.protocol_registry();
     registry
         .register(!net::SESSION_SEED, move |channel, p2p| {
-            let sender = sender2.clone();
-            let seen_privmsg_ids = seen_privmsg_ids2.clone();
-            async move { ProtocolPrivMsg::init(channel, sender, seen_privmsg_ids, p2p).await }
+            let sender = sender_clone.clone();
+            let seen_privmsg_ids = seen_privmsg_ids_clone.clone();
+            async move { ProtocolPrivmsg::init(channel, sender, seen_privmsg_ids, p2p).await }
         })
         .await;
 
     //
-    // p2p network main instance
+    // P2P network main instance
     //
-    // Performs seed session
     p2p.clone().start(executor.clone()).await?;
-    // Actual main p2p session
-    let ex2 = executor.clone();
-    let p2p2 = p2p.clone();
+    let executor_clone = executor.clone();
+    let p2p_clone = p2p.clone();
     executor
         .spawn(async move {
-            if let Err(err) = p2p2.run(ex2).await {
-                error!("Error: p2p run failed {}", err);
+            if let Err(e) = p2p_clone.run(executor_clone).await {
+                error!("P2P run failed: {}", e);
             }
         })
         .detach();
 
     //
     // RPC interface
-    //
-    let ex2 = executor.clone();
-    let ex3 = ex2.clone();
-    let rpc_interface =
-        Arc::new(JsonRpcInterface { p2p: p2p.clone(), rpc_listen_addr: options.rpc_listen_addr });
+    let executor_clone = executor.clone();
+    let rpc_interface = Arc::new(JsonRpcInterface { p2p: p2p.clone(), addr: args.rpc });
     executor
-        .spawn(async move { listen_and_serve(server_config, rpc_interface, ex3).await })
+        .spawn(async move { listen_and_serve(rpc_config, rpc_interface, executor_clone.clone()).await })
         .detach();
 
     //
@@ -176,81 +187,51 @@ async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<(
     loop {
         let (stream, peer_addr) = match listener.accept().await {
             Ok((s, a)) => (s, a),
-            Err(err) => {
-                error!("Error listening for connections: {}", err);
+            Err(e) => {
+                error!("Failed listening for connections: {}", e);
                 return Err(Error::ServiceStopped)
             }
         };
+
         info!("Accepted client: {}", peer_addr);
 
-        let p2p2 = p2p.clone();
-        let ex2 = executor.clone();
+        let p2p_clone = p2p.clone();
         executor
-            .spawn(process(recvr.clone(), stream, peer_addr, p2p2, seen_privmsg_ids.clone(), ex2))
+            .spawn(process(
+                receiver.clone(),
+                stream,
+                peer_addr,
+                p2p_clone,
+                seen_privmsg_ids.clone(),
+            ))
             .detach();
     }
 }
 
-struct JsonRpcInterface {
-    p2p: net::P2pPtr,
-    rpc_listen_addr: SocketAddr,
-}
-
-#[async_trait]
-impl RequestHandler for JsonRpcInterface {
-    async fn handle_request(&self, req: JsonRequest, _executor: Arc<Executor<'_>>) -> JsonResult {
-        if req.params.as_array().is_none() {
-            return JsonResult::Err(jsonerr(InvalidParams, None, req.id))
-        }
-
-        debug!(target: "RPC", "--> {}", serde_json::to_string(&req).unwrap());
-
-        match req.method.as_str() {
-            Some("ping") => return self.pong(req.id, req.params).await,
-            Some("get_info") => return self.get_info(req.id, req.params).await,
-            Some(_) | None => return JsonResult::Err(jsonerr(MethodNotFound, None, req.id)),
-        }
-    }
-}
-
-impl JsonRpcInterface {
-    // --> {"jsonrpc": "2.0", "method": "ping", "params": [], "id": 42}
-    // <-- {"jsonrpc": "2.0", "result": "pong", "id": 42}
-    async fn pong(&self, id: Value, _params: Value) -> JsonResult {
-        JsonResult::Resp(jsonresp(json!("pong"), id))
-    }
-
-    //--> {"jsonrpc": "2.0", "method": "poll", "params": [], "id": 42}
-    // <-- {"jsonrpc": "2.0", "result": {"nodeID": [], "nodeinfo" [], "id": 42}
-    async fn get_info(&self, id: Value, _params: Value) -> JsonResult {
-        let resp = self.p2p.get_info().await;
-        JsonResult::Resp(jsonresp(resp, id))
-    }
-}
-
 fn main() -> Result<()> {
-    let options = ProgramOptions::load()?;
-
-    let verbosity_level = options.app.occurrences_of("verbose");
-
-    let (lvl, cfg) = log_config(verbosity_level)?;
-
-    TermLogger::init(lvl, cfg, TerminalMode::Mixed, ColorChoice::Auto)?;
+    let args = Args::parse();
+
+    let (lvl, conf) = log_config(args.verbose.into())?;
+    TermLogger::init(lvl, conf, TerminalMode::Mixed, ColorChoice::Auto)?;
+
+    let net_settings = net::Settings {
+        inbound: args.accept,
+        outbound_connections: args.slots,
+        external_addr: args.external,
+        peers: args.connect.clone(),
+        seeds: args.seed.clone(),
+        ..Default::default()
+    };
 
     let ex = Arc::new(Executor::new());
+    let ex_clone = ex.clone();
     let (signal, shutdown) = async_channel::unbounded::<()>();
-
-    let ex2 = ex.clone();
-
-    // let nthreads = num_cpus::get();
-    // debug!(target: "IRC DAEMON", "Run {} executor threads", nthreads);
-
     let (_, result) = Parallel::new()
         .each(0..4, |_| smol::future::block_on(ex.run(shutdown.recv())))
         // Run the main future on the current thread.
         .finish(|| {
             smol::future::block_on(async move {
-                start(ex2.clone(), options).await?;
+                start(ex_clone.clone(), args, net_settings).await?;
                 drop(signal);
                 Ok::<(), darkfi::Error>(())
             })

+ 0 - 68
bin/ircd/src/privmsg.rs

@@ -1,68 +0,0 @@
-use async_std::sync::Mutex;
-use std::{io, sync::Arc};
-
-use fxhash::FxHashSet;
-
-use darkfi::{
-    net,
-    util::serial::{Decodable, Encodable},
-    Result,
-};
-
-pub type PrivMsgId = u32;
-
-#[derive(Debug, Clone)]
-pub struct PrivMsg {
-    pub id: PrivMsgId,
-    pub nickname: String,
-    pub channel: String,
-    pub message: String,
-}
-
-impl net::Message for PrivMsg {
-    fn name() -> &'static str {
-        "privmsg"
-    }
-}
-
-impl Encodable for PrivMsg {
-    fn encode<S: io::Write>(&self, mut s: S) -> Result<usize> {
-        let mut len = 0;
-        len += self.id.encode(&mut s)?;
-        len += self.nickname.encode(&mut s)?;
-        len += self.channel.encode(&mut s)?;
-        len += self.message.encode(&mut s)?;
-        Ok(len)
-    }
-}
-
-impl Decodable for PrivMsg {
-    fn decode<D: io::Read>(mut d: D) -> Result<Self> {
-        Ok(Self {
-            id: Decodable::decode(&mut d)?,
-            nickname: Decodable::decode(&mut d)?,
-            channel: Decodable::decode(&mut d)?,
-            message: Decodable::decode(&mut d)?,
-        })
-    }
-}
-
-pub struct SeenPrivMsgIds {
-    privmsg_ids: Mutex<FxHashSet<PrivMsgId>>,
-}
-
-pub type SeenPrivMsgIdsPtr = Arc<SeenPrivMsgIds>;
-
-impl SeenPrivMsgIds {
-    pub fn new() -> Arc<Self> {
-        Arc::new(Self { privmsg_ids: Mutex::new(FxHashSet::default()) })
-    }
-
-    pub async fn add_seen(&self, id: u32) {
-        self.privmsg_ids.lock().await.insert(id);
-    }
-
-    pub async fn is_seen(&self, id: u32) -> bool {
-        self.privmsg_ids.lock().await.contains(&id)
-    }
-}

+ 0 - 158
bin/ircd/src/program_options.rs

@@ -1,158 +0,0 @@
-use clap::{Arg, ArgMatches, Command};
-use std::net::SocketAddr;
-
-use darkfi::{net, Result};
-
-pub struct ProgramOptions {
-    pub network_settings: net::Settings,
-    pub log_path: Box<std::path::PathBuf>,
-    pub irc_accept_addr: SocketAddr,
-    pub rpc_listen_addr: SocketAddr,
-    pub app: ArgMatches,
-}
-
-impl ProgramOptions {
-    pub fn load() -> Result<ProgramOptions> {
-        let app = Command::new("dfi")
-            .version("0.1.0")
-            .author("Amir Taaki <amir@dyne.org>")
-            .about("Dark node")
-            .arg(
-                Arg::new("ACCEPT")
-                    .short('a')
-                    .long("accept")
-                    .value_name("ACCEPT")
-                    .help("Accept address")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("SEED_NODES")
-                    .short('s')
-                    .long("seeds")
-                    .value_name("SEED_NODES")
-                    .help("Seed nodes")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("CONNECTS")
-                    .short('c')
-                    .long("connect")
-                    .value_name("CONNECTS")
-                    .help("Manual connections")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("CONNECT_SLOTS")
-                    .long("slots")
-                    .value_name("CONNECT_SLOTS")
-                    .help("Connection slots")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("EXTERNAL_ADDR")
-                    .short('e')
-                    .long("external")
-                    .value_name("EXTERNAL_ADDR")
-                    .help("External address")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("LOG_PATH")
-                    .long("log")
-                    .value_name("LOG_PATH")
-                    .help("Logfile path")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("IRC_ACCEPT")
-                    .short('r')
-                    .long("irc")
-                    .value_name("IRC_ACCEPT")
-                    .help("IRC accept address")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("RPC_LISTEN")
-                    .long("rpc")
-                    .value_name("RPC_LISTEN")
-                    .help("RPC listen address")
-                    .takes_value(true),
-            )
-            .arg(
-                Arg::new("verbose")
-                    .short('v')
-                    .long("verbose")
-                    .multiple_occurrences(true)
-                    .help("Sets the level of verbosity"),
-            )
-            .get_matches();
-
-        let accept_addr = if let Some(accept_addr) = app.value_of("ACCEPT") {
-            Some(accept_addr.parse()?)
-        } else {
-            None
-        };
-
-        let mut seed_addrs: Vec<SocketAddr> = vec![];
-        if let Some(seeds) = app.values_of("SEED_NODES") {
-            for seed in seeds {
-                seed_addrs.push(seed.parse()?);
-            }
-        }
-
-        let mut manual_connects: Vec<SocketAddr> = vec![];
-        if let Some(connections) = app.values_of("CONNECTS") {
-            for connect in connections {
-                manual_connects.push(connect.parse()?);
-            }
-        }
-
-        let connection_slots = if let Some(connection_slots) = app.value_of("CONNECT_SLOTS") {
-            connection_slots.parse()?
-        } else {
-            0
-        };
-
-        let external_addr = if let Some(external_addr) = app.value_of("EXTERNAL_ADDR") {
-            Some(external_addr.parse()?)
-        } else {
-            None
-        };
-
-        let log_path = Box::new(
-            if let Some(log_path) = app.value_of("LOG_PATH") {
-                std::path::Path::new(log_path)
-            } else {
-                std::path::Path::new("/tmp/darkfid.log")
-            }
-            .to_path_buf(),
-        );
-
-        let irc_accept_addr = if let Some(accept_addr) = app.value_of("IRC_ACCEPT") {
-            accept_addr.parse()?
-        } else {
-            ([127, 0, 0, 1], 6667).into()
-        };
-
-        let rpc_listen_addr = if let Some(rpc_addr) = app.value_of("RPC_LISTEN") {
-            rpc_addr.parse()?
-        } else {
-            ([127, 0, 0, 1], 8000).into()
-        };
-
-        Ok(ProgramOptions {
-            network_settings: net::Settings {
-                inbound: accept_addr,
-                outbound_connections: connection_slots,
-                external_addr,
-                peers: manual_connects,
-                seeds: seed_addrs,
-                ..Default::default()
-            },
-            log_path,
-            irc_accept_addr,
-            rpc_listen_addr,
-            app,
-        })
-    }
-}

+ 1 - 0
bin/ircd/src/proto/mod.rs

@@ -0,0 +1 @@
+pub(crate) mod privmsg;

+ 121 - 0
bin/ircd/src/proto/privmsg.rs

@@ -0,0 +1,121 @@
+use std::sync::Arc;
+
+use async_channel::Sender;
+use async_executor::Executor;
+use async_std::sync::Mutex;
+use async_trait::async_trait;
+use fxhash::FxHashSet;
+use log::debug;
+
+use darkfi::{
+    net,
+    util::serial::{SerialDecodable, SerialEncodable},
+    Result,
+};
+
+pub type PrivmsgId = u32;
+
+#[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
+pub struct Privmsg {
+    pub id: PrivmsgId,
+    pub nickname: String,
+    pub channel: String,
+    pub message: String,
+}
+
+impl net::Message for Privmsg {
+    fn name() -> &'static str {
+        "privmsg"
+    }
+}
+
+pub struct SeenPrivmsgIds {
+    ids: Mutex<FxHashSet<PrivmsgId>>,
+}
+
+pub type SeenPrivmsgIdsPtr = Arc<SeenPrivmsgIds>;
+
+impl SeenPrivmsgIds {
+    pub fn new() -> Arc<Self> {
+        Arc::new(Self { ids: Mutex::new(FxHashSet::default()) })
+    }
+
+    pub async fn add_seen(&self, id: u32) {
+        self.ids.lock().await.insert(id);
+    }
+
+    pub async fn is_seen(&self, id: u32) -> bool {
+        self.ids.lock().await.contains(&id)
+    }
+}
+
+pub struct ProtocolPrivmsg {
+    notify_queue_sender: Sender<Arc<Privmsg>>,
+    privmsg_sub: net::MessageSubscription<Privmsg>,
+    jobsman: net::ProtocolJobsManagerPtr,
+    seen_ids: SeenPrivmsgIdsPtr,
+    p2p: net::P2pPtr,
+}
+
+#[async_trait]
+impl net::ProtocolBase for ProtocolPrivmsg {
+    /// Starts ping-pong keep-alive messages exchange. Runs ping-pong in the
+    /// protocol task manager, then queues the reply. Sends out a ping and
+    /// waits for pong reply. Waits for ping and replies with a pong.
+    async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
+        debug!(target: "ircd", "ProtocolPrivMsg::start() [START]");
+        self.jobsman.clone().start(executor.clone());
+        self.jobsman.clone().spawn(self.clone().handle_receive_privmsg(), executor.clone()).await;
+        debug!(target: "ircd", "ProtocolPrivmsg::start() [END]");
+        Ok(())
+    }
+
+    fn name(&self) -> &'static str {
+        "ProtocolPrivMsg"
+    }
+}
+
+impl ProtocolPrivmsg {
+    pub async fn init(
+        channel: net::ChannelPtr,
+        notify_queue_sender: Sender<Arc<Privmsg>>,
+        seen_ids: SeenPrivmsgIdsPtr,
+        p2p: net::P2pPtr,
+    ) -> net::ProtocolBasePtr {
+        let message_subsystem = channel.get_message_subsystem();
+        message_subsystem.add_dispatch::<Privmsg>().await;
+
+        let sub = channel.subscribe_msg::<Privmsg>().await.expect("Missing Privmsg dispatcher!");
+
+        Arc::new(Self {
+            notify_queue_sender,
+            privmsg_sub: sub,
+            jobsman: net::ProtocolJobsManager::new("PrivmsgProtocol", channel),
+            seen_ids,
+            p2p,
+        })
+    }
+
+    async fn handle_receive_privmsg(self: Arc<Self>) -> Result<()> {
+        debug!(target: "ircd", "ProtocolPrivmsg::handle_receive_privmsg() [START]");
+
+        loop {
+            let privmsg = self.privmsg_sub.receive().await?;
+
+            debug!(target: "ircd", "ProtocolPrivmsg::handle_receive_privmsg() received {:?}", privmsg);
+
+            // Do we already have this message?
+            if self.seen_ids.is_seen(privmsg.id).await {
+                continue
+            }
+
+            self.seen_ids.add_seen(privmsg.id).await;
+
+            // If not, then broadcast to network.
+            let privmsg_copy = (*privmsg).clone();
+            self.p2p.broadcast(privmsg_copy).await?;
+
+            self.notify_queue_sender.send(privmsg).await.expect("notify_queue_sender send failed!");
+        }
+    }
+}

+ 0 - 83
bin/ircd/src/protocol_privmsg.rs

@@ -1,83 +0,0 @@
-use async_executor::Executor;
-use async_trait::async_trait;
-
-use darkfi::{net, Result};
-use log::debug;
-use std::sync::Arc;
-
-use crate::privmsg::{PrivMsg, SeenPrivMsgIdsPtr};
-
-pub struct ProtocolPrivMsg {
-    notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
-    privmsg_sub: net::MessageSubscription<PrivMsg>,
-    jobsman: net::ProtocolJobsManagerPtr,
-    seen_privmsg_ids: SeenPrivMsgIdsPtr,
-    p2p: net::P2pPtr,
-}
-
-impl ProtocolPrivMsg {
-    pub async fn init(
-        channel: net::ChannelPtr,
-        notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
-        seen_privmsg_ids: SeenPrivMsgIdsPtr,
-        p2p: net::P2pPtr,
-    ) -> net::ProtocolBasePtr {
-        let message_subsytem = channel.get_message_subsystem();
-        message_subsytem.add_dispatch::<PrivMsg>().await;
-
-        let privmsg_sub =
-            channel.subscribe_msg::<PrivMsg>().await.expect("Missing PrivMsg dispatcher!");
-
-        Arc::new(Self {
-            notify_queue_sender,
-            privmsg_sub,
-            jobsman: net::ProtocolJobsManager::new("PrivMsgProtocol", channel),
-            seen_privmsg_ids,
-            p2p,
-        })
-    }
-
-    async fn handle_receive_privmsg(self: Arc<Self>) -> Result<()> {
-        debug!(target: "ircd", "ProtocolPrivMsg::handle_receive_privmsg() [START]");
-        loop {
-            let privmsg = self.privmsg_sub.receive().await?;
-
-            debug!(
-                target: "ircd",
-                "ProtocolPrivMsg::handle_receive_privmsg() received {:?}",
-                privmsg
-            );
-
-            // Do we already have this message?
-            if self.seen_privmsg_ids.is_seen(privmsg.id).await {
-                continue
-            }
-
-            self.seen_privmsg_ids.add_seen(privmsg.id).await;
-
-            // If not then broadcast to everybody else
-            let privmsg_copy = (*privmsg).clone();
-            self.p2p.broadcast(privmsg_copy).await?;
-
-            self.notify_queue_sender.send(privmsg).await.expect("notify_queue_sender send failed!");
-        }
-    }
-}
-
-#[async_trait]
-impl net::ProtocolBase for ProtocolPrivMsg {
-    /// Starts ping-pong keep-alive messages exchange. Runs ping-pong in the
-    /// protocol task manager, then queues the reply. Sends out a ping and
-    /// waits for pong reply. Waits for ping and replies with a pong.
-    async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
-        debug!(target: "ircd", "ProtocolPrivMsg::start() [START]");
-        self.jobsman.clone().start(executor.clone());
-        self.jobsman.clone().spawn(self.clone().handle_receive_privmsg(), executor.clone()).await;
-        debug!(target: "ircd", "ProtocolPrivMsg::start() [END]");
-        Ok(())
-    }
-
-    fn name(&self) -> &'static str {
-        "ProtocolPrivMsg"
-    }
-}

+ 56 - 0
bin/ircd/src/rpc.rs

@@ -0,0 +1,56 @@
+use std::{net::SocketAddr, sync::Arc};
+
+use async_executor::Executor;
+use async_trait::async_trait;
+use log::debug;
+use serde_json::{json, Value};
+
+use darkfi::{
+    net,
+    rpc::{
+        jsonrpc,
+        jsonrpc::{ErrorCode, JsonRequest, JsonResult},
+        rpcserver::RequestHandler,
+    },
+};
+
+pub struct JsonRpcInterface {
+    pub p2p: net::P2pPtr,
+    pub addr: SocketAddr,
+}
+
+#[async_trait]
+impl RequestHandler for JsonRpcInterface {
+    async fn handle_request(&self, req: JsonRequest, _executor: Arc<Executor<'_>>) -> JsonResult {
+        if req.params.as_array().is_none() {
+            return jsonrpc::error(ErrorCode::InvalidRequest, None, req.id).into()
+        }
+
+        debug!(target: "RPC", "--> {}", serde_json::to_string(&req).unwrap());
+
+        match req.method.as_str() {
+            Some("ping") => self.pong(req.id, req.params).await,
+            Some("get_info") => self.get_info(req.id, req.params).await,
+            Some(_) | None => jsonrpc::error(ErrorCode::MethodNotFound, None, req.id).into(),
+        }
+    }
+}
+
+impl JsonRpcInterface {
+    // RPCAPI:
+    // Replies to a ping method.
+    // --> {"jsonrpc": "2.0", "method": "ping", "params": [], "id": 42}
+    // <-- {"jsonrpc": "2.0", "result": "pong", "id": 42}
+    async fn pong(&self, id: Value, _params: Value) -> JsonResult {
+        jsonrpc::response(json!("pong"), id).into()
+    }
+
+    // RPCAPI:
+    // Retrieves P2P network information.
+    // --> {"jsonrpc": "2.0", "method": "get_info", "params": [], "id": 42}
+    // <-- {"jsonrpc": "2.0", result": {"nodeID": [], "nodeinfo": [], "id": 42}
+    async fn get_info(&self, id: Value, _params: Value) -> JsonResult {
+        let resp = self.p2p.get_info().await;
+        jsonrpc::response(resp, id).into()
+    }
+}

+ 40 - 49
bin/ircd/src/irc_server.rs → bin/ircd/src/server.rs

@@ -1,33 +1,15 @@
+use async_std::net::TcpStream;
 use futures::{io::WriteHalf, AsyncWriteExt};
-use log::{debug, info};
+use log::{debug, info, warn};
 use rand::{rngs::OsRng, RngCore};
-use smol::Async;
-use std::net::TcpStream;
 
 use darkfi::{net, Error, Result};
 
-use crate::privmsg::{PrivMsg, SeenPrivMsgIdsPtr};
-
-/*
-NICK fifififif
-USER username 0 * :Real
-:behemoth 001 fifififif :Hi, welcome to IRC
-:behemoth 002 fifififif :Your host is behemoth, running version miniircd-2.1
-:behemoth 003 fifififif :This server was created sometime
-:behemoth 004 fifififif behemoth miniircd-2.1 o o
-:behemoth 251 fifififif :There are 1 users and 0 services on 1 server
-:behemoth 422 fifififif :MOTD File is missing
-JOIN #dev
-:fifififif!username@127.0.0.1 JOIN #dev
-:behemoth 331 fifififif #dev :No topic is set
-:behemoth 353 fifififif = #dev :fifififif
-:behemoth 366 fifififif #dev :End of NAMES list
-PRIVMSG #dev hihi
-*/
+use crate::proto::privmsg::{Privmsg, SeenPrivmsgIdsPtr};
 
 pub struct IrcServerConnection {
-    write_stream: WriteHalf<Async<TcpStream>>,
-    seen_privmsg_ids: SeenPrivMsgIdsPtr,
+    write_stream: WriteHalf<TcpStream>,
+    seen_privmsg_ids: SeenPrivmsgIdsPtr,
     is_nick_init: bool,
     is_user_init: bool,
     is_registered: bool,
@@ -36,13 +18,10 @@ pub struct IrcServerConnection {
 }
 
 impl IrcServerConnection {
-    pub fn new(
-        write_stream: WriteHalf<Async<TcpStream>>,
-        seen_privmsg_ids: SeenPrivMsgIdsPtr,
-    ) -> Self {
+    pub fn new(write_stream: WriteHalf<TcpStream>, seen_ids: SeenPrivmsgIdsPtr) -> Self {
         Self {
             write_stream,
-            seen_privmsg_ids,
+            seen_privmsg_ids: seen_ids,
             is_nick_init: false,
             is_user_init: false,
             is_registered: false,
@@ -53,83 +32,95 @@ impl IrcServerConnection {
 
     pub async fn update(&mut self, line: String, p2p: net::P2pPtr) -> Result<()> {
         let mut tokens = line.split_ascii_whitespace();
-        // Commands can begin with :garbage but we will reject clients doing that for now
-        // to keep the protocol simple and focused.
+        // Commands can begin with :garbage but we will reject clients doing
+        // that for now to keep the protocol simple and focused.
         let command = tokens.next().ok_or(Error::MalformedPacket)?;
 
         debug!("Received command: {}", command);
 
         match command {
+            "USER" => {
+                // We can stuff any extra things like public keys in here.
+                // Ignore it for now.
+                self.is_user_init = true;
+            }
             "NICK" => {
                 let nickname = tokens.next().ok_or(Error::MalformedPacket)?;
                 self.is_nick_init = true;
                 let old_nick = std::mem::replace(&mut self.nickname, nickname.to_string());
 
-                let nick_reply = format!(":{}!darkfi@127.0.0.1 NICK {}\n", old_nick, self.nickname);
+                let nick_reply = format!(":{}!anon@dark.fi NICK {}\r\n", old_nick, self.nickname);
                 self.reply(&nick_reply).await?;
             }
-            "USER" => {
-                // We can stuff any extra things like public keys in here
-                // Ignore it for now
-                self.is_user_init = true;
-            }
             "JOIN" => {
                 // Ignore since channels are all autojoin
-                //let channel = tokens.next().ok_or(Error::MalformedPacket)?;
-                //self.channels.push(channel.to_string());
+                // let channel = tokens.next().ok_or(Error::MalformedPacket)?;
+                // self.channels.push(channel.to_string());
 
-                //let join_reply = format!(":{}!darkfi@127.0.0.1 JOIN {}\n", self.nickname,
-                // channel); self.reply(&join_reply).await?;
+                // let join_reply = format!(":{}!anon@dark.fi JOIN {}\r\n", self.nickname, channel);
+                // self.reply(&join_reply).await?;
 
-                //self.write_stream.write_all(b":f00!f00@127.0.0.1 PRIVMSG #dev :y0\n").await?;
+                // self.write_stream.write_all(b":f00!f00@127.0.01 PRIVMSG #dev :y0\r\n").await?;
             }
             "PING" => {
                 let line_clone = line.clone();
                 let split_line: Vec<&str> = line_clone.split_whitespace().collect();
                 if split_line.len() > 1 && split_line[0] == "PING" {
-                    let pong = format!("PONG {}\n", split_line[1]);
+                    let pong = format!("PONG {}\r\n", split_line[1]);
                     self.reply(&pong).await?;
                 }
             }
             "PRIVMSG" => {
                 let channel = tokens.next().ok_or(Error::MalformedPacket)?;
-
                 let substr_idx = line.find(':').ok_or(Error::MalformedPacket)?;
+
                 if substr_idx >= line.len() {
                     return Err(Error::MalformedPacket)
                 }
+
                 let message = &line[substr_idx + 1..];
                 info!("Message {}: {}", channel, message);
 
                 let random_id = OsRng.next_u32();
                 self.seen_privmsg_ids.add_seen(random_id).await;
 
-                let protocol_msg = PrivMsg {
+                let protocol_msg = Privmsg {
                     id: random_id,
                     nickname: self.nickname.clone(),
                     channel: channel.to_string(),
                     message: message.to_string(),
                 };
+
                 p2p.broadcast(protocol_msg).await?;
             }
             "QUIT" => {
                 // Close the connection
                 return Err(Error::ServiceStopped)
             }
-            _ => {}
+            _ => {
+                warn!("Unimplemented `{}` command", command);
+            }
         }
 
         if !self.is_registered && self.is_nick_init && self.is_user_init {
             debug!("Initializing peer connection");
-            let register_reply = format!(":darkfi 001 {} :Let there be dark\n", self.nickname);
+            let register_reply = format!(":darkfi 001 {} :Let there be dark\r\n", self.nickname);
             self.reply(&register_reply).await?;
             self.is_registered = true;
 
             // Auto-joins
-            for channel in ["#dev", "#markets", "#welcome"] {
-                let join_reply = format!(":{}!darkfi@127.0.0.1 JOIN {}\n", self.nickname, channel);
-                self.reply(&join_reply).await?;
+            macro_rules! autojoin {
+                ($channel:expr,$topic:expr) => {
+                    let j = format!(":{}!anon@dark.fi JOIN {}\r\n", self.nickname, $channel);
+                    let t = format!(":DarkFi TOPIC {} :{}\r\n", $channel, $topic);
+                    self.reply(&j).await?;
+                    self.reply(&t).await?;
+                };
             }
+
+            autojoin!("#dev", "Development of DarkFi");
+            autojoin!("#markets", "Markets, trading, DeFi, algo, biz, finance, and economics");
+            autojoin!("#memes", "Memetic engineering");
         }
 
         Ok(())