use std::time::Duration;
use chrono::{DateTime, Utc};
use optionstratlib::OptionStyle;
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::{Decimal, Positive};
use tokio::sync::mpsc;
use crate::app::{App, EventBridge, LiveScreen, LiveState, Mode, ScreenLoad, SourceBinding};
use crate::chain::{
AliasCatalog, ChainFetch, ChainSource, ChainStore, ContractSpecFingerprint, ExerciseStyle,
ExpirySource, Instrument, InstrumentKey, MarketUpdate, MergeOutcome, ProviderId,
SettlementStyle, StreamHealth,
};
use crate::config::ThemeChoice;
use crate::event::Command;
use crate::providers::deribit::{BenchProducerStaging, bench_stream_burst, deribit_capabilities};
const EXPIRY_MS: i64 = 1_751_011_200_000;
const STRIKE_BASE: f64 = 10_000.0;
const STRIKE_STEP: f64 = 250.0;
const UNDERLYING: &str = "BTC";
fn pos(value: f64) -> Positive {
match Positive::new(value) {
Ok(p) => p,
Err(_) => Positive::ZERO,
}
}
fn deribit_pid() -> ProviderId {
match ProviderId::new("deribit") {
Ok(id) => id,
Err(_) => unreachable!("`deribit` is a valid, reserved provider id literal"),
}
}
fn expiry_utc() -> DateTime<Utc> {
DateTime::<Utc>::from_timestamp_millis(EXPIRY_MS).unwrap_or(DateTime::<Utc>::MIN_UTC)
}
fn base_time() -> DateTime<Utc> {
DateTime::<Utc>::from_timestamp(1_700_000_000, 0).unwrap_or(DateTime::<Utc>::MIN_UTC)
}
fn strike_at(index: usize) -> f64 {
let offset = u32::try_from(index).map(f64::from).unwrap_or(0.0);
STRIKE_BASE + offset * STRIKE_STEP
}
fn synthetic_spec() -> ContractSpecFingerprint {
ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: UNDERLYING.to_owned(),
}
}
fn synthetic_chain(strikes: usize) -> OptionChain {
let spot = pos(STRIKE_BASE + (strike_span(strikes) / 2.0));
let mut chain = OptionChain::new(UNDERLYING, spot, expiry_utc().to_rfc3339(), None, None);
for index in 0..strikes {
chain.add_option(
pos(strike_at(index)),
Some(pos(1.05)), Some(pos(1.25)), Some(pos(2.05)), Some(pos(2.35)), pos(0.4922), Some(Decimal::new(55, 2)), Some(Decimal::new(-45, 2)), Some(Decimal::new(1, 2)), Some(pos(12.0)), Some(100), None, );
}
chain
}
fn strike_span(strikes: usize) -> f64 {
match strikes.checked_sub(1) {
Some(last) => strike_at(last) - STRIKE_BASE,
None => 0.0,
}
}
fn synthetic_legs(strikes: usize) -> (Vec<Instrument>, AliasCatalog) {
let provider = deribit_pid();
let mut legs = Vec::new();
let mut aliases = AliasCatalog::new();
for index in 0..strikes {
for style in [OptionStyle::Call, OptionStyle::Put] {
let instrument = Instrument {
key: InstrumentKey {
underlying: UNDERLYING.to_owned(),
expiration_utc: expiry_utc(),
strike: pos(strike_at(index)),
style,
},
provider: provider.clone(),
native_symbol: format!("{UNDERLYING}-{index}-{}", style.as_str()),
stream_symbol: None,
spec: synthetic_spec(),
};
aliases.insert(instrument.clone());
legs.push(instrument);
}
}
(legs, aliases)
}
#[must_use]
pub fn seeded_store(strikes: usize) -> ChainStore {
let (_legs, aliases) = synthetic_legs(strikes);
ChainStore::seed(
ChainFetch::new(
synthetic_chain(strikes),
ExpirySource::new(UNDERLYING, expiry_utc(), deribit_pid()),
aliases,
),
ChainSource::Merged,
Duration::from_secs(2),
base_time(),
)
}
#[must_use]
pub fn live_ready_app(strikes: usize) -> App {
let (tx, _rx) = mpsc::channel::<Command>(64);
let mut live = LiveState::new(
SourceBinding::new(deribit_pid(), deribit_capabilities(), StreamHealth::Live),
seeded_store(strikes),
);
live.screen = LiveScreen::Chain;
live.load = ScreenLoad::Ready;
let mut app = App::new(Mode::Live(live), ThemeChoice::Auto, tx);
app.mark_drawn();
app
}
#[must_use]
pub fn market_burst(strikes: usize, round: u64) -> Vec<MarketUpdate> {
let (legs, _aliases) = synthetic_legs(strikes);
bench_stream_burst(&legs, round, base_time())
}
fn fold_into_store(store: &mut ChainStore, update: MarketUpdate) -> bool {
match update {
MarketUpdate::Quote(quote) => matches!(store.apply_quote("e), MergeOutcome::Applied),
MarketUpdate::Greeks(greeks) => {
matches!(store.apply_greeks(&greeks), MergeOutcome::Applied)
}
MarketUpdate::Depth(_) | MarketUpdate::Chain(_) | MarketUpdate::Health(_, _) => false,
}
}
pub struct ChainMergeHarness {
store: ChainStore,
bridge: EventBridge,
producer: BenchProducerStaging,
legs: Vec<Instrument>,
received: DateTime<Utc>,
}
impl ChainMergeHarness {
#[must_use]
pub fn new(strikes: usize, channel_capacity: usize) -> Self {
let (legs, _aliases) = synthetic_legs(strikes);
let store = seeded_store(strikes);
let (bridge, senders) = EventBridge::new(channel_capacity);
let producer = BenchProducerStaging::new(senders.market_update_sink());
Self {
store,
bridge,
producer,
legs,
received: base_time(),
}
}
#[must_use]
pub fn legs(&self) -> usize {
self.legs.len()
}
pub fn run_burst(&mut self, round: u64) -> usize {
self.publish(round);
let store = &mut self.store;
let mut applied: usize = 0;
self.bridge.pump_into(
|update| {
if fold_into_store(store, update) {
applied = applied.checked_add(1).unwrap_or(applied);
}
},
|_command| {},
);
applied
}
pub fn staging_bound(&mut self, round: u64) -> (usize, usize) {
self.publish(round);
self.bridge.coalesce_pending();
let bound = (self.bridge.staged_len(), self.bridge.staged_capacity());
let store = &mut self.store;
self.bridge.pump_into(
|update| {
let _ = fold_into_store(store, update);
},
|_command| {},
);
bound
}
#[must_use]
pub fn store_pending(&self) -> usize {
self.store.pending_len()
}
fn publish(&mut self, round: u64) {
let _ = self
.producer
.publish_burst(&self.legs, round, self.received);
}
}
#[cfg(test)]
mod tests {
use std::collections::{HashMap, HashSet};
use optionstratlib::prelude::Positive;
use tokio::sync::mpsc;
use super::{base_time, pos, synthetic_legs};
use crate::chain::{InstrumentKey, MarketUpdate};
use crate::providers::MarketUpdateSink;
use crate::providers::deribit::BenchProducerStaging;
fn round_base(round: u64) -> f64 {
let step = u32::try_from(round % 16).unwrap_or(0);
1.0 + f64::from(step) * 0.05
}
#[test]
fn producer_overwrite_on_full_keeps_newest_under_saturation() {
const STRIKES: usize = 8; const CAPACITY: usize = 4; let (legs, _aliases) = synthetic_legs(STRIKES);
let (tx, mut rx) = mpsc::channel::<MarketUpdate>(CAPACITY);
let mut producer = BenchProducerStaging::new(MarketUpdateSink::new(tx.clone(), tx.clone()));
let received = base_time();
let old = 3_u64; let new = 9_u64;
assert!(producer.publish_burst(&legs, old, received));
assert!(producer.publish_burst(&legs, new, received));
let mut quote_bid: HashMap<InstrumentKey, Positive> = HashMap::new();
let mut depth_bid: HashMap<InstrumentKey, Positive> = HashMap::new();
let mut greeks_seen: HashSet<InstrumentKey> = HashSet::new();
loop {
assert!(producer.flush(), "the consumer receiver stays alive");
let mut drained = false;
while let Ok(update) = rx.try_recv() {
drained = true;
match update {
MarketUpdate::Quote(quote) => {
if let Some(bid) = quote.bid {
let _ = quote_bid.insert(quote.instrument.key, bid);
}
}
MarketUpdate::Depth(depth) => {
if let Some(best) = depth.bids.first().map(|level| level.price) {
let _ = depth_bid.insert(depth.instrument.key, best);
}
}
MarketUpdate::Greeks(greeks) => {
let _ = greeks_seen.insert(greeks.instrument.key);
}
MarketUpdate::Chain(_) | MarketUpdate::Health(_, _) => {}
}
}
if !producer.has_pending() && !drained {
break;
}
}
assert_eq!(quote_bid.len(), legs.len(), "every leg's quote survived");
assert_eq!(depth_bid.len(), legs.len(), "every leg's depth survived");
assert_eq!(greeks_seen.len(), legs.len(), "every leg's Greeks survived");
let newest = pos(round_base(new));
let stale = pos(round_base(old));
assert_ne!(newest, stale, "the two rounds carry distinct values");
for bid in quote_bid.values() {
assert_eq!(*bid, newest, "the freshest quote won under saturation");
}
for bid in depth_bid.values() {
assert_eq!(*bid, newest, "the freshest depth won under saturation");
}
}
}