Skip to main content

nautilus_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//! Async event loop runner for live and sandbox trading nodes.
17//!
18//! `AsyncRunner` owns seven tokio mpsc channel pairs plus a shutdown
19//! signal channel. Construction creates the channels without side
20//! effects. The sender halves are placed into thread-local storage
21//! via [`AsyncRunner::bind_senders`] so that adapters and engine
22//! components can resolve them through the `get_*_sender()` accessors
23//! in `nautilus_common::runner` and `nautilus_common::live::runner`.
24//!
25//! Channel pairs:
26//!
27//! - **Time events**: timer callbacks dispatched by the clock.
28//! - **System events**: system notifications handled by the live node.
29//! - **System commands**: control requests handled by the live node.
30//! - **Execution events**: fills, order updates, and account state from
31//!   execution clients to the execution engine.
32//! - **Trading commands**: deferred order actions routed to their direct endpoint.
33//! - **Data events**: market data from adapters to the data engine.
34//! - **Data commands**: subscribe/unsubscribe requests to data clients.
35//!
36//! Both `AsyncRunner::run` and `LiveNode::run` use a `biased;` select with
37//! system and execution branches polled ahead of data branches. Within each
38//! channel pair, events are polled before commands.
39//!
40//! The runner can drive the event loop in two ways:
41//!
42//! - **Standalone**: call [`AsyncRunner::run`], which binds senders and
43//!   enters a `tokio::select!` loop internally.
44//! - **Integrated**: call [`AsyncRunner::take_channels`] to extract the
45//!   receivers and run the `select!` loop directly inside `LiveNode::run`,
46//!   where it is interleaved with startup, reconciliation, and shutdown
47//!   phases.
48//!
49//! # Invariants
50//!
51//! - `bind_senders` must be called before any code that reads from TLS.
52//!   This includes adapter constructors, clock initialization, and
53//!   execution client start methods. Every path from construction to
54//!   the event loop must bind before the first TLS read.
55//! - The event loop and all TLS consumers must execute on the same
56//!   thread. Senders are cloneable and `Send`, but the `RefCell`-backed
57//!   TLS slots are not accessible from other threads.
58//! - Only one runner at a time should own the TLS slots on a given
59//!   thread. `bind_senders` overwrites any existing TLS contents on the
60//!   thread, so the last caller wins.
61
62use std::{
63    fmt::Debug,
64    sync::Arc,
65    thread::{self, ThreadId},
66};
67
68use nautilus_common::{
69    live::{
70        dispatch::DispatchMessage,
71        runner::{
72            replace_data_event_sender, replace_exec_event_sender, replace_system_command_sender,
73            replace_system_event_sender,
74        },
75        sender::{DispatchSender, EventSender},
76    },
77    messages::{
78        DataEvent, ExecutionEvent, ExecutionReport, SystemCommand, SystemEvent, data::DataCommand,
79        execution::TradingCommand,
80    },
81    msgbus::{self, MessagingSwitchboard},
82    runner::{
83        DataCommandSender, TimeEventMessage, TimeEventSender, TradingCommandMessage,
84        TradingCommandSender, replace_data_cmd_sender, replace_exec_cmd_sender,
85        replace_time_event_sender,
86    },
87};
88use nautilus_model::events::OrderEventAny;
89
90use crate::dispatch::drain_callbacks;
91#[cfg(feature = "node")]
92use crate::node::{LiveNodeHandle, NodeState};
93
94/// Asynchronous implementation of `DataCommandSender` for live environments.
95#[derive(Debug)]
96pub struct AsyncDataCommandSender {
97    cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<DataCommand>>,
98    owner: ThreadId,
99    #[cfg(feature = "node")]
100    node_handle: Option<LiveNodeHandle>,
101}
102
103impl AsyncDataCommandSender {
104    /// Creates a sender owned by the calling thread.
105    ///
106    /// Construct it on the runtime thread so command sends capture callback roots.
107    #[must_use]
108    pub fn new(cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<DataCommand>>) -> Self {
109        Self {
110            cmd_tx,
111            owner: thread::current().id(),
112            #[cfg(feature = "node")]
113            node_handle: None,
114        }
115    }
116}
117
118impl DataCommandSender for AsyncDataCommandSender {
119    fn execute(&self, command: DataCommand) {
120        if let Err(e) = self.cmd_tx.send(DispatchMessage::new(command, self.owner)) {
121            // Disposal releases retained subscriptions after the node drops its receivers
122            #[cfg(feature = "node")]
123            if self
124                .node_handle
125                .as_ref()
126                .is_some_and(|handle| handle.state() == NodeState::Stopped)
127            {
128                return;
129            }
130
131            log::error!("Failed to send data command: {e}");
132        }
133    }
134}
135
136/// Asynchronous implementation of `TimeEventSender` for live environments.
137#[derive(Debug, Clone)]
138pub struct AsyncTimeEventSender {
139    time_tx: DispatchSender<TimeEventMessage>,
140}
141
142impl AsyncTimeEventSender {
143    /// Creates a sender owned by the calling runtime thread.
144    #[must_use]
145    pub fn new(
146        time_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<TimeEventMessage>>,
147    ) -> Self {
148        Self {
149            time_tx: DispatchSender::new(time_tx),
150        }
151    }
152}
153
154impl TimeEventSender for AsyncTimeEventSender {
155    fn send(&self, message: TimeEventMessage) {
156        if let Err(e) = self.time_tx.send(message) {
157            log::error!("Failed to send time event message: {e}");
158        }
159    }
160}
161
162/// Asynchronous implementation of `TradingCommandSender` for live environments.
163#[derive(Debug)]
164pub struct AsyncTradingCommandSender {
165    cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<TradingCommandMessage>>,
166    owner: ThreadId,
167}
168
169impl AsyncTradingCommandSender {
170    /// Creates a sender owned by the calling thread.
171    ///
172    /// Construct it on the runtime thread so command sends capture callback roots.
173    #[must_use]
174    pub fn new(
175        cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<TradingCommandMessage>>,
176    ) -> Self {
177        Self {
178            cmd_tx,
179            owner: thread::current().id(),
180        }
181    }
182}
183
184impl TradingCommandSender for AsyncTradingCommandSender {
185    fn execute(&self, message: TradingCommandMessage) {
186        if let Err(e) = self.cmd_tx.send(DispatchMessage::new(message, self.owner)) {
187            log::error!("Failed to send trading command: {e}");
188        }
189    }
190}
191
192pub trait Runner {
193    fn run(&mut self);
194}
195
196/// Channel receivers for the async event loop.
197///
198/// These can be extracted from `AsyncRunner` via `take_channels()` to drive
199/// the event loop directly on the same thread as the msgbus endpoints.
200#[derive(Debug)]
201pub struct AsyncRunnerChannels {
202    pub time_evt_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<TimeEventMessage>>,
203    pub system_evt_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<SystemEvent>>,
204    pub system_cmd_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<SystemCommand>>,
205    pub exec_evt_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<ExecutionEvent>>,
206    pub exec_cmd_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<TradingCommandMessage>>,
207    pub data_evt_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<DataEvent>>,
208    pub data_cmd_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<DataCommand>>,
209}
210
211#[cfg(feature = "node")]
212#[allow(
213    clippy::large_enum_variant,
214    reason = "runner events are consumed immediately; boxing would add routing allocations"
215)]
216pub(crate) enum PendingRunnerEvent {
217    TimeEvent(DispatchMessage<TimeEventMessage>),
218    SystemEvent(DispatchMessage<SystemEvent>),
219    SystemCommand(DispatchMessage<SystemCommand>),
220    ExecEvent(DispatchMessage<ExecutionEvent>),
221    ExecCommand(DispatchMessage<TradingCommandMessage>),
222    DataEvent(DispatchMessage<DataEvent>),
223    DataCommand(DispatchMessage<DataCommand>),
224}
225
226pub struct AsyncRunner {
227    channels: AsyncRunnerChannels,
228    time_evt_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<TimeEventMessage>>,
229    system_evt_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<SystemEvent>>,
230    system_cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<SystemCommand>>,
231    signal_rx: tokio::sync::mpsc::UnboundedReceiver<()>,
232    signal_tx: tokio::sync::mpsc::UnboundedSender<()>,
233    exec_evt_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<ExecutionEvent>>,
234    exec_cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<TradingCommandMessage>>,
235    data_evt_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<DataEvent>>,
236    data_cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<DataCommand>>,
237}
238
239/// Handle for stopping the `AsyncRunner` from another context.
240#[derive(Clone, Debug)]
241pub struct AsyncRunnerHandle {
242    signal_tx: tokio::sync::mpsc::UnboundedSender<()>,
243}
244
245impl AsyncRunnerHandle {
246    /// Signals the runner to stop.
247    pub fn stop(&self) {
248        if let Err(e) = self.signal_tx.send(()) {
249            log::error!("Failed to send shutdown signal: {e}");
250        }
251    }
252}
253
254impl Default for AsyncRunner {
255    fn default() -> Self {
256        Self::new()
257    }
258}
259
260impl Debug for AsyncRunner {
261    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
262        f.debug_struct(stringify!(AsyncRunner)).finish()
263    }
264}
265
266impl AsyncRunner {
267    /// Creates a new [`AsyncRunner`] instance.
268    ///
269    /// Creates channels but does not bind senders to thread-local storage.
270    /// Call [`bind_senders`](Self::bind_senders) before creating clients that
271    /// read from TLS, and again before entering the event loop.
272    #[must_use]
273    #[rustfmt::skip]
274    pub fn new() -> Self {
275        use tokio::sync::mpsc::unbounded_channel; // tokio-import-ok
276
277        let (time_evt_tx, time_evt_rx) = unbounded_channel::<DispatchMessage<TimeEventMessage>>();
278        let (system_evt_tx, system_evt_rx) = unbounded_channel::<DispatchMessage<SystemEvent>>();
279        let (system_cmd_tx, system_cmd_rx) = unbounded_channel::<DispatchMessage<SystemCommand>>();
280        let (signal_tx, signal_rx) = unbounded_channel::<()>();
281        let (exec_evt_tx, exec_evt_rx) = unbounded_channel::<DispatchMessage<ExecutionEvent>>();
282        let (exec_cmd_tx, exec_cmd_rx) = unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
283        let (data_evt_tx, data_evt_rx) = unbounded_channel::<DispatchMessage<DataEvent>>();
284        let (data_cmd_tx, data_cmd_rx) = unbounded_channel::<DispatchMessage<DataCommand>>();
285
286        Self {
287            channels: AsyncRunnerChannels {
288                time_evt_rx,
289                system_evt_rx,
290                system_cmd_rx,
291                exec_evt_rx,
292                exec_cmd_rx,
293                data_evt_rx,
294                data_cmd_rx,
295            },
296            time_evt_tx,
297            system_evt_tx,
298            system_cmd_tx,
299            signal_rx,
300            signal_tx,
301            exec_evt_tx,
302            exec_cmd_tx,
303            data_evt_tx,
304            data_cmd_tx,
305        }
306    }
307
308    /// Binds this runner's channel senders to thread-local storage.
309    ///
310    /// Call before creating clients that read from TLS (e.g., in the builder),
311    /// and again before entering the event loop to reclaim ownership if another
312    /// runner was constructed on this thread in the interim.
313    pub fn bind_senders(&self) {
314        self.bind_senders_with_data_sender(AsyncDataCommandSender::new(self.data_cmd_tx.clone()));
315    }
316
317    #[cfg(feature = "node")]
318    pub(crate) fn bind_senders_for_node(&self, handle: LiveNodeHandle) {
319        self.bind_senders_with_data_sender(AsyncDataCommandSender {
320            cmd_tx: self.data_cmd_tx.clone(),
321            owner: thread::current().id(),
322            node_handle: Some(handle),
323        });
324    }
325
326    #[rustfmt::skip]
327    fn bind_senders_with_data_sender(&self, sender: AsyncDataCommandSender) {
328        replace_time_event_sender(Arc::new(AsyncTimeEventSender::new(self.time_evt_tx.clone())));
329        replace_system_event_sender(EventSender::new(self.system_evt_tx.clone()));
330        replace_system_command_sender(DispatchSender::new(self.system_cmd_tx.clone()));
331        replace_exec_event_sender(EventSender::new(self.exec_evt_tx.clone()));
332        replace_exec_cmd_sender(Arc::new(AsyncTradingCommandSender::new(self.exec_cmd_tx.clone())));
333        replace_data_event_sender(EventSender::new(self.data_evt_tx.clone()));
334        replace_data_cmd_sender(Arc::new(sender));
335    }
336
337    /// Stops the runner with an internal shutdown signal.
338    pub fn stop(&self) {
339        if let Err(e) = self.signal_tx.send(()) {
340            log::error!("Failed to send shutdown signal: {e}");
341        }
342    }
343
344    /// Returns a handle that can be used to stop the runner from another context.
345    #[must_use]
346    pub fn handle(&self) -> AsyncRunnerHandle {
347        AsyncRunnerHandle {
348            signal_tx: self.signal_tx.clone(),
349        }
350    }
351
352    /// Consumes the runner and returns the channel receivers for direct event loop driving.
353    ///
354    /// This is used when the event loop needs to run on the same thread as the msgbus
355    /// endpoints (which use thread-local storage).
356    #[must_use]
357    pub fn take_channels(self) -> AsyncRunnerChannels {
358        self.channels
359    }
360
361    /// Flushes all pending data events and commands from the channels.
362    ///
363    /// Loops until both data channels are empty, processing each item
364    /// into the cache immediately. Used in `start()` where channels are
365    /// not extracted.
366    pub fn flush_pending_data(&mut self) {
367        let mut total = 0;
368
369        loop {
370            let mut progressed = false;
371
372            // Events drain before commands here even though the runtime select
373            // prefers the opposite for everything-else: `LiveNode::start()`
374            // calls this after `connect_data_clients()` to push queued
375            // `DataEvent::Instrument` items into the cache. A pending
376            // subscription command (e.g. `SubscribeBars`) processed before the
377            // matching instrument lands would be rejected by the data engine.
378            while let Ok(evt) = self.channels.data_evt_rx.try_recv() {
379                Self::dispatch_data_event(evt);
380                progressed = true;
381                total += 1;
382            }
383
384            while let Ok(cmd) = self.channels.data_cmd_rx.try_recv() {
385                Self::handle_data_command(cmd);
386                progressed = true;
387                total += 1;
388            }
389
390            if !progressed {
391                break;
392            }
393        }
394
395        if total > 0 {
396            log::debug!("Flushed {total} pending data events/commands");
397        }
398    }
399
400    #[cfg(feature = "node")]
401    pub(crate) fn drain_pending_system_events(&mut self) -> Vec<DispatchMessage<SystemEvent>> {
402        let mut events = Vec::new();
403
404        while let Ok(event) = self.channels.system_evt_rx.try_recv() {
405            events.push(event);
406        }
407
408        events
409    }
410
411    #[cfg(feature = "node")]
412    pub(crate) fn drain_pending_system_commands(&mut self) -> Vec<DispatchMessage<SystemCommand>> {
413        let mut commands = Vec::new();
414
415        while let Ok(command) = self.channels.system_cmd_rx.try_recv() {
416            commands.push(command);
417        }
418
419        commands
420    }
421
422    /// Runs the async runner event loop.
423    ///
424    /// This method processes time, system, execution, and data events in an async loop.
425    /// It will run until a signal is received or the event streams are closed.
426    ///
427    /// # Errors
428    ///
429    /// Returns a callback dispatch failure, retaining pending channel messages and the failure latch.
430    pub async fn run(&mut self) -> anyhow::Result<()> {
431        self.bind_senders();
432
433        log::info!("AsyncRunner starting");
434
435        loop {
436            let callbacks_pending = drain_callbacks().await?;
437
438            tokio::select! {
439                biased;
440
441                Some(()) = self.signal_rx.recv() => {
442                    log::info!("AsyncRunner received signal, shutting down");
443                    return Ok(());
444                },
445                () = std::future::ready(()), if callbacks_pending => {},
446                Some(handler) = self.channels.time_evt_rx.recv() => {
447                    let _ = Self::handle_time_event(handler);
448                },
449                Some(event) = self.channels.system_evt_rx.recv() => {
450                    log::error!("System event {event} requires the LiveNode runner");
451                },
452                Some(command) = self.channels.system_cmd_rx.recv() => {
453                    log::error!("System command {command} requires the LiveNode runner");
454                },
455                Some(evt) = self.channels.exec_evt_rx.recv() => {
456                    Self::dispatch_exec_event(evt);
457                },
458                Some(cmd) = self.channels.exec_cmd_rx.recv() => {
459                    Self::handle_trading_command(cmd);
460                },
461                Some(evt) = self.channels.data_evt_rx.recv() => {
462                    Self::dispatch_data_event(evt);
463                },
464                Some(cmd) = self.channels.data_cmd_rx.recv() => {
465                    Self::handle_data_command(cmd);
466                },
467                else => {
468                    log::debug!("AsyncRunner all channels closed, exiting");
469                    return Ok(());
470                }
471            };
472        }
473    }
474
475    /// Handles a time event under its originating callback root.
476    ///
477    /// # Panics
478    ///
479    /// Panics if a rooted message is processed outside its owner thread.
480    #[inline]
481    #[must_use]
482    pub fn handle_time_event(message: DispatchMessage<TimeEventMessage>) -> bool {
483        message.dispatch(TimeEventMessage::dispatch)
484    }
485
486    /// Handles a data command under its captured callback root by sending to the `DataEngine`.
487    ///
488    /// # Panics
489    ///
490    /// Panics if a rooted command is processed outside its owner thread.
491    #[inline]
492    pub fn handle_data_command(cmd: DispatchMessage<DataCommand>) {
493        cmd.dispatch(|cmd| {
494            msgbus::send_data_command(MessagingSwitchboard::data_engine_execute(), cmd);
495        });
496    }
497
498    /// Dispatches a received data event under its originating callback root.
499    ///
500    /// # Panics
501    ///
502    /// Panics if a rooted message is processed outside its owner thread.
503    pub fn dispatch_data_event(event: DispatchMessage<DataEvent>) {
504        event.dispatch(Self::handle_data_event);
505    }
506
507    /// Dispatches a received execution event under its originating callback root.
508    ///
509    /// # Panics
510    ///
511    /// Panics if a rooted message is processed outside its owner thread.
512    pub fn dispatch_exec_event(event: DispatchMessage<ExecutionEvent>) {
513        event.dispatch(Self::handle_exec_event);
514    }
515
516    /// Handles a data event by sending to the appropriate `DataEngine` endpoint.
517    #[inline]
518    pub fn handle_data_event(event: DataEvent) {
519        match event {
520            DataEvent::Data(data) => {
521                msgbus::send_data(MessagingSwitchboard::data_engine_process_data(), data);
522            }
523            DataEvent::Instrument(data) => {
524                msgbus::send_any(MessagingSwitchboard::data_engine_process(), &data);
525            }
526            DataEvent::Response(resp) => {
527                msgbus::send_data_response(MessagingSwitchboard::data_engine_response(), resp);
528            }
529            DataEvent::FundingRate(funding_rate) => {
530                msgbus::send_any(MessagingSwitchboard::data_engine_process(), &funding_rate);
531            }
532            DataEvent::InstrumentStatus(status) => {
533                msgbus::send_any(MessagingSwitchboard::data_engine_process(), &status);
534            }
535            DataEvent::OptionGreeks(greeks) => {
536                msgbus::send_any(MessagingSwitchboard::data_engine_process(), &greeks);
537            }
538            #[cfg(feature = "defi")]
539            DataEvent::DeFi(data) => {
540                msgbus::send_defi_data(MessagingSwitchboard::data_engine_process_defi_data(), data);
541            }
542        }
543    }
544
545    /// Dispatches an internal execution command directly to the execution engine.
546    #[inline]
547    pub fn handle_exec_command(cmd: TradingCommand) {
548        msgbus::send_trading_command(MessagingSwitchboard::exec_engine_execute(), cmd);
549    }
550
551    /// Dispatches a trading command and its deferred children under their captured callback roots.
552    ///
553    /// # Panics
554    ///
555    /// Panics if a rooted command is processed outside its owner thread.
556    #[inline]
557    pub fn handle_trading_command(message: DispatchMessage<TradingCommandMessage>) {
558        message.dispatch_trading(|_| {});
559    }
560
561    /// Handles an execution event by sending to the appropriate engine endpoint.
562    #[inline]
563    pub fn handle_exec_event(event: ExecutionEvent) {
564        match event {
565            ExecutionEvent::Order(order_event) => {
566                msgbus::send_order_event(MessagingSwitchboard::exec_engine_process(), order_event);
567            }
568            ExecutionEvent::OrderSubmittedBatch(batch) => {
569                for submitted in batch {
570                    msgbus::send_order_event(
571                        MessagingSwitchboard::exec_engine_process(),
572                        OrderEventAny::Submitted(submitted),
573                    );
574                }
575            }
576            ExecutionEvent::OrderAcceptedBatch(batch) => {
577                for accepted in batch {
578                    msgbus::send_order_event(
579                        MessagingSwitchboard::exec_engine_process(),
580                        OrderEventAny::Accepted(accepted),
581                    );
582                }
583            }
584            ExecutionEvent::OrderCanceledBatch(batch) => {
585                for canceled in batch {
586                    msgbus::send_order_event(
587                        MessagingSwitchboard::exec_engine_process(),
588                        OrderEventAny::Canceled(canceled),
589                    );
590                }
591            }
592            ExecutionEvent::Report(report) => {
593                Self::handle_exec_report(report);
594            }
595            ExecutionEvent::Account(ref account) => {
596                msgbus::send_account_state(
597                    MessagingSwitchboard::portfolio_update_account(),
598                    account,
599                );
600            }
601        }
602    }
603
604    #[inline]
605    pub fn handle_exec_report(report: ExecutionReport) {
606        let endpoint = MessagingSwitchboard::exec_engine_reconcile_execution_report();
607        msgbus::send_execution_report(endpoint, report);
608    }
609}
610
611#[cfg(feature = "node")]
612impl AsyncRunner {
613    pub(crate) fn poll_pending(&mut self, mut process: impl FnMut(PendingRunnerEvent)) -> usize {
614        self.bind_senders();
615
616        let pending = (
617            self.channels.time_evt_rx.len(),
618            self.channels.system_evt_rx.len(),
619            self.channels.system_cmd_rx.len(),
620            self.channels.exec_evt_rx.len(),
621            self.channels.exec_cmd_rx.len(),
622            self.channels.data_evt_rx.len(),
623            self.channels.data_cmd_rx.len(),
624        );
625
626        let mut processed = 0;
627        processed += poll_channel(
628            &mut self.channels.time_evt_rx,
629            pending.0,
630            PendingRunnerEvent::TimeEvent,
631            &mut process,
632        );
633        processed += poll_channel(
634            &mut self.channels.system_evt_rx,
635            pending.1,
636            PendingRunnerEvent::SystemEvent,
637            &mut process,
638        );
639        processed += poll_channel(
640            &mut self.channels.system_cmd_rx,
641            pending.2,
642            PendingRunnerEvent::SystemCommand,
643            &mut process,
644        );
645        processed += poll_channel(
646            &mut self.channels.exec_evt_rx,
647            pending.3,
648            PendingRunnerEvent::ExecEvent,
649            &mut process,
650        );
651        processed += poll_channel(
652            &mut self.channels.exec_cmd_rx,
653            pending.4,
654            PendingRunnerEvent::ExecCommand,
655            &mut process,
656        );
657        processed += poll_channel(
658            &mut self.channels.data_evt_rx,
659            pending.5,
660            PendingRunnerEvent::DataEvent,
661            &mut process,
662        );
663        processed += poll_channel(
664            &mut self.channels.data_cmd_rx,
665            pending.6,
666            PendingRunnerEvent::DataCommand,
667            &mut process,
668        );
669        processed
670    }
671
672    pub(crate) async fn recv(&mut self) -> Option<PendingRunnerEvent> {
673        tokio::select! {
674            biased;
675
676            Some(message) = self.channels.time_evt_rx.recv() => {
677                Some(PendingRunnerEvent::TimeEvent(message))
678            }
679            Some(event) = self.channels.system_evt_rx.recv() => {
680                Some(PendingRunnerEvent::SystemEvent(event))
681            }
682            Some(command) = self.channels.system_cmd_rx.recv() => {
683                Some(PendingRunnerEvent::SystemCommand(command))
684            }
685            Some(event) = self.channels.exec_evt_rx.recv() => {
686                Some(PendingRunnerEvent::ExecEvent(event))
687            }
688            Some(command) = self.channels.exec_cmd_rx.recv() => {
689                Some(PendingRunnerEvent::ExecCommand(command))
690            }
691            Some(event) = self.channels.data_evt_rx.recv() => {
692                Some(PendingRunnerEvent::DataEvent(event))
693            }
694            Some(command) = self.channels.data_cmd_rx.recv() => {
695                Some(PendingRunnerEvent::DataCommand(command))
696            }
697            else => None,
698        }
699    }
700}
701
702#[cfg(feature = "node")]
703fn poll_channel<T>(
704    receiver: &mut tokio::sync::mpsc::UnboundedReceiver<T>,
705    pending: usize,
706    event: impl Fn(T) -> PendingRunnerEvent,
707    process: &mut impl FnMut(PendingRunnerEvent),
708) -> usize {
709    let mut processed = 0;
710
711    for _ in 0..pending {
712        let Ok(message) = receiver.try_recv() else {
713            break;
714        };
715
716        process(event(message));
717        processed += 1;
718    }
719
720    processed
721}
722
723#[cfg(test)]
724mod tests {
725    #[cfg(not(all(feature = "simulation", madsim)))]
726    use std::num::NonZeroU64;
727    use std::{cell::RefCell, rc::Rc, sync::Arc, time::Duration};
728
729    #[cfg(not(all(feature = "simulation", madsim)))]
730    use nautilus_common::live::LiveTimer;
731    use nautilus_common::{
732        actor,
733        cache::Cache,
734        clock::VirtualClock,
735        live::{
736            dispatch::DispatchMessage,
737            runner::{
738                get_data_event_sender, get_exec_event_sender, get_system_command_sender,
739                get_system_event_sender, try_get_system_command_sender,
740                try_get_system_event_sender,
741            },
742        },
743        messages::{
744            ExecutionEvent, ExecutionReport,
745            data::{SubscribeCommand, SubscribeCustomData},
746            execution::{CancelAllOrders, QueryAccount, TradingCommand},
747            system::{ReconnectSocket, SocketState, SocketStateChange},
748        },
749        msgbus::{TypedIntoHandler, stubs::get_typed_into_message_saving_handler},
750        runner::{
751            SyncTradingCommandSender, TimeEventMessage, drain_trading_cmd_queue,
752            get_data_cmd_sender, get_time_event_sender, get_trading_cmd_sender,
753            replace_exec_cmd_sender, try_get_time_event_sender, try_get_trading_cmd_sender,
754        },
755        timer::{TimeEvent, TimeEventCallback},
756    };
757    #[cfg(not(all(feature = "simulation", madsim)))]
758    use nautilus_core::time::get_atomic_clock_realtime;
759    use nautilus_core::{UUID4, UnixNanos};
760    use nautilus_execution::engine::ExecutionEngine;
761    use nautilus_model::{
762        data::{Data, DataType, quote::QuoteTick},
763        enums::{
764            AccountType, LiquiditySide, OrderSide, OrderStatus, OrderType, PositionSide,
765            TimeInForce,
766        },
767        events::{
768            OrderAcceptedBatch, OrderCanceledBatch, OrderEvent, OrderEventAny, OrderSubmittedBatch,
769            account::state::AccountState,
770            order::spec::{OrderAcceptedSpec, OrderCanceledSpec, OrderSubmittedSpec},
771        },
772        identifiers::{
773            AccountId, ClientId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId,
774            TraderId, Venue, VenueOrderId,
775        },
776        reports::{FillReport, OrderStatusReport, PositionStatusReport},
777        types::{Money, Price, Quantity},
778    };
779    use rstest::rstest;
780    use ustr::Ustr;
781
782    use super::*;
783
784    #[rstest]
785    #[case(false)]
786    #[case(true)]
787    #[tokio::test]
788    async fn test_runner_callback_failure_stops_before_next_command(#[case] before_run: bool) {
789        actor::clear_callbacks().unwrap();
790        let mut runner = AsyncRunner::new();
791        let received = Rc::new(RefCell::new(Vec::new()));
792        let observed = received.clone();
793        msgbus::register_trading_command_endpoint(
794            MessagingSwitchboard::risk_engine_execute(),
795            TypedIntoHandler::from(move |command: TradingCommand| {
796                observed.borrow_mut().push(command.ts_init());
797                crate::dispatch::tests::latch_callback_failure();
798            }),
799        );
800
801        let sender = AsyncTradingCommandSender::new(runner.exec_cmd_tx.clone());
802        let seed_endpoint = ustr::Ustr::from("test.callback-root-seed");
803        msgbus::register_trading_command_endpoint(
804            seed_endpoint.into(),
805            TypedIntoHandler::from(move |command: TradingCommand| {
806                sender.execute(TradingCommandMessage::new(
807                    MessagingSwitchboard::risk_engine_execute(),
808                    command,
809                ));
810            }),
811        );
812
813        for timestamp in [17, 23] {
814            SyncTradingCommandSender.execute(TradingCommandMessage::new(
815                seed_endpoint.into(),
816                TradingCommand::QueryAccount(QueryAccount::new(
817                    "TRADER-001".into(),
818                    None,
819                    "SIM-001".into(),
820                    UUID4::new(),
821                    timestamp.into(),
822                    None,
823                    None,
824                )),
825            ));
826        }
827
828        drain_trading_cmd_queue();
829        let first = runner.channels.exec_cmd_rx.try_recv().unwrap();
830        let second = runner.channels.exec_cmd_rx.try_recv().unwrap();
831        assert!(first.is_rooted());
832        assert!(second.is_rooted());
833        runner.exec_cmd_tx.send(first).unwrap();
834        runner.exec_cmd_tx.send(second).unwrap();
835
836        if before_run {
837            crate::dispatch::tests::latch_callback_failure();
838        }
839
840        let error = runner.run().await.unwrap_err();
841
842        assert_eq!(
843            error.downcast_ref::<actor::CallbackDispatchError>(),
844            Some(&actor::CallbackDispatchError::DeliveryUnwound)
845        );
846        assert_eq!(
847            *received.borrow(),
848            if before_run {
849                vec![]
850            } else {
851                vec![UnixNanos::from(17)]
852            }
853        );
854        assert_eq!(
855            runner.channels.exec_cmd_rx.len(),
856            if before_run { 2 } else { 1 }
857        );
858        assert!(!runner.exec_cmd_tx.is_closed());
859        assert_eq!(
860            actor::callback_failure(),
861            Some(actor::CallbackDispatchError::DeliveryUnwound)
862        );
863        drop(runner);
864        assert_eq!(actor::clear_callbacks(), Ok(()));
865    }
866
867    #[tokio::test]
868    async fn test_runner_stop_preserves_messages_and_channel_only_scheduling() {
869        let mut runner = AsyncRunner::new();
870        let received = Rc::new(RefCell::new(Vec::new()));
871        let stop = runner.signal_tx.clone();
872
873        for timestamp in 1..=65 {
874            let received = received.clone();
875            let stop = stop.clone();
876            runner
877                .time_evt_tx
878                .send(
879                    TimeEventMessage::new(
880                        TimeEvent::new(
881                            "runner-resume".into(),
882                            UUID4::new(),
883                            timestamp.into(),
884                            timestamp.into(),
885                        ),
886                        TimeEventCallback::RustLocal(Rc::new(move |_| {
887                            received.borrow_mut().push(timestamp);
888                            if timestamp == 65 {
889                                stop.send(()).unwrap();
890                            }
891                        })),
892                    )
893                    .into(),
894                )
895                .unwrap();
896        }
897
898        stop.send(()).unwrap();
899        runner.run().await.unwrap();
900
901        assert!(received.borrow().is_empty());
902        assert_eq!(runner.channels.time_evt_rx.len(), 65);
903        assert!(!runner.time_evt_tx.is_closed());
904
905        let (result, at_yield) = tokio::join!(biased;
906            runner.run(),
907            async { received.borrow().clone() },
908        );
909        result.unwrap();
910
911        let expected: Vec<u64> = (1..=65).collect();
912        assert_eq!(at_yield, expected);
913        assert_eq!(*received.borrow(), expected);
914        assert_eq!(runner.channels.time_evt_rx.len(), 0);
915        assert!(!runner.time_evt_tx.is_closed());
916    }
917
918    // Test fixture for creating test quotes
919    fn test_quote() -> QuoteTick {
920        QuoteTick {
921            instrument_id: InstrumentId::from("EUR/USD.SIM"),
922            bid_price: Price::from("1.10000"),
923            ask_price: Price::from("1.10001"),
924            bid_size: Quantity::from(1_000_000),
925            ask_size: Quantity::from(1_000_000),
926            ts_event: UnixNanos::default(),
927            ts_init: UnixNanos::default(),
928        }
929    }
930
931    fn test_system_event() -> SystemEvent {
932        SystemEvent::SocketState(SocketStateChange::new(
933            ClientId::from("BINANCE"),
934            Some(Venue::from("BINANCE")),
935            Ustr::from("binance-futures-market-streams"),
936            SocketState::Connected,
937        ))
938    }
939
940    fn test_system_command() -> SystemCommand {
941        SystemCommand::ReconnectSocket(ReconnectSocket::new(
942            TraderId::from("TRADER-001"),
943            ClientId::from("POLYMARKET"),
944            Ustr::from("polymarket-market-streams"),
945            UnixNanos::from(3),
946        ))
947    }
948
949    // Test fixture to create AsyncRunner with manual channels.
950    // Sender halves are dummies (not connected to the test receivers) since
951    // these tests exercise the event loop, not TLS binding.
952    fn create_test_runner(
953        time_evt_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<TimeEventMessage>>,
954        data_evt_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<DataEvent>>,
955        data_cmd_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<DataCommand>>,
956        exec_evt_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<ExecutionEvent>>,
957        exec_cmd_rx: tokio::sync::mpsc::UnboundedReceiver<DispatchMessage<TradingCommandMessage>>,
958        signal_rx: tokio::sync::mpsc::UnboundedReceiver<()>,
959        signal_tx: tokio::sync::mpsc::UnboundedSender<()>,
960    ) -> AsyncRunner {
961        let (time_evt_tx, _) = tokio::sync::mpsc::unbounded_channel();
962        let (system_evt_tx, system_evt_rx) = tokio::sync::mpsc::unbounded_channel();
963        let (system_cmd_tx, system_cmd_rx) = tokio::sync::mpsc::unbounded_channel();
964        let (data_evt_tx, _) = tokio::sync::mpsc::unbounded_channel();
965        let (data_cmd_tx, _) = tokio::sync::mpsc::unbounded_channel();
966        let (exec_evt_tx, _) = tokio::sync::mpsc::unbounded_channel();
967        let (exec_cmd_tx, _) = tokio::sync::mpsc::unbounded_channel();
968
969        AsyncRunner {
970            channels: AsyncRunnerChannels {
971                time_evt_rx,
972                system_evt_rx,
973                system_cmd_rx,
974                exec_evt_rx,
975                exec_cmd_rx,
976                data_evt_rx,
977                data_cmd_rx,
978            },
979            time_evt_tx,
980            system_evt_tx,
981            system_cmd_tx,
982            exec_evt_tx,
983            exec_cmd_tx,
984            data_evt_tx,
985            data_cmd_tx,
986            signal_rx,
987            signal_tx,
988        }
989    }
990
991    #[cfg(feature = "node")]
992    #[rstest]
993    fn test_poll_pending_processes_entry_snapshot_across_channels() {
994        let (time_evt_tx, time_evt_rx) = tokio::sync::mpsc::unbounded_channel();
995        let (data_evt_tx, data_evt_rx) = tokio::sync::mpsc::unbounded_channel();
996        let (data_cmd_tx, data_cmd_rx) = tokio::sync::mpsc::unbounded_channel();
997        let (exec_evt_tx, exec_evt_rx) = tokio::sync::mpsc::unbounded_channel();
998        let (exec_cmd_tx, exec_cmd_rx) = tokio::sync::mpsc::unbounded_channel();
999        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel();
1000
1001        let time_event = TimeEvent::new(
1002            Ustr::from("test"),
1003            UUID4::new(),
1004            UnixNanos::from(1),
1005            UnixNanos::from(2),
1006        );
1007        time_evt_tx
1008            .send(
1009                (TimeEventMessage::new(time_event, TimeEventCallback::from(|_: TimeEvent| {})))
1010                    .into(),
1011            )
1012            .unwrap();
1013        exec_evt_tx
1014            .send(
1015                (ExecutionEvent::Order(OrderEventAny::Submitted(
1016                    OrderSubmittedSpec::builder()
1017                        .client_order_id(ClientOrderId::from("O-POLL-001"))
1018                        .build(),
1019                )))
1020                .into(),
1021            )
1022            .unwrap();
1023        exec_cmd_tx
1024            .send(
1025                TradingCommandMessage::new(
1026                    MessagingSwitchboard::exec_engine_execute(),
1027                    TradingCommand::CancelAllOrders(CancelAllOrders::new(
1028                        TraderId::from("TRADER-001"),
1029                        None,
1030                        StrategyId::from("S-POLL-001"),
1031                        InstrumentId::from("EUR/USD.SIM"),
1032                        Some(OrderSide::Buy),
1033                        UUID4::new(),
1034                        UnixNanos::from(3),
1035                        None,
1036                        None,
1037                    )),
1038                )
1039                .into(),
1040            )
1041            .unwrap();
1042        data_evt_tx
1043            .send((DataEvent::Data(Data::Quote(test_quote()))).into())
1044            .unwrap();
1045        data_cmd_tx
1046            .send(
1047                DataCommand::Subscribe(SubscribeCommand::Data(SubscribeCustomData {
1048                    client_id: Some(ClientId::from("POLL")),
1049                    venue: None,
1050                    data_type: DataType::new("QuoteTick", None, None),
1051                    command_id: UUID4::new(),
1052                    ts_init: UnixNanos::from(4),
1053                    correlation_id: None,
1054                    params: None,
1055                }))
1056                .into(),
1057            )
1058            .unwrap();
1059
1060        let mut runner = create_test_runner(
1061            time_evt_rx,
1062            data_evt_rx,
1063            data_cmd_rx,
1064            exec_evt_rx,
1065            exec_cmd_rx,
1066            signal_rx,
1067            signal_tx,
1068        );
1069        runner.bind_senders();
1070        get_system_command_sender()
1071            .send(test_system_command())
1072            .unwrap();
1073        get_system_event_sender().send(test_system_event()).unwrap();
1074        get_system_event_sender().send(test_system_event()).unwrap();
1075        let mut processed_by_channel = [0; 7];
1076        let mut processed_order = Vec::new();
1077
1078        let first = runner.poll_pending(|event| match event {
1079            PendingRunnerEvent::TimeEvent(_) => {
1080                processed_by_channel[0] += 1;
1081                processed_order.push("time");
1082            }
1083            PendingRunnerEvent::SystemEvent(_) => {
1084                processed_by_channel[1] += 1;
1085                processed_order.push("system_event");
1086            }
1087            PendingRunnerEvent::SystemCommand(_) => {
1088                processed_by_channel[2] += 1;
1089                processed_order.push("system_command");
1090            }
1091            PendingRunnerEvent::ExecEvent(_) => {
1092                processed_by_channel[3] += 1;
1093                processed_order.push("exec_event");
1094            }
1095            PendingRunnerEvent::ExecCommand(_) => {
1096                processed_by_channel[4] += 1;
1097                processed_order.push("exec_command");
1098            }
1099            PendingRunnerEvent::DataEvent(_) => {
1100                processed_by_channel[5] += 1;
1101                processed_order.push("data_event");
1102                data_evt_tx
1103                    .send((DataEvent::Data(Data::Quote(test_quote()))).into())
1104                    .unwrap();
1105            }
1106            PendingRunnerEvent::DataCommand(_) => {
1107                processed_by_channel[6] += 1;
1108                processed_order.push("data_command");
1109            }
1110        });
1111
1112        let second = runner.poll_pending(|event| match event {
1113            PendingRunnerEvent::DataEvent(_) => {
1114                processed_by_channel[5] += 1;
1115                processed_order.push("data_event");
1116            }
1117            _ => panic!("Unexpected runner event"),
1118        });
1119
1120        assert_eq!(first, 8);
1121        assert_eq!(second, 1);
1122        assert_eq!(processed_by_channel, [1, 2, 1, 1, 1, 2, 1]);
1123        assert_eq!(
1124            processed_order,
1125            [
1126                "time",
1127                "system_event",
1128                "system_event",
1129                "system_command",
1130                "exec_event",
1131                "exec_command",
1132                "data_event",
1133                "data_command",
1134                "data_event",
1135            ]
1136        );
1137    }
1138
1139    #[cfg(feature = "node")]
1140    #[tokio::test]
1141    async fn test_recv_processes_system_event_before_command() {
1142        let (_time_evt_tx, time_evt_rx) = tokio::sync::mpsc::unbounded_channel();
1143        let (_data_evt_tx, data_evt_rx) = tokio::sync::mpsc::unbounded_channel();
1144        let (_data_cmd_tx, data_cmd_rx) = tokio::sync::mpsc::unbounded_channel();
1145        let (_exec_evt_tx, exec_evt_rx) = tokio::sync::mpsc::unbounded_channel();
1146        let (_exec_cmd_tx, exec_cmd_rx) = tokio::sync::mpsc::unbounded_channel();
1147        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel();
1148        let mut runner = create_test_runner(
1149            time_evt_rx,
1150            data_evt_rx,
1151            data_cmd_rx,
1152            exec_evt_rx,
1153            exec_cmd_rx,
1154            signal_rx,
1155            signal_tx,
1156        );
1157
1158        runner
1159            .system_cmd_tx
1160            .send((test_system_command()).into())
1161            .unwrap();
1162        runner
1163            .system_evt_tx
1164            .send((test_system_event()).into())
1165            .unwrap();
1166
1167        assert!(matches!(
1168            runner.recv().await,
1169            Some(PendingRunnerEvent::SystemEvent(_))
1170        ));
1171        assert!(matches!(
1172            runner.recv().await,
1173            Some(PendingRunnerEvent::SystemCommand(_))
1174        ));
1175    }
1176
1177    #[rstest]
1178    fn test_async_data_command_sender_creation() {
1179        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
1180        let sender = AsyncDataCommandSender::new(tx);
1181        assert!(format!("{sender:?}").contains("AsyncDataCommandSender"));
1182    }
1183
1184    #[cfg(feature = "node")]
1185    #[rstest]
1186    fn test_data_command_sender_shutdown_logging() {
1187        struct ErrorCapture(std::sync::Mutex<Vec<String>>);
1188
1189        impl log::Log for ErrorCapture {
1190            fn enabled(&self, metadata: &log::Metadata<'_>) -> bool {
1191                metadata.level() == log::Level::Error
1192                    && metadata.target() == "nautilus_live::runner"
1193            }
1194
1195            fn log(&self, record: &log::Record<'_>) {
1196                if self.enabled(record.metadata()) {
1197                    self.0.lock().unwrap().push(record.args().to_string());
1198                }
1199            }
1200
1201            fn flush(&self) {}
1202        }
1203
1204        static ERRORS: ErrorCapture = ErrorCapture(std::sync::Mutex::new(Vec::new()));
1205        log::set_logger(&ERRORS).unwrap();
1206        log::set_max_level(log::LevelFilter::Error);
1207
1208        let runner = AsyncRunner::new();
1209        let handle = LiveNodeHandle::new();
1210        handle.set_starting();
1211        runner.bind_senders_for_node(handle.clone());
1212        let sender = get_data_cmd_sender();
1213        let mut channels = runner.take_channels();
1214
1215        let command = DataCommand::Subscribe(SubscribeCommand::Data(SubscribeCustomData {
1216            client_id: Some(ClientId::from("TEST")),
1217            venue: None,
1218            data_type: DataType::new("QuoteTick", None, None),
1219            command_id: UUID4::new(),
1220            ts_init: UnixNanos::default(),
1221            correlation_id: None,
1222            params: None,
1223        }));
1224
1225        // Stopping and final draining must still deliver commands while the receiver is alive
1226        for stopped in [false, true] {
1227            if stopped {
1228                handle.set_stopped();
1229            } else {
1230                handle.set_shutting_down();
1231            }
1232
1233            sender.execute(command.clone());
1234            assert_eq!(
1235                channels
1236                    .data_cmd_rx
1237                    .try_recv()
1238                    .unwrap()
1239                    .dispatch(|command| command),
1240                command
1241            );
1242        }
1243
1244        drop(channels);
1245        sender.execute(command.clone());
1246        assert_eq!(*ERRORS.0.lock().unwrap(), Vec::<String>::new());
1247
1248        // A stopped previous node must not hide an unexpected closure in its replacement
1249        let runner = AsyncRunner::new();
1250        let handle = LiveNodeHandle::new();
1251        runner.bind_senders_for_node(handle.clone());
1252        let sender = get_data_cmd_sender();
1253        drop(runner);
1254
1255        for shutting_down in [false, true] {
1256            if shutting_down {
1257                handle.set_shutting_down();
1258            } else {
1259                handle.set_starting();
1260            }
1261
1262            sender.execute(command.clone());
1263        }
1264
1265        assert_eq!(
1266            *ERRORS.0.lock().unwrap(),
1267            vec!["Failed to send data command: channel closed"; 2],
1268        );
1269
1270        ERRORS.0.lock().unwrap().clear();
1271        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1272        drop(rx);
1273        AsyncDataCommandSender::new(tx).execute(command);
1274        assert_eq!(
1275            *ERRORS.0.lock().unwrap(),
1276            vec!["Failed to send data command: channel closed"],
1277        );
1278    }
1279
1280    #[rstest]
1281    fn test_async_time_event_sender_creation() {
1282        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
1283        let sender = AsyncTimeEventSender::new(tx);
1284        assert!(format!("{sender:?}").contains("AsyncTimeEventSender"));
1285    }
1286
1287    #[tokio::test]
1288    async fn test_async_data_command_sender_execute() {
1289        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1290        let sender = AsyncDataCommandSender::new(tx);
1291
1292        let command = DataCommand::Subscribe(SubscribeCommand::Data(SubscribeCustomData {
1293            client_id: Some(ClientId::from("TEST")),
1294            venue: None,
1295            data_type: DataType::new("QuoteTick", None, None),
1296            command_id: UUID4::new(),
1297            ts_init: UnixNanos::default(),
1298            correlation_id: None,
1299            params: None,
1300        }));
1301
1302        sender.execute(command.clone());
1303
1304        let received = rx.recv().await.unwrap().dispatch(|command| command);
1305        match (received, command) {
1306            (
1307                DataCommand::Subscribe(SubscribeCommand::Data(r)),
1308                DataCommand::Subscribe(SubscribeCommand::Data(c)),
1309            ) => {
1310                assert_eq!(r.client_id, c.client_id);
1311                assert_eq!(r.data_type, c.data_type);
1312            }
1313            _ => panic!("Command mismatch"),
1314        }
1315    }
1316
1317    #[rstest]
1318    #[case(false)]
1319    #[case(true)]
1320    fn test_async_command_senders_capture_owner_root(#[case] rooted: bool) {
1321        let runner = AsyncRunner::new();
1322        runner.bind_senders();
1323        let mut channels = runner.take_channels();
1324
1325        let data = DataCommand::Subscribe(SubscribeCommand::Data(SubscribeCustomData {
1326            client_id: Some(ClientId::from("ROOT")),
1327            venue: None,
1328            data_type: DataType::new("QuoteTick", None, None),
1329            command_id: UUID4::new(),
1330            ts_init: UnixNanos::from(7),
1331            correlation_id: None,
1332            params: None,
1333        }));
1334
1335        let trading = TradingCommand::CancelAllOrders(CancelAllOrders::new(
1336            TraderId::from("TRADER-001"),
1337            None,
1338            StrategyId::from("ROOT-001"),
1339            InstrumentId::from("EUR/USD.SIM"),
1340            Some(OrderSide::Sell),
1341            UUID4::new(),
1342            UnixNanos::from(11),
1343            None,
1344            None,
1345        ));
1346        let expected_data = data.clone();
1347        let expected_trading = trading.clone();
1348        let trigger = trading.clone();
1349
1350        let send = move || {
1351            get_data_cmd_sender().execute(data.clone());
1352            get_trading_cmd_sender().execute(TradingCommandMessage::new(
1353                MessagingSwitchboard::exec_engine_execute(),
1354                trading.clone(),
1355            ));
1356        };
1357
1358        if rooted {
1359            msgbus::register_trading_command_endpoint(
1360                MessagingSwitchboard::risk_engine_execute(),
1361                TypedIntoHandler::from(move |_| send()),
1362            );
1363            SyncTradingCommandSender.execute(TradingCommandMessage::new(
1364                MessagingSwitchboard::risk_engine_execute(),
1365                trigger,
1366            ));
1367            drain_trading_cmd_queue();
1368        } else {
1369            send();
1370        }
1371
1372        let data = channels.data_cmd_rx.try_recv().unwrap();
1373        let trading = channels.exec_cmd_rx.try_recv().unwrap();
1374
1375        let results = thread::spawn(move || {
1376            let data = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1377                data.dispatch(|command| command)
1378            }));
1379
1380            let trading = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1381                trading.dispatch(|message| (message.endpoint(), message.command().clone()))
1382            }));
1383
1384            (data, trading)
1385        })
1386        .join()
1387        .unwrap();
1388
1389        if rooted {
1390            for error in [results.0.unwrap_err(), results.1.unwrap_err()] {
1391                let message = error.downcast_ref::<String>().unwrap();
1392                assert!(message.contains("command context dispatched outside its owner thread"));
1393            }
1394        } else {
1395            assert_eq!(results.0.unwrap(), expected_data);
1396            assert_eq!(
1397                results.1.unwrap(),
1398                (
1399                    MessagingSwitchboard::exec_engine_execute(),
1400                    expected_trading
1401                ),
1402            );
1403        }
1404
1405        assert!(channels.data_cmd_rx.is_empty());
1406        assert!(channels.exec_cmd_rx.is_empty());
1407    }
1408
1409    #[rstest]
1410    #[case(false)]
1411    #[case(true)]
1412    fn test_event_senders_capture_owner_root(
1413        #[case] rooted: bool,
1414        #[values(false, true)] dispatch: bool,
1415    ) {
1416        let runner = AsyncRunner::new();
1417        runner.bind_senders();
1418        let mut channels = runner.take_channels();
1419        let expected_quote = test_quote();
1420        let expected_order = OrderSubmittedSpec::builder()
1421            .client_order_id(ClientOrderId::from("EVENT-017"))
1422            .build();
1423        let order = expected_order.clone();
1424
1425        let send = move || {
1426            get_data_event_sender()
1427                .send(DataEvent::Data(Data::Quote(expected_quote)))
1428                .unwrap();
1429            get_exec_event_sender()
1430                .send(ExecutionEvent::Order(OrderEventAny::Submitted(
1431                    order.clone(),
1432                )))
1433                .unwrap();
1434        };
1435
1436        if rooted {
1437            msgbus::register_trading_command_endpoint(
1438                MessagingSwitchboard::risk_engine_execute(),
1439                TypedIntoHandler::from(move |_| send()),
1440            );
1441            SyncTradingCommandSender.execute(TradingCommandMessage::new(
1442                MessagingSwitchboard::risk_engine_execute(),
1443                TradingCommand::QueryAccount(QueryAccount::new(
1444                    "TRADER-001".into(),
1445                    None,
1446                    "SIM-001".into(),
1447                    UUID4::new(),
1448                    19.into(),
1449                    None,
1450                    None,
1451                )),
1452            ));
1453            drain_trading_cmd_queue();
1454        } else {
1455            send();
1456        }
1457
1458        if dispatch {
1459            msgbus::register_data_endpoint(
1460                MessagingSwitchboard::data_engine_process_data(),
1461                TypedIntoHandler::from(move |data| {
1462                    assert_eq!(data, Data::Quote(expected_quote));
1463                    get_data_event_sender().send(DataEvent::Data(data)).unwrap();
1464                }),
1465            );
1466
1467            let order = expected_order.clone();
1468            msgbus::register_order_event_endpoint(
1469                MessagingSwitchboard::exec_engine_process(),
1470                TypedIntoHandler::from(move |event| {
1471                    assert_eq!(event, OrderEventAny::Submitted(order.clone()));
1472                    get_exec_event_sender()
1473                        .send(ExecutionEvent::Order(event))
1474                        .unwrap();
1475                }),
1476            );
1477
1478            AsyncRunner::dispatch_data_event(channels.data_evt_rx.try_recv().unwrap());
1479            AsyncRunner::dispatch_exec_event(channels.exec_evt_rx.try_recv().unwrap());
1480        }
1481
1482        let data = channels.data_evt_rx.try_recv().unwrap();
1483        let exec = channels.exec_evt_rx.try_recv().unwrap();
1484
1485        let (data, exec) = thread::spawn(move || {
1486            (
1487                std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1488                    data.dispatch(|event| event)
1489                })),
1490                std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1491                    exec.dispatch(|event| event)
1492                })),
1493            )
1494        })
1495        .join()
1496        .unwrap();
1497
1498        if rooted || dispatch {
1499            for error in [data.unwrap_err(), exec.unwrap_err()] {
1500                assert!(
1501                    error
1502                        .downcast_ref::<String>()
1503                        .unwrap()
1504                        .contains("command context dispatched outside its owner thread")
1505                );
1506            }
1507        } else {
1508            let DataEvent::Data(Data::Quote(quote)) = data.unwrap() else {
1509                panic!("expected quote")
1510            };
1511
1512            let ExecutionEvent::Order(OrderEventAny::Submitted(order)) = exec.unwrap() else {
1513                panic!("expected submitted order")
1514            };
1515
1516            assert_eq!(quote, expected_quote);
1517            assert_eq!(order, expected_order);
1518        }
1519
1520        assert!(channels.data_evt_rx.is_empty());
1521        assert!(channels.exec_evt_rx.is_empty());
1522    }
1523
1524    #[rstest]
1525    fn test_system_and_time_senders_capture_owner_root(#[values(false, true)] rooted: bool) {
1526        let runner = AsyncRunner::new();
1527        runner.bind_senders();
1528        let mut channels = runner.take_channels();
1529        let event = TimeEvent::new("system-time".into(), UUID4::new(), 17.into(), 23.into());
1530        let observed = Rc::new(RefCell::new(Vec::new()));
1531        let received = observed.clone();
1532        let expected_event = event.clone();
1533
1534        let callback = TimeEventCallback::RustLocal(Rc::new(move |event| {
1535            received.borrow_mut().push(event);
1536            get_system_command_sender()
1537                .send(test_system_command())
1538                .unwrap();
1539        }));
1540
1541        let send = move || {
1542            get_system_event_sender().send(test_system_event()).unwrap();
1543            get_system_command_sender()
1544                .send(test_system_command())
1545                .unwrap();
1546            get_time_event_sender().send(TimeEventMessage::new(event.clone(), callback.clone()));
1547        };
1548
1549        if rooted {
1550            msgbus::register_trading_command_endpoint(
1551                MessagingSwitchboard::risk_engine_execute(),
1552                TypedIntoHandler::from(move |_| send()),
1553            );
1554            SyncTradingCommandSender.execute(TradingCommandMessage::new(
1555                MessagingSwitchboard::risk_engine_execute(),
1556                TradingCommand::QueryAccount(QueryAccount::new(
1557                    "TRADER-001".into(),
1558                    None,
1559                    "SIM-001".into(),
1560                    UUID4::new(),
1561                    31.into(),
1562                    None,
1563                    None,
1564                )),
1565            ));
1566            drain_trading_cmd_queue();
1567        } else {
1568            send();
1569        }
1570
1571        let system_event = channels.system_evt_rx.try_recv().unwrap();
1572        let system_command = channels.system_cmd_rx.try_recv().unwrap();
1573        let time_event = channels.time_evt_rx.try_recv().unwrap();
1574        assert_eq!(system_event.is_rooted(), rooted);
1575        assert_eq!(system_command.is_rooted(), rooted);
1576        assert_eq!(time_event.is_rooted(), rooted);
1577        assert!(AsyncRunner::handle_time_event(time_event));
1578        let child = channels.system_cmd_rx.try_recv().unwrap();
1579
1580        assert_eq!(system_event.dispatch(|event| event), test_system_event());
1581        assert_eq!(
1582            system_command.dispatch(|command| command),
1583            test_system_command()
1584        );
1585        assert_eq!(*observed.borrow(), vec![expected_event]);
1586        assert!(child.is_rooted());
1587        assert_eq!(child.dispatch(|command| command), test_system_command());
1588        assert!(channels.time_evt_rx.is_empty());
1589        assert!(channels.system_evt_rx.is_empty());
1590        assert!(channels.system_cmd_rx.is_empty());
1591    }
1592
1593    #[cfg(not(all(feature = "simulation", madsim)))]
1594    #[tokio::test]
1595    async fn test_live_timer_started_under_root_emits_independent_event() {
1596        let runner = AsyncRunner::new();
1597        runner.bind_senders();
1598        let mut channels = runner.take_channels();
1599        let now = get_atomic_clock_realtime().get_time_ns();
1600        let name = Ustr::from("ROOTED_TIMER");
1601        let owner = thread::current().id();
1602        let observed = Rc::new(RefCell::new(Vec::new()));
1603        let received = observed.clone();
1604
1605        let callback = TimeEventCallback::RustLocal(Rc::new(move |event| {
1606            received
1607                .borrow_mut()
1608                .push((event.name, event.ts_event, thread::current().id()));
1609        }));
1610
1611        let timer = Rc::new(RefCell::new(None));
1612        let started = timer.clone();
1613        msgbus::register_trading_command_endpoint(
1614            MessagingSwitchboard::risk_engine_execute(),
1615            TypedIntoHandler::from(move |_| {
1616                get_system_command_sender()
1617                    .send(test_system_command())
1618                    .unwrap();
1619
1620                let mut timer = LiveTimer::new(
1621                    name,
1622                    NonZeroU64::new(1_000_000).unwrap(),
1623                    now,
1624                    Some(now),
1625                    callback.clone(),
1626                    true,
1627                    Some(get_time_event_sender()),
1628                );
1629                timer.start();
1630                *started.borrow_mut() = Some(timer);
1631            }),
1632        );
1633
1634        SyncTradingCommandSender.execute(TradingCommandMessage::new(
1635            MessagingSwitchboard::risk_engine_execute(),
1636            TradingCommand::QueryAccount(QueryAccount::new(
1637                "TRADER-001".into(),
1638                None,
1639                "SIM-001".into(),
1640                UUID4::new(),
1641                31.into(),
1642                None,
1643                None,
1644            )),
1645        ));
1646        drain_trading_cmd_queue();
1647        let witness = channels.system_cmd_rx.try_recv().unwrap();
1648        let message = tokio::time::timeout(Duration::from_secs(2), channels.time_evt_rx.recv())
1649            .await
1650            .unwrap()
1651            .unwrap();
1652        timer.borrow_mut().take().unwrap().cancel();
1653
1654        assert!(witness.is_rooted());
1655        assert!(!message.is_rooted());
1656        assert!(observed.borrow().is_empty());
1657        assert!(AsyncRunner::handle_time_event(message));
1658        assert_eq!(*observed.borrow(), vec![(name, now, owner)]);
1659        assert!(channels.time_evt_rx.is_empty());
1660    }
1661
1662    #[tokio::test]
1663    async fn test_async_time_event_sender_send() {
1664        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1665        let sender = AsyncTimeEventSender::new(tx);
1666
1667        let event = TimeEvent::new(
1668            Ustr::from("test"),
1669            UUID4::new(),
1670            UnixNanos::from(1),
1671            UnixNanos::from(2),
1672        );
1673        let callback = TimeEventCallback::from(|_: TimeEvent| {});
1674        let message = TimeEventMessage::new(event, callback);
1675
1676        sender.send(message);
1677
1678        assert!(rx.recv().await.is_some());
1679    }
1680
1681    #[tokio::test]
1682    async fn test_runner_shutdown_signal() {
1683        // Create runner with manual channels to avoid global state
1684        let (_data_tx, data_evt_rx) =
1685            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
1686        let (_cmd_tx, data_cmd_rx) =
1687            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
1688        let (_time_tx, time_evt_rx) =
1689            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
1690        let (_exec_evt_tx, exec_evt_rx) =
1691            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
1692        let (_exec_cmd_tx, exec_cmd_rx) =
1693            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
1694        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
1695
1696        let mut runner = create_test_runner(
1697            time_evt_rx,
1698            data_evt_rx,
1699            data_cmd_rx,
1700            exec_evt_rx,
1701            exec_cmd_rx,
1702            signal_rx,
1703            signal_tx.clone(),
1704        );
1705
1706        // Start runner
1707        let runner_handle = tokio::spawn(async move {
1708            runner.run().await.unwrap();
1709        });
1710
1711        // Send shutdown signal
1712        signal_tx.send(()).unwrap();
1713
1714        // Runner should stop quickly
1715        let result = tokio::time::timeout(Duration::from_millis(100), runner_handle).await;
1716        assert!(result.is_ok(), "Runner should stop on signal");
1717    }
1718
1719    #[tokio::test]
1720    async fn test_runner_closes_on_channel_drop() {
1721        let (data_tx, data_evt_rx) =
1722            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
1723        let (_cmd_tx, data_cmd_rx) =
1724            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
1725        let (_time_tx, time_evt_rx) =
1726            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
1727        let (_exec_evt_tx, exec_evt_rx) =
1728            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
1729        let (_exec_cmd_tx, exec_cmd_rx) =
1730            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
1731        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
1732
1733        let mut runner = create_test_runner(
1734            time_evt_rx,
1735            data_evt_rx,
1736            data_cmd_rx,
1737            exec_evt_rx,
1738            exec_cmd_rx,
1739            signal_rx,
1740            signal_tx.clone(),
1741        );
1742
1743        // Start runner
1744        let runner_handle = tokio::spawn(async move {
1745            runner.run().await.unwrap();
1746        });
1747
1748        drop(data_tx);
1749
1750        // Yield to let runner enter event loop before stop signal
1751        tokio::task::yield_now().await;
1752        signal_tx.send(()).ok();
1753
1754        // Runner should stop when channels close or on signal
1755        let result = tokio::time::timeout(Duration::from_millis(200), runner_handle).await;
1756        assert!(
1757            result.is_ok(),
1758            "Runner should stop when channels close or on signal"
1759        );
1760    }
1761
1762    #[tokio::test]
1763    async fn test_concurrent_event_sending() {
1764        let (data_evt_tx, data_evt_rx) =
1765            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
1766        let (_data_cmd_tx, data_cmd_rx) =
1767            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
1768        let (_time_evt_tx, time_evt_rx) =
1769            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
1770        let (_exec_evt_tx, exec_evt_rx) =
1771            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
1772        let (_exec_cmd_tx, exec_cmd_rx) =
1773            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
1774        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
1775
1776        // Setup runner
1777        let mut runner = create_test_runner(
1778            time_evt_rx,
1779            data_evt_rx,
1780            data_cmd_rx,
1781            exec_evt_rx,
1782            exec_cmd_rx,
1783            signal_rx,
1784            signal_tx.clone(),
1785        );
1786
1787        // Spawn multiple concurrent senders
1788        let mut handles = vec![];
1789
1790        for _ in 0..5 {
1791            let tx_clone = data_evt_tx.clone();
1792
1793            let handle = tokio::spawn(async move {
1794                for _ in 0..20 {
1795                    let quote = test_quote();
1796                    tx_clone
1797                        .send(DataEvent::Data(Data::Quote(quote)).into())
1798                        .unwrap();
1799                    tokio::task::yield_now().await;
1800                }
1801            });
1802
1803            handles.push(handle);
1804        }
1805
1806        // Start runner in background
1807        let runner_handle = tokio::spawn(async move {
1808            runner.run().await.unwrap();
1809        });
1810
1811        // Wait for all senders
1812        for handle in handles {
1813            handle.await.unwrap();
1814        }
1815
1816        // Yield to let runner enter event loop before stop signal
1817        tokio::task::yield_now().await;
1818        signal_tx.send(()).unwrap();
1819
1820        let _ = tokio::time::timeout(Duration::from_millis(200), runner_handle).await;
1821    }
1822
1823    #[rstest]
1824    #[case(10)]
1825    #[case(100)]
1826    #[case(1000)]
1827    fn test_channel_send_performance(#[case] count: usize) {
1828        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
1829        let quote = test_quote();
1830
1831        // Send events
1832        for _ in 0..count {
1833            tx.send(DataEvent::Data(Data::Quote(quote))).unwrap();
1834        }
1835
1836        // Verify all received
1837        let mut received = 0;
1838        while rx.try_recv().is_ok() {
1839            received += 1;
1840        }
1841
1842        assert_eq!(received, count);
1843    }
1844
1845    #[rstest]
1846    fn test_async_trading_command_sender_creation() {
1847        let (tx, _rx) = tokio::sync::mpsc::unbounded_channel();
1848        let sender = AsyncTradingCommandSender::new(tx);
1849        assert!(format!("{sender:?}").contains("AsyncTradingCommandSender"));
1850    }
1851
1852    #[rstest]
1853    fn test_async_trading_command_sender_preserves_target_endpoints() {
1854        std::thread::spawn(|| {
1855            msgbus::get_message_bus().borrow_mut().dispose();
1856            let (risk_handler, risk_saving_handler) =
1857                get_typed_into_message_saving_handler::<TradingCommand>(Some(Ustr::from(
1858                    "RiskEngine.execute",
1859                )));
1860            msgbus::register_trading_command_endpoint(
1861                MessagingSwitchboard::risk_engine_execute(),
1862                risk_handler,
1863            );
1864            let (exec_handler, exec_saving_handler) =
1865                get_typed_into_message_saving_handler::<TradingCommand>(Some(Ustr::from(
1866                    "ExecEngine.execute",
1867                )));
1868            msgbus::register_trading_command_endpoint(
1869                MessagingSwitchboard::exec_engine_execute(),
1870                exec_handler,
1871            );
1872
1873            let (tx, mut rx) =
1874                tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
1875            let sender = AsyncTradingCommandSender::new(tx);
1876            sender.execute(TradingCommandMessage::new(
1877                MessagingSwitchboard::risk_engine_execute(),
1878                TradingCommand::CancelAllOrders(CancelAllOrders::new(
1879                    TraderId::from("TRADER-001"),
1880                    None,
1881                    StrategyId::from("RISK-001"),
1882                    InstrumentId::from("EUR/USD.SIM"),
1883                    Some(OrderSide::Buy),
1884                    UUID4::new(),
1885                    UnixNanos::default(),
1886                    None,
1887                    None,
1888                )),
1889            ));
1890            sender.execute(TradingCommandMessage::new(
1891                MessagingSwitchboard::exec_engine_execute(),
1892                TradingCommand::CancelAllOrders(CancelAllOrders::new(
1893                    TraderId::from("TRADER-001"),
1894                    None,
1895                    StrategyId::from("EXEC-001"),
1896                    InstrumentId::from("EUR/USD.SIM"),
1897                    Some(OrderSide::Sell),
1898                    UUID4::new(),
1899                    UnixNanos::default(),
1900                    None,
1901                    None,
1902                )),
1903            ));
1904
1905            AsyncRunner::handle_trading_command(rx.try_recv().unwrap());
1906            AsyncRunner::handle_trading_command(rx.try_recv().unwrap());
1907
1908            let risk_commands = risk_saving_handler.get_messages();
1909            let exec_commands = exec_saving_handler.get_messages();
1910            assert!(rx.try_recv().is_err());
1911            assert_eq!(risk_commands.len(), 1);
1912            assert_eq!(
1913                risk_commands[0].strategy_id(),
1914                Some(StrategyId::from("RISK-001"))
1915            );
1916            assert_eq!(exec_commands.len(), 1);
1917            assert_eq!(
1918                exec_commands[0].strategy_id(),
1919                Some(StrategyId::from("EXEC-001"))
1920            );
1921        })
1922        .join()
1923        .unwrap();
1924    }
1925
1926    #[rstest]
1927    fn test_async_runner_preserves_deferred_follow_up_order() {
1928        std::thread::spawn(|| {
1929            msgbus::get_message_bus().borrow_mut().dispose();
1930            let clock = Rc::new(RefCell::new(VirtualClock::new()));
1931            let cache = Rc::new(RefCell::new(Cache::default()));
1932            let exec_engine = Rc::new(RefCell::new(ExecutionEngine::new(clock, cache, None)));
1933            ExecutionEngine::register_msgbus_handlers(&exec_engine);
1934            msgbus::register_trading_command_endpoint(
1935                MessagingSwitchboard::risk_engine_execute(),
1936                TypedIntoHandler::from(|command: TradingCommand| {
1937                    msgbus::send_trading_command(
1938                        MessagingSwitchboard::exec_engine_queue_execute(),
1939                        command,
1940                    );
1941                }),
1942            );
1943
1944            let (exec_handler, exec_saving_handler) =
1945                get_typed_into_message_saving_handler::<TradingCommand>(Some(Ustr::from(
1946                    "ExecEngine.execute",
1947                )));
1948            msgbus::register_trading_command_endpoint(
1949                MessagingSwitchboard::exec_engine_execute(),
1950                exec_handler,
1951            );
1952
1953            let (tx, mut rx) =
1954                tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
1955            let sender = Arc::new(AsyncTradingCommandSender::new(tx));
1956            replace_exec_cmd_sender(sender.clone());
1957            sender.execute(TradingCommandMessage::new(
1958                MessagingSwitchboard::risk_engine_execute(),
1959                TradingCommand::CancelAllOrders(CancelAllOrders::new(
1960                    TraderId::from("TRADER-001"),
1961                    None,
1962                    StrategyId::from("FIRST-001"),
1963                    InstrumentId::from("EUR/USD.SIM"),
1964                    Some(OrderSide::Buy),
1965                    UUID4::new(),
1966                    UnixNanos::default(),
1967                    None,
1968                    None,
1969                )),
1970            ));
1971            sender.execute(TradingCommandMessage::new(
1972                MessagingSwitchboard::exec_engine_execute(),
1973                TradingCommand::CancelAllOrders(CancelAllOrders::new(
1974                    TraderId::from("TRADER-001"),
1975                    None,
1976                    StrategyId::from("SECOND-001"),
1977                    InstrumentId::from("EUR/USD.SIM"),
1978                    Some(OrderSide::Sell),
1979                    UUID4::new(),
1980                    UnixNanos::default(),
1981                    None,
1982                    None,
1983                )),
1984            ));
1985
1986            AsyncRunner::handle_trading_command(rx.try_recv().unwrap());
1987            AsyncRunner::handle_trading_command(rx.try_recv().unwrap());
1988
1989            let commands = exec_saving_handler.get_messages();
1990            let strategy_ids = commands
1991                .iter()
1992                .map(TradingCommand::strategy_id)
1993                .collect::<Vec<_>>();
1994            assert!(rx.try_recv().is_err());
1995            assert_eq!(commands.len(), 2);
1996            assert_eq!(
1997                strategy_ids,
1998                vec![
1999                    Some(StrategyId::from("FIRST-001")),
2000                    Some(StrategyId::from("SECOND-001"))
2001                ]
2002            );
2003        })
2004        .join()
2005        .unwrap();
2006    }
2007
2008    #[rstest]
2009    fn test_async_runner_dispatches_deferred_exec_command_once() {
2010        std::thread::spawn(|| {
2011            msgbus::get_message_bus().borrow_mut().dispose();
2012            let clock = Rc::new(RefCell::new(VirtualClock::new()));
2013            let cache = Rc::new(RefCell::new(Cache::default()));
2014            let exec_engine = Rc::new(RefCell::new(ExecutionEngine::new(clock, cache, None)));
2015            ExecutionEngine::register_msgbus_handlers(&exec_engine);
2016
2017            let (tx, mut rx) =
2018                tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
2019            replace_exec_cmd_sender(Arc::new(AsyncTradingCommandSender::new(tx)));
2020            let command = TradingCommand::CancelAllOrders(CancelAllOrders::new(
2021                TraderId::from("TRADER-001"),
2022                None,
2023                StrategyId::from("EXEC-001"),
2024                InstrumentId::from("EUR/USD.SIM"),
2025                Some(OrderSide::Buy),
2026                UUID4::new(),
2027                UnixNanos::default(),
2028                None,
2029                None,
2030            ));
2031
2032            msgbus::send_trading_command(
2033                MessagingSwitchboard::exec_engine_queue_execute(),
2034                command,
2035            );
2036            assert_eq!(exec_engine.borrow().command_count(), 0);
2037
2038            AsyncRunner::handle_trading_command(rx.try_recv().unwrap());
2039
2040            assert!(rx.try_recv().is_err());
2041            assert_eq!(exec_engine.borrow().command_count(), 1);
2042        })
2043        .join()
2044        .unwrap();
2045    }
2046
2047    #[tokio::test]
2048    async fn test_runner_processes_trading_commands() {
2049        let (_data_evt_tx, data_evt_rx) =
2050            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
2051        let (_data_cmd_tx, data_cmd_rx) =
2052            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
2053        let (_time_evt_tx, time_evt_rx) =
2054            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
2055        let (_exec_evt_tx, exec_evt_rx) =
2056            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
2057        let (exec_cmd_tx, exec_cmd_rx) =
2058            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
2059        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
2060
2061        let mut runner = create_test_runner(
2062            time_evt_rx,
2063            data_evt_rx,
2064            data_cmd_rx,
2065            exec_evt_rx,
2066            exec_cmd_rx,
2067            signal_rx,
2068            signal_tx.clone(),
2069        );
2070
2071        let runner_handle = tokio::spawn(async move {
2072            runner.run().await.unwrap();
2073        });
2074
2075        let command = TradingCommand::CancelAllOrders(CancelAllOrders::new(
2076            TraderId::from("TRADER-001"),
2077            None,
2078            StrategyId::from("S-001"),
2079            InstrumentId::from("EUR/USD.SIM"),
2080            Some(OrderSide::Buy),
2081            UUID4::new(),
2082            UnixNanos::default(),
2083            None,
2084            None, // correlation_id
2085        ));
2086        exec_cmd_tx
2087            .send(
2088                TradingCommandMessage::new(MessagingSwitchboard::exec_engine_execute(), command)
2089                    .into(),
2090            )
2091            .unwrap();
2092
2093        tokio::task::yield_now().await;
2094        signal_tx.send(()).unwrap();
2095
2096        let result = tokio::time::timeout(Duration::from_millis(100), runner_handle).await;
2097        assert!(result.is_ok(), "Runner should process command and stop");
2098    }
2099
2100    #[tokio::test]
2101    async fn test_runner_processes_multiple_trading_commands() {
2102        let (_data_evt_tx, data_evt_rx) =
2103            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
2104        let (_data_cmd_tx, data_cmd_rx) =
2105            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
2106        let (_time_evt_tx, time_evt_rx) =
2107            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
2108        let (_exec_evt_tx, exec_evt_rx) =
2109            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
2110        let (exec_cmd_tx, exec_cmd_rx) =
2111            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
2112        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
2113
2114        let mut runner = create_test_runner(
2115            time_evt_rx,
2116            data_evt_rx,
2117            data_cmd_rx,
2118            exec_evt_rx,
2119            exec_cmd_rx,
2120            signal_rx,
2121            signal_tx.clone(),
2122        );
2123
2124        let runner_handle = tokio::spawn(async move {
2125            runner.run().await.unwrap();
2126        });
2127
2128        for i in 0..10 {
2129            let strategy_id = format!("S-{i:03}");
2130            let command = TradingCommand::CancelAllOrders(CancelAllOrders::new(
2131                TraderId::from("TRADER-001"),
2132                None,
2133                StrategyId::from(strategy_id.as_str()),
2134                InstrumentId::from("EUR/USD.SIM"),
2135                Some(OrderSide::Buy),
2136                UUID4::new(),
2137                UnixNanos::default(),
2138                None,
2139                None, // correlation_id
2140            ));
2141            exec_cmd_tx
2142                .send(
2143                    TradingCommandMessage::new(
2144                        MessagingSwitchboard::exec_engine_execute(),
2145                        command,
2146                    )
2147                    .into(),
2148                )
2149                .unwrap();
2150        }
2151
2152        tokio::task::yield_now().await;
2153        signal_tx.send(()).unwrap();
2154
2155        let result = tokio::time::timeout(Duration::from_millis(100), runner_handle).await;
2156        assert!(
2157            result.is_ok(),
2158            "Runner should process all commands and stop"
2159        );
2160    }
2161
2162    #[tokio::test]
2163    async fn test_execution_event_order_channel() {
2164        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2165
2166        let event = OrderSubmittedSpec::builder()
2167            .client_order_id(ClientOrderId::from("O-001"))
2168            .build();
2169
2170        tx.send(ExecutionEvent::Order(OrderEventAny::Submitted(event)))
2171            .unwrap();
2172
2173        let received = rx.recv().await.unwrap();
2174        match received {
2175            ExecutionEvent::Order(OrderEventAny::Submitted(e)) => {
2176                assert_eq!(e.client_order_id(), ClientOrderId::from("O-001"));
2177            }
2178            _ => panic!("Expected OrderSubmitted event"),
2179        }
2180    }
2181
2182    #[tokio::test]
2183    async fn test_execution_report_order_status_channel() {
2184        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2185
2186        let report = OrderStatusReport::new(
2187            AccountId::from("SIM-001"),
2188            InstrumentId::from("EUR/USD.SIM"),
2189            Some(ClientOrderId::from("O-001")),
2190            VenueOrderId::from("V-001"),
2191            OrderSide::Buy.into(),
2192            OrderType::Market,
2193            TimeInForce::Gtc,
2194            OrderStatus::Accepted,
2195            Quantity::from(100_000),
2196            Quantity::from(100_000),
2197            UnixNanos::from(1),
2198            UnixNanos::from(2),
2199            UnixNanos::from(3),
2200            None,
2201        );
2202
2203        tx.send(ExecutionEvent::Report(ExecutionReport::Order(Box::new(
2204            report,
2205        ))))
2206        .unwrap();
2207
2208        let received = rx.recv().await.unwrap();
2209        match received {
2210            ExecutionEvent::Report(ExecutionReport::Order(r)) => {
2211                assert_eq!(r.venue_order_id.as_str(), "V-001");
2212                assert_eq!(r.order_status, OrderStatus::Accepted);
2213            }
2214            _ => panic!("Expected OrderStatusReport"),
2215        }
2216    }
2217
2218    #[tokio::test]
2219    async fn test_execution_report_fill() {
2220        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2221
2222        let report = FillReport::new(
2223            AccountId::from("SIM-001"),
2224            InstrumentId::from("EUR/USD.SIM"),
2225            VenueOrderId::from("V-001"),
2226            TradeId::from("T-001"),
2227            OrderSide::Buy,
2228            Quantity::from(100_000),
2229            Price::from("1.10000"),
2230            Money::from("10 USD"),
2231            LiquiditySide::Taker,
2232            Some(ClientOrderId::from("O-001")),
2233            None,
2234            UnixNanos::from(1),
2235            UnixNanos::from(2),
2236            None,
2237        );
2238
2239        tx.send(ExecutionEvent::Report(ExecutionReport::Fill(Box::new(
2240            report,
2241        ))))
2242        .unwrap();
2243
2244        let received = rx.recv().await.unwrap();
2245        match received {
2246            ExecutionEvent::Report(ExecutionReport::Fill(r)) => {
2247                assert_eq!(r.venue_order_id.as_str(), "V-001");
2248                assert_eq!(r.trade_id.to_string(), "T-001");
2249            }
2250            _ => panic!("Expected FillReport"),
2251        }
2252    }
2253
2254    #[tokio::test]
2255    async fn test_execution_report_position() {
2256        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2257
2258        let report = PositionStatusReport::new(
2259            AccountId::from("SIM-001"),
2260            InstrumentId::from("EUR/USD.SIM"),
2261            PositionSide::Long,
2262            Quantity::from(100_000),
2263            UnixNanos::from(1),
2264            UnixNanos::from(2),
2265            None,
2266            Some(PositionId::from("P-001")),
2267            None,
2268        );
2269
2270        tx.send(ExecutionEvent::Report(ExecutionReport::Position(Box::new(
2271            report,
2272        ))))
2273        .unwrap();
2274
2275        let received = rx.recv().await.unwrap();
2276        match received {
2277            ExecutionEvent::Report(ExecutionReport::Position(r)) => {
2278                assert_eq!(r.venue_position_id.unwrap().as_str(), "P-001");
2279            }
2280            _ => panic!("Expected PositionStatusReport"),
2281        }
2282    }
2283
2284    #[tokio::test]
2285    async fn test_execution_event_account() {
2286        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2287
2288        let account_state = AccountState::new(
2289            AccountId::from("SIM-001"),
2290            AccountType::Cash,
2291            vec![],
2292            vec![],
2293            true,
2294            UUID4::new(),
2295            UnixNanos::from(1),
2296            UnixNanos::from(2),
2297            None,
2298        );
2299
2300        tx.send(ExecutionEvent::Account(account_state)).unwrap();
2301
2302        let received = rx.recv().await.unwrap();
2303        match received {
2304            ExecutionEvent::Account(r) => {
2305                assert_eq!(r.account_id.as_str(), "SIM-001");
2306            }
2307            _ => panic!("Expected AccountState"),
2308        }
2309    }
2310
2311    #[tokio::test]
2312    async fn test_runner_stop_method() {
2313        let (_data_tx, data_evt_rx) =
2314            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
2315        let (_cmd_tx, data_cmd_rx) =
2316            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
2317        let (_time_tx, time_evt_rx) =
2318            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
2319        let (_exec_evt_tx, exec_evt_rx) =
2320            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
2321        let (_exec_cmd_tx, exec_cmd_rx) =
2322            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
2323        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
2324
2325        let mut runner = create_test_runner(
2326            time_evt_rx,
2327            data_evt_rx,
2328            data_cmd_rx,
2329            exec_evt_rx,
2330            exec_cmd_rx,
2331            signal_rx,
2332            signal_tx.clone(),
2333        );
2334
2335        let runner_handle = tokio::spawn(async move {
2336            runner.run().await.unwrap();
2337        });
2338
2339        // Use stop via signal_tx directly
2340        signal_tx.send(()).unwrap();
2341
2342        let result = tokio::time::timeout(Duration::from_millis(100), runner_handle).await;
2343        assert!(result.is_ok(), "Runner should stop when stop() is called");
2344    }
2345
2346    #[tokio::test]
2347    async fn test_all_event_types_integration() {
2348        let (data_evt_tx, data_evt_rx) =
2349            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
2350        let (data_cmd_tx, data_cmd_rx) =
2351            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
2352        let (time_evt_tx, time_evt_rx) =
2353            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
2354        let (exec_evt_tx, exec_evt_rx) =
2355            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
2356        let (_exec_cmd_tx, exec_cmd_rx) =
2357            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
2358        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
2359
2360        let mut runner = create_test_runner(
2361            time_evt_rx,
2362            data_evt_rx,
2363            data_cmd_rx,
2364            exec_evt_rx,
2365            exec_cmd_rx,
2366            signal_rx,
2367            signal_tx.clone(),
2368        );
2369
2370        let runner_handle = tokio::spawn(async move {
2371            runner.run().await.unwrap();
2372        });
2373
2374        // Send data event
2375        let quote = test_quote();
2376        data_evt_tx
2377            .send((DataEvent::Data(Data::Quote(quote))).into())
2378            .unwrap();
2379
2380        // Send data command
2381        let command = DataCommand::Subscribe(SubscribeCommand::Data(SubscribeCustomData {
2382            client_id: Some(ClientId::from("TEST")),
2383            venue: None,
2384            data_type: DataType::new("QuoteTick", None, None),
2385            command_id: UUID4::new(),
2386            ts_init: UnixNanos::default(),
2387            correlation_id: None,
2388            params: None,
2389        }));
2390
2391        data_cmd_tx.send(command.into()).unwrap();
2392
2393        // Send time event
2394        let event = TimeEvent::new(
2395            Ustr::from("test"),
2396            UUID4::new(),
2397            UnixNanos::from(1),
2398            UnixNanos::from(2),
2399        );
2400        let callback = TimeEventCallback::from(|_: TimeEvent| {});
2401        let message = TimeEventMessage::new(event, callback);
2402        time_evt_tx.send((message).into()).unwrap();
2403
2404        // Send execution order event
2405        let order_event = OrderSubmittedSpec::builder()
2406            .client_order_id(ClientOrderId::from("O-001"))
2407            .build();
2408        exec_evt_tx
2409            .send((ExecutionEvent::Order(OrderEventAny::Submitted(order_event))).into())
2410            .unwrap();
2411
2412        // Send execution report (OrderStatus)
2413        let order_status = OrderStatusReport::new(
2414            AccountId::from("SIM-001"),
2415            InstrumentId::from("EUR/USD.SIM"),
2416            Some(ClientOrderId::from("O-001")),
2417            VenueOrderId::from("V-001"),
2418            OrderSide::Buy.into(),
2419            OrderType::Market,
2420            TimeInForce::Gtc,
2421            OrderStatus::Accepted,
2422            Quantity::from(100_000),
2423            Quantity::from(100_000),
2424            UnixNanos::from(1),
2425            UnixNanos::from(2),
2426            UnixNanos::from(3),
2427            None,
2428        );
2429        exec_evt_tx
2430            .send((ExecutionEvent::Report(ExecutionReport::Order(Box::new(order_status)))).into())
2431            .unwrap();
2432
2433        // Send execution report (Fill)
2434        let fill = FillReport::new(
2435            AccountId::from("SIM-001"),
2436            InstrumentId::from("EUR/USD.SIM"),
2437            VenueOrderId::from("V-001"),
2438            TradeId::from("T-001"),
2439            OrderSide::Buy,
2440            Quantity::from(100_000),
2441            Price::from("1.10000"),
2442            Money::from("10 USD"),
2443            LiquiditySide::Taker,
2444            Some(ClientOrderId::from("O-001")),
2445            None,
2446            UnixNanos::from(1),
2447            UnixNanos::from(2),
2448            None,
2449        );
2450        exec_evt_tx
2451            .send((ExecutionEvent::Report(ExecutionReport::Fill(Box::new(fill)))).into())
2452            .unwrap();
2453
2454        // Send execution report (Position)
2455        let position = PositionStatusReport::new(
2456            AccountId::from("SIM-001"),
2457            InstrumentId::from("EUR/USD.SIM"),
2458            PositionSide::Long,
2459            Quantity::from(100_000),
2460            UnixNanos::from(1),
2461            UnixNanos::from(2),
2462            None,
2463            Some(PositionId::from("P-001")),
2464            None,
2465        );
2466        exec_evt_tx
2467            .send((ExecutionEvent::Report(ExecutionReport::Position(Box::new(position)))).into())
2468            .unwrap();
2469
2470        // Send account event
2471        let account_state = AccountState::new(
2472            AccountId::from("SIM-001"),
2473            AccountType::Cash,
2474            vec![],
2475            vec![],
2476            true,
2477            UUID4::new(),
2478            UnixNanos::from(1),
2479            UnixNanos::from(2),
2480            None,
2481        );
2482        exec_evt_tx
2483            .send((ExecutionEvent::Account(account_state)).into())
2484            .unwrap();
2485
2486        // Yield to let runner enter event loop before stop signal
2487        tokio::task::yield_now().await;
2488        signal_tx.send(()).unwrap();
2489
2490        let result = tokio::time::timeout(Duration::from_millis(200), runner_handle).await;
2491        assert!(
2492            result.is_ok(),
2493            "Runner should process all event types and stop cleanly"
2494        );
2495    }
2496
2497    #[tokio::test]
2498    async fn test_runner_handle_stops_runner() {
2499        let (_data_tx, data_evt_rx) =
2500            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
2501        let (_cmd_tx, data_cmd_rx) =
2502            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
2503        let (_time_tx, time_evt_rx) =
2504            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
2505        let (_exec_evt_tx, exec_evt_rx) =
2506            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
2507        let (_exec_cmd_tx, exec_cmd_rx) =
2508            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
2509        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
2510
2511        let mut runner = create_test_runner(
2512            time_evt_rx,
2513            data_evt_rx,
2514            data_cmd_rx,
2515            exec_evt_rx,
2516            exec_cmd_rx,
2517            signal_rx,
2518            signal_tx.clone(),
2519        );
2520
2521        // Get handle before moving runner
2522        let handle = runner.handle();
2523
2524        let runner_handle = tokio::spawn(async move {
2525            runner.run().await.unwrap();
2526        });
2527
2528        // Use handle to stop
2529        handle.stop();
2530
2531        let result = tokio::time::timeout(Duration::from_millis(100), runner_handle).await;
2532        assert!(result.is_ok(), "Runner should stop via handle");
2533    }
2534
2535    #[tokio::test]
2536    async fn test_runner_handle_is_cloneable() {
2537        let (signal_tx, _signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
2538        let handle = AsyncRunnerHandle { signal_tx };
2539
2540        let handle2 = handle.clone();
2541
2542        // Both handles should be able to send stop signals
2543        assert!(handle.signal_tx.send(()).is_ok());
2544        assert!(handle2.signal_tx.send(()).is_ok());
2545    }
2546
2547    #[tokio::test]
2548    async fn test_runner_processes_events_before_stop() {
2549        let (data_evt_tx, data_evt_rx) =
2550            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataEvent>>();
2551        let (_cmd_tx, data_cmd_rx) =
2552            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<DataCommand>>();
2553        let (_time_tx, time_evt_rx) =
2554            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TimeEventMessage>>();
2555        let (_exec_evt_tx, exec_evt_rx) =
2556            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<ExecutionEvent>>();
2557        let (_exec_cmd_tx, exec_cmd_rx) =
2558            tokio::sync::mpsc::unbounded_channel::<DispatchMessage<TradingCommandMessage>>();
2559        let (signal_tx, signal_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
2560
2561        let mut runner = create_test_runner(
2562            time_evt_rx,
2563            data_evt_rx,
2564            data_cmd_rx,
2565            exec_evt_rx,
2566            exec_cmd_rx,
2567            signal_rx,
2568            signal_tx.clone(),
2569        );
2570
2571        let handle = runner.handle();
2572
2573        // Send events before starting runner
2574        for _ in 0..10 {
2575            let quote = test_quote();
2576            data_evt_tx
2577                .send((DataEvent::Data(Data::Quote(quote))).into())
2578                .unwrap();
2579        }
2580
2581        let runner_handle = tokio::spawn(async move {
2582            runner.run().await.unwrap();
2583        });
2584
2585        // Yield to let runner enter event loop before stop signal
2586        tokio::task::yield_now().await;
2587        handle.stop();
2588
2589        let result = tokio::time::timeout(Duration::from_millis(200), runner_handle).await;
2590        assert!(result.is_ok(), "Runner should process events and stop");
2591    }
2592
2593    #[rstest]
2594    fn test_new_does_not_bind_tls() {
2595        std::thread::spawn(|| {
2596            let _runner = AsyncRunner::new();
2597            assert!(try_get_time_event_sender().is_none());
2598            assert!(try_get_system_command_sender().is_none());
2599            assert!(try_get_system_event_sender().is_none());
2600            assert!(try_get_trading_cmd_sender().is_none());
2601        })
2602        .join()
2603        .unwrap();
2604    }
2605
2606    #[rstest]
2607    fn test_bind_senders_routes_to_runner_channels() {
2608        std::thread::spawn(|| {
2609            let mut runner = AsyncRunner::new();
2610            runner.bind_senders();
2611
2612            get_data_cmd_sender().execute(DataCommand::Subscribe(SubscribeCommand::Data(
2613                SubscribeCustomData {
2614                    client_id: Some(ClientId::from("TEST")),
2615                    venue: None,
2616                    data_type: DataType::new("test", None, None),
2617                    command_id: UUID4::new(),
2618                    ts_init: UnixNanos::default(),
2619                    correlation_id: None,
2620                    params: None,
2621                },
2622            )));
2623
2624            assert!(runner.channels.data_cmd_rx.try_recv().is_ok());
2625
2626            get_trading_cmd_sender().execute(TradingCommandMessage::new(
2627                MessagingSwitchboard::exec_engine_execute(),
2628                TradingCommand::CancelAllOrders(CancelAllOrders::new(
2629                    TraderId::from("TRADER-001"),
2630                    None,
2631                    StrategyId::from("S-001"),
2632                    InstrumentId::from("EUR/USD.SIM"),
2633                    Some(OrderSide::Buy),
2634                    UUID4::new(),
2635                    UnixNanos::default(),
2636                    None,
2637                    None, // correlation_id
2638                )),
2639            ));
2640            assert!(runner.channels.exec_cmd_rx.try_recv().is_ok());
2641
2642            let event = TimeEvent::new(
2643                Ustr::from("test"),
2644                UUID4::new(),
2645                UnixNanos::from(1),
2646                UnixNanos::from(2),
2647            );
2648            let callback = TimeEventCallback::from(|_: TimeEvent| {});
2649            get_time_event_sender().send(TimeEventMessage::new(event, callback));
2650            assert!(runner.channels.time_evt_rx.try_recv().is_ok());
2651
2652            get_system_event_sender().send(test_system_event()).unwrap();
2653            assert_eq!(
2654                runner
2655                    .channels
2656                    .system_evt_rx
2657                    .try_recv()
2658                    .unwrap()
2659                    .dispatch(|value| value),
2660                test_system_event()
2661            );
2662
2663            get_system_command_sender()
2664                .send(test_system_command())
2665                .unwrap();
2666            assert_eq!(
2667                runner
2668                    .channels
2669                    .system_cmd_rx
2670                    .try_recv()
2671                    .unwrap()
2672                    .dispatch(|value| value),
2673                test_system_command()
2674            );
2675
2676            get_data_event_sender()
2677                .send(DataEvent::Data(Data::Quote(test_quote())))
2678                .unwrap();
2679            assert!(runner.channels.data_evt_rx.try_recv().is_ok());
2680
2681            let account = AccountState::new(
2682                AccountId::from("SIM-001"),
2683                AccountType::Cash,
2684                vec![],
2685                vec![],
2686                true,
2687                UUID4::new(),
2688                UnixNanos::from(1),
2689                UnixNanos::from(2),
2690                None,
2691            );
2692            get_exec_event_sender()
2693                .send(ExecutionEvent::Account(account))
2694                .unwrap();
2695            assert!(runner.channels.exec_evt_rx.try_recv().is_ok());
2696        })
2697        .join()
2698        .unwrap();
2699    }
2700
2701    #[cfg(feature = "node")]
2702    #[rstest]
2703    fn test_drain_pending_system_events_keeps_data_events_separate() {
2704        std::thread::spawn(|| {
2705            let mut runner = AsyncRunner::new();
2706            runner.bind_senders();
2707            let system_event = test_system_event();
2708
2709            get_system_event_sender().send(system_event).unwrap();
2710            get_data_event_sender()
2711                .send(DataEvent::Data(Data::Quote(test_quote())))
2712                .unwrap();
2713
2714            let system_events = runner
2715                .drain_pending_system_events()
2716                .into_iter()
2717                .map(|message| message.dispatch(|value| value))
2718                .collect::<Vec<_>>();
2719
2720            assert_eq!(system_events, vec![system_event]);
2721            assert!(runner.channels.system_evt_rx.try_recv().is_err());
2722            assert!(runner.channels.data_evt_rx.try_recv().is_ok());
2723        })
2724        .join()
2725        .unwrap();
2726    }
2727
2728    #[cfg(feature = "node")]
2729    #[rstest]
2730    fn test_drain_pending_system_commands_keeps_events_separate() {
2731        std::thread::spawn(|| {
2732            let mut runner = AsyncRunner::new();
2733            runner.bind_senders();
2734            let system_command = test_system_command();
2735
2736            get_system_command_sender().send(system_command).unwrap();
2737            get_system_event_sender().send(test_system_event()).unwrap();
2738
2739            let system_commands = runner
2740                .drain_pending_system_commands()
2741                .into_iter()
2742                .map(|message| message.dispatch(|value| value))
2743                .collect::<Vec<_>>();
2744
2745            assert_eq!(system_commands, vec![system_command]);
2746            assert!(runner.channels.system_cmd_rx.try_recv().is_err());
2747            assert!(runner.channels.system_evt_rx.try_recv().is_ok());
2748        })
2749        .join()
2750        .unwrap();
2751    }
2752
2753    #[rstest]
2754    fn test_bind_senders_reclaims_tls_from_previous_runner() {
2755        std::thread::spawn(|| {
2756            let mut runner1 = AsyncRunner::new();
2757            runner1.bind_senders();
2758
2759            let mut runner2 = AsyncRunner::new();
2760            runner2.bind_senders();
2761
2762            get_data_cmd_sender().execute(DataCommand::Subscribe(SubscribeCommand::Data(
2763                SubscribeCustomData {
2764                    client_id: Some(ClientId::from("TEST")),
2765                    venue: None,
2766                    data_type: DataType::new("test", None, None),
2767                    command_id: UUID4::new(),
2768                    ts_init: UnixNanos::default(),
2769                    correlation_id: None,
2770                    params: None,
2771                },
2772            )));
2773
2774            assert!(runner2.channels.data_cmd_rx.try_recv().is_ok());
2775            assert!(runner1.channels.data_cmd_rx.try_recv().is_err());
2776        })
2777        .join()
2778        .unwrap();
2779    }
2780
2781    #[tokio::test]
2782    async fn test_execution_event_order_submitted_batch_channel() {
2783        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2784
2785        let events = vec![
2786            OrderSubmittedSpec::builder()
2787                .client_order_id(ClientOrderId::from("O-001"))
2788                .build(),
2789            OrderSubmittedSpec::builder()
2790                .client_order_id(ClientOrderId::from("O-002"))
2791                .build(),
2792        ];
2793
2794        let batch = OrderSubmittedBatch::new(events);
2795        tx.send(ExecutionEvent::OrderSubmittedBatch(batch)).unwrap();
2796
2797        let received = rx.recv().await.unwrap();
2798        match received {
2799            ExecutionEvent::OrderSubmittedBatch(b) => {
2800                assert_eq!(b.len(), 2);
2801                assert_eq!(b.events[0].client_order_id, ClientOrderId::from("O-001"));
2802                assert_eq!(b.events[1].client_order_id, ClientOrderId::from("O-002"));
2803            }
2804            _ => panic!("Expected OrderSubmittedBatch event"),
2805        }
2806    }
2807
2808    #[tokio::test]
2809    async fn test_execution_event_order_accepted_batch_channel() {
2810        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2811
2812        let events = vec![
2813            OrderAcceptedSpec::builder()
2814                .client_order_id(ClientOrderId::from("O-001"))
2815                .build(),
2816            OrderAcceptedSpec::builder()
2817                .client_order_id(ClientOrderId::from("O-002"))
2818                .build(),
2819        ];
2820
2821        let batch = OrderAcceptedBatch::new(events);
2822        tx.send(ExecutionEvent::OrderAcceptedBatch(batch)).unwrap();
2823
2824        let received = rx.recv().await.unwrap();
2825        match received {
2826            ExecutionEvent::OrderAcceptedBatch(b) => {
2827                assert_eq!(b.len(), 2);
2828                assert_eq!(b.events[0].client_order_id, ClientOrderId::from("O-001"));
2829                assert_eq!(b.events[1].client_order_id, ClientOrderId::from("O-002"));
2830            }
2831            _ => panic!("Expected OrderAcceptedBatch event"),
2832        }
2833    }
2834
2835    #[tokio::test]
2836    async fn test_execution_event_order_canceled_batch_channel() {
2837        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
2838
2839        let events = vec![
2840            OrderCanceledSpec::builder()
2841                .client_order_id(ClientOrderId::from("O-001"))
2842                .build(),
2843            OrderCanceledSpec::builder()
2844                .client_order_id(ClientOrderId::from("O-002"))
2845                .build(),
2846        ];
2847
2848        let batch = OrderCanceledBatch::new(events);
2849        tx.send(ExecutionEvent::OrderCanceledBatch(batch)).unwrap();
2850
2851        let received = rx.recv().await.unwrap();
2852        match received {
2853            ExecutionEvent::OrderCanceledBatch(b) => {
2854                assert_eq!(b.len(), 2);
2855                assert_eq!(b.events[0].client_order_id, ClientOrderId::from("O-001"));
2856                assert_eq!(b.events[1].client_order_id, ClientOrderId::from("O-002"));
2857            }
2858            _ => panic!("Expected OrderCanceledBatch event"),
2859        }
2860    }
2861}