Skip to main content

forest/tool/subcommands/api_cmd/
generate_test_snapshot.rs

1// Copyright 2019-2026 ChainSafe Systems
2// SPDX-License-Identifier: Apache-2.0, MIT
3
4use super::*;
5use crate::{
6    KeyStore, KeyStoreConfig,
7    blocks::TipsetKey,
8    chain::ChainStore,
9    chain_sync::{SyncStatusReport, network_context::SyncNetworkContext},
10    daemon::{bundle::load_actor_bundles, db_util::load_all_forest_cars},
11    db::{
12        CAR_DB_DIR_NAME, EthMappingsStore, HeaviestTipsetKeyProvider, MemoryDB, SettingsStore,
13        SettingsStoreExt, db_engine::open_db, parity_db::ParityDb,
14    },
15    genesis::read_genesis_header,
16    libp2p::{NetworkMessage, PeerManager},
17    libp2p_bitswap::{BitswapStoreRead, BitswapStoreReadWrite, Block64},
18    message_pool::{MessagePool, MpoolLocker, NonceTracker},
19    networks::ChainConfig,
20    prelude::*,
21    shim::address::CurrentNetwork,
22    state_manager::StateManager,
23};
24use api_compare_tests::TestDump;
25use arc_swap::ArcSwap;
26use fvm_shared4::address::Network;
27use openrpc_types::ParamStructure;
28use parking_lot::RwLock;
29use rpc::{RPCState, RpcMethod as _, eth::filter::EthEventHandler};
30use std::sync::atomic::{AtomicBool, Ordering};
31use tokio::{sync::mpsc, task::JoinSet};
32
33pub async fn run_test_with_dump(
34    test_dump: &TestDump,
35    db: Arc<ReadOpsTrackingStore<ManyCar<ParityDb>>>,
36    chain: &NetworkChain,
37    allow_response_mismatch: bool,
38    allow_failure: bool,
39) -> anyhow::Result<()> {
40    if chain.is_testnet() {
41        CurrentNetwork::set_global(Network::Testnet);
42    }
43    let mut run = false;
44    let chain_config = Arc::new(ChainConfig::from_chain(chain));
45    let (ctx, _, _) = ctx(db, chain_config).await?;
46    let params_raw = Some(serde_json::to_string(&test_dump.request.params)?);
47    let mut ext = http::Extensions::new();
48    ext.insert(test_dump.path);
49    macro_rules! run_test {
50        ($ty:ty) => {
51            if (test_dump.request.method_name.as_ref() == <$ty>::NAME
52                || Some(test_dump.request.method_name.as_ref()) == <$ty>::NAME_ALIAS)
53                && <$ty>::API_PATHS.contains(test_dump.path)
54            {
55                let params = <$ty>::parse_params(params_raw.clone(), ParamStructure::Either)?;
56                match <$ty>::handle(ctx.clone(), params, &ext).await {
57                    Ok(result) => {
58                        anyhow::ensure!(
59                            allow_response_mismatch
60                                || test_dump.forest_response == Ok(result.into_lotus_json_value()?),
61                            "Response mismatch between Forest and Lotus"
62                        );
63                    }
64                    Err(_) if allow_failure => {
65                        // If we allow failure, we do not check the error
66                    }
67                    Err(e) => {
68                        bail!("Error running test {}: {}", <$ty>::NAME, e);
69                    }
70                }
71                run = true;
72            }
73        };
74    }
75    crate::for_each_rpc_method!(run_test);
76    anyhow::ensure!(run, "RPC method not found");
77    Ok(())
78}
79
80pub async fn load_db(
81    db_root: &Path,
82    chain: Option<&NetworkChain>,
83) -> anyhow::Result<Arc<ReadOpsTrackingStore<ManyCar<ParityDb>>>> {
84    let db_writer = open_db(db_root.into(), &Default::default())?;
85    let db = ManyCar::new(db_writer);
86    let forest_car_db_dir = db_root.join(CAR_DB_DIR_NAME);
87    load_all_forest_cars(&db, &forest_car_db_dir)?;
88    if let Some(chain) = chain {
89        load_actor_bundles(&db, chain).await?;
90    }
91    Ok(Arc::new(ReadOpsTrackingStore::new(db)))
92}
93
94pub(super) fn build_index(db: Arc<ReadOpsTrackingStore<ManyCar<ParityDb>>>) -> Option<Index> {
95    let mut index = Index::default();
96    let reader = db.tracker.eth_mappings_db.read();
97    for (k, v) in reader.iter() {
98        index
99            .eth_mappings
100            .get_or_insert_with(Default::default)
101            .insert(k.to_string(), Payload(v.clone()));
102    }
103    if index == Index::default() {
104        None
105    } else {
106        Some(index)
107    }
108}
109
110async fn ctx(
111    db: Arc<ReadOpsTrackingStore<ManyCar<ParityDb>>>,
112    chain_config: Arc<ChainConfig>,
113) -> anyhow::Result<(
114    Arc<RPCState>,
115    flume::Receiver<NetworkMessage>,
116    tokio::sync::mpsc::Receiver<()>,
117)> {
118    let (network_send, network_rx) = flume::bounded(5);
119    let (tipset_send, _) = flume::bounded(5);
120    let genesis_header =
121        read_genesis_header(None, chain_config.genesis_bytes(&db).await?.as_deref(), &db).await?;
122    let chain_store = ChainStore::new(db.clone(), chain_config, genesis_header)?;
123    let state_manager = StateManager::new(chain_store.shallow_clone())?
124        // cache must be disabled to avoid flakiness in RPC regression tests
125        .with_id_address_cache_disabled();
126    let mut services: JoinSet<anyhow::Result<()>> = JoinSet::new();
127    let message_pool = MessagePool::new(
128        chain_store,
129        network_send.clone(),
130        Default::default(),
131        state_manager.chain_config().clone(),
132        &mut services,
133    )?;
134    // See `super::test_snapshot::drain_mpool_services` for rationale.
135    services.abort_all();
136    tokio::spawn(super::test_snapshot::drain_mpool_services(services));
137
138    let peer_manager = Arc::new(PeerManager::default());
139    let sync_network_context =
140        SyncNetworkContext::new(network_send, peer_manager, state_manager.db_owned());
141    let (shutdown, shutdown_recv) = mpsc::channel(1);
142    let nonce_tracker = NonceTracker::new();
143    let eth_event_handler = Arc::new(EthEventHandler::from_config(
144        &crate::cli_shared::cli::EventsConfig::default(),
145        state_manager.chain_config().eth_chain_id,
146        message_pool.subscriber(),
147    ));
148    let rpc_state = Arc::new(RPCState {
149        state_manager,
150        keystore: Arc::new(RwLock::new(KeyStore::new(KeyStoreConfig::Memory)?)),
151        mpool: message_pool,
152        bad_blocks: Default::default(),
153        sync_status: Arc::new(ArcSwap::from_pointee(SyncStatusReport::init())),
154        eth_event_handler,
155        eth_logs_feed: Default::default(),
156        sync_network_context,
157        start_time: chrono::Utc::now(),
158        shutdown,
159        tipset_send,
160        snapshot_progress_tracker: Default::default(),
161        mpool_locker: MpoolLocker::new(),
162        nonce_tracker,
163        temp_dir: Arc::new(std::env::temp_dir()),
164    });
165    Ok((rpc_state, network_rx, shutdown_recv))
166}
167
168/// A [`Blockstore`] wrapper that tracks read operations to the inner [`Blockstore`] with an [`MemoryDB`]
169pub struct ReadOpsTrackingStore<T> {
170    inner: T,
171    pub tracker: Arc<MemoryDB>,
172    tracking: AtomicBool,
173}
174
175impl<T> ReadOpsTrackingStore<T> {
176    pub fn resume_tracking(&self) {
177        self.tracking.store(true, Ordering::Relaxed);
178    }
179
180    pub fn pause_tracking(&self) {
181        self.tracking.store(false, Ordering::Relaxed);
182    }
183
184    fn tracking(&self) -> bool {
185        self.tracking.load(Ordering::Relaxed)
186    }
187}
188
189impl<T> ReadOpsTrackingStore<T>
190where
191    T: Blockstore + SettingsStore + HeaviestTipsetKeyProvider,
192{
193    fn is_chain_head_tracked(&self) -> anyhow::Result<bool> {
194        SettingsStore::exists(&self.tracker, crate::db::setting_keys::HEAD_KEY)
195    }
196
197    pub fn ensure_chain_head_is_tracked(&self) -> anyhow::Result<()> {
198        if !self.is_chain_head_tracked()? {
199            SettingsStoreExt::write_obj(
200                &self.tracker,
201                crate::db::setting_keys::HEAD_KEY,
202                &self
203                    .inner
204                    .heaviest_tipset_key()?
205                    .context("heaviest tipset key not found")?,
206            )?;
207        }
208
209        Ok(())
210    }
211}
212
213impl<T> ReadOpsTrackingStore<T> {
214    pub fn new(inner: T) -> Self {
215        Self {
216            inner,
217            tracker: Arc::new(Default::default()),
218            tracking: AtomicBool::new(true),
219        }
220    }
221
222    pub async fn export_forest_car<W: tokio::io::AsyncWrite + Unpin>(
223        &self,
224        writer: &mut W,
225    ) -> anyhow::Result<()> {
226        self.tracker.export_forest_car(writer).await
227    }
228}
229
230impl<T: HeaviestTipsetKeyProvider> HeaviestTipsetKeyProvider for ReadOpsTrackingStore<T> {
231    fn heaviest_tipset_key(&self) -> anyhow::Result<Option<TipsetKey>> {
232        self.inner.heaviest_tipset_key()
233    }
234
235    fn set_heaviest_tipset_key(&self, tsk: &TipsetKey) -> anyhow::Result<()> {
236        self.inner.set_heaviest_tipset_key(tsk)
237    }
238}
239
240impl<T: Blockstore> Blockstore for ReadOpsTrackingStore<T> {
241    fn get(&self, k: &Cid) -> anyhow::Result<Option<Vec<u8>>> {
242        let result = self.inner.get(k)?;
243        if self.tracking()
244            && let Some(v) = &result
245        {
246            self.tracker.put_keyed(k, v.as_slice())?;
247        }
248        Ok(result)
249    }
250
251    fn put_keyed(&self, k: &Cid, block: &[u8]) -> anyhow::Result<()> {
252        self.inner.put_keyed(k, block)
253    }
254}
255
256impl<T: SettingsStore> SettingsStore for ReadOpsTrackingStore<T> {
257    fn read_bin(&self, key: &str) -> anyhow::Result<Option<Vec<u8>>> {
258        let result = self.inner.read_bin(key)?;
259        if self.tracking()
260            && let Some(v) = &result
261        {
262            SettingsStore::write_bin(&self.tracker, key, v.as_slice())?;
263        }
264        Ok(result)
265    }
266
267    fn write_bin(&self, key: &str, value: &[u8]) -> anyhow::Result<()> {
268        self.inner.write_bin(key, value)
269    }
270
271    fn exists(&self, key: &str) -> anyhow::Result<bool> {
272        let result = self.inner.read_bin(key)?;
273        if self.tracking()
274            && let Some(v) = &result
275        {
276            SettingsStore::write_bin(&self.tracker, key, v.as_slice())?;
277        }
278        Ok(result.is_some())
279    }
280
281    fn setting_keys(&self) -> anyhow::Result<Vec<String>> {
282        self.inner.setting_keys()
283    }
284}
285
286impl<T: BitswapStoreRead> BitswapStoreRead for ReadOpsTrackingStore<T> {
287    fn contains(&self, cid: &Cid) -> anyhow::Result<bool> {
288        let result = self.inner.get(cid)?;
289        if self.tracking()
290            && let Some(v) = &result
291        {
292            Blockstore::put_keyed(&self.tracker, cid, v.as_slice())?;
293        }
294        Ok(result.is_some())
295    }
296
297    fn get(&self, cid: &Cid) -> anyhow::Result<Option<Vec<u8>>> {
298        let result = self.inner.get(cid)?;
299        if self.tracking()
300            && let Some(v) = &result
301        {
302            Blockstore::put_keyed(&self.tracker, cid, v.as_slice())?;
303        }
304        Ok(result)
305    }
306}
307
308impl<T: BitswapStoreReadWrite> BitswapStoreReadWrite for ReadOpsTrackingStore<T> {
309    type Hashes = <T as BitswapStoreReadWrite>::Hashes;
310
311    fn insert(&self, block: &Block64<Self::Hashes>) -> anyhow::Result<()> {
312        self.inner.insert(block)
313    }
314}
315
316impl<T: EthMappingsStore> EthMappingsStore for ReadOpsTrackingStore<T> {
317    fn read_bin(&self, key: &EthHash) -> anyhow::Result<Option<Vec<u8>>> {
318        let result = self.inner.read_bin(key)?;
319        if self.tracking()
320            && let Some(v) = &result
321        {
322            EthMappingsStore::write_bin(&self.tracker, key, v.as_slice())?;
323        }
324        Ok(result)
325    }
326
327    fn write_bin(&self, key: &EthHash, value: &[u8]) -> anyhow::Result<()> {
328        self.inner.write_bin(key, value)
329    }
330
331    fn exists(&self, key: &EthHash) -> anyhow::Result<bool> {
332        self.inner.exists(key)
333    }
334
335    fn get_message_cids(&self) -> anyhow::Result<Vec<(Cid, u64)>> {
336        self.inner.get_message_cids()
337    }
338
339    fn delete(&self, keys: Vec<EthHash>) -> anyhow::Result<()> {
340        self.inner.delete(keys)
341    }
342
343    fn tipset_key_by_epoch(&self, epoch: ChainEpoch) -> anyhow::Result<Option<TipsetKey>> {
344        let result = self.inner.tipset_key_by_epoch(epoch)?;
345        if self.tracking()
346            && let Some(tsk) = &result
347        {
348            EthMappingsStore::set_tipset_key_at_epoch_raw(&self.tracker, epoch, tsk)?;
349        }
350        Ok(result)
351    }
352
353    fn delete_tipset_key_at_epoch(&self, epoch: ChainEpoch) -> anyhow::Result<()> {
354        EthMappingsStore::delete_tipset_key_at_epoch(&self.tracker, epoch)
355    }
356
357    fn set_tipset_key_at_epoch_raw(
358        &self,
359        epoch: ChainEpoch,
360        tsk: &TipsetKey,
361    ) -> anyhow::Result<()> {
362        self.inner.set_tipset_key_at_epoch_raw(epoch, tsk)
363    }
364}