Support client_trusts_lsp on LSPS2
What changed, and why it matters
This commit adds a new LSPS2 trust model called client_trusts_lsp. In this model, a Lightning Service Provider (LSP) delays broadcasting the Bitcoin funding transaction for a new channel until the client has actually paid the channel-opening fee via a forwarded payment. The commit also adds a new ChannelManager method that lets LDK validate a funding transaction without immediately broadcasting it. The changes are primarily a feature addition with safety checks to avoid broadcasting funding transactions after a channel has closed or force-closed.
Review as a normal feature/security enhancement. No immediate patch required. Operators using LSPS2 should ensure they call store_funding_transaction, set_funding_tx_broadcast_safe, and payment_forwarded correctly, and understand that client_trusts_lsp changes the trust assumptions and capital lock-up timing for the LSP.
Security signals we found
New manual-broadcast funding path in ChannelManager with validation but no automatic broadcast
LSPS2 service defers funding broadcast until opening fee is collected
Guard that skips broadcast if ChannelManager channel is no longer ready
New state machine fields persisted via TLV serialization
Integration tests assert no broadcast after force-close or partial fee
Evidence from the diff
The patch implements the client_trusts_lsp flow for LSPS2 JIT channels. It introduces a TrustModel enum, stores the funding transaction and a broadcast-safe flag in OutboundJITChannel, and only broadcasts via BroadcasterInterface once (a) the client_trusts_lsp flag is set, (b) the opening fee has been fully skimmed/paid, (c) the funding tx is marked broadcast-safe, and (d) ChannelManager still reports the channel as ready. It also adds ChannelManager::funding_transaction_generated_manual_broadcast and FundingType::CheckedManualBroadcast so LDK validates the funding tx but does not auto-broadcast. Extensive integration tests cover end-to-end success, partial-fee non-broadcast, late broadcast-safe after force-close, and HTLC timeout paths.
Changed components
lightning-liquidity/src/lsps2/service.rslightning-liquidity/src/lsps2/msgs.rslightning-liquidity/src/manager.rslightning/src/ln/channelmanager.rslightning-background-processor/src/lib.rsfuzz/src/lsps_message.rslightning-liquidity/tests/lsps2_integration_tests.rslightning-liquidity/tests/common/mod.rslightning-liquidity/tests/lsps5_integration_tests.rsInspect captured patch +1971 / −92
diff --git a/fuzz/src/lsps_message.rs b/fuzz/src/lsps_message.rs
index 2bc83c3..420fed7 100644
--- a/fuzz/src/lsps_message.rs
+++ b/fuzz/src/lsps_message.rs
@@ -84,6 +84,7 @@ pub fn do_test(data: &[u8]) {
None::<Arc<dyn Filter + Send + Sync>>,
None,
kv_store,
+ Arc::clone(&tx_broadcaster),
None,
None,
).unwrap());
diff --git a/lightning-background-processor/src/lib.rs b/lightning-background-processor/src/lib.rs
index be99f03..288e228 100644
--- a/lightning-background-processor/src/lib.rs
+++ b/lightning-background-processor/src/lib.rs
@@ -426,6 +426,8 @@ pub const NO_LIQUIDITY_MANAGER: Option<
K = &dyn lightning::util::persist::KVStore,
TimeProvider = dyn lightning_liquidity::utils::time::TimeProvider,
TP = &dyn lightning_liquidity::utils::time::TimeProvider,
+ BroadcasterInterface = dyn lightning::chain::chaininterface::BroadcasterInterface,
+ T = &dyn BroadcasterInterface,
> + Send
+ Sync,
>,
@@ -449,6 +451,8 @@ pub const NO_LIQUIDITY_MANAGER_SYNC: Option<
KS = &dyn lightning::util::persist::KVStoreSync,
TimeProvider = dyn lightning_liquidity::utils::time::TimeProvider,
TP = &dyn lightning_liquidity::utils::time::TimeProvider,
+ BroadcasterInterface = dyn lightning::chain::chaininterface::BroadcasterInterface,
+ T = &dyn BroadcasterInterface,
> + Send
+ Sync,
>,
@@ -766,7 +770,7 @@ use futures_util::{dummy_waker, Joiner, OptionalSelector, Selector, SelectorOutp
/// # type P2PGossipSync<UL> = lightning::routing::gossip::P2PGossipSync<Arc<NetworkGraph>, Arc<UL>, Arc<Logger>>;
/// # type ChannelManager<B, F, FE> = lightning::ln::channelmanager::SimpleArcChannelManager<ChainMonitor<B, F, FE>, B, FE, Logger>;
/// # type OnionMessenger<B, F, FE> = lightning::onion_message::messenger::OnionMessenger<Arc<lightning::sign::KeysManager>, Arc<lightning::sign::KeysManager>, Arc<Logger>, Arc<ChannelManager<B, F, FE>>, Arc<lightning::onion_message::messenger::DefaultMessageRouter<Arc<NetworkGraph>, Arc<Logger>, Arc<lightning::sign::KeysManager>>>, Arc<ChannelManager<B, F, FE>>, lightning::ln::peer_handler::IgnoringMessageHandler, lightning::ln::peer_handler::IgnoringMessageHandler, lightning::ln::peer_handler::IgnoringMessageHandler>;
-/// # type LiquidityManager<B, F, FE> = lightning_liquidity::LiquidityManager<Arc<lightning::sign::KeysManager>, Arc<lightning::sign::KeysManager>, Arc<ChannelManager<B, F, FE>>, Arc<F>, Arc<Store>, Arc<DefaultTimeProvider>>;
+/// # type LiquidityManager<B, F, FE> = lightning_liquidity::LiquidityManager<Arc<lightning::sign::KeysManager>, Arc<lightning::sign::KeysManager>, Arc<ChannelManager<B, F, FE>>, Arc<F>, Arc<Store>, Arc<DefaultTimeProvider>, Arc<B>>;
/// # type Scorer = RwLock<lightning::routing::scoring::ProbabilisticScorer<Arc<NetworkGraph>, Arc<Logger>>>;
/// # type PeerManager<B, F, FE, UL> = lightning::ln::peer_handler::SimpleArcPeerManager<SocketDescriptor, ChainMonitor<B, F, FE>, B, FE, Arc<UL>, Logger, F, StoreSync>;
/// # type OutputSweeper<B, D, FE, F, O> = lightning::util::sweep::OutputSweeper<Arc<B>, Arc<D>, Arc<FE>, Arc<F>, Arc<Store>, Arc<Logger>, Arc<O>>;
@@ -1950,6 +1954,7 @@ mod tests {
Arc<dyn Filter + Sync + Send>,
Arc<Persister>,
DefaultTimeProvider,
+ Arc<test_utils::TestBroadcaster>,
>;
struct Node {
@@ -2404,6 +2409,7 @@ mod tests {
None,
None,
Arc::clone(&kv_store),
+ Arc::clone(&tx_broadcaster),
None,
None,
)
@@ -2778,10 +2784,10 @@ mod tests {
let kv_store = KVStoreSyncWrapper(kv_store_sync);
// Yes, you can unsafe { turn off the borrow checker }
- let lm_async: &'static LiquidityManager<_, _, _, _, _, _> = unsafe {
+ let lm_async: &'static LiquidityManager<_, _, _, _, _, _, _> = unsafe {
&*(nodes[0].liquidity_manager.get_lm_async()
- as *const LiquidityManager<_, _, _, _, _, _>)
- as &'static LiquidityManager<_, _, _, _, _, _>
+ as *const LiquidityManager<_, _, _, _, _, _, _>)
+ as &'static LiquidityManager<_, _, _, _, _, _, _>
};
let sweeper_async: &'static OutputSweeper<_, _, _, _, _, _, _> = unsafe {
&*(nodes[0].sweeper.sweeper_async() as *const OutputSweeper<_, _, _, _, _, _, _>)
@@ -3297,10 +3303,10 @@ mod tests {
let kv_store = KVStoreSyncWrapper(kv_store_sync);
// Yes, you can unsafe { turn off the borrow checker }
- let lm_async: &'static LiquidityManager<_, _, _, _, _, _> = unsafe {
+ let lm_async: &'static LiquidityManager<_, _, _, _, _, _, _> = unsafe {
&*(nodes[0].liquidity_manager.get_lm_async()
- as *const LiquidityManager<_, _, _, _, _, _>)
- as &'static LiquidityManager<_, _, _, _, _, _>
+ as *const LiquidityManager<_, _, _, _, _, _, _>)
+ as &'static LiquidityManager<_, _, _, _, _, _, _>
};
let sweeper_async: &'static OutputSweeper<_, _, _, _, _, _, _> = unsafe {
&*(nodes[0].sweeper.sweeper_async() as *const OutputSweeper<_, _, _, _, _, _, _>)
@@ -3524,10 +3530,10 @@ mod tests {
let (exit_sender, exit_receiver) = tokio::sync::watch::channel(());
// Yes, you can unsafe { turn off the borrow checker }
- let lm_async: &'static LiquidityManager<_, _, _, _, _, _> = unsafe {
+ let lm_async: &'static LiquidityManager<_, _, _, _, _, _, _> = unsafe {
&*(nodes[0].liquidity_manager.get_lm_async()
- as *const LiquidityManager<_, _, _, _, _, _>)
- as &'static LiquidityManager<_, _, _, _, _, _>
+ as *const LiquidityManager<_, _, _, _, _, _, _>)
+ as &'static LiquidityManager<_, _, _, _, _, _, _>
};
let sweeper_async: &'static OutputSweeper<_, _, _, _, _, _, _> = unsafe {
&*(nodes[0].sweeper.sweeper_async() as *const OutputSweeper<_, _, _, _, _, _, _>)
diff --git a/lightning-liquidity/src/lsps2/msgs.rs b/lightning-liquidity/src/lsps2/msgs.rs
index 21e1af8..ba4d0fe 100644
--- a/lightning-liquidity/src/lsps2/msgs.rs
+++ b/lightning-liquidity/src/lsps2/msgs.rs
@@ -183,7 +183,14 @@ pub struct LSPS2BuyResponse {
pub jit_channel_scid: LSPS2InterceptScid,
/// The locktime expiry delta the lsp requires.
pub lsp_cltv_expiry_delta: u32,
- /// A flag that indicates who is trusting who.
+ /// Trust model flag (default: false).
+ ///
+ /// false => "LSP trusts client": LSP immediately (or as soon as safe) broadcasts the
+ /// funding transaction; client may wait for broadcast / confirmations
+ /// before revealing the preimage.
+ /// true => "Client trusts LSP": LSP may defer broadcasting until after the client
+ /// reveals the preimage; client MUST send the preimage once HTLC(s) are
+ /// irrevocably committed.
#[serde(default)]
pub client_trusts_lsp: bool,
}
diff --git a/lightning-liquidity/src/lsps2/service.rs b/lightning-liquidity/src/lsps2/service.rs
index 2d727ac..4837b73 100644
--- a/lightning-liquidity/src/lsps2/service.rs
+++ b/lightning-liquidity/src/lsps2/service.rs
@@ -40,6 +40,7 @@ use crate::prelude::{new_hash_map, HashMap};
use crate::sync::{Arc, Mutex, MutexGuard, RwLock};
use crate::utils::async_poll::dummy_waker;
+use lightning::chain::chaininterface::BroadcasterInterface;
use lightning::events::HTLCHandlingFailureType;
use lightning::ln::channelmanager::{AChannelManager, FailureCode, InterceptId};
use lightning::ln::msgs::{ErrorAction, LightningError};
@@ -52,6 +53,7 @@ use lightning::{impl_writeable_tlv_based, impl_writeable_tlv_based_enum};
use lightning_types::payment::PaymentHash;
use bitcoin::secp256k1::PublicKey;
+use bitcoin::Transaction;
use crate::lsps2::msgs::{
LSPS2BuyRequest, LSPS2BuyResponse, LSPS2GetInfoRequest, LSPS2GetInfoResponse, LSPS2Message,
@@ -117,6 +119,77 @@ struct ForwardPaymentAction(ChannelId, FeePayment);
#[derive(Debug, PartialEq)]
struct ForwardHTLCsAction(ChannelId, Vec<InterceptedHTLC>);
+#[derive(Debug, Clone)]
+enum TrustModel {
+ ClientTrustsLsp { funding_tx_broadcast_safe: bool, funding_tx: Option<Transaction> },
+ LspTrustsClient,
+}
+
+impl TrustModel {
+ fn should_manually_broadcast(&self, state_is_payment_forward: bool) -> bool {
+ match self {
+ TrustModel::ClientTrustsLsp { funding_tx_broadcast_safe, funding_tx } => {
+ *funding_tx_broadcast_safe && state_is_payment_forward && funding_tx.is_some()
+ },
+ // in lsp-trusts-client, the broadcast is automatic, so we never need to manually broadcast.
+ TrustModel::LspTrustsClient => false,
+ }
+ }
+
+ fn new(client_trusts_lsp: bool) -> Self {
+ if client_trusts_lsp {
+ TrustModel::ClientTrustsLsp { funding_tx_broadcast_safe: false, funding_tx: None }
+ } else {
+ TrustModel::LspTrustsClient
+ }
+ }
+
+ fn set_funding_tx(&mut self, funding_tx: Transaction) {
+ match self {
+ TrustModel::ClientTrustsLsp { funding_tx: tx, .. } => {
+ *tx = Some(funding_tx);
+ },
+ TrustModel::LspTrustsClient => {
+ // No-op
+ },
+ }
+ }
+
+ fn set_funding_tx_broadcast_safe(&mut self, funding_tx_broadcast_safe: bool) {
+ match self {
+ TrustModel::ClientTrustsLsp { funding_tx_broadcast_safe: safe, .. } => {
+ *safe = funding_tx_broadcast_safe;
+ },
+ TrustModel::LspTrustsClient => {
+ // No-op
+ },
+ }
+ }
+
+ fn get_funding_tx(&self) -> Option<Transaction> {
+ match self {
+ TrustModel::ClientTrustsLsp { funding_tx, .. } => funding_tx.clone(),
+ _ => None,
+ }
+ }
+
+ fn is_client_trusts_lsp(&self) -> bool {
+ match self {
+ TrustModel::ClientTrustsLsp { .. } => true,
+ TrustModel::LspTrustsClient => false,
+ }
+ }
+}
+
+impl_writeable_tlv_based_enum!(TrustModel,
+ (0, ClientTrustsLsp) => {
+ (0, funding_tx_broadcast_safe, required),
+ (2, funding_tx, option),
+ },
+ (2, LspTrustsClient) => {
+ },
+);
+
/// The different states a requested JIT channel can be in.
#[derive(Debug)]
enum OutboundJITChannelState {
@@ -360,15 +433,28 @@ impl OutboundJITChannelState {
}
}
- fn payment_forwarded(&mut self) -> Result<Option<ForwardHTLCsAction>, ChannelStateError> {
+ fn payment_forwarded(
+ &mut self, skimmed_fee_msat: u64,
+ ) -> Result<Option<ForwardHTLCsAction>, ChannelStateError> {
match self {
OutboundJITChannelState::PendingPaymentForward {
- payment_queue, channel_id, ..
+ payment_queue,
+ channel_id,
+ opening_fee_msat,
} => {
- let mut payment_queue = core::mem::take(payment_queue);
- let forward_htlcs = ForwardHTLCsAction(*channel_id, payment_queue.clear());
- *self = OutboundJITChannelState::PaymentForwarded { channel_id: *channel_id };
- Ok(Some(forward_htlcs))
+ if skimmed_fee_msat >= *opening_fee_msat {
+ let mut pq = core::mem::take(payment_queue);
+ let forward_htlcs = ForwardHTLCsAction(*channel_id, pq.clear());
+ *self = OutboundJITChannelState::PaymentForwarded { channel_id: *channel_id };
+ Ok(Some(forward_htlcs))
+ } else {
+ *self = OutboundJITChannelState::PendingPaymentForward {
+ payment_queue: core::mem::take(payment_queue),
+ opening_fee_msat: *opening_fee_msat,
+ channel_id: *channel_id,
+ };
+ Ok(None)
+ }
},
OutboundJITChannelState::PaymentForwarded { channel_id } => {
*self = OutboundJITChannelState::PaymentForwarded { channel_id: *channel_id };
@@ -410,6 +496,7 @@ struct OutboundJITChannel {
user_channel_id: u128,
opening_fee_params: LSPS2OpeningFeeParams,
payment_size_msat: Option<u64>,
+ trust_model: TrustModel,
}
impl_writeable_tlv_based!(OutboundJITChannel, {
@@ -417,18 +504,20 @@ impl_writeable_tlv_based!(OutboundJITChannel, {
(2, user_channel_id, required),
(4, opening_fee_params, required),
(6, payment_size_msat, option),
+ (8, trust_model, required),
});
impl OutboundJITChannel {
fn new(
payment_size_msat: Option<u64>, opening_fee_params: LSPS2OpeningFeeParams,
- user_channel_id: u128,
+ user_channel_id: u128, client_trusts_lsp: bool,
) -> Self {
Self {
user_channel_id,
state: OutboundJITChannelState::new(),
opening_fee_params,
payment_size_msat,
+ trust_model: TrustModel::new(client_trusts_lsp),
}
}
@@ -452,8 +541,10 @@ impl OutboundJITChannel {
Ok(action)
}
- fn payment_forwarded(&mut self) -> Result<Option<ForwardHTLCsAction>, LightningError> {
- let action = self.state.payment_forwarded()?;
+ fn payment_forwarded(
+ &mut self, skimmed_fee_msat: u64,
+ ) -> Result<Option<ForwardHTLCsAction>, LightningError> {
+ let action = self.state.payment_forwarded(skimmed_fee_msat)?;
Ok(action)
}
@@ -467,6 +558,36 @@ impl OutboundJITChannel {
let is_expired = is_expired_opening_fee_params(&self.opening_fee_params);
self.is_pending_initial_payment() && is_expired
}
+
+ fn set_funding_tx(&mut self, funding_tx: Transaction) {
+ self.trust_model.set_funding_tx(funding_tx);
+ }
+
+ fn set_funding_tx_broadcast_safe(&mut self, funding_tx_broadcast_safe: bool) {
+ self.trust_model.set_funding_tx_broadcast_safe(funding_tx_broadcast_safe);
+ }
+
+ fn should_broadcast_funding_transaction(&self) -> bool {
+ self.trust_model.should_manually_broadcast(matches!(
+ self.state,
+ OutboundJITChannelState::PaymentForwarded { .. }
+ ))
+ }
+
+ fn get_channel_id(&self) -> Option<ChannelId> {
+ match self.state {
+ OutboundJITChannelState::PaymentForwarded { channel_id } => Some(channel_id),
+ _ => None,
+ }
+ }
+
+ fn get_funding_tx(&self) -> Option<Transaction> {
+ self.trust_model.get_funding_tx()
+ }
+
+ fn is_client_trusts_lsp(&self) -> bool {
+ self.trust_model.is_client_trusts_lsp()
+ }
}
pub(crate) struct PeerState {
@@ -579,13 +700,15 @@ macro_rules! get_or_insert_peer_state_entry {
}
/// The main object allowing to send and receive bLIP-52 / LSPS2 messages.
-pub struct LSPS2ServiceHandler<CM: Deref, K: Deref + Clone>
+pub struct LSPS2ServiceHandler<CM: Deref, K: Deref + Clone, T: Deref>
where
CM::Target: AChannelManager,
K::Target: KVStore,
+ T::Target: BroadcasterInterface,
{
channel_manager: CM,
kv_store: K,
+ tx_broadcaster: T,
pending_messages: Arc<MessageQueue>,
pending_events: Arc<EventQueue<K>>,
per_peer_state: RwLock<HashMap<PublicKey, Mutex<PeerState>>>,
@@ -595,15 +718,16 @@ where
config: LSPS2ServiceConfig,
}
-impl<CM: Deref, K: Deref + Clone> LSPS2ServiceHandler<CM, K>
+impl<CM: Deref, K: Deref + Clone, T: Deref + Clone> LSPS2ServiceHandler<CM, K, T>
where
CM::Target: AChannelManager,
K::Target: KVStore,
+ T::Target: BroadcasterInterface,
{
/// Constructs a `LSPS2ServiceHandler`.
pub(crate) fn new(
per_peer_state: HashMap<PublicKey, Mutex<PeerState>>, pending_messages: Arc<MessageQueue>,
- pending_events: Arc<EventQueue<K>>, channel_manager: CM, kv_store: K,
+ pending_events: Arc<EventQueue<K>>, channel_manager: CM, kv_store: K, tx_broadcaster: T,
config: LSPS2ServiceConfig,
) -> Result<Self, lightning::io::Error> {
let mut peer_by_intercept_scid = new_hash_map();
@@ -642,6 +766,7 @@ where
total_pending_requests: AtomicUsize::new(0),
channel_manager,
kv_store,
+ tx_broadcaster,
config,
})
}
@@ -768,6 +893,12 @@ where
///
/// Should be called in response to receiving a [`LSPS2ServiceEvent::BuyRequest`] event.
///
+ /// `client_trusts_lsp`:
+ /// * false (default) => "LSP trusts client": LSP broadcasts the funding
+ /// transaction as soon as it is safe and forwards the payment normally.
+ /// * true => "Client trusts LSP": LSP may defer broadcasting the funding
+ /// transaction until after the client claims the forwarded HTLC(s).
+ ///
/// [`ChannelManager::create_channel`]: lightning::ln::channelmanager::ChannelManager::create_channel
/// [`ChannelManager::get_intercept_scid`]: lightning::ln::channelmanager::ChannelManager::get_intercept_scid
/// [`LSPS2ServiceEvent::BuyRequest`]: crate::lsps2::event::LSPS2ServiceEvent::BuyRequest
@@ -794,6 +925,7 @@ where
buy_request.payment_size_msat,
buy_request.opening_fee_params,
user_channel_id,
+ client_trusts_lsp,
);
peer_state_lock
@@ -1033,17 +1165,21 @@ where
/// Forward [`Event::PaymentForwarded`] event parameter into this function.
///
/// Will register the forwarded payment as having paid the JIT channel fee, and forward any held
- /// and future HTLCs for the SCID of the initial invoice. In the future, this will verify the
- /// `skimmed_fee_msat` in [`Event::PaymentForwarded`].
+ /// and future HTLCs for the SCID of the initial invoice.
+ ///
+ /// When the reported skimmed fee equals or exceeds the promised opening fee, any HTLCs that
+ /// were being held for that JIT channel are forwarded. In a `client_trusts_lsp` flow, once
+ /// the fee has been fully paid, the channel's funding transaction will be broadcasted.
///
- /// Note that `next_channel_id` is required to be provided. Therefore, the corresponding
- /// [`Event::PaymentForwarded`] events need to be generated and serialized by LDK versions
- /// greater or equal to 0.0.107.
+ /// Note that `next_channel_id` and `skimmed_fee_msat` are required to be provided.
+ /// Therefore, the corresponding [`Event::PaymentForwarded`] events need to be generated and
+ /// serialized by LDK versions greater or equal to 0.0.122.
///
/// [`Event::PaymentForwarded`]: lightning::events::Event::PaymentForwarded
- pub async fn payment_forwarded(&self, next_channel_id: ChannelId) -> Result<(), APIError> {
+ pub async fn payment_forwarded(
+ &self, next_channel_id: ChannelId, skimmed_fee_msat: u64,
+ ) -> Result<(), APIError> {
let mut should_persist = None;
-
if let Some(counterparty_node_id) =
self.peer_by_channel_id.read().unwrap().get(&next_channel_id)
{
@@ -1059,7 +1195,7 @@ where
if let Some(jit_channel) =
peer_state.outbound_channels_by_intercept_scid.get_mut(&intercept_scid)
{
- match jit_channel.payment_forwarded() {
+ match jit_channel.payment_forwarded(skimmed_fee_msat) {
Ok(Some(ForwardHTLCsAction(channel_id, htlcs))) => {
for htlc in htlcs {
self.channel_manager.get_cm().forward_intercepted_htlc(
@@ -1080,6 +1216,8 @@ where
})
},
}
+
+ self.broadcast_funding_transaction_if_applies(jit_channel);
}
} else {
return Err(APIError::APIMisuseError {
@@ -1690,12 +1828,172 @@ where
peer_state_lock.prune_expired_request_state();
}
}
+
+ /// Checks if the JIT channel with the given `user_channel_id` needs manual broadcast.
+ ///
+ /// Will be `true` if `client_trusts_lsp` is set to `true`.
+ pub fn channel_needs_manual_broadcast(
+ &self, user_channel_id: u128, counterparty_node_id: &PublicKey,
+ ) -> Result<bool, APIError> {
+ let outer_state_lock = self.per_peer_state.read().unwrap();
+ let inner_state_lock =
+ outer_state_lock.get(counterparty_node_id).ok_or_else(|| APIError::APIMisuseError {
+ err: format!("No counterparty state for: {}", counterparty_node_id),
+ })?;
+ let peer_state = inner_state_lock.lock().unwrap();
+
+ let intercept_scid = peer_state
+ .intercept_scid_by_user_channel_id
+ .get(&user_channel_id)
+ .copied()
+ .ok_or_else(|| APIError::APIMisuseError {
+ err: format!("Could not find a channel with user_channel_id {}", user_channel_id),
+ })?;
+
+ let jit_channel = peer_state
+ .outbound_channels_by_intercept_scid
+ .get(&intercept_scid)
+ .ok_or_else(|| APIError::APIMisuseError {
+ err: format!(
+ "Failed to map intercept_scid {} for user_channel_id {} to a channel.",
+ intercept_scid, user_channel_id,
+ ),
+ })?;
+
+ Ok(jit_channel.is_client_trusts_lsp())
+ }
+
+ /// Stores the funding transaction for a JIT channel.
+ ///
+ /// Call this when the funding transaction is created.
+ ///
+ /// In `client_trusts_lsp` the broadcasting of the funding transaction will be handled internally
+ /// after you also mark it as broadcast-safe via
+ /// [`set_funding_tx_broadcast_safe`] and once the opening fee has been collected. You do not need
+ /// to broadcast the funding transaction yourself in this flow.
+ ///
+ /// [`set_funding_tx_broadcast_safe`]: Self::set_funding_tx_broadcast_safe
+ pub fn store_funding_transaction(
+ &self, user_channel_id: u128, counterparty_node_id: &PublicKey, funding_tx: Transaction,
+ ) -> Result<(), APIError> {
+ let outer_state_lock = self.per_peer_state.read().unwrap();
+ let inner_state_lock =
+ outer_state_lock.get(counterparty_node_id).ok_or_else(|| APIError::APIMisuseError {
+ err: format!("No counterparty state for: {}", counterparty_node_id),
+ })?;
+ let mut peer_state = inner_state_lock.lock().unwrap();
+
+ let intercept_scid = peer_state
+ .intercept_scid_by_user_channel_id
+ .get(&user_channel_id)
+ .copied()
+ .ok_or_else(|| APIError::APIMisuseError {
+ err: format!("Could not find a channel with user_channel_id {}", user_channel_id),
+ })?;
+
+ let jit_channel = peer_state
+ .outbound_channels_by_intercept_scid
+ .get_mut(&intercept_scid)
+ .ok_or_else(|| APIError::APIMisuseError {
+ err: format!(
+ "Failed to map intercept_scid {} for user_channel_id {} to a channel.",
+ intercept_scid, user_channel_id,
+ ),
+ })?;
+
+ jit_channel.set_funding_tx(funding_tx);
+
+ self.broadcast_funding_transaction_if_applies(jit_channel);
+ Ok(())
+ }
+
+ /// Marks that the funding transaction for the JIT channel identified by `user_channel_id`
+ /// is now safe to broadcast.
+ ///
+ /// In LDK call this when you receive [`Event::FundingTxBroadcastSafe`]. In other Lightning
+ /// backends call it once the funding transaction is fully negotiated and signed (all
+ /// signatures verified), your channel state machine will now proceed assuming the funding
+ /// transaction will confirm, and you are intentionally deferring the actual broadcast so
+ /// the LSPS2 flow (when `client_trusts_lsp = true`) can first collect the opening fee from
+ /// the intercepted payment.
+ ///
+ /// In a `client_trusts_lsp` flow, after this is set and the opening fee has been fully skimmed,
+ /// the channel's funding transaction will be broadcasted if the channel is still usable.
+ /// If the channel has been closed or force-closed before this point, the funding transaction will not be broadcasted.
+ ///
+ /// [`Event::FundingTxBroadcastSafe`]: lightning::events::Event::FundingTxBroadcastSafe
+ pub fn set_funding_tx_broadcast_safe(
+ &self, user_channel_id: u128, counterparty_node_id: &PublicKey,
+ ) -> Result<(), APIError> {
+ let outer_state_lock = self.per_peer_state.read().unwrap();
+ let inner_state_lock =
+ outer_state_lock.get(counterparty_node_id).ok_or_else(|| APIError::APIMisuseError {
+ err: format!("No counterparty state for: {}", counterparty_node_id),
+ })?;
+ let mut peer_state = inner_state_lock.lock().unwrap();
+
+ let intercept_scid = peer_state
+ .intercept_scid_by_user_channel_id
+ .get(&user_channel_id)
+ .copied()
+ .ok_or_else(|| APIError::APIMisuseError {
+ err: format!("Could not find a channel with user_channel_id {}", user_channel_id),
+ })?;
+
+ let jit_channel = peer_state
+ .outbound_channels_by_intercept_scid
+ .get_mut(&intercept_scid)
+ .ok_or_else(|| APIError::APIMisuseError {
+ err: format!(
+ "Failed to map intercept_scid {} for user_channel_id {} to a channel.",
+ intercept_scid, user_channel_id,
+ ),
+ })?;
+
+ jit_channel.set_funding_tx_broadcast_safe(true);
+
+ self.broadcast_funding_transaction_if_applies(jit_channel);
+ Ok(())
+ }
+
+ fn broadcast_funding_transaction_if_applies(&self, jit_channel: &OutboundJITChannel) {
+ if !jit_channel.should_broadcast_funding_transaction() {
+ return;
+ }
+
+ // Broadcast the funding transaction only if the LDK channel is still usable. In
+ // the `client_trusts_lsp` flow we delay funding broadcast until the opening fee is
+ // collected. Before that happens, LDK may force-close the not‑yet‑funded channel
+ // (for example when a forwarded HTLC nears expiry). Broadcasting funding after a
+ // close could then confirm the commitment and trigger unintended on‑chain handling.
+ // To avoid this, we check ChannelManager’s view (`is_channel_ready`) before broadcasting.
+ let channel_id_opt = jit_channel.get_channel_id();
+ if let Some(ch_id) = channel_id_opt {
+ let is_channel_ready = self
+ .channel_manager
+ .get_cm()
+ .list_channels()
+ .into_iter()
+ .any(|cd| cd.channel_id == ch_id && cd.is_channel_ready);
+ if !is_channel_ready {
+ return;
+ }
+ } else {
+ return;
+ }
+
+ if let Some(funding_tx) = jit_channel.get_funding_tx() {
+ self.tx_broadcaster.broadcast_transactions(&[&funding_tx]);
+ }
+ }
}
-impl<CM: Deref, K: Deref + Clone> LSPSProtocolMessageHandler for LSPS2ServiceHandler<CM, K>
+impl<CM: Deref, K: Deref + Clone, T: Deref + Clone> LSPSProtocolMessageHandler
+ for LSPS2ServiceHandler<CM, K, T>
where
CM::Target: AChannelManager,
K::Target: KVStore,
+ T::Target: BroadcasterInterface,
{
type ProtocolMessage = LSPS2Message;
const PROTOCOL_NUMBER: Option<u16> = Some(2);
@@ -1765,20 +2063,22 @@ fn calculate_amount_to_forward_per_htlc(
/// A synchroneous wrapper around [`LSPS2ServiceHandler`] to be used in contexts where async is not
/// available.
-pub struct LSPS2ServiceHandlerSync<'a, CM: Deref, K: Deref + Clone>
+pub struct LSPS2ServiceHandlerSync<'a, CM: Deref, K: Deref + Clone, T: Deref + Clone>
where
CM::Target: AChannelManager,
K::Target: KVStore,
+ T::Target: BroadcasterInterface,
{
- inner: &'a LSPS2ServiceHandler<CM, K>,
+ inner: &'a LSPS2ServiceHandler<CM, K, T>,
}
-impl<'a, CM: Deref, K: Deref + Clone> LSPS2ServiceHandlerSync<'a, CM, K>
+impl<'a, CM: Deref, K: Deref + Clone, T: Deref + Clone> LSPS2ServiceHandlerSync<'a, CM, K, T>
where
CM::Target: AChannelManager,
K::Target: KVStore,
+ T::Target: BroadcasterInterface,
{
- pub(crate) fn from_inner(inner: &'a LSPS2ServiceHandler<CM, K>) -> Self {
+ pub(crate) fn from_inner(inner: &'a LSPS2ServiceHandler<CM, K, T>) -> Self {
Self { inner }
}
@@ -1893,8 +2193,10 @@ where
/// Wraps [`LSPS2ServiceHandler::payment_forwarded`].
///
/// [`Event::PaymentForwarded`]: lightning::events::Event::PaymentForwarded
- pub fn payment_forwarded(&self, next_channel_id: ChannelId) -> Result<(), APIError> {
- let mut fut = Box::pin(self.inner.payment_forwarded(next_channel_id));
+ pub fn payment_forwarded(
+ &self, next_channel_id: ChannelId, skimmed_fee_msat: u64,
+ ) -> Result<(), APIError> {
+ let mut fut = Box::pin(self.inner.payment_forwarded(next_channel_id, skimmed_fee_msat));
let mut waker = dummy_waker();
let mut ctx = task::Context::from_waker(&mut waker);
@@ -1907,6 +2209,27 @@ where
}
}
+ /// Wraps [`LSPS2ServiceHandler::channel_needs_manual_broadcast`].
+ pub fn channel_needs_manual_broadcast(
+ &self, user_channel_id: u128, counterparty_node_id: &PublicKey,
+ ) -> Result<bool, APIError> {
+ self.inner.channel_needs_manual_broadcast(user_channel_id, counterparty_node_id)
+ }
+
+ /// Wraps [`LSPS2ServiceHandler::store_funding_transaction`].
+ pub fn store_funding_transaction(
+ &self, user_channel_id: u128, counterparty_node_id: &PublicKey, funding_tx: Transaction,
+ ) -> Result<(), APIError> {
+ self.inner.store_funding_transaction(user_channel_id, counterparty_node_id, funding_tx)
+ }
+
+ /// Wraps [`LSPS2ServiceHandler::set_funding_tx_broadcast_safe`].
+ pub fn set_funding_tx_broadcast_safe(
+ &self, user_channel_id: u128, counterparty_node_id: &PublicKey,
+ ) -> Result<(), APIError> {
+ self.inner.set_funding_tx_broadcast_safe(user_channel_id, counterparty_node_id)
+ }
+
/// Abandons a pending JIT‐open flow for `user_channel_id`, removing all local state.
///
/// Wraps [`LSPS2ServiceHandler::channel_open_abandoned`].
@@ -1978,6 +2301,7 @@ mod tests {
use proptest::prelude::*;
+ use bitcoin::{absolute::LockTime, transaction::Version};
use core::str::FromStr;
const MAX_VALUE_MSAT: u64 = 21_000_000_0000_0000_000;
@@ -2214,7 +2538,7 @@ mod tests {
}
// Payment completes, queued payments get forwarded.
{
- let action = state.payment_forwarded().unwrap();
+ let action = state.payment_forwarded(100000000000).unwrap();
assert!(matches!(state, OutboundJITChannelState::PaymentForwarded { .. }));
match action {
Some(ForwardHTLCsAction(channel_id, htlcs)) => {
@@ -2349,7 +2673,7 @@ mod tests {
}
// Payment completes, queued payments get forwarded.
{
- let action = state.payment_forwarded().unwrap();
+ let action = state.payment_forwarded(10000000000).unwrap();
assert!(matches!(state, OutboundJITChannelState::PaymentForwarded { .. }));
match action {
Some(ForwardHTLCsAction(channel_id, htlcs)) => {
@@ -2385,4 +2709,105 @@ mod tests {
);
}
}
+
+ #[test]
+ fn broadcast_not_allowed_after_non_paying_fee_payment_claimed() {
+ let min_fee_msat: u64 = 12345;
+ let opening_fee_params = LSPS2OpeningFeeParams {
+ min_fee_msat,
+ proportional: 0,
+ valid_until: LSPSDateTime::from_str("2035-05-20T08:30:45Z").unwrap(),
+ min_lifetime: 144,
+ max_client_to_self_delay: 128,
+ min_payment_size_msat: 1,
+ max_payment_size_msat: 10_000_000_000,
+ promise: "ignore".to_string(),
+ };
+
+ let payment_size_msat = Some(1_000_000);
+ let user_channel_id = 4242u128;
+ let mut jit_channel = OutboundJITChannel::new(
+ payment_size_msat,
+ opening_fee_params.clone(),
+ user_channel_id,
+ true,
+ );
+
+ let opening_payment_hash = PaymentHash([42; 32]);
+ let htlcs_for_opening = [
+ InterceptedHTLC {
+ intercept_id: InterceptId([0; 32]),
+ expected_outbound_amount_msat: 400_000,
+ payment_hash: opening_payment_hash,
+ },
+ InterceptedHTLC {
+ intercept_id: InterceptId([1; 32]),
+ expected_outbound_amount_msat: 600_000,
+ payment_hash: opening_payment_hash,
+ },
+ ];
+
+ assert!(jit_channel.htlc_intercepted(htlcs_for_opening[0].clone()).unwrap().is_none());
+ let action = jit_channel.htlc_intercepted(htlcs_for_opening[1].clone()).unwrap();
+ match action {
+ Some(HTLCInterceptedAction::OpenChannel(_)) => {},
+ other => panic!("Expected OpenChannel action, got {:?}", other),
+ }
+
+ let channel_id = ChannelId([7; 32]);
+ let ForwardPaymentAction(_, fee_payment) = jit_channel.channel_ready(channel_id).unwrap();
+ assert_eq!(fee_payment.opening_fee_msat, min_fee_msat);
+
+ let followup = jit_channel.htlc_handling_failed().unwrap();
+ assert!(followup.is_none());
+
+ let dummy_tx = Transaction {
+ version: Version(2),
+ lock_time: LockTime::ZERO,
+ input: vec![],
+ output: vec![],
+ };
+ jit_channel.set_funding_tx(dummy_tx);
+ jit_channel.set_funding_tx_broadcast_safe(true);
+ assert!(
+ !jit_channel.should_broadcast_funding_transaction(),
+ "Should not broadcast before any successful payment is claimed"
+ );
+
+ let second_payment_hash = PaymentHash([99; 32]);
+ let second_htlc = InterceptedHTLC {
+ intercept_id: InterceptId([2; 32]),
+ expected_outbound_amount_msat: min_fee_msat,
+ payment_hash: second_payment_hash,
+ };
+ let action2 = jit_channel.htlc_intercepted(second_htlc).unwrap();
+ let (forwarded_channel_id, fee_payment2) = match action2 {
+ Some(HTLCInterceptedAction::ForwardPayment(cid, fp)) => (cid, fp),
+ other => panic!("Expected ForwardPayment for second HTLC, got {:?}", other),
+ };
+ assert_eq!(forwarded_channel_id, channel_id);
+ assert_eq!(fee_payment2.opening_fee_msat, min_fee_msat);
+
+ assert!(
+ !jit_channel.should_broadcast_funding_transaction(),
+ "Should not broadcast before any successful payment is claimed"
+ );
+
+ // Forward a payment that is not enough to cover the fees
+ let _ = jit_channel.payment_forwarded(min_fee_msat - 1).unwrap();
+
+ assert!(
+ !jit_channel.should_broadcast_funding_transaction(),
+ "Should not broadcast before all the fees are collected"
+ );
+
+ let _ = jit_channel.payment_forwarded(min_fee_msat).unwrap();
+
+ let broadcast_allowed = jit_channel.should_broadcast_funding_transaction();
+
+ assert!(
+ broadcast_allowed,
+ "Broadcast was not allowed even though all the skimmed fees were collected"
+ );
+ }
}
diff --git a/lightning-liquidity/src/manager.rs b/lightning-liquidity/src/manager.rs
index 9b45233..5d95d32 100644
--- a/lightning-liquidity/src/manager.rs
+++ b/lightning-liquidity/src/manager.rs
@@ -43,6 +43,7 @@ use crate::utils::async_poll::dummy_waker;
use crate::utils::time::DefaultTimeProvider;
use crate::utils::time::TimeProvider;
+use lightning::chain::chaininterface::BroadcasterInterface;
use lightning::chain::{self, BestBlock, Confirm, Filter, Listen};
use lightning::ln::channelmanager::{AChannelManager, ChainParameters};
use lightning::ln::msgs::{ErrorAction, LightningError};
@@ -68,6 +69,7 @@ const LSPS_FEATURE_BIT: usize = 729;
///
/// Allows end-users to configure options when using the [`LiquidityManager`]
/// to provide liquidity services to clients.
+#[derive(Clone)]
pub struct LiquidityServiceConfig {
/// Optional server-side configuration for LSPS1 channel requests.
#[cfg(lsps1_service)]
@@ -86,6 +88,7 @@ pub struct LiquidityServiceConfig {
///
/// Allows end-user to configure options when using the [`LiquidityManager`]
/// to access liquidity services from a provider.
+#[derive(Clone)]
pub struct LiquidityClientConfig {
/// Optional client-side configuration for LSPS1 channel requests.
pub lsps1_client_config: Option<LSPS1ClientConfig>,
@@ -124,9 +127,14 @@ pub trait ALiquidityManager {
type TimeProvider: TimeProvider + ?Sized;
/// A type that may be dereferenced to [`Self::TimeProvider`].
type TP: Deref<Target = Self::TimeProvider> + Clone;
+ /// A type implementing [`BroadcasterInterface`].
+ type BroadcasterInterface: BroadcasterInterface + ?Sized;
+ /// A type that may be dereferenced to [`Self::BroadcasterInterface`].
+ type T: Deref<Target = Self::BroadcasterInterface> + Clone;
/// Returns a reference to the actual [`LiquidityManager`] object.
- fn get_lm(&self)
- -> &LiquidityManager<Self::ES, Self::NS, Self::CM, Self::C, Self::K, Self::TP>;
+ fn get_lm(
+ &self,
+ ) -> &LiquidityManager<Self::ES, Self::NS, Self::CM, Self::C, Self::K, Self::TP, Self::T>;
}
impl<
@@ -136,7 +144,8 @@ impl<
C: Deref + Clone,
K: Deref + Clone,
TP: Deref + Clone,
- > ALiquidityManager for LiquidityManager<ES, NS, CM, C, K, TP>
+ T: Deref + Clone,
+ > ALiquidityManager for LiquidityManager<ES, NS, CM, C, K, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -144,6 +153,7 @@ where
C::Target: Filter,
K::Target: KVStore,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
type EntropySource = ES::Target;
type ES = ES;
@@ -157,7 +167,9 @@ where
type K = K;
type TimeProvider = TP::Target;
type TP = TP;
- fn get_lm(&self) -> &LiquidityManager<ES, NS, CM, C, K, TP> {
+ type BroadcasterInterface = T::Target;
+ type T = T;
+ fn get_lm(&self) -> &LiquidityManager<ES, NS, CM, C, K, TP, T> {
self
}
}
@@ -191,6 +203,10 @@ pub trait ALiquidityManagerSync {
type TimeProvider: TimeProvider + ?Sized;
/// A type that may be dereferenced to [`Self::TimeProvider`].
type TP: Deref<Target = Self::TimeProvider> + Clone;
+ /// A type implementing [`BroadcasterInterface`].
+ type BroadcasterInterface: BroadcasterInterface + ?Sized;
+ /// A type that may be dereferenced to [`Self::BroadcasterInterface`].
+ type T: Deref<Target = Self::BroadcasterInterface> + Clone;
/// Returns the inner async [`LiquidityManager`] for testing purposes.
#[cfg(any(test, feature = "_test_utils"))]
fn get_lm_async(
@@ -202,11 +218,12 @@ pub trait ALiquidityManagerSync {
Self::C,
KVStoreSyncWrapper<Self::KS>,
Self::TP,
+ Self::T,
>;
/// Returns a reference to the actual [`LiquidityManager`] object.
fn get_lm(
&self,
- ) -> &LiquidityManagerSync<Self::ES, Self::NS, Self::CM, Self::C, Self::KS, Self::TP>;
+ ) -> &LiquidityManagerSync<Self::ES, Self::NS, Self::CM, Self::C, Self::KS, Self::TP, Self::T>;
}
impl<
@@ -216,7 +233,8 @@ impl<
C: Deref + Clone,
KS: Deref + Clone,
TP: Deref + Clone,
- > ALiquidityManagerSync for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+ T: Deref + Clone,
+ > ALiquidityManagerSync for LiquidityManagerSync<ES, NS, CM, C, KS, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -224,6 +242,7 @@ where
C::Target: Filter,
KS::Target: KVStoreSync,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
type EntropySource = ES::Target;
type ES = ES;
@@ -237,6 +256,8 @@ where
type KS = KS;
type TimeProvider = TP::Target;
type TP = TP;
+ type BroadcasterInterface = T::Target;
+ type T = T;
/// Returns the inner async [`LiquidityManager`] for testing purposes.
#[cfg(any(test, feature = "_test_utils"))]
fn get_lm_async(
@@ -248,10 +269,11 @@ where
Self::C,
KVStoreSyncWrapper<Self::KS>,
Self::TP,
+ Self::T,
> {
&self.inner
}
- fn get_lm(&self) -> &LiquidityManagerSync<ES, NS, CM, C, KS, TP> {
+ fn get_lm(&self) -> &LiquidityManagerSync<ES, NS, CM, C, KS, TP, T> {
self
}
}
@@ -282,6 +304,7 @@ pub struct LiquidityManager<
C: Deref + Clone,
K: Deref + Clone,
TP: Deref + Clone,
+ T: Deref + Clone,
> where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -289,6 +312,7 @@ pub struct LiquidityManager<
C::Target: Filter,
K::Target: KVStore,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
pending_messages: Arc<MessageQueue>,
pending_events: Arc<EventQueue<K>>,
@@ -300,7 +324,7 @@ pub struct LiquidityManager<
#[cfg(lsps1_service)]
lsps1_service_handler: Option<LSPS1ServiceHandler<ES, CM, C, K>>,
lsps1_client_handler: Option<LSPS1ClientHandler<ES, K>>,
- lsps2_service_handler: Option<LSPS2ServiceHandler<CM, K>>,
+ lsps2_service_handler: Option<LSPS2ServiceHandler<CM, K, T>>,
lsps2_client_handler: Option<LSPS2ClientHandler<ES, K>>,
lsps5_service_handler: Option<LSPS5ServiceHandler<CM, NS, K, TP>>,
lsps5_client_handler: Option<LSPS5ClientHandler<ES, K>>,
@@ -318,20 +342,22 @@ impl<
CM: Deref + Clone,
C: Deref + Clone,
K: Deref + Clone,
- > LiquidityManager<ES, NS, CM, C, K, DefaultTimeProvider>
+ T: Deref + Clone,
+ > LiquidityManager<ES, NS, CM, C, K, DefaultTimeProvider, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
C::Target: Filter,
K::Target: KVStore,
+ T::Target: BroadcasterInterface,
{
/// Constructor for the [`LiquidityManager`] using the default system clock
///
/// Will read persisted service states from the given [`KVStore`].
pub async fn new(
entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
- chain_params: Option<ChainParameters>, kv_store: K,
+ chain_params: Option<ChainParameters>, kv_store: K, transaction_broadcaster: T,
service_config: Option<LiquidityServiceConfig>,
client_config: Option<LiquidityClientConfig>,
) -> Result<Self, lightning::io::Error> {
@@ -339,6 +365,7 @@ where
entropy_source,
node_signer,
channel_manager,
+ transaction_broadcaster,
chain_source,
chain_params,
kv_store,
@@ -357,7 +384,8 @@ impl<
C: Deref + Clone,
K: Deref + Clone,
TP: Deref + Clone,
- > LiquidityManager<ES, NS, CM, C, K, TP>
+ T: Deref + Clone,
+ > LiquidityManager<ES, NS, CM, C, K, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -365,6 +393,7 @@ where
C::Target: Filter,
K::Target: KVStore,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
/// Constructor for the [`LiquidityManager`] with a custom time provider.
///
@@ -375,8 +404,8 @@ where
/// Sets up the required protocol message handlers based on the given
/// [`LiquidityClientConfig`] and [`LiquidityServiceConfig`].
pub async fn new_with_custom_time_provider(
- entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
- chain_params: Option<ChainParameters>, kv_store: K,
+ entropy_source: ES, node_signer: NS, channel_manager: CM, transaction_broadcaster: T,
+ chain_source: Option<C>, chain_params: Option<ChainParameters>, kv_store: K,
service_config: Option<LiquidityServiceConfig>,
client_config: Option<LiquidityClientConfig>, time_provider: TP,
) -> Result<Self, lightning::io::Error> {
@@ -407,7 +436,7 @@ where
let lsps2_service_handler = if let Some(service_config) = service_config.as_ref() {
if let Some(lsps2_service_config) = service_config.lsps2_service_config.as_ref() {
if let Some(number) =
- <LSPS2ServiceHandler<CM, K> as LSPSProtocolMessageHandler>::PROTOCOL_NUMBER
+ <LSPS2ServiceHandler<CM, K, T> as LSPSProtocolMessageHandler>::PROTOCOL_NUMBER
{
supported_protocols.push(number);
}
@@ -419,6 +448,7 @@ where
Arc::clone(&pending_events),
channel_manager.clone(),
kv_store.clone(),
+ transaction_broadcaster.clone(),
lsps2_service_config.clone(),
)?)
} else {
@@ -565,7 +595,7 @@ where
/// Returns a reference to the LSPS2 server-side handler.
///
/// The returned hendler allows to initiate the LSPS2 service-side flow.
- pub fn lsps2_service_handler(&self) -> Option<&LSPS2ServiceHandler<CM, K>> {
+ pub fn lsps2_service_handler(&self) -> Option<&LSPS2ServiceHandler<CM, K, T>> {
self.lsps2_service_handler.as_ref()
}
@@ -777,7 +807,8 @@ impl<
C: Deref + Clone,
K: Deref + Clone,
TP: Deref + Clone,
- > CustomMessageReader for LiquidityManager<ES, NS, CM, C, K, TP>
+ T: Deref + Clone,
+ > CustomMessageReader for LiquidityManager<ES, NS, CM, C, K, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -785,6 +816,7 @@ where
C::Target: Filter,
K::Target: KVStore,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
type CustomMessage = RawLSPSMessage;
@@ -807,7 +839,8 @@ impl<
C: Deref + Clone,
K: Deref + Clone,
TP: Deref + Clone,
- > CustomMessageHandler for LiquidityManager<ES, NS, CM, C, K, TP>
+ T: Deref + Clone,
+ > CustomMessageHandler for LiquidityManager<ES, NS, CM, C, K, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -815,6 +848,7 @@ where
C::Target: Filter,
K::Target: KVStore,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
fn handle_custom_message(
&self, msg: Self::CustomMessage, sender_node_id: PublicKey,
@@ -939,7 +973,8 @@ impl<
C: Deref + Clone,
K: Deref + Clone,
TP: Deref + Clone,
- > Listen for LiquidityManager<ES, NS, CM, C, K, TP>
+ T: Deref + Clone,
+ > Listen for LiquidityManager<ES, NS, CM, C, K, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -947,6 +982,7 @@ where
C::Target: Filter,
K::Target: KVStore,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
fn filtered_block_connected(
&self, header: &bitcoin::block::Header, txdata: &chain::transaction::TransactionData,
@@ -983,7 +1019,8 @@ impl<
C: Deref + Clone,
K: Deref + Clone,
TP: Deref + Clone,
- > Confirm for LiquidityManager<ES, NS, CM, C, K, TP>
+ T: Deref + Clone,
+ > Confirm for LiquidityManager<ES, NS, CM, C, K, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -991,6 +1028,7 @@ where
C::Target: Filter,
K::Target: KVStore,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
fn transactions_confirmed(
&self, _header: &bitcoin::block::Header, _txdata: &chain::transaction::TransactionData,
@@ -1027,6 +1065,7 @@ pub struct LiquidityManagerSync<
C: Deref + Clone,
KS: Deref + Clone,
TP: Deref + Clone,
+ T: Deref + Clone,
> where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -1034,8 +1073,9 @@ pub struct LiquidityManagerSync<
C::Target: Filter,
KS::Target: KVStoreSync,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
- inner: LiquidityManager<ES, NS, CM, C, KVStoreSyncWrapper<KS>, TP>,
+ inner: LiquidityManager<ES, NS, CM, C, KVStoreSyncWrapper<KS>, TP, T>,
}
#[cfg(feature = "time")]
@@ -1045,20 +1085,22 @@ impl<
CM: Deref + Clone,
C: Deref + Clone,
KS: Deref + Clone,
- > LiquidityManagerSync<ES, NS, CM, C, KS, DefaultTimeProvider>
+ T: Deref + Clone,
+ > LiquidityManagerSync<ES, NS, CM, C, KS, DefaultTimeProvider, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
CM::Target: AChannelManager,
KS::Target: KVStoreSync,
C::Target: Filter,
+ T::Target: BroadcasterInterface,
{
/// Constructor for the [`LiquidityManagerSync`] using the default system clock
///
/// Wraps [`LiquidityManager::new`].
pub fn new(
entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
- chain_params: Option<ChainParameters>, kv_store_sync: KS,
+ chain_params: Option<ChainParameters>, kv_store_sync: KS, transaction_broadcaster: T,
service_config: Option<LiquidityServiceConfig>,
client_config: Option<LiquidityClientConfig>,
) -> Result<Self, lightning::io::Error> {
@@ -1071,6 +1113,7 @@ where
chain_source,
chain_params,
kv_store,
+ transaction_broadcaster,
service_config,
client_config,
));
@@ -1095,7 +1138,8 @@ impl<
C: Deref + Clone,
KS: Deref + Clone,
TP: Deref + Clone,
- > LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+ T: Deref + Clone,
+ > LiquidityManagerSync<ES, NS, CM, C, KS, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -1103,13 +1147,14 @@ where
C::Target: Filter,
KS::Target: KVStoreSync,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
/// Constructor for the [`LiquidityManagerSync`] with a custom time provider.
///
/// Wraps [`LiquidityManager::new_with_custom_time_provider`].
pub fn new_with_custom_time_provider(
entropy_source: ES, node_signer: NS, channel_manager: CM, chain_source: Option<C>,
- chain_params: Option<ChainParameters>, kv_store_sync: KS,
+ chain_params: Option<ChainParameters>, kv_store_sync: KS, transaction_broadcaster: T,
service_config: Option<LiquidityServiceConfig>,
client_config: Option<LiquidityClientConfig>, time_provider: TP,
) -> Result<Self, lightning::io::Error> {
@@ -1118,6 +1163,7 @@ where
entropy_source,
node_signer,
channel_manager,
+ transaction_broadcaster,
chain_source,
chain_params,
kv_store,
@@ -1181,7 +1227,7 @@ where
/// Wraps [`LiquidityManager::lsps2_service_handler`].
pub fn lsps2_service_handler<'a>(
&'a self,
- ) -> Option<LSPS2ServiceHandlerSync<'a, CM, KVStoreSyncWrapper<KS>>> {
+ ) -> Option<LSPS2ServiceHandlerSync<'a, CM, KVStoreSyncWrapper<KS>, T>> {
self.inner.lsps2_service_handler.as_ref().map(|r| LSPS2ServiceHandlerSync::from_inner(r))
}
@@ -1260,7 +1306,8 @@ impl<
C: Deref + Clone,
KS: Deref + Clone,
TP: Deref + Clone,
- > CustomMessageReader for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+ T: Deref + Clone,
+ > CustomMessageReader for LiquidityManagerSync<ES, NS, CM, C, KS, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -1268,6 +1315,7 @@ where
C::Target: Filter,
KS::Target: KVStoreSync,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
type CustomMessage = RawLSPSMessage;
@@ -1285,7 +1333,8 @@ impl<
C: Deref + Clone,
KS: Deref + Clone,
TP: Deref + Clone,
- > CustomMessageHandler for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+ T: Deref + Clone,
+ > CustomMessageHandler for LiquidityManagerSync<ES, NS, CM, C, KS, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -1293,6 +1342,7 @@ where
C::Target: Filter,
KS::Target: KVStoreSync,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
fn handle_custom_message(
&self, msg: Self::CustomMessage, sender_node_id: PublicKey,
@@ -1330,7 +1380,8 @@ impl<
C: Deref + Clone,
KS: Deref + Clone,
TP: Deref + Clone,
- > Listen for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+ T: Deref + Clone,
+ > Listen for LiquidityManagerSync<ES, NS, CM, C, KS, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -1338,6 +1389,7 @@ where
C::Target: Filter,
KS::Target: KVStoreSync,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
fn filtered_block_connected(
&self, header: &bitcoin::block::Header, txdata: &chain::transaction::TransactionData,
@@ -1358,7 +1410,8 @@ impl<
C: Deref + Clone,
KS: Deref + Clone,
TP: Deref + Clone,
- > Confirm for LiquidityManagerSync<ES, NS, CM, C, KS, TP>
+ T: Deref + Clone,
+ > Confirm for LiquidityManagerSync<ES, NS, CM, C, KS, TP, T>
where
ES::Target: EntropySource,
NS::Target: NodeSigner,
@@ -1366,6 +1419,7 @@ where
C::Target: Filter,
KS::Target: KVStoreSync,
TP::Target: TimeProvider,
+ T::Target: BroadcasterInterface,
{
fn transactions_confirmed(
&self, header: &bitcoin::block::Header, txdata: &chain::transaction::TransactionData,
diff --git a/lightning-liquidity/tests/common/mod.rs b/lightning-liquidity/tests/common/mod.rs
index 08d705c..dea9875 100644
--- a/lightning-liquidity/tests/common/mod.rs
+++ b/lightning-liquidity/tests/common/mod.rs
@@ -6,7 +6,7 @@ use lightning_liquidity::{LiquidityClientConfig, LiquidityManagerSync, Liquidity
use lightning::chain::{BestBlock, Filter};
use lightning::ln::channelmanager::ChainParameters;
use lightning::ln::functional_test_utils::{Node, TestChannelManager};
-use lightning::util::test_utils::{TestKeysInterface, TestStore};
+use lightning::util::test_utils::{TestBroadcaster, TestKeysInterface, TestStore};
use bitcoin::Network;
@@ -19,22 +19,31 @@ pub(crate) struct LSPSNodes<'a, 'b, 'c> {
pub client_node: LiquidityNode<'a, 'b, 'c>,
}
-pub(crate) fn create_service_and_client_nodes_with_kv_stores<'a, 'b, 'c>(
+fn build_service_and_client_nodes<'a, 'b, 'c>(
nodes: Vec<Node<'a, 'b, 'c>>, service_config: LiquidityServiceConfig,
client_config: LiquidityClientConfig, time_provider: Arc<dyn TimeProvider + Send + Sync>,
service_kv_store: Arc<TestStore>, client_kv_store: Arc<TestStore>,
-) -> LSPSNodes<'a, 'b, 'c> {
+) -> (LiquidityNode<'a, 'b, 'c>, LiquidityNode<'a, 'b, 'c>, Option<Node<'a, 'b, 'c>>) {
+ assert!(nodes.len() >= 2, "Need at least two nodes (service and client)");
+
let chain_params = ChainParameters {
network: Network::Testnet,
best_block: BestBlock::from_network(Network::Testnet),
};
+
+ let mut nodes_iter = nodes.into_iter();
+ let service_inner = nodes_iter.next().expect("missing service node");
+ let client_inner = nodes_iter.next().expect("missing client node");
+ let leftover = nodes_iter.next();
+
let service_lm = LiquidityManagerSync::new_with_custom_time_provider(
- nodes[0].keys_manager,
- nodes[0].keys_manager,
- nodes[0].node,
+ service_inner.keys_manager,
+ service_inner.keys_manager,
+ service_inner.node,
None::<Arc<dyn Filter + Send + Sync>>,
Some(chain_params.clone()),
service_kv_store,
+ service_inner.tx_broadcaster,
Some(service_config),
None,
Arc::clone(&time_provider),
@@ -42,26 +51,49 @@ pub(crate) fn create_service_and_client_nodes_with_kv_stores<'a, 'b, 'c>(
.unwrap();
let client_lm = LiquidityManagerSync::new_with_custom_time_provider(
- nodes[1].keys_manager,
- nodes[1].keys_manager,
- nodes[1].node,
+ client_inner.keys_manager,
+ client_inner.keys_manager,
+ client_inner.node,
None::<Arc<dyn Filter + Send + Sync>>,
Some(chain_params),
client_kv_store,
+ client_inner.tx_broadcaster,
None,
Some(client_config),
time_provider,
)
.unwrap();
- let mut iter = nodes.into_iter();
- let service_node = LiquidityNode::new(iter.next().unwrap(), service_lm);
- let client_node = LiquidityNode::new(iter.next().unwrap(), client_lm);
+ let service_node = LiquidityNode::new(service_inner, service_lm);
+ let client_node = LiquidityNode::new(client_inner, client_lm);
+ (service_node, client_node, leftover)
+}
+// this is ONLY used on LSPS2 so it says it's not used but it is
+#[allow(dead_code)]
+pub(crate) struct LSPSNodesWithPayer<'a, 'b, 'c> {
+ pub service_node: LiquidityNode<'a, 'b, 'c>,
+ pub client_node: LiquidityNode<'a, 'b, 'c>,
+ pub payer_node: Node<'a, 'b, 'c>,
+}
+
+pub(crate) fn create_service_and_client_nodes_with_kv_stores<'a, 'b, 'c>(
+ nodes: Vec<Node<'a, 'b, 'c>>, service_config: LiquidityServiceConfig,
+ client_config: LiquidityClientConfig, time_provider: Arc<dyn TimeProvider + Send + Sync>,
+ service_kv_store: Arc<TestStore>, client_kv_store: Arc<TestStore>,
+) -> LSPSNodes<'a, 'b, 'c> {
+ let (service_node, client_node, _) = build_service_and_client_nodes(
+ nodes,
+ service_config,
+ client_config,
+ time_provider,
+ service_kv_store,
+ client_kv_store,
+ );
LSPSNodes { service_node, client_node }
}
-#[allow(unused)]
+#[allow(dead_code)]
pub(crate) fn create_service_and_client_nodes<'a, 'b, 'c>(
nodes: Vec<Node<'a, 'b, 'c>>, service_config: LiquidityServiceConfig,
client_config: LiquidityClientConfig, time_provider: Arc<dyn TimeProvider + Send + Sync>,
@@ -78,6 +110,27 @@ pub(crate) fn create_service_and_client_nodes<'a, 'b, 'c>(
)
}
+// this is ONLY used on LSPS2 so it says it's not used but it is
+#[allow(dead_code)]
+pub(crate) fn create_service_client_and_payer_nodes<'a, 'b, 'c>(
+ nodes: Vec<Node<'a, 'b, 'c>>, service_config: LiquidityServiceConfig,
+ client_config: LiquidityClientConfig, time_provider: Arc<dyn TimeProvider + Send + Sync>,
+) -> LSPSNodesWithPayer<'a, 'b, 'c> {
+ assert!(nodes.len() >= 3, "Need three nodes (service, client, payer)");
+ let service_kv_store = Arc::new(TestStore::new(false));
+ let client_kv_store = Arc::new(TestStore::new(false));
+ let (service_node, client_node, payer_opt) = build_service_and_client_nodes(
+ nodes,
+ service_config,
+ client_config,
+ time_provider,
+ service_kv_store,
+ client_kv_store,
+ );
+ let payer_node = payer_opt.expect("payer node missing");
+ LSPSNodesWithPayer { service_node, client_node, payer_node }
+}
+
pub(crate) struct LiquidityNode<'a, 'b, 'c> {
pub inner: Node<'a, 'b, 'c>,
pub liquidity_manager: LiquidityManagerSync<
@@ -87,6 +140,7 @@ pub(crate) struct LiquidityNode<'a, 'b, 'c> {
Arc<dyn Filter + Send + Sync>,
Arc<TestStore>,
Arc<dyn TimeProvider + Send + Sync>,
+ &'c TestBroadcaster,
>,
}
@@ -100,6 +154,7 @@ impl<'a, 'b, 'c> LiquidityNode<'a, 'b, 'c> {
Arc<dyn Filter + Send + Sync>,
Arc<TestStore>,
Arc<dyn TimeProvider + Send + Sync>,
+ &'c TestBroadcaster,
>,
) -> Self {
Self { inner: node, liquidity_manager }
diff --git a/lightning-liquidity/tests/lsps2_integration_tests.rs b/lightning-liquidity/tests/lsps2_integration_tests.rs
index da884c7..82f93b5 100644
--- a/lightning-liquidity/tests/lsps2_integration_tests.rs
+++ b/lightning-liquidity/tests/lsps2_integration_tests.rs
@@ -3,9 +3,28 @@
mod common;
use common::{
- create_service_and_client_nodes_with_kv_stores, get_lsps_message, LSPSNodes, LiquidityNode,
+ create_service_and_client_nodes_with_kv_stores, create_service_client_and_payer_nodes,
+ get_lsps_message, LSPSNodes, LSPSNodesWithPayer, LiquidityNode,
};
+use lightning::check_added_monitors;
+use lightning::events::{ClosureReason, Event};
+use lightning::get_event_msg;
+use lightning::ln::channelmanager::PaymentId;
+use lightning::ln::channelmanager::Retry;
+use lightning::ln::functional_test_utils::create_funding_transaction;
+use lightning::ln::functional_test_utils::do_commitment_signed_dance;
+use lightning::ln::functional_test_utils::expect_channel_pending_event;
+use lightning::ln::functional_test_utils::expect_channel_ready_event;
+use lightning::ln::functional_test_utils::expect_payment_sent;
+use lightning::ln::functional_test_utils::test_default_channel_config;
+use lightning::ln::functional_test_utils::SendEvent;
+use lightning::ln::functional_test_utils::{connect_blocks, create_chan_between_nodes_with_value};
+use lightning::ln::msgs::BaseMessageHandler;
+use lightning::ln::msgs::ChannelMessageHandler;
+use lightning::ln::msgs::MessageSendEvent;
+use lightning::ln::types::ChannelId;
+
use lightning_liquidity::events::LiquidityEvent;
use lightning_liquidity::lsps0::ser::LSPSDateTime;
use lightning_liquidity::lsps2::client::LSPS2ClientConfig;
@@ -29,7 +48,7 @@ use lightning::routing::router::{RouteHint, RouteHintHop};
use lightning::sign::NodeSigner;
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
-use lightning::util::test_utils::TestStore;
+use lightning::util::test_utils::{TestBroadcaster, TestStore};
use lightning_invoice::{Bolt11Invoice, InvoiceBuilder, RoutingFees};
@@ -38,6 +57,7 @@ use lightning_types::payment::PaymentHash;
use bitcoin::hashes::{sha256, Hash};
use bitcoin::secp256k1::{PublicKey, Secp256k1, SecretKey};
use bitcoin::Network;
+use lightning_types::payment::PaymentPreimage;
use std::str::FromStr;
use std::sync::Arc;
@@ -46,9 +66,7 @@ use std::time::Duration;
const MAX_PENDING_REQUESTS_PER_PEER: usize = 10;
const MAX_TOTAL_PENDING_REQUESTS: usize = 1000;
-fn setup_test_lsps2_nodes_with_kv_stores<'a, 'b, 'c>(
- nodes: Vec<Node<'a, 'b, 'c>>, service_kv_store: Arc<TestStore>, client_kv_store: Arc<TestStore>,
-) -> (LSPSNodes<'a, 'b, 'c>, [u8; 32]) {
+fn build_lsps2_configs() -> ([u8; 32], LiquidityServiceConfig, LiquidityClientConfig) {
let promise_secret = [42; 32];
let lsps2_service_config = LSPS2ServiceConfig { promise_secret };
let service_config = LiquidityServiceConfig {
@@ -65,6 +83,14 @@ fn setup_test_lsps2_nodes_with_kv_stores<'a, 'b, 'c>(
lsps2_client_config: Some(lsps2_client_config),
lsps5_client_config: None,
};
+
+ (promise_secret, service_config, client_config)
+}
+
+fn setup_test_lsps2_nodes_with_kv_stores<'a, 'b, 'c>(
+ nodes: Vec<Node<'a, 'b, 'c>>, service_kv_store: Arc<TestStore>, client_kv_store: Arc<TestStore>,
+) -> (LSPSNodes<'a, 'b, 'c>, [u8; 32]) {
+ let (promise_secret, service_config, client_config) = build_lsps2_configs();
let lsps_nodes = create_service_and_client_nodes_with_kv_stores(
nodes,
service_config,
@@ -73,7 +99,6 @@ fn setup_test_lsps2_nodes_with_kv_stores<'a, 'b, 'c>(
service_kv_store,
client_kv_store,
);
-
(lsps_nodes, promise_secret)
}
@@ -85,6 +110,19 @@ fn setup_test_lsps2_nodes<'a, 'b, 'c>(
setup_test_lsps2_nodes_with_kv_stores(nodes, service_kv_store, client_kv_store)
}
+fn setup_test_lsps2_nodes_with_payer<'a, 'b, 'c>(
+ nodes: Vec<Node<'a, 'b, 'c>>,
+) -> (LSPSNodesWithPayer<'a, 'b, 'c>, [u8; 32]) {
+ let (promise_secret, service_config, client_config) = build_lsps2_configs();
+ let lsps_nodes = create_service_client_and_payer_nodes(
+ nodes,
+ service_config,
+ client_config,
+ Arc::new(DefaultTimeProvider),
+ );
+ (lsps_nodes, promise_secret)
+}
+
fn create_jit_invoice(
node: &LiquidityNode<'_, '_, '_>, service_node_id: PublicKey, intercept_scid: u64,
cltv_expiry_delta: u32, payment_size_msat: Option<u64>, description: &str, expiry_secs: u32,
@@ -132,11 +170,13 @@ fn create_jit_invoice(
let sign_fn =
node.inner.keys_manager.sign_invoice(&raw_invoice, lightning::sign::Recipient::Node);
- raw_invoice.sign(|_| sign_fn).and_then(|signed_raw| {
+ let invoice = raw_invoice.sign(|_| sign_fn).and_then(|signed_raw| {
Bolt11Invoice::from_signed(signed_raw).map_err(|e| {
log_error!(node.inner.logger, "Failed to create invoice from signed raw: {:?}", e);
})
- })
+ })?;
+
+ Ok(invoice)
}
#[test]
@@ -1049,6 +1089,8 @@ fn lsps2_service_handler_persistence_across_restarts() {
best_block: BestBlock::from_network(Network::Testnet),
};
+ let transaction_broadcaster = Arc::new(TestBroadcaster::new(Network::Testnet));
+
let restarted_service_lm = LiquidityManagerSync::new_with_custom_time_provider(
nodes_restart[0].keys_manager,
nodes_restart[0].keys_manager,
@@ -1056,6 +1098,7 @@ fn lsps2_service_handler_persistence_across_restarts() {
None::<Arc<dyn Filter + Send + Sync>>,
Some(chain_params),
service_kv_store,
+ transaction_broadcaster,
Some(service_config),
None,
time_provider,
@@ -1097,3 +1140,1243 @@ fn lsps2_service_handler_persistence_across_restarts() {
}
}
}
+
+#[test]
+fn client_trusts_lsp_end_to_end_test() {
+ // There are 3 nodes. Payer, service and client.
+ // client_trusts_lsp=true, that means that funding transaction broadcast will need to happen manually
+ // after the client claims the HTLC.
+ //
+ // 1. Create a channel between payer and service
+ // 2. Do the LSPS2 ceremony between client and service, to prepare the service to intercept an htlc and eventually create a JIT channel
+ // 3. Make the client create a JIT invoice and make the payer pay it
+ // 4. Assert that the service intercepts the HTLC
+ // 5. Assert that the service emits a LiquidityEvent::OpenChannel. This means that the intercepted HTLC was enough
+ // and that it's ready to proceed with channel creation.
+ // 6. Proceed with the JIT channel creation (we create it with funding_transaction_generated_manual_broadcast because
+ // client_trusts_lsp=true).
+ // 7. Call the service's channel_ready function
+ // 8. The service will now forward the intercepted HTLC to the client on the new JIT channel
+ // 9. The client will see the PaymentClaimable event
+ // 10. Assert that the service has not broadcasted the funding transaction yet, because the client has not claimed the HTLC yet
+ // 11. Make the client claim the HTLC
+ // 12. Assert that the service has broadcasted the funding tx
+ // 13. Assert that the payer received the PaymentSent event
+ let chanmon_cfgs = create_chanmon_cfgs(3);
+ let node_cfgs = create_node_cfgs(3, &chanmon_cfgs);
+ let mut service_node_config = test_default_channel_config();
+ service_node_config.accept_intercept_htlcs = true;
+
+ let mut client_node_config = test_default_channel_config();
+ client_node_config.manually_accept_inbound_channels = true;
+ client_node_config.channel_config.accept_underpaying_htlcs = true;
+ let node_chanmgrs = create_node_chanmgrs(
+ 3,
+ &node_cfgs,
+ &[Some(service_node_config), Some(client_node_config), None],
+ );
+ let nodes = create_network(3, &node_cfgs, &node_chanmgrs);
+ let (lsps_nodes, promise_secret) = setup_test_lsps2_nodes_with_payer(nodes);
+ let LSPSNodesWithPayer { ref service_node, ref client_node, ref payer_node } = lsps_nodes;
+
+ let payer_node_id = payer_node.node.get_our_node_id();
+ let service_node_id = service_node.inner.node.get_our_node_id();
+ let client_node_id = client_node.inner.node.get_our_node_id();
+
+ let service_handler = service_node.liquidity_manager.lsps2_service_handler().unwrap();
+
+ create_chan_between_nodes_with_value(&payer_node, &service_node.inner, 2000000, 100000);
+
+ let intercept_scid = service_node.node.get_intercept_scid();
+ let user_channel_id = 42;
+ let cltv_expiry_delta: u32 = 144;
+ let payment_size_msat = Some(1_000_000);
+
+ let fee_base_msat = 1000;
+
+ execute_lsps2_dance(
+ &lsps_nodes,
+ intercept_scid,
+ user_channel_id,
+ cltv_expiry_delta,
+ promise_secret,
+ payment_size_msat,
+ fee_base_msat,
+ );
+
+ let invoice = create_jit_invoice(
+ &client_node,
+ service_node_id,
+ intercept_scid,
+ cltv_expiry_delta,
+ payment_size_msat,
+ "asdf",
+ 3600,
+ )
+ .unwrap();
+
+ payer_node
+ .node
+ .pay_for_bolt11_invoice(
+ &invoice,
+ PaymentId(invoice.payment_hash().to_byte_array()),
+ None,
+ Default::default(),
+ Retry::Attempts(3),
+ )
+ .unwrap();
+
+ check_added_monitors!(payer_node, 1);
+ let events = payer_node.node.get_and_clear_pending_msg_events();
+ let ev = SendEvent::from_event(events[0].clone());
+ service_node.inner.node.handle_update_add_htlc(payer_node_id, &ev.msgs[0]);
+ do_commitment_signed_dance(&service_node.inner, &payer_node, &ev.commitment_msg, false, true);
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ let events = service_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(events.len(), 1);
+ let (payment_hash, expected_outbound_amount_msat) = match &events[0] {
+ Event::HTLCIntercepted {
+ intercept_id,
+ requested_next_hop_scid,
+ payment_hash,
+ expected_outbound_amount_msat,
+ ..
+ } => {
+ assert_eq!(*requested_next_hop_scid, intercept_scid);
+
+ service_handler
+ .htlc_intercepted(
+ *requested_next_hop_scid,
+ *intercept_id,
+ *expected_outbound_amount_msat,
+ *payment_hash,
+ )
+ .unwrap();
+ (*payment_hash, expected_outbound_amount_msat)
+ },
+ other => panic!("Expected HTLCIntercepted event, got: {:?}", other),
+ };
+
+ let open_channel_event = service_node.liquidity_manager.next_event().unwrap();
+
+ match open_channel_event {
+ LiquidityEvent::LSPS2Service(LSPS2ServiceEvent::OpenChannel {
+ their_network_key,
+ amt_to_forward_msat,
+ opening_fee_msat,
+ user_channel_id,
+ intercept_scid: iscd,
+ }) => {
+ assert_eq!(their_network_key, client_node_id);
+ assert_eq!(amt_to_forward_msat, payment_size_msat.unwrap() - fee_base_msat);
+ assert_eq!(opening_fee_msat, fee_base_msat);
+ assert_eq!(user_channel_id, 42);
+ assert_eq!(iscd, intercept_scid);
+ },
+ other => panic!("Expected OpenChannel event, got: {:?}", other),
+ };
+
+ let result =
+ service_handler.channel_needs_manual_broadcast(user_channel_id, &client_node_id).unwrap();
+ assert!(result, "Channel should require manual broadcast");
+
+ let (channel_id, funding_tx) = create_channel_with_manual_broadcast(
+ &service_node_id,
+ &client_node_id,
+ &service_node,
+ &client_node,
+ user_channel_id,
+ expected_outbound_amount_msat,
+ true,
+ );
+
+ service_handler.channel_ready(user_channel_id, &channel_id, &client_node_id).unwrap();
+
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ let pay_event = {
+ {
+ let mut added_monitors =
+ service_node.inner.chain_monitor.added_monitors.lock().unwrap();
+ assert_eq!(added_monitors.len(), 1);
+ added_monitors.clear();
+ }
+ let mut events = service_node.inner.node.get_and_clear_pending_msg_events();
+ assert_eq!(events.len(), 1);
+ SendEvent::from_event(events.remove(0))
+ };
+
+ client_node.inner.node.handle_update_add_htlc(service_node_id, &pay_event.msgs[0]);
+ do_commitment_signed_dance(
+ &client_node.inner,
+ &service_node.inner,
+ &pay_event.commitment_msg,
+ false,
+ true,
+ );
+ client_node.inner.node.process_pending_htlc_forwards();
+
+ let client_events = client_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(client_events.len(), 1);
+ let preimage = match &client_events[0] {
+ Event::PaymentClaimable { payment_hash: ph, purpose, .. } => {
+ assert_eq!(*ph, payment_hash);
+ purpose.preimage()
+ },
+ other => panic!("Expected PaymentClaimable event on client, got: {:?}", other),
+ };
+
+ // Check that before the client claims, the service node has not broadcasted anything
+ let broadcasted = service_node.inner.tx_broadcaster.txn_broadcasted.lock().unwrap();
+ assert!(broadcasted.is_empty(), "There should be no broadcasted txs yet");
+ drop(broadcasted);
+
+ client_node.inner.node.claim_funds(preimage.unwrap());
+
+ claim_and_assert_forwarded_only(
+ &payer_node,
+ &service_node.inner,
+ &client_node.inner,
+ preimage.unwrap(),
+ );
+
+ let service_events = service_node.node.get_and_clear_pending_events();
+ assert_eq!(service_events.len(), 1);
+
+ let total_fee_msat = match service_events[0].clone() {
+ Event::PaymentForwarded {
+ prev_node_id,
+ next_node_id,
+ skimmed_fee_msat,
+ total_fee_earned_msat,
+ ..
+ } => {
+ assert_eq!(prev_node_id, Some(payer_node_id));
+ assert_eq!(next_node_id, Some(client_node_id));
+ service_handler.payment_forwarded(channel_id, skimmed_fee_msat.unwrap_or(0)).unwrap();
+ Some(total_fee_earned_msat.unwrap() - skimmed_fee_msat.unwrap())
+ },
+ _ => panic!("Expected PaymentForwarded event, got: {:?}", service_events[0]),
+ };
+
+ let broadcasted = service_node.inner.tx_broadcaster.txn_broadcasted.lock().unwrap();
+ assert!(broadcasted.iter().any(|b| b.compute_txid() == funding_tx.compute_txid()));
+
+ expect_payment_sent(&payer_node, preimage.unwrap(), Some(total_fee_msat), true, true);
+}
+
+fn execute_lsps2_dance(
+ lsps_nodes: &LSPSNodesWithPayer, intercept_scid: u64, user_channel_id: u128,
+ cltv_expiry_delta: u32, promise_secret: [u8; 32], payment_size_msat: Option<u64>,
+ fee_base_msat: u64,
+) {
+ let service_node = &lsps_nodes.service_node;
+ let client_node = &lsps_nodes.client_node;
+
+ let service_node_id = service_node.inner.node.get_our_node_id();
+ let client_node_id = client_node.inner.node.get_our_node_id();
+
+ let client_handler = client_node.liquidity_manager.lsps2_client_handler().unwrap();
+ let service_handler = service_node.liquidity_manager.lsps2_service_handler().unwrap();
+
+ let get_info_request_id = client_handler.request_opening_params(service_node_id, None);
+ let get_info_request = get_lsps_message!(client_node, service_node_id);
+
+ service_node.liquidity_manager.handle_custom_message(get_info_request, client_node_id).unwrap();
+
+ let get_info_event = service_node.liquidity_manager.next_event().unwrap();
+ if let LiquidityEvent::LSPS2Service(LSPS2ServiceEvent::GetInfo {
+ request_id,
+ counterparty_node_id,
+ token,
+ }) = get_info_event
+ {
+ assert_eq!(request_id, get_info_request_id);
+ assert_eq!(counterparty_node_id, client_node_id);
+ assert_eq!(token, None);
+ } else {
+ panic!("Unexpected event");
+ }
+
+ let raw_opening_params = LSPS2RawOpeningFeeParams {
+ min_fee_msat: fee_base_msat,
+ proportional: 0,
+ valid_until: LSPSDateTime::from_str("2035-05-20T08:30:45Z").unwrap(),
+ min_lifetime: 144,
+ max_client_to_self_delay: 128,
+ min_payment_size_msat: 1,
+ max_payment_size_msat: 100_000_000,
+ };
+
+ service_handler
+ .opening_fee_params_generated(
+ &client_node_id,
+ get_info_request_id.clone(),
+ vec![raw_opening_params],
+ )
+ .unwrap();
+ let get_info_response = get_lsps_message!(service_node, client_node_id);
+
+ client_node
+ .liquidity_manager
+ .handle_custom_message(get_info_response, service_node_id)
+ .unwrap();
+
+ let opening_params_event = client_node.liquidity_manager.next_event().unwrap();
+ let opening_fee_params = match opening_params_event {
+ LiquidityEvent::LSPS2Client(LSPS2ClientEvent::OpeningParametersReady {
+ request_id,
+ counterparty_node_id,
+ opening_fee_params_menu,
+ }) => {
+ assert_eq!(request_id, get_info_request_id);
+ assert_eq!(counterparty_node_id, service_node_id);
+ let opening_fee_params = opening_fee_params_menu.first().unwrap().clone();
+ assert!(is_valid_opening_fee_params(
+ &opening_fee_params,
+ &promise_secret,
+ &client_node_id
+ ));
+ opening_fee_params
+ },
+ _ => panic!("Unexpected event"),
+ };
+
+ let buy_request_id = client_handler
+ .select_opening_params(service_node_id, payment_size_msat, opening_fee_params.clone())
+ .unwrap();
+
+ let buy_request = get_lsps_message!(client_node, service_node_id);
+ service_node.liquidity_manager.handle_custom_message(buy_request, client_node_id).unwrap();
+
+ let buy_event = service_node.liquidity_manager.next_event().unwrap();
+ if let LiquidityEvent::LSPS2Service(LSPS2ServiceEvent::BuyRequest {
+ request_id,
+ counterparty_node_id,
+ opening_fee_params: ofp,
+ payment_size_msat: psm,
+ }) = buy_event
+ {
+ assert_eq!(request_id, buy_request_id);
+ assert_eq!(counterparty_node_id, client_node_id);
+ assert_eq!(opening_fee_params, ofp);
+ assert_eq!(payment_size_msat, psm);
+ } else {
+ panic!("Unexpected event");
+ }
+
+ let client_trusts_lsp = true;
+
+ service_handler
+ .invoice_parameters_generated(
+ &client_node_id,
+ buy_request_id.clone(),
+ intercept_scid,
+ cltv_expiry_delta,
+ client_trusts_lsp,
+ user_channel_id,
+ )
+ .unwrap();
+
+ let buy_response = get_lsps_message!(service_node, client_node_id);
+ client_node.liquidity_manager.handle_custom_message(buy_response, service_node_id).unwrap();
+
+ let invoice_params_event = client_node.liquidity_manager.next_event().unwrap();
+ if let LiquidityEvent::LSPS2Client(LSPS2ClientEvent::InvoiceParametersReady {
+ request_id,
+ counterparty_node_id,
+ intercept_scid: iscid,
+ cltv_expiry_delta: ced,
+ payment_size_msat: psm,
+ }) = invoice_params_event
+ {
+ assert_eq!(request_id, buy_request_id);
+ assert_eq!(counterparty_node_id, service_node_id);
+ assert_eq!(intercept_scid, iscid);
+ assert_eq!(cltv_expiry_delta, ced);
+ assert_eq!(payment_size_msat, psm);
+ } else {
+ panic!("Unexpected event");
+ }
+}
+
+fn create_channel_with_manual_broadcast(
+ service_node_id: &PublicKey, client_node_id: &PublicKey, service_node: &LiquidityNode,
+ client_node: &LiquidityNode, user_channel_id: u128, expected_outbound_amount_msat: &u64,
+ mark_broadcast_safe: bool,
+) -> (ChannelId, bitcoin::Transaction) {
+ assert!(service_node
+ .node
+ .create_channel(
+ *client_node_id,
+ *expected_outbound_amount_msat,
+ 0,
+ user_channel_id,
+ None,
+ None
+ )
+ .is_ok());
+ let open_channel =
+ get_event_msg!(service_node, MessageSendEvent::SendOpenChannel, *client_node_id);
+
+ client_node.node.handle_open_channel(*service_node_id, &open_channel);
+
+ let events = client_node.node.get_and_clear_pending_events();
+ assert_eq!(events.len(), 1);
+ match events[0] {
+ Event::OpenChannelRequest { temporary_channel_id, .. } => {
+ client_node
+ .node
+ .accept_inbound_channel_from_trusted_peer_0conf(
+ &temporary_channel_id,
+ &service_node_id,
+ user_channel_id,
+ None,
+ )
+ .unwrap();
+ },
+ _ => panic!("Unexpected event"),
+ };
+
+ let accept_channel =
+ get_event_msg!(client_node, MessageSendEvent::SendAcceptChannel, *service_node_id);
+ assert_eq!(accept_channel.common_fields.minimum_depth, 0);
+
+ service_node.node.handle_accept_channel(*client_node_id, &accept_channel);
+ let (temp_channel_id, funding_tx, funding_outpoint) = create_funding_transaction(
+ &service_node,
+ &client_node_id,
+ *expected_outbound_amount_msat,
+ user_channel_id,
+ );
+ let service_handler = service_node.liquidity_manager.lsps2_service_handler().unwrap();
+ service_handler
+ .store_funding_transaction(user_channel_id, &client_node_id, funding_tx.clone())
+ .unwrap();
+ service_node
+ .node
+ .funding_transaction_generated_manual_broadcast(
+ temp_channel_id,
+ *client_node_id,
+ funding_tx.clone(),
+ )
+ .unwrap();
+
+ let funding_created =
+ get_event_msg!(service_node, MessageSendEvent::SendFundingCreated, *client_node_id);
+ client_node.node.handle_funding_created(*service_node_id, &funding_created);
+ check_added_monitors!(client_node.inner, 1);
+
+ let bs_signed_locked = client_node.node.get_and_clear_pending_msg_events();
+ assert_eq!(bs_signed_locked.len(), 2);
+
+ let as_channel_ready;
+ match &bs_signed_locked[0] {
+ MessageSendEvent::SendFundingSigned { node_id, msg } => {
+ assert_eq!(*node_id, *service_node_id);
+ service_node.node.handle_funding_signed(*client_node_id, &msg);
+ let events = &service_node.node.get_and_clear_pending_events();
+ assert_eq!(events.len(), 2);
+ match &events[0] {
+ Event::FundingTxBroadcastSafe {
+ funding_txo,
+ user_channel_id,
+ counterparty_node_id,
+ ..
+ } => {
+ assert_eq!(funding_txo.txid, funding_outpoint.txid);
+ assert_eq!(funding_txo.vout, funding_outpoint.index as u32);
+ if mark_broadcast_safe {
+ service_handler
+ .set_funding_tx_broadcast_safe(*user_channel_id, counterparty_node_id)
+ .unwrap();
+ }
+ },
+ _ => panic!("Unexpected event"),
+ };
+ match &events[1] {
+ Event::ChannelPending { counterparty_node_id, .. } => {
+ assert_eq!(counterparty_node_id, client_node_id);
+ },
+ _ => panic!("Unexpected event"),
+ }
+ expect_channel_pending_event(&client_node, &service_node_id);
+ check_added_monitors!(service_node.inner, 1);
+
+ as_channel_ready =
+ get_event_msg!(service_node, MessageSendEvent::SendChannelReady, *client_node_id);
+ },
+ _ => panic!("Unexpected event"),
+ }
+
+ match &bs_signed_locked[1] {
+ MessageSendEvent::SendChannelReady { node_id, msg } => {
+ assert_eq!(*node_id, *service_node_id);
+ service_node.node.handle_channel_ready(*client_node_id, &msg);
+ expect_channel_ready_event(&service_node, &client_node_id);
+ },
+ _ => panic!("Unexpected event"),
+ }
+
+ client_node.node.handle_channel_ready(*service_node_id, &as_channel_ready);
+ expect_channel_ready_event(&client_node, &service_node_id);
+
+ let as_channel_update =
+ get_event_msg!(service_node, MessageSendEvent::SendChannelUpdate, *client_node_id);
+ let bs_channel_update =
+ get_event_msg!(client_node, MessageSendEvent::SendChannelUpdate, *service_node_id);
+
+ service_node.node.handle_channel_update(*client_node_id, &bs_channel_update);
+ client_node.node.handle_channel_update(*service_node_id, &as_channel_update);
+
+ (as_channel_ready.channel_id, funding_tx)
+}
+
+#[test]
+fn late_payment_forwarded_and_safe_after_force_close_does_not_broadcast() {
+ let chanmon_cfgs = create_chanmon_cfgs(3);
+ let node_cfgs = create_node_cfgs(3, &chanmon_cfgs);
+ let mut service_node_config = test_default_channel_config();
+ service_node_config.accept_intercept_htlcs = true;
+
+ let mut client_node_config = test_default_channel_config();
+ client_node_config.manually_accept_inbound_channels = true;
+ client_node_config.channel_config.accept_underpaying_htlcs = true;
+
+ let node_chanmgrs = create_node_chanmgrs(
+ 3,
+ &node_cfgs,
+ &[Some(service_node_config), Some(client_node_config), None],
+ );
+ let nodes = create_network(3, &node_cfgs, &node_chanmgrs);
+ let (lsps_nodes, promise_secret) = setup_test_lsps2_nodes_with_payer(nodes);
+ let LSPSNodesWithPayer { ref service_node, ref client_node, ref payer_node } = lsps_nodes;
+
+ let payer_node_id = payer_node.node.get_our_node_id();
+ let service_node_id = service_node.inner.node.get_our_node_id();
+ let client_node_id = client_node.inner.node.get_our_node_id();
+
+ let service_handler = service_node.liquidity_manager.lsps2_service_handler().unwrap();
+
+ create_chan_between_nodes_with_value(&payer_node, &service_node.inner, 2_000_000, 100_000);
+
+ let intercept_scid = service_node.node.get_intercept_scid();
+ let user_channel_id = 43u128;
+ let cltv_expiry_delta: u32 = 144;
+ let payment_size_msat = Some(1_000_000);
+ let fee_base_msat: u64 = 10_000;
+
+ execute_lsps2_dance(
+ &lsps_nodes,
+ intercept_scid,
+ user_channel_id,
+ cltv_expiry_delta,
+ promise_secret,
+ payment_size_msat,
+ fee_base_msat,
+ );
+
+ let invoice = create_jit_invoice(
+ &client_node,
+ service_node_id,
+ intercept_scid,
+ cltv_expiry_delta,
+ payment_size_msat,
+ "late-safe",
+ 3600,
+ )
+ .unwrap();
+
+ payer_node
+ .node
+ .pay_for_bolt11_invoice(
+ &invoice,
+ PaymentId(invoice.payment_hash().to_byte_array()),
+ None,
+ Default::default(),
+ Retry::Attempts(3),
+ )
+ .unwrap();
+
+ check_added_monitors!(payer_node, 1);
+ let events = payer_node.node.get_and_clear_pending_msg_events();
+ let ev = SendEvent::from_event(events[0].clone());
+ service_node.inner.node.handle_update_add_htlc(payer_node_id, &ev.msgs[0]);
+ do_commitment_signed_dance(&service_node.inner, &payer_node, &ev.commitment_msg, false, true);
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ let events = service_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(events.len(), 1);
+ match &events[0] {
+ Event::HTLCIntercepted {
+ intercept_id,
+ requested_next_hop_scid,
+ payment_hash: _,
+ expected_outbound_amount_msat,
+ ..
+ } => {
+ assert_eq!(*requested_next_hop_scid, intercept_scid);
+ service_handler
+ .htlc_intercepted(
+ *requested_next_hop_scid,
+ *intercept_id,
+ *expected_outbound_amount_msat,
+ PaymentHash(invoice.payment_hash().to_byte_array()),
+ )
+ .unwrap();
+ },
+ other => panic!("Expected HTLCIntercepted, got {:?}", other),
+ }
+
+ // Create channel but DO NOT mark broadcast safe yet
+ let (channel_id, funding_tx) = create_channel_with_manual_broadcast(
+ &service_node_id,
+ &client_node_id,
+ &service_node,
+ &client_node,
+ user_channel_id,
+ &(payment_size_msat.unwrap() - fee_base_msat),
+ false,
+ );
+
+ service_handler.channel_ready(user_channel_id, &channel_id, &client_node_id).unwrap();
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ // Run forward to client and let client claim. do not notify service handler yet.
+ let pay_event = {
+ {
+ let mut added_monitors =
+ service_node.inner.chain_monitor.added_monitors.lock().unwrap();
+ assert_eq!(added_monitors.len(), 1);
+ added_monitors.clear();
+ }
+ let mut msg_events = service_node.inner.node.get_and_clear_pending_msg_events();
+ assert_eq!(msg_events.len(), 1);
+ SendEvent::from_event(msg_events.remove(0))
+ };
+
+ client_node.inner.node.handle_update_add_htlc(service_node_id, &pay_event.msgs[0]);
+ do_commitment_signed_dance(
+ &client_node.inner,
+ &service_node.inner,
+ &pay_event.commitment_msg,
+ false,
+ true,
+ );
+ client_node.inner.node.process_pending_htlc_forwards();
+
+ let client_events = client_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(client_events.len(), 1);
+ let preimage = match &client_events[0] {
+ Event::PaymentClaimable { purpose, .. } => purpose.preimage().unwrap(),
+ other => panic!("Expected PaymentClaimable, got {:?}", other),
+ };
+
+ client_node.inner.node.claim_funds(preimage);
+ claim_and_assert_forwarded_only(&payer_node, &service_node.inner, &client_node.inner, preimage);
+
+ // Service now has PaymentForwarded. Record in JIT state but still not safe to broadcast.
+ let events = service_node.node.get_and_clear_pending_events();
+ let skimmed = match events[0].clone() {
+ Event::PaymentForwarded { skimmed_fee_msat, .. } => skimmed_fee_msat.unwrap_or(0),
+ other => panic!("Expected PaymentForwarded, got {:?}", other),
+ };
+ service_handler.payment_forwarded(channel_id, skimmed).unwrap();
+
+ // Force-close the service->client channel
+ service_node
+ .inner
+ .node
+ .force_close_broadcasting_latest_txn(&channel_id, &client_node_id, "test fc".to_string())
+ .unwrap();
+
+ service_node.inner.node.get_and_clear_pending_msg_events();
+ client_node.inner.node.get_and_clear_pending_msg_events();
+ payer_node.node.get_and_clear_pending_msg_events();
+ service_node.inner.node.get_and_clear_pending_events();
+ client_node.inner.node.get_and_clear_pending_events();
+ payer_node.node.get_and_clear_pending_events();
+ service_node.inner.chain_monitor.added_monitors.lock().unwrap().clear();
+ client_node.inner.chain_monitor.added_monitors.lock().unwrap().clear();
+ payer_node.chain_monitor.added_monitors.lock().unwrap().clear();
+
+ // Simulate late FundingTxBroadcastSafe arrival after close. ensure no broadcast of funding tx.
+ service_handler.set_funding_tx_broadcast_safe(user_channel_id, &client_node_id).unwrap();
+ {
+ let broadcasted = service_node.inner.tx_broadcaster.txn_broadcasted.lock().unwrap();
+ assert!(
+ broadcasted.iter().all(|tx| tx.compute_txid() != funding_tx.compute_txid()),
+ "Funding tx must not be broadcast after close"
+ );
+ }
+
+ // Also simulate re-storing the funding tx late. still must not broadcast.
+ service_handler
+ .store_funding_transaction(user_channel_id, &client_node_id, funding_tx.clone())
+ .unwrap();
+ {
+ let broadcasted = service_node.inner.tx_broadcaster.txn_broadcasted.lock().unwrap();
+ assert!(
+ broadcasted.iter().all(|tx| tx.compute_txid() != funding_tx.compute_txid()),
+ "Funding tx must not be broadcast after close (late store)"
+ );
+ }
+}
+
+#[test]
+fn htlc_timeout_before_client_claim_results_in_handling_failed() {
+ let chanmon_cfgs = create_chanmon_cfgs(3);
+ let node_cfgs = create_node_cfgs(3, &chanmon_cfgs);
+ let mut service_node_config = test_default_channel_config();
+ service_node_config.accept_intercept_htlcs = true;
+
+ let mut client_node_config = test_default_channel_config();
+ client_node_config.manually_accept_inbound_channels = true;
+ client_node_config.channel_config.accept_underpaying_htlcs = true;
+
+ let node_chanmgrs = create_node_chanmgrs(
+ 3,
+ &node_cfgs,
+ &[Some(service_node_config), Some(client_node_config), None],
+ );
+ let nodes = create_network(3, &node_cfgs, &node_chanmgrs);
+ let (lsps_nodes, promise_secret) = setup_test_lsps2_nodes_with_payer(nodes);
+ let LSPSNodesWithPayer { ref service_node, ref client_node, ref payer_node } = lsps_nodes;
+
+ let payer_node_id = payer_node.node.get_our_node_id();
+ let service_node_id = service_node.inner.node.get_our_node_id();
+ let client_node_id = client_node.inner.node.get_our_node_id();
+
+ let service_handler = service_node.liquidity_manager.lsps2_service_handler().unwrap();
+
+ create_chan_between_nodes_with_value(&payer_node, &service_node.inner, 2_000_000, 100_000);
+
+ let intercept_scid = service_node.node.get_intercept_scid();
+ let user_channel_id = 44u128;
+ let cltv_expiry_delta: u32 = 144;
+ let payment_size_msat = Some(1_000_000);
+ let fee_base_msat: u64 = 10_000;
+
+ execute_lsps2_dance(
+ &lsps_nodes,
+ intercept_scid,
+ user_channel_id,
+ cltv_expiry_delta,
+ promise_secret,
+ payment_size_msat,
+ fee_base_msat,
+ );
+
+ let invoice = create_jit_invoice(
+ &client_node,
+ service_node_id,
+ intercept_scid,
+ cltv_expiry_delta,
+ payment_size_msat,
+ "timeout-before-claim",
+ 3600,
+ )
+ .unwrap();
+
+ payer_node
+ .node
+ .pay_for_bolt11_invoice(
+ &invoice,
+ PaymentId(invoice.payment_hash().to_byte_array()),
+ None,
+ Default::default(),
+ Retry::Attempts(3),
+ )
+ .unwrap();
+
+ check_added_monitors!(payer_node, 1);
+ let events = payer_node.node.get_and_clear_pending_msg_events();
+ let ev = SendEvent::from_event(events[0].clone());
+ service_node.inner.node.handle_update_add_htlc(payer_node_id, &ev.msgs[0]);
+ do_commitment_signed_dance(&service_node.inner, &payer_node, &ev.commitment_msg, false, true);
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ let events = service_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(events.len(), 1);
+ match &events[0] {
+ Event::HTLCIntercepted {
+ intercept_id,
+ requested_next_hop_scid,
+ payment_hash: _,
+ expected_outbound_amount_msat,
+ ..
+ } => {
+ assert_eq!(*requested_next_hop_scid, intercept_scid);
+ service_handler
+ .htlc_intercepted(
+ *requested_next_hop_scid,
+ *intercept_id,
+ *expected_outbound_amount_msat,
+ PaymentHash(invoice.payment_hash().to_byte_array()),
+ )
+ .unwrap();
+ },
+ other => panic!("Expected HTLCIntercepted, got {:?}", other),
+ }
+
+ // Create and mark broadcast safe so the channel is fully ready
+ let expected_outbound_amount_msat = payment_size_msat.unwrap() - fee_base_msat;
+ let (channel_id, _funding_tx) = create_channel_with_manual_broadcast(
+ &service_node_id,
+ &client_node_id,
+ &service_node,
+ &client_node,
+ user_channel_id,
+ &expected_outbound_amount_msat,
+ true,
+ );
+
+ service_handler.channel_ready(user_channel_id, &channel_id, &client_node_id).unwrap();
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ // Forward to client, but do not claim yet
+ let pay_event = {
+ {
+ let mut added_monitors =
+ service_node.inner.chain_monitor.added_monitors.lock().unwrap();
+ assert_eq!(added_monitors.len(), 1);
+ added_monitors.clear();
+ }
+ let mut msg_events = service_node.inner.node.get_and_clear_pending_msg_events();
+ assert_eq!(msg_events.len(), 1);
+ SendEvent::from_event(msg_events.remove(0))
+ };
+
+ client_node.inner.node.handle_update_add_htlc(service_node_id, &pay_event.msgs[0]);
+ do_commitment_signed_dance(
+ &client_node.inner,
+ &service_node.inner,
+ &pay_event.commitment_msg,
+ false,
+ true,
+ );
+ client_node.inner.node.process_pending_htlc_forwards();
+
+ let client_events = client_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(client_events.len(), 1);
+ let preimage = match &client_events[0] {
+ Event::PaymentClaimable { purpose, .. } => purpose.preimage().unwrap(),
+ other => panic!("Expected PaymentClaimable, got {:?}", other),
+ };
+
+ // Advance blocks past CLTV expiry before the client attempts to claim
+ const SOME_EXTRA_BLOCKS: u32 = 3;
+ let client_htlc_cltv_expiry = pay_event.msgs[0].cltv_expiry;
+ let target_height = client_htlc_cltv_expiry.saturating_add(SOME_EXTRA_BLOCKS);
+ let cur_height = service_node.inner.best_block_info().1;
+ let d = target_height - cur_height;
+ connect_blocks(&service_node.inner, d);
+ connect_blocks(&client_node.inner, d);
+ connect_blocks(&payer_node, d);
+
+ service_node.inner.node.process_pending_htlc_forwards();
+ client_node.inner.node.process_pending_htlc_forwards();
+
+ // Service->client channel should close due to HTLC timeout
+ let svc_events = service_node.inner.node.get_and_clear_pending_events();
+ let closed_on_service = svc_events.iter().any(|ev| {
+ matches!(ev, Event::ChannelClosed { reason: ClosureReason::HTLCsTimedOut { .. }, .. })
+ });
+ assert!(closed_on_service, "Expected service->client channel to close due to HTLC timeout");
+
+ // Client tries to claim but should fail since HTLC timed out
+ client_node.inner.node.claim_funds(preimage);
+ let client_events = client_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(client_events.len(), 1);
+ match &client_events[0] {
+ Event::HTLCHandlingFailed { failure_type, .. } => match failure_type {
+ lightning::events::HTLCHandlingFailureType::Receive { payment_hash } => {
+ assert_eq!(*payment_hash, PaymentHash(invoice.payment_hash().to_byte_array()));
+ },
+ _ => panic!("Unexpected failure_type: {:?}", failure_type),
+ },
+ other => panic!("Expected HTLCHandlingFailed after timeout, got {:?}", other),
+ }
+
+ // Payer->service channel should remain open
+ {
+ let chans = service_node.inner.node.list_channels();
+ assert!(chans
+ .iter()
+ .any(|cd| cd.counterparty.node_id == payer_node_id && cd.is_channel_ready));
+ }
+
+ service_node.inner.node.get_and_clear_pending_msg_events();
+ client_node.inner.node.get_and_clear_pending_msg_events();
+ payer_node.node.get_and_clear_pending_msg_events();
+ service_node.inner.chain_monitor.added_monitors.lock().unwrap().clear();
+ client_node.inner.chain_monitor.added_monitors.lock().unwrap().clear();
+ payer_node.chain_monitor.added_monitors.lock().unwrap().clear();
+}
+
+fn claim_and_assert_forwarded_only<'a, 'b, 'c>(
+ payer_node: &lightning::ln::functional_test_utils::Node<'a, 'b, 'c>,
+ service_node: &lightning::ln::functional_test_utils::Node<'a, 'b, 'c>,
+ client_node: &lightning::ln::functional_test_utils::Node<'a, 'b, 'c>,
+ preimage: PaymentPreimage,
+) {
+ let payer_node_id = payer_node.node.get_our_node_id();
+ let service_node_id = service_node.node.get_our_node_id();
+ let client_node_id = client_node.node.get_our_node_id();
+
+ let client_events = client_node.node.get_and_clear_pending_events();
+ assert_eq!(client_events.len(), 1);
+ match &client_events[0] {
+ Event::PaymentClaimed { purpose, .. } => {
+ assert_eq!(purpose.preimage().unwrap(), preimage);
+ },
+ other => panic!("Expected PaymentClaimed, got {:?}", other),
+ }
+
+ let mut client_msg_events = client_node.node.get_and_clear_pending_msg_events();
+ assert_eq!(client_msg_events.len(), 1);
+ let (fulfill_msg, client_commitment_signed) = match client_msg_events.remove(0) {
+ MessageSendEvent::UpdateHTLCs { node_id, updates, .. } => {
+ assert_eq!(node_id, service_node_id);
+ assert_eq!(updates.update_fulfill_htlcs.len(), 1);
+ (updates.update_fulfill_htlcs[0].clone(), updates.commitment_signed.clone())
+ },
+ other => panic!("Unexpected client msg event: {:?}", other),
+ };
+
+ service_node.node.handle_update_fulfill_htlc(client_node_id, fulfill_msg);
+ service_node
+ .node
+ .handle_commitment_signed_batch_test(client_node_id, &client_commitment_signed);
+ service_node.chain_monitor.added_monitors.lock().unwrap().clear();
+
+ let service_msg_events = service_node.node.get_and_clear_pending_msg_events();
+ assert!(
+ service_msg_events.len() >= 2 && service_msg_events.len() <= 3,
+ "Unexpected service msg events len = {}",
+ service_msg_events.len()
+ );
+
+ let mut revoke_and_ack_to_client = None;
+ let mut upstream_updates = None;
+ let mut client_commitment_update = None;
+
+ for ev in service_msg_events {
+ match ev {
+ MessageSendEvent::SendRevokeAndACK { node_id, msg } => {
+ assert_eq!(node_id, client_node_id);
+ revoke_and_ack_to_client = Some(msg);
+ },
+ MessageSendEvent::UpdateHTLCs { node_id, updates, .. } => {
+ if node_id == payer_node_id {
+ assert_eq!(updates.update_fulfill_htlcs.len(), 1, "Expected upstream fulfill");
+ upstream_updates = Some(updates);
+ } else if node_id == client_node_id {
+ assert!(updates.update_fulfill_htlcs.is_empty());
+ client_commitment_update = Some(updates);
+ } else {
+ panic!("Unexpected UpdateHTLCs destination");
+ }
+ },
+ other => panic!("Unexpected service msg event: {:?}", other),
+ }
+ }
+
+ let revoke_and_ack_to_client =
+ revoke_and_ack_to_client.expect("Missing RevokeAndACK to client");
+ let upstream_updates = upstream_updates.expect("Missing upstream fulfill updates");
+
+ client_node.node.handle_revoke_and_ack(service_node_id, &revoke_and_ack_to_client);
+ client_node.chain_monitor.added_monitors.lock().unwrap().clear();
+
+ if let Some(cu) = client_commitment_update {
+ client_node
+ .node
+ .handle_commitment_signed_batch_test(service_node_id, &cu.commitment_signed);
+ client_node.chain_monitor.added_monitors.lock().unwrap().clear();
+
+ let raa_back =
+ get_event_msg!(client_node, MessageSendEvent::SendRevokeAndACK, service_node_id);
+ service_node.node.handle_revoke_and_ack(client_node_id, &raa_back);
+ service_node.chain_monitor.added_monitors.lock().unwrap().clear();
+ }
+
+ payer_node.node.handle_update_fulfill_htlc(
+ service_node_id,
+ upstream_updates.update_fulfill_htlcs[0].clone(),
+ );
+ payer_node
+ .node
+ .handle_commitment_signed_batch_test(service_node_id, &upstream_updates.commitment_signed);
+ payer_node.chain_monitor.added_monitors.lock().unwrap().clear();
+
+ let payer_msg_events = payer_node.node.get_and_clear_pending_msg_events();
+
+ for ev in payer_msg_events {
+ match ev {
+ MessageSendEvent::SendRevokeAndACK { node_id, msg } => {
+ assert_eq!(node_id, service_node_id);
+ service_node.node.handle_revoke_and_ack(payer_node_id, &msg);
+ service_node.chain_monitor.added_monitors.lock().unwrap().clear();
+ },
+ MessageSendEvent::UpdateHTLCs { updates, .. } => {
+ service_node
+ .node
+ .handle_commitment_signed_batch_test(payer_node_id, &updates.commitment_signed);
+ service_node.chain_monitor.added_monitors.lock().unwrap().clear();
+ let mut svc_resp = service_node.node.get_and_clear_pending_msg_events();
+ for resp in svc_resp.drain(..) {
+ match resp {
+ MessageSendEvent::SendRevokeAndACK { msg, .. } => {
+ payer_node.node.handle_revoke_and_ack(service_node_id, &msg);
+ payer_node.chain_monitor.added_monitors.lock().unwrap().clear();
+ },
+ MessageSendEvent::UpdateHTLCs { updates, .. } => {
+ payer_node.node.handle_commitment_signed_batch_test(
+ service_node_id,
+ &updates.commitment_signed,
+ );
+ payer_node.chain_monitor.added_monitors.lock().unwrap().clear();
+ let maybe_final = payer_node.node.get_and_clear_pending_msg_events();
+ for final_ev in maybe_final {
+ if let MessageSendEvent::SendRevokeAndACK { msg, .. } = final_ev {
+ service_node.node.handle_revoke_and_ack(payer_node_id, &msg);
+ service_node
+ .chain_monitor
+ .added_monitors
+ .lock()
+ .unwrap()
+ .clear();
+ }
+ }
+ },
+ _ => {},
+ }
+ }
+ },
+ other => panic!("Unexpected payer msg event: {:?}", other),
+ }
+ }
+}
+
+#[test]
+fn client_trusts_lsp_partial_fee_does_not_trigger_broadcast() {
+ let chanmon_cfgs = create_chanmon_cfgs(3);
+ let node_cfgs = create_node_cfgs(3, &chanmon_cfgs);
+ let mut service_node_config = test_default_channel_config();
+ service_node_config.accept_intercept_htlcs = true;
+
+ let mut client_node_config = test_default_channel_config();
+ client_node_config.manually_accept_inbound_channels = true;
+ client_node_config.channel_config.accept_underpaying_htlcs = true;
+
+ let node_chanmgrs = create_node_chanmgrs(
+ 3,
+ &node_cfgs,
+ &[Some(service_node_config), Some(client_node_config), None],
+ );
+ let nodes = create_network(3, &node_cfgs, &node_chanmgrs);
+ let (lsps_nodes, promise_secret) = setup_test_lsps2_nodes_with_payer(nodes);
+ let LSPSNodesWithPayer { ref service_node, ref client_node, ref payer_node } = lsps_nodes;
+
+ let payer_node_id = payer_node.node.get_our_node_id();
+ let service_node_id = service_node.inner.node.get_our_node_id();
+ let client_node_id = client_node.inner.node.get_our_node_id();
+
+ let service_handler = service_node.liquidity_manager.lsps2_service_handler().unwrap();
+
+ create_chan_between_nodes_with_value(&payer_node, &service_node.inner, 2_000_000, 100_000);
+
+ let intercept_scid = service_node.node.get_intercept_scid();
+ let user_channel_id = 42;
+ let cltv_expiry_delta: u32 = 144;
+ let payment_size_msat = Some(1_000_000);
+
+ let fee_base_msat: u64 = 10_000;
+
+ execute_lsps2_dance(
+ &lsps_nodes,
+ intercept_scid,
+ user_channel_id,
+ cltv_expiry_delta,
+ promise_secret,
+ payment_size_msat,
+ fee_base_msat,
+ );
+
+ let invoice = create_jit_invoice(
+ &client_node,
+ service_node_id,
+ intercept_scid,
+ cltv_expiry_delta,
+ payment_size_msat,
+ "test partial fee",
+ 3600,
+ )
+ .unwrap();
+
+ payer_node
+ .node
+ .pay_for_bolt11_invoice(
+ &invoice,
+ PaymentId(invoice.payment_hash().to_byte_array()),
+ None,
+ Default::default(),
+ Retry::Attempts(3),
+ )
+ .unwrap();
+
+ check_added_monitors!(payer_node, 1);
+ let events = payer_node.node.get_and_clear_pending_msg_events();
+ let ev = SendEvent::from_event(events[0].clone());
+ service_node.inner.node.handle_update_add_htlc(payer_node_id, &ev.msgs[0]);
+ do_commitment_signed_dance(&service_node.inner, &payer_node, &ev.commitment_msg, false, true);
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ let events = service_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(events.len(), 1);
+ let (payment_hash, expected_outbound_amount_msat) = match &events[0] {
+ Event::HTLCIntercepted {
+ intercept_id,
+ requested_next_hop_scid,
+ payment_hash,
+ expected_outbound_amount_msat,
+ ..
+ } => {
+ assert_eq!(*requested_next_hop_scid, intercept_scid);
+ service_handler
+ .htlc_intercepted(
+ *requested_next_hop_scid,
+ *intercept_id,
+ *expected_outbound_amount_msat,
+ *payment_hash,
+ )
+ .unwrap();
+ (*payment_hash, expected_outbound_amount_msat)
+ },
+ other => panic!("Expected HTLCIntercepted, got {:?}", other),
+ };
+
+ match service_node.liquidity_manager.next_event().unwrap() {
+ LiquidityEvent::LSPS2Service(LSPS2ServiceEvent::OpenChannel {
+ their_network_key,
+ amt_to_forward_msat,
+ opening_fee_msat,
+ user_channel_id: u,
+ intercept_scid: sc,
+ }) => {
+ assert_eq!(their_network_key, client_node_id);
+ assert_eq!(u, user_channel_id);
+ assert_eq!(sc, intercept_scid);
+ assert_eq!(opening_fee_msat, fee_base_msat);
+ assert_eq!(amt_to_forward_msat, payment_size_msat.unwrap() - fee_base_msat);
+ },
+ other => panic!("Unexpected event: {:?}", other),
+ };
+
+ assert!(service_handler
+ .channel_needs_manual_broadcast(user_channel_id, &client_node_id)
+ .unwrap());
+
+ let (channel_id, _) = create_channel_with_manual_broadcast(
+ &service_node_id,
+ &client_node_id,
+ &service_node,
+ &client_node,
+ user_channel_id,
+ expected_outbound_amount_msat,
+ true,
+ );
+
+ service_handler.channel_ready(user_channel_id, &channel_id, &client_node_id).unwrap();
+
+ service_node.inner.node.process_pending_htlc_forwards();
+
+ let pay_event = {
+ {
+ let mut added_monitors =
+ service_node.inner.chain_monitor.added_monitors.lock().unwrap();
+ assert_eq!(added_monitors.len(), 1);
+ added_monitors.clear();
+ }
+ let mut msg_events = service_node.inner.node.get_and_clear_pending_msg_events();
+ assert_eq!(msg_events.len(), 1);
+ SendEvent::from_event(msg_events.remove(0))
+ };
+
+ client_node.inner.node.handle_update_add_htlc(service_node_id, &pay_event.msgs[0]);
+ do_commitment_signed_dance(
+ &client_node.inner,
+ &service_node.inner,
+ &pay_event.commitment_msg,
+ false,
+ true,
+ );
+ client_node.inner.node.process_pending_htlc_forwards();
+
+ let client_events = client_node.inner.node.get_and_clear_pending_events();
+ assert_eq!(client_events.len(), 1);
+ match &client_events[0] {
+ Event::PaymentClaimable { payment_hash: ph, .. } => assert_eq!(*ph, payment_hash),
+ other => panic!("Expected PaymentClaimable, got {:?}", other),
+ };
+
+ assert!(service_node.liquidity_manager.get_and_clear_pending_events().is_empty());
+
+ let partial_skim_msat = fee_base_msat - 1; // less than promised fee
+ service_handler.payment_forwarded(channel_id, partial_skim_msat).unwrap();
+
+ let broadcasted = service_node.inner.tx_broadcaster.txn_broadcasted.lock().unwrap();
+ assert!(broadcasted.is_empty(), "There should be no broadcasted txs yet");
+ drop(broadcasted);
+
+ // before mining blocks, service node should have 2 channels
+ {
+ let chans = service_node.inner.node.list_channels();
+ assert_eq!(chans.len(), 2);
+ assert!(chans.iter().any(|cd| cd.counterparty.node_id == payer_node_id));
+ assert!(chans.iter().any(|cd| cd.counterparty.node_id == client_node_id));
+ }
+
+ const SOME_EXTRA_BLOCKS: u32 = 3;
+ let client_htlc_cltv_expiry = pay_event.msgs[0].cltv_expiry;
+ let target_height = client_htlc_cltv_expiry.saturating_add(SOME_EXTRA_BLOCKS);
+ let cur_height = service_node.inner.best_block_info().1;
+ let d = target_height - cur_height;
+ connect_blocks(&service_node.inner, d);
+ connect_blocks(&client_node.inner, d);
+ connect_blocks(&payer_node, d);
+
+ service_node.inner.node.process_pending_htlc_forwards();
+ client_node.inner.node.process_pending_htlc_forwards();
+
+ let svc_events = service_node.inner.node.get_and_clear_pending_events();
+ let _ = client_node.inner.node.get_and_clear_pending_events();
+ let closed_on_service = svc_events.iter().any(|ev| {
+ matches!(ev, Event::ChannelClosed { reason: ClosureReason::HTLCsTimedOut { .. }, .. })
+ });
+ assert!(
+ closed_on_service,
+ "Expected service->client channel to be force-closed due to HTLC timeout. svc_events = {:?}",
+ svc_events
+ );
+
+ // now check the service->payer channel
+ {
+ let chans = service_node.inner.node.list_channels();
+ assert!(chans.len() == 1);
+ assert!(
+ chans.iter().any(|cd| cd.counterparty.node_id == payer_node_id && cd.is_channel_ready),
+ "Expected payer->service channel to remain open. channels: {:?}",
+ chans
+ );
+ }
+
+ service_node.inner.node.get_and_clear_pending_msg_events();
+ client_node.inner.node.get_and_clear_pending_msg_events();
+ payer_node.node.get_and_clear_pending_msg_events();
+ service_node.inner.chain_monitor.added_monitors.lock().unwrap().clear();
+ client_node.inner.chain_monitor.added_monitors.lock().unwrap().clear();
+ payer_node.chain_monitor.added_monitors.lock().unwrap().clear();
+}
diff --git a/lightning-liquidity/tests/lsps5_integration_tests.rs b/lightning-liquidity/tests/lsps5_integration_tests.rs
index a3d2ecf..41af2e8 100644
--- a/lightning-liquidity/tests/lsps5_integration_tests.rs
+++ b/lightning-liquidity/tests/lsps5_integration_tests.rs
@@ -1610,6 +1610,7 @@ fn lsps5_service_handler_persistence_across_restarts() {
None::<Arc<dyn Filter + Send + Sync>>,
Some(chain_params),
service_kv_store,
+ nodes_restart[0].tx_broadcaster,
Some(service_config),
None,
Arc::clone(&time_provider),
diff --git a/lightning/src/ln/channelmanager.rs b/lightning/src/ln/channelmanager.rs
index 4365407..1d87ecc 100644
--- a/lightning/src/ln/channelmanager.rs
+++ b/lightning/src/ln/channelmanager.rs
@@ -1098,6 +1098,11 @@ enum FundingType {
///
/// This is the normal flow.
Checked(Transaction),
+ /// This variant is useful when we want LDK to validate the funding transaction and
+ /// broadcast it manually.
+ ///
+ /// Used in LSPS2 on a client_trusts_lsp model
+ CheckedManualBroadcast(Transaction),
/// This variant is useful when we want to loosen the validation checks and allow to
/// manually broadcast the funding transaction, leaving the responsibility to the caller.
///
@@ -1112,6 +1117,7 @@ impl FundingType {
fn txid(&self) -> Txid {
match self {
FundingType::Checked(tx) => tx.compute_txid(),
+ FundingType::CheckedManualBroadcast(tx) => tx.compute_txid(),
FundingType::Unchecked(outp) => outp.txid,
}
}
@@ -1119,6 +1125,7 @@ impl FundingType {
fn transaction_or_dummy(&self) -> Transaction {
match self {
FundingType::Checked(tx) => tx.clone(),
+ FundingType::CheckedManualBroadcast(tx) => tx.clone(),
FundingType::Unchecked(_) => Transaction {
version: bitcoin::transaction::Version::TWO,
lock_time: bitcoin::absolute::LockTime::ZERO,
@@ -1131,6 +1138,7 @@ impl FundingType {
fn is_manual_broadcast(&self) -> bool {
match self {
FundingType::Checked(_) => false,
+ FundingType::CheckedManualBroadcast(_) => true,
FundingType::Unchecked(_) => true,
}
}
@@ -6040,6 +6048,43 @@ where
self.batch_funding_transaction_generated_intern(temporary_chans, funding_type)
}
+ /// Call this upon creation of a funding transaction for the given channel.
+ ///
+ /// This method executes the same checks as [`ChannelManager::funding_transaction_generated`],
+ /// but it does not automatically broadcast the funding transaction.
+ ///
+ /// Call this in response to a [`Event::FundingGenerationReady`] event, only in a context where you want to manually
+ /// control the broadcast of the funding transaction.
+ ///
+ /// The associated [`ChannelMonitor`] likewise avoids broadcasting holder commitment or CPFP
+ /// transactions until the funding has been observed on chain. This
+ /// prevents attempting to broadcast unconfirmable commitment transactions before the channel's
+ /// funding exists in a block.
+ ///
+ /// If HTLCs would otherwise approach timeout while the funding transaction has not yet appeared
+ /// on chain, the monitor avoids broadcasting force-close transactions in manual-broadcast
+ /// mode until the funding is seen. It may still close the channel off-chain (emitting a
+ /// `ChannelClosed` event) to avoid accepting further updates. Ensure your application either
+ /// broadcasts the funding transaction in a timely manner or avoids forwarding HTLCs that could
+ /// approach timeout during this interim state.
+ ///
+ /// See also [`ChannelMonitor::broadcast_latest_holder_commitment_txn`]. For channels using
+ /// manual-broadcast, calling that method has no effect until the funding has been observed
+ /// on-chain.
+ ///
+ /// [`ChannelManager::funding_transaction_generated`]: crate::ln::channelmanager::ChannelManager::funding_transaction_generated
+ /// [`Event::FundingGenerationReady`]: crate::events::Event::FundingGenerationReady
+ pub fn funding_transaction_generated_manual_broadcast(
+ &self, temporary_channel_id: ChannelId, counterparty_node_id: PublicKey,
+ funding_transaction: Transaction,
+ ) -> Result<(), APIError> {
+ let _persistence_guard = PersistenceNotifierGuard::notify_on_drop(self);
+ self.batch_funding_transaction_generated_intern(
+ &[(&temporary_channel_id, &counterparty_node_id)],
+ FundingType::CheckedManualBroadcast(funding_transaction),
+ )
+ }
+
/// Call this upon creation of a batch funding transaction for the given channels.
///
/// Return values are identical to [`Self::funding_transaction_generated`], respective to
@@ -6061,7 +6106,9 @@ where
#[rustfmt::skip]
fn batch_funding_transaction_generated_intern(&self, temporary_channels: &[(&ChannelId, &PublicKey)], funding: FundingType) -> Result<(), APIError> {
let mut result = Ok(());
- if let FundingType::Checked(funding_transaction) = &funding {
+ if let FundingType::Checked(funding_transaction) |
+ FundingType::CheckedManualBroadcast(funding_transaction) = &funding
+ {
if !funding_transaction.is_coinbase() {
for inp in funding_transaction.input.iter() {
if inp.witness.is_empty() {
@@ -6121,7 +6168,7 @@ where
let mut output_index = None;
let expected_spk = chan.funding.get_funding_redeemscript().to_p2wsh();
let outpoint = match &funding {
- FundingType::Checked(tx) => {
+ FundingType::Checked(tx) | FundingType::CheckedManualBroadcast(tx) => {
for (idx, outp) in tx.output.iter().enumerate() {
if outp.script_pubkey == expected_spk && outp.value.to_sat() == chan.funding.get_value_satoshis() {
if output_index.is_some() {
Why this scored 34/100
Community notes
Notes can correct, qualify, or add evidence to the AI analysis. Every note shown here has been validated by a human moderator.
The AI analysis stands alone for now. Submit a note if you can add evidence or important context.