1pub mod bundle;
5mod context;
6pub mod db_util;
7pub mod main;
8
9use crate::blocks::TipsetKey;
10use crate::chain::ChainStore;
11use crate::chain_sync::ChainFollower;
12use crate::chain_sync::network_context::SyncNetworkContext;
13use crate::cli_shared::snapshot;
14use crate::cli_shared::{
15 chain_path,
16 cli::{CliOpts, Config},
17 delete_chain_data,
18};
19use crate::daemon::{context::AppContext, db_util::import_chain_as_forest_car};
20use crate::db::gc::SnapshotGarbageCollector;
21use crate::db::ttl::EthMappingCollector;
22use crate::libp2p::{Libp2pService, PeerManager};
23use crate::message_pool::{MessagePool, MpoolConfig, MpoolLocker, NonceTracker};
24use crate::networks::{self, ChainConfig};
25use crate::prelude::*;
26use crate::rpc::RPCState;
27use crate::rpc::eth::filter::EthEventHandler;
28use crate::rpc::start_rpc;
29use crate::shim::address::Address;
30use crate::shim::clock::ChainEpoch;
31use crate::shim::state_tree::StateTree;
32use crate::shim::version::NetworkVersion;
33use crate::state_manager::StateManager;
34use crate::utils::misc::env::is_env_truthy;
35use crate::utils::{self};
36use crate::utils::{proofs_api::ensure_proof_params_downloaded, version::FOREST_VERSION_STRING};
37use anyhow::{Context as _, bail};
38use backon::{ExponentialBuilder, Retryable};
39use dialoguer::theme::ColorfulTheme;
40use futures::{Future, FutureExt};
41use std::path::Path;
42use std::sync::Arc;
43use std::sync::OnceLock;
44use std::time::{Duration, Instant};
45use tokio::sync::broadcast::error::RecvError;
46use tokio::{
47 signal::{
48 ctrl_c,
49 unix::{SignalKind, signal},
50 },
51 sync::mpsc,
52 task::JoinSet,
53};
54use tokio_util::sync::CancellationToken;
55use tracing::{debug, info, warn};
56
57pub static GLOBAL_SNAPSHOT_GC: OnceLock<Arc<SnapshotGarbageCollector>> = OnceLock::new();
58
59fn maybe_increase_fd_limit() -> anyhow::Result<()> {
63 static DESIRED_SOFT_LIMIT: u64 = 8192;
64 let (soft_before, _) = rlimit::Resource::NOFILE.get()?;
65
66 let soft_after = rlimit::increase_nofile_limit(DESIRED_SOFT_LIMIT)?;
67 if soft_before < soft_after {
68 debug!("Increased file descriptor limit from {soft_before} to {soft_after}");
69 }
70 if soft_after < DESIRED_SOFT_LIMIT {
71 warn!(
72 "File descriptor limit is too low: {soft_after} < {DESIRED_SOFT_LIMIT}. \
73 You may encounter 'too many open files' errors.",
74 );
75 }
76
77 Ok(())
78}
79
80pub async fn start_interruptable(opts: CliOpts, config: Config) -> anyhow::Result<()> {
82 let start_time = chrono::Utc::now();
83 let mut terminate = signal(SignalKind::terminate())?;
84 let (shutdown_send, mut shutdown_recv) = mpsc::channel(1);
85 let (rpc_stop_handle, rpc_server_handle) = jsonrpsee::server::stop_channel();
86 let result = tokio::select! {
87 ret = start(start_time, opts, config, shutdown_send, rpc_stop_handle) => ret,
88 _ = ctrl_c() => {
89 info!("Keyboard interrupt.");
90 Ok(())
91 },
92 _ = terminate.recv() => {
93 info!("Received SIGTERM.");
94 Ok(())
95 },
96 _ = shutdown_recv.recv() => {
97 info!("Client requested a shutdown.");
98 Ok(())
99 },
100 };
101 _ = rpc_server_handle.stop();
102 crate::utils::io::terminal_cleanup();
103 result
104}
105
106fn startup_init(config: &Config) -> anyhow::Result<()> {
111 maybe_increase_fd_limit()?;
112 crate::utils::proofs_api::maybe_set_proofs_parameter_cache_dir_env(&config.client.data_dir);
115 info!(
116 "Starting Forest daemon, version {}",
117 FOREST_VERSION_STRING.as_str()
118 );
119 info!("Using data directory: {}", config.client.data_dir.display());
120 Ok(())
121}
122
123async fn maybe_import_snapshot(
124 opts: &CliOpts,
125 config: &mut Config,
126 ctx: &AppContext,
127) -> anyhow::Result<()> {
128 let chain_config = ctx.state_manager.chain_config();
129 if config.client.snapshot_path.is_none() && !opts.stateless {
131 maybe_set_snapshot_path(
132 config,
133 chain_config,
134 ctx.state_manager.chain_store().heaviest_tipset().epoch(),
135 opts.auto_download_snapshot,
136 &ctx.db_meta_data.get_root_dir(),
137 )
138 .await?;
139 }
140
141 let snapshot_tracker = ctx.snapshot_progress_tracker.clone();
142 if !opts.skip_load.unwrap_or_default()
144 && let Some(path) = &config.client.snapshot_path
145 {
146 let (car_db_path, ts) = import_chain_as_forest_car(
147 path,
148 &ctx.db_meta_data.get_forest_car_db_dir(),
149 config.client.import_mode,
150 config.client.rpc_v1_endpoint()?,
151 &crate::f3::get_f3_root(config),
152 ctx.chain_config(),
153 &snapshot_tracker,
154 )
155 .await?;
156 ctx.db
157 .read_only_files(std::iter::once(car_db_path.clone()))?;
158 let ts_epoch = ts.epoch();
159 ctx.state_manager.chain_store().set_heaviest_tipset(ts)?;
162 debug!(
163 "Loaded car DB at {} and set current head to epoch {ts_epoch}",
164 car_db_path.display(),
165 );
166 }
167
168 if !snapshot_tracker.is_completed() {
171 snapshot_tracker.not_required();
172 }
173
174 if let Some(validate_from) = config.client.snapshot_height {
175 ensure_proof_params_downloaded().await?;
177 let current_height = config
179 .client
180 .snapshot_head
181 .unwrap_or_else(|| ctx.state_manager.chain_store().heaviest_tipset().epoch());
182
183 let validation_range = validation_range(current_height, validate_from)?;
184 let state_manager = ctx.state_manager.shallow_clone();
187 tokio::task::spawn_blocking(move || {
188 state_manager.validate_range_blocking(validation_range)
189 })
190 .await??;
191 }
192
193 Ok(())
194}
195
196fn validation_range(
199 current: ChainEpoch,
200 from: ChainEpoch,
201) -> anyhow::Result<std::ops::RangeInclusive<ChainEpoch>> {
202 anyhow::ensure!(
203 current.is_positive(),
204 "current head epoch {current} is invalid"
205 );
206
207 let start = if from.is_negative() {
211 current.saturating_add(from).max(0)
212 } else {
213 from
214 };
215
216 anyhow::ensure!(
219 start <= current,
220 "requested validation start epoch {start} is beyond the current head at epoch {current}",
221 );
222
223 Ok(start..=current)
224}
225
226async fn maybe_start_metrics_service(
227 services: &mut JoinSet<anyhow::Result<()>>,
228 config: &Config,
229 ctx: &AppContext,
230) -> anyhow::Result<()> {
231 if config.client.enable_metrics_endpoint {
232 let prometheus_listener =
233 crate::utils::net::bind_tcp_listener(config.client.metrics_address, 0).await?;
234 info!(
235 "Prometheus server started at {}",
236 config.client.metrics_address
237 );
238 let db_directory = crate::db::db_engine::db_root(&chain_path(config))?;
239 let db = ctx.db.writer().clone();
240
241 let get_chain_head_height = Arc::new({
242 let cs = ctx.state_manager.chain_store().shallow_clone();
243 move || cs.heaviest_tipset().epoch()
244 });
245 let get_chain_head_actor_version = Arc::new({
246 let cs = ctx.state_manager.chain_store().shallow_clone();
247 move || {
248 if let Ok(state) =
249 StateTree::new_from_root(cs.db(), cs.heaviest_tipset().parent_state())
250 && let Ok(bundle_meta) = state.get_actor_bundle_metadata()
251 && let Ok(actor_version) = bundle_meta.actor_major_version()
252 {
253 actor_version
254 } else {
255 0
256 }
257 }
258 });
259 services.spawn({
260 let chain_config = ctx.chain_config().clone();
261 let get_chain_head_height = get_chain_head_height.clone();
262 async {
263 crate::metrics::init_prometheus(
264 prometheus_listener,
265 db_directory,
266 db,
267 chain_config,
268 get_chain_head_height,
269 get_chain_head_actor_version,
270 )
271 .await
272 .context("Failed to initiate prometheus server")
273 }
274 });
275
276 crate::metrics::register_collector(Box::new(
277 networks::metrics::NetworkHeightCollector::new(
278 ctx.state_manager.chain_config().block_delay_secs,
279 ctx.state_manager
280 .chain_store()
281 .genesis_block_header()
282 .timestamp,
283 get_chain_head_height,
284 ),
285 ));
286 }
287 Ok(())
288}
289
290async fn create_p2p_service(
291 services: &mut JoinSet<anyhow::Result<()>>,
292 config: &mut Config,
293 ctx: &AppContext,
294) -> anyhow::Result<Libp2pService> {
295 if config.network.bootstrap_peers.is_empty() {
297 config.network.bootstrap_peers = ctx.state_manager.chain_config().bootstrap_peers.clone();
298 }
299
300 let peer_manager = Arc::new(PeerManager::default());
301 services.spawn(peer_manager.clone().peer_operation_event_loop_task());
302 let p2p_service = Libp2pService::new(
304 config.network.clone(),
305 ctx.state_manager.chain_store().shallow_clone(),
306 peer_manager.clone(),
307 ctx.net_keypair.clone(),
308 config.chain.genesis_name(),
309 *ctx.state_manager.chain_store().genesis_block_header().cid(),
310 )
311 .await?;
312 Ok(p2p_service)
313}
314
315fn create_mpool(
316 services: &mut JoinSet<anyhow::Result<()>>,
317 p2p_service: &Libp2pService,
318 ctx: &AppContext,
319) -> anyhow::Result<MessagePool<ChainStore>> {
320 Ok(MessagePool::new(
321 ctx.state_manager.chain_store().shallow_clone(),
322 p2p_service.network_sender().clone(),
323 MpoolConfig::load_config(ctx.db.writer().as_ref())?,
324 ctx.state_manager.chain_config().clone(),
325 services,
326 )?)
327}
328
329fn create_chain_follower(
330 opts: &CliOpts,
331 p2p_service: &Libp2pService,
332 mpool: MessagePool<ChainStore>,
333 ctx: &AppContext,
334) -> anyhow::Result<ChainFollower> {
335 let network_send = p2p_service.network_sender().clone();
336 let peer_manager = p2p_service.peer_manager().clone();
337 let network = SyncNetworkContext::new(network_send, peer_manager, ctx.db.clone().into());
338 Ok(ChainFollower::new(
339 ctx.state_manager.shallow_clone(),
340 network,
341 ctx.state_manager.chain_store().genesis_tipset(),
342 p2p_service.network_receiver(),
343 opts.stateless,
344 mpool,
345 ))
346}
347
348fn start_chain_follower_service(
349 services: &mut JoinSet<anyhow::Result<()>>,
350 opts: &CliOpts,
351 config: &Config,
352 chain_follower: ChainFollower,
353) {
354 services.spawn({
355 let chain_follower = chain_follower.shallow_clone();
356 async move { chain_follower.run().await }
357 });
358 maybe_prefill_rpc_caches(services, opts, config, chain_follower);
359}
360
361fn maybe_prefill_rpc_caches(
362 services: &mut JoinSet<anyhow::Result<()>>,
363 opts: &CliOpts,
364 config: &Config,
365 chain_follower: ChainFollower,
366) {
367 if config.client.enable_rpc && !opts.stateless {
369 let sync_status = chain_follower.sync_status.shallow_clone();
370 let state_manager = chain_follower.state_manager.shallow_clone();
371 let mut validated_tipset_rx = chain_follower.subscribe_validated_tipset();
372 services.spawn(async move {
373 let cancellation_token = CancellationToken::new();
374 let _cancellation_token_drop_guard = cancellation_token.drop_guard_ref();
375 loop {
376 match validated_tipset_rx.recv().await {
377 Ok(_) if !sync_status.load().is_synced() => {
378 continue;
380 }
381 Ok(tsk) => {
382 let state_manager = state_manager.shallow_clone();
383 let cancellation_token = cancellation_token.clone();
384 tokio::spawn(async move {
385 cancellation_token
386 .run_until_cancelled(prefill_rpc_caches_for_tipset(
387 state_manager,
388 tsk,
389 ))
390 .await
391 });
392 }
393 Err(RecvError::Lagged(n)) => {
394 warn!("validated tipset broadcast lagged: skipped {n} tipsets")
395 }
396 Err(RecvError::Closed) => break Ok(()),
397 }
398 }
399 });
400 }
401}
402
403async fn prefill_rpc_caches_for_tipset(state_manager: StateManager, tsk: TipsetKey) {
404 match state_manager.chain_index().load_required_tipset(&tsk) {
405 Ok(ts) => {
406 {
407 if let Err(e) = state_manager.load_executed_tipset(&ts).await {
409 warn!("failed to load executed tipset for cache warmup: {e:#}");
410 return; }
412 }
413 for tx_info in [crate::rpc::eth::TxInfo::Full, crate::rpc::eth::TxInfo::Hash] {
414 if let Err(e) = crate::rpc::eth::Block::from_filecoin_tipset(
415 &state_manager,
416 ts.shallow_clone(),
417 tx_info,
418 )
419 .await
420 {
421 warn!("failed to call `Block::from_filecoin_tipset` for cache warmup: {e:#}");
422 }
423 }
424 {
425 if let Err(e) = crate::rpc::eth::eth_trace_block(&state_manager, &ts).await {
428 warn!("failed to call `eth_trace_block` for cache warmup: {e:#}");
429 }
430 }
431 {
432 use crate::rpc::eth::filter::{Matcher, SkipEvent};
433 struct CollectEventsCachePrefillingMatcher;
434 impl Matcher for CollectEventsCachePrefillingMatcher {
435 fn msg_cid_filter(&self) -> Option<&Cid> {
436 None
437 }
438 fn matches(
439 &self,
440 _: &Address,
441 _: &[crate::shim::executor::Entry],
442 ) -> anyhow::Result<bool> {
443 Ok(false)
444 }
445 }
446 let mut collected_events = vec![];
447 if let Err(e) = EthEventHandler::collect_events(
448 &state_manager,
449 &ts,
450 Some(&CollectEventsCachePrefillingMatcher),
451 SkipEvent::OnUnresolvedAddress,
452 &mut collected_events,
453 )
454 .await
455 {
456 warn!("failed to collect events for cache warmup: {e:#}");
457 }
458 }
459 }
460 Err(e) => {
461 warn!("failed to load tipset for cache warmup: {e:#}");
462 }
463 }
464}
465
466async fn maybe_start_health_check_service(
467 services: &mut JoinSet<anyhow::Result<()>>,
468 config: &Config,
469 p2p_service: &Libp2pService,
470 chain_follower: &ChainFollower,
471 ctx: &AppContext,
472) -> anyhow::Result<()> {
473 if config.client.enable_health_check {
474 let forest_state = crate::health::ForestState {
475 config: config.clone(),
476 chain_config: ctx.state_manager.chain_config().clone(),
477 genesis_timestamp: ctx
478 .state_manager
479 .chain_store()
480 .genesis_block_header()
481 .timestamp,
482 sync_status: chain_follower.sync_status.clone(),
483 peer_manager: p2p_service.peer_manager().clone(),
484 };
485 let healthcheck_address = forest_state.config.client.healthcheck_address;
486 info!("Healthcheck endpoint will listen at {healthcheck_address}");
487 let listener = crate::utils::net::bind_tcp_listener(healthcheck_address, 0).await?;
488 services.spawn(async move {
489 crate::health::init_healthcheck_server(forest_state, listener)
490 .await
491 .context("Failed to initiate healthcheck server")
492 });
493 } else {
494 info!("Healthcheck service is disabled");
495 }
496 Ok(())
497}
498
499fn maybe_start_gc_service(
500 services: &mut JoinSet<anyhow::Result<()>>,
501 opts: &CliOpts,
502 config: &Config,
503 chain_follower: ChainFollower,
504) -> anyhow::Result<()> {
505 if opts.stateless {
507 return Ok(());
508 }
509
510 let snap_gc = Arc::new(SnapshotGarbageCollector::new(chain_follower, config)?);
511
512 GLOBAL_SNAPSHOT_GC
513 .set(snap_gc.clone())
514 .ok()
515 .context("failed to set GLOBAL_SNAPSHOT_GC")?;
516
517 services.spawn({
518 let snap_gc = snap_gc.clone();
519 async move {
520 snap_gc.event_loop().await;
521 Ok(())
522 }
523 });
524
525 if !opts.no_gc {
527 services.spawn({
528 let snap_gc = snap_gc.clone();
529 async move {
530 snap_gc.scheduler_loop().await;
531 Ok(())
532 }
533 });
534 }
535
536 Ok(())
537}
538
539#[allow(clippy::too_many_arguments)]
540fn maybe_start_rpc_service(
541 services: &mut JoinSet<anyhow::Result<()>>,
542 config: &Config,
543 mpool: MessagePool<ChainStore>,
544 chain_follower: &ChainFollower,
545 start_time: chrono::DateTime<chrono::Utc>,
546 shutdown: mpsc::Sender<()>,
547 rpc_stop_handle: jsonrpsee::server::StopHandle,
548 ctx: &AppContext,
549) -> anyhow::Result<()> {
550 if config.client.enable_rpc {
551 let rpc_address = config.client.rpc_address;
552 let metrics_mode = crate::rpc::MetricsMode::from(config.client.enable_metrics_endpoint);
553 let filter_list = config
554 .client
555 .rpc_filter_list
556 .as_ref()
557 .map(|path| crate::rpc::FilterList::new_from_file(path).map(Arc::new))
558 .transpose()?;
559 info!("JSON-RPC endpoint will listen at {rpc_address}");
560 let eth_event_handler = Arc::new(EthEventHandler::from_config(
561 &config.events,
562 ctx.chain_config().eth_chain_id,
563 mpool.subscriber(),
564 ));
565 if is_env_truthy("FOREST_JWT_DISABLE_EXP_VALIDATION") {
566 warn!(
567 "JWT expiration validation is disabled; this significantly weakens security and should only be used in tightly controlled environments"
568 );
569 }
570 services.spawn({
571 let state_manager = ctx.state_manager.shallow_clone();
572 let bad_blocks = chain_follower.bad_blocks.shallow_clone();
573 let sync_status = chain_follower.sync_status.shallow_clone();
574 let sync_network_context = chain_follower.network.shallow_clone();
575 let tipset_send = chain_follower.tipset_sender.clone();
576 let keystore = ctx.keystore.shallow_clone();
577 let snapshot_progress_tracker = ctx.snapshot_progress_tracker.clone();
578 let nonce_tracker = NonceTracker::new();
579 let mpool_locker = MpoolLocker::new();
580 let temp_dir = Arc::new(ctx.temp_dir.clone());
581 async move {
582 let rpc_listener = crate::utils::net::bind_tcp_listener(
583 rpc_address,
584 crate::rpc::default_max_connections(),
585 )
586 .await?;
587 start_rpc(
588 RPCState {
589 state_manager,
590 keystore,
591 mpool,
592 bad_blocks,
593 sync_status,
594 eth_event_handler,
595 eth_logs_feed: Default::default(),
596 sync_network_context,
597 start_time,
598 shutdown,
599 tipset_send,
600 snapshot_progress_tracker,
601 mpool_locker,
602 nonce_tracker,
603 temp_dir,
604 },
605 rpc_listener,
606 rpc_stop_handle,
607 filter_list,
608 metrics_mode,
609 )
610 .await
611 }
612 });
613 } else {
614 debug!("RPC disabled.");
615 };
616 Ok(())
617}
618
619fn maybe_start_f3_service(opts: &CliOpts, config: &Config, ctx: &AppContext) -> anyhow::Result<()> {
620 if crate::rpc::f3::F3_LEASE_MANAGER.get().is_some() {
622 return Ok(());
623 }
624
625 if !config.client.enable_rpc {
626 if crate::f3::is_sidecar_ffi_enabled(ctx.state_manager.chain_config()) {
627 tracing::warn!("F3 sidecar is enabled but not run because RPC is disabled. ")
628 }
629 return Ok(());
630 }
631
632 if !opts.halt_after_import && !opts.stateless {
633 let rpc_endpoint = config.client.rpc_v1_endpoint()?;
634 let state_manager = &ctx.state_manager;
635 let p2p_peer_id = ctx.p2p_peer_id;
636 let admin_jwt = ctx.admin_jwt.clone();
637 tokio::task::spawn_blocking({
638 crate::rpc::f3::F3_LEASE_MANAGER
639 .set(crate::rpc::f3::F3LeaseManager::new(
640 state_manager.chain_config().network.clone(),
641 p2p_peer_id,
642 ))
643 .expect("F3 lease manager should not have been initialized before");
644 let chain_config = state_manager.chain_config().clone();
645 let f3_root = crate::f3::get_f3_root(config);
646 let crate::f3::F3Options {
647 chain_finality,
648 bootstrap_epoch,
649 initial_power_table,
650 } = crate::f3::get_f3_sidecar_params(&chain_config);
651 move || {
652 crate::f3::run_f3_sidecar_if_enabled(
653 &chain_config,
654 rpc_endpoint.to_string(),
655 admin_jwt,
656 crate::rpc::f3::get_f3_rpc_endpoint().to_string(),
657 initial_power_table
658 .map(|i| i.to_string())
659 .unwrap_or_default(),
660 bootstrap_epoch,
661 chain_finality,
662 f3_root.display().to_string(),
663 );
664 }
665 });
666 tokio::task::spawn({
667 let chain_store = ctx.chain_store().shallow_clone();
668 async move {
669 tokio::time::sleep(Duration::from_secs(1)).await;
671 match (|| crate::rpc::f3::F3GetLatestCertificate::get())
672 .retry(ExponentialBuilder::default())
673 .await
674 {
675 Ok(f3_finalized_cert) => {
676 let f3_finalized_head = f3_finalized_cert.chain_head();
677 match chain_store
678 .chain_index()
679 .load_required_tipset(&f3_finalized_head.key)
680 {
681 Ok(ts) => {
682 chain_store.set_f3_finalized_tipset(ts);
683 tracing::info!(
684 "Set F3 finalized tipset to epoch {} and key {}",
685 f3_finalized_head.epoch,
686 f3_finalized_head.key,
687 );
688 }
689 Err(e) => {
690 tracing::error!(
691 "Failed to get F3 finalized tipset epoch {} and key {}: {e}",
692 f3_finalized_head.epoch,
693 f3_finalized_head.key
694 );
695 }
696 }
697 }
698 Err(e) => {
699 tracing::error!("Failed to get F3 latest certificate: {e:#}");
700 }
701 }
702 }
703 });
704 }
705
706 Ok(())
707}
708
709fn maybe_start_indexer_service(
710 services: &mut JoinSet<anyhow::Result<()>>,
711 opts: &CliOpts,
712 config: &Config,
713 ctx: &AppContext,
714) {
715 if config.chain_indexer.enable_indexer
716 && !opts.stateless
717 && !ctx.state_manager.chain_config().is_devnet()
718 {
719 let mut head_changes_rx = ctx.state_manager.chain_store().subscribe_head_changes();
720 let chain_store = ctx.state_manager.chain_store().shallow_clone();
721 services.spawn(async move {
722 tracing::info!("Starting indexer service");
723
724 loop {
726 match head_changes_rx.recv().await {
727 Ok(changes) => {
728 for ts in changes.applies {
729 tracing::debug!("Indexing tipset {}", ts.key());
730 let delegated_messages = chain_store
731 .headers_delegated_messages(ts.block_headers().iter())?;
732 chain_store.process_signed_messages(&delegated_messages)?;
733 }
734 }
735 Err(RecvError::Lagged(n)) => {
736 warn!("indexer service lagged: skipping {n} events")
737 }
738 Err(RecvError::Closed) => break Ok(()),
739 }
740 }
741 });
742
743 if let Some(retention_epochs) = config.chain_indexer.gc_retention_epochs {
745 let chain_store = ctx.state_manager.chain_store().shallow_clone();
746 let chain_config = ctx.state_manager.chain_config().clone();
747 services.spawn(async move {
748 tracing::info!("Starting collector for eth_mappings");
749 let mut collector = EthMappingCollector::new(
750 chain_store.db_owned(),
751 chain_config.eth_chain_id,
752 retention_epochs.into(),
753 );
754 collector.run().await
755 });
756 }
757 }
758}
759
760pub(super) async fn start(
762 start_time: chrono::DateTime<chrono::Utc>,
763 opts: CliOpts,
764 config: Config,
765 shutdown_send: mpsc::Sender<()>,
766 rpc_stop_handle: jsonrpsee::server::StopHandle,
767) -> anyhow::Result<()> {
768 startup_init(&config)?;
769 if opts.remove_existing_chain {
770 warn!("Deleting existing chain data for {}", config.chain());
771 delete_chain_data(&config)?;
772 }
773 start_services(
774 start_time,
775 &opts,
776 config.clone(),
777 shutdown_send.clone(),
778 rpc_stop_handle,
779 )
780 .await
781}
782
783pub(super) async fn start_services(
784 start_time: chrono::DateTime<chrono::Utc>,
785 opts: &CliOpts,
786 mut config: Config,
787 shutdown_send: mpsc::Sender<()>,
788 rpc_stop_handle: jsonrpsee::server::StopHandle,
789) -> anyhow::Result<()> {
790 crate::metrics::reset_collector_registry();
792 let mut services = JoinSet::new();
793 let network = config.chain();
794 let ctx = AppContext::init(opts, &config).await?;
795 info!("Using network :: {network}");
796 utils::misc::display_chain_logo(config.chain());
797 if opts.exit_after_init {
798 return Ok(());
799 }
800 if !opts.stateless
801 && !opts.skip_load_actors
802 && let Err(e) = ctx.state_manager.maybe_rewind_heaviest_tipset().await
803 {
804 tracing::warn!("error in maybe_rewind_heaviest_tipset: {e:#}");
805 }
806
807 let p2p_service = create_p2p_service(&mut services, &mut config, &ctx).await?;
808 let mpool = create_mpool(&mut services, &p2p_service, &ctx)?;
809 let chain_follower = create_chain_follower(opts, &p2p_service, mpool.shallow_clone(), &ctx)?;
810
811 maybe_start_rpc_service(
812 &mut services,
813 &config,
814 mpool.shallow_clone(),
815 &chain_follower,
816 start_time,
817 shutdown_send.clone(),
818 rpc_stop_handle,
819 &ctx,
820 )?;
821
822 maybe_import_snapshot(opts, &mut config, &ctx).await?;
823 if opts.halt_after_import {
824 services.shutdown().await;
826 return Ok(());
827 }
828
829 warmup_in_background(&ctx);
830 maybe_start_gc_service(&mut services, opts, &config, chain_follower.shallow_clone())?;
831 maybe_start_metrics_service(&mut services, &config, &ctx).await?;
832 maybe_start_f3_service(opts, &config, &ctx)?;
833 maybe_start_health_check_service(&mut services, &config, &p2p_service, &chain_follower, &ctx)
834 .await?;
835 maybe_start_indexer_service(&mut services, opts, &config, &ctx);
836 if !opts.stateless {
837 ensure_proof_params_downloaded().await?;
838 }
839 services.spawn(p2p_service.run());
840 start_chain_follower_service(&mut services, opts, &config, chain_follower);
841 propagate_error(&mut services)
843 .await
844 .context("services failure")
845 .map(|_| {})
846}
847
848fn warmup_in_background(ctx: &AppContext) {
849 let cs = ctx.chain_store().shallow_clone();
851 tokio::task::spawn_blocking(move || {
852 let start = Instant::now();
853 let head = cs.heaviest_tipset();
854 match cs.chain_index().repair_tipset_lookup_window(
855 &head,
856 head.epoch(),
857 cs.ec_calculator_finalized_epoch(),
858 ) {
859 Ok(n_repaired) => tracing::info!(
860 "Successfully verified tipset lookup table, {n_repaired} entries repaired, took {}.",
861 humantime::format_duration(start.elapsed()),
862 ),
863 Err(e) => warn!("failed to verify tipset lookup table: {e:#?}"),
864 }
865 });
866}
867
868async fn maybe_set_snapshot_path(
873 config: &mut Config,
874 chain_config: &ChainConfig,
875 epoch: ChainEpoch,
876 auto_download_snapshot: bool,
877 download_directory: &Path,
878) -> anyhow::Result<()> {
879 if !download_directory.is_dir() {
880 anyhow::bail!(
881 "`download_directory` does not exist: {}",
882 download_directory.display()
883 );
884 }
885
886 let vendor = snapshot::TrustedVendor::default();
887 let chain = config.chain();
888
889 let network_version = chain_config.network_version(epoch);
891 let network_version_is_small = network_version < NetworkVersion::V16;
892
893 let require_a_snapshot = network_version_is_small;
896 let have_a_snapshot = config.client.snapshot_path.is_some();
897
898 match (require_a_snapshot, have_a_snapshot, auto_download_snapshot) {
899 (false, _, _) => {} (true, true, _) => {} (true, false, true) => {
902 const AUTO_SNAPSHOT_PATH_ENV_KEY: &str = "FOREST_AUTO_DOWNLOAD_SNAPSHOT_PATH";
903 match std::env::var(AUTO_SNAPSHOT_PATH_ENV_KEY) {
904 Ok(path) if !path.is_empty() => {
905 tracing::info!(
906 "importing snapshot from {path} set by `{AUTO_SNAPSHOT_PATH_ENV_KEY}`"
907 );
908 config.client.snapshot_path = Some(path.into());
909 }
910 _ => {
911 let (resolved_url, _num_bytes, filename) =
915 crate::cli_shared::snapshot::peek(vendor, chain).await?;
916 tracing::info!("Downloading snapshot: {filename}");
917 config.client.snapshot_path = Some(resolved_url.to_string().into());
918 }
919 }
920 }
921 (true, false, false) => {
922 let (url, num_bytes, filename) = crate::cli_shared::snapshot::peek(vendor, chain)
924 .await
925 .context("couldn't get snapshot size")?;
926 println!(
929 "Forest requires a snapshot to sync with the network, but automatic fetching is disabled."
930 );
931 let message = format!(
932 "Fetch a {} snapshot? (denying will exit the program). ",
933 indicatif::HumanBytes(num_bytes)
934 );
935 let have_permission = asyncify(|| {
936 dialoguer::Confirm::with_theme(&ColorfulTheme::default())
937 .with_prompt(message)
938 .default(false)
939 .interact()
940 .unwrap_or(false)
942 })
943 .await;
944 if !have_permission {
945 bail!(
946 "Forest requires a snapshot to sync with the network, but automatic fetching is disabled."
947 )
948 }
949 tracing::info!("Downloading snapshot: {filename}");
950 config.client.snapshot_path = Some(url.to_string().into());
951 }
952 };
953
954 Ok(())
955}
956
957async fn propagate_error(
960 services: &mut JoinSet<anyhow::Result<()>>,
961) -> anyhow::Result<std::convert::Infallible> {
962 while let Some(result) = services.join_next().await {
963 if let Ok(Err(error_message)) = result {
964 return Err(error_message);
965 }
966 }
967 std::future::pending().await
968}
969
970fn asyncify<T>(f: impl FnOnce() -> T + Send + 'static) -> impl Future<Output = T>
975where
976 T: Send + 'static,
977{
978 tokio::task::spawn_blocking(f).then(|res| async { res.expect("spawned task panicked") })
979}
980
981#[cfg(test)]
982mod tests {
983 use rstest::rstest;
984
985 use super::*;
986
987 #[rstest]
988 #[case::current_non_positive(0, 1, anyhow::Result::Err(anyhow::anyhow!(
989 "current head epoch 0 is invalid"
990 )))]
991 #[case::current_non_positive(-1, 1, anyhow::Result::Err(anyhow::anyhow!(
992 "current head epoch 0 is invalid"
993 )))]
994 #[case::from_positive_beyond_head(10, 11, anyhow::Result::Err(anyhow::anyhow!(
995 "requested validation start epoch 11 is beyond the current head at epoch 10"
996 )))]
997 #[case::from_positive_within_range(10, 5, anyhow::Result::Ok(5..=10))]
998 #[case::from_zero(10, 0, anyhow::Result::Ok(0..=10))]
999 #[case::from_negative_within_range(10, -5, anyhow::Result::Ok(5..=10))]
1000 #[case::from_negative_beyond_range(10, -15, anyhow::Result::Ok(0..=10))]
1001 fn test_validation_range(
1002 #[case] current: ChainEpoch,
1003 #[case] from: ChainEpoch,
1004 #[case] expected: anyhow::Result<std::ops::RangeInclusive<ChainEpoch>>,
1005 ) {
1006 let result = validation_range(current, from);
1007 match expected {
1008 Ok(expected_range) => {
1009 assert_eq!(result.unwrap(), expected_range);
1010 }
1011 Err(_) => {
1012 assert!(result.is_err());
1013 }
1014 }
1015 }
1016}