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 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
176pub 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}