wedb_embed 0.1.13

Embedded database engine providing Redis-like APIs, built on fjall / 嵌入式数据库引擎,提供类似 Redis 的接口,底层基于 fjall 开发
Documentation
use rapidhash::{HashMapExt, HashSetExt, RapidHashMap as HashMap, RapidHashSet as HashSet};

use crate::{
  api::hash::{
    CachedFieldState,
    r#const::ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
    hfe::{
      apply_setex_field_in_batch, commit_hash_batch, load_field_state, remove_field_in_batch,
      resolve_target_expire,
    },
    meta::{HashFieldStateKind, HashItemKeyComposer, compose_hash_meta_key, is_immediate_expire},
    opt::{HSet, HashFieldSetCondition, HashSetEx, TTLAction},
    prepare_hash_meta_for_write,
  },
  engine::Engine,
  error::{Error, Result},
  meta::current_now_ms,
  wedb::Db,
};

impl<E: Engine> Db<E>
where
  Error: From<E::Error>,
{
  #[inline]
  pub fn hsetex_one<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
    &self,
    key: K,
    field: F,
    val: V,
    opt_li: impl IntoIterator<Item = HSet>,
  ) -> Result<bool> {
    let now_ms = current_now_ms();
    let opts = HashSetEx::from_options(opt_li, now_ms);
    self.set_field_with_expire_one(key, field, val, opts, now_ms)
  }

  #[inline]
  pub fn hsetex<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
    &self,
    key: K,
    field_values: &[(F, V)],
    opt_li: impl IntoIterator<Item = HSet>,
  ) -> Result<bool> {
    let now_ms = current_now_ms();
    let opts = HashSetEx::from_options(opt_li, now_ms);
    self.set_fields_with_expire(key, field_values, opts)
  }

  #[inline]
  pub(crate) fn set_field_with_expire_one<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
    &self,
    key: K,
    field: F,
    val: V,
    options: HashSetEx,
    now_ms: u64,
  ) -> Result<bool> {
    let key_bytes = key.as_ref();
    let kc = self.kc();
    let meta_k = compose_hash_meta_key(&kc, key_bytes);

    let mut batch = self.batch_with_capacity(2);
    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 f_bytes = field.as_ref();
    let v_bytes = val.as_ref();
    let mut composer = HashItemKeyComposer::new(&kc, key_bytes);
    let item_k = composer.key_for_field(f_bytes);

    let entry = if metadata_existed {
      load_field_state(self.data(), &meta, item_k, now_ms)?
    } else {
      CachedFieldState {
        kind: HashFieldStateKind::Missing,
        expire: 0,
        raw: None,
      }
    };

    match options.condition {
      HashFieldSetCondition::None => {}
      HashFieldSetCondition::Fnx => {
        if entry.kind != HashFieldStateKind::Missing
          && entry.kind != HashFieldStateKind::ExpiredTTLPhysical
        {
          return Ok(false);
        }
      }
      HashFieldSetCondition::Fxx => {
        if entry.kind != HashFieldStateKind::Persistent && entry.kind != HashFieldStateKind::LiveTTL
        {
          return Ok(false);
        }
      }
    }

    let is_immediate =
      options.ttl_action == TTLAction::Set && is_immediate_expire(options.expire_at_ms, now_ms);

    if is_immediate {
      if entry.kind != HashFieldStateKind::Missing {
        remove_field_in_batch(&mut meta, item_k, entry.kind, &mut batch);
        commit_hash_batch(&meta_k, &mut meta, batch)?;
      }
      return Ok(true);
    }

    let target_expire = resolve_target_expire(
      options.ttl_action,
      options.expire_at_ms,
      entry.kind,
      entry.expire,
    );
    apply_setex_field_in_batch(
      &mut meta,
      item_k,
      v_bytes,
      entry.kind,
      target_expire,
      &mut batch,
    );
    commit_hash_batch(&meta_k, &mut meta, batch)?;
    Ok(true)
  }

  #[inline]
  pub(crate) fn set_fields_with_expire<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
    &self,
    key: K,
    field_values: &[(F, V)],
    options: HashSetEx,
  ) -> Result<bool> {
    if field_values.is_empty() {
      return Ok(false);
    }
    let now_ms = current_now_ms();
    if field_values.len() == 1 {
      return self.set_field_with_expire_one(
        key,
        &field_values[0].0,
        &field_values[0].1,
        options,
        now_ms,
      );
    }

    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 mut composer = HashItemKeyComposer::new(&kc, key_bytes);

    // 1. 先行校验 Fnx / Fxx 前置条件
    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 state_entry = if metadata_existed {
          load_field_state(data_ks, &meta, item_k, now_ms)?
        } else {
          CachedFieldState {
            kind: HashFieldStateKind::Missing,
            expire: 0,
            raw: None,
          }
        };
        state_cache.insert(f_bytes, state_entry.clone());

        let condition_met = match options.condition {
          HashFieldSetCondition::None => true,
          HashFieldSetCondition::Fnx => {
            state_entry.kind == HashFieldStateKind::Missing
              || state_entry.kind == HashFieldStateKind::ExpiredTTLPhysical
          }
          HashFieldSetCondition::Fxx => {
            state_entry.kind == HashFieldStateKind::Persistent
              || state_entry.kind == HashFieldStateKind::LiveTTL
          }
        };

        if !condition_met {
          return Ok(false);
        }
      }
    }

    // 2. 去重并逆序保留最新值
    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 {
        if entry.kind != HashFieldStateKind::Missing {
          remove_field_in_batch(&mut meta, item_k, entry.kind, &mut batch);
        }
        state_cache.insert(
          f_bytes,
          CachedFieldState {
            kind: HashFieldStateKind::Missing,
            expire: 0,
            raw: None,
          },
        );
        continue;
      }

      let target_expire = resolve_target_expire(
        options.ttl_action,
        options.expire_at_ms,
        entry.kind,
        entry.expire,
      );
      apply_setex_field_in_batch(
        &mut meta,
        item_k,
        v_bytes,
        entry.kind,
        target_expire,
        &mut batch,
      );
      state_cache.insert(
        f_bytes,
        CachedFieldState {
          kind: if target_expire == 0 {
            HashFieldStateKind::Persistent
          } else {
            HashFieldStateKind::LiveTTL
          },
          expire: target_expire,
          raw: None,
        },
      );
    }

    commit_hash_batch(&meta_k, &mut meta, batch)?;
    Ok(true)
  }
}