1use ahash::AHashMap;
19use nautilus_common::messages::DataEvent;
20use nautilus_core::UnixNanos;
21use nautilus_model::{
22 data::InstrumentStatus, enums::MarketStatusAction, identifiers::InstrumentId,
23};
24
25use crate::spot::sbe::generated::symbol_status::SymbolStatus;
26
27impl From<SymbolStatus> for MarketStatusAction {
28 fn from(status: SymbolStatus) -> Self {
29 match status {
30 SymbolStatus::Trading => Self::Trading,
31 SymbolStatus::EndOfDay => Self::Close,
32 SymbolStatus::Halt => Self::Halt,
33 SymbolStatus::Break => Self::Pause,
34 SymbolStatus::CancelOnly => Self::Halt,
35 SymbolStatus::NonRepresentable | SymbolStatus::NullVal => Self::NotAvailableForTrading,
36 }
37 }
38}
39
40pub fn diff_and_emit_statuses(
46 new_statuses: &AHashMap<InstrumentId, MarketStatusAction>,
47 cached_statuses: &mut AHashMap<InstrumentId, MarketStatusAction>,
48 sender: &tokio::sync::mpsc::UnboundedSender<DataEvent>,
49 ts_event: UnixNanos,
50 ts_init: UnixNanos,
51) {
52 for (instrument_id, &new_action) in new_statuses {
53 let changed = cached_statuses
54 .get(instrument_id)
55 .is_none_or(|&prev| prev != new_action);
56
57 if changed {
58 cached_statuses.insert(*instrument_id, new_action);
59 emit_status(sender, *instrument_id, new_action, ts_event, ts_init);
60 }
61 }
62
63 let removed: Vec<InstrumentId> = cached_statuses
65 .keys()
66 .filter(|id| !new_statuses.contains_key(id))
67 .copied()
68 .collect();
69
70 for instrument_id in removed {
71 cached_statuses.remove(&instrument_id);
72 emit_status(
73 sender,
74 instrument_id,
75 MarketStatusAction::NotAvailableForTrading,
76 ts_event,
77 ts_init,
78 );
79 }
80}
81
82fn emit_status(
83 sender: &tokio::sync::mpsc::UnboundedSender<DataEvent>,
84 instrument_id: InstrumentId,
85 action: MarketStatusAction,
86 ts_event: UnixNanos,
87 ts_init: UnixNanos,
88) {
89 let is_trading = Some(matches!(action, MarketStatusAction::Trading));
90 let status = InstrumentStatus::new(
91 instrument_id,
92 action,
93 ts_event,
94 ts_init,
95 None,
96 None,
97 is_trading,
98 None,
99 None,
100 );
101
102 if let Err(e) = sender.send(DataEvent::InstrumentStatus(status)) {
103 log::error!("Failed to emit instrument status event: {e}");
104 }
105}
106
107#[cfg(test)]
108mod tests {
109 use nautilus_model::identifiers::InstrumentId;
110 use rstest::rstest;
111
112 use super::{
113 super::enums::{BinanceContractStatus, BinanceTradingStatus},
114 *,
115 };
116
117 #[rstest]
118 #[case(SymbolStatus::Trading, MarketStatusAction::Trading)]
119 #[case(SymbolStatus::EndOfDay, MarketStatusAction::Close)]
120 #[case(SymbolStatus::Halt, MarketStatusAction::Halt)]
121 #[case(SymbolStatus::Break, MarketStatusAction::Pause)]
122 #[case(SymbolStatus::CancelOnly, MarketStatusAction::Halt)]
123 #[case(
124 SymbolStatus::NonRepresentable,
125 MarketStatusAction::NotAvailableForTrading
126 )]
127 #[case(SymbolStatus::NullVal, MarketStatusAction::NotAvailableForTrading)]
128 fn test_symbol_status_to_market_action(
129 #[case] input: SymbolStatus,
130 #[case] expected: MarketStatusAction,
131 ) {
132 assert_eq!(MarketStatusAction::from(input), expected);
133 }
134
135 #[rstest]
136 #[case(BinanceTradingStatus::Trading, MarketStatusAction::Trading)]
137 #[case(BinanceTradingStatus::PendingTrading, MarketStatusAction::PreOpen)]
138 #[case(BinanceTradingStatus::PreTrading, MarketStatusAction::PreOpen)]
139 #[case(BinanceTradingStatus::PostTrading, MarketStatusAction::PostClose)]
140 #[case(BinanceTradingStatus::EndOfDay, MarketStatusAction::Close)]
141 #[case(BinanceTradingStatus::Halt, MarketStatusAction::Halt)]
142 #[case(BinanceTradingStatus::AuctionMatch, MarketStatusAction::Cross)]
143 #[case(BinanceTradingStatus::Break, MarketStatusAction::Pause)]
144 #[case(BinanceTradingStatus::PreDelivering, MarketStatusAction::PreClose)]
145 #[case(BinanceTradingStatus::Delivering, MarketStatusAction::Close)]
146 #[case(BinanceTradingStatus::Delivered, MarketStatusAction::Close)]
147 #[case(BinanceTradingStatus::PreSettle, MarketStatusAction::PreClose)]
148 #[case(BinanceTradingStatus::Settling, MarketStatusAction::Close)]
149 #[case(BinanceTradingStatus::Close, MarketStatusAction::Close)]
150 #[case(BinanceTradingStatus::TradingHalt, MarketStatusAction::Halt)]
151 #[case(BinanceTradingStatus::TradingCancelOnly, MarketStatusAction::Halt)]
152 #[case(
153 BinanceTradingStatus::Unknown,
154 MarketStatusAction::NotAvailableForTrading
155 )]
156 fn test_trading_status_to_market_action(
157 #[case] input: BinanceTradingStatus,
158 #[case] expected: MarketStatusAction,
159 ) {
160 assert_eq!(MarketStatusAction::from(input), expected);
161 }
162
163 #[rstest]
164 #[case(BinanceContractStatus::Trading, MarketStatusAction::Trading)]
165 #[case(BinanceContractStatus::PendingTrading, MarketStatusAction::PreOpen)]
166 #[case(BinanceContractStatus::PreDelivering, MarketStatusAction::PreClose)]
167 #[case(BinanceContractStatus::Delivering, MarketStatusAction::Close)]
168 #[case(BinanceContractStatus::Delivered, MarketStatusAction::Close)]
169 #[case(BinanceContractStatus::PreSettle, MarketStatusAction::PreClose)]
170 #[case(BinanceContractStatus::Settling, MarketStatusAction::Close)]
171 #[case(BinanceContractStatus::Close, MarketStatusAction::Close)]
172 #[case(BinanceContractStatus::TradingHalt, MarketStatusAction::Halt)]
173 #[case(BinanceContractStatus::TradingCancelOnly, MarketStatusAction::Halt)]
174 #[case(BinanceContractStatus::PreDelisting, MarketStatusAction::PreClose)]
175 #[case(BinanceContractStatus::Delisting, MarketStatusAction::Suspend)]
176 #[case(
177 BinanceContractStatus::Down,
178 MarketStatusAction::NotAvailableForTrading
179 )]
180 #[case(
181 BinanceContractStatus::Unknown,
182 MarketStatusAction::NotAvailableForTrading
183 )]
184 fn test_contract_status_to_market_action(
185 #[case] input: BinanceContractStatus,
186 #[case] expected: MarketStatusAction,
187 ) {
188 assert_eq!(MarketStatusAction::from(input), expected);
189 }
190
191 #[rstest]
192 fn test_diff_emits_on_change() {
193 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
194 let id = InstrumentId::from("BTCUSDT.BINANCE");
195
196 let mut cached = AHashMap::new();
197 cached.insert(id, MarketStatusAction::Trading);
198
199 let mut new_statuses = AHashMap::new();
200 new_statuses.insert(id, MarketStatusAction::Halt);
201
202 diff_and_emit_statuses(
203 &new_statuses,
204 &mut cached,
205 &tx,
206 UnixNanos::default(),
207 UnixNanos::default(),
208 );
209
210 let event = rx.try_recv().expect("expected status event");
211 match event {
212 DataEvent::InstrumentStatus(status) => {
213 assert_eq!(status.instrument_id, id);
214 assert_eq!(status.action, MarketStatusAction::Halt);
215 assert_eq!(status.is_trading, Some(false));
216 }
217 _ => panic!("expected InstrumentStatus event"),
218 }
219
220 assert_eq!(cached.get(&id), Some(&MarketStatusAction::Halt));
221 }
222
223 #[rstest]
224 fn test_diff_no_emit_when_unchanged() {
225 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
226 let id = InstrumentId::from("BTCUSDT.BINANCE");
227
228 let mut cached = AHashMap::new();
229 cached.insert(id, MarketStatusAction::Trading);
230
231 let mut new_statuses = AHashMap::new();
232 new_statuses.insert(id, MarketStatusAction::Trading);
233
234 diff_and_emit_statuses(
235 &new_statuses,
236 &mut cached,
237 &tx,
238 UnixNanos::default(),
239 UnixNanos::default(),
240 );
241
242 rx.try_recv().unwrap_err();
243 }
244
245 #[rstest]
246 fn test_diff_emits_for_new_symbol() {
247 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
248 let id = InstrumentId::from("ETHUSDT.BINANCE");
249
250 let mut cached = AHashMap::new();
251 let mut new_statuses = AHashMap::new();
252 new_statuses.insert(id, MarketStatusAction::Trading);
253
254 diff_and_emit_statuses(
255 &new_statuses,
256 &mut cached,
257 &tx,
258 UnixNanos::default(),
259 UnixNanos::default(),
260 );
261
262 let event = rx.try_recv().expect("expected status event for new symbol");
263 match event {
264 DataEvent::InstrumentStatus(status) => {
265 assert_eq!(status.instrument_id, id);
266 assert_eq!(status.action, MarketStatusAction::Trading);
267 assert_eq!(status.is_trading, Some(true));
268 }
269 _ => panic!("expected InstrumentStatus event"),
270 }
271 }
272
273 #[rstest]
274 fn test_diff_emits_not_available_for_removed_symbol() {
275 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
276 let id = InstrumentId::from("BTCUSDT.BINANCE");
277
278 let mut cached = AHashMap::new();
279 cached.insert(id, MarketStatusAction::Trading);
280
281 let new_statuses = AHashMap::new(); diff_and_emit_statuses(
284 &new_statuses,
285 &mut cached,
286 &tx,
287 UnixNanos::default(),
288 UnixNanos::default(),
289 );
290
291 let event = rx
292 .try_recv()
293 .expect("expected status event for removed symbol");
294 match event {
295 DataEvent::InstrumentStatus(status) => {
296 assert_eq!(status.instrument_id, id);
297 assert_eq!(status.action, MarketStatusAction::NotAvailableForTrading);
298 assert_eq!(status.is_trading, Some(false));
299 }
300 _ => panic!("expected InstrumentStatus event"),
301 }
302
303 assert!(!cached.contains_key(&id));
304 }
305}