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 #[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}