1#![deny(clippy::pedantic)]
2#![allow(clippy::cast_possible_wrap)]
3#![allow(clippy::module_name_repetitions)]
4#![allow(clippy::must_use_candidate)]
5#![allow(clippy::too_many_lines)]
6
7pub mod db;
8use std::collections::{BTreeMap, BTreeSet};
9use std::time::Duration;
10
11use anyhow::{Context, bail};
12use bitcoin_hashes::{Hash as BitcoinHash, sha256};
13use fedimint_api_client::api::{DynModuleApi, FederationApiExt};
14use fedimint_core::config::{
15 ServerModuleConfig, ServerModuleConsensusConfig, TypedServerModuleConfig,
16 TypedServerModuleConsensusConfig,
17};
18use fedimint_core::core::ModuleInstanceId;
19use fedimint_core::db::{DatabaseTransaction, DatabaseValue, IDatabaseTransactionOpsCoreTyped};
20use fedimint_core::encoding::Encodable;
21use fedimint_core::encoding::btc::NetworkLegacyEncodingWrapper;
22use fedimint_core::envs::{FM_ENABLE_MODULE_LNV1_ENV, is_env_var_set_opt, next_poll_delay};
23use fedimint_core::module::audit::Audit;
24use fedimint_core::module::{
25 Amounts, ApiEndpoint, ApiEndpointContext, ApiError, ApiRequestErased, ApiVersion,
26 CoreConsensusVersion, InputMeta, ModuleConsensusVersion, ModuleInit, TransactionItemAmounts,
27 public_api_endpoint,
28};
29use fedimint_core::secp256k1::{Message, PublicKey, SECP256K1};
30use fedimint_core::task::{TaskGroup, sleep};
31use fedimint_core::util::FmtCompact;
32use fedimint_core::{
33 Amount, InPoint, NumPeers, NumPeersExt, OutPoint, PeerId, apply, async_trait_maybe_send,
34 push_db_pair_items,
35};
36pub use fedimint_ln_common as common;
37use fedimint_ln_common::config::{
38 FeeConsensus, LightningClientConfig, LightningConfig, LightningConfigConsensus,
39 LightningConfigPrivate,
40};
41use fedimint_ln_common::contracts::incoming::{IncomingContractAccount, IncomingContractOffer};
42use fedimint_ln_common::contracts::{
43 Contract, ContractId, ContractOutcome, DecryptedPreimage, DecryptedPreimageStatus,
44 EncryptedPreimage, FundedContract, IdentifiableContract, Preimage, PreimageDecryptionShare,
45 PreimageKey,
46};
47use fedimint_ln_common::federation_endpoint_constants::{
48 ACCOUNT_ENDPOINT, AWAIT_ACCOUNT_ENDPOINT, AWAIT_BLOCK_HEIGHT_ENDPOINT, AWAIT_OFFER_ENDPOINT,
49 AWAIT_OUTGOING_CONTRACT_CANCELLED_ENDPOINT, AWAIT_PREIMAGE_DECRYPTION, BLOCK_COUNT_ENDPOINT,
50 GET_DECRYPTED_PREIMAGE_STATUS, LIST_GATEWAYS_ENDPOINT, MODULE_CONSENSUS_VERSION_ENDPOINT,
51 OFFER_ENDPOINT, REGISTER_GATEWAY_ENDPOINT, REMOVE_GATEWAY_CHALLENGE_ENDPOINT,
52 REMOVE_GATEWAY_ENDPOINT, SUPPORTED_MODULE_CONSENSUS_VERSION_ENDPOINT,
53};
54use fedimint_ln_common::{
55 CONTRACT_FUNDED_ONCE_MODULE_CONSENSUS_VERSION, ContractAccount, LightningCommonInit,
56 LightningConsensusItem, LightningGatewayAnnouncement, LightningGatewayRegistration,
57 LightningInput, LightningInputError, LightningModuleTypes, LightningOutput,
58 LightningOutputError, LightningOutputOutcome, LightningOutputOutcomeV0, LightningOutputV0,
59 MODULE_CONSENSUS_VERSION, RemoveGatewayRequest, create_gateway_registration_message,
60 create_gateway_remove_message,
61};
62use fedimint_logging::LOG_MODULE_LN;
63use fedimint_server_core::bitcoin_rpc::ServerBitcoinRpcMonitor;
64use fedimint_server_core::config::PeerHandleOps;
65use fedimint_server_core::{
66 ConfigGenModuleArgs, EnvVarDoc, ServerModule, ServerModuleInit, ServerModuleInitArgs,
67};
68use futures::StreamExt;
69use futures::future::join_all;
70use metrics::{LN_CANCEL_OUTGOING_CONTRACTS, LN_FUNDED_CONTRACT_SATS, LN_INCOMING_OFFER};
71use rand::rngs::OsRng;
72use strum::IntoEnumIterator;
73use threshold_crypto::poly::Commitment;
74use threshold_crypto::serde_impl::SerdeSecret;
75use threshold_crypto::{PublicKeySet, SecretKeyShare};
76use tokio::sync::watch;
77use tracing::{debug, error, info, info_span, trace, warn};
78
79use crate::db::{
80 AgreedDecryptionShareContractIdPrefix, AgreedDecryptionShareKey,
81 AgreedDecryptionShareKeyPrefix, BlockCountVoteKey, BlockCountVotePrefix,
82 ConsensusVersionVoteKey, ConsensusVersionVotePrefix, ContractKey, ContractKeyPrefix,
83 ContractUpdateKey, ContractUpdateKeyPrefix, DbKeyPrefix, EncryptedPreimageIndexKey,
84 EncryptedPreimageIndexKeyPrefix, LightningAuditItemKey, LightningAuditItemKeyPrefix,
85 LightningGatewayKey, LightningGatewayKeyPrefix, OfferKey, OfferKeyPrefix,
86 ProposeDecryptionShareKey, ProposeDecryptionShareKeyPrefix,
87};
88
89mod metrics;
90
91#[derive(Debug, Clone)]
92pub struct LightningInit;
93
94impl ModuleInit for LightningInit {
95 type Common = LightningCommonInit;
96
97 async fn dump_database(
98 &self,
99 dbtx: &mut DatabaseTransaction<'_>,
100 prefix_names: Vec<String>,
101 ) -> Box<dyn Iterator<Item = (String, Box<dyn erased_serde::Serialize + Send>)> + '_> {
102 let mut lightning: BTreeMap<String, Box<dyn erased_serde::Serialize + Send>> =
103 BTreeMap::new();
104 let filtered_prefixes = DbKeyPrefix::iter().filter(|f| {
105 prefix_names.is_empty() || prefix_names.contains(&f.to_string().to_lowercase())
106 });
107 for table in filtered_prefixes {
108 match table {
109 DbKeyPrefix::AgreedDecryptionShare => {
110 push_db_pair_items!(
111 dbtx,
112 AgreedDecryptionShareKeyPrefix,
113 AgreedDecryptionShareKey,
114 PreimageDecryptionShare,
115 lightning,
116 "Accepted Decryption Shares"
117 );
118 }
119 DbKeyPrefix::Contract => {
120 push_db_pair_items!(
121 dbtx,
122 ContractKeyPrefix,
123 ContractKey,
124 ContractAccount,
125 lightning,
126 "Contracts"
127 );
128 }
129 DbKeyPrefix::ContractUpdate => {
130 push_db_pair_items!(
131 dbtx,
132 ContractUpdateKeyPrefix,
133 ContractUpdateKey,
134 LightningOutputOutcomeV0,
135 lightning,
136 "Contract Updates"
137 );
138 }
139 DbKeyPrefix::LightningGateway => {
140 push_db_pair_items!(
141 dbtx,
142 LightningGatewayKeyPrefix,
143 LightningGatewayKey,
144 LightningGatewayRegistration,
145 lightning,
146 "Lightning Gateways"
147 );
148 }
149 DbKeyPrefix::Offer => {
150 push_db_pair_items!(
151 dbtx,
152 OfferKeyPrefix,
153 OfferKey,
154 IncomingContractOffer,
155 lightning,
156 "Offers"
157 );
158 }
159 DbKeyPrefix::ProposeDecryptionShare => {
160 push_db_pair_items!(
161 dbtx,
162 ProposeDecryptionShareKeyPrefix,
163 ProposeDecryptionShareKey,
164 PreimageDecryptionShare,
165 lightning,
166 "Proposed Decryption Shares"
167 );
168 }
169 DbKeyPrefix::BlockCountVote => {
170 push_db_pair_items!(
171 dbtx,
172 BlockCountVotePrefix,
173 BlockCountVoteKey,
174 u64,
175 lightning,
176 "Block Count Votes"
177 );
178 }
179 DbKeyPrefix::EncryptedPreimageIndex => {
180 push_db_pair_items!(
181 dbtx,
182 EncryptedPreimageIndexKeyPrefix,
183 EncryptedPreimageIndexKey,
184 (),
185 lightning,
186 "Encrypted Preimage Hashes"
187 );
188 }
189 DbKeyPrefix::LightningAuditItem => {
190 push_db_pair_items!(
191 dbtx,
192 LightningAuditItemKeyPrefix,
193 LightningAuditItemKey,
194 Amount,
195 lightning,
196 "Lightning Audit Items"
197 );
198 }
199 DbKeyPrefix::ConsensusVersionVote => {
200 push_db_pair_items!(
201 dbtx,
202 ConsensusVersionVotePrefix,
203 ConsensusVersionVoteKey,
204 ModuleConsensusVersion,
205 lightning,
206 "Consensus Version Votes"
207 );
208 }
209 }
210 }
211
212 Box::new(lightning.into_iter())
213 }
214}
215
216#[apply(async_trait_maybe_send!)]
217impl ServerModuleInit for LightningInit {
218 type Module = Lightning;
219
220 fn versions(&self, _core: CoreConsensusVersion) -> &[ModuleConsensusVersion] {
221 &[MODULE_CONSENSUS_VERSION]
222 }
223
224 fn is_enabled_by_default(&self) -> bool {
225 is_env_var_set_opt(FM_ENABLE_MODULE_LNV1_ENV).unwrap_or(false)
226 }
227
228 fn get_documented_env_vars(&self) -> Vec<EnvVarDoc> {
229 vec![EnvVarDoc {
230 name: FM_ENABLE_MODULE_LNV1_ENV,
231 description: "Set to 1/true to enable the LNv1 Lightning module. Disabled by default.",
232 }]
233 }
234
235 async fn init(&self, args: &ServerModuleInitArgs<Self>) -> anyhow::Result<Self::Module> {
236 LN_CANCEL_OUTGOING_CONTRACTS.get();
238
239 let peer_supported_consensus_version =
240 Lightning::spawn_peer_supported_consensus_version_task(
241 args.module_api().clone(),
242 args.task_group(),
243 args.our_peer_id(),
244 );
245
246 Ok(Lightning {
247 cfg: args.cfg().to_typed()?,
248 our_peer_id: args.our_peer_id(),
249 num_peers: args.num_peers(),
250 peer_supported_consensus_version,
251 server_bitcoin_rpc_monitor: args.server_bitcoin_rpc_monitor(),
252 })
253 }
254
255 fn trusted_dealer_gen(
256 &self,
257 peers: &[PeerId],
258 args: &ConfigGenModuleArgs,
259 ) -> BTreeMap<PeerId, ServerModuleConfig> {
260 let sks = threshold_crypto::SecretKeySet::random(peers.to_num_peers().degree(), &mut OsRng);
261 let pks = sks.public_keys();
262
263 peers
264 .iter()
265 .map(|&peer| {
266 let sk = sks.secret_key_share(peer.to_usize());
267
268 (
269 peer,
270 LightningConfig {
271 consensus: LightningConfigConsensus {
272 threshold_pub_keys: pks.clone(),
273 fee_consensus: FeeConsensus::default(),
274 network: NetworkLegacyEncodingWrapper(args.network),
275 },
276 private: LightningConfigPrivate {
277 threshold_sec_key: threshold_crypto::serde_impl::SerdeSecret(sk),
278 },
279 }
280 .to_erased(),
281 )
282 })
283 .collect()
284 }
285
286 async fn distributed_gen(
287 &self,
288 peers: &(dyn PeerHandleOps + Send + Sync),
289 args: &ConfigGenModuleArgs,
290 ) -> anyhow::Result<ServerModuleConfig> {
291 let (polynomial, mut sks) = peers.run_dkg_g1().await?;
292
293 let server = LightningConfig {
294 consensus: LightningConfigConsensus {
295 threshold_pub_keys: PublicKeySet::from(Commitment::from(polynomial)),
296 fee_consensus: FeeConsensus::default(),
297 network: NetworkLegacyEncodingWrapper(args.network),
298 },
299 private: LightningConfigPrivate {
300 threshold_sec_key: SerdeSecret(SecretKeyShare::from_mut(&mut sks)),
301 },
302 };
303
304 Ok(server.to_erased())
305 }
306
307 fn validate_config(&self, identity: &PeerId, config: ServerModuleConfig) -> anyhow::Result<()> {
308 let config = config.to_typed::<LightningConfig>()?;
309 if config.private.threshold_sec_key.public_key_share()
310 != config
311 .consensus
312 .threshold_pub_keys
313 .public_key_share(identity.to_usize())
314 {
315 bail!("Lightning private key doesn't match pubkey share");
316 }
317 Ok(())
318 }
319
320 fn get_client_config(
321 &self,
322 config: &ServerModuleConsensusConfig,
323 ) -> anyhow::Result<LightningClientConfig> {
324 let config = LightningConfigConsensus::from_erased(config)?;
325 Ok(LightningClientConfig {
326 threshold_pub_key: config.threshold_pub_keys.public_key(),
327 fee_consensus: config.fee_consensus,
328 network: config.network,
329 })
330 }
331
332 fn used_db_prefixes(&self) -> Option<BTreeSet<u8>> {
333 Some(DbKeyPrefix::iter().map(|p| p as u8).collect())
334 }
335}
336#[derive(Debug)]
358pub struct Lightning {
359 cfg: LightningConfig,
360 our_peer_id: PeerId,
361 num_peers: NumPeers,
362 peer_supported_consensus_version: watch::Receiver<Option<ModuleConsensusVersion>>,
365 server_bitcoin_rpc_monitor: ServerBitcoinRpcMonitor,
366}
367
368#[apply(async_trait_maybe_send!)]
369impl ServerModule for Lightning {
370 type Common = LightningModuleTypes;
371 type Init = LightningInit;
372
373 async fn consensus_proposal(
374 &self,
375 dbtx: &mut DatabaseTransaction<'_>,
376 ) -> Vec<LightningConsensusItem> {
377 let mut items: Vec<LightningConsensusItem> = dbtx
378 .find_by_prefix(&ProposeDecryptionShareKeyPrefix)
379 .await
380 .map(|(ProposeDecryptionShareKey(contract_id), share)| {
381 LightningConsensusItem::DecryptPreimage(contract_id, share)
382 })
383 .collect()
384 .await;
385
386 if let Ok(block_count_vote) = self.get_block_count() {
387 trace!(target: LOG_MODULE_LN, ?block_count_vote, "Proposing block count");
388 items.push(LightningConsensusItem::BlockCount(block_count_vote));
389 }
390
391 let active_consensus_version = self.consensus_module_consensus_version(dbtx).await;
397
398 if let Some(supported_consensus_version) = *self.peer_supported_consensus_version.borrow()
399 && active_consensus_version < supported_consensus_version
402 {
403 items.push(LightningConsensusItem::ModuleConsensusVersion(
404 supported_consensus_version,
405 ));
406 }
407
408 items
409 }
410
411 async fn process_consensus_item<'a, 'b>(
412 &'a self,
413 dbtx: &mut DatabaseTransaction<'b>,
414 consensus_item: LightningConsensusItem,
415 peer_id: PeerId,
416 ) -> anyhow::Result<()> {
417 let span = info_span!("process decryption share", %peer_id);
418 let _guard = span.enter();
419 trace!(target: LOG_MODULE_LN, ?consensus_item, "Processing consensus item proposal");
420
421 match consensus_item {
422 LightningConsensusItem::DecryptPreimage(contract_id, share) => {
423 if dbtx
424 .get_value(&AgreedDecryptionShareKey(contract_id, peer_id))
425 .await
426 .is_some()
427 {
428 bail!("Already received a valid decryption share for this peer");
429 }
430
431 let account = dbtx
432 .get_value(&ContractKey(contract_id))
433 .await
434 .context("Contract account for this decryption share does not exist")?;
435
436 let (contract, out_point) = match account.contract {
437 FundedContract::Incoming(contract) => (contract.contract, contract.out_point),
438 FundedContract::Outgoing(..) => {
439 bail!("Contract account for this decryption share is outgoing");
440 }
441 };
442
443 if contract.decrypted_preimage != DecryptedPreimage::Pending {
444 bail!("Contract for this decryption share is not pending");
445 }
446
447 if !self.validate_decryption_share(peer_id, &share, &contract.encrypted_preimage) {
448 bail!("Decryption share is invalid");
449 }
450
451 dbtx.insert_new_entry(&AgreedDecryptionShareKey(contract_id, peer_id), &share)
453 .await;
454
455 let decryption_shares = dbtx
457 .find_by_prefix(&AgreedDecryptionShareContractIdPrefix(contract_id))
458 .await
459 .map(|(key, decryption_share)| (key.1, decryption_share))
460 .collect::<Vec<_>>()
461 .await;
462
463 if decryption_shares.len() < self.cfg.consensus.threshold() {
464 return Ok(());
465 }
466
467 debug!(target: LOG_MODULE_LN, "Beginning to decrypt preimage");
468
469 let Ok(preimage_vec) = self.cfg.consensus.threshold_pub_keys.decrypt(
470 decryption_shares
471 .iter()
472 .map(|(peer, share)| (peer.to_usize(), &share.0)),
473 &contract.encrypted_preimage.0,
474 ) else {
475 error!(target: LOG_MODULE_LN, contract_hash = %contract.hash, "Failed to decrypt preimage");
478 return Ok(());
479 };
480
481 dbtx.remove_entry(&ProposeDecryptionShareKey(contract_id))
483 .await;
484
485 dbtx.remove_by_prefix(&AgreedDecryptionShareContractIdPrefix(contract_id))
486 .await;
487
488 let decrypted_preimage = if preimage_vec.len() == 33
489 && contract.hash
490 == sha256::Hash::hash(&sha256::Hash::hash(&preimage_vec).to_byte_array())
491 {
492 let preimage = PreimageKey(
493 preimage_vec
494 .as_slice()
495 .try_into()
496 .expect("Invalid preimage length"),
497 );
498 if preimage.to_public_key().is_ok() {
499 DecryptedPreimage::Some(preimage)
500 } else {
501 DecryptedPreimage::Invalid
502 }
503 } else {
504 DecryptedPreimage::Invalid
505 };
506
507 debug!(target: LOG_MODULE_LN, ?decrypted_preimage);
508
509 let contract_db_key = ContractKey(contract_id);
512 let mut contract_account = dbtx
513 .get_value(&contract_db_key)
514 .await
515 .expect("checked before that it exists");
516 let incoming = match &mut contract_account.contract {
517 FundedContract::Incoming(incoming) => incoming,
518 FundedContract::Outgoing(_) => {
519 unreachable!("previously checked that it's an incoming contract")
520 }
521 };
522 incoming.contract.decrypted_preimage = decrypted_preimage.clone();
523 trace!(?contract_account, "Updating contract account");
524 dbtx.insert_entry(&contract_db_key, &contract_account).await;
525
526 let mut outcome = dbtx
528 .get_value(&ContractUpdateKey(out_point))
529 .await
530 .expect("outcome was created on funding");
531
532 let LightningOutputOutcomeV0::Contract {
533 outcome: ContractOutcome::Incoming(incoming_contract_outcome_preimage),
534 ..
535 } = &mut outcome
536 else {
537 panic!("We are expecting an incoming contract")
538 };
539 *incoming_contract_outcome_preimage = decrypted_preimage.clone();
540 dbtx.insert_entry(&ContractUpdateKey(out_point), &outcome)
541 .await;
542 }
543 LightningConsensusItem::BlockCount(block_count) => {
544 let current_vote = dbtx
545 .get_value(&BlockCountVoteKey(peer_id))
546 .await
547 .unwrap_or(0);
548
549 if block_count < current_vote {
550 bail!("Block count vote decreased");
551 }
552
553 if block_count == current_vote {
554 bail!("Block height vote is redundant");
555 }
556
557 dbtx.insert_entry(&BlockCountVoteKey(peer_id), &block_count)
558 .await;
559 }
560 LightningConsensusItem::ModuleConsensusVersion(module_consensus_version) => {
561 let current_vote = dbtx
562 .get_value(&ConsensusVersionVoteKey(peer_id))
563 .await
564 .unwrap_or(ModuleConsensusVersion::new(2, 0));
565
566 if module_consensus_version <= current_vote {
567 bail!("Module consensus version vote is redundant");
568 }
569
570 dbtx.insert_entry(&ConsensusVersionVoteKey(peer_id), &module_consensus_version)
571 .await;
572
573 assert!(
574 self.consensus_module_consensus_version(dbtx).await <= MODULE_CONSENSUS_VERSION,
575 "Lightning module does not support new consensus version, please upgrade the module"
576 );
577 }
578 LightningConsensusItem::Default { variant, .. } => {
579 bail!("Unknown lightning consensus item received, variant={variant}");
580 }
581 }
582
583 Ok(())
584 }
585
586 async fn process_input<'a, 'b, 'c>(
587 &'a self,
588 dbtx: &mut DatabaseTransaction<'c>,
589 input: &'b LightningInput,
590 _in_point: InPoint,
591 ) -> Result<InputMeta, LightningInputError> {
592 let input = input.ensure_v0_ref()?;
593
594 let mut account = dbtx
595 .get_value(&ContractKey(input.contract_id))
596 .await
597 .ok_or(LightningInputError::UnknownContract(input.contract_id))?;
598
599 if account.amount < input.amount {
600 return Err(LightningInputError::InsufficientFunds(
601 account.amount,
602 input.amount,
603 ));
604 }
605
606 let consensus_block_count = self.consensus_block_count(dbtx).await;
607
608 let pub_key = match &account.contract {
609 FundedContract::Outgoing(outgoing) => {
610 if u64::from(outgoing.timelock) + 1 > consensus_block_count && !outgoing.cancelled {
611 let preimage_hash = bitcoin_hashes::sha256::Hash::hash(
613 &input
614 .witness
615 .as_ref()
616 .ok_or(LightningInputError::MissingPreimage)?
617 .0,
618 );
619
620 if preimage_hash != outgoing.hash {
622 return Err(LightningInputError::InvalidPreimage);
623 }
624
625 outgoing.gateway_key
627 } else {
628 outgoing.user_key
630 }
631 }
632 FundedContract::Incoming(incoming) => match &incoming.contract.decrypted_preimage {
633 DecryptedPreimage::Pending => {
635 return Err(LightningInputError::ContractNotReady);
636 }
637 DecryptedPreimage::Some(preimage) => match preimage.to_public_key() {
639 Ok(pub_key) => pub_key,
640 Err(_) => return Err(LightningInputError::InvalidPreimage),
641 },
642 DecryptedPreimage::Invalid => incoming.contract.gateway_key,
644 },
645 };
646
647 account.amount -= input.amount;
648
649 dbtx.insert_entry(&ContractKey(input.contract_id), &account)
650 .await;
651
652 let audit_key = LightningAuditItemKey::from_funded_contract(&account.contract);
656 if account.amount.msats == 0 {
657 dbtx.remove_entry(&audit_key).await;
658 } else {
659 dbtx.insert_entry(&audit_key, &account.amount).await;
660 }
661
662 Ok(InputMeta {
663 amount: TransactionItemAmounts {
664 amounts: Amounts::new_bitcoin(input.amount),
665 fees: Amounts::new_bitcoin(self.cfg.consensus.fee_consensus.contract_input),
666 },
667 pub_key,
668 })
669 }
670
671 async fn process_output<'a, 'b>(
672 &'a self,
673 dbtx: &mut DatabaseTransaction<'b>,
674 output: &'a LightningOutput,
675 out_point: OutPoint,
676 ) -> Result<TransactionItemAmounts, LightningOutputError> {
677 let output = output.ensure_v0_ref()?;
678
679 match output {
680 LightningOutputV0::Contract(contract) => {
681 if self.is_contract_funded_once_active(dbtx).await
691 && dbtx
692 .get_value(&ContractKey(contract.contract.contract_id()))
693 .await
694 .is_some()
695 {
696 return Err(LightningOutputError::ContractAlreadyFunded(
697 contract.contract.contract_id(),
698 ));
699 }
700
701 if let Contract::Incoming(incoming) = &contract.contract {
703 if dbtx
708 .get_value(&ProposeDecryptionShareKey(incoming.contract_id()))
709 .await
710 .is_some()
711 {
712 return Err(LightningOutputError::ContractAlreadyFunded(
713 incoming.contract_id(),
714 ));
715 }
716
717 let offer = dbtx
718 .get_value(&OfferKey(incoming.hash))
719 .await
720 .ok_or(LightningOutputError::NoOffer(incoming.hash))?;
721
722 if contract.amount < offer.amount {
723 return Err(LightningOutputError::InsufficientIncomingFunding(
725 offer.amount,
726 contract.amount,
727 ));
728 }
729
730 if self.is_contract_funded_once_active(dbtx).await
738 && incoming.encrypted_preimage != offer.encrypted_preimage
739 {
740 return Err(LightningOutputError::EncryptedPreimageMismatch);
741 }
742
743 if self.is_contract_funded_once_active(dbtx).await
751 && incoming.decrypted_preimage != DecryptedPreimage::Pending
752 {
753 return Err(LightningOutputError::PreDecryptedIncomingContract);
754 }
755
756 if !incoming.encrypted_preimage.0.verify() {
761 return Err(LightningOutputError::InvalidEncryptedPreimage);
762 }
763 }
764
765 if contract.amount == Amount::ZERO {
766 return Err(LightningOutputError::ZeroOutput);
767 }
768
769 let contract_db_key = ContractKey(contract.contract.contract_id());
770
771 let updated_contract_account = dbtx.get_value(&contract_db_key).await.map_or_else(
772 || ContractAccount {
773 amount: contract.amount,
774 contract: contract.contract.clone().to_funded(out_point),
775 },
776 |mut value: ContractAccount| {
777 value.amount += contract.amount;
778 value
779 },
780 );
781
782 dbtx.insert_entry(
783 &LightningAuditItemKey::from_funded_contract(
784 &updated_contract_account.contract,
785 ),
786 &updated_contract_account.amount,
787 )
788 .await;
789
790 if dbtx
791 .insert_entry(&contract_db_key, &updated_contract_account)
792 .await
793 .is_none()
794 {
795 dbtx.on_commit(move || {
796 record_funded_contract_metric(&updated_contract_account);
797 });
798 }
799
800 dbtx.insert_new_entry(
801 &ContractUpdateKey(out_point),
802 &LightningOutputOutcomeV0::Contract {
803 id: contract.contract.contract_id(),
804 outcome: contract.contract.to_outcome(),
805 },
806 )
807 .await;
808
809 if let Contract::Incoming(incoming) = &contract.contract {
810 let offer = dbtx
811 .get_value(&OfferKey(incoming.hash))
812 .await
813 .expect("offer exists if output is valid");
814
815 let decryption_share = self
816 .cfg
817 .private
818 .threshold_sec_key
819 .decrypt_share(&incoming.encrypted_preimage.0)
820 .ok_or(LightningOutputError::InvalidEncryptedPreimage)?;
821
822 dbtx.insert_new_entry(
823 &ProposeDecryptionShareKey(contract.contract.contract_id()),
824 &PreimageDecryptionShare(decryption_share),
825 )
826 .await;
827
828 dbtx.remove_entry(&OfferKey(offer.hash)).await;
829 }
830
831 Ok(TransactionItemAmounts {
832 amounts: Amounts::new_bitcoin(contract.amount),
833 fees: Amounts::new_bitcoin(self.cfg.consensus.fee_consensus.contract_output),
834 })
835 }
836 LightningOutputV0::Offer(offer) => {
837 if self.is_contract_funded_once_active(dbtx).await
844 && dbtx
845 .get_value(&ContractKey(offer.contract_id()))
846 .await
847 .is_some()
848 {
849 return Err(LightningOutputError::OfferForFundedContract(
850 offer.contract_id(),
851 ));
852 }
853
854 if !offer.encrypted_preimage.0.verify() {
855 return Err(LightningOutputError::InvalidEncryptedPreimage);
856 }
857
858 if dbtx
860 .insert_entry(
861 &EncryptedPreimageIndexKey(offer.encrypted_preimage.consensus_hash()),
862 &(),
863 )
864 .await
865 .is_some()
866 {
867 return Err(LightningOutputError::DuplicateEncryptedPreimage);
868 }
869
870 dbtx.insert_new_entry(
871 &ContractUpdateKey(out_point),
872 &LightningOutputOutcomeV0::Offer { id: offer.id() },
873 )
874 .await;
875
876 if dbtx
878 .insert_entry(&OfferKey(offer.hash), &(*offer).clone())
879 .await
880 .is_some()
881 {
882 return Err(LightningOutputError::DuplicateEncryptedPreimage);
887 }
888
889 dbtx.on_commit(|| {
890 LN_INCOMING_OFFER.inc();
891 });
892
893 Ok(TransactionItemAmounts::ZERO)
894 }
895 LightningOutputV0::CancelOutgoing {
896 contract,
897 gateway_signature,
898 } => {
899 let contract_account = dbtx
900 .get_value(&ContractKey(*contract))
901 .await
902 .ok_or(LightningOutputError::UnknownContract(*contract))?;
903
904 let outgoing_contract = match &contract_account.contract {
905 FundedContract::Outgoing(contract) => contract,
906 FundedContract::Incoming(_) => {
907 return Err(LightningOutputError::NotOutgoingContract);
908 }
909 };
910
911 SECP256K1
912 .verify_schnorr(
913 gateway_signature,
914 &Message::from_digest(*outgoing_contract.cancellation_message().as_ref()),
915 &outgoing_contract.gateway_key.x_only_public_key().0,
916 )
917 .map_err(|_| LightningOutputError::InvalidCancellationSignature)?;
918
919 let updated_contract_account = {
920 let mut contract_account = dbtx
921 .get_value(&ContractKey(*contract))
922 .await
923 .expect("Contract exists if output is valid");
924
925 let outgoing_contract = match &mut contract_account.contract {
926 FundedContract::Outgoing(contract) => contract,
927 FundedContract::Incoming(_) => {
928 panic!("Contract type was checked in validate_output");
929 }
930 };
931
932 outgoing_contract.cancelled = true;
933
934 contract_account
935 };
936
937 dbtx.insert_entry(&ContractKey(*contract), &updated_contract_account)
938 .await;
939
940 dbtx.insert_new_entry(
941 &ContractUpdateKey(out_point),
942 &LightningOutputOutcomeV0::CancelOutgoingContract { id: *contract },
943 )
944 .await;
945
946 dbtx.on_commit(|| {
947 LN_CANCEL_OUTGOING_CONTRACTS.inc();
948 });
949
950 Ok(TransactionItemAmounts::ZERO)
951 }
952 }
953 }
954
955 async fn output_status(
956 &self,
957 dbtx: &mut DatabaseTransaction<'_>,
958 out_point: OutPoint,
959 ) -> Option<LightningOutputOutcome> {
960 dbtx.get_value(&ContractUpdateKey(out_point))
961 .await
962 .map(LightningOutputOutcome::V0)
963 }
964
965 #[doc(hidden)]
983 async fn verify_output_submission<'a, 'b>(
984 &'a self,
985 dbtx: &mut DatabaseTransaction<'b>,
986 output: &'a LightningOutput,
987 _out_point: OutPoint,
988 ) -> Result<(), LightningOutputError> {
989 match output.ensure_v0_ref()? {
990 LightningOutputV0::Contract(contract) => {
991 let contract_id = contract.contract.contract_id();
992
993 if dbtx.get_value(&ContractKey(contract_id)).await.is_some() {
994 return Err(LightningOutputError::ContractAlreadyFunded(contract_id));
995 }
996 }
997 LightningOutputV0::Offer(offer) => {
998 if dbtx
999 .get_value(&ContractKey(offer.contract_id()))
1000 .await
1001 .is_some()
1002 {
1003 return Err(LightningOutputError::OfferForFundedContract(
1004 offer.contract_id(),
1005 ));
1006 }
1007 }
1008 LightningOutputV0::CancelOutgoing { .. } => {}
1009 }
1010
1011 Ok(())
1012 }
1013
1014 async fn audit(
1015 &self,
1016 dbtx: &mut DatabaseTransaction<'_>,
1017 audit: &mut Audit,
1018 module_instance_id: ModuleInstanceId,
1019 ) {
1020 audit
1021 .add_items(
1022 dbtx,
1023 module_instance_id,
1024 &LightningAuditItemKeyPrefix,
1025 |_, v| -(v.msats as i64),
1028 )
1029 .await;
1030 }
1031
1032 fn api_endpoints(&self) -> Vec<ApiEndpoint<Self>> {
1033 vec![
1034 public_api_endpoint! {
1035 BLOCK_COUNT_ENDPOINT,
1036 ApiVersion::new(0, 0),
1037 async |module: &Lightning, context, _v: ()| -> Option<u64> {
1038 let db = context.db();
1039 let mut dbtx = db.begin_transaction_nc().await;
1040 Ok(Some(module.consensus_block_count(&mut dbtx).await))
1041 }
1042 },
1043 public_api_endpoint! {
1044 MODULE_CONSENSUS_VERSION_ENDPOINT,
1045 ApiVersion::new(0, 1),
1046 async |module: &Lightning, context, _params: ()| -> ModuleConsensusVersion {
1047 let db = context.db();
1048 let mut dbtx = db.begin_transaction_nc().await;
1049 Ok(module.consensus_module_consensus_version(&mut dbtx).await)
1050 }
1051 },
1052 public_api_endpoint! {
1053 SUPPORTED_MODULE_CONSENSUS_VERSION_ENDPOINT,
1054 ApiVersion::new(0, 1),
1055 async |_module: &Lightning, _context, _params: ()| -> ModuleConsensusVersion {
1056 Ok(MODULE_CONSENSUS_VERSION)
1057 }
1058 },
1059 public_api_endpoint! {
1060 ACCOUNT_ENDPOINT,
1061 ApiVersion::new(0, 0),
1062 async |module: &Lightning, context, contract_id: ContractId| -> Option<ContractAccount> {
1063 let db = context.db();
1064 let mut dbtx = db.begin_transaction_nc().await;
1065 Ok(module
1066 .get_contract_account(&mut dbtx, contract_id)
1067 .await)
1068 }
1069 },
1070 public_api_endpoint! {
1071 AWAIT_ACCOUNT_ENDPOINT,
1072 ApiVersion::new(0, 0),
1073 async |module: &Lightning, context, contract_id: ContractId| -> ContractAccount {
1074 Ok(module
1075 .wait_contract_account(context, contract_id)
1076 .await)
1077 }
1078 },
1079 public_api_endpoint! {
1080 AWAIT_BLOCK_HEIGHT_ENDPOINT,
1081 ApiVersion::new(0, 0),
1082 async |module: &Lightning, context, block_height: u64| -> () {
1083 let db = context.db();
1084 let mut dbtx = db.begin_transaction_nc().await;
1085 module.wait_block_height(block_height, &mut dbtx).await;
1086 Ok(())
1087 }
1088 },
1089 public_api_endpoint! {
1090 AWAIT_OUTGOING_CONTRACT_CANCELLED_ENDPOINT,
1091 ApiVersion::new(0, 0),
1092 async |module: &Lightning, context, contract_id: ContractId| -> ContractAccount {
1093 Ok(module.wait_outgoing_contract_account_cancelled(context, contract_id).await)
1094 }
1095 },
1096 public_api_endpoint! {
1097 GET_DECRYPTED_PREIMAGE_STATUS,
1098 ApiVersion::new(0, 0),
1099 async |module: &Lightning, context, contract_id: ContractId| -> (IncomingContractAccount, DecryptedPreimageStatus) {
1100 module.get_decrypted_preimage_status(context, contract_id).await
1101 }
1102 },
1103 public_api_endpoint! {
1104 AWAIT_PREIMAGE_DECRYPTION,
1105 ApiVersion::new(0, 0),
1106 async |module: &Lightning, context, contract_id: ContractId| -> (IncomingContractAccount, Option<Preimage>) {
1107 Ok(module.wait_preimage_decrypted(context, contract_id).await)
1108 }
1109 },
1110 public_api_endpoint! {
1111 OFFER_ENDPOINT,
1112 ApiVersion::new(0, 0),
1113 async |module: &Lightning, context, payment_hash: bitcoin_hashes::sha256::Hash| -> Option<IncomingContractOffer> {
1114 let db = context.db();
1115 let mut dbtx = db.begin_transaction_nc().await;
1116 Ok(module
1117 .get_offer(&mut dbtx, payment_hash)
1118 .await)
1119 }
1120 },
1121 public_api_endpoint! {
1122 AWAIT_OFFER_ENDPOINT,
1123 ApiVersion::new(0, 0),
1124 async |module: &Lightning, context, payment_hash: bitcoin_hashes::sha256::Hash| -> IncomingContractOffer {
1125 Ok(module
1126 .wait_offer(context, payment_hash)
1127 .await)
1128 }
1129 },
1130 public_api_endpoint! {
1131 LIST_GATEWAYS_ENDPOINT,
1132 ApiVersion::new(0, 0),
1133 async |module: &Lightning, context, _v: ()| -> Vec<LightningGatewayAnnouncement> {
1134 let db = context.db();
1135 let mut dbtx = db.begin_transaction_nc().await;
1136 Ok(module.list_gateways(&mut dbtx).await)
1137 }
1138 },
1139 public_api_endpoint! {
1140 REGISTER_GATEWAY_ENDPOINT,
1141 ApiVersion::new(0, 0),
1142 async |module: &Lightning, context, gateway: LightningGatewayAnnouncement| -> () {
1143 let db = context.db();
1144 let mut dbtx = db.begin_transaction().await;
1145 let gateway_id = gateway.info.gateway_id;
1146 module.register_gateway(&mut dbtx.to_ref_nc(), gateway).await.map_err(|err| {
1147 warn!(target: LOG_MODULE_LN, err = %err.fmt_compact(), %gateway_id, "Rejected gateway registration");
1148 ApiError::bad_request(err.to_string())
1149 })?;
1150 dbtx.commit_tx_result().await?;
1151 Ok(())
1152 }
1153 },
1154 public_api_endpoint! {
1155 REMOVE_GATEWAY_CHALLENGE_ENDPOINT,
1156 ApiVersion::new(0, 1),
1157 async |module: &Lightning, context, gateway_id: PublicKey| -> Option<sha256::Hash> {
1158 let db = context.db();
1159 let mut dbtx = db.begin_transaction_nc().await;
1160 Ok(module.get_gateway_remove_challenge(gateway_id, &mut dbtx).await)
1161 }
1162 },
1163 public_api_endpoint! {
1164 REMOVE_GATEWAY_ENDPOINT,
1165 ApiVersion::new(0, 1),
1166 async |module: &Lightning, context, remove_gateway_request: RemoveGatewayRequest| -> bool {
1167 let db = context.db();
1168 let mut dbtx = db.begin_transaction().await;
1169 let result = module.remove_gateway(remove_gateway_request.clone(), &mut dbtx.to_ref_nc()).await;
1170 match result {
1171 Ok(()) => {
1172 dbtx.commit_tx_result().await?;
1173 Ok(true)
1174 },
1175 Err(err) => {
1176 warn!(target: LOG_MODULE_LN, err = %err.fmt_compact(), gateway_id = %remove_gateway_request.gateway_id, "Unable to remove gateway registration");
1177 Ok(false)
1178 },
1179 }
1180 }
1181 },
1182 ]
1183 }
1184}
1185
1186impl Lightning {
1187 fn get_block_count(&self) -> anyhow::Result<u64> {
1188 self.server_bitcoin_rpc_monitor
1189 .status()
1190 .map(|status| status.block_count)
1191 .context("Block count not available yet")
1192 }
1193
1194 async fn consensus_block_count(&self, dbtx: &mut DatabaseTransaction<'_>) -> u64 {
1195 let peer_count = 3 * (self.cfg.consensus.threshold() / 2) + 1;
1196
1197 let mut counts = dbtx
1198 .find_by_prefix(&BlockCountVotePrefix)
1199 .await
1200 .map(|(.., count)| count)
1201 .collect::<Vec<_>>()
1202 .await;
1203
1204 assert!(counts.len() <= peer_count);
1205
1206 while counts.len() < peer_count {
1207 counts.push(0);
1208 }
1209
1210 counts.sort_unstable();
1211
1212 counts[peer_count / 2]
1213 }
1214
1215 async fn wait_block_height(&self, block_height: u64, dbtx: &mut DatabaseTransaction<'_>) {
1216 while block_height >= self.consensus_block_count(dbtx).await {
1217 sleep(Duration::from_secs(5)).await;
1218 }
1219 }
1220
1221 async fn consensus_module_consensus_version(
1222 &self,
1223 dbtx: &mut DatabaseTransaction<'_>,
1224 ) -> ModuleConsensusVersion {
1225 let mut versions = dbtx
1226 .find_by_prefix(&ConsensusVersionVotePrefix)
1227 .await
1228 .map(|entry| entry.1)
1229 .collect::<Vec<ModuleConsensusVersion>>()
1230 .await;
1231
1232 while versions.len() < self.num_peers.total() {
1233 versions.push(ModuleConsensusVersion::new(2, 0));
1234 }
1235
1236 assert_eq!(versions.len(), self.num_peers.total());
1237
1238 versions.sort_unstable();
1239
1240 assert!(versions.first() <= versions.last());
1241
1242 versions[self.num_peers.max_evil()]
1243 }
1244
1245 async fn is_contract_funded_once_active(&self, dbtx: &mut DatabaseTransaction<'_>) -> bool {
1249 CONTRACT_FUNDED_ONCE_MODULE_CONSENSUS_VERSION
1250 <= self.consensus_module_consensus_version(dbtx).await
1251 }
1252
1253 fn spawn_peer_supported_consensus_version_task(
1260 api_client: DynModuleApi,
1261 task_group: &TaskGroup,
1262 our_peer_id: PeerId,
1263 ) -> watch::Receiver<Option<ModuleConsensusVersion>> {
1264 let (sender, receiver) = watch::channel(None);
1265 task_group.spawn_cancellable("fetch-peer-consensus-versions", async move {
1266 loop {
1267 let request_futures = api_client
1268 .all_peers()
1269 .iter()
1270 .filter(|&&peer| peer != our_peer_id)
1271 .map(|&peer| {
1272 let api_client = api_client.clone();
1273
1274 async move {
1275 api_client
1276 .request_single_peer::<ModuleConsensusVersion>(
1277 SUPPORTED_MODULE_CONSENSUS_VERSION_ENDPOINT.to_owned(),
1278 ApiRequestErased::default(),
1279 peer,
1280 )
1281 .await
1282 .inspect_err(|err| warn!(
1283 target: LOG_MODULE_LN,
1284 %peer,
1285 err = %err.fmt_compact(),
1286 "Failed to fetch supported consensus version from peer"
1287 ))
1288 .ok()
1289 }
1290 });
1291
1292 let all_peers_supported_version = join_all(request_futures)
1296 .await
1297 .into_iter()
1298 .collect::<Option<Vec<_>>>()
1299 .map(|peer_versions| {
1300 peer_versions
1301 .into_iter()
1302 .chain(std::iter::once(MODULE_CONSENSUS_VERSION))
1303 .min()
1304 .expect("Our own version is always present")
1305 });
1306
1307 debug!(
1308 target: LOG_MODULE_LN,
1309 ?all_peers_supported_version,
1310 "Fetched supported consensus versions from peers"
1311 );
1312
1313 #[allow(clippy::disallowed_methods)]
1314 if sender.send(all_peers_supported_version).is_err() {
1315 warn!(target: LOG_MODULE_LN, "Failed to send consensus version to watch channel, stopping task");
1316 break;
1317 }
1318
1319 sleep(next_poll_delay(all_peers_supported_version.is_some())).await;
1320 }
1321 });
1322 receiver
1323 }
1324
1325 fn validate_decryption_share(
1326 &self,
1327 peer: PeerId,
1328 share: &PreimageDecryptionShare,
1329 message: &EncryptedPreimage,
1330 ) -> bool {
1331 self.cfg
1332 .consensus
1333 .threshold_pub_keys
1334 .public_key_share(peer.to_usize())
1335 .verify_decryption_share(&share.0, &message.0)
1336 }
1337
1338 async fn get_offer(
1339 &self,
1340 dbtx: &mut DatabaseTransaction<'_>,
1341 payment_hash: bitcoin_hashes::sha256::Hash,
1342 ) -> Option<IncomingContractOffer> {
1343 dbtx.get_value(&OfferKey(payment_hash)).await
1344 }
1345
1346 async fn wait_offer(
1347 &self,
1348 context: &mut ApiEndpointContext,
1349 payment_hash: bitcoin_hashes::sha256::Hash,
1350 ) -> IncomingContractOffer {
1351 let future = context.wait_key_exists(OfferKey(payment_hash));
1352 future.await
1353 }
1354
1355 async fn get_contract_account(
1356 &self,
1357 dbtx: &mut DatabaseTransaction<'_>,
1358 contract_id: ContractId,
1359 ) -> Option<ContractAccount> {
1360 dbtx.get_value(&ContractKey(contract_id)).await
1361 }
1362
1363 async fn wait_contract_account(
1364 &self,
1365 context: &mut ApiEndpointContext,
1366 contract_id: ContractId,
1367 ) -> ContractAccount {
1368 let future = context.wait_key_exists(ContractKey(contract_id));
1370 future.await
1371 }
1372
1373 async fn wait_outgoing_contract_account_cancelled(
1374 &self,
1375 context: &mut ApiEndpointContext,
1376 contract_id: ContractId,
1377 ) -> ContractAccount {
1378 let future =
1379 context.wait_value_matches(ContractKey(contract_id), |contract| {
1380 match &contract.contract {
1381 FundedContract::Outgoing(c) => c.cancelled,
1382 FundedContract::Incoming(_) => false,
1383 }
1384 });
1385 future.await
1386 }
1387
1388 async fn get_decrypted_preimage_status(
1389 &self,
1390 context: &mut ApiEndpointContext,
1391 contract_id: ContractId,
1392 ) -> Result<(IncomingContractAccount, DecryptedPreimageStatus), ApiError> {
1393 let f_contract = context.wait_key_exists(ContractKey(contract_id));
1394 let contract = f_contract.await;
1395 let incoming_contract_account =
1398 Self::get_incoming_contract_account(contract).ok_or_else(|| {
1399 ApiError::bad_request("Contract is not an incoming contract".to_string())
1400 })?;
1401 Ok(
1402 match &incoming_contract_account.contract.decrypted_preimage {
1403 DecryptedPreimage::Some(key) => (
1404 incoming_contract_account.clone(),
1405 DecryptedPreimageStatus::Some(Preimage(
1406 sha256::Hash::hash(&key.0).to_byte_array(),
1407 )),
1408 ),
1409 DecryptedPreimage::Pending => {
1410 (incoming_contract_account, DecryptedPreimageStatus::Pending)
1411 }
1412 DecryptedPreimage::Invalid => {
1413 (incoming_contract_account, DecryptedPreimageStatus::Invalid)
1414 }
1415 },
1416 )
1417 }
1418
1419 async fn wait_preimage_decrypted(
1420 &self,
1421 context: &mut ApiEndpointContext,
1422 contract_id: ContractId,
1423 ) -> (IncomingContractAccount, Option<Preimage>) {
1424 let future =
1425 context.wait_value_matches(ContractKey(contract_id), |contract| {
1426 match &contract.contract {
1427 FundedContract::Incoming(c) => match c.contract.decrypted_preimage {
1428 DecryptedPreimage::Pending => false,
1429 DecryptedPreimage::Some(_) | DecryptedPreimage::Invalid => true,
1430 },
1431 FundedContract::Outgoing(_) => false,
1432 }
1433 });
1434
1435 let decrypt_preimage = future.await;
1436 let incoming_contract_account = Self::get_incoming_contract_account(decrypt_preimage)
1437 .expect("the matcher above only resolves for incoming contracts");
1438 match incoming_contract_account
1439 .clone()
1440 .contract
1441 .decrypted_preimage
1442 {
1443 DecryptedPreimage::Some(key) => (
1444 incoming_contract_account,
1445 Some(Preimage(sha256::Hash::hash(&key.0).to_byte_array())),
1446 ),
1447 _ => (incoming_contract_account, None),
1448 }
1449 }
1450
1451 fn get_incoming_contract_account(contract: ContractAccount) -> Option<IncomingContractAccount> {
1452 match contract.contract {
1453 FundedContract::Incoming(incoming) => Some(IncomingContractAccount {
1454 amount: contract.amount,
1455 contract: incoming.contract,
1456 }),
1457 FundedContract::Outgoing(_) => None,
1458 }
1459 }
1460
1461 async fn list_gateways(
1462 &self,
1463 dbtx: &mut DatabaseTransaction<'_>,
1464 ) -> Vec<LightningGatewayAnnouncement> {
1465 let stream = dbtx.find_by_prefix(&LightningGatewayKeyPrefix).await;
1466 stream
1467 .filter_map(|(_, gw)| async { if gw.is_expired() { None } else { Some(gw) } })
1468 .collect::<Vec<LightningGatewayRegistration>>()
1469 .await
1470 .into_iter()
1471 .map(LightningGatewayRegistration::unanchor)
1472 .collect::<Vec<LightningGatewayAnnouncement>>()
1473 }
1474
1475 async fn register_gateway(
1492 &self,
1493 dbtx: &mut DatabaseTransaction<'_>,
1494 mut gateway: LightningGatewayAnnouncement,
1495 ) -> anyhow::Result<()> {
1496 self.delete_expired_gateways(dbtx).await;
1502
1503 let gateway_id = gateway.info.gateway_id;
1504
1505 anyhow::ensure!(
1506 gateway.info.fees.proportional_millionths <= 1_000_000,
1507 "Gateway registration fee of {} proportional millionths exceeds the payment itself",
1508 gateway.info.fees.proportional_millionths
1509 );
1510
1511 if let Some(auth) = &gateway.auth {
1514 let msg = create_gateway_registration_message(
1515 self.cfg.consensus.threshold_pub_keys.public_key(),
1516 auth.nonce,
1517 &gateway.info,
1518 );
1519
1520 auth.signature
1521 .verify(&msg, &gateway_id.x_only_public_key().0)
1522 .context("Invalid gateway registration signature")?;
1523 }
1524
1525 if let Some(existing) = dbtx.get_value(&LightningGatewayKey(gateway_id)).await
1528 && let Some(existing_auth) = existing.auth
1529 {
1530 let auth = gateway.auth.as_ref().context(
1531 "Gateway registration is signed and cannot be replaced by an unsigned one",
1532 )?;
1533
1534 anyhow::ensure!(
1541 auth.nonce > existing_auth.nonce || gateway.info == existing.info,
1542 "Gateway registration nonce must increase to change settings, got {} but stored {}",
1543 auth.nonce,
1544 existing_auth.nonce
1545 );
1546
1547 if auth.nonce < existing_auth.nonce {
1555 gateway.auth = Some(existing_auth);
1556 }
1557 }
1558
1559 gateway.vetted = false;
1562
1563 dbtx.insert_entry(&LightningGatewayKey(gateway_id), &gateway.anchor())
1564 .await;
1565
1566 Ok(())
1567 }
1568
1569 async fn delete_expired_gateways(&self, dbtx: &mut DatabaseTransaction<'_>) {
1570 let expired_gateway_keys = dbtx
1571 .find_by_prefix(&LightningGatewayKeyPrefix)
1572 .await
1573 .filter_map(|(key, gw)| async move { if gw.is_expired() { Some(key) } else { None } })
1574 .collect::<Vec<LightningGatewayKey>>()
1575 .await;
1576
1577 for key in expired_gateway_keys {
1578 dbtx.remove_entry(&key).await;
1579 }
1580 }
1581
1582 async fn get_gateway_remove_challenge(
1589 &self,
1590 gateway_id: PublicKey,
1591 dbtx: &mut DatabaseTransaction<'_>,
1592 ) -> Option<sha256::Hash> {
1593 match dbtx.get_value(&LightningGatewayKey(gateway_id)).await {
1594 Some(gateway) => {
1595 let mut valid_until_bytes = vec![];
1596 fedimint_core::encoding::encode_legacy_system_time(
1597 &gateway.valid_until,
1598 &mut valid_until_bytes,
1599 )
1600 .expect("encoding to a vector cannot fail");
1601 let mut challenge_bytes = gateway_id.to_bytes();
1602 challenge_bytes.append(&mut valid_until_bytes);
1603 Some(sha256::Hash::hash(&challenge_bytes))
1604 }
1605 _ => None,
1606 }
1607 }
1608
1609 async fn remove_gateway(
1613 &self,
1614 remove_gateway_request: RemoveGatewayRequest,
1615 dbtx: &mut DatabaseTransaction<'_>,
1616 ) -> anyhow::Result<()> {
1617 let fed_public_key = self.cfg.consensus.threshold_pub_keys.public_key();
1618 let gateway_id = remove_gateway_request.gateway_id;
1619 let our_peer_id = self.our_peer_id;
1620 let signature = remove_gateway_request
1621 .signatures
1622 .get(&our_peer_id)
1623 .ok_or_else(|| {
1624 warn!(target: LOG_MODULE_LN, "No signature provided for gateway: {gateway_id}");
1625 anyhow::anyhow!("No signature provided for gateway {gateway_id}")
1626 })?;
1627
1628 let challenge = self
1631 .get_gateway_remove_challenge(gateway_id, dbtx)
1632 .await
1633 .ok_or(anyhow::anyhow!(
1634 "Gateway {gateway_id} is not registered with peer {our_peer_id}"
1635 ))?;
1636
1637 let msg = create_gateway_remove_message(fed_public_key, our_peer_id, challenge);
1639 signature.verify(&msg, &gateway_id.x_only_public_key().0)?;
1640
1641 dbtx.remove_entry(&LightningGatewayKey(gateway_id)).await;
1642 info!(target: LOG_MODULE_LN, "Successfully removed gateway: {gateway_id}");
1643 Ok(())
1644 }
1645}
1646
1647fn record_funded_contract_metric(updated_contract_account: &ContractAccount) {
1648 LN_FUNDED_CONTRACT_SATS
1649 .with_label_values(&[match updated_contract_account.contract {
1650 FundedContract::Incoming(_) => "incoming",
1651 FundedContract::Outgoing(_) => "outgoing",
1652 }])
1653 .observe(updated_contract_account.amount.sats_f64());
1654}
1655
1656#[cfg(test)]
1657mod tests;