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
113 changes: 80 additions & 33 deletions src/chain/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 {

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.

Why maintain separate immediate and retry paths instead of treating every broadcast as scheduled retryable work?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Refactor the duplicated code, but kept the paths separate. Now that we have a bounded queue and deduplication, using the same path would mean we'd drop newer packages.

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(),
);
},
}
}
}
Expand Down
59 changes: 56 additions & 3 deletions src/payment/pending_payment_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,9 +105,19 @@ impl StorableObject for PendingPaymentDetails {
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 {
// 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()
&& self.candidates != update.candidates
&& self.candidates.iter().all(extends_history)
{
self.candidates = update.candidates;
updated = true;
}
Expand Down Expand Up @@ -243,6 +253,49 @@ mod tests {
);
}

/// 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 stale_update = PendingPaymentDetailsUpdate {
id: payment_id,
payment_update: None,
conflicting_txids: None,
candidates: vec![candidate(txid_a, 400)],
};
assert!(!pending.update(stale_update), "a stale history must not shrink the stored one");
assert_eq!(pending.candidates, 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(),
};
assert!(pending.update(fresh_update));
assert_eq!(pending.candidates, extended);
}

#[test]
fn funding_classification_pending_update_preserves_mirrored_confirmation() {
use bitcoin::BlockHash;
Expand Down
Loading
Loading