use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use dashmap::DashMap;
use rustc_hash::FxBuildHasher;
use smallvec::SmallVec;
use crate::cluster::{Cluster, ClusterInner};
use crate::mask::apply_masks;
use crate::options::Options;
use crate::similarity::similarity;
use crate::snapshot::{decode, encode, ClusterSnapshot, SnapshotV1, TokenSnapshot};
use crate::tokenize::{is_numeric_token, split_first_line, tokenize_with, Token};
use crate::{ClusterId, OwnedToken};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum UpdateType {
Created,
TemplateChanged,
None,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct AddResult {
pub cluster_id: ClusterId,
pub update: UpdateType,
}
struct Shard {
root: crate::tree::TreeNode,
}
pub struct Miner {
shards: DashMap<usize, Arc<RwLock<Shard>>, FxBuildHasher>,
clusters_by_id: DashMap<ClusterId, Arc<RwLock<ClusterInner>>, FxBuildHasher>,
options: Arc<Options>,
counter: AtomicU64,
tick: AtomicU64,
}
impl Miner {
pub fn from_options(options: Options) -> Self {
Miner {
shards: DashMap::with_hasher(FxBuildHasher),
clusters_by_id: DashMap::with_hasher(FxBuildHasher),
options: Arc::new(options),
counter: AtomicU64::new(0),
tick: AtomicU64::new(0),
}
}
fn next_tick(&self) -> u64 {
self.tick.fetch_add(1, Ordering::Relaxed)
}
pub fn builder() -> crate::MinerBuilder {
crate::MinerBuilder::new()
}
pub fn len(&self) -> usize {
self.clusters_by_id.len()
}
pub fn is_empty(&self) -> bool {
self.clusters_by_id.is_empty()
}
fn descent_keys<'a>(&'a self, tokens: &'a [Token<'a>]) -> SmallVec<[&'a str; 8]> {
let n = self.options.prefix_len().min(tokens.len());
let mut keys: SmallVec<[&'a str; 8]> = SmallVec::new();
for tok in &tokens[..n] {
if self.options.parametrize_numeric_tokens && is_numeric_token(tok.text) {
keys.push(&self.options.wildcard);
} else {
keys.push(tok.text);
}
}
keys
}
fn shard_for(&self, count: usize) -> Arc<RwLock<Shard>> {
if let Some(s) = self.shards.get(&count) {
return s.clone();
}
self.shards
.entry(count)
.or_insert_with(|| {
Arc::new(RwLock::new(Shard {
root: crate::tree::TreeNode::new_internal(),
}))
})
.clone()
}
pub fn add(&self, line: &str) -> AddResult {
self.add_inner(line, None)
}
pub fn add_with_member(&self, line: &str, member: &str) -> AddResult {
self.add_inner(line, Some(member))
}
fn add_inner(&self, line: &str, member: Option<&str>) -> AddResult {
let masked = apply_masks(line, &self.options.masks);
let (first, suffix) = if self.options.first_line_only {
split_first_line(&masked)
} else {
(masked.as_ref(), None)
};
let tokens = tokenize_with(first, self.options.active_path_delimiters());
let count = tokens.len();
let keys = self.descent_keys(&tokens);
let shard = self.shard_for(count);
let matched = {
let guard = shard.read().expect("shard lock poisoned");
guard
.root
.descend(&keys)
.and_then(|leaf| self.best_match(leaf, &tokens))
.filter(|(_, sim, _)| *sim >= self.options.sim_threshold)
};
if let Some((id, _, arc)) = matched {
return self.apply_match(&arc, id, &tokens, member);
}
let mut guard = shard.write().expect("shard lock poisoned");
let leaf = guard
.root
.descend_or_create(&keys, self.options.max_clusters_per_leaf);
if let Some((id, sim, arc)) = self.best_match(leaf, &tokens) {
if sim >= self.options.sim_threshold {
drop(guard);
return self.apply_match(&arc, id, &tokens, member);
}
}
let id = self.counter.fetch_add(1, Ordering::Relaxed) + 1;
let owned: Vec<OwnedToken> = tokens.iter().map(OwnedToken::from).collect();
let mut inner = ClusterInner::new(id, owned, SystemTime::now(), suffix.map(Arc::from));
inner.last_used.store(self.next_tick(), Ordering::Relaxed); if let Some(m) = member {
inner.add_member(m);
}
self.clusters_by_id.insert(id, Arc::new(RwLock::new(inner)));
self.evict_if_full(leaf);
leaf.insert(id);
AddResult {
cluster_id: id,
update: UpdateType::Created,
}
}
fn best_match(
&self,
leaf: &crate::tree::LeafBucket,
tokens: &[Token<'_>],
) -> Option<(ClusterId, f64, Arc<RwLock<ClusterInner>>)> {
let mut best: Option<(ClusterId, f64, Arc<RwLock<ClusterInner>>)> = None;
for &id in leaf.ids() {
if let Some(entry) = self.clusters_by_id.get(&id) {
let sim = {
let body = entry.read().expect("cluster lock poisoned");
similarity(&body.tokens, tokens, &self.options.wildcard)
};
if best.as_ref().map_or(true, |(_, b, _)| sim > *b) {
best = Some((id, sim, entry.value().clone()));
}
}
}
best
}
fn apply_match(
&self,
arc: &Arc<RwLock<ClusterInner>>,
id: ClusterId,
tokens: &[Token<'_>],
member: Option<&str>,
) -> AddResult {
let needs_generalize = {
let body = arc.read().expect("cluster lock poisoned");
body.touch(self.next_tick());
body.would_generalize(tokens, &self.options.wildcard)
};
if !needs_generalize && member.is_none() {
return AddResult {
cluster_id: id,
update: UpdateType::None,
};
}
let mut body = arc.write().expect("cluster lock poisoned");
let changed = needs_generalize && body.generalize(tokens, &self.options.wildcard);
if let Some(m) = member {
body.add_member(m);
}
AddResult {
cluster_id: id,
update: if changed {
UpdateType::TemplateChanged
} else {
UpdateType::None
},
}
}
fn evict_if_full(&self, leaf: &mut crate::tree::LeafBucket) {
if !leaf.is_full() {
return;
}
let victim = leaf.ids().iter().copied().min_by_key(|id| {
self.clusters_by_id
.get(id)
.map(|a| a.read().expect("cluster lock poisoned").recency())
.unwrap_or(0)
});
if let Some(v) = victim {
leaf.remove(v);
self.clusters_by_id.remove(&v);
}
}
fn tokens_for_query<'a>(&self, masked: &'a str) -> crate::tokenize::Tokens<'a> {
let first = if self.options.first_line_only {
split_first_line(masked).0
} else {
masked
};
tokenize_with(first, self.options.active_path_delimiters())
}
pub fn match_only(&self, line: &str) -> Option<ClusterId> {
let masked = apply_masks(line, &self.options.masks);
let tokens = self.tokens_for_query(&masked);
let count = tokens.len();
let keys = self.descent_keys(&tokens);
let shard = self.shards.get(&count)?.clone();
let guard = shard.read().expect("shard lock poisoned");
let leaf = guard.root.descend(&keys)?;
self.best_match(leaf, &tokens)
.filter(|(_, sim, _)| *sim >= self.options.sim_threshold)
.map(|(id, _, _)| id)
}
pub fn extract(&self, line: &str) -> Option<(ClusterId, Vec<String>)> {
let id = self.match_only(line)?;
let masked = apply_masks(line, &self.options.masks);
let tokens = self.tokens_for_query(&masked);
let arc = self.clusters_by_id.get(&id)?.clone();
let body = arc.read().expect("cluster lock poisoned");
let mut params = Vec::new();
for (stored, tok) in body.tokens.iter().zip(tokens.iter()) {
if stored.text == self.options.wildcard {
params.push(tok.text.to_string());
}
}
Some((id, params))
}
pub fn clusters(&self) -> Vec<Cluster> {
self.clusters_by_id
.iter()
.map(|e| {
e.value()
.read()
.expect("cluster lock poisoned")
.to_public(&self.options.wildcard)
})
.collect()
}
pub fn cluster(&self, id: ClusterId) -> Option<Cluster> {
let arc = self.clusters_by_id.get(&id)?.clone();
let body = arc.read().expect("cluster lock poisoned");
Some(body.to_public(&self.options.wildcard))
}
fn insert_existing(&self, inner: ClusterInner) {
let count = inner.tokens.len();
let id = inner.id;
let shard = self.shard_for(count);
let mut guard = shard.write().expect("shard lock poisoned");
{
let view: SmallVec<[Token<'_>; 16]> = inner
.tokens
.iter()
.map(|t| Token {
text: &t.text,
leading_delim: t.leading_delim,
trailing_delim: t.trailing_delim,
})
.collect();
let keys = self.descent_keys(&view);
let leaf = guard
.root
.descend_or_create(&keys, self.options.max_clusters_per_leaf);
self.evict_if_full(leaf);
leaf.insert(id);
}
self.clusters_by_id.insert(id, Arc::new(RwLock::new(inner)));
}
pub fn snapshot(&self) -> Vec<u8> {
let clusters = self
.clusters_by_id
.iter()
.map(|e| {
let b = e.value().read().expect("cluster lock poisoned");
ClusterSnapshot {
id: b.id,
tokens: b
.tokens
.iter()
.map(|t| TokenSnapshot {
text: t.text.to_string(),
leading_delim: t.leading_delim,
trailing_delim: t.trailing_delim,
})
.collect(),
size: b.size.load(Ordering::Relaxed),
created_at_ms: system_time_to_ms(b.created_at),
updated_at_ms: b.updated_at_ms.load(Ordering::Relaxed),
suffix: b.suffix.as_ref().map(|s| s.to_string()),
members: b.members.iter().map(|m| m.to_string()).collect(),
}
})
.collect();
let body = SnapshotV1 {
options: (*self.options).clone(),
counter: self.counter.load(Ordering::Relaxed),
clusters,
};
encode(&body)
}
pub fn restore(&self, bytes: &[u8]) -> Result<(), crate::LogdrainError> {
let body = decode(bytes)?;
self.shards.clear();
self.clusters_by_id.clear();
self.counter.store(body.counter, Ordering::Relaxed);
for cs in body.clusters {
let tokens: Vec<OwnedToken> = cs
.tokens
.into_iter()
.map(|t| OwnedToken {
text: Arc::from(t.text.as_str()),
leading_delim: t.leading_delim,
trailing_delim: t.trailing_delim,
})
.collect();
let inner = ClusterInner {
id: cs.id,
tokens,
size: AtomicU64::new(cs.size),
created_at: ms_to_system_time(cs.created_at_ms),
updated_at_ms: AtomicU64::new(cs.updated_at_ms),
last_used: AtomicU64::new(0),
suffix: cs.suffix.map(Arc::from),
members: cs.members.into_iter().map(Arc::from).collect(),
};
self.insert_existing(inner);
}
Ok(())
}
pub fn save_state(&self, p: &dyn crate::Persistence) -> Result<(), crate::LogdrainError> {
p.save(&self.snapshot())?;
Ok(())
}
pub fn load_state(&self, p: &dyn crate::Persistence) -> Result<bool, crate::LogdrainError> {
match p.load()? {
Some(bytes) => {
self.restore(&bytes)?;
Ok(true)
}
None => Ok(false),
}
}
}
fn system_time_to_ms(t: SystemTime) -> u64 {
t.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
fn ms_to_system_time(ms: u64) -> SystemTime {
UNIX_EPOCH + Duration::from_millis(ms)
}
impl crate::MinerBuilder {
pub fn build(self) -> Result<Miner, crate::LogdrainError> {
Ok(Miner::from_options(self.build_options()?))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::MinerBuilder;
fn miner() -> Miner {
Miner::from_options(MinerBuilder::new().build_options().unwrap())
}
#[test]
fn first_line_creates_cluster() {
let m = miner();
let r = m.add("user 42 logged in");
assert_eq!(r.update, UpdateType::Created);
assert_eq!(r.cluster_id, 1);
assert_eq!(m.len(), 1);
}
#[test]
fn similar_line_joins_and_generalizes() {
let m = miner();
let a = m.add("user 42 logged in");
let b = m.add("user 99 logged in");
assert_eq!(b.cluster_id, a.cluster_id);
assert_eq!(b.update, UpdateType::TemplateChanged);
assert_eq!(m.len(), 1);
}
#[test]
fn identical_line_is_update_none() {
let m = miner();
m.add("a b c d e");
let r = m.add("a b c d e");
assert_eq!(r.update, UpdateType::None);
}
#[test]
fn different_token_count_makes_new_cluster() {
let m = miner();
let a = m.add("a b c");
let b = m.add("a b c d");
assert_ne!(a.cluster_id, b.cluster_id);
assert_eq!(m.len(), 2);
}
#[test]
fn dissimilar_same_length_makes_new_cluster() {
let m = miner();
let a = m.add("alpha one two three four");
let b = m.add("alpha NINE TEN ELEVEN TWELVE");
assert_ne!(a.cluster_id, b.cluster_id);
assert_eq!(m.len(), 2);
}
#[test]
fn numeric_parametrization_groups_by_wildcard_prefix() {
let m = miner();
let a = m.add("100 ms elapsed for request");
let b = m.add("200 ms elapsed for request");
assert_eq!(a.cluster_id, b.cluster_id);
}
#[test]
fn ids_are_monotonic() {
let m = miner();
let a = m.add("a b c");
let b = m.add("x y z");
assert_eq!(a.cluster_id, 1);
assert_eq!(b.cluster_id, 2);
}
#[test]
fn builder_build_constructs_miner() {
let m = MinerBuilder::new().sim_threshold(0.5).build().unwrap();
assert_eq!(m.len(), 0);
assert!(MinerBuilder::new().depth(1).build().is_err());
}
#[test]
fn match_only_finds_without_learning() {
let m = miner();
let a = m.add("user 42 logged in");
let before = m.len();
let hit = m.match_only("user 7 logged in");
assert_eq!(hit, Some(a.cluster_id));
assert_eq!(m.len(), before); assert_eq!(m.cluster(a.cluster_id).unwrap().size(), 1);
}
#[test]
fn match_only_misses_return_none() {
let m = miner();
m.add("a b c");
assert_eq!(m.match_only("x y z w"), None); assert_eq!(m.match_only("p q r"), None); }
#[test]
fn extract_returns_wildcard_values() {
let m = miner();
m.add("user 42 logged in");
m.add("user 99 logged in"); let (id, params) = m.extract("user 7 logged in").unwrap();
assert_eq!(id, m.match_only("user 7 logged in").unwrap());
assert_eq!(params, vec!["7".to_string()]);
}
#[test]
fn clusters_and_cluster_snapshots() {
let m = miner();
let a = m.add("a b c");
m.add("x y z");
let all = m.clusters();
assert_eq!(all.len(), 2);
let one = m.cluster(a.cluster_id).unwrap();
assert_eq!(one.id(), a.cluster_id);
assert!(m.cluster(99999).is_none());
}
fn miner_with(b: MinerBuilder) -> Miner {
Miner::from_options(b.build_options().unwrap())
}
#[test]
fn path_clustering_preserves_structure() {
let m = miner_with(MinerBuilder::new().path_delimiters(&['/']));
let a = m.add("PUT /servers/409/foo/10.0.0.1");
let b = m.add("PUT /servers/410/foo/10.0.0.2");
assert_eq!(a.cluster_id, b.cluster_id);
assert_eq!(
m.cluster(a.cluster_id).unwrap().template(),
"PUT /servers/<*>/foo/<*>"
);
}
#[test]
fn masks_cluster_high_cardinality_tokens() {
let m = miner_with(MinerBuilder::new().masks([crate::builtin_masks::uuid()]));
let a = m.add("request 550e8400-e29b-41d4-a716-446655440000 ok");
let b = m.add("request 6ba7b810-9dad-11d1-80b4-00c04fd430c8 ok");
assert_eq!(a.cluster_id, b.cluster_id);
assert_eq!(b.update, UpdateType::None); assert_eq!(
m.cluster(a.cluster_id).unwrap().template(),
"request <uuid> ok"
);
}
#[test]
fn first_line_only_captures_suffix() {
let m = miner_with(MinerBuilder::new().first_line_only(true));
let a = m.add("NullPointerException at Foo\n at bar()\n at baz()");
assert_eq!(
m.cluster(a.cluster_id).unwrap().suffix(),
Some(" at bar()\n at baz()")
);
let b = m.add("NullPointerException at Foo\n at other()");
assert_eq!(a.cluster_id, b.cluster_id);
assert_eq!(
m.cluster(a.cluster_id).unwrap().suffix(),
Some(" at bar()\n at baz()")
);
}
#[test]
fn add_with_member_records_deduped_members() {
let m = miner();
let a = m.add_with_member("user 1 logged in", "svc-a");
m.add_with_member("user 2 logged in", "svc-b");
m.add_with_member("user 3 logged in", "svc-a"); let c = m.cluster(a.cluster_id).unwrap();
let members: Vec<&str> = c.members().iter().map(|m| &**m).collect();
assert_eq!(members, vec!["svc-a", "svc-b"]);
let d = m.add("totally different shape here now");
assert!(m.cluster(d.cluster_id).unwrap().members().is_empty());
}
#[test]
fn extract_honors_masks_and_path() {
let m = miner_with(MinerBuilder::new().path_delimiters(&['/']));
m.add("GET /u/409/x");
m.add("GET /u/410/x"); let (_, params) = m.extract("GET /u/777/x").unwrap();
assert_eq!(params, vec!["777".to_string()]);
}
#[test]
fn full_leaf_evicts_least_recently_used() {
let m = miner_with(MinerBuilder::new().max_clusters_per_leaf(2));
m.add("p q a b c d"); m.add("p q e f g h"); assert_eq!(m.len(), 2);
m.add("p q a b c d"); m.add("p q i j k l");
assert_eq!(m.len(), 2, "per-leaf cap bounds cluster count via eviction");
assert!(
m.match_only("p q a b c d").is_some(),
"recently-used A retained"
);
assert!(m.match_only("p q i j k l").is_some(), "new C retained");
assert!(
m.match_only("p q e f g h").is_none(),
"least-recently-used B evicted"
);
}
}