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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion dash-spv/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -1013,7 +1013,6 @@ pub trait SyncManager: Send + Sync + Debug {
fn state(&self) -> SyncState;
fn wanted_message_types(&self) -> &'static [MessageType];

async fn initialize(&mut self) -> SyncResult<()>;
async fn start_sync(&mut self, requests: &RequestSender) -> SyncResult<Vec<SyncEvent>>;
async fn handle_message(&mut self, msg: Message, requests: &RequestSender) -> SyncResult<Vec<SyncEvent>>;
async fn handle_sync_event(&mut self, event: &SyncEvent, requests: &RequestSender) -> SyncResult<Vec<SyncEvent>>;
Expand Down
175 changes: 86 additions & 89 deletions dash-spv/src/client/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,12 +36,16 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
pub async fn new(
config: ClientConfig,
network: N,
storage: S,
mut storage: S,
wallet: Arc<RwLock<W>>,
) -> Result<Self> {
// Validate configuration
config.validate().map_err(SpvError::Config)?;

// Initialize genesis block or checkpoint before creating managers,
// so they can read the tip from storage during construction.
Self::initialize_genesis_block(&config, &mut storage).await?;

let masternode_engine = {
if config.enable_masternodes {
Some(Arc::new(RwLock::new(MasternodeListEngine::default_for_network(
Expand All @@ -68,49 +72,60 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
};
let checkpoint_manager = Arc::new(CheckpointManager::new(checkpoints));
managers.block_headers =
Some(BlockHeadersManager::new(storage.block_headers(), checkpoint_manager));
Some(BlockHeadersManager::new(storage.block_headers(), checkpoint_manager).await?);

if config.enable_filters {
managers.filter_headers =
Some(FilterHeadersManager::new(storage.block_headers(), storage.filter_headers()));
managers.filters = Some(FiltersManager::new(
wallet.clone(),
storage.block_headers(),
storage.filter_headers(),
storage.filters(),
));
managers.blocks =
Some(BlocksManager::new(wallet.clone(), storage.block_headers(), storage.blocks()));
managers.filter_headers = Some(
FilterHeadersManager::new(storage.block_headers(), storage.filter_headers())
.await?,
);
managers.filters = Some(
FiltersManager::new(
wallet.clone(),
storage.block_headers(),
storage.filter_headers(),
storage.filters(),
)
.await,
);
managers.blocks = Some(
BlocksManager::new(wallet.clone(), storage.block_headers(), storage.blocks()).await,
);
}

// Build masternode manager if enabled
if config.enable_masternodes {
let masternode_list_engine = masternode_engine
.clone()
.expect("Masternode list engine must exist if masternodes are enabled");
managers.masternode = Some(MasternodesManager::new(
storage.block_headers(),
masternode_list_engine.clone(),
config.network,
));
managers.chainlock = Some(ChainLockManager::new(
storage.block_headers(),
storage.metadata(),
masternode_list_engine.clone(),
));
managers.masternode = Some(
MasternodesManager::new(
storage.block_headers(),
masternode_list_engine.clone(),
config.network,
)
.await,
);
managers.chainlock = Some(
ChainLockManager::new(
storage.block_headers(),
storage.metadata(),
masternode_list_engine.clone(),
)
.await,
);
managers.instantsend = Some(InstantSendManager::new(masternode_list_engine.clone()));
}

// Create sync coordinator (managers are passed to start() later)
let sync_coordinator = SyncCoordinator::new(managers);
let sync_coordinator = SyncCoordinator::new(managers).await;

// Create mempool state
let mempool_state = Arc::new(RwLock::new(MempoolState::default()));

// Wrap storage in Arc<Mutex>
let storage = Arc::new(Mutex::new(storage));

Ok(Self {
let client = Self {
config: Arc::new(RwLock::new(config)),
network: Arc::new(Mutex::new(network)),
storage,
Expand All @@ -120,56 +135,51 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
running: Arc::new(RwLock::new(false)),
mempool_state,
mempool_filter: Arc::new(RwLock::new(None)),
})
}

/// Start the SPV client.
pub(super) async fn start(&self) -> Result<()> {
{
let running = self.running.read().await;
if *running {
return Err(SpvError::Config("Client already running".to_string()));
}
}
};

// Load wallet data from storage
self.load_wallet_data().await?;

let config = self.config.read().await;
client.load_wallet_data().await?;

// Initialize mempool filter if mempool tracking is enabled
if config.enable_mempool_tracking {
// TODO: Get monitored addresses from wallet
let filter = Arc::new(MempoolFilter::new(
config.mempool_strategy,
config.max_mempool_transactions,
self.mempool_state.clone(),
HashSet::new(), // Will be populated from wallet's monitored addresses
config.network,
));

*self.mempool_filter.write().await = Some(filter);

// Load mempool state from storage if persistence is enabled
if config.persist_mempool {
if let Some(state) = self
.storage
.lock()
.await
.load_mempool_state()
.await
.map_err(SpvError::Storage)?
{
*self.mempool_state.write().await = state;
{
let config = client.config.read().await;
if config.enable_mempool_tracking {
let filter = Arc::new(MempoolFilter::new(
config.mempool_strategy,
config.max_mempool_transactions,
client.mempool_state.clone(),
HashSet::new(), // TODO: populate from wallet's monitored addresses
config.network,
));
*client.mempool_filter.write().await = Some(filter);

// Load mempool state from storage if persistence is enabled
if config.persist_mempool {
if let Some(state) = client
.storage
.lock()
.await
.load_mempool_state()
.await
.map_err(SpvError::Storage)?
{
*client.mempool_state.write().await = state;
}
}
}
}

// Drop config before calling methods that also read it
drop(config);
Ok(client)
}

// Initialize genesis block if not already present
self.initialize_genesis_block().await?;
/// Start the SPV client: spawn sync tasks and connect to the network.
pub(super) async fn start(&self) -> Result<()> {
{
let running = self.running.read().await;
if *running {
return Err(SpvError::Config("Client already running".to_string()));
}
}

// Start all sync tasks before connecting to the network to make sure initial connection
// events are handled correctly in the sync coordinator.
Expand Down Expand Up @@ -230,15 +240,12 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
self.stop().await
}

/// Initialize genesis block or checkpoint.
pub(super) async fn initialize_genesis_block(&self) -> Result<()> {
let config = self.config.read().await;

/// Initialize genesis block or checkpoint in storage.
///
/// Called before creating managers so they can read the tip during construction.
async fn initialize_genesis_block(config: &ClientConfig, storage: &mut S) -> Result<()> {
// Check if we already have any headers in storage
let current_tip = {
let storage = self.storage.lock().await;
storage.get_tip_height().await
};
let current_tip = storage.get_tip_height().await;

if current_tip.is_some() {
// We already have headers, genesis block should be at height 0
Expand Down Expand Up @@ -298,12 +305,9 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
calculated_hash
);
} else {
{
let mut storage = self.storage.lock().await;
storage
.store_headers_at_height(&[checkpoint_header], checkpoint.height)
.await?;
}
storage
.store_headers_at_height(&[checkpoint_header], checkpoint.height)
.await?;

tracing::info!(
"✅ Initialized from checkpoint at height {}, skipping {} headers",
Expand Down Expand Up @@ -343,17 +347,10 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
tracing::debug!("Using genesis block header with hash: {}", calculated_hash);

// Store the genesis header at height 0
let genesis_headers = vec![genesis_header];
{
let mut storage = self.storage.lock().await;
storage.store_headers(&genesis_headers).await.map_err(SpvError::Storage)?;
}
storage.store_headers(&[genesis_header]).await.map_err(SpvError::Storage)?;

// Verify it was stored correctly
let stored_height = {
let storage = self.storage.lock().await;
storage.get_tip_height().await
};
let stored_height = storage.get_tip_height().await;
tracing::info!(
"✅ Genesis block initialized at height 0, storage reports tip height: {:?}",
stored_height
Expand Down
38 changes: 30 additions & 8 deletions dash-spv/src/sync/block_headers/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,13 +52,30 @@ impl<H: BlockHeaderStorage> std::fmt::Debug for BlockHeadersManager<H> {

impl<H: BlockHeaderStorage> BlockHeadersManager<H> {
/// Create a new headers manager with the given storage and checkpoint manager.
pub fn new(header_storage: Arc<RwLock<H>>, checkpoint_manager: Arc<CheckpointManager>) -> Self {
Self {
progress: BlockHeadersProgress::default(),
pub async fn new(
header_storage: Arc<RwLock<H>>,
checkpoint_manager: Arc<CheckpointManager>,
) -> SyncResult<Self> {
let tip = header_storage
.read()
.await
.get_tip()
.await
.ok_or_else(|| SyncError::MissingDependency("No tip in storage".to_string()))?;

let mut initial_progress = BlockHeadersProgress::default();
initial_progress.set_state(SyncState::WaitingForConnections);
initial_progress.update_tip_height(tip.height());
initial_progress.update_target_height(tip.height());

tracing::info!("BlockHeadersManager initialized at height {}", tip.height());

Ok(Self {
progress: initial_progress,
header_storage,
pipeline: HeadersPipeline::new(checkpoint_manager.clone()),
pipeline: HeadersPipeline::new(checkpoint_manager),
pending_announcements: HashMap::new(),
}
})
}

pub(super) async fn tip(&self) -> SyncResult<BlockHeaderTip> {
Expand Down Expand Up @@ -227,16 +244,21 @@ mod tests {
}

async fn create_test_manager() -> TestBlockHeadersManager {
let storage = DiskStorageManager::with_temp_dir().await.unwrap();
let mut storage = DiskStorageManager::with_temp_dir().await.unwrap();
// Store a genesis header so the manager can initialize
let genesis = Header::dummy_batch(0..1);
storage.store_headers(&genesis).await.unwrap();
let checkpoint_manager = create_test_checkpoint_manager();
BlockHeadersManager::new(storage.block_headers(), checkpoint_manager)
.await
.expect("Failed to create BlockHeadersManager")
}

#[tokio::test]
async fn test_block_headers_manager_new() {
let manager = create_test_manager().await;
assert_eq!(manager.identifier(), ManagerIdentifier::BlockHeader);
assert_eq!(manager.state(), SyncState::WaitForEvents);
assert_eq!(manager.state(), SyncState::WaitingForConnections);
assert_eq!(manager.wanted_message_types(), vec![MessageType::Headers, MessageType::Inv]);
}

Expand All @@ -249,7 +271,7 @@ mod tests {

let progress = manager.progress();
if let SyncManagerProgress::BlockHeaders(progress) = progress {
assert_eq!(progress.state(), SyncState::WaitForEvents);
assert_eq!(progress.state(), SyncState::WaitingForConnections);
assert_eq!(progress.tip_height(), 100);
assert_eq!(progress.target_height(), 200);
assert_eq!(progress.processed(), 50);
Expand Down
19 changes: 0 additions & 19 deletions dash-spv/src/sync/block_headers/sync_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@ use crate::sync::{
BlockHeadersManager, ManagerIdentifier, ProgressPercentage, SyncEvent, SyncManager,
SyncManagerProgress, SyncState,
};
use crate::SyncError;
use async_trait::async_trait;
use dashcore::network::message::NetworkMessage;
use dashcore::BlockHash;
Expand Down Expand Up @@ -37,24 +36,6 @@ impl<H: BlockHeaderStorage> SyncManager for BlockHeadersManager<H> {
&[MessageType::Headers, MessageType::Inv]
}

async fn initialize(&mut self) -> SyncResult<()> {
let tip = self
.header_storage
.read()
.await
.get_tip()
.await
.ok_or_else(|| SyncError::MissingDependency("No tip in storage".to_string()))?;

self.progress.set_state(SyncState::WaitingForConnections);
self.progress.update_tip_height(tip.height());
self.progress.update_target_height(tip.height());

tracing::info!("BlockHeadersManager initialized at height {}", tip.height());

Ok(())
}

async fn start_sync(&mut self, requests: &RequestSender) -> SyncResult<Vec<SyncEvent>> {
ensure_not_started(self.state(), self.identifier())?;
self.progress.set_state(SyncState::Syncing);
Expand Down
11 changes: 8 additions & 3 deletions dash-spv/src/sync/blocks/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,13 +43,18 @@ pub struct BlocksManager<H: BlockHeaderStorage, B: BlockStorage, W: WalletInterf

impl<H: BlockHeaderStorage, B: BlockStorage, W: WalletInterface> BlocksManager<H, B, W> {
/// Create a new blocks manager with the given storage references.
pub fn new(
pub async fn new(
wallet: Arc<RwLock<W>>,
header_storage: Arc<RwLock<H>>,
block_storage: Arc<RwLock<B>>,
) -> Self {
let synced_height = wallet.read().await.synced_height();

let mut initial_progress = BlocksProgress::default();
initial_progress.update_last_processed(synced_height);

Self {
progress: BlocksProgress::default(),
progress: initial_progress,
header_storage,
block_storage,
wallet,
Expand Down Expand Up @@ -170,7 +175,7 @@ mod tests {
async fn create_test_manager() -> TestBlocksManager {
let storage = DiskStorageManager::with_temp_dir().await.unwrap();
let wallet = Arc::new(RwLock::new(MockWallet::new()));
BlocksManager::new(wallet, storage.block_headers(), storage.blocks())
BlocksManager::new(wallet, storage.block_headers(), storage.blocks()).await
}

#[tokio::test]
Expand Down
Loading
Loading