diff --git a/src/chain/mod.rs b/src/chain/mod.rs index f01c1c8cb..2d53cf9d6 100644 --- a/src/chain/mod.rs +++ b/src/chain/mod.rs @@ -37,9 +37,15 @@ use crate::config::{BackgroundSyncConfig, Config, WALLET_SYNC_INTERVAL_MINIMUM_S use crate::fee_estimator::OnchainFeeEstimator; use crate::logger::{log_debug, log_error, log_info, log_trace, LdkLogger, Logger}; use crate::runtime::Runtime; +use crate::tx_broadcaster::{BroadcastPackage, RetryQueue, ScheduleOutcome}; use crate::types::{Broadcaster, ChainMonitor, ChannelManager, DynStore, Sweeper, Wallet}; use crate::{Error, PersistedNodeMetrics}; +/// 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); + /// We use this parent-child TRUC package to make sure the configured chain source supports /// broadcasting packages via the `submitpackage` Bitcoin Core RPC. const PARENT_TXID: &str = "9a015f93fac6cb203c2b994e18b85176eb0354a22a468255516f3c6002d3f696"; @@ -562,50 +568,91 @@ impl ChainSource { } } + /// Classifies the package's funding broadcasts into payment records, then broadcasts it. + /// Returns the package back on classification failure so the caller can retry it after a + /// delay: 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. + async fn classify_and_broadcast( + &self, package: BroadcastPackage, + ) -> Result<(), BroadcastPackage> { + if let Err(e) = self.tx_broadcaster.classify_package(&package).await { + log_error!( + self.logger, + "Delaying broadcast: failed to persist payment records, will retry: {:?}", + e, + ); + return Err(package); + } + let package = package.into_sorted_transactions(); + match &self.kind { + #[cfg(feature = "chain-esplora")] + ChainSourceKind::Esplora(esplora_chain_source) => { + esplora_chain_source.process_transaction_broadcast(package).await + }, + #[cfg(feature = "chain-electrum")] + ChainSourceKind::Electrum(electrum_chain_source) => { + electrum_chain_source.process_transaction_broadcast(package).await + }, + #[cfg(feature = "chain-bitcoind")] + ChainSourceKind::Bitcoind(bitcoind_chain_source) => { + bitcoind_chain_source.process_transaction_broadcast(package).await + }, + } + Ok(()) + } + pub(crate) async fn continuously_process_broadcast_queue( &self, mut stop_tx_bcast_receiver: tokio::sync::watch::Receiver<()>, ) { let mut receiver = self.tx_broadcaster.get_broadcast_queue().await; + // Packages whose classification failed, each waiting out FAILED_CLASSIFY_RETRY_DELAY + // before its next attempt. New packages keep flowing while these wait, and pending + // retries die with the loop on shutdown rather than resurfacing after a later start. + let mut retries = RetryQueue::new(); loop { - let tx_bcast_logger = Arc::clone(&self.logger); - tokio::select! { + let next_retry_at = retries.next_retry_at(); + let package = tokio::select! { _ = stop_tx_bcast_receiver.changed() => { log_debug!( - tx_bcast_logger, + self.logger, "Stopping broadcasting transactions.", ); return; } - 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(); - match &self.kind { - #[cfg(feature = "chain-esplora")] - ChainSourceKind::Esplora(esplora_chain_source) => { - esplora_chain_source.process_transaction_broadcast(package).await - }, - #[cfg(feature = "chain-electrum")] - ChainSourceKind::Electrum(electrum_chain_source) => { - electrum_chain_source.process_transaction_broadcast(package).await - }, - #[cfg(feature = "chain-bitcoind")] - ChainSourceKind::Bitcoind(bitcoind_chain_source) => { - bitcoind_chain_source.process_transaction_broadcast(package).await - }, - } + Some(next_package) = receiver.recv() => next_package, + _ = tokio::time::sleep_until( + next_retry_at.unwrap_or_else(tokio::time::Instant::now) + ), if next_retry_at.is_some() => { + retries.pop_next().expect("a retry is queued") + } + }; + if let Err(package) = self.classify_and_broadcast(package).await { + let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY; + match retries.schedule(package, retry_at) { + ScheduleOutcome::Scheduled { dropped: None } => {}, + ScheduleOutcome::Scheduled { dropped: Some(dropped) } => { + log_error!( + self.logger, + "Dropped the oldest package awaiting a classification retry; LDK re-broadcasts its transactions periodically: {:?}", + dropped.sorted_txids(), + ); + }, + ScheduleOutcome::AlreadyQueued(duplicate) => { + log_debug!( + self.logger, + "Dropped a re-broadcast package; an identical one already awaits a classification retry: {:?}", + duplicate.sorted_txids(), + ); + }, + ScheduleOutcome::Refused(package) => { + log_error!( + self.logger, + "Dropped a package failing classification; too many await retries: {:?}", + package.sorted_txids(), + ); + }, } } } diff --git a/src/payment/pending_payment_store.rs b/src/payment/pending_payment_store.rs index 30a113537..e1f1d4f2a 100644 --- a/src/payment/pending_payment_store.rs +++ b/src/payment/pending_payment_store.rs @@ -5,9 +5,13 @@ // http://opensource.org/licenses/MIT>, at your option. You may not use this file except in // accordance with one or both of these licenses. -use bitcoin::Txid; -use lightning::impl_writeable_tlv_based; +use bitcoin::secp256k1::PublicKey; +use bitcoin::{TxOut, Txid}; +use lightning::chain::transaction::OutPoint; use lightning::ln::channelmanager::PaymentId; +use lightning::ln::funding::FundingContribution; +use lightning::ln::types::ChannelId; +use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum}; use crate::data_store::{StorableObject, StorableObjectUpdate}; use crate::payment::store::PaymentDetailsUpdate; @@ -36,37 +40,176 @@ impl_writeable_tlv_based!(FundingTxCandidate, { (4, fee_paid_msat, option), }); -/// Represents a pending payment +/// The parameters of the API call that initiated a splice, recording what was attempted +/// independently of the contribution built from them. #[derive(Clone, Debug, PartialEq, Eq)] -pub struct PendingPaymentDetails { - /// The full payment details - pub details: PaymentDetails, - /// Transaction IDs that have replaced or conflict with this payment. - pub conflicting_txids: Vec, - /// For interactive funding (splices), this node's per-candidate funding figures across the - /// RBF history, keyed by each candidate's txid. Empty for non-funding payments and for - /// records written before per-candidate tracking existed. - pub(crate) candidates: Vec, +pub(crate) enum SpliceKind { + /// [`Node::splice_in`] with a resolved amount. + /// + /// [`Node::splice_in`]: crate::Node::splice_in + In { amount_sats: u64 }, + /// [`Node::splice_out`] to the given outputs. + /// + /// [`Node::splice_out`]: crate::Node::splice_out + Out { outputs: Vec }, + /// [`Node::bump_channel_funding_fee`] of a pending splice. + /// + /// [`Node::bump_channel_funding_fee`]: crate::Node::bump_channel_funding_fee + Rbf {}, +} + +impl_writeable_tlv_based_enum!(SpliceKind, + (0, In) => { + (0, amount_sats, required), + }, + (2, Out) => { + (0, outputs, required_vec), + }, + (4, Rbf) => {}, +); + +/// A user-initiated splice that has been handed to LDK but is not yet guaranteed to survive a +/// restart. LDK only persists a splice once its negotiation reaches `AwaitingSignatures`, and it +/// abandons an in-progress negotiation whenever the peer disconnects (which includes stopping the +/// node). Until the new funding transaction locks we keep enough state to recognize a splice LDK +/// no longer knows about and to describe events about it in terms of the original request. +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) struct SpliceIntent { + /// The channel counterparty. + pub counterparty_node_id: PublicKey, + /// The channel being spliced. + pub channel_id: ChannelId, + /// The channel's funding outpoint when the splice was initiated. It only changes once a splice + /// locks, so a mismatch with the channel's current funding outpoint means the splice (or a + /// replacement) completed and the intent is stale. + pub pre_splice_funding_txo: OutPoint, + /// The contribution handed to [`ChannelManager::funding_contributed`], kept to match later + /// events about the splice back to this intent. + /// + /// [`ChannelManager::funding_contributed`]: lightning::ln::channelmanager::ChannelManager::funding_contributed + pub contribution: FundingContribution, + /// The parameters of the originating API call. + pub kind: SpliceKind, +} + +impl_writeable_tlv_based!(SpliceIntent, { + (0, counterparty_node_id, required), + (2, channel_id, required), + (4, pre_splice_funding_txo, required), + (6, contribution, required), + (8, kind, required), +}); + +/// A pending payment tracked by LDK Node, keyed by [`PaymentId`]. +/// +/// A user-initiated splice is persisted as a [`PendingSplice`] before its contribution is handed +/// to LDK — at which point no funding transaction, and therefore no [`PaymentDetails`], exists yet. +/// Once the splice is broadcast and classified it becomes a [`Tracked`] payment carrying the real +/// [`PaymentDetails`], while retaining its [`SpliceIntent`] until the splice locks. +/// +/// [`PendingSplice`]: Self::PendingSplice +/// [`Tracked`]: Self::Tracked +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) enum PendingPaymentDetails { + /// A user-initiated splice persisted before hand-off to LDK; no funding transaction exists yet. + /// Keyed by the generated [`PaymentId`]; never mirrored into the payment store. + PendingSplice { id: PaymentId, intent: SpliceIntent }, + /// A pending payment tracked toward confirmation, optionally still carrying a live splice + /// intent until the splice locks. + /// + /// Each field is written by a different subsystem: wallet sync records `conflicting_txids` + /// for any wallet transaction (splice fundings included), broadcast-time classification + /// records `candidates` for interactive funding, and `splice_intent` is carried over from a + /// [`PendingSplice`] record when the payment is promoted — nothing persists an intent at + /// splice initiation yet; that lands with the splice tracking built on this. A splice uses + /// all of them; the fields do not partition by payment type. + /// + /// [`PendingSplice`]: Self::PendingSplice + Tracked { + /// The full payment details. + details: PaymentDetails, + /// Transaction IDs wallet sync observed to have replaced or to conflict with this + /// payment, used to map later events about those txids back to this record. This is + /// BDK's view, distinct from `candidates`: it can hold conflicts that were never + /// negotiated candidates, while a candidate replaced between wallet syncs may never + /// appear here (it gets no `TxReplaced` event of its own). + conflicting_txids: Vec, + /// For interactive funding (splices), this node's per-candidate funding figures across the + /// RBF history, keyed by each candidate's txid and recorded as each round's broadcast is + /// classified. Empty for non-funding payments. + candidates: Vec, + /// The live splice intent, or `None` for a non-splice payment or a splice that has + /// locked. It lives here as well as on + /// [`PendingSplice`] because a fee bump — a fresh negotiation LDK likewise abandons if the + /// peer disconnects before signing — would share the broadcast splice's record rather than + /// get one of its own. + /// + /// [`PendingSplice`]: Self::PendingSplice + splice_intent: Option, + }, } impl PendingPaymentDetails { pub(crate) fn new( details: PaymentDetails, conflicting_txids: Vec, candidates: Vec, ) -> Self { - Self { details, conflicting_txids, candidates } + Self::tracked(details, conflicting_txids, candidates, None) + } + + pub(crate) fn tracked( + details: PaymentDetails, conflicting_txids: Vec, candidates: Vec, + splice_intent: Option, + ) -> Self { + Self::Tracked { details, conflicting_txids, candidates, splice_intent } + } + + /// The full payment details, or `None` for a splice not yet broadcast. + pub(crate) fn details(&self) -> Option<&PaymentDetails> { + match self { + Self::PendingSplice { .. } => None, + Self::Tracked { details, .. } => Some(details), + } + } + + /// Transaction IDs that have replaced or conflict with this payment. + pub(crate) fn conflicting_txids(&self) -> &[Txid] { + match self { + Self::PendingSplice { .. } => &[], + Self::Tracked { conflicting_txids, .. } => conflicting_txids, + } + } + + /// The splice intent this record carries, if it is a splice that has not yet locked. + pub(crate) fn splice_intent(&self) -> Option<&SpliceIntent> { + match self { + Self::PendingSplice { intent, .. } => Some(intent), + Self::Tracked { splice_intent, .. } => splice_intent.as_ref(), + } } /// Returns this node's recorded funding figures for the candidate with the given txid, if any. pub(crate) fn candidate(&self, txid: Txid) -> Option<&FundingTxCandidate> { - self.candidates.iter().find(|candidate| candidate.txid == txid) + match self { + Self::PendingSplice { .. } => None, + Self::Tracked { candidates, .. } => { + candidates.iter().find(|candidate| candidate.txid == txid) + }, + } } } -impl_writeable_tlv_based!(PendingPaymentDetails, { - (0, details, required), - (2, conflicting_txids, optional_vec), - (4, candidates, optional_vec), -}); +impl_writeable_tlv_based_enum!(PendingPaymentDetails, + (0, PendingSplice) => { + (0, id, required), + (2, intent, required), + }, + (2, Tracked) => { + (0, details, required), + (2, conflicting_txids, optional_vec), + (4, candidates, optional_vec), + (6, splice_intent, option), + }, +); #[derive(Clone, Debug, PartialEq, Eq)] pub(crate) struct PendingPaymentDetailsUpdate { @@ -74,6 +217,10 @@ pub(crate) struct PendingPaymentDetailsUpdate { pub payment_update: Option, pub conflicting_txids: Option>, pub candidates: Vec, + /// The splice intent to set (`Some(Some(..))`) or clear (`Some(None)`), or `None` to leave it + /// unchanged. Setting it on a [`PendingPaymentDetails::PendingSplice`] replaces the intent; + /// clearing a pre-broadcast splice is done by removing the record, not through this field. + pub splice_intent: Option>, } impl StorableObject for PendingPaymentDetails { @@ -81,38 +228,73 @@ impl StorableObject for PendingPaymentDetails { type Update = PendingPaymentDetailsUpdate; fn id(&self) -> Self::Id { - self.details.id + match self { + Self::PendingSplice { id, .. } => *id, + Self::Tracked { details, .. } => details.id, + } } fn update(&mut self, update: Self::Update) -> bool { - let mut updated = false; - - // Update the underlying payment details if present - if let Some(payment_update) = update.payment_update { - updated |= self.details.update(payment_update); - } - - if let Some(new_conflicting_txids) = update.conflicting_txids { - if self.conflicting_txids != new_conflicting_txids { - self.conflicting_txids = new_conflicting_txids; - updated = true; - } - } - - if let PaymentKind::Onchain { txid, .. } = &self.details.kind { - let conflicts_len = self.conflicting_txids.len(); - self.conflicting_txids.retain(|conflicting_txid| conflicting_txid != txid); - updated |= self.conflicting_txids.len() != conflicts_len; - } - - // Each classify passes the complete candidate history, so a non-empty update replaces the - // stored list. An empty update (e.g. a non-funding payment) leaves it untouched. - if !update.candidates.is_empty() && self.candidates != update.candidates { - self.candidates = update.candidates; - updated = true; + match self { + Self::PendingSplice { intent, .. } => { + // A pre-broadcast record only carries a splice intent; the only meaningful update is + // replacing that intent. Clearing it is done by removing the record. + if let Some(Some(new_intent)) = update.splice_intent { + if *intent != new_intent { + *intent = new_intent; + return true; + } + } + false + }, + Self::Tracked { details, conflicting_txids, candidates, splice_intent } => { + let mut updated = false; + + // Update the underlying payment details if present + if let Some(payment_update) = update.payment_update { + updated |= details.update(payment_update); + } + + if let Some(new_conflicting_txids) = update.conflicting_txids { + if *conflicting_txids != new_conflicting_txids { + *conflicting_txids = new_conflicting_txids; + updated = true; + } + } + + if let PaymentKind::Onchain { txid, .. } = &details.kind { + let conflicts_len = conflicting_txids.len(); + conflicting_txids.retain(|conflicting_txid| conflicting_txid != txid); + updated |= conflicting_txids.len() != conflicts_len; + } + + // Each classify passes the candidate history as of its own broadcast, so a + // non-empty update replaces the stored list. An empty update (e.g. a non-funding + // payment) leaves it untouched — as does an update missing a stored candidate: + // the history only ever grows, so such an update was built before that candidate + // existed (a classification retry running after a newer round classified) and + // replacing would orphan the newer round's transactions. + let extends_history = |stored: &FundingTxCandidate| { + update.candidates.iter().any(|candidate| candidate.txid == stored.txid) + }; + if !update.candidates.is_empty() + && *candidates != update.candidates + && candidates.iter().all(extends_history) + { + *candidates = update.candidates; + updated = true; + } + + if let Some(new_splice_intent) = update.splice_intent { + if *splice_intent != new_splice_intent { + *splice_intent = new_splice_intent; + updated = true; + } + } + + updated + }, } - - updated } fn to_update(&self) -> Self::Update { @@ -128,23 +310,65 @@ impl StorableObjectUpdate for PendingPaymentDetailsUpdate impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate { fn from(value: &PendingPaymentDetails) -> Self { - let conflicting_txids = if value.conflicting_txids.is_empty() { - None - } else { - Some(value.conflicting_txids.clone()) - }; - Self { - id: value.id(), - payment_update: Some(value.details.to_update()), - conflicting_txids, - candidates: value.candidates.clone(), + match value { + PendingPaymentDetails::PendingSplice { id, intent } => Self { + id: *id, + payment_update: None, + conflicting_txids: None, + candidates: Vec::new(), + splice_intent: Some(Some(intent.clone())), + }, + PendingPaymentDetails::Tracked { + details, + conflicting_txids, + candidates, + splice_intent, + } => { + let conflicting_txids = if conflicting_txids.is_empty() { + None + } else { + Some(conflicting_txids.clone()) + }; + Self { + id: details.id, + payment_update: Some(details.to_update()), + conflicting_txids, + candidates: candidates.clone(), + splice_intent: Some(splice_intent.clone()), + } + }, } } } +/// Builds a [`FundingContribution`] for tests through its `Readable` impl — the only path open +/// outside `rust-lightning`, which keeps its builder private. The length-prefixed stream holds +/// just the required TLV records: estimated fee, feerate, max feerate, and the is-splice flag. +#[cfg(test)] +pub(crate) fn test_funding_contribution() -> FundingContribution { + test_funding_contribution_with_feerate(253) +} + +/// Like [`test_funding_contribution`], but with the given input-selection feerate in sat/kwu. +#[cfg(test)] +pub(crate) fn test_funding_contribution_with_feerate(feerate: u64) -> FundingContribution { + let mut tlv_bytes = vec![ + 33u8, // BigSize length prefix over the TLV records below + 1, 8, 0, 0, 0, 0, 0, 0, 0, 0, // (1, estimated_fee: 0 sat) + 9, 8, // (9, feerate) + ]; + tlv_bytes.extend_from_slice(&feerate.to_be_bytes()); + tlv_bytes.extend_from_slice(&[11, 8]); // (11, max_feerate) + tlv_bytes.extend_from_slice(&feerate.to_be_bytes()); + tlv_bytes.extend_from_slice(&[13, 1, 1]); // (13, is_splice: true) + lightning::util::ser::Readable::read(&mut &tlv_bytes[..]) + .expect("hand-built TLV stream must decode") +} + #[cfg(test)] mod tests { use bitcoin::hashes::Hash; + use lightning::util::ser::{Readable, Writeable}; use super::*; use crate::payment::store::ConfirmationStatus; @@ -237,12 +461,61 @@ mod tests { assert!(pending_payment.update(update)); assert_eq!( - pending_payment.conflicting_txids, + pending_payment.conflicting_txids(), Vec::::new(), "current txid must not remain in its own conflict list" ); } + /// The candidate history only ever grows. An update carrying a shorter history was built + /// before the newer candidates existed — a classification retry running after a newer round + /// classified — and must not shrink the stored list, or the newer candidates' transactions + /// could no longer be mapped back to the record. + #[test] + fn candidate_history_never_shrinks() { + let txid_a = test_txid(1); + let txid_b = test_txid(2); + let txid_c = test_txid(3); + let payment_id = PaymentId(txid_a.to_byte_array()); + let candidate = |txid, fee| FundingTxCandidate { + txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(fee), + }; + let history = vec![candidate(txid_a, 400), candidate(txid_b, 500)]; + + let mut pending = PendingPaymentDetails::new( + pending_onchain_payment(payment_id, txid_b), + Vec::new(), + history.clone(), + ); + let stored_candidates = |pending: &PendingPaymentDetails| match pending { + PendingPaymentDetails::Tracked { candidates, .. } => candidates.clone(), + pending => panic!("unexpected variant {:?}", pending), + }; + let stale_update = PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: None, + candidates: vec![candidate(txid_a, 400)], + splice_intent: None, + }; + assert!(!pending.update(stale_update), "a stale history must not shrink the stored one"); + assert_eq!(stored_candidates(&pending), history); + + // A history that extends the stored one still replaces it, refreshed figures included. + let extended = vec![candidate(txid_a, 400), candidate(txid_b, 550), candidate(txid_c, 600)]; + let fresh_update = PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: None, + candidates: extended.clone(), + splice_intent: None, + }; + assert!(pending.update(fresh_update)); + assert_eq!(stored_candidates(&pending), extended); + } + #[test] fn funding_classification_pending_update_preserves_mirrored_confirmation() { use bitcoin::BlockHash; @@ -289,7 +562,7 @@ mod tests { assert!(downgraded.update(full_update)); assert!( matches!( - downgraded.details.kind, + downgraded.details().expect("tracked").kind, PaymentKind::Onchain { status: ConfirmationStatus::Unconfirmed, .. } ), "a full merge of a fresh classification downgrades a mirrored confirmation", @@ -304,17 +577,80 @@ mod tests { payment_update: Some(PaymentDetailsUpdate::funding_reclassification(fresh)), conflicting_txids: None, candidates: candidates.clone(), + splice_intent: None, }; assert!(merged.update(narrow_update)); + let merged_details = merged.details().expect("tracked"); assert!( matches!( - merged.details.kind, + merged_details.kind, PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { .. }, .. } ), "a narrow classification update must not downgrade a mirrored confirmation", ); - assert_eq!(merged.candidates, candidates); - assert_eq!(merged.details.amount_msat, Some(1_000)); - assert_eq!(merged.details.fee_paid_msat, Some(100)); + assert_eq!(merged.candidate(txid), Some(&candidates[0])); + assert_eq!(merged_details.amount_msat, Some(1_000)); + assert_eq!(merged_details.fee_paid_msat, Some(100)); + } + + #[test] + fn splice_kind_round_trips() { + for kind in [ + SpliceKind::In { amount_sats: 500_000 }, + SpliceKind::Out { + outputs: vec![TxOut { + value: bitcoin::Amount::from_sat(400_000), + script_pubkey: bitcoin::ScriptBuf::new(), + }], + }, + SpliceKind::Rbf {}, + ] { + let encoded = kind.encode(); + let decoded = SpliceKind::read(&mut &encoded[..]).unwrap(); + assert_eq!(kind, decoded); + } + } + + #[test] + fn pending_splice_round_trips() { + use std::str::FromStr; + + let id = PaymentId([10u8; 32]); + let intent = SpliceIntent { + counterparty_node_id: PublicKey::from_str( + "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(), + channel_id: ChannelId([11u8; 32]), + pre_splice_funding_txo: OutPoint { txid: test_txid(12), index: 0 }, + contribution: test_funding_contribution(), + kind: SpliceKind::In { amount_sats: 500_000 }, + }; + let record = PendingPaymentDetails::PendingSplice { id, intent }; + + let encoded = record.encode(); + let decoded = PendingPaymentDetails::read(&mut &encoded[..]).unwrap(); + assert_eq!(record, decoded); + assert_eq!(decoded.id(), id); + assert!(decoded.details().is_none()); + } + + #[test] + fn tracked_payment_round_trips() { + // The `PendingSplice` variant round-trips in `pending_splice_round_trips`; here we cover the + // `Tracked` variant and its enum discriminant. + let payment_id = PaymentId([7u8; 32]); + let txid = Txid::from_byte_array([8u8; 32]); + let record = PendingPaymentDetails::new( + pending_onchain_payment(payment_id, txid), + vec![Txid::from_byte_array([9u8; 32])], + vec![FundingTxCandidate { txid, amount_msat: Some(1_000), fee_paid_msat: Some(100) }], + ); + + let encoded = record.encode(); + let decoded = PendingPaymentDetails::read(&mut &encoded[..]).unwrap(); + assert_eq!(record, decoded); + assert_eq!(decoded.id(), payment_id); + assert!(decoded.details().is_some()); } } diff --git a/src/tx_broadcaster.rs b/src/tx_broadcaster.rs index 782112dad..da283dd5d 100644 --- a/src/tx_broadcaster.rs +++ b/src/tx_broadcaster.rs @@ -5,14 +5,16 @@ // http://opensource.org/licenses/MIT>, at your option. You may not use this file except in // accordance with one or both of these licenses. +use std::collections::VecDeque; use std::ops::Deref; use std::sync::{Mutex as StdMutex, Weak}; -use bitcoin::Transaction; +use bitcoin::{Transaction, Txid}; use lightning::chain::chaininterface::{ BroadcasterInterface, TransactionType as LdkTransactionType, }; use tokio::sync::{mpsc, Mutex, MutexGuard}; +use tokio::time::Instant; use crate::logger::{log_error, LdkLogger}; use crate::types::Wallet; @@ -20,6 +22,13 @@ use crate::Error; const BCAST_PACKAGE_QUEUE_SIZE: usize = 256; +/// The most non-funding packages [`RetryQueue`] holds. Claims and sweeps re-enter the +/// broadcast queue on LDK's periodic rebroadcast timers, so one dropped here resurfaces on its +/// own once the store recovers. Funding packages don't count against the bound: nothing +/// re-broadcasts them for us, and they are finite — one per negotiated candidate, since a copy +/// of a waiting package is never queued twice. +const MAX_QUEUED_RETRIES: usize = BCAST_PACKAGE_QUEUE_SIZE; + /// 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 @@ -47,6 +56,100 @@ impl BroadcastPackage { let txs = self.0.into_iter().map(|(tx, _)| tx).collect(); SortedTransactions::sort_parents_child_package_topologically(txs) } + + /// The packaged transactions' txids in sorted order, identifying the package's effect on + /// chain: two packages with the same txids broadcast the same transactions. + pub(crate) fn sorted_txids(&self) -> Vec { + let mut txids: Vec = self.0.iter().map(|(tx, _)| tx.compute_txid()).collect(); + txids.sort_unstable(); + txids + } + + /// Whether the package contains a funding transaction (a channel open or splice), whose + /// classification writes the payment record tracking the funding. + fn contains_funding(&self) -> bool { + self.0.iter().any(|(_, tx_type)| { + matches!( + tx_type, + Some( + LdkTransactionType::Funding { .. } + | LdkTransactionType::InteractiveFunding { .. } + ) + ) + }) + } +} + +/// What [`RetryQueue::schedule`] did with a package, so the caller can log the cases in which +/// the package won't be retried as-is. +pub(crate) enum ScheduleOutcome { + /// The package waits for its retry deadline. When the bound was reached, the oldest waiting + /// non-funding package was dropped to make room and is returned — its transactions resurface + /// with LDK's next periodic rebroadcast. + Scheduled { dropped: Option }, + /// A package broadcasting the same transactions already waits, and its retry covers this + /// one: the incoming package is dropped and returned. + AlreadyQueued(BroadcastPackage), + /// The bound was reached and every waiting package is a funding package, which must not be + /// dropped: the incoming package is refused and returned. + Refused(BroadcastPackage), +} + +/// Packages whose classification failed, each waiting out a retry delay before its next attempt. +/// Deduplicated and bounded: LDK re-broadcasts pending claims every 30 seconds (and sweeps once +/// per block) until they confirm, so while the store is unavailable, copies would otherwise +/// accumulate without bound and replay as a burst on recovery. An identical copy is never queued +/// twice — the waiting entry and its deadline stand; fee-bumped rebroadcast variants carry new +/// txids, so the bound — not the dedup — is what limits their accumulation. +pub(crate) struct RetryQueue(VecDeque<(Instant, Vec, BroadcastPackage)>); + +impl RetryQueue { + pub(crate) fn new() -> Self { + Self(VecDeque::new()) + } + + /// The deadline of the next retry, if a package is waiting. Packages are scheduled with a fixed + /// delay, so the front entry is always the next to retry. + pub(crate) fn next_retry_at(&self) -> Option { + self.0.front().map(|(deadline, _, _)| *deadline) + } + + /// Removes and returns the package scheduled to retry first. + pub(crate) fn pop_next(&mut self) -> Option { + self.0.pop_front().map(|(_, _, package)| package) + } + + /// Schedules a package to retry at `retry_at`, unless a package with the same transactions already + /// waits or accepting it would exceed [`MAX_QUEUED_RETRIES`] with no non-funding package to + /// drop for it; see [`ScheduleOutcome`]. + pub(crate) fn schedule( + &mut self, package: BroadcastPackage, retry_at: Instant, + ) -> ScheduleOutcome { + let txids = package.sorted_txids(); + if self.0.iter().any(|(_, waiting, _)| *waiting == txids) { + // Same transactions, same classification outcome: keep the waiting entry and its + // earlier deadline. The one same-txid package with a *different* type is LDK's + // re-typed generic-funding rebroadcast of a promoted 0conf splice, which always + // arrives after the interactive-funding original (the zero-conf rebroadcast canary + // tests assert that ordering), so the entry kept is the richer of the two — and its + // classification declines the downgrade anyway. + return ScheduleOutcome::AlreadyQueued(package); + } + + let mut dropped = None; + if !package.contains_funding() && self.0.len() >= MAX_QUEUED_RETRIES { + // Drop the oldest non-funding package: LDK re-broadcasts its transactions + // periodically, while the incoming package may carry a fresher fee-bumped variant. + // A funding package is never dropped — nothing would re-broadcast it, and losing it + // leaves its transaction confirming without a recorded candidate. + match self.0.iter().position(|(_, _, waiting)| !waiting.contains_funding()) { + Some(oldest) => dropped = self.0.remove(oldest).map(|(_, _, package)| package), + None => return ScheduleOutcome::Refused(package), + } + } + self.0.push_back((retry_at, txids, package)); + ScheduleOutcome::Scheduled { dropped } + } } pub(crate) struct SortedTransactions(Vec); @@ -133,12 +236,10 @@ 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 retry it later 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 +248,7 @@ where } } } - Ok(package) + Ok(()) } pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) { @@ -173,7 +274,10 @@ mod tests { use bitcoin::hashes::Hash; use bitcoin::{Amount, OutPoint, ScriptBuf, Sequence, Transaction, TxIn, TxOut, Txid, Witness}; - use super::SortedTransactions; + use super::{ + BroadcastPackage, LdkTransactionType, RetryQueue, ScheduleOutcome, SortedTransactions, + MAX_QUEUED_RETRIES, + }; fn txin(txid: Txid, vout: u32) -> TxIn { TxIn { @@ -314,4 +418,142 @@ mod tests { fn topological_sort_accepts_empty_vec() { SortedTransactions::sort_parents_child_package_topologically(Vec::new()); } + + fn funding_package(tx: &Transaction) -> BroadcastPackage { + BroadcastPackage::new(&[(tx, LdkTransactionType::Funding { channels: vec![] })]) + } + + fn deadline(secs: u64) -> tokio::time::Instant { + tokio::time::Instant::now() + std::time::Duration::from_secs(secs) + } + + /// A re-broadcast of the same transactions is not queued again: the waiting entry keeps its + /// earlier deadline and its package — the first arrival carries the richer classification + /// when LDK later re-types a rebroadcast. + #[tokio::test] + async fn retry_queue_queues_identical_transactions_once() { + let tx = parent_tx(1); + let mut retries = RetryQueue::new(); + + let first_deadline = deadline(2); + assert!(matches!( + retries.schedule(funding_package(&tx), first_deadline), + ScheduleOutcome::Scheduled { dropped: None } + )); + assert!(matches!( + retries.schedule(BroadcastPackage::unclassified(tx.clone()), deadline(4)), + ScheduleOutcome::AlreadyQueued(_) + )); + + assert_eq!(retries.next_retry_at(), Some(first_deadline)); + let kept = retries.pop_next().expect("the first package is kept"); + assert!( + matches!(kept.transactions()[0].1, Some(LdkTransactionType::Funding { .. })), + "the first-scheduled package must be kept" + ); + assert!(retries.pop_next().is_none()); + } + + #[tokio::test] + async fn retry_queue_retries_in_schedule_order() { + let (tx_a, tx_b) = (parent_tx(1), parent_tx(2)); + let mut retries = RetryQueue::new(); + + assert!(matches!( + retries.schedule(BroadcastPackage::unclassified(tx_a.clone()), deadline(2)), + ScheduleOutcome::Scheduled { dropped: None } + )); + assert!(matches!( + retries.schedule(BroadcastPackage::unclassified(tx_b.clone()), deadline(2)), + ScheduleOutcome::Scheduled { dropped: None } + )); + + let popped = retries.pop_next().expect("first package"); + assert_eq!(popped.sorted_txids(), vec![tx_a.compute_txid()]); + let popped = retries.pop_next().expect("second package"); + assert_eq!(popped.sorted_txids(), vec![tx_b.compute_txid()]); + } + + /// Distinct transactions (e.g. fee-bumped claim variants during a store outage) are held to + /// the bound: the oldest non-funding package is dropped for an incoming one, never a funding + /// package. + #[tokio::test] + async fn retry_queue_drops_the_oldest_non_funding_package_at_the_bound() { + fn numbered_tx(n: u32) -> Transaction { + Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: bitcoin::absolute::LockTime::ZERO, + input: vec![txin(Txid::from_byte_array([7u8; 32]), n)], + output: vec![txout(1_000)], + } + } + + let mut retries = RetryQueue::new(); + let funding_tx = numbered_tx(0); + assert!(matches!( + retries.schedule(funding_package(&funding_tx), deadline(2)), + ScheduleOutcome::Scheduled { dropped: None } + )); + let oldest_claim = numbered_tx(1); + for n in 1..(MAX_QUEUED_RETRIES as u32) { + assert!(matches!( + retries.schedule(BroadcastPackage::unclassified(numbered_tx(n)), deadline(2)), + ScheduleOutcome::Scheduled { dropped: None } + )); + } + + // At the bound, an incoming non-funding package drops the oldest waiting one — not the + // older funding package. + let new_claim = numbered_tx(MAX_QUEUED_RETRIES as u32); + match retries.schedule(BroadcastPackage::unclassified(new_claim.clone()), deadline(2)) { + ScheduleOutcome::Scheduled { dropped: Some(dropped) } => { + assert_eq!(dropped.sorted_txids(), vec![oldest_claim.compute_txid()]); + }, + _ => panic!("the incoming claim must be scheduled by dropping the oldest one"), + } + + // An incoming funding package is never dropped for the bound. + let new_funding_tx = numbered_tx(MAX_QUEUED_RETRIES as u32 + 1); + assert!(matches!( + retries.schedule(funding_package(&new_funding_tx), deadline(2)), + ScheduleOutcome::Scheduled { dropped: None } + )); + + let mut remaining = Vec::new(); + while let Some(package) = retries.pop_next() { + remaining.extend(package.sorted_txids()); + } + assert!(remaining.contains(&funding_tx.compute_txid()), "funding is never dropped"); + assert!(remaining.contains(&new_claim.compute_txid())); + assert!(!remaining.contains(&oldest_claim.compute_txid())); + } + + /// When only funding packages wait at the bound, an incoming non-funding package is refused: + /// LDK re-broadcasts claims and sweeps periodically, while a dropped funding package would + /// leave its transaction confirming without a recorded candidate. + #[tokio::test] + async fn retry_queue_refuses_a_non_funding_package_over_waiting_funding_packages() { + fn numbered_tx(n: u32) -> Transaction { + Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: bitcoin::absolute::LockTime::ZERO, + input: vec![txin(Txid::from_byte_array([8u8; 32]), n)], + output: vec![txout(1_000)], + } + } + + let mut retries = RetryQueue::new(); + for n in 0..(MAX_QUEUED_RETRIES as u32) { + assert!(matches!( + retries.schedule(funding_package(&numbered_tx(n)), deadline(2)), + ScheduleOutcome::Scheduled { dropped: None } + )); + } + + let claim = numbered_tx(MAX_QUEUED_RETRIES as u32); + assert!(matches!( + retries.schedule(BroadcastPackage::unclassified(claim), deadline(2)), + ScheduleOutcome::Refused(_) + )); + } } diff --git a/src/wallet/mod.rs b/src/wallet/mod.rs index b9c12b4a7..5e2320e69 100644 --- a/src/wallet/mod.rs +++ b/src/wallet/mod.rs @@ -12,6 +12,7 @@ use std::str::FromStr; use std::sync::{Arc, Mutex}; use bdk_chain::spk_client::{FullScanRequest, SyncRequest}; +use bdk_chain::ChainPosition; use bdk_wallet::descriptor::ExtendedDescriptor; use bdk_wallet::error::{BuildFeeBumpError, CreateTxError}; #[allow(deprecated)] @@ -160,9 +161,9 @@ pub(crate) struct Wallet { logger: Arc, pending_payment_store: Arc, // Serializes the writers that must observe the payment record and its pending-store entry - // (candidate history included) as one consistent unit: classification holds it across its - // two-store write pair, and wallet sync's event arms hold it from payment-id resolution - // through their last write. Without it, a confirmation landing between classification's two + // (candidate history included) as one consistent unit: classification and wallet sync's event + // arms each hold it from payment-id resolution through their last write (classification's + // being its two-store pair). Without it, a confirmation landing between classification's two // writes sees the record classified but the candidate history absent — resolving the wrong // payment id or stamping the confirmed candidate with another candidate's figures — and a // classification landing inside an arm's decision sequence gets overwritten by the arm's @@ -346,12 +347,12 @@ 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) .await? .unwrap_or_else(|| PaymentId(txid.to_byte_array())); - if self + match self .apply_funding_status_update_locked( &guard, payment_id, @@ -360,6 +361,23 @@ impl Wallet { ) .await? { + 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()); + }, + } + + // The fallback id belongs to a settled funding payment whose entry is gone: + // skip rather than resurrect it (see `has_funding_record`). + if self.has_funding_record(&payment_id).await? { + log_debug!( + self.logger, + "Skipping wallet event for transaction {} of a settled funding payment", + txid, + ); continue; } @@ -378,35 +396,40 @@ impl Wallet { self.payment_store.insert_or_update(payment.clone()).await?; if payment_status == PaymentStatus::Pending { - let pending_payment = - self.create_pending_payment_from_tx(payment, Vec::new()); - - self.pending_payment_store.insert_or_update(pending_payment).await?; + self.upsert_pending_payment(payment, Vec::new()).await?; } }, WalletEvent::ChainTipChanged { new_tip, .. } => { let pending_payments: Vec = self .pending_payment_store - .list_filter(|p| { - debug_assert!( - p.details.status == PaymentStatus::Pending, - "Non-pending payment {:?} found in pending store", - p.details.id, - ); - p.details.status == PaymentStatus::Pending - && matches!(p.details.kind, PaymentKind::Onchain { .. }) + .list_filter(|p| match p.details() { + // A pre-broadcast splice intent carries no payment yet and cannot graduate. + None => false, + Some(details) => { + debug_assert!( + details.status == PaymentStatus::Pending, + "Non-pending payment {:?} found in pending store", + details.id, + ); + details.status == PaymentStatus::Pending + && matches!(details.kind, PaymentKind::Onchain { .. }) + }, }) .await; let mut unconfirmed_outbound_txids: Vec = Vec::new(); for payment in pending_payments { - match payment.details.kind { + // The filter admits only Tracked funding payments. + let PendingPaymentDetails::Tracked { ref details, .. } = payment else { + continue; + }; + match details.kind { PaymentKind::Onchain { status: ConfirmationStatus::Confirmed { height, .. }, .. } => { - let payment_id = payment.details.id; + let payment_id = details.id; if new_tip.height >= height + ANTI_REORG_DELAY - 1 { // Graduate from the live record, not the snapshot listed // above: a classification landing since then must not have @@ -449,8 +472,16 @@ impl Wallet { txid, status: ConfirmationStatus::Unconfirmed, .. - } if payment.details.direction == PaymentDirection::Outbound => { - unconfirmed_outbound_txids.push(txid); + } => { + if self + .fail_funding_payment_lost_to_conflict(&payment, new_tip.height) + .await? + { + continue; + } + if details.direction == PaymentDirection::Outbound { + unconfirmed_outbound_txids.push(txid); + } }, _ => {}, } @@ -487,12 +518,12 @@ 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) .await? .unwrap_or_else(|| PaymentId(txid.to_byte_array())); - if self + match self .apply_funding_status_update_locked( &guard, payment_id, @@ -501,6 +532,23 @@ impl Wallet { ) .await? { + 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()); + }, + } + + // The fallback id belongs to a settled funding payment whose entry is gone: + // skip rather than resurrect it (see `has_funding_record`). + if self.has_funding_record(&payment_id).await? { + log_debug!( + self.logger, + "Skipping wallet event for transaction {} of a settled funding payment", + txid, + ); continue; } @@ -515,10 +563,8 @@ impl Wallet { ConfirmationStatus::Unconfirmed, ) }; - let pending_payment = - self.create_pending_payment_from_tx(payment.clone(), Vec::new()); - self.payment_store.insert_or_update(payment).await?; - self.pending_payment_store.insert_or_update(pending_payment).await?; + self.payment_store.insert_or_update(payment.clone()).await?; + self.upsert_pending_payment(payment, Vec::new()).await?; }, WalletEvent::TxReplaced { txid, conflicts, .. } => { // See `TxConfirmed`: id resolution and the writes below must not interleave @@ -553,22 +599,31 @@ impl Wallet { payment_id, ); let payment = stored_payment.ok_or(Error::InvalidPaymentId)?; - let pending_payment_details = - self.create_pending_payment_from_tx(payment, conflict_txids.clone()); - self.pending_payment_store.insert_or_update(pending_payment_details).await?; + // A terminal record means the entry is the leftover of an interrupted settle + // — the record write landed, the entry removal was lost to a crash — and this + // event is the restart's replay of the same transition. Re-embedding the + // record would stamp the terminal status into the entry and hide it from the + // pending listing that repairs such leftovers; finish the interrupted removal + // instead. + if payment.status != PaymentStatus::Pending { + self.pending_payment_store.remove(&payment_id).await?; + continue; + } + + self.upsert_pending_payment(payment, conflict_txids).await?; }, WalletEvent::TxDropped { txid, tx } => { // See `TxConfirmed`: id resolution and the writes below must not interleave // 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) .await? .unwrap_or_else(|| PaymentId(txid.to_byte_array())); - if self + match self .apply_funding_status_update_locked( &guard, payment_id, @@ -577,6 +632,23 @@ impl Wallet { ) .await? { + 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()); + }, + } + + // The fallback id belongs to a settled funding payment whose entry is gone: + // skip rather than resurrect it (see `has_funding_record`). + if self.has_funding_record(&payment_id).await? { + log_debug!( + self.logger, + "Skipping wallet event for transaction {} of a settled funding payment", + txid, + ); continue; } @@ -591,10 +663,8 @@ impl Wallet { ConfirmationStatus::Unconfirmed, ) }; - let pending_payment = - self.create_pending_payment_from_tx(payment.clone(), Vec::new()); - self.payment_store.insert_or_update(payment).await?; - self.pending_payment_store.insert_or_update(pending_payment).await?; + self.payment_store.insert_or_update(payment.clone()).await?; + self.upsert_pending_payment(payment, Vec::new()).await?; }, _ => { continue; @@ -605,6 +675,157 @@ impl Wallet { Ok(()) } + /// Whether a funding-classified record exists under the given id. A funding record's id is + /// anchored to its first candidate's txid, so a wallet event for that transaction falls back + /// to this id whenever the pending entry no longer maps it — which only happens once the + /// negotiation settled and the entry was removed. The generic event handling must then skip + /// its write: merging a wallet-view `Pending` payment into the settled record would resurrect + /// it with figures no classification derived. + async fn has_funding_record(&self, payment_id: &PaymentId) -> Result { + Ok(self.payment_store.get(payment_id).await?.is_some_and(|payment| { + matches!( + payment.kind, + PaymentKind::Onchain { + tx_type: Some( + TransactionType::Funding { .. } + | TransactionType::InteractiveFunding { .. } + ), + .. + } + ) + })) + } + + /// Fails a funding payment whose transaction has irrevocably lost a conflict: a transaction + /// outside the record's candidate history — e.g. a channel close double-spending a pending + /// splice's shared input — has confirmed through [`ANTI_REORG_DELAY`] while neither the + /// record's transaction nor any candidate is canonical anymore. Returns whether the payment + /// was failed; failing also removes the pending entry, dropping the dead record from the + /// tip-change pass. (Its transaction was already excluded from rebroadcast by the same + /// canonical-only `get_tx` gate used below.) + /// + /// Only funding-classified records are considered: nothing re-submits a replaced funding + /// transaction under the same record (an RBF round is a new candidate), so a buried foreign + /// conflict is final for them. The liveness check guards the case where the conflict + /// double-spent only one round of the negotiation: as long as some candidate — including one + /// classification hasn't recorded yet — can still confirm, the record must stay pending. + async fn fail_funding_payment_lost_to_conflict( + &self, payment: &PendingPaymentDetails, tip_height: u32, + ) -> Result { + let payment_id = match payment.details() { + Some(details) => match details.kind { + PaymentKind::Onchain { + status: ConfirmationStatus::Unconfirmed, + tx_type: + Some( + TransactionType::Funding { .. } + | TransactionType::InteractiveFunding { .. }, + ), + .. + } => details.id, + _ => return Ok(false), + }, + None => return Ok(false), + }; + if payment.conflicting_txids().is_empty() { + return Ok(false); + } + + // Serialize with classification, whose retries extend the candidate history: the + // decision below must see that history in its settled form, and holding the lock keeps a + // concurrent write from resurrecting the entry removed at the end. + let _guard = self.funding_payment_update_lock.lock().await; + + // Re-read the entry under the lock; the listing snapshot may predate a classification. + let entry = match self.pending_payment_store.get(&payment_id).await? { + Some(entry) => entry, + None => return Ok(false), + }; + let PendingPaymentDetails::Tracked { details, conflicting_txids, candidates, .. } = &entry + else { + return Ok(false); + }; + let record_txid = match details.kind { + PaymentKind::Onchain { + txid, + status: ConfirmationStatus::Unconfirmed, + tx_type: + Some( + TransactionType::Funding { .. } + | TransactionType::InteractiveFunding { .. }, + ), + } => txid, + _ => return Ok(false), + }; + + let foreign_conflicts: Vec = conflicting_txids + .iter() + .copied() + .filter(|conflict| *conflict != record_txid && entry.candidate(*conflict).is_none()) + .collect(); + if foreign_conflicts.is_empty() { + return Ok(false); + } + + let lost = { + let locked_wallet = self.inner.lock().expect("lock"); + // `get_tx` is canonical-only: a transaction that lost to a confirmed conflict + // returns `None`, while one that can still confirm is `Some`. + let a_candidate_is_live = locked_wallet.get_tx(record_txid).is_some() + || candidates.iter().any(|c| locked_wallet.get_tx(c.txid).is_some()); + !a_candidate_is_live + && foreign_conflicts.iter().any(|conflict| { + match locked_wallet.get_tx(*conflict).map(|tx| tx.chain_position) { + Some(ChainPosition::Confirmed { anchor, .. }) => { + tip_height >= anchor.block_id.height + ANTI_REORG_DELAY - 1 + }, + _ => false, + } + }) + }; + if !lost { + return Ok(false); + } + + // As with graduation, decide from the live record and write only the status. A record + // already `Failed` — a prior pass whose entry removal below was lost to a crash — still + // matches, no-ops the update, and gets its lingering entry removed. + let mut failed = false; + self.payment_store + .mutate(&payment_id, |existing| { + let current = existing?; + match current.kind { + PaymentKind::Onchain { + txid, + status: ConfirmationStatus::Unconfirmed, + tx_type: + Some( + TransactionType::Funding { .. } + | TransactionType::InteractiveFunding { .. }, + ), + } if txid == record_txid => { + failed = true; + let mut update = PaymentDetailsUpdate::new(payment_id); + update.status = Some(PaymentStatus::Failed); + let mut updated = current.clone(); + updated.update(update).then_some(updated) + }, + _ => None, + } + }) + .await?; + if failed { + self.pending_payment_store.remove(&payment_id).await?; + log_info!( + self.logger, + "Failed funding payment {}: transaction {} lost to a conflicting transaction confirmed beyond the reorg depth", + payment_id, + record_txid, + ); + } + Ok(failed) + } + #[allow(deprecated)] pub(crate) async fn create_funding_transaction( &self, output_script: ScriptBuf, amount: Amount, confirmation_target: ConfirmationTarget, @@ -1586,7 +1807,15 @@ impl Wallet { return Ok(()); } - let payment_id = PaymentId(txid.to_byte_array()); + // Resolution and the writes below must share one lock acquisition: resolved outside it, + // the id could go stale against a record wallet sync creates for the same transaction, + // and the write below would create a divergent record. + let guard = self.funding_payment_update_lock.lock().await; + + // Adopt the id of a record that already tracks this transaction — e.g. a 0conf splice + // re-broadcast through LDK's generic funding path resolves back to its + // interactive-funding record here — otherwise generate a fresh id. + let payment_id = self.find_payment_by_txid(txid).await?.unwrap_or_else(random_payment_id); // A promoted-but-unconfirmed 0conf splice comes back through this generic path re-typed // and carrying wallet-view figures; `funding_reclassification_update` declines the @@ -1621,7 +1850,7 @@ impl Wallet { direction, PaymentStatus::Pending, ); - self.persist_funding_payment(details, Vec::new()).await?; + self.persist_funding_payment_locked(&guard, details, Vec::new()).await?; log_debug!( self.logger, "Recorded channel-funding broadcast {} for channel {}", @@ -1631,6 +1860,25 @@ impl Wallet { Ok(()) } + /// Returns the `PaymentId` of a user-initiated splice intent for one of the channels in + /// `candidate`, if any, so a classified splice adopts the id chosen at splice time rather than + /// deriving one from the first candidate's txid. A fee bump reuses the channel's existing intent, + /// so at most one in-flight intent matches and the first is unambiguous. + async fn find_splice_payment_id(&self, candidate: &FundingCandidate) -> Option { + self.pending_payment_store + .list_filter(|p| { + p.splice_intent().is_some_and(|intent| { + candidate.channels.iter().any(|channel| { + channel.channel_id == intent.channel_id + && channel.counterparty_node_id == intent.counterparty_node_id + }) + }) + }) + .await + .first() + .map(|p| p.id()) + } + /// Records an interactive-funding broadcast (splice, or a V2 dual-funded open) as a pending /// on-chain payment, tagged with its transaction type. Amount and fee are this node's share, /// derived from the active candidate's contributions; broadcasts we didn't contribute to, or @@ -1644,10 +1892,6 @@ impl Wallet { Some(c) => c, None => return Ok(()), }; - let first = match candidates.first() { - Some(c) => c, - None => return Ok(()), - }; let txid = tx.compute_txid(); debug_assert_eq!(active.txid, txid, "broadcast tx must match the active candidate"); @@ -1681,9 +1925,28 @@ impl Wallet { return Ok(()); } - // Anchor the `PaymentId` to the first negotiated candidate so the record stays stable - // across RBF replacements. - let payment_id = PaymentId(first.txid.to_byte_array()); + // Resolution and the writes below must share one lock acquisition: resolved outside it, + // the id could go stale against a record wallet sync creates for the same transaction, + // and the write below would create a divergent record. + let guard = self.funding_payment_update_lock.lock().await; + + // Adopt the `PaymentId` generated when the splice was initiated so its splice intent, + // funding payment, and candidate history share one record. If the intent is already gone + // (e.g. the splice locked before this classification ran), adopt the id of a record wallet + // sync created for any candidate rather than creating a divergent one; otherwise generate + // a fresh id — an id derived from a txid would tie the record's identity to one round of a + // replaceable transaction, and resolution through the record's txid history is what keeps + // its identity stable across RBF replacements. + let mut resolved_id = self.find_splice_payment_id(active).await; + if resolved_id.is_none() { + for candidate in candidates.iter() { + if let Some(id) = self.find_payment_by_txid(candidate.txid).await? { + resolved_id = Some(id); + break; + } + } + } + let payment_id = resolved_id.unwrap_or_else(random_payment_id); // Record every candidate's figures (`None` for any round we didn't contribute to, e.g. a // counterparty-initiated splice our `splice_in` later joined via RBF) so the confirmed @@ -1713,7 +1976,7 @@ impl Wallet { direction, PaymentStatus::Pending, ); - self.persist_funding_payment(details, candidate_records).await?; + self.persist_funding_payment_locked(&guard, details, candidate_records).await?; log_debug!( self.logger, "Recorded interactive-funding broadcast {} ({} candidates, {} channels)", @@ -1760,13 +2023,31 @@ impl Wallet { /// Writes a freshly-classified funding payment to the authoritative payment store and adds a /// pending-store index entry, so wallet sync graduates it through `ANTI_REORG_DELAY`. + /// + /// Production callers go through [`Self::persist_funding_payment_locked`] because they resolve + /// the record's id under the same lock acquisition; this wrapper models that acquisition for + /// tests entering classification mid-flow. + #[cfg(test)] async fn persist_funding_payment( &self, details: PaymentDetails, candidates: Vec, ) -> Result<(), Error> { // Hold the cross-store lock across both writes so a funding confirmation never observes // the record classified but the candidate history it needs still missing. - let _guard = self.funding_payment_update_lock.lock().await; + let guard = self.funding_payment_update_lock.lock().await; + self.persist_funding_payment_locked(&guard, details, candidates).await + } + /// Writes a freshly-classified funding payment to the authoritative payment store and adds a + /// pending-store index entry, so wallet sync graduates it through `ANTI_REORG_DELAY`. The + /// caller holds the cross-store lock, resolving the record's id and performing both store + /// writes under one acquisition, so a funding confirmation never observes the record + /// classified but the candidate history it needs still missing, and the resolved id never + /// goes stale against a concurrent sync write. + async fn persist_funding_payment_locked( + &self, _guard: &tokio::sync::MutexGuard<'_, ()>, details: PaymentDetails, + candidates: Vec, + ) -> Result<(), Error> { + let merge_candidates = candidates.clone(); // Everything this write does depends on the record's current state, so all of it must be // decided inside the store's critical section. When a record exists — no matter when it // appeared — only the classification (`tx_type`) and the figures of whichever candidate @@ -1811,34 +2092,137 @@ impl Wallet { self.pending_payment_store .mutate_async(&id, move |existing| async move { // The record was written above and payment records are never removed, so absence - // means the write failed out; fall back to the fresh details. + // means the write failed out; fall back to the fresh details. A promoted or + // (re)created entry embeds this post-write record rather than the fresh + // Unconfirmed details, so a confirmation wallet sync already recorded keeps + // driving graduation. let recorded = payment_store.get(&id).await?.unwrap_or(details); + // A candidate history that lacks the record's current txid is stale — a queued + // classification retrying after a newer round classified. The merge arm below + // refuses such a history; creating or promoting an entry from it would smuggle + // it past that refusal, so leave that to a fresh classification (the newer + // round's own write, or its retry) instead. + let stale = match &recorded.kind { + PaymentKind::Onchain { txid, .. } if !candidates.is_empty() => { + !candidates.iter().any(|c| c.txid == *txid) + }, + _ => false, + }; Ok(match existing { - // The inserted entry embeds the post-write record rather than the fresh - // details, so a confirmation wallet sync already recorded keeps driving - // graduation. - None if recorded.status == PaymentStatus::Pending => { - Some(PendingPaymentDetails::new(recorded, Vec::new(), candidates)) + // First time we record this funding payment — or a crash between the two + // store writes left a Pending record with no index entry: (re)create it so + // the payment can graduate and its candidate txids stay mapped. A graduated + // payment is never `Pending`, so absence with an advanced record means the + // graduation path removed the entry and it must not be re-indexed. + None => (recorded.status == PaymentStatus::Pending && !stale).then(|| { + PendingPaymentDetails::tracked(recorded, Vec::new(), candidates, None) + }), + // A user-initiated splice has a pre-broadcast `PendingSplice` intent under + // this id; carry its intent into the `Tracked` record so promotion does + // not drop it (nothing persists or consumes intents yet — that arrives + // with the follow-up that makes splice retries survive restarts). If the + // payment already advanced beyond `Pending` (wallet sync confirmed it + // through `ANTI_REORG_DELAY` first), it must not enter the pending store; + // the leftover intent record stays until that follow-up adds its clearing + // path. + Some(PendingPaymentDetails::PendingSplice { intent, .. }) => { + if recorded.status == PaymentStatus::Pending && !stale { + Some(PendingPaymentDetails::tracked( + recorded, + Vec::new(), + candidates, + Some(intent), + )) + } else { + None + } }, - // The payment already advanced beyond Pending: the graduation path removed - // the entry and it must not be re-created. - None => None, - // The entry predates this classification — wallet sync recorded the - // transaction before it was classified (its arms and this write pair - // serialize on the cross-store lock, so nothing lands in between): merge - // only the classification into the existing entry. - Some(mut entry) => { + // An earlier candidate's classification or wallet sync recorded this payment + // before this classification ran (sync's arms and this write pair serialize + // on the cross-store lock, so nothing lands in between): merge only the + // classification (`tx_type`, candidate history and the figures of whichever + // candidate the record's state makes authoritative) into it. + Some(mut tracked @ PendingPaymentDetails::Tracked { .. }) => { let pending_update = PendingPaymentDetailsUpdate { id, payment_update: Some(update), conflicting_txids: None, candidates, + splice_intent: None, }; - entry.update(pending_update).then_some(entry) + tracked.update(pending_update).then_some(tracked) }, }) }) .await?; + + // With the candidate history recorded, duplicates wallet sync created for rounds that were + // not yet candidates can be folded back into this record. Runs after both writes so the + // funding-status gate accepts the candidates it adopts, and under the same lock + // acquisition, so sync cannot interleave; a failure surfaces to the broadcast queue's + // classification retry, which re-runs this idempotently. + self.merge_duplicate_candidate_records(_guard, id, &merge_candidates).await?; + Ok(()) + } + + /// Merges duplicate records wallet sync created for this funding payment's candidates before + /// they were classified. Sync re-keys an event for a round it cannot attribute to the + /// funding record — not yet a candidate, so the funding-status gate reports it foreign — to + /// the round's txid-derived id, creating an untyped duplicate whose pending entry then + /// shadows the funding record in [`Self::find_payment_by_txid`]'s direct probe. Once the + /// round is a recorded candidate, the duplicate's confirmation (if any) belongs on the + /// funding record: adopt it, then remove the duplicate and its pending entry. + /// + /// The caller must hold [`Self::funding_payment_update_lock`], per + /// [`Self::apply_funding_status_update_locked`]'s contract. + async fn merge_duplicate_candidate_records( + &self, guard: &tokio::sync::MutexGuard<'_, ()>, id: PaymentId, + candidates: &[FundingTxCandidate], + ) -> Result<(), Error> { + for candidate in candidates { + let duplicate_id = PaymentId(candidate.txid.to_byte_array()); + if duplicate_id == id { + continue; + } + let duplicate = match self.payment_store.get(&duplicate_id).await? { + Some(duplicate) => duplicate, + None => continue, + }; + // Only a duplicate view of this candidate's transaction qualifies: an untyped record + // wallet sync created, or one a funding-typed rebroadcast classified onto it. Anything + // else keyed by the txid-derived id is left alone. + let status = match &duplicate.kind { + PaymentKind::Onchain { + txid, + status, + tx_type: None | Some(TransactionType::Funding { .. }), + } if *txid == candidate.txid => status.clone(), + _ => continue, + }; + // Only a confirmation is worth adopting; an unconfirmed duplicate carries nothing the + // record needs — the actively-broadcast candidate stays the record's current txid. + if matches!(status, ConfirmationStatus::Confirmed { .. }) { + let outcome = self + .apply_funding_status_update_locked(guard, id, candidate.txid, status) + .await?; + debug_assert!(matches!(outcome, FundingStatusUpdate::Applied)); + if !matches!(outcome, FundingStatusUpdate::Applied) { + // Adoption declined; keep the duplicate rather than discard its confirmation. + continue; + } + } + log_debug!( + self.logger, + "Merging duplicate payment record for funding transaction {}", + candidate.txid, + ); + // Pending entry first: the retry of a failure between these two removals rediscovers + // the duplicate through its payment record. Removed the other way around, the + // leftover pending entry would be unreachable to the retry yet keep shadowing the + // funding record in `find_payment_by_txid`'s direct probe. + self.pending_payment_store.remove(&duplicate_id).await?; + self.payment_store.remove(&duplicate_id).await?; + } Ok(()) } @@ -1904,10 +2288,52 @@ impl Wallet { PaymentDetails::new(payment_id, kind, amount_msat, fee_paid_msat, direction, payment_status) } - fn create_pending_payment_from_tx( + /// Inserts or refreshes the pending-store entry tracking `payment` toward graduation, + /// atomically with reading the entry's current state. + async fn upsert_pending_payment( &self, payment: PaymentDetails, conflicting_txids: Vec, - ) -> PendingPaymentDetails { - PendingPaymentDetails::new(payment, conflicting_txids, Vec::new()) + ) -> Result<(), Error> { + let id = payment.id; + let payment_store = Arc::clone(&self.payment_store); + self.pending_payment_store + .mutate_async(&id, move |existing| async move { + // Only `Pending` payments belong in the pending store. Like in + // [`Self::persist_funding_payment`], the authoritative status is re-read inside + // the store's critical section, where it cannot go stale against graduation. + let is_pending = payment_store + .get(&id) + .await? + .map_or(payment.status == PaymentStatus::Pending, |recorded| { + recorded.status == PaymentStatus::Pending + }); + if !is_pending { + return Ok(None); + } + Ok(match existing { + None => { + Some(PendingPaymentDetails::new(payment, conflicting_txids, Vec::new())) + }, + // Promote a pre-broadcast splice intent: wallet sync saw the splice + // transaction before its broadcast-time classification recorded it. Carrying + // the intent into the `Tracked` record makes the entry visible to txid + // lookups while preserving the intent. + Some(PendingPaymentDetails::PendingSplice { intent, .. }) => { + Some(PendingPaymentDetails::tracked( + payment, + conflicting_txids, + Vec::new(), + Some(intent), + )) + }, + Some(mut tracked @ PendingPaymentDetails::Tracked { .. }) => { + let fresh = + PendingPaymentDetails::new(payment, conflicting_txids, Vec::new()); + tracked.update(fresh.to_update()).then_some(tracked) + }, + }) + }) + .await?; + Ok(()) } async fn find_payment_by_txid(&self, target_txid: Txid) -> Result, Error> { @@ -1919,17 +2345,38 @@ impl Wallet { if let Some(replaced_details) = self .pending_payment_store .list_filter(|p| { - matches!(p.details.kind, PaymentKind::Onchain { txid, .. } if txid == target_txid) - || p.conflicting_txids.contains(&target_txid) + p.details().is_some_and( + |d| matches!(d.kind, PaymentKind::Onchain { txid, .. } if txid == target_txid), + ) || p.conflicting_txids().contains(&target_txid) // A middle RBF round is not the record's current txid and may never have - // received a `TxReplaced` event of its own, so map any of its candidate - // txids (an earlier RBF round may confirm) back to the record. + // received a `TxReplaced` event of its own, and a splice keyed by a generated + // PaymentId is not found by the txid-derived id above: map any of the + // candidate txids (an earlier RBF round may confirm) back to the record. || p.candidate(target_txid).is_some() }) .await .first() { - return Ok(Some(replaced_details.details.id)); + return Ok(Some(replaced_details.id())); + } + + // The pending store only indexes in-flight records — graduation removes the entry — so a + // graduated record's transaction resolves through the payment store itself. Without this, + // a funding-typed broadcast classified after graduation (e.g. LDK re-broadcasting a promoted + // 0conf splice whose confirmation landed while the node was offline) would create a + // duplicate record, and a post-graduation reorg's events would never reach the record. + let mut page_token = None; + loop { + let page = self.payment_store.list_page(page_token).await?; + if let Some(payment) = page.objects.iter().find( + |p| matches!(p.kind, PaymentKind::Onchain { txid, .. } if txid == target_txid), + ) { + return Ok(Some(payment.id)); + } + match page.next_page_token { + Some(token) => page_token = Some(token), + None => break, + } } Ok(None) @@ -1938,9 +2385,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 @@ -1949,38 +2398,51 @@ 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 caller's wallet-level lock keeps the candidate history stable while we await its - // read. The funding-type gate and write then share the payment store's mutation lock: - // against a separate payment `get`, a classification merging in between would have its - // `tx_type` and contribution figures clobbered by this stale snapshot. + // read. The funding-type gate, the candidate lookup, and the write then share the payment + // store's mutation lock: against a separate payment `get`, a classification merging in + // between would have its `tx_type` and contribution figures clobbered by this stale + // snapshot. let pending_payment = self.pending_payment_store.get(&payment_id).await?; + 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 owns_event_tx = event_txid == current_txid + || pending_payment.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) = pending_payment.as_ref() { - 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_payment.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 }; @@ -1998,17 +2460,16 @@ 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 // the same dual-write the default `TxConfirmed` path performs; an empty conflicting-txids // list leaves any stored conflicts intact (the update treats absent as "unchanged"). if payment.status == PaymentStatus::Pending { - let pending = self.create_pending_payment_from_tx(payment, Vec::new()); - self.pending_payment_store.insert_or_update(pending).await?; + self.upsert_pending_payment(payment, Vec::new()).await?; } - Ok(true) + Ok(FundingStatusUpdate::Applied) } #[allow(deprecated)] @@ -2249,8 +2710,6 @@ impl Wallet { ConfirmationStatus::Unconfirmed, ); - let pending_payment_store = - self.create_pending_payment_from_tx(new_payment.clone(), Vec::new()); let change_set = locked_wallet.take_staged().unwrap_or_default(); drop(locked_wallet); locked_persister.persist_changeset(change_set).await.map_err(|e| { @@ -2258,8 +2717,8 @@ impl Wallet { Error::PersistenceFailed })?; - self.payment_store.insert_or_update(new_payment).await?; - self.pending_payment_store.insert_or_update(pending_payment_store).await?; + self.payment_store.insert_or_update(new_payment.clone()).await?; + self.upsert_pending_payment(new_payment, Vec::new()).await?; self.broadcaster.broadcast_unclassified_transaction(fee_bumped_tx); @@ -2311,6 +2770,29 @@ fn aggregate_local_stakes(candidate: &FundingCandidate) -> LocalStakeAggregate { } } +/// Generates a fresh funding-record [`PaymentId`] from the OS entropy source. A funding record's id +/// carries no meaning beyond uniqueness: the record is found through its transaction history +/// ([`Wallet::find_payment_by_txid`]), never re-derived from a txid. +fn random_payment_id() -> PaymentId { + let mut bytes = [0u8; 32]; + getrandom::fill(&mut bytes).expect("getrandom failed"); + PaymentId(bytes) +} + +/// 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, @@ -2638,7 +3120,7 @@ fn ldk_to_bdk_satisfaction_weight(ldk_satisfaction_weight: u64) -> Weight { /// classification. /// /// `current` is the record as observed inside the payment store's `mutate` critical section — its -/// sole caller, [`Wallet::persist_funding_payment`], builds and applies the update within one +/// sole caller, [`Wallet::persist_funding_payment_locked`], builds and applies the update within one /// closure — so the candidate choice cannot go stale against a concurrent confirmation before the /// update lands. [`PaymentDetails::update`]'s confirmed-figures rule still arbitrates which /// figures may land on the record. @@ -2667,6 +3149,29 @@ fn funding_reclassification_update( return PaymentDetailsUpdate::new(details.id); } + // An interactive-funding classification carries the full candidate history as of its own + // broadcast, and once a record is funding-classified its txid only ever names a candidate + // from that history. A classification whose history lacks such a record's current txid was + // therefore built before that candidate existed — a queued retry running after a newer round + // classified. Applying it would rotate the record backwards; the newer round's + // classification already recorded everything this one knows. A record that is not yet + // funding-classified gives no such signal — wallet sync can have rotated its txid to a + // conflicting transaction that is no candidate at all — so its first classification must + // still land. + if !candidates.is_empty() { + if let Some(PaymentKind::Onchain { + txid: current_txid, + tx_type: + Some(TransactionType::Funding { .. } | TransactionType::InteractiveFunding { .. }), + .. + }) = current.map(|payment| &payment.kind) + { + if !candidates.iter().any(|c| c.txid == *current_txid) { + return PaymentDetailsUpdate::new(details.id); + } + } + } + let mut update = PaymentDetailsUpdate::funding_reclassification(details); if let Some(PaymentKind::Onchain { txid: confirmed_txid, @@ -2687,7 +3192,7 @@ fn funding_reclassification_update( #[cfg(all(test, any(feature = "chain-esplora", feature = "chain-electrum")))] mod tests { - use std::sync::atomic::{AtomicBool, Ordering}; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::time::Duration; use bdk_chain::{BlockId, ConfirmationBlockTime}; @@ -2717,11 +3222,13 @@ mod tests { const EXTERNAL_DESCRIPTOR: &str = "wpkh(tprv8ZgxMBicQKsPdy6LMhUtFHAgpocR8GC6QmwMSFpZs7h6Eziw3SpThFfczTDh5rW2krkqffa11UpX3XkeTTB2FvzZKWXqPY54Y6Rq4AQ5R8L/84'/1'/0'/0/*)"; const INTERNAL_DESCRIPTOR: &str = "wpkh(tprv8ZgxMBicQKsPdy6LMhUtFHAgpocR8GC6QmwMSFpZs7h6Eziw3SpThFfczTDh5rW2krkqffa11UpX3XkeTTB2FvzZKWXqPY54Y6Rq4AQ5R8L/84'/1'/0'/1/*)"; - /// An in-memory store whose writes can be made to fail on demand. + /// An in-memory store whose writes can be made to fail on demand, counting the failures so + /// tests can wait for a write to have actually failed rather than guessing with a sleep. #[derive(Clone)] struct FailSwitchStore { inner: Arc, fail_writes: Arc, + failed_writes: Arc, } impl FailSwitchStore { @@ -2729,6 +3236,7 @@ mod tests { Self { inner: Arc::new(InMemoryStore::new()), fail_writes: Arc::new(AtomicBool::new(false)), + failed_writes: Arc::new(AtomicUsize::new(0)), } } } @@ -2745,11 +3253,13 @@ mod tests { ) -> impl Future> + 'static + Send { let inner = Arc::clone(&self.inner); let fail_writes = Arc::clone(&self.fail_writes); + let failed_writes = Arc::clone(&self.failed_writes); let primary_namespace = primary_namespace.to_string(); let secondary_namespace = secondary_namespace.to_string(); let key = key.to_string(); async move { if fail_writes.load(Ordering::Acquire) { + failed_writes.fetch_add(1, Ordering::AcqRel); return Err(io::Error::new(io::ErrorKind::Other, "writes disabled")); } KVStore::write(&*inner, &primary_namespace, &secondary_namespace, &key, buf).await @@ -2783,6 +3293,86 @@ mod tests { } } + /// An in-memory store that fails the next remove issued against an armed namespace, for + /// exercising cleanup paths that must survive a failure between two removals. + #[derive(Clone)] + struct FailRemoveStore { + inner: Arc, + fail_remove_in: Arc>>, + } + + impl FailRemoveStore { + fn new() -> Self { + Self { + inner: Arc::new(InMemoryStore::new()), + fail_remove_in: Arc::new(std::sync::Mutex::new(None)), + } + } + + fn fail_next_remove_in(&self, primary_namespace: &str) { + *self.fail_remove_in.lock().unwrap() = Some(primary_namespace.to_string()); + } + } + + impl KVStore for FailRemoveStore { + fn read( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, + ) -> impl Future, io::Error>> + 'static + Send { + KVStore::read(&*self.inner, primary_namespace, secondary_namespace, key) + } + + fn write( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, buf: Vec, + ) -> impl Future> + 'static + Send { + KVStore::write(&*self.inner, primary_namespace, secondary_namespace, key, buf) + } + + fn remove( + &self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool, + ) -> impl Future> + 'static + Send { + let inner = Arc::clone(&self.inner); + let armed = Arc::clone(&self.fail_remove_in); + let primary_namespace = primary_namespace.to_string(); + let secondary_namespace = secondary_namespace.to_string(); + let key = key.to_string(); + async move { + let fail = { + let mut armed = armed.lock().unwrap(); + if armed.as_deref() == Some(primary_namespace.as_str()) { + *armed = None; + true + } else { + false + } + }; + if fail { + return Err(io::Error::new(io::ErrorKind::Other, "removes disabled")); + } + KVStore::remove(&*inner, &primary_namespace, &secondary_namespace, &key, lazy).await + } + } + + fn list( + &self, primary_namespace: &str, secondary_namespace: &str, + ) -> impl Future, io::Error>> + 'static + Send { + KVStore::list(&*self.inner, primary_namespace, secondary_namespace) + } + } + + impl PaginatedKVStore for FailRemoveStore { + fn list_paginated( + &self, primary_namespace: &str, secondary_namespace: &str, + page_token: Option, + ) -> impl Future> + 'static + Send { + PaginatedKVStore::list_paginated( + &*self.inner, + primary_namespace, + secondary_namespace, + page_token, + ) + } + } + /// Constructs a `Wallet` around the given store, either creating a fresh BDK wallet or /// loading the one the store already holds. async fn new_test_wallet(store: Arc, load_existing: bool) -> Arc { @@ -3713,6 +4303,57 @@ mod tests { } } + /// Inserts `tx` into the BDK wallet as canonically confirmed at `height`, extending the + /// local chain to that height. + fn insert_confirmed_tx(wallet: &Wallet, tx: Transaction, height: u32) { + let txid = tx.compute_txid(); + let mut locked = wallet.inner.lock().unwrap(); + let block = + BlockId { height, hash: bitcoin::BlockHash::from_byte_array([height as u8; 32]) }; + let chain = locked.latest_checkpoint().insert(block); + let mut tx_update = bdk_chain::TxUpdate::default(); + tx_update.txs = vec![Arc::new(tx)]; + tx_update.anchors = + [(ConfirmationBlockTime { block_id: block, confirmation_time: 100 }, txid)].into(); + locked + .apply_update(Update { tx_update, chain: Some(chain), ..Default::default() }) + .unwrap(); + } + + /// Inserts `tx` into the BDK wallet as canonically unconfirmed (seen in the mempool). + fn insert_unconfirmed_tx(wallet: &Wallet, tx: Transaction) { + let txid = tx.compute_txid(); + let mut locked = wallet.inner.lock().unwrap(); + let mut tx_update = bdk_chain::TxUpdate::default(); + tx_update.txs = vec![Arc::new(tx)]; + tx_update.seen_ats = [(txid, 100)].into(); + locked.apply_update(Update { tx_update, ..Default::default() }).unwrap(); + } + + /// Builds a transaction paying a wallet address, spending an outpoint derived from + /// `input_byte` (distinct bytes yield non-conflicting transactions). + fn wallet_paying_tx(wallet: &Wallet, input_byte: u8) -> Transaction { + let script_pubkey = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: LockTime::ZERO, + input: vec![bitcoin::TxIn { + previous_output: OutPoint { + txid: Txid::from_byte_array([input_byte; 32]), + vout: 0, + }, + ..Default::default() + }], + output: vec![TxOut { value: Amount::from_sat(90_000), script_pubkey }], + } + } + #[test] fn funding_reclassification_update_substitutes_the_confirmed_candidate() { let confirmed_txid = Txid::from_byte_array([1u8; 32]); @@ -3755,38 +4396,100 @@ mod tests { #[test] fn funding_reclassification_update_keeps_the_active_candidate() { + let prior_txid = Txid::from_byte_array([1u8; 32]); let active_txid = Txid::from_byte_array([2u8; 32]); - let candidates = vec![FundingTxCandidate { - txid: active_txid, - amount_msat: Some(1_000_000), - fee_paid_msat: Some(500), - }]; - let details = onchain_details(active_txid, ConfirmationStatus::Unconfirmed); - - // No record yet: the update describes the active candidate. - let update = funding_reclassification_update(details.clone(), &candidates, None); - assert_eq!(update.txid, Some(active_txid)); - assert_eq!(update.amount_msat, Some(Some(1_000_000))); - - // An unconfirmed record: still the active candidate (RBF rotation). - let unconfirmed = - onchain_details(Txid::from_byte_array([1u8; 32]), ConfirmationStatus::Unconfirmed); + let candidates = vec![ + FundingTxCandidate { + txid: prior_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(400), + }, + FundingTxCandidate { + txid: active_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }, + ]; + let details = onchain_details(active_txid, ConfirmationStatus::Unconfirmed); + + // No record yet: the update describes the active candidate. + let update = funding_reclassification_update(details.clone(), &candidates, None); + assert_eq!(update.txid, Some(active_txid)); + assert_eq!(update.amount_msat, Some(Some(1_000_000))); + + // An unconfirmed record on the prior candidate: rotate to the active one (RBF). + let unconfirmed = onchain_details(prior_txid, ConfirmationStatus::Unconfirmed); let update = funding_reclassification_update(details.clone(), &candidates, Some(&unconfirmed)); assert_eq!(update.txid, Some(active_txid)); // The record confirmed the active candidate itself: nothing to substitute. let current = onchain_details(active_txid, confirmed_status()); - let update = funding_reclassification_update(details.clone(), &candidates, Some(¤t)); + let update = funding_reclassification_update(details, &candidates, Some(¤t)); assert_eq!(update.txid, Some(active_txid)); assert_eq!(update.amount_msat, Some(Some(1_000_000))); + } - // A confirmed txid outside the candidate history (e.g. the record is an unrelated - // same-id payment): fall back to the active candidate; `PaymentDetails::update` keeps - // the confirmed figures in place on mismatch. - let foreign = onchain_details(Txid::from_byte_array([9u8; 32]), confirmed_status()); - let update = funding_reclassification_update(details, &candidates, Some(&foreign)); - assert_eq!(update.txid, Some(active_txid)); + /// A classification whose candidate history lacks a funding-classified record's current txid + /// was built before that candidate existed — a queued retry running after a newer round + /// classified — and must move nothing, whatever the record's confirmation state. A record + /// that is not yet funding-classified gives no such signal (wallet sync can have rotated its + /// txid to a conflicting non-candidate), so its first classification must still land. + #[test] + fn funding_reclassification_update_refuses_a_stale_candidate_history() { + let stale_txid = Txid::from_byte_array([1u8; 32]); + let newer_txid = Txid::from_byte_array([2u8; 32]); + let payment_id = PaymentId(stale_txid.to_byte_array()); + let stale_history = vec![FundingTxCandidate { + txid: stale_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(400), + }]; + let details = + interactive_funding_details(payment_id, stale_txid, Some(1_000_000), Some(400)); + + // The record moved on to a newer candidate while this classification was queued. + let unconfirmed = + interactive_funding_details(payment_id, newer_txid, Some(1_000_000), Some(500)); + let update = + funding_reclassification_update(details.clone(), &stale_history, Some(&unconfirmed)); + let mut updated = unconfirmed.clone(); + assert!(!updated.update(update), "a stale retry must not move an unconfirmed record"); + assert_eq!(updated, unconfirmed); + + // Same when the newer candidate has already confirmed. + let mut confirmed = unconfirmed.clone(); + confirmed.kind = PaymentKind::Onchain { + txid: newer_txid, + status: confirmed_status(), + tx_type: Some(TransactionType::InteractiveFunding { channels: vec![] }), + }; + let update = + funding_reclassification_update(details.clone(), &stale_history, Some(&confirmed)); + let mut updated = confirmed.clone(); + assert!(!updated.update(update), "a stale retry must not move a confirmed record"); + assert_eq!(updated, confirmed); + + // A record that was never funding-classified: wallet sync rotated its txid to a + // conflicting transaction, which is no candidate. Its first classification is not stale + // and must land. + let mut unclassified = + interactive_funding_details(payment_id, newer_txid, Some(1_000_000), Some(500)); + unclassified.kind = PaymentKind::Onchain { + txid: newer_txid, + status: ConfirmationStatus::Unconfirmed, + tx_type: None, + }; + let update = funding_reclassification_update(details, &stale_history, Some(&unclassified)); + let mut updated = unclassified.clone(); + assert!(updated.update(update), "a first classification must not be treated as stale"); + match &updated.kind { + PaymentKind::Onchain { txid, tx_type, .. } => { + assert_eq!(*txid, stale_txid); + assert!(matches!(tx_type, Some(TransactionType::InteractiveFunding { .. }))); + }, + kind => panic!("unexpected kind {:?}", kind), + } } /// A funding-typed (re)classification of a record already classified as interactive funding @@ -3960,52 +4663,45 @@ mod tests { assert_eq!(wallet.find_payment_by_txid(txid2).await.unwrap(), Some(payment_id)); } - /// 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 - /// contribution, or none of the moved funds are the wallet's — would otherwise come back as - /// a spurious zero-amount record that nothing ever confirms. + /// A graduated funding record has no pending entry — graduation removes it — so its txid must + /// resolve through the payment store itself. Without that fallback, a funding-typed broadcast + /// classified after graduation (e.g. LDK re-broadcasting a promoted 0conf splice whose + /// confirmation landed while the node was offline) would miss the record and create a duplicate + /// under a fresh id. #[tokio::test] - async fn funding_broadcast_without_wallet_activity_is_not_recorded() { + async fn find_payment_by_txid_resolves_graduated_records() { let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); let wallet = new_test_wallet(store, false).await; - let counterparty_node_id = PublicKey::from_str( - "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", - ) - .unwrap(); - let channels = vec![(counterparty_node_id, ChannelId([7u8; 32]))]; - let tx_type = TransactionType::Funding { channels: vec![] }; + let txid = Txid::from_byte_array([6u8; 32]); + let payment_id = PaymentId([21u8; 32]); + let mut graduated = + interactive_funding_details(payment_id, txid, Some(1_000_000), Some(500)); + graduated.kind = PaymentKind::Onchain { + txid, + status: confirmed_status(), + tx_type: Some(TransactionType::InteractiveFunding { channels: vec![] }), + }; + graduated.status = PaymentStatus::Succeeded; + wallet.payment_store.insert_or_update(graduated).await.unwrap(); - // No inputs or outputs involve the wallet: nothing to record. - wallet.classify_funding(&dummy_tx(), &channels, tx_type.clone()).await.unwrap(); - assert!(wallet.payment_store.list_page(None).await.unwrap().objects.is_empty()); - assert!(wallet.pending_payment_store.list_filter(|_| true).await.is_empty()); + assert_eq!(wallet.find_payment_by_txid(txid).await.unwrap(), Some(payment_id)); + } - // A computable fee is not wallet participation. The wallet can resolve a splice's shared - // input whenever the previous funding transaction touched it (e.g. it funded the original - // channel open), so it derives the splice's fee even when no wallet funds move. - let prev_funding_outpoint = OutPoint { txid: Txid::from_byte_array([8u8; 32]), vout: 0 }; - wallet.inner.lock().unwrap().insert_txout( - prev_funding_outpoint, - TxOut { value: Amount::from_sat(100_000), script_pubkey: ScriptBuf::new() }, - ); - let splice_tx = Transaction { - version: bitcoin::transaction::Version::TWO, - lock_time: LockTime::ZERO, - input: vec![bitcoin::TxIn { - previous_output: prev_funding_outpoint, - ..Default::default() - }], - output: vec![TxOut { - value: Amount::from_sat(99_000), - script_pubkey: ScriptBuf::new(), - }], - }; - wallet.classify_funding(&splice_tx, &channels, tx_type.clone()).await.unwrap(); - assert!(wallet.payment_store.list_page(None).await.unwrap().objects.is_empty()); + /// 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; - // Control: a funding transaction the wallet participates in is still recorded. + 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() @@ -4013,93 +4709,1302 @@ mod tests { .reveal_next_address(KeychainKind::External) .address .script_pubkey(); - let funded_tx = Transaction { + let close_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 }], + 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 }], }; - wallet.classify_funding(&funded_tx, &channels, tx_type).await.unwrap(); - let payments = wallet.payment_store.list_page(None).await.unwrap().objects; - assert_eq!(payments.len(), 1); - assert_eq!(payments[0].id, PaymentId(funded_tx.compute_txid().to_byte_array())); + 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(), + splice_intent: None, + }) + .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).await.unwrap().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())) + .await + .unwrap() + .unwrap(); + match &close.kind { + PaymentKind::Onchain { txid, status, .. } => { + assert_eq!(*txid, close_txid); + assert!(matches!(status, ConfirmationStatus::Confirmed { .. })); + }, + kind => panic!("unexpected kind {:?}", kind), + } } - /// LDK re-broadcasts a promoted-but-unconfirmed 0conf splice through its generic funding - /// path: same txid, but typed as a plain funding transaction with wallet-view figures and no - /// contribution metadata. The rebroadcast must not overwrite the contribution-derived - /// figures or the interactive-funding classification — neither while the record is - /// unconfirmed nor once it confirmed under that same txid, where updates naming the - /// confirmed txid may otherwise move figures. + /// Continues the story above: once the conflicting close confirms through the anti-reorg + /// depth, the splice's funding transaction can never confirm — its shared input is spent for + /// good. The record must fail rather than stay `Pending` forever, and removing the pending + /// entry stops the dead transaction's rebroadcast on every tip change. #[tokio::test] - async fn funding_rebroadcast_keeps_interactive_funding_classification() { + async fn funding_payment_fails_once_a_foreign_conflict_confirms_to_depth() { let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); let wallet = new_test_wallet(store, false).await; - // The rebroadcast passes the wallet-activity guard: a splice-in funds the new channel - // output partly from the wallet, so the wallet sees movement. - 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 txid = tx.compute_txid(); - let payment_id = PaymentId(txid.to_byte_array()); + let close_tx = wallet_paying_tx(&wallet, 3); + 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, + txid: splice_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 details = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); wallet.persist_funding_payment(details, candidates).await.unwrap(); + wallet + .pending_payment_store + .update(PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: Some(vec![close_txid]), + candidates: Vec::new(), + splice_intent: None, + }) + .await + .unwrap(); - let counterparty_node_id = PublicKey::from_str( - "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", - ) - .unwrap(); - let channels = vec![(counterparty_node_id, ChannelId([7u8; 32]))]; - let tx_type = TransactionType::Funding { channels: vec![] }; + // The close is canonically confirmed; the splice transaction, having lost the conflict, + // is no longer canonical (here: never inserted at all). + insert_confirmed_tx(&wallet, close_tx, 5); - async fn assert_unchanged(wallet: &Wallet, payment_id: PaymentId, confirmed: bool) { - let payments = wallet.payment_store.list_page(None).await.unwrap().objects; - assert_eq!(payments.len(), 1, "the rebroadcast must not mint a second record"); - let payment = &payments[0]; - assert_eq!(payment.id, payment_id); - assert_eq!(payment.amount_msat, Some(1_000_000)); - assert_eq!(payment.fee_paid_msat, Some(500)); - match &payment.kind { - PaymentKind::Onchain { - status, - tx_type: Some(TransactionType::InteractiveFunding { .. }), - .. - } => assert_eq!(matches!(status, ConfirmationStatus::Confirmed { .. }), confirmed), - kind => panic!("unexpected kind {:?}", kind), - } + let block_id = + |height| BlockId { height, hash: bitcoin::BlockHash::from_byte_array([7u8; 32]) }; + let event = WalletEvent::ChainTipChanged { + old_tip: block_id(9), + new_tip: block_id(5 + ANTI_REORG_DELAY - 1), + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Failed); + match &payment.kind { + PaymentKind::Onchain { txid, status, tx_type } => { + assert_eq!(*txid, splice_txid, "failing must not adopt the conflict's txid"); + assert!(matches!(status, ConfirmationStatus::Unconfirmed)); + assert!(matches!(tx_type, Some(TransactionType::InteractiveFunding { .. }))); + }, + kind => panic!("unexpected kind {:?}", kind), } + assert_eq!(payment.amount_msat, Some(1_000_000)); + assert_eq!(payment.fee_paid_msat, Some(500)); + assert!( + wallet.pending_payment_store.get(&payment_id).await.unwrap().is_none(), + "the entry must go so the dead transaction stops being rebroadcast" + ); + } - wallet.classify_funding(&tx, &channels, tx_type.clone()).await.unwrap(); - assert_unchanged(&wallet, payment_id, false).await; + /// A confirmed conflict that is one of the record's own candidates is RBF resolution, not a + /// loss: classification adopts it into the record, so the failure pass must leave the record + /// alone. + #[tokio::test] + async fn funding_payment_survives_a_confirmed_conflict_that_is_a_candidate() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; - // Confirm the record, then replay the rebroadcast: a monitor-update completion can race - // wallet sync around confirmation. - let event = WalletEvent::TxConfirmed { - txid, - tx: Arc::new(tx.clone()), - block_time: confirmed_block_time(5), - old_block_time: None, + let bumped_tx = wallet_paying_tx(&wallet, 3); + let bumped_txid = bumped_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), + }, + FundingTxCandidate { + txid: bumped_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(600), + }, + ]; + let details = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + wallet.persist_funding_payment(details, candidates).await.unwrap(); + wallet + .pending_payment_store + .update(PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: Some(vec![bumped_txid]), + candidates: Vec::new(), + splice_intent: None, + }) + .await + .unwrap(); + + insert_confirmed_tx(&wallet, bumped_tx, 5); + + let block_id = + |height| BlockId { height, hash: bitcoin::BlockHash::from_byte_array([7u8; 32]) }; + let event = WalletEvent::ChainTipChanged { + old_tip: block_id(9), + new_tip: block_id(5 + ANTI_REORG_DELAY - 1), }; wallet.update_payment_store(vec![event]).await.unwrap(); - wallet.classify_funding(&tx, &channels, tx_type).await.unwrap(); - assert_unchanged(&wallet, payment_id, true).await; + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Pending); + assert!( + wallet.pending_payment_store.get(&payment_id).await.unwrap().is_some(), + "the entry must survive for classification to adopt the confirmed candidate" + ); + } + + /// A foreign conflict that has confirmed but not yet through the anti-reorg depth may still + /// be reorged out, letting the funding transaction confirm after all; the record must stay + /// pending until the conflict's confirmation is final. + #[tokio::test] + async fn funding_payment_survives_a_foreign_conflict_short_of_depth() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let close_tx = wallet_paying_tx(&wallet, 3); + 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(); + wallet + .pending_payment_store + .update(PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: Some(vec![close_txid]), + candidates: Vec::new(), + splice_intent: None, + }) + .await + .unwrap(); + + insert_confirmed_tx(&wallet, close_tx, 5); + + let block_id = + |height| BlockId { height, hash: bitcoin::BlockHash::from_byte_array([7u8; 32]) }; + let event = WalletEvent::ChainTipChanged { + old_tip: block_id(9), + new_tip: block_id(5 + ANTI_REORG_DELAY - 2), + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Pending); + assert!(wallet.pending_payment_store.get(&payment_id).await.unwrap().is_some()); + } + + /// A conflict may double-spend only one round of the negotiation — e.g. it shares an input + /// with an RBF attempt but not with the original candidate. While any candidate is still + /// canonical it can still confirm, so the record must stay pending. + #[tokio::test] + async fn funding_payment_survives_while_a_candidate_can_still_confirm() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let conflict_tx = wallet_paying_tx(&wallet, 3); + let conflict_txid = conflict_tx.compute_txid(); + // A live candidate: spends a different outpoint, so the conflict didn't kill it. + let live_candidate_tx = wallet_paying_tx(&wallet, 4); + let live_candidate_txid = live_candidate_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), + }, + FundingTxCandidate { + txid: live_candidate_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(600), + }, + ]; + let details = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + wallet.persist_funding_payment(details, candidates).await.unwrap(); + wallet + .pending_payment_store + .update(PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: Some(vec![conflict_txid]), + candidates: Vec::new(), + splice_intent: None, + }) + .await + .unwrap(); + + insert_confirmed_tx(&wallet, conflict_tx, 5); + insert_unconfirmed_tx(&wallet, live_candidate_tx); + + let block_id = + |height| BlockId { height, hash: bitcoin::BlockHash::from_byte_array([7u8; 32]) }; + let event = WalletEvent::ChainTipChanged { + old_tip: block_id(9), + new_tip: block_id(5 + ANTI_REORG_DELAY - 1), + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Pending); + assert!( + wallet.pending_payment_store.get(&payment_id).await.unwrap().is_some(), + "a candidate can still confirm, so the record must stay pending" + ); + } + + /// The failure write pair is record first, entry second: a crash in between leaves a + /// `Failed` record with a lingering entry. The next tip pass must finish the job — remove + /// the entry without disturbing the record. + #[tokio::test] + async fn a_failed_funding_payment_with_a_lingering_entry_is_cleaned_up() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let close_tx = wallet_paying_tx(&wallet, 3); + let close_txid = close_tx.compute_txid(); + + let splice_txid = Txid::from_byte_array([2u8; 32]); + let payment_id = PaymentId([21u8; 32]); + let mut recorded = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + recorded.status = PaymentStatus::Failed; + recorded.latest_update_timestamp = 7; + wallet.payment_store.insert_or_update(recorded).await.unwrap(); + + // The entry embeds the pre-failure snapshot, as a crash between the two writes leaves it. + let snapshot = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + let candidates = vec![FundingTxCandidate { + txid: splice_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }]; + let entry = PendingPaymentDetails::new(snapshot, vec![close_txid], candidates); + wallet.pending_payment_store.insert_or_update(entry).await.unwrap(); + + insert_confirmed_tx(&wallet, close_tx, 5); + + let block_id = + |height| BlockId { height, hash: bitcoin::BlockHash::from_byte_array([7u8; 32]) }; + let event = WalletEvent::ChainTipChanged { + old_tip: block_id(9), + new_tip: block_id(5 + ANTI_REORG_DELAY - 1), + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Failed); + assert_eq!(payment.latest_update_timestamp, 7, "the repair pass must not rewrite"); + assert!( + wallet.pending_payment_store.get(&payment_id).await.unwrap().is_none(), + "the lingering entry must be removed" + ); + } + + /// A crash between the failure's record write and its entry removal loses the wallet + /// changeset too, so the restart's catch-up sync replays the same events: `TxReplaced` for + /// the dead funding transaction resolves through the lingering entry to the already-`Failed` + /// record. Re-embedding that record would stamp `Failed` into the entry and hide it from the + /// pending listing that repairs it; the replay must instead finish the interrupted removal. + #[tokio::test] + async fn replayed_replacement_finishes_an_interrupted_failure() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let close_tx = wallet_paying_tx(&wallet, 3); + let close_txid = close_tx.compute_txid(); + + let splice_txid = Txid::from_byte_array([2u8; 32]); + let payment_id = PaymentId([21u8; 32]); + let mut recorded = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + recorded.status = PaymentStatus::Failed; + recorded.latest_update_timestamp = 7; + wallet.payment_store.insert_or_update(recorded).await.unwrap(); + + let snapshot = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + let candidates = vec![FundingTxCandidate { + txid: splice_txid, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }]; + let entry = PendingPaymentDetails::new(snapshot, vec![close_txid], candidates); + wallet.pending_payment_store.insert_or_update(entry).await.unwrap(); + + insert_confirmed_tx(&wallet, close_tx, 5); + + let block_id = + |height| BlockId { height, hash: bitcoin::BlockHash::from_byte_array([7u8; 32]) }; + let events = vec![ + WalletEvent::TxReplaced { + txid: splice_txid, + tx: Arc::new(dummy_tx()), + conflicts: vec![(0, close_txid)], + }, + WalletEvent::ChainTipChanged { + old_tip: block_id(9), + new_tip: block_id(5 + ANTI_REORG_DELAY - 1), + }, + ]; + wallet.update_payment_store(events).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Failed); + assert_eq!(payment.latest_update_timestamp, 7, "the replay must not rewrite the record"); + assert!( + wallet.pending_payment_store.get(&payment_id).await.unwrap().is_none(), + "the replay must finish the interrupted entry removal" + ); + } + + /// A funding record's id is anchored to its first candidate's txid. Once the payment settles + /// and its entry is removed, a wallet event for that candidate no longer resolves through the + /// candidate history — the fallback keys it by its own txid, colliding with the record's id. + /// Recording the event there would merge a fresh wallet-view `Pending` payment into the + /// terminal record; such events must be skipped. + #[tokio::test] + async fn candidate_event_does_not_resurrect_a_settled_funding_payment() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + // The record's id derives from the first candidate r1; its txid rotated to the RBF round + // r2. The payment failed and its pending entry is gone. + let r1 = Txid::from_byte_array([2u8; 32]); + let r2 = Txid::from_byte_array([4u8; 32]); + let payment_id = PaymentId(r1.to_byte_array()); + let mut recorded = interactive_funding_details(payment_id, r2, Some(1_000_000), Some(600)); + recorded.status = PaymentStatus::Failed; + recorded.latest_update_timestamp = 7; + wallet.payment_store.insert_or_update(recorded).await.unwrap(); + + // r1 reappears in the mempool after the failure... + let event = + WalletEvent::TxUnconfirmed { txid: r1, tx: Arc::new(dummy_tx()), old_block_time: None }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Failed, "the record must not resurrect"); + assert!(matches!(payment.kind, PaymentKind::Onchain { txid, .. } if txid == r2)); + assert_eq!(payment.latest_update_timestamp, 7); + assert!(wallet.pending_payment_store.get(&payment_id).await.unwrap().is_none()); + + // ...and even confirms: the record settled as `Failed` and must stay that way. + let event = WalletEvent::TxConfirmed { + txid: r1, + tx: Arc::new(dummy_tx()), + block_time: confirmed_block_time(5), + old_block_time: None, + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Failed, "the record must not resurrect"); + assert!(matches!(payment.kind, PaymentKind::Onchain { txid, .. } if txid == r2)); + assert_eq!(payment.latest_update_timestamp, 7); + assert!(wallet.pending_payment_store.get(&payment_id).await.unwrap().is_none()); + } + + /// The failure transition must apply regardless of the payment's direction: a splice-out + /// records as `Inbound` (funds return to the wallet) and dies to a conflicting close the + /// same way an outbound one does. + #[tokio::test] + async fn inbound_funding_payment_fails_once_a_foreign_conflict_confirms_to_depth() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let close_tx = wallet_paying_tx(&wallet, 3); + 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 mut details = + interactive_funding_details(payment_id, splice_txid, Some(1_000_000), Some(500)); + details.direction = PaymentDirection::Inbound; + wallet.persist_funding_payment(details, candidates).await.unwrap(); + wallet + .pending_payment_store + .update(PendingPaymentDetailsUpdate { + id: payment_id, + payment_update: None, + conflicting_txids: Some(vec![close_txid]), + candidates: Vec::new(), + splice_intent: None, + }) + .await + .unwrap(); + + insert_confirmed_tx(&wallet, close_tx, 5); + + let block_id = + |height| BlockId { height, hash: bitcoin::BlockHash::from_byte_array([7u8; 32]) }; + let event = WalletEvent::ChainTipChanged { + old_tip: block_id(9), + new_tip: block_id(5 + ANTI_REORG_DELAY - 1), + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + + let payment = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + assert_eq!(payment.status, PaymentStatus::Failed); + assert!(wallet.pending_payment_store.get(&payment_id).await.unwrap().is_none()); + } + + /// 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 + /// contribution, or none of the moved funds are the wallet's — would otherwise come back as + /// a spurious zero-amount record that nothing ever confirms. + #[tokio::test] + async fn funding_broadcast_without_wallet_activity_is_not_recorded() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let counterparty_node_id = PublicKey::from_str( + "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(); + let channels = vec![(counterparty_node_id, ChannelId([7u8; 32]))]; + let tx_type = TransactionType::Funding { channels: vec![] }; + + // No inputs or outputs involve the wallet: nothing to record. + wallet.classify_funding(&dummy_tx(), &channels, tx_type.clone()).await.unwrap(); + assert!(wallet.payment_store.list_page(None).await.unwrap().objects.is_empty()); + assert!(wallet.pending_payment_store.list_filter(|_| true).await.is_empty()); + + // A computable fee is not wallet participation. The wallet can resolve a splice's shared + // input whenever the previous funding transaction touched it (e.g. it funded the original + // channel open), so it derives the splice's fee even when no wallet funds move. + let prev_funding_outpoint = OutPoint { txid: Txid::from_byte_array([8u8; 32]), vout: 0 }; + wallet.inner.lock().unwrap().insert_txout( + prev_funding_outpoint, + TxOut { value: Amount::from_sat(100_000), script_pubkey: ScriptBuf::new() }, + ); + let splice_tx = Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: LockTime::ZERO, + input: vec![bitcoin::TxIn { + previous_output: prev_funding_outpoint, + ..Default::default() + }], + output: vec![TxOut { + value: Amount::from_sat(99_000), + script_pubkey: ScriptBuf::new(), + }], + }; + wallet.classify_funding(&splice_tx, &channels, tx_type.clone()).await.unwrap(); + assert!(wallet.payment_store.list_page(None).await.unwrap().objects.is_empty()); + + // Control: a funding transaction the wallet participates in is still recorded. + let script_pubkey = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + let funded_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 }], + }; + wallet.classify_funding(&funded_tx, &channels, tx_type).await.unwrap(); + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1); + match &payments[0].kind { + PaymentKind::Onchain { txid, .. } => assert_eq!(*txid, funded_tx.compute_txid()), + kind => panic!("unexpected kind {:?}", kind), + } + } + + /// A funding record's PaymentId is generated at record creation instead of being derived from a + /// txid: a replaceable transaction's txid is no stable identity for the record. Every lookup + /// resolves the record through its txid history (current txid, candidates, conflicts) rather + /// than re-deriving the id, so nothing may rely on the id and the txid coinciding. + #[tokio::test] + async fn funding_record_is_keyed_by_a_generated_id() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let counterparty_node_id = PublicKey::from_str( + "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(); + let channels = vec![(counterparty_node_id, ChannelId([7u8; 32]))]; + let tx_type = TransactionType::Funding { channels: vec![] }; + + let script_pubkey = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + let funded_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 txid = funded_tx.compute_txid(); + wallet.classify_funding(&funded_tx, &channels, tx_type).await.unwrap(); + + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1); + let record = &payments[0]; + assert_ne!(record.id, PaymentId(txid.to_byte_array()), "the id must not be the txid"); + match &record.kind { + PaymentKind::Onchain { txid: kind_txid, .. } => assert_eq!(*kind_txid, txid), + kind => panic!("unexpected kind {:?}", kind), + } + // The pending entry shares the id, and txid lookups resolve to the record. + assert!(wallet.pending_payment_store.get(&record.id).await.unwrap().is_some()); + assert_eq!(wallet.find_payment_by_txid(txid).await.unwrap(), Some(record.id)); + } + + /// A funding transaction classified again — e.g. a 0conf splice re-broadcast through LDK's + /// generic funding path after a restart — must resolve to the record's generated id rather + /// than create a second record for the same transaction. + #[tokio::test] + async fn funding_rebroadcast_resolves_to_the_generated_id() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let counterparty_node_id = PublicKey::from_str( + "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(); + let channels = vec![(counterparty_node_id, ChannelId([7u8; 32]))]; + + let script_pubkey = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + let funded_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 txid = funded_tx.compute_txid(); + + // The record the interactive-funding classification created, keyed by a generated id. + let payment_id = PaymentId([42u8; 32]); + 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)); + wallet.persist_funding_payment(details, candidates).await.unwrap(); + + // The re-typed rebroadcast comes back through the generic funding path. + wallet + .classify_funding(&funded_tx, &channels, TransactionType::Funding { channels: vec![] }) + .await + .unwrap(); + + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1, "the rebroadcast must not create a second record"); + assert_eq!(payments[0].id, payment_id); + // The interactive classification and contribution figures survive the generic + // wallet-view update (`funding_reclassification_update` declines the downgrade). + assert!(matches!( + payments[0].kind, + PaymentKind::Onchain { tx_type: Some(TransactionType::InteractiveFunding { .. }), .. } + )); + assert_eq!(payments[0].amount_msat, Some(1_000_000)); + } + + /// LDK re-broadcasts a promoted-but-unconfirmed 0conf splice through its generic funding + /// path: same txid, but typed as a plain funding transaction with wallet-view figures and no + /// contribution metadata. The rebroadcast must not overwrite the contribution-derived + /// figures or the interactive-funding classification — neither while the record is + /// unconfirmed nor once it confirmed under that same txid, where updates naming the + /// confirmed txid may otherwise move figures. + #[tokio::test] + async fn funding_rebroadcast_keeps_interactive_funding_classification() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + // The rebroadcast passes the wallet-activity guard: a splice-in funds the new channel + // output partly from the wallet, so the wallet sees movement. + 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 txid = tx.compute_txid(); + let payment_id = PaymentId(txid.to_byte_array()); + + 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)); + wallet.persist_funding_payment(details, candidates).await.unwrap(); + + let counterparty_node_id = PublicKey::from_str( + "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(); + let channels = vec![(counterparty_node_id, ChannelId([7u8; 32]))]; + let tx_type = TransactionType::Funding { channels: vec![] }; + + async fn assert_unchanged(wallet: &Wallet, payment_id: PaymentId, confirmed: bool) { + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1, "the rebroadcast must not mint a second record"); + let payment = &payments[0]; + assert_eq!(payment.id, payment_id); + assert_eq!(payment.amount_msat, Some(1_000_000)); + assert_eq!(payment.fee_paid_msat, Some(500)); + match &payment.kind { + PaymentKind::Onchain { + status, + tx_type: Some(TransactionType::InteractiveFunding { .. }), + .. + } => assert_eq!(matches!(status, ConfirmationStatus::Confirmed { .. }), confirmed), + kind => panic!("unexpected kind {:?}", kind), + } + } + + wallet.classify_funding(&tx, &channels, tx_type.clone()).await.unwrap(); + assert_unchanged(&wallet, payment_id, false).await; + + // Confirm the record, then replay the rebroadcast: a monitor-update completion can race + // wallet sync around confirmation. + let event = WalletEvent::TxConfirmed { + txid, + tx: Arc::new(tx.clone()), + block_time: confirmed_block_time(5), + old_block_time: None, + }; + wallet.update_payment_store(vec![event]).await.unwrap(); + wallet.classify_funding(&tx, &channels, tx_type).await.unwrap(); + assert_unchanged(&wallet, payment_id, true).await; + } + + /// A user-initiated splice's record is keyed by the PaymentId chosen at splice time, not by + /// its funding txid. The generic funding path must resolve a rebroadcast of that funding tx + /// back to the existing record rather than creating a duplicate under the txid-derived id. + #[tokio::test] + async fn classify_funding_resolves_the_splice_time_payment_id() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + 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 txid = tx.compute_txid(); + + let payment_id = PaymentId([21u8; 32]); + 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)); + wallet.persist_funding_payment(details, candidates).await.unwrap(); + + let counterparty_node_id = PublicKey::from_str( + "0279be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798", + ) + .unwrap(); + let channels = vec![(counterparty_node_id, ChannelId([7u8; 32]))]; + let tx_type = TransactionType::Funding { channels: vec![] }; + wallet.classify_funding(&tx, &channels, tx_type).await.unwrap(); + + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1, "the rebroadcast must not create a second record"); + assert_eq!(payments[0].id, payment_id); + assert_eq!(payments[0].amount_msat, Some(1_000_000)); + assert_eq!(payments[0].fee_paid_msat, Some(500)); + } + + /// 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 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]))], + }, + )]); + + // Wait until the loop has actually failed a classification write; re-enabling writes + // before the first attempt would let the first attempt succeed and the test pass + // without any retry happening. A failed classification must not leave a partial + // record behind. + let mut failed_writes = 0; + for _ in 0..100 { + tokio::time::sleep(Duration::from_millis(100)).await; + failed_writes = fail_store.failed_writes.load(Ordering::Acquire); + if failed_writes > 0 { + break; + } + } + assert!(failed_writes > 0, "classification never attempted a payment-store write"); + assert!(wallet.payment_store.list_page(None).await.unwrap().objects.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_page(None).await.unwrap().objects; + 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(); + } + + /// A package awaiting a classification retry must die when the node stops. When the retry + /// was a detached task, it outlived the broadcast loop: its re-send into the still-open + /// queue succeeded after `stop()`, so a later `start()` would classify and broadcast the + /// stale package. + #[tokio::test] + async fn failed_classification_retry_dies_at_stop() { + 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)); + + 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 + }); + + 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 and wait for the loop to + // fail a classification attempt, leaving a retry pending. + fail_store.fail_writes.store(true, Ordering::Release); + wallet.broadcaster.broadcast_transactions(&[( + &tx, + LdkTransactionType::Funding { + channels: vec![(counterparty_node_id, ChannelId([7u8; 32]))], + }, + )]); + let mut failed_writes = 0; + for _ in 0..100 { + tokio::time::sleep(Duration::from_millis(100)).await; + failed_writes = fail_store.failed_writes.load(Ordering::Acquire); + if failed_writes > 0 { + break; + } + } + assert!(failed_writes > 0, "classification never attempted a payment-store write"); + + // Stop the node with the retry still pending, then bring the loop back up with + // working persistence, as a stop()/start() cycle would. + stop_sender.send(()).unwrap(); + loop_task.await.unwrap(); + fail_store.fail_writes.store(false, Ordering::Release); + + 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 + }); + + // Watch well past the retry delay: the package from before the stop must not be + // classified or broadcast by the restarted loop. + for _ in 0..40 { + tokio::time::sleep(Duration::from_millis(100)).await; + assert!( + wallet.payment_store.list_page(None).await.unwrap().objects.is_empty(), + "a package from before stop() resurfaced after restart" + ); + } + + stop_sender.send(()).unwrap(); + loop_task.await.unwrap(); + } + + /// A queued classification can retry after a newer candidate of the same funding already + /// classified: the retry carries the candidate history as of its own broadcast, which no + /// longer includes the newer candidate. Applying it would rotate the record's txid backwards + /// and shrink the stored candidate history, after which the newer transaction can no longer + /// be mapped back to the record and wallet sync would file it as a foreign duplicate. + #[tokio::test] + async fn stale_classification_retry_keeps_the_newer_candidate() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let txid_a = Txid::from_byte_array([1u8; 32]); + let txid_b = Txid::from_byte_array([2u8; 32]); + // The record's id is anchored to the first negotiated candidate, so the stale retry + // resolves to the same record. + let payment_id = PaymentId(txid_a.to_byte_array()); + let candidate_a = FundingTxCandidate { + txid: txid_a, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }; + let candidate_b = FundingTxCandidate { + txid: txid_b, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(999), + }; + + // The bump candidate B classifies first, carrying the full history [A, B]. + let fresh = interactive_funding_details(payment_id, txid_b, Some(1_000_000), Some(999)); + wallet + .persist_funding_payment(fresh, vec![candidate_a.clone(), candidate_b.clone()]) + .await + .unwrap(); + + // The queued classification of A retries, carrying the history as of A's broadcast. + let stale = interactive_funding_details(payment_id, txid_a, Some(1_000_000), Some(500)); + wallet.persist_funding_payment(stale, vec![candidate_a.clone()]).await.unwrap(); + + let record = wallet.payment_store.get(&payment_id).await.unwrap().unwrap(); + match &record.kind { + PaymentKind::Onchain { txid, .. } => { + assert_eq!(*txid, txid_b, "the stale retry must not rotate the record back"); + }, + kind => panic!("unexpected kind {:?}", kind), + } + assert_eq!(record.fee_paid_msat, Some(999)); + + let pending = wallet.pending_payment_store.get(&payment_id).await.unwrap().unwrap(); + let PendingPaymentDetails::Tracked { candidates, .. } = &pending else { + panic!("unexpected variant {:?}", pending); + }; + assert_eq!( + *candidates, + vec![candidate_a, candidate_b], + "the stale retry must not shrink the candidate history" + ); + + // The consequence the history protects against: B must stay mapped to the record, or + // wallet sync would file it as a foreign duplicate. + assert_eq!(wallet.find_payment_by_txid(txid_b).await.unwrap(), Some(payment_id)); + } + + /// A missing pending entry is normally recreated from the incoming classification — but not + /// from a stale retry, whose truncated candidate history would otherwise slip past the merge + /// path's refusal. Recreation is left to a fresh classification instead. + #[tokio::test] + async fn stale_classification_retry_does_not_recreate_the_pending_entry() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let txid_a = Txid::from_byte_array([1u8; 32]); + let txid_b = Txid::from_byte_array([2u8; 32]); + let payment_id = PaymentId(txid_a.to_byte_array()); + let candidate_a = FundingTxCandidate { + txid: txid_a, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }; + let candidate_b = FundingTxCandidate { + txid: txid_b, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(999), + }; + + // The newer round B classified, but its write pair was torn by the same store failure + // that queued this retry: the record exists, the pending entry does not. + let recorded = interactive_funding_details(payment_id, txid_b, Some(1_000_000), Some(999)); + wallet.payment_store.insert(recorded).await.unwrap(); + + // The queued classification of A retries with its pre-B history. + let stale = interactive_funding_details(payment_id, txid_a, Some(1_000_000), Some(500)); + wallet.persist_funding_payment(stale, vec![candidate_a.clone()]).await.unwrap(); + assert!( + wallet.pending_payment_store.get(&payment_id).await.unwrap().is_none(), + "a stale retry must not recreate the pending entry from its truncated history" + ); + + // B's own retry recreates the entry with the full history. + let fresh = interactive_funding_details(payment_id, txid_b, Some(1_000_000), Some(999)); + wallet + .persist_funding_payment(fresh, vec![candidate_a.clone(), candidate_b.clone()]) + .await + .unwrap(); + let pending = wallet.pending_payment_store.get(&payment_id).await.unwrap().unwrap(); + let PendingPaymentDetails::Tracked { candidates, .. } = &pending else { + panic!("unexpected variant {:?}", pending); + }; + assert_eq!(*candidates, vec![candidate_a, candidate_b]); + } + + /// Wallet sync can record a genuine replacement round before classification records it as a + /// candidate — e.g. the counterparty broadcast a round whose classification failed here and + /// is still being retried. The funding-status gate then routes the round's confirmation to a + /// duplicate record keyed by the round's txid, whose pending entry shadows the funding + /// record in `find_payment_by_txid`'s direct probe. Once the round's classification lands, + /// it must merge the duplicate — adopt its confirmation and remove it — so a single record + /// tracks the splice. + #[tokio::test] + async fn classification_merges_duplicate_records_for_its_candidates() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let funding_id = PaymentId([21u8; 32]); + let txid1 = Txid::from_byte_array([1u8; 32]); + let txid2 = Txid::from_byte_array([2u8; 32]); + + // Round 1 classified normally. + let round1 = vec![FundingTxCandidate { + txid: txid1, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }]; + let details = interactive_funding_details(funding_id, txid1, Some(1_000_000), Some(500)); + wallet.persist_funding_payment(details, round1).await.unwrap(); + + // Wallet sync recorded round 2's confirmation while the round was not yet a candidate: a + // duplicate untyped record under the txid-derived id, plus its pending entry. + let duplicate_id = PaymentId(txid2.to_byte_array()); + let duplicate = PaymentDetails::new( + duplicate_id, + PaymentKind::Onchain { txid: txid2, status: confirmed_status(), tx_type: None }, + Some(999_000), + Some(999), + PaymentDirection::Outbound, + PaymentStatus::Pending, + ); + wallet.payment_store.insert_or_update(duplicate.clone()).await.unwrap(); + wallet + .pending_payment_store + .insert_or_update(PendingPaymentDetails::new(duplicate, Vec::new(), Vec::new())) + .await + .unwrap(); + assert_eq!(wallet.find_payment_by_txid(txid2).await.unwrap(), Some(duplicate_id)); + + // Round 2's classification lands (e.g. retried after a persistence failure). + let rounds = vec![ + FundingTxCandidate { + txid: txid1, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }, + FundingTxCandidate { + txid: txid2, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(400), + }, + ]; + let details = interactive_funding_details(funding_id, txid2, Some(1_000_000), Some(400)); + wallet.persist_funding_payment(details, rounds).await.unwrap(); + + // One record: the funding record carries the duplicate's confirmation and the confirmed + // candidate's figures; the duplicate and its pending entry are gone, so the round's txid + // resolves to the funding record again. + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1, "the duplicate must be merged away"); + let payment = &payments[0]; + assert_eq!(payment.id, funding_id); + assert_eq!(payment.amount_msat, Some(1_000_000)); + assert_eq!(payment.fee_paid_msat, Some(400)); + match &payment.kind { + PaymentKind::Onchain { + txid, + status: ConfirmationStatus::Confirmed { .. }, + tx_type: Some(TransactionType::InteractiveFunding { .. }), + } => assert_eq!(*txid, txid2), + kind => panic!("unexpected kind {:?}", kind), + } + assert!(wallet.pending_payment_store.get(&duplicate_id).await.unwrap().is_none()); + assert_eq!(wallet.find_payment_by_txid(txid2).await.unwrap(), Some(funding_id)); + } + + /// A duplicate for an *unconfirmed* round carries no state the funding record needs: the + /// merge removes it without touching the record's active txid or figures, and the round's + /// txid maps back to the funding record through its candidate history. + #[tokio::test] + async fn classification_drops_unconfirmed_duplicates_without_adopting_their_txid() { + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let wallet = new_test_wallet(store, false).await; + + let funding_id = PaymentId([21u8; 32]); + let txid1 = Txid::from_byte_array([1u8; 32]); + let txid2 = Txid::from_byte_array([2u8; 32]); + + // Wallet sync saw round 1 — still unconfirmed — before any classification ran. + let duplicate_id = PaymentId(txid1.to_byte_array()); + let duplicate = PaymentDetails::new( + duplicate_id, + PaymentKind::Onchain { + txid: txid1, + status: ConfirmationStatus::Unconfirmed, + tx_type: None, + }, + Some(999_000), + Some(999), + PaymentDirection::Outbound, + PaymentStatus::Pending, + ); + wallet.payment_store.insert_or_update(duplicate.clone()).await.unwrap(); + wallet + .pending_payment_store + .insert_or_update(PendingPaymentDetails::new(duplicate, Vec::new(), Vec::new())) + .await + .unwrap(); + + // Round 2 is the active broadcast; its classification lists both rounds. + let rounds = vec![ + FundingTxCandidate { + txid: txid1, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }, + FundingTxCandidate { + txid: txid2, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(400), + }, + ]; + let details = interactive_funding_details(funding_id, txid2, Some(1_000_000), Some(400)); + wallet.persist_funding_payment(details, rounds).await.unwrap(); + + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1, "the duplicate must be merged away"); + let payment = &payments[0]; + assert_eq!(payment.id, funding_id); + // The record keeps tracking the actively-broadcast round; a duplicate that never confirmed + // has nothing to adopt. + match &payment.kind { + PaymentKind::Onchain { txid, status: ConfirmationStatus::Unconfirmed, .. } => { + assert_eq!(*txid, txid2) + }, + kind => panic!("unexpected kind {:?}", kind), + } + assert_eq!(payment.fee_paid_msat, Some(400)); + assert_eq!(wallet.find_payment_by_txid(txid1).await.unwrap(), Some(funding_id)); + } + + /// Removing the duplicate is two store writes, and the failure between them must leave a + /// state the classification retry can finish cleaning up. If the payment record went first, + /// a failure on the pending-entry removal would orphan that entry where the retry can no + /// longer discover it (the record lookup misses), and it would keep shadowing the funding + /// record in `find_payment_by_txid`'s direct probe — re-creating the duplicate problem with + /// no further classification pass coming to fix it. + #[tokio::test] + async fn classification_retry_completes_a_partially_failed_duplicate_removal() { + let fail_store = FailRemoveStore::new(); + let store: Arc = Arc::new(DynStoreWrapper(fail_store.clone())); + let wallet = new_test_wallet(store, false).await; + + let funding_id = PaymentId([21u8; 32]); + let txid1 = Txid::from_byte_array([1u8; 32]); + let txid2 = Txid::from_byte_array([2u8; 32]); + + // Round 1 classified normally. + let round1 = vec![FundingTxCandidate { + txid: txid1, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }]; + let details = interactive_funding_details(funding_id, txid1, Some(1_000_000), Some(500)); + wallet.persist_funding_payment(details, round1).await.unwrap(); + + // Wallet sync recorded round 2's confirmation while the round was not yet a candidate. + let duplicate_id = PaymentId(txid2.to_byte_array()); + let duplicate = PaymentDetails::new( + duplicate_id, + PaymentKind::Onchain { txid: txid2, status: confirmed_status(), tx_type: None }, + Some(999_000), + Some(999), + PaymentDirection::Outbound, + PaymentStatus::Pending, + ); + wallet.payment_store.insert_or_update(duplicate.clone()).await.unwrap(); + wallet + .pending_payment_store + .insert_or_update(PendingPaymentDetails::new(duplicate, Vec::new(), Vec::new())) + .await + .unwrap(); + + // Round 2's classification lands, but one of the duplicate's two removals fails. + let rounds = vec![ + FundingTxCandidate { + txid: txid1, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(500), + }, + FundingTxCandidate { + txid: txid2, + amount_msat: Some(1_000_000), + fee_paid_msat: Some(400), + }, + ]; + let details = interactive_funding_details(funding_id, txid2, Some(1_000_000), Some(400)); + fail_store.fail_next_remove_in(PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE); + let res = wallet.persist_funding_payment(details.clone(), rounds.clone()).await; + assert!(res.is_err(), "the injected remove failure must surface"); + + // The broadcast loop re-runs a failed classification; the retry must finish the cleanup. + wallet.persist_funding_payment(details, rounds).await.unwrap(); + + let payments = wallet.payment_store.list_page(None).await.unwrap().objects; + assert_eq!(payments.len(), 1, "the duplicate must be merged away"); + assert_eq!(payments[0].id, funding_id); + assert!(wallet.pending_payment_store.get(&duplicate_id).await.unwrap().is_none()); + assert_eq!(wallet.find_payment_by_txid(txid2).await.unwrap(), Some(funding_id)); } /// Barrier test, classification-first ordering: wallet sync's confirmation handling must diff --git a/tests/integration_tests_rust.rs b/tests/integration_tests_rust.rs index 5f4a95b7e..e7af6e8cc 100644 --- a/tests/integration_tests_rust.rs +++ b/tests/integration_tests_rust.rs @@ -81,6 +81,16 @@ async fn wait_for_classified_funding_payment(node: &Node, funding_txid: Txid) { }); } +/// Resolves the payment record whose current transaction is `funding_txid`. Funding records are +/// keyed by a random id generated at creation, so they are found through their transaction history +/// rather than by deriving an id from a txid. +fn funding_payment(node: &Node, funding_txid: Txid) -> PaymentDetails { + node.list_all_payments() + .into_iter() + .find(|p| matches!(p.kind, PaymentKind::Onchain { txid, .. } if txid == funding_txid)) + .expect("funding payment exists") +} + #[derive(Clone)] struct ContendedStore { inner: Arc, @@ -2093,9 +2103,7 @@ async fn splice_channel() { // them to the channel balance since there may not be a change output. let expected_splice_in_lightning_balance_sat = 4_000_002; - let payments = node_b.list_all_payments(); - let payment = - payments.into_iter().find(|p| p.id == PaymentId(txo.txid.to_byte_array())).unwrap(); + let payment = funding_payment(&node_b, txo.txid); assert_eq!(payment.fee_paid_msat, Some(expected_splice_in_fee_sat * 1_000)); assert_eq!( @@ -2146,9 +2154,7 @@ async fn splice_channel() { let expected_splice_out_fee_sat = 183; - let payments = node_a.list_all_payments(); - let payment = - payments.into_iter().find(|p| p.id == PaymentId(txo.txid.to_byte_array())).unwrap(); + let payment = funding_payment(&node_a, txo.txid); assert_eq!(payment.fee_paid_msat, Some(expected_splice_out_fee_sat * 1_000)); // The splice-out graduated to a confirmed interactive-funding payment. Its `direction` is left // unasserted on purpose: the destination is our own address, so it is a self-transfer (channel @@ -2443,15 +2449,20 @@ async fn run_rbf_splice_channel_test(confirm_original: bool) { // Node B contributed to this splice; wait for its classification before syncing so the sync // takes the funding short-circuit rather than racing the broadcaster's queue. wait_for_classified_funding_payment(&node_b, original_txo.txid).await; + // The record's random id is fixed at creation; capture it while the original candidate is + // current so its stability can be asserted across the RBF rounds below. + let splice_payment_id = funding_payment(&node_b, original_txo.txid).id; node_a.sync_wallets().unwrap(); node_b.sync_wallets().unwrap(); // For `confirm_original`, capture the original candidate's fee and raw transaction now, before // the RBF replaces it, so it can be force-confirmed (instead of the RBF) further below. let original_candidate: Option<(Option, String)> = if confirm_original { - let payment_id = PaymentId(original_txo.txid.to_byte_array()); - let fee = - node_b.payment(&payment_id).unwrap().expect("splice payment exists").fee_paid_msat; + let fee = node_b + .payment(&splice_payment_id) + .unwrap() + .expect("splice payment exists") + .fee_paid_msat; let raw_tx: String = bitcoind .client .call("getrawtransaction", &[json!(original_txo.txid.to_string())]) @@ -2484,12 +2495,11 @@ async fn run_rbf_splice_channel_test(confirm_original: bool) { node_b.sync_wallets().unwrap(); // After RBF but before confirmation, node_b (the initiator) should have a single on-chain - // payment covering both candidates: id anchored to the first broadcast, `kind.txid` pointing - // at the latest (RBF) candidate, and the durable interactive-funding `tx_type` preserved across - // the replacement. + // payment covering both candidates: still under the id it was created with, `kind.txid` + // pointing at the latest (RBF) candidate, and the durable interactive-funding `tx_type` + // preserved across the replacement. let rbf_candidate_fee = { - let payment_id = PaymentId(original_txo.txid.to_byte_array()); - let payment = node_b.payment(&payment_id).unwrap().expect("splice payment exists"); + let payment = node_b.payment(&splice_payment_id).unwrap().expect("splice payment exists"); match payment.kind { PaymentKind::Onchain { txid, @@ -2563,8 +2573,8 @@ async fn run_rbf_splice_channel_test(confirm_original: bool) { // channel-lifecycle signal, not what drives payment status. Its `kind.txid` reflects the // winning RBF candidate, and `fee_paid_msat` carries this node's `FundingContribution` fee. { - let payment_id = PaymentId(original_txo.txid.to_byte_array()); - let payment = node_b.payment(&payment_id).unwrap().expect("splice payment graduated"); + let payment = + node_b.payment(&splice_payment_id).unwrap().expect("splice payment graduated"); assert_eq!(payment.status, PaymentStatus::Succeeded); match payment.kind { PaymentKind::Onchain { txid, status: ConfirmationStatus::Confirmed { .. }, .. } => { @@ -2624,8 +2634,7 @@ async fn funding_payment_graduates_without_channel_ready() { // The funding payment is `Succeeded` purely from wallet sync reaching `ANTI_REORG_DELAY` // confirmations, asserted before draining any LDK event — so graduation is not driven by the // Lightning `ChannelReady` signal. - let payment_id = PaymentId(funding_txo.txid.to_byte_array()); - let payment = node_a.payment(&payment_id).unwrap().expect("funding payment exists"); + let payment = funding_payment(&node_a, funding_txo.txid); assert_eq!(payment.status, PaymentStatus::Succeeded); match payment.kind { PaymentKind::Onchain { @@ -2687,8 +2696,8 @@ async fn splice_payment_reorged_to_unconfirmed() { generate_blocks_and_wait(&bitcoind.client, &electrsd.client, 1).await; node_b.sync_wallets().unwrap(); - let payment_id = PaymentId(splice_txo.txid.to_byte_array()); - let payment = node_b.payment(&payment_id).unwrap().expect("splice payment exists"); + let payment = funding_payment(&node_b, splice_txo.txid); + let payment_id = payment.id; assert_eq!(payment.status, PaymentStatus::Pending); assert!(matches!( payment.kind,