1use std::collections::{BTreeMap, HashMap};
2use std::path::Path;
3use std::str::FromStr;
4use std::sync::Arc;
5use std::time::{Duration, UNIX_EPOCH};
6
7use async_trait::async_trait;
8use bitcoin::hashes::{Hash, sha256};
9use bitcoin::{FeeRate, Network, OutPoint};
10use fedimint_bip39::Mnemonic;
11use fedimint_core::envs::is_running_in_test_env;
12use fedimint_core::task::{TaskGroup, TaskHandle, block_in_place};
13use fedimint_core::util::{FmtCompact, SafeUrl};
14use fedimint_core::{Amount, BitcoinAmountOrAll, crit};
15use fedimint_gateway_common::{
16 ChainSource, ConnectPeerRequest, GetInvoiceRequest, GetInvoiceResponse,
17 ListTransactionsResponse, NodeAddress, SetChannelFeesRequest,
18};
19use fedimint_ln_common::contracts::Preimage;
20use fedimint_logging::{LOG_LIGHTNING, LOG_LIGHTNING_LDK};
21use ldk_node::config::ChannelConfig;
22use ldk_node::lightning::ln::msgs::SocketAddress;
23use ldk_node::lightning::routing::gossip::{NodeAlias, NodeId};
24use ldk_node::logger::{LogLevel, LogRecord, LogWriter};
25use ldk_node::payment::{
26 PaymentDetails, PaymentDirection, PaymentKind, PaymentStatus, SendingParameters,
27};
28use lightning::ln::channelmanager::PaymentId;
29use lightning::offers::offer::{Offer, OfferId};
30use lightning::types::payment::{PaymentHash, PaymentPreimage};
31use lightning_invoice::{Bolt11Invoice, Bolt11InvoiceDescription, Description};
32use thiserror::Error;
33use tokio::sync::mpsc::Sender;
34use tokio::sync::{RwLock, oneshot};
35use tokio_stream::wrappers::ReceiverStream;
36use tracing::{debug, error, info, trace, warn};
37
38use super::{ChannelInfo, ILnRpcClient, LightningRpcError, ListChannelsResponse, RouteHtlcStream};
39use crate::{
40 CloseChannelsWithPeerRequest, CloseChannelsWithPeerResponse, CreateInvoiceRequest,
41 CreateInvoiceResponse, GetBalancesResponse, GetLnOnchainAddressResponse, GetNodeInfoResponse,
42 GetRouteHintsResponse, InterceptPaymentRequest, InterceptPaymentResponse, InvoiceDescription,
43 NO_INCOMING_CIRCUIT, OpenChannelRequest, OpenChannelResponse, PayInvoiceResponse,
44 PaymentAction, SendOnchainRequest, SendOnchainResponse,
45};
46
47struct LdkTracingLogger {
55 in_test_env: bool,
59}
60
61impl LdkTracingLogger {
62 fn downgraded_level(&self, record: &LogRecord<'_>) -> LogLevel {
72 if self.in_test_env
73 && record.level == LogLevel::Error
74 && record.module_path == "ldk_node::chain"
75 && format!("{}", record.args).contains("Failed to retrieve fee rate estimates")
76 {
77 LogLevel::Debug
78 } else {
79 record.level
80 }
81 }
82}
83
84impl LogWriter for LdkTracingLogger {
85 fn log(&self, record: LogRecord<'_>) {
86 match self.downgraded_level(&record) {
88 LogLevel::Gossip | LogLevel::Trace => trace!(
89 target: LOG_LIGHTNING_LDK,
90 ldk_module = record.module_path, line = record.line, "{}", record.args,
91 ),
92 LogLevel::Debug => debug!(
93 target: LOG_LIGHTNING_LDK,
94 ldk_module = record.module_path, line = record.line, "{}", record.args,
95 ),
96 LogLevel::Info => info!(
97 target: LOG_LIGHTNING_LDK,
98 ldk_module = record.module_path, line = record.line, "{}", record.args,
99 ),
100 LogLevel::Warn => warn!(
101 target: LOG_LIGHTNING_LDK,
102 ldk_module = record.module_path, line = record.line, "{}", record.args,
103 ),
104 LogLevel::Error => error!(
105 target: LOG_LIGHTNING_LDK,
106 ldk_module = record.module_path, line = record.line, "{}", record.args,
107 ),
108 }
109 }
110}
111
112pub struct GatewayLdkClient {
113 node: Arc<ldk_node::Node>,
115
116 task_group: TaskGroup,
117
118 htlc_stream_receiver_or: Option<tokio::sync::mpsc::Receiver<InterceptPaymentRequest>>,
121
122 outbound_lightning_payment_lock_pool: lockable::LockPool<PaymentId>,
126
127 outbound_offer_lock_pool: lockable::LockPool<LdkOfferId>,
132
133 pending_channels: Arc<RwLock<BTreeMap<UserChannelId, PendingChannelSender>>>,
138
139 pending_payments: Arc<RwLock<HashMap<PaymentId, oneshot::Sender<()>>>>,
145}
146
147type PendingChannelSender = oneshot::Sender<Result<OutPoint, String>>;
150
151impl std::fmt::Debug for GatewayLdkClient {
152 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
153 f.debug_struct("GatewayLdkClient").finish_non_exhaustive()
154 }
155}
156
157impl GatewayLdkClient {
158 pub fn new(
170 data_dir: &Path,
171 chain_source: ChainSource,
172 network: Network,
173 lightning_port: u16,
174 alias: String,
175 mnemonic: Mnemonic,
176 runtime: Arc<tokio::runtime::Runtime>,
177 ) -> Result<Self, LdkClientInitError> {
178 let mut bytes = [0u8; 32];
179 let alias = if alias.is_empty() {
180 "LDK Gateway".to_string()
181 } else {
182 alias
183 };
184 let alias_bytes = alias.as_bytes();
185 let truncated = &alias_bytes[..alias_bytes.len().min(32)];
186 bytes[..truncated.len()].copy_from_slice(truncated);
187 let node_alias = Some(NodeAlias(bytes));
188
189 let mut node_builder = ldk_node::Builder::from_config(ldk_node::config::Config {
190 network,
191 listening_addresses: Some(vec![SocketAddress::TcpIpV4 {
192 addr: [0, 0, 0, 0],
193 port: lightning_port,
194 }]),
195 node_alias,
196 ..Default::default()
197 });
198
199 node_builder.set_custom_logger(Arc::new(LdkTracingLogger {
203 in_test_env: is_running_in_test_env(),
204 }));
205
206 node_builder.set_entropy_bip39_mnemonic(mnemonic, None);
207
208 match chain_source.clone() {
209 ChainSource::Bitcoind {
210 username,
211 password,
212 server_url,
213 } => {
214 node_builder.set_chain_source_bitcoind_rpc(
215 server_url
216 .host_str()
217 .expect("Could not retrieve host from bitcoind RPC url")
218 .to_string(),
219 server_url
220 .port()
221 .expect("Could not retrieve port from bitcoind RPC url"),
222 username,
223 password,
224 );
225 }
226 ChainSource::Esplora { server_url } => {
227 node_builder.set_chain_source_esplora(get_esplora_url(server_url)?, None);
228 }
229 };
230 let Some(data_dir_str) = data_dir.to_str() else {
231 return Err(LdkClientInitError::InvalidDataDir);
232 };
233 node_builder.set_storage_dir_path(data_dir_str.to_string());
234
235 info!(chain_source = %chain_source, data_dir = %data_dir_str, alias = %alias, "Starting LDK Node...");
236 let node = Arc::new(node_builder.build().map_err(LdkClientInitError::Build)?);
237 node.start_with_runtime(runtime).map_err(|err| {
238 crit!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Failed to start LDK Node");
239 LdkClientInitError::Start(err)
240 })?;
241
242 let (htlc_stream_sender, htlc_stream_receiver) = tokio::sync::mpsc::channel(1024);
243 let task_group = TaskGroup::new();
244
245 let node_clone = node.clone();
246 let pending_channels = Arc::new(RwLock::new(BTreeMap::new()));
247 let pending_channels_clone = pending_channels.clone();
248 let pending_payments = Arc::new(RwLock::new(HashMap::new()));
249 let pending_payments_clone = pending_payments.clone();
250 task_group.spawn("ldk lightning node event handler", |handle| async move {
251 loop {
252 Self::handle_next_event(
253 &node_clone,
254 &htlc_stream_sender,
255 &handle,
256 pending_channels_clone.clone(),
257 pending_payments_clone.clone(),
258 )
259 .await;
260 }
261 });
262
263 info!("Successfully started LDK Gateway");
264 Ok(GatewayLdkClient {
265 node,
266 task_group,
267 htlc_stream_receiver_or: Some(htlc_stream_receiver),
268 outbound_lightning_payment_lock_pool: lockable::LockPool::new(),
269 outbound_offer_lock_pool: lockable::LockPool::new(),
270 pending_channels,
271 pending_payments,
272 })
273 }
274
275 async fn handle_next_event(
276 node: &ldk_node::Node,
277 htlc_stream_sender: &Sender<InterceptPaymentRequest>,
278 handle: &TaskHandle,
279 pending_channels: Arc<RwLock<BTreeMap<UserChannelId, PendingChannelSender>>>,
280 pending_payments: Arc<RwLock<HashMap<PaymentId, oneshot::Sender<()>>>>,
281 ) {
282 let event = tokio::select! {
286 event = node.next_event_async() => {
287 event
288 }
289 () = handle.make_shutdown_rx() => {
290 return;
291 }
292 };
293
294 match event {
295 ldk_node::Event::PaymentClaimable {
296 payment_id: _,
297 payment_hash,
298 claimable_amount_msat,
299 claim_deadline,
300 custom_records: _,
301 } => {
302 if let Err(err) = htlc_stream_sender
303 .send(InterceptPaymentRequest {
304 payment_hash: Hash::from_slice(&payment_hash.0)
305 .expect("Failed to create Hash"),
306 amount_msat: claimable_amount_msat,
309 incoming_amount_msat: claimable_amount_msat,
310 expiry: claim_deadline.unwrap_or_default(),
311 short_channel_id: None,
312 incoming_chan_id: NO_INCOMING_CIRCUIT.0,
315 htlc_id: NO_INCOMING_CIRCUIT.1,
316 })
317 .await
318 {
319 warn!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Failed send InterceptHtlcRequest to stream");
320 }
321 }
322 ldk_node::Event::ChannelPending {
323 channel_id,
324 user_channel_id,
325 former_temporary_channel_id: _,
326 counterparty_node_id: _,
327 funding_txo,
328 } => {
329 info!(target: LOG_LIGHTNING, %channel_id, "LDK Channel is pending");
330 let mut channels = pending_channels.write().await;
331 if let Some(sender) = channels.remove(&UserChannelId(user_channel_id)) {
332 let _ = sender.send(Ok(funding_txo));
333 } else {
334 debug!(
335 ?user_channel_id,
336 "No channel pending channel open for user channel id"
337 );
338 }
339 }
340 ldk_node::Event::ChannelClosed {
341 channel_id,
342 user_channel_id,
343 counterparty_node_id: _,
344 reason,
345 } => {
346 info!(target: LOG_LIGHTNING, %channel_id, "LDK Channel is closed");
347 let mut channels = pending_channels.write().await;
348 if let Some(sender) = channels.remove(&UserChannelId(user_channel_id)) {
349 let reason = if let Some(reason) = reason {
350 reason.to_string()
351 } else {
352 "Channel has been closed".to_string()
353 };
354 let _ = sender.send(Err(reason));
355 } else {
356 debug!(
357 ?user_channel_id,
358 "No channel pending channel open for user channel id"
359 );
360 }
361 }
362 ldk_node::Event::PaymentSuccessful {
363 payment_id: Some(payment_id),
364 ..
365 }
366 | ldk_node::Event::PaymentFailed {
367 payment_id: Some(payment_id),
368 ..
369 } => {
370 Self::wake_pending_payment(&pending_payments, payment_id).await;
371 }
372 _ => {}
373 }
374
375 if let Err(err) = node.event_handled() {
380 warn!(err = %err.fmt_compact(), "LDK could not mark event handled");
381 }
382 }
383
384 async fn wake_pending_payment(
388 pending_payments: &Arc<RwLock<HashMap<PaymentId, oneshot::Sender<()>>>>,
389 payment_id: PaymentId,
390 ) -> PendingPaymentWakeup {
391 let Some(sender) = pending_payments.write().await.remove(&payment_id) else {
392 return PendingPaymentWakeup::NoWaiter;
393 };
394
395 if sender.send(()).is_ok() {
396 PendingPaymentWakeup::Woken
397 } else {
398 PendingPaymentWakeup::ReceiverDropped
399 }
400 }
401
402 fn outbound_payment(&self, payment_id: PaymentId) -> Option<PaymentDetails> {
412 self.node
413 .payment(&payment_id)
414 .filter(|details| details.direction == PaymentDirection::Outbound)
415 }
416
417 fn ldk_payment_result(
422 &self,
423 payment_id: PaymentId,
424 ) -> Option<Result<PayInvoiceResponse, LightningRpcError>> {
425 let payment_details = self.outbound_payment(payment_id)?;
426 match payment_details.status {
427 PaymentStatus::Pending => None,
428 PaymentStatus::Succeeded => {
429 if let PaymentKind::Bolt11 {
430 preimage: Some(preimage),
431 ..
432 } = payment_details.kind
433 {
434 Some(Ok(PayInvoiceResponse {
435 preimage: Preimage(preimage.0),
436 }))
437 } else {
438 Some(Err(LightningRpcError::FailedPayment {
439 failure_reason: "LDK payment succeeded without preimage".to_string(),
440 }))
441 }
442 }
443 PaymentStatus::Failed => Some(Err(LightningRpcError::FailedPayment {
444 failure_reason: "LDK payment failed".to_string(),
445 })),
446 }
447 }
448}
449
450#[derive(Debug, Error)]
455#[non_exhaustive]
456pub enum LdkClientInitError {
457 #[error("Missing esplora host")]
459 MissingEsploraHost,
460
461 #[error("Invalid data dir path")]
464 InvalidDataDir,
465
466 #[error("The LDK node could not be built")]
469 Build(#[source] ldk_node::BuildError),
470
471 #[error("The LDK node could not be started")]
473 Start(#[source] ldk_node::NodeError),
474}
475
476#[derive(Debug, Clone, Copy, PartialEq, Eq)]
478enum InboundRegistrationRefusal {
479 OutboundInFlight,
482 OutboundRecorded,
484}
485
486impl std::fmt::Display for InboundRegistrationRefusal {
487 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
488 match self {
489 Self::OutboundInFlight => {
490 write!(f, "an outbound payment for this hash is in flight")
491 }
492 Self::OutboundRecorded => {
493 write!(f, "the node holds an outbound payment for this hash")
494 }
495 }
496 }
497}
498
499fn check_inbound_registration(
516 outbound_lock_acquired: bool,
517 existing_direction: Option<PaymentDirection>,
518) -> Result<(), InboundRegistrationRefusal> {
519 if !outbound_lock_acquired {
520 return Err(InboundRegistrationRefusal::OutboundInFlight);
521 }
522
523 match existing_direction {
524 Some(PaymentDirection::Outbound) => Err(InboundRegistrationRefusal::OutboundRecorded),
525 Some(PaymentDirection::Inbound) | None => Ok(()),
526 }
527}
528
529fn htlc_completion_error(err: &ldk_node::NodeError, payment_hash: &str) -> LightningRpcError {
541 match err {
542 ldk_node::NodeError::InvalidPaymentHash
543 | ldk_node::NodeError::InvalidPaymentPreimage
544 | ldk_node::NodeError::InvalidAmount => LightningRpcError::HtlcCompletionRejected {
545 failure_reason: format!(
546 "LDK rejected completion of payment with hash {payment_hash}: {err}"
547 ),
548 },
549 _ => LightningRpcError::FailedToCompleteHtlc {
550 failure_reason: format!(
551 "Failed to complete LDK payment with hash {payment_hash}: {err}"
552 ),
553 },
554 }
555}
556
557impl Drop for GatewayLdkClient {
558 fn drop(&mut self) {
559 self.task_group.shutdown();
560
561 info!(target: LOG_LIGHTNING, "Stopping LDK Node...");
562 match self.node.stop() {
563 Err(err) => {
564 warn!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Failed to stop LDK Node");
565 }
566 _ => {
567 info!(target: LOG_LIGHTNING, "LDK Node stopped.");
568 }
569 }
570 }
571}
572
573#[async_trait]
574impl ILnRpcClient for GatewayLdkClient {
575 async fn info(&self) -> Result<GetNodeInfoResponse, LightningRpcError> {
576 let node_status = self.node.status();
577 let ldk_block_height = node_status.current_best_block.height;
578 let onchain_sync = node_status.latest_onchain_wallet_sync_timestamp;
579 let lightning_sync = node_status.latest_lightning_wallet_sync_timestamp;
580 let is_running = node_status.is_running;
581 debug!(target: LOG_LIGHTNING, ?onchain_sync, ?lightning_sync, ?is_running, "LDK Sync Status");
582
583 Ok(GetNodeInfoResponse {
584 pub_key: self.node.node_id(),
585 alias: match self.node.node_alias() {
586 Some(alias) => alias.to_string(),
587 None => format!("LDK Fedimint Gateway Node {}", self.node.node_id()),
588 },
589 network: self.node.config().network.to_string(),
590 block_height: ldk_block_height,
591 synced_to_chain: lightning_sync.is_some(),
594 })
595 }
596
597 async fn routehints(
598 &self,
599 _num_route_hints: usize,
600 ) -> Result<GetRouteHintsResponse, LightningRpcError> {
601 Ok(GetRouteHintsResponse {
607 route_hints: vec![],
608 })
609 }
610
611 async fn pay(
612 &self,
613 invoice: Bolt11Invoice,
614 max_delay: u64,
615 max_fee: Amount,
616 ) -> Result<PayInvoiceResponse, LightningRpcError> {
617 let payment_id = PaymentId(*invoice.payment_hash().as_byte_array());
618
619 let _payment_lock_guard = self
625 .outbound_lightning_payment_lock_pool
626 .async_lock(payment_id)
627 .await;
628
629 let (payment_sender, payment_receiver) = oneshot::channel();
633 self.pending_payments
634 .write()
635 .await
636 .insert(payment_id, payment_sender);
637
638 if self.outbound_payment(payment_id).is_none() {
649 let sent_payment_id = match self.node.bolt11_payment().send(
650 &invoice,
651 Some(SendingParameters {
652 max_total_routing_fee_msat: Some(Some(max_fee.msats)),
653 max_total_cltv_expiry_delta: Some(max_delay as u32),
654 max_path_count: None,
655 max_channel_saturation_power_of_half: None,
656 }),
657 ) {
658 Ok(sent_payment_id) => sent_payment_id,
659 Err(err) => {
660 self.pending_payments.write().await.remove(&payment_id);
661 return Err(LightningRpcError::FailedPayment {
664 failure_reason: format!("LDK payment failed to initialize: {err:?}"),
665 });
666 }
667 };
668 assert_eq!(sent_payment_id, payment_id);
669 }
670
671 if let Some(result) = self.ldk_payment_result(payment_id) {
675 self.pending_payments.write().await.remove(&payment_id);
676 return result;
677 }
678
679 let _ = payment_receiver.await;
684
685 self.pending_payments.write().await.remove(&payment_id);
686 self.ldk_payment_result(payment_id).unwrap_or_else(|| {
687 Err(LightningRpcError::FailedPayment {
688 failure_reason: "LDK payment event fired without terminal payment status"
689 .to_string(),
690 })
691 })
692 }
693
694 async fn outbound_payment_exists(
695 &self,
696 payment_hash: sha256::Hash,
697 ) -> Result<bool, LightningRpcError> {
698 Ok(self
699 .outbound_payment(PaymentId(payment_hash.to_byte_array()))
700 .is_some())
701 }
702
703 async fn route_htlcs<'a>(
704 mut self: Box<Self>,
705 _task_group: &TaskGroup,
706 ) -> Result<(RouteHtlcStream<'a>, Arc<dyn ILnRpcClient>), LightningRpcError> {
707 let route_htlc_stream = match self.htlc_stream_receiver_or.take() {
708 Some(stream) => Ok(Box::pin(ReceiverStream::new(stream))),
709 None => Err(LightningRpcError::FailedToRouteHtlcs {
710 failure_reason:
711 "Stream does not exist. Likely was already taken by calling `route_htlcs()`."
712 .to_string(),
713 }),
714 }?;
715
716 Ok((route_htlc_stream, Arc::new(*self)))
717 }
718
719 async fn complete_htlc(&self, htlc: InterceptPaymentResponse) -> Result<(), LightningRpcError> {
720 let InterceptPaymentResponse {
721 action,
722 payment_hash,
723 incoming_chan_id: _,
724 htlc_id: _,
725 } = htlc;
726
727 let ph = PaymentHash(*payment_hash.clone().as_byte_array());
728
729 let claimable_amount_msat = 999_999_999_999_999;
735
736 let ph_hex_str = hex::encode(payment_hash);
737
738 if let PaymentAction::Settle(preimage) = action {
739 self.node
740 .bolt11_payment()
741 .claim_for_hash(ph, claimable_amount_msat, PaymentPreimage(preimage.0))
742 .map_err(|err| htlc_completion_error(&err, &ph_hex_str))?;
743 } else {
744 warn!(target: LOG_LIGHTNING, payment_hash = %ph_hex_str, "Unwinding payment because the action was not `Settle`");
745 self.node
746 .bolt11_payment()
747 .fail_for_hash(ph)
748 .map_err(|err| htlc_completion_error(&err, &ph_hex_str))?;
749 }
750
751 return Ok(());
752 }
753
754 async fn create_invoice(
755 &self,
756 create_invoice_request: CreateInvoiceRequest,
757 ) -> Result<CreateInvoiceResponse, LightningRpcError> {
758 let payment_hash_or = if let Some(payment_hash) = create_invoice_request.payment_hash {
759 let ph = PaymentHash(*payment_hash.as_byte_array());
760 Some(ph)
761 } else {
762 None
763 };
764
765 let description = match create_invoice_request.description {
766 Some(InvoiceDescription::Direct(desc)) => {
767 Bolt11InvoiceDescription::Direct(Description::new(desc).map_err(|_| {
768 LightningRpcError::FailedToGetInvoice {
769 failure_reason: "Invalid description".to_string(),
770 }
771 })?)
772 }
773 Some(InvoiceDescription::Hash(hash)) => {
774 Bolt11InvoiceDescription::Hash(lightning_invoice::Sha256(hash))
775 }
776 None => Bolt11InvoiceDescription::Direct(Description::empty()),
777 };
778
779 let invoice = match payment_hash_or {
780 Some(payment_hash) => {
781 let payment_id = PaymentId(payment_hash.0);
788 let outbound_lock_guard = self
789 .outbound_lightning_payment_lock_pool
790 .try_lock(payment_id);
791 let existing_direction = self
792 .node
793 .payment(&payment_id)
794 .map(|details| details.direction);
795 if let Err(refusal) =
796 check_inbound_registration(outbound_lock_guard.is_some(), existing_direction)
797 {
798 warn!(
799 target: LOG_LIGHTNING,
800 payment_hash = %hex::encode(payment_hash.0),
801 %refusal,
802 "Refusing to register an invoice for a payment hash we pay outbound"
803 );
804 return Err(LightningRpcError::FailedToGetInvoice {
805 failure_reason: format!(
806 "Payment hash cannot be registered for an invoice: {refusal}"
807 ),
808 });
809 }
810
811 self.node.bolt11_payment().receive_for_hash(
812 create_invoice_request.amount_msat,
813 &description,
814 create_invoice_request.expiry_secs,
815 payment_hash,
816 )
817 }
818 None => self.node.bolt11_payment().receive(
819 create_invoice_request.amount_msat,
820 &description,
821 create_invoice_request.expiry_secs,
822 ),
823 }
824 .map_err(|e| LightningRpcError::FailedToGetInvoice {
825 failure_reason: e.to_string(),
826 })?;
827
828 Ok(CreateInvoiceResponse {
829 invoice: invoice.to_string(),
830 })
831 }
832
833 async fn get_ln_onchain_address(
834 &self,
835 ) -> Result<GetLnOnchainAddressResponse, LightningRpcError> {
836 self.node
837 .onchain_payment()
838 .new_address()
839 .map(|address| GetLnOnchainAddressResponse {
840 address: address.to_string(),
841 })
842 .map_err(|e| LightningRpcError::FailedToGetLnOnchainAddress {
843 failure_reason: e.to_string(),
844 })
845 }
846
847 async fn send_onchain(
848 &self,
849 SendOnchainRequest {
850 address,
851 amount,
852 fee_rate_sats_per_vbyte,
853 }: SendOnchainRequest,
854 ) -> Result<SendOnchainResponse, LightningRpcError> {
855 let onchain = self.node.onchain_payment();
856
857 let retain_reserves = false;
858 let txid = match amount {
859 BitcoinAmountOrAll::All => onchain.send_all_to_address(
860 &address.assume_checked(),
861 retain_reserves,
862 FeeRate::from_sat_per_vb(fee_rate_sats_per_vbyte),
863 ),
864 BitcoinAmountOrAll::Amount(amount_sats) => onchain.send_to_address(
865 &address.assume_checked(),
866 amount_sats.to_sat(),
867 FeeRate::from_sat_per_vb(fee_rate_sats_per_vbyte),
868 ),
869 }
870 .map_err(|e| LightningRpcError::FailedToWithdrawOnchain {
871 failure_reason: e.to_string(),
872 })?;
873
874 Ok(SendOnchainResponse {
875 txid: txid.to_string(),
876 })
877 }
878
879 async fn open_channel(
880 &self,
881 OpenChannelRequest {
882 pubkey,
883 host,
884 channel_size_sats,
885 push_amount_sats,
886 fee_rate_sats_per_vbyte,
887 base_fee_msat,
888 parts_per_million,
889 }: OpenChannelRequest,
890 ) -> Result<OpenChannelResponse, LightningRpcError> {
891 let push_amount_msats_or = if push_amount_sats == 0 {
892 None
893 } else {
894 Some(push_amount_sats * 1000)
895 };
896
897 if fee_rate_sats_per_vbyte.is_some() {
898 warn!(
901 target: LOG_LIGHTNING,
902 "Ignoring fee_rate_sats_per_vbyte on LDK channel open; LDK uses its built-in fee estimator"
903 );
904 }
905
906 let channel_config = match (base_fee_msat, parts_per_million) {
907 (None, None) => None,
908 (base, ppm) => {
909 let mut config = ChannelConfig::default();
910 if let Some(base) = base {
911 config.forwarding_fee_base_msat = u32::try_from(base).map_err(|_| {
912 LightningRpcError::FailedToOpenChannel {
913 failure_reason: format!(
914 "base_fee_msat {base} does not fit in u32 (LDK limit)"
915 ),
916 }
917 })?;
918 }
919 if let Some(ppm) = ppm {
920 config.forwarding_fee_proportional_millionths =
921 u32::try_from(ppm).map_err(|_| LightningRpcError::FailedToOpenChannel {
922 failure_reason: format!(
923 "parts_per_million {ppm} does not fit in u32 (LDK limit)"
924 ),
925 })?;
926 }
927 Some(config)
928 }
929 };
930
931 let (tx, rx) = oneshot::channel::<Result<OutPoint, String>>();
932
933 {
934 let mut channels = self.pending_channels.write().await;
935 let user_channel_id = self
936 .node
937 .open_announced_channel(
938 pubkey,
939 SocketAddress::from_str(&host).map_err(|e| {
940 LightningRpcError::FailedToConnectToPeer {
941 failure_reason: e.to_string(),
942 }
943 })?,
944 channel_size_sats,
945 push_amount_msats_or,
946 channel_config,
947 )
948 .map_err(|e| LightningRpcError::FailedToOpenChannel {
949 failure_reason: e.to_string(),
950 })?;
951
952 channels.insert(UserChannelId(user_channel_id), tx);
953 }
954
955 match rx
956 .await
957 .map_err(|err| LightningRpcError::FailedToOpenChannel {
958 failure_reason: err.to_string(),
959 })? {
960 Ok(outpoint) => {
961 let funding_txid = outpoint.txid;
962
963 Ok(OpenChannelResponse {
964 funding_txid: funding_txid.to_string(),
965 })
966 }
967 Err(failure_reason) => Err(LightningRpcError::FailedToOpenChannel { failure_reason }),
968 }
969 }
970
971 async fn connect_peer(&self, payload: ConnectPeerRequest) -> Result<(), LightningRpcError> {
972 let NodeAddress { pubkey, address } = payload.node_address;
973 self.node.connect(pubkey, address, true).map_err(|e| {
979 LightningRpcError::FailedToConnectToPeer {
980 failure_reason: e.to_string(),
981 }
982 })
983 }
984
985 async fn close_channels_with_peer(
986 &self,
987 CloseChannelsWithPeerRequest {
988 pubkey,
989 force,
990 sats_per_vbyte: _,
991 }: CloseChannelsWithPeerRequest,
992 ) -> Result<CloseChannelsWithPeerResponse, LightningRpcError> {
993 let mut num_channels_closed = 0;
994
995 info!(%pubkey, "Closing all channels with peer");
996 for channel_with_peer in self
997 .node
998 .list_channels()
999 .iter()
1000 .filter(|channel| channel.counterparty_node_id == pubkey)
1001 {
1002 if force {
1003 match self.node.force_close_channel(
1004 &channel_with_peer.user_channel_id,
1005 pubkey,
1006 Some("User initiated force close".to_string()),
1007 ) {
1008 Ok(()) => num_channels_closed += 1,
1009 Err(err) => {
1010 error!(%pubkey, err = %err.fmt_compact(), "Could not force close channel");
1011 }
1012 }
1013 } else {
1014 match self
1015 .node
1016 .close_channel(&channel_with_peer.user_channel_id, pubkey)
1017 {
1018 Ok(()) => {
1019 num_channels_closed += 1;
1020 }
1021 Err(err) => {
1022 error!(%pubkey, err = %err.fmt_compact(), "Could not close channel");
1023 }
1024 }
1025 }
1026 }
1027
1028 Ok(CloseChannelsWithPeerResponse {
1029 num_channels_closed,
1030 })
1031 }
1032
1033 async fn list_channels(&self) -> Result<ListChannelsResponse, LightningRpcError> {
1034 let mut channels = Vec::new();
1035 let network_graph = self.node.network_graph();
1036
1037 let peer_addresses: std::collections::HashMap<_, _> = self
1039 .node
1040 .list_peers()
1041 .into_iter()
1042 .map(|peer| (peer.node_id, peer.address.to_string()))
1043 .collect();
1044
1045 for channel_details in self.node.list_channels().iter() {
1046 let node_id = NodeId::from_pubkey(&channel_details.counterparty_node_id);
1047 let node_info = network_graph.node(&node_id);
1048
1049 let remote_node_alias = node_info.as_ref().and_then(|info| {
1051 info.announcement_info.as_ref().and_then(|announcement| {
1052 let alias = announcement.alias().to_string();
1053 if alias.is_empty() { None } else { Some(alias) }
1054 })
1055 });
1056
1057 let remote_address = peer_addresses
1058 .get(&channel_details.counterparty_node_id)
1059 .cloned();
1060
1061 channels.push(ChannelInfo {
1062 remote_pubkey: channel_details.counterparty_node_id,
1063 channel_size_sats: channel_details.channel_value_sats,
1064 outbound_liquidity_sats: channel_details.outbound_capacity_msat / 1000,
1065 inbound_liquidity_sats: channel_details.inbound_capacity_msat / 1000,
1066 is_active: channel_details.is_usable,
1067 funding_outpoint: channel_details.funding_txo,
1068 remote_node_alias,
1069 remote_address,
1070 base_fee_msat: Some(u64::from(channel_details.config.forwarding_fee_base_msat)),
1071 parts_per_million: Some(u64::from(
1072 channel_details
1073 .config
1074 .forwarding_fee_proportional_millionths,
1075 )),
1076 });
1077 }
1078
1079 Ok(ListChannelsResponse { channels })
1080 }
1081
1082 async fn set_channel_fees(
1083 &self,
1084 payload: SetChannelFeesRequest,
1085 ) -> Result<(), LightningRpcError> {
1086 let channel = self
1090 .node
1091 .list_channels()
1092 .into_iter()
1093 .find(|c| c.funding_txo == Some(payload.funding_outpoint))
1094 .ok_or_else(|| LightningRpcError::FailedToSetChannelFees {
1095 failure_reason: format!(
1096 "No channel found with funding outpoint {}",
1097 payload.funding_outpoint,
1098 ),
1099 })?;
1100
1101 let forwarding_fee_base_msat = u32::try_from(payload.base_fee_msat).map_err(|_| {
1102 LightningRpcError::FailedToSetChannelFees {
1103 failure_reason: format!(
1104 "base_fee_msat {} does not fit in u32 (LDK limit)",
1105 payload.base_fee_msat,
1106 ),
1107 }
1108 })?;
1109 let forwarding_fee_proportional_millionths = u32::try_from(payload.parts_per_million)
1110 .map_err(|_| LightningRpcError::FailedToSetChannelFees {
1111 failure_reason: format!(
1112 "parts_per_million {} does not fit in u32 (LDK limit)",
1113 payload.parts_per_million,
1114 ),
1115 })?;
1116
1117 let new_config = ChannelConfig {
1120 forwarding_fee_base_msat,
1121 forwarding_fee_proportional_millionths,
1122 ..channel.config
1123 };
1124
1125 self.node
1126 .update_channel_config(
1127 &channel.user_channel_id,
1128 channel.counterparty_node_id,
1129 new_config,
1130 )
1131 .map_err(|e| LightningRpcError::FailedToSetChannelFees {
1132 failure_reason: e.to_string(),
1133 })?;
1134
1135 Ok(())
1136 }
1137
1138 async fn get_balances(&self) -> Result<GetBalancesResponse, LightningRpcError> {
1139 let balances = self.node.list_balances();
1140 let channel_lists = self
1141 .node
1142 .list_channels()
1143 .into_iter()
1144 .filter(|chan| chan.is_usable)
1145 .collect::<Vec<_>>();
1146 let total_inbound_liquidity_balance_msat: u64 = channel_lists
1148 .iter()
1149 .map(|channel| channel.inbound_capacity_msat)
1150 .sum();
1151
1152 Ok(GetBalancesResponse {
1153 onchain_balance_sats: balances.total_onchain_balance_sats,
1154 lightning_balance_msats: balances.total_lightning_balance_sats * 1000,
1155 inbound_lightning_liquidity_msats: total_inbound_liquidity_balance_msat,
1156 })
1157 }
1158
1159 async fn get_invoice(
1160 &self,
1161 get_invoice_request: GetInvoiceRequest,
1162 ) -> Result<Option<GetInvoiceResponse>, LightningRpcError> {
1163 let invoices = self
1164 .node
1165 .list_payments_with_filter(|details| {
1166 details.direction == PaymentDirection::Inbound
1167 && details.id == PaymentId(get_invoice_request.payment_hash.to_byte_array())
1168 && !matches!(details.kind, PaymentKind::Onchain { .. })
1169 })
1170 .iter()
1171 .map(|details| {
1172 let (preimage, payment_hash, _) = get_preimage_and_payment_hash(&details.kind);
1173 let status = match details.status {
1174 PaymentStatus::Failed => fedimint_gateway_common::PaymentStatus::Failed,
1175 PaymentStatus::Succeeded => fedimint_gateway_common::PaymentStatus::Succeeded,
1176 PaymentStatus::Pending => fedimint_gateway_common::PaymentStatus::Pending,
1177 };
1178 GetInvoiceResponse {
1179 preimage: preimage.map(|p| p.to_string()),
1180 payment_hash,
1181 amount: Amount::from_msats(
1182 details
1183 .amount_msat
1184 .expect("amountless invoices are not supported"),
1185 ),
1186 created_at: UNIX_EPOCH + Duration::from_secs(details.latest_update_timestamp),
1187 status,
1188 }
1189 })
1190 .collect::<Vec<_>>();
1191
1192 Ok(invoices.first().cloned())
1193 }
1194
1195 async fn list_transactions(
1196 &self,
1197 start_secs: u64,
1198 end_secs: u64,
1199 ) -> Result<ListTransactionsResponse, LightningRpcError> {
1200 let transactions = self
1201 .node
1202 .list_payments_with_filter(|details| {
1203 !matches!(details.kind, PaymentKind::Onchain { .. })
1204 && details.latest_update_timestamp >= start_secs
1205 && details.latest_update_timestamp < end_secs
1206 })
1207 .iter()
1208 .map(|details| {
1209 let (preimage, payment_hash, payment_kind) =
1210 get_preimage_and_payment_hash(&details.kind);
1211 let direction = match details.direction {
1212 PaymentDirection::Outbound => {
1213 fedimint_gateway_common::PaymentDirection::Outbound
1214 }
1215 PaymentDirection::Inbound => fedimint_gateway_common::PaymentDirection::Inbound,
1216 };
1217 let status = match details.status {
1218 PaymentStatus::Failed => fedimint_gateway_common::PaymentStatus::Failed,
1219 PaymentStatus::Succeeded => fedimint_gateway_common::PaymentStatus::Succeeded,
1220 PaymentStatus::Pending => fedimint_gateway_common::PaymentStatus::Pending,
1221 };
1222 fedimint_gateway_common::PaymentDetails {
1223 payment_hash,
1224 preimage: preimage.map(|p| p.to_string()),
1225 payment_kind,
1226 amount: Amount::from_msats(
1227 details
1228 .amount_msat
1229 .expect("amountless invoices are not supported"),
1230 ),
1231 direction,
1232 status,
1233 timestamp_secs: details.latest_update_timestamp,
1234 }
1235 })
1236 .collect::<Vec<_>>();
1237 Ok(ListTransactionsResponse { transactions })
1238 }
1239
1240 fn create_offer(
1241 &self,
1242 amount: Option<Amount>,
1243 description: Option<String>,
1244 expiry_secs: Option<u32>,
1245 quantity: Option<u64>,
1246 ) -> Result<String, LightningRpcError> {
1247 let description = description.unwrap_or_default();
1248 let offer = if let Some(amount) = amount {
1249 self.node
1250 .bolt12_payment()
1251 .receive(amount.msats, &description, expiry_secs, quantity)
1252 .map_err(|err| LightningRpcError::Bolt12Error {
1253 failure_reason: err.to_string(),
1254 })?
1255 } else {
1256 self.node
1257 .bolt12_payment()
1258 .receive_variable_amount(&description, expiry_secs)
1259 .map_err(|err| LightningRpcError::Bolt12Error {
1260 failure_reason: err.to_string(),
1261 })?
1262 };
1263
1264 Ok(offer.to_string())
1265 }
1266
1267 async fn pay_offer(
1268 &self,
1269 offer: String,
1270 quantity: Option<u64>,
1271 amount: Option<Amount>,
1272 payer_note: Option<String>,
1273 ) -> Result<Preimage, LightningRpcError> {
1274 let offer = Offer::from_str(&offer).map_err(|_| LightningRpcError::Bolt12Error {
1275 failure_reason: "Failed to parse Bolt12 Offer".to_string(),
1276 })?;
1277
1278 let _offer_lock_guard = self
1279 .outbound_offer_lock_pool
1280 .blocking_lock(LdkOfferId(offer.id()));
1281
1282 let payment_id = if let Some(amount) = amount {
1283 self.node
1284 .bolt12_payment()
1285 .send_using_amount(&offer, amount.msats, quantity, payer_note)
1286 .map_err(|err| LightningRpcError::Bolt12Error {
1287 failure_reason: err.to_string(),
1288 })?
1289 } else {
1290 self.node
1291 .bolt12_payment()
1292 .send(&offer, quantity, payer_note)
1293 .map_err(|err| LightningRpcError::Bolt12Error {
1294 failure_reason: err.to_string(),
1295 })?
1296 };
1297
1298 loop {
1299 if let Some(payment_details) = self.node.payment(&payment_id) {
1300 match payment_details.status {
1301 PaymentStatus::Pending => {}
1302 PaymentStatus::Succeeded => match payment_details.kind {
1303 PaymentKind::Bolt12Offer {
1304 preimage: Some(preimage),
1305 ..
1306 } => {
1307 info!(target: LOG_LIGHTNING, offer = %offer, payment_id = %payment_id, preimage = %preimage, "Successfully paid offer");
1308 return Ok(Preimage(preimage.0));
1309 }
1310 _ => {
1311 return Err(LightningRpcError::FailedPayment {
1312 failure_reason: "Unexpected payment kind".to_string(),
1313 });
1314 }
1315 },
1316 PaymentStatus::Failed => {
1317 return Err(LightningRpcError::FailedPayment {
1318 failure_reason: "Bolt12 payment failed".to_string(),
1319 });
1320 }
1321 }
1322 }
1323 fedimint_core::runtime::sleep(Duration::from_millis(100)).await;
1324 }
1325 }
1326
1327 fn sync_wallet(&self) -> Result<(), LightningRpcError> {
1328 block_in_place(|| {
1329 let _ = self.node.sync_wallets();
1330 });
1331 Ok(())
1332 }
1333}
1334
1335fn get_preimage_and_payment_hash(
1338 kind: &PaymentKind,
1339) -> (
1340 Option<Preimage>,
1341 Option<sha256::Hash>,
1342 fedimint_gateway_common::PaymentKind,
1343) {
1344 match kind {
1345 PaymentKind::Bolt11 {
1346 hash,
1347 preimage,
1348 secret: _,
1349 } => (
1350 preimage.map(|p| Preimage(p.0)),
1351 Some(sha256::Hash::from_slice(&hash.0).expect("Failed to convert payment hash")),
1352 fedimint_gateway_common::PaymentKind::Bolt11,
1353 ),
1354 PaymentKind::Bolt11Jit {
1355 hash,
1356 preimage,
1357 secret: _,
1358 lsp_fee_limits: _,
1359 ..
1360 } => (
1361 preimage.map(|p| Preimage(p.0)),
1362 Some(sha256::Hash::from_slice(&hash.0).expect("Failed to convert payment hash")),
1363 fedimint_gateway_common::PaymentKind::Bolt11,
1364 ),
1365 PaymentKind::Bolt12Offer {
1366 hash,
1367 preimage,
1368 secret: _,
1369 offer_id: _,
1370 payer_note: _,
1371 quantity: _,
1372 } => (
1373 preimage.map(|p| Preimage(p.0)),
1374 hash.map(|h| sha256::Hash::from_slice(&h.0).expect("Failed to convert payment hash")),
1375 fedimint_gateway_common::PaymentKind::Bolt12Offer,
1376 ),
1377 PaymentKind::Bolt12Refund {
1378 hash,
1379 preimage,
1380 secret: _,
1381 payer_note: _,
1382 quantity: _,
1383 } => (
1384 preimage.map(|p| Preimage(p.0)),
1385 hash.map(|h| sha256::Hash::from_slice(&h.0).expect("Failed to convert payment hash")),
1386 fedimint_gateway_common::PaymentKind::Bolt12Refund,
1387 ),
1388 PaymentKind::Spontaneous { hash, preimage } => (
1389 preimage.map(|p| Preimage(p.0)),
1390 Some(sha256::Hash::from_slice(&hash.0).expect("Failed to convert payment hash")),
1391 fedimint_gateway_common::PaymentKind::Bolt11,
1392 ),
1393 PaymentKind::Onchain { .. } => (None, None, fedimint_gateway_common::PaymentKind::Onchain),
1394 }
1395}
1396
1397fn get_esplora_url(server_url: SafeUrl) -> Result<String, LdkClientInitError> {
1405 let host = server_url
1407 .host_str()
1408 .ok_or(LdkClientInitError::MissingEsploraHost)?;
1409 let server_url = if let Some(port) = server_url.port() {
1410 format!("{}://{}:{}", server_url.scheme(), host, port)
1411 } else {
1412 server_url.to_string()
1413 };
1414 Ok(server_url)
1415}
1416
1417#[derive(Debug, Clone, Copy, Eq, PartialEq)]
1419enum PendingPaymentWakeup {
1420 NoWaiter,
1422 Woken,
1424 ReceiverDropped,
1426}
1427
1428#[derive(Debug, Clone, Copy, Eq, PartialEq)]
1429struct LdkOfferId(OfferId);
1430
1431impl std::hash::Hash for LdkOfferId {
1432 fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
1433 state.write(&self.0.0);
1434 }
1435}
1436
1437#[derive(Debug, Copy, Clone, PartialEq, Eq)]
1438pub struct UserChannelId(pub ldk_node::UserChannelId);
1439
1440impl PartialOrd for UserChannelId {
1441 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
1442 Some(self.cmp(other))
1443 }
1444}
1445
1446impl Ord for UserChannelId {
1447 fn cmp(&self, other: &Self) -> std::cmp::Ordering {
1448 self.0.0.cmp(&other.0.0)
1449 }
1450}
1451
1452#[cfg(test)]
1453mod tests;