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 iroh_server;
24mod metrics;
25pub mod rpc_server;
26mod types;
27
28use std::collections::{BTreeMap, BTreeSet};
29use std::env;
30use std::fmt::Display;
31use std::net::SocketAddr;
32use std::str::FromStr;
33use std::sync::Arc;
34use std::time::{Duration, UNIX_EPOCH};
35
36use anyhow::{Context, anyhow, ensure};
37use async_trait::async_trait;
38use bitcoin::hashes::sha256;
39use bitcoin::{Address, Network, Txid, secp256k1};
40use clap::Parser;
41use client::GatewayClientBuilder;
42pub use config::GatewayParameters;
43use config::{DatabaseBackend, GatewayOpts};
44use envs::FM_GATEWAY_SKIP_WAIT_FOR_SYNC_ENV;
45use error::FederationNotConnected;
46use events::ALL_GATEWAY_EVENTS;
47use federation_manager::FederationManager;
48use fedimint_bip39::{Bip39RootSecretStrategy, Language, Mnemonic};
49use fedimint_bitcoind::bitcoincore::BitcoindClient;
50use fedimint_bitcoind::{EsploraClient, IBitcoindRpc};
51use fedimint_client::module_init::ClientModuleInitRegistry;
52use fedimint_client::secret::RootSecretStrategy;
53use fedimint_client::{Client, ClientHandleArc};
54use fedimint_core::base32::{self, FEDIMINT_PREFIX};
55use fedimint_core::config::FederationId;
56use fedimint_core::core::OperationId;
57use fedimint_core::db::{Committable, Database, DatabaseTransaction, apply_migrations};
58use fedimint_core::envs::is_env_var_set;
59use fedimint_core::invite_code::InviteCode;
60use fedimint_core::module::CommonModuleInit;
61use fedimint_core::module::registry::ModuleDecoderRegistry;
62use fedimint_core::rustls::install_crypto_provider;
63use fedimint_core::secp256k1::PublicKey;
64use fedimint_core::secp256k1::schnorr::Signature;
65use fedimint_core::task::{TaskGroup, TaskHandle, TaskShutdownToken, sleep};
66use fedimint_core::time::duration_since_epoch;
67use fedimint_core::util::backoff_util::fibonacci_max_one_hour;
68use fedimint_core::util::{FmtCompact, FmtCompactAnyhow, SafeUrl, Spanned, retry};
69use fedimint_core::{
70 Amount, BitcoinAmountOrAll, PeerId, TieredCounts, crit, fedimint_build_code_version_env,
71 get_network_for_address,
72};
73use fedimint_eventlog::{DBTransactionEventLogExt, EventLogId, StructuredPaymentEvents};
74use fedimint_gateway_common::{
75 BackupPayload, ChainSource, CloseChannelsWithPeerRequest, CloseChannelsWithPeerResponse,
76 ConnectFedPayload, ConnectPeerRequest, ConnectorType, CreateInvoiceForOperatorPayload,
77 CreateOfferPayload, CreateOfferResponse, DepositAddressPayload, DepositAddressRecheckPayload,
78 FederationBalanceInfo, FederationConfig, FederationInfo, GatewayBalances, GatewayFedConfig,
79 GatewayInfo, GetInvoiceRequest, GetInvoiceResponse, LeaveFedPayload, LightningInfo,
80 LightningMode, ListTransactionsPayload, ListTransactionsResponse, MnemonicResponse,
81 OpenChannelRequest, PayInvoiceForOperatorPayload, PayOfferPayload, PayOfferResponse,
82 PaymentLogPayload, PaymentLogResponse, PaymentStats, PaymentSummaryPayload,
83 PaymentSummaryResponse, PeginFromOnchainPayload, ReceiveEcashPayload, ReceiveEcashResponse,
84 RegisteredProtocol, SendOnchainRequest, SetChannelFeesRequest, SetFeesPayload,
85 SetMnemonicPayload, SpendEcashPayload, SpendEcashResponse, V1_API_ENDPOINT, WithdrawPayload,
86 WithdrawPreviewPayload, WithdrawPreviewResponse, WithdrawResponse, WithdrawToOnchainPayload,
87};
88use fedimint_gateway_server_db::{GatewayDbtxNcExt as _, get_gatewayd_database_migrations};
89pub use fedimint_gateway_ui::IAdminGateway;
90use fedimint_gw_client::events::compute_lnv1_stats;
91use fedimint_gw_client::pay::{OutgoingPaymentError, OutgoingPaymentErrorType};
92use fedimint_gw_client::{
93 GatewayClientModule, GatewayExtPayStates, GatewayExtReceiveStates, IGatewayClientV1,
94 SwapParameters,
95};
96use fedimint_gwv2_client::events::compute_lnv2_stats;
97use fedimint_gwv2_client::{
98 EXPIRATION_DELTA_MINIMUM_V2, FinalReceiveState, GatewayClientModuleV2, IGatewayClientV2,
99};
100use fedimint_lightning::lnd::GatewayLndClient;
101use fedimint_lightning::{
102 CreateInvoiceRequest, ILnRpcClient, InterceptPaymentRequest, InterceptPaymentResponse,
103 InvoiceDescription, LightningContext, LightningRpcError, LnRpcTracked, Lnv2HoldInvoiceFilter,
104 PayInvoiceResponse, PaymentAction, RouteHtlcStream, ldk,
105};
106use fedimint_ln_client::pay::PaymentData;
107use fedimint_ln_common::LightningCommonInit;
108use fedimint_ln_common::config::LightningClientConfig;
109use fedimint_ln_common::contracts::outgoing::OutgoingContractAccount;
110use fedimint_ln_common::contracts::{IdentifiableContract, Preimage};
111use fedimint_lnurl::VerifyResponse;
112use fedimint_lnv2_common::Bolt11InvoiceDescription;
113use fedimint_lnv2_common::contracts::{IncomingContract, PaymentImage};
114use fedimint_lnv2_common::gateway_api::{
115 CreateBolt11InvoicePayload, PaymentFee, RoutingInfo, SendPaymentPayload,
116};
117use fedimint_logging::LOG_GATEWAY;
118use fedimint_mint_client::{MintClientInit, MintClientModule, OOBNotes};
119use fedimint_mintv2_client::{
120 MintClientInit as MintV2ClientInit, MintClientModule as MintV2ClientModule,
121};
122use fedimint_wallet_client::{PegOutFees, WalletClientInit, WalletClientModule, WithdrawState};
123use futures::stream::StreamExt;
124use lightning_invoice::{Bolt11Invoice, RoutingFees};
125use rand::rngs::OsRng;
126use tokio::sync::RwLock;
127use tracing::{debug, info, info_span, warn};
128
129use crate::envs::FM_GATEWAY_MNEMONIC_ENV;
130use crate::error::{AdminGatewayError, LNv1Error, LNv2Error, PublicGatewayError};
131use crate::events::get_events_for_duration;
132use crate::rpc_server::run_webserver;
133use crate::types::PrettyInterceptPaymentRequest;
134
135const GW_ANNOUNCEMENT_TTL: Duration = Duration::from_mins(10);
137
138const DEFAULT_NUM_ROUTE_HINTS: u32 = 1;
141
142pub const DEFAULT_NETWORK: Network = Network::Regtest;
144
145pub type Result<T> = std::result::Result<T, PublicGatewayError>;
146pub type AdminResult<T> = std::result::Result<T, AdminGatewayError>;
147
148const DB_FILE: &str = "gatewayd.db";
151
152const LDK_NODE_DB_FOLDER: &str = "ldk_node";
155
156#[cfg_attr(doc, aquamarine::aquamarine)]
157#[derive(Clone, Debug)]
170pub enum GatewayState {
171 NotConfigured {
172 mnemonic_sender: tokio::sync::broadcast::Sender<()>,
175 },
176 Disconnected,
177 Syncing,
178 Connected,
179 Running {
180 lightning_context: LightningContext,
181 },
182 ShuttingDown {
183 lightning_context: LightningContext,
184 },
185}
186
187impl Display for GatewayState {
188 fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
189 match self {
190 GatewayState::NotConfigured { .. } => write!(f, "NotConfigured"),
191 GatewayState::Disconnected => write!(f, "Disconnected"),
192 GatewayState::Syncing => write!(f, "Syncing"),
193 GatewayState::Connected => write!(f, "Connected"),
194 GatewayState::Running { .. } => write!(f, "Running"),
195 GatewayState::ShuttingDown { .. } => write!(f, "ShuttingDown"),
196 }
197 }
198}
199
200#[derive(Debug, Clone)]
203struct Registration {
204 endpoint_url: SafeUrl,
206
207 keypair: secp256k1::Keypair,
209}
210
211impl Registration {
212 pub async fn new(db: &Database, endpoint_url: SafeUrl, protocol: RegisteredProtocol) -> Self {
213 let keypair = Gateway::load_or_create_gateway_keypair(db, protocol).await;
214 Self {
215 endpoint_url,
216 keypair,
217 }
218 }
219}
220
221#[bon::bon]
222impl Gateway {
223 #[builder(start_fn = builder, finish_fn = build)]
238 pub async fn new_with_builder(
239 #[builder(start_fn)] lightning_mode: LightningMode,
240 #[builder(start_fn)] client_builder: GatewayClientBuilder,
241 #[builder(start_fn)] gateway_db: Database,
242 bcrypt_password_hash: bcrypt::HashParts,
243 bcrypt_liquidity_manager_password_hash: Option<bcrypt::HashParts>,
244 gateway_state: GatewayState,
245 chain_source: ChainSource,
246 #[builder(default = ([127, 0, 0, 1], 80).into())] listen: SocketAddr,
247 api_addr: Option<SafeUrl>,
248 #[builder(default = DEFAULT_NETWORK)] network: Network,
249 #[builder(default = DEFAULT_NUM_ROUTE_HINTS)] num_route_hints: u32,
250 #[builder(default = PaymentFee::TRANSACTION_FEE_DEFAULT)] default_routing_fees: PaymentFee,
251 #[builder(default = PaymentFee::TRANSACTION_FEE_DEFAULT)]
252 default_transaction_fees: PaymentFee,
253 iroh_listen: Option<SocketAddr>,
254 iroh_dns: Option<SafeUrl>,
255 #[builder(default)] iroh_relays: Vec<SafeUrl>,
256 metrics_listen: Option<SocketAddr>,
257 ) -> anyhow::Result<Gateway> {
258 let versioned_api = api_addr.map(|addr| {
259 addr.join(V1_API_ENDPOINT)
260 .expect("Failed to version gateway API address")
261 });
262
263 let metrics_listen = metrics_listen.unwrap_or_else(|| {
264 SocketAddr::new(
265 std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
266 listen.port() + 1,
267 )
268 });
269
270 Gateway::new(
271 lightning_mode,
272 GatewayParameters {
273 listen,
274 versioned_api,
275 bcrypt_password_hash,
276 bcrypt_liquidity_manager_password_hash,
277 network,
278 num_route_hints,
279 default_routing_fees,
280 default_transaction_fees,
281 iroh_listen,
282 iroh_dns,
283 iroh_relays,
284 skip_setup: true,
285 metrics_listen,
286 },
287 gateway_db,
288 client_builder,
289 gateway_state,
290 chain_source,
291 )
292 .await
293 }
294}
295
296enum ReceivePaymentStreamAction {
298 RetryAfterDelay,
299 NoRetry,
300}
301
302#[derive(Clone)]
303pub struct Gateway {
304 federation_manager: Arc<RwLock<FederationManager>>,
306
307 lightning_mode: LightningMode,
309
310 state: Arc<RwLock<GatewayState>>,
312
313 client_builder: GatewayClientBuilder,
316
317 gateway_db: Database,
319
320 listen: SocketAddr,
322
323 metrics_listen: SocketAddr,
325
326 task_group: TaskGroup,
328
329 bcrypt_password_hash: String,
331
332 bcrypt_liquidity_manager_password_hash: Option<String>,
335
336 num_route_hints: u32,
338
339 network: Network,
341
342 chain_source: ChainSource,
344
345 default_routing_fees: PaymentFee,
347
348 default_transaction_fees: PaymentFee,
350
351 iroh_sk: iroh::SecretKey,
353
354 iroh_listen: Option<SocketAddr>,
356
357 iroh_dns: Option<SafeUrl>,
359
360 iroh_relays: Vec<SafeUrl>,
363
364 registrations: BTreeMap<RegisteredProtocol, Registration>,
367}
368
369impl std::fmt::Debug for Gateway {
370 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
371 f.debug_struct("Gateway")
372 .field("federation_manager", &self.federation_manager)
373 .field("state", &self.state)
374 .field("client_builder", &self.client_builder)
375 .field("gateway_db", &self.gateway_db)
376 .field("listen", &self.listen)
377 .field("registrations", &self.registrations)
378 .finish_non_exhaustive()
379 }
380}
381
382struct WithdrawDetails {
384 amount: Amount,
385 mint_fees: Option<Amount>,
386 peg_out_fees: PegOutFees,
387}
388
389async fn withdraw_v2(
391 client: &ClientHandleArc,
392 wallet_module: &fedimint_walletv2_client::WalletClientModule,
393 address: &Address,
394 amount: BitcoinAmountOrAll,
395) -> AdminResult<WithdrawResponse> {
396 let fee = wallet_module
397 .send_fee()
398 .await
399 .map_err(|e| AdminGatewayError::WithdrawError {
400 failure_reason: e.to_string(),
401 })?;
402
403 let withdraw_amount = match amount {
404 BitcoinAmountOrAll::All => {
405 let balance = bitcoin::Amount::from_sat(
406 client
407 .get_balance_for_btc()
408 .await
409 .map_err(|err| {
410 AdminGatewayError::Unexpected(anyhow!(
411 "Balance not available: {}",
412 err.fmt_compact_anyhow()
413 ))
414 })?
415 .msats
416 / 1000,
417 );
418 balance
419 .checked_sub(fee)
420 .ok_or_else(|| AdminGatewayError::WithdrawError {
421 failure_reason: format!("Insufficient funds. Balance: {balance} Fee: {fee}"),
422 })?
423 }
424 BitcoinAmountOrAll::Amount(a) => a,
425 };
426
427 let operation_id = wallet_module
428 .send(
429 address.as_unchecked().clone(),
430 withdraw_amount,
431 Some(fee),
432 serde_json::Value::Null,
433 )
434 .await
435 .map_err(|e| AdminGatewayError::WithdrawError {
436 failure_reason: e.to_string(),
437 })?;
438
439 let result = wallet_module
440 .await_final_send_operation_state(operation_id)
441 .await
442 .map_err(|e| AdminGatewayError::WithdrawError {
443 failure_reason: e.to_string(),
444 })?;
445
446 let fees = PegOutFees::from_amount(fee);
447
448 match result {
449 fedimint_walletv2_client::FinalSendOperationState::Success(txid) => {
450 info!(target: LOG_GATEWAY, amount = %withdraw_amount, address = %address, "Sent funds via walletv2");
451 Ok(WithdrawResponse { txid, fees })
452 }
453 fedimint_walletv2_client::FinalSendOperationState::Aborted => {
454 Err(AdminGatewayError::WithdrawError {
455 failure_reason: "Withdrawal transaction was aborted".to_string(),
456 })
457 }
458 fedimint_walletv2_client::FinalSendOperationState::Failure => {
459 Err(AdminGatewayError::WithdrawError {
460 failure_reason: "Withdrawal failed".to_string(),
461 })
462 }
463 }
464}
465
466async fn calculate_max_withdrawable(
468 client: &ClientHandleArc,
469 address: &Address,
470) -> AdminResult<WithdrawDetails> {
471 let balance = client.get_balance_for_btc().await.map_err(|err| {
472 AdminGatewayError::Unexpected(anyhow!(
473 "Balance not available: {}",
474 err.fmt_compact_anyhow()
475 ))
476 })?;
477
478 let peg_out_fees = if let Ok(wallet_module) = client.get_first_module::<WalletClientModule>() {
479 wallet_module
480 .get_withdraw_fees(
481 address,
482 bitcoin::Amount::from_sat(balance.sats_round_down()),
483 )
484 .await?
485 } else if let Ok(wallet_module) =
486 client.get_first_module::<fedimint_walletv2_client::WalletClientModule>()
487 {
488 let fee = wallet_module
489 .send_fee()
490 .await
491 .map_err(|e| AdminGatewayError::WithdrawError {
492 failure_reason: e.to_string(),
493 })?;
494 PegOutFees::from_amount(fee)
495 } else {
496 return Err(AdminGatewayError::Unexpected(anyhow!(
497 "No wallet module found"
498 )));
499 };
500
501 let max_withdrawable_before_mint_fees = balance
502 .checked_sub(peg_out_fees.amount().into())
503 .ok_or_else(|| AdminGatewayError::WithdrawError {
504 failure_reason: "Insufficient balance to cover peg-out fees".to_string(),
505 })?;
506
507 let mint_fees = if let Ok(mint_module) = client.get_first_module::<MintClientModule>() {
509 mint_module.estimate_spend_all_fees().await
510 } else {
511 Amount::ZERO
512 };
513
514 let max_withdrawable = max_withdrawable_before_mint_fees.saturating_sub(mint_fees);
515
516 Ok(WithdrawDetails {
517 amount: max_withdrawable,
518 mint_fees: Some(mint_fees),
519 peg_out_fees,
520 })
521}
522
523impl Gateway {
524 fn get_bitcoind_client(
527 opts: &GatewayOpts,
528 network: bitcoin::Network,
529 gateway_id: &PublicKey,
530 ) -> anyhow::Result<(BitcoindClient, ChainSource)> {
531 let bitcoind_username = opts
532 .bitcoind_username
533 .clone()
534 .expect("FM_BITCOIND_URL is set but FM_BITCOIND_USERNAME is not");
535 let url = opts.bitcoind_url.clone().expect("No bitcoind url set");
536 let password = opts
537 .bitcoind_password
538 .clone()
539 .expect("FM_BITCOIND_URL is set but FM_BITCOIND_PASSWORD is not");
540
541 let chain_source = ChainSource::Bitcoind {
542 username: bitcoind_username.clone(),
543 password: password.clone(),
544 server_url: url.clone(),
545 };
546 let wallet_name = format!("gatewayd-{gateway_id}");
547 let client = BitcoindClient::new(&url, bitcoind_username, password, &wallet_name, network)?;
548 Ok((client, chain_source))
549 }
550
551 pub async fn new_with_default_modules(
554 mnemonic_sender: tokio::sync::broadcast::Sender<()>,
555 ) -> anyhow::Result<Gateway> {
556 let opts = GatewayOpts::parse();
557 let gateway_parameters = opts.to_gateway_parameters()?;
558 let decoders = ModuleDecoderRegistry::default();
559
560 let db_path = opts.data_dir.join(DB_FILE);
561 let gateway_db = match opts.db_backend {
562 DatabaseBackend::RocksDb => {
563 debug!(target: LOG_GATEWAY, "Using RocksDB database backend");
564 Database::new(
565 fedimint_rocksdb::RocksDb::build(db_path).open().await?,
566 decoders,
567 )
568 }
569 DatabaseBackend::CursedRedb => {
570 debug!(target: LOG_GATEWAY, "Using CursedRedb database backend");
571 Database::new(
572 fedimint_cursed_redb::MemAndRedb::new(db_path).await?,
573 decoders,
574 )
575 }
576 };
577
578 apply_migrations(
581 &gateway_db,
582 (),
583 "gatewayd".to_string(),
584 get_gatewayd_database_migrations(),
585 None,
586 None,
587 )
588 .await?;
589
590 let http_id = Self::load_or_create_gateway_keypair(&gateway_db, RegisteredProtocol::Http)
593 .await
594 .public_key();
595 let (dyn_bitcoin_rpc, chain_source) =
596 match (opts.bitcoind_url.as_ref(), opts.esplora_url.as_ref()) {
597 (Some(_), None) => {
598 let (client, chain_source) =
599 Self::get_bitcoind_client(&opts, gateway_parameters.network, &http_id)?;
600 (client.into_dyn(), chain_source)
601 }
602 (None, Some(url)) => {
603 let client = EsploraClient::new(url)
604 .expect("Could not create EsploraClient")
605 .into_dyn();
606 let chain_source = ChainSource::Esplora {
607 server_url: url.clone(),
608 };
609 (client, chain_source)
610 }
611 (Some(_), Some(_)) => {
612 let (client, chain_source) =
614 Self::get_bitcoind_client(&opts, gateway_parameters.network, &http_id)?;
615 (client.into_dyn(), chain_source)
616 }
617 _ => unreachable!("ArgGroup already enforced XOR relation"),
618 };
619
620 let mut registry = ClientModuleInitRegistry::new();
623 registry.attach(MintClientInit);
624 registry.attach(MintV2ClientInit);
625 registry.attach(WalletClientInit::new(dyn_bitcoin_rpc));
626 registry.attach(fedimint_walletv2_client::WalletClientInit);
627
628 let client_builder =
629 GatewayClientBuilder::new(opts.data_dir.clone(), registry, opts.db_backend).await?;
630
631 let gateway_state = if Self::load_mnemonic(&gateway_db).await.is_some() {
632 GatewayState::Disconnected
633 } else {
634 if gateway_parameters.skip_setup {
637 let mnemonic = if let Ok(words) = std::env::var(FM_GATEWAY_MNEMONIC_ENV) {
638 info!(target: LOG_GATEWAY, "Using provided mnemonic from environment variable");
639 Mnemonic::parse_in_normalized(Language::English, words.as_str()).map_err(
640 |e| {
641 AdminGatewayError::MnemonicError(anyhow!(format!(
642 "Seed phrase provided in environment was invalid {e:?}"
643 )))
644 },
645 )?
646 } else {
647 debug!(target: LOG_GATEWAY, "Generating mnemonic and writing entropy to client storage");
648 Bip39RootSecretStrategy::<12>::random(&mut OsRng)
649 };
650
651 Client::store_encodable_client_secret(&gateway_db, mnemonic.to_entropy())
652 .await
653 .map_err(AdminGatewayError::MnemonicError)?;
654 GatewayState::Disconnected
655 } else {
656 GatewayState::NotConfigured { mnemonic_sender }
657 }
658 };
659
660 info!(
661 target: LOG_GATEWAY,
662 version = %fedimint_build_code_version_env!(),
663 "Starting gatewayd",
664 );
665
666 Gateway::new(
667 opts.mode,
668 gateway_parameters,
669 gateway_db,
670 client_builder,
671 gateway_state,
672 chain_source,
673 )
674 .await
675 }
676
677 async fn new(
680 lightning_mode: LightningMode,
681 gateway_parameters: GatewayParameters,
682 gateway_db: Database,
683 client_builder: GatewayClientBuilder,
684 gateway_state: GatewayState,
685 chain_source: ChainSource,
686 ) -> anyhow::Result<Gateway> {
687 let num_route_hints = gateway_parameters.num_route_hints;
688 let network = gateway_parameters.network;
689
690 let task_group = TaskGroup::new();
691 task_group.install_kill_handler();
692
693 let mut registrations = BTreeMap::new();
694 if let Some(http_url) = gateway_parameters.versioned_api {
695 registrations.insert(
696 RegisteredProtocol::Http,
697 Registration::new(&gateway_db, http_url, RegisteredProtocol::Http).await,
698 );
699 }
700
701 let iroh_sk = Self::load_or_create_iroh_key(&gateway_db).await;
702 if gateway_parameters.iroh_listen.is_some() {
703 let endpoint_url = SafeUrl::parse(&format!("iroh://{}", iroh_sk.public()))?;
704 registrations.insert(
705 RegisteredProtocol::Iroh,
706 Registration::new(&gateway_db, endpoint_url, RegisteredProtocol::Iroh).await,
707 );
708 }
709
710 Ok(Self {
711 federation_manager: Arc::new(RwLock::new(FederationManager::new())),
712 lightning_mode,
713 state: Arc::new(RwLock::new(gateway_state)),
714 client_builder,
715 gateway_db: gateway_db.clone(),
716 listen: gateway_parameters.listen,
717 metrics_listen: gateway_parameters.metrics_listen,
718 task_group,
719 bcrypt_password_hash: gateway_parameters.bcrypt_password_hash.to_string(),
720 bcrypt_liquidity_manager_password_hash: gateway_parameters
721 .bcrypt_liquidity_manager_password_hash
722 .map(|h| h.to_string()),
723 num_route_hints,
724 network,
725 chain_source,
726 default_routing_fees: gateway_parameters.default_routing_fees,
727 default_transaction_fees: gateway_parameters.default_transaction_fees,
728 iroh_sk,
729 iroh_dns: gateway_parameters.iroh_dns,
730 iroh_relays: gateway_parameters.iroh_relays,
731 iroh_listen: gateway_parameters.iroh_listen,
732 registrations,
733 })
734 }
735
736 async fn load_or_create_gateway_keypair(
737 gateway_db: &Database,
738 protocol: RegisteredProtocol,
739 ) -> secp256k1::Keypair {
740 let mut dbtx = gateway_db.begin_transaction().await;
741 let keypair = dbtx.load_or_create_gateway_keypair(protocol).await;
742 dbtx.commit_tx().await;
743 keypair
744 }
745
746 async fn load_or_create_iroh_key(gateway_db: &Database) -> iroh::SecretKey {
749 let mut dbtx = gateway_db.begin_transaction().await;
750 let iroh_sk = dbtx.load_or_create_iroh_key().await;
751 dbtx.commit_tx().await;
752 iroh_sk
753 }
754
755 pub async fn http_gateway_id(&self) -> PublicKey {
756 Self::load_or_create_gateway_keypair(&self.gateway_db, RegisteredProtocol::Http)
757 .await
758 .public_key()
759 }
760
761 async fn get_state(&self) -> GatewayState {
762 self.state.read().await.clone()
763 }
764
765 pub async fn dump_database(
768 dbtx: &mut DatabaseTransaction<'_>,
769 prefix_names: Vec<String>,
770 ) -> BTreeMap<String, Box<dyn erased_serde::Serialize + Send>> {
771 dbtx.dump_database(prefix_names).await
772 }
773
774 pub async fn run(
779 self,
780 runtime: Arc<tokio::runtime::Runtime>,
781 mnemonic_receiver: tokio::sync::broadcast::Receiver<()>,
782 ) -> anyhow::Result<TaskShutdownToken> {
783 install_crypto_provider().await;
784 self.register_clients_timer();
785 self.load_clients().await?;
786 self.start_gateway(runtime, mnemonic_receiver.resubscribe());
787 self.spawn_backup_task();
788 fedimint_metrics::spawn_api_server(self.metrics_listen, self.task_group.clone()).await?;
790 let handle = self.task_group.make_handle();
792 run_webserver(Arc::new(self), mnemonic_receiver.resubscribe()).await?;
793 let shutdown_receiver = handle.make_shutdown_rx();
794 Ok(shutdown_receiver)
795 }
796
797 fn spawn_backup_task(&self) {
800 let self_copy = self.clone();
801 self.task_group
802 .spawn_cancellable_silent("backup ecash", async move {
803 const BACKUP_UPDATE_INTERVAL: Duration = Duration::from_hours(1);
804 let mut interval = tokio::time::interval(BACKUP_UPDATE_INTERVAL);
805 interval.tick().await;
806 loop {
807 {
808 let mut dbtx = self_copy.gateway_db.begin_transaction().await;
809 self_copy.backup_all_federations(&mut dbtx).await;
810 dbtx.commit_tx().await;
811 interval.tick().await;
812 }
813 }
814 });
815 }
816
817 pub async fn backup_all_federations(&self, dbtx: &mut DatabaseTransaction<'_, Committable>) {
821 const BACKUP_THRESHOLD_DURATION: Duration = Duration::from_hours(24);
824
825 let now = fedimint_core::time::now();
826 let threshold = now
827 .checked_sub(BACKUP_THRESHOLD_DURATION)
828 .expect("Cannot be negative");
829 for (id, last_backup) in dbtx.load_backup_records().await {
830 match last_backup {
831 Some(backup_time) if backup_time < threshold => {
832 let fed_manager = self.federation_manager.read().await;
833 fed_manager.backup_federation(&id, dbtx, now).await;
834 }
835 None => {
836 let fed_manager = self.federation_manager.read().await;
837 fed_manager.backup_federation(&id, dbtx, now).await;
838 }
839 _ => {}
840 }
841 }
842 }
843
844 fn start_gateway(
847 &self,
848 runtime: Arc<tokio::runtime::Runtime>,
849 mut mnemonic_receiver: tokio::sync::broadcast::Receiver<()>,
850 ) {
851 const PAYMENT_STREAM_RETRY_SECONDS: u64 = 60;
852
853 let self_copy = self.clone();
854 let tg = self.task_group.clone();
855 self.task_group.spawn(
856 "Subscribe to intercepted lightning payments in stream",
857 |handle| async move {
858 loop {
860 if handle.is_shutting_down() {
861 info!(target: LOG_GATEWAY, "Gateway lightning payment stream handler loop is shutting down");
862 break;
863 }
864
865 if let GatewayState::NotConfigured{ .. } = self_copy.get_state().await {
866 info!(
867 target: LOG_GATEWAY,
868 "Waiting for the mnemonic to be set before starting lightning receive loop."
869 );
870 info!(
871 target: LOG_GATEWAY,
872 "You might need to provide it from the UI or refer to documentation w.r.t how to initialize it."
873 );
874
875 let _ = mnemonic_receiver.recv().await;
876 info!(
877 target: LOG_GATEWAY,
878 "Received mnemonic, attempting to start lightning receive loop"
879 );
880 }
881
882 let payment_stream_task_group = tg.make_subgroup();
883 let lnrpc_route = self_copy.create_lightning_client(runtime.clone()).await;
884
885 debug!(target: LOG_GATEWAY, "Establishing lightning payment stream...");
886 let (stream, ln_client) = match lnrpc_route.route_htlcs(&payment_stream_task_group).await
887 {
888 Ok((stream, ln_client)) => (stream, ln_client),
889 Err(err) => {
890 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failed to open lightning payment stream");
891 if let Err(err) = payment_stream_task_group.shutdown_join_all(None).await {
899 crit!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Lightning payment stream task group shutdown");
900 }
901 sleep(Duration::from_secs(PAYMENT_STREAM_RETRY_SECONDS)).await;
902 continue
903 }
904 };
905
906 self_copy.set_gateway_state(GatewayState::Connected).await;
908 info!(target: LOG_GATEWAY, "Established lightning payment stream");
909
910 let route_payments_response =
911 self_copy.route_lightning_payments(&handle, stream, ln_client).await;
912
913 self_copy.set_gateway_state(GatewayState::Disconnected).await;
914 if let Err(err) = payment_stream_task_group.shutdown_join_all(None).await {
915 crit!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Lightning payment stream task group shutdown");
916 }
917
918 self_copy.unannounce_from_all_federations().await;
919
920 match route_payments_response {
921 ReceivePaymentStreamAction::RetryAfterDelay => {
922 warn!(target: LOG_GATEWAY, retry_interval = %PAYMENT_STREAM_RETRY_SECONDS, "Disconnected from lightning node");
923 sleep(Duration::from_secs(PAYMENT_STREAM_RETRY_SECONDS)).await;
924 }
925 ReceivePaymentStreamAction::NoRetry => break,
926 }
927 }
928 },
929 );
930 }
931
932 async fn route_lightning_payments<'a>(
936 &'a self,
937 handle: &TaskHandle,
938 mut stream: RouteHtlcStream<'a>,
939 ln_client: Arc<dyn ILnRpcClient>,
940 ) -> ReceivePaymentStreamAction {
941 let LightningInfo::Connected {
942 public_key: lightning_public_key,
943 alias: lightning_alias,
944 network: lightning_network,
945 block_height: _,
946 synced_to_chain,
947 } = ln_client.parsed_node_info().await
948 else {
949 warn!(target: LOG_GATEWAY, "Failed to retrieve Lightning info");
950 return ReceivePaymentStreamAction::RetryAfterDelay;
951 };
952
953 assert!(
954 self.network == lightning_network,
955 "Lightning node network does not match Gateway's network. LN: {lightning_network} Gateway: {}",
956 self.network
957 );
958
959 if synced_to_chain || is_env_var_set(FM_GATEWAY_SKIP_WAIT_FOR_SYNC_ENV) {
960 info!(target: LOG_GATEWAY, "Gateway is already synced to chain");
961 } else {
962 self.set_gateway_state(GatewayState::Syncing).await;
963 info!(target: LOG_GATEWAY, "Waiting for chain sync");
964 if let Err(err) = ln_client.wait_for_chain_sync().await {
965 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failed to wait for chain sync");
966 return ReceivePaymentStreamAction::RetryAfterDelay;
967 }
968 }
969
970 let lightning_context = LightningContext {
971 lnrpc: LnRpcTracked::new(ln_client, "gateway"),
972 lightning_public_key,
973 lightning_alias,
974 lightning_network,
975 };
976 self.set_gateway_state(GatewayState::Running { lightning_context })
977 .await;
978 info!(target: LOG_GATEWAY, "Gateway is running");
979
980 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
981 let mut dbtx = self.gateway_db.begin_transaction_nc().await;
984 let all_federations_configs =
985 dbtx.load_federation_configs().await.into_iter().collect();
986 self.register_federations(&all_federations_configs, &self.task_group)
987 .await;
988 }
989
990 let htlc_task_group = self.task_group.make_subgroup();
993 if handle
994 .cancel_on_shutdown(async move {
995 loop {
996 let payment_request_or = tokio::select! {
997 payment_request_or = stream.next() => {
998 payment_request_or
999 }
1000 () = self.is_shutting_down_safely() => {
1001 break;
1002 }
1003 };
1004
1005 let Some(payment_request) = payment_request_or else {
1006 warn!(
1007 target: LOG_GATEWAY,
1008 "Unexpected response from incoming lightning payment stream. Shutting down payment processor"
1009 );
1010 break;
1011 };
1012
1013 let state_guard = self.state.read().await;
1014 if let GatewayState::Running { ref lightning_context } = *state_guard {
1015 let gateway = self.clone();
1017 let lightning_context = lightning_context.clone();
1018 htlc_task_group.spawn_cancellable_silent(
1019 "handle_lightning_payment",
1020 async move {
1021 let start = fedimint_core::time::now();
1022 let outcome = gateway
1023 .handle_lightning_payment(payment_request, &lightning_context)
1024 .await;
1025 metrics::HTLC_HANDLING_DURATION_SECONDS
1026 .with_label_values(&[outcome])
1027 .observe(
1028 fedimint_core::time::now()
1029 .duration_since(start)
1030 .unwrap_or_default()
1031 .as_secs_f64(),
1032 );
1033 },
1034 );
1035 } else {
1036 warn!(
1037 target: LOG_GATEWAY,
1038 state = %state_guard,
1039 "Gateway isn't in a running state, cannot handle incoming payments."
1040 );
1041 break;
1042 }
1043 }
1044 })
1045 .await
1046 .is_ok()
1047 {
1048 warn!(target: LOG_GATEWAY, "Lightning payment stream connection broken. Gateway is disconnected");
1049 ReceivePaymentStreamAction::RetryAfterDelay
1050 } else {
1051 info!(target: LOG_GATEWAY, "Received shutdown signal");
1052 ReceivePaymentStreamAction::NoRetry
1053 }
1054 }
1055
1056 async fn is_shutting_down_safely(&self) {
1059 loop {
1060 if let GatewayState::ShuttingDown { .. } = self.get_state().await {
1061 return;
1062 }
1063
1064 fedimint_core::task::sleep(Duration::from_secs(1)).await;
1065 }
1066 }
1067
1068 async fn handle_lightning_payment(
1078 &self,
1079 payment_request: InterceptPaymentRequest,
1080 lightning_context: &LightningContext,
1081 ) -> &'static str {
1082 info!(
1083 target: LOG_GATEWAY,
1084 lightning_payment = %PrettyInterceptPaymentRequest(&payment_request),
1085 "Intercepting lightning payment",
1086 );
1087
1088 let lnv2_start = fedimint_core::time::now();
1089 let lnv2_result = self
1090 .try_handle_lightning_payment_lnv2(&payment_request, lightning_context)
1091 .await;
1092 let lnv2_outcome = if lnv2_result.is_ok() {
1093 "success"
1094 } else {
1095 "error"
1096 };
1097 metrics::HTLC_LNV2_ATTEMPT_DURATION_SECONDS
1098 .with_label_values(&[lnv2_outcome])
1099 .observe(
1100 fedimint_core::time::now()
1101 .duration_since(lnv2_start)
1102 .unwrap_or_default()
1103 .as_secs_f64(),
1104 );
1105 if lnv2_result.is_ok() {
1106 return "lnv2";
1107 }
1108
1109 let lnv1_start = fedimint_core::time::now();
1110 let lnv1_result = self
1111 .try_handle_lightning_payment_ln_legacy(&payment_request)
1112 .await;
1113 let lnv1_outcome = if lnv1_result.is_ok() {
1114 "success"
1115 } else {
1116 "error"
1117 };
1118 metrics::HTLC_LNV1_ATTEMPT_DURATION_SECONDS
1119 .with_label_values(&[lnv1_outcome])
1120 .observe(
1121 fedimint_core::time::now()
1122 .duration_since(lnv1_start)
1123 .unwrap_or_default()
1124 .as_secs_f64(),
1125 );
1126 if lnv1_result.is_ok() {
1127 return "lnv1";
1128 }
1129
1130 let is_federation_scid = match payment_request.short_channel_id {
1135 Some(scid) => self
1136 .federation_manager
1137 .read()
1138 .await
1139 .get_client_for_index(scid)
1140 .is_some(),
1141 None => false,
1142 };
1143
1144 if is_federation_scid {
1145 warn!(
1152 target: LOG_GATEWAY,
1153 payment_hash = %payment_request.payment_hash,
1154 short_channel_id = ?payment_request.short_channel_id,
1155 amount_msat = payment_request.amount_msat,
1156 incoming_chan_id = payment_request.incoming_chan_id,
1157 htlc_id = payment_request.htlc_id,
1158 lnv2_err = ?lnv2_result.as_ref().err(),
1159 lnv1_err = ?lnv1_result.as_ref().err(),
1160 "Unmatched lightning payment for federation scid: cancelling HTLC",
1161 );
1162 Self::cancel_unmatched_lightning_payment(payment_request, lightning_context).await;
1163 "cancel"
1164 } else {
1165 Self::forward_lightning_payment(payment_request, lightning_context).await;
1170 "forward"
1171 }
1172 }
1173
1174 async fn try_handle_lightning_payment_lnv2(
1177 &self,
1178 htlc_request: &InterceptPaymentRequest,
1179 lightning_context: &LightningContext,
1180 ) -> Result<()> {
1181 let (contract, client) = self
1187 .get_registered_incoming_contract_and_client_v2(
1188 PaymentImage::Hash(htlc_request.payment_hash),
1189 htlc_request.amount_msat,
1190 )
1191 .await?;
1192
1193 if let Err(err) = client
1194 .get_first_module::<GatewayClientModuleV2>()
1195 .expect("Must have client module")
1196 .relay_incoming_htlc(
1197 htlc_request.payment_hash,
1198 htlc_request.incoming_chan_id,
1199 htlc_request.htlc_id,
1200 contract,
1201 htlc_request.amount_msat,
1202 )
1203 .await
1204 {
1205 warn!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Error relaying incoming lightning payment");
1206
1207 let outcome = InterceptPaymentResponse {
1208 action: PaymentAction::Cancel,
1209 payment_hash: htlc_request.payment_hash,
1210 incoming_chan_id: htlc_request.incoming_chan_id,
1211 htlc_id: htlc_request.htlc_id,
1212 };
1213
1214 if let Err(err) = lightning_context.lnrpc.complete_htlc(outcome).await {
1215 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Error sending HTLC response to lightning node");
1216 }
1217 }
1218
1219 Ok(())
1220 }
1221
1222 async fn try_handle_lightning_payment_ln_legacy(
1225 &self,
1226 htlc_request: &InterceptPaymentRequest,
1227 ) -> Result<()> {
1228 let Some(federation_index) = htlc_request.short_channel_id else {
1230 return Err(PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1231 "Incoming payment has not last hop short channel id".to_string(),
1232 )));
1233 };
1234
1235 let Some(client) = self
1236 .federation_manager
1237 .read()
1238 .await
1239 .get_client_for_index(federation_index)
1240 else {
1241 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())));
1242 };
1243
1244 client
1245 .borrow()
1246 .with(|client| async {
1247 let htlc = htlc_request.clone().try_into();
1248 match htlc {
1249 Ok(htlc) => {
1250 let lnv1 =
1251 client
1252 .get_first_module::<GatewayClientModule>()
1253 .map_err(|_| {
1254 PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1255 "Federation does not have LNv1 module".to_string(),
1256 ))
1257 })?;
1258 match lnv1.gateway_handle_intercepted_htlc(htlc).await {
1259 Ok(_) => Ok(()),
1260 Err(e) => Err(PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1261 format!("Error intercepting lightning payment {e:?}"),
1262 ))),
1263 }
1264 }
1265 _ => Err(PublicGatewayError::LNv1(LNv1Error::IncomingPayment(
1266 "Could not convert InterceptHtlcResult into an HTLC".to_string(),
1267 ))),
1268 }
1269 })
1270 .await
1271 }
1272
1273 async fn cancel_unmatched_lightning_payment(
1287 htlc_request: InterceptPaymentRequest,
1288 lightning_context: &LightningContext,
1289 ) {
1290 let outcome = InterceptPaymentResponse {
1291 action: PaymentAction::Cancel,
1292 payment_hash: htlc_request.payment_hash,
1293 incoming_chan_id: htlc_request.incoming_chan_id,
1294 htlc_id: htlc_request.htlc_id,
1295 };
1296
1297 if let Err(err) = lightning_context.lnrpc.complete_htlc(outcome).await {
1298 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Error sending lightning payment response to lightning node");
1299 }
1300 }
1301
1302 async fn forward_lightning_payment(
1307 htlc_request: InterceptPaymentRequest,
1308 lightning_context: &LightningContext,
1309 ) {
1310 let outcome = InterceptPaymentResponse {
1311 action: PaymentAction::Forward,
1312 payment_hash: htlc_request.payment_hash,
1313 incoming_chan_id: htlc_request.incoming_chan_id,
1314 htlc_id: htlc_request.htlc_id,
1315 };
1316
1317 if let Err(err) = lightning_context.lnrpc.complete_htlc(outcome).await {
1318 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Error sending lightning payment response to lightning node");
1319 }
1320 }
1321
1322 async fn set_gateway_state(&self, state: GatewayState) {
1324 let mut lock = self.state.write().await;
1325 *lock = state;
1326 }
1327
1328 pub async fn handle_get_federation_config(
1331 &self,
1332 federation_id_or: Option<FederationId>,
1333 ) -> AdminResult<GatewayFedConfig> {
1334 if !matches!(self.get_state().await, GatewayState::Running { .. }) {
1335 return Ok(GatewayFedConfig {
1336 federations: BTreeMap::new(),
1337 });
1338 }
1339
1340 let federations = if let Some(federation_id) = federation_id_or {
1341 let mut federations = BTreeMap::new();
1342 federations.insert(
1343 federation_id,
1344 self.federation_manager
1345 .read()
1346 .await
1347 .get_federation_config(federation_id)
1348 .await?,
1349 );
1350 federations
1351 } else {
1352 self.federation_manager
1353 .read()
1354 .await
1355 .get_all_federation_configs()
1356 .await
1357 };
1358
1359 Ok(GatewayFedConfig { federations })
1360 }
1361
1362 pub async fn handle_address_msg(&self, payload: DepositAddressPayload) -> AdminResult<Address> {
1365 let client = self.select_client(payload.federation_id).await?;
1366
1367 if let Ok(wallet_module) = client.value().get_first_module::<WalletClientModule>() {
1368 let address = wallet_module
1369 .allocate_deposit_address_expert_only(())
1370 .await?
1371 .address;
1372 Ok(address)
1373 } else if let Ok(wallet_module) = client
1374 .value()
1375 .get_first_module::<fedimint_walletv2_client::WalletClientModule>()
1376 {
1377 Ok(wallet_module.receive().await)
1378 } else {
1379 Err(AdminGatewayError::Unexpected(anyhow!(
1380 "No wallet module found"
1381 )))
1382 }
1383 }
1384
1385 async fn handle_pay_invoice_msg(
1388 &self,
1389 payload: fedimint_ln_client::pay::PayInvoicePayload,
1390 ) -> Result<Preimage> {
1391 let GatewayState::Running { .. } = self.get_state().await else {
1392 return Err(PublicGatewayError::Lightning(
1393 LightningRpcError::FailedToConnect,
1394 ));
1395 };
1396
1397 debug!(target: LOG_GATEWAY, "Handling pay invoice message");
1398 let client = self.select_client(payload.federation_id).await?;
1399 let contract_id = payload.contract_id;
1400 let gateway_module = &client
1401 .value()
1402 .get_first_module::<GatewayClientModule>()
1403 .map_err(LNv1Error::OutgoingPayment)
1404 .map_err(PublicGatewayError::LNv1)?;
1405 let operation_id = gateway_module
1406 .gateway_pay_bolt11_invoice(payload)
1407 .await
1408 .map_err(LNv1Error::OutgoingPayment)
1409 .map_err(PublicGatewayError::LNv1)?;
1410 let mut updates = gateway_module
1411 .gateway_subscribe_ln_pay(operation_id)
1412 .await
1413 .map_err(LNv1Error::OutgoingPayment)
1414 .map_err(PublicGatewayError::LNv1)?
1415 .into_stream();
1416 while let Some(update) = updates.next().await {
1417 match update {
1418 GatewayExtPayStates::Success { preimage, .. } => {
1419 debug!(target: LOG_GATEWAY, contract_id = %contract_id, "Successfully paid invoice");
1420 return Ok(preimage);
1421 }
1422 GatewayExtPayStates::Fail {
1423 error,
1424 error_message,
1425 } => {
1426 return Err(PublicGatewayError::LNv1(LNv1Error::OutgoingContract {
1427 error: Box::new(error),
1428 message: format!(
1429 "{error_message} while paying invoice with contract id {contract_id}"
1430 ),
1431 }));
1432 }
1433 GatewayExtPayStates::Canceled { error } => {
1434 return Err(PublicGatewayError::LNv1(LNv1Error::OutgoingContract {
1435 error: Box::new(error.clone()),
1436 message: format!(
1437 "Cancelled with {error} while paying invoice with contract id {contract_id}"
1438 ),
1439 }));
1440 }
1441 GatewayExtPayStates::Created => {
1442 debug!(target: LOG_GATEWAY, contract_id = %contract_id, "Start pay invoice state machine");
1443 }
1444 other => {
1445 debug!(target: LOG_GATEWAY, state = ?other, contract_id = %contract_id, "Got state while paying invoice");
1446 }
1447 }
1448 }
1449
1450 Err(PublicGatewayError::LNv1(LNv1Error::OutgoingPayment(
1451 anyhow!("Ran out of state updates while paying invoice"),
1452 )))
1453 }
1454
1455 pub async fn handle_backup_msg(
1458 &self,
1459 BackupPayload { federation_id }: BackupPayload,
1460 ) -> AdminResult<()> {
1461 let federation_manager = self.federation_manager.read().await;
1462 let client = federation_manager
1463 .client(&federation_id)
1464 .ok_or(AdminGatewayError::ClientCreationError(anyhow::anyhow!(
1465 format!("Gateway has not connected to {federation_id}")
1466 )))?
1467 .value();
1468 let metadata: BTreeMap<String, String> = BTreeMap::new();
1469 #[allow(deprecated)]
1470 client
1471 .backup_to_federation(fedimint_client::backup::Metadata::from_json_serialized(
1472 metadata,
1473 ))
1474 .await?;
1475 Ok(())
1476 }
1477
1478 pub async fn handle_recheck_address_msg(
1480 &self,
1481 payload: DepositAddressRecheckPayload,
1482 ) -> AdminResult<()> {
1483 let client = self.select_client(payload.federation_id).await?;
1484
1485 if let Ok(wallet_module) = client.value().get_first_module::<WalletClientModule>() {
1486 wallet_module
1487 .recheck_pegin_address_by_address(payload.address)
1488 .await?;
1489 Ok(())
1490 } else if client
1491 .value()
1492 .get_first_module::<fedimint_walletv2_client::WalletClientModule>()
1493 .is_ok()
1494 {
1495 Ok(())
1497 } else {
1498 Err(AdminGatewayError::Unexpected(anyhow!(
1499 "No wallet module found"
1500 )))
1501 }
1502 }
1503
1504 pub async fn handle_receive_ecash_msg(
1506 &self,
1507 payload: ReceiveEcashPayload,
1508 ) -> Result<ReceiveEcashResponse> {
1509 let federation_id_prefix = base32::decode_prefixed::<fedimint_mintv2_client::ECash>(
1511 FEDIMINT_PREFIX,
1512 &payload.notes,
1513 )
1514 .ok()
1515 .and_then(|e| e.mint())
1516 .map(|id| id.to_prefix())
1517 .or_else(|| {
1518 OOBNotes::from_str(&payload.notes)
1519 .ok()
1520 .map(|n| n.federation_id_prefix())
1521 })
1522 .ok_or_else(|| PublicGatewayError::ReceiveEcashError {
1523 failure_reason: "Invalid ecash format: could not parse as ECash or OOBNotes"
1524 .to_string(),
1525 })?;
1526
1527 let client = self
1528 .federation_manager
1529 .read()
1530 .await
1531 .get_client_for_federation_id_prefix(federation_id_prefix)
1532 .ok_or(FederationNotConnected {
1533 federation_id_prefix,
1534 })?;
1535
1536 if let Ok(mint) = client.value().get_first_module::<MintClientModule>() {
1538 let notes = OOBNotes::from_str(&payload.notes).map_err(|e| {
1539 PublicGatewayError::ReceiveEcashError {
1540 failure_reason: format!("Expected OOBNotes for MintV1 federation: {e}"),
1541 }
1542 })?;
1543 let amount = notes.total_amount();
1544
1545 let operation_id = mint.reissue_external_notes(notes, ()).await.map_err(|e| {
1546 PublicGatewayError::ReceiveEcashError {
1547 failure_reason: e.to_string(),
1548 }
1549 })?;
1550 if payload.wait {
1551 let mut updates = mint
1552 .subscribe_reissue_external_notes(operation_id)
1553 .await
1554 .unwrap()
1555 .into_stream();
1556
1557 while let Some(update) = updates.next().await {
1558 if let fedimint_mint_client::ReissueExternalNotesState::Failed(e) = update {
1559 return Err(PublicGatewayError::ReceiveEcashError {
1560 failure_reason: e.clone(),
1561 });
1562 }
1563 }
1564 }
1565
1566 Ok(ReceiveEcashResponse { amount })
1567 } else if let Ok(mint) = client.value().get_first_module::<MintV2ClientModule>() {
1568 let ecash: fedimint_mintv2_client::ECash =
1569 base32::decode_prefixed(FEDIMINT_PREFIX, &payload.notes).map_err(|e| {
1570 PublicGatewayError::ReceiveEcashError {
1571 failure_reason: format!("Expected ECash for MintV2 federation: {e}"),
1572 }
1573 })?;
1574 let amount = ecash.amount();
1575
1576 let operation_id = mint
1577 .receive(ecash, serde_json::Value::Null)
1578 .await
1579 .map_err(|e| PublicGatewayError::ReceiveEcashError {
1580 failure_reason: e.to_string(),
1581 })?;
1582
1583 if payload.wait {
1584 let final_state = mint
1585 .await_final_receive_operation_state(operation_id)
1586 .await
1587 .map_err(|e| PublicGatewayError::ReceiveEcashError {
1588 failure_reason: e.to_string(),
1589 })?;
1590 match final_state {
1591 fedimint_mintv2_client::FinalReceiveOperationState::Success => {}
1592 fedimint_mintv2_client::FinalReceiveOperationState::Rejected => {
1593 return Err(PublicGatewayError::ReceiveEcashError {
1594 failure_reason: "ECash receive was rejected".to_string(),
1595 });
1596 }
1597 }
1598 }
1599
1600 Ok(ReceiveEcashResponse { amount })
1601 } else {
1602 Err(PublicGatewayError::ReceiveEcashError {
1603 failure_reason: "No mint module found".to_string(),
1604 })
1605 }
1606 }
1607
1608 pub async fn handle_get_invoice_msg(
1611 &self,
1612 payload: GetInvoiceRequest,
1613 ) -> AdminResult<Option<GetInvoiceResponse>> {
1614 let lightning_context = self.get_lightning_context().await?;
1615 let invoice = lightning_context.lnrpc.get_invoice(payload).await?;
1616 Ok(invoice)
1617 }
1618
1619 pub async fn handle_withdraw_to_onchain_msg(
1622 &self,
1623 payload: WithdrawToOnchainPayload,
1624 ) -> AdminResult<WithdrawResponse> {
1625 let address = self.handle_get_ln_onchain_address_msg().await?;
1626 let withdraw = WithdrawPayload {
1627 address: address.into_unchecked(),
1628 federation_id: payload.federation_id,
1629 amount: payload.amount,
1630 quoted_fees: None,
1631 };
1632 self.handle_withdraw_msg(withdraw).await
1633 }
1634
1635 pub async fn handle_pegin_from_onchain_msg(
1638 &self,
1639 payload: PeginFromOnchainPayload,
1640 ) -> AdminResult<Txid> {
1641 let deposit = DepositAddressPayload {
1642 federation_id: payload.federation_id,
1643 };
1644 let address = self.handle_address_msg(deposit).await?;
1645 let send_onchain = SendOnchainRequest {
1646 address: address.into_unchecked(),
1647 amount: payload.amount,
1648 fee_rate_sats_per_vbyte: payload.fee_rate_sats_per_vbyte,
1649 };
1650 let txid = self.handle_send_onchain_msg(send_onchain).await?;
1651
1652 Ok(txid)
1653 }
1654
1655 async fn register_federations(
1657 &self,
1658 federations: &BTreeMap<FederationId, FederationConfig>,
1659 register_task_group: &TaskGroup,
1660 ) {
1661 if let Ok(lightning_context) = self.get_lightning_context().await {
1662 let route_hints = lightning_context
1663 .lnrpc
1664 .parsed_route_hints(self.num_route_hints)
1665 .await;
1666 if route_hints.is_empty() {
1667 warn!(target: LOG_GATEWAY, "Gateway did not retrieve any route hints, may reduce receive success rate.");
1668 }
1669
1670 for (federation_id, federation_config) in federations {
1671 let fed_manager = self.federation_manager.read().await;
1672 if let Some(client) = fed_manager.client(federation_id) {
1673 let client_arc = client.clone().into_value();
1674 let route_hints = route_hints.clone();
1675 let lightning_context = lightning_context.clone();
1676 let federation_config = federation_config.clone();
1677 let registrations =
1678 self.registrations.clone().into_values().collect::<Vec<_>>();
1679
1680 register_task_group.spawn_cancellable_silent(
1681 "register federation",
1682 async move {
1683 let Ok(gateway_client) =
1684 client_arc.get_first_module::<GatewayClientModule>()
1685 else {
1686 return;
1687 };
1688
1689 for registration in registrations {
1690 gateway_client
1691 .try_register_with_federation(
1692 route_hints.clone(),
1693 GW_ANNOUNCEMENT_TTL,
1694 federation_config.lightning_fee.into(),
1695 lightning_context.clone(),
1696 registration.endpoint_url,
1697 registration.keypair.public_key(),
1698 )
1699 .await;
1700 }
1701 },
1702 );
1703 }
1704 }
1705 }
1706 }
1707
1708 pub async fn select_client(
1711 &self,
1712 federation_id: FederationId,
1713 ) -> std::result::Result<Spanned<fedimint_client::ClientHandleArc>, FederationNotConnected>
1714 {
1715 self.federation_manager
1716 .read()
1717 .await
1718 .client(&federation_id)
1719 .cloned()
1720 .ok_or(FederationNotConnected {
1721 federation_id_prefix: federation_id.to_prefix(),
1722 })
1723 }
1724
1725 async fn load_mnemonic(gateway_db: &Database) -> Option<Mnemonic> {
1726 let secret = Client::load_decodable_client_secret::<Vec<u8>>(gateway_db)
1727 .await
1728 .ok()?;
1729 Mnemonic::from_entropy(&secret).ok()
1730 }
1731
1732 async fn load_clients(&self) -> AdminResult<()> {
1736 if let GatewayState::NotConfigured { .. } = self.get_state().await {
1737 return Ok(());
1738 }
1739
1740 let mut federation_manager = self.federation_manager.write().await;
1741
1742 let configs = {
1743 let mut dbtx = self.gateway_db.begin_transaction_nc().await;
1744 dbtx.load_federation_configs().await
1745 };
1746
1747 if let Some(max_federation_index) = configs.values().map(|cfg| cfg.federation_index).max() {
1748 federation_manager.set_next_index(max_federation_index + 1);
1749 }
1750
1751 let mnemonic = Self::load_mnemonic(&self.gateway_db)
1752 .await
1753 .expect("mnemonic should be set");
1754
1755 for (federation_id, config) in configs {
1756 let federation_index = config.federation_index;
1757 match Box::pin(Spanned::try_new(
1758 info_span!(target: LOG_GATEWAY, "client", federation_id = %federation_id.clone()),
1759 self.client_builder
1760 .build(config, Arc::new(self.clone()), &mnemonic),
1761 ))
1762 .await
1763 {
1764 Ok(client) => {
1765 federation_manager.add_client(federation_index, client);
1766 }
1767 _ => {
1768 warn!(target: LOG_GATEWAY, federation_id = %federation_id, "Failed to load client");
1769 }
1770 }
1771 }
1772
1773 Ok(())
1774 }
1775
1776 fn register_clients_timer(&self) {
1782 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
1784 info!(target: LOG_GATEWAY, "Spawning register task...");
1785 let gateway = self.clone();
1786 let register_task_group = self.task_group.make_subgroup();
1787 self.task_group.spawn_cancellable("register clients", async move {
1788 loop {
1789 let gateway_state = gateway.get_state().await;
1790 if let GatewayState::Running { .. } = &gateway_state {
1791 let mut dbtx = gateway.gateway_db.begin_transaction_nc().await;
1792 let all_federations_configs = dbtx.load_federation_configs().await.into_iter().collect();
1793 gateway.register_federations(&all_federations_configs, ®ister_task_group).await;
1794 } else {
1795 const NOT_RUNNING_RETRY: Duration = Duration::from_secs(10);
1797 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");
1798 sleep(NOT_RUNNING_RETRY).await;
1799 continue;
1800 }
1801
1802 sleep(GW_ANNOUNCEMENT_TTL.mul_f32(0.85)).await;
1805 }
1806 });
1807 }
1808 }
1809
1810 async fn check_federation_network(
1813 client: &ClientHandleArc,
1814 network: Network,
1815 ) -> AdminResult<()> {
1816 let federation_id = client.federation_id();
1817 let config = client.config().await;
1818
1819 let lnv1_cfg = config
1820 .modules
1821 .values()
1822 .find(|m| LightningCommonInit::KIND == m.kind);
1823
1824 let lnv2_cfg = config
1825 .modules
1826 .values()
1827 .find(|m| fedimint_lnv2_common::LightningCommonInit::KIND == m.kind);
1828
1829 if lnv1_cfg.is_none() && lnv2_cfg.is_none() {
1831 return Err(AdminGatewayError::ClientCreationError(anyhow!(
1832 "Federation {federation_id} does not have any lightning module (LNv1 or LNv2)"
1833 )));
1834 }
1835
1836 if let Some(cfg) = lnv1_cfg {
1838 let ln_cfg: &LightningClientConfig = cfg.cast()?;
1839
1840 if ln_cfg.network.0 != network {
1841 crit!(
1842 target: LOG_GATEWAY,
1843 federation_id = %federation_id,
1844 network = %network,
1845 "Incorrect LNv1 network for federation",
1846 );
1847 return Err(AdminGatewayError::ClientCreationError(anyhow!(format!(
1848 "Unsupported LNv1 network {}",
1849 ln_cfg.network
1850 ))));
1851 }
1852 }
1853
1854 if let Some(cfg) = lnv2_cfg {
1856 let ln_cfg: &fedimint_lnv2_common::config::LightningClientConfig = cfg.cast()?;
1857
1858 if ln_cfg.network != network {
1859 crit!(
1860 target: LOG_GATEWAY,
1861 federation_id = %federation_id,
1862 network = %network,
1863 "Incorrect LNv2 network for federation",
1864 );
1865 return Err(AdminGatewayError::ClientCreationError(anyhow!(format!(
1866 "Unsupported LNv2 network {}",
1867 ln_cfg.network
1868 ))));
1869 }
1870 }
1871
1872 Ok(())
1873 }
1874
1875 pub async fn get_lightning_context(
1879 &self,
1880 ) -> std::result::Result<LightningContext, LightningRpcError> {
1881 match self.get_state().await {
1882 GatewayState::Running { lightning_context }
1883 | GatewayState::ShuttingDown { lightning_context } => Ok(lightning_context),
1884 _ => Err(LightningRpcError::FailedToConnect),
1885 }
1886 }
1887
1888 pub async fn unannounce_from_all_federations(&self) {
1891 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
1892 for registration in self.registrations.values() {
1893 self.federation_manager
1894 .read()
1895 .await
1896 .unannounce_from_all_federations(registration.keypair)
1897 .await;
1898 }
1899 }
1900 }
1901
1902 async fn create_lightning_client(
1903 &self,
1904 runtime: Arc<tokio::runtime::Runtime>,
1905 ) -> Box<dyn ILnRpcClient> {
1906 match self.lightning_mode.clone() {
1907 LightningMode::Lnd {
1908 lnd_rpc_addr,
1909 lnd_tls_cert,
1910 lnd_macaroon,
1911 lnd_time_pref,
1912 lnd_payment_timeout_secs,
1913 } => {
1914 let gateway_db = self.gateway_db.clone();
1919 let lnv2_filter: Lnv2HoldInvoiceFilter = Arc::new(move |hash| {
1920 let gateway_db = gateway_db.clone();
1921 Box::pin(async move {
1922 gateway_db
1923 .begin_transaction_nc()
1924 .await
1925 .load_registered_incoming_contract(PaymentImage::Hash(hash))
1926 .await
1927 .is_some()
1928 })
1929 });
1930
1931 Box::new(GatewayLndClient::new(
1932 lnd_rpc_addr,
1933 lnd_tls_cert,
1934 lnd_macaroon,
1935 lnd_time_pref,
1936 lnd_payment_timeout_secs,
1937 None,
1938 lnv2_filter,
1939 ))
1940 }
1941 LightningMode::Ldk {
1942 lightning_port,
1943 alias,
1944 } => {
1945 let mnemonic = Self::load_mnemonic(&self.gateway_db)
1946 .await
1947 .expect("mnemonic should be set");
1948 retry("create LDK Node", fibonacci_max_one_hour(), || async {
1952 ldk::GatewayLdkClient::new(
1953 &self.client_builder.data_dir().join(LDK_NODE_DB_FOLDER),
1954 self.chain_source.clone(),
1955 self.network,
1956 lightning_port,
1957 alias.clone(),
1958 mnemonic.clone(),
1959 runtime.clone(),
1960 )
1961 .map(Box::new)
1962 })
1963 .await
1964 .expect("Could not create LDK Node")
1965 }
1966 }
1967 }
1968}
1969
1970#[async_trait]
1971impl IAdminGateway for Gateway {
1972 type Error = AdminGatewayError;
1973
1974 async fn handle_get_info(&self) -> AdminResult<GatewayInfo> {
1977 let GatewayState::Running { lightning_context } = self.get_state().await else {
1978 return Ok(GatewayInfo {
1979 federations: vec![],
1980 federation_fake_scids: None,
1981 version_hash: fedimint_build_code_version_env!().to_string(),
1982 gateway_state: self.state.read().await.to_string(),
1983 lightning_info: LightningInfo::NotConnected,
1984 lightning_mode: self.lightning_mode.clone(),
1985 registrations: self
1986 .registrations
1987 .iter()
1988 .map(|(k, v)| (k.clone(), (v.endpoint_url.clone(), v.keypair.public_key())))
1989 .collect(),
1990 });
1991 };
1992
1993 let dbtx = self.gateway_db.begin_transaction_nc().await;
1994 let federations = self
1995 .federation_manager
1996 .read()
1997 .await
1998 .federation_info_all_federations(dbtx)
1999 .await;
2000
2001 let channels: BTreeMap<u64, FederationId> = federations
2002 .iter()
2003 .map(|federation_info| {
2004 (
2005 federation_info.config.federation_index,
2006 federation_info.federation_id,
2007 )
2008 })
2009 .collect();
2010
2011 let lightning_info = lightning_context.lnrpc.parsed_node_info().await;
2012
2013 Ok(GatewayInfo {
2014 federations,
2015 federation_fake_scids: Some(channels),
2016 version_hash: fedimint_build_code_version_env!().to_string(),
2017 gateway_state: self.state.read().await.to_string(),
2018 lightning_info,
2019 lightning_mode: self.lightning_mode.clone(),
2020 registrations: self
2021 .registrations
2022 .iter()
2023 .map(|(k, v)| (k.clone(), (v.endpoint_url.clone(), v.keypair.public_key())))
2024 .collect(),
2025 })
2026 }
2027
2028 async fn handle_list_channels_msg(
2031 &self,
2032 ) -> AdminResult<Vec<fedimint_gateway_common::ChannelInfo>> {
2033 let context = self.get_lightning_context().await?;
2034 let response = context.lnrpc.list_channels().await?;
2035 Ok(response.channels)
2036 }
2037
2038 async fn handle_payment_summary_msg(
2041 &self,
2042 PaymentSummaryPayload {
2043 start_millis,
2044 end_millis,
2045 }: PaymentSummaryPayload,
2046 ) -> AdminResult<PaymentSummaryResponse> {
2047 let federation_manager = self.federation_manager.read().await;
2048 let fed_configs = federation_manager.get_all_federation_configs().await;
2049 let federation_ids = fed_configs.keys().collect::<Vec<_>>();
2050 let start = UNIX_EPOCH + Duration::from_millis(start_millis);
2051 let end = UNIX_EPOCH + Duration::from_millis(end_millis);
2052
2053 if start > end {
2054 return Err(AdminGatewayError::Unexpected(anyhow!("Invalid time range")));
2055 }
2056
2057 let mut outgoing = StructuredPaymentEvents::default();
2058 let mut incoming = StructuredPaymentEvents::default();
2059 for fed_id in federation_ids {
2060 let client = federation_manager
2061 .client(fed_id)
2062 .expect("No client available")
2063 .value();
2064 let all_events = &get_events_for_duration(client, start, end).await;
2065
2066 let (mut lnv1_outgoing, mut lnv1_incoming) = compute_lnv1_stats(all_events);
2067 let (mut lnv2_outgoing, mut lnv2_incoming) = compute_lnv2_stats(all_events);
2068 outgoing.combine(&mut lnv1_outgoing);
2069 incoming.combine(&mut lnv1_incoming);
2070 outgoing.combine(&mut lnv2_outgoing);
2071 incoming.combine(&mut lnv2_incoming);
2072 }
2073
2074 Ok(PaymentSummaryResponse {
2075 outgoing: PaymentStats::compute(&outgoing),
2076 incoming: PaymentStats::compute(&incoming),
2077 })
2078 }
2079
2080 async fn handle_leave_federation(
2085 &self,
2086 payload: LeaveFedPayload,
2087 ) -> AdminResult<FederationInfo> {
2088 let mut federation_manager = self.federation_manager.write().await;
2091 let mut dbtx = self.gateway_db.begin_transaction().await;
2092
2093 let federation_info = federation_manager
2094 .leave_federation(
2095 payload.federation_id,
2096 &mut dbtx.to_ref_nc(),
2097 self.registrations.values().collect(),
2098 )
2099 .await?;
2100
2101 dbtx.remove_federation_config(payload.federation_id).await;
2102 dbtx.commit_tx().await;
2103 Ok(federation_info)
2104 }
2105
2106 async fn handle_connect_federation(
2111 &self,
2112 payload: ConnectFedPayload,
2113 ) -> AdminResult<FederationInfo> {
2114 let GatewayState::Running { lightning_context } = self.get_state().await else {
2115 return Err(AdminGatewayError::Lightning(
2116 LightningRpcError::FailedToConnect,
2117 ));
2118 };
2119
2120 let invite_code = InviteCode::from_str(&payload.invite_code).map_err(|e| {
2121 AdminGatewayError::ClientCreationError(anyhow!(format!(
2122 "Invalid federation member string {e:?}"
2123 )))
2124 })?;
2125
2126 let federation_id = invite_code.federation_id();
2127
2128 let mut federation_manager = self.federation_manager.write().await;
2129
2130 if federation_manager.has_federation(federation_id) {
2132 return Err(AdminGatewayError::ClientCreationError(anyhow!(
2133 "Federation has already been registered"
2134 )));
2135 }
2136
2137 let federation_index = federation_manager.pop_next_index()?;
2140
2141 let federation_config = FederationConfig {
2142 invite_code,
2143 federation_index,
2144 lightning_fee: self.default_routing_fees,
2145 transaction_fee: self.default_transaction_fees,
2146 _connector: ConnectorType::Tcp,
2148 };
2149
2150 let mnemonic = Self::load_mnemonic(&self.gateway_db)
2151 .await
2152 .expect("mnemonic should be set");
2153 let recover = payload.recover.unwrap_or(false);
2154 if recover {
2155 self.client_builder
2156 .recover(federation_config.clone(), Arc::new(self.clone()), &mnemonic)
2157 .await?;
2158 }
2159
2160 let client = self
2161 .client_builder
2162 .build(federation_config.clone(), Arc::new(self.clone()), &mnemonic)
2163 .await?;
2164
2165 if recover {
2166 client.wait_for_all_active_state_machines().await?;
2167 }
2168
2169 let federation_info = FederationInfo {
2172 federation_id,
2173 federation_name: federation_manager.federation_name(&client).await,
2174 balance_msat: client.get_balance_for_btc().await.unwrap_or_else(|err| {
2175 warn!(
2176 target: LOG_GATEWAY,
2177 err = %err.fmt_compact_anyhow(),
2178 %federation_id,
2179 "Balance not immediately available after joining/recovering."
2180 );
2181 Amount::default()
2182 }),
2183 config: federation_config.clone(),
2184 last_backup_time: None,
2185 };
2186
2187 Self::check_federation_network(&client, self.network).await?;
2188 if matches!(self.lightning_mode, LightningMode::Lnd { .. })
2189 && let Ok(lnv1) = client.get_first_module::<GatewayClientModule>()
2190 {
2191 for registration in self.registrations.values() {
2192 lnv1.try_register_with_federation(
2193 Vec::new(),
2195 GW_ANNOUNCEMENT_TTL,
2196 federation_config.lightning_fee.into(),
2197 lightning_context.clone(),
2198 registration.endpoint_url.clone(),
2199 registration.keypair.public_key(),
2200 )
2201 .await;
2202 }
2203 }
2204
2205 federation_manager.add_client(
2207 federation_index,
2208 Spanned::new(
2209 info_span!(target: LOG_GATEWAY, "client", federation_id=%federation_id.clone()),
2210 async { client },
2211 )
2212 .await,
2213 );
2214
2215 let mut dbtx = self.gateway_db.begin_transaction().await;
2216 dbtx.save_federation_config(&federation_config).await;
2217 dbtx.save_federation_backup_record(federation_id, None)
2218 .await;
2219 dbtx.commit_tx().await;
2220 debug!(
2221 target: LOG_GATEWAY,
2222 federation_id = %federation_id,
2223 federation_index = %federation_index,
2224 "Federation connected"
2225 );
2226
2227 Ok(federation_info)
2228 }
2229
2230 async fn handle_set_fees_msg(
2233 &self,
2234 SetFeesPayload {
2235 federation_id,
2236 lightning_base,
2237 lightning_parts_per_million,
2238 transaction_base,
2239 transaction_parts_per_million,
2240 }: SetFeesPayload,
2241 ) -> AdminResult<()> {
2242 let mut dbtx = self.gateway_db.begin_transaction().await;
2243 let mut fed_configs = if let Some(fed_id) = federation_id {
2244 dbtx.load_federation_configs()
2245 .await
2246 .into_iter()
2247 .filter(|(id, _)| *id == fed_id)
2248 .collect::<BTreeMap<_, _>>()
2249 } else {
2250 dbtx.load_federation_configs().await
2251 };
2252
2253 let federation_manager = self.federation_manager.read().await;
2254
2255 for (federation_id, config) in &mut fed_configs {
2256 let mut lightning_fee = config.lightning_fee;
2257 if let Some(lightning_base) = lightning_base {
2258 lightning_fee.base = lightning_base;
2259 }
2260
2261 if let Some(lightning_ppm) = lightning_parts_per_million {
2262 lightning_fee.parts_per_million = lightning_ppm;
2263 }
2264
2265 let mut transaction_fee = config.transaction_fee;
2266 if let Some(transaction_base) = transaction_base {
2267 transaction_fee.base = transaction_base;
2268 }
2269
2270 if let Some(transaction_ppm) = transaction_parts_per_million {
2271 transaction_fee.parts_per_million = transaction_ppm;
2272 }
2273
2274 let client =
2275 federation_manager
2276 .client(federation_id)
2277 .ok_or(FederationNotConnected {
2278 federation_id_prefix: federation_id.to_prefix(),
2279 })?;
2280 let client_config = client.value().config().await;
2281 let contains_lnv2 = client_config
2282 .modules
2283 .values()
2284 .any(|m| fedimint_lnv2_common::LightningCommonInit::KIND == m.kind);
2285
2286 let send_fees = lightning_fee + transaction_fee;
2288 if contains_lnv2 && send_fees.gt(&PaymentFee::SEND_FEE_LIMIT) {
2289 return Err(AdminGatewayError::GatewayConfigurationError(format!(
2290 "Total Send fees exceeded {}",
2291 PaymentFee::SEND_FEE_LIMIT
2292 )));
2293 }
2294
2295 if contains_lnv2 && transaction_fee.gt(&PaymentFee::RECEIVE_FEE_LIMIT) {
2297 return Err(AdminGatewayError::GatewayConfigurationError(format!(
2298 "Transaction fees exceeded RECEIVE LIMIT {}",
2299 PaymentFee::RECEIVE_FEE_LIMIT
2300 )));
2301 }
2302
2303 config.lightning_fee = lightning_fee;
2304 config.transaction_fee = transaction_fee;
2305 dbtx.save_federation_config(config).await;
2306 }
2307
2308 dbtx.commit_tx().await;
2309
2310 if matches!(self.lightning_mode, LightningMode::Lnd { .. }) {
2311 let register_task_group = TaskGroup::new();
2312
2313 self.register_federations(&fed_configs, ®ister_task_group)
2314 .await;
2315 }
2316
2317 Ok(())
2318 }
2319
2320 async fn handle_mnemonic_msg(&self) -> AdminResult<MnemonicResponse> {
2324 let mnemonic = Self::load_mnemonic(&self.gateway_db)
2325 .await
2326 .expect("mnemonic should be set");
2327 let words = mnemonic
2328 .words()
2329 .map(std::string::ToString::to_string)
2330 .collect::<Vec<_>>();
2331 let all_federations = self
2332 .federation_manager
2333 .read()
2334 .await
2335 .get_all_federation_configs()
2336 .await
2337 .keys()
2338 .copied()
2339 .collect::<BTreeSet<_>>();
2340 let legacy_federations = self.client_builder.legacy_federations(all_federations);
2341 let mnemonic_response = MnemonicResponse {
2342 mnemonic: words,
2343 legacy_federations,
2344 };
2345 Ok(mnemonic_response)
2346 }
2347
2348 async fn handle_open_channel_msg(&self, payload: OpenChannelRequest) -> AdminResult<Txid> {
2351 info!(target: LOG_GATEWAY, pubkey = %payload.pubkey, host = %payload.host, amount = %payload.channel_size_sats, "Opening Lightning channel...");
2352 let context = self.get_lightning_context().await?;
2353 let res = context.lnrpc.open_channel(payload).await?;
2354 info!(target: LOG_GATEWAY, txid = %res.funding_txid, "Initiated channel open");
2355 Txid::from_str(&res.funding_txid).map_err(|e| {
2356 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2357 failure_reason: format!("Received invalid channel funding txid string {e}"),
2358 })
2359 })
2360 }
2361
2362 async fn handle_connect_peer_msg(&self, payload: ConnectPeerRequest) -> AdminResult<()> {
2365 info!(
2366 target: LOG_GATEWAY,
2367 pubkey = %payload.node_address.pubkey,
2368 host = %payload.node_address.host_with_port(),
2369 "Connecting to Lightning peer..."
2370 );
2371 let context = self.get_lightning_context().await?;
2372 context.lnrpc.connect_peer(payload).await?;
2373 info!(target: LOG_GATEWAY, "Connected to Lightning peer");
2374 Ok(())
2375 }
2376
2377 async fn handle_close_channels_with_peer_msg(
2380 &self,
2381 payload: CloseChannelsWithPeerRequest,
2382 ) -> AdminResult<CloseChannelsWithPeerResponse> {
2383 info!(target: LOG_GATEWAY, close_channel_request = %payload, "Closing lightning channel...");
2384 let context = self.get_lightning_context().await?;
2385 let response = context
2386 .lnrpc
2387 .close_channels_with_peer(payload.clone())
2388 .await?;
2389 info!(target: LOG_GATEWAY, close_channel_request = %payload, "Initiated channel closure");
2390 Ok(response)
2391 }
2392
2393 async fn handle_set_channel_fees_msg(&self, payload: SetChannelFeesRequest) -> AdminResult<()> {
2396 info!(
2397 target: LOG_GATEWAY,
2398 funding_outpoint = %payload.funding_outpoint,
2399 base_fee_msat = payload.base_fee_msat,
2400 parts_per_million = payload.parts_per_million,
2401 "Updating channel fees..."
2402 );
2403 let context = self.get_lightning_context().await?;
2404 context.lnrpc.set_channel_fees(payload).await?;
2405 Ok(())
2406 }
2407
2408 async fn handle_get_balances_msg(&self) -> AdminResult<GatewayBalances> {
2411 let dbtx = self.gateway_db.begin_transaction_nc().await;
2412 let federation_infos = self
2413 .federation_manager
2414 .read()
2415 .await
2416 .federation_info_all_federations(dbtx)
2417 .await;
2418
2419 let ecash_balances: Vec<FederationBalanceInfo> = federation_infos
2420 .iter()
2421 .map(|federation_info| FederationBalanceInfo {
2422 federation_id: federation_info.federation_id,
2423 ecash_balance_msats: Amount {
2424 msats: federation_info.balance_msat.msats,
2425 },
2426 })
2427 .collect();
2428
2429 let context = self.get_lightning_context().await?;
2430 let lightning_node_balances = context.lnrpc.get_balances().await?;
2431
2432 Ok(GatewayBalances {
2433 onchain_balance_sats: lightning_node_balances.onchain_balance_sats,
2434 lightning_balance_msats: lightning_node_balances.lightning_balance_msats,
2435 ecash_balances,
2436 inbound_lightning_liquidity_msats: lightning_node_balances
2437 .inbound_lightning_liquidity_msats,
2438 })
2439 }
2440
2441 async fn handle_send_onchain_msg(&self, payload: SendOnchainRequest) -> AdminResult<Txid> {
2443 let context = self.get_lightning_context().await?;
2444 let response = context.lnrpc.send_onchain(payload.clone()).await?;
2445 let txid =
2446 Txid::from_str(&response.txid).map_err(|e| AdminGatewayError::WithdrawError {
2447 failure_reason: format!("Failed to parse withdrawal TXID: {e}"),
2448 })?;
2449 info!(onchain_request = %payload, txid = %txid, "Sent onchain transaction");
2450 Ok(txid)
2451 }
2452
2453 async fn handle_get_ln_onchain_address_msg(&self) -> AdminResult<Address> {
2455 let context = self.get_lightning_context().await?;
2456 let response = context.lnrpc.get_ln_onchain_address().await?;
2457
2458 let address = Address::from_str(&response.address).map_err(|e| {
2459 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2460 failure_reason: e.to_string(),
2461 })
2462 })?;
2463
2464 address.require_network(self.network).map_err(|e| {
2465 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2466 failure_reason: e.to_string(),
2467 })
2468 })
2469 }
2470
2471 async fn handle_deposit_address_msg(
2472 &self,
2473 payload: DepositAddressPayload,
2474 ) -> AdminResult<Address> {
2475 self.handle_address_msg(payload).await
2476 }
2477
2478 async fn handle_receive_ecash_msg(
2479 &self,
2480 payload: ReceiveEcashPayload,
2481 ) -> AdminResult<ReceiveEcashResponse> {
2482 Self::handle_receive_ecash_msg(self, payload)
2483 .await
2484 .map_err(|e| AdminGatewayError::Unexpected(anyhow::anyhow!("{e}")))
2485 }
2486
2487 async fn handle_create_invoice_for_operator_msg(
2490 &self,
2491 payload: CreateInvoiceForOperatorPayload,
2492 ) -> AdminResult<Bolt11Invoice> {
2493 let GatewayState::Running { lightning_context } = self.get_state().await else {
2494 return Err(AdminGatewayError::Lightning(
2495 LightningRpcError::FailedToConnect,
2496 ));
2497 };
2498
2499 Bolt11Invoice::from_str(
2500 &lightning_context
2501 .lnrpc
2502 .create_invoice(CreateInvoiceRequest {
2503 payment_hash: None, amount_msat: payload.amount_msats,
2506 expiry_secs: payload.expiry_secs.unwrap_or(3600),
2507 description: payload.description.map(InvoiceDescription::Direct),
2508 })
2509 .await?
2510 .invoice,
2511 )
2512 .map_err(|e| {
2513 AdminGatewayError::Lightning(LightningRpcError::InvalidMetadata {
2514 failure_reason: e.to_string(),
2515 })
2516 })
2517 }
2518
2519 async fn handle_pay_invoice_for_operator_msg(
2522 &self,
2523 payload: PayInvoiceForOperatorPayload,
2524 ) -> AdminResult<Preimage> {
2525 const BASE_FEE: u64 = 50;
2527 const FEE_DENOMINATOR: u64 = 100;
2528 const MAX_DELAY: u64 = 1008;
2529
2530 let GatewayState::Running { lightning_context } = self.get_state().await else {
2531 return Err(AdminGatewayError::Lightning(
2532 LightningRpcError::FailedToConnect,
2533 ));
2534 };
2535
2536 let max_fee = BASE_FEE
2537 + payload
2538 .invoice
2539 .amount_milli_satoshis()
2540 .context("Invoice is missing amount")?
2541 .saturating_div(FEE_DENOMINATOR);
2542
2543 let res = lightning_context
2544 .lnrpc
2545 .pay(payload.invoice, MAX_DELAY, Amount::from_msats(max_fee))
2546 .await?;
2547 Ok(res.preimage)
2548 }
2549
2550 async fn handle_list_transactions_msg(
2552 &self,
2553 payload: ListTransactionsPayload,
2554 ) -> AdminResult<ListTransactionsResponse> {
2555 let lightning_context = self.get_lightning_context().await?;
2556 let response = lightning_context
2557 .lnrpc
2558 .list_transactions(payload.start_secs, payload.end_secs)
2559 .await?;
2560 Ok(response)
2561 }
2562
2563 async fn handle_spend_ecash_msg(
2565 &self,
2566 payload: SpendEcashPayload,
2567 ) -> AdminResult<SpendEcashResponse> {
2568 let client = self
2569 .select_client(payload.federation_id)
2570 .await?
2571 .into_value();
2572
2573 if let Ok(mint_module) = client.get_first_module::<MintClientModule>() {
2574 let notes = mint_module.send_oob_notes(payload.amount, ()).await?;
2575 debug!(target: LOG_GATEWAY, ?notes, "Spend ecash notes");
2576 Ok(SpendEcashResponse {
2577 notes: notes.to_string(),
2578 })
2579 } else if let Ok(mint_module) = client.get_first_module::<MintV2ClientModule>() {
2580 let (_, ecash) = mint_module
2581 .send(payload.amount, serde_json::Value::Null, true)
2582 .await
2583 .map_err(|e| AdminGatewayError::Unexpected(e.into()))?;
2584
2585 Ok(SpendEcashResponse {
2586 notes: base32::encode_prefixed(FEDIMINT_PREFIX, &ecash),
2587 })
2588 } else {
2589 Err(AdminGatewayError::Unexpected(anyhow::anyhow!(
2590 "No mint module available"
2591 )))
2592 }
2593 }
2594
2595 async fn handle_shutdown_msg(&self, task_group: TaskGroup) -> AdminResult<()> {
2598 let mut state_guard = self.state.write().await;
2600 if let GatewayState::Running { lightning_context } = state_guard.clone() {
2601 *state_guard = GatewayState::ShuttingDown { lightning_context };
2602
2603 self.federation_manager
2604 .read()
2605 .await
2606 .wait_for_incoming_payments()
2607 .await?;
2608 }
2609
2610 let tg = task_group.clone();
2611 tg.spawn("Kill Gateway", |_task_handle| async {
2612 if let Err(err) = task_group.shutdown_join_all(Duration::from_mins(3)).await {
2613 warn!(target: LOG_GATEWAY, err = %err.fmt_compact_anyhow(), "Error shutting down gateway");
2614 }
2615 });
2616 Ok(())
2617 }
2618
2619 fn get_task_group(&self) -> TaskGroup {
2620 self.task_group.clone()
2621 }
2622
2623 async fn handle_withdraw_msg(&self, payload: WithdrawPayload) -> AdminResult<WithdrawResponse> {
2626 let WithdrawPayload {
2627 amount,
2628 address,
2629 federation_id,
2630 quoted_fees,
2631 } = payload;
2632
2633 let address_network = get_network_for_address(&address);
2634 let gateway_network = self.network;
2635 let Ok(address) = address.require_network(gateway_network) else {
2636 return Err(AdminGatewayError::WithdrawError {
2637 failure_reason: format!(
2638 "Gateway is running on network {gateway_network}, but provided withdraw address is for network {address_network}"
2639 ),
2640 });
2641 };
2642
2643 let client = self.select_client(federation_id).await?;
2644
2645 if let Ok(wallet_module) = client
2646 .value()
2647 .get_first_module::<fedimint_walletv2_client::WalletClientModule>()
2648 {
2649 return withdraw_v2(client.value(), &wallet_module, &address, amount).await;
2650 }
2651
2652 let wallet_module = client.value().get_first_module::<WalletClientModule>()?;
2653
2654 let (withdraw_amount, fees) = match quoted_fees {
2657 Some(fees) => {
2659 let amt = match amount {
2660 BitcoinAmountOrAll::Amount(a) => a,
2661 BitcoinAmountOrAll::All => {
2662 return Err(AdminGatewayError::WithdrawError {
2664 failure_reason:
2665 "Cannot use 'all' with quoted fees - amount must be resolved first"
2666 .to_string(),
2667 });
2668 }
2669 };
2670 (amt, fees)
2671 }
2672 None => match amount {
2674 BitcoinAmountOrAll::All => {
2677 let balance = bitcoin::Amount::from_sat(
2678 client
2679 .value()
2680 .get_balance_for_btc()
2681 .await
2682 .map_err(|err| {
2683 AdminGatewayError::Unexpected(anyhow!(
2684 "Balance not available: {}",
2685 err.fmt_compact_anyhow()
2686 ))
2687 })?
2688 .msats
2689 / 1000,
2690 );
2691 let fees = wallet_module.get_withdraw_fees(&address, balance).await?;
2692 let withdraw_amount = balance.checked_sub(fees.amount());
2693 if withdraw_amount.is_none() {
2694 return Err(AdminGatewayError::WithdrawError {
2695 failure_reason: format!(
2696 "Insufficient funds. Balance: {balance} Fees: {fees:?}"
2697 ),
2698 });
2699 }
2700 (withdraw_amount.expect("checked above"), fees)
2701 }
2702 BitcoinAmountOrAll::Amount(amount) => (
2703 amount,
2704 wallet_module.get_withdraw_fees(&address, amount).await?,
2705 ),
2706 },
2707 };
2708
2709 let operation_id = wallet_module
2710 .withdraw(&address, withdraw_amount, fees, ())
2711 .await?;
2712 let mut updates = wallet_module
2713 .subscribe_withdraw_updates(operation_id)
2714 .await?
2715 .into_stream();
2716
2717 while let Some(update) = updates.next().await {
2718 match update {
2719 WithdrawState::Succeeded(txid) => {
2720 info!(target: LOG_GATEWAY, amount = %withdraw_amount, address = %address, "Sent funds");
2721 return Ok(WithdrawResponse { txid, fees });
2722 }
2723 WithdrawState::Failed(e) => {
2724 return Err(AdminGatewayError::WithdrawError { failure_reason: e });
2725 }
2726 WithdrawState::Created => {}
2727 }
2728 }
2729
2730 Err(AdminGatewayError::WithdrawError {
2731 failure_reason: "Ran out of state updates while withdrawing".to_string(),
2732 })
2733 }
2734
2735 async fn handle_withdraw_preview_msg(
2738 &self,
2739 payload: WithdrawPreviewPayload,
2740 ) -> AdminResult<WithdrawPreviewResponse> {
2741 let gateway_network = self.network;
2742 let address_checked = payload
2743 .address
2744 .clone()
2745 .require_network(gateway_network)
2746 .map_err(|_| AdminGatewayError::WithdrawError {
2747 failure_reason: "Address network mismatch".to_string(),
2748 })?;
2749
2750 let client = self.select_client(payload.federation_id).await?;
2751
2752 let WithdrawDetails {
2753 amount,
2754 mint_fees,
2755 peg_out_fees,
2756 } = match payload.amount {
2757 BitcoinAmountOrAll::All => {
2758 calculate_max_withdrawable(client.value(), &address_checked).await?
2759 }
2760 BitcoinAmountOrAll::Amount(btc_amount) => {
2761 if let Ok(wallet_module) = client.value().get_first_module::<WalletClientModule>() {
2762 WithdrawDetails {
2763 amount: btc_amount.into(),
2764 mint_fees: None,
2765 peg_out_fees: wallet_module
2766 .get_withdraw_fees(&address_checked, btc_amount)
2767 .await?,
2768 }
2769 } else if let Ok(wallet_module) = client
2770 .value()
2771 .get_first_module::<fedimint_walletv2_client::WalletClientModule>(
2772 ) {
2773 let fee = wallet_module.send_fee().await.map_err(|e| {
2774 AdminGatewayError::WithdrawError {
2775 failure_reason: e.to_string(),
2776 }
2777 })?;
2778 WithdrawDetails {
2779 amount: btc_amount.into(),
2780 mint_fees: None,
2781 peg_out_fees: PegOutFees::from_amount(fee),
2782 }
2783 } else {
2784 return Err(AdminGatewayError::Unexpected(anyhow!(
2785 "No wallet module found"
2786 )));
2787 }
2788 }
2789 };
2790
2791 let total_cost = amount
2792 .checked_add(peg_out_fees.amount().into())
2793 .and_then(|a| a.checked_add(mint_fees.unwrap_or(Amount::ZERO)))
2794 .ok_or_else(|| AdminGatewayError::Unexpected(anyhow!("Total cost overflow")))?;
2795
2796 Ok(WithdrawPreviewResponse {
2797 withdraw_amount: amount,
2798 address: payload.address.assume_checked().to_string(),
2799 peg_out_fees,
2800 total_cost,
2801 mint_fees,
2802 })
2803 }
2804
2805 async fn handle_payment_log_msg(
2817 &self,
2818 PaymentLogPayload {
2819 end_position,
2820 pagination_size,
2821 federation_id,
2822 event_kinds,
2823 }: PaymentLogPayload,
2824 ) -> AdminResult<PaymentLogResponse> {
2825 const BATCH_SIZE: u64 = 10_000;
2826 let federation_manager = self.federation_manager.read().await;
2827 let client = federation_manager
2828 .client(&federation_id)
2829 .ok_or(FederationNotConnected {
2830 federation_id_prefix: federation_id.to_prefix(),
2831 })?
2832 .value();
2833
2834 let event_kinds = if event_kinds.is_empty() {
2838 ALL_GATEWAY_EVENTS.to_vec()
2839 } else {
2840 event_kinds
2841 };
2842
2843 let end_position = if let Some(position) = end_position {
2844 position
2845 } else {
2846 let mut dbtx = client.db().begin_transaction_nc().await;
2847 dbtx.get_next_event_log_id().await
2848 };
2849
2850 let mut start_position = end_position.saturating_sub(BATCH_SIZE);
2851
2852 let mut payment_log = Vec::new();
2853
2854 while payment_log.len() < pagination_size {
2855 let batch = client.get_event_log(Some(start_position), BATCH_SIZE).await;
2856 let mut filtered_batch = batch
2857 .into_iter()
2858 .filter(|e| e.id() <= end_position && event_kinds.contains(&e.as_raw().kind))
2859 .collect::<Vec<_>>();
2860 filtered_batch.reverse();
2861 payment_log.extend(filtered_batch);
2862
2863 start_position = start_position.saturating_sub(BATCH_SIZE);
2865
2866 if start_position == EventLogId::LOG_START {
2867 break;
2868 }
2869 }
2870
2871 payment_log.truncate(pagination_size);
2873
2874 Ok(PaymentLogResponse(payment_log))
2875 }
2876
2877 async fn handle_set_mnemonic_msg(&self, payload: SetMnemonicPayload) -> AdminResult<()> {
2880 let mut state_guard = self.state.write().await;
2885
2886 let GatewayState::NotConfigured { mnemonic_sender } = state_guard.clone() else {
2888 return Err(AdminGatewayError::MnemonicError(anyhow!(
2889 "Gateway is not is NotConfigured state"
2890 )));
2891 };
2892
2893 let mnemonic = if let Some(words) = payload.words {
2894 info!(target: LOG_GATEWAY, "Using user provided mnemonic");
2895 Mnemonic::parse_in_normalized(Language::English, words.as_str()).map_err(|e| {
2896 AdminGatewayError::MnemonicError(anyhow!(format!(
2897 "Seed phrase provided in environment was invalid {e:?}"
2898 )))
2899 })?
2900 } else {
2901 debug!(target: LOG_GATEWAY, "Generating mnemonic and writing entropy to client storage");
2902 Bip39RootSecretStrategy::<12>::random(&mut OsRng)
2903 };
2904
2905 Client::store_encodable_client_secret(&self.gateway_db, mnemonic.to_entropy())
2906 .await
2907 .map_err(AdminGatewayError::MnemonicError)?;
2908
2909 *state_guard = GatewayState::Disconnected;
2910 drop(state_guard);
2911
2912 let _ = mnemonic_sender.send(());
2914
2915 Ok(())
2916 }
2917
2918 async fn handle_create_offer_for_operator_msg(
2920 &self,
2921 payload: CreateOfferPayload,
2922 ) -> AdminResult<CreateOfferResponse> {
2923 let lightning_context = self.get_lightning_context().await?;
2924 let offer = lightning_context.lnrpc.create_offer(
2925 payload.amount,
2926 payload.description,
2927 payload.expiry_secs,
2928 payload.quantity,
2929 )?;
2930 Ok(CreateOfferResponse { offer })
2931 }
2932
2933 async fn handle_pay_offer_for_operator_msg(
2935 &self,
2936 payload: PayOfferPayload,
2937 ) -> AdminResult<PayOfferResponse> {
2938 let lightning_context = self.get_lightning_context().await?;
2939 let preimage = lightning_context
2940 .lnrpc
2941 .pay_offer(
2942 payload.offer,
2943 payload.quantity,
2944 payload.amount,
2945 payload.payer_note,
2946 )
2947 .await?;
2948 Ok(PayOfferResponse {
2949 preimage: preimage.to_string(),
2950 })
2951 }
2952
2953 async fn handle_export_invite_codes(
2956 &self,
2957 ) -> BTreeMap<FederationId, BTreeMap<PeerId, (String, InviteCode)>> {
2958 let fed_manager = self.federation_manager.read().await;
2959 fed_manager.all_invite_codes().await
2960 }
2961
2962 async fn handle_get_note_summary_msg(
2965 &self,
2966 federation_id: &FederationId,
2967 ) -> AdminResult<TieredCounts> {
2968 let fed_manager = self.federation_manager.read().await;
2969 fed_manager.get_note_summary(federation_id).await
2970 }
2971
2972 fn get_password_hash(&self) -> String {
2973 self.bcrypt_password_hash.clone()
2974 }
2975
2976 fn gatewayd_version(&self) -> String {
2977 let gatewayd_version = env!("CARGO_PKG_VERSION");
2978 gatewayd_version.to_string()
2979 }
2980
2981 async fn get_chain_source(&self) -> (ChainSource, Network) {
2982 (self.chain_source.clone(), self.network)
2983 }
2984
2985 fn lightning_mode(&self) -> LightningMode {
2986 self.lightning_mode.clone()
2987 }
2988
2989 async fn is_configured(&self) -> bool {
2990 !matches!(self.get_state().await, GatewayState::NotConfigured { .. })
2991 }
2992}
2993
2994impl Gateway {
2996 async fn public_key_v2(&self, federation_id: &FederationId) -> Option<PublicKey> {
3000 self.federation_manager
3001 .read()
3002 .await
3003 .client(federation_id)
3004 .map(|client| {
3005 client
3006 .value()
3007 .get_first_module::<GatewayClientModuleV2>()
3008 .expect("Must have client module")
3009 .keypair
3010 .public_key()
3011 })
3012 }
3013
3014 pub async fn routing_info_v2(
3017 &self,
3018 federation_id: &FederationId,
3019 ) -> Result<Option<RoutingInfo>> {
3020 let context = self.get_lightning_context().await?;
3021
3022 let mut dbtx = self.gateway_db.begin_transaction_nc().await;
3023 let fed_config = dbtx.load_federation_config(*federation_id).await.ok_or(
3024 PublicGatewayError::FederationNotConnected(FederationNotConnected {
3025 federation_id_prefix: federation_id.to_prefix(),
3026 }),
3027 )?;
3028
3029 let lightning_fee = fed_config.lightning_fee;
3030 let transaction_fee = fed_config.transaction_fee;
3031
3032 Ok(self
3033 .public_key_v2(federation_id)
3034 .await
3035 .map(|module_public_key| RoutingInfo {
3036 lightning_public_key: context.lightning_public_key,
3037 lightning_alias: Some(context.lightning_alias.clone()),
3038 module_public_key,
3039 send_fee_default: lightning_fee + transaction_fee,
3040 send_fee_minimum: transaction_fee,
3044 expiration_delta_default: 1440,
3045 expiration_delta_minimum: EXPIRATION_DELTA_MINIMUM_V2,
3046 receive_fee: transaction_fee,
3049 }))
3050 }
3051
3052 async fn send_payment_v2(
3055 &self,
3056 payload: SendPaymentPayload,
3057 ) -> Result<std::result::Result<[u8; 32], Signature>> {
3058 self.select_client(payload.federation_id)
3059 .await?
3060 .value()
3061 .get_first_module::<GatewayClientModuleV2>()
3062 .expect("Must have client module")
3063 .send_payment(payload)
3064 .await
3065 .map_err(LNv2Error::OutgoingPayment)
3066 .map_err(PublicGatewayError::LNv2)
3067 }
3068
3069 async fn create_bolt11_invoice_v2(
3074 &self,
3075 payload: CreateBolt11InvoicePayload,
3076 ) -> Result<Bolt11Invoice> {
3077 if !payload.contract.verify() {
3078 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3079 "The contract is invalid".to_string(),
3080 )));
3081 }
3082
3083 let payment_info = self.routing_info_v2(&payload.federation_id).await?.ok_or(
3084 LNv2Error::IncomingPayment(format!(
3085 "Federation {} does not exist",
3086 payload.federation_id
3087 )),
3088 )?;
3089
3090 if payload.contract.commitment.refund_pk != payment_info.module_public_key {
3091 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3092 "The incoming contract is keyed to another gateway".to_string(),
3093 )));
3094 }
3095
3096 let contract_amount = payment_info.receive_fee.subtract_from(payload.amount.msats);
3097
3098 if contract_amount == Amount::ZERO {
3099 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3100 "Zero amount incoming contracts are not supported".to_string(),
3101 )));
3102 }
3103
3104 if contract_amount != payload.contract.commitment.amount {
3105 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3106 "The contract amount does not pay the correct amount of fees".to_string(),
3107 )));
3108 }
3109
3110 if payload.contract.commitment.expiration_or_fee <= duration_since_epoch().as_secs() {
3111 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3112 "The contract has already expired".to_string(),
3113 )));
3114 }
3115
3116 let payment_hash = match payload.contract.commitment.payment_image {
3117 PaymentImage::Hash(payment_hash) => payment_hash,
3118 PaymentImage::Point(..) => {
3119 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3120 "PaymentImage is not a payment hash".to_string(),
3121 )));
3122 }
3123 };
3124
3125 let invoice = self
3126 .create_invoice_via_lnrpc_v2(
3127 payment_hash,
3128 payload.amount,
3129 payload.description.clone(),
3130 payload.expiry_secs,
3131 )
3132 .await?;
3133
3134 let mut dbtx = self.gateway_db.begin_transaction().await;
3135
3136 if dbtx
3137 .save_registered_incoming_contract(
3138 payload.federation_id,
3139 payload.amount,
3140 payload.contract,
3141 )
3142 .await
3143 .is_some()
3144 {
3145 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3146 "PaymentHash is already registered".to_string(),
3147 )));
3148 }
3149
3150 dbtx.commit_tx_result().await.map_err(|_| {
3151 PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3152 "Payment hash is already registered".to_string(),
3153 ))
3154 })?;
3155
3156 Ok(invoice)
3157 }
3158
3159 pub async fn create_invoice_via_lnrpc_v2(
3162 &self,
3163 payment_hash: sha256::Hash,
3164 amount: Amount,
3165 description: Bolt11InvoiceDescription,
3166 expiry_time: u32,
3167 ) -> std::result::Result<Bolt11Invoice, LightningRpcError> {
3168 let lnrpc = self.get_lightning_context().await?.lnrpc;
3169
3170 let response = match description {
3171 Bolt11InvoiceDescription::Direct(description) => {
3172 lnrpc
3173 .create_invoice(CreateInvoiceRequest {
3174 payment_hash: Some(payment_hash),
3175 amount_msat: amount.msats,
3176 expiry_secs: expiry_time,
3177 description: Some(InvoiceDescription::Direct(description)),
3178 })
3179 .await?
3180 }
3181 Bolt11InvoiceDescription::Hash(hash) => {
3182 lnrpc
3183 .create_invoice(CreateInvoiceRequest {
3184 payment_hash: Some(payment_hash),
3185 amount_msat: amount.msats,
3186 expiry_secs: expiry_time,
3187 description: Some(InvoiceDescription::Hash(hash)),
3188 })
3189 .await?
3190 }
3191 };
3192
3193 Bolt11Invoice::from_str(&response.invoice).map_err(|e| {
3194 LightningRpcError::FailedToGetInvoice {
3195 failure_reason: e.to_string(),
3196 }
3197 })
3198 }
3199
3200 pub async fn verify_bolt11_preimage_v2(
3201 &self,
3202 payment_hash: sha256::Hash,
3203 wait: bool,
3204 ) -> std::result::Result<VerifyResponse, String> {
3205 let registered_contract = self
3206 .gateway_db
3207 .begin_transaction_nc()
3208 .await
3209 .load_registered_incoming_contract(PaymentImage::Hash(payment_hash))
3210 .await
3211 .ok_or("Unknown payment hash".to_string())?;
3212
3213 let client = self
3214 .select_client(registered_contract.federation_id)
3215 .await
3216 .map_err(|_| "Not connected to federation".to_string())?
3217 .into_value();
3218
3219 let operation_id = OperationId::from_encodable(®istered_contract.contract);
3220
3221 if !(wait || client.operation_exists(operation_id).await) {
3222 return Ok(VerifyResponse {
3223 settled: false,
3224 preimage: None,
3225 });
3226 }
3227
3228 let state = client
3229 .get_first_module::<GatewayClientModuleV2>()
3230 .expect("Must have client module")
3231 .await_receive(operation_id)
3232 .await;
3233
3234 let preimage = match state {
3235 FinalReceiveState::Success(preimage) => Ok(preimage),
3236 FinalReceiveState::Failure => Err("Payment has failed".to_string()),
3237 FinalReceiveState::Refunded => Err("Payment has been refunded".to_string()),
3238 FinalReceiveState::Rejected => Err("Payment has been rejected".to_string()),
3239 }?;
3240
3241 Ok(VerifyResponse {
3242 settled: true,
3243 preimage: Some(preimage),
3244 })
3245 }
3246
3247 pub async fn get_registered_incoming_contract_and_client_v2(
3251 &self,
3252 payment_image: PaymentImage,
3253 amount_msats: u64,
3254 ) -> Result<(IncomingContract, ClientHandleArc)> {
3255 let registered_incoming_contract = self
3256 .gateway_db
3257 .begin_transaction_nc()
3258 .await
3259 .load_registered_incoming_contract(payment_image)
3260 .await
3261 .ok_or(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3262 "No corresponding decryption contract available".to_string(),
3263 )))?;
3264
3265 if registered_incoming_contract.incoming_amount_msats != amount_msats {
3266 return Err(PublicGatewayError::LNv2(LNv2Error::IncomingPayment(
3267 "The available decryption contract's amount is not equal to the requested amount"
3268 .to_string(),
3269 )));
3270 }
3271
3272 let client = self
3273 .select_client(registered_incoming_contract.federation_id)
3274 .await?
3275 .into_value();
3276
3277 Ok((registered_incoming_contract.contract, client))
3278 }
3279}
3280
3281#[async_trait]
3282impl IGatewayClientV2 for Gateway {
3283 async fn complete_htlc(&self, htlc_response: InterceptPaymentResponse) {
3284 loop {
3285 match self.get_lightning_context().await {
3286 Ok(lightning_context) => {
3287 match lightning_context
3288 .lnrpc
3289 .complete_htlc(htlc_response.clone())
3290 .await
3291 {
3292 Ok(..) => return,
3293 Err(err) => {
3294 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failure trying to complete payment");
3295 }
3296 }
3297 }
3298 Err(err) => {
3299 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failure trying to complete payment");
3300 }
3301 }
3302
3303 sleep(Duration::from_secs(5)).await;
3304 }
3305 }
3306
3307 async fn is_direct_swap(
3308 &self,
3309 invoice: &Bolt11Invoice,
3310 ) -> anyhow::Result<Option<(IncomingContract, ClientHandleArc)>> {
3311 let lightning_context = self.get_lightning_context().await?;
3312 if lightning_context.lightning_public_key == invoice.get_payee_pub_key() {
3313 let (contract, client) = self
3314 .get_registered_incoming_contract_and_client_v2(
3315 PaymentImage::Hash(*invoice.payment_hash()),
3316 invoice
3317 .amount_milli_satoshis()
3318 .expect("The amount invoice has been previously checked"),
3319 )
3320 .await?;
3321 Ok(Some((contract, client)))
3322 } else {
3323 Ok(None)
3324 }
3325 }
3326
3327 async fn pay(
3328 &self,
3329 invoice: Bolt11Invoice,
3330 max_delay: u64,
3331 max_fee: Amount,
3332 ) -> std::result::Result<[u8; 32], LightningRpcError> {
3333 let lightning_context = self.get_lightning_context().await?;
3334 lightning_context
3335 .lnrpc
3336 .pay(invoice, max_delay, max_fee)
3337 .await
3338 .map(|response| response.preimage.0)
3339 }
3340
3341 async fn min_contract_amount(
3342 &self,
3343 federation_id: &FederationId,
3344 amount: u64,
3345 ) -> anyhow::Result<Amount> {
3346 Ok(self
3347 .routing_info_v2(federation_id)
3348 .await?
3349 .ok_or(anyhow!("Routing Info not available"))?
3350 .send_fee_minimum
3351 .add_to(amount))
3352 }
3353
3354 async fn is_lnv1_invoice(&self, invoice: &Bolt11Invoice) -> Option<Spanned<ClientHandleArc>> {
3355 let rhints = invoice.route_hints();
3356 match rhints.first().and_then(|rh| rh.0.last()) {
3357 None => None,
3358 Some(hop) => match self.get_lightning_context().await {
3359 Ok(lightning_context) => {
3360 if hop.src_node_id != lightning_context.lightning_public_key {
3361 return None;
3362 }
3363
3364 self.federation_manager
3365 .read()
3366 .await
3367 .get_client_for_index(hop.short_channel_id)
3368 }
3369 Err(_) => None,
3370 },
3371 }
3372 }
3373
3374 async fn relay_lnv1_swap(
3375 &self,
3376 client: &ClientHandleArc,
3377 invoice: &Bolt11Invoice,
3378 ) -> anyhow::Result<FinalReceiveState> {
3379 let swap_params = SwapParameters {
3380 payment_hash: *invoice.payment_hash(),
3381 amount_msat: Amount::from_msats(
3382 invoice
3383 .amount_milli_satoshis()
3384 .ok_or(anyhow!("Amountless invoice not supported"))?,
3385 ),
3386 };
3387 let lnv1 = client
3388 .get_first_module::<GatewayClientModule>()
3389 .expect("No LNv1 module");
3390 let operation_id = lnv1.gateway_handle_direct_swap(swap_params).await?;
3391 let mut stream = lnv1
3392 .gateway_subscribe_ln_receive(operation_id)
3393 .await?
3394 .into_stream();
3395 let mut final_state = FinalReceiveState::Failure;
3396 while let Some(update) = stream.next().await {
3397 match update {
3398 GatewayExtReceiveStates::Funding => {}
3399 GatewayExtReceiveStates::FundingFailed { error: _ } => {
3400 final_state = FinalReceiveState::Rejected;
3401 }
3402 GatewayExtReceiveStates::Preimage(preimage) => {
3403 final_state = FinalReceiveState::Success(preimage.0);
3404 }
3405 GatewayExtReceiveStates::RefundError {
3406 error_message: _,
3407 error: _,
3408 } => {
3409 final_state = FinalReceiveState::Failure;
3410 }
3411 GatewayExtReceiveStates::RefundSuccess {
3412 out_points: _,
3413 error: _,
3414 } => {
3415 final_state = FinalReceiveState::Refunded;
3416 }
3417 }
3418 }
3419
3420 Ok(final_state)
3421 }
3422
3423 async fn claim_payment_image(
3424 &self,
3425 payment_image: &PaymentImage,
3426 operation_id: OperationId,
3427 ) -> bool {
3428 self.gateway_db
3432 .autocommit(
3433 |dbtx, _| {
3434 let payment_image = payment_image.clone();
3435 Box::pin(async move {
3436 let claimer = dbtx
3437 .claim_outgoing_payment_image(payment_image, operation_id)
3438 .await;
3439 Ok::<_, std::convert::Infallible>(claimer == operation_id)
3440 })
3441 },
3442 None,
3443 )
3444 .await
3445 .expect("Retries until the transaction commits")
3446 }
3447}
3448
3449#[async_trait]
3450impl IGatewayClientV1 for Gateway {
3451 async fn verify_preimage_authentication(
3452 &self,
3453 payment_hash: sha256::Hash,
3454 preimage_auth: sha256::Hash,
3455 contract: OutgoingContractAccount,
3456 ) -> std::result::Result<(), OutgoingPaymentError> {
3457 let mut dbtx = self.gateway_db.begin_transaction().await;
3458 if let Some(secret_hash) = dbtx.load_preimage_authentication(payment_hash).await {
3459 if secret_hash != preimage_auth {
3460 return Err(OutgoingPaymentError {
3461 error_type: OutgoingPaymentErrorType::InvalidInvoicePreimage,
3462 contract_id: contract.contract.contract_id(),
3463 contract: Some(contract),
3464 });
3465 }
3466 } else {
3467 dbtx.save_new_preimage_authentication(payment_hash, preimage_auth)
3470 .await;
3471 return dbtx
3472 .commit_tx_result()
3473 .await
3474 .map_err(|_| OutgoingPaymentError {
3475 error_type: OutgoingPaymentErrorType::InvoiceAlreadyPaid,
3476 contract_id: contract.contract.contract_id(),
3477 contract: Some(contract),
3478 });
3479 }
3480
3481 Ok(())
3482 }
3483
3484 async fn verify_pruned_invoice(&self, payment_data: PaymentData) -> anyhow::Result<()> {
3485 let lightning_context = self.get_lightning_context().await?;
3486
3487 if matches!(payment_data, PaymentData::PrunedInvoice { .. }) {
3488 ensure!(
3489 lightning_context.lnrpc.supports_private_payments(),
3490 "Private payments are not supported by the lightning node"
3491 );
3492 }
3493
3494 Ok(())
3495 }
3496
3497 async fn get_routing_fees(&self, federation_id: FederationId) -> Option<RoutingFees> {
3498 let mut gateway_dbtx = self.gateway_db.begin_transaction_nc().await;
3499 gateway_dbtx
3500 .load_federation_config(federation_id)
3501 .await
3502 .map(|c| c.lightning_fee.into())
3503 }
3504
3505 async fn get_client(&self, federation_id: &FederationId) -> Option<Spanned<ClientHandleArc>> {
3506 self.federation_manager
3507 .read()
3508 .await
3509 .client(federation_id)
3510 .cloned()
3511 }
3512
3513 async fn get_client_for_invoice(
3514 &self,
3515 payment_data: PaymentData,
3516 ) -> Option<Spanned<ClientHandleArc>> {
3517 let rhints = payment_data.route_hints();
3518 match rhints.first().and_then(|rh| rh.0.last()) {
3519 None => None,
3520 Some(hop) => match self.get_lightning_context().await {
3521 Ok(lightning_context) => {
3522 if hop.src_node_id != lightning_context.lightning_public_key {
3523 return None;
3524 }
3525
3526 self.federation_manager
3527 .read()
3528 .await
3529 .get_client_for_index(hop.short_channel_id)
3530 }
3531 Err(_) => None,
3532 },
3533 }
3534 }
3535
3536 async fn pay(
3537 &self,
3538 payment_data: PaymentData,
3539 max_delay: u64,
3540 max_fee: Amount,
3541 ) -> std::result::Result<PayInvoiceResponse, LightningRpcError> {
3542 let lightning_context = self.get_lightning_context().await?;
3543
3544 match payment_data {
3545 PaymentData::Invoice(invoice) => {
3546 lightning_context
3547 .lnrpc
3548 .pay(invoice, max_delay, max_fee)
3549 .await
3550 }
3551 PaymentData::PrunedInvoice(invoice) => {
3552 lightning_context
3553 .lnrpc
3554 .pay_private(invoice, max_delay, max_fee)
3555 .await
3556 }
3557 }
3558 }
3559
3560 async fn complete_htlc(
3561 &self,
3562 htlc: InterceptPaymentResponse,
3563 ) -> std::result::Result<(), LightningRpcError> {
3564 let lightning_context = loop {
3566 match self.get_lightning_context().await {
3567 Ok(lightning_context) => break lightning_context,
3568 Err(err) => {
3569 warn!(target: LOG_GATEWAY, err = %err.fmt_compact(), "Failure trying to complete payment");
3570 sleep(Duration::from_secs(5)).await;
3571 }
3572 }
3573 };
3574
3575 lightning_context.lnrpc.complete_htlc(htlc).await
3576 }
3577
3578 async fn is_lnv2_direct_swap(
3579 &self,
3580 payment_hash: sha256::Hash,
3581 amount: Amount,
3582 ) -> anyhow::Result<
3583 Option<(
3584 fedimint_lnv2_common::contracts::IncomingContract,
3585 ClientHandleArc,
3586 )>,
3587 > {
3588 let (contract, client) = self
3589 .get_registered_incoming_contract_and_client_v2(
3590 PaymentImage::Hash(payment_hash),
3591 amount.msats,
3592 )
3593 .await?;
3594 Ok(Some((contract, client)))
3595 }
3596}