massa_event_cache/
worker.rs

1// std
2use std::sync::Arc;
3use std::thread;
4// third-party
5// use massa_time::MassaTime;
6use parking_lot::{Condvar, Mutex, RwLock};
7use tracing::{debug, info};
8// internal
9use crate::config::EventCacheConfig;
10use crate::controller::{
11    EventCacheController, EventCacheControllerImpl, EventCacheWriterInputData,
12};
13use crate::event_cache::EventCache;
14
15/// Structure gathering all elements needed by the event cache thread
16pub(crate) struct EventCacheWriterThread {
17    // A copy of the input data allowing access to incoming requests
18    input_data: Arc<(Condvar, Mutex<EventCacheWriterInputData>)>,
19    /// Event cache
20    cache: Arc<RwLock<EventCache>>,
21}
22
23impl EventCacheWriterThread {
24    fn new(
25        input_data: Arc<(Condvar, Mutex<EventCacheWriterInputData>)>,
26        event_cache: Arc<RwLock<EventCache>>,
27    ) -> Self {
28        Self {
29            input_data,
30            cache: event_cache,
31        }
32    }
33
34    /// Waits until there is something to do: events to flush and/or a stop request.
35    ///
36    /// The input is not consumed here: see [`Self::main_loop`] for why the batch is only taken once
37    /// the cache is locked.
38    fn wait_for_input(&self) {
39        let mut input_data_lock = self.input_data.1.lock();
40        while input_data_lock.events.is_empty() && !input_data_lock.stop {
41            self.input_data.0.wait(&mut input_data_lock);
42        }
43    }
44
45    /// Main loop of the worker
46    pub fn main_loop(&mut self) {
47        loop {
48            self.wait_for_input();
49
50            // Lock the cache before taking the batch out of the queue. Readers look for events in
51            // the queue, then in the cache. Taking the batch first left a window where the events
52            // were in neither, so a reader could miss events that execution had already marked as
53            // final. With the cache locked first, a reader that no longer finds them in the queue
54            // waits on the cache lock until they are inserted. Readers never hold both locks, so
55            // this order cannot deadlock.
56            let mut cache = self.cache.write();
57
58            // Take the whole input, resetting it. The stop flag comes with the events so that a
59            // final queued batch is still flushed before the loop terminates.
60            let input_data: EventCacheWriterInputData =
61                std::mem::take(&mut *self.input_data.1.lock());
62            debug!(
63                "Event cache writer loop triggered, {} events, stop = {}",
64                input_data.events.len(),
65                input_data.stop
66            );
67
68            if !input_data.events.is_empty() {
69                cache.insert_multi_it(input_data.events.into_iter());
70            }
71            // drop the lock as early as possible
72            drop(cache);
73
74            if input_data.stop {
75                // we need to stop
76                break;
77            }
78        }
79    }
80}
81
82/// Event cache manager trait used to stop the event cache thread
83pub trait EventCacheManager {
84    /// Stop the event cache thread
85    /// Note that we do not take self by value to consume it
86    /// because it is not allowed to move out of `Box<dyn ExecutionManager>`
87    /// This will improve if the `unsized_fn_params` feature stabilizes enough to be safely usable.
88    fn stop(&mut self);
89}
90
91/// ... manager
92/// Allows stopping the ... worker
93pub struct EventCacheWriterManagerImpl {
94    /// input data to process in the VM loop
95    /// with a wake-up condition variable that needs to be triggered when the data changes
96    pub(crate) input_data: Arc<(Condvar, Mutex<EventCacheWriterInputData>)>,
97    /// handle used to join the worker thread
98    pub(crate) thread_handle: Option<std::thread::JoinHandle<()>>,
99}
100
101impl EventCacheManager for EventCacheWriterManagerImpl {
102    /// stops the worker
103    fn stop(&mut self) {
104        info!("Stopping Execution controller...");
105        // notify the worker thread to stop
106        {
107            let mut input_wlock = self.input_data.1.lock();
108            input_wlock.stop = true;
109            self.input_data.0.notify_one();
110        }
111        // join the thread
112        if let Some(join_handle) = self.thread_handle.take() {
113            join_handle.join().expect("VM controller thread panicked");
114        }
115        info!("Execution controller stopped");
116    }
117}
118
119pub fn start_event_cache_writer_worker(
120    cfg: EventCacheConfig,
121) -> (Box<dyn EventCacheManager>, Box<dyn EventCacheController>) {
122    let event_cache = Arc::new(RwLock::new(EventCache::new(
123        cfg.event_cache_path.as_path(),
124        cfg.max_event_cache_length,
125        cfg.snip_amount,
126        cfg.thread_count,
127        cfg.max_call_stack_length,
128        cfg.max_event_data_length,
129        cfg.max_events_per_operation,
130        cfg.max_operations_per_block,
131        cfg.max_events_per_query,
132    )));
133
134    // define the input data interface
135    let input_data = Arc::new((Condvar::new(), Mutex::new(EventCacheWriterInputData::new())));
136    let input_data_clone = input_data.clone();
137
138    // create a controller
139    let controller = EventCacheControllerImpl {
140        input_data: input_data.clone(),
141        cache: event_cache.clone(),
142    };
143
144    let thread_builder = thread::Builder::new().name("event_cache".into());
145    let thread_handle = thread_builder
146        .spawn(move || {
147            EventCacheWriterThread::new(input_data_clone, event_cache).main_loop();
148        })
149        .expect("failed to spawn thread : event_cache");
150
151    // create a manager
152    let manager = EventCacheWriterManagerImpl {
153        input_data,
154        thread_handle: Some(thread_handle),
155    };
156
157    // return the manager and controller pair
158    (Box::new(manager), Box::new(controller))
159}
160
161#[cfg(test)]
162mod tests {
163    use super::*;
164    use crate::controller::EventCacheWriterInputData;
165    use crate::event_cache::EventCache;
166    use massa_models::config::{
167        MAX_EVENT_DATA_SIZE, MAX_EVENT_PER_OPERATION, MAX_OPERATIONS_PER_BLOCK,
168        MAX_RECURSIVE_CALLS_DEPTH, THREAD_COUNT,
169    };
170    use massa_models::output_event::{EventExecutionContext, SCOutputEvent};
171    use massa_models::slot::Slot;
172    use std::time::{Duration, Instant};
173    use tempfile::TempDir;
174
175    fn sample_event() -> SCOutputEvent {
176        SCOutputEvent {
177            context: EventExecutionContext {
178                slot: Slot::new(1, 0),
179                block: None,
180                read_only: false,
181                index_in_slot: 0,
182                call_stack: Default::default(),
183                origin_operation_id: None,
184                is_final: true,
185                is_error: false,
186                deferred_call_id: None,
187                async_msg_id: None,
188            },
189            data: "shutdown-race event".to_string(),
190        }
191    }
192
193    /// Reproduces the lost-stop-signal race: a final batch of events is queued
194    /// *together with* a stop request. The writer must flush the batch and then
195    /// terminate, instead of consuming the stop flag with the batch and waiting
196    /// on the condvar forever.
197    #[test]
198    fn stop_is_not_lost_when_a_final_batch_is_queued() {
199        let tmp = TempDir::new().unwrap();
200        let cache = Arc::new(RwLock::new(EventCache::new(
201            tmp.path(),
202            1000,
203            300,
204            THREAD_COUNT,
205            MAX_RECURSIVE_CALLS_DEPTH,
206            MAX_EVENT_DATA_SIZE as u64,
207            MAX_EVENT_PER_OPERATION as u64,
208            MAX_OPERATIONS_PER_BLOCK as u64,
209            5000,
210        )));
211
212        let mut input = EventCacheWriterInputData::new();
213        input.events.push_back(sample_event());
214        input.stop = true;
215        let input_data = Arc::new((Condvar::new(), Mutex::new(input)));
216
217        let mut worker = EventCacheWriterThread::new(input_data, cache);
218        let handle = thread::spawn(move || worker.main_loop());
219
220        let deadline = Instant::now() + Duration::from_secs(10);
221        while !handle.is_finished() {
222            assert!(
223                Instant::now() < deadline,
224                "event cache writer thread did not stop after a stop request"
225            );
226            thread::sleep(Duration::from_millis(10));
227        }
228        handle.join().expect("event cache writer thread panicked");
229    }
230
231    /// Reproduces the handoff race between the writer and readers: once the writer has taken a batch
232    /// out of the queue, the events must stay visible to readers until they are in the cache.
233    ///
234    /// The test holds a read lock on the cache, so the writer cannot insert. A reader at that point
235    /// must still find the saved event, either in the queue or in the cache. The writer used to empty
236    /// the queue first and lock the cache afterwards: the event was in neither.
237    #[test]
238    fn saved_events_stay_visible_while_the_writer_waits_for_the_cache() {
239        let tmp = TempDir::new().unwrap();
240        let cache = Arc::new(RwLock::new(EventCache::new(
241            tmp.path(),
242            1000,
243            300,
244            THREAD_COUNT,
245            MAX_RECURSIVE_CALLS_DEPTH,
246            MAX_EVENT_DATA_SIZE as u64,
247            MAX_EVENT_PER_OPERATION as u64,
248            MAX_OPERATIONS_PER_BLOCK as u64,
249            5000,
250        )));
251        let input_data = Arc::new((Condvar::new(), Mutex::new(EventCacheWriterInputData::new())));
252        let controller = EventCacheControllerImpl {
253            input_data: input_data.clone(),
254            cache: cache.clone(),
255        };
256
257        let mut worker = EventCacheWriterThread::new(input_data.clone(), cache.clone());
258        let handle = thread::spawn(move || worker.main_loop());
259
260        let filter = massa_models::execution::EventFilter::default();
261        {
262            // keep the writer from inserting
263            let cache_read = cache.read();
264
265            controller.save_events([sample_event()].into());
266            // let the writer wake up and process the batch as far as it can
267            thread::sleep(Duration::from_millis(200));
268
269            let in_queue = !input_data.1.lock().events.is_empty();
270            let in_cache = !cache_read
271                .get_filtered_sc_output_events(&filter)
272                .1
273                .is_empty();
274            assert!(
275                in_queue || in_cache,
276                "a saved event is neither in the queue nor in the cache"
277            );
278        }
279
280        // once the cache is free, the writer inserts the event
281        let deadline = Instant::now() + Duration::from_secs(10);
282        while cache
283            .read()
284            .get_filtered_sc_output_events(&filter)
285            .1
286            .is_empty()
287        {
288            assert!(Instant::now() < deadline, "the event was never inserted");
289            thread::sleep(Duration::from_millis(10));
290        }
291
292        {
293            let mut input = input_data.1.lock();
294            input.stop = true;
295            input_data.0.notify_one();
296        }
297        handle.join().expect("event cache writer thread panicked");
298    }
299}