1use std::{collections::HashSet, fmt::Debug, future::Future, pin::Pin, time::Duration};
81
82use anyhow::Context;
83use indexmap::IndexSet;
84use nautilus_common::{
85 actor::{Actor, DataActor, DataActorNative},
86 cache::database::{CacheDatabaseAdapter, CacheDatabaseFactory},
87 clients::ExecutionClient,
88 component::Component,
89 enums::{Environment, LogColor},
90 live::dst,
91 log_info,
92 messages::{
93 DataEvent, ExecutionEvent, ExecutionReport, SystemCommand, SystemEvent,
94 data::DataCommand,
95 execution::{GenerateOrderStatusReports, GeneratePositionStatusReports, TradingCommand},
96 system::{QueueStateChanged, SocketStateChange, SocketStateChanged},
97 },
98 msgbus::{self, BusMessage, MessagingSwitchboard},
99 runner::{SystemChannel, TimeEventMessage, TradingCommandMessage},
100};
101use nautilus_core::{
102 UUID4,
103 datetime::{NANOSECONDS_IN_MILLISECOND, mins_to_secs, secs_to_nanos_unchecked},
104};
105use nautilus_execution::engine::ExecutionEngine;
106#[cfg(test)]
107use nautilus_model::reports::OrderStatusReport;
108use nautilus_model::{
109 events::OrderEventAny,
110 identifiers::{ClientId, ClientOrderId, InstrumentId, StrategyId, TraderId},
111 orders::Order,
112 reports::PositionStatusReport,
113};
114#[cfg(feature = "python")]
115use nautilus_system::trader::Trader;
116use nautilus_system::{config::NautilusKernelConfig, kernel::NautilusKernel};
117use nautilus_trading::{
118 ExecutionAlgorithm, ExecutionAlgorithmNative,
119 strategy::{Strategy, StrategyNative},
120};
121use tabled::{builder::Builder, settings::Style};
122
123use crate::{
124 execution::{
125 client::LiveExecutionClient,
126 manager::{
127 ExecutionManager, ExecutionManagerConfig, OpenOrderReportCheck, PositionReportCheck,
128 SourcedOrderStatusReport, TargetedOrderQuery, TargetedOrderReportResult,
129 request_targeted_order_reports,
130 },
131 },
132 runner::{AsyncRunner, AsyncRunnerChannels, PendingRunnerEvent},
133};
134
135pub mod builder;
136pub mod config;
137
138#[cfg(feature = "plugin")]
139pub mod plugin;
140
141mod metrics;
142mod queue;
143mod state;
144
145use builder::ExternalMessageBusIngress;
146pub use builder::LiveNodeBuilder;
147use config::{LiveNodeConfig, PluginConfig, validate_live_environment};
148pub use metrics::{RunnerChannelMetricsSnapshot, RunnerMetricsDelta, RunnerMetricsSnapshot};
149use metrics::{RunnerChannelQueueDepths, RunnerMetrics};
150use queue::{QueueMonitor, QueueStateTransition};
151use state::{EngineConnectionStatus, RunningTransition};
152pub use state::{LiveNodeHandle, NodeRunMode, NodeState};
153
154const DISPATCHES_PER_YIELD: usize = 64;
160
161#[derive(Debug)]
166pub struct LiveNode {
167 kernel: NautilusKernel,
168 runner: Option<AsyncRunner>,
169 config: LiveNodeConfig,
170 handle: LiveNodeHandle,
171 exec_manager: ExecutionManager,
172 exec_clients: Vec<LiveExecutionClient>,
173 cache_database_factory: Option<Box<dyn CacheDatabaseFactory>>,
174 external_msgbus: Option<ExternalMessageBusIngress>,
175 shutdown_deadline: Option<dst::time::Instant>,
176 #[cfg(feature = "plugin")]
177 plugins: plugin::NodePlugins,
178}
179
180impl LiveNode {
181 #[must_use]
185 pub(crate) fn new_from_builder(
186 kernel: NautilusKernel,
187 runner: AsyncRunner,
188 config: LiveNodeConfig,
189 exec_manager: ExecutionManager,
190 exec_clients: Vec<LiveExecutionClient>,
191 cache_database_factory: Option<Box<dyn CacheDatabaseFactory>>,
192 external_msgbus: Option<ExternalMessageBusIngress>,
193 ) -> Self {
194 Self {
195 kernel,
196 runner: Some(runner),
197 config,
198 handle: LiveNodeHandle::new(),
199 exec_manager,
200 exec_clients,
201 cache_database_factory,
202 external_msgbus,
203 shutdown_deadline: None,
204 #[cfg(feature = "plugin")]
205 plugins: plugin::NodePlugins,
206 }
207 }
208
209 pub fn builder(
215 trader_id: TraderId,
216 environment: Environment,
217 ) -> anyhow::Result<LiveNodeBuilder> {
218 LiveNodeBuilder::new(trader_id, environment)
219 }
220
221 pub fn build(name: String, config: Option<LiveNodeConfig>) -> anyhow::Result<Self> {
231 let config = config.unwrap_or_default();
232 validate_live_environment(config.environment())?;
233
234 config.validate_runtime_support()?;
235
236 if config.event_store.is_some() {
237 anyhow::bail!(
238 "LiveNodeConfig.event_store is set but LiveNode::build cannot install a factory; \
239 use LiveNodeBuilder::with_event_store(...) instead"
240 );
241 }
242
243 let runner = AsyncRunner::new();
244 runner.bind_senders();
245
246 let kernel = NautilusKernel::new(name, config.clone())?;
247 #[cfg(feature = "python")]
248 if let Some(controller) = config.controller.as_ref() {
249 Trader::add_controller_from_importable_config(&kernel.trader, controller)?;
250 }
251 #[cfg(not(feature = "python"))]
252 if let Some(controller) = config.controller.as_ref() {
253 anyhow::bail!(
254 "LiveNodeConfig.controller for importable controller '{}' requires the python feature",
255 controller.controller_path
256 );
257 }
258
259 let exec_manager_config =
260 ExecutionManagerConfig::from(&config.exec_engine).with_trader_id(config.trader_id);
261 let exec_manager = ExecutionManager::new(
262 kernel.clock.clone(),
263 kernel.cache.clone(),
264 exec_manager_config,
265 );
266
267 let node = Self {
268 kernel,
269 runner: Some(runner),
270 config,
271 handle: LiveNodeHandle::new(),
272 exec_manager,
273 exec_clients: Vec::new(),
274 cache_database_factory: None,
275 external_msgbus: None,
276 shutdown_deadline: None,
277 #[cfg(feature = "plugin")]
278 plugins: plugin::NodePlugins,
279 };
280 node.load_configured_plugins()?;
281
282 log::info!("LiveNode built successfully with kernel config");
283
284 Ok(node)
285 }
286
287 pub(crate) fn load_configured_plugins(&self) -> anyhow::Result<()> {
293 if self.config.plugins.is_empty() {
294 return Ok(());
295 }
296
297 anyhow::bail!(
298 "LiveNodeConfig.plugins requires host-side plug-in support; nautilus-plugin is the guest SDK only"
299 )
300 }
301
302 #[expect(
308 clippy::needless_pass_by_value,
309 reason = "signature mirrors the host-enabled API"
310 )]
311 pub fn add_plugin(&mut self, config: PluginConfig) -> anyhow::Result<()> {
312 #[cfg(feature = "plugin")]
313 config.validate_runtime_support(self.config.plugins.len())?;
314 #[cfg(not(feature = "plugin"))]
315 let _ = config;
316
317 anyhow::bail!(
318 "LiveNode::add_plugin requires host-side plug-in support; nautilus-plugin is the guest SDK only"
319 )
320 }
321
322 #[must_use]
324 pub fn handle(&self) -> LiveNodeHandle {
325 self.handle.clone()
326 }
327
328 pub async fn start(&mut self) -> anyhow::Result<()> {
339 if self.state().is_running() {
340 anyhow::bail!("Already running");
341 }
342
343 if self.external_msgbus.is_some() {
344 log::warn!(
345 "External message bus ingress is configured but LiveNode::start() does not service it; use LiveNode::run()"
346 );
347 }
348
349 self.prepare_cache().await?;
350
351 if let Some(runner) = self.runner.as_ref() {
352 runner.bind_senders();
353 }
354
355 self.handle.set_starting();
356
357 self.kernel.reset_shutdown_flag();
358 self.kernel.start_async().await;
359
360 if self.kernel.is_event_store_replay() {
361 log::info!(
362 "Event-store replay loaded; skipping live client connection and reconciliation",
363 );
364
365 if !self.finish_startup_replay().await? {
366 return Ok(());
367 }
368 return Ok(());
369 }
370
371 if self.kernel.is_event_store_replay_configured() {
372 self.abort_startup("Event-store replay did not start")
373 .await?;
374 return Ok(());
375 }
376
377 let connection_deadline = dst::time::Instant::now() + self.config.timeout_connection;
378
379 if let Err(e) = self.connect_data_phase(connection_deadline).await {
381 return self
382 .abort_startup_with_error("Data client connection timed out", e)
383 .await;
384 }
385
386 let (startup_system_events, startup_system_commands) =
387 if let Some(runner) = self.runner.as_mut() {
388 runner.flush_pending_data();
389 (
390 runner.drain_pending_system_events(),
391 runner.drain_pending_system_commands(),
392 )
393 } else {
394 (Vec::new(), Vec::new())
395 };
396
397 if let Err(e) = self.connect_exec_clients(connection_deadline).await {
398 return self
399 .abort_startup_with_error("Execution client connection timed out", e)
400 .await;
401 }
402
403 if let Some(reason) = self.startup_abort_reason() {
404 self.abort_startup(reason).await?;
405 return Ok(());
406 }
407
408 match self.await_engines_connected(connection_deadline).await {
409 EngineConnectionStatus::Connected => {}
410 EngineConnectionStatus::TimedOut => {
411 return self
412 .abort_startup_with_error(
413 "Engine readiness timed out",
414 anyhow::anyhow!("readiness timeout while waiting for engine connections"),
415 )
416 .await;
417 }
418 EngineConnectionStatus::StopRequested => {
419 self.abort_startup("Stop signal received during startup")
420 .await?;
421 return Ok(());
422 }
423 EngineConnectionStatus::ShutdownRequested => {
424 self.abort_startup("Shutdown signal received during startup")
425 .await?;
426 return Ok(());
427 }
428 }
429
430 if let Err(e) = self.perform_startup_reconciliation().await {
431 if let Err(finalize_err) = self.abort_startup("Startup reconciliation failed").await {
432 anyhow::bail!(
433 "startup reconciliation failed: {e}; failed to finalize startup abort: {finalize_err}"
434 );
435 }
436
437 return Err(e);
438 }
439
440 if let Some(reason) = self.startup_abort_reason() {
441 self.abort_startup(reason).await?;
442 return Ok(());
443 }
444
445 if let Err(e) = self.kernel.start_trader() {
446 return self.abort_after_trader_start_failure(e).await;
447 }
448 #[cfg(feature = "plugin")]
449 if let Err(e) = self.plugins.start_controllers() {
450 return self.abort_after_trader_start_failure(e).await;
451 }
452
453 self.process_system_events(startup_system_events);
454 self.process_system_commands(startup_system_commands);
455
456 if !self.finish_startup_trader(None).await? {
457 return Ok(());
458 }
459
460 Ok(())
461 }
462
463 pub async fn stop(&mut self) -> anyhow::Result<()> {
472 if !self.state().is_running() {
473 anyhow::bail!("Not running");
474 }
475
476 self.handle.set_shutting_down();
477
478 #[cfg(feature = "plugin")]
479 let controller_stop_result = self.plugins.stop_controllers();
480 #[cfg(not(feature = "plugin"))]
481 let controller_stop_result: anyhow::Result<()> = Ok(());
482
483 self.kernel.stop_trader();
484 let delay = self.kernel.delay_post_stop();
485 log::info!("Awaiting residual events ({delay:?})...");
486
487 let residual_events = self.process_runner_for(delay).await;
488 if residual_events > 0 {
489 log::debug!("Processed {residual_events} residual events during shutdown");
490 }
491
492 let stop_result = self.finalize_stop().await;
493 let drained_events = self.drain_runner_pending();
494 if drained_events > 0 {
495 log::info!("Drained {drained_events} remaining events during shutdown");
496 }
497
498 match (controller_stop_result, stop_result) {
499 (Ok(()), Ok(())) => Ok(()),
500 (Err(controller_err), Ok(())) => Err(controller_err),
501 (Ok(()), Err(stop_err)) => Err(stop_err),
502 (Err(controller_err), Err(stop_err)) => {
503 log::error!("Error stopping plug-in controllers: {controller_err}");
504 Err(stop_err)
505 }
506 }
507 }
508
509 pub fn dispose(&mut self) {
511 self.close_external_ingress();
512 self.kernel.dispose();
513 self.handle.set_stopped();
514 }
515
516 async fn process_runner_for(&mut self, duration: Duration) -> usize {
517 let Some(mut runner) = self.runner.take() else {
518 dst::time::sleep(duration).await;
519 return 0;
520 };
521
522 runner.bind_senders();
523 let deadline = dst::time::Instant::now() + duration;
524 let mut processed = 0;
525
526 loop {
527 tokio::select! {
528 biased;
529
530 () = dst::time::sleep_until(deadline) => break,
531 event = runner.recv() => {
532 let Some(event) = event else {
533 dst::time::sleep_until(deadline).await;
534 break;
535 };
536
537 self.process_runner_event(event);
538 processed += 1;
539 }
540 }
541 }
542
543 self.runner = Some(runner);
544 processed
545 }
546
547 fn drain_runner_pending(&mut self) -> usize {
548 let Some(mut runner) = self.runner.take() else {
549 return 0;
550 };
551
552 let processed = runner.poll_pending(|event| self.process_runner_event(event));
553 self.runner = Some(runner);
554 processed
555 }
556
557 fn process_runner_event(&mut self, event: PendingRunnerEvent) {
558 match event {
559 PendingRunnerEvent::TimeEvent(message) => {
560 let _ = AsyncRunner::handle_time_event(message);
561 }
562 PendingRunnerEvent::SystemEvent(event) => self.process_system_event(event),
563 PendingRunnerEvent::SystemCommand(command) => self.process_system_command(command),
564 PendingRunnerEvent::ExecEvent(event) => self.process_exec_event(event),
565 PendingRunnerEvent::ExecCommand(command) => self.process_exec_command(command),
566 PendingRunnerEvent::DataEvent(event) => AsyncRunner::handle_data_event(event),
567 PendingRunnerEvent::DataCommand(command) => AsyncRunner::handle_data_command(command),
568 }
569 }
570
571 fn process_system_events(&self, events: Vec<SystemEvent>) {
572 for event in events {
573 self.process_system_event(event);
574 }
575 }
576
577 fn process_system_commands(&self, commands: Vec<SystemCommand>) {
578 for command in commands {
579 self.process_system_command(command);
580 }
581 }
582
583 fn process_system_command(&self, command: SystemCommand) {
584 match command {
585 SystemCommand::ReconnectSocket(command) => {
586 self.kernel.process_socket_reconnect(command);
587 }
588 }
589 }
590
591 fn process_system_event(&self, event: SystemEvent) {
592 match event {
593 SystemEvent::SocketState(change) => self.publish_socket_state_change(change),
594 }
595 }
596
597 fn publish_socket_state_change(&self, change: SocketStateChange) {
598 let timestamp = self.kernel.generate_timestamp_ns();
599 let event = SocketStateChanged::new(
600 self.config.trader_id,
601 change.client_id,
602 change.venue,
603 change.endpoint,
604 change.state,
605 UUID4::new(),
606 timestamp,
607 timestamp,
608 );
609
610 msgbus::publish_any(
611 MessagingSwitchboard::socket_state_changed_topic(),
612 event.as_any(),
613 );
614 }
615
616 async fn await_engines_connected(
620 &self,
621 deadline: dst::time::Instant,
622 ) -> EngineConnectionStatus {
623 log::info!(
624 "Awaiting engine connections ({:?} timeout)...",
625 self.config.timeout_connection
626 );
627
628 let interval = Duration::from_millis(100);
629
630 loop {
631 if self.handle.should_stop() {
632 log::warn!("Stop signal received, aborting connection wait");
633 return EngineConnectionStatus::StopRequested;
634 }
635
636 if self.kernel.is_shutdown_requested() {
637 log::warn!("Shutdown signal received, aborting connection wait");
638 return EngineConnectionStatus::ShutdownRequested;
639 }
640
641 if self.kernel.check_engines_connected() {
642 log::info!("All engine clients connected");
643 return EngineConnectionStatus::Connected;
644 }
645
646 let now = dst::time::Instant::now();
647 if now >= deadline {
648 break;
649 }
650
651 dst::time::sleep(interval.min(deadline - now)).await;
652 }
653
654 self.log_connection_status();
655 EngineConnectionStatus::TimedOut
656 }
657
658 async fn await_engines_disconnected(&self, deadline: dst::time::Instant) -> anyhow::Result<()> {
662 log::info!(
663 "Awaiting engine disconnections ({:?} timeout)...",
664 self.config.timeout_disconnection
665 );
666
667 let timeout = self.config.timeout_disconnection;
668 let interval = Duration::from_millis(100);
669
670 loop {
671 if self.kernel.check_engines_disconnected() {
672 log::info!("All engine clients disconnected");
673 return Ok(());
674 }
675
676 let now = dst::time::Instant::now();
677 if now >= deadline {
678 break;
679 }
680
681 dst::time::sleep(interval.min(deadline - now)).await;
682 }
683
684 log::error!(
685 "Timed out ({:?}) waiting for engines to disconnect\n\
686 DataEngine.check_disconnected() == {}\n\
687 ExecEngine.check_disconnected() == {}",
688 timeout,
689 self.kernel.data_engine().check_disconnected(),
690 self.kernel.exec_engine().borrow().check_disconnected(),
691 );
692 anyhow::bail!("disconnect readiness timeout while waiting for engine disconnections")
693 }
694
695 fn log_connection_status(&self) {
696 let data_status = self.kernel.data_client_connection_status();
697 let exec_status = self.kernel.exec_client_connection_status();
698
699 let mut rows: Vec<ClientStatus> = Vec::new();
700
701 for (client_id, connected) in data_status {
702 rows.push(ClientStatus {
703 client: client_id.to_string(),
704 client_type: "Data",
705 connected,
706 });
707 }
708
709 for (client_id, connected) in exec_status {
710 rows.push(ClientStatus {
711 client: client_id.to_string(),
712 client_type: "Execution",
713 connected,
714 });
715 }
716
717 let table = render_client_statuses(rows);
718
719 log::warn!(
720 "Timed out ({:?}) waiting for engines to connect\n\n{table}\n\n\
721 DataEngine.check_connected() == {}\n\
722 ExecEngine.check_connected() == {}",
723 self.config.timeout_connection,
724 self.kernel.data_engine().check_connected(),
725 self.kernel.exec_engine().borrow().check_connected(),
726 );
727 }
728
729 #[expect(clippy::await_holding_refcell_ref)] async fn perform_startup_reconciliation(&mut self) -> anyhow::Result<()> {
739 if !self.config.exec_engine.reconciliation {
740 log::info!("Startup reconciliation disabled");
741 self.kernel
742 .portfolio
743 .borrow_mut()
744 .initialize_wallet_orders()?;
745 return Ok(());
746 }
747
748 log_info!(
749 "Starting execution state reconciliation...",
750 color = LogColor::Blue
751 );
752
753 let lookback_mins = self
754 .config
755 .exec_engine
756 .reconciliation_lookback_mins
757 .map(u64::from);
758
759 let timeout = self.config.timeout_reconciliation;
760 let start = dst::time::Instant::now();
761 let client_ids = self.kernel.exec_engine.borrow().client_ids();
762
763 for client_id in client_ids {
764 let elapsed = start.elapsed();
765 if elapsed >= timeout {
766 anyhow::bail!("Startup reconciliation timeout reached");
767 }
768 let remaining = timeout
769 .checked_sub(elapsed)
770 .expect("elapsed checked against reconciliation timeout");
771
772 log_info!(
773 "Requesting mass status from {}...",
774 client_id,
775 color = LogColor::Blue
776 );
777
778 let mass_status_result = dst::time::timeout(remaining, async {
779 self.kernel
780 .exec_engine
781 .borrow_mut()
782 .generate_mass_status(&client_id, lookback_mins)
783 .await
784 })
785 .await
786 .map_err(|_| {
787 anyhow::anyhow!(
788 "Startup reconciliation timeout reached while requesting mass status from {client_id}"
789 )
790 })?;
791
792 match mass_status_result {
793 Ok(Some(mass_status)) => {
794 log_info!(
795 "Reconciling ExecutionMassStatus for {}",
796 client_id,
797 color = LogColor::Blue
798 );
799
800 let exec_engine_rc = self.kernel.exec_engine.clone();
801
802 let result = self
803 .exec_manager
804 .reconcile_execution_mass_status(mass_status, exec_engine_rc)
805 .await;
806
807 anyhow::ensure!(
808 self.kernel
809 .exec_engine
810 .borrow()
811 .get_client(&client_id)
812 .is_some(),
813 "Execution client {client_id} disappeared during startup reconciliation",
814 );
815
816 if result.events.is_empty() {
817 log_info!(
818 "Reconciliation for {} succeeded",
819 client_id,
820 color = LogColor::Blue
821 );
822 } else {
823 log::info!(
824 color = LogColor::Blue as u8;
825 "Reconciliation for {} processed {} events",
826 client_id,
827 result.events.len()
828 );
829 }
830
831 if !result.external_orders.is_empty() {
833 let exec_engine = self.kernel.exec_engine.borrow();
834 let source_client = exec_engine.get_client(&client_id).ok_or_else(|| {
835 anyhow::anyhow!(
836 "Execution client {client_id} disappeared during startup reconciliation"
837 )
838 })?;
839
840 for external in result.external_orders {
841 source_client.register_external_order(
842 external.client_order_id,
843 external.venue_order_id,
844 external.instrument_id,
845 external.strategy_id,
846 external.ts_init,
847 );
848 }
849 }
850 }
851 Ok(None) => {
852 log::warn!(
853 "No mass status available from {client_id} \
854 (likely adapter error when generating reports)"
855 );
856 }
857 Err(e) => {
858 return Err(e).context(format!("Failed to get mass status from {client_id}"));
859 }
860 }
861 }
862
863 self.kernel.portfolio.borrow_mut().initialize_orders();
864 self.kernel.portfolio.borrow_mut().initialize_positions();
865 self.kernel
866 .portfolio
867 .borrow_mut()
868 .initialize_wallet_orders()?;
869
870 let elapsed_secs = start.elapsed().as_secs_f64();
871 log_info!(
872 "Startup reconciliation completed in {:.2}s",
873 elapsed_secs,
874 color = LogColor::Blue
875 );
876
877 Ok(())
878 }
879
880 pub async fn run(&mut self) -> anyhow::Result<()> {
902 self.run_with_mode(NodeRunMode::Owned).await
903 }
904
905 pub async fn run_with_mode(&mut self, mode: NodeRunMode) -> anyhow::Result<()> {
915 if self.state().is_running() {
916 anyhow::bail!("Already running");
917 }
918
919 if self.runner.is_none() {
920 anyhow::bail!("Runner already consumed - run() called twice");
921 }
922
923 self.prepare_cache().await?;
924
925 let Some(runner) = self.runner.take() else {
926 anyhow::bail!("Runner already consumed - run() called twice");
927 };
928 runner.bind_senders();
929
930 let AsyncRunnerChannels {
931 mut time_evt_rx,
932 mut system_evt_rx,
933 mut system_cmd_rx,
934 mut exec_evt_rx,
935 mut exec_cmd_rx,
936 mut data_evt_rx,
937 mut data_cmd_rx,
938 } = runner.take_channels();
939
940 log::info!("Event loop starting");
941
942 self.handle.set_starting();
943 self.kernel.reset_shutdown_flag();
944 self.kernel.start_async().await;
945
946 if self.kernel.is_event_store_replay() {
947 log::info!(
948 "Event-store replay loaded; skipping live client connection and reconciliation",
949 );
950
951 if !self.finish_startup_replay().await? {
952 return Ok(());
953 }
954 return Ok(());
955 }
956
957 if self.kernel.is_event_store_replay_configured() {
958 self.abort_startup("Event-store replay did not start")
959 .await?;
960 return Ok(());
961 }
962
963 let mut external_msgbus_rx = match self.take_external_ingress_receiver() {
964 Ok(rx) => rx,
965 Err(e) => {
966 let result = self
967 .abort_startup("External message bus ingress failed to start")
968 .await;
969 Self::drain_channels(
970 &mut time_evt_rx,
971 &mut system_evt_rx,
972 &mut system_cmd_rx,
973 &mut exec_evt_rx,
974 &mut exec_cmd_rx,
975 &mut data_evt_rx,
976 &mut data_cmd_rx,
977 );
978 log::info!("Event loop stopped");
979
980 if let Err(finalize_err) = result {
981 anyhow::bail!(
982 "failed to start external message bus ingress: {e}; failed to finalize startup abort: {finalize_err}"
983 );
984 }
985
986 return Err(e);
987 }
988 };
989
990 let stop_handle = self.handle.clone();
991 let mut pending = PendingEvents::default();
992 let mut startup_system_events = Vec::new();
993 let mut startup_system_commands = Vec::new();
994 let connection_deadline = dst::time::Instant::now() + self.config.timeout_connection;
995
996 let data_connect_result = drive_with_event_buffering(
999 self.connect_data_phase(connection_deadline),
1000 &mut pending,
1001 &mut time_evt_rx,
1002 &mut system_evt_rx,
1003 &mut system_cmd_rx,
1004 &mut exec_evt_rx,
1005 &mut exec_cmd_rx,
1006 &mut data_evt_rx,
1007 &mut data_cmd_rx,
1008 )
1009 .await;
1010
1011 if let Err(e) = data_connect_result {
1012 flush_all_pending(
1013 &mut pending,
1014 &mut time_evt_rx,
1015 &mut system_evt_rx,
1016 &mut system_cmd_rx,
1017 &mut exec_evt_rx,
1018 &mut exec_cmd_rx,
1019 &mut data_evt_rx,
1020 &mut data_cmd_rx,
1021 );
1022 let result = self
1023 .abort_startup_with_error("Data client connection timed out", e)
1024 .await;
1025 Self::drain_channels(
1026 &mut time_evt_rx,
1027 &mut system_evt_rx,
1028 &mut system_cmd_rx,
1029 &mut exec_evt_rx,
1030 &mut exec_cmd_rx,
1031 &mut data_evt_rx,
1032 &mut data_cmd_rx,
1033 );
1034 log::info!("Event loop stopped");
1035 return result;
1036 }
1037
1038 flush_pending_data(&mut pending, &mut data_evt_rx, &mut data_cmd_rx);
1042 startup_system_events.extend(pending.take_system_events());
1043 startup_system_commands.extend(pending.take_system_commands());
1044 debug_assert!(
1045 pending.data_evts.is_empty() && pending.data_cmds.is_empty(),
1046 "data must be drained into cache before exec clients connect",
1047 );
1048
1049 let engine_connection_result = drive_with_event_buffering(
1051 self.connect_exec_phase(connection_deadline),
1052 &mut pending,
1053 &mut time_evt_rx,
1054 &mut system_evt_rx,
1055 &mut system_cmd_rx,
1056 &mut exec_evt_rx,
1057 &mut exec_cmd_rx,
1058 &mut data_evt_rx,
1059 &mut data_cmd_rx,
1060 )
1061 .await;
1062
1063 flush_all_pending(
1065 &mut pending,
1066 &mut time_evt_rx,
1067 &mut system_evt_rx,
1068 &mut system_cmd_rx,
1069 &mut exec_evt_rx,
1070 &mut exec_cmd_rx,
1071 &mut data_evt_rx,
1072 &mut data_cmd_rx,
1073 );
1074 startup_system_events.extend(pending.take_system_events());
1075 startup_system_commands.extend(pending.take_system_commands());
1076 debug_assert!(
1077 pending.is_empty(),
1078 "all startup events must be processed before reconciliation",
1079 );
1080
1081 let engine_connection_status = match engine_connection_result {
1082 Ok(status) => status,
1083 Err(e) => {
1084 let result = self
1085 .abort_startup_with_error("Execution client connection timed out", e)
1086 .await;
1087 Self::drain_channels(
1088 &mut time_evt_rx,
1089 &mut system_evt_rx,
1090 &mut system_cmd_rx,
1091 &mut exec_evt_rx,
1092 &mut exec_cmd_rx,
1093 &mut data_evt_rx,
1094 &mut data_cmd_rx,
1095 );
1096 log::info!("Event loop stopped");
1097 return result;
1098 }
1099 };
1100
1101 if engine_connection_status == EngineConnectionStatus::TimedOut {
1102 let result = self
1103 .abort_startup_with_error(
1104 "Engine readiness timed out",
1105 anyhow::anyhow!("readiness timeout while waiting for engine connections"),
1106 )
1107 .await;
1108 Self::drain_channels(
1109 &mut time_evt_rx,
1110 &mut system_evt_rx,
1111 &mut system_cmd_rx,
1112 &mut exec_evt_rx,
1113 &mut exec_cmd_rx,
1114 &mut data_evt_rx,
1115 &mut data_cmd_rx,
1116 );
1117 log::info!("Event loop stopped");
1118 return result;
1119 }
1120
1121 if let Some(reason) = engine_connection_status
1122 .abort_reason()
1123 .or_else(|| self.startup_abort_reason())
1124 {
1125 self.abort_startup(reason).await?;
1126 Self::drain_channels(
1127 &mut time_evt_rx,
1128 &mut system_evt_rx,
1129 &mut system_cmd_rx,
1130 &mut exec_evt_rx,
1131 &mut exec_cmd_rx,
1132 &mut data_evt_rx,
1133 &mut data_cmd_rx,
1134 );
1135 log::info!("Event loop stopped");
1136 return Ok(());
1137 }
1138
1139 debug_assert_eq!(engine_connection_status, EngineConnectionStatus::Connected);
1140
1141 if let Err(e) = self.perform_startup_reconciliation().await {
1143 let result = self.abort_startup("Startup reconciliation failed").await;
1144 Self::drain_channels(
1145 &mut time_evt_rx,
1146 &mut system_evt_rx,
1147 &mut system_cmd_rx,
1148 &mut exec_evt_rx,
1149 &mut exec_cmd_rx,
1150 &mut data_evt_rx,
1151 &mut data_cmd_rx,
1152 );
1153 log::info!("Event loop stopped");
1154
1155 if let Err(finalize_err) = result {
1156 anyhow::bail!(
1157 "startup reconciliation failed: {e}; failed to finalize startup abort: {finalize_err}"
1158 );
1159 }
1160
1161 return Err(e);
1162 }
1163
1164 if let Some(reason) = self.startup_abort_reason() {
1165 let result = self.abort_startup(reason).await;
1166 Self::drain_channels(
1167 &mut time_evt_rx,
1168 &mut system_evt_rx,
1169 &mut system_cmd_rx,
1170 &mut exec_evt_rx,
1171 &mut exec_cmd_rx,
1172 &mut data_evt_rx,
1173 &mut data_cmd_rx,
1174 );
1175 log::info!("Event loop stopped");
1176 return result;
1177 }
1178
1179 if let Err(e) = self.kernel.start_trader() {
1180 let result = self.abort_after_trader_start_failure(e).await;
1181 Self::drain_channels(
1182 &mut time_evt_rx,
1183 &mut system_evt_rx,
1184 &mut system_cmd_rx,
1185 &mut exec_evt_rx,
1186 &mut exec_cmd_rx,
1187 &mut data_evt_rx,
1188 &mut data_cmd_rx,
1189 );
1190 log::info!("Event loop stopped");
1191 return result;
1192 }
1193 #[cfg(feature = "plugin")]
1194 if let Err(e) = self.plugins.start_controllers() {
1195 let result = self.abort_after_trader_start_failure(e).await;
1196 Self::drain_channels(
1197 &mut time_evt_rx,
1198 &mut system_evt_rx,
1199 &mut system_cmd_rx,
1200 &mut exec_evt_rx,
1201 &mut exec_cmd_rx,
1202 &mut data_evt_rx,
1203 &mut data_cmd_rx,
1204 );
1205 log::info!("Event loop stopped");
1206 return result;
1207 }
1208
1209 self.process_system_events(startup_system_events);
1210 self.process_system_commands(startup_system_commands);
1211
1212 let finish_result = {
1213 let mut receivers = RunnerReceivers {
1214 time_evt: &mut time_evt_rx,
1215 system_evt: &mut system_evt_rx,
1216 system_cmd: &mut system_cmd_rx,
1217 data_evt: &mut data_evt_rx,
1218 data_cmd: &mut data_cmd_rx,
1219 exec_evt: &mut exec_evt_rx,
1220 exec_cmd: &mut exec_cmd_rx,
1221 };
1222 self.finish_startup_trader(Some(&mut receivers)).await
1223 };
1224
1225 match finish_result {
1226 Ok(true) => {}
1227 result => {
1228 log::info!("Event loop stopped");
1229 return result.map(|_| ());
1230 }
1231 }
1232
1233 let exec_config = &self.config.exec_engine;
1234 let inflight_interval_ns =
1235 u64::from(exec_config.inflight_check_interval_ms) * NANOSECONDS_IN_MILLISECOND;
1236 let open_interval_ns = exec_config
1237 .open_check_interval_secs
1238 .filter(|&s| s > 0.0)
1239 .map_or(0, secs_to_nanos_unchecked);
1240 let position_interval_ns = exec_config
1241 .position_check_interval_secs
1242 .filter(|&s| s > 0.0)
1243 .map_or(0, secs_to_nanos_unchecked);
1244 let has_clients = !self
1245 .kernel
1246 .exec_engine
1247 .borrow()
1248 .get_all_clients()
1249 .is_empty();
1250 let recon_enabled = has_clients
1251 && (inflight_interval_ns > 0 || open_interval_ns > 0 || position_interval_ns > 0);
1252
1253 let recon_min_interval = if recon_enabled {
1254 let mut intervals = Vec::new();
1255
1256 if exec_config.inflight_check_interval_ms > 0 {
1257 intervals.push(Duration::from_millis(u64::from(
1258 exec_config.inflight_check_interval_ms,
1259 )));
1260 }
1261
1262 if let Some(s) = exec_config.open_check_interval_secs.filter(|&s| s > 0.0) {
1263 intervals.push(Duration::from_secs_f64(s));
1264 }
1265
1266 if let Some(s) = exec_config
1267 .position_check_interval_secs
1268 .filter(|&s| s > 0.0)
1269 {
1270 intervals.push(Duration::from_secs_f64(s));
1271 }
1272
1273 intervals
1274 .into_iter()
1275 .min()
1276 .unwrap_or(Duration::from_secs(1))
1277 } else {
1278 Duration::from_secs(1) };
1280
1281 let startup_delay = if self.config.exec_engine.reconciliation {
1286 Duration::from_secs_f64(exec_config.reconciliation_startup_delay_secs)
1287 } else {
1288 Duration::ZERO
1289 };
1290
1291 let recon_start = dst::time::Instant::now() + startup_delay;
1292
1293 let mut last_inflight_check = dst::time::Instant::now();
1294 let mut last_open_check = last_inflight_check;
1295 let mut last_position_check = last_inflight_check;
1296
1297 let far_future = Duration::from_hours(24 * 365 * 100);
1300
1301 let make_schedule = |opt_dur: Option<Duration>| -> (Duration, dst::time::Instant) {
1302 let dur = opt_dur.unwrap_or(far_future);
1303 (dur, recon_start + dur)
1304 };
1305
1306 let (recon_interval, mut recon_next) = make_schedule(if recon_enabled {
1307 Some(recon_min_interval)
1308 } else {
1309 None
1310 });
1311
1312 let (purge_orders_interval, mut purge_orders_next) = make_schedule(
1313 exec_config
1314 .purge_closed_orders_interval_mins
1315 .filter(|&m| m > 0)
1316 .map(|m| Duration::from_secs(mins_to_secs(u64::from(m)))),
1317 );
1318
1319 let (purge_positions_interval, mut purge_positions_next) = make_schedule(
1320 exec_config
1321 .purge_closed_positions_interval_mins
1322 .filter(|&m| m > 0)
1323 .map(|m| Duration::from_secs(mins_to_secs(u64::from(m)))),
1324 );
1325
1326 let (purge_account_interval, mut purge_account_next) = make_schedule(
1327 exec_config
1328 .purge_account_events_interval_mins
1329 .filter(|&m| m > 0)
1330 .map(|m| Duration::from_secs(mins_to_secs(u64::from(m)))),
1331 );
1332
1333 let (own_books_interval, mut own_books_next) = make_schedule(
1334 exec_config
1335 .own_books_audit_interval_secs
1336 .filter(|&s| s > 0.0)
1337 .map(Duration::from_secs_f64),
1338 );
1339
1340 let (prune_fills_interval, mut prune_fills_next) =
1341 make_schedule(Some(Duration::from_mins(1)));
1342
1343 let mut maintenance_timer = dst::time::interval(Duration::from_millis(100));
1344 maintenance_timer.set_missed_tick_behavior(dst::time::MissedTickBehavior::Skip);
1345
1346 let mut stop_check_timer = dst::time::interval(Duration::from_millis(100));
1351 stop_check_timer.set_missed_tick_behavior(dst::time::MissedTickBehavior::Skip);
1352
1353 let mut residual_events = 0usize;
1355 let mut open_order_report_task: Option<OpenOrderReportTask> = None;
1356 let mut targeted_order_report_task: Option<TargetedOrderReportTask> = None;
1357 let mut position_report_task: Option<PositionReportTask> = None;
1358
1359 let owns_signals = mode.owns_signals();
1362
1363 let ctrl_c = async move {
1364 if owns_signals {
1365 dst::signal::ctrl_c().await
1366 } else {
1367 std::future::pending::<std::io::Result<()>>().await
1368 }
1369 };
1370
1371 let terminate = async move {
1372 if owns_signals {
1373 dst::signal::terminate().await
1374 } else {
1375 std::future::pending::<std::io::Result<()>>().await
1376 }
1377 };
1378
1379 tokio::pin!(ctrl_c);
1380 tokio::pin!(terminate);
1381
1382 let metrics = self.handle.metrics.clone();
1383 let metrics_start = dst::time::Instant::now();
1384 metrics.reset();
1385
1386 let mut queue_monitor = self
1387 .config
1388 .queue_monitor
1389 .as_ref()
1390 .map(|config| QueueMonitor::new(config, metrics.snapshot()));
1391 let mut dispatches_since_yield = 0usize;
1392
1393 loop {
1394 let shutdown_deadline = self.shutdown_deadline;
1395 let is_shutting_down = self.state() == NodeState::ShuttingDown;
1396 let is_running = self.state() == NodeState::Running;
1397
1398 tokio::select! {
1399 biased;
1400
1401 result = &mut ctrl_c, if is_running => {
1403 match result {
1404 Ok(()) => log::info!("Received SIGINT, shutting down"),
1405 Err(e) => log::error!("Failed to listen for SIGINT: {e}"),
1406 }
1407 self.initiate_shutdown();
1408 }
1409 result = &mut terminate, if is_running => {
1410 match result {
1411 Ok(()) => log::info!("Received SIGTERM, shutting down"),
1412 Err(e) => log::error!("Failed to listen for SIGTERM: {e}"),
1413 }
1414 self.initiate_shutdown();
1415 }
1416 _ = stop_check_timer.tick(), if is_running => {
1417 if stop_handle.should_stop() {
1418 log::info!("Received stop signal from handle");
1419 self.initiate_shutdown();
1420 } else if self.kernel.is_shutdown_requested() {
1421 log::info!("Received ShutdownSystem command, shutting down");
1422 self.initiate_shutdown();
1423 }
1424 }
1425 () = async {
1426 match shutdown_deadline {
1427 Some(deadline) => dst::time::sleep_until(deadline).await,
1428 None => std::future::pending::<()>().await,
1429 }
1430 }, if self.state() == NodeState::ShuttingDown => {
1431 break;
1432 }
1433 result = async {
1434 match open_order_report_task.as_mut() {
1435 Some(task) => task.future.as_mut().await,
1436 None => std::future::pending::<ReportTaskOutcome<OpenOrderReportResult>>().await,
1437 }
1438 }, if open_order_report_task.is_some() => {
1439 let maintenance_start = dst::time::Instant::now();
1440
1441 drop(open_order_report_task.take());
1442
1443 match result {
1444 ReportTaskOutcome::Completed(result) => {
1445 let client_refs = self
1446 .exec_clients
1447 .iter()
1448 .map(|client| client as &dyn ExecutionClient)
1449 .collect::<Vec<_>>();
1450 let reconciliation = self.exec_manager.reconcile_open_order_reports(
1451 &result.check,
1452 result.reports,
1453 &result.queried_clients,
1454 &result.failed_clients,
1455 &client_refs,
1456 );
1457 self.process_reconciliation_events(&reconciliation.events);
1458 if !reconciliation.targeted_queries.is_empty() {
1459 targeted_order_report_task = Some(
1460 self.start_targeted_order_report_check(
1461 reconciliation.targeted_queries,
1462 ),
1463 );
1464 }
1465 }
1466 ReportTaskOutcome::TimedOut => {
1467 self.cleanup_cancelled_report_tasks(&[]);
1468 log::warn!(
1469 "Open-order report collection expired after {:?}",
1470 self.config.timeout_reconciliation,
1471 );
1472 }
1473 }
1474 record_runner_maintenance(&metrics, maintenance_start, metrics_start);
1475 }
1476 result = async {
1477 match targeted_order_report_task.as_mut() {
1478 Some(task) => task.future.as_mut().await,
1479 None => std::future::pending::<ReportTaskOutcome<Vec<TargetedOrderReportResult>>>().await,
1480 }
1481 }, if targeted_order_report_task.is_some() => {
1482 let maintenance_start = dst::time::Instant::now();
1483
1484 let planned_client_order_ids = targeted_order_report_task
1485 .as_ref()
1486 .map(|task| task.planned_client_order_ids.clone())
1487 .unwrap_or_default();
1488 drop(targeted_order_report_task.take());
1489
1490 match result {
1491 ReportTaskOutcome::Completed(result) => {
1492 let client_refs = self
1493 .exec_clients
1494 .iter()
1495 .map(|client| client as &dyn ExecutionClient)
1496 .collect::<Vec<_>>();
1497 let events = self
1498 .exec_manager
1499 .reconcile_targeted_order_reports(result, &client_refs);
1500 self.process_reconciliation_events(&events);
1501 }
1502 ReportTaskOutcome::TimedOut => {
1503 self.cleanup_cancelled_report_tasks(&planned_client_order_ids);
1504 log::warn!(
1505 "Targeted order report collection expired after {:?}",
1506 self.config.timeout_reconciliation,
1507 );
1508 }
1509 }
1510 record_runner_maintenance(&metrics, maintenance_start, metrics_start);
1511 }
1512 result = async {
1513 match position_report_task.as_mut() {
1514 Some(task) => task.future.as_mut().await,
1515 None => std::future::pending::<ReportTaskOutcome<PositionReportResult>>().await,
1516 }
1517 }, if position_report_task.is_some() => {
1518 let maintenance_start = dst::time::Instant::now();
1519
1520 drop(position_report_task.take());
1521
1522 match result {
1523 ReportTaskOutcome::Completed(result) => {
1524 let events = self.exec_manager.reconcile_position_reports(
1525 &result.check,
1526 result.reports,
1527 &result.queried_clients,
1528 &result.failed_clients,
1529 );
1530 self.process_reconciliation_events(&events);
1531 }
1532 ReportTaskOutcome::TimedOut => {
1533 self.cleanup_cancelled_report_tasks(&[]);
1534 log::warn!(
1535 "Position report collection expired after {:?}",
1536 self.config.timeout_reconciliation,
1537 );
1538 }
1539 }
1540 record_runner_maintenance(&metrics, maintenance_start, metrics_start);
1541 }
1542
1543 _ = maintenance_timer.tick(), if is_running => {
1546 let maintenance_start = dst::time::Instant::now();
1547 metrics.publish_queue_depths(
1548 RunnerChannelQueueDepths::from_receivers(
1549 &time_evt_rx,
1550 &exec_evt_rx,
1551 &exec_cmd_rx,
1552 &data_evt_rx,
1553 &data_cmd_rx,
1554 ),
1555 metrics_start.elapsed(),
1556 );
1557
1558 if let Some(queue_monitor) = queue_monitor.as_mut() {
1559 let transitions = queue_monitor.evaluate(metrics.snapshot());
1560 self.publish_queue_state_transitions(&transitions);
1561 }
1562
1563 let mut now = dst::time::Instant::now();
1564
1565 if recon_enabled && now >= recon_next {
1566 let recon_intervals = ReconciliationCheckIntervals {
1567 inflight: Duration::from_nanos(inflight_interval_ns),
1568 open: Duration::from_nanos(open_interval_ns),
1569 position: Duration::from_nanos(position_interval_ns),
1570 };
1571 let mut recon_state = ReconciliationCheckState {
1572 last_inflight_check: &mut last_inflight_check,
1573 last_open_check: &mut last_open_check,
1574 last_position_check: &mut last_position_check,
1575 open_order_report_task: &mut open_order_report_task,
1576 targeted_order_report_task: &mut targeted_order_report_task,
1577 position_report_task: &mut position_report_task,
1578 };
1579
1580 self.run_reconciliation_checks(
1581 now,
1582 recon_intervals,
1583 &mut recon_state,
1584 );
1585
1586 now = dst::time::Instant::now();
1587 recon_next = now + recon_interval;
1588 }
1589
1590 if now >= purge_orders_next {
1591 self.exec_manager.purge_closed_orders();
1592 purge_orders_next = now + purge_orders_interval;
1593 }
1594
1595 if now >= purge_positions_next {
1596 self.exec_manager.purge_closed_positions();
1597 purge_positions_next = now + purge_positions_interval;
1598 }
1599
1600 if now >= purge_account_next {
1601 self.exec_manager.purge_account_events();
1602 purge_account_next = now + purge_account_interval;
1603 }
1604
1605 if now >= own_books_next {
1606 self.kernel.cache().borrow_mut().audit_own_order_books();
1607 own_books_next = now + own_books_interval;
1608 }
1609
1610 if now >= prune_fills_next {
1611 self.exec_manager.prune_recent_fills_cache(60.0);
1612 self.exec_manager.prune_processed_fills();
1613 self.exec_manager.prune_order_local_activity();
1614 prune_fills_next = now + prune_fills_interval;
1615 }
1616
1617 record_runner_maintenance(&metrics, maintenance_start, metrics_start);
1618 }
1619
1620 Some(handler) = time_evt_rx.recv() => {
1625 let dispatch_start = dst::time::Instant::now();
1626 let dispatched = AsyncRunner::handle_time_event(handler);
1627
1628 if dispatched && is_shutting_down {
1629 log::debug!("Residual time event");
1630 residual_events += 1;
1631 }
1632
1633 if dispatched {
1634 record_runner_dispatch(
1635 &metrics,
1636 SystemChannel::TimeEvents,
1637 dispatch_start,
1638 metrics_start,
1639 );
1640 }
1641 }
1642 Some(event) = system_evt_rx.recv() => {
1643 if is_shutting_down {
1644 log::debug!("Residual system event: {event:?}");
1645 residual_events += 1;
1646 }
1647 self.process_system_event(event);
1648 }
1649 Some(command) = system_cmd_rx.recv() => {
1650 if is_shutting_down {
1651 log::debug!("Residual system command: {command:?}");
1652 residual_events += 1;
1653 }
1654 self.process_system_command(command);
1655 }
1656 Some(evt) = exec_evt_rx.recv() => {
1657 let dispatch_start = dst::time::Instant::now();
1658
1659 if is_shutting_down {
1660 log::debug!("Residual exec event: {evt:?}");
1661 residual_events += 1;
1662 }
1663
1664 self.process_exec_event(evt);
1665 record_runner_dispatch(
1666 &metrics,
1667 SystemChannel::ExecEvents,
1668 dispatch_start,
1669 metrics_start,
1670 );
1671 }
1672 Some(cmd) = exec_cmd_rx.recv() => {
1673 let dispatch_start = dst::time::Instant::now();
1674
1675 if is_shutting_down {
1676 log::debug!("Residual exec command: {cmd:?}");
1677 residual_events += 1;
1678 }
1679
1680 self.process_exec_command(cmd);
1681 record_runner_dispatch(
1682 &metrics,
1683 SystemChannel::ExecCommands,
1684 dispatch_start,
1685 metrics_start,
1686 );
1687 }
1688 message = recv_external_msgbus_message(&mut external_msgbus_rx) => {
1689 let external_msgbus_start = dst::time::Instant::now();
1690
1691 match message {
1692 Some(message) => {
1693 if is_shutting_down {
1694 log::debug!("Residual external message bus message: {message}");
1695 residual_events += 1;
1696 }
1697 Self::republish_external_msgbus_message(&message);
1698 }
1699 None => {
1700 log::info!("External message bus ingress closed");
1701 external_msgbus_rx = None;
1702 self.close_external_ingress();
1703 }
1704 }
1705
1706 record_runner_external_msgbus(
1707 &metrics,
1708 external_msgbus_start,
1709 metrics_start,
1710 );
1711 }
1712 Some(evt) = data_evt_rx.recv() => {
1713 let dispatch_start = dst::time::Instant::now();
1714
1715 if is_shutting_down {
1716 log::debug!("Residual data event: {evt:?}");
1717 residual_events += 1;
1718 }
1719 AsyncRunner::handle_data_event(evt);
1720 record_runner_dispatch(
1721 &metrics,
1722 SystemChannel::DataEvents,
1723 dispatch_start,
1724 metrics_start,
1725 );
1726 }
1727 Some(cmd) = data_cmd_rx.recv() => {
1728 let dispatch_start = dst::time::Instant::now();
1729
1730 if is_shutting_down {
1731 log::debug!("Residual data command: {cmd:?}");
1732 residual_events += 1;
1733 }
1734 AsyncRunner::handle_data_command(cmd);
1735 record_runner_dispatch(
1736 &metrics,
1737 SystemChannel::DataCommands,
1738 dispatch_start,
1739 metrics_start,
1740 );
1741 }
1742 }
1743
1744 dispatches_since_yield += 1;
1745 if dispatches_since_yield >= DISPATCHES_PER_YIELD {
1746 dispatches_since_yield = 0;
1747 tokio::task::yield_now().await;
1748 }
1749 }
1750
1751 if residual_events > 0 {
1752 log::debug!("Processed {residual_events} residual events during shutdown");
1753 }
1754
1755 self.cancel_report_tasks(
1756 &mut open_order_report_task,
1757 &mut targeted_order_report_task,
1758 &mut position_report_task,
1759 );
1760 drop(external_msgbus_rx.take());
1761 let _ = self.kernel.cache().borrow().check_residuals();
1762
1763 let stop_result = self.finalize_stop().await;
1764
1765 Self::drain_channels(
1767 &mut time_evt_rx,
1768 &mut system_evt_rx,
1769 &mut system_cmd_rx,
1770 &mut exec_evt_rx,
1771 &mut exec_cmd_rx,
1772 &mut data_evt_rx,
1773 &mut data_cmd_rx,
1774 );
1775
1776 log::info!("Event loop stopped");
1777
1778 stop_result
1779 }
1780
1781 fn publish_queue_state_transitions(&self, transitions: &[QueueStateTransition]) {
1782 let topic = MessagingSwitchboard::queue_state_changed_topic();
1783
1784 for transition in transitions {
1785 let timestamp = self.kernel.generate_timestamp_ns();
1786 let event = QueueStateChanged::new(
1787 self.config.trader_id,
1788 transition.channel,
1789 transition.condition,
1790 transition.state,
1791 transition.queue_depth,
1792 transition.mean_dispatch_ns,
1793 UUID4::new(),
1794 timestamp,
1795 timestamp,
1796 );
1797
1798 msgbus::publish_any(topic, event.as_any());
1799 }
1800 }
1801
1802 #[expect(
1803 clippy::await_holding_refcell_ref,
1804 reason = "cache loading is serialized before the single-threaded live node starts"
1805 )]
1806 async fn prepare_cache(&mut self) -> anyhow::Result<()> {
1807 self.install_cache_database().await?;
1808
1809 let cache = self.kernel.cache();
1810 if !cache.borrow().has_backing() {
1811 return Ok(());
1812 }
1813
1814 if self
1815 .config
1816 .cache
1817 .as_ref()
1818 .is_some_and(|config| config.flush_on_start)
1819 {
1820 cache.borrow_mut().flush_db();
1821 return Ok(());
1822 }
1823
1824 if self.config.exec_engine.load_cache {
1825 self.kernel
1826 .exec_engine()
1827 .borrow_mut()
1828 .load_cache()
1829 .await
1830 .context("Failed to load persistent cache")?;
1831 }
1832
1833 Ok(())
1834 }
1835
1836 #[must_use]
1838 pub const fn has_pending_cache_database(&self) -> bool {
1839 self.cache_database_factory.is_some()
1840 }
1841
1842 async fn install_cache_database(&mut self) -> anyhow::Result<()> {
1847 let Some(factory) = self.cache_database_factory.as_ref() else {
1848 return Ok(());
1849 };
1850
1851 let config = self.config.cache.clone().unwrap_or_default();
1852
1853 let database = factory
1856 .create(self.config.trader_id, self.kernel.instance_id, config)
1857 .await
1858 .context("failed to create cache database backing")?;
1859 self.kernel.cache().borrow_mut().set_database(database);
1860 self.cache_database_factory = None;
1861
1862 Ok(())
1863 }
1864
1865 fn take_external_ingress_receiver(
1866 &mut self,
1867 ) -> anyhow::Result<Option<tokio::sync::mpsc::Receiver<BusMessage>>> {
1868 let Some(external_ingress) = self.external_msgbus.as_mut() else {
1869 return Ok(None);
1870 };
1871
1872 let receiver = external_ingress.take_receiver()?;
1873 log::info!("External message bus ingress started");
1874 Ok(Some(receiver))
1875 }
1876
1877 fn republish_external_msgbus_message(message: &BusMessage) {
1878 if let Err(e) = msgbus::republish_external_message(message) {
1879 log::error!(
1880 "Failed to republish external message bus topic '{}': {e:#}",
1881 message.topic
1882 );
1883 }
1884 }
1885
1886 fn close_external_ingress(&mut self) {
1887 if let Some(external_ingress) = self.external_msgbus.as_mut()
1888 && !external_ingress.is_closed()
1889 {
1890 external_ingress.close();
1891 }
1892 }
1893
1894 fn process_reconciliation_events(&mut self, events: &[OrderEventAny]) {
1895 if events.is_empty() {
1896 return;
1897 }
1898
1899 log::info!(
1900 "Processing {} reconciliation event{}",
1901 events.len(),
1902 if events.len() == 1 { "" } else { "s" }
1903 );
1904
1905 for event in events {
1906 self.exec_manager
1907 .record_local_activity(event.client_order_id());
1908 if let OrderEventAny::Filled(fill) = event {
1909 self.exec_manager
1910 .record_position_activity(fill.instrument_id, fill.account_id);
1911 }
1912 self.kernel.exec_engine.borrow_mut().process(event);
1913 if let OrderEventAny::Filled(fill) = event {
1914 self.exec_manager.commit_recent_fill_if_applied(fill);
1915 }
1916 }
1917 }
1918
1919 fn process_exec_event(&mut self, event: ExecutionEvent) {
1920 let Some(close_ids) = self.observe_exec_event_before_dispatch(&event) else {
1921 return;
1922 };
1923
1924 self.dispatch_exec_event_and_commit_fill(event);
1925
1926 for client_order_id in &close_ids {
1927 let is_closed = self
1928 .kernel
1929 .cache()
1930 .borrow()
1931 .order(client_order_id)
1932 .is_some_and(|order| order.is_closed());
1933 if is_closed {
1934 self.exec_manager
1935 .clear_recon_tracking(client_order_id, true);
1936 }
1937 }
1938 }
1939
1940 fn process_exec_command(&mut self, message: TradingCommandMessage) {
1941 let mut messages = vec![message];
1942 while let Some(message) = messages.pop() {
1943 if message.endpoint() == MessagingSwitchboard::exec_engine_execute() {
1944 self.observe_exec_command_before_dispatch(message.command());
1945 }
1946 messages.extend(message.dispatch().into_iter().rev());
1947 }
1948 }
1949
1950 fn dispatch_exec_event_and_commit_fill(&mut self, evt: ExecutionEvent) {
1959 let recent_fill_candidate = match &evt {
1960 ExecutionEvent::Order(OrderEventAny::Filled(fill)) => Some(fill.clone()),
1961 _ => None,
1962 };
1963
1964 AsyncRunner::handle_exec_event(evt);
1965
1966 if let Some(fill) = &recent_fill_candidate {
1967 self.exec_manager.commit_recent_fill_if_applied(fill);
1968 }
1969 }
1970
1971 async fn connect_data_phase(&mut self, deadline: dst::time::Instant) -> anyhow::Result<()> {
1972 let remaining = deadline.saturating_duration_since(dst::time::Instant::now());
1978 dst::time::timeout(remaining, self.kernel.connect_data_clients())
1979 .await
1980 .map_err(|_| anyhow::anyhow!("data-connect timeout"))
1981 }
1982
1983 async fn connect_exec_clients(&mut self, deadline: dst::time::Instant) -> anyhow::Result<()> {
1984 let remaining = deadline.saturating_duration_since(dst::time::Instant::now());
1985 dst::time::timeout(remaining, self.kernel.connect_exec_clients())
1986 .await
1987 .map_err(|_| anyhow::anyhow!("exec-connect timeout"))
1988 }
1989
1990 async fn connect_exec_phase(
1995 &mut self,
1996 deadline: dst::time::Instant,
1997 ) -> anyhow::Result<EngineConnectionStatus> {
1998 self.connect_exec_clients(deadline).await?;
1999 Ok(self.await_engines_connected(deadline).await)
2000 }
2001
2002 fn startup_abort_reason(&self) -> Option<&'static str> {
2003 if self.handle.should_stop() {
2004 Some("Stop signal received during startup")
2005 } else if self.kernel.is_shutdown_requested() {
2006 Some("Shutdown signal received during startup")
2007 } else {
2008 None
2009 }
2010 }
2011
2012 async fn finish_startup_replay(&mut self) -> anyhow::Result<bool> {
2013 match self.handle.try_set_running() {
2014 RunningTransition::Entered => Ok(true),
2015 RunningTransition::StopRequested => {
2016 self.abort_startup("Stop signal received during startup")
2017 .await?;
2018 Ok(false)
2019 }
2020 RunningTransition::Invalid(control) => {
2021 self.abort_startup_with_error(
2022 "Invalid lifecycle state during startup",
2023 anyhow::anyhow!(
2024 "Invalid LiveNode control state {control:#04x} while entering Running"
2025 ),
2026 )
2027 .await?;
2028 Ok(false)
2029 }
2030 }
2031 }
2032
2033 async fn finish_startup_trader(
2034 &mut self,
2035 receivers: Option<&mut RunnerReceivers<'_>>,
2036 ) -> anyhow::Result<bool> {
2037 match self.handle.try_set_running() {
2038 RunningTransition::Entered => Ok(true),
2039 RunningTransition::StopRequested => {
2040 self.abort_started_trader("Stop signal received during startup", receivers)
2041 .await?;
2042 Ok(false)
2043 }
2044 RunningTransition::Invalid(control) => {
2045 let state_err = anyhow::anyhow!(
2046 "Invalid LiveNode control state {control:#04x} while entering Running"
2047 );
2048
2049 match self
2050 .abort_started_trader("Invalid lifecycle state during startup", receivers)
2051 .await
2052 {
2053 Ok(()) => Err(state_err),
2054 Err(finalize_err) => {
2055 anyhow::bail!(
2056 "{state_err}; failed to finalize startup abort: {finalize_err}"
2057 )
2058 }
2059 }
2060 }
2061 }
2062 }
2063
2064 async fn abort_startup(&mut self, reason: &str) -> anyhow::Result<()> {
2065 log::info!("{reason}, aborting startup");
2066 self.handle.set_shutting_down();
2067 self.finalize_stop().await
2068 }
2069
2070 async fn abort_startup_with_error(
2071 &mut self,
2072 reason: &str,
2073 startup_err: anyhow::Error,
2074 ) -> anyhow::Result<()> {
2075 match self.abort_startup(reason).await {
2076 Ok(()) => Err(startup_err),
2077 Err(finalize_err) => {
2078 anyhow::bail!("{startup_err}; failed to finalize startup abort: {finalize_err}")
2079 }
2080 }
2081 }
2082
2083 async fn abort_started_trader(
2084 &mut self,
2085 reason: &str,
2086 mut receivers: Option<&mut RunnerReceivers<'_>>,
2087 ) -> anyhow::Result<()> {
2088 log::info!("{reason}, aborting startup");
2089 self.handle.set_shutting_down();
2090
2091 #[cfg(feature = "plugin")]
2092 let controller_stop_result = self.plugins.stop_controllers();
2093 #[cfg(not(feature = "plugin"))]
2094 let controller_stop_result: anyhow::Result<()> = Ok(());
2095
2096 let trader_stop_result = self.kernel.stop_trader_after_start_failure();
2097 let delay = self.kernel.delay_post_stop();
2098 log::info!("Awaiting residual events ({delay:?})...");
2099
2100 let residual_events = match receivers.as_mut() {
2101 Some(receivers) => self.process_receivers_for(delay, receivers).await,
2102 None => self.process_runner_for(delay).await,
2103 };
2104
2105 if residual_events > 0 {
2106 log::debug!("Processed {residual_events} residual events during shutdown");
2107 }
2108
2109 let finalize_result = self.finalize_stop().await;
2110
2111 if let Some(receivers) = receivers {
2112 Self::drain_channels(
2113 receivers.time_evt,
2114 receivers.system_evt,
2115 receivers.system_cmd,
2116 receivers.exec_evt,
2117 receivers.exec_cmd,
2118 receivers.data_evt,
2119 receivers.data_cmd,
2120 );
2121 } else {
2122 let drained_events = self.drain_runner_pending();
2123 if drained_events > 0 {
2124 log::info!("Drained {drained_events} remaining events during shutdown");
2125 }
2126 }
2127
2128 let mut errors = Vec::new();
2129
2130 if let Err(e) = controller_stop_result {
2131 errors.push(format!("Failed to stop plug-in controllers: {e}"));
2132 }
2133
2134 if let Err(e) = trader_stop_result {
2135 errors.push(format!("Failed to stop trader: {e}"));
2136 }
2137
2138 if let Err(e) = finalize_result {
2139 errors.push(format!("Failed to finalize startup abort: {e}"));
2140 }
2141
2142 if errors.is_empty() {
2143 Ok(())
2144 } else {
2145 anyhow::bail!("{}", errors.join("; "))
2146 }
2147 }
2148
2149 async fn process_receivers_for(
2150 &mut self,
2151 duration: Duration,
2152 receivers: &mut RunnerReceivers<'_>,
2153 ) -> usize {
2154 let deadline = dst::time::Instant::now() + duration;
2155 let mut processed = 0;
2156
2157 loop {
2158 tokio::select! {
2159 biased;
2160
2161 () = dst::time::sleep_until(deadline) => break,
2162 Some(message) = receivers.time_evt.recv() => {
2163 let _ = AsyncRunner::handle_time_event(message);
2164 processed += 1;
2165 }
2166 Some(event) = receivers.system_evt.recv() => {
2167 self.process_system_event(event);
2168 processed += 1;
2169 }
2170 Some(command) = receivers.system_cmd.recv() => {
2171 self.process_system_command(command);
2172 processed += 1;
2173 }
2174 Some(event) = receivers.exec_evt.recv() => {
2175 self.process_exec_event(event);
2176 processed += 1;
2177 }
2178 Some(command) = receivers.exec_cmd.recv() => {
2179 self.process_exec_command(command);
2180 processed += 1;
2181 }
2182 Some(event) = receivers.data_evt.recv() => {
2183 AsyncRunner::handle_data_event(event);
2184 processed += 1;
2185 }
2186 Some(command) = receivers.data_cmd.recv() => {
2187 AsyncRunner::handle_data_command(command);
2188 processed += 1;
2189 }
2190 }
2191 }
2192
2193 processed
2194 }
2195
2196 async fn abort_after_trader_start_failure(
2197 &mut self,
2198 start_err: anyhow::Error,
2199 ) -> anyhow::Result<()> {
2200 log::info!("Trader startup failed, aborting startup");
2201 self.handle.set_shutting_down();
2202 let stop_result = self.kernel.stop_trader_after_start_failure();
2203 let finalize_result = self.finalize_stop().await;
2204
2205 match (stop_result, finalize_result) {
2206 (Ok(()), Ok(())) => Err(start_err),
2207 (Err(stop_err), Ok(())) => anyhow::bail!(
2208 "Failed during trader startup: {start_err}; failed to stop partial trader start: \
2209 {stop_err}"
2210 ),
2211 (Ok(()), Err(finalize_err)) => anyhow::bail!(
2212 "Failed during trader startup: {start_err}; failed to finalize startup abort: \
2213 {finalize_err}"
2214 ),
2215 (Err(stop_err), Err(finalize_err)) => anyhow::bail!(
2216 "Failed during trader startup: {start_err}; failed to stop partial trader start: \
2217 {stop_err}; failed to finalize startup abort: {finalize_err}"
2218 ),
2219 }
2220 }
2221
2222 fn initiate_shutdown(&mut self) {
2223 #[cfg(feature = "plugin")]
2224 if let Err(e) = self.plugins.stop_controllers() {
2225 log::error!("Error stopping plug-in controllers: {e}");
2226 }
2227 self.kernel.stop_trader();
2228 let delay = self.kernel.delay_post_stop();
2229 log::info!("Awaiting residual events ({delay:?})...");
2230
2231 self.shutdown_deadline = Some(dst::time::Instant::now() + delay);
2232 self.handle.set_shutting_down();
2233 }
2234
2235 async fn finalize_stop(&mut self) -> anyhow::Result<()> {
2236 self.close_external_ingress();
2237
2238 let timeout = self.config.timeout_disconnection;
2239 let deadline = dst::time::Instant::now() + timeout;
2240 let disconnect_result =
2241 match dst::time::timeout(timeout, self.kernel.disconnect_clients()).await {
2242 Ok(result) => result,
2243 Err(_) => Err(anyhow::anyhow!(
2244 "disconnect timeout while disconnecting clients"
2245 )),
2246 };
2247
2248 if let Err(ref e) = disconnect_result {
2249 log::error!("Error disconnecting clients: {e}");
2250 }
2251
2252 let readiness_result = self.await_engines_disconnected(deadline).await;
2253 let kernel_result = self.kernel.finalize_stop().await;
2254
2255 self.handle.set_stopped();
2256
2257 let mut errors = Vec::new();
2258 if let Err(e) = disconnect_result {
2259 errors.push(e.to_string());
2260 }
2261
2262 if let Err(e) = readiness_result {
2263 errors.push(format!("failed while awaiting engine disconnection: {e}"));
2264 }
2265
2266 if let Err(e) = kernel_result {
2267 errors.push(format!("failed while finalizing kernel shutdown: {e}"));
2268 }
2269
2270 if errors.is_empty() {
2271 Ok(())
2272 } else {
2273 anyhow::bail!("{}", errors.join("; "))
2274 }
2275 }
2276
2277 fn drain_channels(
2278 time_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<TimeEventMessage>,
2279 system_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<SystemEvent>,
2280 system_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<SystemCommand>,
2281 exec_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2282 exec_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<TradingCommandMessage>,
2283 data_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
2284 data_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataCommand>,
2285 ) {
2286 let mut drained = 0;
2287
2288 while let Ok(handler) = time_evt_rx.try_recv() {
2289 let _ = AsyncRunner::handle_time_event(handler);
2290 drained += 1;
2291 }
2292
2293 while system_evt_rx.try_recv().is_ok() {
2294 drained += 1;
2295 }
2296
2297 while system_cmd_rx.try_recv().is_ok() {
2298 drained += 1;
2299 }
2300
2301 while let Ok(evt) = data_evt_rx.try_recv() {
2302 AsyncRunner::handle_data_event(evt);
2303 drained += 1;
2304 }
2305
2306 while let Ok(cmd) = data_cmd_rx.try_recv() {
2307 AsyncRunner::handle_data_command(cmd);
2308 drained += 1;
2309 }
2310
2311 while let Ok(evt) = exec_evt_rx.try_recv() {
2312 AsyncRunner::handle_exec_event(evt);
2313 drained += 1;
2314 }
2315
2316 while let Ok(cmd) = exec_cmd_rx.try_recv() {
2317 AsyncRunner::handle_trading_command(cmd);
2318 drained += 1;
2319 }
2320
2321 if drained > 0 {
2322 log::info!("Drained {drained} remaining events during shutdown");
2323 }
2324 }
2325
2326 fn observe_exec_event_before_dispatch(
2327 &mut self,
2328 evt: &ExecutionEvent,
2329 ) -> Option<Vec<ClientOrderId>> {
2330 let mut close_ids = Vec::new();
2331
2332 match evt {
2333 ExecutionEvent::Order(order_evt) => {
2334 self.exec_manager.observe_order_event(order_evt);
2335 close_ids.push(order_evt.client_order_id());
2336 }
2337 ExecutionEvent::OrderSubmittedBatch(batch) => {
2338 for submitted in &batch.events {
2339 self.exec_manager
2340 .record_local_activity(submitted.client_order_id);
2341 }
2342 }
2343 ExecutionEvent::OrderAcceptedBatch(batch) => {
2344 for accepted in &batch.events {
2345 self.exec_manager
2346 .clear_recon_tracking(&accepted.client_order_id, true);
2347 self.exec_manager
2348 .record_local_activity(accepted.client_order_id);
2349 }
2350 }
2351 ExecutionEvent::OrderCanceledBatch(batch) => {
2352 for canceled in &batch.events {
2353 self.exec_manager
2354 .clear_recon_tracking(&canceled.client_order_id, true);
2355 self.exec_manager
2356 .record_local_activity(canceled.client_order_id);
2357 close_ids.push(canceled.client_order_id);
2358 }
2359 }
2360 ExecutionEvent::Report(report) => {
2361 if let ExecutionReport::Fill(fill_report) = report
2362 && self.exec_manager.is_fill_recently_processed(
2363 fill_report.account_id,
2364 fill_report.instrument_id,
2365 fill_report.trade_id,
2366 )
2367 {
2368 log::debug!(
2369 "Skipping recently processed fill report: {}",
2370 fill_report.trade_id,
2371 );
2372 return None;
2373 }
2374 self.exec_manager.observe_execution_report(report);
2375
2376 if let Some(client_order_id) = Self::closed_order_report_client_order_id(report) {
2377 close_ids.push(client_order_id);
2378 }
2379 }
2380 ExecutionEvent::Account(_) => {}
2381 }
2382
2383 Some(close_ids)
2384 }
2385
2386 fn closed_order_report_client_order_id(report: &ExecutionReport) -> Option<ClientOrderId> {
2387 match report {
2388 ExecutionReport::Order(order_report)
2389 | ExecutionReport::OrderWithFills(order_report, _)
2390 if order_report.order_status.is_closed() =>
2391 {
2392 order_report.client_order_id
2393 }
2394 _ => None,
2395 }
2396 }
2397
2398 fn observe_exec_command_before_dispatch(&mut self, cmd: &TradingCommand) {
2399 match cmd {
2400 TradingCommand::SubmitOrder(submit) => {
2401 self.exec_manager.register_inflight(submit.client_order_id);
2402 }
2403 TradingCommand::SubmitOrderList(submit) => {
2404 for order_init in &submit.order_inits {
2405 self.exec_manager
2406 .register_inflight(order_init.client_order_id);
2407 }
2408 }
2409 TradingCommand::ModifyOrder(modify) => {
2410 self.exec_manager.register_inflight(modify.client_order_id);
2411 }
2412 TradingCommand::ModifyOrders(modify) => {
2413 for child in &modify.modifies {
2414 self.exec_manager.register_inflight(child.client_order_id);
2415 }
2416 }
2417 TradingCommand::CancelOrder(cancel) => {
2418 self.exec_manager.register_inflight(cancel.client_order_id);
2419 }
2420 TradingCommand::CancelOrders(cancel) => {
2421 for child in &cancel.cancels {
2422 self.exec_manager.register_inflight(child.client_order_id);
2423 }
2424 }
2425 _ => {}
2426 }
2427 }
2428
2429 #[must_use]
2431 pub fn environment(&self) -> Environment {
2432 self.kernel.environment()
2433 }
2434
2435 #[must_use]
2437 pub const fn kernel(&self) -> &NautilusKernel {
2438 &self.kernel
2439 }
2440
2441 #[must_use]
2443 pub const fn kernel_mut(&mut self) -> &mut NautilusKernel {
2444 &mut self.kernel
2445 }
2446
2447 #[must_use]
2449 pub fn trader_id(&self) -> TraderId {
2450 self.kernel.trader_id()
2451 }
2452
2453 #[must_use]
2455 pub const fn instance_id(&self) -> UUID4 {
2456 self.kernel.instance_id()
2457 }
2458
2459 #[must_use]
2461 pub fn state(&self) -> NodeState {
2462 self.handle.state()
2463 }
2464
2465 #[must_use]
2467 pub fn is_running(&self) -> bool {
2468 self.state().is_running()
2469 }
2470
2471 pub fn set_cache_database(
2481 &mut self,
2482 database: Box<dyn CacheDatabaseAdapter>,
2483 ) -> anyhow::Result<()> {
2484 if self.state() != NodeState::Idle {
2485 anyhow::bail!(
2486 "Cannot set cache database while node is running, set it before running the node"
2487 );
2488 }
2489
2490 self.kernel.cache().borrow_mut().set_database(database);
2491 Ok(())
2492 }
2493
2494 #[must_use]
2496 pub fn exec_manager(&self) -> &ExecutionManager {
2497 &self.exec_manager
2498 }
2499
2500 #[must_use]
2502 pub fn exec_manager_mut(&mut self) -> &mut ExecutionManager {
2503 &mut self.exec_manager
2504 }
2505
2506 pub fn add_actor<T>(&mut self, actor: T) -> anyhow::Result<()>
2519 where
2520 T: DataActor + DataActorNative + Component + Actor + 'static,
2521 {
2522 if self.state() != NodeState::Idle {
2523 anyhow::bail!(
2524 "Cannot add actor while node is running, add actors before running the node"
2525 );
2526 }
2527
2528 self.kernel.trader.borrow_mut().add_actor(actor)
2529 }
2530
2531 pub fn add_actor_from_factory<F, T>(&mut self, factory: F) -> anyhow::Result<()>
2543 where
2544 F: FnOnce() -> anyhow::Result<T>,
2545 T: DataActor + DataActorNative + Component + Actor + 'static,
2546 {
2547 if self.state() != NodeState::Idle {
2548 anyhow::bail!(
2549 "Cannot add actor while node is running, add actors before running the node"
2550 );
2551 }
2552
2553 self.kernel
2554 .trader
2555 .borrow_mut()
2556 .add_actor_from_factory(factory)
2557 }
2558
2559 pub fn add_strategy<T>(&mut self, mut strategy: T) -> anyhow::Result<()>
2575 where
2576 T: Strategy + StrategyNative + DataActorNative + Component + Debug + 'static,
2577 {
2578 if self.state() != NodeState::Idle {
2579 anyhow::bail!(
2580 "Cannot add strategy while node is running, add strategies before running the node"
2581 );
2582 }
2583
2584 let strategy_id = self
2586 .kernel
2587 .trader
2588 .borrow()
2589 .prepare_strategy_for_registration(&mut strategy)?;
2590 let oms_type = StrategyNative::strategy_core(&strategy).config.oms_type;
2591 let claims = strategy.external_order_claims().unwrap_or_default();
2592
2593 let mut exec_engine = if claims.is_empty() && oms_type.is_none() {
2596 None
2597 } else {
2598 Some(self.kernel.exec_engine.try_borrow_mut().map_err(|e| {
2599 anyhow::anyhow!("Cannot register external order claims or OMS type: {e}")
2600 })?)
2601 };
2602 let instrument_ids = match &exec_engine {
2603 Some(exec_engine) => Self::preflight_external_order_claims(
2604 &self.exec_manager,
2605 exec_engine,
2606 strategy_id,
2607 &claims,
2608 )?,
2609 None => HashSet::new(),
2610 };
2611
2612 self.kernel.trader.borrow_mut().add_strategy(strategy)?;
2613
2614 if let Some(exec_engine) = &mut exec_engine {
2617 exec_engine.commit_external_order_claims(strategy_id, &instrument_ids);
2618 self.exec_manager
2619 .register_external_order_claims(strategy_id, &instrument_ids);
2620
2621 if let Some(oms_type) = oms_type {
2622 exec_engine.register_oms_type(strategy_id, oms_type);
2623 }
2624 }
2625
2626 Ok(())
2627 }
2628
2629 pub fn register_external_order_claims(
2641 &mut self,
2642 strategy_id: StrategyId,
2643 claims: &[InstrumentId],
2644 ) -> anyhow::Result<()> {
2645 let mut exec_engine = self
2646 .kernel
2647 .exec_engine
2648 .try_borrow_mut()
2649 .map_err(|e| anyhow::anyhow!("Cannot register external order claims: {e}"))?;
2650 let instrument_ids = Self::preflight_external_order_claims(
2651 &self.exec_manager,
2652 &exec_engine,
2653 strategy_id,
2654 claims,
2655 )?;
2656
2657 exec_engine.commit_external_order_claims(strategy_id, &instrument_ids);
2658 self.exec_manager
2659 .register_external_order_claims(strategy_id, &instrument_ids);
2660
2661 Ok(())
2662 }
2663
2664 fn preflight_external_order_claims(
2665 exec_manager: &ExecutionManager,
2666 exec_engine: &ExecutionEngine,
2667 strategy_id: StrategyId,
2668 claims: &[InstrumentId],
2669 ) -> anyhow::Result<HashSet<InstrumentId>> {
2670 let mut instrument_ids = HashSet::new();
2671
2672 for instrument_id in claims {
2673 if !instrument_ids.insert(*instrument_id) {
2674 anyhow::bail!(
2675 "External order claim for {instrument_id} already exists for {strategy_id}"
2676 );
2677 }
2678 }
2679
2680 for instrument_id in &instrument_ids {
2681 if let Some(existing) = exec_manager.get_external_order_claim(instrument_id) {
2682 anyhow::bail!(
2683 "External order claim for {instrument_id} already exists for {existing}"
2684 );
2685 }
2686
2687 if let Some(existing) = exec_engine.get_external_order_claim(instrument_id) {
2688 anyhow::bail!(
2689 "External order claim for {instrument_id} already exists for {existing}"
2690 );
2691 }
2692 }
2693
2694 Ok(instrument_ids)
2695 }
2696
2697 pub fn deregister_external_order_claims(
2709 &mut self,
2710 strategy_id: StrategyId,
2711 ) -> anyhow::Result<()> {
2712 let mut exec_engine = self
2713 .kernel
2714 .exec_engine
2715 .try_borrow_mut()
2716 .map_err(|e| anyhow::anyhow!("Cannot deregister external order claims: {e}"))?;
2717 let manager_instruments = self
2718 .exec_manager
2719 .get_external_order_claims_for_strategy(strategy_id);
2720 let engine_instruments = exec_engine.get_external_order_claims_for_strategy(strategy_id);
2721
2722 if manager_instruments != engine_instruments {
2723 anyhow::bail!(
2724 "External order claims for {strategy_id} differ between the execution manager and engine"
2725 );
2726 }
2727
2728 exec_engine.deregister_external_order_claims(strategy_id);
2729 self.exec_manager
2730 .deregister_external_order_claims(strategy_id);
2731
2732 Ok(())
2733 }
2734
2735 pub fn add_exec_algorithm<T>(&mut self, exec_algorithm: T) -> anyhow::Result<()>
2746 where
2747 T: ExecutionAlgorithm + ExecutionAlgorithmNative + Component + Debug + 'static,
2748 {
2749 if self.state() != NodeState::Idle {
2750 anyhow::bail!(
2751 "Cannot add exec algorithm while node is running, add exec algorithms before running the node"
2752 );
2753 }
2754
2755 self.kernel
2756 .trader
2757 .borrow_mut()
2758 .add_exec_algorithm(exec_algorithm)
2759 }
2760
2761 fn run_reconciliation_checks(
2765 &mut self,
2766 now: dst::time::Instant,
2767 intervals: ReconciliationCheckIntervals,
2768 state: &mut ReconciliationCheckState<'_>,
2769 ) {
2770 if reconciliation_check_due(now, *state.last_inflight_check, intervals.inflight) {
2771 if self.state() == NodeState::ShuttingDown {
2772 return;
2773 }
2774 let result = self.exec_manager.check_inflight_orders();
2775 self.process_reconciliation_events(&result.events);
2776 for cmd in result.queries {
2777 AsyncRunner::handle_exec_command(cmd);
2778 }
2779 *state.last_inflight_check = now;
2780 }
2781
2782 let open_due = reconciliation_check_due(now, *state.last_open_check, intervals.open);
2783 let position_due =
2784 reconciliation_check_due(now, *state.last_position_check, intervals.position);
2785
2786 if (open_due || position_due) && self.state() == NodeState::ShuttingDown {
2787 return;
2788 }
2789
2790 if state.open_order_report_task.is_some() || state.targeted_order_report_task.is_some() {
2791 if open_due {
2792 log::debug!("Open-order reconciliation already in progress");
2793 *state.last_open_check = now;
2794 }
2795
2796 if position_due {
2797 log::debug!(
2798 "Position reconciliation delayed: open-order reconciliation in progress"
2799 );
2800 }
2801
2802 return;
2803 }
2804
2805 if state.position_report_task.is_some() {
2806 if position_due {
2807 log::debug!("Position reconciliation already in progress");
2808 *state.last_position_check = now;
2809 }
2810
2811 if open_due {
2812 log::debug!(
2813 "Open-order reconciliation delayed: position reconciliation in progress"
2814 );
2815 }
2816
2817 return;
2818 }
2819
2820 if position_due && (!open_due || *state.last_position_check < *state.last_open_check) {
2821 *state.position_report_task = self.start_position_report_check();
2822 *state.last_position_check = now;
2823 } else if open_due {
2824 *state.open_order_report_task = self.start_open_order_report_check();
2825 *state.last_open_check = now;
2826 }
2827 }
2828
2829 fn start_open_order_report_check(&mut self) -> Option<OpenOrderReportTask> {
2830 if self.exec_clients.is_empty() {
2831 log::debug!("No execution clients to check orders consistency");
2832 return None;
2833 }
2834
2835 let client_refs = self
2836 .exec_clients
2837 .iter()
2838 .map(|client| client as &dyn ExecutionClient)
2839 .collect::<Vec<_>>();
2840 let check = self
2841 .exec_manager
2842 .prepare_open_order_report_check(UUID4::new(), &client_refs);
2843 let command = check.command.clone();
2844 let clients = self.exec_clients.clone();
2845 let deadline = dst::time::Instant::now() + self.config.timeout_reconciliation;
2846
2847 Some(OpenOrderReportTask {
2848 future: Box::pin(async move {
2849 let remaining = deadline.saturating_duration_since(dst::time::Instant::now());
2850 match dst::time::timeout(remaining, request_open_order_reports(clients, command))
2851 .await
2852 {
2853 Ok(result) => ReportTaskOutcome::Completed(OpenOrderReportResult {
2854 check,
2855 reports: result.reports,
2856 queried_clients: result.queried_clients,
2857 failed_clients: result.failed_clients,
2858 }),
2859 Err(_) => ReportTaskOutcome::TimedOut,
2860 }
2861 }),
2862 })
2863 }
2864
2865 fn start_targeted_order_report_check(
2866 &self,
2867 queries: Vec<TargetedOrderQuery>,
2868 ) -> TargetedOrderReportTask {
2869 let clients = self.exec_clients.clone();
2870 let query_delay = Duration::from_millis(u64::from(
2871 self.config.exec_engine.single_order_query_delay_ms,
2872 ));
2873 let planned_client_order_ids = queries
2874 .iter()
2875 .map(TargetedOrderQuery::client_order_id)
2876 .collect();
2877 let deadline = dst::time::Instant::now() + self.config.timeout_reconciliation;
2878
2879 TargetedOrderReportTask {
2880 future: Box::pin(async move {
2881 let client_refs = clients
2882 .iter()
2883 .map(|client| client as &dyn ExecutionClient)
2884 .collect::<Vec<_>>();
2885 let remaining = deadline.saturating_duration_since(dst::time::Instant::now());
2886 match dst::time::timeout(
2887 remaining,
2888 request_targeted_order_reports(&client_refs, queries, query_delay),
2889 )
2890 .await
2891 {
2892 Ok(result) => ReportTaskOutcome::Completed(result),
2893 Err(_) => ReportTaskOutcome::TimedOut,
2894 }
2895 }),
2896 planned_client_order_ids,
2897 }
2898 }
2899
2900 fn start_position_report_check(&self) -> Option<PositionReportTask> {
2901 if self.exec_clients.is_empty() {
2902 log::debug!("No execution clients to check positions consistency");
2903 return None;
2904 }
2905
2906 let client_refs = self
2907 .exec_clients
2908 .iter()
2909 .map(|client| client as &dyn ExecutionClient)
2910 .collect::<Vec<_>>();
2911 let check = self
2912 .exec_manager
2913 .prepare_position_report_check(UUID4::new(), &client_refs);
2914 let command = check.command.clone();
2915 let clients = self.exec_clients.clone();
2916 let deadline = dst::time::Instant::now() + self.config.timeout_reconciliation;
2917
2918 Some(PositionReportTask {
2919 future: Box::pin(async move {
2920 let remaining = deadline.saturating_duration_since(dst::time::Instant::now());
2921 match dst::time::timeout(remaining, request_position_reports(clients, command))
2922 .await
2923 {
2924 Ok(result) => ReportTaskOutcome::Completed(PositionReportResult {
2925 check,
2926 reports: result.reports,
2927 queried_clients: result.queried_clients,
2928 failed_clients: result.failed_clients,
2929 }),
2930 Err(_) => ReportTaskOutcome::TimedOut,
2931 }
2932 }),
2933 })
2934 }
2935
2936 fn flush_pending_exec_client_instruments(&self) {
2937 for client in &self.exec_clients {
2938 client.flush_pending_instruments();
2939 }
2940 }
2941
2942 fn cleanup_cancelled_report_tasks(&mut self, planned_client_order_ids: &[ClientOrderId]) {
2943 self.flush_pending_exec_client_instruments();
2944 self.exec_manager
2945 .remove_targeted_order_queries(planned_client_order_ids);
2946 }
2947
2948 fn cancel_report_tasks(
2949 &mut self,
2950 open_order_report_task: &mut Option<OpenOrderReportTask>,
2951 targeted_order_report_task: &mut Option<TargetedOrderReportTask>,
2952 position_report_task: &mut Option<PositionReportTask>,
2953 ) {
2954 let planned_client_order_ids = targeted_order_report_task
2955 .as_ref()
2956 .map(|task| task.planned_client_order_ids.clone())
2957 .unwrap_or_default();
2958
2959 drop(open_order_report_task.take());
2960 drop(targeted_order_report_task.take());
2961 drop(position_report_task.take());
2962 self.cleanup_cancelled_report_tasks(&planned_client_order_ids);
2963 }
2964}
2965
2966fn record_runner_dispatch(
2967 metrics: &RunnerMetrics,
2968 channel: SystemChannel,
2969 dispatch_start: dst::time::Instant,
2970 metrics_start: dst::time::Instant,
2971) {
2972 let dispatch_end = dst::time::Instant::now();
2973 metrics.record_dispatch(
2974 channel,
2975 dispatch_end.duration_since(dispatch_start),
2976 dispatch_end.duration_since(metrics_start),
2977 );
2978}
2979
2980fn record_runner_maintenance(
2981 metrics: &RunnerMetrics,
2982 work_start: dst::time::Instant,
2983 metrics_start: dst::time::Instant,
2984) {
2985 let work_end = dst::time::Instant::now();
2986 metrics.record_maintenance(
2987 work_end.duration_since(work_start),
2988 work_end.duration_since(metrics_start),
2989 );
2990}
2991
2992fn record_runner_external_msgbus(
2993 metrics: &RunnerMetrics,
2994 work_start: dst::time::Instant,
2995 metrics_start: dst::time::Instant,
2996) {
2997 let work_end = dst::time::Instant::now();
2998 metrics.record_external_msgbus(
2999 work_end.duration_since(work_start),
3000 work_end.duration_since(metrics_start),
3001 );
3002}
3003
3004async fn recv_external_msgbus_message(
3005 rx: &mut Option<tokio::sync::mpsc::Receiver<BusMessage>>,
3006) -> Option<BusMessage> {
3007 match rx {
3008 Some(rx) => rx.recv().await,
3009 None => std::future::pending::<Option<BusMessage>>().await,
3010 }
3011}
3012
3013async fn request_open_order_reports(
3014 clients: Vec<LiveExecutionClient>,
3015 command: GenerateOrderStatusReports,
3016) -> OpenOrderReportQueryResult {
3017 let mut all_reports = Vec::new();
3018 let mut queried_clients = IndexSet::new();
3019 let mut failed_clients = IndexSet::new();
3020
3021 for client in clients {
3022 let client_id = client.client_id();
3023 queried_clients.insert(client_id);
3024
3025 match client.generate_order_status_reports(&command).await {
3026 Ok(reports) => {
3027 all_reports.extend(
3028 reports
3029 .into_iter()
3030 .map(|report| SourcedOrderStatusReport { client_id, report }),
3031 );
3032 }
3033 Err(e) => {
3034 failed_clients.insert(client_id);
3035 log::warn!(
3036 "Failed to generate order status reports from {}: {e}",
3037 client.client_id()
3038 );
3039 }
3040 }
3041 }
3042
3043 OpenOrderReportQueryResult {
3044 reports: all_reports,
3045 queried_clients,
3046 failed_clients,
3047 }
3048}
3049
3050async fn request_position_reports(
3051 clients: Vec<LiveExecutionClient>,
3052 command: GeneratePositionStatusReports,
3053) -> PositionReportQueryResult {
3054 let mut all_reports = Vec::new();
3055 let mut queried_clients = IndexSet::new();
3056 let mut failed_clients = IndexSet::new();
3057
3058 for client in clients {
3059 let client_id = client.client_id();
3060 queried_clients.insert(client_id);
3061
3062 match client.generate_position_status_reports(&command).await {
3063 Ok(reports) => {
3064 all_reports.extend(reports);
3065 }
3066 Err(e) => {
3067 failed_clients.insert(client_id);
3068 log::warn!(
3069 "Failed to generate position status reports from {}: {e}",
3070 client.client_id()
3071 );
3072 }
3073 }
3074 }
3075
3076 PositionReportQueryResult {
3077 reports: all_reports,
3078 queried_clients,
3079 failed_clients,
3080 }
3081}
3082
3083fn reconciliation_check_due(
3084 now: dst::time::Instant,
3085 last: dst::time::Instant,
3086 interval: Duration,
3087) -> bool {
3088 interval > Duration::ZERO
3089 && now
3090 .checked_duration_since(last)
3091 .is_some_and(|elapsed| elapsed >= interval)
3092}
3093
3094#[derive(Clone, Copy)]
3095struct ReconciliationCheckIntervals {
3096 inflight: Duration,
3097 open: Duration,
3098 position: Duration,
3099}
3100
3101struct ReconciliationCheckState<'a> {
3102 last_inflight_check: &'a mut dst::time::Instant,
3103 last_open_check: &'a mut dst::time::Instant,
3104 last_position_check: &'a mut dst::time::Instant,
3105 open_order_report_task: &'a mut Option<OpenOrderReportTask>,
3106 targeted_order_report_task: &'a mut Option<TargetedOrderReportTask>,
3107 position_report_task: &'a mut Option<PositionReportTask>,
3108}
3109
3110enum ReportTaskOutcome<T> {
3111 Completed(T),
3112 TimedOut,
3113}
3114
3115type OpenOrderReportFuture =
3116 Pin<Box<dyn Future<Output = ReportTaskOutcome<OpenOrderReportResult>>>>;
3117
3118struct OpenOrderReportTask {
3119 future: OpenOrderReportFuture,
3120}
3121
3122struct OpenOrderReportResult {
3123 check: OpenOrderReportCheck,
3124 reports: Vec<SourcedOrderStatusReport>,
3125 queried_clients: IndexSet<ClientId>,
3126 failed_clients: IndexSet<ClientId>,
3127}
3128
3129type TargetedOrderReportFuture =
3130 Pin<Box<dyn Future<Output = ReportTaskOutcome<Vec<TargetedOrderReportResult>>>>>;
3131
3132struct TargetedOrderReportTask {
3133 future: TargetedOrderReportFuture,
3134 planned_client_order_ids: Vec<ClientOrderId>,
3135}
3136
3137struct OpenOrderReportQueryResult {
3138 reports: Vec<SourcedOrderStatusReport>,
3139 queried_clients: IndexSet<ClientId>,
3140 failed_clients: IndexSet<ClientId>,
3141}
3142
3143type PositionReportFuture = Pin<Box<dyn Future<Output = ReportTaskOutcome<PositionReportResult>>>>;
3144
3145struct PositionReportTask {
3146 future: PositionReportFuture,
3147}
3148
3149struct PositionReportResult {
3150 check: PositionReportCheck,
3151 reports: Vec<PositionStatusReport>,
3152 queried_clients: IndexSet<ClientId>,
3153 failed_clients: IndexSet<ClientId>,
3154}
3155
3156struct PositionReportQueryResult {
3157 reports: Vec<PositionStatusReport>,
3158 queried_clients: IndexSet<ClientId>,
3159 failed_clients: IndexSet<ClientId>,
3160}
3161
3162struct RunnerReceivers<'a> {
3163 time_evt: &'a mut tokio::sync::mpsc::UnboundedReceiver<TimeEventMessage>,
3164 system_evt: &'a mut tokio::sync::mpsc::UnboundedReceiver<SystemEvent>,
3165 system_cmd: &'a mut tokio::sync::mpsc::UnboundedReceiver<SystemCommand>,
3166 exec_evt: &'a mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3167 exec_cmd: &'a mut tokio::sync::mpsc::UnboundedReceiver<TradingCommandMessage>,
3168 data_evt: &'a mut tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
3169 data_cmd: &'a mut tokio::sync::mpsc::UnboundedReceiver<DataCommand>,
3170}
3171
3172fn flush_pending_data(
3179 pending: &mut PendingEvents,
3180 data_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
3181 data_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataCommand>,
3182) {
3183 loop {
3184 let mut progressed = pending.drain_data();
3185
3186 while let Ok(evt) = data_evt_rx.try_recv() {
3187 AsyncRunner::handle_data_event(evt);
3188 progressed = true;
3189 }
3190
3191 while let Ok(cmd) = data_cmd_rx.try_recv() {
3192 AsyncRunner::handle_data_command(cmd);
3193 progressed = true;
3194 }
3195
3196 if !progressed {
3197 break;
3198 }
3199 }
3200}
3201
3202#[expect(
3208 clippy::too_many_arguments,
3209 reason = "all runner receivers are drained together"
3210)]
3211fn flush_all_pending(
3212 pending: &mut PendingEvents,
3213 time_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<TimeEventMessage>,
3214 system_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<SystemEvent>,
3215 system_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<SystemCommand>,
3216 exec_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3217 exec_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<TradingCommandMessage>,
3218 data_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
3219 data_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataCommand>,
3220) {
3221 while let Ok(handler) = time_evt_rx.try_recv() {
3223 let _ = AsyncRunner::handle_time_event(handler);
3224 }
3225
3226 while let Ok(event) = system_evt_rx.try_recv() {
3227 pending.system_events.push(event);
3228 }
3229
3230 while let Ok(command) = system_cmd_rx.try_recv() {
3231 pending.system_commands.push(command);
3232 }
3233
3234 while let Ok(evt) = data_evt_rx.try_recv() {
3235 pending.data_evts.push(evt);
3236 }
3237
3238 while let Ok(cmd) = data_cmd_rx.try_recv() {
3239 pending.data_cmds.push(cmd);
3240 }
3241
3242 while let Ok(evt) = exec_evt_rx.try_recv() {
3243 match evt {
3244 ExecutionEvent::Account(_) => {
3245 AsyncRunner::handle_exec_event(evt);
3246 }
3247 ExecutionEvent::Report(report) => {
3248 pending.exec_reports.push(report);
3249 }
3250 ExecutionEvent::Order(order_evt) => {
3251 pending.order_evts.push(order_evt);
3252 }
3253 ExecutionEvent::OrderSubmittedBatch(batch) => {
3254 for submitted in batch {
3255 pending.order_evts.push(OrderEventAny::Submitted(submitted));
3256 }
3257 }
3258 ExecutionEvent::OrderAcceptedBatch(batch) => {
3259 for accepted in batch {
3260 pending.order_evts.push(OrderEventAny::Accepted(accepted));
3261 }
3262 }
3263 ExecutionEvent::OrderCanceledBatch(batch) => {
3264 for canceled in batch {
3265 pending.order_evts.push(OrderEventAny::Canceled(canceled));
3266 }
3267 }
3268 }
3269 }
3270
3271 while let Ok(cmd) = exec_cmd_rx.try_recv() {
3272 pending.exec_cmds.push(cmd);
3273 }
3274
3275 pending.drain();
3276}
3277
3278#[expect(
3283 clippy::too_many_arguments,
3284 reason = "startup buffering owns one future plus the pending state and all runner receivers"
3285)]
3286async fn drive_with_event_buffering<F: std::future::Future>(
3287 future: F,
3288 pending: &mut PendingEvents,
3289 time_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<TimeEventMessage>,
3290 system_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<SystemEvent>,
3291 system_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<SystemCommand>,
3292 exec_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3293 exec_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<TradingCommandMessage>,
3294 data_evt_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
3295 data_cmd_rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataCommand>,
3296) -> F::Output {
3297 tokio::pin!(future);
3298
3299 loop {
3300 tokio::select! {
3301 biased;
3302
3303 result = &mut future => {
3304 break result;
3305 }
3306 Some(handler) = time_evt_rx.recv() => {
3307 let _ = AsyncRunner::handle_time_event(handler);
3308 }
3309 Some(event) = system_evt_rx.recv() => {
3310 pending.system_events.push(event);
3311 }
3312 Some(command) = system_cmd_rx.recv() => {
3313 pending.system_commands.push(command);
3314 }
3315 Some(evt) = exec_evt_rx.recv() => {
3316 match evt {
3320 ExecutionEvent::Account(_) => {
3321 AsyncRunner::handle_exec_event(evt);
3322 }
3323 ExecutionEvent::Report(report) => {
3324 pending.exec_reports.push(report);
3325 }
3326 ExecutionEvent::Order(order_evt) => {
3327 pending.order_evts.push(order_evt);
3328 }
3329 ExecutionEvent::OrderSubmittedBatch(batch) => {
3330 for submitted in batch {
3331 pending.order_evts.push(OrderEventAny::Submitted(submitted));
3332 }
3333 }
3334 ExecutionEvent::OrderAcceptedBatch(batch) => {
3335 for accepted in batch {
3336 pending.order_evts.push(OrderEventAny::Accepted(accepted));
3337 }
3338 }
3339 ExecutionEvent::OrderCanceledBatch(batch) => {
3340 for canceled in batch {
3341 pending.order_evts.push(OrderEventAny::Canceled(canceled));
3342 }
3343 }
3344 }
3345 }
3346 Some(cmd) = exec_cmd_rx.recv() => {
3347 pending.exec_cmds.push(cmd);
3348 }
3349 Some(evt) = data_evt_rx.recv() => {
3350 pending.data_evts.push(evt);
3351 }
3352 Some(cmd) = data_cmd_rx.recv() => {
3353 pending.data_cmds.push(cmd);
3354 }
3355 }
3356 }
3357}
3358
3359#[derive(Default)]
3360struct PendingEvents {
3361 system_events: Vec<SystemEvent>,
3362 system_commands: Vec<SystemCommand>,
3363 data_evts: Vec<DataEvent>,
3364 data_cmds: Vec<DataCommand>,
3365 exec_reports: Vec<ExecutionReport>,
3366 order_evts: Vec<OrderEventAny>,
3367 exec_cmds: Vec<TradingCommandMessage>,
3368}
3369
3370impl PendingEvents {
3371 fn is_empty(&self) -> bool {
3372 self.system_events.is_empty()
3373 && self.system_commands.is_empty()
3374 && self.data_evts.is_empty()
3375 && self.data_cmds.is_empty()
3376 && self.exec_reports.is_empty()
3377 && self.order_evts.is_empty()
3378 && self.exec_cmds.is_empty()
3379 }
3380
3381 fn drain_data(&mut self) -> bool {
3385 let total = self.data_evts.len() + self.data_cmds.len();
3386
3387 if total > 0 {
3388 log::debug!(
3389 "Draining {total} data events/commands into cache \
3390 (data_evts={}, data_cmds={})",
3391 self.data_evts.len(),
3392 self.data_cmds.len(),
3393 );
3394 }
3395
3396 for evt in self.data_evts.drain(..) {
3397 AsyncRunner::handle_data_event(evt);
3398 }
3399
3400 for cmd in self.data_cmds.drain(..) {
3401 AsyncRunner::handle_data_command(cmd);
3402 }
3403
3404 total > 0
3405 }
3406
3407 fn drain(&mut self) {
3409 let total = self.data_evts.len()
3410 + self.data_cmds.len()
3411 + self.exec_reports.len()
3412 + self.order_evts.len()
3413 + self.exec_cmds.len();
3414
3415 if total > 0 {
3416 log::debug!(
3417 "Processing {total} events/commands queued during startup \
3418 (data_evts={}, data_cmds={}, exec_reports={}, order_evts={}, exec_cmds={})",
3419 self.data_evts.len(),
3420 self.data_cmds.len(),
3421 self.exec_reports.len(),
3422 self.order_evts.len(),
3423 self.exec_cmds.len()
3424 );
3425 }
3426
3427 for evt in self.data_evts.drain(..) {
3428 AsyncRunner::handle_data_event(evt);
3429 }
3430
3431 for cmd in self.data_cmds.drain(..) {
3432 AsyncRunner::handle_data_command(cmd);
3433 }
3434
3435 for report in self.exec_reports.drain(..) {
3436 AsyncRunner::handle_exec_event(ExecutionEvent::Report(report));
3437 }
3438
3439 for evt in self.order_evts.drain(..) {
3440 AsyncRunner::handle_exec_event(ExecutionEvent::Order(evt));
3441 }
3442
3443 for cmd in self.exec_cmds.drain(..) {
3444 AsyncRunner::handle_trading_command(cmd);
3445 }
3446 }
3447
3448 fn take_system_events(&mut self) -> Vec<SystemEvent> {
3449 std::mem::take(&mut self.system_events)
3450 }
3451
3452 fn take_system_commands(&mut self) -> Vec<SystemCommand> {
3453 std::mem::take(&mut self.system_commands)
3454 }
3455}
3456
3457struct ClientStatus {
3458 client: String,
3459 client_type: &'static str,
3460 connected: bool,
3461}
3462
3463fn render_client_statuses(rows: Vec<ClientStatus>) -> String {
3464 let mut builder = Builder::with_capacity(rows.len() + 1, 3);
3465 builder.push_record(["Client", "Type", "Connected"]);
3466
3467 for row in rows {
3468 builder.push_record([
3469 row.client,
3470 row.client_type.to_string(),
3471 row.connected.to_string(),
3472 ]);
3473 }
3474
3475 builder.build().with(Style::rounded()).to_string()
3476}
3477
3478#[cfg(test)]
3479mod tests {
3480 use std::{
3481 cell::{Cell, RefCell},
3482 fmt::Debug,
3483 rc::Rc,
3484 sync::{
3485 Arc, Mutex,
3486 atomic::{AtomicBool, Ordering},
3487 },
3488 };
3489
3490 use bytes::Bytes;
3491 use indexmap::IndexMap;
3492 use log::{Level, LevelFilter, Log, Metadata, Record};
3493 #[cfg(feature = "python")]
3494 use nautilus_common::runner::{
3495 SyncDataCommandSender, SyncTradingCommandSender, replace_data_cmd_sender,
3496 replace_exec_cmd_sender,
3497 };
3498 use nautilus_common::{
3499 actor::{DataActor, DataActorCore, data_actor::DataActorConfig},
3500 cache::Cache,
3501 clock::{Clock, TestClock},
3502 enums::SerializationEncoding,
3503 live::runner::{get_data_event_sender, get_exec_event_sender, get_system_event_sender},
3504 messages::{
3505 execution::{QueryAccount, SubmitOrder, TradingCommand},
3506 system::{
3507 QueueCondition, QueueState, ReconnectSocket, SocketState, SocketStateChanged,
3508 },
3509 },
3510 msgbus::{
3511 self, BusMessage, BusPayloadType, MessageBusBacking, MessageBusBackingFactory,
3512 MessageBusConfig, MessageBusExternalEgress, MessageBusExternalIngress,
3513 MessagingSwitchboard, ShareableMessageHandler, TypedHandler, TypedIntoHandler,
3514 },
3515 nautilus_actor,
3516 testing::wait_until_async,
3517 };
3518 use nautilus_core::{UUID4, UnixNanos};
3519 use nautilus_execution::{
3520 engine::{ExecutionEngine, SnapshotAnchorer, stubs::StubExecutionClient},
3521 reconciliation::create_inferred_fill_for_qty,
3522 };
3523 use nautilus_model::{
3524 accounts::{AccountAny, MarginAccount},
3525 data::QuoteTick,
3526 enums::{
3527 AccountType, LiquiditySide, OmsType, OrderSide, OrderStatus, OrderType, TimeInForce,
3528 },
3529 events::{
3530 AccountState, OrderAcceptedBatch, OrderFilled,
3531 order::spec::{OrderAcceptedSpec, OrderPendingUpdateSpec, OrderUpdatedSpec},
3532 },
3533 identifiers::{
3534 AccountId, ActorId, ClientId, InstrumentId, PositionId, StrategyId, TradeId, TraderId,
3535 Venue, VenueOrderId,
3536 },
3537 instruments::{Instrument, InstrumentAny, stubs::crypto_perpetual_ethusdt},
3538 orders::{OrderTestBuilder, stubs::TestOrderEventStubs},
3539 reports::FillReport,
3540 types::{AccountBalance, Currency, Money, Price, Quantity},
3541 };
3542 use nautilus_system::{KernelEventStore, RegisteredComponents, event_store::EventStoreConfig};
3543 use nautilus_testkit::{
3544 cache::TestCacheDatabaseControl,
3545 components::{StateActor, StateStrategy},
3546 };
3547 use nautilus_trading::{
3548 nautilus_strategy,
3549 strategy::{config::StrategyConfig, core::StrategyCore},
3550 };
3551 use rstest::*;
3552 use rust_decimal_macros::dec;
3553 use ustr::Ustr;
3554
3555 use super::*;
3556
3557 struct ExternalIngressLogCapture {
3558 messages: Mutex<Vec<String>>,
3559 }
3560
3561 static EXTERNAL_INGRESS_LOG_CAPTURE: ExternalIngressLogCapture = ExternalIngressLogCapture {
3562 messages: Mutex::new(Vec::new()),
3563 };
3564
3565 #[derive(Debug)]
3566 struct StartupSocketActor {
3567 core: DataActorCore,
3568 received: Rc<RefCell<Vec<SocketStateChanged>>>,
3569 }
3570
3571 impl StartupSocketActor {
3572 fn new(received: Rc<RefCell<Vec<SocketStateChanged>>>) -> Self {
3573 Self {
3574 core: DataActorCore::new(DataActorConfig {
3575 actor_id: Some(ActorId::from("SOCKET-STARTUP-ACTOR")),
3576 ..Default::default()
3577 }),
3578 received,
3579 }
3580 }
3581 }
3582
3583 impl DataActor for StartupSocketActor {
3584 fn on_start(&mut self) -> anyhow::Result<()> {
3585 self.subscribe_socket_state(None);
3586 Ok(())
3587 }
3588
3589 fn on_socket_state(&mut self, event: &SocketStateChanged) -> anyhow::Result<()> {
3590 self.received.borrow_mut().push(event.clone());
3591 Ok(())
3592 }
3593 }
3594
3595 nautilus_actor!(StartupSocketActor);
3596
3597 impl Log for ExternalIngressLogCapture {
3598 fn enabled(&self, metadata: &Metadata<'_>) -> bool {
3599 metadata.level() == Level::Error && metadata.target() == "nautilus_live::node"
3600 }
3601
3602 fn log(&self, record: &Record<'_>) {
3603 if self.enabled(record.metadata()) {
3604 self.messages
3605 .lock()
3606 .unwrap()
3607 .push(record.args().to_string());
3608 }
3609 }
3610
3611 fn flush(&self) {}
3612 }
3613
3614 #[rstest]
3615 fn test_render_client_statuses() {
3616 let rows = vec![
3617 ClientStatus {
3618 client: "BINANCE".to_string(),
3619 client_type: "Data",
3620 connected: true,
3621 },
3622 ClientStatus {
3623 client: "SIM".to_string(),
3624 client_type: "Execution",
3625 connected: false,
3626 },
3627 ];
3628
3629 let output = render_client_statuses(rows);
3630 let expected = "â•─────────┬───────────┬───────────╮\n\
3631│ Client │ Type │ Connected │\n\
3632├─────────┼───────────┼───────────┤\n\
3633│ BINANCE │ Data │ true │\n\
3634│ SIM │ Execution │ false │\n\
3635╰─────────┴───────────┴───────────╯";
3636
3637 assert_eq!(output, expected);
3638 }
3639
3640 #[rstest]
3641 fn test_republish_external_msgbus_message_logs_topic_and_error_chain() {
3642 log::set_logger(&EXTERNAL_INGRESS_LOG_CAPTURE).expect("test logger already installed");
3643 log::set_max_level(LevelFilter::Error);
3644 EXTERNAL_INGRESS_LOG_CAPTURE
3645 .messages
3646 .lock()
3647 .unwrap()
3648 .clear();
3649 let message = BusMessage::with_str_topic(
3650 "data.quotes.AUDUSD.SIM*",
3651 BusPayloadType::Custom(Ustr::from("UnregisteredCustomData")),
3652 Bytes::new(),
3653 SerializationEncoding::Json,
3654 );
3655
3656 LiveNode::republish_external_msgbus_message(&message);
3657
3658 assert_eq!(
3659 *EXTERNAL_INGRESS_LOG_CAPTURE.messages.lock().unwrap(),
3660 vec![
3661 "Failed to republish external message bus topic 'data.quotes.AUDUSD.SIM*': invalid \
3662 external message topic: Topic `value` contained invalid characters, was \
3663 data.quotes.AUDUSD.SIM*"
3664 .to_string()
3665 ],
3666 );
3667 }
3668
3669 #[rstest]
3670 fn test_publish_queue_state_transitions_reaches_typed_subscriber() {
3671 let config = LiveNodeConfig {
3672 trader_id: TraderId::from("QUEUE-001"),
3673 exec_engine: crate::config::LiveExecEngineConfig {
3674 reconciliation: false,
3675 ..Default::default()
3676 },
3677 ..Default::default()
3678 };
3679 let node = LiveNode::build("QueuePublicationNode".to_string(), Some(config)).unwrap();
3680 let received = Rc::new(RefCell::new(Vec::<QueueStateChanged>::new()));
3681
3682 let handler = ShareableMessageHandler::from_typed({
3683 let received = received.clone();
3684 move |event: &QueueStateChanged| received.borrow_mut().push(event.clone())
3685 });
3686
3687 msgbus::subscribe_any(
3688 MessagingSwitchboard::queue_state_changed_topic().into(),
3689 handler,
3690 None,
3691 );
3692 let transitions = [
3693 QueueStateTransition {
3694 channel: SystemChannel::DataEvents,
3695 condition: QueueCondition::Backlogged,
3696 state: QueueState::Triggered,
3697 queue_depth: 17,
3698 mean_dispatch_ns: 23,
3699 },
3700 QueueStateTransition {
3701 channel: SystemChannel::DataEvents,
3702 condition: QueueCondition::Slow,
3703 state: QueueState::Triggered,
3704 queue_depth: 17,
3705 mean_dispatch_ns: 23,
3706 },
3707 ];
3708
3709 node.publish_queue_state_transitions(&transitions);
3710
3711 let events = received.borrow();
3712 assert_eq!(events.len(), 2);
3713
3714 for (event, transition) in events.iter().zip(transitions) {
3715 assert_eq!(event.trader_id, TraderId::from("QUEUE-001"));
3716 assert_eq!(event.channel, transition.channel);
3717 assert_eq!(event.condition, transition.condition);
3718 assert_eq!(event.state, transition.state);
3719 assert_eq!(event.queue_depth, transition.queue_depth);
3720 assert_eq!(event.mean_dispatch_ns, transition.mean_dispatch_ns);
3721 assert_ne!(event.event_id, UUID4::default());
3722 assert_ne!(event.ts_event, UnixNanos::default());
3723 assert_eq!(event.ts_init, event.ts_event);
3724 }
3725 assert_ne!(events[0].event_id, events[1].event_id);
3726 drop(events);
3727 msgbus::get_message_bus().borrow_mut().dispose();
3728 }
3729
3730 #[rstest]
3731 fn test_process_socket_state_change_reaches_typed_subscriber() {
3732 let config = LiveNodeConfig {
3733 trader_id: TraderId::from("SOCKET-001"),
3734 exec_engine: crate::config::LiveExecEngineConfig {
3735 reconciliation: false,
3736 ..Default::default()
3737 },
3738 ..Default::default()
3739 };
3740 let node = LiveNode::build("SocketPublicationNode".to_string(), Some(config)).unwrap();
3741 let received = Rc::new(RefCell::new(Vec::<SocketStateChanged>::new()));
3742 let handler = ShareableMessageHandler::from_typed({
3743 let received = received.clone();
3744 move |event: &SocketStateChanged| received.borrow_mut().push(event.clone())
3745 });
3746 msgbus::subscribe_any(
3747 MessagingSwitchboard::socket_state_changed_topic().into(),
3748 handler,
3749 None,
3750 );
3751 let change = SocketStateChange::new(
3752 ClientId::from("BINANCE"),
3753 Some(Venue::from("BINANCE")),
3754 Ustr::from("binance-futures-market-streams"),
3755 SocketState::Disconnected,
3756 );
3757
3758 node.process_system_event(SystemEvent::SocketState(change));
3759
3760 let events = received.borrow();
3761 assert_eq!(events.len(), 1);
3762 assert_eq!(events[0].trader_id, TraderId::from("SOCKET-001"));
3763 assert_eq!(events[0].client_id, ClientId::from("BINANCE"));
3764 assert_eq!(events[0].venue, Some(Venue::from("BINANCE")));
3765 assert_eq!(
3766 events[0].endpoint,
3767 Ustr::from("binance-futures-market-streams")
3768 );
3769 assert_eq!(events[0].state, SocketState::Disconnected);
3770 assert_ne!(events[0].event_id, UUID4::default());
3771 assert_ne!(events[0].ts_event, UnixNanos::default());
3772 assert_eq!(events[0].ts_init, events[0].ts_event);
3773 drop(events);
3774 msgbus::get_message_bus().borrow_mut().dispose();
3775 }
3776
3777 #[rstest]
3778 #[tokio::test]
3779 async fn test_start_publishes_socket_change_after_actor_subscribes() {
3780 let config = LiveNodeConfig {
3781 trader_id: TraderId::from("SOCKET-STARTUP-001"),
3782 exec_engine: crate::config::LiveExecEngineConfig {
3783 reconciliation: false,
3784 ..Default::default()
3785 },
3786 timeout_connection: Duration::ZERO,
3787 timeout_reconciliation: Duration::ZERO,
3788 timeout_portfolio: Duration::ZERO,
3789 timeout_disconnection: Duration::ZERO,
3790 delay_post_stop: Duration::ZERO,
3791 timeout_shutdown: Duration::ZERO,
3792 ..Default::default()
3793 };
3794 let mut node = LiveNode::build("SocketStartupNode".to_string(), Some(config)).unwrap();
3795 let received = Rc::new(RefCell::new(Vec::new()));
3796 node.add_actor(StartupSocketActor::new(Rc::clone(&received)))
3797 .unwrap();
3798 node.runner.as_ref().unwrap().bind_senders();
3799 let change = SocketStateChange::new(
3800 ClientId::from("BINANCE"),
3801 Some(Venue::from("BINANCE")),
3802 Ustr::from("binance-futures-market-streams"),
3803 SocketState::Connected,
3804 );
3805 get_system_event_sender()
3806 .send(SystemEvent::SocketState(change))
3807 .unwrap();
3808
3809 node.start().await.unwrap();
3810
3811 {
3812 let events = received.borrow();
3813 assert_eq!(events.len(), 1);
3814 assert_eq!(events[0].trader_id, TraderId::from("SOCKET-STARTUP-001"));
3815 assert_eq!(events[0].client_id, change.client_id);
3816 assert_eq!(events[0].venue, change.venue);
3817 assert_eq!(events[0].endpoint, change.endpoint);
3818 assert_eq!(events[0].state, change.state);
3819 assert_ne!(events[0].event_id, UUID4::default());
3820 assert_ne!(events[0].ts_event, UnixNanos::default());
3821 assert_eq!(events[0].ts_init, events[0].ts_event);
3822 }
3823
3824 node.stop().await.unwrap();
3825 node.dispose();
3826 }
3827
3828 #[rstest]
3829 #[tokio::test(flavor = "current_thread")]
3830 async fn test_run_publishes_queue_state_after_dispatch_sample() {
3831 let config = LiveNodeConfig {
3832 trader_id: TraderId::from("QUEUE-RUN-001"),
3833 queue_monitor: Some(crate::config::QueueMonitorConfig {
3834 queue_depth_trigger: usize::MAX,
3835 queue_depth_clear: 0,
3836 mean_dispatch_ns_trigger: 1,
3837 mean_dispatch_ns_clear: 0,
3838 }),
3839 exec_engine: crate::config::LiveExecEngineConfig {
3840 reconciliation: false,
3841 ..Default::default()
3842 },
3843 delay_post_stop: Duration::ZERO,
3844 ..Default::default()
3845 };
3846 let mut node = LiveNode::build("QueueMonitorRunNode".to_string(), Some(config)).unwrap();
3847 let handle = node.handle();
3848 let received = Rc::new(RefCell::new(Vec::<QueueStateChanged>::new()));
3849
3850 let handler = ShareableMessageHandler::from_typed({
3851 let received = received.clone();
3852 let stop_handle = handle.clone();
3853
3854 move |event: &QueueStateChanged| {
3855 received.borrow_mut().push(event.clone());
3856 stop_handle.stop();
3857 }
3858 });
3859
3860 msgbus::subscribe_any(
3861 MessagingSwitchboard::queue_state_changed_topic().into(),
3862 handler,
3863 None,
3864 );
3865 let drive_handle = handle.clone();
3866
3867 let result = tokio::time::timeout(Duration::from_secs(5), async {
3868 let run = node.run();
3869 tokio::pin!(run);
3870
3871 let drive = async move {
3872 wait_until_async(
3873 || async { drive_handle.is_running() },
3874 Duration::from_secs(2),
3875 )
3876 .await;
3877 get_data_event_sender().send(stub_data_event()).unwrap();
3878 };
3879
3880 let (run_result, ()) = tokio::join!(run, drive);
3881 run_result
3882 })
3883 .await;
3884
3885 assert!(
3886 result.is_ok(),
3887 "queue state event should arrive before timeout"
3888 );
3889 assert!(result.unwrap().is_ok(), "run() should succeed");
3890 let events = received.borrow();
3891 assert_eq!(events.len(), 1);
3892 assert_eq!(events[0].trader_id, TraderId::from("QUEUE-RUN-001"));
3893 assert_eq!(events[0].channel, SystemChannel::DataEvents);
3894 assert_eq!(events[0].condition, QueueCondition::Slow);
3895 assert_eq!(events[0].state, QueueState::Triggered);
3896 assert_eq!(events[0].queue_depth, 0);
3897 assert!(events[0].mean_dispatch_ns > 0);
3898 assert_ne!(events[0].event_id, UUID4::default());
3899 assert_ne!(events[0].ts_event, UnixNanos::default());
3900 assert_eq!(events[0].ts_init, events[0].ts_event);
3901 drop(events);
3902 msgbus::get_message_bus().borrow_mut().dispose();
3903 }
3904
3905 #[rstest]
3906 fn test_observe_exec_event_before_dispatch_skips_recent_fill_report() {
3907 let config = LiveNodeConfig {
3908 exec_engine: crate::config::LiveExecEngineConfig {
3909 reconciliation: false,
3910 ..Default::default()
3911 },
3912 ..Default::default()
3913 };
3914 let mut node = LiveNode::build("FillSkipNode".to_string(), Some(config)).unwrap();
3915 let event = stub_exec_event();
3916 let account_id = AccountId::from("TEST-001");
3917 let instrument_id = InstrumentId::from("TEST.VENUE");
3918 let trade_id = TradeId::from("T-001");
3919
3920 let close_ids = node.observe_exec_event_before_dispatch(&event);
3921 assert_eq!(close_ids, Some(Vec::new()));
3922 assert!(
3923 !node
3924 .exec_manager
3925 .is_fill_recently_processed(account_id, instrument_id, trade_id)
3926 );
3927
3928 node.exec_manager
3929 .mark_fill_processed(account_id, instrument_id, trade_id);
3930
3931 let close_ids = node.observe_exec_event_before_dispatch(&event);
3932 assert_eq!(close_ids, None);
3933 }
3934
3935 #[rstest]
3936 #[case(false, false, OrderStatus::Canceled, 1)]
3937 #[case(false, true, OrderStatus::Accepted, 0)]
3938 #[case(true, false, OrderStatus::Canceled, 1)]
3939 #[case(true, true, OrderStatus::Accepted, 0)]
3940 fn test_process_exec_event_clears_terminal_activity_only_after_cached_order_closes(
3941 #[case] with_fills: bool,
3942 #[case] superseded: bool,
3943 #[case] expected_status: OrderStatus,
3944 #[case] expected_query_count: usize,
3945 ) {
3946 let config = LiveNodeConfig {
3947 exec_engine: crate::config::LiveExecEngineConfig {
3948 reconciliation: true,
3949 open_check_threshold_ms: 5_000,
3950 single_order_query_delay_ms: 0,
3951 ..Default::default()
3952 },
3953 ..Default::default()
3954 };
3955 let mut node = LiveNode::build("TerminalReportNode".to_string(), Some(config)).unwrap();
3956 let client_order_id = ClientOrderId::from("O-TERMINAL-REPORT");
3957 let old_venue_order_id = VenueOrderId::from("V-TERMINAL-REPORT-OLD");
3958 let new_venue_order_id = VenueOrderId::from("V-TERMINAL-REPORT-NEW");
3959 let account_id = AccountId::from("TEST-001");
3960 let client_id = ClientId::from("TEST");
3961 let instrument = crypto_perpetual_ethusdt();
3962 let instrument_id = instrument.id();
3963 let account = AccountAny::Margin(MarginAccount::new(
3964 AccountState::new(
3965 account_id,
3966 AccountType::Margin,
3967 vec![AccountBalance::new(
3968 Money::from("1000000 USDT"),
3969 Money::from("0 USDT"),
3970 Money::from("1000000 USDT"),
3971 )],
3972 Vec::new(),
3973 true,
3974 UUID4::new(),
3975 UnixNanos::default(),
3976 UnixNanos::default(),
3977 Some(Currency::USDT()),
3978 ),
3979 true,
3980 ));
3981 node.kernel.cache.borrow_mut().add_account(account).unwrap();
3982 node.kernel
3983 .cache
3984 .borrow_mut()
3985 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
3986 .unwrap();
3987 insert_accepted_limit_order_in_node(
3988 &node,
3989 account_id,
3990 client_id,
3991 instrument_id,
3992 client_order_id,
3993 old_venue_order_id,
3994 );
3995
3996 if superseded {
3997 let order = node
3998 .kernel
3999 .cache
4000 .borrow()
4001 .order_owned(&client_order_id)
4002 .unwrap();
4003 let pending_update = OrderPendingUpdateSpec::builder()
4004 .trader_id(order.trader_id())
4005 .strategy_id(order.strategy_id())
4006 .instrument_id(order.instrument_id())
4007 .client_order_id(client_order_id)
4008 .account_id(account_id)
4009 .venue_order_id(old_venue_order_id)
4010 .build();
4011 node.kernel
4012 .cache
4013 .borrow_mut()
4014 .update_order(&OrderEventAny::PendingUpdate(pending_update))
4015 .unwrap();
4016 let order = node
4017 .kernel
4018 .cache
4019 .borrow()
4020 .order_owned(&client_order_id)
4021 .unwrap();
4022 let updated = OrderUpdatedSpec::builder()
4023 .trader_id(order.trader_id())
4024 .strategy_id(order.strategy_id())
4025 .instrument_id(order.instrument_id())
4026 .client_order_id(client_order_id)
4027 .quantity(order.quantity())
4028 .venue_order_id(new_venue_order_id)
4029 .account_id(account_id)
4030 .build();
4031 node.kernel
4032 .cache
4033 .borrow_mut()
4034 .update_order(&OrderEventAny::Updated(updated))
4035 .unwrap();
4036 }
4037
4038 let report = OrderStatusReport::new(
4039 account_id,
4040 instrument_id,
4041 Some(client_order_id),
4042 old_venue_order_id,
4043 OrderSide::Buy,
4044 OrderType::Limit,
4045 TimeInForce::Gtc,
4046 OrderStatus::Canceled,
4047 Quantity::from("10.0"),
4048 Quantity::from("0.0"),
4049 UnixNanos::from(1_000),
4050 UnixNanos::from(2_000),
4051 UnixNanos::from(3_000),
4052 None,
4053 );
4054 let report = if with_fills {
4055 ExecutionReport::OrderWithFills(Box::new(report), Vec::new())
4056 } else {
4057 ExecutionReport::Order(Box::new(report))
4058 };
4059 let event = ExecutionEvent::Report(report);
4060
4061 node.process_exec_event(event);
4062
4063 let order = node
4064 .kernel
4065 .cache
4066 .borrow()
4067 .order_owned(&client_order_id)
4068 .unwrap();
4069 let expected_venue_order_id = if superseded {
4070 new_venue_order_id
4071 } else {
4072 old_venue_order_id
4073 };
4074
4075 assert_eq!(order.status(), expected_status);
4076 assert_eq!(order.venue_order_id(), Some(expected_venue_order_id));
4077
4078 if !superseded {
4079 let replacement = OrderTestBuilder::new(OrderType::Limit)
4080 .client_order_id(client_order_id)
4081 .instrument_id(instrument_id)
4082 .quantity(Quantity::from("20.0"))
4083 .price(Price::from("200.0"))
4084 .build();
4085 let submitted = TestOrderEventStubs::submitted(&replacement, account_id);
4086 node.kernel
4087 .cache
4088 .borrow_mut()
4089 .add_order(replacement, None, Some(client_id), true)
4090 .unwrap();
4091 let replacement = node
4092 .kernel
4093 .cache
4094 .borrow_mut()
4095 .update_order(&submitted)
4096 .unwrap();
4097 let accepted =
4098 TestOrderEventStubs::accepted(&replacement, account_id, old_venue_order_id);
4099 node.kernel
4100 .cache
4101 .borrow_mut()
4102 .update_order(&accepted)
4103 .unwrap();
4104 }
4105
4106 assert_eq!(
4107 node.exec_manager.check_open_order_queries().len(),
4108 expected_query_count,
4109 );
4110 }
4111
4112 #[rstest]
4113 fn test_rejected_direct_fill_stays_eligible_for_later_report() {
4114 let (mut node, mut fill_event, _) = recent_fill_test_fixture("RejectedDirectFillNode");
4115 let OrderEventAny::Filled(fill) = &mut fill_event else {
4116 unreachable!();
4117 };
4118 fill.client_order_id = ClientOrderId::from("O-UNKNOWN");
4119 fill.venue_order_id = VenueOrderId::from("V-UNKNOWN");
4120 let fill = fill.clone();
4121 let report_event = fill_report_event(&fill);
4122 let event = ExecutionEvent::Order(OrderEventAny::Filled(fill.clone()));
4123
4124 assert!(node.observe_exec_event_before_dispatch(&event).is_some());
4125 let marked_before_dispatch = is_recent_fill(&node, &fill);
4126
4127 node.dispatch_exec_event_and_commit_fill(event);
4128
4129 let marked_after_dispatch = is_recent_fill(&node, &fill);
4130 let later_report_is_eligible = node
4131 .observe_exec_event_before_dispatch(&report_event)
4132 .is_some();
4133 assert_eq!(
4134 (
4135 marked_before_dispatch,
4136 marked_after_dispatch,
4137 later_report_is_eligible,
4138 ),
4139 (false, false, true),
4140 );
4141 }
4142
4143 #[rstest]
4144 fn test_applied_direct_fill_commits_and_skips_later_report() {
4145 let (mut node, fill_event, _) = recent_fill_test_fixture("AppliedDirectFillNode");
4146 let OrderEventAny::Filled(fill) = &fill_event else {
4147 unreachable!();
4148 };
4149 let fill = fill.clone();
4150 let report_event = fill_report_event(&fill);
4151 let event = ExecutionEvent::Order(fill_event);
4152
4153 assert!(node.observe_exec_event_before_dispatch(&event).is_some());
4154 assert!(!is_recent_fill(&node, &fill));
4155
4156 node.dispatch_exec_event_and_commit_fill(event);
4157
4158 assert!(is_recent_fill(&node, &fill));
4159 assert_eq!(node.observe_exec_event_before_dispatch(&report_event), None);
4160 }
4161
4162 #[rstest]
4163 fn test_canonical_duplicate_fill_counts_as_applied() {
4164 let (mut node, fill_event, _) = recent_fill_test_fixture("DuplicateDirectFillNode");
4165 let OrderEventAny::Filled(fill) = &fill_event else {
4166 unreachable!();
4167 };
4168 let mut fill = fill.clone();
4169 node.kernel
4170 .cache
4171 .borrow_mut()
4172 .update_order(&fill_event)
4173 .unwrap();
4174 fill.client_order_id = ClientOrderId::from("O-DUPLICATE-UNKNOWN");
4175
4176 node.exec_manager.commit_recent_fill_if_applied(&fill);
4177
4178 assert!(is_recent_fill(&node, &fill));
4179 }
4180
4181 #[rstest]
4182 #[case(false)]
4183 #[case(true)]
4184 fn test_continuous_reconciliation_commits_only_applied_fill(#[case] applied: bool) {
4185 let (mut node, mut fill_event, _) = recent_fill_test_fixture(if applied {
4186 "AppliedContinuousFillNode"
4187 } else {
4188 "RejectedContinuousFillNode"
4189 });
4190
4191 if !applied {
4192 let OrderEventAny::Filled(fill) = &mut fill_event else {
4193 unreachable!();
4194 };
4195 fill.client_order_id = ClientOrderId::from("O-CONTINUOUS-UNKNOWN");
4196 fill.venue_order_id = VenueOrderId::from("V-CONTINUOUS-UNKNOWN");
4197 }
4198 let OrderEventAny::Filled(fill) = &fill_event else {
4199 unreachable!();
4200 };
4201 let fill = fill.clone();
4202
4203 node.process_reconciliation_events(&[fill_event]);
4204
4205 assert_eq!(is_recent_fill(&node, &fill), applied);
4206 }
4207
4208 #[rstest]
4209 fn test_recent_fill_commit_requires_account_and_instrument_match() {
4210 let (mut node, fill_event, _) = recent_fill_test_fixture("MismatchedDirectFillNode");
4211 node.kernel
4212 .cache
4213 .borrow_mut()
4214 .update_order(&fill_event)
4215 .unwrap();
4216 let OrderEventAny::Filled(fill) = fill_event else {
4217 unreachable!();
4218 };
4219 let mut account_mismatch = fill.clone();
4220 account_mismatch.account_id = AccountId::from("OTHER-001");
4221 let mut instrument_mismatch = fill;
4222 instrument_mismatch.instrument_id = InstrumentId::from("OTHER.VENUE");
4223
4224 node.exec_manager
4225 .commit_recent_fill_if_applied(&account_mismatch);
4226 node.exec_manager
4227 .commit_recent_fill_if_applied(&instrument_mismatch);
4228
4229 assert!(!is_recent_fill(&node, &account_mismatch));
4230 assert!(!is_recent_fill(&node, &instrument_mismatch));
4231 }
4232
4233 #[rstest]
4234 fn test_applied_inferred_fill_remains_recently_processed() {
4235 let (mut node, _, instrument) = recent_fill_test_fixture("InferredFillNode");
4236 let client_order_id = ClientOrderId::from("O-RECENT-FILL");
4237 let venue_order_id = VenueOrderId::from("V-RECENT-FILL");
4238 let account_id = AccountId::from("TEST-001");
4239 let order = node
4240 .kernel
4241 .cache
4242 .borrow()
4243 .order_owned(&client_order_id)
4244 .unwrap();
4245 let report = OrderStatusReport::new(
4246 account_id,
4247 instrument.id(),
4248 Some(client_order_id),
4249 venue_order_id,
4250 OrderSide::Buy,
4251 OrderType::Limit,
4252 TimeInForce::Gtc,
4253 OrderStatus::PartiallyFilled,
4254 Quantity::from("10.0"),
4255 Quantity::from("1.0"),
4256 UnixNanos::from(1_000),
4257 UnixNanos::from(1_000),
4258 UnixNanos::from(1_000),
4259 None,
4260 )
4261 .with_avg_px(dec!(100.0));
4262 let inferred = create_inferred_fill_for_qty(
4263 &order,
4264 &report,
4265 &account_id,
4266 &instrument,
4267 Quantity::from("1.0"),
4268 UnixNanos::from(1_000),
4269 None,
4270 )
4271 .unwrap();
4272 let OrderEventAny::Filled(fill) = &inferred else {
4273 unreachable!();
4274 };
4275 let fill = fill.clone();
4276
4277 node.process_reconciliation_events(&[inferred]);
4278
4279 assert!(fill.reconciliation);
4280 assert!(is_recent_fill(&node, &fill));
4281 }
4282
4283 #[rstest]
4284 fn test_observe_exec_event_before_dispatch_accepted_batch_stamps_local_activity() {
4285 let config = LiveNodeConfig {
4286 exec_engine: crate::config::LiveExecEngineConfig {
4287 reconciliation: true,
4288 open_check_threshold_ms: 5_000,
4289 single_order_query_delay_ms: 0,
4290 ..Default::default()
4291 },
4292 ..Default::default()
4293 };
4294 let mut node = LiveNode::build("AcceptedBatchNode".to_string(), Some(config)).unwrap();
4295 let account_id = AccountId::from("TEST-ACCEPTED-BATCH-001");
4296 let client_id = ClientId::from("TEST-ACCEPTED-BATCH");
4297 let instrument = crypto_perpetual_ethusdt();
4298 let instrument_id = instrument.id();
4299 let client_order_id = ClientOrderId::from("O-ACCEPTED-BATCH");
4300 let venue_order_id = VenueOrderId::from("V-ACCEPTED-BATCH");
4301
4302 node.kernel
4303 .cache
4304 .borrow_mut()
4305 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
4306 .unwrap();
4307 insert_accepted_limit_order_in_node(
4308 &node,
4309 account_id,
4310 client_id,
4311 instrument_id,
4312 client_order_id,
4313 venue_order_id,
4314 );
4315
4316 assert_eq!(node.exec_manager.check_open_order_queries().len(), 1);
4317
4318 let accepted = OrderAcceptedSpec::builder()
4319 .instrument_id(instrument_id)
4320 .client_order_id(client_order_id)
4321 .venue_order_id(venue_order_id)
4322 .account_id(account_id)
4323 .build();
4324 let event = ExecutionEvent::OrderAcceptedBatch(OrderAcceptedBatch::new(vec![accepted]));
4325
4326 let close_ids = node.observe_exec_event_before_dispatch(&event);
4327
4328 assert_eq!(close_ids, Some(Vec::new()));
4329 assert!(node.exec_manager.check_open_order_queries().is_empty());
4330 }
4331
4332 #[rstest]
4333 #[cfg_attr(
4334 not(all(feature = "simulation", madsim)),
4335 tokio::test(start_paused = true)
4336 )]
4337 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
4338 async fn test_batch_cancel_command_registers_each_child_for_inflight_timeout() {
4339 use nautilus_common::messages::execution::{BatchCancelOrders, CancelOrder};
4340 use nautilus_model::{events::OrderPendingCancel, identifiers::ClientOrderId};
4341
4342 let config = LiveNodeConfig {
4343 exec_engine: crate::config::LiveExecEngineConfig {
4344 reconciliation: true,
4345 inflight_check_threshold_ms: 100,
4346 inflight_check_retries: 1,
4347 ..Default::default()
4348 },
4349 ..Default::default()
4350 };
4351 let mut node = LiveNode::build("BatchCancelNode".to_string(), Some(config)).unwrap();
4352 let trader_id = TraderId::from("TESTER-001");
4353 let strategy_id = StrategyId::from("S-BATCH-CANCEL");
4354 let account_id = AccountId::from("TEST-001");
4355 let instrument = crypto_perpetual_ethusdt();
4356 let instrument_id = instrument.id();
4357 let child_ids = [
4358 ClientOrderId::from("O-BATCH-CANCEL-1"),
4359 ClientOrderId::from("O-BATCH-CANCEL-2"),
4360 ];
4361 node.kernel
4362 .cache
4363 .borrow_mut()
4364 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
4365 .unwrap();
4366
4367 for client_order_id in child_ids {
4368 let order = OrderTestBuilder::new(OrderType::Limit)
4369 .trader_id(trader_id)
4370 .strategy_id(strategy_id)
4371 .client_order_id(client_order_id)
4372 .instrument_id(instrument_id)
4373 .quantity(Quantity::from("10.0"))
4374 .price(Price::from("100.0"))
4375 .build();
4376 let venue_order_id = VenueOrderId::from(format!("V-{client_order_id}").as_str());
4377 let submitted = TestOrderEventStubs::submitted(&order, account_id);
4381 let accepted = TestOrderEventStubs::accepted(&order, account_id, venue_order_id);
4382 let pending_cancel = OrderEventAny::PendingCancel(OrderPendingCancel::new(
4383 trader_id,
4384 strategy_id,
4385 instrument_id,
4386 client_order_id,
4387 Some(account_id),
4388 UUID4::new(),
4389 UnixNanos::default(),
4390 UnixNanos::default(),
4391 false,
4392 Some(venue_order_id),
4393 ));
4394 let mut cache = node.kernel.cache.borrow_mut();
4395 cache.add_order(order, None, None, false).unwrap();
4396 cache.update_order(&submitted).unwrap();
4397 cache.update_order(&accepted).unwrap();
4398 cache.update_order(&pending_cancel).unwrap();
4399 }
4400 let cancels = child_ids
4401 .into_iter()
4402 .map(|client_order_id| {
4403 CancelOrder::new(
4404 trader_id,
4405 None,
4406 strategy_id,
4407 instrument_id,
4408 client_order_id,
4409 None,
4410 UUID4::new(),
4411 UnixNanos::default(),
4412 None,
4413 None,
4414 )
4415 })
4416 .collect();
4417 let command = TradingCommand::CancelOrders(BatchCancelOrders::new(
4418 trader_id,
4419 None,
4420 strategy_id,
4421 instrument_id,
4422 cancels,
4423 UUID4::new(),
4424 UnixNanos::default(),
4425 None,
4426 None,
4427 ));
4428
4429 node.observe_exec_command_before_dispatch(&command);
4430 advance_clock(Duration::from_millis(101)).await;
4431 let result = node.exec_manager.check_inflight_orders();
4432 let timed_out_ids = result
4433 .events
4434 .iter()
4435 .map(OrderEventAny::client_order_id)
4436 .collect::<IndexSet<_>>();
4437
4438 assert_eq!(timed_out_ids, IndexSet::from(child_ids));
4439 assert_eq!(result.events.len(), child_ids.len());
4440 assert!(
4441 result
4442 .events
4443 .iter()
4444 .all(|event| matches!(event, OrderEventAny::Canceled(_))),
4445 "batch-cancel children must time out as Canceled events",
4446 );
4447 }
4448
4449 #[rstest]
4450 #[cfg_attr(
4451 not(all(feature = "simulation", madsim)),
4452 tokio::test(start_paused = true)
4453 )]
4454 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
4455 async fn test_risk_bound_command_does_not_register_inflight() {
4456 let config = LiveNodeConfig {
4457 exec_engine: crate::config::LiveExecEngineConfig {
4458 reconciliation: true,
4459 inflight_check_threshold_ms: 100,
4460 inflight_check_retries: 2,
4461 ..Default::default()
4462 },
4463 ..Default::default()
4464 };
4465 let mut node = LiveNode::build("RiskBoundNode".to_string(), Some(config)).unwrap();
4466 msgbus::register_trading_command_endpoint(
4467 MessagingSwitchboard::risk_engine_execute(),
4468 TypedIntoHandler::from(|_: TradingCommand| {}),
4469 );
4470 let instrument = crypto_perpetual_ethusdt();
4471 let instrument_id = instrument.id();
4472 let order = OrderTestBuilder::new(OrderType::Limit)
4473 .trader_id(node.trader_id())
4474 .strategy_id(StrategyId::from("S-RISK-DENIED"))
4475 .instrument_id(instrument_id)
4476 .side(OrderSide::NoOrderSide)
4477 .quantity(Quantity::from("1.000"))
4478 .price(Price::from("100.00"))
4479 .build();
4480 let client_order_id = order.client_order_id();
4481
4482 {
4483 let mut cache = node.kernel.cache.borrow_mut();
4484 cache
4485 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
4486 .unwrap();
4487 cache.add_order(order.clone(), None, None, false).unwrap();
4488 }
4489
4490 let submit_order = SubmitOrder::new(
4491 order.trader_id(),
4492 None,
4493 order.strategy_id(),
4494 instrument_id,
4495 client_order_id,
4496 order.init_event().clone(),
4497 None,
4498 None,
4499 None,
4500 UUID4::new(),
4501 UnixNanos::default(),
4502 None,
4503 );
4504 node.process_exec_command(TradingCommandMessage::new(
4505 MessagingSwitchboard::risk_engine_execute(),
4506 TradingCommand::SubmitOrder(submit_order),
4507 ));
4508
4509 advance_clock(Duration::from_millis(101)).await;
4510 let result = node.exec_manager.check_inflight_orders();
4511 let status = node
4512 .kernel
4513 .cache
4514 .borrow()
4515 .order(&client_order_id)
4516 .unwrap()
4517 .status();
4518
4519 assert_eq!(status, OrderStatus::Initialized);
4520 assert_eq!(
4521 node.exec_manager.recon_check_retry_count(&client_order_id),
4522 0
4523 );
4524 assert!(result.events.is_empty());
4525 assert!(result.queries.is_empty());
4526 }
4527
4528 #[rstest]
4529 #[cfg_attr(
4530 not(all(feature = "simulation", madsim)),
4531 tokio::test(start_paused = true)
4532 )]
4533 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
4534 async fn test_risk_approved_command_registers_inflight() {
4535 let config = LiveNodeConfig {
4536 risk_engine: crate::config::LiveRiskEngineConfig {
4537 bypass: true,
4538 ..Default::default()
4539 },
4540 exec_engine: crate::config::LiveExecEngineConfig {
4541 reconciliation: true,
4542 inflight_check_threshold_ms: 100,
4543 inflight_check_retries: 2,
4544 ..Default::default()
4545 },
4546 ..Default::default()
4547 };
4548 let mut node = LiveNode::build("RiskApprovedNode".to_string(), Some(config)).unwrap();
4549 msgbus::register_trading_command_endpoint(
4550 MessagingSwitchboard::exec_engine_execute(),
4551 TypedIntoHandler::from(|_: TradingCommand| {}),
4552 );
4553 let instrument = crypto_perpetual_ethusdt();
4554 let instrument_id = instrument.id();
4555 let order = OrderTestBuilder::new(OrderType::Limit)
4556 .trader_id(node.trader_id())
4557 .strategy_id(StrategyId::from("S-RISK-APPROVED"))
4558 .instrument_id(instrument_id)
4559 .side(OrderSide::Buy)
4560 .quantity(Quantity::from("1.000"))
4561 .price(Price::from("100.00"))
4562 .build();
4563 let client_order_id = order.client_order_id();
4564
4565 {
4566 let mut cache = node.kernel.cache.borrow_mut();
4567 cache
4568 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
4569 .unwrap();
4570 cache.add_order(order.clone(), None, None, false).unwrap();
4571 }
4572
4573 let submit_order = SubmitOrder::new(
4574 order.trader_id(),
4575 None,
4576 order.strategy_id(),
4577 instrument_id,
4578 client_order_id,
4579 order.init_event().clone(),
4580 None,
4581 None,
4582 None,
4583 UUID4::new(),
4584 UnixNanos::default(),
4585 None,
4586 );
4587 node.process_exec_command(TradingCommandMessage::new(
4588 MessagingSwitchboard::risk_engine_execute(),
4589 TradingCommand::SubmitOrder(submit_order),
4590 ));
4591
4592 advance_clock(Duration::from_millis(101)).await;
4593 let result = node.exec_manager.check_inflight_orders();
4594 let [TradingCommand::QueryOrder(query)] = result.queries.as_slice() else {
4595 panic!("expected one query order command");
4596 };
4597
4598 assert_eq!(query.client_order_id, client_order_id);
4599 assert_eq!(
4600 node.exec_manager.recon_check_retry_count(&client_order_id),
4601 1
4602 );
4603 assert!(result.events.is_empty());
4604 }
4605
4606 #[rstest]
4607 fn test_live_node_builder_clock_factory_drives_kernel_clock() {
4608 let calls = Rc::new(Cell::new(0usize));
4609 let calls_in_factory = calls.clone();
4610 let sentinel = UnixNanos::from(123_456_789_u64);
4611
4612 let node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4613 .unwrap()
4614 .with_reconciliation(false)
4615 .with_clock_factory(move || {
4616 calls_in_factory.set(calls_in_factory.get() + 1);
4617 let mut clock = TestClock::new();
4618 clock.advance_time(sentinel, true);
4619 Rc::new(RefCell::new(clock)) as Rc<RefCell<dyn Clock>>
4620 })
4621 .build()
4622 .unwrap();
4623
4624 assert_eq!(node.kernel().clock().borrow().timestamp_ns(), sentinel);
4625 assert_eq!(calls.get(), 1);
4626 }
4627
4628 #[derive(Debug)]
4629 struct ReplayKernelEventStore {
4630 fail_restore: bool,
4631 }
4632
4633 impl KernelEventStore for ReplayKernelEventStore {
4634 fn restore_parent_cache(
4635 &mut self,
4636 _instance_id: UUID4,
4637 _cache: &mut Cache,
4638 ) -> anyhow::Result<()> {
4639 if self.fail_restore {
4640 anyhow::bail!("replay restore failed");
4641 }
4642
4643 Ok(())
4644 }
4645
4646 fn open(
4647 &mut self,
4648 _instance_id: UUID4,
4649 _components: &RegisteredComponents,
4650 _environment: Environment,
4651 ) -> anyhow::Result<()> {
4652 Ok(())
4653 }
4654
4655 fn snapshot_anchorer(&self) -> Option<SnapshotAnchorer> {
4656 None
4657 }
4658
4659 fn seal(&mut self, _ts_init: UnixNanos) {}
4660
4661 fn run_id(&self) -> Option<&str> {
4662 Some("replay-child")
4663 }
4664
4665 fn parent_run_id(&self) -> Option<&str> {
4666 Some("seed-run")
4667 }
4668
4669 fn is_event_store_replay_configured(&self) -> bool {
4670 true
4671 }
4672
4673 fn is_halted(&self) -> bool {
4674 false
4675 }
4676 }
4677
4678 #[derive(Debug)]
4679 struct TestStrategy {
4680 core: StrategyCore,
4681 }
4682
4683 impl TestStrategy {
4684 fn new(config: StrategyConfig) -> Self {
4685 Self {
4686 core: StrategyCore::new(config),
4687 }
4688 }
4689 }
4690
4691 impl DataActor for TestStrategy {}
4692
4693 nautilus_strategy!(TestStrategy, {
4694 fn external_order_claims(&self) -> Option<Vec<InstrumentId>> {
4695 self.core.config.external_order_claims.clone()
4696 }
4697 });
4698
4699 fn live_node_with_replay_store(fail_restore: bool) -> LiveNode {
4700 let builder = LiveNodeBuilder::new(TraderId::default(), Environment::Live)
4703 .unwrap()
4704 .with_exec_engine_config(crate::config::LiveExecEngineConfig {
4705 reconciliation: false,
4706 ..Default::default()
4707 })
4708 .with_load_state(true)
4709 .with_name("TestKernel")
4710 .with_event_store(move |_instance_id: UUID4, _clock: Rc<RefCell<dyn Clock>>| {
4711 Ok(Box::new(ReplayKernelEventStore { fail_restore }) as Box<dyn KernelEventStore>)
4712 });
4713
4714 builder.build().unwrap()
4715 }
4716
4717 #[rstest]
4718 fn test_add_strategy_registers_external_order_claims_with_manager_and_engine() {
4719 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4720 .unwrap()
4721 .with_reconciliation(false)
4722 .with_delay_post_stop_secs(0)
4723 .with_timeout_connection(1)
4724 .build()
4725 .unwrap();
4726 let instrument_id = InstrumentId::from("AUDUSD.SIM");
4727 let strategy_id = StrategyId::from("CLAIMS-001");
4728
4729 node.add_strategy(TestStrategy::new(StrategyConfig {
4730 strategy_id: Some(strategy_id),
4731 external_order_claims: Some(vec![instrument_id]),
4732 ..Default::default()
4733 }))
4734 .unwrap();
4735
4736 assert_eq!(
4737 node.exec_manager.get_external_order_claim(&instrument_id),
4738 Some(strategy_id)
4739 );
4740
4741 {
4742 let exec_engine = node.kernel().exec_engine.borrow();
4743 assert_eq!(
4744 exec_engine.get_external_order_claim(&instrument_id),
4745 Some(strategy_id)
4746 );
4747 }
4748 }
4749
4750 #[rstest]
4751 fn test_register_external_order_claims_after_build_reaches_both_tiers() {
4752 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4753 .unwrap()
4754 .with_reconciliation(false)
4755 .build()
4756 .unwrap();
4757 let instrument_id = InstrumentId::from("AUDUSD.SIM");
4758 let strategy_id = StrategyId::from("CLAIMS-001");
4759
4760 node.register_external_order_claims(strategy_id, &[instrument_id])
4761 .unwrap();
4762
4763 assert_eq!(
4764 node.exec_manager.get_external_order_claim(&instrument_id),
4765 Some(strategy_id)
4766 );
4767 assert_eq!(
4768 node.kernel
4769 .exec_engine
4770 .borrow()
4771 .get_external_order_claim(&instrument_id),
4772 Some(strategy_id)
4773 );
4774 }
4775
4776 #[rstest]
4777 #[tokio::test]
4778 async fn test_register_external_order_claims_while_running_reaches_both_tiers() {
4779 let mut node = live_node_with_replay_store(false);
4780 let instrument_id = InstrumentId::from("AUDUSD.SIM");
4781 let strategy_id = StrategyId::from("CLAIMS-001");
4782
4783 node.start().await.unwrap();
4784 assert_eq!(node.state(), NodeState::Running);
4785
4786 node.register_external_order_claims(strategy_id, &[instrument_id])
4787 .unwrap();
4788
4789 assert_eq!(
4790 node.exec_manager.get_external_order_claim(&instrument_id),
4791 Some(strategy_id)
4792 );
4793 assert_eq!(
4794 node.kernel
4795 .exec_engine
4796 .borrow()
4797 .get_external_order_claim(&instrument_id),
4798 Some(strategy_id)
4799 );
4800 }
4801
4802 #[rstest]
4803 fn test_register_external_order_claims_conflicting_batch_leaves_new_claims_absent() {
4804 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4805 .unwrap()
4806 .with_reconciliation(false)
4807 .build()
4808 .unwrap();
4809 let existing_instrument = InstrumentId::from("AUDUSD.SIM");
4810 let new_instruments = [
4811 InstrumentId::from("EURUSD.SIM"),
4812 InstrumentId::from("GBPUSD.SIM"),
4813 ];
4814 let existing_strategy_id = StrategyId::from("CLAIMS-001");
4815 let new_strategy_id = StrategyId::from("CLAIMS-002");
4816 node.register_external_order_claims(existing_strategy_id, &[existing_instrument])
4817 .unwrap();
4818
4819 let result = node.register_external_order_claims(
4820 new_strategy_id,
4821 &[new_instruments[0], existing_instrument, new_instruments[1]],
4822 );
4823
4824 assert!(result.is_err());
4825 assert_eq!(
4826 node.exec_manager
4827 .get_external_order_claim(&existing_instrument),
4828 Some(existing_strategy_id)
4829 );
4830
4831 for instrument_id in new_instruments {
4832 assert_eq!(
4833 node.exec_manager.get_external_order_claim(&instrument_id),
4834 None
4835 );
4836 assert_eq!(
4837 node.kernel
4838 .exec_engine
4839 .borrow()
4840 .get_external_order_claim(&instrument_id),
4841 None
4842 );
4843 }
4844 }
4845
4846 #[rstest]
4847 fn test_register_external_order_claims_one_tier_conflict_changes_neither_tier() {
4848 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4849 .unwrap()
4850 .with_reconciliation(false)
4851 .build()
4852 .unwrap();
4853 let conflicting_instrument = InstrumentId::from("AUDUSD.SIM");
4854 let new_instrument = InstrumentId::from("EURUSD.SIM");
4855 let existing_strategy_id = StrategyId::from("CLAIMS-001");
4856 let new_strategy_id = StrategyId::from("CLAIMS-002");
4857 node.exec_manager
4858 .claim_external_orders(conflicting_instrument, existing_strategy_id)
4859 .unwrap();
4860
4861 let result = node.register_external_order_claims(
4862 new_strategy_id,
4863 &[new_instrument, conflicting_instrument],
4864 );
4865
4866 assert!(result.is_err());
4867 assert_eq!(
4868 node.exec_manager
4869 .get_external_order_claim(&conflicting_instrument),
4870 Some(existing_strategy_id)
4871 );
4872 assert_eq!(
4873 node.kernel
4874 .exec_engine
4875 .borrow()
4876 .get_external_order_claim(&conflicting_instrument),
4877 None
4878 );
4879 assert_eq!(
4880 node.exec_manager.get_external_order_claim(&new_instrument),
4881 None
4882 );
4883 assert_eq!(
4884 node.kernel
4885 .exec_engine
4886 .borrow()
4887 .get_external_order_claim(&new_instrument),
4888 None
4889 );
4890 }
4891
4892 #[rstest]
4893 fn test_register_external_order_claims_engine_only_conflict_changes_neither_tier() {
4894 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4895 .unwrap()
4896 .with_reconciliation(false)
4897 .build()
4898 .unwrap();
4899 let conflicting_instrument = InstrumentId::from("AUDUSD.SIM");
4900 let new_instrument = InstrumentId::from("EURUSD.SIM");
4901 let existing_strategy_id = StrategyId::from("CLAIMS-001");
4902 let new_strategy_id = StrategyId::from("CLAIMS-002");
4903
4904 node.kernel
4906 .exec_engine
4907 .borrow_mut()
4908 .register_external_order_claims(
4909 existing_strategy_id,
4910 &HashSet::from([conflicting_instrument]),
4911 )
4912 .unwrap();
4913
4914 let result = node.register_external_order_claims(
4915 new_strategy_id,
4916 &[new_instrument, conflicting_instrument],
4917 );
4918
4919 assert!(result.is_err());
4920 assert_eq!(
4921 node.kernel
4922 .exec_engine
4923 .borrow()
4924 .get_external_order_claim(&conflicting_instrument),
4925 Some(existing_strategy_id)
4926 );
4927 assert_eq!(
4928 node.exec_manager
4929 .get_external_order_claim(&conflicting_instrument),
4930 None
4931 );
4932 assert_eq!(
4933 node.kernel
4934 .exec_engine
4935 .borrow()
4936 .get_external_order_claim(&new_instrument),
4937 None
4938 );
4939 assert_eq!(
4940 node.exec_manager.get_external_order_claim(&new_instrument),
4941 None
4942 );
4943 }
4944
4945 #[rstest]
4946 fn test_deregister_external_order_claims_allows_successor_to_claim() {
4947 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4948 .unwrap()
4949 .with_reconciliation(false)
4950 .build()
4951 .unwrap();
4952 let instruments = [
4953 InstrumentId::from("AUDUSD.SIM"),
4954 InstrumentId::from("EURUSD.SIM"),
4955 ];
4956 let first_strategy_id = StrategyId::from("CLAIMS-001");
4957 let successor_strategy_id = StrategyId::from("CLAIMS-002");
4958 node.register_external_order_claims(first_strategy_id, &instruments)
4959 .unwrap();
4960
4961 node.deregister_external_order_claims(first_strategy_id)
4962 .unwrap();
4963 node.register_external_order_claims(successor_strategy_id, &instruments)
4964 .unwrap();
4965
4966 for instrument_id in instruments {
4967 assert_eq!(
4968 node.exec_manager.get_external_order_claim(&instrument_id),
4969 Some(successor_strategy_id)
4970 );
4971 assert_eq!(
4972 node.kernel
4973 .exec_engine
4974 .borrow()
4975 .get_external_order_claim(&instrument_id),
4976 Some(successor_strategy_id)
4977 );
4978 }
4979 }
4980
4981 #[rstest]
4982 fn test_deregister_external_order_claims_divergence_changes_neither_tier() {
4983 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
4984 .unwrap()
4985 .with_reconciliation(false)
4986 .build()
4987 .unwrap();
4988 let manager_instrument = InstrumentId::from("AUDUSD.SIM");
4989 let engine_instrument = InstrumentId::from("EURUSD.SIM");
4990 let strategy_id = StrategyId::from("CLAIMS-001");
4991 node.exec_manager
4992 .claim_external_orders(manager_instrument, strategy_id)
4993 .unwrap();
4994 node.kernel
4995 .exec_engine
4996 .borrow_mut()
4997 .register_external_order_claims(strategy_id, &HashSet::from([engine_instrument]))
4998 .unwrap();
4999
5000 let result = node.deregister_external_order_claims(strategy_id);
5001
5002 assert!(result.is_err());
5003 assert_eq!(
5004 node.exec_manager
5005 .get_external_order_claim(&manager_instrument),
5006 Some(strategy_id)
5007 );
5008 assert_eq!(
5009 node.kernel
5010 .exec_engine
5011 .borrow()
5012 .get_external_order_claim(&engine_instrument),
5013 Some(strategy_id)
5014 );
5015 }
5016
5017 #[rstest]
5018 fn test_deregister_external_order_claims_without_claims_is_idempotent() {
5019 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
5020 .unwrap()
5021 .with_reconciliation(false)
5022 .build()
5023 .unwrap();
5024 let strategy_id = StrategyId::from("CLAIMS-001");
5025
5026 node.deregister_external_order_claims(strategy_id).unwrap();
5027 node.deregister_external_order_claims(strategy_id).unwrap();
5028 }
5029
5030 #[rstest]
5031 fn test_add_strategy_rejects_duplicate_external_order_claim_without_overwriting() {
5032 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
5033 .unwrap()
5034 .with_reconciliation(false)
5035 .with_delay_post_stop_secs(0)
5036 .with_timeout_connection(1)
5037 .build()
5038 .unwrap();
5039 let instrument_id = InstrumentId::from("AUDUSD.SIM");
5040 let strategy_id = StrategyId::from("CLAIMS-001");
5041 let duplicate_strategy_id = StrategyId::from("OTHER-002");
5042
5043 node.add_strategy(TestStrategy::new(StrategyConfig {
5044 strategy_id: Some(strategy_id),
5045 external_order_claims: Some(vec![instrument_id]),
5046 ..Default::default()
5047 }))
5048 .unwrap();
5049
5050 let result = node.add_strategy(TestStrategy::new(StrategyConfig {
5051 strategy_id: Some(duplicate_strategy_id),
5052 external_order_claims: Some(vec![instrument_id]),
5053 ..Default::default()
5054 }));
5055
5056 assert!(result.is_err());
5057 assert!(
5058 result
5059 .unwrap_err()
5060 .to_string()
5061 .contains("already exists for CLAIMS-001")
5062 );
5063 assert_eq!(
5064 node.exec_manager.get_external_order_claim(&instrument_id),
5065 Some(strategy_id)
5066 );
5067
5068 {
5069 let exec_engine = node.kernel().exec_engine.borrow();
5070 assert_eq!(
5071 exec_engine.get_external_order_claim(&instrument_id),
5072 Some(strategy_id)
5073 );
5074 }
5075 }
5076
5077 #[rstest]
5078 fn test_add_strategy_rejects_repeated_external_order_claim_without_registering() {
5079 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
5080 .unwrap()
5081 .with_reconciliation(false)
5082 .with_delay_post_stop_secs(0)
5083 .with_timeout_connection(1)
5084 .build()
5085 .unwrap();
5086 let instrument_id = InstrumentId::from("AUDUSD.SIM");
5087 let strategy_id = StrategyId::from("CLAIMS-001");
5088
5089 let result = node.add_strategy(TestStrategy::new(StrategyConfig {
5090 strategy_id: Some(strategy_id),
5091 external_order_claims: Some(vec![instrument_id, instrument_id]),
5092 ..Default::default()
5093 }));
5094
5095 assert!(result.is_err());
5096 assert!(
5097 result
5098 .unwrap_err()
5099 .to_string()
5100 .contains("already exists for CLAIMS-001")
5101 );
5102 assert_eq!(
5103 node.exec_manager.get_external_order_claim(&instrument_id),
5104 None
5105 );
5106
5107 {
5108 let exec_engine = node.kernel().exec_engine.borrow();
5109 assert_eq!(exec_engine.get_external_order_claim(&instrument_id), None);
5110 }
5111 }
5112
5113 #[rstest]
5114 fn test_add_strategy_failure_does_not_register_external_order_claims() {
5115 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
5116 .unwrap()
5117 .with_reconciliation(false)
5118 .with_delay_post_stop_secs(0)
5119 .with_timeout_connection(1)
5120 .build()
5121 .unwrap();
5122 let instrument_id = InstrumentId::from("AUDUSD.SIM");
5123 let strategy_id = StrategyId::from("CLAIMS-001");
5124 let mut strategy = TestStrategy::new(StrategyConfig {
5125 strategy_id: Some(strategy_id),
5126 external_order_claims: Some(vec![instrument_id]),
5127 ..Default::default()
5128 });
5129
5130 strategy
5131 .core
5132 .register(
5133 node.trader_id(),
5134 node.kernel.clock(),
5135 node.kernel.cache.clone(),
5136 node.kernel.portfolio.clone(),
5137 )
5138 .unwrap();
5139
5140 let result = node.add_strategy(strategy);
5141
5142 assert!(result.is_err());
5143 assert!(
5144 result
5145 .unwrap_err()
5146 .to_string()
5147 .contains("already registered with trader")
5148 );
5149 assert_eq!(
5150 node.exec_manager.get_external_order_claim(&instrument_id),
5151 None
5152 );
5153 assert_eq!(
5154 node.kernel
5155 .exec_engine
5156 .borrow()
5157 .get_external_order_claim(&instrument_id),
5158 None
5159 );
5160 }
5161
5162 #[rstest]
5163 fn test_add_strategy_without_claims_or_oms_type_does_not_require_engine_borrow() {
5164 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
5165 .unwrap()
5166 .with_reconciliation(false)
5167 .build()
5168 .unwrap();
5169 let exec_engine = node.kernel.exec_engine.clone();
5170 let _engine_borrow = exec_engine.borrow_mut();
5171
5172 node.add_strategy(TestStrategy::new(StrategyConfig {
5173 strategy_id: Some(StrategyId::from("NOCLAIMS-001")),
5174 ..Default::default()
5175 }))
5176 .unwrap();
5177 }
5178
5179 #[rstest]
5180 fn test_add_strategy_registers_configured_hedging_oms_type() {
5181 let mut node = LiveNode::builder(TraderId::from("TESTER-001"), Environment::Sandbox)
5182 .unwrap()
5183 .with_reconciliation(false)
5184 .with_delay_post_stop_secs(0)
5185 .with_timeout_connection(1)
5186 .build()
5187 .unwrap();
5188 let strategy_id = StrategyId::from("FUNDING_ARBITRAGE-001");
5189
5190 node.add_strategy(TestStrategy::new(StrategyConfig {
5191 strategy_id: Some(strategy_id),
5192 oms_type: Some(OmsType::Hedging),
5193 ..Default::default()
5194 }))
5195 .unwrap();
5196
5197 let instrument = crypto_perpetual_ethusdt();
5198 let instrument_id = instrument.id();
5199 let client_id = ClientId::from("STUB");
5200
5201 node.kernel
5202 .cache
5203 .borrow_mut()
5204 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
5205 .unwrap();
5206 node.kernel
5207 .exec_engine
5208 .borrow_mut()
5209 .register_client(Box::new(StubExecutionClient::new(
5210 client_id,
5211 AccountId::from("TEST-ACCOUNT"),
5212 instrument_id.venue,
5213 OmsType::Netting,
5214 None,
5215 )))
5216 .unwrap();
5217
5218 let order = OrderTestBuilder::new(OrderType::Market)
5219 .trader_id(node.trader_id())
5220 .strategy_id(strategy_id)
5221 .instrument_id(instrument_id)
5222 .quantity(Quantity::from("1.000"))
5223 .build();
5224 let position_id = PositionId::new("CUSTOM-POSITION-001");
5225
5226 node.kernel
5227 .cache
5228 .borrow_mut()
5229 .add_order(order.clone(), Some(position_id), Some(client_id), true)
5230 .unwrap();
5231
5232 let submit_order = SubmitOrder::new(
5233 order.trader_id(),
5234 Some(client_id),
5235 strategy_id,
5236 instrument_id,
5237 order.client_order_id(),
5238 order.init_event().clone(),
5239 order.exec_algorithm_id(),
5240 Some(position_id),
5241 None,
5242 UUID4::new(),
5243 UnixNanos::default(),
5244 None,
5245 );
5246
5247 node.kernel
5248 .exec_engine
5249 .borrow()
5250 .execute(TradingCommand::SubmitOrder(submit_order));
5251
5252 let exec_engine = node.kernel.exec_engine.borrow();
5253 let cache = exec_engine.cache().borrow();
5254 let cached_order = cache
5255 .order(&order.client_order_id())
5256 .expect("Order should be cached");
5257
5258 assert_eq!(cached_order.status(), OrderStatus::Initialized);
5259 }
5260
5261 #[cfg(all(feature = "simulation", madsim))]
5262 async fn advance_clock(d: Duration) {
5263 madsim::time::advance(d);
5264 madsim::task::yield_now().await;
5265 }
5266
5267 #[cfg(not(all(feature = "simulation", madsim)))]
5268 async fn advance_clock(d: Duration) {
5269 tokio::time::advance(d).await;
5270 }
5271
5272 #[cfg_attr(
5273 not(all(feature = "simulation", madsim)),
5274 tokio::test(start_paused = true)
5275 )]
5276 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
5277 async fn test_reconciliation_check_due_uses_monotonic_elapsed_time() {
5278 let last = dst::time::Instant::now();
5279 let interval = Duration::from_millis(100);
5280
5281 assert!(!reconciliation_check_due(last, last, Duration::ZERO));
5282 assert!(!reconciliation_check_due(last, last, interval));
5283
5284 advance_clock(Duration::from_millis(99)).await;
5285 let before_interval = dst::time::Instant::now();
5286 assert!(!reconciliation_check_due(before_interval, last, interval));
5287
5288 advance_clock(Duration::from_millis(1)).await;
5289 let at_interval = dst::time::Instant::now();
5290 assert!(reconciliation_check_due(at_interval, last, interval));
5291
5292 assert!(!reconciliation_check_due(last, at_interval, interval));
5293 }
5294
5295 #[cfg_attr(
5296 not(all(feature = "simulation", madsim)),
5297 tokio::test(start_paused = true)
5298 )]
5299 #[cfg_attr(all(feature = "simulation", madsim), madsim::test)]
5300 async fn test_run_reconciliation_checks_does_not_publish_open_order_queries() {
5301 let config = LiveNodeConfig {
5302 exec_engine: crate::config::LiveExecEngineConfig {
5303 reconciliation: true,
5304 open_check_interval_secs: Some(1.0),
5305 position_check_interval_secs: Some(1.0),
5306 max_single_order_queries_per_cycle: 5,
5307 ..Default::default()
5308 },
5309 ..Default::default()
5310 };
5311 let mut node =
5312 LiveNode::build("ReconciliationFallbackNode".to_string(), Some(config)).unwrap();
5313 let client_id = ClientId::from("TEST-QUERY");
5314 let account_id = AccountId::from("TEST-QUERY-001");
5315
5316 let trading_commands = Rc::new(RefCell::new(Vec::new()));
5317 msgbus::register_trading_command_endpoint(
5318 MessagingSwitchboard::exec_engine_execute(),
5319 TypedIntoHandler::from({
5320 let trading_commands = trading_commands.clone();
5321 move |command: TradingCommand| {
5322 trading_commands.borrow_mut().push(command);
5323 }
5324 }),
5325 );
5326
5327 let venue_order_id = VenueOrderId::from("V-NODE-QUERY-001");
5328 let instrument = crypto_perpetual_ethusdt();
5329 let instrument_id = instrument.id();
5330 let client_order_id = ClientOrderId::from("O-NODE-QUERY-001");
5331
5332 node.kernel
5333 .cache
5334 .borrow_mut()
5335 .add_instrument(InstrumentAny::CryptoPerpetual(instrument))
5336 .unwrap();
5337 insert_accepted_limit_order_in_node(
5338 &node,
5339 account_id,
5340 client_id,
5341 instrument_id,
5342 client_order_id,
5343 venue_order_id,
5344 );
5345
5346 let last = dst::time::Instant::now();
5347 advance_clock(Duration::from_nanos(1)).await;
5348 let now = dst::time::Instant::now();
5349 let mut last_inflight_check = last;
5350 let mut last_open_check = last;
5351 let mut last_position_check = last;
5352 let mut open_order_report_task = None;
5353 let mut targeted_order_report_task = None;
5354 let mut position_report_task = None;
5355
5356 node.run_reconciliation_checks(
5357 now,
5358 ReconciliationCheckIntervals {
5359 inflight: Duration::ZERO,
5360 open: Duration::from_nanos(1),
5361 position: Duration::ZERO,
5362 },
5363 &mut ReconciliationCheckState {
5364 last_inflight_check: &mut last_inflight_check,
5365 last_open_check: &mut last_open_check,
5366 last_position_check: &mut last_position_check,
5367 open_order_report_task: &mut open_order_report_task,
5368 targeted_order_report_task: &mut targeted_order_report_task,
5369 position_report_task: &mut position_report_task,
5370 },
5371 );
5372
5373 let commands = trading_commands.borrow();
5374
5375 assert!(commands.is_empty());
5376 assert!(open_order_report_task.is_none());
5377 assert!(targeted_order_report_task.is_none());
5378 assert!(position_report_task.is_none());
5379
5380 ExecutionEngine::register_msgbus_handlers(&node.kernel.exec_engine);
5381 }
5382
5383 fn insert_accepted_limit_order_in_node(
5384 node: &LiveNode,
5385 account_id: AccountId,
5386 client_id: ClientId,
5387 instrument_id: InstrumentId,
5388 client_order_id: ClientOrderId,
5389 venue_order_id: VenueOrderId,
5390 ) {
5391 let order = OrderTestBuilder::new(OrderType::Limit)
5392 .client_order_id(client_order_id)
5393 .instrument_id(instrument_id)
5394 .quantity(Quantity::from("10.0"))
5395 .price(Price::from("100.0"))
5396 .build();
5397 let submitted = TestOrderEventStubs::submitted(&order, account_id);
5398 node.kernel
5399 .cache
5400 .borrow_mut()
5401 .add_order(order, None, Some(client_id), false)
5402 .unwrap();
5403 let order = node
5404 .kernel
5405 .cache
5406 .borrow_mut()
5407 .update_order(&submitted)
5408 .unwrap();
5409 let accepted = TestOrderEventStubs::accepted(&order, account_id, venue_order_id);
5410 node.kernel
5411 .cache
5412 .borrow_mut()
5413 .update_order(&accepted)
5414 .unwrap();
5415 }
5416
5417 fn recent_fill_test_fixture(name: &str) -> (LiveNode, OrderEventAny, InstrumentAny) {
5418 let config = LiveNodeConfig {
5419 exec_engine: crate::config::LiveExecEngineConfig {
5420 reconciliation: true,
5421 ..Default::default()
5422 },
5423 ..Default::default()
5424 };
5425 let node = LiveNode::build(name.to_string(), Some(config)).unwrap();
5426 let account_id = AccountId::from("TEST-001");
5427 let client_id = ClientId::from("TEST-RECENT-FILL");
5428 let client_order_id = ClientOrderId::from("O-RECENT-FILL");
5429 let venue_order_id = VenueOrderId::from("V-RECENT-FILL");
5430 let instrument = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
5431 node.kernel
5432 .cache
5433 .borrow_mut()
5434 .add_instrument(instrument.clone())
5435 .unwrap();
5436 insert_accepted_limit_order_in_node(
5437 &node,
5438 account_id,
5439 client_id,
5440 instrument.id(),
5441 client_order_id,
5442 venue_order_id,
5443 );
5444 let order = node
5445 .kernel
5446 .cache
5447 .borrow()
5448 .order_owned(&client_order_id)
5449 .unwrap();
5450 let fill = TestOrderEventStubs::filled(
5451 &order,
5452 &instrument,
5453 Some(TradeId::from("T-RECENT-FILL")),
5454 None,
5455 Some(Price::from("100.0")),
5456 Some(Quantity::from("1.0")),
5457 Some(LiquiditySide::Taker),
5458 None,
5459 None,
5460 Some(account_id),
5461 );
5462
5463 (node, fill, instrument)
5464 }
5465
5466 fn fill_report_event(fill: &OrderFilled) -> ExecutionEvent {
5467 ExecutionEvent::Report(ExecutionReport::Fill(Box::new(FillReport::new(
5468 fill.account_id,
5469 fill.instrument_id,
5470 fill.venue_order_id,
5471 fill.trade_id,
5472 fill.order_side,
5473 fill.last_qty,
5474 fill.last_px,
5475 fill.commission
5476 .unwrap_or_else(|| Money::zero(fill.currency)),
5477 fill.liquidity_side,
5478 Some(fill.client_order_id),
5479 fill.position_id,
5480 fill.ts_event,
5481 fill.ts_init,
5482 None,
5483 ))))
5484 }
5485
5486 fn is_recent_fill(node: &LiveNode, fill: &OrderFilled) -> bool {
5487 node.exec_manager.is_fill_recently_processed(
5488 fill.account_id,
5489 fill.instrument_id,
5490 fill.trade_id,
5491 )
5492 }
5493
5494 #[rstest]
5495 #[case(0, NodeState::Idle)]
5496 #[case(1, NodeState::Starting)]
5497 #[case(2, NodeState::Running)]
5498 #[case(3, NodeState::ShuttingDown)]
5499 #[case(4, NodeState::Stopped)]
5500 fn test_node_state_from_u8_valid(#[case] value: u8, #[case] expected: NodeState) {
5501 assert_eq!(NodeState::from_u8(value), expected);
5502 }
5503
5504 #[rstest]
5505 #[case(5)]
5506 #[case(255)]
5507 #[should_panic(expected = "Invalid NodeState value")]
5508 fn test_node_state_from_u8_invalid_panics(#[case] value: u8) {
5509 let _ = NodeState::from_u8(value);
5510 }
5511
5512 #[rstest]
5513 fn test_node_state_roundtrip() {
5514 for state in [
5515 NodeState::Idle,
5516 NodeState::Starting,
5517 NodeState::Running,
5518 NodeState::ShuttingDown,
5519 NodeState::Stopped,
5520 ] {
5521 assert_eq!(NodeState::from_u8(state.as_u8()), state);
5522 }
5523 }
5524
5525 #[rstest]
5526 fn test_node_state_is_running_only_for_running() {
5527 assert!(!NodeState::Idle.is_running());
5528 assert!(!NodeState::Starting.is_running());
5529 assert!(NodeState::Running.is_running());
5530 assert!(!NodeState::ShuttingDown.is_running());
5531 assert!(!NodeState::Stopped.is_running());
5532 }
5533
5534 #[rstest]
5535 #[tokio::test]
5536 async fn test_await_engines_connected_returns_stop_requested() {
5537 let node = LiveNode::build("TestNode".to_string(), None).unwrap();
5538 let handle = node.handle();
5539
5540 handle.stop();
5541
5542 let deadline = dst::time::Instant::now() + Duration::from_secs(1);
5543 let status = node.await_engines_connected(deadline).await;
5544
5545 assert_eq!(status, EngineConnectionStatus::StopRequested);
5546 assert!(handle.should_stop());
5547 }
5548
5549 #[rstest]
5550 #[tokio::test]
5551 async fn test_await_engines_connected_returns_shutdown_requested() {
5552 let node = LiveNode::build("TestNode".to_string(), None).unwrap();
5553
5554 node.kernel().shutdown_flag().set(true);
5555
5556 let deadline = dst::time::Instant::now() + Duration::from_secs(1);
5557 let status = node.await_engines_connected(deadline).await;
5558
5559 assert_eq!(status, EngineConnectionStatus::ShutdownRequested);
5560 }
5561
5562 #[rstest]
5563 #[tokio::test]
5564 async fn test_start_stop_request_aborts_startup_without_running() {
5565 let config = LiveNodeConfig {
5566 exec_engine: crate::config::LiveExecEngineConfig {
5567 reconciliation: false,
5568 ..Default::default()
5569 },
5570 timeout_disconnection: Duration::from_millis(50),
5571 ..Default::default()
5572 };
5573 let mut node = LiveNode::build("TestNode".to_string(), Some(config)).unwrap();
5574 let handle = node.handle();
5575
5576 handle.stop();
5577 node.start().await.unwrap();
5578
5579 assert_eq!(handle.state(), NodeState::Stopped);
5580 assert!(handle.should_stop());
5581 assert!(!handle.is_running());
5582 }
5583
5584 #[rstest]
5585 #[tokio::test(start_paused = true)]
5586 async fn test_stop_processes_residual_exec_event_during_grace_period() {
5587 let config = LiveNodeConfig {
5588 exec_engine: crate::config::LiveExecEngineConfig {
5589 reconciliation: false,
5590 ..Default::default()
5591 },
5592 timeout_connection: Duration::ZERO,
5593 timeout_reconciliation: Duration::ZERO,
5594 timeout_portfolio: Duration::ZERO,
5595 timeout_disconnection: Duration::ZERO,
5596 delay_post_stop: Duration::from_millis(20),
5597 timeout_shutdown: Duration::ZERO,
5598 ..Default::default()
5599 };
5600 let mut node = LiveNode::build("TestNode".to_string(), Some(config)).unwrap();
5601 let order = OrderTestBuilder::new(OrderType::Market)
5602 .instrument_id(InstrumentId::from("EUR/USD.SIM"))
5603 .quantity(Quantity::from("2"))
5604 .build();
5605 let client_order_id = order.client_order_id();
5606 let submitted = TestOrderEventStubs::submitted(&order, AccountId::from("POLL-STOP-001"));
5607
5608 node.kernel
5609 .cache()
5610 .borrow_mut()
5611 .add_order(order, None, None, false)
5612 .unwrap();
5613
5614 node.start().await.unwrap();
5615 let exec_event_sender = get_exec_event_sender();
5616
5617 let send_residual = tokio::spawn(async move {
5618 tokio::time::sleep(Duration::from_millis(1)).await;
5619 exec_event_sender
5620 .send(ExecutionEvent::Order(submitted))
5621 .unwrap();
5622 });
5623
5624 node.stop().await.unwrap();
5625 send_residual.await.unwrap();
5626
5627 assert_eq!(
5628 node.kernel
5629 .cache()
5630 .borrow()
5631 .order(&client_order_id)
5632 .unwrap()
5633 .status(),
5634 OrderStatus::Submitted
5635 );
5636
5637 node.dispose();
5638 }
5639
5640 #[rstest]
5641 #[tokio::test]
5642 async fn test_live_state_persistence_loads_before_start_and_saves_after_stop() {
5643 let actor_id = ActorId::from("LIVE-STATE-ACTOR");
5644 let strategy_id = StrategyId::from("LIVE-STATE-STRATEGY-001");
5645 let actor_load = IndexMap::from([("actor-load".to_string(), b"actor-loaded".to_vec())]);
5646 let strategy_load =
5647 IndexMap::from([("strategy-load".to_string(), b"strategy-loaded".to_vec())]);
5648 let actor_save = IndexMap::from([("actor-save".to_string(), b"actor-saved".to_vec())]);
5649 let strategy_save =
5650 IndexMap::from([("strategy-save".to_string(), b"strategy-saved".to_vec())]);
5651 let (database, control) = TestCacheDatabaseControl::create();
5652 control.set_actor_state(actor_id, &actor_load);
5653 control.set_strategy_state(strategy_id, &strategy_load);
5654 let config = LiveNodeConfig {
5655 load_state: true,
5656 save_state: true,
5657 exec_engine: crate::config::LiveExecEngineConfig {
5658 reconciliation: false,
5659 ..Default::default()
5660 },
5661 timeout_connection: Duration::ZERO,
5662 timeout_reconciliation: Duration::ZERO,
5663 timeout_portfolio: Duration::ZERO,
5664 timeout_disconnection: Duration::ZERO,
5665 delay_post_stop: Duration::ZERO,
5666 timeout_shutdown: Duration::ZERO,
5667 ..Default::default()
5668 };
5669 let mut node = LiveNode::build("StatePersistenceNode".to_string(), Some(config)).unwrap();
5670 node.set_cache_database(Box::new(database)).unwrap();
5671 node.add_actor(StateActor::new(
5672 actor_id,
5673 control.clone(),
5674 actor_save.clone(),
5675 ))
5676 .unwrap();
5677 node.add_strategy(StateStrategy::new(
5678 strategy_id,
5679 control.clone(),
5680 strategy_save.clone(),
5681 ))
5682 .unwrap();
5683
5684 node.start().await.unwrap();
5685 node.stop().await.unwrap();
5686 node.dispose();
5687
5688 assert_eq!(
5689 control.events(),
5690 vec![
5691 "actor.load:LIVE-STATE-ACTOR",
5692 "actor.on_load",
5693 "strategy.load:LIVE-STATE-STRATEGY-001",
5694 "strategy.on_load",
5695 "actor.on_start",
5696 "strategy.on_start",
5697 "actor.on_stop",
5698 "strategy.on_stop",
5699 "actor.on_save",
5700 "actor.update:LIVE-STATE-ACTOR",
5701 "strategy.on_save",
5702 "strategy.update:LIVE-STATE-STRATEGY-001",
5703 "database.close",
5704 ]
5705 );
5706 assert_eq!(control.actor_state(&actor_id), Some(actor_save));
5707 assert_eq!(control.strategy_state(&strategy_id), Some(strategy_save));
5708 assert_eq!(node.state(), NodeState::Stopped);
5709 }
5710
5711 #[rstest]
5712 #[tokio::test]
5713 async fn test_live_state_persistence_reports_callback_errors_after_shutdown() {
5714 let actor_id = ActorId::from("LIVE-FAIL-SAVE-ACTOR");
5715 let strategy_id = StrategyId::from("LIVE-FAIL-SAVE-STRATEGY-001");
5716 let (database, control) = TestCacheDatabaseControl::create();
5717 let config = LiveNodeConfig {
5718 save_state: true,
5719 exec_engine: crate::config::LiveExecEngineConfig {
5720 reconciliation: false,
5721 ..Default::default()
5722 },
5723 timeout_connection: Duration::ZERO,
5724 timeout_reconciliation: Duration::ZERO,
5725 timeout_portfolio: Duration::ZERO,
5726 timeout_disconnection: Duration::ZERO,
5727 delay_post_stop: Duration::ZERO,
5728 timeout_shutdown: Duration::ZERO,
5729 ..Default::default()
5730 };
5731 let mut node =
5732 LiveNode::build("StatePersistenceErrorNode".to_string(), Some(config)).unwrap();
5733 node.set_cache_database(Box::new(database)).unwrap();
5734 node.add_actor(
5735 StateActor::new(actor_id, control.clone(), IndexMap::new()).with_fail_save(),
5736 )
5737 .unwrap();
5738 node.add_strategy(
5739 StateStrategy::new(strategy_id, control.clone(), IndexMap::new()).with_fail_save(),
5740 )
5741 .unwrap();
5742
5743 node.start().await.unwrap();
5744 let error = node.stop().await.unwrap_err();
5745 node.dispose();
5746
5747 assert_eq!(
5748 error.to_string(),
5749 "failed while finalizing kernel shutdown: Failed to save component state: actor \
5750 LIVE-FAIL-SAVE-ACTOR callback: test actor on_save failure; strategy \
5751 LIVE-FAIL-SAVE-STRATEGY-001 callback: test strategy on_save failure"
5752 );
5753 assert_eq!(
5754 control.events(),
5755 vec![
5756 "actor.on_start",
5757 "strategy.on_start",
5758 "actor.on_stop",
5759 "strategy.on_stop",
5760 "actor.on_save",
5761 "strategy.on_save",
5762 "database.close",
5763 ]
5764 );
5765 assert_eq!(node.state(), NodeState::Stopped);
5766 }
5767
5768 #[rstest]
5769 #[tokio::test]
5770 async fn test_stop_drains_queued_exec_event_after_zero_grace() {
5771 let config = LiveNodeConfig {
5772 exec_engine: crate::config::LiveExecEngineConfig {
5773 reconciliation: false,
5774 ..Default::default()
5775 },
5776 timeout_connection: Duration::ZERO,
5777 timeout_reconciliation: Duration::ZERO,
5778 timeout_portfolio: Duration::ZERO,
5779 timeout_disconnection: Duration::ZERO,
5780 delay_post_stop: Duration::ZERO,
5781 timeout_shutdown: Duration::ZERO,
5782 ..Default::default()
5783 };
5784 let mut node = LiveNode::build("TestNode".to_string(), Some(config)).unwrap();
5785 let order = OrderTestBuilder::new(OrderType::Market)
5786 .instrument_id(InstrumentId::from("GBP/USD.SIM"))
5787 .quantity(Quantity::from("3"))
5788 .build();
5789 let client_order_id = order.client_order_id();
5790 let submitted = TestOrderEventStubs::submitted(&order, AccountId::from("POLL-DRAIN-001"));
5791
5792 node.kernel
5793 .cache()
5794 .borrow_mut()
5795 .add_order(order, None, None, false)
5796 .unwrap();
5797
5798 node.start().await.unwrap();
5799 get_exec_event_sender()
5800 .send(ExecutionEvent::Order(submitted))
5801 .unwrap();
5802
5803 node.stop().await.unwrap();
5804
5805 assert_eq!(
5806 node.kernel
5807 .cache()
5808 .borrow()
5809 .order(&client_order_id)
5810 .unwrap()
5811 .status(),
5812 OrderStatus::Submitted
5813 );
5814
5815 node.dispose();
5816 }
5817
5818 #[rstest]
5819 #[tokio::test]
5820 async fn test_start_event_store_replay_skips_live_connections() {
5821 let mut node = live_node_with_replay_store(false);
5822 let handle = node.handle();
5823
5824 node.start().await.unwrap();
5825
5826 assert_eq!(handle.state(), NodeState::Running);
5827 assert!(handle.is_running());
5828 assert!(node.kernel.is_event_store_replay());
5829 assert!(node.runner.is_some());
5830 }
5831
5832 #[rstest]
5833 #[tokio::test]
5834 async fn test_start_event_store_replay_preserves_stop_request() {
5835 let mut node = live_node_with_replay_store(false);
5836 let handle = node.handle();
5837 handle.stop();
5838
5839 node.start().await.unwrap();
5840
5841 assert_eq!(handle.state(), NodeState::Stopped);
5842 assert!(handle.should_stop());
5843 assert!(node.kernel.is_event_store_replay());
5844 assert!(node.runner.is_some());
5845 }
5846
5847 #[rstest]
5848 #[tokio::test]
5849 async fn test_start_event_store_replay_config_failure_aborts_startup() {
5850 let mut node = live_node_with_replay_store(true);
5851 let handle = node.handle();
5852
5853 node.start().await.unwrap();
5854
5855 assert_eq!(handle.state(), NodeState::Stopped);
5856 assert!(!handle.is_running());
5857 assert!(node.kernel.is_event_store_replay_configured());
5858 assert!(!node.kernel.is_event_store_replay());
5859 assert!(node.runner.is_some());
5860 }
5861
5862 #[rstest]
5863 #[tokio::test]
5864 async fn test_run_event_store_replay_consumes_runner_and_stops_before_connections() {
5865 let mut node = live_node_with_replay_store(false);
5866 let handle = node.handle();
5867
5868 node.run().await.unwrap();
5869
5870 assert_eq!(handle.state(), NodeState::Running);
5871 assert!(handle.is_running());
5872 assert!(node.kernel.is_event_store_replay());
5873 assert!(node.runner.is_none());
5874 }
5875
5876 #[rstest]
5877 #[tokio::test]
5878 async fn test_run_event_store_replay_preserves_stop_request() {
5879 let mut node = live_node_with_replay_store(false);
5880 let handle = node.handle();
5881 handle.stop();
5882
5883 node.run().await.unwrap();
5884
5885 assert_eq!(handle.state(), NodeState::Stopped);
5886 assert!(handle.should_stop());
5887 assert!(node.kernel.is_event_store_replay());
5888 assert!(node.runner.is_none());
5889 }
5890
5891 #[rstest]
5892 #[tokio::test]
5893 async fn test_run_event_store_replay_config_failure_aborts_startup() {
5894 let mut node = live_node_with_replay_store(true);
5895 let handle = node.handle();
5896
5897 node.run().await.unwrap();
5898
5899 assert_eq!(handle.state(), NodeState::Stopped);
5900 assert!(!handle.is_running());
5901 assert!(node.kernel.is_event_store_replay_configured());
5902 assert!(!node.kernel.is_event_store_replay());
5903 assert!(node.runner.is_none());
5904 }
5905
5906 #[rstest]
5907 fn test_build_rejects_event_store_config_without_factory() {
5908 let config = LiveNodeConfig {
5909 event_store: Some(EventStoreConfig::default()),
5910 exec_engine: crate::config::LiveExecEngineConfig {
5911 reconciliation: false,
5912 ..Default::default()
5913 },
5914 ..Default::default()
5915 };
5916
5917 let err = LiveNodeBuilder::from_config(config)
5918 .expect("builder")
5919 .build()
5920 .expect_err("should reject event_store config without factory");
5921
5922 assert!(
5923 err.to_string().contains("with_event_store"),
5924 "error message should mention with_event_store, was: {err}"
5925 );
5926 }
5927
5928 #[rstest]
5929 fn test_direct_build_rejects_event_store_config() {
5930 let config = LiveNodeConfig {
5931 event_store: Some(EventStoreConfig::default()),
5932 exec_engine: crate::config::LiveExecEngineConfig {
5933 reconciliation: false,
5934 ..Default::default()
5935 },
5936 ..Default::default()
5937 };
5938
5939 let err = LiveNode::build("TestNode".to_string(), Some(config))
5940 .expect_err("LiveNode::build should reject event_store config");
5941
5942 assert!(
5943 err.to_string().contains("with_event_store"),
5944 "error message should mention with_event_store, was: {err}"
5945 );
5946 }
5947
5948 #[rstest]
5949 fn test_dispose_before_start_is_idempotent() {
5950 let mut node = LiveNode::build("TestNode".to_string(), None).unwrap();
5951 node.add_strategy(TestStrategy::new(StrategyConfig {
5952 strategy_id: Some(StrategyId::from("DISPOSAL-001")),
5953 ..Default::default()
5954 }))
5955 .unwrap();
5956
5957 node.dispose();
5958 node.dispose();
5959
5960 assert!(node.kernel.trader().borrow().is_disposed());
5961 assert_eq!(node.kernel.trader().borrow().component_count(), 0);
5962 assert_eq!(node.state(), NodeState::Stopped);
5963 }
5964
5965 #[rstest]
5966 fn test_handle_initial_state() {
5967 let handle = LiveNodeHandle::new();
5968
5969 assert_eq!(handle.state(), NodeState::Idle);
5970 assert!(!handle.should_stop());
5971 assert!(!handle.is_running());
5972 }
5973
5974 #[rstest]
5975 fn test_handle_initial_metrics_snapshot_is_zero() {
5976 let handle = LiveNodeHandle::new();
5977
5978 assert_eq!(handle.metrics_snapshot(), RunnerMetricsSnapshot::default());
5979 }
5980
5981 #[rstest]
5982 fn test_record_runner_dispatch_updates_selected_channel() {
5983 let metrics = RunnerMetrics::default();
5984 let dispatch_start = dst::time::Instant::now();
5985 let metrics_start = dispatch_start
5986 .checked_sub(Duration::from_micros(1))
5987 .expect("test instant should support a one-microsecond lookback");
5988
5989 record_runner_dispatch(
5990 &metrics,
5991 SystemChannel::DataCommands,
5992 dispatch_start,
5993 metrics_start,
5994 );
5995 let snapshot = metrics.snapshot();
5996
5997 assert_eq!(snapshot.time_events.dispatched, 0);
5998 assert_eq!(snapshot.exec_events.dispatched, 0);
5999 assert_eq!(snapshot.exec_commands.dispatched, 0);
6000 assert_eq!(snapshot.data_events.dispatched, 0);
6001 assert_eq!(snapshot.data_commands.dispatched, 1);
6002 assert_eq!(
6003 snapshot.data_commands.last_dispatch_at_ns,
6004 snapshot.elapsed_ns
6005 );
6006 assert_eq!(
6007 snapshot.data_commands.dispatch_busy_ns,
6008 snapshot.dispatch_busy_ns
6009 );
6010 assert!(snapshot.dispatch_busy_ns < snapshot.elapsed_ns);
6011 }
6012
6013 #[rstest]
6014 fn test_handle_stop_sets_flag() {
6015 let handle = LiveNodeHandle::new();
6016
6017 handle.stop();
6018
6019 assert!(handle.should_stop());
6020 }
6021
6022 #[rstest]
6023 fn test_handle_stop_blocks_running_transition() {
6024 let handle = LiveNodeHandle::new();
6025 handle.set_starting();
6026 handle.stop();
6027
6028 let transition = handle.try_set_running();
6029
6030 assert_eq!(transition, RunningTransition::StopRequested);
6031 assert_eq!(handle.state(), NodeState::Starting);
6032 assert!(handle.should_stop());
6033 assert!(!handle.is_running());
6034 }
6035
6036 #[rstest]
6037 fn test_handle_stop_after_running_transition_remains_pending() {
6038 let handle = LiveNodeHandle::new();
6039 handle.set_starting();
6040
6041 let transition = handle.try_set_running();
6042 handle.stop();
6043
6044 assert_eq!(transition, RunningTransition::Entered);
6045 assert_eq!(handle.state(), NodeState::Running);
6046 assert!(handle.should_stop());
6047 assert!(handle.is_running());
6048 }
6049
6050 #[rstest]
6051 fn test_handle_node_state_transitions() {
6052 let handle = LiveNodeHandle::new();
6053 assert_eq!(handle.state(), NodeState::Idle);
6054
6055 handle.set_starting();
6056 assert_eq!(handle.state(), NodeState::Starting);
6057 assert!(!handle.is_running());
6058
6059 assert_eq!(handle.try_set_running(), RunningTransition::Entered);
6060 assert_eq!(handle.state(), NodeState::Running);
6061 assert!(handle.is_running());
6062
6063 handle.set_shutting_down();
6064 assert_eq!(handle.state(), NodeState::ShuttingDown);
6065 assert!(!handle.is_running());
6066
6067 handle.set_stopped();
6068 assert_eq!(handle.state(), NodeState::Stopped);
6069 assert!(!handle.is_running());
6070 }
6071
6072 #[rstest]
6073 fn test_handle_clone_shares_state_bidirectionally() {
6074 let handle1 = LiveNodeHandle::new();
6075 let handle2 = handle1.clone();
6076
6077 handle1.set_starting();
6078 let transition = handle2.try_set_running();
6079 handle1.stop();
6080
6081 assert_eq!(transition, RunningTransition::Entered);
6082 assert_eq!(handle1.state(), NodeState::Running);
6083 assert!(handle2.should_stop());
6084 }
6085
6086 #[rstest]
6087 fn test_handle_stop_flag_survives_non_running_state_changes() {
6088 let handle = LiveNodeHandle::new();
6089
6090 handle.set_starting();
6091 handle.stop();
6092 handle.set_shutting_down();
6093 handle.set_stopped();
6094
6095 assert_eq!(handle.state(), NodeState::Stopped);
6096 assert!(handle.should_stop());
6097 }
6098
6099 #[rstest]
6100 fn test_builder_creation() {
6101 let result = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox);
6102
6103 assert!(result.is_ok());
6104 }
6105
6106 #[rstest]
6107 fn test_builder_rejects_backtest() {
6108 let result = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Backtest);
6109
6110 assert!(result.is_err());
6111 assert!(result.unwrap_err().to_string().contains("Backtest"));
6112 }
6113
6114 #[rstest]
6115 fn test_builder_accepts_live_environment() {
6116 let result = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Live);
6117
6118 assert!(result.is_ok());
6119 }
6120
6121 #[rstest]
6122 fn test_builder_accepts_sandbox_environment() {
6123 let result = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox);
6124
6125 assert!(result.is_ok());
6126 }
6127
6128 #[rstest]
6129 fn test_builder_fluent_api_chaining() {
6130 let builder = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Live)
6131 .unwrap()
6132 .with_name("TestNode")
6133 .with_instance_id(UUID4::new())
6134 .with_load_state(false)
6135 .with_save_state(true)
6136 .with_timeout_connection(30)
6137 .with_timeout_reconciliation(60)
6138 .with_reconciliation(true)
6139 .with_reconciliation_lookback_mins(120)
6140 .with_timeout_portfolio(10)
6141 .with_timeout_disconnection_secs(5)
6142 .with_delay_post_stop_secs(3)
6143 .with_delay_shutdown_secs(10);
6144
6145 assert_eq!(builder.name(), "TestNode");
6146 }
6147
6148 #[rstest]
6149 fn test_builder_with_external_msgbus_egress_uses_configured_encoding() {
6150 let (external_egress, publications, closed) = CapturingExternalEgress::new();
6151 let msgbus_config = MessageBusConfig {
6152 encoding: SerializationEncoding::Json,
6153 ..Default::default()
6154 };
6155 let node = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox)
6156 .unwrap()
6157 .with_msgbus_config(msgbus_config)
6158 .with_external_msgbus_egress(Box::new(external_egress))
6159 .build()
6160 .expect("node builds with external message bus egress");
6161 let quote = QuoteTick::default();
6162
6163 msgbus::publish_quote("data.quotes.TEST".into(), "e);
6164
6165 let publications = publications.borrow();
6166 assert_eq!(publications.len(), 1);
6167 assert_eq!(publications[0].topic, "data.quotes.TEST");
6168 assert_eq!(
6169 serde_json::from_slice::<QuoteTick>(&publications[0].payload)
6170 .expect("JSON payload must decode as QuoteTick"),
6171 quote
6172 );
6173 drop(publications);
6174
6175 msgbus::get_message_bus().borrow_mut().dispose();
6176 assert!(closed.get());
6177 drop(node);
6178 }
6179
6180 #[rstest]
6181 #[tokio::test(flavor = "current_thread")]
6182 async fn test_builder_with_external_msgbus_factory_installs_egress_and_ingress() {
6183 let quote = QuoteTick::default();
6184 let (tx, rx) = tokio::sync::mpsc::channel::<BusMessage>(1);
6185 let publications = Arc::new(Mutex::new(Vec::new()));
6186 let closed = Arc::new(AtomicBool::new(false));
6187 let factory = CapturingBackingFactory::new(publications.clone(), closed.clone(), Some(rx));
6188 let msgbus_config = MessageBusConfig {
6189 external_streams: Some(vec!["stream".to_string()]),
6190 ..Default::default()
6191 };
6192 let config = LiveNodeConfig {
6193 environment: Environment::Sandbox,
6194 msgbus: Some(msgbus_config),
6195 exec_engine: crate::config::LiveExecEngineConfig {
6196 reconciliation: false,
6197 ..Default::default()
6198 },
6199 delay_post_stop: Duration::ZERO,
6200 timeout_connection: Duration::from_millis(500),
6201 timeout_disconnection: Duration::from_millis(500),
6202 ..Default::default()
6203 };
6204 let mut node = LiveNodeBuilder::from_config(config)
6205 .unwrap()
6206 .with_external_msgbus_factory(Box::new(factory))
6207 .build()
6208 .expect("node builds with external message bus factory");
6209
6210 msgbus::publish_quote("data.quotes.TEST".into(), "e);
6211 {
6212 let publications = publications.lock().unwrap();
6213 assert_eq!(publications.len(), 1);
6214 assert_eq!(publications[0].topic, "data.quotes.TEST");
6215 assert_eq!(
6216 serde_json::from_slice::<QuoteTick>(&publications[0].payload)
6217 .expect("JSON payload must decode as QuoteTick"),
6218 quote
6219 );
6220 }
6221
6222 let received = Rc::new(RefCell::new(Vec::<QuoteTick>::new()));
6223 let handle = node.handle();
6224 let handler = TypedHandler::from({
6227 let received = received.clone();
6228 move |quote: &QuoteTick| {
6229 received.borrow_mut().push(*quote);
6230 }
6231 });
6232 msgbus::subscribe_quotes("data.quotes.*".into(), handler, None);
6233 msgbus::get_message_bus()
6234 .borrow_mut()
6235 .add_streaming_type(BusPayloadType::QuoteTick);
6236
6237 let payload =
6238 Bytes::from(serde_json::to_vec("e).expect("QuoteTick should serialize as JSON"));
6239 let message = BusMessage::with_str_topic(
6240 "data.quotes.TEST",
6241 BusPayloadType::QuoteTick,
6242 payload,
6243 SerializationEncoding::Json,
6244 );
6245
6246 tokio::time::timeout(Duration::from_secs(30), async {
6247 let run = node.run();
6248 tokio::pin!(run);
6249
6250 let drive = async {
6251 wait_until_async(|| async { handle.is_running() }, Duration::from_secs(10)).await;
6252
6253 tx.send(message)
6254 .await
6255 .expect("external ingress receiver should be open");
6256
6257 wait_until_async(
6258 || async { received.borrow().len() == 1 },
6259 Duration::from_secs(10),
6260 )
6261 .await;
6262 assert_eq!(*received.borrow(), vec![quote]);
6263 handle.stop();
6264 };
6265
6266 tokio::select! {
6267 biased;
6268
6269 () = drive => {}
6270 result = &mut run => {
6271 panic!("node stopped before factory ingress was republished: {result:?}");
6272 }
6273 }
6274
6275 run.await.expect("node should stop cleanly");
6276 })
6277 .await
6278 .expect("live node should republish factory ingress and stop before timeout");
6279
6280 assert_eq!(handle.state(), NodeState::Stopped);
6281 assert!(closed.load(Ordering::Relaxed));
6282 msgbus::get_message_bus().borrow_mut().dispose();
6283 }
6284
6285 #[rstest]
6286 #[tokio::test(flavor = "current_thread")]
6287 async fn test_builder_with_external_msgbus_factory_without_streams_runs_without_ingress() {
6288 let quote = QuoteTick::default();
6289 let publications = Arc::new(Mutex::new(Vec::new()));
6290 let closed = Arc::new(AtomicBool::new(false));
6291 let factory = CapturingBackingFactory::new(publications.clone(), closed.clone(), None);
6292 let config = LiveNodeConfig {
6293 environment: Environment::Sandbox,
6294 msgbus: Some(MessageBusConfig::default()),
6295 exec_engine: crate::config::LiveExecEngineConfig {
6296 reconciliation: false,
6297 ..Default::default()
6298 },
6299 delay_post_stop: Duration::ZERO,
6300 timeout_connection: Duration::from_millis(500),
6301 timeout_disconnection: Duration::from_millis(500),
6302 ..Default::default()
6303 };
6304 let mut node = LiveNodeBuilder::from_config(config)
6305 .unwrap()
6306 .with_external_msgbus_factory(Box::new(factory))
6307 .build()
6308 .expect("node builds with egress-only message bus factory");
6309 let handle = node.handle();
6310
6311 msgbus::publish_quote("data.quotes.TEST".into(), "e);
6312 {
6313 let publications = publications.lock().unwrap();
6314 assert_eq!(publications.len(), 1);
6315 assert_eq!(publications[0].topic, "data.quotes.TEST");
6316 }
6317
6318 tokio::time::timeout(Duration::from_secs(30), async {
6319 let run = node.run();
6320 tokio::pin!(run);
6321
6322 let drive = async {
6323 wait_until_async(|| async { handle.is_running() }, Duration::from_secs(10)).await;
6324 handle.stop();
6325 };
6326
6327 tokio::select! {
6328 biased;
6329
6330 () = drive => {}
6331 result = &mut run => {
6332 panic!("node stopped before egress-only factory run was observed: {result:?}");
6333 }
6334 }
6335
6336 run.await.expect("node should stop cleanly");
6337 })
6338 .await
6339 .expect("live node should run without external ingress before timeout");
6340
6341 assert_eq!(handle.state(), NodeState::Stopped);
6342 msgbus::get_message_bus().borrow_mut().dispose();
6343 assert!(closed.load(Ordering::Relaxed));
6344 }
6345
6346 #[rstest]
6347 fn test_builder_with_external_msgbus_factory_rejects_injected_surfaces() {
6348 let (external_egress, _publications, _closed) = CapturingExternalEgress::new();
6349 let egress_factory = CapturingBackingFactory::new(
6350 Arc::new(Mutex::new(Vec::new())),
6351 Arc::new(AtomicBool::new(false)),
6352 None,
6353 );
6354 let egress_error = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox)
6355 .unwrap()
6356 .with_external_msgbus_factory(Box::new(egress_factory))
6357 .with_external_msgbus_egress(Box::new(external_egress))
6358 .build()
6359 .expect_err("builder should reject factory plus injected egress");
6360
6361 assert!(
6362 egress_error
6363 .to_string()
6364 .contains("cannot be combined with injected egress or ingress")
6365 );
6366
6367 let (_tx, rx) = tokio::sync::mpsc::channel::<BusMessage>(1);
6368 let ingress_factory = CapturingBackingFactory::new(
6369 Arc::new(Mutex::new(Vec::new())),
6370 Arc::new(AtomicBool::new(false)),
6371 None,
6372 );
6373 let ingress = CapturingExternalIngress::new(rx, Rc::new(Cell::new(false)));
6374 let ingress_error = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox)
6375 .unwrap()
6376 .with_external_msgbus_factory(Box::new(ingress_factory))
6377 .with_external_ingress(Box::new(ingress))
6378 .build()
6379 .expect_err("builder should reject factory plus injected ingress");
6380
6381 assert!(
6382 ingress_error
6383 .to_string()
6384 .contains("cannot be combined with injected egress or ingress")
6385 );
6386 }
6387
6388 #[rstest]
6389 #[tokio::test(flavor = "current_thread")]
6390 async fn test_run_republishes_external_ingress_on_local_msgbus() {
6391 let quote = QuoteTick::default();
6392 let received = Rc::new(RefCell::new(Vec::<QuoteTick>::new()));
6393 let payload =
6394 Bytes::from(serde_json::to_vec("e).expect("QuoteTick should serialize as JSON"));
6395 let message = BusMessage::with_str_topic(
6396 "data.quotes.TEST",
6397 BusPayloadType::QuoteTick,
6398 payload,
6399 SerializationEncoding::Json,
6400 );
6401 let (tx, rx) = tokio::sync::mpsc::channel::<BusMessage>(1);
6402 let closed = Rc::new(Cell::new(false));
6403 let ingress = CapturingExternalIngress::new(rx, closed.clone());
6404 let config = LiveNodeConfig {
6405 environment: Environment::Sandbox,
6406 exec_engine: crate::config::LiveExecEngineConfig {
6407 reconciliation: false,
6408 ..Default::default()
6409 },
6410 delay_post_stop: Duration::ZERO,
6411 timeout_connection: Duration::from_millis(500),
6412 timeout_disconnection: Duration::from_millis(500),
6413 ..Default::default()
6414 };
6415 let mut node = LiveNodeBuilder::from_config(config)
6416 .unwrap()
6417 .with_external_ingress(Box::new(ingress))
6418 .build()
6419 .expect("node builds with external message bus ingress");
6420 let handle = node.handle();
6421 let handler = TypedHandler::from({
6422 let received = received.clone();
6423 move |quote: &QuoteTick| {
6424 received.borrow_mut().push(*quote);
6425 }
6426 });
6427 msgbus::subscribe_quotes("data.quotes.*".into(), handler, None);
6428 msgbus::get_message_bus()
6429 .borrow_mut()
6430 .add_streaming_type(BusPayloadType::QuoteTick);
6431
6432 tokio::time::timeout(Duration::from_secs(30), async {
6433 let run = node.run();
6434 tokio::pin!(run);
6435
6436 let drive = async {
6437 wait_until_async(|| async { handle.is_running() }, Duration::from_secs(10)).await;
6438
6439 tx.send(message)
6440 .await
6441 .expect("external ingress receiver should be open");
6442
6443 wait_until_async(
6444 || async { received.borrow().len() == 1 },
6445 Duration::from_secs(10),
6446 )
6447 .await;
6448 assert_eq!(*received.borrow(), vec![quote]);
6449 handle.stop();
6450 };
6451
6452 tokio::select! {
6453 biased;
6454
6455 () = drive => {}
6456 result = &mut run => {
6457 panic!("node stopped before external message was republished: {result:?}");
6458 }
6459 }
6460
6461 run.await.expect("node should stop cleanly");
6462 })
6463 .await
6464 .expect("live node should republish ingress and stop before timeout");
6465
6466 assert_eq!(handle.state(), NodeState::Stopped);
6467 assert!(closed.get());
6468 msgbus::get_message_bus().borrow_mut().dispose();
6469 }
6470
6471 #[rstest]
6472 #[tokio::test(flavor = "current_thread")]
6473 async fn test_run_closes_external_ingress_when_receiver_closes() {
6474 let (tx, rx) = tokio::sync::mpsc::channel::<BusMessage>(1);
6475 let closed = Rc::new(Cell::new(false));
6476 let ingress = CapturingExternalIngress::new(rx, closed.clone());
6477 let config = LiveNodeConfig {
6478 environment: Environment::Sandbox,
6479 exec_engine: crate::config::LiveExecEngineConfig {
6480 reconciliation: false,
6481 ..Default::default()
6482 },
6483 delay_post_stop: Duration::ZERO,
6484 timeout_connection: Duration::from_millis(500),
6485 timeout_disconnection: Duration::from_millis(500),
6486 ..Default::default()
6487 };
6488 let mut node = LiveNodeBuilder::from_config(config)
6489 .unwrap()
6490 .with_external_ingress(Box::new(ingress))
6491 .build()
6492 .expect("node builds with external message bus ingress");
6493 let handle = node.handle();
6494
6495 tokio::time::timeout(Duration::from_secs(30), async {
6496 let run = node.run();
6497 tokio::pin!(run);
6498
6499 let drive = async {
6500 wait_until_async(|| async { handle.is_running() }, Duration::from_secs(10)).await;
6501
6502 drop(tx);
6503
6504 wait_until_async(|| async { closed.get() }, Duration::from_secs(10)).await;
6505 assert!(
6506 handle.is_running(),
6507 "node should keep running after ingress closes"
6508 );
6509 handle.stop();
6510 };
6511
6512 tokio::select! {
6513 biased;
6514
6515 () = drive => {}
6516 result = &mut run => {
6517 panic!("node stopped before ingress close was observed: {result:?}");
6518 }
6519 }
6520
6521 run.await.expect("node should stop cleanly");
6522 })
6523 .await
6524 .expect("live node should close ingress and stop before timeout");
6525
6526 assert_eq!(handle.state(), NodeState::Stopped);
6527 }
6528
6529 #[rstest]
6530 #[tokio::test(flavor = "current_thread")]
6531 async fn test_run_aborts_startup_when_external_ingress_receiver_unavailable() {
6532 let closed = Rc::new(Cell::new(false));
6533 let ingress = FailingExternalIngress::new(closed.clone());
6534 let config = LiveNodeConfig {
6535 environment: Environment::Sandbox,
6536 exec_engine: crate::config::LiveExecEngineConfig {
6537 reconciliation: false,
6538 ..Default::default()
6539 },
6540 delay_post_stop: Duration::ZERO,
6541 timeout_connection: Duration::from_millis(500),
6542 timeout_disconnection: Duration::from_millis(500),
6543 ..Default::default()
6544 };
6545 let mut node = LiveNodeBuilder::from_config(config)
6546 .unwrap()
6547 .with_external_ingress(Box::new(ingress))
6548 .build()
6549 .expect("node builds with external message bus ingress");
6550 let handle = node.handle();
6551
6552 let err = node.run().await.expect_err("run should fail");
6553
6554 assert!(
6555 err.to_string()
6556 .contains("external ingress receiver unavailable")
6557 );
6558 assert_eq!(handle.state(), NodeState::Stopped);
6559 assert!(closed.get());
6560 }
6561
6562 #[cfg(feature = "python")]
6563 #[rstest]
6564 fn test_node_build_and_initial_state() {
6565 let node = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox)
6566 .unwrap()
6567 .with_name("TestNode")
6568 .build()
6569 .unwrap();
6570
6571 assert_eq!(node.state(), NodeState::Idle);
6572 assert!(!node.is_running());
6573 assert_eq!(node.environment(), Environment::Sandbox);
6574 assert_eq!(node.trader_id(), TraderId::from("TRADER-001"));
6575 }
6576
6577 #[cfg(feature = "python")]
6578 #[rstest]
6579 fn test_node_build_replaces_stale_runner_senders() {
6580 replace_data_cmd_sender(Arc::new(SyncDataCommandSender));
6581 replace_exec_cmd_sender(Arc::new(SyncTradingCommandSender));
6582
6583 let first = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox)
6584 .unwrap()
6585 .with_name("FirstNode")
6586 .build()
6587 .unwrap();
6588
6589 assert_eq!(first.state(), NodeState::Idle);
6590 drop(first);
6591
6592 let second = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox)
6593 .unwrap()
6594 .with_name("SecondNode")
6595 .build()
6596 .unwrap();
6597
6598 assert_eq!(second.state(), NodeState::Idle);
6599 assert!(!second.is_running());
6600 }
6601
6602 #[cfg(feature = "python")]
6603 #[rstest]
6604 fn test_node_handle_reflects_node_state() {
6605 let node = LiveNode::builder(TraderId::from("TRADER-001"), Environment::Sandbox)
6606 .unwrap()
6607 .with_name("TestNode")
6608 .build()
6609 .unwrap();
6610
6611 let handle = node.handle();
6612
6613 assert_eq!(handle.state(), NodeState::Idle);
6614 assert!(!handle.is_running());
6615 }
6616
6617 #[rstest]
6618 fn test_pending_drain_data_returns_false_when_empty() {
6619 let mut pending = PendingEvents::default();
6620
6621 assert!(!pending.drain_data());
6622 }
6623
6624 #[rstest]
6625 fn test_pending_drain_data_returns_true_when_non_empty() {
6626 use nautilus_model::instruments::{InstrumentAny, stubs::crypto_perpetual_ethusdt};
6627
6628 let mut pending = PendingEvents::default();
6629 pending
6630 .data_evts
6631 .push(DataEvent::Instrument(InstrumentAny::CryptoPerpetual(
6632 crypto_perpetual_ethusdt(),
6633 )));
6634
6635 assert!(pending.drain_data());
6636 assert!(pending.data_evts.is_empty());
6637 }
6638
6639 fn stub_data_event() -> DataEvent {
6640 use nautilus_model::instruments::{InstrumentAny, stubs::crypto_perpetual_ethusdt};
6641
6642 DataEvent::Instrument(InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt()))
6643 }
6644
6645 fn stub_data_command() -> DataCommand {
6646 use nautilus_common::messages::data::{SubscribeCommand, subscribe::SubscribeInstruments};
6647 use nautilus_core::{UUID4, UnixNanos};
6648 use nautilus_model::identifiers::Venue;
6649
6650 DataCommand::Subscribe(SubscribeCommand::Instruments(SubscribeInstruments::new(
6651 None,
6652 Venue::from("TEST"),
6653 UUID4::new(),
6654 UnixNanos::default(),
6655 None,
6656 None,
6657 )))
6658 }
6659
6660 fn stub_system_command() -> SystemCommand {
6661 SystemCommand::ReconnectSocket(ReconnectSocket::new(
6662 TraderId::from("TRADER-001"),
6663 ClientId::from("POLYMARKET"),
6664 Ustr::from("polymarket-market-streams"),
6665 UnixNanos::default(),
6666 ))
6667 }
6668
6669 #[rstest]
6670 fn test_flush_pending_data_drains_events_and_commands() {
6671 let (evt_tx, mut evt_rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
6672 let (cmd_tx, mut cmd_rx) = tokio::sync::mpsc::unbounded_channel::<DataCommand>();
6673
6674 let mut pending = PendingEvents::default();
6675
6676 pending.data_evts.push(stub_data_event());
6678 pending.data_cmds.push(stub_data_command());
6679
6680 evt_tx.send(stub_data_event()).unwrap();
6682 cmd_tx.send(stub_data_command()).unwrap();
6683
6684 flush_pending_data(&mut pending, &mut evt_rx, &mut cmd_rx);
6685
6686 assert!(pending.data_evts.is_empty());
6687 assert!(pending.data_cmds.is_empty());
6688 assert!(evt_rx.try_recv().is_err());
6689 assert!(cmd_rx.try_recv().is_err());
6690 }
6691
6692 #[rstest]
6693 fn test_flush_pending_data_drains_mixed_sources() {
6694 let (evt_tx, mut evt_rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
6695 let (cmd_tx, mut cmd_rx) = tokio::sync::mpsc::unbounded_channel::<DataCommand>();
6696
6697 let mut pending = PendingEvents::default();
6698
6699 pending.data_evts.push(stub_data_event());
6701 cmd_tx.send(stub_data_command()).unwrap();
6702
6703 evt_tx.send(stub_data_event()).unwrap();
6705 evt_tx.send(stub_data_event()).unwrap();
6706 cmd_tx.send(stub_data_command()).unwrap();
6707
6708 flush_pending_data(&mut pending, &mut evt_rx, &mut cmd_rx);
6709
6710 assert!(pending.data_evts.is_empty());
6711 assert!(pending.data_cmds.is_empty());
6712 assert!(evt_rx.try_recv().is_err());
6713 assert!(cmd_rx.try_recv().is_err());
6714 }
6715
6716 #[rstest]
6717 fn test_pending_system_events_stay_separate_from_data() {
6718 let mut pending = PendingEvents::default();
6719 let change = SocketStateChange::new(
6720 ClientId::from("BINANCE"),
6721 Some(Venue::from("BINANCE")),
6722 ustr::Ustr::from("binance-futures-market-streams"),
6723 SocketState::Connected,
6724 );
6725
6726 pending.system_events.push(SystemEvent::SocketState(change));
6727 let system_events = pending.take_system_events();
6728
6729 assert_eq!(system_events, vec![SystemEvent::SocketState(change)]);
6730 assert!(pending.is_empty());
6731 }
6732
6733 #[rstest]
6734 fn test_pending_system_commands_stay_separate_from_data() {
6735 let mut pending = PendingEvents::default();
6736 let command = stub_system_command();
6737
6738 pending.system_commands.push(command);
6739 let system_commands = pending.take_system_commands();
6740
6741 assert_eq!(system_commands, vec![command]);
6742 assert!(pending.is_empty());
6743 }
6744
6745 fn stub_time_event_handler() -> TimeEventMessage {
6746 use std::rc::Rc;
6747
6748 use nautilus_common::{
6749 runner::TimeEventMessage,
6750 timer::{TimeEvent, TimeEventCallback},
6751 };
6752 use nautilus_core::{UUID4, UnixNanos};
6753 use ustr::Ustr;
6754
6755 TimeEventMessage::new(
6756 TimeEvent::new(
6757 Ustr::from("test-timer"),
6758 UUID4::new(),
6759 UnixNanos::default(),
6760 UnixNanos::default(),
6761 ),
6762 TimeEventCallback::RustLocal(Rc::new(|_| {})),
6763 )
6764 }
6765
6766 fn stub_trading_command_message() -> TradingCommandMessage {
6767 use nautilus_common::messages::execution::query::QueryAccount;
6768 use nautilus_core::{UUID4, UnixNanos};
6769 use nautilus_model::identifiers::AccountId;
6770
6771 TradingCommandMessage::new(
6772 MessagingSwitchboard::exec_engine_execute(),
6773 TradingCommand::QueryAccount(QueryAccount::new(
6774 TraderId::from("TESTER-001"),
6775 None,
6776 AccountId::from("TEST-001"),
6777 UUID4::new(),
6778 UnixNanos::default(),
6779 None,
6780 None, )),
6782 )
6783 }
6784
6785 fn stub_exec_event() -> ExecutionEvent {
6786 use nautilus_model::{
6787 enums::{LiquiditySide, OrderSide},
6788 identifiers::{AccountId, InstrumentId, TradeId, VenueOrderId},
6789 reports::FillReport,
6790 types::{Money, Price, Quantity},
6791 };
6792
6793 ExecutionEvent::Report(ExecutionReport::Fill(Box::new(FillReport::new(
6794 AccountId::from("TEST-001"),
6795 InstrumentId::from("TEST.VENUE"),
6796 VenueOrderId::from("V-001"),
6797 TradeId::from("T-001"),
6798 OrderSide::Buy,
6799 Quantity::from("1.0"),
6800 Price::from("100.0"),
6801 Money::from("0.01 USD"),
6802 LiquiditySide::Maker,
6803 None,
6804 None,
6805 nautilus_core::UnixNanos::default(),
6806 nautilus_core::UnixNanos::default(),
6807 None,
6808 ))))
6809 }
6810
6811 #[rstest]
6812 fn test_flush_all_pending_drains_buffered_channels() {
6813 let (time_tx, mut time_rx) = tokio::sync::mpsc::unbounded_channel::<TimeEventMessage>();
6814 let (system_evt_tx, mut system_evt_rx) =
6815 tokio::sync::mpsc::unbounded_channel::<SystemEvent>();
6816 let (system_cmd_tx, mut system_cmd_rx) =
6817 tokio::sync::mpsc::unbounded_channel::<SystemCommand>();
6818 let (data_evt_tx, mut data_evt_rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
6819 let (data_cmd_tx, mut data_cmd_rx) = tokio::sync::mpsc::unbounded_channel::<DataCommand>();
6820 let (exec_evt_tx, mut exec_evt_rx) =
6821 tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
6822 let (exec_cmd_tx, mut exec_cmd_rx) =
6823 tokio::sync::mpsc::unbounded_channel::<TradingCommandMessage>();
6824
6825 let mut pending = PendingEvents::default();
6826
6827 pending.data_evts.push(stub_data_event());
6829 pending.data_cmds.push(stub_data_command());
6830
6831 time_tx.send(stub_time_event_handler()).unwrap();
6833 let change = SocketStateChange::new(
6834 ClientId::from("BINANCE"),
6835 Some(Venue::from("BINANCE")),
6836 Ustr::from("binance-futures-market-streams"),
6837 SocketState::Connected,
6838 );
6839 system_evt_tx
6840 .send(SystemEvent::SocketState(change))
6841 .unwrap();
6842 system_cmd_tx.send(stub_system_command()).unwrap();
6843 data_evt_tx.send(stub_data_event()).unwrap();
6844 data_cmd_tx.send(stub_data_command()).unwrap();
6845 exec_evt_tx.send(stub_exec_event()).unwrap();
6846 exec_cmd_tx.send(stub_trading_command_message()).unwrap();
6847
6848 flush_all_pending(
6849 &mut pending,
6850 &mut time_rx,
6851 &mut system_evt_rx,
6852 &mut system_cmd_rx,
6853 &mut exec_evt_rx,
6854 &mut exec_cmd_rx,
6855 &mut data_evt_rx,
6856 &mut data_cmd_rx,
6857 );
6858
6859 let system_events = pending.take_system_events();
6860 let system_commands = pending.take_system_commands();
6861 assert_eq!(system_events, vec![SystemEvent::SocketState(change)]);
6862 assert_eq!(system_commands, vec![stub_system_command()]);
6863 assert!(pending.data_evts.is_empty());
6864 assert!(pending.data_cmds.is_empty());
6865 assert!(pending.exec_reports.is_empty());
6866 assert!(pending.exec_cmds.is_empty());
6867 assert!(pending.order_evts.is_empty());
6868 assert!(time_rx.try_recv().is_err());
6869 assert!(system_evt_rx.try_recv().is_err());
6870 assert!(system_cmd_rx.try_recv().is_err());
6871 assert!(data_evt_rx.try_recv().is_err());
6872 assert!(data_cmd_rx.try_recv().is_err());
6873 assert!(exec_evt_rx.try_recv().is_err());
6874 assert!(exec_cmd_rx.try_recv().is_err());
6875 }
6876
6877 fn stub_order_event() -> ExecutionEvent {
6878 use nautilus_model::events::order::spec::OrderSubmittedSpec;
6879
6880 ExecutionEvent::Order(OrderEventAny::Submitted(
6881 OrderSubmittedSpec::builder().build(),
6882 ))
6883 }
6884
6885 fn stub_account_event() -> ExecutionEvent {
6886 use nautilus_core::{UUID4, UnixNanos};
6887 use nautilus_model::{
6888 enums::AccountType, events::account::state::AccountState, identifiers::AccountId,
6889 };
6890
6891 ExecutionEvent::Account(AccountState::new(
6892 AccountId::from("TEST-001"),
6893 AccountType::Cash,
6894 vec![],
6895 vec![],
6896 true,
6897 UUID4::new(),
6898 UnixNanos::default(),
6899 UnixNanos::default(),
6900 None,
6901 ))
6902 }
6903
6904 #[rstest]
6905 fn test_flush_all_pending_routes_order_event_to_order_evts() {
6906 let (_time_tx, mut time_rx) = tokio::sync::mpsc::unbounded_channel::<TimeEventMessage>();
6907 let (_system_evt_tx, mut system_evt_rx) =
6908 tokio::sync::mpsc::unbounded_channel::<SystemEvent>();
6909 let (_system_cmd_tx, mut system_cmd_rx) =
6910 tokio::sync::mpsc::unbounded_channel::<SystemCommand>();
6911 let (_data_evt_tx, mut data_evt_rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
6912 let (_data_cmd_tx, mut data_cmd_rx) = tokio::sync::mpsc::unbounded_channel::<DataCommand>();
6913 let (exec_evt_tx, mut exec_evt_rx) =
6914 tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
6915 let (_exec_cmd_tx, mut exec_cmd_rx) =
6916 tokio::sync::mpsc::unbounded_channel::<TradingCommandMessage>();
6917
6918 let mut pending = PendingEvents::default();
6919
6920 exec_evt_tx.send(stub_order_event()).unwrap();
6921 exec_evt_tx.send(stub_exec_event()).unwrap();
6922
6923 flush_all_pending(
6924 &mut pending,
6925 &mut time_rx,
6926 &mut system_evt_rx,
6927 &mut system_cmd_rx,
6928 &mut exec_evt_rx,
6929 &mut exec_cmd_rx,
6930 &mut data_evt_rx,
6931 &mut data_cmd_rx,
6932 );
6933
6934 assert!(pending.order_evts.is_empty());
6936 assert!(pending.exec_reports.is_empty());
6937 assert!(exec_evt_rx.try_recv().is_err());
6938 }
6939
6940 #[rstest]
6941 fn test_flush_all_pending_routes_account_event_immediately() {
6942 let (_time_tx, mut time_rx) = tokio::sync::mpsc::unbounded_channel::<TimeEventMessage>();
6943 let (_system_evt_tx, mut system_evt_rx) =
6944 tokio::sync::mpsc::unbounded_channel::<SystemEvent>();
6945 let (_system_cmd_tx, mut system_cmd_rx) =
6946 tokio::sync::mpsc::unbounded_channel::<SystemCommand>();
6947 let (_data_evt_tx, mut data_evt_rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
6948 let (_data_cmd_tx, mut data_cmd_rx) = tokio::sync::mpsc::unbounded_channel::<DataCommand>();
6949 let (exec_evt_tx, mut exec_evt_rx) =
6950 tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
6951 let (_exec_cmd_tx, mut exec_cmd_rx) =
6952 tokio::sync::mpsc::unbounded_channel::<TradingCommandMessage>();
6953
6954 let mut pending = PendingEvents::default();
6955
6956 exec_evt_tx.send(stub_account_event()).unwrap();
6957
6958 flush_all_pending(
6959 &mut pending,
6960 &mut time_rx,
6961 &mut system_evt_rx,
6962 &mut system_cmd_rx,
6963 &mut exec_evt_rx,
6964 &mut exec_cmd_rx,
6965 &mut data_evt_rx,
6966 &mut data_cmd_rx,
6967 );
6968
6969 assert!(pending.exec_reports.is_empty());
6971 assert!(pending.order_evts.is_empty());
6972 assert!(pending.exec_cmds.is_empty());
6973 assert!(exec_evt_rx.try_recv().is_err());
6974 }
6975
6976 #[rstest]
6977 fn test_pending_is_empty_when_default() {
6978 let pending = PendingEvents::default();
6979
6980 assert!(pending.is_empty());
6981 }
6982
6983 #[rstest]
6984 fn test_pending_is_empty_false_with_data_evt() {
6985 let mut pending = PendingEvents::default();
6986 pending.data_evts.push(stub_data_event());
6987
6988 assert!(!pending.is_empty());
6989 }
6990
6991 #[rstest]
6992 fn test_pending_is_empty_false_with_data_cmd() {
6993 let mut pending = PendingEvents::default();
6994 pending.data_cmds.push(stub_data_command());
6995
6996 assert!(!pending.is_empty());
6997 }
6998
6999 #[rstest]
7000 fn test_pending_is_empty_false_with_exec_cmd() {
7001 let mut pending = PendingEvents::default();
7002 pending.exec_cmds.push(stub_trading_command_message());
7003
7004 assert!(!pending.is_empty());
7005 }
7006
7007 #[rstest]
7008 fn test_pending_drain_preserves_trading_command_target() {
7009 std::thread::spawn(|| {
7010 msgbus::get_message_bus().borrow_mut().dispose();
7011 let risk_commands = Rc::new(RefCell::new(Vec::new()));
7012 let exec_commands = Rc::new(RefCell::new(Vec::new()));
7013
7014 let risk_commands_handler = risk_commands.clone();
7015 msgbus::register_trading_command_endpoint(
7016 MessagingSwitchboard::risk_engine_execute(),
7017 TypedIntoHandler::from(move |command: TradingCommand| {
7018 risk_commands_handler.borrow_mut().push(command);
7019 }),
7020 );
7021 let exec_commands_handler = exec_commands.clone();
7022 msgbus::register_trading_command_endpoint(
7023 MessagingSwitchboard::exec_engine_execute(),
7024 TypedIntoHandler::from(move |command: TradingCommand| {
7025 exec_commands_handler.borrow_mut().push(command);
7026 }),
7027 );
7028
7029 let mut pending = PendingEvents::default();
7030 pending.exec_cmds.push(TradingCommandMessage::new(
7031 MessagingSwitchboard::risk_engine_execute(),
7032 TradingCommand::QueryAccount(QueryAccount::new(
7033 TraderId::from("TESTER-001"),
7034 None,
7035 AccountId::from("TEST-001"),
7036 UUID4::new(),
7037 UnixNanos::default(),
7038 None,
7039 None,
7040 )),
7041 ));
7042
7043 pending.drain();
7044
7045 assert!(pending.is_empty());
7046 assert_eq!(risk_commands.borrow().len(), 1);
7047 assert!(matches!(
7048 &risk_commands.borrow()[0],
7049 TradingCommand::QueryAccount(_)
7050 ));
7051 assert_eq!(exec_commands.borrow().as_slice(), &[]);
7052 })
7053 .join()
7054 .unwrap();
7055 }
7056
7057 #[rstest]
7058 fn test_pending_is_empty_false_with_exec_report() {
7059 let mut pending = PendingEvents::default();
7060
7061 if let ExecutionEvent::Report(report) = stub_exec_event() {
7062 pending.exec_reports.push(report);
7063 }
7064
7065 assert!(!pending.is_empty());
7066 }
7067
7068 #[rstest]
7069 fn test_pending_is_empty_false_with_order_evt() {
7070 let mut pending = PendingEvents::default();
7071
7072 if let ExecutionEvent::Order(order_evt) = stub_order_event() {
7073 pending.order_evts.push(order_evt);
7074 }
7075
7076 assert!(!pending.is_empty());
7077 }
7078
7079 fn stub_submitted_batch_event() -> ExecutionEvent {
7080 use nautilus_model::{
7081 events::{OrderSubmittedBatch, order::spec::OrderSubmittedSpec},
7082 identifiers::ClientOrderId,
7083 };
7084
7085 let events = vec![
7086 OrderSubmittedSpec::builder()
7087 .client_order_id(ClientOrderId::from("O-001"))
7088 .build(),
7089 OrderSubmittedSpec::builder()
7090 .client_order_id(ClientOrderId::from("O-002"))
7091 .build(),
7092 ];
7093
7094 ExecutionEvent::OrderSubmittedBatch(OrderSubmittedBatch::new(events))
7095 }
7096
7097 fn stub_canceled_batch_event() -> ExecutionEvent {
7098 use nautilus_model::{
7099 events::{OrderCanceledBatch, order::spec::OrderCanceledSpec},
7100 identifiers::ClientOrderId,
7101 };
7102
7103 let events = vec![
7104 OrderCanceledSpec::builder()
7105 .client_order_id(ClientOrderId::from("O-001"))
7106 .build(),
7107 OrderCanceledSpec::builder()
7108 .client_order_id(ClientOrderId::from("O-002"))
7109 .build(),
7110 ];
7111
7112 ExecutionEvent::OrderCanceledBatch(OrderCanceledBatch::new(events))
7113 }
7114
7115 #[rstest]
7116 fn test_flush_all_pending_buffers_submitted_batch_as_individual_events() {
7117 let (_time_tx, mut time_rx) = tokio::sync::mpsc::unbounded_channel::<TimeEventMessage>();
7118 let (_system_evt_tx, mut system_evt_rx) =
7119 tokio::sync::mpsc::unbounded_channel::<SystemEvent>();
7120 let (_system_cmd_tx, mut system_cmd_rx) =
7121 tokio::sync::mpsc::unbounded_channel::<SystemCommand>();
7122 let (_data_evt_tx, mut data_evt_rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
7123 let (_data_cmd_tx, mut data_cmd_rx) = tokio::sync::mpsc::unbounded_channel::<DataCommand>();
7124 let (exec_evt_tx, mut exec_evt_rx) =
7125 tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
7126 let (_exec_cmd_tx, mut exec_cmd_rx) =
7127 tokio::sync::mpsc::unbounded_channel::<TradingCommandMessage>();
7128
7129 let mut pending = PendingEvents::default();
7130
7131 exec_evt_tx.send(stub_submitted_batch_event()).unwrap();
7132
7133 flush_all_pending(
7134 &mut pending,
7135 &mut time_rx,
7136 &mut system_evt_rx,
7137 &mut system_cmd_rx,
7138 &mut exec_evt_rx,
7139 &mut exec_cmd_rx,
7140 &mut data_evt_rx,
7141 &mut data_cmd_rx,
7142 );
7143
7144 assert!(pending.order_evts.is_empty());
7146 assert!(exec_evt_rx.try_recv().is_err());
7147 }
7148
7149 #[rstest]
7150 fn test_flush_all_pending_buffers_canceled_batch_as_individual_events() {
7151 let (_time_tx, mut time_rx) = tokio::sync::mpsc::unbounded_channel::<TimeEventMessage>();
7152 let (_system_evt_tx, mut system_evt_rx) =
7153 tokio::sync::mpsc::unbounded_channel::<SystemEvent>();
7154 let (_system_cmd_tx, mut system_cmd_rx) =
7155 tokio::sync::mpsc::unbounded_channel::<SystemCommand>();
7156 let (_data_evt_tx, mut data_evt_rx) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
7157 let (_data_cmd_tx, mut data_cmd_rx) = tokio::sync::mpsc::unbounded_channel::<DataCommand>();
7158 let (exec_evt_tx, mut exec_evt_rx) =
7159 tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
7160 let (_exec_cmd_tx, mut exec_cmd_rx) =
7161 tokio::sync::mpsc::unbounded_channel::<TradingCommandMessage>();
7162
7163 let mut pending = PendingEvents::default();
7164
7165 exec_evt_tx.send(stub_canceled_batch_event()).unwrap();
7166
7167 flush_all_pending(
7168 &mut pending,
7169 &mut time_rx,
7170 &mut system_evt_rx,
7171 &mut system_cmd_rx,
7172 &mut exec_evt_rx,
7173 &mut exec_cmd_rx,
7174 &mut data_evt_rx,
7175 &mut data_cmd_rx,
7176 );
7177
7178 assert!(pending.order_evts.is_empty());
7180 assert!(exec_evt_rx.try_recv().is_err());
7181 }
7182
7183 #[rstest]
7184 fn test_flush_all_pending_expands_batch_into_order_evts_before_drain() {
7185 use nautilus_model::identifiers::ClientOrderId;
7186
7187 let (exec_evt_tx, mut exec_evt_rx) =
7188 tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
7189
7190 exec_evt_tx.send(stub_canceled_batch_event()).unwrap();
7191
7192 let mut pending = PendingEvents::default();
7193
7194 while let Ok(evt) = exec_evt_rx.try_recv() {
7196 match evt {
7197 ExecutionEvent::Account(_) => {
7198 AsyncRunner::handle_exec_event(evt);
7199 }
7200 ExecutionEvent::Report(report) => {
7201 pending.exec_reports.push(report);
7202 }
7203 ExecutionEvent::Order(order_evt) => {
7204 pending.order_evts.push(order_evt);
7205 }
7206 ExecutionEvent::OrderSubmittedBatch(batch) => {
7207 for submitted in batch {
7208 pending.order_evts.push(OrderEventAny::Submitted(submitted));
7209 }
7210 }
7211 ExecutionEvent::OrderAcceptedBatch(batch) => {
7212 for accepted in batch {
7213 pending.order_evts.push(OrderEventAny::Accepted(accepted));
7214 }
7215 }
7216 ExecutionEvent::OrderCanceledBatch(batch) => {
7217 for canceled in batch {
7218 pending.order_evts.push(OrderEventAny::Canceled(canceled));
7219 }
7220 }
7221 }
7222 }
7223
7224 assert_eq!(pending.order_evts.len(), 2);
7225 assert!(
7226 matches!(&pending.order_evts[0], OrderEventAny::Canceled(c) if c.client_order_id == ClientOrderId::from("O-001"))
7227 );
7228 assert!(
7229 matches!(&pending.order_evts[1], OrderEventAny::Canceled(c) if c.client_order_id == ClientOrderId::from("O-002"))
7230 );
7231 }
7232
7233 #[derive(Debug)]
7234 struct CapturedEgressMessage {
7235 topic: String,
7236 payload: Bytes,
7237 }
7238
7239 type CapturedEgressMessages = Rc<RefCell<Vec<CapturedEgressMessage>>>;
7240 type SharedClosed = Rc<Cell<bool>>;
7241
7242 #[derive(Debug)]
7243 struct CapturingExternalIngress {
7244 rx: Option<tokio::sync::mpsc::Receiver<BusMessage>>,
7245 closed: SharedClosed,
7246 }
7247
7248 impl CapturingExternalIngress {
7249 fn new(rx: tokio::sync::mpsc::Receiver<BusMessage>, closed: SharedClosed) -> Self {
7250 Self {
7251 rx: Some(rx),
7252 closed,
7253 }
7254 }
7255 }
7256
7257 impl MessageBusExternalIngress for CapturingExternalIngress {
7258 fn is_closed(&self) -> bool {
7259 self.closed.get()
7260 }
7261
7262 fn take_receiver(&mut self) -> anyhow::Result<tokio::sync::mpsc::Receiver<BusMessage>> {
7263 self.rx
7264 .take()
7265 .ok_or_else(|| anyhow::anyhow!("external ingress receiver already taken"))
7266 }
7267
7268 fn close(&mut self) {
7269 self.closed.set(true);
7270 }
7271 }
7272
7273 #[derive(Debug)]
7274 struct FailingExternalIngress {
7275 closed: SharedClosed,
7276 }
7277
7278 impl FailingExternalIngress {
7279 fn new(closed: SharedClosed) -> Self {
7280 Self { closed }
7281 }
7282 }
7283
7284 impl MessageBusExternalIngress for FailingExternalIngress {
7285 fn is_closed(&self) -> bool {
7286 self.closed.get()
7287 }
7288
7289 fn take_receiver(&mut self) -> anyhow::Result<tokio::sync::mpsc::Receiver<BusMessage>> {
7290 anyhow::bail!("external ingress receiver unavailable")
7291 }
7292
7293 fn close(&mut self) {
7294 self.closed.set(true);
7295 }
7296 }
7297
7298 struct CapturingExternalEgress {
7299 publications: CapturedEgressMessages,
7300 closed: SharedClosed,
7301 }
7302
7303 impl CapturingExternalEgress {
7304 fn new() -> (Self, CapturedEgressMessages, SharedClosed) {
7305 let publications = Rc::new(RefCell::new(Vec::new()));
7306 let closed = Rc::new(Cell::new(false));
7307 (
7308 Self {
7309 publications: publications.clone(),
7310 closed: closed.clone(),
7311 },
7312 publications,
7313 closed,
7314 )
7315 }
7316 }
7317
7318 impl MessageBusExternalEgress for CapturingExternalEgress {
7319 fn is_closed(&self) -> bool {
7320 self.closed.get()
7321 }
7322
7323 fn publish(&self, message: BusMessage) {
7324 self.publications.borrow_mut().push(CapturedEgressMessage {
7325 topic: message.topic.to_string(),
7326 payload: message.payload,
7327 });
7328 }
7329
7330 fn close(&mut self) {
7331 self.closed.set(true);
7332 }
7333 }
7334
7335 struct CapturingBackingFactory {
7336 publications: Arc<Mutex<Vec<CapturedEgressMessage>>>,
7337 closed: Arc<AtomicBool>,
7338 rx: Mutex<Option<tokio::sync::mpsc::Receiver<BusMessage>>>,
7339 }
7340
7341 impl CapturingBackingFactory {
7342 fn new(
7343 publications: Arc<Mutex<Vec<CapturedEgressMessage>>>,
7344 closed: Arc<AtomicBool>,
7345 rx: Option<tokio::sync::mpsc::Receiver<BusMessage>>,
7346 ) -> Self {
7347 Self {
7348 publications,
7349 closed,
7350 rx: Mutex::new(rx),
7351 }
7352 }
7353 }
7354
7355 impl Debug for CapturingBackingFactory {
7356 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
7357 f.debug_struct(stringify!(CapturingBackingFactory))
7358 .finish_non_exhaustive()
7359 }
7360 }
7361
7362 impl MessageBusBackingFactory for CapturingBackingFactory {
7363 fn create(
7364 &self,
7365 _trader_id: TraderId,
7366 _instance_id: UUID4,
7367 _config: MessageBusConfig,
7368 ) -> anyhow::Result<Box<dyn MessageBusBacking>> {
7369 let rx = self.rx.lock().unwrap().take();
7370 Ok(Box::new(CapturingBacking {
7371 publications: self.publications.clone(),
7372 closed: self.closed.clone(),
7373 rx,
7374 }))
7375 }
7376 }
7377
7378 struct CapturingBacking {
7379 publications: Arc<Mutex<Vec<CapturedEgressMessage>>>,
7380 closed: Arc<AtomicBool>,
7381 rx: Option<tokio::sync::mpsc::Receiver<BusMessage>>,
7382 }
7383
7384 impl MessageBusBacking for CapturingBacking {
7385 fn is_closed(&self) -> bool {
7386 self.closed.load(Ordering::Relaxed)
7387 }
7388
7389 fn publish(&self, message: BusMessage) {
7390 self.publications
7391 .lock()
7392 .unwrap()
7393 .push(CapturedEgressMessage {
7394 topic: message.topic.to_string(),
7395 payload: message.payload,
7396 });
7397 }
7398
7399 fn take_receiver(&mut self) -> anyhow::Result<tokio::sync::mpsc::Receiver<BusMessage>> {
7400 self.rx
7401 .take()
7402 .ok_or_else(|| anyhow::anyhow!("external ingress receiver unavailable"))
7403 }
7404
7405 fn close(&mut self) {
7406 self.closed.store(true, Ordering::Relaxed);
7407 }
7408 }
7409}