1use 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 }
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 .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 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
177pub 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}