use bincode::config::Configuration;
use kivis::{BufferOp, BufferOverflowError, Repository, Storage};
use serde::{Deserialize, Serialize};
use std::ops::Range;
use thiserror::Error;
#[derive(Debug, Clone)]
pub struct Client {
base_url: String,
client: reqwest::blocking::Client,
}
#[derive(Debug, Error)]
pub enum ClientError {
#[error("HTTP error: {0}")]
Http(String),
#[error("Serialization error: {0:?}")]
Serialization(#[from] bincode::error::EncodeError),
#[error("Deserialization error: {0:?}")]
Deserialization(#[from] bincode::error::DecodeError),
#[error("JSON error: {0}")]
Json(String),
#[error("Server error: {0}")]
Server(String),
#[error("Buffer overflow error")]
BufferOverflow(#[from] BufferOverflowError),
}
#[derive(Debug, Serialize, Deserialize)]
struct InsertRequest {
key: String,
value: String,
}
#[derive(Debug, Serialize, Deserialize)]
struct GetResponse {
value: Option<String>,
}
#[derive(Debug, Serialize, Deserialize)]
struct RemoveResponse {
value: Option<String>,
}
#[derive(Debug, Serialize, Deserialize)]
struct KeysResponse {
keys: Vec<String>,
}
impl Client {
pub fn new(base_url: u16) -> Self {
Self {
base_url: format!("http://127.0.0.1:{}", base_url),
client: reqwest::blocking::Client::new(),
}
}
}
impl Storage for Client {
type Repo = Self;
type KeyUnifier = Configuration;
type ValueUnifier = Configuration;
type Container = Vec<BufferOp>;
fn repository(&self) -> &Self::Repo {
self
}
fn repository_mut(&mut self) -> &mut Self::Repo {
self
}
}
impl Repository for Client {
type K = Vec<u8>;
type V = Vec<u8>;
type Error = ClientError;
fn insert_entry(&mut self, key: &[u8], value: &[u8]) -> Result<(), Self::Error> {
let request = InsertRequest {
key: hex::encode(key),
value: hex::encode(value),
};
let response = self
.client
.post(format!("{}/insert", self.base_url))
.json(&request)
.send()
.map_err(|e| ClientError::Http(e.to_string()))?;
if response.status().is_success() {
Ok(())
} else {
Err(ClientError::Server(format!(
"Insert failed with status: {}",
response.status()
)))
}
}
fn get_entry(&self, key: &[u8]) -> Result<Option<Self::V>, Self::Error> {
let key_hex = hex::encode(key);
let response = self
.client
.get(format!("{}/get/{}", self.base_url, key_hex))
.send()
.map_err(|e| ClientError::Http(e.to_string()))?;
if response.status().is_success() {
let get_response: GetResponse = response
.json()
.map_err(|e| ClientError::Json(e.to_string()))?;
Ok(get_response
.value
.and_then(|hex_val| hex::decode(&hex_val).ok()))
} else if response.status() == reqwest::StatusCode::NOT_FOUND {
Ok(None)
} else {
Err(ClientError::Server(format!(
"Get failed with status: {}",
response.status()
)))
}
}
fn remove_entry(&mut self, key: &[u8]) -> Result<Option<Self::V>, Self::Error> {
let key_hex = hex::encode(key);
let response = self
.client
.delete(format!("{}/remove/{}", self.base_url, key_hex))
.send()
.map_err(|e| ClientError::Http(e.to_string()))?;
if response.status().is_success() {
let remove_response: RemoveResponse = response
.json()
.map_err(|e| ClientError::Json(e.to_string()))?;
Ok(remove_response
.value
.and_then(|hex_val| hex::decode(&hex_val).ok()))
} else if response.status() == reqwest::StatusCode::NOT_FOUND {
Ok(None)
} else {
Err(ClientError::Server(format!(
"Remove failed with status: {}",
response.status()
)))
}
}
fn scan_range(
&self,
range: Range<Self::K>,
) -> Result<impl Iterator<Item = Result<Self::K, Self::Error>>, Self::Error> {
let start = hex::encode(&range.start);
let end = hex::encode(&range.end);
let response = self
.client
.get(format!("{}/keys/{}/{}", self.base_url, start, end))
.send()
.map_err(|e| ClientError::Http(e.to_string()))?;
if response.status().is_success() {
let keys_response: KeysResponse = response
.json()
.map_err(|e| ClientError::Json(e.to_string()))?;
let keys: Vec<Result<Vec<u8>, ClientError>> = keys_response
.keys
.into_iter()
.filter_map(|k| hex::decode(&k).ok())
.map(Ok)
.collect();
Ok(keys.into_iter())
} else {
Err(ClientError::Server(format!(
"Keys iteration failed with status: {}",
response.status()
)))
}
}
}