massa_event_cache/
event_cache.rs

1// std
2use std::cmp::max;
3use std::collections::{BTreeMap, BTreeSet};
4use std::path::Path;
5// third-party
6use num_enum::IntoPrimitive;
7use rocksdb::{IteratorMode, Options, WriteBatch, DB};
8use tracing::{debug, warn};
9// internal
10use 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";
23// const COUNTER_INIT_ERROR: &str = "critical: cannot init rocksdb counters";
24const 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/// Prefix u8 used to identify rocksdb keys
33#[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
45/// Key for this type of data that we want to get
46enum 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    /// Operation Id Serializer
63    op_id_ser: OperationIdSerializer,
64}
65
66impl DbKeyBuilder {
67    fn new() -> Self {
68        Self {
69            op_id_ser: OperationIdSerializer::new(),
70        }
71    }
72
73    /// Low level key builder function
74    /// There is no guarantees that the key will be unique
75    /// Recommended to use high level function like:
76    /// `key_from_event`, `prefix_key_from_indent`,
77    /// `prefix_key_from_filter_item` or `counter_key_from_filter_item`
78    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    /// Key usually used to populate the DB
117    fn key_from_event(
118        &self,
119        event: &SCOutputEvent,
120        indent: &KeyIndent,
121        key_kind: &KeyKind,
122    ) -> Option<Vec<u8>> {
123        // Db format:
124        // * Regular keys:
125        //   * Event key: [Event Indent][Slot][Index] -> Event value: Event serialized
126        //   * Emitter address key: [Emitter Address Indent][Addr][Addr len][Event key] -> []
127        // * Prefix keys:
128        //   * Emitter address prefix key: [Emitter Address Indent][Addr][Addr len]
129        // * Counter keys:
130        //   * Emitter address counter key: [Counter indent][Emitter Address Indent][Addr][Addr len][Event key] -> u64
131
132        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    /// Prefix key to iterate over all events / emitter_address / ...
195    fn prefix_key_from_indent(&self, indent: &KeyIndent) -> Vec<u8> {
196        self.key(indent, KeyBuilderType::None, &KeyKind::Regular)
197    }
198
199    /// Prefix key to iterate over specific emitter_address / operation_id / ...
200    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    /// Counter key for specific emitter_address / operation_id / ...
230    fn counter_key_from_filter_item(
231        &self,
232        filter_item: &FilterItem,
233        indent: &KeyIndent,
234    ) -> Vec<u8> {
235        // High level key builder function
236        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
267/// Disk based event cache db (rocksdb based)
268pub(crate) struct EventCache {
269    /// RocksDB database
270    db: DB,
271    /// How many entries are in the db. Count is initialized at creation time by iterating
272    /// over all the entries in the db then it is maintained in memory
273    entry_count: usize,
274    /// Maximum number of entries we want to keep in the db.
275    /// When this maximum is reached `snip_amount` entries are removed
276    max_entry_count: usize,
277    /// How many entries are removed when `entry_count` reaches `max_entry_count`
278    snip_amount: usize,
279    /// Event serializer
280    event_ser: SCOutputEventSerializer,
281    /// Event deserializer
282    event_deser: SCOutputEventDeserializer,
283    /// Key builder
284    key_builder: DbKeyBuilder,
285    /// First event slot in db
286    first_slot: Slot,
287    /// Last event slot in db
288    last_slot: Slot,
289    /// Thread count
290    thread_count: u8,
291    /// Maximum number of events per operation
292    max_events_per_operation: u64,
293    /// Maximum number of operations per block
294    max_operations_per_block: u64,
295    /// Max number of events returned by a query
296    max_events_per_query: usize,
297}
298
299impl EventCache {
300    #[allow(clippy::too_many_arguments)]
301    /// Create a new EventCache
302    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        // Clear the db
314        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    /// From an event add keys & values into a rocksdb batch
350    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        // Keep track of last slot (and start slot) of events in the DB
420        // Help for event filtering
421        self.last_slot = max(self.last_slot, event.context.slot);
422    }
423
424    #[allow(dead_code)]
425    /// Insert a new event in the cache
426    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        // Note:
436        // This assumes that events are always added, never overwritten
437        self.entry_count = self.entry_count.saturating_add(1);
438
439        debug!("(Event insert) entry_count is: {}", self.entry_count);
440    }
441
442    /// Insert new events in the cache
443    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        // Note:
460        // This assumes that events are always added, never overwritten
461        self.entry_count = self.entry_count.saturating_add(events_len);
462
463        debug!("(Events insert) entry_count is: {}", self.entry_count);
464    }
465
466    /// Get events filtered by the given argument
467    pub(crate) fn get_filtered_sc_output_events(
468        &self,
469        filter: &EventFilter,
470    ) -> (Vec<u64>, Vec<SCOutputEvent>) {
471        // All events stored in this cache are final: a query for non-final
472        // events can never match anything here
473        if filter.is_final == Some(false) {
474            return (vec![], vec![]);
475        }
476
477        // Step 1
478        // Build a (sorted) map with key: (counter value, indent), value: filter
479        // Will be used to iterate from the lower count index to the highest count index
480        // e.g. if index for emitter address is 10 (index count), and origin operation id is 20
481        //      iter over emitter address index then origin operation id index
482
483        let mut filter_items = from_event_filter(filter);
484
485        if filter_items.is_empty() {
486            // No criteria (is_final is None or Some(true)): scan the whole event
487            // column, bounded below by max_events_per_query
488            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        // Step 2: apply filter from the lowest counter to the highest counter
508
509        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        // Step 3: get values & deserialize
524
525        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                // An index key without its event (orphan index) must not kill
547                // the node: skip it
548                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                        // Stop as soon as our key does not start with the right indent
608                        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                        // We have already found as many items as in a previous search
620                        // As we search from the lowest count to the highest count, we will never add new items
621                        // in our result, so we can break here
622                        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                        // Stop as soon as our key does not start with the right indent
652                        break;
653                    }
654
655                    if !kvb.0.starts_with(prefix_filter.as_slice()) {
656                        // Stop as soon as our key does not start with our current prefix
657                        break;
658                    }
659
660                    let found = kvb
661                        .0
662                        .strip_prefix(prefix_filter.as_slice())
663                        .unwrap() // safe to unwrap() - already tested
664                        .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                        // We have already found as many items as in a previous search
674                        // As we search from the lowest count to the highest count, we will never add new items
675                        // in our result, so we can break here
676                        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    /// Estimate for a given KeyIndent & FilterItem the number of row to iterate
690    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); // If start is after self.last_slot, we have no events in the filter
701                                   // Note: Pessimistic estimation - should we keep an average count of events per slot
702                                   //       and use that instead?
703                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                // If start is after end, we have no events in the filter as it is inconsistent
709                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                // If end is after self.first_slot, we have no events in the filter
716                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    /// Try to remove some entries from the db
767    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() // safe to unwrap - just tested it
782                .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            // delete all associated key
795            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        // delete the key and reduce entry_count
853        self.db.write(&batch).expect(CRUD_ERROR);
854        self.entry_count = self.entry_count.saturating_sub(snipped_count);
855
856        // delete key counters where value == 0
857        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        // Update first_slot / last_slot in the DB
869        if self.entry_count == 0 {
870            // Reset
871            self.first_slot = Slot::new(0, 0);
872            self.last_slot = Slot::new(0, 0);
873        } else {
874            // Get the first event in the db
875            // By using a prefix iterator this should be fast
876
877            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/// A filter parameter - used to decompose an EventFilter in multiple filters
893#[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
904/// Convert a EventFilter into a list of (KeyIndent, FilterItem)
905fn 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    /// Iterate over all keys & values in the db - test only
947    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    // std
961    use std::collections::VecDeque;
962    use std::str::FromStr;
963    // third-party
964    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    // internal
970    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, // MAX_EVENTS_PER_QUERY,
990        )
991    }
992
993    #[test]
994    #[serial]
995    fn test_db_insert_order() {
996        // Test that the data will be correctly ordered (when iterated from start) in db
997
998        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        // Randomize the events so we insert in random orders in DB
1043        events.shuffle(&mut thread_rng());
1044
1045        for event in events {
1046            cache.insert(event);
1047        }
1048
1049        // Now check that we are going to iter in correct order
1050        // let db_it = cache.db_iter(Some(IteratorMode::Start));
1051        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        // Test insert (and snip) so we do not store too many event in the cache
1083
1084        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        // fill the db: add cache.max_entry_count entries
1102        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        // insert one more entry
1110        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        // Test insert_multi_it (and snip) so we do not store too many event in the cache
1124
1125        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        // insert one more entry
1152        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        // Test snip so we enforce that all db keys are removed
1166
1167        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        // Test snip so we enforce that all db keys are removed
1205
1206        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        // Check counters key length
1242        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        // println!("kvbs: {:#?}", kvbs);
1256
1257        // Expected 4 counters:
1258        // 2 for emitter address (emit_addr_1 & emit_addr_2)
1259        // 1 for original caller address (dummy_addr)
1260        // 1 for is_error(false)
1261        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        // println!("v1: {:?} - v2: {:?}", v1, v2);
1276        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        // println!("v1: {:?} - v2: {:?}", v1, v2);
1285        assert_eq!(v1, Ok(None)); // counter has been removed
1286        assert_eq!(v2, Ok(Some(1u64.to_be_bytes().to_vec())));
1287    }
1288
1289    #[test]
1290    #[serial]
1291    fn test_event_filter() {
1292        // Test that the data will be correctly ordered (when filtered) in db
1293
1294        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        // Randomize the events so we insert in random orders in the DB
1339        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        // Test get_filtered_sc_output_events + op id
1361
1362        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        // Randomize the events so we insert in random orders in the DB
1418        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        // Test get_filtered_sc_output_events + emitter address
1459
1460        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        // Randomize the events so we insert in random orders in the DB
1526        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        // filter with emit_addr_2
1544        {
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        // filter with dummy_addr
1553        {
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        // filter with address that is not in the DB
1559        {
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        // Test get_filtered_sc_output_events + original caller addr
1570
1571        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        // Randomize the events so we insert in random orders in the DB
1637        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        // Test get_filtered_sc_output_events + is error
1677
1678        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        // Randomize the events so we insert in random orders in the DB
1743        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        // An index key whose event is missing (e.g. event key overwritten by a
1762        // duplicate (slot, index) then snipped) used to panic on unwrap()
1763
1764        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        // Remove the event value but keep its index keys
1785        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        // Test we iterate over the right number of rows when filtering
1803
1804        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.context.is_error = true;
1864            event
1865        };
1866        events.push(event_slot_2.clone());
1867        events.push(event_slot_2_2.clone());
1868        // Randomize the events so we insert in random orders in the DB
1869        events.shuffle(&mut thread_rng());
1870        cache.insert_multi_it(events.into_iter());
1871
1872        // Check if we correctly count the number of events in the DB with emit_addr_1 & emit_addr_2
1873        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        // Check if we correctly count the number of events in the DB with slot related filters
1890        // First with filters that should give a count of 0 (too early, too late, or inconsistent start/end)
1891        let slot_start_1 = Slot::new(3, 0); // start is after the last slot
1892        let slot_end_1 = Slot::new(0, 0); // end is before the first slot
1893        let slot_start_end_1 = Slot::new(2, 0)..Slot::new(1, 0); // start is after end
1894
1895        // Then with filters that should give a positive count
1896        // num_thread * max_events_per_operation * max_operations_per_block
1897        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); // 2 periods as first_slot is (0,0)
1935        assert_eq!(slot_start_end_count_2, expected_count_per_period);
1936
1937        // Check if we query first by emitter address then is_error
1938
1939        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        // Check that we iter no more than needed (here: only 2 (== threshold) event with emit addr 1)
1950        assert_eq!(query_counts[0], threshold as u64);
1951        // For second filter (is_error) we could have iter more (all events have is_error = true)
1952        // but as soon as we found 2 items we could return (as the previous filter already limit the final count)
1953        assert_eq!(query_counts[1], threshold as u64);
1954    }
1955}