use crate::error::{Error, Result};
use crate::hash::conf::HExpire;
use crate::meta::{KeyMeta, RedisType, generate_version};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[repr(u8)]
pub enum HashSubkeyEncodingMode {
#[default]
Legacy = 0,
FieldExpiration = 1,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum HashFieldStateKind {
#[default]
Missing,
Persistent,
LiveTTL,
ExpiredTTLPhysical,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct HashFieldState<'a> {
pub kind: HashFieldStateKind,
pub expire: u64,
pub value: &'a [u8],
}
#[inline]
pub fn hexpire_condition_passes(
condition: HExpire,
kind: HashFieldStateKind,
current_expire_at: u64,
target_expire_at: u64,
) -> bool {
match kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => false,
HashFieldStateKind::Persistent => match condition {
HExpire::None | HExpire::Nx | HExpire::Lt => true,
HExpire::Xx | HExpire::Gt => false,
},
HashFieldStateKind::LiveTTL => match condition {
HExpire::None | HExpire::Xx => true,
HExpire::Nx => false,
HExpire::Gt => target_expire_at > current_expire_at,
HExpire::Lt => target_expire_at < current_expire_at,
},
}
}
#[inline]
pub fn decode_field_state<'a>(
meta: &HashMeta,
raw_value: &'a [u8],
now_ms: u64,
) -> Option<HashFieldState<'a>> {
let (expire, value) = meta.decode_subkey_value(raw_value)?;
let kind = if expire == 0 {
HashFieldStateKind::Persistent
} else if is_field_expired(expire, now_ms) {
HashFieldStateKind::ExpiredTTLPhysical
} else {
HashFieldStateKind::LiveTTL
};
Some(HashFieldState {
kind,
expire,
value,
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct HashMeta {
pub base: KeyMeta,
pub mode: HashSubkeyEncodingMode,
pub persist: u64,
pub lower: u64,
pub upper: u64,
}
impl HashMeta {
pub const FIELD_EXPIRATION_PREFIX_SIZE: usize = 8;
pub const ENCODED_SIZE: usize = KeyMeta::ENCODED_SIZE + 1 + 8 + 8 + 8;
#[inline]
pub fn new(expire_at: u64, version: u64, size: u64) -> Self {
Self {
base: KeyMeta::new(RedisType::Hash, expire_at, version, size),
mode: HashSubkeyEncodingMode::FieldExpiration,
persist: size,
lower: 0,
upper: 0,
}
}
#[inline]
pub fn new_with_version(expire_at: u64, size: u64) -> Self {
Self {
base: KeyMeta::new(RedisType::Hash, expire_at, generate_version(), size),
mode: HashSubkeyEncodingMode::FieldExpiration,
persist: size,
lower: 0,
upper: 0,
}
}
#[inline]
pub fn new_with_mode(
mode: HashSubkeyEncodingMode,
expire_at: u64,
version: u64,
size: u64,
) -> Self {
Self {
base: KeyMeta::new(RedisType::Hash, expire_at, version, size),
mode,
persist: size,
lower: 0,
upper: 0,
}
}
#[inline]
pub fn is_expired(&self, now_ms: u64) -> bool {
self.base.is_expired(now_ms)
}
#[inline]
pub fn is_legacy_subkey_encoding(&self) -> bool {
self.mode == HashSubkeyEncodingMode::Legacy
}
#[inline]
pub fn is_field_expiration_encoding(&self) -> bool {
self.mode == HashSubkeyEncodingMode::FieldExpiration
}
#[inline]
pub fn validate_metadata(&self) -> Result<()> {
if self.persist > self.base.size {
return Err(Error::invalid_data(
"invalid hash field expiration metadata: persist exceeds size",
));
}
Ok(())
}
#[inline]
pub fn validate_missing_field_transition(&self) -> Result<()> {
self.validate_metadata()?;
if self.base.size == u64::MAX {
return Err(Error::invalid_data(
"invalid hash field expiration metadata: size overflow",
));
}
Ok(())
}
#[inline]
pub fn validate_persistent_field_transition(&self) -> Result<()> {
self.validate_metadata()?;
if self.base.size == 0 || self.persist == 0 {
return Err(Error::invalid_data(
"invalid hash field expiration metadata: no persistent field to update",
));
}
Ok(())
}
#[inline]
pub fn validate_ttl_field_transition(&self) -> Result<()> {
self.validate_metadata()?;
if self.base.size == 0 || self.persist == self.base.size {
return Err(Error::invalid_data(
"invalid hash field expiration metadata: no TTL field to update",
));
}
Ok(())
}
#[inline]
pub fn encode_subkey_value(&self, value: &[u8], expire_at_ms: u64) -> Vec<u8> {
if self.is_legacy_subkey_encoding() {
value.to_vec()
} else {
let mut out = Vec::with_capacity(Self::FIELD_EXPIRATION_PREFIX_SIZE + value.len());
out.extend_from_slice(&expire_at_ms.to_be_bytes());
out.extend_from_slice(value);
out
}
}
#[inline]
pub fn decode_subkey_value<'a>(&self, raw: &'a [u8]) -> Option<(u64, &'a [u8])> {
if self.is_legacy_subkey_encoding() {
Some((0, raw))
} else {
if raw.len() < Self::FIELD_EXPIRATION_PREFIX_SIZE {
return None;
}
let mut exp_buf = [0u8; 8];
exp_buf.copy_from_slice(&raw[..Self::FIELD_EXPIRATION_PREFIX_SIZE]);
let expire_at = u64::from_be_bytes(exp_buf);
let payload = &raw[Self::FIELD_EXPIRATION_PREFIX_SIZE..];
Some((expire_at, payload))
}
}
#[inline]
pub fn clear_bounds_if_no_ttl_candidates(&mut self) {
if self.is_field_expiration_encoding() && self.base.size == self.persist {
self.lower = 0;
self.upper = 0;
}
}
#[inline]
pub fn expand_expire_bounds(&mut self, expire_at: u64) {
if !self.is_field_expiration_encoding() || expire_at == 0 {
return;
}
if self.base.size == self.persist {
self.lower = expire_at;
self.upper = expire_at;
return;
}
if self.lower == 0 {
self.lower = expire_at;
} else {
self.lower = self.lower.min(expire_at);
}
self.upper = self.upper.max(expire_at);
}
#[inline]
pub fn apply_missing_to_persistent(&mut self) {
self.base.size = self.base.size.saturating_add(1);
if self.is_field_expiration_encoding() {
self.persist = self.persist.saturating_add(1);
self.clear_bounds_if_no_ttl_candidates();
}
}
#[inline]
pub fn apply_missing_to_ttl(&mut self, expire_at: u64) {
self.expand_expire_bounds(expire_at);
self.base.size = self.base.size.saturating_add(1);
}
#[inline]
pub fn apply_persistent_to_ttl(&mut self, expire_at: u64) {
self.expand_expire_bounds(expire_at);
self.persist = self.persist.saturating_sub(1);
}
#[inline]
pub fn apply_ttl_to_ttl(&mut self, expire_at: u64) {
self.expand_expire_bounds(expire_at);
}
#[inline]
pub fn apply_ttl_to_persistent(&mut self) {
self.persist = self.persist.saturating_add(1).min(self.base.size);
self.clear_bounds_if_no_ttl_candidates();
}
#[inline]
pub fn apply_persistent_to_deleted(&mut self) {
self.base.size = self.base.size.saturating_sub(1);
self.persist = self.persist.saturating_sub(1);
self.clear_bounds_if_no_ttl_candidates();
}
#[inline]
pub fn apply_ttl_to_deleted(&mut self) {
self.base.size = self.base.size.saturating_sub(1);
if self.persist > self.base.size {
self.persist = self.base.size;
}
self.clear_bounds_if_no_ttl_candidates();
}
#[inline]
pub fn encode(&self) -> Vec<u8> {
if self.is_legacy_subkey_encoding() {
return self.base.encode().to_vec();
}
let mut buf = Vec::with_capacity(Self::ENCODED_SIZE);
buf.extend_from_slice(&self.base.encode());
buf.push(self.mode as u8);
buf.extend_from_slice(&self.persist.to_be_bytes());
buf.extend_from_slice(&self.lower.to_be_bytes());
buf.extend_from_slice(&self.upper.to_be_bytes());
buf
}
#[inline]
pub fn decode(bytes: &[u8]) -> Option<Self> {
if bytes.len() < KeyMeta::KVROCKS_COMPLEX_ENCODED_SIZE {
return None;
}
let base = KeyMeta::decode(bytes)?;
let base_len = if bytes.len() >= KeyMeta::ENCODED_SIZE && bytes[0] <= 14 {
KeyMeta::ENCODED_SIZE
} else {
KeyMeta::KVROCKS_COMPLEX_ENCODED_SIZE
};
if bytes.len() <= base_len {
return Some(Self {
base,
mode: HashSubkeyEncodingMode::Legacy,
persist: base.size,
lower: 0,
upper: 0,
});
}
let remain = &bytes[base_len..];
if remain.len() < 1 + 8 + 8 + 8 {
return Some(Self {
base,
mode: HashSubkeyEncodingMode::Legacy,
persist: base.size,
lower: 0,
upper: 0,
});
}
let mode = match remain[0] {
1 => HashSubkeyEncodingMode::FieldExpiration,
_ => HashSubkeyEncodingMode::Legacy,
};
let mut persist_buf = [0u8; 8];
persist_buf.copy_from_slice(&remain[1..9]);
let persist = u64::from_be_bytes(persist_buf);
let mut lower_buf = [0u8; 8];
lower_buf.copy_from_slice(&remain[9..17]);
let lower = u64::from_be_bytes(lower_buf);
let mut upper_buf = [0u8; 8];
upper_buf.copy_from_slice(&remain[17..25]);
let upper = u64::from_be_bytes(upper_buf);
Some(Self {
base,
mode,
persist,
lower,
upper,
})
}
}
pub const FIELD_EXPIRE_PREFIX_LEN: usize = 8;
#[inline]
pub fn encode_hash_value(val: &[u8], expire_at_ms: u64) -> Vec<u8> {
let mut buf = Vec::with_capacity(FIELD_EXPIRE_PREFIX_LEN + val.len());
buf.extend_from_slice(&expire_at_ms.to_be_bytes());
buf.extend_from_slice(val);
buf
}
#[inline]
pub fn decode_hash_value(bytes: &[u8]) -> (u64, &[u8]) {
if bytes.len() >= FIELD_EXPIRE_PREFIX_LEN {
let mut exp_buf = [0u8; 8];
exp_buf.copy_from_slice(&bytes[..FIELD_EXPIRE_PREFIX_LEN]);
let expire_at = u64::from_be_bytes(exp_buf);
(expire_at, &bytes[FIELD_EXPIRE_PREFIX_LEN..])
} else {
(0, bytes)
}
}
#[inline]
pub fn is_field_expired(expire_at: u64, now_ms: u64) -> bool {
expire_at > 0 && expire_at < now_ms
}
#[inline]
pub fn is_immediate_expire(expire_at: u64, now_ms: u64) -> bool {
expire_at <= now_ms
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_hash_meta_field_expiration_roundtrip() {
let mut meta = HashMeta::new(1_900_000_000_000, 1001, 3);
meta.apply_persistent_to_ttl(1_950_000_000_000);
meta.apply_persistent_to_ttl(1_920_000_000_000);
let enc = meta.encode();
assert_eq!(enc.len(), HashMeta::ENCODED_SIZE);
let decoded = HashMeta::decode(&enc).expect("decode failed");
assert_eq!(decoded.mode, HashSubkeyEncodingMode::FieldExpiration);
assert_eq!(decoded.base.size, 3);
assert_eq!(decoded.persist, 1);
assert_eq!(decoded.lower, 1_920_000_000_000);
assert_eq!(decoded.upper, 1_950_000_000_000);
}
#[test]
fn test_hash_meta_legacy_roundtrip() {
let meta = HashMeta::new_with_mode(HashSubkeyEncodingMode::Legacy, 0, 1002, 5);
let enc = meta.encode();
assert_eq!(enc.len(), KeyMeta::ENCODED_SIZE);
let decoded = HashMeta::decode(&enc).expect("decode legacy failed");
assert_eq!(decoded.mode, HashSubkeyEncodingMode::Legacy);
assert_eq!(decoded.base.size, 5);
}
#[test]
fn test_subkey_value_encoding() {
let meta_exp = HashMeta::new(0, 1, 1);
let enc_val = meta_exp.encode_subkey_value(b"hello", 123456);
let (exp, payload) = meta_exp
.decode_subkey_value(&enc_val)
.expect("decode subkey failed");
assert_eq!(exp, 123456);
assert_eq!(payload, b"hello");
let meta_legacy = HashMeta::new_with_mode(HashSubkeyEncodingMode::Legacy, 0, 1, 1);
let enc_legacy = meta_legacy.encode_subkey_value(b"world", 0);
assert_eq!(enc_legacy, b"world");
let (exp2, payload2) = meta_legacy
.decode_subkey_value(&enc_legacy)
.expect("decode legacy subkey failed");
assert_eq!(exp2, 0);
assert_eq!(payload2, b"world");
}
#[test]
fn test_state_machine_transitions() {
let mut meta = HashMeta::new(0, 100, 0);
assert_eq!(meta.persist, 0);
assert_eq!(meta.base.size, 0);
meta.apply_missing_to_persistent();
assert_eq!(meta.base.size, 1);
assert_eq!(meta.persist, 1);
assert_eq!(meta.lower, 0);
assert_eq!(meta.upper, 0);
meta.apply_missing_to_ttl(2000);
assert_eq!(meta.base.size, 2);
assert_eq!(meta.persist, 1);
assert_eq!(meta.lower, 2000);
assert_eq!(meta.upper, 2000);
meta.apply_missing_to_ttl(1000);
assert_eq!(meta.base.size, 3);
assert_eq!(meta.persist, 1);
assert_eq!(meta.lower, 1000);
assert_eq!(meta.upper, 2000);
meta.apply_persistent_to_ttl(3000);
assert_eq!(meta.base.size, 3);
assert_eq!(meta.persist, 0);
assert_eq!(meta.lower, 1000);
assert_eq!(meta.upper, 3000);
meta.apply_ttl_to_persistent();
assert_eq!(meta.base.size, 3);
assert_eq!(meta.persist, 1);
meta.apply_ttl_to_deleted();
assert_eq!(meta.base.size, 2);
assert_eq!(meta.persist, 1);
meta.apply_persistent_to_deleted();
assert_eq!(meta.base.size, 1);
assert_eq!(meta.persist, 0);
}
#[test]
fn test_hexpire_conditions() {
assert!(hexpire_condition_passes(
HExpire::None,
HashFieldStateKind::Persistent,
0,
5000
));
assert!(hexpire_condition_passes(
HExpire::Nx,
HashFieldStateKind::Persistent,
0,
5000
));
assert!(!hexpire_condition_passes(
HExpire::Xx,
HashFieldStateKind::Persistent,
0,
5000
));
assert!(!hexpire_condition_passes(
HExpire::Gt,
HashFieldStateKind::Persistent,
0,
5000
));
assert!(hexpire_condition_passes(
HExpire::Lt,
HashFieldStateKind::Persistent,
0,
5000
));
assert!(hexpire_condition_passes(
HExpire::None,
HashFieldStateKind::LiveTTL,
3000,
5000
));
assert!(!hexpire_condition_passes(
HExpire::Nx,
HashFieldStateKind::LiveTTL,
3000,
5000
));
assert!(hexpire_condition_passes(
HExpire::Xx,
HashFieldStateKind::LiveTTL,
3000,
5000
));
assert!(hexpire_condition_passes(
HExpire::Gt,
HashFieldStateKind::LiveTTL,
3000,
5000
));
assert!(!hexpire_condition_passes(
HExpire::Gt,
HashFieldStateKind::LiveTTL,
3000,
2000
));
assert!(hexpire_condition_passes(
HExpire::Lt,
HashFieldStateKind::LiveTTL,
3000,
2000
));
assert!(!hexpire_condition_passes(
HExpire::Lt,
HashFieldStateKind::LiveTTL,
3000,
5000
));
assert!(!hexpire_condition_passes(
HExpire::None,
HashFieldStateKind::ExpiredTTLPhysical,
1000,
5000
));
}
}