Skip to content

Commit 0ed48ff

Browse files
ZocoLiniclaude
andcommitted
refactor(dash-spv): rewrite the peer-to-peer network module
Replaces the old network layer with a self-coordinating PeerNetworkManager (broker): it owns request de-duplication, pacing, per-peer in-flight sizing, timeouts, retries and peer hot-swap, so the sync pipelines just declare what they want. Highlights: - Restore sync/ from dev and re-wire it to the broker; delete the per-pipeline DownloadCoordinator. - Minimal `NetworkManager` trait + in-memory `MockNetworkManager`; the client is generic over the network (`DashSpvClient<W, N, S>`) with the network injected. - Strict-priority message scheduler: control, then blocks, filters, filter headers, block headers. - Drop peer storage and the peer-reputation system (unused; to be redone). - Un-Box the `NetworkMessage` variants to stay close to dev. - Storage kept as on dev (minus the removed peer storage). - Restore dev's integration tests and the network-dependent unit tests via the mock. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CUb3bkX9C1gBFA65GFsN53
1 parent 9c87ccb commit 0ed48ff

75 files changed

Lines changed: 4479 additions & 9266 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎dash-spv-bench/src/main.rs‎

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -152,9 +152,7 @@ async fn main() -> Result<()> {
152152
let wallet_probe = wallet.clone();
153153

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

159157
let client = DashSpvClient::new(
160158
config,

‎dash-spv-ffi/src/callbacks.rs‎

Lines changed: 4 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -349,6 +349,7 @@ impl FFISyncEventCallbacks {
349349
start_height,
350350
end_height,
351351
tip_height,
352+
..
352353
} => {
353354
if let Some(cb) = self.on_filter_headers_stored {
354355
cb(*start_height, *end_height, *tip_height, self.user_data);
@@ -550,17 +551,13 @@ impl FFINetworkEventCallbacks {
550551
use dash_spv::network::NetworkEvent;
551552

552553
match event {
553-
NetworkEvent::PeerConnected {
554-
address,
555-
} => {
554+
NetworkEvent::PeerConnected(address) => {
556555
if let Some(cb) = self.on_peer_connected {
557556
let c_addr = CString::new(address.to_string()).unwrap_or_default();
558557
cb(c_addr.as_ptr(), self.user_data);
559558
}
560559
}
561-
NetworkEvent::PeerDisconnected {
562-
address,
563-
} => {
560+
NetworkEvent::PeerDisconnected(address) => {
564561
if let Some(cb) = self.on_peer_disconnected {
565562
let c_addr = CString::new(address.to_string()).unwrap_or_default();
566563
cb(c_addr.as_ptr(), self.user_data);
@@ -569,10 +566,9 @@ impl FFINetworkEventCallbacks {
569566
NetworkEvent::PeersUpdated {
570567
connected_count,
571568
best_height,
572-
..
573569
} => {
574570
if let Some(cb) = self.on_peers_updated {
575-
cb(*connected_count as u32, best_height.unwrap_or(0), self.user_data);
571+
cb(*connected_count, *best_height, self.user_data);
576572
}
577573
}
578574
}

‎dash-spv-ffi/src/client.rs‎

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -82,15 +82,15 @@ pub unsafe extern "C" fn dash_spv_ffi_client_new(
8282

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

92-
match (network, storage) {
93-
(Ok(network), Ok(storage)) => {
91+
match storage {
92+
Ok(storage) => {
93+
let network = dash_spv::network::PeerNetworkManager::new(&client_config).await;
9494
DashSpvClient::new(
9595
client_config,
9696
network,
@@ -100,8 +100,7 @@ pub unsafe extern "C" fn dash_spv_ffi_client_new(
100100
)
101101
.await
102102
}
103-
(Err(e), _) => Err(e),
104-
(_, Err(e)) => Err(dash_spv::SpvError::Storage(e)),
103+
Err(e) => Err(dash_spv::SpvError::Storage(e)),
105104
}
106105
});
107106

‎dash-spv/Cargo.toml‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ rust-version = "1.89"
1010

1111
[dependencies]
1212
# Core Dash libraries
13-
dashcore = { path = "../dash", features = ["serde", "core-block-hash-use-x11", "message_verification", "bls", "quorum_validation"] }
13+
dashcore = { path = "../dash", features = ["serde", "core-block-hash-use-x11", "message_verification", "bls", "quorum_validation", "tokio"] }
1414
dashcore_hashes = { path = "../hashes" }
1515
dash-network-seeds = { path = "../dash-network-seeds" }
1616
key-wallet = { path = "../key-wallet" }
@@ -22,6 +22,7 @@ clap = { version = "4.0", features = ["derive", "env"] }
2222

2323
# Async runtime
2424
tokio = { version = "1.0", features = ["full"] }
25+
parking_lot = "0.12"
2526
tokio-util = "0.7"
2627
tokio-stream = { version = "0.1", features = ["sync"] }
2728
async-trait = "0.1"

‎dash-spv/examples/filter_sync.rs‎

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,5 @@
11
//! BIP157 filter synchronization example.
22
3-
use dash_spv::network::PeerNetworkManager;
43
use dash_spv::storage::DiskStorageManager;
54
use dash_spv::{init_console_logging, ClientConfig, DashSpvClient, LevelFilter};
65
use dashcore::Address;
@@ -25,18 +24,15 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
2524
.with_storage_path("./.tmp/filter-sync-example-storage")
2625
.without_masternodes(); // Skip masternode sync for this example
2726

28-
// Create network manager
29-
let network_manager = PeerNetworkManager::new(&config).await?;
30-
3127
// Create storage manager
3228
let storage_manager = DiskStorageManager::new(&config).await?;
3329

3430
// Create wallet manager
3531
let wallet = Arc::new(RwLock::new(WalletManager::<ManagedWalletInfo>::new(config.network)));
3632

3733
// Create the client
38-
let client =
39-
DashSpvClient::new(config, network_manager, storage_manager, wallet, vec![]).await?;
34+
let network = dash_spv::network::PeerNetworkManager::new(&config).await;
35+
let client = DashSpvClient::new(config, network, storage_manager, wallet, vec![]).await?;
4036

4137
println!("Starting synchronization with filter support...");
4238
println!("Watching address: {:?}", watch_address);

‎dash-spv/examples/simple_sync.rs‎

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,5 @@
11
//! Simple header synchronization example.
22
3-
use dash_spv::network::PeerNetworkManager;
43
use dash_spv::storage::DiskStorageManager;
54
use dash_spv::{init_console_logging, ClientConfig, DashSpvClient, LevelFilter};
65
use key_wallet::wallet::managed_wallet_info::ManagedWalletInfo;
@@ -20,18 +19,15 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
2019
.without_filters() // Skip filter sync for this example
2120
.without_masternodes(); // Skip masternode sync for this example
2221

23-
// Create network manager
24-
let network_manager = PeerNetworkManager::new(&config).await?;
25-
2622
// Create storage manager
2723
let storage_manager = DiskStorageManager::new(&config).await?;
2824

2925
// Create wallet manager
3026
let wallet = Arc::new(RwLock::new(WalletManager::<ManagedWalletInfo>::new(config.network)));
3127

3228
// Create the client
33-
let client =
34-
DashSpvClient::new(config, network_manager, storage_manager, wallet, vec![]).await?;
29+
let network = dash_spv::network::PeerNetworkManager::new(&config).await;
30+
let client = DashSpvClient::new(config, network, storage_manager, wallet, vec![]).await?;
3531

3632
println!("Starting header synchronization...");
3733

‎dash-spv/examples/spv_with_wallet.rs‎

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,6 @@
22
//!
33
//! This example shows how to integrate the SPV client with a wallet manager.
44
5-
use dash_spv::network::PeerNetworkManager;
65
use dash_spv::storage::DiskStorageManager;
76
use dash_spv::{ClientConfig, DashSpvClient, LevelFilter};
87
use key_wallet::wallet::managed_wallet_info::ManagedWalletInfo;
@@ -20,18 +19,15 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
2019
.with_storage_path("./.tmp/spv-with-wallet-example-storage")
2120
.with_validation_mode(dash_spv::ValidationMode::Full);
2221

23-
// Create network manager
24-
let network_manager = PeerNetworkManager::new(&config).await?;
25-
2622
// Create storage manager - use disk storage for persistence
2723
let storage_manager = DiskStorageManager::new(&config).await?;
2824

2925
// Create wallet manager
3026
let wallet = Arc::new(RwLock::new(WalletManager::<ManagedWalletInfo>::new(config.network)));
3127

3228
// Create the SPV client with all components
33-
let client =
34-
DashSpvClient::new(config, network_manager, storage_manager, wallet, vec![]).await?;
29+
let network = dash_spv::network::PeerNetworkManager::new(&config).await;
30+
let client = DashSpvClient::new(config, network, storage_manager, wallet, vec![]).await?;
3531

3632
// The wallet will automatically be notified of:
3733
// - New blocks via process_block()

‎dash-spv/src/client/core.rs‎

Lines changed: 12 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -35,8 +35,8 @@ pub(super) type PersistentSyncCoordinator<W> = SyncCoordinator<
3535
///
3636
/// # Generic Design Philosophy
3737
///
38-
/// This struct uses three generic parameters (`W`, `N`, `S`) instead of concrete types or
39-
/// trait objects. This design choice provides significant benefits for a library:
38+
/// This struct uses two generic parameters (`W`, `S`) instead of concrete types or
39+
/// trait objects. This design choice provides significant benefits for a library.
4040
///
4141
/// ## Benefits of Generic Architecture
4242
///
@@ -57,8 +57,7 @@ pub(super) type PersistentSyncCoordinator<W> = SyncCoordinator<
5757
/// - Essential for a reusable library
5858
///
5959
/// ### 4. **Testing Without Mocks** 🧪
60-
/// - Test implementations (`MockNetworkManager`) are
61-
/// first-class types, not runtime injections
60+
/// - Test implementations are first-class types, not runtime injections
6261
/// - No conditional compilation or feature flags needed for tests
6362
/// - Type system ensures test and production code are compatible
6463
///
@@ -70,8 +69,8 @@ pub(super) type PersistentSyncCoordinator<W> = SyncCoordinator<
7069
/// ## Type Parameters
7170
///
7271
/// - `W: WalletInterface` - Handles UTXO tracking, address management, transaction processing
73-
/// - `N: NetworkManager` - Manages peer connections, message routing, network protocol
7472
/// - `S: StorageManager` - Persistent storage for headers, filters, chain state
73+
/// - Networking is the concrete `crate::network::PeerNetworkManager`
7574
/// - Event handlers are stored as `Vec<Arc<dyn EventHandler>>`
7675
///
7776
/// ## Common Configurations
@@ -82,14 +81,6 @@ pub(super) type PersistentSyncCoordinator<W> = SyncCoordinator<
8281
/// // Production configuration
8382
/// type StandardSpvClient = DashSpvClient<
8483
/// WalletManager,
85-
/// PeerNetworkManager,
86-
/// DiskStorageManager,
87-
/// >;
88-
///
89-
/// // Test configuration
90-
/// type TestSpvClient = DashSpvClient<
91-
/// WalletManager,
92-
/// MockNetworkManager,
9384
/// DiskStorageManager,
9485
/// >;
9586
/// ```
@@ -105,7 +96,7 @@ pub(super) type PersistentSyncCoordinator<W> = SyncCoordinator<
10596
/// The generic design is an intentional, beneficial architectural choice for a library.
10697
pub struct DashSpvClient<W: WalletInterface, N: NetworkManager, S: StorageManager> {
10798
pub(super) config: Arc<RwLock<ClientConfig>>,
108-
pub(super) network: Arc<Mutex<N>>,
99+
pub(super) network: Arc<N>,
109100
pub(super) storage: Arc<Mutex<S>>,
110101
/// External wallet implementation (required)
111102
pub(super) wallet: Arc<RwLock<W>>,
@@ -114,6 +105,12 @@ pub struct DashSpvClient<W: WalletInterface, N: NetworkManager, S: StorageManage
114105
/// `true` while running, `false` once a stop is requested. Stored as a
115106
/// `watch` so a stop is observed immediately rather than polled.
116107
pub(super) running: Arc<watch::Sender<bool>>,
108+
/// Set by `stop()` before it flips `running` false; read by `start()` under
109+
/// the `running` watch lock so a stop that races startup is not lost. Without
110+
/// it, a `stop()` arriving while `start()` is still connecting would no-op
111+
/// (running not yet true), then `start()` would flip running true and the run
112+
/// task would sync forever — hanging any `run_handle.await`.
113+
pub(super) stop_requested: Arc<std::sync::atomic::AtomicBool>,
117114
pub(super) event_handlers: Arc<Vec<Arc<dyn super::EventHandler>>>,
118115
}
119116

@@ -127,6 +124,7 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> Clone for DashSpv
127124
masternode_engine: self.masternode_engine.clone(),
128125
sync_coordinator: Arc::clone(&self.sync_coordinator),
129126
running: Arc::clone(&self.running),
127+
stop_requested: Arc::clone(&self.stop_requested),
130128
event_handlers: Arc::clone(&self.event_handlers),
131129
}
132130
}

‎dash-spv/src/client/event_handler.rs‎

Lines changed: 4 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -326,9 +326,8 @@ mod tests {
326326
};
327327
handler.on_sync_event(&event);
328328
handler.on_network_event(&NetworkEvent::PeersUpdated {
329-
connected_count: 0,
330-
addresses: vec![],
331-
best_height: None,
329+
connected_count: 1,
330+
best_height: 100,
332331
});
333332
handler.on_progress(&SyncProgress::default());
334333
handler.on_error("test error");
@@ -533,14 +532,8 @@ mod tests {
533532
);
534533

535534
let addr: SocketAddr = "127.0.0.1:9999".parse().unwrap();
536-
tx.send(NetworkEvent::PeerConnected {
537-
address: addr,
538-
})
539-
.unwrap();
540-
tx.send(NetworkEvent::PeerDisconnected {
541-
address: addr,
542-
})
543-
.unwrap();
535+
tx.send(NetworkEvent::PeerConnected(addr)).unwrap();
536+
tx.send(NetworkEvent::PeerDisconnected(addr)).unwrap();
544537
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
545538

546539
shutdown.cancel();

‎dash-spv/src/client/events.rs‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,10 @@
44
//! - Event receiver management
55
//! - Event emission
66
7+
use crate::network::NetworkManager;
78
use tokio::sync::watch;
89

9-
use crate::network::{NetworkEvent, NetworkManager};
10+
use crate::network::NetworkEvent;
1011
use crate::storage::StorageManager;
1112
use crate::sync::{SyncEvent, SyncProgress};
1213
use key_wallet_manager::WalletInterface;
@@ -32,6 +33,6 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
3233

3334
/// Subscribe to network events.
3435
pub(crate) async fn subscribe_network_events(&self) -> broadcast::Receiver<NetworkEvent> {
35-
self.network.lock().await.subscribe_network_events()
36+
self.network.events()
3637
}
3738
}

0 commit comments

Comments
 (0)