Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 54 additions & 9 deletions tests/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
#![allow(dead_code)]

use std::collections::HashMap;
use std::sync::Arc;
use std::sync::{Arc, Condvar, Mutex};

use alloy_eips::BlockId;
use alloy_primitives::{Address, Bytes, U256, hex};
Expand Down Expand Up @@ -128,19 +128,64 @@ pub fn failing_fetcher() -> StorageBatchFetchFn {
})
}

/// Build a stub [`StorageBatchFetchFn`] that reports chosen values *and* flips a
/// shared flag the first time it is called.
/// A one-shot synchronous gate: a holder blocks in [`wait`](Gate::wait) until
/// some other thread calls [`release`](Gate::release). Cloning shares the same
/// underlying state, and `release` is sticky — once released, every present and
/// future `wait` returns immediately.
///
/// Used by the Drop-abort test to prove the background validator was cancelled
/// before it ever fetched (so it could not have queued a correction). The
/// returned values otherwise behave exactly like [`stub_fetcher`].
pub fn tracking_fetcher(
/// Used by the Drop-abort test to make the background validator's fetch
/// deterministically ordered *after* the drop. The fetcher (running on a worker
/// thread) cannot return — and therefore the validator cannot reach its
/// post-fetch checkpoint — until the test has dropped the `SpeculativeSim` and
/// released the gate, eliminating the spawn/poll race regardless of how the
/// multi-thread scheduler interleaves the two threads.
///
/// Built on a `Mutex<bool>` + `Condvar` so the whole thing is `Send + Sync`,
/// which a [`StorageBatchFetchFn`] closure must be.
#[derive(Clone, Default)]
pub struct Gate {
inner: Arc<(Mutex<bool>, Condvar)>,
}

impl Gate {
pub fn new() -> Self {
Self::default()
}

/// Block until [`release`](Gate::release) has been called (returns
/// immediately if it already has).
pub fn wait(&self) {
let (lock, cv) = &*self.inner;
let mut released = lock.lock().unwrap_or_else(|e| e.into_inner());
while !*released {
released = cv.wait(released).unwrap_or_else(|e| e.into_inner());
}
}

/// Wake any current waiter and let all future waiters pass.
pub fn release(&self) {
let (lock, cv) = &*self.inner;
*lock.lock().unwrap_or_else(|e| e.into_inner()) = true;
cv.notify_all();
}
}

/// Build a stub [`StorageBatchFetchFn`] that reports chosen values but blocks on
/// `gate` before returning.
///
/// Used by the Drop-abort test: by releasing the gate only *after* dropping the
/// `SpeculativeSim`, the test guarantees the validator's fetch completes (and so
/// its post-fetch, correction-queuing checkpoint runs) strictly after the
/// cancel flag is set — so the dropped speculation can never queue a correction,
/// no matter how the scheduler races the two threads. The returned values
/// otherwise behave exactly like [`stub_fetcher`].
pub fn gated_tracking_fetcher(
values: HashMap<(Address, U256), U256>,
called: Arc<std::sync::atomic::AtomicBool>,
gate: Gate,
) -> StorageBatchFetchFn {
Arc::new(
move |requests: Vec<(Address, U256)>, _block: Option<BlockId>| {
called.store(true, std::sync::atomic::Ordering::SeqCst);
gate.wait();
requests
.into_iter()
.map(|(addr, slot)| {
Expand Down
38 changes: 24 additions & 14 deletions tests/freshness.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ use alloy_sol_types::SolCall;
use anyhow::Result;

use common::{
MOCK_ERC20_BALANCE_SLOT, MockERC20, failing_fetcher, install_default_account,
install_mock_erc20, panicking_fetcher, setup_cache, stub_fetcher, tracking_fetcher,
Gate, MOCK_ERC20_BALANCE_SLOT, MockERC20, failing_fetcher, gated_tracking_fetcher,
install_default_account, install_mock_erc20, panicking_fetcher, setup_cache, stub_fetcher,
};
use evm_fork_cache::cache::{
EvmCache, EvmOverlay, SimStatus, SlotObservationTracker, StorageBatchFetchFn,
Expand Down Expand Up @@ -1224,20 +1224,31 @@ async fn run_into_optimistic_aborts_validation() -> Result<()> {

// T3 (part 1): dropping the SpeculativeSim (no validate/into_optimistic) aborts
// the validation task before it can push a correction. The fetcher reports a
// CHANGED value and flips a "called" flag; after the drop + settle we assert the
// pending queue is empty (and, robustly, that the fetcher was never even
// reached) — proving the abort beat the push.
// CHANGED value, so an *uncancelled* validator would queue a correction and bump
// the re-run count; we assert neither happens after the drop.
//
// Determinism: the validator's only correction-queuing path runs *after* its
// fetch returns (the post-fetch cancel checkpoint in `run_validator` gates it).
// We make that ordering race-free with a gate the test controls — the fetcher
// blocks until `gate.release()`, and we release only *after* `drop(sim)` has set
// the cancel flag. So however the multi-thread scheduler interleaves the spawned
// task and this thread, the fetch (and thus the post-fetch checkpoint) can only
// complete once cancellation is already observable, and the correction is
// suppressed. We deliberately do NOT assert the fetcher was never reached: the
// product only guarantees a cancel seen at a checkpoint suppresses side effects,
// not that an in-flight fetch is skipped — asserting the latter was the original
// over-strict, racy condition.
#[tokio::test(flavor = "multi_thread")]
async fn dropping_speculative_sim_aborts_before_queueing_correction() -> Result<()> {
let token = Address::repeat_byte(0x44);
let owner = Address::repeat_byte(0x55);
let recipient = Address::repeat_byte(0x66);

let mut cache = cache_with_balance(token, owner, U256::from(1000)).await?;
let called = Arc::new(std::sync::atomic::AtomicBool::new(false));
cache.set_storage_batch_fetcher(tracking_fetcher(
let gate = Gate::new();
cache.set_storage_batch_fetcher(gated_tracking_fetcher(
HashMap::from([((token, balance_slot_for(owner)), U256::from(50))]),
Arc::clone(&called),
gate.clone(),
));

let mut controller = FreshnessController::new(FreshnessRegistry::new(), AlwaysVerify);
Expand All @@ -1249,9 +1260,12 @@ async fn dropping_speculative_sim_aborts_before_queueing_correction() -> Result<
transfer_calldata(recipient, U256::from(100)),
)],
)?;
// Drop immediately, with NO intervening await, so the abort flag is set
// before the spawned task is ever polled.
// Drop with NO intervening await, then release the gate. Releasing only after
// the drop guarantees the validator's fetch (if it even reaches it) returns
// strictly after the cancel flag is set, so its post-fetch checkpoint bails
// out before queuing anything.
drop(sim);
gate.release();

settle().await;

Expand All @@ -1260,10 +1274,6 @@ async fn dropping_speculative_sim_aborts_before_queueing_correction() -> Result<
0,
"dropping the sim must abort validation before it queues a correction"
);
assert!(
!called.load(std::sync::atomic::Ordering::SeqCst),
"the aborted validator should never have reached the fetcher"
);
assert_eq!(controller.rerun_count(), 0, "no re-run after abort");
Ok(())
}
Expand Down
Loading