darkirc2.rs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405
  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_channel::{Receiver, Sender};
  19. use async_lock::Mutex as AsyncMutex;
  20. use darkfi::{
  21. event_graph::{self},
  22. net::transport::{Dialer, PtStream},
  23. system::{sleep, ExecutorPtr},
  24. util::path::expand_path,
  25. Error, Result,
  26. };
  27. use darkfi_serial::{
  28. async_trait, deserialize_async_partial, serialize_async, AsyncDecodable, AsyncEncodable,
  29. Encodable, SerialDecodable, SerialEncodable,
  30. };
  31. use evgrd::{
  32. FetchEventsMessage, LocalEventGraph, LocalEventGraphPtr, VersionMessage, MSG_EVENT,
  33. MSG_FETCHEVENTS, MSG_SENDEVENT,
  34. };
  35. use futures::{select, AsyncWriteExt, FutureExt};
  36. use log::{error, info};
  37. use sled_overlay::sled;
  38. use smol::{
  39. fs,
  40. io::{ReadHalf, WriteHalf},
  41. };
  42. use std::{
  43. sync::{
  44. atomic::{AtomicBool, Ordering},
  45. Arc, Mutex as SyncMutex, Weak,
  46. },
  47. time::UNIX_EPOCH,
  48. };
  49. use url::Url;
  50. use crate::{
  51. prop::{PropertyBool, PropertyFloat32, PropertyStr, Role},
  52. scene::{SceneNodePtr, Slot},
  53. ui::chatview::MessageId,
  54. };
  55. #[cfg(target_os = "android")]
  56. const EVGRDB_PATH: &str = "/data/data/darkfi.darkwallet/evgr/";
  57. #[cfg(not(target_os = "android"))]
  58. const EVGRDB_PATH: &str = "~/.local/darkfi/darkwallet/evgr/";
  59. //const ENDPOINT: &str = "tcp://agorism.dev:25588";
  60. const ENDPOINT: &str = "tor://obbc5rgtsqtscnph7yxrbsgsm5axbppfn552yr5lrrd2ocgkdcsjcnyd.onion:25589";
  61. const CHANNEL: &str = "#random";
  62. /// Due to drift between different machine's clocks, if the message timestamp is recent
  63. /// then we will just correct it to the current time so messages appear sequential in the UI.
  64. const RECENT_TIME_DIST: u64 = 10_000;
  65. #[derive(Clone, Debug, SerialEncodable, SerialDecodable)]
  66. pub struct Privmsg {
  67. pub channel: String,
  68. pub nick: String,
  69. pub msg: String,
  70. }
  71. impl Privmsg {
  72. pub fn new(channel: String, nick: String, msg: String) -> Self {
  73. Self { channel, nick, msg }
  74. }
  75. pub fn msg_id(&self, timest: u64) -> MessageId {
  76. let mut hasher = blake3::Hasher::new();
  77. timest.encode(&mut hasher).unwrap();
  78. self.channel.encode(&mut hasher).unwrap();
  79. self.nick.encode(&mut hasher).unwrap();
  80. self.msg.encode(&mut hasher).unwrap();
  81. MessageId(hasher.finalize().into())
  82. }
  83. }
  84. pub type LocalDarkIRCPtr = Arc<LocalDarkIRC>;
  85. pub struct LocalDarkIRC {
  86. stream: AsyncMutex<Option<Box<dyn PtStream>>>,
  87. evgr: LocalEventGraphPtr,
  88. tasks: SyncMutex<Vec<smol::Task<()>>>,
  89. send_sender: Sender<(u64, Privmsg)>,
  90. send_recvr: Receiver<(u64, Privmsg)>,
  91. chatview_node: SceneNodePtr,
  92. sendbtn_node: SceneNodePtr,
  93. editbox_node: SceneNodePtr,
  94. editbox_text: PropertyStr,
  95. chatview_scroll: PropertyFloat32,
  96. upgrade_popup_is_visible: PropertyBool,
  97. seen_msgs: SyncMutex<Vec<MessageId>>,
  98. }
  99. impl LocalDarkIRC {
  100. pub async fn new(sg_root: SceneNodePtr, ex: ExecutorPtr) -> Result<Arc<Self>> {
  101. let chatview_node = sg_root.clone().lookup_node("/window/view/chatty").unwrap();
  102. let sendbtn_node = sg_root.clone().lookup_node("/window/view/send_btn").unwrap();
  103. let editbox_node = sg_root.clone().lookup_node("/window/view/editz").unwrap();
  104. let editbox_text = PropertyStr::wrap(&editbox_node, Role::App, "text", 0).unwrap();
  105. let chatview_scroll =
  106. PropertyFloat32::wrap(&chatview_node, Role::Internal, "scroll", 0).unwrap();
  107. let upgrade_popup_node = sg_root.clone().lookup_node("/window/view/upgrade_popup").unwrap();
  108. let upgrade_popup_is_visible =
  109. PropertyBool::wrap(&upgrade_popup_node, Role::App, "is_visible", 0).unwrap();
  110. info!(target: "darkirc", "Instantiating DarkIRC event DAG");
  111. let datastore = expand_path(EVGRDB_PATH)?;
  112. fs::create_dir_all(&datastore).await?;
  113. let sled_db = sled::open(datastore)?;
  114. let evgr = LocalEventGraph::new(sled_db.clone(), "darkirc_dag", 1, ex.clone()).await?;
  115. let (send_sender, send_recvr) = async_channel::unbounded();
  116. Ok(Arc::new(Self {
  117. stream: AsyncMutex::new(None),
  118. evgr,
  119. tasks: SyncMutex::new(vec![]),
  120. send_sender,
  121. send_recvr,
  122. chatview_node,
  123. sendbtn_node,
  124. editbox_node,
  125. editbox_text,
  126. chatview_scroll,
  127. upgrade_popup_is_visible,
  128. seen_msgs: SyncMutex::new(vec![]),
  129. }))
  130. }
  131. pub async fn start(self: Arc<Self>, ex: ExecutorPtr) -> Result<()> {
  132. debug!(target: "darkirc", "LocalDarkIRC::start()");
  133. let me = Arc::downgrade(&self);
  134. let mainloop_task = ex.spawn(Self::run_mainloop(me));
  135. let (slot, recvr) = Slot::new("send_button_clicked");
  136. self.sendbtn_node.register("click", slot).unwrap();
  137. let me = Arc::downgrade(&self);
  138. let send_task = ex.spawn(async move {
  139. while let Some(self_) = me.upgrade() {
  140. let Ok(_) = recvr.recv().await else {
  141. error!(target: "ui::win", "Button click recvr closed");
  142. break
  143. };
  144. self_.handle_send().await;
  145. }
  146. });
  147. let (slot, recvr) = Slot::new("enter_pressed");
  148. self.editbox_node.register("enter_pressed", slot).unwrap();
  149. let me = Arc::downgrade(&self);
  150. let enter_task = ex.spawn(async move {
  151. while let Some(self_) = me.upgrade() {
  152. let Ok(_) = recvr.recv().await else {
  153. error!(target: "ui::win", "EditBox enter_pressed recvr closed");
  154. break
  155. };
  156. self_.handle_send().await;
  157. }
  158. });
  159. let mut tasks = self.tasks.lock().unwrap();
  160. assert!(tasks.is_empty());
  161. *tasks = vec![mainloop_task, send_task, enter_task];
  162. Ok(())
  163. }
  164. async fn run_mainloop(me: Weak<Self>) {
  165. let mut send_queue: Vec<(u64, Privmsg)> = vec![];
  166. 'reconnect: loop {
  167. loop {
  168. debug!(target: "darkirc", "Connecting to evgrd...");
  169. let Some(self_) = me.upgrade() else { return };
  170. while let Err(e) = self_.connect().await {
  171. error!(target: "darkirc", "Unable to connect to evgrd backend: {e}");
  172. sleep(2).await;
  173. }
  174. debug!(target: "darkirc", "Attempting version exchange...");
  175. let Err(e) = self_.version_exchange().await else { break };
  176. error!(target: "darkirc", "Version exchange with evgrd failed: {e}");
  177. }
  178. info!(target: "darkirc", "Connected to evgrd backend");
  179. let Some(self_) = me.upgrade() else { return };
  180. if !send_queue.is_empty() {
  181. info!(target: "darkirc", "Resending {} messages", send_queue.len());
  182. }
  183. while let Some((timest, privmsg)) = send_queue.pop() {
  184. if let Err(e) = self_.send_msg(timest, privmsg.clone()).await {
  185. error!(target: "darkirc", "Send failed");
  186. send_queue.push((timest, privmsg));
  187. continue 'reconnect
  188. }
  189. }
  190. drop(self_);
  191. loop {
  192. let Some(self_) = me.upgrade() else { return };
  193. select! {
  194. res = self_.receive_msg().fuse() => {
  195. if let Err(e) = res {
  196. error!(target: "darkirc", "Receive failed: {e}");
  197. continue 'reconnect
  198. };
  199. }
  200. res = self_.send_recvr.recv().fuse() => {
  201. let (timest, privmsg) = res.unwrap();
  202. info!(target: "darkirc", "Sending msg: {timest} {privmsg:?}");
  203. if let Err(e) = self_.send_msg(timest, privmsg.clone()).await {
  204. error!(target: "darkirc", "Send failed");
  205. send_queue.push((timest, privmsg));
  206. continue 'reconnect
  207. }
  208. }
  209. }
  210. }
  211. }
  212. }
  213. async fn connect(&self) -> Result<()> {
  214. let endpoint = Url::parse(ENDPOINT)?;
  215. let dialer = Dialer::new(endpoint.clone(), None).await?;
  216. let timeout = std::time::Duration::from_secs(60);
  217. let stream = dialer.dial(Some(timeout)).await?;
  218. info!(target: "darkirc", "Connected to the backend: {endpoint}");
  219. *self.stream.lock().await = Some(stream);
  220. Ok(())
  221. }
  222. async fn version_exchange(&self) -> Result<()> {
  223. let Some(stream) = &mut *self.stream.lock().await else { return Err(Error::ConnectFailed) };
  224. let version = VersionMessage::new();
  225. debug!(target: "darkirc", "Sending version: {version:?}");
  226. version.encode_async(stream).await?;
  227. stream.flush().await?;
  228. debug!(target: "darkirc", "Receiving version...");
  229. let server_version = VersionMessage::decode_async(stream).await?;
  230. info!(target: "darkirc", "Backend server version: {}", server_version.protocol_version);
  231. if server_version.protocol_version > evgrd::PROTOCOL_VERSION {
  232. self.upgrade_popup_is_visible.set(true);
  233. }
  234. let unref_tips = self.evgr.unreferenced_tips.read().await.clone();
  235. let fetchevs = FetchEventsMessage::new(unref_tips);
  236. MSG_FETCHEVENTS.encode_async(stream).await?;
  237. stream.flush().await?;
  238. fetchevs.encode_async(stream).await?;
  239. stream.flush().await?;
  240. Ok(())
  241. }
  242. async fn send_msg(&self, timestamp: u64, msg: Privmsg) -> Result<()> {
  243. let Some(stream) = &mut *self.stream.lock().await else { return Err(Error::ConnectFailed) };
  244. MSG_SENDEVENT.encode_async(stream).await?;
  245. stream.flush().await?;
  246. timestamp.encode_async(stream).await?;
  247. stream.flush().await?;
  248. let content: Vec<u8> = serialize_async(&msg).await;
  249. content.encode_async(stream).await?;
  250. stream.flush().await?;
  251. Ok(())
  252. }
  253. async fn receive_msg(&self) -> Result<()> {
  254. debug!(target: "darkirc", "Receiving message...");
  255. let Some(stream) = &mut *self.stream.lock().await else { return Err(Error::ConnectFailed) };
  256. let msg_type = u8::decode_async(stream).await?;
  257. debug!(target: "darkirc", "Received: {msg_type:?}");
  258. if msg_type != MSG_EVENT {
  259. error!(target: "darkirc", "Received invalid msg_type: {msg_type}");
  260. //return Err(Error::MalformedPacket)
  261. return Ok(())
  262. }
  263. let ev = event_graph::Event::decode_async(stream).await?;
  264. let privmsg: Privmsg = match deserialize_async_partial(ev.content()).await {
  265. Ok((v, _)) => v,
  266. Err(e) => {
  267. error!(target: "darkirc", "Failed deserializing incoming Privmsg event: {e}");
  268. return Ok(())
  269. }
  270. };
  271. let mut timest = ev.timestamp;
  272. if timest < 6047051717 {
  273. timest *= 1000;
  274. }
  275. debug!(target: "darkirc", "Recv privmsg: <{timest}> {privmsg:?}");
  276. let genesis_timestamp = self.evgr.current_genesis.read().await.clone().timestamp;
  277. let ev_id = ev.id();
  278. if self.evgr.dag.contains_key(ev_id.as_bytes()).unwrap() ||
  279. !ev.validate(&self.evgr.dag, genesis_timestamp, self.evgr.days_rotation, None)
  280. .await?
  281. {
  282. error!(target: "darkirc", "Event is invalid! {ev:?}");
  283. return Ok(())
  284. }
  285. self.evgr.dag_insert(&[ev.clone()]).await.unwrap();
  286. if privmsg.channel != CHANNEL {
  287. //debug!(target: "darkirc", "{} != {CHANNEL}", privmsg.channel);
  288. return Ok(())
  289. }
  290. // This is a hack to make messages appear sequentially in the UI
  291. let mut adj_timest = timest;
  292. let now_timest = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  293. if timest.abs_diff(now_timest) < RECENT_TIME_DIST {
  294. debug!(target: "darkirc", "Applied timestamp correction: <{timest}> => <{now_timest}>");
  295. adj_timest = now_timest;
  296. }
  297. let msg_id = privmsg.msg_id(timest);
  298. {
  299. let mut seen = self.seen_msgs.lock().unwrap();
  300. if seen.contains(&msg_id) {
  301. warn!(target: "darkirc", "Skipping duplicate seen message: {msg_id}");
  302. return Ok(())
  303. }
  304. seen.push(msg_id.clone());
  305. }
  306. let mut arg_data = vec![];
  307. adj_timest.encode_async(&mut arg_data).await.unwrap();
  308. msg_id.encode_async(&mut arg_data).await.unwrap();
  309. privmsg.nick.encode_async(&mut arg_data).await.unwrap();
  310. privmsg.msg.encode_async(&mut arg_data).await.unwrap();
  311. self.chatview_node.call_method("insert_line", arg_data).await.unwrap();
  312. Ok(())
  313. }
  314. async fn handle_send(&self) {
  315. // Get text from editbox
  316. let text = self.editbox_text.get();
  317. if text.is_empty() {
  318. return
  319. }
  320. // Clear editbox
  321. self.editbox_text.set("");
  322. self.chatview_scroll.set(0.);
  323. // Send text to channel
  324. let timest = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  325. debug!(target: "darkirc", "Sending privmsg: <{timest}> {text}");
  326. let msg = Privmsg::new(CHANNEL.to_string(), "anon".to_string(), text);
  327. let mut arg_data = vec![];
  328. timest.encode_async(&mut arg_data).await.unwrap();
  329. msg.msg_id(timest).encode_async(&mut arg_data).await.unwrap();
  330. msg.nick.encode_async(&mut arg_data).await.unwrap();
  331. msg.msg.encode_async(&mut arg_data).await.unwrap();
  332. self.send_sender.send((timest, msg)).await.unwrap();
  333. self.chatview_node.call_method("insert_unconf_line", arg_data).await.unwrap();
  334. }
  335. }