use std::collections::{BTreeMap, HashMap};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use chrono::{DateTime, NaiveDate, Utc};
use deribit_http::model::instrument::{
Instrument as DeribitInstrument, OptionType as DeribitOptionType,
};
use deribit_http::model::other::OptionInstrument;
use deribit_http::{DeribitHttpClient, HttpConfig, HttpError};
use deribit_websocket::install_default_crypto_provider;
use deribit_websocket::prelude::{
DeribitWebSocketClient, NotificationHandler, SubscriptionChannel, WebSocketConfig,
};
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::{Decimal, Positive};
use optionstratlib::{ExpirationDate, OptionStyle};
use serde::Deserialize;
#[cfg(test)]
use tokio::sync::mpsc;
use tokio::task::JoinSet;
use tokio_util::sync::CancellationToken;
use super::{
AuthKind, ChainCapability, ChainPollCapability, GreeksCapability, MarketUpdateSink,
OptionStreamCapability, Provider, ProviderCapabilities, SendState, SubscriptionHandle,
SubscriptionRequest, UnderlyingRef,
};
use crate::chain::{
AliasCatalog, ChainFetch, ChainSnapshot, ChainSource, ContractSpecFingerprint, DepthLadder,
DepthLevel, ExerciseStyle, ExpirySource, GreeksOrigin, GreeksRow, Instrument, InstrumentKey,
MarketUpdate, PremiumNumeraire, ProviderId, QuoteUpdate, SettlementStyle, StreamHealth,
};
use crate::error::{NormalizeKind, ProviderError, TransportDetail, TransportKind};
const DERIBIT_ID: &str = "deribit";
const REFRESH_HINT_SECS: u32 = 2;
const DEFAULT_QUOTE_CURRENCY: &str = "USD";
const MULTIPLIER_MAX_F64: f64 = 4_294_967_295.0;
const OI_MAX_F64: f64 = 9_007_199_254_740_992.0;
const MAX_CONCURRENT_TICKERS: usize = 16;
const BACKOFF_BASE_MS: f64 = 250.0;
const BACKOFF_MAX_MS: f64 = 30_000.0;
const JITTER_MAGNITUDE: f64 = 0.2;
const BACKOFF_MAX_SHIFT: u32 = 20;
const STAGING_FLUSH_INTERVAL: Duration = Duration::from_millis(10);
#[derive(Clone)]
pub(crate) struct DeribitAdapter {
client: DeribitHttpClient,
id: ProviderId,
}
impl DeribitAdapter {
#[must_use]
pub(crate) fn new() -> Self {
Self {
client: DeribitHttpClient::with_config(HttpConfig::production()),
id: deribit_provider_id(),
}
}
async fn hydrate_legs(&self, selected: Vec<DeribitInstrument>) -> Hydration {
let mut pending = selected.into_iter();
let mut join_set: JoinSet<LegOutcome> = JoinSet::new();
for _ in 0..MAX_CONCURRENT_TICKERS {
let Some(instrument) = pending.next() else {
break;
};
self.spawn_ticker(&mut join_set, instrument);
}
let mut outcomes = Vec::new();
while let Some(joined) = join_set.join_next().await {
let outcome = match joined {
Ok(outcome) => outcome,
Err(_) => LegOutcome::Dropped,
};
outcomes.push(outcome);
if let Some(instrument) = pending.next() {
self.spawn_ticker(&mut join_set, instrument);
}
}
collect_outcomes(outcomes)
}
fn spawn_ticker(&self, join_set: &mut JoinSet<LegOutcome>, instrument: DeribitInstrument) {
let client = self.client.clone();
let _ = join_set.spawn(async move {
let ticker = match client.get_ticker(&instrument.instrument_name).await {
Ok(ticker) => ticker,
Err(_) => return LegOutcome::TransportFailed,
};
let option = OptionInstrument { instrument, ticker };
match normalize_leg(&option) {
Ok(leg) => LegOutcome::Hydrated(Box::new(leg)),
Err(_) => LegOutcome::Dropped,
}
});
}
}
#[async_trait]
impl Provider for DeribitAdapter {
fn id(&self) -> ProviderId {
self.id.clone()
}
fn capabilities(&self) -> ProviderCapabilities {
deribit_capabilities()
}
async fn discover(&self) -> Result<Vec<UnderlyingRef>, ProviderError> {
let currencies = self
.client
.get_currencies()
.await
.map_err(|err| transport_error(&err))?;
Ok(currencies
.into_iter()
.map(|currency| UnderlyingRef::new(currency.currency))
.collect())
}
async fn fetch_chain(
&self,
underlying: &str,
expiration: &ExpirationDate,
) -> Result<ChainFetch, ProviderError> {
let currency = underlying.to_ascii_uppercase();
let target = expiration
.get_date()
.map_err(|_| ProviderError::Normalize {
kind: NormalizeKind::UnparseableExpiry,
})?;
let target_day = target.date_naive();
let instruments = self
.client
.get_instruments(¤cy, Some("option"), Some(false))
.await
.map_err(|err| transport_error(&err))?;
let selected: Vec<DeribitInstrument> = instruments
.into_iter()
.filter(|instrument| {
instrument.is_option()
&& instrument
.expiration_timestamp
.and_then(DateTime::<Utc>::from_timestamp_millis)
.is_some_and(|expiry| expiry.date_naive() == target_day)
})
.collect();
let selected_count = selected.len();
let Hydration {
legs,
transport_failures,
} = self.hydrate_legs(selected).await;
if legs.is_empty() {
return Err(empty_expiry_outcome(
selected_count,
transport_failures,
¤cy,
target,
));
}
let expiration_utc = legs.first().map_or(target, |leg| leg.key.expiration_utc);
let spot =
legs.iter()
.find_map(|leg| leg.underlying_price)
.ok_or(ProviderError::Normalize {
kind: NormalizeKind::MissingField("underlying_price"),
})?;
Ok(assemble_chain(
¤cy,
spot,
expiration_utc,
&legs,
&self.id,
))
}
async fn subscribe(
&self,
req: SubscriptionRequest,
sink: MarketUpdateSink,
) -> Result<SubscriptionHandle, ProviderError> {
let transport = LiveTransport::new(self.clone(), WebSocketConfig::default());
let id = self.id.clone();
let SubscriptionRequest {
underlying,
expiration_utc,
instruments,
cancel,
} = req;
let loop_cancel = cancel.clone();
let handle = tokio::spawn(run_reconnect_loop(
transport,
id,
underlying,
expiration_utc,
instruments,
sink,
loop_cancel,
));
Ok(SubscriptionHandle::spawned(cancel, handle))
}
}
enum LegOutcome {
Hydrated(Box<NormalizedLeg>),
Dropped,
TransportFailed,
}
struct Hydration {
legs: Vec<NormalizedLeg>,
transport_failures: usize,
}
fn collect_outcomes(outcomes: impl IntoIterator<Item = LegOutcome>) -> Hydration {
let mut legs = Vec::new();
let mut transport_failures = 0usize;
for outcome in outcomes {
match outcome {
LegOutcome::Hydrated(leg) => legs.push(*leg),
LegOutcome::TransportFailed => transport_failures += 1,
LegOutcome::Dropped => {}
}
}
Hydration {
legs,
transport_failures,
}
}
fn empty_expiry_outcome(
selected_count: usize,
transport_failures: usize,
underlying: &str,
expiration: DateTime<Utc>,
) -> ProviderError {
if selected_count > 0 && transport_failures > 0 {
return transport(TransportKind::Closed);
}
ProviderError::NoChain {
underlying: underlying.to_owned(),
expiration: expiration.to_rfc3339(),
}
}
fn deribit_provider_id() -> ProviderId {
match ProviderId::new(DERIBIT_ID) {
Ok(id) => id,
Err(_) => unreachable!("`deribit` is a valid, reserved provider id literal"),
}
}
#[must_use]
pub(crate) fn deribit_capabilities() -> ProviderCapabilities {
ProviderCapabilities::builder()
.chain(ChainCapability::Assemble)
.depth(true)
.greeks(GreeksCapability::Provided)
.option_stream(OptionStreamCapability::ChainQuotes { verified: false })
.underlying_stream(true)
.chain_poll(ChainPollCapability::Poll {
interval_hint_secs: REFRESH_HINT_SECS,
})
.trades_tape(false)
.auth(AuthKind::None)
.build()
}
fn positive_or_drop(value: f64) -> Option<Positive> {
Positive::new(value).ok()
}
fn strike_positive(value: f64) -> Result<Positive, NormalizeKind> {
if !value.is_finite() {
return Err(NormalizeKind::NonFinite("strike"));
}
if value <= 0.0 {
return Err(NormalizeKind::OutOfRange("strike"));
}
let strike = Positive::new(value).map_err(|_| NormalizeKind::OutOfRange("strike"))?;
if strike == Positive::ZERO {
return Err(NormalizeKind::OutOfRange("strike"));
}
Ok(strike)
}
fn normalize_iv(mark_iv: f64) -> Result<Positive, NormalizeKind> {
if !mark_iv.is_finite() {
return Err(NormalizeKind::NonFinite("iv"));
}
let fraction = mark_iv / 100.0;
Positive::new(fraction).map_err(|_| NormalizeKind::OutOfRange("iv"))
}
fn greek_or_drop(value: Option<f64>) -> Option<Decimal> {
let raw = value?;
if !raw.is_finite() {
return None;
}
Decimal::try_from(raw).ok()
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
struct NormalizedQuote {
bid: Option<Positive>,
ask: Option<Positive>,
}
fn normalize_quote(bid: Option<f64>, ask: Option<f64>) -> Result<NormalizedQuote, NormalizeKind> {
let bid = bid.and_then(positive_or_drop);
let ask = ask.and_then(positive_or_drop);
if let (Some(bid_value), Some(ask_value)) = (bid, ask) {
if ask_value < bid_value {
return Err(NormalizeKind::OutOfRange("ask"));
}
}
Ok(NormalizedQuote { bid, ask })
}
#[derive(Debug, Clone)]
struct ParsedName {
underlying: String,
expiry_code: String,
strike: f64,
style: OptionStyle,
}
fn parse_instrument_name(name: &str) -> Result<ParsedName, NormalizeKind> {
let mut parts = name.split('-');
let underlying = parts
.next()
.filter(|segment| !segment.is_empty())
.ok_or(NormalizeKind::MissingField("instrument_name"))?;
let expiry_code = parts
.next()
.filter(|segment| !segment.is_empty())
.ok_or(NormalizeKind::UnparseableExpiry)?;
let strike_segment = parts
.next()
.filter(|segment| !segment.is_empty())
.ok_or(NormalizeKind::MissingField("strike"))?;
let style_segment = parts
.next()
.filter(|segment| !segment.is_empty())
.ok_or(NormalizeKind::UnknownStyle)?;
if parts.next().is_some() {
return Err(NormalizeKind::MissingField("instrument_name"));
}
let strike = strike_segment
.parse::<f64>()
.map_err(|_| NormalizeKind::OutOfRange("strike"))?;
let style = match style_segment.to_ascii_uppercase().as_str() {
"C" => OptionStyle::Call,
"P" => OptionStyle::Put,
_ => return Err(NormalizeKind::UnknownStyle),
};
Ok(ParsedName {
underlying: underlying.to_ascii_uppercase(),
expiry_code: expiry_code.to_ascii_uppercase(),
strike,
style,
})
}
fn expiry_code_to_utc(code: &str) -> Result<DateTime<Utc>, NormalizeKind> {
if !code.is_ascii() || code.len() < 6 {
return Err(NormalizeKind::UnparseableExpiry);
}
let year_at = code
.len()
.checked_sub(2)
.ok_or(NormalizeKind::UnparseableExpiry)?;
let (head, year_str) = code.split_at(year_at);
let month_at = head
.len()
.checked_sub(3)
.ok_or(NormalizeKind::UnparseableExpiry)?;
let (day_str, month_str) = head.split_at(month_at);
if day_str.is_empty() {
return Err(NormalizeKind::UnparseableExpiry);
}
let day = day_str
.parse::<u32>()
.map_err(|_| NormalizeKind::UnparseableExpiry)?;
let year_two = year_str
.parse::<i32>()
.map_err(|_| NormalizeKind::UnparseableExpiry)?;
let year = 2000 + year_two;
let month = month_from_code(month_str)?;
let date = NaiveDate::from_ymd_opt(year, month, day).ok_or(NormalizeKind::UnparseableExpiry)?;
let naive = date
.and_hms_opt(8, 0, 0)
.ok_or(NormalizeKind::UnparseableExpiry)?;
Ok(DateTime::<Utc>::from_naive_utc_and_offset(naive, Utc))
}
fn month_from_code(code: &str) -> Result<u32, NormalizeKind> {
let month = match code {
"JAN" => 1,
"FEB" => 2,
"MAR" => 3,
"APR" => 4,
"MAY" => 5,
"JUN" => 6,
"JUL" => 7,
"AUG" => 8,
"SEP" => 9,
"OCT" => 10,
"NOV" => 11,
"DEC" => 12,
_ => return Err(NormalizeKind::UnparseableExpiry),
};
Ok(month)
}
fn utc_from_millis(millis: i64) -> Result<DateTime<Utc>, NormalizeKind> {
DateTime::<Utc>::from_timestamp_millis(millis).ok_or(NormalizeKind::UnparseableExpiry)
}
fn instrument_key_from_name(name: &str) -> Result<InstrumentKey, NormalizeKind> {
let parsed = parse_instrument_name(name)?;
let expiration_utc = expiry_code_to_utc(&parsed.expiry_code)?;
let strike = strike_positive(parsed.strike)?;
Ok(InstrumentKey {
underlying: parsed.underlying,
expiration_utc,
strike,
style: parsed.style,
})
}
fn style_of(option_type: DeribitOptionType) -> OptionStyle {
match option_type {
DeribitOptionType::Call => OptionStyle::Call,
DeribitOptionType::Put => OptionStyle::Put,
}
}
#[derive(Debug, Clone)]
struct NormalizedLeg {
key: InstrumentKey,
native_symbol: String,
spec: ContractSpecFingerprint,
bid: Option<Positive>,
ask: Option<Positive>,
iv: Option<Positive>,
delta: Option<Decimal>,
gamma: Option<Decimal>,
volume: Option<Positive>,
open_interest: Option<u64>,
underlying_price: Option<Positive>,
style: OptionStyle,
inverse: bool,
}
fn normalize_leg(option: &OptionInstrument) -> Result<NormalizedLeg, NormalizeKind> {
let instrument = &option.instrument;
let ticker = &option.ticker;
let mut key = instrument_key_from_name(&instrument.instrument_name)?;
if let Some(strike) = instrument.strike {
key.strike = strike_positive(strike)?;
}
if let Some(option_type) = instrument.option_type.clone() {
key.style = style_of(option_type);
}
if let Some(millis) = instrument.expiration_timestamp {
key.expiration_utc = utc_from_millis(millis)?;
}
let quote = normalize_quote(ticker.best_bid_price, ticker.best_ask_price).unwrap_or_default();
let iv = ticker.mark_iv.and_then(|value| normalize_iv(value).ok());
let (delta, gamma) = match &ticker.greeks {
Some(greeks) => (greek_or_drop(greeks.delta), greek_or_drop(greeks.gamma)),
None => (None, None),
};
let volume = positive_or_drop(ticker.stats.volume);
let open_interest = ticker.open_interest.and_then(oi_to_u64);
let underlying_price = ticker
.underlying_price
.or(ticker.index_price)
.and_then(positive_or_drop);
let spec = deribit_fingerprint(instrument, &key.underlying);
let style = key.style;
Ok(NormalizedLeg {
key,
native_symbol: instrument.instrument_name.clone(),
spec,
bid: quote.bid,
ask: quote.ask,
iv,
delta,
gamma,
volume,
open_interest,
underlying_price,
style,
inverse: is_inverse_contract(instrument),
})
}
fn is_inverse_contract(instrument: &DeribitInstrument) -> bool {
match (&instrument.settlement_currency, &instrument.quote_currency) {
(Some(settlement), Some(quote)) => !settlement.eq_ignore_ascii_case(quote),
_ => false,
}
}
fn deribit_fingerprint(
instrument: &DeribitInstrument,
underlying: &str,
) -> ContractSpecFingerprint {
ContractSpecFingerprint {
contract_multiplier: contract_multiplier_of(instrument),
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: instrument
.quote_currency
.clone()
.unwrap_or_else(|| DEFAULT_QUOTE_CURRENCY.to_owned()),
venue_product_code: underlying.to_owned(),
}
}
fn contract_multiplier_of(instrument: &DeribitInstrument) -> u32 {
match instrument.contract_size {
Some(size) if size.is_finite() && (1.0..=MULTIPLIER_MAX_F64).contains(&size) => {
size.trunc() as u32
}
_ => 1,
}
}
fn oi_to_u64(value: f64) -> Option<u64> {
if value.is_finite() && (0.0..=OI_MAX_F64).contains(&value) {
Some(value.trunc() as u64)
} else {
None
}
}
#[derive(Debug, Default)]
struct StrikePair<'a> {
call: Option<&'a NormalizedLeg>,
put: Option<&'a NormalizedLeg>,
}
fn assemble_chain(
underlying: &str,
spot: Positive,
expiration_utc: DateTime<Utc>,
legs: &[NormalizedLeg],
provider: &ProviderId,
) -> ChainFetch {
let mut aliases = AliasCatalog::new();
for leg in legs {
aliases.insert(Instrument {
key: leg.key.clone(),
provider: provider.clone(),
native_symbol: leg.native_symbol.clone(),
stream_symbol: None,
spec: leg.spec.clone(),
});
}
let mut by_strike: BTreeMap<Positive, StrikePair<'_>> = BTreeMap::new();
for leg in legs {
let entry = by_strike.entry(leg.key.strike).or_default();
match leg.style {
OptionStyle::Call => entry.call = Some(leg),
OptionStyle::Put => entry.put = Some(leg),
}
}
let mut chain = OptionChain::new(underlying, spot, expiration_utc.to_rfc3339(), None, None);
for (strike, pair) in by_strike {
let iv = pair
.call
.and_then(|leg| leg.iv)
.or_else(|| pair.put.and_then(|leg| leg.iv))
.unwrap_or(Positive::ZERO);
chain.add_option(
strike,
pair.call.and_then(|leg| leg.bid),
pair.call.and_then(|leg| leg.ask),
pair.put.and_then(|leg| leg.bid),
pair.put.and_then(|leg| leg.ask),
iv,
pair.call.and_then(|leg| leg.delta),
pair.put.and_then(|leg| leg.delta),
pair.call
.and_then(|leg| leg.gamma)
.or_else(|| pair.put.and_then(|leg| leg.gamma)),
pair.call
.and_then(|leg| leg.volume)
.or_else(|| pair.put.and_then(|leg| leg.volume)),
pair.call
.and_then(|leg| leg.open_interest)
.or_else(|| pair.put.and_then(|leg| leg.open_interest)),
None,
);
}
let greeks_seed = venue_greeks_seed(legs, expiration_utc, provider);
let premium_numeraire = if legs.iter().any(|leg| leg.inverse) {
PremiumNumeraire::UnderlyingCoin
} else {
PremiumNumeraire::QuoteCurrency
};
ChainFetch::new(
chain,
ExpirySource::new(underlying, expiration_utc, provider.clone()),
aliases,
)
.with_greeks_seed(greeks_seed)
.with_premium_numeraire(premium_numeraire)
}
fn venue_greeks_seed(
legs: &[NormalizedLeg],
expiration_utc: DateTime<Utc>,
provider: &ProviderId,
) -> Vec<GreeksRow> {
legs.iter()
.filter(|leg| leg.iv.is_some() || leg.gamma.is_some())
.map(|leg| GreeksRow {
instrument: Instrument {
key: leg.key.clone(),
provider: provider.clone(),
native_symbol: leg.native_symbol.clone(),
stream_symbol: None,
spec: leg.spec.clone(),
},
iv: leg.iv,
delta: None,
gamma: leg.gamma,
theta: None,
vega: None,
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: expiration_utc,
})
.collect()
}
#[cfg(test)]
pub(crate) fn fixture_btc_chain_fetch_named(underlying: &str) -> ChainFetch {
use deribit_http::model::ticker::TickerData;
use deribit_websocket::prelude::Value;
const INSTRUMENTS_JSON: &str =
include_str!("../../tests/fixtures/deribit/instruments/instruments_btc.json");
const TICKER_60000_CALL_JSON: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_normal.json");
const TICKER_60000_PUT_JSON: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_put.json");
const TICKER_61000_CALL_JSON: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_61000_call.json");
fn ticker_json_for(instrument_name: &str) -> &'static str {
match instrument_name {
"BTC-27JUN25-60000-P" => TICKER_60000_PUT_JSON,
"BTC-27JUN25-61000-C" => TICKER_61000_CALL_JSON,
_ => TICKER_60000_CALL_JSON,
}
}
fn deserialize_ticker(json: &str) -> TickerData {
let value: Value = match json.parse() {
Ok(value) => value,
Err(e) => panic!("ticker fixture must parse: {e}"),
};
match TickerData::deserialize(&value) {
Ok(ticker) => ticker,
Err(e) => panic!("ticker fixture must deserialize: {e}"),
}
}
let instruments_value: Value = match INSTRUMENTS_JSON.parse() {
Ok(value) => value,
Err(e) => panic!("instruments fixture must parse: {e}"),
};
let instruments = match Vec::<DeribitInstrument>::deserialize(&instruments_value) {
Ok(list) => list,
Err(e) => panic!("instruments fixture must deserialize: {e}"),
};
let legs: Vec<NormalizedLeg> = instruments
.into_iter()
.filter(|instrument| instrument.is_option())
.filter_map(|instrument| {
let ticker = deserialize_ticker(ticker_json_for(&instrument.instrument_name));
normalize_leg(&OptionInstrument { instrument, ticker }).ok()
})
.collect();
let spot = match legs.iter().find_map(|leg| leg.underlying_price) {
Some(spot) => spot,
None => panic!("the recorded tickers carry an underlying price"),
};
let expiration_utc = match legs.first() {
Some(leg) => leg.key.expiration_utc,
None => panic!("the recorded fixture yields at least one normalized leg"),
};
assemble_chain(
underlying,
spot,
expiration_utc,
&legs,
&deribit_provider_id(),
)
}
#[cfg(test)]
pub(crate) fn fixture_btc_stream_updates(received: DateTime<Utc>) -> Vec<MarketUpdate> {
use deribit_websocket::prelude::Value;
const TICKER_60000_CALL_JSON: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_normal.json");
const TICKER_60000_PUT_JSON: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_put.json");
const TICKER_61000_CALL_JSON: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_61000_call.json");
fn ticker_json_for(instrument_name: &str) -> &'static str {
match instrument_name {
"BTC-27JUN25-60000-P" => TICKER_60000_PUT_JSON,
"BTC-27JUN25-61000-C" => TICKER_61000_CALL_JSON,
_ => TICKER_60000_CALL_JSON,
}
}
let fetch = fixture_btc_chain_fetch_named("BTC");
let mut out = Vec::new();
for instrument in fetch.aliases.instruments() {
let json = ticker_json_for(&instrument.native_symbol);
let value: Value = match json.parse() {
Ok(value) => value,
Err(e) => panic!("ticker fixture must parse: {e}"),
};
let payload = match TickerPayload::deserialize(&value) {
Ok(payload) => payload,
Err(e) => panic!("ticker fixture must deserialize as a TickerPayload: {e}"),
};
let (quote, greeks) = normalize_ticker(instrument, &payload, received);
out.push(MarketUpdate::Quote(quote));
out.push(MarketUpdate::Greeks(greeks));
}
out
}
#[cfg(test)]
pub(crate) fn fixture_btc_depth_ladder(key: InstrumentKey, received: DateTime<Utc>) -> DepthLadder {
use deribit_websocket::prelude::Value;
const BOOK_GROUPED_JSON: &str =
include_str!("../../tests/fixtures/deribit/book/book_grouped_snapshot.json");
let value: Value = match BOOK_GROUPED_JSON.parse() {
Ok(value) => value,
Err(e) => panic!("grouped book fixture must parse: {e}"),
};
let payload = match BookPayload::deserialize(&value) {
Ok(payload) => payload,
Err(e) => panic!("grouped book fixture must deserialize as BookPayload: {e}"),
};
let instrument = Instrument {
key,
provider: deribit_provider_id(),
native_symbol: "BTC-27JUN25-60000-C".to_owned(),
stream_symbol: None,
spec: ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: "BTC".to_owned(),
},
};
normalize_book(&instrument, &payload, received)
}
#[cfg(feature = "bench")]
pub(crate) fn bench_stream_burst(
legs: &[Instrument],
round: u64,
received: DateTime<Utc>,
) -> Vec<MarketUpdate> {
let step = u32::try_from(round % 16).unwrap_or(0);
let base = 1.0 + f64::from(step) * 0.05;
let base_ms = 1_751_011_200_000_i64; let event_ms = base_ms
.checked_add(i64::try_from(round).unwrap_or(0))
.unwrap_or(i64::MAX);
let mut out = Vec::new();
for leg in legs {
let ticker = TickerPayload {
best_bid_price: Some(base),
best_ask_price: Some(base + 0.2),
best_bid_amount: Some(10.0),
best_ask_amount: Some(12.0),
last_price: Some(base + 0.1),
mark_iv: Some(49.22),
timestamp: Some(event_ms),
greeks: Some(GreeksPayload {
delta: Some(0.5),
gamma: Some(0.01),
}),
};
let (quote, greeks) = normalize_ticker(leg, &ticker, received);
out.push(MarketUpdate::Quote(quote));
out.push(MarketUpdate::Greeks(greeks));
let book = BookPayload {
change_id: Some(round),
timestamp: Some(event_ms),
bids: vec![
BookLevel::Priced([base, 10.0]),
BookLevel::Priced([base - 0.1, 20.0]),
],
asks: vec![
BookLevel::Priced([base + 0.2, 12.0]),
BookLevel::Priced([base + 0.3, 22.0]),
],
};
out.push(MarketUpdate::Depth(normalize_book(leg, &book, received)));
}
out
}
#[cfg(feature = "bench")]
#[derive(Debug)]
pub(crate) struct BenchProducerStaging {
sink: MarketUpdateSink,
}
#[cfg(feature = "bench")]
impl BenchProducerStaging {
pub(crate) fn new(sink: MarketUpdateSink) -> Self {
Self { sink }
}
pub(crate) fn publish_burst(
&mut self,
legs: &[Instrument],
round: u64,
received: DateTime<Utc>,
) -> bool {
for update in bench_stream_burst(legs, round, received) {
if self.sink.publish_coalesced(update) == SendState::Closed {
return false;
}
}
true
}
#[cfg(test)]
pub(crate) fn flush(&mut self) -> bool {
self.sink.flush() == SendState::Open
}
#[cfg(test)]
pub(crate) fn has_pending(&self) -> bool {
self.sink.has_pending()
}
}
#[cfg(feature = "fuzz")]
fn fuzz_instrument() -> Instrument {
Instrument {
key: InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: DateTime::<Utc>::from_timestamp_millis(1_751_011_200_000)
.unwrap_or(DateTime::<Utc>::MIN_UTC),
strike: Positive::new(60_000.0).unwrap_or(Positive::ZERO),
style: OptionStyle::Call,
},
provider: deribit_provider_id(),
native_symbol: "BTC-27JUN25-60000-C".to_owned(),
stream_symbol: None,
spec: ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: "BTC".to_owned(),
},
}
}
#[cfg(feature = "fuzz")]
fn fuzz_received() -> DateTime<Utc> {
DateTime::<Utc>::from_timestamp(1_700_000_000, 0).unwrap_or(DateTime::<Utc>::MIN_UTC)
}
#[cfg(feature = "fuzz")]
fn assert_quote_valid(quote: &QuoteUpdate) {
for p in [
quote.bid,
quote.ask,
quote.last,
quote.bid_size,
quote.ask_size,
]
.into_iter()
.flatten()
{
assert!(
p >= Positive::ZERO,
"a fuzzed quote produced a non-domain Positive"
);
}
if let (Some(bid), Some(ask)) = (quote.bid, quote.ask) {
assert!(
ask >= bid,
"a fuzzed quote must never yield a crossed bid/ask"
);
}
}
#[cfg(feature = "fuzz")]
fn assert_greeks_valid(greeks: &GreeksRow) {
if let Some(iv) = greeks.iv {
assert!(
iv >= Positive::ZERO,
"a fuzzed IV produced a non-domain Positive"
);
}
assert!(
greeks.theta.is_none() && greeks.vega.is_none() && greeks.rho.is_none(),
"the Deribit ticker seam must never emit theta/vega/rho"
);
}
#[cfg(feature = "fuzz")]
pub(crate) fn fuzz_normalize_ticker(bytes: &[u8]) {
let Ok(payload) = serde_json::from_slice::<TickerPayload>(bytes) else {
return;
};
let instrument = fuzz_instrument();
let (quote, greeks) = normalize_ticker(&instrument, &payload, fuzz_received());
assert_quote_valid("e);
assert_greeks_valid(&greeks);
}
#[cfg(feature = "fuzz")]
pub(crate) fn fuzz_normalize_book(bytes: &[u8]) {
let Ok(payload) = serde_json::from_slice::<BookPayload>(bytes) else {
return;
};
let ladder = normalize_book(&fuzz_instrument(), &payload, fuzz_received());
for level in ladder.bids.iter().chain(ladder.asks.iter()) {
assert!(
level.price >= Positive::ZERO && level.size >= Positive::ZERO,
"a fuzzed depth level produced a non-domain Positive"
);
}
assert!(
ladder.bids.len() <= payload.bids.len() && ladder.asks.len() <= payload.asks.len(),
"normalize_book must never fabricate a depth level"
);
}
#[cfg(feature = "fuzz")]
pub(crate) fn fuzz_instrument_key_from_name(bytes: &[u8]) {
let name = String::from_utf8_lossy(bytes);
if let Ok(key) = instrument_key_from_name(&name) {
assert!(
key.strike > Positive::ZERO,
"a parsed instrument key must carry a positive strike"
);
}
}
fn transport_error(err: &HttpError) -> ProviderError {
match err {
HttpError::AuthenticationFailed(_) => ProviderError::Auth,
HttpError::RateLimitExceeded => ProviderError::RateLimited(None),
HttpError::NetworkError(_) => transport(TransportKind::Closed),
HttpError::RequestFailed(_) | HttpError::ConfigError(_) => transport(TransportKind::Http),
HttpError::InvalidResponse(_) | HttpError::ParseError(_) => {
transport(TransportKind::Decode)
}
}
}
fn transport(kind: TransportKind) -> ProviderError {
ProviderError::Transport(Box::new(TransportDetail::new(kind, None)))
}
#[must_use]
fn backoff_delay(attempt: u32, jitter: f64) -> Duration {
let exponent = attempt.min(BACKOFF_MAX_SHIFT);
let uncapped = BACKOFF_BASE_MS * 2.0_f64.powi(exponent as i32);
let capped = uncapped.min(BACKOFF_MAX_MS);
let jitter = jitter.clamp(-JITTER_MAGNITUDE, JITTER_MAGNITUDE);
let millis = capped * (1.0 + jitter);
Duration::from_secs_f64(millis / 1000.0)
}
fn sample_jitter() -> f64 {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |elapsed| elapsed.subsec_nanos());
let unit = f64::from(nanos) / 1_000_000_000.0; (unit * 2.0 - 1.0) * JITTER_MAGNITUDE }
fn now_utc() -> DateTime<Utc> {
let since = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or(Duration::ZERO);
let secs = i64::try_from(since.as_secs()).unwrap_or(i64::MAX);
DateTime::<Utc>::from_timestamp(secs, since.subsec_nanos()).unwrap_or(DateTime::<Utc>::MIN_UTC)
}
fn millis_to_event_time(millis: i64) -> Option<DateTime<Utc>> {
DateTime::<Utc>::from_timestamp_millis(millis)
}
#[derive(Debug, Clone, Deserialize)]
struct TickerPayload {
#[serde(default)]
best_bid_price: Option<f64>,
#[serde(default)]
best_ask_price: Option<f64>,
#[serde(default)]
best_bid_amount: Option<f64>,
#[serde(default)]
best_ask_amount: Option<f64>,
#[serde(default)]
last_price: Option<f64>,
#[serde(default)]
mark_iv: Option<f64>,
#[serde(default)]
timestamp: Option<i64>,
#[serde(default)]
greeks: Option<GreeksPayload>,
}
#[derive(Debug, Clone, Deserialize)]
struct GreeksPayload {
#[serde(default)]
delta: Option<f64>,
#[serde(default)]
gamma: Option<f64>,
}
#[derive(Debug, Clone, Deserialize)]
struct BookPayload {
#[serde(default)]
change_id: Option<u64>,
#[serde(default)]
timestamp: Option<i64>,
#[serde(default)]
bids: Vec<BookLevel>,
#[serde(default)]
asks: Vec<BookLevel>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(untagged)]
enum BookLevel {
Priced([f64; 2]),
Actioned(String, f64, f64),
}
impl BookLevel {
fn price_size(&self) -> (f64, f64) {
match self {
BookLevel::Priced([price, amount]) => (*price, *amount),
BookLevel::Actioned(_action, price, amount) => (*price, *amount),
}
}
}
fn normalize_ticker(
instrument: &Instrument,
payload: &TickerPayload,
received: DateTime<Utc>,
) -> (QuoteUpdate, GreeksRow) {
let quote = normalize_quote(payload.best_bid_price, payload.best_ask_price).unwrap_or_default();
let event_time = payload.timestamp.and_then(millis_to_event_time);
let quote_update = QuoteUpdate {
instrument: instrument.clone(),
bid: quote.bid,
ask: quote.ask,
last: payload.last_price.and_then(positive_or_drop),
bid_size: payload.best_bid_amount.and_then(positive_or_drop),
ask_size: payload.best_ask_amount.and_then(positive_or_drop),
event_time,
received_time: received,
};
let iv = payload.mark_iv.and_then(|value| normalize_iv(value).ok());
let (delta, gamma) = match &payload.greeks {
Some(greeks) => (greek_or_drop(greeks.delta), greek_or_drop(greeks.gamma)),
None => (None, None),
};
let greeks_row = GreeksRow {
instrument: instrument.clone(),
iv,
delta,
gamma,
theta: None,
vega: None,
rho: None,
origin: GreeksOrigin::Provider,
event_time,
received_time: received,
};
(quote_update, greeks_row)
}
fn normalize_book(
instrument: &Instrument,
payload: &BookPayload,
received: DateTime<Utc>,
) -> DepthLadder {
DepthLadder {
instrument: instrument.clone(),
bids: payload.bids.iter().filter_map(depth_level).collect(),
asks: payload.asks.iter().filter_map(depth_level).collect(),
event_time: payload.timestamp.and_then(millis_to_event_time),
received_time: received,
change_id: payload.change_id,
}
}
fn depth_level(level: &BookLevel) -> Option<DepthLevel> {
let (price, size) = level.price_size();
Some(DepthLevel {
price: positive_or_drop(price)?,
size: positive_or_drop(size)?,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TransportGone;
#[async_trait]
trait DeribitTransport: Send {
async fn connect_and_subscribe(&mut self, channels: Vec<String>) -> Result<(), TransportGone>;
async fn receive(&mut self) -> Result<String, TransportGone>;
async fn refetch(
&mut self,
underlying: &str,
expiration: &ExpirationDate,
) -> Option<ChainFetch>;
}
struct LiveTransport {
adapter: DeribitAdapter,
ws_config: WebSocketConfig,
session: Option<DeribitWebSocketClient>,
}
impl LiveTransport {
fn new(adapter: DeribitAdapter, ws_config: WebSocketConfig) -> Self {
Self {
adapter,
ws_config,
session: None,
}
}
}
#[async_trait]
impl DeribitTransport for LiveTransport {
async fn connect_and_subscribe(&mut self, channels: Vec<String>) -> Result<(), TransportGone> {
let _ = install_default_crypto_provider();
let client = DeribitWebSocketClient::new(&self.ws_config).map_err(|_| TransportGone)?;
client.connect().await.map_err(|_| TransportGone)?;
client
.subscribe(channels)
.await
.map_err(|_| TransportGone)?;
self.session = Some(client);
Ok(())
}
async fn receive(&mut self) -> Result<String, TransportGone> {
match self.session.as_ref() {
Some(client) => client.receive_message().await.map_err(|_| TransportGone),
None => Err(TransportGone),
}
}
async fn refetch(
&mut self,
underlying: &str,
expiration: &ExpirationDate,
) -> Option<ChainFetch> {
self.adapter.fetch_chain(underlying, expiration).await.ok()
}
}
enum StreamExit {
Reconnect,
Shutdown,
}
async fn run_reconnect_loop<T: DeribitTransport>(
mut transport: T,
id: ProviderId,
underlying: String,
expiration_utc: DateTime<Utc>,
mut instruments: Vec<Instrument>,
mut sink: MarketUpdateSink,
cancel: CancellationToken,
) {
let mut attempt: u32 = 0;
loop {
if cancel.is_cancelled() || sink.is_closed() {
return;
}
let exit = tokio::select! {
biased;
() = cancel.cancelled() => return,
exit = connect_stream_once(&mut transport, &id, &instruments, &mut sink, &cancel, &mut attempt) => exit,
};
if matches!(exit, StreamExit::Shutdown) || cancel.is_cancelled() {
return;
}
sink.epoch();
attempt = attempt.checked_add(1).unwrap_or(attempt);
let health = MarketUpdate::Health(id.clone(), StreamHealth::Reconnecting { attempt });
let health_sent = tokio::select! {
biased;
() = cancel.cancelled() => return,
state = sink.send_control(health) => state,
};
if health_sent == SendState::Closed {
return; }
let delay = backoff_delay(attempt, sample_jitter());
tokio::select! {
biased;
() = cancel.cancelled() => return,
() = tokio::time::sleep(delay) => {}
}
if let Some(fresh) = refetch(
&mut transport,
&id,
&underlying,
expiration_utc,
&mut sink,
&cancel,
)
.await
&& !fresh.is_empty()
{
instruments = fresh;
}
}
}
async fn connect_stream_once<T: DeribitTransport>(
transport: &mut T,
id: &ProviderId,
instruments: &[Instrument],
sink: &mut MarketUpdateSink,
cancel: &CancellationToken,
attempt: &mut u32,
) -> StreamExit {
let channels = subscription_channels(instruments);
let subscribed = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
result = transport.connect_and_subscribe(channels) => result,
};
if subscribed.is_err() {
return StreamExit::Reconnect;
}
*attempt = 0;
let live = MarketUpdate::Health(id.clone(), StreamHealth::Live);
if sink.send_control(live).await == SendState::Closed {
return StreamExit::Shutdown;
}
let lookup = instrument_lookup(instruments);
let mut flush_tick = tokio::time::interval(STAGING_FLUSH_INTERVAL);
flush_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
let message = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
_ = flush_tick.tick(), if sink.has_pending() => {
if sink.flush() == SendState::Closed {
return StreamExit::Shutdown; }
continue;
}
message = transport.receive() => message,
};
let text = match message {
Ok(text) => text,
Err(_) => return StreamExit::Reconnect, };
if route_message(&text, &lookup, sink) == SendState::Closed {
return StreamExit::Shutdown; }
}
}
async fn refetch<T: DeribitTransport>(
transport: &mut T,
id: &ProviderId,
underlying: &str,
expiration_utc: DateTime<Utc>,
sink: &mut MarketUpdateSink,
cancel: &CancellationToken,
) -> Option<Vec<Instrument>> {
let expiration = ExpirationDate::DateTime(expiration_utc);
let fetched = tokio::select! {
biased;
() = cancel.cancelled() => return None,
result = transport.refetch(underlying, &expiration) => result,
};
let fetch = fetched?;
let snapshot = MarketUpdate::Chain(chain_snapshot(&fetch, now_utc()));
let snapshot_sent = tokio::select! {
biased;
() = cancel.cancelled() => return None,
state = sink.send_control(snapshot) => state,
};
if snapshot_sent == SendState::Closed {
return None;
}
let instruments: Vec<Instrument> = fetch
.aliases
.instruments()
.filter(|instrument| instrument.provider == *id)
.cloned()
.collect();
Some(instruments)
}
fn chain_snapshot(fetch: &ChainFetch, last_poll: DateTime<Utc>) -> ChainSnapshot {
ChainSnapshot {
chain_key: (
fetch.expiry_source.provider.clone(),
fetch.expiry_source.underlying.clone(),
fetch.expiry_source.expiration_utc,
),
chain: fetch.chain.clone(),
aliases: fetch.aliases.clone(),
source: ChainSource::Merged,
health: StreamHealth::Live,
last_full_poll: Some(last_poll),
}
}
const BOOK_GROUP: &str = "none";
const BOOK_DEPTH: &str = "20";
const BOOK_INTERVAL: &str = "100ms";
fn subscription_channels(instruments: &[Instrument]) -> Vec<String> {
let hint = instruments
.len()
.checked_mul(2)
.unwrap_or(instruments.len());
let mut channels = Vec::with_capacity(hint);
for instrument in instruments {
let native = instrument.native_symbol.clone();
channels.push(SubscriptionChannel::Ticker(native.clone()).channel_name());
channels.push(
SubscriptionChannel::GroupedOrderBook {
instrument: native,
group: BOOK_GROUP.to_owned(),
depth: BOOK_DEPTH.to_owned(),
interval: BOOK_INTERVAL.to_owned(),
}
.channel_name(),
);
}
channels
}
fn instrument_lookup(instruments: &[Instrument]) -> HashMap<String, Instrument> {
instruments
.iter()
.map(|instrument| (instrument.native_symbol.clone(), instrument.clone()))
.collect()
}
fn route_message(
text: &str,
lookup: &HashMap<String, Instrument>,
sink: &mut MarketUpdateSink,
) -> SendState {
let handler = NotificationHandler::new();
let Ok(notification) = handler.parse_notification(text) else {
return SendState::Open;
};
if !handler.is_subscription_notification(¬ification) {
return SendState::Open;
}
let (Some(channel), Some(data)) = (
handler.extract_channel(¬ification),
handler.extract_data(¬ification),
) else {
return SendState::Open;
};
let received = now_utc();
let instrument_segment = channel.split('.').nth(1);
if channel.starts_with("ticker.") {
let Some(symbol) = instrument_segment else {
return SendState::Open;
};
let Some(instrument) = lookup.get(symbol) else {
return SendState::Open; };
let Ok(payload) = TickerPayload::deserialize(&data) else {
return SendState::Open;
};
let (quote, greeks) = normalize_ticker(instrument, &payload, received);
if sink.publish_coalesced(MarketUpdate::Quote(quote)) == SendState::Closed {
return SendState::Closed;
}
return sink.publish_coalesced(MarketUpdate::Greeks(greeks));
}
if channel.starts_with("book.") {
let Some(symbol) = instrument_segment else {
return SendState::Open;
};
let Some(instrument) = lookup.get(symbol) else {
return SendState::Open; };
let Ok(payload) = BookPayload::deserialize(&data) else {
return SendState::Open;
};
let ladder = normalize_book(instrument, &payload, received);
return sink.publish_coalesced(MarketUpdate::Depth(ladder));
}
SendState::Open }
#[cfg(test)]
mod tests {
use std::sync::{Arc, Mutex as StdMutex};
use deribit_http::model::instrument::InstrumentKind;
use deribit_http::model::other::Greeks;
use deribit_http::model::ticker::{TickerData, TickerStats};
use deribit_websocket::prelude::Value;
use proptest::prelude::*;
use super::*;
#[track_caller]
fn pos(value: f64) -> Positive {
match Positive::new(value) {
Ok(p) => p,
Err(e) => panic!("invalid test positive `{value}`: {e}"),
}
}
#[track_caller]
fn utc_millis(millis: i64) -> DateTime<Utc> {
match DateTime::<Utc>::from_timestamp_millis(millis) {
Some(t) => t,
None => panic!("invalid test millis: {millis}"),
}
}
fn assert_send_sync<T: Send + Sync>() {}
fn sink_from(tx: &mpsc::Sender<MarketUpdate>) -> MarketUpdateSink {
MarketUpdateSink::new(tx.clone(), tx.clone())
}
fn deribit_instrument(
name: &str,
strike: Option<f64>,
option_type: Option<DeribitOptionType>,
expiration_ms: Option<i64>,
) -> DeribitInstrument {
DeribitInstrument {
instrument_name: name.to_owned(),
kind: Some(InstrumentKind::Option),
currency: Some("BTC".to_owned()),
is_active: Some(true),
expiration_timestamp: expiration_ms,
strike,
option_type,
contract_size: Some(1.0),
quote_currency: Some("USD".to_owned()),
..DeribitInstrument::default()
}
}
fn ticker(
name: &str,
best_bid_price: Option<f64>,
best_ask_price: Option<f64>,
mark_iv: Option<f64>,
greeks: Option<Greeks>,
) -> TickerData {
TickerData {
instrument_name: name.to_owned(),
last_price: None,
mark_price: 0.05,
best_bid_price,
best_ask_price,
best_bid_amount: 0.0,
best_ask_amount: 0.0,
volume: Some(12.0),
volume_usd: None,
open_interest: Some(34.0),
high: None,
low: None,
price_change: None,
price_change_percentage: None,
bid_iv: None,
ask_iv: None,
mark_iv,
timestamp: 0,
state: "open".to_owned(),
settlement_price: None,
stats: TickerStats {
volume: 12.0,
volume_usd: None,
price_change: None,
high: None,
low: None,
},
greeks,
index_price: Some(61_000.0),
min_price: None,
max_price: None,
interest_rate: None,
underlying_price: Some(60_500.0),
underlying_index: None,
estimated_delivery_price: None,
}
}
fn greeks(delta: Option<f64>, gamma: Option<f64>) -> Greeks {
Greeks {
delta,
gamma,
vega: None,
theta: None,
rho: None,
}
}
fn sample_option() -> OptionInstrument {
OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25-60000-C",
Some(60_000.0),
Some(DeribitOptionType::Call),
Some(1_751_011_200_000), ),
ticker: ticker(
"BTC-27JUN25-60000-C",
Some(0.05),
Some(0.06),
Some(49.22),
Some(greeks(Some(0.55), Some(0.0001))),
),
}
}
#[test]
fn test_deribit_capabilities_match_section_8_row() {
let caps = deribit_capabilities();
assert_eq!(caps.chain, ChainCapability::Assemble);
assert!(caps.depth);
assert_eq!(caps.greeks, GreeksCapability::Provided);
assert_eq!(
caps.option_stream,
OptionStreamCapability::ChainQuotes { verified: false }
);
assert!(caps.underlying_stream);
assert_eq!(
caps.chain_poll,
ChainPollCapability::Poll {
interval_hint_secs: REFRESH_HINT_SECS
}
);
assert!(!caps.trades_tape);
assert_eq!(caps.auth, AuthKind::None);
}
#[test]
fn test_deribit_id_is_valid_and_reserved() {
let id = deribit_provider_id();
assert_eq!(id.as_str(), "deribit");
assert!(id.is_reserved());
assert!(ProviderId::new(DERIBIT_ID).is_ok());
}
#[test]
fn test_deribit_normalize_iv_divides_by_100() {
match normalize_iv(49.22) {
Ok(iv) => assert_eq!(iv, pos(0.4922)),
Err(e) => panic!("expected 49.22% -> 0.4922, got {e}"),
}
}
#[test]
fn test_deribit_normalize_iv_zero_is_valid() {
match normalize_iv(0.0) {
Ok(iv) => assert_eq!(iv, Positive::ZERO),
Err(e) => panic!("expected zero IV to be valid, got {e}"),
}
}
#[test]
fn test_deribit_normalize_iv_rejects_non_finite() {
assert_eq!(normalize_iv(f64::NAN), Err(NormalizeKind::NonFinite("iv")));
assert_eq!(
normalize_iv(f64::INFINITY),
Err(NormalizeKind::NonFinite("iv"))
);
}
#[test]
fn test_deribit_normalize_iv_rejects_negative() {
assert_eq!(normalize_iv(-1.0), Err(NormalizeKind::OutOfRange("iv")));
}
#[test]
fn test_deribit_normalize_quote_keeps_zero_bid() {
match normalize_quote(Some(0.0), Some(1.0)) {
Ok(quote) => {
assert_eq!(quote.bid, Some(Positive::ZERO));
assert_eq!(quote.ask, Some(pos(1.0)));
}
Err(e) => panic!("a zero bid is valid, got {e}"),
}
}
#[test]
fn test_deribit_normalize_quote_rejects_zero_ask_on_nonzero_bid() {
assert_eq!(
normalize_quote(Some(5.0), Some(0.0)),
Err(NormalizeKind::OutOfRange("ask"))
);
}
#[test]
fn test_deribit_normalize_quote_rejects_crossed() {
assert_eq!(
normalize_quote(Some(5.0), Some(3.0)),
Err(NormalizeKind::OutOfRange("ask"))
);
}
#[test]
fn test_deribit_normalize_quote_drops_negative_price_field() {
match normalize_quote(Some(-1.0), Some(2.0)) {
Ok(quote) => {
assert_eq!(quote.bid, None);
assert_eq!(quote.ask, Some(pos(2.0)));
}
Err(e) => panic!("a negative bid drops only that field, got {e}"),
}
}
#[test]
fn test_deribit_normalize_quote_drops_non_finite_price_field() {
match normalize_quote(Some(f64::NAN), Some(2.0)) {
Ok(quote) => {
assert_eq!(quote.bid, None);
assert_eq!(quote.ask, Some(pos(2.0)));
}
Err(e) => panic!("a NaN bid drops only that field, got {e}"),
}
}
#[test]
fn test_deribit_normalize_quote_both_zero_is_valid() {
match normalize_quote(Some(0.0), Some(0.0)) {
Ok(quote) => {
assert_eq!(quote.bid, Some(Positive::ZERO));
assert_eq!(quote.ask, Some(Positive::ZERO));
}
Err(e) => panic!("a zero bid AND zero ask is valid, got {e}"),
}
}
#[test]
fn test_deribit_greek_decimal_keeps_negative() {
match greek_or_drop(Some(-0.45)) {
Some(delta) => assert_eq!(delta, Decimal::new(-45, 2)),
None => panic!("a negative Greek must be kept"),
}
}
#[test]
fn test_deribit_greek_decimal_drops_non_finite() {
assert_eq!(greek_or_drop(Some(f64::NAN)), None);
assert_eq!(greek_or_drop(Some(f64::INFINITY)), None);
assert_eq!(greek_or_drop(None), None);
}
#[test]
fn test_deribit_parse_instrument_name_maps_fields() {
match parse_instrument_name("BTC-27JUN25-60000-C") {
Ok(parsed) => {
assert_eq!(parsed.underlying, "BTC");
assert_eq!(parsed.expiry_code, "27JUN25");
assert_eq!(parsed.strike, 60_000.0);
assert_eq!(parsed.style, OptionStyle::Call);
}
Err(e) => panic!("expected a clean parse, got {e}"),
}
}
#[test]
fn test_deribit_instrument_name_maps_to_instrument_key() {
match instrument_key_from_name("BTC-27JUN25-60000-P") {
Ok(key) => {
assert_eq!(key.underlying, "BTC");
assert_eq!(key.strike, pos(60_000.0));
assert_eq!(key.style, OptionStyle::Put);
assert_eq!(key.expiration_utc, utc_millis(1_751_011_200_000));
}
Err(e) => panic!("expected a clean key mapping, got {e}"),
}
}
#[test]
fn test_deribit_expiry_code_resolves_to_0800_utc() {
match expiry_code_to_utc("27JUN25") {
Ok(instant) => {
assert_eq!(instant.to_rfc3339(), "2025-06-27T08:00:00+00:00");
}
Err(e) => panic!("expected 08:00 UTC settlement, got {e}"),
}
}
#[test]
fn test_deribit_expiry_code_single_digit_day() {
match expiry_code_to_utc("3JAN25") {
Ok(instant) => assert_eq!(instant.to_rfc3339(), "2025-01-03T08:00:00+00:00"),
Err(e) => panic!("expected a single-digit day to parse, got {e}"),
}
}
#[test]
fn test_deribit_normalize_uses_direct_utc_expiry() {
match normalize_leg(&sample_option()) {
Ok(leg) => assert_eq!(leg.key.expiration_utc, utc_millis(1_751_011_200_000)),
Err(e) => panic!("expected a normalized leg, got {e}"),
}
}
#[test]
fn test_deribit_normalize_rejects_unparseable_expiry() {
let option = OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25-60000-C",
Some(60_000.0),
Some(DeribitOptionType::Call),
Some(i64::MAX),
),
ticker: ticker("BTC-27JUN25-60000-C", Some(0.05), Some(0.06), None, None),
};
match normalize_leg(&option) {
Err(kind) => assert_eq!(kind, NormalizeKind::UnparseableExpiry),
Ok(_) => panic!("an out-of-range expiry must reject the row"),
}
}
#[test]
fn test_deribit_parse_instrument_name_rejects_unparseable_expiry() {
assert_eq!(
instrument_key_from_name("BTC-99XYZ25-60000-C"),
Err(NormalizeKind::UnparseableExpiry)
);
}
#[test]
fn test_deribit_normalize_rejects_missing_strike() {
let option = OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25--C",
None,
Some(DeribitOptionType::Call),
Some(1_751_011_200_000),
),
ticker: ticker("BTC-27JUN25--C", Some(0.05), Some(0.06), None, None),
};
match normalize_leg(&option) {
Err(kind) => assert_eq!(kind, NormalizeKind::MissingField("strike")),
Ok(_) => panic!("a row with no strike must be rejected"),
}
}
#[test]
fn test_deribit_normalize_rejects_zero_strike() {
let option = OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25-0-C",
Some(0.0),
Some(DeribitOptionType::Call),
Some(1_751_011_200_000),
),
ticker: ticker("BTC-27JUN25-0-C", Some(0.05), Some(0.06), None, None),
};
match normalize_leg(&option) {
Err(kind) => assert_eq!(kind, NormalizeKind::OutOfRange("strike")),
Ok(_) => panic!("a zero strike must reject the row"),
}
}
#[test]
fn test_deribit_normalize_rejects_unknown_style() {
assert_eq!(
instrument_key_from_name("BTC-27JUN25-60000-X"),
Err(NormalizeKind::UnknownStyle)
);
}
#[test]
fn test_deribit_normalize_leg_iv_is_decimal_fraction() {
match normalize_leg(&sample_option()) {
Ok(leg) => {
assert_eq!(leg.iv, Some(pos(0.4922)));
assert_eq!(leg.bid, Some(pos(0.05)));
assert_eq!(leg.ask, Some(pos(0.06)));
assert_eq!(leg.delta, Some(Decimal::new(55, 2)));
assert_eq!(leg.key.underlying, "BTC");
assert_eq!(leg.style, OptionStyle::Call);
}
Err(e) => panic!("expected a normalized leg, got {e}"),
}
}
#[test]
fn test_deribit_normalize_leg_drops_crossed_quote_keeps_row() {
let option = OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25-60000-C",
Some(60_000.0),
Some(DeribitOptionType::Call),
Some(1_751_011_200_000),
),
ticker: ticker("BTC-27JUN25-60000-C", Some(0.06), Some(0.03), None, None),
};
match normalize_leg(&option) {
Ok(leg) => {
assert_eq!(leg.bid, None);
assert_eq!(leg.ask, None);
assert_eq!(leg.key.strike, pos(60_000.0));
}
Err(e) => panic!("a crossed quote drops only the quote, got {e}"),
}
}
#[test]
fn test_deribit_assemble_chain_merges_call_and_put_into_row() {
let call = match normalize_leg(&sample_option()) {
Ok(leg) => leg,
Err(e) => panic!("call leg should normalize, got {e}"),
};
let put_option = OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25-60000-P",
Some(60_000.0),
Some(DeribitOptionType::Put),
Some(1_751_011_200_000),
),
ticker: ticker(
"BTC-27JUN25-60000-P",
Some(0.04),
Some(0.05),
Some(50.0),
Some(greeks(Some(-0.45), Some(0.0001))),
),
};
let put = match normalize_leg(&put_option) {
Ok(leg) => leg,
Err(e) => panic!("put leg should normalize, got {e}"),
};
let provider = deribit_provider_id();
let fetch = assemble_chain(
"BTC",
pos(60_500.0),
utc_millis(1_751_011_200_000),
&[call, put],
&provider,
);
assert_eq!(fetch.chain.options.len(), 1);
assert_eq!(fetch.chain.symbol, "BTC");
assert_eq!(fetch.aliases.len(), 2);
assert!(
fetch
.aliases
.resolve_symbol("BTC-27JUN25-60000-C")
.is_some()
);
assert!(
fetch
.aliases
.resolve_symbol("BTC-27JUN25-60000-P")
.is_some()
);
assert_eq!(fetch.expiry_source.underlying, "BTC");
assert_eq!(
fetch.expiry_source.expiration_utc,
utc_millis(1_751_011_200_000)
);
match fetch.chain.options.iter().next() {
Some(row) => {
assert_eq!(row.strike_price, pos(60_000.0));
assert_eq!(row.call_bid, Some(pos(0.05)));
assert_eq!(row.put_bid, Some(pos(0.04)));
}
None => panic!("expected exactly one strike row"),
}
}
#[test]
fn test_deribit_assemble_chain_is_order_independent() {
let call = match normalize_leg(&sample_option()) {
Ok(leg) => leg,
Err(e) => panic!("call leg should normalize, got {e}"),
};
let put_option = OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25-60000-P",
Some(60_000.0),
Some(DeribitOptionType::Put),
Some(1_751_011_200_000),
),
ticker: ticker(
"BTC-27JUN25-60000-P",
Some(0.04),
Some(0.05),
Some(50.0),
Some(greeks(Some(-0.45), Some(0.0001))),
),
};
let put = match normalize_leg(&put_option) {
Ok(leg) => leg,
Err(e) => panic!("put leg should normalize, got {e}"),
};
let provider = deribit_provider_id();
let expiry = utc_millis(1_751_011_200_000);
let forward = assemble_chain(
"BTC",
pos(60_500.0),
expiry,
&[call.clone(), put.clone()],
&provider,
);
let reversed = assemble_chain("BTC", pos(60_500.0), expiry, &[put, call], &provider);
assert_eq!(forward.chain.options.len(), reversed.chain.options.len());
assert_eq!(forward.aliases.len(), reversed.aliases.len());
let forward_row = forward.chain.options.iter().next();
let reversed_row = reversed.chain.options.iter().next();
match (forward_row, reversed_row) {
(Some(a), Some(b)) => {
assert_eq!(a.strike_price, b.strike_price);
assert_eq!(a.call_bid, b.call_bid);
assert_eq!(a.call_ask, b.call_ask);
assert_eq!(a.put_bid, b.put_bid);
assert_eq!(a.put_ask, b.put_ask);
}
_ => panic!("both orderings must yield exactly one strike row"),
}
}
#[track_caller]
fn sample_leg() -> NormalizedLeg {
match normalize_leg(&sample_option()) {
Ok(leg) => leg,
Err(e) => panic!("sample option should normalize, got {e}"),
}
}
#[test]
fn test_deribit_collect_outcomes_partial_keeps_hydrated_legs() {
let hydration = collect_outcomes(vec![
LegOutcome::Hydrated(Box::new(sample_leg())),
LegOutcome::TransportFailed,
LegOutcome::Dropped,
]);
assert_eq!(hydration.legs.len(), 1);
assert_eq!(hydration.transport_failures, 1);
}
#[test]
fn test_deribit_collect_outcomes_counts_all_transport_failures() {
let hydration = collect_outcomes(vec![
LegOutcome::TransportFailed,
LegOutcome::TransportFailed,
LegOutcome::TransportFailed,
]);
assert!(hydration.legs.is_empty());
assert_eq!(hydration.transport_failures, 3);
}
#[test]
fn test_deribit_all_tickers_fail_surfaces_transport_outage() {
let err = empty_expiry_outcome(8, 8, "BTC", utc_millis(1_751_011_200_000));
assert!(matches!(err, ProviderError::Transport(_)));
assert_eq!(err.to_string(), "upstream transport: transport: closed");
}
#[test]
fn test_deribit_partial_hydration_is_not_an_outage() {
let leg_count = collect_outcomes(vec![
LegOutcome::Hydrated(Box::new(sample_leg())),
LegOutcome::TransportFailed,
])
.legs
.len();
assert_eq!(
leg_count, 1,
"a partial hydration keeps its successful legs"
);
}
#[test]
fn test_deribit_empty_instrument_list_is_no_chain() {
let err = empty_expiry_outcome(0, 0, "BTC", utc_millis(1_751_011_200_000));
match err {
ProviderError::NoChain {
underlying,
expiration,
} => {
assert_eq!(underlying, "BTC");
assert_eq!(expiration, utc_millis(1_751_011_200_000).to_rfc3339());
}
other => panic!("expected NoChain for an empty instrument list, got {other:?}"),
}
}
#[test]
fn test_deribit_tickers_answered_but_unnormalizable_is_no_chain() {
let err = empty_expiry_outcome(5, 0, "ETH", utc_millis(1_751_011_200_000));
assert!(matches!(err, ProviderError::NoChain { .. }));
}
#[test]
fn test_deribit_transport_error_maps_by_category() {
assert!(matches!(
transport_error(&HttpError::AuthenticationFailed(
"secret-bearing".to_owned()
)),
ProviderError::Auth
));
assert!(matches!(
transport_error(&HttpError::RateLimitExceeded),
ProviderError::RateLimited(None)
));
let rendered = transport_error(&HttpError::RequestFailed(
"https://user:pass@example/secret".to_owned(),
))
.to_string();
assert!(!rendered.contains("example"));
assert!(!rendered.contains("pass"));
assert_eq!(rendered, "upstream transport: transport: http");
}
#[tokio::test]
async fn test_deribit_subscribe_spawns_cancellable_loop() {
let adapter = DeribitAdapter::new();
let (tx, _rx) = mpsc::channel::<MarketUpdate>(4);
let sink = sink_from(&tx);
let request = SubscriptionRequest::new(
"BTC",
utc_millis(1_751_011_200_000),
Vec::new(),
CancellationToken::new(),
);
match adapter.subscribe(request, sink).await {
Ok(mut handle) => {
let join = handle.take_join_handle();
assert!(
join.is_some(),
"a spawned subscription carries a join handle"
);
if let Some(join) = join {
join.abort();
}
drop(handle);
}
Err(e) => panic!("subscribe should spawn the reconnect loop, got {e:?}"),
}
}
#[test]
fn test_deribit_adapter_reports_id_and_capabilities() {
let adapter = DeribitAdapter::new();
assert_eq!(adapter.id().as_str(), "deribit");
assert_eq!(adapter.capabilities().chain, ChainCapability::Assemble);
}
#[test]
fn test_deribit_adapter_is_send_sync() {
assert_send_sync::<DeribitAdapter>();
}
proptest! {
#![proptest_config(ProptestConfig { cases: 512, ..ProptestConfig::default() })]
#[test]
fn prop_normalize_quote_is_total(
bid in proptest::num::f64::ANY,
ask in proptest::num::f64::ANY,
) {
match normalize_quote(Some(bid), Some(ask)) {
Ok(quote) => {
if let (Some(b), Some(a)) = (quote.bid, quote.ask) {
prop_assert!(a >= b);
}
}
Err(kind) => prop_assert_eq!(kind, NormalizeKind::OutOfRange("ask")),
}
}
#[test]
fn prop_normalize_iv_is_total(raw in proptest::num::f64::ANY) {
match normalize_iv(raw) {
Ok(iv) => prop_assert!(iv >= Positive::ZERO),
Err(kind) => prop_assert!(
kind == NormalizeKind::NonFinite("iv")
|| kind == NormalizeKind::OutOfRange("iv")
),
}
}
#[test]
fn prop_parse_instrument_name_is_total(name in ".{0,24}") {
let _ = parse_instrument_name(&name);
}
#[test]
fn prop_normalize_leg_is_total(
name in prop_oneof![
Just("BTC-27JUN25-60000-C".to_owned()),
Just("ETH-3JAN25-2000-P".to_owned()),
Just("BTC-27JUN25--C".to_owned()),
Just("garbage".to_owned()),
".{0,16}",
],
strike in prop_oneof![Just(None), proptest::num::f64::ANY.prop_map(Some)],
option_type in prop_oneof![
Just(None),
Just(Some(DeribitOptionType::Call)),
Just(Some(DeribitOptionType::Put)),
],
expiry_ms in prop_oneof![Just(None), any::<i64>().prop_map(Some)],
bid in prop_oneof![Just(None), proptest::num::f64::ANY.prop_map(Some)],
ask in prop_oneof![Just(None), proptest::num::f64::ANY.prop_map(Some)],
iv in prop_oneof![Just(None), proptest::num::f64::ANY.prop_map(Some)],
) {
let option = OptionInstrument {
instrument: deribit_instrument(&name, strike, option_type, expiry_ms),
ticker: ticker(&name, bid, ask, iv, Some(greeks(Some(0.5), Some(0.01)))),
};
if let Ok(leg) = normalize_leg(&option) {
prop_assert!(leg.key.strike > Positive::ZERO);
if let Some(value) = leg.iv {
prop_assert!(value >= Positive::ZERO);
}
}
}
}
fn sample_instrument() -> Instrument {
Instrument {
key: InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: utc_millis(1_751_011_200_000),
strike: pos(60_000.0),
style: OptionStyle::Call,
},
provider: deribit_provider_id(),
native_symbol: "BTC-27JUN25-60000-C".to_owned(),
stream_symbol: None,
spec: ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: "BTC".to_owned(),
},
}
}
fn instrument_at(strike: f64) -> Instrument {
Instrument {
key: InstrumentKey {
strike: pos(strike),
..sample_instrument().key
},
native_symbol: format!("BTC-27JUN25-{strike}-C"),
..sample_instrument()
}
}
fn greeks_payload(delta: Option<f64>, gamma: Option<f64>) -> GreeksPayload {
GreeksPayload { delta, gamma }
}
fn ticker_payload(
bid: Option<f64>,
ask: Option<f64>,
mark_iv: Option<f64>,
greeks: Option<GreeksPayload>,
) -> TickerPayload {
TickerPayload {
best_bid_price: bid,
best_ask_price: ask,
best_bid_amount: Some(5.0),
best_ask_amount: Some(4.0),
last_price: Some(0.055),
mark_iv,
timestamp: Some(1_751_011_200_000),
greeks,
}
}
fn quote_update(strike: f64, bid: f64) -> MarketUpdate {
let payload = ticker_payload(Some(bid), Some(bid + 0.1), Some(50.0), None);
let (quote, _greeks) = normalize_ticker(&instrument_at(strike), &payload, utc_millis(0));
MarketUpdate::Quote(quote)
}
fn quote_strike_bid(update: &MarketUpdate) -> Option<(Positive, Option<Positive>)> {
match update {
MarketUpdate::Quote(quote) => Some((quote.instrument.key.strike, quote.bid)),
MarketUpdate::Greeks(_)
| MarketUpdate::Depth(_)
| MarketUpdate::Chain(_)
| MarketUpdate::Health(_, _) => None,
}
}
fn drain_channel(rx: &mut mpsc::Receiver<MarketUpdate>) -> Vec<MarketUpdate> {
let mut out = Vec::new();
while let Ok(update) = rx.try_recv() {
out.push(update);
}
out
}
#[track_caller]
fn assert_delay_ms(delay: Duration, expected_ms: f64) {
let got = delay.as_secs_f64() * 1000.0;
assert!(
(got - expected_ms).abs() < 1.0,
"expected ~{expected_ms}ms, got {got}ms"
);
}
fn opt_f64() -> impl Strategy<Value = Option<f64>> {
prop_oneof![Just(None), proptest::num::f64::ANY.prop_map(Some)]
}
#[test]
fn test_deribit_backoff_attempt_zero_is_base() {
assert_delay_ms(backoff_delay(0, 0.0), 250.0);
}
#[test]
fn test_deribit_backoff_doubles_per_attempt() {
assert_delay_ms(backoff_delay(1, 0.0), 500.0);
assert_delay_ms(backoff_delay(2, 0.0), 1000.0);
assert_delay_ms(backoff_delay(3, 0.0), 2000.0);
}
#[test]
fn test_deribit_backoff_caps_at_max() {
assert_delay_ms(backoff_delay(100, 0.0), 30_000.0);
assert_delay_ms(backoff_delay(u32::MAX, 0.0), 30_000.0);
}
#[test]
fn test_deribit_backoff_jitter_widens_range() {
assert_delay_ms(backoff_delay(5, -0.2), 8000.0 * 0.8);
assert_delay_ms(backoff_delay(5, 0.2), 8000.0 * 1.2);
assert!(backoff_delay(5, -0.2) < backoff_delay(5, 0.2));
}
#[test]
fn test_deribit_backoff_clamps_out_of_range_jitter() {
assert_eq!(backoff_delay(0, 1.0), backoff_delay(0, 0.2));
assert_eq!(backoff_delay(0, -1.0), backoff_delay(0, -0.2));
}
#[test]
fn test_deribit_backoff_never_exceeds_max_plus_jitter() {
for attempt in 0..40u32 {
let delay = backoff_delay(attempt, 0.2).as_secs_f64();
assert!(
delay <= 36.0 + 1e-6,
"attempt {attempt} exceeded 36 s: {delay}"
);
let low = backoff_delay(attempt, -0.2).as_secs_f64();
assert!(low >= 0.2 - 1e-6, "attempt {attempt} below 200 ms: {low}");
}
}
#[test]
fn test_deribit_normalize_ticker_maps_quote_and_greeks() {
let payload = ticker_payload(
Some(0.05),
Some(0.06),
Some(49.22),
Some(greeks_payload(Some(0.55), Some(0.0001))),
);
let (quote, greeks) = normalize_ticker(
&sample_instrument(),
&payload,
utc_millis(1_751_011_200_000),
);
assert_eq!(quote.bid, Some(pos(0.05)));
assert_eq!(quote.ask, Some(pos(0.06)));
assert_eq!(quote.last, Some(pos(0.055)));
assert_eq!(quote.bid_size, Some(pos(5.0)));
assert_eq!(quote.ask_size, Some(pos(4.0)));
assert_eq!(quote.event_time, Some(utc_millis(1_751_011_200_000)));
assert_eq!(greeks.iv, Some(pos(0.4922)));
assert_eq!(greeks.delta, Some(Decimal::new(55, 2)));
assert!(greeks.gamma.is_some());
assert_eq!(greeks.origin, GreeksOrigin::Provider);
assert_eq!(greeks.event_time, Some(utc_millis(1_751_011_200_000)));
}
#[test]
fn test_deribit_normalize_ticker_discards_theta_vega_rho() {
use deribit_websocket::prelude::json;
let data = json!({
"best_bid_price": 0.05,
"best_ask_price": 0.06,
"mark_iv": 50.0,
"greeks": { "delta": 0.5, "gamma": 0.001, "theta": -9.9, "vega": 8.8, "rho": 7.7 }
});
let payload = match TickerPayload::deserialize(&data) {
Ok(payload) => payload,
Err(e) => panic!("ticker payload should deserialize: {e}"),
};
let (_quote, greeks) = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(
greeks.delta,
Some(Decimal::new(5, 1)),
"venue delta forwarded"
);
assert!(greeks.gamma.is_some(), "venue gamma forwarded");
assert!(greeks.theta.is_none(), "streamed theta must be discarded");
assert!(greeks.vega.is_none(), "streamed vega must be discarded");
assert!(greeks.rho.is_none(), "streamed rho must be discarded");
}
#[test]
fn test_deribit_normalize_ticker_crossed_quote_drops_bid_ask() {
let payload = ticker_payload(Some(0.06), Some(0.03), None, None);
let (quote, _greeks) = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(quote.bid, None);
assert_eq!(quote.ask, None);
assert_eq!(quote.last, Some(pos(0.055)));
}
#[test]
fn test_deribit_normalize_ticker_missing_greeks_are_none() {
let payload = ticker_payload(Some(0.05), Some(0.06), None, None);
let (_quote, greeks) = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
assert!(greeks.iv.is_none());
assert!(greeks.delta.is_none());
assert!(greeks.gamma.is_none());
}
#[test]
fn test_deribit_ticker_payload_deserializes_from_json() {
use deribit_websocket::prelude::json;
let data = json!({
"best_bid_price": 0.05,
"best_ask_price": 0.06,
"mark_iv": 49.22,
"timestamp": 1_751_011_200_000i64,
"greeks": { "delta": 0.55, "gamma": 0.0001, "theta": -1.0, "vega": 2.0, "rho": 3.0 }
});
match TickerPayload::deserialize(&data) {
Ok(payload) => {
let (quote, greeks) =
normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(quote.bid, Some(pos(0.05)));
assert_eq!(greeks.iv, Some(pos(0.4922)));
assert!(greeks.theta.is_none());
}
Err(e) => panic!("ticker payload should deserialize: {e}"),
}
}
#[test]
fn test_deribit_normalize_book_captures_change_id_and_levels() {
let payload = BookPayload {
change_id: Some(770),
timestamp: Some(1_751_011_200_000),
bids: vec![
BookLevel::Priced([60_000.0, 2.0]),
BookLevel::Priced([59_990.0, 5.0]),
],
asks: vec![BookLevel::Priced([60_010.0, 1.5])],
};
let ladder = normalize_book(
&sample_instrument(),
&payload,
utc_millis(1_751_011_200_000),
);
assert_eq!(ladder.change_id, Some(770));
assert_eq!(ladder.event_time, Some(utc_millis(1_751_011_200_000)));
assert_eq!(ladder.bids.len(), 2);
assert_eq!(ladder.asks.len(), 1);
match ladder.bids.first() {
Some(level) => {
assert_eq!(level.price, pos(60_000.0));
assert_eq!(level.size, pos(2.0));
}
None => panic!("expected the best bid at index 0"),
}
}
#[test]
fn test_deribit_normalize_book_drops_invalid_levels() {
let payload = BookPayload {
change_id: Some(1),
timestamp: None,
bids: vec![
BookLevel::Priced([f64::NAN, 1.0]),
BookLevel::Priced([60_000.0, 2.0]),
BookLevel::Priced([59_000.0, -1.0]),
],
asks: Vec::new(),
};
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(ladder.bids.len(), 1, "NaN price and negative size dropped");
assert_eq!(ladder.event_time, None, "no timestamp -> no event_time");
}
#[test]
fn test_deribit_normalize_book_decodes_raw_action_levels() {
let payload = BookPayload {
change_id: Some(2),
timestamp: Some(1_751_011_200_000),
bids: vec![BookLevel::Actioned("new".to_owned(), 60_000.0, 3.0)],
asks: vec![BookLevel::Actioned("delete".to_owned(), 60_010.0, 0.0)],
};
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
match ladder.bids.first() {
Some(level) => assert_eq!(level.price, pos(60_000.0)),
None => panic!("the raw `new` bid should decode"),
}
assert_eq!(ladder.asks.len(), 1);
}
#[test]
fn test_deribit_book_payload_deserializes_both_level_encodings() {
use deribit_websocket::prelude::json;
let aggregated =
json!({ "change_id": 5, "bids": [[60_000.0, 2.0]], "asks": [[60_010.0, 1.0]] });
match BookPayload::deserialize(&aggregated) {
Ok(payload) => {
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(ladder.change_id, Some(5));
assert_eq!(ladder.bids.len(), 1);
}
Err(e) => panic!("aggregated `[price, amount]` book should deserialize: {e}"),
}
let raw = json!({ "change_id": 6, "bids": [["new", 60_000.0, 2.0]], "asks": [] });
match BookPayload::deserialize(&raw) {
Ok(payload) => {
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(ladder.change_id, Some(6));
assert_eq!(ladder.bids.len(), 1);
}
Err(e) => panic!("raw `[action, price, amount]` book should deserialize: {e}"),
}
}
#[test]
fn test_deribit_producer_staging_overwrites_on_full_channel() {
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(1);
let mut sink = sink_from(&tx);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 1.0)),
SendState::Open
);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 2.0)),
SendState::Open
);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 3.0)),
SendState::Open
);
assert_eq!(sink.staged_len(), 1, "one slot per instrument, overwritten");
let sent = drain_channel(&mut rx);
assert_eq!(sent.len(), 1);
assert_eq!(
sent.first()
.and_then(quote_strike_bid)
.and_then(|(_, bid)| bid),
Some(pos(1.0))
);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 4.0)),
SendState::Open
);
let after = drain_channel(&mut rx);
assert_eq!(after.len(), 1, "flush-on-space delivers the staged value");
assert_eq!(
after
.first()
.and_then(quote_strike_bid)
.and_then(|(_, bid)| bid),
Some(pos(3.0)),
"the freshest staged value survived, not the intermediate 2.0"
);
}
#[test]
fn test_deribit_producer_staging_flush_delivers_freshest_after_quiet() {
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(1);
let mut sink = sink_from(&tx);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 1.0)),
SendState::Open
);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 2.0)),
SendState::Open
);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 3.0)),
SendState::Open
);
assert!(sink.has_pending(), "the burst left a value staged");
let sent = drain_channel(&mut rx);
assert_eq!(sent.len(), 1);
assert_eq!(sink.flush(), SendState::Open);
let after = drain_channel(&mut rx);
assert_eq!(
after.len(),
1,
"flush-on-tick delivers the staged value with no further publish"
);
assert_eq!(
after
.first()
.and_then(quote_strike_bid)
.and_then(|(_, bid)| bid),
Some(pos(3.0)),
"the freshest staged value (3.0) reached the channel after the feed went quiet"
);
assert!(
!sink.has_pending(),
"nothing remains staged once the tick flush drains the slot"
);
}
#[test]
fn test_deribit_producer_staging_keeps_quote_and_greeks_independently() {
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(2);
let mut sink = sink_from(&tx);
let _ = sink.publish_coalesced(quote_update(200.0, 9.0));
let _ = sink.publish_coalesced(quote_update(201.0, 9.0));
let payload = ticker_payload(
Some(0.05),
Some(0.06),
Some(50.0),
Some(greeks_payload(Some(0.5), Some(0.01))),
);
let (quote, greeks) = normalize_ticker(&instrument_at(100.0), &payload, utc_millis(0));
let _ = sink.publish_coalesced(MarketUpdate::Quote(quote));
let _ = sink.publish_coalesced(MarketUpdate::Greeks(greeks));
assert_eq!(sink.staged_len(), 1, "one slot holds both kinds");
let _ = drain_channel(&mut rx);
let _ = sink.publish_coalesced(quote_update(300.0, 1.0));
let flushed = drain_channel(&mut rx);
assert!(
flushed.iter().any(|update| matches!(
update,
MarketUpdate::Quote(quote) if quote.instrument.key.strike == pos(100.0)
)),
"the staged quote flushed"
);
assert!(
flushed
.iter()
.any(|update| matches!(update, MarketUpdate::Greeks(_))),
"the staged greeks flushed"
);
}
#[test]
fn test_deribit_producer_staging_is_bounded_by_instruments() {
let (tx, _rx) = mpsc::channel::<MarketUpdate>(1);
let mut sink = sink_from(&tx);
let _ = sink.publish_coalesced(quote_update(1.0, 1.0));
for round in 0..200u32 {
for strike in [1.0, 2.0, 3.0] {
assert_eq!(
sink.publish_coalesced(quote_update(strike, f64::from(round) + 1.0)),
SendState::Open
);
}
assert!(
sink.staged_len() <= 3,
"staging is O(N=3 instruments), not O(burst): round {round}"
);
}
}
#[test]
fn test_deribit_producer_staging_reports_closed_channel() {
let (tx, rx) = mpsc::channel::<MarketUpdate>(4);
drop(rx);
let mut sink = sink_from(&tx);
assert_eq!(
sink.publish_coalesced(quote_update(100.0, 1.0)),
SendState::Closed,
"a closed consumer channel stops the loop, never a silent buffer"
);
}
#[test]
fn test_deribit_route_message_ticker_publishes_quote_and_greeks() {
use deribit_websocket::prelude::json;
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(8);
let lookup = instrument_lookup(&[sample_instrument()]);
let mut sink = sink_from(&tx);
let text = json!({
"jsonrpc": "2.0",
"method": "subscription",
"params": {
"channel": "ticker.BTC-27JUN25-60000-C",
"data": { "best_bid_price": 0.05, "best_ask_price": 0.06, "mark_iv": 49.22,
"greeks": { "delta": 0.55, "gamma": 0.0001, "theta": -1.0 } }
}
})
.to_string();
assert_eq!(route_message(&text, &lookup, &mut sink), SendState::Open);
let out = drain_channel(&mut rx);
assert!(out.iter().any(|u| matches!(u, MarketUpdate::Quote(_))));
assert!(
out.iter()
.any(|u| matches!(u, MarketUpdate::Greeks(g) if g.theta.is_none())),
"greeks published with theta discarded"
);
}
#[test]
fn test_deribit_route_message_ticker_with_interval_suffix_still_routes() {
use deribit_websocket::prelude::json;
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(8);
let lookup = instrument_lookup(&[sample_instrument()]);
let mut sink = sink_from(&tx);
let text = json!({
"jsonrpc": "2.0",
"method": "subscription",
"params": {
"channel": "ticker.BTC-27JUN25-60000-C.100ms",
"data": { "best_bid_price": 0.05, "best_ask_price": 0.06 }
}
})
.to_string();
assert_eq!(route_message(&text, &lookup, &mut sink), SendState::Open);
let out = drain_channel(&mut rx);
match out.iter().find_map(quote_strike_bid) {
Some((strike, bid)) => {
assert_eq!(strike, pos(60_000.0));
assert_eq!(bid, Some(pos(0.05)));
}
None => panic!("a ticker frame with a trailing interval must still route"),
}
}
#[test]
fn test_deribit_route_message_book_publishes_depth_with_change_id() {
use deribit_websocket::prelude::json;
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(8);
let lookup = instrument_lookup(&[sample_instrument()]);
let mut sink = sink_from(&tx);
let text = json!({
"jsonrpc": "2.0",
"method": "subscription",
"params": {
"channel": "book.BTC-27JUN25-60000-C.none.20.100ms",
"data": {
"change_id": 770,
"bids": [[60_000.0, 2.0], [59_990.0, 5.0]],
"asks": [[60_010.0, 1.0]]
}
}
})
.to_string();
assert_eq!(route_message(&text, &lookup, &mut sink), SendState::Open);
match drain_channel(&mut rx).first() {
Some(MarketUpdate::Depth(ladder)) => {
assert_eq!(ladder.change_id, Some(770));
assert_eq!(ladder.bids.len(), 2, "the grouped frame is a full ladder");
assert_eq!(ladder.asks.len(), 1);
}
other => panic!("expected a Depth update with change_id, got {other:?}"),
}
}
#[test]
fn test_deribit_route_message_unknown_symbol_is_dropped() {
use deribit_websocket::prelude::json;
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(8);
let lookup = instrument_lookup(&[sample_instrument()]);
let mut sink = sink_from(&tx);
let text = json!({
"jsonrpc": "2.0",
"method": "subscription",
"params": {
"channel": "ticker.BTC-27JUN25-99999-C",
"data": { "best_bid_price": 0.05, "best_ask_price": 0.06 }
}
})
.to_string();
assert_eq!(route_message(&text, &lookup, &mut sink), SendState::Open);
assert!(
drain_channel(&mut rx).is_empty(),
"an update for an unsubscribed symbol is dropped, never keyed blindly"
);
}
#[test]
fn test_deribit_route_message_ignores_non_subscription_and_malformed() {
use deribit_websocket::prelude::json;
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(8);
let lookup = instrument_lookup(&[sample_instrument()]);
let mut sink = sink_from(&tx);
let heartbeat =
json!({ "jsonrpc": "2.0", "method": "heartbeat", "params": {} }).to_string();
assert_eq!(
route_message(&heartbeat, &lookup, &mut sink),
SendState::Open
);
assert_eq!(
route_message("{ not json", &lookup, &mut sink),
SendState::Open
);
assert!(drain_channel(&mut rx).is_empty());
}
#[test]
fn test_deribit_subscription_channels_ticker_and_book_never_trades() {
let channels = subscription_channels(&[sample_instrument()]);
assert!(channels.contains(&"ticker.BTC-27JUN25-60000-C".to_owned()));
assert!(
channels.contains(&"book.BTC-27JUN25-60000-C.none.20.100ms".to_owned()),
"the book leg is the grouped full-snapshot channel (group/depth/interval), \
so the coalesce + overwrite fold is correct (#48)"
);
assert!(
!channels
.iter()
.any(|channel| channel == "book.BTC-27JUN25-60000-C.raw"),
"the raw delta book channel is never subscribed — it would tear the book"
);
assert!(
!channels
.iter()
.any(|channel| channel.starts_with("trades.")),
"the trades tape is deferred, never subscribed"
);
assert_eq!(channels.len(), 2, "exactly ticker + book per leg");
}
#[test]
fn test_deribit_chain_snapshot_from_fetch_is_merged_live_and_carries_aliases() {
let call = match normalize_leg(&sample_option()) {
Ok(leg) => leg,
Err(e) => panic!("call leg should normalize, got {e}"),
};
let fetch = assemble_chain(
"BTC",
pos(60_500.0),
utc_millis(1_751_011_200_000),
&[call],
&deribit_provider_id(),
);
let snapshot = chain_snapshot(&fetch, utc_millis(1_751_011_200_001));
assert_eq!(snapshot.source, ChainSource::Merged);
assert!(matches!(snapshot.health, StreamHealth::Live));
assert_eq!(snapshot.chain_key.1, "BTC");
assert_eq!(snapshot.chain_key.2, utc_millis(1_751_011_200_000));
assert_eq!(snapshot.last_full_poll, Some(utc_millis(1_751_011_200_001)));
assert!(
snapshot
.aliases
.resolve_symbol("BTC-27JUN25-60000-C")
.is_some()
);
}
proptest! {
#![proptest_config(ProptestConfig { cases: 256, ..ProptestConfig::default() })]
#[test]
fn prop_deribit_backoff_is_bounded(attempt in 0u32..64, jitter in -5.0f64..5.0) {
let delay = backoff_delay(attempt, jitter).as_secs_f64();
prop_assert!(delay >= 0.2 - 1e-9);
prop_assert!(delay <= 36.0 + 1e-9);
}
#[test]
fn prop_deribit_normalize_ticker_is_total(
bid in opt_f64(),
ask in opt_f64(),
last in opt_f64(),
bid_amt in opt_f64(),
ask_amt in opt_f64(),
iv in opt_f64(),
delta in opt_f64(),
gamma in opt_f64(),
ts in prop_oneof![Just(None), any::<i64>().prop_map(Some)],
) {
let payload = TickerPayload {
best_bid_price: bid,
best_ask_price: ask,
best_bid_amount: bid_amt,
best_ask_amount: ask_amt,
last_price: last,
mark_iv: iv,
timestamp: ts,
greeks: Some(GreeksPayload { delta, gamma }),
};
let (quote, greeks) = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
if let (Some(b), Some(a)) = (quote.bid, quote.ask) {
prop_assert!(a >= b);
}
if let Some(value) = greeks.iv {
prop_assert!(value >= Positive::ZERO);
}
prop_assert!(greeks.theta.is_none());
prop_assert!(greeks.vega.is_none());
prop_assert!(greeks.rho.is_none());
}
#[test]
fn prop_deribit_normalize_book_is_total(
levels in proptest::collection::vec(
(proptest::num::f64::ANY, proptest::num::f64::ANY),
0..8,
),
change_id in prop_oneof![Just(None), any::<u64>().prop_map(Some)],
) {
let payload = BookPayload {
change_id,
timestamp: None,
bids: levels.iter().map(|(p, a)| BookLevel::Priced([*p, *a])).collect(),
asks: Vec::new(),
};
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
prop_assert_eq!(ladder.change_id, change_id);
for level in &ladder.bids {
prop_assert!(level.price >= Positive::ZERO);
prop_assert!(level.size >= Positive::ZERO);
}
}
}
const FIXTURE_INSTRUMENTS_BTC: &str =
include_str!("../../tests/fixtures/deribit/instruments/instruments_btc.json");
const FIXTURE_INSTRUMENTS_MISSING_STRIKE: &str =
include_str!("../../tests/fixtures/deribit/instruments/instruments_missing_strike.json");
const FIXTURE_TICKER_NORMAL: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_normal.json");
const FIXTURE_TICKER_ZERO_BID: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_zero_bid.json");
const FIXTURE_TICKER_CROSSED: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_crossed.json");
const FIXTURE_TICKER_NEGATIVE: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_negative.json");
const FIXTURE_TICKER_NON_FINITE: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_non_finite.json");
const FIXTURE_BOOK_SNAPSHOT: &str =
include_str!("../../tests/fixtures/deribit/book/book_snapshot.json");
const FIXTURE_BOOK_DELTA: &str =
include_str!("../../tests/fixtures/deribit/book/book_delta.json");
const FIXTURE_BOOK_GROUPED_SNAPSHOT: &str =
include_str!("../../tests/fixtures/deribit/book/book_grouped_snapshot.json");
#[track_caller]
fn fixture_value(json: &str) -> Value {
match json.parse::<Value>() {
Ok(value) => value,
Err(e) => panic!("fixture JSON should parse: {e}"),
}
}
#[track_caller]
fn instruments_fixture(json: &str) -> Vec<DeribitInstrument> {
let value = fixture_value(json);
match Vec::<DeribitInstrument>::deserialize(&value) {
Ok(instruments) => instruments,
Err(e) => panic!("instruments fixture should deserialize: {e}"),
}
}
#[track_caller]
fn ticker_data_fixture(json: &str) -> TickerData {
let value = fixture_value(json);
match TickerData::deserialize(&value) {
Ok(ticker) => ticker,
Err(e) => panic!("ticker fixture should deserialize as TickerData: {e}"),
}
}
#[track_caller]
fn ticker_payload_fixture(json: &str) -> TickerPayload {
let value = fixture_value(json);
match TickerPayload::deserialize(&value) {
Ok(payload) => payload,
Err(e) => panic!("ticker fixture should deserialize as TickerPayload: {e}"),
}
}
#[track_caller]
fn book_payload_fixture(json: &str) -> BookPayload {
let value = fixture_value(json);
match BookPayload::deserialize(&value) {
Ok(payload) => payload,
Err(e) => panic!("book fixture should deserialize as BookPayload: {e}"),
}
}
fn ticker_frame(symbol: &str, data_json: &str) -> String {
subscription_frame(&format!("ticker.{symbol}"), data_json)
}
fn subscription_frame(channel: &str, data_json: &str) -> String {
format!(
"{{\"jsonrpc\":\"2.0\",\"method\":\"subscription\",\
\"params\":{{\"channel\":\"{channel}\",\"data\":{data_json}}}}}"
)
}
#[test]
fn test_deribit_fixture_instruments_assemble_to_option_chain() {
let instruments = instruments_fixture(FIXTURE_INSTRUMENTS_BTC);
let options: Vec<DeribitInstrument> = instruments
.iter()
.filter(|i| i.is_option())
.cloned()
.collect();
assert_eq!(options.len(), 3, "the perpetual future is filtered out");
let ticker = ticker_data_fixture(FIXTURE_TICKER_NORMAL);
let legs: Vec<NormalizedLeg> = options
.into_iter()
.filter_map(|instrument| {
normalize_leg(&OptionInstrument {
instrument,
ticker: ticker.clone(),
})
.ok()
})
.collect();
assert_eq!(legs.len(), 3, "every option leg normalizes");
let fetch = assemble_chain(
"BTC",
pos(60_500.0),
utc_millis(1_751_011_200_000),
&legs,
&deribit_provider_id(),
);
assert_eq!(fetch.chain.symbol, "BTC");
assert_eq!(fetch.chain.options.len(), 2);
assert_eq!(
fetch.aliases.len(),
3,
"three native aliases, one per option"
);
assert!(
fetch
.aliases
.resolve_symbol("BTC-27JUN25-60000-P")
.is_some()
);
match fetch
.chain
.options
.iter()
.find(|row| row.strike_price == pos(60_000.0))
{
Some(row) => {
assert_eq!(row.call_bid, Some(pos(0.05)));
assert_eq!(row.call_ask, Some(pos(0.06)));
assert_eq!(row.put_bid, Some(pos(0.05)));
assert_eq!(row.put_ask, Some(pos(0.06)));
}
None => panic!("expected the 60000 strike row"),
}
match legs.iter().find(|leg| leg.style == OptionStyle::Call) {
Some(leg) => {
assert_eq!(leg.iv, Some(pos(0.4922)));
assert_eq!(leg.delta, Some(Decimal::new(55, 2)));
assert_eq!(leg.underlying_price, Some(pos(60_500.0)));
}
None => panic!("expected a normalized call leg"),
}
}
#[test]
fn test_deribit_fixture_missing_strike_and_style_reject_rows() {
let instruments = instruments_fixture(FIXTURE_INSTRUMENTS_MISSING_STRIKE);
let ticker = ticker_data_fixture(FIXTURE_TICKER_NORMAL);
let mut kinds = Vec::new();
for instrument in instruments {
let option = OptionInstrument {
instrument,
ticker: ticker.clone(),
};
match normalize_leg(&option) {
Err(kind) => kinds.push(kind),
Ok(_) => panic!("a missing-strike / unknown-style row must be rejected"),
}
}
assert!(
kinds.contains(&NormalizeKind::MissingField("strike")),
"the empty strike segment is a missing-field reject"
);
assert!(
kinds.contains(&NormalizeKind::UnknownStyle),
"the `-X` style is an unknown-style reject"
);
}
#[test]
fn test_deribit_fixture_ticker_normal_normalizes_quote_and_greeks() {
let payload = ticker_payload_fixture(FIXTURE_TICKER_NORMAL);
let (quote, greeks) = normalize_ticker(
&sample_instrument(),
&payload,
utc_millis(1_751_011_200_000),
);
assert_eq!(quote.bid, Some(pos(0.05)));
assert_eq!(quote.ask, Some(pos(0.06)));
assert_eq!(quote.last, Some(pos(0.055)));
assert_eq!(quote.bid_size, Some(pos(5.0)));
assert_eq!(quote.ask_size, Some(pos(4.0)));
assert_eq!(quote.event_time, Some(utc_millis(1_751_011_200_000)));
assert_eq!(greeks.iv, Some(pos(0.4922)));
assert_eq!(greeks.delta, Some(Decimal::new(55, 2)));
assert!(greeks.gamma.is_some());
assert!(greeks.theta.is_none());
assert!(greeks.vega.is_none());
assert!(greeks.rho.is_none());
}
#[test]
fn test_deribit_fixture_ticker_zero_bid_keeps_quote() {
let payload = ticker_payload_fixture(FIXTURE_TICKER_ZERO_BID);
let (quote, _greeks) = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(quote.bid, Some(Positive::ZERO));
assert_eq!(quote.ask, Some(pos(0.06)));
}
#[test]
fn test_deribit_fixture_ticker_crossed_drops_bid_ask() {
let payload = ticker_payload_fixture(FIXTURE_TICKER_CROSSED);
let (quote, _greeks) = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(quote.bid, None);
assert_eq!(quote.ask, None);
assert_eq!(
quote.last,
Some(pos(0.055)),
"last survives a crossed quote"
);
}
#[test]
fn test_deribit_fixture_ticker_negative_drops_bid_keeps_ask() {
let payload = ticker_payload_fixture(FIXTURE_TICKER_NEGATIVE);
let (quote, _greeks) = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(quote.bid, None);
assert_eq!(quote.ask, Some(pos(0.06)));
}
#[test]
fn test_deribit_fixture_ticker_non_finite_refuses_frame() {
let value = fixture_value(FIXTURE_TICKER_NON_FINITE);
assert!(
TickerPayload::deserialize(&value).is_err(),
"a non-numeric price field refuses the whole frame"
);
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(4);
let lookup = instrument_lookup(&[sample_instrument()]);
let mut sink = sink_from(&tx);
let frame = ticker_frame("BTC-27JUN25-60000-C", FIXTURE_TICKER_NON_FINITE);
assert_eq!(
route_message(&frame, &lookup, &mut sink),
SendState::Open,
"a degraded frame is skipped, never a panic"
);
assert!(
drain_channel(&mut rx).is_empty(),
"the degraded frame produces no update, never a fabricated value"
);
}
#[test]
fn test_deribit_fixture_book_grouped_snapshot_normalizes_full_ladder() {
let payload = book_payload_fixture(FIXTURE_BOOK_GROUPED_SNAPSHOT);
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(ladder.change_id, Some(1_419_600_010));
assert_eq!(ladder.event_time, Some(utc_millis(1_751_011_200_000)));
assert_eq!(
ladder.bids.len(),
3,
"the grouped snapshot is a full ladder"
);
assert_eq!(ladder.asks.len(), 2);
match ladder.bids.first() {
Some(level) => {
assert_eq!(level.price, pos(0.05));
assert_eq!(level.size, pos(5.0));
}
None => panic!("expected the best bid at index 0"),
}
}
#[test]
fn test_deribit_fixture_book_snapshot_normalizes_ladder() {
let payload = book_payload_fixture(FIXTURE_BOOK_SNAPSHOT);
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(ladder.change_id, Some(1_419_600_001));
assert_eq!(ladder.event_time, Some(utc_millis(1_751_011_200_000)));
assert_eq!(
ladder.bids.len(),
3,
"raw-book `new` bids decode best-first"
);
assert_eq!(ladder.asks.len(), 2);
match ladder.bids.first() {
Some(level) => {
assert_eq!(level.price, pos(0.05));
assert_eq!(level.size, pos(5.0));
}
None => panic!("expected the best bid at index 0"),
}
}
#[test]
fn test_deribit_fixture_book_delta_captures_change_and_delete() {
let payload = book_payload_fixture(FIXTURE_BOOK_DELTA);
let ladder = normalize_book(&sample_instrument(), &payload, utc_millis(0));
assert_eq!(ladder.change_id, Some(1_419_600_002));
assert_eq!(ladder.bids.len(), 2, "both change and delete levels decode");
match ladder.bids.iter().find(|level| level.price == pos(0.049)) {
Some(level) => assert_eq!(level.size, Positive::ZERO),
None => panic!("the delete level should decode with size 0"),
}
assert_eq!(ladder.asks.len(), 1);
}
#[test]
fn test_deribit_fixture_corpus_normalizes_without_panic() {
for json in [
FIXTURE_TICKER_NORMAL,
FIXTURE_TICKER_ZERO_BID,
FIXTURE_TICKER_CROSSED,
FIXTURE_TICKER_NEGATIVE,
FIXTURE_TICKER_NON_FINITE,
] {
let value = fixture_value(json);
if let Ok(payload) = TickerPayload::deserialize(&value) {
let _ = normalize_ticker(&sample_instrument(), &payload, utc_millis(0));
}
}
for json in [
FIXTURE_BOOK_GROUPED_SNAPSHOT,
FIXTURE_BOOK_SNAPSHOT,
FIXTURE_BOOK_DELTA,
] {
let payload = book_payload_fixture(json);
let _ = normalize_book(&sample_instrument(), &payload, utc_millis(0));
}
}
enum MockFrame {
Text(String),
Drop,
}
#[derive(Default)]
struct MockState {
connect_calls: usize,
refetch_calls: usize,
channel_sets: Vec<Vec<String>>,
}
struct MockTransport {
state: Arc<StdMutex<MockState>>,
frames: mpsc::UnboundedReceiver<MockFrame>,
canned_fetch: ChainFetch,
}
#[async_trait]
impl DeribitTransport for MockTransport {
async fn connect_and_subscribe(
&mut self,
channels: Vec<String>,
) -> Result<(), TransportGone> {
if let Ok(mut state) = self.state.lock() {
state.connect_calls += 1;
state.channel_sets.push(channels);
}
Ok(())
}
async fn receive(&mut self) -> Result<String, TransportGone> {
match self.frames.recv().await {
Some(MockFrame::Text(text)) => Ok(text),
Some(MockFrame::Drop) => Err(TransportGone),
None => std::future::pending().await,
}
}
async fn refetch(
&mut self,
_underlying: &str,
_expiration: &ExpirationDate,
) -> Option<ChainFetch> {
if let Ok(mut state) = self.state.lock() {
state.refetch_calls += 1;
}
Some(self.canned_fetch.clone())
}
}
fn canned_reconnect_fetch() -> ChainFetch {
let call_60k = match normalize_leg(&sample_option()) {
Ok(leg) => leg,
Err(e) => panic!("the 60000 call should normalize: {e}"),
};
let option_61k = OptionInstrument {
instrument: deribit_instrument(
"BTC-27JUN25-61000-C",
Some(61_000.0),
Some(DeribitOptionType::Call),
Some(1_751_011_200_000),
),
ticker: ticker(
"BTC-27JUN25-61000-C",
Some(0.03),
Some(0.04),
Some(45.0),
Some(greeks(Some(0.4), Some(0.0001))),
),
};
let call_61k = match normalize_leg(&option_61k) {
Ok(leg) => leg,
Err(e) => panic!("the 61000 call should normalize: {e}"),
};
assemble_chain(
"BTC",
pos(60_500.0),
utc_millis(1_751_011_200_000),
&[call_60k, call_61k],
&deribit_provider_id(),
)
}
fn mock_transport() -> (
MockTransport,
mpsc::UnboundedSender<MockFrame>,
Arc<StdMutex<MockState>>,
) {
let (script_tx, script_rx) = mpsc::unbounded_channel::<MockFrame>();
let state = Arc::new(StdMutex::new(MockState::default()));
let transport = MockTransport {
state: Arc::clone(&state),
frames: script_rx,
canned_fetch: canned_reconnect_fetch(),
};
(transport, script_tx, state)
}
fn spawn_loop(
transport: MockTransport,
tx: mpsc::Sender<MarketUpdate>,
) -> (tokio::task::JoinHandle<()>, CancellationToken) {
let cancel = CancellationToken::new();
let sink = sink_from(&tx);
let join = tokio::spawn(run_reconnect_loop(
transport,
deribit_provider_id(),
"BTC".to_owned(),
utc_millis(1_751_011_200_000),
vec![sample_instrument()],
sink,
cancel.clone(),
));
(join, cancel)
}
async fn drain_until<F>(rx: &mut mpsc::Receiver<MarketUpdate>, mut stop: F) -> Vec<MarketUpdate>
where
F: FnMut(&MarketUpdate) -> bool,
{
let mut collected = Vec::new();
loop {
match tokio::time::timeout(Duration::from_secs(5), rx.recv()).await {
Ok(Some(update)) => {
let done = stop(&update);
collected.push(update);
if done {
return collected;
}
}
Ok(None) | Err(_) => return collected,
}
}
}
async fn loop_stopped(join: tokio::task::JoinHandle<()>) -> bool {
matches!(
tokio::time::timeout(Duration::from_secs(5), join).await,
Ok(Ok(()))
)
}
fn is_live(update: &MarketUpdate) -> bool {
matches!(update, MarketUpdate::Health(_, StreamHealth::Live))
}
fn is_reconnecting(update: &MarketUpdate) -> bool {
matches!(
update,
MarketUpdate::Health(_, StreamHealth::Reconnecting { .. })
)
}
fn is_chain(update: &MarketUpdate) -> bool {
matches!(update, MarketUpdate::Chain(_))
}
fn reconnect_attempt(update: &MarketUpdate) -> Option<u32> {
match update {
MarketUpdate::Health(_, StreamHealth::Reconnecting { attempt }) => Some(*attempt),
MarketUpdate::Quote(_)
| MarketUpdate::Greeks(_)
| MarketUpdate::Depth(_)
| MarketUpdate::Chain(_)
| MarketUpdate::Health(_, StreamHealth::Live | StreamHealth::Stale { .. }) => None,
}
}
fn burst_ticker_json(round: u32) -> String {
let bid = 0.01 + f64::from(round) * 0.001;
let ask = bid + 0.01;
format!(
"{{\"best_bid_price\":{bid},\"best_ask_price\":{ask},\
\"best_bid_amount\":5.0,\"best_ask_amount\":4.0,\"mark_iv\":50.0}}"
)
}
#[tokio::test(start_paused = true)]
async fn test_deribit_lifecycle_socket_close_emits_reconnecting() {
let (transport, script_tx, _state) = mock_transport();
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(16);
let (join, cancel) = spawn_loop(transport, tx);
let _ = script_tx.send(MockFrame::Text(ticker_frame(
"BTC-27JUN25-60000-C",
FIXTURE_TICKER_NORMAL,
)));
let _ = script_tx.send(MockFrame::Drop);
let updates = drain_until(&mut rx, is_reconnecting).await;
assert!(
updates.iter().any(is_live),
"the stream first surfaced Live"
);
assert!(
updates.iter().any(|u| matches!(u, MarketUpdate::Quote(_))),
"the live frame produced a quote before the drop"
);
match updates.last().and_then(reconnect_attempt) {
Some(attempt) => assert_eq!(attempt, 1, "the first drop surfaces Reconnecting{{1}}"),
None => panic!("a socket close must surface Health(Reconnecting)"),
}
cancel.cancel();
assert!(loop_stopped(join).await, "the cancelled loop stops");
}
#[tokio::test(start_paused = true)]
async fn test_deribit_lifecycle_stream_error_reconnects_without_panic() {
let (transport, script_tx, state) = mock_transport();
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(16);
let (join, cancel) = spawn_loop(transport, tx);
let _ = script_tx.send(MockFrame::Drop);
let mut lives = 0;
let updates = drain_until(&mut rx, |u| {
if is_live(u) {
lives += 1;
}
lives >= 2
})
.await;
assert!(
updates.iter().any(is_reconnecting),
"the stream error surfaced Reconnecting, no panic"
);
assert_eq!(
updates.iter().filter(|u| is_live(u)).count(),
2,
"the loop reconnected to Live after the error"
);
let refetches = state.lock().map(|s| s.refetch_calls).unwrap_or(0);
assert!(
refetches >= 1,
"the reconnect re-fetched the chain (backfill)"
);
cancel.cancel();
assert!(loop_stopped(join).await);
}
#[tokio::test(start_paused = true)]
async fn test_deribit_lifecycle_reconnect_refetches_and_resubscribes_fresh() {
let (transport, script_tx, state) = mock_transport();
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(32);
let (join, cancel) = spawn_loop(transport, tx);
let _ = script_tx.send(MockFrame::Drop);
let _ = script_tx.send(MockFrame::Drop);
let mut reconnects = 0;
let updates = drain_until(&mut rx, |u| {
if is_reconnecting(u) {
reconnects += 1;
}
reconnects >= 2
})
.await;
let attempts: Vec<u32> = updates.iter().filter_map(reconnect_attempt).collect();
assert_eq!(
attempts,
vec![1, 1],
"attempt resets to 1 after each successful resubscribe"
);
assert!(
updates.iter().any(is_chain),
"the reconnect backfill emitted a Chain snapshot"
);
let (connects, refetches, sets) = match state.lock() {
Ok(guard) => (
guard.connect_calls,
guard.refetch_calls,
guard.channel_sets.clone(),
),
Err(_) => panic!("mock state lock poisoned"),
};
assert!(connects >= 2, "the loop reconnected");
assert!(refetches >= 1, "the reconnect re-issued fetch_chain");
match sets.first() {
Some(first) => assert_eq!(first.len(), 2, "initial subscribe = ticker + book"),
None => panic!("expected an initial subscribe set"),
}
assert!(
sets.iter()
.skip(1)
.any(|set| set.iter().any(|c| c == "ticker.BTC-27JUN25-61000-C")),
"a resubscribe used the fresh 61000-C alias from the re-fetch"
);
cancel.cancel();
assert!(loop_stopped(join).await);
}
#[test]
fn test_deribit_lifecycle_saturation_coalesces_with_flat_memory() {
let (tx, _rx) = mpsc::channel::<MarketUpdate>(1);
let instruments = [
instrument_at(60_000.0),
instrument_at(61_000.0),
instrument_at(62_000.0),
];
let lookup = instrument_lookup(&instruments);
let mut sink = sink_from(&tx);
for round in 0..500u32 {
for strike in [60_000.0, 61_000.0, 62_000.0] {
let symbol = format!("BTC-27JUN25-{strike}-C");
let frame = ticker_frame(&symbol, &burst_ticker_json(round));
assert_eq!(
route_message(&frame, &lookup, &mut sink),
SendState::Open,
"a saturated bridge never drops the loop"
);
}
assert!(
sink.staged_len() <= 3,
"staging stays O(N = 3 instruments), not O(burst): round {round}"
);
}
drop(_rx);
}
#[tokio::test(start_paused = true)]
async fn test_deribit_lifecycle_lag_does_not_stall_control_updates() {
let (transport, script_tx, _state) = mock_transport();
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(1);
let (join, cancel) = spawn_loop(transport, tx);
for round in 0..50u32 {
let _ = script_tx.send(MockFrame::Text(ticker_frame(
"BTC-27JUN25-60000-C",
&burst_ticker_json(round),
)));
}
let _ = script_tx.send(MockFrame::Drop);
let updates = drain_until(&mut rx, is_chain).await;
assert!(
updates.iter().any(is_live),
"Health(Live) delivered despite the lagging consumer"
);
assert!(
updates.iter().any(is_reconnecting),
"Health(Reconnecting) not stalled by the quote backlog"
);
assert!(
updates.iter().any(is_chain),
"the Chain backfill reached the lagging consumer"
);
cancel.cancel();
assert!(loop_stopped(join).await);
}
#[tokio::test(start_paused = true)]
async fn test_deribit_lifecycle_shutdown_on_handle_drop_stops_loop() {
let (transport, script_tx, _state) = mock_transport();
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(16);
let (join, cancel) = spawn_loop(transport, tx);
let handle = SubscriptionHandle::new(move || cancel.cancel());
let _ = script_tx.send(MockFrame::Text(ticker_frame(
"BTC-27JUN25-60000-C",
FIXTURE_TICKER_NORMAL,
)));
let updates = drain_until(&mut rx, is_live).await;
assert!(
updates.iter().any(is_live),
"the stream is live before shutdown"
);
drop(handle);
assert!(
loop_stopped(join).await,
"dropping the SubscriptionHandle stops the reconnect loop"
);
drop(script_tx);
}
mod greeks_parity {
use chrono::{DateTime, Utc};
use deribit_websocket::prelude::Value;
use optionstratlib::chains::OptionData;
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::{Positive, ToPrimitive};
use optionstratlib::{ExpirationDate, OptionStyle};
use crate::chain::{
ContractSpecFingerprint, ExerciseStyle, GreeksOrigin, GreeksRow, GreeksSidecar,
Instrument, InstrumentKey, PremiumNumeraire, PricingInputs, ProviderId, QuoteClocks,
SettlementStyle, compute_leg_greeks,
};
const FIXTURE_TICKER: &str =
include_str!("../../tests/fixtures/deribit/ticker/ticker_normal.json");
const STRIKE: f64 = 60_000.0;
const SPOT: f64 = 60_500.0;
const AS_OF: i64 = 1_750_530_600;
const DELTA_TOLERANCE: f64 = 0.02;
const GAMMA_TOLERANCE: f64 = 0.00001;
#[track_caller]
fn pos(value: f64) -> Positive {
match Positive::new(value) {
Ok(p) => p,
Err(e) => panic!("invalid test positive `{value}`: {e}"),
}
}
#[track_caller]
fn utc(secs: i64) -> DateTime<Utc> {
match DateTime::<Utc>::from_timestamp(secs, 0) {
Some(t) => t,
None => panic!("invalid test timestamp: {secs}"),
}
}
#[track_caller]
fn pid() -> ProviderId {
match ProviderId::new("deribit") {
Ok(p) => p,
Err(e) => panic!("invalid provider id: {e}"),
}
}
#[track_caller]
fn f64_field(value: &Value, key: &str) -> f64 {
match value.get(key).and_then(Value::as_f64) {
Some(v) => v,
None => panic!("fixture ticker missing numeric field `{key}`"),
}
}
fn fixture_chain() -> OptionChain {
let mut chain = OptionChain::new("BTC", pos(SPOT), "2025-06-27".to_owned(), None, None);
let mut od = OptionData {
strike_price: pos(STRIKE),
call_bid: Some(pos(0.05)),
call_ask: Some(pos(0.06)),
implied_volatility: pos(0.4922),
..Default::default()
};
od.set_mid_prices();
let _ = chain.options.insert(od);
chain
}
#[track_caller]
fn resolved_expiry() -> DateTime<Utc> {
match fixture_chain().get_expiration() {
Some(ExpirationDate::DateTime(dt)) => dt,
other => panic!("expected an absolute-UTC chain expiry, got {other:?}"),
}
}
fn call_instrument() -> Instrument {
Instrument {
key: InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: resolved_expiry(),
strike: pos(STRIKE),
style: OptionStyle::Call,
},
provider: pid(),
native_symbol: "BTC-27JUN25-60000-C".to_owned(),
stream_symbol: None,
spec: ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: "BTC".to_owned(),
},
}
}
#[test]
fn test_deribit_local_greeks_match_venue_delta_gamma_within_tolerance() {
let ticker: Value = match FIXTURE_TICKER.parse() {
Ok(v) => v,
Err(e) => panic!("fixture ticker must parse: {e}"),
};
let mark_iv_pct = f64_field(&ticker, "mark_iv");
let venue_greeks = match ticker.get("greeks") {
Some(g) => g,
None => panic!("fixture ticker missing `greeks`"),
};
let venue_delta = f64_field(venue_greeks, "delta");
let venue_gamma = f64_field(venue_greeks, "gamma");
let venue_iv = pos(mark_iv_pct / 100.0);
let mut sink = GreeksSidecar::new();
sink.apply_venue_greeks(&GreeksRow {
instrument: call_instrument(),
iv: Some(venue_iv),
delta: None,
gamma: None,
theta: None,
vega: None,
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: utc(AS_OF),
});
let ctx = PricingInputs::new(pos(SPOT), utc(AS_OF), 1);
let chain = fixture_chain();
let clocks = crate::chain::QuoteClocks::new();
match compute_leg_greeks(&chain, &ctx, &clocks, &mut sink) {
Ok(()) => {}
Err(e) => panic!("compute_leg_greeks failed: {e}"),
}
let leg = match sink.get(&call_instrument().key) {
Some(g) => *g,
None => panic!("expected a sidecar entry for the 60000 call"),
};
assert_eq!(
leg.iv_origin,
GreeksOrigin::Provider,
"the local Greeks priced off the venue mark_iv, not a local inversion",
);
assert_eq!(leg.iv, Some(venue_iv), "the venue IV is preserved verbatim");
let local_delta = match leg.delta.and_then(|d| d.to_f64()) {
Some(d) => d,
None => panic!("expected a local delta"),
};
assert!(
(local_delta - venue_delta).abs() <= DELTA_TOLERANCE,
"local delta {local_delta} vs venue {venue_delta} (tol {DELTA_TOLERANCE})",
);
let local_gamma = match leg.gamma.and_then(|g| g.to_f64()) {
Some(g) => g,
None => panic!("expected a local gamma"),
};
assert!(
(local_gamma - venue_gamma).abs() <= GAMMA_TOLERANCE,
"local gamma {local_gamma} vs venue {venue_gamma} (tol {GAMMA_TOLERANCE})",
);
let local_theta = match leg.theta.and_then(|t| t.to_f64()) {
Some(t) => t,
None => panic!("expected a local theta"),
};
let local_vega = match leg.vega.and_then(|v| v.to_f64()) {
Some(v) => v,
None => panic!("expected a local vega"),
};
assert!(
(-200.0..-50.0).contains(&local_theta),
"local theta within the engine's magnitude band (finite, negative, \
~-125 at this snapshot; USD-denominated, no venue-coin parity): {local_theta}",
);
assert!(
(10.0..60.0).contains(&local_vega),
"local vega within the engine's magnitude band (finite, positive, \
~30.5 at this snapshot; USD-denominated, no venue-coin parity): {local_vega}",
);
}
const AS_OF_25D: i64 = 1_748_889_000;
const INVERSE_IV_TOLERANCE: f64 = 0.01;
#[test]
fn test_deribit_inverse_contract_local_iv_matches_venue_mark_iv() {
let ticker: Value = match FIXTURE_TICKER.parse() {
Ok(v) => v,
Err(e) => panic!("fixture ticker must parse: {e}"),
};
let mark_iv = f64_field(&ticker, "mark_iv") / 100.0;
let mut ctx = PricingInputs::new(pos(SPOT), utc(AS_OF_25D), 1);
ctx.premium_numeraire = PremiumNumeraire::UnderlyingCoin;
let chain = fixture_chain();
let clocks = QuoteClocks::new();
let mut sink = GreeksSidecar::new();
match compute_leg_greeks(&chain, &ctx, &clocks, &mut sink) {
Ok(()) => {}
Err(e) => panic!("compute_leg_greeks failed: {e}"),
}
let leg = match sink.get(&call_instrument().key) {
Some(g) => *g,
None => panic!("expected a sidecar entry for the 60000 call"),
};
assert_eq!(
leg.iv_origin,
GreeksOrigin::ComputedLocally,
"the IV was inverted locally (no venue seed), exercising the #83 path",
);
let local_iv = match leg.iv.map(|iv| iv.to_f64()) {
Some(iv) => iv,
None => panic!("expected a locally inverted IV for the inverse contract"),
};
assert!(
(local_iv - mark_iv).abs() <= INVERSE_IV_TOLERANCE,
"local inverse-contract IV {local_iv} vs venue mark_iv {mark_iv} \
(tol {INVERSE_IV_TOLERANCE}) — the coin premium was priced in the \
strike currency, not near-zero USD garbage",
);
assert!(
local_iv > 0.005,
"the inverted IV clears the plausibility floor: {local_iv}",
);
}
}
}