use std::collections::{HashMap, HashSet, VecDeque};
use std::time::Duration;
use chrono::{DateTime, Utc};
use optionstratlib::OptionStyle;
use optionstratlib::chains::OptionData;
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::Positive;
use super::events::{
ChainSnapshot, ChainSource, DIRECTION_DECAY, FEED_DELAY_WARN, GREEKS_STALE_AFTER, GreeksRow,
QUOTE_STALE_AFTER, QuoteUpdate, StreamHealth, chain_stale_after,
};
use super::fetch::{AliasCatalog, ChainFetch};
use super::greeks::{
GreeksSidecar, LegGreeks, PremiumNumeraire, PricingInputs, QuoteClocks, compute_dirty_legs,
compute_leg_greeks,
};
use super::identity::{InstrumentKey, ProviderId};
pub const MAX_PENDING: usize = 256;
#[must_use]
pub fn pending_ttl(refresh_interval: Duration) -> Duration {
chain_stale_after(refresh_interval)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[repr(u8)]
pub enum TickDir {
Up,
Down,
#[default]
Flat,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Freshness {
Fresh,
Delayed {
by: Duration,
},
Stale {
since: DateTime<Utc>,
},
Absent,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum MergeOutcome {
Applied,
Buffered,
DroppedOutOfOrder,
DroppedTombstoned,
DroppedCrossed,
OverlayRefused,
}
#[derive(Debug, Clone, Default)]
struct InstrumentState {
watermark: Option<DateTime<Utc>>,
receipt_watermark: Option<DateTime<Utc>>,
quote_received: Option<DateTime<Utc>>,
quote_event_time: Option<DateTime<Utc>>,
greeks_received: Option<DateTime<Utc>>,
greeks_event_time: Option<DateTime<Utc>>,
prev_bid: Option<Positive>,
prev_ask: Option<Positive>,
bid_dir: TickDir,
ask_dir: TickDir,
bid_changed_at: Option<DateTime<Utc>>,
ask_changed_at: Option<DateTime<Utc>>,
dropped_stale: u64,
}
#[derive(Debug, Clone)]
enum PendingUpdate {
Quote(QuoteUpdate),
Greeks(GreeksRow),
}
#[derive(Debug, Clone)]
struct PendingEntry {
update: PendingUpdate,
strike: Positive,
inserted_at: DateTime<Utc>,
}
#[derive(Debug)]
pub struct ChainStore {
chain_key: (ProviderId, String, DateTime<Utc>),
chain: OptionChain,
aliases: AliasCatalog,
source: ChainSource,
health: StreamHealth,
last_full_poll: Option<DateTime<Utc>>,
refresh_interval: Duration,
generation: u64,
tombstones: HashSet<Positive>,
pending: VecDeque<PendingEntry>,
dropped_overflow: u64,
instruments: HashMap<InstrumentKey, InstrumentState>,
overlay_refused: HashSet<InstrumentKey>,
sidecar: GreeksSidecar,
input_generation: u64,
analytics_as_of: DateTime<Utc>,
premium_numeraire: PremiumNumeraire,
}
impl ChainStore {
#[must_use]
pub fn seed(
fetch: ChainFetch,
source: ChainSource,
refresh_interval: Duration,
now: DateTime<Utc>,
) -> Self {
let ChainFetch {
chain,
expiry_source,
aliases,
greeks_seed,
premium_numeraire,
} = fetch;
let strikes = chain.options.len();
let instrument_cap = strikes.checked_mul(2).unwrap_or(strikes);
let chain_key = (
expiry_source.provider,
expiry_source.underlying,
expiry_source.expiration_utc,
);
let mut store = Self {
chain_key,
chain,
aliases,
source,
health: StreamHealth::Live,
last_full_poll: Some(now),
refresh_interval,
generation: 1,
tombstones: HashSet::new(),
pending: VecDeque::new(),
dropped_overflow: 0,
instruments: HashMap::with_capacity(instrument_cap),
overlay_refused: HashSet::new(),
sidecar: GreeksSidecar::new(),
input_generation: 1,
analytics_as_of: now,
premium_numeraire,
};
store.apply_greeks_seed(&greeks_seed);
store.recompute_sidecar();
store
}
pub fn apply_poll(&mut self, fetch: ChainFetch, now: DateTime<Utc>) {
self.generation = self.generation.checked_add(1).unwrap_or(self.generation);
let ChainFetch {
chain,
expiry_source,
aliases,
greeks_seed,
premium_numeraire: _,
} = fetch;
let new_strikes: HashSet<Positive> = chain.options.iter().map(|o| o.strike_price).collect();
let delisted: HashSet<Positive> = self
.chain
.options
.iter()
.map(|o| o.strike_price)
.filter(|strike| !new_strikes.contains(strike))
.collect();
for strike in &delisted {
let _ = self.tombstones.insert(*strike);
}
for strike in &new_strikes {
let _ = self.tombstones.remove(strike);
}
if !delisted.is_empty() {
self.instruments
.retain(|key, _| !delisted.contains(&key.strike));
self.overlay_refused
.retain(|key| !delisted.contains(&key.strike));
}
self.chain = chain;
self.aliases = aliases;
self.chain_key = (
expiry_source.provider,
expiry_source.underlying,
expiry_source.expiration_utc,
);
self.last_full_poll = Some(now);
self.drain_pending(now);
self.apply_greeks_seed(&greeks_seed);
self.on_option_data_changed(now);
}
pub fn apply_quote(&mut self, update: &QuoteUpdate) -> MergeOutcome {
let key = &update.instrument.key;
if let Some(outcome) = self.gate_overlay(key, &update.instrument.provider) {
return outcome;
}
if self.is_out_of_order(key, update.event_time, update.received_time) {
self.count_dropped_stale(key);
return MergeOutcome::DroppedOutOfOrder;
}
let strike = key.strike;
if self.contains_strike(strike) {
if self.apply_quote_to_row(update) {
let _ = self.overlay_refused.remove(key);
self.on_leg_data_changed(update.received_time, key);
MergeOutcome::Applied
} else {
MergeOutcome::DroppedCrossed
}
} else if self.tombstones.contains(&strike) {
MergeOutcome::DroppedTombstoned
} else if update.bid.is_some() && update.ask.is_some() && is_crossed(update.bid, update.ask)
{
MergeOutcome::DroppedCrossed
} else {
self.buffer_pending(
PendingUpdate::Quote(update.clone()),
strike,
update.received_time,
);
MergeOutcome::Buffered
}
}
pub fn apply_greeks(&mut self, update: &GreeksRow) -> MergeOutcome {
let key = &update.instrument.key;
if let Some(outcome) = self.gate_overlay(key, &update.instrument.provider) {
return outcome;
}
if self.is_out_of_order(key, update.event_time, update.received_time) {
self.count_dropped_stale(key);
return MergeOutcome::DroppedOutOfOrder;
}
let strike = key.strike;
if self.contains_strike(strike) {
self.apply_greeks_to_row(update);
let _ = self.overlay_refused.remove(key);
self.on_leg_data_changed(update.received_time, key);
MergeOutcome::Applied
} else if self.tombstones.contains(&strike) {
MergeOutcome::DroppedTombstoned
} else {
self.buffer_pending(
PendingUpdate::Greeks(update.clone()),
strike,
update.received_time,
);
MergeOutcome::Buffered
}
}
pub fn apply_health(&mut self, health: StreamHealth) {
match health {
StreamHealth::Stale { .. } | StreamHealth::Reconnecting { .. } => {
for state in self.instruments.values_mut() {
state.bid_dir = TickDir::Flat;
state.ask_dir = TickDir::Flat;
}
}
StreamHealth::Live => {}
}
self.health = health;
}
#[must_use]
pub fn chain(&self) -> &OptionChain {
&self.chain
}
#[must_use]
pub fn leg_greeks(&self, key: &InstrumentKey) -> Option<&LegGreeks> {
self.sidecar.get(key)
}
#[must_use]
pub fn aliases(&self) -> &AliasCatalog {
&self.aliases
}
#[must_use]
pub fn source(&self) -> ChainSource {
self.source
}
#[must_use]
pub fn health(&self) -> &StreamHealth {
&self.health
}
#[must_use]
pub fn chain_key(&self) -> &(ProviderId, String, DateTime<Utc>) {
&self.chain_key
}
#[must_use]
pub fn last_full_poll(&self) -> Option<DateTime<Utc>> {
self.last_full_poll
}
#[must_use]
pub fn pending_len(&self) -> usize {
self.pending.len()
}
#[must_use]
pub fn dropped_overflow(&self) -> u64 {
self.dropped_overflow
}
#[must_use]
pub fn dropped_stale(&self, key: &InstrumentKey) -> u64 {
self.instruments.get(key).map_or(0, |s| s.dropped_stale)
}
#[must_use]
pub fn is_tombstoned(&self, strike: Positive) -> bool {
self.tombstones.contains(&strike)
}
#[must_use]
pub fn contains_strike(&self, strike: Positive) -> bool {
self.chain.options.contains(&probe_row(strike))
}
#[must_use]
pub fn is_overlay_refused(&self, key: &InstrumentKey) -> bool {
self.overlay_refused.contains(key)
}
#[must_use]
pub fn bid_dir(&self, key: &InstrumentKey, now: DateTime<Utc>) -> TickDir {
match self.instruments.get(key) {
Some(state) => decayed(state.bid_dir, state.bid_changed_at, now),
None => TickDir::Flat,
}
}
#[must_use]
pub fn ask_dir(&self, key: &InstrumentKey, now: DateTime<Utc>) -> TickDir {
match self.instruments.get(key) {
Some(state) => decayed(state.ask_dir, state.ask_changed_at, now),
None => TickDir::Flat,
}
}
#[must_use]
pub fn quote_freshness(&self, key: &InstrumentKey, now: DateTime<Utc>) -> Freshness {
match self.instruments.get(key) {
Some(state) => classify(
state.quote_received,
state.quote_event_time,
now,
QUOTE_STALE_AFTER,
),
None => Freshness::Absent,
}
}
#[must_use]
pub fn greeks_freshness(&self, key: &InstrumentKey, now: DateTime<Utc>) -> Freshness {
match self.instruments.get(key) {
Some(state) => classify(
state.greeks_received,
state.greeks_event_time,
now,
GREEKS_STALE_AFTER,
),
None => Freshness::Absent,
}
}
#[must_use]
pub fn chain_freshness(&self, now: DateTime<Utc>) -> Freshness {
let threshold = chain_stale_after(self.refresh_interval);
match self.last_full_poll {
Some(polled) => {
if age_between(polled, now) > threshold {
Freshness::Stale { since: polled }
} else {
Freshness::Fresh
}
}
None => Freshness::Absent,
}
}
#[must_use]
pub fn snapshot(&self) -> ChainSnapshot {
ChainSnapshot {
chain_key: self.chain_key.clone(),
chain: self.chain.clone(),
aliases: self.aliases.clone(),
source: self.source,
health: self.health.clone(),
last_full_poll: self.last_full_poll,
}
}
fn gate_overlay(&mut self, key: &InstrumentKey, overlay: &ProviderId) -> Option<MergeOutcome> {
let source = &self.chain_key.0;
if overlay == source {
return None;
}
if self
.aliases
.overlay_compatible(key, source, overlay)
.is_err()
{
let _ = self.overlay_refused.insert(key.clone());
return Some(MergeOutcome::OverlayRefused);
}
None
}
fn instrument_state_mut(&mut self, key: &InstrumentKey) -> &mut InstrumentState {
if !self.instruments.contains_key(key) {
self.instruments
.insert(key.clone(), InstrumentState::default());
}
self.instruments
.get_mut(key)
.unwrap_or_else(|| unreachable!("instrument state was just inserted"))
}
fn is_out_of_order(
&self,
key: &InstrumentKey,
event_time: Option<DateTime<Utc>>,
received_time: DateTime<Utc>,
) -> bool {
let Some(state) = self.instruments.get(key) else {
return false;
};
match event_time {
Some(event_time) => state.watermark.is_some_and(|w| event_time < w),
None => state.receipt_watermark.is_some_and(|w| received_time < w),
}
}
fn count_dropped_stale(&mut self, key: &InstrumentKey) {
if let Some(state) = self.instruments.get_mut(key) {
state.dropped_stale = state
.dropped_stale
.checked_add(1)
.unwrap_or(state.dropped_stale);
}
}
fn apply_quote_to_row(&mut self, update: &QuoteUpdate) -> bool {
if self.patch_quote_row(update) {
self.record_quote(&update.instrument.key, update);
true
} else {
false
}
}
fn patch_quote_row(&mut self, update: &QuoteUpdate) -> bool {
let key = &update.instrument.key;
let Some(existing) = self.chain.options.get(&probe_row(key.strike)) else {
return false;
};
let mut row = existing.clone();
let (prior_bid, prior_ask) = match key.style {
OptionStyle::Call => (row.call_bid, row.call_ask),
OptionStyle::Put => (row.put_bid, row.put_ask),
};
let new_bid = update.bid.or(prior_bid);
let new_ask = update.ask.or(prior_ask);
if is_crossed(new_bid, new_ask) {
return false;
}
match key.style {
OptionStyle::Call => {
row.call_bid = new_bid;
row.call_ask = new_ask;
row.call_middle = midpoint(new_bid, new_ask);
}
OptionStyle::Put => {
row.put_bid = new_bid;
row.put_ask = new_ask;
row.put_middle = midpoint(new_bid, new_ask);
}
}
let _ = self.chain.options.replace(row);
true
}
fn apply_greeks_to_row(&mut self, update: &GreeksRow) {
let key = &update.instrument.key;
let Some(existing) = self.chain.options.get(&probe_row(key.strike)) else {
return;
};
let mut row = existing.clone();
match key.style {
OptionStyle::Call => {
if let Some(delta) = update.delta {
row.delta_call = Some(delta);
}
}
OptionStyle::Put => {
if let Some(delta) = update.delta {
row.delta_put = Some(delta);
}
}
}
if let Some(gamma) = update.gamma {
row.gamma = Some(gamma);
}
if let Some(iv) = update.iv {
row.implied_volatility = iv;
}
let _ = self.chain.options.replace(row);
self.sidecar.apply_venue_greeks(update);
self.record_greeks(key, update);
}
fn record_quote(&mut self, key: &InstrumentKey, update: &QuoteUpdate) {
let state = self.instrument_state_mut(key);
state.quote_received = Some(update.received_time);
commit_watermark(
&mut state.watermark,
&mut state.receipt_watermark,
update.event_time,
update.received_time,
);
if update.event_time.is_some() {
state.quote_event_time = update.event_time;
}
if let Some(bid) = update.bid {
update_dir(
&mut state.bid_dir,
&mut state.prev_bid,
&mut state.bid_changed_at,
bid,
update.received_time,
);
}
if let Some(ask) = update.ask {
update_dir(
&mut state.ask_dir,
&mut state.prev_ask,
&mut state.ask_changed_at,
ask,
update.received_time,
);
}
}
fn record_greeks(&mut self, key: &InstrumentKey, update: &GreeksRow) {
let state = self.instrument_state_mut(key);
state.greeks_received = Some(update.received_time);
commit_watermark(
&mut state.watermark,
&mut state.receipt_watermark,
update.event_time,
update.received_time,
);
if update.event_time.is_some() {
state.greeks_event_time = update.event_time;
}
}
fn on_option_data_changed(&mut self, as_of: DateTime<Utc>) {
self.analytics_as_of = as_of;
self.input_generation = self
.input_generation
.checked_add(1)
.unwrap_or(self.input_generation);
self.recompute_sidecar();
}
fn on_leg_data_changed(&mut self, as_of: DateTime<Utc>, key: &InstrumentKey) {
self.analytics_as_of = as_of;
self.input_generation = self
.input_generation
.checked_add(1)
.unwrap_or(self.input_generation);
self.recompute_leg(key);
}
fn recompute_leg(&mut self, key: &InstrumentKey) {
let ctx = self.pricing_inputs();
let mut clocks = QuoteClocks::new();
if let Some(received) = self.instruments.get(key).and_then(|s| s.quote_received) {
clocks.insert(key.clone(), received);
}
let dirty = [(key.strike, key.style)];
let _ = compute_dirty_legs(&self.chain, &ctx, &clocks, &dirty, &mut self.sidecar);
}
#[must_use]
pub(crate) fn quote_clocks(&self) -> QuoteClocks {
let mut clocks = QuoteClocks::new();
for (key, state) in &self.instruments {
if let Some(received) = state.quote_received {
clocks.insert(key.clone(), received);
}
}
clocks
}
#[must_use]
pub(crate) fn analytics_as_of(&self) -> DateTime<Utc> {
self.analytics_as_of
}
fn apply_greeks_seed(&mut self, greeks_seed: &[GreeksRow]) {
for row in greeks_seed {
self.sidecar.apply_venue_greeks(row);
}
}
fn pricing_inputs(&self) -> PricingInputs {
let mut ctx = PricingInputs::new(
self.chain.underlying_price,
self.analytics_as_of,
self.input_generation,
);
ctx.premium_numeraire = self.premium_numeraire;
ctx
}
fn recompute_sidecar(&mut self) {
let ctx = self.pricing_inputs();
let mut clocks = QuoteClocks::new();
for (key, state) in &self.instruments {
if let Some(received) = state.quote_received {
clocks.insert(key.clone(), received);
}
}
let _ = compute_leg_greeks(&self.chain, &ctx, &clocks, &mut self.sidecar);
}
fn buffer_pending(
&mut self,
update: PendingUpdate,
strike: Positive,
inserted_at: DateTime<Utc>,
) {
if self.pending.len() >= MAX_PENDING {
let _ = self.pending.pop_front();
self.dropped_overflow = self
.dropped_overflow
.checked_add(1)
.unwrap_or(self.dropped_overflow);
}
self.pending.push_back(PendingEntry {
update,
strike,
inserted_at,
});
}
fn drain_pending(&mut self, now: DateTime<Utc>) {
let ttl = pending_ttl(self.refresh_interval);
let entries = std::mem::take(&mut self.pending);
let mut retained: VecDeque<PendingEntry> = VecDeque::with_capacity(entries.len());
for entry in entries {
if self.contains_strike(entry.strike) {
match &entry.update {
PendingUpdate::Quote(quote) => {
let _ = self.apply_quote_to_row(quote);
}
PendingUpdate::Greeks(greeks) => {
self.apply_greeks_to_row(greeks);
}
}
} else if self.tombstones.contains(&entry.strike)
|| age_between(entry.inserted_at, now) > ttl
{
} else {
retained.push_back(entry);
}
}
self.pending = retained;
}
}
fn probe_row(strike: Positive) -> OptionData {
OptionData {
strike_price: strike,
..Default::default()
}
}
fn is_crossed(bid: Option<Positive>, ask: Option<Positive>) -> bool {
match (bid, ask) {
(Some(bid), Some(ask)) => ask < bid || (ask.is_zero() && !bid.is_zero()),
_ => false,
}
}
fn midpoint(bid: Option<Positive>, ask: Option<Positive>) -> Option<Positive> {
match (bid, ask) {
(Some(bid), Some(ask)) => Some(((bid + ask) / Positive::TWO).round_to(4)),
_ => None,
}
}
fn age_between(earlier: DateTime<Utc>, now: DateTime<Utc>) -> Duration {
now.signed_duration_since(earlier)
.to_std()
.unwrap_or(Duration::ZERO)
}
fn classify(
received: Option<DateTime<Utc>>,
event_time: Option<DateTime<Utc>>,
now: DateTime<Utc>,
threshold: Duration,
) -> Freshness {
let Some(received) = received else {
return Freshness::Absent;
};
if age_between(received, now) > threshold {
return Freshness::Stale { since: received };
}
if let Some(event_time) = event_time {
let delay = age_between(event_time, now);
if delay > FEED_DELAY_WARN {
return Freshness::Delayed { by: delay };
}
}
Freshness::Fresh
}
fn decayed(dir: TickDir, changed_at: Option<DateTime<Utc>>, now: DateTime<Utc>) -> TickDir {
match changed_at {
Some(changed_at) if age_between(changed_at, now) > DIRECTION_DECAY => TickDir::Flat,
Some(_) => dir,
None => TickDir::Flat,
}
}
fn update_dir(
dir: &mut TickDir,
prev: &mut Option<Positive>,
changed_at: &mut Option<DateTime<Utc>>,
value: Positive,
at: DateTime<Utc>,
) {
match *prev {
None => {
*dir = TickDir::Flat;
*prev = Some(value);
*changed_at = Some(at);
}
Some(previous) => {
if value > previous {
*dir = TickDir::Up;
*prev = Some(value);
*changed_at = Some(at);
} else if value < previous {
*dir = TickDir::Down;
*prev = Some(value);
*changed_at = Some(at);
}
}
}
}
fn commit_watermark(
event_watermark: &mut Option<DateTime<Utc>>,
receipt_watermark: &mut Option<DateTime<Utc>>,
event_time: Option<DateTime<Utc>>,
received_time: DateTime<Utc>,
) {
match event_time {
Some(event_time) => {
*event_watermark = Some((*event_watermark).map_or(event_time, |w| w.max(event_time)));
}
None => {
*receipt_watermark =
Some((*receipt_watermark).map_or(received_time, |w| w.max(received_time)));
}
}
}
#[cfg(test)]
mod tests {
use std::collections::BTreeSet;
use std::time::Duration;
use optionstratlib::OptionStyle;
use optionstratlib::chains::OptionData;
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::{Decimal, Positive};
use proptest::prelude::*;
use super::{ChainStore, Freshness, MAX_PENDING, MergeOutcome, TickDir, pending_ttl};
use crate::chain::events::{ChainSource, GreeksOrigin, GreeksRow, QuoteUpdate, StreamHealth};
use crate::chain::fetch::{AliasCatalog, ChainFetch, ExpirySource};
use crate::chain::greeks::{LegGreeks, LegStatus};
use crate::chain::identity::{
ContractSpecFingerprint, ExerciseStyle, Instrument, InstrumentKey, ProviderId,
SettlementStyle,
};
const EXP: i64 = 1_700_000_000;
#[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) -> chrono::DateTime<chrono::Utc> {
match chrono::DateTime::<chrono::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 dec(mantissa: i64, scale: u32) -> Decimal {
Decimal::new(mantissa, scale)
}
fn refresh() -> Duration {
Duration::from_secs(2)
}
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 ikey(strike: f64, style: OptionStyle) -> InstrumentKey {
InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: utc(EXP),
strike: pos(strike),
style,
}
}
fn instrument_spec(
provider: &str,
strike: f64,
style: OptionStyle,
multiplier: u32,
) -> Instrument {
Instrument {
key: ikey(strike, style),
provider: pid(provider),
native_symbol: format!("BTC-{strike}-{}", style.as_str()),
stream_symbol: None,
spec: spec(multiplier),
}
}
fn instrument(provider: &str, strike: f64, style: OptionStyle) -> Instrument {
instrument_spec(provider, strike, style, 1)
}
fn quote(
provider: &str,
strike: f64,
style: OptionStyle,
bid: Option<f64>,
ask: Option<f64>,
event: Option<i64>,
received: i64,
) -> QuoteUpdate {
QuoteUpdate {
instrument: instrument(provider, strike, style),
bid: bid.map(pos),
ask: ask.map(pos),
last: None,
bid_size: None,
ask_size: None,
event_time: event.map(utc),
received_time: utc(received),
}
}
fn greeks(
provider: &str,
strike: f64,
style: OptionStyle,
iv: Option<f64>,
delta: Option<Decimal>,
gamma: Option<Decimal>,
) -> GreeksRow {
GreeksRow {
instrument: instrument(provider, strike, style),
iv: iv.map(pos),
delta,
gamma,
theta: None,
vega: None,
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: utc(EXP + 100),
}
}
#[track_caller]
fn resolved_exp(store: &ChainStore) -> chrono::DateTime<chrono::Utc> {
match store.chain().get_expiration() {
Some(optionstratlib::ExpirationDate::DateTime(dt)) => dt,
other => panic!("expected an absolute-UTC chain expiry, got {other:?}"),
}
}
fn leg_key(store: &ChainStore, strike: f64, style: OptionStyle) -> InstrumentKey {
InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: resolved_exp(store),
strike: pos(strike),
style,
}
}
fn greeks_exp(
exp: chrono::DateTime<chrono::Utc>,
strike: f64,
style: OptionStyle,
iv: Option<f64>,
delta: Option<Decimal>,
gamma: Option<Decimal>,
) -> GreeksRow {
GreeksRow {
instrument: Instrument {
key: InstrumentKey {
underlying: "BTC".to_owned(),
expiration_utc: exp,
strike: pos(strike),
style,
},
provider: pid("deribit"),
native_symbol: format!("BTC-{strike}-{}", style.as_str()),
stream_symbol: None,
spec: spec(1),
},
iv: iv.map(pos),
delta,
gamma,
theta: Some(dec(-9, 1)),
vega: Some(dec(8, 1)),
rho: Some(dec(7, 1)),
origin: GreeksOrigin::Provider,
event_time: None,
received_time: utc(EXP + 100),
}
}
fn row(
strike: f64,
call_bid: Option<f64>,
call_ask: Option<f64>,
put_bid: Option<f64>,
put_ask: Option<f64>,
) -> OptionData {
let mut od = OptionData {
strike_price: pos(strike),
call_bid: call_bid.map(pos),
call_ask: call_ask.map(pos),
put_bid: put_bid.map(pos),
put_ask: put_ask.map(pos),
implied_volatility: pos(0.5),
..Default::default()
};
od.set_mid_prices();
od
}
fn chain_with(rows: &[OptionData]) -> OptionChain {
let mut chain = OptionChain::new("BTC", pos(60_000.0), "2025-06-27".to_owned(), None, None);
for od in rows {
let _ = chain.options.insert(od.clone());
}
chain
}
fn fetch_for(chain: OptionChain, provider: &str, aliases: AliasCatalog) -> ChainFetch {
ChainFetch::new(
chain,
ExpirySource::new("BTC", utc(EXP), pid(provider)),
aliases,
)
}
fn seed_single() -> ChainStore {
let chain = chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]);
ChainStore::seed(
fetch_for(chain, "deribit", AliasCatalog::new()),
ChainSource::Merged,
refresh(),
utc(EXP),
)
}
fn seed_two() -> ChainStore {
let chain = chain_with(&[
row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4)),
row(61_000.0, Some(1.5), Some(1.8), Some(2.5), Some(2.9)),
]);
ChainStore::seed(
fetch_for(chain, "deribit", AliasCatalog::new()),
ChainSource::Merged,
refresh(),
utc(EXP),
)
}
#[track_caller]
fn leg_of(store: &ChainStore, strike: f64, style: OptionStyle) -> LegGreeks {
match store.leg_greeks(&leg_key(store, strike, style)) {
Some(g) => *g,
None => panic!("expected a sidecar entry for strike {strike} {style:?}"),
}
}
#[track_caller]
fn find_row(store: &ChainStore, strike: f64) -> OptionData {
let target = pos(strike);
match store
.chain()
.options
.iter()
.find(|o| o.strike_price == target)
{
Some(o) => o.clone(),
None => panic!("expected a row at strike {strike}"),
}
}
#[test]
fn test_seed_carries_alias_catalog_and_key_forward() {
let mut catalog = AliasCatalog::new();
catalog.insert(instrument("deribit", 60_000.0, OptionStyle::Call));
let chain = chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]);
let store = ChainStore::seed(
fetch_for(chain, "deribit", catalog),
ChainSource::Merged,
refresh(),
utc(EXP),
);
assert_eq!(store.chain_key().0.as_str(), "deribit");
assert_eq!(store.chain_key().1, "BTC");
assert_eq!(store.chain_key().2, utc(EXP));
assert_eq!(store.source(), ChainSource::Merged);
assert_eq!(store.aliases().len(), 1);
assert!(matches!(store.health(), StreamHealth::Live));
assert_eq!(store.last_full_poll(), Some(utc(EXP)));
}
#[test]
fn test_snapshot_reflects_current_state() {
let store = seed_single();
let snapshot = store.snapshot();
assert_eq!(snapshot.chain_key.0.as_str(), "deribit");
assert_eq!(snapshot.source, ChainSource::Merged);
assert_eq!(snapshot.chain.symbol, "BTC");
assert!(snapshot.aliases.is_empty());
assert_eq!(snapshot.last_full_poll, Some(utc(EXP)));
}
#[test]
fn test_apply_quote_call_then_put_preserve_opposite_leg_and_middles() {
let mut store = seed_single();
let key_call = ikey(60_000.0, OptionStyle::Call);
let key_put = ikey(60_000.0, OptionStyle::Put);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.4),
Some(1.6),
None,
EXP + 100
)),
MergeOutcome::Applied
);
let after_call = find_row(&store, 60_000.0);
assert_eq!(after_call.call_bid, Some(pos(1.4)));
assert_eq!(after_call.call_ask, Some(pos(1.6)));
assert_eq!(after_call.call_middle, Some(pos(1.5)));
assert_eq!(after_call.put_bid, Some(pos(2.0)));
assert_eq!(after_call.put_ask, Some(pos(2.4)));
assert_eq!(after_call.put_middle, Some(pos(2.2)));
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Put,
Some(2.6),
Some(3.0),
None,
EXP + 101
)),
MergeOutcome::Applied
);
let after_put = find_row(&store, 60_000.0);
assert_eq!(after_put.put_bid, Some(pos(2.6)));
assert_eq!(after_put.put_ask, Some(pos(3.0)));
assert_eq!(after_put.put_middle, Some(pos(2.8)));
assert_eq!(after_put.call_bid, Some(pos(1.4)));
assert_eq!(after_put.call_middle, Some(pos(1.5)));
assert_eq!(store.bid_dir(&key_call, utc(EXP + 100)), TickDir::Flat);
assert_eq!(store.bid_dir(&key_put, utc(EXP + 101)), TickDir::Flat);
}
#[test]
fn test_apply_quote_put_then_call_yields_same_row_as_call_then_put() {
let mut a = seed_single();
let _ = a.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.4),
Some(1.6),
None,
EXP + 100,
));
let _ = a.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Put,
Some(2.6),
Some(3.0),
None,
EXP + 101,
));
let mut b = seed_single();
let _ = b.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Put,
Some(2.6),
Some(3.0),
None,
EXP + 100,
));
let _ = b.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.4),
Some(1.6),
None,
EXP + 101,
));
assert_eq!(find_row(&a, 60_000.0), find_row(&b, 60_000.0));
}
#[test]
fn test_apply_quote_crossed_ask_below_bid_rejects_whole_update() {
let mut store = seed_single();
let outcome = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(2.0),
Some(1.5),
None,
EXP + 100,
));
assert_eq!(outcome, MergeOutcome::DroppedCrossed);
let after = find_row(&store, 60_000.0);
assert_eq!(after.call_bid, Some(pos(1.0)));
assert_eq!(after.call_ask, Some(pos(1.2)));
}
#[test]
fn test_apply_quote_zero_ask_on_nonzero_bid_rejects_whole_update() {
let mut store = seed_single();
let outcome = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(0.0),
None,
EXP + 100,
));
assert_eq!(outcome, MergeOutcome::DroppedCrossed);
let after = find_row(&store, 60_000.0);
assert_eq!(after.call_bid, Some(pos(1.0)));
assert_eq!(after.call_ask, Some(pos(1.2)));
}
#[test]
fn test_apply_quote_zero_bid_is_valid() {
let mut store = seed_single();
let outcome = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(0.0),
Some(0.5),
None,
EXP + 100,
));
assert_eq!(outcome, MergeOutcome::Applied);
let after = find_row(&store, 60_000.0);
assert_eq!(after.call_bid, Some(pos(0.0)));
assert_eq!(after.call_ask, Some(pos(0.5)));
assert_eq!(after.call_middle, Some(pos(0.25)));
}
#[test]
fn test_apply_quote_missing_bid_keeps_prior_bid() {
let mut store = seed_single();
let outcome = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
None,
Some(1.6),
None,
EXP + 100,
));
assert_eq!(outcome, MergeOutcome::Applied);
let after = find_row(&store, 60_000.0);
assert_eq!(after.call_bid, Some(pos(1.0))); assert_eq!(after.call_ask, Some(pos(1.6)));
assert_eq!(after.call_middle, Some(pos(1.3)));
}
#[test]
fn test_apply_greeks_call_sets_delta_call_keeps_delta_put() {
let mut od = row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4));
od.delta_put = Some(dec(-3, 1)); let chain = chain_with(&[od]);
let mut store = ChainStore::seed(
fetch_for(chain, "deribit", AliasCatalog::new()),
ChainSource::Merged,
refresh(),
utc(EXP),
);
let outcome = store.apply_greeks(&greeks(
"deribit",
60_000.0,
OptionStyle::Call,
Some(0.4),
Some(dec(-6, 1)), Some(dec(1, 2)), ));
assert_eq!(outcome, MergeOutcome::Applied);
let after = find_row(&store, 60_000.0);
assert_eq!(after.delta_call, Some(dec(-6, 1)));
assert_eq!(after.delta_put, Some(dec(-3, 1))); assert_eq!(after.gamma, Some(dec(1, 2)));
assert_eq!(after.implied_volatility, pos(0.4));
}
#[test]
fn test_apply_greeks_missing_field_keeps_prior() {
let mut od = row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4));
od.gamma = Some(dec(2, 2)); let chain = chain_with(&[od]);
let mut store = ChainStore::seed(
fetch_for(chain, "deribit", AliasCatalog::new()),
ChainSource::Merged,
refresh(),
utc(EXP),
);
let _ = store.apply_greeks(&greeks(
"deribit",
60_000.0,
OptionStyle::Call,
None,
Some(dec(-5, 1)),
None,
));
let after = find_row(&store, 60_000.0);
assert_eq!(after.delta_call, Some(dec(-5, 1)));
assert_eq!(after.gamma, Some(dec(2, 2))); assert_eq!(after.implied_volatility, pos(0.5)); }
#[test]
fn test_sidecar_seeded_fills_local_theta_vega_rho() {
let store = seed_single();
let leg = match store.leg_greeks(&leg_key(&store, 60_000.0, OptionStyle::Call)) {
Some(g) => *g,
None => panic!("expected a seeded sidecar entry"),
};
assert!(leg.theta.is_some());
assert!(leg.vega.is_some());
assert!(leg.rho.is_some());
assert_eq!(leg.theta_origin, GreeksOrigin::ComputedLocally);
assert_eq!(leg.vega_origin, GreeksOrigin::ComputedLocally);
assert_eq!(leg.iv_origin, GreeksOrigin::ComputedLocally);
assert_eq!(leg.gamma_origin, GreeksOrigin::ComputedLocally);
}
#[test]
fn test_apply_greeks_folds_venue_iv_gamma_into_sidecar_per_style() {
let mut store = seed_single();
let exp = resolved_exp(&store);
let outcome = store.apply_greeks(&greeks_exp(
exp,
60_000.0,
OptionStyle::Call,
Some(0.42),
Some(dec(-6, 1)),
Some(dec(1, 2)),
));
assert_eq!(outcome, MergeOutcome::Applied);
let call = match store.leg_greeks(&leg_key(&store, 60_000.0, OptionStyle::Call)) {
Some(g) => *g,
None => panic!("expected a call sidecar entry"),
};
assert_eq!(call.iv, Some(pos(0.42)));
assert_eq!(call.iv_origin, GreeksOrigin::Provider);
assert_eq!(call.gamma, Some(dec(1, 2)));
assert_eq!(call.gamma_origin, GreeksOrigin::Provider);
assert!(call.theta.is_some());
assert_eq!(call.theta_origin, GreeksOrigin::ComputedLocally);
let put = match store.leg_greeks(&leg_key(&store, 60_000.0, OptionStyle::Put)) {
Some(g) => *g,
None => panic!("expected a put sidecar entry"),
};
assert_eq!(put.gamma_origin, GreeksOrigin::ComputedLocally);
}
#[test]
fn test_input_generation_bumps_only_on_applied_data_change() {
let mut store = seed_single();
let gen_seed = store.input_generation;
assert_eq!(store.sidecar.computed_generation(), Some(gen_seed));
let before = match store.leg_greeks(&leg_key(&store, 60_000.0, OptionStyle::Call)) {
Some(g) => *g,
None => panic!("expected a seed entry"),
};
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(2.0),
Some(1.5),
None,
EXP + 100
)),
MergeOutcome::DroppedCrossed
);
assert_eq!(
store.input_generation, gen_seed,
"a dropped update does not bump the pricing generation"
);
let after = match store.leg_greeks(&leg_key(&store, 60_000.0, OptionStyle::Call)) {
Some(g) => *g,
None => panic!("expected a seed entry"),
};
assert_eq!(
before, after,
"a no-op fold leaves the cached analytics intact"
);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.4),
Some(1.6),
None,
EXP + 101
)),
MergeOutcome::Applied
);
assert!(store.input_generation > gen_seed);
assert_eq!(
store.sidecar.computed_generation(),
Some(store.input_generation)
);
}
#[test]
fn test_applied_quote_recomputes_only_the_dirty_leg() {
let mut store = seed_two();
let k1_call_before = leg_of(&store, 60_000.0, OptionStyle::Call);
let k1_put_before = leg_of(&store, 60_000.0, OptionStyle::Put);
let k2_call_before = leg_of(&store, 61_000.0, OptionStyle::Call);
let k2_put_before = leg_of(&store, 61_000.0, OptionStyle::Put);
let gen_before = store.input_generation;
let outcome = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.4),
Some(1.6),
None,
EXP + 40_000_000,
));
assert_eq!(outcome, MergeOutcome::Applied);
assert!(store.input_generation > gen_before);
assert_eq!(
store.sidecar.computed_generation(),
Some(store.input_generation)
);
assert_ne!(leg_of(&store, 60_000.0, OptionStyle::Call), k1_call_before);
assert_eq!(leg_of(&store, 60_000.0, OptionStyle::Put), k1_put_before);
assert_eq!(leg_of(&store, 61_000.0, OptionStyle::Call), k2_call_before);
assert_eq!(leg_of(&store, 61_000.0, OptionStyle::Put), k2_put_before);
}
#[test]
fn test_poll_still_recomputes_every_leg() {
let mut store = seed_two();
let k2_call_before = leg_of(&store, 61_000.0, OptionStyle::Call);
let k2_put_before = leg_of(&store, 61_000.0, OptionStyle::Put);
store.apply_poll(
fetch_for(
chain_with(&[
row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4)),
row(61_000.0, Some(1.5), Some(1.8), Some(2.5), Some(2.9)),
]),
"deribit",
AliasCatalog::new(),
),
utc(EXP + 40_000_000),
);
assert_ne!(leg_of(&store, 61_000.0, OptionStyle::Call), k2_call_before);
assert_ne!(leg_of(&store, 61_000.0, OptionStyle::Put), k2_put_before);
}
#[test]
fn test_scoped_recompute_matches_full_recompute_for_dirty_leg() {
let received = EXP + 1_000;
let mut scoped = seed_two();
assert_eq!(
scoped.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.4),
Some(1.6),
None,
received,
)),
MergeOutcome::Applied
);
let scoped_leg = leg_of(&scoped, 60_000.0, OptionStyle::Call);
assert_eq!(scoped_leg.status, LegStatus::Computed);
assert!(scoped_leg.iv.is_some());
assert!(scoped_leg.delta.is_some());
assert!(scoped_leg.theta.is_some());
assert!(scoped_leg.vega.is_some());
assert!(scoped_leg.rho.is_some());
let mut full = seed_two();
full.apply_poll(
fetch_for(
chain_with(&[
row(60_000.0, Some(1.4), Some(1.6), Some(2.0), Some(2.4)),
row(61_000.0, Some(1.5), Some(1.8), Some(2.5), Some(2.9)),
]),
"deribit",
AliasCatalog::new(),
),
utc(received),
);
assert_eq!(scoped_leg, leg_of(&full, 60_000.0, OptionStyle::Call));
assert_eq!(
leg_of(&scoped, 61_000.0, OptionStyle::Call).status,
LegStatus::Computed
);
assert_eq!(
leg_of(&scoped, 60_000.0, OptionStyle::Put).status,
LegStatus::Computed
);
}
#[test]
fn test_quote_freshness_absent_when_never_received() {
let store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
assert_eq!(
store.quote_freshness(&key, utc(EXP + 10)),
Freshness::Absent
);
}
#[test]
fn test_quote_freshness_fresh_within_threshold() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 100,
));
assert_eq!(
store.quote_freshness(&key, utc(EXP + 103)),
Freshness::Fresh
);
}
#[test]
fn test_quote_freshness_stale_past_threshold() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 100,
));
match store.quote_freshness(&key, utc(EXP + 106)) {
Freshness::Stale { since } => assert_eq!(since, utc(EXP + 100)),
other => panic!("expected Stale, got {other:?}"),
}
}
#[test]
fn test_quote_freshness_delayed_when_event_time_lags() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
Some(EXP + 97),
EXP + 100,
));
match store.quote_freshness(&key, utc(EXP + 100)) {
Freshness::Delayed { by } => assert_eq!(by, Duration::from_secs(3)),
other => panic!("expected Delayed, got {other:?}"),
}
}
#[test]
fn test_quote_freshness_negative_skew_clamped_to_zero_is_fresh() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
Some(EXP + 110),
EXP + 100,
));
assert_eq!(
store.quote_freshness(&key, utc(EXP + 101)),
Freshness::Fresh
);
}
#[test]
fn test_greeks_freshness_stale_after_ten_seconds() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_greeks(&greeks(
"deribit",
60_000.0,
OptionStyle::Call,
Some(0.4),
Some(dec(-5, 1)),
None,
));
assert_eq!(
store.greeks_freshness(&key, utc(EXP + 109)),
Freshness::Fresh
);
assert!(matches!(
store.greeks_freshness(&key, utc(EXP + 111)),
Freshness::Stale { .. }
));
}
#[test]
fn test_chain_freshness_stale_past_refresh_plus_slack() {
let store = seed_single();
assert_eq!(store.chain_freshness(utc(EXP + 3)), Freshness::Fresh);
assert!(matches!(
store.chain_freshness(utc(EXP + 5)),
Freshness::Stale { .. }
));
}
#[test]
fn test_out_of_order_update_dropped_and_counted() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.5),
Some(1.7),
Some(EXP + 10),
EXP + 100
)),
MergeOutcome::Applied
);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(2.0),
Some(2.2),
Some(EXP + 9),
EXP + 101
)),
MergeOutcome::DroppedOutOfOrder
);
let after = find_row(&store, 60_000.0);
assert_eq!(after.call_bid, Some(pos(1.5))); assert_eq!(store.dropped_stale(&key), 1);
}
#[test]
fn test_none_event_time_applies_without_advancing_watermark() {
let mut store = seed_single();
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
Some(EXP + 10),
EXP + 100,
));
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(2.0),
Some(2.2),
None,
EXP + 101
)),
MergeOutcome::Applied
);
assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(2.0)));
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(3.0),
Some(3.2),
Some(EXP + 9),
EXP + 102
)),
MergeOutcome::DroppedOutOfOrder
);
assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(2.0)));
}
#[test]
fn test_rejected_update_does_not_advance_watermark() {
let mut store = seed_single();
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
Some(EXP + 50),
EXP + 100
)),
MergeOutcome::Applied
);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(2.0),
Some(1.5),
Some(EXP + 60),
EXP + 101
)),
MergeOutcome::DroppedCrossed
);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.3),
Some(1.5),
Some(EXP + 55),
EXP + 102
)),
MergeOutcome::Applied
);
let after = find_row(&store, 60_000.0);
assert_eq!(after.call_bid, Some(pos(1.3)));
assert_eq!(after.call_ask, Some(pos(1.5)));
}
#[test]
fn test_no_timestamp_updates_apply_in_receipt_order() {
let mut store = seed_single();
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.5),
Some(1.7),
None,
EXP + 102
)),
MergeOutcome::Applied
);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(2.0),
Some(2.2),
None,
EXP + 101
)),
MergeOutcome::DroppedOutOfOrder
);
assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(1.5)));
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(2.5),
Some(2.7),
None,
EXP + 103
)),
MergeOutcome::Applied
);
assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(2.5)));
}
#[test]
fn test_direction_first_ever_is_flat() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 100,
));
assert_eq!(store.bid_dir(&key, utc(EXP + 100)), TickDir::Flat);
}
#[test]
fn test_direction_up_then_down_then_equal_keeps_prior() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 100,
));
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.5),
Some(1.7),
None,
EXP + 101,
));
assert_eq!(store.bid_dir(&key, utc(EXP + 101)), TickDir::Up);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.2),
Some(1.4),
None,
EXP + 102,
));
assert_eq!(store.bid_dir(&key, utc(EXP + 102)), TickDir::Down);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.2),
Some(1.4),
None,
EXP + 103,
));
assert_eq!(store.bid_dir(&key, utc(EXP + 103)), TickDir::Down);
}
#[test]
fn test_direction_decays_to_flat_after_threshold() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 100,
));
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.5),
Some(1.7),
None,
EXP + 101,
));
assert_eq!(store.bid_dir(&key, utc(EXP + 103)), TickDir::Up);
assert_eq!(store.bid_dir(&key, utc(EXP + 105)), TickDir::Flat);
}
#[test]
fn test_direction_cleared_to_flat_on_stale_health() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 100,
));
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.5),
Some(1.7),
None,
EXP + 101,
));
assert_eq!(store.bid_dir(&key, utc(EXP + 101)), TickDir::Up);
store.apply_health(StreamHealth::Stale {
since: utc(EXP + 102),
});
assert_eq!(store.bid_dir(&key, utc(EXP + 101)), TickDir::Flat);
}
#[test]
fn test_direction_cleared_to_flat_on_reconnecting_health() {
let mut store = seed_single();
let key = ikey(60_000.0, OptionStyle::Call);
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 100,
));
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(0.5),
Some(0.7),
None,
EXP + 101,
));
assert_eq!(store.bid_dir(&key, utc(EXP + 101)), TickDir::Down);
store.apply_health(StreamHealth::Reconnecting { attempt: 1 });
assert_eq!(store.bid_dir(&key, utc(EXP + 101)), TickDir::Flat);
}
#[test]
fn test_tombstoned_strike_update_dropped_not_resurrected() {
let mut store = ChainStore::seed(
fetch_for(
chain_with(&[
row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4)),
row(61_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4)),
]),
"deribit",
AliasCatalog::new(),
),
ChainSource::Merged,
refresh(),
utc(EXP),
);
store.apply_poll(
fetch_for(
chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]),
"deribit",
AliasCatalog::new(),
),
utc(EXP + 2),
);
assert!(store.is_tombstoned(pos(61_000.0)));
let outcome = store.apply_quote("e(
"deribit",
61_000.0,
OptionStyle::Call,
Some(9.0),
Some(9.2),
None,
EXP + 3,
));
assert_eq!(outcome, MergeOutcome::DroppedTombstoned);
assert!(!store.contains_strike(pos(61_000.0)));
assert_eq!(store.pending_len(), 0);
}
#[test]
fn test_new_listing_from_pending_applied_on_poll() {
let mut store = seed_single(); let outcome = store.apply_quote("e(
"deribit",
61_000.0,
OptionStyle::Call,
Some(5.0),
Some(5.2),
None,
EXP + 1,
));
assert_eq!(outcome, MergeOutcome::Buffered);
assert_eq!(store.pending_len(), 1);
store.apply_poll(
fetch_for(
chain_with(&[
row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4)),
row(61_000.0, Some(0.1), Some(0.2), Some(0.3), Some(0.4)),
]),
"deribit",
AliasCatalog::new(),
),
utc(EXP + 2),
);
assert!(store.contains_strike(pos(61_000.0)));
assert_eq!(store.pending_len(), 0);
let after = find_row(&store, 61_000.0);
assert_eq!(after.call_bid, Some(pos(5.0)));
assert_eq!(after.call_ask, Some(pos(5.2)));
}
#[test]
fn test_pending_expired_dropped_on_ttl() {
let mut store = seed_single();
let _ = store.apply_quote("e(
"deribit",
61_000.0,
OptionStyle::Call,
Some(5.0),
Some(5.2),
None,
EXP,
));
assert_eq!(store.pending_len(), 1);
let ttl_secs = pending_ttl(refresh()).as_secs() as i64;
store.apply_poll(
fetch_for(
chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]),
"deribit",
AliasCatalog::new(),
),
utc(EXP + ttl_secs + 1),
);
assert_eq!(store.pending_len(), 0);
assert!(!store.contains_strike(pos(61_000.0)));
}
#[test]
fn test_pending_overflow_drops_oldest_counted() {
let mut store = seed_single();
for i in 0u32..300 {
let outcome = store.apply_quote("e(
"deribit",
70_000.0 + f64::from(i),
OptionStyle::Call,
Some(1.0),
Some(1.2),
None,
EXP + 200,
));
assert_eq!(outcome, MergeOutcome::Buffered);
}
assert_eq!(store.pending_len(), MAX_PENDING);
assert_eq!(store.dropped_overflow(), 300 - MAX_PENDING as u64);
}
#[test]
fn test_rejected_key_burst_does_not_grow_instrument_sidecar() {
let mut store = seed_single(); for i in 0u32..300 {
let outcome = store.apply_quote("e(
"deribit",
70_000.0 + f64::from(i),
OptionStyle::Call,
Some(1.0),
Some(1.2),
Some(EXP + 200 + i64::from(i)),
EXP + 200 + i64::from(i),
));
assert_eq!(outcome, MergeOutcome::Buffered);
}
assert_eq!(store.pending_len(), MAX_PENDING);
assert_eq!(store.instruments.len(), 0);
assert_eq!(
store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(1.0),
Some(1.2),
Some(EXP + 600),
EXP + 600
)),
MergeOutcome::Applied
);
assert_eq!(store.instruments.len(), 1);
}
#[test]
fn test_delisted_strike_prunes_instrument_sidecar() {
let mut store = ChainStore::seed(
fetch_for(
chain_with(&[
row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4)),
row(61_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4)),
]),
"deribit",
AliasCatalog::new(),
),
ChainSource::Merged,
refresh(),
utc(EXP),
);
assert_eq!(
store.apply_quote("e(
"deribit",
61_000.0,
OptionStyle::Call,
Some(1.5),
Some(1.7),
Some(EXP + 10),
EXP + 100
)),
MergeOutcome::Applied
);
assert_eq!(store.instruments.len(), 1);
store.apply_poll(
fetch_for(
chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]),
"deribit",
AliasCatalog::new(),
),
utc(EXP + 2),
);
assert!(store.is_tombstoned(pos(61_000.0)));
assert_eq!(store.instruments.len(), 0);
assert_eq!(
store.apply_quote("e(
"deribit",
61_000.0,
OptionStyle::Call,
Some(9.0),
Some(9.2),
Some(EXP + 11),
EXP + 103
)),
MergeOutcome::DroppedTombstoned
);
assert_eq!(store.instruments.len(), 0);
}
#[test]
fn test_relisted_strike_clears_tombstone() {
let mut store = ChainStore::seed(
fetch_for(
chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]),
"deribit",
AliasCatalog::new(),
),
ChainSource::Merged,
refresh(),
utc(EXP),
);
store.apply_poll(
fetch_for(chain_with(&[]), "deribit", AliasCatalog::new()),
utc(EXP + 2),
);
assert!(store.is_tombstoned(pos(60_000.0)));
store.apply_poll(
fetch_for(
chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]),
"deribit",
AliasCatalog::new(),
),
utc(EXP + 4),
);
assert!(!store.is_tombstoned(pos(60_000.0)));
assert!(store.contains_strike(pos(60_000.0)));
}
fn overlay_store(source_mult: u32, overlay_mult: u32) -> ChainStore {
let mut catalog = AliasCatalog::new();
catalog.insert(instrument_spec(
"deribit",
60_000.0,
OptionStyle::Call,
source_mult,
));
catalog.insert(instrument_spec(
"dxlink",
60_000.0,
OptionStyle::Call,
overlay_mult,
));
let chain = chain_with(&[row(60_000.0, Some(1.0), Some(1.2), Some(2.0), Some(2.4))]);
ChainStore::seed(
fetch_for(chain, "deribit", catalog),
ChainSource::Merged,
refresh(),
utc(EXP),
)
}
#[test]
fn test_overlay_refused_on_spec_mismatch_keeps_source() {
let mut store = overlay_store(1, 100);
let key = ikey(60_000.0, OptionStyle::Call);
let outcome = store.apply_quote("e(
"dxlink",
60_000.0,
OptionStyle::Call,
Some(5.0),
Some(5.2),
None,
EXP + 10,
));
assert_eq!(outcome, MergeOutcome::OverlayRefused);
assert!(store.is_overlay_refused(&key));
assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(1.0)));
}
#[test]
fn test_overlay_applied_on_spec_match() {
let mut store = overlay_store(1, 1);
let key = ikey(60_000.0, OptionStyle::Call);
let outcome = store.apply_quote("e(
"dxlink",
60_000.0,
OptionStyle::Call,
Some(5.0),
Some(5.2),
None,
EXP + 10,
));
assert_eq!(outcome, MergeOutcome::Applied);
assert!(!store.is_overlay_refused(&key));
assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(5.0)));
}
#[test]
fn test_within_provider_merge_bypasses_gate_with_empty_catalog() {
let mut store = seed_single();
let outcome = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(5.0),
Some(5.2),
None,
EXP + 10,
));
assert_eq!(outcome, MergeOutcome::Applied);
assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(5.0)));
}
fn strike_of(idx: u32) -> f64 {
60_000.0 + f64::from(idx) * 1_000.0
}
fn chain_of(indices: &BTreeSet<u32>) -> OptionChain {
let mut chain = OptionChain::new("BTC", pos(60_000.0), "2025-06-27".to_owned(), None, None);
for &idx in indices {
let _ = chain.options.insert(row(
strike_of(idx),
Some(1.0),
Some(1.2),
Some(2.0),
Some(2.4),
));
}
chain
}
#[derive(Debug, Clone)]
enum Op {
Poll(Vec<u32>),
Quote(u32),
}
fn op_strategy() -> impl Strategy<Value = Op> {
prop_oneof![
proptest::collection::vec(0u32..5, 0..5).prop_map(Op::Poll),
(0u32..5).prop_map(Op::Quote),
]
}
proptest! {
#![proptest_config(ProptestConfig { cases: 512, max_shrink_iters: 20_000, ..ProptestConfig::default() })]
#[test]
fn prop_chain_merge_idempotent(
bid_ticks in 1u32..80,
spread in 1u32..30,
event in 0i64..100,
) {
let bid = f64::from(bid_ticks) * 0.1;
let ask = bid + f64::from(spread) * 0.1;
let q = quote(
"deribit", 60_000.0, OptionStyle::Call, Some(bid), Some(ask), Some(EXP + event), EXP + 200,
);
let mut once = seed_single();
let _ = once.apply_quote(&q);
let mut twice = seed_single();
let _ = twice.apply_quote(&q);
let _ = twice.apply_quote(&q);
prop_assert_eq!(find_row(&once, 60_000.0), find_row(&twice, 60_000.0));
}
#[test]
fn prop_overlay_spec_gate(overlay_mult in 1u32..250) {
let mut store = overlay_store(1, overlay_mult);
let outcome = store.apply_quote("e(
"dxlink", 60_000.0, OptionStyle::Call, Some(5.0), Some(5.2), None, EXP + 10,
));
let after = find_row(&store, 60_000.0);
if overlay_mult == 1 {
prop_assert_eq!(outcome, MergeOutcome::Applied);
prop_assert_eq!(after.call_bid, Some(pos(5.0)));
} else {
prop_assert_eq!(outcome, MergeOutcome::OverlayRefused);
prop_assert_eq!(after.call_bid, Some(pos(1.0)));
}
}
#[test]
fn prop_no_resurrection_and_bounded_memory(ops in proptest::collection::vec(op_strategy(), 0..40)) {
let all: BTreeSet<u32> = (0u32..5).collect();
let mut store = ChainStore::seed(
fetch_for(chain_of(&all), "deribit", AliasCatalog::new()),
ChainSource::Merged,
refresh(),
utc(EXP),
);
let mut last_poll = all.clone();
for (step, op) in (1i64..).zip(ops) {
match op {
Op::Poll(indices) => {
let set: BTreeSet<u32> = indices.into_iter().collect();
store.apply_poll(
fetch_for(chain_of(&set), "deribit", AliasCatalog::new()),
utc(EXP + step),
);
last_poll = set;
}
Op::Quote(idx) => {
let _ = store.apply_quote("e(
"deribit",
strike_of(idx),
OptionStyle::Call,
Some(1.5),
Some(1.7),
None,
EXP + step,
));
}
}
}
for idx in 0u32..5 {
let strike = pos(strike_of(idx));
let present = store.contains_strike(strike);
prop_assert_eq!(present, last_poll.contains(&idx));
if present {
prop_assert!(!store.is_tombstoned(strike));
}
}
prop_assert!(store.pending_len() <= MAX_PENDING);
}
#[test]
fn prop_freshness_out_of_order_keeps_max_event_value(
order in Just((0u32..6).collect::<Vec<u32>>()).prop_shuffle()
) {
let mut store = seed_single();
for i in order {
let bid = 1.0 + f64::from(i) * 0.5;
let ask = bid + 0.5;
let _ = store.apply_quote("e(
"deribit",
60_000.0,
OptionStyle::Call,
Some(bid),
Some(ask),
Some(EXP + i64::from(i)),
EXP + 200,
));
}
prop_assert_eq!(find_row(&store, 60_000.0).call_bid, Some(pos(1.0 + 5.0 * 0.5)));
}
}
}