Skip to main content

nautilus_live/node/
mod.rs

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