1use std::sync::Arc;
3use std::thread;
4use parking_lot::{Condvar, Mutex, RwLock};
7use tracing::{debug, info};
8use crate::config::EventCacheConfig;
10use crate::controller::{
11 EventCacheController, EventCacheControllerImpl, EventCacheWriterInputData,
12};
13use crate::event_cache::EventCache;
14
15pub(crate) struct EventCacheWriterThread {
17 input_data: Arc<(Condvar, Mutex<EventCacheWriterInputData>)>,
19 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 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 pub fn main_loop(&mut self) {
47 loop {
48 self.wait_for_input();
49
50 let mut cache = self.cache.write();
57
58 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(cache);
73
74 if input_data.stop {
75 break;
77 }
78 }
79 }
80}
81
82pub trait EventCacheManager {
84 fn stop(&mut self);
89}
90
91pub struct EventCacheWriterManagerImpl {
94 pub(crate) input_data: Arc<(Condvar, Mutex<EventCacheWriterInputData>)>,
97 pub(crate) thread_handle: Option<std::thread::JoinHandle<()>>,
99}
100
101impl EventCacheManager for EventCacheWriterManagerImpl {
102 fn stop(&mut self) {
104 info!("Stopping Execution controller...");
105 {
107 let mut input_wlock = self.input_data.1.lock();
108 input_wlock.stop = true;
109 self.input_data.0.notify_one();
110 }
111 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 let input_data = Arc::new((Condvar::new(), Mutex::new(EventCacheWriterInputData::new())));
136 let input_data_clone = input_data.clone();
137
138 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 let manager = EventCacheWriterManagerImpl {
153 input_data,
154 thread_handle: Some(thread_handle),
155 };
156
157 (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 #[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 #[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 let cache_read = cache.read();
264
265 controller.save_events([sample_event()].into());
266 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 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}