Skip to main content

fedimint_gwv2_client/
lib.rs

1mod api;
2mod complete_sm;
3pub mod events;
4mod receive_sm;
5mod send_sm;
6
7use std::collections::BTreeMap;
8use std::fmt;
9use std::fmt::Debug;
10use std::sync::Arc;
11use std::time::Duration;
12
13use async_trait::async_trait;
14use bitcoin::hashes::sha256;
15use bitcoin::secp256k1::Message;
16use events::{IncomingPaymentStarted, OutgoingPaymentStarted};
17use fedimint_api_client::api::{DynModuleApi, FederationError};
18use fedimint_client::ClientHandleArc;
19use fedimint_client_module::error::ClientModuleError;
20use fedimint_client_module::module::init::{ClientModuleInit, ClientModuleInitArgs};
21use fedimint_client_module::module::recovery::NoModuleBackup;
22use fedimint_client_module::module::{ClientContext, ClientModule, IClientModule, OutPointRange};
23use fedimint_client_module::sm::{Context, DynState, ModuleNotifier, State, StateTransition};
24use fedimint_client_module::transaction::{
25    ClientOutput, ClientOutputBundle, ClientOutputSM, TransactionBuilder,
26};
27use fedimint_client_module::{
28    DynGlobalClientContext, TransactionSubmitError, sm_enum_variant_translation,
29};
30use fedimint_core::config::FederationId;
31use fedimint_core::core::{Decoder, IntoDynInstance, ModuleInstanceId, ModuleKind, OperationId};
32use fedimint_core::db::DatabaseTransaction;
33use fedimint_core::encoding::{Decodable, Encodable};
34use fedimint_core::module::{
35    Amounts, ApiVersion, CommonModuleInit, ModuleCommon, ModuleInit, MultiApiVersion,
36};
37use fedimint_core::secp256k1::Keypair;
38use fedimint_core::task::timeout;
39use fedimint_core::time::now;
40use fedimint_core::util::{FmtCompact, Spanned, backoff_util, retry};
41use fedimint_core::{Amount, OutPoint, PeerId, apply, async_trait_maybe_send, secp256k1};
42use fedimint_lightning::{InterceptPaymentResponse, LightningRpcError};
43use fedimint_lnv2_common::config::LightningClientConfig;
44use fedimint_lnv2_common::contracts::{IncomingContract, PaymentImage};
45use fedimint_lnv2_common::gateway_api::SendPaymentPayload;
46use fedimint_lnv2_common::{
47    LightningCommonInit, LightningInvoice, LightningModuleTypes, LightningOutput, LightningOutputV0,
48};
49use futures::StreamExt;
50use lightning_invoice::Bolt11Invoice;
51use receive_sm::{ReceiveSMState, ReceiveStateMachine};
52use secp256k1::schnorr::Signature;
53use send_sm::{SendSMState, SendStateMachine};
54use serde::{Deserialize, Serialize};
55use thiserror::Error;
56use tpe::{AggregatePublicKey, PublicKeyShare};
57use tracing::{info, warn};
58
59use crate::api::GatewayFederationApi;
60pub use crate::complete_sm::IncomingCircuitKey;
61use crate::complete_sm::{
62    CircuitCompleteSMCommon, CircuitCompleteStateMachine, CompleteSMState, CompleteStateMachine,
63};
64use crate::receive_sm::ReceiveSMCommon;
65use crate::send_sm::SendSMCommon;
66
67/// Bound on the federation liveness probe that gates funding a fresh incoming
68/// contract. A healthy federation answers well within this; the bound only
69/// decides how quickly an unreachable one fails the HTLC back.
70const FEDERATION_LIVENESS_TIMEOUT: Duration = Duration::from_secs(10);
71
72/// Minimum number of blocks between the current height and an incoming HTLC's
73/// claim deadline for the gateway to fund its incoming contract.
74const LNV2_CLAIM_DEADLINE_MARGIN: u32 = 2;
75
76/// LNv2 CLTV Delta in blocks
77pub const EXPIRATION_DELTA_MINIMUM_V2: u64 = 144;
78
79fn incoming_circuit_operation_id(
80    receive_operation_id: OperationId,
81    circuit: IncomingCircuitKey,
82) -> OperationId {
83    OperationId::from_encodable(&(
84        "gateway-lnv2-incoming-circuit",
85        receive_operation_id,
86        circuit,
87    ))
88}
89
90fn is_legacy_completion_for_circuit(
91    state: &GatewayClientStateMachinesV2,
92    circuit: IncomingCircuitKey,
93) -> bool {
94    let IncomingCircuitKey {
95        incoming_chan_id,
96        htlc_id,
97    } = circuit;
98    matches!(
99        state,
100        GatewayClientStateMachinesV2::Complete(CompleteStateMachine {
101            common,
102            ..
103        }) if common.incoming_chan_id == incoming_chan_id && common.htlc_id == htlc_id
104    )
105}
106
107fn legacy_completion_in_states(
108    active: &[GatewayClientStateMachinesV2],
109    inactive: &[GatewayClientStateMachinesV2],
110    circuit: IncomingCircuitKey,
111) -> bool {
112    active
113        .iter()
114        .chain(inactive)
115        .any(|state| is_legacy_completion_for_circuit(state, circuit))
116}
117
118#[derive(Debug, Clone, Copy, Eq, PartialEq)]
119enum IncomingRelayPlan {
120    Replay,
121    AddCompletion,
122    CreateReceiveAndCompletion,
123}
124
125fn incoming_relay_plan(
126    receive_exists: bool,
127    completion_exists: bool,
128    legacy_completion_exists: bool,
129) -> IncomingRelayPlan {
130    if completion_exists || legacy_completion_exists {
131        IncomingRelayPlan::Replay
132    } else if receive_exists {
133        IncomingRelayPlan::AddCompletion
134    } else {
135        IncomingRelayPlan::CreateReceiveAndCompletion
136    }
137}
138
139fn operation_creation_failed_permanently(
140    creation_failed: bool,
141    operation_exists_after_failure: bool,
142) -> bool {
143    creation_failed && !operation_exists_after_failure
144}
145
146/// Identifies the role of an LNv2 gateway operation.
147#[derive(Debug, Clone, Serialize, Deserialize)]
148#[serde(untagged)]
149pub enum GatewayOperationMetaV2 {
150    /// Metadata written before operation roles were distinguished.
151    Legacy(()),
152    /// Metadata written for a role-specific operation.
153    Role {
154        /// Role used when recovering active operations.
155        role: GatewayOperationRoleV2,
156    },
157}
158
159/// Role of an LNv2 gateway operation.
160#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
161#[serde(rename_all = "snake_case")]
162pub enum GatewayOperationRoleV2 {
163    /// An outgoing payment operation.
164    Send,
165    /// A federation receive operation.
166    Receive,
167    /// An incoming Lightning circuit completion operation.
168    CircuitCompletion,
169}
170
171impl GatewayOperationMetaV2 {
172    /// Constructs metadata for `role`.
173    pub fn role(role: GatewayOperationRoleV2) -> Self {
174        Self::Role { role }
175    }
176
177    /// Returns whether shutdown must await Lightning circuit completion.
178    pub fn waits_for_completion(&self) -> bool {
179        matches!(
180            self,
181            Self::Legacy(())
182                | Self::Role {
183                    role: GatewayOperationRoleV2::CircuitCompletion
184                }
185        )
186    }
187}
188
189#[derive(Debug, Clone)]
190pub struct GatewayClientInitV2 {
191    pub gateway: Arc<dyn IGatewayClientV2>,
192}
193
194impl ModuleInit for GatewayClientInitV2 {
195    type Common = LightningCommonInit;
196
197    async fn dump_database(
198        &self,
199        _dbtx: &mut DatabaseTransaction<'_>,
200        _prefix_names: Vec<String>,
201    ) -> Box<dyn Iterator<Item = (String, Box<dyn erased_serde::Serialize + Send>)> + '_> {
202        Box::new(vec![].into_iter())
203    }
204}
205
206#[apply(async_trait_maybe_send!)]
207impl ClientModuleInit for GatewayClientInitV2 {
208    type Module = GatewayClientModuleV2;
209
210    fn supported_api_versions(&self) -> MultiApiVersion {
211        MultiApiVersion::try_from_iter([ApiVersion { major: 0, minor: 0 }])
212            .expect("no version conflicts")
213    }
214
215    async fn init(
216        &self,
217        args: &ClientModuleInitArgs<Self>,
218    ) -> Result<Self::Module, ClientModuleError> {
219        Ok(GatewayClientModuleV2 {
220            federation_id: *args.federation_id(),
221            cfg: args.cfg().clone(),
222            notifier: args.notifier().clone(),
223            client_ctx: args.context(),
224            module_api: args.module_api().clone(),
225            keypair: args
226                .module_root_secret()
227                .clone()
228                .to_secp_key(fedimint_core::secp256k1::SECP256K1),
229            gateway: self.gateway.clone(),
230        })
231    }
232}
233
234#[derive(Debug, Clone)]
235pub struct GatewayClientModuleV2 {
236    pub federation_id: FederationId,
237    pub cfg: LightningClientConfig,
238    pub notifier: ModuleNotifier<GatewayClientStateMachinesV2>,
239    pub client_ctx: ClientContext<Self>,
240    pub module_api: DynModuleApi,
241    pub keypair: Keypair,
242    pub gateway: Arc<dyn IGatewayClientV2>,
243}
244
245#[derive(Debug, Clone)]
246pub struct GatewayClientContextV2 {
247    pub module: GatewayClientModuleV2,
248    pub decoder: Decoder,
249    pub tpe_agg_pk: AggregatePublicKey,
250    pub tpe_pks: BTreeMap<PeerId, PublicKeyShare>,
251    pub gateway: Arc<dyn IGatewayClientV2>,
252}
253
254impl Context for GatewayClientContextV2 {
255    const KIND: Option<ModuleKind> = Some(fedimint_lnv2_common::KIND);
256}
257
258impl ClientModule for GatewayClientModuleV2 {
259    type Init = GatewayClientInitV2;
260    type Common = LightningModuleTypes;
261    type Backup = NoModuleBackup;
262    type ModuleStateMachineContext = GatewayClientContextV2;
263    type States = GatewayClientStateMachinesV2;
264
265    fn context(&self) -> Self::ModuleStateMachineContext {
266        GatewayClientContextV2 {
267            module: self.clone(),
268            decoder: self.decoder(),
269            tpe_agg_pk: self.cfg.tpe_agg_pk,
270            tpe_pks: self.cfg.tpe_pks.clone(),
271            gateway: self.gateway.clone(),
272        }
273    }
274    fn input_fee(
275        &self,
276        amount: &Amounts,
277        _input: &<Self::Common as ModuleCommon>::Input,
278    ) -> Option<Amounts> {
279        Some(Amounts::new_bitcoin(
280            self.cfg.fee_consensus.fee(amount.expect_only_bitcoin()),
281        ))
282    }
283
284    fn output_fee(
285        &self,
286        _amount: &Amounts,
287        output: &<Self::Common as ModuleCommon>::Output,
288    ) -> Option<Amounts> {
289        let amount = match output.ensure_v0_ref().ok()? {
290            LightningOutputV0::Outgoing(contract) => contract.amount,
291            LightningOutputV0::Incoming(contract) => contract.commitment.amount,
292        };
293
294        Some(Amounts::new_bitcoin(self.cfg.fee_consensus.fee(amount)))
295    }
296}
297
298#[derive(Debug, Clone, Eq, PartialEq, Hash, Decodable, Encodable)]
299pub enum GatewayClientStateMachinesV2 {
300    Send(SendStateMachine),
301    Receive(ReceiveStateMachine),
302    /// Legacy completion state retained to decode and resume combined receive
303    /// operations created by older clients.
304    Complete(CompleteStateMachine),
305    /// Completes one incoming circuit independently of other circuits carrying
306    /// the same payment hash and amount.
307    CircuitComplete(CircuitCompleteStateMachine),
308}
309
310impl fmt::Display for GatewayClientStateMachinesV2 {
311    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
312        match self {
313            GatewayClientStateMachinesV2::Send(send) => {
314                write!(f, "{send}")
315            }
316            GatewayClientStateMachinesV2::Receive(receive) => {
317                write!(f, "{receive}")
318            }
319            GatewayClientStateMachinesV2::Complete(complete) => {
320                write!(f, "{complete}")
321            }
322            GatewayClientStateMachinesV2::CircuitComplete(complete) => {
323                write!(f, "{complete}")
324            }
325        }
326    }
327}
328
329impl IntoDynInstance for GatewayClientStateMachinesV2 {
330    type DynType = DynState;
331
332    fn into_dyn(self, instance_id: ModuleInstanceId) -> Self::DynType {
333        DynState::from_typed(instance_id, self)
334    }
335}
336
337impl State for GatewayClientStateMachinesV2 {
338    type ModuleContext = GatewayClientContextV2;
339
340    fn transitions(
341        &self,
342        context: &Self::ModuleContext,
343        global_context: &DynGlobalClientContext,
344    ) -> Vec<StateTransition<Self>> {
345        match self {
346            GatewayClientStateMachinesV2::Send(state) => {
347                sm_enum_variant_translation!(
348                    state.transitions(context, global_context),
349                    GatewayClientStateMachinesV2::Send
350                )
351            }
352            GatewayClientStateMachinesV2::Receive(state) => {
353                sm_enum_variant_translation!(
354                    state.transitions(context, global_context),
355                    GatewayClientStateMachinesV2::Receive
356                )
357            }
358            GatewayClientStateMachinesV2::Complete(state) => {
359                sm_enum_variant_translation!(
360                    state.transitions(context, global_context),
361                    GatewayClientStateMachinesV2::Complete
362                )
363            }
364            GatewayClientStateMachinesV2::CircuitComplete(state) => {
365                sm_enum_variant_translation!(
366                    state.transitions(context, global_context),
367                    GatewayClientStateMachinesV2::CircuitComplete
368                )
369            }
370        }
371    }
372
373    fn operation_id(&self) -> OperationId {
374        match self {
375            GatewayClientStateMachinesV2::Send(state) => state.operation_id(),
376            GatewayClientStateMachinesV2::Receive(state) => state.operation_id(),
377            GatewayClientStateMachinesV2::Complete(state) => state.operation_id(),
378            GatewayClientStateMachinesV2::CircuitComplete(state) => state.operation_id(),
379        }
380    }
381}
382
383#[derive(Debug, Clone, Eq, PartialEq, Hash, Serialize, Deserialize, Decodable, Encodable)]
384pub enum FinalReceiveState {
385    Rejected,
386    Success([u8; 32]),
387    Refunded,
388    Failure,
389}
390
391impl GatewayClientModuleV2 {
392    /// Starts paying the invoice of an LNv2 outgoing contract, or joins the
393    /// payment already under way for it, and waits for its outcome.
394    ///
395    /// The `Ok` value is that outcome: the preimage if the payment succeeded,
396    /// or the gateway's forfeit signature, which lets the sender reclaim the
397    /// contract, if it was cancelled.
398    ///
399    /// # Errors
400    ///
401    /// Fails with a [`GatewaySendPaymentError`] if the request is refused
402    /// before any payment starts: the gateway is not connected to its
403    /// lightning node, the contract is another gateway's, the
404    /// request's signature does not verify, the federation has not confirmed
405    /// this contract at that outpoint or cannot be asked, the invoice has no
406    /// amount or does not match the contract, the gateway cannot price the
407    /// payment, or another contract already claimed the payment image.
408    pub async fn send_payment(
409        &self,
410        payload: SendPaymentPayload,
411    ) -> Result<Result<[u8; 32], Signature>, GatewaySendPaymentError> {
412        let operation_start = now();
413
414        // The operation id is equal to the contract id which also doubles as the
415        // message signed by the gateway via the forfeit signature to forfeit
416        // the gateways claim to a contract in case of cancellation. We only create a
417        // forfeit signature after we have started the send state machine to
418        // prevent replay attacks with a previously cancelled outgoing contract
419        let operation_id = OperationId::from_encodable(&payload.contract.clone());
420
421        // Since the following checks may only fail due to client side
422        // programming error we do not have to enable cancellation and can check
423        // them before we start the state machine.
424        if payload.contract.claim_pk != self.keypair.public_key() {
425            return Err(GatewaySendPaymentError::NotOurContract);
426        }
427
428        // This prevents DOS attacks where an attacker submits a different invoice.
429        if secp256k1::SECP256K1
430            .verify_schnorr(
431                &payload.auth,
432                &Message::from_digest(*payload.invoice.consensus_hash::<sha256::Hash>().as_ref()),
433                &payload.contract.refund_pk.x_only_public_key().0,
434            )
435            .is_err()
436        {
437            return Err(GatewaySendPaymentError::InvalidAuthSignature);
438        }
439
440        // The operation id is derived from the contract, which is public in the
441        // funding transaction, and joining the operation yields its preimage. So
442        // the join belongs behind the signature above.
443        if self.client_ctx.operation_exists(operation_id).await {
444            return Ok(self.subscribe_send(operation_id).await);
445        }
446
447        // A send accepted while the lightning node is unreachable would sit
448        // waiting for it, and nothing about the payment is known until then.
449        // Refusing up front keeps the sender free to use another gateway. Only
450        // new sends are refused: joining an existing operation above resumes
451        // one already started.
452        if !self.gateway.is_lightning_connected().await {
453            return Err(GatewaySendPaymentError::LightningNotConnected);
454        }
455
456        // We need to check that the contract has been confirmed by the federation
457        // before we start the state machine to prevent DOS attacks.
458        let (contract_id, expiration) = self
459            .module_api
460            .outgoing_contract_expiration(payload.outpoint)
461            .await?
462            .ok_or(GatewaySendPaymentError::ContractNotConfirmed)?;
463
464        if contract_id != payload.contract.contract_id() {
465            return Err(GatewaySendPaymentError::ContractIdMismatch);
466        }
467
468        let (payment_hash, amount) = match &payload.invoice {
469            LightningInvoice::Bolt11(invoice) => (
470                invoice.payment_hash(),
471                invoice
472                    .amount_milli_satoshis()
473                    .ok_or(GatewaySendPaymentError::MissingInvoiceAmount)?,
474            ),
475        };
476
477        if PaymentImage::Hash(*payment_hash) != payload.contract.payment_image {
478            return Err(GatewaySendPaymentError::PaymentHashMismatch);
479        }
480
481        let min_contract_amount = self
482            .gateway
483            .min_contract_amount(&payload.federation_id, amount)
484            .await
485            .map_err(GatewaySendPaymentError::MinContractAmount)?;
486
487        // Contracts for different invoices, on LNv1 or in other federations, can
488        // share a payment image, and paying one pays them all. So only the first
489        // operation to claim the image may ever pay it out; the claim is never
490        // released. It is taken before the state machine exists, so a refused
491        // contract has not been paid for. A retry of this operation finds its own
492        // claim.
493        //
494        // TODO(joschisan): review whether answering with a forfeit signature is
495        // safe here, so the sender is refunded immediately instead of at the
496        // contract's expiration. It should be: a refused contract can never take
497        // the claim later, so we never pay out on it.
498        if !self
499            .gateway
500            .claim_payment_image(&payload.contract.payment_image, operation_id)
501            .await
502        {
503            return Err(GatewaySendPaymentError::PaymentImageAlreadyClaimed);
504        }
505
506        let send_sm = GatewayClientStateMachinesV2::Send(SendStateMachine {
507            common: SendSMCommon {
508                operation_id,
509                outpoint: payload.outpoint,
510                contract: payload.contract.clone(),
511                max_delay: expiration.saturating_sub(EXPIRATION_DELTA_MINIMUM_V2),
512                min_contract_amount,
513                invoice: payload.invoice,
514                claim_keypair: self.keypair,
515            },
516            state: SendSMState::Sending,
517        });
518
519        let mut dbtx = self.client_ctx.module_db().begin_transaction().await;
520        self.client_ctx
521            .manual_operation_start_dbtx(
522                &mut dbtx.to_ref_nc(),
523                operation_id,
524                LightningCommonInit::KIND.as_str(),
525                GatewayOperationMetaV2::role(GatewayOperationRoleV2::Send),
526                vec![self.client_ctx.make_dyn_state(send_sm)],
527            )
528            .await
529            .ok();
530
531        self.client_ctx
532            .log_event(
533                &mut dbtx,
534                OutgoingPaymentStarted {
535                    operation_start,
536                    outgoing_contract: payload.contract.clone(),
537                    min_contract_amount,
538                    invoice_amount: Amount::from_msats(amount),
539                    max_delay: expiration.saturating_sub(EXPIRATION_DELTA_MINIMUM_V2),
540                },
541            )
542            .await;
543        dbtx.commit_tx().await;
544
545        Ok(self.subscribe_send(operation_id).await)
546    }
547
548    /// Returns the timelock budget, in blocks, that the outgoing contract at
549    /// `outpoint` leaves for a payment dispatched now, retrying until the
550    /// federation answers. A contract the federation no longer knows leaves
551    /// no budget.
552    async fn await_outgoing_contract_max_delay(&self, outpoint: OutPoint) -> u64 {
553        let expiration = retry(
554            "outgoing contract expiration",
555            backoff_util::background_backoff(),
556            || self.module_api.outgoing_contract_expiration(outpoint),
557        )
558        .await
559        .expect("Retries until the federation answers");
560
561        expiration.map_or(0, |(_, expiration)| {
562            expiration.saturating_sub(EXPIRATION_DELTA_MINIMUM_V2)
563        })
564    }
565
566    pub async fn subscribe_send(&self, operation_id: OperationId) -> Result<[u8; 32], Signature> {
567        let mut stream = self.notifier.subscribe(operation_id).await;
568
569        loop {
570            if let Some(GatewayClientStateMachinesV2::Send(state)) = stream.next().await {
571                match state.state {
572                    SendSMState::Sending => {}
573                    SendSMState::Claiming(claiming) => {
574                        // The preimage is proof the payment succeeded, so return it to
575                        // the sender as soon as it is available rather than waiting for
576                        // an additional ordering. The gateway's claim of the outgoing
577                        // contract has already been submitted by the send state machine
578                        // and finalizes in the background.
579                        return Ok(claiming.preimage);
580                    }
581                    SendSMState::Cancelled(cancelled) => {
582                        warn!("Outgoing lightning payment is cancelled {:?}", cancelled);
583
584                        let signature = self
585                            .keypair
586                            .sign_schnorr(state.common.contract.forfeit_message());
587
588                        assert!(state.common.contract.verify_forfeit_signature(&signature));
589
590                        return Err(signature);
591                    }
592                }
593            }
594        }
595    }
596
597    async fn legacy_completion_exists(
598        &self,
599        operation_id: OperationId,
600        circuit: IncomingCircuitKey,
601    ) -> bool {
602        let active = self
603            .client_ctx
604            .get_own_operation_active_states(operation_id)
605            .await
606            .into_iter()
607            .map(|(state, _)| state)
608            .collect::<Vec<_>>();
609        let inactive = self
610            .client_ctx
611            .get_own_operation_inactive_states(operation_id)
612            .await
613            .into_iter()
614            .map(|(state, _)| state)
615            .collect::<Vec<_>>();
616
617        legacy_completion_in_states(&active, &inactive, circuit)
618    }
619
620    /// Refuses to fund a fresh incoming contract unless a threshold of
621    /// guardians is answering. LNv1 gets this check implicitly by fetching the
622    /// offer from the federation first; LNv2 reads the contract from the
623    /// gateway's own database, so it has to probe explicitly.
624    async fn ensure_federation_responsive(
625        &self,
626        payment_hash: sha256::Hash,
627    ) -> Result<(), RelayIncomingHtlcError> {
628        match timeout(
629            FEDERATION_LIVENESS_TIMEOUT,
630            self.module_api.consensus_block_count(),
631        )
632        .await
633        {
634            Ok(Ok(_consensus_block_count)) => Ok(()),
635            Ok(Err(err)) => {
636                warn!(
637                    %payment_hash,
638                    err = %err.fmt_compact(),
639                    "Federation liveness probe failed, refusing to fund incoming contract"
640                );
641                Err(RelayIncomingHtlcError::FederationUnreachable(Box::new(err)))
642            }
643            Err(_elapsed) => {
644                warn!(
645                    %payment_hash,
646                    timeout_secs = FEDERATION_LIVENESS_TIMEOUT.as_secs(),
647                    "Federation liveness probe timed out, refusing to fund incoming contract"
648                );
649                Err(RelayIncomingHtlcError::FederationTimeout {
650                    timeout_secs: FEDERATION_LIVENESS_TIMEOUT.as_secs(),
651                })
652            }
653        }
654    }
655
656    /// Funds the incoming contract of an intercepted LNv2 HTLC and starts the
657    /// operation that completes the HTLC's circuit.
658    ///
659    /// A contract or circuit that is already being handled is joined rather
660    /// than started twice.
661    ///
662    /// # Errors
663    ///
664    /// Fails with a [`RelayIncomingHtlcError`] if a fresh contract is not
665    /// funded because the HTLC's claim deadline is too close or the federation
666    /// does not answer, or if the funding
667    /// transaction or the completion operation could not be started and no
668    /// earlier attempt had started it.
669    pub async fn relay_incoming_htlc(
670        &self,
671        payment_hash: sha256::Hash,
672        incoming_chan_id: u64,
673        htlc_id: u64,
674        contract: IncomingContract,
675        amount_msat: u64,
676        blocks_to_claim_deadline: u32,
677    ) -> Result<(), RelayIncomingHtlcError> {
678        let operation_start = now();
679        let receive_operation_id = OperationId::from_encodable(&contract);
680        let circuit = IncomingCircuitKey {
681            incoming_chan_id,
682            htlc_id,
683        };
684        let completion_operation_id = incoming_circuit_operation_id(receive_operation_id, circuit);
685        let receive_exists = self.client_ctx.operation_exists(receive_operation_id).await;
686        let completion_exists = self
687            .client_ctx
688            .operation_exists(completion_operation_id)
689            .await;
690        let legacy_completion_exists = self
691            .legacy_completion_exists(receive_operation_id, circuit)
692            .await;
693        let plan = incoming_relay_plan(receive_exists, completion_exists, legacy_completion_exists);
694        if plan == IncomingRelayPlan::Replay {
695            return Ok(());
696        }
697
698        let commitment = contract.commitment.clone();
699        if plan == IncomingRelayPlan::CreateReceiveAndCompletion {
700            // Only gate fresh funding: the other plans resume an already
701            // funded contract and must not be cancelled.
702            //
703            // Funding is irreversible, and the gateway is only reimbursed by
704            // settling the HTLC before its claim deadline, so do not fund an
705            // HTLC that is about to expire. The margin is small because LDK's
706            // default invoices leave only 3 blocks before the claim deadline.
707            if blocks_to_claim_deadline < LNV2_CLAIM_DEADLINE_MARGIN {
708                return Err(RelayIncomingHtlcError::ClaimDeadlineTooClose {
709                    blocks_to_claim_deadline,
710                });
711            }
712            self.ensure_federation_responsive(payment_hash).await?;
713
714            let refund_keypair = self.keypair;
715            let client_output = ClientOutput::<LightningOutput> {
716                output: LightningOutput::V0(LightningOutputV0::Incoming(contract.clone())),
717                amounts: Amounts::new_bitcoin(contract.commitment.amount),
718            };
719            let client_output_sm = ClientOutputSM::<GatewayClientStateMachinesV2> {
720                state_machines: Arc::new(move |range: OutPointRange| {
721                    assert_eq!(range.count(), 1);
722
723                    vec![GatewayClientStateMachinesV2::Receive(ReceiveStateMachine {
724                        common: ReceiveSMCommon {
725                            operation_id: receive_operation_id,
726                            contract: contract.clone(),
727                            outpoint: range.into_iter().next().unwrap(),
728                            refund_keypair,
729                        },
730                        state: ReceiveSMState::Funding,
731                    })]
732                }),
733            };
734
735            let client_output = self.client_ctx.make_client_outputs(ClientOutputBundle::new(
736                vec![client_output],
737                vec![client_output_sm],
738            ));
739            let transaction = TransactionBuilder::new().with_outputs(client_output);
740
741            let creation_result = self
742                .client_ctx
743                .finalize_and_submit_transaction(
744                    receive_operation_id,
745                    LightningCommonInit::KIND.as_str(),
746                    |_| GatewayOperationMetaV2::role(GatewayOperationRoleV2::Receive),
747                    transaction,
748                )
749                .await;
750            if let Err(error) = creation_result {
751                let operation_exists = self.client_ctx.operation_exists(receive_operation_id).await;
752                if operation_creation_failed_permanently(true, operation_exists) {
753                    return Err(error.into());
754                }
755            } else {
756                let mut dbtx = self.client_ctx.module_db().begin_transaction().await;
757                self.client_ctx
758                    .log_event(
759                        &mut dbtx,
760                        IncomingPaymentStarted {
761                            operation_start,
762                            incoming_contract_commitment: commitment,
763                            invoice_amount: Amount::from_msats(amount_msat),
764                        },
765                    )
766                    .await;
767                dbtx.commit_tx().await;
768            }
769        }
770
771        let completion =
772            GatewayClientStateMachinesV2::CircuitComplete(CircuitCompleteStateMachine {
773                common: CircuitCompleteSMCommon {
774                    operation_id: completion_operation_id,
775                    receive_operation_id,
776                    payment_hash,
777                    circuit,
778                },
779                state: CompleteSMState::Pending,
780            });
781        let creation_result = self
782            .client_ctx
783            .manual_operation_start(
784                completion_operation_id,
785                LightningCommonInit::KIND.as_str(),
786                GatewayOperationMetaV2::role(GatewayOperationRoleV2::CircuitCompletion),
787                vec![self.client_ctx.make_dyn_state(completion)],
788            )
789            .await;
790        if let Err(error) = creation_result {
791            let operation_exists = self
792                .client_ctx
793                .operation_exists(completion_operation_id)
794                .await;
795            if operation_creation_failed_permanently(true, operation_exists) {
796                return Err(error.into());
797            }
798        }
799
800        Ok(())
801    }
802
803    /// Funds the incoming contract of a direct swap and waits for its final
804    /// state.
805    ///
806    /// A swap that was already started resumes unconditionally: its operation
807    /// is keyed on the contract, so a re-entrant call — a restarted state
808    /// machine, or a concurrent payer — joins the operation in progress
809    /// instead of funding twice. `allow_fresh_dispatch` is consulted only
810    /// when no such operation exists: callers pass `false` when a wall-clock
811    /// gate such as invoice expiry forbids starting a new swap, and receive
812    /// `Ok(None)` to signal that nothing was started.
813    ///
814    /// # Errors
815    ///
816    /// Fails with a [`TransactionSubmitError`] if the funding transaction of a
817    /// fresh swap could not be submitted.
818    pub async fn relay_direct_swap(
819        &self,
820        contract: IncomingContract,
821        amount_msat: u64,
822        allow_fresh_dispatch: bool,
823    ) -> Result<Option<FinalReceiveState>, TransactionSubmitError> {
824        let operation_start = now();
825
826        let operation_id = OperationId::from_encodable(&contract);
827
828        if self.client_ctx.operation_exists(operation_id).await {
829            return Ok(Some(self.await_receive(operation_id).await));
830        }
831
832        if !allow_fresh_dispatch {
833            return Ok(None);
834        }
835
836        let refund_keypair = self.keypair;
837
838        let client_output = ClientOutput::<LightningOutput> {
839            output: LightningOutput::V0(LightningOutputV0::Incoming(contract.clone())),
840            amounts: Amounts::new_bitcoin(contract.commitment.amount),
841        };
842        let commitment = contract.commitment.clone();
843        let client_output_sm = ClientOutputSM::<GatewayClientStateMachinesV2> {
844            state_machines: Arc::new(move |range| {
845                assert_eq!(range.count(), 1);
846
847                vec![GatewayClientStateMachinesV2::Receive(ReceiveStateMachine {
848                    common: ReceiveSMCommon {
849                        operation_id,
850                        contract: contract.clone(),
851                        outpoint: range.into_iter().next().unwrap(),
852                        refund_keypair,
853                    },
854                    state: ReceiveSMState::Funding,
855                })]
856            }),
857        };
858
859        let client_output = self.client_ctx.make_client_outputs(ClientOutputBundle::new(
860            vec![client_output],
861            vec![client_output_sm],
862        ));
863
864        let transaction = TransactionBuilder::new().with_outputs(client_output);
865
866        self.client_ctx
867            .finalize_and_submit_transaction(
868                operation_id,
869                LightningCommonInit::KIND.as_str(),
870                |_| GatewayOperationMetaV2::role(GatewayOperationRoleV2::Receive),
871                transaction,
872            )
873            .await?;
874
875        let mut dbtx = self.client_ctx.module_db().begin_transaction().await;
876        self.client_ctx
877            .log_event(
878                &mut dbtx,
879                IncomingPaymentStarted {
880                    operation_start,
881                    incoming_contract_commitment: commitment,
882                    invoice_amount: Amount::from_msats(amount_msat),
883                },
884            )
885            .await;
886        dbtx.commit_tx().await;
887
888        Ok(Some(self.await_receive(operation_id).await))
889    }
890
891    pub async fn await_receive(&self, operation_id: OperationId) -> FinalReceiveState {
892        let mut stream = self.notifier.subscribe(operation_id).await;
893
894        loop {
895            if let Some(GatewayClientStateMachinesV2::Receive(state)) = stream.next().await {
896                match state.state {
897                    ReceiveSMState::Funding => {}
898                    ReceiveSMState::Rejected(..) => return FinalReceiveState::Rejected,
899                    ReceiveSMState::Success(preimage) => {
900                        return FinalReceiveState::Success(preimage);
901                    }
902                    ReceiveSMState::Refunding(out_points) => {
903                        if self
904                            .client_ctx
905                            .await_primary_module_outputs(operation_id, out_points)
906                            .await
907                            .is_err()
908                        {
909                            return FinalReceiveState::Failure;
910                        }
911
912                        return FinalReceiveState::Refunded;
913                    }
914                    ReceiveSMState::Failure => return FinalReceiveState::Failure,
915                }
916            }
917        }
918    }
919
920    /// Waits for a legacy combined receive operation or a circuit-completion
921    /// operation to finish.
922    ///
923    /// `operation_id` must identify either a legacy operation containing
924    /// [`GatewayClientStateMachinesV2::Complete`] or a role-specific operation
925    /// containing [`GatewayClientStateMachinesV2::CircuitComplete`]. Both
926    /// successful completion and a durable permanent outcome conflict terminate
927    /// the wait.
928    pub async fn await_completion(&self, operation_id: OperationId) {
929        let mut stream = self.notifier.subscribe(operation_id).await;
930
931        loop {
932            match stream.next().await {
933                Some(GatewayClientStateMachinesV2::Complete(state)) => {
934                    if matches!(
935                        state.state,
936                        CompleteSMState::Completed | CompleteSMState::CompletionFailed(_)
937                    ) {
938                        info!(%state, "LNv2 completion state machine finished");
939                        return;
940                    }
941
942                    info!(%state, "Waiting for LNv2 completion state machine");
943                }
944                Some(GatewayClientStateMachinesV2::Receive(state)) => {
945                    info!(%state, "Waiting for LNv2 completion state machine");
946                    continue;
947                }
948                Some(GatewayClientStateMachinesV2::CircuitComplete(state)) => {
949                    if matches!(
950                        state.state,
951                        CompleteSMState::Completed | CompleteSMState::CompletionFailed(_)
952                    ) {
953                        info!(%state, "LNv2 circuit completion state machine finished");
954                        return;
955                    }
956
957                    info!(%state, "Waiting for LNv2 circuit completion state machine");
958                }
959                Some(state) => {
960                    warn!(%state, "Operation is not an LNv2 completion state machine");
961                    return;
962                }
963                None => return,
964            }
965        }
966    }
967}
968
969/// A refusal to pay an invoice for an LNv2 client.
970///
971/// These are the checks the gateway makes before it starts paying. A payment
972/// that starts and is later cancelled is not a failure: it is reported as the
973/// forfeit signature in the `Ok` value.
974#[derive(Debug, Error)]
975#[non_exhaustive]
976pub enum GatewaySendPaymentError {
977    /// The gateway is not connected to its lightning node, so a new send
978    /// would sit waiting for it.
979    #[error("The gateway is not connected to its lightning node")]
980    LightningNotConnected,
981
982    /// The outgoing contract names another gateway's key, so this gateway
983    /// could never claim it.
984    #[error("The outgoing contract is keyed to another gateway")]
985    NotOurContract,
986
987    /// The request's signature over the invoice does not verify against the
988    /// contract's refund key.
989    #[error("Invalid auth signature for the invoice data")]
990    InvalidAuthSignature,
991
992    /// The federation could not be asked about the outgoing contract.
993    #[error("The gateway can not reach the federation")]
994    Federation(#[source] Box<FederationError>),
995
996    /// The federation has not confirmed the outgoing contract yet.
997    #[error("The outgoing contract has not yet been confirmed")]
998    ContractNotConfirmed,
999
1000    /// The contract the federation confirmed at the request's outpoint is not
1001    /// the one in the request.
1002    #[error("Contract Id returned by the federation does not match contract in request")]
1003    ContractIdMismatch,
1004
1005    /// The invoice carries no amount.
1006    #[error("Invoice is missing amount")]
1007    MissingInvoiceAmount,
1008
1009    /// The invoice's payment hash is not the one the contract is locked to.
1010    #[error("The invoices payment hash does not match the contracts payment hash")]
1011    PaymentHashMismatch,
1012
1013    /// The gateway could not work out the smallest contract amount it accepts
1014    /// for this payment.
1015    #[error("The minimum contract amount could not be computed")]
1016    MinContractAmount(#[source] GatewayClientV2Error),
1017
1018    /// Another operation already claimed the contract's payment image, so this
1019    /// gateway will never pay it out for this contract.
1020    #[error("Another contract for this payment image was already accepted")]
1021    PaymentImageAlreadyClaimed,
1022}
1023
1024/// Why the gateway did not relay an intercepted LNv2 HTLC.
1025#[derive(Debug, Error)]
1026pub enum RelayIncomingHtlcError {
1027    /// The federation liveness probe failed, so a fresh incoming contract is
1028    /// not funded.
1029    #[error("The federation did not answer the liveness probe")]
1030    FederationUnreachable(#[source] Box<FederationError>),
1031
1032    /// The federation liveness probe did not complete in time, so a fresh
1033    /// incoming contract is not funded.
1034    #[error("The federation did not answer the liveness probe within {timeout_secs}s")]
1035    FederationTimeout {
1036        /// How long the probe waited.
1037        timeout_secs: u64,
1038    },
1039
1040    /// The HTLC's claim deadline is too close to fund a fresh incoming
1041    /// contract and still settle the HTLC in time.
1042    #[error("HTLC claim deadline is only {blocks_to_claim_deadline} blocks away")]
1043    ClaimDeadlineTooClose {
1044        /// Blocks left until the HTLC's claim deadline.
1045        blocks_to_claim_deadline: u32,
1046    },
1047
1048    /// The funding transaction or the completion operation could not be
1049    /// started.
1050    #[error("The incoming contract could not be funded")]
1051    Transaction(#[from] TransactionSubmitError),
1052}
1053
1054impl From<FederationError> for GatewaySendPaymentError {
1055    fn from(source: FederationError) -> Self {
1056        Self::Federation(Box::new(source))
1057    }
1058}
1059
1060/// An interface between module implementation and the general `Gateway`
1061///
1062/// To abstract away and decouple the core gateway from the modules, the
1063/// interface between the is expressed as a trait. The core gateway handles
1064/// LNv2 operations that require access to the database or lightning node.
1065#[async_trait]
1066pub trait IGatewayClientV2: Debug + Send + Sync {
1067    /// Uses the gateway's Lightning node to complete a payment.
1068    ///
1069    /// Implementations must absorb and retry every transient node or
1070    /// connectivity failure. They return `Err` only when Lightning has reached
1071    /// a permanent state that makes the requested outcome impossible. The
1072    /// future may block while retrying and must remain cancellation-safe.
1073    /// Completion state machines persist any returned error as terminal
1074    /// `CompletionFailed`.
1075    async fn complete_htlc(
1076        &self,
1077        htlc_response: InterceptPaymentResponse,
1078    ) -> Result<(), LightningRpcError>;
1079
1080    /// Determines if the payment can be completed using a direct swap to
1081    /// another federation.
1082    ///
1083    /// A direct swap is determined by checking the gateway's connected
1084    /// lightning node against the invoice's payee lightning node. If they
1085    /// are the same, then the gateway can use another client to complete
1086    /// the payment be swapping ecash instead of a payment over the
1087    /// Lightning network.
1088    ///
1089    /// # Errors
1090    ///
1091    /// Fails with a [`GatewayClientV2Error`] if the invoice is payable by a
1092    /// direct swap but the gateway holds no incoming contract it can use for
1093    /// it. The send state machine cancels the payment on a failure.
1094    async fn is_direct_swap(
1095        &self,
1096        invoice: &Bolt11Invoice,
1097    ) -> Result<Option<(IncomingContract, ClientHandleArc)>, GatewayClientV2Error>;
1098
1099    /// Initiates a payment over the Lightning network.
1100    async fn pay(
1101        &self,
1102        invoice: Bolt11Invoice,
1103        max_delay: u64,
1104        max_fee: Amount,
1105    ) -> Result<[u8; 32], LightningRpcError>;
1106
1107    /// Returns whether the gateway's Lightning node has any record of an
1108    /// outbound payment for `payment_hash`, whatever its state.
1109    ///
1110    /// The send state machine consults this when it resumes after a restart:
1111    /// a payment the node already knows was dispatched before the crash and
1112    /// must be resolved through [`IGatewayClientV2::pay`]'s idempotent resume
1113    /// path rather than cancelled by pre-dispatch checks. A wrong `false`
1114    /// forfeits a contract whose payment may still settle, so implementations
1115    /// must absorb transient node failures and only answer once the node's
1116    /// payment store could actually be consulted.
1117    async fn outbound_payment_exists(&self, payment_hash: sha256::Hash) -> bool;
1118
1119    /// Computes the minimum contract amount necessary for making an outgoing
1120    /// payment.
1121    ///
1122    /// The minimum contract amount must contain transaction fees to cover the
1123    /// gateway's transaction fee and optionally additional fee to cover the
1124    /// gateway's Lightning fee if the payment goes over the Lightning
1125    /// network.
1126    ///
1127    /// # Errors
1128    ///
1129    /// Fails with a [`GatewayClientV2Error`] if the gateway cannot price a
1130    /// payment for this federation.
1131    async fn min_contract_amount(
1132        &self,
1133        federation_id: &FederationId,
1134        amount: u64,
1135    ) -> Result<Amount, GatewayClientV2Error>;
1136
1137    /// Check if this invoice was created using LNv1 and if the gateway is
1138    /// connected to the target federation.
1139    async fn is_lnv1_invoice(&self, invoice: &Bolt11Invoice) -> Option<Spanned<ClientHandleArc>>;
1140
1141    /// Perform a swap from an LNv2 `OutgoingContract` to an LNv1
1142    /// `IncomingContract`.
1143    ///
1144    /// A swap that was already started resumes unconditionally;
1145    /// `allow_fresh_dispatch` is consulted only when no operation for this
1146    /// swap exists yet. Callers pass `false` when a wall-clock gate such as
1147    /// invoice expiry forbids starting a new swap, and receive `Ok(None)` to
1148    /// signal that nothing was started.
1149    ///
1150    /// # Errors
1151    ///
1152    /// Fails with a [`GatewayClientV2Error`] if the swap could not be started
1153    /// or followed. The send state machine cancels the payment on a failure.
1154    async fn relay_lnv1_swap(
1155        &self,
1156        client: &ClientHandleArc,
1157        invoice: &Bolt11Invoice,
1158        allow_fresh_dispatch: bool,
1159    ) -> Result<Option<FinalReceiveState>, GatewayClientV2Error>;
1160
1161    /// Claims the given payment image for `operation_id` in the gateway's
1162    /// global database, returning `true` if this operation may claim the
1163    /// outgoing contract (the image was unclaimed, or already claimed by
1164    /// this same operation) and `false` if another operation already
1165    /// claimed it.
1166    ///
1167    /// A single Lightning payment yields a single preimage, so the gateway may
1168    /// claim at most one outgoing contract per payment image. Unlike the
1169    /// per-federation module database, this spans all of the gateway's
1170    /// federations, so two clients in different federations paying the same
1171    /// invoice cannot both be claimed.
1172    async fn claim_payment_image(
1173        &self,
1174        payment_image: &PaymentImage,
1175        operation_id: OperationId,
1176    ) -> bool;
1177
1178    /// Returns whether the gateway currently holds a connection to its
1179    /// lightning node. Only suitable for refusing new work: the answer is
1180    /// local, so it must not decide the fate of a payment already started.
1181    async fn is_lightning_connected(&self) -> bool;
1182
1183    /// Waits until the gateway holds a connection to its lightning node.
1184    async fn await_lightning_connected(&self);
1185}
1186
1187/// A failure reported by the gateway behind [`IGatewayClientV2`].
1188///
1189/// The trait is implemented by the gateway, not by this module, so the causes
1190/// are the gateway's own. This type carries them unchanged: its `Display` and
1191/// its `source()` are the cause's.
1192#[derive(Debug, Error)]
1193#[error(transparent)]
1194pub struct GatewayClientV2Error(Box<dyn std::error::Error + Send + Sync>);
1195
1196impl GatewayClientV2Error {
1197    /// Wraps a failure of the gateway's [`IGatewayClientV2`] implementation,
1198    /// which may be any error value or a plain message.
1199    pub fn new<E>(source: E) -> Self
1200    where
1201        E: Into<Box<dyn std::error::Error + Send + Sync>>,
1202    {
1203        Self(source.into())
1204    }
1205}
1206
1207#[cfg(test)]
1208mod tests;