1use humantime::format_duration;
2use massa_consensus_exports::bootstrapable_graph::BootstrapableGraph;
3use massa_db_exports::DBBatch;
4use massa_final_state::{FinalStateController, FinalStateError};
5use massa_logging::massa_trace;
6use massa_metrics::MassaMetrics;
7use massa_models::{
8 block_id::BlockId, node::NodeId, prehash::PreHashSet, slot::Slot,
9 streaming_step::StreamingStep, timeslots::get_block_slot_timestamp, version::Version,
10};
11use massa_signature::PublicKey;
12use massa_time::MassaTime;
13use massa_versioning::versioning::{ComponentStateTypeId, MipInfo, MipState, StateAtError};
14use parking_lot::RwLock;
15use rand::{
16 prelude::{SliceRandom, StdRng},
17 SeedableRng,
18};
19use std::collections::BTreeMap;
20use std::{
21 collections::HashSet,
22 io,
23 net::{SocketAddr, TcpStream},
24 sync::{Arc, Condvar, Mutex},
25 time::Duration,
26};
27use tracing::{debug, info, warn};
28
29use crate::{
30 bindings::BootstrapClientBinder,
31 error::BootstrapError,
32 messages::{BootstrapClientMessage, BootstrapServerMessage},
33 settings::IpType,
34 BootstrapConfig, GlobalBootstrapState,
35};
36
37#[cfg_attr(test, mockall::automock)]
39pub trait BSConnector {
40 fn connect_timeout(
43 &self,
44 addr: SocketAddr,
45 duration: Option<MassaTime>,
46 ) -> io::Result<TcpStream>;
47}
48
49#[derive(Debug)]
51pub struct DefaultConnector;
52
53impl BSConnector for DefaultConnector {
54 fn connect_timeout(
59 &self,
60 addr: SocketAddr,
61 duration: Option<MassaTime>,
62 ) -> io::Result<TcpStream> {
63 let Some(duration) = duration else {
64 return TcpStream::connect(addr);
65 };
66 TcpStream::connect_timeout(&addr, duration.to_duration())
67 }
68}
69pub(crate) fn stream_final_state_and_consensus(
77 cfg: &BootstrapConfig,
78 client: &mut BootstrapClientBinder,
79 next_bootstrap_message: &mut BootstrapClientMessage,
80 global_bootstrap_state: &mut GlobalBootstrapState,
81) -> Result<(), BootstrapError> {
82 align_consensus_resume_state_before_ask(next_bootstrap_message, global_bootstrap_state);
83
84 if let BootstrapClientMessage::AskBootstrapPart {
85 send_last_start_period: false,
86 ..
87 } = &next_bootstrap_message
88 {
89 client.set_last_start_period(Some(
92 global_bootstrap_state
93 .final_state
94 .read()
95 .get_last_start_period(),
96 ));
97 } else if let BootstrapClientMessage::AskBootstrapPart {
98 send_last_start_period: true,
99 ..
100 } = &next_bootstrap_message
101 {
102 client.set_last_start_period(None);
103 }
104
105 if let BootstrapClientMessage::AskBootstrapPart { .. } = &next_bootstrap_message {
106 client.send_timeout(
107 next_bootstrap_message,
108 Some(cfg.write_timeout.to_duration()),
109 )?;
110
111 loop {
112 match client.next_timeout(Some(cfg.read_timeout.to_duration()))? {
113 BootstrapServerMessage::BootstrapPart {
114 slot,
115 state_part,
116 versioning_part,
117 consensus_part,
118 consensus_outdated_ids,
119 last_start_period,
120 last_slot_before_downtime,
121 } => {
122 let mut write_final_state = global_bootstrap_state.final_state.write();
124
125 check_restart_metadata(cfg, &last_start_period, &last_slot_before_downtime)?;
130
131 if let Some(last_start_period) = last_start_period {
132 write_final_state.set_last_start_period(last_start_period);
133 client.set_last_start_period(Some(last_start_period));
134 }
135 if let Some(last_slot_before_downtime) = last_slot_before_downtime {
136 write_final_state.set_last_slot_before_downtime(last_slot_before_downtime);
137 }
138
139 let (last_state_step, last_versioning_step) = write_final_state
140 .get_database()
141 .write()
142 .write_batch_bootstrap_client(state_part, versioning_part)
143 .map_err(|e| {
144 BootstrapError::GeneralError(format!(
145 "Cannot write received stream batch to disk: {}",
146 e
147 ))
148 })?;
149
150 if let Some(graph) = global_bootstrap_state.graph.as_mut() {
152 graph.final_blocks.extend(consensus_part.final_blocks);
154 graph.final_blocks.retain(|block_export| {
156 !consensus_outdated_ids.contains(&block_export.block.id)
157 });
158 } else {
159 global_bootstrap_state.graph = Some(consensus_part);
160 }
161
162 let last_consensus_step = capped_consensus_resume_step(
164 global_bootstrap_state.graph.as_ref(),
165 cfg.max_consensus_block_ids as usize,
166 );
167
168 *next_bootstrap_message = BootstrapClientMessage::AskBootstrapPart {
170 last_slot: Some(slot),
171 last_state_step,
172 last_versioning_step,
173 last_consensus_step,
174 send_last_start_period: false,
175 };
176
177 debug!(
179 "client final state bootstrap cursors: {:?}",
180 next_bootstrap_message
181 );
182 }
183 BootstrapServerMessage::BootstrapFinished => {
184 info!("State bootstrap complete");
185
186 let missing = {
193 let guard = global_bootstrap_state.final_state.read();
194 let bootstrapped_at = get_block_slot_timestamp(
195 cfg.thread_count,
196 cfg.t0,
197 cfg.genesis_timestamp,
198 guard.get_slot(),
199 )?;
200 guard.get_mip_store().started_mips_missing_from_db(
201 &guard.get_database().clone(),
202 bootstrapped_at,
203 )
204 };
205 if !missing.is_empty() {
206 let names: Vec<_> = missing.iter().map(|mi| mi.name.as_str()).collect();
207 restart_state_streaming(
208 client,
209 next_bootstrap_message,
210 global_bootstrap_state,
211 );
212 return Err(BootstrapError::GeneralError(format!(
213 "bootstrap server does not know {}, whose vote has already started: \
214 it cannot provide its state",
215 names.join(", ")
216 )));
217 }
218
219 let mut guard = global_bootstrap_state.final_state.write();
221 let db = guard.get_database().clone();
222 let (updated, added) = guard
223 .get_mip_store_mut()
224 .extend_from_db(db)
225 .map_err(|e| BootstrapError::from(FinalStateError::from(e)))?;
226
227 warn_user_about_versioning_updates(updated, added);
228
229 if let Some(last_slot_before_downtime) = *guard.get_last_slot_before_downtime()
233 {
234 let (shutdown_start, shutdown_end) = shutdown_range(
235 cfg,
236 guard.get_last_start_period(),
237 last_slot_before_downtime,
238 )?;
239 guard
240 .get_mip_store()
241 .is_consistent_with_shutdown_period(
242 shutdown_start,
243 shutdown_end,
244 cfg.thread_count,
245 cfg.t0,
246 cfg.genesis_timestamp,
247 )
248 .map_err(|e| {
249 BootstrapError::GeneralError(format!(
250 "bootstrapped MIP store is not consistent with the shutdown period announced by the server: {}",
251 e
252 ))
253 })?;
254 }
255
256 *next_bootstrap_message = BootstrapClientMessage::AskBootstrapPeers;
260
261 return Ok(());
262 }
263 BootstrapServerMessage::SlotTooOld => {
264 info!("Slot is too old retry bootstrap from scratch");
265 restart_state_streaming(client, next_bootstrap_message, global_bootstrap_state);
266 return Err(BootstrapError::GeneralError(String::from("Slot too old")));
267 }
268 BootstrapServerMessage::BootstrapError { error } => {
270 return Err(BootstrapError::GeneralError(error))
271 }
272 _ => {
273 return Err(BootstrapError::GeneralError(
274 "unexpected message".to_string(),
275 ))
276 }
277 }
278 }
279 } else {
280 Err(BootstrapError::GeneralError(format!(
281 "Try to stream the final state but the message to send to the server was {:#?}",
282 next_bootstrap_message
283 )))
284 }
285}
286
287fn restart_state_streaming(
290 client: &mut BootstrapClientBinder,
291 next_bootstrap_message: &mut BootstrapClientMessage,
292 global_bootstrap_state: &mut GlobalBootstrapState,
293) {
294 *next_bootstrap_message = BootstrapClientMessage::AskBootstrapPart {
295 last_slot: None,
296 last_state_step: StreamingStep::Started,
297 last_versioning_step: StreamingStep::Started,
298 last_consensus_step: StreamingStep::Started,
299 send_last_start_period: true,
300 };
301 let mut write_final_state = global_bootstrap_state.final_state.write();
302 write_final_state.reset();
303 write_final_state.set_last_start_period(0);
306 write_final_state.set_last_slot_before_downtime(None);
307 drop(write_final_state);
308 client.set_last_start_period(None);
309 global_bootstrap_state.graph = None;
313 global_bootstrap_state.peers = None;
314}
315
316pub(crate) fn bootstrap_from_server(
319 cfg: &BootstrapConfig,
320 client: &mut BootstrapClientBinder,
321 next_bootstrap_message: &mut BootstrapClientMessage,
322 global_bootstrap_state: &mut GlobalBootstrapState,
323 our_version: Version,
324) -> Result<(), BootstrapError> {
325 massa_trace!("bootstrap.lib.bootstrap_from_server", {});
326
327 match client.next_timeout(Some(cfg.read_error_timeout.to_duration())) {
330 Err(BootstrapError::TimedOut(_)) => {
331 massa_trace!(
332 "bootstrap.lib.bootstrap_from_server: No error sent at connection",
333 {}
334 );
335 }
336 Err(e) => return Err(e),
337 Ok(BootstrapServerMessage::BootstrapError { error: err }) => {
338 return Err(BootstrapError::ReceivedError(err))
339 }
340 Ok(msg) => return Err(BootstrapError::UnexpectedServerMessage(msg)),
341 };
342
343 let send_time_uncompensated = MassaTime::now();
345 client.handshake(our_version)?;
347
348 let ping = MassaTime::now().saturating_sub(send_time_uncompensated);
350 if ping > cfg.max_ping {
351 return Err(BootstrapError::GeneralError(
352 "bootstrap ping too high".into(),
353 ));
354 }
355
356 let server_time = match client.next_timeout(Some(cfg.read_timeout.into())) {
359 Err(e) => return Err(e),
360 Ok(BootstrapServerMessage::BootstrapTime {
361 server_time,
362 version,
363 }) => {
364 if !our_version.is_compatible(&version) {
365 return Err(BootstrapError::IncompatibleVersionError(format!(
366 "remote is running incompatible version: {} (local node version: {})",
367 version, our_version
368 )));
369 }
370 server_time
371 }
372 Ok(BootstrapServerMessage::BootstrapError { error }) => {
373 return Err(BootstrapError::ReceivedError(error))
374 }
375 Ok(msg) => return Err(BootstrapError::UnexpectedServerMessage(msg)),
376 };
377
378 let recv_time = MassaTime::now();
380
381 let ping = recv_time.saturating_sub(send_time_uncompensated);
383 if ping > cfg.max_ping {
384 return Err(BootstrapError::GeneralError(
385 "bootstrap ping too high".into(),
386 ));
387 }
388
389 let adjusted_server_time = server_time.checked_add(ping.checked_div_u64(2)?)?;
393 let clock_delta = adjusted_server_time.abs_diff(recv_time);
394
395 if clock_delta > cfg.max_clock_delta {
397 warn!("client and server clocks differ too much, please check your clock");
398 let message = format!(
399 "client = {}, server = {}, ping = {}, max_delta = {}",
400 recv_time, server_time, ping, cfg.max_clock_delta
401 );
402 return Err(BootstrapError::ClockError(message));
403 }
404
405 let write_timeout: std::time::Duration = cfg.write_timeout.into();
406 loop {
408 match next_bootstrap_message {
409 BootstrapClientMessage::AskBootstrapPart { .. } => {
410 stream_final_state_and_consensus(
411 cfg,
412 client,
413 next_bootstrap_message,
414 global_bootstrap_state,
415 )?;
416 }
417 BootstrapClientMessage::AskBootstrapPeers => {
418 let peers = match send_client_message(
419 next_bootstrap_message,
420 client,
421 write_timeout,
422 cfg.read_timeout.into(),
423 "ask bootstrap peers timed out",
424 )? {
425 BootstrapServerMessage::BootstrapPeers { peers } => peers,
426 BootstrapServerMessage::BootstrapError { error } => {
427 return Err(BootstrapError::ReceivedError(error))
428 }
429 other => return Err(BootstrapError::UnexpectedServerMessage(other)),
430 };
431 global_bootstrap_state.peers = Some(peers);
432 *next_bootstrap_message = BootstrapClientMessage::BootstrapSuccess;
433 }
434 BootstrapClientMessage::BootstrapSuccess => {
435 client.send_timeout(next_bootstrap_message, Some(write_timeout))?;
436 break;
437 }
438 BootstrapClientMessage::BootstrapError { error: _ } => {
439 panic!("The next message to send shouldn't be BootstrapError");
440 }
441 };
442 }
443 info!("Successful bootstrap");
444 Ok(())
445}
446
447fn check_restart_metadata(
455 cfg: &BootstrapConfig,
456 last_start_period: &Option<u64>,
457 last_slot_before_downtime: &Option<Option<Slot>>,
458) -> Result<(), BootstrapError> {
459 let (last_start_period, last_slot_before_downtime) =
460 match (last_start_period, last_slot_before_downtime) {
461 (None, None) => return Ok(()),
462 (Some(period), Some(slot)) => (*period, *slot),
463 _ => {
464 return Err(BootstrapError::GeneralError(
465 "the server sent only one half of the network restart metadata".into(),
466 ))
467 }
468 };
469
470 get_block_slot_timestamp(
472 cfg.thread_count,
473 cfg.t0,
474 cfg.genesis_timestamp,
475 Slot::new(last_start_period, cfg.thread_count.saturating_sub(1)),
476 )
477 .map_err(|e| {
478 BootstrapError::GeneralError(format!(
479 "the server sent an out of range last_start_period {}: {}",
480 last_start_period, e
481 ))
482 })?;
483
484 let Some(last_slot_before_downtime) = last_slot_before_downtime else {
485 return Ok(());
486 };
487
488 if last_slot_before_downtime >= Slot::new(last_start_period, 0) {
490 return Err(BootstrapError::GeneralError(format!(
491 "the server announced a downtime starting at {} but a restart at period {}",
492 last_slot_before_downtime, last_start_period
493 )));
494 }
495
496 shutdown_range(cfg, last_start_period, last_slot_before_downtime).map(|_| ())
497}
498
499fn shutdown_range(
502 cfg: &BootstrapConfig,
503 last_start_period: u64,
504 last_slot_before_downtime: Slot,
505) -> Result<(Slot, Slot), BootstrapError> {
506 let shutdown_start = last_slot_before_downtime
507 .get_next_slot(cfg.thread_count)
508 .map_err(|e| {
509 BootstrapError::GeneralError(format!(
510 "the server sent an out of range last_slot_before_downtime {}: {}",
511 last_slot_before_downtime, e
512 ))
513 })?;
514 let shutdown_end = Slot::new(last_start_period, cfg.thread_count.saturating_sub(1));
518 Ok((shutdown_start, shutdown_end))
519}
520
521fn send_client_message(
522 message_to_send: &BootstrapClientMessage,
523 client: &mut BootstrapClientBinder,
524 write_timeout: Duration,
525 read_timeout: Duration,
526 error: &str,
527) -> Result<BootstrapServerMessage, BootstrapError> {
528 client.send_timeout(message_to_send, Some(write_timeout))?;
529
530 client
531 .next_timeout(Some(read_timeout))
532 .map_err(|e| match e {
533 BootstrapError::TimedOut(_) => {
534 BootstrapError::TimedOut(std::io::Error::new(std::io::ErrorKind::TimedOut, error))
535 }
536 _ => e,
537 })
538}
539
540pub(crate) fn connect_to_server(
541 connector: &mut impl BSConnector,
542 bootstrap_config: &BootstrapConfig,
543 addr: &SocketAddr,
544 pub_key: &PublicKey,
545 rw_limit: Option<u64>,
546) -> Result<BootstrapClientBinder, BootstrapError> {
547 let socket = connector.connect_timeout(*addr, Some(bootstrap_config.connect_timeout))?;
548 socket.set_nonblocking(false)?;
549 Ok(BootstrapClientBinder::new(
550 socket,
551 *pub_key,
552 bootstrap_config.into(),
553 rw_limit,
554 ))
555}
556
557fn filter_bootstrap_list(
558 bootstrap_list: Vec<(SocketAddr, NodeId)>,
559 ip_type: IpType,
560) -> Vec<(SocketAddr, NodeId)> {
561 let ip_filter: fn(&(SocketAddr, NodeId)) -> bool = match ip_type {
562 IpType::IPv4 => |&(addr, _)| addr.is_ipv4(),
563 IpType::IPv6 => |&(addr, _)| addr.is_ipv6(),
564 IpType::Both => |_| true,
565 };
566
567 let prev_bootstrap_list_len = bootstrap_list.len();
568
569 let filtered_bootstrap_list: Vec<_> = bootstrap_list.into_iter().filter(ip_filter).collect();
570
571 let new_bootstrap_list_len = filtered_bootstrap_list.len();
572
573 debug!(
574 "Keeping {:?} bootstrap ip types. Filtered out {} bootstrap addresses out of a total of {} bootstrap servers.",
575 ip_type,
576 prev_bootstrap_list_len as i32 - new_bootstrap_list_len as i32,
577 prev_bootstrap_list_len
578 );
579
580 filtered_bootstrap_list
581}
582
583#[allow(clippy::too_many_arguments)]
587pub fn get_state(
588 bootstrap_config: &BootstrapConfig,
589 final_state: Arc<RwLock<dyn FinalStateController>>,
590 mut connector: impl BSConnector,
591 version: Version,
592 genesis_timestamp: MassaTime,
593 end_timestamp: Option<MassaTime>,
594 restart_from_snapshot_at_period: Option<u64>,
595 interrupted: Arc<(Mutex<bool>, Condvar)>,
596 massa_metrics: MassaMetrics,
597) -> Result<GlobalBootstrapState, BootstrapError> {
598 massa_trace!("bootstrap.lib.get_state", {});
599
600 if restart_from_snapshot_at_period.is_some() {
602 massa_trace!("bootstrap.lib.get_state.init_from_snapshot", {});
603 return Ok(GlobalBootstrapState::new(final_state));
604 }
605
606 if MassaTime::now() < genesis_timestamp {
608 massa_trace!("bootstrap.lib.get_state.init_from_scratch", {});
609 {
611 let mut final_state_guard = final_state.write();
612
613 if !bootstrap_config.keep_ledger {
614 final_state_guard
616 .get_ledger_mut()
617 .load_initial_ledger()
618 .map_err(|err| {
619 BootstrapError::GeneralError(format!(
620 "could not load initial ledger: {}",
621 err
622 ))
623 })?;
624 }
625
626 let slot = Slot::new(
627 final_state_guard.get_last_start_period(),
628 bootstrap_config.thread_count.saturating_sub(1),
629 );
630
631 let mut batch = DBBatch::new();
633 let mut db_versioning_batch: BTreeMap<Vec<u8>, Option<Vec<u8>>> = DBBatch::new();
634 final_state_guard
635 .get_pos_state_mut()
636 .create_initial_cycle(&mut batch);
637
638 final_state_guard.init_execution_trail_hash_to_batch(&mut batch);
640
641 final_state_guard
643 .load_initial_deferred_credits(&mut batch)
644 .map_err(|err| {
645 BootstrapError::GeneralError(format!(
646 "could not load initial deferred credits: {}",
647 err
648 ))
649 })?;
650
651 final_state_guard
653 .get_mip_store()
654 .update_batches(&mut batch, &mut db_versioning_batch, None)
655 .map_err(|e| BootstrapError::GeneralError(e.to_string()))?;
656
657 final_state_guard.get_database().write().write_batch(
658 batch,
659 db_versioning_batch,
660 Some(slot),
661 );
662 }
663 return Ok(GlobalBootstrapState::new(final_state));
664 }
665
666 let filtered_bootstrap_list = get_bootstrap_list_iter(bootstrap_config)?;
669
670 let mut next_bootstrap_message: BootstrapClientMessage =
671 BootstrapClientMessage::AskBootstrapPart {
672 last_slot: None,
673 last_state_step: StreamingStep::Started,
674 last_versioning_step: StreamingStep::Started,
675 last_consensus_step: StreamingStep::Started,
676 send_last_start_period: true,
677 };
678 let mut global_bootstrap_state = GlobalBootstrapState::new(final_state);
679
680 let limit = bootstrap_config.rate_limit;
681 loop {
682 if *interrupted
684 .0
685 .lock()
686 .expect("double-lock on interrupt-mutex")
687 {
688 return Err(BootstrapError::Interrupted(
689 "Sig INT received while getting state".to_string(),
690 ));
691 }
692 for (addr, node_id) in filtered_bootstrap_list.iter() {
693 if let Some(end) = end_timestamp {
694 if MassaTime::now() > end {
695 panic!("This episode has come to an end, please get the latest testnet node version to continue");
696 }
697 }
698 info!("Start bootstrapping from {}", addr);
699 let conn = connect_to_server(
700 &mut connector,
701 bootstrap_config,
702 addr,
703 &node_id.get_public_key(),
704 Some(limit),
705 );
706 match conn {
707 Ok(mut client) => {
708 massa_metrics.inc_bootstrap_counter();
709 let bs = bootstrap_from_server(
710 bootstrap_config,
711 &mut client,
712 &mut next_bootstrap_message,
713 &mut global_bootstrap_state,
714 version,
715 );
716 match bs {
718 Err(BootstrapError::ReceivedError(error)) => {
719 warn!("Error received from bootstrap server: {}", error)
720 }
721 Err(e) => {
722 warn!("Error while bootstrapping: {}", &e);
723 let _ = client.send_timeout(
725 &BootstrapClientMessage::BootstrapError {
726 error: e.to_string(),
727 },
728 Some(bootstrap_config.write_error_timeout.into()),
729 );
730 }
731 Ok(()) => return Ok(global_bootstrap_state),
732 }
733 }
734 Err(e) => {
735 warn!("Error while connecting to bootstrap server: {}", e);
736 }
737 };
738
739 info!("Bootstrap from server {} failed. Your node will try to bootstrap from another server in {}.", addr, format_duration(bootstrap_config.retry_delay.to_duration()).to_string());
740
741 let int_sig = interrupted
759 .0
760 .lock()
761 .expect("double-lock() on interrupted signal mutex");
762 let wake = interrupted
763 .1
764 .wait_timeout(int_sig, bootstrap_config.retry_delay.to_duration())
765 .expect("interrupt signal mutex poisoned");
766 if *wake.0 {
767 return Err(BootstrapError::Interrupted(
768 "Sig INT during bootstrap retry-wait".to_string(),
769 ));
770 }
771 }
772 }
773}
774
775fn get_bootstrap_list_iter(
776 bootstrap_config: &BootstrapConfig,
777) -> Result<Vec<(SocketAddr, NodeId)>, BootstrapError> {
778 let mut filtered_bootstrap_list = filter_bootstrap_list(
779 bootstrap_config.bootstrap_list.clone(),
780 bootstrap_config.bootstrap_protocol,
781 );
782
783 massa_trace!("bootstrap.lib.get_state.init_from_others", {});
785 if filtered_bootstrap_list.is_empty() {
786 return Err(BootstrapError::GeneralError(
787 "no bootstrap nodes found in list".into(),
788 ));
789 }
790
791 filtered_bootstrap_list.shuffle(&mut StdRng::from_entropy());
793
794 let mut unique_node_ids: HashSet<NodeId> = HashSet::new();
796 filtered_bootstrap_list.retain(|e| unique_node_ids.insert(e.1));
797 Ok(filtered_bootstrap_list)
798}
799
800fn warn_user_about_versioning_updates(updated: Vec<MipInfo>, added: BTreeMap<MipInfo, MipState>) {
801 if !added.is_empty() {
802 for (mip_info, mip_state) in added.iter() {
803 let now = MassaTime::now();
804 match mip_state.state_at(
805 now,
806 mip_info.start,
807 mip_info.timeout,
808 mip_info.activation_delay,
809 ) {
810 Ok(st_id) => {
811 if st_id == ComponentStateTypeId::LockedIn {
812 warn!(
814 "A new MIP has been locked in: {}, version: {}",
815 mip_info.name, mip_info.version
816 );
817 let activation_at = mip_state.activation_at(mip_info).unwrap();
819
820 warn!(
821 "Please update your Massa node before: {}",
822 activation_at.format_instant()
823 );
824 } else if st_id == ComponentStateTypeId::Active {
825 warn!(
827 "A new MIP has become active {:?}, version: {:?}",
828 mip_info.name, mip_info.version
829 );
830 panic!(
831 "Please update your Massa node to support MIP version {} ({})",
832 mip_info.version, mip_info.name
833 );
834 } else if st_id == ComponentStateTypeId::Defined {
835 warn!(
838 "A new MIP has been defined: {}, version: {}",
839 mip_info.name, mip_info.version
840 );
841 debug!("MIP state: {:?}", mip_state);
842
843 warn!("Please update your node between: {} and {} if you want to support this update", mip_info.start.format_instant(), mip_info.timeout.format_instant());
844 } else {
845 warn!(
848 "A new MIP has been received: {}, version: {}",
849 mip_info.name, mip_info.version
850 );
851 debug!("MIP state: {:?}", mip_state);
852 warn!("Please update your Massa node to support it");
853 }
854 }
855 Err(StateAtError::Unpredictable) => {
856 warn!(
857 "A new MIP has started: {}, version: {}",
858 mip_info.name, mip_info.version
859 );
860 debug!("MIP state: {:?}", mip_state);
861
862 warn!("Please update your node between: {} and {} if you want to support this update", mip_info.start.format_instant(), mip_info.timeout.format_instant());
863 }
864 Err(e) => {
865 panic!(
867 "Unable to get state at {} of mip info: {:?}, error: {}",
868 now, mip_info, e
869 )
870 }
871 }
872 }
873 }
874
875 debug!("MIP store got {} MIP updated from bootstrap", updated.len());
876}
877
878fn align_consensus_resume_state_before_ask(
879 next_bootstrap_message: &mut BootstrapClientMessage,
880 global_bootstrap_state: &mut GlobalBootstrapState,
881) {
882 let mut must_restart_consensus = false;
883 if let BootstrapClientMessage::AskBootstrapPart {
884 last_consensus_step: StreamingStep::Ongoing(ids),
885 ..
886 } = next_bootstrap_message
887 {
888 if let Some(graph) = global_bootstrap_state.graph.as_mut() {
889 graph
890 .final_blocks
891 .retain(|block_export| ids.contains(&block_export.block.id));
892 ids.clear();
893 ids.extend(
894 graph
895 .final_blocks
896 .iter()
897 .map(|block_export| block_export.block.id),
898 );
899 if graph.final_blocks.is_empty() {
900 global_bootstrap_state.graph = None;
901 }
902 } else {
903 ids.clear();
904 }
905 must_restart_consensus = ids.is_empty();
906 }
907 if must_restart_consensus {
908 if let BootstrapClientMessage::AskBootstrapPart {
909 last_consensus_step,
910 ..
911 } = next_bootstrap_message
912 {
913 *last_consensus_step = StreamingStep::Started;
914 }
915 }
916}
917
918fn capped_consensus_resume_step(
919 graph: Option<&BootstrapableGraph>,
920 max_ids: usize,
921) -> StreamingStep<PreHashSet<BlockId>> {
922 let Some(graph) = graph else {
923 return StreamingStep::Started;
924 };
925 let ids = capped_consensus_cursor_ids(graph, max_ids);
926 if ids.is_empty() {
927 StreamingStep::Started
928 } else {
929 StreamingStep::Ongoing(ids)
930 }
931}
932
933fn capped_consensus_cursor_ids(graph: &BootstrapableGraph, max_ids: usize) -> PreHashSet<BlockId> {
935 if max_ids == 0 || graph.final_blocks.is_empty() {
936 return PreHashSet::default();
937 }
938 if graph.final_blocks.len() <= max_ids {
939 return graph
940 .final_blocks
941 .iter()
942 .map(|block_export| block_export.block.id)
943 .collect();
944 }
945 let mut ordered: Vec<_> = graph
946 .final_blocks
947 .iter()
948 .map(|block_export| {
949 (
950 block_export.block.content.header.content.slot,
951 block_export.block.id,
952 )
953 })
954 .collect();
955 ordered.sort_unstable_by(|a, b| a.0.cmp(&b.0).then_with(|| a.1.cmp(&b.1)));
956 ordered[ordered.len() - max_ids..]
957 .iter()
958 .map(|(_, id)| *id)
959 .collect()
960}
961
962#[cfg(test)]
963mod restart_metadata_tests {
964 use super::*;
965 use crate::tests::tools::{gen_export_active_blocks, get_bootstrap_config};
966 use massa_consensus_exports::bootstrapable_graph::BootstrapableGraph;
967 use massa_signature::KeyPair;
968
969 fn config() -> BootstrapConfig {
970 get_bootstrap_config(NodeId::new(KeyPair::generate(0).unwrap().get_public_key()))
971 }
972
973 #[test]
974 fn accepts_plausible_restart_metadata() {
975 let cfg = config();
976 let tc = cfg.thread_count;
977 for (last_start_period, last_slot_before_downtime) in [
978 (None, None),
980 (Some(0), Some(None)),
982 (Some(1), Some(Some(Slot::new(0, tc - 1)))),
984 (Some(100), Some(Some(Slot::new(42, 3)))),
986 ] {
987 check_restart_metadata(&cfg, &last_start_period, &last_slot_before_downtime).unwrap();
988 }
989 }
990
991 #[test]
992 fn rejects_impossible_restart_metadata() {
993 let cfg = config();
994 for (last_start_period, last_slot_before_downtime) in [
995 (Some(2), None),
997 (None, Some(Some(Slot::new(1, 0)))),
998 (Some(0), Some(Some(Slot::new(0, 0)))),
1000 (Some(1), Some(Some(Slot::new(5, 0)))),
1002 (Some(u64::MAX), Some(None)),
1004 ] {
1005 check_restart_metadata(&cfg, &last_start_period, &last_slot_before_downtime)
1006 .expect_err(&format!(
1007 "accepted {:?} / {:?}",
1008 last_start_period, last_slot_before_downtime
1009 ));
1010 }
1011 }
1012
1013 #[test]
1014 fn capped_consensus_cursor_returns_started_when_empty() {
1015 assert_eq!(
1016 capped_consensus_resume_step(None, 10),
1017 StreamingStep::Started
1018 );
1019 let graph = BootstrapableGraph {
1020 final_blocks: vec![],
1021 };
1022 assert_eq!(
1023 capped_consensus_resume_step(Some(&graph), 10),
1024 StreamingStep::Started
1025 );
1026 }
1027
1028 #[test]
1029 fn capped_consensus_cursor_is_bounded() {
1030 let mut rng = rand::thread_rng();
1031 let graph = BootstrapableGraph {
1032 final_blocks: (0..10)
1033 .map(|_| gen_export_active_blocks(&mut rng))
1034 .collect(),
1035 };
1036 match capped_consensus_resume_step(Some(&graph), 3) {
1037 StreamingStep::Ongoing(ids) => assert_eq!(ids.len(), 3),
1038 other => panic!("expected Ongoing, got {other:?}"),
1039 }
1040 }
1041
1042 #[test]
1043 fn align_reconnect_cursor_prunes_graph_to_claimed_ids() {
1044 let mut rng = rand::thread_rng();
1045 let blocks: Vec<_> = (0..8).map(|_| gen_export_active_blocks(&mut rng)).collect();
1046 let claimed: PreHashSet<_> = blocks
1047 .iter()
1048 .take(3)
1049 .map(|block_export| block_export.block.id)
1050 .collect();
1051
1052 let mut state = GlobalBootstrapState {
1053 final_state: Arc::new(RwLock::new(
1054 massa_final_state::MockFinalStateController::new(),
1055 )),
1056 graph: Some(BootstrapableGraph {
1057 final_blocks: blocks,
1058 }),
1059 peers: None,
1060 };
1061 let mut msg = BootstrapClientMessage::AskBootstrapPart {
1062 last_slot: Some(Slot::new(1, 0)),
1063 last_state_step: StreamingStep::Ongoing(vec![0]),
1064 last_versioning_step: StreamingStep::Ongoing(vec![0]),
1065 last_consensus_step: StreamingStep::Ongoing(claimed.clone()),
1066 send_last_start_period: false,
1067 };
1068
1069 align_consensus_resume_state_before_ask(&mut msg, &mut state);
1070
1071 let graph = state.graph.expect("graph should remain");
1072 assert_eq!(graph.final_blocks.len(), claimed.len());
1073 for b in &graph.final_blocks {
1074 assert!(claimed.contains(&b.block.id));
1075 }
1076 match msg {
1077 BootstrapClientMessage::AskBootstrapPart {
1078 last_consensus_step: StreamingStep::Ongoing(ids),
1079 ..
1080 } => assert_eq!(ids, claimed),
1081 other => panic!("unexpected message {other:?}"),
1082 }
1083 }
1084}