1#![doc = include_str!("../README.md")]
19#![cfg_attr(docsrs, feature(doc_cfg))]
20
21#[cfg(feature = "mimalloc")]
22use mimalloc as _;
23
24#[cfg(all(
25 feature = "jemalloc",
26 not(any(target_env = "msvc", target_os = "freebsd"))
27))]
28use tikv_jemallocator as _;
29
30pub mod blockchain;
31pub mod config;
32pub mod constants;
33pub mod logging;
34pub mod monitor;
35pub mod version;
36
37mod p2p;
38mod rpc;
39mod tor;
40mod txpool;
41
42use std::sync::Arc;
43
44use anyhow::Context;
45use tokio::sync::{oneshot, RwLock};
46use tower::{Service, ServiceExt};
47use tracing::error;
48
49use cuprate_p2p::NetworkInterface;
50use cuprate_p2p_core::{ClearNet, Tor};
51use cuprate_txpool::service::TxpoolReadHandle;
52use cuprate_types::blockchain::BlockchainWriteRequest;
53
54use crate::{
55 blockchain::{BlockchainInterface, BlockchainManagerHandle, BlockchainSyncerHandle},
56 config::Config,
57 constants::DATABASE_CORRUPT_MSG,
58 monitor::TaskExecutor,
59 tor::initialize_tor_if_enabled,
60 txpool::IncomingTxHandler,
61};
62
63#[derive(Clone)]
70pub(crate) struct LaunchContext {
71 pub config: Arc<Config>,
73
74 pub reorg_lock: Arc<RwLock<()>>,
82
83 pub blockchain: BlockchainInterface,
85
86 pub txpool_read: TxpoolReadHandle,
88
89 pub task_executor: TaskExecutor,
91}
92
93#[must_use]
97pub struct Node {
98 pub blockchain: BlockchainInterface,
100
101 pub txpool: TxpoolReadHandle,
103
104 pub clearnet: NetworkInterface<ClearNet>,
106
107 pub tor: Option<oneshot::Receiver<NetworkInterface<Tor>>>,
109
110 pub config: Arc<Config>,
112
113 pub task_executor: TaskExecutor,
115}
116
117impl Drop for Node {
118 fn drop(&mut self) {
119 self.shutdown();
120 }
121}
122
123impl Node {
124 pub async fn launch(config: impl Into<Arc<Config>>) -> Result<Self, anyhow::Error> {
139 let config: Arc<Config> = config.into();
140 let task_executor = TaskExecutor::new();
141 let launch_executor = task_executor.clone();
142
143 let node: Result<Self, anyhow::Error> = async move {
144 let db_thread_pool = Arc::new(
146 rayon::ThreadPoolBuilder::new()
147 .num_threads(config.storage.reader_threads)
148 .build()
149 .context("failed to build rayon database thread pool")?,
150 );
151
152 let fjall_db = fjall::Database::builder(config.fjall_directory())
154 .cache_size(config.fjall_cache_size())
155 .open()
156 .context(DATABASE_CORRUPT_MSG)?;
157
158 let (mut blockchain_read_handle, mut blockchain_write_handle, _) =
159 cuprate_blockchain::service::init_with_pool(
160 &config.blockchain_config(),
161 fjall_db.clone(),
162 Arc::clone(&db_thread_pool),
163 )
164 .context(DATABASE_CORRUPT_MSG)?;
165
166 let (txpool_read_handle, txpool_write_handle) =
167 cuprate_txpool::service::init_with_pool(fjall_db, db_thread_pool)
168 .context(DATABASE_CORRUPT_MSG)?;
169
170 blockchain_write_handle
172 .ready()
173 .await?
174 .call(BlockchainWriteRequest::FlushAltBlocks)
175 .await?;
176
177 blockchain::check_add_genesis(
179 &mut blockchain_read_handle,
180 &mut blockchain_write_handle,
181 config.network(),
182 )
183 .await?;
184
185 let context_svc =
187 blockchain::init_consensus(blockchain_read_handle.clone(), config.context_config())
188 .await
189 .map_err(anyhow::Error::from_boxed)?;
190
191 let (blockchain_syncer_handle, synced_tx) = BlockchainSyncerHandle::new();
193
194 let (blockchain_manager_handle, command_rx) = BlockchainManagerHandle::new();
196
197 let blockchain_interface = BlockchainInterface::new(
199 blockchain_read_handle,
200 context_svc,
201 blockchain_manager_handle,
202 blockchain_syncer_handle,
203 );
204
205 let launch_ctx = LaunchContext {
207 config,
208 reorg_lock: Arc::new(RwLock::new(())),
209 blockchain: blockchain_interface,
210 txpool_read: txpool_read_handle,
211 task_executor,
212 };
213
214 let tor_enabled = !launch_ctx.config.offline && launch_ctx.config.p2p.tor_net.enabled;
216 let tor_context = initialize_tor_if_enabled(&launch_ctx).await;
217
218 let (clearnet_interface, clearnet_tx_handler_subscriber) =
220 p2p::initialize_clearnet_p2p(&launch_ctx, &tor_context).await?;
221
222 let (tor_router_tx, tor_router_rx) = tor_enabled.then(oneshot::channel).unzip();
224
225 let tx_handler = IncomingTxHandler::init(
227 &launch_ctx,
228 clearnet_interface.clone(),
229 tor_router_rx,
230 txpool_write_handle,
231 )
232 .await?;
233
234 if let Some(subscriber) = clearnet_tx_handler_subscriber {
236 if subscriber.send(tx_handler.clone()).is_err() {
237 unreachable!()
238 }
239 }
240
241 let (tor_tx, tor_rx) = oneshot::channel();
243
244 blockchain::init_blockchain_manager(
246 &launch_ctx,
247 clearnet_interface.clone(),
248 blockchain_write_handle,
249 tx_handler.txpool_manager.clone(),
250 synced_tx,
251 command_rx,
252 )
253 .await?;
254
255 rpc::init_rpc_servers(&launch_ctx, &tx_handler).await?;
257
258 if tor_enabled {
260 p2p::initialize_tor_p2p(
261 launch_ctx.clone(),
262 tor_context,
263 tx_handler,
264 tor_tx,
265 tor_router_tx,
266 );
267 }
268
269 let LaunchContext {
270 blockchain,
271 txpool_read,
272 config,
273 task_executor,
274 ..
275 } = launch_ctx;
276
277 Ok(Self {
278 blockchain,
279 txpool: txpool_read,
280 clearnet: clearnet_interface,
281 tor: if tor_enabled { Some(tor_rx) } else { None },
282 config,
283 task_executor,
284 })
285 }
286 .await;
287
288 if let Err(e) = &node {
289 error!("Failed to launch node: {e:#}");
290 launch_executor.cancellation_token().cancel();
291 drop(launch_executor.wait_for_shutdown().await);
292 }
293 node
294 }
295
296 pub fn shutdown(&self) {
298 self.task_executor.trigger_shutdown();
299 }
300}