Skip to main content

apple_quant_algorithmic/
frontend.rs

1mod cli;
2
3use std::{path::PathBuf, range::RangeInclusive};
4
5use apple_quant_core::{
6	log::{error, info, initialize_tracing, instrument, trace_span, CtxExt},
7	AnyhowResultExt,
8};
9
10use clap::Parser;
11
12use time::{format_description::well_known::Iso8601, Date, Duration, Time, UtcDateTime};
13
14use tokio::{fs::File, io::BufReader};
15
16#[allow(unused_imports)]
17use tracing::trace;
18
19use crate::{
20	backend::{
21		DataBackend, HistoricalDataBackend, MarketDataDecoderProvider, OrderIdGenerator,
22		OrdersBackend, OrdersBackendUpdate, OrdersBackendUpdateRecycle,
23	},
24	instrument::{InstrumentData, InstrumentKind, InstrumentSpec, InstrumentTicker},
25	order::{ApplyOrderStateUpdate, DeferredOrderActions, OrderStateUpdate, StateGoal},
26	order_manager::{OrderManager, OrdersCapacitySpec},
27	timestamp::{TickTimestamp, TimestampRangeExclusive, Timestamped, UtcNs},
28	strategy::Strategy, volume::DirectionalExposure,
29};
30
31use cli::*;
32
33pub async fn run<
34	'instrument_data: 'aggregated_data,
35	'aggregated_data,
36	IS: InstrumentSpec,
37	OB: OrdersBackend<IS>,
38	OrdersCS: OrdersCapacitySpec,
39	S: Strategy<IS, OB, OrdersCS>,
40	DB: DataBackend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
41>(
42	strategy: S,
43) {
44	initialize_tracing();
45
46	let CLI {
47		key,
48		command,
49	} = CLI::parse();
50
51	#[cfg(feature = "log-trace-frontend")]
52	trace!("CLI parsed with key: `{key:?}` command: `{command:?}`.");
53
54	match command {
55		Commands::Backtest {
56			start,
57			last,
58		} => {
59			let date_range = parse_inclusive_date_range(&start, &last).unwrap();
60
61			#[cfg(feature = "log-trace-frontend")]
62			trace!("Parsed date range: `{date_range:?}`.");
63
64			backtest::<IS, OB, OrdersCS, S, DB>(
65				strategy,
66				key.as_ref().map(|string| string.as_str()),
67				date_range,
68			).await;
69		},
70		Commands::Fetch {
71			start,
72			last,
73		} => {
74			let date_range = parse_inclusive_date_range(&start, &last).unwrap();
75
76			#[cfg(feature = "log-trace-frontend")]
77			trace!("Parsed date range: `{date_range:?}`.");
78
79			fetch::<IS, OB, OrdersCS, S, DB>(
80				strategy,
81				key.as_ref().map(|string| string.as_str()),
82				date_range,
83			).await;
84		},
85	}
86}
87
88#[instrument(skip_all)]
89async fn backtest<
90	'instrument_data: 'aggregated_data,
91	'aggregated_data,
92	IS: InstrumentSpec,
93	OB: OrdersBackend<IS>,
94	OrdersCS: OrdersCapacitySpec,
95	S: Strategy<IS, OB, OrdersCS>,
96	DB: DataBackend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
97>(
98	strategy: S,
99	key: Option<&str>,
100	date_range_inclusive: RangeInclusive<Date>,
101) {
102	let start = unsafe {
103		UtcNs::from_utc_date_time(
104			&UtcDateTime::new(date_range_inclusive.start, Time::MIDNIGHT),
105		).into_timestamp()
106	};
107
108	let end = unsafe {
109		UtcNs::from_utc_date_time(
110			&UtcDateTime::new(
111				date_range_inclusive.last + Duration::days(1),
112				Time::MIDNIGHT,
113			),
114		).into_timestamp()
115	};
116
117	let walk_range = TimestampRangeExclusive { start, end };
118
119	#[cfg(feature = "log-trace-frontend")]
120	trace!("Backtesting with walk range: `{walk_range:?}`.");
121
122	let RangeInclusive {
123		start: start_julian_day,
124		last: last_julian_day,
125	} = julian_days_from_date(&date_range_inclusive).unwrap();
126
127	let day_count = (last_julian_day - start_julian_day) as usize + 1;
128	let mut buf_readers = Vec::with_capacity(day_count);
129
130	for julian_day in start_julian_day..=last_julian_day {
131		let mut iso_string_date = iso_string_date_from_julian(julian_day).unwrap();
132
133		iso_string_date.push_str(DB::file_name_postpend());
134
135		let mut path_buf = PathBuf::new();
136
137		path_buf.push("market-data");
138		path_buf.push(IS::INSTRUMENT_KIND.as_str());
139		path_buf.push(IS::BARE_SYMBOL);
140		path_buf.push("trades");
141		path_buf.push(iso_string_date);
142
143		#[cfg(feature = "log-trace-frontend")]
144		trace!("Opening market data file \"{}\".", path_buf.display());
145
146		let Ok(
147			file,
148		) = File::open(&path_buf).await else {
149			error!("Unable to open market data file: \"{:?}\".", path_buf);
150			return;
151		};
152
153		#[cfg(feature = "log-trace-frontend")]
154		trace!("File handle loaded successfully");
155
156		buf_readers.push(BufReader::new(file));
157	}
158
159	let market_data_decoder = DB::decode_market_data(
160		buf_readers.into_iter(),
161	).await;
162
163	let (
164		frontend,
165		orders_backend_update_receiver,
166	) = Frontend::<'_, '_, IS, OB, OrdersCS, S>::new::<DB>(
167		None,
168		strategy,
169		Some(market_data_decoder),
170	).await;
171
172	let mut data_backend = DB::new(key, frontend, orders_backend_update_receiver);
173	data_backend.backtest(walk_range).await;
174}
175
176#[instrument(skip_all)]
177async fn fetch<
178	'instrument_data: 'aggregated_data,
179	'aggregated_data,
180	IS: InstrumentSpec,
181	OB: OrdersBackend<IS>,
182	OrdersCS: OrdersCapacitySpec,
183	S: Strategy<IS, OB, OrdersCS>,
184	DB: DataBackend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
185>(
186	strategy: S,
187	key: Option<&str>,
188	date_range_inclusive: RangeInclusive<Date>,
189) {
190	let (
191		frontend,
192		orders_backend_update_receiver,
193	) = Frontend::<'_, '_, IS, OB, OrdersCS, S>::new::<DB>(
194		None,
195		strategy,
196		None,
197	).await;
198
199	let mut data_backend = DB::new(key, frontend, orders_backend_update_receiver);
200	let julian_days = julian_days_from_date(&date_range_inclusive).unwrap();
201
202	for julian_day in julian_days.start..=julian_days.last {
203		let Ok(
204			date,
205		) = date_from_julian(julian_day) else {
206			return;
207		};
208
209		let instrument_ticker = InstrumentTicker::new(
210			IS::BARE_SYMBOL,
211			InstrumentKind::Future,
212		);
213
214		let historical_data_backend = data_backend.historical_mut().unwrap().unwrap();
215
216		historical_data_backend
217			.fetch_save_day(&instrument_ticker, date).await
218			.into_iter().for_each(drop);
219	}
220}
221
222pub struct Frontend<
223	'instrument_data,
224	'aggregated_data,
225	IS: InstrumentSpec,
226	OB: OrdersBackend<IS>,
227	OrdersCS: OrdersCapacitySpec,
228	S: Strategy<IS, OB, OrdersCS>,
229> {
230	instrument_data: InstrumentData<'instrument_data, 'aggregated_data, IS>,
231	orders_backend: OB,
232	order_manager: OrderManager<IS, OrdersCS>,
233	order_id_generator: OrderIdGenerator,
234	deferred_order_actions: DeferredOrderActions<IS>,
235	directional_exposure: DirectionalExposure<IS>,
236	state_goal: StateGoal,
237	strategy: S,
238}
239
240impl<
241	'instrument_data,
242	'aggregated_data,
243	IS: InstrumentSpec,
244	OB: OrdersBackend<IS>,
245	OrdersCS: OrdersCapacitySpec,
246	S: Strategy<IS, OB, OrdersCS>,
247> Frontend<
248	'instrument_data,
249	'aggregated_data,
250	IS,
251	OB,
252	OrdersCS,
253	S,
254> {
255	async fn new<MDDP: MarketDataDecoderProvider<IS>>(
256		orders_key: Option<&str>,
257		strategy: S,
258		market_data_decoder: Option<MDDP::MarketDataDecoder>,
259	) -> (
260		Self,
261		thingbuf::mpsc::Receiver<OrdersBackendUpdate<IS>, OrdersBackendUpdateRecycle>,
262	) {
263		let (
264			orders_backend_updates,
265			orders_backend,
266		) = OB::new::<OrdersCS, MDDP>(market_data_decoder).await;
267
268		let frontend = Self {
269			instrument_data: InstrumentData::default(),
270			orders_backend,
271			order_manager: OrderManager::default(),
272			order_id_generator: OrderIdGenerator::new(),
273			deferred_order_actions: DeferredOrderActions::default(),
274			directional_exposure: DirectionalExposure::default(),
275			state_goal: StateGoal::default(),
276			strategy,
277		};
278
279		(frontend, orders_backend_updates)
280	}
281
282	pub(crate) async fn initialize(
283		&mut self,
284		tick_timestamp: &TickTimestamp,
285	) {
286		#[cfg(feature = "log-trace-frontend")]
287		trace!("Initializing walk.");
288
289		self.strategy
290			.initialize(
291				tick_timestamp,
292				&self.instrument_data,
293				&mut self.deferred_order_actions,
294				&mut self.order_id_generator,
295			).await
296			.ctx("While initializing strategy.")
297			.now_checked().unwrap();
298	}
299
300	#[inline(always)]
301	pub(crate) fn tick(
302		&mut self,
303		orders_backend_updates: &mut thingbuf::mpsc::Receiver<
304			OrdersBackendUpdate<IS>,
305			OrdersBackendUpdateRecycle,
306		>,
307		tick_timestamp: &mut TickTimestamp,
308	) -> impl Future<Output = Result<(), ()>> {
309		tick_frontend(self, orders_backend_updates, tick_timestamp)
310	}
311
312	#[instrument(skip_all)]
313	#[inline(always)]
314	async fn apply_order_state_update(
315		&mut self,
316		order_state_update: &OrderStateUpdate<IS>,
317		tick_timestamp: &TickTimestamp,
318	) {
319		order_state_update
320			.apply(
321				tick_timestamp,
322				&mut self.order_manager,
323				&mut self.orders_backend,
324				&mut self.order_id_generator,
325				&mut self.directional_exposure,
326				&mut self.deferred_order_actions,
327				&mut self.state_goal,
328			).await
329			.unwrap_or_else(|e| {
330				error!("{e:?}");
331				panic!();
332			});
333	}
334
335	/// Returns true if polling needs to be enabled.
336	#[instrument(skip_all)]
337	#[inline(always)]
338	async fn drain_deferred_order_actions(
339		&mut self,
340		tick_timestamp: &TickTimestamp,
341	) -> bool {
342		self.deferred_order_actions
343			.trigger(
344				tick_timestamp,
345				&mut self.order_manager,
346				&mut self.orders_backend,
347				&mut self.order_id_generator,
348				&mut self.directional_exposure,
349				&mut self.state_goal,
350			).await
351			.unwrap_or_else(|e| {
352				error!("{e:?}");
353				panic!()
354			})
355	}
356
357	#[instrument(skip_all)]
358	#[inline(always)]
359	async fn tick_strategy(
360		&mut self,
361		tick_timestamp: &TickTimestamp,
362	) {
363		unsafe {
364			self.strategy
365				.data_update(
366					tick_timestamp,
367					&self.instrument_data,
368					&self.order_manager,
369					&self.directional_exposure,
370					&self.state_goal,
371					&mut self.deferred_order_actions,
372					&mut self.order_id_generator,
373				).await
374				.ctx("While strategy data update.").now_unchecked_release();
375		}
376	}
377}
378
379#[instrument(skip_all)]
380#[inline(always)]
381async fn tick_frontend<
382	IS: InstrumentSpec,
383	OB: OrdersBackend<IS>,
384	OrdersCS: OrdersCapacitySpec,
385	S: Strategy<IS, OB, OrdersCS>,
386>(
387	frontend: &mut Frontend<'_, '_, IS, OB, OrdersCS, S>,
388	orders_backend_updates: &mut thingbuf::mpsc::Receiver<
389		OrdersBackendUpdate<IS>,
390		OrdersBackendUpdateRecycle,
391	>,
392	tick_timestamp: &mut TickTimestamp,
393) -> Result<(), ()> {
394	let span = trace_span!("polling orders backend updates");
395	let span = span.enter();
396
397	#[cfg(feature = "log-trace-frontend")]
398	trace!("polling orders backend updates");
399
400	let recv_ref = match orders_backend_updates.try_recv_ref() {
401		Ok(
402			recv_ref,
403		) => recv_ref,
404		Err(thingbuf::mpsc::errors::TryRecvError::Empty) => return Ok(()),
405		Err(thingbuf::mpsc::errors::TryRecvError::Closed) => {
406			info!("Orders backend updates channel closed. Terminating program.");
407			panic!();
408		},
409		_ => return Err(()),
410	};
411
412	drop(span);
413
414	let span = trace_span!("processing new trades");
415	let span = span.enter();
416
417	#[cfg(feature = "log-trace-frontend")]
418	trace!("processing new trades");
419
420	let OrdersBackendUpdate {
421		trades,
422		order_state_updates,
423	} = &*recv_ref;
424
425	frontend.instrument_data.walk_new_trades(trades.iter().cloned());
426
427	if let Some(
428		timestamp,
429	) = trades.last().map(|
430		trade_trade_timestamp,
431	| trade_trade_timestamp.timestamp()) {
432		*tick_timestamp = TickTimestamp::from_timestamp(timestamp);
433	}
434
435	drop(span);
436
437	let span = trace_span!("order state updates");
438	let span = span.enter();
439
440	#[cfg(feature = "log-trace-frontend")]
441	trace!("order state updates");
442
443	for order_state_update in order_state_updates {
444		frontend.apply_order_state_update(
445			order_state_update,
446			tick_timestamp,
447		).await;
448	}
449
450	drop(span);
451
452	frontend.drain_deferred_order_actions(tick_timestamp).await;
453
454	let span = trace_span!("processing client orders");
455	let span = span.enter();
456
457	#[cfg(feature = "log-trace-frontend")]
458	trace!("processing client orders");
459
460	frontend.order_manager.pending_client_orders.tick_client_orders(
461		&mut frontend.deferred_order_actions,
462		&mut frontend.directional_exposure,
463		frontend.instrument_data.iter_recent_latest_trades_forward(),
464	);
465
466	frontend.drain_deferred_order_actions(tick_timestamp).await;
467
468	drop(span);
469
470	let span = trace_span!("processing strategy");
471	let span = span.enter();
472
473	#[cfg(feature = "log-trace-frontend")]
474	trace!("processing strategy");
475
476	frontend.tick_strategy(&tick_timestamp).await;
477	frontend.drain_deferred_order_actions(tick_timestamp).await;
478
479	drop(span);
480
481	Ok(())
482}
483
484pub(crate) fn julian_days_from_date(
485	date_range_inclusive: &RangeInclusive<Date>,
486) -> Result<RangeInclusive<i32>, ()> {
487	let start_julian_day = date_range_inclusive.start.to_julian_day();
488	let last_julian_day = date_range_inclusive.last.to_julian_day();
489
490	#[cfg(feature = "log-trace-frontend")]
491	trace!("Using julian days start: `{start_julian_day}` last: `{last_julian_day}`.");
492
493	if start_julian_day > last_julian_day {
494		error!(
495			"Invalid start date \"{}\" (inclusive date, inclusive time at midnight) and/or last date \"{}\" (inclusive date, exclusive time at midnight).",
496			date_range_inclusive.start,
497			date_range_inclusive.last,
498		);
499
500		return Err(());
501	}
502
503	let julian_day_range = RangeInclusive {
504		start: start_julian_day,
505		last: last_julian_day,
506	};
507
508	Ok(julian_day_range)
509}
510
511pub(crate) fn date_from_julian(
512	julian_day: i32,
513) -> Result<Date, ()> {
514	let Ok(
515		date,
516	) = Date::from_julian_day(julian_day) else {
517		error!("Invalid start date or end date.");
518		return Err(());
519	};
520
521	Ok(date)
522}
523
524pub(crate) fn iso_string_date_from_julian(
525	julian_day: i32,
526) -> Result<String, ()> {
527	let date = date_from_julian(julian_day)?;
528	date.format(&Iso8601::DATE).map_err(|_| ())
529}
530
531pub(crate) fn iso_string_date(
532	date: Date,
533) -> Result<String, ()> {
534	date.format(&Iso8601::DATE).map_err(|_| ())
535}
536
537pub(crate) fn parse_inclusive_date_range(
538	start: &str,
539	last: &str,
540) -> Result<RangeInclusive<Date>, ()> {
541	let Ok(
542		start,
543	) = Date::parse(&start, &Iso8601::DATE) else {
544		error!("Invalid start date (inclusive): {start:?}.");
545		return Err(());
546	};
547
548	let Ok(
549		last,
550	) = Date::parse(&last, &Iso8601::DATE) else {
551		error!("Invalid last date (inclusive): {last:?}.");
552		return Err(());
553	};
554
555	let date_range = RangeInclusive { start, last };
556
557	Ok(date_range)
558}