extern crate kucoin_api;
use chrono::prelude::Local;
use kucoin_arbitrage::model::order::OrderSide;
use kucoin_arbitrage::translator::translator::OrderBookChangeTranslator;
use kucoin_api::failure;
use kucoin_api::futures::TryStreamExt;
use kucoin_api::{
client::{Kucoin, KucoinEnv},
model::websocket::{KucoinWebsocketMsg, WSTopic, WSType},
};
use uuid::Uuid;
#[tokio::main]
async fn main() -> Result<(), failure::Error> {
kucoin_arbitrage::logger::log_init();
log::info!("Testing Kucoin REST-to-WS latency");
let credentials = kucoin_arbitrage::globals::config::credentials();
log::info!("{credentials:#?}");
let api = Kucoin::new(KucoinEnv::Live, Some(credentials))?;
let url = api.get_socket_endpoint(WSType::Public).await?;
let mut ws = api.websocket();
let subs = vec![WSTopic::OrderBook(vec!["BTC-USDT".to_string()])];
let id: Uuid = Uuid::new_v4();
let test_symbol: &str = "BTC-USDT";
let test_price: f64 = 1.0; let test_volume: f64 = 0.1;
let dt_order_placed = Local::now();
api.cancel_all_orders(None, None).await.unwrap();
api.post_limit_order(
id.to_string().as_str(),
test_symbol,
OrderSide::Buy.as_ref(),
test_price.to_string().as_str(),
test_volume.to_string().as_str(),
None,
)
.await?;
log::info!("Order placed {dt_order_placed}");
ws.subscribe(url, subs).await?;
log::info!("Async polling");
let serial = 0;
while let Some(msg) = ws.try_next().await? {
match msg {
KucoinWebsocketMsg::OrderBookMsg(msg) => {
let (symbol, data) = msg.data.to_internal(serial);
if symbol.ne(test_symbol) {
continue;
}
if let Some(_) = data.bid.get(&ordered_float::OrderedFloat(test_price)) {
log::info!("data: {:#?}", data);
let dt_order_reported = Local::now();
let delta = dt_order_reported - dt_order_placed;
log::info!("REST-to-WS: {}ms", delta.num_milliseconds());
return Ok(());
}
}
KucoinWebsocketMsg::PongMsg(_) => continue,
KucoinWebsocketMsg::WelcomeMsg(_) => continue,
_ => {
panic!("unexpected msgs received: {msg:?}")
}
}
}
Ok(())
}