1#![deny(clippy::pedantic)]
2#![allow(clippy::cast_possible_truncation)]
3#![allow(clippy::missing_errors_doc)]
4#![allow(clippy::missing_panics_doc)]
5#![allow(clippy::module_name_repetitions)]
6#![allow(clippy::must_use_candidate)]
7#![allow(clippy::return_self_not_must_use)]
8#![allow(clippy::too_many_lines)]
9
10pub use fedimint_mintv2_common as common;
11
12mod api;
13#[cfg(feature = "cli")]
14mod cli;
15pub mod client_db;
16mod ecash;
17mod events;
18mod input;
19pub mod issuance;
20mod output;
21mod receive;
22
23use std::collections::{BTreeMap, BTreeSet};
24use std::convert::Infallible;
25use std::sync::Arc;
26use std::time::Duration;
27
28use bitcoin_hashes::sha256;
29use client_db::{RecoveryState, RecoveryStateKey, SpendableNoteAmountPrefix, SpendableNotePrefix};
30pub use events::*;
31use fedimint_api_client::api::DynModuleApi;
32use fedimint_client::module::ClientModule;
33use fedimint_client::transaction::{
34 ClientInput, ClientInputBundle, ClientInputSM, ClientOutput, ClientOutputBundle,
35 ClientOutputSM, FeeQuote, FeeQuoteRequest, TransactionBuilder,
36};
37use fedimint_client_module::db::ClientModuleMigrationFn;
38use fedimint_client_module::error::{
39 ClientModuleError, InsufficientBalanceError, OperationLookupError, TransactionSubmitError,
40};
41use fedimint_client_module::module::init::{
42 ClientModuleInit, ClientModuleInitArgs, ClientModuleRecoverArgs,
43 ClientModuleRecoveryPrepareArgs, RecoveryMode,
44};
45use fedimint_client_module::module::recovery::{NoModuleBackup, RecoveryProgress};
46use fedimint_client_module::module::{
47 ClientContext, OutPointRange, PrimaryModulePriority, PrimaryModuleSupport,
48};
49use fedimint_client_module::sm::{Context, DynState, ModuleNotifier, State, StateTransition};
50use fedimint_client_module::{DynGlobalClientContext, sm_enum_variant_translation};
51use fedimint_core::base32::{self, FEDIMINT_PREFIX};
52use fedimint_core::config::FederationId;
53use fedimint_core::core::{IntoDynInstance, ModuleInstanceId, ModuleKind, OperationId};
54use fedimint_core::db::{DatabaseTransaction, DatabaseVersion, IDatabaseTransactionOpsCoreTyped};
55use fedimint_core::encoding::{Decodable, Encodable};
56use fedimint_core::module::{
57 AmountUnit, Amounts, ApiVersion, CommonModuleInit, ModuleCommon, ModuleInit, MultiApiVersion,
58};
59use fedimint_core::secp256k1::rand::{Rng, thread_rng};
60use fedimint_core::secp256k1::{Keypair, PublicKey};
61use fedimint_core::util::backoff_util::custom_backoff;
62use fedimint_core::util::{BoxStream, NextOrPending};
63use fedimint_core::{Amount, OutPoint, PeerId, apply, async_trait_maybe_send};
64use fedimint_derive_secret::DerivableSecret;
65use fedimint_mintv2_common::config::{FeeConsensus, MintClientConfig, client_denominations};
66use fedimint_mintv2_common::{
67 Denomination, KIND, MintCommonInit, MintInput, MintModuleTypes, MintOutput, Note, RecoveryItem,
68};
69use futures::{StreamExt, pin_mut};
70use itertools::Itertools;
71use serde::{Deserialize, Serialize};
72use serde_json::Value;
73use tbs::AggregatePublicKey;
74use thiserror::Error;
75
76use crate::api::MintV2ModuleApi;
77use crate::client_db::SpendableNoteKey;
78pub use crate::ecash::ECash;
79use crate::input::{InputSMCommon, InputSMState, InputStateMachine};
80use crate::issuance::NoteIssuanceRequest;
81use crate::output::{MintOutputStateMachine, OutputSMCommon, OutputSMState};
82use crate::receive::{ReceiveSMState, ReceiveStateMachine};
83
84const TARGET_PER_DENOMINATION: usize = 3;
85const SLICE_SIZE: u64 = 10000;
86const PEER_READMISSION: Duration = Duration::from_secs(60);
88const SLICE_TIMEOUT: Duration = Duration::from_secs(10);
96const MAX_SLICE_TIMEOUT: Duration = Duration::from_secs(30);
99const PARALLEL_HASH_REQUESTS: usize = 10;
100const PARALLEL_SLICE_REQUESTS: usize = 10;
101
102#[derive(Debug, Clone, PartialEq, Eq, Hash, Encodable, Decodable)]
103pub struct SpendableNote {
104 pub denomination: Denomination,
105 pub keypair: Keypair,
106 pub signature: tbs::Signature,
107}
108
109impl SpendableNote {
110 pub fn amount(&self) -> Amount {
111 self.denomination.amount()
112 }
113}
114
115impl SpendableNote {
116 fn nonce(&self) -> PublicKey {
117 self.keypair.public_key()
118 }
119
120 fn note(&self) -> Note {
121 Note {
122 denomination: self.denomination,
123 nonce: self.nonce(),
124 signature: self.signature,
125 }
126 }
127}
128
129#[derive(Debug, Clone, Serialize, Deserialize)]
130pub enum MintOperationMeta {
131 Send {
132 ecash: String,
133 custom_meta: Value,
134 },
135 Reissue {
136 change_outpoint_range: OutPointRange,
137 amount: Amount,
138 custom_meta: Value,
139 },
140 Receive {
141 change_outpoint_range: OutPointRange,
142 ecash: String,
143 custom_meta: Value,
144 },
145}
146
147#[derive(Debug, Clone)]
148pub struct MintClientInit;
149
150impl ModuleInit for MintClientInit {
151 type Common = MintCommonInit;
152
153 async fn dump_database(
154 &self,
155 _dbtx: &mut DatabaseTransaction<'_>,
156 _prefix_names: Vec<String>,
157 ) -> Box<dyn Iterator<Item = (String, Box<dyn erased_serde::Serialize + Send>)> + '_> {
158 Box::new(BTreeMap::new().into_iter())
159 }
160}
161
162#[apply(async_trait_maybe_send!)]
163impl ClientModuleInit for MintClientInit {
164 type Module = MintClientModule;
165
166 fn supported_api_versions(&self) -> MultiApiVersion {
167 MultiApiVersion::try_from_iter([ApiVersion { major: 0, minor: 1 }])
168 .expect("no version conflicts")
169 }
170
171 fn recovery_mode(&self) -> RecoveryMode {
172 RecoveryMode::Usable
173 }
174
175 async fn prepare_recovery(
176 &self,
177 args: &ClientModuleRecoveryPrepareArgs,
178 ) -> Result<(), ClientModuleError> {
179 if args
180 .db()
181 .begin_transaction_nc()
182 .await
183 .get_value(&RecoveryStateKey)
184 .await
185 .is_some()
186 {
187 return Ok(());
188 }
189
190 let state = RecoveryState {
197 next_index: 0,
198 total_items: args
199 .module_api()
200 .fetch_recovery_count()
201 .await
202 .map_err(ClientModuleError::other)?,
203 requests: BTreeMap::new(),
204 nonces: BTreeSet::new(),
205 };
206
207 let mut dbtx = args.db().begin_transaction().await;
208
209 dbtx.insert_entry(&RecoveryStateKey, &state).await;
210
211 dbtx.commit_tx().await;
212
213 Ok(())
214 }
215
216 async fn recover(
217 &self,
218 args: &ClientModuleRecoverArgs<Self>,
219 _snapshot: Option<&NoModuleBackup>,
220 ) -> Result<Option<Amount>, ClientModuleError> {
221 let mut state = args
222 .db()
223 .begin_transaction_nc()
224 .await
225 .get_value(&RecoveryStateKey)
226 .await
227 .expect("Prepare recovery commits the state before the recovery is started");
228
229 if state.next_index == state.total_items {
230 return Ok(None);
231 }
232
233 let peer_pool = PeerPool::new(args.api().all_peers());
234
235 let mut recovery_stream = futures::stream::iter(
236 (state.next_index..state.total_items).step_by(SLICE_SIZE as usize),
237 )
238 .map(|start| {
239 let api = args.module_api().clone();
240 let end = std::cmp::min(start + SLICE_SIZE, state.total_items);
241
242 async move { (start, end, api.fetch_recovery_slice_hash(start, end).await) }
243 })
244 .buffered(PARALLEL_HASH_REQUESTS)
245 .map(|(start, end, hash)| {
246 let module_api = args.module_api().clone();
247 let peer_pool = peer_pool.clone();
248
249 async move {
250 (
251 start,
252 download_slice(module_api, peer_pool, start, end, hash).await,
253 )
254 }
255 })
256 .buffer_unordered(PARALLEL_SLICE_REQUESTS);
262
263 let tweak_filter = issuance::tweak_filter(args.module_root_secret());
264
265 let mut pending: BTreeMap<u64, Vec<RecoveryItem>> = BTreeMap::new();
270
271 loop {
272 let items = loop {
273 if let Some(items) = pending.remove(&state.next_index) {
274 break items;
275 }
276
277 let (start, items) = recovery_stream.next().await.ok_or_else(|| {
278 ClientModuleError::other("Recovery stream finished before recovery is complete")
279 })?;
280
281 pending.insert(start, items);
282 };
283
284 for item in &items {
285 match item {
286 RecoveryItem::Output {
287 denomination,
288 nonce_hash,
289 tweak,
290 } => {
291 if !issuance::check_tweak(*tweak, tweak_filter) {
292 continue;
293 }
294 let output_secret = issuance::output_secret(
295 *denomination,
296 *tweak,
297 args.module_root_secret(),
298 );
299
300 if !issuance::check_nonce(&output_secret, *nonce_hash) {
301 continue;
302 }
303
304 let computed_nonce_hash = issuance::nonce(&output_secret).consensus_hash();
305
306 if !state.nonces.insert(computed_nonce_hash) {
308 continue;
309 }
310
311 state.requests.insert(
312 computed_nonce_hash,
313 NoteIssuanceRequest::new(
314 *denomination,
315 *tweak,
316 args.module_root_secret(),
317 ),
318 );
319 }
320 RecoveryItem::Input { nonce_hash } => {
321 state.requests.remove(nonce_hash);
322 state.nonces.remove(nonce_hash);
323 }
324 }
325 }
326
327 state.next_index += items.len() as u64;
328
329 let mut dbtx = args.db().begin_transaction().await;
330
331 dbtx.insert_entry(&RecoveryStateKey, &state).await;
332
333 if state.next_index == state.total_items {
334 let recovered_amount = state
336 .requests
337 .values()
338 .map(|request| request.denomination.amount())
339 .sum::<Amount>();
340
341 let state_machines = args
342 .context()
343 .map_dyn(vec![MintClientStateMachines::Output(
344 MintOutputStateMachine {
345 common: OutputSMCommon {
346 operation_id: OperationId::new_random(),
347 range: None,
348 issuance_requests: state.requests.into_values().collect(),
349 },
350 state: OutputSMState::Pending,
351 },
352 )])
353 .collect();
354
355 args.context()
356 .add_state_machines_dbtx(&mut dbtx.to_ref_nc(), state_machines)
357 .await
358 .expect("state machine is valid");
359
360 dbtx.commit_tx().await;
361
362 return Ok(Some(recovered_amount));
363 }
364
365 dbtx.commit_tx().await;
366
367 args.update_recovery_progress(RecoveryProgress {
368 complete: state.next_index.try_into().unwrap_or(u32::MAX),
369 total: state.total_items.try_into().unwrap_or(u32::MAX),
370 });
371 }
372 }
373
374 async fn init(
375 &self,
376 args: &ClientModuleInitArgs<Self>,
377 ) -> Result<Self::Module, ClientModuleError> {
378 let (tweak_sender, tweak_receiver) = async_channel::bounded(50);
379
380 let filter = issuance::tweak_filter(args.module_root_secret());
381
382 fedimint_core::task::spawn("mintv2-tweak-grinder", async move {
389 loop {
390 let tweak: [u8; 16] = thread_rng().r#gen();
391
392 if !issuance::check_tweak(tweak, filter) {
393 continue;
394 }
395
396 if tweak_sender.send(tweak).await.is_err() {
397 return;
398 }
399
400 fedimint_core::task::sleep(Duration::ZERO).await;
401 }
402 });
403
404 Ok(MintClientModule {
405 federation_id: *args.federation_id(),
406 cfg: args.cfg().clone(),
407 root_secret: args.module_root_secret().clone(),
408 notifier: args.notifier().clone(),
409 client_ctx: args.context(),
410 balance_update_sender: tokio::sync::watch::channel(()).0,
411 tweak_receiver,
412 })
413 }
414
415 fn get_database_migrations(&self) -> BTreeMap<DatabaseVersion, ClientModuleMigrationFn> {
416 BTreeMap::new()
417 }
418}
419
420#[derive(Debug)]
421pub struct MintClientModule {
422 federation_id: FederationId,
423 cfg: MintClientConfig,
424 root_secret: DerivableSecret,
425 notifier: ModuleNotifier<MintClientStateMachines>,
426 client_ctx: ClientContext<Self>,
427 balance_update_sender: tokio::sync::watch::Sender<()>,
428 tweak_receiver: async_channel::Receiver<[u8; 16]>,
429}
430
431#[derive(Debug, Clone)]
432pub struct MintClientContext {
433 client_ctx: ClientContext<MintClientModule>,
434 tbs_agg_pks: BTreeMap<Denomination, AggregatePublicKey>,
435 tbs_pks: BTreeMap<Denomination, BTreeMap<PeerId, tbs::PublicKeyShare>>,
436 pub balance_update_sender: tokio::sync::watch::Sender<()>,
437}
438
439impl Context for MintClientContext {
440 const KIND: Option<ModuleKind> = Some(KIND);
441}
442
443#[apply(async_trait_maybe_send!)]
444impl ClientModule for MintClientModule {
445 type Init = MintClientInit;
446 type Common = MintModuleTypes;
447 type Backup = NoModuleBackup;
448 type ModuleStateMachineContext = MintClientContext;
449 type States = MintClientStateMachines;
450
451 fn context(&self) -> Self::ModuleStateMachineContext {
452 MintClientContext {
453 client_ctx: self.client_ctx.clone(),
454 tbs_agg_pks: self.cfg.tbs_agg_pks.clone(),
455 tbs_pks: self.cfg.tbs_pks.clone(),
456 balance_update_sender: self.balance_update_sender.clone(),
457 }
458 }
459
460 fn input_fee(
461 &self,
462 amounts: &Amounts,
463 _input: &<Self::Common as ModuleCommon>::Input,
464 ) -> Option<Amounts> {
465 let unit = self.cfg.amount_unit;
466 let amount = amounts.get(&unit).copied().unwrap_or_default();
467 let fee = self.cfg.fee_consensus.fee(amount);
468
469 Some(Amounts::new_custom(unit, fee))
470 }
471
472 fn output_fee(
473 &self,
474 amounts: &Amounts,
475 _output: &<Self::Common as ModuleCommon>::Output,
476 ) -> Option<Amounts> {
477 let unit = self.cfg.amount_unit;
478 let amount = amounts.get(&unit).copied().unwrap_or_default();
479 let fee = self.cfg.fee_consensus.fee(amount);
480
481 Some(Amounts::new_custom(unit, fee))
482 }
483
484 #[cfg(feature = "cli")]
485 async fn handle_cli_command(
486 &self,
487 args: &[std::ffi::OsString],
488 ) -> Result<serde_json::Value, ClientModuleError> {
489 cli::handle_cli_command(self, args)
490 .await
491 .map_err(ClientModuleError::other)
492 }
493
494 fn supports_being_primary(&self) -> PrimaryModuleSupport {
495 PrimaryModuleSupport::selected(PrimaryModulePriority::HIGH, [self.cfg.amount_unit])
496 }
497
498 async fn create_final_inputs_and_outputs(
499 &self,
500 dbtx: &mut DatabaseTransaction<'_>,
501 operation_id: OperationId,
502 unit: AmountUnit,
503 mut input_amount: Amount,
504 mut output_amount: Amount,
505 ) -> Result<
506 (
507 ClientInputBundle<MintInput, MintClientStateMachines>,
508 ClientOutputBundle<MintOutput, MintClientStateMachines>,
509 ),
510 ClientModuleError,
511 > {
512 if unit != self.cfg.amount_unit {
513 return Err(ClientModuleError::other(
514 "Module can only handle its configured amount unit",
515 ));
516 }
517
518 let requested_amount = output_amount.saturating_sub(input_amount);
519 let Some(funding_notes) = self.select_funding_input(dbtx, requested_amount).await else {
522 let total_amount = self.get_balance(dbtx, unit).await;
523 return Err(InsufficientBalanceError {
524 requested_amount,
525 total_amount,
526 }
527 .into());
528 };
529
530 for note in &funding_notes {
531 self.remove_spendable_note(dbtx, note).await;
532 }
533
534 input_amount += funding_notes.iter().map(SpendableNote::amount).sum();
535
536 output_amount += funding_notes
537 .iter()
538 .map(|input| self.cfg.fee_consensus.fee(input.amount()))
539 .sum();
540
541 assert!(output_amount <= input_amount);
542
543 let (input_notes, output_amounts) = self
544 .rebalance(dbtx, &self.cfg.fee_consensus, input_amount - output_amount)
545 .await;
546
547 for note in &input_notes {
548 self.remove_spendable_note(dbtx, note).await;
549 }
550
551 input_amount += input_notes.iter().map(SpendableNote::amount).sum();
552
553 output_amount += input_notes
554 .iter()
555 .map(|note| self.cfg.fee_consensus.fee(note.amount()))
556 .sum();
557
558 output_amount += output_amounts
559 .iter()
560 .map(|denomination| {
561 denomination.amount() + self.cfg.fee_consensus.fee(denomination.amount())
562 })
563 .sum();
564
565 assert!(output_amount <= input_amount);
566
567 let mut spendable_notes = funding_notes
568 .into_iter()
569 .chain(input_notes)
570 .collect::<Vec<SpendableNote>>();
571
572 spendable_notes.sort_by_key(|note| note.denomination);
574
575 let input_bundle =
576 Self::create_input_bundle(operation_id, spendable_notes, false, self.cfg.amount_unit);
577
578 let mut denominations = represent_amount_with_fees(
579 input_amount.saturating_sub(output_amount),
580 &self.cfg.fee_consensus,
581 )
582 .into_iter()
583 .chain(output_amounts)
584 .collect::<Vec<Denomination>>();
585
586 denominations.sort();
588
589 let output_bundle = self.create_output_bundle(operation_id, denominations).await;
590
591 let sender = self.balance_update_sender.clone();
592 dbtx.on_commit(move || sender.send_replace(()));
593
594 Ok((input_bundle, output_bundle))
595 }
596
597 async fn await_primary_module_output(
598 &self,
599 operation_id: OperationId,
600 outpoint: OutPoint,
601 ) -> Result<(), ClientModuleError> {
602 self.await_output_sm_success(operation_id, outpoint)
603 .await
604 .map_err(ClientModuleError::other)
605 }
606
607 async fn get_balance(&self, dbtx: &mut DatabaseTransaction<'_>, unit: AmountUnit) -> Amount {
608 if unit != self.cfg.amount_unit {
609 return Amount::ZERO;
610 }
611
612 self.get_count_by_denomination_dbtx(dbtx)
613 .await
614 .into_iter()
615 .map(|(denomination, count)| denomination.amount().mul_u64(count))
616 .sum()
617 }
618
619 async fn subscribe_balance_changes(&self) -> BoxStream<'static, ()> {
620 Box::pin(tokio_stream::wrappers::WatchStream::new(
621 self.balance_update_sender.subscribe(),
622 ))
623 }
624}
625
626impl MintClientModule {
627 async fn select_funding_input(
628 &self,
629 dbtx: &mut DatabaseTransaction<'_>,
630 mut excess_output: Amount,
631 ) -> Option<Vec<SpendableNote>> {
632 let mut selected_notes = Vec::new();
633 let mut target_notes = Vec::new();
634 let mut excess_notes = Vec::new();
635
636 for amount in client_denominations().rev() {
637 let notes_amount = dbtx
638 .find_by_prefix(&SpendableNoteAmountPrefix(amount))
639 .await
640 .map(|entry| entry.0.0)
641 .collect::<Vec<SpendableNote>>()
642 .await;
643
644 target_notes.extend(notes_amount.iter().take(TARGET_PER_DENOMINATION).cloned());
645
646 if notes_amount.len() > 2 * TARGET_PER_DENOMINATION {
647 for note in notes_amount.into_iter().skip(TARGET_PER_DENOMINATION) {
648 let note_fee = self.cfg.fee_consensus.fee(note.amount());
649
650 let note_value = note
651 .amount()
652 .checked_sub(note_fee)
653 .expect("All our notes are economical");
654
655 excess_output = excess_output.saturating_sub(note_value);
656
657 selected_notes.push(note);
658 }
659 } else {
660 excess_notes.extend(notes_amount.into_iter().skip(TARGET_PER_DENOMINATION));
661 }
662 }
663
664 if excess_output == Amount::ZERO {
665 return Some(selected_notes);
666 }
667
668 for note in excess_notes.into_iter().chain(target_notes) {
669 let note_amount = note.amount();
670 let note_value = note_amount
671 .checked_sub(self.cfg.fee_consensus.fee(note_amount))
672 .expect("All our notes are economical");
673
674 excess_output = excess_output.saturating_sub(note_value);
675
676 selected_notes.push(note);
677
678 if excess_output == Amount::ZERO {
679 return Some(selected_notes);
680 }
681 }
682
683 None
684 }
685
686 async fn rebalance(
687 &self,
688 dbtx: &mut DatabaseTransaction<'_>,
689 fee: &FeeConsensus,
690 mut excess_input: Amount,
691 ) -> (Vec<SpendableNote>, Vec<Denomination>) {
692 let n_denominations = self.get_count_by_denomination_dbtx(dbtx).await;
693
694 let mut notes = dbtx
695 .find_by_prefix_sorted_descending(&SpendableNotePrefix)
696 .await
697 .map(|entry| entry.0.0)
698 .fuse();
699
700 let mut input_notes = Vec::new();
701 let mut output_denominations = Vec::new();
702
703 for d in client_denominations() {
704 let n_denomination = n_denominations.get(&d).copied().unwrap_or(0);
705
706 let n_missing = TARGET_PER_DENOMINATION.saturating_sub(n_denomination as usize);
707
708 for _ in 0..n_missing {
709 match excess_input.checked_sub(d.amount() + fee.fee(d.amount())) {
710 Some(remaining_excess) => excess_input = remaining_excess,
711 None => match notes.next().await {
712 Some(note) => {
713 if note.amount() <= d.amount() + fee.fee(d.amount()) {
714 break;
715 }
716
717 excess_input += note.amount() - (d.amount() + fee.fee(d.amount()));
718
719 input_notes.push(note);
720 }
721 None => break,
722 },
723 }
724
725 output_denominations.push(d);
726 }
727 }
728
729 (input_notes, output_denominations)
730 }
731
732 fn create_input_bundle(
733 operation_id: OperationId,
734 notes: Vec<SpendableNote>,
735 include_receive_sm: bool,
736 amount_unit: AmountUnit,
737 ) -> ClientInputBundle<MintInput, MintClientStateMachines> {
738 let inputs = notes
739 .iter()
740 .map(|spendable_note| ClientInput {
741 input: MintInput::new_v0(spendable_note.note()),
742 keys: vec![spendable_note.keypair],
743 amounts: Amounts::new_custom(amount_unit, spendable_note.amount()),
744 })
745 .collect();
746
747 let input_sms = vec![ClientInputSM {
748 state_machines: Arc::new(move |range: OutPointRange| {
749 let mut sms = vec![MintClientStateMachines::Input(InputStateMachine {
750 common: InputSMCommon {
751 operation_id,
752 txid: range.txid(),
753 spendable_notes: notes.clone(),
754 },
755 state: InputSMState::Pending,
756 })];
757
758 if include_receive_sm {
759 sms.push(MintClientStateMachines::Receive(ReceiveStateMachine {
760 common: crate::receive::ReceiveSMCommon {
761 operation_id,
762 txid: range.txid(),
763 },
764 state: crate::receive::ReceiveSMState::Pending,
765 }));
766 }
767
768 sms
769 }),
770 }];
771
772 ClientInputBundle::new(inputs, input_sms)
773 }
774
775 async fn create_output_bundle(
786 &self,
787 operation_id: OperationId,
788 requested_denominations: Vec<Denomination>,
789 ) -> ClientOutputBundle<MintOutput, MintClientStateMachines> {
790 let issuance_requests = futures::stream::iter(requested_denominations)
791 .zip(self.tweak_receiver.clone())
792 .map(|(d, tweak)| NoteIssuanceRequest::new(d, tweak, &self.root_secret))
793 .collect::<Vec<NoteIssuanceRequest>>()
794 .await;
795
796 let amount_unit = self.cfg.amount_unit;
797 let outputs = issuance_requests
798 .iter()
799 .map(|request| ClientOutput {
800 output: request.output(),
801 amounts: Amounts::new_custom(amount_unit, request.denomination.amount()),
802 })
803 .collect();
804
805 let output_sms = vec![ClientOutputSM {
806 state_machines: Arc::new(move |range: OutPointRange| {
807 vec![MintClientStateMachines::Output(MintOutputStateMachine {
808 common: OutputSMCommon {
809 operation_id,
810 range: Some(range),
811 issuance_requests: issuance_requests.clone(),
812 },
813 state: OutputSMState::Pending,
814 })]
815 }),
816 }];
817
818 ClientOutputBundle::new(outputs, output_sms)
819 }
820
821 async fn await_output_sm_success(
831 &self,
832 operation_id: OperationId,
833 outpoint: OutPoint,
834 ) -> Result<(), AwaitOutputError> {
835 let stream = self
836 .notifier
837 .subscribe(operation_id)
838 .await
839 .filter_map(|state| async {
840 let MintClientStateMachines::Output(state) = state else {
841 return None;
842 };
843
844 if !state.common.range?.into_iter().contains(&outpoint) {
845 return None;
846 }
847
848 match state.state {
849 OutputSMState::Pending => None,
850 OutputSMState::Success => Some(Ok(())),
851 OutputSMState::Aborted => Some(Err(AwaitOutputError::Rejected)),
852 OutputSMState::Failure => Some(Err(AwaitOutputError::Failed)),
853 }
854 });
855
856 pin_mut!(stream);
857
858 stream.next_or_pending().await
859 }
860
861 pub async fn get_count_by_denomination(&self) -> BTreeMap<Denomination, u64> {
863 self.get_count_by_denomination_dbtx(
864 &mut self.client_ctx.module_db().begin_transaction_nc().await,
865 )
866 .await
867 }
868
869 async fn get_count_by_denomination_dbtx(
870 &self,
871 dbtx: &mut DatabaseTransaction<'_>,
872 ) -> BTreeMap<Denomination, u64> {
873 dbtx.find_by_prefix(&SpendableNotePrefix)
874 .await
875 .fold(BTreeMap::new(), |mut acc, entry| async move {
876 acc.entry(entry.0.0.denomination)
877 .and_modify(|count| *count += 1)
878 .or_insert(1);
879
880 acc
881 })
882 .await
883 }
884
885 pub async fn send(
911 &self,
912 amount: Amount,
913 custom_meta: Value,
914 include_invite: bool,
915 ) -> Result<(OperationId, ECash), SendECashError> {
916 let amount = round_to_multiple(amount, client_denominations().next().unwrap().amount());
917
918 if let Some((operation_id, ecash)) = self
919 .client_ctx
920 .module_db()
921 .autocommit(
922 |dbtx, _| {
923 Box::pin(self.send_ecash_dbtx(
924 dbtx,
925 amount,
926 custom_meta.clone(),
927 include_invite,
928 ))
929 },
930 Some(100),
931 )
932 .await
933 .expect("Failed to commit dbtx after 100 retries")
934 {
935 return Ok((operation_id, ecash));
936 }
937
938 self.client_ctx
939 .global_api()
940 .session_count()
941 .await
942 .map_err(|_| SendECashError::Offline)?;
943
944 let operation_id = OperationId::new_random();
945
946 let output = self
947 .create_output_bundle(operation_id, represent_amount(amount))
948 .await;
949 let output = self.client_ctx.make_client_outputs(output);
950 let cm = custom_meta.clone();
951
952 let range = self
953 .client_ctx
954 .finalize_and_submit_transaction(
955 operation_id,
956 MintCommonInit::KIND.as_str(),
957 move |change_outpoint_range| MintOperationMeta::Reissue {
958 change_outpoint_range,
959 amount,
960 custom_meta: cm.clone(),
961 },
962 TransactionBuilder::new().with_outputs(output),
963 )
964 .await
965 .map_err(|error| match error {
966 TransactionSubmitError::InsufficientFunds(_) => SendECashError::InsufficientBalance,
967 other => SendECashError::Failed(other),
968 })?;
969
970 for outpoint in range {
971 self.await_output_sm_success(operation_id, outpoint)
972 .await
973 .map_err(|_| SendECashError::Failure)?;
974 }
975
976 Box::pin(self.send(amount, custom_meta, include_invite)).await
977 }
978
979 async fn send_ecash_dbtx(
990 &self,
991 dbtx: &mut DatabaseTransaction<'_>,
992 remaining_amount: Amount,
993 custom_meta: Value,
994 include_invite: bool,
995 ) -> Result<Option<(OperationId, ECash)>, Infallible> {
996 let Some(notes) = Self::select_exact_change(&mut dbtx.to_ref_nc(), remaining_amount).await
997 else {
998 return Ok(None);
999 };
1000
1001 for spendable_note in ¬es {
1002 self.remove_spendable_note(dbtx, spendable_note).await;
1003 }
1004
1005 let ecash = if include_invite {
1006 let invite = self.client_ctx.get_invite_code().await;
1007 ECash::new_with_invite(notes, &invite)
1008 } else {
1009 ECash::new(self.federation_id, notes)
1010 }
1011 .with_unit(self.cfg.amount_unit);
1012 let amount = ecash.amount();
1013 let operation_id = OperationId::new_random();
1014
1015 self.client_ctx
1016 .add_operation_log_entry_dbtx(
1017 dbtx,
1018 operation_id,
1019 MintCommonInit::KIND.as_str(),
1020 MintOperationMeta::Send {
1021 ecash: base32::encode_prefixed(FEDIMINT_PREFIX, &ecash),
1022 custom_meta,
1023 },
1024 )
1025 .await;
1026
1027 self.client_ctx
1028 .log_event(
1029 dbtx,
1030 SendPaymentEvent {
1031 operation_id,
1032 amount,
1033 ecash: base32::encode_prefixed(FEDIMINT_PREFIX, &ecash),
1034 },
1035 )
1036 .await;
1037
1038 let sender = self.balance_update_sender.clone();
1039 dbtx.on_commit(move || sender.send_replace(()));
1040
1041 Ok(Some((operation_id, ecash)))
1042 }
1043
1044 pub async fn receive(
1046 &self,
1047 ecash: ECash,
1048 custom_meta: Value,
1049 ) -> Result<OperationId, ReceiveECashError> {
1050 let operation_id = OperationId::from_encodable(&ecash);
1051
1052 if ecash.mint() != Some(self.federation_id) {
1053 return Err(ReceiveECashError::WrongFederation);
1054 }
1055
1056 if ecash
1057 .notes()
1058 .iter()
1059 .any(|note| note.amount() <= self.cfg.fee_consensus.base_fee())
1060 {
1061 return Err(ReceiveECashError::UneconomicalDenomination);
1062 }
1063
1064 let input =
1065 Self::create_input_bundle(operation_id, ecash.notes(), true, self.cfg.amount_unit);
1066 let input = self.client_ctx.make_client_inputs(input);
1067 let ec = base32::encode_prefixed(FEDIMINT_PREFIX, &ecash);
1068
1069 self.client_ctx
1070 .finalize_and_submit_transaction(
1071 operation_id,
1072 MintCommonInit::KIND.as_str(),
1073 move |change_outpoint_range| MintOperationMeta::Receive {
1074 change_outpoint_range,
1075 ecash: ec.clone(),
1076 custom_meta: custom_meta.clone(),
1077 },
1078 TransactionBuilder::new().with_inputs(input),
1079 )
1080 .await
1081 .map_err(|error| match error {
1082 TransactionSubmitError::OperationAlreadyExists(_) => {
1083 ReceiveECashError::AlreadyReceived
1084 }
1085 TransactionSubmitError::InsufficientFunds(_) => {
1086 ReceiveECashError::InsufficientFunds
1087 }
1088 other => ReceiveECashError::Failed(other),
1089 })?;
1090
1091 let mut dbtx = self.client_ctx.module_db().begin_transaction().await;
1092
1093 self.client_ctx
1094 .log_event(
1095 &mut dbtx,
1096 ReceivePaymentEvent {
1097 operation_id,
1098 amount: ecash.amount(),
1099 },
1100 )
1101 .await;
1102
1103 dbtx.commit_tx().await;
1104
1105 Ok(operation_id)
1106 }
1107
1108 pub async fn receive_fee_quote(
1118 &self,
1119 ecash: &ECash,
1120 ) -> Result<FeeQuote, TransactionSubmitError> {
1121 let notes = ecash.notes();
1125 let input_amount: Amount = notes.iter().map(SpendableNote::amount).sum();
1126 let input_fee: Amount = notes
1127 .iter()
1128 .map(|note| self.cfg.fee_consensus.fee(note.amount()))
1129 .sum();
1130
1131 self.client_ctx
1132 .fee_quote(
1133 OperationId::new_random(),
1134 FeeQuoteRequest {
1135 input_amount: Amounts::new_custom(self.cfg.amount_unit, input_amount),
1136 output_amount: Amounts::ZERO,
1137 input_fee: Amounts::new_custom(self.cfg.amount_unit, input_fee),
1138 output_fee: Amounts::ZERO,
1139 },
1140 )
1141 .await
1142 }
1143
1144 pub async fn send_fee_quote(&self, amount: Amount) -> Result<FeeQuote, TransactionSubmitError> {
1158 let amount = round_to_multiple(amount, client_denominations().next().unwrap().amount());
1159
1160 if self.can_make_exact_change(amount).await {
1162 return Ok(FeeQuote::ZERO);
1163 }
1164
1165 let denominations = represent_amount(amount);
1169 let output_amount: Amount = denominations.iter().map(|d| d.amount()).sum();
1170 let output_fee: Amount = denominations
1171 .iter()
1172 .map(|d| self.cfg.fee_consensus.fee(d.amount()))
1173 .sum();
1174
1175 self.client_ctx
1176 .fee_quote(
1177 OperationId::new_random(),
1178 FeeQuoteRequest {
1179 input_amount: Amounts::ZERO,
1180 output_amount: Amounts::new_custom(self.cfg.amount_unit, output_amount),
1181 input_fee: Amounts::ZERO,
1182 output_fee: Amounts::new_custom(self.cfg.amount_unit, output_fee),
1183 },
1184 )
1185 .await
1186 }
1187
1188 async fn can_make_exact_change(&self, remaining_amount: Amount) -> bool {
1193 let mut dbtx = self.client_ctx.module_db().begin_transaction_nc().await;
1194
1195 Self::select_exact_change(&mut dbtx, remaining_amount)
1196 .await
1197 .is_some()
1198 }
1199
1200 async fn select_exact_change(
1206 dbtx: &mut DatabaseTransaction<'_>,
1207 mut remaining_amount: Amount,
1208 ) -> Option<Vec<SpendableNote>> {
1209 let mut stream = dbtx
1210 .find_by_prefix_sorted_descending(&SpendableNotePrefix)
1211 .await
1212 .map(|entry| entry.0.0);
1213
1214 let mut notes = vec![];
1215
1216 while let Some(spendable_note) = stream.next().await {
1217 remaining_amount = match remaining_amount.checked_sub(spendable_note.amount()) {
1218 Some(amount) => amount,
1219 None => continue,
1220 };
1221
1222 notes.push(spendable_note);
1223
1224 if remaining_amount == Amount::ZERO {
1225 break;
1226 }
1227 }
1228
1229 (remaining_amount == Amount::ZERO).then_some(notes)
1230 }
1231
1232 pub async fn await_final_receive_operation_state(
1234 &self,
1235 operation_id: OperationId,
1236 ) -> Result<FinalReceiveOperationState, OperationLookupError> {
1237 let operation = self.client_ctx.get_operation(operation_id).await?;
1238 let mut stream = self.notifier.subscribe(operation_id).await;
1239
1240 let mut stream = self
1241 .client_ctx
1242 .outcome_or_updates(&operation, operation_id, |_| true, move || {
1243 async_stream::stream! {
1244 loop {
1245 if let Some(MintClientStateMachines::Receive(state)) = stream.next().await {
1246 match state.state {
1247 ReceiveSMState::Pending => {}
1248 ReceiveSMState::Success => {
1249 yield FinalReceiveOperationState::Success;
1250 return;
1251 }
1252 ReceiveSMState::Rejected(..) => {
1253 yield FinalReceiveOperationState::Rejected;
1254 return;
1255 }
1256 }
1257 }
1258 }
1259 }
1260 })
1261 .into_stream();
1262
1263 let mut final_state = None;
1264
1265 while let Some(state) = stream.next().await {
1266 final_state = Some(state);
1267 }
1268
1269 Ok(final_state.expect("Stream contains one final state"))
1270 }
1271
1272 async fn remove_spendable_note(
1273 &self,
1274 dbtx: &mut DatabaseTransaction<'_>,
1275 spendable_note: &SpendableNote,
1276 ) {
1277 dbtx.remove_entry(&SpendableNoteKey(spendable_note.clone()))
1278 .await
1279 .expect("Must delete existing spendable note");
1280 }
1281}
1282
1283#[derive(Clone)]
1290struct PeerPool {
1291 receiver: async_channel::Receiver<PeerId>,
1292 sender: async_channel::Sender<PeerId>,
1293}
1294
1295impl PeerPool {
1296 fn new(peers: &BTreeSet<PeerId>) -> Self {
1297 let (sender, receiver) = async_channel::bounded(peers.len().max(1));
1298
1299 for peer in peers {
1300 sender
1301 .try_send(*peer)
1302 .expect("Capacity was sized to hold every peer");
1303 }
1304
1305 Self { receiver, sender }
1306 }
1307
1308 async fn acquire(&self) -> PeerId {
1310 self.receiver
1311 .recv()
1312 .await
1313 .expect("The sender is held for as long as the receiver")
1314 }
1315
1316 fn retire(&self, peer: PeerId) {
1323 let pool = self.clone();
1324
1325 fedimint_core::runtime::spawn("mintv2 recovery peer readmission", async move {
1326 fedimint_core::runtime::sleep(PEER_READMISSION).await;
1327
1328 pool.release(peer);
1329 });
1330 }
1331
1332 fn release(&self, peer: PeerId) {
1334 self.sender
1335 .try_send(peer)
1336 .expect("Only peers taken from the pool are put back");
1337 }
1338}
1339
1340async fn download_slice(
1343 module_api: DynModuleApi,
1344 peers: PeerPool,
1345 start: u64,
1346 end: u64,
1347 expected_hash: sha256::Hash,
1348) -> Vec<RecoveryItem> {
1349 let mut timeouts = custom_backoff(SLICE_TIMEOUT, MAX_SLICE_TIMEOUT, None);
1350
1351 loop {
1352 let peer = peers.acquire().await;
1353
1354 let timeout = timeouts.next().expect("The backoff never gives up");
1355
1356 let result = module_api
1357 .fetch_recovery_slice(peer, timeout, start, end)
1358 .await;
1359
1360 match result {
1361 Ok(data) if data.consensus_hash::<sha256::Hash>() == expected_hash => {
1362 peers.release(peer);
1363
1364 return data;
1365 }
1366 Ok(_) | Err(_) => peers.retire(peer),
1373 }
1374 }
1375}
1376
1377#[derive(Error, Debug)]
1379#[non_exhaustive]
1380pub enum SendECashError {
1381 #[error("We need to reissue notes but the client is offline")]
1384 Offline,
1385 #[error("The clients balance is insufficient")]
1387 InsufficientBalance,
1388 #[error("A non-recoverable error has occurred")]
1391 Failure,
1392 #[error("The reissue transaction could not be submitted")]
1395 Failed(#[source] TransactionSubmitError),
1396}
1397
1398#[derive(Error, Debug)]
1400#[non_exhaustive]
1401pub enum ReceiveECashError {
1402 #[error("The ECash is from a different federation")]
1404 WrongFederation,
1405
1406 #[error("ECash contains an uneconomical denomination")]
1408 UneconomicalDenomination,
1409
1410 #[error("Receiving ecash requires additional funds")]
1412 InsufficientFunds,
1413
1414 #[error("The ECash was already received")]
1417 AlreadyReceived,
1418
1419 #[error("The reissue transaction could not be submitted")]
1422 Failed(#[source] TransactionSubmitError),
1423}
1424
1425#[derive(Debug, Error)]
1427enum AwaitOutputError {
1428 #[error("Transaction was rejected")]
1430 Rejected,
1431
1432 #[error("Failed to finalize notes")]
1434 Failed,
1435}
1436
1437#[derive(Debug, Clone, Eq, PartialEq, Serialize, Deserialize)]
1438pub enum FinalReceiveOperationState {
1439 Success,
1441 Rejected,
1443}
1444
1445#[derive(Debug, Clone, Eq, PartialEq, Hash, Decodable, Encodable)]
1446pub enum MintClientStateMachines {
1447 Input(InputStateMachine),
1448 Output(MintOutputStateMachine),
1449 Receive(ReceiveStateMachine),
1450}
1451
1452impl IntoDynInstance for MintClientStateMachines {
1453 type DynType = DynState;
1454
1455 fn into_dyn(self, instance_id: ModuleInstanceId) -> Self::DynType {
1456 DynState::from_typed(instance_id, self)
1457 }
1458}
1459
1460impl State for MintClientStateMachines {
1461 type ModuleContext = MintClientContext;
1462
1463 fn transitions(
1464 &self,
1465 context: &Self::ModuleContext,
1466 global_context: &DynGlobalClientContext,
1467 ) -> Vec<StateTransition<Self>> {
1468 match self {
1469 MintClientStateMachines::Input(redemption_state) => {
1470 sm_enum_variant_translation!(
1471 redemption_state.transitions(context, global_context),
1472 MintClientStateMachines::Input
1473 )
1474 }
1475 MintClientStateMachines::Output(issuance_state) => {
1476 sm_enum_variant_translation!(
1477 issuance_state.transitions(context, global_context),
1478 MintClientStateMachines::Output
1479 )
1480 }
1481 MintClientStateMachines::Receive(receive_state) => {
1482 sm_enum_variant_translation!(
1483 receive_state.transitions(context, global_context),
1484 MintClientStateMachines::Receive
1485 )
1486 }
1487 }
1488 }
1489
1490 fn operation_id(&self) -> OperationId {
1491 match self {
1492 MintClientStateMachines::Input(redemption_state) => redemption_state.operation_id(),
1493 MintClientStateMachines::Output(issuance_state) => issuance_state.operation_id(),
1494 MintClientStateMachines::Receive(receive_state) => receive_state.operation_id(),
1495 }
1496 }
1497}
1498
1499fn round_to_multiple(amount: Amount, min_denomiation: Amount) -> Amount {
1500 Amount::from_msats(amount.msats.next_multiple_of(min_denomiation.msats))
1501}
1502
1503fn represent_amount_with_fees(
1504 mut remaining_amount: Amount,
1505 fee_consensus: &FeeConsensus,
1506) -> Vec<Denomination> {
1507 let mut denominations = Vec::new();
1508
1509 for denomination in client_denominations().rev() {
1511 let n_add =
1512 remaining_amount / (denomination.amount() + fee_consensus.fee(denomination.amount()));
1513
1514 denominations.extend(std::iter::repeat_n(denomination, n_add as usize));
1515
1516 remaining_amount -=
1517 n_add * (denomination.amount() + fee_consensus.fee(denomination.amount()));
1518 }
1519
1520 denominations.sort();
1522
1523 denominations
1524}
1525
1526fn represent_amount(mut remaining_amount: Amount) -> Vec<Denomination> {
1527 let mut denominations = Vec::new();
1528
1529 for denomination in client_denominations().rev() {
1531 let n_add = remaining_amount / denomination.amount();
1532
1533 denominations.extend(std::iter::repeat_n(denomination, n_add as usize));
1534
1535 remaining_amount -= n_add * denomination.amount();
1536 }
1537
1538 denominations
1539}