main.rs 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205
  1. use async_std::sync::{Arc, Mutex};
  2. use std::{
  3. fs::{create_dir_all, read_dir},
  4. path::{Path, PathBuf},
  5. };
  6. use async_executor::Executor;
  7. use futures::{select, FutureExt};
  8. use fxhash::FxHashMap;
  9. use log::{error, warn};
  10. use serde::Deserialize;
  11. use smol::future;
  12. use structopt::StructOpt;
  13. use structopt_toml::StructOptToml;
  14. use url::Url;
  15. use darkfi::{
  16. async_daemonize,
  17. net::{self, settings::SettingsOpt},
  18. raft::{NetMsg, ProtocolRaft, Raft, RaftSettings},
  19. rpc::server::listen_and_serve,
  20. util::{
  21. cli::{get_log_config, get_log_level, spawn_config},
  22. expand_path,
  23. file::{load_file, load_json_file, save_file, save_json_file},
  24. gen_id,
  25. path::get_config_path,
  26. },
  27. Error, Result,
  28. };
  29. mod error;
  30. mod jsonrpc;
  31. mod sequence;
  32. use error::DarkWikiResult;
  33. use jsonrpc::JsonRpcInterface;
  34. use sequence::{Operation, Sequence};
  35. pub const CONFIG_FILE: &str = "darkwiki.toml";
  36. pub const CONFIG_FILE_CONTENTS: &str = include_str!("../darkwiki.toml");
  37. pub const DOCS_PATH: &str = "~/darkwiki";
  38. /// darkwiki cli
  39. #[derive(Clone, Debug, Deserialize, StructOpt, StructOptToml)]
  40. #[serde(default)]
  41. #[structopt(name = "darkwiki")]
  42. pub struct Args {
  43. /// Sets a custom config file
  44. #[structopt(long)]
  45. pub config: Option<String>,
  46. /// Sets Datastore Path
  47. #[structopt(long, default_value = "~/.config/darkfi/darkwiki")]
  48. pub datastore: String,
  49. /// JSON-RPC listen URL
  50. #[structopt(long = "rpc", default_value = "tcp://127.0.0.1:13055")]
  51. pub rpc_listen: Url,
  52. #[structopt(flatten)]
  53. pub net: SettingsOpt,
  54. /// Increase verbosity
  55. #[structopt(short, parse(from_occurrences))]
  56. pub verbose: u8,
  57. }
  58. fn on_receive_operation(op: Operation, datastore_path: &Path) -> DarkWikiResult<()> {
  59. let json_files_path = datastore_path.join("files");
  60. let docs_path = PathBuf::from(expand_path(&DOCS_PATH)?);
  61. let id_path = json_files_path.join(op.id());
  62. let json_file = load_json_file::<Sequence>(&id_path);
  63. let mut seq: Sequence = if let Ok(file) = json_file { file } else { Sequence::new(&op.id()) };
  64. seq.add_op(&op)?;
  65. save_json_file::<Sequence>(&id_path, &seq)?;
  66. //let st = seq.apply();
  67. //save_file(&docs_path.join(op.id()), &st)?;
  68. Ok(())
  69. }
  70. fn on_receive_update(datastore_path: &Path) -> DarkWikiResult<Vec<Operation>> {
  71. let ret = vec![];
  72. let json_files_path = datastore_path.join("files");
  73. let docs_path = PathBuf::from(expand_path(&DOCS_PATH)?);
  74. let files = read_dir(&docs_path).unwrap();
  75. for file in files {
  76. let file_path = file.unwrap().path();
  77. let path = docs_path.join(&file_path);
  78. let edit = load_file(&path)?;
  79. if let Ok(seq) = load_json_file::<Sequence>(&json_files_path.join(&file_path)) {
  80. //
  81. // TODO the transformation should happen here
  82. //
  83. } else {
  84. let mut seq = Sequence::new(&gen_id(30));
  85. //seq.insert(0, &edit)?;
  86. save_json_file(&json_files_path.join(file_path), &seq)?;
  87. }
  88. }
  89. Ok(ret)
  90. }
  91. async fn start(
  92. update_notifier_rv: async_channel::Receiver<()>,
  93. raft_sender: async_channel::Sender<Operation>,
  94. raft_receiver: async_channel::Receiver<Operation>,
  95. datastore_path: PathBuf,
  96. ) -> DarkWikiResult<()> {
  97. loop {
  98. select! {
  99. _ = update_notifier_rv.recv().fuse() => {
  100. let ops = on_receive_update(&datastore_path)?;
  101. for op in ops {
  102. raft_sender.send(op).await.map_err(Error::from)?;
  103. }
  104. }
  105. op = raft_receiver.recv().fuse() => {
  106. let op = op.map_err(Error::from)?;
  107. on_receive_operation(op, &datastore_path)?;
  108. }
  109. }
  110. }
  111. }
  112. async_daemonize!(realmain);
  113. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  114. let datastore_path = expand_path(&settings.datastore)?;
  115. create_dir_all(expand_path(&DOCS_PATH)?)?;
  116. create_dir_all(datastore_path.join("files"))?;
  117. let (update_notifier_sx, update_notifier_rv) = async_channel::unbounded::<()>();
  118. //
  119. // RPC
  120. //
  121. let rpc_interface = Arc::new(JsonRpcInterface::new(update_notifier_sx));
  122. executor.spawn(listen_and_serve(settings.rpc_listen.clone(), rpc_interface)).detach();
  123. //
  124. // Raft
  125. //
  126. let net_settings = settings.net;
  127. let seen_net_msgs = Arc::new(Mutex::new(FxHashMap::default()));
  128. let datastore_raft = datastore_path.join("darkwiki.db");
  129. let raft_settings = RaftSettings { datastore_path: datastore_raft, ..RaftSettings::default() };
  130. let mut raft = Raft::<Operation>::new(raft_settings, seen_net_msgs.clone())?;
  131. executor
  132. .spawn(start(update_notifier_rv, raft.sender(), raft.receiver(), datastore_path.clone()))
  133. .detach();
  134. //
  135. // P2p setup
  136. //
  137. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<NetMsg>();
  138. let p2p = net::P2p::new(net_settings.into()).await;
  139. let p2p = p2p.clone();
  140. let registry = p2p.protocol_registry();
  141. let raft_node_id = raft.id();
  142. registry
  143. .register(net::SESSION_ALL, move |channel, p2p| {
  144. let raft_node_id = raft_node_id.clone();
  145. let sender = p2p_send_channel.clone();
  146. let seen_net_msgs_cloned = seen_net_msgs.clone();
  147. async move {
  148. ProtocolRaft::init(raft_node_id, channel, sender, p2p, seen_net_msgs_cloned).await
  149. }
  150. })
  151. .await;
  152. p2p.clone().start(executor.clone()).await?;
  153. executor.spawn(p2p.clone().run(executor.clone())).detach();
  154. //
  155. // Waiting Exit signal
  156. //
  157. let (signal, shutdown) = async_channel::bounded::<()>(1);
  158. ctrlc_async::set_async_handler(async move {
  159. warn!(target: "darkwiki", "Catch exit signal");
  160. // cleaning up tasks running in the background
  161. if let Err(e) = signal.send(()).await {
  162. error!("Error on sending exit signal: {}", e);
  163. }
  164. })
  165. .unwrap();
  166. raft.run(p2p.clone(), p2p_recv_channel.clone(), executor.clone(), shutdown.clone()).await?;
  167. Ok(())
  168. }