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