use std::sync::Arc;
use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderValue, USER_AGENT};
use reqwest_middleware::{ClientBuilder, ClientWithMiddleware};
use reqwest_retry::{RetryTransientMiddleware, policies::ExponentialBackoff};
use reqwest_tracing::TracingMiddleware;
use crate::auth::{CredentialsProvider, IncreasingNonce, NonceProvider};
use crate::error::KrakenError;
use crate::futures::auth::sign_futures_request;
use crate::futures::rest::endpoints::{FUTURES_BASE_URL, private, public};
use crate::futures::rest::types::*;
use crate::futures::types::*;
#[derive(Clone)]
pub struct FuturesRestClient {
http_client: ClientWithMiddleware,
base_url: String,
credentials: Option<Arc<dyn CredentialsProvider>>,
nonce_provider: Arc<dyn NonceProvider>,
}
impl FuturesRestClient {
pub fn new() -> Self {
Self::builder().build()
}
pub fn builder() -> FuturesRestClientBuilder {
FuturesRestClientBuilder::new()
}
pub(crate) async fn public_get<T>(&self, endpoint: &str) -> Result<T, KrakenError>
where
T: serde::de::DeserializeOwned,
{
let url = format!("{}{}", self.base_url, endpoint);
let response = self.http_client.get(&url).send().await?;
self.parse_futures_response(response).await
}
pub(crate) async fn public_get_with_params<T, Q>(
&self,
endpoint: &str,
params: &Q,
) -> Result<T, KrakenError>
where
T: serde::de::DeserializeOwned,
Q: serde::Serialize + ?Sized,
{
let query_string = serde_urlencoded::to_string(params)
.map_err(|e| KrakenError::InvalidResponse(e.to_string()))?;
let url = if query_string.is_empty() {
format!("{}{}", self.base_url, endpoint)
} else {
format!("{}{}?{}", self.base_url, endpoint, query_string)
};
let response = self.http_client.get(&url).send().await?;
self.parse_futures_response(response).await
}
pub(crate) async fn private_get<T>(&self, endpoint: &str) -> Result<T, KrakenError>
where
T: serde::de::DeserializeOwned,
{
let credentials = self
.credentials
.as_ref()
.ok_or(KrakenError::MissingCredentials)?;
let nonce = self.nonce_provider.next_nonce();
let creds = credentials.get_credentials();
let signature = sign_futures_request(creds, endpoint, nonce, "")?;
let url = format!("{}{}", self.base_url, endpoint);
let response = self
.http_client
.get(&url)
.header("APIKey", &creds.api_key)
.header("Authent", signature)
.header("Nonce", nonce.to_string())
.send()
.await?;
self.parse_futures_response(response).await
}
pub(crate) async fn private_post<T, P>(
&self,
endpoint: &str,
params: &P,
) -> Result<T, KrakenError>
where
T: serde::de::DeserializeOwned,
P: serde::Serialize,
{
let credentials = self
.credentials
.as_ref()
.ok_or(KrakenError::MissingCredentials)?;
let nonce = self.nonce_provider.next_nonce();
let creds = credentials.get_credentials();
let form_data = serde_urlencoded::to_string(params)
.map_err(|e| KrakenError::InvalidResponse(e.to_string()))?;
let signature = sign_futures_request(creds, endpoint, nonce, &form_data)?;
let url = format!("{}{}", self.base_url, endpoint);
let response = self
.http_client
.post(&url)
.header("APIKey", &creds.api_key)
.header("Authent", signature)
.header("Nonce", nonce.to_string())
.header(CONTENT_TYPE, "application/x-www-form-urlencoded")
.body(form_data)
.send()
.await?;
self.parse_futures_response(response).await
}
async fn parse_futures_response<T>(&self, response: reqwest::Response) -> Result<T, KrakenError>
where
T: serde::de::DeserializeOwned,
{
let status = response.status();
let body = response.text().await?;
if let Ok(error_response) = serde_json::from_str::<FuturesErrorResponse>(&body) {
if error_response.result == "error" {
return Err(KrakenError::Api(crate::error::ApiError::new(
"EFutures",
error_response
.error
.unwrap_or_else(|| "Unknown error".to_string()),
)));
}
}
serde_json::from_str::<T>(&body).map_err(|e| {
if !status.is_success() {
KrakenError::InvalidResponse(format!("HTTP {}: {}", status, body))
} else {
KrakenError::InvalidResponse(format!(
"Failed to parse response: {}. Body: {}",
e, body
))
}
})
}
pub async fn get_tickers(&self) -> Result<Vec<FuturesTicker>, KrakenError> {
let response: TickersResponse = self.public_get(public::TICKERS).await?;
Ok(response.tickers)
}
pub async fn get_ticker(&self, symbol: &str) -> Result<Option<FuturesTicker>, KrakenError> {
let tickers = self.get_tickers().await?;
Ok(tickers.into_iter().find(|t| t.symbol == symbol))
}
pub async fn get_orderbook(&self, symbol: &str) -> Result<FuturesOrderBook, KrakenError> {
#[derive(serde::Serialize)]
struct Params<'a> {
symbol: &'a str,
}
let response: OrderBookResponse = self
.public_get_with_params(public::ORDERBOOK, &Params { symbol })
.await?;
Ok(response.order_book)
}
pub async fn get_trade_history(
&self,
symbol: &str,
last_time: Option<&str>,
) -> Result<Vec<FuturesTrade>, KrakenError> {
#[derive(serde::Serialize)]
struct Params<'a> {
symbol: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
#[serde(rename = "lastTime")]
last_time: Option<&'a str>,
}
let response: TradeHistoryResponse = self
.public_get_with_params(public::HISTORY, &Params { symbol, last_time })
.await?;
Ok(response.history)
}
pub async fn get_instruments(&self) -> Result<Vec<FuturesInstrument>, KrakenError> {
let response: InstrumentsResponse = self.public_get(public::INSTRUMENTS).await?;
Ok(response.instruments)
}
pub async fn get_accounts(&self) -> Result<AccountsResponse, KrakenError> {
self.private_get(private::ACCOUNTS).await
}
pub async fn get_open_positions(&self) -> Result<Vec<FuturesPosition>, KrakenError> {
let response: OpenPositionsResponse = self.private_get(private::OPEN_POSITIONS).await?;
Ok(response.open_positions)
}
pub async fn get_open_orders(&self) -> Result<Vec<FuturesOrder>, KrakenError> {
let response: OpenOrdersResponse = self.private_get(private::OPEN_ORDERS).await?;
Ok(response.open_orders)
}
pub async fn get_fills(
&self,
request: Option<&FillsRequest>,
) -> Result<Vec<FuturesFill>, KrakenError> {
let response: FillsResponse = match request {
Some(req) => self.public_get_with_params(private::FILLS, req).await?,
None => self.private_get(private::FILLS).await?,
};
Ok(response.fills)
}
pub async fn send_order(
&self,
request: &SendOrderRequest,
) -> Result<SendOrderResponse, KrakenError> {
self.private_post(private::SEND_ORDER, request).await
}
pub async fn edit_order(
&self,
request: &EditOrderRequest,
) -> Result<EditOrderResponse, KrakenError> {
self.private_post(private::EDIT_ORDER, request).await
}
pub async fn cancel_order(&self, order_id: &str) -> Result<CancelOrderResponse, KrakenError> {
#[derive(serde::Serialize)]
struct Params<'a> {
order_id: &'a str,
}
self.private_post(private::CANCEL_ORDER, &Params { order_id })
.await
}
pub async fn cancel_order_by_cli_ord_id(
&self,
cli_ord_id: &str,
) -> Result<CancelOrderResponse, KrakenError> {
#[derive(serde::Serialize)]
struct Params<'a> {
#[serde(rename = "cliOrdId")]
cli_ord_id: &'a str,
}
self.private_post(private::CANCEL_ORDER, &Params { cli_ord_id })
.await
}
pub async fn cancel_all_orders(&self) -> Result<CancelAllOrdersResponse, KrakenError> {
#[derive(serde::Serialize)]
struct Empty {}
self.private_post(private::CANCEL_ALL_ORDERS, &Empty {})
.await
}
pub async fn cancel_all_orders_for_symbol(
&self,
symbol: &str,
) -> Result<CancelAllOrdersResponse, KrakenError> {
#[derive(serde::Serialize)]
struct Params<'a> {
symbol: &'a str,
}
self.private_post(private::CANCEL_ALL_ORDERS, &Params { symbol })
.await
}
pub async fn cancel_all_orders_after(
&self,
timeout_seconds: u32,
) -> Result<CancelAllOrdersAfterResponse, KrakenError> {
#[derive(serde::Serialize)]
struct Params {
timeout: u32,
}
self.private_post(
private::CANCEL_ALL_ORDERS_AFTER,
&Params {
timeout: timeout_seconds,
},
)
.await
}
pub async fn batch_order(
&self,
request: &BatchOrderRequest,
) -> Result<BatchOrderResponse, KrakenError> {
self.private_post(private::BATCH_ORDER, request).await
}
}
impl Default for FuturesRestClient {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Debug for FuturesRestClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FuturesRestClient")
.field("base_url", &self.base_url)
.field("has_credentials", &self.credentials.is_some())
.finish()
}
}
pub struct FuturesRestClientBuilder {
base_url: String,
credentials: Option<Arc<dyn CredentialsProvider>>,
nonce_provider: Option<Arc<dyn NonceProvider>>,
user_agent: Option<String>,
max_retries: u32,
}
impl FuturesRestClientBuilder {
pub fn new() -> Self {
Self {
base_url: FUTURES_BASE_URL.to_string(),
credentials: None,
nonce_provider: None,
user_agent: None,
max_retries: 3,
}
}
pub fn base_url(mut self, url: impl Into<String>) -> Self {
self.base_url = url.into();
self
}
pub fn use_demo(mut self) -> Self {
self.base_url = crate::futures::rest::endpoints::FUTURES_DEMO_URL.to_string();
self
}
pub fn credentials(mut self, credentials: Arc<dyn CredentialsProvider>) -> Self {
self.credentials = Some(credentials);
self
}
pub fn nonce_provider(mut self, provider: Arc<dyn NonceProvider>) -> Self {
self.nonce_provider = Some(provider);
self
}
pub fn user_agent(mut self, user_agent: impl Into<String>) -> Self {
self.user_agent = Some(user_agent.into());
self
}
pub fn max_retries(mut self, retries: u32) -> Self {
self.max_retries = retries;
self
}
pub fn build(self) -> FuturesRestClient {
let mut headers = HeaderMap::new();
let user_agent = self
.user_agent
.unwrap_or_else(|| format!("kraken-api-client/{}", env!("CARGO_PKG_VERSION")));
let header_value = HeaderValue::from_str(&user_agent)
.unwrap_or_else(|_| HeaderValue::from_static("kraken-api-client"));
headers.insert(USER_AGENT, header_value);
let reqwest_client = reqwest::Client::builder()
.default_headers(headers)
.build()
.unwrap_or_else(|_| reqwest::Client::new());
let retry_policy = ExponentialBackoff::builder().build_with_max_retries(self.max_retries);
let client = ClientBuilder::new(reqwest_client)
.with(TracingMiddleware::default())
.with(RetryTransientMiddleware::new_with_policy(retry_policy))
.build();
let nonce_provider = self
.nonce_provider
.unwrap_or_else(|| Arc::new(IncreasingNonce::new()));
FuturesRestClient {
http_client: client,
base_url: self.base_url,
credentials: self.credentials,
nonce_provider,
}
}
}
impl Default for FuturesRestClientBuilder {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, serde::Deserialize)]
struct FuturesErrorResponse {
result: String,
error: Option<String>,
#[serde(rename = "serverTime")]
#[allow(dead_code)]
server_time: Option<String>,
}