massa_bootstrap/
server.rs

1//! start the bootstrapping system using [`start_bootstrap_server`]
2//! Once your node will be ready, you may want other to bootstrap from you.
3//!
4//! # Listener
5//!
6//! Runs in the server-dedication tokio async runtime
7//! Accepts bootstrap connections in an async-loop
8//! Upon connection, pushes the accepted connection onto a channel for the worker loop to consume
9//!
10//! # Updater
11//!
12//! Runs on a dedicated thread. Signal sent my manager stop method terminates the thread.
13//! Shares an `Arc<RwLock>>` guarded list of white and blacklists with the main worker.
14//! Periodically does a read-only check to see if list needs updating.
15//! Creates an updated list then swaps it out with write-locked list
16//! Assuming no errors in code, this is the only write occurrence, and is only a pointer-swap
17//! under the hood, making write contention virtually non-existent.
18//!
19//! # Worker loop
20//!
21//! 1. Checks if the stopper has been invoked.
22//! 2. Checks if the client is permitted under the white/black list rules
23//! 3. Checks if there are not too many active sessions already
24//! 4. Checks if the client has attempted too recently
25//! 5. All checks have passed: spawn a thread on which to run the bootstrap session
26//!    This thread creates a new tokio runtime, and runs it with `block_on`
27
28use crossbeam::channel::tick;
29use humantime::format_duration;
30use massa_consensus_exports::{bootstrapable_graph::BootstrapableGraph, ConsensusController};
31use massa_db_exports::{MassaDBError, CHANGE_ID_DESER_ERROR};
32use massa_final_state::FinalStateController;
33use massa_logging::massa_trace;
34use massa_metrics::MassaMetrics;
35use massa_models::{
36    block_id::BlockId, prehash::PreHashSet, slot::Slot, streaming_step::StreamingStep,
37    version::Version,
38};
39
40use massa_protocol_exports::ProtocolController;
41use massa_signature::KeyPair;
42use massa_time::MassaTime;
43
44use parking_lot::RwLock;
45use std::{
46    collections::HashMap,
47    net::{IpAddr, SocketAddr},
48    sync::Arc,
49    thread,
50    time::{Duration, Instant},
51};
52use tracing::{debug, error, info, warn};
53
54#[cfg(not(test))]
55use crate::listener::BootstrapTcpListener;
56#[cfg(test)]
57use crate::listener::MockBootstrapTcpListener as BootstrapTcpListener;
58use crate::{
59    bindings::BootstrapServerBinder,
60    error::BootstrapError,
61    listener::{BootstrapListenerStopHandle, PollEvent},
62    messages::{BootstrapClientMessage, BootstrapServerMessage},
63    white_black_list::SharedWhiteBlackList,
64    BootstrapConfig,
65};
66
67// Abstraction layer over data produced by the listener, and transported
68// over to the worker via a channel
69
70/// handle on the bootstrap server
71pub struct BootstrapManager {
72    update_handle: thread::JoinHandle<Result<(), BootstrapError>>,
73    // need to preserve the listener handle up to here to prevent it being destroyed
74    #[allow(clippy::type_complexity)]
75    main_handle: thread::JoinHandle<Result<(), BootstrapError>>,
76    listener_stopper: BootstrapListenerStopHandle,
77    update_stopper_tx: crossbeam::channel::Sender<()>,
78    /// shared white/black list
79    pub white_black_list: SharedWhiteBlackList<'static>,
80}
81
82impl BootstrapManager {
83    /// create a new bootstrap manager, but no means of stopping the listener
84    /// use [`set_listen_stop_handle`] to set the handle
85    pub(crate) fn new(
86        update_handle: thread::JoinHandle<Result<(), BootstrapError>>,
87        main_handle: thread::JoinHandle<Result<(), BootstrapError>>,
88        update_stopper_tx: crossbeam::channel::Sender<()>,
89        listener_stopper: BootstrapListenerStopHandle,
90        white_black_list: SharedWhiteBlackList<'static>,
91    ) -> Self {
92        Self {
93            update_handle,
94            main_handle,
95            update_stopper_tx,
96            listener_stopper,
97            white_black_list,
98        }
99    }
100
101    /// stop the bootstrap server
102    pub fn stop(self) -> Result<(), BootstrapError> {
103        massa_trace!("bootstrap.lib.stop", {});
104        // TODO: Refactor the waker so that its existence is tied to the life of the event-loop
105        if self.listener_stopper.stop().is_err() {
106            warn!("bootstrap server already dropped");
107        }
108        if self.update_stopper_tx.send(()).is_err() {
109            warn!("bootstrap ip-list-updater already dropped");
110        }
111        // TODO?: handle join errors.
112
113        // when the runtime is dropped at the end of this stop, the listener is auto-aborted
114
115        self.update_handle
116            .join()
117            .expect("in BootstrapManager::stop() joining on updater thread")?;
118
119        let res = self
120            .main_handle
121            .join()
122            .expect("in BootstrapManager::stop() joining on bootstrap main-loop thread");
123        info!("bootstrap server stopped");
124        res
125    }
126}
127
128/// See module level documentation for details
129#[allow(clippy::too_many_arguments)]
130pub fn start_bootstrap_server(
131    ev_poller: BootstrapTcpListener,
132    listener_stopper: BootstrapListenerStopHandle,
133    consensus_controller: Box<dyn ConsensusController>,
134    protocol_controller: Box<dyn ProtocolController>,
135    final_state: Arc<RwLock<dyn FinalStateController>>,
136    config: BootstrapConfig,
137    keypair: KeyPair,
138    version: Version,
139    massa_metrics: MassaMetrics,
140) -> Result<BootstrapManager, BootstrapError> {
141    massa_trace!("bootstrap.lib.start_bootstrap_server", {});
142
143    // TODO(low prio): See if a zero capacity channel model can work
144    let (update_stopper_tx, update_stopper_rx) = crossbeam::channel::bounded::<()>(1);
145
146    let Ok(max_bootstraps) = config.max_simultaneous_bootstraps.try_into() else {
147        return Err(BootstrapError::GeneralError(
148            "Fail to convert u32 to usize".to_string(),
149        ));
150    };
151
152    let white_black_list = SharedWhiteBlackList::new(
153        config.bootstrap_whitelist_path.clone(),
154        config.bootstrap_blacklist_path.clone(),
155    )?;
156
157    let updater_lists = white_black_list.clone();
158    let update_handle = thread::Builder::new()
159        .name("wb_list_updater".to_string())
160        .spawn(move || {
161            let res = BootstrapServer::run_updater(
162                updater_lists,
163                config.cache_duration.into(),
164                update_stopper_rx,
165            );
166            match res {
167                Ok(_) => info!("ip white/blacklist updater exited cleanly"),
168                Err(ref er) => error!("updater exited with error: {}", er),
169            };
170            res
171        })
172        .expect("in `start_bootstrap_server`, OS failed to spawn list-updater thread");
173
174    let w_b_list = white_black_list.clone();
175    let main_handle = thread::Builder::new()
176        .name("bs-main-loop".to_string())
177        .spawn(move || {
178            BootstrapServer {
179                consensus_controller,
180                protocol_controller,
181                final_state,
182                ev_poller,
183                white_black_list: w_b_list,
184                keypair,
185                version,
186                ip_hist_map: HashMap::with_capacity(config.ip_list_max_size),
187                bootstrap_config: config,
188                massa_metrics,
189            }
190            .event_loop(max_bootstraps)
191        })
192        .expect("in `start_bootstrap_server`, OS failed to spawn main-loop thread");
193    // Give the runtime to the bootstrap manager, otherwise it will be dropped, forcibly aborting the spawned tasks.
194    // TODO: make the tasks sync, so the runtime is redundant
195    Ok(BootstrapManager::new(
196        update_handle,
197        main_handle,
198        update_stopper_tx,
199        listener_stopper,
200        white_black_list,
201    ))
202}
203
204struct BootstrapServer<'a> {
205    consensus_controller: Box<dyn ConsensusController>,
206    protocol_controller: Box<dyn ProtocolController>,
207    final_state: Arc<RwLock<dyn FinalStateController>>,
208    ev_poller: BootstrapTcpListener,
209    white_black_list: SharedWhiteBlackList<'a>,
210    keypair: KeyPair,
211    bootstrap_config: BootstrapConfig,
212    version: Version,
213    ip_hist_map: HashMap<IpAddr, Instant>,
214    massa_metrics: MassaMetrics,
215}
216
217impl BootstrapServer<'_> {
218    fn run_updater(
219        mut list: SharedWhiteBlackList<'_>,
220        interval: Duration,
221        stopper: crossbeam::channel::Receiver<()>,
222    ) -> Result<(), BootstrapError> {
223        let ticker = tick(interval);
224
225        loop {
226            crossbeam::select! {
227                recv(stopper) -> res => {
228                    match res {
229                        Ok(()) => return Ok(()),
230                        Err(e) => return Err(BootstrapError::GeneralError(format!("update stopper error : {}", e))),
231                    }
232                },
233                recv(ticker) -> _ => {
234                    // A failed refresh (e.g. malformed JSON in an ACL file)
235                    // must not kill the updater: log it and keep the previously
236                    // loaded lists, retrying on the next tick. Propagating the
237                    // error here would stop all future reloads while the event
238                    // loop kept enforcing a stale ACL snapshot.
239                    if let Err(e) = list.update() {
240                        warn!("failed to refresh bootstrap white/black list, keeping previous lists: {}", e);
241                    }
242                },
243            }
244        }
245    }
246
247    fn event_loop(mut self, max_bootstraps: usize) -> Result<(), BootstrapError> {
248        // Use the strong-count of this variable to track the session count
249        let bootstrap_sessions_counter: Arc<()> = Arc::new(());
250        let per_ip_min_interval = self.bootstrap_config.per_ip_min_interval.to_duration();
251        // TODO: Work out how to integration-test this
252        let limit = self.bootstrap_config.rate_limit;
253        loop {
254            // block until we have a connection to work with, or break out of main-loop
255            let connections = match self.ev_poller.poll() {
256                Ok(PollEvent::Stop) => return Ok(()),
257                Ok(PollEvent::NewConnections(connections)) => connections,
258                Err(e) => {
259                    error!("bootstrap listener error: {}", e);
260                    // Intuitively, there would be no connection at this point, However an `nc` that
261                    // leads to this scope doesn't exit client-side. This depends on a timeout error
262                    // client-side
263                    continue;
264                }
265            };
266
267            for (dplx, remote_addr) in connections {
268                // claim a slot in the max_bootstrap_sessions
269                let server_binding = BootstrapServerBinder::new(
270                    dplx,
271                    self.keypair.clone(),
272                    (&self.bootstrap_config).into(),
273                    Some(limit),
274                );
275
276                // the `- 1` is to account for the top-level Arc that is created at the top
277                // of this method. subsequent counts correspond to each `clone` that is passed
278                // into a thread
279                // TODO: If we don't find a way to handle the counting automagically, make
280                //       a dedicated wrapper-type with doc-comments, manual drop impl that
281                //       integrates logging, etc...
282                if Arc::strong_count(&bootstrap_sessions_counter) - 1 < max_bootstraps {
283                    let bootstrap_count_token = bootstrap_sessions_counter.clone();
284                    // check whether incoming peer IP is allowed.
285                    if let Err(error_msg) = self.white_black_list.is_ip_allowed(&remote_addr) {
286                        server_binding.close_and_send_error(
287                            error_msg.to_string(),
288                            remote_addr,
289                            move || {},
290                        );
291                        self.massa_metrics.inc_bootstrap_peers_failed();
292                        continue;
293                    };
294                    massa_trace!("bootstrap.lib.run.select.accept", {
295                        "remote_addr": remote_addr
296                    });
297                    let now = Instant::now();
298
299                    // clear IP history if necessary
300                    if self.ip_hist_map.len() > self.bootstrap_config.ip_list_max_size {
301                        self.ip_hist_map
302                            .retain(|_k, v| now.duration_since(*v) <= per_ip_min_interval);
303                        if self.ip_hist_map.len() > self.bootstrap_config.ip_list_max_size {
304                            // too many IPs are spamming us: clear cache
305                            warn!("high bootstrap load: at least {} different IPs attempted bootstrap in the last {}", self.ip_hist_map.len(),format_duration(self.bootstrap_config.per_ip_min_interval.to_duration()).to_string());
306                            self.ip_hist_map.clear();
307                        }
308                    }
309
310                    // check IP's bootstrap attempt history
311                    if let Err(msg) = BootstrapServer::greedy_client_check(
312                        &mut self.ip_hist_map,
313                        remote_addr,
314                        now,
315                        per_ip_min_interval,
316                    ) {
317                        // Client has been too greedy: send out the bad-news :(
318                        let msg = format!(
319                            "Your last bootstrap on this server was {} ago and you have to wait {} before retrying.",
320                            format_duration(msg),
321                            format_duration(per_ip_min_interval.saturating_sub(msg))
322                        );
323                        let tracer = move || {
324                            massa_trace!("bootstrap.lib.run.select.accept.refuse_limit", {
325                                "remote_addr": remote_addr
326                            })
327                        };
328                        server_binding.close_and_send_error(msg, remote_addr, tracer);
329                        self.massa_metrics.inc_bootstrap_peers_failed();
330                        continue;
331                    };
332
333                    // Clients Option<last-attempt> is good, and has been updated
334                    massa_trace!("bootstrap.lib.run.select.accept.cache_available", {});
335
336                    // launch bootstrap
337                    let version = self.version;
338                    let data_execution = self.final_state.clone();
339                    let consensus_command_sender = self.consensus_controller.clone();
340                    let protocol_controller = self.protocol_controller.clone();
341                    let config = self.bootstrap_config.clone();
342
343                    let massa_metrics = self.massa_metrics.clone();
344
345                    let _ = thread::Builder::new()
346                        .name(format!("bootstrap thread, peer: {}", remote_addr))
347                        .spawn(move || {
348                            run_bootstrap_session(
349                                server_binding,
350                                bootstrap_count_token,
351                                config,
352                                remote_addr,
353                                data_execution,
354                                version,
355                                consensus_command_sender,
356                                protocol_controller,
357                                massa_metrics,
358                            )
359                        });
360
361                    massa_trace!("bootstrap.session.started", {
362                        "active_count": Arc::strong_count(&bootstrap_sessions_counter) - 1
363                    });
364                } else {
365                    server_binding.close_and_send_error(
366                        "Bootstrap failed because the bootstrap server currently has no slots available.".to_string(),
367                        remote_addr,
368                        move || debug!("did not bootstrap {}: no available slots", remote_addr),
369                    );
370                    self.massa_metrics.inc_bootstrap_peers_failed();
371                }
372            }
373        }
374    }
375
376    /// Checks latest attempt. If too recent, provides the bad news (as an error).
377    /// Updates the latest attempt to "now" if it's all good.
378    ///
379    /// # Error
380    /// The elapsed time which is insufficient
381    fn greedy_client_check(
382        ip_hist_map: &mut HashMap<IpAddr, Instant>,
383        remote_addr: SocketAddr,
384        now: Instant,
385        per_ip_min_interval: Duration,
386    ) -> Result<(), Duration> {
387        let mut res = Ok(());
388        ip_hist_map
389            .entry(remote_addr.ip())
390            .and_modify(|occ| {
391                // Well, let's only update the latest
392                if now.duration_since(*occ) <= per_ip_min_interval {
393                    res = Err(occ.elapsed());
394                } else {
395                    // in list, expired
396                    *occ = now;
397                }
398            })
399            .or_insert(now);
400        res
401    }
402}
403
404/// To be called from a `thread::spawn` invocation
405///
406/// Runs the bootstrap management in a dedicated thread, handling the async by using
407/// a multi-thread-aware tokio runtime (the bs-main-loop runtime, to be exact). When this
408/// function blocks in the `block_on`, it should thread-block, and switch to another session
409///
410/// The arc_counter variable is used as a proxy to keep track the number of active bootstrap
411/// sessions.
412#[allow(clippy::too_many_arguments)]
413fn run_bootstrap_session(
414    mut server: BootstrapServerBinder,
415    arc_counter: Arc<()>,
416    config: BootstrapConfig,
417    remote_addr: SocketAddr,
418    data_execution: Arc<RwLock<dyn FinalStateController>>,
419    version: Version,
420    consensus_command_sender: Box<dyn ConsensusController>,
421    protocol_controller: Box<dyn ProtocolController>,
422    massa_metrics: MassaMetrics,
423) {
424    debug!("running bootstrap for peer {}", remote_addr);
425    let deadline = Instant::now() + config.bootstrap_timeout.to_duration();
426    // TODO: reinstate prevention of bootstrap slot camping. Deadline cancellation is one option
427    let res = manage_bootstrap(
428        &config,
429        &mut server,
430        data_execution,
431        version,
432        consensus_command_sender,
433        protocol_controller,
434        deadline,
435    );
436
437    // This drop allows the server to accept new connections before having to complete the error notifications
438    // account for this session being finished, as well as the root-instance
439    massa_trace!("bootstrap.session.finished", {
440        "sessions_remaining": Arc::strong_count(&arc_counter) - 2
441    });
442    drop(arc_counter);
443    match res {
444        Err(BootstrapError::TimedOut(_)) => {
445            debug!("bootstrap timeout for peer {}", remote_addr);
446            // We allow unused result because we don't care if an error is thrown when
447            // sending the error message to the server we will close the socket anyway.
448            let _ = server.send_error_timeout(format!(
449                "Bootstrap process timedout ({})",
450                format_duration(config.bootstrap_timeout.to_duration())
451            ));
452            massa_metrics.inc_bootstrap_peers_failed();
453        }
454        Err(BootstrapError::ReceivedError(error)) => {
455            debug!(
456                "bootstrap serving error received from peer {}: {}",
457                remote_addr, error
458            );
459            massa_metrics.inc_bootstrap_peers_failed();
460        }
461        Err(err) => {
462            debug!("bootstrap serving error for peer {}: {}", remote_addr, err);
463            // We allow unused result because we don't care if an error is thrown when
464            // sending the error message to the server we will close the socket anyway.
465            let _ = server.send_error_timeout(err.to_string());
466            massa_metrics.inc_bootstrap_peers_failed();
467        }
468        Ok(_) => {
469            info!("bootstrapped peer {}", remote_addr);
470            massa_metrics.inc_bootstrap_peers_success();
471        }
472    }
473}
474
475#[allow(clippy::too_many_arguments)]
476pub fn stream_bootstrap_information(
477    server: &mut BootstrapServerBinder,
478    final_state: Arc<RwLock<dyn FinalStateController>>,
479    consensus_controller: Box<dyn ConsensusController>,
480    mut last_slot: Option<Slot>,
481    mut last_state_step: StreamingStep<Vec<u8>>,
482    mut last_versioning_step: StreamingStep<Vec<u8>>,
483    mut last_consensus_step: StreamingStep<PreHashSet<BlockId>>,
484    mut send_last_start_period: bool,
485    bs_deadline: &Instant,
486    write_timeout: Duration,
487) -> Result<(), BootstrapError> {
488    loop {
489        let current_slot;
490        let state_part;
491        let versioning_part;
492        let last_start_period;
493        let last_slot_before_downtime;
494
495        // Scope of the final state read
496        {
497            let final_state_read = final_state.read();
498
499            last_start_period = if send_last_start_period {
500                Some(final_state_read.get_last_start_period())
501            } else {
502                None
503            };
504            last_slot_before_downtime = if send_last_start_period {
505                Some(*final_state_read.get_last_slot_before_downtime())
506            } else {
507                None
508            };
509
510            let state_part_res = final_state_read
511                .get_database()
512                .read()
513                .get_batch_to_stream(&last_state_step, last_slot);
514
515            match state_part_res {
516                    Ok(part) => state_part = part,
517                    Err(MassaDBError::CacheMissError(s)) if s.as_str() == "all our changes are strictly after last_change_id, we can't be sure we did not miss any" => {
518                        return server.send_msg(write_timeout, BootstrapServerMessage::SlotTooOld);
519                    }
520                    Err(e) => return Err(BootstrapError::GeneralError(format!("Error get_batch_to_stream: {}", e))),
521                };
522
523            let new_state_step = match (&last_state_step, state_part.is_empty()) {
524                // We already finished streaming the state
525                (StreamingStep::Finished(_), _) => StreamingStep::Finished(None),
526
527                // We receive our first empty state batch
528                (StreamingStep::Ongoing(_), true) => StreamingStep::Finished(None),
529
530                // We receive our first empty state batch, but we've just started streaming: warn the user
531                (StreamingStep::Started, true) => {
532                    warn!("State bootstrap is finished but nothing has been streamed yet");
533                    StreamingStep::Finished(None)
534                }
535
536                // We still need to stream the state, we update the current reference to the last_key if needed
537                (StreamingStep::Ongoing(last_key), false) => {
538                    match state_part.new_elements.last_key_value() {
539                        Some((new_last_key, _)) => StreamingStep::Ongoing(new_last_key.clone()), // We received new elements
540                        None => StreamingStep::Ongoing(last_key.clone()), // We only received changes
541                    }
542                }
543
544                // We still need to stream the state
545                (StreamingStep::Started, false) => match state_part.new_elements.last_key_value() {
546                    Some((new_last_key, _)) => StreamingStep::Ongoing(new_last_key.clone()), // We received new elements
547                    None => {
548                        // We only received changes
549                        return Err(BootstrapError::GeneralError(String::from(
550                            "State bootstrap started but we have no new elements to stream",
551                        )));
552                    }
553                },
554            };
555
556            let versioning_part_res = final_state_read
557                .get_database()
558                .read()
559                .get_versioning_batch_to_stream(&last_versioning_step, last_slot);
560
561            match versioning_part_res {
562               Ok(part) => versioning_part = part,
563               Err(MassaDBError::CacheMissError(s)) if s.as_str() == "all our changes are strictly after last_change_id, we can't be sure we did not miss any" => {
564                   return server.send_msg(write_timeout, BootstrapServerMessage::SlotTooOld);
565               }
566               Err(e) => return Err(BootstrapError::GeneralError(format!("Error get_batch_to_stream: {}", e))),
567           };
568
569            let new_versioning_step = match (&last_versioning_step, versioning_part.is_empty()) {
570                // We already finished streaming the versioning
571                (StreamingStep::Finished(_), _) => StreamingStep::Finished(None),
572
573                // We receive our first empty versioning batch
574                (StreamingStep::Ongoing(_), true) => StreamingStep::Finished(None),
575
576                // We receive our first empty versioning batch, but we've just started streaming: warn the user
577                (StreamingStep::Started, true) => {
578                    warn!("Versioning bootstrap is finished but nothing has been streamed yet");
579                    StreamingStep::Finished(None)
580                }
581
582                // We still need to stream the versioning, we update the current reference to the last_key if needed
583                (StreamingStep::Ongoing(last_key), false) => {
584                    match versioning_part.new_elements.last_key_value() {
585                        Some((new_last_key, _)) => StreamingStep::Ongoing(new_last_key.clone()), // We received new elements
586                        None => StreamingStep::Ongoing(last_key.clone()), // We only received changes
587                    }
588                }
589
590                // We still need to stream the versioning
591                (StreamingStep::Started, false) => {
592                    match versioning_part.new_elements.last_key_value() {
593                        Some((new_last_key, _)) => StreamingStep::Ongoing(new_last_key.clone()), // We received new elements
594                        None => {
595                            // We only received changes
596                            return Err(BootstrapError::GeneralError(String::from(
597                                "Versioning bootstrap started but we have no new elements to stream",
598                            )));
599                        }
600                    }
601                }
602            };
603
604            let db_slot = final_state_read
605                .get_database()
606                .read()
607                .get_change_id()
608                .expect(CHANGE_ID_DESER_ERROR);
609
610            if let Some(slot) = last_slot {
611                if slot > db_slot {
612                    return Err(BootstrapError::GeneralError(
613                        "Bootstrap cursor set to future slot".to_string(),
614                    ));
615                }
616            }
617
618            // Update cursors for next turn
619            last_state_step = new_state_step;
620            last_versioning_step = new_versioning_step;
621            last_slot = Some(db_slot);
622            current_slot = db_slot;
623            send_last_start_period = false;
624        }
625
626        // Setup final state global cursor
627        let final_state_global_step =
628            if last_state_step.finished() && last_versioning_step.finished() {
629                StreamingStep::Finished(Some(current_slot))
630            } else {
631                StreamingStep::Ongoing(current_slot)
632            };
633
634        // Stream consensus blocks if final state base bootstrap is finished
635        let mut consensus_part = BootstrapableGraph {
636            final_blocks: Default::default(),
637        };
638        let mut consensus_outdated_ids: PreHashSet<BlockId> = PreHashSet::default();
639
640        if final_state_global_step.finished() {
641            let (part, outdated_ids, new_consensus_step) = consensus_controller
642                .get_bootstrap_part(last_consensus_step, final_state_global_step)?;
643            consensus_part = part;
644            consensus_outdated_ids = outdated_ids;
645            last_consensus_step = new_consensus_step;
646        }
647
648        // Logs for an easier diagnostic if needed
649        debug!(
650            "Final state bootstrap cursor: {:?}",
651            final_state_global_step
652        );
653        debug!(
654            "Consensus blocks bootstrap cursor: {:?}",
655            last_consensus_step
656        );
657        if let StreamingStep::Ongoing(ids) = &last_consensus_step {
658            debug!("Consensus bootstrap cursor length: {}", ids.len());
659        }
660
661        // Exit only when all cursors are finished and there are no pending state/versioning
662        // deltas for this iteration (execution catch-up after consensus may still produce some).
663        // We don't bother with the bs-deadline, as this is the last step of the bootstrap process - defer to general write-timeout
664        if final_state_global_step.finished()
665            && last_consensus_step.finished()
666            && state_part.is_empty()
667            && versioning_part.is_empty()
668        {
669            server.send_msg(write_timeout, BootstrapServerMessage::BootstrapFinished)?;
670            break;
671        }
672
673        let Some(write_timeout) = step_timeout_duration(bs_deadline, &write_timeout) else {
674            return Err(BootstrapError::Interrupted(
675                "insufficient time left to provide next bootstrap part".to_string(),
676            ));
677        };
678        // At this point we know that consensus, final state or both are not finished
679        server.send_msg(
680            write_timeout,
681            BootstrapServerMessage::BootstrapPart {
682                slot: current_slot,
683                state_part,
684                versioning_part,
685                consensus_part,
686                consensus_outdated_ids,
687                last_start_period,
688                last_slot_before_downtime,
689            },
690        )?;
691    }
692    Ok(())
693}
694
695// derives the duration allowed for a step in the bootstrap process.
696// Returns None if the deadline for the entire bs-process has been reached
697fn step_timeout_duration(bs_deadline: &Instant, step_timeout: &Duration) -> Option<Duration> {
698    let now = Instant::now();
699    if now >= *bs_deadline {
700        return None;
701    }
702
703    let remaining = *bs_deadline - now;
704    Some(std::cmp::min(remaining, *step_timeout))
705}
706#[allow(clippy::too_many_arguments)]
707pub(crate) fn manage_bootstrap(
708    bootstrap_config: &BootstrapConfig,
709    server: &mut BootstrapServerBinder,
710    final_state: Arc<RwLock<dyn FinalStateController>>,
711    version: Version,
712    consensus_controller: Box<dyn ConsensusController>,
713    protocol_controller: Box<dyn ProtocolController>,
714    deadline: Instant,
715) -> Result<(), BootstrapError> {
716    massa_trace!("bootstrap.lib.manage_bootstrap", {});
717    let read_error_timeout: Duration = bootstrap_config.read_error_timeout.into();
718
719    let Some(hs_timeout) =
720        step_timeout_duration(&deadline, &bootstrap_config.read_timeout.to_duration())
721    else {
722        return Err(BootstrapError::Interrupted(
723            "insufficient time left to begin handshake".to_string(),
724        ));
725    };
726
727    server.handshake_timeout(version, Some(hs_timeout))?;
728
729    // Check for error from client
730    if Instant::now() + read_error_timeout >= deadline {
731        return Err(BootstrapError::Interrupted(
732            "insufficient time to check for error from client".to_string(),
733        ));
734    };
735    match server.next_timeout(Some(read_error_timeout)) {
736        Err(BootstrapError::TimedOut(_)) => {}
737        Err(e) => return Err(e),
738        Ok(BootstrapClientMessage::BootstrapError { error }) => {
739            return Err(BootstrapError::GeneralError(error));
740        }
741        Ok(msg) => return Err(BootstrapError::UnexpectedClientMessage(Box::new(msg))),
742    };
743
744    // Sync clocks
745    let send_time_timeout =
746        step_timeout_duration(&deadline, &bootstrap_config.write_timeout.to_duration());
747    let Some(next_step_timeout) = send_time_timeout else {
748        return Err(BootstrapError::Interrupted(
749            "insufficient time left to send server time".to_string(),
750        ));
751    };
752    server.send_msg(
753        next_step_timeout,
754        BootstrapServerMessage::BootstrapTime {
755            server_time: MassaTime::now(),
756            version,
757        },
758    )?;
759
760    // A well-behaved client asks for the peers exactly once, right before sending
761    // `BootstrapSuccess`. Serving the request more than once lets a client hold a bootstrap
762    // slot until the deadline for free, and each request costs a round-trip to the peer
763    // management thread, so the repetition is refused.
764    let mut peers_already_sent = false;
765
766    loop {
767        let Some(read_timeout) =
768            step_timeout_duration(&deadline, &bootstrap_config.read_timeout.to_duration())
769        else {
770            return Err(BootstrapError::Interrupted(
771                "insufficient time left to process next message".to_string(),
772            ));
773        };
774        match server.next_timeout(Some(read_timeout)) {
775            Err(BootstrapError::TimedOut(_)) => break Ok(()),
776            Err(e) => break Err(e),
777            Ok(msg) => match msg {
778                BootstrapClientMessage::AskBootstrapPeers => {
779                    if peers_already_sent {
780                        break Err(BootstrapError::UnexpectedClientMessage(Box::new(
781                            BootstrapClientMessage::AskBootstrapPeers,
782                        )));
783                    }
784
785                    let Some(write_timeout) = step_timeout_duration(
786                        &deadline,
787                        &bootstrap_config.write_timeout.to_duration(),
788                    ) else {
789                        return Err(BootstrapError::Interrupted(
790                            "insufficient time left to respond te request for peers".to_string(),
791                        ));
792                    };
793
794                    server.send_msg(
795                        write_timeout,
796                        BootstrapServerMessage::BootstrapPeers {
797                            peers: protocol_controller.get_bootstrap_peers()?,
798                        },
799                    )?;
800
801                    peers_already_sent = true;
802                }
803                BootstrapClientMessage::AskBootstrapPart {
804                    last_slot,
805                    last_state_step,
806                    last_versioning_step,
807                    last_consensus_step,
808                    send_last_start_period,
809                } => {
810                    stream_bootstrap_information(
811                        server,
812                        final_state.clone(),
813                        consensus_controller.clone(),
814                        last_slot,
815                        last_state_step,
816                        last_versioning_step,
817                        last_consensus_step,
818                        send_last_start_period,
819                        &deadline,
820                        bootstrap_config.write_timeout.to_duration(),
821                    )?;
822                }
823                BootstrapClientMessage::BootstrapSuccess => break Ok(()),
824                BootstrapClientMessage::BootstrapError { error } => {
825                    break Err(BootstrapError::ReceivedError(error));
826                }
827            },
828        };
829    }
830}
831
832#[cfg(test)]
833mod updater_tests {
834    use super::*;
835    use std::time::Duration;
836    use tempfile::TempDir;
837
838    #[test]
839    fn updater_survives_a_malformed_acl_reload() {
840        let dir = TempDir::new().unwrap();
841        let white = dir.path().join("whitelist.json");
842        let black = dir.path().join("blacklist.json");
843        // Start from a valid (empty) blacklist so init succeeds.
844        std::fs::write(&black, "[]").unwrap();
845
846        let list = SharedWhiteBlackList::new(white, black.clone()).unwrap();
847        let (stop_tx, stop_rx) = crossbeam::channel::unbounded::<()>();
848        let handle = std::thread::spawn(move || {
849            BootstrapServer::run_updater(list, Duration::from_millis(20), stop_rx)
850        });
851
852        // Corrupt the file so subsequent periodic reloads fail to parse.
853        std::fs::write(&black, "{ not valid json ]").unwrap();
854        // Let several ticks elapse; the updater must keep running.
855        std::thread::sleep(Duration::from_millis(200));
856        assert!(
857            !handle.is_finished(),
858            "updater thread must survive a malformed ACL reload instead of dying"
859        );
860
861        // It must still stop cleanly on request.
862        stop_tx.send(()).unwrap();
863        let res = handle.join().expect("updater thread panicked");
864        assert!(res.is_ok(), "clean stop should return Ok, got {:?}", res);
865    }
866}