wedb_embed 0.1.1

Embedded database engine providing Redis-like APIs, built on fjall / 嵌入式数据库引擎,提供类似 Redis 的接口,底层基于 fjall 开发
Documentation
pub mod algo;
pub mod core;
pub mod dense;
pub mod meta;
pub mod sparse;

pub use algo::{
    HLL_ALPHA_INF, HLL_DENSE_SIZE, HLL_HASH_BIT_COUNT, HLL_HASH_SEED, HLL_REGISTER_BITS,
    HLL_REGISTER_COUNT_MASK, HLL_REGISTER_COUNT_POW, HLL_REGISTER_MAX, HLL_REGISTERS,
    HLL_SEGMENT_BYTES, HLL_SEGMENT_COUNT, HLL_SEGMENT_REGISTERS, extract_dense_hll_result,
    hll_estimate_from_histo, hll_murmur_hash_64a, hll_sigma, hll_tau, murmur_hash_64a, rapid_hash,
};
pub use core::HyperLogLog;
pub use dense::{
    dense_estimate, get_register, hll_dense_estimate, hll_dense_estimate_segments,
    hll_dense_get_register, hll_dense_reg_histo, hll_dense_set_register, hll_merge_bytes,
    hll_merge_segments, set_register,
};
pub use meta::{HllEncodeType, HyperLogLogMeta};
pub use sparse::{
    HLL_SPARSE_MAX_BYTES, HLL_SPARSE_VAL_MAX_LEN, HLL_SPARSE_VAL_MAX_VALUE,
    HLL_SPARSE_XZERO_MAX_LEN, HLL_SPARSE_ZERO_MAX_LEN, HllSparseOp, decode_sparse_op,
    encode_sparse_val, encode_sparse_zero, hll_dense_to_sparse, hll_merge_sparse_into_dense,
    hll_sparse_estimate, hll_sparse_get_register, hll_sparse_is_valid, hll_sparse_new,
    hll_sparse_reg_histo, hll_sparse_set_register, hll_sparse_to_dense,
};

use crate::db::WeDb;
use crate::error::Result;
use crate::key_composer::{KeyComposer, KeyTag};
use crate::meta::current_now_ms;
use rapidhash::RapidHashSet as HashSet;

/// 构造 HLL 元数据键字节序列(二进制安全,支持任意字节序列)
#[inline]
fn compose_hll_meta_key(kc: &KeyComposer, key: &[u8]) -> Vec<u8> {
    kc.compose_meta_key(KeyTag::HllMeta.as_slice(), key)
}

/// 构造 HLL 数据键字节序列(二进制安全,支持任意字节序列)
#[inline]
fn compose_hll_data_key(kc: &KeyComposer, key: &[u8]) -> Vec<u8> {
    kc.compose_meta_key(KeyTag::HllRaw.as_slice(), key)
}

impl WeDb {
    #[inline]
    fn get_hll_meta_checked(
        &self,
        kc: &KeyComposer<'_>,
        k: &[u8],
        mk: &[u8],
        now: u64,
    ) -> Result<Option<HyperLogLogMeta>> {
        self.get_meta_checked(kc, k, mk, now)
    }

    /// PFADD key element [element ...](带命名空间隔离,对标 Apache Kvrocks HyperLogLog::Add)
    pub fn pfadd_ns<K: AsRef<[u8]>, E: AsRef<[u8]>>(
        &self,
        ns: &str,
        key: K,
        elements: &[E],
    ) -> Result<bool> {
        let k_bytes = key.as_ref();
        let kc = KeyComposer::new(ns);
        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 self.get_hll_meta_checked(&kc, k_bytes, &meta_k, now_ms)? {
                Some(meta) => (meta, true),
                None => (HyperLogLogMeta::new_with_version(0), false),
            };

        if elements.is_empty() {
            // 对标 Kvrocks: 空 elements 时若 key 不存在则不写入,若 key 存在且类型正确则返回 false
            return Ok(false);
        }

        let mut registers = if exists_valid {
            self.data
                .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(_) => {
                            // 稀疏超限,自动晋升为 Dense
                            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.db.batch();
            batch.insert(&self.data, &data_k, &registers);
            batch.insert(&self.meta, &meta_k, meta.encode());
            batch.commit()?;
        }

        Ok(updated)
    }

    /// PFADD key element [element ...]
    #[inline]
    pub fn pfadd<K: AsRef<[u8]>, E: AsRef<[u8]>>(&self, key: K, elements: &[E]) -> Result<bool> {
        self.pfadd_ns("default", key, elements)
    }

    /// PFCOUNT key [key ...](带命名空间隔离,对标 Apache Kvrocks HyperLogLog::Count / CountMultiple)
    pub fn pfcount_ns<K: AsRef<[u8]>>(&self, ns: &str, keys: &[K]) -> Result<u64> {
        if keys.is_empty() {
            return Ok(0);
        }

        let kc = KeyComposer::new(ns);
        let now_ms = current_now_ms();

        if keys.len() == 1 {
            let k_bytes = keys[0].as_ref();
            let meta_k = compose_hll_meta_key(&kc, k_bytes);
            let meta = match self.get_hll_meta_checked(&kc, 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 self.data.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))),
            };
        }

        // 合并多个 HLL 寄存器进行多键基数统计(过滤重复 key 与过期 key,零中间堆分配)
        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 self.get_hll_meta_checked(&kc, 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) = self.data.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))
    }

    /// PFCOUNT key [key ...]
    #[inline]
    pub fn pfcount<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<u64> {
        self.pfcount_ns("default", keys)
    }

    /// PFMERGE destkey sourcekey [sourcekey ...](带命名空间隔离,对标 Apache Kvrocks HyperLogLog::Merge)
    pub fn pfmerge_ns<K: AsRef<[u8]>>(&self, ns: &str, dest: K, sources: &[K]) -> Result<()> {
        if sources.is_empty() {
            return Ok(());
        }

        let kc = KeyComposer::new(ns);
        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 (mut dest_meta, dest_valid) =
            match self.get_hll_meta_checked(&kc, 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 = self
                .data
                .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 self.get_hll_meta_checked(&kc, 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) = self.data.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.db.batch();
        batch.insert(&self.data, &dest_data_k, &merged);
        batch.insert(&self.meta, &dest_meta_k, dest_meta.encode());
        batch.commit()?;

        Ok(())
    }

    /// PFMERGE destkey sourcekey [sourcekey ...]
    #[inline]
    pub fn pfmerge<K: AsRef<[u8]>>(&self, dest: K, sources: &[K]) -> Result<()> {
        self.pfmerge_ns("default", dest, sources)
    }

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