diff --git a/CHANGELOG.md b/CHANGELOG.md index 2df28716145..3a20a3c324c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ ### 2025-05-22 - Add block snapshots with diff layers. [2664](https://github.com/lambdaclass/ethrex/pull/2664) +- Make disk layer in snapshots use the database. [2848](https://github.com/lambdaclass/ethrex/pull/2848) ### 2025-05-20 diff --git a/crates/l2/prover/bench/src/rpc/mod.rs b/crates/l2/prover/bench/src/rpc/mod.rs index 4321f53ce29..d0588c2d7b3 100644 --- a/crates/l2/prover/bench/src/rpc/mod.rs +++ b/crates/l2/prover/bench/src/rpc/mod.rs @@ -208,7 +208,7 @@ pub async fn get_account( let trie = Trie::from_nodes(Some(root), &other) .map_err(|err| format!("failed to build account proof trie: {err}"))?; if trie - .get(&hash_address(address)) + .get(hash_address(address)) .map_err(|err| format!("failed get account from proof trie: {err}"))? .is_none() { diff --git a/crates/storage/api.rs b/crates/storage/api.rs index 7386e850d0e..6e5046ca4e4 100644 --- a/crates/storage/api.rs +++ b/crates/storage/api.rs @@ -316,6 +316,13 @@ pub trait StoreEngine: Debug + Send + Sync + RefUnwindSafe { /// Clears all checkpoint data created during the last snap sync async fn clear_snap_state(&self) -> Result<(), StoreError>; + /// Write an account batch into the current state snapshot. Blocking non-async version. + fn write_snapshot_account_batch_blocking( + &self, + account_hashes: Vec, + account_states: Vec, + ) -> Result<(), StoreError>; + /// Write an account batch into the current state snapshot async fn write_snapshot_account_batch( &self, @@ -339,6 +346,14 @@ pub trait StoreEngine: Debug + Send + Sync + RefUnwindSafe { storage_values: Vec>, ) -> Result<(), StoreError>; + /// Write multiple storage batches belonging to different accounts into the current storage snapshot. Blocking non-async version. + fn write_snapshot_storage_batches_blocking( + &self, + account_hashes: Vec, + storage_keys: Vec>, + storage_values: Vec>, + ) -> Result<(), StoreError>; + /// Set the latest root of the rebuilt state trie and the last downloaded hashes from each segment async fn set_state_trie_rebuild_checkpoint( &self, @@ -374,6 +389,16 @@ pub trait StoreEngine: Debug + Send + Sync + RefUnwindSafe { account_hash: H256, ) -> Result, StoreError>; + /// Gets a single account from the snapshot state. + fn get_account_snapshot(&self, account_hash: H256) -> Result, StoreError>; + + /// Gets a single storage value from the snapshot state. + fn get_storage_snapshot( + &self, + account_hash: H256, + storage_hash: H256, + ) -> Result, StoreError>; + /// The `forkchoice_update` and `new_payload` methods require the `latest_valid_hash` /// when processing an invalid payload. To provide this, we must track invalid chains. /// diff --git a/crates/storage/snapshot/difflayer.rs b/crates/storage/snapshot/difflayer.rs index 9c9a2c03932..d2e5b55968a 100644 --- a/crates/storage/snapshot/difflayer.rs +++ b/crates/storage/snapshot/difflayer.rs @@ -20,7 +20,7 @@ pub struct DiffLayer { state_root: H256, stale: bool, accounts: HashMap>, // None if deleted - storage: HashMap>, + storage: HashMap>>, /// tracks all diffed items up to disk layer pub(crate) diffed: Bloom, } @@ -32,7 +32,7 @@ impl DiffLayer { block_hash: BlockHash, state_root: H256, accounts: HashMap>, - storage: HashMap>, + storage: HashMap>>, ) -> Self { DiffLayer { origin: origin.clone(), @@ -97,7 +97,7 @@ impl DiffLayer { // If bloom misses we can skip diff layers if !hit { - return self.origin.get_account(hash, layers); + return self.origin.get_account(hash); } // Start traversing layers. @@ -119,7 +119,7 @@ impl DiffLayer { .get(&account_hash) .and_then(|x| x.get(&storage_hash)) { - return Ok(Some(*value)); + return Ok(*value); } let bloom_hash = account_hash ^ storage_hash; @@ -129,7 +129,7 @@ impl DiffLayer { // If bloom misses we can skip diff layers if !hit { - return self.origin.get_storage(account_hash, storage_hash, layers); + return self.origin.get_storage(account_hash, storage_hash); } // Start traversing layers. @@ -155,7 +155,7 @@ impl DiffLayer { block: BlockHash, state_root: H256, accounts: HashMap>, - storage: HashMap>, + storage: HashMap>>, ) -> DiffLayer { let mut layer = DiffLayer::new( self.block_hash, @@ -196,7 +196,7 @@ impl DiffLayer { // delegate to parent match &layers[&self.parent] { - Layer::DiskLayer(disk_layer) => disk_layer.get_account(hash, layers), + Layer::DiskLayer(disk_layer) => disk_layer.get_account(hash), Layer::DiffLayer(diff_layer) => diff_layer .read() .map_err(|error| SnapshotError::LockError(error.to_string()))? @@ -220,14 +220,12 @@ impl DiffLayer { .get(&account_hash) .and_then(|x| x.get(&storage_hash)) { - return Ok(Some(*value)); + return Ok(*value); } // delegate to parent match &layers[&self.parent] { - Layer::DiskLayer(disk_layer) => { - disk_layer.get_storage(account_hash, storage_hash, layers) - } + Layer::DiskLayer(disk_layer) => disk_layer.get_storage(account_hash, storage_hash), Layer::DiffLayer(diff_layer) => diff_layer .read() .map_err(|error| SnapshotError::LockError(error.to_string()))? @@ -239,7 +237,7 @@ impl DiffLayer { self.accounts.extend(accounts); } - pub fn add_storage(&mut self, storage: HashMap>) { + pub fn add_storage(&mut self, storage: HashMap>>) { for (address, st) in storage.iter() { let entry = self.storage.entry(*address).or_default(); entry.extend(st); @@ -250,7 +248,7 @@ impl DiffLayer { self.accounts.clone() } - pub fn storage(&self) -> HashMap> { + pub fn storage(&self) -> HashMap>> { self.storage.clone() } diff --git a/crates/storage/snapshot/disklayer.rs b/crates/storage/snapshot/disklayer.rs index 5eba966b103..34c22db38a2 100644 --- a/crates/storage/snapshot/disklayer.rs +++ b/crates/storage/snapshot/disklayer.rs @@ -5,6 +5,7 @@ use std::{ atomic::{AtomicBool, Ordering}, Arc, }, + time::Instant, }; use crate::api::StoreEngine; @@ -13,8 +14,9 @@ use ethrex_common::{ H256, U256, }; use ethrex_rlp::decode::RLPDecode; +use tracing::info; -use super::{cache::DiskCache, difflayer::DiffLayer, error::SnapshotError, tree::Layers}; +use super::{cache::DiskCache, difflayer::DiffLayer, error::SnapshotError}; /// A disk layer is the bottom most layer. /// @@ -27,6 +29,7 @@ pub struct DiskLayer { pub(super) block_hash: BlockHash, pub(super) state_root: H256, pub(super) stale: Arc, + pub(super) generating: Arc, } impl fmt::Debug for DiskLayer { @@ -34,7 +37,8 @@ impl fmt::Debug for DiskLayer { f.debug_struct("DiskLayer") .field("db", &self.db) .field("cache", &self.cache) - .field("root", &self.state_root) + .field("block_hash", &self.block_hash) + .field("state_root", &self.state_root) .field("stale", &self.stale) .finish_non_exhaustive() } @@ -46,92 +50,48 @@ impl DiskLayer { block_hash, state_root, db, - cache: DiskCache::new(10000, 10000), + cache: DiskCache::new(20000, 40000), stale: Arc::new(AtomicBool::new(false)), + generating: Arc::new(AtomicBool::new(false)), } } -} -impl DiskLayer { pub fn root(&self) -> H256 { self.state_root } - pub fn get_account( - &self, - hash: H256, - _layers: &Layers, - ) -> Result, SnapshotError> { + pub fn get_account(&self, hash: H256) -> Result, SnapshotError> { // Try to get the account from the cache. if let Some(value) = self.cache.accounts.get(&hash) { return Ok(value.clone()); } - // TODO: Right now we use the state trie, but the disk layer should use - // it's own database table of snapshots for faster lookup. - let state_trie = self.db.open_state_trie(self.state_root); - - let value = if let Some(value) = state_trie - .get(hash) - .ok() - .flatten() - .map(|x| AccountState::decode(&x)) - { - value - } else { - self.cache.accounts.insert(hash, None); - return Ok(None); - }; + // TODO: check that snapshot is done to make sure None is None? + let account = self.db.get_account_snapshot(hash)?; - let value: AccountState = value?; + self.cache.accounts.insert(hash, account.clone()); - self.cache.accounts.insert(hash, value.clone().into()); - - Ok(Some(value)) + Ok(account) } pub fn get_storage( &self, account_hash: H256, storage_hash: H256, - layers: &Layers, ) -> Result, SnapshotError> { // Look into the cache first. if let Some(value) = self.cache.storages.get(&(account_hash, storage_hash)) { return Ok(value); } - let account = if let Some(account) = self.get_account(account_hash, layers)? { - account - } else { - self.cache - .storages - .insert((account_hash, storage_hash), None); - return Ok(None); - }; - - // TODO: Right now we use the storage trie, but the disk layer should use - // it's own database table of snapshots for faster lookup. - - let storage_trie = self - .db - .open_storage_trie(account_hash, account.storage_root); - - let value = if let Some(value) = storage_trie.get(storage_hash).ok().flatten() { - value - } else { - self.cache - .storages - .insert((account_hash, storage_hash), None); - return Ok(None); - }; - let value: U256 = U256::decode(&value)?; + // TODO: check that snapshot is done to make sure None is None? + let value = self.db.get_storage_snapshot(account_hash, storage_hash)?; self.cache .storages - .insert((account_hash, storage_hash), Some(value)); + .insert((account_hash, storage_hash), value); - Ok(Some(value)) + Ok(value) } pub fn block_hash(&self) -> H256 { @@ -143,7 +103,7 @@ impl DiskLayer { block_hash: BlockHash, state_root: H256, accounts: HashMap>, - storage: HashMap>, + storage: HashMap>>, ) -> DiffLayer { let mut layer = DiffLayer::new( self.block_hash, @@ -166,4 +126,102 @@ impl DiskLayer { pub fn mark_stale(&self) -> bool { self.stale.swap(true, Ordering::SeqCst) } + + // Starts in a blocking task the disk layer generation. + pub fn start_generating(self: &Arc) { + let layer = (*self).clone(); + tokio::task::spawn_blocking(move || layer.generate()); + } + + fn generate(self: Arc) { + // Note: this method can call blocking methods because it's run outside the main thread. + + // todo: we should be able to stop mid generation in case disk layer changes? + self.generating.store(true, Ordering::SeqCst); + info!("Disk layer generating"); + let start = Instant::now(); + + let state_trie = self.db.open_state_trie(self.state_root); + + let account_iter = state_trie.into_iter().content().map_while(|(path, value)| { + Some((H256::from_slice(&path), AccountState::decode(&value).ok()?)) + }); + + let mut account_hashes = Vec::with_capacity(1024); + let mut account_states = Vec::with_capacity(1024); + let mut storage_keys: Vec> = Vec::with_capacity(64); + let mut storage_values: Vec> = Vec::with_capacity(64); + + // buffers + let mut keys = Vec::with_capacity(32); + let mut values = Vec::with_capacity(32); + + // TODO: figure out optimal + // Write to db every ACCOUNT_BATCH's accounts processed. + const ACCOUNT_BATCH: usize = 100; + + for (hash, state) in account_iter { + keys.clear(); + values.clear(); + + account_hashes.push(hash); + let storage_root = state.storage_root; + account_states.push(state); + + let storage_trie = self.db.open_storage_trie(hash, storage_root); + let storage_iter = storage_trie.into_iter().content(); + + for (storage_hash, value) in storage_iter { + keys.push(H256::from_slice(&storage_hash)); + values.push(U256::from_big_endian(&value)); + } + + storage_keys.push(keys.clone()); + storage_values.push(values.clone()); + + if account_hashes.len() >= ACCOUNT_BATCH { + self.db + .write_snapshot_account_batch_blocking( + account_hashes.clone(), + account_states.clone(), + ) + .expect("convert into a error"); + self.db + .write_snapshot_storage_batches_blocking( + account_hashes.clone(), + storage_keys.clone(), + storage_values.clone(), + ) + .expect("convert into a error"); + account_hashes.clear(); + storage_keys.clear(); + storage_values.clear(); + } + } + + if !account_hashes.is_empty() { + self.db + .write_snapshot_account_batch_blocking( + account_hashes.clone(), + account_states.clone(), + ) + .expect("convert into a error"); + self.db + .write_snapshot_storage_batches_blocking( + account_hashes.clone(), + storage_keys.clone(), + storage_values.clone(), + ) + .expect("convert into a error"); + account_hashes.clear(); + storage_keys.clear(); + storage_values.clear(); + } + + self.generating.store(false, Ordering::SeqCst); + info!( + "Disk layer generation complete, done in {:?}", + start.elapsed() + ); + } } diff --git a/crates/storage/snapshot/error.rs b/crates/storage/snapshot/error.rs index 04436204f59..6a4768406a2 100644 --- a/crates/storage/snapshot/error.rs +++ b/crates/storage/snapshot/error.rs @@ -1,6 +1,8 @@ use ethrex_common::H256; use ethrex_rlp::error::RLPDecodeError; +use crate::error::StoreError; + #[derive(Debug, thiserror::Error)] pub enum SnapshotError { #[error("SnapShot with root {0} not found.")] @@ -19,4 +21,6 @@ pub enum SnapshotError { RLPDecodeError(#[from] RLPDecodeError), #[error("Error getting a lock: {0}")] LockError(String), + #[error(transparent)] + StoreError(#[from] StoreError), } diff --git a/crates/storage/snapshot/tree.rs b/crates/storage/snapshot/tree.rs index 43b7a6d6704..43b23fb4a10 100644 --- a/crates/storage/snapshot/tree.rs +++ b/crates/storage/snapshot/tree.rs @@ -45,32 +45,54 @@ impl SnapshotTree { } } + pub fn add_data( + &self, + account_hashes: Vec, + account_states: Vec, + storage_keys: Vec>, + storage_values: Vec>, + ) { + self.db + .write_snapshot_account_batch_blocking(account_hashes.clone(), account_states) + .expect("convert into a error"); + self.db + .write_snapshot_storage_batches_blocking(account_hashes, storage_keys, storage_values) + .expect("convert into a error"); + } + /// Rebuilds the tree, marking all current layers stale, creating a new base disk layer from the given root. - pub fn rebuild(&self, block_hash: BlockHash, state_root: H256) -> Result<(), SnapshotError> { - let mut layers = self - .layers - .write() - .map_err(|error| SnapshotError::LockError(error.to_string()))?; + pub fn rebuild( + &self, + block_hash: BlockHash, + state_root: H256, + generate: bool, + ) -> Result<(), SnapshotError> { + let disk = { + let mut layers = self + .layers + .write() + .map_err(|error| SnapshotError::LockError(error.to_string()))?; - for layer in layers.values() { - match layer { - Layer::DiskLayer(disk_layer) => disk_layer.mark_stale(), - Layer::DiffLayer(diff_layer) => diff_layer - .write() - .map_err(|error| SnapshotError::LockError(error.to_string()))? - .mark_stale(), - }; + for layer in layers.values() { + match layer { + Layer::DiskLayer(disk_layer) => disk_layer.mark_stale(), + Layer::DiffLayer(diff_layer) => diff_layer + .write() + .map_err(|error| SnapshotError::LockError(error.to_string()))? + .mark_stale(), + }; + } + + layers.clear(); + let disk = Arc::new(DiskLayer::new(self.db.clone(), block_hash, state_root)); + layers.insert(block_hash, Layer::DiskLayer(disk.clone())); + disk + }; + + if generate { + disk.start_generating(); } - layers.clear(); - layers.insert( - block_hash, - Layer::DiskLayer(Arc::new(DiskLayer::new( - self.db.clone(), - block_hash, - state_root, - ))), - ); Ok(()) } @@ -85,7 +107,7 @@ impl SnapshotTree { block_state_root: H256, parent_block_hash: H256, accounts: HashMap>, - storage: HashMap>, + storage: HashMap>>, ) -> Result<(), SnapshotError> { info!("Creating new diff snapshot"); if block_hash == parent_block_hash { @@ -342,6 +364,11 @@ impl SnapshotTree { /// /// Returns Err if the current disk layer is already marked stale. fn save_diff(&self, diff: Layer) -> Result, SnapshotError> { + // Note: Interacting with the db should be made through non-async methods, converting this method + // into an async one, will give problems, mainly with holding RwLocks across await points. + // + // This method is called from `cap` and related mainly, which they themselves are called outside the main thread, so blocking is not bad. + let diff = match diff { Layer::DiskLayer(disk_layer) => { return Err(SnapshotError::SnapshotIsdiskLayer(disk_layer.state_root)) @@ -360,32 +387,68 @@ impl SnapshotTree { // TODO: here we should save the diff to the db (in the future snapshots table) too. let accounts = diff_value.accounts(); - // TODO: Need to make sure it's correct to leave the cache as is - prev_disk.cache.accounts.clear(); + let mut account_hashes = Vec::with_capacity(accounts.len()); + let mut account_states = Vec::with_capacity(accounts.len()); for (hash, acc) in accounts.iter() { - prev_disk.cache.accounts.insert(*hash, acc.clone()); + if let Some(acc) = acc { + // TODO: Important, if acc is None it means it comes from a account update + // with the removed flag, should we remove it from db too? + account_hashes.push(*hash); + account_states.push(acc.clone()); + prev_disk.cache.accounts.insert(*hash, Some(acc.clone())); + } else { + prev_disk.cache.accounts.remove(hash); + } } - // TODO: Need to make sure it's correct to leave the cache as is - prev_disk.cache.storages.clear(); + prev_disk + .db + .write_snapshot_account_batch_blocking(account_hashes, account_states)?; let storage = diff_value.storage(); + + let mut account_hashes = Vec::with_capacity(storage.len()); + let mut storage_keys = Vec::with_capacity(storage.len()); + let mut storage_values = Vec::with_capacity(storage.len()); + for (account_hash, storage) in storage.iter() { + account_hashes.push(*account_hash); + let mut keys = Vec::new(); + let mut values = Vec::new(); for (storage_hash, value) in storage.iter() { - prev_disk - .cache - .storages - .insert((*account_hash, *storage_hash), Some(*value)); + // TODO: Important, if acc is None it means it had a value of zero should we remove it from db too? + if let Some(value) = &value { + values.push(*value); + keys.push(*storage_hash); + prev_disk + .cache + .storages + .insert((*account_hash, *storage_hash), Some(*value)); + } else { + prev_disk + .cache + .storages + .remove(&(*account_hash, *storage_hash)); + } } + storage_values.push(values); + storage_keys.push(keys); } + prev_disk.db.write_snapshot_storage_batches_blocking( + account_hashes, + storage_keys, + storage_values, + )?; + let disk = DiskLayer { db: self.db.clone(), cache: prev_disk.cache.clone(), block_hash: diff_value.block_hash(), state_root: diff_value.root(), stale: Arc::new(AtomicBool::new(false)), + generating: prev_disk.generating.clone(), }; Ok(Arc::new(disk)) } @@ -412,7 +475,7 @@ impl SnapshotTree { let address = hash_address_fixed(&address); match snapshot { - Layer::DiskLayer(snapshot) => snapshot.get_account(address, &layers), + Layer::DiskLayer(snapshot) => snapshot.get_account(address), Layer::DiffLayer(snapshot) => snapshot .read() .map_err(|error| SnapshotError::LockError(error.to_string()))? @@ -432,7 +495,7 @@ impl SnapshotTree { block_hash: BlockHash, address: Address, storage_key: H256, - ) -> Result>, SnapshotError> { + ) -> Result, SnapshotError> { debug!( "called get_storage_at_hash with block {} address {} key {}", block_hash, address, storage_key @@ -444,27 +507,18 @@ impl SnapshotTree { .map_err(|error| SnapshotError::LockError(error.to_string()))?; let address = hash_address_fixed(&address); - let (value, origin_block_hash) = match snapshot { - Layer::DiskLayer(snapshot) => ( - snapshot.get_storage(address, storage_key, &layers)?, - snapshot.block_hash, - ), + let value = match snapshot { + Layer::DiskLayer(snapshot) => snapshot.get_storage(address, storage_key)?, Layer::DiffLayer(snapshot) => { let snapshot = snapshot .read() .map_err(|error| SnapshotError::LockError(error.to_string()))?; - ( - snapshot.get_storage(address, storage_key, &layers)?, - snapshot.origin().block_hash, - ) + + snapshot.get_storage(address, storage_key, &layers)? } }; - if value.is_none() && block_hash != origin_block_hash { - return Ok(None); - } else { - return Ok(Some(value)); - } + return Ok(value); } Err(SnapshotError::SnapshotNotFound(block_hash)) @@ -530,7 +584,6 @@ impl SnapshotTree { parent_value.storage(), ); - // TODO: should we rebloom here? layer.diffed = layer_value.diffed(); Ok(Layer::DiffLayer(Arc::new(RwLock::new(layer)))) @@ -551,8 +604,8 @@ mod tests { SnapshotTree::new(db) } - #[test] - fn test_add_single_account_in_single_difflayer() { + #[tokio::test] + async fn test_add_single_account_in_single_difflayer() { let tree = create_mock_tree(); let root = H256::from_low_u64_be(1); @@ -567,7 +620,7 @@ mod tests { }; // Add a disklayer to the tree - tree.rebuild(H256::zero(), H256::zero()).unwrap(); + tree.rebuild(H256::zero(), H256::zero(), false).unwrap(); // Add a single account in a single difflayer tree.update( @@ -584,10 +637,10 @@ mod tests { assert_eq!(retrieved_account, Some(account_state)); } - #[test] - fn test_add_two_accounts_in_different_difflayers() { + #[tokio::test] + async fn test_add_two_accounts_in_different_difflayers() { let tree = create_mock_tree(); - tree.rebuild(H256::zero(), H256::zero()).unwrap(); + tree.rebuild(H256::zero(), H256::zero(), false).unwrap(); let root1 = H256::from_low_u64_be(1); let root2 = H256::from_low_u64_be(2); @@ -638,10 +691,10 @@ mod tests { assert_eq!(retrieved_account2, Some(account2_state)); } - #[test] - fn test_override_account_in_second_difflayer() { + #[tokio::test] + async fn test_override_account_in_second_difflayer() { let tree = create_mock_tree(); - tree.rebuild(H256::zero(), H256::zero()).unwrap(); + tree.rebuild(H256::zero(), H256::zero(), false).unwrap(); let root1 = H256::from_low_u64_be(1); let root2 = H256::from_low_u64_be(2); let address = Address::from_low_u64_be(1); @@ -690,10 +743,10 @@ mod tests { assert_eq!(retrieved_account, Some(account_state1)); } - #[test] - fn test_override_account_storage_flattening() { + #[tokio::test] + async fn test_override_account_storage_flattening() { let tree = create_mock_tree(); - tree.rebuild(H256::zero(), H256::zero()).unwrap(); + tree.rebuild(H256::zero(), H256::zero(), false).unwrap(); let root1 = H256::from_low_u64_be(1); let root2 = H256::from_low_u64_be(2); @@ -724,8 +777,8 @@ mod tests { H256::zero(), HashMap::from([(account_hash, Some(account_state1.clone()))]), HashMap::from([(account_hash, { - let mut map: HashMap = HashMap::new(); - map.insert(H256::zero(), U256::one()); + let mut map: HashMap> = HashMap::new(); + map.insert(H256::zero(), Some(U256::one())); map })]), ) @@ -737,8 +790,8 @@ mod tests { root1, HashMap::from([(account_hash, Some(account_state2.clone()))]), HashMap::from([(account_hash, { - let mut map: HashMap = HashMap::new(); - map.insert(H256::zero(), U256::zero()); + let mut map: HashMap> = HashMap::new(); + map.insert(H256::zero(), Some(U256::zero())); map })]), ) @@ -754,7 +807,7 @@ mod tests { let value = tree .get_storage_at_hash(root2, address, H256::zero()) .unwrap(); - assert_eq!(value, Some(Some(U256::zero()))); + assert_eq!(value, Some(U256::zero())); // Retrieve it from the first hash and check it returns the first value let retrieved_account = tree.get_account_state(root1, address).unwrap(); @@ -763,13 +816,13 @@ mod tests { let value = tree .get_storage_at_hash(root1, address, H256::zero()) .unwrap(); - assert_eq!(value, Some(Some(U256::one()))); + assert_eq!(value, Some(U256::one())); } - #[test] - fn test_override_account_storage_in_second_difflayer() { + #[tokio::test] + async fn test_override_account_storage_in_second_difflayer() { let tree = create_mock_tree(); - tree.rebuild(H256::zero(), H256::zero()).unwrap(); + tree.rebuild(H256::zero(), H256::zero(), false).unwrap(); let root1 = H256::from_low_u64_be(1); let root2 = H256::from_low_u64_be(2); @@ -800,8 +853,8 @@ mod tests { H256::zero(), HashMap::from([(account_hash, Some(account_state1.clone()))]), HashMap::from([(account_hash, { - let mut map: HashMap = HashMap::new(); - map.insert(H256::zero(), U256::one()); + let mut map: HashMap> = HashMap::new(); + map.insert(H256::zero(), Some(U256::one())); map })]), ) @@ -814,8 +867,8 @@ mod tests { root1, HashMap::from([(account_hash, Some(account_state2.clone()))]), HashMap::from([(account_hash, { - let mut map: HashMap = HashMap::new(); - map.insert(H256::zero(), U256::zero()); + let mut map: HashMap> = HashMap::new(); + map.insert(H256::zero(), Some(U256::zero())); map })]), ) @@ -828,7 +881,7 @@ mod tests { let value = tree .get_storage_at_hash(root2, address, H256::zero()) .unwrap(); - assert_eq!(value, Some(Some(U256::zero()))); + assert_eq!(value, Some(U256::zero())); // Retrieve it from the first hash and check it returns the first value let retrieved_account = tree.get_account_state(root1, address).unwrap(); @@ -837,6 +890,6 @@ mod tests { let value = tree .get_storage_at_hash(root1, address, H256::zero()) .unwrap(); - assert_eq!(value, Some(Some(U256::one()))); + assert_eq!(value, Some(U256::one())); } } diff --git a/crates/storage/store.rs b/crates/storage/store.rs index fc4ff6ccd41..4f10c4803d7 100644 --- a/crates/storage/store.rs +++ b/crates/storage/store.rs @@ -112,7 +112,9 @@ impl Store { nonce: account_state.nonce, })) } - Ok(None) => {} + Ok(None) => { + return Ok(None); + } Err(snapshot_error) => { debug!("failed to fetch snapshot (state): {}", snapshot_error); } @@ -393,8 +395,18 @@ impl Store { genesis_accounts: BTreeMap, ) -> Result { let mut genesis_state_trie = self.engine.open_state_trie(*EMPTY_TRIE_HASH); + + // For snapshots + let mut account_hashes = Vec::with_capacity(1024); + let mut account_states = Vec::with_capacity(1024); + let mut storage_keys: Vec> = Vec::with_capacity(64); + let mut storage_values: Vec> = Vec::with_capacity(64); + for (address, account) in genesis_accounts { + let mut keys = Vec::with_capacity(32); + let mut values = Vec::with_capacity(32); let hashed_address = hash_address(&address); + let hashed_address_fixed = hash_address_fixed(&address); // Store account code (as this won't be stored in the trie) let code_hash = code_hash(&account.code); self.add_account_code(code_hash, account.code).await?; @@ -406,6 +418,9 @@ impl Store { if !storage_value.is_zero() { let hashed_key = hash_key(&H256(storage_key.to_big_endian())); storage_trie.insert(hashed_key, storage_value.encode_to_vec())?; + + keys.push(H256(storage_key.to_big_endian())); + values.push(storage_value); } } let storage_root = storage_trie.hash()?; @@ -416,8 +431,17 @@ impl Store { storage_root, code_hash, }; + account_hashes.push(hashed_address_fixed); + account_states.push(account_state.clone()); + storage_keys.push(keys); + storage_values.push(values); genesis_state_trie.insert(hashed_address, account_state.encode_to_vec())?; } + + // Add the initial genesis data to the snapshot db. + self.snapshots + .add_data(account_hashes, account_states, storage_keys, storage_values); + genesis_state_trie.hash().map_err(StoreError::Trie) } @@ -501,6 +525,10 @@ impl Store { self.set_canonical_block(genesis_block_number, genesis_hash) .await?; + self.snapshots + .rebuild(genesis_hash, genesis_state_root, false) + .unwrap(); + // Set chain config self.set_chain_config(&genesis.config).await } @@ -549,9 +577,7 @@ impl Store { .snapshots .get_storage_at_hash(block_hash, address, storage_key) { - Ok(Some(Some(value))) => return Ok(Some(value)), - // It may be none, but until the disk layer is fixed we can't trust it. - Ok(_) => {} + Ok(value) => return Ok(value), // snapshot errors are non-fatal Err(snapshot_error) => { debug!("failed to fetch snapshot (storage): {}", snapshot_error); @@ -711,8 +737,7 @@ impl Store { #[cfg(feature = "snapshots")] match self.snapshots.get_account_state(block_hash, address) { - Ok(Some(value)) => return Ok(Some(value)), - Ok(None) => {} + Ok(value) => return Ok(value), Err(snapshot_error) => { debug!("failed to fetch snapshot (state): {}", snapshot_error); } @@ -729,6 +754,14 @@ impl Store { block_hash: BlockHash, address: Address, ) -> Result, StoreError> { + #[cfg(feature = "snapshots")] + match self.snapshots.get_account_state(block_hash, address) { + Ok(value) => return Ok(value), + Err(snapshot_error) => { + debug!("failed to fetch snapshot (state): {}", snapshot_error); + } + } + let Some(state_trie) = self.state_trie(block_hash)? else { return Ok(None); }; @@ -1155,7 +1188,7 @@ impl Store { if store.snapshots.len() == 0 { // There are no snapshots yet, use this block as root // TODO: find if there is a better place to create the initial "disk layer". - store.snapshots.rebuild(hash, state_root)?; + store.snapshots.rebuild(hash, state_root, true)?; info!( "Snapshot (disk layer) created for {} with parent {}", hash, parent_hash @@ -1165,14 +1198,12 @@ impl Store { let mut accounts = HashMap::new(); let state_trie = store.open_state_trie(state_root); - let mut storage: HashMap> = HashMap::new(); + let mut storage: HashMap>> = HashMap::new(); for update in account_updates.iter() { let hashed_address = hash_address_fixed(&update.address); - if update.removed { - accounts.insert(hashed_address, None); - } else { + if !update.removed { let account_state = match state_trie.get(hashed_address).unwrap() { Some(encoded_state) => AccountState::decode(&encoded_state).unwrap(), None => AccountState::default(), @@ -1181,8 +1212,14 @@ impl Store { for (storage_key, storage_value) in &update.added_storage { let slots = storage.entry(hashed_address).or_default(); - slots.insert(*storage_key, *storage_value); + if !storage_value.is_zero() { + slots.insert(*storage_key, Some(*storage_value)); + } else { + slots.insert(*storage_key, None); + } } + } else { + accounts.insert(hashed_address, None); } } diff --git a/crates/storage/store_db/in_memory.rs b/crates/storage/store_db/in_memory.rs index bab541ab0e8..a8deb31449c 100644 --- a/crates/storage/store_db/in_memory.rs +++ b/crates/storage/store_db/in_memory.rs @@ -47,6 +47,7 @@ struct StoreInner { state_snapshot: BTreeMap, // Stores Storage trie leafs from the last downloaded tries storage_snapshot: HashMap>, + // Snapshot (diff layers) } #[derive(Default, Debug)] @@ -565,6 +566,15 @@ impl StoreEngine for Store { &self, account_hashes: Vec, account_states: Vec, + ) -> Result<(), StoreError> { + self.write_snapshot_account_batch_blocking(account_hashes, account_states) + } + + #[inline] + fn write_snapshot_account_batch_blocking( + &self, + account_hashes: Vec, + account_states: Vec, ) -> Result<(), StoreError> { self.inner() .state_snapshot @@ -590,6 +600,15 @@ impl StoreEngine for Store { account_hashes: Vec, storage_keys: Vec>, storage_values: Vec>, + ) -> Result<(), StoreError> { + self.write_snapshot_storage_batches_blocking(account_hashes, storage_keys, storage_values) + } + + fn write_snapshot_storage_batches_blocking( + &self, + account_hashes: Vec, + storage_keys: Vec>, + storage_values: Vec>, ) -> Result<(), StoreError> { for (account_hash, (storage_keys, storage_values)) in account_hashes .into_iter() @@ -686,6 +705,23 @@ impl StoreEngine for Store { .insert(bad_block, latest_valid); Ok(()) } + + fn get_account_snapshot(&self, account_hash: H256) -> Result, StoreError> { + Ok(self.inner().state_snapshot.get(&account_hash).cloned()) + } + + fn get_storage_snapshot( + &self, + account_hash: H256, + storage_hash: H256, + ) -> Result, StoreError> { + Ok(self + .inner() + .storage_snapshot + .get(&account_hash) + .and_then(|x| x.get(&storage_hash)) + .cloned()) + } } impl Debug for Store { diff --git a/crates/storage/store_db/libmdbx.rs b/crates/storage/store_db/libmdbx.rs index 44c456abe74..04e4cda0b7a 100644 --- a/crates/storage/store_db/libmdbx.rs +++ b/crates/storage/store_db/libmdbx.rs @@ -33,6 +33,7 @@ use std::fmt::{Debug, Formatter}; use std::path::Path; use std::sync::Arc; +#[derive(Clone)] pub struct Store { db: Arc, } @@ -57,24 +58,34 @@ impl Store { } // Helper method to write into a libmdbx table in batch - async fn write_batch( + #[inline] + fn write_batch_sync( &self, key_values: Vec<(T::Key, T::Value)>, ) -> Result<(), StoreError> { - let db = self.db.clone(); - tokio::task::spawn_blocking(move || { - let txn = db.begin_readwrite().map_err(StoreError::LibmdbxError)?; + let txn = self + .db + .begin_readwrite() + .map_err(StoreError::LibmdbxError)?; - let mut cursor = txn.cursor::().map_err(StoreError::LibmdbxError)?; - for (key, value) in key_values { - cursor - .upsert(key, value) - .map_err(StoreError::LibmdbxError)?; - } - txn.commit().map_err(StoreError::LibmdbxError) - }) - .await - .map_err(|e| StoreError::Custom(format!("task panicked: {e}")))? + let mut cursor = txn.cursor::().map_err(StoreError::LibmdbxError)?; + for (key, value) in key_values { + cursor + .upsert(key, value) + .map_err(StoreError::LibmdbxError)?; + } + txn.commit().map_err(StoreError::LibmdbxError) + } + + // Helper method to write into a libmdbx table in batch, async version. + async fn write_batch( + &self, + key_values: Vec<(T::Key, T::Value)>, + ) -> Result<(), StoreError> { + let db = (*self).clone(); + tokio::task::spawn_blocking(move || db.write_batch_sync::(key_values)) + .await + .map_err(|e| StoreError::Custom(format!("task panicked: {e}")))? } // Helper method to read from a libmdbx table @@ -803,6 +814,20 @@ impl StoreEngine for Store { .await } + fn write_snapshot_account_batch_blocking( + &self, + account_hashes: Vec, + account_states: Vec, + ) -> Result<(), StoreError> { + self.write_batch_sync::( + account_hashes + .into_iter() + .map(|h| h.into()) + .zip(account_states.into_iter().map(|a| a.into())) + .collect(), + ) + } + async fn write_snapshot_storage_batch( &self, account_hash: H256, @@ -830,24 +855,40 @@ impl StoreEngine for Store { storage_keys: Vec>, storage_values: Vec>, ) -> Result<(), StoreError> { - let db = self.db.clone(); + let store = self.clone(); tokio::task::spawn_blocking(move || { - let txn = db.begin_readwrite().map_err(StoreError::LibmdbxError)?; - for (account_hash, (storage_keys, storage_values)) in account_hashes - .into_iter() - .zip(storage_keys.into_iter().zip(storage_values.into_iter())) - { - for (key, value) in storage_keys.into_iter().zip(storage_values.into_iter()) { - txn.upsert::(account_hash.into(), (key.into(), value.into())) - .map_err(StoreError::LibmdbxError)?; - } - } - txn.commit().map_err(StoreError::LibmdbxError) + store.write_snapshot_storage_batches_blocking( + account_hashes, + storage_keys, + storage_values, + ) }) .await .map_err(|e| StoreError::Custom(format!("task panicked: {e}")))? } + fn write_snapshot_storage_batches_blocking( + &self, + account_hashes: Vec, + storage_keys: Vec>, + storage_values: Vec>, + ) -> Result<(), StoreError> { + let txn = self + .db + .begin_readwrite() + .map_err(StoreError::LibmdbxError)?; + for (account_hash, (storage_keys, storage_values)) in account_hashes + .into_iter() + .zip(storage_keys.into_iter().zip(storage_values.into_iter())) + { + for (key, value) in storage_keys.into_iter().zip(storage_values.into_iter()) { + txn.upsert::(account_hash.into(), (key.into(), value.into())) + .map_err(StoreError::LibmdbxError)?; + } + } + txn.commit().map_err(StoreError::LibmdbxError) + } + async fn set_state_trie_rebuild_checkpoint( &self, checkpoint: (H256, [H256; STATE_TRIE_SEGMENTS]), @@ -962,6 +1003,28 @@ impl StoreEngine for Store { self.write::(bad_block.into(), latest_valid.into()) .await } + + fn get_account_snapshot(&self, account_hash: H256) -> Result, StoreError> { + let txn = self.db.begin_read().map_err(StoreError::LibmdbxError)?; + txn.get::(account_hash.into()) + .map_err(StoreError::LibmdbxError) + .map(|x| x.map(|y| y.to())) + } + + fn get_storage_snapshot( + &self, + account_hash: H256, + storage_hash: H256, + ) -> Result, StoreError> { + let txn = self.db.begin_read().map_err(StoreError::LibmdbxError)?; + let mut cursor = txn + .cursor::() + .map_err(StoreError::LibmdbxError)?; + Ok(cursor + .seek_value(account_hash.into(), storage_hash.into()) + .map_err(StoreError::LibmdbxError)? + .map(|x| U256::from_big_endian(&x.1 .0))) + } } impl Debug for Store { @@ -1154,12 +1217,12 @@ table!( ); table!( - /// State Snapshot used by an ongoing sync process + /// State Snapshot used by sync and disk layer snapshots ( StateSnapShot ) AccountHashRLP => AccountStateRLP ); dupsort!( - /// Storage Snapshot used by an ongoing sync process + /// Storage Snapshot used by sync and disk layer snapshots ( StorageSnapShot ) AccountHashRLP => (AccountStorageKeyBytes, AccountStorageValueBytes)[AccountStorageKeyBytes] ); diff --git a/crates/storage/store_db/redb.rs b/crates/storage/store_db/redb.rs index 49e38a5f009..62642ea8623 100644 --- a/crates/storage/store_db/redb.rs +++ b/crates/storage/store_db/redb.rs @@ -66,7 +66,7 @@ const STORAGE_SNAPSHOT_TABLE: MultimapTableDefinition = TableDefinition::new("StorageHealPaths"); -#[derive(Debug)] +#[derive(Debug, Clone)] pub struct RedBStore { db: Arc, } @@ -152,21 +152,39 @@ impl RedBStore { 'k: 'static, 'v: 'static, { - let db = self.db.clone(); - tokio::task::spawn_blocking(move || { - let write_txn = db.begin_write()?; - { - let mut table = write_txn.open_table(table)?; - for (key, value) in key_values { - table.insert(key, value)?; - } + let store = self.clone(); + tokio::task::spawn_blocking(move || store.write_batch_sync(table, key_values)) + .await + .map_err(|e| StoreError::Custom(format!("task panicked: {e}")))? + } + + // Helper method to write into a redb table. Sync version + #[inline] + fn write_batch_sync<'k, 'v, 'a, K, V>( + &self, + table: TableDefinition<'a, K, V>, + key_values: Vec<(K::SelfType<'k>, V::SelfType<'v>)>, + ) -> Result<(), StoreError> + where + K: Key + Send + 'static, + V: Value + Send + 'static, + K::SelfType<'k>: Send, + V::SelfType<'v>: Send, + TableDefinition<'a, K, V>: Send, + 'a: 'static, + 'k: 'static, + 'v: 'static, + { + let write_txn = self.db.begin_write()?; + { + let mut table = write_txn.open_table(table)?; + for (key, value) in key_values { + table.insert(key, value)?; } - write_txn.commit()?; + } + write_txn.commit()?; - Ok(()) - }) - .await - .map_err(|e| StoreError::Custom(format!("task panicked: {e}")))? + Ok(()) } // Helper method to write into a redb table @@ -1063,6 +1081,25 @@ impl StoreEngine for RedBStore { .await } + fn write_snapshot_account_batch_blocking( + &self, + account_hashes: Vec, + account_states: Vec, + ) -> Result<(), StoreError> { + self.write_batch_sync( + STATE_SNAPSHOT_TABLE, + account_hashes + .into_iter() + .map(>::into) + .zip( + account_states + .into_iter() + .map(>::into), + ) + .collect::>(), + ) + } + async fn write_snapshot_storage_batch( &self, account_hash: H256, @@ -1087,6 +1124,15 @@ impl StoreEngine for RedBStore { account_hashes: Vec, storage_keys: Vec>, storage_values: Vec>, + ) -> Result<(), StoreError> { + self.write_snapshot_storage_batches_blocking(account_hashes, storage_keys, storage_values) + } + + fn write_snapshot_storage_batches_blocking( + &self, + account_hashes: Vec, + storage_keys: Vec>, + storage_values: Vec>, ) -> Result<(), StoreError> { let write_tx = self.db.begin_write()?; { @@ -1232,6 +1278,33 @@ impl StoreEngine for RedBStore { ) .await } + + fn get_account_snapshot(&self, account_hash: H256) -> Result, StoreError> { + let read_tx = self.db.begin_read()?; + let table = read_tx.open_table(STATE_SNAPSHOT_TABLE)?; + Ok(table + .get(&account_hash.into())? + .map(|elem| elem.value().to())) + } + + fn get_storage_snapshot( + &self, + account_hash: H256, + storage_hash: H256, + ) -> Result, StoreError> { + let read_tx = self.db.begin_read()?; + let table = read_tx.open_multimap_table(STORAGE_SNAPSHOT_TABLE)?; + Ok(table.get(&account_hash.into())?.find_map(|elem| { + elem.ok().and_then(|x| { + let val = x.value(); + if H256(val.0) == storage_hash { + Some(U256::from_big_endian(&val.1)) + } else { + None + } + }) + })) + } } impl redb::Value for ChainDataIndex {