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

bin/taud: run tau on raft network

ghassmo 4 лет назад
Родитель
Сommit
5a7258d7f1
5 измененных файлов с 76 добавлено и 18 удалено
  1. 3 5
      bin/taud/src/jsonrpc.rs
  2. 52 12
      bin/taud/src/main.rs
  3. 4 0
      bin/taud/src/task_info.rs
  4. 16 1
      bin/taud/src/util.rs
  5. 1 0
      bin/taud/taud_config.toml

+ 3 - 5
bin/taud/src/jsonrpc.rs

@@ -91,8 +91,6 @@ impl JsonRpcInterface {
         new_task.set_project(&task.project);
         new_task.set_assign(&task.assign);
 
-        new_task.save()?;
-        new_task.activate()?;
         self.notify_queue_sender.send(new_task).await.map_err(Error::from)?;
 
         Ok(json!(true))
@@ -119,7 +117,7 @@ impl JsonRpcInterface {
         }
 
         let task = self.check_data_for_update(&args[0], &args[1])?;
-        task.save()?;
+
         self.notify_queue_sender.send(task).await.map_err(Error::from)?;
 
         Ok(json!(true))
@@ -156,7 +154,7 @@ impl JsonRpcInterface {
 
         let mut task: TaskInfo = self.load_task_by_id(&args[0])?;
         task.set_state(&state);
-        task.save()?;
+
         self.notify_queue_sender.send(task).await.map_err(Error::from)?;
 
         Ok(json!(true))
@@ -178,7 +176,7 @@ impl JsonRpcInterface {
 
         let mut task: TaskInfo = self.load_task_by_id(&args[0])?;
         task.set_comment(Comment::new(&comment_content, &comment_author));
-        task.save()?;
+
         self.notify_queue_sender.send(task).await.map_err(Error::from)?;
         Ok(json!(true))
     }

+ 52 - 12
bin/taud/src/main.rs

@@ -1,16 +1,18 @@
-use std::{fs::create_dir_all, sync::Arc};
+use std::{fs::create_dir_all, path::PathBuf, sync::Arc};
 
 use async_executor::Executor;
-use clap::{IntoApp, Parser};
+use clap::Parser;
 use simplelog::{ColorChoice, TermLogger, TerminalMode};
 
 use darkfi::{
     net::Settings as P2pSettings,
+    raft::Raft,
     rpc::rpcserver::{listen_and_serve, RpcServerConfig},
     util::{
         cli::{log_config, spawn_config, Config},
         expand_path,
         path::get_config_path,
+        sleep,
     },
     Error, Result,
 };
@@ -22,12 +24,13 @@ mod task_info;
 mod util;
 
 use crate::{
+    error::TaudResult,
     jsonrpc::JsonRpcInterface,
     task_info::TaskInfo,
     util::{CliTaud, Settings, TauConfig, CONFIG_FILE_CONTENTS},
 };
 
-async fn start(config: TauConfig, executor: Arc<Executor<'_>>) -> Result<()> {
+async fn start(config: TauConfig, args: CliTaud, executor: Arc<Executor<'_>>) -> Result<()> {
     if config.dataset_path.is_empty() {
         return Err(Error::ParseFailed("Failed to parse dataset_path"))
     }
@@ -40,7 +43,23 @@ async fn start(config: TauConfig, executor: Arc<Executor<'_>>) -> Result<()> {
 
     let settings = Settings { dataset_path };
 
-    let _p2p_settings = P2pSettings::default();
+    let p2p_settings = P2pSettings {
+        inbound: args.accept,
+        outbound_connections: args.slots,
+        external_addr: args.accept,
+        peers: args.connect.clone(),
+        seeds: args.seed.clone(),
+        ..Default::default()
+    };
+
+    //
+    //Raft
+    //
+    let mut raft =
+        Raft::<TaskInfo>::new(p2p_settings.inbound, PathBuf::from(config.datastore_raft))?;
+
+    let raft_sender = raft.get_broadcast().clone();
+    let commits = raft.get_commits().clone();
 
     //
     // RPC
@@ -59,33 +78,54 @@ async fn start(config: TauConfig, executor: Arc<Executor<'_>>) -> Result<()> {
 
     let recv_update_from_rpc: smol::Task<Result<()>> = executor.spawn(async move {
         loop {
-            let _task_info = rcv.recv().await?;
-            // XXX
+            let task_info = rcv.recv().await?;
+            raft_sender.send(task_info).await?;
         }
     });
 
-    listen_and_serve(server_config, rpc_interface, executor).await?;
+    let recv_update_from_raft: smol::Task<TaudResult<()>> = executor.spawn(async move {
+        loop {
+            // FIXME TODO
+            // this should update once receive rpc request from the tau-cli
+            sleep(1).await;
+            let recv_commits = commits.lock().await;
+
+            for task_info in recv_commits.iter() {
+                task_info.save()?;
+                if task_info.get_state() == "open" {
+                    task_info.activate()?;
+                } else {
+                    let mut mt = task_info.get_month_task()?;
+                    mt.remove(&task_info.get_ref_id());
+                }
+            }
+        }
+    });
+
+    let ex2 = executor.clone();
+    ex2.spawn(listen_and_serve(server_config, rpc_interface, executor.clone())).detach();
+
+    raft.start(p2p_settings.clone(), executor.clone()).await?;
 
     recv_update_from_rpc.cancel().await;
+    recv_update_from_raft.cancel().await;
     Ok(())
 }
 
 #[async_std::main]
 async fn main() -> Result<()> {
     let args = CliTaud::parse();
-    let matches = CliTaud::command().get_matches();
 
-    let verbosity_level = matches.occurrences_of("verbose");
-    let (lvl, conf) = log_config(verbosity_level)?;
+    let (lvl, conf) = log_config(args.verbose.into())?;
     TermLogger::init(lvl, conf, TerminalMode::Mixed, ColorChoice::Auto)?;
 
-    let config_path = get_config_path(args.config, "taud_config.toml")?;
+    let config_path = get_config_path(args.config.clone(), "taud_config.toml")?;
     spawn_config(&config_path, CONFIG_FILE_CONTENTS)?;
 
     let config: TauConfig = Config::<TauConfig>::load(config_path)?;
 
     let ex = Arc::new(Executor::new());
-    smol::block_on(ex.run(start(config, ex.clone())))
+    smol::block_on(ex.run(start(config, args, ex.clone())))
 }
 
 #[cfg(test)]

+ 4 - 0
bin/taud/src/task_info.rs

@@ -120,6 +120,10 @@ impl TaskInfo {
         mt.save()
     }
 
+    pub fn get_month_task(&self) -> TaudResult<MonthTasks> {
+        MonthTasks::load_or_create(&self.created_at, &self.settings)
+    }
+
     pub fn get_state(&self) -> String {
         if let Some(ev) = self.events.0.last() {
             ev.action.clone()

+ 16 - 1
bin/taud/src/util.rs

@@ -1,6 +1,7 @@
 use std::{
     fs::File,
     io::BufReader,
+    net::SocketAddr,
     path::{Path, PathBuf},
 };
 
@@ -71,8 +72,20 @@ pub struct Timestamp(pub i64);
 #[clap(name = "taud")]
 pub struct CliTaud {
     /// Sets a custom config file
-    #[clap(short, long)]
+    #[clap(long)]
     pub config: Option<String>,
+    /// Raft Accept address
+    #[clap(short, long)]
+    pub accept: Option<SocketAddr>,
+    /// Raft Seed node (repeatable)
+    #[clap(short, long)]
+    pub seed: Vec<SocketAddr>,
+    /// Raft Manual connection (repeatable)
+    #[clap(short, long)]
+    pub connect: Vec<SocketAddr>,
+    /// Raft Connection slots
+    #[clap(long, default_value_t = 0)]
+    pub slots: u32,
     /// Increase verbosity
     #[clap(short, parse(from_occurrences))]
     pub verbose: u8,
@@ -82,6 +95,8 @@ pub struct CliTaud {
 pub struct TauConfig {
     /// path to dataset
     pub dataset_path: String,
+    /// path to datastore  for raft
+    pub datastore_raft: String,
     /// Path to DER-formatted PKCS#12 archive. (used only with tls listener url)
     pub tls_identity_path: String,
     /// The address where taud should bind its RPC socket

+ 1 - 0
bin/taud/taud_config.toml

@@ -5,6 +5,7 @@
 
 # Path to the dataset 
 dataset_path = "~/.config/tau"
+datastore_raft = "~/.config/tau.db"
 
 # Path to DER-formatted PKCS#12 archive. (used only with tls url)
 # This can be created using openssl: