Skip to main content

apple_quant_algorithmic/backends/
simulated.rs

1mod orders;
2mod remote;
3mod requests;
4mod specs;
5
6use std::{
7	sync::{
8		atomic::{self, AtomicBool},
9		Arc,
10	},
11	marker::PhantomData, time::Instant,
12};
13
14use apple_quant_core::{
15	log::{error, info},
16	UnwrapOptionExt,
17};
18
19use databento::dbn::decode::AsyncDbnDecoder;
20use rand::rngs::SmallRng;
21use rand_distr::Normal;
22use smallvec::SmallVec;
23use time::Duration;
24
25use tokio::{fs::File, io::BufReader, task::JoinHandle};
26
27use tracing::instrument;
28
29use crate::{
30	aggregation::{TradeTradeTimestamp, TradeTradeTimestampListInline},
31	backend::{
32		DataBackend, LocalOrderId, MarketDataDecoder, MarketDataDecoderProvider,
33		OrdersBackend, OrdersBackendSubmitRemoteOrderError, OrdersBackendUpdate,
34		OrdersBackendUpdateRecycle, RemoteOrderId,
35	},
36	order::{OrderStateUpdate, RemoteOrder},
37	order_manager::{CancelRemoteOrderError, OrdersCapacitySpec},
38	timestamp::{Timestamp, TradeTimestamp, TradeTimestamped},
39	instrument::InstrumentSpec, strategy::Strategy, volume::DirectionlessVolume,
40};
41
42pub use specs::*;
43use orders::*;
44use remote::*;
45use requests::*;
46
47pub(self) struct MarketDataStream<
48	IS: InstrumentSpec,
49	MDD: MarketDataDecoder<IS>,
50> {
51	market_data_decoder: Option<MDD>,
52	_is: PhantomData<IS>,
53}
54
55impl<
56	IS: InstrumentSpec,
57	MDD: MarketDataDecoder<IS>,
58> MarketDataStream<IS, MDD> {
59	pub(self) async fn decode_tick(
60		&mut self,
61	) -> Result<
62		Option<(
63			TradeTradeTimestampListInline<IS, 8>,
64			TradeTimestamp,
65		)>,
66		(),
67	> {
68		let Some(
69			market_data_decoder,
70		) = &mut self.market_data_decoder else {
71			return Ok(None);
72		};
73
74		let Some((
75			trades,
76			newest_trade_timestamp,
77		)) = market_data_decoder.decode_trade_list().await else {
78			info!("No trades left to decode.");
79			return Err(());
80		};
81
82		debug_assert!(!trades.is_empty());
83		Ok(Some((trades, newest_trade_timestamp)))
84	}
85}
86
87pub struct SimulatedOrdersBackend<
88	IS: InstrumentSpec,
89	CapacitySpec: specs::CapacitySpec = CapacitySpecLowFrequencyTrades,
90	PhysicalLocationLD: specs::PhysicalLocationLD = PhysicalLocationLDProximityBareMetal,
91	OrderRouterLD: specs::OrderRouterLD = OrderRounterLDRithmicFull,
92	ExchangeLD: specs::ExchangeLD = ExchangeLD_CME,
93> {
94	rng: SmallRng,
95	normal_ns: Normal<f64>,
96	min_ns: u32,
97	max_ns: u32,
98	link_request_sender: thingbuf::mpsc::Sender<RemoteRequest<IS>>,
99	remote_join_handle: JoinHandle<()>,
100	_capacity_spec: PhantomData<CapacitySpec>,
101	_physical_location_ld: PhantomData<PhysicalLocationLD>,
102	_order_router_ld: PhantomData<OrderRouterLD>,
103	_exchange_ld: PhantomData<ExchangeLD>,
104}
105
106impl<
107	IS: InstrumentSpec,
108	CapacitySpec: specs::CapacitySpec,
109	PhysicalLocationLD: specs::PhysicalLocationLD,
110	OrderRouterLD: specs::OrderRouterLD,
111	ExchangeLD: specs::ExchangeLD,
112> OrdersBackend<IS> for SimulatedOrdersBackend<
113	IS,
114	CapacitySpec,
115	PhysicalLocationLD,
116	OrderRouterLD,
117	ExchangeLD,
118> {
119	async fn new<
120		OrdersCS: OrdersCapacitySpec,
121		MDDP: MarketDataDecoderProvider<IS>,
122	>(
123		market_data_decoder: Option<MDDP::MarketDataDecoder>,
124	) -> (
125		thingbuf::mpsc::Receiver<OrdersBackendUpdate<IS>, OrdersBackendUpdateRecycle>,
126		Self,
127	) {
128		let mean_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.mean_ns +
129			OrderRouterLD::LATENCY_DISTRIBUTION.mean_ns +
130			ExchangeLD::LATENCY_DISTRIBUTION.mean_ns;
131
132		let std_dev_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.std_dev_ns +
133			OrderRouterLD::LATENCY_DISTRIBUTION.std_dev_ns +
134			ExchangeLD::LATENCY_DISTRIBUTION.std_dev_ns;
135
136		let min_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.min_ns +
137			OrderRouterLD::LATENCY_DISTRIBUTION.min_ns +
138			ExchangeLD::LATENCY_DISTRIBUTION.min_ns;
139
140		let max_ns = PhysicalLocationLD::LATENCY_DISTRIBUTION.max_ns +
141			OrderRouterLD::LATENCY_DISTRIBUTION.max_ns +
142			ExchangeLD::LATENCY_DISTRIBUTION.max_ns;
143
144		let (
145			link_request_sender,
146			link_request_receiver,
147		) = thingbuf::mpsc::channel(64);
148
149		let (
150			orders_backend_update_sender,
151			orders_backend_update_receiver,
152		) = thingbuf::mpsc::with_recycle(1, OrdersBackendUpdateRecycle);
153
154		let market_data_stream = MarketDataStream {
155			market_data_decoder,
156			_is: PhantomData::default(),
157		};
158
159		let remote_join_handle = tokio::spawn(async move {
160			remote_link::<IS, MDDP::MarketDataDecoder, CapacitySpec>(
161				link_request_receiver,
162				orders_backend_update_sender,
163				market_data_stream,
164			).await;
165		});
166
167		let simulated_orders_backend = Self {
168			rng: rand::make_rng(),
169			normal_ns: Normal::new(mean_ns as f64, std_dev_ns as f64).unwrap(),
170			min_ns,
171			max_ns,
172			link_request_sender,
173			remote_join_handle,
174			_capacity_spec: PhantomData::default(),
175			_physical_location_ld: PhantomData::default(),
176			_order_router_ld: PhantomData::default(),
177			_exchange_ld: PhantomData::default(),
178		};
179
180		(orders_backend_update_receiver, simulated_orders_backend)
181	}
182
183	#[instrument(skip_all)]
184	async fn submit_order(
185		&mut self,
186		local_order_id: &LocalOrderId,
187		client_submission_timestamp: &Timestamp,
188		remote_order: RemoteOrder<IS>,
189	) -> Result<(), OrdersBackendSubmitRemoteOrderError> {
190		let remote_request = RemoteRequest::SubmitOrder {
191			local_order_id: *local_order_id,
192			remote_order,
193		};
194
195		self.link_request_sender.send(remote_request).await.unwrap();
196		Ok(())
197	}
198
199	#[instrument(skip_all)]
200	async fn modify_order_volume(
201		&mut self,
202		remote_order_id: &RemoteOrderId,
203		volume: &DirectionlessVolume<IS>,
204	) {
205		let instant = Instant::now() + Duration::milliseconds(1);
206
207		let remote_request = RemoteRequest::ModifyVolume {
208			remote_order_id: *remote_order_id,
209			volume: *volume,
210		};
211
212		self.link_request_sender.send(remote_request).await.unwrap();
213	}
214
215	#[instrument(skip_all)]
216	async fn initiate_cancel_order(
217		&mut self,
218		remote_order_id: &RemoteOrderId,
219	) -> Result<(), CancelRemoteOrderError> {
220		let instant = Instant::now() + Duration::milliseconds(1);
221		let remote_request = RemoteRequest::CancelOrder(*remote_order_id);
222
223		self.link_request_sender.send(remote_request).await.unwrap();
224
225		Ok(())
226	}
227
228	#[instrument(skip_all)]
229	async fn initiate_cancel_all(
230		&mut self,
231	) {
232		let instant = Instant::now() + Duration::milliseconds(1);
233		let remote_request = RemoteRequest::CancelAll;
234
235		self.link_request_sender.send(remote_request).await.unwrap();
236	}
237
238	#[instrument(skip_all)]
239	async fn initiate_cancel_flatten_all(
240		&mut self,
241	) {
242		let instant = Instant::now() + Duration::milliseconds(1);
243		let remote_request = RemoteRequest::CancelFlattenAll;
244
245		self.link_request_sender.send(remote_request).await.unwrap();
246	}
247}