use std::{fmt, io};
use gxhash::HashMap as GxHashMap;
use wdev::Device;
use super::super::storage_session::StorageSession;
pub(crate) const TAG_STRING: u8 = 0x00;
pub(crate) const TAG_META: u8 = 0x01;
pub(crate) const TAG_TTL: u8 = 0x09;
pub const CLUSTER_SLOTS: u16 = 16384;
pub(crate) fn scan_err(e: impl fmt::Display) -> wkv::Error {
wkv::Error::Io(io::Error::other(e.to_string()))
}
impl<'a, D: Device> StorageSession<'a, D> {
pub(crate) fn phys_prefix_len(&self) -> usize {
self.batch.session_prefix().as_slice().len()
}
pub async fn db_scan(
&self,
pattern: &[u8],
all_keys: bool,
cursor: &[u8],
count: usize,
) -> wkv::Result<(Vec<u8>, Vec<Vec<u8>>)> {
let (_map, keys) = self.string_snapshot().await?;
let start = keys.partition_point(|k| !cursor.is_empty() && k.as_slice() <= cursor);
let count = count.max(1);
let mut items: Vec<Vec<u8>> = Vec::new();
let mut last: Option<Vec<u8>> = None;
let mut truncated = false;
for key in keys.into_iter().skip(start) {
if items.len() >= count {
truncated = true;
break;
}
if all_keys || glob_match(pattern, &key) {
items.push(key.clone());
}
last = Some(key);
}
let next = if truncated {
last.unwrap_or_default()
} else {
Vec::new()
};
Ok((next, items))
}
pub async fn iterate_store(
&self,
mut on_record: impl FnMut(&[u8], &[u8]) -> bool,
) -> wkv::Result<usize> {
let (map, keys) = self.string_snapshot().await?;
let mut n = 0usize;
for key in keys {
let value = map.get(&key).cloned().flatten().unwrap_or_default();
n += 1;
if !on_record(&key, &value) {
break;
}
}
Ok(n)
}
pub async fn delete_slot_keys(&self, slots: &[u16]) -> wkv::Result<u64> {
let (_, keys) = self.string_snapshot().await?;
let mut deleted = 0u64;
for key in keys {
if slots.contains(&cluster_slot(&key)) && self.delete_string(&key).await? {
deleted += 1;
}
}
Ok(deleted)
}
pub async fn db_keys(&self, pattern: &[u8]) -> wkv::Result<Vec<Vec<u8>>> {
let (_, keys) = self.string_snapshot().await?;
Ok(
keys
.into_iter()
.filter(|k| glob_match(pattern, k))
.collect(),
)
}
pub async fn db_size(&self) -> wkv::Result<usize> {
let (_, keys) = self.string_snapshot().await?;
Ok(keys.len())
}
pub async fn delete_if_expired_in_memory(&self, key: &[u8]) -> wkv::Result<bool> {
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
if matches!(self.batch.probe_ttl(key, now_ms), wkv::TtlProbe::Due) {
let _ = self.delete_string(key).await?;
return Ok(true);
}
Ok(false)
}
pub fn is_internal_record(&self, physical_key: &[u8]) -> bool {
let prefix_len = self.phys_prefix_len();
physical_key.len() > prefix_len
&& (physical_key[prefix_len] == TAG_META || physical_key[prefix_len] == TAG_TTL)
}
pub(crate) async fn collect_records(&self) -> wkv::Result<GxHashMap<Vec<u8>, Option<Vec<u8>>>> {
let prefix = self.batch.session_prefix().as_slice().to_vec();
let mut map = GxHashMap::default();
self
.batch
.store
.hlog()
.scan(
self.batch.store.begin_address(),
self.batch.store.tail_address(),
|_addr, rec| {
let key = rec.key();
let Some(rest) = key.strip_prefix(prefix.as_slice()) else {
return Ok(true);
};
let Some(user_key) = rest.strip_prefix(&[TAG_STRING][..]) else {
return Ok(true);
};
if rec.is_tombstone() {
map.insert(user_key.to_vec(), None);
} else {
map.insert(user_key.to_vec(), Some(rec.value().to_vec()));
}
Ok(true)
},
)
.await
.map_err(scan_err)?;
Ok(map)
}
pub(crate) async fn string_snapshot(
&self,
) -> wkv::Result<(GxHashMap<Vec<u8>, Option<Vec<u8>>>, Vec<Vec<u8>>)> {
let map = self.collect_records().await?;
let now_ms = coarsetime::Clock::now_since_epoch().as_millis();
let mut keys: Vec<Vec<u8>> = map
.iter()
.filter(|(k, v)| {
v.is_some() && !matches!(self.batch.probe_ttl(k, now_ms), wkv::TtlProbe::Due)
})
.map(|(k, _)| k.clone())
.collect();
keys.sort_unstable();
Ok((map, keys))
}
}
pub(crate) fn glob_match(pattern: &[u8], text: &[u8]) -> bool {
let (mut p, mut t) = (0usize, 0usize);
let (mut star_p, mut star_t) = (usize::MAX, 0usize);
while t < text.len() {
if p < pattern.len() && pattern[p] == b'*' {
star_p = p;
star_t = t;
p += 1;
} else if p < pattern.len() && pattern[p] == b'[' {
let (hit, next) = match_bracket(pattern, p, text[t]);
if hit {
p = next;
t += 1;
} else if star_p != usize::MAX {
star_t += 1;
t = star_t;
p = star_p + 1;
} else {
return false;
}
} else if p < pattern.len() && (pattern[p] == b'?' || pattern[p] == text[t]) {
p += 1;
t += 1;
} else if star_p != usize::MAX {
star_t += 1;
t = star_t;
p = star_p + 1;
} else {
return false;
}
}
while p < pattern.len() && pattern[p] == b'*' {
p += 1;
}
p == pattern.len()
}
fn match_bracket(pattern: &[u8], start: usize, c: u8) -> (bool, usize) {
let mut i = start + 1;
let mut negate = false;
if i < pattern.len() && (pattern[i] == b'^' || pattern[i] == b'!') {
negate = true;
i += 1;
}
let mut hit = false;
while i < pattern.len() && pattern[i] != b']' {
if i + 2 < pattern.len() && pattern[i + 1] == b'-' && pattern[i + 2] != b']' {
if pattern[i] <= c && c <= pattern[i + 2] {
hit = true;
}
i += 3;
} else {
if pattern[i] == c {
hit = true;
}
i += 1;
}
}
if i < pattern.len() {
i += 1;
} else {
hit = negate; return (hit, i);
}
(hit != negate, i)
}
pub fn cluster_slot(key: &[u8]) -> u16 {
let hashed = hash_tag_of(key);
(crc16_xmodem(hashed) % u32::from(CLUSTER_SLOTS)) as u16
}
fn hash_tag_of(key: &[u8]) -> &[u8] {
let Some(open) = key.iter().position(|&b| b == b'{') else {
return key;
};
let Some(close_rel) = key[open + 1..].iter().position(|&b| b == b'}') else {
return key;
};
let inner = &key[open + 1..open + 1 + close_rel];
if inner.is_empty() {
return key;
}
inner
}
pub fn crc16_xmodem(data: &[u8]) -> u32 {
let mut crc = 0u32;
for &b in data {
crc ^= u32::from(b) << 8;
for _ in 0..8 {
crc = if crc & 0x8000 != 0 {
(crc << 1) ^ 0x1021
} else {
crc << 1
};
crc &= 0xFFFF;
}
}
crc
}