mod.rs 114 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944
  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. //! Multi-DAG Event Graph with bidirectional sync, RLN rate limiting,
  19. //! and periodic DAG rotation.
  20. use std::{
  21. collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque},
  22. path::PathBuf,
  23. str::FromStr,
  24. sync::{
  25. atomic::{AtomicBool, Ordering},
  26. Arc,
  27. },
  28. };
  29. use darkfi_sdk::{crypto::pasta_prelude::PrimeField, pasta::pallas};
  30. use darkfi_serial::{deserialize_async, deserialize_async_partial, serialize_async};
  31. use futures::{stream::FuturesUnordered, StreamExt};
  32. use sled_overlay::{sled, SledTreeOverlay};
  33. use smol::{
  34. lock::{OnceCell, RwLock},
  35. Executor,
  36. };
  37. use tracing::{error, info, warn};
  38. use url::Url;
  39. use crate::{
  40. net::{channel::Channel, P2pPtr},
  41. system::{msleep, Publisher, PublisherPtr, StoppableTask, StoppableTaskPtr, Subscription},
  42. Error, Result,
  43. };
  44. pub mod event;
  45. pub use event::{display_order, Event, Header};
  46. pub mod proto;
  47. use proto::{
  48. cap_layer_tips, count_layer_tips, EventRep, EventReq, HeaderRep, HeaderReq, StaticPut,
  49. SyncDirection, TipRep, TipReq, MAX_EVENT_REP_EVENTS, MAX_EVENT_REQ_IDS, MAX_HEADER_REP_HEADERS,
  50. MAX_HEADER_REQ_TIPS, MAX_RANGE_PAGE_SIZE, MAX_TIP_REP_TIPS,
  51. };
  52. pub mod rln;
  53. use rln::{IdentityState, RlnState, ZkKeys};
  54. pub mod util;
  55. use util::{
  56. generate_genesis, millis_until_next_rotation, next_hour_timestamp, next_rotation_timestamp,
  57. replayer_log,
  58. };
  59. pub mod deg;
  60. use deg::DegEvent;
  61. #[cfg(test)]
  62. mod tests;
  63. #[cfg(test)]
  64. mod tests_rln;
  65. #[cfg(test)]
  66. mod test_helpers;
  67. /// Number of parent references each event carries.
  68. pub const N_EVENT_PARENTS: usize = 5;
  69. /// Allowed timestamp drift in milliseconds.
  70. const EVENT_TIME_DRIFT: u64 = 60_000;
  71. /// The null event ID (32 zero bytes).
  72. pub const NULL_ID: blake3::Hash = blake3::Hash::from_bytes([0x00; blake3::OUT_LEN]);
  73. /// Array of null parents (used by genesis events).
  74. pub const NULL_PARENTS: [blake3::Hash; N_EVENT_PARENTS] = [NULL_ID; N_EVENT_PARENTS];
  75. /// Maximum number of static-DAG events `static_sync` will pull in
  76. /// one invocation. Defends against malicious deep-ancestry chains.
  77. const SYNC_MAX_STATIC_EVENTS: usize = 100_000;
  78. /// Runtime configuration for an Event Graph instance.
  79. #[derive(Clone, Debug)]
  80. pub struct EventGraphConfig {
  81. /// Epoch origin timestamp in millis.
  82. /// All rotation boundaries are computed as offsets from this point.
  83. /// Should be UTC midnight for clean hourly alignment.
  84. pub initial_genesis: u64,
  85. /// How often the DAG rotates, in hours. 0 = no rotation.
  86. pub hours_rotation: u64,
  87. /// Unique payload embedded in genesis events.
  88. /// Different protocols must use different values.
  89. pub genesis_contents: Vec<u8>,
  90. /// App-provided pregenerated RLN identity commitments.
  91. ///
  92. /// EventGraph treats these as the only proof-less registration
  93. /// commitments accepted with [`rln::GENESIS_BLOB_GUARD`]. Apps
  94. /// that do not use pregenerated RLN identities should leave this
  95. /// empty.
  96. pub pregenerated_identity_commitments: Vec<[u8; 32]>,
  97. /// Maximum number of DAGs to keep in the rolling window.
  98. ///
  99. /// * `Some(n)` - keep n rotation periods.
  100. /// When the n+1 period is created, the oldest is permanently
  101. /// deleted from sled. This is the normal mode for end-user nodes.
  102. /// * `None` - never prune. Every DAG ever created is kept in sled
  103. /// and loaded at startup. This is archive mode for nodes that want
  104. /// complete history.
  105. ///
  106. /// With `hours_rotation = 1` and `max_dags = Some(24)`, events
  107. /// older than 24 hours are lost. With `hours_rotation = 6` and
  108. /// `max_dags = Some(24)`, the window is 6 days.
  109. pub max_dags: Option<usize>,
  110. }
  111. impl EventGraphConfig {
  112. /// Validate consensus-critical event graph configuration.
  113. pub fn validate(&self) -> Result<()> {
  114. if self.max_dags == Some(0) {
  115. return Err(Error::Custom("event graph max_dags must be greater than 0".into()))
  116. }
  117. self.rotation_period_millis()?;
  118. Ok(())
  119. }
  120. /// Rotation period in milliseconds, or `None` for non-rotating graphs.
  121. pub(crate) fn rotation_period_millis(&self) -> Result<Option<u64>> {
  122. if self.hours_rotation == 0 {
  123. return Ok(None)
  124. }
  125. let rotation_ms = self.hours_rotation.checked_mul(util::HOUR_MS).ok_or_else(|| {
  126. Error::Custom("event graph rotation period overflows milliseconds".into())
  127. })?;
  128. if self.initial_genesis.checked_add(rotation_ms).is_none() {
  129. return Err(Error::Custom(
  130. "event graph initial genesis plus one rotation overflows".into(),
  131. ))
  132. }
  133. Ok(Some(rotation_ms))
  134. }
  135. }
  136. pub type EventGraphPtr = Arc<EventGraph>;
  137. /// Unreferenced tips grouped by layer.
  138. pub type LayerUTips = BTreeMap<u64, HashSet<blake3::Hash>>;
  139. /// Generate the deterministic genesis event for the static DAG.
  140. fn generate_static_genesis(config: &EventGraphConfig) -> Event {
  141. let header = Header {
  142. timestamp: config.initial_genesis,
  143. parents: NULL_PARENTS,
  144. layer: 0,
  145. content_hash: blake3::hash(&config.genesis_contents),
  146. };
  147. Event { header, content: config.genesis_contents.clone() }
  148. }
  149. fn validate_pregenerated_identity_commitments(
  150. config: &EventGraphConfig,
  151. ) -> Result<(Vec<pallas::Base>, HashSet<[u8; 32]>)> {
  152. let mut commitments = Vec::with_capacity(config.pregenerated_identity_commitments.len());
  153. let mut reprs = HashSet::with_capacity(config.pregenerated_identity_commitments.len());
  154. for (index, repr) in config.pregenerated_identity_commitments.iter().enumerate() {
  155. if !reprs.insert(*repr) {
  156. return Err(Error::Custom(format!(
  157. "duplicate pregenerated identity commitment at index {index}"
  158. )))
  159. }
  160. let commitment: Option<pallas::Base> = pallas::Base::from_repr(*repr).into();
  161. let Some(commitment) = commitment else {
  162. return Err(Error::Custom(format!(
  163. "invalid pregenerated identity commitment at index {index}"
  164. )))
  165. };
  166. commitments.push(commitment);
  167. }
  168. Ok((commitments, reprs))
  169. }
  170. /// Bidirectional timestamp -> event-ID index.
  171. #[derive(Clone, Debug, Default)]
  172. pub struct TimeIndex {
  173. index: BTreeMap<u64, Vec<blake3::Hash>>,
  174. count: usize,
  175. }
  176. impl TimeIndex {
  177. pub fn new() -> Self {
  178. Self::default()
  179. }
  180. pub async fn from_header_dag(tree: &sled::Tree) -> Self {
  181. let mut idx = Self::new();
  182. for item in tree.iter() {
  183. let (id, hdr) = item.unwrap();
  184. let id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  185. let hdr: Header = deserialize_async(&hdr).await.unwrap();
  186. idx.insert(hdr.timestamp, id);
  187. }
  188. idx
  189. }
  190. pub fn insert(&mut self, ts: u64, id: blake3::Hash) {
  191. let ids = self.index.entry(ts).or_default();
  192. if ids.contains(&id) {
  193. return
  194. }
  195. ids.push(id);
  196. self.count += 1;
  197. }
  198. pub fn newest(&self, n: usize) -> Vec<blake3::Hash> {
  199. self.rev(u64::MAX, n)
  200. }
  201. pub fn oldest(&self, n: usize) -> Vec<blake3::Hash> {
  202. self.fwd(0, n)
  203. }
  204. pub fn before(&self, cursor: u64, n: usize) -> Vec<blake3::Hash> {
  205. self.rev(cursor.saturating_sub(1), n)
  206. }
  207. pub fn after(&self, cursor: u64, n: usize) -> Vec<blake3::Hash> {
  208. self.fwd(cursor.saturating_add(1), n)
  209. }
  210. fn rev(&self, start: u64, n: usize) -> Vec<blake3::Hash> {
  211. let mut out = Vec::with_capacity(n);
  212. for (_, ids) in self.index.range(..=start).rev() {
  213. for id in ids {
  214. out.push(*id);
  215. if out.len() >= n {
  216. return out
  217. }
  218. }
  219. }
  220. out
  221. }
  222. fn fwd(&self, start: u64, n: usize) -> Vec<blake3::Hash> {
  223. let mut out = Vec::with_capacity(n);
  224. for (_, ids) in self.index.range(start..) {
  225. for id in ids {
  226. out.push(*id);
  227. if out.len() >= n {
  228. return out
  229. }
  230. }
  231. }
  232. out
  233. }
  234. pub fn len(&self) -> usize {
  235. self.count
  236. }
  237. pub fn is_empty(&self) -> bool {
  238. self.count == 0
  239. }
  240. }
  241. /// All per-DAG state: trees, tips, and the timestamp index.
  242. pub struct DagSlot {
  243. pub header_tree: sled::Tree,
  244. pub main_tree: sled::Tree,
  245. pub tips: LayerUTips,
  246. pub time_index: TimeIndex,
  247. }
  248. /// Full-scan tip computation.
  249. /// Compute unreferenced tips - events that exist in the DAG but are
  250. /// not referenced as a parent by any other event - grouped by layer.
  251. pub(crate) async fn compute_unreferenced_tips(dag: &sled::Tree) -> LayerUTips {
  252. let mut candidates: HashMap<blake3::Hash, u64> = HashMap::new();
  253. let mut referenced: HashSet<blake3::Hash> = HashSet::new();
  254. for item in dag.iter() {
  255. let (id_bytes, val_bytes) = item.unwrap();
  256. let id = blake3::Hash::from_bytes((&id_bytes as &[u8]).try_into().unwrap());
  257. let ev: Event = deserialize_async(&val_bytes).await.unwrap();
  258. candidates.insert(id, ev.header.layer);
  259. for p in ev.header.parents.iter() {
  260. if *p != NULL_ID {
  261. referenced.insert(*p);
  262. }
  263. }
  264. }
  265. // Bucket the unreferenced candidates by their layer
  266. let mut map: LayerUTips = BTreeMap::new();
  267. for (id, layer) in candidates {
  268. if !referenced.contains(&id) {
  269. map.entry(layer).or_default().insert(id);
  270. }
  271. }
  272. map
  273. }
  274. /// Pick up to N_EVENT_PARENTS tips from the highest layers.
  275. ///
  276. /// If the highest local tip is already at `u64::MAX`, no valid child
  277. /// layer exists. Return `u64::MAX` rather than wrapping; header
  278. /// validation rejects any attempted child of that saturated layer.
  279. fn select_parents_from_tips(tips: &LayerUTips) -> (u64, [blake3::Hash; N_EVENT_PARENTS]) {
  280. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  281. let mut i = 0;
  282. 'outer: for (_, layer_tips) in tips.iter().rev() {
  283. for t in layer_tips {
  284. parents[i] = *t;
  285. i += 1;
  286. if i >= N_EVENT_PARENTS {
  287. break 'outer
  288. }
  289. }
  290. }
  291. let layer =
  292. tips.last_key_value().and_then(|(layer, _)| layer.checked_add(1)).unwrap_or(u64::MAX);
  293. (layer, parents)
  294. }
  295. /// Storage layer for all rotating DAGs.
  296. pub struct DagStore {
  297. db: sled::Db,
  298. dags: BTreeMap<u64, DagSlot>,
  299. }
  300. impl DagStore {
  301. /// Create or open DAG slots.
  302. ///
  303. /// * **Bounded mode** (`max_dags = Some(n)`): create a rolling
  304. /// window of the most recent `n` DAGs. Old trees already in
  305. /// sled outside this window are left untouched (they're just
  306. /// not loaded into memory).
  307. /// * **Archive mode** (`max_dags = None`): discover *all*
  308. /// existing DAG trees in sled and load them, plus ensure the
  309. /// recent window exists. Nothing is ever dropped.
  310. pub async fn new(sled_db: sled::Db, config: &EventGraphConfig) -> Result<Self> {
  311. config.validate()?;
  312. let mut dags = BTreeMap::new();
  313. if config.hours_rotation == 0 {
  314. let genesis = generate_genesis(config);
  315. dags.insert(genesis.header.timestamp, Self::create_slot(&sled_db, &genesis).await);
  316. return Ok(Self { db: sled_db, dags })
  317. }
  318. // Determine how many recent DAGs to create/ensure exist.
  319. let window = config.max_dags.unwrap_or(24);
  320. // In archive mode, first discover and load any existing DAG
  321. // trees that are already in sled from previous runs.
  322. //
  323. // A DAG is stored across two trees: `<timestamp>` for events
  324. // and `headers_<timestamp>` for headers. We walk every tree
  325. // name in sled and pick out the ones whose name is a valid u64
  326. // timestamp.
  327. if config.max_dags.is_none() {
  328. for name in sled_db.tree_names() {
  329. let name_str = String::from_utf8_lossy(&name);
  330. if let Ok(ts) = name_str.parse::<u64>() {
  331. // Reconstruct the genesis for this timestamp
  332. let hdr = Header {
  333. timestamp: ts,
  334. parents: NULL_PARENTS,
  335. layer: 0,
  336. content_hash: blake3::hash(&config.genesis_contents),
  337. };
  338. let genesis = Event { header: hdr, content: config.genesis_contents.clone() };
  339. let slot = Self::create_slot(&sled_db, &genesis).await;
  340. dags.insert(ts, slot);
  341. }
  342. }
  343. }
  344. // Ensure the recent window of DAGs exists.
  345. // Creates them if they're not already loaded from the discovery step.
  346. for i in 1..=window {
  347. let ts = next_hour_timestamp((i as i64) - (window as i64));
  348. if dags.contains_key(&ts) {
  349. // Already loaded from sled discovery
  350. continue
  351. }
  352. let hdr = Header {
  353. timestamp: ts,
  354. parents: NULL_PARENTS,
  355. layer: 0,
  356. content_hash: blake3::hash(&config.genesis_contents),
  357. };
  358. let genesis = Event { header: hdr, content: config.genesis_contents.clone() };
  359. dags.insert(ts, Self::create_slot(&sled_db, &genesis).await);
  360. }
  361. Ok(Self { db: sled_db, dags })
  362. }
  363. async fn create_slot(db: &sled::Db, genesis: &Event) -> DagSlot {
  364. let name = genesis.header.timestamp.to_string();
  365. let ht = db.open_tree(format!("headers_{name}")).unwrap();
  366. let mt = db.open_tree(&name).unwrap();
  367. for (tree, data) in
  368. [(&ht, serialize_async(&genesis.header).await), (&mt, serialize_async(genesis).await)]
  369. {
  370. if tree.is_empty() {
  371. let mut ov = SledTreeOverlay::new(tree);
  372. ov.insert(genesis.id().as_bytes(), &data).unwrap();
  373. if let Some(b) = ov.aggregate() {
  374. tree.apply_batch(b).unwrap();
  375. }
  376. }
  377. }
  378. DagSlot {
  379. tips: compute_unreferenced_tips(&mt).await,
  380. time_index: TimeIndex::from_header_dag(&ht).await,
  381. header_tree: ht,
  382. main_tree: mt,
  383. }
  384. }
  385. /// Add a new DAG on rotation. In bounded mode, drops the oldest DAG
  386. /// when the limit is reached. In archive mode, never drops.
  387. pub async fn add_dag(&mut self, genesis: &Event, max_dags: Option<usize>) -> Result<()> {
  388. if let Some(limit) = max_dags {
  389. if limit == 0 {
  390. return Err(Error::Custom("event graph max_dags must be greater than 0".into()))
  391. }
  392. if self.dags.len() >= limit {
  393. let Some((_, old)) = self.dags.pop_first() else {
  394. return Err(Error::Custom("event graph DAG store is empty".into()))
  395. };
  396. self.db.drop_tree(old.header_tree.name()).unwrap();
  397. self.db.drop_tree(old.main_tree.name()).unwrap();
  398. }
  399. }
  400. let slot = Self::create_slot(&self.db, genesis).await;
  401. self.dags.insert(genesis.header.timestamp, slot);
  402. Ok(())
  403. }
  404. pub fn get_slot(&self, ts: &u64) -> Option<&DagSlot> {
  405. self.dags.get(ts)
  406. }
  407. pub fn get_slot_mut(&mut self, ts: &u64) -> Option<&mut DagSlot> {
  408. self.dags.get_mut(ts)
  409. }
  410. pub fn get_header_tree(&self, dag_name: &str) -> sled::Tree {
  411. self.db.open_tree(format!("headers_{dag_name}")).unwrap()
  412. }
  413. pub fn dag_timestamps(&self) -> Vec<u64> {
  414. self.dags.keys().cloned().collect()
  415. }
  416. }
  417. enum PeerStatus {
  418. Free,
  419. Busy,
  420. Failed,
  421. }
  422. /// Match an `EventRep` against the exact IDs requested for one sync chunk.
  423. ///
  424. /// The response may be partial, but every returned event must be unique and
  425. /// must belong to the outstanding request. Returned events and blobs are
  426. /// reordered to match the request order, and missing IDs are returned for
  427. /// retry with another peer.
  428. pub(crate) fn filter_requested_event_rep(
  429. requested: &[blake3::Hash],
  430. events: Vec<Event>,
  431. blobs: Vec<Vec<u8>>,
  432. ) -> Result<(Vec<Event>, Vec<Vec<u8>>, Vec<blake3::Hash>)> {
  433. if events.len() != blobs.len() {
  434. return Err(Error::DagSyncFailed)
  435. }
  436. let requested_set: HashSet<blake3::Hash> = requested.iter().copied().collect();
  437. let mut by_id = HashMap::with_capacity(events.len());
  438. for (event, blob) in events.into_iter().zip(blobs.into_iter()) {
  439. let event_id = event.id();
  440. if !requested_set.contains(&event_id) || by_id.insert(event_id, (event, blob)).is_some() {
  441. return Err(Error::DagSyncFailed)
  442. }
  443. }
  444. let mut matched_events = Vec::with_capacity(by_id.len());
  445. let mut matched_blobs = Vec::with_capacity(by_id.len());
  446. let mut missing = Vec::new();
  447. for id in requested {
  448. if let Some((event, blob)) = by_id.remove(id) {
  449. matched_events.push(event);
  450. matched_blobs.push(blob);
  451. } else {
  452. missing.push(*id);
  453. }
  454. }
  455. Ok((matched_events, matched_blobs, missing))
  456. }
  457. /// Merge one static-sync `EventRep` into the current batch state.
  458. ///
  459. /// Returns the number of still-pending requested IDs satisfied by this
  460. /// response. Invalid responses are rejected before any state is mutated.
  461. pub(crate) fn merge_static_sync_event_rep(
  462. requested: &[blake3::Hash],
  463. pending: &mut HashSet<blake3::Hash>,
  464. known: &mut HashSet<blake3::Hash>,
  465. want: &mut HashSet<blake3::Hash>,
  466. fetched: &mut Vec<(Event, Vec<u8>)>,
  467. events: Vec<Event>,
  468. blobs: Vec<Vec<u8>>,
  469. ) -> Result<usize> {
  470. let (matched_events, matched_blobs, _) = filter_requested_event_rep(requested, events, blobs)?;
  471. let mut matched = 0;
  472. for (ev, blob) in matched_events.into_iter().zip(matched_blobs) {
  473. let eid = ev.id();
  474. if !pending.remove(&eid) {
  475. continue
  476. }
  477. matched += 1;
  478. if known.insert(eid) {
  479. for p in ev.header.parents.iter() {
  480. if *p != NULL_ID && !known.contains(p) {
  481. want.insert(*p);
  482. }
  483. }
  484. fetched.push((ev, blob));
  485. }
  486. }
  487. Ok(matched)
  488. }
  489. /// The main Event Graph instance.
  490. ///
  491. /// Manages a rolling window of DAGs (one per rotation period), a
  492. /// static DAG for long-lived state (RLN identities), and the P2P
  493. /// protocol for syncing with peers.
  494. ///
  495. /// # Sync model
  496. ///
  497. /// Headers are synced eagerly (complete DAG skeleton in seconds).
  498. /// Event content is fetched lazily in the direction the application
  499. /// needs.
  500. ///
  501. /// The [`TimeIndex`] in each [`DagSlot`] enables O(log n)
  502. /// bidirectional pagination that crosses DAG boundaries
  503. /// transparently - the caller sees a flat chronological stream.
  504. pub struct EventGraph {
  505. pub(crate) p2p: P2pPtr,
  506. pub(crate) dag_store: RwLock<DagStore>,
  507. /// Side-table mapping `event_id -> original RLN signal blob` for
  508. /// rotating-DAG events. Mirror of [`Self::static_dag_blobs`] but
  509. /// for the rotating DAGs.
  510. ///
  511. /// Populated by `handle_event_put` after successful RLN
  512. /// verification, and by `dag_insert_with_blobs` during sync when
  513. /// the serving peer included the blob in its `EventRep`. Read
  514. /// by `handle_event_req` to forward blobs to syncing peers.
  515. /// Pruned by `dag_prune` when the corresponding rotating DAG
  516. /// rolls out of the retention window.
  517. pub(crate) dag_blobs: sled::Tree,
  518. /// Historical SMT roots, in canonical apply order.
  519. ///
  520. /// Key: `(layer:u64_be, event_id:32) = 40 bytes`. Value:
  521. /// `(root:32, timestamp:u64_be:8) = 40 bytes`.
  522. ///
  523. /// Big-endian layer encoding makes lexicographic byte order
  524. /// match canonical apply order, so `Tree::range` iterates
  525. /// chronologically and `Tree::get_lt` / `get_gt` give cheap
  526. /// neighbor lookups (used to find the timestamp interval during
  527. /// which a given root was the live root).
  528. ///
  529. /// See [`Self::apply_rln_static_event`] for the canonical-order
  530. /// rationale and [`Self::is_root_valid_at`] for how this is
  531. /// consulted during signal verification.
  532. pub(crate) rln_historical_roots_ordered: sled::Tree,
  533. /// Reverse index: `(root:32, ordered_key:40) -> []`.
  534. ///
  535. /// A root can appear more than once when static events are no-ops
  536. /// (duplicate registrations, idempotent slashes), so the value
  537. /// index stores every canonical occurrence rather than a single
  538. /// root-to-key mapping. [`Self::is_root_valid_at`] scans this
  539. /// prefix and accepts if any interval for the root matches.
  540. pub(crate) rln_historical_roots_by_value: sled::Tree,
  541. pub(crate) static_dag: sled::Tree,
  542. /// Side-table mapping `event_id -> original RLN blob` for static
  543. /// events. Used by [`Self::static_sync`] to re-verify the ZK
  544. /// proof of historical events at sync time. Every static-DAG
  545. /// event MUST have a corresponding entry - `static_sync` rejects
  546. /// events whose blob isn't available rather than falling through.
  547. pub(crate) static_dag_blobs: sled::Tree,
  548. datastore: PathBuf,
  549. replay_mode: bool,
  550. pub(crate) broadcasted_ids: RwLock<HashSet<blake3::Hash>>,
  551. pub prune_task: OnceCell<StoppableTaskPtr>,
  552. pub event_pub: PublisherPtr<Event>,
  553. pub static_pub: PublisherPtr<Event>,
  554. pub current_genesis: RwLock<Event>,
  555. pub config: EventGraphConfig,
  556. /// Decoded app-provided pregenerated RLN commitments.
  557. pregenerated_identity_commitments: Vec<pallas::Base>,
  558. /// Canonical byte representations for fast admission checks.
  559. pregenerated_identity_commitment_reprs: HashSet<[u8; 32]>,
  560. pub synced: AtomicBool,
  561. pub deg_enabled: AtomicBool,
  562. deg_publisher: PublisherPtr<DegEvent>,
  563. pub sled_db: sled::Db,
  564. pub zk_keys: Arc<ZkKeys>,
  565. pub identity_state: RwLock<IdentityState>,
  566. pub rln_state: RwLock<RlnState>,
  567. /// App identifier mixed into the RLN external nullifier. Derived
  568. /// from `config.genesis_contents` so two deployments using the
  569. /// same circuit cannot collide on internal_nullifiers.
  570. rln_app_id: rln::RlnAppId,
  571. }
  572. impl EventGraph {
  573. /// Create a new Event Graph.
  574. pub async fn new(
  575. p2p: P2pPtr,
  576. sled_db: sled::Db,
  577. datastore: PathBuf,
  578. replay_mode: bool,
  579. config: EventGraphConfig,
  580. ex: Arc<Executor<'_>>,
  581. ) -> Result<EventGraphPtr> {
  582. config.validate()?;
  583. let zk_keys = Arc::new(ZkKeys::build_and_load(&sled_db)?);
  584. Self::with_zk_keys(p2p, sled_db, datastore, replay_mode, config, zk_keys, ex).await
  585. }
  586. /// Same as [`Self::new`] but accepts a pre-built [`ZkKeys`].
  587. ///
  588. /// Production always wants `Self::new`, which builds keys once
  589. /// against its own sled DB. Tests use this variant to share a
  590. /// single [`Arc<ZkKeys>`] across many `EventGraph` instances -
  591. /// proving keys are large (hundreds of MB each) and copying
  592. /// them per-test would blow out RAM and `/dev/shm`.
  593. pub async fn with_zk_keys(
  594. p2p: P2pPtr,
  595. sled_db: sled::Db,
  596. datastore: PathBuf,
  597. replay_mode: bool,
  598. config: EventGraphConfig,
  599. zk_keys: Arc<ZkKeys>,
  600. ex: Arc<Executor<'_>>,
  601. ) -> Result<EventGraphPtr> {
  602. config.validate()?;
  603. let identity_state = IdentityState::new(&sled_db)?;
  604. let rln_app_id = rln::RlnAppId::from_genesis(&config.genesis_contents);
  605. let current_genesis = generate_genesis(&config);
  606. let (pregenerated_identity_commitments, pregenerated_identity_commitment_reprs) =
  607. validate_pregenerated_identity_commitments(&config)?;
  608. let dag_store = DagStore::new(sled_db.clone(), &config).await?;
  609. let static_dag = Self::static_new(&sled_db, &config).await?;
  610. let static_dag_blobs = sled_db.open_tree("static-dag-blobs")?;
  611. let dag_blobs = sled_db.open_tree("dag-blobs")?;
  612. // Historical-roots side-tables. See the design comment on
  613. // `EventGraph::apply_rln_static_event` for the full rationale.
  614. // In short: every static-DAG mutation produces a new SMT root,
  615. // and we need to recognize *any* historical root for sync-time
  616. // signal verification, not just the most recent N. The
  617. // `ordered` tree gives us canonical replay (and successor
  618. // lookup for the time-window check), the `by_value` tree
  619. // gives us O(log n) "is this root historical?" queries.
  620. let rln_historical_roots_ordered = sled_db.open_tree("rln-historical-roots-ordered")?;
  621. let rln_historical_roots_by_value = sled_db.open_tree("rln-historical-roots-by-value")?;
  622. // Check whether the current genesis event is already in the
  623. // store. If not, we need to prune (create a fresh slot).
  624. let dag_ts = current_genesis.header.timestamp;
  625. let need_prune = dag_store
  626. .get_slot(&dag_ts)
  627. .map(|s| !s.main_tree.contains_key(current_genesis.id().as_bytes()).unwrap_or(false))
  628. .unwrap_or(true);
  629. let self_ = Arc::new(Self {
  630. p2p,
  631. sled_db: sled_db.clone(),
  632. dag_store: RwLock::new(dag_store),
  633. static_dag,
  634. static_dag_blobs,
  635. dag_blobs,
  636. rln_historical_roots_ordered,
  637. rln_historical_roots_by_value,
  638. datastore,
  639. replay_mode,
  640. broadcasted_ids: RwLock::new(HashSet::new()),
  641. prune_task: OnceCell::new(),
  642. event_pub: Publisher::new(),
  643. static_pub: Publisher::new(),
  644. current_genesis: RwLock::new(current_genesis.clone()),
  645. config: config.clone(),
  646. pregenerated_identity_commitments,
  647. pregenerated_identity_commitment_reprs,
  648. synced: AtomicBool::new(false),
  649. deg_enabled: AtomicBool::new(false),
  650. deg_publisher: Publisher::new(),
  651. zk_keys,
  652. identity_state: RwLock::new(identity_state),
  653. rln_state: RwLock::new(RlnState::new()),
  654. rln_app_id,
  655. });
  656. if need_prune {
  657. info!(
  658. target: "event_graph::new",
  659. "[EVENTGRAPH] Pruning: current genesis not found",
  660. );
  661. self_.dag_prune(current_genesis).await?;
  662. }
  663. // Reconcile persisted RLN state before bootstrapping. If an
  664. // earlier process crashed after writing identity leaves but before
  665. // inserting the corresponding static event, bootstrapping must see
  666. // the corrected leaf set rather than skip the configured identity.
  667. self_.rebuild_historical_roots_if_needed().await?;
  668. // Init genesis registration events after recovery has made the
  669. // static DAG authoritative for the current identity tree.
  670. if config.hours_rotation > 0 {
  671. self_.bootstrap_genesis_identities().await?;
  672. }
  673. self_.audit_static_blobs().await?;
  674. if config.hours_rotation > 0 {
  675. let task = StoppableTask::new();
  676. let _ = self_.prune_task.set(task.clone()).await;
  677. task.clone().start(
  678. self_.clone().dag_prune_task(),
  679. |res| async move {
  680. if let Err(e) = res {
  681. if !matches!(e, Error::DetachedTaskStopped) {
  682. error!("Prune: {e}");
  683. }
  684. }
  685. },
  686. Error::DetachedTaskStopped,
  687. ex,
  688. );
  689. }
  690. Ok(self_)
  691. }
  692. /// Rebuild the RLN state side-tables from the static DAG.
  693. ///
  694. /// Called once at startup. No-op if the historical-root indexes match
  695. /// the canonical static-DAG event sequence and the persisted identity
  696. /// leaves match the commitment set obtained by replaying that sequence.
  697. /// Otherwise resets the identity SMT and root indexes, then replays every
  698. /// parseable static-DAG event in canonical `(layer, event_id)` order.
  699. ///
  700. /// **Side effect.** The static DAG is authoritative. The in-memory SMT,
  701. /// the persistent `rln-identity-leaves` tree, and both historical-root
  702. /// indexes are derived from it so crashes between the old split write
  703. /// steps cannot leave stale leaves or unusable root indexes behind.
  704. async fn rebuild_historical_roots_if_needed(self: &Arc<Self>) -> Result<()> {
  705. let mut events: Vec<(Event, rln::RLNNode)> = vec![];
  706. for item in self.static_dag.iter() {
  707. let (_, val) = item?;
  708. let ev: Event = deserialize_async(&val).await?;
  709. if ev.header.parents == NULL_PARENTS {
  710. continue
  711. }
  712. let Ok((node, _)) = deserialize_async_partial::<rln::RLNNode>(ev.content()).await
  713. else {
  714. continue
  715. };
  716. events.push((ev, node));
  717. }
  718. events.sort_by(|(a, _), (b, _)| {
  719. a.header
  720. .layer
  721. .cmp(&b.header.layer)
  722. .then_with(|| a.id().as_bytes().cmp(b.id().as_bytes()))
  723. });
  724. let mut expected_commitments = BTreeSet::new();
  725. for (_, node) in &events {
  726. match node {
  727. rln::RLNNode::Registration(commitment) => {
  728. expected_commitments.insert(commitment.to_repr());
  729. }
  730. rln::RLNNode::Slashing(commitment) => {
  731. expected_commitments.remove(&commitment.to_repr());
  732. }
  733. }
  734. }
  735. let expected_leaves = expected_commitments.len();
  736. let actual_commitments = self.identity_state.read().await.commitment_reprs();
  737. let (actual_leaves, leaves_consistent) = match actual_commitments {
  738. Ok(commitments) => {
  739. let len = commitments.len();
  740. (len, commitments == expected_commitments)
  741. }
  742. Err(e) => {
  743. warn!(
  744. target: "event_graph::new",
  745. "[EVENTGRAPH] RLN identity leaf audit failed: {e}; rebuilding",
  746. );
  747. (0, false)
  748. }
  749. };
  750. let static_count = events.len();
  751. let historical_roots_consistent = self.historical_roots_index_consistent(static_count)?;
  752. let recorded_count = self.rln_historical_roots_ordered.len();
  753. let by_value_count = self.rln_historical_roots_by_value.len();
  754. let consistent = historical_roots_consistent && leaves_consistent;
  755. info!(
  756. target: "event_graph::new",
  757. concat!(
  758. "[EVENTGRAPH] RLN state audit: static_count={} recorded_count={} ",
  759. "by_value_count={} actual_leaves={} expected_leaves={} consistent={}",
  760. ),
  761. static_count, recorded_count, by_value_count, actual_leaves, expected_leaves,
  762. consistent,
  763. );
  764. if consistent {
  765. return Ok(())
  766. }
  767. info!(
  768. target: "event_graph::new",
  769. concat!(
  770. "[EVENTGRAPH] Rebuilding RLN state: {} static events, {} recorded roots, ",
  771. "{} by-value roots, {} leaves (expected {})",
  772. ),
  773. static_count, recorded_count, by_value_count, actual_leaves, expected_leaves,
  774. );
  775. self.rln_historical_roots_ordered.clear()?;
  776. self.rln_historical_roots_by_value.clear()?;
  777. {
  778. let mut state = self.identity_state.write().await;
  779. state.clear_for_rebuild()?;
  780. }
  781. for (ev, rln_node) in events {
  782. let _ = self.apply_rln_static_event(&ev, &rln_node).await?;
  783. }
  784. info!(
  785. target: "event_graph::new",
  786. "[EVENTGRAPH] RLN state rebuild complete",
  787. );
  788. Ok(())
  789. }
  790. fn historical_roots_index_consistent(&self, expected_count: usize) -> Result<bool> {
  791. if self.rln_historical_roots_ordered.len() != expected_count {
  792. return Ok(false)
  793. }
  794. if self.rln_historical_roots_by_value.len() != expected_count {
  795. return Ok(false)
  796. }
  797. for item in self.rln_historical_roots_ordered.iter() {
  798. let (ordered_key_bytes, value_bytes) = item?;
  799. if ordered_key_bytes.len() != 40 {
  800. return Ok(false)
  801. }
  802. let Ok((root, _)) = decode_historical_root_value(&value_bytes) else {
  803. return Ok(false)
  804. };
  805. let mut ordered_key = [0u8; 40];
  806. ordered_key.copy_from_slice(&ordered_key_bytes);
  807. let by_value_key = encode_historical_root_by_value_key(&root, &ordered_key);
  808. if !self.rln_historical_roots_by_value.contains_key(by_value_key)? {
  809. return Ok(false)
  810. }
  811. }
  812. Ok(true)
  813. }
  814. /// After header sync, event content can be fetched lazily via
  815. /// [`fetch_page`] or peer [`RangeReq`] - the application pulls
  816. /// the events it actually wants to display or process, without
  817. /// downloading the entire content on every sync.
  818. pub async fn dag_sync_headers(&self, dag_ts: u64) -> Result<()> {
  819. self.sync_impl(dag_ts, false).await
  820. }
  821. /// Full sync: headers plus all event content currently in the DAG.
  822. ///
  823. /// Use this when the application wants the complete historical
  824. /// content (e.g. an archive node, or a node rebuilding local state
  825. /// from the full event stream).
  826. pub async fn dag_sync(&self, dag_ts: u64) -> Result<()> {
  827. self.sync_impl(dag_ts, true).await
  828. }
  829. async fn sync_impl(&self, dag_ts: u64, fetch_content: bool) -> Result<()> {
  830. let dag_name = dag_ts.to_string();
  831. let channels = self.p2p.hosts().peers();
  832. // We need at least one peer to ask
  833. if channels.is_empty() {
  834. return Err(Error::DagSyncFailed)
  835. }
  836. let timeout = self.p2p.settings().read().await.outbound_connect_timeout_max();
  837. // Parallel tip collection
  838. let mut futs = FuturesUnordered::new();
  839. for ch in channels.iter() {
  840. futs.push(request_tips(ch, dag_name.clone(), timeout));
  841. }
  842. let mut tips: HashMap<blake3::Hash, (u64, usize)> = HashMap::new();
  843. let mut responded = 0usize;
  844. while let Some(res) = futs.next().await {
  845. if let Ok(peer_tips) = res {
  846. responded += 1;
  847. for (layer, hashes) in &peer_tips {
  848. for h in hashes {
  849. tips.entry(*h).and_modify(|e| e.1 += 1).or_insert((*layer, 1));
  850. }
  851. }
  852. }
  853. }
  854. if tips.is_empty() {
  855. return Err(Error::DagSyncFailed)
  856. }
  857. // 2/3 quorum
  858. let threshold = (responded * 2).div_ceil(3);
  859. let accepted: HashSet<blake3::Hash> = tips
  860. .iter()
  861. .filter(|(h, (_, n))| **h != NULL_ID && *n >= threshold)
  862. .map(|(h, _)| *h)
  863. .collect();
  864. let store = self.dag_store.read().await;
  865. let slot = store.get_slot(&dag_ts).ok_or(Error::DagSyncFailed)?;
  866. let missing: HashSet<blake3::Hash> = accepted
  867. .iter()
  868. .filter(|h| !slot.main_tree.contains_key(h.as_bytes()).unwrap_or(true))
  869. .cloned()
  870. .collect();
  871. if missing.is_empty() {
  872. return Ok(())
  873. }
  874. let our_tips = slot.tips.clone();
  875. drop(store);
  876. // Parallel header sync
  877. let mut hfuts = FuturesUnordered::new();
  878. for ch in channels.iter() {
  879. hfuts.push(request_header(ch, dag_name.clone(), our_tips.clone(), timeout));
  880. }
  881. while let Some(res) = hfuts.next().await {
  882. if let Ok(hdrs) = res {
  883. self.header_dag_insert(hdrs, &dag_name).await?;
  884. }
  885. }
  886. if fetch_content {
  887. self.fetch_missing_events(dag_ts, &dag_name, timeout).await?;
  888. }
  889. Ok(())
  890. }
  891. async fn fetch_missing_events(&self, dag_ts: u64, dag_name: &str, timeout: u64) -> Result<()> {
  892. let store = self.dag_store.read().await;
  893. let slot = store.get_slot(&dag_ts).ok_or(Error::DagSyncFailed)?;
  894. let mut sorted = vec![];
  895. for item in slot.header_tree.iter() {
  896. let (hb, val) = item.unwrap();
  897. let hdr: Header = deserialize_async(&val).await.unwrap();
  898. if hdr.parents != NULL_PARENTS && !slot.main_tree.contains_key(hb)? {
  899. sorted.push(hdr);
  900. }
  901. }
  902. sorted.sort_by_key(|h| h.layer);
  903. drop(store);
  904. if sorted.is_empty() {
  905. return Ok(())
  906. }
  907. let batch = 20;
  908. let mut chunks: BTreeMap<usize, Vec<blake3::Hash>> = BTreeMap::new();
  909. for (i, c) in sorted.chunks(batch).enumerate() {
  910. chunks.insert(i, c.iter().map(|h| h.id()).collect());
  911. }
  912. let mut remaining: BTreeSet<usize> = chunks.keys().cloned().collect();
  913. let mut peer_st: HashMap<Url, PeerStatus> = HashMap::new();
  914. let mut count = 0;
  915. let mut fs = FuturesUnordered::new();
  916. // Collected by event ID so partial chunk retries cannot disturb
  917. // the final layer-sorted insertion order. Empty `Vec<u8>` entries
  918. // mean "this event has no blob from the serving peer".
  919. let mut received: HashMap<blake3::Hash, (Event, Vec<u8>)> = HashMap::new();
  920. while count < sorted.len() {
  921. let mut free = vec![];
  922. let mut busy = 0;
  923. self.p2p.hosts().peers().iter().for_each(|ch| match peer_st.get(ch.address()) {
  924. Some(PeerStatus::Free) | None => {
  925. free.push(ch.clone());
  926. }
  927. Some(PeerStatus::Busy) => {
  928. busy += 1;
  929. }
  930. _ => {}
  931. });
  932. if free.is_empty() && busy == 0 {
  933. return Err(Error::DagSyncFailed)
  934. }
  935. if remaining.is_empty() && fs.is_empty() {
  936. return Err(Error::DagSyncFailed)
  937. }
  938. let n = std::cmp::min(free.len(), remaining.len());
  939. let ids: Vec<usize> = remaining.iter().take(n).copied().collect();
  940. for (i, cid) in ids.iter().enumerate() {
  941. fs.push(request_event(free[i].clone(), chunks[cid].clone(), *cid, timeout));
  942. remaining.remove(cid);
  943. peer_st.insert(free[i].address().clone(), PeerStatus::Busy);
  944. }
  945. if let Some((evts, cid, ch)) = fs.next().await {
  946. if let Ok((events, blobs)) = evts {
  947. let Some(requested) = chunks.get(&cid) else {
  948. peer_st.insert(ch.address().clone(), PeerStatus::Failed);
  949. continue
  950. };
  951. match filter_requested_event_rep(requested, events, blobs) {
  952. Ok((matched_events, matched_blobs, missing)) => {
  953. let matched = matched_events.len();
  954. for (event, blob) in matched_events.into_iter().zip(matched_blobs) {
  955. let event_id = event.id();
  956. if received.insert(event_id, (event, blob)).is_none() {
  957. count += 1;
  958. }
  959. }
  960. if missing.is_empty() {
  961. peer_st.insert(ch.address().clone(), PeerStatus::Free);
  962. } else {
  963. chunks.insert(cid, missing);
  964. remaining.insert(cid);
  965. let status = if matched == 0 {
  966. PeerStatus::Failed
  967. } else {
  968. PeerStatus::Free
  969. };
  970. peer_st.insert(ch.address().clone(), status);
  971. }
  972. }
  973. Err(_) => {
  974. remaining.insert(cid);
  975. peer_st.insert(ch.address().clone(), PeerStatus::Failed);
  976. }
  977. }
  978. } else {
  979. remaining.insert(cid);
  980. peer_st.insert(ch.address().clone(), PeerStatus::Failed);
  981. }
  982. }
  983. }
  984. let mut events = Vec::with_capacity(sorted.len());
  985. let mut blobs = Vec::with_capacity(sorted.len());
  986. for hdr in sorted {
  987. let event_id = hdr.id();
  988. let Some((event, blob)) = received.remove(&event_id) else {
  989. return Err(Error::DagSyncFailed)
  990. };
  991. events.push(event);
  992. blobs.push(blob);
  993. }
  994. // dag_insert_with_blobs handles RLN re-verification per event when
  995. // blobs are present, and falls through to the trust-the-quorum path
  996. // when they're not. See dag_insert_with_blobs's docstring for the
  997. // policy.
  998. self.dag_insert_with_blobs(&events, &blobs, dag_name).await?;
  999. Ok(())
  1000. }
  1001. /// Sync the `count` most recent DAGs (full content).
  1002. ///
  1003. /// Iterates oldest-first so that later syncs build on earlier
  1004. /// ones (parent events exist before children reference them).
  1005. pub async fn sync_selected(&self, count: usize) -> Result<()> {
  1006. let ts: Vec<u64> =
  1007. self.dag_store.read().await.dag_timestamps().into_iter().rev().take(count).collect();
  1008. for t in ts.into_iter().rev() {
  1009. self.dag_sync(t).await?;
  1010. }
  1011. self.synced.store(true, Ordering::Release);
  1012. Ok(())
  1013. }
  1014. /// Sync only headers for the `count` most recent DAGs.
  1015. ///
  1016. /// Fast variant - gives a full DAG skeleton without downloading
  1017. /// event bodies. Pair with [`fetch_page`] to pull content on-demand.
  1018. pub async fn sync_selected_headers(&self, count: usize) -> Result<()> {
  1019. let ts: Vec<u64> =
  1020. self.dag_store.read().await.dag_timestamps().into_iter().rev().take(count).collect();
  1021. for t in ts.into_iter().rev() {
  1022. self.dag_sync_headers(t).await?;
  1023. }
  1024. self.synced.store(true, Ordering::Release);
  1025. Ok(())
  1026. }
  1027. /// Sync the static DAG from peers.
  1028. ///
  1029. /// The static DAG holds RLN identity events (registrations and
  1030. /// slashes). It is *persistent* across rotation windows - unlike
  1031. /// rotating DAGs, events are never pruned - and has no separate
  1032. /// `header_tree`, so it uses a different sync strategy:
  1033. ///
  1034. /// 1. Ask every peer for their `"static-dag"` tips.
  1035. /// 2. Take the tips that reach a 2/3 quorum.
  1036. /// 3. BFS-fetch the events and their ancestors directly via
  1037. /// `EventReq` until the entire reachable subgraph is local.
  1038. ///
  1039. /// Peers serve static-DAG event requests even when the IDs are
  1040. /// not in their `broadcasted_ids` set (see the relaxation in
  1041. /// `handle_event_req`), because static-DAG state is public
  1042. /// consensus information. Registration-event proof verification,
  1043. /// duplicate detection, and commitment-tree updates are all done
  1044. /// through the normal `StaticPut` ingestion path
  1045. /// (`handle_static_put`) - but `static_sync` uses direct-insert
  1046. /// via [`Self::static_insert`] plus on-the-fly identity-state
  1047. /// application, because we're catching up rather than processing
  1048. /// broadcasts.
  1049. ///
  1050. /// Note: for security, this method ONLY applies events whose
  1051. /// blob/RLN verification passes. We do not trust peers blindly
  1052. /// on historical state - proofs are re-verified locally for
  1053. /// every single event before its effect is merged into the
  1054. /// identity tree. This is the same discipline `handle_static_put`
  1055. /// uses; see [`Self::rln_verify_static_event`].
  1056. pub async fn static_sync(&self) -> Result<()> {
  1057. static DAG_NAME: &str = "static-dag";
  1058. let channels = self.p2p.hosts().peers();
  1059. if channels.is_empty() {
  1060. return Err(Error::DagSyncFailed)
  1061. }
  1062. let timeout = self.p2p.settings().read().await.outbound_connect_timeout_max();
  1063. // Step 1: gather tips from every peer in parallel.
  1064. let mut tip_futs = FuturesUnordered::new();
  1065. for ch in channels.iter() {
  1066. tip_futs.push(request_tips(ch, DAG_NAME.to_string(), timeout));
  1067. }
  1068. let mut tip_counts: HashMap<blake3::Hash, usize> = HashMap::new();
  1069. let mut responded = 0usize;
  1070. while let Some(res) = tip_futs.next().await {
  1071. if let Ok(peer_tips) = res {
  1072. responded += 1;
  1073. for hashes in peer_tips.values() {
  1074. for h in hashes {
  1075. *tip_counts.entry(*h).or_insert(0) += 1;
  1076. }
  1077. }
  1078. }
  1079. }
  1080. // If no peer answered we have nothing to do. An empty
  1081. // network-side static DAG is a valid state (brand new app
  1082. // deployment), so we return Ok rather than error.
  1083. if responded == 0 {
  1084. info!(
  1085. target: "event_graph::static_sync",
  1086. "[STATIC_SYNC] no peer responded to TipReq; nothing to sync"
  1087. );
  1088. return Ok(())
  1089. }
  1090. // Step 2: take tips at 2/3 quorum. This matches the
  1091. // threshold used in `sync_impl`.
  1092. let threshold = (responded * 2).div_ceil(3);
  1093. let total_distinct_tips = tip_counts.len();
  1094. let tip_ids: HashSet<blake3::Hash> = tip_counts
  1095. .into_iter()
  1096. .filter(|(h, n)| *h != NULL_ID && *n >= threshold)
  1097. .map(|(h, _)| h)
  1098. .collect();
  1099. // What's already local?
  1100. let mut known: HashSet<blake3::Hash> = HashSet::new();
  1101. for item in self.static_dag.iter() {
  1102. let (k, _) = item?;
  1103. if let Ok(bytes) = <[u8; 32]>::try_from(&k as &[u8]) {
  1104. known.insert(blake3::Hash::from_bytes(bytes));
  1105. }
  1106. }
  1107. info!(
  1108. target: "event_graph::static_sync",
  1109. "[STATIC_SYNC] peers_responded={} threshold={} distinct_tips_seen={} \
  1110. tip_ids_quorum={} known_local={}",
  1111. responded, threshold, total_distinct_tips, tip_ids.len(), known.len(),
  1112. );
  1113. // Step 3: BFS from the quorum tips, fetching events we
  1114. // don't have. Any event we pull in may reference ancestors
  1115. // we ALSO don't have; enqueue them and keep going until the
  1116. // frontier is empty.
  1117. //
  1118. // Bounded at SYNC_MAX_STATIC_EVENTS (defined at module level)
  1119. // to defend against a malicious peer who serves a fabricated
  1120. // deep-ancestry chain. In practice static DAGs are small (one
  1121. // event per registration / slash), so this bound is
  1122. // comfortably above any real deployment's size.
  1123. let mut want: HashSet<blake3::Hash> = tip_ids.difference(&known).copied().collect();
  1124. // Events fetched during BFS, paired with their blobs (empty
  1125. // Vec if the peer didn't have the blob - see EventRep
  1126. // docstring). Index alignment is preserved through the
  1127. // entire pipeline up to the apply loop.
  1128. let mut fetched: Vec<(Event, Vec<u8>)> = vec![];
  1129. while !want.is_empty() {
  1130. want.retain(|id| !known.contains(id));
  1131. if want.is_empty() {
  1132. break
  1133. }
  1134. if fetched.len() >= SYNC_MAX_STATIC_EVENTS {
  1135. error!(
  1136. target: "event_graph::static_sync",
  1137. "[STATIC_SYNC] reached {} event cap; aborting",
  1138. SYNC_MAX_STATIC_EVENTS,
  1139. );
  1140. return Err(Error::DagSyncFailed)
  1141. }
  1142. let batch: Vec<blake3::Hash> = want.iter().take(MAX_EVENT_REQ_IDS).copied().collect();
  1143. for id in &batch {
  1144. want.remove(id);
  1145. }
  1146. let mut pending: HashSet<blake3::Hash> = batch.iter().copied().collect();
  1147. // Ask every peer for the same batch. We keep consuming
  1148. // responses until the batch is complete or every peer has
  1149. // failed to help. Irrelevant, duplicate, or blob-misaligned
  1150. // replies do not satisfy the request.
  1151. let mut req_futs = FuturesUnordered::new();
  1152. for (i, ch) in channels.iter().enumerate() {
  1153. req_futs.push(request_event(ch.clone(), batch.clone(), i, timeout));
  1154. }
  1155. let mut made_progress = false;
  1156. while !pending.is_empty() {
  1157. let Some((res, _, _)) = req_futs.next().await else { break };
  1158. let Ok((evs, blobs)) = res else { continue };
  1159. if evs.is_empty() {
  1160. continue
  1161. }
  1162. let Ok(matched) = merge_static_sync_event_rep(
  1163. &batch,
  1164. &mut pending,
  1165. &mut known,
  1166. &mut want,
  1167. &mut fetched,
  1168. evs,
  1169. blobs,
  1170. ) else {
  1171. continue
  1172. };
  1173. if matched > 0 {
  1174. made_progress = true;
  1175. }
  1176. }
  1177. if !pending.is_empty() {
  1178. want.extend(pending.iter().copied().filter(|id| !known.contains(id)));
  1179. }
  1180. want.retain(|id| !known.contains(id));
  1181. if !made_progress {
  1182. // Nobody responded usefully. Give up so we don't
  1183. // loop forever on an unreachable ancestor.
  1184. error!(
  1185. target: "event_graph::static_sync",
  1186. "[STATIC_SYNC] no peer served requested events; aborting",
  1187. );
  1188. return Err(Error::DagSyncFailed)
  1189. }
  1190. }
  1191. // Step 4: canonical-order the fetched events so all nodes
  1192. // produce the same intermediate SMT roots. Primary key:
  1193. // layer (matches DAG topology). Secondary key: event_id
  1194. // (32-byte hash, lexicographic byte order is total). Without
  1195. // the tie-breaker, two events at the same layer could be
  1196. // applied in different orders on different nodes, producing
  1197. // different intermediate roots and breaking sync-time signal
  1198. // verification. See the design comment on
  1199. // `apply_rln_static_event` for the full rationale.
  1200. fetched.sort_by(|(a, _), (b, _)| {
  1201. a.header
  1202. .layer
  1203. .cmp(&b.header.layer)
  1204. .then_with(|| a.id().as_bytes().cmp(b.id().as_bytes()))
  1205. });
  1206. // Track the apply-loop outcome for the summary log.
  1207. let mut applied = 0usize;
  1208. let mut already_present = 0usize;
  1209. let mut blob_missing = 0usize;
  1210. let mut rejected = 0usize;
  1211. let mut structural_invalid = 0usize;
  1212. let mut content_unparseable = 0usize;
  1213. let total_to_consider = fetched.len();
  1214. for (ev, blob) in fetched {
  1215. // Skip if someone else inserted it concurrently.
  1216. if self.static_dag.contains_key(ev.id().as_bytes())? {
  1217. already_present += 1;
  1218. continue
  1219. }
  1220. // Structural validation always runs. Static-DAG events
  1221. // are persistent and may be far older than the 60s drift
  1222. // window allowed by `validate_new`; use the static
  1223. // sibling that omits the freshness check while keeping
  1224. // the structural ones.
  1225. if !ev.validate_new_static() {
  1226. structural_invalid += 1;
  1227. continue
  1228. }
  1229. let rln_node: rln::RLNNode = match deserialize_async_partial(ev.content()).await {
  1230. Ok((v, _)) => v,
  1231. Err(_) => {
  1232. content_unparseable += 1;
  1233. continue
  1234. }
  1235. };
  1236. // RLN verification is mandatory. A non-genesis static
  1237. // event without a blob during sync is treated as
  1238. // misbehavior: either the serving peer is buggy or
  1239. // adversarial, or the originator never persisted the blob
  1240. // (which itself is a protocol violation). Skip with a
  1241. // loud log - we don't strike here because static_sync
  1242. // doesn't have a single peer to attribute the failure
  1243. // to (the quorum collected blobs from multiple peers).
  1244. if blob.is_empty() {
  1245. blob_missing += 1;
  1246. error!(
  1247. target: "event_graph::static_sync",
  1248. "[STATIC_SYNC] no blob available for static event {}; skipping. \
  1249. Every static-DAG event must carry an RLN blob.",
  1250. ev.id(),
  1251. );
  1252. continue
  1253. }
  1254. let outcome = self.rln_verify_static_event(&rln_node, &blob, ev.header.timestamp).await;
  1255. match outcome {
  1256. rln::StaticEventCheck::AcceptedRegistration(_) |
  1257. rln::StaticEventCheck::AcceptedSlash(_) => {
  1258. self.commit_verified_static_event(&ev, &blob, &rln_node).await?;
  1259. applied += 1;
  1260. }
  1261. rln::StaticEventCheck::Rejected | rln::StaticEventCheck::Malicious => {
  1262. // A historical event whose blob fails
  1263. // re-verification despite being held by the 2/3
  1264. // quorum is a serious finding - either the blob
  1265. // was tampered with, the quorum was compromised,
  1266. // or our verifying keys diverged. Log loudly and
  1267. // skip.
  1268. rejected += 1;
  1269. error!(
  1270. target: "event_graph::static_sync",
  1271. "[STATIC_SYNC] historical blob FAILED re-verification for event {}: {:?}; \
  1272. skipping event despite quorum inclusion",
  1273. ev.id(),
  1274. outcome,
  1275. );
  1276. }
  1277. }
  1278. }
  1279. info!(
  1280. target: "event_graph::static_sync",
  1281. "[STATIC_SYNC] complete: fetched={} applied={} already_present={} \
  1282. blob_missing={} verification_rejected={} structural_invalid={} \
  1283. unparseable={}",
  1284. total_to_consider, applied, already_present, blob_missing, rejected,
  1285. structural_invalid, content_unparseable,
  1286. );
  1287. Ok(())
  1288. }
  1289. /// Fetch a page of events, crossing DAG boundaries transparently.
  1290. pub async fn fetch_page(
  1291. &self,
  1292. cursor_ts: u64,
  1293. dir: SyncDirection,
  1294. limit: usize,
  1295. ) -> Result<Vec<Event>> {
  1296. let limit = limit.min(MAX_RANGE_PAGE_SIZE);
  1297. let mut out = vec![];
  1298. let store = self.dag_store.read().await;
  1299. let slots: Vec<_> = match dir {
  1300. SyncDirection::Forward => store.dags.iter().collect(),
  1301. SyncDirection::Backward => store.dags.iter().rev().collect(),
  1302. };
  1303. for (_, slot) in slots {
  1304. if out.len() >= limit {
  1305. break
  1306. }
  1307. let rem = limit - out.len();
  1308. let ids = match dir {
  1309. SyncDirection::Forward => slot.time_index.after(cursor_ts, rem),
  1310. SyncDirection::Backward => slot.time_index.before(cursor_ts, rem),
  1311. };
  1312. for id in ids {
  1313. if let Some(bytes) = slot.main_tree.get(id.as_bytes())? {
  1314. out.push(deserialize_async(&bytes).await?);
  1315. }
  1316. }
  1317. }
  1318. out.truncate(limit);
  1319. Ok(out)
  1320. }
  1321. async fn dag_prune(&self, genesis: Event) -> Result<()> {
  1322. let mut bcast = self.broadcasted_ids.write().await;
  1323. let mut cur = self.current_genesis.write().await;
  1324. // Before the DAG store evicts the oldest DAG (which would
  1325. // drop its main_tree), enumerate the about-to-be-dropped
  1326. // event IDs so we can remove their blobs from `dag_blobs`.
  1327. // Without this, blob entries would orphan and accumulate
  1328. // forever - the side-table is not bounded by the rotation
  1329. // window on its own.
  1330. if let Some(limit) = self.config.max_dags {
  1331. let store = self.dag_store.read().await;
  1332. if store.dags.len() >= limit {
  1333. if let Some((_, oldest)) = store.dags.iter().next() {
  1334. for item in oldest.main_tree.iter() {
  1335. let (eid, _) = match item {
  1336. Ok(v) => v,
  1337. Err(_) => continue,
  1338. };
  1339. let _ = self.dag_blobs.remove(&eid);
  1340. }
  1341. }
  1342. }
  1343. }
  1344. self.dag_store.write().await.add_dag(&genesis, self.config.max_dags).await?;
  1345. *cur = genesis;
  1346. *bcast = HashSet::new();
  1347. Ok(())
  1348. }
  1349. async fn dag_prune_task(self: Arc<Self>) -> Result<()> {
  1350. loop {
  1351. let next =
  1352. next_rotation_timestamp(self.config.initial_genesis, self.config.hours_rotation)?;
  1353. let hdr = Header {
  1354. timestamp: next,
  1355. parents: NULL_PARENTS,
  1356. layer: 0,
  1357. content_hash: blake3::hash(&self.config.genesis_contents),
  1358. };
  1359. let genesis = Event { header: hdr, content: self.config.genesis_contents.clone() };
  1360. msleep(millis_until_next_rotation(next)?).await;
  1361. self.dag_prune(genesis).await?;
  1362. }
  1363. }
  1364. /// Public insertion path for a rotating-DAG signal event.
  1365. ///
  1366. /// Non-genesis rotating events must carry an RLN signal blob. This method
  1367. /// inserts the header, re-verifies the blob, records RLN metadata, stores
  1368. /// the blob for future sync, and only then commits the event body. External
  1369. /// applications should use this instead of the unchecked post-verification
  1370. /// insertion path.
  1371. pub async fn insert_signal_with_blob(
  1372. &self,
  1373. event: &Event,
  1374. blob: &[u8],
  1375. dag_name: &str,
  1376. ) -> Result<Vec<blake3::Hash>> {
  1377. if event.header.parents != NULL_PARENTS && blob.is_empty() {
  1378. return Err(Error::Custom("rotating-DAG signal event blob must not be empty".into()))
  1379. }
  1380. let dag_ts = u64::from_str(dag_name)?;
  1381. let already_known = if event.header.parents == NULL_PARENTS {
  1382. false
  1383. } else {
  1384. let store = self.dag_store.read().await;
  1385. match store.get_slot(&dag_ts) {
  1386. Some(slot) => slot.main_tree.contains_key(event.id().as_bytes())?,
  1387. None => false,
  1388. }
  1389. };
  1390. self.header_dag_insert(vec![event.header.clone()], dag_name).await?;
  1391. let blobs = if blob.is_empty() { vec![] } else { vec![blob.to_vec()] };
  1392. let ids = self.dag_insert_with_blobs(std::slice::from_ref(event), &blobs, dag_name).await?;
  1393. let accepted = ids.contains(&event.id()) || already_known || {
  1394. let store = self.dag_store.read().await;
  1395. match store.get_slot(&dag_ts) {
  1396. Some(slot) => slot.main_tree.contains_key(event.id().as_bytes())?,
  1397. None => false,
  1398. }
  1399. };
  1400. if event.header.parents != NULL_PARENTS && !accepted {
  1401. return Err(Error::Custom("rotating-DAG signal event was not accepted".into()))
  1402. }
  1403. Ok(ids)
  1404. }
  1405. /// Insert events into a rotating DAG **without RLN verification**.
  1406. ///
  1407. /// This is the crate-internal post-verification entry point for callers
  1408. /// that have already verified the proof separately and recorded RLN
  1409. /// metadata. Public callers must use [`Self::insert_signal_with_blob`] or
  1410. /// [`Self::dag_insert_with_blobs`] so non-genesis events cannot be inserted
  1411. /// without their proof blob.
  1412. pub(crate) async fn dag_insert(
  1413. &self,
  1414. events: &[Event],
  1415. dag_name: &str,
  1416. ) -> Result<Vec<blake3::Hash>> {
  1417. self.dag_insert_inner(events, &[], /* require_blobs */ false, dag_name).await
  1418. }
  1419. /// Commit a rotating-DAG signal event whose RLN proof has already been
  1420. /// verified and recorded by the caller.
  1421. ///
  1422. /// Used by live protocol ingestion after `verify_rln_signal()` accepts the
  1423. /// event. The blob is persisted before the event body so late joiners never
  1424. /// observe a locally committed non-genesis event without its proof blob.
  1425. pub(crate) async fn insert_verified_signal(
  1426. &self,
  1427. event: &Event,
  1428. blob: &[u8],
  1429. dag_name: &str,
  1430. ) -> Result<Vec<blake3::Hash>> {
  1431. if event.header.parents != NULL_PARENTS && blob.is_empty() {
  1432. return Err(Error::Custom("verified signal event blob must not be empty".into()))
  1433. }
  1434. self.header_dag_insert(vec![event.header.clone()], dag_name).await?;
  1435. if event.header.parents != NULL_PARENTS {
  1436. self.dag_blob_store(&event.id(), blob)?;
  1437. }
  1438. self.dag_insert(std::slice::from_ref(event), dag_name).await
  1439. }
  1440. /// Insert events into a rotating DAG, with mandatory RLN
  1441. /// verification.
  1442. ///
  1443. /// `blobs` is index-aligned with `events`. Every non-genesis
  1444. /// event MUST have a non-empty `blobs[i]`; events that don't
  1445. /// (whether `blobs` is empty, shorter, or has an empty entry
  1446. /// at position `i`) are rejected with a loud log. This is the
  1447. /// strict policy required for sync paths - a peer that serves
  1448. /// an event without its blob is buggy or adversarial.
  1449. ///
  1450. /// On `Slashable`, this method does NOT broadcast a slash -
  1451. /// that's the protocol layer's job (see
  1452. /// `proto::handle_event_put::verify_rln_signal`). Sync-time
  1453. /// detection of a slashable conflict simply skips the event.
  1454. /// We don't want a node coming online to flood the network
  1455. /// with stale slash broadcasts.
  1456. pub async fn dag_insert_with_blobs(
  1457. &self,
  1458. events: &[Event],
  1459. blobs: &[Vec<u8>],
  1460. dag_name: &str,
  1461. ) -> Result<Vec<blake3::Hash>> {
  1462. self.dag_insert_inner(events, blobs, /* require_blobs */ true, dag_name).await
  1463. }
  1464. /// Inner implementation shared by both insert paths. The
  1465. /// `require_blobs` flag selects strict (sync) vs. lenient
  1466. /// (post-verified) semantics.
  1467. async fn dag_insert_inner(
  1468. &self,
  1469. events: &[Event],
  1470. blobs: &[Vec<u8>],
  1471. require_blobs: bool,
  1472. dag_name: &str,
  1473. ) -> Result<Vec<blake3::Hash>> {
  1474. if events.is_empty() {
  1475. return Ok(vec![])
  1476. }
  1477. // Pre-flight structural validation and RLN verification. Done
  1478. // BEFORE acquiring the DAG-store write lock so slow proof work does
  1479. // not hold up other inserts. Cheap structural checks run first, so
  1480. // malformed events cannot force proof verification or mutate RLN
  1481. // metadata.
  1482. //
  1483. // Events we already have are skipped without verification. This
  1484. // matters because `rln_verify_signal` records the share on `Accepted`,
  1485. // and re-running it for an already-seen event would trip its
  1486. // duplicate-share check.
  1487. let dag_ts = u64::from_str(dag_name)?;
  1488. let (already_have, structurally_valid): (Vec<bool>, Vec<bool>) = {
  1489. let store = self.dag_store.read().await;
  1490. let slot = store.get_slot(&dag_ts);
  1491. let mut already_have = Vec::with_capacity(events.len());
  1492. let mut structurally_valid = Vec::with_capacity(events.len());
  1493. for ev in events {
  1494. let eid = ev.id();
  1495. let have = match slot {
  1496. Some(s) => s.main_tree.contains_key(eid.as_bytes())?,
  1497. None => false,
  1498. };
  1499. already_have.push(have);
  1500. if have || ev.header.parents == NULL_PARENTS {
  1501. structurally_valid.push(true);
  1502. continue
  1503. }
  1504. let Some(slot) = slot else {
  1505. structurally_valid.push(false);
  1506. continue
  1507. };
  1508. if !slot.header_tree.contains_key(eid.as_bytes())? {
  1509. structurally_valid.push(false);
  1510. continue
  1511. }
  1512. structurally_valid
  1513. .push(ev.dag_validate(&slot.header_tree, &self.config, dag_ts).await?);
  1514. }
  1515. (already_have, structurally_valid)
  1516. };
  1517. let mut accepted: Vec<usize> = Vec::with_capacity(events.len());
  1518. for (i, ev) in events.iter().enumerate() {
  1519. if !structurally_valid[i] {
  1520. error!(
  1521. target: "event_graph::dag_insert",
  1522. "[DAG_INSERT] event {} failed structural validation before RLN verification; skipping",
  1523. ev.id(),
  1524. );
  1525. continue
  1526. }
  1527. // Already-known events go through structurally (the downstream
  1528. // `contains_key` check will skip them) but skip the RLN verifier to
  1529. // avoid double-recording the share for the same
  1530. // (epoch, internal_nullifier, x, y) tuple.
  1531. if already_have[i] {
  1532. accepted.push(i);
  1533. continue
  1534. }
  1535. // Genesis-shaped events have no blob and no proof - they're
  1536. // consensus inputs, not user signals.
  1537. if ev.header.parents == NULL_PARENTS {
  1538. accepted.push(i);
  1539. continue
  1540. }
  1541. let blob = blobs.get(i).cloned().unwrap_or_default();
  1542. if blob.is_empty() {
  1543. if require_blobs {
  1544. error!(
  1545. target: "event_graph::dag_insert",
  1546. "[DAG_INSERT] sync event {} arrived without an RLN blob; rejecting. \
  1547. Every non-genesis rotating-DAG event must carry a blob.",
  1548. ev.id(),
  1549. );
  1550. continue
  1551. }
  1552. // Lenient path: caller pre-verified. Accept the event
  1553. // structurally without running the RLN verifier on it.
  1554. accepted.push(i);
  1555. continue
  1556. }
  1557. match self.rln_verify_signal(ev, &blob).await {
  1558. rln::SignalCheck::Accepted => accepted.push(i),
  1559. rln::SignalCheck::Rejected => {
  1560. error!(
  1561. target: "event_graph::dag_insert",
  1562. "[DAG_INSERT] sync event {} failed RLN re-verification; skipping",
  1563. ev.id(),
  1564. );
  1565. }
  1566. rln::SignalCheck::Slashable(_) => {
  1567. // The conflicting share is recorded inside
  1568. // `rln_verify_signal` ONLY on `Accepted`. On `Slashable` it
  1569. // returns the conflicting shares without mutating metadata,
  1570. // so we don't double-record. We don't broadcast a slash
  1571. // here - that's the live broadcast handler's job. We just
  1572. // skip the event.
  1573. error!(
  1574. target: "event_graph::dag_insert",
  1575. "[DAG_INSERT] sync event {} is slashable (slot reuse); skipping",
  1576. ev.id(),
  1577. );
  1578. }
  1579. }
  1580. }
  1581. let mut bcast = self.broadcasted_ids.write().await;
  1582. let mut store = self.dag_store.write().await;
  1583. let slot = store.get_slot_mut(&dag_ts).ok_or(Error::DagSyncFailed)?;
  1584. let mut ids = Vec::with_capacity(accepted.len());
  1585. let mut overlay = SledTreeOverlay::new(&slot.main_tree);
  1586. for &i in &accepted {
  1587. let ev = &events[i];
  1588. let eid = ev.id();
  1589. if ev.header.parents == NULL_PARENTS {
  1590. continue
  1591. }
  1592. if slot.main_tree.contains_key(eid.as_bytes())? {
  1593. continue
  1594. }
  1595. if !slot.header_tree.contains_key(eid.as_bytes())? {
  1596. continue
  1597. }
  1598. if !ev.dag_validate(&slot.header_tree, &self.config, dag_ts).await? {
  1599. return Err(Error::EventIsInvalid)
  1600. }
  1601. let se = serialize_async(ev).await;
  1602. overlay.insert(eid.as_bytes(), &se)?;
  1603. if self.replay_mode {
  1604. replayer_log(&self.datastore, "insert".into(), se)?;
  1605. }
  1606. // Persist the blob alongside the event for future
  1607. // sync-time re-verification by other late-joiners.
  1608. if let Some(blob) = blobs.get(i) {
  1609. if !blob.is_empty() {
  1610. let _ = self.dag_blob_store(&eid, blob);
  1611. }
  1612. }
  1613. ids.push(eid);
  1614. }
  1615. if let Some(b) = overlay.aggregate() {
  1616. slot.main_tree.apply_batch(b).unwrap();
  1617. } else {
  1618. return Ok(vec![])
  1619. }
  1620. for &i in &accepted {
  1621. let ev = &events[i];
  1622. let eid = ev.id();
  1623. if ev.header.parents == NULL_PARENTS {
  1624. continue
  1625. }
  1626. for pid in ev.header.parents.iter() {
  1627. if *pid != NULL_ID {
  1628. for (layer, tips) in slot.tips.iter_mut() {
  1629. if *layer < ev.header.layer {
  1630. tips.remove(pid);
  1631. }
  1632. }
  1633. bcast.insert(*pid);
  1634. }
  1635. }
  1636. slot.tips.retain(|_, t| !t.is_empty());
  1637. slot.tips.entry(ev.header.layer).or_default().insert(eid);
  1638. self.event_pub.notify(ev.clone()).await;
  1639. }
  1640. Ok(ids)
  1641. }
  1642. pub async fn header_dag_insert(&self, headers: Vec<Header>, dag_name: &str) -> Result<()> {
  1643. let dag_ts = u64::from_str(dag_name)?;
  1644. // The genesis ID we expect any layer-1 header in this slot
  1645. // to reference. Computed locally from config - two networks
  1646. // with different `genesis_contents` (or any other config
  1647. // mismatch) produce different genesis ids, so a peer whose
  1648. // layer-1 headers reference something else is on a different
  1649. // network. Catching this explicitly here is strictly a
  1650. // defense-in-depth and diagnostics improvement: the existing
  1651. // parent-existence check in `Header::validate` already
  1652. // rejects these (genesis headers are filtered from
  1653. // `header_tree` on insert, so a foreign genesis id never
  1654. // lands in the local tree). The explicit boundary check just
  1655. // turns "HeaderIsInvalid" into a logged, named condition, so
  1656. // an operator debugging a misconfigured deployment sees
  1657. // "peer is on a different network" instead of a generic
  1658. // header rejection.
  1659. //
  1660. // Why layer 1 is sufficient: `select_parents_from_tips` puts
  1661. // an event at layer N+1 where N is the highest layer with
  1662. // tips. For layer = 1, the highest tip layer must be 0, and
  1663. // the only layer-0 entry in any slot is the genesis (the
  1664. // single event placed by `DagStore::create_slot`). So every
  1665. // layer-1 event's non-NULL parents are equal to that slot's
  1666. // genesis id. Higher layers don't need the check because
  1667. // their parent chains transitively pass through layer 1; if
  1668. // the layer-1 events get rejected, layer-2+ events lose
  1669. // their referenced parents and fail the existing parent-
  1670. // existence check.
  1671. let local_genesis_id = Header {
  1672. timestamp: dag_ts,
  1673. parents: NULL_PARENTS,
  1674. layer: 0,
  1675. content_hash: blake3::hash(&self.config.genesis_contents),
  1676. }
  1677. .id();
  1678. let mut store = self.dag_store.write().await;
  1679. let slot = store.get_slot_mut(&dag_ts).ok_or(Error::DagSyncFailed)?;
  1680. let mut overlay = SledTreeOverlay::new(&slot.header_tree);
  1681. let mut hdrs = headers;
  1682. hdrs.sort_by_key(|h| h.layer);
  1683. for hdr in &hdrs {
  1684. if hdr.parents == NULL_PARENTS {
  1685. continue
  1686. }
  1687. // Cross-network detection at the layer-1 boundary.
  1688. if hdr.layer == 1 {
  1689. for pid in hdr.parents.iter() {
  1690. if *pid != NULL_ID && *pid != local_genesis_id {
  1691. error!(
  1692. target: "event_graph::header_dag_insert",
  1693. "[HEADER_DAG_INSERT] layer-1 header for dag {dag_ts} \
  1694. references foreign genesis: claimed parent {pid:?}, \
  1695. local genesis is {local_genesis_id:?}. Peer is on a \
  1696. different network.",
  1697. );
  1698. return Err(Error::HeaderIsInvalid)
  1699. }
  1700. }
  1701. }
  1702. let hid = hdr.id();
  1703. if overlay.get(hid.as_bytes())?.is_some() {
  1704. continue
  1705. }
  1706. if !hdr.validate(&slot.header_tree, &self.config, dag_ts, Some(&overlay)).await? {
  1707. return Err(Error::HeaderIsInvalid)
  1708. }
  1709. overlay.insert(hid.as_bytes(), &serialize_async(hdr).await)?;
  1710. slot.time_index.insert(hdr.timestamp, hid);
  1711. }
  1712. if let Some(b) = overlay.aggregate() {
  1713. slot.header_tree.apply_batch(b).unwrap();
  1714. }
  1715. Ok(())
  1716. }
  1717. pub async fn fetch_event_from_dags(&self, eid: &blake3::Hash) -> Result<Option<Event>> {
  1718. for (_, slot) in self.dag_store.read().await.dags.iter() {
  1719. if let Some(b) = slot.main_tree.get(eid.as_bytes())? {
  1720. return Ok(Some(deserialize_async(&b).await?))
  1721. }
  1722. }
  1723. // Also check the static DAG. Static events (RLN registrations
  1724. // and slashes) are public consensus state, so they're served
  1725. // alongside rotating-DAG events through the same EventReq
  1726. // path. This is what lets a fresh peer's `static_sync` walk
  1727. // ancestry through EventReq after discovering tips.
  1728. if let Some(b) = self.static_dag.get(eid.as_bytes())? {
  1729. return Ok(Some(deserialize_async(&b).await?))
  1730. }
  1731. Ok(None)
  1732. }
  1733. pub(crate) async fn get_next_layer_with_parents(
  1734. &self,
  1735. dag_ts: &u64,
  1736. ) -> (u64, [blake3::Hash; N_EVENT_PARENTS]) {
  1737. select_parents_from_tips(&self.dag_store.read().await.get_slot(dag_ts).unwrap().tips)
  1738. }
  1739. pub(crate) async fn get_next_layer_with_parents_static(
  1740. &self,
  1741. ) -> (u64, [blake3::Hash; N_EVENT_PARENTS]) {
  1742. select_parents_from_tips(&compute_unreferenced_tips(&self.static_dag).await)
  1743. }
  1744. pub async fn order_events(&self) -> Vec<Event> {
  1745. let mut all = vec![];
  1746. for (_, slot) in self.dag_store.read().await.dags.iter() {
  1747. for item in slot.main_tree.iter() {
  1748. let (_, b) = item.unwrap();
  1749. let ev: Event = deserialize_async(&b).await.unwrap();
  1750. if ev.header.parents != NULL_PARENTS {
  1751. all.push(ev);
  1752. }
  1753. }
  1754. }
  1755. all.sort_unstable_by(display_order);
  1756. all
  1757. }
  1758. pub async fn fetch_headers_with_tips(
  1759. &self,
  1760. dag_name: &str,
  1761. tips: &LayerUTips,
  1762. ) -> Result<Vec<Header>> {
  1763. if count_layer_tips(tips) > MAX_HEADER_REQ_TIPS {
  1764. return Err(Error::DagSyncFailed)
  1765. }
  1766. let dag_ts = u64::from_str(dag_name)?;
  1767. let store = self.dag_store.read().await;
  1768. let slot = store.get_slot(&dag_ts).ok_or(Error::DagSyncFailed)?;
  1769. let mut ancestors = HashSet::new();
  1770. for hashes in tips.values() {
  1771. for h in hashes {
  1772. ancestors.insert(*h);
  1773. if let Some(v) = slot.header_tree.get(h.as_bytes())? {
  1774. self.get_ancestors(
  1775. &mut ancestors,
  1776. deserialize_async(&v).await?,
  1777. &slot.header_tree,
  1778. )
  1779. .await?;
  1780. }
  1781. }
  1782. }
  1783. let mut out = Vec::with_capacity(MAX_HEADER_REP_HEADERS);
  1784. let sort_headers = |headers: &mut Vec<Header>| {
  1785. headers.sort_unstable_by(|a, b| {
  1786. a.layer.cmp(&b.layer).then_with(|| a.id().as_bytes().cmp(b.id().as_bytes()))
  1787. });
  1788. };
  1789. for item in slot.header_tree.iter() {
  1790. let (id, v) = item?;
  1791. let h = blake3::Hash::from_bytes((&id as &[u8]).try_into()?);
  1792. if ancestors.contains(&h) {
  1793. continue
  1794. }
  1795. out.push(deserialize_async(&v).await?);
  1796. if out.len() >= MAX_HEADER_REP_HEADERS * 2 {
  1797. sort_headers(&mut out);
  1798. out.truncate(MAX_HEADER_REP_HEADERS);
  1799. }
  1800. }
  1801. sort_headers(&mut out);
  1802. out.truncate(MAX_HEADER_REP_HEADERS);
  1803. Ok(out)
  1804. }
  1805. pub(crate) async fn get_ancestors(
  1806. &self,
  1807. visited: &mut HashSet<blake3::Hash>,
  1808. hdr: Header,
  1809. tree: &sled::Tree,
  1810. ) -> Result<()> {
  1811. let mut stack = VecDeque::new();
  1812. stack.push_back(hdr);
  1813. while let Some(h) = stack.pop_back() {
  1814. for p in h.parents {
  1815. if p != NULL_ID && visited.insert(p) {
  1816. if let Some(v) = tree.get(p.as_bytes())? {
  1817. stack.push_back(deserialize_async(&v).await?);
  1818. }
  1819. }
  1820. }
  1821. }
  1822. Ok(())
  1823. }
  1824. async fn static_new(sled_db: &sled::Db, config: &EventGraphConfig) -> Result<sled::Tree> {
  1825. let tree = sled_db.open_tree("static-dag")?;
  1826. let genesis = generate_static_genesis(config);
  1827. let mut ov = SledTreeOverlay::new(&tree);
  1828. ov.insert(genesis.id().as_bytes(), &serialize_async(&genesis).await).unwrap();
  1829. if let Some(b) = ov.aggregate() {
  1830. tree.apply_batch(b).unwrap();
  1831. }
  1832. Ok(tree)
  1833. }
  1834. pub async fn static_broadcast(&self, ev: Event, blob: Vec<u8>) -> Result<()> {
  1835. self.p2p.broadcast(&StaticPut(ev, blob)).await;
  1836. Ok(())
  1837. }
  1838. async fn static_persist(&self, ev: &Event) -> Result<()> {
  1839. let mut ov = SledTreeOverlay::new(&self.static_dag);
  1840. ov.insert(ev.id().as_bytes(), &serialize_async(ev).await).unwrap();
  1841. if let Some(b) = ov.aggregate() {
  1842. self.static_dag.apply_batch(b).unwrap();
  1843. }
  1844. Ok(())
  1845. }
  1846. #[cfg(test)]
  1847. pub(crate) async fn static_insert(&self, ev: &Event) -> Result<()> {
  1848. self.static_persist(ev).await?;
  1849. self.static_pub.notify(ev.clone()).await;
  1850. Ok(())
  1851. }
  1852. /// Durably commit a verified static RLN event.
  1853. ///
  1854. /// The write order is intentional: blob first, static DAG second, RLN
  1855. /// state last. If a process crashes after the static event becomes
  1856. /// durable but before the identity tree or historical-root indexes are
  1857. /// updated, startup recovery can rebuild those RLN side tables from the
  1858. /// static DAG. Subscribers are notified only after the RLN apply step, so
  1859. /// applications observe the same post-state semantics as the receive path.
  1860. pub async fn commit_verified_static_event(
  1861. &self,
  1862. ev: &Event,
  1863. blob: &[u8],
  1864. rln_node: &rln::RLNNode,
  1865. ) -> Result<pallas::Base> {
  1866. if blob.is_empty() {
  1867. return Err(Error::Custom("static RLN event blob must not be empty".into()))
  1868. }
  1869. self.static_blob_store(&ev.id(), blob)?;
  1870. self.static_persist(ev).await?;
  1871. let root = self.apply_rln_static_event(ev, rln_node).await?;
  1872. self.static_pub.notify(ev.clone()).await;
  1873. Ok(root)
  1874. }
  1875. pub async fn static_fetch(&self, eid: &blake3::Hash) -> Result<Option<Event>> {
  1876. Ok(match self.static_dag.get(eid.as_bytes())? {
  1877. Some(b) => Some(deserialize_async(&b).await?),
  1878. None => None,
  1879. })
  1880. }
  1881. pub async fn static_unreferenced_tips(&self) -> LayerUTips {
  1882. compute_unreferenced_tips(&self.static_dag).await
  1883. }
  1884. /// Audit static-DAG blob coverage and repair deterministic guard blobs.
  1885. ///
  1886. /// A static DAG event without its RLN blob cannot be served to late
  1887. /// joiners because they must re-verify historical static events. The only
  1888. /// blob we can safely reconstruct is the pregenerated-registration guard:
  1889. /// it is valid exactly for commitments supplied by this app config. Slash
  1890. /// blobs and future staked registration proofs are not reconstructible and
  1891. /// are logged for operator intervention.
  1892. async fn audit_static_blobs(&self) -> Result<()> {
  1893. let mut repaired = 0usize;
  1894. let mut unrecoverable = 0usize;
  1895. let mut malformed = 0usize;
  1896. for item in self.static_dag.iter() {
  1897. let (_, val) = item?;
  1898. let ev: Event = deserialize_async(&val).await?;
  1899. if ev.header.parents == NULL_PARENTS {
  1900. continue
  1901. }
  1902. if matches!(self.static_blob_fetch(&ev.id())?, Some(blob) if !blob.is_empty()) {
  1903. continue
  1904. }
  1905. let rln_node: rln::RLNNode = match deserialize_async_partial(ev.content()).await {
  1906. Ok((node, _)) => node,
  1907. Err(_) => {
  1908. malformed += 1;
  1909. continue
  1910. }
  1911. };
  1912. match rln_node {
  1913. rln::RLNNode::Registration(commitment)
  1914. if self
  1915. .pregenerated_identity_commitment_reprs
  1916. .contains(&commitment.to_repr()) =>
  1917. {
  1918. self.static_blob_store(&ev.id(), rln::GENESIS_BLOB_GUARD)?;
  1919. repaired += 1;
  1920. }
  1921. _ => {
  1922. unrecoverable += 1;
  1923. warn!(
  1924. target: "event_graph::new",
  1925. "[EVENTGRAPH] static event {} is missing its RLN blob and cannot be reconstructed",
  1926. ev.id(),
  1927. );
  1928. }
  1929. }
  1930. }
  1931. if repaired > 0 || unrecoverable > 0 || malformed > 0 {
  1932. info!(
  1933. target: "event_graph::new",
  1934. "[EVENTGRAPH] static blob audit: repaired={} unrecoverable={} malformed={}",
  1935. repaired, unrecoverable, malformed,
  1936. );
  1937. }
  1938. Ok(())
  1939. }
  1940. /// Persist the original RLN blob for a static-DAG event. The
  1941. /// blob is the wire payload from the originating `StaticPut` -
  1942. /// proof + public inputs + attestation - needed to re-verify
  1943. /// the proof at sync time by late-joining peers.
  1944. ///
  1945. /// Writing the same `(eid, blob)` repeatedly is safe.
  1946. pub fn static_blob_store(&self, eid: &blake3::Hash, blob: &[u8]) -> Result<()> {
  1947. self.static_dag_blobs.insert(eid.as_bytes(), blob)?;
  1948. Ok(())
  1949. }
  1950. /// Look up the original RLN blob for a static-DAG event.
  1951. ///
  1952. /// Returns `Ok(None)` only for legitimate "not stored" cases -
  1953. /// a peer that's never seen the event, or an event that pre-dates
  1954. /// the side-table. The verification path in `static_sync` treats
  1955. /// missing blobs on non-genesis events as a sync failure, not as
  1956. /// a fall-through.
  1957. pub fn static_blob_fetch(&self, eid: &blake3::Hash) -> Result<Option<Vec<u8>>> {
  1958. Ok(self.static_dag_blobs.get(eid.as_bytes())?.map(|ivec| ivec.to_vec()))
  1959. }
  1960. /// Persist the original RLN signal blob for a rotating-DAG event.
  1961. ///
  1962. /// Mirror of [`Self::static_blob_store`] but for rotating-DAG
  1963. /// events. Idempotent. Called after successful RLN verification
  1964. /// in `handle_event_put`, and during sync by
  1965. /// `dag_insert_with_blobs` when the peer included a blob in
  1966. /// `EventRep`.
  1967. pub fn dag_blob_store(&self, eid: &blake3::Hash, blob: &[u8]) -> Result<()> {
  1968. self.dag_blobs.insert(eid.as_bytes(), blob)?;
  1969. Ok(())
  1970. }
  1971. /// Look up the original RLN signal blob for a rotating-DAG
  1972. /// event. Returns `Ok(None)` if not stored - see
  1973. /// [`Self::static_blob_fetch`] for the exhaustive list of
  1974. /// reasons a blob may legitimately be missing.
  1975. ///
  1976. /// Note: rotating-DAG blobs are pruned alongside their DAGs
  1977. /// (see `dag_blobs_prune`). Older-than-window events therefore
  1978. /// don't accumulate blobs in this side-table.
  1979. pub fn dag_blob_fetch(&self, eid: &blake3::Hash) -> Result<Option<Vec<u8>>> {
  1980. Ok(self.dag_blobs.get(eid.as_bytes())?.map(|ivec| ivec.to_vec()))
  1981. }
  1982. /// Apply a static-DAG event (registration or slash) to the
  1983. /// identity-state SMT, and record the resulting root in the
  1984. /// historical-roots side-tables.
  1985. ///
  1986. /// **This is the single canonical entry point** for SMT
  1987. /// mutation. All callers - live broadcast (`handle_static_put`),
  1988. /// originator (`nickserv.rs::handle_register`), and sync
  1989. /// (`static_sync` apply loop) - go through here. Bypassing it
  1990. /// will desynchronize the SMT from the historical-roots tables,
  1991. /// which silently breaks signal verification.
  1992. ///
  1993. /// **Canonical order requirement.** Two nodes processing the
  1994. /// same set of static events must produce the same sequence of
  1995. /// intermediate roots. SMTs are commutative under set-of-leaves
  1996. /// (final root is order-independent) but the *intermediate*
  1997. /// roots produced during application are order-dependent. We
  1998. /// pin the order with `(layer, event_id)`: layer is the primary
  1999. /// key (defined by the event's parent links and consensus-agreed),
  2000. /// event_id is the tie-breaker within a layer (32-byte hash,
  2001. /// total-ordered lexicographically).
  2002. ///
  2003. /// In live broadcast and originator paths, events arrive one at
  2004. /// a time; the canonical-order requirement is automatically
  2005. /// satisfied because each event's layer is greater than its
  2006. /// parents'. In sync, the caller must sort by `(layer, event_id)`
  2007. /// before invoking this method (see `static_sync`).
  2008. ///
  2009. /// **Returns** the post-mutation SMT root, or an error if the
  2010. /// SMT mutation itself fails. A duplicate-registration or
  2011. /// slash-of-nonexistent are both treated as soft no-ops at the
  2012. /// SMT layer, but we still record the root (which equals the
  2013. /// pre-call root in that case) - this preserves the invariant
  2014. /// that "every static-DAG event has a corresponding entry in
  2015. /// rln-historical-roots-ordered" without complicating the
  2016. /// caller's logic.
  2017. pub async fn apply_rln_static_event(
  2018. &self,
  2019. ev: &Event,
  2020. node: &rln::RLNNode,
  2021. ) -> Result<pallas::Base> {
  2022. let mut state = self.identity_state.write().await;
  2023. match node {
  2024. rln::RLNNode::Registration(commitment) => {
  2025. // Soft-fail on duplicate (race with another peer).
  2026. let _ = state.register(*commitment);
  2027. }
  2028. rln::RLNNode::Slashing(commitment) => {
  2029. // Idempotent - slashing a non-present identity is
  2030. // a no-op.
  2031. let _ = state.slash(*commitment);
  2032. }
  2033. }
  2034. let new_root = state.root();
  2035. drop(state);
  2036. // Record the root in both side-tables. We do this even if
  2037. // the SMT mutation was a no-op (duplicate register, slash of
  2038. // missing) so the historical-roots table has one entry per
  2039. // static-DAG event. This makes canonical-order replay simple:
  2040. // every event has exactly one entry, no conditional skips.
  2041. let key = encode_historical_root_key(ev.header.layer, &ev.id());
  2042. let value = encode_historical_root_value(&new_root, ev.header.timestamp);
  2043. self.rln_historical_roots_ordered.insert(key, value.as_slice())?;
  2044. let by_value_key = encode_historical_root_by_value_key(&new_root, &key);
  2045. self.rln_historical_roots_by_value.insert(by_value_key, &[])?;
  2046. Ok(new_root)
  2047. }
  2048. /// Check whether `root` is a valid SMT root for a signal whose
  2049. /// `signal_timestamp` is given (in millis-since-epoch).
  2050. ///
  2051. /// A root is valid if it was the live root at any time in the
  2052. /// drift window `[signal_timestamp - DRIFT, signal_timestamp +
  2053. /// DRIFT]`. The drift symmetry handles two distinct concerns:
  2054. ///
  2055. /// * **Forward drift** (signal sees a slightly stale root): the
  2056. /// originator built a proof against root `R_n`, then someone
  2057. /// else registered, producing `R_{n+1}`, before the signal
  2058. /// reached the verifier. The signal's claimed `R_n` is older
  2059. /// than the verifier's current root by an amount up to the
  2060. /// propagation delay. Accept if R_n was current within DRIFT
  2061. /// of the signal's timestamp.
  2062. ///
  2063. /// * **Backward drift** (signal arrives before its root): rare
  2064. /// but possible if the static-DAG broadcast is racing the
  2065. /// rotating-DAG broadcast. The originator's machine knew about
  2066. /// a registration that hadn't fully propagated yet. We tolerate
  2067. /// up to DRIFT of backward skew.
  2068. ///
  2069. /// The check uses `rln_historical_roots_by_value` to find every
  2070. /// canonical position where `root` appears, then
  2071. /// `rln_historical_roots_ordered` to bracket each interval during
  2072. /// which `root` was live. Each interval starts at the timestamp
  2073. /// of an event that produced `root` and ends just before the next
  2074. /// event timestamp (or `u64::MAX` if `root` is currently live).
  2075. ///
  2076. /// This subsumes the old `recent_roots` window as a special case
  2077. /// (live verification = signal_timestamp ~= now). The in-memory
  2078. /// `recent_roots` cache in `IdentityState` remains as a fast
  2079. /// path for the live-broadcast hot loop; this method is consulted
  2080. /// when the cache misses.
  2081. pub fn is_root_valid_at(&self, root: &pallas::Base, signal_timestamp: u64) -> Result<bool> {
  2082. let drift = EVENT_TIME_DRIFT;
  2083. let lo = signal_timestamp.saturating_sub(drift);
  2084. let hi = signal_timestamp.saturating_add(drift);
  2085. for item in self.rln_historical_roots_by_value.scan_prefix(root.to_repr()) {
  2086. let (by_value_key, _) = item?;
  2087. if by_value_key.len() != 72 {
  2088. continue
  2089. }
  2090. let ordered_key = &by_value_key[32..];
  2091. let Some(value_bytes) = self.rln_historical_roots_ordered.get(ordered_key)? else {
  2092. continue
  2093. };
  2094. let (recorded_root, root_timestamp) = decode_historical_root_value(&value_bytes)?;
  2095. if &recorded_root != root {
  2096. continue
  2097. }
  2098. let next_timestamp: u64 = {
  2099. use std::ops::Bound::{Excluded, Unbounded};
  2100. match self
  2101. .rln_historical_roots_ordered
  2102. .range::<&[u8], _>((Excluded(ordered_key), Unbounded))
  2103. .next()
  2104. {
  2105. Some(Ok((_, val))) => decode_historical_root_value(&val)?.1,
  2106. Some(Err(e)) => return Err(e.into()),
  2107. None => u64::MAX,
  2108. }
  2109. };
  2110. if root_timestamp <= hi && next_timestamp > lo {
  2111. return Ok(true)
  2112. }
  2113. }
  2114. Ok(false)
  2115. }
  2116. }
  2117. fn encode_historical_root_key(layer: u64, event_id: &blake3::Hash) -> [u8; 40] {
  2118. let mut buf = [0u8; 40];
  2119. buf[..8].copy_from_slice(&layer.to_be_bytes());
  2120. buf[8..].copy_from_slice(event_id.as_bytes());
  2121. buf
  2122. }
  2123. fn encode_historical_root_value(root: &pallas::Base, timestamp: u64) -> [u8; 40] {
  2124. let mut buf = [0u8; 40];
  2125. buf[..32].copy_from_slice(&root.to_repr());
  2126. buf[32..].copy_from_slice(&timestamp.to_be_bytes());
  2127. buf
  2128. }
  2129. fn encode_historical_root_by_value_key(root: &pallas::Base, ordered_key: &[u8; 40]) -> [u8; 72] {
  2130. let mut buf = [0u8; 72];
  2131. buf[..32].copy_from_slice(&root.to_repr());
  2132. buf[32..].copy_from_slice(ordered_key);
  2133. buf
  2134. }
  2135. fn decode_historical_root_value(bytes: &[u8]) -> Result<(pallas::Base, u64)> {
  2136. if bytes.len() != 40 {
  2137. return Err(Error::Custom(format!(
  2138. "historical-root value must be 40 bytes, got {}",
  2139. bytes.len()
  2140. )))
  2141. }
  2142. let mut root_repr = [0u8; 32];
  2143. root_repr.copy_from_slice(&bytes[..32]);
  2144. let root: pallas::Base = match pallas::Base::from_repr(root_repr).into() {
  2145. Some(r) => r,
  2146. None => return Err(Error::Custom("invalid root encoding".into())),
  2147. };
  2148. let mut ts_bytes = [0u8; 8];
  2149. ts_bytes.copy_from_slice(&bytes[32..]);
  2150. Ok((root, u64::from_be_bytes(ts_bytes)))
  2151. }
  2152. impl EventGraph {
  2153. /// Return a JSON-RPC response representing the current state of
  2154. /// the event graph.
  2155. ///
  2156. /// Shape (matches [`util::recreate_from_replayer_log`] so clients
  2157. /// can reuse their parsers):
  2158. ///
  2159. /// ```json
  2160. /// {
  2161. /// "eventgraph_info": {
  2162. /// "dag": {
  2163. /// "<event-id-hex>": <event>,
  2164. /// ...
  2165. /// }
  2166. /// }
  2167. /// }
  2168. /// ```
  2169. ///
  2170. /// Walks every event currently held in every rotating DAG *and*
  2171. /// every event in the static DAG. Genesis events are included.
  2172. #[cfg(feature = "rpc")]
  2173. pub async fn eventgraph_info(
  2174. &self,
  2175. id: u16,
  2176. _params: crate::rpc::util::JsonValue,
  2177. ) -> crate::rpc::jsonrpc::JsonResult {
  2178. use crate::rpc::{
  2179. jsonrpc::{JsonResponse, JsonResult},
  2180. util::{json_map, JsonValue},
  2181. };
  2182. let mut dag = HashMap::new();
  2183. // Walk every rotating DAG.
  2184. for (_, slot) in self.dag_store.read().await.dags.iter() {
  2185. for item in slot.main_tree.iter() {
  2186. let (eid, val) = match item {
  2187. Ok(v) => v,
  2188. Err(_) => continue,
  2189. };
  2190. let Ok(ev) = deserialize_async::<Event>(&val).await else { continue };
  2191. let key = blake3::Hash::from_bytes(match (&eid as &[u8]).try_into() {
  2192. Ok(b) => b,
  2193. Err(_) => continue,
  2194. });
  2195. dag.insert(key.to_string(), JsonValue::from(ev));
  2196. }
  2197. }
  2198. // And the static DAG.
  2199. for item in self.static_dag.iter() {
  2200. let (eid, val) = match item {
  2201. Ok(v) => v,
  2202. Err(_) => continue,
  2203. };
  2204. let Ok(ev) = deserialize_async::<Event>(&val).await else { continue };
  2205. let key = blake3::Hash::from_bytes(match (&eid as &[u8]).try_into() {
  2206. Ok(b) => b,
  2207. Err(_) => continue,
  2208. });
  2209. dag.insert(key.to_string(), JsonValue::from(ev));
  2210. }
  2211. let values = json_map([("dag", JsonValue::Object(dag))]);
  2212. let result = JsonValue::Object(HashMap::from([("eventgraph_info".into(), values)]));
  2213. JsonResult::Response(JsonResponse::new(result, id))
  2214. }
  2215. pub fn deg_enable(&self) {
  2216. self.deg_enabled.store(true, Ordering::Release);
  2217. }
  2218. pub fn deg_disable(&self) {
  2219. self.deg_enabled.store(false, Ordering::Release);
  2220. }
  2221. pub fn is_deg_enabled(&self) -> bool {
  2222. self.deg_enabled.load(Ordering::Acquire)
  2223. }
  2224. pub fn is_synced(&self) -> bool {
  2225. self.synced.load(Ordering::Acquire)
  2226. }
  2227. pub async fn deg_subscribe(&self) -> Subscription<DegEvent> {
  2228. self.deg_publisher.clone().subscribe().await
  2229. }
  2230. pub async fn deg_notify(&self, ev: DegEvent) {
  2231. self.deg_publisher.notify(ev).await;
  2232. }
  2233. /// Subscribe to rotating-DAG event insertions.
  2234. ///
  2235. /// Each subscriber receives a clone of every [`Event`] that
  2236. /// passes validation and is committed via `dag_insert`. The
  2237. /// publisher fires *after* state mutation, so subscribers can
  2238. /// rely on the event being durably present in the DAG by the
  2239. /// time they observe it.
  2240. ///
  2241. /// Used to build live JSON-RPC subscription endpoints (e.g.
  2242. /// the Gource-feeding endpoint in DarkIRC). Bridge a
  2243. /// [`Subscription`] to a `JsonSubscriber` in the application
  2244. /// layer; this method itself contains no JSON-RPC logic.
  2245. pub async fn event_subscribe(&self) -> Subscription<Event> {
  2246. self.event_pub.clone().subscribe().await
  2247. }
  2248. /// Subscribe to static-DAG event insertions (RLN registrations
  2249. /// and slashes). Mirrors [`Self::event_subscribe`] but for the
  2250. /// static DAG. See that method for semantics.
  2251. pub async fn static_subscribe(&self) -> Subscription<Event> {
  2252. self.static_pub.clone().subscribe().await
  2253. }
  2254. /// The app identifier mixed into RLN external nullifiers.
  2255. /// Exposed so the proto layer (and clients constructing signal
  2256. /// proofs) can use the same value the verifier uses.
  2257. pub fn rln_app_id(&self) -> rln::RlnAppId {
  2258. self.rln_app_id
  2259. }
  2260. /// Build a membership proof for an identity commitment, plus the
  2261. /// current root.
  2262. ///
  2263. /// This is the *only* sanctioned path for clients to produce a
  2264. /// signal proof - they should not be holding their own copy of
  2265. /// the SMT (the previous `event_graph.rln_identity_tree`
  2266. /// pattern). Centralising this keeps client and verifier in
  2267. /// agreement on the root they're proving against.
  2268. pub async fn rln_membership_path(
  2269. &self,
  2270. commitment: &darkfi_sdk::pasta::pallas::Base,
  2271. ) -> (darkfi_sdk::pasta::pallas::Base, darkfi_sdk::crypto::smt::PathFp) {
  2272. let s = self.identity_state.read().await;
  2273. (s.root(), s.prove_membership(commitment))
  2274. }
  2275. /// True if the given identity commitment is registered.
  2276. pub async fn rln_contains(&self, commitment: &darkfi_sdk::pasta::pallas::Base) -> bool {
  2277. self.identity_state.read().await.contains(commitment)
  2278. }
  2279. /// Verify an RLN signal blob against this event graph's state.
  2280. ///
  2281. /// The returned variant tells the caller what to do:
  2282. /// * `Accepted` - proof valid, no conflict, share recorded.
  2283. /// * `Rejected` - drop silently (bad proof, bad bounds, bad
  2284. /// root, exact duplicate).
  2285. /// * `Slashable(shares)` - different `(x, y)` for the same
  2286. /// internal nullifier; the caller (protocol layer) should
  2287. /// build and broadcast a slash from these shares.
  2288. ///
  2289. /// Critical invariant: `Rejected` and `Slashable` outcomes never
  2290. /// mutate `metadata`. `Accepted` is the only mutating outcome.
  2291. /// This is what prevents share-poisoning by an adversary who
  2292. /// observes an honest internal_nullifier and tries to forge a
  2293. /// share against it.
  2294. pub async fn rln_verify_signal(&self, event: &Event, blob: &[u8]) -> rln::SignalCheck {
  2295. use darkfi_sdk::{crypto::poseidon_hash, pasta::pallas};
  2296. use rln::{epoch_of, hash_event, Blob, SignalCheck, MAX_MSG_LIMIT};
  2297. let rcvd: Blob = match deserialize_async_partial(blob).await {
  2298. Ok((v, _)) => v,
  2299. Err(_) => return SignalCheck::Rejected,
  2300. };
  2301. // Defensive bounds. The proof PI binds these too, but
  2302. // checking up front lets us skip an expensive verify() call
  2303. // for trivially malformed blobs.
  2304. if rcvd.user_msg_limit == 0 || rcvd.user_msg_limit > MAX_MSG_LIMIT {
  2305. return SignalCheck::Rejected
  2306. }
  2307. let epoch_n = epoch_of(event.header.timestamp);
  2308. let epoch_field = pallas::Base::from(epoch_n);
  2309. let app_id = self.rln_app_id();
  2310. let ext_null = poseidon_hash([epoch_field, app_id.as_field()]);
  2311. let x = hash_event(event);
  2312. // 1) The merkle root must be valid for a signal at this
  2313. // timestamp. We accept any root that was the live SMT
  2314. // root at any time within EVENT_TIME_DRIFT of the signal's
  2315. // timestamp. This subsumes the old "recent_roots window"
  2316. // behaviour as a special case (live verification, signal
  2317. // timestamp ~= now) AND supports sync of historical signals
  2318. // (signal timestamp = signing time, root corresponds to
  2319. // that historical state).
  2320. //
  2321. // Hot-path optimization: check the in-memory recent_roots
  2322. // cache first. For live broadcasts (the overwhelming
  2323. // majority of verifications) the root will be in the
  2324. // cache and we skip the sled lookup.
  2325. {
  2326. let id_state = self.identity_state.read().await;
  2327. if !id_state.is_known_root(&rcvd.merkle_root) {
  2328. drop(id_state);
  2329. match self.is_root_valid_at(&rcvd.merkle_root, event.header.timestamp) {
  2330. Ok(true) => {}
  2331. Ok(false) => {
  2332. // Useful diagnostic: this rejection path
  2333. // catches both "garbage root" (attacker
  2334. // submitted a forged root) and "slashed-user
  2335. // replay" (slashed identity claiming a
  2336. // pre-slash root after the propagation
  2337. // window expired). The two are
  2338. // indistinguishable from the verifier's
  2339. // perspective by design - RLN-V2 privacy
  2340. // guarantees prevent identifying the
  2341. // signer. But observed in aggregate, a
  2342. // burst of these from a single peer is a
  2343. // strong signal of replay-after-slash
  2344. // misbehavior, useful for operators
  2345. // debugging "why are my messages being
  2346. // rejected" or investigating peer abuse.
  2347. let local_root = self.identity_state.read().await.root();
  2348. let historical_count = self.rln_historical_roots_ordered.len();
  2349. warn!(
  2350. target: "event_graph::rln_verify_signal",
  2351. "[RLN] Signal rejected: merkle_root not valid at signal \
  2352. timestamp {}. Possible causes: (1) forged or out-of-sync \
  2353. root, (2) slashed identity replaying against a pre-slash \
  2354. root after the propagation window expired. event_id={}, \
  2355. received_root={:?}, local_current_root={:?}, \
  2356. historical_root_count={}",
  2357. event.header.timestamp,
  2358. event.id(),
  2359. rcvd.merkle_root,
  2360. local_root,
  2361. historical_count,
  2362. );
  2363. return SignalCheck::Rejected
  2364. }
  2365. Err(e) => {
  2366. error!(
  2367. target: "event_graph::rln_verify_signal",
  2368. "[RLN] is_root_valid_at lookup failed for event {}: {e}",
  2369. event.id(),
  2370. );
  2371. return SignalCheck::Rejected
  2372. }
  2373. }
  2374. }
  2375. }
  2376. // 2) Verify the ZK proof. PI order MUST match
  2377. // constrain_instance() in rlnv2-diff-signal.zk:
  2378. // root, external_nullifier, user_message_limit, x, y, internal_nullifier
  2379. let pi = vec![
  2380. rcvd.merkle_root,
  2381. ext_null,
  2382. pallas::Base::from(rcvd.user_msg_limit),
  2383. x,
  2384. rcvd.y,
  2385. rcvd.internal_nullifier,
  2386. ];
  2387. if rcvd.proof.verify(&self.zk_keys.signal_vk, &pi).is_err() {
  2388. return SignalCheck::Rejected
  2389. }
  2390. // 3) Now consult the metadata table. Any share we look at
  2391. // here is guaranteed to be from a valid proof.
  2392. let mut state = self.rln_state.write().await;
  2393. // Prune metadata relative to THIS signal's epoch, not
  2394. // wall-clock. The retention window is conceptually
  2395. // "epochs near the signal we're processing", and the only
  2396. // entries that matter for reuse detection are siblings
  2397. // within `METADATA_RETAIN_EPOCHS` of `epoch_n`.
  2398. //
  2399. // In production, signals carry roughly-current wall-clock
  2400. // timestamps, so `epoch_n ~= current_epoch()` and this is
  2401. // equivalent to the previous `current_epoch()`-based prune.
  2402. // The change matters in three places:
  2403. //
  2404. // * Sync of historical signals: a late-arriving signal
  2405. // keeps its epoch's metadata visible long enough for
  2406. // the verifier to detect reuse. With wall-clock prune,
  2407. // a signal old enough to be outside the retention
  2408. // window would always silently lose its sibling shares
  2409. // before they could be matched.
  2410. //
  2411. // * Tests with deterministic event timestamps: the
  2412. // verifier and the metadata stay consistent regardless
  2413. // of when the test runs.
  2414. //
  2415. // * Minor DoS surface: a peer with a far-future system
  2416. // clock no longer causes mass-wipe of real metadata
  2417. // before consultation.
  2418. state.metadata.prune_old(epoch_n);
  2419. if state.metadata.is_duplicate(epoch_n, &rcvd.internal_nullifier, &x, &rcvd.y) {
  2420. return SignalCheck::Rejected
  2421. }
  2422. if state.metadata.is_reused(epoch_n, &rcvd.internal_nullifier) {
  2423. let mut shares = state.metadata.get_shares(epoch_n, &rcvd.internal_nullifier);
  2424. shares.push((x, rcvd.y));
  2425. return SignalCheck::Slashable(shares)
  2426. }
  2427. state.metadata.add_share(epoch_n, rcvd.internal_nullifier, x, rcvd.y);
  2428. SignalCheck::Accepted
  2429. }
  2430. /// Verify a static-DAG event (RLN registration or slashing)
  2431. /// against this event graph's state.
  2432. ///
  2433. /// This is the testable core of the protocol-layer
  2434. /// `handle_static_put`. It performs all checks that the
  2435. /// protocol layer does - bounds, attestation, proof, root
  2436. /// recency - and returns a [`rln::StaticEventCheck`] outcome
  2437. /// telling the caller what to do next.
  2438. ///
  2439. /// **This method does NOT mutate state.** The caller is
  2440. /// responsible for invoking `IdentityState::register` /
  2441. /// `slash` on `Accepted*` outcomes, and for striking the peer
  2442. /// on `Malicious`. Separating decision from action makes the
  2443. /// behaviour fully testable and lets the protocol layer keep
  2444. /// its mutation under a single locked critical section.
  2445. pub async fn rln_verify_static_event(
  2446. &self,
  2447. rln_node: &rln::RLNNode,
  2448. blob: &[u8],
  2449. event_timestamp: u64,
  2450. ) -> rln::StaticEventCheck {
  2451. use darkfi_sdk::crypto::poseidon_hash;
  2452. use rln::{RLNNode, SlashBlob, StaticEventCheck};
  2453. match rln_node {
  2454. RLNNode::Registration(commitment) => {
  2455. // Current admission policy is pregenerated identities only.
  2456. // The guard blob is valid exclusively for commitments supplied
  2457. // by the app config; pairing it with any other commitment is
  2458. // an unambiguous forgery attempt.
  2459. if blob == rln::GENESIS_BLOB_GUARD {
  2460. let repr = commitment.to_repr();
  2461. if self.pregenerated_identity_commitment_reprs.contains(&repr) {
  2462. if self.identity_state.read().await.contains(commitment) {
  2463. return StaticEventCheck::Rejected
  2464. }
  2465. return StaticEventCheck::AcceptedRegistration(*commitment)
  2466. } else {
  2467. return StaticEventCheck::Malicious
  2468. }
  2469. }
  2470. // Free non-pregenerated registration is intentionally disabled:
  2471. // it is a sybil attack surface. Keep the proof scaffolding
  2472. // below for the future staked tier, where acceptance must be
  2473. // backed by a DarkFi smart-contract attestation.
  2474. StaticEventCheck::Rejected
  2475. /*
  2476. #[allow(unreachable_code)]
  2477. let reg: RegistrationBlob = match deserialize_async_partial(blob).await {
  2478. Ok((v, _)) => v,
  2479. Err(_) => return StaticEventCheck::Rejected,
  2480. };
  2481. // Bounds. Out-of-range limits are unambiguous misbehavior.
  2482. if reg.user_message_limit == 0 ||
  2483. reg.user_message_limit > MAX_MSG_LIMIT ||
  2484. reg.max_message_limit != MAX_MSG_LIMIT
  2485. {
  2486. return StaticEventCheck::Malicious
  2487. }
  2488. // Attestation must permit the claimed limit.
  2489. // (Free-tier cap; `Staked` rejected until the stake
  2490. // contract is online.)
  2491. if !reg.attestation.permits(reg.user_message_limit) {
  2492. return StaticEventCheck::Malicious
  2493. }
  2494. // Duplicate registration is a soft Reject (we may
  2495. // have raced a peer), checked here so we don't
  2496. // pay the proof-verification cost for known leaves.
  2497. if self.identity_state.read().await.contains(commitment) {
  2498. return StaticEventCheck::Rejected
  2499. }
  2500. // Proof.
  2501. let pi = vec![
  2502. *commitment,
  2503. pallas::Base::from(reg.user_message_limit),
  2504. pallas::Base::from(reg.max_message_limit),
  2505. ];
  2506. if reg.proof.verify(&self.zk_keys.register_vk, &pi).is_err() {
  2507. return StaticEventCheck::Rejected
  2508. }
  2509. StaticEventCheck::AcceptedRegistration(*commitment)
  2510. */
  2511. }
  2512. RLNNode::Slashing(commitment) => {
  2513. let sl: SlashBlob = match deserialize_async_partial(blob).await {
  2514. Ok((v, _)) => v,
  2515. Err(_) => return StaticEventCheck::Rejected,
  2516. };
  2517. let pi = vec![sl.identity_secret_hash, sl.merkle_root];
  2518. if sl.proof.verify(&self.zk_keys.slash_vk, &pi).is_err() {
  2519. return StaticEventCheck::Rejected
  2520. }
  2521. // Slash event names a specific commitment; verify
  2522. // that the recovered identity_secret_hash actually
  2523. // maps to it. Mismatch is unambiguous misbehavior.
  2524. let rebuilt = poseidon_hash([sl.identity_secret_hash]);
  2525. if *commitment != rebuilt {
  2526. return StaticEventCheck::Malicious
  2527. }
  2528. // The proof's root must be a valid SMT root at the
  2529. // slash event's timestamp. Same logic as signal
  2530. // verification: use the time-window check, which
  2531. // accepts any root that was live within DRIFT of the
  2532. // slash timestamp.
  2533. {
  2534. let id_state = self.identity_state.read().await;
  2535. if !id_state.is_known_root(&sl.merkle_root) {
  2536. drop(id_state);
  2537. match self.is_root_valid_at(&sl.merkle_root, event_timestamp) {
  2538. Ok(true) => {}
  2539. Ok(false) | Err(_) => return StaticEventCheck::Rejected,
  2540. }
  2541. }
  2542. }
  2543. StaticEventCheck::AcceptedSlash(rebuilt)
  2544. }
  2545. }
  2546. }
  2547. /// Insert proof-less pregenerated identity commitments into the
  2548. /// static DAG, called once at startup after the static genesis
  2549. /// event itself is inserted. Idempotent - skips any commitment
  2550. /// already present in the identity tree.
  2551. pub async fn bootstrap_genesis_identities(&self) -> Result<()> {
  2552. // Deterministic for configured pregenerated identities.
  2553. let genesis_event = generate_static_genesis(&self.config);
  2554. let genesis_id = genesis_event.id();
  2555. if !self.static_dag.contains_key(genesis_id.as_bytes())? {
  2556. return Err(Error::Custom("static DAG genesis missing during bootstrap".into()))
  2557. }
  2558. for commitment in self.pregenerated_identity_commitments.iter() {
  2559. if self.identity_state.read().await.contains(commitment) {
  2560. continue
  2561. }
  2562. let rln_node = rln::RLNNode::Registration(*commitment);
  2563. let content = serialize_async(&rln_node).await;
  2564. // Inserted at layer 1
  2565. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  2566. parents[0] = genesis_id;
  2567. let header = Header {
  2568. timestamp: genesis_event.header.timestamp + 1,
  2569. parents,
  2570. layer: 1,
  2571. content_hash: blake3::hash(&content),
  2572. };
  2573. let event = Event { header, content };
  2574. if self.static_dag.contains_key(event.id().as_bytes())? {
  2575. continue
  2576. }
  2577. let blob = rln::GENESIS_BLOB_GUARD.to_vec();
  2578. self.commit_verified_static_event(&event, &blob, &rln_node).await?;
  2579. }
  2580. Ok(())
  2581. }
  2582. }
  2583. async fn request_tips(
  2584. peer: &Channel,
  2585. dag: String,
  2586. timeout: u64,
  2587. ) -> Result<BTreeMap<u64, HashSet<blake3::Hash>>> {
  2588. let sub = peer.subscribe_msg::<TipRep>().await?;
  2589. peer.send(&TipReq(dag)).await?;
  2590. let r = sub
  2591. .receive_with_timeout(timeout)
  2592. .await
  2593. .map_err(|_| Error::EventNotFound("tip timeout".into()))?;
  2594. sub.unsubscribe().await;
  2595. if count_layer_tips(&r.0) > MAX_TIP_REP_TIPS {
  2596. return Err(Error::DagSyncFailed)
  2597. }
  2598. Ok(r.0.clone())
  2599. }
  2600. async fn request_header(
  2601. peer: &Channel,
  2602. name: String,
  2603. tips: LayerUTips,
  2604. timeout: u64,
  2605. ) -> Result<Vec<Header>> {
  2606. let sub = peer.subscribe_msg::<HeaderRep>().await?;
  2607. let tips = cap_layer_tips(&tips, MAX_HEADER_REQ_TIPS);
  2608. peer.send(&HeaderReq(name, tips)).await?;
  2609. let r = sub
  2610. .receive_with_timeout(timeout)
  2611. .await
  2612. .map_err(|_| Error::EventNotFound("hdr timeout".into()))?;
  2613. sub.unsubscribe().await;
  2614. if r.0.len() > MAX_HEADER_REP_HEADERS {
  2615. return Err(Error::DagSyncFailed)
  2616. }
  2617. Ok(r.0.to_vec())
  2618. }
  2619. async fn request_event(
  2620. peer: Arc<Channel>,
  2621. ids: Vec<blake3::Hash>,
  2622. cid: usize,
  2623. timeout: u64,
  2624. ) -> (Result<(Vec<Event>, Vec<Vec<u8>>)>, usize, Arc<Channel>) {
  2625. if ids.len() > MAX_EVENT_REQ_IDS {
  2626. return (Err(Error::DagSyncFailed), cid, peer)
  2627. }
  2628. let sub = match peer.subscribe_msg::<EventRep>().await {
  2629. Ok(s) => s,
  2630. Err(e) => return (Err(e), cid, peer),
  2631. };
  2632. if let Err(e) = peer.send(&EventReq(ids)).await {
  2633. return (Err(e), cid, peer)
  2634. }
  2635. match sub.receive_with_timeout(timeout).await {
  2636. Ok(r) => {
  2637. sub.unsubscribe().await;
  2638. if r.0.len() > MAX_EVENT_REP_EVENTS || r.1.len() > MAX_EVENT_REP_EVENTS {
  2639. return (Err(Error::DagSyncFailed), cid, peer)
  2640. }
  2641. (Ok((r.0.clone(), r.1.clone())), cid, peer)
  2642. }
  2643. Err(_) => (Err(Error::EventNotFound("ev timeout".into())), cid, peer),
  2644. }
  2645. }