use wdev::Device;
use wkv::{BatchStoreSession, Result};
const TTL_VAL_LEN: usize = 8;
#[inline]
pub fn now_unix_ms() -> u64 {
coarsetime::Clock::now_since_epoch().as_millis()
}
#[inline]
const fn ttl_val_decode(v: &[u8]) -> Option<u64> {
match v.split_first_chunk::<TTL_VAL_LEN>() {
Some((&arr, _)) => Some(u64::from_be_bytes(arr)),
None => None,
}
}
pub fn ttl_of_sync<D: wdev::Device>(
session: &BatchStoreSession<'_, D>,
key: &[u8],
) -> Result<Option<Option<u64>>> {
let ttl_k = session.ttl_key(key);
if !has_ttl_record(session, &ttl_k)? {
return Ok(Some(None));
}
match session.try_read_raw_in_memory(&ttl_k, ttl_val_decode)? {
None => Ok(None),
Some(v) => Ok(Some(v.flatten())),
}
}
pub fn put_ttl_sync<D: wdev::Device>(
session: &BatchStoreSession<'_, D>,
key: &[u8],
expire_at_ms: u64,
) -> Result<bool> {
let bytes = expire_at_ms.to_be_bytes();
let ttl_k = session.ttl_key(key);
let in_place = session
.try_modify_raw_in_place_unprotected(&ttl_k, |slot| {
if slot.len() == TTL_VAL_LEN {
slot.copy_from_slice(&bytes);
Some(())
} else {
None
}
})?
.is_some();
if in_place {
return Ok(true);
}
Ok(session.try_upsert_raw_sync(&ttl_k, &bytes)?.is_ok())
}
pub fn del_ttl_sync<D: wdev::Device>(
session: &BatchStoreSession<'_, D>,
key: &[u8],
) -> Result<bool> {
let ttl_k = session.ttl_key(key);
if !has_ttl_record(session, &ttl_k)? {
return Ok(true);
}
Ok(session.try_delete_raw_sync(&ttl_k)?.is_ok())
}
#[inline]
fn has_ttl_record<D: wdev::Device>(
session: &BatchStoreSession<'_, D>,
ttl_k: &[u8],
) -> Result<bool> {
Ok(session.store.index.find_tag(ttl_k).is_some())
}
pub fn data_alive_sync<D: wdev::Device>(
session: &BatchStoreSession<'_, D>,
key: &[u8],
) -> Result<Option<bool>> {
Ok(
session
.try_read_in_memory_unprotected(key, |_| ())?
.map(|found| found.is_some()),
)
}
pub fn read_adjudicated_sync<R, D: Device>(
session: &BatchStoreSession<'_, D>,
key: &[u8],
f: impl FnOnce(&[u8]) -> R,
) -> Result<Option<Option<R>>> {
let raw = session.try_read_in_memory_unprotected(key, f)?;
if raw.as_ref().is_some_and(|v| v.is_some()) {
match ttl_of_sync(session, key)? {
None => return Ok(None),
Some(Some(exp)) if exp <= now_unix_ms() => return Ok(None),
_ => {}
}
}
Ok(raw)
}
pub fn probe_alive<D: Device>(
session: &BatchStoreSession<'_, D>,
key: &[u8],
) -> Result<Option<bool>> {
let alive = match data_alive_sync(session, key)? {
None => return Ok(None),
Some(alive) => alive,
};
if !alive {
return Ok(Some(false));
}
match ttl_of_sync(session, key)? {
None => Ok(None),
Some(Some(exp)) if exp <= now_unix_ms() => Ok(None),
Some(_) => Ok(Some(true)),
}
}
#[cfg(test)]
mod tests {
use super::ttl_val_decode;
#[test]
fn ttl_val_decode_roundtrip() {
let bytes = 1_700_000_000_123_u64.to_be_bytes();
assert_eq!(ttl_val_decode(&bytes), Some(1_700_000_000_123));
assert_eq!(ttl_val_decode(&[]), None);
assert_eq!(ttl_val_decode(&[0, 0, 0]), None);
}
}