Skip to content
Merged
Show file tree
Hide file tree
Changes from 23 commits
Commits
Show all changes
25 commits
Select commit Hold shift + click to select a range
9375a66
fix(spv): exit manager task loop on network errors and signal shutdow…
lklimek Feb 16, 2026
d8bc066
fix(dash-spv): shutdown token not checked when waiting for peer conne…
lklimek Feb 16, 2026
cdccab6
fix(spv): address review findings for network error handling
lklimek Feb 16, 2026
0a6dff1
refactor(spv): replace signal_shutdown with full shutdown in stop()
lklimek Feb 16, 2026
9f3970c
fix(spv): resolve deadlock in PeerNetworkManager shutdown
lklimek Feb 16, 2026
86ebbab
doc: document logging
lklimek Feb 16, 2026
85cf2bd
feat(spv): surface FatalNetwork errors to API consumers
lklimek Feb 16, 2026
8541836
fix(spv): self-recover on FatalNetwork instead of exiting manager loop
lklimek Feb 16, 2026
cb84ccc
fix(spv): remove take() in PeerNetworkManager::shutdown to close race
lklimek Feb 16, 2026
f4a52dd
fix(spv): add missing shutdown checks in maintenance loop
lklimek Feb 16, 2026
7356e2b
fix(spv): address PR #440 audit findings
lklimek Feb 16, 2026
f48a3a5
fmt
lklimek Feb 16, 2026
04a87a7
Merge remote-tracking branch 'origin/v0.42-dev' into fix/block-header…
lklimek Feb 17, 2026
8e22175
refactor(spv): extract network error cooldown into named constant
lklimek Feb 17, 2026
8872062
fmt
lklimek Feb 17, 2026
c3fa668
fix(spv): restore shutdown checks dropped during upstream refactor
lklimek Feb 17, 2026
ddc539f
fmt
lklimek Feb 17, 2026
531eb11
revert: CLAUDE.md logging info
lklimek Feb 17, 2026
6affdaa
revert logging changes
lklimek Feb 17, 2026
ebdfbe3
Update dash-spv/src/network/manager.rs
lklimek Feb 17, 2026
e6c01e8
chore: fmt
lklimek Feb 17, 2026
d3d0305
fix: don't set SyncState::WaitingForConnections) on network error
lklimek Feb 17, 2026
b239c4e
refactor: remove cooldown
lklimek Feb 17, 2026
9501aed
Merge branch 'v0.42-dev' into fix/block-header-error-retry-loop
lklimek Feb 18, 2026
8f7eacc
refactor: strip non-shutdown changes from PR 440
lklimek Feb 18, 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
6 changes: 6 additions & 0 deletions dash-spv/src/client/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,12 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,
}
}

// 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.shutdown().await {
log::warn!("Error shutting down sync coordinator: {}", e);
}

// Disconnect from network
self.network.disconnect().await?;

Expand Down
2 changes: 1 addition & 1 deletion dash-spv/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,7 @@ pub enum SyncError {
#[error("Timeout error: {0}")]
Timeout(String),

/// Network-related errors (e.g., connection failures, protocol errors)
/// Network-related errors (e.g., connection failures, protocol errors).
#[error("Network error: {0}")]
Network(String),

Expand Down
32 changes: 27 additions & 5 deletions dash-spv/src/network/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -241,12 +241,27 @@ impl PeerNetworkManager {
let message_dispatcher = self.message_dispatcher.clone();
let network_event_sender = self.network_event_sender.clone();

// Spawn connection task
let mut tasks = self.tasks.lock().await;
// Spawn connection task — use select to avoid blocking on the lock during shutdown
let mut tasks = tokio::select! {
guard = self.tasks.lock() => guard,
_ = self.shutdown_token.cancelled() => {
self.pool.remove_peer(&addr).await;
return;
}
};
tasks.spawn(async move {
log::debug!("Attempting to connect to {}", addr);

match Peer::connect(addr, CONNECTION_TIMEOUT.as_secs(), network).await {
let connect_result = tokio::select! {
result = Peer::connect(addr, CONNECTION_TIMEOUT.as_secs(), network) => result,
_ = shutdown_token.cancelled() => {
log::debug!("Connection to {} cancelled by shutdown", addr);
pool.remove_peer(&addr).await;
return;
}
};

match connect_result {
Ok(mut peer) => {
// Perform handshake
let mut handshake_manager =
Expand Down Expand Up @@ -800,6 +815,10 @@ impl PeerNetworkManager {
}
}

if self.shutdown_token.is_cancelled() {
return;
}

// Send ping to all peers if needed
for (addr, peer) in self.pool.get_all_peers().await {
let mut peer_guard = peer.write().await;
Expand Down Expand Up @@ -863,7 +882,8 @@ impl PeerNetworkManager {
let mut tasks = self.tasks.lock().await;
tasks.spawn(async move {
// Periodic DNS discovery check (only active in non-exclusive mode)
let mut dns_interval = time::interval_at(Instant::now() + DNS_DISCOVERY_DELAY, DNS_DISCOVERY_DELAY);
let mut dns_interval =
time::interval_at(Instant::now() + DNS_DISCOVERY_DELAY, DNS_DISCOVERY_DELAY);
// Periodic reconnection check (active in both modes)
let mut maintenance_interval = time::interval(MAINTENANCE_INTERVAL);

Expand Down Expand Up @@ -1235,7 +1255,9 @@ impl PeerNetworkManager {
log::warn!("Failed to save reputation data on shutdown: {}", e);
}

// Wait for tasks to complete
// Drain tasks while holding the lock. connect_to_peer() already uses
// `select!` with the cancellation token when acquiring this lock, so no
// deadlock can occur once the shutdown token is cancelled above.
let mut tasks = self.tasks.lock().await;
while let Some(result) = tasks.join_next().await {
if let Err(e) = result {
Expand Down
2 changes: 1 addition & 1 deletion dash-spv/src/sync/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ pub enum SyncEvent {
height: u32,
},

/// A manager encountered a recoverable error.
/// A manager encountered an error during sync.
///
/// Emitted by: Any manager
/// Consumed by: Coordinator (for logging/monitoring)
Expand Down
168 changes: 167 additions & 1 deletion dash-spv/src/sync/sync_manager.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use crate::error::SyncResult;
use crate::error::{SyncError, SyncResult};
use crate::network::{Message, MessageType, NetworkEvent, RequestSender};
use crate::sync::{
BlockHeadersProgress, BlocksProgress, ChainLockProgress, FilterHeadersProgress,
Expand Down Expand Up @@ -194,6 +194,26 @@ pub trait SyncManager: Send + Sync + std::fmt::Debug {
}
}

/// Log a network error and emit a `ManagerError` event.
///
/// State is intentionally left unchanged: the `PeersUpdated` event path
/// handles the transition to `WaitingForConnections` when all peers are lost.
fn recover_from_network_error(
&mut self,
context: &SyncManagerTaskContext,
source: &str,
msg: &str,
) {
let identifier = self.identifier();
log::warn!("{} {} network error: {}", identifier, source, msg);
let progress = self.progress();
context.progress_sender.send(progress).ok();
context.emit_sync_event(SyncEvent::ManagerError {
manager: identifier,
error: format!("Network error ({}): {}", source, msg),
});
}

/// Run the manager task, processing messages, events, and periodic ticks.
///
/// This consumes the manager and runs until shutdown is signaled.
Expand Down Expand Up @@ -235,6 +255,9 @@ pub trait SyncManager: Send + Sync + std::fmt::Debug {
}
self.try_emit_progress(progress_before, &context.progress_sender);
}
Err(SyncError::Network(ref msg)) => {
self.recover_from_network_error(&context, "message handler", msg);
}
Err(e) => {
tracing::error!("{} error handling message: {}", identifier, e);
let error_event = SyncEvent::ManagerError {
Expand All @@ -261,6 +284,9 @@ pub trait SyncManager: Send + Sync + std::fmt::Debug {
}
self.try_emit_progress(progress_before, &context.progress_sender);
}
Err(SyncError::Network(ref msg)) => {
self.recover_from_network_error(&context, "sync event handler", msg);
}
Err(e) => {
tracing::error!("{} error handling event: {}", identifier, e);
}
Expand Down Expand Up @@ -288,6 +314,9 @@ pub trait SyncManager: Send + Sync + std::fmt::Debug {
}
self.try_emit_progress(progress_before, &context.progress_sender);
}
Err(SyncError::Network(ref msg)) => {
self.recover_from_network_error(&context, "network event handler", msg);
}
Err(e) => {
tracing::error!("{} error handling network event: {}", identifier, e);
}
Expand All @@ -309,6 +338,9 @@ pub trait SyncManager: Send + Sync + std::fmt::Debug {
}
self.try_emit_progress(progress_before, &context.progress_sender);
}
Err(SyncError::Network(ref msg)) => {
self.recover_from_network_error(&context, "tick", msg);
}
Err(e) => {
tracing::error!("{} tick error: {}", identifier, e);
}
Expand Down Expand Up @@ -447,4 +479,138 @@ mod tests {
// Verify tick was called multiple times
assert!(tick_count.load(Ordering::Relaxed) > 0);
}

/// Mock manager whose tick() returns SyncError::Network after a threshold.
struct NetworkErrorManager {
identifier: ManagerIdentifier,
state: SyncState,
tick_count: Arc<AtomicU32>,
error_after: u32,
}

impl std::fmt::Debug for NetworkErrorManager {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("NetworkErrorManager").field("identifier", &self.identifier).finish()
}
}

#[async_trait]
impl SyncManager for NetworkErrorManager {
fn identifier(&self) -> ManagerIdentifier {
self.identifier
}
fn state(&self) -> SyncState {
self.state
}
fn set_state(&mut self, state: SyncState) {
self.state = state;
}
fn wanted_message_types(&self) -> &'static [MessageType] {
&[]
}
fn progress(&self) -> SyncManagerProgress {
let mut progress = BlockHeadersProgress::default();
progress.set_state(self.state);
SyncManagerProgress::BlockHeaders(progress)
}

async fn handle_message(
&mut self,
_msg: Message,
_requests: &RequestSender,
) -> SyncResult<Vec<SyncEvent>> {
Ok(vec![])
}

async fn handle_sync_event(
&mut self,
_event: &SyncEvent,
_requests: &RequestSender,
) -> SyncResult<Vec<SyncEvent>> {
Ok(vec![])
}

async fn tick(&mut self, _requests: &RequestSender) -> SyncResult<Vec<SyncEvent>> {
let count = self.tick_count.fetch_add(1, Ordering::Relaxed);
if count >= self.error_after {
Err(SyncError::Network("channel closed".into()))
} else {
Ok(vec![])
}
}
}

/// Given a manager whose tick() returns SyncError::Network after a few calls,
/// When the task loop processes the error,
/// Then it stays in its current state and keeps running.
#[tokio::test]
async fn test_manager_resets_on_fatal_network_error() {
let tick_count = Arc::new(AtomicU32::new(0));

let manager = NetworkErrorManager {
identifier: ManagerIdentifier::BlockHeader,
state: SyncState::Initializing,
tick_count: tick_count.clone(),
error_after: 3,
};

// Create channels
let (_, message_receiver) = mpsc::unbounded_channel();
let sync_event_sender = broadcast::Sender::<SyncEvent>::new(100);
let network_event_sender = broadcast::Sender::<NetworkEvent>::new(100);
let (req_tx, _req_rx) = mpsc::unbounded_channel::<NetworkRequest>();
let requests = RequestSender::new(req_tx);
let shutdown = CancellationToken::new();
let (progress_sender, progress_rx) = watch::channel(manager.progress());

let context = SyncManagerTaskContext {
message_receiver,
sync_event_sender: sync_event_sender.clone(),
network_event_receiver: network_event_sender.subscribe(),
requests,
shutdown: shutdown.clone(),
progress_sender,
};

// Subscribe to sync events to verify ManagerError is emitted
let mut event_rx = sync_event_sender.subscribe();

// Spawn the task — it should keep running after the Network error
let handle = tokio::spawn(async move { manager.run(context).await });

// Wait long enough for the error to fire (tick 3 at ~300ms) plus a
// few more ticks so we can verify the loop keeps running.
tokio::time::sleep(Duration::from_millis(600)).await;

// State stays at WaitingForConnections (set by initialize(), not changed by error recovery)
assert_eq!(progress_rx.borrow().state(), SyncState::WaitingForConnections);

// Verify tick was called more than the error threshold (manager kept
// running after the Network error).
assert!(
tick_count.load(Ordering::Relaxed) > 4,
"manager should keep ticking after Network error"
);

// Verify ManagerError event was emitted
let mut found_error = false;
while let Ok(event) = event_rx.try_recv() {
if matches!(event, SyncEvent::ManagerError { .. }) {
found_error = true;
break;
}
}
assert!(found_error, "ManagerError event should have been emitted");

// Shut down the manager via the shutdown token
shutdown.cancel();

let result = tokio::time::timeout(Duration::from_secs(2), handle)
.await
.expect("task should exit after shutdown signal")
.unwrap();

assert!(result.is_ok());
assert_eq!(result.unwrap(), ManagerIdentifier::BlockHeader);
}
}
Loading