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