/* This file is part of DarkFi (https://dark.fi) * * Copyright (C) 2020-2024 Dyne.org foundation * * This program is free software: you can redistribute it and/or modify * it under the terms of the GNU Affero General Public License as * published by the Free Software Foundation, either version 3 of the * License, or (at your option) any later version. * * This program is distributed in the hope that it will be useful, * but WITHOUT ANY WARRANTY; without even the implied warranty of * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU Affero General Public License for more details. * * You should have received a copy of the GNU Affero General Public License * along with this program. If not, see . */ use async_channel::{Receiver, Sender}; use async_lock::Mutex as AsyncMutex; use darkfi::{ event_graph::{self}, net::transport::{Dialer, PtStream}, system::{sleep, ExecutorPtr}, util::path::expand_path, Error, Result, }; use darkfi_serial::{ async_trait, deserialize_async_partial, serialize_async, AsyncDecodable, AsyncEncodable, Encodable, SerialDecodable, SerialEncodable, }; use evgrd::{ FetchEventsMessage, LocalEventGraph, LocalEventGraphPtr, VersionMessage, MSG_EVENT, MSG_FETCHEVENTS, MSG_SENDEVENT, }; use futures::{select, FutureExt}; use log::{error, info}; use sled_overlay::sled; use smol::{ fs, io::{ReadHalf, WriteHalf}, }; use std::{ sync::{ atomic::{AtomicBool, Ordering}, Arc, Mutex as SyncMutex, Weak, }, time::UNIX_EPOCH, }; use url::Url; use crate::{ prop::{PropertyBool, PropertyStr, Role}, scene::{SceneNodePtr, Slot}, }; #[cfg(target_os = "android")] const EVGRDB_PATH: &str = "/data/data/darkfi.darkwallet/evgr/"; #[cfg(target_os = "linux")] const EVGRDB_PATH: &str = "~/.local/darkfi/darkwallet/evgr/"; const ENDPOINT: &str = "tcp://agorism.dev:5588"; const CHANNEL: &str = "#random"; #[derive(Clone, Debug, SerialEncodable, SerialDecodable)] pub struct Privmsg { pub channel: String, pub nick: String, pub msg: String, } impl Privmsg { pub fn new(channel: String, nick: String, msg: String) -> Self { Self { channel, nick, msg } } } pub type LocalDarkIRCPtr = Arc; pub struct LocalDarkIRC { stream: AsyncMutex>>, evgr: LocalEventGraphPtr, tasks: SyncMutex>>, send_sender: Sender<(u64, Privmsg)>, send_recvr: Receiver<(u64, Privmsg)>, chatview_node: SceneNodePtr, sendbtn_node: SceneNodePtr, editbox_node: SceneNodePtr, editbox_text: PropertyStr, upgrade_popup_is_visible: PropertyBool, } impl LocalDarkIRC { pub async fn new(sg_root: SceneNodePtr, ex: ExecutorPtr) -> Result> { let chatview_node = sg_root.clone().lookup_node("/window/view/chatty").unwrap(); let sendbtn_node = sg_root.clone().lookup_node("/window/view/send_btn").unwrap(); let editbox_node = sg_root.clone().lookup_node("/window/view/editz").unwrap(); let editbox_text = PropertyStr::wrap(&editbox_node, Role::App, "text", 0).unwrap(); let upgrade_popup_node = sg_root.clone().lookup_node("/window/view/upgrade_popup").unwrap(); let upgrade_popup_is_visible = PropertyBool::wrap(&upgrade_popup_node, Role::App, "is_visible", 0).unwrap(); info!(target: "darkirc", "Instantiating DarkIRC event DAG"); let datastore = expand_path(EVGRDB_PATH)?; fs::create_dir_all(&datastore).await?; let sled_db = sled::open(datastore)?; let evgr = LocalEventGraph::new(sled_db.clone(), "darkirc_dag", 1, ex.clone()).await?; let (send_sender, send_recvr) = async_channel::unbounded(); Ok(Arc::new(Self { stream: AsyncMutex::new(None), evgr, tasks: SyncMutex::new(vec![]), send_sender, send_recvr, chatview_node, sendbtn_node, editbox_node, editbox_text, upgrade_popup_is_visible, })) } pub async fn start(self: Arc, ex: ExecutorPtr) -> Result<()> { debug!(target: "darkirc", "LocalDarkIRC::start()"); let me = Arc::downgrade(&self); let mainloop_task = ex.spawn(Self::run_mainloop(me)); let (slot, recvr) = Slot::new("send_button_clicked"); self.sendbtn_node.register("click", slot).unwrap(); let me = Arc::downgrade(&self); let send_task = ex.spawn(async move { while let Some(self_) = me.upgrade() { let Ok(_) = recvr.recv().await else { error!(target: "ui::win", "Button click recvr closed"); break }; self_.handle_send().await; } }); let (slot, recvr) = Slot::new("enter_pressed"); self.editbox_node.register("enter_pressed", slot).unwrap(); let me = Arc::downgrade(&self); let enter_task = ex.spawn(async move { while let Some(self_) = me.upgrade() { let Ok(_) = recvr.recv().await else { error!(target: "ui::win", "EditBox enter_pressed recvr closed"); break }; self_.handle_send().await; } }); let mut tasks = self.tasks.lock().unwrap(); assert!(tasks.is_empty()); *tasks = vec![mainloop_task, send_task, enter_task]; Ok(()) } async fn run_mainloop(me: Weak) { let mut send_queue = vec![]; 'reconnect: loop { loop { debug!(target: "darkirc", "Connecting to evgrd..."); let Some(self_) = me.upgrade() else { return }; while let Err(e) = self_.connect().await { error!(target: "darkirc", "Unable to connect to evgrd backend: {e}"); sleep(2).await; } let Err(e) = self_.version_exchange().await else { break }; error!(target: "darkirc", "Version exchange with evgrd failed: {e}"); } info!(target: "darkirc", "Connected to evgrd backend"); let Some(self_) = me.upgrade() else { return }; if !send_queue.is_empty() { info!(target: "darkirc", "Resending {} messages", send_queue.len()); } for (timest, privmsg) in std::mem::take(&mut send_queue) { if let Err(e) = self_.send_msg(timest, privmsg).await { error!(target: "darkirc", "Send failed"); continue 'reconnect } } drop(self_); loop { let Some(self_) = me.upgrade() else { return }; select! { res = self_.receive_msg().fuse() => { let Ok(msg) = res else { continue 'reconnect }; } res = self_.send_recvr.recv().fuse() => { let (timest, privmsg) = res.unwrap(); info!(target: "darkirc", "Sending msg: {timest} {privmsg:?}"); if let Err(e) = self_.send_msg(timest, privmsg.clone()).await { error!(target: "darkirc", "Send failed"); send_queue.push((timest, privmsg)); continue 'reconnect } } } } } } async fn connect(&self) -> Result<()> { let endpoint = Url::parse(ENDPOINT)?; let dialer = Dialer::new(endpoint.clone(), None).await?; let timeout = std::time::Duration::from_secs(60); let stream = dialer.dial(Some(timeout)).await?; info!(target: "darkirc", "Connected to the backend: {endpoint}"); *self.stream.lock().await = Some(stream); Ok(()) } async fn version_exchange(&self) -> Result<()> { let Some(stream) = &mut *self.stream.lock().await else { return Err(Error::ConnectFailed) }; let version = VersionMessage::new(); version.encode_async(stream).await?; let server_version = VersionMessage::decode_async(stream).await?; info!(target: "darkirc", "Backend server version: {}", server_version.protocol_version); if server_version.protocol_version > evgrd::PROTOCOL_VERSION { self.upgrade_popup_is_visible.set(true); } let unref_tips = self.evgr.unreferenced_tips.read().await.clone(); let fetchevs = FetchEventsMessage::new(unref_tips); MSG_FETCHEVENTS.encode_async(stream).await?; fetchevs.encode_async(stream).await?; Ok(()) } async fn send_msg(&self, timestamp: u64, msg: Privmsg) -> Result<()> { let Some(stream) = &mut *self.stream.lock().await else { return Err(Error::ConnectFailed) }; MSG_SENDEVENT.encode_async(stream).await?; timestamp.encode_async(stream).await?; let content: Vec = serialize_async(&msg).await; content.encode_async(stream).await?; Ok(()) } async fn receive_msg(&self) -> Result<()> { debug!(target: "darkirc", "Receiving message..."); let Some(stream) = &mut *self.stream.lock().await else { return Err(Error::ConnectFailed) }; let msg_type = u8::decode_async(stream).await?; debug!(target: "darkirc", "Received: {msg_type:?}"); if msg_type != MSG_EVENT { error!(target: "darkirc", "Received invalid msg_type: {msg_type}"); //return Err(Error::MalformedPacket) return Ok(()) } let ev = event_graph::Event::decode_async(stream).await?; let genesis_timestamp = self.evgr.current_genesis.read().await.clone().timestamp; let ev_id = ev.id(); if self.evgr.dag.contains_key(ev_id.as_bytes()).unwrap() || !ev.validate(&self.evgr.dag, genesis_timestamp, self.evgr.days_rotation, None) .await? { error!(target: "darkirc", "Event is invalid! {ev:?}"); return Ok(()) } debug!(target: "darkirc", "got {ev:?}"); self.evgr.dag_insert(&[ev.clone()]).await.unwrap(); let privmsg: Privmsg = match deserialize_async_partial(ev.content()).await { Ok((v, _)) => v, Err(e) => { error!(target: "darkirc", "Failed deserializing incoming Privmsg event: {e}"); return Ok(()) } }; debug!(target: "darkirc", "privmsg: {privmsg:?}"); let mut timest = ev.timestamp; if timest < 6047051717 { timest *= 1000; } if privmsg.channel != CHANNEL { return Ok(()) } let mut arg_data = vec![]; timest.encode_async(&mut arg_data).await.unwrap(); ev.id().as_bytes().encode_async(&mut arg_data).await.unwrap(); privmsg.nick.encode_async(&mut arg_data).await.unwrap(); privmsg.msg.encode_async(&mut arg_data).await.unwrap(); self.chatview_node.call_method("insert_line", arg_data).await.unwrap(); Ok(()) } async fn handle_send(&self) { // Get text from editbox let text = self.editbox_text.get(); // Clear editbox self.editbox_text.set(""); // Send text to channel debug!(target: "darkirc", "Sending privmsg: {text}"); let msg = Privmsg::new(CHANNEL.to_string(), "anon".to_string(), text); let timestamp = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64; self.send_sender.send((timestamp, msg)).await.unwrap(); } }