Просмотр исходного кода

wallet: add LocalEventGraph backend

darkfi 1 год назад
Родитель
Сommit
b852e893c4

+ 1 - 0
bin/darkwallet/Cargo.toml

@@ -24,6 +24,7 @@ glam = "0.29.0"
 #async_zmq = "0.4.0"
 zeromq = { version = "0.4.0", default-features = false, features = ["async-std-runtime", "all-transport"] }
 darkfi = {path = "../../", features = ["async-daemonize", "event-graph", "net", "util", "system", "zk"]}
+evgrd = {path = "../../script/evgrd/"}
 #darkfi-sdk = {path = "../../src/sdk", features = ["async"]}
 darkfi-serial = {path = "../../src/serial", features = ["async"]}
 thiserror = "1.0.63"

+ 10 - 66
bin/darkwallet/src/app/mod.rs

@@ -18,6 +18,7 @@
 
 use async_recursion::async_recursion;
 use chrono::{Local, NaiveDate, NaiveDateTime, TimeZone};
+use darkfi::system::CondVar;
 use darkfi_serial::Encodable;
 use futures::{stream::FuturesUnordered, StreamExt};
 use sled_overlay::sled;
@@ -115,13 +116,14 @@ impl AsyncRuntime {
 pub type AppPtr = Arc<App>;
 
 pub struct App {
-    pub(self) sg_root: SceneNodePtr,
-    pub(self) render_api: RenderApiPtr,
-    pub(self) event_pub: GraphicsEventPublisherPtr,
-    pub(self) text_shaper: TextShaperPtr,
+    pub sg_root: SceneNodePtr,
+    pub is_started: Arc<CondVar>,
+    pub render_api: RenderApiPtr,
+    pub event_pub: GraphicsEventPublisherPtr,
+    pub text_shaper: TextShaperPtr,
     //pub(self) darkirc_backend: DarkIrcBackendPtr,
-    pub(self) tasks: SyncMutex<Vec<Task<()>>>,
-    pub(self) ex: ExecutorPtr,
+    pub tasks: SyncMutex<Vec<Task<()>>>,
+    pub ex: ExecutorPtr,
 }
 
 impl App {
@@ -135,6 +137,7 @@ impl App {
     ) -> Arc<Self> {
         Arc::new(Self {
             sg_root,
+            is_started: Arc::new(CondVar::new()),
             ex,
             render_api,
             event_pub,
@@ -147,61 +150,6 @@ impl App {
     pub async fn start(self: Arc<Self>) {
         debug!(target: "app", "App::start()");
 
-        ////////////////////////////////////////////////////////////////////////////////////
-        // OLD
-        ////////////////////////////////////////////////////////////////////////////////////
-
-        /*
-        // Setup UI
-        let mut sg = self.sg.lock().await;
-
-        let window = sg.add_node("window", SceneNodeType::Window);
-
-        let mut prop = Property::new("screen_size", PropertyType::Float32, PropertySubType::Pixel);
-        prop.set_array_len(2);
-        // Window not yet initialized so we can't set these.
-        //prop.set_f32(Role::App, 0, screen_width);
-        //prop.set_f32(Role::App, 1, screen_height);
-        window.add_property(prop).unwrap();
-
-        let mut prop = Property::new("scale", PropertyType::Float32, PropertySubType::Pixel);
-        prop.set_defaults_f32(vec![1.]).unwrap();
-        window.add_property(prop).unwrap();
-
-        let window_id = window.id;
-
-        // Create Window
-        // Window::new(window, weak sg)
-        drop(sg);
-        let pimpl = Window::new(
-            self.ex.clone(),
-            self.sg.clone(),
-            window_id,
-            self.render_api.clone(),
-            self.event_pub.clone(),
-        )
-        .await;
-        // -> reads any props it needs
-        // -> starts procs
-        let mut sg = self.sg.lock().await;
-        let node = sg.get_node_mut(window_id).unwrap();
-        node.pimpl = pimpl;
-
-        sg.link(window_id, SceneGraph::ROOT_ID).unwrap();
-
-        // Testing
-        let node = sg.get_node(window_id).unwrap();
-        node.set_property_f32(Role::App, "scale", 2.).unwrap();
-
-        drop(sg);
-
-        //schema::make(&self).await;
-        debug!(target: "app", "Schema loaded");
-        */
-
-        ////////////////////////////////////////////////////////////////////////////////////
-        // NEW
-        ////////////////////////////////////////////////////////////////////////////////////
         let mut window = SceneNode3::new("window", SceneNodeType3::Window);
 
         let mut prop = Property::new("screen_size", PropertyType::Float32, PropertySubType::Pixel);
@@ -228,12 +176,8 @@ impl App {
         // Access drawable in window node and call draw()
         self.trigger_draw().await;
 
-        // Start the backend
-        //if let Err(err) = self.darkirc_backend.start(self.sg.clone(), self.ex.clone()).await {
-        //    error!(target: "app", "backend error: {err}");
-        //}
-
         debug!(target: "app", "App started");
+        self.is_started.notify();
     }
 
     pub fn stop(&self) {

+ 3 - 3
bin/darkwallet/src/app/schema.rs

@@ -618,9 +618,9 @@ pub(super) async fn make(app: &App, window: SceneNodePtr) {
 
     let db = sled::open(CHATDB_PATH).expect("cannot open sleddb");
     let chat_tree = db.open_tree(b"chat").unwrap();
-    if chat_tree.is_empty() {
-        populate_tree(&chat_tree);
-    }
+    //if chat_tree.is_empty() {
+    //    populate_tree(&chat_tree);
+    //}
     debug!(target: "app", "db has {} lines", chat_tree.len());
     let node = node
         .setup(|me| {

+ 126 - 0
bin/darkwallet/src/darkirc2.rs

@@ -0,0 +1,126 @@
+/* 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 <https://www.gnu.org/licenses/>.
+ */
+
+use darkfi::{
+    event_graph::{self},
+    net::transport::Dialer,
+    system::ExecutorPtr,
+    util::path::expand_path,
+    Error, Result,
+};
+use darkfi_serial::{
+    async_trait, deserialize_async_partial, AsyncDecodable, AsyncEncodable, Encodable,
+    SerialDecodable, SerialEncodable,
+};
+use evgrd::{FetchEventsMessage, LocalEventGraph, VersionMessage, MSG_EVENT, MSG_FETCHEVENTS};
+use log::{error, info};
+use sled_overlay::sled;
+use smol::fs;
+use url::Url;
+
+use crate::scene::SceneNodePtr;
+
+#[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/";
+
+#[derive(Clone, Debug, SerialEncodable, SerialDecodable)]
+pub struct Privmsg {
+    pub channel: String,
+    pub nick: String,
+    pub msg: String,
+}
+
+pub async fn receive_msgs(sg_root: SceneNodePtr, ex: ExecutorPtr) -> Result<()> {
+    let chatview_node = sg_root.lookup_node("/window/view/chatty").ok_or(Error::ConnectFailed)?;
+
+    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 endpoint = "tcp://127.0.0.1:5588";
+    let endpoint = Url::parse(endpoint)?;
+
+    let dialer = Dialer::new(endpoint.clone(), None).await?;
+    let timeout = std::time::Duration::from_secs(60);
+
+    let mut stream = dialer.dial(Some(timeout)).await?;
+    info!(target: "darkirc", "Connected to the backend: {endpoint}");
+
+    let version = VersionMessage::new();
+    version.encode_async(&mut stream).await?;
+
+    let server_version = VersionMessage::decode_async(&mut stream).await?;
+    info!(target: "darkirc", "Backend server version: {}", server_version.protocol_version);
+
+    let unref_tips = evgr.unreferenced_tips.read().await.clone();
+    let fetchevs = FetchEventsMessage::new(unref_tips);
+    MSG_FETCHEVENTS.encode_async(&mut stream).await?;
+    fetchevs.encode_async(&mut stream).await?;
+
+    loop {
+        let msg_type = u8::decode_async(&mut 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)
+        }
+
+        let ev = event_graph::Event::decode_async(&mut stream).await?;
+
+        let genesis_timestamp = evgr.current_genesis.read().await.clone().timestamp;
+        let ev_id = ev.id();
+        if evgr.dag.contains_key(ev_id.as_bytes()).unwrap() ||
+            !ev.validate(&evgr.dag, genesis_timestamp, evgr.days_rotation, None).await?
+        {
+            error!(target: "darkirc", "Event is invalid! {ev:?}");
+            continue
+        }
+
+        debug!(target: "darkirc", "got {ev:?}");
+        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}");
+                continue
+            }
+        };
+
+        debug!(target: "darkirc", "privmsg: {privmsg:?}");
+
+        if privmsg.channel != "random" {
+            continue
+        }
+
+        let response_fn = Box::new(|_| {});
+
+        let mut arg_data = vec![];
+        ev.timestamp.encode(&mut arg_data).unwrap();
+        ev.id().as_bytes().encode(&mut arg_data).unwrap();
+        privmsg.nick.encode(&mut arg_data).unwrap();
+        privmsg.msg.encode(&mut arg_data).unwrap();
+
+        chatview_node.call_method("insert_line", arg_data, response_fn).unwrap();
+    }
+}

+ 16 - 2
bin/darkwallet/src/main.rs

@@ -47,6 +47,7 @@ use log::LevelFilter;
 
 mod app;
 //mod darkirc;
+mod darkirc2;
 mod error;
 mod expr;
 mod gfx;
@@ -127,8 +128,8 @@ fn main() {
 
     //let darkirc_backend = DarkIrcBackend::new();
     let app = app::App::new(
-        sg_root.clone(),
-        render_api.clone(),
+        sg_root,
+        render_api,
         event_pub.clone(),
         text_shaper,
         //darkirc_backend,
@@ -137,6 +138,18 @@ fn main() {
     let app_task = ex.spawn(app.clone().start());
     async_runtime.push_task(app_task);
 
+    let cv_started = app.is_started.clone();
+    let sg_root = app.sg_root.clone();
+    let ex2 = ex.clone();
+    let darkirc_task = ex.spawn(async move {
+        cv_started.wait().await;
+        if let Err(e) = darkirc2::receive_msgs(sg_root, ex2).await {
+            error!("DarkIRC error: {e}")
+        }
+    });
+    async_runtime.push_task(darkirc_task);
+
+    /*
     // Nice to see which events exist
     let ev_sub = event_pub.subscribe_key_down();
     let ev_relay_task = ex.spawn(async move {
@@ -186,6 +199,7 @@ fn main() {
         }
     });
     async_runtime.push_task(ev_relay_task);
+    */
 
     //let stage = gfx::Stage::new(method_rep, event_pub);
     gfx::run_gui(app, async_runtime, method_rep, event_pub);

+ 15 - 0
bin/darkwallet/src/scene.rs

@@ -293,6 +293,21 @@ impl SceneNode {
     fn has_method(&self, name: &str) -> bool {
         self.methods.iter().any(|sig| sig.name == name)
     }
+
+    pub fn get_method(&self, name: &str) -> Option<&Method> {
+        self.methods.iter().find(|method| method.name == name)
+    }
+
+    pub fn call_method(
+        &self,
+        name: &str,
+        arg_data: Vec<u8>,
+        response_fn: MethodResponseFn,
+    ) -> Result<()> {
+        let method = self.get_method(name).ok_or(Error::MethodNotFound)?;
+        (method.method_fn)(arg_data, response_fn);
+        Ok(())
+    }
 }
 
 impl std::fmt::Debug for SceneNode {