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