-
Notifications
You must be signed in to change notification settings - Fork 161
Keep funding payment records accurate #1057
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
c7e8123
0a9f121
0206dee
6b438c6
a459cfa
c159844
bab0312
6730351
cdaeda7
c25e57d
7b60504
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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. | ||
| pub(crate) 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,97 @@ 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 keep a funding transaction off-chain until LDK re-hands | ||
| /// it when the channel next resumes — no timer re-broadcasts it, and the wallet's tip-change | ||
| /// re-broadcast covers recorded transactions only. | ||
| 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(); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Rather than adding a second queue on top, can we make the single queue do prioritization so that it skips entries that failed classification but keeps them in the queue or similar?
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. See #1057 (comment). The mpsc channel only allows push and pop, so we wouldn't be able to de-dup against LDK's periodic claim broadcasts if we re-queued. Or did you have something else in mind? |
||
| loop { | ||
| let tx_bcast_logger = Arc::clone(&self.logger); | ||
| tokio::select! { | ||
| let next_retry_at = retries.next_retry_at(); | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Not sure we need a
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Added |
||
| let package = tokio::select! { | ||
| // Polled in order: a stop request first, then a fresh package, and a due retry | ||
| // only when neither is ready, so retries never hold back the broadcasts that | ||
| // keep arriving during a store outage. The retry keeps its delay regardless: | ||
| // without it, an empty channel would retry a fast-failing store back to back, | ||
| // logging an error each time. | ||
| biased; | ||
| _ = 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.txids(), | ||
| ); | ||
| }, | ||
| ScheduleOutcome::AlreadyQueued(duplicate) => { | ||
| log_debug!( | ||
| self.logger, | ||
| "Dropped a re-broadcast package; an identical one already awaits a classification retry: {:?}", | ||
| duplicate.txids(), | ||
| ); | ||
| }, | ||
| ScheduleOutcome::Refused(package) => { | ||
| log_error!( | ||
| self.logger, | ||
| "Dropped a package failing classification; too many await retries: {:?}", | ||
| package.txids(), | ||
| ); | ||
| }, | ||
| } | ||
| } | ||
| } | ||
|
|
||
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.