Skip to main content

fedimint_meta_client/
lib.rs

1#![deny(clippy::pedantic)]
2#![allow(clippy::missing_errors_doc)]
3#![allow(clippy::module_name_repetitions)]
4
5pub mod api;
6#[cfg(feature = "cli")]
7pub mod cli;
8pub mod db;
9pub mod states;
10
11use std::collections::BTreeMap;
12use std::time::Duration;
13
14use api::MetaFederationApi;
15use common::{KIND, MetaConsensusValue, MetaKey, MetaValue};
16use db::DbKeyPrefix;
17use fedimint_api_client::api::{DynGlobalApi, DynModuleApi, FederationError};
18use fedimint_client_module::db::ClientModuleMigrationFn;
19use fedimint_client_module::error::{ClientModuleError, MetaFetchError};
20use fedimint_client_module::meta::{FetchKind, LegacyMetaSource, MetaSource, MetaValues};
21use fedimint_client_module::module::init::{ClientModuleInit, ClientModuleInitArgs};
22use fedimint_client_module::module::recovery::NoModuleBackup;
23use fedimint_client_module::module::{ClientModule, IClientModule};
24use fedimint_client_module::sm::Context;
25use fedimint_core::config::ClientConfig;
26use fedimint_core::core::{Decoder, ModuleKind};
27use fedimint_core::db::{DatabaseTransaction, DatabaseVersion};
28use fedimint_core::module::{
29    Amounts, ApiAuth, ApiVersion, ModuleCommon, ModuleInit, MultiApiVersion,
30};
31use fedimint_core::util::backoff_util::FibonacciBackoff;
32use fedimint_core::util::{BoxStream, backoff_util, retry};
33use fedimint_core::{PeerId, apply, async_trait_maybe_send};
34use fedimint_logging::LOG_CLIENT_MODULE_META;
35pub use fedimint_meta_common as common;
36use fedimint_meta_common::{DEFAULT_META_KEY, MetaCommonInit, MetaModuleTypes};
37use futures::{TryStreamExt as _, stream};
38use serde::Deserialize;
39use serde_json::json;
40use states::MetaStateMachine;
41use strum::IntoEnumIterator;
42use thiserror::Error;
43use tracing::{debug, warn};
44
45#[derive(Debug)]
46pub struct MetaClientModule {
47    module_api: DynModuleApi,
48    admin_auth: Option<ApiAuth>,
49}
50
51impl MetaClientModule {
52    /// The admin credentials this client was built with.
53    ///
54    /// # Errors
55    ///
56    /// Fails with [`MetaAdminError::AdminAuthMissing`] if the client was built
57    /// without them.
58    fn admin_auth(&self) -> Result<ApiAuth, MetaAdminError> {
59        self.admin_auth
60            .clone()
61            .ok_or(MetaAdminError::AdminAuthMissing)
62    }
63
64    /// Submit a meta consensus value
65    ///
66    /// When *threshold* amount of peers submits the exact same value it
67    /// becomes a new consensus value.
68    ///
69    /// To "cancel" previous vote, peer can submit a value equal to the current
70    /// consensus value.
71    ///
72    /// # Errors
73    ///
74    /// Fails with [`MetaAdminError::AdminAuthMissing`] if this client holds no
75    /// admin credentials, and with [`MetaAdminError::Federation`] if the
76    /// guardians could not be asked to record the vote.
77    pub async fn submit(&self, key: MetaKey, value: MetaValue) -> Result<(), MetaAdminError> {
78        self.module_api
79            .submit(key, value, self.admin_auth()?)
80            .await?;
81
82        Ok(())
83    }
84
85    /// Get the current meta consensus value along with it's revision
86    ///
87    /// See [`Self::get_consensus_value_rev`] to use when checking for updates.
88    ///
89    /// # Errors
90    ///
91    /// Fails with a [`FederationError`] if the federation could not be asked
92    /// for the value.
93    pub async fn get_consensus_value(
94        &self,
95        key: MetaKey,
96    ) -> Result<Option<MetaConsensusValue>, FederationError> {
97        self.module_api.get_consensus(key).await
98    }
99
100    /// Get the current meta consensus value revision
101    ///
102    /// Each time a meta consensus value changes, the revision increases,
103    /// so checking just the revision can save a lot of bandwidth in periodic
104    /// checks.
105    ///
106    /// # Errors
107    ///
108    /// Fails with a [`FederationError`] if the federation could not be asked
109    /// for the revision.
110    pub async fn get_consensus_value_rev(
111        &self,
112        key: MetaKey,
113    ) -> Result<Option<u64>, FederationError> {
114        self.module_api.get_consensus_rev(key).await
115    }
116
117    /// Get current submissions to change the meta consensus value.
118    ///
119    /// Upon changing the consensus
120    ///
121    /// # Errors
122    ///
123    /// Fails with [`MetaAdminError::AdminAuthMissing`] if this client holds no
124    /// admin credentials, and with [`MetaAdminError::Federation`] if the
125    /// guardians could not be asked for their submissions.
126    pub async fn get_submissions(
127        &self,
128        key: MetaKey,
129    ) -> Result<BTreeMap<PeerId, MetaValue>, MetaAdminError> {
130        Ok(self
131            .module_api
132            .get_submissions(key, self.admin_auth()?)
133            .await?)
134    }
135}
136
137/// A failure of one of the meta client's guardian-only operations.
138///
139/// These endpoints can fail before any request leaves the client, when it
140/// holds no admin credentials, as well as while talking to the federation.
141#[derive(Debug, Error)]
142#[non_exhaustive]
143pub enum MetaAdminError {
144    /// This client was not built with admin credentials, so it cannot call a
145    /// guardian-only endpoint.
146    #[error("Admin auth not set")]
147    AdminAuthMissing,
148
149    /// The federation rejected the request, or could not be reached.
150    #[error("The federation request failed")]
151    Federation(#[source] Box<FederationError>),
152}
153
154impl From<FederationError> for MetaAdminError {
155    fn from(source: FederationError) -> Self {
156        Self::Federation(Box::new(source))
157    }
158}
159
160#[derive(Debug, Deserialize)]
161struct GetConsensusValueRequest {
162    key: MetaKey,
163}
164
165fn format_rpc_consensus_value_response(
166    maybe_consensus_value: Option<MetaConsensusValue>,
167) -> Result<serde_json::Value, RpcError> {
168    Ok(match maybe_consensus_value {
169        Some(MetaConsensusValue { revision, value }) => {
170            let value = value
171                .to_json_lossy()
172                .map_err(RpcError::ConsensusValueJson)?;
173
174            json!({
175                "revision": revision,
176                "value": value,
177            })
178        }
179        None => serde_json::Value::Null,
180    })
181}
182
183/// Data needed by the state machine
184#[derive(Debug, Clone)]
185pub struct MetaClientContext {
186    pub meta_decoder: Decoder,
187}
188
189// TODO: Boiler-plate
190impl Context for MetaClientContext {
191    const KIND: Option<ModuleKind> = Some(KIND);
192}
193
194#[apply(async_trait_maybe_send!)]
195impl ClientModule for MetaClientModule {
196    type Init = MetaClientInit;
197    type Common = MetaModuleTypes;
198    type Backup = NoModuleBackup;
199    type ModuleStateMachineContext = MetaClientContext;
200    type States = MetaStateMachine;
201
202    fn context(&self) -> Self::ModuleStateMachineContext {
203        MetaClientContext {
204            meta_decoder: self.decoder(),
205        }
206    }
207
208    fn input_fee(
209        &self,
210        _amount: &Amounts,
211        _input: &<Self::Common as ModuleCommon>::Input,
212    ) -> Option<Amounts> {
213        unreachable!()
214    }
215
216    fn output_fee(
217        &self,
218        _amount: &Amounts,
219        _output: &<Self::Common as ModuleCommon>::Output,
220    ) -> Option<Amounts> {
221        unreachable!()
222    }
223
224    async fn handle_rpc(
225        &self,
226        method: String,
227        request: serde_json::Value,
228    ) -> BoxStream<'_, Result<serde_json::Value, ClientModuleError>> {
229        Box::pin(
230            stream::once(async move {
231                match method.as_str() {
232                    "get_consensus_value" => {
233                        let req: GetConsensusValueRequest = serde_json::from_value(request)?;
234                        let maybe_consensus_value = self.get_consensus_value(req.key).await?;
235                        format_rpc_consensus_value_response(maybe_consensus_value)
236                    }
237                    _ => Err(RpcError::UnknownMethod {
238                        method: method.clone(),
239                    }),
240                }
241            })
242            .map_err(ClientModuleError::other),
243        )
244    }
245
246    #[cfg(feature = "cli")]
247    async fn handle_cli_command(
248        &self,
249        args: &[std::ffi::OsString],
250    ) -> Result<serde_json::Value, ClientModuleError> {
251        cli::handle_cli_command(self, args)
252            .await
253            .map_err(ClientModuleError::other)
254    }
255}
256
257/// A failure of a meta module RPC request.
258#[derive(Debug, Error)]
259enum RpcError {
260    /// The request's parameters do not fit the method.
261    #[error(transparent)]
262    Json(#[from] serde_json::Error),
263
264    /// The federation did not serve the request.
265    #[error(transparent)]
266    Federation(#[from] FederationError),
267
268    /// The consensus value is not valid JSON.
269    #[error("deserializing consensus value as json")]
270    ConsensusValueJson(#[source] serde_json::Error),
271
272    /// The request names a method the module does not have.
273    #[error("Unknown method: {method}")]
274    UnknownMethod { method: String },
275}
276
277#[derive(Debug, Clone)]
278pub struct MetaClientInit;
279
280// TODO: Boilerplate-code
281impl ModuleInit for MetaClientInit {
282    type Common = MetaCommonInit;
283
284    async fn dump_database(
285        &self,
286        _dbtx: &mut DatabaseTransaction<'_>,
287        prefix_names: Vec<String>,
288    ) -> Box<dyn Iterator<Item = (String, Box<dyn erased_serde::Serialize + Send>)> + '_> {
289        let items: BTreeMap<String, Box<dyn erased_serde::Serialize + Send>> = BTreeMap::new();
290        let filtered_prefixes = DbKeyPrefix::iter().filter(|f| {
291            prefix_names.is_empty() || prefix_names.contains(&f.to_string().to_lowercase())
292        });
293
294        #[allow(clippy::never_loop)]
295        for table in filtered_prefixes {
296            match table {}
297        }
298
299        Box::new(items.into_iter())
300    }
301}
302
303/// Generates the client module
304#[apply(async_trait_maybe_send!)]
305impl ClientModuleInit for MetaClientInit {
306    type Module = MetaClientModule;
307
308    fn supported_api_versions(&self) -> MultiApiVersion {
309        MultiApiVersion::try_from_iter([ApiVersion { major: 0, minor: 0 }])
310            .expect("no version conflicts")
311    }
312
313    async fn init(
314        &self,
315        args: &ClientModuleInitArgs<Self>,
316    ) -> Result<Self::Module, ClientModuleError> {
317        Ok(MetaClientModule {
318            module_api: args.module_api().clone(),
319            admin_auth: args.admin_auth().cloned(),
320        })
321    }
322
323    fn get_database_migrations(&self) -> BTreeMap<DatabaseVersion, ClientModuleMigrationFn> {
324        BTreeMap::new()
325    }
326}
327
328/// Meta source fetching meta values from the meta module if available or the
329/// legacy meta source otherwise.
330#[derive(Clone, Debug, Default)]
331pub struct MetaModuleMetaSourceWithFallback<S = LegacyMetaSource> {
332    legacy: S,
333}
334
335impl<S> MetaModuleMetaSourceWithFallback<S> {
336    pub fn new(legacy: S) -> Self {
337        Self { legacy }
338    }
339}
340
341#[apply(async_trait_maybe_send!)]
342impl<S: MetaSource> MetaSource for MetaModuleMetaSourceWithFallback<S> {
343    async fn wait_for_update(&self) {
344        fedimint_core::runtime::sleep(Duration::from_mins(10)).await;
345    }
346
347    async fn fetch(
348        &self,
349        client_config: &ClientConfig,
350        api: &DynGlobalApi,
351        fetch_kind: fedimint_client_module::meta::FetchKind,
352        last_revision: Option<u64>,
353    ) -> Result<fedimint_client_module::meta::MetaValues, MetaFetchError> {
354        let backoff = match fetch_kind {
355            // need to be fast the first time.
356            FetchKind::Initial => backoff_util::aggressive_backoff(),
357            FetchKind::Background => backoff_util::background_backoff(),
358        };
359
360        let maybe_meta_module_meta = get_meta_module_value(client_config, api, backoff)
361            .await
362            .map(|meta| {
363                Result::<_, MetaFetchError>::Ok(MetaValues {
364                    values: serde_json::from_slice(meta.value.as_slice())?,
365                    revision: meta.revision,
366                })
367            })
368            .transpose()?;
369
370        // If we couldn't fetch valid meta values from the meta module for any reason,
371        // fall back to the legacy meta source
372        if let Some(maybe_meta_module_meta) = maybe_meta_module_meta {
373            Ok(maybe_meta_module_meta)
374        } else {
375            self.legacy
376                .fetch(client_config, api, fetch_kind, last_revision)
377                .await
378        }
379    }
380}
381
382async fn get_meta_module_value(
383    client_config: &ClientConfig,
384    api: &DynGlobalApi,
385    backoff: FibonacciBackoff,
386) -> Option<MetaConsensusValue> {
387    match client_config.get_first_module_by_kind_cfg(KIND) {
388        Ok((instance_id, _)) => {
389            let meta_api = api.with_module(instance_id);
390
391            let overrides_res = retry("fetch_meta_values", backoff, || async {
392                meta_api.get_consensus(DEFAULT_META_KEY).await
393            })
394            .await;
395
396            match overrides_res {
397                Ok(Some(consensus)) => Some(consensus),
398                Ok(None) => {
399                    debug!(target: LOG_CLIENT_MODULE_META, "Meta module returned no consensus value");
400                    None
401                }
402                Err(e) => {
403                    warn!(target: LOG_CLIENT_MODULE_META, "Failed to fetch meta module consensus value: {}", e);
404                    None
405                }
406            }
407        }
408        _ => None,
409    }
410}
411
412#[cfg(test)]
413mod tests;