#![allow(clippy::unwrap_used, clippy::expect_used)]
mod common;
use alloy_primitives::{Address, U256};
use bal_archive::{Archive, ArchiveConfig, ArchiveError, NotAvailable, Provenance};
use common::*;
use std::collections::BTreeMap;
const A: Address = Address::repeat_byte(0xAA);
const B: Address = Address::repeat_byte(0xBB);
fn world() -> (World, Chain) {
let mut states: BTreeMap<u64, Storage> = BTreeMap::new();
let mut st = Storage::new();
st.insert(slot(1), U256::from(7));
st.insert(slot(3), U256::from(42));
for b in 8..=9 {
states.insert(b, st.clone());
}
st.insert(slot(1), U256::from(100));
for b in 10..=11 {
states.insert(b, st.clone());
}
st.insert(slot(1), U256::from(200));
st.insert(slot(2), U256::from(5));
for b in 12..=13 {
states.insert(b, st.clone());
}
st.insert(slot(1), U256::from(300));
states.insert(14, st.clone());
let world = World { addr: A, states };
let chain = Chain::new();
chain.push(8, A, &[], world.root_at(8), 0);
chain.push(9, A, &[], world.root_at(9), 0);
chain.push(10, A, &[(1, 100)], world.root_at(10), 0);
chain.push(11, A, &[], world.root_at(11), 0);
chain.push(12, A, &[(1, 200), (2, 5)], world.root_at(12), 0);
chain.push(13, A, &[], world.root_at(13), 0);
chain.push(14, A, &[(1, 300)], world.root_at(14), 0);
(world, chain)
}
fn open(cfg: ArchiveConfig) -> (Archive, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
let a = Archive::open_with(dir.path().join("a.redb"), cfg).unwrap();
(a, dir)
}
#[tokio::test]
async fn sync_reads_and_bounds() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
let rep = ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(
(rep.from, rep.to, rep.blocks_applied),
(Some(10), Some(14), 5)
);
assert_eq!(rep.bootstrapped, 2, "s1 at 9 and s2 at 11 proven");
assert_eq!(rep.bootstrap_pending, 0);
assert_eq!(ar.head().unwrap().map(|h| h.0), Some(14));
assert_eq!(
ar.storage_at(B, slot(1), 12).unwrap_err(),
NotAvailable::NotWatched(B)
);
assert_eq!(
ar.storage_at(A, slot(1), 9).unwrap_err(),
NotAvailable::BeforeStart {
requested: 9,
start: 10
}
);
assert_eq!(
ar.storage_at(A, slot(1), 15).unwrap_err(),
NotAvailable::AfterHead {
requested: 15,
head: 14
}
);
let v = |b| ar.storage_at(A, slot(1), b).unwrap();
assert_eq!(
(v(10).value, v(10).provenance, v(10).set_at),
(val(100), Provenance::Bal, 10)
);
assert_eq!((v(11).value, v(11).set_at), (val(100), 10));
assert_eq!(v(12).value, val(200));
assert_eq!(v(13).value, val(200));
assert_eq!(v(14).value, val(300));
let z = ar.storage_at(A, slot(2), 11).unwrap();
assert_eq!(
(z.value, z.provenance, z.set_at),
(val(0), Provenance::Proof, 9)
);
assert_eq!(ar.storage_at(A, slot(2), 12).unwrap().value, val(5));
assert_eq!(
ar.storage_at(A, slot(3), 10).unwrap_err(),
NotAvailable::NotBootstrapped
);
ar.bootstrap_slot(&world, A, slot(3)).await.unwrap();
let s3 = ar.storage_at(A, slot(3), 10).unwrap();
assert_eq!((s3.value, s3.provenance), (val(42), Provenance::Proof));
assert_eq!(ar.storage_at(A, slot(3), 14).unwrap().value, val(42));
ar.bootstrap_slot(&world, A, slot(9)).await.unwrap();
assert_eq!(ar.storage_at(A, slot(9), 12).unwrap().value, val(0));
assert_eq!(ar.changed_slots(A, 12).unwrap(), vec![slot(1), slot(2)]);
assert!(ar.changed_slots(A, 13).unwrap().is_empty());
let h: Vec<(u64, _)> = ar
.history(A, slot(1), 10..15)
.unwrap()
.into_iter()
.map(|e| (e.block, e.value))
.collect();
assert_eq!(h, vec![(10, val(100)), (12, val(200)), (14, val(300))]);
assert_eq!(
ar.history(A, slot(1), 9..12).unwrap_err(),
NotAvailable::BeforeStart {
requested: 9,
start: 10
}
);
let rep2 = ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(rep2.blocks_applied, 0);
}
#[tokio::test]
async fn pending_then_resolved_then_lost() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig {
bootstrap_window: 3,
..Default::default()
});
ar.watch(A, 10).unwrap();
let rep = ar.sync(&chain, None).await.unwrap();
assert_eq!(rep.bootstrap_pending, 2);
assert_eq!(
ar.storage_at(A, slot(2), 11).unwrap_err(),
NotAvailable::BootstrapPending { first_seen: 12 }
);
let rep = ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(
(rep.bootstrapped, rep.bootstrap_lost, rep.bootstrap_pending),
(1, 1, 0)
);
assert_eq!(ar.storage_at(A, slot(2), 11).unwrap().value, val(0));
assert_eq!(
ar.boot_state(A, slot(1)).unwrap(),
Some(bal_archive::BootState::Lost { first_seen: 10 })
);
}
#[tokio::test]
async fn reorg_rolls_back_and_reapplies() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(ar.storage_at(A, slot(1), 14).unwrap().value, val(300));
chain.truncate(12);
chain.push(13, A, &[(1, 999)], world.root_at(13), 1);
chain.push(14, A, &[(2, 6)], world.root_at(14), 1);
let rep = ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(rep.reorged_to, Some(12));
assert_eq!(rep.blocks_applied, 2);
assert_eq!(ar.storage_at(A, slot(1), 13).unwrap().value, val(999));
assert_eq!(ar.storage_at(A, slot(1), 14).unwrap().value, val(999));
assert_eq!(ar.storage_at(A, slot(2), 14).unwrap().value, val(6));
assert_eq!(ar.storage_at(A, slot(2), 13).unwrap().value, val(5));
assert_eq!(ar.changed_slots(A, 14).unwrap(), vec![slot(2)]);
let h: Vec<u64> = ar
.history(A, slot(1), 10..15)
.unwrap()
.into_iter()
.map(|e| e.block)
.collect();
assert_eq!(h, vec![10, 12, 13]);
}
#[tokio::test]
async fn watch_rules_and_verification() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
assert!(matches!(ar.watch(A, 0), Err(ArchiveError::InvalidStart(0))));
ar.watch(A, 10).unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
assert!(matches!(
ar.watch(B, 12),
Err(ArchiveError::StartInPast {
from_block: 12,
head: 14
})
));
ar.watch(B, 15).unwrap();
{
let mut blocks = chain.blocks.lock().unwrap();
let mut b = blocks[&14].clone();
b.header.number = 15;
b.header.parent_hash = blocks[&14].header.hash;
b.header.hash = alloy_primitives::B256::repeat_byte(0x15);
b.header.block_access_list_hash = Some(alloy_primitives::B256::repeat_byte(0xEE));
blocks.insert(15, b);
}
let err = ar.sync(&chain, Some(&world)).await.unwrap_err();
assert!(
matches!(err, ArchiveError::Verification { block: 15, .. }),
"{err}"
);
assert_eq!(
ar.head().unwrap().map(|h| h.0),
Some(14),
"nothing applied past the bad block"
);
}
#[tokio::test]
#[allow(clippy::reversed_empty_ranges)]
async fn inverted_history_range_is_an_error() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(
ar.history(A, slot(1), 13..11).unwrap_err(),
NotAvailable::InvalidRange { start: 13, end: 11 }
);
assert_eq!(
ar.history(A, slot(1), 12..12).unwrap_err(),
NotAvailable::InvalidRange { start: 12, end: 12 }
);
}
struct LyingWorld<'a>(&'a World);
#[async_trait::async_trait]
impl bal_source::StateSource for LyingWorld<'_> {
async fn proof(
&self,
addr: Address,
slots: &[alloy_primitives::B256],
block: u64,
) -> bal_source::Result<bal_source::AccountProof> {
let mut with_extra = slots.to_vec();
with_extra.push(slot(1)); self.0.proof(addr, &with_extra, block).await
}
}
#[tokio::test]
async fn proof_with_unrequested_slot_is_rejected() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
let before = ar.storage_at(A, slot(1), 10).unwrap();
let err = ar
.bootstrap_slot(&LyingWorld(&world), A, slot(3))
.await
.unwrap_err();
assert!(
matches!(
err,
ArchiveError::Proof(bal_source::ProofError::UnexpectedSlot(_))
),
"{err}"
);
assert_eq!(
ar.storage_at(A, slot(1), 10).unwrap(),
before,
"s1 untouched"
);
assert_eq!(
ar.storage_at(A, slot(3), 10).unwrap_err(),
NotAvailable::NotBootstrapped
);
}
#[tokio::test]
async fn lazy_bootstrap_never_stores_a_post_value() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig {
bootstrap_window: 0,
..Default::default()
});
ar.watch(A, 10).unwrap();
ar.sync(&chain, None).await.unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
assert!(matches!(
ar.storage_at(A, slot(2), 11),
Err(NotAvailable::BootstrapLost { .. })
));
ar.bootstrap_slot(&world, A, slot(2)).await.unwrap();
assert!(matches!(
ar.storage_at(A, slot(2), 11),
Err(NotAvailable::BootstrapLost { .. })
));
}
#[tokio::test]
async fn watch_below_in_flight_block_is_refused() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
assert!(matches!(
ar.watch(B, 14),
Err(ArchiveError::StartInPast { head: 14, .. })
));
ar.watch(B, 15).unwrap();
}
#[tokio::test]
async fn reorg_below_start_drops_proofs() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(B, 10).unwrap();
ar.watch(A, 12).unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(ar.storage_at(A, slot(1), 12).unwrap().value, val(200));
ar.bootstrap_slot(&world, A, slot(3)).await.unwrap();
assert_eq!(
ar.storage_at(A, slot(3), 12).unwrap().provenance,
Provenance::Proof
);
chain.truncate(10);
for b in 11..=14 {
chain.push(b, A, &[], world.root_at(b), 7);
}
let rep = ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(rep.reorged_to, Some(10));
assert_eq!(
ar.boot_state(A, slot(3)).unwrap(),
None,
"orphaned-branch proof dropped"
);
assert_eq!(
ar.storage_at(A, slot(3), 12).unwrap_err(),
NotAvailable::NotBootstrapped
);
}
struct SlowChain<'a> {
inner: &'a Chain,
gate: tokio::sync::Semaphore,
}
#[async_trait::async_trait]
impl bal_source::BalSource for SlowChain<'_> {
async fn head(&self) -> bal_source::Result<u64> {
self.inner.head().await
}
async fn finalized(&self) -> bal_source::Result<u64> {
self.inner.finalized().await
}
async fn block(&self, n: u64) -> bal_source::Result<bal_source::SourcedBlock> {
let _p = self.gate.acquire().await.expect("semaphore open");
self.inner.block(n).await
}
}
#[tokio::test]
async fn concurrent_sync_is_refused_then_allowed() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
let slow = SlowChain {
inner: &chain,
gate: tokio::sync::Semaphore::new(0),
};
let first = ar.sync(&slow, Some(&world));
tokio::pin!(first);
tokio::select! {
_ = &mut first => panic!("first pass must be parked"),
_ = tokio::time::sleep(std::time::Duration::from_millis(50)) => {}
}
let second = ar.sync(&chain, Some(&world)).await;
assert!(
matches!(second, Err(ArchiveError::SyncInProgress)),
"{second:?}"
);
slow.gate.add_permits(100);
first.await.unwrap();
let rep = ar.sync(&chain, Some(&world)).await.unwrap();
assert_eq!(rep.blocks_applied, 0);
assert_eq!(ar.head().unwrap().map(|h| h.0), Some(14));
}
struct NoProofs;
#[async_trait::async_trait]
impl bal_source::StateSource for NoProofs {
async fn proof(
&self,
_: Address,
_: &[alloy_primitives::B256],
block: u64,
) -> bal_source::Result<bal_source::AccountProof> {
Err(bal_source::SourceError::Rpc {
code: -32602,
message: format!("distance to target block {block} exceeds maximum proof window"),
})
}
}
#[tokio::test]
async fn backup_source_rescues_bootstrap() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
let rep = ar.sync(&chain, Some(&NoProofs)).await.unwrap();
assert_eq!((rep.bootstrapped, rep.bootstrap_pending), (0, 2));
let (ar2, _d2) = open(ArchiveConfig::default());
ar2.watch(A, 10).unwrap();
let with_backup = bal_source::Fallback::new(NoProofs, &world);
let rep = ar2.sync(&chain, Some(&with_backup)).await.unwrap();
assert_eq!((rep.bootstrapped, rep.bootstrap_pending), (2, 0));
let z = ar2.storage_at(A, slot(2), 11).unwrap();
assert_eq!((z.value, z.provenance), (val(0), Provenance::Proof));
let lying = bal_source::Fallback::new(NoProofs, LyingWorld(&world));
let err = ar2.bootstrap_slot(&lying, A, slot(3)).await.unwrap_err();
assert!(matches!(err, ArchiveError::Proof(_)), "{err}");
}
#[tokio::test]
async fn watch_twice_and_config_mismatch() {
let (world, chain) = world();
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("a.redb");
{
let ar = Archive::open_with(path.clone(), ArchiveConfig::default()).unwrap();
ar.watch(A, 10).unwrap();
ar.watch(A, 10).unwrap(); assert!(matches!(
ar.watch(A, 11),
Err(ArchiveError::AlreadyWatched { from_block: 10, .. })
));
ar.sync(&chain, Some(&world)).await.unwrap();
}
let err = Archive::open_with(
path.clone(),
ArchiveConfig {
full_detail: true,
..Default::default()
},
)
.err()
.expect("full_detail mismatch must be refused");
assert!(
matches!(
err,
ArchiveError::ConfigMismatch {
option: "full_detail",
..
}
),
"{err}"
);
let ar = Archive::open_with(path, ArchiveConfig::default()).unwrap();
assert_eq!(ar.storage_at(A, slot(1), 14).unwrap().value, val(300));
}
#[tokio::test]
async fn unwatch_removes_everything() {
let (world, chain) = world();
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 10).unwrap();
ar.sync(&chain, Some(&world)).await.unwrap();
ar.unwatch(A).unwrap();
assert_eq!(
ar.storage_at(A, slot(1), 12).unwrap_err(),
NotAvailable::NotWatched(A)
);
assert!(ar.watchlist().unwrap().is_empty());
assert_eq!(ar.boot_state(A, slot(1)).unwrap(), None);
}
#[tokio::test]
async fn creation_settles_every_pre_value_without_proofs() {
let chain = Chain::new();
chain.push(8, A, &[], val(0), 0);
chain.push(9, A, &[], val(0), 0);
chain.push_with(10, A, &[(1, 100)], val(0), 0, true);
chain.push(11, A, &[], val(0), 0);
chain.push(12, A, &[(2, 5)], val(0), 0);
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 9).unwrap();
let rep = ar.sync(&chain, None).await.unwrap();
assert_eq!((rep.bootstrap_pending, rep.bootstrap_lost), (0, 0));
assert_eq!(ar.created_at(A).unwrap(), Some(10));
assert_eq!(ar.storage_at(A, slot(1), 12).unwrap().value, val(100));
let s2 = ar.storage_at(A, slot(2), 11).unwrap();
assert_eq!(
(s2.value, s2.provenance, s2.set_at),
(val(0), Provenance::Bal, 10)
);
assert_eq!(ar.storage_at(A, slot(2), 12).unwrap().value, val(5));
assert_eq!(ar.storage_at(A, slot(77), 9).unwrap().value, val(0));
let st = ar.stats().unwrap();
assert_eq!(st.created, vec![(A, 10)]);
assert_eq!((st.slots_pending, st.slots_lost, st.slots_done), (0, 0, 2));
}
#[tokio::test]
async fn backfill_walks_back_to_creation() {
use bal_archive::{BackfillOpts, BackfillStop};
let chain = Chain::new();
chain.push_with(8, A, &[(1, 7), (3, 42)], val(0), 0, true);
chain.push(9, A, &[], val(0), 0);
chain.push(10, A, &[(1, 100)], val(0), 0);
chain.push(11, A, &[], val(0), 0);
chain.push(12, A, &[(1, 200), (2, 5)], val(0), 0);
chain.push(13, A, &[], val(0), 0);
chain.push(14, A, &[(1, 300)], val(0), 0);
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 13).unwrap();
ar.sync(&chain, None).await.unwrap();
assert_eq!(
ar.storage_at(A, slot(1), 13).unwrap_err(),
NotAvailable::BootstrapPending { first_seen: 14 }
);
assert!(matches!(
ar.storage_at(A, slot(3), 13).unwrap_err(),
NotAvailable::NotBootstrapped
));
let rep = ar
.backfill(
&chain,
A,
BackfillOpts {
resolve_only: true,
..Default::default()
},
)
.await
.unwrap();
assert_eq!(rep.stopped, BackfillStop::Resolved);
assert_eq!(
(rep.from, rep.to, rep.blocks_scanned, rep.slots_resolved),
(13, 12, 1, 1)
);
let v = ar.storage_at(A, slot(1), 13).unwrap();
assert_eq!(
(v.value, v.provenance, v.set_at),
(val(200), Provenance::Bal, 12)
);
assert_eq!(ar.storage_at(A, slot(2), 12).unwrap().value, val(5));
assert_eq!(
ar.storage_at(A, slot(1), 12).unwrap().value,
val(200),
"start moved down to 12"
);
assert_eq!(
ar.storage_at(A, slot(1), 11).unwrap_err(),
NotAvailable::BeforeStart {
requested: 11,
start: 12
}
);
let h = ar.history(A, slot(1), 12..15).unwrap();
assert_eq!(
h.iter().map(|e| (e.block, e.value)).collect::<Vec<_>>(),
vec![(12, val(200)), (14, val(300))]
);
let rep = ar
.backfill(
&chain,
A,
BackfillOpts {
max_blocks: Some(2),
..Default::default()
},
)
.await
.unwrap();
assert_eq!((rep.stopped, rep.to), (BackfillStop::Budget, 10));
assert_eq!(ar.storage_at(A, slot(1), 10).unwrap().value, val(100));
assert_eq!(
ar.storage_at(A, slot(2), 10).unwrap_err(),
NotAvailable::BootstrapPending { first_seen: 12 }
);
let rep = ar
.backfill(&chain, A, BackfillOpts::default())
.await
.unwrap();
assert_eq!(
(rep.stopped, rep.to, rep.created_at),
(BackfillStop::Creation(8), 8, Some(8))
);
assert_eq!(rep.unresolved, 0);
assert_eq!(ar.storage_at(A, slot(1), 9).unwrap().value, val(7));
assert_eq!(ar.storage_at(A, slot(3), 13).unwrap().value, val(42));
let z = ar.storage_at(A, slot(2), 11).unwrap();
assert_eq!((z.value, z.provenance), (val(0), Provenance::Bal));
assert_eq!(ar.stats().unwrap().slots_pending, 0);
let rep = ar
.backfill(&chain, A, BackfillOpts::default())
.await
.unwrap();
assert_eq!(rep.stopped, BackfillStop::Nothing);
}
#[tokio::test]
async fn backfill_refuses_a_broken_chain() {
use bal_archive::BackfillOpts;
let chain = Chain::new();
chain.push(8, A, &[(1, 1)], val(0), 0);
chain.push(9, A, &[(1, 2)], val(0), 0);
chain.push(10, A, &[(1, 3)], val(0), 0);
chain.push(11, A, &[], val(0), 0);
chain.push(12, A, &[(1, 4)], val(0), 0);
let (ar, _d) = open(ArchiveConfig::default());
ar.watch(A, 11).unwrap();
ar.sync(&chain, None).await.unwrap();
chain.push(9, A, &[(1, 999)], val(0), 7);
let err = ar
.backfill(&chain, A, BackfillOpts::default())
.await
.unwrap_err();
assert!(matches!(err, ArchiveError::InconsistentSource(10)), "{err}");
assert_eq!(ar.storage_at(A, slot(1), 10).unwrap().value, val(3));
assert_eq!(ar.watchlist().unwrap(), vec![(A, 10)]);
assert!(matches!(
ar.backfill(&chain, B, BackfillOpts::default())
.await
.unwrap_err(),
ArchiveError::NotWatched(_)
));
}