mod.rs 48 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2026 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. // use async_std::stream::from_iter;
  19. use std::{
  20. collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque},
  21. path::PathBuf,
  22. str::FromStr,
  23. sync::Arc,
  24. };
  25. // use futures::stream::FuturesOrdered;
  26. use blake3::Hash;
  27. use darkfi_serial::{deserialize_async, serialize_async};
  28. use event::Header;
  29. use futures::{
  30. // future,
  31. stream::FuturesUnordered,
  32. StreamExt,
  33. };
  34. use num_bigint::BigUint;
  35. use sled_overlay::{sled, SledTreeOverlay};
  36. use smol::{
  37. lock::{OnceCell, RwLock},
  38. Executor,
  39. };
  40. use tracing::{debug, error, info, warn};
  41. use url::Url;
  42. use crate::{
  43. event_graph::util::{midnight_timestamp, replayer_log},
  44. net::{channel::Channel, P2pPtr},
  45. system::{msleep, Publisher, PublisherPtr, StoppableTask, StoppableTaskPtr, Subscription},
  46. Error, Result,
  47. };
  48. #[cfg(feature = "rpc")]
  49. use {
  50. crate::rpc::{
  51. jsonrpc::{JsonResponse, JsonResult},
  52. util::json_map,
  53. },
  54. tinyjson::JsonValue::{self},
  55. };
  56. /// An event graph event
  57. pub mod event;
  58. pub use event::Event;
  59. /// P2P protocol implementation for the Event Graph
  60. pub mod proto;
  61. use proto::{EventRep, EventReq, HeaderRep, HeaderReq, TipRep, TipReq};
  62. /// Utility functions
  63. pub mod util;
  64. use util::{generate_genesis, millis_until_next_rotation, next_rotation_timestamp};
  65. // Debugging event graph
  66. pub mod deg;
  67. use deg::DegEvent;
  68. #[cfg(test)]
  69. mod tests;
  70. /// Initial genesis timestamp in millis (07 Sep 2023, 00:00:00 UTC)
  71. /// Must always be UTC midnight.
  72. pub const INITIAL_GENESIS: u64 = 1_694_044_800_000;
  73. /// Genesis event contents
  74. pub const GENESIS_CONTENTS: &[u8] = &[0x47, 0x45, 0x4e, 0x45, 0x53, 0x49, 0x53];
  75. /// The number of parents an event is supposed to have.
  76. pub const N_EVENT_PARENTS: usize = 5;
  77. /// Allowed timestamp drift in milliseconds
  78. const EVENT_TIME_DRIFT: u64 = 60_000;
  79. /// Null event ID
  80. pub const NULL_ID: Hash = Hash::from_bytes([0x00; blake3::OUT_LEN]);
  81. /// Maximum number of DAGs to store, this should be configurable
  82. pub const DAGS_MAX_NUMBER: i8 = 5;
  83. /// Atomic pointer to an [`EventGraph`] instance.
  84. pub type EventGraphPtr = Arc<EventGraph>;
  85. pub type LayerUTips = BTreeMap<u64, HashSet<blake3::Hash>>;
  86. #[derive(Clone)]
  87. pub struct DAGStore {
  88. db: sled::Db,
  89. header_dags: HashMap<Hash, (sled::Tree, LayerUTips)>,
  90. main_dags: HashMap<Hash, (sled::Tree, LayerUTips)>,
  91. }
  92. impl DAGStore {
  93. pub async fn new(&self, sled_db: sled::Db, days_rotation: u64) -> Self {
  94. let mut considered_trees = HashMap::new();
  95. let mut considered_header_trees = HashMap::new();
  96. if days_rotation > 0 {
  97. // Create previous genesises if not existing, since they are deterministic.
  98. for i in 1..=DAGS_MAX_NUMBER {
  99. let i_days_ago = midnight_timestamp((i - DAGS_MAX_NUMBER).into());
  100. let header =
  101. Header { timestamp: i_days_ago, parents: [NULL_ID; N_EVENT_PARENTS], layer: 0 };
  102. let genesis = Event { header, content: GENESIS_CONTENTS.to_vec() };
  103. let tree_name = genesis.id().to_string();
  104. let hdr_tree_name = format!("headers_{tree_name}");
  105. let hdr_dag = sled_db.open_tree(hdr_tree_name).unwrap();
  106. let dag = sled_db.open_tree(tree_name).unwrap();
  107. if hdr_dag.is_empty() {
  108. let mut overlay = SledTreeOverlay::new(&hdr_dag);
  109. let header_se = serialize_async(&genesis.header).await;
  110. // Add the header to the overlay
  111. overlay.insert(genesis.id().as_bytes(), &header_se).unwrap();
  112. // Aggregate changes into a single batch
  113. let batch = overlay.aggregate().unwrap();
  114. // Atomically apply the batch.
  115. // Panic if something is corrupted.
  116. if let Err(e) = hdr_dag.apply_batch(batch) {
  117. panic!("Failed applying header_dag_insert batch to sled: {}", e);
  118. }
  119. }
  120. if dag.is_empty() {
  121. let mut overlay = SledTreeOverlay::new(&dag);
  122. let event_se = serialize_async(&genesis).await;
  123. // Add the event to the overlay
  124. overlay.insert(genesis.id().as_bytes(), &event_se).unwrap();
  125. // Aggregate changes into a single batch
  126. let batch = overlay.aggregate().unwrap();
  127. // Atomically apply the batch.
  128. // Panic if something is corrupted.
  129. if let Err(e) = dag.apply_batch(batch) {
  130. panic!("Failed applying dag_insert batch to sled: {}", e);
  131. }
  132. }
  133. let utips = self.find_unreferenced_tips(&dag).await;
  134. considered_header_trees.insert(genesis.id(), (hdr_dag, utips.clone()));
  135. considered_trees.insert(genesis.id(), (dag, utips));
  136. }
  137. } else {
  138. let genesis = generate_genesis(0);
  139. let tree_name = genesis.id().to_string();
  140. let hdr_tree_name = format!("headers_{tree_name}");
  141. let hdr_dag = sled_db.open_tree(hdr_tree_name).unwrap();
  142. let dag = sled_db.open_tree(tree_name).unwrap();
  143. if hdr_dag.is_empty() {
  144. let mut overlay = SledTreeOverlay::new(&hdr_dag);
  145. let header_se = serialize_async(&genesis.header).await;
  146. // Add the header to the overlay
  147. overlay.insert(genesis.id().as_bytes(), &header_se).unwrap();
  148. // Aggregate changes into a single batch
  149. let batch = overlay.aggregate().unwrap();
  150. // Atomically apply the batch.
  151. // Panic if something is corrupted.
  152. if let Err(e) = hdr_dag.apply_batch(batch) {
  153. panic!("Failed applying header_dag_insert batch to sled: {}", e);
  154. }
  155. }
  156. if dag.is_empty() {
  157. let mut overlay = SledTreeOverlay::new(&dag);
  158. let event_se = serialize_async(&genesis).await;
  159. // Add the event to the overlay
  160. overlay.insert(genesis.id().as_bytes(), &event_se).unwrap();
  161. // Aggregate changes into a single batch
  162. let batch = overlay.aggregate().unwrap();
  163. // Atomically apply the batch.
  164. // Panic if something is corrupted.
  165. if let Err(e) = dag.apply_batch(batch) {
  166. panic!("Failed applying dag_insert batch to sled: {}", e);
  167. }
  168. }
  169. let utips = self.find_unreferenced_tips(&dag).await;
  170. considered_header_trees.insert(genesis.id(), (hdr_dag, utips.clone()));
  171. considered_trees.insert(genesis.id(), (dag, utips));
  172. }
  173. Self { db: sled_db, header_dags: considered_header_trees, main_dags: considered_trees }
  174. }
  175. /// Adds a DAG into the set of DAGs and drops the oldest one if exeeding DAGS_MAX_NUMBER,
  176. /// This is called if prune_task activates.
  177. pub async fn add_dag(&mut self, dag_name: &str, genesis_event: &Event) {
  178. debug!("add_dag::dags: {}", self.main_dags.len());
  179. if self.main_dags.len() != self.header_dags.len() {
  180. panic!("main dags length is not the same as header dags")
  181. }
  182. // TODO: sort dags by timestamp and drop the oldest
  183. if self.main_dags.len() > DAGS_MAX_NUMBER.try_into().unwrap() {
  184. while self.main_dags.len() >= DAGS_MAX_NUMBER.try_into().unwrap() {
  185. debug!("[EVENTGRAPH] dropping oldest dag");
  186. let sorted_dags = self.sort_dags().await;
  187. // since dags are sorted in reverse
  188. let oldest_tree = sorted_dags.last().unwrap().name();
  189. let oldest_key = String::from_utf8_lossy(&oldest_tree);
  190. let oldest_key = blake3::Hash::from_str(&oldest_key).unwrap();
  191. let oldest_hdr_tree = self.header_dags.remove(&oldest_key).unwrap();
  192. let oldest_tree = self.main_dags.remove(&oldest_key).unwrap();
  193. self.db.drop_tree(oldest_hdr_tree.0.name()).unwrap();
  194. self.db.drop_tree(oldest_tree.0.name()).unwrap();
  195. }
  196. }
  197. // Insert genesis
  198. let hdr_tree_name = format!("headers_{dag_name}");
  199. let hdr_dag = self.get_dag(&hdr_tree_name);
  200. hdr_dag
  201. .insert(genesis_event.id().as_bytes(), serialize_async(&genesis_event.header).await)
  202. .unwrap();
  203. let dag = self.get_dag(dag_name);
  204. dag.insert(genesis_event.id().as_bytes(), serialize_async(genesis_event).await).unwrap();
  205. let utips = self.find_unreferenced_tips(&dag).await;
  206. self.header_dags.insert(genesis_event.id(), (hdr_dag, utips.clone()));
  207. self.main_dags.insert(genesis_event.id(), (dag, utips));
  208. }
  209. // Get a DAG providing its name.
  210. pub fn get_dag(&self, dag_name: &str) -> sled::Tree {
  211. self.db.open_tree(dag_name).unwrap()
  212. }
  213. /// Get {count} many DAGs.
  214. pub async fn get_dags(&self, count: usize) -> Vec<sled::Tree> {
  215. let sorted_dags = self.sort_dags().await;
  216. sorted_dags.into_iter().take(count).collect()
  217. }
  218. /// Sort DAGs chronologically
  219. async fn sort_dags(&self) -> Vec<sled::Tree> {
  220. let mut vec_dags = vec![];
  221. let dags = self
  222. .main_dags
  223. .iter()
  224. .map(|x| {
  225. let trees = x.1;
  226. trees.0.clone()
  227. })
  228. .collect::<Vec<_>>();
  229. for dag in dags {
  230. let genesis = dag.first().unwrap().unwrap().1;
  231. let genesis_event: Event = deserialize_async(&genesis).await.unwrap();
  232. vec_dags.push((genesis_event.header.timestamp, dag));
  233. }
  234. vec_dags.sort_by_key(|&(ts, _)| ts);
  235. vec_dags.reverse();
  236. vec_dags.into_iter().map(|(_, dag)| dag).collect()
  237. }
  238. /// Find the unreferenced tips in the current DAG state, mapped by their layers.
  239. async fn find_unreferenced_tips(&self, dag: &sled::Tree) -> LayerUTips {
  240. // First get all the event IDs
  241. let mut tips = HashSet::new();
  242. for iter_elem in dag.iter() {
  243. let (id, _) = iter_elem.unwrap();
  244. let id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  245. tips.insert(id);
  246. }
  247. // Iterate again to find unreferenced IDs
  248. for iter_elem in dag.iter() {
  249. let (_, event) = iter_elem.unwrap();
  250. let event: Event = deserialize_async(&event).await.unwrap();
  251. for parent in event.header.parents.iter() {
  252. tips.remove(parent);
  253. }
  254. }
  255. // Build the layers map
  256. let mut map: LayerUTips = BTreeMap::new();
  257. for tip in tips {
  258. let event = self.fetch_event_from_dag(&tip, &dag).await.unwrap().unwrap();
  259. if let Some(layer_tips) = map.get_mut(&event.header.layer) {
  260. layer_tips.insert(tip);
  261. } else {
  262. let mut layer_tips = HashSet::new();
  263. layer_tips.insert(tip);
  264. map.insert(event.header.layer, layer_tips);
  265. }
  266. }
  267. map
  268. }
  269. /// Fetch an event from the DAG
  270. pub async fn fetch_event_from_dag(
  271. &self,
  272. event_id: &blake3::Hash,
  273. dag: &sled::Tree,
  274. ) -> Result<Option<Event>> {
  275. let Some(bytes) = dag.get(event_id.as_bytes())? else {
  276. return Ok(None);
  277. };
  278. let event: Event = deserialize_async(&bytes).await?;
  279. return Ok(Some(event))
  280. }
  281. }
  282. enum PeerStatus {
  283. Free,
  284. Busy,
  285. Failed,
  286. }
  287. /// An Event Graph instance
  288. pub struct EventGraph {
  289. /// Pointer to the P2P network instance
  290. p2p: P2pPtr,
  291. /// Sled tree containing the headers
  292. dag_store: RwLock<DAGStore>,
  293. /// Replay logs path.
  294. datastore: PathBuf,
  295. /// Run in replay_mode where if set we log Sled DB instructions
  296. /// into `datastore`, useful to reacreate a faulty DAG to debug.
  297. replay_mode: bool,
  298. /// A `HashSet` containg event IDs and their 1-level parents.
  299. /// These come from the events we've sent out using `EventPut`.
  300. /// They are used with `EventReq` to decide if we should reply
  301. /// or not. Additionally it is also used when we broadcast the
  302. /// `TipRep` message telling peers about our unreferenced tips.
  303. broadcasted_ids: RwLock<HashSet<Hash>>,
  304. /// DAG Pruning Task
  305. pub prune_task: OnceCell<StoppableTaskPtr>,
  306. /// Event publisher, this notifies whenever an event is
  307. /// inserted into the DAG
  308. pub event_pub: PublisherPtr<Event>,
  309. /// Current genesis event
  310. pub current_genesis: RwLock<Event>,
  311. /// Currently configured DAG rotation, in days
  312. days_rotation: u64,
  313. /// Flag signalling DAG has finished initial sync
  314. pub synced: RwLock<bool>,
  315. /// Enable graph debugging
  316. pub deg_enabled: RwLock<bool>,
  317. /// The publisher for which we can give deg info over
  318. deg_publisher: PublisherPtr<DegEvent>,
  319. /// Run in replay_mode where if set we log Sled DB instructions
  320. /// into `datastore`, useful to reacreate a faulty DAG to debug.
  321. fast_mode: bool,
  322. }
  323. impl EventGraph {
  324. /// Create a new [`EventGraph`] instance, creates a new Genesis
  325. /// event and checks if it
  326. /// is containd in DAG, if not prunes DAG, may also start a pruning
  327. /// task based on `days_rotation`, and return an atomic instance of
  328. /// `Self`
  329. /// * `p2p` atomic pointer to p2p.
  330. /// * `sled_db` sled DB instance.
  331. /// * `datastore` path where we should log db instrucion if run in
  332. /// replay mode.
  333. /// * `replay_mode` set the flag to keep a log of db instructions.
  334. /// * `days_rotation` marks the lifetime of the DAG before it's
  335. /// pruned.
  336. pub async fn new(
  337. p2p: P2pPtr,
  338. sled_db: sled::Db,
  339. datastore: PathBuf,
  340. replay_mode: bool,
  341. fast_mode: bool,
  342. days_rotation: u64,
  343. ex: Arc<Executor<'_>>,
  344. ) -> Result<EventGraphPtr> {
  345. let broadcasted_ids = RwLock::new(HashSet::new());
  346. let event_pub = Publisher::new();
  347. // Create the current genesis event based on the `days_rotation`
  348. let current_genesis = generate_genesis(days_rotation);
  349. let current_dag_tree_name = current_genesis.id().to_string();
  350. let dag_store = DAGStore {
  351. db: sled_db.clone(),
  352. header_dags: HashMap::default(),
  353. main_dags: HashMap::default(),
  354. }
  355. .new(sled_db, days_rotation)
  356. .await;
  357. let self_ = Arc::new(Self {
  358. p2p,
  359. dag_store: RwLock::new(dag_store.clone()),
  360. datastore,
  361. replay_mode,
  362. fast_mode,
  363. broadcasted_ids,
  364. prune_task: OnceCell::new(),
  365. event_pub,
  366. current_genesis: RwLock::new(current_genesis.clone()),
  367. days_rotation,
  368. synced: RwLock::new(false),
  369. deg_enabled: RwLock::new(false),
  370. deg_publisher: Publisher::new(),
  371. });
  372. // Check if we have it in our DAG.
  373. // If not, we can prune the DAG and insert this new genesis event.
  374. let dag = dag_store.get_dag(&current_dag_tree_name);
  375. if !dag.contains_key(current_genesis.id().as_bytes())? {
  376. info!(
  377. target: "event_graph::new",
  378. "[EVENTGRAPH] DAG does not contain current genesis, pruning existing data",
  379. );
  380. self_.dag_prune(current_genesis).await?;
  381. }
  382. // Spawn the DAG pruning task
  383. if days_rotation > 0 {
  384. let prune_task = StoppableTask::new();
  385. let _ = self_.prune_task.set(prune_task.clone()).await;
  386. prune_task.clone().start(
  387. self_.clone().dag_prune_task(days_rotation),
  388. |res| async move {
  389. match res {
  390. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  391. Err(e) => error!(target: "event_graph::_handle_stop", "[EVENTGRAPH] Failed stopping prune task: {e}")
  392. }
  393. },
  394. Error::DetachedTaskStopped,
  395. ex.clone(),
  396. );
  397. }
  398. Ok(self_)
  399. }
  400. pub fn days_rotation(&self) -> u64 {
  401. self.days_rotation
  402. }
  403. /// Sync the DAG from connected peers
  404. pub async fn dag_sync(&self, dag: sled::Tree, fast_mode: bool) -> Result<()> {
  405. // We do an optimistic sync where we ask all our connected peers for
  406. // the latest layer DAG tips (unreferenced events) and then we accept
  407. // the ones we see the most times.
  408. // * Compare received tips with local ones, identify which we are missing.
  409. // * Request these from peers
  410. // * Recursively request these backward
  411. //
  412. // Verification:
  413. // * Timestamps should go backwards
  414. // * Cross-check with multiple peers, this means we should request the
  415. // same event from multiple peers and make sure it is the same.
  416. // * Since we should be pruning, if we're not synced after some reasonable
  417. // amount of iterations, these could be faulty peers and we can try again
  418. // from the beginning
  419. let dag_name = String::from_utf8_lossy(&dag.name()).to_string();
  420. // Get references to all our peers.
  421. let channels = self.p2p.hosts().peers();
  422. let mut communicated_peers = channels.len();
  423. info!(
  424. target: "event_graph::dag_sync",
  425. "[EVENTGRAPH] Syncing DAG from {communicated_peers} peers..."
  426. );
  427. let comms_timeout = self.p2p.settings().read().await.outbound_connect_timeout_max();
  428. // Here we keep track of the tips, their layers and how many time we've seen them.
  429. let mut tips: HashMap<Hash, (u64, usize)> = HashMap::new();
  430. // Let's first ask all of our peers for their tips and collect them
  431. // in our hashmap above.
  432. for channel in channels.iter() {
  433. let url = channel.display_address();
  434. let tip_rep_sub = match channel.subscribe_msg::<TipRep>().await {
  435. Ok(v) => v,
  436. Err(e) => {
  437. error!(
  438. target: "event_graph::dag_sync",
  439. "[EVENTGRAPH] Sync: Couldn't subscribe TipReq for peer {url}, skipping ({e})"
  440. );
  441. communicated_peers -= 1;
  442. continue
  443. }
  444. };
  445. if let Err(e) = channel.send(&TipReq(dag_name.clone())).await {
  446. error!(
  447. target: "event_graph::dag_sync",
  448. "[EVENTGRAPH] Sync: Couldn't contact peer {url}, skipping ({e})"
  449. );
  450. communicated_peers -= 1;
  451. continue
  452. };
  453. // Node waits for response
  454. let Ok(peer_tips) = tip_rep_sub.receive_with_timeout(comms_timeout).await else {
  455. error!(
  456. target: "event_graph::dag_sync",
  457. "[EVENTGRAPH] Sync: Peer {url} didn't reply with tips in time, skipping"
  458. );
  459. communicated_peers -= 1;
  460. continue
  461. };
  462. let peer_tips: &BTreeMap<u64, HashSet<Hash>> = &peer_tips.0;
  463. // Note down the seen tips
  464. for (layer, layer_tips) in peer_tips {
  465. for tip in layer_tips {
  466. if let Some(seen_tip) = tips.get_mut(tip) {
  467. seen_tip.1 += 1;
  468. } else {
  469. tips.insert(*tip, (*layer, 1));
  470. }
  471. }
  472. }
  473. }
  474. // After we've communicated all the peers, let's see what happened.
  475. if tips.is_empty() {
  476. error!(
  477. target: "event_graph::dag_sync",
  478. "[EVENTGRAPH] Sync: Could not find any DAG tips",
  479. );
  480. return Err(Error::DagSyncFailed)
  481. }
  482. // We know the number of peers we've communicated with,
  483. // so we will consider events we saw at more than 2/3 of
  484. // those peers.
  485. let consideration_threshold = communicated_peers * 2 / 3;
  486. let mut considered_tips = HashSet::new();
  487. for (tip, (_, amount)) in tips.iter() {
  488. if amount > &consideration_threshold {
  489. considered_tips.insert(*tip);
  490. }
  491. }
  492. drop(tips);
  493. if fast_mode {
  494. // Now begin fetching the events backwards.
  495. let mut missing_parents = HashSet::new();
  496. for tip in considered_tips.iter() {
  497. assert!(tip != &NULL_ID);
  498. if !dag.contains_key(tip.as_bytes()).unwrap() {
  499. missing_parents.insert(*tip);
  500. }
  501. }
  502. if missing_parents.is_empty() {
  503. *self.synced.write().await = true;
  504. info!(target: "event_graph::dag_sync", "[EVENTGRAPH] DAG synced successfully!");
  505. return Ok(())
  506. }
  507. }
  508. // Header sync first
  509. // TODO: requesting headers should be in a way that we wouldn't
  510. // recieve the same header(s) again, by sending our tip, other
  511. // nodes should send back the ones after it
  512. let hdr_tree_name = format!("headers_{dag_name}");
  513. let header_dag = self.dag_store.read().await.get_dag(&hdr_tree_name);
  514. let mut headers_requests = FuturesUnordered::new();
  515. for channel in channels.iter() {
  516. headers_requests.push(request_header(&channel, dag_name.clone(), comms_timeout))
  517. }
  518. while let Some(peer_headers) = headers_requests.next().await {
  519. self.header_dag_insert(peer_headers?, &dag_name).await?
  520. }
  521. // start download payload
  522. if !fast_mode {
  523. info!(target: "event_graph::dag_sync()", "[EVENTGRAPH] Fetching events");
  524. let mut header_sorted = vec![];
  525. for iter_elem in header_dag.iter() {
  526. let (_, val) = iter_elem.unwrap();
  527. let val: Header = deserialize_async(&val).await.unwrap();
  528. header_sorted.push(val);
  529. }
  530. header_sorted.sort_by(|x, y| y.layer.cmp(&x.layer));
  531. info!(target: "event_graph::dag_sync()", "[EVENTGRAPH] Retrieving {} Events", header_sorted.len());
  532. // Implement parallel download of events with a batch size
  533. let batch = 20;
  534. // Mapping of the chunk group id to the chunk, using a BTreeMap help us to
  535. // prioritize the older headers when our request fails and we retry
  536. let mut chunks: BTreeMap<usize, Vec<blake3::Hash>> = BTreeMap::new();
  537. for (i, chunk) in header_sorted.chunks(batch).enumerate() {
  538. chunks.insert(i, chunk.iter().map(|h| h.id()).collect());
  539. }
  540. let mut remaining_chunk_ids: BTreeSet<usize> = chunks.keys().cloned().collect();
  541. // Mapping of the chunk group id to the received events, using a BTreeMap help
  542. // us to verify and insert the events in order
  543. let mut received_events: BTreeMap<usize, Vec<Event>> = BTreeMap::new();
  544. // Track peers status so that we don't send a new request to the same peer before they
  545. // finish the first or send to a failed peer
  546. let mut peer_status: HashMap<Url, PeerStatus> = HashMap::new();
  547. let mut retrieved_count = 0;
  548. let mut futures = FuturesUnordered::new();
  549. while retrieved_count < header_sorted.len() {
  550. // Retrieve peers in each loop so we don't send requests to a closed channel
  551. let mut free_channels = vec![];
  552. let mut busy_channels = 0;
  553. self.p2p.hosts().peers().iter().for_each(|channel| {
  554. if let Some(status) = peer_status.get(channel.address()) {
  555. match status {
  556. PeerStatus::Free => free_channels.push(channel.clone()),
  557. PeerStatus::Busy => busy_channels += 1,
  558. _ => {}
  559. }
  560. } else {
  561. peer_status.insert(channel.address().clone(), PeerStatus::Free);
  562. free_channels.push(channel.clone());
  563. }
  564. });
  565. // We don't have any channels we can assign to or wait to get response from
  566. if free_channels.is_empty() && busy_channels == 0 {
  567. return Err(Error::DagSyncFailed);
  568. }
  569. // We will distribute the remaining chunks to each channel
  570. let requested_chunks_len =
  571. std::cmp::min(free_channels.len(), remaining_chunk_ids.len());
  572. let requested_chunk_ids: Vec<usize> =
  573. remaining_chunk_ids.iter().take(requested_chunks_len).copied().collect();
  574. for (i, chunk_id) in requested_chunk_ids.iter().enumerate() {
  575. futures.push(request_event(
  576. free_channels[i].clone(),
  577. chunks.get(chunk_id).unwrap().clone(),
  578. *chunk_id,
  579. comms_timeout,
  580. ));
  581. remaining_chunk_ids.remove(chunk_id);
  582. peer_status.insert(free_channels[i].address().clone(), PeerStatus::Busy);
  583. }
  584. info!(target: "event_graph::dag_sync()", "[EVENTGRAPH] Retrieving Events from {} peers", futures.len());
  585. if let Some(resp) = futures.next().await {
  586. let (events, chunk_id, channel) = resp;
  587. if let Ok(events) = events {
  588. retrieved_count += events.len();
  589. received_events.insert(chunk_id, events.clone());
  590. peer_status.insert(channel.address().clone(), PeerStatus::Free);
  591. } else {
  592. remaining_chunk_ids.insert(chunk_id);
  593. peer_status.insert(channel.address().clone(), PeerStatus::Failed);
  594. }
  595. info!(target: "event_graph::dag_sync()", "[EVENTGRAPH] Retrieved Events: {}/{}", retrieved_count, header_sorted.len());
  596. }
  597. }
  598. let mut verified_count = 0;
  599. for (_, chunk) in received_events {
  600. verified_count += chunk.len();
  601. self.dag_insert(&chunk, &dag_name).await?;
  602. info!(target: "event_graph::dag_sync()", "[EVENTGRAPH] Verified Events: {}/{}", verified_count, retrieved_count);
  603. }
  604. }
  605. // <-- end download payload
  606. *self.synced.write().await = true;
  607. info!(target: "event_graph::dag_sync()", "[EVENTGRAPH] DAG synced successfully!");
  608. Ok(())
  609. }
  610. /// Choose how many dags to sync
  611. pub async fn sync_selected(&self, count: usize, fast_mode: bool) -> Result<()> {
  612. let mut dags_to_sync = self.dag_store.read().await.get_dags(count).await;
  613. // Since get_dags() return sorted dags in reverse
  614. dags_to_sync.reverse();
  615. for dag in dags_to_sync {
  616. match self.dag_sync(dag, fast_mode).await {
  617. Ok(()) => continue,
  618. Err(e) => {
  619. return Err(e);
  620. }
  621. }
  622. }
  623. Ok(())
  624. }
  625. /// Atomically prune the DAG and insert the given event as genesis.
  626. async fn dag_prune(&self, genesis_event: Event) -> Result<()> {
  627. debug!(target: "event_graph::dag_prune", "Pruning DAG...");
  628. // Acquire exclusive locks to unreferenced_tips, broadcasted_ids and
  629. // current_genesis while this operation is happening. We do this to
  630. // ensure that during the pruning operation, no other operations are
  631. // able to access the intermediate state which could lead to producing
  632. // the wrong state after pruning.
  633. let mut broadcasted_ids = self.broadcasted_ids.write().await;
  634. let mut current_genesis = self.current_genesis.write().await;
  635. let dag_name = genesis_event.id().to_string();
  636. self.dag_store.write().await.add_dag(&dag_name, &genesis_event).await;
  637. // Clear bcast ids
  638. *current_genesis = genesis_event;
  639. *broadcasted_ids = HashSet::new();
  640. drop(broadcasted_ids);
  641. drop(current_genesis);
  642. debug!(target: "event_graph::dag_prune", "DAG pruned successfully");
  643. Ok(())
  644. }
  645. /// Background task periodically pruning the DAG.
  646. async fn dag_prune_task(self: Arc<Self>, days_rotation: u64) -> Result<()> {
  647. // The DAG should periodically be pruned. This can be a configurable
  648. // parameter. By pruning, we should deterministically replace the
  649. // genesis event (can use a deterministic timestamp) and drop everything
  650. // in the DAG, leaving just the new genesis event.
  651. debug!(target: "event_graph::dag_prune_task", "Spawned background DAG pruning task");
  652. loop {
  653. // Find the next rotation timestamp:
  654. let next_rotation = next_rotation_timestamp(INITIAL_GENESIS, days_rotation);
  655. let header =
  656. Header { timestamp: next_rotation, parents: [NULL_ID; N_EVENT_PARENTS], layer: 0 };
  657. // Prepare the new genesis event
  658. let current_genesis = Event { header, content: GENESIS_CONTENTS.to_vec() };
  659. // Sleep until it's time to rotate.
  660. let s = millis_until_next_rotation(next_rotation);
  661. debug!(target: "event_graph::dag_prune_task", "Sleeping {s}ms until next DAG prune");
  662. msleep(s).await;
  663. debug!(target: "event_graph::dag_prune_task", "Rotation period reached");
  664. // Trigger DAG prune
  665. self.dag_prune(current_genesis).await?;
  666. }
  667. }
  668. /// Atomically insert given events into the DAG and return the event IDs.
  669. /// All provided events must be valid. An overlay is used over the DAG tree,
  670. /// temporary writting each event in order. After all events have been
  671. /// validated and inserted successfully, we write the overlay to sled.
  672. /// This will append the new events into the unreferenced tips set, and
  673. /// remove the events' parents from it. It will also append the events'
  674. /// level-1 parents to the `broadcasted_ids` set, so the P2P protocol
  675. /// knows that any requests for them are actually legitimate.
  676. /// TODO: The `broadcasted_ids` set should periodically be pruned, when
  677. /// some sensible time has passed after broadcasting the event.
  678. pub async fn dag_insert(&self, events: &[Event], dag_name: &str) -> Result<Vec<Hash>> {
  679. // Sanity check
  680. if events.is_empty() {
  681. return Ok(vec![])
  682. }
  683. // Acquire exclusive locks to `broadcasted_ids`
  684. let dag_name_hash = blake3::Hash::from_str(dag_name).unwrap();
  685. let mut broadcasted_ids = self.broadcasted_ids.write().await;
  686. let main_dag = self.dag_store.read().await.get_dag(dag_name);
  687. let hdr_tree_name = format!("headers_{dag_name}");
  688. let header_dag = self.dag_store.read().await.get_dag(&hdr_tree_name);
  689. // Here we keep the IDs to return
  690. let mut ids = Vec::with_capacity(events.len());
  691. // Create an overlay over the DAG tree
  692. let mut overlay = SledTreeOverlay::new(&main_dag);
  693. // Grab genesis timestamp
  694. // let genesis_timestamp = self.current_genesis.read().await.header.timestamp;
  695. // Iterate over given events to validate them and
  696. // write them to the overlay
  697. for event in events {
  698. let event_id = event.id();
  699. if event.header.layer == 0 {
  700. break
  701. }
  702. debug!(
  703. target: "event_graph::dag_insert",
  704. "Inserting event {event_id} into the DAG layer: {}", event.header.layer
  705. );
  706. // check if we already have the event
  707. if main_dag.contains_key(event_id.as_bytes())? {
  708. continue
  709. }
  710. // check if its header is in header's store
  711. if !header_dag.contains_key(event_id.as_bytes())? {
  712. continue
  713. }
  714. if !event.dag_validate(&header_dag).await? {
  715. error!(target: "event_graph::dag_insert()", "Event {} is invalid!", event_id);
  716. return Err(Error::EventIsInvalid)
  717. }
  718. let event_se = serialize_async(event).await;
  719. // Add the event to the overlay
  720. overlay.insert(event_id.as_bytes(), &event_se)?;
  721. if self.replay_mode {
  722. replayer_log(&self.datastore, "insert".to_owned(), event_se)?;
  723. }
  724. // Note down the event ID to return
  725. ids.push(event_id);
  726. }
  727. // Aggregate changes into a single batch
  728. let batch = match overlay.aggregate() {
  729. Some(x) => x,
  730. None => return Ok(vec![]),
  731. };
  732. // Atomically apply the batch.
  733. // Panic if something is corrupted.
  734. if let Err(e) = main_dag.apply_batch(batch) {
  735. panic!("Failed applying dag_insert batch to sled: {e}");
  736. }
  737. drop(main_dag);
  738. drop(header_dag);
  739. let mut dag_store = self.dag_store.write().await;
  740. let (_, unreferenced_tips) = &mut dag_store.main_dags.get_mut(&dag_name_hash).unwrap();
  741. // Iterate over given events to update references and
  742. // send out notifications about them
  743. for event in events {
  744. let event_id = event.id();
  745. // Update the unreferenced DAG tips set
  746. debug!(
  747. target: "event_graph::dag_insert",
  748. "Event {event_id} parents {:#?}", event.header.parents,
  749. );
  750. for parent_id in event.header.parents.iter() {
  751. if parent_id != &NULL_ID {
  752. debug!(
  753. target: "event_graph::dag_insert",
  754. "Removing {parent_id} from unreferenced_tips"
  755. );
  756. // Iterate over unreferenced tips in previous layers
  757. // and remove the parent
  758. // NOTE: this might be too exhaustive, but the
  759. // assumption is that previous layers unreferenced
  760. // tips will be few.
  761. for (layer, tips) in unreferenced_tips.iter_mut() {
  762. if layer >= &event.header.layer {
  763. continue
  764. }
  765. tips.remove(parent_id);
  766. }
  767. broadcasted_ids.insert(*parent_id);
  768. }
  769. }
  770. unreferenced_tips.retain(|_, tips| !tips.is_empty());
  771. debug!(
  772. target: "event_graph::dag_insert",
  773. "Adding {event_id} to unreferenced tips"
  774. );
  775. if let Some(layer_tips) = unreferenced_tips.get_mut(&event.header.layer) {
  776. layer_tips.insert(event_id);
  777. } else {
  778. let mut layer_tips = HashSet::new();
  779. layer_tips.insert(event_id);
  780. unreferenced_tips.insert(event.header.layer, layer_tips);
  781. }
  782. // Send out notifications about the new event
  783. self.event_pub.notify(event.clone()).await;
  784. }
  785. // Drop the exclusive locks
  786. drop(dag_store);
  787. drop(broadcasted_ids);
  788. let mut dag_store = self.dag_store.write().await;
  789. dag_store.header_dags.get_mut(&dag_name_hash).unwrap().1 =
  790. dag_store.main_dags.get(&dag_name_hash).unwrap().1.clone();
  791. drop(dag_store);
  792. Ok(ids)
  793. }
  794. pub async fn header_dag_insert(&self, headers: Vec<Header>, dag_name: &str) -> Result<()> {
  795. let hdr_tree_name = format!("headers_{dag_name}");
  796. let header_dag = self.dag_store.read().await.get_dag(&hdr_tree_name);
  797. // Create an overlay over the DAG tree
  798. let mut overlay = SledTreeOverlay::new(&header_dag);
  799. // Acquire exclusive locks to `unreferenced_tips and broadcasted_ids`
  800. // let mut unreferenced_header = self.unreferenced_tips.write().await;
  801. // let mut broadcasted_ids = self.broadcasted_ids.write().await;
  802. let mut hdrs = headers;
  803. hdrs.sort_by(|x, y| x.layer.cmp(&y.layer));
  804. // Iterate over given events to validate them and
  805. // write them to the overlay
  806. for header in hdrs {
  807. let header_id = header.id();
  808. if header.layer == 0 {
  809. continue
  810. }
  811. debug!(
  812. target: "event_graph::header_dag_insert()",
  813. "Inserting header {} into the DAG", header_id,
  814. );
  815. if !header.validate(&header_dag, self.days_rotation, Some(&overlay)).await? {
  816. error!(target: "event_graph::header_dag_insert()", "Header {} is invalid!", header_id);
  817. return Err(Error::HeaderIsInvalid)
  818. }
  819. let header_se = serialize_async(&header).await;
  820. // Add the event to the overlay
  821. overlay.insert(header_id.as_bytes(), &header_se)?;
  822. }
  823. // Aggregate changes into a single batch
  824. let batch = match overlay.aggregate() {
  825. Some(x) => x,
  826. None => return Ok(()),
  827. };
  828. // Atomically apply the batch.
  829. // Panic if something is corrupted.
  830. if let Err(e) = header_dag.apply_batch(batch) {
  831. panic!("Failed applying dag_insert batch to sled: {}", e);
  832. }
  833. Ok(())
  834. }
  835. /// Search and fetch an event through all DAGs
  836. pub async fn fetch_event_from_dags(&self, event_id: &blake3::Hash) -> Result<Option<Event>> {
  837. let store = self.dag_store.read().await;
  838. for tree_elem in store.main_dags.clone() {
  839. let dag_name = tree_elem.0.to_string();
  840. let Some(bytes) = store.get_dag(&dag_name).get(event_id.as_bytes())? else {
  841. continue;
  842. };
  843. let event: Event = deserialize_async(&bytes).await?;
  844. return Ok(Some(event))
  845. }
  846. Ok(None)
  847. }
  848. /// Get next layer along with its N_EVENT_PARENTS from the unreferenced
  849. /// tips of the DAG. Since tips are mapped by their layer, we go backwards
  850. /// until we fill the vector, ensuring we always use latest layers tips as
  851. /// parents.
  852. async fn get_next_layer_with_parents(
  853. &self,
  854. dag_name: &Hash,
  855. ) -> (u64, [blake3::Hash; N_EVENT_PARENTS]) {
  856. let store = self.dag_store.read().await;
  857. let (_, unreferenced_tips) = store.header_dags.get(dag_name).unwrap();
  858. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  859. let mut index = 0;
  860. 'outer: for (_, tips) in unreferenced_tips.iter().rev() {
  861. for tip in tips.iter() {
  862. parents[index] = *tip;
  863. index += 1;
  864. if index >= N_EVENT_PARENTS {
  865. break 'outer;
  866. }
  867. }
  868. }
  869. let next_layer = unreferenced_tips.last_key_value().unwrap().0 + 1;
  870. assert!(parents.iter().any(|x| x != &NULL_ID));
  871. (next_layer, parents)
  872. }
  873. /// Internal function used for DAG sorting.
  874. async fn get_unreferenced_tips_sorted(&self) -> Vec<[blake3::Hash; N_EVENT_PARENTS]> {
  875. let mut vec_tips = vec![];
  876. let mut tips_sorted = [NULL_ID; N_EVENT_PARENTS];
  877. for (i, _) in self.dag_store.read().await.header_dags.iter() {
  878. let (_, tips) = self.get_next_layer_with_parents(&i).await;
  879. // Convert the hash to BigUint for sorting
  880. let mut sorted: Vec<_> =
  881. tips.iter().map(|x| BigUint::from_bytes_be(x.as_bytes())).collect();
  882. sorted.sort_unstable();
  883. // Convert back to blake3
  884. for (i, id) in sorted.iter().enumerate() {
  885. let mut bytes = id.to_bytes_be();
  886. // Ensure we have 32 bytes
  887. while bytes.len() < blake3::OUT_LEN {
  888. bytes.insert(0, 0);
  889. }
  890. tips_sorted[i] = blake3::Hash::from_bytes(bytes.try_into().unwrap());
  891. }
  892. vec_tips.push(tips_sorted);
  893. }
  894. vec_tips
  895. }
  896. // TODO: Fix fetching all events from all dags and then order and retrun them
  897. /// Perform a topological sort of the DAG.
  898. pub async fn order_events(&self) -> Vec<Event> {
  899. let mut ordered_events = VecDeque::new();
  900. let mut visited = HashSet::new();
  901. for i in self.get_unreferenced_tips_sorted().await {
  902. for tip in i {
  903. if !visited.contains(&tip) && tip != NULL_ID {
  904. let tip = self.fetch_event_from_dags(&tip).await.unwrap().unwrap();
  905. ordered_events.extend(self.dfs_topological_sort(tip, &mut visited).await);
  906. }
  907. }
  908. }
  909. let mut ord_events_vec = ordered_events.make_contiguous().to_vec();
  910. // Order events by timestamp.
  911. ord_events_vec.sort_unstable_by(|a, b| a.1.header.timestamp.cmp(&b.1.header.timestamp));
  912. ord_events_vec.iter().map(|a| a.1.clone()).collect::<Vec<Event>>()
  913. }
  914. /// We do a non-recursive DFS (<https://en.wikipedia.org/wiki/Depth-first_search>),
  915. /// and additionally we consider the timestamps.
  916. async fn dfs_topological_sort(
  917. &self,
  918. event: Event,
  919. visited: &mut HashSet<Hash>,
  920. ) -> VecDeque<(u64, Event)> {
  921. let mut ordered_events = VecDeque::new();
  922. let mut stack = VecDeque::new();
  923. let event_id = event.id();
  924. stack.push_back(event_id);
  925. while let Some(event_id) = stack.pop_front() {
  926. if !visited.contains(&event_id) && event_id != NULL_ID {
  927. visited.insert(event_id);
  928. if let Some(event) = self.fetch_event_from_dags(&event_id).await.unwrap() {
  929. for parent in event.header.parents.iter() {
  930. stack.push_back(*parent);
  931. }
  932. ordered_events.push_back((event.header.layer, event))
  933. }
  934. }
  935. }
  936. ordered_events
  937. }
  938. /// Enable graph debugging
  939. pub async fn deg_enable(&self) {
  940. *self.deg_enabled.write().await = true;
  941. warn!("[EVENTGRAPH] Graph debugging enabled!");
  942. }
  943. /// Disable graph debugging
  944. pub async fn deg_disable(&self) {
  945. *self.deg_enabled.write().await = false;
  946. warn!("[EVENTGRAPH] Graph debugging disabled!");
  947. }
  948. /// Subscribe to deg events
  949. pub async fn deg_subscribe(&self) -> Subscription<DegEvent> {
  950. self.deg_publisher.clone().subscribe().await
  951. }
  952. /// Send a deg notification over the publisher
  953. pub async fn deg_notify(&self, event: DegEvent) {
  954. self.deg_publisher.notify(event).await;
  955. }
  956. #[cfg(feature = "rpc")]
  957. pub async fn eventgraph_info(&self, id: u16, _params: JsonValue) -> JsonResult {
  958. let current_genesis = self.current_genesis.read().await;
  959. let dag_name = current_genesis.id().to_string();
  960. let mut graph = HashMap::new();
  961. for iter_elem in self.dag_store.read().await.get_dag(&dag_name).iter() {
  962. let (id, val) = iter_elem.unwrap();
  963. let id = Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  964. let val: Event = deserialize_async(&val).await.unwrap();
  965. graph.insert(id, val);
  966. }
  967. let json_graph = graph
  968. .into_iter()
  969. .map(|(k, v)| {
  970. let key = k.to_string();
  971. let value = JsonValue::from(v);
  972. (key, value)
  973. })
  974. .collect();
  975. let values = json_map([("dag", JsonValue::Object(json_graph))]);
  976. let result = JsonValue::Object(HashMap::from([("eventgraph_info".to_string(), values)]));
  977. JsonResponse::new(result, id).into()
  978. }
  979. /// Fetch all the events that are on a higher layers than the
  980. /// provided ones.
  981. pub async fn fetch_successors_of(&self, tips: LayerUTips) -> Result<Vec<Event>> {
  982. debug!(
  983. target: "event_graph::fetch_successors_of",
  984. "fetching successors of {tips:?}"
  985. );
  986. let current_genesis = self.current_genesis.read().await;
  987. let dag_name = current_genesis.id().to_string();
  988. let mut graph = HashMap::new();
  989. for iter_elem in self.dag_store.read().await.get_dag(&dag_name).iter() {
  990. let (id, val) = iter_elem.unwrap();
  991. let hash = Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  992. let event: Event = deserialize_async(&val).await.unwrap();
  993. graph.insert(hash, event);
  994. }
  995. let mut result = vec![];
  996. 'outer: for tip in tips.iter() {
  997. for i in tip.1.iter() {
  998. if !graph.contains_key(i) {
  999. continue 'outer;
  1000. }
  1001. }
  1002. for (_, ev) in graph.iter() {
  1003. if ev.header.layer > *tip.0 && !result.contains(ev) {
  1004. result.push(ev.clone())
  1005. }
  1006. }
  1007. }
  1008. result.sort_by(|a, b| a.header.layer.cmp(&b.header.layer));
  1009. Ok(result)
  1010. }
  1011. }
  1012. async fn request_header(
  1013. peer: &Channel,
  1014. tree_name: String,
  1015. comms_timeout: u64,
  1016. ) -> Result<Vec<Header>> {
  1017. let url = peer.address();
  1018. let hdr_rep_sub = match peer.subscribe_msg::<HeaderRep>().await {
  1019. Ok(v) => v,
  1020. Err(e) => {
  1021. error!(
  1022. target: "event_graph::dag_sync()",
  1023. "[EVENTGRAPH] Sync: Couldn't subscribe HeaderReq for peer {}, skipping ({})",
  1024. url, e,
  1025. );
  1026. return Err(Error::EventNotFound("Couldn't subscribe HeaderReq".to_owned()));
  1027. }
  1028. };
  1029. if let Err(e) = peer.send(&HeaderReq(tree_name)).await {
  1030. error!(
  1031. target: "event_graph::dag_sync()",
  1032. "[EVENTGRAPH] Sync: Couldn't contact peer {}, skipping ({})", url, e,
  1033. );
  1034. return Err(Error::EventNotFound("Couldn't contact peer".to_owned()));
  1035. };
  1036. // Node waits for response
  1037. let Ok(peer_headers) = hdr_rep_sub.receive_with_timeout(comms_timeout).await else {
  1038. error!(
  1039. target: "event_graph::dag_sync()",
  1040. "[EVENTGRAPH] Sync: Peer {} didn't reply with headers in time, skipping", url,
  1041. );
  1042. // communicated_peers -= 1;
  1043. return Err(Error::EventNotFound("Peer didn't reply with headers in time".to_owned()));
  1044. };
  1045. let peer_headers = &peer_headers.0;
  1046. Ok(peer_headers.to_vec())
  1047. }
  1048. async fn request_event(
  1049. peer: Arc<Channel>,
  1050. headers: Vec<Hash>,
  1051. chunk_id: usize,
  1052. comms_timeout: u64,
  1053. ) -> (Result<Vec<Event>>, usize, Arc<Channel>) {
  1054. let url = peer.address();
  1055. debug!(
  1056. target: "event_graph::dag_sync()",
  1057. "Requesting {:?} from {}...", headers, url,
  1058. );
  1059. let ev_rep_sub = match peer.subscribe_msg::<EventRep>().await {
  1060. Ok(v) => v,
  1061. Err(e) => {
  1062. error!(
  1063. target: "event_graph::dag_sync()",
  1064. "[EVENTGRAPH] Sync: Couldn't subscribe EventRep for peer {}, skipping ({})",
  1065. url, e,
  1066. );
  1067. return (
  1068. Err(Error::EventNotFound("Couldn't subscribe EventRep".to_owned())),
  1069. chunk_id,
  1070. peer,
  1071. );
  1072. }
  1073. };
  1074. // let request_missing_events = missing_parents.clone().into_iter().collect();
  1075. if let Err(e) = peer.send(&EventReq(headers.clone())).await {
  1076. error!(
  1077. target: "event_graph::dag_sync()",
  1078. "[EVENTGRAPH] Sync: Failed communicating EventReq({:?}) to {}: {}",
  1079. headers, url, e,
  1080. );
  1081. return (
  1082. Err(Error::EventNotFound("Failed communicating EventReq".to_owned())),
  1083. chunk_id,
  1084. peer,
  1085. );
  1086. }
  1087. // Node waits for response
  1088. let Ok(event) = ev_rep_sub.receive_with_timeout(comms_timeout).await else {
  1089. error!(
  1090. target: "event_graph::dag_sync()",
  1091. "[EVENTGRAPH] Sync: Timeout waiting for parents {:?} from {}",
  1092. headers, url,
  1093. );
  1094. return (
  1095. Err(Error::EventNotFound("Timeout waiting for parents".to_owned())),
  1096. chunk_id,
  1097. peer,
  1098. );
  1099. };
  1100. (Ok(event.0.clone()), chunk_id, peer)
  1101. }