/* This file is part of DarkFi (https://dark.fi)
*
* Copyright (C) 2020-2026 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, VecDeque},
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 sled_overlay::sled;
use smol::{
io::{self, AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader},
lock::{OnceCell, RwLock},
net::SocketAddr,
prelude::{AsyncRead, AsyncWrite},
};
use tracing::{debug, error, warn};
use super::{
server::{IrcServer, RlnMessageReservation, MAX_MSG_LEN},
NickServ, SERVER_NAME,
};
use crate::Privmsg;
const PENALTY_LIMIT: usize = 5;
const MAX_IRC_LINE_LEN: usize = 1024;
const MAX_PENDING_PRIVMSGS: usize = 128;
/// Read one IRC line without allowing unbounded buffer growth.
async fn read_bounded_line(reader: &mut R, line: &mut String) -> Result
where
R: AsyncBufRead + Unpin,
{
line.clear();
let mut bytes = Vec::new();
loop {
let (consumed, complete) = {
let available = reader.fill_buf().await?;
if available.is_empty() {
if bytes.is_empty() {
return Ok(0)
}
*line = String::from_utf8(bytes)?;
return Ok(line.len())
}
let newline = available.iter().position(|b| *b == b'\n');
let take = newline.map_or(available.len(), |idx| idx + 1);
if bytes.len().saturating_add(take) > MAX_IRC_LINE_LEN {
return Err(Error::ParseFailed("IRC line too long"))
}
bytes.extend_from_slice(&available[..take]);
(take, newline.is_some())
};
reader.consume(consumed);
if complete {
*line = String::from_utf8(bytes)?;
return Ok(line.len())
}
}
}
fn enqueue_pending_privmsg(args_queue: &mut VecDeque, privmsg: Privmsg) -> bool {
if args_queue.len() >= MAX_PENDING_PRIVMSGS {
return false
}
args_queue.push_back(privmsg);
true
}
/// 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),
/// NOTICE reply (from, to, what)
Notice((String, String, String)),
}
/// Stateful IRC client handler, used for each client connection
pub struct Client {
/// Pointer to parent `IrcServer`
pub server: Arc,
/// Subscription for incoming events
pub incoming: Subscription,
/// Subscription for incoming static events
pub incoming_st: Subscription,
/// Client socket addr
pub addr: SocketAddr,
/// ID of the last sent event
pub last_sent: RwLock,
/// Active (joined) channels for this client
pub channels: RwLock>,
/// Penalty counter, when limit is reached, disconnect client
pub penalty: AtomicUsize,
/// Registration marker
pub registered: AtomicBool,
/// Registration pause marker
pub reg_paused: AtomicBool,
/// CAP END marker
pub is_cap_end: AtomicBool,
/// Password setup marker
pub is_pass_set: AtomicBool,
/// Client username
pub username: Arc>,
/// Client nickname
pub nickname: Arc>,
/// Client realname
pub realname: RwLock,
/// Client caps
pub caps: RwLock>,
/// Set of seen messages for the user
/// TODO: It grows indefinitely, needs to be pruned.
pub seen: OnceCell,
/// NickServ instance
pub nickserv: Arc,
}
impl Client {
/// Instantiate a new Client.
pub async fn new(
server: Arc,
incoming: Subscription,
incoming_st: Subscription,
addr: SocketAddr,
) -> Result {
let caps =
HashMap::from([("no-history".to_string(), false), ("no-autojoin".to_string(), false)]);
let username = Arc::new(RwLock::new(String::from("*")));
let nickname = Arc::new(RwLock::new(String::from("*")));
Ok(Self {
server: server.clone(),
incoming,
incoming_st,
addr,
last_sent: RwLock::new(NULL_ID),
channels: RwLock::new(HashSet::new()),
penalty: AtomicUsize::new(0),
registered: AtomicBool::new(false),
reg_paused: AtomicBool::new(false),
is_cap_end: AtomicBool::new(false),
is_pass_set: AtomicBool::new(false),
username: username.clone(),
nickname: nickname.clone(),
realname: RwLock::new(String::from("*")),
caps: RwLock::new(caps),
seen: OnceCell::new(),
nickserv: Arc::new(
NickServ::new(username.clone(), nickname.clone(), server.clone()).await?,
),
})
}
/// 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();
let mut args_queue: VecDeque<_> = VecDeque::new();
loop {
futures::select! {
// Process message from the IRC client
r = read_bounded_line(&mut reader, &mut line).fuse() => {
// If client closed unexpectedly, we disconnect.
if let Ok(0) = r {
error!("[IRC CLIENT] Read failed for {}: Client disconnected", self.addr);
self.incoming.unsubscribe().await;
self.incoming_st.unsubscribe().await;
return Err(Error::ChannelStopped)
}
// If something failed during reading, we disconnect.
if let Err(e) = r {
error!("[IRC CLIENT] Read failed for {}: {e}", self.addr);
self.incoming.unsubscribe().await;
self.incoming_st.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;
self.incoming_st.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, &mut args_queue).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(events)) => {
for event in events {
// Update the last sent event.
let event_id = event.header.id();
*self.last_sent.write().await = event_id;
let current_genesis = self.server.darkirc.event_graph.current_genesis.read().await;
let dag_name = current_genesis.header.timestamp.to_string();
drop(current_genesis);
// Build the RLN signal blob before touching the local
// DAG when RLN is enabled. With RLN disabled, outbound
// events deliberately carry no proof blob.
let blob = if self.server.darkirc.event_graph.rln_enabled() {
let (rln_identity, mid) = match self
.server
.reserve_rln_message_id(event.header.timestamp)
.await?
{
RlnMessageReservation::Reserved {
identity,
message_id,
} => (identity, message_id),
RlnMessageReservation::MissingIdentity => {
warn!(
"[IRC CLIENT] No RLN identity registered; \
refusing to send. Use \
`/msg NickServ REGISTER ...` to register."
);
continue
}
RlnMessageReservation::BudgetExhausted => {
warn!(
"[IRC CLIENT] RLN message budget \
exhausted for this epoch; dropping \
message to avoid slash"
);
continue
}
};
match rln_identity
.create_signal(
&event,
mid,
&self.server.darkirc.event_graph,
)
.await
{
Ok(blob) => serialize_async(&blob).await,
Err(e) => {
error!(
"[IRC CLIENT] Failed creating RLN \
signal proof: {e}"
);
return Err(e)
}
}
} else {
Vec::new()
};
// Commit our outbound signal through
// the safe public API. It inserts the
// header, verifies and stores the RLN
// blob, then commits the event body.
if let Err(e) = self
.server
.darkirc
.event_graph
.insert_signal_with_blob(&event, &blob, &dag_name)
.await
{
error!(
"[IRC CLIENT] Failed inserting verified \
signal event: {e}"
);
continue
}
// We sent this, so it should be considered seen.
if let Err(e) = self.mark_seen(&event_id).await {
error!("[IRC CLIENT] (multiplex_connection) self.mark_seen({event_id}) failed: {e}");
return Err(e)
}
if let Err(e) =
self.server.darkirc.p2p.broadcast(&EventPut(event, blob)).await
{
error!("[IRC CLIENT] Event broadcast was not admitted: {e}");
}
}
}
// If we got nothing, we just pass.
Ok(None) => {}
// If we got an error, we disconnect the client.
Err(e) => {
self.incoming.unsubscribe().await;
self.incoming_st.unsubscribe().await;
return Err(e)
}
}
// Clear the line buffer
line = String::new();
}
// Process message from the network. These should only be PRIVMSG.
//
// N.b. handling "historical messages", i.e. outstanding messages
// which have occured when darkirc is offline are handled in
// ) -> Result> {>
// for which the logic for delivery should be kept in sync
r = self.incoming.receive().fuse() => {
// We will skip this if it's our own message.
let event_id = r.header.id();
if *self.last_sent.read().await == event_id {
continue
}
// If this event was seen, skip it
match self.is_seen(&event_id).await {
Ok(true) => continue,
Ok(false) => {},
Err(e) => {
error!("[IRC CLIENT] (multiplex_connection) self.is_seen({event_id}) failed: {e}");
return Err(e)
}
}
// Try to deserialize the `Event`'s content into a `Privmsg`
let mut privmsg = match deserialize_async_partial(r.content()).await {
Ok((v, _)) => v,
Err(e) => {
error!(target: "irc::client", "[IRC CLIENT] Failed deserializing event: {e}");
continue
}
};
// If successful, potentially decrypt it:
self.server.try_decrypt(&mut privmsg, self.nickname.read().await.as_ref()).await;
// We should skip any attempts to contact services from the network.
if ["nickserv", "chanserv"].contains(&privmsg.nick.to_lowercase().as_str()) {
continue
}
// If the privmsg is not intented for any of the given
// channels or contacts, ignore it
// otherwise add it as a reply and mark it as seen
// in the seen_events tree.
let channels = self.channels.read().await;
let contacts = self.server.contacts.read().await;
if !channels.contains(&privmsg.channel) &&
!contacts.contains_key(&privmsg.channel)
{
continue
}
// Add the nickname to the list of nicks on the channel, if it's a channel.
let mut chans_lock = self.server.channels.write().await;
if let Some(chan) = chans_lock.get_mut(&privmsg.channel) {
chan.nicks.insert(privmsg.nick.clone());
}
drop(chans_lock);
// Handle message lines individually
for line in privmsg.msg.lines() {
// Skip empty lines
if line.is_empty() {
continue
}
// Format the message
let msg = format!("PRIVMSG {} :{line}", privmsg.channel);
// Send it to the client
let reply = ReplyType::Client((privmsg.nick.clone(), msg));
if let Err(e) = self.reply(&mut writer, &reply).await {
error!("[IRC CLIENT] Failed writing PRIVMSG to client: {e}");
continue
}
}
// Mark the message as seen for this USER
if let Err(e) = self.mark_seen(&event_id).await {
error!("[IRC CLIENT] (multiplex_connection) self.mark_seen({}) failed: {}", event_id, e);
return Err(e)
}
}
// Process message from the network. These should only be RLN identities.
r = self.incoming_st.receive().fuse() => {
// We will skip this if it's our own message.
let event_id = r.header.id();
if *self.last_sent.read().await == event_id {
continue
}
// If this event was seen, skip it
match self.is_seen(&event_id).await {
Ok(true) => continue,
Ok(false) => {},
Err(e) => {
error!("[IRC CLIENT] (multiplex_connection) self.is_seen({}) failed: {}", event_id, e);
return Err(e)
}
}
// Static-event arrival path. EventGraph notifies
// `static_pub` only after `commit_verified_static_event`
// has durably stored the event/blob and applied the RLN
// state change. So all we need to do is bookkeeping for
// this client's seen-set.
// Mark the message as seen for this USER
if let Err(e) = self.mark_seen(&event_id).await {
error!("[IRC CLIENT] (multiplex_connection) self.mark_seen({event_id}) failed: {e}");
return Err(e)
}
}
}
}
}
/// 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!(":{SERVER_NAME} {rpl:03} {msg}"),
ReplyType::Client((nick, msg)) => format!(":{nick}!~anon@darkirc {msg}"),
ReplyType::Pong(origin) => format!(":{SERVER_NAME} PONG :{origin}"),
ReplyType::Cap(msg) => format!(":{SERVER_NAME} {msg}"),
ReplyType::Notice((src, dst, msg)) => {
format!(":{src}!~anon@darkirc NOTICE {dst} :{msg}")
}
};
debug!("[{}] <-- {r}", self.addr);
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,
args_queue: &mut VecDeque,
) -> Result