use wdev::Device;
use super::super::storage_session::StorageSession;
use crate::{api::garnet_status::GarnetStatus, objects::types::object_output::ObjectOutput};
pub(crate) const OBJ_TAG_SORTED_SET: u8 = 1;
pub(crate) const OBJ_TAG_LIST: u8 = 2;
pub(crate) const OBJ_TAG_HASH: u8 = 3;
pub(crate) const OBJ_TAG_SET: u8 = 4;
pub(crate) enum ObjState {
Absent,
WrongType,
Present(Vec<u8>),
}
pub(crate) fn obj_encode(tag: u8, payload: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(payload.len() + 1);
out.push(tag);
out.extend_from_slice(payload);
out
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RmwOutcome<R> {
Written(R),
Aborted,
WrongType,
}
impl<R> RmwOutcome<R> {
#[inline]
pub(crate) fn unwrap_or(self, fallback: R) -> R {
match self {
Self::Written(r) => r,
Self::Aborted | Self::WrongType => fallback,
}
}
}
pub(crate) fn obj_decode(raw: &[u8], want: u8) -> Option<&[u8]> {
raw
.split_first()
.filter(|(t, _)| **t == want)
.map(|(_, p)| p)
}
impl<'a, D: Device> StorageSession<'a, D> {
pub(crate) async fn obj_load(&self, key: &[u8], tag: u8) -> wkv::Result<ObjState> {
match self.read_string(key).await? {
None => Ok(ObjState::Absent),
Some(raw) => Ok(match obj_decode(&raw, tag) {
Some(p) => ObjState::Present(p.to_vec()),
None => ObjState::WrongType,
}),
}
}
pub(crate) async fn obj_save(&self, key: &[u8], tag: u8, payload: &[u8]) -> wkv::Result<()> {
self.upsert_string(key, &obj_encode(tag, payload)).await
}
pub(crate) async fn rmw_object_store_operation<R>(
&self,
key: &[u8],
tag: u8,
on_load: impl FnOnce(Option<Vec<u8>>) -> Option<(Vec<u8>, R)>,
) -> wkv::Result<RmwOutcome<R>> {
match self.read_string(key).await? {
Some(raw) if obj_decode(&raw, tag).is_none() => Ok(RmwOutcome::WrongType),
current => {
let input = current.and_then(|raw| obj_decode(&raw, tag).map(<[u8]>::to_vec));
if let Some((payload, r)) = on_load(input) {
self.obj_save(key, tag, &payload).await?;
Ok(RmwOutcome::Written(r))
} else {
Ok(RmwOutcome::Aborted)
}
}
}
}
pub(crate) async fn read_object_store_operation(
&self,
key: &[u8],
tag: u8,
) -> wkv::Result<(GarnetStatus, Option<Vec<u8>>)> {
match self.obj_load(key, tag).await? {
ObjState::Absent => Ok((GarnetStatus::NotFound, None)),
ObjState::WrongType => Ok((GarnetStatus::WrongType, None)),
ObjState::Present(p) => Ok((GarnetStatus::Ok, Some(p))),
}
}
pub(crate) async fn finalize_removal<R>(
&self,
key: &[u8],
removed: RmwOutcome<(R, bool)>,
fallback: R,
) -> wkv::Result<R> {
match removed {
RmwOutcome::Written((r, true)) => {
let _ = self.delete_string(key).await?;
Ok(r)
}
RmwOutcome::Written((r, false)) => Ok(r),
RmwOutcome::Aborted | RmwOutcome::WrongType => Ok(fallback),
}
}
pub(crate) async fn object_scan(
&self,
key: &[u8],
tag: u8,
pattern: &[u8],
cursor: &[u8],
count: usize,
members_of: impl Fn(&[u8]) -> Option<Vec<Vec<u8>>>,
) -> wkv::Result<(GarnetStatus, Vec<u8>, Vec<Vec<u8>>)> {
let payload = match self.obj_load(key, tag).await? {
ObjState::Absent => return Ok((GarnetStatus::NotFound, Vec::new(), Vec::new())),
ObjState::WrongType => return Ok((GarnetStatus::WrongType, Vec::new(), Vec::new())),
ObjState::Present(p) => p,
};
let Some(mut members) = members_of(&payload) else {
return Ok((GarnetStatus::WrongType, Vec::new(), Vec::new()));
};
members.sort();
let count = count.max(1);
let mut items = Vec::new();
let mut last: Option<Vec<u8>> = None;
let mut truncated = false;
for m in members {
if !cursor.is_empty() && m.as_slice() <= cursor {
continue;
}
if items.len() >= count {
truncated = true;
break;
}
if pattern.is_empty()
|| super::super::common::array_key_iteration_functions::glob_match(pattern, &m)
{
items.push(m.clone());
}
last = Some(m);
}
let next = if truncated {
last.unwrap_or_default()
} else {
Vec::new()
};
Ok((GarnetStatus::Ok, next, items))
}
pub async fn delete_object_store(&self, key: &[u8]) -> wkv::Result<GarnetStatus> {
if self.delete_string(key).await? {
Ok(GarnetStatus::Ok)
} else {
Ok(GarnetStatus::NotFound)
}
}
pub async fn object_collect(
&self,
mut on_object: impl FnMut(u8, &[u8]) -> bool,
) -> wkv::Result<usize> {
let map = self.collect_records().await?;
let mut n = 0usize;
let mut keys: Vec<&Vec<u8>> = map.keys().collect();
keys.sort();
for key in keys {
if let Some(Some(v)) = map.get(key)
&& let Some(&tag) = v.first()
&& (OBJ_TAG_SORTED_SET..=OBJ_TAG_SET).contains(&tag)
{
n += 1;
if !on_object(tag, key) {
break;
}
}
}
Ok(n)
}
pub fn throw_object_store_uninitialized_exception(&self) -> bool {
true
}
pub fn complete_pending_and_get_garnet_status(&self, status: GarnetStatus) -> GarnetStatus {
status
}
pub fn copy_pinned_span_byte_array_to_scratch_buffer(items: &[Vec<u8>]) -> Vec<Vec<u8>> {
items.to_vec()
}
pub fn copy_pinned_span_byte_pairs_to_scratch_buffer(
pairs: &[(Vec<u8>, Vec<u8>)],
) -> Vec<(Vec<u8>, Vec<u8>)> {
pairs.to_vec()
}
pub fn is_supported_array_type(tag: u8) -> bool {
(OBJ_TAG_SORTED_SET..=OBJ_TAG_SET).contains(&tag)
}
pub fn process_resp_array_output(&self, output: &mut Vec<u8>, items: &[Vec<u8>]) {
push_resp_array(output, &items.iter().map(Vec::as_slice).collect::<Vec<_>>());
}
pub fn process_resp2_array_output(&self, output: &mut Vec<u8>, items: &[Vec<u8>]) {
push_resp_array(output, &items.iter().map(Vec::as_slice).collect::<Vec<_>>());
}
pub fn process_resp3_array_output(&self, output: &mut Vec<u8>, items: &[Vec<u8>]) {
self.process_resp2_array_output(output, items);
}
pub fn process_resp_integer_array_output(&self, output: &mut Vec<u8>, items: &[i64]) {
output.extend_from_slice(format!("*{}\r\n", items.len()).as_bytes());
for v in items {
output.push(b':');
output.extend_from_slice(itoa::Buffer::new().format(*v).as_bytes());
output.extend_from_slice(b"\r\n");
}
}
pub fn process_resp_int64_array_output(&self, output: &mut Vec<u8>, items: &[i64]) {
self.process_resp_integer_array_output(output, items);
}
pub fn process_resp_array_output_as_pairs(&self, output: &mut Vec<u8>, pairs: &[(Vec<u8>, f64)]) {
let mut flat = Vec::with_capacity(pairs.len() * 2);
for (m, s) in pairs {
flat.push(m.clone());
flat.push(ObjectOutput::format_double(*s).into_bytes());
}
self.process_resp2_array_output(output, &flat);
}
pub fn process_resp_single_token_output(&self, output: &mut Vec<u8>, token: &[u8]) {
output.push(b'$');
output.extend_from_slice(itoa::Buffer::new().format(token.len()).as_bytes());
output.extend_from_slice(b"\r\n");
output.extend_from_slice(token);
output.extend_from_slice(b"\r\n");
}
pub fn try_process_resp_simple64_int_output(&self, output: &mut Vec<u8>, value: i64) -> bool {
output.push(b':');
output.extend_from_slice(itoa::Buffer::new().format(value).as_bytes());
output.extend_from_slice(b"\r\n");
true
}
}
pub(crate) fn push_resp_array(output: &mut Vec<u8>, items: &[&[u8]]) {
output.extend_from_slice(format!("*{}\r\n", items.len()).as_bytes());
for item in items {
output.push(b'$');
output.extend_from_slice(itoa::Buffer::new().format(item.len()).as_bytes());
output.extend_from_slice(b"\r\n");
output.extend_from_slice(item);
output.extend_from_slice(b"\r\n");
}
}