1use std::collections::{BTreeMap, BTreeSet};
2use std::fmt::Debug;
3use std::num::NonZeroUsize;
4use std::sync::Arc;
5
6use bitcoin::secp256k1;
7use fedimint_connectors::{DynGuaridianConnection, PeerStatus, ServerResult};
8use fedimint_core::admin_client::{GuardianConfigBackup, SetLocalParamsRequest, SetupStatus};
9use fedimint_core::backup::{BackupStatistics, ClientBackupSnapshot};
10use fedimint_core::core::backup::SignedBackupRequest;
11use fedimint_core::core::{ModuleInstanceId, ModuleKind};
12use fedimint_core::endpoint_constants::{
13 ADD_PEER_SETUP_CODE_ENDPOINT, API_ANNOUNCEMENTS_ENDPOINT, AUDIT_ENDPOINT, AUTH_ENDPOINT,
14 AWAIT_SESSION_OUTCOME_ENDPOINT, AWAIT_TRANSACTION_ENDPOINT, BACKUP_ENDPOINT,
15 BACKUP_STATISTICS_ENDPOINT, CHAIN_ID_ENDPOINT, FEDIMINTD_VERSION_ENDPOINT,
16 GET_SETUP_CODE_ENDPOINT, GUARDIAN_CONFIG_BACKUP_ENDPOINT, GUARDIAN_METADATA_ENDPOINT,
17 INVITE_CODE_ENDPOINT, RECOVER_ENDPOINT, RESET_PEER_SETUP_CODES_ENDPOINT,
18 RESTART_FEDERATION_SETUP_ENDPOINT, SESSION_COUNT_ENDPOINT, SESSION_STATUS_ENDPOINT,
19 SESSION_STATUS_V2_ENDPOINT, SET_LOCAL_PARAMS_ENDPOINT, SETUP_STATUS_ENDPOINT,
20 SHUTDOWN_ENDPOINT, SIGN_API_ANNOUNCEMENT_ENDPOINT, SIGN_GUARDIAN_METADATA_ENDPOINT,
21 START_DKG_ENDPOINT, STATUS_ENDPOINT, SUBMIT_API_ANNOUNCEMENT_ENDPOINT,
22 SUBMIT_GUARDIAN_METADATA_ENDPOINT, SUBMIT_TRANSACTION_ENDPOINT,
23};
24use fedimint_core::invite_code::InviteCode;
25use fedimint_core::module::audit::AuditSummary;
26use fedimint_core::module::registry::ModuleDecoderRegistry;
27use fedimint_core::module::{
28 ApiAuth, ApiRequestErased, ApiVersion, SerdeModuleEncoding, SerdeModuleEncodingBase64,
29};
30use fedimint_core::net::api_announcement::{
31 SignedApiAnnouncement, SignedApiAnnouncementSubmission,
32};
33use fedimint_core::session_outcome::{
34 AcceptedItem, SessionOutcome, SessionStatus, SessionStatusV2,
35};
36use fedimint_core::task::{MaybeSend, MaybeSync};
37use fedimint_core::transaction::{SerdeTransaction, Transaction, TransactionSubmissionOutcome};
38use fedimint_core::util::SafeUrl;
39use fedimint_core::{ChainId, NumPeersExt, PeerId, TransactionId, apply, async_trait_maybe_send};
40use fedimint_logging::LOG_CLIENT_NET_API;
41use futures::future::join_all;
42use futures::stream::BoxStream;
43use itertools::Itertools;
44use rand::seq::SliceRandom;
45use serde_json::Value;
46use tokio::sync::OnceCell;
47use tracing::{debug, trace};
48
49use super::super::{DynModuleApi, IGlobalFederationApi, IRawFederationApi, StatusResponse};
50use crate::api::{
51 FederationApiExt, FederationError, FederationGeneralError, FederationResult,
52 VERSION_THAT_INTRODUCED_GET_SESSION_STATUS_V2,
53};
54use crate::query::FilterMapThreshold;
55
56pub trait GlobalFederationApiWithCacheExt
59where
60 Self: Sized,
61{
62 fn with_cache(self) -> GlobalFederationApiWithCache<Self>;
63}
64
65impl<T> GlobalFederationApiWithCacheExt for T
66where
67 T: IRawFederationApi + MaybeSend + MaybeSync + 'static,
68{
69 fn with_cache(self) -> GlobalFederationApiWithCache<T> {
70 GlobalFederationApiWithCache::new(self)
71 }
72}
73
74#[derive(Debug)]
79pub struct GlobalFederationApiWithCache<T> {
80 pub(crate) inner: T,
81 pub(crate) await_session_lru:
91 Arc<tokio::sync::Mutex<lru::LruCache<u64, Arc<OnceCell<SessionOutcome>>>>>,
92
93 pub(crate) get_session_status_lru:
101 Arc<tokio::sync::Mutex<lru::LruCache<u64, Arc<OnceCell<SessionOutcome>>>>>,
102}
103
104impl<T> GlobalFederationApiWithCache<T> {
105 pub fn new(inner: T) -> GlobalFederationApiWithCache<T> {
106 Self {
107 inner,
108 await_session_lru: Arc::new(tokio::sync::Mutex::new(lru::LruCache::new(
109 NonZeroUsize::new(512).expect("is non-zero"),
110 ))),
111 get_session_status_lru: Arc::new(tokio::sync::Mutex::new(lru::LruCache::new(
112 NonZeroUsize::new(512).expect("is non-zero"),
113 ))),
114 }
115 }
116}
117
118impl<T> GlobalFederationApiWithCache<T>
119where
120 T: IRawFederationApi + MaybeSend + MaybeSync + 'static,
121{
122 pub(crate) async fn await_block_raw(
123 &self,
124 block_index: u64,
125 decoders: &ModuleDecoderRegistry,
126 ) -> FederationResult<SessionOutcome> {
127 if block_index.is_multiple_of(100) {
128 debug!(target: LOG_CLIENT_NET_API, block_index, "Awaiting block's outcome from Federation");
129 } else {
130 trace!(target: LOG_CLIENT_NET_API, block_index, "Awaiting block's outcome from Federation");
131 }
132 let response = self
133 .request_current_consensus::<SerdeModuleEncoding<SessionOutcome>>(
134 AWAIT_SESSION_OUTCOME_ENDPOINT.to_string(),
135 ApiRequestErased::new(block_index),
136 )
137 .await?;
138
139 response.try_into_inner(decoders).map_err(|e| {
140 FederationError::general(
141 AWAIT_SESSION_OUTCOME_ENDPOINT.to_string(),
142 ApiRequestErased::new(block_index),
143 FederationGeneralError::Decode(e),
144 )
145 })
146 }
147
148 pub(crate) fn select_peers_for_status(&self) -> impl Iterator<Item = PeerId> + '_ {
149 let mut peers = self.all_peers().iter().copied().collect_vec();
150 peers.shuffle(&mut rand::thread_rng());
151 peers.into_iter()
152 }
153
154 pub(crate) async fn get_session_status_raw_v2(
155 &self,
156 block_index: u64,
157 broadcast_public_keys: &BTreeMap<PeerId, secp256k1::PublicKey>,
158 decoders: &ModuleDecoderRegistry,
159 ) -> FederationResult<SessionStatus> {
160 if block_index.is_multiple_of(100) {
161 debug!(target: LOG_CLIENT_NET_API, block_index, "Get session status raw v2");
162 } else {
163 trace!(target: LOG_CLIENT_NET_API, block_index, "Get session status raw v2");
164 }
165 let params = ApiRequestErased::new(block_index);
166 let mut last_error = None;
167 for peer_id in self.select_peers_for_status() {
169 let decoded = self
170 .request_single_peer_federation::<SerdeModuleEncodingBase64<SessionStatusV2>>(
171 SESSION_STATUS_V2_ENDPOINT.to_string(),
172 params.clone(),
173 peer_id,
174 )
175 .await
176 .and_then(|s| {
177 s.try_into_inner(decoders).map_err(|e| {
178 FederationError::general(
179 SESSION_STATUS_V2_ENDPOINT.to_string(),
180 params.clone(),
181 FederationGeneralError::Decode(e),
182 )
183 })
184 });
185
186 match decoded {
187 Ok(SessionStatusV2::Complete(signed_session_outcome)) => {
188 if signed_session_outcome.verify(broadcast_public_keys, block_index) {
189 return Ok(SessionStatus::Complete(
191 signed_session_outcome.session_outcome,
192 ));
193 }
194 last_error = Some(FederationError::general(
195 SESSION_STATUS_V2_ENDPOINT.to_string(),
196 params.clone(),
197 FederationGeneralError::InvalidSignature,
198 ));
199 }
200 Ok(SessionStatusV2::Initial | SessionStatusV2::Pending(..)) => {
201 return self.get_session_status_raw(block_index, decoders).await;
203 }
204 Err(err) => {
205 last_error = Some(err);
206 }
207 }
208 assert!(last_error.is_some());
210 }
211 Err(last_error.expect("must have at least one peer"))
212 }
213
214 pub(crate) async fn get_session_status_raw(
215 &self,
216 block_index: u64,
217 decoders: &ModuleDecoderRegistry,
218 ) -> FederationResult<SessionStatus> {
219 if block_index.is_multiple_of(100) {
220 debug!(target: LOG_CLIENT_NET_API, block_index, "Get session status raw v1");
221 } else {
222 trace!(target: LOG_CLIENT_NET_API, block_index, "Get session status raw v1");
223 }
224 let response = self
225 .request_current_consensus::<SerdeModuleEncoding<SessionStatus>>(
226 SESSION_STATUS_ENDPOINT.to_string(),
227 ApiRequestErased::new(block_index),
228 )
229 .await?;
230
231 response
232 .try_into_inner(&decoders.clone().with_fallback())
233 .map_err(|e| {
234 FederationError::general(
235 SESSION_STATUS_ENDPOINT.to_string(),
236 ApiRequestErased::new(block_index),
237 FederationGeneralError::Decode(e),
238 )
239 })
240 }
241}
242
243#[apply(async_trait_maybe_send!)]
244impl<T> IRawFederationApi for GlobalFederationApiWithCache<T>
245where
246 T: IRawFederationApi + MaybeSend + MaybeSync + 'static,
247{
248 fn all_peers(&self) -> &BTreeSet<PeerId> {
249 self.inner.all_peers()
250 }
251
252 fn self_peer(&self) -> Option<PeerId> {
253 self.inner.self_peer()
254 }
255
256 fn with_module(&self, id: ModuleInstanceId) -> DynModuleApi {
257 self.inner.with_module(id)
258 }
259
260 async fn request_raw(
262 &self,
263 peer_id: PeerId,
264 method: &str,
265 params: &ApiRequestErased,
266 ) -> ServerResult<Value> {
267 self.inner.request_raw(peer_id, method, params).await
268 }
269
270 fn connection_status_stream(&self) -> BoxStream<'static, BTreeMap<PeerId, PeerStatus>> {
271 self.inner.connection_status_stream()
272 }
273
274 async fn wait_for_initialized_connections(&self) {
275 self.inner.wait_for_initialized_connections().await;
276 }
277
278 async fn get_peer_connection(&self, peer_id: PeerId) -> ServerResult<DynGuaridianConnection> {
279 self.inner.get_peer_connection(peer_id).await
280 }
281}
282
283#[apply(async_trait_maybe_send!)]
284impl<T> IGlobalFederationApi for GlobalFederationApiWithCache<T>
285where
286 T: IRawFederationApi + MaybeSend + MaybeSync + 'static,
287{
288 async fn await_block(
289 &self,
290 session_idx: u64,
291 decoders: &ModuleDecoderRegistry,
292 ) -> FederationResult<SessionOutcome> {
293 let mut lru_lock = self.await_session_lru.lock().await;
294
295 let entry_arc = lru_lock
296 .get_or_insert(session_idx, || Arc::new(OnceCell::new()))
297 .clone();
298
299 drop(lru_lock);
301
302 entry_arc
303 .get_or_try_init(|| self.await_block_raw(session_idx, decoders))
304 .await
305 .cloned()
306 }
307
308 async fn get_session_status(
309 &self,
310 session_idx: u64,
311 decoders: &ModuleDecoderRegistry,
312 core_api_version: ApiVersion,
313 broadcast_public_keys: Option<&BTreeMap<PeerId, secp256k1::PublicKey>>,
314 ) -> FederationResult<SessionStatus> {
315 enum NoCacheErr {
316 Initial,
317 Pending(Vec<AcceptedItem>),
318 Err(FederationError),
319 }
320
321 let mut lru_lock = self.get_session_status_lru.lock().await;
322
323 let entry_arc = lru_lock
324 .get_or_insert(session_idx, || Arc::new(OnceCell::new()))
325 .clone();
326
327 drop(lru_lock);
329
330 match entry_arc
331 .get_or_try_init(|| async {
332 let session_status =
333 if core_api_version < VERSION_THAT_INTRODUCED_GET_SESSION_STATUS_V2 {
334 self.get_session_status_raw(session_idx, decoders).await
335 } else if let Some(broadcast_public_keys) = broadcast_public_keys {
336 self.get_session_status_raw_v2(session_idx, broadcast_public_keys, decoders)
337 .await
338 } else {
339 self.get_session_status_raw(session_idx, decoders).await
340 };
341 match session_status {
342 Err(e) => Err(NoCacheErr::Err(e)),
343 Ok(SessionStatus::Initial) => Err(NoCacheErr::Initial),
344 Ok(SessionStatus::Pending(s)) => Err(NoCacheErr::Pending(s)),
345 Ok(SessionStatus::Complete(s)) => Ok(s),
347 }
348 })
349 .await
350 .cloned()
351 {
352 Ok(s) => Ok(SessionStatus::Complete(s)),
353 Err(NoCacheErr::Initial) => Ok(SessionStatus::Initial),
354 Err(NoCacheErr::Pending(s)) => Ok(SessionStatus::Pending(s)),
355 Err(NoCacheErr::Err(e)) => Err(e),
356 }
357 }
358
359 async fn submit_transaction(
360 &self,
361 tx: Transaction,
362 ) -> SerdeModuleEncoding<TransactionSubmissionOutcome> {
363 self.request_current_consensus_retry(
364 SUBMIT_TRANSACTION_ENDPOINT.to_owned(),
365 ApiRequestErased::new(SerdeTransaction::from(&tx)),
366 )
367 .await
368 }
369
370 async fn session_count(&self) -> FederationResult<u64> {
371 self.request_current_consensus(
372 SESSION_COUNT_ENDPOINT.to_owned(),
373 ApiRequestErased::default(),
374 )
375 .await
376 }
377
378 async fn await_transaction(&self, txid: TransactionId) -> TransactionId {
379 self.request_current_consensus_retry(
380 AWAIT_TRANSACTION_ENDPOINT.to_owned(),
381 ApiRequestErased::new(txid),
382 )
383 .await
384 }
385
386 async fn upload_backup(&self, request: &SignedBackupRequest) -> FederationResult<()> {
387 self.request_current_consensus(BACKUP_ENDPOINT.to_owned(), ApiRequestErased::new(request))
388 .await
389 }
390
391 async fn download_backup(
392 &self,
393 id: &secp256k1::PublicKey,
394 ) -> FederationResult<BTreeMap<PeerId, Option<ClientBackupSnapshot>>> {
395 self.request_with_strategy(
396 FilterMapThreshold::new(|_, snapshot| Ok(snapshot), self.all_peers().to_num_peers()),
397 RECOVER_ENDPOINT.to_owned(),
398 ApiRequestErased::new(id),
399 )
400 .await
401 }
402
403 async fn setup_status(&self, auth: ApiAuth) -> FederationResult<SetupStatus> {
404 self.request_admin(SETUP_STATUS_ENDPOINT, ApiRequestErased::default(), auth)
405 .await
406 }
407
408 async fn set_local_params(
409 &self,
410 name: String,
411 federation_name: Option<String>,
412 disable_base_fees: Option<bool>,
413 enabled_modules: Option<BTreeSet<ModuleKind>>,
414 federation_size: Option<u32>,
415 auth: ApiAuth,
416 ) -> FederationResult<String> {
417 self.request_admin(
418 SET_LOCAL_PARAMS_ENDPOINT,
419 ApiRequestErased::new(SetLocalParamsRequest {
420 name,
421 federation_name,
422 disable_base_fees,
423 enabled_modules,
424 federation_size,
425 }),
426 auth,
427 )
428 .await
429 }
430
431 async fn add_peer_connection_info(
432 &self,
433 info: String,
434 auth: ApiAuth,
435 ) -> FederationResult<String> {
436 self.request_admin(
437 ADD_PEER_SETUP_CODE_ENDPOINT,
438 ApiRequestErased::new(info),
439 auth,
440 )
441 .await
442 }
443
444 async fn reset_peer_setup_codes(&self, auth: ApiAuth) -> FederationResult<()> {
445 self.request_admin(
446 RESET_PEER_SETUP_CODES_ENDPOINT,
447 ApiRequestErased::default(),
448 auth,
449 )
450 .await
451 }
452
453 async fn get_setup_code(&self, auth: ApiAuth) -> FederationResult<Option<String>> {
454 self.request_admin(GET_SETUP_CODE_ENDPOINT, ApiRequestErased::default(), auth)
455 .await
456 }
457
458 async fn start_dkg(&self, auth: ApiAuth) -> FederationResult<()> {
459 self.request_admin(START_DKG_ENDPOINT, ApiRequestErased::default(), auth)
460 .await
461 }
462
463 async fn status(&self) -> FederationResult<StatusResponse> {
464 self.request_admin_no_auth(STATUS_ENDPOINT, ApiRequestErased::default())
465 .await
466 }
467
468 async fn audit(&self, auth: ApiAuth) -> FederationResult<AuditSummary> {
469 self.request_admin(AUDIT_ENDPOINT, ApiRequestErased::default(), auth)
470 .await
471 }
472
473 async fn guardian_config_backup(
474 &self,
475 auth: ApiAuth,
476 ) -> FederationResult<GuardianConfigBackup> {
477 self.request_admin(
478 GUARDIAN_CONFIG_BACKUP_ENDPOINT,
479 ApiRequestErased::default(),
480 auth,
481 )
482 .await
483 }
484
485 async fn auth(&self, auth: ApiAuth) -> FederationResult<()> {
486 self.request_admin(AUTH_ENDPOINT, ApiRequestErased::default(), auth)
487 .await
488 }
489
490 async fn restart_federation_setup(&self, auth: ApiAuth) -> FederationResult<()> {
491 self.request_admin(
492 RESTART_FEDERATION_SETUP_ENDPOINT,
493 ApiRequestErased::default(),
494 auth,
495 )
496 .await
497 }
498
499 async fn submit_api_announcement(
500 &self,
501 announcement_peer_id: PeerId,
502 announcement: SignedApiAnnouncement,
503 ) -> FederationResult<()> {
504 let peer_errors = join_all(self.all_peers().iter().map(|&peer_id| {
505 let announcement_inner = announcement.clone();
506 async move {
507 (
508 peer_id,
509 self.request_single_peer::<()>(
510 SUBMIT_API_ANNOUNCEMENT_ENDPOINT.into(),
511 ApiRequestErased::new(SignedApiAnnouncementSubmission {
512 signed_api_announcement: announcement_inner,
513 peer_id: announcement_peer_id,
514 }),
515 peer_id,
516 )
517 .await,
518 )
519 }
520 }))
521 .await
522 .into_iter()
523 .filter_map(|(peer_id, result)| match result {
524 Ok(()) => None,
525 Err(e) => Some((peer_id, e)),
526 })
527 .collect::<BTreeMap<_, _>>();
528
529 if peer_errors.is_empty() {
530 Ok(())
531 } else {
532 Err(FederationError {
533 method: SUBMIT_API_ANNOUNCEMENT_ENDPOINT.to_string(),
534 params: serde_json::to_value(announcement).expect("can be serialized"),
535 general: None,
536 peer_errors,
537 })
538 }
539 }
540
541 async fn api_announcements(
542 &self,
543 guardian: PeerId,
544 ) -> ServerResult<BTreeMap<PeerId, SignedApiAnnouncement>> {
545 self.request_single_peer(
546 API_ANNOUNCEMENTS_ENDPOINT.to_owned(),
547 ApiRequestErased::default(),
548 guardian,
549 )
550 .await
551 }
552
553 async fn sign_api_announcement(
554 &self,
555 api_url: SafeUrl,
556 auth: ApiAuth,
557 ) -> FederationResult<SignedApiAnnouncement> {
558 self.request_admin(
559 SIGN_API_ANNOUNCEMENT_ENDPOINT,
560 ApiRequestErased::new(api_url),
561 auth,
562 )
563 .await
564 }
565
566 async fn submit_guardian_metadata(
567 &self,
568 announcement_peer_id: PeerId,
569 metadata: fedimint_core::net::guardian_metadata::SignedGuardianMetadata,
570 ) -> FederationResult<()> {
571 use fedimint_core::net::guardian_metadata::SignedGuardianMetadataSubmission;
572 let peer_errors = join_all(self.all_peers().iter().map(|&peer_id| {
573 let metadata_inner = metadata.clone();
574 async move {
575 (
576 peer_id,
577 self.request_single_peer::<()>(
578 SUBMIT_GUARDIAN_METADATA_ENDPOINT.into(),
579 ApiRequestErased::new(SignedGuardianMetadataSubmission {
580 signed_guardian_metadata: metadata_inner,
581 peer_id: announcement_peer_id,
582 }),
583 peer_id,
584 )
585 .await,
586 )
587 }
588 }))
589 .await
590 .into_iter()
591 .filter_map(|(peer_id, result)| match result {
592 Ok(()) => None,
593 Err(e) => Some((peer_id, e)),
594 })
595 .collect::<BTreeMap<_, _>>();
596
597 if peer_errors.is_empty() {
598 Ok(())
599 } else {
600 Err(FederationError {
601 method: SUBMIT_GUARDIAN_METADATA_ENDPOINT.to_string(),
602 params: serde_json::to_value(&metadata).expect("can be serialized"),
603 general: None,
604 peer_errors,
605 })
606 }
607 }
608
609 async fn guardian_metadata(
610 &self,
611 guardian: PeerId,
612 ) -> ServerResult<BTreeMap<PeerId, fedimint_core::net::guardian_metadata::SignedGuardianMetadata>>
613 {
614 self.request_single_peer(
615 GUARDIAN_METADATA_ENDPOINT.to_owned(),
616 ApiRequestErased::default(),
617 guardian,
618 )
619 .await
620 }
621
622 async fn sign_guardian_metadata(
623 &self,
624 metadata: fedimint_core::net::guardian_metadata::GuardianMetadata,
625 auth: ApiAuth,
626 ) -> FederationResult<fedimint_core::net::guardian_metadata::SignedGuardianMetadata> {
627 self.request_admin(
628 SIGN_GUARDIAN_METADATA_ENDPOINT,
629 ApiRequestErased::new(metadata),
630 auth,
631 )
632 .await
633 }
634
635 async fn shutdown(&self, session: Option<u64>, auth: ApiAuth) -> FederationResult<()> {
636 self.request_admin(SHUTDOWN_ENDPOINT, ApiRequestErased::new(session), auth)
637 .await
638 }
639
640 async fn backup_statistics(&self, auth: ApiAuth) -> FederationResult<BackupStatistics> {
641 self.request_admin(
642 BACKUP_STATISTICS_ENDPOINT,
643 ApiRequestErased::default(),
644 auth,
645 )
646 .await
647 }
648
649 async fn fedimintd_version(&self, peer_id: PeerId) -> ServerResult<String> {
650 self.request_single_peer(
651 FEDIMINTD_VERSION_ENDPOINT.to_owned(),
652 ApiRequestErased::default(),
653 peer_id,
654 )
655 .await
656 }
657
658 async fn get_invite_code(&self, guardian: PeerId) -> ServerResult<InviteCode> {
659 self.request_single_peer(
660 INVITE_CODE_ENDPOINT.to_owned(),
661 ApiRequestErased::default(),
662 guardian,
663 )
664 .await
665 }
666
667 async fn chain_id(&self) -> FederationResult<ChainId> {
668 self.request_current_consensus(CHAIN_ID_ENDPOINT.to_owned(), ApiRequestErased::default())
669 .await
670 }
671}