Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions src/chain/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -577,15 +577,15 @@ impl ChainSource {
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.
// Prepare wallet transactions and classify funding broadcasts before sending.
// If either fails, broadcasting could race another spend or leave an on-chain
// transaction without a payment record.
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: {:?}",
"Skipping broadcast: failed to prepare transaction: {:?}",
e,
);
continue;
Expand Down
117 changes: 81 additions & 36 deletions src/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -551,6 +551,28 @@ impl Future for EventFuture {
}
}

fn discarded_funding_transaction(funding_info: FundingInfo) -> Option<bitcoin::Transaction> {
match funding_info {
FundingInfo::Tx { transaction } => Some(transaction),
FundingInfo::Contribution { inputs, outputs } => Some(bitcoin::Transaction {
version: bitcoin::transaction::Version::TWO,
lock_time: bitcoin::absolute::LockTime::ZERO,
input: inputs
.into_iter()
.map(|previous_output| bitcoin::TxIn {
previous_output,
..bitcoin::TxIn::default()
})
.collect(),
output: outputs
.into_iter()
.map(|script_pubkey| bitcoin::TxOut { value: bitcoin::Amount::ZERO, script_pubkey })
.collect(),
}),
FundingInfo::OutPoint { .. } => None,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Claude:

An LSPS2 manual-broadcast channel closed before broadcast lands here, so the inputs locked in create_funding_transaction are never released. The handler can fetch the transaction the service stored, or the wallet can keep the created transaction in the graph and cancel by txid.

}
}

pub(crate) struct EventHandler<L: Deref + Clone + Sync + Send + 'static>
where
L::Target: LdkLogger,
Expand Down Expand Up @@ -763,6 +785,7 @@ where
.await;
match funding_transaction {
Ok(final_tx) => {
let final_tx_for_cancel = final_tx.clone();
let needs_manual_broadcast = self
.liquidity_source
.lsps2_service()
Expand Down Expand Up @@ -793,27 +816,39 @@ where

match result {
Ok(()) => {},
Err(APIError::APIMisuseError { err }) => {
log_error!(
self.logger,
"Encountered APIMisuseError, this should never happen: {}",
err
);
debug_assert!(false, "APIMisuseError: {}", err);
},
Err(APIError::ChannelUnavailable { err }) => {
log_error!(
self.logger,
"Failed to process funding transaction as channel went away before we could fund it: {}",
err
);
},
Err(err) => {
log_error!(
self.logger,
"Failed to process funding transaction: {:?}",
err
);
if let Err(e) = self.wallet.cancel_tx(final_tx_for_cancel).await {
log_error!(
self.logger,
"Failed to release funding inputs: {}",
e
);
return Err(ReplayEvent());
}
match err {
APIError::APIMisuseError { err } => {
log_error!(
self.logger,
"Encountered APIMisuseError, this should never happen: {}",
err
);
debug_assert!(false, "APIMisuseError: {}", err);
},
APIError::ChannelUnavailable { err } => {
log_error!(
self.logger,
"Failed to process funding transaction as channel went away before we could fund it: {}",
err
);
},
err => {
log_error!(
self.logger,
"Failed to process funding transaction: {:?}",
err
);
},
}
},
}
},
Expand Down Expand Up @@ -1992,27 +2027,14 @@ where
}
},
LdkEvent::DiscardFunding { channel_id, funding_info } => {
if let FundingInfo::Contribution { inputs: _, outputs } = funding_info {
if let Some(tx) = discarded_funding_transaction(funding_info) {
log_info!(
self.logger,
"Reclaiming unused addresses from channel {} funding",
"Reclaiming unused wallet state from channel {} funding",
channel_id,
);

let tx = bitcoin::Transaction {
version: bitcoin::transaction::Version::TWO,
lock_time: bitcoin::absolute::LockTime::ZERO,
input: vec![],
output: outputs
.into_iter()
.map(|script_pubkey| bitcoin::TxOut {
value: bitcoin::Amount::ZERO,
script_pubkey,
})
.collect(),
};
if let Err(e) = self.wallet.cancel_tx(tx).await {
log_error!(self.logger, "Failed reclaiming unused addresses: {}", e);
log_error!(self.logger, "Failed reclaiming unused wallet state: {}", e);
return Err(ReplayEvent());
}
}
Expand Down Expand Up @@ -2283,6 +2305,7 @@ mod tests {
use std::sync::atomic::{AtomicU16, Ordering};
use std::time::Duration;

use bitcoin::hashes::Hash;
use lightning::util::test_utils::TestLogger;

use super::*;
Expand All @@ -2299,6 +2322,28 @@ mod tests {
}
}

#[test]
fn discarded_contribution_preserves_inputs_and_outputs() {
let inputs = vec![
OutPoint::new(bitcoin::Txid::from_byte_array([1; 32]), 2),
OutPoint::new(bitcoin::Txid::from_byte_array([3; 32]), 4),
];
let outputs =
vec![bitcoin::ScriptBuf::from_bytes(vec![5]), bitcoin::ScriptBuf::from_bytes(vec![6])];

let tx = discarded_funding_transaction(FundingInfo::Contribution {
inputs: inputs.clone(),
outputs: outputs.clone(),
})
.unwrap();

assert_eq!(tx.input.iter().map(|txin| txin.previous_output).collect::<Vec<_>>(), inputs,);
assert_eq!(
tx.output.iter().map(|txout| txout.script_pubkey.clone()).collect::<Vec<_>>(),
outputs,
);
}

#[test]
fn lsps2_payment_metadata_decodes_total_fee_limit() {
let metadata = PaymentMetadata {
Expand Down
7 changes: 7 additions & 0 deletions src/io/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,13 @@ pub(crate) const BDK_WALLET_TX_GRAPH_PRIMARY_NAMESPACE: &str = "bdk_wallet";
pub(crate) const BDK_WALLET_TX_GRAPH_SECONDARY_NAMESPACE: &str = "";
pub(crate) const BDK_WALLET_TX_GRAPH_KEY: &str = "tx_graph";

/// The BDK wallet's [`ChangeSet::locked_outpoints`] will be persisted under this key.
///
/// [`ChangeSet::locked_outpoints`]: bdk_wallet::ChangeSet::locked_outpoints
pub(crate) const BDK_WALLET_LOCKED_OUTPOINTS_PRIMARY_NAMESPACE: &str = "bdk_wallet";
pub(crate) const BDK_WALLET_LOCKED_OUTPOINTS_SECONDARY_NAMESPACE: &str = "";
pub(crate) const BDK_WALLET_LOCKED_OUTPOINTS_KEY: &str = "locked_outpoints";

/// The BDK wallet's [`ChangeSet::indexer`] will be persisted under this key.
///
/// [`ChangeSet::indexer`]: bdk_wallet::ChangeSet::indexer
Expand Down
13 changes: 13 additions & 0 deletions src/io/utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ use bdk_chain::local_chain::ChangeSet as BdkLocalChainChangeSet;
use bdk_chain::miniscript::{Descriptor, DescriptorPublicKey};
use bdk_chain::tx_graph::ChangeSet as BdkTxGraphChangeSet;
use bdk_chain::ConfirmationBlockTime;
use bdk_wallet::locked_outpoints::ChangeSet as BdkLockedOutpointsChangeSet;
use bdk_wallet::ChangeSet as BdkWalletChangeSet;
use bitcoin::Network;
use lightning::ln::msgs::DecodeError;
Expand Down Expand Up @@ -721,6 +722,15 @@ impl_read_write_change_set_type!(
BDK_WALLET_TX_GRAPH_KEY
);

impl_read_write_change_set_type!(
read_bdk_wallet_locked_outpoints,
write_bdk_wallet_locked_outpoints,
BdkLockedOutpointsChangeSet,
BDK_WALLET_LOCKED_OUTPOINTS_PRIMARY_NAMESPACE,
BDK_WALLET_LOCKED_OUTPOINTS_SECONDARY_NAMESPACE,
BDK_WALLET_LOCKED_OUTPOINTS_KEY
);

impl_read_write_change_set_type!(
read_bdk_wallet_indexer,
write_bdk_wallet_indexer,
Expand Down Expand Up @@ -763,6 +773,9 @@ pub(crate) async fn read_bdk_wallet_change_set(
read_bdk_wallet_tx_graph(&*kv_store, logger)
.await?
.map(|tx_graph| change_set.tx_graph = tx_graph);
read_bdk_wallet_locked_outpoints(&*kv_store, logger)
.await?
.map(|locked_outpoints| change_set.locked_outpoints = locked_outpoints);
read_bdk_wallet_indexer(&*kv_store, logger).await?.map(|indexer| change_set.indexer = indexer);
Ok(Some(change_set))
}
Expand Down
21 changes: 18 additions & 3 deletions src/tx_broadcaster.rs
Original file line number Diff line number Diff line change
Expand Up @@ -133,15 +133,30 @@ 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.
/// Prepares a queued package in the wallet, classifies it into payment records, and returns the
/// package ready for the chain client. Returns `Err` if preparation or classification fails;
/// callers must not broadcast the package in that case.
pub(crate) async fn classify_package(
&self, package: BroadcastPackage,
) -> Result<BroadcastPackage, 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() {
let should_broadcast = match tx_type {
Some(LdkTransactionType::Funding { .. }) => {
wallet.prepare_funding_broadcast(tx).await?
},
Comment on lines +146 to +148

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note that #1057 removes classify_package and wallet` from the broadcaster, so we may want to consider landing that first. This would move to the event handler, IIUC.

None => wallet.prepare_unclassified_broadcast(tx).await?,
_ => true,
};
if !should_broadcast {
log_error!(
self.logger,
"Skipping broadcast of {} because an input is no longer available",
tx.compute_txid(),
);
return Err(Error::WalletOperationFailed);
}
if let Some(tx_type) = tx_type {
wallet.classify_broadcast(tx, tx_type).await?;
}
Expand Down
Loading
Loading