Skip to main content

fedimint_lightning/
ldk.rs

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
47/// Forwards `ldk-node`'s log records into the gateway's `tracing` subscriber.
48///
49/// By default `ldk-node` writes to its own append-only `ldk_node/ldk_node.log`
50/// file, which is invisible to stdout/stderr log collectors and grows without
51/// bound. Routing the records through `tracing` (under the
52/// [`LOG_LIGHTNING_LDK`] target) puts them alongside the rest of gatewayd's
53/// logs and makes them filterable via `RUST_LOG`.
54struct LdkTracingLogger {
55    /// Whether we're running under devimint/tests. When set, some benign LDK
56    /// error logs that are expected in regtest are downgraded to avoid spamming
57    /// the test output. See [`Self::downgraded_level`].
58    in_test_env: bool,
59}
60
61impl LdkTracingLogger {
62    /// Returns the level to emit `record` at, downgrading benign-but-noisy LDK
63    /// errors when running under devimint/tests.
64    ///
65    /// In regtest there is no fee-rate history, so `ldk-node` logs "Failed to
66    /// retrieve fee rate estimates ... Falling back to default" at `Error` on
67    /// essentially every sync. This is harmless (LDK falls back to a default
68    /// feerate), so in test environments we emit it at `Debug` instead. In
69    /// production the original `Error` level is preserved, since a persistent
70    /// failure there can indicate a real problem.
71    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        // `tracing` requires a static level per call-site, so match each LDK level.
87        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    /// The underlying lightning node.
114    node: Arc<ldk_node::Node>,
115
116    task_group: TaskGroup,
117
118    /// The HTLC stream, until it is taken by calling
119    /// `ILnRpcClient::route_htlcs`.
120    htlc_stream_receiver_or: Option<tokio::sync::mpsc::Receiver<InterceptPaymentRequest>>,
121
122    /// Lock pool used to ensure that our implementation of `ILnRpcClient::pay`
123    /// doesn't allow for multiple simultaneous calls with the same invoice to
124    /// execute in parallel. This helps ensure that the function is idempotent.
125    outbound_lightning_payment_lock_pool: lockable::LockPool<PaymentId>,
126
127    /// Lock pool used to ensure that our implementation of
128    /// `ILnRpcClient::pay_offer` doesn't allow for multiple simultaneous
129    /// calls with the same offer to execute in parallel. This helps ensure
130    /// that the function is idempotent.
131    outbound_offer_lock_pool: lockable::LockPool<LdkOfferId>,
132
133    /// A map keyed by the `UserChannelId` of a channel that is currently
134    /// opening. The `Sender` is used to communicate the `OutPoint` back to
135    /// the API handler from the event handler when the channel has been
136    /// opened and is now pending, or the reason it closed before that.
137    pending_channels: Arc<RwLock<BTreeMap<UserChannelId, PendingChannelSender>>>,
138
139    /// Waiters for outgoing LDK payments that are woken by terminal payment
140    /// events (`PaymentSuccessful` / `PaymentFailed`). This lets `pay()` block
141    /// until the payment resolves instead of polling `node.payment()`. The
142    /// actual result is still read from `node.payment()` after the wakeup; this
143    /// map only signals that a terminal event has arrived.
144    pending_payments: Arc<RwLock<HashMap<PaymentId, oneshot::Sender<()>>>>,
145}
146
147/// Sends a channel's funding outpoint once the channel is pending, or the
148/// reason the channel closed before that.
149type 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    /// Creates a new `GatewayLdkClient` instance and starts the underlying
159    /// lightning node. All resources, including the lightning node, will be
160    /// cleaned up when the returned `GatewayLdkClient` instance is dropped.
161    /// There's no need to manually stop the node.
162    ///
163    /// # Errors
164    ///
165    /// Fails with [`LdkClientInitError::MissingEsploraHost`] or
166    /// [`LdkClientInitError::InvalidDataDir`] if the configuration cannot be
167    /// used, with [`LdkClientInitError::Build`] if LDK cannot build the node,
168    /// and with [`LdkClientInitError::Start`] if it cannot start it.
169    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        // Route LDK's logs into the gateway's `tracing` subscriber so they land in
200        // the same place (stderr / log file) and honor `RUST_LOG`, instead of LDK's
201        // default append-only `ldk_node/ldk_node.log` file.
202        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        // We manually check for task termination in case we receive a payment while the
283        // task is shutting down. In that case, we want to finish the payment
284        // before shutting this task down.
285        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                        // LDK reports the real claimable amount, so the two
307                        // amounts coincide here.
308                        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                        // LDK claims payments through its own payment store,
313                        // so it never intercepts forwards for the gateway.
314                        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        // `PaymentClaimable`, `ChannelPending`/`ChannelClosed`, and terminal
376        // outgoing payment events (`PaymentSuccessful` / `PaymentFailed`) are the
377        // only event types that we are interested in. We can safely ignore all
378        // other events.
379        if let Err(err) = node.event_handled() {
380            warn!(err = %err.fmt_compact(), "LDK could not mark event handled");
381        }
382    }
383
384    /// Wakes the `pay()` waiter (if any) for `payment_id` once a terminal
385    /// payment event has been observed. The actual payment result is read from
386    /// `node.payment()` by the woken waiter.
387    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    /// Returns the node's payment record for `payment_id`, but only when it is
403    /// one of our own outbound attempts.
404    ///
405    /// LDK keys BOLT11 payments by `PaymentId(payment_hash)` in both
406    /// directions, so a registered invoice or a claimed inbound payment for the
407    /// same hash shares a slot with our outbound send. Without this direction
408    /// check such an inbound record could be mistaken for the result of our
409    /// `pay()`: reported as a spurious success (a preimage we never sent for),
410    /// a spurious failure, or -- while still pending -- block `pay()` forever.
411    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    /// Reads the result of an outgoing payment from `node.payment()`.
418    ///
419    /// Returns `None` while the payment is still pending (or not yet known to
420    /// the node), and `Some` once it has reached a terminal status.
421    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/// A failure to create the LDK Lightning node behind a [`GatewayLdkClient`].
451///
452/// Creating the client configures an LDK node from the gateway's settings,
453/// builds it on top of the data directory and starts it.
454#[derive(Debug, Error)]
455#[non_exhaustive]
456pub enum LdkClientInitError {
457    /// The Esplora chain source's URL names no host to connect to.
458    #[error("Missing esplora host")]
459    MissingEsploraHost,
460
461    /// The data directory's path is not valid UTF-8, which LDK needs for its
462    /// storage path.
463    #[error("Invalid data dir path")]
464    InvalidDataDir,
465
466    /// LDK could not build the node, for example because its storage could
467    /// not be read or was created for another network.
468    #[error("The LDK node could not be built")]
469    Build(#[source] ldk_node::BuildError),
470
471    /// LDK built the node but could not start it.
472    #[error("The LDK node could not be started")]
473    Start(#[source] ldk_node::NodeError),
474}
475
476/// Why an invoice must not be registered for a payment hash on the node.
477#[derive(Debug, Clone, Copy, PartialEq, Eq)]
478enum InboundRegistrationRefusal {
479    /// `pay()` holds the per-hash lock, so an outbound payment for the hash
480    /// is being dispatched or awaited right now.
481    OutboundInFlight,
482    /// The node holds an outbound record for the hash, pending or terminal.
483    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
499/// Decides whether an invoice may be registered for a payment hash the node
500/// may already know as one of our outbound payments.
501///
502/// `ldk-node` keys BOLT11 payments by `PaymentId(payment_hash)` in both
503/// directions, and registering an invoice overwrites whatever record shares
504/// the key. Registering one for a hash we pay would erase the outbound record
505/// `pay()` reads its result from, so a settled payment would be reported as
506/// failed and the outgoing contract forfeited while the payee keeps the funds.
507///
508/// `outbound_lock_acquired` is whether the caller took the per-hash lock
509/// `pay()` holds for the whole life of a payment, and `existing_direction` is
510/// the direction of the node's record for the hash, if any. A terminal
511/// outbound record is refused as well: a restarted state machine re-runs
512/// `pay()` and needs that record to recover the payment's result instead of
513/// re-dispatching. An existing inbound record is left to the caller, whose
514/// own reservation already rejects duplicate registrations.
515fn 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
529/// Classifies an `ldk-node` claim or fail error for the gateway's completion
530/// retry loop.
531///
532/// `claim_for_hash` and `fail_for_hash` fail deterministically for an unknown
533/// payment hash, a preimage that does not hash to it, or an amount below the
534/// registered one. Retrying cannot change any of those, so they are reported
535/// as [`LightningRpcError::HtlcCompletionRejected`] and the completion state
536/// machine records the outcome instead of retrying forever. Everything else,
537/// today only a failed store write, is treated as transient: the incoming
538/// contract is already funded by the time this runs, so an error this list
539/// does not know must keep retrying rather than be recorded as final.
540fn 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` is used for determining if the Lightning node is ready, so we care
592            // about the `lightning_sync` status.
593            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        // `ILnRpcClient::routehints()` is currently only ever used for LNv1 payment
602        // receives and will be removed when we switch to LNv2. The LDK gateway will
603        // never support LNv1 payment receives, only LNv2 payment receives, which
604        // require that the gateway's lightning node generates invoices rather than the
605        // fedimint client, so it is able to insert the proper route hints on its own.
606        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        // Lock by the payment hash to prevent multiple simultaneous calls with the same
620        // invoice from executing. This prevents `ldk-node::Bolt11Payment::send()` from
621        // being called multiple times with the same invoice. This is important because
622        // `ldk-node::Bolt11Payment::send()` is not idempotent, but this function must
623        // be idempotent.
624        let _payment_lock_guard = self
625            .outbound_lightning_payment_lock_pool
626            .async_lock(payment_id)
627            .await;
628
629        // Register a waiter before initiating the payment so that a terminal
630        // payment event firing immediately after `send()` returns still wakes
631        // us, rather than racing ahead of the registration.
632        let (payment_sender, payment_receiver) = oneshot::channel();
633        self.pending_payments
634            .write()
635            .await
636            .insert(payment_id, payment_sender);
637
638        // If no outbound attempt of ours is known to the node we can initiate
639        // it, and if one is known we can skip calling
640        // `ldk-node::Bolt11Payment::send()` and wait for the payment to
641        // complete. Checking specifically for an outbound record matters because
642        // an inbound payment shares `PaymentId(payment_hash)` with our send: a
643        // registered invoice for the same hash must not make us skip `send()`.
644        // The lock guard above guarantees that this block is only executed once
645        // at a time for a given payment hash, ensuring that there is no race
646        // condition between checking if a payment is known and initiating a new
647        // payment if it isn't.
648        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                    // TODO: Investigate whether all error types returned by
662                    // `Bolt11Payment::send()` result in idempotency.
663                    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        // The payment may already be in a terminal state (a known/resumed
672        // payment, or an event that fired before we registered the waiter), so
673        // check once up front before waiting.
674        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        // Otherwise wait for the event handler to wake us when a terminal
680        // `PaymentSuccessful` / `PaymentFailed` event arrives, instead of
681        // polling. A wakeup is delivered exactly once; the payment status is
682        // terminal by the time it fires.
683        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        // TODO: Get the actual amount from the LDK node. Probably makes the
730        // most sense to pipe it through the `InterceptHtlcResponse` struct.
731        // This value is only used by `ldk-node` to ensure that the amount
732        // claimed isn't less than the amount expected, but we've already
733        // verified that the amount is correct when we intercepted the payment.
734        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                // `pay()` holds this lock for the whole life of an outbound
782                // payment, so failing to take it means a payment for this hash
783                // is in flight right now. Keeping it until `receive_for_hash()`
784                // has run means `pay()` cannot write its outbound record between
785                // our check and the registration's insert, which would
786                // overwrite it.
787                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            // LDK manages its own fee estimation for funding transactions; the
899            // user-supplied rate cannot be applied here.
900            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        // Persist the peer so ldk-node automatically reconnects after restarts
974        // and connection drops. Without this a gateway whose only channel was
975        // opened inbound (e.g. bought from an LSP) has no stored peer address
976        // and the channel stays inactive after any disconnect until an
977        // operator manually reconnects.
978        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        // Build a map of peer pubkey -> address from connected/known peers
1038        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            // Look up peer alias from network graph
1050            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        // ldk-node's `update_channel_config` is keyed by `UserChannelId` +
1087        // counterparty pubkey, so resolve the funding outpoint to those by
1088        // scanning the live channel list.
1089        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        // Copy the channel's current config so we only change the two fee
1118        // fields and preserve cltv_expiry_delta, dust limits, etc.
1119        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        // map and get the total inbound_capacity_msat in the channels
1147        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
1335/// Maps LDK's `PaymentKind` to an optional preimage and an optional payment
1336/// hash depending on the type of payment.
1337fn 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
1397/// When a port is specified in the Esplora URL, the esplora client inside LDK
1398/// node cannot connect to the lightning node when there is a trailing slash.
1399/// The `SafeUrl::Display` function will always serialize the `SafeUrl` with a
1400/// trailing slash, which causes the connection to fail.
1401///
1402/// To handle this, we explicitly construct the esplora URL when a port is
1403/// specified.
1404fn get_esplora_url(server_url: SafeUrl) -> Result<String, LdkClientInitError> {
1405    // Esplora client cannot handle trailing slashes
1406    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/// Outcome of attempting to wake a `pay()` waiter for a terminal payment event.
1418#[derive(Debug, Clone, Copy, Eq, PartialEq)]
1419enum PendingPaymentWakeup {
1420    /// No waiter was registered for the payment id.
1421    NoWaiter,
1422    /// A waiter was registered and successfully woken.
1423    Woken,
1424    /// A waiter was registered but its receiver had already been dropped.
1425    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;