1use 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
57pub 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 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 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 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 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 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}