use std::collections::BTreeSet;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use chrono::{DateTime, Datelike, NaiveDate, NaiveDateTime, TimeDelta, Utc, Weekday};
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::{Decimal, Positive};
use optionstratlib::{ExpirationDate, OptionStyle};
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use ibapi::contracts::tick_types::TickType;
use ibapi::contracts::{
Contract, ContractDetails, OptionChain as IbOptionChain, OptionComputation, OptionRight,
SecurityType,
};
use ibapi::market_data::realtime::{TickPrice, TickSize, TickTypes};
use ibapi::prelude::StreamExt;
use ibapi::subscriptions::SubscriptionItem;
use super::{
AuthKind, ChainCapability, ChainPollCapability, GreeksCapability, MarketUpdateSink,
OptionStreamCapability, Provider, ProviderCapabilities, SendState, SubscriptionHandle,
SubscriptionRequest, UnderlyingRef,
};
use crate::chain::{
AliasCatalog, ChainFetch, ChainSnapshot, ChainSource, ContractSpecFingerprint, ExerciseStyle,
ExpirySource, GreeksOrigin, GreeksRow, Instrument, InstrumentKey, MarketUpdate, ProviderId,
QuoteUpdate, SettlementStyle, StreamHealth,
};
use crate::config::{EnvSource, provider_env_var};
use crate::error::{ConfigError, NormalizeKind, ProviderError, TransportDetail, TransportKind};
const IBKR_ID: &str = "ibkr";
const ENDPOINT_KEY: &str = "ENDPOINT";
const CLIENT_ID_KEY: &str = "CLIENT_ID";
const DEFAULT_CLIENT_ID: i32 = 1;
const MAX_SUBSCRIPTIONS: usize = 100;
const MAX_STRIKES: usize = 4_096;
const STRIKE_CAP: &str = "ibkr strike cap";
const REFRESH_HINT_SECS: u32 = 5;
const DEFAULT_CONTRACT_MULTIPLIER: u32 = 100;
const DEFAULT_QUOTE_CURRENCY: &str = "USD";
const CANDIDATE_UNDERLYINGS: [&str; 10] = [
"SPX", "SPY", "QQQ", "IWM", "AAPL", "MSFT", "NVDA", "AMZN", "TSLA", "META",
];
const INDEX_UNDERLYINGS: [&str; 5] = ["SPX", "NDX", "RUT", "VIX", "XSP"];
const OPTION_PARAMS_TIMEOUT_SECS: u64 = 10;
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 IbkrAdapter {
id: ProviderId,
endpoint: String,
client_id: i32,
}
impl std::fmt::Debug for IbkrAdapter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("IbkrAdapter")
.field("id", &self.id)
.field("endpoint", &"<redacted>")
.field("client_id", &"<redacted>")
.finish()
}
}
impl IbkrAdapter {
pub(crate) fn from_env(env: &dyn EnvSource) -> Result<Self, ConfigError> {
let id = ibkr_provider_id();
let endpoint = env
.get(&provider_env_var(id.as_str(), ENDPOINT_KEY))
.map(|value| value.trim().to_owned())
.filter(|value| !value.is_empty());
let Some(endpoint) = endpoint else {
return Err(ConfigError::MissingCredential(id));
};
let Some(endpoint) = normalize_endpoint(&endpoint) else {
return Err(ConfigError::InvalidValue {
field: "ibkr endpoint".to_owned(),
reason: "CHAINVIEW_IBKR_ENDPOINT must be host:port or scheme://host:port \
(e.g. tcp://127.0.0.1:7497)"
.to_owned(),
});
};
let client_id = match env
.get(&provider_env_var(id.as_str(), CLIENT_ID_KEY))
.map(|value| value.trim().to_owned())
.filter(|value| !value.is_empty())
{
Some(raw) => raw.parse::<i32>().map_err(|_| ConfigError::InvalidValue {
field: "ibkr client id".to_owned(),
reason: "CHAINVIEW_IBKR_CLIENT_ID must be an integer".to_owned(),
})?,
None => DEFAULT_CLIENT_ID,
};
Ok(Self {
id,
endpoint,
client_id,
})
}
}
#[async_trait]
impl Provider for IbkrAdapter {
fn id(&self) -> ProviderId {
self.id.clone()
}
fn capabilities(&self) -> ProviderCapabilities {
ibkr_capabilities()
}
async fn discover(&self) -> Result<Vec<UnderlyingRef>, ProviderError> {
Ok(CANDIDATE_UNDERLYINGS
.iter()
.map(|symbol| UnderlyingRef::new(*symbol))
.collect())
}
async fn fetch_chain(
&self,
underlying: &str,
expiration: &ExpirationDate,
) -> Result<ChainFetch, ProviderError> {
let source = LiveChainSource::connect(self).await?;
let assembled = compose_chain(&source, underlying, expiration, &self.id, now_utc()).await?;
Ok(assembled.fetch)
}
async fn subscribe(
&self,
req: SubscriptionRequest,
sink: MarketUpdateSink,
) -> Result<SubscriptionHandle, ProviderError> {
let transport = LiveStreamTransport::new(self.clone());
let id = self.id.clone();
let SubscriptionRequest {
underlying,
expiration_utc,
instruments,
cancel,
} = req;
let mut aliases = AliasCatalog::new();
for instrument in &instruments {
aliases.insert(instrument.clone());
}
let loop_cancel = cancel.clone();
let handle = tokio::spawn(run_reconnect_loop(
transport,
id,
underlying,
expiration_utc,
aliases,
sink,
loop_cancel,
));
Ok(SubscriptionHandle::spawned(cancel, handle))
}
}
fn ibkr_provider_id() -> ProviderId {
match ProviderId::new(IBKR_ID) {
Ok(id) => id,
Err(_) => unreachable!("`ibkr` is a valid, reserved provider id literal"),
}
}
#[must_use]
pub(crate) fn ibkr_capabilities() -> ProviderCapabilities {
ProviderCapabilities::builder()
.chain(ChainCapability::Assemble)
.depth(false)
.greeks(GreeksCapability::Provided)
.option_stream(OptionStreamCapability::ChainQuotes { verified: false })
.underlying_stream(false)
.chain_poll(ChainPollCapability::Poll {
interval_hint_secs: REFRESH_HINT_SECS,
})
.trades_tape(false)
.auth(AuthKind::None)
.build()
}
fn is_valid_endpoint(endpoint: &str) -> bool {
let Some((host, port)) = endpoint.rsplit_once(':') else {
return false;
};
!host.is_empty()
&& !host.contains('/')
&& !host.contains('@')
&& !port.is_empty()
&& port.parse::<u16>().is_ok()
}
fn normalize_endpoint(raw: &str) -> Option<String> {
let socket = match raw.split_once("://") {
Some((scheme, authority)) => {
let scheme_ok = matches!(scheme.chars().next(), Some(c) if c.is_ascii_alphabetic())
&& scheme
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '+' | '-' | '.'));
if !scheme_ok {
return None;
}
authority
}
None => raw,
};
is_valid_endpoint(socket).then(|| socket.to_owned())
}
fn is_index_underlying(underlying: &str) -> bool {
INDEX_UNDERLYINGS.contains(&underlying)
}
fn security_type_for(underlying: &str) -> SecurityType {
if is_index_underlying(underlying) {
SecurityType::Index
} else {
SecurityType::Stock
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SessionZone {
UsEastern,
UkLondon,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct SessionClose {
zone: SessionZone,
hour: u32,
minute: u32,
}
impl SessionClose {
const fn new(zone: SessionZone, hour: u32, minute: u32) -> Self {
Self { zone, hour, minute }
}
}
fn default_session() -> SessionClose {
SessionClose::new(SessionZone::UsEastern, 16, 0)
}
fn session_from_zone_id(zone_id: &str) -> SessionClose {
let upper = zone_id.to_ascii_uppercase();
if upper.contains("LONDON") || upper.contains("GMT") || upper.contains("BST") {
SessionClose::new(SessionZone::UkLondon, 16, 30)
} else {
default_session()
}
}
fn resolve_expiry(
last_trade_time: Option<&str>,
ymd: &str,
session: &SessionClose,
) -> Result<DateTime<Utc>, NormalizeKind> {
let date = parse_ymd(ymd).ok_or(NormalizeKind::UnparseableExpiry)?;
let naive = match last_trade_time.map(str::trim).filter(|s| !s.is_empty()) {
Some(stamp) => {
parse_ib_last_trade_time(stamp, date).ok_or(NormalizeKind::UnparseableExpiry)?
}
None => date
.and_hms_opt(session.hour, session.minute, 0)
.ok_or(NormalizeKind::UnparseableExpiry)?,
};
local_to_utc(naive, session.zone).ok_or(NormalizeKind::UnparseableExpiry)
}
fn parse_ymd(s: &str) -> Option<NaiveDate> {
let t = s.trim();
if t.len() != 8 || !t.chars().all(|c| c.is_ascii_digit()) {
return None;
}
let year = t.get(0..4)?.parse::<i32>().ok()?;
let month = t.get(4..6)?.parse::<u32>().ok()?;
let day = t.get(6..8)?.parse::<u32>().ok()?;
NaiveDate::from_ymd_opt(year, month, day)
}
fn parse_ib_last_trade_time(s: &str, date: NaiveDate) -> Option<NaiveDateTime> {
let t = s.trim();
if let Some((maybe_date, time_part)) = t.split_once('-').or_else(|| t.split_once(' '))
&& let Some(full_date) = parse_ymd(maybe_date)
{
return combine_date_time(full_date, time_part);
}
combine_date_time(date, t)
}
fn combine_date_time(date: NaiveDate, time_str: &str) -> Option<NaiveDateTime> {
let mut parts = time_str.trim().split(':');
let hour = parts.next()?.trim().parse::<u32>().ok()?;
let minute = parts.next()?.trim().parse::<u32>().ok()?;
let second = match parts.next() {
Some(sec) => sec.trim().parse::<u32>().ok()?,
None => 0,
};
if parts.next().is_some() {
return None;
}
date.and_hms_opt(hour, minute, second)
}
fn zone_offset_hours(date: NaiveDate, zone: SessionZone) -> i64 {
match zone {
SessionZone::UsEastern => {
if us_eastern_dst(date) {
-4
} else {
-5
}
}
SessionZone::UkLondon => {
if uk_dst(date) {
1
} else {
0
}
}
}
}
fn local_to_utc(naive: NaiveDateTime, zone: SessionZone) -> Option<DateTime<Utc>> {
let offset = zone_offset_hours(naive.date(), zone);
let utc_naive = naive.checked_sub_signed(TimeDelta::hours(offset))?;
Some(DateTime::<Utc>::from_naive_utc_and_offset(utc_naive, Utc))
}
fn utc_to_zone_date(utc: DateTime<Utc>, zone: SessionZone) -> NaiveDate {
let offset = zone_offset_hours(utc.date_naive(), zone);
let local = utc
.checked_add_signed(TimeDelta::hours(offset))
.unwrap_or(utc);
local.date_naive()
}
fn format_ymd_in_zone(utc: DateTime<Utc>, zone: SessionZone) -> String {
utc_to_zone_date(utc, zone).format("%Y%m%d").to_string()
}
fn uk_dst(date: NaiveDate) -> bool {
let year = date.year();
match (
last_weekday_of_month(year, 3, Weekday::Sun),
last_weekday_of_month(year, 10, Weekday::Sun),
) {
(Some(start), Some(end)) => date >= start && date < end,
_ => false,
}
}
fn us_eastern_dst(date: NaiveDate) -> bool {
let year = date.year();
match (
nth_weekday_of_month(year, 3, Weekday::Sun, 2),
nth_weekday_of_month(year, 11, Weekday::Sun, 1),
) {
(Some(start), Some(end)) => date >= start && date < end,
_ => false,
}
}
fn nth_weekday_of_month(year: i32, month: u32, weekday: Weekday, n: u32) -> Option<NaiveDate> {
let first = NaiveDate::from_ymd_opt(year, month, 1)?;
let first_dow = first.weekday().num_days_from_sunday();
let target_dow = weekday.num_days_from_sunday();
let offset = (target_dow + 7 - first_dow) % 7;
let day = 1u32
.checked_add(offset)?
.checked_add(n.checked_sub(1)?.checked_mul(7)?)?;
NaiveDate::from_ymd_opt(year, month, day)
}
fn last_weekday_of_month(year: i32, month: u32, weekday: Weekday) -> Option<NaiveDate> {
let (next_year, next_month) = if month == 12 {
(year.checked_add(1)?, 1)
} else {
(year, month.checked_add(1)?)
};
let first_next = NaiveDate::from_ymd_opt(next_year, next_month, 1)?;
let mut cursor = first_next.pred_opt()?;
let target = weekday.num_days_from_sunday();
for _ in 0..7 {
if cursor.weekday().num_days_from_sunday() == target {
return Some(cursor);
}
cursor = cursor.pred_opt()?;
}
None
}
fn positive_or_drop(value: f64) -> Option<Positive> {
if !value.is_finite() {
return None;
}
Positive::new(value).ok()
}
fn iv_or_drop(value: f64) -> Option<Positive> {
if !value.is_finite() {
return None;
}
Positive::new(value).ok()
}
fn greek_or_drop(value: Option<f64>) -> Option<Decimal> {
let raw = value?;
if !raw.is_finite() {
return None;
}
Decimal::try_from(raw).ok()
}
fn strike_positive(value: f64) -> Result<Positive, NormalizeKind> {
let strike = positive_or_drop(value).ok_or(NormalizeKind::OutOfRange("strike"))?;
if strike == Positive::ZERO {
return Err(NormalizeKind::OutOfRange("strike"));
}
Ok(strike)
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
struct NormalizedQuote {
bid: Option<Positive>,
ask: Option<Positive>,
}
fn normalize_quote(bid: Option<f64>, ask: Option<f64>) -> Result<NormalizedQuote, NormalizeKind> {
let bid = bid.and_then(positive_or_drop);
let ask = ask.and_then(positive_or_drop);
if let (Some(bid_value), Some(ask_value)) = (bid, ask)
&& (ask_value < bid_value || (ask_value == Positive::ZERO && bid_value != Positive::ZERO))
{
return Err(NormalizeKind::OutOfRange("ask"));
}
Ok(NormalizedQuote { bid, ask })
}
fn ibkr_fingerprint(multiplier: u32, trading_class: &str) -> ContractSpecFingerprint {
ContractSpecFingerprint {
contract_multiplier: multiplier,
settlement: SettlementStyle::Physical,
exercise: ExerciseStyle::American,
quote_currency: DEFAULT_QUOTE_CURRENCY.to_owned(),
venue_product_code: trading_class.to_owned(),
}
}
fn ibkr_native_symbol(underlying: &str, ymd: &str, strike: Positive, style: OptionStyle) -> String {
format!("{underlying}-{ymd}-{strike}-{}", style_char(style))
}
fn style_char(style: OptionStyle) -> char {
match style {
OptionStyle::Call => 'C',
OptionStyle::Put => 'P',
}
}
#[derive(Debug, Clone)]
pub(crate) struct RawOptionParams {
trading_class: String,
multiplier: String,
#[allow(dead_code)]
exchange: String,
expirations: Vec<String>,
strikes: Vec<f64>,
#[allow(dead_code)]
underlying_conid: i32,
}
#[derive(Debug, Clone)]
pub(crate) struct RawContractDetail {
last_trade_time: Option<String>,
time_zone_id: String,
}
#[derive(Debug, Clone)]
enum RawTick {
Quote(RawQuoteTick),
Greeks(RawGreeksTick),
}
#[derive(Debug, Clone, Default, PartialEq)]
struct RawQuoteTick {
symbol: String,
bid: Option<f64>,
ask: Option<f64>,
last: Option<f64>,
bid_size: Option<f64>,
ask_size: Option<f64>,
}
#[derive(Debug, Clone, Default, PartialEq)]
struct RawGreeksTick {
symbol: String,
iv: Option<f64>,
delta: Option<f64>,
gamma: Option<f64>,
theta: Option<f64>,
vega: Option<f64>,
}
fn map_option_chain(chain: &IbOptionChain) -> RawOptionParams {
RawOptionParams {
trading_class: chain.trading_class.clone(),
multiplier: chain.multiplier.clone(),
exchange: chain.exchange.clone(),
expirations: chain.expirations.clone(),
strikes: chain.strikes.clone(),
underlying_conid: chain.underlying_contract_id,
}
}
fn map_contract_details(details: &ContractDetails) -> RawContractDetail {
let last_trade_time = if details.last_trade_time.trim().is_empty() {
None
} else {
Some(details.last_trade_time.clone())
};
RawContractDetail {
last_trade_time,
time_zone_id: details.time_zone_id.clone(),
}
}
fn map_option_computation(symbol: &str, oc: &OptionComputation) -> RawGreeksTick {
RawGreeksTick {
symbol: symbol.to_owned(),
iv: oc.implied_volatility,
delta: oc.delta,
gamma: oc.gamma,
theta: oc.theta,
vega: oc.vega,
}
}
#[derive(Debug, Clone, Default)]
struct QuoteAccumulator {
bid: Option<f64>,
ask: Option<f64>,
last: Option<f64>,
bid_size: Option<f64>,
ask_size: Option<f64>,
}
impl QuoteAccumulator {
fn apply_price(&mut self, tick_type: TickType, price: f64) -> bool {
match tick_type {
TickType::Bid | TickType::BidOption => self.bid = Some(price),
TickType::Ask | TickType::AskOption => self.ask = Some(price),
TickType::Last | TickType::LastOption => self.last = Some(price),
_ => return false,
}
true
}
fn apply_size(&mut self, tick_type: TickType, size: f64) -> bool {
match tick_type {
TickType::BidSize => self.bid_size = Some(size),
TickType::AskSize => self.ask_size = Some(size),
_ => return false,
}
true
}
fn snapshot(&self, symbol: &str) -> RawQuoteTick {
RawQuoteTick {
symbol: symbol.to_owned(),
bid: self.bid,
ask: self.ask,
last: self.last,
bid_size: self.bid_size,
ask_size: self.ask_size,
}
}
}
fn quote_update(
tick: &RawQuoteTick,
aliases: &AliasCatalog,
provider: &ProviderId,
received: DateTime<Utc>,
) -> Option<QuoteUpdate> {
let key = aliases.resolve_symbol(&tick.symbol)?.clone();
let instrument = aliases.instrument(&key, provider)?.clone();
let quote = normalize_quote(tick.bid, tick.ask).ok()?;
if quote.bid.is_none() && quote.ask.is_none() {
return None;
}
Some(QuoteUpdate {
instrument,
bid: quote.bid,
ask: quote.ask,
last: tick.last.and_then(positive_or_drop),
bid_size: tick.bid_size.and_then(positive_or_drop),
ask_size: tick.ask_size.and_then(positive_or_drop),
event_time: None,
received_time: received,
})
}
fn greeks_row(
tick: &RawGreeksTick,
aliases: &AliasCatalog,
provider: &ProviderId,
received: DateTime<Utc>,
) -> Option<GreeksRow> {
let key = aliases.resolve_symbol(&tick.symbol)?.clone();
let instrument = aliases.instrument(&key, provider)?.clone();
let iv = tick.iv.and_then(iv_or_drop);
let delta = greek_or_drop(tick.delta);
let gamma = greek_or_drop(tick.gamma);
let theta = greek_or_drop(tick.theta);
let vega = greek_or_drop(tick.vega);
if iv.is_none() && delta.is_none() && gamma.is_none() && theta.is_none() && vega.is_none() {
return None;
}
Some(GreeksRow {
instrument,
iv,
delta,
gamma,
theta,
vega,
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: received,
})
}
fn pacing_window(mut instruments: Vec<Instrument>, cap: usize) -> Vec<Instrument> {
if instruments.len() <= cap {
return instruments;
}
instruments.sort_by_key(|instrument| instrument.key.strike);
let mid = instruments.len() / 2;
let start = mid
.saturating_sub(cap / 2)
.min(instruments.len().saturating_sub(cap));
instruments.into_iter().skip(start).take(cap).collect()
}
fn instruments_of(aliases: &AliasCatalog, provider: &ProviderId) -> Vec<Instrument> {
aliases
.instruments()
.filter(|instrument| &instrument.provider == provider)
.cloned()
.collect()
}
#[async_trait]
trait IbkrChainSource: Send + Sync {
async fn option_params(&self, underlying: &str) -> Result<Vec<RawOptionParams>, ProviderError>;
async fn expiry_detail(
&self,
underlying: &str,
expiration_ymd: &str,
strike: f64,
style: OptionStyle,
) -> Option<RawContractDetail>;
}
struct LiveChainSource {
client: ibapi::Client,
}
impl LiveChainSource {
async fn connect(adapter: &IbkrAdapter) -> Result<Self, ProviderError> {
let client = ibapi::Client::connect(&adapter.endpoint, adapter.client_id)
.await
.map_err(ibkr_error)?;
Ok(Self { client })
}
}
#[async_trait]
impl IbkrChainSource for LiveChainSource {
async fn option_params(&self, underlying: &str) -> Result<Vec<RawOptionParams>, ProviderError> {
let security_type = security_type_for(underlying);
let mut subscription = self
.client
.option_chain(underlying, security_type, 0)
.subscribe()
.await
.map_err(ibkr_error)?;
let chains = subscription
.collect_for(Duration::from_secs(OPTION_PARAMS_TIMEOUT_SECS))
.await;
Ok(chains.iter().map(map_option_chain).collect())
}
async fn expiry_detail(
&self,
underlying: &str,
expiration_ymd: &str,
strike: f64,
style: OptionStyle,
) -> Option<RawContractDetail> {
let right = match style {
OptionStyle::Call => OptionRight::Call,
OptionStyle::Put => OptionRight::Put,
};
let contract = Contract::option(underlying, expiration_ymd, strike, right);
let details = self.client.contract_details(&contract).await.ok()?;
details.first().map(map_contract_details)
}
}
#[derive(Debug, Clone)]
struct AssembledChain {
fetch: ChainFetch,
}
fn target_expiry(
expiration: &ExpirationDate,
received: DateTime<Utc>,
) -> Result<DateTime<Utc>, ProviderError> {
match expiration {
ExpirationDate::DateTime(dt) => Ok(*dt),
ExpirationDate::Days(days) => {
let seconds = (days.to_f64() * 86_400.0).round();
if !seconds.is_finite() {
return Err(normalize_err(NormalizeKind::UnparseableExpiry));
}
let seconds = seconds as i64;
let instant = received
.checked_add_signed(TimeDelta::seconds(seconds))
.ok_or_else(|| normalize_err(NormalizeKind::UnparseableExpiry))?;
let session = default_session();
let naive = instant
.date_naive()
.and_hms_opt(session.hour, session.minute, 0)
.ok_or_else(|| normalize_err(NormalizeKind::UnparseableExpiry))?;
local_to_utc(naive, session.zone)
.ok_or_else(|| normalize_err(NormalizeKind::UnparseableExpiry))
}
}
}
fn select_expiration<'a>(
params_list: &'a [RawOptionParams],
target: DateTime<Utc>,
session: &SessionClose,
) -> Option<(&'a RawOptionParams, String)> {
let target_date = utc_to_zone_date(target, session.zone);
for params in params_list {
for ymd in ¶ms.expirations {
if parse_ymd(ymd) == Some(target_date) {
return Some((params, ymd.clone()));
}
}
}
None
}
async fn compose_chain<S: IbkrChainSource + ?Sized>(
source: &S,
underlying: &str,
expiration: &ExpirationDate,
provider: &ProviderId,
received: DateTime<Utc>,
) -> Result<AssembledChain, ProviderError> {
let symbol = underlying.to_ascii_uppercase();
let target = target_expiry(expiration, received)?;
let params_list = source.option_params(&symbol).await?;
let session = default_session();
let (params, ymd) = match select_expiration(¶ms_list, target, &session) {
Some((params, ymd)) => (params.clone(), ymd),
None => return Err(no_chain(&symbol, target)),
};
let repr_strike = representative_strike(¶ms.strikes).unwrap_or(1.0);
let detail = source
.expiry_detail(&symbol, &ymd, repr_strike, OptionStyle::Call)
.await;
let session = detail
.as_ref()
.map(|d| session_from_zone_id(&d.time_zone_id))
.unwrap_or(session);
let last_trade_time = detail.as_ref().and_then(|d| d.last_trade_time.clone());
let expiry_utc =
resolve_expiry(last_trade_time.as_deref(), &ymd, &session).map_err(normalize_err)?;
assemble_chain(¶ms, &ymd, expiry_utc, &symbol, provider)
}
fn representative_strike(strikes: &[f64]) -> Option<f64> {
let mut sorted: Vec<f64> = strikes
.iter()
.copied()
.filter(|s| s.is_finite() && *s > 0.0)
.collect();
sorted.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
sorted.get(sorted.len() / 2).copied()
}
fn assemble_chain(
params: &RawOptionParams,
expiration_ymd: &str,
expiration_utc: DateTime<Utc>,
underlying: &str,
provider: &ProviderId,
) -> Result<AssembledChain, ProviderError> {
if params.strikes.len() > MAX_STRIKES {
return Err(normalize_err(NormalizeKind::LimitExceeded(STRIKE_CAP)));
}
let multiplier = params
.multiplier
.trim()
.parse::<u32>()
.ok()
.filter(|m| *m > 0)
.unwrap_or(DEFAULT_CONTRACT_MULTIPLIER);
let spec = ibkr_fingerprint(multiplier, ¶ms.trading_class);
let mut aliases = AliasCatalog::new();
let mut strikes: BTreeSet<Positive> = BTreeSet::new();
for raw_strike in ¶ms.strikes {
let Ok(strike) = strike_positive(*raw_strike) else {
continue;
};
let _ = strikes.insert(strike);
for style in [OptionStyle::Call, OptionStyle::Put] {
let native_symbol = ibkr_native_symbol(underlying, expiration_ymd, strike, style);
let key = InstrumentKey {
underlying: underlying.to_owned(),
expiration_utc,
strike,
style,
};
aliases.insert(Instrument {
key,
provider: provider.clone(),
native_symbol,
stream_symbol: None,
spec: spec.clone(),
});
}
}
if strikes.is_empty() {
return Err(no_chain(underlying, expiration_utc));
}
let ordered: Vec<Positive> = strikes.iter().copied().collect();
let spot = ordered
.get(ordered.len() / 2)
.copied()
.unwrap_or(Positive::ONE);
let mut chain = OptionChain::new(underlying, spot, expiration_utc.to_rfc3339(), None, None);
for strike in &ordered {
chain.add_option(
*strike,
None,
None,
None,
None,
Positive::ZERO,
None,
None,
None,
None,
None,
None,
);
}
let fetch = ChainFetch::new(
chain,
ExpirySource::new(underlying, expiration_utc, provider.clone()),
aliases,
);
Ok(AssembledChain { fetch })
}
fn normalize_err(kind: NormalizeKind) -> ProviderError {
ProviderError::Normalize { kind }
}
fn no_chain(underlying: &str, expiration_utc: DateTime<Utc>) -> ProviderError {
ProviderError::NoChain {
underlying: underlying.to_owned(),
expiration: expiration_utc.to_rfc3339(),
}
}
fn ibkr_error(err: ibapi::Error) -> ProviderError {
match err {
ibapi::Error::ConnectionFailed
| ibapi::Error::ConnectionRejected(_)
| ibapi::Error::ConnectionReset
| ibapi::Error::Io(_) => transport(TransportKind::Closed),
ibapi::Error::Parse(_, _, _)
| ibapi::Error::ParseInt(_)
| ibapi::Error::ParseTime(_)
| ibapi::Error::FromUtf8(_) => transport(TransportKind::Decode),
_ => transport(TransportKind::Closed),
}
}
fn transport(kind: TransportKind) -> ProviderError {
ProviderError::Transport(Box::new(TransportDetail::new(kind, None)))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TransportGone;
#[async_trait]
trait IbkrStreamTransport: Send {
async fn connect_and_subscribe(
&mut self,
instruments: &[Instrument],
) -> Result<(), TransportGone>;
async fn receive(&mut self) -> Result<RawTick, TransportGone>;
async fn poll(
&mut self,
underlying: &str,
expiration: &ExpirationDate,
received: DateTime<Utc>,
) -> Option<AssembledChain>;
}
struct LiveStreamTransport {
adapter: IbkrAdapter,
client: Option<ibapi::Client>,
tasks: Vec<JoinHandle<()>>,
receiver: Option<tokio::sync::mpsc::UnboundedReceiver<RawTick>>,
}
impl LiveStreamTransport {
fn new(adapter: IbkrAdapter) -> Self {
Self {
adapter,
client: None,
tasks: Vec::new(),
receiver: None,
}
}
fn teardown(&mut self) {
for task in self.tasks.drain(..) {
task.abort();
}
self.receiver = None;
self.client = None;
}
}
impl Drop for LiveStreamTransport {
fn drop(&mut self) {
self.teardown();
}
}
fn build_contract(instrument: &Instrument) -> Contract {
let key = &instrument.key;
let ymd = format_ymd_in_zone(key.expiration_utc, default_session().zone);
let right = match key.style {
OptionStyle::Call => OptionRight::Call,
OptionStyle::Put => OptionRight::Put,
};
let mut contract = Contract::option(&key.underlying, &ymd, key.strike.to_f64(), right);
contract.trading_class = instrument.spec.venue_product_code.clone();
contract.multiplier = instrument.spec.contract_multiplier.to_string();
contract
}
async fn drive_subscription(
mut subscription: ibapi::subscriptions::Subscription<TickTypes>,
symbol: String,
tx: tokio::sync::mpsc::UnboundedSender<RawTick>,
) {
let mut quote = QuoteAccumulator::default();
while let Some(item) = subscription.next().await {
let tick = match item {
Ok(SubscriptionItem::Data(tick)) => tick,
Ok(SubscriptionItem::Notice(_)) => continue,
Err(_) => break,
};
let sent = match tick {
TickTypes::Price(TickPrice {
tick_type, price, ..
}) => {
if quote.apply_price(tick_type, price) {
tx.send(RawTick::Quote(quote.snapshot(&symbol)))
} else {
Ok(())
}
}
TickTypes::Size(TickSize { tick_type, size }) => {
if quote.apply_size(tick_type, size) {
tx.send(RawTick::Quote(quote.snapshot(&symbol)))
} else {
Ok(())
}
}
TickTypes::OptionComputation(oc) => {
tx.send(RawTick::Greeks(map_option_computation(&symbol, &oc)))
}
_ => Ok(()),
};
if sent.is_err() {
break;
}
}
}
#[async_trait]
impl IbkrStreamTransport for LiveStreamTransport {
async fn connect_and_subscribe(
&mut self,
instruments: &[Instrument],
) -> Result<(), TransportGone> {
self.teardown();
if instruments.is_empty() {
return Err(TransportGone);
}
let client = ibapi::Client::connect(&self.adapter.endpoint, self.adapter.client_id)
.await
.map_err(|_| TransportGone)?;
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let mut tasks = Vec::new();
for instrument in instruments.iter().take(MAX_SUBSCRIPTIONS) {
let contract = build_contract(instrument);
let subscription = match client.market_data(&contract).subscribe().await {
Ok(subscription) => subscription,
Err(_) => continue,
};
let tx = tx.clone();
let symbol = instrument.native_symbol.clone();
tasks.push(tokio::spawn(drive_subscription(subscription, symbol, tx)));
}
if tasks.is_empty() {
return Err(TransportGone);
}
self.client = Some(client);
self.tasks = tasks;
self.receiver = Some(rx);
Ok(())
}
async fn receive(&mut self) -> Result<RawTick, TransportGone> {
match self.receiver.as_mut() {
Some(receiver) => match receiver.recv().await {
Some(tick) => Ok(tick),
None => Err(TransportGone),
},
None => Err(TransportGone),
}
}
async fn poll(
&mut self,
underlying: &str,
expiration: &ExpirationDate,
received: DateTime<Utc>,
) -> Option<AssembledChain> {
let source = LiveChainSource::connect(&self.adapter).await.ok()?;
compose_chain(&source, underlying, expiration, &self.adapter.id, received)
.await
.ok()
}
}
enum StreamExit {
Reconnect,
Shutdown,
}
async fn run_reconnect_loop<T: IbkrStreamTransport>(
mut transport: T,
id: ProviderId,
underlying: String,
expiration_utc: DateTime<Utc>,
mut aliases: AliasCatalog,
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, &underlying, expiration_utc,
&mut aliases, &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) => {}
}
}
}
#[allow(clippy::too_many_arguments)]
async fn connect_stream_once<T: IbkrStreamTransport>(
transport: &mut T,
id: &ProviderId,
underlying: &str,
expiration_utc: DateTime<Utc>,
aliases: &mut AliasCatalog,
sink: &mut MarketUpdateSink,
cancel: &CancellationToken,
attempt: &mut u32,
) -> StreamExit {
let window = pacing_window(instruments_of(aliases, id), MAX_SUBSCRIPTIONS);
let subscribed = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
result = transport.connect_and_subscribe(&window) => result,
};
if subscribed.is_err() {
return StreamExit::Reconnect;
}
*attempt = 0;
if go_live_and_backfill(
transport,
id,
underlying,
expiration_utc,
aliases,
sink,
cancel,
)
.await
== SendState::Closed
{
return StreamExit::Shutdown;
}
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,
};
let update = match event {
RawTick::Quote(tick) => {
quote_update(&tick, aliases, id, now_utc()).map(MarketUpdate::Quote)
}
RawTick::Greeks(tick) => {
greeks_row(&tick, aliases, id, now_utc()).map(MarketUpdate::Greeks)
}
};
let Some(update) = update else {
continue;
};
let sent = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
outcome = sink.send(update) => outcome,
};
if sent == SendState::Closed {
return StreamExit::Shutdown;
}
}
}
async fn go_live_and_backfill<T: IbkrStreamTransport>(
transport: &mut T,
id: &ProviderId,
underlying: &str,
expiration_utc: DateTime<Utc>,
aliases: &mut AliasCatalog,
sink: &mut MarketUpdateSink,
cancel: &CancellationToken,
) -> SendState {
let live = MarketUpdate::Health(id.clone(), StreamHealth::Live);
let sent = tokio::select! {
biased;
() = cancel.cancelled() => return SendState::Open,
outcome = sink.send(live) => outcome,
};
if sent == SendState::Closed {
return SendState::Closed;
}
let expiration = ExpirationDate::DateTime(expiration_utc);
let composed = tokio::select! {
biased;
() = cancel.cancelled() => return SendState::Open,
result = transport.poll(underlying, &expiration, now_utc()) => result,
};
let Some(composed) = composed else {
return SendState::Open;
};
*aliases = composed.fetch.aliases.clone();
let snapshot = MarketUpdate::Chain(chain_snapshot(&composed.fetch, now_utc()));
tokio::select! {
biased;
() = cancel.cancelled() => SendState::Open,
outcome = sink.send(snapshot) => outcome,
}
}
fn chain_snapshot(fetch: &ChainFetch, last_poll: DateTime<Utc>) -> ChainSnapshot {
ChainSnapshot {
chain_key: (
fetch.expiry_source.provider.clone(),
fetch.expiry_source.underlying.clone(),
fetch.expiry_source.expiration_utc,
),
chain: fetch.chain.clone(),
aliases: fetch.aliases.clone(),
source: ChainSource::Merged,
health: StreamHealth::Live,
last_full_poll: Some(last_poll),
}
}
#[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::HashMap as StdHashMap;
use std::sync::Arc;
use std::sync::Mutex as StdMutex;
use tokio::sync::mpsc;
use super::*;
use crate::chain::{ChainStore, MergeOutcome};
#[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_rfc3339(s: &str) -> DateTime<Utc> {
match DateTime::parse_from_rfc3339(s) {
Ok(dt) => dt.with_timezone(&Utc),
Err(e) => panic!("invalid rfc3339 `{s}`: {e}"),
}
}
#[track_caller]
fn date(y: i32, m: u32, d: u32) -> NaiveDate {
match NaiveDate::from_ymd_opt(y, m, d) {
Some(date) => date,
None => panic!("invalid test date {y}-{m}-{d}"),
}
}
struct MapEnv(StdHashMap<String, String>);
impl EnvSource for MapEnv {
fn get(&self, key: &str) -> Option<String> {
self.0.get(key).cloned()
}
}
fn endpoint_env() -> MapEnv {
let mut env = StdHashMap::new();
let _ = env.insert(
"CHAINVIEW_IBKR_ENDPOINT".to_owned(),
"127.0.0.1:7497".to_owned(),
);
MapEnv(env)
}
#[track_caller]
fn sample_adapter() -> IbkrAdapter {
match IbkrAdapter::from_env(&endpoint_env()) {
Ok(adapter) => adapter,
Err(e) => panic!("from_env should succeed with the endpoint 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,
)
}
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 OPTION_PARAMS_FIXTURE: &str =
include_str!("../../tests/fixtures/ibkr/sec_def_opt_params_spx.json");
#[derive(serde::Deserialize)]
struct FixtureParams {
trading_class: String,
multiplier: String,
exchange: String,
underlying_conid: i32,
expirations: Vec<String>,
strikes: Vec<f64>,
}
#[derive(serde::Deserialize)]
struct FixtureDetail {
expiration_ymd: String,
last_trade_time: String,
time_zone_id: String,
}
#[derive(serde::Deserialize)]
struct Fixture {
params: FixtureParams,
contract_detail: FixtureDetail,
}
#[track_caller]
fn load_fixture() -> Fixture {
match serde_json::from_str(OPTION_PARAMS_FIXTURE) {
Ok(fixture) => fixture,
Err(e) => panic!("the ibkr fixture must deserialize: {e}"),
}
}
#[track_caller]
fn fixture_params() -> RawOptionParams {
let fixture = load_fixture();
RawOptionParams {
trading_class: fixture.params.trading_class,
multiplier: fixture.params.multiplier,
exchange: fixture.params.exchange,
expirations: fixture.params.expirations,
strikes: fixture.params.strikes,
underlying_conid: fixture.params.underlying_conid,
}
}
#[track_caller]
fn fixture_detail() -> (RawContractDetail, String) {
let fixture = load_fixture();
(
RawContractDetail {
last_trade_time: Some(fixture.contract_detail.last_trade_time),
time_zone_id: fixture.contract_detail.time_zone_id,
},
fixture.contract_detail.expiration_ymd,
)
}
#[track_caller]
fn fixture_expiry_utc() -> DateTime<Utc> {
let (detail, ymd) = fixture_detail();
let session = session_from_zone_id(&detail.time_zone_id);
match resolve_expiry(detail.last_trade_time.as_deref(), &ymd, &session) {
Ok(dt) => dt,
Err(e) => panic!("fixture expiry must resolve: {e}"),
}
}
#[track_caller]
fn fixture_assembled() -> AssembledChain {
let params = fixture_params();
let (_, ymd) = fixture_detail();
match assemble_chain(¶ms, &ymd, fixture_expiry_utc(), "SPX", &pid("ibkr")) {
Ok(a) => a,
Err(e) => panic!("assembly should succeed, got: {e}"),
}
}
#[test]
fn test_ibkr_id_is_valid_and_reserved() {
let id = ibkr_provider_id();
assert_eq!(id.as_str(), "ibkr");
assert!(id.is_reserved());
assert!(ProviderId::new(IBKR_ID).is_ok());
}
#[test]
fn test_ibkr_capabilities_match_section_8_row() {
let caps = ibkr_capabilities();
assert_eq!(caps.chain, ChainCapability::Assemble);
assert!(
!caps.depth,
"IBKR claims no option depth (no recorded fixture)"
);
assert_eq!(caps.greeks, GreeksCapability::Provided);
assert_eq!(
caps.option_stream,
OptionStreamCapability::ChainQuotes { verified: false }
);
assert!(
!caps.underlying_stream,
"IBKR subscribes only option contracts (no underlying fold; #40/#41 honesty)"
);
assert_eq!(
caps.chain_poll,
ChainPollCapability::Poll {
interval_hint_secs: REFRESH_HINT_SECS
}
);
assert!(!caps.trades_tape, "no public trade tape on this path");
assert_eq!(
caps.auth,
AuthKind::None,
"the TWS/Gateway holds the session"
);
}
#[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(), "ibkr");
assert_eq!(adapter.capabilities().chain, ChainCapability::Assemble);
assert_eq!(adapter.capabilities().greeks, GreeksCapability::Provided);
assert_eq!(adapter.capabilities().auth, AuthKind::None);
}
#[test]
fn test_from_env_reads_endpoint_and_default_client_id() {
let adapter = sample_adapter();
assert_eq!(adapter.endpoint, "127.0.0.1:7497");
assert_eq!(adapter.client_id, DEFAULT_CLIENT_ID);
}
#[test]
fn test_from_env_reads_client_id_when_set() {
let mut env = StdHashMap::new();
let _ = env.insert("CHAINVIEW_IBKR_ENDPOINT".to_owned(), "gw:4002".to_owned());
let _ = env.insert("CHAINVIEW_IBKR_CLIENT_ID".to_owned(), "42".to_owned());
match IbkrAdapter::from_env(&MapEnv(env)) {
Ok(adapter) => {
assert_eq!(adapter.endpoint, "gw:4002");
assert_eq!(adapter.client_id, 42);
}
Err(e) => panic!("from_env should succeed: {e}"),
}
}
#[test]
fn test_from_env_missing_endpoint_is_missing_credential() {
match IbkrAdapter::from_env(&MapEnv(StdHashMap::new())) {
Err(ConfigError::MissingCredential(id)) => assert_eq!(id.as_str(), "ibkr"),
other => panic!("expected MissingCredential(ibkr), got: {other:?}"),
}
}
#[test]
fn test_from_env_malformed_endpoint_is_invalid_value() {
let mut env = StdHashMap::new();
let _ = env.insert(
"CHAINVIEW_IBKR_ENDPOINT".to_owned(),
"not-a-socket".to_owned(),
);
match IbkrAdapter::from_env(&MapEnv(env)) {
Err(ConfigError::InvalidValue { field, .. }) => assert_eq!(field, "ibkr endpoint"),
other => panic!("expected InvalidValue on ibkr endpoint, got: {other:?}"),
}
}
#[test]
fn test_from_env_malformed_client_id_is_invalid_value() {
let mut env = StdHashMap::new();
let _ = env.insert(
"CHAINVIEW_IBKR_ENDPOINT".to_owned(),
"127.0.0.1:7497".to_owned(),
);
let _ = env.insert("CHAINVIEW_IBKR_CLIENT_ID".to_owned(), "abc".to_owned());
match IbkrAdapter::from_env(&MapEnv(env)) {
Err(ConfigError::InvalidValue { field, .. }) => assert_eq!(field, "ibkr client id"),
other => panic!("expected InvalidValue on ibkr client id, got: {other:?}"),
}
}
#[test]
fn test_endpoint_validation() {
assert!(is_valid_endpoint("127.0.0.1:7497"));
assert!(is_valid_endpoint("gateway.local:4002"));
assert!(!is_valid_endpoint("no-port"));
assert!(!is_valid_endpoint(":7497"));
assert!(!is_valid_endpoint("host:"));
assert!(!is_valid_endpoint("host:notaport"));
assert!(!is_valid_endpoint("host:99999999"));
assert!(!is_valid_endpoint("127.0.0.1:7497/a:80"));
assert!(!is_valid_endpoint("user:pass@host:7497"));
assert!(is_valid_endpoint("[::1]:7497"));
}
#[test]
fn test_normalize_endpoint_accepts_both_grammars() {
assert_eq!(
normalize_endpoint("tcp://127.0.0.1:7497").as_deref(),
Some("127.0.0.1:7497")
);
assert_eq!(
normalize_endpoint("ibkr+tws://gateway.local:4002").as_deref(),
Some("gateway.local:4002")
);
assert_eq!(
normalize_endpoint("127.0.0.1:7497").as_deref(),
Some("127.0.0.1:7497")
);
assert_eq!(
normalize_endpoint("tcp://[::1]:7497").as_deref(),
Some("[::1]:7497")
);
}
#[test]
fn test_normalize_endpoint_rejects_malformed() {
assert_eq!(normalize_endpoint("tcp://"), None);
assert_eq!(normalize_endpoint("://127.0.0.1:7497"), None);
assert_eq!(normalize_endpoint("1tcp://127.0.0.1:7497"), None);
assert_eq!(normalize_endpoint("tcp://127.0.0.1"), None);
assert_eq!(normalize_endpoint("tcp://127.0.0.1:7497/path"), None);
assert_eq!(normalize_endpoint("tcp://host:notaport"), None);
assert_eq!(normalize_endpoint("no-port"), None);
assert_eq!(normalize_endpoint("tcp://127.0.0.1:7497/a:80"), None);
assert_eq!(normalize_endpoint("tcp://user:pass@host:7497"), None);
assert_eq!(normalize_endpoint("tcp://tcp://127.0.0.1:7497"), None);
}
#[test]
fn test_from_env_accepts_absolute_url_endpoint() {
let mut env = StdHashMap::new();
let _ = env.insert(
"CHAINVIEW_IBKR_ENDPOINT".to_owned(),
"tcp://127.0.0.1:7497".to_owned(),
);
match IbkrAdapter::from_env(&MapEnv(env)) {
Ok(adapter) => assert_eq!(adapter.endpoint, "127.0.0.1:7497"),
Err(e) => panic!("from_env should accept the absolute-URL form: {e}"),
}
}
#[derive(Clone, Default)]
struct LogBuffer(Arc<StdMutex<Vec<u8>>>);
impl LogBuffer {
fn contents(&self) -> String {
match self.0.lock() {
Ok(bytes) => String::from_utf8_lossy(&bytes).into_owned(),
Err(poisoned) => String::from_utf8_lossy(&poisoned.into_inner()).into_owned(),
}
}
}
impl std::io::Write for LogBuffer {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
if let Ok(mut bytes) = self.0.lock() {
bytes.extend_from_slice(buf);
}
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for LogBuffer {
type Writer = LogBuffer;
fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}
const CONTROL_CANARY: &str = "chainview-ibkr-canary-7c1d-present";
const SECRET_SHAPED_ENDPOINT: &str = "10.9.8.7:65001";
#[test]
fn test_construction_and_errors_never_log_secret_material() {
let logs = LogBuffer::default();
let subscriber = tracing_subscriber::fmt()
.with_max_level(tracing::Level::TRACE)
.with_ansi(false)
.with_writer(logs.clone())
.finish();
let _guard = tracing::subscriber::set_default(subscriber);
tracing::debug!("{CONTROL_CANARY}");
let mut env = StdHashMap::new();
let _ = env.insert(
"CHAINVIEW_IBKR_ENDPOINT".to_owned(),
SECRET_SHAPED_ENDPOINT.to_owned(),
);
let _ = env.insert("CHAINVIEW_IBKR_CLIENT_ID".to_owned(), "31337".to_owned());
let adapter = match IbkrAdapter::from_env(&MapEnv(env)) {
Ok(adapter) => adapter,
Err(e) => panic!("from_env should succeed: {e}"),
};
let _ = &adapter;
let err = ibkr_error(ibapi::Error::ConnectionFailed);
tracing::debug!(error = %err, "ibkr error rendered");
let decode = ibkr_error(ibapi::Error::NotImplemented);
tracing::debug!(error = %decode, "ibkr error rendered");
let output = logs.contents();
assert!(output.contains(CONTROL_CANARY), "sink must capture content");
assert!(
!output.contains(SECRET_SHAPED_ENDPOINT),
"the endpoint leaked into logs:\n{output}"
);
assert!(
!output.contains("31337"),
"the client id leaked into logs:\n{output}"
);
assert!(
!format!("{err} {err:?}").contains(SECRET_SHAPED_ENDPOINT),
"a ProviderError must never carry the endpoint"
);
}
#[test]
fn test_expiry_date_only_resolves_via_us_session_close() {
let session = default_session();
match resolve_expiry(None, "20260717", &session) {
Ok(utc) => assert_eq!(utc.to_rfc3339(), "2026-07-17T20:00:00+00:00"),
Err(e) => panic!("date-only US expiry should resolve, got: {e}"),
}
}
#[test]
fn test_expiry_authoritative_last_trade_time_wins() {
let session = default_session();
let resolved = match resolve_expiry(Some("16:15:00"), "20260717", &session) {
Ok(dt) => dt,
Err(e) => panic!("timestamped expiry should resolve, got: {e}"),
};
assert_eq!(resolved.to_rfc3339(), "2026-07-17T20:15:00+00:00");
let date_only = match resolve_expiry(None, "20260717", &session) {
Ok(dt) => dt,
Err(e) => panic!("date-only should resolve, got: {e}"),
};
assert_ne!(
resolved, date_only,
"the authoritative last-trade time must win"
);
}
#[test]
fn test_expiry_authoritative_date_and_time_form() {
let session = default_session();
match resolve_expiry(Some("20260717-16:15:00"), "20260717", &session) {
Ok(utc) => assert_eq!(utc.to_rfc3339(), "2026-07-17T20:15:00+00:00"),
Err(e) => panic!("date+time form should resolve, got: {e}"),
}
}
#[test]
fn test_expiry_ambiguous_and_unparseable_are_normalize_errors() {
let session = default_session();
assert_eq!(
resolve_expiry(None, "notadate", &session),
Err(NormalizeKind::UnparseableExpiry)
);
assert_eq!(
resolve_expiry(Some("garbage"), "20260717", &session),
Err(NormalizeKind::UnparseableExpiry)
);
assert_eq!(
resolve_expiry(None, "2026071", &session),
Err(NormalizeKind::UnparseableExpiry)
);
}
#[test]
fn test_dst_boundaries() {
assert!(us_eastern_dst(date(2026, 3, 8)));
assert!(!us_eastern_dst(date(2026, 3, 7)));
assert!(!us_eastern_dst(date(2026, 11, 1)));
assert!(us_eastern_dst(date(2026, 10, 31)));
assert!(uk_dst(date(2026, 3, 29)));
assert!(!uk_dst(date(2026, 3, 28)));
}
#[test]
fn test_parse_ymd_forms() {
assert_eq!(parse_ymd("20260717"), Some(date(2026, 7, 17)));
assert_eq!(parse_ymd("2026-07-17"), None);
assert_eq!(parse_ymd("2026071"), None);
assert_eq!(parse_ymd("2026071X"), None);
}
#[test]
fn test_strike_positive_parses_and_rejects() {
match strike_positive(7500.0) {
Ok(strike) => assert_eq!(strike, pos(7500.0)),
Err(e) => panic!("7500 should parse, got: {e}"),
}
assert_eq!(
strike_positive(0.0),
Err(NormalizeKind::OutOfRange("strike"))
);
assert_eq!(
strike_positive(-5.0),
Err(NormalizeKind::OutOfRange("strike"))
);
assert_eq!(
strike_positive(f64::NAN),
Err(NormalizeKind::OutOfRange("strike"))
);
}
#[test]
fn test_normalize_quote_rules() {
match normalize_quote(Some(0.0), Some(1.0)) {
Ok(q) => {
assert_eq!(q.bid, Some(Positive::ZERO));
assert_eq!(q.ask, Some(pos(1.0)));
}
Err(e) => panic!("zero bid valid, got: {e}"),
}
assert_eq!(
normalize_quote(Some(5.0), Some(3.0)),
Err(NormalizeKind::OutOfRange("ask"))
);
match normalize_quote(Some(f64::NAN), Some(2.0)) {
Ok(q) => {
assert_eq!(q.bid, None);
assert_eq!(q.ask, Some(pos(2.0)));
}
Err(e) => panic!("NaN bid drops only that field, got: {e}"),
}
}
#[test]
fn test_greek_or_drop_guards_non_finite() {
assert_eq!(greek_or_drop(Some(0.55)), Some(Decimal::new(55, 2)));
assert_eq!(greek_or_drop(Some(-0.45)), Some(Decimal::new(-45, 2)));
assert_eq!(greek_or_drop(Some(f64::NAN)), None);
assert_eq!(greek_or_drop(Some(f64::INFINITY)), None);
assert_eq!(greek_or_drop(None), None);
}
#[test]
fn test_map_option_chain_from_real_ibapi_struct() {
let chain = IbOptionChain {
underlying_contract_id: 416_904,
trading_class: "SPX".to_owned(),
multiplier: "100".to_owned(),
exchange: "SMART".to_owned(),
expirations: vec!["20260717".to_owned(), "20260821".to_owned()],
strikes: vec![7400.0, 7500.0, 7600.0],
};
let raw = map_option_chain(&chain);
assert_eq!(raw.trading_class, "SPX");
assert_eq!(raw.multiplier, "100");
assert_eq!(raw.expirations, vec!["20260717", "20260821"]);
assert_eq!(raw.strikes, vec![7400.0, 7500.0, 7600.0]);
assert_eq!(raw.underlying_conid, 416_904);
}
#[test]
fn test_map_option_computation_from_real_ibapi_struct() {
let oc = OptionComputation {
implied_volatility: Some(0.25),
delta: Some(0.55),
gamma: Some(0.01),
vega: Some(1.2),
theta: Some(-0.5),
underlying_price: Some(7500.0),
..Default::default()
};
let raw = map_option_computation("SPX-20260717-7500-C", &oc);
assert_eq!(raw.iv, Some(0.25));
assert_eq!(raw.delta, Some(0.55));
assert_eq!(raw.gamma, Some(0.01));
assert_eq!(raw.vega, Some(1.2));
assert_eq!(raw.theta, Some(-0.5));
}
fn one_leg_aliases() -> (AliasCatalog, String) {
let assembled = fixture_assembled();
let symbol = match assembled
.fetch
.aliases
.instruments()
.find(|i| i.provider == pid("ibkr"))
{
Some(instrument) => instrument.native_symbol.clone(),
None => panic!("the fixture must carry at least one leg"),
};
(assembled.fetch.aliases, symbol)
}
#[test]
fn test_greeks_row_is_tagged_provider_and_checked_into_decimal() {
let (aliases, symbol) = one_leg_aliases();
let tick = RawGreeksTick {
symbol,
iv: Some(0.25),
delta: Some(0.55),
gamma: Some(0.01),
theta: Some(-0.5),
vega: Some(1.2),
};
match greeks_row(&tick, &aliases, &pid("ibkr"), now_utc()) {
Some(row) => {
assert_eq!(row.origin, GreeksOrigin::Provider);
assert_eq!(row.iv, Some(pos(0.25)));
assert_eq!(row.delta, Some(Decimal::new(55, 2)));
assert_eq!(row.theta, Some(Decimal::new(-5, 1)));
assert!(row.rho.is_none(), "IBKR OptionComputation carries no rho");
}
None => panic!("expected a Provider GreeksRow"),
}
}
#[test]
fn test_greeks_row_nan_is_guarded() {
let (aliases, symbol) = one_leg_aliases();
let tick = RawGreeksTick {
symbol,
iv: Some(f64::NAN),
delta: Some(f64::INFINITY),
gamma: Some(0.01),
theta: None,
vega: None,
};
match greeks_row(&tick, &aliases, &pid("ibkr"), now_utc()) {
Some(row) => {
assert!(row.iv.is_none(), "NaN IV dropped");
assert!(row.delta.is_none(), "Inf delta dropped");
assert_eq!(row.gamma, Some(Decimal::new(1, 2)));
}
None => panic!("a partially-valid computation still yields a row"),
}
}
#[test]
fn test_crossed_or_stale_computation_leaves_leg_none() {
let (aliases, symbol) = one_leg_aliases();
let tick = RawGreeksTick {
symbol,
iv: Some(f64::NAN),
delta: Some(f64::NAN),
gamma: None,
theta: None,
vega: None,
};
assert!(
greeks_row(&tick, &aliases, &pid("ibkr"), now_utc()).is_none(),
"a stale computation whose analytics all drop leaves the leg None"
);
}
#[test]
fn test_greeks_row_unknown_symbol_is_dropped() {
let (aliases, _) = one_leg_aliases();
let tick = RawGreeksTick {
symbol: "SPX-20260717-99999-C".to_owned(),
iv: Some(0.25),
delta: Some(0.5),
gamma: None,
theta: None,
vega: None,
};
assert!(
greeks_row(&tick, &aliases, &pid("ibkr"), now_utc()).is_none(),
"a tick for an unknown symbol must be dropped, never keyed"
);
}
#[test]
fn test_quote_accumulator_folds_partial_ticks() {
let mut acc = QuoteAccumulator::default();
assert!(acc.apply_price(TickType::Bid, 12.5));
assert!(acc.apply_price(TickType::Ask, 13.5));
assert!(acc.apply_size(TickType::BidSize, 4.0));
assert!(
!acc.apply_price(TickType::Open, 10.0),
"non-quote tick ignored"
);
let snapshot = acc.snapshot("SPX-20260717-7500-C");
assert_eq!(snapshot.bid, Some(12.5));
assert_eq!(snapshot.ask, Some(13.5));
assert_eq!(snapshot.bid_size, Some(4.0));
assert_eq!(snapshot.ask_size, None);
}
#[test]
fn test_quote_update_from_accumulated_tick() {
let (aliases, symbol) = one_leg_aliases();
let tick = RawQuoteTick {
symbol,
bid: Some(12.5),
ask: Some(13.5),
last: Some(13.0),
bid_size: Some(4.0),
ask_size: Some(6.0),
};
match quote_update(&tick, &aliases, &pid("ibkr"), now_utc()) {
Some(update) => {
assert_eq!(update.bid, Some(pos(12.5)));
assert_eq!(update.ask, Some(pos(13.5)));
assert_eq!(update.last, Some(pos(13.0)));
}
None => panic!("the quote should normalize"),
}
}
#[test]
fn test_assemble_chain_from_recorded_shape() {
let assembled = fixture_assembled();
let strikes: Vec<Positive> = assembled
.fetch
.chain
.options
.iter()
.map(|o| o.strike_price)
.collect();
assert_eq!(strikes, vec![pos(7400.0), pos(7500.0), pos(7600.0)]);
let leg_count = assembled
.fetch
.aliases
.instruments()
.filter(|i| i.provider == pid("ibkr"))
.count();
assert_eq!(leg_count, 6);
for instrument in assembled.fetch.aliases.instruments() {
assert_eq!(instrument.spec.contract_multiplier, 100);
assert_eq!(instrument.spec.venue_product_code, "SPX");
}
assert_eq!(
assembled.fetch.expiry_source.expiration_utc,
utc_rfc3339("2026-07-17T20:15:00+00:00")
);
}
#[test]
fn test_assemble_chain_strike_cap_exceeded() {
let params = RawOptionParams {
trading_class: "SPX".to_owned(),
multiplier: "100".to_owned(),
exchange: "SMART".to_owned(),
expirations: vec!["20260717".to_owned()],
strikes: vec![1.0; MAX_STRIKES + 1],
underlying_conid: 1,
};
match assemble_chain(
¶ms,
"20260717",
fixture_expiry_utc(),
"SPX",
&pid("ibkr"),
) {
Err(ProviderError::Normalize {
kind: NormalizeKind::LimitExceeded(cap),
}) => assert_eq!(cap, STRIKE_CAP),
other => panic!("expected LimitExceeded, got: {other:?}"),
}
}
fn window_instruments(count: usize) -> Vec<Instrument> {
let spec = ibkr_fingerprint(100, "SPX");
(0..count)
.map(|i| {
let strike = pos(1000.0 + i as f64);
Instrument {
key: InstrumentKey {
underlying: "SPX".to_owned(),
expiration_utc: fixture_expiry_utc(),
strike,
style: OptionStyle::Call,
},
provider: pid("ibkr"),
native_symbol: format!("SPX-20260717-{strike}-C"),
stream_symbol: None,
spec: spec.clone(),
}
})
.collect()
}
#[test]
fn test_pacing_window_never_exceeds_cap() {
for count in [0usize, 1, 50, MAX_SUBSCRIPTIONS, MAX_SUBSCRIPTIONS + 1, 300] {
let window = pacing_window(window_instruments(count), MAX_SUBSCRIPTIONS);
assert!(
window.len() <= MAX_SUBSCRIPTIONS,
"count {count}: window {} exceeds the pacing cap {MAX_SUBSCRIPTIONS}",
window.len()
);
assert_eq!(window.len(), count.min(MAX_SUBSCRIPTIONS));
}
}
#[test]
fn test_pacing_window_is_atm_centered() {
let window = pacing_window(window_instruments(300), MAX_SUBSCRIPTIONS);
let strikes: Vec<f64> = window.iter().map(|i| i.key.strike.to_f64()).collect();
let lo = strikes.iter().copied().fold(f64::INFINITY, f64::min);
let hi = strikes.iter().copied().fold(f64::NEG_INFINITY, f64::max);
assert!(lo > 1000.0, "the low wing is excluded, got {lo}");
assert!(hi < 1299.0, "the high wing is excluded, got {hi}");
assert!(
lo <= 1150.0 && hi >= 1150.0,
"the window straddles the median"
);
}
struct MockChainSource {
params: Vec<RawOptionParams>,
detail: Option<RawContractDetail>,
}
#[async_trait]
impl IbkrChainSource for MockChainSource {
async fn option_params(
&self,
_underlying: &str,
) -> Result<Vec<RawOptionParams>, ProviderError> {
Ok(self.params.clone())
}
async fn expiry_detail(
&self,
_underlying: &str,
_expiration_ymd: &str,
_strike: f64,
_style: OptionStyle,
) -> Option<RawContractDetail> {
self.detail.clone()
}
}
#[tokio::test]
async fn test_discover_returns_candidate_underlyings() {
let adapter = sample_adapter();
match adapter.discover().await {
Ok(refs) => {
let names: Vec<String> = refs.into_iter().map(|r| r.underlying).collect();
assert!(names.contains(&"SPX".to_owned()));
assert!(names.contains(&"SPY".to_owned()));
assert_eq!(names.len(), CANDIDATE_UNDERLYINGS.len());
}
Err(e) => panic!("discover should succeed, got: {e}"),
}
}
#[tokio::test]
async fn test_compose_chain_over_mock_source() {
let (detail, _) = fixture_detail();
let source = MockChainSource {
params: vec![fixture_params()],
detail: Some(detail),
};
let expiration = ExpirationDate::DateTime(fixture_expiry_utc());
let received = utc_rfc3339("2026-07-01T14:00:00+00:00");
match compose_chain(&source, "spx", &expiration, &pid("ibkr"), received).await {
Ok(assembled) => {
assert_eq!(assembled.fetch.chain.symbol, "SPX");
assert_eq!(assembled.fetch.chain.options.len(), 3);
assert_eq!(
assembled.fetch.expiry_source.expiration_utc,
fixture_expiry_utc()
);
}
Err(e) => panic!("compose_chain should succeed, got: {e}"),
}
}
#[tokio::test]
async fn test_compose_chain_mismatched_expiry_is_no_chain() {
let source = MockChainSource {
params: vec![fixture_params()],
detail: None,
};
let other = ExpirationDate::DateTime(utc_rfc3339("2027-01-15T20:00:00+00:00"));
let received = utc_rfc3339("2026-07-01T14:00:00+00:00");
match compose_chain(&source, "spx", &other, &pid("ibkr"), received).await {
Err(ProviderError::NoChain { underlying, .. }) => assert_eq!(underlying, "SPX"),
other => panic!("expected NoChain for a mismatched expiry, got {other:?}"),
}
}
struct MockTransport {
attempts: Vec<Vec<RawTick>>,
attempt_idx: usize,
cursor: usize,
backfill: Option<AssembledChain>,
connects: Arc<StdMutex<u32>>,
max_window: Arc<StdMutex<usize>>,
cancel: CancellationToken,
}
#[async_trait]
impl IbkrStreamTransport for MockTransport {
async fn connect_and_subscribe(
&mut self,
instruments: &[Instrument],
) -> Result<(), TransportGone> {
if let Ok(mut count) = self.connects.lock() {
*count += 1;
}
if let Ok(mut max) = self.max_window.lock() {
*max = (*max).max(instruments.len());
}
self.cursor = 0;
Ok(())
}
async fn receive(&mut self) -> Result<RawTick, TransportGone> {
let Some(events) = self.attempts.get(self.attempt_idx) else {
self.cancel.cancel();
return Err(TransportGone);
};
if let Some(event) = events.get(self.cursor) {
self.cursor = self.cursor.saturating_add(1);
return Ok(event.clone());
}
self.attempt_idx = self.attempt_idx.saturating_add(1);
self.cursor = 0;
if self.attempt_idx >= self.attempts.len() {
self.cancel.cancel();
}
Err(TransportGone)
}
async fn poll(
&mut self,
_underlying: &str,
_expiration: &ExpirationDate,
_received: DateTime<Utc>,
) -> Option<AssembledChain> {
self.backfill.clone()
}
}
struct PendingTransport;
#[async_trait]
impl IbkrStreamTransport for PendingTransport {
async fn connect_and_subscribe(
&mut self,
_instruments: &[Instrument],
) -> Result<(), TransportGone> {
Ok(())
}
async fn receive(&mut self) -> Result<RawTick, TransportGone> {
std::future::pending::<()>().await;
Err(TransportGone)
}
async fn poll(
&mut self,
_underlying: &str,
_expiration: &ExpirationDate,
_received: DateTime<Utc>,
) -> Option<AssembledChain> {
None
}
}
#[tokio::test(start_paused = true)]
async fn test_reconnect_loop_emits_quotes_greeks_backfill_without_panic() {
let assembled = fixture_assembled();
let cancel = CancellationToken::new();
let connects = Arc::new(StdMutex::new(0));
let max_window = Arc::new(StdMutex::new(0));
let symbol = match assembled
.fetch
.aliases
.instruments()
.find(|i| i.provider == pid("ibkr"))
{
Some(instrument) => instrument.native_symbol.clone(),
None => panic!("the fixture must carry at least one leg"),
};
let quote = RawTick::Quote(RawQuoteTick {
symbol: symbol.clone(),
bid: Some(12.5),
ask: Some(13.5),
last: None,
bid_size: None,
ask_size: None,
});
let greeks = RawTick::Greeks(RawGreeksTick {
symbol,
iv: Some(0.25),
delta: Some(0.55),
gamma: Some(0.01),
theta: Some(-0.5),
vega: Some(1.2),
});
let transport = MockTransport {
attempts: vec![vec![quote.clone(), greeks.clone()], vec![quote, greeks]],
attempt_idx: 0,
cursor: 0,
backfill: Some(assembled.clone()),
connects: Arc::clone(&connects),
max_window: Arc::clone(&max_window),
cancel: cancel.clone(),
};
let (sink, mut rx_control, mut rx_coalesced) = test_sink(256);
run_reconnect_loop(
transport,
pid("ibkr"),
"SPX".to_owned(),
fixture_expiry_utc(),
assembled.fetch.aliases.clone(),
sink,
cancel,
)
.await;
assert_eq!(
*connects.lock().unwrap_or_else(|e| e.into_inner()),
2,
"connected twice"
);
assert!(
*max_window.lock().unwrap_or_else(|e| e.into_inner()) <= MAX_SUBSCRIPTIONS,
"the subscribe window never exceeds the pacing cap"
);
let control = drain(&mut rx_control);
assert!(
control.iter().any(|u| matches!(u, MarketUpdate::Chain(_))),
"a Chain backfill was emitted"
);
assert!(
control
.iter()
.any(|u| matches!(u, MarketUpdate::Health(_, StreamHealth::Live))),
"a Live health was emitted"
);
assert!(
control.iter().any(|u| matches!(
u,
MarketUpdate::Health(_, StreamHealth::Reconnecting { .. })
)),
"a Reconnecting health was emitted on the stream drop"
);
let coalesced = drain(&mut rx_coalesced);
assert!(
coalesced
.iter()
.any(|u| matches!(u, MarketUpdate::Quote(_))),
"the normalized quote reached the coalesced channel"
);
assert!(
coalesced
.iter()
.any(|u| matches!(u, MarketUpdate::Greeks(_))),
"the native Greeks reached the coalesced channel"
);
}
#[tokio::test(start_paused = true)]
async fn test_reconnect_loop_drops_unknown_symbol_tick() {
let assembled = fixture_assembled();
let cancel = CancellationToken::new();
let connects = Arc::new(StdMutex::new(0));
let max_window = Arc::new(StdMutex::new(0));
let stray = RawTick::Quote(RawQuoteTick {
symbol: "SPX-20260717-99999-C".to_owned(),
bid: Some(1.0),
ask: Some(2.0),
last: None,
bid_size: None,
ask_size: None,
});
let transport = MockTransport {
attempts: vec![vec![stray]],
attempt_idx: 0,
cursor: 0,
backfill: Some(assembled.clone()),
connects: Arc::clone(&connects),
max_window: Arc::clone(&max_window),
cancel: cancel.clone(),
};
let (sink, mut _rx_control, mut rx_coalesced) = test_sink(256);
run_reconnect_loop(
transport,
pid("ibkr"),
"SPX".to_owned(),
fixture_expiry_utc(),
assembled.fetch.aliases.clone(),
sink,
cancel,
)
.await;
let coalesced = drain(&mut rx_coalesced);
assert!(
!coalesced
.iter()
.any(|u| matches!(u, MarketUpdate::Quote(_))),
"a tick for an unknown symbol must be dropped, never keyed"
);
}
#[tokio::test]
async fn test_reconnect_loop_stops_on_cancel() {
let cancel = CancellationToken::new();
let (sink, _rx_control, _rx_coalesced) = test_sink(8);
let loop_cancel = cancel.clone();
let handle = tokio::spawn(run_reconnect_loop(
PendingTransport,
pid("ibkr"),
"SPX".to_owned(),
fixture_expiry_utc(),
AliasCatalog::new(),
sink,
loop_cancel,
));
tokio::task::yield_now().await;
cancel.cancel();
match handle.await {
Ok(()) => {}
Err(e) => panic!("the loop task should join cleanly on cancel, got: {e}"),
}
}
#[track_caller]
fn seeded_store() -> ChainStore {
let assembled = fixture_assembled();
let seeded_at = utc_rfc3339("2026-07-01T14:00:00+00:00");
ChainStore::seed(
assembled.fetch,
ChainSource::Merged,
Duration::from_secs(2),
seeded_at,
)
}
#[test]
fn test_venue_greeks_and_quote_merge_into_seeded_store() {
let mut store = seeded_store();
let (aliases, symbol) = one_leg_aliases();
let key = match aliases.resolve_symbol(&symbol) {
Some(k) => k.clone(),
None => panic!("the fixture symbol must resolve"),
};
let instrument = match aliases.instrument(&key, &pid("ibkr")) {
Some(i) => i.clone(),
None => panic!("the fixture instrument must resolve"),
};
let greeks = GreeksRow {
instrument: instrument.clone(),
iv: Some(pos(0.25)),
delta: Some(Decimal::new(55, 2)),
gamma: Some(Decimal::new(1, 2)),
theta: Some(Decimal::new(-5, 1)),
vega: Some(Decimal::new(12, 1)),
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: now_utc(),
};
assert_eq!(store.apply_greeks(&greeks), MergeOutcome::Applied);
let quote = QuoteUpdate {
instrument,
bid: Some(pos(12.5)),
ask: Some(pos(13.5)),
last: None,
bid_size: None,
ask_size: None,
event_time: None,
received_time: now_utc(),
};
assert_eq!(store.apply_quote("e), MergeOutcome::Applied);
}
proptest::proptest! {
#[test]
fn prop_resolve_expiry_total(
ltt in "\\PC{0,20}",
ymd in "\\PC{0,12}",
) {
let session = default_session();
let ltt_opt = if ltt.is_empty() { None } else { Some(ltt.as_str()) };
let _ = resolve_expiry(ltt_opt, &ymd, &session);
}
#[test]
fn prop_strike_positive_total(raw in -1.0e9f64..1.0e9) {
match strike_positive(raw) {
Ok(strike) => proptest::prop_assert!(strike > Positive::ZERO),
Err(kind) => proptest::prop_assert_eq!(kind, NormalizeKind::OutOfRange("strike")),
}
}
#[test]
fn prop_normalize_quote_total(bid in -1.0e6f64..1.0e6, ask in -1.0e6f64..1.0e6) {
match normalize_quote(Some(bid), Some(ask)) {
Ok(quote) => {
if let (Some(b), Some(a)) = (quote.bid, quote.ask) {
proptest::prop_assert!(a >= b, "an accepted quote is never crossed");
}
}
Err(kind) => proptest::prop_assert_eq!(kind, NormalizeKind::OutOfRange("ask")),
}
}
}
}