massa_executed_ops/
executed_ops.rs

1//! Copyright (c) 2022 MASSA LABS <info@massa.net>
2
3//! This file defines a structure to list and prune previously executed operations.
4//! Used to detect operation reuse.
5
6use crate::ExecutedOpsConfig;
7use massa_db_exports::{
8    DBBatch, ShareableMassaDBController, CRUD_ERROR, EXECUTED_OPS_ID_DESER_ERROR,
9    EXECUTED_OPS_ID_SER_ERROR, EXECUTED_OPS_PREFIX, STATE_CF,
10};
11use massa_models::{
12    operation::{OperationId, OperationIdDeserializer, OperationIdSerializer},
13    prehash::{PreHashMap, PreHashSet},
14    slot::{Slot, SlotDeserializer, SlotSerializer},
15};
16use massa_serialization::{
17    BoolDeserializer, BoolSerializer, DeserializeError, Deserializer, Serializer,
18};
19use std::{
20    collections::{BTreeMap, HashMap, HashSet},
21    ops::Bound::{Excluded, Included},
22};
23
24/// Changes for ExecutedOps (was_successful, op_expiry_slot)
25pub type ExecutedOpsChanges = PreHashMap<OperationId, (bool, Slot)>;
26
27/// Op id key formatting macro
28#[macro_export]
29macro_rules! op_id_key {
30    ($id:expr) => {
31        [&EXECUTED_OPS_PREFIX.as_bytes(), &$id[..]].concat()
32    };
33}
34
35/// A structure to list and prune previously executed operations
36#[derive(Clone)]
37pub struct ExecutedOps {
38    /// Executed operations configuration
39    config: ExecutedOpsConfig,
40    /// RocksDB Instance
41    pub db: ShareableMassaDBController,
42    /// Executed operations btreemap with slot as index for better pruning complexity
43    pub sorted_ops: BTreeMap<Slot, PreHashSet<OperationId>>,
44    /// execution status of operations (true: success, false: fail)
45    pub op_exec_status: HashMap<OperationId, bool>,
46    operation_id_deserializer: OperationIdDeserializer,
47    operation_id_serializer: OperationIdSerializer,
48    bool_deserializer: BoolDeserializer,
49    bool_serializer: BoolSerializer,
50    slot_deserializer: SlotDeserializer,
51    slot_serializer: SlotSerializer,
52}
53
54impl ExecutedOps {
55    /// Creates a new `ExecutedOps`
56    pub fn new(config: ExecutedOpsConfig, db: ShareableMassaDBController) -> Self {
57        let slot_deserializer = SlotDeserializer::new(
58            (Included(u64::MIN), Included(u64::MAX)),
59            (Included(0), Excluded(config.thread_count)),
60        );
61        Self {
62            config,
63            db,
64            sorted_ops: BTreeMap::new(),
65            op_exec_status: HashMap::new(),
66            operation_id_deserializer: OperationIdDeserializer::new(),
67            operation_id_serializer: OperationIdSerializer::new(),
68            bool_deserializer: BoolDeserializer::new(),
69            bool_serializer: BoolSerializer::new(),
70            slot_deserializer,
71            slot_serializer: SlotSerializer::new(),
72        }
73    }
74
75    /// Get the execution statuses of a set of operations.
76    /// Returns a list where each element is None if no execution was found for that op,
77    /// or a boolean indicating whether the execution was successful (true) or had an error (false).
78    pub fn get_ops_exec_status(&self, batch: &[OperationId]) -> Vec<Option<bool>> {
79        batch
80            .iter()
81            .map(|op_id| self.op_exec_status.get(op_id).copied())
82            .collect()
83    }
84
85    /// Recomputes the local caches after bootstrap or loading the state from disk
86    pub fn recompute_sorted_ops_and_op_exec_status(&mut self) {
87        self.sorted_ops.clear();
88        self.op_exec_status.clear();
89
90        let db = self.db.read();
91
92        for (serialized_op_id, serialized_value) in
93            db.prefix_iterator_cf(STATE_CF, EXECUTED_OPS_PREFIX.as_bytes())
94        {
95            if !serialized_op_id.starts_with(EXECUTED_OPS_PREFIX.as_bytes()) {
96                break;
97            }
98
99            let (_, op_id) = self
100                .operation_id_deserializer
101                .deserialize::<DeserializeError>(&serialized_op_id[EXECUTED_OPS_PREFIX.len()..])
102                .expect(EXECUTED_OPS_ID_DESER_ERROR);
103
104            let (rest, op_exec_status) = self
105                .bool_deserializer
106                .deserialize::<DeserializeError>(&serialized_value)
107                .expect(EXECUTED_OPS_ID_DESER_ERROR);
108            let (_, slot) = self
109                .slot_deserializer
110                .deserialize::<DeserializeError>(rest)
111                .expect(EXECUTED_OPS_ID_DESER_ERROR);
112
113            self.sorted_ops
114                .entry(slot)
115                .and_modify(|ids| {
116                    ids.insert(op_id);
117                })
118                .or_insert_with(|| {
119                    let mut new = HashSet::default();
120                    new.insert(op_id);
121                    new
122                });
123            self.op_exec_status.insert(op_id, op_exec_status);
124        }
125    }
126
127    /// Reset the executed operations
128    ///
129    /// USED FOR BOOTSTRAP ONLY
130    pub fn reset(&mut self) {
131        self.db
132            .write()
133            .delete_prefix(EXECUTED_OPS_PREFIX, STATE_CF, None);
134
135        self.recompute_sorted_ops_and_op_exec_status();
136    }
137
138    /// Apply speculative operations changes to the final executed operations state
139    pub fn apply_changes_to_batch(
140        &mut self,
141        changes: ExecutedOpsChanges,
142        slot: Slot,
143        batch: &mut DBBatch,
144    ) {
145        for (id, value) in changes.iter() {
146            self.put_entry(id, value, batch);
147        }
148
149        for (op_id, (op_exec_success, slot)) in changes {
150            self.sorted_ops
151                .entry(slot)
152                .and_modify(|ids| {
153                    ids.insert(op_id);
154                })
155                .or_insert_with(|| {
156                    let mut new = PreHashSet::default();
157                    new.insert(op_id);
158                    new
159                });
160            self.op_exec_status.insert(op_id, op_exec_success);
161        }
162
163        self.prune_to_batch(slot, batch);
164    }
165
166    /// Check if an operation was executed
167    pub fn contains(&self, op_id: &OperationId) -> bool {
168        let db = self.db.read();
169
170        let mut serialized_op_id = Vec::new();
171        self.operation_id_serializer
172            .serialize(op_id, &mut serialized_op_id)
173            .expect(EXECUTED_OPS_ID_SER_ERROR);
174
175        db.get_cf(STATE_CF, op_id_key!(serialized_op_id))
176            .expect(CRUD_ERROR)
177            .is_some()
178    }
179
180    /// Prune all expired operations
181    fn prune_to_batch(&mut self, slot: Slot, batch: &mut DBBatch) {
182        // Force-keep `keep_executed_history_extra_periods` for API polling safety
183        let cutoff_slot = match slot
184            .period
185            .checked_sub(self.config.keep_executed_history_extra_periods)
186        {
187            Some(cutoff_slot) => Slot::new(cutoff_slot, slot.thread),
188            None => return,
189        };
190
191        let kept = self.sorted_ops.split_off(&cutoff_slot);
192        let removed = std::mem::take(&mut self.sorted_ops);
193        for (_, ids) in removed {
194            for op_id in ids {
195                self.op_exec_status.remove(&op_id);
196                self.delete_entry(&op_id, batch);
197            }
198        }
199        self.sorted_ops = kept;
200    }
201
202    /// Add an executed_op to the DB
203    ///
204    /// # Arguments
205    /// * `op_id`
206    /// * `value`: execution status and validity slot
207    /// * `batch`: the given operation batch to update
208    fn put_entry(&self, op_id: &OperationId, value: &(bool, Slot), batch: &mut DBBatch) {
209        let db = self.db.read();
210
211        let mut serialized_op_id = Vec::new();
212        self.operation_id_serializer
213            .serialize(op_id, &mut serialized_op_id)
214            .expect(EXECUTED_OPS_ID_SER_ERROR);
215
216        let mut serialized_op_value = Vec::new();
217        self.bool_serializer
218            .serialize(&value.0, &mut serialized_op_value)
219            .expect(EXECUTED_OPS_ID_SER_ERROR);
220        self.slot_serializer
221            .serialize(&value.1, &mut serialized_op_value)
222            .expect(EXECUTED_OPS_ID_SER_ERROR);
223
224        db.put_or_update_entry_value(batch, op_id_key!(serialized_op_id), &serialized_op_value);
225    }
226
227    /// Remove a op_id from the DB
228    ///
229    /// # Arguments
230    /// * batch: the given operation batch to update
231    fn delete_entry(&self, op_id: &OperationId, batch: &mut DBBatch) {
232        let db = self.db.read();
233
234        let mut serialized_op_id = Vec::new();
235        self.operation_id_serializer
236            .serialize(op_id, &mut serialized_op_id)
237            .expect(EXECUTED_OPS_ID_SER_ERROR);
238
239        db.delete_key(batch, op_id_key!(serialized_op_id));
240    }
241
242    /// Deserializes the key and value, useful after bootstrap
243    pub fn is_key_value_valid(&self, serialized_key: &[u8], serialized_value: &[u8]) -> bool {
244        if !serialized_key.starts_with(EXECUTED_OPS_PREFIX.as_bytes()) {
245            return false;
246        }
247
248        let Ok((rest, _id)): Result<(&[u8], OperationId), nom::Err<DeserializeError>> = self
249            .operation_id_deserializer
250            .deserialize::<DeserializeError>(&serialized_key[EXECUTED_OPS_PREFIX.len()..])
251        else {
252            return false;
253        };
254        if !rest.is_empty() {
255            return false;
256        }
257
258        let Ok((rest, _bool)) = self
259            .bool_deserializer
260            .deserialize::<DeserializeError>(serialized_value)
261        else {
262            return false;
263        };
264        let Ok((rest, _slot)) = self.slot_deserializer.deserialize::<DeserializeError>(rest) else {
265            return false;
266        };
267        if !rest.is_empty() {
268            return false;
269        }
270
271        true
272    }
273}
274#[cfg(test)]
275mod test {
276    use std::sync::Arc;
277
278    use parking_lot::RwLock;
279    use tempfile::{tempdir, TempDir};
280
281    use massa_db_exports::{MassaDBConfig, MassaDBController, STATE_HASH_INITIAL_BYTES};
282    use massa_db_worker::MassaDB;
283    use massa_hash::{Hash, HashXof};
284    use massa_models::config::{KEEP_EXECUTED_HISTORY_EXTRA_PERIODS, THREAD_COUNT};
285    use massa_models::prehash::PreHashMap;
286    use massa_models::secure_share::Id;
287
288    use super::*;
289
290    #[test]
291    fn test_executed_ops_cache() {
292        // initialize the executed ops config
293        // let thread_count = 2;
294        let config = ExecutedOpsConfig {
295            thread_count: THREAD_COUNT,
296            keep_executed_history_extra_periods: KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
297        };
298
299        // Db init
300        let temp_dir = tempdir().expect("Unable to create a temp folder");
301        let db_config = MassaDBConfig {
302            path: temp_dir.path().to_path_buf(),
303            max_history_length: 100,
304            max_final_state_elements_size: 100,
305            max_versioning_elements_size: 100,
306            thread_count: THREAD_COUNT,
307            max_ledger_backups: 10,
308            enable_metrics: false,
309        };
310        let db = Arc::new(RwLock::new(
311            Box::new(MassaDB::new(db_config.clone())) as Box<dyn MassaDBController + 'static>
312        ));
313
314        let mut exec_ops = ExecutedOps::new(config.clone(), db.clone());
315
316        let mut changes = PreHashMap::default();
317
318        let slot_1 = Slot::new(1, 0);
319        let op_id_1 = OperationId::new(Hash::compute_from(&[0]));
320        changes.insert(op_id_1, (true, slot_1));
321        let slot_2 = Slot::new(KEEP_EXECUTED_HISTORY_EXTRA_PERIODS + 2, 3);
322        let op_id_2 = OperationId::new(Hash::compute_from(&[1]));
323        changes.insert(op_id_2, (true, slot_2));
324
325        let mut batch = DBBatch::new();
326        exec_ops.apply_changes_to_batch(changes, slot_2, &mut batch);
327        db.write().write_batch(batch, Default::default(), None);
328
329        // cache len expected is 1, expect op_id_1 to be discarded (as it is tool old)
330        assert_eq!(exec_ops.sorted_ops.len(), 1);
331        assert!(!exec_ops.contains(&op_id_1));
332        assert!(exec_ops.contains(&op_id_2));
333
334        let sorted_ops_1 = exec_ops.sorted_ops.clone();
335        drop(db);
336        drop(exec_ops);
337
338        let db2 = Arc::new(RwLock::new(
339            Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>
340        ));
341
342        // After an init from disk, cache is empty, so recompute it and compare with original
343        let mut exec_ops2 = ExecutedOps::new(config, db2.clone());
344        exec_ops2.recompute_sorted_ops_and_op_exec_status();
345        assert_eq!(exec_ops2.sorted_ops, sorted_ops_1);
346
347        // Reset cache
348        exec_ops2.reset();
349        assert_eq!(exec_ops2.sorted_ops.len(), 0);
350    }
351
352    #[test]
353    fn test_executed_ops_hash_computing() {
354        // initialize the executed ops config
355        let thread_count = 2;
356        let config = ExecutedOpsConfig {
357            thread_count,
358            keep_executed_history_extra_periods: 2,
359        };
360        let tempdir_a = TempDir::new().expect("cannot create temp directory");
361        let tempdir_c = TempDir::new().expect("cannot create temp directory");
362        let db_a_config = MassaDBConfig {
363            path: tempdir_a.path().to_path_buf(),
364            max_history_length: 10,
365            max_final_state_elements_size: 100,
366            max_versioning_elements_size: 100,
367            thread_count,
368            max_ledger_backups: 10,
369            enable_metrics: false,
370        };
371        let db_c_config = MassaDBConfig {
372            path: tempdir_c.path().to_path_buf(),
373            max_history_length: 10,
374            max_final_state_elements_size: 100,
375            max_versioning_elements_size: 100,
376            thread_count,
377            max_ledger_backups: 10,
378            enable_metrics: false,
379        };
380
381        let db_a = Arc::new(RwLock::new(
382            Box::new(MassaDB::new(db_a_config)) as Box<dyn MassaDBController + 'static>
383        ));
384        let db_c = Arc::new(RwLock::new(
385            Box::new(MassaDB::new(db_c_config)) as Box<dyn MassaDBController + 'static>
386        ));
387
388        // initialize the executed ops and executed ops changes
389        let mut a = ExecutedOps::new(config.clone(), db_a.clone());
390        let mut c = ExecutedOps::new(config.clone(), db_c.clone());
391        let mut change_a = PreHashMap::default();
392        let mut change_b = PreHashMap::default();
393        let mut change_c = PreHashMap::default();
394        for i in 0u8..20 {
395            let expiration_slot = Slot {
396                period: (i as u64).saturating_sub(config.keep_executed_history_extra_periods),
397                thread: 0,
398            };
399            if i < 12 {
400                change_a.insert(
401                    OperationId::new(Hash::compute_from(&[i])),
402                    (true, expiration_slot),
403                );
404            }
405            if i > 8 {
406                change_b.insert(
407                    OperationId::new(Hash::compute_from(&[i])),
408                    (true, expiration_slot),
409                );
410            }
411            change_c.insert(
412                OperationId::new(Hash::compute_from(&[i])),
413                (true, expiration_slot),
414            );
415        }
416
417        // apply change_b to a which performs a.hash ^ $(change_b)
418        let apply_slot = Slot {
419            period: 0,
420            thread: 0,
421        };
422
423        let mut batch_a = DBBatch::new();
424        a.apply_changes_to_batch(change_a, apply_slot, &mut batch_a);
425        db_a.write().write_batch(batch_a, Default::default(), None);
426
427        let mut batch_b = DBBatch::new();
428        a.apply_changes_to_batch(change_b, apply_slot, &mut batch_b);
429        db_a.write().write_batch(batch_b, Default::default(), None);
430
431        let mut batch_c = DBBatch::new();
432        c.apply_changes_to_batch(change_c, apply_slot, &mut batch_c);
433        db_c.write().write_batch(batch_c, Default::default(), None);
434
435        // check that a.hash ^ $(change_b) = c.hash
436        assert_ne!(
437            db_a.read().get_xof_db_hash(),
438            HashXof(*STATE_HASH_INITIAL_BYTES)
439        );
440        assert_eq!(
441            db_a.read().get_xof_db_hash(),
442            db_c.read().get_xof_db_hash(),
443            "'a' and 'c' hashes are not equal"
444        );
445
446        // prune every element
447        let prune_slot = Slot {
448            period: 20,
449            thread: 0,
450        };
451        let mut batch_a = DBBatch::new();
452        a.prune_to_batch(prune_slot, &mut batch_a);
453        db_a.write().write_batch(batch_a, Default::default(), None);
454
455        // at this point the hash should have been reset to its original value
456        assert_eq!(
457            db_a.read().get_xof_db_hash(),
458            HashXof(*STATE_HASH_INITIAL_BYTES),
459            "'a' was not reset to its initial value"
460        );
461    }
462}