1use std::collections::HashMap;
12use std::hash::{Hash, Hasher};
13
14pub fn hash64<T: Hash + ?Sized>(v: &T) -> u64 {
16 let mut h = std::hash::DefaultHasher::new();
17 v.hash(&mut h);
18 h.finish()
19}
20
21#[derive(Debug, Clone)]
23pub struct HyperLogLog {
24 registers: Vec<u8>,
25}
26
27const HLL_P: u32 = 12;
28const HLL_M: usize = 1 << HLL_P;
29
30impl Default for HyperLogLog {
31 fn default() -> Self {
32 Self::new()
33 }
34}
35
36impl HyperLogLog {
37 pub fn new() -> Self {
38 Self {
39 registers: vec![0; HLL_M],
40 }
41 }
42
43 pub fn insert_hash(&mut self, h: u64) {
45 let idx = (h >> (64 - HLL_P)) as usize;
46 let rest = h << HLL_P;
47 let rank = if rest == 0 {
50 (64 - HLL_P) as u8 + 1
51 } else {
52 rest.leading_zeros() as u8 + 1
53 };
54 if rank > self.registers[idx] {
55 self.registers[idx] = rank;
56 }
57 }
58
59 pub fn insert<T: Hash + ?Sized>(&mut self, v: &T) {
61 self.insert_hash(hash64(v));
62 }
63
64 pub fn merge(&mut self, other: &HyperLogLog) {
66 for (a, b) in self.registers.iter_mut().zip(&other.registers) {
67 if *b > *a {
68 *a = *b;
69 }
70 }
71 }
72
73 pub fn estimate(&self) -> u64 {
75 let m = HLL_M as f64;
76 let alpha = 0.7213 / (1.0 + 1.079 / m);
77 let sum: f64 = self.registers.iter().map(|&r| 2f64.powi(-(r as i32))).sum();
78 let raw = alpha * m * m / sum;
79 let zeros = self.registers.iter().filter(|&&r| r == 0).count();
80 let est = if raw <= 2.5 * m && zeros > 0 {
81 m * (m / zeros as f64).ln()
83 } else {
84 raw
85 };
86 est.round().max(0.0) as u64
87 }
88}
89
90#[derive(Debug, Clone)]
93pub struct TopK {
94 capacity: usize,
95 counters: HashMap<String, (u64, u64)>,
96}
97
98impl TopK {
99 pub fn new(k: usize) -> Self {
102 Self {
103 capacity: (k * 4).max(1),
104 counters: HashMap::new(),
105 }
106 }
107
108 pub fn insert(&mut self, value: &str) {
109 if let Some(c) = self.counters.get_mut(value) {
110 c.0 += 1;
111 return;
112 }
113 if self.counters.len() < self.capacity {
114 self.counters.insert(value.to_string(), (1, 0));
115 return;
116 }
117 let (victim, min) = self
119 .counters
120 .iter()
121 .map(|(k, (c, _))| (k.clone(), *c))
122 .min_by(|a, b| a.1.cmp(&b.1).then_with(|| a.0.cmp(&b.0)))
123 .expect("capacity >= 1");
124 self.counters.remove(&victim);
125 self.counters.insert(value.to_string(), (min + 1, min));
126 }
127
128 pub fn top(&self, k: usize) -> Vec<(String, u64, u64)> {
131 let mut all: Vec<(String, u64, u64)> = self
132 .counters
133 .iter()
134 .map(|(v, (c, e))| (v.clone(), *c, *e))
135 .collect();
136 all.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
137 all.truncate(k);
138 all
139 }
140}
141
142#[derive(Debug, Clone)]
144pub struct Reservoir {
145 capacity: usize,
146 seen: u64,
147 values: Vec<f64>,
148 rng: u64,
149}
150
151impl Reservoir {
152 pub fn new(capacity: usize) -> Self {
153 Self {
154 capacity: capacity.max(1),
155 seen: 0,
156 values: Vec::new(),
157 rng: 0x9E37_79B9_7F4A_7C15,
158 }
159 }
160
161 fn next_u64(&mut self) -> u64 {
162 let mut x = self.rng;
164 x ^= x >> 12;
165 x ^= x << 25;
166 x ^= x >> 27;
167 self.rng = x;
168 x.wrapping_mul(0x2545_F491_4F6C_DD1D)
169 }
170
171 pub fn insert(&mut self, v: f64) {
172 self.seen += 1;
173 if self.values.len() < self.capacity {
174 self.values.push(v);
175 return;
176 }
177 let j = self.next_u64() % self.seen;
178 if (j as usize) < self.capacity {
179 self.values[j as usize] = v;
180 }
181 }
182
183 pub fn quantile(&self, q: f64) -> Option<f64> {
185 if self.values.is_empty() {
186 return None;
187 }
188 let mut sorted = self.values.clone();
189 sorted.sort_by(f64::total_cmp);
190 Some(crate::anomaly::quantile(&sorted, q))
191 }
192
193 pub fn len(&self) -> usize {
194 self.values.len()
195 }
196
197 pub fn is_empty(&self) -> bool {
198 self.values.is_empty()
199 }
200}
201
202#[derive(Debug, Clone, Default)]
204pub struct Welford {
205 pub count: u64,
206 mean: f64,
207 m2: f64,
208 pub min: f64,
209 pub max: f64,
210}
211
212impl Welford {
213 pub fn insert(&mut self, x: f64) {
214 if self.count == 0 {
215 self.min = x;
216 self.max = x;
217 } else {
218 self.min = self.min.min(x);
219 self.max = self.max.max(x);
220 }
221 self.count += 1;
222 let delta = x - self.mean;
223 self.mean += delta / self.count as f64;
224 self.m2 += delta * (x - self.mean);
225 }
226
227 pub fn mean(&self) -> Option<f64> {
228 (self.count > 0).then_some(self.mean)
229 }
230
231 pub fn stddev(&self) -> Option<f64> {
233 (self.count > 0).then(|| (self.m2 / self.count as f64).sqrt())
234 }
235}
236
237#[cfg(test)]
238mod tests {
239 use super::*;
240
241 #[test]
242 fn hll_estimates_within_tolerance() {
243 for &n in &[1u64, 10, 100, 1_000, 50_000] {
244 let mut h = HyperLogLog::new();
245 for i in 0..n {
246 h.insert(&format!("value-{i}"));
247 }
248 let est = h.estimate() as f64;
249 let err = (est - n as f64).abs() / n as f64;
250 assert!(err < 0.05, "n {n} est {est} err {err}");
251 }
252 assert_eq!(HyperLogLog::new().estimate(), 0);
253 }
254
255 #[test]
256 fn hll_merge_is_register_max_and_duplicates_do_not_count() {
257 let mut a = HyperLogLog::new();
258 let mut b = HyperLogLog::new();
259 for i in 0..1000 {
260 a.insert(&i);
261 a.insert(&i);
262 b.insert(&(i + 500));
263 }
264 a.merge(&b);
265 let est = a.estimate() as f64;
266 assert!((est - 1500.0).abs() / 1500.0 < 0.05, "{est}");
267 let mut z = HyperLogLog::new();
268 z.insert_hash(0);
269 assert_eq!(z.registers.iter().filter(|&&r| r > 0).count(), 1);
270 }
271
272 #[test]
273 fn topk_keeps_the_heavy_hitters_with_error_bounds() {
274 let mut t = TopK::new(2);
275 for _ in 0..100 {
276 t.insert("a");
277 }
278 for _ in 0..50 {
279 t.insert("b");
280 }
281 for i in 0..40 {
282 t.insert(&format!("noise-{i}"));
283 }
284 let top = t.top(2);
285 assert_eq!(top[0].0, "a");
286 assert_eq!(top[0].1, 100);
287 assert_eq!(top[1].0, "b");
288 assert!(top[1].1 >= 50);
289 assert!(t.counters.len() <= 8);
290 assert!(t.counters.values().any(|(_, e)| *e > 0));
292 }
293
294 #[test]
295 fn topk_orders_ties_by_value() {
296 let mut t = TopK::new(3);
297 for v in ["z", "y", "x"] {
298 t.insert(v);
299 }
300 let names: Vec<_> = t.top(3).into_iter().map(|(v, _, _)| v).collect();
301 assert_eq!(names, vec!["x", "y", "z"]);
302 assert!(TopK::new(0).top(1).is_empty());
303 }
304
305 #[test]
306 fn reservoir_is_uniform_enough_and_deterministic() {
307 let mut r = Reservoir::new(200);
308 for i in 0..10_000 {
309 r.insert(i as f64);
310 }
311 assert_eq!(r.len(), 200);
312 let median = r.quantile(0.5).unwrap();
313 assert!((3500.0..=6500.0).contains(&median), "{median}");
314 let mut r2 = Reservoir::new(200);
315 for i in 0..10_000 {
316 r2.insert(i as f64);
317 }
318 assert_eq!(r.quantile(0.9), r2.quantile(0.9));
319 assert!(Reservoir::new(4).quantile(0.5).is_none());
320 assert!(Reservoir::new(4).is_empty());
321 }
322
323 #[test]
324 fn welford_matches_closed_form() {
325 let mut w = Welford::default();
326 assert!(w.mean().is_none());
327 assert!(w.stddev().is_none());
328 for x in [2.0, 4.0, 4.0, 4.0, 5.0, 5.0, 7.0, 9.0] {
329 w.insert(x);
330 }
331 assert_eq!(w.count, 8);
332 assert_eq!(w.mean(), Some(5.0));
333 assert!((w.stddev().unwrap() - 2.0).abs() < 1e-12);
334 assert_eq!(w.min, 2.0);
335 assert_eq!(w.max, 9.0);
336 }
337}