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 price_filter: Some(BinancePriceFilterSbe {
440 price_exponent: -8,
441 min_price: 1_000_000,
442 max_price: 100_000_000_000_000,
443 tick_size: 1_000_000,
444 }),
445 lot_size_filter: Some(BinanceLotSizeFilterSbe {
446 qty_exponent: -8,
447 min_qty: 10_000,
448 max_qty: 900_000_000_000,
449 step_size: 10_000,
450 }),
451 },
452 permissions: vec![vec!["SPOT".to_string()]],
453 };
454
455 let ts = UnixNanos::from(1_700_000_000_000_000_000u64);
456 parse_spot_instrument_sbe(&symbol, ts, ts).unwrap()
457 }
458
459 #[rstest]
460 fn test_decode_empty_buffer() {
461 let err = decode_market_data(&[]).unwrap_err();
462 assert!(matches!(err, StreamDecodeError::BufferTooShort { .. }));
463 }
464
465 #[rstest]
466 fn test_decode_short_buffer() {
467 let buf = [0u8; 5];
468 let err = decode_market_data(&buf).unwrap_err();
469 assert!(matches!(err, StreamDecodeError::BufferTooShort { .. }));
470 }
471
472 #[rstest]
473 fn test_decode_wrong_schema() {
474 let mut buf = [0u8; 100];
475 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());
477 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();
481 assert!(matches!(err, StreamDecodeError::SchemaMismatch { .. }));
482 }
483
484 #[rstest]
485 fn test_decode_unknown_template() {
486 let mut buf = [0u8; 100];
487 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());
490 buf[6..8].copy_from_slice(&0u16.to_le_bytes()); let err = decode_market_data(&buf).unwrap_err();
493 assert!(matches!(err, StreamDecodeError::UnknownTemplateId(9999)));
494 }
495
496 #[rstest]
497 fn test_decode_valid_best_bid_ask() {
498 let buf = make_bbo_buffer();
499 let msg = decode_market_data(&buf).unwrap();
500
501 match msg {
502 MarketDataMessage::BestBidAsk(event) => {
503 assert_eq!(event.event_time_us, 1_000_000);
504 assert_eq!(event.symbol, Ustr::from("BTCUSDT"));
505 }
506 _ => panic!("Expected BestBidAsk"),
507 }
508 }
509
510 #[rstest]
511 fn test_parse_trades_event() {
512 let instrument = sample_instrument();
513 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
514 let event = TradesStreamEvent {
515 event_time_us: 1_700_000_000_000_000,
516 transact_time_us: 1_700_000_000_100_000,
517 price_exponent: -2,
518 qty_exponent: -4,
519 trades: vec![
520 Trade {
521 id: 1,
522 price_mantissa: 12_345,
523 qty_mantissa: 25_000,
524 is_buyer_maker: false,
525 },
526 Trade {
527 id: 2,
528 price_mantissa: 12_340,
529 qty_mantissa: 10_000,
530 is_buyer_maker: true,
531 },
532 ],
533 symbol: Ustr::from("ETHUSDT"),
534 };
535
536 let data = parse_trades_event(&event, &instrument, ts_init);
537
538 assert_eq!(data.len(), 2);
539 match &data[0] {
540 Data::Trade(trade) => {
541 assert_eq!(trade.instrument_id, instrument.id());
542 assert_eq!(trade.price, Price::new(123.45, 2));
543 assert_eq!(trade.size, Quantity::new(2.5, 4));
544 assert_eq!(trade.aggressor_side, AggressorSide::Buy);
545 assert_eq!(trade.trade_id, TradeId::new("1"));
546 assert_eq!(
547 trade.ts_event,
548 UnixNanos::from(1_700_000_000_100_000_000u64)
549 );
550 assert_eq!(trade.ts_init, ts_init);
551 }
552 other => panic!("Expected trade data, was {other:?}"),
553 }
554
555 match &data[1] {
556 Data::Trade(trade) => {
557 assert_eq!(trade.aggressor_side, AggressorSide::Sell);
558 assert_eq!(trade.ts_init, ts_init);
559 }
560 other => panic!("Expected trade data, was {other:?}"),
561 }
562 }
563
564 #[rstest]
565 fn test_parse_bbo_event() {
566 let instrument = sample_instrument();
567 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
568 let event = BestBidAskStreamEvent {
569 event_time_us: 1_700_000_000_000_000,
570 book_update_id: 123,
571 price_exponent: -2,
572 qty_exponent: -4,
573 bid_price_mantissa: 12_345,
574 bid_qty_mantissa: 25_000,
575 ask_price_mantissa: 12_350,
576 ask_qty_mantissa: 30_000,
577 symbol: Ustr::from("ETHUSDT"),
578 };
579
580 let quote = parse_bbo_event(&event, &instrument, ts_init);
581
582 assert_eq!(quote.instrument_id, instrument.id());
583 assert_eq!(quote.bid_price, Price::new(123.45, 2));
584 assert_eq!(quote.ask_price, Price::new(123.50, 2));
585 assert_eq!(quote.bid_size, Quantity::new(2.5, 4));
586 assert_eq!(quote.ask_size, Quantity::new(3.0, 4));
587 assert_eq!(
588 quote.ts_event,
589 UnixNanos::from(1_700_000_000_000_000_000u64)
590 );
591 assert_eq!(quote.ts_init, ts_init);
592 }
593
594 #[rstest]
595 #[case::negative(-1)]
596 #[case::overflow(i64::MAX)]
597 fn test_parse_bbo_event_falls_back_for_invalid_timestamp(#[case] event_time_us: i64) {
598 let instrument = sample_instrument();
599 let event = BestBidAskStreamEvent {
600 event_time_us,
601 book_update_id: 123,
602 price_exponent: -2,
603 qty_exponent: -4,
604 bid_price_mantissa: 12_345,
605 bid_qty_mantissa: 25_000,
606 ask_price_mantissa: 12_350,
607 ask_qty_mantissa: 30_000,
608 symbol: Ustr::from("ETHUSDT"),
609 };
610
611 let ts_init = UnixNanos::from(1);
612 let quote = parse_bbo_event(&event, &instrument, ts_init);
613
614 assert_eq!(quote.ts_event, ts_init);
615 assert_eq!(quote.ts_init, ts_init);
616 }
617
618 #[rstest]
619 fn test_parse_depth_snapshot() {
620 let instrument = sample_instrument();
621 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
622 let event = DepthSnapshotStreamEvent {
623 event_time_us: 1_700_000_000_000_000,
624 book_update_id: 123,
625 price_exponent: -2,
626 qty_exponent: -4,
627 bids: vec![PriceLevel {
628 price_mantissa: 12_345,
629 qty_mantissa: 25_000,
630 }],
631 asks: vec![PriceLevel {
632 price_mantissa: 12_350,
633 qty_mantissa: 30_000,
634 }],
635 symbol: Ustr::from("ETHUSDT"),
636 };
637
638 let deltas = parse_depth_snapshot(&event, &instrument, ts_init).unwrap();
639
640 assert_eq!(deltas.instrument_id, instrument.id());
641 assert_eq!(deltas.deltas.len(), 3);
642 assert_eq!(deltas.deltas[0].action, BookAction::Clear);
643 assert_eq!(deltas.deltas[1].action, BookAction::Add);
644 assert_eq!(deltas.deltas[1].order.side, OrderSide::Buy.into());
645 assert_eq!(deltas.deltas[1].order.price, Price::new(123.45, 2));
646 assert_eq!(deltas.deltas[1].order.size, Quantity::new(2.5, 4));
647 assert_eq!(deltas.deltas[2].action, BookAction::Add);
648 assert_eq!(deltas.deltas[2].order.side, OrderSide::Sell.into());
649 assert_eq!(deltas.deltas[2].order.price, Price::new(123.50, 2));
650 assert_eq!(deltas.deltas[2].order.size, Quantity::new(3.0, 4));
651 assert_eq!(deltas.deltas[2].flags, RecordFlag::F_LAST as u8);
652 assert_eq!(deltas.deltas[0].sequence, 123);
653 assert_eq!(deltas.deltas[2].sequence, 123);
654 assert_eq!(
655 deltas.ts_event,
656 UnixNanos::from(1_700_000_000_000_000_000u64)
657 );
658 assert_eq!(deltas.ts_init, ts_init);
659 assert!(deltas.deltas.iter().all(|delta| delta.ts_init == ts_init));
660 }
661
662 #[rstest]
663 fn test_parse_depth_snapshot_empty_returns_none() {
664 let instrument = sample_instrument();
665 let event = DepthSnapshotStreamEvent {
666 event_time_us: 1_700_000_000_000_000,
667 book_update_id: 123,
668 price_exponent: -2,
669 qty_exponent: -4,
670 bids: vec![],
671 asks: vec![],
672 symbol: Ustr::from("ETHUSDT"),
673 };
674
675 let deltas = parse_depth_snapshot(&event, &instrument, UnixNanos::from(2));
676
677 assert!(deltas.is_none());
678 }
679
680 #[rstest]
681 fn test_parse_depth_diff() {
682 let instrument = sample_instrument();
683 let ts_init = UnixNanos::from(1_800_000_000_000_000_000u64);
684 let event = DepthDiffStreamEvent {
685 event_time_us: 1_700_000_000_000_000,
686 first_book_update_id: 100,
687 last_book_update_id: 101,
688 price_exponent: -2,
689 qty_exponent: -4,
690 bids: vec![
691 PriceLevel {
692 price_mantissa: 12_345,
693 qty_mantissa: 25_000,
694 },
695 PriceLevel {
696 price_mantissa: 12_340,
697 qty_mantissa: 0,
698 },
699 ],
700 asks: vec![PriceLevel {
701 price_mantissa: 12_350,
702 qty_mantissa: 30_000,
703 }],
704 symbol: Ustr::from("ETHUSDT"),
705 };
706
707 let deltas = parse_depth_diff(&event, &instrument, ts_init).unwrap();
708
709 assert_eq!(deltas.instrument_id, instrument.id());
710 assert_eq!(deltas.deltas.len(), 3);
711 assert_eq!(deltas.deltas[0].action, BookAction::Update);
712 assert_eq!(deltas.deltas[0].order.side, OrderSide::Buy.into());
713 assert_eq!(deltas.deltas[1].action, BookAction::Delete);
714 assert_eq!(deltas.deltas[1].order.side, OrderSide::Buy.into());
715 assert_eq!(deltas.deltas[2].action, BookAction::Update);
716 assert_eq!(deltas.deltas[2].order.side, OrderSide::Sell.into());
717 assert_eq!(deltas.deltas[2].flags, RecordFlag::F_LAST as u8);
718 assert_eq!(deltas.deltas[0].sequence, 101);
719 assert_eq!(deltas.deltas[2].sequence, 101);
720 assert_eq!(
721 deltas.ts_event,
722 UnixNanos::from(1_700_000_000_000_000_000u64)
723 );
724 assert_eq!(deltas.ts_init, ts_init);
725 assert!(deltas.deltas.iter().all(|delta| delta.ts_init == ts_init));
726 }
727
728 #[rstest]
729 fn test_parse_depth_diff_empty_returns_none() {
730 let instrument = sample_instrument();
731 let event = DepthDiffStreamEvent {
732 event_time_us: 1_700_000_000_000_000,
733 first_book_update_id: 100,
734 last_book_update_id: 101,
735 price_exponent: -2,
736 qty_exponent: -4,
737 bids: vec![],
738 asks: vec![],
739 symbol: Ustr::from("ETHUSDT"),
740 };
741
742 let deltas = parse_depth_diff(&event, &instrument, UnixNanos::from(2));
743
744 assert!(deltas.is_none());
745 }
746}