massa_bootstrap/
messages.rs

1// Copyright (c) 2022 MASSA LABS <info@massa.net>
2
3use crate::settings::BootstrapServerMessageDeserializerArgs;
4use massa_consensus_exports::bootstrapable_graph::{
5    BootstrapableGraph, BootstrapableGraphDeserializer, BootstrapableGraphSerializer,
6};
7
8use massa_db_exports::StreamBatch;
9
10use massa_models::block::BlockDeserializerArgs;
11use massa_models::block_id::{BlockId, BlockIdDeserializer, BlockIdSerializer};
12
13use massa_models::config::{
14    bootstrap_batch_allocation_budget, BOOTSTRAP_BATCH_ENTRY_OVERHEAD,
15    MAX_BOOTSTRAP_MESSAGE_FROM_SERVER_SIZE,
16};
17use massa_models::prehash::PreHashSet;
18use massa_models::serialization::{
19    PreHashSetDeserializer, PreHashSetSerializer, VecU8Deserializer, VecU8Serializer,
20};
21use massa_models::slot::{Slot, SlotDeserializer, SlotSerializer};
22use massa_models::streaming_step::{
23    StreamingStep, StreamingStepDeserializer, StreamingStepSerializer,
24};
25use massa_models::version::{Version, VersionDeserializer, VersionSerializer};
26use massa_protocol_exports::{
27    BootstrapPeers, BootstrapPeersDeserializer, BootstrapPeersSerializer,
28};
29use massa_serialization::{
30    BoolDeserializer, BoolSerializer, Deserializer, OptionDeserializer, OptionSerializer,
31    SerializeError, Serializer, U32VarIntDeserializer, U32VarIntSerializer, U64VarIntDeserializer,
32    U64VarIntSerializer,
33};
34use std::collections::BTreeMap;
35
36use massa_time::{MassaTime, MassaTimeDeserializer, MassaTimeSerializer};
37use nom::error::context;
38use nom::multi::{length_data, length_value};
39use nom::sequence::tuple;
40use nom::Parser;
41use nom::{
42    error::{ContextError, ParseError},
43    IResult,
44};
45use num_enum::{IntoPrimitive, TryFromPrimitive};
46use std::convert::TryInto;
47use std::ops::Bound::{Excluded, Included};
48
49/// Messages used during bootstrap by server
50#[derive(Debug, Clone)]
51#[allow(clippy::large_enum_variant)]
52pub enum BootstrapServerMessage {
53    /// Sync clocks
54    BootstrapTime {
55        /// The current time on the bootstrap server.
56        server_time: MassaTime,
57        /// The version of the bootstrap server.
58        version: Version,
59    },
60    /// Bootstrap peers
61    BootstrapPeers {
62        /// Server peers
63        peers: BootstrapPeers,
64    },
65    /// Part of final state and consensus
66    BootstrapPart {
67        /// Slot the state changes are attached to
68        slot: Slot,
69        /// Part of the state in a serialized way
70        state_part: StreamBatch<Slot>,
71        /// Part of the state (specific to versioning) in a serialized way
72        versioning_part: StreamBatch<Slot>,
73        /// Part of the consensus graph
74        consensus_part: BootstrapableGraph,
75        /// Outdated block ids in the current consensus graph bootstrap
76        consensus_outdated_ids: PreHashSet<BlockId>,
77        /// Last Start Period for network restart management
78        last_start_period: Option<u64>,
79        /// Last Slot before downtime for network restart management
80        last_slot_before_downtime: Option<Option<Slot>>,
81    },
82    /// Message sent when the final state and consensus bootstrap are finished
83    BootstrapFinished,
84    /// Slot sent to get state changes is too old
85    SlotTooOld,
86    /// Bootstrap error
87    BootstrapError {
88        /// Error message
89        error: String,
90    },
91}
92
93#[allow(clippy::to_string_trait_impl)]
94impl ToString for BootstrapServerMessage {
95    fn to_string(&self) -> String {
96        match self {
97            BootstrapServerMessage::BootstrapTime { .. } => "BootstrapTime".to_string(),
98            BootstrapServerMessage::BootstrapPeers { .. } => "BootstrapPeers".to_string(),
99            BootstrapServerMessage::BootstrapPart { .. } => "BootstrapPart".to_string(),
100            BootstrapServerMessage::BootstrapFinished => "BootstrapFinished".to_string(),
101            BootstrapServerMessage::SlotTooOld => "SlotTooOld".to_string(),
102            BootstrapServerMessage::BootstrapError { error } => {
103                format!("BootstrapError {{ error: {} }}", error)
104            }
105        }
106    }
107}
108
109#[derive(IntoPrimitive, Debug, Eq, PartialEq, TryFromPrimitive)]
110#[repr(u32)]
111enum MessageServerTypeId {
112    BootstrapTime = 0u32,
113    Peers = 1u32,
114    FinalStatePart = 2u32,
115    FinalStateFinished = 3u32,
116    SlotTooOld = 4u32,
117    BootstrapError = 5u32,
118}
119
120/// Serializer for `BootstrapServerMessage`
121pub struct BootstrapServerMessageSerializer {
122    u32_serializer: U32VarIntSerializer,
123    u64_serializer: U64VarIntSerializer,
124    time_serializer: MassaTimeSerializer,
125    version_serializer: VersionSerializer,
126    peers_serializer: BootstrapPeersSerializer,
127    bootstrapable_graph_serializer: BootstrapableGraphSerializer,
128    block_id_set_serializer: PreHashSetSerializer<BlockId, BlockIdSerializer>,
129    vec_u8_serializer: VecU8Serializer,
130    opt_vec_u8_serializer: OptionSerializer<Vec<u8>, VecU8Serializer>,
131    slot_serializer: SlotSerializer,
132    opt_last_start_period_serializer: OptionSerializer<u64, U64VarIntSerializer>,
133    opt_last_slot_before_downtime_serializer:
134        OptionSerializer<Option<Slot>, OptionSerializer<Slot, SlotSerializer>>,
135}
136
137impl Default for BootstrapServerMessageSerializer {
138    fn default() -> Self {
139        Self::new()
140    }
141}
142
143impl BootstrapServerMessageSerializer {
144    /// Creates a new `BootstrapServerMessageSerializer`
145    pub fn new() -> Self {
146        Self {
147            u32_serializer: U32VarIntSerializer::new(),
148            u64_serializer: U64VarIntSerializer::new(),
149            time_serializer: MassaTimeSerializer::new(),
150            version_serializer: VersionSerializer::new(),
151            peers_serializer: BootstrapPeersSerializer::new(),
152            bootstrapable_graph_serializer: BootstrapableGraphSerializer::new(),
153            block_id_set_serializer: PreHashSetSerializer::new(BlockIdSerializer::new()),
154            vec_u8_serializer: VecU8Serializer::new(),
155            opt_vec_u8_serializer: OptionSerializer::new(VecU8Serializer::new()),
156            slot_serializer: SlotSerializer::new(),
157            opt_last_start_period_serializer: OptionSerializer::new(U64VarIntSerializer::new()),
158            opt_last_slot_before_downtime_serializer: OptionSerializer::new(OptionSerializer::new(
159                SlotSerializer::new(),
160            )),
161        }
162    }
163}
164
165impl Serializer<BootstrapServerMessage> for BootstrapServerMessageSerializer {
166    /// ## Example
167    /// ```rust
168    /// use massa_bootstrap::{BootstrapServerMessage, BootstrapServerMessageSerializer};
169    /// use massa_serialization::Serializer;
170    /// use massa_time::MassaTime;
171    /// use massa_models::version::Version;
172    /// use std::str::FromStr;
173    ///
174    /// let message_serializer = BootstrapServerMessageSerializer::new();
175    /// let bootstrap_server_message = BootstrapServerMessage::BootstrapTime {
176    ///    server_time: MassaTime::from_millis(0),
177    ///    version: Version::from_str("TEST.1.10").unwrap(),
178    /// };
179    /// let mut message_serialized = Vec::new();
180    /// message_serializer.serialize(&bootstrap_server_message, &mut message_serialized).unwrap();
181    /// ```
182    fn serialize(
183        &self,
184        value: &BootstrapServerMessage,
185        buffer: &mut Vec<u8>,
186    ) -> Result<(), SerializeError> {
187        match value {
188            BootstrapServerMessage::BootstrapTime {
189                server_time,
190                version,
191            } => {
192                self.u32_serializer
193                    .serialize(&u32::from(MessageServerTypeId::BootstrapTime), buffer)?;
194                self.time_serializer.serialize(server_time, buffer)?;
195                self.version_serializer.serialize(version, buffer)?;
196            }
197            BootstrapServerMessage::BootstrapPeers { peers } => {
198                self.u32_serializer
199                    .serialize(&u32::from(MessageServerTypeId::Peers), buffer)?;
200                self.peers_serializer.serialize(peers, buffer)?;
201            }
202            BootstrapServerMessage::BootstrapPart {
203                slot,
204                state_part,
205                versioning_part,
206                consensus_part,
207                consensus_outdated_ids,
208                last_start_period,
209                last_slot_before_downtime,
210            } => {
211                // message type
212                self.u32_serializer
213                    .serialize(&u32::from(MessageServerTypeId::FinalStatePart), buffer)?;
214                // slot
215                self.slot_serializer.serialize(slot, buffer)?;
216                // state new_elements
217                let mut state_new_element_buffer: Vec<u8> = Vec::new();
218                for (key, value) in state_part.new_elements.iter() {
219                    self.vec_u8_serializer
220                        .serialize(key, &mut state_new_element_buffer)?;
221                    self.vec_u8_serializer
222                        .serialize(value, &mut state_new_element_buffer)?;
223                }
224                self.u64_serializer.serialize(
225                    &state_new_element_buffer
226                        .len()
227                        .try_into()
228                        .expect("Overflow of state new_elements len"),
229                    buffer,
230                )?;
231                buffer.extend(state_new_element_buffer);
232                // state updates
233                let mut state_updates_buffer: Vec<u8> = Vec::new();
234                for (key, value) in state_part.updates_on_previous_elements.iter() {
235                    self.vec_u8_serializer
236                        .serialize(key, &mut state_updates_buffer)?;
237                    self.opt_vec_u8_serializer
238                        .serialize(value, &mut state_updates_buffer)?;
239                }
240                self.u64_serializer.serialize(
241                    &state_updates_buffer
242                        .len()
243                        .try_into()
244                        .expect("Overflow of state updates len"),
245                    buffer,
246                )?;
247                buffer.extend(state_updates_buffer);
248                self.slot_serializer
249                    .serialize(&state_part.change_id, buffer)?;
250                // versioning new_elements
251                let mut versioning_new_element_buffer: Vec<u8> = Vec::new();
252                for (key, value) in versioning_part.new_elements.iter() {
253                    self.vec_u8_serializer
254                        .serialize(key, &mut versioning_new_element_buffer)?;
255                    self.vec_u8_serializer
256                        .serialize(value, &mut versioning_new_element_buffer)?;
257                }
258                self.u64_serializer.serialize(
259                    &versioning_new_element_buffer
260                        .len()
261                        .try_into()
262                        .expect("Overflow of versioning new_elements len"),
263                    buffer,
264                )?;
265                buffer.extend(versioning_new_element_buffer);
266                // versioning updates
267                let mut versioning_updates_buffer: Vec<u8> = Vec::new();
268                for (key, value) in versioning_part.updates_on_previous_elements.iter() {
269                    self.vec_u8_serializer
270                        .serialize(key, &mut versioning_updates_buffer)?;
271                    self.opt_vec_u8_serializer
272                        .serialize(value, &mut versioning_updates_buffer)?;
273                }
274                self.u64_serializer.serialize(
275                    &versioning_updates_buffer
276                        .len()
277                        .try_into()
278                        .expect("Overflow of versioning updates len"),
279                    buffer,
280                )?;
281                buffer.extend(versioning_updates_buffer);
282                self.slot_serializer
283                    .serialize(&versioning_part.change_id, buffer)?;
284                // consensus graph
285                self.bootstrapable_graph_serializer
286                    .serialize(consensus_part, buffer)?;
287                // consensus outdated ids
288                self.block_id_set_serializer
289                    .serialize(consensus_outdated_ids, buffer)?;
290                // initial state
291                self.opt_last_start_period_serializer
292                    .serialize(last_start_period, buffer)?;
293                // initial state
294                self.opt_last_slot_before_downtime_serializer
295                    .serialize(last_slot_before_downtime, buffer)?;
296            }
297            BootstrapServerMessage::BootstrapFinished => {
298                self.u32_serializer
299                    .serialize(&u32::from(MessageServerTypeId::FinalStateFinished), buffer)?;
300            }
301            BootstrapServerMessage::SlotTooOld => {
302                self.u32_serializer
303                    .serialize(&u32::from(MessageServerTypeId::SlotTooOld), buffer)?;
304            }
305            BootstrapServerMessage::BootstrapError { error } => {
306                self.u32_serializer
307                    .serialize(&u32::from(MessageServerTypeId::BootstrapError), buffer)?;
308                self.u32_serializer.serialize(
309                    &error.len().try_into().map_err(|_| {
310                        SerializeError::GeneralError("Fail to convert usize to u32".to_string())
311                    })?,
312                    buffer,
313                )?;
314                buffer.extend(error.as_bytes())
315            }
316        }
317        Ok(())
318    }
319}
320
321/// Deserializer for `BootstrapServerMessage`
322pub struct BootstrapServerMessageDeserializer {
323    message_id_deserializer: U32VarIntDeserializer,
324    time_deserializer: MassaTimeDeserializer,
325    version_deserializer: VersionDeserializer,
326    peers_deserializer: BootstrapPeersDeserializer,
327    state_new_elements_length_deserializer: U64VarIntDeserializer,
328    versioning_part_new_elements_length_deserializer: U64VarIntDeserializer,
329    stream_batch_updates_length_deserializer: U64VarIntDeserializer,
330    datastore_key_deserializer: VecU8Deserializer,
331    datastore_val_deserializer: VecU8Deserializer,
332    opt_vec_u8_deserializer: OptionDeserializer<Vec<u8>, VecU8Deserializer>,
333    bootstrapable_graph_deserializer: BootstrapableGraphDeserializer,
334    block_id_set_deserializer: PreHashSetDeserializer<BlockId, BlockIdDeserializer>,
335    length_bootstrap_error: U64VarIntDeserializer,
336    slot_deserializer: SlotDeserializer,
337    opt_last_start_period_deserializer: OptionDeserializer<u64, U64VarIntDeserializer>,
338    opt_last_slot_before_downtime_deserializer:
339        OptionDeserializer<Option<Slot>, OptionDeserializer<Slot, SlotDeserializer>>,
340    state_batch_budget: usize,
341    versioning_batch_budget: usize,
342    updates_batch_budget: usize,
343}
344
345impl BootstrapServerMessageDeserializer {
346    /// Builds a deserializer that applies `last_start_period` to bootstrap block header
347    /// deserialization so genesis vs non-genesis parent-count rules are enforced while
348    /// parsing. Pass `None` for the first part; the server sends restart metadata once on
349    /// that part (empty consensus graph), then the client reuses the cached value.
350    pub fn with_last_start_period(
351        args: BootstrapServerMessageDeserializerArgs,
352        last_start_period: Option<u64>,
353    ) -> Self {
354        let mut block_args: BlockDeserializerArgs = (&args).into();
355        block_args.last_start_period = last_start_period;
356        Self {
357            message_id_deserializer: U32VarIntDeserializer::new(Included(0), Included(u32::MAX)),
358            time_deserializer: MassaTimeDeserializer::new((
359                Included(MassaTime::from_millis(0)),
360                Included(MassaTime::from_millis(u64::MAX)),
361            )),
362            version_deserializer: VersionDeserializer::new(),
363            peers_deserializer: BootstrapPeersDeserializer::new(
364                args.max_advertise_length,
365                args.max_listeners_per_peer,
366            ),
367            datastore_key_deserializer: VecU8Deserializer::new(
368                Included(0),
369                Included(args.max_datastore_key_length.into()),
370            ),
371            datastore_val_deserializer: VecU8Deserializer::new(
372                Included(0),
373                Included(args.max_datastore_value_length),
374            ),
375            opt_vec_u8_deserializer: OptionDeserializer::new(VecU8Deserializer::new(
376                Included(0),
377                Included(args.max_datastore_value_length),
378            )),
379            bootstrapable_graph_deserializer: BootstrapableGraphDeserializer::new(
380                block_args,
381                args.max_bootstrap_blocks_length,
382            ),
383            block_id_set_deserializer: PreHashSetDeserializer::new(
384                BlockIdDeserializer::new(),
385                Included(0),
386                Included(args.max_bootstrap_blocks_length.into()),
387            ),
388            length_bootstrap_error: U64VarIntDeserializer::new(
389                Included(0),
390                Included(args.max_bootstrap_error_length),
391            ),
392            versioning_part_new_elements_length_deserializer: U64VarIntDeserializer::new(
393                Included(0),
394                Included(args.max_versioning_elements_size.into()),
395            ),
396            state_new_elements_length_deserializer: U64VarIntDeserializer::new(
397                Included(0),
398                Included(args.max_final_state_elements_size.into()),
399            ),
400            // The updates section is deliberately allowed to be large (see the
401            // `MAX_BOOTSTRAP_MESSAGE_FROM_SERVER_SIZE` doc: update sizes are not capped by
402            // the per-part limits), but it can never exceed the whole message, so bound it
403            // to the maximum message size instead of `u64::MAX`.
404            stream_batch_updates_length_deserializer: U64VarIntDeserializer::new(
405                Included(0),
406                Included(MAX_BOOTSTRAP_MESSAGE_FROM_SERVER_SIZE as u64),
407            ),
408            slot_deserializer: SlotDeserializer::new(
409                (Included(0), Included(u64::MAX)),
410                (Included(0), Excluded(args.thread_count)),
411            ),
412            opt_last_start_period_deserializer: OptionDeserializer::new(
413                U64VarIntDeserializer::new(Included(u64::MIN), Included(u64::MAX)),
414            ),
415            opt_last_slot_before_downtime_deserializer: OptionDeserializer::new(
416                OptionDeserializer::new(SlotDeserializer::new(
417                    (Included(0), Included(u64::MAX)),
418                    (Included(0), Excluded(args.thread_count)),
419                )),
420            ),
421            state_batch_budget: args.max_final_state_batch_allocation as usize,
422            versioning_batch_budget: args.max_versioning_batch_allocation as usize,
423            // The updates sections are deliberately not capped by the per-part size limits (see
424            // MAX_BOOTSTRAP_MESSAGE_FROM_SERVER_SIZE), so their budget derives from the only
425            // bound they do have: the whole message.
426            updates_batch_budget: bootstrap_batch_allocation_budget(
427                MAX_BOOTSTRAP_MESSAGE_FROM_SERVER_SIZE as usize,
428            ),
429        }
430    }
431
432    /// Deserialize a map of `(key, value)` pairs, bounding the in-memory footprint of the result.
433    ///
434    /// Charges every parsed pair the bytes it occupied on the wire plus
435    /// [`BOOTSTRAP_BATCH_ENTRY_OVERHEAD`], and fails once `budget` is spent. Bounding the section
436    /// in wire bytes alone is not enough: an entry can cost two bytes on the wire and two orders
437    /// of magnitude more as a map node, so a flood of tiny entries amplifies well past the size
438    /// the section announced. Pairs are charged as parsed rather than as inserted, so a
439    /// duplicate-key flood is charged too.
440    ///
441    /// The map is folded during parsing rather than collected into an intermediate `Vec` first:
442    /// these sections can be large, so skipping the extra full-size allocation roughly halves the
443    /// peak memory of parsing one.
444    fn deserialize_kv_map<'a, V, E, F>(
445        &self,
446        mut input: &'a [u8],
447        mut parse_entry: F,
448        budget: usize,
449    ) -> IResult<&'a [u8], BTreeMap<Vec<u8>, V>, E>
450    where
451        E: ParseError<&'a [u8]> + ContextError<&'a [u8]>,
452        F: FnMut(&'a [u8]) -> IResult<&'a [u8], (Vec<u8>, V), E>,
453    {
454        let mut acc = BTreeMap::new();
455        let mut footprint = 0usize;
456        while !input.is_empty() {
457            let (rest, (key, value)) = parse_entry(input)?;
458            footprint = footprint
459                .saturating_add(input.len().saturating_sub(rest.len()))
460                .saturating_add(BOOTSTRAP_BATCH_ENTRY_OVERHEAD);
461            if footprint > budget {
462                return Err(nom::Err::Failure(ContextError::add_context(
463                    input,
464                    "bootstrap batch over its allocation budget",
465                    ParseError::from_error_kind(input, nom::error::ErrorKind::Count),
466                )));
467            }
468            acc.insert(key, value);
469            input = rest;
470        }
471        Ok((input, acc))
472    }
473}
474
475impl Deserializer<BootstrapServerMessage> for BootstrapServerMessageDeserializer {
476    /// ## Example
477    /// ```rust
478    /// use massa_bootstrap::{BootstrapServerMessage, BootstrapServerMessageSerializer, BootstrapServerMessageDeserializer};
479    /// use massa_bootstrap::BootstrapServerMessageDeserializerArgs;
480    /// use massa_serialization::{Serializer, Deserializer, DeserializeError};
481    /// use massa_time::MassaTime;
482    /// use massa_models::version::Version;
483    /// use std::str::FromStr;
484    /// use massa_models::config::CHAINID;
485    ///
486    /// let message_serializer = BootstrapServerMessageSerializer::new();
487    /// let args = BootstrapServerMessageDeserializerArgs {
488    ///     thread_count: 32, endorsement_count: 16,
489    ///     max_listeners_per_peer: 1000,
490    ///     max_advertise_length: 1000, max_bootstrap_blocks_length: 1000,
491    ///     max_operations_per_block: 1000, max_versioning_elements_size: 1000,
492    ///     max_ledger_changes_count: 1000, max_datastore_key_length: 255,
493    ///     max_datastore_value_length: 1000,
494    ///     max_final_state_elements_size: 1000,
495    ///     max_final_state_batch_allocation: 1_000_000, max_versioning_batch_allocation: 1_000_000,
496    ///     max_datastore_entry_count: 1000, max_bootstrap_error_length: 1000, max_changes_slot_count: 1000,
497    ///     max_rolls_length: 1000, max_production_stats_length: 1000, max_credits_length: 1000,
498    ///     max_executed_ops_length: 1000, max_ops_changes_length: 1000,
499    ///     mip_store_stats_block_considered: 100,
500    ///     max_denunciations_per_block_header: 128, max_denunciation_changes_length: 1000,
501    ///     chain_id: *CHAINID
502    /// };
503    /// let message_deserializer = BootstrapServerMessageDeserializer::with_last_start_period(args, None);
504    /// let bootstrap_server_message = BootstrapServerMessage::BootstrapTime {
505    ///    server_time: MassaTime::from_millis(0),
506    ///    version: Version::from_str("TEST.1.10").unwrap(),
507    /// };
508    /// let mut message_serialized = Vec::new();
509    /// message_serializer.serialize(&bootstrap_server_message, &mut message_serialized).unwrap();
510    /// let (rest, message_deserialized) = message_deserializer.deserialize::<DeserializeError>(&message_serialized).unwrap();
511    /// match message_deserialized {
512    ///     BootstrapServerMessage::BootstrapTime {
513    ///        server_time,
514    ///        version,
515    ///    } => {
516    ///     assert_eq!(server_time, MassaTime::from_millis(0));
517    ///     assert_eq!(version, Version::from_str("TEST.1.10").unwrap());
518    ///   }
519    ///   _ => panic!("Unexpected message"),
520    /// }
521    /// assert_eq!(rest.len(), 0);
522    /// ```
523    fn deserialize<'a, E: nom::error::ParseError<&'a [u8]> + nom::error::ContextError<&'a [u8]>>(
524        &self,
525        buffer: &'a [u8],
526    ) -> IResult<&'a [u8], BootstrapServerMessage, E> {
527        context("Failed BootstrapServerMessage deserialization", |buffer| {
528            let (input, id) = context("Failed id deserialization", |input| {
529                self.message_id_deserializer.deserialize(input)
530            })
531            .map(|id| {
532                MessageServerTypeId::try_from(id).map_err(|_| {
533                    nom::Err::Error(ParseError::from_error_kind(
534                        buffer,
535                        nom::error::ErrorKind::Eof,
536                    ))
537                })
538            })
539            .parse(buffer)?;
540            match id? {
541                MessageServerTypeId::BootstrapTime => tuple((
542                    context("Failed server_time deserialization", |input| {
543                        self.time_deserializer.deserialize(input)
544                    }),
545                    context("Failed version deserialization", |input| {
546                        self.version_deserializer.deserialize(input)
547                    }),
548                ))
549                .map(
550                    |(server_time, version)| BootstrapServerMessage::BootstrapTime {
551                        server_time,
552                        version,
553                    },
554                )
555                .parse(input),
556                MessageServerTypeId::Peers => context("Failed peers deserialization", |input| {
557                    self.peers_deserializer.deserialize(input)
558                })
559                .map(|peers| BootstrapServerMessage::BootstrapPeers { peers })
560                .parse(input),
561                MessageServerTypeId::FinalStatePart => tuple((
562                    context("Failed slot deserialization", |input| {
563                        self.slot_deserializer.deserialize(input)
564                    }),
565                    context(
566                        "Failed state_part deserialization",
567                        tuple((
568                            context(
569                                "Failed new_elements deserialization",
570                                length_value(
571                                    context("Failed length deserialization", |input| {
572                                        self.state_new_elements_length_deserializer
573                                            .deserialize(input)
574                                    }),
575                                    |input| {
576                                        self.deserialize_kv_map(
577                                            input,
578                                            |input| {
579                                                tuple((
580                                                    |input| {
581                                                        self.datastore_key_deserializer
582                                                            .deserialize(input)
583                                                    },
584                                                    |input| {
585                                                        self.datastore_val_deserializer
586                                                            .deserialize(input)
587                                                    },
588                                                ))
589                                                .parse(input)
590                                            },
591                                            self.state_batch_budget,
592                                        )
593                                    },
594                                ),
595                            ),
596                            context(
597                                "Failed updates deserialization",
598                                length_value(
599                                    context("Failed length deserialization", |input| {
600                                        self.stream_batch_updates_length_deserializer
601                                            .deserialize(input)
602                                    }),
603                                    |input| {
604                                        self.deserialize_kv_map(
605                                            input,
606                                            |input| {
607                                                tuple((
608                                                    |input| {
609                                                        self.datastore_key_deserializer
610                                                            .deserialize(input)
611                                                    },
612                                                    |input| {
613                                                        self.opt_vec_u8_deserializer
614                                                            .deserialize(input)
615                                                    },
616                                                ))
617                                                .parse(input)
618                                            },
619                                            self.updates_batch_budget,
620                                        )
621                                    },
622                                ),
623                            ),
624                            context("Failed slot deserialization", |input| {
625                                self.slot_deserializer.deserialize(input)
626                            }),
627                        )),
628                    ),
629                    context(
630                        "Failed versioning_part deserialization",
631                        tuple((
632                            context(
633                                "Failed new_elements deserialization",
634                                length_value(
635                                    context("Failed length deserialization", |input| {
636                                        self.versioning_part_new_elements_length_deserializer
637                                            .deserialize(input)
638                                    }),
639                                    |input| {
640                                        self.deserialize_kv_map(
641                                            input,
642                                            |input| {
643                                                tuple((
644                                                    |input| {
645                                                        self.datastore_key_deserializer
646                                                            .deserialize(input)
647                                                    },
648                                                    |input| {
649                                                        self.datastore_val_deserializer
650                                                            .deserialize(input)
651                                                    },
652                                                ))
653                                                .parse(input)
654                                            },
655                                            self.versioning_batch_budget,
656                                        )
657                                    },
658                                ),
659                            ),
660                            context(
661                                "Failed updates deserialization",
662                                length_value(
663                                    context("Failed length deserialization", |input| {
664                                        self.stream_batch_updates_length_deserializer
665                                            .deserialize(input)
666                                    }),
667                                    |input| {
668                                        self.deserialize_kv_map(
669                                            input,
670                                            |input| {
671                                                tuple((
672                                                    |input| {
673                                                        self.datastore_key_deserializer
674                                                            .deserialize(input)
675                                                    },
676                                                    |input| {
677                                                        self.opt_vec_u8_deserializer
678                                                            .deserialize(input)
679                                                    },
680                                                ))
681                                                .parse(input)
682                                            },
683                                            self.updates_batch_budget,
684                                        )
685                                    },
686                                ),
687                            ),
688                            context("Failed slot deserialization", |input| {
689                                self.slot_deserializer.deserialize(input)
690                            }),
691                        )),
692                    ),
693                    context("Failed consensus_part deserialization", |input| {
694                        self.bootstrapable_graph_deserializer.deserialize(input)
695                    }),
696                    context("Failed consensus_outdated_ids deserialization", |input| {
697                        self.block_id_set_deserializer.deserialize(input)
698                    }),
699                    context("Failed last_start_period deserialization", |input| {
700                        self.opt_last_start_period_deserializer.deserialize(input)
701                    }),
702                    context(
703                        "Failed last_slot_before_downtime deserialization",
704                        |input| {
705                            self.opt_last_slot_before_downtime_deserializer
706                                .deserialize(input)
707                        },
708                    ),
709                ))
710                .map(
711                    |(
712                        slot,
713                        (state_part_new_elems, state_part_updates, state_part_change_id),
714                        (
715                            versioning_part_new_elems,
716                            versioning_part_updates,
717                            versioning_part_change_id,
718                        ),
719                        consensus_part,
720                        consensus_outdated_ids,
721                        last_start_period,
722                        last_slot_before_downtime,
723                    )| {
724                        let state_part = StreamBatch::<Slot> {
725                            // both already `BTreeMap`s (folded during parsing)
726                            new_elements: state_part_new_elems,
727                            updates_on_previous_elements: state_part_updates,
728                            change_id: state_part_change_id,
729                        };
730                        let versioning_part = StreamBatch::<Slot> {
731                            // both already `BTreeMap`s (folded during parsing)
732                            new_elements: versioning_part_new_elems,
733                            updates_on_previous_elements: versioning_part_updates,
734                            change_id: versioning_part_change_id,
735                        };
736
737                        BootstrapServerMessage::BootstrapPart {
738                            slot,
739                            state_part,
740                            versioning_part,
741                            consensus_part,
742                            consensus_outdated_ids,
743                            last_start_period,
744                            last_slot_before_downtime,
745                        }
746                    },
747                )
748                .parse(input),
749                MessageServerTypeId::FinalStateFinished => {
750                    Ok((input, BootstrapServerMessage::BootstrapFinished))
751                }
752                MessageServerTypeId::SlotTooOld => Ok((input, BootstrapServerMessage::SlotTooOld)),
753                MessageServerTypeId::BootstrapError => context(
754                    "Failed BootstrapError deserialization",
755                    length_data(context("Failed length deserialization", |input| {
756                        self.length_bootstrap_error.deserialize(input)
757                    })),
758                )
759                .map(|error| BootstrapServerMessage::BootstrapError {
760                    error: String::from_utf8_lossy(error).into_owned(),
761                })
762                .parse(input),
763            }
764        })
765        .parse(buffer)
766    }
767}
768
769/// Messages used during bootstrap by client
770#[derive(Debug, Clone)]
771#[allow(clippy::large_enum_variant)]
772pub enum BootstrapClientMessage {
773    /// Ask for bootstrap peers
774    AskBootstrapPeers,
775    /// Ask for a final state and consensus part
776    AskBootstrapPart {
777        /// Slot we are attached to for changes
778        last_slot: Option<Slot>,
779        /// Last received state key
780        last_state_step: StreamingStep<Vec<u8>>,
781        /// Last received versioning key
782        last_versioning_step: StreamingStep<Vec<u8>>,
783        /// Last received consensus block slot
784        last_consensus_step: StreamingStep<PreHashSet<BlockId>>,
785        /// Should be true only for the first part, false later
786        send_last_start_period: bool,
787    },
788    /// Bootstrap error
789    BootstrapError {
790        /// Error message
791        error: String,
792    },
793    /// Bootstrap succeed
794    BootstrapSuccess,
795}
796
797#[derive(IntoPrimitive, Debug, Eq, PartialEq, TryFromPrimitive)]
798#[repr(u32)]
799enum MessageClientTypeId {
800    AskBootstrapPeers = 0u32,
801    AskFinalStatePart = 1u32,
802    BootstrapError = 2u32,
803    BootstrapSuccess = 3u32,
804}
805
806/// Serializer for `BootstrapClientMessage`
807pub struct BootstrapClientMessageSerializer {
808    u32_serializer: U32VarIntSerializer,
809    slot_serializer: SlotSerializer,
810    state_step_serializer: StreamingStepSerializer<Vec<u8>, VecU8Serializer>,
811    block_ids_step_serializer: StreamingStepSerializer<
812        PreHashSet<BlockId>,
813        PreHashSetSerializer<BlockId, BlockIdSerializer>,
814    >,
815    bool_serializer: BoolSerializer,
816}
817
818impl BootstrapClientMessageSerializer {
819    /// Creates a new `BootstrapClientMessageSerializer`
820    pub fn new() -> Self {
821        Self {
822            u32_serializer: U32VarIntSerializer::new(),
823            slot_serializer: SlotSerializer::new(),
824            state_step_serializer: StreamingStepSerializer::new(VecU8Serializer::new()),
825            block_ids_step_serializer: StreamingStepSerializer::new(PreHashSetSerializer::new(
826                BlockIdSerializer::new(),
827            )),
828            bool_serializer: BoolSerializer::new(),
829        }
830    }
831}
832
833impl Default for BootstrapClientMessageSerializer {
834    fn default() -> Self {
835        Self::new()
836    }
837}
838
839impl Serializer<BootstrapClientMessage> for BootstrapClientMessageSerializer {
840    /// ## Example
841    /// ```rust
842    /// use massa_bootstrap::{BootstrapClientMessage, BootstrapClientMessageSerializer};
843    /// use massa_serialization::Serializer;
844    /// use massa_time::MassaTime;
845    /// use massa_models::version::Version;
846    /// use std::str::FromStr;
847    ///
848    /// let message_serializer = BootstrapClientMessageSerializer::new();
849    /// let bootstrap_server_message = BootstrapClientMessage::AskBootstrapPeers;
850    /// let mut message_serialized = Vec::new();
851    /// message_serializer.serialize(&bootstrap_server_message, &mut message_serialized).unwrap();
852    /// ```
853    fn serialize(
854        &self,
855        value: &BootstrapClientMessage,
856        buffer: &mut Vec<u8>,
857    ) -> Result<(), SerializeError> {
858        match value {
859            BootstrapClientMessage::AskBootstrapPeers => {
860                self.u32_serializer
861                    .serialize(&u32::from(MessageClientTypeId::AskBootstrapPeers), buffer)?;
862            }
863            BootstrapClientMessage::AskBootstrapPart {
864                last_slot,
865                last_state_step,
866                last_versioning_step,
867                last_consensus_step,
868                send_last_start_period,
869            } => {
870                self.u32_serializer
871                    .serialize(&u32::from(MessageClientTypeId::AskFinalStatePart), buffer)?;
872                if let Some(slot) = last_slot {
873                    self.slot_serializer.serialize(slot, buffer)?;
874                    self.state_step_serializer
875                        .serialize(last_state_step, buffer)?;
876                    self.state_step_serializer
877                        .serialize(last_versioning_step, buffer)?;
878                    self.block_ids_step_serializer
879                        .serialize(last_consensus_step, buffer)?;
880                    self.bool_serializer
881                        .serialize(send_last_start_period, buffer)?;
882                }
883            }
884            BootstrapClientMessage::BootstrapError { error } => {
885                self.u32_serializer
886                    .serialize(&u32::from(MessageClientTypeId::BootstrapError), buffer)?;
887                self.u32_serializer.serialize(
888                    &error.len().try_into().map_err(|_| {
889                        SerializeError::GeneralError("Fail to convert usize to u32".to_string())
890                    })?,
891                    buffer,
892                )?;
893                buffer.extend(error.as_bytes())
894            }
895            BootstrapClientMessage::BootstrapSuccess => {
896                self.u32_serializer
897                    .serialize(&u32::from(MessageClientTypeId::BootstrapSuccess), buffer)?;
898            }
899        }
900        Ok(())
901    }
902}
903
904/// Deserializer for `BootstrapClientMessage`
905pub struct BootstrapClientMessageDeserializer {
906    id_deserializer: U32VarIntDeserializer,
907    length_error_deserializer: U32VarIntDeserializer,
908    slot_deserializer: SlotDeserializer,
909    state_step_deserializer: StreamingStepDeserializer<Vec<u8>, VecU8Deserializer>,
910    block_ids_step_deserializer: StreamingStepDeserializer<
911        PreHashSet<BlockId>,
912        PreHashSetDeserializer<BlockId, BlockIdDeserializer>,
913    >,
914    bool_deserializer: BoolDeserializer,
915}
916
917impl BootstrapClientMessageDeserializer {
918    /// Creates a new `BootstrapClientMessageDeserializer`
919    pub fn new(
920        thread_count: u8,
921        max_datastore_key_length: u8,
922        max_consensus_block_ids: u64,
923    ) -> Self {
924        Self {
925            id_deserializer: U32VarIntDeserializer::new(Included(0), Included(u32::MAX)),
926            length_error_deserializer: U32VarIntDeserializer::new(Included(0), Included(u32::MAX)),
927            slot_deserializer: SlotDeserializer::new(
928                (Included(0), Included(u64::MAX)),
929                (Included(0), Excluded(thread_count)),
930            ),
931            state_step_deserializer: StreamingStepDeserializer::new(VecU8Deserializer::new(
932                Included(0),
933                Included(max_datastore_key_length.into()),
934            )),
935            block_ids_step_deserializer: StreamingStepDeserializer::new(
936                PreHashSetDeserializer::new(
937                    BlockIdDeserializer::new(),
938                    Included(0),
939                    Included(max_consensus_block_ids),
940                ),
941            ),
942            bool_deserializer: BoolDeserializer::new(),
943        }
944    }
945}
946
947impl Deserializer<BootstrapClientMessage> for BootstrapClientMessageDeserializer {
948    /// ## Example
949    /// ```rust
950    /// use massa_bootstrap::{BootstrapClientMessage, BootstrapClientMessageSerializer, BootstrapClientMessageDeserializer};
951    /// use massa_serialization::{Serializer, Deserializer, DeserializeError};
952    /// use massa_time::MassaTime;
953    /// use massa_models::version::Version;
954    /// use std::str::FromStr;
955    ///
956    /// let message_serializer = BootstrapClientMessageSerializer::new();
957    /// let message_deserializer = BootstrapClientMessageDeserializer::new(32, 255, 50);
958    /// let bootstrap_server_message = BootstrapClientMessage::AskBootstrapPeers;
959    /// let mut message_serialized = Vec::new();
960    /// message_serializer.serialize(&bootstrap_server_message, &mut message_serialized).unwrap();
961    /// let (rest, message_deserialized) = message_deserializer.deserialize::<DeserializeError>(&message_serialized).unwrap();
962    /// match message_deserialized {
963    ///     BootstrapClientMessage::AskBootstrapPeers => (),
964    ///   _ => panic!("Unexpected message"),
965    /// };
966    /// assert_eq!(rest.len(), 0);
967    /// ```
968    fn deserialize<'a, E: ParseError<&'a [u8]> + ContextError<&'a [u8]>>(
969        &self,
970        buffer: &'a [u8],
971    ) -> IResult<&'a [u8], BootstrapClientMessage, E> {
972        context("Failed BootstrapClientMessage deserialization", |buffer| {
973            let (input, id) = context("Failed id deserialization", |input| {
974                self.id_deserializer.deserialize(input)
975            })
976            .map(|id| {
977                MessageClientTypeId::try_from(id).map_err(|_| {
978                    nom::Err::Error(ParseError::from_error_kind(
979                        buffer,
980                        nom::error::ErrorKind::Eof,
981                    ))
982                })
983            })
984            .parse(buffer)?;
985            match id? {
986                MessageClientTypeId::AskBootstrapPeers => {
987                    Ok((input, BootstrapClientMessage::AskBootstrapPeers))
988                }
989                MessageClientTypeId::AskFinalStatePart => {
990                    if input.is_empty() {
991                        Ok((
992                            input,
993                            BootstrapClientMessage::AskBootstrapPart {
994                                last_slot: None,
995                                last_state_step: StreamingStep::Started,
996                                last_versioning_step: StreamingStep::Started,
997                                last_consensus_step: StreamingStep::Started,
998                                send_last_start_period: true,
999                            },
1000                        ))
1001                    } else {
1002                        tuple((
1003                            context("Failed last_slot deserialization", |input| {
1004                                self.slot_deserializer.deserialize(input)
1005                            }),
1006                            context("Failed last_state_step deserialization", |input| {
1007                                self.state_step_deserializer.deserialize(input)
1008                            }),
1009                            context("Failed last_versioning_step deserialization", |input| {
1010                                self.state_step_deserializer.deserialize(input)
1011                            }),
1012                            context("Failed last_consensus_step deserialization", |input| {
1013                                self.block_ids_step_deserializer.deserialize(input)
1014                            }),
1015                            context("Failed send_last_start_period deserialization", |input| {
1016                                self.bool_deserializer.deserialize(input)
1017                            }),
1018                        ))
1019                        .map(
1020                            |(
1021                                last_slot,
1022                                last_state_step,
1023                                last_versioning_step,
1024                                last_consensus_step,
1025                                send_last_start_period,
1026                            )| {
1027                                BootstrapClientMessage::AskBootstrapPart {
1028                                    last_slot: Some(last_slot),
1029                                    last_state_step,
1030                                    last_versioning_step,
1031                                    last_consensus_step,
1032                                    send_last_start_period,
1033                                }
1034                            },
1035                        )
1036                        .parse(input)
1037                    }
1038                }
1039                MessageClientTypeId::BootstrapError => context(
1040                    "Failed BootstrapError deserialization",
1041                    length_data(context("Failed length deserialization", |input| {
1042                        self.length_error_deserializer.deserialize(input)
1043                    })),
1044                )
1045                .map(|error| BootstrapClientMessage::BootstrapError {
1046                    error: String::from_utf8_lossy(error).into_owned(),
1047                })
1048                .parse(input),
1049                MessageClientTypeId::BootstrapSuccess => {
1050                    Ok((input, BootstrapClientMessage::BootstrapSuccess))
1051                }
1052            }
1053        })
1054        .parse(buffer)
1055    }
1056}