use ahash::{AHashMap, AHashSet};
use nautilus_common::messages::DataEvent;
use nautilus_model::{
identifiers::InstrumentId,
instruments::{Instrument, InstrumentAny},
};
fn economics_differ(a: &InstrumentAny, b: &InstrumentAny) -> bool {
a.maker_fee() != b.maker_fee()
|| a.taker_fee() != b.taker_fee()
|| a.margin_init() != b.margin_init()
|| a.margin_maint() != b.margin_maint()
|| a.price_precision() != b.price_precision()
|| a.size_precision() != b.size_precision()
|| a.price_increment() != b.price_increment()
|| a.size_increment() != b.size_increment()
|| a.multiplier() != b.multiplier()
|| a.lot_size() != b.lot_size()
|| a.min_quantity() != b.min_quantity()
|| a.max_quantity() != b.max_quantity()
|| a.min_notional() != b.min_notional()
|| a.max_notional() != b.max_notional()
|| a.min_price() != b.min_price()
|| a.max_price() != b.max_price()
}
pub fn diff_and_emit_instruments(
new_instruments: &[InstrumentAny],
cached: &mut AHashMap<InstrumentId, InstrumentAny>,
subscriptions: Option<&AHashSet<InstrumentId>>,
sender: &tokio::sync::mpsc::UnboundedSender<DataEvent>,
) {
let is_subscribed = |id: &InstrumentId| subscriptions.is_none_or(|subs| subs.contains(id));
for instrument in new_instruments {
let id = instrument.id();
let changed = cached
.get(&id)
.is_none_or(|prev| economics_differ(prev, instrument));
if changed {
cached.insert(id, instrument.clone());
if is_subscribed(&id)
&& let Err(e) = sender.send(DataEvent::Instrument(instrument.clone()))
{
log::error!("Failed to emit instrument event: {e}");
}
}
}
}
#[cfg(test)]
mod tests {
use nautilus_core::UnixNanos;
use nautilus_model::{
identifiers::{InstrumentId, Symbol},
instruments::{CryptoPerpetual, InstrumentAny},
types::{Currency, Money, Price, Quantity},
};
use rstest::rstest;
use rust_decimal::Decimal;
use rust_decimal_macros::dec;
use super::*;
fn perp(
maker_fee: Decimal,
taker_fee: Decimal,
size_increment: Quantity,
min_notional: Option<Money>,
) -> InstrumentAny {
InstrumentAny::CryptoPerpetual(CryptoPerpetual::new(
InstrumentId::from("BTCUSDT-LINEAR.BYBIT"),
Symbol::from("BTCUSDT-LINEAR"),
Currency::BTC(),
Currency::USDT(),
Currency::USDT(),
false, 1, 3, Price::from("0.1"),
size_increment,
None, None, None, Some(Quantity::from("0.001")), None, min_notional,
None, None, None, None, Some(maker_fee),
Some(taker_fee),
None, None, UnixNanos::default(), UnixNanos::default(), ))
}
fn default_perp() -> InstrumentAny {
perp(dec!(0.0001), dec!(0.00055), Quantity::from("0.001"), None)
}
#[rstest]
fn test_emits_for_new_instrument() {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let instrument = default_perp();
let mut cached = AHashMap::new();
diff_and_emit_instruments(std::slice::from_ref(&instrument), &mut cached, None, &tx);
match rx.try_recv().expect("expected instrument event") {
DataEvent::Instrument(emitted) => assert_eq!(emitted.id(), instrument.id()),
_ => panic!("expected Instrument event"),
}
assert!(cached.contains_key(&instrument.id()));
}
#[rstest]
fn test_no_emit_when_unchanged() {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let instrument = default_perp();
let mut cached = AHashMap::new();
cached.insert(instrument.id(), default_perp());
diff_and_emit_instruments(&[instrument], &mut cached, None, &tx);
assert!(rx.try_recv().is_err());
}
#[rstest]
fn test_emits_on_fee_change() {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let id = default_perp().id();
let mut cached = AHashMap::new();
cached.insert(id, default_perp());
let updated = perp(dec!(0.0001), dec!(0.0008), Quantity::from("0.001"), None);
diff_and_emit_instruments(&[updated], &mut cached, None, &tx);
match rx
.try_recv()
.expect("expected instrument event on fee change")
{
DataEvent::Instrument(emitted) => assert_eq!(emitted.taker_fee(), dec!(0.0008)),
_ => panic!("expected Instrument event"),
}
assert_eq!(cached.get(&id).unwrap().taker_fee(), dec!(0.0008));
}
#[rstest]
fn test_emits_on_size_increment_and_min_notional_change() {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let id = default_perp().id();
let mut cached = AHashMap::new();
cached.insert(id, default_perp());
let bigger_step = perp(dec!(0.0001), dec!(0.00055), Quantity::from("0.002"), None);
diff_and_emit_instruments(&[bigger_step], &mut cached, None, &tx);
assert!(rx.try_recv().is_ok(), "size_increment change should emit");
let mut cached = AHashMap::new();
cached.insert(id, default_perp());
let with_min = perp(
dec!(0.0001),
dec!(0.00055),
Quantity::from("0.001"),
Some(Money::new(5.0, Currency::USDT())),
);
diff_and_emit_instruments(&[with_min], &mut cached, None, &tx);
assert!(rx.try_recv().is_ok(), "min_notional change should emit");
}
#[rstest]
fn test_subscription_gating_updates_cache_but_only_emits_for_subscribed() {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let id = default_perp().id();
let empty_subs = AHashSet::new();
let mut cached = AHashMap::new();
cached.insert(id, default_perp());
let updated = perp(dec!(0.0001), dec!(0.0008), Quantity::from("0.001"), None);
diff_and_emit_instruments(&[updated], &mut cached, Some(&empty_subs), &tx);
assert!(rx.try_recv().is_err(), "unsubscribed should not emit");
assert_eq!(
cached.get(&id).unwrap().taker_fee(),
dec!(0.0008),
"cache should update regardless of subscription"
);
let mut subs = AHashSet::new();
subs.insert(id);
let mut cached = AHashMap::new();
cached.insert(id, default_perp());
let updated = perp(dec!(0.0001), dec!(0.0008), Quantity::from("0.001"), None);
diff_and_emit_instruments(&[updated], &mut cached, Some(&subs), &tx);
assert!(rx.try_recv().is_ok(), "subscribed should emit");
}
}