apple-quant-algorithmic 0.1.0

Apple Quant's algorithmic trading api
Documentation
use std::{
	sync::{
		Arc,
		atomic::{self, AtomicBool},
	},
	time::{Duration, Instant},
};

use bevy::log::error;
use smallvec::SmallVec;
use tokio::{sync::mpsc, time::sleep_until};
use tokio_stream::StreamExt;
use tokio_util::time::DelayQueue;

use crate::{
	aggregation::TradeTradeTimestamp,
	backend::RemoteOrderId,
	instrument::InstrumentSpec,
	liquidity::LiquidityEstimation,
	order::{
		DesiredVolumeOrder, MatchableOrder, NewCanceledState, NewPartialFillState, NewWorkingState,
		OrderStateUpdate, PartialOrderFill, RemoteOrder,
	},
	timestamp::Timestamp,
	volume::{DirectionlessVolumeSubError, ZeroableDirectionalIntentVolume},
};

use super::{OrderInWorking, RemoteRequest, specs};

pub(super) async fn remote_link<IS: InstrumentSpec + 'static, CapacitySpec: specs::CapacitySpec>(
	mut link_request_receiver: mpsc::Receiver<(Instant, RemoteRequest<IS>)>,
	order_state_update_sender: mpsc::Sender<OrderStateUpdate<IS>>,
	tick_trades_receiver: mpsc::Receiver<SmallVec<[TradeTradeTimestamp<IS>; 10]>>,
	send_trades_flag: Arc<AtomicBool>,
) {
	let (remote_request_sender, remote_request_receiver) = mpsc::channel(8);

	let (remote_response_sender, mut remote_response_receiver) = mpsc::channel(8);

	let _remote_join_handle = tokio::spawn(async move {
		remote::<IS, CapacitySpec>(
			remote_request_receiver,
			remote_response_sender,
			tick_trades_receiver,
			send_trades_flag,
		)
		.await;
	});

	let mut request_delay_queue = DelayQueue::with_capacity(32);
	let mut response_delay_queue = DelayQueue::with_capacity(32);

	'remote_link: loop {
		let now = Instant::now();
		let tokio_now = tokio::time::Instant::from_std(now);

		sleep_until(tokio_now + Duration::from_micros(100)).await;

		// LINK REQUEST RECEIVE

		// Exit condition continues and breaks 'remote_loop.
		#[allow(irrefutable_let_patterns)]
		'receiver: while let (instant, remote_request) = match link_request_receiver.try_recv() {
			Ok((instant, remote_request)) => (instant, remote_request),
			Err(mpsc::error::TryRecvError::Empty) => {
				break 'receiver;
			}
			Err(mpsc::error::TryRecvError::Disconnected) => break 'remote_link,
		} {
			// Bypass delay for market data.
			if let RemoteRequest::<IS>::MarketDataUpdate = &remote_request {
				remote_request_sender
					.send(RemoteRequest::<IS>::MarketDataUpdate)
					.await
					.unwrap();
			}

			request_delay_queue.insert(
				remote_request,
				Duration::from_millis(1),
			);
		}

		// REMOTE RESPONSE RECEIVE

		#[allow(irrefutable_let_patterns)]
		'receiver: while let remote_response = match remote_response_receiver.try_recv() {
			Ok(remote_response) => remote_response,
			Err(mpsc::error::TryRecvError::Empty) => {
				break 'receiver;
			}
			Err(mpsc::error::TryRecvError::Disconnected) => break 'remote_link,
		} {
			response_delay_queue.insert(
				remote_response,
				Duration::from_millis(1),
			);
		}

		// REMOTE REQUEST SEND

		while let Some(remote_request) = request_delay_queue
			.next()
			.await
		{
			remote_request_sender
				.send(remote_request.into_inner())
				.await
				.unwrap();
		}

		// LINK RESPONSE SEND

		while let Some(remote_response) = response_delay_queue
			.next()
			.await
		{
			order_state_update_sender
				.send(remote_response.into_inner())
				.await
				.unwrap();
		}
	}
}

async fn remote<IS: InstrumentSpec, CapacitySpec: specs::CapacitySpec>(
	mut request_receiver: mpsc::Receiver<RemoteRequest<IS>>,
	response_sender: mpsc::Sender<OrderStateUpdate<IS>>,
	mut tick_trades_receiver: mpsc::Receiver<SmallVec<[TradeTradeTimestamp<IS>; 10]>>,
	send_trades_flag: Arc<AtomicBool>,
) {
	let mut orders_in_working: SmallVec<[OrderInWorking<IS>; CapacitySpec::WORKING]> =
		SmallVec::default();

	let mut remote_order_id_counter = 0;

	while let Some(remote_request) = request_receiver.recv().await {
		match remote_request {
			RemoteRequest::SubmitOrder { local_order_id, remote_order } => {
				remote_order_id_counter += 1;
				let remote_order_id = RemoteOrderId::new_from_u64(remote_order_id_counter);

				let order_in_working = OrderInWorking {
					remote_order_id,
					remote_order,
					liquidity_estimation: LiquidityEstimation::default(),
					partial_order_fills: SmallVec::default(),
				};

				orders_in_working.push(order_in_working);

				let remote_response = OrderStateUpdate::NewWorkingState(NewWorkingState {
					local_order_id,
					remote_order_id,
					booking_timestamp: Timestamp::now(),
				});

				// Pending order is always guaranteed because it ties the [`LocalOrderId`] and [`RemoteOrderId`] together.
				response_sender
					.send(remote_response)
					.await
					.unwrap();

				// Relaxed becasue no instruction or memory order requirement.
				send_trades_flag.store(
					true,
					atomic::Ordering::Relaxed,
				);
			}
			RemoteRequest::ModifyVolume { remote_order_id, volume } => {
				let Some((order_in_working_idx, order_in_working)) = orders_in_working
					.iter_mut()
					.enumerate()
					.find(|(_, order_in_working)| {
						order_in_working.remote_order_id == remote_order_id
					})
				else {
					error!(
						"Could not find {:?} in working orders.",
						remote_order_id
					);

					continue;
				};

				let RemoteOrder::Limit(limit_order) = &mut order_in_working.remote_order else {
					continue;
				};

				let mut filled_volume = ZeroableDirectionalIntentVolume::ZERO;

				for partial_order_fill in order_in_working
					.partial_order_fills
					.iter()
				{
					filled_volume += partial_order_fill
						.directional_intent_volume
						.as_zeroable();
				}

				let Some(filled_volume) = filled_volume
					.directional_intent_volume()
					.map(|directional_intent_volume| {
						directional_intent_volume.directionless_volume
					})
				else {
					limit_order
						.resting_volume
						.directionless_volume = volume;

					continue;
				};

				let will_have_resting_volume = match volume - filled_volume {
					Ok(new_resting_volume) => new_resting_volume.is_some(),
					Err(DirectionlessVolumeSubError::NegativeVolume { lhs: _, rhs: _ }) => false,
				};

				if !will_have_resting_volume {
					orders_in_working.swap_remove(order_in_working_idx);

					let remote_response =
						OrderStateUpdate::NewCanceledState(NewCanceledState { remote_order_id });

					response_sender
						.send(remote_response)
						.await
						.unwrap();

					if orders_in_working.is_empty() {
						// Relaxed becasue no instruction or memory order requirement.
						send_trades_flag.store(
							false,
							atomic::Ordering::Relaxed,
						);
					}

					continue;
				};

				limit_order
					.resting_volume
					.directionless_volume = volume;
			}
			RemoteRequest::CancelOrder(remote_order_id) => {
				let Some(orders_in_working_idx) = orders_in_working
					.iter()
					.position(|order_in_working| {
						order_in_working.remote_order_id == remote_order_id
					})
				else {
					error!(
						"Could not find {:?} in working orders.",
						remote_order_id
					);

					continue;
				};

				let _ = orders_in_working.swap_remove(orders_in_working_idx);

				let remote_response =
					OrderStateUpdate::NewCanceledState(NewCanceledState { remote_order_id });

				response_sender
					.send(remote_response)
					.await
					.unwrap();

				if orders_in_working.is_empty() {
					// Relaxed becasue no instruction or memory order requirement.
					send_trades_flag.store(
						false,
						atomic::Ordering::Relaxed,
					);
				}
			}
			RemoteRequest::CancelAll | RemoteRequest::CancelFlattenAll => {
				for remote_order_id in orders_in_working
					.drain(..)
					.map(|order_in_working| order_in_working.remote_order_id)
				{
					let remote_response =
						OrderStateUpdate::NewCanceledState(NewCanceledState { remote_order_id });

					response_sender
						.send(remote_response)
						.await
						.unwrap();
				}

				// Relaxed becasue no instruction or memory order requirement.
				send_trades_flag.store(
					false,
					atomic::Ordering::Relaxed,
				);
			}
			RemoteRequest::MarketDataUpdate => {
				while let Ok(trade_trade_timestamps) = tick_trades_receiver.try_recv() {
					for order_in_working in orders_in_working.iter_mut() {
						order_in_working
							.liquidity_estimation
							.walk_trades(trade_trade_timestamps.iter());
					}
				}

				let mut responses: SmallVec<[OrderStateUpdate<IS>; 4]> = SmallVec::default();

				let _: SmallVec<[OrderInWorking<IS>; 4]> = orders_in_working
					.drain_filter(|order_in_working| {
						let Some(price) = order_in_working
							.remote_order
							.is_liquidable(&order_in_working.liquidity_estimation)
						else {
							return false;
						};

						let directional_intent_volume = match order_in_working
							.remote_order
							.desired_directional_intent_volume()
						{
							Ok(directional_intent_volume) => directional_intent_volume,
							Err(error) => {
								error!("{error:?}");
								return true;
							}
						};

						let partial_order_fill = PartialOrderFill {
							price,
							directional_intent_volume,
						};

						let remote_response =
							OrderStateUpdate::NewPartialFillState(NewPartialFillState {
								remote_order_id: order_in_working.remote_order_id,
								partial_order_fill,
							});

						responses.push(remote_response);

						true
					})
					.collect();

				for remote_response in responses {
					response_sender
						.send(remote_response)
						.await
						.unwrap();
				}

				if orders_in_working.is_empty() {
					// Relaxed becasue no instruction or memory order requirement.
					send_trades_flag.store(
						false,
						atomic::Ordering::Relaxed,
					);
				}
			}
		}
	}
}