Skip to main content

nautilus_tardis/machine/
parse.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16use 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        // Derivative ticker messages are handled through a separate callback path
159        // for FundingRateUpdate since they're not part of the Data enum.
160        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/// Parses a Tardis option summary message into a Nautilus `OptionGreeks`.
202///
203/// Greeks absent from the exchange feed default to `0.0` (matching the existing exchange-greeks
204/// producers); implied volatilities, underlying price and open interest stay `None` when the
205/// exchange does not provide them.
206#[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
231/// Parses a Tardis option summary best bid/offer into a Nautilus `QuoteTick`.
232///
233/// Returns `Ok(None)` when any best bid/offer field is absent.
234///
235/// # Errors
236///
237/// Returns an error if a provided bid or ask price or size is invalid.
238pub 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/// Parse a Tardis WebSocket message specifically for funding rate updates.
274/// Returns `Some(FundingRateUpdate)` if the message contains funding rate data, `None` otherwise.
275#[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, // Only derivative ticker messages can contain funding rates
291    }
292}
293
294/// Parse a book change message into order book deltas, returning an error if timestamps invalid.
295/// Parse a book change message into order book deltas.
296///
297/// # Errors
298///
299/// Returns an error if timestamp fields cannot be converted to nanoseconds.
300pub 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
318/// Parse a book snapshot message into order book deltas, returning an error if timestamps invalid.
319/// Parse a book snapshot message into order book deltas.
320///
321/// # Errors
322///
323/// Returns an error if timestamp fields cannot be converted to nanoseconds.
324pub 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
342/// Parse a book snapshot message into an [`OrderBookDepth10`].
343///
344/// # Errors
345///
346/// Returns an error if timestamp fields cannot be converted to nanoseconds.
347pub 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, // Sequence not available from Tardis
389        ts_event,
390        ts_init,
391    ))
392}
393
394/// Parse raw book levels into order book deltas, returning error for invalid timestamps.
395#[expect(clippy::too_many_arguments)]
396/// Parse raw book levels into order book deltas.
397///
398/// # Errors
399///
400/// Returns an error if timestamp fields cannot be converted to nanoseconds.
401pub 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/// Parse a single book level into an order book delta.
465///
466/// # Errors
467///
468/// Returns an error if a non-delete action has a zero size after normalization.
469#[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; // Not applicable for L2 data
485    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; // Not available
492
493    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
509/// Parse a book snapshot message into a quote tick, returning an error on invalid data.
510/// Parse a book snapshot message into a quote tick.
511///
512/// # Errors
513///
514/// Returns an error if missing bid/ask levels or invalid sizes.
515pub 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
551/// Parse a trade message into a trade tick, returning an error on invalid data.
552/// Parse a trade message into a trade tick.
553///
554/// # Errors
555///
556/// Returns an error if invalid trade size is encountered.
557pub 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
591/// Parse a bar message into a Bar.
592///
593/// # Errors
594///
595/// Returns an error if the bar specification cannot be parsed.
596pub 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
618/// Extracts event and init timestamps from a derivative ticker message.
619fn 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
628/// Parses a derivative ticker message into a funding rate update.
629///
630/// # Errors
631///
632/// Returns an error if timestamp conversion or decimal conversion fails.
633pub 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
660/// Parses a derivative ticker message into a mark price update.
661///
662/// # Errors
663///
664/// Returns an error if timestamp conversion fails.
665pub 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
685/// Parses a derivative ticker message into an index price update.
686///
687/// # Errors
688///
689/// Returns an error if timestamp conversion fails.
690pub 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        // CLEAR delta
786        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        // First bid delta
794        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        // First ask delta
806        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        // Check first bid level
842        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        // Check second bid level
849        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        // Check first ask level
855        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        // Check second ask level
862        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        // Check empty levels are NULL_ORDER
868        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        // Absent greeks default to 0.0; absent IVs, underlying, and open interest stay None.
1296        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        // Test with minimal data (only funding_rate, no mark/index)
1366        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        // OKX X-Perps publish no predicted rate, so the funding timestamp is the only forward
1426        // reference carried through normalization
1427        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        // Captured either side of the 2023-04-10T08:40Z index migration, which switches the index
1468        // feed Tardis reads from BTC-USD to BTC-USDC but leaves the contract symbol alone, so both
1469        // must still resolve to one instrument. The index prices are a month apart and only pin
1470        // each capture, they do not themselves demonstrate the switch.
1471        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}