Skip to main content

cuprated/
lib.rs

1//! `cuprated` library.
2//!
3//! Call [`Node::launch`] to initialize and run the node. Returns a [`Node`]
4//! with handles to node services.
5//!
6//! # Example
7//!
8//! ```ignore
9//! use cuprated::{config::Config, Node};
10//!
11//! let config = Config::read_from_path("cuprated.toml")?;
12//! cuprated::logging::init_logging(&config);
13//!
14//! let mut node = Node::launch(config).await;
15//! let height = node.blockchain.context().chain_height;
16//! ```
17
18#![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/// Captures the necessary context for launching the node.
64///
65/// A field belongs here if it is `Clone + Send + Sync`, used by
66/// multiple subsystems, and available before subsystem init begins.
67/// Write handles, single-consumer channels, `!Sync` types,
68/// and late-constructed services do _not_ belong here.
69#[derive(Clone)]
70pub(crate) struct LaunchContext {
71    /// The configuration this node was launched with.
72    pub config: Arc<Config>,
73
74    /// Reorg lock.
75    ///
76    /// A [`RwLock`] where a write lock is taken during a reorg and a read lock can be taken
77    /// for any operation which must complete without a reorg happening.
78    ///
79    /// Currently, the only operation that needs to take a read lock is adding txs to the tx-pool,
80    /// this can potentially be removed in the future, see: <https://github.com/Cuprate/cuprate/issues/305>
81    pub reorg_lock: Arc<RwLock<()>>,
82
83    /// Interface to the blockchain (database reads, cached state, mutations).
84    pub blockchain: BlockchainInterface,
85
86    /// Read handle to the transaction pool.
87    pub txpool_read: TxpoolReadHandle,
88
89    /// Task spawning and shutdown coordination.
90    pub task_executor: TaskExecutor,
91}
92
93/// An active `cuprated` node.
94///
95/// Returned by [`Node::launch`]. Use this to interact with the running node.
96#[must_use]
97pub struct Node {
98    /// Interface to the blockchain.
99    pub blockchain: BlockchainInterface,
100
101    /// Transaction pool queries.
102    pub txpool: TxpoolReadHandle,
103
104    /// Clearnet P2P interface.
105    pub clearnet: NetworkInterface<ClearNet>,
106
107    /// Tor P2P interface (available after sync).
108    pub tor: Option<oneshot::Receiver<NetworkInterface<Tor>>>,
109
110    /// The configuration this node was launched with.
111    pub config: Arc<Config>,
112
113    /// Task spawning and shutdown executor.
114    pub task_executor: TaskExecutor,
115}
116
117impl Drop for Node {
118    fn drop(&mut self) {
119        self.shutdown();
120    }
121}
122
123impl Node {
124    /// Launch a new `cuprated` process.
125    ///
126    /// Sets up thread pools, databases, P2P networking, the blockchain manager,
127    /// and RPC servers.
128    ///
129    /// The caller should set up the following before calling this:
130    /// - Tracing/logging (the node emits tracing events during initialization)
131    /// - Global rayon thread pool (optional, uses rayon defaults if not set)
132    /// - Memory resolution (call [`resolve_max_memory`](crate::config::resolve_max_memory))
133    ///
134    /// # Errors
135    ///
136    /// Returns an error if the database is corrupt, critical services fail to start,
137    /// or `target_max_memory` is unresolved.
138    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            // Initialize the database thread pool.
145            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            // Start the blockchain & tx-pool databases.
153            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            // TODO: Add an argument/option for keeping alt blocks between restart.
171            blockchain_write_handle
172                .ready()
173                .await?
174                .call(BlockchainWriteRequest::FlushAltBlocks)
175                .await?;
176
177            // Check add the genesis block to the blockchain.
178            blockchain::check_add_genesis(
179                &mut blockchain_read_handle,
180                &mut blockchain_write_handle,
181                config.network(),
182            )
183            .await?;
184
185            // Start the context service and the block/tx verifier.
186            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            // Create the blockchain syncer handle and synced signal sender.
192            let (blockchain_syncer_handle, synced_tx) = BlockchainSyncerHandle::new();
193
194            // Create the blockchain manager handle and command receiver.
195            let (blockchain_manager_handle, command_rx) = BlockchainManagerHandle::new();
196
197            // Create the blockchain interface.
198            let blockchain_interface = BlockchainInterface::new(
199                blockchain_read_handle,
200                context_svc,
201                blockchain_manager_handle,
202                blockchain_syncer_handle,
203            );
204
205            // Create the launch context.
206            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            // Bootstrap or configure Tor if enabled.
215            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            // Start clearnet P2P zone
219            let (clearnet_interface, clearnet_tx_handler_subscriber) =
220                p2p::initialize_clearnet_p2p(&launch_ctx, &tor_context).await?;
221
222            // Create Tor router delivery channel.
223            let (tor_router_tx, tor_router_rx) = tor_enabled.then(oneshot::channel).unzip();
224
225            // Create the incoming tx handler service.
226            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            // Send tx handler sender to clearnet zone, offline zones have no subscriber.
235            if let Some(subscriber) = clearnet_tx_handler_subscriber {
236                if subscriber.send(tx_handler.clone()).is_err() {
237                    unreachable!()
238                }
239            }
240
241            // Tor interface channel - populated when Tor starts after sync.
242            let (tor_tx, tor_rx) = oneshot::channel();
243
244            // Initialize the blockchain manager.
245            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            // Initialize the RPC server(s).
256            rpc::init_rpc_servers(&launch_ctx, &tx_handler).await?;
257
258            // Start Tor P2P zone after sync completes.
259            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    /// Trigger a graceful shutdown.
297    pub fn shutdown(&self) {
298        self.task_executor.trigger_shutdown();
299    }
300}