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        snapshot_progress_tracker: Default::default(),
169        mpool_locker: MpoolLocker::new(),
170        nonce_tracker,
171        temp_dir: Arc::new(std::env::temp_dir()),
172    });
173    Ok((rpc_state, network_rx, shutdown_recv))
174}
175
176/// A [`Blockstore`] wrapper that tracks read operations to the inner [`Blockstore`] with an [`MemoryDB`]
177pub struct ReadOpsTrackingStore<T> {
178    inner: T,
179    pub tracker: Arc<MemoryDB>,
180    tracking: AtomicBool,
181}
182
183impl<T> ReadOpsTrackingStore<T> {
184    pub fn resume_tracking(&self) {
185        self.tracking.store(true, Ordering::Relaxed);
186    }
187
188    pub fn pause_tracking(&self) {
189        self.tracking.store(false, Ordering::Relaxed);
190    }
191
192    fn tracking(&self) -> bool {
193        self.tracking.load(Ordering::Relaxed)
194    }
195}
196
197impl<T> ReadOpsTrackingStore<T>
198where
199    T: Blockstore + SettingsStore + HeaviestTipsetKeyProvider,
200{
201    fn is_chain_head_tracked(&self) -> anyhow::Result<bool> {
202        SettingsStore::exists(&self.tracker, crate::db::setting_keys::HEAD_KEY)
203    }
204
205    pub fn ensure_chain_head_is_tracked(&self) -> anyhow::Result<()> {
206        if !self.is_chain_head_tracked()? {
207            SettingsStoreExt::write_obj(
208                &self.tracker,
209                crate::db::setting_keys::HEAD_KEY,
210                &self
211                    .inner
212                    .heaviest_tipset_key()?
213                    .context("heaviest tipset key not found")?,
214            )?;
215        }
216
217        Ok(())
218    }
219}
220
221impl<T> ReadOpsTrackingStore<T> {
222    pub fn new(inner: T) -> Self {
223        Self {
224            inner,
225            tracker: Arc::new(Default::default()),
226            tracking: AtomicBool::new(true),
227        }
228    }
229
230    pub async fn export_forest_car<W: tokio::io::AsyncWrite + Unpin>(
231        &self,
232        writer: &mut W,
233    ) -> anyhow::Result<()> {
234        self.tracker.export_forest_car(writer).await
235    }
236}
237
238impl<T: HeaviestTipsetKeyProvider> HeaviestTipsetKeyProvider for ReadOpsTrackingStore<T> {
239    fn heaviest_tipset_key(&self) -> anyhow::Result<Option<TipsetKey>> {
240        self.inner.heaviest_tipset_key()
241    }
242
243    fn set_heaviest_tipset_key(&self, tsk: &TipsetKey) -> anyhow::Result<()> {
244        self.inner.set_heaviest_tipset_key(tsk)
245    }
246}
247
248impl<T: Blockstore> Blockstore for ReadOpsTrackingStore<T> {
249    fn get(&self, k: &Cid) -> anyhow::Result<Option<Vec<u8>>> {
250        let result = self.inner.get(k)?;
251        if self.tracking()
252            && let Some(v) = &result
253        {
254            self.tracker.put_keyed(k, v.as_slice())?;
255        }
256        Ok(result)
257    }
258
259    fn put_keyed(&self, k: &Cid, block: &[u8]) -> anyhow::Result<()> {
260        self.inner.put_keyed(k, block)
261    }
262}
263
264impl<T: SettingsStore> SettingsStore for ReadOpsTrackingStore<T> {
265    fn read_bin(&self, key: &str) -> anyhow::Result<Option<Vec<u8>>> {
266        let result = self.inner.read_bin(key)?;
267        if self.tracking()
268            && let Some(v) = &result
269        {
270            SettingsStore::write_bin(&self.tracker, key, v.as_slice())?;
271        }
272        Ok(result)
273    }
274
275    fn write_bin(&self, key: &str, value: &[u8]) -> anyhow::Result<()> {
276        self.inner.write_bin(key, value)
277    }
278
279    fn exists(&self, key: &str) -> anyhow::Result<bool> {
280        let result = self.inner.read_bin(key)?;
281        if self.tracking()
282            && let Some(v) = &result
283        {
284            SettingsStore::write_bin(&self.tracker, key, v.as_slice())?;
285        }
286        Ok(result.is_some())
287    }
288
289    fn setting_keys(&self) -> anyhow::Result<Vec<String>> {
290        self.inner.setting_keys()
291    }
292}
293
294impl<T: BitswapStoreRead> BitswapStoreRead for ReadOpsTrackingStore<T> {
295    fn contains(&self, cid: &Cid) -> anyhow::Result<bool> {
296        let result = self.inner.get(cid)?;
297        if self.tracking()
298            && let Some(v) = &result
299        {
300            Blockstore::put_keyed(&self.tracker, cid, v.as_slice())?;
301        }
302        Ok(result.is_some())
303    }
304
305    fn get(&self, cid: &Cid) -> anyhow::Result<Option<Vec<u8>>> {
306        let result = self.inner.get(cid)?;
307        if self.tracking()
308            && let Some(v) = &result
309        {
310            Blockstore::put_keyed(&self.tracker, cid, v.as_slice())?;
311        }
312        Ok(result)
313    }
314}
315
316impl<T: BitswapStoreReadWrite> BitswapStoreReadWrite for ReadOpsTrackingStore<T> {
317    type Hashes = <T as BitswapStoreReadWrite>::Hashes;
318
319    fn insert(&self, block: &Block64<Self::Hashes>) -> anyhow::Result<()> {
320        self.inner.insert(block)
321    }
322}
323
324impl<T: EthMappingsStore> EthMappingsStore for ReadOpsTrackingStore<T> {
325    fn read_bin(&self, key: &EthHash) -> anyhow::Result<Option<Vec<u8>>> {
326        let result = self.inner.read_bin(key)?;
327        if self.tracking()
328            && let Some(v) = &result
329        {
330            EthMappingsStore::write_bin(&self.tracker, key, v.as_slice())?;
331        }
332        Ok(result)
333    }
334
335    fn write_bin(&self, key: &EthHash, value: &[u8]) -> anyhow::Result<()> {
336        self.inner.write_bin(key, value)
337    }
338
339    fn exists(&self, key: &EthHash) -> anyhow::Result<bool> {
340        self.inner.exists(key)
341    }
342
343    fn get_message_cids(&self) -> anyhow::Result<Vec<(Cid, u64)>> {
344        self.inner.get_message_cids()
345    }
346
347    fn delete(&self, keys: Vec<EthHash>) -> anyhow::Result<()> {
348        self.inner.delete(keys)
349    }
350
351    fn tipset_key_by_epoch(&self, epoch: ChainEpoch) -> anyhow::Result<Option<TipsetKey>> {
352        let result = self.inner.tipset_key_by_epoch(epoch)?;
353        if self.tracking()
354            && let Some(tsk) = &result
355        {
356            EthMappingsStore::set_tipset_key_at_epoch_raw(&self.tracker, epoch, tsk)?;
357        }
358        Ok(result)
359    }
360
361    fn delete_tipset_key_at_epoch(&self, epoch: ChainEpoch) -> anyhow::Result<()> {
362        EthMappingsStore::delete_tipset_key_at_epoch(&self.tracker, epoch)
363    }
364
365    fn set_tipset_key_at_epoch_raw(
366        &self,
367        epoch: ChainEpoch,
368        tsk: &TipsetKey,
369    ) -> anyhow::Result<()> {
370        self.inner.set_tipset_key_at_epoch_raw(epoch, tsk)
371    }
372}
373
374impl<T: EthBlockBloomStore> EthBlockBloomStore for ReadOpsTrackingStore<T> {
375    fn read_bloom(&self, key: &Cid) -> anyhow::Result<Option<[u8; BLOCK_BLOOM_LEN]>> {
376        let result = self.inner.read_bloom(key)?;
377        if self.tracking()
378            && let Some(bloom) = &result
379        {
380            self.tracker.write_bloom(key, 0, bloom)?;
381        }
382        Ok(result)
383    }
384
385    fn write_bloom(
386        &self,
387        key: &Cid,
388        height: ChainEpoch,
389        bloom: &[u8; BLOCK_BLOOM_LEN],
390    ) -> anyhow::Result<()> {
391        self.inner.write_bloom(key, height, bloom)
392    }
393
394    fn delete_blooms_before_height(&self, height: ChainEpoch) -> anyhow::Result<()> {
395        self.inner.delete_blooms_before_height(height)
396    }
397}