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
9 changes: 1 addition & 8 deletions dash-spv-bench/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() => {}
Expand All @@ -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")]
Expand Down
2 changes: 1 addition & 1 deletion dash-spv-ffi/FFI_API.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion dash-spv-ffi/src/bin/ffi_cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()));
Expand Down
78 changes: 14 additions & 64 deletions dash-spv-ffi/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<
Expand All @@ -24,7 +23,6 @@ type InnerClient = DashSpvClient<
pub struct FFIDashSpvClient {
pub(crate) inner: InnerClient,
pub(crate) runtime: Arc<Runtime>,
run_task: Mutex<Option<JoinHandle<()>>>,
}

impl FFIDashSpvClient {
Expand Down Expand Up @@ -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))
}
Expand All @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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.
Expand All @@ -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.
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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");
}
}
Expand Down
51 changes: 5 additions & 46 deletions dash-spv-ffi/tests/unit/test_client_lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -213,57 +212,17 @@ mod tests {

#[test]
#[serial]
fn test_client_error_callback_fires_on_start_failure() {
let (tx, rx) = mpsc::channel::<String>();
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<String>) };
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);
}
Expand Down
4 changes: 4 additions & 0 deletions dash-spv/examples/filter_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,10 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {

client.run().await?;

// Sync in the background until Ctrl-C.
tokio::signal::ctrl_c().await?;
client.stop().await?;

println!("Done!");
Ok(())
}
4 changes: 4 additions & 0 deletions dash-spv/examples/simple_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,10 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {

client.run().await?;

// Sync in the background until Ctrl-C.
tokio::signal::ctrl_c().await?;
client.stop().await?;

println!("Done!");
Ok(())
}
4 changes: 4 additions & 0 deletions dash-spv/examples/spv_with_wallet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,10 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {

client.run().await?;

// Sync in the background until Ctrl-C.
tokio::signal::ctrl_c().await?;
client.stop().await?;

println!("Done!");
Ok(())
}
26 changes: 18 additions & 8 deletions dash-spv/src/client/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -111,12 +113,20 @@ pub struct DashSpvClient<W: WalletInterface, N: NetworkManager, S: StorageManage
pub(super) wallet: Arc<RwLock<W>>,
pub(super) masternode_engine: Option<Arc<RwLock<MasternodeListEngine>>>,
pub(super) sync_coordinator: Arc<Mutex<SyncCoordinator>>,
/// `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>>,
/// 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<Mutex<Option<SyncLoop>>>,
pub(super) event_handlers: Arc<Vec<Arc<dyn super::EventHandler>>>,
}

/// 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<W: WalletInterface, N: NetworkManager, S: StorageManager> Clone for DashSpvClient<W, N, S> {
fn clone(&self) -> Self {
Self {
Expand All @@ -126,7 +136,7 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> 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),
}
}
Expand All @@ -152,9 +162,9 @@ impl<W: WalletInterface, N: NetworkManager, S: StorageManager> DashSpvClient<W,

// ============ State Queries ============

/// Check if the client is running.
pub fn is_running(&self) -> 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.as_ref().is_some_and(|running| !running.shutdown.is_cancelled())
}

/// Returns the current chain tip hash if available.
Expand Down
Loading
Loading