Skip to main content

forest/daemon/
mod.rs

1// Copyright 2019-2026 ChainSafe Systems
2// SPDX-License-Identifier: Apache-2.0, MIT
3
4pub 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
59/// Increase the file descriptor limit to a reasonable number.
60/// This prevents the node from failing if the default soft limit is too low.
61/// Note that the value is only increased, never decreased.
62fn 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
80// Start the daemon and abort if we're interrupted by ctrl-c, SIGTERM, or `forest-cli shutdown`.
81pub 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
106/// This function initialize Forest with below steps
107/// - increase file descriptor limit (for parity-db)
108/// - setup proofs parameter cache directory
109/// - prints Forest version
110fn startup_init(config: &Config) -> anyhow::Result<()> {
111    maybe_increase_fd_limit()?;
112    // Sets proof parameter file download path early, the files will be checked and
113    // downloaded later right after snapshot import step
114    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    // Sets the latest snapshot if needed for downloading later
130    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    // Import chain if needed
143    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        // Explicitly set heaviest tipset here in case HEAD_KEY has already been set
160        // in the current setting store
161        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 the snapshot progress state is not completed,
169    // set the state to not required
170    if !snapshot_tracker.is_completed() {
171        snapshot_tracker.not_required();
172    }
173
174    if let Some(validate_from) = config.client.snapshot_height {
175        // We've been provided a snapshot and asked to validate it
176        ensure_proof_params_downloaded().await?;
177        // Use the specified HEAD, otherwise take the current HEAD.
178        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        // `validate_range` is CPU-bound (drives rayon-parallel VM execution) and
185        // can run for minutes. Safer to spawn it on a blocking thread.
186        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
196/// Returns the range of epochs to validate. This includes special handling for negative `from`
197/// values, which are interpreted as offsets from the current epoch.
198fn 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    // Negative values scroll back from the current head (e.g. --height=-1000).
208    // `saturating_add` + `.max(0)` keeps extreme negatives from underflowing or
209    // wrapping to a huge positive (which would silently produce an empty range).
210    let start = if from.is_negative() {
211        current.saturating_add(from).max(0)
212    } else {
213        from
214    };
215
216    // An absolute `--height` past the head would otherwise produce an empty
217    // range and silently succeed without validating anything.
218    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 bootstrap peers are not set, set them
296    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    // Libp2p service setup
303    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    // Prefill RPC method caches for newly validated tipsets to speed up subsequent RPC calls.
368    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                        // Skip if the node is catching up to avoid unnecessary work, as the head may be changing rapidly.
379                        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                // First, compute state for the ts as it's disallowed for RPC methods by default
408                if let Err(e) = state_manager.load_executed_tipset(&ts).await {
409                    warn!("failed to load executed tipset for cache warmup: {e:#}");
410                    return; // Skip when state computation fails
411                }
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                // Warms both the FVM-replay cache and the parity-trace cache,
426                // since `eth_trace_block` calls `execution_trace` internally.
427                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 the node is stateless, GC shouldn't get triggered even on demand.
506    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    // GC shouldn't run periodically if the node is stateless or if the user has disabled it.
526    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    // already running
621    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                // wait 1s to let F3 RPC server start
670                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            // Continuously listen for head changes
725            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        // Run the collector only if chain indexer is enabled
744        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
760/// Starts daemon process
761pub(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    // Cleanup the collector prometheus metrics registry on start
791    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        // Cancel all async services
825        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    // blocking until any of the services returns an error,
842    propagate_error(&mut services)
843        .await
844        .context("services failure")
845        .map(|_| {})
846}
847
848fn warmup_in_background(ctx: &AppContext) {
849    // Verify and re-populate the `tipset_by_height` lookup table over the whole chain.
850    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
868/// If our current chain is below a supported height, we need a snapshot to bring it up
869/// to a supported height. If we've not been given a snapshot by the user, get one.
870///
871/// An [`Err`] should be considered fatal.
872async 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    // What height is our chain at right now, and what network version does that correspond to?
890    let network_version = chain_config.network_version(epoch);
891    let network_version_is_small = network_version < NetworkVersion::V16;
892
893    // We don't support small network versions (we can't validate from e.g genesis).
894    // So we need a snapshot (which will be from a recent network version)
895    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, _, _) => {}   // noop - don't need a snapshot
900        (true, true, _) => {} // noop - we need a snapshot, and we have one
901        (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                    // Resolve the redirect URL to get the actual snapshot URL
912                    // This ensures all chunks download from the same snapshot even if
913                    // a new snapshot is published during the download
914                    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            // we need a snapshot, don't have one, and don't have permission to download one, so ask the user
923            let (url, num_bytes, filename) = crate::cli_shared::snapshot::peek(vendor, chain)
924                .await
925                .context("couldn't get snapshot size")?;
926            // dialoguer will double-print long lines, so manually print the first clause ourselves,
927            // then let `Confirm` handle the second.
928            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                    // e.g not a tty (or some other error), so haven't got permission.
941                    .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
957/// returns the first error with which any of the services end, or never returns at all
958// This should return anyhow::Result<!> once the `Never` type is stabilized
959async 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
970/// Run the closure on a thread where blocking is allowed
971///
972/// # Panics
973/// If the closure panics
974fn 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}