Skip to main content

fedimint_lightning/
lnd.rs

1use std::collections::BTreeMap;
2use std::fmt::{self, Display};
3use std::str::FromStr;
4use std::sync::Arc;
5use std::time::{Duration, UNIX_EPOCH};
6
7use async_trait::async_trait;
8use bitcoin::OutPoint;
9use bitcoin::hashes::{Hash, sha256};
10use fedimint_core::encoding::Encodable;
11use fedimint_core::task::{TaskGroup, sleep};
12use fedimint_core::util::FmtCompact;
13use fedimint_core::{Amount, BitcoinAmountOrAll, crit, secp256k1};
14use fedimint_gateway_common::{
15    ConnectPeerRequest, ListTransactionsResponse, PaymentDetails, PaymentDirection, PaymentKind,
16};
17use fedimint_ln_common::PrunedInvoice;
18use fedimint_ln_common::contracts::Preimage;
19use fedimint_ln_common::route_hints::{RouteHint, RouteHintHop};
20use fedimint_logging::LOG_LIGHTNING;
21use hex::ToHex;
22use secp256k1::PublicKey;
23use tokio::sync::mpsc;
24use tokio_stream::wrappers::ReceiverStream;
25use tonic_lnd::invoicesrpc::lookup_invoice_msg::InvoiceRef;
26use tonic_lnd::invoicesrpc::{
27    AddHoldInvoiceRequest, CancelInvoiceMsg, LookupInvoiceMsg, SettleInvoiceMsg,
28    SubscribeSingleInvoiceRequest,
29};
30use tonic_lnd::lnrpc::channel_point::FundingTxid;
31use tonic_lnd::lnrpc::failure::FailureCode;
32use tonic_lnd::lnrpc::invoice::InvoiceState;
33use tonic_lnd::lnrpc::payment::PaymentStatus;
34use tonic_lnd::lnrpc::policy_update_request::Scope as PolicyUpdateScope;
35use tonic_lnd::lnrpc::{
36    ChanInfoRequest, ChannelBalanceRequest, ChannelPoint, CloseChannelRequest,
37    ConnectPeerRequest as LndConnectPeerRequest, FeeReportRequest, GetInfoRequest, Invoice,
38    InvoiceHtlc, InvoiceHtlcState, InvoiceSubscription, LightningAddress, ListChannelsRequest,
39    ListInvoiceRequest, ListPaymentsRequest, ListPeersRequest, OpenChannelRequest,
40    PolicyUpdateRequest, SendCoinsRequest, UpdateFailure, WalletBalanceRequest,
41};
42use tonic_lnd::routerrpc::{
43    CircuitKey, ForwardHtlcInterceptResponse, ResolveHoldForwardAction, SendPaymentRequest,
44    TrackPaymentRequest,
45};
46use tonic_lnd::tonic::Code;
47use tonic_lnd::walletrpc::AddrRequest;
48use tonic_lnd::{Client as LndClient, connect};
49use tracing::{debug, info, trace, warn};
50
51use super::{
52    ChannelInfo, ILnRpcClient, LightningRpcError, ListChannelsResponse, Lnv2HoldInvoiceFilter,
53    MAX_LIGHTNING_RETRIES, RouteHtlcStream,
54};
55use crate::{
56    CloseChannelsWithPeerRequest, CloseChannelsWithPeerResponse, CreateInvoiceRequest,
57    CreateInvoiceResponse, GetBalancesResponse, GetInvoiceRequest, GetInvoiceResponse,
58    GetLnOnchainAddressResponse, GetNodeInfoResponse, GetRouteHintsResponse,
59    InterceptPaymentRequest, InterceptPaymentResponse, InvoiceDescription, NO_INCOMING_CIRCUIT,
60    OpenChannelResponse, PayInvoiceResponse, PaymentAction, SendOnchainRequest,
61    SendOnchainResponse, SetChannelFeesRequest,
62};
63
64type HtlcSubscriptionSender = mpsc::Sender<InterceptPaymentRequest>;
65
66/// Final CLTV delta of the HOLD invoices created for LNv2 receives. Left unset,
67/// LND would use its routing `bitcoin.timelockdelta`, which operators may lower
68/// to 24 (18 before v0.21), leaving only a few blocks between accepting a
69/// payment and LND cancelling it by itself `invoices.holdexpirydelta` blocks
70/// before expiry. LND rejects any HTLC expiring sooner than this many blocks
71/// after it arrives, so it bounds every accepted HTLC, not just honest payers.
72const LNV2_HOLD_INVOICE_CLTV_EXPIRY: u64 = 144;
73
74/// How many blocks before its earliest HTLC expiry LND is assumed to cancel an
75/// accepted HOLD invoice by itself. LND's `invoices.holdexpirydelta` defaults
76/// to 18 since v0.21 (16 + 2) and 12 before, has no upper bound, and cannot be
77/// read cheaply before v0.21 (`GetDebugInfo` returns the whole log), so this
78/// assumes twice the current default. The excess also absorbs blocks found
79/// between the deadline check and settlement. It costs nothing for a payment
80/// handled promptly, which still has about 108 blocks left under
81/// `LNV2_HOLD_INVOICE_CLTV_EXPIRY`.
82const LND_ASSUMED_HOLD_EXPIRY_DELTA: u32 = 36;
83
84/// Whether a HOLD invoice accepted an HTLC without an MPP total. For the
85/// non-blinded HOLD invoices the gateway creates, standard LND builds require
86/// the payment address, which travels in the MPP record, from every HTLC
87/// except one carrying a keysend record. Since LND v0.20.3 and v0.21.2 that
88/// record must hold the invoice's preimage; before, any value passes unless
89/// `accept-keysend` is on.
90fn has_accepted_keysend_htlc(htlcs: &[InvoiceHtlc]) -> bool {
91    htlcs
92        .iter()
93        .any(|htlc| htlc.state() == InvoiceHtlcState::Accepted && htlc.mpp_total_amt_msat == 0)
94}
95
96/// Block height at which LND is assumed to cancel an accepted HOLD invoice by
97/// itself, `LND_ASSUMED_HOLD_EXPIRY_DELTA` blocks before its earliest accepted
98/// HTLC expires, after which the gateway can no longer settle it. Without an
99/// accepted HTLC there is no deadline to trust, so it is `0`.
100///
101/// The same holds once a keysend HTLC is accepted (see
102/// [`has_accepted_keysend_htlc`]): the invoice then keeps accepting further
103/// keysend HTLCs, and until it restarts LND neither reports them nor moves its
104/// cancel height for them.
105fn hold_invoice_claim_deadline(htlcs: &[InvoiceHtlc]) -> u32 {
106    if has_accepted_keysend_htlc(htlcs) {
107        return 0;
108    }
109
110    htlcs
111        .iter()
112        .filter(|htlc| htlc.state() == InvoiceHtlcState::Accepted)
113        .map(|htlc| (htlc.expiry_height as u32).saturating_sub(LND_ASSUMED_HOLD_EXPIRY_DELTA))
114        .min()
115        .unwrap_or_default()
116}
117
118/// The hex preimage of an LND invoice, if LND knows it. A HOLD invoice has no
119/// preimage until it is settled, so `r_preimage` is empty while one is pending
120/// or after it was canceled.
121fn invoice_preimage_hex(invoice: &Invoice) -> Option<String> {
122    <[u8; 32]>::try_from(invoice.r_preimage.as_slice())
123        .ok()
124        .map(|preimage| preimage.consensus_encode_to_hex())
125}
126
127#[derive(Debug, Clone, Copy, Eq, PartialEq)]
128enum HoldInvoiceAction {
129    Complete,
130    AlreadyComplete,
131}
132
133#[derive(Debug, Clone, Copy, Eq, PartialEq)]
134struct HoldInvoiceStateError {
135    failure_reason: &'static str,
136    permanent: bool,
137}
138
139fn hold_invoice_action(
140    requested_action: PaymentActionKind,
141    invoice_state: Option<InvoiceState>,
142) -> Result<HoldInvoiceAction, HoldInvoiceStateError> {
143    match (requested_action, invoice_state) {
144        (PaymentActionKind::Settle, Some(InvoiceState::Accepted))
145        | (PaymentActionKind::Cancel, Some(InvoiceState::Open | InvoiceState::Accepted)) => {
146            Ok(HoldInvoiceAction::Complete)
147        }
148        (PaymentActionKind::Settle, Some(InvoiceState::Settled))
149        | (PaymentActionKind::Cancel, Some(InvoiceState::Canceled)) => {
150            Ok(HoldInvoiceAction::AlreadyComplete)
151        }
152        (PaymentActionKind::Settle, Some(InvoiceState::Canceled)) => Err(HoldInvoiceStateError {
153            failure_reason: "HOLD invoice was canceled instead of settled",
154            permanent: true,
155        }),
156        (PaymentActionKind::Cancel, Some(InvoiceState::Settled)) => Err(HoldInvoiceStateError {
157            failure_reason: "HOLD invoice was settled instead of canceled",
158            permanent: true,
159        }),
160        (PaymentActionKind::Settle, Some(InvoiceState::Open)) => Err(HoldInvoiceStateError {
161            failure_reason: "HOLD invoice is open and has no accepted HTLC to settle",
162            permanent: false,
163        }),
164        (_, None) => Err(HoldInvoiceStateError {
165            failure_reason: "HOLD invoice does not exist",
166            permanent: true,
167        }),
168    }
169}
170
171#[derive(Debug, Clone, Copy, Eq, PartialEq)]
172enum PaymentActionKind {
173    Settle,
174    Cancel,
175}
176
177#[derive(Clone)]
178pub struct GatewayLndClient {
179    /// LND client
180    address: String,
181    tls_cert: String,
182    macaroon: String,
183    time_pref: f64,
184    /// How long (in seconds) LND keeps trying to route an outgoing payment
185    /// before giving up. Passed as `timeout_seconds` in `SendPaymentRequest`.
186    payment_timeout_secs: i32,
187    lnd_sender: Option<mpsc::Sender<ForwardHtlcInterceptResponse>>,
188    /// Predicate used to distinguish HOLD invoices the gateway created
189    /// (federation-bound) from unrelated HOLD invoices on the same LND node.
190    /// Without this, every HOLD invoice on a shared LND would be intercepted
191    /// as if it were federation-bound, producing invalid LNv1 responses that
192    /// crash LND's htlc_interceptor stream.
193    lnv2_filter: Lnv2HoldInvoiceFilter,
194}
195
196impl GatewayLndClient {
197    pub fn new(
198        address: String,
199        tls_cert: String,
200        macaroon: String,
201        time_pref: f64,
202        payment_timeout_secs: i32,
203        lnd_sender: Option<mpsc::Sender<ForwardHtlcInterceptResponse>>,
204        lnv2_filter: Lnv2HoldInvoiceFilter,
205    ) -> Self {
206        info!(
207            target: LOG_LIGHTNING,
208            address = %address,
209            tls_cert_path = %tls_cert,
210            macaroon = %macaroon,
211            time_pref,
212            payment_timeout_secs,
213            "Gateway configured to connect to LND LnRpcClient",
214        );
215        GatewayLndClient {
216            address,
217            tls_cert,
218            macaroon,
219            time_pref,
220            payment_timeout_secs,
221            lnd_sender,
222            lnv2_filter,
223        }
224    }
225
226    async fn connect(&self) -> Result<LndClient, LightningRpcError> {
227        let mut retries = 0;
228        let client = loop {
229            if retries >= MAX_LIGHTNING_RETRIES {
230                return Err(LightningRpcError::FailedToConnect);
231            }
232
233            retries += 1;
234
235            match connect(
236                self.address.clone(),
237                self.tls_cert.clone(),
238                self.macaroon.clone(),
239            )
240            .await
241            {
242                Ok(client) => break client,
243                Err(err) => {
244                    debug!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Couldn't connect to LND, retrying in 1 second...");
245                    sleep(Duration::from_secs(1)).await;
246                }
247            }
248        };
249
250        Ok(client)
251    }
252
253    async fn connect_peer_if_needed(
254        &self,
255        client: &mut LndClient,
256        pubkey: PublicKey,
257        host: String,
258    ) -> Result<(), LightningRpcError> {
259        let peers = client
260            .lightning()
261            .list_peers(ListPeersRequest { latest_error: true })
262            .await
263            .map_err(|e| LightningRpcError::FailedToConnectToPeer {
264                failure_reason: format!("Could not list peers: {e:?}"),
265            })?
266            .into_inner();
267
268        if peers.peers.into_iter().any(|peer| {
269            PublicKey::from_str(&peer.pub_key).expect("LND returned invalid peer public key")
270                == pubkey
271        }) {
272            return Ok(());
273        }
274
275        client
276            .lightning()
277            .connect_peer(LndConnectPeerRequest {
278                addr: Some(LightningAddress {
279                    pubkey: pubkey.to_string(),
280                    host,
281                }),
282                perm: false,
283                timeout: 10,
284            })
285            .await
286            .map_err(|e| LightningRpcError::FailedToConnectToPeer {
287                failure_reason: format!("Failed to connect to peer {e:?}"),
288            })?;
289
290        Ok(())
291    }
292
293    /// Spawns a new background task that subscribes to updates of a specific
294    /// HOLD invoice. When the HOLD invoice is ACCEPTED, we can request the
295    /// preimage from the Gateway. A new task is necessary because LND's
296    /// global `subscribe_invoices` does not currently emit updates for HOLD invoices: <https://github.com/lightningnetwork/lnd/issues/3120>
297    async fn spawn_lnv2_hold_invoice_subscription(
298        &self,
299        task_group: &TaskGroup,
300        payment_stream_group: TaskGroup,
301        gateway_sender: HtlcSubscriptionSender,
302        payment_hash: Vec<u8>,
303    ) -> Result<(), LightningRpcError> {
304        let mut client = self.connect().await?;
305
306        let self_copy = self.clone();
307        let r_hash = payment_hash.clone();
308        task_group.spawn("LND HOLD Invoice Subscription", |handle| async move {
309            let future_stream =
310                client
311                    .invoices()
312                    .subscribe_single_invoice(SubscribeSingleInvoiceRequest {
313                        r_hash: r_hash.clone(),
314                    });
315
316            let mut hold_stream = tokio::select! {
317                stream = future_stream => {
318                    match stream {
319                        Ok(stream) => stream.into_inner(),
320                        Err(err) => {
321                            crit!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Failed to subscribe to hold invoice updates, shutting down payment-stream subgroup to trigger gateway reconnect");
322                            payment_stream_group.shutdown();
323                            return;
324                        }
325                    }
326                },
327                () = handle.make_shutdown_rx() => {
328                    info!(target: LOG_LIGHTNING, "LND HOLD Invoice Subscription received shutdown signal");
329                    return;
330                }
331            };
332
333            loop {
334                let hold = tokio::select! {
335                    () = handle.make_shutdown_rx() => {
336                        info!(target: LOG_LIGHTNING, "LND HOLD Invoice Subscription received shutdown signal");
337                        break;
338                    }
339                    hold_update = hold_stream.message() => {
340                        match hold_update {
341                            Ok(Some(hold)) => hold,
342                            Ok(None) => {
343                                // LND closed the stream because the invoice
344                                // reached a terminal state (settled, canceled,
345                                // or expired).
346                                break;
347                            }
348                            Err(err) => {
349                                crit!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Error received over hold invoice update stream, shutting down payment-stream subgroup to trigger gateway reconnect");
350                                payment_stream_group.shutdown();
351                                break;
352                            }
353                        }
354                    }
355                };
356
357                debug!(
358                    target: LOG_LIGHTNING,
359                    payment_hash = %PrettyPaymentHash(&r_hash),
360                    state = %hold.state,
361                    "LND HOLD Invoice Update",
362                );
363
364                if hold.state() == InvoiceState::Accepted {
365                    // Only forward HOLD invoices that the gateway created on
366                    // behalf of a federation. We check here (rather than at
367                    // the new-invoice add-event) because the contract is
368                    // saved to gateway_db *after* the HOLD invoice is created
369                    // on LND, so the add-event races registration. By the
370                    // time `Accepted` fires the HTLC has arrived, which means
371                    // the BOLT11 invoice was published and the contract is
372                    // committed.
373                    let hash = sha256::Hash::from_slice(&hold.r_hash)
374                        .expect("LND payment hashes are 32 bytes");
375                    if !(self_copy.lnv2_filter)(hash).await {
376                        trace!(
377                            target: LOG_LIGHTNING,
378                            payment_hash = %PrettyPaymentHash(&hold.r_hash),
379                            "Ignoring HOLD invoice not created by this gateway",
380                        );
381                        continue;
382                    }
383
384                    if has_accepted_keysend_htlc(&hold.htlcs) {
385                        warn!(
386                            target: LOG_LIGHTNING,
387                            payment_hash = %PrettyPaymentHash(&hold.r_hash),
388                            "LNv2 HOLD invoice accepted a keysend HTLC, its claim deadline is not trusted",
389                        );
390                    }
391
392                    let (incoming_chan_id, htlc_id) = NO_INCOMING_CIRCUIT;
393                    let intercept = InterceptPaymentRequest {
394                        payment_hash: Hash::from_slice(&hold.r_hash.clone())
395                            .expect("Failed to convert to Hash"),
396                        // A HOLD invoice reports the real paid amount, so the
397                        // two amounts coincide here.
398                        amount_msat: hold.amt_paid_msat as u64,
399                        incoming_amount_msat: hold.amt_paid_msat as u64,
400                        expiry: hold_invoice_claim_deadline(&hold.htlcs),
401                        short_channel_id: Some(0),
402                        // The payment is held by a HOLD invoice on our own
403                        // node rather than by an intercepted forward, which is
404                        // how `complete_htlc` knows to resolve it by settling
405                        // or canceling that invoice.
406                        incoming_chan_id,
407                        htlc_id,
408                    };
409
410                    match gateway_sender.send(intercept).await {
411                        Ok(()) => {}
412                        Err(err) => {
413                            warn!(
414                                target: LOG_LIGHTNING,
415                                err = %err.fmt_compact(),
416                                "Hold Invoice Subscription failed to send Intercept to gateway"
417                            );
418                            let _ = self_copy.cancel_hold_invoice(hold.r_hash).await;
419                        }
420                    }
421                }
422            }
423        });
424
425        Ok(())
426    }
427
428    /// Spawns a new background task that subscribes to "add" updates for all
429    /// invoices. This is used to detect when a new invoice has been
430    /// created. If this invoice is a HOLD invoice, it is potentially destined
431    /// for a federation. At this point, we spawn a separate task to monitor the
432    /// status of the HOLD invoice.
433    async fn spawn_lnv2_invoice_subscription(
434        &self,
435        task_group: &TaskGroup,
436        gateway_sender: HtlcSubscriptionSender,
437    ) -> Result<(), LightningRpcError> {
438        let mut client = self.connect().await?;
439
440        let list_response = client
441            .lightning()
442            .list_invoices(ListInvoiceRequest {
443                pending_only: true,
444                index_offset: 0,
445                num_max_invoices: u64::MAX,
446                reversed: false,
447                ..Default::default()
448            })
449            .await
450            .map_err(|status| {
451                warn!(target: LOG_LIGHTNING, status = %status, "Failed to list all invoices");
452                LightningRpcError::FailedToRouteHtlcs {
453                    failure_reason: "Failed to list all invoices".to_string(),
454                }
455            })?
456            .into_inner();
457
458        let self_copy = self.clone();
459        let hold_group = task_group.make_subgroup();
460        // See the matching comment in `spawn_lnv1_htlc_interceptor`: if this
461        // task exits unexpectedly we shut down the payment-stream subgroup so
462        // the gateway transitions to `Disconnected` and reconnects.
463        let subgroup = task_group.clone();
464
465        // The `SubscribeInvoices` backlog cannot replay pre-existing pending
466        // HOLD invoices reliably: an `add_index` of 0 means "no backlog" to
467        // LND, so the oldest pending invoice is skipped whenever it is the
468        // first invoice ever created on the node (`add_index == 1`), and
469        // invoices already in the `Accepted` state are never delivered to
470        // all-invoice subscribers at all. Instead, spawn a monitor task for
471        // every pending HOLD invoice (`r_preimage` empty) from the listing
472        // directly. Invoices not created by this gateway are filtered out by
473        // the monitor task once they reach the `Accepted` state.
474        for invoice in &list_response.invoices {
475            if invoice.r_preimage.is_empty() {
476                info!(
477                    target: LOG_LIGHTNING,
478                    payment_hash = %PrettyPaymentHash(&invoice.r_hash),
479                    "Monitoring pre-existing pending LNv2 invoice",
480                );
481                self.spawn_lnv2_hold_invoice_subscription(
482                    &hold_group,
483                    subgroup.clone(),
484                    gateway_sender.clone(),
485                    invoice.r_hash.clone(),
486                )
487                .await?;
488            }
489        }
490
491        // The listing above covers all pending invoices, so the subscription
492        // only needs to replay invoices added after the listing was taken.
493        // Add indices increase monotonically, so any such invoice has an
494        // `add_index` strictly greater than the listing's `last_index_offset`
495        // (0 when no invoices were listed) and `SubscribeInvoices` replays
496        // exactly those, without duplicating the invoices handled above.
497        let add_index = list_response.last_index_offset;
498        task_group.spawn("LND Invoice Subscription", move |handle| async move {
499            let future_stream = client.lightning().subscribe_invoices(InvoiceSubscription {
500                add_index,
501                settle_index: u64::MAX, // we do not need settle invoice events
502            });
503            let mut invoice_stream = tokio::select! {
504                stream = future_stream => {
505                    match stream {
506                        Ok(stream) => stream.into_inner(),
507                        Err(err) => {
508                            warn!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Failed to subscribe to all invoice updates");
509                            subgroup.shutdown();
510                            return;
511                        }
512                    }
513                },
514                () = handle.make_shutdown_rx() => {
515                    info!(target: LOG_LIGHTNING, "LND Invoice Subscription received shutdown signal");
516                    return;
517                }
518            };
519
520            info!(target: LOG_LIGHTNING, "LND Invoice Subscription: starting to process invoice updates");
521            while let Some(invoice) = tokio::select! {
522                () = handle.make_shutdown_rx() => {
523                    info!(target: LOG_LIGHTNING, "LND Invoice Subscription task received shutdown signal");
524                    None
525                }
526                invoice_update = invoice_stream.message() => {
527                    match invoice_update {
528                        Ok(invoice) => invoice,
529                        Err(err) => {
530                            warn!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Error received over invoice update stream");
531                            None
532                        }
533                    }
534                }
535            } {
536                // If the `r_preimage` is empty and the invoice is OPEN, this means a new HOLD
537                // invoice has been created, which is potentially an invoice destined for a
538                // federation. We will spawn a new task to monitor the status of
539                // the HOLD invoice.
540                let payment_hash = invoice.r_hash.clone();
541
542                debug!(
543                    target: LOG_LIGHTNING,
544                    payment_hash = %PrettyPaymentHash(&payment_hash),
545                    state = %invoice.state,
546                    "LND HOLD Invoice Update",
547                );
548
549                if invoice.r_preimage.is_empty() && invoice.state() == InvoiceState::Open {
550                    info!(
551                        target: LOG_LIGHTNING,
552                        payment_hash = %PrettyPaymentHash(&payment_hash),
553                        "Monitoring new LNv2 invoice",
554                    );
555                    if let Err(err) = self_copy
556                        .spawn_lnv2_hold_invoice_subscription(
557                            &hold_group,
558                            subgroup.clone(),
559                            gateway_sender.clone(),
560                            payment_hash.clone(),
561                        )
562                        .await
563                    {
564                        // Spawning failed because `connect()` exhausted its
565                        // retries, a strong signal that LND is unreachable. We
566                        // can no longer observe this invoice's `Accepted`
567                        // update, so shut down the payment-stream subgroup to
568                        // force a gateway reconnect.
569                        warn!(
570                            target: LOG_LIGHTNING,
571                            err = %err.fmt_compact(),
572                            payment_hash = %PrettyPaymentHash(&payment_hash),
573                            "Failed to spawn HOLD invoice subscription task, shutting down payment-stream subgroup to trigger gateway reconnect",
574                        );
575                        subgroup.shutdown();
576                    }
577                }
578            }
579
580            if !handle.is_shutting_down() {
581                warn!(target: LOG_LIGHTNING, "LND Invoice Subscription exited unexpectedly, shutting down payment-stream subgroup to trigger gateway reconnect");
582                subgroup.shutdown();
583            }
584        });
585
586        Ok(())
587    }
588
589    /// Spawns a new background task that intercepts HTLCs from the LND node. In
590    /// the LNv1 protocol, this is used as a trigger mechanism for
591    /// requesting the Gateway to retrieve the preimage for a payment.
592    async fn spawn_lnv1_htlc_interceptor(
593        &self,
594        task_group: &TaskGroup,
595        lnd_sender: mpsc::Sender<ForwardHtlcInterceptResponse>,
596        lnd_rx: mpsc::Receiver<ForwardHtlcInterceptResponse>,
597        gateway_sender: HtlcSubscriptionSender,
598    ) -> Result<(), LightningRpcError> {
599        let mut client = self.connect().await?;
600
601        // Verify that LND is reachable via RPC before attempting to spawn a new thread
602        // that will intercept HTLCs.
603        client
604            .lightning()
605            .get_info(GetInfoRequest {})
606            .await
607            .map_err(|status| LightningRpcError::FailedToGetNodeInfo {
608                failure_reason: format!("Failed to get node info {status:?}"),
609            })?;
610
611        // If the HTLC interceptor exits unexpectedly we shut down the
612        // payment-stream subgroup. That cascades to the lnv2 invoice
613        // subscription (and its hold-invoice subtasks), which drop their
614        // `gateway_sender` clones, closing the gateway's HTLC stream and
615        // driving the gateway back to `Disconnected` so it reconnects.
616        let subgroup = task_group.clone();
617        task_group.spawn("LND HTLC Subscription", |handle| async move {
618                let future_stream = client
619                    .router()
620                    .htlc_interceptor(ReceiverStream::new(lnd_rx));
621                let mut htlc_stream = tokio::select! {
622                    stream = future_stream => {
623                        match stream {
624                            Ok(stream) => stream.into_inner(),
625                            Err(e) => {
626                                crit!(target: LOG_LIGHTNING, err = %e.fmt_compact(), "Failed to establish htlc stream");
627                                subgroup.shutdown();
628                                return;
629                            }
630                        }
631                    },
632                    () = handle.make_shutdown_rx() => {
633                        info!(target: LOG_LIGHTNING, "LND HTLC Subscription received shutdown signal while trying to intercept HTLC stream, exiting...");
634                        return;
635                    }
636                };
637
638                debug!(target: LOG_LIGHTNING, "LND HTLC Subscription: starting to process stream");
639                // To gracefully handle shutdown signals, we need to be able to receive signals
640                // while waiting for the next message from the HTLC stream.
641                //
642                // If we're in the middle of processing a message from the stream, we need to
643                // finish before stopping the spawned task. Checking if the task group is
644                // shutting down at the start of each iteration will cause shutdown signals to
645                // not process until another message arrives from the HTLC stream, which may
646                // take a long time, or never.
647                while let Some(htlc) = tokio::select! {
648                    () = handle.make_shutdown_rx() => {
649                        info!(target: LOG_LIGHTNING, "LND HTLC Subscription task received shutdown signal");
650                        None
651                    }
652                    htlc_message = htlc_stream.message() => {
653                        match htlc_message {
654                            Ok(htlc) => htlc,
655                            Err(err) => {
656                                warn!(target: LOG_LIGHTNING, err = %err.fmt_compact(), "Error received over HTLC stream");
657                                None
658                            }
659                    }}
660                } {
661                    trace!(target: LOG_LIGHTNING, ?htlc, "LND Handling HTLC");
662
663                    let Some(incoming_circuit_key) = htlc.incoming_circuit_key else {
664                        // We have no circuit key, so the HTLC cannot be cancelled
665                        // either; it will time out at LND. Log enough context to
666                        // correlate with the sender's invoice and the target
667                        // federation.
668                        warn!(
669                            target: LOG_LIGHTNING,
670                            payment_hash = %PrettyPaymentHash(&htlc.payment_hash),
671                            scid = htlc.outgoing_requested_chan_id,
672                            amount_msat = htlc.outgoing_amount_msat,
673                            "Cannot route HTLC: incoming_circuit_key is None"
674                        );
675                        continue;
676                    };
677
678                    let chan_id = incoming_circuit_key.chan_id;
679                    let htlc_id = incoming_circuit_key.htlc_id;
680
681                    // Forward all HTLCs to gatewayd, gatewayd will filter them based on scid
682                    let intercept = InterceptPaymentRequest {
683                        payment_hash: Hash::from_slice(&htlc.payment_hash).expect("Failed to convert payment Hash"),
684                        // `outgoing_amount_msat` is the sender-written onion
685                        // `amt_to_forward`; `incoming_amount_msat` is the amount
686                        // actually locked in the incoming HTLC. Carry both so
687                        // downstream funding is decided against the real value.
688                        amount_msat: htlc.outgoing_amount_msat,
689                        incoming_amount_msat: htlc.incoming_amount_msat,
690                        expiry: htlc.incoming_expiry,
691                        short_channel_id: Some(htlc.outgoing_requested_chan_id),
692                        incoming_chan_id: chan_id,
693                        htlc_id,
694                    };
695
696                    match gateway_sender.send(intercept).await {
697                        Ok(()) => {}
698                        Err(err) => {
699                            warn!(
700                                target: LOG_LIGHTNING,
701                                err = %err.fmt_compact(),
702                                payment_hash = %PrettyPaymentHash(&htlc.payment_hash),
703                                scid = htlc.outgoing_requested_chan_id,
704                                amount_msat = htlc.outgoing_amount_msat,
705                                "Failed to send HTLC to gatewayd for processing"
706                            );
707                            let _ = Self::cancel_htlc(incoming_circuit_key, lnd_sender.clone())
708                                .await
709                                .map_err(|err| {
710                                    warn!(
711                                        target: LOG_LIGHTNING,
712                                        err = %err.fmt_compact(),
713                                        payment_hash = %PrettyPaymentHash(&htlc.payment_hash),
714                                        chan_id,
715                                        htlc_id,
716                                        "Failed to cancel HTLC"
717                                    );
718                                });
719                        }
720                    }
721                }
722
723                // Loop exited because of an HTLC stream error or end-of-stream
724                // (the expected-shutdown case is handled above).
725                if !handle.is_shutting_down() {
726                    warn!(target: LOG_LIGHTNING, "LND HTLC Subscription exited unexpectedly, shutting down payment-stream subgroup to trigger gateway reconnect");
727                    subgroup.shutdown();
728                }
729            });
730
731        Ok(())
732    }
733
734    /// Spawns background tasks for monitoring the status of incoming payments.
735    async fn spawn_interceptor(
736        &self,
737        task_group: &TaskGroup,
738        lnd_sender: mpsc::Sender<ForwardHtlcInterceptResponse>,
739        lnd_rx: mpsc::Receiver<ForwardHtlcInterceptResponse>,
740        gateway_sender: HtlcSubscriptionSender,
741    ) -> Result<(), LightningRpcError> {
742        self.spawn_lnv1_htlc_interceptor(task_group, lnd_sender, lnd_rx, gateway_sender.clone())
743            .await?;
744
745        self.spawn_lnv2_invoice_subscription(task_group, gateway_sender)
746            .await?;
747
748        Ok(())
749    }
750
751    async fn cancel_htlc(
752        key: CircuitKey,
753        lnd_sender: mpsc::Sender<ForwardHtlcInterceptResponse>,
754    ) -> Result<(), LightningRpcError> {
755        // TODO: Specify a failure code and message
756        let response = ForwardHtlcInterceptResponse {
757            incoming_circuit_key: Some(key),
758            action: ResolveHoldForwardAction::Fail.into(),
759            preimage: vec![],
760            failure_message: vec![],
761            failure_code: FailureCode::TemporaryChannelFailure.into(),
762            ..Default::default()
763        };
764        Self::send_lnd_response(lnd_sender, response).await
765    }
766
767    async fn send_lnd_response(
768        lnd_sender: mpsc::Sender<ForwardHtlcInterceptResponse>,
769        response: ForwardHtlcInterceptResponse,
770    ) -> Result<(), LightningRpcError> {
771        // TODO: Consider retrying this if the send fails
772        lnd_sender.send(response).await.map_err(|send_error| {
773            LightningRpcError::FailedToCompleteHtlc {
774                failure_reason: format!(
775                    "Failed to send ForwardHtlcInterceptResponse to LND {send_error:?}"
776                ),
777            }
778        })
779    }
780
781    async fn lookup_payment(
782        &self,
783        payment_hash: Vec<u8>,
784        client: &mut LndClient,
785    ) -> Result<Option<String>, LightningRpcError> {
786        // Loop until we successfully get the status of the payment, or determine that
787        // the payment has not been made yet.
788        loop {
789            let payments = client
790                .router()
791                .track_payment_v2(TrackPaymentRequest {
792                    payment_hash: payment_hash.clone(),
793                    no_inflight_updates: true,
794                })
795                .await;
796
797            match payments {
798                Ok(payments) => {
799                    let mut updates = payments.into_inner();
800
801                    // Block until LND reports a terminal status. `no_inflight_updates`
802                    // asks LND to hold back everything else, but a payment still in
803                    // flight must never be read as failed, so only `Failed` counts as
804                    // a failure here: `InFlight`, `Initiated` (with
805                    // `routerrpc.usestatusinitiated`) and any status this build does
806                    // not know, which prost decodes as `Unknown`, keep us waiting.
807                    let outcome = loop {
808                        match updates.message().await {
809                            Ok(Some(payment)) if payment.status() == PaymentStatus::Succeeded => {
810                                return Ok(Some(payment.payment_preimage));
811                            }
812                            Ok(Some(payment)) if payment.status() == PaymentStatus::Failed => {
813                                let failure_reason = payment.failure_reason();
814                                return Err(LightningRpcError::FailedPayment {
815                                    failure_reason: format!("{failure_reason:?}"),
816                                });
817                            }
818                            Ok(Some(payment)) => {
819                                debug!(
820                                    target: LOG_LIGHTNING,
821                                    payment_hash = %PrettyPaymentHash(&payment_hash),
822                                    status = ?payment.status(),
823                                    "Tracked payment is not terminal yet, waiting",
824                                );
825                            }
826                            outcome => break outcome,
827                        }
828                    };
829
830                    // A premature end of stream (`Ok(None)`) or a transport fault
831                    // (`Err`) is not a payment outcome. Retry rather than reporting a
832                    // failure the node never produced.
833                    warn!(
834                        target: LOG_LIGHTNING,
835                        payment_hash = %PrettyPaymentHash(&payment_hash),
836                        outcome = ?outcome,
837                        "Payment tracking stream ended or faulted. Trying again in 5 seconds"
838                    );
839                    sleep(Duration::from_secs(5)).await;
840                }
841                Err(err) => {
842                    // Break if we got a response back from the LND node that indicates the payment
843                    // hash was not found.
844                    if err.code() == Code::NotFound {
845                        return Ok(None);
846                    }
847
848                    warn!(
849                        target: LOG_LIGHTNING,
850                        payment_hash = %PrettyPaymentHash(&payment_hash),
851                        err = %err.fmt_compact(),
852                        "Could not get the status of payment. Trying again in 5 seconds"
853                    );
854                    sleep(Duration::from_secs(5)).await;
855                }
856            }
857        }
858    }
859
860    /// Looks up the invoice carrying `payment_hash`, returning `None` if the
861    /// node has no such invoice.
862    ///
863    /// Any other lookup failure is reported as an error, since callers retry on
864    /// error and an unreachable node is a condition a retry can clear.
865    async fn lookup_invoice(
866        client: &mut LndClient,
867        payment_hash: &[u8],
868    ) -> Result<Option<Invoice>, LightningRpcError> {
869        match client
870            .invoices()
871            .lookup_invoice_v2(LookupInvoiceMsg {
872                invoice_ref: Some(InvoiceRef::PaymentHash(payment_hash.to_vec())),
873                lookup_modifier: 0,
874            })
875            .await
876        {
877            Ok(invoice) => Ok(Some(invoice.into_inner())),
878            Err(err) if err.code() == Code::NotFound => Ok(None),
879            Err(err) => Err(LightningRpcError::FailedToCompleteHtlc {
880                failure_reason: format!("Failed to look up invoice: {}", err.fmt_compact()),
881            }),
882        }
883    }
884
885    /// Settles the LNv2 HOLD invoice carrying `payment_hash` with `preimage`.
886    ///
887    /// Only an already-settled invoice is an idempotent success. Missing,
888    /// nonterminal, and canceled invoices fail so callers cannot record a
889    /// settle outcome that Lightning did not produce.
890    async fn settle_hold_invoice(
891        &self,
892        payment_hash: Vec<u8>,
893        preimage: Preimage,
894    ) -> Result<(), LightningRpcError> {
895        let mut client = self.connect().await?;
896        let invoice = Self::lookup_invoice(&mut client, &payment_hash).await?;
897        match hold_invoice_action(
898            PaymentActionKind::Settle,
899            invoice.as_ref().map(Invoice::state),
900        ) {
901            Ok(HoldInvoiceAction::Complete) => {}
902            Ok(HoldInvoiceAction::AlreadyComplete) => {
903                info!(
904                    target: LOG_LIGHTNING,
905                    payment_hash = %PrettyPaymentHash(&payment_hash),
906                    "HOLD invoice was already settled",
907                );
908                return Ok(());
909            }
910            Err(error) => {
911                warn!(
912                    target: LOG_LIGHTNING,
913                    state = ?invoice.as_ref().map(Invoice::state),
914                    payment_hash = %PrettyPaymentHash(&payment_hash),
915                    failure_reason = error.failure_reason,
916                    "Cannot settle HOLD invoice",
917                );
918                return Err(if error.permanent {
919                    LightningRpcError::HtlcCompletionRejected {
920                        failure_reason: error.failure_reason.to_owned(),
921                    }
922                } else {
923                    LightningRpcError::FailedToCompleteHtlc {
924                        failure_reason: error.failure_reason.to_owned(),
925                    }
926                });
927            }
928        }
929
930        client
931            .invoices()
932            .settle_invoice(SettleInvoiceMsg {
933                preimage: preimage.0.to_vec(),
934            })
935            .await
936            .map_err(|err| {
937                warn!(
938                    target: LOG_LIGHTNING,
939                    err = %err.fmt_compact(),
940                    payment_hash = %PrettyPaymentHash(&payment_hash),
941                    "Failed to settle HOLD invoice",
942                );
943                LightningRpcError::FailedToCompleteHtlc {
944                    failure_reason: "Failed to settle HOLD invoice".to_string(),
945                }
946            })?;
947
948        info!(
949            target: LOG_LIGHTNING,
950            payment_hash = %PrettyPaymentHash(&payment_hash),
951            "Successfully settled HOLD invoice",
952        );
953
954        Ok(())
955    }
956
957    /// Cancels the LNv2 HOLD invoice carrying `payment_hash`, failing back any
958    /// HTLC it holds.
959    ///
960    /// Only an already-canceled invoice is an idempotent success. A settled
961    /// invoice fails so callers cannot record a cancel outcome after a racing
962    /// settle won.
963    async fn cancel_hold_invoice(&self, payment_hash: Vec<u8>) -> Result<(), LightningRpcError> {
964        let mut client = self.connect().await?;
965        let invoice = Self::lookup_invoice(&mut client, &payment_hash).await?;
966        match hold_invoice_action(
967            PaymentActionKind::Cancel,
968            invoice.as_ref().map(Invoice::state),
969        ) {
970            Ok(HoldInvoiceAction::Complete) => {}
971            Ok(HoldInvoiceAction::AlreadyComplete) => {
972                info!(
973                    target: LOG_LIGHTNING,
974                    payment_hash = %PrettyPaymentHash(&payment_hash),
975                    "HOLD invoice was already canceled",
976                );
977                return Ok(());
978            }
979            Err(error) => {
980                warn!(
981                    target: LOG_LIGHTNING,
982                    state = ?invoice.as_ref().map(Invoice::state),
983                    payment_hash = %PrettyPaymentHash(&payment_hash),
984                    failure_reason = error.failure_reason,
985                    "Cannot cancel HOLD invoice",
986                );
987                return Err(LightningRpcError::HtlcCompletionRejected {
988                    failure_reason: error.failure_reason.to_owned(),
989                });
990            }
991        }
992
993        client
994            .invoices()
995            .cancel_invoice(CancelInvoiceMsg {
996                payment_hash: payment_hash.clone(),
997            })
998            .await
999            .map_err(|err| {
1000                warn!(
1001                    target: LOG_LIGHTNING,
1002                    err = %err.fmt_compact(),
1003                    payment_hash = %PrettyPaymentHash(&payment_hash),
1004                    "Failed to cancel HOLD invoice",
1005                );
1006                LightningRpcError::FailedToCompleteHtlc {
1007                    failure_reason: "Failed to cancel HOLD invoice".to_string(),
1008                }
1009            })?;
1010
1011        info!(
1012            target: LOG_LIGHTNING,
1013            payment_hash = %PrettyPaymentHash(&payment_hash),
1014            "Successfully canceled HOLD invoice",
1015        );
1016
1017        Ok(())
1018    }
1019}
1020
1021impl fmt::Debug for GatewayLndClient {
1022    fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
1023        write!(f, "LndClient")
1024    }
1025}
1026
1027#[async_trait]
1028impl ILnRpcClient for GatewayLndClient {
1029    async fn info(&self) -> Result<GetNodeInfoResponse, LightningRpcError> {
1030        let mut client = self.connect().await?;
1031        let info = client
1032            .lightning()
1033            .get_info(GetInfoRequest {})
1034            .await
1035            .map_err(|status| LightningRpcError::FailedToGetNodeInfo {
1036                failure_reason: format!("Failed to get node info {status:?}"),
1037            })?
1038            .into_inner();
1039
1040        let pub_key: PublicKey =
1041            info.identity_pubkey
1042                .parse()
1043                .map_err(|e| LightningRpcError::FailedToGetNodeInfo {
1044                    failure_reason: format!("Failed to parse public key {e:?}"),
1045                })?;
1046
1047        let network = match info
1048            .chains
1049            .first()
1050            .ok_or_else(|| LightningRpcError::FailedToGetNodeInfo {
1051                failure_reason: "Failed to parse node network".to_string(),
1052            })?
1053            .network
1054            .as_str()
1055        {
1056            // LND uses "mainnet", but rust-bitcoin uses "bitcoin".
1057            // TODO: create a fedimint `Network` type that understands "mainnet"
1058            "mainnet" => "bitcoin",
1059            other => other,
1060        }
1061        .to_string();
1062
1063        return Ok(GetNodeInfoResponse {
1064            pub_key,
1065            alias: info.alias,
1066            network,
1067            block_height: info.block_height,
1068            synced_to_chain: info.synced_to_chain,
1069        });
1070    }
1071
1072    async fn routehints(
1073        &self,
1074        num_route_hints: usize,
1075    ) -> Result<GetRouteHintsResponse, LightningRpcError> {
1076        let mut client = self.connect().await?;
1077        let mut channels = client
1078            .lightning()
1079            .list_channels(ListChannelsRequest {
1080                active_only: true,
1081                inactive_only: false,
1082                public_only: false,
1083                private_only: false,
1084                peer: vec![],
1085                peer_alias_lookup: false,
1086            })
1087            .await
1088            .map_err(|status| LightningRpcError::FailedToGetRouteHints {
1089                failure_reason: format!("Failed to list channels {status:?}"),
1090            })?
1091            .into_inner()
1092            .channels;
1093
1094        // Take the channels with the largest incoming capacity
1095        channels.sort_by_key(|b| std::cmp::Reverse(b.remote_balance));
1096        channels.truncate(num_route_hints);
1097
1098        let mut route_hints: Vec<RouteHint> = vec![];
1099        for chan in &channels {
1100            let info = client
1101                .lightning()
1102                .get_chan_info(ChanInfoRequest {
1103                    chan_id: chan.chan_id,
1104                    ..Default::default()
1105                })
1106                .await
1107                .map_err(|status| LightningRpcError::FailedToGetRouteHints {
1108                    failure_reason: format!("Failed to get channel info {status:?}"),
1109                })?
1110                .into_inner();
1111
1112            // The route hint hop describes the remote peer forwarding a payment
1113            // *into* the gateway, so it must carry that peer's advertised channel
1114            // policy. `node1`/`node2` are ordered by pubkey (BOLT 7), not by which
1115            // side is the gateway, so select the policy belonging to the remote
1116            // peer rather than always taking node1's.
1117            let policy = if info.node1_pub == chan.remote_pubkey {
1118                info.node1_policy
1119            } else {
1120                info.node2_policy
1121            };
1122            let Some(policy) = policy else {
1123                continue;
1124            };
1125            let src_node_id =
1126                PublicKey::from_str(&chan.remote_pubkey).expect("Failed to parse pubkey");
1127            let short_channel_id = chan.chan_id;
1128            let base_msat = policy.fee_base_msat as u32;
1129            let proportional_millionths = policy.fee_rate_milli_msat as u32;
1130            let cltv_expiry_delta = policy.time_lock_delta;
1131            let htlc_maximum_msat = Some(policy.max_htlc_msat);
1132            let htlc_minimum_msat = Some(policy.min_htlc as u64);
1133
1134            let route_hint_hop = RouteHintHop {
1135                src_node_id,
1136                short_channel_id,
1137                base_msat,
1138                proportional_millionths,
1139                cltv_expiry_delta: cltv_expiry_delta as u16,
1140                htlc_minimum_msat,
1141                htlc_maximum_msat,
1142            };
1143            route_hints.push(RouteHint(vec![route_hint_hop]));
1144        }
1145
1146        Ok(GetRouteHintsResponse { route_hints })
1147    }
1148
1149    async fn pay_private(
1150        &self,
1151        invoice: PrunedInvoice,
1152        max_delay: u64,
1153        max_fee: Amount,
1154    ) -> Result<PayInvoiceResponse, LightningRpcError> {
1155        let payment_hash = invoice.payment_hash.to_byte_array().to_vec();
1156        info!(
1157            target: LOG_LIGHTNING,
1158            payment_hash = %PrettyPaymentHash(&payment_hash),
1159            "LND Paying invoice",
1160        );
1161        let mut client = self.connect().await?;
1162
1163        debug!(
1164            target: LOG_LIGHTNING,
1165            payment_hash = %PrettyPaymentHash(&payment_hash),
1166            "pay_private checking if payment for invoice exists"
1167        );
1168
1169        // If the payment exists, that means we've already tried to pay the invoice
1170        let preimage: Vec<u8> = match self
1171            .lookup_payment(invoice.payment_hash.to_byte_array().to_vec(), &mut client)
1172            .await?
1173        {
1174            Some(preimage) => {
1175                info!(
1176                    target: LOG_LIGHTNING,
1177                    payment_hash = %PrettyPaymentHash(&payment_hash),
1178                    "LND payment already exists for invoice",
1179                );
1180                hex::FromHex::from_hex(preimage.as_str()).map_err(|error| {
1181                    LightningRpcError::FailedPayment {
1182                        failure_reason: format!("Failed to convert preimage {error:?}"),
1183                    }
1184                })?
1185            }
1186            _ => 'dispatch: {
1187                // LND API allows fee limits in the `i64` range, but we use `u64` for
1188                // max_fee_msat. This means we can only set an enforceable fee limit
1189                // between 0 and i64::MAX
1190                let fee_limit_msat: i64 =
1191                    max_fee
1192                        .msats
1193                        .try_into()
1194                        .map_err(|error| LightningRpcError::FailedPayment {
1195                            failure_reason: format!(
1196                                "max_fee_msat exceeds valid LND fee limit ranges {error:?}"
1197                            ),
1198                        })?;
1199
1200                let amt_msat = invoice.amount.msats.try_into().map_err(|error| {
1201                    LightningRpcError::FailedPayment {
1202                        failure_reason: format!("amount exceeds valid LND amount ranges {error:?}"),
1203                    }
1204                })?;
1205                let final_cltv_delta =
1206                    invoice.min_final_cltv_delta.try_into().map_err(|error| {
1207                        LightningRpcError::FailedPayment {
1208                            failure_reason: format!(
1209                                "final cltv delta exceeds valid LND range {error:?}"
1210                            ),
1211                        }
1212                    })?;
1213                // LND reads a `cltv_limit` of zero as "no limit set" and
1214                // enforces its `--max-cltv-expiry` default instead, silently
1215                // lifting the caller's timelock cap.
1216                if max_delay == 0 {
1217                    return Err(LightningRpcError::FailedPayment {
1218                        failure_reason: "a max delay of zero would disable LND's CLTV limit"
1219                            .to_string(),
1220                    });
1221                }
1222                let cltv_limit =
1223                    max_delay
1224                        .try_into()
1225                        .map_err(|error| LightningRpcError::FailedPayment {
1226                            failure_reason: format!("max delay exceeds valid LND range {error:?}"),
1227                        })?;
1228
1229                let dest_features = wire_features_to_lnd_feature_vec(&invoice.destination_features)
1230                    .map_err(|e| LightningRpcError::FailedPayment {
1231                        failure_reason: e.to_string(),
1232                    })?;
1233
1234                debug!(
1235                    target: LOG_LIGHTNING,
1236                    payment_hash = %PrettyPaymentHash(&payment_hash),
1237                    "LND payment does not exist, will attempt to pay",
1238                );
1239                let payments = match client
1240                    .router()
1241                    .send_payment_v2(SendPaymentRequest {
1242                        amt_msat,
1243                        dest: invoice.destination.serialize().to_vec(),
1244                        dest_features,
1245                        payment_hash: invoice.payment_hash.to_byte_array().to_vec(),
1246                        payment_addr: invoice.payment_secret.to_vec(),
1247                        route_hints: route_hints_to_lnd(&invoice.route_hints),
1248                        final_cltv_delta,
1249                        cltv_limit,
1250                        no_inflight_updates: false,
1251                        timeout_seconds: self.payment_timeout_secs,
1252                        fee_limit_msat,
1253                        time_pref: self.time_pref,
1254                        ..Default::default()
1255                    })
1256                    .await
1257                {
1258                    Ok(payments) => payments,
1259                    // An error here is not proof that nothing was paid. LND
1260                    // registers the payment before its first status update,
1261                    // so a transport fault on the way back can surface as an
1262                    // error for a payment that is now in flight. Reporting
1263                    // failure would forfeit the outgoing contract while the
1264                    // payment may still settle, so only the node's own record
1265                    // decides.
1266                    Err(status) => {
1267                        warn!(
1268                            target: LOG_LIGHTNING,
1269                            status = %status,
1270                            payment_hash = %PrettyPaymentHash(&payment_hash),
1271                            "LND payment request failed, checking whether LND registered the payment",
1272                        );
1273
1274                        let Some(preimage) = self
1275                            .lookup_payment(payment_hash.clone(), &mut client)
1276                            .await?
1277                        else {
1278                            return Err(LightningRpcError::FailedPayment {
1279                                failure_reason: format!(
1280                                    "Failed to make outgoing payment {status:?}"
1281                                ),
1282                            });
1283                        };
1284
1285                        break 'dispatch hex::FromHex::from_hex(preimage.as_str()).map_err(
1286                            |error| LightningRpcError::FailedPayment {
1287                                failure_reason: format!("Failed to convert preimage {error:?}"),
1288                            },
1289                        )?;
1290                    }
1291                };
1292
1293                debug!(
1294                    target: LOG_LIGHTNING,
1295                    payment_hash = %PrettyPaymentHash(&payment_hash),
1296                    "LND payment request sent, waiting for payment status...",
1297                );
1298                let mut messages = payments.into_inner();
1299                loop {
1300                    match messages.message().await {
1301                        Ok(Some(payment)) if payment.status() == PaymentStatus::Succeeded => {
1302                            info!(
1303                                target: LOG_LIGHTNING,
1304                                payment_hash = %PrettyPaymentHash(&payment_hash),
1305                                "LND payment succeeded for invoice",
1306                            );
1307                            break hex::FromHex::from_hex(payment.payment_preimage.as_str())
1308                                .map_err(|error| LightningRpcError::FailedPayment {
1309                                    failure_reason: format!("Failed to convert preimage {error:?}"),
1310                                })?;
1311                        }
1312                        Ok(Some(payment)) if payment.status() == PaymentStatus::Failed => {
1313                            // The one terminal status besides `Succeeded`: a
1314                            // definitive failure that is safe to report.
1315                            warn!(
1316                                target: LOG_LIGHTNING,
1317                                payment_hash = %PrettyPaymentHash(&payment_hash),
1318                                status = ?payment.status(),
1319                                "LND payment failed",
1320                            );
1321                            let failure_reason = payment.failure_reason();
1322                            return Err(LightningRpcError::FailedPayment {
1323                                failure_reason: format!("{failure_reason:?}"),
1324                            });
1325                        }
1326                        // `InFlight`, `Initiated` (delivered before the first HTLC
1327                        // when `routerrpc.usestatusinitiated` is set) and any status
1328                        // this build does not know, which prost decodes as `Unknown`,
1329                        // are not outcomes. Keep waiting; a stream that ends without
1330                        // a terminal status resumes through `lookup_payment` below.
1331                        Ok(Some(payment)) => {
1332                            debug!(
1333                                target: LOG_LIGHTNING,
1334                                payment_hash = %PrettyPaymentHash(&payment_hash),
1335                                status = ?payment.status(),
1336                                "LND payment is in flight",
1337                            );
1338                            continue;
1339                        }
1340                        // `Ok(None)` is a premature end of the update stream and `Err`
1341                        // is a tonic/HTTP2 transport fault. Neither is a payment
1342                        // outcome: the HTLC may still settle, so reporting failure here
1343                        // would forfeit the outgoing contract for a payment that is
1344                        // still in flight, violating the idempotency contract on
1345                        // `ILnRpcClient::pay`. Resume tracking with `track_payment_v2`
1346                        // to await the real terminal result instead.
1347                        stream_end_or_fault => {
1348                            warn!(
1349                                target: LOG_LIGHTNING,
1350                                payment_hash = %PrettyPaymentHash(&payment_hash),
1351                                outcome = ?stream_end_or_fault,
1352                                "LND payment status stream ended or faulted before a terminal status; resuming tracking",
1353                            );
1354                            match self
1355                                .lookup_payment(payment_hash.clone(), &mut client)
1356                                .await?
1357                            {
1358                                Some(preimage) => {
1359                                    break hex::FromHex::from_hex(preimage.as_str()).map_err(
1360                                        |error| LightningRpcError::FailedPayment {
1361                                            failure_reason: format!(
1362                                                "Failed to convert preimage {error:?}"
1363                                            ),
1364                                        },
1365                                    )?;
1366                                }
1367                                None => {
1368                                    return Err(LightningRpcError::FailedPayment {
1369                                        failure_reason: format!(
1370                                            "LND has no record of dispatched payment for hash {:?}",
1371                                            invoice.payment_hash
1372                                        ),
1373                                    });
1374                                }
1375                            }
1376                        }
1377                    }
1378                }
1379            }
1380        };
1381        Ok(PayInvoiceResponse {
1382            preimage: Preimage(preimage.try_into().expect("Failed to create preimage")),
1383        })
1384    }
1385
1386    /// Returns true if the lightning backend supports payments without full
1387    /// invoices
1388    fn supports_private_payments(&self) -> bool {
1389        true
1390    }
1391
1392    async fn outbound_payment_exists(
1393        &self,
1394        payment_hash: sha256::Hash,
1395    ) -> Result<bool, LightningRpcError> {
1396        let payment_hash_bytes = payment_hash.to_byte_array().to_vec();
1397        let mut client = self.connect().await?;
1398
1399        // Subscribe with in-flight updates enabled so any known payment,
1400        // pending or terminal, yields an immediate first message instead of
1401        // blocking until the payment resolves; an unknown hash fails with
1402        // `NotFound`.
1403        let stream = match client
1404            .router()
1405            .track_payment_v2(TrackPaymentRequest {
1406                payment_hash: payment_hash_bytes.clone(),
1407                no_inflight_updates: false,
1408            })
1409            .await
1410        {
1411            Ok(stream) => stream,
1412            Err(status) if status.code() == Code::NotFound => return Ok(false),
1413            Err(status) => {
1414                return Err(LightningRpcError::FailedPayment {
1415                    failure_reason: format!(
1416                        "Failed to look up payment {}: {status:?}",
1417                        PrettyPaymentHash(&payment_hash_bytes),
1418                    ),
1419                });
1420            }
1421        };
1422
1423        match stream.into_inner().message().await {
1424            Ok(Some(_)) => Ok(true),
1425            Err(status) if status.code() == Code::NotFound => Ok(false),
1426            // A premature end of stream or a transport fault is not an
1427            // answer. Report an error so the caller retries, rather than
1428            // letting it mistake the payment for never having been
1429            // dispatched.
1430            outcome => Err(LightningRpcError::FailedPayment {
1431                failure_reason: format!(
1432                    "Payment lookup stream gave no answer for {}: {outcome:?}",
1433                    PrettyPaymentHash(&payment_hash_bytes),
1434                ),
1435            }),
1436        }
1437    }
1438
1439    async fn route_htlcs<'a>(
1440        self: Box<Self>,
1441        task_group: &TaskGroup,
1442    ) -> Result<(RouteHtlcStream<'a>, Arc<dyn ILnRpcClient>), LightningRpcError> {
1443        const CHANNEL_SIZE: usize = 100;
1444
1445        // Channel to send intercepted htlc to the gateway for processing
1446        let (gateway_sender, gateway_receiver) =
1447            mpsc::channel::<InterceptPaymentRequest>(CHANNEL_SIZE);
1448
1449        let (lnd_sender, lnd_rx) = mpsc::channel::<ForwardHtlcInterceptResponse>(CHANNEL_SIZE);
1450
1451        self.spawn_interceptor(
1452            task_group,
1453            lnd_sender.clone(),
1454            lnd_rx,
1455            gateway_sender.clone(),
1456        )
1457        .await?;
1458        let new_client = Arc::new(Self {
1459            address: self.address.clone(),
1460            tls_cert: self.tls_cert.clone(),
1461            macaroon: self.macaroon.clone(),
1462            time_pref: self.time_pref,
1463            payment_timeout_secs: self.payment_timeout_secs,
1464            lnd_sender: Some(lnd_sender.clone()),
1465            lnv2_filter: self.lnv2_filter.clone(),
1466        });
1467        Ok((Box::pin(ReceiverStream::new(gateway_receiver)), new_client))
1468    }
1469
1470    async fn complete_htlc(&self, htlc: InterceptPaymentResponse) -> Result<(), LightningRpcError> {
1471        let incoming_circuit = htlc.incoming_circuit();
1472        let InterceptPaymentResponse {
1473            action,
1474            payment_hash,
1475            incoming_chan_id: _,
1476            htlc_id: _,
1477        } = htlc;
1478
1479        let (action, preimage) = match action {
1480            PaymentAction::Settle(preimage) => (ResolveHoldForwardAction::Settle, preimage),
1481            PaymentAction::Cancel => (ResolveHoldForwardAction::Fail, Preimage([0; 32])),
1482            PaymentAction::Forward => (ResolveHoldForwardAction::Resume, Preimage([0; 32])),
1483        };
1484
1485        // Resolve the payment the way it arrived. Deciding instead by probing
1486        // LND for a HOLD invoice carrying the payment hash would conflate the
1487        // two ways, because the hash is chosen by whoever is being paid: an
1488        // attacker can register an LNv2 receive and an LNv1 offer for the same
1489        // hash, and the completion for the intercepted LNv1 HTLC would then
1490        // settle or cancel the unrelated LNv2 HOLD invoice.
1491        let Some((chan_id, htlc_id)) = incoming_circuit else {
1492            // LNv2: the payment is held by a HOLD invoice on our own node, so
1493            // there is no forward to resolve.
1494            return match action {
1495                ResolveHoldForwardAction::Settle => {
1496                    self.settle_hold_invoice(payment_hash.to_byte_array().to_vec(), preimage)
1497                        .await
1498                }
1499                // Neither `Fail` nor `Resume` has a meaning for a HOLD invoice
1500                // beyond "the gateway could not claim this payment": there is
1501                // no next hop to resume towards, so fail it back to the payer.
1502                _ => {
1503                    self.cancel_hold_invoice(payment_hash.to_byte_array().to_vec())
1504                        .await
1505                }
1506            };
1507        };
1508
1509        // LNv1: hand the interceptor its response for this exact circuit.
1510        let Some(lnd_sender) = self.lnd_sender.clone() else {
1511            crit!("Gatewayd has not started to route HTLCs");
1512            return Err(LightningRpcError::FailedToCompleteHtlc {
1513                failure_reason: "Gatewayd has not started to route HTLCs".to_string(),
1514            });
1515        };
1516
1517        let response = ForwardHtlcInterceptResponse {
1518            incoming_circuit_key: Some(CircuitKey { chan_id, htlc_id }),
1519            action: action.into(),
1520            preimage: preimage.0.to_vec(),
1521            failure_message: vec![],
1522            failure_code: FailureCode::TemporaryChannelFailure.into(),
1523            ..Default::default()
1524        };
1525
1526        Self::send_lnd_response(lnd_sender, response).await
1527    }
1528
1529    async fn create_invoice(
1530        &self,
1531        create_invoice_request: CreateInvoiceRequest,
1532    ) -> Result<CreateInvoiceResponse, LightningRpcError> {
1533        let mut client = self.connect().await?;
1534        let description = create_invoice_request
1535            .description
1536            .unwrap_or(InvoiceDescription::Direct(String::new()));
1537
1538        if let Some(payment_hash_value) = create_invoice_request.payment_hash {
1539            let payment_hash = payment_hash_value.to_byte_array().to_vec();
1540            let hold_invoice_request = match description {
1541                InvoiceDescription::Direct(description) => AddHoldInvoiceRequest {
1542                    memo: description,
1543                    hash: payment_hash.clone(),
1544                    value_msat: create_invoice_request.amount_msat as i64,
1545                    expiry: i64::from(create_invoice_request.expiry_secs),
1546                    cltv_expiry: LNV2_HOLD_INVOICE_CLTV_EXPIRY,
1547                    ..Default::default()
1548                },
1549                InvoiceDescription::Hash(desc_hash) => AddHoldInvoiceRequest {
1550                    description_hash: desc_hash.to_byte_array().to_vec(),
1551                    hash: payment_hash.clone(),
1552                    value_msat: create_invoice_request.amount_msat as i64,
1553                    expiry: i64::from(create_invoice_request.expiry_secs),
1554                    cltv_expiry: LNV2_HOLD_INVOICE_CLTV_EXPIRY,
1555                    ..Default::default()
1556                },
1557            };
1558
1559            let hold_invoice_response = client
1560                .invoices()
1561                .add_hold_invoice(hold_invoice_request)
1562                .await
1563                .map_err(|e| LightningRpcError::FailedToGetInvoice {
1564                    failure_reason: e.to_string(),
1565                })?;
1566
1567            let invoice = hold_invoice_response.into_inner().payment_request;
1568            Ok(CreateInvoiceResponse { invoice })
1569        } else {
1570            let invoice = match description {
1571                InvoiceDescription::Direct(description) => Invoice {
1572                    memo: description,
1573                    value_msat: create_invoice_request.amount_msat as i64,
1574                    expiry: i64::from(create_invoice_request.expiry_secs),
1575                    ..Default::default()
1576                },
1577                InvoiceDescription::Hash(desc_hash) => Invoice {
1578                    description_hash: desc_hash.to_byte_array().to_vec(),
1579                    value_msat: create_invoice_request.amount_msat as i64,
1580                    expiry: i64::from(create_invoice_request.expiry_secs),
1581                    ..Default::default()
1582                },
1583            };
1584
1585            let add_invoice_response =
1586                client.lightning().add_invoice(invoice).await.map_err(|e| {
1587                    LightningRpcError::FailedToGetInvoice {
1588                        failure_reason: e.to_string(),
1589                    }
1590                })?;
1591
1592            let invoice = add_invoice_response.into_inner().payment_request;
1593            Ok(CreateInvoiceResponse { invoice })
1594        }
1595    }
1596
1597    async fn get_ln_onchain_address(
1598        &self,
1599    ) -> Result<GetLnOnchainAddressResponse, LightningRpcError> {
1600        let mut client = self.connect().await?;
1601
1602        match client
1603            .wallet()
1604            .next_addr(AddrRequest {
1605                account: String::new(), // Default wallet account.
1606                r#type: 4,              // Taproot address.
1607                change: false,
1608            })
1609            .await
1610        {
1611            Ok(response) => Ok(GetLnOnchainAddressResponse {
1612                address: response.into_inner().addr,
1613            }),
1614            Err(e) => Err(LightningRpcError::FailedToGetLnOnchainAddress {
1615                failure_reason: format!("Failed to get funding address {e:?}"),
1616            }),
1617        }
1618    }
1619
1620    async fn send_onchain(
1621        &self,
1622        SendOnchainRequest {
1623            address,
1624            amount,
1625            fee_rate_sats_per_vbyte,
1626        }: SendOnchainRequest,
1627    ) -> Result<SendOnchainResponse, LightningRpcError> {
1628        #[allow(deprecated)]
1629        let request = match amount {
1630            BitcoinAmountOrAll::All => SendCoinsRequest {
1631                addr: address.assume_checked().to_string(),
1632                amount: 0,
1633                target_conf: 0,
1634                sat_per_vbyte: fee_rate_sats_per_vbyte,
1635                sat_per_byte: 0,
1636                send_all: true,
1637                label: String::new(),
1638                min_confs: 0,
1639                spend_unconfirmed: true,
1640                ..Default::default()
1641            },
1642            BitcoinAmountOrAll::Amount(amount) => SendCoinsRequest {
1643                addr: address.assume_checked().to_string(),
1644                amount: amount.to_sat() as i64,
1645                target_conf: 0,
1646                sat_per_vbyte: fee_rate_sats_per_vbyte,
1647                sat_per_byte: 0,
1648                send_all: false,
1649                label: String::new(),
1650                min_confs: 0,
1651                spend_unconfirmed: true,
1652                ..Default::default()
1653            },
1654        };
1655
1656        match self.connect().await?.lightning().send_coins(request).await {
1657            Ok(res) => Ok(SendOnchainResponse {
1658                txid: res.into_inner().txid,
1659            }),
1660            Err(e) => Err(LightningRpcError::FailedToWithdrawOnchain {
1661                failure_reason: format!("Failed to withdraw funds on-chain {e:?}"),
1662            }),
1663        }
1664    }
1665
1666    async fn open_channel(
1667        &self,
1668        crate::OpenChannelRequest {
1669            pubkey,
1670            host,
1671            channel_size_sats,
1672            push_amount_sats,
1673            fee_rate_sats_per_vbyte,
1674            base_fee_msat,
1675            parts_per_million,
1676        }: crate::OpenChannelRequest,
1677    ) -> Result<OpenChannelResponse, LightningRpcError> {
1678        let mut client = self.connect().await?;
1679
1680        self.connect_peer_if_needed(&mut client, pubkey, host)
1681            .await?;
1682
1683        // Build the request, leaving unspecified fee fields at their
1684        // protobuf defaults so LND falls back to its own configuration.
1685        let mut open_request = OpenChannelRequest {
1686            node_pubkey: pubkey.serialize().to_vec(),
1687            local_funding_amount: channel_size_sats.try_into().expect("u64 -> i64"),
1688            push_sat: push_amount_sats.try_into().expect("u64 -> i64"),
1689            ..Default::default()
1690        };
1691        if let Some(rate) = fee_rate_sats_per_vbyte {
1692            open_request.sat_per_vbyte = rate;
1693        }
1694        if let Some(base_fee) = base_fee_msat {
1695            open_request.base_fee = base_fee;
1696            open_request.use_base_fee = true;
1697        }
1698        if let Some(ppm) = parts_per_million {
1699            open_request.fee_rate = ppm;
1700            open_request.use_fee_rate = true;
1701        }
1702
1703        // Open the channel
1704        match client.lightning().open_channel_sync(open_request).await {
1705            Ok(res) => Ok(OpenChannelResponse {
1706                funding_txid: match res.into_inner().funding_txid {
1707                    Some(txid) => match txid {
1708                        FundingTxid::FundingTxidBytes(mut bytes) => {
1709                            bytes.reverse();
1710                            hex::encode(bytes)
1711                        }
1712                        FundingTxid::FundingTxidStr(str) => str,
1713                    },
1714                    None => String::new(),
1715                },
1716            }),
1717            Err(e) => Err(LightningRpcError::FailedToOpenChannel {
1718                failure_reason: format!("Failed to open channel {e:?}"),
1719            }),
1720        }
1721    }
1722
1723    async fn connect_peer(&self, payload: ConnectPeerRequest) -> Result<(), LightningRpcError> {
1724        let mut client = self.connect().await?;
1725        self.connect_peer_if_needed(
1726            &mut client,
1727            payload.node_address.pubkey,
1728            payload.node_address.host_with_port(),
1729        )
1730        .await
1731    }
1732
1733    async fn close_channels_with_peer(
1734        &self,
1735        CloseChannelsWithPeerRequest {
1736            pubkey,
1737            force,
1738            sats_per_vbyte,
1739        }: CloseChannelsWithPeerRequest,
1740    ) -> Result<CloseChannelsWithPeerResponse, LightningRpcError> {
1741        let mut client = self.connect().await?;
1742
1743        let channels_with_peer = client
1744            .lightning()
1745            .list_channels(ListChannelsRequest {
1746                active_only: false,
1747                inactive_only: false,
1748                public_only: false,
1749                private_only: false,
1750                peer: pubkey.serialize().to_vec(),
1751                peer_alias_lookup: false,
1752            })
1753            .await
1754            .map_err(|e| LightningRpcError::FailedToCloseChannelsWithPeer {
1755                failure_reason: format!("Failed to list channels {e:?}"),
1756            })?
1757            .into_inner()
1758            .channels;
1759
1760        for channel in &channels_with_peer {
1761            let channel_point =
1762                bitcoin::OutPoint::from_str(&channel.channel_point).map_err(|e| {
1763                    LightningRpcError::FailedToCloseChannelsWithPeer {
1764                        failure_reason: format!("Failed to parse channel point {e:?}"),
1765                    }
1766                })?;
1767
1768            if force {
1769                client
1770                    .lightning()
1771                    .close_channel(CloseChannelRequest {
1772                        channel_point: Some(ChannelPoint {
1773                            funding_txid: Some(
1774                                tonic_lnd::lnrpc::channel_point::FundingTxid::FundingTxidBytes(
1775                                    <bitcoin::Txid as AsRef<[u8]>>::as_ref(&channel_point.txid)
1776                                        .to_vec(),
1777                                ),
1778                            ),
1779                            output_index: channel_point.vout,
1780                        }),
1781                        force,
1782                        ..Default::default()
1783                    })
1784                    .await
1785                    .map_err(|e| LightningRpcError::FailedToCloseChannelsWithPeer {
1786                        failure_reason: format!("Failed to close channel {e:?}"),
1787                    })?;
1788            } else {
1789                client
1790                    .lightning()
1791                    .close_channel(CloseChannelRequest {
1792                        channel_point: Some(ChannelPoint {
1793                            funding_txid: Some(
1794                                tonic_lnd::lnrpc::channel_point::FundingTxid::FundingTxidBytes(
1795                                    <bitcoin::Txid as AsRef<[u8]>>::as_ref(&channel_point.txid)
1796                                        .to_vec(),
1797                                ),
1798                            ),
1799                            output_index: channel_point.vout,
1800                        }),
1801                        force,
1802                        sat_per_vbyte: sats_per_vbyte.unwrap_or_default(),
1803                        ..Default::default()
1804                    })
1805                    .await
1806                    .map_err(|e| LightningRpcError::FailedToCloseChannelsWithPeer {
1807                        failure_reason: format!("Failed to close channel {e:?}"),
1808                    })?;
1809            }
1810        }
1811
1812        Ok(CloseChannelsWithPeerResponse {
1813            num_channels_closed: channels_with_peer.len() as u32,
1814        })
1815    }
1816
1817    async fn list_channels(&self) -> Result<ListChannelsResponse, LightningRpcError> {
1818        let mut client = self.connect().await?;
1819
1820        // Fetch peer addresses so we can populate remote_address on each channel
1821        let peer_addresses: BTreeMap<String, String> = client
1822            .lightning()
1823            .list_peers(ListPeersRequest {
1824                latest_error: false,
1825            })
1826            .await
1827            .map(|resp| {
1828                resp.into_inner()
1829                    .peers
1830                    .into_iter()
1831                    .filter_map(|peer| {
1832                        if peer.address.is_empty() {
1833                            None
1834                        } else {
1835                            Some((peer.pub_key, peer.address))
1836                        }
1837                    })
1838                    .collect()
1839            })
1840            .unwrap_or_default();
1841
1842        // Fetch the local fee policy for every channel in one call so we can
1843        // join it onto each ChannelInfo below.
1844        let fee_report: BTreeMap<u64, (u64, u64)> = client
1845            .lightning()
1846            .fee_report(FeeReportRequest {})
1847            .await
1848            .map(|resp| {
1849                resp.into_inner()
1850                    .channel_fees
1851                    .into_iter()
1852                    .map(|report| {
1853                        let base_fee_msat = u64::try_from(report.base_fee_msat).unwrap_or_default();
1854                        let parts_per_million =
1855                            u64::try_from(report.fee_per_mil).unwrap_or_default();
1856                        (report.chan_id, (base_fee_msat, parts_per_million))
1857                    })
1858                    .collect()
1859            })
1860            .unwrap_or_default();
1861
1862        match client
1863            .lightning()
1864            .list_channels(ListChannelsRequest {
1865                active_only: false,
1866                inactive_only: false,
1867                public_only: false,
1868                private_only: false,
1869                peer: vec![],
1870                peer_alias_lookup: true,
1871            })
1872            .await
1873        {
1874            Ok(response) => Ok(ListChannelsResponse {
1875                channels: response
1876                    .into_inner()
1877                    .channels
1878                    .into_iter()
1879                    .map(|channel| {
1880                        let channel_size_sats = channel.capacity.try_into().expect("i64 -> u64");
1881
1882                        let local_balance_sats: u64 =
1883                            channel.local_balance.try_into().expect("i64 -> u64");
1884                        let local_channel_reserve_sats: u64 = match channel.local_constraints {
1885                            Some(constraints) => constraints.chan_reserve_sat,
1886                            None => 0,
1887                        };
1888
1889                        let outbound_liquidity_sats =
1890                            local_balance_sats.saturating_sub(local_channel_reserve_sats);
1891
1892                        let remote_balance_sats: u64 =
1893                            channel.remote_balance.try_into().expect("i64 -> u64");
1894                        let remote_channel_reserve_sats: u64 = match channel.remote_constraints {
1895                            Some(constraints) => constraints.chan_reserve_sat,
1896                            None => 0,
1897                        };
1898
1899                        let inbound_liquidity_sats =
1900                            remote_balance_sats.saturating_sub(remote_channel_reserve_sats);
1901
1902                        let funding_outpoint = OutPoint::from_str(&channel.channel_point).ok();
1903
1904                        let remote_address = peer_addresses.get(&channel.remote_pubkey).cloned();
1905
1906                        let (base_fee_msat, parts_per_million) =
1907                            match fee_report.get(&channel.chan_id) {
1908                                Some((base, ppm)) => (Some(*base), Some(*ppm)),
1909                                None => (None, None),
1910                            };
1911
1912                        ChannelInfo {
1913                            remote_pubkey: PublicKey::from_str(&channel.remote_pubkey)
1914                                .expect("Lightning node returned invalid remote channel pubkey"),
1915                            channel_size_sats,
1916                            outbound_liquidity_sats,
1917                            inbound_liquidity_sats,
1918                            is_active: channel.active,
1919                            funding_outpoint,
1920                            remote_node_alias: if channel.peer_alias.is_empty() {
1921                                None
1922                            } else {
1923                                Some(channel.peer_alias.clone())
1924                            },
1925                            remote_address,
1926                            base_fee_msat,
1927                            parts_per_million,
1928                        }
1929                    })
1930                    .collect(),
1931            }),
1932            Err(e) => Err(LightningRpcError::FailedToListChannels {
1933                failure_reason: format!("Failed to list active channels {e:?}"),
1934            }),
1935        }
1936    }
1937
1938    async fn set_channel_fees(
1939        &self,
1940        payload: SetChannelFeesRequest,
1941    ) -> Result<(), LightningRpcError> {
1942        let mut client = self.connect().await?;
1943
1944        // LND's `PolicyUpdateRequest` applies every field it receives, so we
1945        // need the channel's current `time_lock_delta` (and htlc min/max) to
1946        // avoid clobbering them when only base + ppm are being changed. To
1947        // look those up we first resolve the funding outpoint to LND's
1948        // numeric `chan_id`, then call `get_chan_info`.
1949        let target = format!(
1950            "{}:{}",
1951            payload.funding_outpoint.txid, payload.funding_outpoint.vout
1952        );
1953        let channel = client
1954            .lightning()
1955            .list_channels(ListChannelsRequest::default())
1956            .await
1957            .map_err(|e| LightningRpcError::FailedToSetChannelFees {
1958                failure_reason: format!("Failed to list channels: {e:?}"),
1959            })?
1960            .into_inner()
1961            .channels
1962            .into_iter()
1963            .find(|c| c.channel_point == target)
1964            .ok_or_else(|| LightningRpcError::FailedToSetChannelFees {
1965                failure_reason: format!("No channel found with funding outpoint {target}"),
1966            })?;
1967
1968        let our_pubkey = client
1969            .lightning()
1970            .get_info(GetInfoRequest {})
1971            .await
1972            .map_err(|e| LightningRpcError::FailedToSetChannelFees {
1973                failure_reason: format!("Failed to get node info: {e:?}"),
1974            })?
1975            .into_inner()
1976            .identity_pubkey;
1977
1978        let edge = client
1979            .lightning()
1980            .get_chan_info(ChanInfoRequest {
1981                chan_id: channel.chan_id,
1982                ..Default::default()
1983            })
1984            .await
1985            .map_err(|e| LightningRpcError::FailedToSetChannelFees {
1986                failure_reason: format!("Failed to get channel info: {e:?}"),
1987            })?
1988            .into_inner();
1989
1990        // Pick the policy advertised by our node (the local side); fall back
1991        // to node1_policy when neither pubkey matches, which only happens if
1992        // the gossip data has not propagated yet.
1993        let current_policy = if edge.node1_pub == our_pubkey {
1994            edge.node1_policy
1995        } else if edge.node2_pub == our_pubkey {
1996            edge.node2_policy
1997        } else {
1998            edge.node1_policy
1999        };
2000
2001        let fee_rate_ppm = u32::try_from(payload.parts_per_million).map_err(|_| {
2002            LightningRpcError::FailedToSetChannelFees {
2003                failure_reason: format!(
2004                    "parts_per_million {} does not fit in u32",
2005                    payload.parts_per_million,
2006                ),
2007            }
2008        })?;
2009
2010        let base_fee_msat = i64::try_from(payload.base_fee_msat).map_err(|_| {
2011            LightningRpcError::FailedToSetChannelFees {
2012                failure_reason: format!(
2013                    "base_fee_msat {} does not fit in i64",
2014                    payload.base_fee_msat,
2015                ),
2016            }
2017        })?;
2018
2019        // Default time_lock_delta of 40 matches LND's CLI default for
2020        // `lncli updatechanpolicy` when the channel's existing CLTV is not
2021        // discoverable. max_htlc_msat == 0 means "no max" in LND.
2022        let time_lock_delta = current_policy
2023            .as_ref()
2024            .map(|p| p.time_lock_delta)
2025            .unwrap_or(40);
2026        let max_htlc_msat = current_policy
2027            .as_ref()
2028            .map(|p| p.max_htlc_msat)
2029            .unwrap_or(0);
2030        let min_htlc_msat = current_policy
2031            .as_ref()
2032            .map(|p| p.min_htlc as u64)
2033            .unwrap_or(0);
2034
2035        let chan_point = ChannelPoint {
2036            funding_txid: Some(FundingTxid::FundingTxidBytes(
2037                <bitcoin::Txid as AsRef<[u8]>>::as_ref(&payload.funding_outpoint.txid).to_vec(),
2038            )),
2039            output_index: payload.funding_outpoint.vout,
2040        };
2041
2042        let request = PolicyUpdateRequest {
2043            base_fee_msat,
2044            fee_rate_ppm,
2045            time_lock_delta,
2046            max_htlc_msat,
2047            min_htlc_msat,
2048            min_htlc_msat_specified: false,
2049            scope: Some(PolicyUpdateScope::ChanPoint(chan_point)),
2050            ..Default::default()
2051        };
2052
2053        let response = client
2054            .lightning()
2055            .update_channel_policy(request)
2056            .await
2057            .map_err(|e| LightningRpcError::FailedToSetChannelFees {
2058                failure_reason: format!("update_channel_policy failed: {e:?}"),
2059            })?
2060            .into_inner();
2061
2062        if !response.failed_updates.is_empty() {
2063            let details = response
2064                .failed_updates
2065                .iter()
2066                .map(|f| {
2067                    let outpoint = f
2068                        .outpoint
2069                        .as_ref()
2070                        .map(|op| format!("{}:{}", op.txid_str, op.output_index))
2071                        .unwrap_or_else(|| "<unknown outpoint>".to_string());
2072                    let reason = UpdateFailure::try_from(f.reason)
2073                        .map(|r| r.as_str_name())
2074                        .unwrap_or("UPDATE_FAILURE_UNKNOWN");
2075                    format!("{outpoint}: {reason} ({})", f.update_error)
2076                })
2077                .collect::<Vec<_>>()
2078                .join("; ");
2079            return Err(LightningRpcError::FailedToSetChannelFees {
2080                failure_reason: format!("update_channel_policy reported failures: {details}"),
2081            });
2082        }
2083
2084        Ok(())
2085    }
2086
2087    async fn get_balances(&self) -> Result<GetBalancesResponse, LightningRpcError> {
2088        let mut client = self.connect().await?;
2089
2090        let wallet_balance_response = client
2091            .lightning()
2092            .wallet_balance(WalletBalanceRequest {
2093                ..Default::default()
2094            })
2095            .await
2096            .map_err(|e| LightningRpcError::FailedToGetBalances {
2097                failure_reason: format!("Failed to get on-chain balance {e:?}"),
2098            })?
2099            .into_inner();
2100
2101        let channel_balance_response = client
2102            .lightning()
2103            .channel_balance(ChannelBalanceRequest {})
2104            .await
2105            .map_err(|e| LightningRpcError::FailedToGetBalances {
2106                failure_reason: format!("Failed to get lightning balance {e:?}"),
2107            })?
2108            .into_inner();
2109        let total_outbound = channel_balance_response.local_balance.unwrap_or_default();
2110        let unsettled_outbound = channel_balance_response
2111            .unsettled_local_balance
2112            .unwrap_or_default();
2113        let pending_outbound = channel_balance_response
2114            .pending_open_local_balance
2115            .unwrap_or_default();
2116        let lightning_balance_msats = total_outbound
2117            .msat
2118            .saturating_sub(unsettled_outbound.msat)
2119            .saturating_sub(pending_outbound.msat);
2120
2121        let total_inbound = channel_balance_response.remote_balance.unwrap_or_default();
2122        let unsettled_inbound = channel_balance_response
2123            .unsettled_remote_balance
2124            .unwrap_or_default();
2125        let pending_inbound = channel_balance_response
2126            .pending_open_remote_balance
2127            .unwrap_or_default();
2128        let inbound_lightning_liquidity_msats = total_inbound
2129            .msat
2130            .saturating_sub(unsettled_inbound.msat)
2131            .saturating_sub(pending_inbound.msat);
2132
2133        Ok(GetBalancesResponse {
2134            onchain_balance_sats: (wallet_balance_response.total_balance
2135                + wallet_balance_response.reserved_balance_anchor_chan)
2136                as u64,
2137            lightning_balance_msats,
2138            inbound_lightning_liquidity_msats,
2139        })
2140    }
2141
2142    async fn get_invoice(
2143        &self,
2144        get_invoice_request: GetInvoiceRequest,
2145    ) -> Result<Option<GetInvoiceResponse>, LightningRpcError> {
2146        let mut client = self.connect().await?;
2147        let invoice = client
2148            .invoices()
2149            .lookup_invoice_v2(LookupInvoiceMsg {
2150                invoice_ref: Some(InvoiceRef::PaymentHash(
2151                    get_invoice_request.payment_hash.consensus_encode_to_vec(),
2152                )),
2153                ..Default::default()
2154            })
2155            .await;
2156        let invoice = match invoice {
2157            Ok(invoice) => invoice.into_inner(),
2158            Err(_) => return Ok(None),
2159        };
2160        let status = match &invoice.state() {
2161            InvoiceState::Settled => fedimint_gateway_common::PaymentStatus::Succeeded,
2162            InvoiceState::Canceled => fedimint_gateway_common::PaymentStatus::Failed,
2163            _ => fedimint_gateway_common::PaymentStatus::Pending,
2164        };
2165
2166        Ok(Some(GetInvoiceResponse {
2167            preimage: invoice_preimage_hex(&invoice),
2168            payment_hash: Some(
2169                sha256::Hash::from_slice(&invoice.r_hash).expect("Could not convert payment hash"),
2170            ),
2171            amount: Amount::from_msats(invoice.value_msat as u64),
2172            created_at: UNIX_EPOCH + Duration::from_secs(invoice.creation_date as u64),
2173            status,
2174        }))
2175    }
2176
2177    async fn list_transactions(
2178        &self,
2179        start_secs: u64,
2180        end_secs: u64,
2181    ) -> Result<ListTransactionsResponse, LightningRpcError> {
2182        let mut client = self.connect().await?;
2183        let payments = client
2184            .lightning()
2185            .list_payments(ListPaymentsRequest {
2186                // On higher versions on LND, we can filter on the time range directly in the query
2187                ..Default::default()
2188            })
2189            .await
2190            .map_err(|err| LightningRpcError::FailedToListTransactions {
2191                failure_reason: err.to_string(),
2192            })?
2193            .into_inner();
2194
2195        let mut payments = payments
2196            .payments
2197            .iter()
2198            .filter_map(|payment| {
2199                let timestamp_secs = (payment.creation_time_ns / 1_000_000_000) as u64;
2200                if timestamp_secs < start_secs || timestamp_secs >= end_secs {
2201                    return None;
2202                }
2203                let payment_hash = sha256::Hash::from_str(&payment.payment_hash).ok();
2204                let preimage = (!payment.payment_preimage.is_empty())
2205                    .then_some(payment.payment_preimage.clone());
2206                let status = match &payment.status() {
2207                    PaymentStatus::Succeeded => fedimint_gateway_common::PaymentStatus::Succeeded,
2208                    PaymentStatus::Failed => fedimint_gateway_common::PaymentStatus::Failed,
2209                    _ => fedimint_gateway_common::PaymentStatus::Pending,
2210                };
2211                Some(PaymentDetails {
2212                    payment_hash,
2213                    preimage,
2214                    payment_kind: PaymentKind::Bolt11,
2215                    amount: Amount::from_msats(payment.value_msat as u64),
2216                    direction: PaymentDirection::Outbound,
2217                    status,
2218                    timestamp_secs,
2219                })
2220            })
2221            .collect::<Vec<_>>();
2222
2223        let invoices = client
2224            .lightning()
2225            .list_invoices(ListInvoiceRequest {
2226                pending_only: false,
2227                // On higher versions on LND, we can filter on the time range directly in the query
2228                ..Default::default()
2229            })
2230            .await
2231            .map_err(|err| LightningRpcError::FailedToListTransactions {
2232                failure_reason: err.to_string(),
2233            })?
2234            .into_inner();
2235
2236        let mut incoming_payments = invoices
2237            .invoices
2238            .iter()
2239            .filter_map(|invoice| {
2240                let timestamp_secs = invoice.settle_date as u64;
2241                if timestamp_secs < start_secs || timestamp_secs >= end_secs {
2242                    return None;
2243                }
2244                let status = match &invoice.state() {
2245                    InvoiceState::Settled => fedimint_gateway_common::PaymentStatus::Succeeded,
2246                    InvoiceState::Canceled => fedimint_gateway_common::PaymentStatus::Failed,
2247                    _ => return None,
2248                };
2249                let preimage = (!invoice.r_preimage.is_empty())
2250                    .then_some(invoice.r_preimage.encode_hex::<String>());
2251                Some(PaymentDetails {
2252                    payment_hash: Some(
2253                        sha256::Hash::from_slice(&invoice.r_hash)
2254                            .expect("Could not convert payment hash"),
2255                    ),
2256                    preimage,
2257                    payment_kind: PaymentKind::Bolt11,
2258                    amount: Amount::from_msats(invoice.value_msat as u64),
2259                    direction: PaymentDirection::Inbound,
2260                    status,
2261                    timestamp_secs,
2262                })
2263            })
2264            .collect::<Vec<_>>();
2265
2266        payments.append(&mut incoming_payments);
2267        payments.sort_by_key(|p| p.timestamp_secs);
2268
2269        Ok(ListTransactionsResponse {
2270            transactions: payments,
2271        })
2272    }
2273
2274    fn create_offer(
2275        &self,
2276        _amount_msat: Option<Amount>,
2277        _description: Option<String>,
2278        _expiry_secs: Option<u32>,
2279        _quantity: Option<u64>,
2280    ) -> Result<String, LightningRpcError> {
2281        Err(LightningRpcError::Bolt12Error {
2282            failure_reason: "LND Does not support Bolt12".to_string(),
2283        })
2284    }
2285
2286    async fn pay_offer(
2287        &self,
2288        _offer: String,
2289        _quantity: Option<u64>,
2290        _amount: Option<Amount>,
2291        _payer_note: Option<String>,
2292    ) -> Result<Preimage, LightningRpcError> {
2293        Err(LightningRpcError::Bolt12Error {
2294            failure_reason: "LND Does not support Bolt12".to_string(),
2295        })
2296    }
2297
2298    fn sync_wallet(&self) -> Result<(), LightningRpcError> {
2299        // There is nothing explicit needed to do for syncing an LND node
2300        Ok(())
2301    }
2302}
2303
2304fn route_hints_to_lnd(
2305    route_hints: &[fedimint_ln_common::route_hints::RouteHint],
2306) -> Vec<tonic_lnd::lnrpc::RouteHint> {
2307    route_hints
2308        .iter()
2309        .map(|hint| tonic_lnd::lnrpc::RouteHint {
2310            hop_hints: hint
2311                .0
2312                .iter()
2313                .map(|hop| tonic_lnd::lnrpc::HopHint {
2314                    node_id: hop.src_node_id.serialize().encode_hex(),
2315                    chan_id: hop.short_channel_id,
2316                    fee_base_msat: hop.base_msat,
2317                    fee_proportional_millionths: hop.proportional_millionths,
2318                    cltv_expiry_delta: u32::from(hop.cltv_expiry_delta),
2319                })
2320                .collect(),
2321        })
2322        .collect()
2323}
2324
2325fn wire_features_to_lnd_feature_vec(
2326    features_wire_encoded: &[u8],
2327) -> Result<Vec<i32>, FeatureVectorTooLargeError> {
2328    if features_wire_encoded.len() > 1_000 {
2329        return Err(FeatureVectorTooLargeError);
2330    }
2331
2332    let lnd_features = features_wire_encoded
2333        .iter()
2334        .rev()
2335        .enumerate()
2336        .flat_map(|(byte_idx, &feature_byte)| {
2337            (0..8).filter_map(move |bit_idx| {
2338                if (feature_byte & (1u8 << bit_idx)) != 0 {
2339                    Some(
2340                        i32::try_from(byte_idx * 8 + bit_idx)
2341                            .expect("Index will never exceed i32::MAX for feature vectors <8MB"),
2342                    )
2343                } else {
2344                    None
2345                }
2346            })
2347        })
2348        .collect::<Vec<_>>();
2349
2350    Ok(lnd_features)
2351}
2352
2353/// Utility struct for logging payment hashes. Useful for debugging.
2354struct PrettyPaymentHash<'a>(&'a Vec<u8>);
2355
2356impl Display for PrettyPaymentHash<'_> {
2357    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2358        write!(f, "payment_hash={}", self.0.encode_hex::<String>())
2359    }
2360}
2361
2362/// An invoice's destination feature bits are too long to convert for LND.
2363#[derive(Debug, thiserror::Error)]
2364#[error("Will not process feature bit vectors larger than 1000 byte")]
2365struct FeatureVectorTooLargeError;
2366
2367#[cfg(test)]
2368mod tests;