Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
bf23546
refactor(dash-spv): network module refactor
ZocoLini Jul 29, 2026
2cd965b
test(dash-spv): recovered removed tests during the refactor
ZocoLini Jul 29, 2026
01f9377
tests(dash-spv): network manager unit tests
ZocoLini Jul 29, 2026
6cbd3ce
coderabbit comments
ZocoLini Jul 30, 2026
67e5525
fix(dash-spv): cfilters timeout loop
ZocoLini Aug 2, 2026
f71cac0
test(dash-spv): rescan validation sugested by coderabbit
ZocoLini Aug 2, 2026
9389e92
feat(dash-spv): send GetHeaders2 when possible
ZocoLini Aug 10, 2026
d5d9ec0
feat(dah-spv): copied some logic from Kevins approach when doing peer…
ZocoLini Aug 11, 2026
0f93dec
feat(dash-spv): prune spent single-use CoinJoin addresses from the fi…
QuantumExplorer Aug 11, 2026
429d529
feat(dash-spv): learn peer addresses from addrv2 gossip
ZocoLini Aug 11, 2026
ecdc329
Merge remote-tracking branch 'origin/dev' into refactor/network-mod
ZocoLini Aug 11, 2026
18298bc
refactor(dash-spv): drop the unimplemented on_peer_disconnect hook
ZocoLini Aug 11, 2026
a8f5189
test(dash-spv): cover requeuing a departed peer's in-flight requests
ZocoLini Aug 11, 2026
e349dc3
feat(dash-spv): ask peers for addresses with getaddr
ZocoLini Aug 11, 2026
d28a574
fix(dash-spv): keep socket writes out of the peer-set lock
ZocoLini Aug 11, 2026
a5ff048
fix(dash-spv): recover from a rejected header batch instead of wedging
ZocoLini Aug 11, 2026
d0b30b2
fix(dash-spv): measure peer latency before asking for addresses
ZocoLini Aug 11, 2026
589584d
fix(dash-spv): size peer in-flight by its bandwidth-delay product
ZocoLini Aug 12, 2026
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
4 changes: 1 addition & 3 deletions dash-spv-bench/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -154,9 +154,7 @@ async fn main() -> Result<()> {
let wallet_probe = wallet.clone();

let handler = Arc::new(BenchEventHandler::new(dashboard.clone()));
let network = dash_spv::network::PeerNetworkManager::new(&config)
.await
.map_err(|e| anyhow!("network new: {e}"))?;
let network = dash_spv::network::PeerNetworkManager::new(&config).await;

let client = DashSpvClient::new(
config,
Expand Down
12 changes: 4 additions & 8 deletions dash-spv-ffi/src/callbacks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -349,6 +349,7 @@ impl FFISyncEventCallbacks {
start_height,
end_height,
tip_height,
..
} => {
if let Some(cb) = self.on_filter_headers_stored {
cb(*start_height, *end_height, *tip_height, self.user_data);
Expand Down Expand Up @@ -550,17 +551,13 @@ impl FFINetworkEventCallbacks {
use dash_spv::network::NetworkEvent;

match event {
NetworkEvent::PeerConnected {
address,
} => {
NetworkEvent::PeerConnected(address) => {
if let Some(cb) = self.on_peer_connected {
let c_addr = CString::new(address.to_string()).unwrap_or_default();
cb(c_addr.as_ptr(), self.user_data);
}
}
NetworkEvent::PeerDisconnected {
address,
} => {
NetworkEvent::PeerDisconnected(address) => {
if let Some(cb) = self.on_peer_disconnected {
let c_addr = CString::new(address.to_string()).unwrap_or_default();
cb(c_addr.as_ptr(), self.user_data);
Expand All @@ -569,10 +566,9 @@ impl FFINetworkEventCallbacks {
NetworkEvent::PeersUpdated {
connected_count,
best_height,
..
} => {
if let Some(cb) = self.on_peers_updated {
cb(*connected_count as u32, best_height.unwrap_or(0), self.user_data);
cb(*connected_count, *best_height, self.user_data);
}
}
}
Expand Down
9 changes: 4 additions & 5 deletions dash-spv-ffi/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,15 +82,15 @@ pub unsafe extern "C" fn dash_spv_ffi_client_new(

let client_result = runtime.block_on(async move {
// Construct concrete implementations for generics
let network = dash_spv::network::PeerNetworkManager::new(&client_config).await;
let storage = DiskStorageManager::new(&client_config).await;
let wallet = key_wallet_manager::WalletManager::<
key_wallet::wallet::managed_wallet_info::ManagedWalletInfo,
>::new(client_config.network);
let wallet = std::sync::Arc::new(tokio::sync::RwLock::new(wallet));

match (network, storage) {
(Ok(network), Ok(storage)) => {
match storage {
Ok(storage) => {
let network = dash_spv::network::PeerNetworkManager::new(&client_config).await;
DashSpvClient::new(
client_config,
network,
Expand All @@ -100,8 +100,7 @@ pub unsafe extern "C" fn dash_spv_ffi_client_new(
)
.await
}
(Err(e), _) => Err(e),
(_, Err(e)) => Err(dash_spv::SpvError::Storage(e)),
Err(e) => Err(dash_spv::SpvError::Storage(e)),
}
});

Expand Down
2 changes: 1 addition & 1 deletion dash-spv/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ rust-version = "1.89"

[dependencies]
# Core Dash libraries
dashcore = { path = "../dash", features = ["serde", "core-block-hash-use-x11", "message_verification", "bls", "quorum_validation"] }
dashcore = { path = "../dash", features = ["serde", "core-block-hash-use-x11", "message_verification", "bls", "quorum_validation", "tokio"] }
dashcore_hashes = { path = "../hashes" }
dash-network-seeds = { path = "../dash-network-seeds" }
key-wallet = { path = "../key-wallet" }
Expand Down
2 changes: 1 addition & 1 deletion dash-spv/examples/filter_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.without_masternodes(); // Skip masternode sync for this example

// Create network manager
let network_manager = PeerNetworkManager::new(&config).await?;
let network_manager = PeerNetworkManager::new(&config).await;

// Create storage manager
let storage_manager = DiskStorageManager::new(&config).await?;
Expand Down
2 changes: 1 addition & 1 deletion dash-spv/examples/simple_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.without_masternodes(); // Skip masternode sync for this example

// Create network manager
let network_manager = PeerNetworkManager::new(&config).await?;
let network_manager = PeerNetworkManager::new(&config).await;

// Create storage manager
let storage_manager = DiskStorageManager::new(&config).await?;
Expand Down
2 changes: 1 addition & 1 deletion dash-spv/examples/spv_with_wallet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.with_validation_mode(dash_spv::ValidationMode::Full);

// Create network manager
let network_manager = PeerNetworkManager::new(&config).await?;
let network_manager = PeerNetworkManager::new(&config).await;

// Create storage manager - use disk storage for persistence
let storage_manager = DiskStorageManager::new(&config).await?;
Expand Down
4 changes: 3 additions & 1 deletion dash-spv/src/client/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ pub(super) type PersistentSyncCoordinator<W> = SyncCoordinator<
/// The generic design is an intentional, beneficial architectural choice for a library.
pub struct DashSpvClient<W: WalletInterface, N: NetworkManager, S: StorageManager> {
pub(super) config: Arc<RwLock<ClientConfig>>,
pub(super) network: Arc<Mutex<N>>,
pub(super) network: Arc<N>,
pub(super) storage: Arc<Mutex<S>>,
/// External wallet implementation (required)
pub(super) wallet: Arc<RwLock<W>>,
Expand All @@ -114,6 +114,7 @@ pub struct DashSpvClient<W: WalletInterface, N: NetworkManager, S: StorageManage
/// `true` while running, `false` once a stop is requested. Stored as a
/// `watch` so a stop is observed immediately rather than polled.
pub(super) running: Arc<watch::Sender<bool>>,
pub(super) stop_requested: Arc<std::sync::atomic::AtomicBool>,
pub(super) event_handlers: Arc<Vec<Arc<dyn super::EventHandler>>>,
}

Expand All @@ -127,6 +128,7 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> Clone for DashSpv
masternode_engine: self.masternode_engine.clone(),
sync_coordinator: Arc::clone(&self.sync_coordinator),
running: Arc::clone(&self.running),
stop_requested: Arc::clone(&self.stop_requested),
event_handlers: Arc::clone(&self.event_handlers),
}
}
Expand Down
13 changes: 3 additions & 10 deletions dash-spv/src/client/event_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -327,8 +327,7 @@ mod tests {
handler.on_sync_event(&event);
handler.on_network_event(&NetworkEvent::PeersUpdated {
connected_count: 0,
addresses: vec![],
best_height: None,
best_height: 0,
});
handler.on_progress(&SyncProgress::default());
handler.on_error("test error");
Expand Down Expand Up @@ -533,14 +532,8 @@ mod tests {
);

let addr: SocketAddr = "127.0.0.1:9999".parse().unwrap();
tx.send(NetworkEvent::PeerConnected {
address: addr,
})
.unwrap();
tx.send(NetworkEvent::PeerDisconnected {
address: addr,
})
.unwrap();
tx.send(NetworkEvent::PeerConnected(addr)).unwrap();
tx.send(NetworkEvent::PeerDisconnected(addr)).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(50)).await;

shutdown.cancel();
Expand Down
2 changes: 1 addition & 1 deletion dash-spv/src/client/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,6 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,

/// Subscribe to network events.
pub(crate) async fn subscribe_network_events(&self) -> broadcast::Receiver<NetworkEvent> {
self.network.lock().await.subscribe_network_events()
self.network.events()
}
}
69 changes: 50 additions & 19 deletions dash-spv/src/client/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -156,12 +156,13 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,

let client = Self {
config: Arc::new(RwLock::new(config)),
network: Arc::new(Mutex::new(network)),
network: Arc::new(network),
storage,
wallet,
masternode_engine,
sync_coordinator: Arc::new(Mutex::new(sync_coordinator)),
running: Arc::new(watch::Sender::new(false)),
stop_requested: Arc::new(std::sync::atomic::AtomicBool::new(false)),
event_handlers: Arc::new(event_handlers),
};

Expand All @@ -183,47 +184,77 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
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.
if let Err(e) =
self.sync_coordinator.lock().await.start(&mut *self.network.lock().await).await
{
// Fresh lifecycle: clear any stop request left over from a previous run so
// this start can complete. A `stop()` arriving *after* this point (while we
// connect below) is recorded and honored by the guarded transition at the end.
self.stop_requested.store(false, std::sync::atomic::Ordering::SeqCst);

// Spawn every sync task BEFORE connecting, so each manager has subscribed by the
// time the network announces its peers. The network manager used to connect inside
// `new()`, which meant `PeersUpdated` and every `PeerConnected` fired into the void
// — and the mempool manager, whose peer set is built from those events, stayed empty
// and never sent the `filterload` that enables transaction relay.
let network: Arc<dyn NetworkManager> = self.network.clone();
if let Err(e) = self.sync_coordinator.lock().await.start(&network).await {
tracing::error!("Failed to start sync coordinator: {}", e);
return Err(SpvError::Sync(e));
}

// Connect to network
self.network.lock().await.connect().await?;

// Only mark as running after all startup operations succeed.
// `send_replace` always stores the value regardless of receiver count,
// so this is correct even when `run()` has not subscribed yet.
self.running.send_replace(true);
self.network.start();

// Only mark as running after all startup operations succeed — and only if
// no `stop()` raced in while we were connecting. The check runs inside the
// watch lock (via `send_if_modified`), and `stop()` sets `stop_requested`
// before it flips `running`, so the two orderings are both safe:
// - we win the lock first: set running=true; a later stop() flips it false.
// - stop() won: `stop_requested` is already true here, so we leave running
// false and the run loop tears down immediately instead of syncing forever.
// `send_if_modified` stores the value regardless of receiver count, so this
// is correct even when `run()` has not subscribed yet.
self.running.send_if_modified(|running| {
if self.stop_requested.load(std::sync::atomic::Ordering::SeqCst) {
false
} else {
*running = true;
true
}
});

Ok(())
}

/// Stop the SPV client.
pub async fn stop(&self) -> Result<()> {
// Check if already stopped
if !*self.running.borrow() {
return Ok(());
}
// Record the stop request BEFORE flipping `running`, so a `start()` still
// connecting observes it (under the watch lock) and declines to mark the
// client running. Otherwise a stop that arrives mid-startup would be lost:
// `start()` would flip running true afterwards and the run task would sync
// forever, hanging `run_handle.await`.
self.stop_requested.store(true, std::sync::atomic::Ordering::SeqCst);

// Flip the running state before tearing anything down so a concurrent
// `run()` loop wakes immediately and breaks out before it can lock the
// sync coordinator again. This prevents a tick from racing against the
// shutdown below.
self.running.send_replace(false);

// Always tear down, even if `running` was never true: a stop that raced
// startup leaves `running` false yet `start()` may already have spawned
// the coordinator's manager tasks and connected the network, so those must
// still be stopped. Every step below is idempotent (cancelling an already
// cancelled token, draining an empty task set, persisting again), so the
// normal double call — this stop() plus the run loop's own final stop() —
// is safe.

// Shut down sync coordinator: signals cancellation and waits for manager
// tasks to drain before we tear down the network and storage layers.
if let Err(e) = self.sync_coordinator.lock().await.shutdown().await {
tracing::warn!("Error shutting down sync coordinator: {}", e);
}

// Disconnect from network
self.network.lock().await.disconnect().await?;
// Tear down the network layer: stops the router/pump and cancels every
// peer reader so no more messages arrive after stop() returns.
self.network.stop();

// Shutdown storage to ensure all data is persisted
{
Expand Down
7 changes: 1 addition & 6 deletions dash-spv/src/client/queries.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,12 +24,7 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,

/// Get the number of connected peers.
pub async fn peer_count(&self) -> usize {
self.network.lock().await.peer_count()
}

/// Disconnect a specific peer.
pub async fn disconnect_peer(&self, addr: &std::net::SocketAddr, reason: &str) -> Result<()> {
Ok(self.network.lock().await.disconnect_peer(addr, reason).await?)
self.network.connected_count().await as usize
}

// ============ Masternode Queries ============
Expand Down
8 changes: 3 additions & 5 deletions dash-spv/src/client/transactions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,21 +30,19 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
/// With mempool tracking disabled there is no acceptance tracking and the
/// transaction is simply sent to all connected peers.
pub async fn broadcast_transaction(&self, tx: &dashcore::Transaction) -> Result<()> {
let network_guard = self.network.lock().await;

if network_guard.peer_count() == 0 {
if self.network.connected_count().await == 0 {
return Err(SpvError::Network(NetworkError::NotConnected));
}

if !self.config.read().await.enable_mempool_tracking {
// Legacy untracked path: fan out to every peer.
network_guard.broadcast(NetworkMessage::Tx(tx.clone())).await?;
self.network.broadcast(NetworkMessage::Tx(tx.clone()));
}

// Inject locally so the mempool manager picks it up through handle_tx.
// With tracking enabled the manager performs the actual (targeted)
// network send when it processes this message.
network_guard.dispatch_local(NetworkMessage::Tx(tx.clone())).await;
self.network.dispatch_local(NetworkMessage::Tx(tx.clone())).await;

Ok(())
}
Expand Down
2 changes: 1 addition & 1 deletion dash-spv/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@
//! .with_storage_path("./.tmp/example-storage");
//!
//! // Create the required components
//! let network = PeerNetworkManager::new(&config).await?;
//! let network = PeerNetworkManager::new(&config).await;
//! let storage = DiskStorageManager::new(&config).await?;
//! let wallet = Arc::new(RwLock::new(WalletManager::<ManagedWalletInfo>::new(config.network)));
//!
Expand Down
12 changes: 3 additions & 9 deletions dash-spv/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -259,13 +259,7 @@ async fn run() -> Result<(), Box<dyn std::error::Error>> {
let wallet = Arc::new(tokio::sync::RwLock::new(wallet_manager));

// Create network manager
let network_manager = match dash_spv::network::manager::PeerNetworkManager::new(&config).await {
Ok(nm) => nm,
Err(e) => {
eprintln!("Failed to create network manager: {}", e);
process::exit(1);
}
};
let network_manager = dash_spv::network::PeerNetworkManager::new(&config).await;

let storage_manager = match dash_spv::storage::DiskStorageManager::new(&config).await {
Ok(sm) => sm,
Expand Down Expand Up @@ -382,15 +376,15 @@ fn parse_llmq_devnet_params(raw: &str) -> Result<LlmqDevnetParams, String> {

async fn run_client<S: dash_spv::storage::StorageManager>(
config: ClientConfig,
network_manager: dash_spv::network::manager::PeerNetworkManager,
network_manager: dash_spv::network::PeerNetworkManager,
storage_manager: S,
wallet: Arc<tokio::sync::RwLock<WalletManager<ManagedWalletInfo>>>,
) -> Result<(), Box<dyn std::error::Error>> {
// Create and start the client
let client =
match DashSpvClient::<
WalletManager<ManagedWalletInfo>,
dash_spv::network::manager::PeerNetworkManager,
dash_spv::network::PeerNetworkManager,
S,
>::new(
config.clone(), network_manager, storage_manager, wallet.clone(), Vec::new()
Expand Down
Loading