use std::collections::{HashMap, VecDeque};
use serde::Serialize;
use tokio::sync::RwLock;
use crate::detect::Finding;
type FoldKey<'a> = (Option<(&'a str, &'a str)>, &'a str);
#[non_exhaustive]
#[derive(Debug, Clone, Serialize, serde::Deserialize)]
pub struct StoredFinding {
pub finding: Finding,
pub stored_at_ms: u64,
#[serde(default)]
pub first_seen_ms: u64,
#[serde(default = "default_seen_count")]
pub seen_count: u64,
}
fn default_seen_count() -> u64 {
1
}
#[must_use]
pub fn coalesce_by_signature(entries: &[StoredFinding]) -> Vec<StoredFinding> {
fold_entries(entries.iter())
}
fn fold_entries<'a>(entries: impl Iterator<Item = &'a StoredFinding>) -> Vec<StoredFinding> {
let mut out: Vec<StoredFinding> = Vec::new();
let mut index: HashMap<FoldKey<'_>, usize> = HashMap::new();
for entry in entries {
if entry.finding.signature.is_empty() {
out.push(entry.clone());
continue;
}
let grouping = entry.finding.grouping_identity();
let key = (grouping, entry.finding.signature.as_str());
if let Some(&i) = index.get(&key) {
let kept: &mut StoredFinding = &mut out[i];
let seen = kept.seen_count + entry.seen_count;
let first = kept.first_seen_ms.min(entry.first_seen_ms);
let last = kept.stored_at_ms.max(entry.stored_at_ms);
if entry.finding.severity < kept.finding.severity {
kept.finding = entry.finding.clone();
}
kept.seen_count = seen;
kept.first_seen_ms = first;
kept.stored_at_ms = last;
} else {
index.insert(key, out.len());
out.push(entry.clone());
}
}
out
}
#[non_exhaustive]
#[derive(Debug, Default)]
pub struct FindingsFilter {
pub service: Option<String>,
pub finding_type: Option<String>,
pub severity: Option<String>,
pub limit: usize,
}
#[derive(Debug)]
pub struct FindingsStore {
inner: RwLock<VecDeque<StoredFinding>>,
max_size: usize,
}
impl FindingsStore {
#[must_use]
pub fn new(max_size: usize) -> Self {
const INITIAL_CAPACITY_CEILING: usize = 4096;
let capacity = max_size.min(INITIAL_CAPACITY_CEILING);
Self {
inner: RwLock::new(VecDeque::with_capacity(capacity)),
max_size,
}
}
pub async fn push_batch(&self, findings: &[Finding], now_ms: u64) {
if findings.is_empty() || self.max_size == 0 {
return;
}
let new_entries: Vec<StoredFinding> = findings
.iter()
.map(|f| StoredFinding {
finding: f.clone(),
stored_at_ms: now_ms,
first_seen_ms: now_ms,
seen_count: 1,
})
.collect();
let mut buf = self.inner.write().await;
buf.extend(new_entries);
if buf.len() > self.max_size {
let excess = buf.len() - self.max_size;
buf.drain(..excess);
}
}
pub async fn query(&self, filter: &FindingsFilter) -> Vec<StoredFinding> {
let buf = self.inner.read().await;
let limit = filter.limit;
buf.iter()
.rev()
.filter(|sf| {
if let Some(ref svc) = filter.service
&& sf.finding.service != *svc
{
return false;
}
if let Some(ref ft) = filter.finding_type
&& sf.finding.finding_type.as_str() != ft.as_str()
{
return false;
}
if let Some(ref sev) = filter.severity
&& sf.finding.severity.as_str() != sev.as_str()
{
return false;
}
true
})
.take(limit)
.cloned()
.collect()
}
pub async fn query_coalesced(&self, filter: &FindingsFilter) -> Vec<StoredFinding> {
let buf = self.inner.read().await;
let mut folded = fold_entries(buf.iter().rev().filter(|sf| {
if let Some(ref svc) = filter.service
&& sf.finding.service != *svc
{
return false;
}
if let Some(ref ft) = filter.finding_type
&& sf.finding.finding_type.as_str() != ft.as_str()
{
return false;
}
true
}));
drop(buf);
if let Some(ref sev) = filter.severity {
folded.retain(|sf| sf.finding.severity.as_str() == sev.as_str());
}
folded.truncate(filter.limit);
folded
}
pub async fn by_trace_id(&self, trace_id: &str) -> Vec<StoredFinding> {
let buf = self.inner.read().await;
buf.iter()
.rev()
.filter(|sf| sf.finding.trace_id == trace_id)
.cloned()
.collect()
}
pub async fn len(&self) -> usize {
self.inner.read().await.len()
}
pub async fn is_empty(&self) -> bool {
self.inner.read().await.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::acknowledgments::enrich_with_signatures;
use crate::detect::{Confidence, FindingType, Pattern, Severity};
fn make_finding(service: &str, finding_type: FindingType) -> Finding {
make_finding_with_template(service, finding_type, "SELECT 1")
}
fn make_finding_with_template(
service: &str,
finding_type: FindingType,
template: &str,
) -> Finding {
Finding {
finding_type,
severity: Severity::Warning,
trace_id: "trace-1".to_string(),
service: service.to_string(),
grouping: Vec::new(),
source_endpoint: "POST /api/test".to_string(),
pattern: Pattern {
template: template.to_string(),
occurrences: 5,
window_ms: 200,
distinct_params: 5,
..Default::default()
},
suggestion: "batch".to_string(),
first_timestamp: "2025-07-10T14:32:01.000Z".to_string(),
last_timestamp: "2025-07-10T14:32:01.200Z".to_string(),
green_impact: None,
confidence: Confidence::default(),
classification_method: None,
code_location: None,
instrumentation_scopes: Vec::new(),
suggested_fix: None,
signature: String::new(),
}
}
#[tokio::test]
async fn max_size_zero_disables_store() {
let store = FindingsStore::new(0);
let f = make_finding("svc", FindingType::NPlusOneSql);
store.push_batch(&[f], 1000).await;
assert_eq!(store.len().await, 0);
assert!(store.is_empty().await);
}
#[tokio::test]
async fn push_batch_respects_capacity() {
let store = FindingsStore::new(3);
let findings: Vec<Finding> = (0..5)
.map(|i| {
let mut f = make_finding("svc", FindingType::NPlusOneSql);
f.trace_id = format!("trace-{i}");
f
})
.collect();
store.push_batch(&findings, 1000).await;
assert_eq!(store.len().await, 3);
let all = store
.query(&FindingsFilter {
limit: 100,
..Default::default()
})
.await;
let trace_ids: Vec<&str> = all.iter().map(|sf| sf.finding.trace_id.as_str()).collect();
assert!(trace_ids.contains(&"trace-4"));
assert!(trace_ids.contains(&"trace-3"));
assert!(trace_ids.contains(&"trace-2"));
assert!(!trace_ids.contains(&"trace-0"));
}
#[tokio::test]
async fn query_keeps_instances_and_coalesced_folds_them() {
let store = FindingsStore::new(100);
for (trace, ts) in [("trace-a", 1000u64), ("trace-b", 2000)] {
let mut f = make_finding("svc", FindingType::RedundantSql);
f.trace_id = trace.to_string();
enrich_with_signatures(std::slice::from_mut(&mut f));
store.push_batch(&[f], ts).await;
}
let filter = FindingsFilter {
limit: 100,
..Default::default()
};
assert_eq!(store.query(&filter).await.len(), 2, "instances retained");
let folded = store.query_coalesced(&filter).await;
assert_eq!(folded.len(), 1);
assert_eq!(folded[0].seen_count, 2);
assert_eq!(folded[0].first_seen_ms, 1000);
assert_eq!(folded[0].stored_at_ms, 2000);
assert_eq!(
folded[0].finding.trace_id, "trace-b",
"the latest instance is the one kept"
);
}
#[tokio::test]
async fn coalescing_keeps_identical_signatures_separate_by_grouping_value() {
let mut prod = make_finding("svc", FindingType::RedundantSql);
prod.grouping = crate::test_helpers::k8s_grouping("prod-eu");
prod.grouping = crate::test_helpers::grouping("service.namespace", "payments");
let mut staging = make_finding("svc", FindingType::RedundantSql);
staging.grouping = crate::test_helpers::grouping("service.namespace", "staging");
enrich_with_signatures(std::slice::from_mut(&mut prod));
enrich_with_signatures(std::slice::from_mut(&mut staging));
assert_eq!(prod.signature, staging.signature);
let store = FindingsStore::new(100);
store.push_batch(&[prod], 1000).await;
store.push_batch(&[staging], 2000).await;
let folded = store
.query_coalesced(&FindingsFilter {
limit: 100,
..Default::default()
})
.await;
assert_eq!(folded.len(), 2);
}
#[tokio::test]
async fn coalescing_keeps_equal_grouping_values_separate_by_key() {
let mut tenant = make_finding("svc", FindingType::RedundantSql);
tenant.grouping = crate::test_helpers::grouping("tenant.id", "prod");
let mut namespace = make_finding("svc", FindingType::RedundantSql);
namespace.grouping = crate::test_helpers::grouping("k8s.namespace.name", "prod");
enrich_with_signatures(std::slice::from_mut(&mut tenant));
enrich_with_signatures(std::slice::from_mut(&mut namespace));
let store = FindingsStore::new(100);
store.push_batch(&[tenant], 1000).await;
store.push_batch(&[namespace], 2000).await;
let folded = store
.query_coalesced(&FindingsFilter {
limit: 100,
..Default::default()
})
.await;
assert_eq!(folded.len(), 2);
}
#[tokio::test]
async fn folded_row_evidence_matches_its_severity() {
let mut critical = make_finding("svc", FindingType::NPlusOneSql);
critical.severity = Severity::Critical;
critical.trace_id = "trace-hot".to_string();
critical.pattern.occurrences = 12;
let mut warning = make_finding("svc", FindingType::NPlusOneSql);
warning.severity = Severity::Warning;
warning.trace_id = "trace-quiet".to_string();
warning.pattern.occurrences = 6;
enrich_with_signatures(std::slice::from_mut(&mut critical));
enrich_with_signatures(std::slice::from_mut(&mut warning));
let store = FindingsStore::new(100);
store.push_batch(&[critical], 1000).await;
store.push_batch(&[warning], 2000).await;
let folded = store
.query_coalesced(&FindingsFilter {
limit: 100,
..Default::default()
})
.await;
assert_eq!(folded.len(), 1);
let row = &folded[0];
assert_eq!(row.finding.severity, Severity::Critical);
assert_eq!(
row.finding.trace_id, "trace-hot",
"the critical row must point at the trace that earned it"
);
assert_eq!(row.finding.pattern.occurrences, 12);
assert_eq!(row.first_seen_ms, 1000);
assert_eq!(row.stored_at_ms, 2000);
assert_eq!(row.seen_count, 2);
}
#[tokio::test]
async fn severity_filter_applies_after_folding() {
let mut critical = make_finding("svc", FindingType::NPlusOneSql);
critical.severity = Severity::Critical;
let mut warning = make_finding("svc", FindingType::NPlusOneSql);
warning.severity = Severity::Warning;
enrich_with_signatures(std::slice::from_mut(&mut critical));
enrich_with_signatures(std::slice::from_mut(&mut warning));
let store = FindingsStore::new(100);
store.push_batch(&[critical], 1000).await;
store.push_batch(&[warning], 2000).await;
let unfiltered = store
.query_coalesced(&FindingsFilter {
limit: 100,
..Default::default()
})
.await;
let filtered = store
.query_coalesced(&FindingsFilter {
severity: Some("critical".to_string()),
limit: 100,
..Default::default()
})
.await;
assert_eq!(filtered.len(), 1, "the group's worst severity matches");
assert_eq!(
filtered[0].seen_count, unfiltered[0].seen_count,
"the same problem must not report two different counts"
);
}
#[tokio::test]
async fn coalesced_limit_applies_after_folding() {
let store = FindingsStore::new(1000);
for i in 0..50u64 {
let mut hot =
make_finding_with_template("svc", FindingType::RedundantSql, "SELECT hot");
enrich_with_signatures(std::slice::from_mut(&mut hot));
store.push_batch(&[hot], 1000 + i).await;
}
let mut cold = make_finding_with_template("svc", FindingType::NPlusOneSql, "SELECT cold");
enrich_with_signatures(std::slice::from_mut(&mut cold));
store.push_batch(&[cold], 2000).await;
let folded = store
.query_coalesced(&FindingsFilter {
limit: 2,
..Default::default()
})
.await;
let templates: Vec<&str> = folded
.iter()
.map(|sf| sf.finding.pattern.template.as_str())
.collect();
assert_eq!(templates, ["SELECT cold", "SELECT hot"]);
assert_eq!(folded[1].seen_count, 50);
}
#[tokio::test]
async fn by_trace_id_answers_for_an_older_recurrence() {
let store = FindingsStore::new(100);
for (trace, ts) in [("trace-old", 1000u64), ("trace-new", 2000)] {
let mut f = make_finding("svc", FindingType::RedundantSql);
f.trace_id = trace.to_string();
enrich_with_signatures(std::slice::from_mut(&mut f));
store.push_batch(&[f], ts).await;
}
let hits = store.by_trace_id("trace-old").await;
assert_eq!(hits.len(), 1, "the older instance is still retrievable");
assert_eq!(hits[0].finding.trace_id, "trace-old");
}
#[tokio::test]
async fn query_filters_by_service() {
let store = FindingsStore::new(100);
let f1 = make_finding("order-svc", FindingType::NPlusOneSql);
let f2 = make_finding("payment-svc", FindingType::NPlusOneSql);
store.push_batch(&[f1, f2], 1000).await;
let results = store
.query(&FindingsFilter {
service: Some("order-svc".to_string()),
limit: 100,
..Default::default()
})
.await;
assert_eq!(results.len(), 1);
assert_eq!(results[0].finding.service, "order-svc");
}
#[tokio::test]
async fn query_filters_by_type() {
let store = FindingsStore::new(100);
let f1 = make_finding("svc", FindingType::NPlusOneSql);
let f2 = make_finding("svc", FindingType::RedundantSql);
store.push_batch(&[f1, f2], 1000).await;
let results = store
.query(&FindingsFilter {
finding_type: Some("n_plus_one_sql".to_string()),
limit: 100,
..Default::default()
})
.await;
assert_eq!(results.len(), 1);
assert_eq!(results[0].finding.finding_type, FindingType::NPlusOneSql);
}
#[tokio::test]
async fn by_trace_id_filters_correctly() {
let store = FindingsStore::new(100);
let mut f1 = make_finding_with_template("svc", FindingType::NPlusOneSql, "SELECT a");
f1.trace_id = "trace-a".to_string();
let mut f2 = make_finding_with_template("svc", FindingType::NPlusOneSql, "SELECT b");
f2.trace_id = "trace-b".to_string();
store.push_batch(&[f1, f2], 1000).await;
let results = store.by_trace_id("trace-a").await;
assert_eq!(results.len(), 1);
assert_eq!(results[0].finding.trace_id, "trace-a");
}
#[tokio::test]
async fn query_respects_limit() {
let store = FindingsStore::new(100);
let findings: Vec<Finding> = (0..10)
.map(|i| {
make_finding_with_template("svc", FindingType::NPlusOneSql, &format!("SELECT {i}"))
})
.collect();
store.push_batch(&findings, 1000).await;
let results = store
.query(&FindingsFilter {
limit: 3,
..Default::default()
})
.await;
assert_eq!(results.len(), 3);
}
}