1use 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
49pub 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
73pub 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
87pub 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
101pub 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
117pub 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
188pub 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
242pub 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
272pub 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
343pub 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
382pub 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
508pub 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
530pub 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
552pub 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
578pub 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
600pub 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
641pub 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
690pub 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 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 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 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}