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
67const FEDERATION_LIVENESS_TIMEOUT: Duration = Duration::from_secs(10);
71
72const LNV2_CLAIM_DEADLINE_MARGIN: u32 = 2;
75
76pub 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#[derive(Debug, Clone, Serialize, Deserialize)]
148#[serde(untagged)]
149pub enum GatewayOperationMetaV2 {
150 Legacy(()),
152 Role {
154 role: GatewayOperationRoleV2,
156 },
157}
158
159#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
161#[serde(rename_all = "snake_case")]
162pub enum GatewayOperationRoleV2 {
163 Send,
165 Receive,
167 CircuitCompletion,
169}
170
171impl GatewayOperationMetaV2 {
172 pub fn role(role: GatewayOperationRoleV2) -> Self {
174 Self::Role { role }
175 }
176
177 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 Complete(CompleteStateMachine),
305 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 pub async fn send_payment(
409 &self,
410 payload: SendPaymentPayload,
411 ) -> Result<Result<[u8; 32], Signature>, GatewaySendPaymentError> {
412 let operation_start = now();
413
414 let operation_id = OperationId::from_encodable(&payload.contract.clone());
420
421 if payload.contract.claim_pk != self.keypair.public_key() {
425 return Err(GatewaySendPaymentError::NotOurContract);
426 }
427
428 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 if self.client_ctx.operation_exists(operation_id).await {
444 return Ok(self.subscribe_send(operation_id).await);
445 }
446
447 if !self.gateway.is_lightning_connected().await {
453 return Err(GatewaySendPaymentError::LightningNotConnected);
454 }
455
456 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 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 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 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 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 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 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 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 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#[derive(Debug, Error)]
975#[non_exhaustive]
976pub enum GatewaySendPaymentError {
977 #[error("The gateway is not connected to its lightning node")]
980 LightningNotConnected,
981
982 #[error("The outgoing contract is keyed to another gateway")]
985 NotOurContract,
986
987 #[error("Invalid auth signature for the invoice data")]
990 InvalidAuthSignature,
991
992 #[error("The gateway can not reach the federation")]
994 Federation(#[source] Box<FederationError>),
995
996 #[error("The outgoing contract has not yet been confirmed")]
998 ContractNotConfirmed,
999
1000 #[error("Contract Id returned by the federation does not match contract in request")]
1003 ContractIdMismatch,
1004
1005 #[error("Invoice is missing amount")]
1007 MissingInvoiceAmount,
1008
1009 #[error("The invoices payment hash does not match the contracts payment hash")]
1011 PaymentHashMismatch,
1012
1013 #[error("The minimum contract amount could not be computed")]
1016 MinContractAmount(#[source] GatewayClientV2Error),
1017
1018 #[error("Another contract for this payment image was already accepted")]
1021 PaymentImageAlreadyClaimed,
1022}
1023
1024#[derive(Debug, Error)]
1026pub enum RelayIncomingHtlcError {
1027 #[error("The federation did not answer the liveness probe")]
1030 FederationUnreachable(#[source] Box<FederationError>),
1031
1032 #[error("The federation did not answer the liveness probe within {timeout_secs}s")]
1035 FederationTimeout {
1036 timeout_secs: u64,
1038 },
1039
1040 #[error("HTLC claim deadline is only {blocks_to_claim_deadline} blocks away")]
1043 ClaimDeadlineTooClose {
1044 blocks_to_claim_deadline: u32,
1046 },
1047
1048 #[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#[async_trait]
1066pub trait IGatewayClientV2: Debug + Send + Sync {
1067 async fn complete_htlc(
1076 &self,
1077 htlc_response: InterceptPaymentResponse,
1078 ) -> Result<(), LightningRpcError>;
1079
1080 async fn is_direct_swap(
1095 &self,
1096 invoice: &Bolt11Invoice,
1097 ) -> Result<Option<(IncomingContract, ClientHandleArc)>, GatewayClientV2Error>;
1098
1099 async fn pay(
1101 &self,
1102 invoice: Bolt11Invoice,
1103 max_delay: u64,
1104 max_fee: Amount,
1105 ) -> Result<[u8; 32], LightningRpcError>;
1106
1107 async fn outbound_payment_exists(&self, payment_hash: sha256::Hash) -> bool;
1118
1119 async fn min_contract_amount(
1132 &self,
1133 federation_id: &FederationId,
1134 amount: u64,
1135 ) -> Result<Amount, GatewayClientV2Error>;
1136
1137 async fn is_lnv1_invoice(&self, invoice: &Bolt11Invoice) -> Option<Spanned<ClientHandleArc>>;
1140
1141 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 async fn claim_payment_image(
1173 &self,
1174 payment_image: &PaymentImage,
1175 operation_id: OperationId,
1176 ) -> bool;
1177
1178 async fn is_lightning_connected(&self) -> bool;
1182
1183 async fn await_lightning_connected(&self);
1185}
1186
1187#[derive(Debug, Error)]
1193#[error(transparent)]
1194pub struct GatewayClientV2Error(Box<dyn std::error::Error + Send + Sync>);
1195
1196impl GatewayClientV2Error {
1197 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;