Selaa lähdekoodia

event_graph: dag replay logs and db to be configurable, also user should choose to run in replay mode or not

dasman 2 vuotta sitten
vanhempi
sitoutus
37f85f49cd
4 muutettua tiedostoa jossa 33 lisäystä ja 20 poistoa
  1. 1 1
      src/event_graph/event.rs
  2. 12 2
      src/event_graph/mod.rs
  3. 4 1
      src/event_graph/tests.rs
  4. 16 16
      src/event_graph/util.rs

+ 1 - 1
src/event_graph/event.rs

@@ -216,7 +216,7 @@ mod tests {
         let ex = Arc::new(Executor::new());
         let p2p = P2p::new(Settings::default(), ex.clone()).await;
         let sled_db = sled::Config::new().temporary(true).open().unwrap();
-        EventGraph::new(p2p, sled_db, "dag", 1, ex).await
+        EventGraph::new(p2p, sled_db, "/tmp".into(), false, "dag", 1, ex).await
     }
 
     #[test]

+ 12 - 2
src/event_graph/mod.rs

@@ -19,6 +19,7 @@
 use std::{
     cmp::Ordering,
     collections::{BTreeMap, HashMap, HashSet, VecDeque},
+    path::PathBuf,
     sync::Arc,
     time::Duration,
 };
@@ -89,6 +90,10 @@ pub struct EventGraph {
     p2p: P2pPtr,
     /// Sled tree containing the DAG
     dag: sled::Tree,
+
+    datastore: PathBuf,
+
+    replay_mode: bool,
     /// The set of unreferenced DAG tips
     unreferenced_tips: RwLock<BTreeMap<u64, HashSet<blake3::Hash>>>,
     /// A `HashSet` containg event IDs and their 1-level parents.
@@ -120,6 +125,8 @@ impl EventGraph {
     pub async fn new(
         p2p: P2pPtr,
         sled_db: sled::Db,
+        datastore: PathBuf,
+        replay_mode: bool,
         dag_tree_name: &str,
         days_rotation: u64,
         ex: Arc<Executor<'_>>,
@@ -134,6 +141,8 @@ impl EventGraph {
         let self_ = Arc::new(Self {
             p2p,
             dag: dag.clone(),
+            datastore,
+            replay_mode,
             unreferenced_tips,
             broadcasted_ids,
             prune_task: OnceCell::new(),
@@ -569,8 +578,9 @@ impl EventGraph {
             // Add the event to the overlay
             overlay.insert(event_id.as_bytes(), &event_se)?;
 
-            replayer_log("insert".to_owned(), event_se).unwrap();
-
+            if self.replay_mode {
+                replayer_log(&self.datastore, "insert".to_owned(), event_se)?;
+            }
             // Note down the event ID to return
             ids.push(event_id);
         }

+ 4 - 1
src/event_graph/tests.rs

@@ -91,7 +91,10 @@ async fn spawn_node(
 
     let p2p = P2p::new(settings, ex.clone()).await;
     let sled_db = sled::Config::new().temporary(true).open().unwrap();
-    let event_graph = EventGraph::new(p2p.clone(), sled_db, "dag", 1, ex.clone()).await.unwrap();
+    let event_graph =
+        EventGraph::new(p2p.clone(), sled_db, "/tmp".into(), false, "dag", 1, ex.clone())
+            .await
+            .unwrap();
     *event_graph.synced.write().await = true;
     let event_graph_ = event_graph.clone();
 

+ 16 - 16
src/event_graph/util.rs

@@ -20,6 +20,7 @@ use std::{
     collections::HashMap,
     fs::{File, OpenOptions},
     io::Write,
+    path::PathBuf,
     time::UNIX_EPOCH,
 };
 
@@ -33,7 +34,7 @@ use crate::{
         jsonrpc::{ErrorCode, JsonError, JsonResponse, JsonResult},
         util::json_map,
     },
-    util::{encoding::base64, file::load_file, path::expand_path},
+    util::{encoding::base64, file::load_file},
     Result,
 };
 
@@ -131,14 +132,13 @@ pub(super) fn generate_genesis(days_rotation: u64) -> Event {
     }
 }
 
-pub(super) fn replayer_log(cmd: String, value: Vec<u8>) -> Result<()> {
-    let mut replayer_log_file = expand_path("/tmp")?;
-    replayer_log_file.push("replayer.log");
-    if !replayer_log_file.exists() {
-        File::create(&replayer_log_file)?;
+pub(super) fn replayer_log(datastore: &PathBuf, cmd: String, value: Vec<u8>) -> Result<()> {
+    let datastore = datastore.join("replayer.log");
+    if !datastore.exists() {
+        File::create(&datastore)?;
     };
 
-    let mut file = OpenOptions::new().append(true).open(&replayer_log_file)?;
+    let mut file = OpenOptions::new().append(true).open(&datastore)?;
     let v = base64::encode(&value);
     let f = format!("{cmd} {v}");
     writeln!(file, "{}", f)?;
@@ -146,22 +146,22 @@ pub(super) fn replayer_log(cmd: String, value: Vec<u8>) -> Result<()> {
     Ok(())
 }
 
-pub async fn recreate_from_replayer_log() -> JsonResult {
-    let mut replayer_log_file = expand_path("/tmp").unwrap();
-    replayer_log_file.push("replayer.log");
-    if !replayer_log_file.exists() {
-        error!("Error loading replaied log");
+pub async fn recreate_from_replayer_log(datastore: &PathBuf) -> JsonResult {
+    let log_path = datastore.join("replayer.log");
+    if !log_path.exists() {
+        error!("Error loading replayed log");
         return JsonResult::Error(JsonError::new(
             ErrorCode::ParseError,
-            Some("Error loading replaied log".to_string()),
+            Some("Error loading replayed log".to_string()),
             1,
         ))
     };
 
-    let reader = load_file(&replayer_log_file).unwrap();
+    let reader = load_file(&log_path).unwrap();
 
-    let datastore = expand_path("/tmp/replayed_db").unwrap();
-    let sled_db = sled::open(datastore).unwrap();
+    let db_datastore = datastore.join("replayed_db");
+
+    let sled_db = sled::open(db_datastore).unwrap();
     let dag = sled_db.open_tree("replayer").unwrap();
 
     for line in reader.lines() {