Skip to main content

nautilus_persistence/writer/
subscription.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//! Message-bus subscriptions for streaming writer sinks.
17
18use std::{
19    any::Any,
20    cell::RefCell,
21    fmt::Debug,
22    rc::Rc,
23    sync::{
24        Arc,
25        atomic::{AtomicU64, Ordering},
26    },
27};
28
29use nautilus_common::{
30    clock::Clock,
31    msgbus::{
32        MStr, ShareableMessageHandler, TypedHandler, subscribe_account_state, subscribe_any,
33        subscribe_bars, subscribe_book_deltas, subscribe_book_depth as subscribe_book_depths,
34        subscribe_funding_rates, subscribe_index_prices, subscribe_instruments,
35        subscribe_mark_prices, subscribe_option_greeks, subscribe_order_events,
36        subscribe_position_events, subscribe_quotes, subscribe_trades, unsubscribe_account_state,
37        unsubscribe_any, unsubscribe_bars, unsubscribe_book_deltas,
38        unsubscribe_book_depth as unsubscribe_book_depths, unsubscribe_funding_rates,
39        unsubscribe_index_prices, unsubscribe_instruments, unsubscribe_mark_prices,
40        unsubscribe_option_greeks, unsubscribe_order_events, unsubscribe_position_events,
41        unsubscribe_quotes, unsubscribe_trades,
42    },
43};
44use nautilus_model::{
45    data::{
46        Bar, CustomData, FundingRateUpdate, IndexPriceUpdate, MarkPriceUpdate, OptionGreeks,
47        OrderBookDeltas, OrderBookDepth, QuoteTick, TradeTick,
48    },
49    events::{AccountState, OrderEventAny, PositionEvent},
50    instruments::{Instrument, InstrumentAny},
51};
52
53use super::traits::StreamingDataSink;
54
55type ClockBridge = Option<(Rc<RefCell<dyn Clock>>, Arc<AtomicU64>)>;
56
57/// Message-bus subscriptions forwarding supported messages into a streaming sink.
58pub struct StreamingSinkSubscription {
59    sink: Rc<RefCell<StreamingDataSink>>,
60    clock_bridge: ClockBridge,
61    any_handler: ShareableMessageHandler,
62    quotes_handler: TypedHandler<QuoteTick>,
63    trades_handler: TypedHandler<TradeTick>,
64    bars_handler: TypedHandler<Bar>,
65    deltas_handler: TypedHandler<OrderBookDeltas>,
66    depths_handler: TypedHandler<OrderBookDepth>,
67    mark_prices_handler: TypedHandler<MarkPriceUpdate>,
68    index_prices_handler: TypedHandler<IndexPriceUpdate>,
69    funding_rates_handler: TypedHandler<FundingRateUpdate>,
70    option_greeks_handler: TypedHandler<OptionGreeks>,
71    instruments_handler: TypedHandler<InstrumentAny>,
72    account_state_handler: TypedHandler<AccountState>,
73    order_events_handler: TypedHandler<OrderEventAny>,
74    position_events_handler: TypedHandler<PositionEvent>,
75}
76
77impl Debug for StreamingSinkSubscription {
78    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
79        f.debug_struct(stringify!(StreamingSinkSubscription))
80            .finish_non_exhaustive()
81    }
82}
83
84impl StreamingSinkSubscription {
85    /// Subscribes a streaming sink to typed and dynamic message-bus routes.
86    pub fn subscribe(sink: Rc<RefCell<StreamingDataSink>>, clock_bridge: ClockBridge) -> Self {
87        Self::subscribe_named(sink, clock_bridge, "streaming writer".to_string())
88    }
89
90    /// Subscribes a sink and includes its name in write-failure diagnostics.
91    pub fn subscribe_named(
92        sink: Rc<RefCell<StreamingDataSink>>,
93        clock_bridge: ClockBridge,
94        name: String,
95    ) -> Self {
96        let name: Rc<str> = Rc::from(name);
97        let any_handler = any_handler(Rc::clone(&sink), clock_bridge.clone(), name.clone());
98        let quotes_handler =
99            typed_handler::<QuoteTick>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
100        let trades_handler =
101            typed_handler::<TradeTick>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
102        let bars_handler =
103            typed_handler::<Bar>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
104        let deltas_handler =
105            typed_handler::<OrderBookDeltas>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
106        let depths_handler =
107            typed_handler::<OrderBookDepth>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
108        let mark_prices_handler =
109            typed_handler::<MarkPriceUpdate>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
110        let index_prices_handler =
111            typed_handler::<IndexPriceUpdate>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
112        let funding_rates_handler = typed_handler::<FundingRateUpdate>(
113            Rc::clone(&sink),
114            clock_bridge.clone(),
115            name.clone(),
116        );
117        let option_greeks_handler =
118            typed_handler::<OptionGreeks>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
119        let instruments_handler =
120            typed_handler::<InstrumentAny>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
121        let account_state_handler =
122            typed_handler::<AccountState>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
123        let order_events_handler =
124            typed_handler::<OrderEventAny>(Rc::clone(&sink), clock_bridge.clone(), name.clone());
125        let position_events_handler =
126            typed_handler::<PositionEvent>(Rc::clone(&sink), clock_bridge.clone(), name);
127        let pattern = MStr::pattern("*");
128        // Order events publish once on `events.order.{strategy_id}` and fills re-publish on
129        // `events.order_filled.{instrument_id}`; a `*` pattern captures both and duplicates
130        // every fill row, so the sink subscribes to the strategy-scoped topic only.
131        let order_events_pattern = MStr::pattern("events.order.*");
132
133        subscribe_any(pattern, any_handler.clone(), None);
134        subscribe_quotes(pattern, quotes_handler.clone(), None);
135        subscribe_trades(pattern, trades_handler.clone(), None);
136        subscribe_bars(pattern, bars_handler.clone(), None);
137        subscribe_book_deltas(pattern, deltas_handler.clone(), None);
138        subscribe_book_depths(pattern, depths_handler.clone(), None);
139        subscribe_mark_prices(pattern, mark_prices_handler.clone(), None);
140        subscribe_index_prices(pattern, index_prices_handler.clone(), None);
141        subscribe_funding_rates(pattern, funding_rates_handler.clone(), None);
142        subscribe_option_greeks(pattern, option_greeks_handler.clone(), None);
143        subscribe_instruments(pattern, instruments_handler.clone(), None);
144        subscribe_account_state(pattern, account_state_handler.clone(), None);
145        subscribe_order_events(order_events_pattern, order_events_handler.clone(), None);
146        subscribe_position_events(pattern, position_events_handler.clone(), None);
147
148        Self {
149            sink,
150            clock_bridge,
151            any_handler,
152            quotes_handler,
153            trades_handler,
154            bars_handler,
155            deltas_handler,
156            depths_handler,
157            mark_prices_handler,
158            index_prices_handler,
159            funding_rates_handler,
160            option_greeks_handler,
161            instruments_handler,
162            account_state_handler,
163            order_events_handler,
164            position_events_handler,
165        }
166    }
167
168    /// Unsubscribes and closes the sink.
169    /// # Errors
170    ///
171    /// Returns an error if the underlying sink cannot be closed.
172    pub fn close(&self) -> anyhow::Result<()> {
173        self.unsubscribe();
174        refresh_writer_clock(&self.clock_bridge);
175        self.sink.borrow_mut().close()
176    }
177
178    /// Removes all message-bus subscriptions without closing the sink.
179    pub fn unsubscribe(&self) {
180        let pattern = MStr::pattern("*");
181        let order_events_pattern = MStr::pattern("events.order.*");
182
183        unsubscribe_any(pattern, &self.any_handler);
184        unsubscribe_quotes(pattern, &self.quotes_handler);
185        unsubscribe_trades(pattern, &self.trades_handler);
186        unsubscribe_bars(pattern, &self.bars_handler);
187        unsubscribe_book_deltas(pattern, &self.deltas_handler);
188        unsubscribe_book_depths(pattern, &self.depths_handler);
189        unsubscribe_mark_prices(pattern, &self.mark_prices_handler);
190        unsubscribe_index_prices(pattern, &self.index_prices_handler);
191        unsubscribe_funding_rates(pattern, &self.funding_rates_handler);
192        unsubscribe_option_greeks(pattern, &self.option_greeks_handler);
193        unsubscribe_instruments(pattern, &self.instruments_handler);
194        unsubscribe_account_state(pattern, &self.account_state_handler);
195        unsubscribe_order_events(order_events_pattern, &self.order_events_handler);
196        unsubscribe_position_events(pattern, &self.position_events_handler);
197    }
198}
199
200fn typed_handler<T: 'static>(
201    sink: Rc<RefCell<StreamingDataSink>>,
202    clock_bridge: ClockBridge,
203    name: Rc<str>,
204) -> TypedHandler<T> {
205    TypedHandler::from(move |message: &T| {
206        write_bus_message(&sink, &clock_bridge, &name, message);
207    })
208}
209
210fn any_handler(
211    sink: Rc<RefCell<StreamingDataSink>>,
212    clock_bridge: ClockBridge,
213    name: Rc<str>,
214) -> ShareableMessageHandler {
215    ShareableMessageHandler::from_any(move |message: &dyn Any| {
216        write_bus_message(&sink, &clock_bridge, &name, message);
217    })
218}
219
220fn write_bus_message(
221    sink: &Rc<RefCell<StreamingDataSink>>,
222    clock_bridge: &ClockBridge,
223    name: &str,
224    message: &dyn Any,
225) {
226    refresh_writer_clock(clock_bridge);
227
228    if let Err(e) = sink.borrow_mut().write_any(message) {
229        log::warn!(
230            "Failed to write {name} {} message: {e}",
231            streaming_message_label(message)
232        );
233    }
234}
235
236fn refresh_writer_clock(clock_bridge: &ClockBridge) {
237    if let Some((clock, shared)) = clock_bridge {
238        shared.store(clock.borrow().timestamp_ns().as_u64(), Ordering::Relaxed);
239    }
240}
241
242fn streaming_message_label(message: &dyn Any) -> String {
243    if let Some(custom) = message.downcast_ref::<CustomData>() {
244        return format!(
245            "CustomData({}, identifier={:?})",
246            custom.data.type_name(),
247            custom.data_type.identifier()
248        );
249    }
250
251    if let Some(instrument) = message.downcast_ref::<InstrumentAny>() {
252        return format!("InstrumentAny({})", instrument.id());
253    }
254
255    if let Some(bar) = message.downcast_ref::<Bar>() {
256        return format!("Bar({})", bar.bar_type);
257    }
258
259    if let Some(quotes) = message.downcast_ref::<QuoteTick>() {
260        return format!("QuoteTick({})", quotes.instrument_id);
261    }
262
263    if let Some(trade) = message.downcast_ref::<TradeTick>() {
264        return format!("TradeTick({})", trade.instrument_id);
265    }
266
267    "unknown".to_string()
268}