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_sigma, hll_tau, 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_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 rapidhash::RapidHashSet;
use std::str;
use crate::db::WeDb;
use crate::error::Result;
use crate::key_composer::KeyComposer;
impl WeDb {
pub fn pfadd_ns<K: AsRef<[u8]>, E: AsRef<[u8]>>(
&self,
ns: &str,
key: K,
elements: &[E],
) -> Result<bool> {
let kc = KeyComposer::new(ns);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hll_meta(k_str);
let data_k = kc.hll_key(k_str);
let now_ms = ts_::sec() * 1000;
let (mut meta, exists_valid) = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => {
let m = HyperLogLogMeta::decode(&m_bytes);
match m {
Some(meta) if !meta.is_expired(now_ms) => (meta, true),
_ => (HyperLogLogMeta::new_with_version(0), false),
}
}
None => (HyperLogLogMeta::new_with_version(0), false),
};
let mut registers = if exists_valid {
self.data_ks
.get(data_k.as_bytes())?
.map(|v| v.to_vec())
.unwrap_or_else(|| vec![0u8; HLL_DENSE_SIZE])
} else {
vec![0u8; HLL_DENSE_SIZE]
};
if 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);
let cur = hll_dense_get_register(®isters, reg_idx);
if count > cur {
hll_dense_set_register(&mut registers, reg_idx, count);
updated = true;
}
}
if updated || (elements.is_empty() && !exists_valid) {
meta.base.size = HLL_DENSE_SIZE as u64;
meta.encode_type = HllEncodeType::Dense;
let mut batch = self.db.batch();
batch.insert(&self.data_ks, data_k.as_bytes(), ®isters);
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
}
Ok(updated)
}
#[inline]
pub fn pfadd<K: AsRef<[u8]>, E: AsRef<[u8]>>(&self, key: K, elements: &[E]) -> Result<bool> {
self.pfadd_ns("default", key, elements)
}
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 = ts_::sec() * 1000;
if keys.len() == 1 {
let k_str = str::from_utf8(keys[0].as_ref()).unwrap_or("");
let meta_k = kc.hll_meta(k_str);
let meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => match HyperLogLogMeta::decode(&m_bytes) {
Some(meta) if !meta.is_expired(now_ms) => meta,
_ => return Ok(0),
},
None => return Ok(0),
};
let data_k = kc.hll_key(k_str);
let registers = match self.data_ks.get(data_k.as_bytes())? {
Some(v) => v,
None => return Ok(0),
};
return match meta.encode_type {
HllEncodeType::Dense => Ok(hll_dense_estimate(®isters)),
HllEncodeType::Sparse => Ok(hll_sparse_estimate(®isters)
.unwrap_or_else(|_| hll_dense_estimate(®isters))),
};
}
let mut seen = RapidHashSet::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 k_str = str::from_utf8(k_bytes).unwrap_or("");
let meta_k = kc.hll_meta(k_str);
let meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => match HyperLogLogMeta::decode(&m_bytes) {
Some(meta) if !meta.is_expired(now_ms) => meta,
_ => continue,
},
None => continue,
};
let data_k = kc.hll_key(k_str);
if let Some(reg) = self.data_ks.get(data_k.as_bytes())? {
match meta.encode_type {
HllEncodeType::Dense => {
hll_merge_bytes(&mut merged, ®);
}
HllEncodeType::Sparse => {
let mut dense_buf = vec![0u8; HLL_DENSE_SIZE];
if hll_sparse_to_dense(®, &mut dense_buf).is_ok() {
hll_merge_bytes(&mut merged, &dense_buf);
} else {
hll_merge_bytes(&mut merged, ®);
}
}
}
has_any = true;
}
}
if !has_any {
return Ok(0);
}
Ok(hll_dense_estimate(&merged))
}
#[inline]
pub fn pfcount<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<u64> {
self.pfcount_ns("default", keys)
}
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_str = str::from_utf8(dest.as_ref()).unwrap_or("");
let dest_meta_k = kc.hll_meta(dest_str);
let dest_data_k = kc.hll_key(dest_str);
let now_ms = ts_::sec() * 1000;
let (mut dest_meta, dest_valid) = match self.meta_ks.get(dest_meta_k.as_bytes())? {
Some(m_bytes) => {
let m = HyperLogLogMeta::decode(&m_bytes);
match m {
Some(meta) if !meta.is_expired(now_ms) => (meta, true),
_ => (HyperLogLogMeta::new_with_version(0), false),
}
}
None => (HyperLogLogMeta::new_with_version(0), false),
};
let mut merged = if dest_valid {
let data = self
.data_ks
.get(dest_data_k.as_bytes())?
.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 = RapidHashSet::default();
seen.insert(dest.as_ref());
for k in sources {
let k_bytes = k.as_ref();
if !seen.insert(k_bytes) {
continue;
}
let k_str = str::from_utf8(k_bytes).unwrap_or("");
let meta_k = kc.hll_meta(k_str);
let meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => match HyperLogLogMeta::decode(&m_bytes) {
Some(meta) if !meta.is_expired(now_ms) => meta,
_ => continue,
},
None => continue,
};
let data_k = kc.hll_key(k_str);
if let Some(reg) = self.data_ks.get(data_k.as_bytes())? {
match meta.encode_type {
HllEncodeType::Dense => {
hll_merge_bytes(&mut merged, ®);
}
HllEncodeType::Sparse => {
let mut dense_buf = vec![0u8; HLL_DENSE_SIZE];
if hll_sparse_to_dense(®, &mut dense_buf).is_ok() {
hll_merge_bytes(&mut merged, &dense_buf);
} else {
hll_merge_bytes(&mut merged, ®);
}
}
}
}
}
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_ks, dest_data_k.as_bytes(), &merged);
batch.insert(&self.meta_ks, dest_meta_k.as_bytes(), dest_meta.encode());
batch.commit()?;
Ok(())
}
#[inline]
pub fn pfmerge<K: AsRef<[u8]>>(&self, dest: K, sources: &[K]) -> Result<()> {
self.pfmerge_ns("default", dest, sources)
}
#[inline]
pub fn pfselftest(&self) -> bool {
HyperLogLog::selftest()
}
}