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
73 changes: 18 additions & 55 deletions src/chain/cbf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@

use std::collections::{BTreeMap, HashMap};
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::{Arc, Mutex, RwLock};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};

Expand Down Expand Up @@ -43,6 +43,11 @@ const MIN_FEERATE_SAT_PER_KWU: u64 = 250;
/// Number of recent blocks to look back for per-target fee rate estimation.
const FEE_RATE_LOOKBACK_BLOCKS: usize = 6;

/// Number of blocks to walk back from a component's persisted best block height
/// for reorg safety when computing the incremental scan skip height.
/// Matches bdk-kyoto's `IMPOSSIBLE_REORG_DEPTH`.
const REORG_SAFETY_BLOCKS: u32 = 7;

/// The fee estimation back-end used by the CBF chain source.
enum FeeSource {
/// Derive fee rates from the coinbase reward of recent blocks.
Expand Down Expand Up @@ -82,12 +87,6 @@ pub(super) struct CbfChainSource {
scan_lock: tokio::sync::Mutex<()>,
/// Scripts registered by LDK's Filter trait for lightning channel monitoring.
registered_scripts: Mutex<Vec<ScriptBuf>>,
/// Set when new scripts are registered; forces a full rescan on next lightning sync.
lightning_scripts_dirty: Arc<AtomicBool>,
/// Last block height reached by on-chain wallet sync, used for incremental scans.
last_onchain_synced_height: Arc<Mutex<Option<u32>>>,
/// Last block height reached by lightning wallet sync, used for incremental scans.
last_lightning_synced_height: Arc<Mutex<Option<u32>>>,
/// Deduplicates concurrent on-chain wallet sync requests.
onchain_wallet_sync_status: Mutex<WalletSyncStatus>,
/// Deduplicates concurrent lightning wallet sync requests.
Expand Down Expand Up @@ -116,9 +115,6 @@ struct CbfEventState {
matched_block_hashes: Arc<Mutex<Vec<(u32, BlockHash)>>>,
sync_completion_tx: Arc<Mutex<Option<oneshot::Sender<SyncUpdate>>>>,
filter_skip_height: Arc<AtomicU32>,
last_onchain_synced_height: Arc<Mutex<Option<u32>>>,
last_lightning_synced_height: Arc<Mutex<Option<u32>>>,
lightning_scripts_dirty: Arc<AtomicBool>,
}

impl CbfChainSource {
Expand Down Expand Up @@ -149,10 +145,7 @@ impl CbfChainSource {
let sync_completion_tx = Arc::new(Mutex::new(None));
let filter_skip_height = Arc::new(AtomicU32::new(0));
let registered_scripts = Mutex::new(Vec::new());
let lightning_scripts_dirty = Arc::new(AtomicBool::new(true));
let scan_lock = tokio::sync::Mutex::new(());
let last_onchain_synced_height = Arc::new(Mutex::new(None));
let last_lightning_synced_height = Arc::new(Mutex::new(None));
let onchain_wallet_sync_status = Mutex::new(WalletSyncStatus::Completed);
let lightning_wallet_sync_status = Mutex::new(WalletSyncStatus::Completed);
Ok(Self {
Expand All @@ -166,10 +159,7 @@ impl CbfChainSource {
sync_completion_tx,
filter_skip_height,
registered_scripts,
lightning_scripts_dirty,
scan_lock,
last_onchain_synced_height,
last_lightning_synced_height,
onchain_wallet_sync_status,
lightning_wallet_sync_status,
fee_estimator,
Expand Down Expand Up @@ -249,9 +239,6 @@ impl CbfChainSource {
matched_block_hashes: Arc::clone(&self.matched_block_hashes),
sync_completion_tx: Arc::clone(&self.sync_completion_tx),
filter_skip_height: Arc::clone(&self.filter_skip_height),
last_onchain_synced_height: Arc::clone(&self.last_onchain_synced_height),
last_lightning_synced_height: Arc::clone(&self.last_lightning_synced_height),
lightning_scripts_dirty: Arc::clone(&self.lightning_scripts_dirty),
};
let event_logger = Arc::clone(&self.logger);
runtime.spawn_cancellable_background_task(Self::process_events(
Expand Down Expand Up @@ -320,24 +307,9 @@ impl CbfChainSource {
accepted.len(),
);

// Reset synced heights to just before the earliest reorganized
// block so the next incremental scan covers the affected range.
if let Some(min_reorg_height) = reorganized.iter().map(|h| h.height).min() {
let reset_height = if min_reorg_height > 0 {
Some(min_reorg_height - 1)
} else {
None
};
*state.last_onchain_synced_height.lock().unwrap() = reset_height;
*state.last_lightning_synced_height.lock().unwrap() = reset_height;
state.lightning_scripts_dirty.store(true, Ordering::Release);
log_debug!(
logger,
"Reset synced heights to {:?} due to reorg at height {}.",
reset_height,
min_reorg_height,
);
}
// No height reset needed: skip heights are derived from
// BDK's checkpoint (on-chain) and LDK's best block
// (lightning), both walked back by REORG_SAFETY_BLOCKS.
},
BlockHeaderChanges::Connected(header) => {
log_trace!(logger, "CBF block connected at height {}", header.height,);
Expand Down Expand Up @@ -382,13 +354,11 @@ impl CbfChainSource {
/// Register a transaction script for Lightning channel monitoring.
pub(crate) fn register_tx(&self, _txid: &Txid, script_pubkey: &Script) {
self.registered_scripts.lock().unwrap().push(script_pubkey.to_owned());
self.lightning_scripts_dirty.store(true, Ordering::Release);
}

/// Register a watched output script for Lightning channel monitoring.
pub(crate) fn register_output(&self, output: WatchedOutput) {
self.registered_scripts.lock().unwrap().push(output.script_pubkey.clone());
self.lightning_scripts_dirty.store(true, Ordering::Release);
}

/// Run a CBF filter scan: set watched scripts, trigger a rescan, wait for
Expand Down Expand Up @@ -458,7 +428,7 @@ impl CbfChainSource {
Duration::from_secs(
self.sync_config.timeouts_config.onchain_wallet_sync_timeout_secs,
),
self.sync_onchain_wallet_op(requester, scripts),
self.sync_onchain_wallet_op(requester, &onchain_wallet, scripts),
);

let (tx_update, sync_update) = match timeout_fut.await {
Expand Down Expand Up @@ -513,12 +483,11 @@ impl CbfChainSource {
}

async fn sync_onchain_wallet_op(
&self, requester: Requester, scripts: Vec<ScriptBuf>,
&self, requester: Requester, onchain_wallet: &Wallet, scripts: Vec<ScriptBuf>,
) -> Result<(TxUpdate<ConfirmationBlockTime>, SyncUpdate), Error> {
// Always do a full scan (skip_height=None) for the on-chain wallet.
// Unlike the Lightning wallet which can rely on reorg_queue events,
// the on-chain wallet needs to see all blocks to correctly detect
// reorgs via checkpoint comparison in the caller.
// Derive skip height from BDK's persisted checkpoint, walked back by
// REORG_SAFETY_BLOCKS for reorg safety (same approach as bdk-kyoto).
// This survives restarts since BDK persists its checkpoint chain.
//
// We include LDK-registered scripts (e.g., channel funding output
// scripts) alongside the wallet scripts. This ensures the on-chain
Expand All @@ -530,9 +499,10 @@ impl CbfChainSource {
// unknown. This mirrors what the Bitcoind chain source does in
// `Wallet::block_connected` by inserting registered tx outputs.
let mut all_scripts = scripts;
// we query all registered scripts, not only BDK-related
all_scripts.extend(self.registered_scripts.lock().unwrap().iter().cloned());
let (sync_update, matched) = self.run_filter_scan(all_scripts, None).await?;
let skip_height =
onchain_wallet.latest_checkpoint().height().checked_sub(REORG_SAFETY_BLOCKS);
let (sync_update, matched) = self.run_filter_scan(all_scripts, skip_height).await?;

log_debug!(
self.logger,
Expand Down Expand Up @@ -561,9 +531,6 @@ impl CbfChainSource {
}
}

let tip = sync_update.tip();
*self.last_onchain_synced_height.lock().unwrap() = Some(tip.height);

Ok((tx_update, sync_update))
}

Expand Down Expand Up @@ -644,9 +611,8 @@ impl CbfChainSource {
&self, requester: Requester, channel_manager: Arc<ChannelManager>,
chain_monitor: Arc<ChainMonitor>, output_sweeper: Arc<Sweeper>, scripts: Vec<ScriptBuf>,
) -> Result<(), Error> {
let scripts_dirty = self.lightning_scripts_dirty.load(Ordering::Acquire);
let skip_height =
if scripts_dirty { None } else { *self.last_lightning_synced_height.lock().unwrap() };
channel_manager.current_best_block().height.checked_sub(REORG_SAFETY_BLOCKS);
let (sync_update, matched) = self.run_filter_scan(scripts, skip_height).await?;

log_debug!(
Expand Down Expand Up @@ -681,9 +647,6 @@ impl CbfChainSource {
}
}

*self.last_lightning_synced_height.lock().unwrap() = Some(tip.height);
self.lightning_scripts_dirty.store(false, Ordering::Release);

Ok(())
}

Expand Down
1 change: 1 addition & 0 deletions tests/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -509,6 +509,7 @@ pub(crate) fn setup_node(chain_source: &TestChainSource, config: TestConfig) ->
let sync_config = CbfSyncConfig {
background_sync_config: None,
timeouts_config,
required_peers: 1,
..Default::default()
};
builder.set_chain_source_cbf(vec![peer_addr], Some(sync_config), None);
Expand Down
Loading