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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Fix two funding payment record bugs by jkczyz · Pull Request #1057 · lightningdevkit/ldk-node · GitHub
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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix two funding payment record bugs by jkczyz · Pull Request #1057 · lightningdevkit/ldk-node · GitHub
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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix two funding payment record bugs by jkczyz · Pull Request #1057 · lightningdevkit/ldk-node · GitHub
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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' Fix two funding payment record bugs by jkczyz · Pull Request #1057 · lightningdevkit/ldk-node · GitHub
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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix two funding payment record bugs by jkczyz · Pull Request #1057 · lightningdevkit/ldk-node · GitHub
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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' Fix two funding payment record bugs by jkczyz · Pull Request #1057 · lightningdevkit/ldk-node · GitHub
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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); Fix two funding payment record bugs by jkczyz · Pull Request #1057 · lightningdevkit/ldk-node · GitHub
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
87 changes: 59 additions & 28 deletions src/chain/mod.rs
Original file line numberDiff line numberDiff 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;
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 DownExpand Up@@ -562,12 +568,53 @@ 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 parked: Vec<(tokio::time::Instant, BroadcastPackage)> = Vec::new();
loop {
let tx_bcast_logger = Arc::clone(&self.logger);
// Entries are appended with a fixed delay, so the first is always the next due.
let next_retry_at = parked.first().map(|(deadline, _)| *deadline);
tokio::select! {
_ = stop_tx_bcast_receiver.changed() => {
log_debug!(
Expand All@@ -577,34 +624,18 @@ 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.
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
},
if let Err(package) = self.classify_and_broadcast(next_package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
_ = tokio::time::sleep_until(
next_retry_at.unwrap_or_else(tokio::time::Instant::now)
), if next_retry_at.is_some() => {
let (_, package) = parked.remove(0);
if let Err(package) = self.classify_and_broadcast(package).await {
let retry_at = tokio::time::Instant::now() + FAILED_CLASSIFY_RETRY_DELAY;
parked.push((retry_at, package));
}
}
}
Expand Down
12 changes: 5 additions & 7 deletions src/tx_broadcaster.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -133,12 +133,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<BroadcastPackage, Error> {
/// 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() {
Expand All@@ -147,7 +145,7 @@ where
}
}
}
Ok(package)
Ok(())
}

pub(crate) fn broadcast_unclassified_transaction(&self, tx: Transaction) {
Expand Down
Loading
Loading