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;
#[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,
} {
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),
);
}
#[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),
);
}
while let Some(remote_request) = request_delay_queue
.next()
.await
{
remote_request_sender
.send(remote_request.into_inner())
.await
.unwrap();
}
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(),
});
response_sender
.send(remote_response)
.await
.unwrap();
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() {
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() {
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();
}
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() {
send_trades_flag.store(
false,
atomic::Ordering::Relaxed,
);
}
}
}
}
}