diff --git a/dash-spv/src/client/lifecycle.rs b/dash-spv/src/client/lifecycle.rs index 46e26f71c..11452e20f 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, + MasternodeStateStorage, PersistentBlockHeaderStorage, PersistentBlockStorage, + PersistentFilterHeaderStorage, PersistentFilterStorage, PersistentMetadataStorage, + StorageManager, }; use crate::sync::{ BlockHeadersManager, BlocksManager, ChainLockManager, FilterHeadersManager, FiltersManager, @@ -65,11 +66,17 @@ impl DashSpvClient DashSpvClient StorageResult<()>; - - async fn load_masternode_state(&self) -> StorageResult>; + async fn store_engine( + &mut self, + engine: &MasternodeListEngine, + height: u32, + ) -> StorageResult<()>; + + /// Always yields an engine: with nothing persisted yet, the network's + /// default, which is what a first run starts from anyway. + async fn load_engine(&self, network: Network) -> StorageResult; } pub struct PersistentMasternodeStateStorage { @@ -39,13 +53,31 @@ impl PersistentStorage for PersistentMasternodeStateStorage { #[async_trait] impl MasternodeStateStorage for PersistentMasternodeStateStorage { - async fn store_masternode_state(&mut self, state: &MasternodeState) -> StorageResult<()> { + async fn store_engine( + &mut self, + engine: &MasternodeListEngine, + height: u32, + ) -> 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?; - let json = serde_json::to_string_pretty(state).map_err(|e| { + let state = MasternodeState { + last_height: height, + engine_state: serde_json::to_vec(engine).map_err(|e| { + crate::error::StorageError::Serialization(format!( + "Failed to serialize masternode engine: {}", + e + )) + })?, + last_update: std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0), + }; + + let json = serde_json::to_string_pretty(&state).map_err(|e| { crate::error::StorageError::Serialization(format!( "Failed to serialize masternode state: {}", e @@ -56,21 +88,29 @@ impl MasternodeStateStorage for PersistentMasternodeStateStorage { Ok(()) } - async fn load_masternode_state(&self) -> StorageResult> { + async fn load_engine(&self, network: Network) -> StorageResult { let path = self.storage_path.join(Self::FOLDER_NAME).join(Self::MASTERNODE_FILE_NAME); if !path.exists() { - return Ok(None); + tracing::debug!("No persisted masternode state, starting from the network default"); + return Ok(MasternodeListEngine::default_for_network(network)); } let content = tokio::fs::read_to_string(path).await?; - let state = serde_json::from_str(&content).map_err(|e| { + let state: MasternodeState = serde_json::from_str(&content).map_err(|e| { crate::error::StorageError::Serialization(format!( "Failed to deserialize masternode state: {}", e )) })?; + let engine = serde_json::from_slice(&state.engine_state).map_err(|e| { + crate::error::StorageError::Serialization(format!( + "Failed to deserialize masternode engine: {}", + e + )) + })?; - Ok(Some(state)) + tracing::debug!("Loaded masternode engine from height {}", state.last_height); + Ok(engine) } } diff --git a/dash-spv/src/storage/mod.rs b/dash-spv/src/storage/mod.rs index 70a851acb..cfe974b77 100644 --- a/dash-spv/src/storage/mod.rs +++ b/dash-spv/src/storage/mod.rs @@ -19,6 +19,8 @@ use crate::ClientConfig; use async_trait::async_trait; use dashcore::hash_types::FilterHeader; use dashcore::prelude::CoreBlockHeight; +use dashcore::sml::masternode_list_engine::MasternodeListEngine; +use dashcore::Network; use std::ops::Range; use std::path::{Path, PathBuf}; use std::sync::Arc; @@ -78,6 +80,8 @@ pub trait StorageManager: /// Returns shared access to the metadata storage. fn metadata(&self) -> Arc>; + + fn masternodestate(&self) -> Arc>; } /// Disk-based storage manager with segmented files and async background saving. @@ -282,6 +286,10 @@ impl StorageManager for DiskStorageManager { fn metadata(&self) -> Arc> { Arc::clone(&self.metadata) } + + fn masternodestate(&self) -> Arc> { + Arc::clone(&self.masternodestate) + } } #[async_trait] @@ -432,12 +440,16 @@ 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 + async fn store_engine( + &mut self, + engine: &MasternodeListEngine, + height: u32, + ) -> StorageResult<()> { + self.masternodestate.write().await.store_engine(engine, height).await } - async fn load_masternode_state(&self) -> StorageResult> { - self.masternodestate.read().await.load_masternode_state().await + async fn load_engine(&self, network: Network) -> StorageResult { + self.masternodestate.read().await.load_engine(network).await } } diff --git a/dash-spv/src/sync/masternodes/manager.rs b/dash-spv/src/sync/masternodes/manager.rs index 428673535..0571c1d7c 100644 --- a/dash-spv/src/sync/masternodes/manager.rs +++ b/dash-spv/src/sync/masternodes/manager.rs @@ -14,7 +14,9 @@ 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, MasternodeStateStorage, PersistentMasternodeStateStorage, +}; use crate::sync::{MasternodesProgress, SyncEvent, SyncManager, SyncState}; use dashcore::network::message_qrinfo::QRInfo; use dashcore::BlockHash; @@ -299,6 +301,8 @@ pub struct MasternodesManager { network: dashcore::Network, /// Sync state tracking. pub(super) sync_state: MasternodeSyncState, + /// `None` leaves the list in memory only. + pub(super) state_storage: Option>>, } impl MasternodesManager { @@ -307,6 +311,7 @@ impl MasternodesManager { header_storage: Arc>, engine: Arc>, network: dashcore::Network, + state_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,6 +342,21 @@ impl MasternodesManager { engine, network, sync_state, + state_storage, + } + } + + /// Best effort: an unwritten list costs a rebuild next start, a failed sync + /// costs the list now. + pub(super) async fn persist_engine(&self, height: u32) { + let Some(storage) = &self.state_storage else { + return; + }; + let engine = self.engine.read().await; + if let Err(e) = storage.write().await.store_engine(&engine, height).await { + tracing::warn!("Could not persist masternode state at {height}: {e}"); + } else { + tracing::debug!("Persisted masternode state at height {height}"); } } @@ -559,6 +579,7 @@ impl MasternodesManager { self.sync_state.last_synced_block_hash = Some(latest_block_hash); self.progress.update_current_height(height); + self.persist_engine(height).await; tracing::debug!("Incremental MnListDiff complete at height {}", height); Ok(vec![SyncEvent::MasternodeStateUpdated { height, @@ -662,6 +683,10 @@ impl MasternodesManager { drop(engine); + if !events.is_empty() { + self.persist_engine(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 +721,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 +758,7 @@ mod tests { block_headers, Arc::new(RwLock::new(engine)), dashcore::Network::Regtest, + None, ) .await; manager.set_state(SyncState::Synced); @@ -964,6 +990,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..d59b2b00b 100644 --- a/dash-spv/src/sync/masternodes/sync_manager.rs +++ b/dash-spv/src/sync/masternodes/sync_manager.rs @@ -1097,9 +1097,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..8166df9bb 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,94 @@ use super::setup::{TestContext, SYNC_TIMEOUT}; /// Mine a DKG cycle and wait for the SPV to surface a `MasternodeStateUpdated` /// event above `baseline_height`. +/// Files held under each immediate subdirectory of the storage root, keyed by +/// directory name. +/// +/// A sync writes into these and never removes a whole class of state, so across +/// a restart every directory must still be there and hold at least as much — +/// see [`assert_storage_did_not_shrink`]. +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() +} + +/// Directories that must hold state once this test's first session has run, and +/// why. `filters` and `blocks` are deliberately absent: the client is stopped +/// as soon as the masternode phase reports `Synced`, which is before the filter +/// phase leaves `WaitForEvents`, so those stay legitimately empty here. +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"), + ("masternodestate", "the masternode list this session built"), +]; + +/// Assert every directory in [`EXPECTED_STORAGE`] exists and holds at least one +/// file, reporting all of them at once rather than the first to fail. +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"), + ); +} + +/// Every directory present before a restart must still be present after, with +/// at least as many files. A directory that vanishes or shrinks means a restart +/// threw away state that the previous session had already earned. +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..b38d28783 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,29 @@ 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(); + + // Control: the first session really built a list, so the persistence + // assertion below cannot be satisfied by a client that synced nothing. + 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); + // What the first session earned and wrote down. A clean shutdown of a + // fully-synced client must leave every sync phase's state on disk. + 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 +144,11 @@ async fn test_masternode_list_sync_with_restart() { "Should reach Synced state after restart" ); + // A restart re-syncs on top of what it restored; it never discards a whole + // class of state it already had. + 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,