1use nautilus_core::nanos::UnixNanos;
19use nautilus_model::{
20 data::{BookOrder, Data, OrderBookDelta, OrderBookDeltas, QuoteTick, TradeTick},
21 enums::{AggressorSide, BookAction, OrderSide, RecordFlag},
22 identifiers::TradeId,
23 instruments::{Instrument, InstrumentAny},
24 types::{Price, Quantity},
25};
26
27use crate::{
28 common::parse::parse_micros_or_init,
29 spot::sbe::stream::{
30 BestBidAskStreamEvent, DepthDiffStreamEvent, DepthSnapshotStreamEvent, MessageHeader,
31 StreamDecodeError, TradesStreamEvent, template_id,
32 },
33};
34
35#[derive(Debug)]
37pub enum MarketDataMessage {
38 Trades(TradesStreamEvent),
40 BestBidAsk(BestBidAskStreamEvent),
42 DepthSnapshot(DepthSnapshotStreamEvent),
44 DepthDiff(DepthDiffStreamEvent),
46}
47
48pub fn decode_market_data(buf: &[u8]) -> Result<MarketDataMessage, StreamDecodeError> {
58 let header = MessageHeader::decode(buf)?;
59 header.validate_schema()?;
60
61 match header.template_id {
62 template_id::TRADES_STREAM_EVENT => Ok(MarketDataMessage::Trades(
63 TradesStreamEvent::decode_validated(buf)?,
64 )),
65 template_id::BEST_BID_ASK_STREAM_EVENT => Ok(MarketDataMessage::BestBidAsk(
66 BestBidAskStreamEvent::decode_validated(buf)?,
67 )),
68 template_id::DEPTH_SNAPSHOT_STREAM_EVENT => Ok(MarketDataMessage::DepthSnapshot(
69 DepthSnapshotStreamEvent::decode_validated(buf)?,
70 )),
71 template_id::DEPTH_DIFF_STREAM_EVENT => Ok(MarketDataMessage::DepthDiff(
72 DepthDiffStreamEvent::decode_validated(buf)?,
73 )),
74 _ => Err(StreamDecodeError::UnknownTemplateId(header.template_id)),
75 }
76}
77
78pub fn parse_trades_event(
81 event: &TradesStreamEvent,
82 instrument: &InstrumentAny,
83 ts_init: UnixNanos,
84) -> Vec<Data> {
85 let instrument_id = instrument.id();
86 let price_precision = instrument.price_precision();
87 let size_precision = instrument.size_precision();
88
89 event
90 .trades
91 .iter()
92 .map(|t| {
93 let price = Price::from_mantissa_exponent(
94 t.price_mantissa,
95 event.price_exponent,
96 price_precision,
97 );
98 let size = Quantity::from_mantissa_exponent(
99 t.qty_mantissa as u64,
100 event.qty_exponent,
101 size_precision,
102 );
103 let ts_event = parse_micros_or_init(
104 event.transact_time_us,
105 "Spot SBE stream transaction time",
106 ts_init,
107 );
108
109 let trade = TradeTick::new(
110 instrument_id,
111 price,
112 size,
113 if t.is_buyer_maker {
114 AggressorSide::Sell
115 } else {
116 AggressorSide::Buy
117 },
118 TradeId::new(t.id.to_string()),
119 ts_event,
120 ts_init,
121 );
122 Data::from(trade)
123 })
124 .collect()
125}
126
127pub fn parse_bbo_event(
130 event: &BestBidAskStreamEvent,
131 instrument: &InstrumentAny,
132 ts_init: UnixNanos,
133) -> QuoteTick {
134 let instrument_id = instrument.id();
135 let price_precision = instrument.price_precision();
136 let size_precision = instrument.size_precision();
137
138 let bid_price = Price::from_mantissa_exponent(
139 event.bid_price_mantissa,
140 event.price_exponent,
141 price_precision,
142 );
143 let bid_size = Quantity::from_mantissa_exponent(
144 event.bid_qty_mantissa as u64,
145 event.qty_exponent,
146 size_precision,
147 );
148 let ask_price = Price::from_mantissa_exponent(
149 event.ask_price_mantissa,
150 event.price_exponent,
151 price_precision,
152 );
153 let ask_size = Quantity::from_mantissa_exponent(
154 event.ask_qty_mantissa as u64,
155 event.qty_exponent,
156 size_precision,
157 );
158 let ts_event = parse_micros_or_init(event.event_time_us, "Spot SBE BBO event time", ts_init);
159
160 QuoteTick::new(
161 instrument_id,
162 bid_price,
163 ask_price,
164 bid_size,
165 ask_size,
166 ts_event,
167 ts_init,
168 )
169}
170
171pub fn parse_depth_snapshot(
176 event: &DepthSnapshotStreamEvent,
177 instrument: &InstrumentAny,
178 ts_init: UnixNanos,
179) -> Option<OrderBookDeltas> {
180 let instrument_id = instrument.id();
181 let price_precision = instrument.price_precision();
182 let size_precision = instrument.size_precision();
183 let ts_event = parse_micros_or_init(
184 event.event_time_us,
185 "Spot SBE depth snapshot event time",
186 ts_init,
187 );
188 let sequence = event.book_update_id as u64;
189
190 let mut deltas = Vec::with_capacity(event.bids.len() + event.asks.len() + 1);
191
192 deltas.push(OrderBookDelta::clear(
194 instrument_id,
195 sequence,
196 ts_event,
197 ts_init,
198 ));
199
200 for (i, level) in event.bids.iter().enumerate() {
202 let price = Price::from_mantissa_exponent(
203 level.price_mantissa,
204 event.price_exponent,
205 price_precision,
206 );
207 let size = Quantity::from_mantissa_exponent(
208 level.qty_mantissa as u64,
209 event.qty_exponent,
210 size_precision,
211 );
212 let flags = if i == event.bids.len() - 1 && event.asks.is_empty() {
213 RecordFlag::F_LAST as u8
214 } else {
215 0
216 };
217
218 let order = BookOrder::new(OrderSide::Buy, price, size, 0);
219
220 deltas.push(OrderBookDelta::new(
221 instrument_id,
222 BookAction::Add,
223 order,
224 flags,
225 sequence,
226 ts_event,
227 ts_init,
228 ));
229 }
230
231 for (i, level) in event.asks.iter().enumerate() {
233 let price = Price::from_mantissa_exponent(
234 level.price_mantissa,
235 event.price_exponent,
236 price_precision,
237 );
238 let size = Quantity::from_mantissa_exponent(
239 level.qty_mantissa as u64,
240 event.qty_exponent,
241 size_precision,
242 );
243 let flags = if i == event.asks.len() - 1 {
244 RecordFlag::F_LAST as u8
245 } else {
246 0
247 };
248
249 let order = BookOrder::new(OrderSide::Sell, price, size, 0);
250
251 deltas.push(OrderBookDelta::new(
252 instrument_id,
253 BookAction::Add,
254 order,
255 flags,
256 sequence,
257 ts_event,
258 ts_init,
259 ));
260 }
261
262 if deltas.len() <= 1 {
265 return None;
266 }
267
268 Some(OrderBookDeltas::new(instrument_id, deltas))
269}
270
271pub fn parse_depth_diff(
276 event: &DepthDiffStreamEvent,
277 instrument: &InstrumentAny,
278 ts_init: UnixNanos,
279) -> Option<OrderBookDeltas> {
280 let instrument_id = instrument.id();
281 let price_precision = instrument.price_precision();
282 let size_precision = instrument.size_precision();
283 let ts_event = parse_micros_or_init(
284 event.event_time_us,
285 "Spot SBE depth diff event time",
286 ts_init,
287 );
288 let sequence = event.last_book_update_id as u64;
289
290 let mut deltas = Vec::with_capacity(event.bids.len() + event.asks.len());
291
292 for (i, level) in event.bids.iter().enumerate() {
294 let price = Price::from_mantissa_exponent(
295 level.price_mantissa,
296 event.price_exponent,
297 price_precision,
298 );
299 let size = Quantity::from_mantissa_exponent(
300 level.qty_mantissa as u64,
301 event.qty_exponent,
302 size_precision,
303 );
304
305 let action = if level.qty_mantissa == 0 {
307 BookAction::Delete
308 } else {
309 BookAction::Update
310 };
311
312 let flags = if i == event.bids.len() - 1 && event.asks.is_empty() {
313 RecordFlag::F_LAST as u8
314 } else {
315 0
316 };
317
318 let order = BookOrder::new(OrderSide::Buy, price, size, 0);
319
320 deltas.push(OrderBookDelta::new(
321 instrument_id,
322 action,
323 order,
324 flags,
325 sequence,
326 ts_event,
327 ts_init,
328 ));
329 }
330
331 for (i, level) in event.asks.iter().enumerate() {
333 let price = Price::from_mantissa_exponent(
334 level.price_mantissa,
335 event.price_exponent,
336 price_precision,
337 );
338 let size = Quantity::from_mantissa_exponent(
339 level.qty_mantissa as u64,
340 event.qty_exponent,
341 size_precision,
342 );
343
344 let action = if level.qty_mantissa == 0 {
345 BookAction::Delete
346 } else {
347 BookAction::Update
348 };
349
350 let flags = if i == event.asks.len() - 1 {
351 RecordFlag::F_LAST as u8
352 } else {
353 0
354 };
355
356 let order = BookOrder::new(OrderSide::Sell, price, size, 0);
357
358 deltas.push(OrderBookDelta::new(
359 instrument_id,
360 action,
361 order,
362 flags,
363 sequence,
364 ts_event,
365 ts_init,
366 ));
367 }
368
369 if deltas.is_empty() {
370 return None;
371 }
372
373 Some(OrderBookDeltas::new(instrument_id, deltas))
374}
375
376#[cfg(test)]
377mod tests {
378 use rstest::rstest;
379 use ustr::Ustr;
380
381 use super::*;
382 use crate::{
383 common::parse::parse_spot_instrument_sbe,
384 spot::{
385 http::models::{
386 BinanceLotSizeFilterSbe, BinancePriceFilterSbe, BinanceSymbolFiltersSbe,
387 BinanceSymbolSbe,
388 },
389 sbe::stream::{PriceLevel, STREAM_SCHEMA_ID, Trade},
390 },
391 };
392
393 fn make_bbo_buffer() -> Vec<u8> {
394 let mut buf = vec![0u8; 70];
395
396 buf[0..2].copy_from_slice(&50u16.to_le_bytes()); buf[2..4].copy_from_slice(&template_id::BEST_BID_ASK_STREAM_EVENT.to_le_bytes());
399 buf[4..6].copy_from_slice(&STREAM_SCHEMA_ID.to_le_bytes());
400 buf[6..8].copy_from_slice(&0u16.to_le_bytes()); let body = &mut buf[8..];
404 body[0..8].copy_from_slice(&1000000i64.to_le_bytes()); body[8..16].copy_from_slice(&12345i64.to_le_bytes()); body[16] = (-2i8) as u8; body[17] = (-8i8) as u8; body[18..26].copy_from_slice(&4200000i64.to_le_bytes()); body[26..34].copy_from_slice(&100000000i64.to_le_bytes()); body[34..42].copy_from_slice(&4200100i64.to_le_bytes()); body[42..50].copy_from_slice(&200000000i64.to_le_bytes()); body[50] = 7;
415 body[51..58].copy_from_slice(b"BTCUSDT");
416
417 buf
418 }
419
420 fn sample_instrument() -> InstrumentAny {
421 let symbol = BinanceSymbolSbe {
422 symbol: "ETHUSDT".to_string(),
423 base_asset: "ETH".to_string(),
424 quote_asset: "USDT".to_string(),
425 base_asset_precision: 8,
426 quote_asset_precision: 8,
427 status: 0,
428 order_types: 0,
429 iceberg_allowed: true,
430 oco_allowed: true,
431 oto_allowed: false,
432 quote_order_qty_market_allowed: true,
433 allow_trailing_stop: true,
434 cancel_replace_allowed: true,
435 amend_allowed: true,
436 is_spot_trading_allowed: true,
437 is_margin_trading_allowed: false,
438 filters: BinanceSymbolFiltersSbe {
439 notional_filters: Vec::new(),
440 price_filter: Some(BinancePriceFilterSbe {
441 price_exponent: -8,
442 min_price: 1_000_000,
443 max_price: 100_000_000_000_000,
444 tick_size: 1_000_000,
445 }),
446 lot_size_filter: Some(BinanceLotSizeFilterSbe {
447 qty_exponent: -8,
448 min_qty: 10_000,
449 max_qty: 900_000_000_000,
450 step_size: 10_000,
451 }),
452 },
453 permissions: vec![vec!["SPOT".to_string()]],
454 };
455
456 let ts = UnixNanos::from(1_700_000_000_000_000_000u64);
457 parse_spot_instrument_sbe(&symbol, ts, ts).unwrap()
458 }
459
460 #[rstest]
461 fn test_decode_empty_buffer() {
462 let err = decode_market_data(&[]).unwrap_err();
463 assert!(matches!(err, StreamDecodeError::BufferTooShort { .. }));
464 }
465
466 #[rstest]
467 fn test_decode_short_buffer() {
468 let buf = [0u8; 5];
469 let err = decode_market_data(&buf).unwrap_err();
470 assert!(matches!(err, StreamDecodeError::BufferTooShort { .. }));
471 }
472
473 #[rstest]
474 fn test_decode_wrong_schema() {
475 let mut buf = [0u8; 100];
476 buf[0..2].copy_from_slice(&50u16.to_le_bytes()); buf[2..4].copy_from_slice(&template_id::BEST_BID_ASK_STREAM_EVENT.to_le_bytes());
478 buf[4..6].copy_from_slice(&99u16.to_le_bytes()); buf[6..8].copy_from_slice(&0u16.to_le_bytes()); let err = decode_market_data(&buf).unwrap_err();
482 assert!(matches!(err, StreamDecodeError::SchemaMismatch { .. }));
483 }
484
485 #[rstest]
486 fn test_decode_unknown_template() {
487 let mut buf = [0u8; 100];
488 buf[0..2].copy_from_slice(&50u16.to_le_bytes()); buf[2..4].copy_from_slice(&9999u16.to_le_bytes()); buf[4..6].copy_from_slice(&STREAM_SCHEMA_ID.to_le_bytes());
491 buf[6..8].copy_from_slice(&0u16.to_le_bytes()); let err = decode_market_data(&buf).unwrap_err();
494 assert!(matches!(err, StreamDecodeError::UnknownTemplateId(9999)));
495 }
496
497 #[rstest]
498 fn test_decode_valid_best_bid_ask() {
499 let buf = make_bbo_buffer();
500 let msg = decode_market_data(&buf).unwrap();
501
502 match msg {
503 MarketDataMessage::BestBidAsk(event) => {
504 assert_eq!(event.event_time_us, 1_000_000);
505 assert_eq!(event.symbol, Ustr::from("BTCUSDT"));
506 }
507 _ => panic!("Expected BestBidAsk"),
508 }
509 }
510
511 #[rstest]
512 fn test_parse_trades_event() {
513 let instrument = sample_instrument();
514 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
515 let event = TradesStreamEvent {
516 event_time_us: 1_700_000_000_000_000,
517 transact_time_us: 1_700_000_000_100_000,
518 price_exponent: -2,
519 qty_exponent: -4,
520 trades: vec![
521 Trade {
522 id: 1,
523 price_mantissa: 12_345,
524 qty_mantissa: 25_000,
525 is_buyer_maker: false,
526 },
527 Trade {
528 id: 2,
529 price_mantissa: 12_340,
530 qty_mantissa: 10_000,
531 is_buyer_maker: true,
532 },
533 ],
534 symbol: Ustr::from("ETHUSDT"),
535 };
536
537 let data = parse_trades_event(&event, &instrument, ts_init);
538
539 assert_eq!(data.len(), 2);
540 match &data[0] {
541 Data::Trade(trade) => {
542 assert_eq!(trade.instrument_id, instrument.id());
543 assert_eq!(trade.price, Price::new(123.45, 2));
544 assert_eq!(trade.size, Quantity::new(2.5, 4));
545 assert_eq!(trade.aggressor_side, AggressorSide::Buy);
546 assert_eq!(trade.trade_id, TradeId::new("1"));
547 assert_eq!(
548 trade.ts_event,
549 UnixNanos::from(1_700_000_000_100_000_000u64)
550 );
551 assert_eq!(trade.ts_init, ts_init);
552 }
553 other => panic!("Expected trade data, was {other:?}"),
554 }
555
556 match &data[1] {
557 Data::Trade(trade) => {
558 assert_eq!(trade.aggressor_side, AggressorSide::Sell);
559 assert_eq!(trade.ts_init, ts_init);
560 }
561 other => panic!("Expected trade data, was {other:?}"),
562 }
563 }
564
565 #[rstest]
566 fn test_parse_bbo_event() {
567 let instrument = sample_instrument();
568 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
569 let event = BestBidAskStreamEvent {
570 event_time_us: 1_700_000_000_000_000,
571 book_update_id: 123,
572 price_exponent: -2,
573 qty_exponent: -4,
574 bid_price_mantissa: 12_345,
575 bid_qty_mantissa: 25_000,
576 ask_price_mantissa: 12_350,
577 ask_qty_mantissa: 30_000,
578 symbol: Ustr::from("ETHUSDT"),
579 };
580
581 let quote = parse_bbo_event(&event, &instrument, ts_init);
582
583 assert_eq!(quote.instrument_id, instrument.id());
584 assert_eq!(quote.bid_price, Price::new(123.45, 2));
585 assert_eq!(quote.ask_price, Price::new(123.50, 2));
586 assert_eq!(quote.bid_size, Quantity::new(2.5, 4));
587 assert_eq!(quote.ask_size, Quantity::new(3.0, 4));
588 assert_eq!(
589 quote.ts_event,
590 UnixNanos::from(1_700_000_000_000_000_000u64)
591 );
592 assert_eq!(quote.ts_init, ts_init);
593 }
594
595 #[rstest]
596 #[case::negative(-1)]
597 #[case::overflow(i64::MAX)]
598 fn test_parse_bbo_event_falls_back_for_invalid_timestamp(#[case] event_time_us: i64) {
599 let instrument = sample_instrument();
600 let event = BestBidAskStreamEvent {
601 event_time_us,
602 book_update_id: 123,
603 price_exponent: -2,
604 qty_exponent: -4,
605 bid_price_mantissa: 12_345,
606 bid_qty_mantissa: 25_000,
607 ask_price_mantissa: 12_350,
608 ask_qty_mantissa: 30_000,
609 symbol: Ustr::from("ETHUSDT"),
610 };
611
612 let ts_init = UnixNanos::from(1);
613 let quote = parse_bbo_event(&event, &instrument, ts_init);
614
615 assert_eq!(quote.ts_event, ts_init);
616 assert_eq!(quote.ts_init, ts_init);
617 }
618
619 #[rstest]
620 fn test_parse_depth_snapshot() {
621 let instrument = sample_instrument();
622 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
623 let event = DepthSnapshotStreamEvent {
624 event_time_us: 1_700_000_000_000_000,
625 book_update_id: 123,
626 price_exponent: -2,
627 qty_exponent: -4,
628 bids: vec![PriceLevel {
629 price_mantissa: 12_345,
630 qty_mantissa: 25_000,
631 }],
632 asks: vec![PriceLevel {
633 price_mantissa: 12_350,
634 qty_mantissa: 30_000,
635 }],
636 symbol: Ustr::from("ETHUSDT"),
637 };
638
639 let deltas = parse_depth_snapshot(&event, &instrument, ts_init).unwrap();
640
641 assert_eq!(deltas.instrument_id, instrument.id());
642 assert_eq!(deltas.deltas.len(), 3);
643 assert_eq!(deltas.deltas[0].action, BookAction::Clear);
644 assert_eq!(deltas.deltas[1].action, BookAction::Add);
645 assert_eq!(deltas.deltas[1].order.side, OrderSide::Buy.into());
646 assert_eq!(deltas.deltas[1].order.price, Price::new(123.45, 2));
647 assert_eq!(deltas.deltas[1].order.size, Quantity::new(2.5, 4));
648 assert_eq!(deltas.deltas[2].action, BookAction::Add);
649 assert_eq!(deltas.deltas[2].order.side, OrderSide::Sell.into());
650 assert_eq!(deltas.deltas[2].order.price, Price::new(123.50, 2));
651 assert_eq!(deltas.deltas[2].order.size, Quantity::new(3.0, 4));
652 assert_eq!(deltas.deltas[2].flags, RecordFlag::F_LAST as u8);
653 assert_eq!(deltas.deltas[0].sequence, 123);
654 assert_eq!(deltas.deltas[2].sequence, 123);
655 assert_eq!(
656 deltas.ts_event,
657 UnixNanos::from(1_700_000_000_000_000_000u64)
658 );
659 assert_eq!(deltas.ts_init, ts_init);
660 assert!(deltas.deltas.iter().all(|delta| delta.ts_init == ts_init));
661 }
662
663 #[rstest]
664 fn test_parse_depth_snapshot_empty_returns_none() {
665 let instrument = sample_instrument();
666 let event = DepthSnapshotStreamEvent {
667 event_time_us: 1_700_000_000_000_000,
668 book_update_id: 123,
669 price_exponent: -2,
670 qty_exponent: -4,
671 bids: vec![],
672 asks: vec![],
673 symbol: Ustr::from("ETHUSDT"),
674 };
675
676 let deltas = parse_depth_snapshot(&event, &instrument, UnixNanos::from(2));
677
678 assert!(deltas.is_none());
679 }
680
681 #[rstest]
682 fn test_parse_depth_diff() {
683 let instrument = sample_instrument();
684 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
685 let event = DepthDiffStreamEvent {
686 event_time_us: 1_700_000_000_000_000,
687 first_book_update_id: 100,
688 last_book_update_id: 101,
689 price_exponent: -2,
690 qty_exponent: -4,
691 bids: vec![
692 PriceLevel {
693 price_mantissa: 12_345,
694 qty_mantissa: 25_000,
695 },
696 PriceLevel {
697 price_mantissa: 12_340,
698 qty_mantissa: 0,
699 },
700 ],
701 asks: vec![PriceLevel {
702 price_mantissa: 12_350,
703 qty_mantissa: 30_000,
704 }],
705 symbol: Ustr::from("ETHUSDT"),
706 };
707
708 let deltas = parse_depth_diff(&event, &instrument, ts_init).unwrap();
709
710 assert_eq!(deltas.instrument_id, instrument.id());
711 assert_eq!(deltas.deltas.len(), 3);
712 assert_eq!(deltas.deltas[0].action, BookAction::Update);
713 assert_eq!(deltas.deltas[0].order.side, OrderSide::Buy.into());
714 assert_eq!(deltas.deltas[1].action, BookAction::Delete);
715 assert_eq!(deltas.deltas[1].order.side, OrderSide::Buy.into());
716 assert_eq!(deltas.deltas[2].action, BookAction::Update);
717 assert_eq!(deltas.deltas[2].order.side, OrderSide::Sell.into());
718 assert_eq!(deltas.deltas[2].flags, RecordFlag::F_LAST as u8);
719 assert_eq!(deltas.deltas[0].sequence, 101);
720 assert_eq!(deltas.deltas[2].sequence, 101);
721 assert_eq!(
722 deltas.ts_event,
723 UnixNanos::from(1_700_000_000_000_000_000u64)
724 );
725 assert_eq!(deltas.ts_init, ts_init);
726 assert!(deltas.deltas.iter().all(|delta| delta.ts_init == ts_init));
727 }
728
729 #[rstest]
730 fn test_parse_depth_diff_empty_returns_none() {
731 let instrument = sample_instrument();
732 let event = DepthDiffStreamEvent {
733 event_time_us: 1_700_000_000_000_000,
734 first_book_update_id: 100,
735 last_book_update_id: 101,
736 price_exponent: -2,
737 qty_exponent: -4,
738 bids: vec![],
739 asks: vec![],
740 symbol: Ustr::from("ETHUSDT"),
741 };
742
743 let deltas = parse_depth_diff(&event, &instrument, UnixNanos::from(2));
744
745 assert!(deltas.is_none());
746 }
747}