1use crate::controller_trait::FinalStateController;
9use crate::{config::FinalStateConfig, error::FinalStateError, state_changes::StateChanges};
10
11use anyhow::{anyhow, Result as AnyResult};
12use massa_async_pool::AsyncPool;
13use massa_db_exports::{
14 DBBatch, MassaIteratorMode, ShareableMassaDBController, ASYNC_POOL_PREFIX,
15 CYCLE_HISTORY_PREFIX, DEFERRED_CALLS_PREFIX, DEFERRED_CALL_TOTAL_GAS, DEFERRED_CREDITS_PREFIX,
16 EXECUTED_DENUNCIATIONS_PREFIX, EXECUTED_OPS_PREFIX, LEDGER_PREFIX, MIP_STORE_PREFIX, STATE_CF,
17};
18use massa_db_exports::{EXECUTION_TRAIL_HASH_PREFIX, MIP_STORE_STATS_PREFIX, VERSIONING_CF};
19use massa_deferred_calls::DeferredCallRegistry;
20use massa_executed_ops::ExecutedDenunciations;
21use massa_executed_ops::ExecutedOps;
22use massa_hash::Hash;
23use massa_ledger_exports::LedgerController;
24use massa_models::operation::OperationId;
25use massa_models::slot::Slot;
26use massa_models::timeslots::get_block_slot_timestamp;
27use massa_models::types::SetOrKeep;
28use massa_pos_exports::{PoSFinalState, SelectorController};
29use massa_versioning::versioning::MipStore;
30use tracing::{debug, info, warn};
31
32pub struct FinalState {
34 pub(crate) config: FinalStateConfig,
36 pub ledger: Box<dyn LedgerController>,
38 pub async_pool: AsyncPool,
40 pub deferred_call_registry: DeferredCallRegistry,
42 pub pos_state: PoSFinalState,
44 pub executed_ops: ExecutedOps,
46 pub executed_denunciations: ExecutedDenunciations,
48 pub mip_store: MipStore,
50 pub last_start_period: u64,
55 pub last_slot_before_downtime: Option<Slot>,
60 pub db: ShareableMassaDBController,
62}
63
64impl FinalState {
65 pub fn new(
73 db: ShareableMassaDBController,
74 config: FinalStateConfig,
75 ledger: Box<dyn LedgerController>,
76 selector: Box<dyn SelectorController>,
77 mip_store: MipStore,
78 reset_final_state: bool,
79 ) -> Result<Self, FinalStateError> {
80 let db_slot = db
81 .read()
82 .get_change_id()
83 .map_err(|_| FinalStateError::InvalidSlot(String::from("Could not get slot in db")))?;
84
85 let pos_state = PoSFinalState::new(
87 config.pos_config.clone(),
88 &config.initial_seed_string,
89 &config.initial_rolls_path,
90 selector,
91 db.clone(),
92 )
93 .map_err(|err| FinalStateError::PosError(format!("PoS final state init error: {}", err)))?;
94
95 let slot = if reset_final_state {
97 Slot::new(0, config.thread_count.saturating_sub(1))
98 } else {
99 db_slot
100 };
101
102 let async_pool = AsyncPool::new(config.async_pool_config.clone(), db.clone());
104
105 let executed_ops = ExecutedOps::new(config.executed_ops_config.clone(), db.clone());
107
108 let executed_denunciations =
110 ExecutedDenunciations::new(config.executed_denunciations_config.clone(), db.clone());
111
112 let deferred_call_registry =
113 DeferredCallRegistry::new(db.clone(), config.deferred_calls_config);
114
115 let mut final_state = FinalState {
116 ledger,
117 async_pool,
118 deferred_call_registry,
119 pos_state,
120 config,
121 executed_ops,
122 executed_denunciations,
123 mip_store,
124 last_start_period: 0,
125 last_slot_before_downtime: None,
126 db,
127 };
128
129 if reset_final_state {
130 final_state.db.read().set_initial_change_id(slot);
131 final_state
133 .db
134 .write()
135 .delete_prefix(EXECUTION_TRAIL_HASH_PREFIX, STATE_CF, None);
136 final_state.async_pool.reset();
137 final_state.pos_state.reset();
138 final_state.executed_ops.reset();
139 final_state.executed_denunciations.reset();
140 }
141
142 info!(
143 "final_state hash at slot {}: {}",
144 slot,
145 final_state.db.read().get_xof_db_hash()
146 );
147
148 Ok(final_state)
150 }
151
152 fn interpolate_downtime(&mut self) -> Result<(), FinalStateError> {
155 let current_slot =
156 self.db.read().get_change_id().map_err(|_| {
157 FinalStateError::InvalidSlot(String::from("Could not get slot in db"))
158 })?;
159 let current_slot_cycle = current_slot.get_cycle(self.config.periods_per_cycle);
160
161 let end_slot = Slot::new(
162 self.last_start_period,
163 self.config.thread_count.saturating_sub(1),
164 );
165 let end_slot_cycle = end_slot.get_cycle(self.config.periods_per_cycle);
166
167 debug!(
168 "Interpolating downtime between slots {} and {}",
169 current_slot, end_slot
170 );
171
172 if current_slot_cycle == end_slot_cycle {
173 self.interpolate_single_cycle(current_slot, end_slot)?;
175 } else {
176 self.interpolate_multiple_cycles(
178 current_slot,
179 end_slot,
180 current_slot_cycle,
181 end_slot_cycle,
182 )?;
183 }
184
185 let final_state_hash = self.db.read().get_xof_db_hash();
187
188 info!(
189 "final_state hash at slot {}: {}",
190 end_slot, final_state_hash
191 );
192
193 let cycle = end_slot.get_cycle(self.config.periods_per_cycle);
195
196 self.pos_state
197 .feed_cycle_state_hash(cycle, final_state_hash);
198
199 Ok(())
200 }
201
202 fn interpolate_single_cycle(
204 &mut self,
205 current_slot: Slot,
206 end_slot: Slot,
207 ) -> Result<(), FinalStateError> {
208 let latest_snapshot_cycle = self.pos_state.cycle_history_cache.back().cloned().ok_or(
209 FinalStateError::SnapshotError(String::from(
210 "Impossible to interpolate the downtime: no cycle in the given snapshot",
211 )),
212 )?;
213
214 let latest_snapshot_cycle_info = self
215 .pos_state
216 .get_cycle_info(latest_snapshot_cycle.0)
217 .ok_or_else(|| FinalStateError::SnapshotError(String::from("Missing cycle info")))?;
218
219 let mut batch = DBBatch::new();
220
221 self.pos_state
222 .cycle_history_cache
223 .pop_back()
224 .ok_or(FinalStateError::SnapshotError(String::from(
225 "Impossible to interpolate the downtime: no cycle in the given snapshot",
226 )))?;
227 self.pos_state
228 .delete_cycle_info(latest_snapshot_cycle.0, &mut batch);
229
230 self.pos_state
231 .create_new_cycle_from_last(
232 &latest_snapshot_cycle_info,
233 current_slot
234 .get_next_slot(self.config.thread_count)
235 .expect("Cannot get next slot"),
236 end_slot,
237 &mut batch,
238 )
239 .map_err(|err| FinalStateError::PosError(format!("{}", err)))?;
240
241 self.pos_state
242 .db
243 .write()
244 .write_batch(batch, Default::default(), Some(end_slot));
245
246 if end_slot.is_last_of_cycle(self.config.periods_per_cycle, self.config.thread_count) {
247 self.feed_cycle_hash_and_selector_for_interpolation(
249 end_slot.get_cycle(self.config.periods_per_cycle),
250 )?;
251 }
252
253 Ok(())
254 }
255
256 fn interpolate_multiple_cycles(
258 &mut self,
259 current_slot: Slot,
260 end_slot: Slot,
261 current_slot_cycle: u64,
262 end_slot_cycle: u64,
263 ) -> Result<(), FinalStateError> {
264 let latest_snapshot_cycle = self.pos_state.cycle_history_cache.back().cloned().ok_or(
265 FinalStateError::SnapshotError(String::from(
266 "Impossible to interpolate the downtime: no cycle in the given snapshot",
267 )),
268 )?;
269
270 let latest_snapshot_cycle_info = self
271 .pos_state
272 .get_cycle_info(latest_snapshot_cycle.0)
273 .ok_or_else(|| FinalStateError::SnapshotError(String::from("Missing cycle info")))?;
274
275 if !current_slot.is_last_of_cycle(self.config.periods_per_cycle, self.config.thread_count) {
277 let mut batch = DBBatch::new();
278
279 self.pos_state
280 .cycle_history_cache
281 .pop_back()
282 .ok_or(FinalStateError::SnapshotError(String::from(
283 "Impossible to interpolate the downtime: no cycle in the given snapshot",
284 )))?;
285 self.pos_state
286 .delete_cycle_info(latest_snapshot_cycle.0, &mut batch);
287
288 let last_slot = Slot::new_last_of_cycle(
289 current_slot_cycle,
290 self.config.periods_per_cycle,
291 self.config.thread_count,
292 )
293 .map_err(|err| {
294 FinalStateError::InvalidSlot(format!(
295 "Cannot create slot for interpolating downtime: {}",
296 err
297 ))
298 })?;
299
300 self.pos_state
301 .create_new_cycle_from_last(
302 &latest_snapshot_cycle_info,
303 current_slot
304 .get_next_slot(self.config.thread_count)
305 .expect("Cannot get next slot"),
306 last_slot,
307 &mut batch,
308 )
309 .map_err(|err| FinalStateError::PosError(format!("{}", err)))?;
310
311 self.pos_state
312 .db
313 .write()
314 .write_batch(batch, Default::default(), Some(last_slot));
315
316 self.feed_cycle_hash_and_selector_for_interpolation(current_slot_cycle)?;
318 }
319
320 let current_slot_cycle = current_slot_cycle + 1;
325
326 for cycle in current_slot_cycle..end_slot_cycle {
327 let first_slot = Slot::new_first_of_cycle(cycle, self.config.periods_per_cycle)
328 .map_err(|err| {
329 FinalStateError::InvalidSlot(format!(
330 "Cannot create slot for interpolating downtime: {}",
331 err
332 ))
333 })?;
334
335 let last_slot = Slot::new_last_of_cycle(
336 cycle,
337 self.config.periods_per_cycle,
338 self.config.thread_count,
339 )
340 .map_err(|err| {
341 FinalStateError::InvalidSlot(format!(
342 "Cannot create slot for interpolating downtime: {}",
343 err
344 ))
345 })?;
346
347 let mut batch = DBBatch::new();
348
349 self.pos_state
350 .create_new_cycle_from_last(
351 &latest_snapshot_cycle_info,
352 first_slot,
353 last_slot,
354 &mut batch,
355 )
356 .map_err(|err| FinalStateError::PosError(format!("{}", err)))?;
357
358 self.pos_state
359 .db
360 .write()
361 .write_batch(batch, Default::default(), Some(last_slot));
362
363 self.feed_cycle_hash_and_selector_for_interpolation(cycle)?;
365 }
366
367 let first_slot = Slot::new_first_of_cycle(end_slot_cycle, self.config.periods_per_cycle)
369 .map_err(|err| {
370 FinalStateError::InvalidSlot(format!(
371 "Cannot create slot for interpolating downtime: {}",
372 err
373 ))
374 })?;
375
376 let mut batch = DBBatch::new();
377
378 self.pos_state
379 .create_new_cycle_from_last(
380 &latest_snapshot_cycle_info,
381 first_slot,
382 end_slot,
383 &mut batch,
384 )
385 .map_err(|err| FinalStateError::PosError(format!("{}", err)))?;
386
387 if end_slot.is_last_of_cycle(self.config.periods_per_cycle, self.config.thread_count) {
389 self.feed_cycle_hash_and_selector_for_interpolation(end_slot_cycle)?;
391 }
392
393 while self.pos_state.cycle_history_cache.len() > self.pos_state.config.cycle_history_length
395 {
396 if let Some((cycle, _)) = self.pos_state.cycle_history_cache.pop_front() {
397 self.pos_state.delete_cycle_info(cycle, &mut batch);
398 }
399 }
400
401 self.db
402 .write()
403 .write_batch(batch, Default::default(), Some(end_slot));
404
405 Ok(())
406 }
407
408 fn feed_cycle_hash_and_selector_for_interpolation(
410 &mut self,
411 cycle: u64,
412 ) -> Result<(), FinalStateError> {
413 let final_state_hash = self.db.read().get_xof_db_hash();
414
415 self.pos_state
416 .feed_cycle_state_hash(cycle, final_state_hash);
417
418 self.pos_state
419 .feed_selector(cycle.checked_add(2).ok_or_else(|| {
420 FinalStateError::PosError("cycle overflow when feeding selector".into())
421 })?)
422 .map_err(|_| {
423 FinalStateError::PosError("cycle overflow when feeding selector".into())
424 })?;
425 Ok(())
426 }
427
428 fn _finalize(
429 &mut self,
430 slot: Slot,
431 changes: StateChanges,
432 network_versions: Option<(u32, Option<u32>)>,
433 ) -> AnyResult<()> {
434 let cur_slot = self.db.read().get_change_id()?;
435 let next_slot = cur_slot.get_next_slot(self.config.thread_count)?;
437
438 if slot != next_slot {
439 return Err(anyhow!(
440 "attempting to apply execution state changes at slot {} while the current slot is {}",
441 slot, cur_slot
442 ));
443 }
444
445 let mut db_batch = DBBatch::new();
446 let mut db_versioning_batch = DBBatch::new();
447
448 self.async_pool
451 .apply_changes_to_batch(&changes.async_pool_changes, &mut db_batch);
452 self.pos_state
453 .apply_changes_to_batch(changes.pos_changes, slot, true, &mut db_batch)?;
454
455 self.ledger
458 .apply_changes_to_batch(changes.ledger_changes, &mut db_batch);
459 self.executed_ops
460 .apply_changes_to_batch(changes.executed_ops_changes, slot, &mut db_batch);
461
462 self.deferred_call_registry
463 .apply_changes_to_batch(changes.deferred_call_changes, &mut db_batch);
464
465 self.executed_denunciations.apply_changes_to_batch(
466 changes.executed_denunciations_changes,
467 slot,
468 &mut db_batch,
469 );
470
471 let slot_ts = get_block_slot_timestamp(
472 self.config.thread_count,
473 self.config.t0,
474 self.config.genesis_timestamp,
475 slot,
476 )?;
477
478 let slot_prev_ts = get_block_slot_timestamp(
479 self.config.thread_count,
480 self.config.t0,
481 self.config.genesis_timestamp,
482 slot.get_prev_slot(self.config.thread_count)?,
483 )?;
484
485 let mut mip_guard = self.mip_store.0.write();
489 mip_guard.update_network_version_stats(slot_ts, network_versions);
490 mip_guard.update_batches(
491 &mut db_batch,
492 &mut db_versioning_batch,
493 Some((&slot_prev_ts, &slot_ts)),
494 )?;
495
496 if let SetOrKeep::Set(new_hash) = changes.execution_trail_hash_change {
498 db_batch.insert(
499 EXECUTION_TRAIL_HASH_PREFIX.as_bytes().to_vec(),
500 Some(new_hash.to_bytes().to_vec()),
501 );
502 }
503
504 self.db
505 .write()
506 .write_batch(db_batch, db_versioning_batch, Some(slot));
507 drop(mip_guard);
508
509 let final_state_hash = self.db.read().get_xof_db_hash();
510
511 info!("final_state hash at slot {}: {}", slot, final_state_hash);
513
514 #[cfg(feature = "bootstrap_server")]
516 if slot.period % self.config.ledger_backup_periods_interval == 0
517 && slot.period != 0
518 && slot.thread == 0
519 {
520 let state_slot = self.db.read().get_change_id();
521 match state_slot {
522 Ok(slot) => {
523 info!(
524 "Backuping db for slot {}, state slot: {}, state hash: {}",
525 slot, slot, final_state_hash
526 );
527 }
528 Err(e) => {
529 info!("{}", e);
530 info!(
531 "Backuping db for unknown state slot, state hash: {}",
532 final_state_hash
533 );
534 }
535 }
536
537 self.db.read().backup_db(slot);
538 }
539
540 let cycle = slot.get_cycle(self.config.periods_per_cycle);
542 self.pos_state
543 .feed_cycle_state_hash(cycle, final_state_hash);
544
545 Ok(())
546 }
547
548 fn _is_db_valid(&self) -> AnyResult<()> {
550 let db = self.db.read();
551
552 {
554 let execution_trail_hash_serialized =
555 match db.get_cf(STATE_CF, EXECUTION_TRAIL_HASH_PREFIX.as_bytes().to_vec()) {
556 Ok(Some(v)) => v,
557 Ok(None) => {
558 return Err(anyhow!("No execution trail hash found in DB"));
559 }
560 Err(err) => {
561 return Err(err.into());
562 }
563 };
564 if let Err(err) = massa_hash::Hash::try_from(&execution_trail_hash_serialized[..]) {
565 warn!("Invalid execution trail hash found in DB: {}", err);
566 return Err(err.into());
567 }
568 }
569
570 for (serialized_key, serialized_value) in
571 db.iterator_cf_for_full_db_traversal(STATE_CF, MassaIteratorMode::Start)
572 {
573 #[allow(clippy::if_same_then_else)]
574 if serialized_key.starts_with(CYCLE_HISTORY_PREFIX.as_bytes()) {
575 if !self
576 .pos_state
577 .is_cycle_history_key_value_valid(&serialized_key, &serialized_value)
578 {
579 warn!(
580 "Wrong key/value for CYCLE_HISTORY_KEY PREFIX serialized_key: {:?}, serialized_value: {:?}",
581 serialized_key, serialized_value
582 );
583 return Err(anyhow!(
584 "Wrong key/value for CYCLE_HISTORY_KEY PREFIX serialized_key: {:?}, serialized_value: {:?}",
585 serialized_key, serialized_value
586 ));
587 }
588 } else if serialized_key.starts_with(DEFERRED_CREDITS_PREFIX.as_bytes()) {
589 if !self
590 .pos_state
591 .is_deferred_credits_key_value_valid(&serialized_key, &serialized_value)
592 {
593 warn!(
594 "Wrong key/value for DEFERRED_CREDITS PREFIX serialized_key: {:?}, serialized_value: {:?}",
595 serialized_key, serialized_value
596 );
597 return Err(anyhow!(
598 "Wrong key/value for DEFERRED_CREDITS PREFIX serialized_key: {:?}, serialized_value: {:?}",
599 serialized_key, serialized_value
600 ));
601 }
602 } else if serialized_key.starts_with(ASYNC_POOL_PREFIX.as_bytes()) {
603 if !self
604 .async_pool
605 .is_key_value_valid(&serialized_key, &serialized_value)
606 {
607 warn!(
608 "Wrong key/value for ASYNC_POOL PREFIX serialized_key: {:?}, serialized_value: {:?}",
609 serialized_key, serialized_value
610 );
611 return Err(anyhow!(
612 "Wrong key/value for ASYNC_POOL PREFIX serialized_key: {:?}, serialized_value: {:?}",
613 serialized_key, serialized_value
614 ));
615 }
616 } else if serialized_key.starts_with(EXECUTED_OPS_PREFIX.as_bytes()) {
617 if !self
618 .executed_ops
619 .is_key_value_valid(&serialized_key, &serialized_value)
620 {
621 warn!(
622 "Wrong key/value for EXECUTED_OPS PREFIX serialized_key: {:?}, serialized_value: {:?}",
623 serialized_key, serialized_value
624 );
625 return Err(anyhow!(
626 "Wrong key/value for EXECUTED_OPS PREFIX serialized_key: {:?}, serialized_value: {:?}",
627 serialized_key, serialized_value
628 ));
629 }
630 } else if serialized_key.starts_with(EXECUTED_DENUNCIATIONS_PREFIX.as_bytes()) {
631 if !self
632 .executed_denunciations
633 .is_key_value_valid(&serialized_key, &serialized_value)
634 {
635 warn!("Wrong key/value for EXECUTED_DENUNCIATIONS PREFIX serialized_key: {:?}, serialized_value: {:?}", serialized_key, serialized_value);
636 return Err(anyhow!(
637 "Wrong key/value for EXECUTED_DENUNCIATIONS PREFIX serialized_key: {:?}, serialized_value: {:?}",
638 serialized_key, serialized_value
639 ));
640 }
641 } else if serialized_key.starts_with(LEDGER_PREFIX.as_bytes()) {
642 if !self
643 .ledger
644 .is_key_value_valid(&serialized_key, &serialized_value)
645 {
646 warn!("Wrong key/value for LEDGER PREFIX serialized_key: {:?}, serialized_value: {:?}", serialized_key, serialized_value);
647 return Err(anyhow!(
648 "Wrong key/value for LEDGER PREFIX serialized_key: {:?}, serialized_value: {:?}",
649 serialized_key, serialized_value
650 ));
651 }
652 } else if serialized_key.starts_with(MIP_STORE_PREFIX.as_bytes()) {
653 if !self
654 .mip_store
655 .is_key_value_valid(&serialized_key, &serialized_value)
656 {
657 warn!("Wrong key/value for MIP Store");
658 return Err(anyhow!(
659 "Wrong key/value for MIP Store serialized_key: {:?}, serialized_value: {:?}",
660 serialized_key, serialized_value
661 ));
662 }
663 } else if serialized_key.starts_with(EXECUTION_TRAIL_HASH_PREFIX.as_bytes()) {
664 } else if serialized_key.starts_with(DEFERRED_CALLS_PREFIX.as_bytes())
666 || serialized_key.eq(DEFERRED_CALL_TOTAL_GAS.as_bytes())
667 {
668 if !self
669 .deferred_call_registry
670 .is_key_value_valid(&serialized_key, &serialized_value)
671 {
672 warn!("Wrong key/value for DEFERRED_CALLS_PREFIX serialized_key: {:?}, serialized_value: {:?}", serialized_key, serialized_value);
673 return Err(anyhow!(
674 "Wrong key/value for DEFERRED_CALLS_PREFIX serialized_key: {:?}, serialized_value: {:?}",
675 serialized_key, serialized_value
676 ));
677 }
678 } else {
679 warn!(
680 "Key/value does not correspond to any prefix: serialized_key: {:?}, serialized_value: {:?}",
681 serialized_key, serialized_value
682 );
683 return Err(anyhow!(
684 "Key/value does not correspond to any prefix: serialized_key: {:?}, serialized_value: {:?}",
685 serialized_key, serialized_value
686 ));
687 }
688 }
689
690 for (serialized_key, serialized_value) in
691 db.iterator_cf(VERSIONING_CF, MassaIteratorMode::Start)
692 {
693 if serialized_key.starts_with(MIP_STORE_PREFIX.as_bytes())
694 || serialized_key.starts_with(MIP_STORE_STATS_PREFIX.as_bytes())
695 {
696 if !self
697 .mip_store
698 .is_key_value_valid(&serialized_key, &serialized_value)
699 {
700 warn!("Wrong key/value for MIP Store");
701 return Err(anyhow!(
702 "Wrong key/value for MIP Store serialized_key: {:?}, serialized_value: {:?}",
703 serialized_key, serialized_value
704 ));
705 }
706 } else {
707 warn!(
708 "Key/value does not correspond to any prefix: serialized_key: {:?}, serialized_value: {:?}",
709 serialized_key, serialized_value
710 );
711 return Err(anyhow!(
712 "Key/value does not correspond to any prefix: serialized_key: {:?}, serialized_value: {:?}",
713 serialized_key, serialized_value
714 ));
715 }
716 }
717
718 if let Err(err) = self.pos_state.validate_selector_history() {
720 warn!("Invalid PoS cycle history for selector feeding: {}", err);
721 return Err(anyhow!(
722 "Invalid PoS cycle history for selector feeding: {}",
723 err
724 ));
725 }
726
727 Ok(())
728 }
729
730 pub fn new_derived_from_snapshot(
739 db: ShareableMassaDBController,
740 config: FinalStateConfig,
741 ledger: Box<dyn LedgerController>,
742 selector: Box<dyn SelectorController>,
743 mip_store: MipStore,
744 last_start_period: u64,
745 ) -> Result<Self, FinalStateError> {
746 info!("Restarting from snapshot");
747
748 let mut final_state =
749 FinalState::new(db, config.clone(), ledger, selector, mip_store, false)?;
750
751 let recovered_slot =
752 final_state.db.read().get_change_id().map_err(|_| {
753 FinalStateError::InvalidSlot(String::from("Could not get slot in db"))
754 })?;
755
756 if cfg!(feature = "test-exports") {
758 let mut batch = DBBatch::new();
759 final_state.pos_state.create_initial_cycle(&mut batch);
760 final_state
761 .db
762 .write()
763 .write_batch(batch, Default::default(), Some(recovered_slot));
764 }
765
766 final_state.last_slot_before_downtime = Some(recovered_slot);
767
768 let shutdown_start = recovered_slot
771 .get_next_slot(config.thread_count)
772 .map_err(|e| {
773 FinalStateError::InvalidSlot(format!(
774 "Unable to get next slot from recovered slot: {:?}",
775 e
776 ))
777 })?;
778 let shutdown_end = Slot::new(last_start_period, config.thread_count.saturating_sub(1));
782 debug!(
783 "Checking if MIP store is consistent against shutdown period: {} - {}",
784 shutdown_start, shutdown_end
785 );
786
787 final_state
788 .mip_store
789 .is_consistent_with_shutdown_period(
790 shutdown_start,
791 shutdown_end,
792 config.thread_count,
793 config.t0,
794 config.genesis_timestamp,
795 )
796 .map_err(FinalStateError::from)?;
797
798 debug!(
799 "Latest consistent slot found in snapshot data: {}",
800 recovered_slot
801 );
802
803 info!(
804 "final_state hash at slot {}: {}",
805 recovered_slot,
806 final_state.db.read().get_xof_db_hash()
807 );
808
809 final_state.last_start_period = last_start_period;
811
812 final_state.recompute_caches();
813
814 final_state.compute_initial_draws()?;
816
817 final_state.interpolate_downtime()?;
818
819 Ok(final_state)
820 }
821}
822
823impl FinalStateController for FinalState {
824 fn compute_initial_draws(&mut self) -> Result<(), FinalStateError> {
825 self.pos_state
826 .compute_initial_draws()
827 .map_err(|err| FinalStateError::PosError(err.to_string()))
828 }
829
830 fn finalize(
831 &mut self,
832 slot: Slot,
833 changes: StateChanges,
834 network_versions: Option<(u32, Option<u32>)>,
835 ) {
836 self._finalize(slot, changes, network_versions).unwrap()
837 }
838
839 fn get_execution_trail_hash(&self) -> Hash {
840 let hash_bytes = self
841 .db
842 .read()
843 .get_cf(STATE_CF, EXECUTION_TRAIL_HASH_PREFIX.as_bytes().to_vec())
844 .expect("could not read execution trail hash from state DB")
845 .expect("could not find execution trail hash in state DB");
846 Hash::from_bytes(
847 hash_bytes
848 .as_slice()
849 .try_into()
850 .expect("invalid execution trail hash in state DB"),
851 )
852 }
853
854 fn get_fingerprint(&self) -> Hash {
855 let internal_hash = self.db.read().get_xof_db_hash();
856 Hash::compute_from(internal_hash.to_bytes())
857 }
858
859 fn get_slot(&self) -> Slot {
860 self.db
861 .read()
862 .get_change_id()
863 .expect("Critical error: Final state has no slot attached")
864 }
865
866 fn init_execution_trail_hash_to_batch(&mut self, batch: &mut DBBatch) {
867 batch.insert(
868 EXECUTION_TRAIL_HASH_PREFIX.as_bytes().to_vec(),
869 Some(massa_hash::Hash::zero().to_bytes().to_vec()),
870 );
871 }
872
873 fn is_db_valid(&self) -> bool {
874 self._is_db_valid().is_ok()
875 }
876
877 fn recompute_caches(&mut self) {
878 self.async_pool.recompute_message_cache();
879 self.executed_ops.recompute_sorted_ops_and_op_exec_status();
880 self.executed_denunciations.recompute_sorted_denunciations();
881 self.pos_state.recompute_pos_state_caches();
882 }
883
884 fn reset(&mut self) {
885 let slot = Slot::new(0, self.config.thread_count.saturating_sub(1));
886 self.db.write().reset_slot_and_history(slot);
887 self.ledger.reset();
888 self.async_pool.reset();
889 self.pos_state.reset();
890 self.executed_ops.reset();
891 self.executed_denunciations.reset();
892 self.mip_store.reset_db(self.db.clone());
893 self.db.write().delete_prefix("", STATE_CF, None);
895 self.db.write().delete_prefix("", VERSIONING_CF, None);
896 }
897
898 fn get_ledger(&self) -> &Box<dyn LedgerController> {
899 &self.ledger
900 }
901
902 fn get_ledger_mut(&mut self) -> &mut Box<dyn LedgerController> {
903 &mut self.ledger
904 }
905
906 fn get_async_pool(&self) -> &AsyncPool {
907 &self.async_pool
908 }
909
910 fn get_pos_state(&self) -> &PoSFinalState {
911 &self.pos_state
912 }
913
914 fn get_pos_state_mut(&mut self) -> &mut PoSFinalState {
915 &mut self.pos_state
916 }
917
918 fn load_initial_deferred_credits(
919 &mut self,
920 batch: &mut DBBatch,
921 ) -> Result<(), massa_pos_exports::PosError> {
922 self.pos_state
923 .load_initial_deferred_credits(batch, self.ledger.as_ref())
924 }
925
926 fn executed_ops_contains(&self, op_id: &OperationId) -> bool {
927 self.executed_ops.contains(op_id)
928 }
929
930 fn get_ops_exec_status(&self, batch: &[OperationId]) -> Vec<Option<bool>> {
931 self.executed_ops.get_ops_exec_status(batch)
932 }
933
934 fn get_executed_denunciations(&self) -> &ExecutedDenunciations {
935 &self.executed_denunciations
936 }
937
938 fn get_database(&self) -> &ShareableMassaDBController {
939 &self.db
940 }
941
942 fn get_last_start_period(&self) -> u64 {
943 self.last_start_period
944 }
945
946 fn set_last_start_period(&mut self, last_start_period: u64) {
947 self.last_start_period = last_start_period;
948 }
949
950 fn get_last_slot_before_downtime(&self) -> &Option<Slot> {
951 &self.last_slot_before_downtime
952 }
953
954 fn set_last_slot_before_downtime(&mut self, last_slot_before_downtime: Option<Slot>) {
955 self.last_slot_before_downtime = last_slot_before_downtime;
956 }
957
958 fn get_mip_store_mut(&mut self) -> &mut MipStore {
959 &mut self.mip_store
960 }
961
962 fn get_mip_store(&self) -> &MipStore {
963 &self.mip_store
964 }
965
966 fn get_deferred_call_registry(&self) -> &DeferredCallRegistry {
967 &self.deferred_call_registry
968 }
969}
970
971#[cfg(test)]
972mod test {
973 use std::collections::BTreeMap;
974 use std::path::PathBuf;
975 use std::str::FromStr;
976 use std::sync::Arc;
977
978 use massa_deferred_calls::config::DeferredCallsConfig;
979 use num::rational::Ratio;
980 use parking_lot::RwLock;
981 use tempfile::tempdir;
982
983 use massa_async_pool::{AsyncPoolChanges, AsyncPoolConfig};
984 use massa_db_exports::{MassaDBConfig, MassaDBController, STATE_HASH_INITIAL_BYTES};
985 use massa_db_worker::MassaDB;
986 use massa_executed_ops::{ExecutedDenunciationsConfig, ExecutedOpsConfig};
987 use massa_hash::Hash;
988 use massa_ledger_exports::{LedgerChanges, LedgerConfig, LedgerEntryUpdate};
989 use massa_ledger_worker::FinalLedger;
990 use massa_models::address::Address;
991 use massa_models::{
992 amount::Amount,
993 async_msg::AsyncMessage,
994 bytecode::Bytecode,
995 config::{
996 DENUNCIATION_EXPIRE_PERIODS, ENDORSEMENT_COUNT, KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
997 MAX_ASYNC_POOL_LENGTH, MAX_BYTECODE_LENGTH, MAX_DATASTORE_KEY_LENGTH,
998 MAX_DATASTORE_VALUE_LENGTH, MAX_DEFERRED_CREDITS_LENGTH,
999 MAX_DENUNCIATIONS_PER_BLOCK_HEADER, MAX_DENUNCIATION_CHANGES_LENGTH,
1000 MAX_FUNCTION_NAME_LENGTH, MAX_PARAMETERS_SIZE, MAX_PRODUCTION_STATS_LENGTH,
1001 MAX_ROLLS_COUNT_LENGTH, MIP_STORE_STATS_BLOCK_CONSIDERED, PERIODS_PER_CYCLE,
1002 POS_SAVED_CYCLES, T0, THREAD_COUNT,
1003 },
1004 types::SetUpdateOrDelete,
1005 };
1006 use massa_pos_exports::MockSelectorController;
1007 use massa_pos_exports::{PoSChanges, PoSConfig, PosError};
1008 use massa_time::MassaTime;
1009 use massa_versioning::versioning::MipStatsConfig;
1010
1011 use super::*;
1012
1013 fn get_final_state_config() -> (FinalStateConfig, LedgerConfig) {
1014 let massa_node_base = PathBuf::from("../massa-node");
1015
1016 let genesis_timestamp = MassaTime::from_millis(0);
1017 let ledger_config = LedgerConfig {
1018 thread_count: THREAD_COUNT,
1019 initial_ledger_path: massa_node_base.join("base_config/initial_ledger.json"),
1020 max_key_length: MAX_DATASTORE_KEY_LENGTH,
1021 max_datastore_value_length: MAX_DATASTORE_VALUE_LENGTH,
1022 max_bytecode_size: MAX_BYTECODE_LENGTH,
1023 };
1024 let async_pool_config = AsyncPoolConfig {
1025 max_length: MAX_ASYNC_POOL_LENGTH,
1026 max_function_length: MAX_FUNCTION_NAME_LENGTH,
1027 max_function_params_length: MAX_PARAMETERS_SIZE as u64,
1028 thread_count: THREAD_COUNT,
1029 max_key_length: MAX_DATASTORE_KEY_LENGTH as u32,
1030 };
1031 let pos_config = PoSConfig {
1032 periods_per_cycle: PERIODS_PER_CYCLE,
1033 thread_count: THREAD_COUNT,
1034 cycle_history_length: POS_SAVED_CYCLES,
1035 max_rolls_length: MAX_ROLLS_COUNT_LENGTH,
1036 max_production_stats_length: MAX_PRODUCTION_STATS_LENGTH,
1037 max_credit_length: MAX_DEFERRED_CREDITS_LENGTH,
1038 initial_deferred_credits_path: Some(
1039 massa_node_base.join("base_config/deferred_credits.json"),
1040 ),
1041 };
1042 let executed_ops_config = ExecutedOpsConfig {
1043 thread_count: THREAD_COUNT,
1044 keep_executed_history_extra_periods: KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
1045 };
1046 let executed_denunciations_config = ExecutedDenunciationsConfig {
1047 denunciation_expire_periods: DENUNCIATION_EXPIRE_PERIODS,
1048 thread_count: THREAD_COUNT,
1049 endorsement_count: ENDORSEMENT_COUNT,
1050 keep_executed_history_extra_periods: KEEP_EXECUTED_HISTORY_EXTRA_PERIODS,
1051 };
1052
1053 let deferred_calls_config = DeferredCallsConfig::default();
1054 let final_state_config = FinalStateConfig {
1055 ledger_config: ledger_config.clone(),
1056 async_pool_config,
1057 deferred_calls_config,
1058 pos_config,
1059 executed_ops_config,
1060 executed_denunciations_config,
1061 final_history_length: 100, thread_count: THREAD_COUNT,
1063 periods_per_cycle: PERIODS_PER_CYCLE,
1064 initial_seed_string: "test".to_string(),
1065 initial_rolls_path: massa_node_base.join("base_config/initial_rolls.json"), endorsement_count: ENDORSEMENT_COUNT,
1067 max_executed_denunciations_length: MAX_DENUNCIATION_CHANGES_LENGTH,
1068 max_denunciations_per_block_header: MAX_DENUNCIATIONS_PER_BLOCK_HEADER,
1069 t0: T0,
1070 ledger_backup_periods_interval: 10,
1071 genesis_timestamp,
1072 };
1073
1074 (final_state_config, ledger_config)
1075 }
1076
1077 fn get_final_state() -> FinalState {
1078 let (final_state_config, ledger_config) = get_final_state_config();
1079
1080 let temp_dir_db = tempdir().expect("Unable to create a temp folder");
1081 let db_config = MassaDBConfig {
1084 path: temp_dir_db.path().to_path_buf(),
1085 max_history_length: 100,
1086 max_final_state_elements_size: 100,
1087 max_versioning_elements_size: 100,
1088 thread_count: THREAD_COUNT,
1089 max_ledger_backups: 10,
1090 enable_metrics: false,
1091 };
1092 let db = Arc::new(RwLock::new(
1093 Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>
1094 ));
1095
1096 let mip_stats_config = MipStatsConfig {
1097 block_count_considered: MIP_STORE_STATS_BLOCK_CONSIDERED,
1098 warn_announced_version_ratio: Ratio::new_raw(30, 100), };
1100 let mip_store =
1101 MipStore::try_from(([], mip_stats_config)).expect("Cannot create an empty MIP store");
1102
1103 let selector_controller = Box::new(MockSelectorController::new());
1104 let ledger = FinalLedger::new(ledger_config, db.clone());
1105
1106 FinalState::new(
1107 db,
1108 final_state_config,
1109 Box::new(ledger),
1110 selector_controller,
1111 mip_store,
1112 false,
1113 )
1114 .expect("Cannot init final state")
1115 }
1116
1117 fn get_state_changes() -> StateChanges {
1118 let mut state_changes = StateChanges::default();
1119 let message = AsyncMessage::new(
1120 Slot::new(1, 0),
1121 0,
1122 Address::from_str("AU12dG5xP1RDEB5ocdHkymNVvvSJmUL9BgHwCksDowqmGWxfpm93x").unwrap(),
1123 Address::from_str("AU12htxRWiEm8jDJpJptr6cwEhWNcCSFWstN1MLSa96DDkVM9Y42G").unwrap(),
1124 String::from("test"),
1125 10000000,
1126 Amount::from_str("1").unwrap(),
1127 Amount::from_str("1").unwrap(),
1128 Slot::new(2, 0),
1129 Slot::new(3, 0),
1130 vec![1, 2, 3, 4],
1131 None,
1132 None,
1133 );
1134 let mut async_pool_changes = AsyncPoolChanges::default();
1135 async_pool_changes
1136 .0
1137 .insert(message.compute_id(), SetUpdateOrDelete::Set(message));
1138 state_changes.async_pool_changes = async_pool_changes;
1139
1140 let amount = Amount::from_str("1").unwrap();
1141 let bytecode = Bytecode(vec![1, 2, 3]);
1142 let ledger_entry = LedgerEntryUpdate {
1143 balance: SetOrKeep::Set(amount),
1144 bytecode: SetOrKeep::Set(bytecode),
1145 datastore: BTreeMap::default(),
1146 };
1147 let mut ledger_changes = LedgerChanges::default();
1148 ledger_changes.0.insert(
1149 Address::from_str("AU12dG5xP1RDEB5ocdHkymNVvvSJmUL9BgHwCksDowqmGWxfpm93x").unwrap(),
1150 SetUpdateOrDelete::Update(ledger_entry),
1151 );
1152 state_changes.ledger_changes = ledger_changes;
1153
1154 let mut pos_changes = PoSChanges::default();
1155 pos_changes.roll_changes.insert(
1156 Address::from_str("AU12r1iM79EcS3sa4dmtUp28TiaPxK1weQcLsATcFoynPdukjdMqM").unwrap(),
1157 0,
1158 );
1159 state_changes.pos_changes = pos_changes;
1160
1161 state_changes
1162 }
1163
1164 #[test]
1165 fn test_final_state_finalize() {
1166 let mut fstate = get_final_state();
1171 let initial_slot = Slot::new(0, 0);
1172 assert_eq!(fstate.get_slot(), initial_slot);
1173
1174 let wrong_next_slot = Slot::new(1, 1);
1175 let ok_next_slot = Slot::new(0, 1);
1176 let changes = get_state_changes();
1177
1178 let res = fstate._finalize(wrong_next_slot, changes.clone(), None);
1179 assert!(res
1180 .err()
1181 .unwrap()
1182 .to_string()
1183 .contains("while the current slot is"));
1184
1185 assert_eq!(fstate.get_slot(), initial_slot);
1186
1187 let res = fstate._finalize(ok_next_slot, changes.clone(), None);
1189
1190 assert!(res.is_err());
1191 match res {
1192 Ok(_) => unreachable!(),
1193 Err(e) => match e.downcast_ref() {
1194 Some(PosError::ContainerInconsistency(_)) => {
1195 println!("Received correct error: PosError:ContainerInconsistency...");
1196 }
1197 Some(_) => {
1198 panic!("Unknown error received: {}", e);
1199 }
1200 None => unreachable!(),
1201 },
1202 }
1203
1204 let mut batch = DBBatch::new();
1205 fstate.pos_state.create_initial_cycle(&mut batch);
1206 let res = fstate._finalize(ok_next_slot, changes, None);
1207 assert!(res.is_ok());
1208 assert_eq!(fstate.get_slot(), ok_next_slot);
1209 }
1210
1211 #[test]
1212 fn test_final_state_from_snapshot_1() {
1213 let mut fstate = get_final_state();
1218 let ok_next_slot = Slot::new(0, 1);
1219 let changes = get_state_changes();
1220 let mut batch = DBBatch::new();
1221 fstate.pos_state.create_initial_cycle(&mut batch);
1222 let res = fstate._finalize(ok_next_slot, changes, None);
1223 assert!(res.is_ok());
1224 assert_eq!(fstate.get_slot(), ok_next_slot);
1225
1226 let db_2 = fstate.db.clone();
1227 let config_2 = fstate.config.clone();
1228 let ledger_2 = FinalLedger::new(fstate.config.ledger_config.clone(), db_2.clone());
1229 let mut selector_controller = Box::new(MockSelectorController::new());
1230 selector_controller
1232 .expect_feed_cycle()
1233 .returning(|_, _, _| Ok(()));
1234 selector_controller
1235 .expect_wait_for_draws()
1236 .returning(|_| Ok(1));
1237 let mip_stats_config = MipStatsConfig {
1238 block_count_considered: MIP_STORE_STATS_BLOCK_CONSIDERED,
1239 warn_announced_version_ratio: Ratio::new_raw(30, 100), };
1241 let mip_store =
1242 MipStore::try_from(([], mip_stats_config)).expect("Cannot create an empty MIP store");
1243
1244 let last_start_period_2 = 2;
1245 let fstate2_ = FinalState::new_derived_from_snapshot(
1246 db_2,
1247 config_2,
1248 Box::new(ledger_2),
1249 selector_controller,
1250 mip_store,
1251 last_start_period_2,
1252 );
1253
1254 assert!(fstate2_.is_ok());
1255 let mut fstate2 = fstate2_.unwrap();
1256 assert_eq!(fstate2.last_slot_before_downtime, Some(ok_next_slot));
1257 assert_eq!(
1258 fstate2.get_slot(),
1259 Slot::new(last_start_period_2, THREAD_COUNT - 1)
1260 );
1261
1262 assert!(!fstate2.is_db_valid()); fstate2.init_execution_trail_hash_to_batch(&mut batch);
1264 fstate2
1265 .db
1266 .write()
1267 .write_batch(batch, Default::default(), None);
1268 assert!(fstate2.is_db_valid());
1269 }
1270
1271 #[test]
1272 fn test_final_state_from_snapshot_2() {
1273 let mut fstate = get_final_state();
1278 let ok_next_slot = Slot::new(0, 1);
1279 let changes = get_state_changes();
1280 let mut batch = DBBatch::new();
1281 fstate.pos_state.create_initial_cycle(&mut batch);
1282 let res = fstate._finalize(ok_next_slot, changes, None);
1283 assert!(res.is_ok());
1284 assert_eq!(fstate.get_slot(), ok_next_slot);
1285
1286 let db_2 = fstate.db.clone();
1287 let config_2 = fstate.config.clone();
1288 let ledger_2 = FinalLedger::new(fstate.config.ledger_config.clone(), db_2.clone());
1289 let mut selector_controller = Box::new(MockSelectorController::new());
1290 selector_controller
1291 .expect_feed_cycle()
1292 .times(4)
1293 .returning(|_, _, _| Ok(()));
1294 selector_controller
1295 .expect_wait_for_draws()
1296 .returning(|_| Ok(1));
1297 let mip_stats_config = MipStatsConfig {
1298 block_count_considered: MIP_STORE_STATS_BLOCK_CONSIDERED,
1299 warn_announced_version_ratio: Ratio::new_raw(30, 100), };
1301 let mip_store =
1302 MipStore::try_from(([], mip_stats_config)).expect("Cannot create an empty MIP store");
1303
1304 let last_start_period_2 = 2 + (PERIODS_PER_CYCLE * 2);
1305 let fstate2_ = FinalState::new_derived_from_snapshot(
1306 db_2,
1307 config_2,
1308 Box::new(ledger_2),
1309 selector_controller,
1310 mip_store,
1311 last_start_period_2,
1312 );
1313
1314 assert!(fstate2_.is_ok());
1315 let mut fstate2 = fstate2_.unwrap();
1316 assert_eq!(fstate2.last_slot_before_downtime, Some(ok_next_slot));
1317 assert_eq!(
1318 fstate2.get_slot(),
1319 Slot::new(last_start_period_2, THREAD_COUNT - 1)
1320 );
1321
1322 assert!(!fstate2.is_db_valid()); fstate2.init_execution_trail_hash_to_batch(&mut batch);
1324 fstate2
1325 .db
1326 .write()
1327 .write_batch(batch, Default::default(), None);
1328 assert!(fstate2.is_db_valid());
1329 }
1330
1331 #[test]
1332 fn test_final_state_reset() {
1333 let mut fstate = get_final_state();
1341
1342 let db_valid = fstate._is_db_valid();
1343 assert!(db_valid.is_err());
1344 assert!(db_valid
1345 .err()
1346 .unwrap()
1347 .to_string()
1348 .starts_with("No execution trail hash"));
1349 assert_eq!(
1351 fstate.get_fingerprint(),
1352 Hash::compute_from(STATE_HASH_INITIAL_BYTES)
1353 );
1354
1355 let mut batch = DBBatch::new();
1356 fstate.init_execution_trail_hash_to_batch(&mut batch);
1357 fstate
1358 .db
1359 .write()
1360 .write_batch(batch, Default::default(), None);
1361
1362 let db_valid = fstate._is_db_valid();
1363 assert!(db_valid.is_ok());
1364
1365 assert_ne!(
1367 fstate.get_fingerprint(),
1368 Hash::compute_from(STATE_HASH_INITIAL_BYTES)
1369 );
1370
1371 fstate.reset();
1372
1373 let db_valid = fstate._is_db_valid();
1374 assert!(db_valid.is_err());
1376 assert_eq!(fstate.get_slot().period, 0);
1377 assert_eq!(
1379 fstate.get_fingerprint(),
1380 Hash::compute_from(STATE_HASH_INITIAL_BYTES)
1381 );
1382 }
1383}