1#![deny(clippy::pedantic)]
2#![allow(clippy::cast_possible_wrap)]
3#![allow(clippy::module_name_repetitions)]
4
5pub use fedimint_lnv2_common as common;
6
7pub mod db;
8
9use std::collections::{BTreeMap, BTreeSet};
10use std::time::Duration;
11
12use anyhow::{Context, anyhow, ensure};
13use bls12_381::{G1Projective, Scalar};
14use fedimint_core::bitcoin::hashes::sha256;
15use fedimint_core::config::{
16 ServerModuleConfig, ServerModuleConsensusConfig, TypedServerModuleConfig,
17 TypedServerModuleConsensusConfig,
18};
19use fedimint_core::core::ModuleInstanceId;
20use fedimint_core::db::{
21 Database, DatabaseTransaction, DatabaseVersion, IDatabaseTransactionOpsCoreTyped,
22};
23use fedimint_core::encoding::Encodable;
24use fedimint_core::envs::{FM_ENABLE_MODULE_LNV2_ENV, is_env_var_set_opt};
25use fedimint_core::module::audit::Audit;
26use fedimint_core::module::{
27 Amounts, ApiEndpoint, ApiError, ApiVersion, CoreConsensusVersion, InputMeta,
28 ModuleConsensusVersion, ModuleInit, TransactionItemAmounts, admin_api_endpoint,
29 public_api_endpoint,
30};
31use fedimint_core::task::timeout;
32use fedimint_core::time::duration_since_epoch;
33use fedimint_core::util::SafeUrl;
34use fedimint_core::{
35 BitcoinHash, InPoint, NumPeers, NumPeersExt, OutPoint, PeerId, apply, async_trait_maybe_send,
36 push_db_pair_items,
37};
38use fedimint_lnv2_common::config::{
39 FeeConsensus, LightningClientConfig, LightningConfig, LightningConfigConsensus,
40 LightningConfigPrivate,
41};
42use fedimint_lnv2_common::contracts::{IncomingContract, OutgoingContract};
43use fedimint_lnv2_common::endpoint_constants::{
44 ADD_GATEWAY_ENDPOINT, AWAIT_INCOMING_CONTRACT_ENDPOINT, AWAIT_INCOMING_CONTRACTS_ENDPOINT,
45 AWAIT_PREIMAGE_ENDPOINT, CONSENSUS_BLOCK_COUNT_ENDPOINT, DECRYPTION_KEY_SHARE_ENDPOINT,
46 GATEWAYS_ENDPOINT, OUTGOING_CONTRACT_EXPIRATION_ENDPOINT, REMOVE_GATEWAY_ENDPOINT,
47};
48use fedimint_lnv2_common::{
49 ContractId, LightningCommonInit, LightningConsensusItem, LightningInput, LightningInputError,
50 LightningInputV0, LightningModuleTypes, LightningOutput, LightningOutputError,
51 LightningOutputOutcome, LightningOutputV0, MODULE_CONSENSUS_VERSION, OutgoingWitness,
52};
53use fedimint_logging::LOG_MODULE_LNV2;
54use fedimint_server_core::bitcoin_rpc::ServerBitcoinRpcMonitor;
55use fedimint_server_core::config::{PeerHandleOps, eval_poly_g1};
56use fedimint_server_core::migration::ServerModuleDbMigrationFn;
57use fedimint_server_core::{
58 ConfigGenModuleArgs, EnvVarDoc, ServerModule, ServerModuleInit, ServerModuleInitArgs,
59};
60use futures::StreamExt;
61use group::Curve;
62use group::ff::Field;
63use rand::SeedableRng;
64use rand_chacha::ChaChaRng;
65use strum::IntoEnumIterator;
66use tpe::{
67 AggregatePublicKey, DecryptionKeyShare, PublicKeyShare, SecretKeyShare, derive_pk_share,
68};
69use tracing::trace;
70
71use crate::db::{
72 BlockCountVoteKey, BlockCountVotePrefix, DbKeyPrefix, DecryptionKeyShareKey,
73 DecryptionKeySharePrefix, GatewayKey, GatewayPrefix, IncomingContractIndexKey,
74 IncomingContractIndexPrefix, IncomingContractKey, IncomingContractOutpointKey,
75 IncomingContractOutpointPrefix, IncomingContractPrefix, IncomingContractStreamIndexKey,
76 IncomingContractStreamKey, IncomingContractStreamPrefix, OutgoingContractKey,
77 OutgoingContractPrefix, PreimageKey, PreimagePrefix, UnixTimeVoteKey, UnixTimeVotePrefix,
78};
79
80const MAX_INCOMING_CONTRACTS_BATCH: usize = 1024;
87
88#[derive(Debug, Clone)]
89pub struct LightningInit;
90
91impl ModuleInit for LightningInit {
92 type Common = LightningCommonInit;
93
94 #[allow(clippy::too_many_lines)]
95 async fn dump_database(
96 &self,
97 dbtx: &mut DatabaseTransaction<'_>,
98 prefix_names: Vec<String>,
99 ) -> Box<dyn Iterator<Item = (String, Box<dyn erased_serde::Serialize + Send>)> + '_> {
100 let mut lightning: BTreeMap<String, Box<dyn erased_serde::Serialize + Send>> =
101 BTreeMap::new();
102
103 let filtered_prefixes = DbKeyPrefix::iter().filter(|f| {
104 prefix_names.is_empty() || prefix_names.contains(&f.to_string().to_lowercase())
105 });
106
107 for table in filtered_prefixes {
108 match table {
109 DbKeyPrefix::BlockCountVote => {
110 push_db_pair_items!(
111 dbtx,
112 BlockCountVotePrefix,
113 BlockCountVoteKey,
114 u64,
115 lightning,
116 "Lightning Block Count Votes"
117 );
118 }
119 DbKeyPrefix::UnixTimeVote => {
120 push_db_pair_items!(
121 dbtx,
122 UnixTimeVotePrefix,
123 UnixTimeVoteKey,
124 u64,
125 lightning,
126 "Lightning Unix Time Votes"
127 );
128 }
129 DbKeyPrefix::OutgoingContract => {
130 push_db_pair_items!(
131 dbtx,
132 OutgoingContractPrefix,
133 LightningOutgoingContractKey,
134 OutgoingContract,
135 lightning,
136 "Lightning Outgoing Contracts"
137 );
138 }
139 DbKeyPrefix::IncomingContract => {
140 push_db_pair_items!(
141 dbtx,
142 IncomingContractPrefix,
143 LightningIncomingContractKey,
144 IncomingContract,
145 lightning,
146 "Lightning Incoming Contracts"
147 );
148 }
149 DbKeyPrefix::IncomingContractOutpoint => {
150 push_db_pair_items!(
151 dbtx,
152 IncomingContractOutpointPrefix,
153 LightningIncomingContractOutpointKey,
154 OutPoint,
155 lightning,
156 "Lightning Incoming Contracts Outpoints"
157 );
158 }
159 DbKeyPrefix::DecryptionKeyShare => {
160 push_db_pair_items!(
161 dbtx,
162 DecryptionKeySharePrefix,
163 DecryptionKeyShareKey,
164 DecryptionKeyShare,
165 lightning,
166 "Lightning Decryption Key Share"
167 );
168 }
169 DbKeyPrefix::Preimage => {
170 push_db_pair_items!(
171 dbtx,
172 PreimagePrefix,
173 LightningPreimageKey,
174 [u8; 32],
175 lightning,
176 "Lightning Preimages"
177 );
178 }
179 DbKeyPrefix::Gateway => {
180 push_db_pair_items!(
181 dbtx,
182 GatewayPrefix,
183 GatewayKey,
184 (),
185 lightning,
186 "Lightning Gateways"
187 );
188 }
189 DbKeyPrefix::IncomingContractStreamIndex => {
190 push_db_pair_items!(
191 dbtx,
192 IncomingContractStreamIndexKey,
193 IncomingContractStreamIndexKey,
194 u64,
195 lightning,
196 "Lightning Incoming Contract Stream Index"
197 );
198 }
199 DbKeyPrefix::IncomingContractStream => {
200 push_db_pair_items!(
201 dbtx,
202 IncomingContractStreamPrefix(0),
203 IncomingContractStreamKey,
204 IncomingContract,
205 lightning,
206 "Lightning Incoming Contract Stream"
207 );
208 }
209 DbKeyPrefix::IncomingContractIndex => {
210 push_db_pair_items!(
211 dbtx,
212 IncomingContractIndexPrefix,
213 IncomingContractIndexKey,
214 u64,
215 lightning,
216 "Lightning Incoming Contract Index"
217 );
218 }
219 }
220 }
221
222 Box::new(lightning.into_iter())
223 }
224}
225
226#[apply(async_trait_maybe_send!)]
227impl ServerModuleInit for LightningInit {
228 type Module = Lightning;
229
230 fn versions(&self, _core: CoreConsensusVersion) -> &[ModuleConsensusVersion] {
231 &[MODULE_CONSENSUS_VERSION]
232 }
233
234 fn is_enabled_by_default(&self) -> bool {
235 is_env_var_set_opt(FM_ENABLE_MODULE_LNV2_ENV).unwrap_or(true)
236 }
237
238 fn get_documented_env_vars(&self) -> Vec<EnvVarDoc> {
239 vec![EnvVarDoc {
240 name: FM_ENABLE_MODULE_LNV2_ENV,
241 description: "Set to 0/false to disable the LNv2 Lightning module. Enabled by default.",
242 }]
243 }
244
245 async fn init(&self, args: &ServerModuleInitArgs<Self>) -> anyhow::Result<Self::Module> {
246 Ok(Lightning {
247 cfg: args.cfg().to_typed()?,
248 db: args.db().clone(),
249 server_bitcoin_rpc_monitor: args.server_bitcoin_rpc_monitor(),
250 })
251 }
252
253 fn trusted_dealer_gen(
254 &self,
255 peers: &[PeerId],
256 args: &ConfigGenModuleArgs,
257 ) -> BTreeMap<PeerId, ServerModuleConfig> {
258 let tpe_pks = peers
259 .iter()
260 .map(|peer| (*peer, dealer_pk(peers.to_num_peers(), *peer)))
261 .collect::<BTreeMap<PeerId, PublicKeyShare>>();
262
263 peers
264 .iter()
265 .map(|peer| {
266 let cfg = LightningConfig {
267 consensus: LightningConfigConsensus {
268 tpe_agg_pk: dealer_agg_pk(),
269 tpe_pks: tpe_pks.clone(),
270 fee_consensus: if args.disable_base_fees {
271 FeeConsensus::zero()
272 } else {
273 FeeConsensus::new(0).expect("Relative fee is within range")
274 },
275 network: args.network,
276 },
277 private: LightningConfigPrivate {
278 sk: dealer_sk(peers.to_num_peers(), *peer),
279 },
280 };
281
282 (*peer, cfg.to_erased())
283 })
284 .collect()
285 }
286
287 async fn distributed_gen(
288 &self,
289 peers: &(dyn PeerHandleOps + Send + Sync),
290 args: &ConfigGenModuleArgs,
291 ) -> anyhow::Result<ServerModuleConfig> {
292 let (polynomial, sks) = peers.run_dkg_g1().await?;
293
294 let server = LightningConfig {
295 consensus: LightningConfigConsensus {
296 tpe_agg_pk: tpe::AggregatePublicKey(polynomial[0].to_affine()),
297 tpe_pks: peers
298 .num_peers()
299 .peer_ids()
300 .map(|peer| (peer, PublicKeyShare(eval_poly_g1(&polynomial, &peer))))
301 .collect(),
302 fee_consensus: if args.disable_base_fees {
303 FeeConsensus::zero()
304 } else {
305 FeeConsensus::new(0).expect("Relative fee is within range")
306 },
307 network: args.network,
308 },
309 private: LightningConfigPrivate {
310 sk: SecretKeyShare(sks),
311 },
312 };
313
314 Ok(server.to_erased())
315 }
316
317 fn validate_config(&self, identity: &PeerId, config: ServerModuleConfig) -> anyhow::Result<()> {
318 let config = config.to_typed::<LightningConfig>()?;
319
320 ensure!(
321 tpe::derive_pk_share(&config.private.sk)
322 == *config
323 .consensus
324 .tpe_pks
325 .get(identity)
326 .context("Public key set has no key for our identity")?,
327 "Preimge encryption secret key share does not match our public key share"
328 );
329
330 Ok(())
331 }
332
333 fn get_client_config(
334 &self,
335 config: &ServerModuleConsensusConfig,
336 ) -> anyhow::Result<LightningClientConfig> {
337 let config = LightningConfigConsensus::from_erased(config)?;
338 Ok(LightningClientConfig {
339 tpe_agg_pk: config.tpe_agg_pk,
340 tpe_pks: config.tpe_pks,
341 fee_consensus: config.fee_consensus,
342 network: config.network,
343 })
344 }
345
346 fn get_database_migrations(
347 &self,
348 ) -> BTreeMap<DatabaseVersion, ServerModuleDbMigrationFn<Lightning>> {
349 let mut migrations: BTreeMap<DatabaseVersion, ServerModuleDbMigrationFn<Lightning>> =
350 BTreeMap::new();
351
352 migrations.insert(
353 DatabaseVersion(0),
354 Box::new(move |ctx| Box::pin(crate::db::migrate_to_v1(ctx))),
355 );
356
357 migrations
358 }
359
360 fn used_db_prefixes(&self) -> Option<BTreeSet<u8>> {
361 Some(DbKeyPrefix::iter().map(|p| p as u8).collect())
362 }
363}
364
365fn dealer_agg_pk() -> AggregatePublicKey {
366 AggregatePublicKey((G1Projective::generator() * coefficient(0)).to_affine())
367}
368
369fn dealer_pk(num_peers: NumPeers, peer: PeerId) -> PublicKeyShare {
370 derive_pk_share(&dealer_sk(num_peers, peer))
371}
372
373fn dealer_sk(num_peers: NumPeers, peer: PeerId) -> SecretKeyShare {
374 let x = Scalar::from(peer.to_usize() as u64 + 1);
375
376 let y = (0..num_peers.threshold())
380 .map(|index| coefficient(index as u64))
381 .rev()
382 .reduce(|accumulator, c| accumulator * x + c)
383 .expect("We have at least one coefficient");
384
385 SecretKeyShare(y)
386}
387
388fn coefficient(index: u64) -> Scalar {
389 Scalar::random(&mut ChaChaRng::from_seed(
390 *index.consensus_hash::<sha256::Hash>().as_byte_array(),
391 ))
392}
393
394#[derive(Debug)]
395pub struct Lightning {
396 cfg: LightningConfig,
397 db: Database,
398 server_bitcoin_rpc_monitor: ServerBitcoinRpcMonitor,
399}
400
401#[apply(async_trait_maybe_send!)]
402impl ServerModule for Lightning {
403 type Common = LightningModuleTypes;
404 type Init = LightningInit;
405
406 async fn consensus_proposal(
407 &self,
408 _dbtx: &mut DatabaseTransaction<'_>,
409 ) -> Vec<LightningConsensusItem> {
410 let mut items = vec![LightningConsensusItem::UnixTimeVote(
413 60 * (duration_since_epoch().as_secs() / 60),
414 )];
415
416 if let Ok(block_count) = self.get_block_count() {
417 trace!(target: LOG_MODULE_LNV2, ?block_count, "Proposing block count");
418 items.push(LightningConsensusItem::BlockCountVote(block_count));
419 }
420
421 items
422 }
423
424 async fn process_consensus_item<'a, 'b>(
425 &'a self,
426 dbtx: &mut DatabaseTransaction<'b>,
427 consensus_item: LightningConsensusItem,
428 peer: PeerId,
429 ) -> anyhow::Result<()> {
430 trace!(target: LOG_MODULE_LNV2, ?consensus_item, "Processing consensus item proposal");
431
432 match consensus_item {
433 LightningConsensusItem::BlockCountVote(vote) => {
434 let current_vote = dbtx
435 .insert_entry(&BlockCountVoteKey(peer), &vote)
436 .await
437 .unwrap_or(0);
438
439 ensure!(current_vote < vote, "Block count vote is redundant");
440
441 Ok(())
442 }
443 LightningConsensusItem::UnixTimeVote(vote) => {
444 let current_vote = dbtx
445 .insert_entry(&UnixTimeVoteKey(peer), &vote)
446 .await
447 .unwrap_or(0);
448
449 ensure!(current_vote < vote, "Unix time vote is redundant");
450
451 Ok(())
452 }
453 LightningConsensusItem::Default { variant, .. } => Err(anyhow!(
454 "Received lnv2 consensus item with unknown variant {variant}"
455 )),
456 }
457 }
458
459 async fn process_input<'a, 'b, 'c>(
460 &'a self,
461 dbtx: &mut DatabaseTransaction<'c>,
462 input: &'b LightningInput,
463 _in_point: InPoint,
464 ) -> Result<InputMeta, LightningInputError> {
465 let (pub_key, amount) = match input.ensure_v0_ref()? {
466 LightningInputV0::Outgoing(outpoint, outgoing_witness) => {
467 let contract = dbtx
468 .remove_entry(&OutgoingContractKey(*outpoint))
469 .await
470 .ok_or(LightningInputError::UnknownContract)?;
471
472 let pub_key = match outgoing_witness {
473 OutgoingWitness::Claim(preimage) => {
474 if contract.expiration <= self.consensus_block_count(dbtx).await {
475 return Err(LightningInputError::Expired);
476 }
477
478 if !contract.verify_preimage(preimage) {
479 return Err(LightningInputError::InvalidPreimage);
480 }
481
482 dbtx.insert_entry(&PreimageKey(*outpoint), preimage).await;
483
484 contract.claim_pk
485 }
486 OutgoingWitness::Refund => {
487 if contract.expiration > self.consensus_block_count(dbtx).await {
488 return Err(LightningInputError::NotExpired);
489 }
490
491 contract.refund_pk
492 }
493 OutgoingWitness::Cancel(forfeit_signature) => {
494 if !contract.verify_forfeit_signature(forfeit_signature) {
495 return Err(LightningInputError::InvalidForfeitSignature);
496 }
497
498 contract.refund_pk
499 }
500 };
501
502 (pub_key, contract.amount)
503 }
504 LightningInputV0::Incoming(outpoint, agg_decryption_key) => {
505 let contract = dbtx
506 .remove_entry(&IncomingContractKey(*outpoint))
507 .await
508 .ok_or(LightningInputError::UnknownContract)?;
509
510 let index = dbtx
511 .remove_entry(&IncomingContractIndexKey(*outpoint))
512 .await
513 .expect("Incoming contract index should exist");
514
515 dbtx.remove_entry(&IncomingContractStreamKey(index)).await;
516
517 if !contract
518 .verify_agg_decryption_key(&self.cfg.consensus.tpe_agg_pk, agg_decryption_key)
519 {
520 return Err(LightningInputError::InvalidDecryptionKey);
521 }
522
523 let pub_key = match contract.decrypt_preimage(agg_decryption_key) {
524 Some(..) => contract.commitment.claim_pk,
525 None => contract.commitment.refund_pk,
526 };
527
528 (pub_key, contract.commitment.amount)
529 }
530 };
531
532 Ok(InputMeta {
533 amount: TransactionItemAmounts {
534 amounts: Amounts::new_bitcoin(amount),
535 fees: Amounts::new_bitcoin(self.cfg.consensus.fee_consensus.fee(amount)),
536 },
537 pub_key,
538 })
539 }
540
541 async fn process_output<'a, 'b>(
542 &'a self,
543 dbtx: &mut DatabaseTransaction<'b>,
544 output: &'a LightningOutput,
545 outpoint: OutPoint,
546 ) -> Result<TransactionItemAmounts, LightningOutputError> {
547 let amount = match output.ensure_v0_ref()? {
548 LightningOutputV0::Outgoing(contract) => {
549 dbtx.insert_new_entry(&OutgoingContractKey(outpoint), contract)
550 .await;
551
552 contract.amount
553 }
554 LightningOutputV0::Incoming(contract) => {
555 if !contract.verify() {
556 return Err(LightningOutputError::InvalidContract);
557 }
558
559 if contract.commitment.expiration_or_fee <= self.consensus_unix_time(dbtx).await {
560 return Err(LightningOutputError::ContractExpired);
561 }
562
563 dbtx.insert_new_entry(&IncomingContractKey(outpoint), contract)
564 .await;
565
566 dbtx.insert_entry(
567 &IncomingContractOutpointKey(contract.contract_id()),
568 &outpoint,
569 )
570 .await;
571
572 let stream_index = dbtx
573 .get_value(&IncomingContractStreamIndexKey)
574 .await
575 .unwrap_or(0);
576
577 dbtx.insert_entry(&IncomingContractStreamKey(stream_index), contract)
578 .await;
579
580 dbtx.insert_entry(&IncomingContractIndexKey(outpoint), &stream_index)
581 .await;
582
583 dbtx.insert_entry(&IncomingContractStreamIndexKey, &(stream_index + 1))
584 .await;
585
586 let dk_share = contract.create_decryption_key_share(&self.cfg.private.sk);
587
588 dbtx.insert_entry(&DecryptionKeyShareKey(outpoint), &dk_share)
589 .await;
590
591 contract.commitment.amount
592 }
593 };
594
595 Ok(TransactionItemAmounts {
596 amounts: Amounts::new_bitcoin(amount),
597 fees: Amounts::new_bitcoin(self.cfg.consensus.fee_consensus.fee(amount)),
598 })
599 }
600
601 async fn output_status(
602 &self,
603 _dbtx: &mut DatabaseTransaction<'_>,
604 _out_point: OutPoint,
605 ) -> Option<LightningOutputOutcome> {
606 None
607 }
608
609 async fn audit(
610 &self,
611 dbtx: &mut DatabaseTransaction<'_>,
612 audit: &mut Audit,
613 module_instance_id: ModuleInstanceId,
614 ) {
615 audit
618 .add_items(
619 dbtx,
620 module_instance_id,
621 &OutgoingContractPrefix,
622 |_, contract| -(contract.amount.msats as i64),
623 )
624 .await;
625
626 audit
627 .add_items(
628 dbtx,
629 module_instance_id,
630 &IncomingContractPrefix,
631 |_, contract| -(contract.commitment.amount.msats as i64),
632 )
633 .await;
634 }
635
636 fn api_endpoints(&self) -> Vec<ApiEndpoint<Self>> {
637 vec![
638 public_api_endpoint! {
639 CONSENSUS_BLOCK_COUNT_ENDPOINT,
640 ApiVersion::new(0, 0),
641 async |module: &Lightning, context, _params : () | -> u64 {
642 let db = context.db();
643 let mut dbtx = db.begin_transaction_nc().await;
644
645 Ok(module.consensus_block_count(&mut dbtx).await)
646 }
647 },
648 public_api_endpoint! {
649 AWAIT_INCOMING_CONTRACT_ENDPOINT,
650 ApiVersion::new(0, 0),
651 async |module: &Lightning, context, params: (ContractId, u64) | -> Option<OutPoint> {
652 let db = context.db();
653
654 Ok(module.await_incoming_contract(db, params.0, params.1).await)
655 }
656 },
657 public_api_endpoint! {
658 AWAIT_PREIMAGE_ENDPOINT,
659 ApiVersion::new(0, 0),
660 async |module: &Lightning, context, params: (OutPoint, u64)| -> Option<[u8; 32]> {
661 let db = context.db();
662
663 Ok(module.await_preimage(db, params.0, params.1).await)
664 }
665 },
666 public_api_endpoint! {
667 DECRYPTION_KEY_SHARE_ENDPOINT,
668 ApiVersion::new(0, 0),
669 async |_module: &Lightning, context, params: OutPoint| -> DecryptionKeyShare {
670 Ok(context
676 .db()
677 .wait_key_exists(&DecryptionKeyShareKey(params))
678 .await)
679 }
680 },
681 public_api_endpoint! {
682 OUTGOING_CONTRACT_EXPIRATION_ENDPOINT,
683 ApiVersion::new(0, 0),
684 async |module: &Lightning, context, outpoint: OutPoint| -> Option<(ContractId, u64)> {
685 let db = context.db();
686
687 Ok(module.outgoing_contract_expiration(db, outpoint).await)
688 }
689 },
690 public_api_endpoint! {
691 AWAIT_INCOMING_CONTRACTS_ENDPOINT,
692 ApiVersion::new(0, 0),
693 async |module: &Lightning, context, params: (u64, usize)| -> (Vec<IncomingContract>, u64) {
694 let db = context.db();
695
696 if params.1 == 0 || params.1 > MAX_INCOMING_CONTRACTS_BATCH {
697 return Err(ApiError::bad_request(format!(
698 "Batch size must be in 1..={MAX_INCOMING_CONTRACTS_BATCH}"
699 )));
700 }
701
702 Ok(module.await_incoming_contracts(db, params.0, params.1).await)
703 }
704 },
705 admin_api_endpoint! {
706 ADD_GATEWAY_ENDPOINT,
707 ApiVersion::new(0, 0),
708 async |_module: &Lightning, context, gateway: SafeUrl| -> bool {
709
710 let db = context.db();
711
712 Ok(Lightning::add_gateway(db, gateway).await)
713 }
714 },
715 admin_api_endpoint! {
716 REMOVE_GATEWAY_ENDPOINT,
717 ApiVersion::new(0, 0),
718 async |_module: &Lightning, context, gateway: SafeUrl| -> bool {
719
720 let db = context.db();
721
722 Ok(Lightning::remove_gateway(db, gateway).await)
723 }
724 },
725 public_api_endpoint! {
726 GATEWAYS_ENDPOINT,
727 ApiVersion::new(0, 0),
728 async |_module: &Lightning, context, _params : () | -> Vec<SafeUrl> {
729 let db = context.db();
730
731 Ok(Lightning::gateways(db).await)
732 }
733 },
734 ]
735 }
736}
737
738impl Lightning {
739 fn get_block_count(&self) -> anyhow::Result<u64> {
740 self.server_bitcoin_rpc_monitor
741 .status()
742 .map(|status| status.block_count)
743 .context("Block count not available yet")
744 }
745
746 async fn consensus_block_count(&self, dbtx: &mut DatabaseTransaction<'_>) -> u64 {
747 let num_peers = self.cfg.consensus.tpe_pks.to_num_peers();
748
749 let mut counts = dbtx
750 .find_by_prefix(&BlockCountVotePrefix)
751 .await
752 .map(|entry| entry.1)
753 .collect::<Vec<u64>>()
754 .await;
755
756 counts.sort_unstable();
757
758 counts.reverse();
759
760 assert!(counts.last() <= counts.first());
761
762 counts.get(num_peers.threshold() - 1).copied().unwrap_or(0)
767 }
768
769 async fn consensus_unix_time(&self, dbtx: &mut DatabaseTransaction<'_>) -> u64 {
770 let num_peers = self.cfg.consensus.tpe_pks.to_num_peers();
771
772 let mut times = dbtx
773 .find_by_prefix(&UnixTimeVotePrefix)
774 .await
775 .map(|entry| entry.1)
776 .collect::<Vec<u64>>()
777 .await;
778
779 times.sort_unstable();
780
781 times.reverse();
782
783 assert!(times.last() <= times.first());
784
785 times.get(num_peers.threshold() - 1).copied().unwrap_or(0)
790 }
791
792 async fn await_incoming_contract(
793 &self,
794 db: Database,
795 contract_id: ContractId,
796 expiration: u64,
797 ) -> Option<OutPoint> {
798 loop {
799 timeout(
800 Duration::from_secs(10),
801 db.wait_key_exists(&IncomingContractOutpointKey(contract_id)),
802 )
803 .await
804 .ok();
805
806 let mut dbtx = db.begin_transaction_nc().await;
809
810 if let Some(outpoint) = dbtx
811 .get_value(&IncomingContractOutpointKey(contract_id))
812 .await
813 {
814 return Some(outpoint);
815 }
816
817 if expiration <= self.consensus_unix_time(&mut dbtx).await {
818 return None;
819 }
820 }
821 }
822
823 async fn await_preimage(
824 &self,
825 db: Database,
826 outpoint: OutPoint,
827 expiration: u64,
828 ) -> Option<[u8; 32]> {
829 loop {
830 timeout(
831 Duration::from_secs(10),
832 db.wait_key_exists(&PreimageKey(outpoint)),
833 )
834 .await
835 .ok();
836
837 let mut dbtx = db.begin_transaction_nc().await;
840
841 if let Some(preimage) = dbtx.get_value(&PreimageKey(outpoint)).await {
842 return Some(preimage);
843 }
844
845 if expiration <= self.consensus_block_count(&mut dbtx).await {
846 return None;
847 }
848 }
849 }
850
851 async fn outgoing_contract_expiration(
852 &self,
853 db: Database,
854 outpoint: OutPoint,
855 ) -> Option<(ContractId, u64)> {
856 let mut dbtx = db.begin_transaction_nc().await;
857
858 let contract = dbtx.get_value(&OutgoingContractKey(outpoint)).await?;
859
860 let consensus_block_count = self.consensus_block_count(&mut dbtx).await;
861
862 let expiration = contract.expiration.saturating_sub(consensus_block_count);
863
864 Some((contract.contract_id(), expiration))
865 }
866
867 async fn await_incoming_contracts(
868 &self,
869 db: Database,
870 start: u64,
871 n: usize,
872 ) -> (Vec<IncomingContract>, u64) {
873 let filter = |next_index: Option<u64>| next_index.filter(|i| *i > start);
874
875 let (mut next_index, mut dbtx) = db
876 .wait_key_check(&IncomingContractStreamIndexKey, filter)
877 .await;
878
879 let mut contracts = Vec::with_capacity(n.min(MAX_INCOMING_CONTRACTS_BATCH));
883
884 let range = IncomingContractStreamKey(start)..IncomingContractStreamKey(u64::MAX);
885
886 for (key, contract) in dbtx
887 .find_by_range(range)
888 .await
889 .take(n)
890 .collect::<Vec<(IncomingContractStreamKey, IncomingContract)>>()
891 .await
892 {
893 contracts.push(contract.clone());
894 next_index = key.0 + 1;
895 }
896
897 (contracts, next_index)
898 }
899
900 async fn add_gateway(db: Database, gateway: SafeUrl) -> bool {
901 let mut dbtx = db.begin_transaction().await;
902
903 let is_new_entry = dbtx.insert_entry(&GatewayKey(gateway), &()).await.is_none();
904
905 dbtx.commit_tx().await;
906
907 is_new_entry
908 }
909
910 async fn remove_gateway(db: Database, gateway: SafeUrl) -> bool {
911 let mut dbtx = db.begin_transaction().await;
912
913 let entry_existed = dbtx.remove_entry(&GatewayKey(gateway)).await.is_some();
914
915 dbtx.commit_tx().await;
916
917 entry_existed
918 }
919
920 async fn gateways(db: Database) -> Vec<SafeUrl> {
921 db.begin_transaction_nc()
922 .await
923 .find_by_prefix(&GatewayPrefix)
924 .await
925 .map(|entry| entry.0.0)
926 .collect()
927 .await
928 }
929
930 pub async fn consensus_block_count_ui(&self) -> u64 {
931 self.consensus_block_count(&mut self.db.begin_transaction_nc().await)
932 .await
933 }
934
935 pub async fn consensus_unix_time_ui(&self) -> u64 {
936 self.consensus_unix_time(&mut self.db.begin_transaction_nc().await)
937 .await
938 }
939
940 pub async fn add_gateway_ui(&self, gateway: SafeUrl) -> bool {
941 Self::add_gateway(self.db.clone(), gateway).await
942 }
943
944 pub async fn remove_gateway_ui(&self, gateway: SafeUrl) -> bool {
945 Self::remove_gateway(self.db.clone(), gateway).await
946 }
947
948 pub async fn gateways_ui(&self) -> Vec<SafeUrl> {
949 Self::gateways(self.db.clone()).await
950 }
951}