diff --git a/dash-spv/src/client/lifecycle.rs b/dash-spv/src/client/lifecycle.rs index 46e26f71c..37904d9c3 100644 --- a/dash-spv/src/client/lifecycle.rs +++ b/dash-spv/src/client/lifecycle.rs @@ -13,8 +13,9 @@ use crate::chain::checkpoints::CheckpointManager; use crate::error::{Result, SpvError}; use crate::network::NetworkManager; use crate::storage::{ - PersistentBlockHeaderStorage, PersistentBlockStorage, PersistentFilterHeaderStorage, - PersistentFilterStorage, PersistentMetadataStorage, StorageManager, + MasternodeStorage, PersistentBlockHeaderStorage, PersistentBlockStorage, + PersistentFilterHeaderStorage, PersistentFilterStorage, PersistentMetadataStorage, + StorageManager, }; use crate::sync::{ BlockHeadersManager, BlocksManager, ChainLockManager, FilterHeadersManager, FiltersManager, @@ -67,9 +68,13 @@ impl DashSpvClient DashSpvClient DashSpvClient; + +struct CachedList { + from: CoreBlockHeight, + until: Option, + list: Option, +} #[async_trait] -pub trait MasternodeStateStorage { - async fn store_masternode_state(&mut self, state: &MasternodeState) -> StorageResult<()>; +pub trait MasternodeStorage: Send + Sync + 'static { + async fn store_diff(&mut self, height: CoreBlockHeight, diff: &MnListDiff) + -> StorageResult<()>; + + async fn store_qr_info( + &mut self, + height: CoreBlockHeight, + qr_info: &QRInfo, + ) -> StorageResult<()>; - async fn load_masternode_state(&self) -> StorageResult>; + async fn load_engine(&self, network: Network) -> StorageResult; + + async fn masternode_list_at_or_before( + &self, + network: Network, + height: CoreBlockHeight, + ) -> StorageResult>; } -pub struct PersistentMasternodeStateStorage { +pub struct PersistentMasternodeStorage { storage_path: PathBuf, + headers: Arc>, + diffs: IndexMap, + qr_infos: IndexMap, + cached_list: Mutex>, } -impl PersistentMasternodeStateStorage { - const FOLDER_NAME: &str = "masternodestate"; - const MASTERNODE_FILE_NAME: &str = "masternodestate.json"; -} +impl PersistentMasternodeStorage { + const FOLDER_NAME: &str = "masternodes"; + const DIFF_PREFIX: &str = "diff_"; + const QRINFO_PREFIX: &str = "qrinfo_"; + const EXTENSION: &str = "dat"; -#[async_trait] -impl PersistentStorage for PersistentMasternodeStateStorage { - async fn open(storage_path: impl Into + Send) -> StorageResult { - Ok(PersistentMasternodeStateStorage { - storage_path: storage_path.into(), + pub async fn open( + storage_path: impl Into + Send, + headers: Arc>, + ) -> StorageResult { + let storage_path = storage_path.into(); + let (diffs, qr_infos) = Self::index_folder(&storage_path.join(Self::FOLDER_NAME)).await?; + + Ok(PersistentMasternodeStorage { + storage_path, + headers, + diffs, + qr_infos, + cached_list: Mutex::new(None), }) } - async fn persist(&mut self, _storage_path: impl Into + Send) -> StorageResult<()> { - // Current implementation persists data everytime data is stored + fn folder(&self) -> PathBuf { + self.storage_path.join(Self::FOLDER_NAME) + } + + fn file_name(prefix: &str, height: CoreBlockHeight) -> String { + format!("{prefix}{height}.{}", Self::EXTENSION) + } + + fn height_from_file_name(name: &str, prefix: &str) -> Option { + name.strip_prefix(prefix)?.strip_suffix(&format!(".{}", Self::EXTENSION))?.parse().ok() + } + + async fn index_folder(folder: &Path) -> StorageResult<(IndexMap, IndexMap)> { + let mut diffs = BTreeMap::new(); + let mut qr_infos = BTreeMap::new(); + + if !folder.exists() { + return Ok((diffs, qr_infos)); + } + + let mut entries = tokio::fs::read_dir(folder).await?; + while let Some(entry) = entries.next_entry().await? { + let path = entry.path(); + let Some(name) = path.file_name().and_then(|n| n.to_str()) else { + continue; + }; + if let Some(height) = Self::height_from_file_name(name, Self::DIFF_PREFIX) { + diffs.insert(height, path); + } else if let Some(height) = Self::height_from_file_name(name, Self::QRINFO_PREFIX) { + qr_infos.insert(height, path); + } + } + + Ok((diffs, qr_infos)) + } + + async fn store_message( + folder: &Path, + index: &mut IndexMap, + prefix: &str, + height: CoreBlockHeight, + message: &T, + ) -> StorageResult<()> { + if index.contains_key(&height) { + return Ok(()); + } + tokio::fs::create_dir_all(folder).await?; + let path = folder.join(Self::file_name(prefix, height)); + atomic_write(&path, &serialize(message)).await?; + index.insert(height, path); Ok(()) } + + async fn read_message(path: &Path) -> StorageResult { + let bytes = tokio::fs::read(path).await?; + deserialize(&bytes).map_err(|e| { + StorageError::Corruption(format!("Failed to decode {}: {e}", path.display())) + }) + } + + async fn load_pending( + index: &IndexMap, + label: &str, + ) -> Vec<(CoreBlockHeight, T)> { + let mut pending = Vec::new(); + for (height, path) in index { + match Self::read_message::(path).await { + Ok(message) => pending.push((*height, message)), + Err(e) => tracing::warn!("Skipping unreadable {label} at {height}: {e}"), + } + } + pending + } + + async fn cached_list_at(&self, height: CoreBlockHeight) -> Option> { + let cached = self.cached_list.lock().await; + let cached = cached.as_ref()?; + (height >= cached.from && cached.until.is_none_or(|until| height < until)) + .then(|| cached.list.clone()) + } + + fn invalidate_cached_list(&mut self) { + *self.cached_list.get_mut() = None; + } + + async fn replay(&self, network: Network) -> StorageResult { + let mut engine = MasternodeListEngine::default_for_network(network); + + let mut pending_qr_infos: Vec<(CoreBlockHeight, QRInfo)> = + Self::load_pending(&self.qr_infos, "QRInfo").await; + let mut pending_diffs: Vec<(CoreBlockHeight, MnListDiff)> = + Self::load_pending(&self.diffs, "MnListDiff").await; + + { + let headers = self.headers.read().await; + for (_, qr_info) in &pending_qr_infos { + feed_qrinfo_heights_to_engine(&mut engine, qr_info, &*headers).await; + } + for (height, diff) in &pending_diffs { + engine.feed_block_height(*height, diff.block_hash); + if let Ok(Some(base_height)) = + headers.get_header_height_by_hash(&diff.base_block_hash).await + { + engine.feed_block_height(base_height, diff.base_block_hash); + } + } + } + + let qr_info_count = pending_qr_infos.len(); + let diff_count = pending_diffs.len(); + + loop { + let remaining = pending_qr_infos.len() + pending_diffs.len(); + + pending_qr_infos + .retain(|(_, qr_info)| engine.feed_qr_info(qr_info.clone(), true, true).is_err()); + + pending_diffs.retain(|(height, diff)| { + engine.apply_diff(diff.clone(), Some(*height), false, None).is_err() + }); + + if pending_qr_infos.len() + pending_diffs.len() == remaining { + break; + } + } + + for (height, _) in &pending_qr_infos { + tracing::warn!("QRInfo at {height} has no reachable base, leaving it to the network"); + } + + for (height, _) in &pending_diffs { + tracing::warn!( + "MnListDiff at {height} has no reachable base, leaving it to the network" + ); + } + + tracing::debug!( + "Replayed {}/{} QRInfo and {}/{} MnListDiff messages into {} masternode lists", + qr_info_count - pending_qr_infos.len(), + qr_info_count, + diff_count - pending_diffs.len(), + diff_count, + engine.masternode_lists.len() + ); + + Ok(engine) + } } #[async_trait] -impl MasternodeStateStorage for PersistentMasternodeStateStorage { - async fn store_masternode_state(&mut self, state: &MasternodeState) -> StorageResult<()> { - let masternodestate_folder = self.storage_path.join(Self::FOLDER_NAME); - let path = masternodestate_folder.join(Self::MASTERNODE_FILE_NAME); - - tokio::fs::create_dir_all(masternodestate_folder).await?; +impl MasternodeStorage for PersistentMasternodeStorage { + async fn store_diff( + &mut self, + height: CoreBlockHeight, + diff: &MnListDiff, + ) -> StorageResult<()> { + let folder = self.folder(); + self.invalidate_cached_list(); + Self::store_message(&folder, &mut self.diffs, Self::DIFF_PREFIX, height, diff).await + } - let json = serde_json::to_string_pretty(state).map_err(|e| { - crate::error::StorageError::Serialization(format!( - "Failed to serialize masternode state: {}", - e - )) - })?; + async fn store_qr_info( + &mut self, + height: CoreBlockHeight, + qr_info: &QRInfo, + ) -> StorageResult<()> { + let folder = self.folder(); + self.invalidate_cached_list(); + Self::store_message(&folder, &mut self.qr_infos, Self::QRINFO_PREFIX, height, qr_info).await + } - atomic_write(&path, json.as_bytes()).await?; - Ok(()) + async fn load_engine(&self, network: Network) -> StorageResult { + self.replay(network).await } - async fn load_masternode_state(&self) -> StorageResult> { - let path = self.storage_path.join(Self::FOLDER_NAME).join(Self::MASTERNODE_FILE_NAME); + async fn masternode_list_at_or_before( + &self, + network: Network, + height: CoreBlockHeight, + ) -> StorageResult> { + if let Some(hit) = self.cached_list_at(height).await { + return Ok(hit); + } + + let engine = self.replay(network).await?; + let (before, after) = engine.masternode_lists_around_height(height); + let list = before.cloned(); + + *self.cached_list.lock().await = Some(CachedList { + from: before.map_or(0, |list| list.known_height), + until: after.map(|next| next.known_height), + list: list.clone(), + }); - if !path.exists() { - return Ok(None); + Ok(list) + } +} + +/// Feed QRInfo block heights to the engine from the header storage. +/// +/// Resolves heights for every hash enumerated by +/// [`MasternodeListEngine::qr_info_referenced_block_hashes`], plus the cycle boundary +/// block for each work-block diff (`work_height + WORK_DIFF_DEPTH`), which is needed +/// for rotated quorum storage key calculation. +pub(crate) async fn feed_qrinfo_heights_to_engine( + engine: &mut MasternodeListEngine, + qr_info: &QRInfo, + storage: &S, +) { + let mut fed_count = 0; + for block_hash in MasternodeListEngine::qr_info_referenced_block_hashes(qr_info) { + if let Ok(Some(height)) = storage.get_header_height_by_hash(&block_hash).await { + engine.feed_block_height(height, block_hash); + fed_count += 1; + tracing::trace!("Fed height {} for block {}", height, block_hash); } + } - let content = tokio::fs::read_to_string(path).await?; - let state = serde_json::from_str(&content).map_err(|e| { - crate::error::StorageError::Serialization(format!( - "Failed to deserialize masternode state: {}", - e - )) - })?; + // Feed cycle boundary heights for all diffs (current and historical cycles). + // Each diff's block_hash is at the "work block" height; the cycle boundary is + // WORK_DIFF_DEPTH higher. + let mut work_block_hashes = vec![ + qr_info.mn_list_diff_h.block_hash, + qr_info.mn_list_diff_at_h_minus_c.block_hash, + qr_info.mn_list_diff_at_h_minus_2c.block_hash, + qr_info.mn_list_diff_at_h_minus_3c.block_hash, + ]; - Ok(Some(state)) + if let Some((_, diff)) = &qr_info.quorum_snapshot_and_mn_list_diff_at_h_minus_4c { + work_block_hashes.push(diff.block_hash); } + + for work_block_hash in work_block_hashes { + if let Ok(Some(work_block_height)) = + storage.get_header_height_by_hash(&work_block_hash).await + { + let cycle_boundary_height = work_block_height + WORK_DIFF_DEPTH; + if let Ok(Some(cycle_boundary_header)) = storage.get_header(cycle_boundary_height).await + { + let cycle_boundary_hash = *cycle_boundary_header.hash(); + engine.feed_block_height(cycle_boundary_height, cycle_boundary_hash); + fed_count += 1; + tracing::debug!( + "Fed cycle boundary height {} for block {}", + cycle_boundary_height, + cycle_boundary_hash + ); + } + } + } + + tracing::info!("Fed {} block heights to engine", fed_count); } diff --git a/dash-spv/src/storage/mod.rs b/dash-spv/src/storage/mod.rs index 70a851acb..8a76c463c 100644 --- a/dash-spv/src/storage/mod.rs +++ b/dash-spv/src/storage/mod.rs @@ -1,7 +1,5 @@ //! Storage abstraction for the Dash SPV client. -pub mod types; - mod block_headers; mod blocks; mod filter_headers; @@ -18,7 +16,12 @@ use crate::types::{HashedBlock, HashedBlockHeader}; use crate::ClientConfig; use async_trait::async_trait; use dashcore::hash_types::FilterHeader; +use dashcore::network::message_qrinfo::QRInfo; +use dashcore::network::message_sml::MnListDiff; use dashcore::prelude::CoreBlockHeight; +use dashcore::sml::masternode_list::MasternodeList; +use dashcore::sml::masternode_list_engine::MasternodeListEngine; +use dashcore::Network; use std::ops::Range; use std::path::{Path, PathBuf}; use std::sync::Arc; @@ -31,12 +34,11 @@ pub use crate::storage::block_headers::{ pub use crate::storage::blocks::{BlockStorage, PersistentBlockStorage}; pub use crate::storage::filter_headers::{FilterHeaderStorage, PersistentFilterHeaderStorage}; pub use crate::storage::filters::{FilterStorage, PersistentFilterStorage}; -pub use crate::storage::masternode::{MasternodeStateStorage, PersistentMasternodeStateStorage}; +pub(crate) use crate::storage::masternode::feed_qrinfo_heights_to_engine; +pub use crate::storage::masternode::{MasternodeStorage, PersistentMasternodeStorage}; pub use crate::storage::metadata::{MetadataStorage, PersistentMetadataStorage}; pub use crate::storage::peers::{PeerStorage, PersistentPeerStorage}; -pub use types::*; - #[async_trait] pub trait PersistentStorage: Sized { /// If the storage_path contains persisted data the storage will use it, if not, @@ -53,7 +55,7 @@ pub trait StorageManager: + FilterStorage + BlockStorage + MetadataStorage - + MasternodeStateStorage + + MasternodeStorage + Send + Sync + 'static @@ -78,6 +80,9 @@ pub trait StorageManager: /// Returns shared access to the metadata storage. fn metadata(&self) -> Arc>; + + fn masternodes(&self) + -> Arc>>; } /// Disk-based storage manager with segmented files and async background saving. @@ -91,7 +96,7 @@ pub struct DiskStorageManager { filters: Arc>, blocks: Arc>, metadata: Arc>, - masternodestate: Arc>, + masternodes: Arc>>, // Background worker worker_handle: Option>, @@ -132,21 +137,23 @@ impl DiskStorageManager { let lock_file = LockFile::new(lock_file)?; + let block_headers = + Arc::new(RwLock::new(PersistentBlockHeaderStorage::open(&storage_path).await?)); + let mut storage = Self { storage_path: storage_path.clone(), - block_headers: Arc::new(RwLock::new( - PersistentBlockHeaderStorage::open(&storage_path).await?, - )), filter_headers: Arc::new(RwLock::new( PersistentFilterHeaderStorage::open(&storage_path).await?, )), filters: Arc::new(RwLock::new(PersistentFilterStorage::open(&storage_path).await?)), blocks: Arc::new(RwLock::new(PersistentBlockStorage::open(&storage_path).await?)), metadata: Arc::new(RwLock::new(PersistentMetadataStorage::open(&storage_path).await?)), - masternodestate: Arc::new(RwLock::new( - PersistentMasternodeStateStorage::open(&storage_path).await?, + masternodes: Arc::new(RwLock::new( + PersistentMasternodeStorage::open(&storage_path, Arc::clone(&block_headers)) + .await?, )), + block_headers, worker_handle: None, @@ -173,7 +180,6 @@ impl DiskStorageManager { let filters = Arc::clone(&self.filters); let blocks = Arc::clone(&self.blocks); let metadata = Arc::clone(&self.metadata); - let masternodestate = Arc::clone(&self.masternodestate); let storage_path = self.storage_path.clone(); @@ -188,7 +194,6 @@ impl DiskStorageManager { let _ = filters.write().await.persist(&storage_path).await; let _ = blocks.write().await.persist(&storage_path).await; let _ = metadata.write().await.persist(&storage_path).await; - let _ = masternodestate.write().await.persist(&storage_path).await; } }); @@ -210,7 +215,6 @@ impl DiskStorageManager { let _ = self.filters.write().await.persist(storage_path).await; let _ = self.blocks.write().await.persist(storage_path).await; let _ = self.metadata.write().await.persist(storage_path).await; - let _ = self.masternodestate.write().await.persist(storage_path).await; } } @@ -247,8 +251,10 @@ impl StorageManager for DiskStorageManager { self.filters = Arc::new(RwLock::new(PersistentFilterStorage::open(storage_path).await?)); self.blocks = Arc::new(RwLock::new(PersistentBlockStorage::open(storage_path).await?)); self.metadata = Arc::new(RwLock::new(PersistentMetadataStorage::open(storage_path).await?)); - self.masternodestate = - Arc::new(RwLock::new(PersistentMasternodeStateStorage::open(storage_path).await?)); + self.masternodes = Arc::new(RwLock::new( + PersistentMasternodeStorage::open(storage_path, Arc::clone(&self.block_headers)) + .await?, + )); // Restart the background worker for future operations self.start_worker().await; @@ -282,6 +288,12 @@ impl StorageManager for DiskStorageManager { fn metadata(&self) -> Arc> { Arc::clone(&self.metadata) } + + fn masternodes( + &self, + ) -> Arc>> { + Arc::clone(&self.masternodes) + } } #[async_trait] @@ -431,13 +443,33 @@ impl metadata::MetadataStorage for DiskStorageManager { } #[async_trait] -impl masternode::MasternodeStateStorage for DiskStorageManager { - async fn store_masternode_state(&mut self, state: &MasternodeState) -> StorageResult<()> { - self.masternodestate.write().await.store_masternode_state(state).await +impl masternode::MasternodeStorage for DiskStorageManager { + async fn store_diff( + &mut self, + height: CoreBlockHeight, + diff: &MnListDiff, + ) -> StorageResult<()> { + self.masternodes.write().await.store_diff(height, diff).await + } + + async fn store_qr_info( + &mut self, + height: CoreBlockHeight, + qr_info: &QRInfo, + ) -> StorageResult<()> { + self.masternodes.write().await.store_qr_info(height, qr_info).await + } + + async fn load_engine(&self, network: Network) -> StorageResult { + self.masternodes.read().await.load_engine(network).await } - async fn load_masternode_state(&self) -> StorageResult> { - self.masternodestate.read().await.load_masternode_state().await + async fn masternode_list_at_or_before( + &self, + network: Network, + height: CoreBlockHeight, + ) -> StorageResult> { + self.masternodes.read().await.masternode_list_at_or_before(network, height).await } } diff --git a/dash-spv/src/storage/types.rs b/dash-spv/src/storage/types.rs deleted file mode 100644 index 678553caa..000000000 --- a/dash-spv/src/storage/types.rs +++ /dev/null @@ -1,16 +0,0 @@ -//! Storage-related types and structures. - -use serde::{Deserialize, Serialize}; - -/// Masternode state for storage. -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct MasternodeState { - /// Last processed height. - pub last_height: u32, - - /// Serialized masternode list engine state. - pub engine_state: Vec, - - /// Last update timestamp. - pub last_update: u64, -} diff --git a/dash-spv/src/sync/chainlock/manager.rs b/dash-spv/src/sync/chainlock/manager.rs index c211919df..0cd868c24 100644 --- a/dash-spv/src/sync/chainlock/manager.rs +++ b/dash-spv/src/sync/chainlock/manager.rs @@ -10,11 +10,14 @@ use std::sync::Arc; use dashcore::ephemerealdata::chain_lock::ChainLock; use dashcore::hash_types::ChainLockHash; use dashcore::sml::masternode_list_engine::MasternodeListEngine; +use dashcore::Network; use std::collections::HashSet; use tokio::sync::RwLock; use crate::error::SyncResult; -use crate::storage::{BlockHeaderStorage, MetadataStorage}; +use crate::storage::{ + BlockHeaderStorage, MasternodeStorage, MetadataStorage, PersistentMasternodeStorage, +}; use crate::sync::{ChainLockProgress, SyncEvent}; /// Metadata key for persisting the best validated ChainLock. @@ -36,6 +39,8 @@ pub struct ChainLockManager { metadata_storage: Arc>, /// Masternode engine for BLS signature validation. masternode_engine: Arc>, + masternode_storage: Option>>>, + network: Network, /// The best (highest height) validated ChainLock. best_chainlock: Option, /// ChainLock hashes that have been requested (to avoid duplicate requests). @@ -56,12 +61,16 @@ impl ChainLockManager { header_storage: Arc>, metadata_storage: Arc>, masternode_engine: Arc>, + masternode_storage: Option>>>, + network: Network, ) -> Self { let mut manager = Self { progress: ChainLockProgress::default(), header_storage, metadata_storage, masternode_engine, + masternode_storage, + network, best_chainlock: None, requested_chainlocks: HashSet::new(), masternode_ready: false, @@ -254,6 +263,55 @@ impl ChainLockManager { "ChainLock signature verified for height {}", chainlock.block_height ); + return true; + } + Err(e) => tracing::debug!( + "ChainLock at height {} not verifiable against the retained lists: {}", + chainlock.block_height, + e + ), + } + drop(engine); + + self.validate_signature_from_storage(chainlock).await + } + + async fn validate_signature_from_storage(&self, chainlock: &ChainLock) -> bool { + let Some(storage) = &self.masternode_storage else { + return false; + }; + + let signing_height = chainlock.block_height.saturating_sub(8); + let list = match storage + .read() + .await + .masternode_list_at_or_before(self.network, signing_height) + .await + { + Ok(Some(list)) => list, + Ok(None) => return false, + Err(e) => { + tracing::warn!( + "Could not rebuild the masternode list for height {}: {}", + signing_height, + e + ); + return false; + } + }; + + let engine = self.masternode_engine.read().await; + let Ok(request_id) = chainlock.request_id() else { + return false; + }; + + match engine.verify_chain_lock_with_masternode_list(chainlock, &list, &request_id) { + Ok(()) => { + tracing::info!( + "ChainLock signature verified for height {} from a rebuilt list at {}", + chainlock.block_height, + list.known_height + ); true } Err(e) => { @@ -309,7 +367,14 @@ mod tests { let storage = DiskStorageManager::with_temp_dir().await.unwrap(); let engine = Arc::new(RwLock::new(MasternodeListEngine::default_for_network(Network::Testnet))); - ChainLockManager::new(storage.block_headers(), storage.metadata(), engine).await + ChainLockManager::new( + storage.block_headers(), + storage.metadata(), + engine, + None, + Network::Testnet, + ) + .await } async fn create_test_manager_with_storage( @@ -317,7 +382,14 @@ mod tests { ) -> TestChainLockManager { let engine = Arc::new(RwLock::new(MasternodeListEngine::default_for_network(Network::Testnet))); - ChainLockManager::new(storage.block_headers(), storage.metadata(), engine).await + ChainLockManager::new( + storage.block_headers(), + storage.metadata(), + engine, + None, + Network::Testnet, + ) + .await } fn create_test_chainlock(height: u32) -> ChainLock { diff --git a/dash-spv/src/sync/masternodes/manager.rs b/dash-spv/src/sync/masternodes/manager.rs index 428673535..234cf5ca7 100644 --- a/dash-spv/src/sync/masternodes/manager.rs +++ b/dash-spv/src/sync/masternodes/manager.rs @@ -14,9 +14,10 @@ use tokio::sync::RwLock; use super::pipeline::MnListDiffPipeline; use crate::error::{SyncError, SyncResult}; use crate::network::RequestSender; -use crate::storage::BlockHeaderStorage; +use crate::storage::{BlockHeaderStorage, MasternodeStorage, PersistentMasternodeStorage}; use crate::sync::{MasternodesProgress, SyncEvent, SyncManager, SyncState}; use dashcore::network::message_qrinfo::QRInfo; +use dashcore::network::message_sml::MnListDiff; use dashcore::BlockHash; use std::collections::BTreeSet; @@ -299,6 +300,7 @@ pub struct MasternodesManager { network: dashcore::Network, /// Sync state tracking. pub(super) sync_state: MasternodeSyncState, + pub(super) message_storage: Option>>>, } impl MasternodesManager { @@ -307,6 +309,7 @@ impl MasternodesManager { header_storage: Arc>, engine: Arc>, network: dashcore::Network, + message_storage: Option>>>, ) -> Self { // Recover sync state from the engine's stored masternode lists so that a // restart can resume from where the previous run left off. @@ -337,9 +340,37 @@ impl MasternodesManager { engine, network, sync_state, + message_storage, } } + pub(super) async fn store_diff(&self, height: u32, diff: &MnListDiff) { + let Some(storage) = &self.message_storage else { + return; + }; + if let Err(e) = storage.write().await.store_diff(height, diff).await { + tracing::warn!("Could not store MnListDiff at {height}: {e}"); + } + } + + pub(super) async fn store_qr_info(&self, height: u32, qr_info: &QRInfo) { + let Some(storage) = &self.message_storage else { + return; + }; + if let Err(e) = storage.write().await.store_qr_info(height, qr_info).await { + tracing::warn!("Could not store QRInfo at {height}: {e}"); + } + } + + pub(super) async fn prune_retained_lists(&self, tip: u32) { + if self.message_storage.is_none() { + return; + } + + let pruned = self.engine.write().await.prune_masternode_lists(tip); + tracing::debug!("Pruned {pruned} in-memory masternode lists at {tip}"); + } + /// Decide which [`PipelineMode`] to use when a new header lands at `tip_height` /// and masternode sync needs to catch up. The rule is: /// @@ -559,6 +590,7 @@ impl MasternodesManager { self.sync_state.last_synced_block_hash = Some(latest_block_hash); self.progress.update_current_height(height); + self.prune_retained_lists(height).await; tracing::debug!("Incremental MnListDiff complete at height {}", height); Ok(vec![SyncEvent::MasternodeStateUpdated { height, @@ -662,6 +694,10 @@ impl MasternodesManager { drop(engine); + if !events.is_empty() { + self.prune_retained_lists(self.progress.current_height()).await; + } + if is_initial_sync { self.set_state(SyncState::Synced); tracing::info!("Masternode sync complete at height {}", self.progress.current_height()); @@ -696,7 +732,7 @@ mod tests { async fn create_test_manager_for(network: dashcore::Network) -> TestMasternodesManager { let storage = DiskStorageManager::with_temp_dir().await.unwrap(); let engine = Arc::new(RwLock::new(MasternodeListEngine::default_for_network(network))); - MasternodesManager::new(storage.block_headers(), engine, network).await + MasternodesManager::new(storage.block_headers(), engine, network, None).await } async fn create_test_manager() -> TestMasternodesManager { @@ -733,6 +769,7 @@ mod tests { block_headers, Arc::new(RwLock::new(engine)), dashcore::Network::Regtest, + None, ) .await; manager.set_state(SyncState::Synced); @@ -964,6 +1001,7 @@ mod tests { storage.block_headers(), Arc::new(RwLock::new(engine)), dashcore::Network::Testnet, + None, ) .await; diff --git a/dash-spv/src/sync/masternodes/sync_manager.rs b/dash-spv/src/sync/masternodes/sync_manager.rs index 1a077a8a2..3d81236ac 100644 --- a/dash-spv/src/sync/masternodes/sync_manager.rs +++ b/dash-spv/src/sync/masternodes/sync_manager.rs @@ -1,15 +1,13 @@ use super::manager::PipelineMode; use crate::error::SyncResult; use crate::network::{Message, MessageType, RequestSender}; -use crate::storage::BlockHeaderStorage; +use crate::storage::{feed_qrinfo_heights_to_engine, BlockHeaderStorage}; use crate::sync::{ ManagerIdentifier, MasternodesManager, SyncEvent, SyncManager, SyncManagerProgress, SyncState, }; use crate::SyncError; use async_trait::async_trait; use dashcore::network::message::NetworkMessage; -use dashcore::network::message_qrinfo::QRInfo; -use dashcore::sml::masternode_list_engine::{MasternodeListEngine, WORK_DIFF_DEPTH}; use dashcore::{BlockHash, QuorumHash}; use dashcore_hashes::Hash; use std::collections::{BTreeSet, HashSet}; @@ -159,63 +157,6 @@ pub(super) async fn build_mnlistdiff_request_pairs( Ok(pairs_with_height.into_iter().map(|(_, base, target)| (base, target)).collect()) } -/// Feed QRInfo block heights to the engine from storage. -/// -/// Resolves heights for every hash enumerated by -/// [`MasternodeListEngine::qr_info_referenced_block_hashes`], plus the cycle boundary -/// block for each work-block diff (`work_height + WORK_DIFF_DEPTH`), which is needed -/// for rotated quorum storage key calculation. -pub(super) async fn feed_qrinfo_heights_to_engine( - engine: &mut MasternodeListEngine, - qr_info: &QRInfo, - storage: &S, -) -> SyncResult { - let mut fed_count = 0; - for block_hash in MasternodeListEngine::qr_info_referenced_block_hashes(qr_info) { - if let Ok(Some(height)) = storage.get_header_height_by_hash(&block_hash).await { - engine.feed_block_height(height, block_hash); - fed_count += 1; - tracing::trace!("Fed height {} for block {}", height, block_hash); - } - } - - // Feed cycle boundary heights for all diffs (current and historical cycles). - // Each diff's block_hash is at the "work block" height; the cycle boundary is - // WORK_DIFF_DEPTH higher. - let mut work_block_hashes = vec![ - qr_info.mn_list_diff_h.block_hash, - qr_info.mn_list_diff_at_h_minus_c.block_hash, - qr_info.mn_list_diff_at_h_minus_2c.block_hash, - qr_info.mn_list_diff_at_h_minus_3c.block_hash, - ]; - - if let Some((_, diff)) = &qr_info.quorum_snapshot_and_mn_list_diff_at_h_minus_4c { - work_block_hashes.push(diff.block_hash); - } - - for work_block_hash in work_block_hashes { - if let Ok(Some(work_block_height)) = - storage.get_header_height_by_hash(&work_block_hash).await - { - let cycle_boundary_height = work_block_height + WORK_DIFF_DEPTH; - if let Ok(Some(cycle_boundary_header)) = storage.get_header(cycle_boundary_height).await - { - let cycle_boundary_hash = *cycle_boundary_header.hash(); - engine.feed_block_height(cycle_boundary_height, cycle_boundary_hash); - fed_count += 1; - tracing::debug!( - "Fed cycle boundary height {} for block {}", - cycle_boundary_height, - cycle_boundary_hash - ); - } - } - } - - tracing::info!("Fed {} block heights to engine", fed_count); - Ok(fed_count) -} - #[async_trait] impl SyncManager for MasternodesManager { fn identifier(&self) -> ManagerIdentifier { @@ -269,9 +210,8 @@ impl SyncManager for MasternodesManager { // Feed block heights to engine using internal storage let storage = self.header_storage.read().await; let mut engine = self.engine.write().await; - let fed = feed_qrinfo_heights_to_engine(&mut engine, qr_info, &*storage).await?; + feed_qrinfo_heights_to_engine(&mut engine, qr_info, &*storage).await; drop(storage); - tracing::info!("Fed {} block heights to engine", fed); // Feed QRInfo to engine first to populate masternode lists let qr_info_result = match engine.feed_qr_info(qr_info.clone(), true, true) { @@ -295,6 +235,9 @@ impl SyncManager for MasternodesManager { } }; + let qr_info_height = + engine.block_container.get_height(&qr_info.mn_list_diff_tip.block_hash); + // Populate known_mn_list_heights from engine after QRInfo processing self.sync_state.known_mn_list_heights = engine.masternode_lists.keys().copied().collect(); @@ -318,6 +261,14 @@ impl SyncManager for MasternodesManager { drop(engine); drop(storage); + match qr_info_height { + Some(height) => self.store_qr_info(height, qr_info).await, + None => tracing::warn!( + "QRInfo tip {} has no known height, rotated quorums will not survive a restart", + qr_info.mn_list_diff_tip.block_hash + ), + } + if let Some(ref qr_info_result) = qr_info_result { tracing::info!( "QRInfo processed: stored_cycle_height={:?}, rotated_quorum_count={}/{}, fully_verified_count={}, newly_qualified_count={}, cycle_key_unresolved={}, previous_cycle_invalid_count={}", @@ -431,6 +382,10 @@ impl SyncManager for MasternodesManager { }; drop(engine); + if apply_ok { + self.store_diff(target_height, diff).await; + } + self.progress.add_diffs_processed(1); self.sync_state.mnlistdiff_pipeline.receive(diff); self.sync_state.mnlistdiff_pipeline.send_pending(requests)?; @@ -747,14 +702,13 @@ impl SyncManager for MasternodesManager { mod tests { use super::super::manager::{MasternodeSyncState, QRInfoInFlight}; use super::{ - feed_qrinfo_heights_to_engine, qrinfo_timeout_for, MAX_RETRY_ATTEMPTS, - QRINFO_STALL_WATCHDOG, QRINFO_TIMEOUT_SCHEDULE_SECS, + qrinfo_timeout_for, MAX_RETRY_ATTEMPTS, QRINFO_STALL_WATCHDOG, QRINFO_TIMEOUT_SCHEDULE_SECS, }; use crate::error::StorageResult; use crate::network::{Message, NetworkRequest, RequestSender}; use crate::storage::{ - BlockHeaderStorage, BlockHeaderTip, DiskStorageManager, PersistentBlockHeaderStorage, - StorageManager, + feed_qrinfo_heights_to_engine, BlockHeaderStorage, BlockHeaderTip, DiskStorageManager, + PersistentBlockHeaderStorage, StorageManager, }; use crate::sync::{MasternodesManager, SyncManager, SyncState}; use crate::types::HashedBlockHeader; @@ -920,9 +874,7 @@ mod tests { network: Network::Testnet, ..Default::default() }; - feed_qrinfo_heights_to_engine(&mut engine, &qr_info, &MockHeaderStorage(height_map)) - .await - .unwrap(); + feed_qrinfo_heights_to_engine(&mut engine, &qr_info, &MockHeaderStorage(height_map)).await; for &b in expected_hashes { let hash = BlockHash::from_slice(&[b; 32]).unwrap(); @@ -1097,9 +1049,13 @@ mod tests { .await .unwrap(); let engine = MasternodeListEngine::default_for_network(Network::Regtest); - let mut manager = - MasternodesManager::new(block_headers, Arc::new(RwLock::new(engine)), Network::Regtest) - .await; + let mut manager = MasternodesManager::new( + block_headers, + Arc::new(RwLock::new(engine)), + Network::Regtest, + None, + ) + .await; manager.progress.update_block_header_tip_height(tip); let (tx, mut rx) = mpsc::unbounded_channel(); diff --git a/dash-spv/tests/dashd_masternode/helpers.rs b/dash-spv/tests/dashd_masternode/helpers.rs index aa27b7a06..d4b6d0269 100644 --- a/dash-spv/tests/dashd_masternode/helpers.rs +++ b/dash-spv/tests/dashd_masternode/helpers.rs @@ -1,3 +1,6 @@ +use std::collections::BTreeMap; +use std::path::Path; + use dash_spv::sync::{MasternodesProgress, SyncEvent, SyncProgress, SyncState}; use dashcore::ephemerealdata::instant_lock::InstantLock; use dashcore::sml::llmq_entry_verification::LLMQEntryVerificationStatus; @@ -14,6 +17,79 @@ use super::setup::{TestContext, SYNC_TIMEOUT}; /// Mine a DKG cycle and wait for the SPV to surface a `MasternodeStateUpdated` /// event above `baseline_height`. +pub(super) fn storage_snapshot(root: &Path) -> BTreeMap { + let mut counts = BTreeMap::new(); + let Ok(entries) = std::fs::read_dir(root) else { + return counts; + }; + for entry in entries.flatten() { + if !entry.path().is_dir() { + continue; + } + let files = walkdir_count(&entry.path()); + counts.insert(entry.file_name().to_string_lossy().into_owned(), files); + } + counts +} + +fn walkdir_count(dir: &Path) -> usize { + let Ok(entries) = std::fs::read_dir(dir) else { + return 0; + }; + entries + .flatten() + .map(|e| { + let path = e.path(); + if path.is_dir() { + walkdir_count(&path) + } else { + 1 + } + }) + .sum() +} + +pub(super) const EXPECTED_STORAGE: &[(&str, &str)] = &[ + ("block_headers", "headers synced to the tip"), + ("filter_headers", "filter headers synced to the tip"), + ("metadata", "sync checkpoints"), + ("peers", "peer set and reputations"), + ("masternodes", "the masternode messages this session stored"), +]; + +pub(super) fn assert_storage_persisted(snapshot: &BTreeMap, what: &str) { + let missing: Vec = EXPECTED_STORAGE + .iter() + .filter(|(dir, _)| snapshot.get(*dir).is_none_or(|files| *files == 0)) + .map(|(dir, why)| format!(" {dir}/ — {why}")) + .collect(); + assert!( + missing.is_empty(), + "{what}: {} storage director{} empty or absent after a clean shutdown:\n{}\n\nstorage holds {snapshot:?}", + missing.len(), + if missing.len() == 1 { "y is" } else { "ies are" }, + missing.join("\n"), + ); +} + +pub(super) fn assert_storage_did_not_shrink( + before: &BTreeMap, + after: &BTreeMap, + what: &str, +) { + for (dir, before_count) in before { + match after.get(dir) { + None => panic!( + "{what}: storage directory {dir:?} disappeared across the restart\n before: {before:?}\n after: {after:?}" + ), + Some(after_count) if after_count < before_count => panic!( + "{what}: storage directory {dir:?} shrank across the restart, {before_count} -> {after_count}\n before: {before:?}\n after: {after:?}" + ), + Some(_) => {} + } + } +} + pub(super) async fn mine_dkg_cycle_and_wait( ctx: &mut TestContext, sync_event_receiver: &mut broadcast::Receiver, diff --git a/dash-spv/tests/dashd_masternode/tests_sync.rs b/dash-spv/tests/dashd_masternode/tests_sync.rs index 805e469ad..1647b9193 100644 --- a/dash-spv/tests/dashd_masternode/tests_sync.rs +++ b/dash-spv/tests/dashd_masternode/tests_sync.rs @@ -10,8 +10,9 @@ use dashcore::sml::llmq_entry_verification::LLMQEntryVerificationStatus; use dashcore::sml::llmq_type::LLMQType; use super::helpers::{ - assert_all_rotated_quorums_verified, wait_for_chainlock_height_at_least, - wait_for_masternode_sync, wait_for_mn_state_event, wait_for_mn_state_event_above, + assert_all_rotated_quorums_verified, assert_storage_did_not_shrink, assert_storage_persisted, + storage_snapshot, wait_for_chainlock_height_at_least, wait_for_masternode_sync, + wait_for_mn_state_event, wait_for_mn_state_event_above, wait_for_mn_state_with_stored_cycle_above, }; use super::setup::{ @@ -103,9 +104,25 @@ async fn test_masternode_list_sync_with_restart() { let first_mn_progress = wait_for_masternode_sync(&mut client_handle.progress_receiver, SYNC_TIMEOUT).await; let first_height = first_mn_progress.current_height(); + + let first_masternodes = { + let engine = client_handle.engine.read().await; + engine.masternode_lists.values().map(|list| list.masternodes.len()).max().unwrap_or(0) + }; + assert!( + first_masternodes > 0, + "the first session must have a masternode list before its persistence can be tested" + ); + client_handle.stop().await; drop(client_handle); + let after_first = storage_snapshot(ctx.storage_path()); + assert_storage_persisted( + &after_first, + &format!("after a first session that built {first_masternodes} masternode(s)"), + ); + // Restart with same storage tracing::info!("=== Restarting with same storage ==="); let mut client_handle = create_and_start_client(&config, Arc::clone(&wallet)).await; @@ -123,6 +140,9 @@ async fn test_masternode_list_sync_with_restart() { "Should reach Synced state after restart" ); + let after_second = storage_snapshot(ctx.storage_path()); + assert_storage_did_not_shrink(&after_first, &after_second, "masternode restart"); + tracing::info!( "Restart verified: first_height={}, second_height={}", first_height, diff --git a/dash/src/sml/masternode_list_engine/helpers.rs b/dash/src/sml/masternode_list_engine/helpers.rs index 9dfbaeccc..b226bea80 100644 --- a/dash/src/sml/masternode_list_engine/helpers.rs +++ b/dash/src/sml/masternode_list_engine/helpers.rs @@ -2,6 +2,8 @@ use crate::QuorumHash; use crate::prelude::CoreBlockHeight; use crate::sml::llmq_entry_verification::LLMQEntryVerificationStatus; use crate::sml::llmq_type::LLMQType; +#[cfg(feature = "quorum_validation")] +use crate::sml::llmq_type::network::NetworkLLMQExt; use crate::sml::masternode_list::MasternodeList; use crate::sml::masternode_list_engine::MasternodeListEngine; use crate::sml::quorum_entry::qualified_quorum_entry::QualifiedQuorumEntry; @@ -14,6 +16,20 @@ use crate::sml::quorum_entry::qualified_quorum_entry::QualifiedQuorumEntry; const QUORUM_WALK_BACK_ACTIVE_WINDOWS: u32 = 4; impl MasternodeListEngine { + #[cfg(feature = "quorum_validation")] + pub fn prune_masternode_lists(&mut self, tip: CoreBlockHeight) -> usize { + let params = self.network.chain_locks_type().params(); + let floor = tip.saturating_sub( + params + .signing_active_quorum_count + .saturating_mul(params.dkg_params.interval) + .saturating_mul(QUORUM_WALK_BACK_ACTIVE_WINDOWS), + ); + let before = self.masternode_lists.len(); + self.masternode_lists.retain(|height, _| *height >= floor); + before - self.masternode_lists.len() + } + /// Retrieves the closest masternode lists before and after a given core block height. /// /// This function searches the `masternode_lists` map to find the nearest masternode lists diff --git a/dash/src/sml/masternode_list_engine/message_request_verification.rs b/dash/src/sml/masternode_list_engine/message_request_verification.rs index 662626ace..303a06827 100644 --- a/dash/src/sml/masternode_list_engine/message_request_verification.rs +++ b/dash/src/sml/masternode_list_engine/message_request_verification.rs @@ -383,7 +383,7 @@ impl MasternodeListEngine { } /// Helper function to verify a ChainLock using a specific masternode list. - fn verify_chain_lock_with_masternode_list( + pub fn verify_chain_lock_with_masternode_list( &self, chain_lock: &ChainLock, masternode_list: &MasternodeList,