Parcourir la source

bin/genev: use DAG eventgraph

Dastan-glitch il y a 2 ans
Parent
commit
6245422714

+ 37 - 0
Cargo.lock

@@ -3063,6 +3063,43 @@ dependencies = [
  "zeroize",
 ]
 
+[[package]]
+name = "genev"
+version = "0.4.1"
+dependencies = [
+ "clap 4.4.7",
+ "darkfi",
+ "darkfi-serial",
+ "genevd",
+ "log",
+ "simplelog",
+ "smol",
+ "tinyjson",
+ "url",
+]
+
+[[package]]
+name = "genevd"
+version = "0.4.1"
+dependencies = [
+ "async-trait",
+ "blake3",
+ "darkfi",
+ "darkfi-serial",
+ "easy-parallel",
+ "log",
+ "serde",
+ "signal-hook",
+ "signal-hook-async-std",
+ "simplelog",
+ "sled",
+ "smol",
+ "structopt",
+ "structopt-toml",
+ "tinyjson",
+ "url",
+]
+
 [[package]]
 name = "getrandom"
 version = "0.1.16"

+ 2 - 2
Cargo.toml

@@ -27,8 +27,8 @@ members = [
     "bin/faucetd",
     #"bin/fud/fu",
     "bin/fud/fud",
-    #"bin/genev/genevd",
-    #"bin/genev/genev-cli",
+    "bin/genev/genevd",
+    "bin/genev/genev-cli",
     "bin/darkirc",
     "bin/tau/taud",
     #"bin/tau/tau-cli",

+ 1 - 1
bin/genev/genev-cli/src/main.rs

@@ -86,7 +86,7 @@ fn main() -> Result<()> {
                         println!("=============================");
                         println!(
                             "- nickname: {}, title: {}, text: {}",
-                            event.action.nick, event.action.title, event.action.text
+                            event.nick, event.title, event.text
                         );
                     }
                 }

+ 4 - 5
bin/genev/genev-cli/src/rpc.rs

@@ -17,7 +17,6 @@
  */
 
 use darkfi::{
-    event_graph::model::Event,
     rpc::{client::RpcClient, jsonrpc::JsonRequest},
     util::encoding::base64,
     Result,
@@ -40,7 +39,7 @@ impl Gen {
     pub async fn add(&self, event: GenEvent) -> Result<()> {
         let event = JsonValue::String(base64::encode(&serialize(&event)));
 
-        let req = JsonRequest::new("add", vec![event]);
+        let req = JsonRequest::new("add", JsonValue::Array([event].to_vec()));
         let rep = self.rpc_client.request(req).await?;
 
         debug!("Got reply: {:?}", rep);
@@ -48,14 +47,14 @@ impl Gen {
     }
 
     /// Get current open tasks ids.
-    pub async fn list(&self) -> Result<Vec<Event<GenEvent>>> {
-        let req = JsonRequest::new("list", vec![]);
+    pub async fn list(&self) -> Result<Vec<GenEvent>> {
+        let req = JsonRequest::new("list", JsonValue::Array([].to_vec()));
         let rep = self.rpc_client.request(req).await?;
 
         debug!("reply: {:?}", rep);
 
         let bytes: Vec<u8> = base64::decode(rep.get::<String>().unwrap()).unwrap();
-        let events: Vec<Event<GenEvent>> = deserialize(&bytes)?;
+        let events: Vec<GenEvent> = deserialize(&bytes)?;
 
         Ok(events)
     }

+ 13 - 3
bin/genev/genevd/Cargo.toml

@@ -17,8 +17,18 @@ name = "genevd"
 path = "src/main.rs"
 
 [dependencies]
-darkfi = {path = "../../../", features = ["async-daemonize", "event-graph", "rpc"]}
-darkfi-serial = {path = "../../../src/serial"}
+darkfi = { path = "../../../", features = [
+    "async-daemonize",
+    "event-graph",
+    "rpc",
+] }
+darkfi-serial = { path = "../../../src/serial" }
+
+# Crypto
+blake3 = "1.5.0"
+
+# Event Graph DB
+sled = "0.34.7"
 
 # Misc
 async-trait = "0.1.74"
@@ -34,6 +44,6 @@ simplelog = "0.12.1"
 smol = "1.3.0"
 
 # Argument parsing
-serde = {version = "1.0.192", features = ["derive"]}
+serde = { version = "1.0.192", features = ["derive"] }
 structopt = "0.3.26"
 structopt-toml = "0.5.1"

+ 0 - 11
bin/genev/genevd/src/lib.rs

@@ -16,7 +16,6 @@
  * along with this program.  If not, see <https://www.gnu.org/licenses/>.
  */
 
-use darkfi::event_graph::EventMsg;
 use darkfi_serial::{async_trait, SerialDecodable, SerialEncodable};
 
 #[derive(SerialEncodable, SerialDecodable, Clone, Debug)]
@@ -25,13 +24,3 @@ pub struct GenEvent {
     pub title: String,
     pub text: String,
 }
-
-impl EventMsg for GenEvent {
-    fn new() -> Self {
-        Self {
-            nick: "groot".to_string(),
-            title: "I am groot".to_string(),
-            text: "I am groot!!".to_string(),
-        }
-    }
-}

+ 78 - 51
bin/genev/genevd/src/main.rs

@@ -16,24 +16,22 @@
  * along with this program.  If not, see <https://www.gnu.org/licenses/>.
  */
 
-use std::sync::Arc;
+use std::{
+    fs::remove_dir_all,
+    sync::{Arc, OnceLock},
+};
 
 use darkfi::{
     async_daemonize, cli_desc,
-    event_graph::{
-        events_queue::EventsQueue,
-        model::{Event, EventId, Model},
-        protocol_event::{ProtocolEvent, Seen, SeenPtr},
-        view::{View, ViewPtr},
-    },
-    net::{self, settings::SettingsOpt},
+    event_graph::{proto::ProtocolEventGraph, EventGraph, EventGraphPtr, NULL_ID},
+    net::{settings::SettingsOpt, P2p, SESSION_ALL},
     rpc::server::{listen_and_serve, RequestHandler},
-    system::StoppableTask,
+    system::{sleep, StoppableTask},
+    util::path::expand_path,
     Error, Result,
 };
-use genevd::GenEvent;
-use log::{error, info};
-use smol::{lock::Mutex, stream::StreamExt};
+use log::{debug, error, info};
+use smol::{lock::RwLock, stream::StreamExt};
 use structopt_toml::{serde::Deserialize, structopt::StructOpt, StructOptToml};
 use url::Url;
 
@@ -58,63 +56,65 @@ struct Args {
     #[structopt(flatten)]
     pub net: SettingsOpt,
 
+    /// Sets Datastore Path
+    #[structopt(long, default_value = "~/.local/darkfi/genev")]
+    pub datastore: String,
+
     #[structopt(short, long)]
     /// Set log file to ouput into
     log: Option<String>,
 
+    #[structopt(long)]
+    pub skip_dag_sync: bool,
+
     #[structopt(short, parse(from_occurrences))]
     /// Increase verbosity (-vvv supported)
     verbose: u8,
 }
 
 async fn start_sync_loop(
-    view: ViewPtr<GenEvent>,
-    seen: SeenPtr<EventId>,
-    missed_events: Arc<Mutex<Vec<Event<GenEvent>>>>,
+    event_graph: EventGraphPtr,
+    last_sent: RwLock<blake3::Hash>,
+    seen: OnceLock<sled::Tree>,
 ) -> Result<()> {
+    let incoming = event_graph.event_sub.clone().subscribe().await;
+    let seen_events = seen.get().unwrap();
     loop {
-        let event = view.lock().await.process().await?;
-        if !seen.push(&event.hash()).await {
+        let event = incoming.receive().await;
+        let event_id = event.id();
+        if *last_sent.read().await == event_id {
+            continue
+        }
+
+        if seen_events.contains_key(event_id.as_bytes()).unwrap() {
             continue
         }
 
-        info!("new event: {:?}", event);
-        missed_events.lock().await.push(event.clone());
+        debug!("new event: {:?}", event);
     }
 }
 
 async_daemonize!(realmain);
-async fn realmain(args: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
+async fn realmain(settings: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
     ////////////////////
     // Initialize the base structures
     ////////////////////
-    let events_queue = EventsQueue::<GenEvent>::new();
-    let model = Arc::new(Mutex::new(Model::new(events_queue.clone())));
-    let view = Arc::new(Mutex::new(View::new(events_queue)));
-    let model_clone = model.clone();
+    info!("Instantiating event DAG");
+    // Create datastore path if not there already.
+    let datastore_path = expand_path(&settings.datastore)?;
 
-    ////////////////////
-    // P2p setup
-    ////////////////////
-    // Buffers
-    let seen_event = Seen::new();
-    let seen_inv = Seen::new();
+    let sled_db = sled::open(datastore_path.clone())?;
+    let p2p = P2p::new(settings.net.into(), executor.clone()).await;
+    let event_graph =
+        EventGraph::new(p2p.clone(), sled_db.clone(), "genevd_dag", 1, executor.clone()).await?;
 
-    // Check the version
-    let net_settings = args.net.clone();
-
-    // New p2p
-    let p2p = net::P2p::new(net_settings.into(), executor.clone()).await;
-    let p2p2 = p2p.clone();
-
-    // Register the protocol_event
+    info!("Registering EventGraph P2P protocol");
+    let event_graph_ = Arc::clone(&event_graph);
     let registry = p2p.protocol_registry();
     registry
-        .register(net::SESSION_ALL, move |channel, p2p| {
-            let seen_event = seen_event.clone();
-            let seen_inv = seen_inv.clone();
-            let model = model.clone();
-            async move { ProtocolEvent::init(channel, p2p, model, seen_event, seen_inv).await }
+        .register(SESSION_ALL, move |channel, _| {
+            let event_graph_ = event_graph_.clone();
+            async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
         })
         .await;
 
@@ -122,16 +122,42 @@ async fn realmain(args: Args, executor: Arc<smol::Executor<'static>>) -> Result<
     info!(target: "genevd", "Starting P2P network");
     p2p.clone().start().await?;
 
+    info!(target: "genevd", "Waiting for some P2P connections...");
+    sleep(5).await;
+
+    // We'll attempt to sync 5 times
+    if !settings.skip_dag_sync {
+        for i in 1..=6 {
+            info!("Syncing event DAG (attempt #{})", i);
+            match event_graph.dag_sync().await {
+                Ok(()) => break,
+                Err(e) => {
+                    if i == 6 {
+                        error!("Failed syncing DAG. Exiting.");
+                        p2p.stop().await;
+                        return Err(Error::DagSyncFailed)
+                    } else {
+                        // TODO: Maybe at this point we should prune or something?
+                        // TODO: Or maybe just tell the user to delete the DAG from FS.
+                        error!("Failed syncing DAG ({}), retrying in 10s...", e);
+                        sleep(10).await;
+                    }
+                }
+            }
+        }
+    }
+
     ////////////////////
     // Listner
     ////////////////////
-    let seen_ids = Seen::new();
-    let missed_events = Arc::new(Mutex::new(vec![]));
+    let last_sent = RwLock::new(NULL_ID);
+    let seen = OnceLock::new();
+    seen.set(sled_db.open_tree("genevdb").unwrap()).unwrap();
 
     info!(target: "genevd", "Starting sync loop task");
     let sync_loop_task = StoppableTask::new();
     sync_loop_task.clone().start(
-        start_sync_loop(view, seen_ids.clone(), missed_events.clone()),
+        start_sync_loop(event_graph.clone(), last_sent, seen.clone()),
         |res| async {
             match res {
                 Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
@@ -147,15 +173,14 @@ async fn realmain(args: Args, executor: Arc<smol::Executor<'static>>) -> Result<
     //
     let rpc_interface = Arc::new(JsonRpcInterface::new(
         "Alolymous".to_string(),
-        missed_events.clone(),
-        model_clone,
-        seen_ids.clone(),
+        event_graph.clone(),
+        seen.clone(),
         p2p.clone(),
     ));
     let rpc_task = StoppableTask::new();
     let rpc_interface_ = rpc_interface.clone();
     rpc_task.clone().start(
-        listen_and_serve(args.rpc_listen, rpc_interface, None, executor.clone()),
+        listen_and_serve(settings.rpc_listen, rpc_interface, None, executor.clone()),
         |res| async move {
             match res {
                 Ok(()) | Err(Error::RpcServerStopped) => rpc_interface_.stop_connections().await,
@@ -178,7 +203,9 @@ async fn realmain(args: Args, executor: Arc<smol::Executor<'static>>) -> Result<
     sync_loop_task.stop().await;
 
     // stop p2p
-    p2p2.stop().await;
+    p2p.stop().await;
+
+    remove_dir_all(datastore_path).unwrap_or(());
 
     Ok(())
 }

+ 49 - 35
bin/genev/genevd/src/rpc.rs

@@ -16,34 +16,30 @@
  * along with this program.  If not, see <https://www.gnu.org/licenses/>.
  */
 
-use std::{collections::HashSet, sync::Arc};
+use std::{collections::HashSet, sync::OnceLock};
 
 use async_trait::async_trait;
-use log::debug;
+use log::{debug, error, info};
 use smol::lock::{Mutex, MutexGuard};
 use tinyjson::JsonValue;
 
 use darkfi::{
-    event_graph::{
-        model::{Event, EventId, ModelPtr},
-        protocol_event::SeenPtr,
-    },
+    event_graph::{proto::EventPut, Event, EventGraphPtr},
     net,
     rpc::{
         jsonrpc::{ErrorCode, JsonError, JsonRequest, JsonResponse, JsonResult},
         server::RequestHandler,
     },
     system::StoppableTaskPtr,
-    util::{encoding::base64, time::Timestamp},
+    util::encoding::base64,
 };
-use darkfi_serial::deserialize;
+use darkfi_serial::{deserialize, deserialize_async_partial, serialize_async};
 use genevd::GenEvent;
 
 pub struct JsonRpcInterface {
     _nickname: String,
-    missed_events: Arc<Mutex<Vec<Event<GenEvent>>>>,
-    model: ModelPtr<GenEvent>,
-    seen: SeenPtr<EventId>,
+    event_graph: EventGraphPtr,
+    seen: OnceLock<sled::Tree>,
     p2p: net::P2pPtr,
     rpc_connections: Mutex<HashSet<StoppableTaskPtr>>,
 }
@@ -69,19 +65,11 @@ impl RequestHandler for JsonRpcInterface {
 impl JsonRpcInterface {
     pub fn new(
         _nickname: String,
-        missed_events: Arc<Mutex<Vec<Event<GenEvent>>>>,
-        model: ModelPtr<GenEvent>,
-        seen: SeenPtr<EventId>,
+        event_graph: EventGraphPtr,
+        seen: OnceLock<sled::Tree>,
         p2p: net::P2pPtr,
     ) -> Self {
-        Self {
-            _nickname,
-            missed_events,
-            model,
-            seen,
-            p2p,
-            rpc_connections: Mutex::new(HashSet::new()),
-        }
+        Self { _nickname, event_graph, seen, p2p, rpc_connections: Mutex::new(HashSet::new()) }
     }
 
     // RPCAPI:
@@ -122,19 +110,16 @@ impl JsonRpcInterface {
         let dec = base64::decode(b64).unwrap();
         let genevent: GenEvent = deserialize(&dec).unwrap();
 
-        let event = Event {
-            previous_event_hash: self.model.lock().await.get_head_hash().unwrap(),
-            action: genevent,
-            timestamp: Timestamp::current_time(),
-        };
+        // Build a DAG event and return it.
+        let event = Event::new(serialize_async(&genevent).await, self.event_graph.clone()).await;
 
-        if !self.seen.push(&event.hash()).await {
-            let json = JsonValue::Boolean(false);
-            return JsonResponse::new(json, id).into()
+        if let Err(e) = self.event_graph.dag_insert(event.clone()).await {
+            error!("Failed inserting new event to DAG: {}", e);
+        } else {
+            // Otherwise, broadcast it
+            self.p2p.broadcast(&EventPut(event)).await;
         }
 
-        self.p2p.broadcast(&event).await;
-
         let json = JsonValue::Boolean(true);
         JsonResponse::new(json, id).into()
     }
@@ -144,10 +129,39 @@ impl JsonRpcInterface {
     // --> {"jsonrpc": "2.0", "method": "list", "params": [], "id": 1}
     // <-- {"jsonrpc": "2.0", "result": [task_id, ...], "id": 1}
     async fn list(&self, id: u16, _params: JsonValue) -> JsonResult {
-        debug!("fetching all events");
-        let msd = self.missed_events.lock().await.clone();
+        debug!("Fetching all events");
+        let mut sevents = vec![];
+
+        ////////////////////
+        // get history
+        ////////////////////
+        let dag_events = self.event_graph.order_events().await;
+        let seen_events = self.seen.get().unwrap();
+
+        for event_id in dag_events.iter() {
+            // If it was seen, skip
+            if seen_events.contains_key(event_id.as_bytes()).unwrap() {
+                continue
+            }
+
+            // Get the event from the DAG
+            let event = self.event_graph.dag_get(event_id).await.unwrap().unwrap();
+
+            // Try to deserialize it. (Here we skip errors)
+            let genevent: GenEvent = match deserialize_async_partial(event.content()).await {
+                Ok((v, _)) => v,
+                Err(e) => {
+                    error!("Failed deserializing incoming event: {}", e);
+                    continue
+                }
+            };
+
+            info!("Marking event {} as seen", event_id);
+            seen_events.insert(event_id.as_bytes(), &[]).unwrap();
+            sevents.push(genevent);
+        }
 
-        let ser = darkfi_serial::serialize(&msd);
+        let ser = darkfi_serial::serialize(&sevents);
         let enc = JsonValue::String(base64::encode(&ser));
 
         JsonResponse::new(enc, id).into()