use rapidhash::{HashMapExt, HashSetExt, RapidHashMap as HashMap, RapidHashSet as HashSet};
use crate::{
api::hash::{
CachedFieldState, HashFieldPair, HashRandField, HashScanResult, ceil_div_1000,
conf::{
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING, ERR_INCREMENT_NAN_OR_INFINITY,
ERR_INCREMENT_OVERFLOW, HASH_EXPIRE_COND_FAILED, HASH_EXPIRE_DELETED, HASH_EXPIRE_SET_OK,
HASH_FIELD_NOT_FOUND, HASH_FIELD_PERSISTENT, HExpire, HGetEx, HSetEx, HashFetchType,
HashFieldSetCondition, HashGetExOptions, HashLengthMode, HashSetExOptions, RangeLexSpec,
TTLAction,
},
meta::{
HashFieldStateKind, HashItemKeyComposer, HashMeta, compose_hash_meta_key,
compose_hash_prefix, compose_hash_prefix_stack, decode_field_state, hexpire_condition_passes,
is_field_expired, is_immediate_expire,
},
parse_hash_float, parse_hash_integer,
traits::Hash,
},
error::{Error, Result},
key::{clear_prefix_in_batch, get_meta_checked},
key_composer::matches_glob_bytes,
meta::current_now_ms,
string::format_float_bytes,
traits::DbLike,
};
#[inline]
fn load_field_state(
data_ks: &fjall::Keyspace,
meta: &HashMeta,
item_k: &[u8],
now_ms: u64,
) -> Result<CachedFieldState> {
match data_ks.get(item_k)? {
Some(raw) => {
if let Some(state) = decode_field_state(meta, &raw, now_ms) {
Ok(CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
})
} else {
Ok(CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
})
}
}
None => Ok(CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}),
}
}
pub fn prepare_hash_meta_for_write(
db: &impl DbLike,
k_bytes: &[u8],
meta_k: &[u8],
now_ms: u64,
batch: &mut fjall::OwnedWriteBatch,
) -> Result<(HashMeta, bool)> {
let kc = db.kc();
match get_meta_checked::<HashMeta>(db, k_bytes, meta_k, now_ms)? {
Some(meta) => Ok((meta, true)),
None => {
let prefix = compose_hash_prefix_stack(&kc, k_bytes);
clear_prefix_in_batch(db.data(), &prefix, batch)?;
Ok((HashMeta::new_with_version(0, 0), false))
}
}
}
fn scan_and_repair_hash<T: DbLike + ?Sized>(
db: &T,
k_bytes: &[u8],
meta_k: &[u8],
meta: &mut HashMeta,
now_ms: u64,
) -> Result<usize> {
let kc = db.kc();
let prefix = compose_hash_prefix(&kc, k_bytes);
let mut repaired = *meta;
repaired.base.size = 0;
repaired.persist = 0;
repaired.lower = 0;
repaired.upper = 0;
let mut batch = db.batch();
let data_ks = db.data();
for g in data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(state) = decode_field_state(meta, &v, now_ms) {
match state.kind {
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(data_ks, &*k);
}
HashFieldStateKind::Persistent => {
repaired.base.size += 1;
repaired.persist += 1;
}
HashFieldStateKind::LiveTTL => {
repaired.base.size += 1;
if repaired.lower == 0 || state.expire < repaired.lower {
repaired.lower = state.expire;
}
repaired.upper = repaired.upper.max(state.expire);
}
HashFieldStateKind::Missing => {}
}
}
}
if repaired.base.size == 0 {
batch.remove(db.meta(), meta_k);
} else {
repaired.clear_bounds_if_no_ttl_candidates();
batch.insert(db.meta(), meta_k, repaired.encode());
}
batch.commit()?;
*meta = repaired;
Ok(repaired.base.size as usize)
}
impl<T: DbLike> Hash for T {
#[inline]
fn hget<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F) -> Result<Option<Vec<u8>>> {
let key_bytes = key.as_ref();
let field_bytes = field.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(None),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(None);
}
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let item_k = composer.key_for_field(field_bytes);
if let Some(raw) = self.data().get(item_k)?
&& let Some((exp, payload)) = meta.decode_subkey_value(&raw)
&& !is_field_expired(exp, now_ms)
{
Ok(Some(payload.to_vec()))
} else {
Ok(None)
}
}
#[inline]
fn hset<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
fields: &[(F, V)],
) -> Result<usize> {
if fields.is_empty() {
return Ok(0);
}
let _opts = HashSetExOptions::default();
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut batch = self.batch();
let (mut meta, metadata_existed) =
prepare_hash_meta_for_write(self, key_bytes, &meta_k, now_ms, &mut batch)?;
if metadata_existed && meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut state_cache: HashMap<&[u8], CachedFieldState> = HashMap::with_capacity(fields.len());
let data_ks = self.data();
let meta_ks = self.meta();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let mut seen = HashSet::with_capacity(fields.len());
let mut unique_fields = Vec::with_capacity(fields.len());
for (f, v) in fields.iter().rev() {
let f_bytes = f.as_ref();
if seen.insert(f_bytes) {
unique_fields.push((f_bytes, v.as_ref()));
}
}
unique_fields.reverse();
let mut inserted_count = 0usize;
for (f_bytes, v_bytes) in unique_fields {
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else if metadata_existed {
let state_entry = load_field_state(data_ks, &meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
};
match entry.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
meta.apply_missing_to_persistent();
inserted_count += 1;
}
HashFieldStateKind::LiveTTL => {
meta.apply_ttl_to_persistent();
}
HashFieldStateKind::Persistent => {}
}
let enc = meta.encode_subkey_value(v_bytes, 0);
batch.insert(data_ks, item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Persistent,
expire: 0,
raw: None,
},
);
}
if meta.base.size == 0 {
if metadata_existed {
batch.remove(meta_ks, &meta_k);
batch.commit()?;
}
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
}
Ok(inserted_count)
}
#[inline]
fn hmset<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
fields: &[(F, V)],
) -> Result<()> {
self.hset(key, fields)?;
Ok(())
}
#[inline]
fn hsetnx<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field: F,
val: V,
) -> Result<bool> {
let opts = HashSetExOptions {
condition: HashFieldSetCondition::Fnx,
..Default::default()
};
self.set_fields_with_expire(key, &[(field, val)], opts)
}
fn hdel<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<usize> {
if fields.is_empty() {
return Ok(0);
}
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(0),
};
let mut deleted = 0usize;
let mut batch = self.batch();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let data_ks = self.data();
let mut seen = HashSet::with_capacity(fields.len());
for f in fields {
let f_bytes = f.as_ref();
if !seen.insert(f_bytes) {
continue;
}
let item_k = composer.key_for_field(f_bytes);
if let Some(raw) = data_ks.get(item_k)?
&& let Some((exp, _)) = meta.decode_subkey_value(&raw)
{
batch.remove(data_ks, item_k);
if !is_field_expired(exp, now_ms) {
deleted += 1;
if exp == 0 {
meta.apply_persistent_to_deleted();
} else {
meta.apply_ttl_to_deleted();
}
} else {
meta.apply_ttl_to_deleted();
}
}
}
if deleted > 0 {
if meta.base.size == 0 {
batch.remove(self.meta(), &meta_k);
} else {
meta.clear_bounds_if_no_ttl_candidates();
batch.insert(self.meta(), &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(deleted)
}
#[inline]
fn hexists<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F) -> Result<bool> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(false),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(false);
}
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let item_k = composer.key_for_field(field.as_ref());
if let Some(raw) = self.data().get(item_k)?
&& let Some((exp, _)) = meta.decode_subkey_value(&raw)
{
return Ok(!is_field_expired(exp, now_ms));
}
Ok(false)
}
#[inline]
fn hlen<K: AsRef<[u8]>>(&self, key: K) -> Result<usize> {
self.hlen_with_mode(key, HashLengthMode::Accurate)
}
fn hlen_with_mode<K: AsRef<[u8]>>(&self, key: K, mode: HashLengthMode) -> Result<usize> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let meta_ks = self.meta();
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
if meta.is_legacy_subkey_encoding() || mode == HashLengthMode::Approximate {
return Ok(meta.base.size as usize);
}
if meta.persist > meta.base.size {
return scan_and_repair_hash(self, key_bytes, &meta_k, &mut meta, now_ms);
}
let ttl_candidates = meta.base.size.saturating_sub(meta.persist);
if ttl_candidates == 0 {
return Ok(meta.base.size as usize);
}
if meta.lower != 0 && now_ms < meta.lower {
return Ok(meta.base.size as usize);
}
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
let mut batch = self.batch();
batch.remove(meta_ks, &meta_k);
batch.commit()?;
return Ok(0);
}
scan_and_repair_hash(self, key_bytes, &meta_k, &mut meta, now_ms)
}
fn hmget<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<Option<Vec<u8>>>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(vec![None; fields.len()]),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(vec![None; fields.len()]);
}
let mut results = Vec::with_capacity(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let data_ks = self.data();
for f in fields {
let item_k = composer.key_for_field(f.as_ref());
let val = match data_ks.get(item_k)? {
Some(raw) => match meta.decode_subkey_value(&raw) {
Some((exp, payload)) => {
if is_field_expired(exp, now_ms) {
None
} else {
Some(payload.to_vec())
}
}
None => None,
},
None => None,
};
results.push(val);
}
Ok(results)
}
#[inline]
fn hgetall<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
let mut results = Vec::new();
self.hiter(key, |f, v| {
results.push((f.to_vec(), v.to_vec()));
true
})?;
Ok(results)
}
fn hget_all_with_type<K: AsRef<[u8]>>(
&self,
key: K,
fetch_type: HashFetchType,
) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
match fetch_type {
HashFetchType::All => self.hgetall(key),
HashFetchType::OnlyKey => {
let keys = self.hkeys(key)?;
Ok(keys.into_iter().map(|k| (k, Vec::new())).collect())
}
HashFetchType::OnlyValue => {
let vals = self.hvals(key)?;
Ok(vals.into_iter().map(|v| (Vec::new(), v)).collect())
}
}
}
fn hkeys<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<Vec<u8>>> {
let mut keys = Vec::new();
self.hiter(key, |f, _| {
keys.push(f.to_vec());
true
})?;
Ok(keys)
}
fn hvals<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<Vec<u8>>> {
let mut vals = Vec::new();
self.hiter(key, |_, v| {
vals.push(v.to_vec());
true
})?;
Ok(vals)
}
fn hincrby<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F, step: i64) -> Result<i64> {
let key_bytes = key.as_ref();
let field_bytes = field.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut batch = self.batch();
let (mut meta, metadata_existed) =
prepare_hash_meta_for_write(self, key_bytes, &meta_k, now_ms, &mut batch)?;
if metadata_existed && meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let item_k = composer.key_for_field(field_bytes);
let data_ks = self.data();
let meta_ks = self.meta();
let (cur_val, is_new, target_expire) = if metadata_existed {
match data_ks.get(item_k)? {
Some(raw) => match meta.decode_subkey_value(&raw) {
Some((exp, payload)) => {
if is_field_expired(exp, now_ms) {
(0i64, true, 0u64)
} else {
(parse_hash_integer(payload)?, false, exp)
}
}
None => (0i64, true, 0u64),
},
None => (0i64, true, 0u64),
}
} else {
(0i64, true, 0u64)
};
let new_val = cur_val
.checked_add(step)
.ok_or_else(|| Error::invalid_data(ERR_INCREMENT_OVERFLOW))?;
if is_new {
meta.apply_missing_to_persistent();
}
let mut itoa_buf = itoa::Buffer::new();
let val_bytes = itoa_buf.format(new_val).as_bytes();
let enc = meta.encode_subkey_value(val_bytes, target_expire);
batch.insert(data_ks, item_k, enc);
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
Ok(new_val)
}
fn hincrbyfloat<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
field: F,
step: f64,
) -> Result<f64> {
let key_bytes = key.as_ref();
let field_bytes = field.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut batch = self.batch();
let (mut meta, metadata_existed) =
prepare_hash_meta_for_write(self, key_bytes, &meta_k, now_ms, &mut batch)?;
if metadata_existed && meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let item_k = composer.key_for_field(field_bytes);
let data_ks = self.data();
let meta_ks = self.meta();
let (cur_val, is_new, target_expire) = if metadata_existed {
match data_ks.get(item_k)? {
Some(raw) => match meta.decode_subkey_value(&raw) {
Some((exp, payload)) => {
if is_field_expired(exp, now_ms) {
(0.0f64, true, 0u64)
} else {
(parse_hash_float(payload)?, false, exp)
}
}
None => (0.0f64, true, 0u64),
},
None => (0.0f64, true, 0u64),
}
} else {
(0.0f64, true, 0u64)
};
let new_val = cur_val + step;
if new_val.is_nan() || new_val.is_infinite() {
return Err(Error::invalid_data(ERR_INCREMENT_NAN_OR_INFINITY));
}
if is_new {
meta.apply_missing_to_persistent();
}
let mut f_buf = zmij::Buffer::new();
let val_bytes = format_float_bytes(new_val, &mut f_buf);
let enc = meta.encode_subkey_value(val_bytes, target_expire);
batch.insert(data_ks, item_k, enc);
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
Ok(new_val)
}
fn hstrlen<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F) -> Result<usize> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(0),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(0);
}
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let item_k = composer.key_for_field(field.as_ref());
if let Some(raw) = self.data().get(item_k)?
&& let Some((exp, payload)) = meta.decode_subkey_value(&raw)
&& !is_field_expired(exp, now_ms)
{
return Ok(payload.len());
}
Ok(0)
}
fn hiter<K: AsRef<[u8]>, F>(&self, key: K, mut f: F) -> Result<()>
where
F: FnMut(&[u8], &[u8]) -> bool,
{
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(()),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(());
}
let prefix = compose_hash_prefix(&kc, key_bytes);
let prefix_len = prefix.len();
for guard in self.data().prefix(&prefix) {
let (k, v) = guard.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some((exp, payload)) = meta.decode_subkey_value(&v)
&& !is_field_expired(exp, now_ms)
{
let field = &k[prefix_len..];
if !f(field, payload) {
break;
}
}
}
Ok(())
}
fn hrandfield<K: AsRef<[u8]>>(
&self,
key: K,
count: i64,
with_values: bool,
) -> Result<Vec<HashRandField>> {
if count == 0 {
return Ok(Vec::new());
}
let mut all = self.hgetall(key)?;
let total = all.len();
if total == 0 {
return Ok(Vec::new());
}
if count > 0 {
let sample_cnt = (count as usize).min(total);
for i in 0..sample_cnt {
let j = fastrand::usize(i..total);
all.swap(i, j);
}
all.truncate(sample_cnt);
let out = all
.into_iter()
.map(|(f, v)| (f, if with_values { Some(v) } else { None }))
.collect();
Ok(out)
} else {
let total_sample = count.unsigned_abs() as usize;
let mut out = Vec::with_capacity(total_sample);
for _ in 0..total_sample {
let idx = fastrand::usize(0..total);
let (f, v) = &all[idx];
out.push((f.clone(), if with_values { Some(v.clone()) } else { None }));
}
Ok(out)
}
}
fn hscan<K: AsRef<[u8]>>(
&self,
key: K,
cursor: usize,
limit: usize,
pattern: Option<&[u8]>,
) -> Result<HashScanResult> {
let is_match_all = match pattern {
Some(p) => p == b"*",
None => true,
};
let pat = pattern.unwrap_or(b"*");
let mut skipped = 0;
let mut matched = Vec::with_capacity(limit);
let mut has_more = false;
self.hiter(key, |field, value| {
if is_match_all || matches_glob_bytes(pat, field) {
if skipped < cursor {
skipped += 1;
} else if matched.len() < limit {
matched.push((field.to_vec(), value.to_vec()));
} else {
has_more = true;
return false;
}
}
true
})?;
let next_cursor = if has_more { cursor + matched.len() } else { 0 };
Ok((next_cursor, matched))
}
fn hexpire<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
seconds: i64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
let target_expire_ms = if seconds <= 0 {
0
} else {
now_ms.saturating_add((seconds as u64).saturating_mul(1000))
};
self.expire_fields(key, fields, target_expire_ms, condition, now_ms)
}
fn hpexpire<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
milliseconds: i64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
let target_expire_ms = if milliseconds <= 0 {
0
} else {
now_ms.saturating_add(milliseconds as u64)
};
self.expire_fields(key, fields, target_expire_ms, condition, now_ms)
}
fn hexpireat<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
unix_time_sec: u64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
let target_expire_ms = unix_time_sec.saturating_mul(1000);
self.expire_fields(key, fields, target_expire_ms, condition, now_ms)
}
fn hpexpireat<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
unix_time_ms: u64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
self.expire_fields(key, fields, unix_time_ms, condition, now_ms)
}
fn expire_fields<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
expire_at_ms: u64,
condition: HExpire,
now_ms: u64,
) -> Result<Vec<i64>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let mut meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let is_immediate = is_immediate_expire(expire_at_ms, now_ms);
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> = HashMap::with_capacity(fields.len());
let data_ks = self.data();
let meta_ks = self.meta();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = load_field_state(data_ks, &meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(data_ks, item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::Persistent | HashFieldStateKind::LiveTTL => {
if !hexpire_condition_passes(condition, entry.kind, entry.expire, expire_at_ms) {
results.push(HASH_EXPIRE_COND_FAILED);
continue;
}
if is_immediate {
batch.remove(data_ks, item_k);
if entry.kind == HashFieldStateKind::Persistent {
meta.apply_persistent_to_deleted();
} else {
meta.apply_ttl_to_deleted();
}
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(HASH_EXPIRE_DELETED);
} else {
if entry.kind == HashFieldStateKind::Persistent {
meta.apply_persistent_to_ttl(expire_at_ms);
} else {
meta.apply_ttl_to_ttl(expire_at_ms);
}
meta_changed = true;
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p)
.unwrap_or(b"");
let enc = meta.encode_subkey_value(payload, expire_at_ms);
batch.insert(data_ks, item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::LiveTTL,
expire: expire_at_ms,
raw: entry.raw,
},
);
results.push(HASH_EXPIRE_SET_OK);
}
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(results)
}
fn httl<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<Vec<i64>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let data_ks = self.data();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = composer.key_for_field(f_bytes);
let res = match data_ks.get(item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => {
let remain_ms = s.expire.saturating_sub(now_ms);
ceil_div_1000(remain_ms) as i64
}
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
fn hpttl<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<Vec<i64>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let data_ks = self.data();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = composer.key_for_field(f_bytes);
let res = match data_ks.get(item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => s.expire.saturating_sub(now_ms) as i64,
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
fn hexpiretime<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<Vec<i64>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let data_ks = self.data();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = composer.key_for_field(f_bytes);
let res = match data_ks.get(item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => (s.expire / 1000) as i64,
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
fn hpexpiretime<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<Vec<i64>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let data_ks = self.data();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = composer.key_for_field(f_bytes);
let res = match data_ks.get(item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => s.expire as i64,
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
fn hpersist<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<Vec<i64>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> = HashMap::with_capacity(fields.len());
let data_ks = self.data();
let meta_ks = self.meta();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = load_field_state(data_ks, &meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(data_ks, item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::Persistent => {
results.push(HASH_FIELD_PERSISTENT);
}
HashFieldStateKind::LiveTTL => {
meta.apply_ttl_to_persistent();
meta_changed = true;
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p)
.unwrap_or(b"");
let enc = meta.encode_subkey_value(payload, 0);
batch.insert(data_ks, item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Persistent,
expire: 0,
raw: entry.raw,
},
);
results.push(HASH_EXPIRE_SET_OK);
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(results)
}
fn hgetdel<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<Option<Vec<u8>>>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![None; fields.len()]),
};
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> = HashMap::with_capacity(fields.len());
let data_ks = self.data();
let meta_ks = self.meta();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = load_field_state(data_ks, &meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(None);
}
HashFieldStateKind::Persistent => {
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p.to_vec())
.unwrap_or_default();
results.push(Some(payload));
batch.remove(data_ks, item_k);
meta.apply_persistent_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
HashFieldStateKind::LiveTTL => {
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p.to_vec())
.unwrap_or_default();
results.push(Some(payload));
batch.remove(data_ks, item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(data_ks, item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(None);
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(results)
}
#[inline]
fn hsetex<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field_vals: &[(F, V)],
flags: &[HSetEx],
) -> Result<bool> {
let now_ms = current_now_ms();
let opts = HashSetExOptions::from_flags(flags, now_ms);
self.set_fields_with_expire(key, field_vals, opts)
}
fn set_fields_with_expire<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field_values: &[(F, V)],
options: HashSetExOptions,
) -> Result<bool> {
if field_values.is_empty() {
return Ok(false);
}
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut batch = self.batch();
let (mut meta, metadata_existed) =
prepare_hash_meta_for_write(self, key_bytes, &meta_k, now_ms, &mut batch)?;
if metadata_existed && meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut state_cache: HashMap<&[u8], CachedFieldState> =
HashMap::with_capacity(field_values.len());
let data_ks = self.data();
let meta_ks = self.meta();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
let mut meta_changed = false;
if options.condition != HashFieldSetCondition::None {
for (f, _) in field_values {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else if metadata_existed {
let state_entry = load_field_state(data_ks, &meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
};
if entry.kind == HashFieldStateKind::ExpiredTTLPhysical {
batch.remove(data_ks, item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
let exists = matches!(
entry.kind,
HashFieldStateKind::Persistent | HashFieldStateKind::LiveTTL
);
if (options.condition == HashFieldSetCondition::Fnx && exists)
|| (options.condition == HashFieldSetCondition::Fxx && !exists)
{
if meta_changed {
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
}
return Ok(false);
}
}
}
let mut seen = HashSet::with_capacity(field_values.len());
let mut unique_field_values = Vec::with_capacity(field_values.len());
for (f, v) in field_values.iter().rev() {
let f_bytes = f.as_ref();
if seen.insert(f_bytes) {
unique_field_values.push((f_bytes, v.as_ref()));
}
}
unique_field_values.reverse();
let is_immediate =
options.ttl_action == TTLAction::Set && is_immediate_expire(options.expire_at_ms, now_ms);
for (f_bytes, v_bytes) in unique_field_values {
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else if metadata_existed {
let state_entry = load_field_state(data_ks, &meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
};
if is_immediate {
match entry.kind {
HashFieldStateKind::Missing => continue,
HashFieldStateKind::Persistent => {
meta.apply_persistent_to_deleted();
}
HashFieldStateKind::LiveTTL | HashFieldStateKind::ExpiredTTLPhysical => {
meta.apply_ttl_to_deleted();
}
}
batch.remove(data_ks, item_k);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
continue;
}
let target_expire = match options.ttl_action {
TTLAction::Discard | TTLAction::Persist => 0,
TTLAction::Keep => {
if entry.kind == HashFieldStateKind::LiveTTL
|| entry.kind == HashFieldStateKind::ExpiredTTLPhysical
{
entry.expire
} else {
0
}
}
TTLAction::Set => options.expire_at_ms,
};
match entry.kind {
HashFieldStateKind::Missing => {
if target_expire == 0 {
meta.apply_missing_to_persistent();
} else {
meta.apply_missing_to_ttl(target_expire);
}
}
HashFieldStateKind::Persistent => {
if target_expire != 0 {
meta.apply_persistent_to_ttl(target_expire);
}
}
HashFieldStateKind::LiveTTL => {
if target_expire == 0 {
meta.apply_ttl_to_persistent();
} else {
meta.apply_ttl_to_ttl(target_expire);
}
}
HashFieldStateKind::ExpiredTTLPhysical => {
if target_expire == 0 {
meta.apply_missing_to_persistent();
} else {
meta.apply_missing_to_ttl(target_expire);
}
}
}
let enc = meta.encode_subkey_value(v_bytes, target_expire);
batch.insert(data_ks, item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: if target_expire == 0 {
HashFieldStateKind::Persistent
} else {
HashFieldStateKind::LiveTTL
},
expire: target_expire,
raw: None,
},
);
}
if meta.base.size == 0 {
if metadata_existed {
batch.remove(meta_ks, &meta_k);
batch.commit()?;
}
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
}
Ok(true)
}
#[inline]
fn hgetex<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
field: F,
flag: Option<HGetEx>,
) -> Result<Option<Vec<u8>>> {
let now_ms = current_now_ms();
let opts = match flag {
Some(f) => HashGetExOptions::from_flag(f, now_ms),
None => HashGetExOptions::default(),
};
let res = self.get_fields_with_expire(key, &[field], opts)?;
Ok(res.into_iter().next().flatten())
}
fn get_fields_with_expire<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
options: HashGetExOptions,
) -> Result<Vec<Option<Vec<u8>>>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![None; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let is_immediate =
options.ttl_action == TTLAction::Set && is_immediate_expire(options.expire_at_ms, now_ms);
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> = HashMap::with_capacity(fields.len());
let data_ks = self.data();
let meta_ks = self.meta();
let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
for f in fields {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = load_field_state(data_ks, &meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(None);
}
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(data_ks, item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(None);
}
HashFieldStateKind::Persistent => {
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p)
.unwrap_or(b"");
results.push(Some(payload.to_vec()));
if is_immediate {
batch.remove(data_ks, item_k);
meta.apply_persistent_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
} else if options.ttl_action == TTLAction::Set && options.expire_at_ms != 0 {
meta.apply_persistent_to_ttl(options.expire_at_ms);
meta_changed = true;
let enc = meta.encode_subkey_value(payload, options.expire_at_ms);
batch.insert(data_ks, item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::LiveTTL,
expire: options.expire_at_ms,
raw: entry.raw,
},
);
}
}
HashFieldStateKind::LiveTTL => {
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p)
.unwrap_or(b"");
results.push(Some(payload.to_vec()));
if is_immediate {
batch.remove(data_ks, item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
} else {
match options.ttl_action {
TTLAction::Persist => {
meta.apply_ttl_to_persistent();
meta_changed = true;
let enc = meta.encode_subkey_value(payload, 0);
batch.insert(data_ks, item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Persistent,
expire: 0,
raw: entry.raw,
},
);
}
TTLAction::Set => {
meta.apply_ttl_to_ttl(options.expire_at_ms);
meta_changed = true;
let enc = meta.encode_subkey_value(payload, options.expire_at_ms);
batch.insert(data_ks, item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::LiveTTL,
expire: options.expire_at_ms,
raw: entry.raw,
},
);
}
TTLAction::Keep | TTLAction::Discard => {}
}
}
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(results)
}
fn hrangebylex<K: AsRef<[u8]>>(&self, key: K, spec: RangeLexSpec) -> Result<Vec<HashFieldPair>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<HashMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
if spec.count == Some(0) {
return Ok(Vec::new());
}
let prefix = compose_hash_prefix(&kc, key_bytes);
let data_ks = self.data();
if !spec.reversed {
let limit = spec.count.unwrap_or(usize::MAX);
let mut matching = Vec::with_capacity(limit.min(128));
let mut skipped = 0;
for g in data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
let field_bytes = &k[prefix.len()..];
let min_ok = if spec.min_infinite {
true
} else if spec.minex {
field_bytes > spec.min.as_slice()
} else {
field_bytes >= spec.min.as_slice()
};
if !min_ok {
continue;
}
let max_ok = if spec.max_infinite {
true
} else if spec.maxex {
field_bytes < spec.max.as_slice()
} else {
field_bytes <= spec.max.as_slice()
};
if !max_ok {
break;
}
if let Some((exp, payload)) = meta.decode_subkey_value(&v)
&& !is_field_expired(exp, now_ms)
{
if skipped < spec.offset {
skipped += 1;
continue;
}
matching.push((field_bytes.to_vec(), payload.to_vec()));
if matching.len() >= limit {
break;
}
}
}
return Ok(matching);
}
let mut matching = Vec::new();
for g in data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
let field_bytes = &k[prefix.len()..];
let min_ok = if spec.min_infinite {
true
} else if spec.minex {
field_bytes > spec.min.as_slice()
} else {
field_bytes >= spec.min.as_slice()
};
if !min_ok {
continue;
}
let max_ok = if spec.max_infinite {
true
} else if spec.maxex {
field_bytes < spec.max.as_slice()
} else {
field_bytes <= spec.max.as_slice()
};
if !max_ok {
break;
}
if let Some((exp, payload)) = meta.decode_subkey_value(&v)
&& !is_field_expired(exp, now_ms)
{
matching.push((field_bytes.to_vec(), payload.to_vec()));
}
}
matching.reverse();
let limit = spec.count.unwrap_or(usize::MAX);
let skipped: Vec<(Vec<u8>, Vec<u8>)> =
matching.into_iter().skip(spec.offset).take(limit).collect();
Ok(skipped)
}
}