tests.rs 46 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290
  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::{
  34. cap_layer_tips, count_layer_tips, filter_parent_event_rep, EventPut, SyncDirection,
  35. MAX_HEADER_REP_HEADERS, MAX_RANGE_PAGE_SIZE,
  36. },
  37. rln::epoch_of,
  38. test_helpers::{
  39. archive_config, bounded_dag_store_config, init_logger, make_eg, make_eg_with_config,
  40. make_network, run_multi_node_test, shutdown_network, test_config, TestIdentity,
  41. },
  42. util::{millis_until_next_rotation, next_hour_timestamp, next_rotation_timestamp},
  43. DagStore, Event, EventGraphConfig, EventGraphPtr, LayerUTips, TimeIndex, NULL_ID,
  44. NULL_PARENTS, N_EVENT_PARENTS,
  45. },
  46. system::{sleep, timeout::timeout},
  47. };
  48. fn test_event(content: &[u8], layer: u64) -> Event {
  49. Event {
  50. header: Header {
  51. timestamp: 1_704_067_200_000 + layer,
  52. parents: NULL_PARENTS,
  53. layer,
  54. content_hash: blake3::hash(content),
  55. },
  56. content: content.to_vec(),
  57. }
  58. }
  59. #[test]
  60. fn evgr_event_rep_filter_matches_only_requested_ids() {
  61. let event_a = test_event(b"requested-a", 1);
  62. let event_b = test_event(b"requested-b", 2);
  63. let unrelated = test_event(b"unrelated", 3);
  64. let requested = vec![event_a.id(), event_b.id()];
  65. let (events, blobs, missing) =
  66. filter_requested_event_rep(&requested, vec![event_b.clone()], vec![b"blob-b".to_vec()])
  67. .unwrap();
  68. assert_eq!(events.iter().map(Event::id).collect::<Vec<_>>(), vec![event_b.id()]);
  69. assert_eq!(blobs, vec![b"blob-b".to_vec()]);
  70. assert_eq!(missing, vec![event_a.id()]);
  71. assert!(filter_requested_event_rep(
  72. &requested,
  73. vec![unrelated],
  74. vec![b"unrelated-blob".to_vec()],
  75. )
  76. .is_err());
  77. assert!(filter_requested_event_rep(
  78. &requested,
  79. vec![event_a.clone(), event_a.clone()],
  80. vec![b"blob-a".to_vec(), b"duplicate-blob".to_vec()],
  81. )
  82. .is_err());
  83. assert!(filter_requested_event_rep(&requested, vec![event_a], Vec::new()).is_err());
  84. }
  85. #[test]
  86. fn evgr_parent_event_rep_requires_progress() {
  87. let parent_a = test_event(b"parent-a", 21);
  88. let parent_b = test_event(b"parent-b", 22);
  89. let unrelated = test_event(b"unrelated-parent", 23);
  90. let requested = vec![parent_a.id(), parent_b.id()];
  91. let mut missing: HashSet<_> = requested.iter().copied().collect();
  92. let mut known = HashSet::new();
  93. assert!(filter_parent_event_rep(&requested, &mut missing, &mut known, vec![], vec![]).is_err());
  94. assert_eq!(missing.len(), 2);
  95. assert!(known.is_empty());
  96. assert!(filter_parent_event_rep(
  97. &requested,
  98. &mut missing,
  99. &mut known,
  100. vec![unrelated],
  101. vec![b"blob".to_vec()],
  102. )
  103. .is_err());
  104. assert_eq!(missing.len(), 2);
  105. assert!(known.is_empty());
  106. assert!(filter_parent_event_rep(
  107. &requested,
  108. &mut missing,
  109. &mut known,
  110. vec![parent_a.clone(), parent_a.clone()],
  111. vec![b"blob-a".to_vec(), b"blob-a-duplicate".to_vec()],
  112. )
  113. .is_err());
  114. assert_eq!(missing.len(), 2);
  115. assert!(known.is_empty());
  116. let resolved = filter_parent_event_rep(
  117. &requested,
  118. &mut missing,
  119. &mut known,
  120. vec![parent_a.clone()],
  121. vec![b"blob-a".to_vec()],
  122. )
  123. .unwrap();
  124. assert_eq!(
  125. resolved.iter().map(|(event, _)| event.id()).collect::<Vec<_>>(),
  126. vec![parent_a.id()]
  127. );
  128. assert!(!missing.contains(&parent_a.id()));
  129. assert!(missing.contains(&parent_b.id()));
  130. assert!(known.contains(&parent_a.id()));
  131. let current_request = vec![parent_b.id()];
  132. assert!(filter_parent_event_rep(
  133. &current_request,
  134. &mut missing,
  135. &mut known,
  136. vec![parent_a.clone()],
  137. vec![b"stale-blob-a".to_vec()],
  138. )
  139. .is_err());
  140. assert!(missing.contains(&parent_b.id()));
  141. let mut empty_missing = HashSet::new();
  142. let mut empty_known = HashSet::new();
  143. assert!(filter_parent_event_rep(
  144. &[parent_a.id()],
  145. &mut empty_missing,
  146. &mut empty_known,
  147. vec![parent_a],
  148. vec![b"blob-a".to_vec()],
  149. )
  150. .is_err());
  151. }
  152. #[test]
  153. fn evgr_static_sync_merge_tracks_partial_requested_batches() {
  154. let event_a = test_event(b"static-requested-a", 11);
  155. let event_b = test_event(b"static-requested-b", 12);
  156. let unrelated = test_event(b"static-unrelated", 13);
  157. let requested = vec![event_a.id(), event_b.id()];
  158. let mut pending: HashSet<_> = requested.iter().copied().collect();
  159. let mut known = HashSet::new();
  160. let mut want = HashSet::new();
  161. let mut fetched = Vec::new();
  162. assert!(merge_static_sync_event_rep(
  163. &requested,
  164. &mut pending,
  165. &mut known,
  166. &mut want,
  167. &mut fetched,
  168. vec![unrelated],
  169. vec![b"unrelated-blob".to_vec()],
  170. )
  171. .is_err());
  172. assert_eq!(pending.len(), 2);
  173. assert!(fetched.is_empty());
  174. let matched = merge_static_sync_event_rep(
  175. &requested,
  176. &mut pending,
  177. &mut known,
  178. &mut want,
  179. &mut fetched,
  180. vec![event_b.clone()],
  181. vec![b"blob-b".to_vec()],
  182. )
  183. .unwrap();
  184. assert_eq!(matched, 1);
  185. assert_eq!(pending, HashSet::from([event_a.id()]));
  186. let matched = merge_static_sync_event_rep(
  187. &requested,
  188. &mut pending,
  189. &mut known,
  190. &mut want,
  191. &mut fetched,
  192. vec![event_a.clone()],
  193. vec![b"blob-a".to_vec()],
  194. )
  195. .unwrap();
  196. assert_eq!(matched, 1);
  197. assert!(pending.is_empty());
  198. let fetched_ids: HashSet<_> = fetched.iter().map(|(ev, _)| ev.id()).collect();
  199. assert_eq!(fetched_ids, HashSet::from([event_a.id(), event_b.id()]));
  200. }
  201. #[test]
  202. fn evgr_layer_tip_cap_is_bounded() {
  203. let mut tips = LayerUTips::new();
  204. tips.entry(0).or_default().insert(blake3::hash(b"tip-0"));
  205. tips.entry(1).or_default().insert(blake3::hash(b"tip-1"));
  206. tips.entry(1).or_default().insert(blake3::hash(b"tip-2"));
  207. let capped = cap_layer_tips(&tips, 2);
  208. assert_eq!(count_layer_tips(&capped), 2);
  209. assert!(capped.get(&0).is_some_and(|layer| layer.len() == 1));
  210. }
  211. #[test]
  212. fn evgr_parent_selection_does_not_wrap_saturated_layer() {
  213. let tip = blake3::hash(b"saturated-tip");
  214. let tips = LayerUTips::from([(u64::MAX, HashSet::from([tip]))]);
  215. let (layer, parents) = super::select_parents_from_tips(&tips);
  216. assert_eq!(layer, u64::MAX);
  217. assert_eq!(parents[0], tip);
  218. }
  219. #[test]
  220. fn evgr_config_rejects_invalid_rotation_settings() {
  221. let zero_dags = EventGraphConfig { max_dags: Some(0), ..test_config() };
  222. assert!(zero_dags.validate().is_err());
  223. let overflowing_rotation = EventGraphConfig { hours_rotation: u64::MAX, ..test_config() };
  224. assert!(overflowing_rotation.validate().is_err());
  225. let overflowing_genesis =
  226. EventGraphConfig { initial_genesis: u64::MAX, hours_rotation: 1, ..test_config() };
  227. assert!(overflowing_genesis.validate().is_err());
  228. }
  229. #[test]
  230. fn evgr_invalid_config_does_not_open_dag_trees() {
  231. smol::block_on(async {
  232. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  233. let before = sled_db.tree_names();
  234. let config = EventGraphConfig { hours_rotation: 1, max_dags: Some(0), ..test_config() };
  235. let result = DagStore::new(sled_db.clone(), &config).await;
  236. assert!(matches!(result, Err(crate::Error::Custom(_))));
  237. assert_eq!(sled_db.tree_names(), before);
  238. })
  239. }
  240. #[test]
  241. fn evgr_rotation_helpers_are_total() {
  242. const HOUR_MS: u64 = 3_600_000;
  243. assert!(next_rotation_timestamp(0, 0).is_err());
  244. let now = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  245. let future_start = now + HOUR_MS;
  246. assert_eq!(next_rotation_timestamp(future_start, 1).unwrap(), future_start);
  247. assert!(millis_until_next_rotation(now.saturating_sub(1)).is_err());
  248. assert_eq!(super::util::hours_since(future_start), 0);
  249. }
  250. #[test]
  251. fn evgr_time_index_queries_and_saturating_cursor() {
  252. // Forward, backward, newest, oldest queries plus the saturating
  253. // cursor at u64 boundaries.
  254. let mut idx = TimeIndex::new();
  255. for (i, ts) in [100_u64, 200, 200, 300, 400, 500].into_iter().enumerate() {
  256. let id = blake3::hash(&[ts.to_be_bytes(), (i as u64).to_be_bytes()].concat());
  257. idx.insert(ts, id);
  258. }
  259. assert_eq!(idx.len(), 6);
  260. assert_eq!(idx.newest(3).len(), 3);
  261. assert_eq!(idx.oldest(2).len(), 2);
  262. assert_eq!(idx.before(300, 10).len(), 3);
  263. assert_eq!(idx.after(200, 10).len(), 3);
  264. let duplicate = blake3::hash(b"duplicate-time-index-entry");
  265. idx.insert(600, duplicate);
  266. idx.insert(600, duplicate);
  267. assert_eq!(idx.len(), 7);
  268. assert_eq!(idx.newest(10).iter().filter(|id| **id == duplicate).count(), 1);
  269. // Saturating cursor: before(0) shouldn't underflow,
  270. // after(u64::MAX) shouldn't overflow.
  271. let mut idx2 = TimeIndex::new();
  272. idx2.insert(100, blake3::hash(b"x"));
  273. assert_eq!(idx2.before(0, 10).len(), 0);
  274. assert_eq!(idx2.after(u64::MAX, 10).len(), 0);
  275. }
  276. async fn make_dag_store() -> Result<DagStore> {
  277. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  278. DagStore::new(sled_db, &bounded_dag_store_config()).await
  279. }
  280. #[test]
  281. fn evgr_dag_store_eviction_policy() {
  282. // Bounded vs archive mode in one test:
  283. // (a) bounded: adding a 25th DAG drops the oldest, total stays 24.
  284. // (b) archive: adding 30 DAGs leaves all 30 plus the originals.
  285. smol::block_on(async {
  286. // (a) bounded
  287. let mut store = make_dag_store().await.unwrap();
  288. let oldest_ts = store.dag_timestamps()[0];
  289. let new_ts = next_hour_timestamp(1);
  290. let hdr = Header {
  291. timestamp: new_ts,
  292. parents: NULL_PARENTS,
  293. layer: 0,
  294. content_hash: blake3::hash(b"test-graph-v1"),
  295. };
  296. let genesis = Event { header: hdr, content: b"test-graph-v1".to_vec() };
  297. store.add_dag(&genesis, Some(24)).await.unwrap();
  298. assert_eq!(store.dag_timestamps().len(), 24);
  299. assert!(store.get_slot(&new_ts).is_some());
  300. assert!(store.get_slot(&oldest_ts).is_none());
  301. // (b) archive
  302. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  303. let mut archive = DagStore::new(sled_db, &archive_config()).await.unwrap();
  304. let initial = archive.dag_timestamps().len();
  305. for i in 1..=30i64 {
  306. let ts = next_hour_timestamp(i);
  307. let hdr = Header {
  308. timestamp: ts,
  309. parents: NULL_PARENTS,
  310. layer: 0,
  311. content_hash: blake3::hash(b"test-graph-v1"),
  312. };
  313. let genesis = Event { header: hdr, content: b"test-graph-v1".to_vec() };
  314. archive.add_dag(&genesis, None).await.unwrap();
  315. }
  316. assert_eq!(archive.dag_timestamps().len(), initial + 30);
  317. })
  318. }
  319. #[test]
  320. fn evgr_dag_store_archive_mode_discovers_existing_trees() {
  321. smol::block_on(async {
  322. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  323. let historical_ts = next_hour_timestamp(-100);
  324. {
  325. let mut store = DagStore::new(sled_db.clone(), &archive_config()).await.unwrap();
  326. let hdr = Header {
  327. timestamp: historical_ts,
  328. parents: NULL_PARENTS,
  329. layer: 0,
  330. content_hash: blake3::hash(b"test-graph-v1"),
  331. };
  332. let genesis = Event { header: hdr, content: b"test-graph-v1".to_vec() };
  333. store.add_dag(&genesis, None).await.unwrap();
  334. drop(store);
  335. }
  336. let store = DagStore::new(sled_db, &archive_config()).await.unwrap();
  337. assert!(
  338. store.get_slot(&historical_ts).is_some(),
  339. "Archive mode should discover historical DAGs on restart"
  340. );
  341. })
  342. }
  343. #[test]
  344. fn evgr_dag_store_rejects_corrupt_header_index_on_open() {
  345. smol::block_on(async {
  346. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  347. let config = bounded_dag_store_config();
  348. let store = DagStore::new(sled_db.clone(), &config).await.unwrap();
  349. let ts = *store.dag_timestamps().last().unwrap();
  350. drop(store);
  351. let bad_id = [0u8; 32];
  352. let headers = sled_db.open_tree(format!("headers_{ts}")).unwrap();
  353. headers.insert(bad_id.as_slice(), b"not-a-header".as_slice()).unwrap();
  354. let result = DagStore::new(sled_db, &config).await;
  355. assert!(result.is_err(), "corrupt header bytes should fail DAG startup");
  356. })
  357. }
  358. #[test]
  359. fn evgr_dag_store_rejects_corrupt_event_tree_on_open() {
  360. smol::block_on(async {
  361. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  362. let config = bounded_dag_store_config();
  363. let store = DagStore::new(sled_db.clone(), &config).await.unwrap();
  364. let ts = *store.dag_timestamps().last().unwrap();
  365. drop(store);
  366. let bad_id = [0u8; 32];
  367. let events = sled_db.open_tree(ts.to_string()).unwrap();
  368. events.insert(bad_id.as_slice(), b"not-an-event".as_slice()).unwrap();
  369. let result = DagStore::new(sled_db, &config).await;
  370. assert!(result.is_err(), "corrupt event bytes should fail DAG startup");
  371. })
  372. }
  373. #[test]
  374. fn evgr_rotating_event_creation_rejects_missing_current_dag_slot() {
  375. smol::block_on(async {
  376. let eg = make_eg().await;
  377. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  378. eg.dag_store.write().await.dags.remove(&dag_ts);
  379. let result = Event::new(b"missing-current-slot".to_vec(), &eg).await;
  380. assert!(matches!(result, Err(crate::Error::Custom(_))));
  381. })
  382. }
  383. #[test]
  384. fn evgr_static_event_creation_rejects_corrupt_static_dag() {
  385. smol::block_on(async {
  386. let eg = make_eg().await;
  387. let bad_id = [0u8; 32];
  388. eg.static_dag.insert(bad_id.as_slice(), b"not-an-event".as_slice()).unwrap();
  389. let result = Event::new_static(b"static-after-corruption".to_vec(), &eg).await;
  390. assert!(result.is_err(), "corrupt static DAG should fail static event creation");
  391. })
  392. }
  393. #[test]
  394. fn evgr_compute_unreferenced_tips_single_pass() {
  395. smol::block_on(async {
  396. let store = make_dag_store().await.unwrap();
  397. let ts = *store.dag_timestamps().last().unwrap();
  398. let slot = store.get_slot(&ts).unwrap();
  399. let genesis_hash = *slot.tips.get(&0).unwrap().iter().next().unwrap();
  400. let now = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  401. let mut p = [NULL_ID; N_EVENT_PARENTS];
  402. p[0] = genesis_hash;
  403. let e2 = Event {
  404. header: Header {
  405. timestamp: now,
  406. parents: p,
  407. layer: 1,
  408. content_hash: blake3::hash(b"e2"),
  409. },
  410. content: b"e2".to_vec(),
  411. };
  412. slot.main_tree.insert(e2.id().as_bytes(), serialize_async(&e2).await).unwrap();
  413. let mut p = [NULL_ID; N_EVENT_PARENTS];
  414. p[0] = e2.id();
  415. let e3 = Event {
  416. header: Header {
  417. timestamp: now,
  418. parents: p,
  419. layer: 2,
  420. content_hash: blake3::hash(b"e3"),
  421. },
  422. content: b"e3".to_vec(),
  423. };
  424. slot.main_tree.insert(e3.id().as_bytes(), serialize_async(&e3).await).unwrap();
  425. let mut p = [NULL_ID; N_EVENT_PARENTS];
  426. p[0] = genesis_hash;
  427. let e4 = Event {
  428. header: Header {
  429. timestamp: now,
  430. parents: p,
  431. layer: 1,
  432. content_hash: blake3::hash(b"e4"),
  433. },
  434. content: b"e4".to_vec(),
  435. };
  436. slot.main_tree.insert(e4.id().as_bytes(), serialize_async(&e4).await).unwrap();
  437. assert_ne!(e2.id(), e4.id(), "e2 and e4 must have distinct IDs");
  438. let tips = compute_unreferenced_tips(&slot.main_tree).await.unwrap();
  439. assert!(tips.get(&2).unwrap().contains(&e3.id()));
  440. assert!(tips.get(&1).unwrap().contains(&e4.id()));
  441. assert!(!tips.values().any(|set| set.contains(&e2.id())));
  442. })
  443. }
  444. #[test]
  445. fn evgr_dag_insert_valid_and_duplicate() {
  446. // First insert: returns the id, updates tips, fires the
  447. // subscriber. Second (duplicate) insert: returns an empty list,
  448. // doesn't re-fire.
  449. smol::block_on(async {
  450. let eg = make_eg().await;
  451. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  452. let dag_name = dag_ts.to_string();
  453. let sub = eg.event_pub.clone().subscribe().await;
  454. let event = Event::new(b"hello".to_vec(), &eg).await.unwrap();
  455. eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  456. let ids = eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  457. assert_eq!(ids.len(), 1);
  458. let store = eg.dag_store.read().await;
  459. let slot = store.get_slot(&dag_ts).unwrap();
  460. assert!(slot.tips.get(&1).unwrap().contains(&event.id()));
  461. drop(store);
  462. let Ok(notified) = timeout(Duration::from_secs(1), sub.receive()).await else {
  463. panic!("Event notification not received");
  464. };
  465. assert_eq!(notified.id(), event.id());
  466. // Re-insert is a no-op.
  467. assert!(eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap().is_empty());
  468. })
  469. }
  470. #[test]
  471. fn evgr_duplicate_header_insert_does_not_duplicate_time_index() {
  472. smol::block_on(async {
  473. let eg = make_eg().await;
  474. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  475. let dag_name = dag_ts.to_string();
  476. let event = Event::new(b"time-index-dedup".to_vec(), &eg).await.unwrap();
  477. eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  478. let indexed_after_first = {
  479. let store = eg.dag_store.read().await;
  480. store.get_slot(&dag_ts).unwrap().time_index.len()
  481. };
  482. eg.header_dag_insert(vec![event.header.clone(), event.header.clone()], &dag_name)
  483. .await
  484. .unwrap();
  485. let indexed_after_duplicates = {
  486. let store = eg.dag_store.read().await;
  487. store.get_slot(&dag_ts).unwrap().time_index.len()
  488. };
  489. assert_eq!(indexed_after_duplicates, indexed_after_first);
  490. let ids = eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  491. assert_eq!(ids, vec![event.id()]);
  492. let page = eg.fetch_page(u64::MAX, SyncDirection::Backward, 10).await.unwrap();
  493. assert_eq!(page.iter().filter(|ev| ev.id() == event.id()).count(), 1);
  494. })
  495. }
  496. #[test]
  497. fn evgr_dag_insert_without_header_skipped() {
  498. smol::block_on(async {
  499. let eg = make_eg().await;
  500. let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
  501. let event = Event::new(b"orphan".to_vec(), &eg).await.unwrap();
  502. let ids = eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  503. assert!(ids.is_empty());
  504. })
  505. }
  506. #[test]
  507. fn evgr_header_insert_rejects_layer_jump() {
  508. smol::block_on(async {
  509. let eg = make_eg().await;
  510. let genesis = eg.current_genesis.read().await.clone();
  511. let dag_name = genesis.header.timestamp.to_string();
  512. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  513. parents[0] = genesis.id();
  514. let event = Event {
  515. header: Header {
  516. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
  517. parents,
  518. layer: 2,
  519. content_hash: blake3::hash(b"layer-jump"),
  520. },
  521. content: b"layer-jump".to_vec(),
  522. };
  523. let err = eg.header_dag_insert(vec![event.header], &dag_name).await.unwrap_err();
  524. assert!(matches!(err, crate::Error::HeaderIsInvalid));
  525. })
  526. }
  527. #[test]
  528. fn evgr_header_insert_rejects_duplicate_parents() {
  529. smol::block_on(async {
  530. let eg = make_eg().await;
  531. let genesis = eg.current_genesis.read().await.clone();
  532. let dag_name = genesis.header.timestamp.to_string();
  533. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  534. parents[0] = genesis.id();
  535. parents[1] = genesis.id();
  536. let event = Event {
  537. header: Header {
  538. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
  539. parents,
  540. layer: 1,
  541. content_hash: blake3::hash(b"duplicate-parents"),
  542. },
  543. content: b"duplicate-parents".to_vec(),
  544. };
  545. let err = eg.header_dag_insert(vec![event.header], &dag_name).await.unwrap_err();
  546. assert!(matches!(err, crate::Error::HeaderIsInvalid));
  547. })
  548. }
  549. #[test]
  550. fn evgr_header_validate_rejects_layer_overflow_parent() {
  551. smol::block_on(async {
  552. let db = sled::Config::new().temporary(true).open().unwrap();
  553. let tree = db.open_tree("headers").unwrap();
  554. let timestamp = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  555. let parent = Header {
  556. timestamp,
  557. parents: NULL_PARENTS,
  558. layer: u64::MAX,
  559. content_hash: blake3::hash(b"overflow-parent"),
  560. };
  561. tree.insert(parent.id().as_bytes(), serialize_async(&parent).await).unwrap();
  562. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  563. parents[0] = parent.id();
  564. let child = Header {
  565. timestamp,
  566. parents,
  567. layer: u64::MAX,
  568. content_hash: blake3::hash(b"overflow-child"),
  569. };
  570. assert!(!child.validate(&tree, &test_config(), timestamp, None).await.unwrap());
  571. })
  572. }
  573. #[test]
  574. fn evgr_header_validate_uses_target_slot_bounds() {
  575. smol::block_on(async {
  576. const HOUR_MS: u64 = 3_600_000;
  577. let db = sled::Config::new().temporary(true).open().unwrap();
  578. let tree = db.open_tree("headers").unwrap();
  579. let dag_ts = 1_704_067_200_000;
  580. let drift = crate::event_graph::EVENT_TIME_DRIFT;
  581. let config = EventGraphConfig { hours_rotation: 6, ..test_config() };
  582. let genesis = Header {
  583. timestamp: dag_ts,
  584. parents: NULL_PARENTS,
  585. layer: 0,
  586. content_hash: blake3::hash(&config.genesis_contents),
  587. };
  588. tree.insert(genesis.id().as_bytes(), serialize_async(&genesis).await).unwrap();
  589. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  590. parents[0] = genesis.id();
  591. let make_header = |timestamp, content: &[u8]| Header {
  592. timestamp,
  593. parents,
  594. layer: 1,
  595. content_hash: blake3::hash(content),
  596. };
  597. let lower_edge = make_header(dag_ts.saturating_sub(drift), b"lower-edge");
  598. assert!(lower_edge.validate(&tree, &config, dag_ts, None).await.unwrap());
  599. let upper_edge = make_header(dag_ts + 6 * HOUR_MS + drift - 1, b"upper-edge");
  600. assert!(upper_edge.validate(&tree, &config, dag_ts, None).await.unwrap());
  601. let too_early = make_header(dag_ts - drift - 1, b"too-early");
  602. assert!(!too_early.validate(&tree, &config, dag_ts, None).await.unwrap());
  603. let too_late = make_header(dag_ts + 6 * HOUR_MS + drift, b"too-late");
  604. assert!(!too_late.validate(&tree, &config, dag_ts, None).await.unwrap());
  605. })
  606. }
  607. #[test]
  608. fn evgr_header_validate_ignores_future_initial_genesis() {
  609. smol::block_on(async {
  610. const HOUR_MS: u64 = 3_600_000;
  611. let db = sled::Config::new().temporary(true).open().unwrap();
  612. let tree = db.open_tree("headers").unwrap();
  613. let now = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  614. let dag_ts = now.saturating_sub(HOUR_MS);
  615. let config =
  616. EventGraphConfig { initial_genesis: now + HOUR_MS, hours_rotation: 1, ..test_config() };
  617. let genesis = Header {
  618. timestamp: dag_ts,
  619. parents: NULL_PARENTS,
  620. layer: 0,
  621. content_hash: blake3::hash(&config.genesis_contents),
  622. };
  623. tree.insert(genesis.id().as_bytes(), serialize_async(&genesis).await).unwrap();
  624. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  625. parents[0] = genesis.id();
  626. let child = Header {
  627. timestamp: dag_ts + 1,
  628. parents,
  629. layer: 1,
  630. content_hash: blake3::hash(b"future-initial-genesis"),
  631. };
  632. assert!(child.validate(&tree, &config, dag_ts, None).await.unwrap());
  633. })
  634. }
  635. #[test]
  636. fn evgr_header_validate_no_rotation_rejects_far_future() {
  637. smol::block_on(async {
  638. let db = sled::Config::new().temporary(true).open().unwrap();
  639. let tree = db.open_tree("headers").unwrap();
  640. let config = test_config();
  641. let dag_ts = config.initial_genesis;
  642. let genesis = Header {
  643. timestamp: dag_ts,
  644. parents: NULL_PARENTS,
  645. layer: 0,
  646. content_hash: blake3::hash(&config.genesis_contents),
  647. };
  648. tree.insert(genesis.id().as_bytes(), serialize_async(&genesis).await).unwrap();
  649. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  650. parents[0] = genesis.id();
  651. let old_history = Header {
  652. timestamp: dag_ts + 1,
  653. parents,
  654. layer: 1,
  655. content_hash: blake3::hash(b"old-no-rotation-history"),
  656. };
  657. assert!(old_history.validate(&tree, &config, dag_ts, None).await.unwrap());
  658. let future = Header {
  659. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64 +
  660. crate::event_graph::EVENT_TIME_DRIFT +
  661. 1,
  662. parents,
  663. layer: 1,
  664. content_hash: blake3::hash(b"future-no-rotation-header"),
  665. };
  666. assert!(!future.validate(&tree, &config, dag_ts, None).await.unwrap());
  667. })
  668. }
  669. #[test]
  670. fn evgr_header_insert_rejects_unloaded_dag_slot() {
  671. smol::block_on(async {
  672. let config = EventGraphConfig { hours_rotation: 1, max_dags: Some(2), ..test_config() };
  673. let eg = make_eg_with_config(config).await;
  674. let dag_ts = next_hour_timestamp(-100);
  675. let dag_name = dag_ts.to_string();
  676. let genesis = Header {
  677. timestamp: dag_ts,
  678. parents: NULL_PARENTS,
  679. layer: 0,
  680. content_hash: blake3::hash(&eg.config.genesis_contents),
  681. };
  682. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  683. parents[0] = genesis.id();
  684. let header = Header {
  685. timestamp: dag_ts + 1,
  686. parents,
  687. layer: 1,
  688. content_hash: blake3::hash(b"unloaded-slot"),
  689. };
  690. let err = eg.header_dag_insert(vec![header], &dag_name).await.unwrap_err();
  691. assert!(matches!(err, crate::Error::DagSyncFailed));
  692. })
  693. }
  694. #[test]
  695. fn evgr_fetch_headers_with_tips_is_bounded_and_layer_ordered() {
  696. smol::block_on(async {
  697. let eg = make_eg().await;
  698. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  699. let dag_name = dag_ts.to_string();
  700. let base = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  701. let mut parent = eg.current_genesis.read().await.id();
  702. let mut headers = Vec::with_capacity(MAX_HEADER_REP_HEADERS + 32);
  703. for i in 0..(MAX_HEADER_REP_HEADERS + 32) {
  704. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  705. parents[0] = parent;
  706. let content = format!("bounded-header-{i}");
  707. let header = Header {
  708. timestamp: base + i as u64,
  709. parents,
  710. layer: i as u64 + 1,
  711. content_hash: blake3::hash(content.as_bytes()),
  712. };
  713. parent = header.id();
  714. headers.push(header);
  715. }
  716. eg.header_dag_insert(headers, &dag_name).await.unwrap();
  717. let hostile_empty_tips = LayerUTips::new();
  718. let response = eg.fetch_headers_with_tips(&dag_name, &hostile_empty_tips).await.unwrap();
  719. assert_eq!(response.len(), MAX_HEADER_REP_HEADERS);
  720. for pair in response.windows(2) {
  721. assert!(pair[0].layer <= pair[1].layer);
  722. }
  723. assert_eq!(response.first().unwrap().layer, 0);
  724. assert!(response.last().unwrap().layer < MAX_HEADER_REP_HEADERS as u64);
  725. })
  726. }
  727. #[test]
  728. fn evgr_fetch_page_both_directions() {
  729. smol::block_on(async {
  730. let eg = make_eg().await;
  731. let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
  732. let base = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  733. for i in 0..(MAX_RANGE_PAGE_SIZE as u64 + 10) {
  734. let ev = Event::with_timestamp(base + i, vec![(i % 251) as u8], &eg).await.unwrap();
  735. eg.header_dag_insert(vec![ev.header.clone()], &dag_name).await.unwrap();
  736. eg.dag_insert(slice::from_ref(&ev), &dag_name).await.unwrap();
  737. }
  738. let page = eg.fetch_page(u64::MAX, SyncDirection::Backward, 5).await.unwrap();
  739. assert_eq!(page.len(), 5);
  740. for w in page.windows(2) {
  741. assert!(w[0].header.timestamp >= w[1].header.timestamp);
  742. }
  743. let page = eg.fetch_page(0, SyncDirection::Forward, 5).await.unwrap();
  744. assert!(!page.is_empty());
  745. for w in page.windows(2) {
  746. assert!(w[0].header.timestamp <= w[1].header.timestamp);
  747. }
  748. let capped = eg
  749. .fetch_page(u64::MAX, SyncDirection::Backward, MAX_RANGE_PAGE_SIZE + 10)
  750. .await
  751. .unwrap();
  752. assert_eq!(capped.len(), MAX_RANGE_PAGE_SIZE);
  753. })
  754. }
  755. #[test]
  756. fn evgr_order_events_rejects_corrupt_event_record() {
  757. smol::block_on(async {
  758. let eg = make_eg().await;
  759. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  760. let bad_id = [0u8; 32];
  761. {
  762. let store = eg.dag_store.read().await;
  763. let slot = store.get_slot(&dag_ts).unwrap();
  764. slot.main_tree.insert(bad_id.as_slice(), b"not-an-event".as_slice()).unwrap();
  765. }
  766. let result = eg.order_events().await;
  767. assert!(result.is_err(), "corrupt event records should fail history ordering");
  768. })
  769. }
  770. async fn build_graph() -> Result<(EventGraphPtr, std::collections::HashMap<&'static str, Event>)> {
  771. let eg = make_eg().await;
  772. let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
  773. let genesis_hash = eg.current_genesis.read().await.id();
  774. let base = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
  775. let make = |off: u64, layer: u64, parents: [blake3::Hash; N_EVENT_PARENTS], name: &str| Event {
  776. header: Header {
  777. timestamp: base + off,
  778. layer,
  779. parents,
  780. content_hash: blake3::hash(name.as_bytes()),
  781. },
  782. content: name.as_bytes().to_vec(),
  783. };
  784. let mut p = [NULL_ID; N_EVENT_PARENTS];
  785. p[0] = genesis_hash;
  786. let e1a = make(1, 1, p, "e1a");
  787. let e1b = make(2, 1, p, "e1b");
  788. let e1c = make(3, 1, p, "e1c");
  789. let e1d = make(4, 1, p, "e1d");
  790. let mut p = [NULL_ID; N_EVENT_PARENTS];
  791. p[0] = e1a.id();
  792. let e2a = make(5, 2, p, "e2a");
  793. p[0] = e1b.id();
  794. let e2b = make(6, 2, p, "e2b");
  795. p[0] = e1c.id();
  796. let e2c = make(7, 2, p, "e2c");
  797. p[0] = e1d.id();
  798. let e2d = make(8, 2, p, "e2d");
  799. let l1 = vec![e1a.clone(), e1b.clone(), e1c.clone(), e1d.clone()];
  800. let l2 = vec![e2a.clone(), e2b.clone(), e2c.clone(), e2d.clone()];
  801. eg.header_dag_insert(l1.iter().map(|e| e.header.clone()).collect(), &dag_name).await.unwrap();
  802. eg.dag_insert(&l1, &dag_name).await.unwrap();
  803. eg.header_dag_insert(l2.iter().map(|e| e.header.clone()).collect(), &dag_name).await.unwrap();
  804. eg.dag_insert(&l2, &dag_name).await.unwrap();
  805. let mut map = std::collections::HashMap::new();
  806. map.insert("e1a", e1a);
  807. map.insert("e1b", e1b);
  808. map.insert("e1c", e1c);
  809. map.insert("e1d", e1d);
  810. map.insert("e2a", e2a);
  811. map.insert("e2b", e2b);
  812. map.insert("e2c", e2c);
  813. map.insert("e2d", e2d);
  814. Ok((eg, map))
  815. }
  816. #[test]
  817. fn evgr_ancestor_walk_via_header_tree() {
  818. smol::block_on(async {
  819. let (eg, evs) = build_graph().await.unwrap();
  820. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  821. let store = eg.dag_store.read().await;
  822. let slot = store.get_slot(&dag_ts).unwrap();
  823. let genesis_hash = eg.current_genesis.read().await.id();
  824. for name in ["e1a", "e1b", "e1c", "e1d"] {
  825. let mut ancestors = HashSet::new();
  826. eg.get_ancestors(&mut ancestors, evs[name].header.clone(), &slot.header_tree)
  827. .await
  828. .unwrap();
  829. assert_eq!(ancestors, HashSet::from([genesis_hash]));
  830. }
  831. let mut ancestors = HashSet::new();
  832. eg.get_ancestors(&mut ancestors, evs["e2a"].header.clone(), &slot.header_tree)
  833. .await
  834. .unwrap();
  835. assert_eq!(ancestors, HashSet::from([genesis_hash, evs["e1a"].id()]));
  836. })
  837. }
  838. #[test]
  839. fn evgr_multi_node_propagation_with_real_blob() {
  840. init_logger();
  841. run_multi_node_test(propagation_with_real_blob);
  842. }
  843. async fn propagation_with_real_blob(ex: Arc<Executor<'static>>) {
  844. let nodes = make_network(ex).await;
  845. let mut alice = TestIdentity::new();
  846. for eg in &nodes {
  847. alice.register_directly(eg).await.expect("register alice");
  848. }
  849. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  850. let dag_name = dag_ts.to_string();
  851. let event = Event::new(b"hello-via-rln".to_vec(), &nodes[0]).await.unwrap();
  852. let message_id =
  853. alice.next_message_id(event.header.timestamp).expect("budget available on first signal");
  854. let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
  855. let blob = serialize_async(&blob_struct).await;
  856. nodes[0].header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  857. nodes[0].dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  858. nodes[0].dag_blob_store(&event.id(), &blob).unwrap();
  859. nodes[0].p2p.broadcast(&EventPut(event.clone(), blob.clone())).await;
  860. sleep(5).await;
  861. for (i, eg) in nodes.iter().enumerate() {
  862. let store = eg.dag_store.read().await;
  863. let slot = store.get_slot(&dag_ts).unwrap();
  864. assert!(
  865. slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  866. "node {i} missing event in main_tree",
  867. );
  868. drop(store);
  869. assert!(
  870. eg.dag_blob_fetch(&event.id()).unwrap().is_some(),
  871. "node {i} missing blob in dag_blobs - late-joiners would have nothing to verify against",
  872. );
  873. }
  874. shutdown_network(&nodes).await;
  875. }
  876. #[test]
  877. fn evgr_multi_node_empty_blob_rejected() {
  878. init_logger();
  879. run_multi_node_test(empty_blob_rejected);
  880. }
  881. async fn empty_blob_rejected(ex: Arc<Executor<'static>>) {
  882. // An attacker broadcasts `EventPut(ev, vec![])` for a non-genesis event.
  883. // Every recipient must strike the sender and refuse to insert.
  884. let nodes = make_network(ex).await;
  885. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  886. let event = Event::new(b"unauthenticated".to_vec(), &nodes[0]).await.unwrap();
  887. nodes[0].p2p.broadcast(&EventPut(event.clone(), vec![])).await;
  888. sleep(5).await;
  889. for (i, eg) in nodes.iter().enumerate().skip(1) {
  890. let store = eg.dag_store.read().await;
  891. let slot = store.get_slot(&dag_ts).unwrap();
  892. assert!(
  893. !slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  894. "Vector-1 breach: node {i} accepted an empty-blob non-genesis event",
  895. );
  896. }
  897. shutdown_network(&nodes).await;
  898. }
  899. #[test]
  900. fn evgr_multi_node_malformed_event_rejected_before_rln() {
  901. init_logger();
  902. run_multi_node_test(malformed_event_rejected_before_rln);
  903. }
  904. async fn malformed_event_rejected_before_rln(ex: Arc<Executor<'static>>) {
  905. let nodes = make_network(ex).await;
  906. let mut alice = TestIdentity::new();
  907. for eg in &nodes {
  908. alice.register_directly(eg).await.unwrap();
  909. }
  910. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  911. let event = Event::new(b"preflight-live".to_vec(), &nodes[0]).await.unwrap();
  912. let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
  913. let blob = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
  914. let internal_nullifier = blob.internal_nullifier;
  915. let blob = serialize_async(&blob).await;
  916. let mut malformed = event.clone();
  917. malformed.content.extend_from_slice(b"-tampered");
  918. assert!(!malformed.content_matches_header());
  919. nodes[0].p2p.broadcast(&EventPut(malformed.clone(), blob)).await;
  920. sleep(5).await;
  921. let epoch = epoch_of(malformed.header.timestamp);
  922. for (i, eg) in nodes.iter().enumerate().skip(1) {
  923. let store = eg.dag_store.read().await;
  924. let slot = store.get_slot(&dag_ts).unwrap();
  925. assert!(
  926. !slot.main_tree.contains_key(malformed.id().as_bytes()).unwrap(),
  927. "node {i} accepted a structurally invalid event",
  928. );
  929. drop(store);
  930. let state = eg.rln_state.read().await;
  931. assert!(
  932. !state.metadata.is_reused(epoch, &internal_nullifier),
  933. "node {i} ran RLN verification before structural rejection",
  934. );
  935. }
  936. shutdown_network(&nodes).await;
  937. }
  938. #[test]
  939. fn evgr_multi_node_genesis_with_blob_rejected() {
  940. init_logger();
  941. run_multi_node_test(genesis_with_blob_rejected);
  942. }
  943. async fn genesis_with_blob_rejected(ex: Arc<Executor<'static>>) {
  944. // Symmetric defense: a genesis-shaped event arriving with a
  945. // non-empty blob is also misbehavior. Genesis events are
  946. // deterministic and don't carry signals.
  947. let nodes = make_network(ex).await;
  948. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  949. let header = Header {
  950. timestamp: dag_ts,
  951. parents: NULL_PARENTS,
  952. layer: 0,
  953. content_hash: blake3::hash(b"forged-genesis"),
  954. };
  955. let event = Event { header, content: b"forged-genesis".to_vec() };
  956. let fake_blob = b"this-should-not-be-here".to_vec();
  957. nodes[0].p2p.broadcast(&EventPut(event.clone(), fake_blob)).await;
  958. sleep(5).await;
  959. for (i, eg) in nodes.iter().enumerate().skip(1) {
  960. let store = eg.dag_store.read().await;
  961. let slot = store.get_slot(&dag_ts).unwrap();
  962. assert!(
  963. !slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  964. "node {i} accepted a genesis-shaped event with a blob",
  965. );
  966. }
  967. shutdown_network(&nodes).await;
  968. }
  969. #[test]
  970. fn evgr_fetch_missing_events_rejects_corrupt_header_record() {
  971. smol::block_on(async {
  972. let eg = make_eg().await;
  973. let dag_ts = eg.current_genesis.read().await.header.timestamp;
  974. let dag_name = dag_ts.to_string();
  975. let bad_id = [0u8; 32];
  976. {
  977. let store = eg.dag_store.read().await;
  978. let slot = store.get_slot(&dag_ts).unwrap();
  979. slot.header_tree.insert(bad_id.as_slice(), b"not-a-header".as_slice()).unwrap();
  980. }
  981. let result = eg.fetch_missing_events(dag_ts, &dag_name, 1).await;
  982. assert!(result.is_err(), "corrupt header records should fail DAG content sync");
  983. })
  984. }
  985. #[test]
  986. fn evgr_multi_node_dag_sync_with_blob() {
  987. init_logger();
  988. run_multi_node_test(dag_sync_with_blob);
  989. }
  990. async fn dag_sync_with_blob(ex: Arc<Executor<'static>>) {
  991. // End-to-end Vector-2 propagation test. Nodes 0..3 receive a
  992. // signal via direct insert (with the blob in their dag_blobs
  993. // side-table). Node 4 catches up via dag_sync - it should
  994. // receive both the event and the blob from a peer, and re-verify
  995. // the proof at sync time.
  996. let nodes = make_network(ex).await;
  997. let mut alice = TestIdentity::new();
  998. for eg in &nodes {
  999. alice.register_directly(eg).await.unwrap();
  1000. }
  1001. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  1002. let dag_name = dag_ts.to_string();
  1003. let event = Event::new(b"synced-message".to_vec(), &nodes[0]).await.unwrap();
  1004. let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
  1005. let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
  1006. let blob = serialize_async(&blob_struct).await;
  1007. for eg in nodes.iter().take(4) {
  1008. eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  1009. eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  1010. eg.dag_blob_store(&event.id(), &blob).unwrap();
  1011. }
  1012. {
  1013. let store = nodes[4].dag_store.read().await;
  1014. let slot = store.get_slot(&dag_ts).unwrap();
  1015. assert!(!slot.main_tree.contains_key(event.id().as_bytes()).unwrap());
  1016. }
  1017. nodes[4].dag_sync(dag_ts).await.expect("dag_sync should succeed");
  1018. sleep(2).await;
  1019. let store = nodes[4].dag_store.read().await;
  1020. let slot = store.get_slot(&dag_ts).unwrap();
  1021. assert!(
  1022. slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  1023. "sync should bring the event over",
  1024. );
  1025. drop(store);
  1026. assert!(
  1027. nodes[4].dag_blob_fetch(&event.id()).unwrap().is_some(),
  1028. "sync should bring the blob over too - without it, node 4 can't serve future late-joiners",
  1029. );
  1030. shutdown_network(&nodes).await;
  1031. }
  1032. #[test]
  1033. fn evgr_multi_node_dag_sync_rejects_bad_blob() {
  1034. init_logger();
  1035. run_multi_node_test(dag_sync_rejects_bad_blob);
  1036. }
  1037. async fn dag_sync_rejects_bad_blob(ex: Arc<Executor<'static>>) {
  1038. // Vector-2 defense: a peer in the 2/3 quorum serves a tampered
  1039. // blob during sync. The recipient's dag_insert_with_blobs runs
  1040. // RLN re-verification, which rejects, and the event does NOT
  1041. // end up in the recipient's main_tree.
  1042. let nodes = make_network(ex).await;
  1043. let alice = TestIdentity::new();
  1044. for eg in &nodes {
  1045. alice.register_directly(eg).await.unwrap();
  1046. }
  1047. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  1048. let dag_name = dag_ts.to_string();
  1049. let event = Event::new(b"crafted-injection".to_vec(), &nodes[0]).await.unwrap();
  1050. // Garbage bytes - won't deserialize as a real RLN signal, won't
  1051. // verify. We don't need a real (failing) proof to exercise the
  1052. // rejection path.
  1053. let bad_blob = b"definitely-not-a-real-rln-blob".to_vec();
  1054. for eg in nodes.iter().take(4) {
  1055. eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  1056. eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  1057. eg.dag_blob_store(&event.id(), &bad_blob).unwrap();
  1058. }
  1059. nodes[4].dag_sync(dag_ts).await.unwrap();
  1060. sleep(2).await;
  1061. let store = nodes[4].dag_store.read().await;
  1062. let slot = store.get_slot(&dag_ts).unwrap();
  1063. assert!(
  1064. !slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  1065. "Vector-2 breach: node 4 accepted an event with a tampered blob during sync",
  1066. );
  1067. shutdown_network(&nodes).await;
  1068. }
  1069. #[test]
  1070. fn evgr_multi_node_dormant_user_can_post_after_long_silence() {
  1071. init_logger();
  1072. run_multi_node_test(dormant_user_can_post_after_long_silence);
  1073. }
  1074. async fn dormant_user_can_post_after_long_silence(ex: Arc<Executor<'static>>) {
  1075. // Alice registers, then 17+ other identities also register. By
  1076. // the time Alice tries to send a signal, her registration root
  1077. // has long fallen out of the in-memory recent_roots window. The
  1078. // historical-roots side-table makes verification succeed anyway.
  1079. //
  1080. // Note: in this test Alice's signal still references the CURRENT
  1081. // root (because `rln_membership_path` returns the live root and
  1082. // Alice is still a member). That's fine - what we're validating
  1083. // here is end-to-end: a deeply-historical state of the SMT
  1084. // doesn't break verification. The unit tests in tests_rln.rs
  1085. // (rln_is_root_valid_at_*) cover the predicate's exact semantics
  1086. // for old-root references.
  1087. let nodes = make_network(ex).await;
  1088. let mut alice = TestIdentity::new();
  1089. for eg in &nodes {
  1090. alice.register_directly(eg).await.unwrap();
  1091. }
  1092. // Push the recent_roots window past Alice's registration.
  1093. for seed in 100..117_u64 {
  1094. let other = TestIdentity::with_seed(seed);
  1095. for eg in &nodes {
  1096. other.register_directly(eg).await.unwrap();
  1097. }
  1098. }
  1099. let event = Event::new(b"long-silent-but-still-registered".to_vec(), &nodes[0]).await.unwrap();
  1100. let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
  1101. let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
  1102. let blob = serialize_async(&blob_struct).await;
  1103. let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
  1104. let dag_name = dag_ts.to_string();
  1105. nodes[0].header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  1106. nodes[0].dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  1107. nodes[0].dag_blob_store(&event.id(), &blob).unwrap();
  1108. nodes[0].p2p.broadcast(&EventPut(event.clone(), blob)).await;
  1109. sleep(5).await;
  1110. for (i, eg) in nodes.iter().enumerate() {
  1111. let store = eg.dag_store.read().await;
  1112. let slot = store.get_slot(&dag_ts).unwrap();
  1113. assert!(
  1114. slot.main_tree.contains_key(event.id().as_bytes()).unwrap(),
  1115. "node {i} rejected Alice's signal even though she's a valid registered identity \
  1116. - historical-roots fallback is broken",
  1117. );
  1118. }
  1119. shutdown_network(&nodes).await;
  1120. }