1use std::cmp::max;
3use std::collections::{BTreeMap, BTreeSet};
4use std::path::Path;
5use num_enum::IntoPrimitive;
7use rocksdb::{IteratorMode, Options, WriteBatch, DB};
8use tracing::{debug, warn};
9use crate::rocksdb_operator::counter_merge;
11use crate::ser_deser::{
12 SCOutputEventDeserializer, SCOutputEventDeserializerArgs, SCOutputEventSerializer,
13};
14use massa_models::address::Address;
15use massa_models::error::ModelsError;
16use massa_models::execution::EventFilter;
17use massa_models::operation::{OperationId, OperationIdSerializer};
18use massa_models::output_event::SCOutputEvent;
19use massa_models::slot::Slot;
20use massa_serialization::{DeserializeError, Deserializer, Serializer};
21
22const OPEN_ERROR: &str = "critical: rocksdb open operation failed";
23const DESTROY_ERROR: &str = "critical: rocksdb delete operation failed";
25const CRUD_ERROR: &str = "critical: rocksdb crud operation failed";
26const EVENT_DESER_ERROR: &str = "critical: event deserialization failed";
27const OPERATION_ID_DESER_ERROR: &str = "critical: deserialization failed for op id in rocksdb";
28const COUNTER_ERROR: &str = "critical: cannot get counter";
29const COUNTER_KEY_CREATION_ERROR: &str = "critical: cannot create counter key";
30
31#[allow(dead_code)]
32#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, IntoPrimitive)]
34#[repr(u8)]
35enum KeyIndent {
36 Counter = 0,
37 Event,
38 EmitterAddress,
39 OriginalCallerAddress,
40 OriginalOperationId,
41 IsError,
42 IsFinal,
43}
44
45enum KeyBuilderType<'a> {
47 Slot(&'a Slot),
48 Event(&'a Slot, u64),
49 Address(&'a Address),
50 OperationId(&'a OperationId),
51 Bool(bool),
52 None,
53}
54
55enum KeyKind {
56 Regular,
57 Prefix,
58 Counter,
59}
60
61struct DbKeyBuilder {
62 op_id_ser: OperationIdSerializer,
64}
65
66impl DbKeyBuilder {
67 fn new() -> Self {
68 Self {
69 op_id_ser: OperationIdSerializer::new(),
70 }
71 }
72
73 fn key(&self, indent: &KeyIndent, key_type: KeyBuilderType, key_kind: &KeyKind) -> Vec<u8> {
79 let mut key_base = if matches!(key_kind, KeyKind::Counter) {
80 vec![u8::from(KeyIndent::Counter), u8::from(*indent)]
81 } else {
82 vec![u8::from(*indent)]
83 };
84
85 match key_type {
86 KeyBuilderType::Slot(slot) => {
87 key_base.extend(slot.to_bytes_key());
88 }
89 KeyBuilderType::Event(slot, index) => {
90 key_base.extend(slot.to_bytes_key());
91 key_base.extend(index.to_be_bytes());
92 }
93 KeyBuilderType::Address(addr) => {
94 let addr_bytes = addr.to_prefixed_bytes();
95 let addr_bytes_len = addr_bytes.len();
96 key_base.extend(addr_bytes);
97 key_base.push(addr_bytes_len as u8);
98 }
99 KeyBuilderType::OperationId(op_id) => {
100 let mut buffer = Vec::new();
101 self.op_id_ser
102 .serialize(op_id, &mut buffer)
103 .expect(OPERATION_ID_DESER_ERROR);
104 key_base.extend(&buffer);
105 key_base.extend(u32::to_be_bytes(buffer.len() as u32));
106 }
107 KeyBuilderType::Bool(value) => {
108 key_base.push(u8::from(value));
109 }
110 KeyBuilderType::None => {}
111 }
112
113 key_base
114 }
115
116 fn key_from_event(
118 &self,
119 event: &SCOutputEvent,
120 indent: &KeyIndent,
121 key_kind: &KeyKind,
122 ) -> Option<Vec<u8>> {
123 let key = match indent {
133 KeyIndent::Event => {
134 let item = KeyBuilderType::Event(&event.context.slot, event.context.index_in_slot);
135 Some(self.key(indent, item, key_kind))
136 }
137 KeyIndent::EmitterAddress => {
138 if let Some(addr) = event.context.call_stack.back() {
139 let item = KeyBuilderType::Address(addr);
140 let mut key = self.key(indent, item, key_kind);
141 let item =
142 KeyBuilderType::Event(&event.context.slot, event.context.index_in_slot);
143 if matches!(key_kind, KeyKind::Regular) {
144 key.extend(self.key(&KeyIndent::Event, item, &KeyKind::Regular));
145 }
146 Some(key)
147 } else {
148 None
149 }
150 }
151 KeyIndent::OriginalCallerAddress => {
152 if let Some(addr) = event.context.call_stack.front() {
153 let item = KeyBuilderType::Address(addr);
154 let mut key = self.key(indent, item, key_kind);
155 let item =
156 KeyBuilderType::Event(&event.context.slot, event.context.index_in_slot);
157 if matches!(key_kind, KeyKind::Regular) {
158 key.extend(self.key(&KeyIndent::Event, item, &KeyKind::Regular));
159 }
160 Some(key)
161 } else {
162 None
163 }
164 }
165 KeyIndent::OriginalOperationId => {
166 if let Some(op_id) = event.context.origin_operation_id.as_ref() {
167 let item = KeyBuilderType::OperationId(op_id);
168 let mut key = self.key(indent, item, key_kind);
169 let item =
170 KeyBuilderType::Event(&event.context.slot, event.context.index_in_slot);
171 if matches!(key_kind, KeyKind::Regular) {
172 key.extend(self.key(&KeyIndent::Event, item, &KeyKind::Regular));
173 }
174 Some(key)
175 } else {
176 None
177 }
178 }
179 KeyIndent::IsError => {
180 let item = KeyBuilderType::Bool(event.context.is_error);
181 let mut key = self.key(indent, item, key_kind);
182 let item = KeyBuilderType::Event(&event.context.slot, event.context.index_in_slot);
183 if matches!(key_kind, KeyKind::Regular) {
184 key.extend(self.key(&KeyIndent::Event, item, &KeyKind::Regular));
185 }
186 Some(key)
187 }
188 _ => unreachable!(),
189 };
190
191 key
192 }
193
194 fn prefix_key_from_indent(&self, indent: &KeyIndent) -> Vec<u8> {
196 self.key(indent, KeyBuilderType::None, &KeyKind::Regular)
197 }
198
199 fn prefix_key_from_filter_item(&self, filter_item: &FilterItem, indent: &KeyIndent) -> Vec<u8> {
201 match (indent, filter_item) {
202 (KeyIndent::Event, FilterItem::SlotStartEnd(_start, _end)) => {
203 unimplemented!()
204 }
205 (KeyIndent::Event, FilterItem::SlotStart(start)) => {
206 self.key(indent, KeyBuilderType::Slot(start), &KeyKind::Prefix)
207 }
208 (KeyIndent::Event, FilterItem::SlotEnd(end)) => {
209 self.key(indent, KeyBuilderType::Slot(end), &KeyKind::Prefix)
210 }
211 (KeyIndent::EmitterAddress, FilterItem::EmitterAddress(addr)) => {
212 self.key(indent, KeyBuilderType::Address(addr), &KeyKind::Prefix)
213 }
214 (KeyIndent::OriginalCallerAddress, FilterItem::OriginalCallerAddress(addr)) => {
215 self.key(indent, KeyBuilderType::Address(addr), &KeyKind::Prefix)
216 }
217 (KeyIndent::OriginalOperationId, FilterItem::OriginalOperationId(op_id)) => {
218 self.key(indent, KeyBuilderType::OperationId(op_id), &KeyKind::Prefix)
219 }
220 (KeyIndent::IsError, FilterItem::IsError(v)) => {
221 self.key(indent, KeyBuilderType::Bool(*v), &KeyKind::Prefix)
222 }
223 _ => {
224 unreachable!()
225 }
226 }
227 }
228
229 fn counter_key_from_filter_item(
231 &self,
232 filter_item: &FilterItem,
233 indent: &KeyIndent,
234 ) -> Vec<u8> {
235 match (indent, filter_item) {
237 (KeyIndent::Event, FilterItem::SlotStartEnd(_start, _end)) => {
238 unimplemented!()
239 }
240 (KeyIndent::Event, FilterItem::SlotStart(start)) => {
241 self.key(indent, KeyBuilderType::Slot(start), &KeyKind::Counter)
242 }
243 (KeyIndent::Event, FilterItem::SlotEnd(end)) => {
244 self.key(indent, KeyBuilderType::Slot(end), &KeyKind::Counter)
245 }
246 (KeyIndent::EmitterAddress, FilterItem::EmitterAddress(addr)) => {
247 self.key(indent, KeyBuilderType::Address(addr), &KeyKind::Counter)
248 }
249 (KeyIndent::OriginalCallerAddress, FilterItem::OriginalCallerAddress(addr)) => {
250 self.key(indent, KeyBuilderType::Address(addr), &KeyKind::Counter)
251 }
252 (KeyIndent::OriginalOperationId, FilterItem::OriginalOperationId(op_id)) => self.key(
253 indent,
254 KeyBuilderType::OperationId(op_id),
255 &KeyKind::Counter,
256 ),
257 (KeyIndent::IsError, FilterItem::IsError(v)) => {
258 self.key(indent, KeyBuilderType::Bool(*v), &KeyKind::Counter)
259 }
260 _ => {
261 unreachable!()
262 }
263 }
264 }
265}
266
267pub(crate) struct EventCache {
269 db: DB,
271 entry_count: usize,
274 max_entry_count: usize,
277 snip_amount: usize,
279 event_ser: SCOutputEventSerializer,
281 event_deser: SCOutputEventDeserializer,
283 key_builder: DbKeyBuilder,
285 first_slot: Slot,
287 last_slot: Slot,
289 thread_count: u8,
291 max_events_per_operation: u64,
293 max_operations_per_block: u64,
295 max_events_per_query: usize,
297}
298
299impl EventCache {
300 #[allow(clippy::too_many_arguments)]
301 pub fn new(
303 path: &Path,
304 max_entry_count: usize,
305 snip_amount: usize,
306 thread_count: u8,
307 max_recursive_call_depth: u16,
308 max_event_data_length: u64,
309 max_events_per_operation: u64,
310 max_operations_per_block: u64,
311 max_events_per_query: usize,
312 ) -> Self {
313 if path.exists() {
315 DB::destroy(&Options::default(), path).expect(DESTROY_ERROR);
316 }
317 let options = {
318 let mut opts = Options::default();
319 opts.create_if_missing(true);
320 opts.set_merge_operator_associative("counter merge operator", counter_merge);
321 opts.set_max_open_files(128);
322 opts
323 };
324 let db = DB::open(&options, path).expect(OPEN_ERROR);
325
326 let key_builder_2 = DbKeyBuilder::new();
327
328 Self {
329 db,
330 entry_count: 0,
331 max_entry_count,
332 snip_amount,
333 event_ser: SCOutputEventSerializer::new(),
334 event_deser: SCOutputEventDeserializer::new(SCOutputEventDeserializerArgs {
335 thread_count,
336 max_call_stack_length: max_recursive_call_depth,
337 max_event_data_length,
338 }),
339 key_builder: key_builder_2,
340 first_slot: Slot::new(0, 0),
341 last_slot: Slot::new(0, 0),
342 thread_count,
343 max_events_per_operation,
344 max_operations_per_block,
345 max_events_per_query,
346 }
347 }
348
349 fn insert_into_batch(&mut self, event: SCOutputEvent, batch: &mut WriteBatch) {
351 let mut event_buffer = Vec::new();
352 self.event_ser.serialize(&event, &mut event_buffer).unwrap();
353
354 batch.put(
355 self.key_builder
356 .key_from_event(&event, &KeyIndent::Event, &KeyKind::Regular)
357 .unwrap(),
358 event_buffer,
359 );
360
361 if let Some(key) =
362 self.key_builder
363 .key_from_event(&event, &KeyIndent::EmitterAddress, &KeyKind::Regular)
364 {
365 let key_counter = self.key_builder.key_from_event(
366 &event,
367 &KeyIndent::EmitterAddress,
368 &KeyKind::Counter,
369 );
370 batch.put(key, vec![]);
371 let key_counter = key_counter.expect(COUNTER_KEY_CREATION_ERROR);
372 batch.merge(key_counter, 1i64.to_be_bytes());
373 }
374
375 if let Some(key) = self.key_builder.key_from_event(
376 &event,
377 &KeyIndent::OriginalCallerAddress,
378 &KeyKind::Regular,
379 ) {
380 let key_counter = self.key_builder.key_from_event(
381 &event,
382 &KeyIndent::OriginalCallerAddress,
383 &KeyKind::Counter,
384 );
385 batch.put(key, vec![]);
386 let key_counter = key_counter.expect(COUNTER_KEY_CREATION_ERROR);
387 batch.merge(key_counter, 1i64.to_be_bytes());
388 }
389
390 if let Some(key) = self.key_builder.key_from_event(
391 &event,
392 &KeyIndent::OriginalOperationId,
393 &KeyKind::Regular,
394 ) {
395 let key_counter = self.key_builder.key_from_event(
396 &event,
397 &KeyIndent::OriginalOperationId,
398 &KeyKind::Counter,
399 );
400 batch.put(key, vec![]);
401 let key_counter = key_counter.expect(COUNTER_KEY_CREATION_ERROR);
402 batch.merge(key_counter, 1i64.to_be_bytes());
403 }
404
405 {
406 if let Some(key) =
407 self.key_builder
408 .key_from_event(&event, &KeyIndent::IsError, &KeyKind::Regular)
409 {
410 let key_counter =
411 self.key_builder
412 .key_from_event(&event, &KeyIndent::IsError, &KeyKind::Counter);
413 let key_counter = key_counter.expect(COUNTER_KEY_CREATION_ERROR);
414 batch.put(key, vec![]);
415 batch.merge(key_counter, 1i64.to_be_bytes());
416 }
417 }
418
419 self.last_slot = max(self.last_slot, event.context.slot);
422 }
423
424 #[allow(dead_code)]
425 pub fn insert(&mut self, event: SCOutputEvent) {
427 if self.entry_count >= self.max_entry_count {
428 self.snip(None);
429 }
430
431 let mut batch = WriteBatch::default();
432 self.insert_into_batch(event, &mut batch);
433 self.db.write(&batch).expect(CRUD_ERROR);
434
435 self.entry_count = self.entry_count.saturating_add(1);
438
439 debug!("(Event insert) entry_count is: {}", self.entry_count);
440 }
441
442 pub fn insert_multi_it(
444 &mut self,
445 events: impl ExactSizeIterator<Item = SCOutputEvent> + Clone,
446 ) {
447 let events_len = events.len();
448
449 if self.entry_count + events_len >= self.max_entry_count {
450 let snip_amount = max(self.snip_amount, events_len);
451 self.snip(Some(snip_amount));
452 }
453
454 let mut batch = WriteBatch::default();
455 for event in events {
456 self.insert_into_batch(event, &mut batch);
457 }
458 self.db.write(&batch).expect(CRUD_ERROR);
459 self.entry_count = self.entry_count.saturating_add(events_len);
462
463 debug!("(Events insert) entry_count is: {}", self.entry_count);
464 }
465
466 pub(crate) fn get_filtered_sc_output_events(
468 &self,
469 filter: &EventFilter,
470 ) -> (Vec<u64>, Vec<SCOutputEvent>) {
471 if filter.is_final == Some(false) {
474 return (vec![], vec![]);
475 }
476
477 let mut filter_items = from_event_filter(filter);
484
485 if filter_items.is_empty() {
486 filter_items.push((KeyIndent::Event, FilterItem::SlotStart(Slot::new(0, 0))));
489 }
490
491 let it = filter_items.iter().map(|(key_indent, filter_item)| {
492 let count = self
493 .filter_item_estimate_count(key_indent, filter_item)
494 .unwrap_or_else(|e| {
495 warn!(
496 "Could not estimate count for key indent: {:?} - filter_item: {:?}: {}",
497 key_indent, filter_item, e
498 );
499 self.max_entry_count as u64
500 });
501 ((count, key_indent), filter_item)
502 });
503
504 let map = BTreeMap::from_iter(it);
505 debug!("Filter items map: {:?}", map);
506
507 let mut query_counts = Vec::with_capacity(map.len());
510 let mut filter_res_prev = None;
511 for ((_counter, indent), filter_item) in map.iter() {
512 let mut filter_res = BTreeSet::new();
513 let query_count = self.filter_for(
514 indent,
515 filter_item,
516 &mut filter_res,
517 filter_res_prev.as_ref(),
518 );
519 query_counts.push(query_count);
520 filter_res_prev = Some(filter_res);
521 }
522
523 let multi_args = filter_res_prev
526 .unwrap()
527 .into_iter()
528 .take(self.max_events_per_query)
529 .collect::<Vec<Vec<u8>>>();
530
531 let res = self.db.multi_get(multi_args);
532 debug!(
533 "Filter will try to deserialize to SCOutputEvent {} values",
534 res.len()
535 );
536
537 let events = res
538 .into_iter()
539 .filter_map(|value| match value {
540 Ok(Some(value)) => self
541 .event_deser
542 .deserialize::<DeserializeError>(&value)
543 .map(|(_, event)| event)
544 .map_err(|e| warn!("Event cache: cannot deserialize event: {}", e))
545 .ok(),
546 Ok(None) => {
549 warn!("Event cache: index key points to a missing event, skipping");
550 None
551 }
552 Err(e) => {
553 warn!("Event cache: cannot read event: {}", e);
554 None
555 }
556 })
557 .collect::<Vec<SCOutputEvent>>();
558
559 (query_counts, events)
560 }
561
562 fn filter_for(
563 &self,
564 indent: &KeyIndent,
565 filter_item: &FilterItem,
566 result: &mut BTreeSet<Vec<u8>>,
567 seen: Option<&BTreeSet<Vec<u8>>>,
568 ) -> u64 {
569 let mut query_count: u64 = 0;
570
571 if *indent == KeyIndent::Event {
572 let opts = match filter_item {
573 FilterItem::SlotStart(_start) => {
574 let key_start = self
575 .key_builder
576 .prefix_key_from_filter_item(filter_item, indent);
577 let mut options = rocksdb::ReadOptions::default();
578 options.set_iterate_lower_bound(key_start);
579 options
580 }
581 FilterItem::SlotEnd(_end) => {
582 let key_end = self
583 .key_builder
584 .prefix_key_from_filter_item(filter_item, indent);
585 let mut options = rocksdb::ReadOptions::default();
586 options.set_iterate_upper_bound(key_end);
587 options
588 }
589 FilterItem::SlotStartEnd(start, end) => {
590 let key_start = self
591 .key_builder
592 .prefix_key_from_filter_item(&FilterItem::SlotStart(*start), indent);
593 let key_end = self
594 .key_builder
595 .prefix_key_from_filter_item(&FilterItem::SlotEnd(*end), indent);
596 let mut options = rocksdb::ReadOptions::default();
597 options.set_iterate_range(key_start..key_end);
598 options
599 }
600 _ => unreachable!(),
601 };
602
603 #[allow(clippy::manual_flatten)]
604 for kvb in self.db.iterator_opt(IteratorMode::Start, opts) {
605 if let Ok(kvb) = kvb {
606 if !kvb.0.starts_with(&[*indent as u8]) {
607 break;
609 }
610
611 let found = kvb.0.to_vec();
612 query_count = query_count.saturating_add(1);
613
614 if let Some(filter_set_seen) = seen {
615 if filter_set_seen.contains(&found) {
616 result.insert(found);
617 }
618
619 if filter_set_seen.len() == result.len() {
623 break;
624 }
625 } else {
626 result.insert(found);
627 }
628 }
629 }
630 } else {
631 let prefix_filter = match filter_item {
632 FilterItem::EmitterAddress(_addr) => self
633 .key_builder
634 .prefix_key_from_filter_item(filter_item, indent),
635 FilterItem::OriginalCallerAddress(_addr) => self
636 .key_builder
637 .prefix_key_from_filter_item(filter_item, indent),
638 FilterItem::OriginalOperationId(_op_id) => self
639 .key_builder
640 .prefix_key_from_filter_item(filter_item, indent),
641 FilterItem::IsError(_is_error) => self
642 .key_builder
643 .prefix_key_from_filter_item(filter_item, indent),
644 _ => unreachable!(),
645 };
646
647 #[allow(clippy::manual_flatten)]
648 for kvb in self.db.prefix_iterator(prefix_filter.as_slice()) {
649 if let Ok(kvb) = kvb {
650 if !kvb.0.starts_with(&[*indent as u8]) {
651 break;
653 }
654
655 if !kvb.0.starts_with(prefix_filter.as_slice()) {
656 break;
658 }
659
660 let found = kvb
661 .0
662 .strip_prefix(prefix_filter.as_slice())
663 .unwrap() .to_vec();
665
666 query_count = query_count.saturating_add(1);
667
668 if let Some(filter_set_seen) = seen {
669 if filter_set_seen.contains(&found) {
670 result.insert(found);
671 }
672
673 if filter_set_seen.len() == result.len() {
677 break;
678 }
679 } else {
680 result.insert(found);
681 }
682 }
683 }
684 }
685
686 query_count
687 }
688
689 fn filter_item_estimate_count(
691 &self,
692 key_indent: &KeyIndent,
693 filter_item: &FilterItem,
694 ) -> Result<u64, ModelsError> {
695 match filter_item {
696 FilterItem::SlotStart(start) => {
697 let diff = self
698 .last_slot
699 .slots_since(start, self.thread_count)
700 .unwrap_or(0); Ok(diff
704 .saturating_mul(self.max_events_per_operation)
705 .saturating_mul(self.max_operations_per_block))
706 }
707 FilterItem::SlotStartEnd(start, end) => {
708 let diff = end.slots_since(start, self.thread_count).unwrap_or(0);
710 Ok(diff
711 .saturating_mul(self.max_events_per_operation)
712 .saturating_mul(self.max_operations_per_block))
713 }
714 FilterItem::SlotEnd(end) => {
715 let diff = end
717 .slots_since(&self.first_slot, self.thread_count)
718 .unwrap_or(0);
719 Ok(diff
720 .saturating_mul(self.max_events_per_operation)
721 .saturating_mul(self.max_operations_per_block))
722 }
723 FilterItem::EmitterAddress(_addr) => {
724 let counter_key = self
725 .key_builder
726 .counter_key_from_filter_item(filter_item, key_indent);
727 let counter = self.db.get(counter_key).expect(COUNTER_ERROR);
728 let counter_value = counter
729 .map(|b| u64::from_be_bytes(b.try_into().unwrap()))
730 .unwrap_or(0);
731 Ok(counter_value)
732 }
733 FilterItem::OriginalCallerAddress(_addr) => {
734 let counter_key = self
735 .key_builder
736 .counter_key_from_filter_item(filter_item, key_indent);
737 let counter = self.db.get(counter_key).expect(COUNTER_ERROR);
738 let counter_value = counter
739 .map(|b| u64::from_be_bytes(b.try_into().unwrap()))
740 .unwrap_or(0);
741 Ok(counter_value)
742 }
743 FilterItem::OriginalOperationId(_op_id) => {
744 let counter_key = self
745 .key_builder
746 .counter_key_from_filter_item(filter_item, key_indent);
747 let counter = self.db.get(counter_key).expect(COUNTER_ERROR);
748 let counter_value = counter
749 .map(|b| u64::from_be_bytes(b.try_into().unwrap()))
750 .unwrap_or(0);
751 Ok(counter_value)
752 }
753 FilterItem::IsError(_is_error) => {
754 let counter_key = self
755 .key_builder
756 .counter_key_from_filter_item(filter_item, key_indent);
757 let counter = self.db.get(counter_key).expect(COUNTER_ERROR);
758 let counter_value = counter
759 .map(|b| u64::from_be_bytes(b.try_into().unwrap()))
760 .unwrap_or(0);
761 Ok(counter_value)
762 }
763 }
764 }
765
766 fn snip(&mut self, snip_amount: Option<usize>) {
768 let mut iter = self.db.iterator(IteratorMode::Start);
769 let mut batch = WriteBatch::default();
770 let mut snipped_count: usize = 0;
771 let snip_amount = snip_amount.unwrap_or(self.snip_amount);
772
773 let mut counter_keys = vec![];
774
775 while snipped_count < snip_amount {
776 let key_value = iter.next();
777 if key_value.is_none() {
778 break;
779 }
780 let kvb = key_value
781 .unwrap() .expect(EVENT_DESER_ERROR);
783
784 let key = kvb.0;
785 if !key.starts_with(&[u8::from(KeyIndent::Event)]) {
786 continue;
787 }
788
789 let (_, event) = self
790 .event_deser
791 .deserialize::<DeserializeError>(&kvb.1)
792 .unwrap();
793
794 if let Some(key) = self.key_builder.key_from_event(
796 &event,
797 &KeyIndent::EmitterAddress,
798 &KeyKind::Regular,
799 ) {
800 let key_counter = self
801 .key_builder
802 .key_from_event(&event, &KeyIndent::EmitterAddress, &KeyKind::Counter)
803 .expect(COUNTER_ERROR);
804 batch.delete(key);
805 counter_keys.push(key_counter.clone());
806 batch.merge(key_counter, (-1i64).to_be_bytes());
807 }
808 if let Some(key) = self.key_builder.key_from_event(
809 &event,
810 &KeyIndent::OriginalCallerAddress,
811 &KeyKind::Regular,
812 ) {
813 let key_counter = self
814 .key_builder
815 .key_from_event(&event, &KeyIndent::OriginalCallerAddress, &KeyKind::Counter)
816 .expect(COUNTER_ERROR);
817 batch.delete(key);
818 counter_keys.push(key_counter.clone());
819 batch.merge(key_counter, (-1i64).to_be_bytes());
820 }
821
822 if let Some(key) = self.key_builder.key_from_event(
823 &event,
824 &KeyIndent::OriginalOperationId,
825 &KeyKind::Regular,
826 ) {
827 let key_counter = self
828 .key_builder
829 .key_from_event(&event, &KeyIndent::OriginalOperationId, &KeyKind::Counter)
830 .expect(COUNTER_ERROR);
831 batch.delete(key);
832 counter_keys.push(key_counter.clone());
833 batch.merge(key_counter, (-1i64).to_be_bytes());
834 }
835 if let Some(key) =
836 self.key_builder
837 .key_from_event(&event, &KeyIndent::IsError, &KeyKind::Regular)
838 {
839 let key_counter = self
840 .key_builder
841 .key_from_event(&event, &KeyIndent::IsError, &KeyKind::Counter)
842 .expect(COUNTER_ERROR);
843 batch.delete(key);
844 counter_keys.push(key_counter.clone());
845 batch.merge(key_counter, (-1i64).to_be_bytes());
846 }
847
848 batch.delete(key);
849 snipped_count += 1;
850 }
851
852 self.db.write(&batch).expect(CRUD_ERROR);
854 self.entry_count = self.entry_count.saturating_sub(snipped_count);
855
856 let mut batch_counters = WriteBatch::default();
858 const U64_ZERO_BYTES: [u8; 8] = 0u64.to_be_bytes();
859 for (value, key) in self.db.multi_get(&counter_keys).iter().zip(counter_keys) {
860 if let Ok(Some(value)) = value {
861 if *value == U64_ZERO_BYTES.to_vec() {
862 batch_counters.delete(key);
863 }
864 }
865 }
866 self.db.write(&batch_counters).expect(CRUD_ERROR);
867
868 if self.entry_count == 0 {
870 self.first_slot = Slot::new(0, 0);
872 self.last_slot = Slot::new(0, 0);
873 } else {
874 let key_prefix = self.key_builder.prefix_key_from_indent(&KeyIndent::Event);
878 let mut it_slot = self.db.prefix_iterator(key_prefix);
879
880 let key_value = it_slot.next();
881 let kvb = key_value.unwrap().expect(EVENT_DESER_ERROR);
882
883 let (_, event) = self
884 .event_deser
885 .deserialize::<DeserializeError>(&kvb.1)
886 .unwrap();
887 self.first_slot = event.context.slot;
888 }
889 }
890}
891
892#[derive(Debug)]
894enum FilterItem {
895 SlotStart(Slot),
896 SlotStartEnd(Slot, Slot),
897 SlotEnd(Slot),
898 EmitterAddress(Address),
899 OriginalCallerAddress(Address),
900 OriginalOperationId(OperationId),
901 IsError(bool),
902}
903
904fn from_event_filter(event_filter: &EventFilter) -> Vec<(KeyIndent, FilterItem)> {
906 let mut filter_items = vec![];
907 if event_filter.start.is_some() && event_filter.end.is_some() {
908 let start = event_filter.start.unwrap();
909 let end = event_filter.end.unwrap();
910 filter_items.push((KeyIndent::Event, FilterItem::SlotStartEnd(start, end)));
911 } else if event_filter.start.is_some() {
912 let start = event_filter.start.unwrap();
913 filter_items.push((KeyIndent::Event, FilterItem::SlotStart(start)));
914 } else if event_filter.end.is_some() {
915 let end = event_filter.end.unwrap();
916 filter_items.push((KeyIndent::Event, FilterItem::SlotEnd(end)));
917 }
918
919 if let Some(addr) = event_filter.emitter_address {
920 filter_items.push((KeyIndent::EmitterAddress, FilterItem::EmitterAddress(addr)));
921 }
922
923 if let Some(addr) = event_filter.original_caller_address {
924 filter_items.push((
925 KeyIndent::OriginalCallerAddress,
926 FilterItem::OriginalCallerAddress(addr),
927 ));
928 }
929
930 if let Some(op_id) = event_filter.original_operation_id {
931 filter_items.push((
932 KeyIndent::OriginalOperationId,
933 FilterItem::OriginalOperationId(op_id),
934 ));
935 }
936
937 if let Some(is_error) = event_filter.is_error {
938 filter_items.push((KeyIndent::IsError, FilterItem::IsError(is_error)));
939 }
940
941 filter_items
942}
943
944#[cfg(test)]
945impl EventCache {
946 fn iter_all(
948 &self,
949 mode: Option<IteratorMode>,
950 ) -> impl Iterator<Item = (Box<[u8]>, Box<[u8]>)> + '_ {
951 self.db
952 .iterator(mode.unwrap_or(IteratorMode::Start))
953 .flatten()
954 }
955}
956
957#[cfg(test)]
958mod tests {
959 use super::*;
960 use std::collections::VecDeque;
962 use std::str::FromStr;
963 use more_asserts::assert_gt;
965 use rand::seq::SliceRandom;
966 use rand::thread_rng;
967 use serial_test::serial;
968 use tempfile::TempDir;
969 use massa_models::config::{
971 MAX_EVENT_DATA_SIZE, MAX_EVENT_PER_OPERATION, MAX_OPERATIONS_PER_BLOCK,
972 MAX_RECURSIVE_CALLS_DEPTH, THREAD_COUNT,
973 };
974 use massa_models::operation::OperationId;
975 use massa_models::output_event::EventExecutionContext;
976 use massa_models::slot::Slot;
977
978 fn setup() -> EventCache {
979 let tmp_path = TempDir::new().unwrap().path().to_path_buf();
980 EventCache::new(
981 &tmp_path,
982 1000,
983 300,
984 THREAD_COUNT,
985 MAX_RECURSIVE_CALLS_DEPTH,
986 MAX_EVENT_DATA_SIZE as u64,
987 MAX_EVENT_PER_OPERATION as u64,
988 MAX_OPERATIONS_PER_BLOCK as u64,
989 5000, )
991 }
992
993 #[test]
994 #[serial]
995 fn test_db_insert_order() {
996 let mut cache = setup();
999 let slot_1 = Slot::new(1, 0);
1000 let index_1_0 = 0;
1001 let event = SCOutputEvent {
1002 context: EventExecutionContext {
1003 slot: slot_1,
1004 block: None,
1005 read_only: false,
1006 index_in_slot: index_1_0,
1007 call_stack: Default::default(),
1008 origin_operation_id: None,
1009 is_final: true,
1010 is_error: false,
1011 deferred_call_id: None,
1012 async_msg_id: None,
1013 },
1014 data: "message foo bar".to_string(),
1015 };
1016
1017 let max_entry_count = cache.max_entry_count - 5;
1018 let mut events = (0..max_entry_count)
1019 .map(|i| {
1020 let mut event = event.clone();
1021 event.context.index_in_slot = i as u64;
1022 event
1023 })
1024 .collect::<Vec<SCOutputEvent>>();
1025
1026 let slot_2 = Slot::new(2, 0);
1027 let event_slot_2 = {
1028 let mut event = event.clone();
1029 event.context.slot = slot_2;
1030 event.context.index_in_slot = 0u64;
1031 event
1032 };
1033 let index_2_2 = 256u64;
1034 let event_slot_2_2 = {
1035 let mut event = event.clone();
1036 event.context.slot = slot_2;
1037 event.context.index_in_slot = index_2_2;
1038 event
1039 };
1040 events.push(event_slot_2.clone());
1041 events.push(event_slot_2_2.clone());
1042 events.shuffle(&mut thread_rng());
1044
1045 for event in events {
1046 cache.insert(event);
1047 }
1048
1049 let mut prev_slot = None;
1052 let mut prev_event_index = None;
1053 for kvb in cache.iter_all(None) {
1054 let bytes = kvb.0.iter().as_slice();
1055
1056 if bytes[0] != u8::from(KeyIndent::Event) {
1057 continue;
1058 }
1059
1060 let slot = Slot::from_bytes_key(&bytes[1..=9].try_into().unwrap());
1061 let event_index = u64::from_be_bytes(bytes[10..].try_into().unwrap());
1062 if prev_slot.is_some() && prev_event_index.is_some() {
1063 assert_gt!(
1064 (slot, event_index),
1065 (prev_slot.unwrap(), prev_event_index.unwrap())
1066 );
1067 } else {
1068 assert_eq!(slot, slot_1);
1069 assert_eq!(event_index, index_1_0);
1070 }
1071 prev_slot = Some(slot);
1072 prev_event_index = Some(event_index);
1073 }
1074
1075 assert_eq!(prev_slot, Some(slot_2));
1076 assert_eq!(prev_event_index, Some(index_2_2));
1077 }
1078
1079 #[test]
1080 #[serial]
1081 fn test_insert_more_than_max_entry() {
1082 let mut cache = setup();
1085 let event = SCOutputEvent {
1086 context: EventExecutionContext {
1087 slot: Slot::new(1, 0),
1088 block: None,
1089 read_only: false,
1090 index_in_slot: 0,
1091 call_stack: Default::default(),
1092 origin_operation_id: None,
1093 is_final: true,
1094 is_error: false,
1095 deferred_call_id: None,
1096 async_msg_id: None,
1097 },
1098 data: "message foo bar".to_string(),
1099 };
1100
1101 for count in 0..cache.max_entry_count {
1103 let mut event = event.clone();
1104 event.context.index_in_slot = count as u64;
1105 cache.insert(event.clone());
1106 }
1107 assert_eq!(cache.entry_count, cache.max_entry_count);
1108
1109 let mut event_last = event.clone();
1111 event_last.context.index_in_slot = u64::MAX;
1112 cache.insert(event_last);
1113 assert_eq!(
1114 cache.entry_count,
1115 cache.max_entry_count - cache.snip_amount + 1
1116 );
1117 dbg!(cache.entry_count);
1118 }
1119
1120 #[test]
1121 #[serial]
1122 fn test_insert_more_than_max_entry_2() {
1123 let mut cache = setup();
1126 let event = SCOutputEvent {
1127 context: EventExecutionContext {
1128 slot: Slot::new(1, 0),
1129 block: None,
1130 read_only: false,
1131 index_in_slot: 0,
1132 call_stack: Default::default(),
1133 origin_operation_id: None,
1134 is_final: true,
1135 is_error: false,
1136 deferred_call_id: None,
1137 async_msg_id: None,
1138 },
1139 data: "message foo bar".to_string(),
1140 };
1141
1142 let it = (0..cache.max_entry_count).map(|i| {
1143 let mut event = event.clone();
1144 event.context.index_in_slot = i as u64;
1145 event
1146 });
1147 cache.insert_multi_it(it);
1148
1149 assert_eq!(cache.entry_count, cache.max_entry_count);
1150
1151 let mut event_last = event.clone();
1153 event_last.context.index_in_slot = u64::MAX;
1154 cache.insert(event_last);
1155 assert_eq!(
1156 cache.entry_count,
1157 cache.max_entry_count - cache.snip_amount + 1
1158 );
1159 dbg!(cache.entry_count);
1160 }
1161
1162 #[test]
1163 #[serial]
1164 fn test_snip() {
1165 let mut cache = setup();
1168 cache.max_entry_count = 10;
1169
1170 let event = SCOutputEvent {
1171 context: EventExecutionContext {
1172 slot: Slot::new(1, 0),
1173 block: None,
1174 read_only: false,
1175 index_in_slot: 0,
1176 call_stack: Default::default(),
1177 origin_operation_id: None,
1178 is_final: true,
1179 is_error: false,
1180 deferred_call_id: None,
1181 async_msg_id: None,
1182 },
1183 data: "message foo bar".to_string(),
1184 };
1185
1186 let it = (0..cache.max_entry_count).map(|i| {
1187 let mut event = event.clone();
1188 event.context.index_in_slot = i as u64;
1189 event
1190 });
1191 cache.insert_multi_it(it);
1192
1193 assert_eq!(cache.entry_count, cache.max_entry_count);
1194
1195 cache.snip(Some(cache.entry_count));
1196
1197 assert_eq!(cache.entry_count, 0);
1198 assert_eq!(cache.iter_all(None).count(), 0);
1199 }
1200
1201 #[test]
1202 #[serial]
1203 fn test_counter_0() {
1204 let mut cache = setup();
1207 cache.max_entry_count = 10;
1208
1209 let dummy_addr =
1210 Address::from_str("AU12qePoXhNbYWE1jZuafqJong7bbq1jw3k89RgbMawbrdZpaasoA").unwrap();
1211 let emit_addr_1 =
1212 Address::from_str("AU122Em8qkqegdLb1eyH8rdkSCNEf7RZLeTJve4Q2inRPGiTJ2xNv").unwrap();
1213 let emit_addr_2 =
1214 Address::from_str("AU12WuVR1Td74q9eAbtYZUnk5jnRbUuUacyhQFwm217bV5v1mNqTZ").unwrap();
1215
1216 let event = SCOutputEvent {
1217 context: EventExecutionContext {
1218 slot: Slot::new(1, 0),
1219 block: None,
1220 read_only: false,
1221 index_in_slot: 0,
1222 call_stack: VecDeque::from(vec![dummy_addr, emit_addr_1]),
1223 origin_operation_id: None,
1224 is_final: true,
1225 is_error: false,
1226 deferred_call_id: None,
1227 async_msg_id: None,
1228 },
1229 data: "message foo bar".to_string(),
1230 };
1231
1232 let event_2 = {
1233 let mut evt = event.clone();
1234 evt.context.slot = Slot::new(2, 0);
1235 evt.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1236 evt
1237 };
1238
1239 cache.insert_multi_it([event, event_2].into_iter());
1240
1241 let key_counters = cache
1243 .key_builder
1244 .prefix_key_from_indent(&KeyIndent::Counter);
1245 let kvbs: Result<Vec<(_, _)>, _> = cache
1246 .db
1247 .prefix_iterator(key_counters)
1248 .take_while(|kvb| {
1249 kvb.as_ref()
1250 .unwrap()
1251 .0
1252 .starts_with(&[u8::from(KeyIndent::Counter)])
1253 })
1254 .collect();
1255 assert_eq!(kvbs.unwrap().len(), 4);
1262
1263 let key_counter_1 = cache.key_builder.counter_key_from_filter_item(
1264 &FilterItem::EmitterAddress(emit_addr_1),
1265 &KeyIndent::EmitterAddress,
1266 );
1267 let key_counter_2 = cache.key_builder.counter_key_from_filter_item(
1268 &FilterItem::EmitterAddress(emit_addr_2),
1269 &KeyIndent::EmitterAddress,
1270 );
1271
1272 let v1 = cache.db.get(key_counter_1.clone());
1273 let v2 = cache.db.get(key_counter_2.clone());
1274
1275 assert_eq!(v1, Ok(Some(1u64.to_be_bytes().to_vec())));
1277 assert_eq!(v2, Ok(Some(1u64.to_be_bytes().to_vec())));
1278
1279 cache.snip(Some(1));
1280
1281 let v1 = cache.db.get(key_counter_1);
1282 let v2 = cache.db.get(key_counter_2);
1283
1284 assert_eq!(v1, Ok(None)); assert_eq!(v2, Ok(Some(1u64.to_be_bytes().to_vec())));
1287 }
1288
1289 #[test]
1290 #[serial]
1291 fn test_event_filter() {
1292 let mut cache = setup();
1295 let slot_1 = Slot::new(1, 0);
1296 let index_1_0 = 0;
1297 let event = SCOutputEvent {
1298 context: EventExecutionContext {
1299 slot: slot_1,
1300 block: None,
1301 read_only: false,
1302 index_in_slot: index_1_0,
1303 call_stack: Default::default(),
1304 origin_operation_id: None,
1305 is_final: true,
1306 is_error: false,
1307 deferred_call_id: None,
1308 async_msg_id: None,
1309 },
1310 data: "message foo bar".to_string(),
1311 };
1312
1313 let mut events = (0..cache.max_entry_count - 5)
1314 .map(|i| {
1315 let mut event = event.clone();
1316 event.context.index_in_slot = i as u64;
1317 event
1318 })
1319 .collect::<Vec<SCOutputEvent>>();
1320
1321 let slot_2 = Slot::new(2, 0);
1322 let index_2_1 = 0u64;
1323 let event_slot_2 = {
1324 let mut event = event.clone();
1325 event.context.slot = slot_2;
1326 event.context.index_in_slot = index_2_1;
1327 event
1328 };
1329 let index_2_2 = 256u64;
1330 let event_slot_2_2 = {
1331 let mut event = event.clone();
1332 event.context.slot = slot_2;
1333 event.context.index_in_slot = index_2_2;
1334 event
1335 };
1336 events.push(event_slot_2.clone());
1337 events.push(event_slot_2_2.clone());
1338 events.shuffle(&mut thread_rng());
1340
1341 cache.insert_multi_it(events.into_iter());
1342
1343 let filter_1 = EventFilter {
1344 start: Some(Slot::new(2, 0)),
1345 ..Default::default()
1346 };
1347
1348 let (_, filtered_events_1) = cache.get_filtered_sc_output_events(&filter_1);
1349
1350 assert_eq!(filtered_events_1.len(), 2);
1351 assert_eq!(filtered_events_1[0].context.slot, slot_2);
1352 assert_eq!(filtered_events_1[0].context.index_in_slot, index_2_1);
1353 assert_eq!(filtered_events_1[1].context.slot, slot_2);
1354 assert_eq!(filtered_events_1[1].context.index_in_slot, index_2_2);
1355 }
1356
1357 #[test]
1358 #[serial]
1359 fn test_event_filter_2() {
1360 let mut cache = setup();
1363 cache.max_entry_count = 10;
1364
1365 let slot_1 = Slot::new(1, 0);
1366 let index_1_0 = 0;
1367 let op_id_1 =
1368 OperationId::from_str("O12n1vt8uTLh3H65J4TVuztaWfBh3oumjjVtRCkke7Ba5qWdXdjD").unwrap();
1369 let op_id_2 =
1370 OperationId::from_str("O1p5P691KF672fQ8tQougxzSERBwDKZF8FwtkifMSJbP14sEuGc").unwrap();
1371 let op_id_unknown =
1372 OperationId::from_str("O1kvXTfsnVbQcmDERkC89vqAd2xRTLCb3q5b2E5WaVPHwFd7Qth").unwrap();
1373
1374 let event = SCOutputEvent {
1375 context: EventExecutionContext {
1376 slot: slot_1,
1377 block: None,
1378 read_only: false,
1379 index_in_slot: index_1_0,
1380 call_stack: Default::default(),
1381 origin_operation_id: Some(op_id_1),
1382 is_final: true,
1383 is_error: false,
1384 deferred_call_id: None,
1385 async_msg_id: None,
1386 },
1387 data: "message foo bar".to_string(),
1388 };
1389
1390 let mut events = (0..cache.max_entry_count - 5)
1391 .map(|i| {
1392 let mut event = event.clone();
1393 event.context.index_in_slot = i as u64;
1394 event
1395 })
1396 .collect::<Vec<SCOutputEvent>>();
1397
1398 let slot_2 = Slot::new(2, 0);
1399 let index_2_1 = 0u64;
1400 let event_slot_2 = {
1401 let mut event = event.clone();
1402 event.context.slot = slot_2;
1403 event.context.index_in_slot = index_2_1;
1404 event.context.origin_operation_id = Some(op_id_2);
1405 event
1406 };
1407 let index_2_2 = 256u64;
1408 let event_slot_2_2 = {
1409 let mut event = event.clone();
1410 event.context.slot = slot_2;
1411 event.context.index_in_slot = index_2_2;
1412 event.context.origin_operation_id = Some(op_id_2);
1413 event
1414 };
1415 events.push(event_slot_2.clone());
1416 events.push(event_slot_2_2.clone());
1417 events.shuffle(&mut thread_rng());
1419 cache.insert_multi_it(events.into_iter());
1420
1421 let mut filter_1 = EventFilter {
1422 original_operation_id: Some(op_id_1),
1423 ..Default::default()
1424 };
1425
1426 let (_, filtered_events_1) = cache.get_filtered_sc_output_events(&filter_1);
1427
1428 assert_eq!(filtered_events_1.len(), cache.max_entry_count - 5);
1429 filtered_events_1.iter().enumerate().for_each(|(i, event)| {
1430 assert_eq!(event.context.slot, slot_1);
1431 assert_eq!(event.context.index_in_slot, i as u64);
1432 });
1433
1434 {
1435 filter_1.original_operation_id = Some(op_id_2);
1436 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1437 assert_eq!(filtered_events_2.len(), 2);
1438 filtered_events_2.iter().enumerate().for_each(|(i, event)| {
1439 assert_eq!(event.context.slot, slot_2);
1440 if i == 0 {
1441 assert_eq!(event.context.index_in_slot, i as u64);
1442 } else {
1443 assert_eq!(event.context.index_in_slot, 256u64);
1444 }
1445 });
1446 }
1447
1448 {
1449 filter_1.original_operation_id = Some(op_id_unknown);
1450 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1451 assert_eq!(filtered_events_2.len(), 0);
1452 }
1453 }
1454
1455 #[test]
1456 #[serial]
1457 fn test_event_filter_3() {
1458 let mut cache = setup();
1461 cache.max_entry_count = 10;
1462
1463 let slot_1 = Slot::new(1, 0);
1464 let index_1_0 = 0;
1465
1466 let dummy_addr =
1467 Address::from_str("AU12qePoXhNbYWE1jZuafqJong7bbq1jw3k89RgbMawbrdZpaasoA").unwrap();
1468 let emit_addr_1 =
1469 Address::from_str("AU122Em8qkqegdLb1eyH8rdkSCNEf7RZLeTJve4Q2inRPGiTJ2xNv").unwrap();
1470 let emit_addr_2 =
1471 Address::from_str("AU12WuVR1Td74q9eAbtYZUnk5jnRbUuUacyhQFwm217bV5v1mNqTZ").unwrap();
1472 let emit_addr_unknown =
1473 Address::from_str("AU1zLC4TFUiaKDg7quQyusMPQcHT4ykWVs3FsFpuhdNSmowUG2As").unwrap();
1474
1475 let event = SCOutputEvent {
1476 context: EventExecutionContext {
1477 slot: slot_1,
1478 block: None,
1479 read_only: false,
1480 index_in_slot: index_1_0,
1481 call_stack: Default::default(),
1482 origin_operation_id: None,
1483 is_final: true,
1484 is_error: false,
1485 deferred_call_id: None,
1486 async_msg_id: None,
1487 },
1488 data: "message foo bar".to_string(),
1489 };
1490
1491 let to_insert_count = cache.max_entry_count - 5;
1492 let threshold = to_insert_count / 2;
1493 let mut events = (0..cache.max_entry_count - 5)
1494 .map(|i| {
1495 let mut event = event.clone();
1496 event.context.index_in_slot = i as u64;
1497 if i < threshold {
1498 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_1]);
1499 } else {
1500 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1501 }
1502 event
1503 })
1504 .collect::<Vec<SCOutputEvent>>();
1505
1506 let slot_2 = Slot::new(2, 0);
1507 let index_2_1 = 0u64;
1508 let event_slot_2 = {
1509 let mut event = event.clone();
1510 event.context.slot = slot_2;
1511 event.context.index_in_slot = index_2_1;
1512 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1513 event
1514 };
1515 let index_2_2 = 256u64;
1516 let event_slot_2_2 = {
1517 let mut event = event.clone();
1518 event.context.slot = slot_2;
1519 event.context.index_in_slot = index_2_2;
1520 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1521 event
1522 };
1523 events.push(event_slot_2.clone());
1524 events.push(event_slot_2_2.clone());
1525 events.shuffle(&mut thread_rng());
1527
1528 cache.insert_multi_it(events.into_iter());
1529
1530 let mut filter_1 = EventFilter {
1531 emitter_address: Some(emit_addr_1),
1532 ..Default::default()
1533 };
1534
1535 let (_, filtered_events_1) = cache.get_filtered_sc_output_events(&filter_1);
1536
1537 assert_eq!(filtered_events_1.len(), threshold);
1538 filtered_events_1.iter().for_each(|event| {
1539 assert_eq!(event.context.slot, slot_1);
1540 assert_eq!(*event.context.call_stack.back().unwrap(), emit_addr_1)
1541 });
1542
1543 {
1545 filter_1.emitter_address = Some(emit_addr_2);
1546 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1547 assert_eq!(filtered_events_2.len(), threshold + 1 + 2);
1548 filtered_events_2.iter().for_each(|event| {
1549 assert_eq!(*event.context.call_stack.back().unwrap(), emit_addr_2)
1550 });
1551 }
1552 {
1554 filter_1.emitter_address = Some(dummy_addr);
1555 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1556 assert_eq!(filtered_events_2.len(), 0);
1557 }
1558 {
1560 filter_1.emitter_address = Some(emit_addr_unknown);
1561 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1562 assert_eq!(filtered_events_2.len(), 0);
1563 }
1564 }
1565
1566 #[test]
1567 #[serial]
1568 fn test_event_filter_4() {
1569 let mut cache = setup();
1572 cache.max_entry_count = 10;
1573
1574 let slot_1 = Slot::new(1, 0);
1575 let index_1_0 = 0;
1576
1577 let dummy_addr =
1578 Address::from_str("AU12qePoXhNbYWE1jZuafqJong7bbq1jw3k89RgbMawbrdZpaasoA").unwrap();
1579 let caller_addr_1 =
1580 Address::from_str("AU122Em8qkqegdLb1eyH8rdkSCNEf7RZLeTJve4Q2inRPGiTJ2xNv").unwrap();
1581 let caller_addr_2 =
1582 Address::from_str("AU12WuVR1Td74q9eAbtYZUnk5jnRbUuUacyhQFwm217bV5v1mNqTZ").unwrap();
1583 let caller_addr_unknown =
1584 Address::from_str("AU1zLC4TFUiaKDg7quQyusMPQcHT4ykWVs3FsFpuhdNSmowUG2As").unwrap();
1585
1586 let event = SCOutputEvent {
1587 context: EventExecutionContext {
1588 slot: slot_1,
1589 block: None,
1590 read_only: false,
1591 index_in_slot: index_1_0,
1592 call_stack: Default::default(),
1593 origin_operation_id: None,
1594 is_final: true,
1595 is_error: false,
1596 deferred_call_id: None,
1597 async_msg_id: None,
1598 },
1599 data: "message foo bar".to_string(),
1600 };
1601
1602 let to_insert_count = cache.max_entry_count - 5;
1603 let threshold = to_insert_count / 2;
1604 let mut events = (0..cache.max_entry_count - 5)
1605 .map(|i| {
1606 let mut event = event.clone();
1607 event.context.index_in_slot = i as u64;
1608 if i < threshold {
1609 event.context.call_stack = VecDeque::from(vec![caller_addr_1, dummy_addr]);
1610 } else {
1611 event.context.call_stack = VecDeque::from(vec![caller_addr_2, dummy_addr]);
1612 }
1613 event
1614 })
1615 .collect::<Vec<SCOutputEvent>>();
1616
1617 let slot_2 = Slot::new(2, 0);
1618 let index_2_1 = 0u64;
1619 let event_slot_2 = {
1620 let mut event = event.clone();
1621 event.context.slot = slot_2;
1622 event.context.index_in_slot = index_2_1;
1623 event.context.call_stack = VecDeque::from(vec![caller_addr_2, dummy_addr]);
1624 event
1625 };
1626 let index_2_2 = 256u64;
1627 let event_slot_2_2 = {
1628 let mut event = event.clone();
1629 event.context.slot = slot_2;
1630 event.context.index_in_slot = index_2_2;
1631 event.context.call_stack = VecDeque::from(vec![caller_addr_2, dummy_addr]);
1632 event
1633 };
1634 events.push(event_slot_2.clone());
1635 events.push(event_slot_2_2.clone());
1636 events.shuffle(&mut thread_rng());
1638 cache.insert_multi_it(events.into_iter());
1639
1640 let mut filter_1 = EventFilter {
1641 original_caller_address: Some(caller_addr_1),
1642 ..Default::default()
1643 };
1644
1645 let (_, filtered_events_1) = cache.get_filtered_sc_output_events(&filter_1);
1646
1647 assert_eq!(filtered_events_1.len(), threshold);
1648 filtered_events_1.iter().for_each(|event| {
1649 assert_eq!(event.context.slot, slot_1);
1650 assert_eq!(*event.context.call_stack.front().unwrap(), caller_addr_1);
1651 });
1652
1653 {
1654 filter_1.original_caller_address = Some(caller_addr_2);
1655 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1656 assert_eq!(filtered_events_2.len(), threshold + 1 + 2);
1657 filtered_events_2.iter().for_each(|event| {
1658 assert_eq!(*event.context.call_stack.front().unwrap(), caller_addr_2);
1659 });
1660 }
1661 {
1662 filter_1.original_caller_address = Some(dummy_addr);
1663 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1664 assert_eq!(filtered_events_2.len(), 0);
1665 }
1666 {
1667 filter_1.original_caller_address = Some(caller_addr_unknown);
1668 let (_, filtered_events_2) = cache.get_filtered_sc_output_events(&filter_1);
1669 assert_eq!(filtered_events_2.len(), 0);
1670 }
1671 }
1672
1673 #[test]
1674 #[serial]
1675 fn test_event_filter_5() {
1676 let mut cache = setup();
1679 cache.max_entry_count = 10;
1680
1681 let slot_1 = Slot::new(1, 0);
1682 let index_1_0 = 0;
1683
1684 let dummy_addr =
1685 Address::from_str("AU12qePoXhNbYWE1jZuafqJong7bbq1jw3k89RgbMawbrdZpaasoA").unwrap();
1686 let emit_addr_1 =
1687 Address::from_str("AU122Em8qkqegdLb1eyH8rdkSCNEf7RZLeTJve4Q2inRPGiTJ2xNv").unwrap();
1688 let emit_addr_2 =
1689 Address::from_str("AU12WuVR1Td74q9eAbtYZUnk5jnRbUuUacyhQFwm217bV5v1mNqTZ").unwrap();
1690
1691 let event = SCOutputEvent {
1692 context: EventExecutionContext {
1693 slot: slot_1,
1694 block: None,
1695 read_only: false,
1696 index_in_slot: index_1_0,
1697 call_stack: Default::default(),
1698 origin_operation_id: None,
1699 is_final: true,
1700 is_error: false,
1701 deferred_call_id: None,
1702 async_msg_id: None,
1703 },
1704 data: "message foo bar".to_string(),
1705 };
1706
1707 let to_insert_count = cache.max_entry_count - 5;
1708 let threshold = to_insert_count / 2;
1709 let mut events = (0..cache.max_entry_count - 5)
1710 .map(|i| {
1711 let mut event = event.clone();
1712 event.context.index_in_slot = i as u64;
1713 if i < threshold {
1714 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_1]);
1715 } else {
1716 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1717 }
1718 event
1719 })
1720 .collect::<Vec<SCOutputEvent>>();
1721
1722 let slot_2 = Slot::new(2, 0);
1723 let index_2_1 = 0u64;
1724 let event_slot_2 = {
1725 let mut event = event.clone();
1726 event.context.slot = slot_2;
1727 event.context.index_in_slot = index_2_1;
1728 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1729 event
1730 };
1731 let index_2_2 = 256u64;
1732 let event_slot_2_2 = {
1733 let mut event = event.clone();
1734 event.context.slot = slot_2;
1735 event.context.index_in_slot = index_2_2;
1736 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1737 event.context.is_error = true;
1738 event
1739 };
1740 events.push(event_slot_2.clone());
1741 events.push(event_slot_2_2.clone());
1742 events.shuffle(&mut thread_rng());
1744 cache.insert_multi_it(events.into_iter());
1745
1746 let filter_1 = EventFilter {
1747 is_error: Some(true),
1748 ..Default::default()
1749 };
1750
1751 let (_, filtered_events_1) = cache.get_filtered_sc_output_events(&filter_1);
1752 assert_eq!(filtered_events_1.len(), 1);
1753 assert!(filtered_events_1[0].context.is_error);
1754 assert_eq!(filtered_events_1[0].context.slot, slot_2);
1755 assert_eq!(filtered_events_1[0].context.index_in_slot, index_2_2);
1756 }
1757
1758 #[test]
1759 #[serial]
1760 fn test_orphan_index_key_does_not_panic() {
1761 let mut cache = setup();
1765 let emitter =
1766 Address::from_str("AU12qePoXhNbYWE1jZuafqJong7bbq1jw3k89RgbMawbrdZpaasoA").unwrap();
1767 let event = SCOutputEvent {
1768 context: EventExecutionContext {
1769 slot: Slot::new(1, 0),
1770 block: None,
1771 read_only: false,
1772 index_in_slot: 0,
1773 call_stack: VecDeque::from(vec![emitter]),
1774 origin_operation_id: None,
1775 is_final: true,
1776 is_error: false,
1777 deferred_call_id: None,
1778 async_msg_id: None,
1779 },
1780 data: "message foo bar".to_string(),
1781 };
1782 cache.insert(event.clone());
1783
1784 let event_key = cache
1786 .key_builder
1787 .key_from_event(&event, &KeyIndent::Event, &KeyKind::Regular)
1788 .unwrap();
1789 cache.db.delete(event_key).unwrap();
1790
1791 let filter = EventFilter {
1792 emitter_address: Some(emitter),
1793 ..Default::default()
1794 };
1795 let (_, events) = cache.get_filtered_sc_output_events(&filter);
1796 assert!(events.is_empty());
1797 }
1798
1799 #[test]
1800 #[serial]
1801 fn test_filter_optimisations() {
1802 let mut cache = setup();
1805 cache.max_entry_count = 10;
1806
1807 let slot_1 = Slot::new(1, 0);
1808 let index_1_0 = 0;
1809
1810 let dummy_addr =
1811 Address::from_str("AU12qePoXhNbYWE1jZuafqJong7bbq1jw3k89RgbMawbrdZpaasoA").unwrap();
1812 let emit_addr_1 =
1813 Address::from_str("AU122Em8qkqegdLb1eyH8rdkSCNEf7RZLeTJve4Q2inRPGiTJ2xNv").unwrap();
1814 let emit_addr_2 =
1815 Address::from_str("AU12WuVR1Td74q9eAbtYZUnk5jnRbUuUacyhQFwm217bV5v1mNqTZ").unwrap();
1816
1817 let event = SCOutputEvent {
1818 context: EventExecutionContext {
1819 slot: slot_1,
1820 block: None,
1821 read_only: false,
1822 index_in_slot: index_1_0,
1823 call_stack: Default::default(),
1824 origin_operation_id: None,
1825 is_final: true,
1826 is_error: true,
1827 deferred_call_id: None,
1828 async_msg_id: None,
1829 },
1830 data: "error foo bar".to_string(),
1831 };
1832
1833 let to_insert_count = cache.max_entry_count - 5;
1834 let threshold = to_insert_count / 2;
1835 let mut events = (0..cache.max_entry_count - 5)
1836 .map(|i| {
1837 let mut event = event.clone();
1838 event.context.index_in_slot = i as u64;
1839 if i < threshold {
1840 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_1]);
1841 } else {
1842 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1843 }
1844 event
1845 })
1846 .collect::<Vec<SCOutputEvent>>();
1847
1848 let slot_2 = Slot::new(2, 0);
1849 let index_2_1 = 0u64;
1850 let event_slot_2 = {
1851 let mut event = event.clone();
1852 event.context.slot = slot_2;
1853 event.context.index_in_slot = index_2_1;
1854 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1855 event
1856 };
1857 let index_2_2 = 256u64;
1858 let event_slot_2_2 = {
1859 let mut event = event.clone();
1860 event.context.slot = slot_2;
1861 event.context.index_in_slot = index_2_2;
1862 event.context.call_stack = VecDeque::from(vec![dummy_addr, emit_addr_2]);
1863 event
1865 };
1866 events.push(event_slot_2.clone());
1867 events.push(event_slot_2_2.clone());
1868 events.shuffle(&mut thread_rng());
1870 cache.insert_multi_it(events.into_iter());
1871
1872 let emit_addr_1_count = cache
1874 .filter_item_estimate_count(
1875 &KeyIndent::EmitterAddress,
1876 &FilterItem::EmitterAddress(emit_addr_1),
1877 )
1878 .unwrap();
1879 let emit_addr_2_count = cache
1880 .filter_item_estimate_count(
1881 &KeyIndent::EmitterAddress,
1882 &FilterItem::EmitterAddress(emit_addr_2),
1883 )
1884 .unwrap();
1885
1886 assert_eq!(emit_addr_1_count, (threshold) as u64);
1887 assert_eq!(emit_addr_2_count, (threshold + 1 + 2) as u64);
1888
1889 let slot_start_1 = Slot::new(3, 0); let slot_end_1 = Slot::new(0, 0); let slot_start_end_1 = Slot::new(2, 0)..Slot::new(1, 0); let slot_start_2 = Slot::new(1, 0);
1898 let slot_end_2 = Slot::new(2, 0);
1899 let slot_start_end_2 = Slot::new(1, 0)..Slot::new(2, 0);
1900
1901 let slot_start_count_1 = cache
1902 .filter_item_estimate_count(&KeyIndent::Event, &FilterItem::SlotStart(slot_start_1))
1903 .unwrap();
1904 let slot_end_count_1 = cache
1905 .filter_item_estimate_count(&KeyIndent::Event, &FilterItem::SlotEnd(slot_end_1))
1906 .unwrap();
1907 let slot_start_end_count_1 = cache
1908 .filter_item_estimate_count(
1909 &KeyIndent::Event,
1910 &FilterItem::SlotStartEnd(slot_start_end_1.start, slot_start_end_1.end),
1911 )
1912 .unwrap();
1913
1914 let slot_start_count_2 = cache
1915 .filter_item_estimate_count(&KeyIndent::Event, &FilterItem::SlotStart(slot_start_2))
1916 .unwrap();
1917 let slot_end_count_2 = cache
1918 .filter_item_estimate_count(&KeyIndent::Event, &FilterItem::SlotEnd(slot_end_2))
1919 .unwrap();
1920 let slot_start_end_count_2 = cache
1921 .filter_item_estimate_count(
1922 &KeyIndent::Event,
1923 &FilterItem::SlotStartEnd(slot_start_end_2.start, slot_start_end_2.end),
1924 )
1925 .unwrap();
1926
1927 assert_eq!(slot_start_count_1, 0);
1928 assert_eq!(slot_end_count_1, 0);
1929 assert_eq!(slot_start_end_count_1, 0);
1930
1931 let expected_count_per_period =
1932 THREAD_COUNT as u64 * MAX_EVENT_PER_OPERATION as u64 * MAX_OPERATIONS_PER_BLOCK as u64;
1933 assert_eq!(slot_start_count_2, expected_count_per_period);
1934 assert_eq!(slot_end_count_2, 2 * expected_count_per_period); assert_eq!(slot_start_end_count_2, expected_count_per_period);
1936
1937 let filter_1 = EventFilter {
1940 emitter_address: Some(emit_addr_1),
1941 is_error: Some(true),
1942 ..Default::default()
1943 };
1944
1945 let (query_counts, _filtered_events_1) = cache.get_filtered_sc_output_events(&filter_1);
1946 println!("threshold: {:?}", threshold);
1947 println!("query_counts: {:?}", query_counts);
1948
1949 assert_eq!(query_counts[0], threshold as u64);
1951 assert_eq!(query_counts[1], threshold as u64);
1954 }
1955}