1use 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#[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 #[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 #[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#[derive(Debug, Clone)]
138pub struct AsyncTimeEventSender {
139 time_tx: DispatchSender<TimeEventMessage>,
140}
141
142impl AsyncTimeEventSender {
143 #[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#[derive(Debug)]
164pub struct AsyncTradingCommandSender {
165 cmd_tx: tokio::sync::mpsc::UnboundedSender<DispatchMessage<TradingCommandMessage>>,
166 owner: ThreadId,
167}
168
169impl AsyncTradingCommandSender {
170 #[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#[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#[derive(Clone, Debug)]
241pub struct AsyncRunnerHandle {
242 signal_tx: tokio::sync::mpsc::UnboundedSender<()>,
243}
244
245impl AsyncRunnerHandle {
246 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 #[must_use]
273 #[rustfmt::skip]
274 pub fn new() -> Self {
275 use tokio::sync::mpsc::unbounded_channel; 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 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 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 #[must_use]
346 pub fn handle(&self) -> AsyncRunnerHandle {
347 AsyncRunnerHandle {
348 signal_tx: self.signal_tx.clone(),
349 }
350 }
351
352 #[must_use]
357 pub fn take_channels(self) -> AsyncRunnerChannels {
358 self.channels
359 }
360
361 pub fn flush_pending_data(&mut self) {
367 let mut total = 0;
368
369 loop {
370 let mut progressed = false;
371
372 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 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 #[inline]
481 #[must_use]
482 pub fn handle_time_event(message: DispatchMessage<TimeEventMessage>) -> bool {
483 message.dispatch(TimeEventMessage::dispatch)
484 }
485
486 #[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 pub fn dispatch_data_event(event: DispatchMessage<DataEvent>) {
504 event.dispatch(Self::handle_data_event);
505 }
506
507 pub fn dispatch_exec_event(event: DispatchMessage<ExecutionEvent>) {
513 event.dispatch(Self::handle_exec_event);
514 }
515
516 #[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 #[inline]
547 pub fn handle_exec_command(cmd: TradingCommand) {
548 msgbus::send_trading_command(MessagingSwitchboard::exec_engine_execute(), cmd);
549 }
550
551 #[inline]
557 pub fn handle_trading_command(message: DispatchMessage<TradingCommandMessage>) {
558 message.dispatch_trading(|_| {});
559 }
560
561 #[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 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 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 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 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 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 let runner_handle = tokio::spawn(async move {
1708 runner.run().await.unwrap();
1709 });
1710
1711 signal_tx.send(()).unwrap();
1713
1714 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 let runner_handle = tokio::spawn(async move {
1745 runner.run().await.unwrap();
1746 });
1747
1748 drop(data_tx);
1749
1750 tokio::task::yield_now().await;
1752 signal_tx.send(()).ok();
1753
1754 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 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 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 let runner_handle = tokio::spawn(async move {
1808 runner.run().await.unwrap();
1809 });
1810
1811 for handle in handles {
1813 handle.await.unwrap();
1814 }
1815
1816 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 for _ in 0..count {
1833 tx.send(DataEvent::Data(Data::Quote(quote))).unwrap();
1834 }
1835
1836 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, ));
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, ));
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 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 let quote = test_quote();
2376 data_evt_tx
2377 .send((DataEvent::Data(Data::Quote(quote))).into())
2378 .unwrap();
2379
2380 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 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 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 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 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 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 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 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 let handle = runner.handle();
2523
2524 let runner_handle = tokio::spawn(async move {
2525 runner.run().await.unwrap();
2526 });
2527
2528 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 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 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 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, )),
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}