apple-quant-algorithmic 0.3.0

Apple Quant's algorithmic library.
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();
	}
}