|
@@ -17,7 +17,7 @@ use darkfi::{
|
|
|
Error, Result,
|
|
Error, Result,
|
|
|
};
|
|
};
|
|
|
use easy_parallel::Parallel;
|
|
use easy_parallel::Parallel;
|
|
|
-use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, FutureExt};
|
|
|
|
|
|
|
+use futures::{io::BufReader, AsyncBufReadExt, AsyncReadExt, FutureExt};
|
|
|
use log::{debug, error, info, warn};
|
|
use log::{debug, error, info, warn};
|
|
|
use simplelog::{ColorChoice, TermLogger, TerminalMode};
|
|
use simplelog::{ColorChoice, TermLogger, TerminalMode};
|
|
|
use smol::future;
|
|
use smol::future;
|
|
@@ -74,43 +74,11 @@ async fn broadcast_msg(
|
|
|
Ok(())
|
|
Ok(())
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-async fn plugin_process(
|
|
|
|
|
- recv_for_plugin: Receiver<String>,
|
|
|
|
|
- send_from_plugin: Sender<String>,
|
|
|
|
|
- stream: TcpStream,
|
|
|
|
|
-) -> Result<()> {
|
|
|
|
|
- let peer_addr = stream.peer_addr()?.clone();
|
|
|
|
|
- let (reader, mut writer) = stream.split();
|
|
|
|
|
- let mut reader = BufReader::new(reader);
|
|
|
|
|
-
|
|
|
|
|
- loop {
|
|
|
|
|
- let mut line = String::new();
|
|
|
|
|
- futures::select! {
|
|
|
|
|
- irc_msg = recv_for_plugin.recv().fuse() => {
|
|
|
|
|
- let irc_msg = irc_msg?;
|
|
|
|
|
- info!("Plugin recv {}", irc_msg);
|
|
|
|
|
- writer.write_all(irc_msg.as_bytes()).await?;
|
|
|
|
|
- },
|
|
|
|
|
- err = reader.read_line(&mut line).fuse() => {
|
|
|
|
|
- if let Err(e) = err {
|
|
|
|
|
- warn!("Read line error: {}", e);
|
|
|
|
|
- return Ok(())
|
|
|
|
|
- }
|
|
|
|
|
- let irc_msg = clean_input(line, &peer_addr)?;
|
|
|
|
|
- info!("Plugin send {}", irc_msg);
|
|
|
|
|
- send_from_plugin.send(irc_msg).await?;
|
|
|
|
|
- }
|
|
|
|
|
- };
|
|
|
|
|
- }
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
async fn process(
|
|
async fn process(
|
|
|
raft_receiver: Receiver<Privmsg>,
|
|
raft_receiver: Receiver<Privmsg>,
|
|
|
stream: TcpStream,
|
|
stream: TcpStream,
|
|
|
peer_addr: SocketAddr,
|
|
peer_addr: SocketAddr,
|
|
|
raft_sender: Sender<Privmsg>,
|
|
raft_sender: Sender<Privmsg>,
|
|
|
- send_for_plugin: Sender<String>,
|
|
|
|
|
- recv_from_plugin: Receiver<String>,
|
|
|
|
|
seen_msg_id: SeenMsgIds,
|
|
seen_msg_id: SeenMsgIds,
|
|
|
) -> Result<()> {
|
|
) -> Result<()> {
|
|
|
let (reader, writer) = stream.split();
|
|
let (reader, writer) = stream.split();
|
|
@@ -134,13 +102,6 @@ async fn process(
|
|
|
|
|
|
|
|
let irc_msg = build_irc_msg(&msg);
|
|
let irc_msg = build_irc_msg(&msg);
|
|
|
conn.reply(&irc_msg).await?;
|
|
conn.reply(&irc_msg).await?;
|
|
|
- send_for_plugin.send(irc_msg).await?;
|
|
|
|
|
- }
|
|
|
|
|
- irc_msg = recv_from_plugin.recv().fuse() => {
|
|
|
|
|
- let irc_msg = irc_msg?;
|
|
|
|
|
- info!("Receive msg from plugin");
|
|
|
|
|
- broadcast_msg(irc_msg, peer_addr, &mut conn).await?;
|
|
|
|
|
-
|
|
|
|
|
}
|
|
}
|
|
|
err = reader.read_line(&mut line).fuse() => {
|
|
err = reader.read_line(&mut line).fuse() => {
|
|
|
if let Err(e) = err {
|
|
if let Err(e) = err {
|
|
@@ -188,9 +149,6 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
|
|
|
listen_and_serve(rpc_config, rpc_interface, executor_cloned.clone()).await
|
|
listen_and_serve(rpc_config, rpc_interface, executor_cloned.clone()).await
|
|
|
});
|
|
});
|
|
|
|
|
|
|
|
- let (send_from_plugin, recv_from_plugin) = async_channel::unbounded::<String>();
|
|
|
|
|
- let (send_for_plugin, recv_for_plugin) = async_channel::unbounded::<String>();
|
|
|
|
|
-
|
|
|
|
|
//
|
|
//
|
|
|
// IRC instance
|
|
// IRC instance
|
|
|
//
|
|
//
|
|
@@ -213,47 +171,12 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
|
|
|
stream,
|
|
stream,
|
|
|
peer_addr,
|
|
peer_addr,
|
|
|
raft_sender.clone(),
|
|
raft_sender.clone(),
|
|
|
- send_for_plugin.clone(),
|
|
|
|
|
- recv_from_plugin.clone(),
|
|
|
|
|
seen_msg_id.clone(),
|
|
seen_msg_id.clone(),
|
|
|
))
|
|
))
|
|
|
.detach();
|
|
.detach();
|
|
|
}
|
|
}
|
|
|
});
|
|
});
|
|
|
|
|
|
|
|
- //
|
|
|
|
|
- // Plugin instance
|
|
|
|
|
- //
|
|
|
|
|
- let executor_cloned = executor.clone();
|
|
|
|
|
- let send_from_plugin_cloned = send_from_plugin.clone();
|
|
|
|
|
- let recv_for_plugin_cloned = recv_for_plugin.clone();
|
|
|
|
|
- let plugin_task: smol::Task<Result<()>> = executor.spawn(async move {
|
|
|
|
|
- if settings.plugin_listen.is_none() {
|
|
|
|
|
- return Ok(())
|
|
|
|
|
- }
|
|
|
|
|
- let plugin_listener = TcpListener::bind(settings.plugin_listen.unwrap()).await?;
|
|
|
|
|
- let plugin_local_addr = plugin_listener.local_addr()?;
|
|
|
|
|
- info!("Plugin listening on {}", plugin_local_addr);
|
|
|
|
|
- loop {
|
|
|
|
|
- let (stream, peer_addr) = match plugin_listener.accept().await {
|
|
|
|
|
- Ok((s, a)) => (s, a),
|
|
|
|
|
- Err(e) => {
|
|
|
|
|
- error!("Failed listening for connections: {}", e);
|
|
|
|
|
- return Err(Error::ServiceStopped)
|
|
|
|
|
- }
|
|
|
|
|
- };
|
|
|
|
|
-
|
|
|
|
|
- info!("Plugin Accepted client: {}", peer_addr);
|
|
|
|
|
- executor_cloned
|
|
|
|
|
- .spawn(plugin_process(
|
|
|
|
|
- recv_for_plugin_cloned.clone(),
|
|
|
|
|
- send_from_plugin_cloned.clone(),
|
|
|
|
|
- stream,
|
|
|
|
|
- ))
|
|
|
|
|
- .detach();
|
|
|
|
|
- }
|
|
|
|
|
- });
|
|
|
|
|
-
|
|
|
|
|
let (signal, shutdown) = async_channel::bounded::<()>(1);
|
|
let (signal, shutdown) = async_channel::bounded::<()>(1);
|
|
|
ctrlc_async::set_async_handler(async move {
|
|
ctrlc_async::set_async_handler(async move {
|
|
|
warn!(target: "ircd", "ircd start Exit Signal");
|
|
warn!(target: "ircd", "ircd start Exit Signal");
|
|
@@ -261,7 +184,6 @@ async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
|
|
|
signal.send(()).await.unwrap();
|
|
signal.send(()).await.unwrap();
|
|
|
rpc_task.cancel().await;
|
|
rpc_task.cancel().await;
|
|
|
irc_task.cancel().await;
|
|
irc_task.cancel().await;
|
|
|
- plugin_task.cancel().await;
|
|
|
|
|
})
|
|
})
|
|
|
.unwrap();
|
|
.unwrap();
|
|
|
|
|
|