From d0193e3cfcfc583a51f2eec1f7f6d42a98d1e380 Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Sat, 3 Oct 2026 02:13:33 +0000 Subject: [PATCH 1/2] fix(dash-spv): stop a client that is still starting `run()` blocked until the client stopped, so every caller spawned it in a task of its own. A `stop()` that came in before that task finished starting found the client not running yet, returned without doing anything, and `run()` then started and kept syncing. Through the FFI, `dash_spv_ffi_client_stop` right after `dash_spv_ffi_client_run` aborted the run task after a 5 second timeout and left the managers, the network and the storage running. `run()` now starts the client and returns: it starts the sync managers, the network and the storage, spawns the sync loop and keeps its handle. `stop()` cancels the loop, waits for it and stops the rest. Both hold a lock for the whole call, so a `stop()` during a start waits for it and then stops the client. The loop handle replaces the `running` watch: `is_running()` checks whether there is one and is now async. A second `run()` on a running client returns `Ok`, as `stop()` already does on a stopped one. When the sync loop fails, it reports the error through `on_error` and stops the client. The storage worker now starts last, so a failed start leaves nothing running. Callers no longer spawn `run()`: the FFI drops its run task and returns the startup error from `dash_spv_ffi_client_run`, and the binary, the examples and the bench wait for Ctrl-C or their own condition before calling `stop()`. The restart tests now stop and run the same client several times in a row, and clear its storage between runs. `wallet_integration_test.rs` is removed: it only checked that a client can be created, that an empty wallet manager is empty, and that the running flag flips without peers. Co-Authored-By: Claude Opus 5.5 --- dash-spv-bench/src/main.rs | 9 +- dash-spv-ffi/FFI_API.md | 2 +- dash-spv-ffi/src/bin/ffi_cli.rs | 2 +- dash-spv-ffi/src/client.rs | 78 ++-------- .../tests/unit/test_client_lifecycle.rs | 51 +------ dash-spv/examples/filter_sync.rs | 4 + dash-spv/examples/simple_sync.rs | 4 + dash-spv/examples/spv_with_wallet.rs | 4 + dash-spv/src/client/core.rs | 25 +++- dash-spv/src/client/lifecycle.rs | 51 +++---- dash-spv/src/client/mod.rs | 2 +- dash-spv/src/client/sync_coordinator.rs | 138 +++++++++--------- dash-spv/src/lib.rs | 4 + dash-spv/src/main.rs | 14 +- dash-spv/tests/dashd_masternode/setup.rs | 19 +-- dash-spv/tests/dashd_masternode/tests_sync.rs | 2 +- dash-spv/tests/dashd_sync/setup.rs | 20 +-- dash-spv/tests/dashd_sync/tests_restart.rs | 37 ++--- dash-spv/tests/peer_test.rs | 12 +- dash-spv/tests/wallet_integration_test.rs | 93 ------------ 20 files changed, 189 insertions(+), 382 deletions(-) delete mode 100644 dash-spv/tests/wallet_integration_test.rs diff --git a/dash-spv-bench/src/main.rs b/dash-spv-bench/src/main.rs index 90f262aaf..705b80655 100644 --- a/dash-spv-bench/src/main.rs +++ b/dash-spv-bench/src/main.rs @@ -184,12 +184,7 @@ async fn main() -> Result<()> { .await .map_err(|e| anyhow!("client new: {e}"))?; - let run_client = client.clone(); - let run_handle = tokio::spawn(async move { - if let Err(e) = run_client.run().await { - tracing::error!("client run() exited with error: {e}"); - } - }); + client.run().await.map_err(|e| anyhow!("client run: {e}"))?; tokio::select! { _ = handler.wait_done() => {} @@ -199,8 +194,6 @@ async fn main() -> Result<()> { let m = handler.snapshot(); let _ = client.stop().await; - run_handle.abort(); - let _ = run_handle.await; let peak_rss_kb = proc_status_kb("VmHWM:"); #[cfg(target_os = "linux")] diff --git a/dash-spv-ffi/FFI_API.md b/dash-spv-ffi/FFI_API.md index e644bd59d..bd7b3bb77 100644 --- a/dash-spv-ffi/FFI_API.md +++ b/dash-spv-ffi/FFI_API.md @@ -679,7 +679,7 @@ dash_spv_ffi_client_run(client: *mut FFIDashSpvClient) -> i32 ``` **Description:** -Start the SPV client and begin syncing in the background. Uses the event callbacks provided at client creation time. Returns immediately after spawning the sync task. # Safety - `client` must be a valid, non-null pointer to a created client. # Returns 0 on success, error code on failure. +Start the SPV client and begin syncing in the background. Uses the event callbacks provided at client creation time. Returns once the storage, the sync managers and the network are started. Starting can take a few seconds, e.g. when peers have to be discovered through DNS. If blocking the calling thread that long is a problem (such as a UI thread), call this from another thread. # Safety - `client` must be a valid, non-null pointer to a created client. # Returns 0 on success, error code on failure. **Safety:** - `client` must be a valid, non-null pointer to a created client. diff --git a/dash-spv-ffi/src/bin/ffi_cli.rs b/dash-spv-ffi/src/bin/ffi_cli.rs index a9715bbf5..26cfde573 100644 --- a/dash-spv-ffi/src/bin/ffi_cli.rs +++ b/dash-spv-ffi/src/bin/ffi_cli.rs @@ -640,7 +640,7 @@ fn main() { println!("Event and progress callbacks configured, starting sync..."); - // Run client - starts sync in background and returns immediately + // Run client - starts sync in background and returns once started let rc = dash_spv_ffi_client_run(client); if rc != FFIErrorCode::Success as i32 { eprintln!("Client run failed: {}", ffi_string_to_rust(dash_spv_ffi_get_last_error())); diff --git a/dash-spv-ffi/src/client.rs b/dash-spv-ffi/src/client.rs index 91e96659c..050eedbba 100644 --- a/dash-spv-ffi/src/client.rs +++ b/dash-spv-ffi/src/client.rs @@ -10,9 +10,8 @@ use dash_spv::DashSpvClient; use tracing::dispatcher::{get_default, set_default}; use std::mem::forget; -use std::sync::{Arc, Mutex}; +use std::sync::Arc; use tokio::runtime::Runtime; -use tokio::task::JoinHandle; /// FFI wrapper around `DashSpvClient`. type InnerClient = DashSpvClient< @@ -24,7 +23,6 @@ type InnerClient = DashSpvClient< pub struct FFIDashSpvClient { pub(crate) inner: InnerClient, pub(crate) runtime: Arc, - run_task: Mutex>>, } impl FFIDashSpvClient { @@ -110,7 +108,6 @@ pub unsafe extern "C" fn dash_spv_ffi_client_new( let ffi_client = FFIDashSpvClient { inner: client, runtime, - run_task: Mutex::new(None), }; Box::into_raw(Box::new(ffi_client)) } @@ -121,38 +118,6 @@ pub unsafe extern "C" fn dash_spv_ffi_client_new( } } -/// Maximum time to wait for the run task to exit cooperatively before aborting. -const RUN_TASK_SHUTDOWN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); - -impl FFIDashSpvClient { - /// Wait for the run task to finish cooperatively, aborting only on timeout. - /// - /// `DashSpvClient::stop()` must have been called first (it flips the client's - /// internal running state, which makes `run()` exit its loop and clean up - /// monitor tasks). This only falls back to `abort()` if the task doesn't - /// exit within the timeout. - fn wait_for_run_task(&self) { - let task = self.run_task.lock().unwrap().take(); - if let Some(mut task) = task { - let finished = self.runtime.block_on(async { - tokio::time::timeout(RUN_TASK_SHUTDOWN_TIMEOUT, &mut task).await - }); - match finished { - Ok(Ok(())) => {} - Ok(Err(e)) => tracing::warn!("Run task exited with join error: {}", e), - Err(_) => { - tracing::warn!( - "Run task did not exit within {:?}, aborting", - RUN_TASK_SHUTDOWN_TIMEOUT, - ); - task.abort(); - let _ = self.runtime.block_on(task); - } - } - } - } -} - /// Update the running client's configuration. /// /// # Safety @@ -191,10 +156,7 @@ pub unsafe extern "C" fn dash_spv_ffi_client_stop(client: *mut FFIDashSpvClient) let client = &(*client); - // `stop()` flips the client's internal running state, making `run()` break - // out of its loop. Wait for the spawned run task only after that. let result = client.runtime.block_on(async { client.inner.stop().await }); - client.wait_for_run_task(); match result { Ok(()) => FFIErrorCode::Success as i32, @@ -207,8 +169,12 @@ pub unsafe extern "C" fn dash_spv_ffi_client_stop(client: *mut FFIDashSpvClient) /// Start the SPV client and begin syncing in the background. /// -/// Uses the event callbacks provided at client creation time. Returns -/// immediately after spawning the sync task. +/// Uses the event callbacks provided at client creation time. Returns once +/// the storage, the sync managers and the network are started. +/// +/// Starting can take a few seconds, e.g. when peers have to be discovered +/// through DNS. If blocking the calling thread that long is a problem (such as +/// a UI thread), call this from another thread. /// /// # Safety /// - `client` must be a valid, non-null pointer to a created client. @@ -221,25 +187,15 @@ pub unsafe extern "C" fn dash_spv_ffi_client_run(client: *mut FFIDashSpvClient) let client = &(*client); - tracing::info!("dash_spv_ffi_client_run: starting sync"); - - let spv_client = client.inner.clone(); - - let task = client.runtime.spawn(async move { - tracing::debug!("Sync task: starting run"); + let result = client.runtime.block_on(client.inner.run()); - if let Err(e) = spv_client.run().await { - tracing::error!("Sync task: error: {}", e); + match result { + Ok(()) => FFIErrorCode::Success as i32, + Err(e) => { + set_last_error(&e.to_string()); + FFIErrorCode::from(e) as i32 } - - tracing::debug!("Sync task: exiting"); - }); - - *client.run_task.lock().unwrap() = Some(task); - - tracing::info!("dash_spv_ffi_client_run: background task spawned, returning"); - - FFIErrorCode::Success as i32 + } } /// Get the current sync progress snapshot. @@ -270,7 +226,6 @@ pub unsafe extern "C" fn dash_spv_ffi_client_clear_storage(client: *mut FFIDashS let client = &(*client); let result = client.runtime.block_on(client.inner.clear_storage()); - client.wait_for_run_task(); match result { Ok(_) => FFIErrorCode::Success as i32, @@ -414,15 +369,10 @@ pub unsafe extern "C" fn dash_spv_ffi_client_destroy(client: *mut FFIDashSpvClie if !client.is_null() { let client = Box::from_raw(client); - // Stop the SPV client (run() calls stop() internally, but this - // handles the case where run() was never called or was aborted). client.runtime.block_on(async { let _ = client.inner.stop().await; }); - // Wait for the run task to finish (cooperative, with timeout fallback) - client.wait_for_run_task(); - tracing::info!("FFI client destroyed and all tasks cleaned up"); } } diff --git a/dash-spv-ffi/tests/unit/test_client_lifecycle.rs b/dash-spv-ffi/tests/unit/test_client_lifecycle.rs index bf50abeb7..43c58b8a2 100644 --- a/dash-spv-ffi/tests/unit/test_client_lifecycle.rs +++ b/dash-spv-ffi/tests/unit/test_client_lifecycle.rs @@ -9,7 +9,6 @@ mod tests { use dash_network::ffi::FFINetwork; use serial_test::serial; use std::ffi::CString; - use std::sync::mpsc; use std::sync::{Arc as StdArc, Mutex as StdMutex}; use std::thread; use std::time::Duration; @@ -213,57 +212,17 @@ mod tests { #[test] #[serial] - fn test_client_error_callback_fires_on_start_failure() { - let (tx, rx) = mpsc::channel::(); - let tx_ptr = Box::into_raw(Box::new(tx)); - - extern "C" fn on_error( - error: *const std::os::raw::c_char, - user_data: *mut std::os::raw::c_void, - ) { - let tx = unsafe { &*(user_data as *const mpsc::Sender) }; - let error_str = unsafe { std::ffi::CStr::from_ptr(error) }.to_str().unwrap().to_owned(); - let _ = tx.send(error_str); - } - + fn test_client_run_twice_succeeds() { unsafe { let (config, _temp_dir) = create_test_config_with_dir(); - let callbacks = FFIEventCallbacks { - error: FFIClientErrorCallback { - on_error: Some(on_error), - user_data: tx_ptr as *mut std::os::raw::c_void, - }, - ..FFIEventCallbacks::default() - }; - let client = dash_spv_ffi_client_new(config, callbacks); + let client = dash_spv_ffi_client_new(config, FFIEventCallbacks::default()); assert!(!client.is_null()); - // Call run() twice — the second run's sync thread will call - // start() on the already-running client, triggering "already running" - let run_result = dash_spv_ffi_client_run(client); - assert_eq!(run_result, FFIErrorCode::Success as i32); - - // Brief wait for the first run's sync thread to complete start() - thread::sleep(Duration::from_millis(200)); - - let _run_result2 = dash_spv_ffi_client_run(client); - - // Wait for the error callback to fire (with timeout) - let error_msg = rx - .recv_timeout(Duration::from_secs(5)) - .expect("Error callback should have been called on start failure"); - assert!( - error_msg.contains("already running"), - "Expected 'already running' error, got: {}", - error_msg - ); + // A second `run` on a running client does nothing. + assert_eq!(dash_spv_ffi_client_run(client), FFIErrorCode::Success as i32); + assert_eq!(dash_spv_ffi_client_run(client), FFIErrorCode::Success as i32); dash_spv_ffi_client_stop(client); - - // Free the sender only after stop has joined all threads, - // so no background thread can call on_error with a dangling user_data. - drop(Box::from_raw(tx_ptr)); - dash_spv_ffi_client_destroy(client); dash_spv_ffi_config_destroy(config); } diff --git a/dash-spv/examples/filter_sync.rs b/dash-spv/examples/filter_sync.rs index 203896ce6..dc875da94 100644 --- a/dash-spv/examples/filter_sync.rs +++ b/dash-spv/examples/filter_sync.rs @@ -43,6 +43,10 @@ async fn main() -> Result<(), Box> { client.run().await?; + // Sync in the background until Ctrl-C. + tokio::signal::ctrl_c().await?; + client.stop().await?; + println!("Done!"); Ok(()) } diff --git a/dash-spv/examples/simple_sync.rs b/dash-spv/examples/simple_sync.rs index 0568768fe..c6bab2dac 100644 --- a/dash-spv/examples/simple_sync.rs +++ b/dash-spv/examples/simple_sync.rs @@ -37,6 +37,10 @@ async fn main() -> Result<(), Box> { client.run().await?; + // Sync in the background until Ctrl-C. + tokio::signal::ctrl_c().await?; + client.stop().await?; + println!("Done!"); Ok(()) } diff --git a/dash-spv/examples/spv_with_wallet.rs b/dash-spv/examples/spv_with_wallet.rs index 2d2c661d8..7c73c9334 100644 --- a/dash-spv/examples/spv_with_wallet.rs +++ b/dash-spv/examples/spv_with_wallet.rs @@ -41,6 +41,10 @@ async fn main() -> Result<(), Box> { client.run().await?; + // Sync in the background until Ctrl-C. + tokio::signal::ctrl_c().await?; + client.stop().await?; + println!("Done!"); Ok(()) } diff --git a/dash-spv/src/client/core.rs b/dash-spv/src/client/core.rs index b4445e4a8..4e6e29fe4 100644 --- a/dash-spv/src/client/core.rs +++ b/dash-spv/src/client/core.rs @@ -10,7 +10,9 @@ use dashcore::sml::masternode_list_engine::MasternodeListEngine; use std::sync::Arc; -use tokio::sync::{watch, Mutex, RwLock}; +use tokio::sync::{Mutex, RwLock}; +use tokio::task::JoinHandle; +use tokio_util::sync::CancellationToken; use super::ClientConfig; use crate::error::{Result, SpvError}; @@ -111,12 +113,19 @@ pub struct DashSpvClient>, pub(super) masternode_engine: Option>>, pub(super) sync_coordinator: Arc>, - /// `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>, + /// The running sync loop, `None` while stopped. `run` and `stop` hold the + /// lock throughout, so they never overlap. + pub(super) sync_loop: Arc>>, pub(super) event_handlers: Arc>>, } +/// The background task of a running client. +pub(super) struct SyncLoop { + pub(super) task: JoinHandle<()>, + /// Stops the loop and the event monitors it runs. + pub(super) shutdown: CancellationToken, +} + impl Clone for DashSpvClient { fn clone(&self) -> Self { Self { @@ -126,7 +135,7 @@ impl Clone for DashSpv wallet: Arc::clone(&self.wallet), masternode_engine: self.masternode_engine.clone(), sync_coordinator: Arc::clone(&self.sync_coordinator), - running: Arc::clone(&self.running), + sync_loop: Arc::clone(&self.sync_loop), event_handlers: Arc::clone(&self.event_handlers), } } @@ -152,9 +161,9 @@ impl DashSpvClient bool { - *self.running.borrow() + /// Check if the client is running. Waits for an ongoing `run` or `stop` to finish. + pub async fn is_running(&self) -> bool { + self.sync_loop.lock().await.is_some() } /// Returns the current chain tip hash if available. diff --git a/dash-spv/src/client/lifecycle.rs b/dash-spv/src/client/lifecycle.rs index 4304c268d..ed316c659 100644 --- a/dash-spv/src/client/lifecycle.rs +++ b/dash-spv/src/client/lifecycle.rs @@ -2,12 +2,12 @@ //! //! This module contains: //! - Constructor (`new`) -//! - Startup logic (`start`) //! - Shutdown logic (`stop`) //! - Sync initiation (`start_sync`) //! - Genesis block initialization //! - Wallet data loading +use super::core::SyncLoop; use super::core::SyncManagers; use super::{ClientConfig, DashSpvClient, EventHandler}; use crate::chain::checkpoints::CheckpointManager; @@ -27,7 +27,7 @@ use dashcore::TxMerkleNode; use dashcore_hashes::Hash; use key_wallet_manager::WalletInterface; use std::sync::Arc; -use tokio::sync::{watch, Mutex, RwLock}; +use tokio::sync::{Mutex, RwLock}; impl DashSpvClient { /// Create a new SPV client with the given configuration, network, storage, and wallet. @@ -80,7 +80,7 @@ impl DashSpvClient DashSpvClient Result<()> { - if self.is_running() { - return Err(SpvError::Config("Client already running".to_string())); - } - - self.storage.lock().await.start().await; - if let Err(e) = self.start_sync().await { - self.storage.lock().await.stop().await; - return Err(e); - } - - // 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); - - Ok(()) - } - - /// Start the sync managers and the network on top of a started storage. - async fn start_sync(&self) -> Result<()> { + /// Start the sync managers, the network and the storage worker. + pub(super) async fn start_sync(&self) -> Result<()> { let managers = self.build_sync_managers().await?; // Start all sync tasks before connecting to the network to make sure initial connection @@ -245,21 +225,30 @@ impl DashSpvClient Result<()> { - // Check if already stopped - if !*self.running.borrow() { + let mut sync_loop = self.sync_loop.lock().await; + let Some(SyncLoop { + task, + shutdown, + }) = sync_loop.take() + else { return Ok(()); - } + }; - // Flip the running state before tearing anything down so a concurrent - // `run()` loop wakes immediately and breaks out before it can lock the + // Stop the sync loop before tearing anything down so it cannot lock the // sync coordinator again. This prevents a tick from racing against the // shutdown below. - self.running.send_replace(false); + shutdown.cancel(); + if let Err(e) = task.await { + tracing::warn!("Sync loop task failed: {}", e); + } // Shut down sync coordinator: signals cancellation and waits for manager // tasks to drain before we tear down the network and storage layers. diff --git a/dash-spv/src/client/mod.rs b/dash-spv/src/client/mod.rs index 6a7c8bda5..4e65af1ae 100644 --- a/dash-spv/src/client/mod.rs +++ b/dash-spv/src/client/mod.rs @@ -8,7 +8,7 @@ //! //! - `config.rs` - Client configuration //! - `core.rs` - Core `DashSpvClient` struct definition and simple accessors -//! - `lifecycle.rs` - Client lifecycle (new, start, stop, shutdown) +//! - `lifecycle.rs` - Client lifecycle (new, stop) //! - `events.rs` - Event emission and progress tracking receivers //! - `queries.rs` - Peer, masternode, and balance queries //! - `transactions.rs` - Transaction operations (e.g., broadcast) diff --git a/dash-spv/src/client/sync_coordinator.rs b/dash-spv/src/client/sync_coordinator.rs index 2addac68f..975677aca 100644 --- a/dash-spv/src/client/sync_coordinator.rs +++ b/dash-spv/src/client/sync_coordinator.rs @@ -5,6 +5,7 @@ use std::time::Duration; use tokio::sync::mpsc; use tokio_util::sync::CancellationToken; +use super::core::SyncLoop; use super::event_handler::{ spawn_broadcast_monitor, spawn_chainlock_wallet_dispatch, spawn_progress_monitor, spawn_reservation_sweep, @@ -19,21 +20,27 @@ use key_wallet_manager::WalletInterface; const SYNC_COORDINATOR_TICK_MS: Duration = Duration::from_millis(100); impl DashSpvClient { - /// Start the client and run the sync loop until `stop()` is called. + /// Start the client and run the sync loop in the background until `stop()` is called. /// /// Subscribes to all event channels internally and dispatches events to the - /// event handler provided at construction. Calls `start()` internally, runs - /// continuous network monitoring, and calls `stop()` before returning. + /// event handler provided at construction. Starts the storage, the sync + /// managers and the network, and returns once the client is running. If the + /// sync loop fails, it reports the error through `on_error` and stops the client. + /// Does nothing if the client is already running. + /// + /// Starting can take a few seconds, e.g. when peers have to be discovered + /// through DNS. If blocking the caller that long is a problem, call `run` + /// from another task or thread. pub async fn run(&self) -> Result<()> { + let mut sync_loop = self.sync_loop.lock().await; + if sync_loop.is_some() { + return Ok(()); + } + let handlers = self.event_handlers.clone(); let monitor_shutdown = CancellationToken::new(); let (monitor_failure_tx, mut monitor_failure_rx) = mpsc::channel::(1); - // Subscribe before `start()` so a `stop()` that races startup is never - // missed: the receiver records the version at subscription time, so any - // later state change is observed even if it lands before the loop runs. - let mut stop_rx = self.running.subscribe(); - // Subscribe and spawn monitors before startup so we don't miss early // connection events. let sync_event_rx = self.subscribe_sync_events().await; @@ -83,7 +90,7 @@ impl DashSpvClient DashSpvClient = loop { - if !self.is_running() { - tracing::info!("Stopping network monitoring"); - break None; - } - - let error: Option = tokio::select! { - _ = sync_coordinator_tick_interval.tick() => { - self.sync_coordinator.lock().await.tick().await.err().map(Into::into) - } - _ = stop_rx.changed() => { - tracing::debug!("DashSpvClient run loop stop requested"); - break None - } - Some(msg) = monitor_failure_rx.recv() => { - break Some(crate::SpvError::ChannelFailure( - "event monitor".into(), - msg, - )) + let client = self.clone(); + let shutdown = monitor_shutdown.clone(); + let task = tokio::spawn(async move { + tracing::info!("Starting continuous network monitoring..."); + + // Run the sync loop + let mut sync_coordinator_tick_interval = + tokio::time::interval(SYNC_COORDINATOR_TICK_MS); + + let error: Option = loop { + let error: Option = tokio::select! { + _ = sync_coordinator_tick_interval.tick() => { + client.sync_coordinator.lock().await.tick().await.err().map(Into::into) + } + _ = monitor_shutdown.cancelled() => { + tracing::info!("Stopping network monitoring"); + break None + } + Some(msg) = monitor_failure_rx.recv() => { + break Some(crate::SpvError::ChannelFailure( + "event monitor".into(), + msg, + )) + } + }; + + if error.is_some() { + break error; } }; - if error.is_some() { - break error; + // Signal monitors to shut down before channels close + monitor_shutdown.cancel(); + let _ = tokio::join!( + sync_task, + chainlock_dispatch_task, + network_task, + wallet_task, + progress_task + ); + if let Some(task) = reservation_sweep_task { + let _ = task.await; } - }; - - // Signal monitors to shut down before channels close - monitor_shutdown.cancel(); - let _ = tokio::join!( - sync_task, - chainlock_dispatch_task, - network_task, - wallet_task, - progress_task - ); - if let Some(task) = reservation_sweep_task { - let _ = task.await; - } - if let Some(ref e) = error { - for handler in handlers.iter() { - handler.on_error(&e.to_string()); + if let Some(e) = error { + for handler in handlers.iter() { + handler.on_error(&e.to_string()); + } + // `stop()` waits for this task, so it runs in a task of its own. + tokio::spawn(async move { + if let Err(e) = client.stop().await { + tracing::warn!("Error stopping the client after a sync failure: {}", e); + } + }); } - } - - let stop_result = self.stop().await; + }); + *sync_loop = Some(SyncLoop { + task, + shutdown, + }); - match error { - Some(e) => Err(e), - None => stop_result, - } + Ok(()) } } diff --git a/dash-spv/src/lib.rs b/dash-spv/src/lib.rs index b0263f08e..9c8daeff9 100644 --- a/dash-spv/src/lib.rs +++ b/dash-spv/src/lib.rs @@ -43,6 +43,10 @@ //! //! client.run().await?; //! +//! // Sync in the background until Ctrl-C. +//! tokio::signal::ctrl_c().await?; +//! client.stop().await?; +//! //! Ok(()) //! } //! ``` diff --git a/dash-spv/src/main.rs b/dash-spv/src/main.rs index eca7d7a6a..02c2ede5a 100644 --- a/dash-spv/src/main.rs +++ b/dash-spv/src/main.rs @@ -404,18 +404,12 @@ async fn run_client( } }; - let stop_client = client.clone(); - tokio::spawn(async move { - if tokio::signal::ctrl_c().await.is_ok() { - tracing::debug!("Shutdown signal received"); - if let Err(e) = stop_client.stop().await { - tracing::warn!("Error during ctrl-c stop: {}", e); - } - } - }); - client.run().await?; + // Sync in the background until Ctrl-C. + tokio::signal::ctrl_c().await?; + client.stop().await?; + Ok(()) } diff --git a/dash-spv/tests/dashd_masternode/setup.rs b/dash-spv/tests/dashd_masternode/setup.rs index ff990562f..1cef15c96 100644 --- a/dash-spv/tests/dashd_masternode/setup.rs +++ b/dash-spv/tests/dashd_masternode/setup.rs @@ -1,7 +1,6 @@ use std::net::SocketAddr; use std::time::{Duration, Instant}; -use dash_spv::error::Result as SpvResult; use dash_spv::network::NetworkEvent; use dash_spv::test_utils::{ create_test_wallet, init_test_logging, next_unused_receive_address, retain_test_dir, @@ -21,7 +20,6 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use tempfile::TempDir; use tokio::sync::{broadcast, watch, RwLock}; -use tokio::task::JoinHandle; use tokio::time; /// Timeout for masternode sync tests (masternode sync takes longer than wallet sync). @@ -32,7 +30,6 @@ pub(super) type TestClient = pub(super) struct ClientHandle { pub(super) client: TestClient, - pub(super) run_handle: Option>>, pub(super) progress_receiver: watch::Receiver, pub(super) sync_event_receiver: broadcast::Receiver, pub(super) wallet_event_receiver: broadcast::Receiver, @@ -41,17 +38,14 @@ pub(super) struct ClientHandle { } impl ClientHandle { - pub(super) fn start(&mut self) { - let run_client = self.client.clone(); - self.run_handle = Some(tokio::task::spawn(async move { run_client.run().await })); + pub(super) async fn run(&mut self) { + tracing::info!("Starting client..."); + self.client.run().await.expect("client run failed"); } pub(super) async fn stop(&mut self) { - tracing::info!("Stopping client run loop..."); + tracing::info!("Stopping client..."); self.client.stop().await.expect("client stop failed"); - if let Some(handle) = self.run_handle.take() { - handle.await.expect("Run task panicked").expect("Run task returned error"); - } } } @@ -191,7 +185,6 @@ pub(super) async fn create_client( ClientHandle { client, - run_handle: None, progress_receiver, sync_event_receiver, wallet_event_receiver, @@ -200,13 +193,13 @@ pub(super) async fn create_client( } } -/// Built and started. Use [`create_client`] plus [`ClientHandle::start`] instead when +/// Built and started. Use [`create_client`] plus [`ClientHandle::run`] instead when /// the test needs to look at the client before it touches the network. pub(super) async fn create_and_start_client( config: &ClientConfig, wallet: Arc>>, ) -> ClientHandle { let mut handle = create_client(config, wallet).await; - handle.start(); + handle.run().await; handle } diff --git a/dash-spv/tests/dashd_masternode/tests_sync.rs b/dash-spv/tests/dashd_masternode/tests_sync.rs index 3b34227f4..52e0796a7 100644 --- a/dash-spv/tests/dashd_masternode/tests_sync.rs +++ b/dash-spv/tests/dashd_masternode/tests_sync.rs @@ -150,7 +150,7 @@ async fn test_masternode_list_sync_with_restart() { not default and let a fresh dashd sync cover for it" ); - client_handle.start(); + client_handle.run().await; let second_mn_progress = wait_for_masternode_sync(&mut client_handle.progress_receiver, SYNC_TIMEOUT).await; let second_height = second_mn_progress.current_height(); diff --git a/dash-spv/tests/dashd_sync/setup.rs b/dash-spv/tests/dashd_sync/setup.rs index e7396f98b..b224f70f2 100644 --- a/dash-spv/tests/dashd_sync/setup.rs +++ b/dash-spv/tests/dashd_sync/setup.rs @@ -234,8 +234,6 @@ pub(super) type TestClient = pub(super) struct ClientHandle { /// The underlying SPV client instance. pub(super) client: TestClient, - /// The handle to the client's run loop task. - pub(super) run_handle: Option>>, /// A channel for receiving progress updates. pub(super) progress_receiver: watch::Receiver, /// A channel for receiving sync events. @@ -247,19 +245,16 @@ pub(super) struct ClientHandle { } impl ClientHandle { - /// Spawns the client's run loop. - pub(super) fn spawn_run(&mut self) { - let client = self.client.clone(); - self.run_handle = Some(tokio::task::spawn(async move { client.run().await })); + /// Starts the SPV client. + pub(super) async fn run(&mut self) { + tracing::info!("Starting client..."); + self.client.run().await.expect("client run failed"); } - /// Stops the SPV client and awaits the termination of the background run task. + /// Stops the SPV client. pub(super) async fn stop(&mut self) { - tracing::info!("Stopping client run loop..."); + tracing::info!("Stopping client..."); self.client.stop().await.expect("client stop failed"); - if let Some(handle) = self.run_handle.take() { - handle.await.expect("Run task panicked").expect("Run task returned error"); - } } } @@ -308,13 +303,12 @@ pub(super) async fn create_and_start_client( }; let mut handle = ClientHandle { client, - run_handle: None, progress_receiver, sync_event_receiver, network_event_receiver, wallet_event_receiver, }; - handle.spawn_run(); + handle.run().await; handle } diff --git a/dash-spv/tests/dashd_sync/tests_restart.rs b/dash-spv/tests/dashd_sync/tests_restart.rs index 37dde0a62..429dd6b8d 100644 --- a/dash-spv/tests/dashd_sync/tests_restart.rs +++ b/dash-spv/tests/dashd_sync/tests_restart.rs @@ -214,7 +214,8 @@ async fn test_sync_with_random_restarts() { tracing::info!("Sync completed after {} random restarts (seed={})", num_restarts, seed); } -/// Verify the same client instance syncs again on every run after a stop. +/// Verify the same client syncs again after being stopped, including after +/// runs that were stopped right away. #[tokio::test] async fn test_sync_restarts_same_client() { let Some(ctx) = TestContext::new(TestChain::Full).await else { @@ -222,38 +223,40 @@ async fn test_sync_restarts_same_client() { }; let mut client_handle = ctx.spawn_new_client().await; - for run in 0..3 { - if run > 0 { - client_handle.sync_event_receiver = client_handle.sync_event_receiver.resubscribe(); - client_handle.spawn_run(); - } - wait_for_sync_complete(&mut client_handle.sync_event_receiver, ctx.dashd.initial_height) - .await; + for _ in 0..3 { + client_handle.run().await; client_handle.stop().await; - assert!(!client_handle.client.is_running()); } + + client_handle.run().await; + wait_for_sync_complete(&mut client_handle.sync_event_receiver, ctx.dashd.initial_height).await; + client_handle.stop().await; + + ctx.assert_synced(&client_handle.client.progress().await).await; } -/// Verify clearing the storage of a running client stops it, and its next run -/// syncs from scratch. +/// Verify clearing the storage of a running client stops it, and the client +/// syncs from scratch after its storage was cleared again between runs. #[tokio::test] async fn test_clear_storage_stops_running_client() { - let Some(ctx) = TestContext::new(TestChain::Full).await else { + let Some(ctx) = TestContext::new(TestChain::Minimal).await else { return; }; let mut client_handle = ctx.spawn_new_client().await; wait_for_sync_complete(&mut client_handle.sync_event_receiver, ctx.dashd.initial_height).await; client_handle.client.clear_storage().await.unwrap(); - let run_handle = client_handle.run_handle.take().unwrap(); - run_handle.await.unwrap().unwrap(); - - assert!(!client_handle.client.is_running()); + assert!(!client_handle.client.is_running().await); assert_eq!(client_handle.client.tip_height().await, 0); assert_eq!(client_handle.client.progress().await, SyncProgress::default()); + client_handle.stop().await; + client_handle.run().await; + client_handle.client.clear_storage().await.unwrap(); + client_handle.sync_event_receiver = client_handle.sync_event_receiver.resubscribe(); - client_handle.spawn_run(); + client_handle.run().await; wait_for_sync_complete(&mut client_handle.sync_event_receiver, ctx.dashd.initial_height).await; client_handle.stop().await; + ctx.assert_synced(&client_handle.client.progress().await).await; } diff --git a/dash-spv/tests/peer_test.rs b/dash-spv/tests/peer_test.rs index 8634d4413..9a945af53 100644 --- a/dash-spv/tests/peer_test.rs +++ b/dash-spv/tests/peer_test.rs @@ -52,8 +52,7 @@ async fn test_peer_connection() { let client = DashSpvClient::new(config, network_manager, storage_manager, wallet, vec![]).await.unwrap(); - let run_client = client.clone(); - let handle = tokio::spawn(async move { run_client.run().await }); + client.run().await.expect("Should run"); // Give it time to connect to peers time::sleep(Duration::from_secs(5)).await; @@ -63,7 +62,6 @@ async fn test_peer_connection() { assert!(peer_count > 0, "Should have connected to at least one peer"); client.stop().await.expect("Should stop"); - let _ = handle.await; } #[tokio::test] @@ -89,8 +87,7 @@ async fn test_peer_persistence() { .await .unwrap(); - let run_client = client.clone(); - let handle = tokio::spawn(async move { run_client.run().await }); + client.run().await.expect("Should run"); time::sleep(Duration::from_secs(5)).await; @@ -98,7 +95,6 @@ async fn test_peer_persistence() { assert!(peer_count > 0, "Should have connected to peers"); client.stop().await.expect("Should stop"); - let _ = handle.await; } // Second run: should load saved peers @@ -117,9 +113,8 @@ async fn test_peer_persistence() { .unwrap(); // Should connect faster due to saved peers - let run_client = client.clone(); let start = tokio::time::Instant::now(); - let handle = tokio::spawn(async move { run_client.run().await }); + client.run().await.expect("Should run"); // Wait for connection but with shorter timeout time::sleep(Duration::from_secs(3)).await; @@ -131,7 +126,6 @@ async fn test_peer_persistence() { println!("Connected to {} peers in {:?} (using saved peers)", peer_count, elapsed); client.stop().await.expect("Should stop"); - let _ = handle.await; } } diff --git a/dash-spv/tests/wallet_integration_test.rs b/dash-spv/tests/wallet_integration_test.rs deleted file mode 100644 index ef374ebfc..000000000 --- a/dash-spv/tests/wallet_integration_test.rs +++ /dev/null @@ -1,93 +0,0 @@ -//! Integration tests for wallet functionality. -//! -//! These tests validate end-to-end wallet operations through the SPVWalletManager. - -use std::sync::Arc; -use std::time::Duration; -use tempfile::TempDir; -use tokio::sync::RwLock; - -use dash_spv::network::PeerNetworkManager; -use dash_spv::storage::DiskStorageManager; -use dash_spv::{ClientConfig, DashSpvClient}; -use dashcore::Network; -use key_wallet::wallet::managed_wallet_info::ManagedWalletInfo; -use key_wallet_manager::WalletManager; -/// Create a test SPV client with memory storage for integration testing. -async fn create_test_client( -) -> DashSpvClient, PeerNetworkManager, DiskStorageManager> { - let config = ClientConfig::testnet() - .without_filters() - .with_storage_path(TempDir::new().unwrap().path()) - .without_masternodes() - // Ensure DNS discovery isn't used since it's causing flakiness in CI and not needed for these tests. - .with_restrict_to_configured_peers(true); - - // Create network manager - let network_manager = PeerNetworkManager::new(&config).await.unwrap(); - - // Create storage manager - let storage_manager = DiskStorageManager::new(&config).await.expect("Failed to create storage"); - - // Create wallet manager - let wallet = Arc::new(RwLock::new(WalletManager::::new(config.network))); - - DashSpvClient::new(config, network_manager, storage_manager, wallet, vec![]).await.unwrap() -} - -#[tokio::test] -async fn test_spv_client_creation() { - // Basic test to ensure client can be created - let client = create_test_client().await; - - // Verify client is created - assert_eq!(client.network().await, Network::Testnet); -} - -#[tokio::test] -async fn test_spv_client_run_stop() { - let client = create_test_client().await; - - // Run twice on the same instance: the watch must re-arm after stop(), so a - // second run() succeeds. - for _ in 0..2 { - let run_client = client.clone(); - let handle = tokio::spawn(async move { run_client.run().await }); - - tokio::time::timeout(Duration::from_secs(5), async { - while !client.is_running() { - tokio::time::sleep(Duration::from_millis(10)).await; - } - }) - .await - .expect("client failed to start"); - - client.stop().await.unwrap(); - - // stop() wakes the run loop immediately via the watch rather than only - // on the next 100ms coordinator tick, so the join completes far inside - // this bound. A hang (e.g. a missed wakeup) would trip the timeout. - tokio::time::timeout(Duration::from_secs(2), handle) - .await - .expect("run task did not exit promptly after stop") - .unwrap() - .unwrap(); - - assert!(!client.is_running()); - } -} - -#[tokio::test] -async fn test_wallet_manager_basic_operations() { - // Test basic wallet manager operations - let wallet_manager = WalletManager::::new(Network::Testnet); - - // Test that we can create a wallet manager - // Check wallet count - assert_eq!(wallet_manager.wallet_count(), 0); - - // Test adding a wallet (this would need actual wallet creation logic) - // For now, just verify the manager is working - let balance = wallet_manager.get_total_balance(); - assert_eq!(balance, 0); -} From 4cef65673cc76c51bbcea23b4d92c528420bfaa4 Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Sat, 3 Oct 2026 02:36:53 +0000 Subject: [PATCH 2/2] fix(dash-spv): stop only a sync loop that failed When the sync loop failed, it spawned a `stop()` that waited for the lock. Until that stop got it, a `run()` found the failed loop stored, returned `Ok` and the client was stopped right after. And if a caller stopped the failed loop and started a new one first, the deferred stop stopped the new, healthy loop. A failed loop has its `shutdown` token cancelled while it is still stored, which a stopped loop never is. The deferred stop is now `stop_failed()`, which stops only such a loop, `run()` tears one down before starting again, and `is_running()` no longer counts it. All three go through `stop_locked()`, the teardown `stop()` does under the lock. Co-Authored-By: Claude Opus 5.5 --- dash-spv/src/client/core.rs | 7 ++++--- dash-spv/src/client/lifecycle.rs | 28 +++++++++++++++++++------ dash-spv/src/client/sync_coordinator.rs | 11 ++++++++-- 3 files changed, 35 insertions(+), 11 deletions(-) diff --git a/dash-spv/src/client/core.rs b/dash-spv/src/client/core.rs index 4e6e29fe4..2b84d0037 100644 --- a/dash-spv/src/client/core.rs +++ b/dash-spv/src/client/core.rs @@ -113,8 +113,9 @@ pub struct DashSpvClient>, pub(super) masternode_engine: Option>>, pub(super) sync_coordinator: Arc>, - /// The running sync loop, `None` while stopped. `run` and `stop` hold the - /// lock throughout, so they never overlap. + /// The running sync loop, `None` while stopped. A loop whose `shutdown` is + /// cancelled has failed and waits to be torn down. `run` and `stop` hold + /// the lock throughout, so they never overlap. pub(super) sync_loop: Arc>>, pub(super) event_handlers: Arc>>, } @@ -163,7 +164,7 @@ impl DashSpvClient bool { - self.sync_loop.lock().await.is_some() + self.sync_loop.lock().await.as_ref().is_some_and(|running| !running.shutdown.is_cancelled()) } /// Returns the current chain tip hash if available. diff --git a/dash-spv/src/client/lifecycle.rs b/dash-spv/src/client/lifecycle.rs index ed316c659..5a60a0844 100644 --- a/dash-spv/src/client/lifecycle.rs +++ b/dash-spv/src/client/lifecycle.rs @@ -234,14 +234,30 @@ impl DashSpvClient Result<()> { let mut sync_loop = self.sync_loop.lock().await; - let Some(SyncLoop { + match sync_loop.take() { + Some(running) => self.stop_locked(running).await, + None => Ok(()), + } + } + + /// Stop the client if its sync loop failed. A loop that was stopped or + /// replaced by a later `run` in the meantime is left alone. + pub(super) async fn stop_failed(&self) -> Result<()> { + let mut sync_loop = self.sync_loop.lock().await; + match sync_loop.take_if(|running| running.shutdown.is_cancelled()) { + Some(failed) => self.stop_locked(failed).await, + None => Ok(()), + } + } + + /// Stop `sync_loop` and everything it drives. The caller holds the lock. + pub(super) async fn stop_locked( + &self, + SyncLoop { task, shutdown, - }) = sync_loop.take() - else { - return Ok(()); - }; - + }: SyncLoop, + ) -> Result<()> { // Stop the sync loop before tearing anything down so it cannot lock the // sync coordinator again. This prevents a tick from racing against the // shutdown below. diff --git a/dash-spv/src/client/sync_coordinator.rs b/dash-spv/src/client/sync_coordinator.rs index 975677aca..551bb4e64 100644 --- a/dash-spv/src/client/sync_coordinator.rs +++ b/dash-spv/src/client/sync_coordinator.rs @@ -33,6 +33,13 @@ impl DashSpvClient Result<()> { let mut sync_loop = self.sync_loop.lock().await; + // A loop that failed and still waits for its own stop is done: tear it + // down here, so the client really runs again. + if let Some(failed) = sync_loop.take_if(|running| running.shutdown.is_cancelled()) { + if let Err(e) = self.stop_locked(failed).await { + tracing::warn!("Error stopping the failed sync loop: {}", e); + } + } if sync_loop.is_some() { return Ok(()); } @@ -161,9 +168,9 @@ impl DashSpvClient