1use crate::{changes::AsyncPoolChanges, config::AsyncPoolConfig};
6use massa_db_exports::{
7 DBBatch, MassaDirection, MassaIteratorMode, ShareableMassaDBController, ASYNC_POOL_PREFIX,
8 MESSAGE_ID_DESER_ERROR, MESSAGE_ID_SER_ERROR, MESSAGE_SER_ERROR, STATE_CF,
9};
10use massa_models::{
11 address::Address,
12 async_msg::{
13 AsyncMessage, AsyncMessageDeserializer, AsyncMessageSerializer, AsyncMessageUpdate,
14 },
15 async_msg_id::{AsyncMessageId, AsyncMessageIdSerializer},
16 types::Applicable,
17};
18use massa_models::{
19 async_msg_id::AsyncMessageIdDeserializer,
20 types::{SetOrKeep, SetUpdateOrDelete},
21};
22use massa_serialization::{
23 DeserializeError, Deserializer, SerializeError, Serializer, U64VarIntDeserializer,
24 U64VarIntSerializer,
25};
26use nom::{
27 error::{context, ContextError, ParseError},
28 multi::length_count,
29 sequence::tuple,
30 IResult, Parser,
31};
32use std::collections::BTreeMap;
33use std::ops::Bound::Included;
34
35const EMISSION_SLOT_IDENT: u8 = 0u8;
36const EMISSION_INDEX_IDENT: u8 = 1u8;
37const SENDER_IDENT: u8 = 2u8;
38const DESTINATION_IDENT: u8 = 3u8;
39const FUNCTION_IDENT: u8 = 4u8;
40const MAX_GAS_IDENT: u8 = 5u8;
41const FEE_IDENT: u8 = 6u8;
42const COINS_IDENT: u8 = 7u8;
43const VALIDITY_START_IDENT: u8 = 8u8;
44const VALIDITY_END_IDENT: u8 = 9u8;
45const FUNCTION_PARAMS_IDENT: u8 = 10u8;
46const TRIGGER_IDENT: u8 = 11u8;
47const CAN_BE_EXECUTED_IDENT: u8 = 12u8;
48
49#[macro_export]
51macro_rules! emission_slot_key {
52 ($id:expr) => {
53 [
54 &ASYNC_POOL_PREFIX.as_bytes(),
55 &$id[..],
56 &[EMISSION_SLOT_IDENT],
57 ]
58 .concat()
59 };
60}
61
62#[macro_export]
64macro_rules! emission_index_key {
65 ($id:expr) => {
66 [
67 &ASYNC_POOL_PREFIX.as_bytes(),
68 &$id[..],
69 &[EMISSION_INDEX_IDENT],
70 ]
71 .concat()
72 };
73}
74
75#[macro_export]
77macro_rules! sender_key {
78 ($id:expr) => {
79 [&ASYNC_POOL_PREFIX.as_bytes(), &$id[..], &[SENDER_IDENT]].concat()
80 };
81}
82
83#[macro_export]
85macro_rules! destination_key {
86 ($id:expr) => {
87 [
88 &ASYNC_POOL_PREFIX.as_bytes(),
89 &$id[..],
90 &[DESTINATION_IDENT],
91 ]
92 .concat()
93 };
94}
95
96#[macro_export]
98macro_rules! function_key {
99 ($id:expr) => {
100 [&ASYNC_POOL_PREFIX.as_bytes(), &$id[..], &[FUNCTION_IDENT]].concat()
101 };
102}
103
104#[macro_export]
106macro_rules! max_gas_key {
107 ($id:expr) => {
108 [&ASYNC_POOL_PREFIX.as_bytes(), &$id[..], &[MAX_GAS_IDENT]].concat()
109 };
110}
111
112#[macro_export]
114macro_rules! fee_key {
115 ($id:expr) => {
116 [&ASYNC_POOL_PREFIX.as_bytes(), &$id[..], &[FEE_IDENT]].concat()
117 };
118}
119
120#[macro_export]
122macro_rules! coins_key {
123 ($id:expr) => {
124 [&ASYNC_POOL_PREFIX.as_bytes(), &$id[..], &[COINS_IDENT]].concat()
125 };
126}
127
128#[macro_export]
130macro_rules! validity_start_key {
131 ($id:expr) => {
132 [
133 &ASYNC_POOL_PREFIX.as_bytes(),
134 &$id[..],
135 &[VALIDITY_START_IDENT],
136 ]
137 .concat()
138 };
139}
140
141#[macro_export]
143macro_rules! validity_end_key {
144 ($id:expr) => {
145 [
146 &ASYNC_POOL_PREFIX.as_bytes(),
147 &$id[..],
148 &[VALIDITY_END_IDENT],
149 ]
150 .concat()
151 };
152}
153
154#[macro_export]
156macro_rules! function_params_key {
157 ($id:expr) => {
158 [
159 &ASYNC_POOL_PREFIX.as_bytes(),
160 &$id[..],
161 &[FUNCTION_PARAMS_IDENT],
162 ]
163 .concat()
164 };
165}
166
167#[macro_export]
169macro_rules! trigger_key {
170 ($id:expr) => {
171 [&ASYNC_POOL_PREFIX.as_bytes(), &$id[..], &[TRIGGER_IDENT]].concat()
172 };
173}
174
175#[macro_export]
177macro_rules! can_be_executed_key {
178 ($id:expr) => {
179 [
180 &ASYNC_POOL_PREFIX.as_bytes(),
181 &$id[..],
182 &[CAN_BE_EXECUTED_IDENT],
183 ]
184 .concat()
185 };
186}
187
188#[macro_export]
190macro_rules! message_id_prefix {
191 ($id:expr) => {
192 [&ASYNC_POOL_PREFIX.as_bytes(), &$id[..]].concat()
193 };
194}
195
196#[derive(Clone)]
197pub struct AsyncPool {
201 pub config: AsyncPoolConfig,
203 pub db: ShareableMassaDBController,
204 pub message_cache: BTreeMap<AsyncMessageId, AsyncMessage>,
206 message_id_serializer: AsyncMessageIdSerializer,
207 message_serializer: AsyncMessageSerializer,
208 message_id_deserializer: AsyncMessageIdDeserializer,
209 message_deserializer_db: AsyncMessageDeserializer,
210}
211
212impl AsyncPool {
213 pub fn new(config: AsyncPoolConfig, db: ShareableMassaDBController) -> AsyncPool {
215 AsyncPool {
216 config: config.clone(),
217 db,
218 message_cache: Default::default(),
219 message_id_serializer: AsyncMessageIdSerializer::new(),
220 message_serializer: AsyncMessageSerializer::new(true),
221 message_id_deserializer: AsyncMessageIdDeserializer::new(config.thread_count),
222 message_deserializer_db: AsyncMessageDeserializer::new(
223 config.thread_count,
224 config.max_function_length,
225 config.max_function_params_length,
226 config.max_key_length,
227 true,
228 ),
229 }
230 }
231
232 pub fn recompute_message_cache(&mut self) {
234 self.message_cache.clear();
235
236 let db = self.db.read();
237
238 let mut last_id: Option<Vec<u8>> = None;
240
241 while let Some((serialized_message_id, _)) = match last_id {
242 Some(id) => db
243 .iterator_cf(
244 STATE_CF,
245 MassaIteratorMode::From(&can_be_executed_key!(id), MassaDirection::Forward),
246 )
247 .nth(1),
248 None => db
249 .prefix_iterator_cf(STATE_CF, ASYNC_POOL_PREFIX.as_bytes())
250 .next(),
251 } {
252 if !serialized_message_id.starts_with(ASYNC_POOL_PREFIX.as_bytes()) {
253 break;
254 }
255
256 let (_, message_id) = self
257 .message_id_deserializer
258 .deserialize::<DeserializeError>(&serialized_message_id[ASYNC_POOL_PREFIX.len()..])
259 .expect(MESSAGE_ID_DESER_ERROR);
260
261 if let Some(message) = self.load_message(&message_id) {
262 self.message_cache.insert(message_id, message);
263 }
264
265 last_id = Some(
267 serialized_message_id[ASYNC_POOL_PREFIX.len()..serialized_message_id.len() - 1]
268 .to_vec(),
269 );
270 }
271 }
272
273 pub fn reset(&mut self) {
277 self.db
278 .write()
279 .delete_prefix(ASYNC_POOL_PREFIX, STATE_CF, None);
280 self.recompute_message_cache();
281 }
282
283 pub fn apply_changes_to_batch(&mut self, changes: &AsyncPoolChanges, batch: &mut DBBatch) {
289 for change in changes.0.iter() {
290 match change {
291 (id, SetUpdateOrDelete::Set(message)) => {
292 self.put_entry(id, message.clone(), batch);
293 self.message_cache.insert(*id, message.clone());
294 }
295
296 (id, SetUpdateOrDelete::Update(message_update)) => {
297 self.update_entry(id, message_update.clone(), batch);
298
299 self.message_cache.entry(*id).and_modify(|message_info| {
300 message_info.apply(message_update.clone());
301 });
302 }
303
304 (id, SetUpdateOrDelete::Delete) => {
305 self.delete_entry(id, batch);
306 self.message_cache.remove(id);
307 }
308 }
309 }
310 }
311
312 fn load_message(&self, message_id: &AsyncMessageId) -> Option<AsyncMessage> {
315 let db = self.db.read();
316
317 let mut serialized_message_id = Vec::new();
318 self.message_id_serializer
319 .serialize(message_id, &mut serialized_message_id)
320 .expect(MESSAGE_ID_SER_ERROR);
321
322 let mut serialized_message: Vec<u8> = Vec::new();
323 for (serialized_key, serialized_value) in
324 db.prefix_iterator_cf(STATE_CF, &message_id_prefix!(serialized_message_id))
325 {
326 if !serialized_key.starts_with(&message_id_prefix!(serialized_message_id)) {
327 break;
328 }
329
330 serialized_message.extend(serialized_value.iter());
331 }
332
333 match self
334 .message_deserializer_db
335 .deserialize::<DeserializeError>(&serialized_message)
336 {
337 Ok((_, message)) => Some(message),
338 _ => None,
339 }
340 }
341
342 pub fn fetch_messages(
344 &self,
345 message_ids: &[AsyncMessageId],
346 ) -> Vec<(AsyncMessageId, Option<AsyncMessage>)> {
347 let mut fetched_messages = Vec::with_capacity(message_ids.len());
348
349 for message_id in message_ids {
350 fetched_messages.push((*message_id, self.message_cache.get(message_id).cloned()));
351 }
352
353 fetched_messages
354 }
355
356 pub fn is_key_value_valid(&self, serialized_key: &[u8], serialized_value: &[u8]) -> bool {
358 if !serialized_key.starts_with(ASYNC_POOL_PREFIX.as_bytes()) {
359 return false;
360 }
361
362 let Ok((rest, _id)) = self
363 .message_id_deserializer
364 .deserialize::<DeserializeError>(&serialized_key[ASYNC_POOL_PREFIX.len()..])
365 else {
366 return false;
367 };
368 if rest.len() != 1 {
369 return false;
370 }
371
372 match rest[0] {
373 EMISSION_SLOT_IDENT => {
374 let Ok((rest, _value)) = self
375 .message_deserializer_db
376 .slot_deserializer
377 .deserialize::<DeserializeError>(serialized_value)
378 else {
379 return false;
380 };
381 if !rest.is_empty() {
382 return false;
383 }
384 }
385 EMISSION_INDEX_IDENT => {
386 let Ok((rest, _value)) = self
387 .message_deserializer_db
388 .emission_index_deserializer
389 .deserialize::<DeserializeError>(serialized_value)
390 else {
391 return false;
392 };
393 if !rest.is_empty() {
394 return false;
395 }
396 }
397 SENDER_IDENT => {
398 let Ok((rest, _value)): std::result::Result<
399 (&[u8], Address),
400 nom::Err<massa_serialization::DeserializeError<'_>>,
401 > = self
402 .message_deserializer_db
403 .address_deserializer
404 .deserialize::<DeserializeError>(serialized_value)
405 else {
406 return false;
407 };
408 if !rest.is_empty() {
409 return false;
410 }
411 }
412 DESTINATION_IDENT => {
413 let Ok((rest, _value)): std::result::Result<
414 (&[u8], Address),
415 nom::Err<massa_serialization::DeserializeError<'_>>,
416 > = self
417 .message_deserializer_db
418 .address_deserializer
419 .deserialize::<DeserializeError>(serialized_value)
420 else {
421 return false;
422 };
423 if !rest.is_empty() {
424 return false;
425 }
426 }
427 FUNCTION_IDENT => {
428 let Some(len) = serialized_value.first() else {
429 return false;
430 };
431
432 if serialized_value.len() != *len as usize + 1 {
433 return false;
434 }
435
436 let Ok(_value) = String::from_utf8(serialized_value[1..].to_vec()) else {
437 return false;
438 };
439 }
440 MAX_GAS_IDENT => {
441 let Ok((rest, _value)) = self
442 .message_deserializer_db
443 .max_gas_deserializer
444 .deserialize::<DeserializeError>(serialized_value)
445 else {
446 return false;
447 };
448 if !rest.is_empty() {
449 return false;
450 }
451 }
452 FEE_IDENT => {
453 let Ok((rest, _value)) = self
454 .message_deserializer_db
455 .amount_deserializer
456 .deserialize::<DeserializeError>(serialized_value)
457 else {
458 return false;
459 };
460 if !rest.is_empty() {
461 return false;
462 }
463 }
464 COINS_IDENT => {
465 let Ok((rest, _value)) = self
466 .message_deserializer_db
467 .amount_deserializer
468 .deserialize::<DeserializeError>(serialized_value)
469 else {
470 return false;
471 };
472 if !rest.is_empty() {
473 return false;
474 }
475 }
476 VALIDITY_START_IDENT => {
477 let Ok((rest, _value)) = self
478 .message_deserializer_db
479 .slot_deserializer
480 .deserialize::<DeserializeError>(serialized_value)
481 else {
482 return false;
483 };
484 if !rest.is_empty() {
485 return false;
486 }
487 }
488 VALIDITY_END_IDENT => {
489 let Ok((rest, _value)) = self
490 .message_deserializer_db
491 .slot_deserializer
492 .deserialize::<DeserializeError>(serialized_value)
493 else {
494 return false;
495 };
496 if !rest.is_empty() {
497 return false;
498 }
499 }
500 FUNCTION_PARAMS_IDENT => {
501 let Ok((rest, _value)) = self
502 .message_deserializer_db
503 .function_params_deserializer
504 .deserialize::<DeserializeError>(serialized_value)
505 else {
506 return false;
507 };
508 if !rest.is_empty() {
509 return false;
510 }
511 }
512 TRIGGER_IDENT => {
513 let Ok((rest, _value)) = self
514 .message_deserializer_db
515 .trigger_deserializer
516 .deserialize::<DeserializeError>(serialized_value)
517 else {
518 return false;
519 };
520 if !rest.is_empty() {
521 return false;
522 }
523 }
524 CAN_BE_EXECUTED_IDENT => {
525 let Ok((rest, _value)) = self
526 .message_deserializer_db
527 .bool_deserializer
528 .deserialize::<DeserializeError>(serialized_value)
529 else {
530 return false;
531 };
532 if !rest.is_empty() {
533 return false;
534 }
535 }
536 _ => {
537 return false;
538 }
539 }
540
541 true
542 }
543}
544
545pub struct AsyncPoolSerializer {
547 u64_serializer: U64VarIntSerializer,
548 async_message_id_serializer: AsyncMessageIdSerializer,
549 async_message_serializer: AsyncMessageSerializer,
550}
551
552impl Default for AsyncPoolSerializer {
553 fn default() -> Self {
554 Self::new()
555 }
556}
557
558impl AsyncPoolSerializer {
559 pub fn new() -> Self {
561 Self {
562 u64_serializer: U64VarIntSerializer::new(),
563 async_message_id_serializer: AsyncMessageIdSerializer::new(),
564 async_message_serializer: AsyncMessageSerializer::new(true),
565 }
566 }
567}
568
569impl Serializer<BTreeMap<AsyncMessageId, AsyncMessage>> for AsyncPoolSerializer {
570 fn serialize(
571 &self,
572 value: &BTreeMap<AsyncMessageId, AsyncMessage>,
573 buffer: &mut Vec<u8>,
574 ) -> Result<(), SerializeError> {
575 self.u64_serializer
577 .serialize(&(value.len() as u64), buffer)?;
578 for (message_id, message) in value {
580 self.async_message_id_serializer
581 .serialize(message_id, buffer)?;
582 self.async_message_serializer.serialize(message, buffer)?;
583 }
584 Ok(())
585 }
586}
587
588pub struct AsyncPoolDeserializer {
590 u64_deserializer: U64VarIntDeserializer,
591 async_message_id_deserializer: AsyncMessageIdDeserializer,
592 async_message_deserializer_db: AsyncMessageDeserializer,
593}
594
595impl AsyncPoolDeserializer {
596 pub fn new(
598 thread_count: u8,
599 max_async_pool_length: u64,
600 max_function_length: u16,
601 max_parameters_length: u64,
602 max_key_length: u32,
603 ) -> AsyncPoolDeserializer {
604 AsyncPoolDeserializer {
605 u64_deserializer: U64VarIntDeserializer::new(
606 Included(0),
607 Included(max_async_pool_length),
608 ),
609 async_message_id_deserializer: AsyncMessageIdDeserializer::new(thread_count),
610 async_message_deserializer_db: AsyncMessageDeserializer::new(
611 thread_count,
612 max_function_length,
613 max_parameters_length,
614 max_key_length,
615 true,
616 ),
617 }
618 }
619}
620
621impl Deserializer<BTreeMap<AsyncMessageId, AsyncMessage>> for AsyncPoolDeserializer {
622 fn deserialize<'a, E: ParseError<&'a [u8]> + ContextError<&'a [u8]>>(
623 &self,
624 buffer: &'a [u8],
625 ) -> IResult<&'a [u8], BTreeMap<AsyncMessageId, AsyncMessage>, E> {
626 context(
627 "Failed async_pool_part deserialization",
628 length_count(
629 context("Failed length deserialization", |input| {
630 self.u64_deserializer.deserialize(input)
631 }),
632 tuple((
633 context("Failed async_message_id deserialization", |input| {
634 self.async_message_id_deserializer.deserialize(input)
635 }),
636 context("Failed async_message deserialization", |input| {
637 self.async_message_deserializer_db.deserialize(input)
638 }),
639 )),
640 ),
641 )
642 .map(|vec| vec.into_iter().collect())
643 .parse(buffer)
644 }
645}
646
647impl AsyncPool {
649 fn put_entry(&self, message_id: &AsyncMessageId, message: AsyncMessage, batch: &mut DBBatch) {
656 let db = self.db.read();
657
658 let mut serialized_message_id = Vec::new();
659 self.message_id_serializer
660 .serialize(message_id, &mut serialized_message_id)
661 .expect(MESSAGE_ID_SER_ERROR);
662
663 let mut serialized_emission_slot = Vec::new();
665 self.message_serializer
666 .slot_serializer
667 .serialize(&message.emission_slot, &mut serialized_emission_slot)
668 .expect(MESSAGE_SER_ERROR);
669 db.put_or_update_entry_value(
670 batch,
671 emission_slot_key!(serialized_message_id),
672 &serialized_emission_slot,
673 );
674
675 let mut serialized_emission_index = Vec::new();
677 self.message_serializer
678 .u64_serializer
679 .serialize(&message.emission_index, &mut serialized_emission_index)
680 .expect(MESSAGE_SER_ERROR);
681 db.put_or_update_entry_value(
682 batch,
683 emission_index_key!(serialized_message_id),
684 &serialized_emission_index,
685 );
686
687 let mut serialized_sender = Vec::new();
689 self.message_serializer
690 .address_serializer
691 .serialize(&message.sender, &mut serialized_sender)
692 .expect(MESSAGE_SER_ERROR);
693 db.put_or_update_entry_value(
694 batch,
695 sender_key!(serialized_message_id),
696 &serialized_sender,
697 );
698
699 let mut serialized_destination = Vec::new();
701 self.message_serializer
702 .address_serializer
703 .serialize(&message.destination, &mut serialized_destination)
704 .expect(MESSAGE_SER_ERROR);
705 db.put_or_update_entry_value(
706 batch,
707 destination_key!(serialized_message_id),
708 &serialized_destination,
709 );
710
711 let mut serialized_function = Vec::new();
713 self.message_serializer
714 .function_serializer
715 .serialize(&message.function, &mut serialized_function)
716 .expect(MESSAGE_SER_ERROR);
717 db.put_or_update_entry_value(
718 batch,
719 function_key!(serialized_message_id),
720 &serialized_function,
721 );
722
723 let mut serialized_max_gas = Vec::new();
725 self.message_serializer
726 .u64_serializer
727 .serialize(&message.max_gas, &mut serialized_max_gas)
728 .expect(MESSAGE_SER_ERROR);
729 db.put_or_update_entry_value(
730 batch,
731 max_gas_key!(serialized_message_id),
732 &serialized_max_gas,
733 );
734
735 let mut serialized_fee = Vec::new();
737 self.message_serializer
738 .amount_serializer
739 .serialize(&message.fee, &mut serialized_fee)
740 .expect(MESSAGE_SER_ERROR);
741 db.put_or_update_entry_value(batch, fee_key!(serialized_message_id), &serialized_fee);
742
743 let mut serialized_coins = Vec::new();
745 self.message_serializer
746 .amount_serializer
747 .serialize(&message.coins, &mut serialized_coins)
748 .expect(MESSAGE_SER_ERROR);
749 db.put_or_update_entry_value(batch, coins_key!(serialized_message_id), &serialized_coins);
750
751 let mut serialized_validity_start = Vec::new();
753 self.message_serializer
754 .slot_serializer
755 .serialize(&message.validity_start, &mut serialized_validity_start)
756 .expect(MESSAGE_SER_ERROR);
757 db.put_or_update_entry_value(
758 batch,
759 validity_start_key!(serialized_message_id),
760 &serialized_validity_start,
761 );
762
763 let mut serialized_validity_end = Vec::new();
765 self.message_serializer
766 .slot_serializer
767 .serialize(&message.validity_end, &mut serialized_validity_end)
768 .expect(MESSAGE_SER_ERROR);
769 db.put_or_update_entry_value(
770 batch,
771 validity_end_key!(serialized_message_id),
772 &serialized_validity_end,
773 );
774
775 let mut serialized_params = Vec::new();
777 self.message_serializer
778 .function_params_serializer
779 .serialize(&message.function_params, &mut serialized_params)
780 .expect(MESSAGE_SER_ERROR);
781 db.put_or_update_entry_value(
782 batch,
783 function_params_key!(serialized_message_id),
784 &serialized_params,
785 );
786
787 let mut serialized_trigger = Vec::new();
789 self.message_serializer
790 .trigger_serializer
791 .serialize(&message.trigger, &mut serialized_trigger)
792 .expect(MESSAGE_SER_ERROR);
793 db.put_or_update_entry_value(
794 batch,
795 trigger_key!(serialized_message_id),
796 &serialized_trigger,
797 );
798
799 let mut serialized_can_be_executed = Vec::new();
801 self.message_serializer
802 .bool_serializer
803 .serialize(&message.can_be_executed, &mut serialized_can_be_executed)
804 .expect(MESSAGE_SER_ERROR);
805 db.put_or_update_entry_value(
806 batch,
807 can_be_executed_key!(serialized_message_id),
808 &serialized_can_be_executed,
809 );
810 }
811
812 fn update_entry(
818 &self,
819 message_id: &AsyncMessageId,
820 message_update: AsyncMessageUpdate,
821 batch: &mut DBBatch,
822 ) {
823 let db = self.db.read();
824
825 let mut serialized_message_id = Vec::new();
826 self.message_id_serializer
827 .serialize(message_id, &mut serialized_message_id)
828 .expect(MESSAGE_ID_SER_ERROR);
829
830 if let SetOrKeep::Set(emission_slot) = message_update.emission_slot {
832 let mut serialized_emission_slot = Vec::new();
833 self.message_serializer
834 .slot_serializer
835 .serialize(&emission_slot, &mut serialized_emission_slot)
836 .expect(MESSAGE_SER_ERROR);
837 db.put_or_update_entry_value(
838 batch,
839 emission_slot_key!(serialized_message_id),
840 &serialized_emission_slot,
841 );
842 }
843
844 if let SetOrKeep::Set(emission_index) = message_update.emission_index {
846 let mut serialized_emission_index = Vec::new();
847 self.message_serializer
848 .u64_serializer
849 .serialize(&emission_index, &mut serialized_emission_index)
850 .expect(MESSAGE_SER_ERROR);
851 db.put_or_update_entry_value(
852 batch,
853 emission_index_key!(serialized_message_id),
854 &serialized_emission_index,
855 );
856 }
857
858 if let SetOrKeep::Set(sender) = message_update.sender {
860 let mut serialized_sender = Vec::new();
861 self.message_serializer
862 .address_serializer
863 .serialize(&sender, &mut serialized_sender)
864 .expect(MESSAGE_SER_ERROR);
865 db.put_or_update_entry_value(
866 batch,
867 sender_key!(serialized_message_id),
868 &serialized_sender,
869 );
870 }
871
872 if let SetOrKeep::Set(destination) = message_update.destination {
874 let mut serialized_destination = Vec::new();
875 self.message_serializer
876 .address_serializer
877 .serialize(&destination, &mut serialized_destination)
878 .expect(MESSAGE_SER_ERROR);
879 db.put_or_update_entry_value(
880 batch,
881 destination_key!(serialized_message_id),
882 &serialized_destination,
883 );
884 }
885
886 if let SetOrKeep::Set(function) = message_update.function {
888 let mut serialized_function = Vec::new();
889 self.message_serializer
890 .function_serializer
891 .serialize(&function, &mut serialized_function)
892 .expect(MESSAGE_SER_ERROR);
893 db.put_or_update_entry_value(
894 batch,
895 function_key!(serialized_message_id),
896 &serialized_function,
897 );
898 }
899
900 if let SetOrKeep::Set(max_gas) = message_update.max_gas {
902 let mut serialized_max_gas = Vec::new();
903 self.message_serializer
904 .u64_serializer
905 .serialize(&max_gas, &mut serialized_max_gas)
906 .expect(MESSAGE_SER_ERROR);
907 db.put_or_update_entry_value(
908 batch,
909 max_gas_key!(serialized_message_id),
910 &serialized_max_gas,
911 );
912 }
913
914 if let SetOrKeep::Set(fee) = message_update.fee {
916 let mut serialized_fee = Vec::new();
917 self.message_serializer
918 .amount_serializer
919 .serialize(&fee, &mut serialized_fee)
920 .expect(MESSAGE_SER_ERROR);
921 db.put_or_update_entry_value(batch, fee_key!(serialized_message_id), &serialized_fee);
922 }
923
924 if let SetOrKeep::Set(coins) = message_update.coins {
926 let mut serialized_coins = Vec::new();
927 self.message_serializer
928 .amount_serializer
929 .serialize(&coins, &mut serialized_coins)
930 .expect(MESSAGE_SER_ERROR);
931 db.put_or_update_entry_value(
932 batch,
933 coins_key!(serialized_message_id),
934 &serialized_coins,
935 );
936 }
937
938 if let SetOrKeep::Set(validity_start) = message_update.validity_start {
940 let mut serialized_validity_start = Vec::new();
941 self.message_serializer
942 .slot_serializer
943 .serialize(&validity_start, &mut serialized_validity_start)
944 .expect(MESSAGE_SER_ERROR);
945 db.put_or_update_entry_value(
946 batch,
947 validity_start_key!(serialized_message_id),
948 &serialized_validity_start,
949 );
950 }
951
952 if let SetOrKeep::Set(validity_end) = message_update.validity_end {
954 let mut serialized_validity_end = Vec::new();
955 self.message_serializer
956 .slot_serializer
957 .serialize(&validity_end, &mut serialized_validity_end)
958 .expect(MESSAGE_SER_ERROR);
959 db.put_or_update_entry_value(
960 batch,
961 validity_end_key!(serialized_message_id),
962 &serialized_validity_end,
963 );
964 }
965
966 if let SetOrKeep::Set(params) = message_update.function_params {
968 let mut serialized_function_params = Vec::new();
969 self.message_serializer
970 .function_params_serializer
971 .serialize(¶ms, &mut serialized_function_params)
972 .expect(MESSAGE_SER_ERROR);
973 db.put_or_update_entry_value(
974 batch,
975 function_params_key!(serialized_message_id),
976 &serialized_function_params,
977 );
978 }
979
980 if let SetOrKeep::Set(trigger) = message_update.trigger {
982 let mut serialized_trigger = Vec::new();
983 self.message_serializer
984 .trigger_serializer
985 .serialize(&trigger, &mut serialized_trigger)
986 .expect(MESSAGE_SER_ERROR);
987 db.put_or_update_entry_value(
988 batch,
989 trigger_key!(serialized_message_id),
990 &serialized_trigger,
991 );
992 }
993
994 if let SetOrKeep::Set(can_be_executed) = message_update.can_be_executed {
996 let mut serialized_can_be_executed = Vec::new();
997 self.message_serializer
998 .bool_serializer
999 .serialize(&can_be_executed, &mut serialized_can_be_executed)
1000 .expect(MESSAGE_SER_ERROR);
1001 db.put_or_update_entry_value(
1002 batch,
1003 can_be_executed_key!(serialized_message_id),
1004 &serialized_can_be_executed,
1005 );
1006 }
1007 }
1008
1009 fn delete_entry(&self, message_id: &AsyncMessageId, batch: &mut DBBatch) {
1014 let db = self.db.read();
1015 let mut serialized_message_id = Vec::new();
1016 self.message_id_serializer
1017 .serialize(message_id, &mut serialized_message_id)
1018 .expect(MESSAGE_ID_SER_ERROR);
1019
1020 db.delete_key(batch, emission_slot_key!(serialized_message_id));
1021 db.delete_key(batch, emission_index_key!(serialized_message_id));
1022 db.delete_key(batch, sender_key!(serialized_message_id));
1023 db.delete_key(batch, destination_key!(serialized_message_id));
1024 db.delete_key(batch, function_key!(serialized_message_id));
1025 db.delete_key(batch, max_gas_key!(serialized_message_id));
1026 db.delete_key(batch, fee_key!(serialized_message_id));
1027 db.delete_key(batch, coins_key!(serialized_message_id));
1028 db.delete_key(batch, validity_start_key!(serialized_message_id));
1029 db.delete_key(batch, validity_end_key!(serialized_message_id));
1030 db.delete_key(batch, function_params_key!(serialized_message_id));
1031 db.delete_key(batch, trigger_key!(serialized_message_id));
1032 db.delete_key(batch, can_be_executed_key!(serialized_message_id));
1033 }
1034}
1035
1036#[cfg(test)]
1037mod tests {
1038 use std::str::FromStr;
1039 use std::sync::Arc;
1040
1041 use massa_db_exports::{MassaDBConfig, MassaDBController};
1042 use massa_models::config::{
1043 MAX_ASYNC_POOL_LENGTH, MAX_DATASTORE_KEY_LENGTH, MAX_FUNCTION_NAME_LENGTH,
1044 MAX_PARAMETERS_SIZE, THREAD_COUNT,
1045 };
1046 use massa_models::{address::Address, amount::Amount, slot::Slot};
1047
1048 use massa_models::async_msg::AsyncMessageTrigger;
1049
1050 use massa_db_worker::MassaDB;
1051 use parking_lot::RwLock;
1052 use tempfile::tempdir;
1053
1054 use super::*;
1055
1056 fn dump_column(
1057 db: Arc<RwLock<Box<dyn MassaDBController>>>,
1058 column: &str,
1059 ) -> BTreeMap<Vec<u8>, Vec<u8>> {
1060 db.read()
1061 .iterator_cf_for_full_db_traversal(column, MassaIteratorMode::Start)
1062 .collect()
1063 }
1064
1065 fn create_message() -> AsyncMessage {
1066 AsyncMessage::new(
1067 Slot::new(1, 0),
1068 0,
1069 Address::from_str("AU12dG5xP1RDEB5ocdHkymNVvvSJmUL9BgHwCksDowqmGWxfpm93x").unwrap(),
1070 Address::from_str("AU12htxRWiEm8jDJpJptr6cwEhWNcCSFWstN1MLSa96DDkVM9Y42G").unwrap(),
1071 String::from("test"),
1072 10000000,
1073 Amount::from_str("1").unwrap(),
1074 Amount::from_str("1").unwrap(),
1075 Slot::new(2, 0),
1076 Slot::new(3, 0),
1077 vec![1, 2, 3, 4],
1078 Some(AsyncMessageTrigger {
1079 address: Address::from_str("AU12dG5xP1RDEB5ocdHkymNVvvSJmUL9BgHwCksDowqmGWxfpm93x")
1080 .unwrap(),
1081 datastore_key: Some(vec![1, 2, 3, 4]),
1082 }),
1083 None,
1084 )
1085 }
1086
1087 #[test]
1088 fn test_pool_ser_deser_empty() {
1089 let config = AsyncPoolConfig::default();
1090 let temp_dir = tempdir().expect("Unable to create a temp folder");
1091 let db_config = MassaDBConfig {
1092 path: temp_dir.path().to_path_buf(),
1093 max_history_length: 100,
1094 max_final_state_elements_size: 100,
1095 max_versioning_elements_size: 100,
1096 thread_count: THREAD_COUNT,
1097 max_ledger_backups: 100,
1098 enable_metrics: false,
1099 };
1100 let db: ShareableMassaDBController = Arc::new(RwLock::new(
1101 Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>,
1102 ));
1103 let pool = AsyncPool::new(config, db);
1104
1105 let mut serialized = Vec::new();
1106 let serializer = AsyncPoolSerializer::new();
1107 let deserializer = AsyncPoolDeserializer::new(
1108 THREAD_COUNT,
1109 MAX_ASYNC_POOL_LENGTH,
1110 MAX_FUNCTION_NAME_LENGTH,
1111 MAX_PARAMETERS_SIZE as u64,
1112 MAX_DATASTORE_KEY_LENGTH as u32,
1113 );
1114
1115 let message_ids = vec![];
1116 let to_ser_ = pool.fetch_messages(&message_ids);
1117 let to_ser = to_ser_
1118 .iter()
1119 .map(|(k, v)| {
1120 (
1121 *k,
1122 v.clone().unwrap_or_else(|| {
1123 panic!(
1124 "message_id {:?} should have been inserted in the pool above",
1125 k
1126 )
1127 }),
1128 )
1129 })
1130 .collect();
1131 serializer.serialize(&to_ser, &mut serialized).unwrap();
1132
1133 let (rest, changes_deser) = deserializer
1134 .deserialize::<DeserializeError>(&serialized)
1135 .unwrap();
1136 assert!(rest.is_empty());
1137 assert_eq!(to_ser, changes_deser);
1138 }
1139
1140 #[test]
1141 fn test_pool_ser_deser() {
1142 let config = AsyncPoolConfig::default();
1143 let temp_dir = tempdir().expect("Unable to create a temp folder");
1144 let db_config = MassaDBConfig {
1145 path: temp_dir.path().to_path_buf(),
1146 max_history_length: 100,
1147 max_final_state_elements_size: 100,
1148 max_versioning_elements_size: 100,
1149 thread_count: THREAD_COUNT,
1150 max_ledger_backups: 100,
1151 enable_metrics: false,
1152 };
1153 let db: ShareableMassaDBController = Arc::new(RwLock::new(
1154 Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>,
1155 ));
1156 let mut pool = AsyncPool::new(config, db);
1157
1158 let mut serialized = Vec::new();
1159 let serializer = AsyncPoolSerializer::new();
1160 let deserializer = AsyncPoolDeserializer::new(
1161 THREAD_COUNT,
1162 MAX_ASYNC_POOL_LENGTH,
1163 MAX_FUNCTION_NAME_LENGTH,
1164 MAX_PARAMETERS_SIZE as u64,
1165 MAX_DATASTORE_KEY_LENGTH as u32,
1166 );
1167
1168 let message = create_message();
1169 let message_id = message.compute_id();
1170 let mut message2 = message.clone();
1171 message2.emission_index += 1; message2.function = "test2".to_string();
1173 let message2_id = message2.compute_id();
1174 assert_ne!(message_id, message2_id);
1175
1176 let mut batch = DBBatch::new();
1177 pool.put_entry(&message.compute_id(), message.clone(), &mut batch);
1178 pool.put_entry(&message2.compute_id(), message2.clone(), &mut batch);
1179 let versioning_batch = DBBatch::new();
1180 let slot_1 = Slot::new(1, 0);
1181 pool.db
1182 .write()
1183 .write_batch(batch, versioning_batch, Some(slot_1));
1184 pool.recompute_message_cache();
1185 let message_ids = vec![message_id, message2_id];
1186 let to_ser_ = pool.fetch_messages(&message_ids);
1187 let to_ser = to_ser_
1188 .iter()
1189 .map(|(k, v)| {
1190 (
1191 *k,
1192 v.clone().unwrap_or_else(|| {
1193 panic!(
1194 "message_id {:?} should have been inserted in the pool above",
1195 k
1196 )
1197 }),
1198 )
1199 })
1200 .collect();
1201 serializer.serialize(&to_ser, &mut serialized).unwrap();
1202 assert_eq!(to_ser.len(), 2);
1203
1204 let (rest, changes_deser) = deserializer
1205 .deserialize::<DeserializeError>(&serialized)
1206 .unwrap();
1207 assert!(rest.is_empty());
1208 assert_eq!(to_ser, changes_deser);
1209 }
1210
1211 #[test]
1212 fn test_pool_ser_deser_too_high() {
1213 let config = AsyncPoolConfig::default();
1216 let temp_dir = tempdir().expect("Unable to create a temp folder");
1217 let db_config = MassaDBConfig {
1218 path: temp_dir.path().to_path_buf(),
1219 max_history_length: 100,
1220 max_final_state_elements_size: 100,
1221 max_versioning_elements_size: 100,
1222 thread_count: THREAD_COUNT,
1223 max_ledger_backups: 100,
1224 enable_metrics: false,
1225 };
1226 let db: ShareableMassaDBController = Arc::new(RwLock::new(
1227 Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>,
1228 ));
1229 let mut pool = AsyncPool::new(config, db);
1230
1231 let mut serialized = Vec::new();
1232 let serializer = AsyncPoolSerializer::new();
1233 let deserializer = AsyncPoolDeserializer::new(
1234 THREAD_COUNT,
1235 1,
1236 MAX_FUNCTION_NAME_LENGTH,
1237 MAX_PARAMETERS_SIZE as u64,
1238 MAX_DATASTORE_KEY_LENGTH as u32,
1239 );
1240
1241 let message = create_message();
1242 let message_id = message.compute_id();
1243 let mut message2 = message.clone();
1244 message2.emission_index += 1; message2.function = "test2".to_string();
1246 let message2_id = message2.compute_id();
1247 let mut batch = DBBatch::new();
1248 pool.put_entry(&message.compute_id(), message.clone(), &mut batch);
1249 pool.put_entry(&message2.compute_id(), message2.clone(), &mut batch);
1250 let versioning_batch = DBBatch::new();
1251 let slot_1 = Slot::new(1, 0);
1252 pool.db
1253 .write()
1254 .write_batch(batch, versioning_batch, Some(slot_1));
1255 pool.recompute_message_cache();
1256
1257 let message_ids = vec![message_id, message2_id];
1258 let to_ser_ = pool.fetch_messages(&message_ids);
1259 let to_ser = to_ser_
1260 .iter()
1261 .map(|(k, v)| {
1262 (
1263 *k,
1264 v.clone().unwrap_or_else(|| {
1265 panic!(
1266 "message_id {:?} should have been inserted in the pool above",
1267 k
1268 )
1269 }),
1270 )
1271 })
1272 .collect();
1273 serializer.serialize(&to_ser, &mut serialized).unwrap();
1274 assert_eq!(to_ser.len(), 2);
1275
1276 let res = deserializer.deserialize::<DeserializeError>(&serialized);
1277 assert!(res.is_err());
1278 }
1279
1280 #[test]
1281 fn test_pool_entry() {
1282 let config = AsyncPoolConfig::default();
1285 let temp_dir = tempdir().expect("Unable to create a temp folder");
1286 let db_config = MassaDBConfig {
1287 path: temp_dir.path().to_path_buf(),
1288 max_history_length: 100,
1289 max_final_state_elements_size: 100,
1290 max_versioning_elements_size: 100,
1291 thread_count: THREAD_COUNT,
1292 max_ledger_backups: 100,
1293 enable_metrics: false,
1294 };
1295 let db: ShareableMassaDBController = Arc::new(RwLock::new(
1296 Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>,
1297 ));
1298 let mut pool = AsyncPool::new(config, db);
1299
1300 let message = create_message();
1301 let message_id = message.compute_id();
1302 let mut message2 = message.clone();
1303 message2.emission_index += 1; message2.function = "test2".to_string();
1305 let message2_id = message2.compute_id();
1306 let mut batch = DBBatch::new();
1307 pool.put_entry(&message_id, message.clone(), &mut batch);
1308 pool.put_entry(&message2_id, message2.clone(), &mut batch);
1309
1310 let versioning_batch = DBBatch::new();
1311 let slot_1 = Slot::new(1, 0);
1312 pool.db
1313 .write()
1314 .write_batch(batch, versioning_batch, Some(slot_1));
1315 pool.recompute_message_cache();
1316
1317 let content = dump_column(pool.db.clone(), "state");
1318 assert_eq!(content.len(), 26); let mut batch2 = DBBatch::new();
1321 pool.delete_entry(&message_id, &mut batch2);
1322 let message_update = AsyncMessageUpdate {
1323 function: SetOrKeep::Set("test0".to_string()),
1324 ..Default::default()
1325 };
1326 pool.update_entry(&message2_id, message_update, &mut batch2);
1327
1328 let versioning_batch2 = DBBatch::new();
1329 let slot_2 = Slot::new(2, 0);
1330 pool.db
1331 .write()
1332 .write_batch(batch2, versioning_batch2, Some(slot_2));
1333 pool.recompute_message_cache();
1334
1335 let content = dump_column(pool.db.clone(), "state");
1336 assert_eq!(content.len(), 13);
1337 }
1338
1339 #[test]
1340 fn test_pool_cache_grow() {
1341 let config = AsyncPoolConfig::default();
1345 let temp_dir = tempdir().expect("Unable to create a temp folder");
1346 let db_config = MassaDBConfig {
1347 path: temp_dir.path().to_path_buf(),
1348 max_history_length: 100,
1349 max_final_state_elements_size: 100,
1350 max_versioning_elements_size: 100,
1351 thread_count: THREAD_COUNT,
1352 max_ledger_backups: 100,
1353 enable_metrics: false,
1354 };
1355 let db: ShareableMassaDBController = Arc::new(RwLock::new(
1356 Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>,
1357 ));
1358 let mut pool = AsyncPool::new(config, db);
1359
1360 assert!(pool.message_cache.is_empty());
1361
1362 let message = create_message();
1363 let message_id = message.compute_id();
1364
1365 let mut changes = AsyncPoolChanges::default();
1366 changes
1367 .0
1368 .insert(message_id, SetUpdateOrDelete::Set(message.clone()));
1369
1370 const EXPECT_CACHE_COUNT: u64 = 100;
1371 for i in 0..EXPECT_CACHE_COUNT {
1372 let mut message2 = message.clone();
1373 message2.fee = Amount::from_raw(i);
1374 assert_ne!(message.compute_id(), message2.compute_id());
1375
1376 changes
1377 .0
1378 .insert(message2.compute_id(), SetUpdateOrDelete::Set(message2));
1379 }
1380
1381 let mut batch = DBBatch::new();
1382 pool.apply_changes_to_batch(&changes, &mut batch);
1383 assert_eq!(pool.message_cache.len() as u64, EXPECT_CACHE_COUNT + 1);
1384
1385 pool.reset();
1386 assert!(pool.message_cache.is_empty());
1387 }
1388
1389 #[test]
1390 fn test_pool_recompute_cache() {
1391 let config = AsyncPoolConfig::default();
1396 let temp_dir = tempdir().expect("Unable to create a temp folder");
1397 let db_config = MassaDBConfig {
1398 path: temp_dir.path().to_path_buf(),
1399 max_history_length: 100,
1400 max_final_state_elements_size: 100,
1401 max_versioning_elements_size: 100,
1402 thread_count: THREAD_COUNT,
1403 max_ledger_backups: 100,
1404 enable_metrics: false,
1405 };
1406 let db: ShareableMassaDBController = Arc::new(RwLock::new(Box::new(MassaDB::new(
1407 db_config.clone(),
1408 ))
1409 as Box<dyn MassaDBController + 'static>));
1410 let mut pool = AsyncPool::new(config.clone(), db);
1411
1412 assert!(pool.message_cache.is_empty());
1413
1414 let message = create_message();
1415 let message_id = message.compute_id();
1416
1417 let mut changes = AsyncPoolChanges::default();
1418 changes
1419 .0
1420 .insert(message_id, SetUpdateOrDelete::Set(message.clone()));
1421
1422 const EXPECT_CACHE_COUNT: u64 = 100;
1423 for i in 0..EXPECT_CACHE_COUNT {
1424 let mut message2 = message.clone();
1425 message2.fee = Amount::from_raw(i);
1426 assert_ne!(message.compute_id(), message2.compute_id());
1427
1428 changes
1429 .0
1430 .insert(message2.compute_id(), SetUpdateOrDelete::Set(message2));
1431 }
1432
1433 let mut batch = DBBatch::new();
1434 pool.apply_changes_to_batch(&changes, &mut batch);
1435 assert_eq!(pool.message_cache.len() as u64, EXPECT_CACHE_COUNT + 1);
1436
1437 let message_cache1 = pool.message_cache.clone();
1438
1439 let versioning_batch = DBBatch::new();
1440 let slot_1 = Slot::new(1, 0);
1441 pool.db
1442 .write()
1443 .write_batch(batch, versioning_batch, Some(slot_1));
1444
1445 drop(pool);
1446
1447 let db2: ShareableMassaDBController = Arc::new(RwLock::new(
1448 Box::new(MassaDB::new(db_config)) as Box<dyn MassaDBController + 'static>,
1449 ));
1450 let mut pool2 = AsyncPool::new(config, db2);
1451
1452 pool2.recompute_message_cache();
1453
1454 assert_eq!(pool2.message_cache, message_cache1);
1455 }
1456}