/* This file is part of DarkFi (https://dark.fi)
*
* Copyright (C) 2020-2026 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 .
*/
// use async_std::stream::from_iter;
use std::{
collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque},
path::PathBuf,
str::FromStr,
sync::Arc,
};
// use futures::stream::FuturesOrdered;
use blake3::Hash;
use darkfi_serial::{deserialize_async, serialize_async};
use event::Header;
use futures::{
// future,
stream::FuturesUnordered,
StreamExt,
};
use num_bigint::BigUint;
use sled_overlay::{sled, SledTreeOverlay};
use smol::{
lock::{OnceCell, RwLock},
Executor,
};
use tracing::{debug, error, info, warn};
use url::Url;
use crate::{
event_graph::util::{midnight_timestamp, replayer_log},
net::{channel::Channel, P2pPtr},
system::{msleep, Publisher, PublisherPtr, StoppableTask, StoppableTaskPtr, Subscription},
Error, Result,
};
#[cfg(feature = "rpc")]
use {
crate::rpc::{
jsonrpc::{JsonResponse, JsonResult},
util::json_map,
},
tinyjson::JsonValue::{self},
};
/// An event graph event
pub mod event;
pub use event::Event;
/// P2P protocol implementation for the Event Graph
pub mod proto;
use proto::{EventRep, EventReq, HeaderRep, HeaderReq, TipRep, TipReq};
/// Utility functions
pub mod util;
use util::{generate_genesis, millis_until_next_rotation, next_rotation_timestamp};
// Debugging event graph
pub mod deg;
use deg::DegEvent;
#[cfg(test)]
mod tests;
/// Initial genesis timestamp in millis (07 Sep 2023, 00:00:00 UTC)
/// Must always be UTC midnight.
pub const INITIAL_GENESIS: u64 = 1_694_044_800_000;
/// Genesis event contents
pub const GENESIS_CONTENTS: &[u8] = &[0x47, 0x45, 0x4e, 0x45, 0x53, 0x49, 0x53];
/// The number of parents an event is supposed to have.
pub const N_EVENT_PARENTS: usize = 5;
/// Allowed timestamp drift in milliseconds
const EVENT_TIME_DRIFT: u64 = 60_000;
/// Null event ID
pub const NULL_ID: Hash = Hash::from_bytes([0x00; blake3::OUT_LEN]);
/// Maximum number of DAGs to store, this should be configurable
pub const DAGS_MAX_NUMBER: i8 = 5;
/// Atomic pointer to an [`EventGraph`] instance.
pub type EventGraphPtr = Arc;
pub type LayerUTips = BTreeMap>;
#[derive(Clone)]
pub struct DAGStore {
db: sled::Db,
header_dags: HashMap,
main_dags: HashMap,
}
impl DAGStore {
pub async fn new(&self, sled_db: sled::Db, days_rotation: u64) -> Self {
let mut considered_trees = HashMap::new();
let mut considered_header_trees = HashMap::new();
if days_rotation > 0 {
// Create previous genesises if not existing, since they are deterministic.
for i in 1..=DAGS_MAX_NUMBER {
let i_days_ago = midnight_timestamp((i - DAGS_MAX_NUMBER).into());
let header =
Header { timestamp: i_days_ago, parents: [NULL_ID; N_EVENT_PARENTS], layer: 0 };
let genesis = Event { header, content: GENESIS_CONTENTS.to_vec() };
let tree_name = genesis.id().to_string();
let hdr_tree_name = format!("headers_{tree_name}");
let hdr_dag = sled_db.open_tree(hdr_tree_name).unwrap();
let dag = sled_db.open_tree(tree_name).unwrap();
if hdr_dag.is_empty() {
let mut overlay = SledTreeOverlay::new(&hdr_dag);
let header_se = serialize_async(&genesis.header).await;
// Add the header to the overlay
overlay.insert(genesis.id().as_bytes(), &header_se).unwrap();
// Aggregate changes into a single batch
let batch = overlay.aggregate().unwrap();
// Atomically apply the batch.
// Panic if something is corrupted.
if let Err(e) = hdr_dag.apply_batch(batch) {
panic!("Failed applying header_dag_insert batch to sled: {}", e);
}
}
if dag.is_empty() {
let mut overlay = SledTreeOverlay::new(&dag);
let event_se = serialize_async(&genesis).await;
// Add the event to the overlay
overlay.insert(genesis.id().as_bytes(), &event_se).unwrap();
// Aggregate changes into a single batch
let batch = overlay.aggregate().unwrap();
// Atomically apply the batch.
// Panic if something is corrupted.
if let Err(e) = dag.apply_batch(batch) {
panic!("Failed applying dag_insert batch to sled: {}", e);
}
}
let utips = self.find_unreferenced_tips(&dag).await;
considered_header_trees.insert(genesis.id(), (hdr_dag, utips.clone()));
considered_trees.insert(genesis.id(), (dag, utips));
}
} else {
let genesis = generate_genesis(0);
let tree_name = genesis.id().to_string();
let hdr_tree_name = format!("headers_{tree_name}");
let hdr_dag = sled_db.open_tree(hdr_tree_name).unwrap();
let dag = sled_db.open_tree(tree_name).unwrap();
if hdr_dag.is_empty() {
let mut overlay = SledTreeOverlay::new(&hdr_dag);
let header_se = serialize_async(&genesis.header).await;
// Add the header to the overlay
overlay.insert(genesis.id().as_bytes(), &header_se).unwrap();
// Aggregate changes into a single batch
let batch = overlay.aggregate().unwrap();
// Atomically apply the batch.
// Panic if something is corrupted.
if let Err(e) = hdr_dag.apply_batch(batch) {
panic!("Failed applying header_dag_insert batch to sled: {}", e);
}
}
if dag.is_empty() {
let mut overlay = SledTreeOverlay::new(&dag);
let event_se = serialize_async(&genesis).await;
// Add the event to the overlay
overlay.insert(genesis.id().as_bytes(), &event_se).unwrap();
// Aggregate changes into a single batch
let batch = overlay.aggregate().unwrap();
// Atomically apply the batch.
// Panic if something is corrupted.
if let Err(e) = dag.apply_batch(batch) {
panic!("Failed applying dag_insert batch to sled: {}", e);
}
}
let utips = self.find_unreferenced_tips(&dag).await;
considered_header_trees.insert(genesis.id(), (hdr_dag, utips.clone()));
considered_trees.insert(genesis.id(), (dag, utips));
}
Self { db: sled_db, header_dags: considered_header_trees, main_dags: considered_trees }
}
/// Adds a DAG into the set of DAGs and drops the oldest one if exeeding DAGS_MAX_NUMBER,
/// This is called if prune_task activates.
pub async fn add_dag(&mut self, dag_name: &str, genesis_event: &Event) {
debug!("add_dag::dags: {}", self.main_dags.len());
if self.main_dags.len() != self.header_dags.len() {
panic!("main dags length is not the same as header dags")
}
// TODO: sort dags by timestamp and drop the oldest
if self.main_dags.len() > DAGS_MAX_NUMBER.try_into().unwrap() {
while self.main_dags.len() >= DAGS_MAX_NUMBER.try_into().unwrap() {
debug!("[EVENTGRAPH] dropping oldest dag");
let sorted_dags = self.sort_dags().await;
// since dags are sorted in reverse
let oldest_tree = sorted_dags.last().unwrap().name();
let oldest_key = String::from_utf8_lossy(&oldest_tree);
let oldest_key = blake3::Hash::from_str(&oldest_key).unwrap();
let oldest_hdr_tree = self.header_dags.remove(&oldest_key).unwrap();
let oldest_tree = self.main_dags.remove(&oldest_key).unwrap();
self.db.drop_tree(oldest_hdr_tree.0.name()).unwrap();
self.db.drop_tree(oldest_tree.0.name()).unwrap();
}
}
// Insert genesis
let hdr_tree_name = format!("headers_{dag_name}");
let hdr_dag = self.get_dag(&hdr_tree_name);
hdr_dag
.insert(genesis_event.id().as_bytes(), serialize_async(&genesis_event.header).await)
.unwrap();
let dag = self.get_dag(dag_name);
dag.insert(genesis_event.id().as_bytes(), serialize_async(genesis_event).await).unwrap();
let utips = self.find_unreferenced_tips(&dag).await;
self.header_dags.insert(genesis_event.id(), (hdr_dag, utips.clone()));
self.main_dags.insert(genesis_event.id(), (dag, utips));
}
// Get a DAG providing its name.
pub fn get_dag(&self, dag_name: &str) -> sled::Tree {
self.db.open_tree(dag_name).unwrap()
}
/// Get {count} many DAGs.
pub async fn get_dags(&self, count: usize) -> Vec {
let sorted_dags = self.sort_dags().await;
sorted_dags.into_iter().take(count).collect()
}
/// Sort DAGs chronologically
async fn sort_dags(&self) -> Vec {
let mut vec_dags = vec![];
let dags = self
.main_dags
.iter()
.map(|x| {
let trees = x.1;
trees.0.clone()
})
.collect::>();
for dag in dags {
let genesis = dag.first().unwrap().unwrap().1;
let genesis_event: Event = deserialize_async(&genesis).await.unwrap();
vec_dags.push((genesis_event.header.timestamp, dag));
}
vec_dags.sort_by_key(|&(ts, _)| ts);
vec_dags.reverse();
vec_dags.into_iter().map(|(_, dag)| dag).collect()
}
/// Find the unreferenced tips in the current DAG state, mapped by their layers.
async fn find_unreferenced_tips(&self, dag: &sled::Tree) -> LayerUTips {
// First get all the event IDs
let mut tips = HashSet::new();
for iter_elem in dag.iter() {
let (id, _) = iter_elem.unwrap();
let id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
tips.insert(id);
}
// Iterate again to find unreferenced IDs
for iter_elem in dag.iter() {
let (_, event) = iter_elem.unwrap();
let event: Event = deserialize_async(&event).await.unwrap();
for parent in event.header.parents.iter() {
tips.remove(parent);
}
}
// Build the layers map
let mut map: LayerUTips = BTreeMap::new();
for tip in tips {
let event = self.fetch_event_from_dag(&tip, &dag).await.unwrap().unwrap();
if let Some(layer_tips) = map.get_mut(&event.header.layer) {
layer_tips.insert(tip);
} else {
let mut layer_tips = HashSet::new();
layer_tips.insert(tip);
map.insert(event.header.layer, layer_tips);
}
}
map
}
/// Fetch an event from the DAG
pub async fn fetch_event_from_dag(
&self,
event_id: &blake3::Hash,
dag: &sled::Tree,
) -> Result