1use async_trait::async_trait;
2use slog::{Logger, debug, info, warn};
3use std::{collections::BTreeMap, io, ops::RangeInclusive, path::Path, sync::Arc};
4
5use mithril_common::crypto_helper::{MKTree, MKTreeStoreInMemory};
6use mithril_common::entities::{CardanoDbBeacon, HexEncodedDigest, ImmutableFileNumber};
7use mithril_common::logging::LoggerExtensions;
8
9use crate::{
10 digesters::{
11 ImmutableDigester, ImmutableDigesterError, cache::ImmutableFileDigestCacheProvider,
12 },
13 entities::ImmutableFile,
14};
15
16use super::immutable_digester::ComputedImmutablesDigests;
17
18pub struct CardanoImmutableDigester {
20 cache_provider: Option<Arc<dyn ImmutableFileDigestCacheProvider>>,
22
23 logger: Logger,
25}
26
27impl CardanoImmutableDigester {
28 pub fn new(
30 cache_provider: Option<Arc<dyn ImmutableFileDigestCacheProvider>>,
31 logger: Logger,
32 ) -> Self {
33 Self {
34 cache_provider,
35 logger: logger.new_with_component_name::<Self>(),
36 }
37 }
38
39 async fn process_immutables(
40 &self,
41 immutables: Vec<ImmutableFile>,
42 ) -> Result<ComputedImmutablesDigests, ImmutableDigesterError> {
43 let cached_values = self.fetch_immutables_cached(immutables).await;
44
45 let logger = self.logger.clone();
47 let computed_digests =
48 tokio::task::spawn_blocking(move || -> Result<ComputedImmutablesDigests, io::Error> {
49 ComputedImmutablesDigests::compute_immutables_digests(cached_values, logger)
50 })
51 .await
52 .map_err(|e| ImmutableDigesterError::DigestComputationError(e.into()))??;
53
54 Ok(computed_digests)
55 }
56
57 async fn fetch_immutables_cached(
58 &self,
59 immutables: Vec<ImmutableFile>,
60 ) -> BTreeMap<ImmutableFile, Option<String>> {
61 match self.cache_provider.as_ref() {
62 None => BTreeMap::from_iter(immutables.into_iter().map(|i| (i, None))),
63 Some(cache_provider) => match cache_provider.get(immutables.clone()).await {
64 Ok(values) => values,
65 Err(error) => {
66 warn!(
67 self.logger, "Error while getting cached immutable files digests";
68 "error" => ?error
69 );
70 BTreeMap::from_iter(immutables.into_iter().map(|i| (i, None)))
71 }
72 },
73 }
74 }
75
76 async fn update_cache(&self, computed_immutables_digests: &ComputedImmutablesDigests) {
77 if let Some(cache_provider) = self.cache_provider.as_ref() {
78 let new_cached_entries = computed_immutables_digests
79 .entries
80 .iter()
81 .filter(|(file, _hash)| {
82 computed_immutables_digests
83 .new_cached_entries
84 .contains(&file.filename)
85 })
86 .map(|(file, hash)| (file.filename.clone(), hash.clone()))
87 .collect();
88
89 if let Err(error) = cache_provider.store(new_cached_entries).await {
90 warn!(
91 self.logger, "Error while storing new immutable files digests to cache";
92 "error" => ?error
93 );
94 }
95 }
96 }
97}
98
99#[async_trait]
100impl ImmutableDigester for CardanoImmutableDigester {
101 async fn compute_digests_for_range(
102 &self,
103 dirpath: &Path,
104 range: &RangeInclusive<ImmutableFileNumber>,
105 ) -> Result<ComputedImmutablesDigests, ImmutableDigesterError> {
106 let immutables_to_process = list_immutable_files_to_process_for_range(dirpath, range)?;
107 info!(self.logger, ">> compute_digests_for_range"; "nb_of_immutables" => immutables_to_process.len());
108 let computed_immutables_digests = self.process_immutables(immutables_to_process).await?;
109
110 self.update_cache(&computed_immutables_digests).await;
111
112 debug!(
113 self.logger,
114 "Successfully computed Digests for Cardano database"; "range" => #?range);
115
116 Ok(computed_immutables_digests)
117 }
118
119 async fn compute_merkle_tree(
120 &self,
121 dirpath: &Path,
122 beacon: &CardanoDbBeacon,
123 ) -> Result<MKTree<MKTreeStoreInMemory>, ImmutableDigesterError> {
124 let immutables_to_process =
125 list_immutable_files_to_process(dirpath, beacon.immutable_file_number)?;
126 info!(self.logger, ">> compute_merkle_tree"; "beacon" => #?beacon, "nb_of_immutables" => immutables_to_process.len());
127 let computed_immutables_digests = self.process_immutables(immutables_to_process).await?;
128
129 self.update_cache(&computed_immutables_digests).await;
130
131 let digests: Vec<HexEncodedDigest> =
132 computed_immutables_digests.entries.into_values().collect();
133 let mktree =
134 MKTree::new(&digests).map_err(ImmutableDigesterError::MerkleTreeComputationError)?;
135
136 debug!(
137 self.logger,
138 "Successfully computed Merkle tree for Cardano database"; "beacon" => #?beacon);
139
140 Ok(mktree)
141 }
142}
143
144fn list_immutable_files_to_process(
145 dirpath: &Path,
146 up_to_file_number: ImmutableFileNumber,
147) -> Result<Vec<ImmutableFile>, ImmutableDigesterError> {
148 let immutables: Vec<ImmutableFile> = ImmutableFile::list_all_in_dir(dirpath)?
149 .into_iter()
150 .filter(|f| f.number <= up_to_file_number)
151 .collect();
152
153 match immutables.last() {
154 None => Err(ImmutableDigesterError::NotEnoughImmutable {
155 expected_number: up_to_file_number,
156 found_number: None,
157 db_dir: dirpath.to_owned(),
158 }),
159 Some(last_immutable_file) if last_immutable_file.number < up_to_file_number => {
160 Err(ImmutableDigesterError::NotEnoughImmutable {
161 expected_number: up_to_file_number,
162 found_number: Some(last_immutable_file.number),
163 db_dir: dirpath.to_owned(),
164 })
165 }
166 Some(_) => Ok(immutables),
167 }
168}
169
170fn list_immutable_files_to_process_for_range(
171 dirpath: &Path,
172 range: &RangeInclusive<ImmutableFileNumber>,
173) -> Result<Vec<ImmutableFile>, ImmutableDigesterError> {
174 let immutables: Vec<ImmutableFile> = ImmutableFile::list_all_in_dir(dirpath)?
175 .into_iter()
176 .filter(|f| range.contains(&f.number))
177 .collect();
178
179 Ok(immutables)
180}
181
182#[cfg(test)]
183mod tests {
184 use sha2::Sha256;
185 use std::{collections::BTreeMap, io, sync::Arc};
186
187 use crate::digesters::cache::{
188 ImmutableDigesterCacheGetError, ImmutableDigesterCacheProviderError,
189 ImmutableDigesterCacheStoreError, MemoryImmutableFileDigestCacheProvider,
190 MockImmutableFileDigestCacheProvider,
191 };
192 use crate::test::{DummyCardanoDbBuilder, TestLogger};
193
194 use super::*;
195
196 fn db_builder(dir_name: &str) -> DummyCardanoDbBuilder {
197 DummyCardanoDbBuilder::new(&format!("cardano_immutable_digester/{dir_name}"))
198 }
199
200 #[tokio::test]
201 async fn fail_if_no_file_in_folder() {
202 let cardano_db = db_builder("fail_if_no_file_in_folder").build();
203
204 let result = list_immutable_files_to_process(cardano_db.get_immutable_dir(), 1)
205 .expect_err("list_immutable_files_to_process should have failed");
206
207 assert_eq!(
208 format!(
209 "{:?}",
210 ImmutableDigesterError::NotEnoughImmutable {
211 expected_number: 1,
212 found_number: None,
213 db_dir: cardano_db.get_immutable_dir().to_path_buf(),
214 }
215 ),
216 format!("{result:?}")
217 );
218 }
219
220 #[tokio::test]
221 async fn fail_if_a_invalid_file_is_in_immutable_folder() {
222 let cardano_db = db_builder("fail_if_no_immutable_exist")
223 .with_non_immutables(&["not_immutable"])
224 .build();
225
226 assert!(list_immutable_files_to_process(cardano_db.get_immutable_dir(), 1).is_err());
227 }
228
229 #[tokio::test]
230 async fn can_list_files_to_process_even_if_theres_only_the_uncompleted_immutable_trio() {
231 let cardano_db = db_builder(
232 "can_list_files_to_process_even_if_theres_only_the_uncompleted_immutable_trio",
233 )
234 .with_immutables(&[1])
235 .build();
236
237 let processable_files =
238 list_immutable_files_to_process(cardano_db.get_immutable_dir(), 1).unwrap();
239
240 assert_eq!(
241 vec![
242 "00001.chunk".to_string(),
243 "00001.primary".to_string(),
244 "00001.secondary".to_string()
245 ],
246 processable_files.into_iter().map(|f| f.filename).collect::<Vec<_>>()
247 );
248 }
249
250 #[tokio::test]
251 async fn fail_if_less_immutable_than_what_required_in_beacon() {
252 let cardano_db = db_builder("fail_if_less_immutable_than_what_required_in_beacon")
253 .with_immutables(&[1, 2, 3, 4, 5])
254 .append_immutable_trio()
255 .build();
256
257 let result = list_immutable_files_to_process(cardano_db.get_immutable_dir(), 10)
258 .expect_err("list_immutable_files_to_process should've failed");
259
260 assert_eq!(
261 format!(
262 "{:?}",
263 ImmutableDigesterError::NotEnoughImmutable {
264 expected_number: 10,
265 found_number: Some(6),
266 db_dir: cardano_db.get_immutable_dir().to_path_buf(),
267 }
268 ),
269 format!("{result:?}")
270 );
271 }
272
273 #[tokio::test]
274 async fn can_compute_merkle_tree_of_a_hundred_immutable_file_trio() {
275 let cardano_db = db_builder("can_compute_merkle_tree_of_a_hundred_immutable_file_trio")
276 .with_immutables(&(1..=100).collect::<Vec<ImmutableFileNumber>>())
277 .append_immutable_trio()
278 .build();
279 let logger = TestLogger::stdout();
280 let digester = CardanoImmutableDigester::new(
281 Some(Arc::new(MemoryImmutableFileDigestCacheProvider::default())),
282 logger.clone(),
283 );
284 let beacon = CardanoDbBeacon::new(1, 100);
285
286 let result = digester
287 .compute_merkle_tree(cardano_db.get_immutable_dir(), &beacon)
288 .await
289 .expect("compute_merkle_tree must not fail");
290
291 let expected_merkle_root = result.compute_root().unwrap().to_hex();
292
293 assert_eq!(
294 "8552f75838176c967a33eb6da1fe5f3c9940b706d75a9c2352c0acd8439f3d84".to_string(),
295 expected_merkle_root
296 )
297 }
298
299 #[tokio::test]
300 async fn can_compute_digests_for_range_of_a_hundred_immutable_file_trio() {
301 let immutable_range = 1..=100;
302 let cardano_db =
303 db_builder("can_compute_digests_for_range_of_a_hundred_immutable_file_trio")
304 .with_immutables(&immutable_range.clone().collect::<Vec<ImmutableFileNumber>>())
305 .append_immutable_trio()
306 .build();
307 let logger = TestLogger::stdout();
308 let digester = CardanoImmutableDigester::new(
309 Some(Arc::new(MemoryImmutableFileDigestCacheProvider::default())),
310 logger.clone(),
311 );
312
313 let result = digester
314 .compute_digests_for_range(cardano_db.get_immutable_dir(), &immutable_range)
315 .await
316 .expect("compute_digests_for_range must not fail");
317
318 assert_eq!(cardano_db.get_immutable_files().len(), result.entries.len())
319 }
320
321 #[tokio::test]
322 async fn can_compute_consistent_digests_for_range() {
323 let immutable_range = 1..=1;
324 let cardano_db = db_builder("can_compute_digests_for_range_consistently")
325 .with_immutables(&immutable_range.clone().collect::<Vec<ImmutableFileNumber>>())
326 .append_immutable_trio()
327 .build();
328 let logger = TestLogger::stdout();
329 let digester = CardanoImmutableDigester::new(
330 Some(Arc::new(MemoryImmutableFileDigestCacheProvider::default())),
331 logger.clone(),
332 );
333
334 let result = digester
335 .compute_digests_for_range(cardano_db.get_immutable_dir(), &immutable_range)
336 .await
337 .expect("compute_digests_for_range must not fail");
338
339 assert_eq!(
340 BTreeMap::from([
341 (
342 ImmutableFile {
343 path: cardano_db.get_immutable_dir().join("00001.chunk"),
344 number: 1,
345 filename: "00001.chunk".to_string()
346 },
347 "faebbf47077f68ef57219396ff69edc738978a3eca946ac7df1983dbf11364ec".to_string()
348 ),
349 (
350 ImmutableFile {
351 path: cardano_db.get_immutable_dir().join("00001.primary"),
352 number: 1,
353 filename: "00001.primary".to_string()
354 },
355 "f11bdb991fc7e72970be7d7f666e10333f92c14326d796fed8c2c041675fa826".to_string()
356 ),
357 (
358 ImmutableFile {
359 path: cardano_db.get_immutable_dir().join("00001.secondary"),
360 number: 1,
361 filename: "00001.secondary".to_string()
362 },
363 "b139684b968fa12ce324cce464d000de0e2c2ded0fd3e473a666410821d3fde3".to_string()
364 )
365 ]),
366 result.entries
367 );
368 }
369
370 #[tokio::test]
371 async fn compute_merkle_tree_store_digests_into_cache_provider() {
372 let cardano_db = db_builder("compute_merkle_tree_store_digests_into_cache_provider")
373 .with_immutables(&[1, 2])
374 .append_immutable_trio()
375 .build();
376 let immutables = cardano_db.get_immutable_files().clone();
377 let cache = Arc::new(MemoryImmutableFileDigestCacheProvider::default());
378 let logger = TestLogger::stdout();
379 let digester = CardanoImmutableDigester::new(Some(cache.clone()), logger.clone());
380 let beacon = CardanoDbBeacon::new(1, 2);
381
382 digester
383 .compute_merkle_tree(cardano_db.get_immutable_dir(), &beacon)
384 .await
385 .expect("compute_digest must not fail");
386
387 let cached_entries = cache
388 .get(immutables.clone())
389 .await
390 .expect("Cache read should not fail");
391 let expected: BTreeMap<_, _> = immutables
392 .into_iter()
393 .map(|i| {
394 let digest = hex::encode(i.compute_raw_hash::<Sha256>().unwrap());
395 (i, Some(digest))
396 })
397 .collect();
398
399 assert_eq!(expected, cached_entries);
400 }
401
402 #[tokio::test]
403 async fn compute_digests_for_range_stores_digests_into_cache_provider() {
404 let cardano_db = db_builder("compute_digests_for_range_stores_digests_into_cache_provider")
405 .with_immutables(&[1, 2])
406 .append_immutable_trio()
407 .build();
408 let immutables = cardano_db.get_immutable_files().clone();
409 let cache = Arc::new(MemoryImmutableFileDigestCacheProvider::default());
410 let logger = TestLogger::stdout();
411 let digester = CardanoImmutableDigester::new(Some(cache.clone()), logger.clone());
412 let immutable_range = 1..=2;
413
414 digester
415 .compute_digests_for_range(cardano_db.get_immutable_dir(), &immutable_range)
416 .await
417 .expect("compute_digests_for_range must not fail");
418
419 let cached_entries = cache
420 .get(immutables.clone())
421 .await
422 .expect("Cache read should not fail");
423 let expected: BTreeMap<_, _> = immutables
424 .into_iter()
425 .filter(|i| immutable_range.contains(&i.number))
426 .map(|i| {
427 let digest = hex::encode(i.compute_raw_hash::<Sha256>().unwrap());
428 (i.to_owned(), Some(digest))
429 })
430 .collect();
431
432 assert_eq!(expected, cached_entries);
433 }
434
435 #[tokio::test]
436 async fn computed_merkle_tree_with_cold_or_hot_or_without_any_cache_are_equals() {
437 let cardano_db = DummyCardanoDbBuilder::new(
438 "computed_merkle_tree_with_cold_or_hot_or_without_any_cache_are_equals",
439 )
440 .with_immutables(&[1, 2, 3])
441 .append_immutable_trio()
442 .build();
443 let logger = TestLogger::stdout();
444 let no_cache_digester = CardanoImmutableDigester::new(None, logger.clone());
445 let cache_digester = CardanoImmutableDigester::new(
446 Some(Arc::new(MemoryImmutableFileDigestCacheProvider::default())),
447 logger.clone(),
448 );
449 let beacon = CardanoDbBeacon::new(1, 3);
450
451 let without_cache_digest = no_cache_digester
452 .compute_merkle_tree(cardano_db.get_immutable_dir(), &beacon)
453 .await
454 .expect("compute_merkle_tree must not fail");
455
456 let cold_cache_digest = cache_digester
457 .compute_merkle_tree(cardano_db.get_immutable_dir(), &beacon)
458 .await
459 .expect("compute_merkle_tree must not fail");
460
461 let full_cache_digest = cache_digester
462 .compute_merkle_tree(cardano_db.get_immutable_dir(), &beacon)
463 .await
464 .expect("compute_merkle_tree must not fail");
465
466 let without_cache_merkle_root = without_cache_digest.compute_root().unwrap();
467 let cold_cache_merkle_root = cold_cache_digest.compute_root().unwrap();
468 let full_cache_merkle_root = full_cache_digest.compute_root().unwrap();
469 assert_eq!(
470 without_cache_merkle_root, full_cache_merkle_root,
471 "Merkle roots with or without cache should be the same"
472 );
473
474 assert_eq!(
475 cold_cache_merkle_root, full_cache_merkle_root,
476 "Merkle roots with cold or with hot cache should be the same"
477 );
478 }
479
480 #[tokio::test]
481 async fn computed_digests_for_range_with_cold_or_hot_or_without_any_cache_are_equals() {
482 let cardano_db = DummyCardanoDbBuilder::new(
483 "computed_digests_for_range_with_cold_or_hot_or_without_any_cache_are_equals",
484 )
485 .with_immutables(&[1, 2, 3])
486 .append_immutable_trio()
487 .build();
488 let logger = TestLogger::stdout();
489 let no_cache_digester = CardanoImmutableDigester::new(None, logger.clone());
490 let cache_digester = CardanoImmutableDigester::new(
491 Some(Arc::new(MemoryImmutableFileDigestCacheProvider::default())),
492 logger.clone(),
493 );
494 let immutable_range = 1..=3;
495
496 let without_cache_digests = no_cache_digester
497 .compute_digests_for_range(cardano_db.get_immutable_dir(), &immutable_range)
498 .await
499 .expect("compute_digests_for_range must not fail");
500
501 let cold_cache_digests = cache_digester
502 .compute_digests_for_range(cardano_db.get_immutable_dir(), &immutable_range)
503 .await
504 .expect("compute_digests_for_range must not fail");
505
506 let full_cache_digests = cache_digester
507 .compute_digests_for_range(cardano_db.get_immutable_dir(), &immutable_range)
508 .await
509 .expect("compute_digests_for_range must not fail");
510
511 let without_cache_entries = without_cache_digests.entries;
512 let cold_cache_entries = cold_cache_digests.entries;
513 let full_cache_entries = full_cache_digests.entries;
514 assert_eq!(
515 without_cache_entries, full_cache_entries,
516 "Digests for range with or without cache should be the same"
517 );
518
519 assert_eq!(
520 cold_cache_entries, full_cache_entries,
521 "Digests for range with cold or with hot cache should be the same"
522 );
523 }
524
525 #[tokio::test]
526 async fn cache_read_failure_dont_block_computations() {
527 let cardano_db = db_builder("cache_read_failure_dont_block_computation")
528 .with_immutables(&[1, 2, 3])
529 .append_immutable_trio()
530 .build();
531 let mut cache = MockImmutableFileDigestCacheProvider::new();
532 cache.expect_get().returning(|_| Ok(BTreeMap::new()));
533 cache.expect_store().returning(|_| {
534 Err(ImmutableDigesterCacheProviderError::Store(
535 ImmutableDigesterCacheStoreError::Io(io::Error::other("error")),
536 ))
537 });
538 let logger = TestLogger::stdout();
539 let digester = CardanoImmutableDigester::new(Some(Arc::new(cache)), logger.clone());
540 let beacon = CardanoDbBeacon::new(1, 3);
541
542 digester
543 .compute_merkle_tree(cardano_db.get_immutable_dir(), &beacon)
544 .await
545 .expect("compute_merkle_tree must not fail even with cache write failure");
546 }
547
548 #[tokio::test]
549 async fn cache_write_failure_dont_block_computation() {
550 let cardano_db = db_builder("cache_write_failure_dont_block_computation")
551 .with_immutables(&[1, 2, 3])
552 .append_immutable_trio()
553 .build();
554 let mut cache = MockImmutableFileDigestCacheProvider::new();
555 cache.expect_get().returning(|_| {
556 Err(ImmutableDigesterCacheProviderError::Get(
557 ImmutableDigesterCacheGetError::Io(io::Error::other("error")),
558 ))
559 });
560 cache.expect_store().returning(|_| Ok(()));
561 let logger = TestLogger::stdout();
562 let digester = CardanoImmutableDigester::new(Some(Arc::new(cache)), logger.clone());
563 let beacon = CardanoDbBeacon::new(1, 3);
564
565 digester
566 .compute_merkle_tree(cardano_db.get_immutable_dir(), &beacon)
567 .await
568 .expect("compute_merkle_tree must not fail even with cache read failure");
569 }
570}