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 fn admin_auth(&self) -> Result<ApiAuth, MetaAdminError> {
59 self.admin_auth
60 .clone()
61 .ok_or(MetaAdminError::AdminAuthMissing)
62 }
63
64 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 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 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 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#[derive(Debug, Error)]
142#[non_exhaustive]
143pub enum MetaAdminError {
144 #[error("Admin auth not set")]
147 AdminAuthMissing,
148
149 #[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#[derive(Debug, Clone)]
185pub struct MetaClientContext {
186 pub meta_decoder: Decoder,
187}
188
189impl 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#[derive(Debug, Error)]
259enum RpcError {
260 #[error(transparent)]
262 Json(#[from] serde_json::Error),
263
264 #[error(transparent)]
266 Federation(#[from] FederationError),
267
268 #[error("deserializing consensus value as json")]
270 ConsensusValueJson(#[source] serde_json::Error),
271
272 #[error("Unknown method: {method}")]
274 UnknownMethod { method: String },
275}
276
277#[derive(Debug, Clone)]
278pub struct MetaClientInit;
279
280impl 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#[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#[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 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 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;