use std::{str, sync::atomic::Ordering::Relaxed};
use wdev::Device;
use super::super::storage_session::StorageSession;
use crate::{api::garnet_status::GarnetStatus, objects::types::object_output::ObjectOutput};
#[derive(Debug, Clone, Copy)]
pub enum StringRMWOp<'k> {
Incr { delta: i64 },
IncrFloat { delta: f64 },
Append(&'k [u8]),
SetRange { offset: usize, data: &'k [u8] },
}
pub type RmwResult = (GarnetStatus, Option<i64>);
impl<'a, D: Device> StorageSession<'a, D> {
pub async fn get_with_pending(
&self,
key: &[u8],
) -> wkv::Result<(GarnetStatus, Option<Vec<u8>>, bool)> {
let (status, val) = self.read_with_unsafe_context(key).await?;
Ok((status, val, false))
}
pub fn get_complete_pending(&self) -> bool {
true
}
pub async fn rmw_main_store(&self, key: &[u8], op: StringRMWOp<'_>) -> wkv::Result<RmwResult> {
match op {
StringRMWOp::Incr { delta } => {
let (status, v) = self.increment(key, delta).await?;
Ok((status, (status == GarnetStatus::Ok).then_some(v)))
}
StringRMWOp::IncrFloat { delta } => {
let (status, new_val) = self.increment_by_float(key, delta).await?;
Ok((status, new_val.map(|f| f as i64)))
}
StringRMWOp::Append(data) => {
let (status, len) = self.append(key, data).await?;
Ok((status, Some(len as i64)))
}
StringRMWOp::SetRange { offset, data } => {
let (status, len) = self.setrange(key, offset, data).await?;
Ok((status, Some(len as i64)))
}
}
}
pub async fn read_main_store(&self, key: &[u8]) -> wkv::Result<(GarnetStatus, Option<Vec<u8>>)> {
self.read_with_unsafe_context(key).await
}
pub async fn read_range_index(&self, key: &[u8]) -> wkv::Result<(GarnetStatus, Option<Vec<u8>>)> {
self.read_with_unsafe_context(key).await
}
pub async fn read_with_prefetch<K: AsRef<[u8]>>(
&self,
keys: &[K],
mut on_item: impl FnMut(usize, Option<&[u8]>),
) -> wkv::Result<()> {
for (i, k) in keys.iter().enumerate() {
let v = self.read_string(k.as_ref()).await?;
on_item(i, v.as_deref());
}
Ok(())
}
pub(crate) async fn increment(&self, key: &[u8], delta: i64) -> wkv::Result<(GarnetStatus, i64)> {
let Some(val) = self.read_string(key).await? else {
self
.upsert_string(key, itoa::Buffer::new().format(delta).as_bytes())
.await?;
return Ok((GarnetStatus::Ok, delta));
};
let Some(current) = str::from_utf8(&val)
.ok()
.and_then(|s| s.parse::<i64>().ok())
else {
return Ok((GarnetStatus::WrongType, 0));
};
let Some(updated) = current.checked_add(delta) else {
return Ok((GarnetStatus::WrongType, current));
};
self
.upsert_string(key, itoa::Buffer::new().format(updated).as_bytes())
.await?;
Ok((GarnetStatus::Ok, updated))
}
pub(crate) async fn increment_by_float(
&self,
key: &[u8],
delta: f64,
) -> wkv::Result<(GarnetStatus, Option<f64>)> {
let Some(val) = self.read_string(key).await? else {
self
.upsert_string(key, ObjectOutput::format_double(delta).as_bytes())
.await?;
return Ok((GarnetStatus::Ok, Some(delta)));
};
let Some(current) = parse_f64(&val) else {
return Ok((GarnetStatus::WrongType, None));
};
let updated = current + delta;
self
.upsert_string(key, ObjectOutput::format_double(updated).as_bytes())
.await?;
Ok((GarnetStatus::Ok, Some(updated)))
}
pub(crate) fn note_notfound(&self) {
self.session_notfound.fetch_add(1, Relaxed);
}
}
pub(crate) fn parse_f64(bytes: &[u8]) -> Option<f64> {
str::from_utf8(bytes)
.ok()?
.trim()
.parse::<f64>()
.ok()
.filter(|v| v.is_finite())
}