Skip to main content

nautilus_common/live/
runner.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Tokio-based channel senders for live trading runtime.
17//!
18//! This module provides thread-local storage for tokio mpsc channels used in live trading.
19
20use std::cell::RefCell;
21
22use crate::messages::{DataEvent, ExecutionEvent, SystemCommand, SystemEvent};
23
24/// Gets the global data event sender.
25///
26/// # Panics
27///
28/// Panics if the sender is uninitialized.
29#[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/// Attempts to get the global data event sender without panicking.
41///
42/// Returns `None` if the sender is not initialized (e.g., in Python/v1 bridge environments
43/// before a runner or adapter bridge has registered a sender).
44#[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
49/// Sets the global data event sender.
50///
51/// Can only be called once per thread.
52///
53/// # Panics
54///
55/// Panics if a sender has already been set.
56pub 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
64/// Replaces the global data event sender for the current thread.
65pub 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/// Gets the global system event sender.
72///
73/// # Panics
74///
75/// Panics if the sender is uninitialized.
76#[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/// Attempts to get the global system event sender without panicking.
88///
89/// Returns `None` if the sender is not initialized (e.g., in test environments).
90#[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
95/// Sets the global system event sender.
96///
97/// Can only be called once per thread.
98///
99/// # Panics
100///
101/// Panics if a sender has already been set.
102pub 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
110/// Replaces the global system event sender for the current thread.
111pub 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/// Gets the global system command sender.
118///
119/// # Panics
120///
121/// Panics if the sender is uninitialized.
122#[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/// Attempts to get the global system command sender without panicking.
134///
135/// Returns `None` if the sender is not initialized.
136#[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
142/// Sets the global system command sender.
143///
144/// Can only be called once per thread.
145///
146/// # Panics
147///
148/// Panics if a sender has already been set.
149pub 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
157/// Replaces the global system command sender for the current thread.
158pub 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/// Gets the global execution event sender.
165///
166/// # Panics
167///
168/// Panics if the sender is uninitialized.
169#[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/// Attempts to get the global execution event sender without panicking.
181///
182/// Returns `None` if the sender is not initialized (e.g., in test environments).
183#[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
188/// Sets the global execution event sender.
189///
190/// Can only be called once per thread.
191///
192/// # Panics
193///
194/// Panics if a sender has already been set.
195pub 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
206/// Replaces the global execution event sender for the current thread.
207pub 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}