darkirc.rs 7.0 KB

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