massa_event_cache/
controller.rs

1// std
2use std::collections::{BTreeSet, VecDeque};
3use std::sync::Arc;
4// third-party
5use parking_lot::{Condvar, Mutex, RwLock};
6// internal
7use crate::event_cache::EventCache;
8use massa_models::execution::EventFilter;
9use massa_models::output_event::SCOutputEvent;
10
11/// structure used to communicate with controller
12#[derive(Debug, Default)]
13pub(crate) struct EventCacheWriterInputData {
14    /// set stop to true to stop the thread
15    pub stop: bool,
16    pub(crate) events: VecDeque<SCOutputEvent>,
17}
18
19impl EventCacheWriterInputData {
20    pub fn new() -> Self {
21        Self {
22            stop: Default::default(),
23            events: Default::default(),
24        }
25    }
26}
27
28/// interface that communicates with the worker thread
29#[cfg_attr(feature = "test-exports", mockall_wrap::wrap, mockall::automock)]
30pub trait EventCacheController: Send + Sync {
31    fn save_events(&self, events: VecDeque<SCOutputEvent>);
32
33    fn get_filtered_sc_output_events(&self, filter: &EventFilter) -> Vec<SCOutputEvent>;
34}
35
36#[derive(Clone)]
37/// implementation of the event cache controller
38pub struct EventCacheControllerImpl {
39    /// input data to process in the VM loop
40    /// with a wake-up condition variable that needs to be triggered when the data changes
41    pub(crate) input_data: Arc<(Condvar, Mutex<EventCacheWriterInputData>)>,
42    /// Event cache
43    pub(crate) cache: Arc<RwLock<EventCache>>,
44}
45
46impl EventCacheController for EventCacheControllerImpl {
47    fn save_events(&self, events: VecDeque<SCOutputEvent>) {
48        // lock input data
49        let mut input_data = self.input_data.1.lock();
50        input_data.events.extend(events);
51        massa_metrics::set_event_cache_vec_len(input_data.events.len());
52        // Wake up the condvar in EventCacheWriterThread waiting for events
53        self.input_data.0.notify_all();
54    }
55
56    fn get_filtered_sc_output_events(&self, filter: &EventFilter) -> Vec<SCOutputEvent> {
57        let mut res_0 = {
58            // Read from new events first
59            let lock_0 = self.input_data.1.lock();
60            #[allow(clippy::unnecessary_filter_map)]
61            let it = lock_0.events.iter().filter_map(|event| {
62                if let Some(start) = filter.start {
63                    if event.context.slot < start {
64                        return None;
65                    }
66                }
67                if let Some(end) = filter.end {
68                    if event.context.slot >= end {
69                        return None;
70                    }
71                }
72                if let Some(is_final) = filter.is_final {
73                    if event.context.is_final != is_final {
74                        return None;
75                    }
76                }
77                if let Some(is_error) = filter.is_error {
78                    if event.context.is_error != is_error {
79                        return None;
80                    }
81                }
82                match (
83                    filter.original_caller_address,
84                    event.context.call_stack.front(),
85                ) {
86                    (Some(addr1), Some(addr2)) if addr1 != *addr2 => return None,
87                    (Some(_), None) => return None,
88                    _ => (),
89                }
90                match (filter.emitter_address, event.context.call_stack.back()) {
91                    (Some(addr1), Some(addr2)) if addr1 != *addr2 => return None,
92                    (Some(_), None) => return None,
93                    _ => (),
94                }
95                match (
96                    filter.original_operation_id,
97                    event.context.origin_operation_id,
98                ) {
99                    (Some(addr1), Some(addr2)) if addr1 != addr2 => return None,
100                    (Some(_), None) => return None,
101                    _ => (),
102                }
103                Some(event)
104            });
105
106            let res_0: BTreeSet<SCOutputEvent> = it.cloned().collect();
107            // Drop the lock on the queue as soon as possible to avoid deadlocks
108            drop(lock_0);
109            res_0
110        };
111
112        let res_1 = {
113            // Read from db (on disk) events
114            let lock = self.cache.read();
115            let (_, res_1) = lock.get_filtered_sc_output_events(filter);
116            // Drop the lock on the event cache db asap
117            drop(lock);
118            res_1
119        };
120
121        // Merge results
122        let res_1: BTreeSet<SCOutputEvent> = BTreeSet::from_iter(res_1);
123        res_0.extend(res_1);
124        Vec::from_iter(res_0)
125    }
126}