lib.rs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2026 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use sled_overlay::sled;
  19. use smol::lock::Mutex;
  20. use std::{collections::HashSet, path::PathBuf};
  21. use darkfi::{
  22. event_graph::EventGraphPtr, net::P2pPtr, rpc::jsonrpc::JsonSubscriber, system::StoppableTaskPtr,
  23. };
  24. use darkfi_serial::{async_trait, SerialDecodable, SerialEncodable};
  25. /// IRC server and client handler implementation
  26. pub mod irc;
  27. use irc::server::IrcServer;
  28. use crate::irc::server::MAX_NICK_LEN;
  29. /// Cryptography utilities
  30. pub mod crypto;
  31. /// Pregenerated DarkIRC RLN identity commitments.
  32. pub mod genesis_commits;
  33. /// JSON-RPC methods
  34. pub mod rpc;
  35. /// Settings utilities
  36. pub mod settings;
  37. /// IRC PRIVMSG
  38. #[derive(Clone, Debug, SerialEncodable, SerialDecodable)]
  39. pub struct Privmsg {
  40. pub version: u8,
  41. pub msg_type: u8,
  42. pub channel: String,
  43. pub nick: String,
  44. pub msg: String,
  45. }
  46. pub struct DarkIrc {
  47. /// P2P network pointer
  48. p2p: P2pPtr,
  49. /// Sled DB (also used in event_graph and for RLN)
  50. sled: sled::Db,
  51. /// Event Graph instance
  52. event_graph: EventGraphPtr,
  53. /// JSON-RPC connection tracker
  54. rpc_connections: Mutex<HashSet<StoppableTaskPtr>>,
  55. /// dnet JSON-RPC subscriber
  56. dnet_sub: JsonSubscriber,
  57. /// deg JSON-RPC subscriber
  58. deg_sub: JsonSubscriber,
  59. /// Gource visualization JSON-RPC subscriber
  60. gource_sub: JsonSubscriber,
  61. /// Replay logs (DB) path
  62. replay_datastore: PathBuf,
  63. }
  64. impl DarkIrc {
  65. pub fn new(
  66. p2p: P2pPtr,
  67. sled: sled::Db,
  68. event_graph: EventGraphPtr,
  69. dnet_sub: JsonSubscriber,
  70. deg_sub: JsonSubscriber,
  71. gource_sub: JsonSubscriber,
  72. replay_datastore: PathBuf,
  73. ) -> Self {
  74. Self {
  75. p2p,
  76. sled,
  77. event_graph,
  78. rpc_connections: Mutex::new(HashSet::new()),
  79. dnet_sub,
  80. deg_sub,
  81. gource_sub,
  82. replay_datastore,
  83. }
  84. }
  85. }
  86. pub fn pad(string: &str) -> Vec<u8> {
  87. let mut bytes = string.as_bytes().to_vec();
  88. bytes.resize(MAX_NICK_LEN, 0x00);
  89. bytes
  90. }
  91. pub fn unpad(vec: &mut Vec<u8>) {
  92. if let Some(i) = vec.iter().rposition(|x| *x != 0) {
  93. let new_len = i + 1;
  94. vec.truncate(new_len);
  95. }
  96. }
  97. #[cfg(test)]
  98. mod tests {
  99. use std::{
  100. collections::HashMap,
  101. sync::{
  102. atomic::{AtomicU16, Ordering},
  103. Arc,
  104. },
  105. };
  106. use darkfi::{
  107. event_graph::{
  108. proto::{ProtocolEventGraph, RangeCursor, SyncDirection},
  109. Event, EventGraph, EventGraphConfig, EventGraphPtr, Header, NULL_PARENTS,
  110. },
  111. net::{
  112. session::SESSION_DEFAULT,
  113. settings::{NetworkProfile, Settings},
  114. P2p, P2pPtr,
  115. },
  116. system::sleep,
  117. };
  118. use darkfi_serial::{deserialize_async_partial, serialize_async};
  119. use easy_parallel::Parallel;
  120. use sled_overlay::sled;
  121. use smol::{channel, future, Executor};
  122. use url::Url;
  123. use super::Privmsg;
  124. struct HistoryNode {
  125. p2p: P2pPtr,
  126. event_graph: EventGraphPtr,
  127. }
  128. fn alloc_port_base() -> u16 {
  129. static NEXT: AtomicU16 = AtomicU16::new(24_400);
  130. NEXT.fetch_add(2, Ordering::SeqCst)
  131. }
  132. fn history_test_config() -> EventGraphConfig {
  133. EventGraphConfig {
  134. initial_genesis: 1_704_067_200_000,
  135. hours_rotation: 1,
  136. genesis_contents: b"darkirc-mobile-history-test".to_vec(),
  137. rln_enabled: false,
  138. pregenerated_identity_commitments: Vec::new(),
  139. max_dags: Some(5),
  140. }
  141. }
  142. async fn spawn_history_node(
  143. port_base: u16,
  144. port_offset: u16,
  145. peer_offsets: &[u16],
  146. ex: Arc<Executor<'static>>,
  147. ) -> HistoryNode {
  148. let mut profiles = HashMap::new();
  149. profiles.insert(
  150. "tcp".to_string(),
  151. NetworkProfile { outbound_connect_timeout: 2, ..Default::default() },
  152. );
  153. let inbound =
  154. vec![Url::parse(&format!("tcp://127.0.0.1:{}", port_base + port_offset)).unwrap()];
  155. let peers = peer_offsets
  156. .iter()
  157. .map(|offset| Url::parse(&format!("tcp://127.0.0.1:{}", port_base + *offset)).unwrap())
  158. .collect();
  159. let settings = Settings {
  160. localnet: true,
  161. inbound_addrs: inbound,
  162. outbound_connections: 0,
  163. inbound_connections: usize::MAX,
  164. peers,
  165. active_profiles: vec!["tcp".to_string()],
  166. profiles,
  167. ..Default::default()
  168. };
  169. let p2p = P2p::new(settings, ex.clone()).await.unwrap();
  170. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  171. let event_graph =
  172. EventGraph::new(p2p.clone(), sled_db, "/tmp".into(), false, history_test_config(), ex)
  173. .await
  174. .unwrap();
  175. event_graph.synced.store(true, Ordering::Release);
  176. let event_graph_weak = Arc::downgrade(&event_graph);
  177. p2p.protocol_registry()
  178. .register(SESSION_DEFAULT, move |channel, _| {
  179. let event_graph_weak = event_graph_weak.clone();
  180. async move {
  181. let event_graph = event_graph_weak
  182. .upgrade()
  183. .expect("EventGraph dropped before protocol factory invoked");
  184. ProtocolEventGraph::init(event_graph, channel).await.unwrap()
  185. }
  186. })
  187. .await;
  188. HistoryNode { p2p, event_graph }
  189. }
  190. async fn make_history_network(ex: Arc<Executor<'static>>) -> Vec<HistoryNode> {
  191. let port_base = alloc_port_base();
  192. let nodes = vec![
  193. spawn_history_node(port_base, 0, &[1], ex.clone()).await,
  194. spawn_history_node(port_base, 1, &[0], ex).await,
  195. ];
  196. for node in &nodes {
  197. node.p2p.clone().start().await.unwrap();
  198. }
  199. sleep(5).await;
  200. nodes
  201. }
  202. async fn shutdown_history_network(nodes: &[HistoryNode]) {
  203. for node in nodes {
  204. node.p2p.stop().await;
  205. }
  206. }
  207. fn run_history_test<F, Fut>(body: F)
  208. where
  209. F: FnOnce(Arc<Executor<'static>>) -> Fut,
  210. Fut: std::future::Future<Output = ()>,
  211. {
  212. let ex = Arc::new(Executor::new());
  213. let ex_ = ex.clone();
  214. let (signal, shutdown) = channel::unbounded::<()>();
  215. Parallel::new().each(0..2, |_| future::block_on(ex.run(shutdown.recv()))).finish(|| {
  216. future::block_on(async {
  217. body(ex_).await;
  218. drop(signal);
  219. })
  220. });
  221. }
  222. fn genesis_for(event_graph: &EventGraphPtr, dag_ts: u64) -> Event {
  223. let content = event_graph.config.genesis_contents.clone();
  224. Event {
  225. header: Header {
  226. timestamp: dag_ts,
  227. parents: NULL_PARENTS,
  228. layer: 0,
  229. content_hash: blake3::hash(&content),
  230. },
  231. content,
  232. }
  233. }
  234. async fn select_current_dag(event_graph: &EventGraphPtr, dag_ts: u64) {
  235. *event_graph.current_genesis.write().await = genesis_for(event_graph, dag_ts);
  236. }
  237. async fn append_chat_event(
  238. event_graph: &EventGraphPtr,
  239. dag_ts: u64,
  240. timestamp: u64,
  241. nick: &str,
  242. message: &str,
  243. ) -> Event {
  244. select_current_dag(event_graph, dag_ts).await;
  245. let privmsg = Privmsg {
  246. version: 0,
  247. msg_type: 0,
  248. channel: "#mobile".to_string(),
  249. nick: nick.to_string(),
  250. msg: message.to_string(),
  251. };
  252. let event = Event::with_timestamp(timestamp, serialize_async(&privmsg).await, event_graph)
  253. .await
  254. .unwrap();
  255. event_graph.insert_signal_with_blob(&event, &[], &dag_ts.to_string()).await.unwrap();
  256. event
  257. }
  258. async fn page_messages(events: &[Event]) -> Vec<String> {
  259. let mut messages = Vec::with_capacity(events.len());
  260. for event in events {
  261. let (privmsg, _) = deserialize_async_partial::<Privmsg>(event.content()).await.unwrap();
  262. messages.push(privmsg.msg);
  263. }
  264. messages
  265. }
  266. #[test]
  267. fn mobile_chat_can_scroll_backwards_across_multiple_dags() {
  268. run_history_test(|ex| async move {
  269. const HOUR_MS: u64 = 3_600_000;
  270. const DAG_COUNT: usize = 5;
  271. const EVENTS_PER_DAG: usize = 4;
  272. const PAGE_SIZE: usize = 3;
  273. let nodes = make_history_network(ex).await;
  274. let source = nodes[0].event_graph.clone();
  275. let mobile = nodes[1].event_graph.clone();
  276. let latest_dag = source.current_genesis.read().await.header.timestamp;
  277. let mut dag_timestamps = Vec::with_capacity(DAG_COUNT);
  278. for offset in (0..DAG_COUNT).rev() {
  279. dag_timestamps.push(latest_dag.saturating_sub((offset as u64) * HOUR_MS));
  280. }
  281. let mut seeded = Vec::new();
  282. let mut expected_scrollback = Vec::new();
  283. for (dag_index, dag_ts) in dag_timestamps.iter().enumerate() {
  284. let mut dag_events = Vec::new();
  285. for event_index in 0..EVENTS_PER_DAG {
  286. let message = format!("dag {dag_index} history message {event_index}");
  287. let event = append_chat_event(
  288. &source,
  289. *dag_ts,
  290. dag_ts + 10_000 + event_index as u64,
  291. "alice",
  292. &message,
  293. )
  294. .await;
  295. dag_events.push(event);
  296. }
  297. expected_scrollback.extend(
  298. (0..EVENTS_PER_DAG).rev().map(|event_index| {
  299. format!("dag {dag_index} history message {event_index}")
  300. }),
  301. );
  302. seeded.push(dag_events);
  303. }
  304. expected_scrollback =
  305. expected_scrollback.chunks(EVENTS_PER_DAG).rev().flatten().cloned().collect();
  306. for dag_events in &seeded {
  307. for event in dag_events {
  308. assert!(mobile.fetch_event_from_dags(&event.id()).await.unwrap().is_none());
  309. }
  310. }
  311. mobile.sync_selected_headers(1).await.unwrap();
  312. let mut collected = Vec::new();
  313. for dag_ts in dag_timestamps.iter().rev() {
  314. let mut cursor = RangeCursor::newest();
  315. loop {
  316. let page = mobile
  317. .dag_sync_range(*dag_ts, cursor, SyncDirection::Backward, PAGE_SIZE)
  318. .await
  319. .unwrap();
  320. collected.extend(page_messages(&page.events).await);
  321. if page.exhausted {
  322. break
  323. }
  324. cursor = page.next_cursor;
  325. }
  326. }
  327. assert_eq!(collected, expected_scrollback);
  328. for dag_events in &seeded {
  329. for event in dag_events {
  330. assert!(mobile.fetch_event_from_dags(&event.id()).await.unwrap().is_some());
  331. }
  332. }
  333. shutdown_history_network(&nodes).await;
  334. });
  335. }
  336. }