ircd.rs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448
  1. #[macro_use]
  2. extern crate clap;
  3. use std::{
  4. io,
  5. net::{SocketAddr, TcpListener, TcpStream},
  6. sync::Arc,
  7. };
  8. use async_executor::Executor;
  9. use async_std::io::BufReader;
  10. use futures::{
  11. io::{ReadHalf, WriteHalf},
  12. AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, Future, FutureExt,
  13. };
  14. use log::{debug, error, info, warn};
  15. use simplelog::{ColorChoice, LevelFilter, TermLogger, TerminalMode};
  16. use smol::Async;
  17. use drk::{
  18. net,
  19. serial::{Decodable, Encodable},
  20. Error, Result,
  21. };
  22. /*
  23. NICK fifififif
  24. USER username 0 * :Real
  25. :behemoth 001 fifififif :Hi, welcome to IRC
  26. :behemoth 002 fifififif :Your host is behemoth, running version miniircd-2.1
  27. :behemoth 003 fifififif :This server was created sometime
  28. :behemoth 004 fifififif behemoth miniircd-2.1 o o
  29. :behemoth 251 fifififif :There are 1 users and 0 services on 1 server
  30. :behemoth 422 fifififif :MOTD File is missing
  31. JOIN #dev
  32. :fifififif!username@127.0.0.1 JOIN #dev
  33. :behemoth 331 fifififif #dev :No topic is set
  34. :behemoth 353 fifififif = #dev :fifififif
  35. :behemoth 366 fifififif #dev :End of NAMES list
  36. PRIVMSG #dev hihi
  37. */
  38. struct ServerConnection {
  39. write_stream: WriteHalf<Async<TcpStream>>,
  40. is_nick_init: bool,
  41. is_user_init: bool,
  42. is_registered: bool,
  43. nickname: String,
  44. channels: Vec<String>,
  45. }
  46. impl ServerConnection {
  47. fn new(write_stream: WriteHalf<Async<TcpStream>>) -> Self {
  48. ServerConnection {
  49. write_stream,
  50. is_nick_init: false,
  51. is_user_init: false,
  52. is_registered: false,
  53. nickname: "".to_string(),
  54. channels: vec![],
  55. }
  56. }
  57. async fn update(&mut self, line: String, p2p: net::P2pPtr) -> Result<()> {
  58. let mut tokens = line.split_ascii_whitespace();
  59. // Commands can begin with :garbage but we will reject clients doing that for now
  60. // to keep the protocol simple and focused.
  61. let command = tokens.next().ok_or(Error::MalformedPacket)?;
  62. debug!("Received command: {}", command);
  63. match command {
  64. "NICK" => {
  65. let nickname = tokens.next().ok_or(Error::MalformedPacket)?;
  66. self.is_nick_init = true;
  67. self.nickname = nickname.to_string();
  68. }
  69. "USER" => {
  70. // We can stuff any extra things like public keys in here
  71. // Ignore it for now
  72. self.is_user_init = true;
  73. }
  74. "JOIN" => {
  75. // Ignore since channels are all autojoin
  76. //let channel = tokens.next().ok_or(Error::MalformedPacket)?;
  77. //self.channels.push(channel.to_string());
  78. //let join_reply = format!(":{}!darkfi@127.0.0.1 JOIN {}\n", self.nickname,
  79. // channel); self.reply(&join_reply).await?;
  80. //self.write_stream.write_all(b":f00!f00@127.0.0.1 PRIVMSG #dev :y0\n").await?;
  81. }
  82. "PING" => {
  83. self.reply("PONG").await?;
  84. }
  85. "PRIVMSG" => {
  86. let channel = tokens.next().ok_or(Error::MalformedPacket)?;
  87. let substr_idx = line.find(':').ok_or(Error::MalformedPacket)?;
  88. if substr_idx >= line.len() {
  89. return Err(Error::MalformedPacket)
  90. }
  91. let message = &line[substr_idx + 1..];
  92. info!("Message {}: {}", channel, message);
  93. let protocol_msg = PrivMsg {
  94. nickname: self.nickname.clone(),
  95. channel: channel.to_string(),
  96. message: message.to_string(),
  97. };
  98. p2p.broadcast(protocol_msg).await?;
  99. }
  100. _ => {}
  101. }
  102. if !self.is_registered && self.is_nick_init && self.is_user_init {
  103. debug!("Initializing peer connection");
  104. let register_reply = format!(":darkfi 001 {} :Let there be dark\n", self.nickname);
  105. self.reply(&register_reply).await?;
  106. self.is_registered = true;
  107. // Auto-joins
  108. for channel in ["#dev", "#markets", "#welcome"] {
  109. let join_reply = format!(":{}!darkfi@127.0.0.1 JOIN {}\n", self.nickname, channel);
  110. self.reply(&join_reply).await?;
  111. }
  112. }
  113. Ok(())
  114. }
  115. async fn reply(&mut self, message: &str) -> Result<()> {
  116. self.write_stream.write_all(message.as_bytes()).await?;
  117. debug!("Sent {}", message);
  118. Ok(())
  119. }
  120. }
  121. #[derive(Debug, Clone)]
  122. struct PrivMsg {
  123. nickname: String,
  124. channel: String,
  125. message: String,
  126. }
  127. impl net::Message for PrivMsg {
  128. fn name() -> &'static str {
  129. "privmsg"
  130. }
  131. }
  132. impl Encodable for PrivMsg {
  133. fn encode<S: io::Write>(&self, mut s: S) -> Result<usize> {
  134. let mut len = 0;
  135. len += self.nickname.encode(&mut s)?;
  136. len += self.channel.encode(&mut s)?;
  137. len += self.message.encode(&mut s)?;
  138. Ok(len)
  139. }
  140. }
  141. impl Decodable for PrivMsg {
  142. fn decode<D: io::Read>(mut d: D) -> Result<Self> {
  143. Ok(Self {
  144. nickname: Decodable::decode(&mut d)?,
  145. channel: Decodable::decode(&mut d)?,
  146. message: Decodable::decode(&mut d)?,
  147. })
  148. }
  149. }
  150. async fn process(
  151. recvr: async_channel::Receiver<Arc<PrivMsg>>,
  152. stream: Async<TcpStream>,
  153. peer_addr: SocketAddr,
  154. p2p: net::P2pPtr,
  155. executor: Arc<Executor<'_>>,
  156. ) -> Result<()> {
  157. //stream.write_all(b":behemoth 001 fifififif :Hi, welcome to IRC").await;
  158. //stream.write_all(b"NICK username");
  159. //stream.write_all(b"USER username 0 * :username");
  160. //stream.write_all(b"JOIN #dev");
  161. //stream.write_all(b"PRIVMSG #dev y0");
  162. // PING :behemoth
  163. let (reader, writer) = stream.split();
  164. let mut reader = BufReader::new(reader);
  165. let mut connection = ServerConnection::new(writer);
  166. loop {
  167. let mut line = String::new();
  168. futures::select! {
  169. privmsg = recvr.recv().fuse() => {
  170. let privmsg = privmsg.expect("internal message queue error");
  171. debug!("ABOUT TO SEND {:?}", privmsg);
  172. let irc_msg = format!(
  173. ":{}!darkfi@127.0.0.1 PRIVMSG {} :{}\n",
  174. privmsg.nickname,
  175. privmsg.channel,
  176. privmsg.message
  177. );
  178. connection.reply(&irc_msg).await?;
  179. }
  180. err = reader.read_line(&mut line).fuse() => {
  181. if let Err(err) = err {
  182. warn!("Read line error. Closing stream for {}: {}", peer_addr, err);
  183. return Ok(())
  184. }
  185. process_user_input(line, peer_addr, &mut connection, p2p.clone()).await;
  186. }
  187. };
  188. }
  189. }
  190. async fn process_user_input(
  191. mut line: String,
  192. peer_addr: SocketAddr,
  193. connection: &mut ServerConnection,
  194. p2p: net::P2pPtr,
  195. ) {
  196. if line.len() == 0 {
  197. warn!("Received empty line from {}. Closing connection.", peer_addr);
  198. return
  199. }
  200. assert!(&line[(line.len() - 1)..] == "\n");
  201. // Remove the \n character
  202. line.pop();
  203. debug!("Received '{}' from {}", line, peer_addr);
  204. if let Err(err) = connection.update(line, p2p.clone()).await {
  205. warn!("Connection error: {} for {}", err, peer_addr);
  206. return
  207. }
  208. }
  209. struct ProtocolPrivMsg {
  210. notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
  211. privmsg_sub: net::MessageSubscription<PrivMsg>,
  212. jobsman: net::ProtocolJobsManagerPtr,
  213. }
  214. impl ProtocolPrivMsg {
  215. async fn new(
  216. channel: net::ChannelPtr,
  217. notify_queue_sender: async_channel::Sender<Arc<PrivMsg>>,
  218. ) -> Arc<Self> {
  219. let message_subsytem = channel.get_message_subsystem();
  220. message_subsytem.add_dispatch::<PrivMsg>().await;
  221. debug!("ADDED DISPATCH");
  222. let privmsg_sub =
  223. channel.subscribe_msg::<PrivMsg>().await.expect("Missing PrivMsg dispatcher!");
  224. Arc::new(Self {
  225. notify_queue_sender,
  226. privmsg_sub,
  227. jobsman: net::ProtocolJobsManager::new("PrivMsgProtocol", channel),
  228. })
  229. }
  230. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) {
  231. debug!(target: "ircd", "ProtocolPrivMsg::start() [START]");
  232. self.jobsman.clone().start(executor.clone());
  233. self.jobsman.clone().spawn(self.clone().handle_receive_privmsg(), executor.clone()).await;
  234. debug!(target: "ircd", "ProtocolPrivMsg::start() [END]");
  235. }
  236. async fn handle_receive_privmsg(self: Arc<Self>) -> Result<()> {
  237. debug!(target: "ircd", "ProtocolAddress::handle_receive_privmsg() [START]");
  238. loop {
  239. let privmsg = self.privmsg_sub.receive().await?;
  240. debug!(
  241. target: "ircd",
  242. "ProtocolPrivMsg::handle_receive_privmsg() received {:?}",
  243. privmsg
  244. );
  245. self.notify_queue_sender.send(privmsg).await.expect("notify_queue_sender send failed!");
  246. }
  247. }
  248. }
  249. async fn channel_loop(
  250. p2p: net::P2pPtr,
  251. sender: async_channel::Sender<Arc<PrivMsg>>,
  252. executor: Arc<Executor<'_>>,
  253. ) -> Result<()> {
  254. debug!("CHANNEL SUBS LOOP");
  255. let new_channel_sub = p2p.subscribe_channel().await;
  256. loop {
  257. let channel = new_channel_sub.receive().await?;
  258. debug!("NEWCHANNEL");
  259. let protocol_privmsg = ProtocolPrivMsg::new(channel, sender.clone()).await;
  260. protocol_privmsg.start(executor.clone()).await;
  261. }
  262. }
  263. async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<()> {
  264. let listener = match Async::<TcpListener>::bind(options.irc_accept_addr) {
  265. Ok(listener) => listener,
  266. Err(err) => {
  267. error!("Bind listener failed: {}", err);
  268. return Err(Error::OperationFailed)
  269. }
  270. };
  271. let local_addr = match listener.get_ref().local_addr() {
  272. Ok(addr) => addr,
  273. Err(err) => {
  274. error!("Failed to get local address: {}", err);
  275. return Err(Error::OperationFailed)
  276. }
  277. };
  278. info!("Listening on {}", local_addr);
  279. let p2p = net::P2p::new(options.network_settings);
  280. // Performs seed session
  281. p2p.clone().start(executor.clone()).await?;
  282. // Actual main p2p session
  283. let ex2 = executor.clone();
  284. let p2p2 = p2p.clone();
  285. executor
  286. .spawn(async move {
  287. if let Err(err) = p2p2.run(ex2).await {
  288. error!("Error: p2p run failed {}", err);
  289. }
  290. })
  291. .detach();
  292. let (sender, recvr) = async_channel::unbounded();
  293. // todo: be careful of zombie processes
  294. // for now we just want things to work
  295. executor.spawn(channel_loop(p2p.clone(), sender, executor.clone())).detach();
  296. loop {
  297. let (stream, peer_addr) = match listener.accept().await {
  298. Ok((s, a)) => (s, a),
  299. Err(err) => {
  300. error!("Error listening for connections: {}", err);
  301. return Err(Error::ServiceStopped)
  302. }
  303. };
  304. info!("Accepted client: {}", peer_addr);
  305. let p2p2 = p2p.clone();
  306. let ex2 = executor.clone();
  307. executor.spawn(process(recvr.clone(), stream, peer_addr, p2p2, ex2)).detach();
  308. }
  309. }
  310. struct ProgramOptions {
  311. network_settings: net::Settings,
  312. log_path: Box<std::path::PathBuf>,
  313. irc_accept_addr: SocketAddr,
  314. }
  315. impl ProgramOptions {
  316. fn load() -> Result<ProgramOptions> {
  317. let app = clap_app!(dfi =>
  318. (version: "0.1.0")
  319. (author: "Amir Taaki <amir@dyne.org>")
  320. (about: "Dark node")
  321. (@arg ACCEPT: -a --accept +takes_value "Accept address")
  322. (@arg SEED_NODES: -s --seeds +takes_value ... "Seed nodes")
  323. (@arg CONNECTS: -c --connect +takes_value ... "Manual connections")
  324. (@arg CONNECT_SLOTS: --slots +takes_value "Connection slots")
  325. (@arg LOG_PATH: --log +takes_value "Logfile path")
  326. (@arg IRC_ACCEPT: -r --irc +takes_value "IRC accept address")
  327. )
  328. .get_matches();
  329. let accept_addr = if let Some(accept_addr) = app.value_of("ACCEPT") {
  330. Some(accept_addr.parse()?)
  331. } else {
  332. None
  333. };
  334. let mut seed_addrs: Vec<SocketAddr> = vec![];
  335. if let Some(seeds) = app.values_of("SEED_NODES") {
  336. for seed in seeds {
  337. seed_addrs.push(seed.parse()?);
  338. }
  339. }
  340. let mut manual_connects: Vec<SocketAddr> = vec![];
  341. if let Some(connections) = app.values_of("CONNECTS") {
  342. for connect in connections {
  343. manual_connects.push(connect.parse()?);
  344. }
  345. }
  346. let connection_slots = if let Some(connection_slots) = app.value_of("CONNECT_SLOTS") {
  347. connection_slots.parse()?
  348. } else {
  349. 0
  350. };
  351. let log_path = Box::new(
  352. if let Some(log_path) = app.value_of("LOG_PATH") {
  353. std::path::Path::new(log_path)
  354. } else {
  355. std::path::Path::new("/tmp/darkfid.log")
  356. }
  357. .to_path_buf(),
  358. );
  359. let irc_accept_addr = if let Some(accept_addr) = app.value_of("IRC_ACCEPT") {
  360. accept_addr.parse()?
  361. } else {
  362. ([127, 0, 0, 1], 6667).into()
  363. };
  364. Ok(ProgramOptions {
  365. network_settings: net::Settings {
  366. inbound: accept_addr,
  367. outbound_connections: connection_slots,
  368. external_addr: accept_addr,
  369. peers: manual_connects,
  370. seeds: seed_addrs,
  371. ..Default::default()
  372. },
  373. log_path,
  374. irc_accept_addr,
  375. })
  376. }
  377. }
  378. fn main() -> Result<()> {
  379. TermLogger::init(
  380. LevelFilter::Debug,
  381. simplelog::Config::default(),
  382. TerminalMode::Mixed,
  383. ColorChoice::Auto,
  384. )?;
  385. let options = ProgramOptions::load()?;
  386. let ex = Arc::new(Executor::new());
  387. smol::block_on(ex.run(start(ex.clone(), options)))
  388. }