nautilus_common/live/
runner.rs1use std::cell::RefCell;
21
22use crate::messages::{DataEvent, ExecutionEvent, SystemCommand, SystemEvent};
23
24#[must_use]
30pub fn get_data_event_sender() -> tokio::sync::mpsc::UnboundedSender<DataEvent> {
31 DATA_EVENT_SENDER.with(|sender| {
32 sender
33 .borrow()
34 .as_ref()
35 .expect("Data event sender should be initialized by runner")
36 .clone()
37 })
38}
39
40#[must_use]
45pub fn try_get_data_event_sender() -> Option<tokio::sync::mpsc::UnboundedSender<DataEvent>> {
46 DATA_EVENT_SENDER.with(|sender| sender.borrow().as_ref().cloned())
47}
48
49pub fn set_data_event_sender(sender: tokio::sync::mpsc::UnboundedSender<DataEvent>) {
57 DATA_EVENT_SENDER.with(|s| {
58 let mut slot = s.borrow_mut();
59 assert!(slot.is_none(), "Data event sender can only be set once");
60 *slot = Some(sender);
61 });
62}
63
64pub fn replace_data_event_sender(sender: tokio::sync::mpsc::UnboundedSender<DataEvent>) {
66 DATA_EVENT_SENDER.with(|s| {
67 *s.borrow_mut() = Some(sender);
68 });
69}
70
71#[must_use]
77pub fn get_system_event_sender() -> tokio::sync::mpsc::UnboundedSender<SystemEvent> {
78 SYSTEM_EVENT_SENDER.with(|sender| {
79 sender
80 .borrow()
81 .as_ref()
82 .expect("System event sender should be initialized by runner")
83 .clone()
84 })
85}
86
87#[must_use]
91pub fn try_get_system_event_sender() -> Option<tokio::sync::mpsc::UnboundedSender<SystemEvent>> {
92 SYSTEM_EVENT_SENDER.with(|sender| sender.borrow().as_ref().cloned())
93}
94
95pub fn set_system_event_sender(sender: tokio::sync::mpsc::UnboundedSender<SystemEvent>) {
103 SYSTEM_EVENT_SENDER.with(|s| {
104 let mut slot = s.borrow_mut();
105 assert!(slot.is_none(), "System event sender can only be set once");
106 *slot = Some(sender);
107 });
108}
109
110pub fn replace_system_event_sender(sender: tokio::sync::mpsc::UnboundedSender<SystemEvent>) {
112 SYSTEM_EVENT_SENDER.with(|s| {
113 *s.borrow_mut() = Some(sender);
114 });
115}
116
117#[must_use]
123pub fn get_system_command_sender() -> tokio::sync::mpsc::UnboundedSender<SystemCommand> {
124 SYSTEM_COMMAND_SENDER.with(|sender| {
125 sender
126 .borrow()
127 .as_ref()
128 .expect("System command sender should be initialized by runner")
129 .clone()
130 })
131}
132
133#[must_use]
137pub fn try_get_system_command_sender() -> Option<tokio::sync::mpsc::UnboundedSender<SystemCommand>>
138{
139 SYSTEM_COMMAND_SENDER.with(|sender| sender.borrow().as_ref().cloned())
140}
141
142pub fn set_system_command_sender(sender: tokio::sync::mpsc::UnboundedSender<SystemCommand>) {
150 SYSTEM_COMMAND_SENDER.with(|s| {
151 let mut slot = s.borrow_mut();
152 assert!(slot.is_none(), "System command sender can only be set once");
153 *slot = Some(sender);
154 });
155}
156
157pub fn replace_system_command_sender(sender: tokio::sync::mpsc::UnboundedSender<SystemCommand>) {
159 SYSTEM_COMMAND_SENDER.with(|s| {
160 *s.borrow_mut() = Some(sender);
161 });
162}
163
164#[must_use]
170pub fn get_exec_event_sender() -> tokio::sync::mpsc::UnboundedSender<ExecutionEvent> {
171 EXEC_EVENT_SENDER.with(|sender| {
172 sender
173 .borrow()
174 .as_ref()
175 .expect("Execution event sender should be initialized by runner")
176 .clone()
177 })
178}
179
180#[must_use]
184pub fn try_get_exec_event_sender() -> Option<tokio::sync::mpsc::UnboundedSender<ExecutionEvent>> {
185 EXEC_EVENT_SENDER.with(|sender| sender.borrow().as_ref().cloned())
186}
187
188pub fn set_exec_event_sender(sender: tokio::sync::mpsc::UnboundedSender<ExecutionEvent>) {
196 EXEC_EVENT_SENDER.with(|s| {
197 let mut slot = s.borrow_mut();
198 assert!(
199 slot.is_none(),
200 "Execution event sender can only be set once"
201 );
202 *slot = Some(sender);
203 });
204}
205
206pub fn replace_exec_event_sender(sender: tokio::sync::mpsc::UnboundedSender<ExecutionEvent>) {
208 EXEC_EVENT_SENDER.with(|s| {
209 *s.borrow_mut() = Some(sender);
210 });
211}
212
213thread_local! {
214 static DATA_EVENT_SENDER: RefCell<Option<tokio::sync::mpsc::UnboundedSender<DataEvent>>> = const { RefCell::new(None) };
215 static EXEC_EVENT_SENDER: RefCell<Option<tokio::sync::mpsc::UnboundedSender<ExecutionEvent>>> = const { RefCell::new(None) };
216 static SYSTEM_EVENT_SENDER: RefCell<Option<tokio::sync::mpsc::UnboundedSender<SystemEvent>>> = const { RefCell::new(None) };
217 static SYSTEM_COMMAND_SENDER: RefCell<Option<tokio::sync::mpsc::UnboundedSender<SystemCommand>>> = const { RefCell::new(None) };
218}
219
220#[cfg(test)]
221mod tests {
222 use std::sync::{Arc, Barrier};
223
224 use rstest::rstest;
225
226 use super::*;
227
228 #[rstest]
229 fn test_replace_data_event_sender_overwrites_previous() {
230 assert_sender_replaced(replace_data_event_sender, get_data_event_sender);
231 }
232
233 #[rstest]
234 fn test_replace_exec_event_sender_overwrites_previous() {
235 assert_sender_replaced(replace_exec_event_sender, get_exec_event_sender);
236 }
237
238 #[rstest]
239 fn test_replace_system_event_sender_overwrites_previous() {
240 assert_sender_replaced(replace_system_event_sender, get_system_event_sender);
241 }
242
243 #[rstest]
244 fn test_replace_system_command_sender_overwrites_previous() {
245 assert_sender_replaced(replace_system_command_sender, get_system_command_sender);
246 }
247
248 #[rstest]
249 fn test_event_senders_are_thread_local() {
250 assert_sender_thread_local(replace_data_event_sender, get_data_event_sender);
251 assert_sender_thread_local(replace_exec_event_sender, get_exec_event_sender);
252 assert_sender_thread_local(replace_system_event_sender, get_system_event_sender);
253 assert_sender_thread_local(replace_system_command_sender, get_system_command_sender);
254 }
255
256 #[rstest]
257 fn test_set_data_event_sender_panics_on_double_set() {
258 let result = std::thread::spawn(|| {
259 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
260 let (tx2, _rx2) = tokio::sync::mpsc::unbounded_channel();
261 set_data_event_sender(tx1);
262 set_data_event_sender(tx2);
263 })
264 .join();
265 assert!(result.is_err());
266 }
267
268 #[rstest]
269 fn test_set_exec_event_sender_panics_on_double_set() {
270 let result = std::thread::spawn(|| {
271 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
272 let (tx2, _rx2) = tokio::sync::mpsc::unbounded_channel();
273 set_exec_event_sender(tx1);
274 set_exec_event_sender(tx2);
275 })
276 .join();
277 assert!(result.is_err());
278 }
279
280 #[rstest]
281 fn test_set_system_event_sender_panics_on_double_set() {
282 let result = std::thread::spawn(|| {
283 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
284 let (tx2, _rx2) = tokio::sync::mpsc::unbounded_channel();
285 set_system_event_sender(tx1);
286 set_system_event_sender(tx2);
287 })
288 .join();
289 assert!(result.is_err());
290 }
291
292 #[rstest]
293 fn test_set_system_command_sender_panics_on_double_set() {
294 let result = std::thread::spawn(|| {
295 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
296 let (tx2, _rx2) = tokio::sync::mpsc::unbounded_channel();
297 set_system_command_sender(tx1);
298 set_system_command_sender(tx2);
299 })
300 .join();
301 assert!(result.is_err());
302 }
303
304 #[rstest]
305 fn test_try_get_exec_event_sender_returns_none_when_unset() {
306 let result = std::thread::spawn(try_get_exec_event_sender)
307 .join()
308 .unwrap();
309 assert!(result.is_none());
310 }
311
312 #[rstest]
313 fn test_try_get_system_event_sender_returns_none_when_unset() {
314 let result = std::thread::spawn(try_get_system_event_sender)
315 .join()
316 .unwrap();
317 assert!(result.is_none());
318 }
319
320 #[rstest]
321 fn test_try_get_system_command_sender_returns_none_when_unset() {
322 let result = std::thread::spawn(try_get_system_command_sender)
323 .join()
324 .unwrap();
325 assert!(result.is_none());
326 }
327
328 fn assert_sender_replaced<T: Send + 'static>(
329 replace: fn(tokio::sync::mpsc::UnboundedSender<T>),
330 get: fn() -> tokio::sync::mpsc::UnboundedSender<T>,
331 ) {
332 std::thread::spawn(move || {
333 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
334 let (tx2, _rx2) = tokio::sync::mpsc::unbounded_channel();
335
336 replace(tx1.clone());
337 replace(tx2.clone());
338 let sender = get();
339
340 assert!(!sender.same_channel(&tx1));
341 assert!(sender.same_channel(&tx2));
342 })
343 .join()
344 .expect("sender replacement test thread should join");
345 }
346
347 fn assert_sender_thread_local<T: Send + 'static>(
348 replace: fn(tokio::sync::mpsc::UnboundedSender<T>),
349 get: fn() -> tokio::sync::mpsc::UnboundedSender<T>,
350 ) {
351 let barrier = Arc::new(Barrier::new(2));
352 let (tx1, _rx1) = tokio::sync::mpsc::unbounded_channel();
353 let (tx2, _rx2) = tokio::sync::mpsc::unbounded_channel();
354 let expected1 = tx1.clone();
355 let expected2 = tx2.clone();
356
357 let barrier1 = Arc::clone(&barrier);
358
359 let thread1 = std::thread::spawn(move || {
360 replace(tx1);
361 barrier1.wait();
362 assert!(get().same_channel(&expected1));
363 });
364
365 let thread2 = std::thread::spawn(move || {
366 replace(tx2);
367 barrier.wait();
368 assert!(get().same_channel(&expected2));
369 });
370
371 thread1
372 .join()
373 .expect("first sender isolation test thread should join");
374 thread2
375 .join()
376 .expect("second sender isolation test thread should join");
377 }
378}