Skip to main content

nautilus_common/msgbus/
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//! In-memory message bus for intra-process communication.
17//!
18//! # Messaging patterns
19//!
20//! - **Point-to-point**: Send messages to named endpoints via `send_*` functions.
21//! - **Pub/sub**: Publish messages to topics via `publish_*`, subscribers receive
22//!   all messages matching their pattern.
23//! - **Request/response**: Register correlation IDs for response sequence tracking.
24//!
25//! # Architecture
26//!
27//! The bus uses thread-local storage for single-threaded async runtimes. Each
28//! thread gets its own `MessageBus` instance, avoiding synchronization overhead.
29//!
30//! Two routing mechanisms serve different needs:
31//!
32//! - **Typed routing** (`publish_quote`, `subscribe_quotes`): Zero-cost dispatch
33//!   for known types. Handlers receive `&T` directly with no runtime type checking.
34//! - **Any-based routing** (`publish_any`, `subscribe_any`): Flexible dispatch for
35//!   custom types and Python interop. Handlers receive `&dyn Any`.
36//!
37//! See [`core`] module documentation for design decisions and performance details.
38
39pub mod config;
40pub mod core;
41pub mod matching;
42pub mod message;
43pub mod mstr;
44pub mod stubs;
45pub mod switchboard;
46pub mod typed_endpoints;
47pub mod typed_handler;
48pub mod typed_router;
49
50pub(crate) mod external;
51
52mod api;
53mod backing;
54
55use std::{
56    any::Any,
57    cell::{Cell, RefCell},
58    rc::Rc,
59};
60
61use nautilus_core::UUID4;
62#[cfg(feature = "defi")]
63use nautilus_model::defi::{Block, Pool, PoolFeeCollect, PoolFlash, PoolLiquidityUpdate, PoolSwap};
64use nautilus_model::{
65    data::{
66        Bar, FundingRateUpdate, GreeksData, IndexPriceUpdate, MarkPriceUpdate, OrderBookDeltas,
67        OrderBookDepth10, QuoteTick, TradeTick,
68        option_chain::{OptionChainSlice, OptionGreeks},
69    },
70    events::{AccountState, OrderEventAny, PortfolioSnapshot, PositionEvent},
71    instruments::InstrumentAny,
72    orderbook::OrderBook,
73};
74use smallvec::SmallVec;
75
76#[cfg(feature = "live")]
77pub use self::backing::{
78    MessageBusExternalIngress, MessageBusExternalReceiver, external_io_from_backing,
79};
80pub use self::{
81    api::*,
82    backing::{
83        MessageBusBacking, MessageBusBackingFactory, MessageBusExternalEgress,
84        external_egress_from_backing,
85    },
86    config::MessageBusConfig,
87    core::{MessageBus, Subscription},
88    message::{BusMessage, BusPayloadCategory, BusPayloadType},
89    mstr::{Endpoint, MStr, Pattern, Topic},
90    switchboard::MessagingSwitchboard,
91    typed_endpoints::{EndpointMap, IntoEndpointMap},
92    typed_handler::{
93        CallbackHandler, Handler, IntoHandler, ShareableMessageHandler, TypedHandler,
94        TypedIntoHandler,
95    },
96    typed_router::{TopicRouter, TypedSubscription},
97};
98use crate::timer::TimeEvent;
99
100/// Inline capacity for handler buffers before heap allocation.
101pub(super) const HANDLER_BUFFER_CAP: usize = 64;
102
103// MessageBus is designed for single-threaded use within each async runtime.
104// Thread-local storage ensures each thread gets its own instance, eliminating
105// the need for unsafe Send/Sync implementations.
106//
107// Handler buffers provide zero-allocation publish on hot paths.
108// Each buffer stores up to 64 handlers inline before spilling to heap.
109// Publish functions use move-out/move-back to avoid holding RefCell borrows
110// during handler calls (enabling re-entrant publishes).
111thread_local! {
112    pub(super) static MESSAGE_BUS: RefCell<Option<Rc<RefCell<MessageBus>>>> = const { RefCell::new(None) };
113    pub(super) static HAS_EXTERNAL_EGRESS: Cell<bool> = const { Cell::new(false) };
114    pub(super) static SUPPRESS_EXTERNAL_DEPTH: Cell<u32> = const { Cell::new(0) };
115
116    pub(super) static ANY_HANDLERS: RefCell<SmallVec<[ShareableMessageHandler; HANDLER_BUFFER_CAP]>> =
117        RefCell::new(SmallVec::new());
118
119    pub(super) static DELTAS_HANDLERS: RefCell<SmallVec<[TypedHandler<OrderBookDeltas>; HANDLER_BUFFER_CAP]>> =
120        RefCell::new(SmallVec::new());
121    pub(super) static DEPTH10_HANDLERS: RefCell<SmallVec<[TypedHandler<OrderBookDepth10>; HANDLER_BUFFER_CAP]>> =
122        RefCell::new(SmallVec::new());
123    pub(super) static BOOK_HANDLERS: RefCell<SmallVec<[TypedHandler<OrderBook>; HANDLER_BUFFER_CAP]>> =
124        RefCell::new(SmallVec::new());
125    pub(super) static QUOTE_HANDLERS: RefCell<SmallVec<[TypedHandler<QuoteTick>; HANDLER_BUFFER_CAP]>> =
126        RefCell::new(SmallVec::new());
127    pub(super) static TRADE_HANDLERS: RefCell<SmallVec<[TypedHandler<TradeTick>; HANDLER_BUFFER_CAP]>> =
128        RefCell::new(SmallVec::new());
129    pub(super) static BAR_HANDLERS: RefCell<SmallVec<[TypedHandler<Bar>; HANDLER_BUFFER_CAP]>> =
130        RefCell::new(SmallVec::new());
131    pub(super) static MARK_PRICE_HANDLERS: RefCell<SmallVec<[TypedHandler<MarkPriceUpdate>; HANDLER_BUFFER_CAP]>> =
132        RefCell::new(SmallVec::new());
133    pub(super) static INDEX_PRICE_HANDLERS: RefCell<SmallVec<[TypedHandler<IndexPriceUpdate>; HANDLER_BUFFER_CAP]>> =
134        RefCell::new(SmallVec::new());
135    pub(super) static FUNDING_RATE_HANDLERS: RefCell<SmallVec<[TypedHandler<FundingRateUpdate>; HANDLER_BUFFER_CAP]>> =
136        RefCell::new(SmallVec::new());
137    pub(super) static GREEKS_HANDLERS: RefCell<SmallVec<[TypedHandler<GreeksData>; HANDLER_BUFFER_CAP]>> =
138        RefCell::new(SmallVec::new());
139    pub(super) static OPTION_GREEKS_HANDLERS: RefCell<SmallVec<[TypedHandler<OptionGreeks>; HANDLER_BUFFER_CAP]>> =
140        RefCell::new(SmallVec::new());
141    pub(super) static OPTION_CHAIN_HANDLERS: RefCell<SmallVec<[TypedHandler<OptionChainSlice>; HANDLER_BUFFER_CAP]>> =
142        RefCell::new(SmallVec::new());
143    pub(super) static ACCOUNT_STATE_HANDLERS: RefCell<SmallVec<[TypedHandler<AccountState>; HANDLER_BUFFER_CAP]>> =
144        RefCell::new(SmallVec::new());
145    pub(super) static PORTFOLIO_SNAPSHOT_HANDLERS: RefCell<SmallVec<[TypedHandler<PortfolioSnapshot>; HANDLER_BUFFER_CAP]>> =
146        RefCell::new(SmallVec::new());
147    pub(super) static ORDER_EVENT_HANDLERS: RefCell<SmallVec<[TypedHandler<OrderEventAny>; HANDLER_BUFFER_CAP]>> =
148        RefCell::new(SmallVec::new());
149    pub(super) static POSITION_EVENT_HANDLERS: RefCell<SmallVec<[TypedHandler<PositionEvent>; HANDLER_BUFFER_CAP]>> =
150        RefCell::new(SmallVec::new());
151    pub(super) static INSTRUMENT_HANDLERS: RefCell<SmallVec<[TypedHandler<InstrumentAny>; HANDLER_BUFFER_CAP]>> =
152        RefCell::new(SmallVec::new());
153
154    #[cfg(feature = "defi")]
155    pub(super) static DEFI_BLOCK_HANDLERS: RefCell<SmallVec<[TypedHandler<Block>; HANDLER_BUFFER_CAP]>> =
156        RefCell::new(SmallVec::new());
157    #[cfg(feature = "defi")]
158    pub(super) static DEFI_POOL_HANDLERS: RefCell<SmallVec<[TypedHandler<Pool>; HANDLER_BUFFER_CAP]>> =
159        RefCell::new(SmallVec::new());
160    #[cfg(feature = "defi")]
161    pub(super) static DEFI_SWAP_HANDLERS: RefCell<SmallVec<[TypedHandler<PoolSwap>; HANDLER_BUFFER_CAP]>> =
162        RefCell::new(SmallVec::new());
163    #[cfg(feature = "defi")]
164    pub(super) static DEFI_LIQUIDITY_HANDLERS: RefCell<SmallVec<[TypedHandler<PoolLiquidityUpdate>; HANDLER_BUFFER_CAP]>> =
165        RefCell::new(SmallVec::new());
166    #[cfg(feature = "defi")]
167    pub(super) static DEFI_COLLECT_HANDLERS: RefCell<SmallVec<[TypedHandler<PoolFeeCollect>; HANDLER_BUFFER_CAP]>> =
168        RefCell::new(SmallVec::new());
169    #[cfg(feature = "defi")]
170    pub(super) static DEFI_FLASH_HANDLERS: RefCell<SmallVec<[TypedHandler<PoolFlash>; HANDLER_BUFFER_CAP]>> =
171        RefCell::new(SmallVec::new());
172}
173
174/// Guard that prevents republished external messages from being forwarded again.
175#[derive(Debug)]
176pub struct SuppressExternalGuard;
177
178impl SuppressExternalGuard {
179    #[must_use]
180    pub fn new() -> Self {
181        SUPPRESS_EXTERNAL_DEPTH.with(|depth| depth.set(depth.get().saturating_add(1)));
182        Self
183    }
184}
185
186impl Default for SuppressExternalGuard {
187    fn default() -> Self {
188        Self::new()
189    }
190}
191
192impl Drop for SuppressExternalGuard {
193    fn drop(&mut self) {
194        SUPPRESS_EXTERNAL_DEPTH.with(|depth| depth.set(depth.get().saturating_sub(1)));
195    }
196}
197
198/// Sets the thread-local message bus, replacing any existing one.
199pub fn set_message_bus(msgbus: Rc<RefCell<MessageBus>>) {
200    HAS_EXTERNAL_EGRESS.with(|flag| flag.set(msgbus.borrow().has_external_egress()));
201    MESSAGE_BUS.with(|bus| {
202        *bus.borrow_mut() = Some(msgbus);
203    });
204}
205
206/// Gets the thread-local message bus.
207///
208/// If no message bus has been set for this thread, a default one is created and initialized.
209pub fn get_message_bus() -> Rc<RefCell<MessageBus>> {
210    MESSAGE_BUS.with(|bus| {
211        let mut slot = bus.borrow_mut();
212        let rc = slot.get_or_insert_with(|| Rc::new(RefCell::new(MessageBus::default())));
213        rc.clone()
214    })
215}
216
217/// Returns the thread-local message bus if one has been set for this thread.
218pub fn try_get_message_bus() -> Option<Rc<RefCell<MessageBus>>> {
219    MESSAGE_BUS.with(|bus| bus.borrow().clone())
220}
221
222/// Observes dispatched bus traffic for the durable event store.
223///
224/// The bus invokes the registered tap (when present) before each publish, send, or
225/// correlation response fanout, so subscribers cannot observe a message that has not
226/// yet been handed to the tap. The tap callback runs on the engine thread and must be
227/// cheap; it must not re-enter the bus (the bus is single-threaded and the call site
228/// holds no live borrow of the bus, so any re-entrant publish would deadlock through
229/// downstream `RefCell::borrow_mut` calls inside the registered tap).
230pub trait BusTap: 'static {
231    /// Invoked before a publish fanout dispatches to subscribers on `topic`.
232    fn on_publish(&self, topic: MStr<Topic>, message: &dyn Any);
233
234    /// Invoked before a send dispatch reaches the endpoint handler.
235    fn on_send(&self, endpoint: MStr<Endpoint>, message: &dyn Any);
236
237    /// Invoked before a correlation response dispatch reaches the response handler.
238    fn on_response(&self, _correlation_id: &UUID4, _message: &dyn Any) {}
239}
240
241thread_local! {
242    pub(super) static BUS_TAP: RefCell<Option<Rc<dyn BusTap>>> = const { RefCell::new(None) };
243}
244
245/// Registers `tap` as the thread-local bus tap, replacing any previously installed tap.
246///
247/// The tap fires before each publish, send, and correlation response fanout. Callers
248/// are responsible for clearing the tap on shutdown via [`clear_bus_tap`] so a stale
249/// adapter does not outlive the writer it captures into.
250pub fn set_bus_tap(tap: Rc<dyn BusTap>) {
251    BUS_TAP.with(|slot| {
252        *slot.borrow_mut() = Some(tap);
253    });
254}
255
256/// Clears the registered bus tap on this thread.
257///
258/// A no-op when no tap is installed.
259pub fn clear_bus_tap() {
260    BUS_TAP.with(|slot| {
261        *slot.borrow_mut() = None;
262    });
263}
264
265#[inline]
266pub(super) fn dispatch_tap_publish(topic: MStr<Topic>, message: &dyn Any) {
267    // Clone the Rc so the cell borrow is released before the tap runs. The tap is
268    // single-threaded with the bus; a re-entrant `set_bus_tap` during dispatch would
269    // otherwise panic on RefCell.
270    let tap = BUS_TAP.with(|slot| slot.borrow().clone());
271    if let Some(tap) = tap {
272        tap.on_publish(topic, message);
273    }
274}
275
276#[inline]
277pub(super) fn dispatch_tap_send(endpoint: MStr<Endpoint>, message: &dyn Any) {
278    let tap = BUS_TAP.with(|slot| slot.borrow().clone());
279    if let Some(tap) = tap {
280        tap.on_send(endpoint, message);
281    }
282}
283
284#[inline]
285pub(super) fn dispatch_tap_response(correlation_id: &UUID4, message: &dyn Any) {
286    let tap = BUS_TAP.with(|slot| slot.borrow().clone());
287    if let Some(tap) = tap {
288        tap.on_response(correlation_id, message);
289    }
290}
291
292#[inline]
293pub(crate) fn dispatch_tap_time_event(event: &TimeEvent) {
294    dispatch_tap_publish(MessagingSwitchboard::time_event_topic(), event);
295}