darkirc2.rs 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314
  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 async_lock::Mutex as AsyncMutex;
  19. use darkfi::{
  20. event_graph::{self},
  21. net::transport::{Dialer, PtStream},
  22. system::ExecutorPtr,
  23. util::path::expand_path,
  24. Error, Result,
  25. };
  26. use darkfi_serial::{
  27. async_trait, deserialize_async_partial, serialize_async, AsyncDecodable, AsyncEncodable,
  28. Encodable, SerialDecodable, SerialEncodable,
  29. };
  30. use evgrd::{
  31. FetchEventsMessage, LocalEventGraph, LocalEventGraphPtr, VersionMessage, MSG_EVENT,
  32. MSG_FETCHEVENTS, MSG_SENDEVENT,
  33. };
  34. use log::{error, info};
  35. use sled_overlay::sled;
  36. use smol::{
  37. fs,
  38. io::{ReadHalf, WriteHalf},
  39. };
  40. use std::{
  41. sync::{
  42. atomic::{AtomicBool, Ordering},
  43. Arc, Mutex as SyncMutex, Weak,
  44. },
  45. time::UNIX_EPOCH,
  46. };
  47. use url::Url;
  48. use crate::{
  49. prop::{PropertyBool, PropertyStr, Role},
  50. scene::{SceneNodePtr, Slot},
  51. };
  52. #[cfg(target_os = "android")]
  53. const EVGRDB_PATH: &str = "/data/data/darkfi.darkwallet/evgr/";
  54. #[cfg(target_os = "linux")]
  55. const EVGRDB_PATH: &str = "~/.local/darkfi/darkwallet/evgr/";
  56. #[derive(Clone, Debug, SerialEncodable, SerialDecodable)]
  57. pub struct Privmsg {
  58. pub channel: String,
  59. pub nick: String,
  60. pub msg: String,
  61. }
  62. impl Privmsg {
  63. pub fn new(channel: String, nick: String, msg: String) -> Self {
  64. Self { channel, nick, msg }
  65. }
  66. }
  67. pub type LocalDarkIRCPtr = Arc<LocalDarkIRC>;
  68. pub struct LocalDarkIRC {
  69. is_connected: AtomicBool,
  70. /// The reading half of the transport stream
  71. reader: AsyncMutex<Option<ReadHalf<Box<dyn PtStream>>>>,
  72. /// The writing half of the transport stream
  73. writer: AsyncMutex<Option<WriteHalf<Box<dyn PtStream>>>>,
  74. evgr: LocalEventGraphPtr,
  75. tasks: SyncMutex<Vec<smol::Task<()>>>,
  76. chatview_node: SceneNodePtr,
  77. sendbtn_node: SceneNodePtr,
  78. editbox_node: SceneNodePtr,
  79. editbox_text: PropertyStr,
  80. upgrade_popup_is_visible: PropertyBool,
  81. }
  82. impl LocalDarkIRC {
  83. pub async fn new(sg_root: SceneNodePtr, ex: ExecutorPtr) -> Result<Arc<Self>> {
  84. let chatview_node = sg_root.clone().lookup_node("/window/view/chatty").unwrap();
  85. let sendbtn_node = sg_root.clone().lookup_node("/window/view/send_btn").unwrap();
  86. let editbox_node = sg_root.clone().lookup_node("/window/view/editz").unwrap();
  87. let editbox_text = PropertyStr::wrap(&editbox_node, Role::App, "text", 0).unwrap();
  88. let upgrade_popup_node = sg_root.clone().lookup_node("/window/view/upgrade_popup").unwrap();
  89. let upgrade_popup_is_visible =
  90. PropertyBool::wrap(&upgrade_popup_node, Role::App, "is_visible", 0).unwrap();
  91. info!(target: "darkirc", "Instantiating DarkIRC event DAG");
  92. let datastore = expand_path(EVGRDB_PATH)?;
  93. fs::create_dir_all(&datastore).await?;
  94. let sled_db = sled::open(datastore)?;
  95. let evgr = LocalEventGraph::new(sled_db.clone(), "darkirc_dag", 1, ex.clone()).await?;
  96. Ok(Arc::new(Self {
  97. is_connected: AtomicBool::new(false),
  98. reader: AsyncMutex::new(None),
  99. writer: AsyncMutex::new(None),
  100. evgr,
  101. tasks: SyncMutex::new(vec![]),
  102. chatview_node,
  103. sendbtn_node,
  104. editbox_node,
  105. editbox_text,
  106. upgrade_popup_is_visible,
  107. }))
  108. }
  109. async fn reconnect(&self) -> Result<()> {
  110. let endpoint = "tcp://127.0.0.1:5588";
  111. let endpoint = "tcp://192.168.1.38:5588";
  112. let endpoint = Url::parse(endpoint)?;
  113. let dialer = Dialer::new(endpoint.clone(), None).await?;
  114. let timeout = std::time::Duration::from_secs(60);
  115. let stream = dialer.dial(Some(timeout)).await?;
  116. info!(target: "darkirc", "Connected to the backend: {endpoint}");
  117. let (reader, writer) = smol::io::split(stream);
  118. *self.writer.lock().await = Some(writer);
  119. *self.reader.lock().await = Some(reader);
  120. Ok(())
  121. }
  122. pub async fn start(self: Arc<Self>, ex: ExecutorPtr) -> Result<()> {
  123. debug!(target: "darkirc", "LocalDarkIRC::start()");
  124. //self.reconnect().await?;
  125. self.version_exchange().await?;
  126. let me = Arc::downgrade(&self);
  127. let recv_task = ex.spawn(async move {
  128. while let Some(self_) = me.upgrade() {
  129. self_.receive_msg().await.unwrap();
  130. }
  131. error!(target: "darkirc", "Closing DarkIRC receive loop");
  132. });
  133. let (slot, recvr) = Slot::new("send_button_clicked");
  134. self.sendbtn_node.register("click", slot).unwrap();
  135. let me = Arc::downgrade(&self);
  136. let send_task = ex.spawn(async move {
  137. while let Some(self_) = me.upgrade() {
  138. let Ok(_) = recvr.recv().await else {
  139. error!(target: "ui::win", "Button click recvr closed");
  140. break
  141. };
  142. self_.handle_send().await;
  143. }
  144. });
  145. let (slot, recvr) = Slot::new("enter_pressed");
  146. self.editbox_node.register("enter_pressed", slot).unwrap();
  147. let me = Arc::downgrade(&self);
  148. let enter_task = ex.spawn(async move {
  149. while let Some(self_) = me.upgrade() {
  150. let Ok(_) = recvr.recv().await else {
  151. error!(target: "ui::win", "EditBox enter_pressed recvr closed");
  152. break
  153. };
  154. self_.handle_send().await;
  155. }
  156. });
  157. let mut tasks = self.tasks.lock().unwrap();
  158. assert!(tasks.is_empty());
  159. *tasks = vec![recv_task, send_task, enter_task];
  160. Ok(())
  161. }
  162. async fn version_exchange(&self) -> Result<()> {
  163. if !self.is_connected.load(Ordering::Relaxed) {
  164. self.reconnect().await?;
  165. }
  166. let mut writer = self.writer.lock().await;
  167. let mut reader = self.reader.lock().await;
  168. let writer = writer.as_mut().unwrap();
  169. let reader = reader.as_mut().unwrap();
  170. let version = VersionMessage::new();
  171. version.encode_async(writer).await?;
  172. let server_version = VersionMessage::decode_async(reader).await?;
  173. info!(target: "darkirc", "Backend server version: {}", server_version.protocol_version);
  174. if server_version.protocol_version > evgrd::PROTOCOL_VERSION {
  175. self.upgrade_popup_is_visible.set(true);
  176. }
  177. let unref_tips = self.evgr.unreferenced_tips.read().await.clone();
  178. let fetchevs = FetchEventsMessage::new(unref_tips);
  179. MSG_FETCHEVENTS.encode_async(writer).await?;
  180. fetchevs.encode_async(writer).await?;
  181. self.is_connected.store(true, Ordering::Relaxed);
  182. Ok(())
  183. }
  184. async fn send_msg(&self, timestamp: u64, msg: Privmsg) -> Result<()> {
  185. if !self.is_connected.load(Ordering::Relaxed) {
  186. debug!(target: "darkirc", "send_msg: not connected, reconnecting...");
  187. self.reconnect().await?;
  188. }
  189. let mut writer = self.writer.lock().await;
  190. let writer = writer.as_mut().unwrap();
  191. MSG_SENDEVENT.encode_async(writer).await?;
  192. timestamp.encode_async(writer).await?;
  193. let content: Vec<u8> = serialize_async(&msg).await;
  194. content.encode_async(writer).await?;
  195. Ok(())
  196. }
  197. async fn receive_msg(&self) -> Result<()> {
  198. if !self.is_connected.load(Ordering::Relaxed) {
  199. self.reconnect().await?;
  200. }
  201. debug!(target: "darkirc", "Receiving message...");
  202. let mut reader = self.reader.lock().await;
  203. let reader = reader.as_mut().unwrap();
  204. let msg_type = u8::decode_async(reader).await?;
  205. debug!(target: "darkirc", "Received: {msg_type:?}");
  206. if msg_type != MSG_EVENT {
  207. error!(target: "darkirc", "Received invalid msg_type: {msg_type}");
  208. //return Err(Error::MalformedPacket)
  209. return Ok(())
  210. }
  211. let ev = event_graph::Event::decode_async(reader).await?;
  212. let genesis_timestamp = self.evgr.current_genesis.read().await.clone().timestamp;
  213. let ev_id = ev.id();
  214. if self.evgr.dag.contains_key(ev_id.as_bytes()).unwrap() ||
  215. !ev.validate(&self.evgr.dag, genesis_timestamp, self.evgr.days_rotation, None)
  216. .await?
  217. {
  218. error!(target: "darkirc", "Event is invalid! {ev:?}");
  219. return Ok(())
  220. }
  221. debug!(target: "darkirc", "got {ev:?}");
  222. self.evgr.dag_insert(&[ev.clone()]).await.unwrap();
  223. let privmsg: Privmsg = match deserialize_async_partial(ev.content()).await {
  224. Ok((v, _)) => v,
  225. Err(e) => {
  226. error!(target: "darkirc", "Failed deserializing incoming Privmsg event: {e}");
  227. return Ok(())
  228. }
  229. };
  230. debug!(target: "darkirc", "privmsg: {privmsg:?}");
  231. let mut timest = ev.timestamp;
  232. if timest < 6047051717 {
  233. timest *= 1000;
  234. }
  235. if privmsg.channel != "#random" {
  236. return Ok(())
  237. }
  238. let mut arg_data = vec![];
  239. timest.encode_async(&mut arg_data).await.unwrap();
  240. ev.id().as_bytes().encode_async(&mut arg_data).await.unwrap();
  241. privmsg.nick.encode_async(&mut arg_data).await.unwrap();
  242. privmsg.msg.encode_async(&mut arg_data).await.unwrap();
  243. self.chatview_node.call_method("insert_line", arg_data).await.unwrap();
  244. Ok(())
  245. }
  246. async fn handle_send(&self) {
  247. // Get text from editbox
  248. let text = self.editbox_text.get();
  249. // Clear editbox
  250. self.editbox_text.set("");
  251. // Send text to channel
  252. debug!(target: "darkirc", "Sending privmsg: {text}");
  253. let msg = Privmsg::new("#random".to_string(), "anon".to_string(), text);
  254. let timestamp = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  255. self.send_msg(timestamp, msg).await.unwrap();
  256. }
  257. }