wedb_embed 0.1.2

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

use crate::{
  api::hll::{
    algo::{HLL_DENSE_SIZE, extract_dense_hll_result, rapid_hash},
    compose_hll_data_key, compose_hll_meta_key,
    core::HyperLogLog,
    dense::{hll_dense_estimate, hll_dense_get_register, hll_dense_set_register, hll_merge_bytes},
    meta::{HllEncodeType, HyperLogLogMeta},
    sparse::{
      hll_merge_sparse_into_dense, hll_sparse_estimate, hll_sparse_set_register,
      hll_sparse_to_dense,
    },
    traits::Hll,
  },
  error::Result,
  key::get_meta_checked,
  meta::current_now_ms,
  traits::DbLike,
};

impl<T: DbLike> Hll for T {
  fn pfadd<K: AsRef<[u8]>, E: AsRef<[u8]>>(&self, key: K, elements: &[E]) -> Result<bool> {
    let k_bytes = key.as_ref();
    let kc = self.kc();
    let meta_k = compose_hll_meta_key(&kc, k_bytes);
    let data_k = compose_hll_data_key(&kc, k_bytes);
    let now_ms = current_now_ms();

    let (mut meta, exists_valid) = match get_meta_checked(self, k_bytes, &meta_k, now_ms)? {
      Some(meta) => (meta, true),
      None => (HyperLogLogMeta::new_with_version(0), false),
    };

    if elements.is_empty() {
      return Ok(false);
    }

    let data_ks = self.data();
    let meta_ks = self.meta();

    let mut registers = if exists_valid {
      data_ks
        .get(&data_k)?
        .map(|v| v.to_vec())
        .unwrap_or_else(|| vec![0u8; HLL_DENSE_SIZE])
    } else {
      vec![0u8; HLL_DENSE_SIZE]
    };

    if meta.encode_type == HllEncodeType::Dense && registers.len() < HLL_DENSE_SIZE {
      registers.resize(HLL_DENSE_SIZE, 0);
    }

    let mut updated = false;
    for el in elements {
      let hash = rapid_hash(el.as_ref());
      let (reg_idx, count) = extract_dense_hll_result(hash);
      match meta.encode_type {
        HllEncodeType::Dense => {
          let cur = hll_dense_get_register(&registers, reg_idx);
          if count > cur {
            hll_dense_set_register(&mut registers, reg_idx, count);
            updated = true;
          }
        }
        HllEncodeType::Sparse => match hll_sparse_set_register(&mut registers, reg_idx, count) {
          Ok(changed) => {
            if changed {
              updated = true;
            }
          }
          Err(_) => {
            let mut dense_buf = vec![0u8; HLL_DENSE_SIZE];
            hll_sparse_to_dense(&registers, &mut dense_buf)?;
            let cur = hll_dense_get_register(&dense_buf, reg_idx);
            if count > cur {
              hll_dense_set_register(&mut dense_buf, reg_idx, count);
              updated = true;
            }
            registers = dense_buf;
            meta.encode_type = HllEncodeType::Dense;
          }
        },
      }
    }

    if updated {
      meta.base.size = registers.len() as u64;
      let mut batch = self.batch();
      batch.insert(data_ks, &data_k, &registers);
      batch.insert(meta_ks, &meta_k, meta.encode());
      batch.commit()?;
    }

    Ok(updated)
  }

  fn pfcount<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<u64> {
    if keys.is_empty() {
      return Ok(0);
    }

    let kc = self.kc();
    let now_ms = current_now_ms();
    let data_ks = self.data();

    if keys.len() == 1 {
      let k_bytes = keys[0].as_ref();
      let meta_k = compose_hll_meta_key(&kc, k_bytes);
      let meta = match get_meta_checked::<HyperLogLogMeta>(self, k_bytes, &meta_k, now_ms)? {
        Some(meta) => meta,
        None => return Ok(0),
      };
      let data_k = compose_hll_data_key(&kc, k_bytes);
      let registers = match data_ks.get(&data_k)? {
        Some(v) => v,
        None => return Ok(0),
      };
      return match meta.encode_type {
        HllEncodeType::Dense => Ok(hll_dense_estimate(&registers)),
        HllEncodeType::Sparse => {
          Ok(hll_sparse_estimate(&registers).unwrap_or_else(|_| hll_dense_estimate(&registers)))
        }
      };
    }

    let mut seen = HashSet::default();
    let mut merged = vec![0u8; HLL_DENSE_SIZE];
    let mut has_any = false;

    for k in keys {
      let k_bytes = k.as_ref();
      if !seen.insert(k_bytes) {
        continue;
      }
      let meta_k = compose_hll_meta_key(&kc, k_bytes);
      let meta = match get_meta_checked::<HyperLogLogMeta>(self, k_bytes, &meta_k, now_ms)? {
        Some(meta) => meta,
        None => continue,
      };
      let data_k = compose_hll_data_key(&kc, k_bytes);
      if let Some(reg) = data_ks.get(&data_k)? {
        match meta.encode_type {
          HllEncodeType::Dense => {
            hll_merge_bytes(&mut merged, &reg);
          }
          HllEncodeType::Sparse => {
            hll_merge_sparse_into_dense(&mut merged, &reg);
          }
        }
        has_any = true;
      }
    }

    if !has_any {
      return Ok(0);
    }

    Ok(hll_dense_estimate(&merged))
  }

  fn pfmerge<K: AsRef<[u8]>>(&self, dest: K, sources: &[K]) -> Result<()> {
    if sources.is_empty() {
      return Ok(());
    }

    let kc = self.kc();
    let dest_bytes = dest.as_ref();
    let dest_meta_k = compose_hll_meta_key(&kc, dest_bytes);
    let dest_data_k = compose_hll_data_key(&kc, dest_bytes);
    let now_ms = current_now_ms();

    let data_ks = self.data();
    let meta_ks = self.meta();

    let (mut dest_meta, dest_valid) =
      match get_meta_checked(self, dest_bytes, &dest_meta_k, now_ms)? {
        Some(meta) => (meta, true),
        None => (HyperLogLogMeta::new_with_version(0), false),
      };

    let mut merged = if dest_valid {
      let data = data_ks
        .get(&dest_data_k)?
        .map(|v| v.to_vec())
        .unwrap_or_else(|| vec![0u8; HLL_DENSE_SIZE]);
      if dest_meta.encode_type == HllEncodeType::Sparse {
        let mut dense_buf = vec![0u8; HLL_DENSE_SIZE];
        if hll_sparse_to_dense(&data, &mut dense_buf).is_ok() {
          dense_buf
        } else {
          data
        }
      } else {
        data
      }
    } else {
      vec![0u8; HLL_DENSE_SIZE]
    };

    if merged.len() < HLL_DENSE_SIZE {
      merged.resize(HLL_DENSE_SIZE, 0);
    }

    let mut seen = HashSet::default();
    seen.insert(dest_bytes);

    for k in sources {
      let k_bytes = k.as_ref();
      if !seen.insert(k_bytes) {
        continue;
      }
      let meta_k = compose_hll_meta_key(&kc, k_bytes);
      let meta = match get_meta_checked::<HyperLogLogMeta>(self, k_bytes, &meta_k, now_ms)? {
        Some(meta) => meta,
        None => continue,
      };
      let data_k = compose_hll_data_key(&kc, k_bytes);
      if let Some(reg) = data_ks.get(&data_k)? {
        match meta.encode_type {
          HllEncodeType::Dense => {
            hll_merge_bytes(&mut merged, &reg);
          }
          HllEncodeType::Sparse => {
            hll_merge_sparse_into_dense(&mut merged, &reg);
          }
        }
      }
    }

    dest_meta.base.size = HLL_DENSE_SIZE as u64;
    dest_meta.encode_type = HllEncodeType::Dense;
    let mut batch = self.batch();
    batch.insert(data_ks, &dest_data_k, &merged);
    batch.insert(meta_ks, &dest_meta_k, dest_meta.encode());
    batch.commit()?;

    Ok(())
  }
  #[inline]
  fn pfselftest(&self) -> bool {
    HyperLogLog::selftest()
  }
}