Add YAML config loading with validation, env overrides for secrets, and a web dashboard exposing health, stats, and recent events.
130 lines
2.6 KiB
V
130 lines
2.6 KiB
V
module main
|
|
|
|
import time
|
|
import sources
|
|
|
|
pub struct DedupEntry {
|
|
pub mut:
|
|
last_seen time.Time
|
|
count int
|
|
}
|
|
|
|
pub struct Dedup {
|
|
pub mut:
|
|
enabled bool
|
|
ttl_seconds int
|
|
max_per_minute int
|
|
max_per_source int
|
|
cache map[string]DedupEntry
|
|
topic_history map[string][]time.Time
|
|
source_history map[string][]time.Time
|
|
}
|
|
|
|
pub fn new_dedup(cfg DedupConfig) Dedup {
|
|
return Dedup{
|
|
enabled: cfg.enabled
|
|
ttl_seconds: cfg.ttl_seconds
|
|
max_per_minute: cfg.rate_limit.max_per_minute
|
|
max_per_source: cfg.rate_limit.max_per_source
|
|
cache: map[string]DedupEntry{}
|
|
topic_history: map[string][]time.Time{}
|
|
source_history: map[string][]time.Time{}
|
|
}
|
|
}
|
|
|
|
fn simple_hash(s string) u64 {
|
|
mut h := u64(14695981039346656037)
|
|
for c in s.bytes() {
|
|
h ^= u64(c)
|
|
h *= u64(1099511628211)
|
|
}
|
|
return h
|
|
}
|
|
|
|
fn dedup_key(event sources.Event) string {
|
|
h := simple_hash(event.message)
|
|
return '${event.source}:${event.name}:${event.topic}:${h}'
|
|
}
|
|
|
|
fn (mut d Dedup) cleanup(now time.Time) {
|
|
now_unix := now.unix()
|
|
for key, entry in d.cache {
|
|
if now_unix - entry.last_seen.unix() > d.ttl_seconds {
|
|
d.cache.delete(key)
|
|
}
|
|
}
|
|
for topic, times in d.topic_history {
|
|
mut kept := []time.Time{}
|
|
for t in times {
|
|
if now_unix - t.unix() <= 60 {
|
|
kept << t
|
|
}
|
|
}
|
|
if kept.len == 0 {
|
|
d.topic_history.delete(topic)
|
|
} else {
|
|
d.topic_history[topic] = kept
|
|
}
|
|
}
|
|
for source, times in d.source_history {
|
|
mut kept := []time.Time{}
|
|
for t in times {
|
|
if now_unix - t.unix() <= 60 {
|
|
kept << t
|
|
}
|
|
}
|
|
if kept.len == 0 {
|
|
d.source_history.delete(source)
|
|
} else {
|
|
d.source_history[source] = kept
|
|
}
|
|
}
|
|
}
|
|
|
|
pub fn (mut d Dedup) allow(event sources.Event) bool {
|
|
if !d.enabled {
|
|
return true
|
|
}
|
|
now := time.now()
|
|
d.cleanup(now)
|
|
|
|
// Rate limit per topic
|
|
topic_times := d.topic_history[event.topic] or { []time.Time{} }
|
|
if topic_times.len >= d.max_per_minute && d.max_per_minute > 0 {
|
|
return false
|
|
}
|
|
|
|
// Rate limit per source
|
|
source_key := '${event.source}:${event.name}'
|
|
source_times := d.source_history[source_key] or { []time.Time{} }
|
|
if source_times.len >= d.max_per_source && d.max_per_source > 0 {
|
|
return false
|
|
}
|
|
|
|
// Dedup check
|
|
key := dedup_key(event)
|
|
entry := d.cache[key] or { DedupEntry{} }
|
|
if entry.count > 0 {
|
|
d.cache[key] = DedupEntry{
|
|
last_seen: now
|
|
count: entry.count + 1
|
|
}
|
|
return false
|
|
}
|
|
|
|
d.cache[key] = DedupEntry{
|
|
last_seen: now
|
|
count: 1
|
|
}
|
|
|
|
mut new_topic_times := topic_times.clone()
|
|
new_topic_times << now
|
|
d.topic_history[event.topic] = new_topic_times
|
|
|
|
mut new_source_times := source_times.clone()
|
|
new_source_times << now
|
|
d.source_history[source_key] = new_source_times
|
|
|
|
return true
|
|
}
|