tests.rs 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957
  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 std::{
  19. collections::HashSet,
  20. slice,
  21. sync::Arc,
  22. time::{Duration, UNIX_EPOCH},
  23. };
  24. use darkfi_serial::serialize_async;
  25. use sled_overlay::sled;
  26. use smol::Executor;
  27. use crate::{
  28. error::Result,
  29. event_graph::{
  30. compute_unreferenced_tips,
  31. event::Header,
  32. filter_requested_event_rep, merge_static_sync_event_rep,
  33. proto::{cap_layer_tips, count_layer_tips, EventPut, SyncDirection, MAX_RANGE_PAGE_SIZE},
  34. test_helpers::{
  35. archive_config, bounded_dag_store_config, init_logger, make_eg, make_eg_with_config,
  36. make_network, run_multi_node_test, shutdown_network, test_config, TestIdentity,
  37. },
  38. util::next_hour_timestamp,
  39. DagStore, Event, EventGraphConfig, EventGraphPtr, LayerUTips, TimeIndex, NULL_ID,
  40. NULL_PARENTS, N_EVENT_PARENTS,
  41. },
  42. system::{sleep, timeout::timeout},
  43. };
  44. fn test_event(content: &[u8], layer: u64) -> Event {
  45. Event {
  46. header: Header {
  47. timestamp: 1_704_067_200_000 + layer,
  48. parents: NULL_PARENTS,
  49. layer,
  50. content_hash: blake3::hash(content),
  51. },
  52. content: content.to_vec(),
  53. }
  54. }
  55. #[test]
  56. fn evgr_event_rep_filter_matches_only_requested_ids() {
  57. let event_a = test_event(b"requested-a", 1);
  58. let event_b = test_event(b"requested-b", 2);
  59. let unrelated = test_event(b"unrelated", 3);
  60. let requested = vec![event_a.id(), event_b.id()];
  61. let (events, blobs, missing) =
  62. filter_requested_event_rep(&requested, vec![event_b.clone()], vec![b"blob-b".to_vec()])
  63. .unwrap();
  64. assert_eq!(events.iter().map(Event::id).collect::<Vec<_>>(), vec![event_b.id()]);
  65. assert_eq!(blobs, vec![b"blob-b".to_vec()]);
  66. assert_eq!(missing, vec![event_a.id()]);
  67. assert!(filter_requested_event_rep(
  68. &requested,
  69. vec![unrelated],
  70. vec![b"unrelated-blob".to_vec()],
  71. )
  72. .is_err());
  73. assert!(filter_requested_event_rep(
  74. &requested,
  75. vec![event_a.clone(), event_a.clone()],
  76. vec![b"blob-a".to_vec(), b"duplicate-blob".to_vec()],
  77. )
  78. .is_err());
  79. assert!(filter_requested_event_rep(&requested, vec![event_a], Vec::new()).is_err());
  80. }
  81. #[test]
  82. fn evgr_static_sync_merge_tracks_partial_requested_batches() {
  83. let event_a = test_event(b"static-requested-a", 11);
  84. let event_b = test_event(b"static-requested-b", 12);
  85. let unrelated = test_event(b"static-unrelated", 13);
  86. let requested = vec![event_a.id(), event_b.id()];
  87. let mut pending: HashSet<_> = requested.iter().copied().collect();
  88. let mut known = HashSet::new();
  89. let mut want = HashSet::new();
  90. let mut fetched = Vec::new();
  91. assert!(merge_static_sync_event_rep(
  92. &requested,
  93. &mut pending,
  94. &mut known,
  95. &mut want,
  96. &mut fetched,
  97. vec![unrelated],
  98. vec![b"unrelated-blob".to_vec()],
  99. )
  100. .is_err());
  101. assert_eq!(pending.len(), 2);
  102. assert!(fetched.is_empty());
  103. let matched = merge_static_sync_event_rep(
  104. &requested,
  105. &mut pending,
  106. &mut known,
  107. &mut want,
  108. &mut fetched,
  109. vec![event_b.clone()],
  110. vec![b"blob-b".to_vec()],
  111. )
  112. .unwrap();
  113. assert_eq!(matched, 1);
  114. assert_eq!(pending, HashSet::from([event_a.id()]));
  115. let matched = merge_static_sync_event_rep(
  116. &requested,
  117. &mut pending,
  118. &mut known,
  119. &mut want,
  120. &mut fetched,
  121. vec![event_a.clone()],
  122. vec![b"blob-a".to_vec()],
  123. )
  124. .unwrap();
  125. assert_eq!(matched, 1);
  126. assert!(pending.is_empty());
  127. let fetched_ids: HashSet<_> = fetched.iter().map(|(ev, _)| ev.id()).collect();
  128. assert_eq!(fetched_ids, HashSet::from([event_a.id(), event_b.id()]));
  129. }
  130. #[test]
  131. fn evgr_layer_tip_cap_is_bounded() {
  132. let mut tips = LayerUTips::new();
  133. tips.entry(0).or_default().insert(blake3::hash(b"tip-0"));
  134. tips.entry(1).or_default().insert(blake3::hash(b"tip-1"));
  135. tips.entry(1).or_default().insert(blake3::hash(b"tip-2"));
  136. let capped = cap_layer_tips(&tips, 2);
  137. assert_eq!(count_layer_tips(&capped), 2);
  138. assert!(capped.get(&0).is_some_and(|layer| layer.len() == 1));
  139. }
  140. #[test]
  141. fn evgr_parent_selection_does_not_wrap_saturated_layer() {
  142. let tip = blake3::hash(b"saturated-tip");
  143. let tips = LayerUTips::from([(u64::MAX, HashSet::from([tip]))]);
  144. let (layer, parents) = super::select_parents_from_tips(&tips);
  145. assert_eq!(layer, u64::MAX);
  146. assert_eq!(parents[0], tip);
  147. }
  148. #[test]
  149. fn evgr_time_index_queries_and_saturating_cursor() {
  150. // Forward, backward, newest, oldest queries plus the saturating
  151. // cursor at u64 boundaries.
  152. let mut idx = TimeIndex::new();
  153. for ts in [100_u64, 200, 200, 300, 400, 500] {
  154. let id = blake3::hash(&ts.to_be_bytes());
  155. idx.insert(ts, id);
  156. }
  157. assert_eq!(idx.len(), 6);
  158. assert_eq!(idx.newest(3).len(), 3);
  159. assert_eq!(idx.oldest(2).len(), 2);
  160. assert_eq!(idx.before(300, 10).len(), 3);
  161. assert_eq!(idx.after(200, 10).len(), 3);
  162. // Saturating cursor: before(0) shouldn't underflow,
  163. // after(u64::MAX) shouldn't overflow.
  164. let mut idx2 = TimeIndex::new();
  165. idx2.insert(100, blake3::hash(b"x"));
  166. assert_eq!(idx2.before(0, 10).len(), 0);
  167. assert_eq!(idx2.after(u64::MAX, 10).len(), 0);
  168. }
  169. async fn make_dag_store() -> Result<DagStore> {
  170. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  171. Ok(DagStore::new(sled_db, &bounded_dag_store_config()).await)
  172. }
  173. #[test]
  174. fn evgr_dag_store_eviction_policy() {
  175. // Bounded vs archive mode in one test:
  176. // (a) bounded: adding a 25th DAG drops the oldest, total stays 24.
  177. // (b) archive: adding 30 DAGs leaves all 30 plus the originals.
  178. smol::block_on(async {
  179. // (a) bounded
  180. let mut store = make_dag_store().await.unwrap();
  181. let oldest_ts = store.dag_timestamps()[0];
  182. let new_ts = next_hour_timestamp(1);
  183. let hdr = Header {
  184. timestamp: new_ts,
  185. parents: NULL_PARENTS,
  186. layer: 0,
  187. content_hash: blake3::hash(b"test-graph-v1"),
  188. };
  189. let genesis = Event { header: hdr, content: b"test-graph-v1".to_vec() };
  190. store.add_dag(&genesis, Some(24)).await;
  191. assert_eq!(store.dag_timestamps().len(), 24);
  192. assert!(store.get_slot(&new_ts).is_some());
  193. assert!(store.get_slot(&oldest_ts).is_none());
  194. // (b) archive
  195. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  196. let mut archive = DagStore::new(sled_db, &archive_config()).await;
  197. let initial = archive.dag_timestamps().len();
  198. for i in 1..=30i64 {
  199. let ts = next_hour_timestamp(i);
  200. let hdr = Header {
  201. timestamp: ts,
  202. parents: NULL_PARENTS,
  203. layer: 0,
  204. content_hash: blake3::hash(b"test-graph-v1"),
  205. };
  206. let genesis = Event { header: hdr, content: b"test-graph-v1".to_vec() };
  207. archive.add_dag(&genesis, None).await;
  208. }
  209. assert_eq!(archive.dag_timestamps().len(), initial + 30);
  210. })
  211. }
  212. #[test]
  213. fn evgr_dag_store_archive_mode_discovers_existing_trees() {
  214. smol::block_on(async {
  215. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  216. let historical_ts = next_hour_timestamp(-100);
  217. {
  218. let mut store = DagStore::new(sled_db.clone(), &archive_config()).await;
  219. let hdr = Header {
  220. timestamp: historical_ts,
  221. parents: NULL_PARENTS,
  222. layer: 0,
  223. content_hash: blake3::hash(b"test-graph-v1"),
  224. };
  225. let genesis = Event { header: hdr, content: b"test-graph-v1".to_vec() };
  226. store.add_dag(&genesis, None).await;
  227. drop(store);
  228. }
  229. let store = DagStore::new(sled_db, &archive_config()).await;
  230. assert!(
  231. store.get_slot(&historical_ts).is_some(),
  232. "Archive mode should discover historical DAGs on restart"
  233. );
  234. })
  235. }
  236. #[test]
  237. fn evgr_compute_unreferenced_tips_single_pass() {
  238. smol::block_on(async {
  239. let store = make_dag_store().await.unwrap();
  240. let ts = *store.dag_timestamps().last().unwrap();
  241. let slot = store.get_slot(&ts).unwrap();
  242. let genesis_hash = *slot.tips.get(&0).unwrap().iter().next().unwrap();
  243. let now = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  244. let mut p = [NULL_ID; N_EVENT_PARENTS];
  245. p[0] = genesis_hash;
  246. let e2 = Event {
  247. header: Header {
  248. timestamp: now,
  249. parents: p,
  250. layer: 1,
  251. content_hash: blake3::hash(b"e2"),
  252. },
  253. content: b"e2".to_vec(),
  254. };
  255. slot.main_tree.insert(e2.id().as_bytes(), serialize_async(&e2).await).unwrap();
  256. let mut p = [NULL_ID; N_EVENT_PARENTS];
  257. p[0] = e2.id();
  258. let e3 = Event {
  259. header: Header {
  260. timestamp: now,
  261. parents: p,
  262. layer: 2,
  263. content_hash: blake3::hash(b"e3"),
  264. },
  265. content: b"e3".to_vec(),
  266. };
  267. slot.main_tree.insert(e3.id().as_bytes(), serialize_async(&e3).await).unwrap();
  268. let mut p = [NULL_ID; N_EVENT_PARENTS];
  269. p[0] = genesis_hash;
  270. let e4 = Event {
  271. header: Header {
  272. timestamp: now,
  273. parents: p,
  274. layer: 1,
  275. content_hash: blake3::hash(b"e4"),
  276. },
  277. content: b"e4".to_vec(),
  278. };
  279. slot.main_tree.insert(e4.id().as_bytes(), serialize_async(&e4).await).unwrap();
  280. assert_ne!(e2.id(), e4.id(), "e2 and e4 must have distinct IDs");
  281. let tips = compute_unreferenced_tips(&slot.main_tree).await;
  282. assert!(tips.get(&2).unwrap().contains(&e3.id()));
  283. assert!(tips.get(&1).unwrap().contains(&e4.id()));
  284. assert!(!tips.values().any(|set| set.contains(&e2.id())));
  285. })
  286. }
  287. #[test]
  288. fn evgr_dag_insert_valid_and_duplicate() {
  289. // First insert: returns the id, updates tips, fires the
  290. // subscriber. Second (duplicate) insert: returns an empty list,
  291. // doesn't re-fire.
  292. smol::block_on(async {
  293. let eg = make_eg().await;
  294. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  295. let dag_name = dag_ts.to_string();
  296. let sub = eg.event_pub.clone().subscribe().await;
  297. let event = Event::new(b"hello".to_vec(), &eg).await;
  298. eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  299. let ids = eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  300. assert_eq!(ids.len(), 1);
  301. let store = eg.dag_store.read().await;
  302. let slot = store.get_slot(&dag_ts).unwrap();
  303. assert!(slot.tips.get(&1).unwrap().contains(&event.id()));
  304. drop(store);
  305. let Ok(notified) = timeout(Duration::from_secs(1), sub.receive()).await else {
  306. panic!("Event notification not received");
  307. };
  308. assert_eq!(notified.id(), event.id());
  309. // Re-insert is a no-op.
  310. assert!(eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap().is_empty());
  311. })
  312. }
  313. #[test]
  314. fn evgr_dag_insert_without_header_skipped() {
  315. smol::block_on(async {
  316. let eg = make_eg().await;
  317. let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
  318. let event = Event::new(b"orphan".to_vec(), &eg).await;
  319. let ids = eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  320. assert!(ids.is_empty());
  321. })
  322. }
  323. #[test]
  324. fn evgr_header_insert_rejects_layer_jump() {
  325. smol::block_on(async {
  326. let eg = make_eg().await;
  327. let genesis = eg.current_genesis.read().await.clone();
  328. let dag_name = genesis.header.timestamp.to_string();
  329. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  330. parents[0] = genesis.id();
  331. let event = Event {
  332. header: Header {
  333. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
  334. parents,
  335. layer: 2,
  336. content_hash: blake3::hash(b"layer-jump"),
  337. },
  338. content: b"layer-jump".to_vec(),
  339. };
  340. let err = eg.header_dag_insert(vec![event.header], &dag_name).await.unwrap_err();
  341. assert!(matches!(err, crate::Error::HeaderIsInvalid));
  342. })
  343. }
  344. #[test]
  345. fn evgr_header_insert_rejects_duplicate_parents() {
  346. smol::block_on(async {
  347. let eg = make_eg().await;
  348. let genesis = eg.current_genesis.read().await.clone();
  349. let dag_name = genesis.header.timestamp.to_string();
  350. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  351. parents[0] = genesis.id();
  352. parents[1] = genesis.id();
  353. let event = Event {
  354. header: Header {
  355. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
  356. parents,
  357. layer: 1,
  358. content_hash: blake3::hash(b"duplicate-parents"),
  359. },
  360. content: b"duplicate-parents".to_vec(),
  361. };
  362. let err = eg.header_dag_insert(vec![event.header], &dag_name).await.unwrap_err();
  363. assert!(matches!(err, crate::Error::HeaderIsInvalid));
  364. })
  365. }
  366. #[test]
  367. fn evgr_header_validate_rejects_layer_overflow_parent() {
  368. smol::block_on(async {
  369. let db = sled::Config::new().temporary(true).open().unwrap();
  370. let tree = db.open_tree("headers").unwrap();
  371. let timestamp = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  372. let parent = Header {
  373. timestamp,
  374. parents: NULL_PARENTS,
  375. layer: u64::MAX,
  376. content_hash: blake3::hash(b"overflow-parent"),
  377. };
  378. tree.insert(parent.id().as_bytes(), serialize_async(&parent).await).unwrap();
  379. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  380. parents[0] = parent.id();
  381. let child = Header {
  382. timestamp,
  383. parents,
  384. layer: u64::MAX,
  385. content_hash: blake3::hash(b"overflow-child"),
  386. };
  387. assert!(!child.validate(&tree, &test_config(), timestamp, None).await.unwrap());
  388. })
  389. }
  390. #[test]
  391. fn evgr_header_validate_uses_target_slot_bounds() {
  392. smol::block_on(async {
  393. const HOUR_MS: u64 = 3_600_000;
  394. let db = sled::Config::new().temporary(true).open().unwrap();
  395. let tree = db.open_tree("headers").unwrap();
  396. let dag_ts = 1_704_067_200_000;
  397. let drift = crate::event_graph::EVENT_TIME_DRIFT;
  398. let config = EventGraphConfig { hours_rotation: 6, ..test_config() };
  399. let genesis = Header {
  400. timestamp: dag_ts,
  401. parents: NULL_PARENTS,
  402. layer: 0,
  403. content_hash: blake3::hash(&config.genesis_contents),
  404. };
  405. tree.insert(genesis.id().as_bytes(), serialize_async(&genesis).await).unwrap();
  406. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  407. parents[0] = genesis.id();
  408. let make_header = |timestamp, content: &[u8]| Header {
  409. timestamp,
  410. parents,
  411. layer: 1,
  412. content_hash: blake3::hash(content),
  413. };
  414. let lower_edge = make_header(dag_ts.saturating_sub(drift), b"lower-edge");
  415. assert!(lower_edge.validate(&tree, &config, dag_ts, None).await.unwrap());
  416. let upper_edge = make_header(dag_ts + 6 * HOUR_MS + drift - 1, b"upper-edge");
  417. assert!(upper_edge.validate(&tree, &config, dag_ts, None).await.unwrap());
  418. let too_early = make_header(dag_ts - drift - 1, b"too-early");
  419. assert!(!too_early.validate(&tree, &config, dag_ts, None).await.unwrap());
  420. let too_late = make_header(dag_ts + 6 * HOUR_MS + drift, b"too-late");
  421. assert!(!too_late.validate(&tree, &config, dag_ts, None).await.unwrap());
  422. })
  423. }
  424. #[test]
  425. fn evgr_header_validate_ignores_future_initial_genesis() {
  426. smol::block_on(async {
  427. const HOUR_MS: u64 = 3_600_000;
  428. let db = sled::Config::new().temporary(true).open().unwrap();
  429. let tree = db.open_tree("headers").unwrap();
  430. let now = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  431. let dag_ts = now.saturating_sub(HOUR_MS);
  432. let config =
  433. EventGraphConfig { initial_genesis: now + HOUR_MS, hours_rotation: 1, ..test_config() };
  434. let genesis = Header {
  435. timestamp: dag_ts,
  436. parents: NULL_PARENTS,
  437. layer: 0,
  438. content_hash: blake3::hash(&config.genesis_contents),
  439. };
  440. tree.insert(genesis.id().as_bytes(), serialize_async(&genesis).await).unwrap();
  441. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  442. parents[0] = genesis.id();
  443. let child = Header {
  444. timestamp: dag_ts + 1,
  445. parents,
  446. layer: 1,
  447. content_hash: blake3::hash(b"future-initial-genesis"),
  448. };
  449. assert!(child.validate(&tree, &config, dag_ts, None).await.unwrap());
  450. })
  451. }
  452. #[test]
  453. fn evgr_header_validate_no_rotation_rejects_far_future() {
  454. smol::block_on(async {
  455. let db = sled::Config::new().temporary(true).open().unwrap();
  456. let tree = db.open_tree("headers").unwrap();
  457. let config = test_config();
  458. let dag_ts = config.initial_genesis;
  459. let genesis = Header {
  460. timestamp: dag_ts,
  461. parents: NULL_PARENTS,
  462. layer: 0,
  463. content_hash: blake3::hash(&config.genesis_contents),
  464. };
  465. tree.insert(genesis.id().as_bytes(), serialize_async(&genesis).await).unwrap();
  466. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  467. parents[0] = genesis.id();
  468. let old_history = Header {
  469. timestamp: dag_ts + 1,
  470. parents,
  471. layer: 1,
  472. content_hash: blake3::hash(b"old-no-rotation-history"),
  473. };
  474. assert!(old_history.validate(&tree, &config, dag_ts, None).await.unwrap());
  475. let future = Header {
  476. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64 +
  477. crate::event_graph::EVENT_TIME_DRIFT +
  478. 1,
  479. parents,
  480. layer: 1,
  481. content_hash: blake3::hash(b"future-no-rotation-header"),
  482. };
  483. assert!(!future.validate(&tree, &config, dag_ts, None).await.unwrap());
  484. })
  485. }
  486. #[test]
  487. fn evgr_header_insert_rejects_unloaded_dag_slot() {
  488. smol::block_on(async {
  489. let config = EventGraphConfig { hours_rotation: 1, max_dags: Some(2), ..test_config() };
  490. let eg = make_eg_with_config(config).await;
  491. let dag_ts = next_hour_timestamp(-100);
  492. let dag_name = dag_ts.to_string();
  493. let genesis = Header {
  494. timestamp: dag_ts,
  495. parents: NULL_PARENTS,
  496. layer: 0,
  497. content_hash: blake3::hash(&eg.config.genesis_contents),
  498. };
  499. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  500. parents[0] = genesis.id();
  501. let header = Header {
  502. timestamp: dag_ts + 1,
  503. parents,
  504. layer: 1,
  505. content_hash: blake3::hash(b"unloaded-slot"),
  506. };
  507. let err = eg.header_dag_insert(vec![header], &dag_name).await.unwrap_err();
  508. assert!(matches!(err, crate::Error::DagSyncFailed));
  509. })
  510. }
  511. #[test]
  512. fn evgr_fetch_page_both_directions() {
  513. smol::block_on(async {
  514. let eg = make_eg().await;
  515. let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
  516. let base = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  517. for i in 0..(MAX_RANGE_PAGE_SIZE as u64 + 10) {
  518. let ev = Event::with_timestamp(base + i, vec![(i % 251) as u8], &eg).await;
  519. eg.header_dag_insert(vec![ev.header.clone()], &dag_name).await.unwrap();
  520. eg.dag_insert(slice::from_ref(&ev), &dag_name).await.unwrap();
  521. }
  522. let page = eg.fetch_page(u64::MAX, SyncDirection::Backward, 5).await.unwrap();
  523. assert_eq!(page.len(), 5);
  524. for w in page.windows(2) {
  525. assert!(w[0].header.timestamp >= w[1].header.timestamp);
  526. }
  527. let page = eg.fetch_page(0, SyncDirection::Forward, 5).await.unwrap();
  528. assert!(!page.is_empty());
  529. for w in page.windows(2) {
  530. assert!(w[0].header.timestamp <= w[1].header.timestamp);
  531. }
  532. let capped = eg
  533. .fetch_page(u64::MAX, SyncDirection::Backward, MAX_RANGE_PAGE_SIZE + 10)
  534. .await
  535. .unwrap();
  536. assert_eq!(capped.len(), MAX_RANGE_PAGE_SIZE);
  537. })
  538. }
  539. async fn build_graph() -> Result<(EventGraphPtr, std::collections::HashMap<&'static str, Event>)> {
  540. let eg = make_eg().await;
  541. let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
  542. let genesis_hash = eg.current_genesis.read().await.id();
  543. let base = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  544. let make = |off: u64, layer: u64, parents: [blake3::Hash; N_EVENT_PARENTS], name: &str| Event {
  545. header: Header {
  546. timestamp: base + off,
  547. layer,
  548. parents,
  549. content_hash: blake3::hash(name.as_bytes()),
  550. },
  551. content: name.as_bytes().to_vec(),
  552. };
  553. let mut p = [NULL_ID; N_EVENT_PARENTS];
  554. p[0] = genesis_hash;
  555. let e1a = make(1, 1, p, "e1a");
  556. let e1b = make(2, 1, p, "e1b");
  557. let e1c = make(3, 1, p, "e1c");
  558. let e1d = make(4, 1, p, "e1d");
  559. let mut p = [NULL_ID; N_EVENT_PARENTS];
  560. p[0] = e1a.id();
  561. let e2a = make(5, 2, p, "e2a");
  562. p[0] = e1b.id();
  563. let e2b = make(6, 2, p, "e2b");
  564. p[0] = e1c.id();
  565. let e2c = make(7, 2, p, "e2c");
  566. p[0] = e1d.id();
  567. let e2d = make(8, 2, p, "e2d");
  568. let l1 = vec![e1a.clone(), e1b.clone(), e1c.clone(), e1d.clone()];
  569. let l2 = vec![e2a.clone(), e2b.clone(), e2c.clone(), e2d.clone()];
  570. eg.header_dag_insert(l1.iter().map(|e| e.header.clone()).collect(), &dag_name).await.unwrap();
  571. eg.dag_insert(&l1, &dag_name).await.unwrap();
  572. eg.header_dag_insert(l2.iter().map(|e| e.header.clone()).collect(), &dag_name).await.unwrap();
  573. eg.dag_insert(&l2, &dag_name).await.unwrap();
  574. let mut map = std::collections::HashMap::new();
  575. map.insert("e1a", e1a);
  576. map.insert("e1b", e1b);
  577. map.insert("e1c", e1c);
  578. map.insert("e1d", e1d);
  579. map.insert("e2a", e2a);
  580. map.insert("e2b", e2b);
  581. map.insert("e2c", e2c);
  582. map.insert("e2d", e2d);
  583. Ok((eg, map))
  584. }
  585. #[test]
  586. fn evgr_ancestor_walk_via_header_tree() {
  587. smol::block_on(async {
  588. let (eg, evs) = build_graph().await.unwrap();
  589. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  590. let store = eg.dag_store.read().await;
  591. let slot = store.get_slot(&dag_ts).unwrap();
  592. let genesis_hash = eg.current_genesis.read().await.id();
  593. for name in ["e1a", "e1b", "e1c", "e1d"] {
  594. let mut ancestors = HashSet::new();
  595. eg.get_ancestors(&mut ancestors, evs[name].header.clone(), &slot.header_tree)
  596. .await
  597. .unwrap();
  598. assert_eq!(ancestors, HashSet::from([genesis_hash]));
  599. }
  600. let mut ancestors = HashSet::new();
  601. eg.get_ancestors(&mut ancestors, evs["e2a"].header.clone(), &slot.header_tree)
  602. .await
  603. .unwrap();
  604. assert_eq!(ancestors, HashSet::from([genesis_hash, evs["e1a"].id()]));
  605. })
  606. }
  607. #[test]
  608. fn evgr_multi_node_propagation_with_real_blob() {
  609. init_logger();
  610. run_multi_node_test(propagation_with_real_blob);
  611. }
  612. async fn propagation_with_real_blob(ex: Arc<Executor<'static>>) {
  613. let nodes = make_network(ex).await;
  614. let mut alice = TestIdentity::new();
  615. for eg in &nodes {
  616. alice.register_directly(eg).await.expect("register alice");
  617. }
  618. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  619. let dag_name = dag_ts.to_string();
  620. let event = Event::new(b"hello-via-rln".to_vec(), &nodes[0]).await;
  621. let message_id =
  622. alice.next_message_id(event.header.timestamp).expect("budget available on first signal");
  623. let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
  624. let blob = serialize_async(&blob_struct).await;
  625. nodes[0].header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  626. nodes[0].dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  627. nodes[0].dag_blob_store(&event.id(), &blob).unwrap();
  628. nodes[0].p2p.broadcast(&EventPut(event.clone(), blob.clone())).await;
  629. sleep(5).await;
  630. for (i, eg) in nodes.iter().enumerate() {
  631. let store = eg.dag_store.read().await;
  632. let slot = store.get_slot(&dag_ts).unwrap();
  633. assert!(
  634. slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  635. "node {i} missing event in main_tree",
  636. );
  637. drop(store);
  638. assert!(
  639. eg.dag_blob_fetch(&event.id()).unwrap().is_some(),
  640. "node {i} missing blob in dag_blobs - late-joiners would have nothing to verify against",
  641. );
  642. }
  643. shutdown_network(&nodes).await;
  644. }
  645. #[test]
  646. fn evgr_multi_node_empty_blob_rejected() {
  647. init_logger();
  648. run_multi_node_test(empty_blob_rejected);
  649. }
  650. async fn empty_blob_rejected(ex: Arc<Executor<'static>>) {
  651. // An attacker broadcasts `EventPut(ev, vec![])` for a non-genesis event.
  652. // Every recipient must strike the sender and refuse to insert.
  653. let nodes = make_network(ex).await;
  654. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  655. let event = Event::new(b"unauthenticated".to_vec(), &nodes[0]).await;
  656. nodes[0].p2p.broadcast(&EventPut(event.clone(), vec![])).await;
  657. sleep(5).await;
  658. for (i, eg) in nodes.iter().enumerate().skip(1) {
  659. let store = eg.dag_store.read().await;
  660. let slot = store.get_slot(&dag_ts).unwrap();
  661. assert!(
  662. !slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  663. "Vector-1 breach: node {i} accepted an empty-blob non-genesis event",
  664. );
  665. }
  666. shutdown_network(&nodes).await;
  667. }
  668. #[test]
  669. fn evgr_multi_node_genesis_with_blob_rejected() {
  670. init_logger();
  671. run_multi_node_test(genesis_with_blob_rejected);
  672. }
  673. async fn genesis_with_blob_rejected(ex: Arc<Executor<'static>>) {
  674. // Symmetric defense: a genesis-shaped event arriving with a
  675. // non-empty blob is also misbehavior. Genesis events are
  676. // deterministic and don't carry signals.
  677. let nodes = make_network(ex).await;
  678. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  679. let header = Header {
  680. timestamp: dag_ts,
  681. parents: NULL_PARENTS,
  682. layer: 0,
  683. content_hash: blake3::hash(b"forged-genesis"),
  684. };
  685. let event = Event { header, content: b"forged-genesis".to_vec() };
  686. let fake_blob = b"this-should-not-be-here".to_vec();
  687. nodes[0].p2p.broadcast(&EventPut(event.clone(), fake_blob)).await;
  688. sleep(5).await;
  689. for (i, eg) in nodes.iter().enumerate().skip(1) {
  690. let store = eg.dag_store.read().await;
  691. let slot = store.get_slot(&dag_ts).unwrap();
  692. assert!(
  693. !slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  694. "node {i} accepted a genesis-shaped event with a blob",
  695. );
  696. }
  697. shutdown_network(&nodes).await;
  698. }
  699. #[test]
  700. fn evgr_multi_node_dag_sync_with_blob() {
  701. init_logger();
  702. run_multi_node_test(dag_sync_with_blob);
  703. }
  704. async fn dag_sync_with_blob(ex: Arc<Executor<'static>>) {
  705. // End-to-end Vector-2 propagation test. Nodes 0..3 receive a
  706. // signal via direct insert (with the blob in their dag_blobs
  707. // side-table). Node 4 catches up via dag_sync - it should
  708. // receive both the event and the blob from a peer, and re-verify
  709. // the proof at sync time.
  710. let nodes = make_network(ex).await;
  711. let mut alice = TestIdentity::new();
  712. for eg in &nodes {
  713. alice.register_directly(eg).await.unwrap();
  714. }
  715. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  716. let dag_name = dag_ts.to_string();
  717. let event = Event::new(b"synced-message".to_vec(), &nodes[0]).await;
  718. let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
  719. let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
  720. let blob = serialize_async(&blob_struct).await;
  721. for eg in nodes.iter().take(4) {
  722. eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  723. eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  724. eg.dag_blob_store(&event.id(), &blob).unwrap();
  725. }
  726. {
  727. let store = nodes[4].dag_store.read().await;
  728. let slot = store.get_slot(&dag_ts).unwrap();
  729. assert!(!slot.main_tree.contains_key(event.id().as_bytes()).unwrap());
  730. }
  731. nodes[4].dag_sync(dag_ts).await.expect("dag_sync should succeed");
  732. sleep(2).await;
  733. let store = nodes[4].dag_store.read().await;
  734. let slot = store.get_slot(&dag_ts).unwrap();
  735. assert!(
  736. slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  737. "sync should bring the event over",
  738. );
  739. drop(store);
  740. assert!(
  741. nodes[4].dag_blob_fetch(&event.id()).unwrap().is_some(),
  742. "sync should bring the blob over too - without it, node 4 can't serve future late-joiners",
  743. );
  744. shutdown_network(&nodes).await;
  745. }
  746. #[test]
  747. fn evgr_multi_node_dag_sync_rejects_bad_blob() {
  748. init_logger();
  749. run_multi_node_test(dag_sync_rejects_bad_blob);
  750. }
  751. async fn dag_sync_rejects_bad_blob(ex: Arc<Executor<'static>>) {
  752. // Vector-2 defense: a peer in the 2/3 quorum serves a tampered
  753. // blob during sync. The recipient's dag_insert_with_blobs runs
  754. // RLN re-verification, which rejects, and the event does NOT
  755. // end up in the recipient's main_tree.
  756. let nodes = make_network(ex).await;
  757. let alice = TestIdentity::new();
  758. for eg in &nodes {
  759. alice.register_directly(eg).await.unwrap();
  760. }
  761. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  762. let dag_name = dag_ts.to_string();
  763. let event = Event::new(b"crafted-injection".to_vec(), &nodes[0]).await;
  764. // Garbage bytes - won't deserialize as a real RLN signal, won't
  765. // verify. We don't need a real (failing) proof to exercise the
  766. // rejection path.
  767. let bad_blob = b"definitely-not-a-real-rln-blob".to_vec();
  768. for eg in nodes.iter().take(4) {
  769. eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  770. eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  771. eg.dag_blob_store(&event.id(), &bad_blob).unwrap();
  772. }
  773. nodes[4].dag_sync(dag_ts).await.unwrap();
  774. sleep(2).await;
  775. let store = nodes[4].dag_store.read().await;
  776. let slot = store.get_slot(&dag_ts).unwrap();
  777. assert!(
  778. !slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  779. "Vector-2 breach: node 4 accepted an event with a tampered blob during sync",
  780. );
  781. shutdown_network(&nodes).await;
  782. }
  783. #[test]
  784. fn evgr_multi_node_dormant_user_can_post_after_long_silence() {
  785. init_logger();
  786. run_multi_node_test(dormant_user_can_post_after_long_silence);
  787. }
  788. async fn dormant_user_can_post_after_long_silence(ex: Arc<Executor<'static>>) {
  789. // Alice registers, then 17+ other identities also register. By
  790. // the time Alice tries to send a signal, her registration root
  791. // has long fallen out of the in-memory recent_roots window. The
  792. // historical-roots side-table makes verification succeed anyway.
  793. //
  794. // Note: in this test Alice's signal still references the CURRENT
  795. // root (because `rln_membership_path` returns the live root and
  796. // Alice is still a member). That's fine - what we're validating
  797. // here is end-to-end: a deeply-historical state of the SMT
  798. // doesn't break verification. The unit tests in tests_rln.rs
  799. // (rln_is_root_valid_at_*) cover the predicate's exact semantics
  800. // for old-root references.
  801. let nodes = make_network(ex).await;
  802. let mut alice = TestIdentity::new();
  803. for eg in &nodes {
  804. alice.register_directly(eg).await.unwrap();
  805. }
  806. // Push the recent_roots window past Alice's registration.
  807. for seed in 100..117_u64 {
  808. let other = TestIdentity::with_seed(seed);
  809. for eg in &nodes {
  810. other.register_directly(eg).await.unwrap();
  811. }
  812. }
  813. let event = Event::new(b"long-silent-but-still-registered".to_vec(), &nodes[0]).await;
  814. let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
  815. let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
  816. let blob = serialize_async(&blob_struct).await;
  817. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  818. let dag_name = dag_ts.to_string();
  819. nodes[0].header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  820. nodes[0].dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  821. nodes[0].dag_blob_store(&event.id(), &blob).unwrap();
  822. nodes[0].p2p.broadcast(&EventPut(event.clone(), blob)).await;
  823. sleep(5).await;
  824. for (i, eg) in nodes.iter().enumerate() {
  825. let store = eg.dag_store.read().await;
  826. let slot = store.get_slot(&dag_ts).unwrap();
  827. assert!(
  828. slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  829. "node {i} rejected Alice's signal even though she's a valid registered identity \
  830. - historical-roots fallback is broken",
  831. );
  832. }
  833. shutdown_network(&nodes).await;
  834. }