mod cli;
use std::{path::PathBuf, range::RangeInclusive};
use apple_quant_core::{
log::{error, info, initialize_tracing, instrument, trace_span, CtxExt},
AnyhowResultExt,
};
use clap::Parser;
use time::{format_description::well_known::Iso8601, Date, Duration, Time, UtcDateTime};
use tokio::{fs::File, io::BufReader};
#[allow(unused_imports)]
use tracing::trace;
use crate::{
backend::{
DataBackend, HistoricalDataBackend, MarketDataDecoderProvider, OrderIdGenerator,
OrdersBackend, OrdersBackendUpdate, OrdersBackendUpdateRecycle,
},
instrument::{InstrumentData, InstrumentKind, InstrumentSpec, InstrumentTicker},
order::{ApplyOrderStateUpdate, DeferredOrderActions, OrderStateUpdate, StateGoal},
order_manager::{OrderManager, OrdersCapacitySpec},
timestamp::{TickTimestamp, TimestampRangeExclusive, Timestamped, UtcNs},
strategy::Strategy, volume::DirectionalExposure,
};
use cli::*;
pub async fn run<
'instrument_data: 'aggregated_data,
'aggregated_data,
IS: InstrumentSpec,
OB: OrdersBackend<IS>,
OrdersCS: OrdersCapacitySpec,
S: Strategy<IS, OB, OrdersCS>,
DB: DataBackend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
>(
strategy: S,
) {
initialize_tracing();
let CLI {
key,
command,
} = CLI::parse();
#[cfg(feature = "log-trace-frontend")]
trace!("CLI parsed with key: `{key:?}` command: `{command:?}`.");
match command {
Commands::Backtest {
start,
last,
} => {
let date_range = parse_inclusive_date_range(&start, &last).unwrap();
#[cfg(feature = "log-trace-frontend")]
trace!("Parsed date range: `{date_range:?}`.");
backtest::<IS, OB, OrdersCS, S, DB>(
strategy,
key.as_ref().map(|string| string.as_str()),
date_range,
).await;
},
Commands::Fetch {
start,
last,
} => {
let date_range = parse_inclusive_date_range(&start, &last).unwrap();
#[cfg(feature = "log-trace-frontend")]
trace!("Parsed date range: `{date_range:?}`.");
fetch::<IS, OB, OrdersCS, S, DB>(
strategy,
key.as_ref().map(|string| string.as_str()),
date_range,
).await;
},
}
}
#[instrument(skip_all)]
async fn backtest<
'instrument_data: 'aggregated_data,
'aggregated_data,
IS: InstrumentSpec,
OB: OrdersBackend<IS>,
OrdersCS: OrdersCapacitySpec,
S: Strategy<IS, OB, OrdersCS>,
DB: DataBackend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
>(
strategy: S,
key: Option<&str>,
date_range_inclusive: RangeInclusive<Date>,
) {
let start = unsafe {
UtcNs::from_utc_date_time(
&UtcDateTime::new(date_range_inclusive.start, Time::MIDNIGHT),
).into_timestamp()
};
let end = unsafe {
UtcNs::from_utc_date_time(
&UtcDateTime::new(
date_range_inclusive.last + Duration::days(1),
Time::MIDNIGHT,
),
).into_timestamp()
};
let walk_range = TimestampRangeExclusive { start, end };
#[cfg(feature = "log-trace-frontend")]
trace!("Backtesting with walk range: `{walk_range:?}`.");
let RangeInclusive {
start: start_julian_day,
last: last_julian_day,
} = julian_days_from_date(&date_range_inclusive).unwrap();
let day_count = (last_julian_day - start_julian_day) as usize + 1;
let mut buf_readers = Vec::with_capacity(day_count);
for julian_day in start_julian_day..=last_julian_day {
let mut iso_string_date = iso_string_date_from_julian(julian_day).unwrap();
iso_string_date.push_str(DB::file_name_postpend());
let mut path_buf = PathBuf::new();
path_buf.push("market-data");
path_buf.push(IS::INSTRUMENT_KIND.as_str());
path_buf.push(IS::BARE_SYMBOL);
path_buf.push("trades");
path_buf.push(iso_string_date);
#[cfg(feature = "log-trace-frontend")]
trace!("Opening market data file \"{}\".", path_buf.display());
let Ok(
file,
) = File::open(&path_buf).await else {
error!("Unable to open market data file: \"{:?}\".", path_buf);
return;
};
#[cfg(feature = "log-trace-frontend")]
trace!("File handle loaded successfully");
buf_readers.push(BufReader::new(file));
}
let market_data_decoder = DB::decode_market_data(
buf_readers.into_iter(),
).await;
let (
frontend,
orders_backend_update_receiver,
) = Frontend::<'_, '_, IS, OB, OrdersCS, S>::new::<DB>(
None,
strategy,
Some(market_data_decoder),
).await;
let mut data_backend = DB::new(key, frontend, orders_backend_update_receiver);
data_backend.backtest(walk_range).await;
}
#[instrument(skip_all)]
async fn fetch<
'instrument_data: 'aggregated_data,
'aggregated_data,
IS: InstrumentSpec,
OB: OrdersBackend<IS>,
OrdersCS: OrdersCapacitySpec,
S: Strategy<IS, OB, OrdersCS>,
DB: DataBackend<'instrument_data, 'aggregated_data, IS, OB, OrdersCS, S>,
>(
strategy: S,
key: Option<&str>,
date_range_inclusive: RangeInclusive<Date>,
) {
let (
frontend,
orders_backend_update_receiver,
) = Frontend::<'_, '_, IS, OB, OrdersCS, S>::new::<DB>(
None,
strategy,
None,
).await;
let mut data_backend = DB::new(key, frontend, orders_backend_update_receiver);
let julian_days = julian_days_from_date(&date_range_inclusive).unwrap();
for julian_day in julian_days.start..=julian_days.last {
let Ok(
date,
) = date_from_julian(julian_day) else {
return;
};
let instrument_ticker = InstrumentTicker::new(
IS::BARE_SYMBOL,
InstrumentKind::Future,
);
let historical_data_backend = data_backend.historical_mut().unwrap().unwrap();
historical_data_backend
.fetch_save_day(&instrument_ticker, date).await
.into_iter().for_each(drop);
}
}
pub struct Frontend<
'instrument_data,
'aggregated_data,
IS: InstrumentSpec,
OB: OrdersBackend<IS>,
OrdersCS: OrdersCapacitySpec,
S: Strategy<IS, OB, OrdersCS>,
> {
instrument_data: InstrumentData<'instrument_data, 'aggregated_data, IS>,
orders_backend: OB,
order_manager: OrderManager<IS, OrdersCS>,
order_id_generator: OrderIdGenerator,
deferred_order_actions: DeferredOrderActions<IS>,
directional_exposure: DirectionalExposure<IS>,
state_goal: StateGoal,
strategy: S,
}
impl<
'instrument_data,
'aggregated_data,
IS: InstrumentSpec,
OB: OrdersBackend<IS>,
OrdersCS: OrdersCapacitySpec,
S: Strategy<IS, OB, OrdersCS>,
> Frontend<
'instrument_data,
'aggregated_data,
IS,
OB,
OrdersCS,
S,
> {
async fn new<MDDP: MarketDataDecoderProvider<IS>>(
orders_key: Option<&str>,
strategy: S,
market_data_decoder: Option<MDDP::MarketDataDecoder>,
) -> (
Self,
thingbuf::mpsc::Receiver<OrdersBackendUpdate<IS>, OrdersBackendUpdateRecycle>,
) {
let (
orders_backend_updates,
orders_backend,
) = OB::new::<OrdersCS, MDDP>(market_data_decoder).await;
let frontend = Self {
instrument_data: InstrumentData::default(),
orders_backend,
order_manager: OrderManager::default(),
order_id_generator: OrderIdGenerator::new(),
deferred_order_actions: DeferredOrderActions::default(),
directional_exposure: DirectionalExposure::default(),
state_goal: StateGoal::default(),
strategy,
};
(frontend, orders_backend_updates)
}
pub(crate) async fn initialize(
&mut self,
tick_timestamp: &TickTimestamp,
) {
#[cfg(feature = "log-trace-frontend")]
trace!("Initializing walk.");
self.strategy
.initialize(
tick_timestamp,
&self.instrument_data,
&mut self.deferred_order_actions,
&mut self.order_id_generator,
).await
.ctx("While initializing strategy.")
.now_checked().unwrap();
}
#[inline(always)]
pub(crate) fn tick(
&mut self,
orders_backend_updates: &mut thingbuf::mpsc::Receiver<
OrdersBackendUpdate<IS>,
OrdersBackendUpdateRecycle,
>,
tick_timestamp: &mut TickTimestamp,
) -> impl Future<Output = Result<(), ()>> {
tick_frontend(self, orders_backend_updates, tick_timestamp)
}
#[instrument(skip_all)]
#[inline(always)]
async fn apply_order_state_update(
&mut self,
order_state_update: &OrderStateUpdate<IS>,
tick_timestamp: &TickTimestamp,
) {
order_state_update
.apply(
tick_timestamp,
&mut self.order_manager,
&mut self.orders_backend,
&mut self.order_id_generator,
&mut self.directional_exposure,
&mut self.deferred_order_actions,
&mut self.state_goal,
).await
.unwrap_or_else(|e| {
error!("{e:?}");
panic!();
});
}
#[instrument(skip_all)]
#[inline(always)]
async fn drain_deferred_order_actions(
&mut self,
tick_timestamp: &TickTimestamp,
) -> bool {
self.deferred_order_actions
.trigger(
tick_timestamp,
&mut self.order_manager,
&mut self.orders_backend,
&mut self.order_id_generator,
&mut self.directional_exposure,
&mut self.state_goal,
).await
.unwrap_or_else(|e| {
error!("{e:?}");
panic!()
})
}
#[instrument(skip_all)]
#[inline(always)]
async fn tick_strategy(
&mut self,
tick_timestamp: &TickTimestamp,
) {
unsafe {
self.strategy
.data_update(
tick_timestamp,
&self.instrument_data,
&self.order_manager,
&self.directional_exposure,
&self.state_goal,
&mut self.deferred_order_actions,
&mut self.order_id_generator,
).await
.ctx("While strategy data update.").now_unchecked_release();
}
}
}
#[instrument(skip_all)]
#[inline(always)]
async fn tick_frontend<
IS: InstrumentSpec,
OB: OrdersBackend<IS>,
OrdersCS: OrdersCapacitySpec,
S: Strategy<IS, OB, OrdersCS>,
>(
frontend: &mut Frontend<'_, '_, IS, OB, OrdersCS, S>,
orders_backend_updates: &mut thingbuf::mpsc::Receiver<
OrdersBackendUpdate<IS>,
OrdersBackendUpdateRecycle,
>,
tick_timestamp: &mut TickTimestamp,
) -> Result<(), ()> {
let span = trace_span!("polling orders backend updates");
let span = span.enter();
#[cfg(feature = "log-trace-frontend")]
trace!("polling orders backend updates");
let recv_ref = match orders_backend_updates.try_recv_ref() {
Ok(
recv_ref,
) => recv_ref,
Err(thingbuf::mpsc::errors::TryRecvError::Empty) => return Ok(()),
Err(thingbuf::mpsc::errors::TryRecvError::Closed) => {
info!("Orders backend updates channel closed. Terminating program.");
panic!();
},
_ => return Err(()),
};
drop(span);
let span = trace_span!("processing new trades");
let span = span.enter();
#[cfg(feature = "log-trace-frontend")]
trace!("processing new trades");
let OrdersBackendUpdate {
trades,
order_state_updates,
} = &*recv_ref;
frontend.instrument_data.walk_new_trades(trades.iter().cloned());
if let Some(
timestamp,
) = trades.last().map(|
trade_trade_timestamp,
| trade_trade_timestamp.timestamp()) {
*tick_timestamp = TickTimestamp::from_timestamp(timestamp);
}
drop(span);
let span = trace_span!("order state updates");
let span = span.enter();
#[cfg(feature = "log-trace-frontend")]
trace!("order state updates");
for order_state_update in order_state_updates {
frontend.apply_order_state_update(
order_state_update,
tick_timestamp,
).await;
}
drop(span);
frontend.drain_deferred_order_actions(tick_timestamp).await;
let span = trace_span!("processing client orders");
let span = span.enter();
#[cfg(feature = "log-trace-frontend")]
trace!("processing client orders");
frontend.order_manager.pending_client_orders.tick_client_orders(
&mut frontend.deferred_order_actions,
&mut frontend.directional_exposure,
frontend.instrument_data.iter_recent_latest_trades_forward(),
);
frontend.drain_deferred_order_actions(tick_timestamp).await;
drop(span);
let span = trace_span!("processing strategy");
let span = span.enter();
#[cfg(feature = "log-trace-frontend")]
trace!("processing strategy");
frontend.tick_strategy(&tick_timestamp).await;
frontend.drain_deferred_order_actions(tick_timestamp).await;
drop(span);
Ok(())
}
pub(crate) fn julian_days_from_date(
date_range_inclusive: &RangeInclusive<Date>,
) -> Result<RangeInclusive<i32>, ()> {
let start_julian_day = date_range_inclusive.start.to_julian_day();
let last_julian_day = date_range_inclusive.last.to_julian_day();
#[cfg(feature = "log-trace-frontend")]
trace!("Using julian days start: `{start_julian_day}` last: `{last_julian_day}`.");
if start_julian_day > last_julian_day {
error!(
"Invalid start date \"{}\" (inclusive date, inclusive time at midnight) and/or last date \"{}\" (inclusive date, exclusive time at midnight).",
date_range_inclusive.start,
date_range_inclusive.last,
);
return Err(());
}
let julian_day_range = RangeInclusive {
start: start_julian_day,
last: last_julian_day,
};
Ok(julian_day_range)
}
pub(crate) fn date_from_julian(
julian_day: i32,
) -> Result<Date, ()> {
let Ok(
date,
) = Date::from_julian_day(julian_day) else {
error!("Invalid start date or end date.");
return Err(());
};
Ok(date)
}
pub(crate) fn iso_string_date_from_julian(
julian_day: i32,
) -> Result<String, ()> {
let date = date_from_julian(julian_day)?;
date.format(&Iso8601::DATE).map_err(|_| ())
}
pub(crate) fn iso_string_date(
date: Date,
) -> Result<String, ()> {
date.format(&Iso8601::DATE).map_err(|_| ())
}
pub(crate) fn parse_inclusive_date_range(
start: &str,
last: &str,
) -> Result<RangeInclusive<Date>, ()> {
let Ok(
start,
) = Date::parse(&start, &Iso8601::DATE) else {
error!("Invalid start date (inclusive): {start:?}.");
return Err(());
};
let Ok(
last,
) = Date::parse(&last, &Iso8601::DATE) else {
error!("Invalid last date (inclusive): {last:?}.");
return Err(());
};
let date_range = RangeInclusive { start, last };
Ok(date_range)
}