use rapidhash::v3::rapidhash_v3;
use serde::{Deserialize, Serialize};
use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[repr(u8)]
pub enum RedisType {
None = 0,
String = 1,
Hash = 2,
List = 3,
Set = 4,
ZSet = 5,
Bitmap = 6,
SortedInt = 7,
Stream = 8,
Bloom = 9,
Json = 10,
HyperLogLog = 11,
TDigest = 12,
TimeSeries = 13,
CuckooFilter = 14,
}
impl RedisType {
#[inline]
pub const fn name(&self) -> &'static str {
match self {
Self::None => "none",
Self::String => "string",
Self::Hash => "hash",
Self::List => "list",
Self::Set => "set",
Self::ZSet => "zset",
Self::Bitmap => "bitmap",
Self::SortedInt => "sortedint",
Self::Stream => "stream",
Self::Bloom => "MBbloom--",
Self::Json => "ReJSON-RL",
Self::HyperLogLog => "hyperloglog",
Self::TDigest => "TDIS-TYPE",
Self::TimeSeries => "timeseries",
Self::CuckooFilter => "MBbloomCF",
}
}
#[inline]
pub const fn from_u8(val: u8) -> Self {
match val {
1 => Self::String,
2 => Self::Hash,
3 => Self::List,
4 => Self::Set,
5 => Self::ZSet,
6 => Self::Bitmap,
7 => Self::SortedInt,
8 => Self::Stream,
9 => Self::Bloom,
10 => Self::Json,
11 => Self::HyperLogLog,
12 => Self::TDigest,
13 => Self::TimeSeries,
14 => Self::CuckooFilter,
_ => Self::None,
}
}
#[inline]
pub const fn is_single_kv_type(&self) -> bool {
matches!(self, Self::String | Self::Json)
}
#[inline]
pub const fn is_emptyable_type(&self) -> bool {
matches!(
self,
Self::String
| Self::Json
| Self::Stream
| Self::Bloom
| Self::HyperLogLog
| Self::TDigest
| Self::TimeSeries
| Self::CuckooFilter
)
}
}
pub const VERSION_COUNTER_BITS: u32 = 11;
pub const VERSION_COUNTER_MASK: u64 = (1 << VERSION_COUNTER_BITS) - 1;
static VERSION_COUNTER: AtomicU64 = AtomicU64::new(0);
pub fn init_version_counter() {
let now_nanos = coarsetime::Clock::now_since_epoch().as_nanos();
let seed = rapidhash_v3(&now_nanos.to_be_bytes());
VERSION_COUNTER.store(seed, Ordering::Relaxed);
}
#[inline]
pub fn generate_version() -> u64 {
let ts_us = coarsetime::Clock::now_since_epoch().as_micros();
let counter = VERSION_COUNTER.fetch_add(1, Ordering::Relaxed);
(ts_us << VERSION_COUNTER_BITS) | (counter & VERSION_COUNTER_MASK)
}
#[inline]
pub fn version_to_time(version: u64) -> (u64, u32) {
let ts_us = version >> VERSION_COUNTER_BITS;
let sec = ts_us / 1_000_000;
let usec = (ts_us % 1_000_000) as u32;
(sec, usec)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct KeyMeta {
pub rtype: RedisType,
pub flags: u8,
pub expire_at: u64,
pub version: u64,
pub size: u64,
}
impl KeyMeta {
pub const META_64BIT_ENCODING_MASK: u8 = 0x80;
pub const META_TYPE_MASK: u8 = 0x0F;
pub const ENCODED_SIZE: usize = 26; pub const KVROCKS_COMPLEX_ENCODED_SIZE: usize = 25; pub const KVROCKS_SINGLE_KV_ENCODED_SIZE: usize = 9;
#[inline]
pub const fn new(rtype: RedisType, expire_at: u64, version: u64, size: u64) -> Self {
Self {
rtype,
flags: 0,
expire_at,
version,
size,
}
}
#[inline]
pub fn new_with_version(rtype: RedisType, expire_at: u64, size: u64) -> Self {
Self {
rtype,
flags: 0,
expire_at,
version: generate_version(),
size,
}
}
#[inline]
pub const fn is_expired(&self, now_ms: u64) -> bool {
if !self.is_emptyable_type() && self.size == 0 {
return true;
}
self.expire_at > 0 && self.expire_at <= now_ms
}
#[inline]
pub const fn is_single_kv_type(&self) -> bool {
self.rtype.is_single_kv_type()
}
#[inline]
pub const fn is_emptyable_type(&self) -> bool {
self.rtype.is_emptyable_type()
}
#[inline]
pub const fn ttl(&self, now_ms: u64) -> i64 {
if self.expire_at == 0 {
-1
} else if self.expire_at < now_ms {
-2
} else {
(self.expire_at - now_ms) as i64
}
}
#[inline]
pub fn expire_at_ms_to_sec(ms: u64) -> u64 {
if ms == 0 {
0
} else if ms < 1000 {
1
} else {
(ms + 499) / 1000
}
}
#[inline]
pub const fn is_64bit_encoded_flags(flags: u8) -> bool {
flags & Self::META_64BIT_ENCODING_MASK != 0
}
#[inline]
pub const fn is_64bit_encoded(&self) -> bool {
Self::is_64bit_encoded_flags(self.flags)
}
#[inline]
pub const fn common_encoded_size(&self) -> usize {
if self.is_64bit_encoded() { 8 } else { 4 }
}
#[inline]
pub const fn get_offset_after_expire(flags: u8) -> usize {
if Self::is_64bit_encoded_flags(flags) {
1 + 8 } else {
1 + 4 }
}
#[inline]
pub const fn get_offset_after_size(flags: u8) -> usize {
if Self::is_64bit_encoded_flags(flags) {
1 + 8 + 8 + 8 } else {
1 + 4 + 8 + 4 }
}
#[inline]
pub fn encode(&self) -> [u8; Self::ENCODED_SIZE] {
let mut buf = [0u8; Self::ENCODED_SIZE];
buf[0] = self.rtype as u8;
buf[1] = self.flags;
buf[2..10].copy_from_slice(&self.expire_at.to_be_bytes());
buf[10..18].copy_from_slice(&self.version.to_be_bytes());
buf[18..26].copy_from_slice(&self.size.to_be_bytes());
buf
}
#[inline]
pub fn encode_kvrocks(&self) -> Vec<u8> {
let flags = Self::META_64BIT_ENCODING_MASK | (self.rtype as u8 & Self::META_TYPE_MASK);
if self.is_single_kv_type() {
let mut out = Vec::with_capacity(Self::KVROCKS_SINGLE_KV_ENCODED_SIZE);
out.push(flags);
out.extend_from_slice(&self.expire_at.to_be_bytes());
out
} else {
let mut out = Vec::with_capacity(Self::KVROCKS_COMPLEX_ENCODED_SIZE);
out.push(flags);
out.extend_from_slice(&self.expire_at.to_be_bytes());
out.extend_from_slice(&self.version.to_be_bytes());
out.extend_from_slice(&self.size.to_be_bytes());
out
}
}
#[inline]
pub fn decode(bytes: &[u8]) -> Option<Self> {
if bytes.len() >= Self::ENCODED_SIZE
&& bytes[0] <= 14
&& (bytes[1] == 0 || bytes[1] == 0x80)
{
let rtype = RedisType::from_u8(bytes[0]);
let flags = bytes[1];
let mut exp_buf = [0u8; 8];
exp_buf.copy_from_slice(&bytes[2..10]);
let expire_at = u64::from_be_bytes(exp_buf);
let mut ver_buf = [0u8; 8];
ver_buf.copy_from_slice(&bytes[10..18]);
let version = u64::from_be_bytes(ver_buf);
let mut size_buf = [0u8; 8];
size_buf.copy_from_slice(&bytes[18..26]);
let size = u64::from_be_bytes(size_buf);
return Some(Self {
rtype,
flags,
expire_at,
version,
size,
});
}
if !bytes.is_empty() && (bytes[0] & Self::META_64BIT_ENCODING_MASK != 0) {
let flags = bytes[0];
let rtype = RedisType::from_u8(flags & Self::META_TYPE_MASK);
if rtype.is_single_kv_type() {
if bytes.len() < Self::KVROCKS_SINGLE_KV_ENCODED_SIZE {
return None;
}
let mut exp_buf = [0u8; 8];
exp_buf.copy_from_slice(&bytes[1..9]);
let expire_at = u64::from_be_bytes(exp_buf);
return Some(Self {
rtype,
flags,
expire_at,
version: 0,
size: 0,
});
} else if bytes.len() >= Self::KVROCKS_COMPLEX_ENCODED_SIZE {
let mut exp_buf = [0u8; 8];
exp_buf.copy_from_slice(&bytes[1..9]);
let expire_at = u64::from_be_bytes(exp_buf);
let mut ver_buf = [0u8; 8];
ver_buf.copy_from_slice(&bytes[9..17]);
let version = u64::from_be_bytes(ver_buf);
let mut size_buf = [0u8; 8];
size_buf.copy_from_slice(&bytes[17..25]);
let size = u64::from_be_bytes(size_buf);
return Some(Self {
rtype,
flags,
expire_at,
version,
size,
});
}
}
if bytes.len() >= Self::ENCODED_SIZE {
let rtype = RedisType::from_u8(bytes[0]);
let flags = bytes[1];
let mut exp_buf = [0u8; 8];
exp_buf.copy_from_slice(&bytes[2..10]);
let expire_at = u64::from_be_bytes(exp_buf);
let mut ver_buf = [0u8; 8];
ver_buf.copy_from_slice(&bytes[10..18]);
let version = u64::from_be_bytes(ver_buf);
let mut size_buf = [0u8; 8];
size_buf.copy_from_slice(&bytes[18..26]);
let size = u64::from_be_bytes(size_buf);
Some(Self {
rtype,
flags,
expire_at,
version,
size,
})
} else {
None
}
}
}
#[inline]
pub fn normalize_range(start: i64, stop: i64, len: i64) -> (i64, i64) {
if len <= 0 {
return (0, -1);
}
let mut s = if start < 0 { len + start } else { start };
let mut e = if stop < 0 { len + stop } else { stop };
if s < 0 {
s = 0;
}
if e >= len {
e = len - 1;
}
(s, e)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_version_generation_and_time() {
init_version_counter();
let v1 = generate_version();
let v2 = generate_version();
assert!(v2 > v1);
let (sec, usec) = version_to_time(v1);
assert!(sec > 1_700_000_000);
assert!(usec < 1_000_000);
}
#[test]
fn test_key_meta_encode_decode_roundtrip() {
let meta = KeyMeta::new(RedisType::Hash, 1_800_000_000_000, 123456789, 42);
let encoded = meta.encode();
assert_eq!(encoded.len(), KeyMeta::ENCODED_SIZE);
let decoded = KeyMeta::decode(&encoded).expect("decode failed");
assert_eq!(decoded.rtype, RedisType::Hash);
assert_eq!(decoded.expire_at, 1_800_000_000_000);
assert_eq!(decoded.version, 123456789);
assert_eq!(decoded.size, 42);
}
#[test]
fn test_key_meta_kvrocks_compatibility() {
let meta = KeyMeta::new(RedisType::Set, 2_000_000_000_000, 9999, 10);
let kvrocks_enc = meta.encode_kvrocks();
assert_eq!(kvrocks_enc.len(), KeyMeta::KVROCKS_COMPLEX_ENCODED_SIZE);
let decoded = KeyMeta::decode(&kvrocks_enc).expect("decode kvrocks failed");
assert_eq!(decoded.rtype, RedisType::Set);
assert_eq!(decoded.expire_at, 2_000_000_000_000);
assert_eq!(decoded.version, 9999);
assert_eq!(decoded.size, 10);
}
}