massa_executed_ops/
executed_denunciations.rs

1//! Copyright (c) 2023 MASSA LABS <info@massa.net>
2
3//! This file defines a structure to list and prune previously executed denunciations.
4//! Used to detect denunciation reuse.
5
6use crate::ExecutedDenunciationsConfig;
7use massa_db_exports::{
8    DBBatch, ShareableMassaDBController, CRUD_ERROR, EXECUTED_DENUNCIATIONS_INDEX_DESER_ERROR,
9    EXECUTED_DENUNCIATIONS_INDEX_SER_ERROR, EXECUTED_DENUNCIATIONS_PREFIX, STATE_CF,
10};
11use massa_models::denunciation::Denunciation;
12use massa_models::{
13    denunciation::{DenunciationIndex, DenunciationIndexDeserializer, DenunciationIndexSerializer},
14    slot::Slot,
15};
16use massa_serialization::{DeserializeError, Deserializer, Serializer};
17use std::collections::{BTreeMap, HashSet};
18
19/// Speculative changes for ExecutedDenunciations
20pub type ExecutedDenunciationsChanges = HashSet<DenunciationIndex>;
21
22/// Denunciation index key formatting macro
23#[macro_export]
24macro_rules! denunciation_index_key {
25    ($id:expr) => {
26        [&EXECUTED_DENUNCIATIONS_PREFIX.as_bytes(), &$id[..]].concat()
27    };
28}
29
30/// A structure to list and prune previously executed denunciations
31#[derive(Clone)]
32pub struct ExecutedDenunciations {
33    /// Executed denunciations configuration
34    config: ExecutedDenunciationsConfig,
35    /// Access to the RocksDB database
36    pub db: ShareableMassaDBController,
37    /// for better pruning complexity
38    pub sorted_denunciations: BTreeMap<Slot, HashSet<DenunciationIndex>>,
39    /// for rocksdb serialization
40    denunciation_index_serializer: DenunciationIndexSerializer,
41    /// for rocksdb deserialization
42    denunciation_index_deserializer: DenunciationIndexDeserializer,
43}
44
45impl ExecutedDenunciations {
46    /// Create a new `ExecutedDenunciations`
47    pub fn new(config: ExecutedDenunciationsConfig, db: ShareableMassaDBController) -> Self {
48        let denunciation_index_deserializer =
49            DenunciationIndexDeserializer::new(config.thread_count, config.endorsement_count);
50        Self {
51            config,
52            db,
53            sorted_denunciations: Default::default(),
54            denunciation_index_serializer: DenunciationIndexSerializer::new(),
55            denunciation_index_deserializer,
56        }
57    }
58
59    /// Recomputes the local caches after bootstrap or loading the state from disk
60    pub fn recompute_sorted_denunciations(&mut self) {
61        self.sorted_denunciations.clear();
62
63        let db = self.db.read();
64
65        for (serialized_de_idx, _) in
66            db.prefix_iterator_cf(STATE_CF, EXECUTED_DENUNCIATIONS_PREFIX.as_bytes())
67        {
68            if !serialized_de_idx.starts_with(EXECUTED_DENUNCIATIONS_PREFIX.as_bytes()) {
69                break;
70            }
71            let (_, de_idx) = self
72                .denunciation_index_deserializer
73                .deserialize::<DeserializeError>(
74                    &serialized_de_idx[EXECUTED_DENUNCIATIONS_PREFIX.len()..],
75                )
76                .expect(EXECUTED_DENUNCIATIONS_INDEX_DESER_ERROR);
77
78            self.sorted_denunciations
79                .entry(*de_idx.get_slot())
80                .and_modify(|ids| {
81                    ids.insert(de_idx);
82                })
83                .or_insert_with(|| {
84                    let mut new = HashSet::default();
85                    new.insert(de_idx);
86                    new
87                });
88        }
89    }
90
91    /// Reset the executed denunciations
92    ///
93    /// USED FOR BOOTSTRAP ONLY
94    pub fn reset(&mut self) {
95        {
96            let mut db = self.db.write();
97            db.delete_prefix(EXECUTED_DENUNCIATIONS_PREFIX, STATE_CF, None);
98        }
99
100        self.recompute_sorted_denunciations();
101    }
102
103    /// Check if a denunciation (e.g. a denunciation index) was executed
104    pub fn contains(&self, de_idx: &DenunciationIndex) -> bool {
105        let db = self.db.read();
106
107        let mut serialized_de_idx = Vec::new();
108        self.denunciation_index_serializer
109            .serialize(de_idx, &mut serialized_de_idx)
110            .expect(EXECUTED_DENUNCIATIONS_INDEX_SER_ERROR);
111
112        db.get_cf(STATE_CF, denunciation_index_key!(serialized_de_idx))
113            .expect(CRUD_ERROR)
114            .is_some()
115    }
116
117    /// Apply speculative operations changes to the final executed denunciations state
118    pub fn apply_changes_to_batch(
119        &mut self,
120        changes: ExecutedDenunciationsChanges,
121        slot: Slot,
122        batch: &mut DBBatch,
123    ) {
124        for de_idx in changes {
125            self.put_entry(&de_idx, batch);
126            self.sorted_denunciations
127                .entry(*de_idx.get_slot())
128                .and_modify(|ids| {
129                    ids.insert(de_idx);
130                })
131                .or_insert_with(|| {
132                    let mut new = HashSet::default();
133                    new.insert(de_idx);
134                    new
135                });
136        }
137
138        self.prune_to_batch(slot, batch);
139    }
140
141    /// Prune all denunciations that have expired, assuming the given slot is final
142    fn prune_to_batch(&mut self, slot: Slot, batch: &mut DBBatch) {
143        // Force-keep `keep_executed_history_extra_periods` for API polling safety
144        let effective_expiry_periods = self
145            .config
146            .denunciation_expire_periods
147            .saturating_add(self.config.keep_executed_history_extra_periods);
148        let mut drained: HashSet<DenunciationIndex> = Default::default();
149        self.sorted_denunciations.retain(|de_idx_slot, de_idx| {
150            if Denunciation::is_expired(
151                &de_idx_slot.period,
152                &slot.period,
153                &effective_expiry_periods,
154            ) {
155                drained.extend(de_idx.iter());
156                return false;
157            }
158            true
159        });
160        for de_idx in drained {
161            self.delete_entry(&de_idx, batch);
162        }
163    }
164
165    /// Add a denunciation_index to the DB
166    ///
167    /// # Arguments
168    /// * `de_idx`
169    /// * `batch`: the given operation batch to update
170    fn put_entry(&self, de_idx: &DenunciationIndex, batch: &mut DBBatch) {
171        let db = self.db.read();
172
173        let mut serialized_de_idx = Vec::new();
174        self.denunciation_index_serializer
175            .serialize(de_idx, &mut serialized_de_idx)
176            .expect(EXECUTED_DENUNCIATIONS_INDEX_SER_ERROR);
177
178        db.put_or_update_entry_value(batch, denunciation_index_key!(serialized_de_idx), b"");
179    }
180
181    /// Remove a denunciation_index from the DB
182    ///
183    /// # Arguments
184    /// * `de_idx`: the denunciation index to remove
185    /// * batch: the given operation batch to update
186    fn delete_entry(&self, de_idx: &DenunciationIndex, batch: &mut DBBatch) {
187        let db = self.db.read();
188
189        let mut serialized_de_idx = Vec::new();
190        self.denunciation_index_serializer
191            .serialize(de_idx, &mut serialized_de_idx)
192            .expect(EXECUTED_DENUNCIATIONS_INDEX_SER_ERROR);
193
194        db.delete_key(batch, denunciation_index_key!(serialized_de_idx));
195    }
196
197    /// Deserializes the key and value, useful after bootstrap
198    pub fn is_key_value_valid(&self, serialized_key: &[u8], serialized_value: &[u8]) -> bool {
199        if !serialized_key.starts_with(EXECUTED_DENUNCIATIONS_PREFIX.as_bytes()) {
200            return false;
201        }
202
203        let Ok((rest, _idx)) = self
204            .denunciation_index_deserializer
205            .deserialize::<DeserializeError>(
206                &serialized_key[EXECUTED_DENUNCIATIONS_PREFIX.len()..],
207            )
208        else {
209            return false;
210        };
211        if !rest.is_empty() {
212            return false;
213        }
214
215        if !serialized_value.is_empty() {
216            return false;
217        }
218
219        true
220    }
221}
222
223#[cfg(test)]
224mod test {
225    use super::*;
226    use massa_db_exports::{MassaDBConfig, MassaDBController};
227    use massa_db_worker::MassaDB;
228    use massa_models::config::{
229        DENUNCIATION_EXPIRE_PERIODS, ENDORSEMENT_COUNT, KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
230        THREAD_COUNT,
231    };
232    use parking_lot::RwLock;
233    use std::sync::Arc;
234    use tempfile::tempdir;
235
236    #[test]
237    fn test_exec_de_cache() {
238        // Check executed denunciations cache grow / reset / recompute
239
240        let config = ExecutedDenunciationsConfig {
241            denunciation_expire_periods: DENUNCIATION_EXPIRE_PERIODS,
242            thread_count: THREAD_COUNT,
243            endorsement_count: ENDORSEMENT_COUNT,
244            keep_executed_history_extra_periods: KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
245        };
246        // Db init
247        let temp_dir = tempdir().expect("Unable to create a temp folder");
248        // println!("Using temp dir: {:?}", temp_dir.path());
249        let db_config = MassaDBConfig {
250            path: temp_dir.path().to_path_buf(),
251            max_history_length: 100,
252            max_final_state_elements_size: 100,
253            max_versioning_elements_size: 100,
254            max_ledger_backups: 10,
255            thread_count: THREAD_COUNT,
256            enable_metrics: false,
257        };
258        let db = Arc::new(RwLock::new(
259            Box::new(MassaDB::new(db_config.clone())) as Box<dyn MassaDBController + 'static>
260        ));
261
262        let mut exec_de = ExecutedDenunciations::new(config.clone(), db);
263
264        // Add some data to exec_de
265        let slot_1 = Slot::new(1, 0);
266        let de_idx_1 = DenunciationIndex::Endorsement {
267            slot: slot_1,
268            index: ENDORSEMENT_COUNT - 4,
269        };
270        let slot_2 = Slot::new(
271            DENUNCIATION_EXPIRE_PERIODS + KEEP_EXECUTED_HISTORY_EXTRA_PERIODS + 2,
272            3,
273        );
274        let de_idx_2 = DenunciationIndex::Endorsement {
275            slot: slot_2,
276            index: ENDORSEMENT_COUNT - 1,
277        };
278        let mut changes = ExecutedDenunciationsChanges::new();
279        changes.insert(de_idx_1);
280        changes.insert(de_idx_2);
281        let mut batch = DBBatch::new();
282        exec_de.apply_changes_to_batch(changes, slot_2, &mut batch);
283        exec_de
284            .db
285            .write()
286            .write_batch(batch.clone(), DBBatch::new(), Some(slot_2));
287
288        assert_eq!(exec_de.sorted_denunciations.len(), 1);
289        assert_eq!(
290            exec_de.sorted_denunciations.get(&slot_2),
291            Some(&HashSet::from([de_idx_2]))
292        );
293        assert!(!exec_de.contains(&de_idx_1));
294        assert!(exec_de.contains(&de_idx_2));
295
296        let sorted_deunciations_1 = exec_de.sorted_denunciations.clone();
297        drop(exec_de);
298
299        // Init an exec de from disk
300        let db2 = Arc::new(RwLock::new(
301            Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>
302        ));
303
304        let mut exec_de2 = ExecutedDenunciations::new(config, db2);
305
306        // After init from disk, the cache is empty, so recompute it and check against saved cache
307        exec_de2.recompute_sorted_denunciations();
308        assert_eq!(exec_de2.sorted_denunciations, sorted_deunciations_1);
309
310        // Reset cache
311        exec_de2.reset();
312        assert_eq!(exec_de2.sorted_denunciations.len(), 0);
313    }
314}