use std::collections::{HashMap, HashSet};
use std::time::Duration;
use anyhow::{Context, Result};
use reqwest::Client;
use serde::{Deserialize, Serialize};
use tracing::{debug, info, warn};
use kingfisher_core::ValidationOutcome;
use crate::cli::commands::scan::ConfidenceLevel;
use crate::reporter::FindingReporterRecord;
pub mod discord;
pub mod generic;
pub mod googlechat;
pub mod mattermost;
pub mod slack;
pub mod teams;
#[derive(Copy, Clone, Debug, PartialEq, Eq, Serialize, Deserialize, clap::ValueEnum)]
#[serde(rename_all = "lowercase")]
#[clap(rename_all = "lowercase")]
#[derive(Default)]
pub enum AlertOn {
#[default]
Findings,
Always,
}
#[derive(Copy, Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize, clap::ValueEnum)]
#[serde(rename_all = "lowercase")]
#[clap(rename_all = "lowercase")]
pub enum AlertDetail {
Summary,
Detail,
#[default]
Auto,
}
pub const AUTO_DETAIL_THRESHOLD: usize = 25;
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
struct AccessMapImpact {
entry_id: usize,
resources: usize,
}
type AccessMapImpactIndex = HashMap<String, Vec<AccessMapImpact>>;
#[derive(Clone, Debug)]
#[doc(hidden)]
pub struct AlertAccessMapEntry {
pub(crate) finding_fingerprints: Vec<String>,
pub(crate) impacted_resources: usize,
pub(crate) mapping_succeeded: bool,
}
#[derive(Copy, Clone, Debug, PartialEq, Eq, Serialize, Deserialize, clap::ValueEnum)]
#[serde(rename_all = "lowercase")]
#[clap(rename_all = "lowercase")]
pub enum AlertFormat {
Slack,
Teams,
Generic,
Discord,
Mattermost,
Googlechat,
}
impl AlertFormat {
pub fn infer_from_url(url: &str) -> Self {
let host = url::Url::parse(url).ok().and_then(|u| u.host_str().map(str::to_lowercase));
match host.as_deref() {
Some(h) if host_matches(h, "slack.com") => AlertFormat::Slack,
Some(h)
if host_matches(h, "office.com")
|| host_matches(h, "webhook.office.com")
|| host_matches(h, "webhook.office.net") =>
{
AlertFormat::Teams
}
Some(h) if host_matches(h, "discord.com") || host_matches(h, "discordapp.com") => {
AlertFormat::Discord
}
Some(h) if host_matches(h, "chat.googleapis.com") => AlertFormat::Googlechat,
_ => AlertFormat::Generic,
}
}
}
#[derive(Copy, Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize, clap::ValueEnum)]
#[serde(rename_all = "kebab-case")]
#[clap(rename_all = "kebab-case")]
pub enum AlertFindingFilter {
#[default]
All,
ExcludeInactive,
OnlyActive,
AccessMapOnly,
}
#[derive(Clone, Debug)]
pub struct AlertSink {
pub url: String,
pub format: AlertFormat,
pub on: AlertOn,
pub min_confidence: ConfidenceLevel,
pub include_secret: bool,
pub report_url: Option<String>,
pub detail: AlertDetail,
pub finding_filter: AlertFindingFilter,
pub prevent_empty: bool,
}
#[derive(Clone, Debug, Serialize)]
pub struct AlertSummary {
pub total: usize,
pub active: usize,
pub inactive: usize,
pub unknown: usize,
pub by_rule: Vec<(String, usize)>,
pub kingfisher_version: String,
pub target: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub report_url: Option<String>,
pub detail: AlertDetail,
pub impacted_resources: usize,
pub unfiltered_total: usize,
}
impl AlertSummary {
pub fn from_findings(findings: &[FindingReporterRecord], target: Option<String>) -> Self {
let findings: Vec<_> = findings.iter().collect();
Self::from_filtered_findings(&findings, target, &HashMap::new())
}
fn from_filtered_findings(
findings: &[&FindingReporterRecord],
target: Option<String>,
access_map_impact: &AccessMapImpactIndex,
) -> Self {
let mut active = 0usize;
let mut inactive = 0usize;
let mut unknown = 0usize;
let mut impacted_resources = 0usize;
let mut counted_access_map_entries = HashSet::new();
let mut by_rule_map: HashMap<String, usize> = HashMap::new();
for f in findings {
*by_rule_map.entry(f.rule.id.clone()).or_default() += 1;
match f.finding.validation.outcome {
ValidationOutcome::VerifiedActive => active += 1,
ValidationOutcome::VerifiedInactive => inactive += 1,
_ => unknown += 1,
}
if let Some(impacts) = access_map_impact.get(&f.finding.fingerprint) {
for impact in impacts {
if counted_access_map_entries.insert(impact.entry_id) {
impacted_resources += impact.resources;
}
}
}
}
let mut by_rule: Vec<(String, usize)> = by_rule_map.into_iter().collect();
by_rule.sort_by(|a, b| b.1.cmp(&a.1).then(a.0.cmp(&b.0)));
by_rule.truncate(5);
Self {
total: findings.len(),
active,
inactive,
unknown,
by_rule,
kingfisher_version: env!("CARGO_PKG_VERSION").to_string(),
target,
report_url: None,
detail: AlertDetail::Detail,
impacted_resources,
unfiltered_total: 0,
}
}
}
fn build_client() -> Result<Client> {
Client::builder()
.timeout(Duration::from_secs(15))
.connect_timeout(Duration::from_secs(5))
.user_agent(format!("kingfisher/{}", env!("CARGO_PKG_VERSION")))
.build()
.context("failed to build webhook reqwest::Client")
}
fn host_matches(host: &str, suffix: &str) -> bool {
host == suffix || host.ends_with(&format!(".{suffix}"))
}
pub fn validate_webhook_url(url: &str) -> Result<()> {
let parsed = url::Url::parse(url)
.with_context(|| format!("invalid webhook URL `{}`", redact_for_log(url)))?;
let scheme = parsed.scheme();
let host = parsed.host_str().unwrap_or("");
if host.is_empty() {
anyhow::bail!("webhook URL `{}` has no host", redact_for_log(url));
}
match scheme {
"https" => {}
"http" if is_loopback_host(host) => {}
"http" => {
anyhow::bail!(
"webhook URL `{}` uses cleartext `http://`; webhook tokens and finding \
metadata must not traverse the network unencrypted. Use `https://`, or a \
loopback host (`localhost`/`127.0.0.1`/`::1`) for local testing.",
redact_for_log(url)
);
}
_ => {
anyhow::bail!(
"webhook URL `{}` uses unsupported scheme `{scheme}` (only `https` is \
allowed; `http` is allowed only for loopback hosts)",
redact_for_log(url)
);
}
}
Ok(())
}
fn is_loopback_host(host: &str) -> bool {
if host.eq_ignore_ascii_case("localhost") {
return true;
}
let trimmed = host.strip_prefix('[').and_then(|s| s.strip_suffix(']')).unwrap_or(host);
if let Ok(ip) = trimmed.parse::<std::net::IpAddr>() {
return ip.is_loopback();
}
false
}
fn redact_for_log(url: &str) -> String {
redact_webhook(url)
}
pub fn redact_webhook(url: &str) -> String {
match url::Url::parse(url) {
Ok(u) => {
let scheme = u.scheme();
let host = u.host_str().unwrap_or("");
let port = u.port().map(|p| format!(":{p}")).unwrap_or_default();
format!("{scheme}://{host}{port}/<redacted>")
}
Err(_) => "<unparseable webhook url>".to_string(),
}
}
pub async fn dispatch(
sinks: &[AlertSink],
findings: &[FindingReporterRecord],
target: Option<String>,
) {
dispatch_with_context(sinks, findings, &[], target, false).await;
}
#[doc(hidden)]
pub async fn dispatch_with_context(
sinks: &[AlertSink],
findings: &[FindingReporterRecord],
access_map: &[AlertAccessMapEntry],
target: Option<String>,
dry_run: bool,
) {
if sinks.is_empty() {
return;
}
let mut client = None;
if dry_run && sinks.iter().any(|sink| sink.include_secret) {
warn!("alert dry-run: include_secret is ignored; dry-run payloads are always redacted");
}
let unfiltered_total = findings.len();
let access_map_impact = build_access_map_impact(access_map);
debug!("alert dispatch: total={} sinks={}", unfiltered_total, sinks.len());
for sink in sinks {
if matches!(sink.on, AlertOn::Findings) && unfiltered_total == 0 {
debug!(
"alert dispatch: skipping {} (on=findings, no findings)",
redact_webhook(&sink.url)
);
continue;
}
let filtered: Vec<&FindingReporterRecord> = findings
.iter()
.filter(|f| matches_min_confidence(&f.finding.confidence, sink.min_confidence))
.filter(|f| {
matches_finding_filter(
f.finding.validation.outcome,
&f.finding.fingerprint,
sink.finding_filter,
&access_map_impact,
)
})
.collect();
let is_heartbeat_sink = matches!(sink.on, AlertOn::Always);
if sink.prevent_empty && filtered.is_empty() && !is_heartbeat_sink {
debug!(
"alert dispatch: skipping {} (filters left nothing to report)",
redact_webhook(&sink.url)
);
continue;
}
let resolved_detail = match sink.detail {
AlertDetail::Auto => {
if filtered.len() > AUTO_DETAIL_THRESHOLD {
AlertDetail::Summary
} else {
AlertDetail::Detail
}
}
other => other,
};
let mut summary =
AlertSummary::from_filtered_findings(&filtered, target.clone(), &access_map_impact);
summary.report_url = sink.report_url.clone();
summary.detail = resolved_detail;
summary.unfiltered_total = unfiltered_total;
let payload = build_sink_payload(sink, &summary, &filtered, dry_run);
if dry_run {
info!(
"alert dry-run: would POST to {} ({} finding(s)):\n{}",
redact_webhook(&sink.url),
filtered.len(),
serde_json::to_string_pretty(&payload).unwrap_or_default()
);
continue;
}
if client.is_none() {
client = match build_client() {
Ok(client) => Some(client),
Err(e) => {
warn!("alert dispatch: failed to build HTTP client: {}", e);
return;
}
};
}
match post(client.as_ref().expect("client initialized above"), &sink.url, &payload).await {
Ok(()) => {
info!("alert posted to {}", redact_webhook(&sink.url));
}
Err(e) => {
warn!("alert dispatch failed for {}: {}", redact_webhook(&sink.url), e);
}
}
}
}
fn build_access_map_impact(access_map: &[AlertAccessMapEntry]) -> AccessMapImpactIndex {
let mut index = HashMap::new();
for (entry_id, entry) in access_map.iter().enumerate() {
if !entry.mapping_succeeded {
continue;
}
let fingerprints: HashSet<&str> =
entry.finding_fingerprints.iter().map(String::as_str).collect();
for fingerprint in fingerprints {
index
.entry(fingerprint.to_string())
.or_insert_with(Vec::new)
.push(AccessMapImpact { entry_id, resources: entry.impacted_resources });
}
}
index
}
fn build_sink_payload(
sink: &AlertSink,
summary: &AlertSummary,
findings: &[&FindingReporterRecord],
dry_run: bool,
) -> serde_json::Value {
let include_secret = sink.include_secret && !dry_run;
match sink.format {
AlertFormat::Slack => slack::build_payload(summary, findings, include_secret),
AlertFormat::Teams => teams::build_payload(summary, findings, include_secret),
AlertFormat::Generic => generic::build_payload(summary, findings, include_secret),
AlertFormat::Discord => discord::build_payload(summary, findings, include_secret),
AlertFormat::Mattermost => mattermost::build_payload(summary, findings, include_secret),
AlertFormat::Googlechat => googlechat::build_payload(summary, findings, include_secret),
}
}
fn matches_min_confidence(finding_confidence: &str, threshold: ConfidenceLevel) -> bool {
let level = match finding_confidence {
"Low" => ConfidenceLevel::Low,
"Medium" => ConfidenceLevel::Medium,
"High" => ConfidenceLevel::High,
_ => ConfidenceLevel::Medium,
};
level >= threshold
}
fn matches_finding_filter(
outcome: ValidationOutcome,
fingerprint: &str,
filter: AlertFindingFilter,
access_map_impact: &AccessMapImpactIndex,
) -> bool {
match filter {
AlertFindingFilter::All => true,
AlertFindingFilter::ExcludeInactive => outcome != ValidationOutcome::VerifiedInactive,
AlertFindingFilter::OnlyActive => outcome.is_verified_active(),
AlertFindingFilter::AccessMapOnly => {
outcome.is_verified_active() && access_map_impact.contains_key(fingerprint)
}
}
}
async fn post(client: &Client, url: &str, payload: &serde_json::Value) -> Result<()> {
let resp = client
.post(url)
.json(payload)
.send()
.await
.with_context(|| format!("POST to {} failed", redact_webhook(url)))?;
let status = resp.status();
if !status.is_success() {
let body = resp.text().await.unwrap_or_default();
anyhow::bail!(
"webhook returned HTTP {}: {}",
status,
body.chars().take(200).collect::<String>()
);
}
Ok(())
}
#[cfg(test)]
pub(crate) fn make_test_record(
rule_id: &str,
fingerprint: &str,
) -> crate::reporter::FindingReporterRecord {
use crate::reporter::{FindingRecordData, FindingReporterRecord, RuleMetadata, ValidationInfo};
FindingReporterRecord {
rule: RuleMetadata {
title: format!("{} => [{}]", rule_id.to_uppercase(), rule_id.to_uppercase()),
name: rule_id.to_string(),
id: rule_id.to_string(),
description: rule_id.to_string(),
},
finding: FindingRecordData {
dependent_captures: Default::default(),
ambiguous_dependencies: Default::default(),
snippet: "AKIAEXAMPLE_REDACTED_TOKEN_12345".to_string(),
fingerprint: fingerprint.to_string(),
confidence: "Medium".to_string(),
entropy: "4.5".to_string(),
validation: ValidationInfo {
outcome: ValidationOutcome::VerifiedActive,
status: "Active Credential".to_string(),
response: String::new(),
},
language: "rust".to_string(),
line: 42,
column_start: 10,
column_end: 50,
path: "src/foo.rs".to_string(),
encoding: None,
git_metadata: None,
validate_command: None,
revoke_command: None,
blast_radius_command: None,
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use kingfisher_core::ValidationOutcome as VO;
#[test]
fn redact_webhook_keeps_host() {
let r = redact_webhook("https://hooks.slack.com/services/T0/B0/XXX");
assert_eq!(r, "https://hooks.slack.com/<redacted>");
}
#[test]
fn redact_webhook_unparseable() {
let r = redact_webhook("not a url");
assert_eq!(r, "<unparseable webhook url>");
}
#[test]
fn validate_webhook_accepts_https() {
validate_webhook_url("https://hooks.slack.com/services/T0/B0/XXX").unwrap();
}
#[test]
fn validate_webhook_rejects_remote_http() {
let err = validate_webhook_url("http://example.com/hook").unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("cleartext `http://`"), "got: {msg}");
}
#[test]
fn validate_webhook_allows_http_localhost() {
validate_webhook_url("http://localhost:8080/hook").unwrap();
validate_webhook_url("http://127.0.0.1:9000/hook").unwrap();
validate_webhook_url("http://[::1]:9000/hook").unwrap();
}
#[test]
fn validate_webhook_rejects_unknown_scheme() {
let err = validate_webhook_url("ftp://example.com/hook").unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("unsupported scheme"), "got: {msg}");
}
#[test]
fn validate_webhook_rejects_no_host() {
let err = validate_webhook_url("file:///etc/passwd").unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("no host") || msg.contains("unsupported scheme"), "got: {msg}");
}
#[test]
fn infer_format_slack() {
assert_eq!(
AlertFormat::infer_from_url("https://hooks.slack.com/services/T0/B0/XXX"),
AlertFormat::Slack
);
}
#[test]
fn infer_format_teams() {
assert_eq!(
AlertFormat::infer_from_url(
"https://outlook.office.com/webhook/abc/IncomingWebhook/def"
),
AlertFormat::Teams
);
}
#[test]
fn infer_format_generic_fallback() {
assert_eq!(
AlertFormat::infer_from_url("https://example.com/webhook"),
AlertFormat::Generic
);
}
#[test]
fn infer_format_discord() {
assert_eq!(
AlertFormat::infer_from_url("https://discord.com/api/webhooks/123/abc"),
AlertFormat::Discord
);
assert_eq!(
AlertFormat::infer_from_url("https://discordapp.com/api/webhooks/123/abc"),
AlertFormat::Discord
);
}
#[test]
fn infer_format_googlechat() {
assert_eq!(
AlertFormat::infer_from_url(
"https://chat.googleapis.com/v1/spaces/AAA/messages?key=k&token=t"
),
AlertFormat::Googlechat
);
}
#[test]
fn infer_format_mattermost_falls_back_to_generic_without_override() {
assert_eq!(
AlertFormat::infer_from_url("https://mattermost.example.com/hooks/abcdef"),
AlertFormat::Generic
);
}
#[test]
fn auto_detail_threshold_is_inclusive_at_25() {
assert_eq!(AUTO_DETAIL_THRESHOLD, 25);
}
#[test]
fn finding_filter_all_matches_everything() {
let map = HashMap::new();
for outcome in [VO::VerifiedActive, VO::VerifiedInactive, VO::NotAttempted, VO::Assumed] {
assert!(matches_finding_filter(outcome, "fp1", AlertFindingFilter::All, &map));
}
}
#[test]
fn finding_filter_exclude_inactive_drops_only_inactive() {
let map = HashMap::new();
assert!(matches_finding_filter(
VO::VerifiedActive,
"fp1",
AlertFindingFilter::ExcludeInactive,
&map
));
assert!(matches_finding_filter(
VO::NotAttempted,
"fp1",
AlertFindingFilter::ExcludeInactive,
&map
));
assert!(!matches_finding_filter(
VO::VerifiedInactive,
"fp1",
AlertFindingFilter::ExcludeInactive,
&map
));
}
#[test]
fn finding_filter_only_active_keeps_only_verified_active() {
let map = HashMap::new();
assert!(matches_finding_filter(
VO::VerifiedActive,
"fp1",
AlertFindingFilter::OnlyActive,
&map
));
for outcome in [VO::VerifiedInactive, VO::NotAttempted, VO::Assumed, VO::Unavailable] {
assert!(!matches_finding_filter(outcome, "fp1", AlertFindingFilter::OnlyActive, &map));
}
}
#[test]
fn finding_filter_access_map_only_requires_fingerprint_match() {
let mut map = HashMap::new();
map.insert("fp-mapped".to_string(), vec![AccessMapImpact { entry_id: 0, resources: 3 }]);
assert!(matches_finding_filter(
VO::VerifiedActive,
"fp-mapped",
AlertFindingFilter::AccessMapOnly,
&map
));
assert!(!matches_finding_filter(
VO::VerifiedActive,
"fp-other",
AlertFindingFilter::AccessMapOnly,
&map
));
}
#[test]
fn finding_filter_access_map_only_still_requires_active_outcome() {
let mut map = HashMap::new();
map.insert("fp-gitlab".to_string(), vec![AccessMapImpact { entry_id: 0, resources: 1 }]);
assert!(!matches_finding_filter(
VO::VerifiedInactive,
"fp-gitlab",
AlertFindingFilter::AccessMapOnly,
&map
));
}
fn record_with(
rule_id: &str,
fingerprint: &str,
confidence: &str,
outcome: ValidationOutcome,
) -> crate::reporter::FindingReporterRecord {
use crate::reporter::{
FindingRecordData, FindingReporterRecord, RuleMetadata, ValidationInfo,
};
FindingReporterRecord {
rule: RuleMetadata {
title: format!("{} => [{}]", rule_id.to_uppercase(), rule_id.to_uppercase()),
name: rule_id.to_string(),
id: rule_id.to_string(),
description: rule_id.to_string(),
},
finding: FindingRecordData {
dependent_captures: Default::default(),
ambiguous_dependencies: Default::default(),
snippet: "AKIAEXAMPLE_REDACTED_TOKEN_12345".to_string(),
fingerprint: fingerprint.to_string(),
confidence: confidence.to_string(),
entropy: "4.5".to_string(),
validation: ValidationInfo {
outcome,
status: outcome.display_name().to_string(),
response: String::new(),
},
language: "rust".to_string(),
line: 42,
column_start: 10,
column_end: 50,
path: "src/foo.rs".to_string(),
encoding: None,
git_metadata: None,
validate_command: None,
revoke_command: None,
blast_radius_command: None,
},
}
}
fn test_sink(url: &str) -> AlertSink {
AlertSink {
url: url.to_string(),
format: AlertFormat::Generic,
on: AlertOn::Findings,
min_confidence: ConfidenceLevel::Low,
include_secret: false,
report_url: None,
detail: AlertDetail::Detail,
finding_filter: AlertFindingFilter::All,
prevent_empty: false,
}
}
fn access_map_entry(
fingerprints: &[&str],
resources: &[&str],
mapping_succeeded: bool,
) -> AlertAccessMapEntry {
AlertAccessMapEntry {
finding_fingerprints: fingerprints
.iter()
.map(|fingerprint| (*fingerprint).to_string())
.collect(),
impacted_resources: resources.len(),
mapping_succeeded,
}
}
#[test]
fn dry_run_payload_stays_redacted_when_sink_includes_secrets() {
let mut sink = test_sink("https://example.com/webhook");
sink.include_secret = true;
let finding =
record_with("betterleaks.aws-access-token", "fp1", "High", VO::VerifiedActive);
let findings = vec![&finding];
let summary = AlertSummary::from_filtered_findings(&findings, None, &HashMap::new());
let dry_run_payload = build_sink_payload(&sink, &summary, &findings, true).to_string();
assert!(!dry_run_payload.contains("AKIAEXAMPLE_REDACTED_TOKEN_12345"));
assert!(dry_run_payload.contains("<redacted>"));
let live_payload = build_sink_payload(&sink, &summary, &findings, false).to_string();
assert!(live_payload.contains("AKIAEXAMPLE_REDACTED_TOKEN_12345"));
}
#[test]
fn access_map_impact_ignores_failed_mappings() {
let access_map = vec![access_map_entry(&["fp-mapped"], &["failed identity"], false)];
assert!(build_access_map_impact(&access_map).is_empty());
}
#[test]
fn access_map_impact_accumulates_distinct_entries_for_one_fingerprint() {
let access_map = vec![
access_map_entry(&["fp-mapped"], &["bucket-a", "bucket-b"], true),
access_map_entry(&["fp-mapped"], &["queue-a"], true),
];
let impact = build_access_map_impact(&access_map);
let finding =
record_with("betterleaks.aws-access-token", "fp-mapped", "High", VO::VerifiedActive);
let summary = AlertSummary::from_filtered_findings(&[&finding], None, &impact);
assert_eq!(summary.impacted_resources, 3);
}
mod dispatch_tests {
use wiremock::{Mock, MockServer, ResponseTemplate, matchers::method};
use super::*;
#[tokio::test]
async fn prevent_empty_skips_sink_when_filters_leave_nothing() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.min_confidence = ConfidenceLevel::High;
sink.prevent_empty = true;
let findings = vec![record_with(
"betterleaks.aws-access-token",
"fp1",
"Medium",
VO::VerifiedActive,
)];
dispatch_with_context(&[sink], &findings, &[], None, false).await;
assert_eq!(server.received_requests().await.unwrap().len(), 0);
}
#[tokio::test]
async fn prevent_empty_false_still_posts_when_filters_leave_nothing() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.min_confidence = ConfidenceLevel::High;
sink.prevent_empty = false;
let findings = vec![record_with(
"betterleaks.aws-access-token",
"fp1",
"Medium",
VO::VerifiedActive,
)];
dispatch_with_context(&[sink], &findings, &[], None, false).await;
let requests = server.received_requests().await.unwrap();
assert_eq!(requests.len(), 1);
let body: serde_json::Value = requests[0].body_json().unwrap();
assert_eq!(body["summary"]["total"], 0);
assert_eq!(body["summary"]["unfiltered_total"], 1);
}
#[tokio::test]
async fn always_heartbeat_posts_despite_prevent_empty_on_zero_total() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.on = AlertOn::Always;
sink.prevent_empty = true;
dispatch_with_context(&[sink], &[], &[], None, false).await;
assert_eq!(server.received_requests().await.unwrap().len(), 1);
}
#[tokio::test]
async fn always_heartbeat_posts_even_when_filters_drop_every_finding() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.on = AlertOn::Always;
sink.prevent_empty = true;
sink.finding_filter = AlertFindingFilter::OnlyActive;
let findings = vec![record_with(
"betterleaks.aws-access-token",
"fp1",
"High",
VO::VerifiedInactive,
)];
dispatch_with_context(&[sink], &findings, &[], None, false).await;
let requests = server.received_requests().await.unwrap();
assert_eq!(requests.len(), 1);
let body: serde_json::Value = requests[0].body_json().unwrap();
assert_eq!(body["summary"]["total"], 0);
assert_eq!(body["summary"]["unfiltered_total"], 1);
}
#[tokio::test]
async fn access_map_only_filters_to_matching_fingerprint_and_reports_impact() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.finding_filter = AlertFindingFilter::AccessMapOnly;
sink.min_confidence = ConfidenceLevel::Low;
let mapped = record_with(
"betterleaks.aws-access-token",
"fp-mapped",
"High",
VO::VerifiedActive,
);
let unmapped = record_with(
"betterleaks.aws-secret-access-key",
"fp-other",
"High",
VO::VerifiedActive,
);
let access_map = vec![access_map_entry(
&["fp-mapped"],
&["arn:aws:s3:::bucket-a", "arn:aws:s3:::bucket-b"],
true,
)];
dispatch_with_context(&[sink], &[mapped, unmapped], &access_map, None, false).await;
let requests = server.received_requests().await.unwrap();
assert_eq!(requests.len(), 1);
let body: serde_json::Value = requests[0].body_json().unwrap();
assert_eq!(body["findings"].as_array().unwrap().len(), 1);
assert_eq!(body["findings"][0]["finding"]["fingerprint"], "fp-mapped");
assert_eq!(body["summary"]["impacted_resources"], 2);
}
#[tokio::test]
async fn access_map_only_excludes_finding_with_access_map_entry_but_inactive_status() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.finding_filter = AlertFindingFilter::AccessMapOnly;
sink.min_confidence = ConfidenceLevel::Low;
sink.prevent_empty = true;
let inactive_but_mapped =
record_with("betterleaks.gitlab-pat", "fp-gitlab", "High", VO::VerifiedInactive);
let access_map = vec![access_map_entry(&["fp-gitlab"], &["group/project"], true)];
dispatch_with_context(&[sink], &[inactive_but_mapped], &access_map, None, false).await;
assert_eq!(server.received_requests().await.unwrap().len(), 0);
}
#[tokio::test]
async fn access_map_only_skips_failed_mapping_results() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.finding_filter = AlertFindingFilter::AccessMapOnly;
sink.prevent_empty = true;
let finding = record_with(
"betterleaks.aws-access-token",
"fp-mapped",
"High",
VO::VerifiedActive,
);
let access_map = vec![access_map_entry(&["fp-mapped"], &["mapping failed"], false)];
dispatch_with_context(&[sink], &[finding], &access_map, None, false).await;
assert_eq!(server.received_requests().await.unwrap().len(), 0);
}
#[tokio::test]
async fn repeated_credential_fingerprints_all_match_without_double_counting_impact() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.finding_filter = AlertFindingFilter::AccessMapOnly;
let findings = vec![
record_with("betterleaks.aws-access-token", "fp-first", "High", VO::VerifiedActive),
record_with(
"betterleaks.aws-access-token",
"fp-second",
"High",
VO::VerifiedActive,
),
];
let access_map = vec![access_map_entry(
&["fp-first", "fp-second"],
&["arn:aws:s3:::bucket-a", "arn:aws:s3:::bucket-b"],
true,
)];
dispatch_with_context(&[sink], &findings, &access_map, None, false).await;
let requests = server.received_requests().await.unwrap();
assert_eq!(requests.len(), 1);
let body: serde_json::Value = requests[0].body_json().unwrap();
assert_eq!(body["findings"].as_array().unwrap().len(), 2);
assert_eq!(body["summary"]["impacted_resources"], 2);
}
#[tokio::test]
async fn access_map_only_filters_everything_when_access_map_is_empty() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.finding_filter = AlertFindingFilter::AccessMapOnly;
sink.prevent_empty = true;
let findings = vec![record_with(
"betterleaks.aws-access-token",
"fp1",
"High",
VO::VerifiedActive,
)];
dispatch_with_context(&[sink], &findings, &[], None, false).await;
assert_eq!(server.received_requests().await.unwrap().len(), 0);
}
#[tokio::test]
async fn dry_run_makes_no_http_calls() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let sink = test_sink(&server.uri());
let findings = vec![record_with(
"betterleaks.aws-access-token",
"fp1",
"High",
VO::VerifiedActive,
)];
dispatch_with_context(&[sink], &findings, &[], None, true).await;
assert_eq!(server.received_requests().await.unwrap().len(), 0);
}
#[tokio::test]
async fn header_counts_reflect_only_active_filter_not_whole_scan() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
.await;
let mut sink = test_sink(&server.uri());
sink.finding_filter = AlertFindingFilter::OnlyActive;
let findings = vec![
record_with(
"betterleaks.aws-access-token",
"fp-active",
"High",
VO::VerifiedActive,
),
record_with(
"betterleaks.aws-secret-access-key",
"fp-inactive",
"High",
VO::VerifiedInactive,
),
];
dispatch_with_context(&[sink], &findings, &[], None, false).await;
let requests = server.received_requests().await.unwrap();
assert_eq!(requests.len(), 1);
let body: serde_json::Value = requests[0].body_json().unwrap();
assert_eq!(body["summary"]["total"], 1);
assert_eq!(body["summary"]["active"], 1);
assert_eq!(body["summary"]["inactive"], 0);
assert_eq!(body["findings"].as_array().unwrap().len(), 1);
}
}
}