use crate::error::Error;
use crate::RpcRequest;
use crate::RpcResponse;
use async_std::task::sleep;
use futures::lock::Mutex;
use isahc::Error as IsahcError;
use log::{info, warn};
use serde_json::json;
use std::sync::Arc;
use std::time::Duration;
pub struct RpcClient {
url: String,
user: Option<String>,
password: Option<String>,
id: Arc<Mutex<u64>>,
retry: bool,
backup_urls: Vec<String>,
}
impl RpcClient {
pub fn new(
url: &str,
user: Option<String>,
password: Option<String>,
retry: bool,
backup_urls: Vec<String>,
) -> Self {
RpcClient {
url: url.to_owned(),
user,
password,
id: Arc::new(Mutex::new(0)),
retry,
backup_urls,
}
}
async fn build_request(&self, method: &str, params: &[serde_json::value::Value]) -> RpcRequest {
let mut id = self.id.lock().await;
*id += 1;
RpcRequest {
method: method.to_owned(),
params: params.to_vec(),
id: json!(*id),
jsonrpc: Some("2.0".to_owned()),
}
}
pub async fn execute<T: for<'a> serde::de::Deserialize<'a>>(
&self,
method: &str,
params: &[serde_json::value::Value],
) -> Result<T, Error> {
let request = self.build_request(method, params).await;
let response = self.send_request(&request).await?;
Ok(response.into_result()?)
}
pub async fn send_request(&self, request: &RpcRequest) -> Result<RpcResponse, Error> {
let response: RpcResponse = self.send_raw(&request).await?;
if response.jsonrpc != None && response.jsonrpc != Some(From::from("2.0")) {
return Err(Error::VersionMismatch);
}
if response.id != request.id {
return Err(Error::IdMismatch);
}
Ok(response)
}
async fn send_raw<B, R>(&self, body: &B) -> Result<R, Error>
where
B: serde::ser::Serialize,
R: for<'de> serde::de::Deserialize<'de>,
{
let retry_max = 5;
let mut retries = 0;
let mut current_url = self.url.clone();
let current_backup_url = 0;
loop {
let mut req = surf::post(¤t_url);
if let Some(ref user) = self.user {
let mut auth = user.clone();
auth.push(':');
if let Some(ref pass) = self.password {
auth.push_str(&pass[..]);
}
let value = format!("Basic {}", &base64::encode(auth.as_bytes()));
req = req.header("Authorization", value);
}
let req = req.body(surf::Body::from_json(body)?);
let mut res = req.send().await?;
match res.body_json().await {
Ok(response) => return Ok(response),
Err(e) => {
warn!("RPC Request failed with error: {}", e);
if let Some(err) = &e.downcast_ref::<IsahcError>() {
match err {
IsahcError::Timeout => {
current_url = self.backup_urls[current_backup_url].clone();
}
_ => {}
}
}
if !self.retry {
return Err(Error::HttpError(e));
}
}
}
if self.retry && retries < retry_max {
retries += 1;
info!("Retrying request... Retry count: {}", retries);
sleep(Duration::from_secs(retries)).await;
continue;
} else {
return Err(Error::FailedRetry);
}
}
}
}