use std::{
cell::{Cell, Ref, RefCell},
fmt::Debug,
rc::Rc,
time::Duration,
};
use nautilus_common::{
cache::{Cache, CacheConfig, database::CacheDatabaseAdapter},
clients::{SocketReconnectLookup, SocketReconnectRequestOutcome},
clock::Clock,
component::Component,
enums::{ComponentState, Environment},
logging::{
arm_shutdown_on_error, disarm_shutdown_on_error, headers, init_logging,
logger::{LogGuard, LoggerConfig},
try_drain_shutdown_on_error_trigger,
},
messages::system::{ReconnectSocket, ShutdownSystem},
msgbus::{
self, MessageBus, MessagingSwitchboard, ShareableMessageHandler, get_message_bus,
set_message_bus,
},
};
use nautilus_core::{UUID4, UnixNanos};
use nautilus_data::engine::DataEngine;
use nautilus_execution::{
engine::ExecutionEngine,
order_emulator::{adapter::OrderEmulatorAdapter, emulator::OrderEmulator},
};
use nautilus_model::identifiers::{ClientId, TraderId};
use nautilus_portfolio::portfolio::Portfolio;
use nautilus_risk::engine::RiskEngine;
use ustr::Ustr;
use crate::{
builder::NautilusKernelBuilder,
clock_factory::ClockFactory,
config::NautilusKernelConfig,
event_store::{EventStoreFactory, KernelEventStore, RegisteredComponents},
trader::Trader,
};
#[derive(Debug)]
pub struct NautilusKernel {
pub name: String,
pub instance_id: UUID4,
pub machine_id: String,
pub config: Box<dyn NautilusKernelConfig>,
pub cache: Rc<RefCell<Cache>>,
pub clock: Rc<RefCell<dyn Clock>>,
pub portfolio: Rc<RefCell<Portfolio>>,
pub log_guard: LogGuard,
pub data_engine: Rc<RefCell<DataEngine>>,
pub risk_engine: Rc<RefCell<RiskEngine>>,
pub exec_engine: Rc<RefCell<ExecutionEngine>>,
pub order_emulator: OrderEmulatorAdapter,
pub trader: Rc<RefCell<Trader>>,
pub ts_created: UnixNanos,
pub ts_started: Option<UnixNanos>,
pub ts_shutdown: Option<UnixNanos>,
shutdown_requested: Rc<Cell<bool>>,
event_store: Option<Box<dyn KernelEventStore>>,
event_store_replay: bool,
state_save_armed: bool,
}
#[derive(Default)]
pub struct NautilusKernelDependencies {
clock_factory: Option<ClockFactory>,
cache_database: Option<Box<dyn CacheDatabaseAdapter>>,
event_store_factory: Option<EventStoreFactory>,
}
impl Debug for NautilusKernelDependencies {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct(stringify!(NautilusKernelDependencies))
.field("clock_factory", &self.clock_factory.is_some())
.field("cache_database", &self.cache_database.is_some())
.field("event_store_factory", &self.event_store_factory.is_some())
.finish()
}
}
impl NautilusKernelDependencies {
#[must_use]
pub fn with_clock_factory(mut self, clock_factory: Option<ClockFactory>) -> Self {
self.clock_factory = clock_factory;
self
}
#[must_use]
pub fn with_cache_database(
mut self,
cache_database: Option<Box<dyn CacheDatabaseAdapter>>,
) -> Self {
self.cache_database = cache_database;
self
}
#[must_use]
pub fn with_event_store_factory(
mut self,
event_store_factory: Option<EventStoreFactory>,
) -> Self {
self.event_store_factory = event_store_factory;
self
}
}
impl NautilusKernel {
#[must_use]
pub const fn builder(
name: String,
trader_id: TraderId,
environment: Environment,
) -> NautilusKernelBuilder {
NautilusKernelBuilder::new(name, trader_id, environment)
}
pub fn new<T: NautilusKernelConfig + 'static>(name: String, config: T) -> anyhow::Result<Self> {
Self::new_with(name, config, None, None)
}
pub fn new_with_cache_database<T: NautilusKernelConfig + 'static>(
name: String,
config: T,
cache_database: Option<Box<dyn CacheDatabaseAdapter>>,
) -> anyhow::Result<Self> {
Self::new_with(name, config, cache_database, None)
}
pub fn new_with<T: NautilusKernelConfig + 'static>(
name: String,
config: T,
cache_database: Option<Box<dyn CacheDatabaseAdapter>>,
event_store_factory: Option<EventStoreFactory>,
) -> anyhow::Result<Self> {
Self::new_with_dependencies(
name,
config,
NautilusKernelDependencies::default()
.with_cache_database(cache_database)
.with_event_store_factory(event_store_factory),
)
}
pub fn new_with_dependencies<T: NautilusKernelConfig + 'static>(
name: String,
config: T,
dependencies: NautilusKernelDependencies,
) -> anyhow::Result<Self> {
let NautilusKernelDependencies {
clock_factory,
cache_database,
event_store_factory,
} = dependencies;
let instance_id = config.instance_id().unwrap_or_default();
let machine_id = Self::determine_machine_id()?;
let logger_config = config.logging();
let log_guard = Self::initialize_logging(config.trader_id(), instance_id, logger_config)?;
headers::log_header(
config.trader_id(),
&machine_id,
instance_id,
Ustr::from(&name),
);
log::info!("Building system kernel");
let clock_factory =
clock_factory.unwrap_or_else(|| ClockFactory::for_environment(config.environment()));
let clock = clock_factory.clock();
let event_store = match event_store_factory {
Some(factory) => Some(factory(instance_id, clock.clone())?),
None => None,
};
let cache = Self::initialize_cache(config.cache(), cache_database);
let msgbus = Rc::new(RefCell::new(MessageBus::new(
config.trader_id(),
instance_id,
Some(name.clone()),
None,
)));
set_message_bus(msgbus);
if let Some(config) = config.msgbus()
&& let Some(filter) = config.types_filter
{
get_message_bus().borrow_mut().set_types_filter(filter);
}
let portfolio = Rc::new(RefCell::new(Portfolio::new(
clock.clone(),
cache.clone(),
config.portfolio(),
)));
let risk_engine = RiskEngine::new(
config.risk_engine().unwrap_or_default(),
portfolio.borrow().clone_shallow(),
clock.clone(),
cache.clone(),
);
let risk_engine = Rc::new(RefCell::new(risk_engine));
let exec_engine = ExecutionEngine::new(clock.clone(), cache.clone(), config.exec_engine());
let exec_engine = Rc::new(RefCell::new(exec_engine));
let order_emulator = OrderEmulatorAdapter::new(clock.clone(), cache.clone());
let data_engine = DataEngine::new(clock.clone(), cache.clone(), config.data_engine());
let data_engine = Rc::new(RefCell::new(data_engine));
DataEngine::register_msgbus_handlers(&data_engine);
RiskEngine::register_msgbus_handlers(&risk_engine);
ExecutionEngine::register_msgbus_handlers(&exec_engine);
OrderEmulator::register_msgbus_handlers(&order_emulator.emulator());
let shutdown_requested = Rc::new(Cell::new(false));
Self::register_shutdown_handler(config.trader_id(), shutdown_requested.clone());
let trader = Rc::new(RefCell::new(Trader::new(
config.trader_id(),
instance_id,
config.environment(),
clock_factory,
cache.clone(),
portfolio.clone(),
)));
let ts_created = clock.borrow().timestamp_ns();
Ok(Self {
name,
instance_id,
machine_id,
event_store,
config: Box::new(config),
cache,
clock,
portfolio,
log_guard,
data_engine,
risk_engine,
exec_engine,
order_emulator,
trader,
ts_created,
ts_started: None,
ts_shutdown: None,
shutdown_requested,
event_store_replay: false,
state_save_armed: false,
})
}
fn register_shutdown_handler(trader_id: TraderId, shutdown_requested: Rc<Cell<bool>>) {
let handler = ShareableMessageHandler::from_typed(move |cmd: &ShutdownSystem| {
if cmd.trader_id != trader_id {
log::warn!("Received {cmd} not for this trader {trader_id}, ignoring");
return;
}
if shutdown_requested.get() {
log::debug!("Shutdown already requested, ignoring {cmd}");
return;
}
log::info!("Received {cmd}, requesting shutdown");
shutdown_requested.set(true);
});
let topic = MessagingSwitchboard::shutdown_system_topic();
msgbus::subscribe_any(topic.into(), handler, None);
}
fn determine_machine_id() -> anyhow::Result<String> {
sysinfo::System::host_name().ok_or_else(|| anyhow::anyhow!("Failed to determine hostname"))
}
fn initialize_logging(
trader_id: TraderId,
instance_id: UUID4,
config: LoggerConfig,
) -> anyhow::Result<LogGuard> {
#[cfg(feature = "tracing-bridge")]
let use_tracing = config.use_tracing;
let file_config = config.file_config.clone().unwrap_or_default();
let log_guard = match init_logging(trader_id, instance_id, config, file_config) {
Ok(guard) => guard,
Err(e) => {
if e.downcast_ref::<log::SetLoggerError>().is_some() {
if let Some(guard) = LogGuard::new() {
guard
} else {
return Err(e.context(
"A non-Nautilus logger is already registered; \
cannot initialize Nautilus logging",
));
}
} else {
return Err(e);
}
}
};
#[cfg(feature = "tracing-bridge")]
if use_tracing && !nautilus_common::logging::bridge::tracing_is_initialized() {
nautilus_common::logging::bridge::init_tracing()?;
}
Ok(log_guard)
}
fn initialize_cache(
cache_config: Option<CacheConfig>,
cache_database: Option<Box<dyn CacheDatabaseAdapter>>,
) -> Rc<RefCell<Cache>> {
let cache_config = cache_config.unwrap_or_default();
let cache = Cache::new(Some(cache_config), cache_database);
Rc::new(RefCell::new(cache))
}
fn cancel_timers(&self) {
self.clock.borrow_mut().cancel_timers();
}
#[must_use]
pub fn generate_timestamp_ns(&self) -> UnixNanos {
self.clock.borrow().timestamp_ns()
}
pub fn process_socket_reconnect(&self, command: ReconnectSocket) {
let outcome = if command.trader_id == self.config.trader_id() {
let data = self
.data_engine
.borrow()
.socket_reconnect_lookup(&command.client_id, command.endpoint);
let execution = self
.exec_engine
.borrow()
.socket_reconnect_lookup(&command.client_id, command.endpoint);
Self::request_socket_reconnect(data, execution)
} else {
SocketReconnectDispatchOutcome::InvalidTrader
};
if outcome == SocketReconnectDispatchOutcome::Accepted {
log::info!(
"Requested socket reconnect for client {}",
command.client_id
);
} else {
log::warn!(
"Rejected socket reconnect request for client {}: {outcome:?}",
command.client_id
);
}
}
fn request_socket_reconnect(
data: SocketReconnectLookup,
execution: SocketReconnectLookup,
) -> SocketReconnectDispatchOutcome {
match (data, execution) {
(SocketReconnectLookup::Handle(_), SocketReconnectLookup::Handle(_)) => {
SocketReconnectDispatchOutcome::AmbiguousEndpoint
}
(SocketReconnectLookup::Handle(handle), _)
| (_, SocketReconnectLookup::Handle(handle)) => match handle.request_reconnect() {
SocketReconnectRequestOutcome::Accepted => SocketReconnectDispatchOutcome::Accepted,
SocketReconnectRequestOutcome::AlreadyReconnecting => {
SocketReconnectDispatchOutcome::AlreadyReconnecting
}
SocketReconnectRequestOutcome::Disconnected => {
SocketReconnectDispatchOutcome::Disconnected
}
SocketReconnectRequestOutcome::Closed => SocketReconnectDispatchOutcome::Closed,
SocketReconnectRequestOutcome::Unsupported => {
SocketReconnectDispatchOutcome::Unsupported
}
},
(SocketReconnectLookup::EndpointNotFound, _)
| (_, SocketReconnectLookup::EndpointNotFound) => {
SocketReconnectDispatchOutcome::UnknownEndpoint
}
(SocketReconnectLookup::Unsupported, _) | (_, SocketReconnectLookup::Unsupported) => {
SocketReconnectDispatchOutcome::Unsupported
}
(SocketReconnectLookup::ClientNotFound, SocketReconnectLookup::ClientNotFound) => {
SocketReconnectDispatchOutcome::UnknownClient
}
}
}
#[must_use]
pub fn environment(&self) -> Environment {
self.config.environment()
}
#[must_use]
pub const fn name(&self) -> &str {
self.name.as_str()
}
#[must_use]
pub fn trader_id(&self) -> TraderId {
self.config.trader_id()
}
#[must_use]
pub fn machine_id(&self) -> &str {
&self.machine_id
}
#[must_use]
pub const fn instance_id(&self) -> UUID4 {
self.instance_id
}
#[must_use]
pub fn delay_post_stop(&self) -> Duration {
self.config.delay_post_stop()
}
#[must_use]
pub const fn ts_created(&self) -> UnixNanos {
self.ts_created
}
#[must_use]
pub const fn ts_started(&self) -> Option<UnixNanos> {
self.ts_started
}
#[must_use]
pub const fn ts_shutdown(&self) -> Option<UnixNanos> {
self.ts_shutdown
}
#[must_use]
pub fn is_shutdown_requested(&self) -> bool {
self.drain_shutdown_on_error_trigger();
self.shutdown_requested.get()
}
pub fn reset_shutdown_flag(&self) {
self.shutdown_requested.set(false);
}
#[must_use]
pub fn shutdown_flag(&self) -> Rc<Cell<bool>> {
self.shutdown_requested.clone()
}
fn drain_shutdown_on_error_trigger(&self) {
try_drain_shutdown_on_error_trigger(|trigger| {
let command = ShutdownSystem::new(
self.config.trader_id(),
trigger.component,
Some(format!(
"Error log received from {}: {}",
trigger.component, trigger.message
)),
UUID4::new(),
trigger.timestamp,
None,
);
msgbus::try_publish_any(
MessagingSwitchboard::shutdown_system_topic(),
command.as_any(),
)
});
}
#[must_use]
pub fn load_state(&self) -> bool {
self.config.load_state()
}
#[must_use]
pub fn save_state(&self) -> bool {
self.config.save_state()
}
#[must_use]
pub fn clock(&self) -> Rc<RefCell<dyn Clock>> {
self.clock.clone()
}
#[must_use]
pub fn cache(&self) -> Rc<RefCell<Cache>> {
self.cache.clone()
}
#[must_use]
pub fn portfolio(&self) -> Ref<'_, Portfolio> {
self.portfolio.borrow()
}
#[must_use]
pub fn data_engine(&self) -> Ref<'_, DataEngine> {
self.data_engine.borrow()
}
#[must_use]
pub const fn risk_engine(&self) -> &Rc<RefCell<RiskEngine>> {
&self.risk_engine
}
#[must_use]
pub const fn exec_engine(&self) -> &Rc<RefCell<ExecutionEngine>> {
&self.exec_engine
}
#[must_use]
pub fn trader(&self) -> &Rc<RefCell<Trader>> {
&self.trader
}
pub fn start(&mut self) {
arm_shutdown_on_error(self.config.shutdown_on_error());
log::info!("Starting");
self.event_store_replay = false;
if let Some(event_store) = self.event_store.as_deref_mut() {
self.exec_engine.borrow_mut().set_snapshot_anchorer(None);
let components = Self::collect_registered_components(&self.trader);
let environment = self.config.environment();
let event_store_replay_configured = event_store.is_event_store_replay_configured();
if event_store_replay_configured && !self.config.load_state() {
log::error!("Event-store replay requires load_state=true");
return;
}
if self.config.load_state()
&& let Err(e) =
event_store.restore_parent_cache(self.instance_id, &mut self.cache.borrow_mut())
{
log::error!("Failed to restore cache from event-store replay source: {e}");
return;
}
if let Err(e) = event_store.open(self.instance_id, &components, environment) {
log::error!("Failed to open event-store run: {e}");
return;
}
let anchorer = event_store.snapshot_anchorer();
self.exec_engine
.borrow_mut()
.set_snapshot_anchorer(anchorer);
self.event_store_replay = event_store_replay_configured;
}
if self.event_store_replay {
log::info!(
"Event-store replay loaded; skipping engines, clients, trader startup, and live reconciliation",
);
self.ts_started = Some(self.clock.borrow().timestamp_ns());
log::info!("Started");
return;
}
self.start_engines();
log::info!("Initializing trader");
if let Err(e) = self.trader.borrow_mut().initialize() {
log::error!("Error initializing trader: {e:?}");
return;
}
self.ts_started = Some(self.clock.borrow().timestamp_ns());
log::info!("Started");
}
fn collect_registered_components(trader: &Rc<RefCell<Trader>>) -> RegisteredComponents {
let trader = trader.borrow();
let mut components = RegisteredComponents::default();
for actor_id in trader.actor_ids() {
components
.actors
.insert(actor_id.to_string(), String::new());
}
for strategy_id in trader.strategy_ids() {
components
.strategies
.insert(strategy_id.to_string(), String::new());
}
for algo_id in trader.exec_algorithm_ids() {
components
.algorithms
.insert(algo_id.to_string(), String::new());
}
components
}
#[expect(
clippy::unused_async,
reason = "keeps the public async kernel API shape stable"
)]
pub async fn start_async(&mut self) {
self.start();
}
pub fn start_trader(&mut self) -> anyhow::Result<()> {
log::info!("Starting trader...");
let load_state = self.config.load_state();
let save_state = self.config.save_state();
if (load_state || save_state) && !self.cache.borrow().has_backing() {
log::warn!(
"Cache has no database backing, load_state={load_state} and save_state={save_state} will have no effect"
);
}
if load_state {
Trader::load_state(&self.trader)
.map_err(|e| anyhow::anyhow!("Failed to load actor and strategy state: {e:#}"))?;
}
self.state_save_armed = save_state;
self.order_emulator.start();
if let Err(start_err) = Trader::start_with_component_callbacks(&self.trader) {
let stop_result = self.stop_trader_after_start_failure();
self.order_emulator.stop();
let save_result = self.save_trader_state();
let mut errors = vec![format!("Failed to start trader: {start_err}")];
if let Err(e) = stop_result {
errors.push(format!("failed to stop partial trader start: {e}"));
}
if let Err(e) = save_result {
errors.push(format!("failed to save partial trader state: {e}"));
}
anyhow::bail!("{}", errors.join("; "));
}
log::info!("Trader started");
Ok(())
}
pub fn stop_trader(&mut self) {
disarm_shutdown_on_error();
if !self.trader.borrow().is_running() {
return;
}
log::info!("Stopping trader...");
if let Err(e) = self.trader.borrow_mut().stop() {
log::error!("Error stopping trader: {e}");
}
}
pub fn stop_trader_after_start_failure(&mut self) -> anyhow::Result<()> {
disarm_shutdown_on_error();
if !matches!(
self.trader.borrow().state(),
ComponentState::Starting | ComponentState::Running
) {
return Ok(());
}
log::info!("Stopping trader immediately...");
self.trader.borrow_mut().stop_after_start_failure()
}
#[allow(unknown_lints)]
#[expect(
clippy::unused_async,
clippy::unused_async_trait_impl,
reason = "keeps the public async kernel API shape stable"
)]
pub async fn finalize_stop(&mut self) -> anyhow::Result<()> {
disarm_shutdown_on_error();
let save_result = self.save_trader_state();
self.portfolio.borrow_mut().finalize_equity_curve();
self.stop_engines();
self.cancel_timers();
let ts_shutdown = self.clock.borrow().timestamp_ns();
if let Some(event_store) = self.event_store.as_deref_mut() {
self.exec_engine.borrow_mut().set_snapshot_anchorer(None);
event_store.seal(ts_shutdown);
}
self.ts_shutdown = Some(ts_shutdown);
log::info!("Stopped");
save_result
}
pub fn save_trader_state(&mut self) -> anyhow::Result<()> {
if !std::mem::take(&mut self.state_save_armed) {
return Ok(());
}
Trader::save_state(&self.trader)
}
#[must_use]
pub fn event_store(&self) -> Option<&dyn KernelEventStore> {
self.event_store.as_deref()
}
#[must_use]
pub fn is_event_store_replay(&self) -> bool {
self.event_store_replay
}
#[must_use]
pub fn is_event_store_replay_configured(&self) -> bool {
self.event_store
.as_deref()
.is_some_and(KernelEventStore::is_event_store_replay_configured)
}
pub fn reset(&mut self) {
disarm_shutdown_on_error();
log::info!("Resetting");
if let Err(e) = self.trader.borrow_mut().reset() {
log::error!("Error resetting trader: {e:?}");
}
self.data_engine.borrow_mut().reset();
self.exec_engine.borrow_mut().reset();
self.risk_engine.borrow_mut().reset();
self.order_emulator.reset();
self.portfolio.borrow_mut().reset();
self.ts_started = None;
self.ts_shutdown = None;
self.state_save_armed = false;
log::info!("Reset");
}
pub fn dispose(&mut self) {
disarm_shutdown_on_error();
log::info!("Disposing");
let trader_state = self.trader.borrow().state();
match trader_state {
ComponentState::Running => self.stop_trader(),
ComponentState::Starting => {
if let Err(e) = self.stop_trader_after_start_failure() {
log::error!("Error stopping partial trader start during disposal: {e:?}");
}
}
_ => {}
}
if let Err(e) = self.save_trader_state() {
log::error!("Error saving trader state during disposal: {e:?}");
}
{
let mut trader = self.trader.borrow_mut();
if trader.state() == ComponentState::PreInitialized
&& let Err(e) = trader.initialize()
{
log::error!("Error initializing trader for disposal: {e:?}");
}
if !trader.is_disposed()
&& let Err(e) = trader.dispose()
{
log::error!("Error disposing trader: {e:?}");
}
}
self.stop_engines();
self.portfolio.borrow_mut().reset();
self.cancel_timers();
if let Some(event_store) = self.event_store.as_deref_mut() {
self.exec_engine.borrow_mut().set_snapshot_anchorer(None);
let ts_dispose = self.clock.borrow().timestamp_ns();
event_store.seal(ts_dispose);
}
self.data_engine.borrow_mut().dispose();
self.exec_engine.borrow_mut().dispose();
self.risk_engine.borrow_mut().dispose();
self.order_emulator.dispose();
self.cache.borrow_mut().dispose();
get_message_bus().borrow_mut().dispose();
log::info!("Disposed");
}
fn start_engines(&self) {
self.data_engine.borrow_mut().start();
self.exec_engine.borrow_mut().start();
self.risk_engine.borrow_mut().start();
}
fn stop_engines(&self) {
self.data_engine.borrow_mut().stop();
self.exec_engine.borrow_mut().stop();
self.risk_engine.borrow_mut().stop();
self.order_emulator.stop();
}
#[expect(clippy::await_holding_refcell_ref)] pub async fn connect_data_clients(&mut self) {
log::info!("Connecting data clients...");
self.data_engine.borrow_mut().connect().await;
}
#[expect(clippy::await_holding_refcell_ref)] pub async fn connect_exec_clients(&mut self) {
log::info!("Connecting execution clients...");
self.exec_engine.borrow_mut().connect().await;
}
#[expect(clippy::await_holding_refcell_ref)] pub async fn disconnect_clients(&mut self) -> anyhow::Result<()> {
log::info!("Disconnecting clients...");
let mut data_engine = self.data_engine.borrow_mut();
let mut exec_engine = self.exec_engine.borrow_mut();
let (data_result, exec_result) =
futures::join!(data_engine.disconnect(), exec_engine.disconnect());
match (data_result, exec_result) {
(Ok(()), Ok(())) => Ok(()),
(Err(data_err), Ok(())) => Err(data_err),
(Ok(()), Err(exec_err)) => Err(exec_err),
(Err(data_err), Err(exec_err)) => anyhow::bail!(
"Failed to disconnect data clients: {data_err}; failed to disconnect execution \
clients: {exec_err}"
),
}
}
#[must_use]
pub fn check_engines_connected(&self) -> bool {
self.data_engine.borrow().check_connected() && self.exec_engine.borrow().check_connected()
}
#[must_use]
pub fn check_engines_disconnected(&self) -> bool {
self.data_engine.borrow().check_disconnected()
&& self.exec_engine.borrow().check_disconnected()
}
#[must_use]
pub fn data_client_connection_status(&self) -> Vec<(ClientId, bool)> {
self.data_engine.borrow().client_connection_status()
}
#[must_use]
pub fn exec_client_connection_status(&self) -> Vec<(ClientId, bool)> {
self.exec_engine.borrow().client_connection_status()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum SocketReconnectDispatchOutcome {
Accepted,
AlreadyReconnecting,
Disconnected,
Closed,
Unsupported,
UnknownClient,
UnknownEndpoint,
AmbiguousEndpoint,
InvalidTrader,
}
#[cfg(test)]
mod socket_reconnect_tests {
use std::sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
};
use nautilus_common::clients::{SocketReconnectHandle, SocketReconnectRegistry};
use rstest::rstest;
use super::*;
fn stateful_handle(count: Arc<AtomicUsize>) -> SocketReconnectHandle {
let reconnecting = AtomicBool::new(false);
SocketReconnectHandle::new(move || {
if reconnecting.swap(true, Ordering::SeqCst) {
SocketReconnectRequestOutcome::AlreadyReconnecting
} else {
count.fetch_add(1, Ordering::SeqCst);
SocketReconnectRequestOutcome::Accepted
}
})
}
#[rstest]
fn endpoint_request_and_duplicate_leave_sibling_untouched() {
let registry = SocketReconnectRegistry::default();
let selected = Arc::new(AtomicUsize::new(0));
let sibling = Arc::new(AtomicUsize::new(0));
let _selected_registration = registry.register(
Ustr::from("market-0"),
stateful_handle(Arc::clone(&selected)),
);
let _sibling_registration = registry.register(
Ustr::from("market-1"),
stateful_handle(Arc::clone(&sibling)),
);
let selected_lookup = || {
SocketReconnectLookup::Handle(
registry
.get(Ustr::from("market-0"))
.expect("selected endpoint should be registered"),
)
};
assert_eq!(
NautilusKernel::request_socket_reconnect(
selected_lookup(),
SocketReconnectLookup::ClientNotFound,
),
SocketReconnectDispatchOutcome::Accepted,
);
assert_eq!(
NautilusKernel::request_socket_reconnect(
selected_lookup(),
SocketReconnectLookup::ClientNotFound,
),
SocketReconnectDispatchOutcome::AlreadyReconnecting,
);
assert_eq!(selected.load(Ordering::SeqCst), 1);
assert_eq!(sibling.load(Ordering::SeqCst), 0);
}
#[rstest]
#[case(
SocketReconnectLookup::ClientNotFound,
SocketReconnectLookup::ClientNotFound,
SocketReconnectDispatchOutcome::UnknownClient
)]
#[case(
SocketReconnectLookup::Unsupported,
SocketReconnectLookup::ClientNotFound,
SocketReconnectDispatchOutcome::Unsupported
)]
#[case(
SocketReconnectLookup::EndpointNotFound,
SocketReconnectLookup::ClientNotFound,
SocketReconnectDispatchOutcome::UnknownEndpoint
)]
fn rejected_lookup_does_not_invoke_an_endpoint(
#[case] data: SocketReconnectLookup,
#[case] execution: SocketReconnectLookup,
#[case] expected: SocketReconnectDispatchOutcome,
) {
assert_eq!(
NautilusKernel::request_socket_reconnect(data, execution),
expected,
);
}
#[rstest]
fn ambiguous_lookup_does_not_invoke_either_endpoint() {
let data_count = Arc::new(AtomicUsize::new(0));
let execution_count = Arc::new(AtomicUsize::new(0));
assert_eq!(
NautilusKernel::request_socket_reconnect(
SocketReconnectLookup::Handle(stateful_handle(Arc::clone(&data_count))),
SocketReconnectLookup::Handle(stateful_handle(Arc::clone(&execution_count))),
),
SocketReconnectDispatchOutcome::AmbiguousEndpoint,
);
assert_eq!(data_count.load(Ordering::SeqCst), 0);
assert_eq!(execution_count.load(Ordering::SeqCst), 0);
}
}
#[cfg(all(test, feature = "python"))]
mod tests {
use nautilus_common::messages::system::ShutdownSystem;
use nautilus_core::UUID4;
use rstest::*;
use ustr::Ustr;
use super::*;
use crate::builder::NautilusKernelBuilder;
#[rstest]
fn test_shutdown_system_sets_kernel_flag() {
let kernel = NautilusKernelBuilder::default().build().unwrap();
assert!(!kernel.is_shutdown_requested());
let command = ShutdownSystem::new(
kernel.trader_id(),
Ustr::from("TestComponent"),
Some("unit test".to_string()),
UUID4::new(),
kernel.generate_timestamp_ns(),
None, );
msgbus::publish_any(
MessagingSwitchboard::shutdown_system_topic(),
command.as_any(),
);
assert!(kernel.is_shutdown_requested());
kernel.reset_shutdown_flag();
assert!(!kernel.is_shutdown_requested());
}
#[rstest]
fn test_shutdown_system_idempotent() {
let kernel = NautilusKernelBuilder::default().build().unwrap();
let make_cmd = || {
ShutdownSystem::new(
kernel.trader_id(),
Ustr::from("TestComponent"),
None,
UUID4::new(),
kernel.generate_timestamp_ns(),
None, )
};
let topic = MessagingSwitchboard::shutdown_system_topic();
msgbus::publish_any(topic, make_cmd().as_any());
assert!(kernel.is_shutdown_requested());
msgbus::publish_any(topic, make_cmd().as_any());
assert!(kernel.is_shutdown_requested());
kernel.reset_shutdown_flag();
assert!(!kernel.is_shutdown_requested());
msgbus::publish_any(topic, make_cmd().as_any());
assert!(kernel.is_shutdown_requested());
}
#[rstest]
fn test_shutdown_system_ignores_other_trader() {
let kernel = NautilusKernelBuilder::default().build().unwrap();
let command = ShutdownSystem::new(
TraderId::from("OTHER-TRADER"),
Ustr::from("TestComponent"),
None,
UUID4::new(),
kernel.generate_timestamp_ns(),
None, );
msgbus::publish_any(
MessagingSwitchboard::shutdown_system_topic(),
command.as_any(),
);
assert!(!kernel.is_shutdown_requested());
}
}
#[cfg(test)]
mod lifecycle_tests {
use futures::FutureExt;
use indexmap::IndexMap;
use nautilus_common::{
actor::registry::get_actor_unchecked,
cache::Cache,
messages::data::{DataCommand, SubscribeCommand, UnsubscribeCommand},
msgbus::stubs::{TypedIntoMessageSavingHandler, get_typed_into_message_saving_handler},
};
use nautilus_execution::engine::SnapshotAnchorer;
use nautilus_model::{
enums::{OrderSide, OrderStatus, OrderType, TriggerType},
identifiers::{ActorId, ClientOrderId, StrategyId},
instruments::{
CryptoPerpetual, Instrument, InstrumentAny, stubs::crypto_perpetual_ethusdt,
},
orders::{Order, OrderAny, OrderTestBuilder},
types::{Price, Quantity},
};
use nautilus_testkit::{
cache::TestCacheDatabaseControl,
components::{StateActor, StateStrategy},
};
use rstest::rstest;
use ustr::Ustr;
use super::*;
use crate::{
builder::NautilusKernelBuilder,
event_store::{KernelEventStore, RegisteredComponents},
};
#[derive(Debug)]
struct RecordingEventStore {
control: TestCacheDatabaseControl,
opened: bool,
}
impl KernelEventStore for RecordingEventStore {
fn restore_parent_cache(
&mut self,
_instance_id: UUID4,
_cache: &mut Cache,
) -> anyhow::Result<()> {
self.control.record("event_store.restore");
Ok(())
}
fn open(
&mut self,
_instance_id: UUID4,
_components: &RegisteredComponents,
_environment: Environment,
) -> anyhow::Result<()> {
self.control.record("event_store.open");
self.opened = true;
Ok(())
}
fn snapshot_anchorer(&self) -> Option<SnapshotAnchorer> {
None
}
fn seal(&mut self, _ts_init: UnixNanos) {
if self.opened {
self.control.record("event_store.seal");
self.opened = false;
}
}
fn run_id(&self) -> Option<&str> {
None
}
fn parent_run_id(&self) -> Option<&str> {
None
}
fn is_halted(&self) -> bool {
false
}
}
fn state(key: &str, value: &[u8]) -> IndexMap<String, Vec<u8>> {
IndexMap::from([(key.to_string(), value.to_vec())])
}
fn finalize(kernel: &mut NautilusKernel) -> anyhow::Result<()> {
kernel
.finalize_stop()
.now_or_never()
.expect("kernel finalization must not yield")
}
fn add_state_components(
kernel: &NautilusKernel,
control: &TestCacheDatabaseControl,
actor: StateActor,
strategy: StateStrategy,
) {
kernel.trader.borrow_mut().add_actor(actor).unwrap();
kernel.trader.borrow_mut().add_strategy(strategy).unwrap();
control.record("components.registered");
}
fn create_stop_market_order(instrument: &CryptoPerpetual, client_order_id: &str) -> OrderAny {
OrderTestBuilder::new(OrderType::StopMarket)
.instrument_id(instrument.id())
.client_order_id(ClientOrderId::from(client_order_id))
.side(OrderSide::Buy)
.trigger_price(Price::from("5100.00"))
.quantity(Quantity::from(1))
.emulation_trigger(TriggerType::BidAsk)
.build()
}
fn register_data_command_handler(id: &str) -> TypedIntoMessageSavingHandler<DataCommand> {
let (handler, saving_handler) =
get_typed_into_message_saving_handler::<DataCommand>(Some(Ustr::from(id)));
msgbus::register_data_command_endpoint(
MessagingSwitchboard::data_engine_queue_execute(),
handler,
);
saving_handler
}
#[rstest]
fn test_state_persistence_orders_restore_load_start_stop_save_seal_and_dispose() {
let actor_id = ActorId::from("STATE-ACTOR");
let strategy_id = StrategyId::from("STATE-STRATEGY-001");
let actor_load = state("actor-loaded", b"actor-load-value");
let strategy_load = state("strategy-loaded", b"strategy-load-value");
let actor_save = state("actor-saved", b"actor-save-value");
let strategy_save = state("strategy-saved", b"strategy-save-value");
let (database, control) = TestCacheDatabaseControl::create();
control.set_actor_state(actor_id, &actor_load);
control.set_strategy_state(strategy_id, &strategy_load);
let event_store_control = control.clone();
let mut kernel = NautilusKernelBuilder::default()
.with_cache_database(Box::new(database))
.with_event_store(move |_instance_id, _clock| {
Ok(Box::new(RecordingEventStore {
control: event_store_control,
opened: false,
}))
})
.build()
.unwrap();
let actor = StateActor::new(actor_id, control.clone(), actor_save.clone());
let strategy = StateStrategy::new(strategy_id, control.clone(), strategy_save.clone());
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
kernel.start_trader().unwrap();
let actor_state = get_actor_unchecked::<StateActor>(&actor_id.inner())
.state_load()
.cloned();
let strategy_state = get_actor_unchecked::<StateStrategy>(&strategy_id.inner())
.state_load()
.cloned();
assert_eq!(actor_state, Some(actor_load));
assert_eq!(strategy_state, Some(strategy_load));
kernel.stop_trader();
kernel.stop_trader();
finalize(&mut kernel).unwrap();
finalize(&mut kernel).unwrap();
kernel.dispose();
assert_eq!(
control.events(),
vec![
"components.registered",
"event_store.restore",
"event_store.open",
"actor.load:STATE-ACTOR",
"actor.on_load",
"strategy.load:STATE-STRATEGY-001",
"strategy.on_load",
"actor.on_start",
"strategy.on_start",
"actor.on_stop",
"strategy.on_stop",
"actor.on_save",
"actor.update:STATE-ACTOR",
"strategy.on_save",
"strategy.update:STATE-STRATEGY-001",
"event_store.seal",
"database.close",
]
);
assert_eq!(control.actor_state(&actor_id), Some(actor_save));
assert_eq!(control.strategy_state(&strategy_id), Some(strategy_save));
}
#[rstest]
fn test_state_persistence_skips_callbacks_without_cache_backing() {
let actor_id = ActorId::from("NO-BACKING-ACTOR");
let strategy_id = StrategyId::from("NO-BACKING-STRATEGY-001");
let control = TestCacheDatabaseControl::default();
let mut kernel = NautilusKernelBuilder::default().build().unwrap();
let actor = StateActor::new(actor_id, control.clone(), state("actor", b"save"));
let strategy = StateStrategy::new(strategy_id, control.clone(), state("strategy", b"save"));
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
kernel.start_trader().unwrap();
kernel.stop_trader();
finalize(&mut kernel).unwrap();
kernel.dispose();
assert_eq!(
control.events(),
vec![
"components.registered",
"actor.on_start",
"strategy.on_start",
"actor.on_stop",
"strategy.on_stop",
]
);
}
#[rstest]
fn test_state_persistence_skips_empty_load_and_persists_empty_save() {
let actor_id = ActorId::from("EMPTY-STATE-ACTOR");
let strategy_id = StrategyId::from("EMPTY-STATE-STRATEGY-001");
let (database, control) = TestCacheDatabaseControl::create();
let mut kernel = NautilusKernelBuilder::default()
.with_cache_database(Box::new(database))
.build()
.unwrap();
let actor = StateActor::new(actor_id, control.clone(), IndexMap::new());
let strategy = StateStrategy::new(strategy_id, control.clone(), IndexMap::new());
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
kernel.start_trader().unwrap();
kernel.stop_trader();
finalize(&mut kernel).unwrap();
assert_eq!(
control.events(),
vec![
"components.registered",
"actor.load:EMPTY-STATE-ACTOR",
"strategy.load:EMPTY-STATE-STRATEGY-001",
"actor.on_start",
"strategy.on_start",
"actor.on_stop",
"strategy.on_stop",
"actor.on_save",
"actor.update:EMPTY-STATE-ACTOR",
"strategy.on_save",
"strategy.update:EMPTY-STATE-STRATEGY-001",
]
);
assert_eq!(control.actor_state(&actor_id), Some(IndexMap::new()));
assert_eq!(control.strategy_state(&strategy_id), Some(IndexMap::new()));
kernel.dispose();
}
#[rstest]
fn test_state_save_reports_all_callback_errors_and_continues_shutdown() {
let actor_id = ActorId::from("FAIL-SAVE-ACTOR");
let strategy_id = StrategyId::from("FAIL-SAVE-STRATEGY-001");
let (database, control) = TestCacheDatabaseControl::create();
let event_store_control = control.clone();
let mut kernel = NautilusKernelBuilder::default()
.with_cache_database(Box::new(database))
.with_event_store(move |_instance_id, _clock| {
Ok(Box::new(RecordingEventStore {
control: event_store_control,
opened: false,
}))
})
.build()
.unwrap();
let actor = StateActor::new(actor_id, control.clone(), IndexMap::new()).with_fail_save();
let strategy =
StateStrategy::new(strategy_id, control.clone(), IndexMap::new()).with_fail_save();
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
kernel.start_trader().unwrap();
kernel.stop_trader();
let expected_shutdown = kernel.clock.borrow().timestamp_ns();
let error = finalize(&mut kernel).unwrap_err();
kernel.dispose();
assert_eq!(
error.to_string(),
"Failed to save component state: actor FAIL-SAVE-ACTOR callback: test actor on_save \
failure; strategy FAIL-SAVE-STRATEGY-001 callback: test strategy on_save failure"
);
assert_eq!(kernel.ts_shutdown, Some(expected_shutdown));
assert_eq!(
control.events(),
vec![
"components.registered",
"event_store.restore",
"event_store.open",
"actor.load:FAIL-SAVE-ACTOR",
"strategy.load:FAIL-SAVE-STRATEGY-001",
"actor.on_start",
"strategy.on_start",
"actor.on_stop",
"strategy.on_stop",
"actor.on_save",
"strategy.on_save",
"event_store.seal",
"database.close",
]
);
}
#[rstest]
fn test_state_load_callback_failure_prevents_start_and_save() {
let actor_id = ActorId::from("FAIL-LOAD-ACTOR");
let strategy_id = StrategyId::from("FAIL-LOAD-STRATEGY-001");
let (database, control) = TestCacheDatabaseControl::create();
control.set_actor_state(actor_id, &state("actor", b"load"));
control.set_strategy_state(strategy_id, &state("strategy", b"load"));
let mut kernel = NautilusKernelBuilder::default()
.with_cache_database(Box::new(database))
.build()
.unwrap();
let actor = StateActor::new(actor_id, control.clone(), IndexMap::new()).with_fail_load();
let strategy = StateStrategy::new(strategy_id, control.clone(), IndexMap::new());
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
let error = kernel.start_trader().unwrap_err();
kernel.dispose();
assert_eq!(
error.to_string(),
"Failed to load actor and strategy state: Failed to restore actor FAIL-LOAD-ACTOR \
state: test actor on_load failure"
);
assert_eq!(
control.events(),
vec![
"components.registered",
"actor.load:FAIL-LOAD-ACTOR",
"actor.on_load",
"database.close",
]
);
}
#[rstest]
fn test_state_save_reports_all_persistence_errors() {
let actor_id = ActorId::from("FAIL-UPDATE-ACTOR");
let strategy_id = StrategyId::from("FAIL-UPDATE-STRATEGY-001");
let (database, control) = TestCacheDatabaseControl::create();
control.set_fail_update_actor(true);
control.set_fail_update_strategy(true);
let mut kernel = NautilusKernelBuilder::default()
.with_cache_database(Box::new(database))
.build()
.unwrap();
let actor = StateActor::new(actor_id, control.clone(), state("actor", b"save"));
let strategy = StateStrategy::new(strategy_id, control.clone(), state("strategy", b"save"));
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
kernel.start_trader().unwrap();
kernel.stop_trader();
let error = finalize(&mut kernel).unwrap_err();
kernel.dispose();
assert_eq!(
error.to_string(),
"Failed to save component state: actor FAIL-UPDATE-ACTOR persistence: test actor \
update failure; strategy FAIL-UPDATE-STRATEGY-001 persistence: test strategy update \
failure"
);
assert_eq!(
control.events(),
vec![
"components.registered",
"actor.load:FAIL-UPDATE-ACTOR",
"strategy.load:FAIL-UPDATE-STRATEGY-001",
"actor.on_start",
"strategy.on_start",
"actor.on_stop",
"strategy.on_stop",
"actor.on_save",
"actor.update:FAIL-UPDATE-ACTOR",
"strategy.on_save",
"strategy.update:FAIL-UPDATE-STRATEGY-001",
"database.close",
]
);
}
#[rstest]
fn test_partial_startup_stops_and_saves_once() {
let actor_id = ActorId::from("PARTIAL-ACTOR");
let strategy_id = StrategyId::from("PARTIAL-STRATEGY-001");
let (database, control) = TestCacheDatabaseControl::create();
let mut kernel = NautilusKernelBuilder::default()
.with_cache_database(Box::new(database))
.build()
.unwrap();
let actor = StateActor::new(actor_id, control.clone(), state("actor", b"partial"));
let strategy =
StateStrategy::new(strategy_id, control.clone(), state("strategy", b"partial"))
.with_fail_start();
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
let error = kernel.start_trader().unwrap_err();
kernel.dispose();
assert_eq!(
error.to_string(),
"Failed to start trader: test strategy on_start failure"
);
assert_eq!(
control.events(),
vec![
"components.registered",
"actor.load:PARTIAL-ACTOR",
"strategy.load:PARTIAL-STRATEGY-001",
"actor.on_start",
"strategy.on_start",
"actor.on_stop",
"strategy.on_stop",
"actor.on_save",
"actor.update:PARTIAL-ACTOR",
"strategy.on_save",
"strategy.update:PARTIAL-STRATEGY-001",
"database.close",
]
);
assert_eq!(
control.actor_state(&actor_id),
Some(state("actor", b"partial"))
);
assert_eq!(
control.strategy_state(&strategy_id),
Some(state("strategy", b"partial"))
);
}
#[rstest]
fn test_forced_dispose_stops_and_saves_once() {
let actor_id = ActorId::from("FORCED-ACTOR");
let strategy_id = StrategyId::from("FORCED-STRATEGY-001");
let (database, control) = TestCacheDatabaseControl::create();
let mut kernel = NautilusKernelBuilder::default()
.with_cache_database(Box::new(database))
.build()
.unwrap();
let actor = StateActor::new(actor_id, control.clone(), state("actor", b"forced"));
let strategy =
StateStrategy::new(strategy_id, control.clone(), state("strategy", b"forced"));
add_state_components(&kernel, &control, actor, strategy);
kernel.start();
kernel.start_trader().unwrap();
kernel.dispose();
assert_eq!(
control.events(),
vec![
"components.registered",
"actor.load:FORCED-ACTOR",
"strategy.load:FORCED-STRATEGY-001",
"actor.on_start",
"strategy.on_start",
"actor.on_stop",
"strategy.on_stop",
"actor.on_save",
"actor.update:FORCED-ACTOR",
"strategy.on_save",
"strategy.update:FORCED-STRATEGY-001",
"database.close",
]
);
}
#[rstest]
fn test_start_trader_starts_order_emulator_for_cached_emulated_orders() {
let mut kernel = NautilusKernelBuilder::default().build().unwrap();
let data_commands = register_data_command_handler("DataEngine.queue_execute.kernel_start");
let instrument = crypto_perpetual_ethusdt();
let instrument_id = instrument.id();
let first_order = create_stop_market_order(&instrument, "O-KERNEL-001");
let second_order = create_stop_market_order(&instrument, "O-KERNEL-002");
let first_client_order_id = first_order.client_order_id();
let second_client_order_id = second_order.client_order_id();
kernel
.cache
.borrow_mut()
.add_instrument(InstrumentAny::CryptoPerpetual(instrument))
.unwrap();
kernel
.cache
.borrow_mut()
.add_order(first_order, None, None, false)
.unwrap();
kernel
.cache
.borrow_mut()
.add_order(second_order, None, None, false)
.unwrap();
kernel.start();
assert!(
kernel
.order_emulator
.get_emulator()
.get_matching_core(&instrument_id)
.is_none()
);
kernel.start_trader().unwrap();
let commands = data_commands.get_messages();
let cache = kernel.cache.borrow();
let first_status = cache.order(&first_client_order_id).unwrap().status();
let second_status = cache.order(&second_client_order_id).unwrap().status();
drop(cache);
let emulator = kernel.order_emulator.get_emulator();
assert!(emulator.get_matching_core(&instrument_id).is_some());
assert_eq!(emulator.subscribed_quotes(), vec![instrument_id]);
assert_eq!(first_status, OrderStatus::Emulated);
assert_eq!(second_status, OrderStatus::Emulated);
assert!(commands.iter().any(|command| matches!(
command,
DataCommand::Subscribe(SubscribeCommand::Quotes(command))
if command.instrument_id == instrument_id
)));
data_commands.clear();
drop(emulator);
kernel.stop_trader();
kernel.dispose();
let commands = data_commands.get_messages();
let emulator = kernel.order_emulator.get_emulator();
assert!(emulator.subscribed_quotes().is_empty());
assert!(emulator.get_matching_core(&instrument_id).is_none());
assert!(commands.iter().any(|command| matches!(
command,
DataCommand::Unsubscribe(UnsubscribeCommand::Quotes(command))
if command.instrument_id == instrument_id
)));
}
#[rstest]
fn test_reset_resets_order_emulator_state() {
let mut kernel = NautilusKernelBuilder::default().build().unwrap();
let data_commands = register_data_command_handler("DataEngine.queue_execute.kernel_reset");
let instrument = crypto_perpetual_ethusdt();
let instrument_id = instrument.id();
let order = create_stop_market_order(&instrument, "O-KERNEL-RESET-001");
kernel
.cache
.borrow_mut()
.add_instrument(InstrumentAny::CryptoPerpetual(instrument))
.unwrap();
kernel
.cache
.borrow_mut()
.add_order(order, None, None, false)
.unwrap();
kernel.start();
kernel.start_trader().unwrap();
assert!(
kernel
.order_emulator
.get_emulator()
.get_matching_core(&instrument_id)
.is_some()
);
kernel.stop_trader();
data_commands.clear();
kernel.reset();
let commands = data_commands.get_messages();
let emulator = kernel.order_emulator.get_emulator();
assert!(emulator.subscribed_quotes().is_empty());
assert!(emulator.get_matching_core(&instrument_id).is_none());
assert!(commands.iter().any(|command| matches!(
command,
DataCommand::Unsubscribe(UnsubscribeCommand::Quotes(command))
if command.instrument_id == instrument_id
)));
drop(emulator);
kernel.dispose();
}
}