mod core_error;
mod kline_store;
mod storage;
use chrono::Utc;
pub use core_error::{CoreError, CoreResult};
pub use kline_store::KlineStore;
use std::cmp::min;
#[cfg(test)]
use crate::data::SpecificOrderDetails;
use crate::{
clock::ClockBase,
core::storage::Storage,
data::{Crypto, Order, OrderDraft, OrderStatus},
market::{Balance, Kline, KlinesParams, Market},
};
const MAX_KLINES_REQUESTS: u32 = 5;
pub struct Core<M: Market> {
market: M,
storage: Storage,
kline_store: KlineStore<M>,
clock: <M as Market>::_Clock,
}
impl<M: Market> Core<M> {
pub async fn new(
name: String,
market: M,
kline_store: &KlineStore<M>,
clock: <M as Market>::_Clock,
) -> CoreResult<Self> {
let storage = Storage::new(&name).await?;
let kline_store = kline_store.clone();
Ok(Self {
market,
storage,
kline_store,
clock,
})
}
pub async fn send_order(&self, draft: OrderDraft) -> CoreResult<Order<M::_OrderDetails>> {
let mut order: Order<M::_OrderDetails> = Order::new(
draft.asset,
draft.quote,
draft.side,
draft.order_type,
draft.amount,
draft.price,
)?;
tracing::info!(
"{} sending order {}, {:#?} {} at {}",
self.storage.get_name(),
order.id(),
order.side(),
order.asset(),
self.clock.now()
);
match self.market.send_order(&mut order).await {
Ok(()) => {
tracing::debug!("Order {} sent successfully", order.id());
self.storage.add_order(order.clone()).await?;
Ok(order)
}
Err(e) => {
order.set_status(OrderStatus::Failed);
tracing::warn!("Failed to send order {} : {}", order.id(), e);
self.storage.add_order(order).await?;
Err(e)
}
}
}
pub fn clock(&self) -> &M::_Clock {
&self.clock
}
pub async fn update_opened_orders(&self) -> CoreResult<Vec<Order<M::_OrderDetails>>> {
tracing::debug!("Updating opened orders");
let mut orders: Vec<Order<M::_OrderDetails>> = self.storage.opened_orders().await;
for order in &mut orders {
tracing::trace!("Updating order {}", order.id());
self.update_order(order).await?;
}
Ok(orders)
}
pub async fn cancel_opened_orders(&self) -> CoreResult<Vec<Order<M::_OrderDetails>>> {
let mut orders: Vec<Order<M::_OrderDetails>> = self.storage.opened_orders().await;
for order in &mut orders {
self.cancel_order(order).await?;
}
Ok(orders)
}
pub async fn get_order(&self, id: &str) -> Option<Order<M::_OrderDetails>> {
self.storage.get_order(id).await
}
pub async fn common_get_opened_orders(&self) -> Vec<Order<M::_OrderDetails>> {
self.storage.opened_orders().await
}
pub async fn update_order(&self, order: &mut Order<M::_OrderDetails>) -> CoreResult<bool> {
tracing::debug!("Updating order {}", order.id());
let query_output = self.market.update_order(order).await?;
if *order.status() != query_output.status.into()
|| *order.price() != Some(query_output.price)
|| *order.executed_qty() != query_output.executed_quantity
|| *order.cummulative_quote_qty() != query_output.cummulative_quote_quantity
{
order.set_status(query_output.status.into());
if *order.price() != Some(query_output.price) {
order.set_price(query_output.price)?;
}
order.set_executed(
query_output.executed_quantity,
query_output.cummulative_quote_quantity,
);
self.storage.update_order(order.clone()).await?;
Ok(true)
} else {
Ok(false)
}
}
pub async fn cancel_order(&self, order: &mut Order<M::_OrderDetails>) -> CoreResult<bool> {
tracing::debug!("Canceling order {}", order.id());
let cancel_output = self.market.cancel_order(order).await?;
if *order.status() == cancel_output.status.into() {
return Ok(false);
}
order.set_status(cancel_output.status.into());
if *order.price() != Some(cancel_output.price) {
order.set_price(cancel_output.price)?;
}
if *order.executed_qty() != cancel_output.executed_quantity
|| *order.cummulative_quote_qty() != cancel_output.cummulative_quote_quantity
{
order.set_executed(
cancel_output.executed_quantity,
cancel_output.cummulative_quote_quantity,
);
}
self.storage.update_order(order.clone()).await?;
Ok(true)
}
pub async fn get_klines(&self, mut params: KlinesParams) -> CoreResult<Vec<Kline>> {
let mut safety = 0;
tracing::debug!("Core : Getting klines");
tracing::trace!("Params: {:#?}", params);
loop {
tracing::trace!("Core klines loop round {}", safety);
let (local_klines, complete) = self.kline_store.try_get_klines(¶ms).await?;
if complete {
tracing::debug!("Klines acquired");
return Ok(local_klines);
}
if safety > MAX_KLINES_REQUESTS {
return Err(CoreError::param_error("Max klines requests reached"));
}
if !local_klines.is_empty() {
params = KlinesParams::new(
*params.asset(),
*params.quote(),
*params.interval(),
local_klines.last().unwrap().close_time,
*params.end_time(),
);
}
tracing::debug!("Getting klines from Market");
let max_endtime = *params.start_time()
+ params
.interval()
.time_delta()
.checked_mul(self.market.klines_limit() - 1)
.ok_or(CoreError::comput_error("TimeDelta overflow"))?;
let now = Utc::now();
let klines_params_extended = KlinesParams::new(
*params.asset(),
*params.quote(),
*params.interval(),
*params.start_time(),
min(max_endtime, now),
);
tracing::trace!(
"Params modified end_date: {:#?}",
klines_params_extended.end_time()
);
let klines: Vec<Kline> = self.market.klines(klines_params_extended).await?;
if klines.is_empty() {
return Err(CoreError::UnavailableData);
}
self.kline_store.inject_klines(¶ms, klines).await?;
safety += 1;
}
}
pub async fn get_balance(&self, crypto: Crypto) -> CoreResult<Balance> {
self.market.get_balance(crypto).await
}
#[cfg(test)]
pub async fn clear_storage<D: SpecificOrderDetails>(&self) -> CoreResult<()> {
self.storage.clear::<D>().await
}
}
#[cfg(test)]
mod tests {
use std::future::ready;
use chrono::{DateTime, TimeDelta, TimeZone as _, Utc};
use rust_decimal::Decimal;
use crate::{
clock::{ClockFactory as _, RunningCheatClockFactory},
core::{Core, KlineStore},
data::{
Crypto, EmptySpecificOrderDetails, KlineInterval, OrderDraft, OrderStatus, Quantity,
},
market::{
mexc_enums, mock_context, CancelOrderOutput, Kline, KlinesParams, MockMarket,
QueryOrderOutput,
},
};
pub const ASSET: Crypto = Crypto::BTC;
pub const QUOTE: Crypto = Crypto::USDT;
fn klines_start() -> DateTime<Utc> {
Utc.with_ymd_and_hms(2024, 1, 1, 0, 0, 0).unwrap()
}
fn klines_end() -> DateTime<Utc> {
Utc.with_ymd_and_hms(2024, 1, 1, 0, 1, 0).unwrap()
}
#[tokio::test]
async fn reload_and_update_orders() {
let name = "core-reload_and_update_orders".to_string();
let mut market1 = MockMarket::new();
market1
.expect_send_order()
.returning(|_| Box::pin(ready(Ok(()))));
let mut market2 = MockMarket::new();
market2.expect_update_order().returning(|order| {
Box::pin(ready(Ok(QueryOrderOutput {
symbol: "BTCUSDT".to_string(),
original_client_order_id: None,
order_id: order.id().to_string(),
client_order_id: None,
price: order.price().unwrap(),
original_quantity: order.quantity().get_amount(),
executed_quantity: *order.executed_qty(),
cummulative_quote_quantity: *order.cummulative_quote_qty(),
status: mexc_enums::OrderStatus::New,
time_in_force: None,
order_type: *order.order_type(),
side: *order.side(),
stop_price: order.price().unwrap(),
time: Utc::now(),
update_time: Utc::now(),
is_working: true,
})))
});
let _ctx = mock_context();
let kline_store = KlineStore::test_new(&name).await;
let clock = RunningCheatClockFactory { date: Utc::now() }
.build()
.unwrap()
.1;
let order = {
let core = Core::new(name.clone(), market1, &kline_store, clock.clone())
.await
.unwrap();
core.clear_storage::<EmptySpecificOrderDetails>()
.await
.unwrap();
let order = core
.send_order(OrderDraft {
asset: ASSET,
quote: QUOTE,
side: mexc_enums::OrderSide::Buy,
order_type: mexc_enums::OrderType::Limit,
amount: Quantity::Asset(Decimal::ONE),
price: Some(Decimal::from(10000)),
})
.await
.unwrap();
drop(core);
order
};
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
let core = Core::new(name.clone(), market2, &kline_store, clock.clone())
.await
.unwrap();
let updated_orders = core.update_opened_orders().await.unwrap();
assert_eq!(updated_orders.len(), 1);
assert_eq!(core.get_order(order.id()).await.unwrap(), updated_orders[0]);
}
#[tokio::test]
async fn cancel_orders() {
let name = "core-cancel_orders".to_string();
let mut market = MockMarket::new();
market
.expect_send_order()
.returning(|_| Box::pin(ready(Ok(()))));
market.expect_cancel_order().returning(|order| {
Box::pin(ready(Ok(CancelOrderOutput {
symbol: "BTCUSDT".to_string(),
original_client_order_id: None,
order_id: order.id().to_string(),
client_order_id: None,
price: order.price().unwrap(),
original_quantity: order.quantity().get_amount(),
executed_quantity: *order.executed_qty(),
cummulative_quote_quantity: *order.cummulative_quote_qty(),
status: mexc_enums::OrderStatus::Canceled,
time_in_force: None,
order_type: *order.order_type(),
side: *order.side(),
})))
});
let _ctx = mock_context();
let kline_store = KlineStore::test_new(&name).await;
let clock = RunningCheatClockFactory { date: Utc::now() }
.build()
.unwrap()
.1;
let core = Core::new(name, market, &kline_store, clock).await.unwrap();
core.clear_storage::<EmptySpecificOrderDetails>()
.await
.unwrap();
let mut order = core
.send_order(OrderDraft {
asset: ASSET,
quote: QUOTE,
side: mexc_enums::OrderSide::Buy,
order_type: mexc_enums::OrderType::Limit,
amount: Quantity::Asset(Decimal::ONE),
price: Some(Decimal::from(10000)),
})
.await
.unwrap();
let canceled_orders = core.cancel_opened_orders().await.unwrap();
assert_eq!(canceled_orders.len(), 1);
order.set_status(OrderStatus::Canceled);
assert_eq!(order, canceled_orders[0]);
}
#[tokio::test]
async fn get_klines() {
let name = "core-get_klines".to_string();
let _ctx = mock_context();
let kline_store = KlineStore::test_new(&name).await;
kline_store.clear(ASSET, QUOTE).await.unwrap();
let mut market = MockMarket::new();
market.expect_klines_limit().return_const(100);
market.expect_klines().once().returning(move |params| {
Box::pin(ready(Kline::vec_over(
*params.interval(),
*params.start_time(),
*params.end_time(),
Kline {
open_time: klines_start(),
open: Decimal::ZERO,
high: Decimal::from(42),
low: Decimal::ZERO,
close: Decimal::ONE,
volume: Decimal::ONE,
close_time: klines_end(),
quote_asset_volume: Decimal::ONE,
},
)))
});
let clock = RunningCheatClockFactory { date: Utc::now() }
.build()
.unwrap()
.1;
let core = Core::new(name, market, &kline_store, clock).await.unwrap();
core.clear_storage::<EmptySpecificOrderDetails>()
.await
.unwrap();
let start = klines_start();
let end = klines_end();
let klines = core
.get_klines(KlinesParams::new(
ASSET,
QUOTE,
KlineInterval::OneMinute,
start + TimeDelta::minutes(3),
end + TimeDelta::minutes(3),
))
.await
.unwrap();
let kline = Kline {
open_time: start+ TimeDelta::minutes(3),
open: Decimal::ZERO,
high: Decimal::from(42),
low: Decimal::ZERO,
close: Decimal::ONE,
volume: Decimal::ONE,
close_time: end+ TimeDelta::minutes(3),
quote_asset_volume: Decimal::ONE,
};
assert_eq!(klines.len(), 1);
assert_eq!(klines[0].high, kline.high);
}
}