1use 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#[derive(Debug, Clone)]
51#[allow(clippy::large_enum_variant)]
52pub enum BootstrapServerMessage {
53 BootstrapTime {
55 server_time: MassaTime,
57 version: Version,
59 },
60 BootstrapPeers {
62 peers: BootstrapPeers,
64 },
65 BootstrapPart {
67 slot: Slot,
69 state_part: StreamBatch<Slot>,
71 versioning_part: StreamBatch<Slot>,
73 consensus_part: BootstrapableGraph,
75 consensus_outdated_ids: PreHashSet<BlockId>,
77 last_start_period: Option<u64>,
79 last_slot_before_downtime: Option<Option<Slot>>,
81 },
82 BootstrapFinished,
84 SlotTooOld,
86 BootstrapError {
88 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
120pub 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 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 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 self.u32_serializer
213 .serialize(&u32::from(MessageServerTypeId::FinalStatePart), buffer)?;
214 self.slot_serializer.serialize(slot, buffer)?;
216 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 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 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 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 self.bootstrapable_graph_serializer
286 .serialize(consensus_part, buffer)?;
287 self.block_id_set_serializer
289 .serialize(consensus_outdated_ids, buffer)?;
290 self.opt_last_start_period_serializer
292 .serialize(last_start_period, buffer)?;
293 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
321pub 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 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 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 updates_batch_budget: bootstrap_batch_allocation_budget(
427 MAX_BOOTSTRAP_MESSAGE_FROM_SERVER_SIZE as usize,
428 ),
429 }
430 }
431
432 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 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 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 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#[derive(Debug, Clone)]
771#[allow(clippy::large_enum_variant)]
772pub enum BootstrapClientMessage {
773 AskBootstrapPeers,
775 AskBootstrapPart {
777 last_slot: Option<Slot>,
779 last_state_step: StreamingStep<Vec<u8>>,
781 last_versioning_step: StreamingStep<Vec<u8>>,
783 last_consensus_step: StreamingStep<PreHashSet<BlockId>>,
785 send_last_start_period: bool,
787 },
788 BootstrapError {
790 error: String,
792 },
793 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
806pub 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 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 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
904pub 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 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 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}