use std::collections::HashMap;
use optionstratlib::ExpirationDate;
use tokio::sync::mpsc;
use super::App;
use crate::chain::{DepthLadder, GreeksRow, InstrumentKey, MarketUpdate, QuoteUpdate};
use crate::event::{AppEvent, Command};
use crate::providers::MarketUpdateSink;
pub const CONTROL_CHANNEL_CAPACITY: usize = 64;
pub const COMMAND_CHANNEL_CAPACITY: usize = 64;
#[derive(Debug, Default)]
struct StagedInstrument {
quote: Option<QuoteUpdate>,
greeks: Option<GreeksRow>,
depth: Option<DepthLadder>,
}
#[derive(Debug, Default)]
struct StagingMap {
slots: HashMap<InstrumentKey, StagedInstrument>,
}
impl StagingMap {
fn new() -> Self {
Self {
slots: HashMap::new(),
}
}
fn stage(&mut self, update: MarketUpdate) -> Option<MarketUpdate> {
match update {
MarketUpdate::Quote(quote) => {
if let Some(slot) = self.slot_mut("e.instrument.key) {
slot.quote = Some(quote);
}
None
}
MarketUpdate::Greeks(greeks) => {
if let Some(slot) = self.slot_mut(&greeks.instrument.key) {
slot.greeks = Some(greeks);
}
None
}
MarketUpdate::Depth(depth) => {
if let Some(slot) = self.slot_mut(&depth.instrument.key) {
slot.depth = Some(depth);
}
None
}
control @ (MarketUpdate::Chain(_) | MarketUpdate::Health(_, _)) => Some(control),
}
}
fn slot_mut(&mut self, key: &InstrumentKey) -> Option<&mut StagedInstrument> {
if !self.slots.contains_key(key) {
let _ = self.slots.insert(key.clone(), StagedInstrument::default());
}
self.slots.get_mut(key)
}
fn remove_subscription(&mut self, underlying: &str, expiration: &ExpirationDate) {
self.slots
.retain(|key, _| !key_in_subscription(key, underlying, expiration));
}
fn drain_into<F: FnMut(MarketUpdate)>(&mut self, sink: &mut F) {
for (_key, staged) in self.slots.drain() {
let StagedInstrument {
quote,
greeks,
depth,
} = staged;
if let Some(quote) = quote {
sink(MarketUpdate::Quote(quote));
}
if let Some(greeks) = greeks {
sink(MarketUpdate::Greeks(greeks));
}
if let Some(depth) = depth {
sink(MarketUpdate::Depth(depth));
}
}
}
#[cfg(any(test, feature = "bench"))]
fn len(&self) -> usize {
self.slots.len()
}
#[cfg(any(test, feature = "bench"))]
fn capacity(&self) -> usize {
self.slots.capacity()
}
}
fn key_in_subscription(key: &InstrumentKey, underlying: &str, expiration: &ExpirationDate) -> bool {
if key.underlying != underlying {
return false;
}
match expiration {
ExpirationDate::DateTime(instant) => key.expiration_utc == *instant,
ExpirationDate::Days(_) => true,
}
}
#[derive(Debug, Clone)]
pub struct BridgeSenders {
pub tx_control: mpsc::Sender<MarketUpdate>,
pub tx_coalesced: mpsc::Sender<MarketUpdate>,
pub tx_command: mpsc::Sender<Command>,
}
impl BridgeSenders {
#[must_use]
pub fn market_update_sink(&self) -> MarketUpdateSink {
MarketUpdateSink::new(self.tx_control.clone(), self.tx_coalesced.clone())
}
}
#[derive(Debug)]
pub struct EventBridge {
rx_control: mpsc::Receiver<MarketUpdate>,
rx_coalesced: mpsc::Receiver<MarketUpdate>,
rx_command: mpsc::Receiver<Command>,
staging: StagingMap,
}
impl EventBridge {
#[must_use = "the returned BridgeSenders must be wired to the data layer and App"]
pub fn new(coalesced_capacity: usize) -> (Self, BridgeSenders) {
let coalesced_capacity = coalesced_capacity.max(1);
let (tx_control, rx_control) = mpsc::channel(CONTROL_CHANNEL_CAPACITY);
let (tx_coalesced, rx_coalesced) = mpsc::channel(coalesced_capacity);
let (tx_command, rx_command) = mpsc::channel(COMMAND_CHANNEL_CAPACITY);
let bridge = Self {
rx_control,
rx_coalesced,
rx_command,
staging: StagingMap::new(),
};
let senders = BridgeSenders {
tx_control,
tx_coalesced,
tx_command,
};
(bridge, senders)
}
pub fn pump<R: FnMut(Command)>(&mut self, app: &mut App, route: R) {
self.pump_into(|update| app.on_event(AppEvent::Market(update)), route);
}
pub fn pump_into<S, R>(&mut self, mut sink: S, mut route: R)
where
S: FnMut(MarketUpdate),
R: FnMut(Command),
{
self.drain_control(&mut sink);
self.coalesce(&mut sink);
self.drain_commands(&mut route);
self.flush(&mut sink);
}
fn drain_control<F: FnMut(MarketUpdate)>(&mut self, sink: &mut F) {
while let Ok(update) = self.rx_control.try_recv() {
sink(update);
}
}
fn coalesce<F: FnMut(MarketUpdate)>(&mut self, sink: &mut F) {
while let Ok(update) = self.rx_coalesced.try_recv() {
if let Some(control) = self.staging.stage(update) {
sink(control);
}
}
}
fn drain_commands<R: FnMut(Command)>(&mut self, route: &mut R) {
while let Ok(command) = self.rx_command.try_recv() {
if let Command::Unsubscribe {
underlying,
expiration,
} = &command
{
self.staging.remove_subscription(underlying, expiration);
}
route(command);
}
}
fn flush<F: FnMut(MarketUpdate)>(&mut self, sink: &mut F) {
self.staging.drain_into(sink);
}
}
#[cfg(any(test, feature = "bench"))]
impl EventBridge {
pub(crate) fn staged_len(&self) -> usize {
self.staging.len()
}
pub(crate) fn staged_capacity(&self) -> usize {
self.staging.capacity()
}
pub(crate) fn coalesce_pending(&mut self) {
let mut discard = |_update: MarketUpdate| {};
self.coalesce(&mut discard);
}
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use chrono::{DateTime, Utc};
use optionstratlib::chains::OptionData;
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::Positive;
use optionstratlib::{ExpirationDate, OptionStyle};
use super::{BridgeSenders, EventBridge, StagingMap};
use crate::app::{App, LiveState, Mode, ReplayState, SourceBinding};
use crate::chain::{
AliasCatalog, ChainFetch, ChainSource, ChainStore, ContractSpecFingerprint, DepthLadder,
ExerciseStyle, ExpirySource, GreeksOrigin, GreeksRow, Instrument, InstrumentKey,
MarketUpdate, ProviderId, QuoteUpdate, SettlementStyle, StreamHealth,
};
use crate::config::ThemeChoice;
use crate::event::Command;
use crate::providers::{
ChainCapability, ChainPollCapability, GreeksCapability, ProviderCapabilities,
};
const EXP: i64 = 1_700_000_000;
const EXP2: i64 = 1_700_086_400;
#[track_caller]
fn pid(id: &str) -> ProviderId {
match ProviderId::new(id) {
Ok(p) => p,
Err(e) => panic!("expected a valid provider id `{id}`, got: {e}"),
}
}
#[track_caller]
fn utc(secs: i64) -> DateTime<Utc> {
match DateTime::<Utc>::from_timestamp(secs, 0) {
Some(t) => t,
None => panic!("invalid test timestamp: {secs}"),
}
}
#[track_caller]
fn pos(value: f64) -> Positive {
match Positive::new(value) {
Ok(p) => p,
Err(e) => panic!("invalid test positive `{value}`: {e}"),
}
}
fn spec() -> ContractSpecFingerprint {
ContractSpecFingerprint {
contract_multiplier: 1,
settlement: SettlementStyle::Cash,
exercise: ExerciseStyle::European,
quote_currency: "USD".to_owned(),
venue_product_code: "BTC".to_owned(),
}
}
fn ikey(underlying: &str, exp: i64, strike: f64) -> InstrumentKey {
InstrumentKey {
underlying: underlying.to_owned(),
expiration_utc: utc(exp),
strike: pos(strike),
style: OptionStyle::Call,
}
}
fn instrument(underlying: &str, exp: i64, strike: f64) -> Instrument {
Instrument {
key: ikey(underlying, exp, strike),
provider: pid("deribit"),
native_symbol: format!("{underlying}-{strike}"),
stream_symbol: None,
spec: spec(),
}
}
fn quote(underlying: &str, exp: i64, strike: f64, bid: f64, received: i64) -> MarketUpdate {
MarketUpdate::Quote(QuoteUpdate {
instrument: instrument(underlying, exp, strike),
bid: Some(pos(bid)),
ask: Some(pos(bid + 0.2)),
last: None,
bid_size: None,
ask_size: None,
event_time: None,
received_time: utc(received),
})
}
fn q(strike: f64, bid: f64, received: i64) -> MarketUpdate {
quote("BTC", EXP, strike, bid, received)
}
fn greeks(strike: f64, iv: f64) -> MarketUpdate {
MarketUpdate::Greeks(GreeksRow {
instrument: instrument("BTC", EXP, strike),
iv: Some(pos(iv)),
delta: None,
gamma: None,
theta: None,
vega: None,
rho: None,
origin: GreeksOrigin::Provider,
event_time: None,
received_time: utc(EXP),
})
}
fn depth(strike: f64) -> MarketUpdate {
MarketUpdate::Depth(DepthLadder {
instrument: instrument("BTC", EXP, strike),
bids: Vec::new(),
asks: Vec::new(),
event_time: None,
received_time: utc(EXP),
change_id: None,
})
}
fn health() -> MarketUpdate {
MarketUpdate::Health(pid("deribit"), StreamHealth::Reconnecting { attempt: 1 })
}
fn unsubscribe(underlying: &str, expiration: ExpirationDate) -> Command {
Command::Unsubscribe {
underlying: underlying.to_owned(),
expiration,
}
}
fn quote_bid(update: &MarketUpdate) -> Option<Positive> {
match update {
MarketUpdate::Quote(q) => q.bid,
MarketUpdate::Greeks(_)
| MarketUpdate::Depth(_)
| MarketUpdate::Chain(_)
| MarketUpdate::Health(_, _) => None,
}
}
fn update_strike(update: &MarketUpdate) -> Option<Positive> {
match update {
MarketUpdate::Quote(q) => Some(q.instrument.key.strike),
MarketUpdate::Greeks(g) => Some(g.instrument.key.strike),
MarketUpdate::Depth(d) => Some(d.instrument.key.strike),
MarketUpdate::Chain(_) | MarketUpdate::Health(_, _) => None,
}
}
fn is_health(update: &MarketUpdate) -> bool {
matches!(update, MarketUpdate::Health(_, _))
}
fn drain(map: &mut StagingMap) -> Vec<MarketUpdate> {
let mut out = Vec::new();
{
let mut sink = |update| out.push(update);
map.drain_into(&mut sink);
}
out
}
fn full_caps() -> ProviderCapabilities {
ProviderCapabilities::builder()
.chain(ChainCapability::Assemble)
.depth(true)
.greeks(GreeksCapability::Provided)
.chain_poll(ChainPollCapability::Poll {
interval_hint_secs: 2,
})
.build()
}
fn row(strike: f64) -> OptionData {
let mut od = OptionData {
strike_price: pos(strike),
call_bid: Some(pos(1.0)),
call_ask: Some(pos(1.2)),
put_bid: Some(pos(2.0)),
put_ask: Some(pos(2.4)),
implied_volatility: pos(0.5),
..Default::default()
};
od.set_mid_prices();
od
}
fn chain_with(strikes: &[f64]) -> OptionChain {
let mut chain = OptionChain::new("BTC", pos(60_000.0), "2025-06-27".to_owned(), None, None);
for strike in strikes {
let _ = chain.options.insert(row(*strike));
}
chain
}
fn store(strikes: &[f64]) -> ChainStore {
ChainStore::seed(
ChainFetch::new(
chain_with(strikes),
ExpirySource::new("BTC", utc(EXP), pid("deribit")),
AliasCatalog::new(),
),
ChainSource::Merged,
std::time::Duration::from_secs(2),
utc(EXP),
)
}
fn live_app_with_bridge() -> (App, EventBridge, BridgeSenders) {
let (bridge, senders) = EventBridge::new(64);
let live = LiveState::new(
SourceBinding::new(pid("deribit"), full_caps(), StreamHealth::Live),
store(&[60_000.0]),
);
let app = App::new(
Mode::Live(live),
ThemeChoice::Auto,
senders.tx_command.clone(),
);
(app, bridge, senders)
}
#[test]
fn test_staging_map_second_quote_same_key_overwrites_first() {
let mut map = StagingMap::new();
let _ = map.stage(q(100.0, 1.0, EXP));
let _ = map.stage(q(100.0, 2.0, EXP + 1));
assert_eq!(map.len(), 1, "one slot per instrument, overwrite in place");
let out = drain(&mut map);
assert_eq!(out.len(), 1, "the two quotes coalesced to one");
assert_eq!(
out.first().and_then(quote_bid),
Some(pos(2.0)),
"the second update for the same key wins"
);
}
#[test]
fn test_staging_map_quote_and_greeks_same_key_both_kept() {
let mut map = StagingMap::new();
let _ = map.stage(q(100.0, 1.0, EXP));
let _ = map.stage(greeks(100.0, 0.5));
assert_eq!(map.len(), 1, "still one slot per instrument");
let out = drain(&mut map);
assert_eq!(out.len(), 2, "both the quote and the Greeks flush");
assert!(out.iter().any(|u| matches!(u, MarketUpdate::Quote(_))));
assert!(out.iter().any(|u| matches!(u, MarketUpdate::Greeks(_))));
}
#[test]
fn test_staging_map_depth_and_quote_same_key_both_kept() {
let mut map = StagingMap::new();
let _ = map.stage(q(100.0, 1.0, EXP));
let _ = map.stage(depth(100.0));
assert_eq!(map.len(), 1);
let out = drain(&mut map);
assert_eq!(out.len(), 2);
assert!(out.iter().any(|u| matches!(u, MarketUpdate::Depth(_))));
}
#[test]
fn test_staging_map_stage_returns_control_update_for_direct_fold() {
let mut map = StagingMap::new();
let returned = map.stage(health());
assert!(matches!(returned, Some(MarketUpdate::Health(_, _))));
assert_eq!(map.len(), 0, "control class is never staged");
}
#[test]
fn test_staging_map_remove_subscription_by_datetime_prunes_matching_expiry() {
let mut map = StagingMap::new();
let _ = map.stage(quote("BTC", EXP, 100.0, 1.0, EXP));
let _ = map.stage(quote("BTC", EXP2, 200.0, 1.0, EXP));
assert_eq!(map.len(), 2);
map.remove_subscription("BTC", &ExpirationDate::DateTime(utc(EXP)));
assert_eq!(map.len(), 1, "only the EXP expiry is pruned");
let out = drain(&mut map);
assert_eq!(out.first().and_then(update_strike), Some(pos(200.0)));
}
#[test]
fn test_staging_map_remove_subscription_days_prunes_whole_underlying() {
let mut map = StagingMap::new();
let _ = map.stage(quote("BTC", EXP, 100.0, 1.0, EXP));
let _ = map.stage(quote("BTC", EXP2, 200.0, 1.0, EXP));
let _ = map.stage(quote("ETH", EXP, 300.0, 1.0, EXP));
assert_eq!(map.len(), 3);
map.remove_subscription("BTC", &ExpirationDate::Days(pos(7.0)));
assert_eq!(map.len(), 1, "both BTC expiries pruned, ETH kept");
let out = drain(&mut map);
assert_eq!(out.first().and_then(update_strike), Some(pos(300.0)));
}
#[test]
fn test_staging_map_remove_subscription_other_underlying_is_noop() {
let mut map = StagingMap::new();
let _ = map.stage(q(100.0, 1.0, EXP));
map.remove_subscription("ETH", &ExpirationDate::DateTime(utc(EXP)));
assert_eq!(
map.len(),
1,
"unsubscribing a different underlying keeps BTC"
);
}
#[test]
fn test_staging_map_drain_retains_capacity_across_bursts() {
let mut map = StagingMap::new();
for strike in 1..=8 {
let _ = map.stage(q(f64::from(strike), 1.0, EXP));
}
let capacity_after_first = map.capacity();
assert!(capacity_after_first >= 8);
let flushed = drain(&mut map);
assert_eq!(flushed.len(), 8);
assert_eq!(map.len(), 0);
assert_eq!(
map.capacity(),
capacity_after_first,
"drain must retain the allocation (HP-3)"
);
for round in 0..1_000 {
for strike in 1..=8 {
let _ = map.stage(q(f64::from(strike), f64::from(round % 5) + 1.0, EXP));
}
assert!(map.len() <= 8, "staging is O(N instruments), not O(burst)");
let _ = drain(&mut map);
assert_eq!(
map.capacity(),
capacity_after_first,
"no per-burst reallocation on the hot path (HP-3)"
);
}
}
#[test]
fn test_staging_map_latest_value_wins_over_a_burst() {
let mut map = StagingMap::new();
for tick in 0..500 {
let _ = map.stage(q(100.0, f64::from(tick) + 1.0, EXP + i64::from(tick)));
}
assert_eq!(map.len(), 1);
let out = drain(&mut map);
assert_eq!(out.len(), 1);
assert_eq!(
out.first().and_then(quote_bid),
Some(pos(500.0)),
"the last-staged value is the one delivered"
);
}
#[test]
fn test_event_bridge_drains_control_before_coalesced() {
let (mut bridge, senders) = EventBridge::new(64);
let _ = senders.tx_coalesced.try_send(q(100.0, 1.0, EXP));
let _ = senders.tx_control.try_send(health());
let mut recorded: Vec<MarketUpdate> = Vec::new();
bridge.pump_into(|u| recorded.push(u), |_c| {});
assert!(
recorded.first().is_some_and(is_health),
"control (Health) must be folded before any coalesced quote"
);
assert!(
recorded.iter().any(|u| matches!(u, MarketUpdate::Quote(_))),
"the coalesced quote is still delivered, just after control"
);
}
#[test]
fn test_event_bridge_health_delivered_while_coalesced_saturated() {
let (mut bridge, senders) = EventBridge::new(4);
for tick in 0..64 {
let strike = f64::from(tick % 3) + 100.0;
let _ = senders.tx_coalesced.try_send(q(
strike,
f64::from(tick) + 1.0,
EXP + i64::from(tick),
));
}
let _ = senders.tx_control.try_send(health());
let mut recorded: Vec<MarketUpdate> = Vec::new();
bridge.pump_into(|u| recorded.push(u), |_c| {});
assert!(
recorded.first().is_some_and(is_health),
"Health is delivered promptly, not behind the quote burst"
);
let quote_count = recorded
.iter()
.filter(|u| matches!(u, MarketUpdate::Quote(_)))
.count();
assert!(
quote_count <= 3,
"quotes coalesced to at most N=3 instruments, got {quote_count}"
);
}
#[test]
fn test_event_bridge_burst_beyond_channel_capacity_keeps_memory_flat() {
let (mut bridge, senders) = EventBridge::new(4);
let mut capacity_baseline: Option<usize> = None;
for round in 0..500u32 {
for tick in 0..9u32 {
let strike = f64::from(tick % 3) + 100.0;
let _ = senders.tx_coalesced.try_send(q(
strike,
f64::from(round) + 1.0,
EXP + i64::from(round),
));
}
bridge.coalesce_pending();
assert!(
bridge.staged_len() <= 3,
"staging is O(N=3 instruments), not O(burst): round {round}"
);
let baseline = *capacity_baseline.get_or_insert_with(|| bridge.staged_capacity());
assert_eq!(
bridge.staged_capacity(),
baseline,
"the staging allocation never grows unboundedly: round {round}"
);
let mut sink = |_u: MarketUpdate| {};
bridge.flush(&mut sink);
}
}
#[test]
fn test_event_bridge_every_instrument_receives_latest_value() {
let (mut bridge, senders) = EventBridge::new(64);
for round in 1..=5 {
for strike in [100.0, 200.0, 300.0] {
let _ = senders.tx_coalesced.try_send(q(
strike,
f64::from(round),
EXP + i64::from(round),
));
}
}
let mut recorded: Vec<MarketUpdate> = Vec::new();
bridge.pump_into(|u| recorded.push(u), |_c| {});
assert_eq!(recorded.len(), 3, "one current value per instrument");
for update in &recorded {
assert_eq!(quote_bid(update), Some(pos(5.0)), "the freshest bid wins");
}
let mut strikes: Vec<Positive> = recorded.iter().filter_map(update_strike).collect();
strikes.sort();
assert_eq!(strikes, vec![pos(100.0), pos(200.0), pos(300.0)]);
}
#[test]
fn test_event_bridge_unsubscribe_command_prunes_staged_instrument() {
let (mut bridge, senders) = EventBridge::new(64);
let _ = senders
.tx_coalesced
.try_send(quote("BTC", EXP, 100.0, 1.0, EXP));
let _ = senders
.tx_coalesced
.try_send(quote("BTC", EXP2, 200.0, 1.0, EXP));
let _ = senders
.tx_command
.try_send(unsubscribe("BTC", ExpirationDate::DateTime(utc(EXP))));
let mut recorded: Vec<MarketUpdate> = Vec::new();
let mut routed: Vec<Command> = Vec::new();
bridge.pump_into(|u| recorded.push(u), |c| routed.push(c));
assert_eq!(recorded.len(), 1, "the unsubscribed instrument was pruned");
assert_eq!(recorded.first().and_then(update_strike), Some(pos(200.0)));
assert_eq!(
routed.len(),
1,
"the command is still routed to the data layer"
);
assert!(matches!(routed.first(), Some(Command::Unsubscribe { .. })));
}
#[test]
fn test_event_bridge_command_routed_to_router() {
let (mut bridge, senders) = EventBridge::new(64);
let _ = senders.tx_command.try_send(Command::Reconnect);
let mut routed: Vec<Command> = Vec::new();
bridge.pump_into(|_u| {}, |c| routed.push(c));
assert!(matches!(routed.first(), Some(Command::Reconnect)));
}
#[test]
fn test_event_bridge_pump_folds_control_into_live_app() {
let (mut app, mut bridge, senders) = live_app_with_bridge();
app.mark_drawn();
assert!(!app.dirty);
let _ = senders.tx_control.try_send(health());
bridge.pump(&mut app, |_c| {});
assert!(app.dirty, "the folded Health marked the app dirty");
}
#[test]
fn test_event_bridge_pump_on_replay_app_ignores_market_update() {
let (mut bridge, senders) = EventBridge::new(64);
let mut app = App::new(
Mode::Replay(ReplayState::new(PathBuf::from("/bundle"))),
ThemeChoice::Auto,
senders.tx_command.clone(),
);
app.mark_drawn();
let _ = senders.tx_coalesced.try_send(q(100.0, 1.0, EXP));
bridge.pump(&mut app, |_c| {});
assert!(
!app.dirty,
"a live market update is meaningless in replay mode"
);
}
}