1use 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
24pub type ExecutedOpsChanges = PreHashMap<OperationId, (bool, Slot)>;
26
27#[macro_export]
29macro_rules! op_id_key {
30 ($id:expr) => {
31 [&EXECUTED_OPS_PREFIX.as_bytes(), &$id[..]].concat()
32 };
33}
34
35#[derive(Clone)]
37pub struct ExecutedOps {
38 config: ExecutedOpsConfig,
40 pub db: ShareableMassaDBController,
42 pub sorted_ops: BTreeMap<Slot, PreHashSet<OperationId>>,
44 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 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 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 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 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 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 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 fn prune_to_batch(&mut self, slot: Slot, batch: &mut DBBatch) {
182 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 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 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 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 let config = ExecutedOpsConfig {
295 thread_count: THREAD_COUNT,
296 keep_executed_history_extra_periods: KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
297 };
298
299 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 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 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 exec_ops2.reset();
349 assert_eq!(exec_ops2.sorted_ops.len(), 0);
350 }
351
352 #[test]
353 fn test_executed_ops_hash_computing() {
354 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 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 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 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 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 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}