1use std::sync::Arc;
17
18use anyhow::Context;
19use jiff::Timestamp;
20use nautilus_core::UnixNanos;
21use nautilus_model::{
22 data::{
23 Bar, BarType, BookOrder, Data, FundingRateUpdate, IndexPriceUpdate, MarkPriceUpdate,
24 OptionGreekValues, OptionGreeks, OrderBookDelta, OrderBookDeltas, OrderBookDepth,
25 QuoteTick, TradeTick,
26 },
27 enums::{AggregationSource, BookAction, GreeksConvention, OrderSide, RecordFlag},
28 identifiers::{InstrumentId, TradeId},
29 types::{Price, Quantity},
30};
31
32use super::{
33 message::{
34 BarMsg, BookChangeMsg, BookLevel, BookSnapshotMsg, DerivativeTickerMsg, OptionSummaryMsg,
35 TradeMsg, WsMessage,
36 },
37 types::TardisInstrumentMiniInfo,
38};
39use crate::{
40 common::parse::{
41 derive_trade_id, normalize_amount, parse_aggressor_side, parse_bar_spec, parse_book_action,
42 validate_non_zero_amount,
43 },
44 config::BookSnapshotOutput,
45};
46
47fn timestamp_to_unix_nanos(timestamp: Timestamp, field: &str) -> anyhow::Result<UnixNanos> {
48 let nanos = u64::try_from(timestamp.as_nanosecond())
49 .with_context(|| format!("invalid timestamp: {field} is outside the UnixNanos range"))?;
50 Ok(UnixNanos::from(nanos))
51}
52
53#[must_use]
54pub fn parse_tardis_ws_message(
55 msg: WsMessage,
56 info: &Arc<TardisInstrumentMiniInfo>,
57 book_snapshot_output: &BookSnapshotOutput,
58) -> Option<Data> {
59 match msg {
60 WsMessage::BookChange(msg) => {
61 if msg.bids.is_empty() && msg.asks.is_empty() {
62 let exchange = msg.exchange;
63 let symbol = &msg.symbol;
64 log::error!("Invalid book change for {exchange} {symbol} (empty bids and asks)");
65 return None;
66 }
67
68 match parse_book_change_msg_as_deltas(
69 &msg,
70 info.price_precision,
71 info.size_precision,
72 info.instrument_id,
73 ) {
74 Ok(deltas) => Some(Data::BookDeltas(Box::new(deltas))),
75 Err(e) => {
76 log::error!("Failed to parse book change message: {e}");
77 None
78 }
79 }
80 }
81 WsMessage::BookSnapshot(msg) => match msg.depth {
82 1 => {
83 match parse_book_snapshot_msg_as_quote(
84 &msg,
85 info.price_precision,
86 info.size_precision,
87 info.instrument_id,
88 ) {
89 Ok(quote) => Some(Data::Quote(quote)),
90 Err(e) => {
91 log::error!("Failed to parse book snapshot quote message: {e}");
92 None
93 }
94 }
95 }
96 _ => match book_snapshot_output {
97 BookSnapshotOutput::Depth => {
98 match parse_book_snapshot_msg_as_depth(
99 &msg,
100 info.price_precision,
101 info.size_precision,
102 info.instrument_id,
103 ) {
104 Ok(depth) => Some(Data::BookDepth(Box::new(depth))),
105 Err(e) => {
106 log::error!("Failed to parse book snapshot as depth: {e}");
107 None
108 }
109 }
110 }
111 BookSnapshotOutput::Deltas => {
112 match parse_book_snapshot_msg_as_deltas(
113 &msg,
114 info.price_precision,
115 info.size_precision,
116 info.instrument_id,
117 ) {
118 Ok(deltas) => Some(Data::BookDeltas(Box::new(deltas))),
119 Err(e) => {
120 log::error!("Failed to parse book snapshot as deltas: {e}");
121 None
122 }
123 }
124 }
125 },
126 },
127 WsMessage::Trade(msg) => {
128 match parse_trade_msg(
129 &msg,
130 info.price_precision,
131 info.size_precision,
132 info.instrument_id,
133 ) {
134 Ok(trade) => Some(Data::Trade(trade)),
135 Err(e) => {
136 log::error!("Failed to parse trade message: {e}");
137 None
138 }
139 }
140 }
141 WsMessage::TradeBar(msg) => {
142 match parse_bar_msg(
143 &msg,
144 info.price_precision,
145 info.size_precision,
146 info.instrument_id,
147 ) {
148 Ok(bar) => Some(Data::Bar(bar)),
149 Err(e) => {
150 log::error!("Failed to parse bar message: {e}");
151 None
152 }
153 }
154 }
155 WsMessage::OptionSummary(msg) => Some(Data::OptionGreeks(parse_option_summary_msg(
156 &msg,
157 info.instrument_id,
158 ))),
159 WsMessage::DerivativeTicker(_) => None,
162 WsMessage::Disconnect(_) => None,
163 }
164}
165
166#[must_use]
167pub fn parse_tardis_ws_message_data(
168 msg: WsMessage,
169 info: &Arc<TardisInstrumentMiniInfo>,
170 book_snapshot_output: &BookSnapshotOutput,
171 extract_bbo_as_quotes: bool,
172) -> Vec<Data> {
173 match msg {
174 WsMessage::OptionSummary(msg) if extract_bbo_as_quotes => {
175 let mut data = Vec::with_capacity(2);
176
177 match parse_option_summary_msg_as_quote(
178 &msg,
179 info.price_precision,
180 info.size_precision,
181 info.instrument_id,
182 ) {
183 Ok(Some(quote)) => data.push(Data::Quote(quote)),
184 Ok(None) => {}
185 Err(e) => {
186 log::error!("Failed to parse option summary quote message: {e}");
187 }
188 }
189
190 data.push(Data::OptionGreeks(parse_option_summary_msg(
191 &msg,
192 info.instrument_id,
193 )));
194 data
195 }
196 msg => parse_tardis_ws_message(msg, info, book_snapshot_output)
197 .into_iter()
198 .collect(),
199 }
200}
201
202#[must_use]
208pub fn parse_option_summary_msg(
209 msg: &OptionSummaryMsg,
210 instrument_id: InstrumentId,
211) -> OptionGreeks {
212 OptionGreeks {
213 instrument_id,
214 convention: GreeksConvention::BlackScholes,
215 greeks: OptionGreekValues {
216 delta: msg.delta.unwrap_or(0.0),
217 gamma: msg.gamma.unwrap_or(0.0),
218 vega: msg.vega.unwrap_or(0.0),
219 theta: msg.theta.unwrap_or(0.0),
220 rho: msg.rho.unwrap_or(0.0),
221 },
222 mark_iv: msg.mark_iv,
223 bid_iv: msg.best_bid_iv,
224 ask_iv: msg.best_ask_iv,
225 underlying_price: msg.underlying_price,
226 open_interest: msg.open_interest,
227 ts_event: UnixNanos::from(msg.timestamp),
228 ts_init: UnixNanos::from(msg.local_timestamp),
229 }
230}
231
232pub fn parse_option_summary_msg_as_quote(
240 msg: &OptionSummaryMsg,
241 price_precision: u8,
242 size_precision: u8,
243 instrument_id: InstrumentId,
244) -> anyhow::Result<Option<QuoteTick>> {
245 let (Some(best_bid_price), Some(best_bid_amount), Some(best_ask_price), Some(best_ask_amount)) = (
246 msg.best_bid_price,
247 msg.best_bid_amount,
248 msg.best_ask_price,
249 msg.best_ask_amount,
250 ) else {
251 return Ok(None);
252 };
253
254 let bid_price = Price::new_checked(best_bid_price, price_precision)
255 .with_context(|| format!("invalid option summary bid price for message: {msg:?}"))?;
256 let ask_price = Price::new_checked(best_ask_price, price_precision)
257 .with_context(|| format!("invalid option summary ask price for message: {msg:?}"))?;
258 validate_non_zero_amount(best_bid_amount, size_precision)
259 .with_context(|| format!("invalid option summary bid size for message: {msg:?}"))?;
260 let bid_size = Quantity::new_checked(best_bid_amount, size_precision)
261 .with_context(|| format!("invalid option summary bid size for message: {msg:?}"))?;
262 validate_non_zero_amount(best_ask_amount, size_precision)
263 .with_context(|| format!("invalid option summary ask size for message: {msg:?}"))?;
264 let ask_size = Quantity::new_checked(best_ask_amount, size_precision)
265 .with_context(|| format!("invalid option summary ask size for message: {msg:?}"))?;
266
267 Ok(Some(QuoteTick::new(
268 instrument_id,
269 bid_price,
270 ask_price,
271 bid_size,
272 ask_size,
273 UnixNanos::from(msg.timestamp),
274 UnixNanos::from(msg.local_timestamp),
275 )))
276}
277
278#[must_use]
281pub fn parse_tardis_ws_message_funding_rate(
282 msg: WsMessage,
283 info: &Arc<TardisInstrumentMiniInfo>,
284) -> Option<FundingRateUpdate> {
285 match msg {
286 WsMessage::DerivativeTicker(msg) => {
287 match parse_derivative_ticker_msg(&msg, info.instrument_id) {
288 Ok(funding_rate) => funding_rate,
289 Err(e) => {
290 log::error!("Failed to parse derivative ticker message for funding rate: {e}");
291 None
292 }
293 }
294 }
295 _ => None, }
297}
298
299pub fn parse_book_change_msg_as_deltas(
306 msg: &BookChangeMsg,
307 price_precision: u8,
308 size_precision: u8,
309 instrument_id: InstrumentId,
310) -> anyhow::Result<OrderBookDeltas> {
311 parse_book_msg_as_deltas(
312 &msg.bids,
313 &msg.asks,
314 msg.is_snapshot,
315 price_precision,
316 size_precision,
317 instrument_id,
318 msg.timestamp,
319 msg.local_timestamp,
320 )
321}
322
323pub fn parse_book_snapshot_msg_as_deltas(
330 msg: &BookSnapshotMsg,
331 price_precision: u8,
332 size_precision: u8,
333 instrument_id: InstrumentId,
334) -> anyhow::Result<OrderBookDeltas> {
335 parse_book_msg_as_deltas(
336 &msg.bids,
337 &msg.asks,
338 true,
339 price_precision,
340 size_precision,
341 instrument_id,
342 msg.timestamp,
343 msg.local_timestamp,
344 )
345}
346
347pub fn parse_book_snapshot_msg_as_depth(
353 msg: &BookSnapshotMsg,
354 price_precision: u8,
355 size_precision: u8,
356 instrument_id: InstrumentId,
357) -> anyhow::Result<OrderBookDepth> {
358 let ts_event = timestamp_to_unix_nanos(msg.timestamp, "event timestamp")?;
359 let ts_init = timestamp_to_unix_nanos(msg.local_timestamp, "init timestamp")?;
360
361 let mut bids = Vec::with_capacity(msg.bids.len());
362 let mut asks = Vec::with_capacity(msg.asks.len());
363 let mut bid_counts = Vec::with_capacity(msg.bids.len());
364 let mut ask_counts = Vec::with_capacity(msg.asks.len());
365
366 for level in &msg.bids {
367 bids.push(BookOrder::new(
368 OrderSide::Buy,
369 Price::new(level.price, price_precision),
370 Quantity::new(level.amount, size_precision),
371 0,
372 ));
373 bid_counts.push(1);
374 }
375
376 for level in &msg.asks {
377 asks.push(BookOrder::new(
378 OrderSide::Sell,
379 Price::new(level.price, price_precision),
380 Quantity::new(level.amount, size_precision),
381 0,
382 ));
383 ask_counts.push(1);
384 }
385
386 Ok(OrderBookDepth::new(
387 instrument_id,
388 bids,
389 asks,
390 bid_counts,
391 ask_counts,
392 RecordFlag::F_SNAPSHOT as u8,
393 0, ts_event,
395 ts_init,
396 ))
397}
398
399#[expect(clippy::too_many_arguments)]
401pub fn parse_book_msg_as_deltas(
407 bids: &[BookLevel],
408 asks: &[BookLevel],
409 is_snapshot: bool,
410 price_precision: u8,
411 size_precision: u8,
412 instrument_id: InstrumentId,
413 timestamp: Timestamp,
414 local_timestamp: Timestamp,
415) -> anyhow::Result<OrderBookDeltas> {
416 let ts_event = timestamp_to_unix_nanos(timestamp, "event timestamp")?;
417 let ts_init = timestamp_to_unix_nanos(local_timestamp, "init timestamp")?;
418
419 let capacity = if is_snapshot {
420 bids.len() + asks.len() + 1
421 } else {
422 bids.len() + asks.len()
423 };
424 let mut deltas: Vec<OrderBookDelta> = Vec::with_capacity(capacity);
425
426 if is_snapshot {
427 deltas.push(OrderBookDelta::clear(instrument_id, 0, ts_event, ts_init));
428 }
429
430 for level in bids {
431 match parse_book_level(
432 instrument_id,
433 price_precision,
434 size_precision,
435 OrderSide::Buy,
436 level,
437 is_snapshot,
438 ts_event,
439 ts_init,
440 ) {
441 Ok(delta) => deltas.push(delta),
442 Err(e) => log::warn!("Skipping invalid bid level for {instrument_id}: {e}"),
443 }
444 }
445
446 for level in asks {
447 match parse_book_level(
448 instrument_id,
449 price_precision,
450 size_precision,
451 OrderSide::Sell,
452 level,
453 is_snapshot,
454 ts_event,
455 ts_init,
456 ) {
457 Ok(delta) => deltas.push(delta),
458 Err(e) => log::warn!("Skipping invalid ask level for {instrument_id}: {e}"),
459 }
460 }
461
462 if let Some(last_delta) = deltas.last_mut() {
463 last_delta.flags |= RecordFlag::F_LAST as u8;
464 }
465
466 Ok(OrderBookDeltas::new(instrument_id, deltas))
467}
468
469#[expect(clippy::too_many_arguments)]
475pub fn parse_book_level(
476 instrument_id: InstrumentId,
477 price_precision: u8,
478 size_precision: u8,
479 side: OrderSide,
480 level: &BookLevel,
481 is_snapshot: bool,
482 ts_event: UnixNanos,
483 ts_init: UnixNanos,
484) -> anyhow::Result<OrderBookDelta> {
485 let amount = normalize_amount(level.amount, size_precision);
486 let action = parse_book_action(is_snapshot, amount);
487 let price = Price::new(level.price, price_precision);
488 let size = Quantity::new(amount, size_precision);
489 let order_id = 0; let order = BookOrder::new(side, price, size, order_id);
491 let flags = if is_snapshot {
492 RecordFlag::F_SNAPSHOT as u8
493 } else {
494 0
495 };
496 let sequence = 0; anyhow::ensure!(
499 !(action != BookAction::Delete && size.is_zero()),
500 "Invalid zero size for {action}"
501 );
502
503 Ok(OrderBookDelta::new(
504 instrument_id,
505 action,
506 order,
507 flags,
508 sequence,
509 ts_event,
510 ts_init,
511 ))
512}
513
514pub fn parse_book_snapshot_msg_as_quote(
521 msg: &BookSnapshotMsg,
522 price_precision: u8,
523 size_precision: u8,
524 instrument_id: InstrumentId,
525) -> anyhow::Result<QuoteTick> {
526 let ts_event = UnixNanos::from(msg.timestamp);
527 let ts_init = UnixNanos::from(msg.local_timestamp);
528
529 let best_bid = msg
530 .bids
531 .first()
532 .context("missing best bid level for quote message")?;
533 let bid_price = Price::new(best_bid.price, price_precision);
534 validate_non_zero_amount(best_bid.amount, size_precision)
535 .with_context(|| format!("Invalid bid size for message: {msg:?}"))?;
536 let bid_size = Quantity::new_checked(best_bid.amount, size_precision)
537 .with_context(|| format!("Invalid bid size for message: {msg:?}"))?;
538
539 let best_ask = msg
540 .asks
541 .first()
542 .context("missing best ask level for quote message")?;
543 let ask_price = Price::new(best_ask.price, price_precision);
544 validate_non_zero_amount(best_ask.amount, size_precision)
545 .with_context(|| format!("Invalid ask size for message: {msg:?}"))?;
546 let ask_size = Quantity::new_checked(best_ask.amount, size_precision)
547 .with_context(|| format!("Invalid ask size for message: {msg:?}"))?;
548
549 Ok(QuoteTick::new(
550 instrument_id,
551 bid_price,
552 ask_price,
553 bid_size,
554 ask_size,
555 ts_event,
556 ts_init,
557 ))
558}
559
560pub fn parse_trade_msg(
567 msg: &TradeMsg,
568 price_precision: u8,
569 size_precision: u8,
570 instrument_id: InstrumentId,
571) -> anyhow::Result<TradeTick> {
572 let price = Price::new(msg.price, price_precision);
573 validate_non_zero_amount(msg.amount, size_precision)
574 .with_context(|| format!("Invalid trade size in message: {msg:?}"))?;
575 let size = Quantity::new_checked(msg.amount, size_precision)
576 .with_context(|| format!("Invalid trade size in message: {msg:?}"))?;
577 let aggressor_side = parse_aggressor_side(&msg.side);
578 let ts_event = UnixNanos::from(msg.timestamp);
579 let ts_init = UnixNanos::from(msg.local_timestamp);
580 let trade_id = match msg.id.as_deref() {
581 Some(id) if !id.is_empty() => TradeId::new(id),
582 _ => derive_trade_id(
583 msg.symbol,
584 ts_event.as_u64(),
585 msg.price,
586 msg.amount,
587 &msg.side,
588 ),
589 };
590
591 Ok(TradeTick::new(
592 instrument_id,
593 price,
594 size,
595 aggressor_side,
596 trade_id,
597 ts_event,
598 ts_init,
599 ))
600}
601
602pub fn parse_bar_msg(
609 msg: &BarMsg,
610 price_precision: u8,
611 size_precision: u8,
612 instrument_id: InstrumentId,
613) -> anyhow::Result<Bar> {
614 let spec = parse_bar_spec(&msg.name)?;
615 let bar_type = BarType::new(instrument_id, spec, AggregationSource::External);
616
617 let open = Price::new(msg.open, price_precision);
618 let high = Price::new(msg.high, price_precision);
619 let low = Price::new(msg.low, price_precision);
620 let close = Price::new(msg.close, price_precision);
621 validate_non_zero_amount(msg.volume, size_precision)?;
622 let volume = Quantity::new_checked(msg.volume, size_precision)?;
623 let ts_event = UnixNanos::from(msg.timestamp);
624 let ts_init = UnixNanos::from(msg.local_timestamp);
625
626 Ok(Bar::new(
627 bar_type, open, high, low, close, volume, ts_event, ts_init,
628 ))
629}
630
631fn parse_derivative_ticker_timestamps(
633 msg: &DerivativeTickerMsg,
634) -> anyhow::Result<(UnixNanos, UnixNanos)> {
635 Ok((
636 timestamp_to_unix_nanos(msg.timestamp, "event timestamp")?,
637 timestamp_to_unix_nanos(msg.local_timestamp, "init timestamp")?,
638 ))
639}
640
641pub fn parse_derivative_ticker_msg(
647 msg: &DerivativeTickerMsg,
648 instrument_id: InstrumentId,
649) -> anyhow::Result<Option<FundingRateUpdate>> {
650 let funding_rate = match msg.funding_rate {
651 Some(rate) => rate,
652 None => return Ok(None),
653 };
654
655 let (ts_event, ts_init) = parse_derivative_ticker_timestamps(msg)?;
656 let rate = rust_decimal::Decimal::try_from(funding_rate)
657 .with_context(|| format!("failed to convert funding rate {funding_rate} to Decimal"))?;
658 let next_funding_ns = msg
659 .funding_timestamp
660 .map(|ts| timestamp_to_unix_nanos(ts, "funding timestamp"))
661 .transpose()?;
662
663 Ok(Some(FundingRateUpdate::new(
664 instrument_id,
665 rate,
666 None,
667 next_funding_ns,
668 ts_event,
669 ts_init,
670 )))
671}
672
673pub fn parse_derivative_ticker_mark_price(
679 msg: &DerivativeTickerMsg,
680 instrument_id: InstrumentId,
681 price_precision: u8,
682) -> anyhow::Result<Option<MarkPriceUpdate>> {
683 let mark_price = match msg.mark_price {
684 Some(p) => p,
685 None => return Ok(None),
686 };
687
688 let (ts_event, ts_init) = parse_derivative_ticker_timestamps(msg)?;
689
690 Ok(Some(MarkPriceUpdate::new(
691 instrument_id,
692 Price::new(mark_price, price_precision),
693 ts_event,
694 ts_init,
695 )))
696}
697
698pub fn parse_derivative_ticker_index_price(
704 msg: &DerivativeTickerMsg,
705 instrument_id: InstrumentId,
706 price_precision: u8,
707) -> anyhow::Result<Option<IndexPriceUpdate>> {
708 let index_price = match msg.index_price {
709 Some(p) => p,
710 None => return Ok(None),
711 };
712
713 let (ts_event, ts_init) = parse_derivative_ticker_timestamps(msg)?;
714
715 Ok(Some(IndexPriceUpdate::new(
716 instrument_id,
717 Price::new(index_price, price_precision),
718 ts_event,
719 ts_init,
720 )))
721}
722
723#[cfg(test)]
724mod tests {
725 use nautilus_model::enums::AggressorSide;
726 use rstest::rstest;
727 use rust_decimal_macros::dec;
728
729 use super::*;
730 use crate::common::{
731 enums::TardisExchange, parse::parse_instrument_id, testing::load_test_json,
732 };
733
734 #[rstest]
735 fn test_parse_book_change_message() {
736 let json_data = load_test_json("book_change.json");
737 let msg: BookChangeMsg = serde_json::from_str(&json_data).unwrap();
738
739 let price_precision = 0;
740 let size_precision = 0;
741 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
742 let deltas =
743 parse_book_change_msg_as_deltas(&msg, price_precision, size_precision, instrument_id)
744 .unwrap();
745
746 assert_eq!(deltas.deltas.len(), 1);
747 assert_eq!(deltas.instrument_id, instrument_id);
748 assert_eq!(deltas.flags, RecordFlag::F_LAST as u8);
749 assert_eq!(deltas.sequence, 0);
750 assert_eq!(deltas.ts_event, UnixNanos::from(1571830193469000000));
751 assert_eq!(deltas.ts_init, UnixNanos::from(1571830193469000000));
752 assert_eq!(
753 deltas.deltas[0].instrument_id,
754 InstrumentId::from("XBTUSD.BITMEX")
755 );
756 assert_eq!(deltas.deltas[0].action, BookAction::Update);
757 assert_eq!(deltas.deltas[0].order.price, Price::from("7985"));
758 assert_eq!(deltas.deltas[0].order.size, Quantity::from(283318));
759 assert_eq!(deltas.deltas[0].order.order_id, 0);
760 assert_eq!(deltas.deltas[0].flags, RecordFlag::F_LAST as u8);
761 assert_eq!(deltas.deltas[0].sequence, 0);
762 assert_eq!(
763 deltas.deltas[0].ts_event,
764 UnixNanos::from(1571830193469000000)
765 );
766 assert_eq!(
767 deltas.deltas[0].ts_init,
768 UnixNanos::from(1571830193469000000)
769 );
770 }
771
772 #[rstest]
773 fn test_parse_book_snapshot_message_as_deltas() {
774 let json_data = load_test_json("book_snapshot.json");
775 let msg: BookSnapshotMsg = serde_json::from_str(&json_data).unwrap();
776
777 let price_precision = 1;
778 let size_precision = 0;
779 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
780 let deltas =
781 parse_book_snapshot_msg_as_deltas(&msg, price_precision, size_precision, instrument_id)
782 .unwrap();
783
784 let clear_delta = deltas.deltas[0];
785 let bid_delta = deltas.deltas[1];
786 let ask_delta = deltas.deltas[3];
787
788 assert_eq!(deltas.deltas.len(), 5);
789 assert_eq!(deltas.instrument_id, instrument_id);
790 assert_eq!(
791 deltas.flags,
792 RecordFlag::F_LAST as u8 + RecordFlag::F_SNAPSHOT as u8
793 );
794 assert_eq!(deltas.sequence, 0);
795 assert_eq!(deltas.ts_event, UnixNanos::from(1572010786950000000));
796 assert_eq!(deltas.ts_init, UnixNanos::from(1572010786961000000));
797
798 assert_eq!(clear_delta.instrument_id, instrument_id);
800 assert_eq!(clear_delta.action, BookAction::Clear);
801 assert_eq!(clear_delta.flags, RecordFlag::F_SNAPSHOT as u8);
802 assert_eq!(clear_delta.sequence, 0);
803 assert_eq!(clear_delta.ts_event, UnixNanos::from(1572010786950000000));
804 assert_eq!(clear_delta.ts_init, UnixNanos::from(1572010786961000000));
805
806 assert_eq!(bid_delta.instrument_id, instrument_id);
808 assert_eq!(bid_delta.action, BookAction::Add);
809 assert_eq!(bid_delta.order.side, OrderSide::Buy.into());
810 assert_eq!(bid_delta.order.price, Price::from("7633.5"));
811 assert_eq!(bid_delta.order.size, Quantity::from(1906067));
812 assert_eq!(bid_delta.order.order_id, 0);
813 assert_eq!(bid_delta.flags, RecordFlag::F_SNAPSHOT as u8);
814 assert_eq!(bid_delta.sequence, 0);
815 assert_eq!(bid_delta.ts_event, UnixNanos::from(1572010786950000000));
816 assert_eq!(bid_delta.ts_init, UnixNanos::from(1572010786961000000));
817
818 assert_eq!(ask_delta.instrument_id, instrument_id);
820 assert_eq!(ask_delta.action, BookAction::Add);
821 assert_eq!(ask_delta.order.side, OrderSide::Sell.into());
822 assert_eq!(ask_delta.order.price, Price::from("7634.0"));
823 assert_eq!(ask_delta.order.size, Quantity::from(1467849));
824 assert_eq!(ask_delta.order.order_id, 0);
825 assert_eq!(ask_delta.flags, RecordFlag::F_SNAPSHOT as u8);
826 assert_eq!(ask_delta.sequence, 0);
827 assert_eq!(ask_delta.ts_event, UnixNanos::from(1572010786950000000));
828 assert_eq!(ask_delta.ts_init, UnixNanos::from(1572010786961000000));
829 }
830
831 #[rstest]
832 fn test_parse_book_snapshot_message_as_depth() {
833 let json_data = load_test_json("book_snapshot.json");
834 let msg: BookSnapshotMsg = serde_json::from_str(&json_data).unwrap();
835
836 let price_precision = 1;
837 let size_precision = 0;
838 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
839
840 let depth =
841 parse_book_snapshot_msg_as_depth(&msg, price_precision, size_precision, instrument_id)
842 .unwrap();
843
844 assert_eq!(depth.instrument_id, instrument_id);
845 assert_eq!(depth.flags, RecordFlag::F_SNAPSHOT as u8);
846 assert_eq!(depth.sequence, 0);
847 assert_eq!(depth.ts_event, UnixNanos::from(1572010786950000000));
848 assert_eq!(depth.ts_init, UnixNanos::from(1572010786961000000));
849
850 assert_eq!(depth.bids[0].side, OrderSide::Buy.into());
852 assert_eq!(depth.bids[0].price, Price::from("7633.5"));
853 assert_eq!(depth.bids[0].size, Quantity::from(1906067));
854 assert_eq!(depth.bids[0].order_id, 0);
855 assert_eq!(depth.bid_counts[0], 1);
856
857 assert_eq!(depth.bids[1].side, OrderSide::Buy.into());
859 assert_eq!(depth.bids[1].price, Price::from("7633.0"));
860 assert_eq!(depth.bids[1].size, Quantity::from(65319));
861 assert_eq!(depth.bid_counts[1], 1);
862
863 assert_eq!(depth.asks[0].side, OrderSide::Sell.into());
865 assert_eq!(depth.asks[0].price, Price::from("7634.0"));
866 assert_eq!(depth.asks[0].size, Quantity::from(1467849));
867 assert_eq!(depth.asks[0].order_id, 0);
868 assert_eq!(depth.ask_counts[0], 1);
869
870 assert_eq!(depth.asks[1].side, OrderSide::Sell.into());
872 assert_eq!(depth.asks[1].price, Price::from("7634.5"));
873 assert_eq!(depth.asks[1].size, Quantity::from(67939));
874 assert_eq!(depth.ask_counts[1], 1);
875
876 assert_eq!(depth.bids.len(), 2);
877 assert_eq!(depth.asks.len(), 2);
878 assert_eq!(depth.bid_counts.as_slice(), &[1; 2]);
879 assert_eq!(depth.ask_counts.as_slice(), &[1; 2]);
880 assert_eq!(depth.bids[1].order_id, 0);
881 assert_eq!(depth.asks[1].order_id, 0);
882 }
883
884 #[rstest]
885 fn test_parse_book_snapshot_message_as_depth_keeps_all_levels() {
886 let bids: Vec<BookLevel> = (0..25)
887 .map(|i| BookLevel {
888 price: 7633.5 - f64::from(i) * 0.5,
889 amount: f64::from(1000 + i),
890 })
891 .collect();
892 let asks: Vec<BookLevel> = (0..25)
893 .map(|i| BookLevel {
894 price: 7634.0 + f64::from(i) * 0.5,
895 amount: f64::from(2000 + i),
896 })
897 .collect();
898 let msg = BookSnapshotMsg {
899 symbol: ustr::Ustr::from("XBTUSD"),
900 exchange: TardisExchange::Bitmex,
901 name: "book_snapshot_25_100ms".to_string(),
902 depth: 25,
903 interval: 100,
904 bids,
905 asks,
906 timestamp: "2019-10-25T13:39:46.950Z".parse::<Timestamp>().unwrap(),
907 local_timestamp: "2019-10-25T13:39:46.961Z".parse::<Timestamp>().unwrap(),
908 };
909
910 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
911 let depth = parse_book_snapshot_msg_as_depth(&msg, 1, 0, instrument_id).unwrap();
912
913 assert_eq!(depth.instrument_id, instrument_id);
914 assert_eq!(depth.bids.len(), 25);
915 assert_eq!(depth.asks.len(), 25);
916 assert_eq!(depth.bid_counts.as_slice(), &[1; 25]);
917 assert_eq!(depth.ask_counts.as_slice(), &[1; 25]);
918 assert_eq!(depth.bids[0].price, Price::from("7633.5"));
919 assert_eq!(depth.bids[24].price, Price::from("7621.5"));
920 assert_eq!(depth.asks[0].price, Price::from("7634.0"));
921 assert_eq!(depth.asks[24].price, Price::from("7646.0"));
922 assert_eq!(depth.flags, RecordFlag::F_SNAPSHOT as u8);
923 assert_eq!(depth.sequence, 0);
924 }
925
926 #[rstest]
927 fn test_parse_book_snapshot_message_as_quote() {
928 let json_data = load_test_json("book_snapshot.json");
929 let msg: BookSnapshotMsg = serde_json::from_str(&json_data).unwrap();
930
931 let price_precision = 1;
932 let size_precision = 0;
933 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
934 let quote =
935 parse_book_snapshot_msg_as_quote(&msg, price_precision, size_precision, instrument_id)
936 .expect("Failed to parse book snapshot quote message");
937
938 assert_eq!(quote.instrument_id, instrument_id);
939 assert_eq!(quote.bid_price, Price::from("7633.5"));
940 assert_eq!(quote.bid_size, Quantity::from(1906067));
941 assert_eq!(quote.ask_price, Price::from("7634.0"));
942 assert_eq!(quote.ask_size, Quantity::from(1467849));
943 assert_eq!(quote.ts_event, UnixNanos::from(1572010786950000000));
944 assert_eq!(quote.ts_init, UnixNanos::from(1572010786961000000));
945 }
946
947 #[rstest]
948 fn test_parse_trade_message() {
949 let json_data = load_test_json("trade.json");
950 let msg: TradeMsg = serde_json::from_str(&json_data).unwrap();
951
952 let price_precision = 0;
953 let size_precision = 0;
954 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
955 let trade = parse_trade_msg(&msg, price_precision, size_precision, instrument_id)
956 .expect("Failed to parse trade message");
957
958 assert_eq!(trade.instrument_id, instrument_id);
959 assert_eq!(trade.price, Price::from("7996"));
960 assert_eq!(trade.size, Quantity::from(50));
961 assert_eq!(trade.aggressor_side, AggressorSide::Sell);
962 assert_eq!(trade.ts_event, UnixNanos::from(1571826769669000000));
963 assert_eq!(trade.ts_init, UnixNanos::from(1571826769740000000));
964 }
965
966 fn build_trade_msg_without_id() -> TradeMsg {
967 let json_data = load_test_json("trade.json");
968 let mut msg: TradeMsg = serde_json::from_str(&json_data).unwrap();
969 msg.id = None;
970 msg
971 }
972
973 #[rstest]
974 fn test_parse_trade_message_derives_trade_id_when_missing() {
975 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
976
977 let first = parse_trade_msg(&build_trade_msg_without_id(), 0, 0, instrument_id).unwrap();
978 let second = parse_trade_msg(&build_trade_msg_without_id(), 0, 0, instrument_id).unwrap();
979
980 assert_eq!(first.trade_id, second.trade_id, "derivation must be stable");
981 assert_eq!(first.trade_id.as_str().len(), 16);
982
983 let mut altered = build_trade_msg_without_id();
984 altered.price = 7997.0;
985 let altered_trade = parse_trade_msg(&altered, 0, 0, instrument_id).unwrap();
986 assert_ne!(first.trade_id, altered_trade.trade_id);
987 }
988
989 #[rstest]
990 fn test_parse_trade_message_derives_trade_id_when_empty() {
991 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
992
993 let mut msg = build_trade_msg_without_id();
994 msg.id = Some(String::new());
995
996 let trade = parse_trade_msg(&msg, 0, 0, instrument_id).unwrap();
997 let fallback = parse_trade_msg(&build_trade_msg_without_id(), 0, 0, instrument_id).unwrap();
998 assert_eq!(trade.trade_id, fallback.trade_id);
999 }
1000
1001 #[rstest]
1002 fn test_parse_bar_message() {
1003 let json_data = load_test_json("bar.json");
1004 let msg: BarMsg = serde_json::from_str(&json_data).unwrap();
1005
1006 let price_precision = 1;
1007 let size_precision = 0;
1008 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
1009 let bar = parse_bar_msg(&msg, price_precision, size_precision, instrument_id).unwrap();
1010
1011 assert_eq!(
1012 bar.bar_type,
1013 BarType::from("XBTUSD.BITMEX-10-SECOND-LAST-EXTERNAL")
1014 );
1015 assert_eq!(bar.open, Price::from("7623.5"));
1016 assert_eq!(bar.high, Price::from("7623.5"));
1017 assert_eq!(bar.low, Price::from("7623"));
1018 assert_eq!(bar.close, Price::from("7623.5"));
1019 assert_eq!(bar.volume, Quantity::from(37034));
1020 assert_eq!(bar.ts_event, UnixNanos::from(1572009100000000000));
1021 assert_eq!(bar.ts_init, UnixNanos::from(1572009100369000000));
1022 }
1023
1024 #[rstest]
1025 fn test_parse_tardis_ws_message_book_snapshot_routes_to_depth() {
1026 let json_data = load_test_json("book_snapshot.json");
1027 let msg: BookSnapshotMsg = serde_json::from_str(&json_data).unwrap();
1028 let ws_msg = WsMessage::BookSnapshot(msg);
1029
1030 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
1031 let info = Arc::new(TardisInstrumentMiniInfo::new(
1032 instrument_id,
1033 None,
1034 TardisExchange::Bitmex,
1035 1,
1036 0,
1037 ));
1038
1039 let result = parse_tardis_ws_message(ws_msg, &info, &BookSnapshotOutput::Depth);
1040
1041 assert!(result.is_some());
1042 assert!(matches!(result.unwrap(), Data::BookDepth(_)));
1043 }
1044
1045 #[rstest]
1046 fn test_parse_tardis_ws_message_sparse_book_snapshot_routes_to_depth() {
1047 let json_data = r#"{
1048 "type": "book_snapshot",
1049 "symbol": "ETC",
1050 "exchange": "hyperliquid",
1051 "name": "book_snapshot_20_10s",
1052 "depth": 20,
1053 "interval": 10000,
1054 "bids": [{"price": 20.002, "amount": 5.81}],
1055 "asks": [{"price": 20.003, "amount": 162.45}, {}],
1056 "timestamp": "2025-03-03T10:48:10.000Z",
1057 "localTimestamp": "2025-03-03T10:48:10.596818Z"
1058 }"#;
1059 let msg: BookSnapshotMsg = serde_json::from_str(json_data).unwrap();
1060 let ws_msg = WsMessage::BookSnapshot(msg);
1061
1062 let instrument_id = InstrumentId::from("ETC.HYPERLIQUID");
1063 let info = Arc::new(TardisInstrumentMiniInfo::new(
1064 instrument_id,
1065 None,
1066 TardisExchange::Hyperliquid,
1067 3,
1068 2,
1069 ));
1070
1071 let result = parse_tardis_ws_message(ws_msg, &info, &BookSnapshotOutput::Depth);
1072
1073 assert!(result.is_some());
1074 assert!(matches!(result.unwrap(), Data::BookDepth(_)));
1075 }
1076
1077 #[rstest]
1078 fn test_parse_tardis_ws_message_book_snapshot_routes_to_deltas() {
1079 let json_data = load_test_json("book_snapshot.json");
1080 let msg: BookSnapshotMsg = serde_json::from_str(&json_data).unwrap();
1081 let ws_msg = WsMessage::BookSnapshot(msg);
1082
1083 let instrument_id = InstrumentId::from("XBTUSD.BITMEX");
1084 let info = Arc::new(TardisInstrumentMiniInfo::new(
1085 instrument_id,
1086 None,
1087 TardisExchange::Bitmex,
1088 1,
1089 0,
1090 ));
1091
1092 let result = parse_tardis_ws_message(ws_msg, &info, &BookSnapshotOutput::Deltas);
1093
1094 assert!(result.is_some());
1095 assert!(matches!(result.unwrap(), Data::BookDeltas(_)));
1096 }
1097
1098 #[rstest]
1099 fn test_parse_tardis_ws_message_mexc_snapshot_5000_routes_to_deltas() {
1100 let json_data = load_test_json("mexc_book_snapshot_5000.json");
1101 let msg: BookSnapshotMsg = serde_json::from_str(&json_data).unwrap();
1102 let ws_msg = WsMessage::BookSnapshot(msg);
1103
1104 let instrument_id = InstrumentId::from("BTCUSDT.MEXC");
1105 let info = Arc::new(TardisInstrumentMiniInfo::new(
1106 instrument_id,
1107 None,
1108 TardisExchange::Mexc,
1109 1,
1110 1,
1111 ));
1112
1113 let result = parse_tardis_ws_message(ws_msg, &info, &BookSnapshotOutput::Deltas);
1114
1115 let Some(Data::BookDeltas(deltas)) = result else {
1116 panic!("Expected Data::BookDeltas, was {result:?}");
1117 };
1118 assert_eq!(deltas.instrument_id, instrument_id);
1119 assert_eq!(deltas.deltas.len(), 7);
1120 assert_eq!(deltas.deltas[0].action, BookAction::Clear);
1121 assert_eq!(deltas.deltas[1].order.price, Price::from("100.0"));
1122 assert_eq!(deltas.deltas[4].order.price, Price::from("101.0"));
1123 assert_eq!(
1124 deltas.deltas[6].flags,
1125 RecordFlag::F_LAST as u8 + RecordFlag::F_SNAPSHOT as u8
1126 );
1127 }
1128
1129 #[rstest]
1130 fn test_parse_tardis_ws_message_option_summary_routes_to_option_greeks() {
1131 let json_data = load_test_json("option_summary.json");
1132 let msg: OptionSummaryMsg = serde_json::from_str(&json_data).unwrap();
1133 let ts_event = UnixNanos::from(msg.timestamp);
1134 let ts_init = UnixNanos::from(msg.local_timestamp);
1135 let ws_msg = WsMessage::OptionSummary(msg);
1136
1137 let instrument_id = InstrumentId::from("BTC-28JUN24-70000-C.DERIBIT");
1138 let info = Arc::new(TardisInstrumentMiniInfo::new(
1139 instrument_id,
1140 None,
1141 TardisExchange::Deribit,
1142 4,
1143 1,
1144 ));
1145
1146 let result = parse_tardis_ws_message(ws_msg, &info, &BookSnapshotOutput::Deltas);
1147
1148 let Some(Data::OptionGreeks(greeks)) = result else {
1149 panic!("Expected Data::OptionGreeks, was {result:?}");
1150 };
1151 assert_eq!(greeks.instrument_id, instrument_id);
1152 assert_eq!(greeks.convention, GreeksConvention::BlackScholes);
1153 assert_eq!(greeks.greeks.delta, 0.25);
1154 assert_eq!(greeks.greeks.gamma, 0.00002);
1155 assert_eq!(greeks.greeks.vega, 45.5);
1156 assert_eq!(greeks.greeks.theta, -15.2);
1157 assert_eq!(greeks.greeks.rho, 0.05);
1158 assert_eq!(greeks.mark_iv, Some(0.565));
1159 assert_eq!(greeks.bid_iv, Some(0.55));
1160 assert_eq!(greeks.ask_iv, Some(0.58));
1161 assert_eq!(greeks.underlying_price, Some(63_500.0));
1162 assert_eq!(greeks.open_interest, Some(150.0));
1163 assert_eq!(greeks.ts_event, ts_event);
1164 assert_eq!(greeks.ts_init, ts_init);
1165 }
1166
1167 #[rstest]
1168 fn test_parse_tardis_ws_message_data_extracts_option_summary_quote_when_enabled() {
1169 let json_data = load_test_json("option_summary.json");
1170 let msg: OptionSummaryMsg = serde_json::from_str(&json_data).unwrap();
1171 let ts_event = UnixNanos::from(msg.timestamp);
1172 let ts_init = UnixNanos::from(msg.local_timestamp);
1173 let ws_msg = WsMessage::OptionSummary(msg);
1174
1175 let instrument_id = InstrumentId::from("BTC-28JUN24-70000-C.DERIBIT");
1176 let info = Arc::new(TardisInstrumentMiniInfo::new(
1177 instrument_id,
1178 None,
1179 TardisExchange::Deribit,
1180 4,
1181 1,
1182 ));
1183
1184 let data = parse_tardis_ws_message_data(ws_msg, &info, &BookSnapshotOutput::Deltas, true);
1185
1186 assert_eq!(data.len(), 2);
1187 let Data::Quote(quote) = &data[0] else {
1188 panic!("Expected first data item to be QuoteTick");
1189 };
1190 let Data::OptionGreeks(greeks) = &data[1] else {
1191 panic!("Expected second data item to be OptionGreeks");
1192 };
1193 assert_eq!(quote.instrument_id, instrument_id);
1194 assert_eq!(quote.bid_price, Price::from("0.035"));
1195 assert_eq!(quote.bid_size, Quantity::from("5"));
1196 assert_eq!(quote.ask_price, Price::from("0.04"));
1197 assert_eq!(quote.ask_size, Quantity::from("10"));
1198 assert_eq!(quote.ts_event, ts_event);
1199 assert_eq!(quote.ts_init, ts_init);
1200 assert_eq!(greeks.instrument_id, instrument_id);
1201 assert_eq!(greeks.ts_event, ts_event);
1202 assert_eq!(greeks.ts_init, ts_init);
1203 }
1204
1205 #[rstest]
1206 fn test_parse_tardis_ws_message_data_skips_option_summary_quote_when_disabled() {
1207 let json_data = load_test_json("option_summary.json");
1208 let msg: OptionSummaryMsg = serde_json::from_str(&json_data).unwrap();
1209 let ws_msg = WsMessage::OptionSummary(msg);
1210
1211 let instrument_id = InstrumentId::from("BTC-28JUN24-70000-C.DERIBIT");
1212 let info = Arc::new(TardisInstrumentMiniInfo::new(
1213 instrument_id,
1214 None,
1215 TardisExchange::Deribit,
1216 4,
1217 1,
1218 ));
1219
1220 let data = parse_tardis_ws_message_data(ws_msg, &info, &BookSnapshotOutput::Deltas, false);
1221
1222 assert_eq!(data.len(), 1);
1223 assert!(
1224 matches!(data[0], Data::OptionGreeks(_)),
1225 "Expected OptionGreeks, was {:?}",
1226 data[0]
1227 );
1228 }
1229
1230 #[rstest]
1231 fn test_parse_tardis_ws_message_data_skips_option_summary_quote_when_bbo_missing() {
1232 let json_data = load_test_json("option_summary.json");
1233 let mut msg: OptionSummaryMsg = serde_json::from_str(&json_data).unwrap();
1234 msg.best_ask_price = None;
1235 let ws_msg = WsMessage::OptionSummary(msg);
1236
1237 let instrument_id = InstrumentId::from("BTC-28JUN24-70000-C.DERIBIT");
1238 let info = Arc::new(TardisInstrumentMiniInfo::new(
1239 instrument_id,
1240 None,
1241 TardisExchange::Deribit,
1242 4,
1243 1,
1244 ));
1245
1246 let data = parse_tardis_ws_message_data(ws_msg, &info, &BookSnapshotOutput::Deltas, true);
1247
1248 assert_eq!(data.len(), 1);
1249 assert!(
1250 matches!(data[0], Data::OptionGreeks(_)),
1251 "Expected OptionGreeks, was {:?}",
1252 data[0]
1253 );
1254 }
1255
1256 #[rstest]
1257 fn test_parse_tardis_ws_message_data_keeps_option_summary_when_bbo_size_invalid() {
1258 let json_data = load_test_json("option_summary.json");
1259 let mut msg: OptionSummaryMsg = serde_json::from_str(&json_data).unwrap();
1260 let ts_event = UnixNanos::from(msg.timestamp);
1261 let ts_init = UnixNanos::from(msg.local_timestamp);
1262 msg.best_bid_amount = Some(0.0);
1263 let ws_msg = WsMessage::OptionSummary(msg);
1264
1265 let instrument_id = InstrumentId::from("BTC-28JUN24-70000-C.DERIBIT");
1266 let info = Arc::new(TardisInstrumentMiniInfo::new(
1267 instrument_id,
1268 None,
1269 TardisExchange::Deribit,
1270 4,
1271 1,
1272 ));
1273
1274 let data = parse_tardis_ws_message_data(ws_msg, &info, &BookSnapshotOutput::Deltas, true);
1275
1276 assert_eq!(data.len(), 1);
1277 let Data::OptionGreeks(greeks) = &data[0] else {
1278 panic!("Expected OptionGreeks, was {:?}", data[0]);
1279 };
1280 assert_eq!(greeks.instrument_id, instrument_id);
1281 assert_eq!(greeks.ts_event, ts_event);
1282 assert_eq!(greeks.ts_init, ts_init);
1283 }
1284
1285 #[rstest]
1286 fn test_parse_tardis_ws_message_data_keeps_option_summary_when_bbo_price_invalid() {
1287 let json_data = load_test_json("option_summary.json");
1288 let mut msg: OptionSummaryMsg = serde_json::from_str(&json_data).unwrap();
1289 let ts_event = UnixNanos::from(msg.timestamp);
1290 let ts_init = UnixNanos::from(msg.local_timestamp);
1291 msg.best_bid_price = Some(f64::MAX);
1292 let ws_msg = WsMessage::OptionSummary(msg);
1293
1294 let instrument_id = InstrumentId::from("BTC-28JUN24-70000-C.DERIBIT");
1295 let info = Arc::new(TardisInstrumentMiniInfo::new(
1296 instrument_id,
1297 None,
1298 TardisExchange::Deribit,
1299 4,
1300 1,
1301 ));
1302
1303 let data = parse_tardis_ws_message_data(ws_msg, &info, &BookSnapshotOutput::Deltas, true);
1304
1305 assert_eq!(data.len(), 1);
1306 let Data::OptionGreeks(greeks) = &data[0] else {
1307 panic!("Expected OptionGreeks, was {:?}", data[0]);
1308 };
1309 assert_eq!(greeks.instrument_id, instrument_id);
1310 assert_eq!(greeks.ts_event, ts_event);
1311 assert_eq!(greeks.ts_init, ts_init);
1312 }
1313
1314 #[rstest]
1315 fn test_parse_option_summary_msg_defaults_absent_fields() {
1316 let ts = "2024-01-15T10:30:00.123Z".parse::<Timestamp>().unwrap();
1317 let msg = OptionSummaryMsg {
1318 symbol: ustr::Ustr::from("BTC-28JUN24-70000-C"),
1319 exchange: TardisExchange::Deribit,
1320 option_type: "call".to_string(),
1321 strike_price: 70_000.0,
1322 expiration_date: ts,
1323 best_bid_price: None,
1324 best_bid_amount: None,
1325 best_bid_iv: None,
1326 best_ask_price: None,
1327 best_ask_amount: None,
1328 best_ask_iv: None,
1329 last_price: None,
1330 open_interest: None,
1331 mark_price: None,
1332 mark_iv: None,
1333 delta: None,
1334 gamma: None,
1335 vega: None,
1336 theta: None,
1337 rho: None,
1338 underlying_price: None,
1339 underlying_index: "BTC-USD".to_string(),
1340 timestamp: ts,
1341 local_timestamp: ts,
1342 };
1343
1344 let instrument_id = InstrumentId::from("BTC-28JUN24-70000-C.DERIBIT");
1345 let greeks = parse_option_summary_msg(&msg, instrument_id);
1346
1347 assert_eq!(greeks.greeks.delta, 0.0);
1349 assert_eq!(greeks.greeks.gamma, 0.0);
1350 assert_eq!(greeks.greeks.vega, 0.0);
1351 assert_eq!(greeks.greeks.theta, 0.0);
1352 assert_eq!(greeks.greeks.rho, 0.0);
1353 assert_eq!(greeks.mark_iv, None);
1354 assert_eq!(greeks.bid_iv, None);
1355 assert_eq!(greeks.ask_iv, None);
1356 assert_eq!(greeks.underlying_price, None);
1357 assert_eq!(greeks.open_interest, None);
1358 }
1359
1360 #[rstest]
1361 fn test_parse_derivative_ticker_funding_rate() {
1362 let json_data = load_test_json("derivative_ticker.json");
1363 let msg: DerivativeTickerMsg = serde_json::from_str(&json_data).unwrap();
1364
1365 let instrument_id = InstrumentId::from("BTC-PERPETUAL.DERIBIT");
1366
1367 let result = parse_derivative_ticker_msg(&msg, instrument_id).unwrap();
1368 assert!(result.is_some());
1369
1370 let funding = result.unwrap();
1371 assert_eq!(funding.instrument_id, instrument_id);
1372 assert_eq!(funding.rate.to_string(), "-0.00001568");
1373 assert!(funding.ts_event.as_u64() > 0);
1374 assert!(funding.ts_init.as_u64() > 0);
1375 }
1376
1377 #[rstest]
1378 fn test_parse_derivative_ticker_mark_price() {
1379 let json_data = load_test_json("derivative_ticker.json");
1380 let msg: DerivativeTickerMsg = serde_json::from_str(&json_data).unwrap();
1381
1382 let instrument_id = InstrumentId::from("BTC-PERPETUAL.DERIBIT");
1383 let price_precision = 2;
1384
1385 let result =
1386 parse_derivative_ticker_mark_price(&msg, instrument_id, price_precision).unwrap();
1387 assert!(result.is_some());
1388
1389 let mark = result.unwrap();
1390 assert_eq!(mark.instrument_id, instrument_id);
1391 assert_eq!(mark.value, Price::new(7987.56, price_precision));
1392 assert!(mark.ts_event.as_u64() > 0);
1393 assert!(mark.ts_init.as_u64() > 0);
1394 }
1395
1396 #[rstest]
1397 fn test_parse_derivative_ticker_index_price() {
1398 let json_data = load_test_json("derivative_ticker.json");
1399 let msg: DerivativeTickerMsg = serde_json::from_str(&json_data).unwrap();
1400
1401 let instrument_id = InstrumentId::from("BTC-PERPETUAL.DERIBIT");
1402 let price_precision = 2;
1403
1404 let result =
1405 parse_derivative_ticker_index_price(&msg, instrument_id, price_precision).unwrap();
1406 assert!(result.is_some());
1407
1408 let index = result.unwrap();
1409 assert_eq!(index.instrument_id, instrument_id);
1410 assert_eq!(index.value, Price::new(7989.28, price_precision));
1411 assert!(index.ts_event.as_u64() > 0);
1412 assert!(index.ts_init.as_u64() > 0);
1413 }
1414
1415 #[rstest]
1416 fn test_parse_derivative_ticker_missing_fields() {
1417 let json = r#"{
1419 "type": "derivative_ticker",
1420 "symbol": "BTCUSD",
1421 "exchange": "bitmex",
1422 "lastPrice": null,
1423 "openInterest": null,
1424 "fundingRate": 0.0001,
1425 "indexPrice": null,
1426 "markPrice": null,
1427 "timestamp": "2024-01-01T00:00:00.000Z",
1428 "localTimestamp": "2024-01-01T00:00:00.100Z"
1429 }"#;
1430 let msg: DerivativeTickerMsg = serde_json::from_str(json).unwrap();
1431
1432 let instrument_id = InstrumentId::from("BTCUSD.BITMEX");
1433
1434 let funding = parse_derivative_ticker_msg(&msg, instrument_id).unwrap();
1435 assert_eq!(funding.unwrap().next_funding_ns, None);
1436
1437 let mark = parse_derivative_ticker_mark_price(&msg, instrument_id, 1).unwrap();
1438 assert!(mark.is_none());
1439
1440 let index = parse_derivative_ticker_index_price(&msg, instrument_id, 1).unwrap();
1441 assert!(index.is_none());
1442 }
1443
1444 #[rstest]
1445 fn test_parse_mexc_futures_derivative_ticker() {
1446 let json_data = load_test_json("mexc_futures_derivative_ticker.json");
1447 let msg: DerivativeTickerMsg = serde_json::from_str(&json_data).unwrap();
1448
1449 let instrument_id = InstrumentId::from("BTC_USDT-PERP.MEXC");
1450
1451 let funding = parse_derivative_ticker_msg(&msg, instrument_id).unwrap();
1452 let mark = parse_derivative_ticker_mark_price(&msg, instrument_id, 1).unwrap();
1453 let index = parse_derivative_ticker_index_price(&msg, instrument_id, 1).unwrap();
1454
1455 assert_eq!(msg.exchange, TardisExchange::MexcFutures);
1456 assert_eq!(funding.unwrap().rate.to_string(), "0.0001");
1457 assert_eq!(mark.unwrap().value, Price::new(61235.2, 1));
1458 assert_eq!(index.unwrap().value, Price::new(61230.1, 1));
1459 }
1460
1461 #[rstest]
1462 fn test_parse_okex_xperp_derivative_ticker() {
1463 let json_data = load_test_json("okex_futures_xperp_derivative_ticker.json");
1464 let msg: DerivativeTickerMsg = serde_json::from_str(&json_data).unwrap();
1465
1466 let instrument_id = parse_instrument_id(&msg.exchange, msg.symbol);
1467 let funding = parse_derivative_ticker_msg(&msg, instrument_id)
1468 .unwrap()
1469 .unwrap();
1470 let mark = parse_derivative_ticker_mark_price(&msg, instrument_id, 1)
1471 .unwrap()
1472 .unwrap();
1473 let index = parse_derivative_ticker_index_price(&msg, instrument_id, 1)
1474 .unwrap()
1475 .unwrap();
1476
1477 assert_eq!(msg.exchange, TardisExchange::OkexFutures);
1480 assert_eq!(
1481 instrument_id,
1482 InstrumentId::from("BTC-USD_UM_XPERP-310404.OKEX")
1483 );
1484 assert_eq!(funding.rate, dec!(-0.0004050736802759));
1485 assert_eq!(
1486 funding.next_funding_ns,
1487 Some(UnixNanos::from("2026-08-10T16:00:00Z"))
1488 );
1489 assert_eq!(funding.interval, None);
1490 assert_eq!(
1491 funding.ts_event,
1492 UnixNanos::from("2026-08-10T12:00:13.061Z")
1493 );
1494 assert_eq!(funding.ts_init, UnixNanos::from("2026-08-10T12:00:13.093Z"));
1495 assert_eq!(mark.value, Price::new(65014.7, 1));
1496 assert_eq!(index.value, Price::new(65066.1, 1));
1497 }
1498
1499 #[rstest]
1500 fn test_parse_okex_usdc_derivative_ticker_across_index_migration() {
1501 let pre_json = load_test_json("okex_swap_derivative_ticker_pre_index_migration.json");
1502 let post_json = load_test_json("okex_swap_derivative_ticker_post_index_migration.json");
1503 let pre: DerivativeTickerMsg = serde_json::from_str(&pre_json).unwrap();
1504 let post: DerivativeTickerMsg = serde_json::from_str(&post_json).unwrap();
1505
1506 let pre_id = parse_instrument_id(&pre.exchange, pre.symbol);
1507 let post_id = parse_instrument_id(&post.exchange, post.symbol);
1508 let pre_index = parse_derivative_ticker_index_price(&pre, pre_id, 1)
1509 .unwrap()
1510 .unwrap();
1511 let post_index = parse_derivative_ticker_index_price(&post, post_id, 1)
1512 .unwrap()
1513 .unwrap();
1514 let pre_funding = parse_derivative_ticker_msg(&pre, pre_id).unwrap().unwrap();
1515 let post_funding = parse_derivative_ticker_msg(&post, post_id)
1516 .unwrap()
1517 .unwrap();
1518
1519 assert_eq!(pre.exchange, TardisExchange::OkexSwap);
1524 assert_eq!(post.exchange, TardisExchange::OkexSwap);
1525 assert_eq!(pre_id, InstrumentId::from("BTC-USDC-SWAP.OKEX"));
1526 assert_eq!(post_id, pre_id);
1527 assert_eq!(pre_index.value, Price::new(28379.9, 1));
1528 assert_eq!(post_index.value, Price::new(28525.1, 1));
1529 assert_eq!(pre_funding.rate, dec!(-0.0001546799749051));
1530 assert_eq!(post_funding.rate, dec!(0.0001309027042856));
1531 assert_eq!(
1532 pre_funding.next_funding_ns,
1533 Some(UnixNanos::from("2023-04-01T16:00:00Z"))
1534 );
1535 assert_eq!(
1536 post_funding.next_funding_ns,
1537 Some(UnixNanos::from("2023-05-01T16:00:00Z"))
1538 );
1539 }
1540}