Skip to content
Open
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
1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,7 @@ base64 = { version = "0.22.1", default-features = false, features = ["std"] }
getrandom = { version = "0.3", default-features = false }
chrono = { version = "0.4", default-features = false, features = ["clock"] }
tokio = { version = "1.39", default-features = false, features = [ "rt-multi-thread", "time", "sync", "macros", "net" ] }
tokio-util = { version = "0.7", default-features = false, features = ["rt"] }
esplora-client = { version = "0.12", default-features = false, features = ["tokio", "async-https-rustls"] }
ldk-esplora-client = { package = "esplora-client", version = "0.13", default-features = false, features = ["tokio", "async-https-rustls"] }
electrum-client = { version = "0.25", default-features = false, features = ["proxy", "use-rustls-ring"] }
Expand Down
80 changes: 69 additions & 11 deletions src/chain/bitcoind.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ use lightning_block_sync::{
};
use serde::Serialize;

use super::WalletSyncStatus;
use super::{WalletSyncGuard, WalletSyncStatus};
use crate::config::{
BitcoindRestClientConfig, Config, DEFAULT_FEE_RATE_CACHE_UPDATE_TIMEOUT_SECS,
DEFAULT_TX_BROADCAST_TIMEOUT_SECS,
Expand All @@ -52,6 +52,31 @@ const CHAIN_POLLING_TIMEOUT_SECS: u64 = 10;
type BitcoindSpvClient =
SpvClient<ChainPoller<Arc<BitcoindClient>, BitcoindClient>, Arc<ChainListener>>;

async fn acquire_initial_wallet_sync_guard<'a>(
wallet_polling_status: &'a Mutex<WalletSyncStatus>,
stop_sync_receiver: &mut tokio::sync::watch::Receiver<()>,
) -> Option<WalletSyncGuard<'a>> {
loop {
let mut pending_sync = {
let mut status_lock = wallet_polling_status.lock().expect("lock");
match status_lock.register_or_subscribe_pending_sync() {
Some(pending_sync) => pending_sync,
None => {
return Some(WalletSyncGuard::new(
wallet_polling_status,
Error::WalletOperationFailed,
));
},
}
};
tokio::select! {
biased;
_ = stop_sync_receiver.changed() => return None,
_ = pending_sync.recv() => {},
}
}
}

pub(super) struct BitcoindChainSource {
api_client: Arc<BitcoindClient>,
spv_client: tokio::sync::Mutex<Option<BitcoindSpvClient>>,
Expand Down Expand Up @@ -160,12 +185,13 @@ impl BitcoindChainSource {
) {
// First register for the wallet polling status to make sure `Node::sync_wallets` calls
// wait on the result before proceeding.
{
let mut status_lock = self.wallet_polling_status.lock().expect("lock");
if status_lock.register_or_subscribe_pending_sync().is_some() {
debug_assert!(false, "Sync already in progress. This should never happen.");
}
}
let Some(initial_sync_guard) =
acquire_initial_wallet_sync_guard(&self.wallet_polling_status, &mut stop_sync_receiver)
.await
else {
log_trace!(self.logger, "Stopping initial chain sync.");
return;
};

log_info!(
self.logger,
Expand Down Expand Up @@ -302,7 +328,7 @@ impl BitcoindChainSource {
}

// Now propagate the initial result to unblock waiting subscribers.
self.wallet_polling_status.lock().expect("lock").propagate_result_to_subscribers(Ok(()));
initial_sync_guard.complete(Ok(()));

let mut chain_polling_interval =
tokio::time::interval(Duration::from_secs(CHAIN_POLLING_INTERVAL_SECS));
Expand Down Expand Up @@ -413,6 +439,8 @@ impl BitcoindChainSource {
Error::WalletOperationFailed
})?;
}
let sync_guard =
WalletSyncGuard::new(&self.wallet_polling_status, Error::WalletOperationFailed);

let res = self
.poll_and_update_listeners_inner(
Expand All @@ -423,7 +451,7 @@ impl BitcoindChainSource {
)
.await;

self.wallet_polling_status.lock().expect("lock").propagate_result_to_subscribers(res);
sync_guard.complete(res);

res
}
Expand Down Expand Up @@ -1588,6 +1616,9 @@ impl std::error::Error for BitcoindClientError {}

#[cfg(test)]
mod tests {
use std::sync::Mutex;
use std::time::Duration;

use bitcoin::hashes::Hash;
use bitcoin::{FeeRate, OutPoint, ScriptBuf, Transaction, TxIn, TxOut, Txid, Witness};
use lightning_block_sync::http::JsonResponse;
Expand All @@ -1597,9 +1628,36 @@ mod tests {
use serde_json::json;

use crate::chain::bitcoind::{
FeeResponse, GetMempoolEntryResponse, GetRawMempoolResponse, GetRawTransactionResponse,
MempoolMinFeeResponse,
acquire_initial_wallet_sync_guard, FeeResponse, GetMempoolEntryResponse,
GetRawMempoolResponse, GetRawTransactionResponse, MempoolMinFeeResponse,
};
use crate::chain::{WalletSyncGuard, WalletSyncStatus};
use crate::Error;

#[tokio::test]
async fn initial_sync_waits_for_in_progress_sync() {
let status = Mutex::new(WalletSyncStatus::Completed);
assert!(status.lock().expect("lock").register_or_subscribe_pending_sync().is_none());
let in_progress_guard = WalletSyncGuard::new(&status, Error::WalletOperationFailed);
let (_stop_sender, mut stop_receiver) = tokio::sync::watch::channel(());
let mut acquire_guard =
Box::pin(acquire_initial_wallet_sync_guard(&status, &mut stop_receiver));

let early_result =
tokio::time::timeout(Duration::from_millis(10), acquire_guard.as_mut()).await;
assert!(early_result.is_err(), "background sync should wait for the active sync");

in_progress_guard.complete(Ok(()));
let acquired_guard = tokio::time::timeout(Duration::from_secs(1), acquire_guard)
.await
.expect("background sync should resume")
.expect("background sync should acquire the sync guard");
assert!(
matches!(*status.lock().expect("lock"), WalletSyncStatus::InProgress { .. }),
"background sync should own the next sync"
);
acquired_guard.complete(Ok(()));
}

prop_compose! {
fn arbitrary_witness()(
Expand Down
13 changes: 7 additions & 6 deletions src/chain/electrum.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ use lightning::chain::{Confirm, Filter, WatchedOutput};
use lightning::util::ser::Writeable;
use lightning_transaction_sync::ElectrumSyncClient;

use super::WalletSyncStatus;
use super::{WalletSyncGuard, WalletSyncStatus};
use crate::config::{
clamp_full_scan_stop_gap, Config, ElectrumSyncConfig, MAX_FULL_SCAN_STOP_GAP,
MIN_FULL_SCAN_STOP_GAP,
Expand Down Expand Up @@ -113,10 +113,12 @@ impl ElectrumChainSource {
Error::WalletOperationFailed
})?;
}
let sync_guard =
WalletSyncGuard::new(&self.onchain_wallet_sync_status, Error::WalletOperationFailed);

let res = self.sync_onchain_wallet_inner(onchain_wallet).await;

self.onchain_wallet_sync_status.lock().expect("lock").propagate_result_to_subscribers(res);
sync_guard.complete(res);

res
}
Expand Down Expand Up @@ -223,14 +225,13 @@ impl ElectrumChainSource {
Error::TxSyncFailed
})?;
}
let sync_guard =
WalletSyncGuard::new(&self.lightning_wallet_sync_status, Error::TxSyncFailed);

let res =
self.sync_lightning_wallet_inner(channel_manager, chain_monitor, output_sweeper).await;

self.lightning_wallet_sync_status
.lock()
.expect("lock")
.propagate_result_to_subscribers(res);
sync_guard.complete(res);

res
}
Expand Down
13 changes: 7 additions & 6 deletions src/chain/esplora.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ use lightning::chain::{Confirm, Filter, WatchedOutput};
use lightning::util::ser::Writeable;
use lightning_transaction_sync::EsploraSyncClient;

use super::WalletSyncStatus;
use super::{WalletSyncGuard, WalletSyncStatus};
use crate::config::{
clamp_full_scan_stop_gap, Config, EsploraSyncConfig, BDK_CLIENT_CONCURRENCY,
MAX_FULL_SCAN_STOP_GAP, MIN_FULL_SCAN_STOP_GAP,
Expand Down Expand Up @@ -133,10 +133,12 @@ impl EsploraChainSource {
Error::WalletOperationFailed
})?;
}
let sync_guard =
WalletSyncGuard::new(&self.onchain_wallet_sync_status, Error::WalletOperationFailed);

let res = self.sync_onchain_wallet_inner(onchain_wallet).await;

self.onchain_wallet_sync_status.lock().expect("lock").propagate_result_to_subscribers(res);
sync_guard.complete(res);

res
}
Expand Down Expand Up @@ -283,14 +285,13 @@ impl EsploraChainSource {
Error::WalletOperationFailed
})?;
}
let sync_guard =
WalletSyncGuard::new(&self.lightning_wallet_sync_status, Error::WalletOperationFailed);

let res =
self.sync_lightning_wallet_inner(channel_manager, chain_monitor, output_sweeper).await;

self.lightning_wallet_sync_status
.lock()
.expect("lock")
.propagate_result_to_subscribers(res);
sync_guard.complete(res);

res
}
Expand Down
64 changes: 54 additions & 10 deletions src/chain/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,34 @@ pub(crate) enum WalletSyncStatus {
InProgress { subscribers: tokio::sync::broadcast::Sender<Result<(), Error>> },
}

pub(crate) struct WalletSyncGuard<'a> {
status: &'a Mutex<WalletSyncStatus>,
cancellation_error: Error,
active: bool,
}

impl<'a> WalletSyncGuard<'a> {
pub(crate) fn new(status: &'a Mutex<WalletSyncStatus>, cancellation_error: Error) -> Self {
Self { status, cancellation_error, active: true }
}

pub(crate) fn complete(mut self, res: Result<(), Error>) {
self.status.lock().expect("lock").propagate_result_to_subscribers(res);
self.active = false;
}
}

impl Drop for WalletSyncGuard<'_> {
fn drop(&mut self) {
if self.active {
self.status
.lock()
.expect("lock")
.propagate_result_to_subscribers(Err(self.cancellation_error));
}
}
}

impl WalletSyncStatus {
fn register_or_subscribe_pending_sync(
&mut self,
Expand Down Expand Up @@ -95,16 +123,7 @@ impl WalletSyncStatus {
WalletSyncStatus::InProgress { subscribers } => {
// A sync is in-progress, we notify subscribers.
if subscribers.receiver_count() > 0 {
match subscribers.send(res) {
Ok(_) => (),
Err(e) => {
debug_assert!(
false,
"Failed to send wallet sync result to subscribers: {:?}",
e
);
},
}
let _ = subscribers.send(res);
}
*self = WalletSyncStatus::Completed;
},
Expand Down Expand Up @@ -561,3 +580,28 @@ impl Filter for ChainSource {
}
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn wallet_sync_guard_resets_abandoned_sync() {
let status = Mutex::new(WalletSyncStatus::Completed);
assert!(status.lock().expect("lock").register_or_subscribe_pending_sync().is_none());
let sync_guard = WalletSyncGuard::new(&status, Error::WalletOperationFailed);
let mut subscriber = status
.lock()
.expect("lock")
.register_or_subscribe_pending_sync()
.expect("sync subscriber");

drop(sync_guard);

assert!(
matches!(*status.lock().expect("lock"), WalletSyncStatus::Completed),
"abandoned wallet sync should reset its status"
);
assert_eq!(subscriber.try_recv(), Ok(Err(Error::WalletOperationFailed)));
}
}
Loading
Loading