/* This file is part of DarkFi (https://dark.fi)
*
* Copyright (C) 2020-2023 Dyne.org foundation
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU Affero General Public License as
* published by the Free Software Foundation, either version 3 of the
* License, or (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU Affero General Public License for more details.
*
* You should have received a copy of the GNU Affero General Public License
* along with this program. If not, see .
*/
use std::{
collections::{HashMap, HashSet},
sync::{
atomic::{AtomicBool, AtomicUsize, Ordering::SeqCst},
Arc,
},
};
use darkfi::{
event_graph::{proto::EventPut, Event, NULL_ID},
system::Subscription,
Error, Result,
};
use darkfi_serial::{deserialize_async_partial, serialize_async};
use futures::FutureExt;
use log::{debug, error, warn};
use smol::{
io::{self, AsyncBufReadExt, AsyncWriteExt, BufReader},
lock::Mutex,
net::SocketAddr,
prelude::{AsyncRead, AsyncWrite},
};
use super::{server::IrcServer, Privmsg, SERVER_NAME};
const PENALTY_LIMIT: usize = 5;
/// Reply types, we can either send server replies, or client replies.
pub enum ReplyType {
/// Server reply, we have to use numerics
Server((u16, String)),
/// Client reply, message from someone to some{one,where}
Client((String, String)),
/// Pong reply, we just use server origin
Pong(String),
/// CAP reply
Cap(String),
}
/// Stateful IRC client, used for each client connection
pub struct Client {
/// Pointer to parent `IrcServer`
pub server: Arc,
/// Subscription for incoming events
pub incoming: Subscription,
/// Client socket addr
pub addr: SocketAddr,
/// ID of the last sent event
pub last_sent: Mutex,
/// Active (joined) channels for this client
pub channels: Mutex>,
/// Penalty counter, when limit is reached, disconnect client
pub penalty: AtomicUsize,
/// Registration marker
pub registered: AtomicBool,
/// Registration pause marker
pub reg_paused: AtomicBool,
/// Client username
pub username: Mutex,
/// Client nickname
pub nickname: Mutex,
/// Client realname
pub realname: Mutex,
/// Client caps
pub caps: Mutex>,
}
impl Client {
/// Instantiate a new Client.
pub async fn new(
server: Arc,
incoming: Subscription,
addr: SocketAddr,
) -> Result {
let caps = HashMap::from([("no-history".to_string(), false)]);
Ok(Self {
server,
incoming,
addr,
last_sent: Mutex::new(NULL_ID),
channels: Mutex::new(HashSet::new()),
penalty: AtomicUsize::new(0),
registered: AtomicBool::new(false),
reg_paused: AtomicBool::new(false),
username: Mutex::new(String::from("*")),
nickname: Mutex::new(String::from("*")),
realname: Mutex::new(String::from("*")),
caps: Mutex::new(caps),
})
}
/// This function handles a single IRC client. We listen to messages from the
/// IRC client and relay them to the network, and we also get notified of
/// incoming messages and relay them to the IRC client. The notifications come
/// from events being inserted into the Event Graph.
pub async fn multiplex_connection(&self, stream: S) -> Result<()>
where
S: AsyncRead + AsyncWrite + Unpin + Send + 'static,
{
let (reader, mut writer) = io::split(stream);
let mut reader = BufReader::new(reader);
// Our buffer for the client line
let mut line = String::new();
loop {
futures::select! {
// Process message from the IRC client
r = reader.read_line(&mut line).fuse() => {
// If something failed during reading, we disconnect.
if let Err(e) = r {
error!("[IRC CLIENT] Read failed for {}: {}", self.addr, e);
self.incoming.unsubscribe().await;
return Err(Error::ChannelStopped)
}
// If the penalty limit is reached, disconnect the client.
if self.penalty.load(SeqCst) == PENALTY_LIMIT {
self.incoming.unsubscribe().await;
return Err(Error::ChannelStopped)
}
// We'll be strict here and disconnect the client
// in case line processing failed in any way.
match self.process_client_line(&line, &mut writer).await {
// If we got an event back, we should broadcast it.
// This means we add it to our DAG, and the DAG will
// handle the rest of the propagation.
Ok(Some(event)) => {
// Update the last sent event.
*self.last_sent.lock().await = event.id();
// If it fails for some reason, for now, we just note it
// and pass.
if let Err(e) = self.server.darkirc.event_graph.dag_insert(event.clone()).await {
error!("[IRC CLIENT] Failed inserting new event to DAG: {}", e);
} else {
// Otherwise, broadcast it
self.server.darkirc.p2p.broadcast(&EventPut(event)).await;
}
}
// If we got nothing, we just pass.
Ok(None) => {}
// If we got an error, we disconnect the client.
Err(e) => {
self.incoming.unsubscribe().await;
return Err(e)
}
}
// Clear the line buffer
line = String::new();
continue
}
// Process message from the network. These should only be PRIVMSG.
r = self.incoming.receive().fuse() => {
// We will skip this if it's our own message.
if *self.last_sent.lock().await == r.id() {
continue
}
// Try to deserialize the `Event`'s content into a `Privmsg`
let mut privmsg: Privmsg = match deserialize_async_partial(r.content()).await {
Ok((v, _)) => v,
Err(e) => {
error!("[IRC CLIENT] Failed deserializing incoming Privmsg event: {}", e);
continue
}
};
// If successful, potentially decrypt it:
self.server.try_decrypt(&mut privmsg).await;
// If we have this channel, or it's a DM to our nickname, forward
// it to the client.
let have_channel = self.channels.lock().await.contains(&privmsg.channel);
let msg_for_self = *self.nickname.lock().await == privmsg.channel;
if have_channel || msg_for_self {
// Add the nickname to the list of nicks on the channel
(*self.server.channels.lock().await).get_mut(&privmsg.channel)
.unwrap().nicks.insert(privmsg.nick.clone());
// Format the message
let msg = format!("PRIVMSG {} :{}", privmsg.channel, privmsg.msg);
// Send it to the client
let reply = ReplyType::Client((privmsg.nick, msg));
if let Err(e) = self.reply(&mut writer, &reply).await {
error!("[IRC CLIENT] Failed writing PRIVMSG to client: {}", e);
continue
}
}
}
}
}
}
/// Send a reply to the IRC client. Matches on the reply type.
async fn reply(&self, writer: &mut W, reply: &ReplyType) -> Result<()>
where
W: AsyncWrite + Unpin,
{
let r = match reply {
ReplyType::Server((rpl, msg)) => format!(":{} {:03} {}", SERVER_NAME, rpl, msg),
ReplyType::Client((nick, msg)) => format!(":{}!~anon@darkirc {}", nick, msg),
ReplyType::Pong(origin) => format!(":{} PONG :{}", SERVER_NAME, origin),
ReplyType::Cap(msg) => format!(":{} {}", SERVER_NAME, msg),
};
debug!("[{}] <-- {}", self.addr, r);
writer.write(r.as_bytes()).await?;
writer.write(b"\r\n").await?;
writer.flush().await?;
Ok(())
}
/// Handle the incoming line given sent by the IRC client
async fn process_client_line(&self, line: &str, writer: &mut W) -> Result