Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
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" + '
Use BDK events in `update_payment_store` by Camillarhi · Pull Request #658 · lightningdevkit/ldk-node · GitHub
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
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('^' + ".*" + ' Use BDK events in `update_payment_store` by Camillarhi · Pull Request #658 · lightningdevkit/ldk-node · GitHub
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
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('^' + ".*" + ' Use BDK events in `update_payment_store` by Camillarhi · Pull Request #658 · lightningdevkit/ldk-node · GitHub
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
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" + ' Use BDK events in `update_payment_store` by Camillarhi · Pull Request #658 · lightningdevkit/ldk-node · GitHub
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
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('^' + ".*" + ' Use BDK events in `update_payment_store` by Camillarhi · Pull Request #658 · lightningdevkit/ldk-node · GitHub
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
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('^' + ".*" + ' Use BDK events in `update_payment_store` by Camillarhi · Pull Request #658 · lightningdevkit/ldk-node · GitHub
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
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); } })(); })(); Use BDK events in `update_payment_store` by Camillarhi · Pull Request #658 · lightningdevkit/ldk-node · GitHub
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,7 +54,7 @@ lightning-macros = { git = "https://github.com/lightningdevkit/rust-lightning",
bdk_chain = { version = "0.23.0", default-features = false, features = ["std"] }
bdk_esplora = { version = "0.22.0", default-features = false, features = ["async-https-rustls", "tokio"]}
bdk_electrum = { version = "0.23.0", default-features = false, features = ["use-rustls-ring"]}
bdk_wallet = { version = "2.2.0", default-features = false, features = ["std", "keys-bip39"]}
bdk_wallet = { version = "2.3.0", default-features = false, features = ["std", "keys-bip39"]}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Please make the BDK bump a dedicated commit.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Thanks, this has been updated


bitreq = { version = "0.3", default-features = false, features = ["async-https"] }
rustls = { version = "0.23", default-features = false }
Expand Down
39 changes: 29 additions & 10 deletions src/builder.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -55,12 +55,14 @@ use crate::gossip::GossipSource;
use crate::io::sqlite_store::SqliteStore;
use crate::io::utils::{
read_event_queue, read_external_pathfinding_scores_from_cache, read_network_graph,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_scorer,
write_node_metrics,
read_node_metrics, read_output_sweeper, read_payments, read_peer_info, read_pending_payments,
read_scorer, write_node_metrics,
};
use crate::io::vss_store::VssStoreBuilder;
use crate::io::{
self, PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE, PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
};
use crate::liquidity::{
LSPS1ClientConfig, LSPS2ClientConfig, LSPS2ServiceConfig, LiquiditySourceBuilder,
Expand All@@ -73,8 +75,8 @@ use crate::runtime::{Runtime, RuntimeSpawner};
use crate::tx_broadcaster::TransactionBroadcaster;
use crate::types::{
AsyncPersister, ChainMonitor, ChannelManager, DynStore, DynStoreWrapper, GossipSync, Graph,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, Persister,
SyncAndAsyncKVStore,
KeysManager, MessageRouter, OnionMessenger, PaymentStore, PeerManager, PendingPaymentStore,
Persister, SyncAndAsyncKVStore,
};
use crate::wallet::persist::KVStoreWalletPersister;
use crate::wallet::Wallet;
Expand DownExpand Up@@ -1057,12 +1059,14 @@ fn build_with_store_internal(

let kv_store_ref = Arc::clone(&kv_store);
let logger_ref = Arc::clone(&logger);
let (payment_store_res, node_metris_res) = runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
)
});
let (payment_store_res, node_metris_res, pending_payment_store_res) =
runtime.block_on(async move {
tokio::join!(
read_payments(&*kv_store_ref, Arc::clone(&logger_ref)),
read_node_metrics(&*kv_store_ref, Arc::clone(&logger_ref)),
read_pending_payments(&*kv_store_ref, Arc::clone(&logger_ref))
)
});

// Initialize the status fields.
let node_metrics = match node_metris_res {
Expand DownExpand Up@@ -1243,6 +1247,20 @@ fn build_with_store_internal(
},
};

let pending_payment_store = match pending_payment_store_res {
Ok(pending_payments) => Arc::new(PendingPaymentStore::new(
pending_payments,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE.to_string(),
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE.to_string(),
Arc::clone(&kv_store),
Arc::clone(&logger),
)),
Err(e) => {
log_error!(logger, "Failed to read pending payment data from store: {}", e);
return Err(BuildError::ReadFailed);
},
};

let wallet = Arc::new(Wallet::new(
bdk_wallet,
wallet_persister,
Expand All@@ -1251,6 +1269,7 @@ fn build_with_store_internal(
Arc::clone(&payment_store),
Arc::clone(&config),
Arc::clone(&logger),
Arc::clone(&pending_payment_store),
));

// Initialize the KeysManager
Expand Down
4 changes: 4 additions & 0 deletions src/data_store.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -167,6 +167,10 @@ where
})?;
Ok(())
}

pub(crate) fn contains_key(&self, id: &SO::Id) -> bool {
self.objects.lock().unwrap().contains_key(id)
}
}

#[cfg(test)]
Expand Down
4 changes: 4 additions & 0 deletions src/io/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,3 +78,7 @@ pub(crate) const BDK_WALLET_INDEXER_KEY: &str = "indexer";
///
/// [`StaticInvoice`]: lightning::offers::static_invoice::StaticInvoice
pub(crate) const STATIC_INVOICE_STORE_PRIMARY_NAMESPACE: &str = "static_invoices";

/// The pending payment information will be persisted under this prefix.
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE: &str = "pending_payments";
pub(crate) const PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE: &str = "";
78 changes: 78 additions & 0 deletions src/io/utils.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -46,6 +46,7 @@ use crate::io::{
NODE_METRICS_KEY, NODE_METRICS_PRIMARY_NAMESPACE, NODE_METRICS_SECONDARY_NAMESPACE,
};
use crate::logger::{log_error, LdkLogger, Logger};
use crate::payment::PendingPaymentDetails;
use crate::peer_store::PeerStore;
use crate::types::{Broadcaster, DynStore, KeysManager, Sweeper};
use crate::wallet::ser::{ChangeSetDeserWrapper, ChangeSetSerWrapper};
Expand DownExpand Up@@ -626,6 +627,83 @@ pub(crate) fn read_bdk_wallet_change_set(
Ok(Some(change_set))
}

/// Read previously persisted pending payments information from the store.
pub(crate) async fn read_pending_payments<L: Deref>(
kv_store: &DynStore, logger: L,
) -> Result<Vec<PendingPaymentDetails>, std::io::Error>
where
L::Target: LdkLogger,
{
let mut res = Vec::new();

let mut stored_keys = KVStore::list(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
)
.await?;

const BATCH_SIZE: usize = 50;

let mut set = tokio::task::JoinSet::new();

// Fill JoinSet with tasks if possible
while set.len() < BATCH_SIZE && !stored_keys.is_empty() {
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}
}

while let Some(read_res) = set.join_next().await {
// Exit early if we get an IO error.
let reader = read_res
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?
.map_err(|e| {
log_error!(logger, "Failed to read PendingPaymentDetails: {}", e);
set.abort_all();
e
})?;

// Refill set for every finished future, if we still have something to do.
if let Some(next_key) = stored_keys.pop() {
let fut = KVStore::read(
&*kv_store,
PENDING_PAYMENT_INFO_PERSISTENCE_PRIMARY_NAMESPACE,
PENDING_PAYMENT_INFO_PERSISTENCE_SECONDARY_NAMESPACE,
&next_key,
);
set.spawn(fut);
debug_assert!(set.len() <= BATCH_SIZE);
}

// Handle result.
let pending_payment = PendingPaymentDetails::read(&mut &*reader).map_err(|e| {
log_error!(logger, "Failed to deserialize PendingPaymentDetails: {}", e);
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Failed to deserialize PendingPaymentDetails",
)
})?;
res.push(pending_payment);
}

debug_assert!(set.is_empty());
debug_assert!(stored_keys.is_empty());

Ok(res)
}

#[cfg(test)]
mod tests {
use super::read_or_generate_seed_file;
Expand Down
2 changes: 2 additions & 0 deletions src/payment/mod.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,13 +11,15 @@ pub(crate) mod asynchronous;
mod bolt11;
mod bolt12;
mod onchain;
pub(crate) mod pending_payment_store;
mod spontaneous;
pub(crate) mod store;
mod unified;

pub use bolt11::Bolt11Payment;
pub use bolt12::Bolt12Payment;
pub use onchain::OnchainPayment;
pub use pending_payment_store::PendingPaymentDetails;
pub use spontaneous::SpontaneousPayment;
pub use store::{
ConfirmationStatus, LSPFeeLimits, PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus,
Expand Down
93 changes: 93 additions & 0 deletions src/payment/pending_payment_store.rs
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,93 @@
// This file is Copyright its original authors, visible in version control history.
//
// This file is licensed under the Apache License, Version 2.0 <LICENSE-APACHE or
// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license <LICENSE-MIT or
// 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, ln::channelmanager::PaymentId};

use crate::{
data_store::{StorableObject, StorableObjectUpdate},
payment::{store::PaymentDetailsUpdate, PaymentDetails},
};

/// Represents a pending payment
#[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<Txid>,
}

impl PendingPaymentDetails {
pub(crate) fn new(details: PaymentDetails, conflicting_txids: Vec<Txid>) -> Self {
Self { details, conflicting_txids }
}

/// Convert to finalized payment for the main payment store
pub fn into_payment_details(self) -> PaymentDetails {
self.details
}
}

impl_writeable_tlv_based!(PendingPaymentDetails, {
(0, details, required),
(2, conflicting_txids, optional_vec),
});

#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct PendingPaymentDetailsUpdate {
pub id: PaymentId,
pub payment_update: Option<PaymentDetailsUpdate>,
pub conflicting_txids: Option<Vec<Txid>>,
}

impl StorableObject for PendingPaymentDetails {
type Id = PaymentId;
type Update = PendingPaymentDetailsUpdate;

fn id(&self) -> Self::Id {
self.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.clone();
updated = true;
}
}

updated
}

fn to_update(&self) -> Self::Update {
self.into()
}
}

impl StorableObjectUpdate<PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn id(&self) -> <PendingPaymentDetails as StorableObject>::Id {
self.id
}
}

impl From<&PendingPaymentDetails> for PendingPaymentDetailsUpdate {
fn from(value: &PendingPaymentDetails) -> Self {
Self {
id: value.id(),
payment_update: Some(value.details.to_update()),
conflicting_txids: Some(value.conflicting_txids.clone()),
}
}
}
4 changes: 3 additions & 1 deletion src/types.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -39,7 +39,7 @@ use crate::data_store::DataStore;
use crate::fee_estimator::OnchainFeeEstimator;
use crate::logger::Logger;
use crate::message_handler::NodeCustomMessageHandler;
use crate::payment::PaymentDetails;
use crate::payment::{PaymentDetails, PendingPaymentDetails};
use crate::runtime::RuntimeSpawner;

/// A supertrait that requires that a type implements both [`KVStore`] and [`KVStoreSync`] at the
Expand DownExpand Up@@ -621,3 +621,5 @@ impl From<&(u64, Vec<u8>)> for CustomTlvRecord {
CustomTlvRecord { type_num: tlv.0, value: tlv.1.clone() }
}
}

pub(crate) type PendingPaymentStore = DataStore<PendingPaymentDetails, Arc<Logger>>;
Loading
Loading