use std::collections::{BTreeMap, HashMap, VecDeque};
use rust_decimal::prelude::ToPrimitive;
use rust_decimal::Decimal;
use wickra_core::{self as wc, Indicator};
use wickra_exchange_core::{BookDelta, BookLevel, Event, OrderBookSnapshot, OrderSide, TradePrint};
use crate::candle::{CandleBuilder, Timeframe};
use crate::config::{IndicatorSpec, MAX_RECORDING};
use crate::error::{Error, Result};
use crate::registry::{self, TickIndicator, TickInput};
use crate::source::{DataSource, SourceId, Symbol};
const PNF_BOX_FRACTION: f64 = 0.01;
const PNF_REVERSAL_BOXES: f64 = 3.0;
#[derive(Debug, Default)]
struct PointAndFigure {
rising: bool,
extreme: f64,
previous_high: Option<f64>,
previous_low: Option<f64>,
on_buy_signal: bool,
started: bool,
}
impl PointAndFigure {
fn update(&mut self, close: f64) {
if !close.is_finite() || close <= 0.0 {
return;
}
if !self.started {
self.started = true;
self.rising = true;
self.extreme = close;
return;
}
let box_size = close * PNF_BOX_FRACTION;
let reversal = box_size * PNF_REVERSAL_BOXES;
if self.rising {
if close >= self.extreme + box_size {
self.extreme = close;
if self.previous_high.is_some_and(|high| close > high) {
self.on_buy_signal = true;
}
} else if close <= self.extreme - reversal {
self.previous_high = Some(self.extreme);
self.rising = false;
self.extreme = close;
}
} else if close <= self.extreme - box_size {
self.extreme = close;
if self.previous_low.is_some_and(|low| close < low) {
self.on_buy_signal = false;
}
} else if close >= self.extreme + reversal {
self.previous_low = Some(self.extreme);
self.rising = true;
self.extreme = close;
}
}
}
const BREADTH_MA_PERIOD: usize = 50;
#[derive(Debug, Clone, Copy, Default, PartialEq, serde::Serialize, serde::Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct DerivativesUpdate {
pub funding_rate: Option<f64>,
pub mark_price: Option<f64>,
pub index_price: Option<f64>,
pub futures_price: Option<f64>,
pub open_interest: Option<f64>,
pub long_size: Option<f64>,
pub short_size: Option<f64>,
pub long_liquidation: Option<f64>,
pub short_liquidation: Option<f64>,
pub timestamp: i64,
}
#[derive(Debug, Default)]
struct DerivativesState {
funding_rate: f64,
mark_price: f64,
index_price: f64,
futures_price: f64,
open_interest: f64,
long_size: f64,
short_size: f64,
taker_buy_volume: f64,
taker_sell_volume: f64,
long_liquidation: f64,
short_liquidation: f64,
timestamp: i64,
priced: bool,
}
impl DerivativesState {
fn apply(&mut self, update: &DerivativesUpdate) {
let set = |target: &mut f64, value: Option<f64>| {
if let Some(value) = value {
if value.is_finite() {
*target = value;
}
}
};
set(&mut self.funding_rate, update.funding_rate);
set(&mut self.mark_price, update.mark_price);
set(&mut self.index_price, update.index_price);
set(&mut self.futures_price, update.futures_price);
set(&mut self.open_interest, update.open_interest);
set(&mut self.long_size, update.long_size);
set(&mut self.short_size, update.short_size);
for (target, value) in [
(&mut self.long_liquidation, update.long_liquidation),
(&mut self.short_liquidation, update.short_liquidation),
] {
if let Some(value) = value {
if value.is_finite() && value >= 0.0 {
*target += value;
}
}
}
if update.timestamp != 0 {
self.timestamp = update.timestamp;
}
self.priced = self.mark_price > 0.0 && self.index_price > 0.0 && self.futures_price > 0.0;
}
fn add_trade(&mut self, quantity: f64, aggressor: OrderSide) {
if !quantity.is_finite() || quantity < 0.0 {
return;
}
match aggressor {
OrderSide::Buy => self.taker_buy_volume += quantity,
OrderSide::Sell => self.taker_sell_volume += quantity,
}
}
fn tick(&self) -> Option<wc::DerivativesTick> {
if !self.priced {
return None;
}
wc::DerivativesTick::new(
self.funding_rate,
self.mark_price,
self.index_price,
self.futures_price,
self.open_interest,
self.long_size,
self.short_size,
self.taker_buy_volume,
self.taker_sell_volume,
self.long_liquidation,
self.short_liquidation,
self.timestamp,
)
.ok()
}
}
#[derive(Debug)]
#[allow(
clippy::struct_excessive_bools,
reason = "three per-bar breadth signals plus a readiness flag, each independent; \
`wc::Member` carries the same set for the same reason"
)]
struct BreadthState {
previous_close: Option<f64>,
change: f64,
volume: f64,
high: f64,
low: f64,
new_high: bool,
new_low: bool,
average: wc::Sma,
above_ma: bool,
point_and_figure: PointAndFigure,
ready: bool,
}
impl BreadthState {
fn new() -> Self {
Self {
previous_close: None,
change: 0.0,
volume: 0.0,
high: f64::NEG_INFINITY,
low: f64::INFINITY,
new_high: false,
new_low: false,
average: wc::Sma::new(BREADTH_MA_PERIOD)
.expect("BREADTH_MA_PERIOD is a non-zero constant"),
above_ma: false,
point_and_figure: PointAndFigure::default(),
ready: false,
}
}
fn update(&mut self, candle: &wc::Candle) {
let close = candle.close;
if !close.is_finite() {
return;
}
self.change = self.previous_close.map_or(0.0, |previous| close - previous);
self.previous_close = Some(close);
self.volume = if candle.volume.is_finite() && candle.volume >= 0.0 {
candle.volume
} else {
0.0
};
self.new_high = candle.high > self.high;
self.new_low = candle.low < self.low;
self.high = self.high.max(candle.high);
self.low = self.low.min(candle.low);
self.above_ma = self
.average
.update(close)
.is_some_and(|average| close > average);
self.point_and_figure.update(close);
self.ready = true;
}
fn member(&self) -> Option<wc::Member> {
if !self.ready {
return None;
}
Some(wc::Member::with_signals(
self.change,
self.volume,
self.new_high,
self.new_low,
self.above_ma,
self.point_and_figure.on_buy_signal,
))
}
}
#[derive(Debug, Default, Clone)]
pub struct BookState {
bids: BTreeMap<Decimal, Decimal>,
asks: BTreeMap<Decimal, Decimal>,
}
impl BookState {
pub fn apply_snapshot(&mut self, snap: &OrderBookSnapshot) {
self.bids.clear();
self.asks.clear();
for level in &snap.bids {
self.bids.insert(level.price, level.quantity);
}
for level in &snap.asks {
self.asks.insert(level.price, level.quantity);
}
}
pub fn apply_delta(&mut self, delta: &BookDelta) {
apply_levels(&mut self.bids, &delta.bids);
apply_levels(&mut self.asks, &delta.asks);
}
#[must_use]
pub fn best_bid(&self) -> Option<(Decimal, Decimal)> {
self.bids.iter().next_back().map(|(p, q)| (*p, *q))
}
#[must_use]
pub fn best_ask(&self) -> Option<(Decimal, Decimal)> {
self.asks.iter().next().map(|(p, q)| (*p, *q))
}
#[must_use]
pub fn to_core(&self) -> Option<wc::OrderBook> {
let bids = levels(self.bids.iter().rev())?;
let asks = levels(self.asks.iter())?;
wc::OrderBook::new(bids, asks).ok()
}
#[must_use]
pub fn spread(&self) -> Option<Decimal> {
match (self.best_bid(), self.best_ask()) {
(Some((bid, _)), Some((ask, _))) => Some(ask - bid),
_ => None,
}
}
#[must_use]
pub fn mid(&self) -> Option<Decimal> {
match (self.best_bid(), self.best_ask()) {
(Some((bid, _)), Some((ask, _))) => Some((bid + ask) / Decimal::TWO),
_ => None,
}
}
#[must_use]
pub fn top_bids(&self, n: usize) -> Vec<(Decimal, Decimal)> {
self.bids
.iter()
.rev()
.take(n)
.map(|(p, q)| (*p, *q))
.collect()
}
#[must_use]
pub fn top_asks(&self, n: usize) -> Vec<(Decimal, Decimal)> {
self.asks.iter().take(n).map(|(p, q)| (*p, *q)).collect()
}
}
fn build_spec(spec: &IndicatorSpec) -> Result<Box<dyn TickIndicator>> {
match &spec.reference {
Some(reference) => registry::build_paired(&spec.kind, &spec.params, reference),
None => registry::build(&spec.kind, &spec.params),
}
}
fn core_trade(print: &TradePrint) -> Option<wc::Trade> {
let side = match print.aggressor {
OrderSide::Buy => wc::Side::Buy,
OrderSide::Sell => wc::Side::Sell,
};
wc::Trade::new(
print.price.to_f64()?,
print.quantity.to_f64()?,
side,
print.timestamp,
)
.ok()
}
fn levels<'a>(side: impl Iterator<Item = (&'a Decimal, &'a Decimal)>) -> Option<Vec<wc::Level>> {
side.map(|(price, size)| wc::Level::new(price.to_f64()?, size.to_f64()?).ok())
.collect()
}
fn apply_levels(side: &mut BTreeMap<Decimal, Decimal>, changes: &[BookLevel]) {
for level in changes {
if level.quantity.is_zero() {
side.remove(&level.price);
} else {
side.insert(level.price, level.quantity);
}
}
}
#[derive(Debug, Clone)]
pub struct TapeRing {
prints: VecDeque<TradePrint>,
cap: usize,
}
impl TapeRing {
#[must_use]
pub fn new(cap: usize) -> Self {
Self {
prints: VecDeque::with_capacity(cap),
cap,
}
}
pub fn push(&mut self, print: TradePrint) {
if self.prints.len() == self.cap {
self.prints.pop_front();
}
self.prints.push_back(print);
}
#[must_use]
pub fn recent(&self, n: usize) -> Vec<TradePrint> {
self.prints.iter().rev().take(n).cloned().collect()
}
#[must_use]
pub fn len(&self) -> usize {
self.prints.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.prints.is_empty()
}
}
const TAPE_PRINTS_KEPT: usize = 256;
impl Default for TapeRing {
fn default() -> Self {
Self::new(TAPE_PRINTS_KEPT)
}
}
const MAX_FOOTPRINT_LEVELS: usize = 1024;
#[derive(Debug, Default, Clone)]
pub struct Footprint {
levels: BTreeMap<Decimal, (Decimal, Decimal)>,
}
impl Footprint {
pub fn add(&mut self, print: &TradePrint) {
let entry = self.levels.entry(print.price).or_default();
let side = match print.aggressor {
OrderSide::Buy => &mut entry.0,
OrderSide::Sell => &mut entry.1,
};
*side = side.checked_add(print.quantity).unwrap_or(*side);
while self.levels.len() > MAX_FOOTPRINT_LEVELS {
let lowest = self.levels.keys().next().copied().unwrap_or(print.price);
let highest = self
.levels
.keys()
.next_back()
.copied()
.unwrap_or(print.price);
let furthest = if print.price.saturating_sub(lowest).abs()
>= highest.saturating_sub(print.price).abs()
{
lowest
} else {
highest
};
self.levels.remove(&furthest);
}
}
#[must_use]
pub fn at(&self, price: Decimal) -> Option<(Decimal, Decimal)> {
self.levels.get(&price).copied()
}
#[must_use]
pub fn around(&self, anchor: Decimal, depth: usize) -> Vec<(Decimal, Decimal, Decimal)> {
let mut below = self.levels.range(..anchor).rev().peekable();
let mut above = self.levels.range(anchor..).peekable();
let mut picked = Vec::with_capacity(depth.min(self.levels.len()));
while picked.len() < depth {
match (below.peek().copied(), above.peek().copied()) {
(None, None) => break,
(Some((&price, &(buy, sell))), None) => {
picked.push((price, buy, sell));
let _ = below.next();
}
(None, Some((&price, &(buy, sell)))) => {
picked.push((price, buy, sell));
let _ = above.next();
}
(Some((&low, &(low_buy, low_sell))), Some((&high, &(high_buy, high_sell)))) => {
if high.saturating_sub(anchor) <= anchor.saturating_sub(low) {
picked.push((high, high_buy, high_sell));
let _ = above.next();
} else {
picked.push((low, low_buy, low_sell));
let _ = below.next();
}
}
}
}
picked.sort_by_key(|level| std::cmp::Reverse(level.0));
picked
}
#[must_use]
pub fn len(&self) -> usize {
self.levels.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.levels.is_empty()
}
}
struct IndicatorEntry {
label: String,
indicator: Box<dyn TickIndicator>,
last: Option<f64>,
fields: Vec<(&'static str, f64)>,
series: VecDeque<f64>,
}
const INDICATOR_SERIES: usize = 120;
#[derive(Debug, Clone, PartialEq)]
pub struct IndicatorReading {
pub label: String,
pub value: Option<f64>,
pub fields: Vec<(&'static str, f64)>,
pub series: Vec<f64>,
}
impl std::fmt::Debug for IndicatorEntry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("IndicatorEntry")
.field("label", &self.label)
.field("last", &self.last)
.field("fields", &self.fields)
.field("series", &self.series.len())
.finish_non_exhaustive()
}
}
struct ProfileEntry {
label: String,
profile: Box<dyn registry::ProfileIndicator>,
last: Option<registry::ProfileReading>,
}
impl std::fmt::Debug for ProfileEntry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ProfileEntry")
.field("label", &self.label)
.field("last", &self.last)
.finish_non_exhaustive()
}
}
#[derive(Debug, Default)]
pub struct ProfileSet {
entries: Vec<ProfileEntry>,
}
impl ProfileSet {
pub fn from_specs(specs: &[IndicatorSpec]) -> Result<Self> {
let entries = specs
.iter()
.map(|spec| {
Ok(ProfileEntry {
label: spec.label(),
profile: registry::build_profile(&spec.kind, &spec.params)?,
last: None,
})
})
.collect::<Result<Vec<_>>>()?;
Ok(Self { entries })
}
pub fn update(&mut self, input: &TickInput) {
for entry in &mut self.entries {
if let Some(reading) = entry.profile.update(input) {
entry.last = Some(reading);
}
}
}
#[must_use]
pub fn readings(&self) -> Vec<(&str, Option<®istry::ProfileReading>)> {
self.entries
.iter()
.map(|entry| (entry.label.as_str(), entry.last.as_ref()))
.collect()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
}
const ALT_BARS_KEPT: usize = 256;
const OHLC_HISTORY: usize = 256;
const PRICE_HISTORY: usize = 512;
struct BarEntry {
label: String,
stream: Box<dyn registry::BarStream>,
bars: VecDeque<registry::AltBar>,
}
impl std::fmt::Debug for BarEntry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("BarEntry")
.field("label", &self.label)
.field("bars", &self.bars.len())
.finish_non_exhaustive()
}
}
#[derive(Debug, Default)]
pub struct BarSet {
entries: Vec<BarEntry>,
}
impl BarSet {
pub fn from_specs(specs: &[IndicatorSpec]) -> Result<Self> {
let entries = specs
.iter()
.map(|spec| {
Ok(BarEntry {
label: spec.label(),
stream: registry::build_bars(&spec.kind, &spec.params)?,
bars: VecDeque::with_capacity(ALT_BARS_KEPT),
})
})
.collect::<Result<Vec<_>>>()?;
Ok(Self { entries })
}
pub fn update(&mut self, input: &TickInput) {
for entry in &mut self.entries {
for bar in entry.stream.update(input) {
if entry.bars.len() == ALT_BARS_KEPT {
entry.bars.pop_front();
}
entry.bars.push_back(bar);
}
}
}
#[must_use]
pub fn streams(&self, depth: usize) -> Vec<(&str, Vec<registry::AltBar>)> {
self.entries
.iter()
.map(|entry| {
let take = entry.bars.len().saturating_sub(depth);
(
entry.label.as_str(),
entry.bars.iter().skip(take).copied().collect(),
)
})
.collect()
}
}
#[derive(Debug)]
pub struct IndicatorSet {
entries: Vec<IndicatorEntry>,
}
impl IndicatorSet {
pub fn from_specs(specs: &[IndicatorSpec]) -> Result<Self> {
let entries = specs
.iter()
.map(|spec| {
Ok(IndicatorEntry {
label: spec.label(),
indicator: build_spec(spec)?,
last: None,
fields: Vec::new(),
series: VecDeque::with_capacity(INDICATOR_SERIES),
})
})
.collect::<Result<Vec<_>>>()?;
Ok(Self { entries })
}
pub fn push(&mut self, spec: &IndicatorSpec) -> Result<()> {
self.entries.push(IndicatorEntry {
label: spec.label(),
indicator: build_spec(spec)?,
last: None,
fields: Vec::new(),
series: VecDeque::with_capacity(INDICATOR_SERIES),
});
Ok(())
}
#[must_use]
pub fn wants_book(&self) -> bool {
self.entries.iter().any(|e| e.indicator.wants_book())
}
#[must_use]
pub fn wants_references(&self) -> bool {
self.entries.iter().any(|e| e.indicator.wants_reference())
}
pub fn remove(&mut self, label: &str) -> bool {
let before = self.entries.len();
self.entries.retain(|entry| entry.label != label);
self.entries.len() != before
}
pub fn update(&mut self, input: &TickInput) {
for entry in &mut self.entries {
if let Some(value) = entry.indicator.update(input) {
entry.last = Some(value);
entry.fields = entry.indicator.fields();
}
if let Some(value) = entry.last {
if entry.series.len() == INDICATOR_SERIES {
entry.series.pop_front();
}
entry.series.push_back(value);
}
}
}
#[must_use]
pub fn values(&self) -> Vec<(String, Option<f64>)> {
self.entries
.iter()
.map(|entry| (entry.label.clone(), entry.last))
.collect()
}
#[must_use]
pub fn snapshot(&self) -> Vec<IndicatorReading> {
self.entries
.iter()
.map(|entry| IndicatorReading {
label: entry.label.clone(),
value: entry.last,
fields: entry.fields.clone(),
series: entry.series.iter().copied().collect(),
})
.collect()
}
}
impl Default for IndicatorSet {
fn default() -> Self {
Self::from_specs(&crate::config::default_indicators())
.expect("the default indicator overlay must be constructible")
}
}
#[derive(Debug)]
pub struct SymbolState {
pub book: BookState,
pub tape: TapeRing,
pub footprint: Footprint,
pub indicators: IndicatorSet,
pub profiles: ProfileSet,
pub bars: BarSet,
pub last: Decimal,
pub bid: Decimal,
pub ask: Decimal,
pub volume: Decimal,
pub open: Decimal,
pub history: VecDeque<Decimal>,
pub candles: CandleBuilder,
ohlc: VecDeque<wc::Candle>,
breadth: BreadthState,
derivatives: DerivativesState,
}
impl SymbolState {
pub fn new(
specs: &[IndicatorSpec],
profiles: &[IndicatorSpec],
bars: &[IndicatorSpec],
timeframe: Timeframe,
) -> Result<Self> {
Ok(Self {
book: BookState::default(),
tape: TapeRing::default(),
footprint: Footprint::default(),
indicators: IndicatorSet::from_specs(specs)?,
profiles: ProfileSet::from_specs(profiles)?,
bars: BarSet::from_specs(bars)?,
last: Decimal::ZERO,
bid: Decimal::ZERO,
ask: Decimal::ZERO,
volume: Decimal::ZERO,
open: Decimal::ZERO,
history: VecDeque::with_capacity(PRICE_HISTORY),
candles: CandleBuilder::new(timeframe),
ohlc: VecDeque::with_capacity(OHLC_HISTORY),
breadth: BreadthState::new(),
derivatives: DerivativesState::default(),
})
}
}
impl SymbolState {
pub(crate) fn apply_derivatives(&mut self, update: &DerivativesUpdate) {
self.derivatives.apply(update);
}
}
impl SymbolState {
pub(crate) fn open_at(&mut self, price: Decimal) {
if self.open.is_zero() && !price.is_zero() {
self.open = price;
}
}
pub(crate) fn seed_bars(&mut self, bars: &[wc::Candle]) {
if let Some(first) = bars.first() {
if let Some(open) = Decimal::from_f64_retain(first.open) {
self.open_at(open);
}
}
for bar in bars {
if self.ohlc.len() == OHLC_HISTORY {
self.ohlc.pop_front();
}
self.ohlc.push_back(*bar);
if self.history.len() == PRICE_HISTORY {
self.history.pop_front();
}
if let Some(close) = Decimal::from_f64_retain(bar.close) {
self.history.push_back(close);
self.last = close;
}
self.breadth.update(bar);
let mut tick = TickInput::price(bar.close);
tick.candle = Some(*bar);
self.indicators.update(&tick);
self.profiles.update(&tick);
self.bars.update(&tick);
}
}
#[must_use]
pub fn ohlc(&self, n: usize) -> Vec<wc::Candle> {
let skip = self.ohlc.len().saturating_sub(n);
self.ohlc.iter().skip(skip).copied().collect()
}
#[must_use]
pub fn forming(&self) -> Option<wc::Candle> {
self.candles.partial()
}
#[must_use]
pub fn series(&self, n: usize) -> Vec<f64> {
let skip = self.history.len().saturating_sub(n);
self.history
.iter()
.skip(skip)
.map(|d| d.to_f64().unwrap_or(0.0))
.collect()
}
}
pub type Key = (SourceId, Symbol);
#[derive(Default)]
pub struct AppState {
pub sources: Vec<Box<dyn DataSource>>,
pub symbols: HashMap<Key, SymbolState>,
pub focus: Option<Key>,
pub watchlist: Vec<Key>,
pub indicators: Vec<IndicatorSpec>,
pub profiles: Vec<IndicatorSpec>,
pub bars: Vec<IndicatorSpec>,
pub timeframe: Timeframe,
pub(crate) record_capacity: Option<usize>,
pub(crate) recorded: VecDeque<Event>,
}
impl std::fmt::Debug for AppState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AppState")
.field("sources", &self.sources.len())
.field("symbols", &self.symbols)
.field("focus", &self.focus)
.field("watchlist", &self.watchlist)
.field("indicators", &self.indicators)
.field("profiles", &self.profiles)
.field("bars", &self.bars)
.field("timeframe", &self.timeframe)
.field("record_capacity", &self.record_capacity)
.field("recorded", &self.recorded.len())
.finish()
}
}
impl AppState {
pub fn fold(&mut self, src: SourceId, sym: &Symbol, event: &Event) {
self.fold_scoped(src, sym, event, None);
}
pub(crate) fn fold_scoped(
&mut self,
src: SourceId,
sym: &Symbol,
event: &Event,
reference_scope: Option<SourceId>,
) {
let references = if self.indicators.iter().any(|s| s.reference.is_some()) {
self.reference_prices(reference_scope)
} else {
BTreeMap::new()
};
let cross_section = if self
.indicators
.iter()
.any(|s| registry::is_cross_section(&s.kind))
{
match event {
Event::Trade(print) => self.cross_section(print.timestamp),
_ => None,
}
} else {
None
};
if matches!(event, Event::Trade(print) if print.quantity < Decimal::ZERO) {
return;
}
let state = self.symbols.entry((src, sym.clone())).or_insert_with(|| {
SymbolState::new(&self.indicators, &self.profiles, &self.bars, self.timeframe)
.expect("indicator specs are validated before they reach the state")
});
match event {
Event::Trade(print) => {
let price = print.price.to_f64().unwrap_or(0.0);
state.open_at(print.price);
state.last = print.price;
state.tape.push(print.clone());
state.footprint.add(print);
let closed = state.candles.update(
price,
print.quantity.to_f64().unwrap_or(0.0),
print.timestamp,
);
let book = if state.indicators.wants_book() {
state.book.to_core()
} else {
None
};
if let Some(bar) = closed.as_ref() {
state.breadth.update(bar);
if state.ohlc.len() == OHLC_HISTORY {
state.ohlc.pop_front();
}
state.ohlc.push_back(*bar);
}
state
.derivatives
.add_trade(print.quantity.to_f64().unwrap_or(0.0), print.aggressor);
let mut tick = TickInput::price(price);
tick.candle = closed;
tick.cross_section = cross_section;
tick.derivatives = state.derivatives.tick();
tick.trade = core_trade(print);
tick.trade_quote = tick.trade.and_then(|trade| {
state
.book
.mid()
.and_then(|mid| mid.to_f64())
.and_then(|mid| wc::TradeQuote::new(trade, mid).ok())
});
tick.book = book;
tick.references = references;
state.indicators.update(&tick);
state.profiles.update(&tick);
state.bars.update(&tick);
if state.history.len() == PRICE_HISTORY {
state.history.pop_front();
}
state.history.push_back(print.price);
}
Event::Ticker(ticker) => {
state.open_at(ticker.last);
state.last = ticker.last;
state.bid = ticker.bid;
state.ask = ticker.ask;
state.volume = ticker.volume;
}
Event::BookSnapshot(snap) => state.book.apply_snapshot(snap),
Event::BookDelta(delta) => state.book.apply_delta(delta),
Event::Derivatives(_)
| Event::OrderUpdate(_)
| Event::BalanceUpdate(_)
| Event::Subscribed { .. }
| Event::Disconnected
| Event::Reconnected => {}
}
}
#[must_use]
fn reference_prices(&self, scope: Option<SourceId>) -> BTreeMap<String, f64> {
self.symbols
.iter()
.filter(|((src, _), _)| scope.is_none_or(|id| *src == id))
.filter_map(|((_, symbol), state)| {
state.last.to_f64().map(|price| (symbol.to_string(), price))
})
.filter(|(_, price)| *price > 0.0)
.collect()
}
#[must_use]
fn cross_section(&self, timestamp: i64) -> Option<wc::CrossSection> {
let members: Vec<wc::Member> = self
.symbols
.values()
.filter_map(|state| state.breadth.member())
.collect();
if members.is_empty() {
return None;
}
wc::CrossSection::new(members, timestamp).ok()
}
#[must_use]
pub fn fresh_market(&self) -> SymbolState {
SymbolState::new(&self.indicators, &self.profiles, &self.bars, self.timeframe)
.expect("indicator specs are validated before they reach the state")
}
pub fn add_indicator(&mut self, spec: &IndicatorSpec) -> Result<()> {
let label = spec.label();
if self.indicators.iter().any(|s| s.label() == label) {
return Err(Error::Config(format!("indicator already tracked: {label}")));
}
build_spec(spec)?;
for state in self.symbols.values_mut() {
state.indicators.push(spec)?;
}
self.indicators.push(spec.clone());
Ok(())
}
pub fn set_timeframe(&mut self, timeframe: Timeframe) -> Result<()> {
self.timeframe = timeframe;
for state in self.symbols.values_mut() {
state.candles = CandleBuilder::new(timeframe);
state.indicators = IndicatorSet::from_specs(&self.indicators)?;
state.ohlc.clear();
}
Ok(())
}
pub fn remove_indicator(&mut self, label: &str) -> bool {
let known = self.indicators.iter().any(|s| s.label() == label);
if !known {
return false;
}
self.indicators.retain(|s| s.label() != label);
for state in self.symbols.values_mut() {
state.indicators.remove(label);
}
true
}
pub fn pump(&mut self) -> usize {
let mut batch: Vec<(SourceId, Symbol, Event)> = Vec::new();
for source in &mut self.sources {
let id = source.id();
for (sym, ev) in source.poll() {
batch.push((id, sym, ev));
}
}
let folded = batch.len();
for (id, sym, ev) in batch {
self.record(&ev);
self.fold(id, &sym, &ev);
}
folded
}
fn record(&mut self, event: &Event) {
let Some(capacity) = self.record_capacity else {
return;
};
if self.recorded.len() == capacity {
self.recorded.pop_front();
}
self.recorded.push_back(event.clone());
}
#[must_use]
pub fn recording(&self) -> Vec<Event> {
self.recorded.iter().cloned().collect()
}
pub fn set_recording(&mut self, capacity: Option<usize>) {
self.record_capacity = capacity.map(|c| c.clamp(1, MAX_RECORDING));
self.recorded = VecDeque::new();
}
#[must_use]
pub fn get(&self, key: &Key) -> Option<&SymbolState> {
self.symbols.get(key)
}
pub fn source_mut(&mut self, id: SourceId) -> Option<&mut Box<dyn DataSource>> {
self.sources.iter_mut().find(|s| s.id() == id)
}
pub fn remove_source(&mut self, id: SourceId) {
self.sources.retain(|s| s.id() != id);
self.symbols.retain(|(src, _), _| *src != id);
self.watchlist.retain(|(src, _)| *src != id);
if matches!(&self.focus, Some((src, _)) if *src == id) {
self.focus = self.watchlist.first().cloned();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_ticker_keeps_every_field_it_carries() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
state.fold(
0,
&sym,
&Event::Ticker(Ticker {
symbol: sym.clone(),
last: dec!(20010),
bid: dec!(20009),
ask: dec!(20011),
volume: dec!(1234),
timestamp: 0,
}),
);
let market = state.get(&(0, sym)).expect("the ticker created the market");
assert_eq!(market.last, dec!(20010));
assert_eq!(market.bid, dec!(20009));
assert_eq!(market.ask, dec!(20011));
assert_eq!(market.volume, dec!(1234));
}
#[test]
fn the_open_is_the_first_price_and_stays_there() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
for price in [dec!(100), dec!(110), dec!(90)] {
state.fold(0, &sym, &trade(&sym, price, OrderSide::Buy));
}
let market = state.get(&(0, sym)).expect("the trades created the market");
assert_eq!(market.open, dec!(100));
assert_eq!(market.last, dec!(90));
}
#[test]
fn a_zero_price_does_not_open_the_window() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
state.fold(0, &sym, &trade(&sym, Decimal::ZERO, OrderSide::Buy));
state.fold(0, &sym, &trade(&sym, dec!(100), OrderSide::Buy));
let market = state.get(&(0, sym)).expect("the trades created the market");
assert_eq!(market.open, dec!(100));
}
#[test]
fn seeded_history_sets_the_open_from_its_oldest_bar() {
let bars: Vec<wickra_core::Candle> = [(100.0, 0), (140.0, 1)]
.into_iter()
.map(|(close, ts)| {
wickra_core::Candle::new(close, close, close, close, 1.0, ts)
.expect("an ordered candle")
})
.collect();
let mut state = SymbolState::new(&[], &[], &[], Timeframe::default())
.expect("an empty indicator set is constructible");
state.seed_bars(&bars);
assert!((state.open.to_f64().unwrap_or(0.0) - 100.0).abs() < 1e-9);
}
#[test]
fn the_rings_report_emptiness_before_and_after_a_print() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
let market = SymbolState::new(&[], &[], &[], Timeframe::default())
.expect("an empty indicator set is constructible");
assert!(market.tape.is_empty());
assert!(market.footprint.is_empty());
state.fold(0, &sym, &trade(&sym, dec!(100), OrderSide::Buy));
let market = state.get(&(0, sym)).expect("the trade created the market");
assert!(!market.tape.is_empty());
assert!(!market.footprint.is_empty());
}
#[test]
fn an_indicator_set_says_whether_it_needs_a_reference() {
let plain = IndicatorSet::from_specs(&[IndicatorSpec::new("Sma", vec![3.0])])
.expect("Sma is registered");
assert!(!plain.wants_references());
let paired = IndicatorSet::from_specs(&[IndicatorSpec::paired(
"RollingCorrelation",
vec![20.0],
"ETH/USDT",
)])
.expect("RollingCorrelation is registered");
assert!(paired.wants_references());
}
#[test]
fn a_breadth_bar_with_no_usable_close_is_ignored() {
let mut breadth = BreadthState::new();
let good =
wickra_core::Candle::new(100.0, 101.0, 99.0, 100.0, 5.0, 0).expect("an ordered candle");
breadth.update(&good);
let after_good = breadth.previous_close;
let mut bad = good;
bad.close = f64::NAN;
breadth.update(&bad);
assert_eq!(
breadth.previous_close, after_good,
"a NaN close moved the previous close"
);
assert!(breadth.change.is_finite(), "the change went to NaN");
}
#[test]
fn a_breadth_bar_with_no_usable_volume_reads_as_zero() {
let mut breadth = BreadthState::new();
let mut bar =
wickra_core::Candle::new(100.0, 101.0, 99.0, 100.0, 5.0, 0).expect("an ordered candle");
bar.volume = f64::NEG_INFINITY;
breadth.update(&bar);
assert!(
breadth.volume.abs() < f64::EPSILON,
"volume: {}",
breadth.volume
);
}
#[test]
fn an_indicator_entry_prints_its_label_and_not_its_series() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState {
indicators: vec![IndicatorSpec::new("Sma", vec![2.0])],
..AppState::default()
};
for price in [dec!(100), dec!(101), dec!(102)] {
state.fold(0, &sym, &trade(&sym, price, OrderSide::Buy));
}
let shown = format!("{:?}", state.get(&(0, sym)).expect("the market exists"));
assert!(shown.contains("IndicatorEntry"), "{shown}");
assert!(shown.contains("Sma(2)"), "the label is missing: {shown}");
assert!(
shown.contains("series: 2"),
"the series was not counted: {shown}"
);
}
#[test]
fn seeding_past_the_rings_evicts_the_oldest_bars() {
let bars: Vec<wickra_core::Candle> = (0..600)
.map(|i| {
let close = 100.0 + f64::from(i);
wickra_core::Candle::new(close, close, close, close, 1.0, i64::from(i))
.expect("an ordered candle")
})
.collect();
let mut state = SymbolState::new(&[], &[], &[], Timeframe::default())
.expect("an empty indicator set is constructible");
state.seed_bars(&bars);
assert_eq!(
state.ohlc.len(),
OHLC_HISTORY,
"the bar ring grew past its bound"
);
assert_eq!(
state.history.len(),
512,
"the price ring grew past its bound"
);
let newest = state.ohlc.back().expect("the ring is not empty");
assert!(
(newest.close - 699.0).abs() < 1e-9,
"the newest bar was evicted"
);
let oldest = state.ohlc.front().expect("the ring is not empty");
let first_kept = 700.0 - f64::from(u16::try_from(OHLC_HISTORY).expect("the ring is small"));
assert!((oldest.close - first_kept).abs() < 1e-9);
}
#[test]
fn the_debug_view_reports_the_recording_by_count() {
let mut state = AppState::default();
state.set_recording(Some(16));
let sym = Symbol::new("BTC", "USDT");
let event = trade(&sym, dec!(100), OrderSide::Buy);
state.fold(0, &sym, &event);
state.record(&event);
let shown = format!("{state:?}");
assert!(shown.contains("record_capacity: Some(16)"), "{shown}");
assert!(shown.contains("recorded: 1"), "{shown}");
assert!(shown.contains("sources: 0"), "{shown}");
}
fn candle_tick(open: f64, high: f64, low: f64, close: f64) -> TickInput {
let mut input = TickInput::price(close);
input.candle = Some(
wickra_core::Candle::new(open, high, low, close, 1.0, 0).expect("an ordered candle"),
);
input
}
#[test]
fn a_point_and_figure_column_ignores_a_price_that_is_not_one() {
let mut pnf = PointAndFigure::default();
for bad in [f64::NAN, f64::INFINITY, 0.0, -5.0] {
pnf.update(bad);
assert!(!pnf.started, "{bad} started a column");
}
pnf.update(100.0);
assert!(pnf.started);
}
#[test]
fn taker_flow_ignores_a_quantity_that_is_not_one() {
let mut state = DerivativesState::default();
state.add_trade(f64::NAN, OrderSide::Buy);
state.add_trade(-1.0, OrderSide::Sell);
assert!(state.taker_buy_volume.abs() < 1e-9);
assert!(state.taker_sell_volume.abs() < 1e-9);
state.add_trade(2.5, OrderSide::Buy);
assert!((state.taker_buy_volume - 2.5).abs() < 1e-9);
}
#[test]
fn a_profile_set_knows_whether_it_tracks_anything() {
assert!(ProfileSet::from_specs(&[]).unwrap().is_empty());
let one =
ProfileSet::from_specs(&[IndicatorSpec::new("VolumeProfile", vec![4.0, 8.0])]).unwrap();
assert!(!one.is_empty());
}
#[test]
fn a_profile_entry_debugs_without_printing_its_whole_histogram() {
let set =
ProfileSet::from_specs(&[IndicatorSpec::new("VolumeProfile", vec![4.0, 8.0])]).unwrap();
let shown = format!("{:?}", set.entries[0]);
assert!(shown.contains("ProfileEntry"), "{shown}");
assert!(shown.contains("VolumeProfile(4,8)"), "{shown}");
}
#[test]
fn a_bar_entry_debugs_its_bar_count_rather_than_every_bar() {
let set = BarSet::from_specs(&[IndicatorSpec::new("RenkoBars", vec![3.0])]).unwrap();
let shown = format!("{:?}", set.entries[0]);
assert!(shown.contains("BarEntry"), "{shown}");
assert!(shown.contains("RenkoBars(3)"), "{shown}");
assert!(shown.contains("bars: 0"), "{shown}");
}
#[test]
fn a_bar_stream_keeps_the_most_recent_bars_and_drops_the_oldest() {
let mut set = BarSet::from_specs(&[IndicatorSpec::new("RenkoBars", vec![1.0])]).unwrap();
let mut price = 100.0;
for _ in 0..(ALT_BARS_KEPT * 2) {
price += 2.0;
set.update(&candle_tick(price - 2.0, price, price - 2.0, price));
}
let kept = set.entries[0].bars.len();
assert_eq!(kept, ALT_BARS_KEPT, "kept {kept}");
let last = set.entries[0].bars.back().expect("a brick");
assert!(last.close >= price - 2.0, "{last:?} against {price}");
}
use rust_decimal_macros::dec;
use wickra_exchange_core::{Symbol, Ticker};
fn trade(sym: &Symbol, price: Decimal, side: OrderSide) -> Event {
Event::Trade(stamped(sym, price, side, 0))
}
fn stamped(sym: &Symbol, price: Decimal, side: OrderSide, timestamp: i64) -> TradePrint {
TradePrint {
symbol: sym.clone(),
price,
quantity: dec!(2),
aggressor: side,
timestamp,
}
}
#[test]
fn fold_trade_updates_last_tape_footprint_and_history() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
state.fold(0, &sym, &trade(&sym, dec!(100), OrderSide::Buy));
state.fold(0, &sym, &trade(&sym, dec!(101), OrderSide::Sell));
let st = state.get(&(0, sym.clone())).unwrap();
assert_eq!(st.last, dec!(101));
assert_eq!(st.tape.len(), 2);
assert_eq!(st.footprint.at(dec!(100)), Some((dec!(2), dec!(0))));
assert_eq!(st.footprint.at(dec!(101)), Some((dec!(0), dec!(2))));
assert_eq!(st.series(10), vec![100.0, 101.0]);
}
fn print_at(sym: &Symbol, price: Decimal) -> Event {
Event::Trade(TradePrint {
symbol: sym.clone(),
price,
quantity: Decimal::ONE,
aggressor: OrderSide::Buy,
timestamp: 0,
})
}
#[test]
fn the_footprint_does_not_grow_with_the_session() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
for tick in 0..5_000 {
let price = Decimal::from(100_000 + tick) / dec!(100);
state.fold(0, &sym, &print_at(&sym, price));
}
let st = state.get(&(0, sym)).unwrap();
let levels = st.footprint.len();
assert!(
levels <= MAX_FOOTPRINT_LEVELS,
"the footprint holds {levels} levels"
);
}
#[test]
fn the_footprint_keeps_the_levels_nearest_the_market() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
for tick in 0..5_000 {
let price = Decimal::from(100_000 + tick) / dec!(100);
state.fold(0, &sym, &print_at(&sym, price));
}
let st = state.get(&(0, sym)).unwrap();
let last = st.last;
assert_eq!(last, dec!(1049.99));
assert!(
st.footprint.at(last).is_some(),
"the level just traded was evicted"
);
let ladder = st.footprint.around(last, 3);
assert_eq!(ladder.len(), 3);
for (price, _, _) in &ladder {
let gap = (*price - last).abs();
assert!(gap < dec!(1), "{price} is {gap} away from a last of {last}");
}
}
#[test]
fn around_returns_a_ladder_centred_on_the_anchor() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
for cents in 0..21 {
let price = Decimal::from(10_000 + cents) / dec!(100);
state.fold(0, &sym, &print_at(&sym, price));
}
let st = state.get(&(0, sym)).unwrap();
let ladder = st.footprint.around(dec!(100.10), 5);
let prices: Vec<Decimal> = ladder.iter().map(|(price, _, _)| *price).collect();
assert_eq!(
prices,
vec![
dec!(100.12),
dec!(100.11),
dec!(100.10),
dec!(100.09),
dec!(100.08)
]
);
}
#[test]
fn a_negative_size_print_is_not_folded_at_all() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
state.fold(0, &sym, &trade(&sym, dec!(100), OrderSide::Buy));
let mut bad = trade(&sym, dec!(100), OrderSide::Buy);
if let Event::Trade(print) = &mut bad {
print.quantity = dec!(-2);
}
state.fold(0, &sym, &bad);
let st = state.get(&(0, sym.clone())).unwrap();
assert_eq!(st.footprint.at(dec!(100)), Some((dec!(2), dec!(0))));
assert_eq!(st.tape.len(), 1, "a malformed print reached the tape");
}
#[test]
fn a_rejected_print_does_not_bring_a_market_into_existence() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
let mut bad = trade(&sym, dec!(100), OrderSide::Buy);
if let Event::Trade(print) = &mut bad {
print.quantity = dec!(-1);
}
state.fold(0, &sym, &bad);
assert!(
state.get(&(0, sym)).is_none(),
"an empty market was created"
);
}
#[test]
fn book_snapshot_then_delta_apply() {
let sym = Symbol::new("BTC", "USDT");
let mut book = BookState::default();
book.apply_snapshot(&OrderBookSnapshot {
symbol: sym.clone(),
last_update_id: 1,
bids: vec![
BookLevel::new(dec!(100), dec!(1)),
BookLevel::new(dec!(99), dec!(2)),
],
asks: vec![BookLevel::new(dec!(101), dec!(1))],
timestamp: 0,
});
assert_eq!(book.best_bid(), Some((dec!(100), dec!(1))));
assert_eq!(book.best_ask(), Some((dec!(101), dec!(1))));
assert_eq!(book.spread(), Some(dec!(1)));
book.apply_delta(&BookDelta {
symbol: sym,
first_update_id: 2,
final_update_id: 2,
bids: vec![BookLevel::new(dec!(100), dec!(0))],
asks: vec![BookLevel::new(dec!(102), dec!(3))],
timestamp: 0,
});
assert_eq!(book.best_bid(), Some((dec!(99), dec!(2))));
assert_eq!(
book.top_asks(2),
vec![(dec!(101), dec!(1)), (dec!(102), dec!(3))]
);
}
#[test]
fn tape_ring_respects_cap() {
let sym = Symbol::new("BTC", "USDT");
let mut ring = TapeRing::new(3);
for i in 0..5 {
ring.push(TradePrint {
symbol: sym.clone(),
price: Decimal::from(i),
quantity: dec!(1),
aggressor: OrderSide::Buy,
timestamp: i,
});
}
assert_eq!(ring.len(), 3);
let recent = ring.recent(3);
assert_eq!(recent[0].price, dec!(4));
assert_eq!(recent[2].price, dec!(2));
}
#[test]
fn footprint_add_saturates_on_overflow() {
let sym = Symbol::new("BTC", "USDT");
let mut footprint = Footprint::default();
let huge = |quantity: Decimal| TradePrint {
symbol: sym.clone(),
price: dec!(100),
quantity,
aggressor: OrderSide::Buy,
timestamp: 0,
};
footprint.add(&huge(Decimal::MAX));
footprint.add(&huge(Decimal::MAX));
assert_eq!(footprint.at(dec!(100)), Some((Decimal::MAX, Decimal::ZERO)));
}
#[test]
fn indicator_set_warms_up_then_reports() {
let price = TickInput::price;
let mut set = IndicatorSet::default();
for _ in 0..19 {
set.update(&price(100.0));
}
assert_eq!(set.values()[0].1, None);
set.update(&price(100.0));
assert_eq!(set.values()[0].1, Some(100.0));
}
#[test]
fn an_indicator_series_records_one_point_per_tick_after_warmup() {
let price = TickInput::price;
let mut set = IndicatorSet::from_specs(&[IndicatorSpec::new("Sma", vec![3.0])]).unwrap();
for step in 0..10 {
set.update(&price(100.0 + f64::from(step)));
}
assert_eq!(set.snapshot()[0].series.len(), 8);
}
#[test]
fn a_bar_indicator_series_carries_its_value_forward_between_bars() {
let mut set = IndicatorSet::from_specs(&[IndicatorSpec::new("Atr", vec![2.0])]).unwrap();
let mut builder = CandleBuilder::new(Timeframe::parse("1s").unwrap());
let mut ticks = 0;
for step in 0..40_i64 {
let price = 100.0 + (step % 4) as f64;
let closed = builder.update(price, 1.0, step * 250);
let mut tick = TickInput::price(price);
tick.candle = closed;
set.update(&tick);
ticks += 1;
}
let series = &set.snapshot()[0].series;
assert!(!series.is_empty(), "Atr recorded nothing over ten bars");
let kept = series.len();
let bars = ticks / 4;
assert!(
kept > bars,
"series of {kept} is barely longer than the {bars} bars, so it is not carrying forward"
);
}
#[test]
fn an_indicator_series_is_bounded() {
let price = TickInput::price;
let mut set = IndicatorSet::from_specs(&[IndicatorSpec::new("Sma", vec![2.0])]).unwrap();
for step in 0..500 {
set.update(&price(100.0 + f64::from(step)));
}
assert_eq!(set.snapshot()[0].series.len(), INDICATOR_SERIES);
}
#[test]
fn a_warming_up_indicator_has_no_series_yet() {
let price = TickInput::price;
let mut set = IndicatorSet::from_specs(&[IndicatorSpec::new("Sma", vec![50.0])]).unwrap();
for _ in 0..10 {
set.update(&price(100.0));
}
assert_eq!(set.snapshot()[0].series, Vec::<f64>::new());
}
#[test]
fn indicator_labels_come_from_the_spec() {
let set = IndicatorSet::default();
let labels: Vec<String> = set.values().into_iter().map(|(l, _)| l).collect();
assert_eq!(labels, vec!["Sma(20)".to_string(), "Ema(50)".to_string()]);
}
#[test]
fn an_unknown_indicator_spec_is_rejected_by_name() {
let err = IndicatorSet::from_specs(&[IndicatorSpec::new("NotReal", vec![])])
.expect_err("an unknown indicator must be rejected")
.to_string();
assert!(
err.contains("NotReal"),
"error does not name the spec: {err}"
);
}
#[test]
fn a_configured_indicator_set_replaces_the_default() {
let set = IndicatorSet::from_specs(&[IndicatorSpec::new("Rsi", vec![14.0])]).unwrap();
let labels: Vec<String> = set.values().into_iter().map(|(l, _)| l).collect();
assert_eq!(labels, vec!["Rsi(14)".to_string()]);
}
#[test]
fn a_candle_indicator_advances_only_when_a_bar_closes() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState {
indicators: vec![IndicatorSpec::new("Atr", vec![2.0])],
timeframe: Timeframe::parse("1m").unwrap(),
..AppState::default()
};
for step in 0..30_i64 {
let print = stamped(&sym, dec!(100), OrderSide::Buy, step * 1_000);
state.fold(0, &sym, &Event::Trade(print));
}
let market = state.get(&(0, sym.clone())).unwrap();
assert_eq!(
market.indicators.values()[0].1,
None,
"Atr reported before a bar closed"
);
assert!(
market.candles.partial().is_some(),
"the builder should hold a bar in progress"
);
}
#[test]
fn account_and_lifecycle_events_do_not_change_market_state() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState::default();
state.fold(0, &sym, &trade(&sym, dec!(100), OrderSide::Buy));
let before = state.get(&(0, sym.clone())).unwrap().last;
state.fold(0, &sym, &Event::Disconnected);
state.fold(0, &sym, &Event::BalanceUpdate(vec![]));
let after = state.get(&(0, sym.clone())).unwrap().last;
assert_eq!(before, after);
}
fn breadth_terminal(kind: &str) -> (AppState, Symbol, Symbol) {
let state = AppState {
indicators: vec![IndicatorSpec {
kind: kind.to_string(),
params: Vec::new(),
reference: None,
}],
..Default::default()
};
(
state,
Symbol::new("BTC", "USDT"),
Symbol::new("ETH", "USDT"),
)
}
fn print_in_bar(sym: &Symbol, price: Decimal, bar: i64) -> Event {
Event::Trade(TradePrint {
symbol: sym.clone(),
price,
quantity: dec!(3),
aggressor: OrderSide::Buy,
timestamp: bar * 60_000 + 1,
})
}
#[test]
fn a_breadth_indicator_reads_the_whole_universe() {
let (mut state, btc, eth) = breadth_terminal("AdvanceDecline");
for bar in 0..40 {
let up = Decimal::from(100 + bar);
let down = Decimal::from(100 - bar);
state.fold(0, &btc, &print_in_bar(&btc, up, bar));
state.fold(0, ð, &print_in_bar(ð, down, bar));
}
let reading = state
.get(&(0, btc.clone()))
.expect("BTC is tracked")
.indicators
.values()
.first()
.and_then(|(_, reading)| *reading);
assert!(
reading.is_some(),
"AdvanceDecline produced no reading after 40 bars across two markets"
);
}
#[test]
fn the_universe_is_absent_until_a_bar_closes() {
let (mut state, btc, _eth) = breadth_terminal("AdvanceDecline");
state.fold(0, &btc, &print_in_bar(&btc, dec!(100), 0));
assert!(state.cross_section(0).is_none());
}
#[test]
fn the_universe_is_not_assembled_when_nothing_reads_it() {
let mut state = AppState::default();
let btc = Symbol::new("BTC", "USDT");
for bar in 0..5 {
state.fold(0, &btc, &print_in_bar(&btc, Decimal::from(100 + bar), bar));
}
assert!(
!state
.indicators
.iter()
.any(|s| registry::is_cross_section(&s.kind)),
"the default overlay should carry no breadth indicator"
);
}
#[test]
fn a_market_that_has_not_closed_a_bar_is_left_out_of_the_universe() {
let (mut state, btc, eth) = breadth_terminal("AdvanceDecline");
for bar in 0..4 {
state.fold(0, &btc, &print_in_bar(&btc, Decimal::from(100 + bar), bar));
}
state.fold(0, ð, &print_in_bar(ð, dec!(50), 0));
let universe = state.cross_section(1).expect("BTC has closed bars");
assert_eq!(universe.members.len(), 1);
}
#[test]
fn breadth_is_folded_per_bar_not_per_tick() {
let (mut state, btc, _eth) = breadth_terminal("AdvanceDecline");
for step in 0..4 {
state.fold(0, &btc, &print_in_bar(&btc, Decimal::from(100 + step), 0));
}
assert!(
state.cross_section(0).is_none(),
"no bar has closed yet, so there is nothing to read"
);
state.fold(0, &btc, &print_in_bar(&btc, dec!(110), 1));
assert!(
state.cross_section(1).is_some(),
"the first bar closed, so the market is now a member"
);
}
#[test]
fn the_taker_flow_comes_from_the_tape_not_the_derivatives_feed() {
let mut derivatives = DerivativesState::default();
derivatives.apply(&DerivativesUpdate {
mark_price: Some(20_000.0),
index_price: Some(20_000.0),
futures_price: Some(20_050.0),
timestamp: 1,
..DerivativesUpdate::default()
});
for _ in 0..3 {
derivatives.add_trade(2.0, OrderSide::Buy);
}
derivatives.add_trade(2.0, OrderSide::Sell);
let tick = derivatives.tick().expect("the three prices have arrived");
assert!(
(tick.taker_buy_volume - 6.0).abs() < 1e-9,
"three buys of two should be six, got {}",
tick.taker_buy_volume
);
assert!(
(tick.taker_sell_volume - 2.0).abs() < 1e-9,
"one sell of two should be two, got {}",
tick.taker_sell_volume
);
}
#[test]
fn a_derivatives_tick_needs_its_prices_first() {
let mut derivatives = DerivativesState::default();
derivatives.apply(&DerivativesUpdate {
funding_rate: Some(0.0001),
timestamp: 1,
..DerivativesUpdate::default()
});
assert!(derivatives.tick().is_none());
}
#[test]
fn liquidations_accumulate_and_other_channels_replace() {
let mut derivatives = DerivativesState::default();
let priced = DerivativesUpdate {
mark_price: Some(20_000.0),
index_price: Some(20_000.0),
futures_price: Some(20_050.0),
timestamp: 1,
..DerivativesUpdate::default()
};
derivatives.apply(&priced);
for _ in 0..3 {
derivatives.apply(&DerivativesUpdate {
long_liquidation: Some(1_000.0),
open_interest: Some(500.0),
..DerivativesUpdate::default()
});
}
let tick = derivatives.tick().expect("priced");
assert!((tick.long_liquidation - 3_000.0).abs() < 1e-9);
assert!((tick.open_interest - 500.0).abs() < 1e-9);
}
#[test]
fn a_trade_quote_pairs_the_print_with_the_standing_mid() {
let sym = Symbol::new("BTC", "USDT");
let mut book = BookState::default();
book.apply_snapshot(&OrderBookSnapshot {
symbol: sym.clone(),
bids: vec![BookLevel {
price: dec!(99),
quantity: dec!(5),
}],
asks: vec![BookLevel {
price: dec!(101),
quantity: dec!(5),
}],
last_update_id: 1,
timestamp: 0,
});
assert_eq!(book.mid(), Some(dec!(100)));
}
#[test]
fn a_one_sided_book_has_no_mid_to_measure_against() {
let sym = Symbol::new("BTC", "USDT");
let mut book = BookState::default();
book.apply_snapshot(&OrderBookSnapshot {
symbol: sym,
bids: vec![BookLevel {
price: dec!(99),
quantity: dec!(5),
}],
asks: Vec::new(),
last_update_id: 1,
timestamp: 0,
});
assert!(book.mid().is_none());
}
#[test]
fn a_microstructure_indicator_reads_prints_against_the_book() {
let sym = Symbol::new("BTC", "USDT");
let mut state = AppState {
indicators: vec![IndicatorSpec {
kind: "EffectiveSpread".to_string(),
params: Vec::new(),
reference: None,
}],
..Default::default()
};
for step in 0..40 {
let mid = Decimal::from(100 + step % 3);
state.fold(
0,
&sym,
&Event::BookSnapshot(OrderBookSnapshot {
symbol: sym.clone(),
bids: vec![BookLevel {
price: mid - dec!(1),
quantity: dec!(5),
}],
asks: vec![BookLevel {
price: mid + dec!(1),
quantity: dec!(5),
}],
last_update_id: u64::try_from(step).unwrap_or(0),
timestamp: 0,
}),
);
state.fold(0, &sym, &trade(&sym, mid + dec!(1), OrderSide::Buy));
}
let reading = state
.get(&(0, sym))
.expect("BTC is tracked")
.indicators
.values()
.first()
.and_then(|(_, reading)| *reading);
assert!(
reading.is_some(),
"EffectiveSpread produced no reading after 40 prints against a two-sided book"
);
}
fn pnf(closes: &[f64]) -> PointAndFigure {
let mut chart = PointAndFigure::default();
for close in closes {
chart.update(*close);
}
chart
}
#[test]
fn a_rising_column_clearing_the_previous_one_is_a_buy_signal() {
assert!(
!pnf(&[100.0, 102.0]).on_buy_signal,
"no previous column to clear yet"
);
assert!(pnf(&[100.0, 102.0, 98.0, 102.0, 104.0]).on_buy_signal);
}
#[test]
fn a_falling_column_undercutting_the_previous_one_takes_it_away() {
let chart = pnf(&[100.0, 102.0, 98.0, 102.0, 104.0]);
assert!(chart.on_buy_signal, "the setup should be on a buy signal");
let chart = pnf(&[100.0, 102.0, 98.0, 102.0, 104.0, 100.0, 104.0, 100.0, 96.0]);
assert!(!chart.on_buy_signal);
}
#[test]
fn a_move_smaller_than_a_box_does_not_advance_the_column() {
let chart = pnf(&[100.0, 100.5, 100.9, 100.4, 100.8]);
assert!(
(chart.extreme - 100.0).abs() < 1e-9,
"the column did not move"
);
assert!(chart.rising);
}
#[test]
fn a_counter_move_smaller_than_the_reversal_does_not_start_a_column() {
let chart = pnf(&[100.0, 104.0, 102.0]);
assert!(
chart.rising,
"a two-box pullback must not reverse the column"
);
assert!((chart.extreme - 104.0).abs() < 1e-9);
}
#[test]
fn the_breadth_member_reports_the_point_and_figure_signal() {
let mut breadth = BreadthState::new();
let bar = |close: f64| {
wc::Candle::new(close, close, close, close, 10.0, 0)
.expect("a flat synthetic bar is valid")
};
for close in [100.0, 102.0, 98.0, 102.0, 104.0] {
breadth.update(&bar(close));
}
let member = breadth.member().expect("bars have closed");
assert!(member.on_buy_signal);
}
}