Skip to main content

mithril_cardano_node_internal_database/digesters/
cardano_immutable_digester.rs

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
18/// A digester working directly on a Cardano DB immutables files
19pub struct CardanoImmutableDigester {
20    /// A [ImmutableFileDigestCacheProvider] instance
21    cache_provider: Option<Arc<dyn ImmutableFileDigestCacheProvider>>,
22
23    /// The logger where the logs should be written
24    logger: Logger,
25}
26
27impl CardanoImmutableDigester {
28    /// ImmutableDigester factory
29    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        // The computation of immutable files digests is done in a separate thread because it is blocking the whole task
46        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}