mod orders;
mod remote;
mod requests;
mod specs;
use std::{
sync::{
atomic::{self, AtomicBool},
Arc,
},
marker::PhantomData, time::Instant,
};
use apple_quant_core::{
log::{error, info},
UnwrapOptionExt,
};
use databento::dbn::decode::AsyncDbnDecoder;
use rand::rngs::SmallRng;
use rand_distr::Normal;
use smallvec::SmallVec;
use time::Duration;
use tokio::{fs::File, io::BufReader, task::JoinHandle};
use tracing::instrument;
use crate::{
aggregation::{TradeTradeTimestamp, TradeTradeTimestampListInline},
backend::{
DataBackend, LocalOrderId, MarketDataDecoder, MarketDataDecoderProvider,
OrdersBackend, OrdersBackendSubmitRemoteOrderError, OrdersBackendUpdate,
OrdersBackendUpdateRecycle, RemoteOrderId,
},
order::{OrderStateUpdate, RemoteOrder},
order_manager::{CancelRemoteOrderError, OrdersCapacitySpec},
timestamp::{Timestamp, TradeTimestamp, TradeTimestamped},
instrument::InstrumentSpec, strategy::Strategy, volume::DirectionlessVolume,
};
pub use specs::*;
use orders::*;
use remote::*;
use requests::*;
pub(self) struct MarketDataStream<
IS: InstrumentSpec,
MDD: MarketDataDecoder<IS>,
> {
market_data_decoder: Option<MDD>,
_is: PhantomData<IS>,
}
impl<
IS: InstrumentSpec,
MDD: MarketDataDecoder<IS>,
> MarketDataStream<IS, MDD> {
pub(self) async fn decode_tick(
&mut self,
) -> Result<
Option<(
TradeTradeTimestampListInline<IS, 8>,
TradeTimestamp,
)>,
(),
> {
let Some(
market_data_decoder,
) = &mut self.market_data_decoder else {
return Ok(None);
};
let Some((
trades,
newest_trade_timestamp,
)) = market_data_decoder.decode_trade_list().await else {
info!("No trades left to decode.");
return Err(());
};
debug_assert!(!trades.is_empty());
Ok(Some((trades, newest_trade_timestamp)))
}
}
pub struct SimulatedOrdersBackend<
IS: InstrumentSpec,
CapacitySpec: specs::CapacitySpec = CapacitySpecLowFrequencyTrades,
PhysicalLocationLD: specs::PhysicalLocationLD = PhysicalLocationLDProximityBareMetal,
OrderRouterLD: specs::OrderRouterLD = OrderRounterLDRithmicFull,
ExchangeLD: specs::ExchangeLD = ExchangeLD_CME,
> {
rng: SmallRng,
normal_ns: Normal<f64>,
min_ns: u32,
max_ns: u32,
link_request_sender: thingbuf::mpsc::Sender<RemoteRequest<IS>>,
remote_join_handle: JoinHandle<()>,
_capacity_spec: PhantomData<CapacitySpec>,
_physical_location_ld: PhantomData<PhysicalLocationLD>,
_order_router_ld: PhantomData<OrderRouterLD>,
_exchange_ld: PhantomData<ExchangeLD>,
}
impl<
IS: InstrumentSpec,
CapacitySpec: specs::CapacitySpec,
PhysicalLocationLD: specs::PhysicalLocationLD,
OrderRouterLD: specs::OrderRouterLD,
ExchangeLD: specs::ExchangeLD,
> OrdersBackend<IS> for SimulatedOrdersBackend<
IS,
CapacitySpec,
PhysicalLocationLD,
OrderRouterLD,
ExchangeLD,
> {
async fn new<
OrdersCS: OrdersCapacitySpec,
MDDP: MarketDataDecoderProvider<IS>,
>(
market_data_decoder: Option<MDDP::MarketDataDecoder>,
) -> (
thingbuf::mpsc::Receiver<OrdersBackendUpdate<IS>, OrdersBackendUpdateRecycle>,
Self,
) {
let mean_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.mean_ns +
OrderRouterLD::LATENCY_DISTRIBUTION.mean_ns +
ExchangeLD::LATENCY_DISTRIBUTION.mean_ns;
let std_dev_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.std_dev_ns +
OrderRouterLD::LATENCY_DISTRIBUTION.std_dev_ns +
ExchangeLD::LATENCY_DISTRIBUTION.std_dev_ns;
let min_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.min_ns +
OrderRouterLD::LATENCY_DISTRIBUTION.min_ns +
ExchangeLD::LATENCY_DISTRIBUTION.min_ns;
let max_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.max_ns +
OrderRouterLD::LATENCY_DISTRIBUTION.max_ns +
ExchangeLD::LATENCY_DISTRIBUTION.max_ns;
let (
link_request_sender,
link_request_receiver,
) = thingbuf::mpsc::channel(64);
let (
orders_backend_update_sender,
orders_backend_update_receiver,
) = thingbuf::mpsc::with_recycle(1, OrdersBackendUpdateRecycle);
let market_data_stream = MarketDataStream {
market_data_decoder,
_is: PhantomData::default(),
};
let remote_join_handle = tokio::spawn(async move {
remote_link::<IS, MDDP::MarketDataDecoder, CapacitySpec>(
link_request_receiver,
orders_backend_update_sender,
market_data_stream,
).await;
});
let simulated_orders_backend = Self {
rng: rand::make_rng(),
normal_ns: Normal::new(mean_ns as f64, std_dev_ns as f64).unwrap(),
min_ns,
max_ns,
link_request_sender,
remote_join_handle,
_capacity_spec: PhantomData::default(),
_physical_location_ld: PhantomData::default(),
_order_router_ld: PhantomData::default(),
_exchange_ld: PhantomData::default(),
};
(orders_backend_update_receiver, simulated_orders_backend)
}
#[instrument(skip_all)]
async fn submit_order(
&mut self,
local_order_id: &LocalOrderId,
client_submission_timestamp: &Timestamp,
remote_order: RemoteOrder<IS>,
) -> Result<(), OrdersBackendSubmitRemoteOrderError> {
let remote_request = RemoteRequest::SubmitOrder {
local_order_id: *local_order_id,
remote_order,
};
self.link_request_sender.send(remote_request).await.unwrap();
Ok(())
}
#[instrument(skip_all)]
async fn modify_order_volume(
&mut self,
remote_order_id: &RemoteOrderId,
volume: &DirectionlessVolume<IS>,
) {
let instant = Instant::now() + Duration::milliseconds(1);
let remote_request = RemoteRequest::ModifyVolume {
remote_order_id: *remote_order_id,
volume: *volume,
};
self.link_request_sender.send(remote_request).await.unwrap();
}
#[instrument(skip_all)]
async fn initiate_cancel_order(
&mut self,
remote_order_id: &RemoteOrderId,
) -> Result<(), CancelRemoteOrderError> {
let instant = Instant::now() + Duration::milliseconds(1);
let remote_request = RemoteRequest::CancelOrder(*remote_order_id);
self.link_request_sender.send(remote_request).await.unwrap();
Ok(())
}
#[instrument(skip_all)]
async fn initiate_cancel_all(
&mut self,
) {
let instant = Instant::now() + Duration::milliseconds(1);
let remote_request = RemoteRequest::CancelAll;
self.link_request_sender.send(remote_request).await.unwrap();
}
#[instrument(skip_all)]
async fn initiate_cancel_flatten_all(
&mut self,
) {
let instant = Instant::now() + Duration::milliseconds(1);
let remote_request = RemoteRequest::CancelFlattenAll;
self.link_request_sender.send(remote_request).await.unwrap();
}
}