use std::collections::BTreeMap;
use std::sync::Arc;
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::Positive;
use optionstratlib::{ExpirationDate, OptionStyle};
use tokio::sync::Notify;
use tokio::sync::mpsc::UnboundedReceiver;
use tokio_util::sync::CancellationToken;
use ig_client::application::client::{Client, StreamerClient};
use ig_client::application::config::{Config, Credentials};
use ig_client::application::interfaces::market::MarketService;
use ig_client::error::AppError;
use ig_client::model::streaming::StreamingMarketField;
use ig_client::presentation::instrument::InstrumentType;
use ig_client::presentation::market::MarketData;
use ig_client::presentation::price::PriceData;
use ig_client::utils::parsing::parse_instrument_name;
use super::{
AuthKind, ChainCapability, ChainPollCapability, GreeksCapability, MarketUpdateSink,
OptionStreamCapability, Provider, ProviderCapabilities, SendState, SubscriptionHandle,
SubscriptionRequest, UnderlyingRef,
};
use crate::chain::{
AliasCatalog, ChainFetch, ChainSnapshot, ChainSource, ContractSpecFingerprint, ExerciseStyle,
ExpirySource, Instrument, InstrumentKey, MarketUpdate, ProviderId, QuoteUpdate,
SettlementStyle, StreamHealth,
};
use crate::config::{EnvSource, Secret, require_credentials};
use crate::error::{NormalizeKind, ProviderError, TransportDetail, TransportKind};
const IG_ID: &str = "ig";
const CREDENTIAL_KEYS: [&str; 3] = ["username", "password", "api_key"];
const ACCOUNT_ID_VAR: &str = "CHAINVIEW_IG_ACCOUNT_ID";
const ENVIRONMENT_VAR: &str = "CHAINVIEW_IG_ENVIRONMENT";
const REFRESH_HINT_SECS: u32 = 5;
const DEFAULT_QUOTE_CURRENCY: &str = "USD";
const DEFAULT_CONTRACT_MULTIPLIER: u32 = 1;
const LIVE_REST_URL: &str = "https://api.ig.com/gateway/deal";
const LIVE_WS_URL: &str = "wss://apd.marketdatasystems.com";
const MAX_MARKETS: usize = 8_192;
const MARKET_CAP: &str = "ig market cap";
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(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub(crate) enum IgEnvironment {
#[default]
Demo,
Live,
}
impl IgEnvironment {
fn from_value(value: &str) -> Self {
if value.trim().eq_ignore_ascii_case("live") {
Self::Live
} else {
Self::Demo
}
}
}
#[derive(Clone)]
pub(crate) struct IgAdapter {
id: ProviderId,
username: Secret,
password: Secret,
api_key: Secret,
account_id: String,
environment: IgEnvironment,
}
impl IgAdapter {
pub(crate) fn from_env(env: &dyn EnvSource) -> Result<Self, crate::error::ConfigError> {
let id = ig_provider_id();
let creds = require_credentials(env, &id, &CREDENTIAL_KEYS)?;
let username = creds
.get("USERNAME")
.cloned()
.ok_or_else(|| crate::error::ConfigError::MissingCredential(id.clone()))?;
let password = creds
.get("PASSWORD")
.cloned()
.ok_or_else(|| crate::error::ConfigError::MissingCredential(id.clone()))?;
let api_key = creds
.get("API_KEY")
.cloned()
.ok_or_else(|| crate::error::ConfigError::MissingCredential(id.clone()))?;
let account_id = env.get(ACCOUNT_ID_VAR).unwrap_or_default();
let environment = env
.get(ENVIRONMENT_VAR)
.map(|value| IgEnvironment::from_value(&value))
.unwrap_or_default();
Ok(Self {
id,
username,
password,
api_key,
account_id,
environment,
})
}
fn credentials(&self) -> Credentials {
Credentials::new(
self.username.expose().to_owned(),
self.password.expose().to_owned(),
self.account_id.clone(),
self.api_key.expose().to_owned(),
)
}
fn upstream_config(&self) -> Config {
let mut config = Config::from_credentials(self.credentials());
if self.environment == IgEnvironment::Live {
config.rest_api.base_url = LIVE_REST_URL.to_owned();
config.websocket.url = LIVE_WS_URL.to_owned();
}
config
}
fn client(&self) -> Result<Client, ProviderError> {
Client::with_config(self.upstream_config()).map_err(ig_error)
}
}
#[async_trait]
impl Provider for IgAdapter {
fn id(&self) -> ProviderId {
self.id.clone()
}
fn capabilities(&self) -> ProviderCapabilities {
ig_capabilities()
}
async fn discover(&self) -> Result<Vec<UnderlyingRef>, ProviderError> {
let source = LiveMarketSource {
client: self.client()?,
};
discover_underlyings(&source).await
}
async fn fetch_chain(
&self,
underlying: &str,
expiration: &ExpirationDate,
) -> Result<ChainFetch, ProviderError> {
let source = LiveMarketSource {
client: self.client()?,
};
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 ig_provider_id() -> ProviderId {
match ProviderId::new(IG_ID) {
Ok(id) => id,
Err(_) => unreachable!("`ig` is a valid, reserved provider id literal"),
}
}
#[must_use]
pub(crate) fn ig_capabilities() -> ProviderCapabilities {
ProviderCapabilities::builder()
.chain(ChainCapability::Partial)
.depth(false)
.greeks(GreeksCapability::ComputedLocally)
.option_stream(OptionStreamCapability::ChainQuotes { verified: false })
.underlying_stream(false)
.chain_poll(ChainPollCapability::Poll {
interval_hint_secs: REFRESH_HINT_SECS,
})
.trades_tape(false)
.auth(AuthKind::UserPass)
.build()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SessionZone {
UkLondon,
UsEastern,
}
#[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 session_for(currency: Option<&str>) -> SessionClose {
match currency.map(str::to_ascii_uppercase).as_deref() {
Some("USD") => SessionClose::new(SessionZone::UsEastern, 16, 0),
_ => SessionClose::new(SessionZone::UkLondon, 16, 30),
}
}
enum ParsedStamp {
Absolute(DateTime<Utc>),
Local(NaiveDateTime),
}
fn resolve_expiry(
timestamped: Option<&str>,
date_only: &str,
session: &SessionClose,
) -> Result<DateTime<Utc>, NormalizeKind> {
if let Some(stamp) = timestamped.map(str::trim).filter(|s| !s.is_empty()) {
return match parse_timestamp(stamp) {
Some(ParsedStamp::Absolute(dt)) => Ok(dt),
Some(ParsedStamp::Local(naive)) => {
local_to_utc(naive, session.zone).ok_or(NormalizeKind::UnparseableExpiry)
}
None => Err(NormalizeKind::UnparseableExpiry),
};
}
let date = parse_ig_date(date_only).ok_or(NormalizeKind::UnparseableExpiry)?;
let naive = 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_timestamp(s: &str) -> Option<ParsedStamp> {
if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
return Some(ParsedStamp::Absolute(dt.with_timezone(&Utc)));
}
parse_naive_datetime(s).map(ParsedStamp::Local)
}
fn parse_naive_datetime(s: &str) -> Option<NaiveDateTime> {
let (date_part, time_part) = s.split_once('T').or_else(|| s.split_once(' '))?;
let date = parse_iso_date(date_part)?;
let mut time = time_part.split(':');
let hour = time.next()?.parse::<u32>().ok()?;
let minute = time.next()?.parse::<u32>().ok()?;
let second = match time.next() {
Some(sec) => sec.parse::<u32>().ok()?,
None => 0,
};
if time.next().is_some() {
return None;
}
date.and_hms_opt(hour, minute, second)
}
fn parse_iso_date(s: &str) -> Option<NaiveDate> {
let mut parts = s.split('-');
let year = parts.next()?.parse::<i32>().ok()?;
let month = parts.next()?.parse::<u32>().ok()?;
let day = parts.next()?.parse::<u32>().ok()?;
if parts.next().is_some() {
return None;
}
NaiveDate::from_ymd_opt(year, month, day)
}
fn parse_ig_date(s: &str) -> Option<NaiveDate> {
let trimmed = s.trim();
if let Some(date) = parse_iso_date(trimmed) {
return Some(date);
}
let mut parts = trimmed.split('-');
let day = parts.next()?.trim().parse::<u32>().ok()?;
let month = month_from_abbrev(parts.next()?.trim())?;
let year_two = parts.next()?.trim().parse::<i32>().ok()?;
if parts.next().is_some() {
return None;
}
let year = 2000 + year_two;
NaiveDate::from_ymd_opt(year, month, day)
}
fn month_from_abbrev(abbrev: &str) -> Option<u32> {
match abbrev.to_ascii_uppercase().as_str() {
"JAN" => Some(1),
"FEB" => Some(2),
"MAR" => Some(3),
"APR" => Some(4),
"MAY" => Some(5),
"JUN" => Some(6),
"JUL" => Some(7),
"AUG" => Some(8),
"SEP" => Some(9),
"OCT" => Some(10),
"NOV" => Some(11),
"DEC" => Some(12),
_ => None,
}
}
fn local_to_utc(naive: NaiveDateTime, zone: SessionZone) -> Option<DateTime<Utc>> {
let date = naive.date();
let offset_hours: i64 = match zone {
SessionZone::UkLondon => {
if uk_dst(date) {
1
} else {
0
}
}
SessionZone::UsEastern => {
if us_eastern_dst(date) {
-4
} else {
-5
}
}
};
let utc_naive = naive.checked_sub_signed(TimeDelta::hours(offset_hours))?;
Some(DateTime::<Utc>::from_naive_utc_and_offset(utc_naive, Utc))
}
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> {
Positive::new(value).ok()
}
#[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 strike_positive(value: &str) -> Result<Positive, NormalizeKind> {
let parsed = value
.trim()
.parse::<f64>()
.ok()
.and_then(positive_or_drop)
.ok_or(NormalizeKind::OutOfRange("strike"))?;
if parsed == Positive::ZERO {
return Err(NormalizeKind::OutOfRange("strike"));
}
Ok(parsed)
}
fn ig_fingerprint(epic_root: &str, quote_currency: &str) -> ContractSpecFingerprint {
ContractSpecFingerprint {
contract_multiplier: DEFAULT_CONTRACT_MULTIPLIER,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: quote_currency.to_owned(),
venue_product_code: epic_root.to_owned(),
}
}
fn epic_root(epic: &str) -> String {
epic.split('.').take(3).collect::<Vec<_>>().join(".")
}
#[derive(Debug, Clone)]
pub(crate) struct RawMarket {
epic: String,
instrument_name: String,
expiry: String,
last_dealing_date: Option<String>,
is_option: bool,
bid: Option<f64>,
offer: Option<f64>,
currency: Option<String>,
}
#[derive(Debug, Clone)]
struct RawPriceEvent {
epic: String,
bid: Option<f64>,
offer: Option<f64>,
}
fn is_option_type(instrument_type: &InstrumentType) -> bool {
matches!(
instrument_type,
InstrumentType::OptCommodities
| InstrumentType::OptCurrencies
| InstrumentType::OptIndices
| InstrumentType::OptRates
| InstrumentType::OptShares
)
}
pub(crate) fn map_market_data(market: &MarketData) -> RawMarket {
RawMarket {
epic: market.epic.clone(),
instrument_name: market.instrument_name.clone(),
expiry: market.expiry.clone(),
last_dealing_date: None,
is_option: is_option_type(&market.instrument_type),
bid: market.bid,
offer: market.offer,
currency: None,
}
}
fn map_price_data(price: &PriceData) -> RawPriceEvent {
let epic = price
.item_name
.strip_prefix("MARKET:")
.unwrap_or(&price.item_name)
.to_owned();
RawPriceEvent {
epic,
bid: price.fields.bid,
offer: price.fields.offer,
}
}
#[async_trait]
trait IgMarketSource: Send + Sync {
async fn search(&self, term: &str) -> Result<Vec<RawMarket>, ProviderError>;
async fn navigation_roots(&self) -> Result<Vec<String>, ProviderError>;
}
struct LiveMarketSource {
client: Client,
}
#[async_trait]
impl IgMarketSource for LiveMarketSource {
async fn search(&self, term: &str) -> Result<Vec<RawMarket>, ProviderError> {
let response = self.client.search_markets(term).await.map_err(ig_error)?;
Ok(response.markets.iter().map(map_market_data).collect())
}
async fn navigation_roots(&self) -> Result<Vec<String>, ProviderError> {
let response = self
.client
.get_market_navigation()
.await
.map_err(ig_error)?;
Ok(response.nodes.into_iter().map(|node| node.name).collect())
}
}
async fn discover_underlyings<S: IgMarketSource + ?Sized>(
source: &S,
) -> Result<Vec<UnderlyingRef>, ProviderError> {
let roots = source.navigation_roots().await?;
Ok(roots
.into_iter()
.take(MAX_MARKETS)
.map(UnderlyingRef::new)
.collect())
}
#[derive(Debug, Clone)]
struct AssembledChain {
fetch: ChainFetch,
}
#[derive(Debug, Clone)]
struct NormalizedLeg {
key: InstrumentKey,
native_symbol: String,
spec: ContractSpecFingerprint,
style: OptionStyle,
bid: Option<Positive>,
ask: Option<Positive>,
}
fn normalize_leg(
market: &RawMarket,
underlying: &str,
target_expiry_utc: DateTime<Utc>,
) -> Option<NormalizedLeg> {
if !market.is_option {
return None;
}
let parsed = parse_instrument_name(&market.instrument_name);
let style = match parsed.option_type.as_deref() {
Some("CALL") => OptionStyle::Call,
Some("PUT") => OptionStyle::Put,
_ => return None,
};
let strike = strike_positive(parsed.strike.as_deref()?).ok()?;
let session = session_for(market.currency.as_deref());
let expiration_utc = resolve_expiry(
market.last_dealing_date.as_deref(),
&market.expiry,
&session,
)
.ok()?;
if expiration_utc != target_expiry_utc {
return None;
}
let quote = normalize_quote(market.bid, market.offer).unwrap_or_default();
let spec = ig_fingerprint(&epic_root(&market.epic), DEFAULT_QUOTE_CURRENCY);
let key = InstrumentKey {
underlying: underlying.to_owned(),
expiration_utc,
strike,
style,
};
Some(NormalizedLeg {
key,
native_symbol: market.epic.clone(),
spec,
style,
bid: quote.bid,
ask: quote.ask,
})
}
#[derive(Debug, Default)]
struct StrikePair<'a> {
call: Option<&'a NormalizedLeg>,
put: Option<&'a NormalizedLeg>,
}
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(ProviderError::Normalize {
kind: NormalizeKind::UnparseableExpiry,
});
}
let seconds = seconds as i64;
let instant = received
.checked_add_signed(TimeDelta::seconds(seconds))
.ok_or(ProviderError::Normalize {
kind: NormalizeKind::UnparseableExpiry,
})?;
let session = SessionClose::new(SessionZone::UkLondon, 16, 30);
let naive = instant
.date_naive()
.and_hms_opt(session.hour, session.minute, 0)
.ok_or(ProviderError::Normalize {
kind: NormalizeKind::UnparseableExpiry,
})?;
local_to_utc(naive, session.zone).ok_or(ProviderError::Normalize {
kind: NormalizeKind::UnparseableExpiry,
})
}
}
}
async fn compose_chain<S: IgMarketSource + ?Sized>(
source: &S,
underlying: &str,
expiration: &ExpirationDate,
provider: &ProviderId,
received: DateTime<Utc>,
) -> Result<AssembledChain, ProviderError> {
let symbol = underlying.to_ascii_uppercase();
let expiration_utc = target_expiry(expiration, received)?;
let markets = source.search(underlying).await?;
if markets.len() > MAX_MARKETS {
return Err(ProviderError::Normalize {
kind: NormalizeKind::LimitExceeded(MARKET_CAP),
});
}
assemble_chain(&markets, &symbol, expiration_utc, provider)
}
fn assemble_chain(
markets: &[RawMarket],
underlying: &str,
expiration_utc: DateTime<Utc>,
provider: &ProviderId,
) -> Result<AssembledChain, ProviderError> {
let legs: Vec<NormalizedLeg> = markets
.iter()
.filter_map(|market| normalize_leg(market, underlying, expiration_utc))
.collect();
if legs.is_empty() {
return Err(ProviderError::NoChain {
underlying: underlying.to_owned(),
expiration: expiration_utc.to_rfc3339(),
});
}
let mut aliases = AliasCatalog::new();
for leg in &legs {
aliases.insert(Instrument {
key: leg.key.clone(),
provider: provider.clone(),
native_symbol: leg.native_symbol.clone(),
stream_symbol: Some(format!("MARKET:{}", leg.native_symbol)),
spec: leg.spec.clone(),
});
}
let mut by_strike: BTreeMap<Positive, StrikePair<'_>> = BTreeMap::new();
for leg in &legs {
let entry = by_strike.entry(leg.key.strike).or_default();
match leg.style {
OptionStyle::Call => entry.call = Some(leg),
OptionStyle::Put => entry.put = Some(leg),
}
}
let spot = median_strike(&by_strike);
let mut chain = OptionChain::new(underlying, spot, expiration_utc.to_rfc3339(), None, None);
for (strike, pair) in &by_strike {
chain.add_option(
*strike,
pair.call.and_then(|leg| leg.bid),
pair.call.and_then(|leg| leg.ask),
pair.put.and_then(|leg| leg.bid),
pair.put.and_then(|leg| leg.ask),
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 median_strike(by_strike: &BTreeMap<Positive, StrikePair<'_>>) -> Positive {
let strikes: Vec<Positive> = by_strike.keys().copied().collect();
let mid = strikes.len() / 2;
strikes.get(mid).copied().unwrap_or(Positive::ONE)
}
fn ig_error(err: AppError) -> ProviderError {
match err {
AppError::Auth(_) | AppError::Unauthorized | AppError::OAuthTokenExpired => {
ProviderError::Auth
}
AppError::Unexpected(status) => ProviderError::Transport(Box::new(TransportDetail::new(
TransportKind::Http,
Some(status.as_u16()),
))),
AppError::NotFound => ProviderError::Transport(Box::new(TransportDetail::new(
TransportKind::Http,
Some(404),
))),
AppError::RateLimitExceeded
| AppError::ApiKeyAllowanceExceeded
| AppError::AccountAllowanceExceeded
| AppError::TradingAllowanceExceeded => ProviderError::RateLimited(None),
AppError::HistoricalDataAllowanceExceeded { allowance_expiry } => {
ProviderError::RateLimited(Some(Duration::from_secs(allowance_expiry)))
}
AppError::Network(_) | AppError::Io(_) | AppError::WebSocketError(_) => {
transport(TransportKind::Closed)
}
AppError::Json(_) | AppError::Deserialization(_) | AppError::SerializationError(_) => {
transport(TransportKind::Decode)
}
AppError::Db(_) | AppError::InvalidInput(_) | AppError::Generic(_) => {
transport(TransportKind::Http)
}
AppError::CatalogRequest { .. } | AppError::CatalogPagination { .. } => {
transport(TransportKind::Http)
}
}
}
fn transport(kind: TransportKind) -> ProviderError {
ProviderError::Transport(Box::new(TransportDetail::new(kind, None)))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TransportGone;
#[async_trait]
trait IgStreamTransport: Send {
async fn connect_and_subscribe(&mut self, epics: &[String]) -> Result<(), TransportGone>;
async fn receive(&mut self) -> Result<RawPriceEvent, TransportGone>;
async fn poll(
&mut self,
underlying: &str,
expiration: &ExpirationDate,
received: DateTime<Utc>,
) -> Option<AssembledChain>;
}
struct LiveStreamTransport {
adapter: IgAdapter,
receiver: Option<UnboundedReceiver<PriceData>>,
shutdown: Option<Arc<Notify>>,
conn_task: Option<tokio::task::JoinHandle<()>>,
}
impl LiveStreamTransport {
fn new(adapter: IgAdapter) -> Self {
Self {
adapter,
receiver: None,
shutdown: None,
conn_task: None,
}
}
fn teardown(&mut self) {
if let Some(signal) = self.shutdown.take() {
signal.notify_one();
}
if let Some(task) = self.conn_task.take() {
task.abort();
}
self.receiver = None;
}
}
impl Drop for LiveStreamTransport {
fn drop(&mut self) {
self.teardown();
}
}
#[async_trait]
impl IgStreamTransport for LiveStreamTransport {
async fn connect_and_subscribe(&mut self, epics: &[String]) -> Result<(), TransportGone> {
self.teardown();
if epics.is_empty() {
return Err(TransportGone);
}
let client = self.adapter.client().map_err(|_| TransportGone)?;
let mut streamer = StreamerClient::with_client(&client)
.await
.map_err(|_| TransportGone)?;
let fields = std::collections::HashSet::from([
StreamingMarketField::Bid,
StreamingMarketField::Offer,
StreamingMarketField::UpdateTime,
]);
let receiver = streamer
.market_subscribe(epics.to_vec(), fields)
.await
.map_err(|_| TransportGone)?;
let signal = Arc::new(Notify::new());
let task_signal = Arc::clone(&signal);
let task = tokio::spawn(async move {
let mut streamer = streamer;
let _ = streamer.connect(Some(task_signal)).await;
let _ = streamer.disconnect().await;
});
self.receiver = Some(receiver);
self.shutdown = Some(signal);
self.conn_task = Some(task);
Ok(())
}
async fn receive(&mut self) -> Result<RawPriceEvent, TransportGone> {
match self.receiver.as_mut() {
Some(receiver) => match receiver.recv().await {
Some(price) => Ok(map_price_data(&price)),
None => Err(TransportGone),
},
None => Err(TransportGone),
}
}
async fn poll(
&mut self,
underlying: &str,
expiration: &ExpirationDate,
received: DateTime<Utc>,
) -> Option<AssembledChain> {
let source = LiveMarketSource {
client: self.adapter.client().ok()?,
};
compose_chain(&source, underlying, expiration, &self.adapter.id, received)
.await
.ok()
}
}
enum StreamExit {
Reconnect,
Shutdown,
}
async fn run_reconnect_loop<T: IgStreamTransport>(
mut transport: T,
id: ProviderId,
underlying: String,
expiration_utc: DateTime<Utc>,
mut aliases: AliasCatalog,
mut sink: MarketUpdateSink,
cancel: CancellationToken,
) {
let mut epics = epics_of(&aliases, &id);
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 epics, &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: IgStreamTransport>(
transport: &mut T,
id: &ProviderId,
underlying: &str,
expiration_utc: DateTime<Utc>,
epics: &mut Vec<String>,
aliases: &mut AliasCatalog,
sink: &mut MarketUpdateSink,
cancel: &CancellationToken,
attempt: &mut u32,
) -> StreamExit {
let subscribed = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
result = transport.connect_and_subscribe(epics) => result,
};
if subscribed.is_err() {
return StreamExit::Reconnect;
}
*attempt = 0;
if go_live_and_backfill(
transport,
id,
underlying,
expiration_utc,
epics,
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 Some(update) = quote_update(&event, aliases, id, now_utc()) else {
continue;
};
let sent = tokio::select! {
biased;
() = cancel.cancelled() => return StreamExit::Shutdown,
outcome = sink.send(MarketUpdate::Quote(update)) => outcome,
};
if sent == SendState::Closed {
return StreamExit::Shutdown;
}
}
}
#[allow(clippy::too_many_arguments)]
async fn go_live_and_backfill<T: IgStreamTransport>(
transport: &mut T,
id: &ProviderId,
underlying: &str,
expiration_utc: DateTime<Utc>,
epics: &mut Vec<String>,
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();
*epics = epics_of(aliases, id);
let snapshot = MarketUpdate::Chain(chain_snapshot(&composed.fetch, now_utc()));
tokio::select! {
biased;
() = cancel.cancelled() => SendState::Open,
outcome = sink.send(snapshot) => outcome,
}
}
fn epics_of(aliases: &AliasCatalog, provider: &ProviderId) -> Vec<String> {
aliases
.instruments()
.filter(|instrument| &instrument.provider == provider)
.map(|instrument| instrument.native_symbol.clone())
.collect()
}
fn quote_update(
event: &RawPriceEvent,
aliases: &AliasCatalog,
provider: &ProviderId,
received: DateTime<Utc>,
) -> Option<QuoteUpdate> {
let key = aliases.resolve_symbol(&event.epic)?.clone();
let instrument = aliases.instrument(&key, provider)?.clone();
let quote = normalize_quote(event.bid, event.offer).ok()?;
if quote.bid.is_none() && quote.ask.is_none() {
return None;
}
Some(QuoteUpdate {
instrument,
bid: quote.bid,
ask: quote.ask,
last: None,
bid_size: None,
ask_size: None,
event_time: None,
received_time: received,
})
}
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;
use std::sync::Arc;
use std::sync::Mutex as StdMutex;
use tokio::sync::mpsc;
use super::*;
use crate::chain::{
ChainStore, GreeksOrigin, LegStatus, MergeOutcome, PricingInputs, QuoteClocks,
compute_leg_greeks,
};
#[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(HashMap<String, String>);
impl EnvSource for MapEnv {
fn get(&self, key: &str) -> Option<String> {
self.0.get(key).cloned()
}
}
fn creds_env() -> MapEnv {
let mut env = HashMap::new();
let _ = env.insert("CHAINVIEW_IG_USERNAME".to_owned(), "alice".to_owned());
let _ = env.insert(
"CHAINVIEW_IG_PASSWORD".to_owned(),
"do-not-log-this-password".to_owned(),
);
let _ = env.insert(
"CHAINVIEW_IG_API_KEY".to_owned(),
"do-not-log-this-key".to_owned(),
);
MapEnv(env)
}
#[track_caller]
fn sample_adapter() -> IgAdapter {
match IgAdapter::from_env(&creds_env()) {
Ok(adapter) => adapter,
Err(e) => panic!("from_env should succeed with all creds 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 MARKET_NAVIGATION: &str =
include_str!("../../tests/fixtures/ig/market_navigation_spx.json");
#[track_caller]
fn navigation_markets() -> Vec<RawMarket> {
use ig_client::model::responses::MarketNavigationResponse;
let response: MarketNavigationResponse = match serde_json::from_str(MARKET_NAVIGATION) {
Ok(r) => r,
Err(e) => panic!("the market-navigation fixture must deserialize: {e}"),
};
response.markets.iter().map(map_market_data).collect()
}
#[track_caller]
fn fixture_expiry_utc() -> DateTime<Utc> {
let session = session_for(None);
match resolve_expiry(None, "18-JUL-26", &session) {
Ok(dt) => dt,
Err(e) => panic!("fixture expiry must resolve: {e}"),
}
}
#[test]
fn test_ig_id_is_valid_and_reserved() {
let id = ig_provider_id();
assert_eq!(id.as_str(), "ig");
assert!(id.is_reserved());
assert!(ProviderId::new(IG_ID).is_ok());
}
#[test]
fn test_ig_capabilities_match_section_8_row() {
let caps = ig_capabilities();
assert_eq!(caps.chain, ChainCapability::Partial);
assert!(!caps.depth, "IG populates no option depth ladder (#50)");
assert_eq!(caps.greeks, GreeksCapability::ComputedLocally);
assert_eq!(
caps.option_stream,
OptionStreamCapability::ChainQuotes { verified: false }
);
assert!(
!caps.underlying_stream,
"IG folds no underlying quote (only option epics stream; #40/#41 honesty)"
);
assert_eq!(
caps.chain_poll,
ChainPollCapability::Poll {
interval_hint_secs: REFRESH_HINT_SECS
}
);
assert!(!caps.trades_tape, "IG's trade stream is deal confirmations");
assert_eq!(caps.auth, AuthKind::UserPass);
}
#[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(), "ig");
assert_eq!(adapter.capabilities().chain, ChainCapability::Partial);
assert_eq!(
adapter.capabilities().greeks,
GreeksCapability::ComputedLocally
);
}
#[test]
fn test_from_env_reads_chainview_namespace_only() {
let mut env = HashMap::new();
let _ = env.insert("CHAINVIEW_IG_USERNAME".to_owned(), "alice".to_owned());
let _ = env.insert("CHAINVIEW_IG_PASSWORD".to_owned(), "pw".to_owned());
let _ = env.insert("CHAINVIEW_IG_API_KEY".to_owned(), "key".to_owned());
let _ = env.insert("IG_USERNAME".to_owned(), "foreign".to_owned());
let adapter = match IgAdapter::from_env(&MapEnv(env)) {
Ok(adapter) => adapter,
Err(e) => panic!("from_env should succeed: {e}"),
};
assert_eq!(adapter.username.expose(), "alice");
assert_eq!(adapter.api_key.expose(), "key");
assert_eq!(adapter.environment, IgEnvironment::Demo);
assert_eq!(adapter.account_id, "");
}
#[test]
fn test_from_env_environment_and_account_id() {
let mut env = HashMap::new();
let _ = env.insert("CHAINVIEW_IG_USERNAME".to_owned(), "u".to_owned());
let _ = env.insert("CHAINVIEW_IG_PASSWORD".to_owned(), "p".to_owned());
let _ = env.insert("CHAINVIEW_IG_API_KEY".to_owned(), "k".to_owned());
let _ = env.insert("CHAINVIEW_IG_ACCOUNT_ID".to_owned(), "ABC123".to_owned());
let _ = env.insert("CHAINVIEW_IG_ENVIRONMENT".to_owned(), "LIVE".to_owned());
match IgAdapter::from_env(&MapEnv(env)) {
Ok(adapter) => {
assert_eq!(adapter.environment, IgEnvironment::Live);
assert_eq!(adapter.account_id, "ABC123");
}
Err(e) => panic!("from_env should succeed: {e}"),
}
}
#[test]
fn test_from_env_missing_credential_is_error() {
let mut env = HashMap::new();
let _ = env.insert("CHAINVIEW_IG_USERNAME".to_owned(), "alice".to_owned());
match IgAdapter::from_env(&MapEnv(env)) {
Err(crate::error::ConfigError::MissingCredential(id)) => {
assert_eq!(id.as_str(), "ig");
}
Err(other) => panic!("expected MissingCredential, got: {other}"),
Ok(_) => panic!("expected MissingCredential, got Ok"),
}
}
#[test]
fn test_secret_debug_never_reveals_credentials() {
let adapter = sample_adapter();
let rendered = format!("{:?} {:?}", adapter.password, adapter.api_key);
assert!(!rendered.contains("do-not-log-this-password"));
assert!(!rendered.contains("do-not-log-this-key"));
assert!(rendered.contains("redacted"));
}
#[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-ig-canary-4b2e-present";
#[test]
fn test_construction_and_errors_never_log_credentials() {
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 adapter = sample_adapter();
tracing::debug!(config = ?adapter.upstream_config(), "ig config built");
tracing::debug!(creds = ?adapter.credentials(), "ig credentials built");
let err = ig_error(AppError::Auth(ig_client::error::AuthError::Other(
"boom".to_owned(),
)));
tracing::debug!(error = %err, "ig error rendered");
let output = logs.contents();
assert!(output.contains(CONTROL_CANARY), "sink must capture content");
assert!(
!output.contains("do-not-log-this-password"),
"the password leaked into logs:\n{output}"
);
assert!(
!output.contains("do-not-log-this-key"),
"the api key leaked into logs:\n{output}"
);
assert!(
!format!("{err} {err:?}").contains("do-not-log-this-key"),
"a ProviderError must never carry a credential"
);
}
#[test]
fn test_ig_allowance_variants_map_to_rate_limited_without_a_retry_after() {
for err in [
AppError::RateLimitExceeded,
AppError::ApiKeyAllowanceExceeded,
AppError::AccountAllowanceExceeded,
AppError::TradingAllowanceExceeded,
] {
match ig_error(err) {
ProviderError::RateLimited(None) => {}
other => panic!("expected RateLimited(None), got {other:?}"),
}
}
match ig_error(AppError::HistoricalDataAllowanceExceeded {
allowance_expiry: 604_800,
}) {
ProviderError::RateLimited(Some(delay)) => {
assert_eq!(delay, Duration::from_secs(604_800));
}
other => panic!("expected RateLimited(Some(_)), got {other:?}"),
}
}
#[test]
fn test_expiry_date_only_resolves_via_uk_session_close() {
let session = session_for(None);
match resolve_expiry(None, "18-JUL-26", &session) {
Ok(utc) => assert_eq!(utc.to_rfc3339(), "2026-07-18T15:30:00+00:00"),
Err(e) => panic!("date-only UK expiry should resolve, got: {e}"),
}
}
#[test]
fn test_expiry_date_only_us_session_close() {
let session = session_for(Some("USD"));
match resolve_expiry(None, "18-JUL-26", &session) {
Ok(utc) => assert_eq!(utc.to_rfc3339(), "2026-07-18T20:00:00+00:00"),
Err(e) => panic!("date-only US expiry should resolve, got: {e}"),
}
}
#[test]
fn test_expiry_timestamped_field_wins_over_date_only() {
let session = session_for(None);
let resolved = match resolve_expiry(Some("2026-07-18T20:30"), "18-JUL-26", &session) {
Ok(dt) => dt,
Err(e) => panic!("timestamped expiry should resolve, got: {e}"),
};
assert_eq!(resolved.to_rfc3339(), "2026-07-18T19:30:00+00:00");
let date_only = match resolve_expiry(None, "18-JUL-26", &session) {
Ok(dt) => dt,
Err(e) => panic!("date-only should resolve, got: {e}"),
};
assert_ne!(resolved, date_only, "the timestamped field must win");
}
#[test]
fn test_expiry_timestamped_rfc3339_offset_is_absolute() {
let session = session_for(None);
match resolve_expiry(Some("2026-07-18T16:30:00-04:00"), "18-JUL-26", &session) {
Ok(utc) => assert_eq!(utc.to_rfc3339(), "2026-07-18T20:30:00+00:00"),
Err(e) => panic!("offset-carrying timestamp should resolve, got: {e}"),
}
}
#[test]
fn test_expiry_ambiguous_and_unparseable_are_normalize_errors() {
let session = session_for(None);
assert_eq!(
resolve_expiry(None, "not-a-date", &session),
Err(NormalizeKind::UnparseableExpiry)
);
assert_eq!(
resolve_expiry(Some("garbage-stamp"), "18-JUL-26", &session),
Err(NormalizeKind::UnparseableExpiry)
);
assert_eq!(
resolve_expiry(None, "18-XXX-26", &session),
Err(NormalizeKind::UnparseableExpiry)
);
}
#[test]
fn test_dst_boundaries() {
assert!(uk_dst(date(2026, 3, 29)));
assert!(!uk_dst(date(2026, 3, 28)));
assert!(!uk_dst(date(2026, 10, 25)));
assert!(uk_dst(date(2026, 10, 24)));
assert!(us_eastern_dst(date(2026, 3, 8)));
assert!(!us_eastern_dst(date(2026, 3, 7)));
}
#[test]
fn test_parse_ig_date_forms() {
assert_eq!(parse_ig_date("18-JUL-26"), Some(date(2026, 7, 18)));
assert_eq!(parse_ig_date("2026-07-18"), Some(date(2026, 7, 18)));
assert_eq!(parse_ig_date("18-XXX-26"), None);
assert_eq!(parse_ig_date("garbage"), None);
}
#[test]
fn test_strike_positive_parses_and_rejects() {
match strike_positive("7500") {
Ok(strike) => assert_eq!(strike, pos(7500.0)),
Err(e) => panic!("7500 should parse, got: {e}"),
}
match strike_positive("10.5") {
Ok(strike) => assert_eq!(strike, pos(10.5)),
Err(e) => panic!("10.5 should parse, got: {e}"),
}
assert_eq!(
strike_positive("0"),
Err(NormalizeKind::OutOfRange("strike"))
);
assert_eq!(
strike_positive("-5"),
Err(NormalizeKind::OutOfRange("strike"))
);
assert_eq!(
strike_positive("abc"),
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"))
);
assert_eq!(
normalize_quote(Some(5.0), Some(0.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_assemble_chain_from_recorded_navigation() {
let markets = navigation_markets();
let assembled = match assemble_chain(&markets, "SPX", fixture_expiry_utc(), &pid("ig")) {
Ok(a) => a,
Err(e) => panic!("assembly should succeed, got: {e}"),
};
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("ig"))
.count();
assert_eq!(leg_count, 6);
assert_eq!(
assembled.fetch.expiry_source.expiration_utc,
fixture_expiry_utc()
);
assert!(
!assembled
.fetch
.aliases
.instruments()
.any(|i| i.native_symbol == "IX.D.SPX.DAILY.IP"),
);
}
#[test]
fn test_assemble_chain_wrong_expiry_is_no_chain() {
let markets = navigation_markets();
let other = utc_rfc3339("2027-01-15T20:00:00+00:00");
match assemble_chain(&markets, "SPX", other, &pid("ig")) {
Err(ProviderError::NoChain { underlying, .. }) => assert_eq!(underlying, "SPX"),
other => panic!("expected NoChain for a mismatched expiry, got {other:?}"),
}
}
#[test]
fn test_epic_root_extracts_product_code() {
assert_eq!(epic_root("OP.D.SPX2.7500C.IP"), "OP.D.SPX2");
assert_eq!(epic_root("SIMPLE"), "SIMPLE");
}
#[track_caller]
fn seeded_store() -> ChainStore {
let markets = navigation_markets();
let assembled = match assemble_chain(&markets, "SPX", fixture_expiry_utc(), &pid("ig")) {
Ok(a) => a,
Err(e) => panic!("assembly should succeed, got: {e}"),
};
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_local_greeks_are_computed_and_tagged_computed_locally() {
let store = seeded_store();
let key = InstrumentKey {
underlying: "SPX".to_owned(),
expiration_utc: fixture_expiry_utc(),
strike: pos(7500.0),
style: OptionStyle::Call,
};
match store.leg_greeks(&key) {
Some(leg) => {
assert_eq!(leg.status, LegStatus::Computed);
assert!(
leg.iv.is_some(),
"IV inverted locally from the two-sided quote"
);
assert_eq!(leg.iv_origin, GreeksOrigin::ComputedLocally);
assert!(leg.delta.is_some());
assert_eq!(leg.delta_origin, GreeksOrigin::ComputedLocally);
assert!(leg.theta.is_some());
assert_eq!(leg.theta_origin, GreeksOrigin::ComputedLocally);
assert!(leg.vega.is_some());
assert!(leg.gamma.is_some());
assert_eq!(leg.gamma_origin, GreeksOrigin::ComputedLocally);
}
None => panic!("expected a computed-locally sidecar entry for the 7500 call"),
}
}
#[test]
fn test_crossed_input_leaves_leg_none_never_a_stale_greek() {
let mut chain = OptionChain::new(
"SPX",
pos(7500.0),
fixture_expiry_utc().to_rfc3339(),
None,
None,
);
chain.add_option(
pos(7500.0),
Some(pos(5.0)),
Some(pos(3.0)),
Some(pos(4.0)),
Some(pos(4.5)),
Positive::ZERO,
None,
None,
None,
None,
None,
None,
);
let ctx = PricingInputs::new(pos(7500.0), utc_rfc3339("2026-07-01T14:00:00+00:00"), 1);
let mut sidecar = crate::chain::GreeksSidecar::new();
match compute_leg_greeks(&chain, &ctx, &QuoteClocks::new(), &mut sidecar) {
Ok(()) => {}
Err(e) => panic!("compute_leg_greeks failed: {e}"),
}
let call_key = InstrumentKey {
underlying: "SPX".to_owned(),
expiration_utc: fixture_expiry_utc(),
strike: pos(7500.0),
style: OptionStyle::Call,
};
match sidecar.get(&call_key) {
Some(leg) => {
assert_eq!(leg.status, LegStatus::Crossed);
assert!(
leg.iv.is_none(),
"a crossed input never yields a computed IV"
);
assert!(leg.delta.is_none());
}
None => panic!("expected a cleared sidecar entry for the crossed call"),
}
let put_key = InstrumentKey {
style: OptionStyle::Put,
..call_key
};
match sidecar.get(&put_key) {
Some(leg) => assert_eq!(leg.status, LegStatus::Computed),
None => panic!("expected a computed entry for the uncrossed put"),
}
}
struct MockSource {
markets: Vec<RawMarket>,
roots: Vec<String>,
}
#[async_trait]
impl IgMarketSource for MockSource {
async fn search(&self, _term: &str) -> Result<Vec<RawMarket>, ProviderError> {
Ok(self.markets.clone())
}
async fn navigation_roots(&self) -> Result<Vec<String>, ProviderError> {
Ok(self.roots.clone())
}
}
#[tokio::test]
async fn test_discover_returns_navigation_roots_as_underlyings() {
let source = MockSource {
markets: Vec::new(),
roots: vec!["Indices".to_owned(), "Shares".to_owned()],
};
match discover_underlyings(&source).await {
Ok(refs) => {
let names: Vec<String> = refs.into_iter().map(|r| r.underlying).collect();
assert_eq!(names, vec!["Indices".to_owned(), "Shares".to_owned()]);
}
Err(e) => panic!("discover should succeed, got: {e}"),
}
}
#[tokio::test]
async fn test_compose_chain_over_mock_source() {
let source = MockSource {
markets: navigation_markets(),
roots: Vec::new(),
};
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("ig"), received).await {
Ok(assembled) => {
assert_eq!(assembled.fetch.chain.symbol, "SPX");
assert_eq!(assembled.fetch.chain.options.len(), 3);
}
Err(e) => panic!("compose_chain should succeed, got: {e}"),
}
}
struct MockTransport {
attempts: Vec<Vec<RawPriceEvent>>,
attempt_idx: usize,
cursor: usize,
backfill: Option<AssembledChain>,
connects: Arc<StdMutex<u32>>,
cancel: CancellationToken,
}
#[async_trait]
impl IgStreamTransport for MockTransport {
async fn connect_and_subscribe(&mut self, _epics: &[String]) -> Result<(), TransportGone> {
if let Ok(mut count) = self.connects.lock() {
*count += 1;
}
self.cursor = 0;
Ok(())
}
async fn receive(&mut self) -> Result<RawPriceEvent, 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 IgStreamTransport for PendingTransport {
async fn connect_and_subscribe(&mut self, _epics: &[String]) -> Result<(), TransportGone> {
Ok(())
}
async fn receive(&mut self) -> Result<RawPriceEvent, TransportGone> {
std::future::pending::<()>().await;
Err(TransportGone)
}
async fn poll(
&mut self,
_underlying: &str,
_expiration: &ExpirationDate,
_received: DateTime<Utc>,
) -> Option<AssembledChain> {
None
}
}
#[track_caller]
fn fixture_assembled() -> AssembledChain {
let markets = navigation_markets();
match assemble_chain(&markets, "SPX", fixture_expiry_utc(), &pid("ig")) {
Ok(a) => a,
Err(e) => panic!("assembly should succeed, got: {e}"),
}
}
#[tokio::test(start_paused = true)]
async fn test_reconnect_loop_emits_reconnecting_and_backfills_without_panic() {
let assembled = fixture_assembled();
let cancel = CancellationToken::new();
let connects = Arc::new(StdMutex::new(0));
let epic = match assembled
.fetch
.aliases
.instruments()
.find(|i| i.provider == pid("ig"))
{
Some(instrument) => instrument.native_symbol.clone(),
None => panic!("the fixture must carry at least one epic"),
};
let quote = RawPriceEvent {
epic,
bid: Some(12.5),
offer: Some(13.5),
};
let transport = MockTransport {
attempts: vec![vec![quote.clone()], vec![quote]],
attempt_idx: 0,
cursor: 0,
backfill: Some(assembled.clone()),
connects: Arc::clone(&connects),
cancel: cancel.clone(),
};
let (sink, mut rx_control, mut rx_coalesced) = test_sink(256);
run_reconnect_loop(
transport,
pid("ig"),
"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"
);
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 per-epic quote reached the coalesced channel"
);
}
#[tokio::test(start_paused = true)]
async fn test_reconnect_loop_drops_unknown_epic_quote() {
let assembled = fixture_assembled();
let cancel = CancellationToken::new();
let connects = Arc::new(StdMutex::new(0));
let stray = RawPriceEvent {
epic: "OP.D.UNKNOWN.9999C.IP".to_owned(),
bid: Some(1.0),
offer: Some(2.0),
};
let transport = MockTransport {
attempts: vec![vec![stray]],
attempt_idx: 0,
cursor: 0,
backfill: Some(assembled.clone()),
connects: Arc::clone(&connects),
cancel: cancel.clone(),
};
let (sink, mut _rx_control, mut rx_coalesced) = test_sink(256);
run_reconnect_loop(
transport,
pid("ig"),
"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 quote for an unknown epic 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("ig"),
"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}"),
}
}
#[test]
fn test_stream_quote_merges_into_seeded_store() {
let mut store = seeded_store();
let key = InstrumentKey {
underlying: "SPX".to_owned(),
expiration_utc: fixture_expiry_utc(),
strike: pos(7500.0),
style: OptionStyle::Call,
};
let instrument = Instrument {
key: key.clone(),
provider: pid("ig"),
native_symbol: "OP.D.SPX2.7500C.IP".to_owned(),
stream_symbol: Some("MARKET:OP.D.SPX2.7500C.IP".to_owned()),
spec: ig_fingerprint("OP.D.SPX2", DEFAULT_QUOTE_CURRENCY),
};
let event = RawPriceEvent {
epic: "OP.D.SPX2.7500C.IP".to_owned(),
bid: Some(20.0),
offer: Some(21.0),
};
let mut aliases = AliasCatalog::new();
aliases.insert(instrument);
let update = match quote_update(&event, &aliases, &pid("ig"), now_utc()) {
Some(u) => u,
None => panic!("the quote should normalize"),
};
assert_eq!(store.apply_quote(&update), MergeOutcome::Applied);
let _ = key;
}
proptest::proptest! {
#[test]
fn prop_resolve_expiry_total(
ts in "\\PC{0,20}",
date_str in "\\PC{0,20}",
) {
let session = session_for(None);
let ts_opt = if ts.is_empty() { None } else { Some(ts.as_str()) };
let _ = resolve_expiry(ts_opt, &date_str, &session);
}
#[test]
fn prop_strike_positive_total(raw in "\\PC{0,12}") {
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")),
}
}
}
}