1pub mod api;
2
3use std::collections::BTreeMap;
4use std::fmt::{Display, Formatter};
5use std::future::pending;
6use std::str::FromStr;
7use std::sync::Arc;
8use std::time::{Duration, SystemTime};
9
10use api::{RecurringdApiError, RecurringdClient};
11use async_stream::stream;
12use bitcoin::hashes::sha256;
13use bitcoin::secp256k1::SECP256K1;
14use fedimint_client_module::OperationId;
15use fedimint_client_module::error::OperationAlreadyExistsError;
16use fedimint_client_module::module::ClientContext;
17use fedimint_client_module::oplog::UpdateStreamOrOutcome;
18use fedimint_core::BitcoinHash;
19use fedimint_core::config::FederationId;
20use fedimint_core::core::ModuleKind;
21use fedimint_core::db::{DatabaseError, IDatabaseTransactionOpsCoreTyped};
22use fedimint_core::encoding::{
23 Decodable, DecodeError, Encodable, decode_field_from_finite_reader,
24 decode_legacy_system_time_from_finite_reader, encode_legacy_system_time, with_decoding_context,
25};
26use fedimint_core::module::registry::ModuleDecoderRegistry;
27use fedimint_core::secp256k1::{Keypair, PublicKey};
28use fedimint_core::task::sleep;
29use fedimint_core::util::{BoxFuture, FmtCompact, SafeUrl};
30use fedimint_derive_secret::ChildId;
31use fedimint_eventlog::{Event, EventKind, EventPersistence};
32use futures::StreamExt;
33use futures::future::select_all;
34use lightning_invoice::Bolt11Invoice;
35use serde::{Deserialize, Serialize};
36use thiserror::Error;
37use tokio::select;
38use tokio::sync::Notify;
39use tracing::{debug, trace, warn};
40
41use crate::db::{RecurringPaymentCodeKey, RecurringPaymentCodeKeyPrefix};
42use crate::receive::LightningReceiveError;
43use crate::{
44 CreateBolt11InvoiceError, LightningClientModule, LightningClientStateMachines,
45 LightningOperationMeta, LightningOperationMetaVariant, LnReceiveState, LnSubscribeError,
46 tweak_user_key, tweak_user_secret_key,
47};
48
49const LOG_CLIENT_RECURRING: &str = "fm::client::ln::recurring";
50
51impl LightningClientModule {
52 pub async fn register_recurring_payment_code(
53 &self,
54 protocol: RecurringPaymentProtocol,
55 recurringd_api: SafeUrl,
56 meta: &str,
57 ) -> Result<RecurringPaymentCodeEntry, RecurringdApiError> {
58 self.client_ctx
59 .module_db()
60 .autocommit(
61 |dbtx, _| {
62 let recurringd_api_inner = recurringd_api.clone();
63 let new_recurring_payment_code = self.new_recurring_payment_code.clone();
64 Box::pin(async move {
65 let next_idx = dbtx
66 .find_by_prefix_sorted_descending(&RecurringPaymentCodeKeyPrefix)
67 .await
68 .map(|(k, _)| k.derivation_idx)
69 .next()
70 .await
71 .map_or(0, |last_idx| last_idx + 1);
72
73 let payment_code_root_key = self.get_payment_code_root_key(next_idx);
74
75 let recurringd_client =
76 RecurringdClient::new(&recurringd_api_inner.clone());
77 let register_response = recurringd_client
78 .register_recurring_payment_code(
79 self.client_ctx
80 .get_config()
81 .await
82 .global
83 .calculate_federation_id(),
84 protocol,
85 crate::recurring::PaymentCodeRootKey(
86 payment_code_root_key.public_key(),
87 ),
88 meta,
89 )
90 .await?;
91
92 debug!(
93 target: LOG_CLIENT_RECURRING,
94 ?register_response,
95 "Registered recurring payment code"
96 );
97
98 let payment_code_entry = RecurringPaymentCodeEntry {
99 protocol,
100 root_keypair: payment_code_root_key,
101 code: register_response.recurring_payment_code,
102 recurringd_api: recurringd_api_inner,
103 last_derivation_index: 0,
104 creation_time: fedimint_core::time::now(),
105 meta: meta.to_owned(),
106 };
107 dbtx.insert_new_entry(
108 &crate::db::RecurringPaymentCodeKey {
109 derivation_idx: next_idx,
110 },
111 &payment_code_entry,
112 )
113 .await;
114 dbtx.on_commit(move || new_recurring_payment_code.notify_waiters());
115
116 Ok(payment_code_entry)
117 })
118 },
119 None,
120 )
121 .await
122 .map_err(|e| match e {
123 fedimint_core::db::AutocommitError::ClosureError { error, .. } => error,
124 fedimint_core::db::AutocommitError::CommitFailed { last_error, .. } => {
125 panic!("Commit failed: {last_error}")
126 }
127 })
128 }
129
130 pub async fn get_recurring_payment_codes(&self) -> Vec<(u64, RecurringPaymentCodeEntry)> {
131 Self::get_recurring_payment_codes_static(self.client_ctx.module_db()).await
132 }
133
134 pub async fn get_recurring_payment_codes_static(
135 db: &fedimint_core::db::Database,
136 ) -> Vec<(u64, RecurringPaymentCodeEntry)> {
137 assert!(!db.is_global(), "Needs to run in module context");
138 db.begin_transaction_nc()
139 .await
140 .find_by_prefix(&RecurringPaymentCodeKeyPrefix)
141 .await
142 .map(|(idx, entry)| (idx.derivation_idx, entry))
143 .collect()
144 .await
145 }
146
147 fn get_payment_code_root_key(&self, payment_code_registration_idx: u64) -> Keypair {
148 self.recurring_payment_code_secret
149 .child_key(ChildId(payment_code_registration_idx))
150 .to_secp_key(&self.secp)
151 }
152
153 pub async fn scan_recurring_payment_code_invoices(
154 client: ClientContext<Self>,
155 new_code_registered: Arc<Notify>,
156 ) {
157 const QUERY_RETRY_DELAY: Duration = Duration::from_mins(1);
158
159 loop {
160 let new_code_registered_future = new_code_registered.notified();
164
165 let all_recurring_invoice_futures = Self::get_recurring_payment_codes_static(client.module_db())
167 .await
168 .into_iter()
169 .map(|(payment_code_idx, payment_code)| Box::pin(async move {
170 let client = RecurringdClient::new(&payment_code.recurringd_api.clone());
171 let invoice_index = payment_code.last_derivation_index + 1;
172
173 trace!(
174 target: LOG_CLIENT_RECURRING,
175 root_key=?payment_code.root_keypair.public_key(),
176 %invoice_index,
177 server=%payment_code.recurringd_api,
178 "Waiting for new invoice from recurringd"
179 );
180
181 match client.await_new_invoice(crate::recurring::PaymentCodeRootKey(payment_code.root_keypair.public_key()), invoice_index).await {
182 Ok(invoice) => {Ok((payment_code_idx, payment_code, invoice_index, invoice))}
183 Err(err) => {
184 debug!(
185 target: LOG_CLIENT_RECURRING,
186 err=%err.fmt_compact(),
187 root_key=?payment_code.root_keypair.public_key(),
188 invoice_index=%invoice_index,
189 server=%payment_code.recurringd_api,
190 "Failed querying recurring payment code invoice, will retry in {:?}",
191 QUERY_RETRY_DELAY,
192 );
193 sleep(QUERY_RETRY_DELAY).await;
194 Err(err)
195 }
196 }
197 }))
198 .collect::<Vec<_>>();
199
200 let await_any_invoice: BoxFuture<_> = if all_recurring_invoice_futures.is_empty() {
202 Box::pin(pending())
203 } else {
204 Box::pin(select_all(all_recurring_invoice_futures))
205 };
206
207 let (payment_code_idx, _payment_code, invoice_idx, invoice) = select! {
208 (ret, _, _) = await_any_invoice => match ret {
209 Ok(ret) => ret,
210 Err(_) => {
211 continue;
212 }
213 },
214 () = new_code_registered_future => {
215 continue;
216 }
217 };
218
219 Self::process_recurring_payment_code_invoice(
220 &client,
221 payment_code_idx,
222 invoice_idx,
223 invoice,
224 )
225 .await;
226
227 sleep(Duration::from_secs(1)).await;
229 }
230 }
231
232 async fn process_recurring_payment_code_invoice(
233 client: &ClientContext<Self>,
234 payment_code_idx: u64,
235 invoice_idx: u64,
236 invoice: lightning_invoice::Bolt11Invoice,
237 ) {
238 let mut dbtx = client.module_db().begin_transaction().await;
240 let old_payment_code_entry = dbtx
241 .get_value(&crate::db::RecurringPaymentCodeKey {
242 derivation_idx: payment_code_idx,
243 })
244 .await
245 .expect("We queried it, so it exists in our DB");
246
247 let new_payment_code_entry = RecurringPaymentCodeEntry {
248 last_derivation_index: invoice_idx,
249 ..old_payment_code_entry.clone()
250 };
251 dbtx.insert_entry(
252 &crate::db::RecurringPaymentCodeKey {
253 derivation_idx: payment_code_idx,
254 },
255 &new_payment_code_entry,
256 )
257 .await;
258
259 let mut dbtx_nc = dbtx.to_ref_nc();
263 if let Ok(operation_id) = Self::create_recurring_receive_operation(
264 client,
265 &mut dbtx_nc,
266 &old_payment_code_entry,
267 invoice_idx,
268 invoice,
269 )
270 .await
271 {
272 client
273 .log_event(
274 &mut dbtx_nc,
275 RecurringInvoiceCreatedEvent {
276 payment_code_idx,
277 invoice_idx,
278 operation_id,
279 },
280 )
281 .await;
282 } else {
283 debug_assert!(
284 false,
285 "Recurring invoice operation creation failed, this should never happen"
286 );
287 }
288 drop(dbtx_nc);
289
290 dbtx.commit_tx().await;
291 }
292
293 #[allow(clippy::pedantic)]
294 async fn create_recurring_receive_operation(
295 client: &ClientContext<Self>,
296 dbtx: &mut fedimint_core::db::DatabaseTransaction<'_>,
297 payment_code: &RecurringPaymentCodeEntry,
298 invoice_index: u64,
299 invoice: lightning_invoice::Bolt11Invoice,
300 ) -> Result<OperationId, OperationAlreadyExistsError> {
301 let invoice_key =
303 tweak_user_secret_key(SECP256K1, payment_code.root_keypair, invoice_index);
304
305 let operation_id = OperationId(*invoice.payment_hash().as_ref());
306 debug!(
307 target: LOG_CLIENT_RECURRING,
308 ?operation_id,
309 payment_code_key=?payment_code.root_keypair.public_key(),
310 invoice_index=%invoice_index,
311 "Creating recurring receive operation"
312 );
313 let ln_state =
314 LightningClientStateMachines::Receive(crate::receive::LightningReceiveStateMachine {
315 operation_id,
316 state: crate::receive::LightningReceiveStates::ConfirmedInvoice(
319 crate::receive::LightningReceiveConfirmedInvoice {
320 invoice: invoice.clone(),
321 receiving_key: crate::ReceivingKey::Personal(invoice_key),
322 },
323 ),
324 });
325
326 if let Err(e) = client
327 .manual_operation_start_dbtx(
328 dbtx,
329 operation_id,
330 "ln",
331 LightningOperationMeta {
332 variant: LightningOperationMetaVariant::RecurringPaymentReceive(
333 ReurringPaymentReceiveMeta {
334 payment_code_id: PaymentCodeRootKey(
335 payment_code.root_keypair.public_key(),
336 )
337 .to_payment_code_id(),
338 invoice,
339 },
340 ),
341 extra_meta: serde_json::Value::Null,
342 },
343 vec![client.make_dyn_state(ln_state)],
344 )
345 .await
346 {
347 warn!(
348 target: LOG_CLIENT_RECURRING,
349 ?operation_id,
350 payment_code_key=?payment_code.root_keypair.public_key(),
351 invoice_index=%invoice_index,
352 err = %e.fmt_compact(),
353 "Failed to create recurring receive operation"
354 );
355 Err(e)
356 } else {
357 Ok(operation_id)
358 }
359 }
360
361 pub async fn subscribe_ln_recurring_receive(
362 &self,
363 operation_id: OperationId,
364 ) -> Result<UpdateStreamOrOutcome<LnReceiveState>, LnSubscribeError> {
365 let operation = self.client_ctx.get_operation(operation_id).await?;
366 let LightningOperationMetaVariant::RecurringPaymentReceive(ReurringPaymentReceiveMeta {
367 invoice,
368 ..
369 }) = operation.meta::<LightningOperationMeta>().variant
370 else {
371 return Err(LnSubscribeError::NotARecurringReceive);
372 };
373
374 let client_ctx = self.client_ctx.clone();
375
376 Ok(self.client_ctx.outcome_or_updates(&operation, operation_id, |state| match state {
377 LnReceiveState::Created
378 | LnReceiveState::WaitingForPayment { .. }
379 | LnReceiveState::Funded
380 | LnReceiveState::AwaitingFunds => false,
381 LnReceiveState::Canceled { .. } | LnReceiveState::Claimed => true,
382 }, move || {
383 stream! {
384 let self_ref = client_ctx.self_ref();
385
386 yield LnReceiveState::Created;
387 yield LnReceiveState::WaitingForPayment { invoice: invoice.to_string(), timeout: invoice.expiry_time() };
388
389 match self_ref.await_receive_success(operation_id).await {
390 Ok(()) => {
391 yield LnReceiveState::Funded;
392
393 if let Ok(out_points) = self_ref.await_claim_acceptance(operation_id).await {
394 yield LnReceiveState::AwaitingFunds;
395
396 if client_ctx.await_primary_module_outputs(operation_id, out_points).await.is_ok() {
397 yield LnReceiveState::Claimed;
398 return;
399 }
400 }
401
402 yield LnReceiveState::Canceled { reason: LightningReceiveError::Rejected };
403 }
404 Err(e) => {
405 yield LnReceiveState::Canceled { reason: e };
406 }
407 }
408 }
409 }))
410 }
411
412 pub async fn list_recurring_payment_codes(&self) -> BTreeMap<u64, RecurringPaymentCodeEntry> {
413 self.client_ctx
414 .module_db()
415 .begin_transaction_nc()
416 .await
417 .find_by_prefix(&RecurringPaymentCodeKeyPrefix)
418 .await
419 .map(|(idx, entry)| (idx.derivation_idx, entry))
420 .collect()
421 .await
422 }
423
424 pub async fn get_recurring_payment_code(
425 &self,
426 payment_code_idx: u64,
427 ) -> Option<RecurringPaymentCodeEntry> {
428 self.client_ctx
429 .module_db()
430 .begin_transaction_nc()
431 .await
432 .get_value(&RecurringPaymentCodeKey {
433 derivation_idx: payment_code_idx,
434 })
435 .await
436 }
437
438 pub async fn list_recurring_payment_code_invoices(
439 &self,
440 payment_code_idx: u64,
441 ) -> Option<BTreeMap<u64, OperationId>> {
442 let payment_code = self.get_recurring_payment_code(payment_code_idx).await?;
443
444 let operations = (1..=payment_code.last_derivation_index)
445 .map(|invoice_idx: u64| {
446 let invoice_key = tweak_user_key(
447 SECP256K1,
448 payment_code.root_keypair.public_key(),
449 invoice_idx,
450 );
451 let payment_hash =
452 sha256::Hash::hash(&sha256::Hash::hash(&invoice_key.serialize())[..]);
453 let operation_id = OperationId(*payment_hash.as_ref());
454
455 (invoice_idx, operation_id)
456 })
457 .collect();
458
459 Some(operations)
460 }
461}
462
463#[derive(
464 Debug,
465 Clone,
466 Copy,
467 PartialOrd,
468 Eq,
469 PartialEq,
470 Hash,
471 Encodable,
472 Decodable,
473 Serialize,
474 Deserialize,
475)]
476pub struct PaymentCodeRootKey(pub PublicKey);
477
478#[derive(
479 Debug,
480 Clone,
481 Copy,
482 PartialOrd,
483 Eq,
484 PartialEq,
485 Hash,
486 Encodable,
487 Decodable,
488 Serialize,
489 Deserialize,
490)]
491pub struct PaymentCodeId(sha256::Hash);
492
493impl PaymentCodeRootKey {
494 pub fn to_payment_code_id(&self) -> PaymentCodeId {
495 PaymentCodeId(sha256::Hash::hash(&self.0.serialize()))
496 }
497}
498
499impl Display for PaymentCodeId {
500 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
501 write!(f, "{}", self.0)
502 }
503}
504
505impl FromStr for PaymentCodeId {
506 type Err = bitcoin::hashes::hex::HexToArrayError;
507
508 fn from_str(s: &str) -> Result<Self, Self::Err> {
509 Ok(Self(sha256::Hash::from_str(s)?))
510 }
511}
512
513impl Display for PaymentCodeRootKey {
514 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
515 write!(f, "{}", self.0)
516 }
517}
518
519impl FromStr for PaymentCodeRootKey {
520 type Err = fedimint_core::secp256k1::Error;
521
522 fn from_str(s: &str) -> Result<Self, Self::Err> {
523 Ok(Self(PublicKey::from_str(s)?))
524 }
525}
526
527#[derive(
528 Debug,
529 Clone,
530 Copy,
531 Eq,
532 PartialEq,
533 PartialOrd,
534 Hash,
535 Encodable,
536 Decodable,
537 Serialize,
538 Deserialize,
539)]
540pub enum RecurringPaymentProtocol {
541 LNURL,
542 BOLT12,
543}
544
545#[derive(Debug, Clone, Serialize, Deserialize)]
546pub struct ReurringPaymentReceiveMeta {
547 pub payment_code_id: PaymentCodeId,
548 pub invoice: Bolt11Invoice,
549}
550
551#[derive(Debug, Error)]
552pub enum RecurringPaymentError {
553 #[error("Unsupported protocol: {0:?}")]
554 UnsupportedProtocol(RecurringPaymentProtocol),
555 #[error("Unknown federation ID: {0}")]
556 UnknownFederationId(FederationId),
557 #[error("Unknown payment code: {0:?}")]
558 UnknownPaymentCode(PaymentCodeId),
559 #[error("Unknown lightning receive operation: {0:?}")]
560 UnknownInvoice(OperationId),
561 #[error("No compatible lightning module found")]
562 NoLightningModuleFound,
563 #[error("No gateway found")]
564 NoGatewayFound,
565 #[error("Payment code already exists with different settings: {0:?}")]
566 PaymentCodeAlreadyExists(PaymentCodeRootKey),
567 #[error("Federation already registered: {0}")]
568 FederationAlreadyRegistered(FederationId),
569 #[error("Error joining federation")]
572 JoiningFederationFailed(#[source] Box<dyn std::error::Error + Send + Sync>),
573 #[error("Database error")]
575 Database(#[from] DatabaseError),
576 #[error("The invoice could not be created")]
578 InvoiceCreation(#[source] CreateBolt11InvoiceError),
579 #[error("The lightning receive could not be followed")]
581 Subscribe(#[source] LnSubscribeError),
582 #[error("BOLT11 invoice not confirmed")]
585 InvoiceNotConfirmed,
586}
587
588#[derive(Debug, Clone, Serialize)]
589pub struct RecurringPaymentCodeEntry {
590 pub protocol: RecurringPaymentProtocol,
591 pub root_keypair: Keypair,
592 pub code: String,
593 pub recurringd_api: SafeUrl,
594 pub last_derivation_index: u64,
595 pub creation_time: SystemTime,
596 pub meta: String,
597}
598
599impl Encodable for RecurringPaymentCodeEntry {
600 fn consensus_encode<W: std::io::Write>(&self, writer: &mut W) -> Result<(), std::io::Error> {
601 self.protocol.consensus_encode(writer)?;
602 self.root_keypair.consensus_encode(writer)?;
603 self.code.consensus_encode(writer)?;
604 self.recurringd_api.consensus_encode(writer)?;
605 self.last_derivation_index.consensus_encode(writer)?;
606 encode_legacy_system_time(&self.creation_time, writer)?;
607 self.meta.consensus_encode(writer)
608 }
609}
610
611impl Decodable for RecurringPaymentCodeEntry {
612 fn consensus_decode_partial_from_finite_reader<D: std::io::Read>(
613 decoder: &mut D,
614 modules: &ModuleDecoderRegistry,
615 ) -> Result<Self, DecodeError> {
616 Ok(Self {
617 protocol: decode_field_from_finite_reader(
618 decoder,
619 modules,
620 "Decoding named block field: RecurringPaymentCodeEntry{ ... protocol ... }",
621 )?,
622 root_keypair: decode_field_from_finite_reader(
623 decoder,
624 modules,
625 "Decoding named block field: RecurringPaymentCodeEntry{ ... root_keypair ... }",
626 )?,
627 code: decode_field_from_finite_reader(
628 decoder,
629 modules,
630 "Decoding named block field: RecurringPaymentCodeEntry{ ... code ... }",
631 )?,
632 recurringd_api: decode_field_from_finite_reader(
633 decoder,
634 modules,
635 "Decoding named block field: RecurringPaymentCodeEntry{ ... recurringd_api ... }",
636 )?,
637 last_derivation_index: decode_field_from_finite_reader(
638 decoder,
639 modules,
640 "Decoding named block field: RecurringPaymentCodeEntry{ ... last_derivation_index ... }",
641 )?,
642 creation_time: with_decoding_context(
643 decode_legacy_system_time_from_finite_reader(decoder, modules),
644 "Decoding named block field: RecurringPaymentCodeEntry{ ... creation_time ... }",
645 )?,
646 meta: decode_field_from_finite_reader(
647 decoder,
648 modules,
649 "Decoding named block field: RecurringPaymentCodeEntry{ ... meta ... }",
650 )?,
651 })
652 }
653}
654
655#[derive(Debug, Clone, Serialize, Deserialize)]
663pub struct RecurringInvoiceCreatedEvent {
664 pub payment_code_idx: u64,
665 pub invoice_idx: u64,
666 pub operation_id: OperationId,
667}
668
669impl Event for RecurringInvoiceCreatedEvent {
670 const MODULE: Option<ModuleKind> = Some(fedimint_ln_common::KIND);
671 const KIND: EventKind = EventKind::from_static("recurring_invoice_created");
672 const PERSISTENCE: EventPersistence = EventPersistence::Persistent;
673}
674
675#[cfg(test)]
676mod tests {
677 use std::str::FromStr as _;
678
679 use super::{PaymentCodeId, PaymentCodeRootKey};
680
681 #[test]
684 fn payment_code_id_rejects_non_hex() {
685 let _err: bitcoin::hashes::hex::HexToArrayError =
689 PaymentCodeId::from_str("not hex").expect_err("not hex");
690 }
691
692 #[test]
695 fn payment_code_root_key_rejects_non_key() {
696 let _err: fedimint_core::secp256k1::Error =
700 PaymentCodeRootKey::from_str("not a key").expect_err("not a key");
701 }
702}