use std::{
sync::{
atomic::{self, AtomicBool, AtomicU64, AtomicU8},
Arc,
},
time::{Duration, Instant},
};
use apple_quant_core::{
log::{error, info, warn},
AddUnchecked, SubChecked,
};
use bitflag::bitflag;
use smallvec::SmallVec;
use thingbuf::recycling::DefaultRecycle;
use tokio::task::yield_now;
use tokio_stream::StreamExt;
use tokio_util::time::DelayQueue;
use crate::{
backend::{
DataBackend, MarketDataDecoder, MarketDataDecoderProvider, OrdersBackend,
OrdersBackendUpdate, OrdersBackendUpdateRecycle, RemoteOrderId,
},
order::{
DesiredVolumeOrder, MatchableOrder, NewCanceledState, NewPartialFillState,
NewWorkingState, OrderStateUpdate, OrderStateUpdateListInline, PartialOrderFill,
RemoteOrder,
},
timestamp::{Timestamp, Timestamped, UtcNs},
volume::{DirectionlessVolume, Zeroable, ZeroableExt, ZeroableVolume},
aggregation::TradeTradeTimestamp, instrument::InstrumentSpec,
liquidity::LiquidityEstimation, order_manager::OrdersCapacitySpec, strategy::Strategy,
};
use super::{specs, MarketDataStream, OrderInWorking, RemoteRequest};
pub(super) async fn remote_link<
IS: InstrumentSpec,
MDD: MarketDataDecoder<IS>,
CapacitySpec: specs::CapacitySpec,
>(
mut link_request_receiver: thingbuf::mpsc::Receiver<RemoteRequest<IS>>,
link_orders_backend_update_sender: thingbuf::mpsc::Sender<
OrdersBackendUpdate<IS>,
OrdersBackendUpdateRecycle,
>,
market_data_stream: MarketDataStream<IS, MDD>,
) {
let (
remote_request_sender,
remote_request_receiver,
) = thingbuf::mpsc::channel(64);
let (
orders_backend_update_sender,
mut orders_backend_update_receiver,
) = thingbuf::mpsc::with_recycle(1, OrdersBackendUpdateRecycle);
let (
timestamp_sender,
mut timestamp_receiver,
) = thingbuf::mpsc::channel(1);
let _remote_join_handle = tokio::spawn(async move {
remote::<IS, MDD, CapacitySpec>(
remote_request_receiver,
orders_backend_update_sender,
market_data_stream,
timestamp_sender,
).await;
});
let mut request_delay_queue = Vec::with_capacity(64);
let mut response_delay_queue = Vec::with_capacity(64);
let mut newest_timestamp: Option<Timestamp> = None;
'remote_link: loop {
yield_now().await;
let recv_ref = match timestamp_receiver.try_recv_ref() {
Ok(
recv_ref,
) => recv_ref,
Err(thingbuf::mpsc::errors::TryRecvError::Empty) => continue,
Err(thingbuf::mpsc::errors::TryRecvError::Closed) => {
warn!("Channel closed; remote link exiting.");
return;
},
_ => {
error!("Exiting; unkown error.");
return;
},
};
newest_timestamp = *recv_ref;
#[allow(irrefutable_let_patterns)]
'receiver: while let remote_request = match link_request_receiver.try_recv() {
Ok(
remote_request,
) => remote_request,
Err(thingbuf::mpsc::errors::TryRecvError::Empty) => break 'receiver,
Err(thingbuf::mpsc::errors::TryRecvError::Closed) => {
warn!("Channel closed; remote link exiting.");
return;
},
_ => {
error!("Exiting; unkown error.");
return;
},
} {
request_delay_queue.push((newest_timestamp, remote_request));
}
#[allow(irrefutable_let_patterns)]
'receiver: while let orders_backend_update = match orders_backend_update_receiver.try_recv() {
Ok(
orders_backend_update,
) => orders_backend_update,
Err(thingbuf::mpsc::errors::TryRecvError::Empty) => break 'receiver,
Err(thingbuf::mpsc::errors::TryRecvError::Closed) => {
warn!("Channel closed; remote link exiting.");
return;
},
_ => {
error!("Exiting; unkown error.");
return;
},
} {
response_delay_queue.push((newest_timestamp, orders_backend_update));
}
let process_remote_requests = request_delay_queue
.extract_if(.., |(
timestamp,
remote_request,
)| {
if let Some(
timestamp,
) = timestamp && newest_timestamp.unwrap() < *timestamp {
return false;
}
true
})
.map(|(
_,
remote_request,
)| remote_request);
for remote_request in process_remote_requests {
if let Err(
_,
) = remote_request_sender.send(
remote_request,
).await {
warn!("Channel closed; remote link exiting.");
return;
}
}
let process_response_responses = response_delay_queue
.extract_if(.., |(
timestamp,
remote_response,
)| {
if let Some(
timestamp,
) = timestamp && newest_timestamp.unwrap() < *timestamp {
return false;
}
true
})
.map(|(
_,
remote_response,
)| remote_response);
for orders_backend_update in process_response_responses {
if let Err(
_,
) = link_orders_backend_update_sender.send(
orders_backend_update,
).await {
warn!("Channel closed; remote link exiting.");
return;
}
}
}
}
async fn remote<
IS: InstrumentSpec,
MDD: MarketDataDecoder<IS>,
CapacitySpec: specs::CapacitySpec,
>(
mut request_receiver: thingbuf::mpsc::Receiver<RemoteRequest<IS>>,
orders_backend_update_sender: thingbuf::mpsc::Sender<
OrdersBackendUpdate<IS>,
OrdersBackendUpdateRecycle,
>,
mut market_data_stream: MarketDataStream<IS, MDD>,
timestamp_sender: thingbuf::mpsc::Sender<Option<Timestamp>>,
) {
let mut orders_in_working: SmallVec<[OrderInWorking<IS>; core::direct_const_arg!(
CapacitySpec::WORKING
)]> = SmallVec::default();
let mut remote_order_id_counter = 0;
loop {
yield_now().await;
let mut remote_requests: SmallVec<[RemoteRequest<IS>; 8]> = SmallVec::default();
loop {
let remote_request = match request_receiver.try_recv() {
Ok(
remote_request,
) => remote_request,
Err(thingbuf::mpsc::errors::TryRecvError::Empty) => break,
Err(thingbuf::mpsc::errors::TryRecvError::Closed) => {
warn!("Channel closed; remote exiting.");
return;
},
_ => {
error!("Exiting; unkown error.");
return;
},
};
remote_requests.push(remote_request);
}
let mut order_state_updates = OrderStateUpdateListInline::default();
for remote_request in remote_requests {
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 order_state_update = OrderStateUpdate::NewWorkingState(NewWorkingState {
local_order_id,
remote_order_id,
booking_timestamp: Timestamp::now(),
});
order_state_updates.push(order_state_update);
},
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 {
warn!("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: Zeroable<DirectionlessVolume<IS>> = Zeroable::ZERO;
for partial_order_fill in order_in_working.partial_order_fills.iter() {
filled_volume = filled_volume.add_unchecked(
partial_order_fill.directional_intent_volume.directionless_volume.as_zeroable(),
);
}
let Some(
filled_volume,
) = filled_volume.as_optional_nonzero_ref() else {
limit_order.resting_volume.directionless_volume = volume;
continue;
};
let will_have_resting_volume = volume.sub_checked(filled_volume).is_ok();
if !will_have_resting_volume {
orders_in_working.swap_remove(order_in_working_idx);
let order_state_update = OrderStateUpdate::NewCanceledState(NewCanceledState {
remote_order_id,
});
order_state_updates.push(order_state_update);
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 {
warn!("Could not find {:?} in working orders.", remote_order_id);
continue;
};
let _ = orders_in_working.swap_remove(orders_in_working_idx);
let order_state_update = OrderStateUpdate::NewCanceledState(NewCanceledState {
remote_order_id,
});
order_state_updates.push(order_state_update);
},
RemoteRequest::CancelAll | RemoteRequest::CancelFlattenAll => {
for remote_order_id in orders_in_working
.drain(..)
.map(|
order_in_working,
| order_in_working.remote_order_id)
{
let order_state_update = OrderStateUpdate::NewCanceledState(NewCanceledState {
remote_order_id,
});
order_state_updates.push(order_state_update);
}
},
_ => { warn!("Remote request nop triggered."); },
}
}
let Ok(
decoded_tick,
) = market_data_stream.decode_tick().await else {
warn!("Channel closed; remote exiting.");
return;
};
let Some((
trades,
newest_trade_timestamp,
)) = decoded_tick else {
continue;
};
for order_in_working in orders_in_working.iter_mut() {
order_in_working.liquidity_estimation.walk_trades(trades.iter());
}
let mut orders_backend_update = OrdersBackendUpdate {
order_state_updates,
trades,
};
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,
timestamp: newest_trade_timestamp,
};
let order_state_update = OrderStateUpdate::NewPartialFillState(NewPartialFillState {
remote_order_id: order_in_working.remote_order_id,
partial_order_fill,
});
orders_backend_update.order_state_updates.push(order_state_update);
true
}).for_each(drop);
if let Err(
_,
) = orders_backend_update_sender.send(
orders_backend_update,
).await {
warn!("Channel closed; remote exiting.");
return;
}
if let Err(
_,
) = timestamp_sender.send(
Some(newest_trade_timestamp.timestamp()),
).await {
warn!("Channel closed; remote exiting.");
return;
}
}
}