Skip to main content

fedimint_ln_server/
lib.rs

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        // Eagerly initialize metrics that trigger infrequently
237        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/// The lightning module implements an account system. It does not have the
337/// privacy guarantees of the e-cash mint module but instead allows for smart
338/// contracting. There exist two contract types that can be used to "lock"
339/// accounts:
340///
341///   * [Outgoing]: an account locked with an HTLC-like contract allowing to
342///     incentivize an external Lightning node to make payments for the funder
343///   * [Incoming]: a contract type that represents the acquisition of a
344///     preimage belonging to a hash. Every incoming contract is preceded by an
345///     offer that specifies how much the seller is asking for the preimage to a
346///     particular hash. It also contains some threshold-encrypted data. Once
347///     the contract is funded the data is decrypted. If it is a valid preimage
348///     the contract's funds are now accessible to the creator of the offer, if
349///     not they are accessible to the funder.
350///
351/// These two primitives allow to integrate the federation with the wider
352/// Lightning network through a centralized but untrusted (except for
353/// availability) Lightning gateway server.
354///
355/// [Outgoing]: fedimint_ln_common::contracts::outgoing::OutgoingContract
356/// [Incoming]: fedimint_ln_common::contracts::incoming::IncomingContract
357#[derive(Debug)]
358pub struct Lightning {
359    cfg: LightningConfig,
360    our_peer_id: PeerId,
361    num_peers: NumPeers,
362    /// The highest module consensus version supported by every peer, as
363    /// reported by their APIs, or `None` while any peer has yet to answer.
364    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        // Consensus upgrade activation voting. There is deliberately no manual
392        // override: a new consensus item variant is only understood by upgraded
393        // peers, and peers that predate it skip it rather than fail, which would
394        // silently fork their state. Requiring every peer to report support
395        // before we ever propose the item is what keeps that from happening.
396        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            // Only vote if the commonly supported version is higher than the
400            // currently active one
401            && 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                // we save the first ordered valid decryption share for every peer
452                dbtx.insert_new_entry(&AgreedDecryptionShareKey(contract_id, peer_id), &share)
453                    .await;
454
455                // collect all valid decryption shares previously received for this contract
456                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                    // TODO: check if that can happen even though shares are verified
476                    // before
477                    error!(target: LOG_MODULE_LN, contract_hash = %contract.hash, "Failed to decrypt preimage");
478                    return Ok(());
479                };
480
481                // Delete decryption shares once we've decrypted the preimage
482                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                // TODO: maybe define update helper fn
510                // Update contract
511                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                // Update output outcome
527                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                    // If the timelock hasn't expired yet …
612                    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                    // … and the spender provides a valid preimage …
621                    if preimage_hash != outgoing.hash {
622                        return Err(LightningInputError::InvalidPreimage);
623                    }
624
625                    // … then the contract account can be spent using the gateway key,
626                    outgoing.gateway_key
627                } else {
628                    // otherwise the user can claim the funds back.
629                    outgoing.user_key
630                }
631            }
632            FundedContract::Incoming(incoming) => match &incoming.contract.decrypted_preimage {
633                // Once the preimage has been decrypted …
634                DecryptedPreimage::Pending => {
635                    return Err(LightningInputError::ContractNotReady);
636                }
637                // … either the user may spend the funds since they sold a valid preimage …
638                DecryptedPreimage::Some(preimage) => match preimage.to_public_key() {
639                    Ok(pub_key) => pub_key,
640                    Err(_) => return Err(LightningInputError::InvalidPreimage),
641                },
642                // … or the gateway may claim back funds for not receiving the advertised preimage.
643                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        // When a contract reaches a terminal state, the associated amount will be
653        // updated to 0. At this point, the contract no longer needs to be tracked
654        // for auditing liabilities, so we can safely remove the audit key.
655        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                // From consensus version 2.1 on, a contract account is funded
682                // exactly once. Contract ids do not commit to the full contract
683                // state, so before this version a second funding output for the
684                // same id topped up the existing account while keeping its
685                // state; for an incoming contract whose preimage decryption
686                // already reached a terminal state, the first contract's
687                // gateway or preimage holder could sweep the new funds. The
688                // pre-2.1 top-up path below must remain reachable so historic
689                // sessions replay identically.
690                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                // Incoming contracts are special, they need to match an offer
702                if let Contract::Incoming(incoming) = &contract.contract {
703                    // An incoming contract's id is only its payment hash, so a second
704                    // funding lands on the account created by the first one. While that
705                    // account is still waiting on a decryption proposal, funding it
706                    // again would overwrite the pending proposal, so reject it.
707                    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                        // If the account is not sufficiently funded fail the output
724                        return Err(LightningOutputError::InsufficientIncomingFunding(
725                            offer.amount,
726                            contract.amount,
727                        ));
728                    }
729
730                    // Nothing ties the funding contract's ciphertext to the one
731                    // the offer put up for sale, so a funder other than the
732                    // intended payer can consume an offer with a ciphertext of
733                    // their own. It decrypts to garbage, the `Invalid` arm hands
734                    // the funds back to their own `gateway_key`, and the account
735                    // they leave behind makes the payment hash unofferable and
736                    // unfundable for good.
737                    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                    // A funder who names the decryption outcome themselves takes
744                    // the `Invalid` arm's refund to their own `gateway_key`, and
745                    // the share proposed for the contract is never consumed: the
746                    // contract is not pending, so the decryption share bails, and
747                    // the key it was proposed under is only removed on the path
748                    // that bail skips. `consensus_proposal` then re-emits an item
749                    // for it every second, on every guardian, for good.
750                    if self.is_contract_funded_once_active(dbtx).await
751                        && incoming.decrypted_preimage != DecryptedPreimage::Pending
752                    {
753                        return Err(LightningOutputError::PreDecryptedIncomingContract);
754                    }
755
756                    // The offer's ciphertext is verified when the offer is created, but the
757                    // contract carries its own copy, which consensus decoding only checks
758                    // for valid point encodings. `decrypt_share` below returns `None` for
759                    // exactly the ciphertexts that fail `verify()`.
760                    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                // From consensus version 2.1 on, no offer can be created for a
838                // payment hash whose incoming contract account already exists:
839                // funding it could only top up that account (rejected above
840                // once 2.1 is active), so such an offer is a dead end that
841                // could still lure a gateway into accepting an HTLC it can
842                // never get funded for.
843                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                // Check that each preimage is only offered for sale once, see #1397
859                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                // TODO: sanity-check encrypted preimage size
877                if dbtx
878                    .insert_entry(&OfferKey(offer.hash), &(*offer).clone())
879                    .await
880                    .is_some()
881                {
882                    // Technically the error isn't due to a duplicate encrypted preimage but due to
883                    // a duplicate payment hash, practically it's the same problem though: re-using
884                    // the invoice key. Since we can't eaily extend the error enum we just re-use
885                    // this variant.
886                    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    /// Reject funding a contract that already has an account, and creating an
966    /// offer for a payment hash whose incoming contract account already exists.
967    ///
968    /// Contract ids do not commit to the full contract state — an incoming
969    /// contract's id commits to the payment hash alone — so a second funding
970    /// output for the same id does not create a new account: it tops up the
971    /// existing one, which keeps the first contract's `decrypted_preimage`,
972    /// `encrypted_preimage` and `gateway_key`. If that state is already
973    /// terminal the new funds are immediately spendable by the *first*
974    /// contract's gateway, and no further decryption can take place.
975    ///
976    /// These are the same rules [`ServerModule::process_output`] enforces in
977    /// consensus from module consensus version 2.1 on. Enforcing them here as
978    /// well protects federations that have not activated 2.1 yet, with policy
979    /// strength only: it is only as strong as the weakest guardian and cannot
980    /// see intra-session ordering, so it does not cover two fundings submitted
981    /// in the same session.
982    #[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                // Both incoming and outgoing contracts represent liabilities to the federation
1026                // since they are obligations to issue notes.
1027                |_, 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    /// Whether the funded-exactly-once rules of
1246    /// [`CONTRACT_FUNDED_ONCE_MODULE_CONSENSUS_VERSION`] are active, i.e. the
1247    /// federation has voted that version in.
1248    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    /// Tracks the highest module consensus version supported by *every* peer.
1254    ///
1255    /// A vote is only ever proposed once this reports a version, which requires
1256    /// every peer to have answered. Peers that predate a version do not serve
1257    /// [`SUPPORTED_MODULE_CONSENSUS_VERSION_ENDPOINT`] at all, so they hold the
1258    /// federation back rather than being voted past.
1259    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                // A peer that does not answer runs a binary without version voting, so
1293                // collecting into an `Option` holds the federation back on any absence
1294                // rather than voting that peer past a version it cannot decode.
1295                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        // not using a variable here leads to a !Send error
1369        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        // `ContractKey` holds either contract variant and anyone can fund an
1396        // outgoing contract, so the caller decides which variant we find here.
1397        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    /// Stores a gateway registration, rejecting announcements that are not
1476    /// entitled to overwrite the record currently held for their `gateway_id`.
1477    ///
1478    /// A registration carrying a valid
1479    /// [`fedimint_ln_common::GatewayRegistrationAuth`] outranks an
1480    /// unsigned one. Since only the holder of the secret key behind
1481    /// `gateway_id` can produce a signature, this means:
1482    ///
1483    /// - a gateway that signs cannot have its record replaced by anyone else,
1484    /// - a gateway that does not sign is exactly as exposed as it was before
1485    ///   proofs existed, and no more,
1486    /// - an attacker can never lock a gateway out of its own `gateway_id`,
1487    ///   because unsigned records never block other unsigned registrations.
1488    ///
1489    /// So gateways gain protection individually as they upgrade, with no
1490    /// coordinated rollout and no regression for those that have not.
1491    async fn register_gateway(
1492        &self,
1493        dbtx: &mut DatabaseTransaction<'_>,
1494        mut gateway: LightningGatewayAnnouncement,
1495    ) -> anyhow::Result<()> {
1496        // Garbage collect expired gateways (since we're already writing to the DB)
1497        // Note: A "gotcha" of doing this here is that if two gateways are registered
1498        // at the same time, they will both attempt to delete the same expired gateways
1499        // and one of them will fail. This should be fine, since the other one will
1500        // succeed and the failed one will just try again.
1501        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        // Reject a forged proof outright rather than silently downgrading it to an
1512        // unsigned registration, which would hide a misconfigured gateway.
1513        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        // Registrations are garbage collected above, so anything still present is
1526        // live and its claim on this `gateway_id` has to be honored.
1527        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            // The nonce only has to move for announcements that actually change
1535            // something. Re-registering identical settings is the common case —
1536            // gateways refresh well inside the TTL — and replaying it cannot
1537            // achieve anything beyond extending a lifetime that is clamped
1538            // anyway. Exempting it keeps a gateway whose clock stepped backwards
1539            // from being locked out of refreshing its own registration.
1540            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            // The exemption must not let the ratchet fall back, or it defeats
1548            // itself: gateways refresh every few minutes and `auth` is served
1549            // publicly, so an attacker could replay an old identical-settings
1550            // announcement to lower the stored nonce and then replay an
1551            // intermediate one to restore stale settings. Keep the highest nonce
1552            // seen, along with the signature that goes with it, so the stored
1553            // proof stays self-consistent for clients that verify it.
1554            if auth.nonce < existing_auth.nonce {
1555                gateway.auth = Some(existing_auth);
1556            }
1557        }
1558
1559        // Whether a gateway is vetted is the federation's judgement to make, not a
1560        // property a gateway gets to assert about itself.
1561        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    /// Returns the challenge to the gateway that must be signed by the
1583    /// gateway's private key in order for the gateway registration record
1584    /// to be removed. The challenge is the concatenation of the gateway's
1585    /// public key and the `valid_until` bytes. This ensures that the
1586    /// challenges changes every time the gateway is re-registered and ensures
1587    /// that the challenge is unique per-gateway.
1588    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    /// Removes the gateway registration record. First the signature provided by
1610    /// the gateway is verified by checking if the gateway's challenge has
1611    /// been signed by the gateway's private key.
1612    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        // If there is no challenge, the gateway does not exist in the database and
1629        // there is nothing to do
1630        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        // Verify the supplied schnorr signature is valid
1638        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;