darkirc.rs 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2024 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::sync::{Arc, Mutex as SyncMutex};
  19. use darkfi::{
  20. event_graph::{
  21. self,
  22. proto::{EventPut, ProtocolEventGraph},
  23. EventGraph, EventGraphPtr,
  24. },
  25. net::{session::SESSION_DEFAULT, settings::Settings as NetSettings, P2p, P2pPtr},
  26. system::{sleep, Subscription},
  27. Error,
  28. };
  29. use darkfi_serial::{
  30. async_trait, deserialize_async, serialize_async, Encodable, SerialDecodable, SerialEncodable,
  31. };
  32. use crate::{scene::SceneGraphPtr2, ExecutorPtr};
  33. #[cfg(target_os = "android")]
  34. const EVGRDB_PATH: &str = "/data/data/darkfi.darkwallet/evgrdb/";
  35. #[cfg(target_os = "linux")]
  36. const EVGRDB_PATH: &str = "evgrdb";
  37. #[derive(Clone, Debug, SerialEncodable, SerialDecodable)]
  38. pub struct Privmsg {
  39. pub channel: String,
  40. pub nick: String,
  41. pub msg: String,
  42. }
  43. async fn relay_darkirc_events(sg: SceneGraphPtr2, ev_sub: Subscription<event_graph::Event>) {
  44. loop {
  45. let ev = ev_sub.receive().await;
  46. // Try to deserialize the `Event`'s content into a `Privmsg`
  47. let privmsg: Privmsg = match deserialize_async(ev.content()).await {
  48. Ok(v) => v,
  49. Err(e) => {
  50. error!("[IRC CLIENT] Failed deserializing incoming Privmsg event: {}", e);
  51. continue
  52. }
  53. };
  54. if privmsg.channel != "#random" {
  55. continue
  56. }
  57. info!(target: "darkirc", "ev_id={:?}", ev.id());
  58. info!(target: "darkirc", "ev: {:?}", ev);
  59. info!(target: "darkirc", "privmsg: {:?}", privmsg);
  60. info!(target: "darkirc", "");
  61. let response_fn = Box::new(|_| {});
  62. let mut arg_data = vec![];
  63. ev.timestamp.encode(&mut arg_data).unwrap();
  64. ev.id().as_bytes().encode(&mut arg_data).unwrap();
  65. privmsg.nick.encode(&mut arg_data).unwrap();
  66. privmsg.msg.encode(&mut arg_data).unwrap();
  67. let mut sg = sg.lock().await;
  68. let chatview_node = sg.lookup_node_mut("/window/view/chatty").unwrap();
  69. chatview_node.call_method("insert_line", arg_data, response_fn).unwrap();
  70. drop(sg);
  71. }
  72. }
  73. pub type DarkIrcBackendPtr = Arc<DarkIrcBackend>;
  74. struct DarkIrcData {
  75. p2p: P2pPtr,
  76. event_graph: EventGraphPtr,
  77. #[allow(dead_code)]
  78. ev_task: smol::Task<()>,
  79. db: sled::Db,
  80. }
  81. pub struct DarkIrcBackend(SyncMutex<Option<DarkIrcData>>);
  82. impl DarkIrcBackend {
  83. pub fn new() -> Arc<Self> {
  84. Arc::new(Self(SyncMutex::new(None)))
  85. }
  86. pub async fn start(&self, sg: SceneGraphPtr2, ex: ExecutorPtr) -> darkfi::Result<()> {
  87. info!(target: "darkirc", "Starting DarkIRC backend");
  88. let sled_db = sled::open(EVGRDB_PATH)?;
  89. let mut p2p_settings: NetSettings = Default::default();
  90. p2p_settings.app_version = semver::Version::parse("0.5.0").unwrap();
  91. p2p_settings.seeds.push(url::Url::parse("tcp+tls://lilith1.dark.fi:5262").unwrap());
  92. let p2p = P2p::new(p2p_settings, ex.clone()).await?;
  93. let event_graph = EventGraph::new(
  94. p2p.clone(),
  95. sled_db.clone(),
  96. std::path::PathBuf::new(),
  97. false,
  98. "darkirc_dag",
  99. 1,
  100. ex.clone(),
  101. )
  102. .await?;
  103. //self.prune_task.lock().unwrap() = Some(event_graph.prune_task.get().unwrap());
  104. info!(target: "darkirc", "Registering EventGraph P2P protocol");
  105. let event_graph_ = Arc::clone(&event_graph);
  106. let registry = p2p.protocol_registry();
  107. registry
  108. .register(SESSION_DEFAULT, move |channel, _| {
  109. let event_graph_ = event_graph_.clone();
  110. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  111. })
  112. .await;
  113. let ev_sub = event_graph.event_pub.clone().subscribe().await;
  114. let ev_task = ex.spawn(relay_darkirc_events(sg, ev_sub));
  115. info!(target: "darkirc", "Starting P2P network");
  116. p2p.clone().start().await?;
  117. info!(target: "darkirc", "Waiting for some P2P connections...");
  118. sleep(5).await;
  119. // We'll attempt to sync {sync_attempts} times
  120. let sync_attempts = 4;
  121. for i in 1..=sync_attempts {
  122. info!(target: "darkirc", "Syncing event DAG (attempt #{})", i);
  123. match event_graph.dag_sync().await {
  124. Ok(()) => break,
  125. Err(e) => {
  126. if i == sync_attempts {
  127. error!("Failed syncing DAG. Exiting.");
  128. p2p.stop().await;
  129. return Err(Error::DagSyncFailed)
  130. } else {
  131. // TODO: Maybe at this point we should prune or something?
  132. // TODO: Or maybe just tell the user to delete the DAG from FS.
  133. error!("Failed syncing DAG ({}), retrying in {}s...", e, 4);
  134. sleep(4).await;
  135. }
  136. }
  137. }
  138. }
  139. *self.0.lock().unwrap() = Some(DarkIrcData { p2p, event_graph, ev_task, db: sled_db });
  140. Ok(())
  141. }
  142. pub async fn stop(&self) {
  143. info!(target: "darkirc", "Stopping DarkIRC backend");
  144. let self_ = self.0.lock().unwrap();
  145. let Some(self_) = &*self_ else {
  146. warn!(target: "darkirc", "Backend wasn't started");
  147. return
  148. };
  149. info!(target: "darkirc", "Stopping P2P network");
  150. self_.p2p.stop().await;
  151. info!(target: "darkirc", "Stopping IRC server");
  152. let prune_task = self_.event_graph.prune_task.get().unwrap();
  153. prune_task.stop().await;
  154. info!(target: "darkirc", "Flushing event graph sled database...");
  155. let Ok(flushed_bytes) = self_.db.flush_async().await else {
  156. error!(target: "darkirc", "Flushing event graph db failed");
  157. return
  158. };
  159. info!(target: "darkirc", "Flushed {} bytes", flushed_bytes);
  160. info!(target: "darkirc", "Shut down backend successfully");
  161. }
  162. pub async fn send(&self, privmsg: Privmsg) {
  163. let (p2p, evgr) = {
  164. let self_ = self.0.lock().unwrap();
  165. let data = self_.as_ref().expect("backend wasnt started");
  166. let evgr = data.event_graph.clone();
  167. let p2p = data.p2p.clone();
  168. (p2p, evgr)
  169. };
  170. let event = event_graph::Event::new(serialize_async(&privmsg).await, &evgr).await;
  171. if let Err(e) = evgr.dag_insert(&[event.clone()]).await {
  172. error!(target: "darkirc", "Failed inserting new event to DAG: {}", e);
  173. }
  174. p2p.broadcast(&EventPut(event)).await;
  175. }
  176. }