server.rs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2023 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::{fs::File, sync::Arc};
  19. use async_rustls::{rustls, TlsAcceptor};
  20. use log::{error, info};
  21. use smol::{
  22. io::{self, AsyncRead, AsyncWrite, BufReader},
  23. lock::Mutex,
  24. net::{SocketAddr, TcpListener},
  25. };
  26. use darkfi::{
  27. event_graph::{
  28. model::{Event, EventId, ModelPtr},
  29. protocol_event::{Seen, SeenPtr},
  30. view::ViewPtr,
  31. },
  32. net::P2pPtr,
  33. system::{StoppableTask, SubscriberPtr},
  34. util::{path::expand_path, time::Timestamp},
  35. Error, Result,
  36. };
  37. use super::{ClientSubMsg, IrcClient, IrcConfig, NotifierMsg};
  38. use crate::{settings::Args, PrivMsgEvent};
  39. mod nickserv;
  40. use nickserv::NickServ;
  41. const NICK_NICKSERV: &str = "nickserv";
  42. pub struct IrcServer {
  43. settings: Args,
  44. p2p: P2pPtr,
  45. model: ModelPtr<PrivMsgEvent>,
  46. view: ViewPtr<PrivMsgEvent>,
  47. clients_subscriptions: SubscriberPtr<ClientSubMsg>,
  48. seen: SeenPtr<EventId>,
  49. missed_events: Arc<Mutex<Vec<Event<PrivMsgEvent>>>>,
  50. /// nickserv service
  51. pub nickserv: NickServ,
  52. }
  53. impl IrcServer {
  54. pub async fn new(
  55. settings: Args,
  56. p2p: P2pPtr,
  57. model: ModelPtr<PrivMsgEvent>,
  58. view: ViewPtr<PrivMsgEvent>,
  59. clients_subscriptions: SubscriberPtr<ClientSubMsg>,
  60. ) -> Result<Self> {
  61. let seen = Seen::new();
  62. let missed_events = Arc::new(Mutex::new(vec![]));
  63. Ok(Self {
  64. settings,
  65. p2p,
  66. model,
  67. view,
  68. clients_subscriptions,
  69. seen,
  70. missed_events,
  71. nickserv: NickServ::default(),
  72. })
  73. }
  74. pub async fn start(&self, executor: Arc<smol::Executor<'_>>) -> Result<()> {
  75. let (msg_notifier, msg_recv) = smol::channel::unbounded();
  76. // Listen to msgs from clients
  77. StoppableTask::new().start(
  78. Self::listen_to_msgs(
  79. self.p2p.clone(),
  80. self.model.clone(),
  81. self.seen.clone(),
  82. msg_recv,
  83. self.missed_events.clone(),
  84. self.clients_subscriptions.clone(),
  85. ),
  86. |res| async {
  87. match res {
  88. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  89. Err(e) => error!(target: "darkirc::irc::server::start", "Failed starting listen to msgs: {}", e),
  90. }
  91. },
  92. Error::DetachedTaskStopped,
  93. executor.clone(),
  94. );
  95. // Listen to msgs from View
  96. StoppableTask::new().start(
  97. Self::listen_to_view(
  98. self.view.clone(),
  99. self.seen.clone(),
  100. self.missed_events.clone(),
  101. self.clients_subscriptions.clone(),
  102. ),
  103. |res| async {
  104. match res {
  105. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  106. Err(e) => error!(target: "darkirc::irc::server::start", "Failed starting listen to view: {}", e),
  107. }
  108. },
  109. Error::DetachedTaskStopped,
  110. executor.clone(),
  111. );
  112. // Start listening for new connections
  113. self.listen(msg_notifier, executor).await?;
  114. Ok(())
  115. }
  116. async fn listen_to_view(
  117. view: ViewPtr<PrivMsgEvent>,
  118. seen: SeenPtr<EventId>,
  119. missed_events: Arc<Mutex<Vec<Event<PrivMsgEvent>>>>,
  120. clients_subscriptions: SubscriberPtr<ClientSubMsg>,
  121. ) -> Result<()> {
  122. loop {
  123. let event = view.lock().await.process().await?;
  124. if !seen.push(&event.hash()).await {
  125. continue
  126. }
  127. missed_events.lock().await.push(event.clone());
  128. let msg = event.action.clone();
  129. clients_subscriptions.notify(ClientSubMsg::Privmsg(msg)).await;
  130. }
  131. }
  132. /// Start listening to msgs from irc clients
  133. pub async fn listen_to_msgs(
  134. p2p: P2pPtr,
  135. model: ModelPtr<PrivMsgEvent>,
  136. seen: SeenPtr<EventId>,
  137. recv: smol::channel::Receiver<(NotifierMsg, usize)>,
  138. missed_events: Arc<Mutex<Vec<Event<PrivMsgEvent>>>>,
  139. clients_subscriptions: SubscriberPtr<ClientSubMsg>,
  140. ) -> Result<()> {
  141. loop {
  142. let (msg, subscription_id) = recv.recv().await?;
  143. match msg {
  144. NotifierMsg::Privmsg(msg) => {
  145. // First check if we're communicating with any services.
  146. // If not, then we proceed with behaving like it's a normal
  147. // message.
  148. // TODO: This needs to be protected from adversaries doing
  149. // remote execution.
  150. #[allow(clippy::single_match)]
  151. match msg.target.to_lowercase().as_str() {
  152. NICK_NICKSERV => {
  153. //self.nickserv.act(msg);
  154. continue
  155. }
  156. _ => {} // pass
  157. }
  158. let event = Event {
  159. previous_event_hash: model.lock().await.get_head_hash()?,
  160. action: msg.clone(),
  161. timestamp: Timestamp::current_time(),
  162. };
  163. // Since this will be added to the View directly, other clients connected to irc
  164. // server must get informed about this new msg
  165. clients_subscriptions
  166. .notify_with_exclude(ClientSubMsg::Privmsg(msg), &[subscription_id])
  167. .await;
  168. if !seen.push(&event.hash()).await {
  169. continue
  170. }
  171. missed_events.lock().await.push(event.clone());
  172. p2p.broadcast(&event).await;
  173. }
  174. NotifierMsg::UpdateConfig => {
  175. //
  176. // load and parse the new settings from configuration file and pass it to all
  177. // irc clients
  178. //
  179. // let new_config = IrcConfig::new()?;
  180. // clients_subscriptions.notify(ClientSubMsg::Config(new_config)).await;
  181. }
  182. }
  183. }
  184. }
  185. /// Start listening to new connections from irc clients
  186. pub async fn listen(
  187. &self,
  188. notifier: smol::channel::Sender<(NotifierMsg, usize)>,
  189. executor: Arc<smol::Executor<'_>>,
  190. ) -> Result<()> {
  191. let (listener, acceptor) = self.setup_listener().await?;
  192. info!("[IRC SERVER] listening on {}", self.settings.irc_listen);
  193. loop {
  194. let (stream, peer_addr) = match listener.accept().await {
  195. Ok((s, a)) => (s, a),
  196. Err(e) => {
  197. error!("[IRC SERVER] Failed accepting new connections: {}", e);
  198. continue
  199. }
  200. };
  201. let result = if let Some(acceptor) = acceptor.clone() {
  202. // TLS connection
  203. let stream = match acceptor.accept(stream).await {
  204. Ok(s) => s,
  205. Err(e) => {
  206. error!("[IRC SERVER] Failed accepting TLS connection: {}", e);
  207. continue
  208. }
  209. };
  210. self.process_connection(stream, peer_addr, notifier.clone(), executor.clone()).await
  211. } else {
  212. // TCP connection
  213. self.process_connection(stream, peer_addr, notifier.clone(), executor.clone()).await
  214. };
  215. if let Err(e) = result {
  216. error!("[IRC SERVER] Failed processing connection {}: {}", peer_addr, e);
  217. continue
  218. };
  219. info!("[IRC SERVER] Accept new connection: {}", peer_addr);
  220. }
  221. }
  222. /// On every new connection create new IrcClient
  223. async fn process_connection<C: AsyncRead + AsyncWrite + Send + Unpin + 'static>(
  224. &self,
  225. stream: C,
  226. peer_addr: SocketAddr,
  227. notifier: smol::channel::Sender<(NotifierMsg, usize)>,
  228. executor: Arc<smol::Executor<'_>>,
  229. ) -> Result<()> {
  230. let (reader, writer) = io::split(stream);
  231. let reader = BufReader::new(reader);
  232. // Subscription for the new client
  233. let client_subscription = self.clients_subscriptions.clone().subscribe().await;
  234. // new irc configuration
  235. let irc_config = IrcConfig::new(&self.settings)?;
  236. // New irc client
  237. let mut client = IrcClient::new(
  238. writer,
  239. reader,
  240. peer_addr,
  241. irc_config,
  242. notifier,
  243. client_subscription,
  244. self.missed_events.clone(),
  245. );
  246. // Start listening and detach
  247. StoppableTask::new().start(
  248. // Weird hack to prevent lifetimes hell
  249. async move {client.listen().await; Ok(())},
  250. |res| async {
  251. match res {
  252. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  253. Err(e) => error!(target: "darkirc::irc::server::process_connection", "Failed starting client listen: {}", e),
  254. }
  255. },
  256. Error::DetachedTaskStopped,
  257. executor,
  258. );
  259. Ok(())
  260. }
  261. /// Setup a listener for irc server
  262. async fn setup_listener(&self) -> Result<(TcpListener, Option<TlsAcceptor>)> {
  263. let listenaddr = self.settings.irc_listen.socket_addrs(|| None)?[0];
  264. let listener = TcpListener::bind(listenaddr).await?;
  265. let acceptor = match self.settings.irc_listen.scheme() {
  266. "tcp+tls" => {
  267. // openssl genpkey -algorithm ED25519 > example.com.key
  268. // openssl req -new -out example.com.csr -key example.com.key
  269. // openssl x509 -req -days 700 -in example.com.csr -signkey example.com.key -out example.com.crt
  270. if self.settings.irc_tls_secret.is_none() || self.settings.irc_tls_cert.is_none() {
  271. error!("[IRC SERVER] To listen using TLS, please set irc_tls_secret and irc_tls_cert in your config file.");
  272. return Err(Error::KeypairPathNotFound)
  273. }
  274. let file =
  275. File::open(expand_path(self.settings.irc_tls_secret.as_ref().unwrap())?)?;
  276. let mut reader = std::io::BufReader::new(file);
  277. let secret = &rustls_pemfile::pkcs8_private_keys(&mut reader)?[0];
  278. let secret = rustls::PrivateKey(secret.clone());
  279. let file = File::open(expand_path(self.settings.irc_tls_cert.as_ref().unwrap())?)?;
  280. let mut reader = std::io::BufReader::new(file);
  281. let certificate = &rustls_pemfile::certs(&mut reader)?[0];
  282. let certificate = rustls::Certificate(certificate.clone());
  283. let config = rustls::ServerConfig::builder()
  284. .with_safe_defaults()
  285. .with_no_client_auth()
  286. .with_single_cert(vec![certificate], secret)?;
  287. let acceptor = TlsAcceptor::from(Arc::new(config));
  288. Some(acceptor)
  289. }
  290. _ => None,
  291. };
  292. Ok((listener, acceptor))
  293. }
  294. }