use std::time::Duration;
use chrono::{DateTime, Utc};
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::{Decimal, Positive};
use super::fetch::AliasCatalog;
use super::identity::{Instrument, ProviderId};
#[derive(Debug, Clone)]
pub struct QuoteUpdate {
pub instrument: Instrument,
pub bid: Option<Positive>,
pub ask: Option<Positive>,
pub last: Option<Positive>,
pub bid_size: Option<Positive>,
pub ask_size: Option<Positive>,
pub event_time: Option<DateTime<Utc>>,
pub received_time: DateTime<Utc>,
}
#[derive(Debug, Clone)]
pub struct GreeksRow {
pub instrument: Instrument,
pub iv: Option<Positive>,
pub delta: Option<Decimal>,
pub gamma: Option<Decimal>,
pub theta: Option<Decimal>,
pub vega: Option<Decimal>,
pub rho: Option<Decimal>,
pub origin: GreeksOrigin,
pub event_time: Option<DateTime<Utc>>,
pub received_time: DateTime<Utc>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum GreeksOrigin {
Provider,
ComputedLocally,
}
#[derive(Debug, Clone)]
pub struct DepthLadder {
pub instrument: Instrument,
pub bids: Vec<DepthLevel>,
pub asks: Vec<DepthLevel>,
pub event_time: Option<DateTime<Utc>>,
pub received_time: DateTime<Utc>,
pub change_id: Option<u64>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct DepthLevel {
pub price: Positive,
pub size: Positive,
}
#[derive(Debug, Clone)]
pub enum MarketUpdate {
Quote(QuoteUpdate),
Greeks(GreeksRow),
Depth(DepthLadder),
Chain(ChainSnapshot),
Health(ProviderId, StreamHealth),
}
#[derive(Debug, Clone)]
pub struct ChainSnapshot {
pub chain_key: (ProviderId, String, DateTime<Utc>),
pub chain: OptionChain,
pub aliases: AliasCatalog,
pub source: ChainSource,
pub health: StreamHealth,
pub last_full_poll: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum ChainSource {
Poll,
Stream,
Merged,
}
#[derive(Debug, Clone)]
pub enum StreamHealth {
Live,
Stale {
since: DateTime<Utc>,
},
Reconnecting {
attempt: u32,
},
}
pub const QUOTE_STALE_AFTER: Duration = Duration::from_secs(5);
pub const GREEKS_STALE_AFTER: Duration = Duration::from_secs(10);
pub const CHAIN_STALE_SLACK: Duration = Duration::from_secs(2);
pub const FEED_DELAY_WARN: Duration = Duration::from_secs(2);
pub const DIRECTION_DECAY: Duration = Duration::from_secs(3);
#[must_use]
pub fn chain_stale_after(refresh_interval: Duration) -> Duration {
refresh_interval
.checked_add(CHAIN_STALE_SLACK)
.unwrap_or(Duration::MAX)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::chain::identity::{
ContractSpecFingerprint, ExerciseStyle, InstrumentKey, SettlementStyle,
};
use optionstratlib::OptionStyle;
#[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 utc(secs: i64) -> DateTime<Utc> {
match DateTime::<Utc>::from_timestamp(secs, 0) {
Some(t) => t,
None => panic!("invalid test timestamp: {secs}"),
}
}
#[track_caller]
fn pos(value: f64) -> Positive {
match Positive::new(value) {
Ok(p) => p,
Err(e) => panic!("invalid test positive `{value}`: {e}"),
}
}
fn sample_instrument() -> Instrument {
Instrument {
key: InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: utc(1_700_000_000),
strike: pos(60_000.0),
style: OptionStyle::Call,
},
provider: pid("deribit"),
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 absent_quote() -> QuoteUpdate {
QuoteUpdate {
instrument: sample_instrument(),
bid: None,
ask: None,
last: None,
bid_size: None,
ask_size: None,
event_time: None,
received_time: utc(1_700_000_100),
}
}
fn absent_greeks() -> GreeksRow {
GreeksRow {
instrument: sample_instrument(),
iv: None,
delta: None,
gamma: None,
theta: None,
vega: None,
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: utc(1_700_000_100),
}
}
#[test]
fn test_quote_update_absent_fields_stay_none() {
let q = absent_quote();
assert!(q.bid.is_none());
assert!(q.ask.is_none());
assert!(q.last.is_none());
assert!(q.bid_size.is_none());
assert!(q.ask_size.is_none());
assert!(q.event_time.is_none());
assert_eq!(q.received_time, utc(1_700_000_100));
}
#[test]
fn test_quote_update_present_fields_carry_values() {
let q = QuoteUpdate {
bid: Some(pos(1.5)),
ask: Some(pos(1.7)),
event_time: Some(utc(1_700_000_099)),
..absent_quote()
};
assert_eq!(q.bid, Some(pos(1.5)));
assert_eq!(q.ask, Some(pos(1.7)));
assert_eq!(q.event_time, Some(utc(1_700_000_099)));
}
#[test]
fn test_greeks_row_absent_greeks_stay_none() {
let g = absent_greeks();
assert!(g.iv.is_none());
assert!(g.delta.is_none());
assert!(g.gamma.is_none());
assert!(g.theta.is_none());
assert!(g.vega.is_none());
assert!(g.rho.is_none());
assert!(g.event_time.is_none());
assert_eq!(g.received_time, utc(1_700_000_100));
}
#[test]
fn test_greeks_origin_computed_locally_is_representable() {
let g = GreeksRow {
origin: GreeksOrigin::ComputedLocally,
..absent_greeks()
};
assert_eq!(g.origin, GreeksOrigin::ComputedLocally);
}
#[test]
fn test_greeks_origin_provider_ne_computed_locally() {
assert_ne!(GreeksOrigin::Provider, GreeksOrigin::ComputedLocally);
}
#[test]
fn test_depth_ladder_change_id_carried_when_present() {
let ladder = DepthLadder {
instrument: sample_instrument(),
bids: vec![DepthLevel {
price: pos(60_000.0),
size: pos(2.0),
}],
asks: vec![DepthLevel {
price: pos(60_010.0),
size: pos(1.0),
}],
event_time: Some(utc(1_700_000_099)),
received_time: utc(1_700_000_100),
change_id: Some(42),
};
assert_eq!(ladder.change_id, Some(42u64));
match (ladder.bids.first(), ladder.asks.first()) {
(Some(bid), Some(ask)) => {
assert_eq!(bid.price, pos(60_000.0));
assert_eq!(ask.price, pos(60_010.0));
}
_ => panic!("expected one bid and one ask level"),
}
}
#[test]
fn test_depth_ladder_change_id_none_when_absent() {
let ladder = DepthLadder {
instrument: sample_instrument(),
bids: Vec::new(),
asks: Vec::new(),
event_time: None,
received_time: utc(1_700_000_100),
change_id: None,
};
assert!(ladder.change_id.is_none());
}
#[test]
fn test_market_update_quote_variant_constructs() {
let update = MarketUpdate::Quote(absent_quote());
match update {
MarketUpdate::Quote(q) => assert!(q.bid.is_none()),
other => panic!("expected Quote, got {other:?}"),
}
}
#[test]
fn test_market_update_greeks_variant_constructs() {
let delta = Decimal::new(-25, 2);
let update = MarketUpdate::Greeks(GreeksRow {
delta: Some(delta),
origin: GreeksOrigin::ComputedLocally,
..absent_greeks()
});
match update {
MarketUpdate::Greeks(g) => {
assert_eq!(g.delta, Some(delta));
assert_eq!(g.origin, GreeksOrigin::ComputedLocally);
}
other => panic!("expected Greeks, got {other:?}"),
}
}
#[test]
fn test_market_update_depth_variant_constructs() {
let update = MarketUpdate::Depth(DepthLadder {
instrument: sample_instrument(),
bids: Vec::new(),
asks: Vec::new(),
event_time: None,
received_time: utc(1_700_000_100),
change_id: Some(7),
});
match update {
MarketUpdate::Depth(d) => assert_eq!(d.change_id, Some(7u64)),
other => panic!("expected Depth, got {other:?}"),
}
}
#[test]
fn test_market_update_chain_variant_constructs() {
let snapshot = ChainSnapshot {
chain_key: (pid("deribit"), "BTC".to_owned(), utc(1_700_000_000)),
chain: OptionChain::new("BTC", pos(60_000.0), "2025-06-27".to_owned(), None, None),
aliases: AliasCatalog::new(),
source: ChainSource::Poll,
health: StreamHealth::Live,
last_full_poll: Some(utc(1_700_000_100)),
};
let update = MarketUpdate::Chain(snapshot);
match update {
MarketUpdate::Chain(c) => {
assert_eq!(c.chain.symbol, "BTC");
assert_eq!(c.chain_key.1, "BTC");
assert_eq!(c.source, ChainSource::Poll);
assert!(c.aliases.is_empty());
}
other => panic!("expected Chain, got {other:?}"),
}
}
#[test]
fn test_market_update_health_variant_constructs() {
let update =
MarketUpdate::Health(pid("deribit"), StreamHealth::Reconnecting { attempt: 3 });
match update {
MarketUpdate::Health(provider, StreamHealth::Reconnecting { attempt }) => {
assert_eq!(provider.as_str(), "deribit");
assert_eq!(attempt, 3);
}
other => panic!("expected Health(_, Reconnecting), got {other:?}"),
}
}
#[test]
fn test_stream_health_stale_carries_since_instant() {
let health = StreamHealth::Stale {
since: utc(1_700_000_050),
};
match health {
StreamHealth::Stale { since } => assert_eq!(since, utc(1_700_000_050)),
other => panic!("expected Stale, got {other:?}"),
}
}
#[test]
fn test_quote_stale_after_equals_five_seconds() {
assert_eq!(QUOTE_STALE_AFTER, Duration::from_secs(5));
}
#[test]
fn test_greeks_stale_after_equals_ten_seconds() {
assert_eq!(GREEKS_STALE_AFTER, Duration::from_secs(10));
}
#[test]
fn test_chain_stale_slack_equals_two_seconds() {
assert_eq!(CHAIN_STALE_SLACK, Duration::from_secs(2));
}
#[test]
fn test_feed_delay_warn_equals_two_seconds() {
assert_eq!(FEED_DELAY_WARN, Duration::from_secs(2));
}
#[test]
fn test_direction_decay_equals_three_seconds() {
assert_eq!(DIRECTION_DECAY, Duration::from_secs(3));
}
#[test]
fn test_chain_stale_after_adds_slack_to_refresh() {
let refresh = Duration::from_secs(2);
assert_eq!(chain_stale_after(refresh), refresh + CHAIN_STALE_SLACK);
}
#[test]
fn test_chain_stale_after_saturates_on_overflow() {
assert_eq!(chain_stale_after(Duration::MAX), Duration::MAX);
}
}