1#![deny(clippy::pedantic)]
2#![allow(clippy::cast_possible_truncation)]
3#![allow(clippy::cast_possible_wrap)]
4#![allow(clippy::cast_sign_loss)]
5#![allow(clippy::default_trait_access)]
6#![allow(clippy::doc_markdown)]
7#![allow(clippy::missing_errors_doc)]
8#![allow(clippy::missing_panics_doc)]
9#![allow(clippy::module_name_repetitions)]
10#![allow(clippy::must_use_candidate)]
11#![allow(clippy::return_self_not_must_use)]
12#![allow(clippy::similar_names)]
13#![allow(clippy::too_many_lines)]
14#![allow(clippy::large_futures)]
15#![allow(clippy::struct_field_names)]
16
17pub mod client;
18pub mod config;
19pub mod envs;
20mod error;
21mod events;
22mod federation_manager;
23mod federation_status;
24mod iroh_server;
25mod metrics;
26mod rate_limit;
27mod registration_health;
28pub mod rpc_server;
29mod types;
30
31use std::collections::{BTreeMap, BTreeSet};
32use std::env;
33use std::fmt::Display;
34use std::net::SocketAddr;
35use std::str::FromStr;
36use std::sync::Arc;
37use std::time::{Duration, UNIX_EPOCH};
38
39use anyhow::{Context, anyhow, ensure};
40use async_trait::async_trait;
41use bitcoin::hashes::sha256;
42use bitcoin::{Address, Network, Txid, secp256k1};
43use clap::Parser;
44use client::GatewayClientBuilder;
45pub use config::GatewayParameters;
46use config::{DatabaseBackend, GatewayOpts};
47use envs::FM_GATEWAY_SKIP_WAIT_FOR_SYNC_ENV;
48use error::FederationNotConnected;
49use events::ALL_GATEWAY_EVENTS;
50use federation_manager::FederationManager;
51use fedimint_bip39::{Bip39RootSecretStrategy, Language, Mnemonic};
52use fedimint_bitcoind::bitcoincore::BitcoindClient;
53use fedimint_bitcoind::{EsploraClient, IBitcoindRpc};
54use fedimint_client::module_init::ClientModuleInitRegistry;
55use fedimint_client::secret::RootSecretStrategy;
56use fedimint_client::{Client, ClientHandleArc};
57use fedimint_core::base32::{self, FEDIMINT_PREFIX};
58use fedimint_core::config::FederationId;
59use fedimint_core::core::OperationId;
60use fedimint_core::db::{Committable, Database, DatabaseTransaction, apply_migrations};
61use fedimint_core::envs::is_env_var_set;
62use fedimint_core::invite_code::InviteCode;
63use fedimint_core::module::CommonModuleInit;
64use fedimint_core::module::registry::ModuleDecoderRegistry;
65use fedimint_core::rustls::install_crypto_provider;
66use fedimint_core::secp256k1::PublicKey;
67use fedimint_core::secp256k1::schnorr::Signature;
68use fedimint_core::task::{TaskGroup, TaskHandle, TaskShutdownToken, sleep, timeout};
69use fedimint_core::time::duration_since_epoch;
70use fedimint_core::util::backoff_util::fibonacci_max_one_hour;
71use fedimint_core::util::{FmtCompact, FmtCompactAnyhow, SafeUrl, Spanned, retry};
72use fedimint_core::{
73 Amount, BitcoinAmountOrAll, PeerId, TieredCounts, crit, fedimint_build_code_version_env,
74 get_network_for_address,
75};
76use fedimint_eventlog::{DBTransactionEventLogExt, EventLogId, StructuredPaymentEvents};
77use fedimint_gateway_common::{
78 BackupPayload, ChainSource, CloseChannelsWithPeerRequest, CloseChannelsWithPeerResponse,
79 ConnectFedPayload, ConnectPeerRequest, ConnectorType, CreateInvoiceForOperatorPayload,
80 CreateOfferPayload, CreateOfferResponse, DepositAddressPayload, DepositAddressRecheckPayload,
81 FederationBalanceInfo, FederationConfig, FederationInfo, GatewayBalances, GatewayFedConfig,
82 GatewayInfo, GetInvoiceRequest, GetInvoiceResponse, LeaveFedPayload, LightningInfo,
83 LightningMode, ListTransactionsPayload, ListTransactionsResponse, MnemonicResponse,
84 OpenChannelRequest, PayInvoiceForOperatorPayload, PayOfferPayload, PayOfferResponse,
85 PaymentLogPayload, PaymentLogResponse, PaymentStats, PaymentSummaryPayload,
86 PaymentSummaryResponse, PeginFromOnchainPayload, ReceiveEcashPayload, ReceiveEcashResponse,
87 RegisteredProtocol, SendOnchainRequest, SetChannelFeesRequest, SetFeesPayload,
88 SetMnemonicPayload, SpendEcashPayload, SpendEcashResponse, V1_API_ENDPOINT, WithdrawPayload,
89 WithdrawPreviewPayload, WithdrawPreviewResponse, WithdrawResponse, WithdrawToOnchainPayload,
90};
91use fedimint_gateway_server_db::{GatewayDbtxNcExt as _, get_gatewayd_database_migrations};
92pub use fedimint_gateway_ui::IAdminGateway;
93use fedimint_gw_client::events::compute_lnv1_stats;
94use fedimint_gw_client::pay::{OutgoingPaymentError, OutgoingPaymentErrorType};
95use fedimint_gw_client::{
96 GatewayClientModule, GatewayExtPayStates, GatewayExtReceiveStates, IGatewayClientV1,
97 SwapParameters,
98};
99use fedimint_gwv2_client::events::compute_lnv2_stats;
100use fedimint_gwv2_client::{
101 EXPIRATION_DELTA_MINIMUM_V2, FinalReceiveState, GatewayClientModuleV2, IGatewayClientV2,
102};
103use fedimint_lightning::lnd::GatewayLndClient;
104use fedimint_lightning::{
105 CreateInvoiceRequest, ILnRpcClient, InterceptPaymentRequest, InterceptPaymentResponse,
106 InvoiceDescription, LightningContext, LightningRpcError, LnRpcTracked, Lnv2HoldInvoiceFilter,
107 PayInvoiceResponse, PaymentAction, RouteHtlcStream, ldk,
108};
109use fedimint_ln_client::pay::PaymentData;
110use fedimint_ln_common::LightningCommonInit;
111use fedimint_ln_common::config::LightningClientConfig;
112use fedimint_ln_common::contracts::outgoing::OutgoingContractAccount;
113use fedimint_ln_common::contracts::{IdentifiableContract, Preimage};
114use fedimint_lnurl::VerifyResponse;
115use fedimint_lnv2_common::Bolt11InvoiceDescription;
116use fedimint_lnv2_common::contracts::{IncomingContract, PaymentImage};
117use fedimint_lnv2_common::gateway_api::{
118 CreateBolt11InvoicePayload, MAX_INVOICE_EXPIRY_SECS, PaymentFee, RoutingInfo,
119 SendPaymentPayload,
120};
121use fedimint_logging::LOG_GATEWAY;
122use fedimint_mint_client::{MintClientInit, MintClientModule, OOBNotes, ReissueExternalNotesState};
123use fedimint_mintv2_client::{
124 MintClientInit as MintV2ClientInit, MintClientModule as MintV2ClientModule,
125};
126use fedimint_wallet_client::{PegOutFees, WalletClientInit, WalletClientModule, WithdrawState};
127use futures::stream::StreamExt;
128use lightning_invoice::{Bolt11Invoice, RoutingFees};
129use rand::rngs::OsRng;
130use tokio::sync::RwLock;
131use tracing::{debug, info, info_span, warn};
132
133use crate::envs::FM_GATEWAY_MNEMONIC_ENV;
134use crate::error::{AdminGatewayError, LNv1Error, LNv2Error, PublicGatewayError};
135use crate::events::get_events_for_duration;
136use crate::rate_limit::TokenBucketRateLimiter;
137use crate::registration_health::RegistrationHealthTracker;
138use crate::rpc_server::run_webserver;
139use crate::types::PrettyInterceptPaymentRequest;
140
141const GW_ANNOUNCEMENT_TTL: Duration = Duration::from_mins(10);
143
144const DEFAULT_NUM_ROUTE_HINTS: u32 = 1;
147
148const DEFAULT_INVOICE_RATE_LIMIT_BURST: u32 = 50;
150
151const DEFAULT_INVOICE_RATE_LIMIT_PER_SECOND: u32 = 5;
154
155pub const DEFAULT_NETWORK: Network = Network::Regtest;
157
158const LIGHTNING_CONTEXT_RETRY_INTERVAL: Duration = Duration::from_secs(5);
161
162const VERIFY_WAIT_TIMEOUT: Duration = Duration::from_secs(30);
170
171pub type Result<T> = std::result::Result<T, PublicGatewayError>;
172pub type AdminResult<T> = std::result::Result<T, AdminGatewayError>;
173
174const DB_FILE: &str = "gatewayd.db";
177
178const LDK_NODE_DB_FOLDER: &str = "ldk_node";
181
182#[cfg_attr(doc, aquamarine::aquamarine)]
183#[derive(Clone, Debug)]
196pub enum GatewayState {
197 NotConfigured {
198 mnemonic_sender: tokio::sync::broadcast::Sender<()>,
201 },
202 Disconnected,
203 Syncing,
204 Connected,
205 Running {
206 lightning_context: LightningContext,
207 },
208 ShuttingDown {
209 lightning_context: LightningContext,
210 },
211}
212
213impl Display for GatewayState {
214 fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
215 match self {
216 GatewayState::NotConfigured { .. } => write!(f, "NotConfigured"),
217 GatewayState::Disconnected => write!(f, "Disconnected"),
218 GatewayState::Syncing => write!(f, "Syncing"),
219 GatewayState::Connected => write!(f, "Connected"),
220 GatewayState::Running { .. } => write!(f, "Running"),
221 GatewayState::ShuttingDown { .. } => write!(f, "ShuttingDown"),
222 }
223 }
224}
225
226#[derive(Debug, Clone)]
229struct Registration {
230 endpoint_url: SafeUrl,
232
233 keypair: secp256k1::Keypair,
235}
236
237impl Registration {
238 pub async fn new(db: &Database, endpoint_url: SafeUrl, protocol: RegisteredProtocol) -> Self {
239 let keypair = Gateway::load_or_create_gateway_keypair(db, protocol).await;
240 Self {
241 endpoint_url,
242 keypair,
243 }
244 }
245}
246
247#[bon::bon]
248impl Gateway {
249 #[builder(start_fn = builder, finish_fn = build)]
264 pub async fn new_with_builder(
265 #[builder(start_fn)] lightning_mode: LightningMode,
266 #[builder(start_fn)] client_builder: GatewayClientBuilder,
267 #[builder(start_fn)] gateway_db: Database,
268 bcrypt_password_hash: bcrypt::HashParts,
269 bcrypt_liquidity_manager_password_hash: Option<bcrypt::HashParts>,
270 gateway_state: GatewayState,
271 chain_source: ChainSource,
272 #[builder(default = ([127, 0, 0, 1], 80).into())] listen: SocketAddr,
273 api_addr: Option<SafeUrl>,
274 #[builder(default = DEFAULT_NETWORK)] network: Network,
275 #[builder(default = DEFAULT_NUM_ROUTE_HINTS)] num_route_hints: u32,
276 #[builder(default = PaymentFee::TRANSACTION_FEE_DEFAULT)] default_routing_fees: PaymentFee,
277 #[builder(default = PaymentFee::TRANSACTION_FEE_DEFAULT)]
278 default_transaction_fees: PaymentFee,
279 iroh_listen: Option<SocketAddr>,
280 iroh_dns: Option<SafeUrl>,
281 #[builder(default)] iroh_relays: Vec<SafeUrl>,
282 metrics_listen: Option<SocketAddr>,
283 ) -> anyhow::Result<Gateway> {
284 let versioned_api = api_addr.map(|addr| {
285 addr.join(V1_API_ENDPOINT)
286 .expect("Failed to version gateway API address")
287 });
288
289 let metrics_listen = metrics_listen.unwrap_or_else(|| {
290 SocketAddr::new(
291 std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
292 listen.port() + 1,
293 )
294 });
295
296 Gateway::new(
297 lightning_mode,
298 GatewayParameters {
299 listen,
300 versioned_api,
301 bcrypt_password_hash,
302 bcrypt_liquidity_manager_password_hash,
303 network,
304 num_route_hints,
305 default_routing_fees,
306 default_transaction_fees,
307 iroh_listen,
308 iroh_dns,
309 iroh_relays,
310 skip_setup: true,
311 metrics_listen,
312 invoice_rate_limit_burst: DEFAULT_INVOICE_RATE_LIMIT_BURST,
313 invoice_rate_limit_per_second: DEFAULT_INVOICE_RATE_LIMIT_PER_SECOND,
314 },
315 gateway_db,
316 client_builder,
317 gateway_state,
318 chain_source,
319 )
320 .await
321 }
322}
323
324enum ReceivePaymentStreamAction {
326 RetryAfterDelay,
327 NoRetry,
328}
329
330#[derive(Clone)]
331pub struct Gateway {
332 federation_manager: Arc<RwLock<FederationManager>>,
334
335 lightning_mode: LightningMode,
337
338 state: Arc<RwLock<GatewayState>>,
340
341 client_builder: GatewayClientBuilder,
344
345 gateway_db: Database,
347
348 listen: SocketAddr,
350
351 metrics_listen: SocketAddr,
353
354 task_group: TaskGroup,
356
357 bcrypt_password_hash: String,
359
360 bcrypt_liquidity_manager_password_hash: Option<String>,
363
364 num_route_hints: u32,
366
367 network: Network,
369
370 chain_source: ChainSource,
372
373 default_routing_fees: PaymentFee,
375
376 default_transaction_fees: PaymentFee,
378
379 iroh_sk: iroh::SecretKey,
381
382 iroh_listen: Option<SocketAddr>,
384
385 iroh_dns: Option<SafeUrl>,
387
388 iroh_relays: Vec<SafeUrl>,
391
392 registrations: BTreeMap<RegisteredProtocol, Registration>,
395
396 registration_health: RegistrationHealthTracker,
398
399 invoice_rate_limiter: Arc<TokenBucketRateLimiter>,
401}
402
403impl std::fmt::Debug for Gateway {
404 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
405 f.debug_struct("Gateway")
406 .field("federation_manager", &self.federation_manager)
407 .field("state", &self.state)
408 .field("client_builder", &self.client_builder)
409 .field("gateway_db", &self.gateway_db)
410 .field("listen", &self.listen)
411 .field("registrations", &self.registrations)
412 .finish_non_exhaustive()
413 }
414}
415
416struct WithdrawDetails {
418 amount: Amount,
419 mint_fees: Option<Amount>,
420 peg_out_fees: PegOutFees,
421}
422
423async fn withdraw_v2(
425 client: &ClientHandleArc,
426 wallet_module: &fedimint_walletv2_client::WalletClientModule,
427 address: &Address,
428 amount: BitcoinAmountOrAll,
429) -> AdminResult<WithdrawResponse> {
430 let fee = wallet_module
431 .send_fee()
432 .await
433 .map_err(|e| AdminGatewayError::WithdrawError {
434 failure_reason: e.to_string(),
435 })?;
436
437 let withdraw_amount = match amount {
438 BitcoinAmountOrAll::All => {
439 let balance = client.get_balance_for_btc().await.map_err(|err| {
440 AdminGatewayError::Unexpected(anyhow!(
441 "Balance not available: {}",
442 err.fmt_compact_anyhow()
443 ))
444 })?;
445
446 wallet_module
451 .max_sendable_amount(balance, fee)
452 .await
453 .map_err(|err| AdminGatewayError::WithdrawError {
454 failure_reason: format!(
455 "Insufficient funds. Balance: {balance} Fee: {fee}: {}",
456 err.fmt_compact_anyhow()
457 ),
458 })?
459 }
460 BitcoinAmountOrAll::Amount(a) => a,
461 };
462
463 let operation_id = wallet_module
464 .send(
465 address.as_unchecked().clone(),
466 withdraw_amount,
467 Some(fee),
468 serde_json::Value::Null,
469 )
470 .await
471 .map_err(|e| AdminGatewayError::WithdrawError {
472 failure_reason: e.to_string(),
473 })?;
474
475 let result = wallet_module
476 .await_final_send_operation_state(operation_id)
477 .await
478 .map_err(|e| AdminGatewayError::WithdrawError {
479 failure_reason: e.to_string(),
480 })?;
481
482 let fees = PegOutFees::from_amount(fee);
483
484 match result {
485 fedimint_walletv2_client::FinalSendOperationState::Success(txid) => {
486 info!(target: LOG_GATEWAY, amount = %withdraw_amount, address = %address, "Sent funds via walletv2");
487 Ok(WithdrawResponse { txid, fees })
488 }
489 fedimint_walletv2_client::FinalSendOperationState::Aborted => {
490 Err(AdminGatewayError::WithdrawError {
491 failure_reason: "Withdrawal transaction was aborted".to_string(),
492 })
493 }
494 fedimint_walletv2_client::FinalSendOperationState::Failure => {
495 Err(AdminGatewayError::WithdrawError {
496 failure_reason: "Withdrawal failed".to_string(),
497 })
498 }
499 }
500}
501
502async fn calculate_max_withdrawable(
504 client: &ClientHandleArc,
505 address: &Address,
506) -> AdminResult<WithdrawDetails> {
507 let balance = client.get_balance_for_btc().await.map_err(|err| {
508 AdminGatewayError::Unexpected(anyhow!(
509 "Balance not available: {}",
510 err.fmt_compact_anyhow()
511 ))
512 })?;
513
514 if let Ok(wallet_module) =
515 client.get_first_module::<fedimint_walletv2_client::WalletClientModule>()
516 {
517 let fee = wallet_module
518 .send_fee()
519 .await
520 .map_err(|e| AdminGatewayError::WithdrawError {
521 failure_reason: e.to_string(),
522 })?;
523
524 let max_withdrawable = wallet_module
525 .max_sendable_amount(balance, fee)
526 .await
527 .map_err(|err| AdminGatewayError::WithdrawError {
528 failure_reason: err.fmt_compact_anyhow().to_string(),
529 })?;
530
531 let federation_fees = balance
536 .saturating_sub(Amount::from_sats(max_withdrawable.to_sat()))
537 .saturating_sub(Amount::from_sats(fee.to_sat()));
538
539 return Ok(WithdrawDetails {
540 amount: Amount::from_sats(max_withdrawable.to_sat()),
541 mint_fees: Some(federation_fees),
542 peg_out_fees: PegOutFees::from_amount(fee),
543 });
544 }
545
546 let Ok(wallet_module) = client.get_first_module::<WalletClientModule>() else {
547 return Err(AdminGatewayError::Unexpected(anyhow!(
548 "No wallet module found"
549 )));
550 };
551
552 let (max_withdrawable, peg_out_fees) = wallet_module
553 .max_withdrawable_amount(address, balance)
554 .await
555 .map_err(|err| AdminGatewayError::WithdrawError {
556 failure_reason: err.fmt_compact_anyhow().to_string(),
557 })?;
558
559 let federation_fees = balance
564 .saturating_sub(Amount::from_sats(max_withdrawable.to_sat()))
565 .saturating_sub(Amount::from_sats(peg_out_fees.amount().to_sat()));
566
567 Ok(WithdrawDetails {
568 amount: Amount::from_sats(max_withdrawable.to_sat()),
569 mint_fees: Some(federation_fees),
570 peg_out_fees,
571 })
572}
573
574impl Gateway {
575 fn get_bitcoind_client(
578 opts: &GatewayOpts,
579 network: bitcoin::Network,
580 gateway_id: &PublicKey,
581 ) -> anyhow::Result<(BitcoindClient, ChainSource)> {
582 let bitcoind_username = opts
583 .bitcoind_username
584 .clone()
585 .expect("FM_BITCOIND_URL is set but FM_BITCOIND_USERNAME is not");
586 let url = opts.bitcoind_url.clone().expect("No bitcoind url set");
587 let password = opts
588 .bitcoind_password
589 .clone()
590 .expect("FM_BITCOIND_URL is set but FM_BITCOIND_PASSWORD is not");
591
592 let chain_source = ChainSource::Bitcoind {
593 username: bitcoind_username.clone(),
594 password: password.clone(),
595 server_url: url.clone(),
596 };
597 let wallet_name = format!("gatewayd-{gateway_id}");
598 let client = BitcoindClient::new(&url, bitcoind_username, password, &wallet_name, network)?;
599 Ok((client, chain_source))
600 }
601
602 pub async fn new_with_default_modules(
605 mnemonic_sender: tokio::sync::broadcast::Sender<()>,
606 ) -> anyhow::Result<Gateway> {
607 let opts = GatewayOpts::parse();
608 let gateway_parameters = opts.to_gateway_parameters()?;
609 let decoders = ModuleDecoderRegistry::default();
610
611 let db_path = opts.data_dir.join(DB_FILE);
612 let gateway_db = match opts.db_backend {
613 DatabaseBackend::RocksDb => {
614 debug!(target: LOG_GATEWAY, "Using RocksDB database backend");
615 Database::new(
616 fedimint_rocksdb::RocksDb::build(db_path).open().await?,
617 decoders,
618 )
619 }
620 DatabaseBackend::CursedRedb => {
621 debug!(target: LOG_GATEWAY, "Using CursedRedb database backend");
622 Database::new(
623 fedimint_cursed_redb::MemAndRedb::new(db_path).await?,
624 decoders,
625 )
626 }
627 };
628
629 apply_migrations(
632 &gateway_db,
633 (),
634 "gatewayd".to_string(),
635 get_gatewayd_database_migrations(),
636 None,
637 None,
638 )
639 .await?;
640
641 let http_id = Self::load_or_create_gateway_keypair(&gateway_db, RegisteredProtocol::Http)
644 .await
645 .public_key();
646 let (dyn_bitcoin_rpc, chain_source) =
647 match (opts.bitcoind_url.as_ref(), opts.esplora_url.as_ref()) {
648 (Some(_), None) => {
649 let (client, chain_source) =
650 Self::get_bitcoind_client(&opts, gateway_parameters.network, &http_id)?;
651 (client.into_dyn(), chain_source)
652 }
653 (None, Some(url)) => {
654 let client = EsploraClient::new(url)
655 .expect("Could not create EsploraClient")
656 .into_dyn();
657 let chain_source = ChainSource::Esplora {
658 server_url: url.clone(),
659 };
660 (client, chain_source)
661 }
662 (Some(_), Some(_)) => {
663 let (client, chain_source) =
665 Self::get_bitcoind_client(&opts, gateway_parameters.network, &http_id)?;
666 (client.into_dyn(), chain_source)
667 }
668 _ => unreachable!("ArgGroup already enforced XOR relation"),
669 };
670
671 let mut registry = ClientModuleInitRegistry::new();
674 registry.attach(MintClientInit);
675 registry.attach(MintV2ClientInit);
676 registry.attach(WalletClientInit::new(dyn_bitcoin_rpc));
677 registry.attach(fedimint_walletv2_client::WalletClientInit);
678
679 let client_builder =
680 GatewayClientBuilder::new(opts.data_dir.clone(), registry, opts.db_backend).await?;
681
682 let gateway_state = if Self::load_mnemonic(&gateway_db).await.is_some() {
683 GatewayState::Disconnected
684 } else {
685 if gateway_parameters.skip_setup {
688 let mnemonic = if let Ok(words) = std::env::var(FM_GATEWAY_MNEMONIC_ENV) {
689 info!(target: LOG_GATEWAY, "Using provided mnemonic from environment variable");
690 Mnemonic::parse_in_normalized(Language::English, words.as_str()).map_err(
691 |e| {
692 AdminGatewayError::MnemonicError(anyhow!(format!(
693 "Seed phrase provided in environment was invalid {e:?}"
694 )))
695 },
696 )?
697 } else {
698 debug!(target: LOG_GATEWAY, "Generating mnemonic and writing entropy to client storage");
699 Bip39RootSecretStrategy::<12>::random(&mut OsRng)
700 };
701
702 Client::store_encodable_client_secret(&gateway_db, mnemonic.to_entropy())
703 .await
704 .map_err(AdminGatewayError::MnemonicError)?;
705 GatewayState::Disconnected
706 } else {
707 GatewayState::NotConfigured { mnemonic_sender }
708 }
709 };
710
711 info!(
712 target: LOG_GATEWAY,
713 version = %fedimint_build_code_version_env!(),
714 "Starting gatewayd",
715 );
716
717 Gateway::new(
718 opts.mode,
719 gateway_parameters,
720 gateway_db,
721 client_builder,
722 gateway_state,
723 chain_source,
724 )
725 .await
726 }
727
728 async fn new(
731 lightning_mode: LightningMode,
732 gateway_parameters: GatewayParameters,
733 gateway_db: Database,
734 client_builder: GatewayClientBuilder,
735 gateway_state: GatewayState,
736 chain_source: ChainSource,
737 ) -> anyhow::Result<Gateway> {
738 let num_route_hints = gateway_parameters.num_route_hints;
739 let network = gateway_parameters.network;
740
741 let task_group = TaskGroup::new();
742 task_group.install_kill_handler();
743
744 let mut registrations = BTreeMap::new();
745 if let Some(http_url) = gateway_parameters.versioned_api {
746 registrations.insert(
747 RegisteredProtocol::Http,
748 Registration::new(&gateway_db, http_url, RegisteredProtocol::Http).await,
749 );
750 }
751
752 let iroh_sk = Self::load_or_create_iroh_key(&gateway_db).await;
753 if gateway_parameters.iroh_listen.is_some() {
754 let endpoint_url = SafeUrl::parse(&format!("iroh://{}", iroh_sk.public()))?;
755 registrations.insert(
756 RegisteredProtocol::Iroh,
757 Registration::new(&gateway_db, endpoint_url, RegisteredProtocol::Iroh).await,
758 );
759 }
760
761 Ok(Self {
762 federation_manager: Arc::new(RwLock::new(FederationManager::new())),
763 lightning_mode,
764 state: Arc::new(RwLock::new(gateway_state)),
765 client_builder,
766 gateway_db: gateway_db.clone(),
767 listen: gateway_parameters.listen,
768 metrics_listen: gateway_parameters.metrics_listen,
769 task_group,
770 bcrypt_password_hash: gateway_parameters.bcrypt_password_hash.to_string(),
771 bcrypt_liquidity_manager_password_hash: gateway_parameters
772 .bcrypt_liquidity_manager_password_hash
773 .map(|h| h.to_string()),
774 num_route_hints,
775 network,
776 chain_source,
777 default_routing_fees: gateway_parameters.default_routing_fees,
778 default_transaction_fees: gateway_parameters.default_transaction_fees,
779 iroh_sk,
780 iroh_dns: gateway_parameters.iroh_dns,
781 iroh_relays: gateway_parameters.iroh_relays,
782 iroh_listen: gateway_parameters.iroh_listen,
783 registrations,
784 registration_health: RegistrationHealthTracker::default(),
785 invoice_rate_limiter: Arc::new(TokenBucketRateLimiter::new(
786 gateway_parameters.invoice_rate_limit_burst,
787 gateway_parameters.invoice_rate_limit_per_second,
788 )),
789 })
790 }
791
792 async fn load_or_create_gateway_keypair(
793 gateway_db: &Database,
794 protocol: RegisteredProtocol,
795 ) -> secp256k1::Keypair {
796 let mut dbtx = gateway_db.begin_transaction().await;
797 let keypair = dbtx.load_or_create_gateway_keypair(protocol).await;
798 dbtx.commit_tx().await;
799 keypair
800 }
801
802 async fn load_or_create_iroh_key(gateway_db: &Database) -> iroh::SecretKey {
805 let mut dbtx = gateway_db.begin_transaction().await;
806 let iroh_sk = dbtx.load_or_create_iroh_key().await;
807 dbtx.commit_tx().await;
808 iroh_sk
809 }
810
811 pub async fn http_gateway_id(&self) -> PublicKey {
812 Self::load_or_create_gateway_keypair(&self.gateway_db, RegisteredProtocol::Http)
813 .await
814 .public_key()
815 }
816
817 async fn get_state(&self) -> GatewayState {
818 self.state.read().await.clone()
819 }
820
821 pub async fn dump_database(
824 dbtx: &mut DatabaseTransaction<'_>,
825 prefix_names: Vec<String>,
826 ) -> BTreeMap<String, Box<dyn erased_serde::Serialize + Send>> {
827 dbtx.dump_database(prefix_names).await
828 }
829
830 pub async fn run(
835 self,
836 runtime: Arc<tokio::runtime::Runtime>,
837 mnemonic_receiver: tokio::sync::broadcast::Receiver<()>,
838 ) -> anyhow::Result<TaskShutdownToken> {
839 install_crypto_provider().await;
840 self.register_clients_timer();
841 self.load_clients().await?;
842 self.start_gateway(runtime, mnemonic_receiver.resubscribe());
843 self.spawn_backup_task();
844 self.spawn_prune_registered_contracts_task();
845 fedimint_metrics::spawn_api_server(self.metrics_listen, self.task_group.clone()).await?;
847 let handle = self.task_group.make_handle();
849 run_webserver(Arc::new(self), mnemonic_receiver.resubscribe()).await?;
850 let shutdown_receiver = handle.make_shutdown_rx();
851 Ok(shutdown_receiver)
852 }
853
854 fn spawn_backup_task(&self) {
857 let self_copy = self.clone();
858 self.task_group
859 .spawn_cancellable_silent("backup ecash", async move {
860 const BACKUP_UPDATE_INTERVAL: Duration = Duration::from_hours(1);
861 let mut interval = tokio::time::interval(BACKUP_UPDATE_INTERVAL);
862 interval.tick().await;
863 loop {
864 {
865 let mut dbtx = self_copy.gateway_db.begin_transaction().await;
866 self_copy.backup_all_federations(&mut dbtx).await;
867 dbtx.commit_tx().await;
868 interval.tick().await;
869 }
870 }
871 });
872 }
873
874 fn spawn_prune_registered_contracts_task(&self) {
878 let self_copy = self.clone();
879 self.task_group.spawn_cancellable_silent(
880 "prune registered incoming contracts",
881 async move {
882 const PRUNE_INTERVAL: Duration = Duration::from_hours(1);
883 const RETENTION_AFTER_EXPIRY: Duration = Duration::from_hours(24);
887
888 let mut interval = tokio::time::interval(PRUNE_INTERVAL);
889 loop {
890 interval.tick().await;
891
892 let cutoff_secs = duration_since_epoch()
893 .saturating_sub(RETENTION_AFTER_EXPIRY)
894 .as_secs();
895
896 let mut dbtx = self_copy.gateway_db.begin_transaction().await;
897 let num_pruned = dbtx.prune_registered_incoming_contracts(cutoff_secs).await;
898 match dbtx.commit_tx_result().await {
899 Ok(()) => {
900 if num_pruned > 0 {
901 info!(
902 target: LOG_GATEWAY,
903 num_pruned,
904 "Pruned expired incoming contract records"
905 );
906 }
907 }
908 Err(err) => {
909 warn!(
910 target: LOG_GATEWAY,
911 err = %err.fmt_compact(),
912 "Failed to prune expired incoming contract records"
913 );
914 }
915 }
916 }
917 },
918 );
919 }
920
921 pub async fn backup_all_federations(&self, dbtx: &mut DatabaseTransaction<'_, Committable>) {
925 const BACKUP_THRESHOLD_DURATION: Duration = Duration::from_hours(24);
928
929 let now = fedimint_core::time::now();
930 let threshold = now
931 .checked_sub(BACKUP_THRESHOLD_DURATION)
932 .expect("Cannot be negative");
933 for (id, last_backup) in dbtx.load_backup_records().await {
934 match last_backup {
935 Some(backup_time) if backup_time < threshold => {
936 let fed_manager = self.federation_manager.read().await;
937 fed_manager.backup_federation(&id, dbtx, now).await;
938 }
939 None => {
940 let fed_manager = self.federation_manager.read().await;
941 fed_manager.backup_federation(&id, dbtx, now).await;
942 }
943 _ => {}
944 }
945 }
946 }
947
948 fn start_gateway(
951 &self,
952 runtime: Arc<tokio::runtime::Runtime>,
953 mut mnemonic_receiver: tokio::sync::broadcast::Receiver<()>,
954 ) {
955 const PAYMENT_STREAM_RETRY_SECONDS: u64 = 60;
956
957 let self_copy = self.clone();
958 let tg = self.task_group.clone();
959 self.task_group.spawn(
960 "Subscribe to intercepted lightning payments in stream",
961 |handle| async move {
962 loop {
964 if handle.is_shutting_down() {
965 info!(target: LOG_GATEWAY, "Gateway lightning payment stream handler loop is shutting down");
966 break;
967 }
968
969 if let GatewayState::NotConfigured{ .. } = self_copy.get_state().await {
970 info!(
971 target: LOG_GATEWAY,
972 "Waiting for the mnemonic to be set before starting lightning receive loop."
973 );
974 info!(
975 target: LOG_GATEWAY,
976 "You might need to provide it from the UI or refer to documentation w.r.t how to initialize it."
977 );
978
979 let _ = mnemonic_receiver.recv().await;
980 info!(
981 target: LOG_GATEWAY,
982 "Received mnemonic, attempting to start lightning receive loop"
983 );
984 }
985
986 let payment_stream_task_group = tg.make_subgroup();
987 let lnrpc_route = self_copy.create_lightning_client(runtime.clone()).await;
988
989 debug!(target: LOG_GATEWAY, "Establishing lightning payment stream...");
990 let (stream, ln_client) = match lnrpc_route.route_htlcs(&payment_stream_task_group).await
991 {
992 Ok((stream, ln_client)) => (stream, ln_client),
993 Err(err) => {
994 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failed to open lightning payment stream");
995 if let Err(err) = payment_stream_task_group.shutdown_join_all(None).await {
1003 crit!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Lightning payment stream task group shutdown");
1004 }
1005 sleep(Duration::from_secs(PAYMENT_STREAM_RETRY_SECONDS)).await;
1006 continue
1007 }
1008 };
1009
1010 self_copy.set_gateway_state(GatewayState::Connected).await;
1012 info!(target: LOG_GATEWAY, "Established lightning payment stream");
1013
1014 let route_payments_response =
1015 self_copy.route_lightning_payments(&handle, stream, ln_client).await;
1016
1017 self_copy.set_gateway_state(GatewayState::Disconnected).await;
1018 if let Err(err) = payment_stream_task_group.shutdown_join_all(None).await {
1019 crit!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Lightning payment stream task group shutdown");
1020 }
1021
1022 self_copy.unannounce_from_all_federations().await;
1023
1024 match route_payments_response {
1025 ReceivePaymentStreamAction::RetryAfterDelay => {
1026 warn!(target: LOG_GATEWAY, retry_interval = %PAYMENT_STREAM_RETRY_SECONDS, "Disconnected from lightning node");
1027 sleep(Duration::from_secs(PAYMENT_STREAM_RETRY_SECONDS)).await;
1028 }
1029 ReceivePaymentStreamAction::NoRetry => break,
1030 }
1031 }
1032 },
1033 );
1034 }
1035
1036 async fn route_lightning_payments<'a>(
1040 &'a self,
1041 handle: &TaskHandle,
1042 mut stream: RouteHtlcStream<'a>,
1043 ln_client: Arc<dyn ILnRpcClient>,
1044 ) -> ReceivePaymentStreamAction {
1045 let LightningInfo::Connected {
1046 public_key: lightning_public_key,
1047 alias: lightning_alias,
1048 network: lightning_network,
1049 block_height: _,
1050 synced_to_chain,
1051 } = ln_client.parsed_node_info().await
1052 else {
1053 warn!(target: LOG_GATEWAY, "Failed to retrieve Lightning info");
1054 return ReceivePaymentStreamAction::RetryAfterDelay;
1055 };
1056
1057 assert!(
1058 self.network == lightning_network,
1059 "Lightning node network does not match Gateway's network. LN: {lightning_network} Gateway: {}",
1060 self.network
1061 );
1062
1063 if synced_to_chain || is_env_var_set(FM_GATEWAY_SKIP_WAIT_FOR_SYNC_ENV) {
1064 info!(target: LOG_GATEWAY, "Gateway is already synced to chain");
1065 } else {
1066 self.set_gateway_state(GatewayState::Syncing).await;
1067 info!(target: LOG_GATEWAY, "Waiting for chain sync");
1068 if let Err(err) = ln_client.wait_for_chain_sync().await {
1069 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failed to wait for chain sync");
1070 return ReceivePaymentStreamAction::RetryAfterDelay;
1071 }
1072 }
1073
1074 let lightning_context = LightningContext {
1075 lnrpc: LnRpcTracked::new(ln_client, "gateway"),
1076 lightning_public_key,
1077 lightning_alias,
1078 lightning_network,
1079 };
1080 if let GatewayState::ShuttingDown { .. } = self
1081 .set_gateway_state(GatewayState::Running { lightning_context })
1082 .await
1083 {
1084 info!(
1085 target: LOG_GATEWAY,
1086 "Reconnected to the lightning node while shutting down, not accepting payments"
1087 );
1088 } else {
1089 info!(target: LOG_GATEWAY, "Gateway is running");
1090 }
1091
1092 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
1093 let mut dbtx = self.gateway_db.begin_transaction_nc().await;
1096 let all_federations_configs =
1097 dbtx.load_federation_configs().await.into_iter().collect();
1098 self.register_federations(&all_federations_configs, &self.task_group)
1099 .await;
1100 }
1101
1102 let htlc_task_group = self.task_group.make_subgroup();
1105 if handle
1106 .cancel_on_shutdown(async move {
1107 loop {
1108 let payment_request_or = tokio::select! {
1109 payment_request_or = stream.next() => {
1110 payment_request_or
1111 }
1112 () = self.is_shutting_down_safely() => {
1113 break;
1114 }
1115 };
1116
1117 let Some(payment_request) = payment_request_or else {
1118 warn!(
1119 target: LOG_GATEWAY,
1120 "Unexpected response from incoming lightning payment stream. Shutting down payment processor"
1121 );
1122 break;
1123 };
1124
1125 let state_guard = self.state.read().await;
1126 if let GatewayState::Running { ref lightning_context } = *state_guard {
1127 let gateway = self.clone();
1129 let lightning_context = lightning_context.clone();
1130 htlc_task_group.spawn_cancellable_silent(
1131 "handle_lightning_payment",
1132 async move {
1133 let start = fedimint_core::time::now();
1134 let outcome = gateway
1135 .handle_lightning_payment(payment_request, &lightning_context)
1136 .await;
1137 metrics::HTLC_HANDLING_DURATION_SECONDS
1138 .with_label_values(&[outcome])
1139 .observe(
1140 fedimint_core::time::now()
1141 .duration_since(start)
1142 .unwrap_or_default()
1143 .as_secs_f64(),
1144 );
1145 },
1146 );
1147 } else {
1148 warn!(
1149 target: LOG_GATEWAY,
1150 state = %state_guard,
1151 "Gateway isn't in a running state, cannot handle incoming payments."
1152 );
1153 break;
1154 }
1155 }
1156 })
1157 .await
1158 .is_ok()
1159 {
1160 warn!(target: LOG_GATEWAY, "Lightning payment stream connection broken. Gateway is disconnected");
1161 ReceivePaymentStreamAction::RetryAfterDelay
1162 } else {
1163 info!(target: LOG_GATEWAY, "Received shutdown signal");
1164 ReceivePaymentStreamAction::NoRetry
1165 }
1166 }
1167
1168 async fn is_shutting_down_safely(&self) {
1171 loop {
1172 if let GatewayState::ShuttingDown { .. } = self.get_state().await {
1173 return;
1174 }
1175
1176 fedimint_core::task::sleep(Duration::from_secs(1)).await;
1177 }
1178 }
1179
1180 async fn handle_lightning_payment(
1190 &self,
1191 payment_request: InterceptPaymentRequest,
1192 lightning_context: &LightningContext,
1193 ) -> &'static str {
1194 info!(
1195 target: LOG_GATEWAY,
1196 lightning_payment = %PrettyInterceptPaymentRequest(&payment_request),
1197 "Intercepting lightning payment",
1198 );
1199
1200 let lnv2_start = fedimint_core::time::now();
1201 let lnv2_result = self
1202 .try_handle_lightning_payment_lnv2(&payment_request, lightning_context)
1203 .await;
1204 let lnv2_outcome = if lnv2_result.is_ok() {
1205 "success"
1206 } else {
1207 "error"
1208 };
1209 metrics::HTLC_LNV2_ATTEMPT_DURATION_SECONDS
1210 .with_label_values(&[lnv2_outcome])
1211 .observe(
1212 fedimint_core::time::now()
1213 .duration_since(lnv2_start)
1214 .unwrap_or_default()
1215 .as_secs_f64(),
1216 );
1217 if lnv2_result.is_ok() {
1218 return "lnv2";
1219 }
1220
1221 let lnv1_start = fedimint_core::time::now();
1222 let lnv1_result = self
1223 .try_handle_lightning_payment_ln_legacy(&payment_request, lightning_context)
1224 .await;
1225 let lnv1_outcome = if lnv1_result.is_ok() {
1226 "success"
1227 } else {
1228 "error"
1229 };
1230 metrics::HTLC_LNV1_ATTEMPT_DURATION_SECONDS
1231 .with_label_values(&[lnv1_outcome])
1232 .observe(
1233 fedimint_core::time::now()
1234 .duration_since(lnv1_start)
1235 .unwrap_or_default()
1236 .as_secs_f64(),
1237 );
1238 if lnv1_result.is_ok() {
1239 return "lnv1";
1240 }
1241
1242 let is_federation_scid = match payment_request.short_channel_id {
1247 Some(scid) => self
1248 .federation_manager
1249 .read()
1250 .await
1251 .get_client_for_index(scid)
1252 .is_some(),
1253 None => false,
1254 };
1255
1256 if is_federation_scid {
1257 warn!(
1264 target: LOG_GATEWAY,
1265 payment_hash = %payment_request.payment_hash,
1266 short_channel_id = ?payment_request.short_channel_id,
1267 amount_msat = payment_request.amount_msat,
1268 incoming_chan_id = payment_request.incoming_chan_id,
1269 htlc_id = payment_request.htlc_id,
1270 lnv2_err = ?lnv2_result.as_ref().err(),
1271 lnv1_err = ?lnv1_result.as_ref().err(),
1272 "Unmatched lightning payment for federation scid: cancelling HTLC",
1273 );
1274 Self::cancel_unmatched_lightning_payment(payment_request, lightning_context).await;
1275 "cancel"
1276 } else {
1277 Self::forward_lightning_payment(payment_request, lightning_context).await;
1282 "forward"
1283 }
1284 }
1285
1286 async fn try_handle_lightning_payment_lnv2(
1289 &self,
1290 htlc_request: &InterceptPaymentRequest,
1291 lightning_context: &LightningContext,
1292 ) -> Result<()> {
1293 let (contract, client) = self
1299 .get_registered_incoming_contract_and_client_v2(
1300 PaymentImage::Hash(htlc_request.payment_hash),
1301 htlc_request.amount_msat,
1302 )
1303 .await?;
1304
1305 if let Err(err) = client
1306 .get_first_module::<GatewayClientModuleV2>()
1307 .expect("Must have client module")
1308 .relay_incoming_htlc(
1309 htlc_request.payment_hash,
1310 htlc_request.incoming_chan_id,
1311 htlc_request.htlc_id,
1312 contract,
1313 htlc_request.amount_msat,
1314 )
1315 .await
1316 {
1317 warn!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Error relaying incoming lightning payment");
1318
1319 let outcome = InterceptPaymentResponse {
1320 action: PaymentAction::Cancel,
1321 payment_hash: htlc_request.payment_hash,
1322 incoming_chan_id: htlc_request.incoming_chan_id,
1323 htlc_id: htlc_request.htlc_id,
1324 };
1325
1326 if let Err(err) = lightning_context.lnrpc.complete_htlc(outcome).await {
1327 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Error sending HTLC response to lightning node");
1328 }
1329 }
1330
1331 Ok(())
1332 }
1333
1334 async fn try_handle_lightning_payment_ln_legacy(
1337 &self,
1338 htlc_request: &InterceptPaymentRequest,
1339 lightning_context: &LightningContext,
1340 ) -> Result<()> {
1341 let Some(federation_index) = htlc_request.short_channel_id else {
1343 return Err(PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1344 "Incoming payment has not last hop short channel id".to_string(),
1345 )));
1346 };
1347
1348 let Some(client) = self
1349 .federation_manager
1350 .read()
1351 .await
1352 .get_client_for_index(federation_index)
1353 else {
1354 return Err(PublicGatewayError::LNv1(LNv1Error::IncomingPayment("Incoming payment has a last hop short channel id that does not map to a known federation".to_string())));
1355 };
1356
1357 client
1362 .borrow()
1363 .with(|client| async {
1364 let htlc = htlc_request.clone().try_into();
1365 match htlc {
1366 Ok(htlc) => {
1367 let lnv1 =
1368 client
1369 .get_first_module::<GatewayClientModule>()
1370 .map_err(|_| {
1371 PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1372 "Federation does not have LNv1 module".to_string(),
1373 ))
1374 })?;
1375 match lnv1
1376 .gateway_handle_intercepted_htlc(htlc, async {
1377 Ok(lightning_context.lnrpc.info().await?.block_height)
1378 })
1379 .await
1380 {
1381 Ok(_) => Ok(()),
1382 Err(e) => Err(PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1383 format!("Error intercepting lightning payment {e:?}"),
1384 ))),
1385 }
1386 }
1387 _ => Err(PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1388 "Could not convert InterceptHtlcResult into an HTLC".to_string(),
1389 ))),
1390 }
1391 })
1392 .await
1393 }
1394
1395 async fn cancel_unmatched_lightning_payment(
1409 htlc_request: InterceptPaymentRequest,
1410 lightning_context: &LightningContext,
1411 ) {
1412 let outcome = InterceptPaymentResponse {
1413 action: PaymentAction::Cancel,
1414 payment_hash: htlc_request.payment_hash,
1415 incoming_chan_id: htlc_request.incoming_chan_id,
1416 htlc_id: htlc_request.htlc_id,
1417 };
1418
1419 if let Err(err) = lightning_context.lnrpc.complete_htlc(outcome).await {
1420 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Error sending lightning payment response to lightning node");
1421 }
1422 }
1423
1424 async fn forward_lightning_payment(
1429 htlc_request: InterceptPaymentRequest,
1430 lightning_context: &LightningContext,
1431 ) {
1432 let outcome = InterceptPaymentResponse {
1433 action: PaymentAction::Forward,
1434 payment_hash: htlc_request.payment_hash,
1435 incoming_chan_id: htlc_request.incoming_chan_id,
1436 htlc_id: htlc_request.htlc_id,
1437 };
1438
1439 if let Err(err) = lightning_context.lnrpc.complete_htlc(outcome).await {
1440 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Error sending lightning payment response to lightning node");
1441 }
1442 }
1443
1444 async fn set_gateway_state(&self, state: GatewayState) -> GatewayState {
1455 let mut lock = self.state.write().await;
1456
1457 if let GatewayState::ShuttingDown { .. } = *lock {
1458 match state {
1459 GatewayState::Running { lightning_context } => {
1460 *lock = GatewayState::ShuttingDown { lightning_context };
1461 }
1462 ignored => {
1463 info!(
1464 target: LOG_GATEWAY,
1465 ignored_state = %ignored,
1466 "Gateway is shutting down, ignoring state change"
1467 );
1468 }
1469 }
1470 } else {
1471 *lock = state;
1472 }
1473
1474 lock.clone()
1475 }
1476
1477 #[doc(hidden)]
1485 pub async fn set_gateway_state_out_of_band(&self, state: GatewayState) {
1486 self.set_gateway_state(state).await;
1487 }
1488
1489 pub async fn handle_get_federation_config(
1492 &self,
1493 federation_id_or: Option<FederationId>,
1494 ) -> AdminResult<GatewayFedConfig> {
1495 if !matches!(self.get_state().await, GatewayState::Running { .. }) {
1496 return Ok(GatewayFedConfig {
1497 federations: BTreeMap::new(),
1498 });
1499 }
1500
1501 let federations = if let Some(federation_id) = federation_id_or {
1502 let mut federations = BTreeMap::new();
1503 federations.insert(
1504 federation_id,
1505 self.federation_manager
1506 .read()
1507 .await
1508 .get_federation_config(federation_id)
1509 .await?,
1510 );
1511 federations
1512 } else {
1513 self.federation_manager
1514 .read()
1515 .await
1516 .get_all_federation_configs()
1517 .await
1518 };
1519
1520 Ok(GatewayFedConfig { federations })
1521 }
1522
1523 pub async fn handle_address_msg(&self, payload: DepositAddressPayload) -> AdminResult<Address> {
1526 let client = self.select_client(payload.federation_id).await?;
1527
1528 if let Ok(wallet_module) = client.value().get_first_module::<WalletClientModule>() {
1529 let address = wallet_module
1530 .allocate_deposit_address_expert_only(())
1531 .await?
1532 .address;
1533 Ok(address)
1534 } else if let Ok(wallet_module) = client
1535 .value()
1536 .get_first_module::<fedimint_walletv2_client::WalletClientModule>()
1537 {
1538 Ok(wallet_module.receive().await)
1539 } else {
1540 Err(AdminGatewayError::Unexpected(anyhow!(
1541 "No wallet module found"
1542 )))
1543 }
1544 }
1545
1546 async fn handle_pay_invoice_msg(
1549 &self,
1550 payload: fedimint_ln_client::pay::PayInvoicePayload,
1551 ) -> Result<Preimage> {
1552 let GatewayState::Running { .. } = self.get_state().await else {
1553 return Err(PublicGatewayError::Lightning(
1554 LightningRpcError::FailedToConnect,
1555 ));
1556 };
1557
1558 debug!(target: LOG_GATEWAY, "Handling pay invoice message");
1559 let client = self.select_client(payload.federation_id).await?;
1560 let contract_id = payload.contract_id;
1561 let gateway_module = &client
1562 .value()
1563 .get_first_module::<GatewayClientModule>()
1564 .map_err(LNv1Error::OutgoingPayment)
1565 .map_err(PublicGatewayError::LNv1)?;
1566 let operation_id = gateway_module
1567 .gateway_pay_bolt11_invoice(payload)
1568 .await
1569 .map_err(LNv1Error::OutgoingPayment)
1570 .map_err(PublicGatewayError::LNv1)?;
1571 let mut updates = gateway_module
1572 .gateway_subscribe_ln_pay(operation_id)
1573 .await
1574 .map_err(LNv1Error::OutgoingPayment)
1575 .map_err(PublicGatewayError::LNv1)?
1576 .into_stream();
1577 while let Some(update) = updates.next().await {
1578 match update {
1579 GatewayExtPayStates::Success { preimage, .. } => {
1580 debug!(target: LOG_GATEWAY, contract_id = %contract_id, "Successfully paid invoice");
1581 return Ok(preimage);
1582 }
1583 GatewayExtPayStates::Fail {
1584 error,
1585 error_message,
1586 } => {
1587 return Err(PublicGatewayError::LNv1(LNv1Error::OutgoingContract {
1588 error: Box::new(error),
1589 message: format!(
1590 "{error_message} while paying invoice with contract id {contract_id}"
1591 ),
1592 }));
1593 }
1594 GatewayExtPayStates::Canceled { error } => {
1595 return Err(PublicGatewayError::LNv1(LNv1Error::OutgoingContract {
1596 error: Box::new(error.clone()),
1597 message: format!(
1598 "Cancelled with {error} while paying invoice with contract id {contract_id}"
1599 ),
1600 }));
1601 }
1602 GatewayExtPayStates::Created => {
1603 debug!(target: LOG_GATEWAY, contract_id = %contract_id, "Start pay invoice state machine");
1604 }
1605 other => {
1606 debug!(target: LOG_GATEWAY, state = ?other, contract_id = %contract_id, "Got state while paying invoice");
1607 }
1608 }
1609 }
1610
1611 Err(PublicGatewayError::LNv1(LNv1Error::OutgoingPayment(
1612 anyhow!("Ran out of state updates while paying invoice"),
1613 )))
1614 }
1615
1616 pub async fn handle_backup_msg(
1619 &self,
1620 BackupPayload { federation_id }: BackupPayload,
1621 ) -> AdminResult<()> {
1622 let federation_manager = self.federation_manager.read().await;
1623 let client = federation_manager
1624 .client(&federation_id)
1625 .ok_or(AdminGatewayError::ClientCreationError(anyhow::anyhow!(
1626 format!("Gateway has not connected to {federation_id}")
1627 )))?
1628 .value();
1629 let metadata: BTreeMap<String, String> = BTreeMap::new();
1630 #[allow(deprecated)]
1631 client
1632 .backup_to_federation(fedimint_client::backup::Metadata::from_json_serialized(
1633 metadata,
1634 ))
1635 .await?;
1636 Ok(())
1637 }
1638
1639 pub async fn handle_recheck_address_msg(
1641 &self,
1642 payload: DepositAddressRecheckPayload,
1643 ) -> AdminResult<()> {
1644 let client = self.select_client(payload.federation_id).await?;
1645
1646 if let Ok(wallet_module) = client.value().get_first_module::<WalletClientModule>() {
1647 wallet_module
1648 .recheck_pegin_address_by_address(payload.address)
1649 .await?;
1650 Ok(())
1651 } else if client
1652 .value()
1653 .get_first_module::<fedimint_walletv2_client::WalletClientModule>()
1654 .is_ok()
1655 {
1656 Ok(())
1658 } else {
1659 Err(AdminGatewayError::Unexpected(anyhow!(
1660 "No wallet module found"
1661 )))
1662 }
1663 }
1664
1665 pub async fn handle_receive_ecash_msg(
1667 &self,
1668 payload: ReceiveEcashPayload,
1669 ) -> Result<ReceiveEcashResponse> {
1670 let federation_id_prefix = base32::decode_prefixed::<fedimint_mintv2_client::ECash>(
1672 FEDIMINT_PREFIX,
1673 &payload.notes,
1674 )
1675 .ok()
1676 .and_then(|e| e.mint())
1677 .map(|id| id.to_prefix())
1678 .or_else(|| {
1679 OOBNotes::from_str(&payload.notes)
1680 .ok()
1681 .map(|n| n.federation_id_prefix())
1682 })
1683 .ok_or_else(|| PublicGatewayError::ReceiveEcashError {
1684 failure_reason: "Invalid ecash format: could not parse as ECash or OOBNotes"
1685 .to_string(),
1686 })?;
1687
1688 let client = self
1689 .federation_manager
1690 .read()
1691 .await
1692 .get_client_for_federation_id_prefix(federation_id_prefix)
1693 .ok_or(FederationNotConnected {
1694 federation_id_prefix,
1695 })?;
1696
1697 if let Ok(mint) = client.value().get_first_module::<MintClientModule>() {
1699 let notes = OOBNotes::from_str(&payload.notes).map_err(|e| {
1700 PublicGatewayError::ReceiveEcashError {
1701 failure_reason: format!("Expected OOBNotes for MintV1 federation: {e}"),
1702 }
1703 })?;
1704 let amount = notes.total_amount();
1705
1706 let operation_id = mint.reissue_external_notes(notes, ()).await.map_err(|e| {
1707 PublicGatewayError::ReceiveEcashError {
1708 failure_reason: e.to_string(),
1709 }
1710 })?;
1711
1712 let mut updates = mint
1713 .subscribe_reissue_external_notes(operation_id)
1714 .await
1715 .map_err(|e| PublicGatewayError::ReceiveEcashError {
1716 failure_reason: format!("Could not subscribe to reissue operation: {e}"),
1717 })?
1718 .into_stream();
1719
1720 let mut reissued = false;
1725 while let Some(update) = updates.next().await {
1726 match update {
1727 ReissueExternalNotesState::Failed(failure_reason) => {
1728 return Err(PublicGatewayError::ReceiveEcashError { failure_reason });
1729 }
1730 ReissueExternalNotesState::Done => reissued = true,
1731 ReissueExternalNotesState::Created | ReissueExternalNotesState::Issuing => {}
1732 }
1733 }
1734
1735 if !reissued {
1736 return Err(PublicGatewayError::ReceiveEcashError {
1737 failure_reason: "Reissue operation ended before the notes were reissued"
1738 .to_string(),
1739 });
1740 }
1741
1742 Ok(ReceiveEcashResponse { amount })
1743 } else if let Ok(mint) = client.value().get_first_module::<MintV2ClientModule>() {
1744 let ecash: fedimint_mintv2_client::ECash =
1745 base32::decode_prefixed(FEDIMINT_PREFIX, &payload.notes).map_err(|e| {
1746 PublicGatewayError::ReceiveEcashError {
1747 failure_reason: format!("Expected ECash for MintV2 federation: {e}"),
1748 }
1749 })?;
1750 let amount = ecash.amount();
1751
1752 let operation_id = mint
1753 .receive(ecash, serde_json::Value::Null)
1754 .await
1755 .map_err(|e| PublicGatewayError::ReceiveEcashError {
1756 failure_reason: e.to_string(),
1757 })?;
1758
1759 let final_state = mint
1760 .await_final_receive_operation_state(operation_id)
1761 .await
1762 .map_err(|e| PublicGatewayError::ReceiveEcashError {
1763 failure_reason: e.to_string(),
1764 })?;
1765 match final_state {
1766 fedimint_mintv2_client::FinalReceiveOperationState::Success => {}
1767 fedimint_mintv2_client::FinalReceiveOperationState::Rejected => {
1768 return Err(PublicGatewayError::ReceiveEcashError {
1769 failure_reason: "ECash receive was rejected".to_string(),
1770 });
1771 }
1772 }
1773
1774 Ok(ReceiveEcashResponse { amount })
1775 } else {
1776 Err(PublicGatewayError::ReceiveEcashError {
1777 failure_reason: "No mint module found".to_string(),
1778 })
1779 }
1780 }
1781
1782 pub async fn handle_get_invoice_msg(
1785 &self,
1786 payload: GetInvoiceRequest,
1787 ) -> AdminResult<Option<GetInvoiceResponse>> {
1788 let lightning_context = self.get_lightning_context().await?;
1789 let invoice = lightning_context.lnrpc.get_invoice(payload).await?;
1790 Ok(invoice)
1791 }
1792
1793 pub async fn handle_withdraw_to_onchain_msg(
1796 &self,
1797 payload: WithdrawToOnchainPayload,
1798 ) -> AdminResult<WithdrawResponse> {
1799 let address = self.handle_get_ln_onchain_address_msg().await?;
1800 let withdraw = WithdrawPayload {
1801 address: address.into_unchecked(),
1802 federation_id: payload.federation_id,
1803 amount: payload.amount,
1804 quoted_fees: None,
1805 };
1806 self.handle_withdraw_msg(withdraw).await
1807 }
1808
1809 pub async fn handle_pegin_from_onchain_msg(
1812 &self,
1813 payload: PeginFromOnchainPayload,
1814 ) -> AdminResult<Txid> {
1815 let deposit = DepositAddressPayload {
1816 federation_id: payload.federation_id,
1817 };
1818 let address = self.handle_address_msg(deposit).await?;
1819 let send_onchain = SendOnchainRequest {
1820 address: address.into_unchecked(),
1821 amount: payload.amount,
1822 fee_rate_sats_per_vbyte: payload.fee_rate_sats_per_vbyte,
1823 };
1824 let txid = self.handle_send_onchain_msg(send_onchain).await?;
1825
1826 Ok(txid)
1827 }
1828
1829 async fn register_federations(
1837 &self,
1838 federations: &BTreeMap<FederationId, FederationConfig>,
1839 register_task_group: &TaskGroup,
1840 ) {
1841 if let GatewayState::ShuttingDown { .. } = self.get_state().await {
1842 info!(
1843 target: LOG_GATEWAY,
1844 "Gateway is shutting down, skipping federation registration"
1845 );
1846 return;
1847 }
1848
1849 if let Ok(lightning_context) = self.get_lightning_context().await {
1850 let route_hints = lightning_context
1851 .lnrpc
1852 .parsed_route_hints(self.num_route_hints)
1853 .await;
1854 if route_hints.is_empty() {
1855 warn!(target: LOG_GATEWAY, "Gateway did not retrieve any route hints, may reduce receive success rate.");
1856 }
1857
1858 for (federation_id, federation_config) in federations {
1859 let routing_fees = match RoutingFees::try_from(federation_config.lightning_fee) {
1864 Ok(routing_fees) => routing_fees,
1865 Err(err) => {
1866 warn!(
1867 target: LOG_GATEWAY,
1868 %federation_id,
1869 err = %err.fmt_compact(),
1870 "Skipping registration, the configured lightning fee cannot be announced. Set a smaller fee with `set_fees`."
1871 );
1872 continue;
1873 }
1874 };
1875
1876 let fed_manager = self.federation_manager.read().await;
1877 if let Some(client) = fed_manager.client(federation_id) {
1878 let federation_id = *federation_id;
1879 let client_arc = client.clone().into_value();
1880 let route_hints = route_hints.clone();
1881 let lightning_context = lightning_context.clone();
1882 let registration_health = self.registration_health.clone();
1883 let registrations = self
1884 .registrations
1885 .clone()
1886 .into_iter()
1887 .map(|(protocol, registration)| {
1888 let attempt =
1889 registration_health.begin_lnv1_attempt(federation_id, protocol);
1890 (registration, attempt)
1891 })
1892 .collect::<Vec<_>>();
1893
1894 register_task_group.spawn_cancellable_silent(
1895 "register federation",
1896 async move {
1897 let Ok(gateway_client) =
1898 client_arc.get_first_module::<GatewayClientModule>()
1899 else {
1900 return;
1901 };
1902
1903 for (registration, attempt) in registrations {
1904 let succeeded = gateway_client
1905 .try_register_with_federation(
1906 route_hints.clone(),
1907 GW_ANNOUNCEMENT_TTL,
1908 routing_fees,
1909 lightning_context.clone(),
1910 registration.endpoint_url,
1911 registration.keypair,
1912 )
1913 .await;
1914 registration_health
1915 .complete_attempt(
1916 attempt,
1917 succeeded,
1918 fedimint_core::time::now(),
1919 fedimint_core::runtime::Instant::now(),
1920 )
1921 .await;
1922 }
1923 },
1924 );
1925 }
1926 }
1927 }
1928 }
1929
1930 pub async fn select_client(
1933 &self,
1934 federation_id: FederationId,
1935 ) -> std::result::Result<Spanned<fedimint_client::ClientHandleArc>, FederationNotConnected>
1936 {
1937 self.federation_manager
1938 .read()
1939 .await
1940 .client(&federation_id)
1941 .cloned()
1942 .ok_or(FederationNotConnected {
1943 federation_id_prefix: federation_id.to_prefix(),
1944 })
1945 }
1946
1947 async fn load_mnemonic(gateway_db: &Database) -> Option<Mnemonic> {
1948 let secret = Client::load_decodable_client_secret::<Vec<u8>>(gateway_db)
1949 .await
1950 .ok()?;
1951 Mnemonic::from_entropy(&secret).ok()
1952 }
1953
1954 async fn load_clients(&self) -> AdminResult<()> {
1958 if let GatewayState::NotConfigured { .. } = self.get_state().await {
1959 return Ok(());
1960 }
1961
1962 let mut federation_manager = self.federation_manager.write().await;
1963
1964 let configs = {
1965 let mut dbtx = self.gateway_db.begin_transaction_nc().await;
1966 dbtx.load_federation_configs().await
1967 };
1968
1969 if let Some(max_federation_index) = configs.values().map(|cfg| cfg.federation_index).max() {
1970 federation_manager.set_next_index(max_federation_index + 1);
1971 }
1972
1973 let mnemonic = Self::load_mnemonic(&self.gateway_db)
1974 .await
1975 .expect("mnemonic should be set");
1976
1977 for (federation_id, config) in configs {
1978 let federation_index = config.federation_index;
1979 match Box::pin(Spanned::try_new(
1980 info_span!(target: LOG_GATEWAY, "client", federation_id = %federation_id.clone()),
1981 self.client_builder
1982 .build(config, Arc::new(self.clone()), &mnemonic),
1983 ))
1984 .await
1985 {
1986 Ok(client) => {
1987 federation_manager.add_client(federation_index, client);
1988 }
1989 _ => {
1990 warn!(target: LOG_GATEWAY, federation_id = %federation_id, "Failed to load client");
1991 }
1992 }
1993 }
1994
1995 Ok(())
1996 }
1997
1998 fn register_clients_timer(&self) {
2004 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
2006 info!(target: LOG_GATEWAY, "Spawning register task...");
2007 let gateway = self.clone();
2008 let register_task_group = self.task_group.make_subgroup();
2009 self.task_group.spawn_cancellable("register clients", async move {
2010 loop {
2011 let gateway_state = gateway.get_state().await;
2012 if let GatewayState::Running { .. } = &gateway_state {
2013 let mut dbtx = gateway.gateway_db.begin_transaction_nc().await;
2014 let all_federations_configs = dbtx.load_federation_configs().await.into_iter().collect();
2015 gateway.register_federations(&all_federations_configs, ®ister_task_group).await;
2016 } else {
2017 const NOT_RUNNING_RETRY: Duration = Duration::from_secs(10);
2019 warn!(target: LOG_GATEWAY, gateway_state = %gateway_state, retry_interval = ?NOT_RUNNING_RETRY, "Will not register federation yet because gateway still not in Running state");
2020 sleep(NOT_RUNNING_RETRY).await;
2021 continue;
2022 }
2023
2024 sleep(GW_ANNOUNCEMENT_TTL.mul_f32(0.85)).await;
2027 }
2028 });
2029 }
2030 }
2031
2032 async fn check_federation_network(
2035 client: &ClientHandleArc,
2036 network: Network,
2037 ) -> AdminResult<()> {
2038 let federation_id = client.federation_id();
2039 let config = client.config().await;
2040
2041 let lnv1_cfg = config
2042 .modules
2043 .values()
2044 .find(|m| LightningCommonInit::KIND == m.kind);
2045
2046 let lnv2_cfg = config
2047 .modules
2048 .values()
2049 .find(|m| fedimint_lnv2_common::LightningCommonInit::KIND == m.kind);
2050
2051 if lnv1_cfg.is_none() && lnv2_cfg.is_none() {
2053 return Err(AdminGatewayError::ClientCreationError(anyhow!(
2054 "Federation {federation_id} does not have any lightning module (LNv1 or LNv2)"
2055 )));
2056 }
2057
2058 if let Some(cfg) = lnv1_cfg {
2060 let ln_cfg: &LightningClientConfig = cfg.cast()?;
2061
2062 if ln_cfg.network.0 != network {
2063 crit!(
2064 target: LOG_GATEWAY,
2065 federation_id = %federation_id,
2066 network = %network,
2067 "Incorrect LNv1 network for federation",
2068 );
2069 return Err(AdminGatewayError::ClientCreationError(anyhow!(format!(
2070 "Unsupported LNv1 network {}",
2071 ln_cfg.network
2072 ))));
2073 }
2074 }
2075
2076 if let Some(cfg) = lnv2_cfg {
2078 let ln_cfg: &fedimint_lnv2_common::config::LightningClientConfig = cfg.cast()?;
2079
2080 if ln_cfg.network != network {
2081 crit!(
2082 target: LOG_GATEWAY,
2083 federation_id = %federation_id,
2084 network = %network,
2085 "Incorrect LNv2 network for federation",
2086 );
2087 return Err(AdminGatewayError::ClientCreationError(anyhow!(format!(
2088 "Unsupported LNv2 network {}",
2089 ln_cfg.network
2090 ))));
2091 }
2092 }
2093
2094 Ok(())
2095 }
2096
2097 pub async fn get_lightning_context(
2107 &self,
2108 ) -> std::result::Result<LightningContext, LightningRpcError> {
2109 match self.get_state().await {
2110 GatewayState::Running { lightning_context }
2111 | GatewayState::ShuttingDown { lightning_context } => Ok(lightning_context),
2112 _ => Err(LightningRpcError::FailedToConnect),
2113 }
2114 }
2115
2116 async fn await_lightning_context(&self) -> LightningContext {
2137 loop {
2138 match self.get_lightning_context().await {
2139 Ok(lightning_context) => return lightning_context,
2140 Err(err) => {
2141 let state = self.get_state().await;
2142
2143 warn!(
2144 target: LOG_GATEWAY,
2145 err = %err.fmt_compact(),
2146 %state,
2147 retry_interval_secs = LIGHTNING_CONTEXT_RETRY_INTERVAL.as_secs(),
2148 "Not connected to the lightning node, waiting before asking it again",
2149 );
2150
2151 sleep(LIGHTNING_CONTEXT_RETRY_INTERVAL).await;
2152 }
2153 }
2154 }
2155 }
2156
2157 pub async fn unannounce_from_all_federations(&self) {
2160 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
2161 for registration in self.registrations.values() {
2162 self.federation_manager
2163 .read()
2164 .await
2165 .unannounce_from_all_federations(registration.keypair)
2166 .await;
2167 }
2168 }
2169 }
2170
2171 async fn create_lightning_client(
2172 &self,
2173 runtime: Arc<tokio::runtime::Runtime>,
2174 ) -> Box<dyn ILnRpcClient> {
2175 match self.lightning_mode.clone() {
2176 LightningMode::Lnd {
2177 lnd_rpc_addr,
2178 lnd_tls_cert,
2179 lnd_macaroon,
2180 lnd_time_pref,
2181 lnd_payment_timeout_secs,
2182 } => {
2183 let gateway_db = self.gateway_db.clone();
2188 let lnv2_filter: Lnv2HoldInvoiceFilter = Arc::new(move |hash| {
2189 let gateway_db = gateway_db.clone();
2190 Box::pin(async move {
2191 gateway_db
2192 .begin_transaction_nc()
2193 .await
2194 .load_registered_incoming_contract(PaymentImage::Hash(hash))
2195 .await
2196 .is_some()
2197 })
2198 });
2199
2200 Box::new(GatewayLndClient::new(
2201 lnd_rpc_addr,
2202 lnd_tls_cert,
2203 lnd_macaroon,
2204 lnd_time_pref,
2205 lnd_payment_timeout_secs,
2206 None,
2207 lnv2_filter,
2208 ))
2209 }
2210 LightningMode::Ldk {
2211 lightning_port,
2212 alias,
2213 } => {
2214 let mnemonic = Self::load_mnemonic(&self.gateway_db)
2215 .await
2216 .expect("mnemonic should be set");
2217 retry("create LDK Node", fibonacci_max_one_hour(), || async {
2221 ldk::GatewayLdkClient::new(
2222 &self.client_builder.data_dir().join(LDK_NODE_DB_FOLDER),
2223 self.chain_source.clone(),
2224 self.network,
2225 lightning_port,
2226 alias.clone(),
2227 mnemonic.clone(),
2228 runtime.clone(),
2229 )
2230 .map(Box::new)
2231 })
2232 .await
2233 .expect("Could not create LDK Node")
2234 }
2235 }
2236 }
2237}
2238
2239#[async_trait]
2240impl IAdminGateway for Gateway {
2241 type Error = AdminGatewayError;
2242
2243 async fn handle_get_info(&self) -> AdminResult<GatewayInfo> {
2246 let GatewayState::Running { lightning_context } = self.get_state().await else {
2247 return Ok(GatewayInfo {
2248 federations: vec![],
2249 federation_fake_scids: None,
2250 version_hash: fedimint_build_code_version_env!().to_string(),
2251 gateway_state: self.state.read().await.to_string(),
2252 lightning_info: LightningInfo::NotConnected,
2253 lightning_mode: self.lightning_mode.clone(),
2254 registrations: self
2255 .registrations
2256 .iter()
2257 .map(|(k, v)| (k.clone(), (v.endpoint_url.clone(), v.keypair.public_key())))
2258 .collect(),
2259 });
2260 };
2261
2262 let dbtx = self.gateway_db.begin_transaction_nc().await;
2263 let federations = self
2264 .federation_manager
2265 .read()
2266 .await
2267 .federation_info_all_federations(dbtx)
2268 .await;
2269
2270 let channels: BTreeMap<u64, FederationId> = federations
2271 .iter()
2272 .map(|federation_info| {
2273 (
2274 federation_info.config.federation_index,
2275 federation_info.federation_id,
2276 )
2277 })
2278 .collect();
2279
2280 let lightning_info = lightning_context.lnrpc.parsed_node_info().await;
2281
2282 Ok(GatewayInfo {
2283 federations,
2284 federation_fake_scids: Some(channels),
2285 version_hash: fedimint_build_code_version_env!().to_string(),
2286 gateway_state: self.state.read().await.to_string(),
2287 lightning_info,
2288 lightning_mode: self.lightning_mode.clone(),
2289 registrations: self
2290 .registrations
2291 .iter()
2292 .map(|(k, v)| (k.clone(), (v.endpoint_url.clone(), v.keypair.public_key())))
2293 .collect(),
2294 })
2295 }
2296
2297 async fn handle_list_channels_msg(
2300 &self,
2301 ) -> AdminResult<Vec<fedimint_gateway_common::ChannelInfo>> {
2302 let context = self.get_lightning_context().await?;
2303 let response = context.lnrpc.list_channels().await?;
2304 Ok(response.channels)
2305 }
2306
2307 async fn handle_payment_summary_msg(
2310 &self,
2311 PaymentSummaryPayload {
2312 start_millis,
2313 end_millis,
2314 }: PaymentSummaryPayload,
2315 ) -> AdminResult<PaymentSummaryResponse> {
2316 let federation_manager = self.federation_manager.read().await;
2317 let fed_configs = federation_manager.get_all_federation_configs().await;
2318 let federation_ids = fed_configs.keys().collect::<Vec<_>>();
2319 let start = UNIX_EPOCH + Duration::from_millis(start_millis);
2320 let end = UNIX_EPOCH + Duration::from_millis(end_millis);
2321
2322 if start > end {
2323 return Err(AdminGatewayError::Unexpected(anyhow!("Invalid time range")));
2324 }
2325
2326 let mut outgoing = StructuredPaymentEvents::default();
2327 let mut incoming = StructuredPaymentEvents::default();
2328 for fed_id in federation_ids {
2329 let client = federation_manager
2330 .client(fed_id)
2331 .expect("No client available")
2332 .value();
2333 let all_events = &get_events_for_duration(client, start, end).await;
2334
2335 let (mut lnv1_outgoing, mut lnv1_incoming) = compute_lnv1_stats(all_events);
2336 let (mut lnv2_outgoing, mut lnv2_incoming) = compute_lnv2_stats(all_events);
2337 outgoing.combine(&mut lnv1_outgoing);
2338 incoming.combine(&mut lnv1_incoming);
2339 outgoing.combine(&mut lnv2_outgoing);
2340 incoming.combine(&mut lnv2_incoming);
2341 }
2342
2343 Ok(PaymentSummaryResponse {
2344 outgoing: PaymentStats::compute(&outgoing),
2345 incoming: PaymentStats::compute(&incoming),
2346 })
2347 }
2348
2349 async fn handle_leave_federation(
2354 &self,
2355 payload: LeaveFedPayload,
2356 ) -> AdminResult<FederationInfo> {
2357 let mut federation_manager = self.federation_manager.write().await;
2360 let mut dbtx = self.gateway_db.begin_transaction().await;
2361
2362 let federation_info = federation_manager
2363 .leave_federation(
2364 payload.federation_id,
2365 &mut dbtx.to_ref_nc(),
2366 self.registrations.values().collect(),
2367 )
2368 .await?;
2369
2370 dbtx.remove_federation_config(payload.federation_id).await;
2371 dbtx.commit_tx().await;
2372 self.registration_health
2373 .clear_federation(payload.federation_id)
2374 .await;
2375 Ok(federation_info)
2376 }
2377
2378 async fn handle_connect_federation(
2383 &self,
2384 payload: ConnectFedPayload,
2385 ) -> AdminResult<FederationInfo> {
2386 let GatewayState::Running { lightning_context } = self.get_state().await else {
2387 return Err(AdminGatewayError::Lightning(
2388 LightningRpcError::FailedToConnect,
2389 ));
2390 };
2391
2392 let invite_code = InviteCode::from_str(&payload.invite_code).map_err(|e| {
2393 AdminGatewayError::ClientCreationError(anyhow!(format!(
2394 "Invalid federation member string {e:?}"
2395 )))
2396 })?;
2397
2398 let federation_id = invite_code.federation_id();
2399
2400 let mut federation_manager = self.federation_manager.write().await;
2401
2402 if federation_manager.has_federation(federation_id) {
2404 return Err(AdminGatewayError::ClientCreationError(anyhow!(
2405 "Federation has already been registered"
2406 )));
2407 }
2408
2409 let federation_index = federation_manager.pop_next_index()?;
2412
2413 let federation_config = FederationConfig {
2414 invite_code,
2415 federation_index,
2416 lightning_fee: self.default_routing_fees,
2417 transaction_fee: self.default_transaction_fees,
2418 _connector: ConnectorType::Tcp,
2420 };
2421
2422 let routing_fees = RoutingFees::try_from(federation_config.lightning_fee)
2426 .map_err(|err| AdminGatewayError::GatewayConfigurationError(err.to_string()))?;
2427
2428 let mnemonic = Self::load_mnemonic(&self.gateway_db)
2429 .await
2430 .expect("mnemonic should be set");
2431 let recover = payload.recover.unwrap_or(false);
2432 if recover {
2433 self.client_builder
2434 .recover(federation_config.clone(), Arc::new(self.clone()), &mnemonic)
2435 .await?;
2436 }
2437
2438 let client = self
2439 .client_builder
2440 .build(federation_config.clone(), Arc::new(self.clone()), &mnemonic)
2441 .await?;
2442
2443 if recover {
2444 client.wait_for_all_active_state_machines().await?;
2445 }
2446
2447 let federation_info = FederationInfo {
2450 federation_id,
2451 federation_name: federation_manager.federation_name(&client).await,
2452 balance_msat: client.get_balance_for_btc().await.unwrap_or_else(|err| {
2453 warn!(
2454 target: LOG_GATEWAY,
2455 err = %err.fmt_compact_anyhow(),
2456 %federation_id,
2457 "Balance not immediately available after joining/recovering."
2458 );
2459 Amount::default()
2460 }),
2461 config: federation_config.clone(),
2462 last_backup_time: None,
2463 };
2464
2465 Self::check_federation_network(&client, self.network).await?;
2466 if matches!(self.lightning_mode, LightningMode::Lnd { .. })
2467 && let Ok(lnv1) = client.get_first_module::<GatewayClientModule>()
2468 {
2469 for (protocol, registration) in &self.registrations {
2470 let attempt = self
2471 .registration_health
2472 .begin_lnv1_attempt(federation_id, protocol.clone());
2473 let succeeded = lnv1
2474 .try_register_with_federation(
2475 Vec::new(),
2477 GW_ANNOUNCEMENT_TTL,
2478 routing_fees,
2479 lightning_context.clone(),
2480 registration.endpoint_url.clone(),
2481 registration.keypair,
2482 )
2483 .await;
2484 self.registration_health
2485 .complete_attempt(
2486 attempt,
2487 succeeded,
2488 fedimint_core::time::now(),
2489 fedimint_core::runtime::Instant::now(),
2490 )
2491 .await;
2492 }
2493 }
2494
2495 federation_manager.add_client(
2497 federation_index,
2498 Spanned::new(
2499 info_span!(target: LOG_GATEWAY, "client", federation_id=%federation_id.clone()),
2500 async { client },
2501 )
2502 .await,
2503 );
2504
2505 let mut dbtx = self.gateway_db.begin_transaction().await;
2506 dbtx.save_federation_config(&federation_config).await;
2507 dbtx.save_federation_backup_record(federation_id, None)
2508 .await;
2509 dbtx.commit_tx().await;
2510 debug!(
2511 target: LOG_GATEWAY,
2512 federation_id = %federation_id,
2513 federation_index = %federation_index,
2514 "Federation connected"
2515 );
2516
2517 Ok(federation_info)
2518 }
2519
2520 async fn handle_set_fees_msg(
2523 &self,
2524 SetFeesPayload {
2525 federation_id,
2526 lightning_base,
2527 lightning_parts_per_million,
2528 transaction_base,
2529 transaction_parts_per_million,
2530 }: SetFeesPayload,
2531 ) -> AdminResult<()> {
2532 let mut dbtx = self.gateway_db.begin_transaction().await;
2533 let mut fed_configs = if let Some(fed_id) = federation_id {
2534 dbtx.load_federation_configs()
2535 .await
2536 .into_iter()
2537 .filter(|(id, _)| *id == fed_id)
2538 .collect::<BTreeMap<_, _>>()
2539 } else {
2540 dbtx.load_federation_configs().await
2541 };
2542
2543 let federation_manager = self.federation_manager.read().await;
2544
2545 for (federation_id, config) in &mut fed_configs {
2546 let mut lightning_fee = config.lightning_fee;
2547 if let Some(lightning_base) = lightning_base {
2548 lightning_fee.base = lightning_base;
2549 }
2550
2551 if let Some(lightning_ppm) = lightning_parts_per_million {
2552 lightning_fee.parts_per_million = lightning_ppm;
2553 }
2554
2555 let mut transaction_fee = config.transaction_fee;
2556 if let Some(transaction_base) = transaction_base {
2557 transaction_fee.base = transaction_base;
2558 }
2559
2560 if let Some(transaction_ppm) = transaction_parts_per_million {
2561 transaction_fee.parts_per_million = transaction_ppm;
2562 }
2563
2564 federation_manager
2567 .client(federation_id)
2568 .ok_or(FederationNotConnected {
2569 federation_id_prefix: federation_id.to_prefix(),
2570 })?;
2571
2572 let send_fees = lightning_fee.checked_add(transaction_fee).ok_or_else(|| {
2578 AdminGatewayError::GatewayConfigurationError(format!(
2579 "Total Send fees overflowed, they may not exceed {}",
2580 PaymentFee::SEND_FEE_LIMIT
2581 ))
2582 })?;
2583
2584 if !send_fees.is_within(&PaymentFee::SEND_FEE_LIMIT) {
2586 return Err(AdminGatewayError::GatewayConfigurationError(format!(
2587 "Total Send fees exceeded {}",
2588 PaymentFee::SEND_FEE_LIMIT
2589 )));
2590 }
2591
2592 if !transaction_fee.is_within(&PaymentFee::RECEIVE_FEE_LIMIT) {
2594 return Err(AdminGatewayError::GatewayConfigurationError(format!(
2595 "Transaction fees exceeded RECEIVE LIMIT {}",
2596 PaymentFee::RECEIVE_FEE_LIMIT
2597 )));
2598 }
2599
2600 config.lightning_fee = lightning_fee;
2601 config.transaction_fee = transaction_fee;
2602 dbtx.save_federation_config(config).await;
2603 }
2604
2605 dbtx.commit_tx().await;
2606
2607 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
2608 let register_task_group = TaskGroup::new();
2609
2610 self.register_federations(&fed_configs, ®ister_task_group)
2611 .await;
2612 }
2613
2614 Ok(())
2615 }
2616
2617 async fn handle_mnemonic_msg(&self) -> AdminResult<MnemonicResponse> {
2621 let mnemonic = Self::load_mnemonic(&self.gateway_db)
2622 .await
2623 .expect("mnemonic should be set");
2624 let words = mnemonic
2625 .words()
2626 .map(std::string::ToString::to_string)
2627 .collect::<Vec<_>>();
2628 let all_federations = self
2629 .federation_manager
2630 .read()
2631 .await
2632 .get_all_federation_configs()
2633 .await
2634 .keys()
2635 .copied()
2636 .collect::<BTreeSet<_>>();
2637 let legacy_federations = self.client_builder.legacy_federations(all_federations);
2638 let mnemonic_response = MnemonicResponse {
2639 mnemonic: words,
2640 legacy_federations,
2641 };
2642 Ok(mnemonic_response)
2643 }
2644
2645 async fn handle_open_channel_msg(&self, payload: OpenChannelRequest) -> AdminResult<Txid> {
2648 info!(target: LOG_GATEWAY, pubkey = %payload.pubkey, host = %payload.host, amount = %payload.channel_size_sats, "Opening Lightning channel...");
2649 let context = self.get_lightning_context().await?;
2650 let res = context.lnrpc.open_channel(payload).await?;
2651 info!(target: LOG_GATEWAY, txid = %res.funding_txid, "Initiated channel open");
2652 Txid::from_str(&res.funding_txid).map_err(|e| {
2653 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2654 failure_reason: format!("Received invalid channel funding txid string {e}"),
2655 })
2656 })
2657 }
2658
2659 async fn handle_connect_peer_msg(&self, payload: ConnectPeerRequest) -> AdminResult<()> {
2662 info!(
2663 target: LOG_GATEWAY,
2664 pubkey = %payload.node_address.pubkey,
2665 host = %payload.node_address.host_with_port(),
2666 "Connecting to Lightning peer..."
2667 );
2668 let context = self.get_lightning_context().await?;
2669 context.lnrpc.connect_peer(payload).await?;
2670 info!(target: LOG_GATEWAY, "Connected to Lightning peer");
2671 Ok(())
2672 }
2673
2674 async fn handle_close_channels_with_peer_msg(
2677 &self,
2678 payload: CloseChannelsWithPeerRequest,
2679 ) -> AdminResult<CloseChannelsWithPeerResponse> {
2680 info!(target: LOG_GATEWAY, close_channel_request = %payload, "Closing lightning channel...");
2681 let context = self.get_lightning_context().await?;
2682 let response = context
2683 .lnrpc
2684 .close_channels_with_peer(payload.clone())
2685 .await?;
2686 info!(target: LOG_GATEWAY, close_channel_request = %payload, "Initiated channel closure");
2687 Ok(response)
2688 }
2689
2690 async fn handle_set_channel_fees_msg(&self, payload: SetChannelFeesRequest) -> AdminResult<()> {
2693 info!(
2694 target: LOG_GATEWAY,
2695 funding_outpoint = %payload.funding_outpoint,
2696 base_fee_msat = payload.base_fee_msat,
2697 parts_per_million = payload.parts_per_million,
2698 "Updating channel fees..."
2699 );
2700 let context = self.get_lightning_context().await?;
2701 context.lnrpc.set_channel_fees(payload).await?;
2702 Ok(())
2703 }
2704
2705 async fn handle_get_balances_msg(&self) -> AdminResult<GatewayBalances> {
2708 let dbtx = self.gateway_db.begin_transaction_nc().await;
2709 let federation_infos = self
2710 .federation_manager
2711 .read()
2712 .await
2713 .federation_info_all_federations(dbtx)
2714 .await;
2715
2716 let ecash_balances: Vec<FederationBalanceInfo> = federation_infos
2717 .iter()
2718 .map(|federation_info| FederationBalanceInfo {
2719 federation_id: federation_info.federation_id,
2720 ecash_balance_msats: Amount {
2721 msats: federation_info.balance_msat.msats,
2722 },
2723 })
2724 .collect();
2725
2726 let context = self.get_lightning_context().await?;
2727 let lightning_node_balances = context.lnrpc.get_balances().await?;
2728
2729 Ok(GatewayBalances {
2730 onchain_balance_sats: lightning_node_balances.onchain_balance_sats,
2731 lightning_balance_msats: lightning_node_balances.lightning_balance_msats,
2732 ecash_balances,
2733 inbound_lightning_liquidity_msats: lightning_node_balances
2734 .inbound_lightning_liquidity_msats,
2735 })
2736 }
2737
2738 async fn handle_send_onchain_msg(&self, payload: SendOnchainRequest) -> AdminResult<Txid> {
2740 let context = self.get_lightning_context().await?;
2741 let response = context.lnrpc.send_onchain(payload.clone()).await?;
2742 let txid =
2743 Txid::from_str(&response.txid).map_err(|e| AdminGatewayError::WithdrawError {
2744 failure_reason: format!("Failed to parse withdrawal TXID: {e}"),
2745 })?;
2746 info!(onchain_request = %payload, txid = %txid, "Sent onchain transaction");
2747 Ok(txid)
2748 }
2749
2750 async fn handle_get_ln_onchain_address_msg(&self) -> AdminResult<Address> {
2752 let context = self.get_lightning_context().await?;
2753 let response = context.lnrpc.get_ln_onchain_address().await?;
2754
2755 let address = Address::from_str(&response.address).map_err(|e| {
2756 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2757 failure_reason: e.to_string(),
2758 })
2759 })?;
2760
2761 address.require_network(self.network).map_err(|e| {
2762 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2763 failure_reason: e.to_string(),
2764 })
2765 })
2766 }
2767
2768 async fn handle_deposit_address_msg(
2769 &self,
2770 payload: DepositAddressPayload,
2771 ) -> AdminResult<Address> {
2772 self.handle_address_msg(payload).await
2773 }
2774
2775 async fn handle_receive_ecash_msg(
2776 &self,
2777 payload: ReceiveEcashPayload,
2778 ) -> AdminResult<ReceiveEcashResponse> {
2779 Self::handle_receive_ecash_msg(self, payload)
2780 .await
2781 .map_err(|e| AdminGatewayError::Unexpected(anyhow::anyhow!("{e}")))
2782 }
2783
2784 async fn handle_create_invoice_for_operator_msg(
2787 &self,
2788 payload: CreateInvoiceForOperatorPayload,
2789 ) -> AdminResult<Bolt11Invoice> {
2790 let GatewayState::Running { lightning_context } = self.get_state().await else {
2791 return Err(AdminGatewayError::Lightning(
2792 LightningRpcError::FailedToConnect,
2793 ));
2794 };
2795
2796 Bolt11Invoice::from_str(
2797 &lightning_context
2798 .lnrpc
2799 .create_invoice(CreateInvoiceRequest {
2800 payment_hash: None, amount_msat: payload.amount_msats,
2803 expiry_secs: payload.expiry_secs.unwrap_or(3600),
2804 description: payload.description.map(InvoiceDescription::Direct),
2805 })
2806 .await?
2807 .invoice,
2808 )
2809 .map_err(|e| {
2810 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2811 failure_reason: e.to_string(),
2812 })
2813 })
2814 }
2815
2816 async fn handle_pay_invoice_for_operator_msg(
2819 &self,
2820 payload: PayInvoiceForOperatorPayload,
2821 ) -> AdminResult<Preimage> {
2822 const BASE_FEE: u64 = 50;
2824 const FEE_DENOMINATOR: u64 = 100;
2825 const MAX_DELAY: u64 = 1008;
2826
2827 let GatewayState::Running { lightning_context } = self.get_state().await else {
2828 return Err(AdminGatewayError::Lightning(
2829 LightningRpcError::FailedToConnect,
2830 ));
2831 };
2832
2833 let max_fee = BASE_FEE
2834 + payload
2835 .invoice
2836 .amount_milli_satoshis()
2837 .context("Invoice is missing amount")?
2838 .saturating_div(FEE_DENOMINATOR);
2839
2840 let res = lightning_context
2841 .lnrpc
2842 .pay(payload.invoice, MAX_DELAY, Amount::from_msats(max_fee))
2843 .await?;
2844 Ok(res.preimage)
2845 }
2846
2847 async fn handle_list_transactions_msg(
2849 &self,
2850 payload: ListTransactionsPayload,
2851 ) -> AdminResult<ListTransactionsResponse> {
2852 let lightning_context = self.get_lightning_context().await?;
2853 let response = lightning_context
2854 .lnrpc
2855 .list_transactions(payload.start_secs, payload.end_secs)
2856 .await?;
2857 Ok(response)
2858 }
2859
2860 async fn handle_spend_ecash_msg(
2862 &self,
2863 payload: SpendEcashPayload,
2864 ) -> AdminResult<SpendEcashResponse> {
2865 let client = self
2866 .select_client(payload.federation_id)
2867 .await?
2868 .into_value();
2869
2870 if let Ok(mint_module) = client.get_first_module::<MintClientModule>() {
2871 let notes = mint_module.send_oob_notes(payload.amount, ()).await?;
2872 debug!(target: LOG_GATEWAY, ?notes, "Spend ecash notes");
2873 Ok(SpendEcashResponse {
2874 notes: notes.to_string(),
2875 })
2876 } else if let Ok(mint_module) = client.get_first_module::<MintV2ClientModule>() {
2877 let (_, ecash) = mint_module
2878 .send(payload.amount, serde_json::Value::Null, true)
2879 .await
2880 .map_err(|e| AdminGatewayError::Unexpected(e.into()))?;
2881
2882 Ok(SpendEcashResponse {
2883 notes: base32::encode_prefixed(FEDIMINT_PREFIX, &ecash),
2884 })
2885 } else {
2886 Err(AdminGatewayError::Unexpected(anyhow::anyhow!(
2887 "No mint module available"
2888 )))
2889 }
2890 }
2891
2892 async fn handle_shutdown_msg(&self, task_group: TaskGroup) -> AdminResult<()> {
2895 let was_running = {
2899 let mut state_guard = self.state.write().await;
2900 if let GatewayState::Running { lightning_context } = state_guard.clone() {
2901 *state_guard = GatewayState::ShuttingDown { lightning_context };
2902 true
2903 } else {
2904 false
2905 }
2906 };
2907
2908 if was_running {
2916 self.federation_manager
2917 .read()
2918 .await
2919 .wait_for_incoming_payments()
2920 .await?;
2921 }
2922
2923 let tg = task_group.clone();
2924 tg.spawn("Kill Gateway", |_task_handle| async {
2925 if let Err(err) = task_group.shutdown_join_all(Duration::from_mins(3)).await {
2926 warn!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Error shutting down gateway");
2927 }
2928 });
2929 Ok(())
2930 }
2931
2932 fn get_task_group(&self) -> TaskGroup {
2933 self.task_group.clone()
2934 }
2935
2936 async fn handle_withdraw_msg(&self, payload: WithdrawPayload) -> AdminResult<WithdrawResponse> {
2939 let WithdrawPayload {
2940 amount,
2941 address,
2942 federation_id,
2943 quoted_fees,
2944 } = payload;
2945
2946 let address_network = get_network_for_address(&address);
2947 let gateway_network = self.network;
2948 let Ok(address) = address.require_network(gateway_network) else {
2949 return Err(AdminGatewayError::WithdrawError {
2950 failure_reason: format!(
2951 "Gateway is running on network {gateway_network}, but provided withdraw address is for network {address_network}"
2952 ),
2953 });
2954 };
2955
2956 let client = self.select_client(federation_id).await?;
2957
2958 if let Ok(wallet_module) = client
2959 .value()
2960 .get_first_module::<fedimint_walletv2_client::WalletClientModule>()
2961 {
2962 return withdraw_v2(client.value(), &wallet_module, &address, amount).await;
2963 }
2964
2965 let wallet_module = client.value().get_first_module::<WalletClientModule>()?;
2966
2967 let (withdraw_amount, fees) = match quoted_fees {
2970 Some(fees) => {
2972 let amt = match amount {
2973 BitcoinAmountOrAll::Amount(a) => a,
2974 BitcoinAmountOrAll::All => {
2975 return Err(AdminGatewayError::WithdrawError {
2977 failure_reason:
2978 "Cannot use 'all' with quoted fees - amount must be resolved first"
2979 .to_string(),
2980 });
2981 }
2982 };
2983 (amt, fees)
2984 }
2985 None => match amount {
2987 BitcoinAmountOrAll::All => {
2992 let balance = client.value().get_balance_for_btc().await.map_err(|err| {
2993 AdminGatewayError::Unexpected(anyhow!(
2994 "Balance not available: {}",
2995 err.fmt_compact_anyhow()
2996 ))
2997 })?;
2998
2999 wallet_module
3000 .max_withdrawable_amount(&address, balance)
3001 .await
3002 .map_err(|err| AdminGatewayError::WithdrawError {
3003 failure_reason: format!(
3004 "Insufficient funds. Balance: {balance}: {}",
3005 err.fmt_compact_anyhow()
3006 ),
3007 })?
3008 }
3009 BitcoinAmountOrAll::Amount(amount) => (
3010 amount,
3011 wallet_module.get_withdraw_fees(&address, amount).await?,
3012 ),
3013 },
3014 };
3015
3016 let operation_id = wallet_module
3017 .withdraw(&address, withdraw_amount, fees, ())
3018 .await?;
3019 let mut updates = wallet_module
3020 .subscribe_withdraw_updates(operation_id)
3021 .await?
3022 .into_stream();
3023
3024 while let Some(update) = updates.next().await {
3025 match update {
3026 WithdrawState::Succeeded(txid) => {
3027 info!(target: LOG_GATEWAY, amount = %withdraw_amount, address = %address, "Sent funds");
3028 return Ok(WithdrawResponse { txid, fees });
3029 }
3030 WithdrawState::Failed(e) => {
3031 return Err(AdminGatewayError::WithdrawError { failure_reason: e });
3032 }
3033 WithdrawState::Created => {}
3034 }
3035 }
3036
3037 Err(AdminGatewayError::WithdrawError {
3038 failure_reason: "Ran out of state updates while withdrawing".to_string(),
3039 })
3040 }
3041
3042 async fn handle_withdraw_preview_msg(
3045 &self,
3046 payload: WithdrawPreviewPayload,
3047 ) -> AdminResult<WithdrawPreviewResponse> {
3048 let gateway_network = self.network;
3049 let address_checked = payload
3050 .address
3051 .clone()
3052 .require_network(gateway_network)
3053 .map_err(|_| AdminGatewayError::WithdrawError {
3054 failure_reason: "Address network mismatch".to_string(),
3055 })?;
3056
3057 let client = self.select_client(payload.federation_id).await?;
3058
3059 let WithdrawDetails {
3060 amount,
3061 mint_fees,
3062 peg_out_fees,
3063 } = match payload.amount {
3064 BitcoinAmountOrAll::All => {
3065 calculate_max_withdrawable(client.value(), &address_checked).await?
3066 }
3067 BitcoinAmountOrAll::Amount(btc_amount) => {
3068 if let Ok(wallet_module) = client.value().get_first_module::<WalletClientModule>() {
3069 WithdrawDetails {
3070 amount: btc_amount.into(),
3071 mint_fees: None,
3072 peg_out_fees: wallet_module
3073 .get_withdraw_fees(&address_checked, btc_amount)
3074 .await?,
3075 }
3076 } else if let Ok(wallet_module) = client
3077 .value()
3078 .get_first_module::<fedimint_walletv2_client::WalletClientModule>(
3079 ) {
3080 let fee = wallet_module.send_fee().await.map_err(|e| {
3081 AdminGatewayError::WithdrawError {
3082 failure_reason: e.to_string(),
3083 }
3084 })?;
3085 WithdrawDetails {
3086 amount: btc_amount.into(),
3087 mint_fees: None,
3088 peg_out_fees: PegOutFees::from_amount(fee),
3089 }
3090 } else {
3091 return Err(AdminGatewayError::Unexpected(anyhow!(
3092 "No wallet module found"
3093 )));
3094 }
3095 }
3096 };
3097
3098 let total_cost = amount
3099 .checked_add(peg_out_fees.amount().into())
3100 .and_then(|a| a.checked_add(mint_fees.unwrap_or(Amount::ZERO)))
3101 .ok_or_else(|| AdminGatewayError::Unexpected(anyhow!("Total cost overflow")))?;
3102
3103 Ok(WithdrawPreviewResponse {
3104 withdraw_amount: amount,
3105 address: payload.address.assume_checked().to_string(),
3106 peg_out_fees,
3107 total_cost,
3108 mint_fees,
3109 })
3110 }
3111
3112 async fn handle_payment_log_msg(
3124 &self,
3125 PaymentLogPayload {
3126 end_position,
3127 pagination_size,
3128 federation_id,
3129 event_kinds,
3130 }: PaymentLogPayload,
3131 ) -> AdminResult<PaymentLogResponse> {
3132 const BATCH_SIZE: u64 = 10_000;
3133 let federation_manager = self.federation_manager.read().await;
3134 let client = federation_manager
3135 .client(&federation_id)
3136 .ok_or(FederationNotConnected {
3137 federation_id_prefix: federation_id.to_prefix(),
3138 })?
3139 .value();
3140
3141 let event_kinds = if event_kinds.is_empty() {
3145 ALL_GATEWAY_EVENTS.to_vec()
3146 } else {
3147 event_kinds
3148 };
3149
3150 let end_position = if let Some(position) = end_position {
3151 position
3152 } else {
3153 let mut dbtx = client.db().begin_transaction_nc().await;
3154 dbtx.get_next_event_log_id().await
3155 };
3156
3157 let mut start_position = end_position.saturating_sub(BATCH_SIZE);
3158
3159 let mut payment_log = Vec::new();
3160
3161 while payment_log.len() < pagination_size {
3162 let batch = client.get_event_log(Some(start_position), BATCH_SIZE).await;
3163 let mut filtered_batch = batch
3164 .into_iter()
3165 .filter(|e| e.id() <= end_position && event_kinds.contains(&e.as_raw().kind))
3166 .collect::<Vec<_>>();
3167 filtered_batch.reverse();
3168 payment_log.extend(filtered_batch);
3169
3170 start_position = start_position.saturating_sub(BATCH_SIZE);
3172
3173 if start_position == EventLogId::LOG_START {
3174 break;
3175 }
3176 }
3177
3178 payment_log.truncate(pagination_size);
3180
3181 Ok(PaymentLogResponse(payment_log))
3182 }
3183
3184 async fn handle_set_mnemonic_msg(&self, payload: SetMnemonicPayload) -> AdminResult<()> {
3187 let mut state_guard = self.state.write().await;
3192
3193 let GatewayState::NotConfigured { mnemonic_sender } = state_guard.clone() else {
3195 return Err(AdminGatewayError::MnemonicError(anyhow!(
3196 "Gateway is not is NotConfigured state"
3197 )));
3198 };
3199
3200 let mnemonic = if let Some(words) = payload.words {
3201 info!(target: LOG_GATEWAY, "Using user provided mnemonic");
3202 Mnemonic::parse_in_normalized(Language::English, words.as_str()).map_err(|e| {
3203 AdminGatewayError::MnemonicError(anyhow!(format!(
3204 "Seed phrase provided in environment was invalid {e:?}"
3205 )))
3206 })?
3207 } else {
3208 debug!(target: LOG_GATEWAY, "Generating mnemonic and writing entropy to client storage");
3209 Bip39RootSecretStrategy::<12>::random(&mut OsRng)
3210 };
3211
3212 Client::store_encodable_client_secret(&self.gateway_db, mnemonic.to_entropy())
3213 .await
3214 .map_err(AdminGatewayError::MnemonicError)?;
3215
3216 *state_guard = GatewayState::Disconnected;
3217 drop(state_guard);
3218
3219 let _ = mnemonic_sender.send(());
3221
3222 Ok(())
3223 }
3224
3225 async fn handle_create_offer_for_operator_msg(
3227 &self,
3228 payload: CreateOfferPayload,
3229 ) -> AdminResult<CreateOfferResponse> {
3230 let lightning_context = self.get_lightning_context().await?;
3231 let offer = lightning_context.lnrpc.create_offer(
3232 payload.amount,
3233 payload.description,
3234 payload.expiry_secs,
3235 payload.quantity,
3236 )?;
3237 Ok(CreateOfferResponse { offer })
3238 }
3239
3240 async fn handle_pay_offer_for_operator_msg(
3242 &self,
3243 payload: PayOfferPayload,
3244 ) -> AdminResult<PayOfferResponse> {
3245 let lightning_context = self.get_lightning_context().await?;
3246 let preimage = lightning_context
3247 .lnrpc
3248 .pay_offer(
3249 payload.offer,
3250 payload.quantity,
3251 payload.amount,
3252 payload.payer_note,
3253 )
3254 .await?;
3255 Ok(PayOfferResponse {
3256 preimage: preimage.to_string(),
3257 })
3258 }
3259
3260 async fn handle_export_invite_codes(
3263 &self,
3264 ) -> BTreeMap<FederationId, BTreeMap<PeerId, (String, InviteCode)>> {
3265 let fed_manager = self.federation_manager.read().await;
3266 fed_manager.all_invite_codes().await
3267 }
3268
3269 async fn handle_get_note_summary_msg(
3272 &self,
3273 federation_id: &FederationId,
3274 ) -> AdminResult<TieredCounts> {
3275 let fed_manager = self.federation_manager.read().await;
3276 fed_manager.get_note_summary(federation_id).await
3277 }
3278
3279 fn get_password_hash(&self) -> String {
3280 self.bcrypt_password_hash.clone()
3281 }
3282
3283 fn gatewayd_version(&self) -> String {
3284 let gatewayd_version = env!("CARGO_PKG_VERSION");
3285 gatewayd_version.to_string()
3286 }
3287
3288 async fn get_chain_source(&self) -> (ChainSource, Network) {
3289 (self.chain_source.clone(), self.network)
3290 }
3291
3292 fn lightning_mode(&self) -> LightningMode {
3293 self.lightning_mode.clone()
3294 }
3295
3296 async fn is_configured(&self) -> bool {
3297 !matches!(self.get_state().await, GatewayState::NotConfigured { .. })
3298 }
3299}
3300
3301impl Gateway {
3303 async fn public_key_v2(&self, federation_id: &FederationId) -> Option<PublicKey> {
3307 self.federation_manager
3308 .read()
3309 .await
3310 .client(federation_id)
3311 .and_then(|client| {
3312 client
3315 .value()
3316 .get_first_module::<GatewayClientModuleV2>()
3317 .ok()
3318 .map(|module| module.keypair.public_key())
3319 })
3320 }
3321
3322 pub async fn routing_info_v2(
3325 &self,
3326 federation_id: &FederationId,
3327 ) -> Result<Option<RoutingInfo>> {
3328 let context = self.get_lightning_context().await?;
3329
3330 let mut dbtx = self.gateway_db.begin_transaction_nc().await;
3331 let fed_config = dbtx.load_federation_config(*federation_id).await.ok_or(
3332 PublicGatewayError::FederationNotConnected(FederationNotConnected {
3333 federation_id_prefix: federation_id.to_prefix(),
3334 }),
3335 )?;
3336
3337 let lightning_fee = fed_config.lightning_fee;
3338 let transaction_fee = fed_config.transaction_fee;
3339
3340 let send_fee_default = lightning_fee.checked_add(transaction_fee).ok_or_else(|| {
3343 PublicGatewayError::Unexpected(anyhow!(
3344 "The configured fees of federation {federation_id} cannot be added"
3345 ))
3346 })?;
3347
3348 Ok(self
3349 .public_key_v2(federation_id)
3350 .await
3351 .map(|module_public_key| RoutingInfo {
3352 lightning_public_key: context.lightning_public_key,
3353 lightning_alias: Some(context.lightning_alias.clone()),
3354 module_public_key,
3355 send_fee_default,
3356 send_fee_minimum: transaction_fee,
3360 expiration_delta_default: 1440,
3361 expiration_delta_minimum: EXPIRATION_DELTA_MINIMUM_V2,
3362 receive_fee: transaction_fee,
3365 }))
3366 }
3367
3368 pub async fn send_payment_v2(
3371 &self,
3372 payload: SendPaymentPayload,
3373 ) -> Result<std::result::Result<[u8; 32], Signature>> {
3374 let client = self.select_client(payload.federation_id).await?;
3375 let module = client
3378 .value()
3379 .get_first_module::<GatewayClientModuleV2>()
3380 .map_err(|err| PublicGatewayError::LNv2(LNv2Error::OutgoingPayment(err)))?;
3381
3382 module
3383 .send_payment(payload)
3384 .await
3385 .map_err(LNv2Error::OutgoingPayment)
3386 .map_err(PublicGatewayError::LNv2)
3387 }
3388
3389 async fn create_bolt11_invoice_v2(
3394 &self,
3395 payload: CreateBolt11InvoicePayload,
3396 ) -> Result<Bolt11Invoice> {
3397 if !self.invoice_rate_limiter.try_acquire() {
3401 return Err(PublicGatewayError::RateLimited);
3402 }
3403
3404 if !payload.contract.verify() {
3405 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3406 "The contract is invalid".to_string(),
3407 )));
3408 }
3409
3410 let payment_info = self.routing_info_v2(&payload.federation_id).await?.ok_or(
3411 LNv2Error::IncomingPayment(format!(
3412 "Federation {} does not exist",
3413 payload.federation_id
3414 )),
3415 )?;
3416
3417 if payload.contract.commitment.refund_pk != payment_info.module_public_key {
3418 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3419 "The incoming contract is keyed to another gateway".to_string(),
3420 )));
3421 }
3422
3423 let contract_amount = payment_info.receive_fee.subtract_from(payload.amount.msats);
3424
3425 if contract_amount == Amount::ZERO {
3426 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3427 "Zero amount incoming contracts are not supported".to_string(),
3428 )));
3429 }
3430
3431 if contract_amount != payload.contract.commitment.amount {
3432 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3433 "The contract amount does not pay the correct amount of fees".to_string(),
3434 )));
3435 }
3436
3437 if payload.contract.commitment.expiration_or_fee <= duration_since_epoch().as_secs() {
3438 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3439 "The contract has already expired".to_string(),
3440 )));
3441 }
3442
3443 if payload.expiry_secs > MAX_INVOICE_EXPIRY_SECS {
3444 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3445 "The invoice expiry exceeds the maximum of one day".to_string(),
3446 )));
3447 }
3448
3449 let payment_hash = match payload.contract.commitment.payment_image {
3450 PaymentImage::Hash(payment_hash) => payment_hash,
3451 PaymentImage::Point(..) => {
3452 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3453 "PaymentImage is not a payment hash".to_string(),
3454 )));
3455 }
3456 };
3457
3458 let mut dbtx = self.gateway_db.begin_transaction().await;
3463
3464 let invoice_expires_at_secs = duration_since_epoch()
3465 .as_secs()
3466 .saturating_add(u64::from(payload.expiry_secs));
3467
3468 if dbtx
3469 .save_registered_incoming_contract(
3470 payload.federation_id,
3471 payload.amount,
3472 invoice_expires_at_secs,
3473 payload.contract,
3474 )
3475 .await
3476 .is_some()
3477 {
3478 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3479 "PaymentHash is already registered".to_string(),
3480 )));
3481 }
3482
3483 dbtx.commit_tx_result().await.map_err(|_| {
3484 PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3485 "Payment hash is already registered".to_string(),
3486 ))
3487 })?;
3488
3489 match self
3490 .create_invoice_via_lnrpc_v2(
3491 payment_hash,
3492 payload.amount,
3493 payload.description.clone(),
3494 payload.expiry_secs,
3495 )
3496 .await
3497 {
3498 Ok(invoice) => Ok(invoice),
3499 Err(err) => {
3500 let mut dbtx = self.gateway_db.begin_transaction().await;
3503 dbtx.delete_registered_incoming_contract(PaymentImage::Hash(payment_hash))
3504 .await;
3505 if let Err(db_err) = dbtx.commit_tx_result().await {
3506 warn!(
3507 target: LOG_GATEWAY,
3508 err = %db_err.fmt_compact(),
3509 %payment_hash,
3510 "Failed to release incoming contract reservation after Lightning error"
3511 );
3512 }
3513
3514 Err(err.into())
3515 }
3516 }
3517 }
3518
3519 pub async fn create_invoice_via_lnrpc_v2(
3522 &self,
3523 payment_hash: sha256::Hash,
3524 amount: Amount,
3525 description: Bolt11InvoiceDescription,
3526 expiry_time: u32,
3527 ) -> std::result::Result<Bolt11Invoice, LightningRpcError> {
3528 let lnrpc = self.get_lightning_context().await?.lnrpc;
3529
3530 let response = match description {
3531 Bolt11InvoiceDescription::Direct(description) => {
3532 lnrpc
3533 .create_invoice(CreateInvoiceRequest {
3534 payment_hash: Some(payment_hash),
3535 amount_msat: amount.msats,
3536 expiry_secs: expiry_time,
3537 description: Some(InvoiceDescription::Direct(description)),
3538 })
3539 .await?
3540 }
3541 Bolt11InvoiceDescription::Hash(hash) => {
3542 lnrpc
3543 .create_invoice(CreateInvoiceRequest {
3544 payment_hash: Some(payment_hash),
3545 amount_msat: amount.msats,
3546 expiry_secs: expiry_time,
3547 description: Some(InvoiceDescription::Hash(hash)),
3548 })
3549 .await?
3550 }
3551 };
3552
3553 Bolt11Invoice::from_str(&response.invoice).map_err(|e| {
3554 LightningRpcError::FailedToGetInvoice {
3555 failure_reason: e.to_string(),
3556 }
3557 })
3558 }
3559
3560 pub async fn verify_bolt11_preimage_v2(
3561 &self,
3562 payment_hash: sha256::Hash,
3563 wait: bool,
3564 ) -> std::result::Result<VerifyResponse, String> {
3565 let registered_contract = self
3566 .gateway_db
3567 .begin_transaction_nc()
3568 .await
3569 .load_registered_incoming_contract(PaymentImage::Hash(payment_hash))
3570 .await
3571 .ok_or("Unknown payment hash".to_string())?;
3572
3573 let client = self
3574 .select_client(registered_contract.federation_id)
3575 .await
3576 .map_err(|_| "Not connected to federation".to_string())?
3577 .into_value();
3578
3579 let operation_id = OperationId::from_encodable(®istered_contract.contract);
3580
3581 if !(wait || client.operation_exists(operation_id).await) {
3582 return Ok(VerifyResponse {
3583 settled: false,
3584 preimage: None,
3585 });
3586 }
3587
3588 let module = client
3589 .get_first_module::<GatewayClientModuleV2>()
3590 .expect("Must have client module");
3591
3592 let Ok(state) = timeout(VERIFY_WAIT_TIMEOUT, module.await_receive(operation_id)).await
3593 else {
3594 return Ok(VerifyResponse {
3595 settled: false,
3596 preimage: None,
3597 });
3598 };
3599
3600 let preimage = match state {
3601 FinalReceiveState::Success(preimage) => Ok(preimage),
3602 FinalReceiveState::Failure => Err("Payment has failed".to_string()),
3603 FinalReceiveState::Refunded => Err("Payment has been refunded".to_string()),
3604 FinalReceiveState::Rejected => Err("Payment has been rejected".to_string()),
3605 }?;
3606
3607 Ok(VerifyResponse {
3608 settled: true,
3609 preimage: Some(preimage),
3610 })
3611 }
3612
3613 pub async fn get_registered_incoming_contract_and_client_v2(
3617 &self,
3618 payment_image: PaymentImage,
3619 amount_msats: u64,
3620 ) -> Result<(IncomingContract, ClientHandleArc)> {
3621 let registered_incoming_contract = self
3622 .gateway_db
3623 .begin_transaction_nc()
3624 .await
3625 .load_registered_incoming_contract(payment_image)
3626 .await
3627 .ok_or(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3628 "No corresponding decryption contract available".to_string(),
3629 )))?;
3630
3631 if registered_incoming_contract.incoming_amount_msats != amount_msats {
3632 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3633 "The available decryption contract's amount is not equal to the requested amount"
3634 .to_string(),
3635 )));
3636 }
3637
3638 let client = self
3639 .select_client(registered_incoming_contract.federation_id)
3640 .await?
3641 .into_value();
3642
3643 Ok((registered_incoming_contract.contract, client))
3644 }
3645}
3646
3647#[async_trait]
3648impl IGatewayClientV2 for Gateway {
3649 async fn complete_htlc(
3650 &self,
3651 htlc_response: InterceptPaymentResponse,
3652 ) -> std::result::Result<(), LightningRpcError> {
3653 loop {
3654 let lightning_context = self.await_lightning_context().await;
3655
3656 match lightning_context
3657 .lnrpc
3658 .complete_htlc(htlc_response.clone())
3659 .await
3660 {
3661 Ok(..) => return Ok(()),
3662 Err(err @ LightningRpcError::HtlcCompletionRejected { .. }) => {
3663 warn!(
3664 target: LOG_GATEWAY,
3665 err = %err.fmt_compact(),
3666 "Lightning cannot reach the requested terminal HTLC outcome",
3667 );
3668 return Err(err);
3669 }
3670 Err(err) => {
3671 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failure trying to complete payment");
3672 }
3673 }
3674
3675 sleep(LIGHTNING_CONTEXT_RETRY_INTERVAL).await;
3676 }
3677 }
3678
3679 async fn is_direct_swap(
3680 &self,
3681 invoice: &Bolt11Invoice,
3682 ) -> anyhow::Result<Option<(IncomingContract, ClientHandleArc)>> {
3683 let lightning_context = self.await_lightning_context().await;
3688 if lightning_context.lightning_public_key == invoice.get_payee_pub_key() {
3689 let (contract, client) = self
3690 .get_registered_incoming_contract_and_client_v2(
3691 PaymentImage::Hash(*invoice.payment_hash()),
3692 invoice
3693 .amount_milli_satoshis()
3694 .expect("The amount invoice has been previously checked"),
3695 )
3696 .await?;
3697 Ok(Some((contract, client)))
3698 } else {
3699 Ok(None)
3700 }
3701 }
3702
3703 async fn pay(
3704 &self,
3705 invoice: Bolt11Invoice,
3706 max_delay: u64,
3707 max_fee: Amount,
3708 ) -> std::result::Result<[u8; 32], LightningRpcError> {
3709 let lightning_context = self.await_lightning_context().await;
3712 lightning_context
3713 .lnrpc
3714 .pay(invoice, max_delay, max_fee)
3715 .await
3716 .map(|response| response.preimage.0)
3717 }
3718
3719 async fn min_contract_amount(
3720 &self,
3721 federation_id: &FederationId,
3722 amount: u64,
3723 ) -> anyhow::Result<Amount> {
3724 Ok(self
3725 .routing_info_v2(federation_id)
3726 .await?
3727 .ok_or(anyhow!("Routing Info not available"))?
3728 .send_fee_minimum
3729 .add_to(amount))
3730 }
3731
3732 async fn is_lnv1_invoice(&self, invoice: &Bolt11Invoice) -> Option<Spanned<ClientHandleArc>> {
3733 let rhints = invoice.route_hints();
3734 let hop = rhints.first().and_then(|rh| rh.0.last())?;
3735
3736 let lightning_context = self.await_lightning_context().await;
3739 if hop.src_node_id != lightning_context.lightning_public_key {
3740 return None;
3741 }
3742
3743 self.federation_manager
3744 .read()
3745 .await
3746 .get_client_for_index(hop.short_channel_id)
3747 }
3748
3749 async fn relay_lnv1_swap(
3750 &self,
3751 client: &ClientHandleArc,
3752 invoice: &Bolt11Invoice,
3753 ) -> anyhow::Result<FinalReceiveState> {
3754 let swap_params = SwapParameters {
3755 payment_hash: *invoice.payment_hash(),
3756 amount_msat: Amount::from_msats(
3757 invoice
3758 .amount_milli_satoshis()
3759 .ok_or(anyhow!("Amountless invoice not supported"))?,
3760 ),
3761 };
3762 let lnv1 = client
3763 .get_first_module::<GatewayClientModule>()
3764 .expect("No LNv1 module");
3765 let operation_id = lnv1.gateway_handle_direct_swap(swap_params).await?;
3766 let mut stream = lnv1
3767 .gateway_subscribe_ln_receive(operation_id)
3768 .await?
3769 .into_stream();
3770 let mut final_state = FinalReceiveState::Failure;
3771 while let Some(update) = stream.next().await {
3772 match update {
3773 GatewayExtReceiveStates::Funding => {}
3774 GatewayExtReceiveStates::FundingFailed { error: _ } => {
3775 final_state = FinalReceiveState::Rejected;
3776 }
3777 GatewayExtReceiveStates::Preimage(preimage) => {
3778 final_state = FinalReceiveState::Success(preimage.0);
3779 }
3780 GatewayExtReceiveStates::RefundError {
3781 error_message: _,
3782 error: _,
3783 } => {
3784 final_state = FinalReceiveState::Failure;
3785 }
3786 GatewayExtReceiveStates::RefundSuccess {
3787 out_points: _,
3788 error: _,
3789 } => {
3790 final_state = FinalReceiveState::Refunded;
3791 }
3792 }
3793 }
3794
3795 Ok(final_state)
3796 }
3797
3798 async fn claim_payment_image(
3799 &self,
3800 payment_image: &PaymentImage,
3801 operation_id: OperationId,
3802 ) -> bool {
3803 self.gateway_db
3807 .autocommit(
3808 |dbtx, _| {
3809 let payment_image = payment_image.clone();
3810 Box::pin(async move {
3811 let claimer = dbtx
3812 .claim_outgoing_payment_image(payment_image, operation_id)
3813 .await;
3814 Ok::<_, std::convert::Infallible>(claimer == operation_id)
3815 })
3816 },
3817 None,
3818 )
3819 .await
3820 .expect("Retries until the transaction commits")
3821 }
3822}
3823
3824#[async_trait]
3825impl IGatewayClientV1 for Gateway {
3826 async fn verify_preimage_authentication(
3827 &self,
3828 payment_hash: sha256::Hash,
3829 preimage_auth: sha256::Hash,
3830 contract: OutgoingContractAccount,
3831 ) -> std::result::Result<(), OutgoingPaymentError> {
3832 let mut dbtx = self.gateway_db.begin_transaction().await;
3833 if let Some(secret_hash) = dbtx.load_preimage_authentication(payment_hash).await {
3834 if secret_hash != preimage_auth {
3835 return Err(OutgoingPaymentError {
3836 error_type: OutgoingPaymentErrorType::InvalidInvoicePreimage,
3837 contract_id: contract.contract.contract_id(),
3838 contract: Some(contract),
3839 });
3840 }
3841 } else {
3842 dbtx.save_new_preimage_authentication(payment_hash, preimage_auth)
3845 .await;
3846 return dbtx
3847 .commit_tx_result()
3848 .await
3849 .map_err(|_| OutgoingPaymentError {
3850 error_type: OutgoingPaymentErrorType::InvoiceAlreadyPaid,
3851 contract_id: contract.contract.contract_id(),
3852 contract: Some(contract),
3853 });
3854 }
3855
3856 Ok(())
3857 }
3858
3859 async fn verify_pruned_invoice(&self, payment_data: PaymentData) -> anyhow::Result<()> {
3860 if matches!(payment_data, PaymentData::PrunedInvoice { .. }) {
3861 let lightning_context = self.get_lightning_context().await?;
3862
3863 ensure!(
3864 lightning_context.lnrpc.supports_private_payments(),
3865 "Private payments are not supported by the lightning node"
3866 );
3867 }
3868
3869 Ok(())
3870 }
3871
3872 async fn get_routing_fees(&self, federation_id: FederationId) -> Option<RoutingFees> {
3873 let mut gateway_dbtx = self.gateway_db.begin_transaction_nc().await;
3874 let lightning_fee = gateway_dbtx
3875 .load_federation_config(federation_id)
3876 .await?
3877 .lightning_fee;
3878
3879 RoutingFees::try_from(lightning_fee)
3883 .inspect_err(|err| {
3884 warn!(
3885 target: LOG_GATEWAY,
3886 %federation_id,
3887 err = %err.fmt_compact(),
3888 "Configured lightning fee cannot be used. Set a smaller fee with `set_fees`."
3889 );
3890 })
3891 .ok()
3892 }
3893
3894 async fn get_client(&self, federation_id: &FederationId) -> Option<Spanned<ClientHandleArc>> {
3895 self.federation_manager
3896 .read()
3897 .await
3898 .client(federation_id)
3899 .cloned()
3900 }
3901
3902 async fn get_client_for_invoice(
3903 &self,
3904 payment_data: PaymentData,
3905 ) -> Option<Spanned<ClientHandleArc>> {
3906 let rhints = payment_data.route_hints();
3907 let hop = rhints.first().and_then(|rh| rh.0.last())?;
3908
3909 let lightning_context = self.await_lightning_context().await;
3912 if hop.src_node_id != lightning_context.lightning_public_key {
3913 return None;
3914 }
3915
3916 self.federation_manager
3917 .read()
3918 .await
3919 .get_client_for_index(hop.short_channel_id)
3920 }
3921
3922 async fn pay(
3923 &self,
3924 payment_data: PaymentData,
3925 max_delay: u64,
3926 max_fee: Amount,
3927 ) -> std::result::Result<PayInvoiceResponse, LightningRpcError> {
3928 let lightning_context = self.await_lightning_context().await;
3931
3932 match payment_data {
3933 PaymentData::Invoice(invoice) => {
3934 lightning_context
3935 .lnrpc
3936 .pay(invoice, max_delay, max_fee)
3937 .await
3938 }
3939 PaymentData::PrunedInvoice(invoice) => {
3940 lightning_context
3941 .lnrpc
3942 .pay_private(invoice, max_delay, max_fee)
3943 .await
3944 }
3945 }
3946 }
3947
3948 async fn complete_htlc(
3949 &self,
3950 htlc: InterceptPaymentResponse,
3951 ) -> std::result::Result<(), LightningRpcError> {
3952 let lightning_context = self.await_lightning_context().await;
3954
3955 lightning_context.lnrpc.complete_htlc(htlc).await
3956 }
3957
3958 async fn is_lnv2_direct_swap(
3959 &self,
3960 payment_hash: sha256::Hash,
3961 amount: Amount,
3962 ) -> anyhow::Result<
3963 Option<(
3964 fedimint_lnv2_common::contracts::IncomingContract,
3965 ClientHandleArc,
3966 )>,
3967 > {
3968 let (contract, client) = self
3969 .get_registered_incoming_contract_and_client_v2(
3970 PaymentImage::Hash(payment_hash),
3971 amount.msats,
3972 )
3973 .await?;
3974 Ok(Some((contract, client)))
3975 }
3976}