diff --git a/src/chain/mod.rs b/src/chain/mod.rs index 0f96c409f..58a6586c0 100644 --- a/src/chain/mod.rs +++ b/src/chain/mod.rs @@ -505,20 +505,21 @@ impl ChainSource { } Some(next_package) = receiver.recv() => { // Classify funding broadcasts into payment records before sending. If - // classification fails we skip the broadcast, since broadcasting a tx we - // failed to record would leave it on-chain without a payment. - let package = match self.tx_broadcaster.classify_package(next_package).await { - Ok(package) => package, - Err(e) => { - log_error!( - tx_bcast_logger, - "Skipping broadcast: failed to persist payment records: {:?}", - e, - ); - continue; - }, - }; - let package = package.into_sorted_transactions(); + // classification fails we delay the broadcast and retry, since broadcasting + // a tx we failed to record would leave it on-chain without a payment — + // while dropping the package would not keep an interactively funded tx + // off-chain (the counterparty broadcasts it regardless), only leave it + // confirming without a recorded candidate. + if let Err(e) = self.tx_broadcaster.classify_package(&next_package).await { + log_error!( + tx_bcast_logger, + "Delaying broadcast: failed to persist payment records, will retry: {:?}", + e, + ); + self.tx_broadcaster.requeue_failed_classify(next_package); + continue; + } + let package = next_package.into_sorted_transactions(); match &self.kind { ChainSourceKind::Esplora(esplora_chain_source) => { esplora_chain_source.process_transaction_broadcast(package).await diff --git a/src/tx_broadcaster.rs b/src/tx_broadcaster.rs index 782112dad..248926d45 100644 --- a/src/tx_broadcaster.rs +++ b/src/tx_broadcaster.rs @@ -7,6 +7,7 @@ use std::ops::Deref; use std::sync::{Mutex as StdMutex, Weak}; +use std::time::Duration; use bitcoin::Transaction; use lightning::chain::chaininterface::{ @@ -20,6 +21,11 @@ use crate::Error; const BCAST_PACKAGE_QUEUE_SIZE: usize = 256; +/// How long to wait before re-classifying a package whose classification failed. Long enough to +/// give a struggling store room to recover, short against the ~minutes until the transaction +/// could confirm. +const FAILED_CLASSIFY_RETRY_DELAY: Duration = Duration::from_secs(2); + /// A package of transactions that LDK handed to the broadcaster in one `broadcast_transactions` /// call, along with each transaction's type. Queued until the background task classifies and /// broadcasts it. Built only via [`BroadcastPackage::new`] from such a call, so unrelated @@ -133,12 +139,11 @@ where self.queue_receiver.lock().await } - /// Classifies a queued package into payment records and returns the package ready for the - /// chain client. Returns `Err` if any classification fails; callers must not broadcast the - /// package in that case, since a crash would leave the transaction on-chain without a record. - pub(crate) async fn classify_package( - &self, package: BroadcastPackage, - ) -> Result { + /// Classifies a queued package into payment records. Returns `Err` if any classification + /// fails; callers must not broadcast the package in that case, since a crash would leave the + /// transaction on-chain without a record — but must requeue it via + /// [`Self::requeue_failed_classify`] rather than drop it. + pub(crate) async fn classify_package(&self, package: &BroadcastPackage) -> Result<(), Error> { let wallet_opt = self.wallet.lock().expect("lock").as_ref().and_then(Weak::upgrade); if let Some(wallet) = wallet_opt { for (tx, tx_type) in package.transactions() { @@ -147,7 +152,21 @@ where } } } - Ok(package) + Ok(()) + } + + /// Re-sends a package whose classification failed back into the queue after a delay, so a + /// transient persistence failure delays the broadcast instead of dropping the package. + /// Dropping an interactive-funding package would not even keep its transaction off-chain — + /// the counterparty broadcasts it regardless — it would only leave the transaction + /// confirming without a recorded candidate. If the queue has closed by the time the delay + /// elapses, the node is shutting down and the package is dropped with it. + pub(crate) fn requeue_failed_classify(&self, package: BroadcastPackage) { + let sender = self.queue_sender.clone(); + tokio::spawn(async move { + tokio::time::sleep(FAILED_CLASSIFY_RETRY_DELAY).await; + let _ = sender.send(package).await; + }); } pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) { diff --git a/src/wallet/mod.rs b/src/wallet/mod.rs index df95c11ec..f7dacc4b0 100644 --- a/src/wallet/mod.rs +++ b/src/wallet/mod.rs @@ -341,11 +341,11 @@ impl Wallet { // duplicating) the record classification just wrote. let guard = self.funding_payment_update_lock.lock().await; - let payment_id = self + let mut payment_id = self .find_payment_by_txid(txid) .unwrap_or_else(|| PaymentId(txid.to_byte_array())); - if self + match self .apply_funding_status_update_locked( &guard, payment_id, @@ -354,7 +354,13 @@ impl Wallet { ) .await? { - continue; + FundingStatusUpdate::Applied => continue, + FundingStatusUpdate::NotFunding => {}, + // Not part of the funding payment's history (e.g. a close spending the + // funding outpoint): record it under its own id below instead. + FundingStatusUpdate::Foreign => { + payment_id = PaymentId(txid.to_byte_array()); + }, } let payment = { @@ -479,11 +485,11 @@ impl Wallet { // with classification. let guard = self.funding_payment_update_lock.lock().await; - let payment_id = self + let mut payment_id = self .find_payment_by_txid(txid) .unwrap_or_else(|| PaymentId(txid.to_byte_array())); - if self + match self .apply_funding_status_update_locked( &guard, payment_id, @@ -492,7 +498,13 @@ impl Wallet { ) .await? { - continue; + FundingStatusUpdate::Applied => continue, + FundingStatusUpdate::NotFunding => {}, + // Not part of the funding payment's history (e.g. a close spending the + // funding outpoint): record it under its own id below instead. + FundingStatusUpdate::Foreign => { + payment_id = PaymentId(txid.to_byte_array()); + }, } let payment = { @@ -554,11 +566,11 @@ impl Wallet { // with classification. let guard = self.funding_payment_update_lock.lock().await; - let payment_id = self + let mut payment_id = self .find_payment_by_txid(txid) .unwrap_or_else(|| PaymentId(txid.to_byte_array())); - if self + match self .apply_funding_status_update_locked( &guard, payment_id, @@ -567,7 +579,13 @@ impl Wallet { ) .await? { - continue; + FundingStatusUpdate::Applied => continue, + FundingStatusUpdate::NotFunding => {}, + // Not part of the funding payment's history (e.g. a close spending the + // funding outpoint): record it under its own id below instead. + FundingStatusUpdate::Foreign => { + payment_id = PaymentId(txid.to_byte_array()); + }, } let payment = { @@ -1927,9 +1945,11 @@ impl Wallet { /// If `payment_id` refers to a classified funding payment, refreshes its confirmation status /// and the candidate txid the event refers to, while preserving the contribution-derived /// amount/fee and `tx_type` that wallet sync must not recompute from its own view: the wallet's - /// `sent`/`received` don't capture our contribution to a shared funding output. Returns `true` - /// when it handled the payment, so the caller skips the default on-chain path. Graduation to - /// `Succeeded` is left to `ChainTipChanged` after `ANTI_REORG_DELAY`. + /// `sent`/`received` don't capture our contribution to a shared funding output. Returns + /// [`FundingStatusUpdate::Applied`] when it handled the payment, so the caller skips the + /// default on-chain path — or [`FundingStatusUpdate::Foreign`] when the transaction is not + /// part of the payment's funding history, so the caller records it under its own id. + /// Graduation to `Succeeded` is left to `ChainTipChanged` after `ANTI_REORG_DELAY`. /// /// The caller must hold [`Self::funding_payment_update_lock`] — from resolving `payment_id` /// through its own last write, not just across this call — so that classification's two-store @@ -1938,36 +1958,47 @@ impl Wallet { async fn apply_funding_status_update_locked( &self, _guard: &tokio::sync::MutexGuard<'_, ()>, payment_id: PaymentId, event_txid: Txid, confirmation_status: ConfirmationStatus, - ) -> Result { + ) -> Result { // The funding-type gate, the candidate lookup, and the write share the store's mutation // lock: against a separate `get`, a classification merging in between would have its // `tx_type` and contribution figures clobbered by this stale snapshot. + let mut outcome = FundingStatusUpdate::NotFunding; let mut handled = None; self.payment_store .mutate(&payment_id, |existing| { let payment = existing?; - let tx_type = match &payment.kind { + let (current_txid, tx_type) = match &payment.kind { PaymentKind::Onchain { + txid, tx_type: tx_type @ Some( TransactionType::Funding { .. } | TransactionType::InteractiveFunding { .. }, ), .. - } => tx_type.clone(), + } => (*txid, tx_type.clone()), _ => return None, }; + // Adopt the event's txid only when the transaction is part of this payment's + // funding history: its current txid or a classified candidate. A conflicting + // transaction that is neither — a close also spends the funding outpoint — must + // not overwrite the record. + let pending = self.pending_payment_store.get(&payment_id); + let owns_event_tx = event_txid == current_txid + || pending.as_ref().is_some_and(|p| p.candidate(event_txid).is_some()); + if !owns_event_tx { + outcome = FundingStatusUpdate::Foreign; + return None; + } // Report the figures of the candidate that actually confirmed, which need not be // the last one broadcast (an earlier, lower-fee candidate may win) and may carry // no figures at all (`None`) for a round we didn't contribute to. (`direction` is // invariant across a splice's candidates and cannot be changed through the store // anyway.) let mut target = payment.clone(); - if let Some(pending) = self.pending_payment_store.get(&payment_id) { - if let Some(candidate) = pending.candidate(event_txid) { - target.amount_msat = candidate.amount_msat; - target.fee_paid_msat = candidate.fee_paid_msat; - } + if let Some(candidate) = pending.as_ref().and_then(|p| p.candidate(event_txid)) { + target.amount_msat = candidate.amount_msat; + target.fee_paid_msat = candidate.fee_paid_msat; } target.kind = PaymentKind::Onchain { txid: event_txid, status: confirmation_status, tx_type }; @@ -1985,7 +2016,7 @@ impl Wallet { }) .await?; let Some(payment) = handled else { - return Ok(false); + return Ok(outcome); }; // Mirror the refreshed confirmation status onto the pending entry: `ChainTipChanged` // graduates by reading the pending entry's details, so it must see the new status. This is @@ -1995,13 +2026,20 @@ impl Wallet { let pending = self.create_pending_payment_from_tx(payment, Vec::new()); self.pending_payment_store.insert_or_update(pending).await?; } - Ok(true) + Ok(FundingStatusUpdate::Applied) } #[allow(deprecated)] pub(crate) async fn bump_fee_rbf( &self, payment_id: PaymentId, fee_rate: Option, cur_anchor_reserve_sats: u64, ) -> Result { + let mut locked_persister = self.persister.lock().await; + // Hold the cross-store lock from the record read through the replacement writes: funding + // classification re-types records concurrently, and a classification landing after the + // funding-kind check below would let the RBF replace a funding transaction. Acquired + // after the persister, matching the lock order of the wallet sync paths. + let funding_guard = self.funding_payment_update_lock.lock().await; + let payment = self.payment_store.get(&payment_id).ok_or_else(|| { log_error!(self.logger, "Payment {} not found in payment store", payment_id); Error::InvalidPaymentId @@ -2060,7 +2098,6 @@ impl Wallet { }, }; - let mut locked_persister = self.persister.lock().await; let mut locked_wallet = self.inner.lock().expect("lock"); debug_assert!( @@ -2247,6 +2284,7 @@ impl Wallet { self.payment_store.insert_or_update(new_payment).await?; self.pending_payment_store.insert_or_update(pending_payment_store).await?; + drop(funding_guard); self.broadcaster.broadcast_unclassified_transaction(fee_bumped_tx); @@ -2298,6 +2336,20 @@ fn aggregate_local_stakes(candidate: &FundingCandidate) -> LocalStakeAggregate { } } +/// The outcome of [`Wallet::apply_funding_status_update_locked`]. +enum FundingStatusUpdate { + /// The event's transaction belongs to the funding payment; its refreshed confirmation status + /// was applied (or was already current). + Applied, + /// The resolved payment is not a classified funding payment; the caller's default on-chain + /// handling applies under the resolved id. + NotFunding, + /// The event's transaction is not part of the funding payment's history — e.g. a close + /// spending the same funding outpoint — so the funding record must not adopt it; the caller + /// should record the transaction under its own txid-derived id. + Foreign, +} + impl Listen for Wallet { fn filtered_block_connected( &self, _header: &bitcoin::block::Header, @@ -3929,6 +3981,93 @@ mod tests { assert_eq!(wallet.find_payment_by_txid(txid2), Some(payment_id)); } + /// A cooperative close conflicts with a pending splice's funding transaction — both spend the + /// pre-splice funding outpoint — so sync records the close among the splice record's + /// conflicting txids, and the close's confirmation then resolves to the splice's PaymentId. + /// The funding record must not adopt the close's txid and confirmation as its own: the close + /// is not a round of the splice. It must land on a record keyed by the close's own id. + #[tokio::test] + async fn funding_record_does_not_adopt_a_conflicting_close() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let funding_outpoint = + bitcoin::OutPoint { txid: Txid::from_byte_array([3u8; 32]), vout: 0 }; + + // The close pays the shutdown script, which is a wallet address. + let script_pubkey = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + let close_tx = Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: LockTime::ZERO, + input: vec![bitcoin::TxIn { + previous_output: funding_outpoint, + script_sig: bitcoin::ScriptBuf::new(), + sequence: bitcoin::Sequence::MAX, + witness: bitcoin::Witness::new(), + }], + output: vec![TxOut { value: Amount::from_sat(90_000), script_pubkey }], + }; + let close_txid = close_tx.compute_txid(); + + let splice_txid = Txid::from_byte_array([2u8; 32]); + let payment_id = PaymentId([21u8; 32]); + let candidates = vec![FundingTxCandidate { + txid: splice_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }]; + let details = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + wallet.persist_funding_payment(details, candidates).await.unwrap(); + + // Sync saw the close double-spend the splice's funding transaction. + wallet + .pending_payment_store + .update(PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: Some(vec![close_txid]), + candidates: Vec::new(), + }) + .await + .unwrap(); + + let event = WalletEvent::TxConfirmed { + txid: close_txid, + tx: Arc::new(close_tx), + block_time: confirmed_block_time(5), + old_block_time: None, + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let funding = wallet.payment_store.get(&payment_id).unwrap(); + match &funding.kind { + PaymentKind::Onchain { txid, status, tx_type } => { + assert_eq!(*txid, splice_txid, "the record must not adopt the close's txid"); + assert!(matches!(status, ConfirmationStatus::Unconfirmed)); + assert!(matches!(tx_type, Some(TransactionType::InteractiveFunding { .. }))); + }, + kind => panic!("unexpected kind {:?}", kind), + } + assert_eq!(funding.amount_msat, Some(1_000_000)); + assert_eq!(funding.fee_paid_msat, Some(500)); + + let close = wallet.payment_store.get(&PaymentId(close_txid.to_byte_array())).unwrap(); + match &close.kind { + PaymentKind::Onchain { txid, status, .. } => { + assert_eq!(*txid, close_txid); + assert!(matches!(status, ConfirmationStatus::Confirmed { .. })); + }, + kind => panic!("unexpected kind {:?}", kind), + } + } + /// A funding-typed broadcast that doesn't touch the on-chain wallet must not be recorded. /// LDK re-broadcasts a promoted-but-unconfirmed 0conf splice through its generic funding /// path, so a splice the interactive-funding classification deliberately declined — no local @@ -4071,6 +4210,87 @@ mod tests { assert_unchanged(true); } + /// A funding broadcast whose classification fails must be retried, not dropped: for + /// interactive funding the counterparty broadcasts the same transaction regardless of + /// whether we do, so dropping the package permanently leaves the confirming transaction + /// unrecorded as a candidate — and the funding-status ownership gate then routes its + /// confirmation to a stray duplicate record instead of the funding record. + #[tokio::test] + async fn failed_funding_classification_is_retried_not_dropped() { + use lightning::chain::chaininterface::BroadcasterInterface; + + let fail_store = FailSwitchStore::new(); + let store: Arc = Arc::new(DynStoreWrapper(fail_store.clone())); + let wallet = new_test_wallet(Arc::clone(&store), false).await; + wallet.broadcaster.set_wallet(Arc::downgrade(&wallet)); + + // Run the production broadcast-queue loop. The broadcast itself fails fast against the + // fixture's unroutable Esplora server, which is irrelevant here: the record is written + // during classification, before the broadcast attempt. + let (stop_sender, stop_receiver) = tokio::sync::watch::channel(()); + let chain_source = Arc::clone(&wallet.chain_source); + let loop_task = tokio::spawn(async move { + chain_source.continuously_process_broadcast_queue(stop_receiver).await + }); + + // A funding transaction paying the wallet passes the wallet-activity guard, so its + // classification reaches the payment-store write. + let script_pubkey = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + let tx = Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: LockTime::ZERO, + input: Vec::new(), + output: vec![TxOut { value: Amount::from_sat(10_000), script_pubkey }], + }; + let counterparty_node_id = PublicKey::from_str( + "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(); + + // Queue the broadcast while payment persistence is failing. + fail_store.fail_writes.store(true, Ordering::Release); + wallet.broadcaster.broadcast_transactions(&[( + &tx, + LdkTransactionType::Funding { + channels: vec![(counterparty_node_id, ChannelId([7u8; 32]))], + }, + )]); + + // Let the loop fail at least one classification round; a failed classification must not + // leave a partial record behind. + tokio::time::sleep(Duration::from_secs(3)).await; + assert!(wallet.payment_store.list_filter(|_| true).is_empty()); + + // Once writes recover, the package must still be alive to classify. + fail_store.fail_writes.store(false, Ordering::Release); + let mut recorded = Vec::new(); + for _ in 0..100 { + tokio::time::sleep(Duration::from_millis(100)).await; + recorded = wallet.payment_store.list_filter(|_| true); + if !recorded.is_empty() { + break; + } + } + assert!( + !recorded.is_empty(), + "the failed classification was never retried; the package was dropped" + ); + assert_eq!(recorded.len(), 1); + assert!(matches!( + recorded[0].kind, + PaymentKind::Onchain { tx_type: Some(TransactionType::Funding { .. }), .. } + )); + + stop_sender.send(()).unwrap(); + loop_task.await.unwrap(); + } + /// Barrier test, classification-first ordering: wallet sync's confirmation handling must /// wait for classification's two-store write pair. Classification is parked between its /// payment-store and pending-store writes (the torn window) and only then is the @@ -4214,4 +4434,105 @@ mod tests { PaymentKind::Onchain { tx_type: Some(TransactionType::InteractiveFunding { .. }), .. } )); } + + /// An on-chain RBF bump must not replace a record a concurrent classification re-types as + /// channel funding: the replacement would double-spend the channel's funding transaction. + /// The bump is parked on the persister lock and classification lands while it waits; unless + /// the bump holds the cross-store lock from its funding-kind check through its writes, it + /// proceeds on the stale pre-classification read and retargets the funding record to the + /// replacement it broadcasts. + #[tokio::test] + async fn fee_bump_waits_for_funding_classification() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + // A confirmed parent funds the wallet so it can build and sign a replaceable spend. The + // checkpoint and anchor make the parent canonical: fee bumping resolves the spent + // prevouts and signing reads the full parent transaction from the graph. + let parent_spk = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + let parent = Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: LockTime::ZERO, + input: Vec::new(), + output: vec![TxOut { value: Amount::from_sat(100_000), script_pubkey: parent_spk }], + }; + let parent_txid = parent.compute_txid(); + { + let mut locked_wallet = wallet.inner.lock().unwrap(); + let block_id = + BlockId { height: 100, hash: bitcoin::BlockHash::from_byte_array([1u8; 32]) }; + let chain = locked_wallet.latest_checkpoint().insert(block_id); + let mut tx_update = bdk_chain::TxUpdate::default(); + tx_update.txs = vec![Arc::new(parent)]; + tx_update.anchors = + [(ConfirmationBlockTime { block_id, confirmation_time: 0 }, parent_txid)].into(); + let update = bdk_wallet::Update { chain: Some(chain), tx_update, ..Default::default() }; + locked_wallet.apply_update(update).unwrap(); + } + + // The wallet builds and signs the replaceable transaction itself, which keeps it + // RBF-signaling and its change output claimable for the extra fee. + let tx = { + let mut locked_wallet = wallet.inner.lock().unwrap(); + let foreign_spk = + ScriptBuf::new_p2wpkh(&bitcoin::WPubkeyHash::from_byte_array([0xab; 20])); + let mut builder = locked_wallet.build_tx(); + builder.add_recipient(foreign_spk, Amount::from_sat(20_000)); + let mut psbt = builder.finish().unwrap(); + assert!(locked_wallet.sign(&mut psbt, SignOptions::default()).unwrap()); + psbt.extract_tx().unwrap() + }; + let txid = tx.compute_txid(); + let payment_id = PaymentId(txid.to_byte_array()); + + // Observing the spend in the mempool mints the plain on-chain record — the same state an + // interactively funded round observed by wallet sync before classification leaves behind. + wallet.apply_mempool_txs(vec![(tx, 1_000)], Vec::new()).await.unwrap(); + let seeded = wallet.payment_store.get(&payment_id).expect("record for the mempool tx"); + assert!(matches!(&seeded.kind, PaymentKind::Onchain { tx_type: None, .. })); + assert_eq!(seeded.direction, PaymentDirection::Outbound); + + // Park the bump on the persister lock and classify while it waits. Classification + // serializes on the cross-store lock, not the persister, so it runs to completion while + // the bump is parked. Polling the bump first runs it synchronously to its first await: + // with an unlocked gate that is the persister acquisition after the funding-kind check, + // so the stale decision is already made; code that takes the locks before reading parks + // before the read, so either arrival order converges on the same final state. + let persister_guard = wallet.persister.lock().await; + let candidates = vec![FundingTxCandidate { + txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }]; + let details = interactive_funding_details(payment_id, txid, Some(1_000_000), Some(500)); + let bump = wallet.bump_fee_rbf(payment_id, None, 0); + let classify = async { + wallet.persist_funding_payment(details, candidates).await.unwrap(); + drop(persister_guard); + }; + let (result, ()) = tokio::join!(bump, classify); + + assert!(result.is_err(), "an RBF bump replaced a freshly-classified funding record"); + let payments = wallet.payment_store.list_filter(|_| true); + assert_eq!(payments.len(), 1); + let payment = &payments[0]; + assert_eq!(payment.amount_msat, Some(1_000_000)); + assert_eq!(payment.fee_paid_msat, Some(500)); + match &payment.kind { + PaymentKind::Onchain { + txid: current, + tx_type: Some(TransactionType::InteractiveFunding { .. }), + .. + } => { + assert_eq!(*current, txid, "the funding record must keep its own transaction"); + }, + kind => panic!("unexpected kind {:?}", kind), + } + } }