use std::collections::HashMap;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use optionstratlib::ExpirationDate;
use tokio::sync::mpsc::Receiver;
use tokio_util::sync::CancellationToken;
use dxlink::{DXLinkClient, EventType, FeedSubscription, MarketEvent};
use super::dxfeed_decode::{DxGreeksEvent, DxQuoteEvent, decode_greeks, decode_quote};
use super::{
AuthKind, ChainCapability, ChainPollCapability, GreeksCapability, MarketUpdateSink,
OptionStreamCapability, Provider, ProviderCapabilities, SendState, SubscriptionHandle,
SubscriptionRequest, UnderlyingRef,
};
use crate::chain::{ChainFetch, Instrument, MarketUpdate, ProviderId, StreamHealth};
use crate::config::{EnvSource, Secret, require_credentials};
use crate::error::ProviderError;
const DXLINK_ID: &str = "dxlink";
const CREDENTIAL_KEYS: [&str; 1] = ["token"];
const URL_VAR: &str = "CHAINVIEW_DXLINK_URL";
const DEFAULT_DXLINK_URL: &str = "wss://tasty-openapi-ws.dxfeed.com/realtime";
const FEED_CONTRACT: &str = "AUTO";
const QUOTE_EVENT: &str = "Quote";
const GREEKS_EVENT: &str = "Greeks";
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;
#[derive(Clone)]
pub(crate) struct DxlinkAdapter {
id: ProviderId,
token: Secret,
url: String,
}
impl DxlinkAdapter {
pub(crate) fn from_env(env: &dyn EnvSource) -> Result<Self, crate::error::ConfigError> {
let id = dxlink_provider_id();
let creds = require_credentials(env, &id, &CREDENTIAL_KEYS)?;
let token = creds
.get("TOKEN")
.cloned()
.ok_or_else(|| crate::error::ConfigError::MissingCredential(id.clone()))?;
let url = env
.get(URL_VAR)
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(|| DEFAULT_DXLINK_URL.to_owned());
Ok(Self { id, token, url })
}
fn client(&self) -> DXLinkClient {
DXLinkClient::new(&self.url, self.token.expose())
}
}
#[async_trait]
impl Provider for DxlinkAdapter {
fn id(&self) -> ProviderId {
self.id.clone()
}
fn capabilities(&self) -> ProviderCapabilities {
dxlink_capabilities()
}
async fn discover(&self) -> Result<Vec<UnderlyingRef>, ProviderError> {
Err(ProviderError::Unsupported("chain discovery"))
}
async fn fetch_chain(
&self,
_underlying: &str,
_expiration: &ExpirationDate,
) -> Result<ChainFetch, ProviderError> {
Err(ProviderError::Unsupported("chain assembly"))
}
async fn subscribe(
&self,
req: SubscriptionRequest,
sink: MarketUpdateSink,
) -> Result<SubscriptionHandle, ProviderError> {
let transport = LiveTransport::new(self.clone());
let id = self.id.clone();
let SubscriptionRequest {
underlying: _underlying,
expiration_utc: _expiration_utc,
instruments,
cancel,
} = req;
let loop_cancel = cancel.clone();
let handle = tokio::spawn(run_reconnect_loop(
transport,
id,
instruments,
sink,
loop_cancel,
));
Ok(SubscriptionHandle::spawned(cancel, handle))
}
}
fn dxlink_provider_id() -> ProviderId {
match ProviderId::new(DXLINK_ID) {
Ok(id) => id,
Err(_) => unreachable!("`dxlink` is a valid, reserved provider id literal"),
}
}
#[must_use]
pub(crate) fn dxlink_capabilities() -> ProviderCapabilities {
ProviderCapabilities::builder()
.chain(ChainCapability::None)
.depth(false)
.greeks(GreeksCapability::Provided)
.option_stream(OptionStreamCapability::SymbolOnly { verified: false })
.underlying_stream(false)
.chain_poll(ChainPollCapability::None)
.trades_tape(false)
.auth(AuthKind::Token)
.build()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TransportGone;
#[derive(Debug, Clone)]
enum RawDxEvent {
Quote {
symbol: String,
bid: f64,
ask: f64,
bid_size: f64,
ask_size: f64,
},
Greeks {
symbol: String,
delta: f64,
gamma: f64,
theta: f64,
vega: f64,
rho: f64,
volatility: f64,
},
Ignored,
}
#[async_trait]
trait DxlinkTransport: Send {
async fn connect_and_subscribe(&mut self, symbols: Vec<String>) -> Result<(), TransportGone>;
async fn receive(&mut self) -> Result<RawDxEvent, TransportGone>;
}
struct LiveTransport {
adapter: DxlinkAdapter,
client: Option<DXLinkClient>,
events: Option<Receiver<MarketEvent>>,
}
impl LiveTransport {
fn new(adapter: DxlinkAdapter) -> Self {
Self {
adapter,
client: None,
events: None,
}
}
}
#[async_trait]
impl DxlinkTransport for LiveTransport {
async fn connect_and_subscribe(&mut self, symbols: Vec<String>) -> Result<(), TransportGone> {
if let Some(mut old) = self.client.take() {
let _ = old.disconnect().await;
}
self.events = None;
let mut client = self.adapter.client();
let events = client.connect().await.map_err(|_| TransportGone)?;
let channel_id = client
.create_feed_channel(FEED_CONTRACT)
.await
.map_err(|_| TransportGone)?;
client
.setup_feed(channel_id, &[EventType::Quote, EventType::Greeks])
.await
.map_err(|_| TransportGone)?;
client
.subscribe(channel_id, feed_subscriptions(&symbols))
.await
.map_err(|_| TransportGone)?;
self.events = Some(events);
self.client = Some(client);
Ok(())
}
async fn receive(&mut self) -> Result<RawDxEvent, TransportGone> {
match self.events.as_mut() {
Some(events) => match events.recv().await {
Some(event) => Ok(map_market_event(event)),
None => Err(TransportGone),
},
None => Err(TransportGone),
}
}
}
fn map_market_event(event: MarketEvent) -> RawDxEvent {
match event {
MarketEvent::Quote(quote) => RawDxEvent::Quote {
symbol: quote.event_symbol,
bid: quote.bid_price,
ask: quote.ask_price,
bid_size: quote.bid_size,
ask_size: quote.ask_size,
},
MarketEvent::Greeks(greeks) => RawDxEvent::Greeks {
symbol: greeks.event_symbol,
delta: greeks.delta,
gamma: greeks.gamma,
theta: greeks.theta,
vega: greeks.vega,
rho: greeks.rho,
volatility: greeks.volatility,
},
_ => RawDxEvent::Ignored,
}
}
fn feed_subscriptions(symbols: &[String]) -> Vec<FeedSubscription> {
let cap = symbols.len().checked_mul(2).unwrap_or(symbols.len());
let mut subs = Vec::with_capacity(cap);
for symbol in symbols {
subs.push(FeedSubscription {
event_type: QUOTE_EVENT.to_owned(),
symbol: symbol.clone(),
from_time: None,
source: None,
});
subs.push(FeedSubscription {
event_type: GREEKS_EVENT.to_owned(),
symbol: symbol.clone(),
from_time: None,
source: None,
});
}
subs
}
enum StreamExit {
Reconnect,
Shutdown,
}
async fn run_reconnect_loop<T: DxlinkTransport>(
mut transport: T,
id: ProviderId,
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;
}
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,
outcome = sink.send(health) => outcome,
};
if health_sent == SendState::Closed {
return;
}
let delay = backoff_delay(attempt, sample_jitter());
tokio::select! {
biased;
() = cancel.cancelled() => return,
() = tokio::time::sleep(delay) => {}
}
}
}
async fn connect_stream_once<T: DxlinkTransport>(
transport: &mut T,
id: &ProviderId,
instruments: &[Instrument],
sink: &mut MarketUpdateSink,
cancel: &CancellationToken,
attempt: &mut u32,
) -> StreamExit {
let symbols = subscription_symbols(instruments);
let subscribed = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
result = transport.connect_and_subscribe(symbols) => result,
};
if subscribed.is_err() {
return StreamExit::Reconnect;
}
*attempt = 0;
let live = MarketUpdate::Health(id.clone(), StreamHealth::Live);
if sink.send(live).await == SendState::Closed {
return StreamExit::Shutdown;
}
let lookup = stream_lookup(instruments);
loop {
let event = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
event = transport.receive() => event,
};
let event = match event {
Ok(event) => event,
Err(_) => return StreamExit::Reconnect,
};
if route_event(&event, &lookup, sink).await == SendState::Closed {
return StreamExit::Shutdown;
}
}
}
fn subscription_symbols(instruments: &[Instrument]) -> Vec<String> {
instruments
.iter()
.map(|instrument| {
instrument
.stream_symbol
.clone()
.unwrap_or_else(|| instrument.native_symbol.clone())
})
.collect()
}
fn stream_lookup(instruments: &[Instrument]) -> HashMap<String, Instrument> {
let mut map = HashMap::new();
for instrument in instruments {
if let Some(stream) = &instrument.stream_symbol {
let _ = map
.entry(stream.clone())
.or_insert_with(|| instrument.clone());
}
let _ = map
.entry(instrument.native_symbol.clone())
.or_insert_with(|| instrument.clone());
}
map
}
async fn route_event(
event: &RawDxEvent,
lookup: &HashMap<String, Instrument>,
sink: &mut MarketUpdateSink,
) -> SendState {
let received = now_utc();
match event {
RawDxEvent::Quote {
symbol,
bid,
ask,
bid_size,
ask_size,
} => {
let Some(instrument) = lookup.get(symbol) else {
return SendState::Open;
};
let view = DxQuoteEvent {
symbol: symbol.clone(),
bid: *bid,
ask: *ask,
bid_size: *bid_size,
ask_size: *ask_size,
last: None,
event_time: None,
received_time: received,
};
match decode_quote(&view, instrument) {
Ok(quote) => sink.send(MarketUpdate::Quote(quote)).await,
Err(_) => SendState::Open,
}
}
RawDxEvent::Greeks {
symbol,
delta,
gamma,
theta,
vega,
rho,
volatility,
} => {
let Some(instrument) = lookup.get(symbol) else {
return SendState::Open;
};
let view = DxGreeksEvent {
symbol: symbol.clone(),
delta: *delta,
gamma: *gamma,
theta: *theta,
vega: *vega,
rho: *rho,
volatility: *volatility,
event_time: None,
received_time: received,
};
match decode_greeks(&view, instrument) {
Ok(greeks) => sink.send(MarketUpdate::Greeks(greeks)).await,
Err(_) => SendState::Open,
}
}
RawDxEvent::Ignored => SendState::Open,
}
}
#[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)
}
#[cfg(test)]
mod tests {
use std::collections::VecDeque;
use std::sync::Arc;
use std::sync::Mutex as StdMutex;
use optionstratlib::OptionStyle;
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::Positive;
use proptest::prelude::*;
use tokio::sync::mpsc;
use super::*;
use crate::chain::{
AliasCatalog, ChainSource, ChainStore, ContractSpecFingerprint, ExerciseStyle,
ExpirySource, GreeksOrigin, GreeksRow, InstrumentKey, MergeOutcome, QuoteUpdate,
SettlementStyle,
};
const EXP: i64 = 1_700_000_000;
const STRIKE: f64 = 60_000.0;
#[track_caller]
fn pid(id: &str) -> ProviderId {
match ProviderId::new(id) {
Ok(p) => p,
Err(e) => panic!("expected a valid provider id `{id}`, got: {e}"),
}
}
#[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}"),
}
}
struct MapEnv(HashMap<String, String>);
impl EnvSource for MapEnv {
fn get(&self, key: &str) -> Option<String> {
self.0.get(key).cloned()
}
}
fn token_env() -> MapEnv {
let mut env = HashMap::new();
let _ = env.insert(
"CHAINVIEW_DXLINK_TOKEN".to_owned(),
"do-not-log-this-token".to_owned(),
);
MapEnv(env)
}
#[track_caller]
fn sample_adapter() -> DxlinkAdapter {
match DxlinkAdapter::from_env(&token_env()) {
Ok(adapter) => adapter,
Err(e) => panic!("from_env should succeed with the token present: {e}"),
}
}
fn test_sink(
capacity: usize,
) -> (
MarketUpdateSink,
mpsc::Receiver<MarketUpdate>,
mpsc::Receiver<MarketUpdate>,
) {
let (tx_control, rx_control) = mpsc::channel::<MarketUpdate>(capacity);
let (tx_coalesced, rx_coalesced) = mpsc::channel::<MarketUpdate>(capacity);
(
MarketUpdateSink::new(tx_control, tx_coalesced),
rx_control,
rx_coalesced,
)
}
#[track_caller]
fn block<F: std::future::Future>(future: F) -> F::Output {
match tokio::runtime::Builder::new_current_thread()
.enable_time()
.build()
{
Ok(rt) => rt.block_on(future),
Err(e) => panic!("failed to build a test runtime: {e}"),
}
}
fn drain(rx: &mut mpsc::Receiver<MarketUpdate>) -> Vec<MarketUpdate> {
let mut out = Vec::new();
while let Ok(update) = rx.try_recv() {
out.push(update);
}
out
}
const DX_SYMBOL: &str = ".BTC271700000000C60000";
fn ikey(strike: f64, style: OptionStyle) -> InstrumentKey {
InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: utc(EXP),
strike: pos(strike),
style,
}
}
fn spec(multiplier: u32) -> ContractSpecFingerprint {
ContractSpecFingerprint {
contract_multiplier: multiplier,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: "BTC".to_owned(),
}
}
fn dxlink_leg(strike: f64, style: OptionStyle, multiplier: u32) -> Instrument {
Instrument {
key: ikey(strike, style),
provider: pid("dxlink"),
native_symbol: DX_SYMBOL.to_owned(),
stream_symbol: Some(DX_SYMBOL.to_owned()),
spec: spec(multiplier),
}
}
fn source_leg(strike: f64, style: OptionStyle, multiplier: u32) -> Instrument {
Instrument {
key: ikey(strike, style),
provider: pid("deribit"),
native_symbol: "BTC-27JUN25-60000-C".to_owned(),
stream_symbol: None,
spec: spec(multiplier),
}
}
#[test]
fn test_dxlink_id_is_valid_and_reserved() {
let id = dxlink_provider_id();
assert_eq!(id.as_str(), "dxlink");
assert!(id.is_reserved());
assert!(ProviderId::new(DXLINK_ID).is_ok());
}
#[test]
fn test_dxlink_capabilities_match_section_73_row() {
let caps = dxlink_capabilities();
assert_eq!(caps.chain, ChainCapability::None);
assert!(!caps.depth);
assert_eq!(caps.greeks, GreeksCapability::Provided);
assert_eq!(
caps.option_stream,
OptionStreamCapability::SymbolOnly { verified: false }
);
assert!(!caps.underlying_stream);
assert_eq!(caps.chain_poll, ChainPollCapability::None);
assert!(!caps.trades_tape);
assert_eq!(caps.auth, AuthKind::Token);
}
#[test]
fn test_adapter_reports_capabilities_and_id_via_trait() {
let adapter: Box<dyn Provider> = Box::new(sample_adapter());
assert_eq!(adapter.id().as_str(), "dxlink");
assert_eq!(adapter.capabilities().chain, ChainCapability::None);
assert_eq!(
adapter.capabilities().option_stream,
OptionStreamCapability::SymbolOnly { verified: false }
);
}
#[test]
fn test_capabilities_chain_none_makes_dxlink_unsuitable_as_source() {
assert_eq!(dxlink_capabilities().chain, ChainCapability::None);
}
#[test]
fn test_from_env_reads_chainview_namespace_only() {
let mut env = HashMap::new();
let _ = env.insert("CHAINVIEW_DXLINK_TOKEN".to_owned(), "tok-a".to_owned());
let _ = env.insert("DXLINK_TOKEN".to_owned(), "foreign".to_owned());
let adapter = match DxlinkAdapter::from_env(&MapEnv(env)) {
Ok(adapter) => adapter,
Err(e) => panic!("from_env should succeed: {e}"),
};
assert_eq!(adapter.token.expose(), "tok-a");
assert_eq!(adapter.url, DEFAULT_DXLINK_URL);
}
#[test]
fn test_from_env_url_default_and_override() {
assert_eq!(sample_adapter().url, DEFAULT_DXLINK_URL);
let mut env = HashMap::new();
let _ = env.insert("CHAINVIEW_DXLINK_TOKEN".to_owned(), "tok".to_owned());
let _ = env.insert(
"CHAINVIEW_DXLINK_URL".to_owned(),
"wss://custom.example/realtime".to_owned(),
);
match DxlinkAdapter::from_env(&MapEnv(env)) {
Ok(adapter) => assert_eq!(adapter.url, "wss://custom.example/realtime"),
Err(e) => panic!("from_env should succeed: {e}"),
}
}
#[test]
fn test_from_env_missing_token_is_error() {
match DxlinkAdapter::from_env(&MapEnv(HashMap::new())) {
Err(crate::error::ConfigError::MissingCredential(id)) => {
assert_eq!(id.as_str(), "dxlink");
}
Err(other) => panic!("expected MissingCredential, got: {other}"),
Ok(_) => panic!("expected MissingCredential, got Ok"),
}
}
#[test]
fn test_token_never_appears_in_debug_of_secret() {
let adapter = sample_adapter();
let rendered = format!("{:?}", adapter.token);
assert!(!rendered.contains("do-not-log-this-token"));
assert!(rendered.contains("redacted"));
}
#[test]
fn test_discover_is_unsupported() {
match block(sample_adapter().discover()) {
Err(ProviderError::Unsupported(what)) => assert_eq!(what, "chain discovery"),
other => panic!("expected Unsupported(chain discovery), got {other:?}"),
}
}
#[test]
fn test_fetch_chain_is_unsupported() {
let adapter = sample_adapter();
let result = block(adapter.fetch_chain("BTC", &ExpirationDate::Days(pos(30.0))));
match result {
Err(ProviderError::Unsupported(what)) => assert_eq!(what, "chain assembly"),
other => panic!("expected Unsupported(chain assembly), got {other:?}"),
}
}
#[test]
fn test_map_market_event_quote_carries_f64_sizes() {
let ev = map_market_event(MarketEvent::Quote(dxlink::events::QuoteEvent {
event_type: "Quote".to_owned(),
event_symbol: DX_SYMBOL.to_owned(),
bid_price: 1.5,
ask_price: 1.7,
bid_size: 12.5,
ask_size: 8.0,
}));
match ev {
RawDxEvent::Quote {
symbol,
bid,
ask,
bid_size,
ask_size,
} => {
assert_eq!(symbol, DX_SYMBOL);
assert_eq!(bid, 1.5);
assert_eq!(ask, 1.7);
assert_eq!(bid_size, 12.5);
assert_eq!(ask_size, 8.0);
}
other => panic!("expected a Quote, got {other:?}"),
}
}
#[test]
fn test_map_market_event_trade_is_ignored() {
let ev = map_market_event(MarketEvent::Trade(dxlink::events::TradeEvent {
event_type: "Trade".to_owned(),
event_symbol: DX_SYMBOL.to_owned(),
price: 1.6,
size: 3.0,
day_volume: 100.0,
}));
assert!(matches!(ev, RawDxEvent::Ignored));
}
#[test]
fn test_map_market_event_unsubscribed_kind_is_ignored() {
let ev = map_market_event(MarketEvent::TradeETH(dxlink::events::TradeETHEvent {
event_type: "TradeETH".to_owned(),
event_symbol: DX_SYMBOL.to_owned(),
event_time: 0,
time: 0,
time_nano_part: 0,
sequence: 0,
exchange_code: "Q".to_owned(),
price: 1.6,
change: 0.0,
size: 3.0,
day_id: 20_665,
day_volume: 100.0,
day_turnover: 160.0,
tick_direction: "UNDEFINED".to_owned(),
extended_trading_hours: true,
}));
assert!(matches!(ev, RawDxEvent::Ignored));
}
#[test]
fn test_stream_lookup_resolves_stream_and_native_symbol() {
let legs = vec![dxlink_leg(STRIKE, OptionStyle::Call, 1)];
let lookup = stream_lookup(&legs);
match lookup.get(DX_SYMBOL) {
Some(instrument) => assert_eq!(instrument.key, ikey(STRIKE, OptionStyle::Call)),
None => panic!("the dxfeed symbol should resolve to the leg"),
}
assert!(lookup.contains_key(DX_SYMBOL));
}
#[test]
fn test_subscription_symbols_uses_stream_symbol() {
let legs = vec![dxlink_leg(STRIKE, OptionStyle::Call, 1)];
assert_eq!(subscription_symbols(&legs), vec![DX_SYMBOL.to_owned()]);
}
#[track_caller]
fn route(event: &RawDxEvent) -> Vec<MarketUpdate> {
let legs = vec![dxlink_leg(STRIKE, OptionStyle::Call, 1)];
let lookup = stream_lookup(&legs);
let (mut sink, mut rx_control, mut rx_coalesced) = test_sink(8);
let sent = block(route_event(event, &lookup, &mut sink));
assert_ne!(sent, SendState::Closed);
let mut out = drain(&mut rx_control);
out.extend(drain(&mut rx_coalesced));
out
}
#[test]
fn test_route_quote_decodes_to_quote_update_with_absent_event_time() {
let updates = route(&RawDxEvent::Quote {
symbol: DX_SYMBOL.to_owned(),
bid: 1.5,
ask: 1.7,
bid_size: 10.0,
ask_size: 20.0,
});
match updates.as_slice() {
[MarketUpdate::Quote(quote)] => {
assert_eq!(quote.instrument.key, ikey(STRIKE, OptionStyle::Call));
assert_eq!(quote.instrument.provider.as_str(), "dxlink");
assert_eq!(quote.bid, Some(pos(1.5)));
assert_eq!(quote.ask, Some(pos(1.7)));
assert_eq!(quote.bid_size, Some(pos(10.0)));
assert!(quote.event_time.is_none());
}
other => panic!("expected a single QuoteUpdate, got {other:?}"),
}
}
#[test]
fn test_route_greeks_decodes_to_provider_greeks_row_iv_as_is() {
let updates = route(&RawDxEvent::Greeks {
symbol: DX_SYMBOL.to_owned(),
delta: 0.55,
gamma: 0.01,
theta: -0.05,
vega: 0.2,
rho: 0.03,
volatility: 0.35,
});
match updates.as_slice() {
[MarketUpdate::Greeks(greeks)] => {
assert_eq!(greeks.origin, GreeksOrigin::Provider);
assert_eq!(greeks.iv, Some(pos(0.35)));
assert!(greeks.event_time.is_none());
}
other => panic!("expected a single GreeksRow, got {other:?}"),
}
}
#[test]
fn test_route_unknown_symbol_is_benign_drop() {
let updates = route(&RawDxEvent::Quote {
symbol: ".UNSUBSCRIBED".to_owned(),
bid: 1.5,
ask: 1.7,
bid_size: 10.0,
ask_size: 20.0,
});
assert!(
updates.is_empty(),
"an unknown-symbol event never emits an update: {updates:?}"
);
}
#[test]
fn test_route_crossed_quote_is_benign_drop() {
let updates = route(&RawDxEvent::Quote {
symbol: DX_SYMBOL.to_owned(),
bid: 2.0,
ask: 1.0,
bid_size: 10.0,
ask_size: 20.0,
});
assert!(
updates.is_empty(),
"a crossed quote is a benign per-tick drop: {updates:?}"
);
}
#[test]
fn test_route_ignored_event_emits_nothing() {
assert!(route(&RawDxEvent::Ignored).is_empty());
}
struct MockTransport {
connections: VecDeque<VecDeque<RawDxEvent>>,
current: Option<VecDeque<RawDxEvent>>,
subscribed: Arc<StdMutex<Vec<Vec<String>>>>,
}
impl MockTransport {
fn new(connections: Vec<Vec<RawDxEvent>>) -> Self {
Self {
connections: connections.into_iter().map(VecDeque::from).collect(),
current: None,
subscribed: Arc::new(StdMutex::new(Vec::new())),
}
}
}
#[async_trait]
impl DxlinkTransport for MockTransport {
async fn connect_and_subscribe(
&mut self,
symbols: Vec<String>,
) -> Result<(), TransportGone> {
if let Ok(mut log) = self.subscribed.lock() {
log.push(symbols);
}
match self.connections.pop_front() {
Some(events) => {
self.current = Some(events);
Ok(())
}
None => Err(TransportGone),
}
}
async fn receive(&mut self) -> Result<RawDxEvent, TransportGone> {
match self.current.as_mut().and_then(VecDeque::pop_front) {
Some(event) => Ok(event),
None => {
self.current = None;
Err(TransportGone)
}
}
}
}
#[test]
fn test_subscribe_streams_health_and_decoded_updates() {
block(async {
let (sink, mut rx_control, mut rx_coalesced) = test_sink(32);
let transport = MockTransport::new(vec![vec![
RawDxEvent::Quote {
symbol: DX_SYMBOL.to_owned(),
bid: 1.5,
ask: 1.7,
bid_size: 10.0,
ask_size: 20.0,
},
RawDxEvent::Greeks {
symbol: DX_SYMBOL.to_owned(),
delta: 0.55,
gamma: 0.01,
theta: -0.05,
vega: 0.2,
rho: 0.03,
volatility: 0.35,
},
]]);
let subscribed = Arc::clone(&transport.subscribed);
let cancel = CancellationToken::new();
let handle = tokio::spawn(run_reconnect_loop(
transport,
pid("dxlink"),
vec![dxlink_leg(STRIKE, OptionStyle::Call, 1)],
sink,
cancel.clone(),
));
tokio::time::sleep(Duration::from_millis(30)).await;
cancel.cancel();
let _ = handle.await;
let control = drain(&mut rx_control);
let coalesced = drain(&mut rx_coalesced);
match subscribed.lock() {
Ok(log) => {
assert!(
log.iter().any(|set| set == &vec![DX_SYMBOL.to_owned()]),
"the leg's dxfeed symbol was subscribed: {log:?}"
);
}
Err(_) => panic!("subscribed log lock poisoned"),
}
assert!(
control
.iter()
.any(|u| matches!(u, MarketUpdate::Health(_, StreamHealth::Live))),
"Health(Live) is emitted on connect: {control:?}"
);
assert!(
coalesced
.iter()
.any(|u| matches!(u, MarketUpdate::Quote(_))),
"a QuoteUpdate is emitted: {coalesced:?}"
);
assert!(
coalesced
.iter()
.any(|u| matches!(u, MarketUpdate::Greeks(_))),
"a GreeksRow is emitted: {coalesced:?}"
);
});
}
#[test]
fn test_subscribe_reconnect_lifecycle_emits_reconnecting() {
block(async {
let (sink, mut rx_control, _rx_coalesced) = test_sink(32);
let transport = MockTransport::new(vec![
vec![RawDxEvent::Quote {
symbol: DX_SYMBOL.to_owned(),
bid: 1.5,
ask: 1.7,
bid_size: 10.0,
ask_size: 20.0,
}],
vec![],
]);
let cancel = CancellationToken::new();
let handle = tokio::spawn(run_reconnect_loop(
transport,
pid("dxlink"),
vec![dxlink_leg(STRIKE, OptionStyle::Call, 1)],
sink,
cancel.clone(),
));
tokio::time::sleep(Duration::from_millis(30)).await;
cancel.cancel();
let _ = handle.await;
let control = drain(&mut rx_control);
assert!(
control.iter().any(|u| matches!(
u,
MarketUpdate::Health(_, StreamHealth::Reconnecting { .. })
)),
"a dropped stream emits Health(Reconnecting): {control:?}"
);
});
}
#[test]
fn test_subscribe_cancel_stops_loop_without_connecting_again() {
block(async {
let (sink, _rx_control, _rx_coalesced) = test_sink(8);
let transport = MockTransport::new(vec![vec![]]);
let subscribed = Arc::clone(&transport.subscribed);
let cancel = CancellationToken::new();
cancel.cancel();
let handle = tokio::spawn(run_reconnect_loop(
transport,
pid("dxlink"),
vec![dxlink_leg(STRIKE, OptionStyle::Call, 1)],
sink,
cancel,
));
let _ = handle.await;
match subscribed.lock() {
Ok(log) => assert!(
log.is_empty(),
"a pre-cancelled loop never connects: {log:?}"
),
Err(_) => panic!("subscribed log lock poisoned"),
}
});
}
fn overlay_quote(bid: f64, ask: f64) -> QuoteUpdate {
QuoteUpdate {
instrument: dxlink_leg(STRIKE, OptionStyle::Call, 1),
bid: Some(pos(bid)),
ask: Some(pos(ask)),
last: None,
bid_size: None,
ask_size: None,
event_time: None,
received_time: utc(EXP + 10),
}
}
fn overlay_store(overlay_mult: u32) -> ChainStore {
let mut catalog = AliasCatalog::new();
catalog.insert(source_leg(STRIKE, OptionStyle::Call, 1));
catalog.insert(dxlink_leg(STRIKE, OptionStyle::Call, overlay_mult));
catalog.insert(source_leg(STRIKE + 1_000.0, OptionStyle::Call, 1));
let mut chain = OptionChain::new("BTC", pos(STRIKE), utc(EXP).to_rfc3339(), None, None);
chain.add_option(
pos(STRIKE),
Some(pos(1.0)),
Some(pos(1.2)),
Some(pos(2.0)),
Some(pos(2.4)),
Positive::ZERO,
None,
None,
None,
None,
None,
None,
);
chain.add_option(
pos(STRIKE + 1_000.0),
Some(pos(3.0)),
Some(pos(3.2)),
None,
None,
Positive::ZERO,
None,
None,
None,
None,
None,
None,
);
let fetch = ChainFetch::new(
chain,
ExpirySource::new("BTC", utc(EXP), pid("deribit")),
catalog,
);
ChainStore::seed(fetch, ChainSource::Merged, Duration::from_secs(2), utc(EXP))
}
#[track_caller]
fn call_bid(store: &ChainStore, strike: f64) -> Option<Positive> {
store
.chain()
.options
.iter()
.find(|o| o.strike_price == pos(strike))
.and_then(|o| o.call_bid)
}
#[test]
fn test_matching_fingerprint_overlay_wins() {
let mut store = overlay_store(1);
let outcome = store.apply_quote(&overlay_quote(5.0, 5.2));
assert_eq!(outcome, MergeOutcome::Applied);
assert!(!store.is_overlay_refused(&ikey(STRIKE, OptionStyle::Call)));
assert_eq!(call_bid(&store, STRIKE), Some(pos(5.0)));
}
#[test]
fn test_mismatched_fingerprint_refused_source_kept_others_unaffected() {
let mut store = overlay_store(100);
let outcome = store.apply_quote(&overlay_quote(5.0, 5.2));
assert_eq!(outcome, MergeOutcome::OverlayRefused);
assert!(store.is_overlay_refused(&ikey(STRIKE, OptionStyle::Call)));
assert_eq!(call_bid(&store, STRIKE), Some(pos(1.0)));
assert_eq!(call_bid(&store, STRIKE + 1_000.0), Some(pos(3.0)));
}
#[test]
fn test_mismatched_greeks_overlay_also_refused() {
let mut store = overlay_store(100);
let row = GreeksRow {
instrument: dxlink_leg(STRIKE, OptionStyle::Call, 1),
iv: Some(pos(0.35)),
delta: None,
gamma: None,
theta: None,
vega: None,
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: utc(EXP + 10),
};
assert_eq!(store.apply_greeks(&row), MergeOutcome::OverlayRefused);
assert!(store.is_overlay_refused(&ikey(STRIKE, OptionStyle::Call)));
}
proptest! {
#![proptest_config(ProptestConfig { cases: 256, ..ProptestConfig::default() })]
#[test]
fn prop_overlay_spec_gate(overlay_mult in 1u32..250) {
let mut store = overlay_store(overlay_mult);
let outcome = store.apply_quote(&overlay_quote(5.0, 5.2));
let after = call_bid(&store, STRIKE);
if overlay_mult == 1 {
prop_assert_eq!(outcome, MergeOutcome::Applied);
prop_assert_eq!(after, Some(pos(5.0)));
} else {
prop_assert_eq!(outcome, MergeOutcome::OverlayRefused);
prop_assert_eq!(after, Some(pos(1.0)));
}
}
}
#[test]
fn test_backoff_delay_grows_and_caps() {
let zero = backoff_delay(0, 0.0);
assert_eq!(zero, Duration::from_millis(250));
assert!(backoff_delay(4, 0.0) > backoff_delay(1, 0.0));
assert_eq!(backoff_delay(30, 0.0), Duration::from_millis(30_000));
}
const FIXTURE_QUOTE_SYMBOL: &str =
include_str!("../../tests/fixtures/dxlink/quote/quote_symbol.json");
const FIXTURE_GREEKS_SYMBOL: &str =
include_str!("../../tests/fixtures/dxlink/greeks/greeks_symbol.json");
const FIXTURE_OVERLAY_MATCHING: &str =
include_str!("../../tests/fixtures/dxlink/overlay/overlay_matching.json");
const FIXTURE_OVERLAY_MISMATCHED: &str =
include_str!("../../tests/fixtures/dxlink/overlay/overlay_mismatched.json");
const FIXTURE_DERIBIT_INSTRUMENTS: &str =
include_str!("../../tests/fixtures/deribit/instruments/instruments_btc.json");
const SOURCE_OPTION_NAME: &str = "BTC-27JUN25-60000-C";
const OVERLAY_STREAM_SYMBOL: &str = ".BTC250627C60000";
#[derive(serde::Deserialize)]
struct FixtureInstrument {
instrument_name: String,
kind: String,
strike: Option<f64>,
expiration_timestamp: Option<i64>,
option_type: Option<String>,
quote_currency: Option<String>,
}
#[derive(serde::Deserialize)]
struct OverlayFixtureSpec {
contract_multiplier: u32,
settlement: String,
exercise: String,
quote_currency: String,
venue_product_code: String,
}
#[track_caller]
fn settlement_of(label: &str) -> SettlementStyle {
match label {
"cash" => SettlementStyle::Cash,
"physical" => SettlementStyle::Physical,
other => panic!("unknown settlement label in fixture: {other}"),
}
}
#[track_caller]
fn exercise_of(label: &str) -> ExerciseStyle {
match label {
"european" => ExerciseStyle::European,
"american" => ExerciseStyle::American,
other => panic!("unknown exercise label in fixture: {other}"),
}
}
#[track_caller]
fn overlay_spec_from(json: &str) -> ContractSpecFingerprint {
let dto: OverlayFixtureSpec = match serde_json::from_str(json) {
Ok(dto) => dto,
Err(e) => panic!("overlay-spec fixture must deserialize: {e}"),
};
ContractSpecFingerprint {
contract_multiplier: dto.contract_multiplier,
settlement: settlement_of(&dto.settlement),
exercise: exercise_of(&dto.exercise),
quote_currency: dto.quote_currency,
venue_product_code: dto.venue_product_code,
}
}
#[track_caller]
fn deribit_source_option() -> FixtureInstrument {
let instruments: Vec<FixtureInstrument> =
match serde_json::from_str(FIXTURE_DERIBIT_INSTRUMENTS) {
Ok(list) => list,
Err(e) => panic!("deribit instruments fixture must deserialize: {e}"),
};
match instruments
.into_iter()
.find(|i| i.instrument_name == SOURCE_OPTION_NAME && i.kind == "option")
{
Some(option) => option,
None => panic!("the deribit fixture must carry {SOURCE_OPTION_NAME}"),
}
}
#[track_caller]
fn source_key(option: &FixtureInstrument) -> InstrumentKey {
let millis = match option.expiration_timestamp {
Some(ms) => ms,
None => panic!("the source option fixture carries an expiration_timestamp"),
};
let expiration_utc = match DateTime::<Utc>::from_timestamp_millis(millis) {
Some(dt) => dt,
None => panic!("the fixture expiration_timestamp is a valid instant: {millis}"),
};
let strike = match option.strike {
Some(s) => pos(s),
None => panic!("the source option fixture carries a strike"),
};
let style = match option.option_type.as_deref() {
Some("call") => OptionStyle::Call,
Some("put") => OptionStyle::Put,
other => panic!("unexpected source option_type: {other:?}"),
};
InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc,
strike,
style,
}
}
fn source_spec(option: &FixtureInstrument) -> ContractSpecFingerprint {
ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: option
.quote_currency
.clone()
.unwrap_or_else(|| "USD".to_owned()),
venue_product_code: "BTC".to_owned(),
}
}
#[track_caller]
fn overlay_pair_store(overlay_spec: ContractSpecFingerprint) -> (ChainStore, Instrument) {
let option = deribit_source_option();
let key = source_key(&option);
let source = Instrument {
key: key.clone(),
provider: pid("deribit"),
native_symbol: SOURCE_OPTION_NAME.to_owned(),
stream_symbol: None,
spec: source_spec(&option),
};
let overlay = Instrument {
key: key.clone(),
provider: pid("dxlink"),
native_symbol: OVERLAY_STREAM_SYMBOL.to_owned(),
stream_symbol: Some(OVERLAY_STREAM_SYMBOL.to_owned()),
spec: overlay_spec,
};
let mut catalog = AliasCatalog::new();
catalog.insert(source);
catalog.insert(overlay.clone());
let expiry = key.expiration_utc;
let mut chain = OptionChain::new("BTC", key.strike, expiry.to_rfc3339(), None, None);
chain.add_option(
key.strike,
Some(pos(0.05)),
Some(pos(0.06)),
Some(pos(0.05)),
Some(pos(0.06)),
Positive::ZERO,
None,
None,
None,
None,
None,
None,
);
let fetch = ChainFetch::new(
chain,
ExpirySource::new("BTC", expiry, pid("deribit")),
catalog,
);
let store = ChainStore::seed(fetch, ChainSource::Merged, Duration::from_secs(2), expiry);
(store, overlay)
}
#[track_caller]
fn decode_fixture_quote(overlay: &Instrument) -> QuoteUpdate {
let event: dxlink::events::QuoteEvent = match serde_json::from_str(FIXTURE_QUOTE_SYMBOL) {
Ok(ev) => ev,
Err(e) => panic!("quote fixture must deserialize as dxlink QuoteEvent: {e}"),
};
let raw = map_market_event(MarketEvent::Quote(event));
let lookup = stream_lookup(std::slice::from_ref(overlay));
let (mut sink, mut rx_control, mut rx_coalesced) = test_sink(8);
let sent = block(route_event(&raw, &lookup, &mut sink));
assert_ne!(sent, SendState::Closed);
let mut out = drain(&mut rx_control);
out.extend(drain(&mut rx_coalesced));
match out.as_slice() {
[MarketUpdate::Quote(quote)] => quote.clone(),
other => panic!("expected a single decoded QuoteUpdate, got {other:?}"),
}
}
#[test]
fn test_dxlink_fixture_quote_symbol_normalizes_to_quote_update() {
let (_store, overlay) = overlay_pair_store(overlay_spec_from(FIXTURE_OVERLAY_MATCHING));
let quote = decode_fixture_quote(&overlay);
assert_eq!(quote.instrument.provider.as_str(), "dxlink");
assert_eq!(quote.instrument.key, overlay.key);
assert_eq!(quote.bid, Some(pos(0.062)));
assert_eq!(quote.ask, Some(pos(0.064)));
assert_eq!(quote.bid_size, Some(pos(12.0)));
assert_eq!(quote.ask_size, Some(pos(8.0)));
assert!(quote.event_time.is_none());
}
#[test]
fn test_dxlink_fixture_greeks_symbol_normalizes_to_greeks_row() {
let (_store, overlay) = overlay_pair_store(overlay_spec_from(FIXTURE_OVERLAY_MATCHING));
let event: dxlink::events::GreeksEvent = match serde_json::from_str(FIXTURE_GREEKS_SYMBOL) {
Ok(ev) => ev,
Err(e) => panic!("greeks fixture must deserialize as dxlink GreeksEvent: {e}"),
};
let raw = map_market_event(MarketEvent::Greeks(event));
let lookup = stream_lookup(std::slice::from_ref(&overlay));
let (mut sink, mut rx_control, mut rx_coalesced) = test_sink(8);
let _ = block(route_event(&raw, &lookup, &mut sink));
let mut out = drain(&mut rx_control);
out.extend(drain(&mut rx_coalesced));
match out.as_slice() {
[MarketUpdate::Greeks(greeks)] => {
assert_eq!(greeks.instrument.provider.as_str(), "dxlink");
assert_eq!(greeks.origin, GreeksOrigin::Provider);
assert_eq!(greeks.iv, Some(pos(0.4922)));
assert!(greeks.event_time.is_none());
}
other => panic!("expected a single decoded GreeksRow, got {other:?}"),
}
}
#[track_caller]
fn overlay_call_bid(store: &ChainStore, strike: Positive) -> Option<Positive> {
store
.chain()
.options
.iter()
.find(|o| o.strike_price == strike)
.and_then(|o| o.call_bid)
}
#[test]
fn test_overlay_pair_leg_selection_subscribes_dxlink_stream_legs() {
let (_store, overlay) = overlay_pair_store(overlay_spec_from(FIXTURE_OVERLAY_MATCHING));
assert_eq!(
subscription_symbols(std::slice::from_ref(&overlay)),
vec![OVERLAY_STREAM_SYMBOL.to_owned()],
"the overlay subscribes the dxfeed stream symbol, not the deribit native"
);
assert_ne!(overlay.native_symbol, SOURCE_OPTION_NAME);
}
#[test]
fn test_overlay_pair_matching_fingerprint_merges_overlay_wins() {
let (mut store, overlay) = overlay_pair_store(overlay_spec_from(FIXTURE_OVERLAY_MATCHING));
let key = overlay.key.clone();
match store.aliases().resolve_symbol(OVERLAY_STREAM_SYMBOL) {
Some(resolved) => assert_eq!(resolved, &key),
None => panic!("the overlay dxfeed symbol must resolve to the shared source key"),
}
let quote = decode_fixture_quote(&overlay);
assert_eq!(store.apply_quote("e), MergeOutcome::Applied);
assert!(!store.is_overlay_refused(&key));
assert_eq!(overlay_call_bid(&store, key.strike), Some(pos(0.062)));
}
#[test]
fn test_overlay_pair_mismatched_fingerprint_refused_source_kept() {
let (mut store, overlay) =
overlay_pair_store(overlay_spec_from(FIXTURE_OVERLAY_MISMATCHED));
let key = overlay.key.clone();
let quote = decode_fixture_quote(&overlay);
assert_eq!(store.apply_quote("e), MergeOutcome::OverlayRefused);
assert!(store.is_overlay_refused(&key));
assert_eq!(overlay_call_bid(&store, key.strike), Some(pos(0.05)));
}
}