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