massa_bootstrap/
client.rs

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/// Specifies a common interface that can be used by standard, or mockers
38#[cfg_attr(test, mockall::automock)]
39pub trait BSConnector {
40    /// The client attempts to connect to the given address.
41    /// If a duration is provided, the attempt will be timed out after the given duration.
42    fn connect_timeout(
43        &self,
44        addr: SocketAddr,
45        duration: Option<MassaTime>,
46    ) -> io::Result<TcpStream>;
47}
48
49/// Initiates a connection with given timeout in milliseconds
50#[derive(Debug)]
51pub struct DefaultConnector;
52
53impl BSConnector for DefaultConnector {
54    /// Tries to connect to address
55    ///
56    /// # Argument
57    /// * `addr`: `SocketAddr` we are trying to connect to.
58    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}
69/// This function will send the starting point to receive a stream of the ledger and will receive and process each part until receive a `BootstrapServerMessage::FinalStateFinished` message from the server.
70/// `next_bootstrap_message` passed as parameter must be `BootstrapClientMessage::AskFinalStatePart` enum variant.
71/// `next_bootstrap_message` will be updated after receiving each part so that in case of connection lost we can restart from the last message we processed.
72///
73/// The consensus reconnect cursor is capped to `max_consensus_block_ids` (newest finals by slot).
74/// The full graph is kept while streaming; on the next attempt it is pruned to the ids we claim
75/// so dropped finals can be re-downloaded.
76pub(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        // Continuation / reconnect: metadata was received on the first part (empty
90        // consensus); seed the binder for block header validation while parsing.
91        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                    // Set final state
123                    let mut write_final_state = global_bootstrap_state.final_state.write();
124
125                    // Reject inconsistent network restart metadata before storing it: both
126                    // fields feed slot and timestamp arithmetic at startup, so impossible values
127                    // would only surface much later, as a panic following an otherwise
128                    // successful bootstrap.
129                    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                    // Set consensus blocks
151                    if let Some(graph) = global_bootstrap_state.graph.as_mut() {
152                        // Extend the final blocks with the received part
153                        graph.final_blocks.extend(consensus_part.final_blocks);
154                        // Remove every outdated block
155                        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                    // Cap the reconnect ask only; keep the full in-session graph.
163                    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                    // Set new message in case of disconnection
169                    *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                    // Logs for an easier diagnostic if needed
178                    debug!(
179                        "client final state bootstrap cursors: {:?}",
180                        next_bootstrap_message
181                    );
182                }
183                BootstrapServerMessage::BootstrapFinished => {
184                    info!("State bootstrap complete");
185
186                    // Refuse a server that predates a MIP whose vote has already started: it
187                    // cannot provide that MIP's state, and replaying the MIP locally from our
188                    // first processed slot would record a different history (and, past lock-in,
189                    // a different activation) than the rest of the network -- a history that
190                    // becomes part of the hashed final state once the MIP is Active. Restart
191                    // from scratch with the next server rather than resume on top of its state.
192                    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                    // Update MIP store by reading from the disk
220                    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                    // The downtime range announced by the server must also be consistent with the
230                    // MIP store we just bootstrapped: startup requires it, so checking it here
231                    // lets us retry with another server instead of failing after the fact.
232                    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                    // Only advance to the next phase once the streamed state has been fully
257                    // post-processed: on failure the retry must replay the state streaming
258                    // instead of resuming from the peers phase with a half-updated store.
259                    *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                // At this point, we have successfully received the next message from the server, and it's an error-message String
269                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
287/// Drop everything streamed from the current server so that the next attempt streams the
288/// state again from scratch, instead of resuming on top of what this server sent.
289fn 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    // `reset()` does not clear restart metadata; drop values from the
304    // aborted server so the next one defines them from the wire again.
305    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    // The cursor above restarts the consensus stream from `Started`, for which the
310    // server reports no outdated ids: blocks kept from the aborted attempt would
311    // never be pruned and would be merged into the next attempt's graph.
312    global_bootstrap_state.graph = None;
313    global_bootstrap_state.peers = None;
314}
315
316/// Gets the state from a bootstrap server (internal private function)
317/// needs to be CANCELLABLE
318pub(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    // read error (if sent by the server)
328    // client.next() is not cancel-safe but we drop the whole client object if cancelled => it's OK
329    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    // handshake
344    let send_time_uncompensated = MassaTime::now();
345    // client.handshake() is not cancel-safe but we drop the whole client object if cancelled => it's OK
346    client.handshake(our_version)?;
347
348    // compute ping
349    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    // First, clock and version.
357    // client.next() is not cancel-safe but we drop the whole client object if cancelled => it's OK
358    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    // get the time of reception
379    let recv_time = MassaTime::now();
380
381    // compute ping
382    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    // compute client / server clock delta
390    // div 2 is an approximation of the time it took the message to do server -> client
391    // the complete ping value being client -> server -> client
392    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 is too high warn the user and restart bootstrap
396    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 to ask data to the server depending on the last message we sent
407    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
447/// Checks the network restart metadata announced by a bootstrap server.
448///
449/// The server sends `last_start_period` and `last_slot_before_downtime` together, and only once
450/// per stream (on the first part, before any consensus blocks). Startup derives the network
451/// downtime range and its timestamps from them (see `massa-node`), assuming they describe a
452/// coherent interval, so anything else has to be rejected as a bootstrap failure rather than
453/// stored.
454fn 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    // the timestamp of the last start slot is computed at startup: it must not overflow
471    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    // the last slot executed before the downtime has to precede the restart
489    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
499/// Computes the slot range of the last network shutdown, as startup does: from the slot right
500/// after the last one executed before the downtime, to the slot right before the restart.
501fn 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    // Include the entire last_start_period as downtime: interpolation attaches at
515    // Slot(last_start_period, thread_count - 1) and block production resumes only
516    // at Slot(last_start_period + 1, 0).
517    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/// Uses the cond-var pattern to handle sig-int cancellation.
584/// Make sure that the passed in `interrupted` shares its Arc
585/// with a sig-int handler setup.
586#[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 we restart from a snapshot, do not bootstrap
601    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 we are before genesis, do not bootstrap
607    if MassaTime::now() < genesis_timestamp {
608        massa_trace!("bootstrap.lib.get_state.init_from_scratch", {});
609        // init final state
610        {
611            let mut final_state_guard = final_state.write();
612
613            if !bootstrap_config.keep_ledger {
614                // load ledger from initial ledger file
615                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            // create the initial cycle of PoS cycle_history
632            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            // set initial execution trail hash
639            final_state_guard.init_execution_trail_hash_to_batch(&mut batch);
640
641            // load initial deferred credits
642            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            // Need to write MIP store to Db if we want to bootstrap it to others
652            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    // If the two conditions above are not verified, we need to bootstrap
667    // we filter the bootstrap list to keep only the ip addresses we are compatible with
668    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        // check for interruption
683        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                    // cancellable
717                    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                            // We allow unused result because we don't care if an error is thrown when sending the error message to the server we will close the socket anyway.
724                            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            // Before, we would use a simple sleep(...), and that was fine
742            // in a cancellable async context: the runtime could
743            // catch the interrupt signal, and just cancel this thread:
744            //
745            // let state = tokio::select!{
746            //    /* detect interrupt */ => /* return, cancelling the async get_state */
747            //    get_state(...) => well, we got the state, and it didn't have to worry about interrupts
748            // };
749            //
750            // Without an external system to preempt this context, we use a condvar to manage the sleep.
751            //
752            // Condvar::wait is basically std::thread::sleep(/* until some magic happens */)
753            // Condvar::wait_timeout(..., duration) is much the same, but for a max-len of `duration`
754            //
755            // The _magic_ happens when, somewhere else, a clone of the Arc<(Mutex<bool>, Condvar)>\
756            // calls Condvar::notify_[one | all], which prompts this thread to wake up. Assuming that
757            // the mutex-wrapped variable has been set appropriately before the notify, this thread
758            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    // we are after genesis => bootstrap
784    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    // we shuffle the list
792    filtered_bootstrap_list.shuffle(&mut StdRng::from_entropy());
793
794    // we remove the duplicated node ids (if a bootstrap server appears both with its IPv4 and IPv6 address)
795    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                        // A new MipInfo @ state locked_in - we need to urge the user to update
813                        warn!(
814                            "A new MIP has been locked in: {}, version: {}",
815                            mip_info.name, mip_info.version
816                        );
817                        // Safe to unwrap here (only panic if not LockedIn)
818                        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                        // A new MipInfo @ state active - we are not compatible anymore
826                        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                        // a new MipInfo @ state defined or started (or failed / error)
836                        // warn the user to update its node
837                        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                        // a new MipInfo @ state defined or started (or failed / error)
846                        // warn the user to update its node
847                        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                    // Should never happen
866                    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
933/// Ids for the reconnect ask: all finals if under the cap, otherwise the newest `max_ids` by slot.
934fn 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            // a server that has nothing to say about a restart
979            (None, None),
980            // a network that never restarted
981            (Some(0), Some(None)),
982            // a restart leaving no idle slot behind
983            (Some(1), Some(Some(Slot::new(0, tc - 1)))),
984            // a restart after an actual downtime
985            (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            // only one half of the metadata
996            (Some(2), None),
997            (None, Some(Some(Slot::new(1, 0)))),
998            // period 0 has no previous slot, the downtime cannot end before the restart
999            (Some(0), Some(Some(Slot::new(0, 0)))),
1000            // the downtime starts after the restart
1001            (Some(1), Some(Some(Slot::new(5, 0)))),
1002            // the restart timestamp does not fit in a `MassaTime`
1003            (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}