Skip to main content

nautilus_derive/websocket/
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
16//! Parsers for Derive public WebSocket subscription payloads.
17
18use anyhow::Context;
19use nautilus_core::{
20    UnixNanos,
21    datetime::{NANOSECONDS_IN_MILLISECOND, NANOSECONDS_IN_SECOND},
22};
23use nautilus_model::{
24    data::{
25        Bar, BarType, BookOrder, FundingRateUpdate, IndexPriceUpdate, MarkPriceUpdate,
26        OrderBookDelta, OrderBookDeltas, OrderBookDepth10, QuoteTick, TradeTick,
27        depth::DEPTH10_LEN, greeks::OptionGreekValues, option_chain::OptionGreeks,
28    },
29    enums::{AggressorSide, BarAggregation, BookAction, GreeksConvention, OrderSide, RecordFlag},
30    identifiers::{InstrumentId, TradeId},
31    types::{Price, Quantity},
32};
33use rust_decimal::prelude::ToPrimitive;
34
35use super::messages::{
36    DeriveOrderbookData, DeriveOrderbookLevel, DeriveOrderbookMsg, DerivePublicWsData,
37    DeriveTickerData, DeriveTickerMsg, DeriveTradesMsg, WsSubscriptionPayload,
38};
39use crate::{
40    common::{
41        enums::{DeriveLiquidityRole, DeriveOrderSide},
42        parse::format_instrument_id,
43    },
44    http::models::{
45        DerivePublicCandle, DerivePublicFundingRate, DerivePublicTrade, DeriveTickerSnapshot,
46    },
47};
48
49/// Parses a Derive public subscription payload into a typed market data update.
50///
51/// # Errors
52///
53/// Returns an error when the channel is unsupported or `params.data` does not
54/// match the channel payload shape.
55pub fn parse_public_ws_data(payload: &WsSubscriptionPayload) -> anyhow::Result<DerivePublicWsData> {
56    let channel = payload.channel.as_str();
57
58    if channel.starts_with("orderbook.") {
59        return parse_orderbook_msg(payload).map(DerivePublicWsData::Orderbook);
60    }
61
62    if channel.starts_with("trades.") {
63        return parse_trades_msg(payload).map(DerivePublicWsData::Trades);
64    }
65
66    if channel.starts_with("ticker_slim.") || channel.starts_with("ticker.") {
67        return parse_ticker_msg(payload).map(|msg| DerivePublicWsData::Ticker(Box::new(msg)));
68    }
69
70    anyhow::bail!("unsupported Derive public WS channel `{}`", payload.channel)
71}
72
73/// Parses an order book subscription payload.
74///
75/// # Errors
76///
77/// Returns an error when `params.data` is not a Derive order book snapshot.
78pub fn parse_orderbook_msg(payload: &WsSubscriptionPayload) -> anyhow::Result<DeriveOrderbookMsg> {
79    let data = serde_json::from_str::<DeriveOrderbookData>(payload.data.get())
80        .context("failed to decode Derive orderbook data")?;
81    Ok(DeriveOrderbookMsg {
82        channel: payload.channel,
83        data,
84    })
85}
86
87/// Parses a public trades subscription payload.
88///
89/// # Errors
90///
91/// Returns an error when `params.data` is not a list of Derive public trades.
92pub fn parse_trades_msg(payload: &WsSubscriptionPayload) -> anyhow::Result<DeriveTradesMsg> {
93    let trades = serde_json::from_str::<Vec<DerivePublicTrade>>(payload.data.get())
94        .context("failed to decode Derive trades data")?;
95    Ok(DeriveTradesMsg {
96        channel: payload.channel,
97        trades,
98    })
99}
100
101/// Parses a ticker subscription payload.
102///
103/// # Errors
104///
105/// Returns an error when `params.data` is not a Derive ticker payload.
106pub fn parse_ticker_msg(payload: &WsSubscriptionPayload) -> anyhow::Result<DeriveTickerMsg> {
107    let mut data = serde_json::from_str::<DeriveTickerData>(payload.data.get())
108        .context("failed to decode Derive ticker data")?;
109    data.apply_channel_context(payload.channel.as_str())
110        .map_err(anyhow::Error::msg)?;
111    Ok(DeriveTickerMsg {
112        channel: payload.channel,
113        data,
114    })
115}
116
117/// Parses an order book snapshot message into Nautilus snapshot deltas.
118///
119/// Derive's grouped order book stream sends a full depth snapshot for the
120/// requested grouping and depth, so the output starts with a clear delta and
121/// marks the last add with `F_LAST`. The payload does not include a change ID,
122/// so this uses the feed timestamp as the snapshot sequence.
123///
124/// Pass price and size precision from the instrument definition rather than
125/// inferring them from the wire values, since Derive may trim trailing zeroes.
126///
127/// # Errors
128///
129/// Returns an error when a price, size, or timestamp cannot be converted.
130pub fn parse_orderbook_deltas(
131    msg: &DeriveOrderbookMsg,
132    price_precision: u8,
133    size_precision: u8,
134    ts_init: UnixNanos,
135) -> anyhow::Result<OrderBookDeltas> {
136    let instrument_id = msg.data.instrument_id();
137    let timestamp =
138        u64::try_from(msg.data.timestamp).context("negative Derive orderbook timestamp")?;
139    let ts_event = timestamp_millis_to_nanos(timestamp, "timestamp")?;
140    let sequence = timestamp;
141    let context = BookDeltaContext {
142        instrument_id,
143        sequence,
144        price_precision,
145        size_precision,
146        ts_event,
147        ts_init,
148    };
149
150    let mut deltas = Vec::with_capacity(1 + msg.data.bids.len() + msg.data.asks.len());
151    let clear_flags = if msg.data.bids.is_empty() && msg.data.asks.is_empty() {
152        RecordFlag::F_SNAPSHOT as u8 | RecordFlag::F_LAST as u8
153    } else {
154        RecordFlag::F_SNAPSHOT as u8
155    };
156    deltas.push(OrderBookDelta::new_checked(
157        context.instrument_id,
158        BookAction::Clear,
159        BookOrder::default(),
160        clear_flags,
161        context.sequence,
162        context.ts_event,
163        context.ts_init,
164    )?);
165
166    for (idx, level) in msg.data.bids.iter().enumerate() {
167        push_level_delta(&mut deltas, &context, OrderSide::Buy, level, idx as u64)?;
168    }
169
170    let bid_count = msg.data.bids.len();
171    for (idx, level) in msg.data.asks.iter().enumerate() {
172        push_level_delta(
173            &mut deltas,
174            &context,
175            OrderSide::Sell,
176            level,
177            (bid_count + idx) as u64,
178        )?;
179    }
180
181    if let Some(last) = deltas.last_mut() {
182        last.flags |= RecordFlag::F_LAST as u8;
183    }
184
185    OrderBookDeltas::new_checked(context.instrument_id, deltas)
186}
187
188/// Parses an order book snapshot message into a fixed top-10 depth update.
189///
190/// Derive sends snapshots for the requested order book channel. Missing levels
191/// are filled with zero-size orders so the fixed arrays are always populated.
192///
193/// # Errors
194///
195/// Returns an error when a price, size, or timestamp cannot be converted.
196pub fn parse_orderbook_depth10(
197    msg: &DeriveOrderbookMsg,
198    price_precision: u8,
199    size_precision: u8,
200    ts_init: UnixNanos,
201) -> anyhow::Result<OrderBookDepth10> {
202    let instrument_id = msg.data.instrument_id();
203    let timestamp =
204        u64::try_from(msg.data.timestamp).context("negative Derive orderbook timestamp")?;
205    let ts_event = timestamp_millis_to_nanos(timestamp, "timestamp")?;
206
207    let mut bids = [BookOrder::default(); DEPTH10_LEN];
208    let mut asks = [BookOrder::default(); DEPTH10_LEN];
209    let mut bid_counts = [0; DEPTH10_LEN];
210    let mut ask_counts = [0; DEPTH10_LEN];
211
212    fill_depth_side(
213        &mut bids,
214        &mut bid_counts,
215        &msg.data.bids,
216        OrderSide::Buy,
217        price_precision,
218        size_precision,
219    )?;
220    fill_depth_side(
221        &mut asks,
222        &mut ask_counts,
223        &msg.data.asks,
224        OrderSide::Sell,
225        price_precision,
226        size_precision,
227    )?;
228
229    Ok(OrderBookDepth10::new(
230        instrument_id,
231        bids,
232        asks,
233        bid_counts,
234        ask_counts,
235        RecordFlag::F_SNAPSHOT as u8,
236        timestamp,
237        ts_event,
238        ts_init,
239    ))
240}
241
242/// Parses a public trade message into a Nautilus trade tick.
243///
244/// The public WS feed defines `direction` as the taker's direction, so it maps
245/// directly to the aggressor side.
246///
247/// Pass price and size precision from the instrument definition rather than
248/// inferring them from the wire values, since Derive may trim trailing zeroes.
249///
250/// # Errors
251///
252/// Returns an error when price, size, or timestamp conversion fails.
253pub fn parse_trade_tick(
254    trade: &DerivePublicTrade,
255    price_precision: u8,
256    size_precision: u8,
257    ts_init: UnixNanos,
258) -> anyhow::Result<TradeTick> {
259    let aggressor_side = match trade.direction {
260        DeriveOrderSide::Buy => AggressorSide::Buy,
261        DeriveOrderSide::Sell => AggressorSide::Sell,
262    };
263    build_trade_tick(
264        trade,
265        aggressor_side,
266        price_precision,
267        size_precision,
268        ts_init,
269    )
270}
271
272/// Parses a REST `public/get_trade_history` row into a Nautilus trade tick.
273///
274/// The endpoint returns one maker row and one taker row per trade under the
275/// same `trade_id`, and each row's `direction` is that participant's own side:
276/// the aggressor side is the taker row's direction and the inverse of the maker
277/// row's. A missing role keeps the public WS contract where `direction` already
278/// denotes the taker; an unknown role degrades the same way so the trade is
279/// still emitted.
280///
281/// Pass price and size precision from the instrument definition rather than
282/// inferring them from the wire values, since Derive may trim trailing zeroes.
283///
284/// # Errors
285///
286/// Returns an error when price, size, or timestamp conversion fails.
287pub fn parse_trade_tick_from_rest(
288    trade: &DerivePublicTrade,
289    price_precision: u8,
290    size_precision: u8,
291    ts_init: UnixNanos,
292) -> anyhow::Result<TradeTick> {
293    if trade.liquidity_role == Some(DeriveLiquidityRole::Unknown) {
294        log::warn!(
295            "Unknown Derive liquidity role for trade {}, treating direction as the taker side",
296            trade.trade_id,
297        );
298    }
299
300    let aggressor_side = match (trade.liquidity_role, trade.direction) {
301        (Some(DeriveLiquidityRole::Maker), DeriveOrderSide::Buy) => AggressorSide::Sell,
302        (Some(DeriveLiquidityRole::Maker), DeriveOrderSide::Sell) => AggressorSide::Buy,
303        (_, DeriveOrderSide::Buy) => AggressorSide::Buy,
304        (_, DeriveOrderSide::Sell) => AggressorSide::Sell,
305    };
306
307    build_trade_tick(
308        trade,
309        aggressor_side,
310        price_precision,
311        size_precision,
312        ts_init,
313    )
314}
315
316fn build_trade_tick(
317    trade: &DerivePublicTrade,
318    aggressor_side: AggressorSide,
319    price_precision: u8,
320    size_precision: u8,
321    ts_init: UnixNanos,
322) -> anyhow::Result<TradeTick> {
323    let instrument_id = format_instrument_id(trade.instrument_name.as_str());
324    let price = Price::from_decimal_dp(trade.trade_price, price_precision)
325        .with_context(|| format!("invalid trade price for {}", trade.instrument_name))?;
326    let size = Quantity::from_decimal_dp(trade.trade_amount, size_precision)
327        .with_context(|| format!("invalid trade amount for {}", trade.instrument_name))?;
328    let trade_id = TradeId::new(&trade.trade_id);
329    let timestamp = u64::try_from(trade.timestamp).context("negative Derive trade timestamp")?;
330    let ts_event = timestamp_millis_to_nanos(timestamp, "timestamp")?;
331
332    TradeTick::new_checked(
333        instrument_id,
334        price,
335        size,
336        aggressor_side,
337        trade_id,
338        ts_event,
339        ts_init,
340    )
341}
342
343/// Parses a ticker message into a Nautilus top-of-book quote.
344///
345/// Pass price and size precision from the instrument definition rather than
346/// inferring them from the wire values, since Derive may trim trailing zeroes.
347///
348/// # Errors
349///
350/// Returns an error when price, size, or timestamp conversion fails.
351pub fn parse_ticker_quote(
352    msg: &DeriveTickerMsg,
353    price_precision: u8,
354    size_precision: u8,
355    ts_init: UnixNanos,
356) -> anyhow::Result<QuoteTick> {
357    let instrument_id = msg.data.instrument_id();
358    let instrument_name = msg.data.instrument_name().as_str();
359    let bid_price = Price::from_decimal_dp(msg.data.best_bid_price(), price_precision)
360        .with_context(|| format!("invalid bid price for {instrument_name}"))?;
361    let ask_price = Price::from_decimal_dp(msg.data.best_ask_price(), price_precision)
362        .with_context(|| format!("invalid ask price for {instrument_name}"))?;
363    let bid_size = Quantity::from_decimal_dp(msg.data.best_bid_amount(), size_precision)
364        .with_context(|| format!("invalid bid amount for {instrument_name}"))?;
365    let ask_size = Quantity::from_decimal_dp(msg.data.best_ask_amount(), size_precision)
366        .with_context(|| format!("invalid ask amount for {instrument_name}"))?;
367    let timestamp =
368        u64::try_from(msg.data.timestamp()).context("negative Derive ticker timestamp")?;
369    let ts_event = timestamp_millis_to_nanos(timestamp, "timestamp")?;
370
371    QuoteTick::new_checked(
372        instrument_id,
373        bid_price,
374        ask_price,
375        bid_size,
376        ask_size,
377        ts_event,
378        ts_init,
379    )
380}
381
382/// Parses a REST `public/get_tickers` snapshot into a Nautilus top-of-book quote.
383///
384/// # Errors
385///
386/// Returns an error when price, size, or timestamp conversion fails.
387pub fn parse_ticker_quote_from_rest(
388    ticker: &DeriveTickerSnapshot,
389    price_precision: u8,
390    size_precision: u8,
391    ts_init: UnixNanos,
392) -> anyhow::Result<QuoteTick> {
393    let instrument_id = format_instrument_id(ticker.instrument_name.as_str());
394    let instrument_name = ticker.instrument_name.as_str();
395    let bid_price = Price::from_decimal_dp(ticker.best_bid_price, price_precision)
396        .with_context(|| format!("invalid bid price for {instrument_name}"))?;
397    let ask_price = Price::from_decimal_dp(ticker.best_ask_price, price_precision)
398        .with_context(|| format!("invalid ask price for {instrument_name}"))?;
399    let bid_size = Quantity::from_decimal_dp(ticker.best_bid_amount, size_precision)
400        .with_context(|| format!("invalid bid amount for {instrument_name}"))?;
401    let ask_size = Quantity::from_decimal_dp(ticker.best_ask_amount, size_precision)
402        .with_context(|| format!("invalid ask amount for {instrument_name}"))?;
403    let timestamp = u64::try_from(ticker.timestamp).context("negative Derive ticker timestamp")?;
404    let ts_event = timestamp_millis_to_nanos(timestamp, "timestamp")?;
405
406    QuoteTick::new_checked(
407        instrument_id,
408        bid_price,
409        ask_price,
410        bid_size,
411        ask_size,
412        ts_event,
413        ts_init,
414    )
415}
416
417#[derive(Debug, Clone, Copy)]
418struct BookDeltaContext {
419    instrument_id: InstrumentId,
420    sequence: u64,
421    price_precision: u8,
422    size_precision: u8,
423    ts_event: UnixNanos,
424    ts_init: UnixNanos,
425}
426
427fn push_level_delta(
428    deltas: &mut Vec<OrderBookDelta>,
429    context: &BookDeltaContext,
430    side: OrderSide,
431    level: &DeriveOrderbookLevel,
432    order_id: u64,
433) -> anyhow::Result<()> {
434    if level.amount().is_zero() {
435        return Ok(());
436    }
437
438    let price = Price::from_decimal_dp(level.price(), context.price_precision)
439        .context("invalid Derive orderbook price")?;
440    let size = Quantity::from_decimal_dp(level.amount(), context.size_precision)
441        .context("invalid Derive orderbook amount")?;
442    let order = BookOrder::new(side, price, size, order_id);
443    deltas.push(OrderBookDelta::new_checked(
444        context.instrument_id,
445        BookAction::Add,
446        order,
447        RecordFlag::F_SNAPSHOT as u8,
448        context.sequence,
449        context.ts_event,
450        context.ts_init,
451    )?);
452    Ok(())
453}
454
455fn fill_depth_side(
456    orders: &mut [BookOrder; DEPTH10_LEN],
457    counts: &mut [u32; DEPTH10_LEN],
458    levels: &[DeriveOrderbookLevel],
459    side: OrderSide,
460    price_precision: u8,
461    size_precision: u8,
462) -> anyhow::Result<()> {
463    let mut index = 0;
464
465    for level in levels {
466        let price = Price::from_decimal_dp(level.price(), price_precision)
467            .context("invalid Derive orderbook price")?;
468        let size = Quantity::from_decimal_dp(level.amount(), size_precision)
469            .context("invalid Derive orderbook amount")?;
470
471        if size.is_zero() {
472            continue;
473        }
474
475        orders[index] = BookOrder::new(side, price, size, 0);
476        counts[index] = 1;
477        index += 1;
478
479        if index == DEPTH10_LEN {
480            break;
481        }
482    }
483
484    for order in orders.iter_mut().skip(index) {
485        *order = BookOrder::new(
486            side,
487            Price::zero(price_precision),
488            Quantity::zero(size_precision),
489            0,
490        );
491    }
492
493    Ok(())
494}
495
496fn timestamp_millis_to_nanos(value: u64, field: &str) -> anyhow::Result<UnixNanos> {
497    let nanos = value
498        .checked_mul(NANOSECONDS_IN_MILLISECOND)
499        .with_context(|| format!("Derive {field} overflows nanoseconds"))?;
500    Ok(UnixNanos::from(nanos))
501}
502
503pub(crate) fn ticker_ts_event(timestamp_ms: i64) -> anyhow::Result<UnixNanos> {
504    let timestamp = u64::try_from(timestamp_ms).context("negative Derive ticker timestamp")?;
505    timestamp_millis_to_nanos(timestamp, "timestamp")
506}
507
508/// Parses a ticker payload into a [`MarkPriceUpdate`].
509///
510/// # Errors
511///
512/// Returns an error when the ticker timestamp is negative or overflows.
513pub fn parse_mark_price(
514    msg: &DeriveTickerMsg,
515    price_precision: u8,
516    ts_init: UnixNanos,
517) -> anyhow::Result<Option<MarkPriceUpdate>> {
518    let instrument_id = msg.data.instrument_id();
519    let value = Price::from_decimal_dp(msg.data.mark_price(), price_precision)
520        .with_context(|| format!("invalid Derive mark price for {instrument_id}"))?;
521    let ts_event = ticker_ts_event(msg.data.timestamp())?;
522    Ok(Some(MarkPriceUpdate::new(
523        instrument_id,
524        value,
525        ts_event,
526        ts_init,
527    )))
528}
529
530/// Parses a ticker payload into an [`IndexPriceUpdate`].
531///
532/// # Errors
533///
534/// Returns an error when the ticker timestamp is negative or overflows.
535pub fn parse_index_price(
536    msg: &DeriveTickerMsg,
537    price_precision: u8,
538    ts_init: UnixNanos,
539) -> anyhow::Result<Option<IndexPriceUpdate>> {
540    let instrument_id = msg.data.instrument_id();
541    let value = Price::from_decimal_dp(msg.data.index_price(), price_precision)
542        .with_context(|| format!("invalid Derive index price for {instrument_id}"))?;
543    let ts_event = ticker_ts_event(msg.data.timestamp())?;
544    Ok(Some(IndexPriceUpdate::new(
545        instrument_id,
546        value,
547        ts_event,
548        ts_init,
549    )))
550}
551
552/// Parses a perpetual ticker payload into a [`FundingRateUpdate`].
553///
554/// Returns `Ok(None)` when the ticker does not carry funding.
555///
556/// # Errors
557///
558/// Returns an error when the ticker timestamp is negative or overflows.
559pub fn parse_funding_rate(
560    msg: &DeriveTickerMsg,
561    ts_init: UnixNanos,
562) -> anyhow::Result<Option<FundingRateUpdate>> {
563    let Some(rate) = msg.data.funding_rate() else {
564        return Ok(None);
565    };
566    let instrument_id = msg.data.instrument_id();
567    let ts_event = ticker_ts_event(msg.data.timestamp())?;
568    Ok(Some(FundingRateUpdate::new(
569        instrument_id,
570        rate,
571        None,
572        None,
573        ts_event,
574        ts_init,
575    )))
576}
577
578/// Parses a `public/get_funding_rate_history` record into a [`FundingRateUpdate`].
579///
580/// # Errors
581///
582/// Returns an error when the record timestamp is negative or overflows.
583pub fn parse_funding_rate_history_record(
584    record: &DerivePublicFundingRate,
585    instrument_id: InstrumentId,
586    interval: Option<u16>,
587    ts_init: UnixNanos,
588) -> anyhow::Result<FundingRateUpdate> {
589    let ts_event = ticker_ts_event(record.timestamp)?;
590    Ok(FundingRateUpdate::new(
591        instrument_id,
592        record.funding_rate,
593        interval,
594        None,
595        ts_event,
596        ts_init,
597    ))
598}
599
600/// Parses a `public/get_tradingview_chart_data` record into a Nautilus [`Bar`].
601///
602/// Pass price and size precision from the instrument definition rather than
603/// inferring them from the wire values. The Derive `timestamp_bucket` is the
604/// bucket start in UNIX seconds; the returned bar's `ts_event` marks that
605/// bucket's close.
606///
607/// # Errors
608///
609/// Returns an error when price, size, or timestamp conversion fails.
610pub fn parse_candle_record(
611    record: &DerivePublicCandle,
612    bar_type: BarType,
613    price_precision: u8,
614    size_precision: u8,
615    ts_init: UnixNanos,
616) -> anyhow::Result<Bar> {
617    let open = Price::from_decimal_dp(record.open_price, price_precision)
618        .context("invalid Derive candle open price")?;
619    let high = Price::from_decimal_dp(record.high_price, price_precision)
620        .context("invalid Derive candle high price")?;
621    let low = Price::from_decimal_dp(record.low_price, price_precision)
622        .context("invalid Derive candle low price")?;
623    let close = Price::from_decimal_dp(record.close_price, price_precision)
624        .context("invalid Derive candle close price")?;
625    let volume = Quantity::from_decimal_dp(record.volume_contracts, size_precision)
626        .context("invalid Derive candle volume")?;
627    let timestamp =
628        u64::try_from(record.timestamp_bucket).context("negative Derive candle timestamp")?;
629    let bucket_start = timestamp_seconds_to_nanos(timestamp, "candle timestamp_bucket")?;
630    let interval_ns = bar_type.spec().timedelta().as_nanos();
631    let interval_ns = u64::try_from(interval_ns)
632        .context("bar interval overflowed the u64 range for nanoseconds")?;
633    let ts_event = bucket_start
634        .checked_add(interval_ns)
635        .context("bar timestamp overflowed when adjusting to close time")?;
636
637    Bar::new_checked(bar_type, open, high, low, close, volume, ts_event, ts_init)
638        .context("failed to construct Bar from Derive candle record")
639}
640
641/// Maps a Nautilus bar aggregation and step to the Derive `period` enum value
642/// (bucket size in seconds).
643///
644/// Derive supports the following bucket sizes: 60, 300, 900, 1800, 3600,
645/// 14400, 28800, 86400, 604800.
646///
647/// # Errors
648///
649/// Returns an error if the aggregation or step has no Derive equivalent.
650pub fn bar_spec_to_derive_period(aggregation: BarAggregation, step: u64) -> anyhow::Result<u32> {
651    match aggregation {
652        BarAggregation::Minute => match step {
653            1 => Ok(60),
654            5 => Ok(300),
655            15 => Ok(900),
656            30 => Ok(1800),
657            _ => anyhow::bail!(
658                "Derive only supports minute intervals 1, 5, 15, 30 (use HOUR for >= 60)"
659            ),
660        },
661        BarAggregation::Hour => match step {
662            1 => Ok(3600),
663            4 => Ok(14400),
664            8 => Ok(28800),
665            _ => anyhow::bail!("Derive only supports hour intervals 1, 4, 8"),
666        },
667        BarAggregation::Day => {
668            if step != 1 {
669                anyhow::bail!("Derive only supports 1 DAY interval bars");
670            }
671            Ok(86400)
672        }
673        BarAggregation::Week => {
674            if step != 1 {
675                anyhow::bail!("Derive only supports 1 WEEK interval bars");
676            }
677            Ok(604800)
678        }
679        _ => anyhow::bail!("Derive does not support {aggregation:?} bars"),
680    }
681}
682
683fn timestamp_seconds_to_nanos(value: u64, field: &str) -> anyhow::Result<UnixNanos> {
684    let nanos = value
685        .checked_mul(NANOSECONDS_IN_SECOND)
686        .with_context(|| format!("Derive {field} overflows nanoseconds"))?;
687    Ok(UnixNanos::from(nanos))
688}
689
690/// Parses an option ticker payload into [`OptionGreeks`].
691///
692/// Returns `Ok(None)` when the ticker does not carry option pricing.
693///
694/// # Errors
695///
696/// Returns an error when the ticker timestamp is negative or overflows.
697pub fn parse_option_greeks(
698    msg: &DeriveTickerMsg,
699    ts_init: UnixNanos,
700) -> anyhow::Result<Option<OptionGreeks>> {
701    let Some(pricing) = msg.data.option_pricing() else {
702        return Ok(None);
703    };
704    let instrument_id = msg.data.instrument_id();
705    let ts_event = ticker_ts_event(msg.data.timestamp())?;
706    let to_f64 = |label: &str, value: rust_decimal::Decimal| {
707        value
708            .to_f64()
709            .ok_or_else(|| anyhow::anyhow!("Derive {label} cannot be represented as f64"))
710    };
711
712    Ok(Some(OptionGreeks {
713        instrument_id,
714        convention: GreeksConvention::BlackScholes,
715        greeks: OptionGreekValues {
716            delta: to_f64("delta", pricing.delta)?,
717            gamma: to_f64("gamma", pricing.gamma)?,
718            vega: to_f64("vega", pricing.vega)?,
719            theta: to_f64("theta", pricing.theta)?,
720            rho: to_f64("rho", pricing.rho)?,
721        },
722        mark_iv: Some(to_f64("iv", pricing.iv)?),
723        bid_iv: Some(to_f64("bid_iv", pricing.bid_iv)?),
724        ask_iv: Some(to_f64("ask_iv", pricing.ask_iv)?),
725        underlying_price: Some(to_f64("forward_price", pricing.forward_price)?),
726        open_interest: msg
727            .data
728            .stats()
729            .map(|s| to_f64("open_interest", s.open_interest))
730            .transpose()?,
731        ts_event,
732        ts_init,
733    }))
734}
735
736#[cfg(test)]
737mod tests {
738    use std::{path::PathBuf, str::FromStr};
739
740    use nautilus_model::{
741        enums::{AggressorSide, BookAction, OrderSide, RecordFlag},
742        identifiers::{InstrumentId, TradeId},
743        types::{Price, Quantity},
744    };
745    use rstest::rstest;
746    use rust_decimal::Decimal;
747    use serde_json::{Value, json};
748    use ustr::Ustr;
749
750    use super::*;
751    use crate::websocket::messages::DeriveWsFrame;
752
753    const PRICE_PRECISION: u8 = 2;
754    const SIZE_PRECISION: u8 = 3;
755    const INVALID_PRECISION: u8 = u8::MAX;
756
757    fn data_path() -> PathBuf {
758        PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("test_data")
759    }
760
761    fn load_json(filename: &str) -> Value {
762        let content = std::fs::read_to_string(data_path().join(filename))
763            .unwrap_or_else(|_| panic!("failed to read {filename}"));
764        serde_json::from_str(&content).expect("invalid json")
765    }
766
767    fn subscription_payload(frame: &Value) -> WsSubscriptionPayload {
768        match DeriveWsFrame::parse(&frame.to_string()).unwrap() {
769            DeriveWsFrame::Subscription(payload) => payload,
770            other => panic!("expected subscription frame, was {other:?}"),
771        }
772    }
773
774    fn subscription_data_payload(channel: &str, data: &Value) -> WsSubscriptionPayload {
775        subscription_payload(&json!({
776            "jsonrpc": "2.0",
777            "method": "subscription",
778            "params": {
779                "channel": channel,
780                "data": data
781            }
782        }))
783    }
784
785    fn orderbook_json(timestamp: i64, bids: &Value, asks: &Value) -> Value {
786        let mut value = load_json("perps/ws_orderbook_eth.json");
787        value["timestamp"] = json!(timestamp);
788        value["bids"] = bids.clone();
789        value["asks"] = asks.clone();
790        value
791    }
792
793    fn trade_json(timestamp: i64, direction: &str) -> Value {
794        trade_json_with_values(timestamp, direction, "3500.2", "0.25")
795    }
796
797    fn trade_json_with_values(
798        timestamp: i64,
799        direction: &str,
800        trade_price: &str,
801        trade_amount: &str,
802    ) -> Value {
803        let mut value = load_json("perps/ws_trade_eth.json");
804        value["direction"] = json!(direction);
805        value["timestamp"] = json!(timestamp);
806        value["trade_amount"] = json!(trade_amount);
807        value["trade_id"] = json!("trade-1");
808        value["trade_price"] = json!(trade_price);
809        value
810    }
811
812    fn fixture_trade(filename: &str) -> DerivePublicTrade {
813        serde_json::from_value(load_json(filename)).expect("invalid Derive public trade")
814    }
815
816    fn ticker_json_with_timestamp(timestamp: i64) -> Value {
817        let mut value = load_json("perps/ws_ticker_eth.json");
818        value["best_ask_amount"] = json!("1.20");
819        value["best_ask_price"] = json!("3501.00");
820        value["best_bid_amount"] = json!("0.80");
821        value["best_bid_price"] = json!("3499.50");
822        value["timestamp"] = json!(timestamp);
823        value
824    }
825
826    fn ticker_json() -> Value {
827        ticker_json_with_timestamp(1_700_000_000_000)
828    }
829
830    fn price(value: &str) -> Price {
831        Price::from_decimal_dp(Decimal::from_str(value).unwrap(), PRICE_PRECISION).unwrap()
832    }
833
834    fn quantity(value: &str) -> Quantity {
835        Quantity::from_decimal_dp(Decimal::from_str(value).unwrap(), SIZE_PRECISION).unwrap()
836    }
837
838    #[rstest]
839    fn test_parse_public_orderbook_frame() {
840        let payload = subscription_data_payload(
841            "orderbook.ETH-PERP.1.10",
842            &orderbook_json(
843                1_700_000_000_000,
844                &json!([["3499.50", "1.20"], ["3499.00", "0.40"]]),
845                &json!([["3501.00", "0.80"]]),
846            ),
847        );
848
849        let msg = parse_orderbook_msg(&payload).unwrap();
850        let deltas =
851            parse_orderbook_deltas(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
852                .unwrap();
853
854        assert_eq!(msg.channel.as_str(), "orderbook.ETH-PERP.1.10");
855        assert_eq!(
856            msg.data.instrument_id(),
857            InstrumentId::from("ETH-PERP.DERIVE")
858        );
859        assert_eq!(msg.data.bids[0].price().to_string(), "3499.50");
860        assert_eq!(deltas.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
861        assert_eq!(deltas.deltas.len(), 4);
862        assert_eq!(deltas.deltas[0].action, BookAction::Clear);
863        assert_eq!(deltas.deltas[1].order.side, OrderSide::Buy.into());
864        assert_eq!(deltas.deltas[1].order.price, price("3499.50"));
865        assert_eq!(deltas.deltas[1].order.size, quantity("1.20"));
866        assert_eq!(deltas.deltas[3].order.side, OrderSide::Sell.into());
867        assert_eq!(
868            deltas.deltas[3].flags,
869            RecordFlag::F_SNAPSHOT as u8 | RecordFlag::F_LAST as u8
870        );
871    }
872
873    #[rstest]
874    fn test_parse_public_trades_frame() {
875        let payload = subscription_data_payload(
876            "trades.perp.ETH",
877            &json!([trade_json(1_700_000_000_001, "buy")]),
878        );
879
880        let msg = parse_trades_msg(&payload).unwrap();
881        let tick = parse_trade_tick(
882            &msg.trades[0],
883            PRICE_PRECISION,
884            SIZE_PRECISION,
885            UnixNanos::from(456),
886        )
887        .unwrap();
888
889        assert_eq!(msg.channel.as_str(), "trades.perp.ETH");
890        assert_eq!(msg.trades.len(), 1);
891        assert_eq!(
892            format_instrument_id(msg.trades[0].instrument_name.as_str()),
893            InstrumentId::from("ETH-PERP.DERIVE")
894        );
895        assert_eq!(tick.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
896        assert_eq!(tick.price, price("3500.2"));
897        assert_eq!(tick.size, quantity("0.25"));
898        assert_eq!(tick.aggressor_side, AggressorSide::Buy);
899        assert_eq!(tick.trade_id, TradeId::from("trade-1"));
900        assert_eq!(tick.ts_event, UnixNanos::from(1_700_000_000_001_000_000));
901    }
902
903    #[rstest]
904    fn test_parse_public_ticker_frame() {
905        let payload = subscription_data_payload(
906            "ticker_slim.ETH-PERP.1000",
907            &load_json("perps/ws_ticker_slim_eth.json"),
908        );
909
910        let msg = parse_ticker_msg(&payload).unwrap();
911        let quote = parse_ticker_quote(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(789))
912            .unwrap();
913
914        assert_eq!(msg.channel.as_str(), "ticker_slim.ETH-PERP.1000");
915        assert_eq!(
916            msg.data.instrument_id(),
917            InstrumentId::from("ETH-PERP.DERIVE")
918        );
919        assert_eq!(msg.data.timestamp(), 1_779_953_796_714);
920        assert_eq!(quote.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
921        assert_eq!(quote.bid_price, price("1992.36"));
922        assert_eq!(quote.ask_price, price("1992.37"));
923        assert_eq!(quote.bid_size, quantity("1.505"));
924        assert_eq!(quote.ask_size, quantity("1.505"));
925        assert_eq!(quote.ts_event, UnixNanos::from(1_779_953_796_714_000_000));
926    }
927
928    #[rstest]
929    fn test_parse_spot_orderbook_frame() {
930        let mut data = load_json("spot/ws_orderbook_eth.json");
931        data["bids"] = json!([["2050.0", "1.20"], ["2049.5", "0.40"]]);
932        data["asks"] = json!([["2051.0", "0.80"]]);
933        let payload = subscription_data_payload("orderbook.ETH-USDC.1.10", &data);
934
935        let msg = parse_orderbook_msg(&payload).unwrap();
936        let deltas =
937            parse_orderbook_deltas(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
938                .unwrap();
939
940        assert_eq!(msg.channel.as_str(), "orderbook.ETH-USDC.1.10");
941        assert_eq!(
942            msg.data.instrument_id(),
943            InstrumentId::from("ETH-USDC.DERIVE")
944        );
945        assert_eq!(deltas.instrument_id, InstrumentId::from("ETH-USDC.DERIVE"));
946        assert_eq!(deltas.deltas.len(), 4);
947        assert_eq!(deltas.deltas[0].action, BookAction::Clear);
948        assert_eq!(deltas.deltas[1].order.side, OrderSide::Buy.into());
949        assert_eq!(deltas.deltas[1].order.price, price("2050.0"));
950        assert_eq!(deltas.deltas[1].order.size, quantity("1.20"));
951        assert_eq!(deltas.deltas[3].order.side, OrderSide::Sell.into());
952        assert_eq!(
953            deltas.deltas[3].flags,
954            RecordFlag::F_SNAPSHOT as u8 | RecordFlag::F_LAST as u8
955        );
956    }
957
958    #[rstest]
959    fn test_parse_spot_trades_frame() {
960        let payload = subscription_data_payload(
961            "trades.erc20.ETH",
962            &json!([load_json("spot/ws_trade_eth.json")]),
963        );
964
965        let msg = parse_trades_msg(&payload).unwrap();
966        let tick = parse_trade_tick(
967            &msg.trades[0],
968            PRICE_PRECISION,
969            SIZE_PRECISION,
970            UnixNanos::from(456),
971        )
972        .unwrap();
973
974        assert_eq!(msg.channel.as_str(), "trades.erc20.ETH");
975        assert_eq!(msg.trades.len(), 1);
976        assert_eq!(tick.instrument_id, InstrumentId::from("ETH-USDC.DERIVE"));
977        assert_eq!(tick.price, price("2050"));
978        assert_eq!(tick.size, quantity("0.1"));
979        assert_eq!(tick.aggressor_side, AggressorSide::Sell);
980        assert_eq!(
981            tick.trade_id,
982            TradeId::from("0445f96a-10fb-4fdc-a0f9-eed94a2f32e1")
983        );
984    }
985
986    #[rstest]
987    fn test_parse_spot_ticker_slim_frame_handles_null_funding() {
988        let payload = subscription_data_payload(
989            "ticker_slim.ETH-USDC.1000",
990            &load_json("spot/ws_ticker_slim_eth.json"),
991        );
992
993        let msg = parse_ticker_msg(&payload).unwrap();
994        let quote = parse_ticker_quote(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(789))
995            .unwrap();
996
997        assert_eq!(msg.channel.as_str(), "ticker_slim.ETH-USDC.1000");
998        assert_eq!(
999            msg.data.instrument_id(),
1000            InstrumentId::from("ETH-USDC.DERIVE")
1001        );
1002        assert_eq!(quote.instrument_id, InstrumentId::from("ETH-USDC.DERIVE"));
1003
1004        assert!(
1005            parse_funding_rate(&msg, UnixNanos::from(789))
1006                .unwrap()
1007                .is_none()
1008        );
1009        let mark = parse_mark_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1010            .unwrap()
1011            .expect("spot slim ticker carries mark price");
1012        let index = parse_index_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1013            .unwrap()
1014            .expect("spot slim ticker carries index price");
1015        assert_eq!(mark.instrument_id, InstrumentId::from("ETH-USDC.DERIVE"));
1016        assert_eq!(index.instrument_id, InstrumentId::from("ETH-USDC.DERIVE"));
1017    }
1018
1019    #[rstest]
1020    fn test_parse_public_ticker_direct_payload() {
1021        let payload = subscription_data_payload(
1022            "ticker.ETH-PERP.1000",
1023            &ticker_json_with_timestamp(1_700_000_000_011),
1024        );
1025
1026        let msg = parse_ticker_msg(&payload).unwrap();
1027        let quote = parse_ticker_quote(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(790))
1028            .unwrap();
1029
1030        assert_eq!(msg.channel.as_str(), "ticker.ETH-PERP.1000");
1031        assert_eq!(msg.data.timestamp(), 1_700_000_000_011);
1032        assert_eq!(
1033            msg.data.instrument_id(),
1034            InstrumentId::from("ETH-PERP.DERIVE")
1035        );
1036        assert_eq!(quote.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
1037        assert_eq!(quote.ts_event, UnixNanos::from(1_700_000_000_011_000_000));
1038    }
1039
1040    #[rstest]
1041    fn test_parse_ticker_quote_uses_supplied_precision_when_wire_scale_varies() {
1042        let mut ticker = ticker_json_with_timestamp(1_700_000_000_012);
1043        ticker["best_bid_price"] = json!("3500");
1044        ticker["best_ask_price"] = json!("3501");
1045        ticker["best_bid_amount"] = json!("1");
1046        ticker["best_ask_amount"] = json!("2");
1047        let payload = subscription_data_payload("ticker.ETH-PERP.1000", &ticker);
1048
1049        let msg = parse_ticker_msg(&payload).unwrap();
1050        let quote = parse_ticker_quote(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(790))
1051            .unwrap();
1052
1053        assert_eq!(quote.bid_price, price("3500"));
1054        assert_eq!(quote.ask_price, price("3501"));
1055        assert_eq!(quote.bid_size, quantity("1"));
1056        assert_eq!(quote.ask_size, quantity("2"));
1057        assert_eq!(quote.bid_price.precision, PRICE_PRECISION);
1058        assert_eq!(quote.bid_size.precision, SIZE_PRECISION);
1059    }
1060
1061    #[rstest]
1062    fn test_parse_ticker_quote_from_rest_emits_quote() {
1063        let ticker: DeriveTickerSnapshot =
1064            serde_json::from_value(ticker_json_with_timestamp(1_700_000_000_013)).unwrap();
1065
1066        let quote = parse_ticker_quote_from_rest(
1067            &ticker,
1068            PRICE_PRECISION,
1069            SIZE_PRECISION,
1070            UnixNanos::from(791),
1071        )
1072        .unwrap();
1073
1074        assert_eq!(quote.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
1075        assert_eq!(quote.bid_price, price("3499.50"));
1076        assert_eq!(quote.ask_price, price("3501.00"));
1077        assert_eq!(quote.bid_size, quantity("0.80"));
1078        assert_eq!(quote.ask_size, quantity("1.20"));
1079        assert_eq!(quote.ts_event, UnixNanos::from(1_700_000_000_013_000_000));
1080    }
1081
1082    #[rstest]
1083    fn test_parse_ticker_quote_from_rest_rejects_negative_timestamp() {
1084        let mut value = ticker_json_with_timestamp(1_700_000_000_013);
1085        value["timestamp"] = json!(-1_i64);
1086        let ticker: DeriveTickerSnapshot = serde_json::from_value(value).unwrap();
1087
1088        let err = parse_ticker_quote_from_rest(
1089            &ticker,
1090            PRICE_PRECISION,
1091            SIZE_PRECISION,
1092            UnixNanos::from(791),
1093        )
1094        .expect_err("must reject negative timestamp");
1095        assert!(err.to_string().contains("negative Derive ticker timestamp"));
1096    }
1097
1098    #[rstest]
1099    fn test_parse_orderbook_deltas_empty_book_marks_clear_last() {
1100        let payload = subscription_data_payload(
1101            "orderbook.ETH-PERP.1.10",
1102            &orderbook_json(1_700_000_000_000, &json!([]), &json!([])),
1103        );
1104
1105        let msg = parse_orderbook_msg(&payload).unwrap();
1106        let deltas =
1107            parse_orderbook_deltas(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
1108                .unwrap();
1109
1110        assert_eq!(deltas.deltas.len(), 1);
1111        assert_eq!(deltas.deltas[0].action, BookAction::Clear);
1112        assert_eq!(
1113            deltas.deltas[0].flags,
1114            RecordFlag::F_SNAPSHOT as u8 | RecordFlag::F_LAST as u8
1115        );
1116    }
1117
1118    #[rstest]
1119    fn test_parse_orderbook_deltas_skips_zero_size_levels() {
1120        let payload = subscription_data_payload(
1121            "orderbook.ETH-PERP.1.10",
1122            &orderbook_json(
1123                1_700_000_000_000,
1124                &json!([["3499.50", "0"], ["3499.00", "0.40"]]),
1125                &json!([["3501.00", "0"]]),
1126            ),
1127        );
1128
1129        let msg = parse_orderbook_msg(&payload).unwrap();
1130        let deltas =
1131            parse_orderbook_deltas(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
1132                .unwrap();
1133
1134        assert_eq!(deltas.deltas.len(), 2);
1135        assert_eq!(deltas.deltas[1].order.side, OrderSide::Buy.into());
1136        assert_eq!(deltas.deltas[1].order.price, price("3499.00"));
1137        assert_eq!(deltas.deltas[1].order.size, quantity("0.40"));
1138        assert_eq!(deltas.deltas[1].order.order_id, 1);
1139        assert_eq!(
1140            deltas.deltas[1].flags,
1141            RecordFlag::F_SNAPSHOT as u8 | RecordFlag::F_LAST as u8
1142        );
1143    }
1144
1145    #[rstest]
1146    fn test_parse_orderbook_deltas_uses_supplied_precision_when_wire_scale_varies() {
1147        let payload = subscription_data_payload(
1148            "orderbook.ETH-PERP.1.10",
1149            &orderbook_json(
1150                1_700_000_000_000,
1151                &json!([["3500", "1"]]),
1152                &json!([["3501", "2"]]),
1153            ),
1154        );
1155
1156        let msg = parse_orderbook_msg(&payload).unwrap();
1157        let deltas =
1158            parse_orderbook_deltas(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
1159                .unwrap();
1160
1161        assert_eq!(deltas.deltas[1].order.price, price("3500"));
1162        assert_eq!(deltas.deltas[1].order.size, quantity("1"));
1163        assert_eq!(deltas.deltas[2].order.price, price("3501"));
1164        assert_eq!(deltas.deltas[2].order.size, quantity("2"));
1165        assert_eq!(deltas.deltas[1].order.price.precision, PRICE_PRECISION);
1166        assert_eq!(deltas.deltas[1].order.size.precision, SIZE_PRECISION);
1167    }
1168
1169    #[rstest]
1170    fn test_parse_orderbook_depth10_skips_zero_sizes_caps_and_zero_fills() {
1171        let bids = Value::Array(
1172            (0..12)
1173                .map(|i| {
1174                    let size = if i == 1 { "0" } else { "1" };
1175                    json!([format!("{}", 3500 - i), size])
1176                })
1177                .collect(),
1178        );
1179        let asks = json!([["3501", "2"], ["3502", "0"], ["3503", "3"]]);
1180        let payload = subscription_data_payload(
1181            "orderbook.ETH-PERP.1.10",
1182            &orderbook_json(1_700_000_000_000, &bids, &asks),
1183        );
1184
1185        let msg = parse_orderbook_msg(&payload).unwrap();
1186        let depth =
1187            parse_orderbook_depth10(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
1188                .unwrap();
1189
1190        assert_eq!(depth.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
1191        assert_eq!(depth.bids[0].price, price("3500"));
1192        assert_eq!(depth.bids[1].price, price("3498"));
1193        assert_eq!(depth.bids[9].price, price("3490"));
1194        assert_eq!(depth.bid_counts[0], 1);
1195        assert_eq!(depth.bid_counts[9], 1);
1196        assert_eq!(depth.asks[0].price, price("3501"));
1197        assert_eq!(depth.asks[1].price, price("3503"));
1198        assert_eq!(depth.asks[2].price, Price::zero(PRICE_PRECISION));
1199        assert_eq!(depth.asks[2].size, Quantity::zero(SIZE_PRECISION));
1200        assert_eq!(depth.ask_counts[0], 1);
1201        assert_eq!(depth.ask_counts[1], 1);
1202        assert_eq!(depth.ask_counts[2], 0);
1203        assert_eq!(depth.sequence, 1_700_000_000_000);
1204        assert_eq!(depth.flags, RecordFlag::F_SNAPSHOT as u8);
1205        assert_eq!(depth.ts_event, UnixNanos::from(1_700_000_000_000_000_000));
1206    }
1207
1208    #[rstest]
1209    fn test_parse_trade_tick_maps_sell_direction() {
1210        let payload = subscription_data_payload(
1211            "trades.perp.ETH",
1212            &json!([trade_json(1_700_000_000_001, "sell")]),
1213        );
1214
1215        let msg = parse_trades_msg(&payload).unwrap();
1216        let tick = parse_trade_tick(
1217            &msg.trades[0],
1218            PRICE_PRECISION,
1219            SIZE_PRECISION,
1220            UnixNanos::from(456),
1221        )
1222        .unwrap();
1223
1224        assert_eq!(tick.aggressor_side, AggressorSide::Sell);
1225    }
1226
1227    #[rstest]
1228    fn test_parse_trade_tick_uses_supplied_precision_when_wire_scale_varies() {
1229        let payload = subscription_data_payload(
1230            "trades.perp.ETH",
1231            &json!([trade_json_with_values(
1232                1_700_000_000_001,
1233                "buy",
1234                "3500",
1235                "1"
1236            )]),
1237        );
1238
1239        let msg = parse_trades_msg(&payload).unwrap();
1240        let tick = parse_trade_tick(
1241            &msg.trades[0],
1242            PRICE_PRECISION,
1243            SIZE_PRECISION,
1244            UnixNanos::from(456),
1245        )
1246        .unwrap();
1247
1248        assert_eq!(tick.price, price("3500"));
1249        assert_eq!(tick.size, quantity("1"));
1250        assert_eq!(tick.price.precision, PRICE_PRECISION);
1251        assert_eq!(tick.size.precision, SIZE_PRECISION);
1252    }
1253
1254    #[rstest]
1255    #[case("buy", AggressorSide::Sell)]
1256    #[case("sell", AggressorSide::Buy)]
1257    fn test_parse_trade_tick_from_rest_inverts_maker_row(
1258        #[case] direction: &str,
1259        #[case] expected: AggressorSide,
1260    ) {
1261        let mut value = load_json("perps/http_public_trade_eth_maker.json");
1262        value["direction"] = json!(direction);
1263        let trade: DerivePublicTrade =
1264            serde_json::from_value(value).expect("invalid Derive public trade");
1265
1266        let tick = parse_trade_tick_from_rest(
1267            &trade,
1268            PRICE_PRECISION,
1269            SIZE_PRECISION,
1270            UnixNanos::from(456),
1271        )
1272        .unwrap();
1273
1274        assert_eq!(tick.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
1275        assert_eq!(tick.price, price("3499.0"));
1276        assert_eq!(tick.size, quantity("0.5"));
1277        assert_eq!(tick.aggressor_side, expected);
1278        assert_eq!(tick.trade_id, TradeId::from("trade-1"));
1279    }
1280
1281    #[rstest]
1282    fn test_parse_trade_tick_from_rest_maps_taker_row_directly() {
1283        let trade = fixture_trade("perps/http_public_trade_eth_sell.json");
1284
1285        let tick = parse_trade_tick_from_rest(
1286            &trade,
1287            PRICE_PRECISION,
1288            SIZE_PRECISION,
1289            UnixNanos::from(456),
1290        )
1291        .unwrap();
1292
1293        assert_eq!(tick.aggressor_side, AggressorSide::Sell);
1294        assert_eq!(tick.trade_id, TradeId::from("trade-1"));
1295    }
1296
1297    #[rstest]
1298    fn test_parse_trade_tick_from_rest_degrades_absent_role_to_taker_side() {
1299        let trade = fixture_trade("perps/ws_trade_eth_absent_role.json");
1300
1301        let tick = parse_trade_tick_from_rest(
1302            &trade,
1303            PRICE_PRECISION,
1304            SIZE_PRECISION,
1305            UnixNanos::from(456),
1306        )
1307        .unwrap();
1308
1309        assert_eq!(tick.aggressor_side, AggressorSide::Sell);
1310        assert_eq!(tick.trade_id, TradeId::from("perp-trade-1"));
1311    }
1312
1313    #[rstest]
1314    fn test_parse_trade_tick_from_rest_degrades_unknown_role_to_taker_side() {
1315        let trade = fixture_trade("perps/http_public_trade_eth_unknown_role.json");
1316
1317        let tick = parse_trade_tick_from_rest(
1318            &trade,
1319            PRICE_PRECISION,
1320            SIZE_PRECISION,
1321            UnixNanos::from(456),
1322        )
1323        .unwrap();
1324
1325        assert_eq!(tick.aggressor_side, AggressorSide::Sell);
1326        assert_eq!(tick.trade_id, TradeId::from("trade-1"));
1327    }
1328
1329    #[rstest]
1330    fn test_parse_trade_tick_maps_absent_role_directly() {
1331        let trade = fixture_trade("perps/ws_trade_eth_absent_role.json");
1332
1333        let tick = parse_trade_tick(
1334            &trade,
1335            PRICE_PRECISION,
1336            SIZE_PRECISION,
1337            UnixNanos::from(456),
1338        )
1339        .unwrap();
1340
1341        assert_eq!(tick.aggressor_side, AggressorSide::Sell);
1342        assert_eq!(tick.trade_id, TradeId::from("perp-trade-1"));
1343    }
1344
1345    #[rstest]
1346    fn test_parse_public_ws_data_dispatches_orderbook_channel() {
1347        let payload = subscription_data_payload(
1348            "orderbook.ETH-PERP.1.10",
1349            &orderbook_json(1_700_000_000_000, &json!([]), &json!([])),
1350        );
1351
1352        let parsed = parse_public_ws_data(&payload).unwrap();
1353
1354        match parsed {
1355            DerivePublicWsData::Orderbook(msg) => {
1356                assert_eq!(msg.channel.as_str(), "orderbook.ETH-PERP.1.10");
1357                assert_eq!(
1358                    msg.data.instrument_id(),
1359                    InstrumentId::from("ETH-PERP.DERIVE")
1360                );
1361            }
1362            other => panic!("expected orderbook data, was {other:?}"),
1363        }
1364    }
1365
1366    #[rstest]
1367    fn test_parse_public_ws_data_dispatches_trades_channel() {
1368        let payload = subscription_data_payload("trades.perp.ETH", &json!([]));
1369
1370        let parsed = parse_public_ws_data(&payload).unwrap();
1371
1372        match parsed {
1373            DerivePublicWsData::Trades(msg) => assert!(msg.trades.is_empty()),
1374            other => panic!("expected trades data, was {other:?}"),
1375        }
1376    }
1377
1378    #[rstest]
1379    fn test_parse_public_ws_data_dispatches_ticker_channel() {
1380        let payload = subscription_data_payload(
1381            "ticker_slim.ETH-PERP.1000",
1382            &load_json("perps/ws_ticker_slim_eth.json"),
1383        );
1384
1385        let parsed = parse_public_ws_data(&payload).unwrap();
1386
1387        match parsed {
1388            DerivePublicWsData::Ticker(msg) => {
1389                assert_eq!(msg.channel.as_str(), "ticker_slim.ETH-PERP.1000");
1390                assert_eq!(
1391                    msg.data.instrument_id(),
1392                    InstrumentId::from("ETH-PERP.DERIVE")
1393                );
1394            }
1395            other => panic!("expected ticker data, was {other:?}"),
1396        }
1397    }
1398
1399    #[rstest]
1400    fn test_parse_orderbook_msg_rejects_malformed_payload() {
1401        let payload = subscription_data_payload(
1402            "orderbook.ETH-PERP.1.10",
1403            &json!({
1404                "instrument_name": "ETH-PERP",
1405                "timestamp": 1_700_000_000_000_i64,
1406                "bids": []
1407            }),
1408        );
1409
1410        let err = parse_orderbook_msg(&payload).expect_err("must reject malformed orderbook");
1411
1412        assert!(
1413            err.to_string()
1414                .contains("failed to decode Derive orderbook data")
1415        );
1416    }
1417
1418    #[rstest]
1419    fn test_parse_trades_msg_rejects_malformed_payload() {
1420        let payload = subscription_data_payload("trades.perp.ETH", &json!({}));
1421
1422        let err = parse_trades_msg(&payload).expect_err("must reject malformed trades");
1423
1424        assert!(
1425            err.to_string()
1426                .contains("failed to decode Derive trades data")
1427        );
1428    }
1429
1430    #[rstest]
1431    fn test_parse_ticker_msg_rejects_malformed_payload() {
1432        let payload = subscription_data_payload(
1433            "ticker.ETH-PERP.1000",
1434            &json!({
1435                "timestamp": 1_700_000_000_010_i64
1436            }),
1437        );
1438
1439        let err = parse_ticker_msg(&payload).expect_err("must reject malformed ticker");
1440
1441        assert!(
1442            err.to_string()
1443                .contains("failed to decode Derive ticker data")
1444        );
1445    }
1446
1447    #[rstest]
1448    #[case("ticker_slim.ETH-PERP")]
1449    #[case("ticker_slim..1000")]
1450    fn test_parse_ticker_msg_rejects_malformed_slim_channel(#[case] channel: &str) {
1451        let payload =
1452            subscription_data_payload(channel, &load_json("perps/ws_ticker_slim_eth.json"));
1453
1454        let err = parse_ticker_msg(&payload).expect_err("must reject malformed slim channel");
1455
1456        assert!(err.to_string().contains("invalid Derive ticker channel"));
1457    }
1458
1459    #[rstest]
1460    fn test_parse_orderbook_deltas_rejects_negative_timestamp() {
1461        let payload = subscription_data_payload(
1462            "orderbook.ETH-PERP.1.10",
1463            &orderbook_json(-1, &json!([]), &json!([])),
1464        );
1465
1466        let msg = parse_orderbook_msg(&payload).unwrap();
1467        let err =
1468            parse_orderbook_deltas(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
1469                .expect_err("must reject negative orderbook timestamp");
1470
1471        assert!(
1472            err.to_string()
1473                .contains("negative Derive orderbook timestamp")
1474        );
1475    }
1476
1477    #[rstest]
1478    fn test_parse_orderbook_deltas_rejects_timestamp_overflow() {
1479        let payload = subscription_data_payload(
1480            "orderbook.ETH-PERP.1.10",
1481            &orderbook_json(i64::MAX, &json!([]), &json!([])),
1482        );
1483
1484        let msg = parse_orderbook_msg(&payload).unwrap();
1485        let err =
1486            parse_orderbook_deltas(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(123))
1487                .expect_err("must reject overflowing orderbook timestamp");
1488
1489        assert!(
1490            err.to_string()
1491                .contains("Derive timestamp overflows nanoseconds")
1492        );
1493    }
1494
1495    #[rstest]
1496    fn test_parse_orderbook_deltas_rejects_invalid_size_precision() {
1497        let payload = subscription_data_payload(
1498            "orderbook.ETH-PERP.1.10",
1499            &orderbook_json(
1500                1_700_000_000_000,
1501                &json!([["3500", "1"]]),
1502                &json!([["3501", "2"]]),
1503            ),
1504        );
1505
1506        let msg = parse_orderbook_msg(&payload).unwrap();
1507        let err = parse_orderbook_deltas(
1508            &msg,
1509            PRICE_PRECISION,
1510            INVALID_PRECISION,
1511            UnixNanos::from(123),
1512        )
1513        .expect_err("must reject invalid orderbook size precision");
1514
1515        assert!(err.to_string().contains("invalid Derive orderbook amount"));
1516    }
1517
1518    #[rstest]
1519    fn test_parse_trade_tick_rejects_negative_timestamp() {
1520        let payload = subscription_data_payload("trades.perp.ETH", &json!([trade_json(-1, "buy")]));
1521
1522        let msg = parse_trades_msg(&payload).unwrap();
1523        let err = parse_trade_tick(
1524            &msg.trades[0],
1525            PRICE_PRECISION,
1526            SIZE_PRECISION,
1527            UnixNanos::from(456),
1528        )
1529        .expect_err("must reject negative trade timestamp");
1530
1531        assert!(err.to_string().contains("negative Derive trade timestamp"));
1532    }
1533
1534    #[rstest]
1535    fn test_parse_trade_tick_rejects_timestamp_overflow() {
1536        let payload =
1537            subscription_data_payload("trades.perp.ETH", &json!([trade_json(i64::MAX, "buy")]));
1538
1539        let msg = parse_trades_msg(&payload).unwrap();
1540        let err = parse_trade_tick(
1541            &msg.trades[0],
1542            PRICE_PRECISION,
1543            SIZE_PRECISION,
1544            UnixNanos::from(456),
1545        )
1546        .expect_err("must reject overflowing trade timestamp");
1547
1548        assert!(
1549            err.to_string()
1550                .contains("Derive timestamp overflows nanoseconds")
1551        );
1552    }
1553
1554    #[rstest]
1555    fn test_parse_trade_tick_rejects_invalid_price_precision() {
1556        let payload = subscription_data_payload(
1557            "trades.perp.ETH",
1558            &json!([trade_json(1_700_000_000_001, "buy")]),
1559        );
1560
1561        let msg = parse_trades_msg(&payload).unwrap();
1562        let err = parse_trade_tick(
1563            &msg.trades[0],
1564            INVALID_PRECISION,
1565            SIZE_PRECISION,
1566            UnixNanos::from(456),
1567        )
1568        .expect_err("must reject invalid trade price precision");
1569
1570        assert!(err.to_string().contains("invalid trade price for ETH-PERP"));
1571    }
1572
1573    #[rstest]
1574    fn test_parse_ticker_quote_rejects_negative_timestamp() {
1575        let payload = subscription_data_payload(
1576            "ticker.ETH-PERP.1000",
1577            &json!({
1578                "timestamp": -1_i64,
1579                "instrument_ticker": ticker_json()
1580            }),
1581        );
1582
1583        let msg = parse_ticker_msg(&payload).unwrap();
1584        let err = parse_ticker_quote(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(789))
1585            .expect_err("must reject negative ticker timestamp");
1586
1587        assert!(err.to_string().contains("negative Derive ticker timestamp"));
1588    }
1589
1590    #[rstest]
1591    fn test_parse_ticker_quote_rejects_timestamp_overflow() {
1592        let payload = subscription_data_payload(
1593            "ticker.ETH-PERP.1000",
1594            &json!({
1595                "timestamp": i64::MAX,
1596                "instrument_ticker": ticker_json()
1597            }),
1598        );
1599
1600        let msg = parse_ticker_msg(&payload).unwrap();
1601        let err = parse_ticker_quote(&msg, PRICE_PRECISION, SIZE_PRECISION, UnixNanos::from(789))
1602            .expect_err("must reject overflowing ticker timestamp");
1603
1604        assert!(
1605            err.to_string()
1606                .contains("Derive timestamp overflows nanoseconds")
1607        );
1608    }
1609
1610    #[rstest]
1611    fn test_parse_public_ws_data_rejects_unknown_channel() {
1612        let payload = WsSubscriptionPayload {
1613            channel: Ustr::from("wallet.ETH"),
1614            data: serde_json::value::to_raw_value(&json!({})).unwrap(),
1615        };
1616
1617        let err = parse_public_ws_data(&payload).expect_err("must reject unknown channel");
1618
1619        assert!(
1620            err.to_string()
1621                .contains("unsupported Derive public WS channel")
1622        );
1623    }
1624
1625    fn option_ticker_json(timestamp: i64) -> Value {
1626        let mut value = load_json("options/http_ticker_eth_snapshot.json");
1627        value["timestamp"] = json!(timestamp);
1628        value
1629    }
1630
1631    fn perp_envelope_payload(timestamp: i64) -> WsSubscriptionPayload {
1632        subscription_data_payload(
1633            "ticker.ETH-PERP.1000",
1634            &json!({
1635                "timestamp": timestamp,
1636                "instrument_ticker": ticker_json_with_timestamp(timestamp),
1637            }),
1638        )
1639    }
1640
1641    fn option_envelope_payload(timestamp: i64) -> WsSubscriptionPayload {
1642        let mut option_data = option_ticker_json(timestamp);
1643        option_data["instrument_name"] = json!("ETH-20260627-3500-C");
1644        subscription_data_payload(
1645            "ticker.ETH-20260627-3500-C.1000",
1646            &json!({
1647                "timestamp": timestamp,
1648                "instrument_ticker": option_data,
1649            }),
1650        )
1651    }
1652
1653    fn slim_payload() -> WsSubscriptionPayload {
1654        subscription_data_payload(
1655            "ticker_slim.ETH-PERP.1000",
1656            &load_json("perps/ws_ticker_slim_eth.json"),
1657        )
1658    }
1659
1660    #[rstest]
1661    fn test_parse_mark_price_maps_slim_variant() {
1662        let msg = parse_ticker_msg(&slim_payload()).unwrap();
1663
1664        let update = parse_mark_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1665            .unwrap()
1666            .expect("slim ticker carries mark price");
1667
1668        assert_eq!(update.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
1669        assert_eq!(update.value, price("1992.49"));
1670        assert_eq!(update.ts_event, UnixNanos::from(1_779_953_796_714_000_000));
1671        assert_eq!(update.ts_init, UnixNanos::from(789));
1672    }
1673
1674    #[rstest]
1675    fn test_parse_index_price_maps_slim_variant() {
1676        let msg = parse_ticker_msg(&slim_payload()).unwrap();
1677
1678        let update = parse_index_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1679            .unwrap()
1680            .expect("slim ticker carries index price");
1681
1682        assert_eq!(update.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
1683        assert_eq!(update.value, price("1991.79"));
1684        assert_eq!(update.ts_event, UnixNanos::from(1_779_953_796_714_000_000));
1685        assert_eq!(update.ts_init, UnixNanos::from(789));
1686    }
1687
1688    #[rstest]
1689    fn test_parse_funding_rate_maps_slim_variant() {
1690        let msg = parse_ticker_msg(&slim_payload()).unwrap();
1691
1692        let update = parse_funding_rate(&msg, UnixNanos::from(789))
1693            .unwrap()
1694            .expect("slim ticker carries perp funding");
1695
1696        assert_eq!(update.instrument_id, InstrumentId::from("ETH-PERP.DERIVE"));
1697        assert_eq!(update.rate, Decimal::from_str("0.000012500").unwrap());
1698        assert_eq!(update.ts_event, UnixNanos::from(1_779_953_796_714_000_000));
1699        assert_eq!(update.ts_init, UnixNanos::from(789));
1700    }
1701
1702    #[rstest]
1703    fn test_parse_option_greeks_returns_none_for_slim_variant_without_option_pricing() {
1704        let msg = parse_ticker_msg(&slim_payload()).unwrap();
1705
1706        let result = parse_option_greeks(&msg, UnixNanos::from(789)).unwrap();
1707
1708        assert!(result.is_none());
1709    }
1710
1711    fn option_slim_payload(filename: &str, instrument_name: &str) -> WsSubscriptionPayload {
1712        subscription_data_payload(
1713            &format!("ticker_slim.{instrument_name}.1000"),
1714            &load_json(filename),
1715        )
1716    }
1717
1718    #[rstest]
1719    fn test_parse_option_greeks_maps_slim_variant() {
1720        let msg = parse_ticker_msg(&option_slim_payload(
1721            "options/ws_ticker_slim_eth_call.json",
1722            "ETH-20260612-1600-C",
1723        ))
1724        .unwrap();
1725
1726        let greeks = parse_option_greeks(&msg, UnixNanos::from(789))
1727            .unwrap()
1728            .expect("slim ticker carries option pricing");
1729
1730        assert_eq!(
1731            greeks.instrument_id,
1732            InstrumentId::from("ETH-20260612-1600-C.DERIVE")
1733        );
1734        assert_eq!(greeks.convention, GreeksConvention::BlackScholes);
1735        assert!((greeks.greeks.delta - 0.95222).abs() < 1e-9);
1736        assert!((greeks.greeks.gamma - 0.00036344).abs() < 1e-9);
1737        assert_eq!(greeks.mark_iv, Some(0.67698));
1738        assert_eq!(greeks.bid_iv, Some(0.0));
1739        assert_eq!(greeks.ask_iv, Some(0.88815));
1740        assert_eq!(greeks.underlying_price, Some(1992.6));
1741        assert_eq!(greeks.open_interest, Some(0.0));
1742        assert_eq!(greeks.ts_event, UnixNanos::from(1_779_953_796_231_000_000));
1743        assert_eq!(greeks.ts_init, UnixNanos::from(789));
1744    }
1745
1746    #[rstest]
1747    fn test_parse_option_greeks_maps_slim_put_variant() {
1748        let msg = parse_ticker_msg(&option_slim_payload(
1749            "options/ws_ticker_slim_eth_put.json",
1750            "ETH-20260612-1900-P",
1751        ))
1752        .unwrap();
1753
1754        let greeks = parse_option_greeks(&msg, UnixNanos::from(789))
1755            .unwrap()
1756            .expect("slim ticker carries put option pricing");
1757
1758        assert_eq!(
1759            greeks.instrument_id,
1760            InstrumentId::from("ETH-20260612-1900-P.DERIVE")
1761        );
1762        assert!((greeks.greeks.delta + 0.30438).abs() < 1e-9);
1763        assert!((greeks.greeks.gamma - 0.00169741).abs() < 1e-9);
1764        assert_eq!(greeks.mark_iv, Some(0.51012));
1765        assert_eq!(greeks.bid_iv, Some(0.48229));
1766        assert_eq!(greeks.ask_iv, Some(0.52063));
1767        assert_eq!(greeks.underlying_price, Some(1992.6));
1768        assert_eq!(greeks.open_interest, Some(42.13));
1769        assert_eq!(greeks.ts_event, UnixNanos::from(1_779_953_797_040_000_000));
1770        assert_eq!(greeks.ts_init, UnixNanos::from(789));
1771    }
1772
1773    #[rstest]
1774    fn test_parse_funding_rate_returns_none_for_option_payload() {
1775        let msg = parse_ticker_msg(&option_envelope_payload(1_700_000_000_010)).unwrap();
1776
1777        let result = parse_funding_rate(&msg, UnixNanos::from(789)).unwrap();
1778
1779        assert!(result.is_none());
1780    }
1781
1782    #[rstest]
1783    fn test_parse_option_greeks_returns_none_for_perp_payload() {
1784        let msg = parse_ticker_msg(&perp_envelope_payload(1_700_000_000_010)).unwrap();
1785
1786        let result = parse_option_greeks(&msg, UnixNanos::from(789)).unwrap();
1787
1788        assert!(result.is_none());
1789    }
1790
1791    #[rstest]
1792    fn test_parse_option_greeks_open_interest_none_when_stats_absent() {
1793        // Legacy full ticker payloads may omit `stats`. When the WS path
1794        // receives one without stats, `open_interest` must degrade to
1795        // None while the remaining greek fields still populate normally.
1796        let timestamp = 1_700_000_000_010_i64;
1797        let mut option_data = option_ticker_json(timestamp);
1798        option_data["instrument_name"] = json!("ETH-20260627-3500-C");
1799        option_data["stats"] = json!(null);
1800        let payload = subscription_data_payload(
1801            "ticker.ETH-20260627-3500-C.1000",
1802            &json!({
1803                "timestamp": timestamp,
1804                "instrument_ticker": option_data,
1805            }),
1806        );
1807        let msg = parse_ticker_msg(&payload).unwrap();
1808
1809        let greeks = parse_option_greeks(&msg, UnixNanos::from(789))
1810            .unwrap()
1811            .expect("option greeks present when option_pricing is set");
1812        assert!(greeks.open_interest.is_none());
1813        // The other greek fields must still be populated from option_pricing
1814        assert!((greeks.greeks.delta - 0.55).abs() < 1e-9);
1815        assert!(greeks.mark_iv.is_some());
1816        assert!(greeks.underlying_price.is_some());
1817    }
1818
1819    #[rstest]
1820    fn test_parse_mark_price_rejects_negative_timestamp() {
1821        let payload = subscription_data_payload(
1822            "ticker.ETH-PERP.1000",
1823            &json!({
1824                "timestamp": -1_i64,
1825                "instrument_ticker": ticker_json(),
1826            }),
1827        );
1828        let msg = parse_ticker_msg(&payload).unwrap();
1829
1830        let err = parse_mark_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1831            .expect_err("must reject negative ticker timestamp");
1832
1833        assert!(err.to_string().contains("negative Derive ticker timestamp"));
1834    }
1835
1836    #[rstest]
1837    fn test_parse_mark_price_rejects_timestamp_overflow() {
1838        let payload = subscription_data_payload(
1839            "ticker.ETH-PERP.1000",
1840            &json!({
1841                "timestamp": i64::MAX,
1842                "instrument_ticker": ticker_json(),
1843            }),
1844        );
1845        let msg = parse_ticker_msg(&payload).unwrap();
1846
1847        let err = parse_mark_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1848            .expect_err("must reject overflowing ticker timestamp");
1849
1850        assert!(
1851            err.to_string()
1852                .contains("Derive timestamp overflows nanoseconds")
1853        );
1854    }
1855
1856    #[rstest]
1857    fn test_parse_index_price_rejects_negative_timestamp() {
1858        let payload = subscription_data_payload(
1859            "ticker.ETH-PERP.1000",
1860            &json!({
1861                "timestamp": -1_i64,
1862                "instrument_ticker": ticker_json(),
1863            }),
1864        );
1865        let msg = parse_ticker_msg(&payload).unwrap();
1866
1867        let err = parse_index_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1868            .expect_err("must reject negative ticker timestamp");
1869
1870        assert!(err.to_string().contains("negative Derive ticker timestamp"));
1871    }
1872
1873    #[rstest]
1874    fn test_parse_index_price_rejects_timestamp_overflow() {
1875        let payload = subscription_data_payload(
1876            "ticker.ETH-PERP.1000",
1877            &json!({
1878                "timestamp": i64::MAX,
1879                "instrument_ticker": ticker_json(),
1880            }),
1881        );
1882        let msg = parse_ticker_msg(&payload).unwrap();
1883
1884        let err = parse_index_price(&msg, PRICE_PRECISION, UnixNanos::from(789))
1885            .expect_err("must reject overflowing ticker timestamp");
1886
1887        assert!(
1888            err.to_string()
1889                .contains("Derive timestamp overflows nanoseconds")
1890        );
1891    }
1892
1893    #[rstest]
1894    fn test_parse_funding_rate_rejects_negative_timestamp() {
1895        let payload = subscription_data_payload(
1896            "ticker.ETH-PERP.1000",
1897            &json!({
1898                "timestamp": -1_i64,
1899                "instrument_ticker": ticker_json(),
1900            }),
1901        );
1902        let msg = parse_ticker_msg(&payload).unwrap();
1903
1904        let err = parse_funding_rate(&msg, UnixNanos::from(789))
1905            .expect_err("must reject negative ticker timestamp");
1906
1907        assert!(err.to_string().contains("negative Derive ticker timestamp"));
1908    }
1909
1910    #[rstest]
1911    fn test_parse_funding_rate_rejects_timestamp_overflow() {
1912        let payload = subscription_data_payload(
1913            "ticker.ETH-PERP.1000",
1914            &json!({
1915                "timestamp": i64::MAX,
1916                "instrument_ticker": ticker_json(),
1917            }),
1918        );
1919        let msg = parse_ticker_msg(&payload).unwrap();
1920
1921        let err = parse_funding_rate(&msg, UnixNanos::from(789))
1922            .expect_err("must reject overflowing ticker timestamp");
1923
1924        assert!(
1925            err.to_string()
1926                .contains("Derive timestamp overflows nanoseconds")
1927        );
1928    }
1929
1930    #[rstest]
1931    fn test_parse_funding_rate_history_record_maps_fields() {
1932        let record = DerivePublicFundingRate {
1933            funding_rate: Decimal::from_str("0.00015").unwrap(),
1934            timestamp: 1_700_000_000_000,
1935        };
1936        let instrument_id = InstrumentId::from("ETH-PERP.DERIVE");
1937
1938        let update = parse_funding_rate_history_record(
1939            &record,
1940            instrument_id,
1941            Some(60),
1942            UnixNanos::from(789),
1943        )
1944        .unwrap();
1945
1946        assert_eq!(update.instrument_id, instrument_id);
1947        assert_eq!(update.rate, Decimal::from_str("0.00015").unwrap());
1948        assert_eq!(update.interval, Some(60));
1949        assert!(update.next_funding_ns.is_none());
1950        assert_eq!(update.ts_event, UnixNanos::from(1_700_000_000_000_000_000));
1951        assert_eq!(update.ts_init, UnixNanos::from(789));
1952    }
1953
1954    #[rstest]
1955    fn test_parse_funding_rate_history_record_rejects_negative_timestamp() {
1956        let record = DerivePublicFundingRate {
1957            funding_rate: Decimal::from_str("0.0001").unwrap(),
1958            timestamp: -1,
1959        };
1960        let err = parse_funding_rate_history_record(
1961            &record,
1962            InstrumentId::from("ETH-PERP.DERIVE"),
1963            None,
1964            UnixNanos::from(789),
1965        )
1966        .expect_err("must reject negative timestamp");
1967
1968        assert!(err.to_string().contains("negative Derive ticker timestamp"));
1969    }
1970
1971    #[rstest]
1972    fn test_parse_candle_record_maps_fields() {
1973        // `timestamp` and `timestamp_bucket` differ so a swap from `timestamp_bucket`
1974        // to `timestamp` in the parser would shift ts_event and fail the assertion.
1975        let record = DerivePublicCandle {
1976            open_price: Decimal::from_str("3500.0").unwrap(),
1977            high_price: Decimal::from_str("3501.5").unwrap(),
1978            low_price: Decimal::from_str("3499.0").unwrap(),
1979            close_price: Decimal::from_str("3501.0").unwrap(),
1980            volume_usd: Decimal::from_str("12345.6").unwrap(),
1981            volume_contracts: Decimal::from_str("3.527").unwrap(),
1982            timestamp: 1_700_000_007,
1983            timestamp_bucket: 1_700_000_000,
1984        };
1985        let bar_type = BarType::from("ETH-PERP.DERIVE-1-MINUTE-LAST-EXTERNAL");
1986
1987        let bar = parse_candle_record(
1988            &record,
1989            bar_type,
1990            PRICE_PRECISION,
1991            SIZE_PRECISION,
1992            UnixNanos::from(789),
1993        )
1994        .unwrap();
1995
1996        assert_eq!(bar.bar_type, bar_type);
1997        assert_eq!(bar.open, Price::from_str("3500.00").unwrap());
1998        assert_eq!(bar.high, Price::from_str("3501.50").unwrap());
1999        assert_eq!(bar.low, Price::from_str("3499.00").unwrap());
2000        assert_eq!(bar.close, Price::from_str("3501.00").unwrap());
2001        assert_eq!(bar.volume, Quantity::from_str("3.527").unwrap());
2002        assert_eq!(bar.ts_event, UnixNanos::from(1_700_000_060_000_000_000));
2003        assert_eq!(bar.ts_init, UnixNanos::from(789));
2004    }
2005
2006    #[rstest]
2007    fn test_parse_candle_record_rejects_negative_timestamp() {
2008        let record = DerivePublicCandle {
2009            open_price: Decimal::from_str("1").unwrap(),
2010            high_price: Decimal::from_str("1").unwrap(),
2011            low_price: Decimal::from_str("1").unwrap(),
2012            close_price: Decimal::from_str("1").unwrap(),
2013            volume_usd: Decimal::ZERO,
2014            volume_contracts: Decimal::ZERO,
2015            timestamp: 1_700_000_000,
2016            timestamp_bucket: -1,
2017        };
2018        let err = parse_candle_record(
2019            &record,
2020            BarType::from("ETH-PERP.DERIVE-1-MINUTE-LAST-EXTERNAL"),
2021            PRICE_PRECISION,
2022            SIZE_PRECISION,
2023            UnixNanos::from(789),
2024        )
2025        .expect_err("must reject negative timestamp");
2026
2027        assert!(err.to_string().contains("negative Derive candle timestamp"));
2028    }
2029
2030    #[rstest]
2031    fn test_parse_candle_record_rejects_timestamp_overflow() {
2032        let record = DerivePublicCandle {
2033            open_price: Decimal::from_str("1").unwrap(),
2034            high_price: Decimal::from_str("1").unwrap(),
2035            low_price: Decimal::from_str("1").unwrap(),
2036            close_price: Decimal::from_str("1").unwrap(),
2037            volume_usd: Decimal::ZERO,
2038            volume_contracts: Decimal::ZERO,
2039            timestamp: 1_700_000_000,
2040            timestamp_bucket: i64::MAX,
2041        };
2042        let err = parse_candle_record(
2043            &record,
2044            BarType::from("ETH-PERP.DERIVE-1-MINUTE-LAST-EXTERNAL"),
2045            PRICE_PRECISION,
2046            SIZE_PRECISION,
2047            UnixNanos::from(789),
2048        )
2049        .expect_err("must reject overflowing timestamp");
2050
2051        assert!(
2052            err.to_string()
2053                .contains("Derive candle timestamp_bucket overflows nanoseconds"),
2054            "{err}",
2055        );
2056    }
2057
2058    #[rstest]
2059    fn test_parse_candle_record_rejects_close_timestamp_overflow() {
2060        let record = DerivePublicCandle {
2061            open_price: Decimal::from_str("1").unwrap(),
2062            high_price: Decimal::from_str("1").unwrap(),
2063            low_price: Decimal::from_str("1").unwrap(),
2064            close_price: Decimal::from_str("1").unwrap(),
2065            volume_usd: Decimal::ZERO,
2066            volume_contracts: Decimal::ZERO,
2067            timestamp: 1_700_000_000,
2068            timestamp_bucket: (u64::MAX / NANOSECONDS_IN_SECOND) as i64,
2069        };
2070        let err = parse_candle_record(
2071            &record,
2072            BarType::from("ETH-PERP.DERIVE-1-MINUTE-LAST-EXTERNAL"),
2073            PRICE_PRECISION,
2074            SIZE_PRECISION,
2075            UnixNanos::from(789),
2076        )
2077        .expect_err("must reject close timestamp overflow");
2078
2079        assert!(
2080            err.to_string()
2081                .contains("bar timestamp overflowed when adjusting to close time"),
2082            "{err}",
2083        );
2084    }
2085
2086    #[rstest]
2087    #[case(BarAggregation::Minute, 1, 60)]
2088    #[case(BarAggregation::Minute, 5, 300)]
2089    #[case(BarAggregation::Minute, 15, 900)]
2090    #[case(BarAggregation::Minute, 30, 1800)]
2091    #[case(BarAggregation::Hour, 1, 3600)]
2092    #[case(BarAggregation::Hour, 4, 14400)]
2093    #[case(BarAggregation::Hour, 8, 28800)]
2094    #[case(BarAggregation::Day, 1, 86400)]
2095    #[case(BarAggregation::Week, 1, 604800)]
2096    fn test_bar_spec_to_derive_period_maps_supported_intervals(
2097        #[case] aggregation: BarAggregation,
2098        #[case] step: u64,
2099        #[case] expected: u32,
2100    ) {
2101        assert_eq!(
2102            bar_spec_to_derive_period(aggregation, step).unwrap(),
2103            expected
2104        );
2105    }
2106
2107    #[rstest]
2108    #[case(BarAggregation::Minute, 2, "minute intervals")]
2109    #[case(BarAggregation::Hour, 2, "hour intervals")]
2110    #[case(BarAggregation::Day, 7, "1 DAY interval")]
2111    #[case(BarAggregation::Week, 2, "1 WEEK interval")]
2112    #[case(BarAggregation::Second, 1, "does not support")]
2113    fn test_bar_spec_to_derive_period_rejects_unsupported(
2114        #[case] aggregation: BarAggregation,
2115        #[case] step: u64,
2116        #[case] expected_msg: &str,
2117    ) {
2118        let err =
2119            bar_spec_to_derive_period(aggregation, step).expect_err("must reject unsupported spec");
2120        assert!(
2121            err.to_string().contains(expected_msg),
2122            "expected {expected_msg:?}, was {err}",
2123        );
2124    }
2125
2126    #[rstest]
2127    fn test_parse_option_greeks_rejects_negative_timestamp() {
2128        let mut option_data = option_ticker_json(1_700_000_000_000);
2129        option_data["instrument_name"] = json!("ETH-20260627-3500-C");
2130        let payload = subscription_data_payload(
2131            "ticker.ETH-20260627-3500-C.1000",
2132            &json!({
2133                "timestamp": -1_i64,
2134                "instrument_ticker": option_data,
2135            }),
2136        );
2137        let msg = parse_ticker_msg(&payload).unwrap();
2138
2139        let err = parse_option_greeks(&msg, UnixNanos::from(789))
2140            .expect_err("must reject negative ticker timestamp");
2141
2142        assert!(err.to_string().contains("negative Derive ticker timestamp"));
2143    }
2144
2145    #[rstest]
2146    fn test_parse_option_greeks_rejects_timestamp_overflow() {
2147        let mut option_data = option_ticker_json(1_700_000_000_000);
2148        option_data["instrument_name"] = json!("ETH-20260627-3500-C");
2149        let payload = subscription_data_payload(
2150            "ticker.ETH-20260627-3500-C.1000",
2151            &json!({
2152                "timestamp": i64::MAX,
2153                "instrument_ticker": option_data,
2154            }),
2155        );
2156        let msg = parse_ticker_msg(&payload).unwrap();
2157
2158        let err = parse_option_greeks(&msg, UnixNanos::from(789))
2159            .expect_err("must reject overflowing ticker timestamp");
2160
2161        assert!(
2162            err.to_string()
2163                .contains("Derive timestamp overflows nanoseconds")
2164        );
2165    }
2166}