Skip to content

Commit 01d569a

Browse files
authored
fix: internal trace calls should not compete the semaphore (#7421)
1 parent c6e0744 commit 01d569a

7 files changed

Lines changed: 105 additions & 31 deletions

File tree

src/daemon/mod.rs

Lines changed: 35 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ use crate::networks::{self, ChainConfig};
2525
use crate::prelude::*;
2626
use crate::rpc::RPCState;
2727
use crate::rpc::eth::filter::EthEventHandler;
28+
use crate::rpc::eth::types::CallSource;
2829
use crate::rpc::start_rpc;
2930
use crate::shim::address::Address;
3031
use crate::shim::clock::ChainEpoch;
@@ -378,6 +379,10 @@ fn maybe_prefill_rpc_caches(
378379
// Skip if the node is catching up to avoid unnecessary work, as the head may be changing rapidly.
379380
continue;
380381
}
382+
Ok(tsk) if state_manager.chain_store().heaviest_tipset().key() != &tsk => {
383+
// Skip if the tipset has already been superseded
384+
continue;
385+
}
381386
Ok(tsk) => {
382387
let state_manager = state_manager.shallow_clone();
383388
let cancellation_token = cancellation_token.clone();
@@ -386,6 +391,7 @@ fn maybe_prefill_rpc_caches(
386391
.run_until_cancelled(prefill_rpc_caches_for_tipset(
387392
state_manager,
388393
tsk,
394+
cancellation_token.clone(),
389395
))
390396
.await
391397
});
@@ -400,7 +406,11 @@ fn maybe_prefill_rpc_caches(
400406
}
401407
}
402408

403-
async fn prefill_rpc_caches_for_tipset(state_manager: StateManager, tsk: TipsetKey) {
409+
async fn prefill_rpc_caches_for_tipset(
410+
state_manager: StateManager,
411+
tsk: TipsetKey,
412+
cancellation_token: CancellationToken,
413+
) {
404414
match state_manager.chain_index().load_required_tipset(&tsk) {
405415
Ok(ts) => {
406416
{
@@ -410,6 +420,30 @@ async fn prefill_rpc_caches_for_tipset(state_manager: StateManager, tsk: TipsetK
410420
return; // Skip when state computation fails
411421
}
412422
}
423+
{
424+
// Warms both the FVM-replay cache and the parity-trace cache,
425+
// since `eth_trace_block` calls `execution_trace` internally.
426+
// Note that we do not block the loop here as the trace computation can be expensive.
427+
// Also, we skip this tipset when it has already been superseded
428+
if state_manager.chain_store().heaviest_tipset().key() == ts.key() {
429+
tokio::spawn({
430+
let state_manager = state_manager.shallow_clone();
431+
let ts = ts.shallow_clone();
432+
async move {
433+
if let Some(Err(e)) = cancellation_token
434+
.run_until_cancelled(crate::rpc::eth::eth_trace_block(
435+
&state_manager,
436+
&ts,
437+
CallSource::Internal,
438+
))
439+
.await
440+
{
441+
warn!("failed to call `eth_trace_block` for cache warmup: {e:#}");
442+
}
443+
}
444+
});
445+
}
446+
}
413447
for tx_info in [crate::rpc::eth::TxInfo::Full, crate::rpc::eth::TxInfo::Hash] {
414448
if let Err(e) = crate::rpc::eth::Block::from_filecoin_tipset(
415449
&state_manager,
@@ -421,13 +455,6 @@ async fn prefill_rpc_caches_for_tipset(state_manager: StateManager, tsk: TipsetK
421455
warn!("failed to call `Block::from_filecoin_tipset` for cache warmup: {e:#}");
422456
}
423457
}
424-
{
425-
// Warms both the FVM-replay cache and the parity-trace cache,
426-
// since `eth_trace_block` calls `execution_trace` internally.
427-
if let Err(e) = crate::rpc::eth::eth_trace_block(&state_manager, &ts).await {
428-
warn!("failed to call `eth_trace_block` for cache warmup: {e:#}");
429-
}
430-
}
431458
{
432459
use crate::rpc::eth::filter::{Matcher, SkipEvent};
433460
struct CollectEventsCachePrefillingMatcher;

src/rpc/methods/eth.rs

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3538,7 +3538,7 @@ impl RpcMethod<1> for EthTraceBlock {
35383538
let ts = resolver
35393539
.tipset_by_block_number_or_hash(block_param, ResolveNullTipset::Fail)
35403540
.await?;
3541-
eth_trace_block(&ctx.state_manager, &ts)
3541+
eth_trace_block(&ctx.state_manager, &ts, CallSource::External)
35423542
.await
35433543
.map(NotNullVec)
35443544
}
@@ -3548,8 +3548,9 @@ impl RpcMethod<1> for EthTraceBlock {
35483548
async fn execute_tipset_traces(
35493549
state_manager: &StateManager,
35503550
ts: &Tipset,
3551+
source: CallSource,
35513552
) -> Result<(StateTree<DbImpl>, Vec<trace::TipsetTraceEntry>), ServerError> {
3552-
let (state_root, raw_traces) = state_manager.execution_trace(ts).await?;
3553+
let (state_root, raw_traces) = state_manager.execution_trace(ts, source).await?;
35533554
let state = state_manager.get_state_tree(&state_root)?;
35543555

35553556
// Resolve every non-system message's tx hash in parallel. Each lookup is
@@ -3604,6 +3605,7 @@ fn non_system_traces_with_positions(
36043605
pub(crate) async fn eth_trace_block(
36053606
state_manager: &StateManager,
36063607
ts: &Tipset,
3608+
source: CallSource,
36073609
) -> Result<Vec<EthBlockTrace>, ServerError> {
36083610
// 64 most-recent blocks; bounded by count, not bytes (a few MiB on mainnet,
36093611
// see the `cache_eth_trace_block_size` metric).
@@ -3616,7 +3618,7 @@ pub(crate) async fn eth_trace_block(
36163618
let block_cid = ts.key().cid()?;
36173619
let traces = ETH_TRACE_BLOCK_CACHE
36183620
.get_or_insert_async(&CidWrapper::from(block_cid), async move {
3619-
let (state, entries) = execute_tipset_traces(state_manager, ts).await?;
3621+
let (state, entries) = execute_tipset_traces(state_manager, ts, source).await?;
36203622
let block_hash: EthHash = block_cid.into();
36213623
let mut all_traces = vec![];
36223624

@@ -3665,6 +3667,7 @@ impl RpcMethod<2> for EthDebugTraceTransaction {
36653667
tx_hash,
36663668
opts,
36673669
&cancellation_token,
3670+
CallSource::External,
36683671
)
36693672
.await
36703673
}
@@ -3676,6 +3679,7 @@ async fn debug_trace_transaction(
36763679
tx_hash: String,
36773680
opts: GethDebugTracingOptions,
36783681
cancellation_token: &CancellationToken,
3682+
source: CallSource,
36793683
) -> Result<GethTrace, ServerError> {
36803684
let tracer = match &opts.tracer {
36813685
Some(t) => t.clone(),
@@ -3752,7 +3756,7 @@ async fn debug_trace_transaction(
37523756
return Ok(GethTrace::PreState(frame));
37533757
}
37543758

3755-
let (state, entries) = execute_tipset_traces(&ctx.state_manager, &ts).await?;
3759+
let (state, entries) = execute_tipset_traces(&ctx.state_manager, &ts, source).await?;
37563760
let entry = entries
37573761
.into_iter()
37583762
.find(|e| e.tx_hash == eth_hash)
@@ -3956,7 +3960,7 @@ impl RpcMethod<1> for EthTraceTransaction {
39563960
.tipset_by_block_number_or_hash(eth_txn.block_number, ResolveNullTipset::TakeOlder)
39573961
.await?;
39583962

3959-
let traces = eth_trace_block(&ctx.state_manager, &ts)
3963+
let traces = eth_trace_block(&ctx.state_manager, &ts, CallSource::External)
39603964
.await?
39613965
.into_iter()
39623966
.filter(|trace| trace.transaction_hash == eth_hash)
@@ -3996,7 +4000,7 @@ impl RpcMethod<2> for EthTraceReplayBlockTransactions {
39964000
.tipset_by_block_number_or_hash(block_param, ResolveNullTipset::Fail)
39974001
.await?;
39984002

3999-
eth_trace_replay_block_transactions(&ctx, &ts)
4003+
eth_trace_replay_block_transactions(&ctx, &ts, CallSource::External)
40004004
.await
40014005
.map(NotNullVec)
40024006
}
@@ -4005,8 +4009,9 @@ impl RpcMethod<2> for EthTraceReplayBlockTransactions {
40054009
async fn eth_trace_replay_block_transactions(
40064010
ctx: &Ctx,
40074011
ts: &Tipset,
4012+
source: CallSource,
40084013
) -> Result<Vec<EthReplayBlockTransactionTrace>, ServerError> {
4009-
let (state, entries) = execute_tipset_traces(&ctx.state_manager, ts).await?;
4014+
let (state, entries) = execute_tipset_traces(&ctx.state_manager, ts, source).await?;
40104015

40114016
let mut all_traces = vec![];
40124017
for entry in entries {

src/rpc/methods/eth/types.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,13 @@ pub const METHOD_GET_STORAGE_AT: u64 = 5;
1717

1818
const UNCOMPRESSED_PUBLIC_KEY_SIZE: usize = 65;
1919

20+
/// Source of a method call
21+
#[derive(Debug, Copy, Clone)]
22+
pub enum CallSource {
23+
Internal,
24+
External,
25+
}
26+
2027
#[derive(
2128
Eq,
2229
Hash,

src/rpc/methods/state.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ use crate::libp2p::NetworkMessage;
1717
use crate::lotus_json::{LotusJson, lotus_json_with_self};
1818
use crate::networks::{ChainConfig, NetworkChain};
1919
use crate::prelude::*;
20+
use crate::rpc::eth::types::CallSource;
2021
use crate::rpc::registry::actors_reg::load_and_serialize_actor_state;
2122
use crate::shim::actors::market::DealState;
2223
use crate::shim::actors::market::ext::MarketStateExt as _;
@@ -157,7 +158,10 @@ impl RpcMethod<2> for StateReplay {
157158
_: &http::Extensions,
158159
) -> Result<Self::Ok, ServerError> {
159160
let tipset = ctx.chain_store().load_required_tipset_or_heaviest(&tsk)?;
160-
Ok(ctx.state_manager.replay(tipset, message_cid).await?)
161+
Ok(ctx
162+
.state_manager
163+
.replay(tipset, message_cid, CallSource::External)
164+
.await?)
161165
}
162166
}
163167

src/state_manager/execution.rs

Lines changed: 19 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ use super::state_computation::{
77
use super::utils::structured;
88
use super::*;
99
use crate::interpreter::{CalledAt, VMTrace};
10+
use crate::rpc::eth::types::CallSource;
1011
use crate::rpc::state::{ApiInvocResult, MessageGasCost};
1112
use anyhow::{Context as _, bail};
1213
use num_traits::identities::Zero;
@@ -21,9 +22,14 @@ impl StateManager {
2122
/// Lotus, which halts at the target message, this executes the whole
2223
/// tipset — the coalescing depends on it, don't port the halt back.
2324
/// Consequently, failures after the target message also fail the replay.
24-
pub async fn replay(&self, ts: Tipset, mcid: Cid) -> Result<ApiInvocResult, Error> {
25+
pub async fn replay(
26+
&self,
27+
ts: Tipset,
28+
mcid: Cid,
29+
source: CallSource,
30+
) -> Result<ApiInvocResult, Error> {
2531
let (_, trace) = self
26-
.execution_trace(&ts)
32+
.execution_trace(&ts, source)
2733
.await
2834
.map_err(|e| Error::Other(format!("unexpected error during execution : {e}")))?;
2935
trace
@@ -202,20 +208,29 @@ impl StateManager {
202208
pub async fn execution_trace(
203209
&self,
204210
tipset: &Tipset,
211+
source: CallSource,
205212
) -> anyhow::Result<(Cid, Vec<Arc<ApiInvocResult>>)> {
206213
let key = tipset.key();
207214
let (state_root, invoc_trace) = self
208215
.trace_cache
209-
.get_or_insert_async(key, self.execution_trace_inner(tipset.shallow_clone()))
216+
.get_or_insert_async(
217+
key,
218+
self.execution_trace_inner(tipset.shallow_clone(), source),
219+
)
210220
.await?;
211221
Ok((state_root.into(), invoc_trace))
212222
}
213223

214224
async fn execution_trace_inner(
215225
&self,
216226
tipset: Tipset,
227+
source: CallSource,
217228
) -> anyhow::Result<(CidWrapper, Vec<Arc<ApiInvocResult>>)> {
218-
let permit = self.replay_permit().await;
229+
// Internal calls like cache prefilling should not compete the semaphore
230+
let permit = match source {
231+
CallSource::External => Some(self.replay_permit().await),
232+
CallSource::Internal => None,
233+
};
219234
let this = self.shallow_clone();
220235
tokio::task::spawn_blocking(move || {
221236
let _permit = permit;

src/state_manager/tests.rs

Lines changed: 17 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,10 @@
33

44
use super::*;
55
use crate::db::MemoryDB;
6+
use crate::rpc::eth::types::CallSource;
67
use crate::shim::executor::StampedEvent;
78
use fil_actors_shared::fvm_ipld_amt::Amt;
9+
use rstest::rstest;
810

911
fn create_raw_event_v4(emitter: u64, key: &str) -> fvm_shared4::event::StampedEvent {
1012
fvm_shared4::event::StampedEvent {
@@ -328,8 +330,10 @@ fn state_manager_with_unexecutable_tipset() -> (StateManager, Tipset) {
328330
(sm, ts)
329331
}
330332

331-
#[tokio::test]
332-
async fn replay_is_served_from_the_tipset_trace_cache() {
333+
#[rstest]
334+
#[case(CallSource::External)]
335+
#[case(CallSource::Internal)]
336+
fn replay_is_served_from_the_tipset_trace_cache(#[case] source: CallSource) {
333337
use crate::utils::cid::CidCborExt;
334338

335339
let (sm, ts) = state_manager_with_unexecutable_tipset();
@@ -343,7 +347,7 @@ async fn replay_is_served_from_the_tipset_trace_cache() {
343347
(Cid::default().into(), vec![Arc::new(cached.clone())]),
344348
);
345349

346-
let replayed = sm.replay(ts, mcid).await.unwrap();
350+
let replayed = tokio_test::block_on(sm.replay(ts, mcid, source)).unwrap();
347351
assert_eq!(replayed, cached);
348352
}
349353

@@ -357,18 +361,22 @@ async fn replay_permits_are_sized_by_configured_concurrency() {
357361
assert_eq!(sm.replay_semaphore.available_permits(), permits - 1);
358362
}
359363

360-
#[tokio::test]
361-
async fn replay_of_message_absent_from_cached_trace_fails_without_executing() {
364+
#[rstest]
365+
#[case(CallSource::External)]
366+
#[case(CallSource::Internal)]
367+
fn replay_of_message_absent_from_cached_trace_fails_without_executing(#[case] source: CallSource) {
362368
use crate::utils::cid::CidCborExt;
363369

364370
let (sm, ts) = state_manager_with_unexecutable_tipset();
365371
sm.trace_cache
366372
.insert(ts.key().clone(), (Cid::default().into(), vec![]));
367373

368-
let err = sm
369-
.replay(ts, Cid::from_cbor_blake2b256(&"absent-message").unwrap())
370-
.await
371-
.unwrap_err();
374+
let err = tokio_test::block_on(sm.replay(
375+
ts,
376+
Cid::from_cbor_blake2b256(&"absent-message").unwrap(),
377+
source,
378+
))
379+
.unwrap_err();
372380
// "failed to replay" is the message-not-found contract exposed via RPC.
373381
assert!(
374382
matches!(err, Error::Other(ref s) if s == "failed to replay"),

src/state_manager/utils.rs

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -359,6 +359,8 @@ pub mod state_compute {
359359
use super::*;
360360
#[cfg(feature = "cargo-test")]
361361
use crate::chain_sync::tipset_syncer::validate_tipset;
362+
#[cfg(feature = "cargo-test")]
363+
use crate::rpc::eth::types::CallSource;
362364

363365
#[tokio::test(flavor = "multi_thread")]
364366
async fn test_list_state_snapshot_files() {
@@ -402,7 +404,10 @@ pub mod state_compute {
402404
.expect("test tipset must contain messages")
403405
.cid();
404406

405-
let replayed = sm.replay(ts.clone(), msg_cid).await.unwrap();
407+
let replayed = sm
408+
.replay(ts.clone(), msg_cid, CallSource::External)
409+
.await
410+
.unwrap();
406411
assert_eq!(replayed.msg_cid, msg_cid);
407412

408413
let (_, trace) = sm
@@ -417,7 +422,10 @@ pub mod state_compute {
417422

418423
// A second replay of the same tipset must not re-execute it.
419424
let misses = sm.trace_cache.misses();
420-
let replayed_again = sm.replay(ts.clone(), msg_cid).await.unwrap();
425+
let replayed_again = sm
426+
.replay(ts.clone(), msg_cid, CallSource::External)
427+
.await
428+
.unwrap();
421429
assert_eq!(replayed_again, replayed);
422430
assert_eq!(sm.trace_cache.misses(), misses);
423431
}

0 commit comments

Comments
 (0)