WP-20/21: vulnerability management with Trivy and mail notifications
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled
CI / ui (push) Has been cancelled

Trivy scanner adapter (rootfs + image JSON, parsed and deduplicated),
findings repository, scan diff that keeps first_seen, marks disappeared
findings fixed and skips failed targets, vulnerability_scan job (daily by
default), digest mail for new findings at or above a configurable severity,
/api/vulnerabilities routes, Vulnerabilities page with severity tiles,
filters, details and acknowledge, notification threshold in settings.
Deploy script installs Trivy from the Aqua apt repository.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Dennis Nemec
2026-09-02 22:45:57 +02:00
parent 5c6e09ad10
commit 872f4373ff
13 changed files with 801 additions and 94 deletions

View File

@ -11,29 +11,39 @@ pub mod settings;
pub mod system;
pub mod test_support;
pub mod users;
pub mod vulnerabilities;
use std::sync::Arc;
use application::scheduler::Scheduler;
use application::{
AuthService, ClusterService, InventoryService, JobRunner, PackageRefreshJob, PackageUpgradeJob,
SettingsService, UserService,
SettingsService, UserService, VulnerabilityScanJob, VulnerabilityService,
};
use axum::{routing::get, Json, Router};
use domain::jobs::JobKind;
use domain::ports::Mailer;
use domain::ports::{ClusterGateway, HostInspector, HostUpdater};
use domain::ports::{ClusterGateway, HostInspector, HostUpdater, VulnerabilityScanner};
use infrastructure::{
AesGcmCipher, Argon2Hasher, DbPool, DebianInspector, DebianUpdater, FakeClusterGateway,
FakeHostInspector, FakeHostUpdater, JwtIssuer, KubeGateway, LettreMailer, SqliteAuditLog,
SqliteInventory, SqliteJobRuns, SqliteRefreshTokens, SqliteSettings, SqliteUsers,
SystemCommandRunner,
FakeHostInspector, FakeHostUpdater, FakeScanner, JwtIssuer, KubeGateway, LettreMailer,
SqliteAuditLog, SqliteFindings, SqliteInventory, SqliteJobRuns, SqliteRefreshTokens,
SqliteSettings, SqliteUsers, SystemCommandRunner, TrivyScanner,
};
use tower_http::services::{ServeDir, ServeFile};
use tower_http::trace::TraceLayer;
pub use config::Config;
/// External adapters the app is wired with (real or fake).
pub struct Adapters {
pub mailer: Arc<dyn Mailer>,
pub inspector: Arc<dyn HostInspector>,
pub updater: Arc<dyn HostUpdater>,
pub cluster: Arc<dyn ClusterGateway>,
pub scanner: Arc<dyn VulnerabilityScanner>,
}
#[derive(Clone)]
pub struct AppState {
pub cfg: Config,
@ -43,6 +53,7 @@ pub struct AppState {
pub jobs: Arc<JobRunner>,
pub inventory: Arc<InventoryService>,
pub cluster: Arc<ClusterService>,
pub vulns: Arc<VulnerabilityService>,
pub login_limiter: Arc<rate_limit::RateLimiter>,
}
@ -50,44 +61,42 @@ impl AppState {
/// Wire the services on top of a connected database.
pub fn new(cfg: Config, pool: DbPool) -> anyhow::Result<Self> {
let runner = Arc::new(SystemCommandRunner);
let (inspector, updater, cluster): (
Arc<dyn HostInspector>,
Arc<dyn HostUpdater>,
Arc<dyn ClusterGateway>,
) = if cfg.fake_host {
(
Arc::new(FakeHostInspector),
Arc::new(FakeHostUpdater),
Arc::new(FakeClusterGateway),
)
let adapters = if cfg.fake_host {
Adapters {
mailer: Arc::new(LettreMailer),
inspector: Arc::new(FakeHostInspector),
updater: Arc::new(FakeHostUpdater),
cluster: Arc::new(FakeClusterGateway),
scanner: Arc::new(FakeScanner),
}
} else {
(
Arc::new(DebianInspector::new(runner.clone())),
Arc::new(DebianUpdater::new(runner)),
Arc::new(KubeGateway::new(cfg.kubeconfig.clone())),
)
Adapters {
mailer: Arc::new(LettreMailer),
inspector: Arc::new(DebianInspector::new(runner.clone())),
updater: Arc::new(DebianUpdater::new(runner.clone())),
cluster: Arc::new(KubeGateway::new(cfg.kubeconfig.clone())),
scanner: Arc::new(TrivyScanner::new(runner)),
}
};
Self::with_adapters(
cfg,
pool,
Arc::new(LettreMailer),
inspector,
updater,
cluster,
|r| r,
)
Self::with_adapters(cfg, pool, adapters, |r| r)
}
/// Wiring with replaceable adapters (used by tests and the fake-host mode).
pub fn with_adapters(
cfg: Config,
pool: DbPool,
mailer: Arc<dyn Mailer>,
inspector: Arc<dyn HostInspector>,
updater: Arc<dyn HostUpdater>,
cluster: Arc<dyn ClusterGateway>,
adapters: Adapters,
register_jobs: impl FnOnce(JobRunner) -> JobRunner,
) -> anyhow::Result<Self> {
let Adapters {
mailer,
inspector,
updater,
cluster,
scanner,
} = adapters;
let pool_for_findings = pool.clone();
let cluster_gateway = cluster.clone();
let users = Arc::new(SqliteUsers(pool.clone()));
let hasher = Arc::new(Argon2Hasher);
let auth = AuthService::new(
@ -119,8 +128,17 @@ impl AppState {
inventory: inventory.clone(),
}),
);
let jobs = Arc::new(register_jobs(runner));
let cluster = Arc::new(ClusterService::new(cluster));
let cluster = Arc::new(ClusterService::new(cluster.clone()));
let vulns = Arc::new(VulnerabilityService::new(
scanner,
Arc::new(SqliteFindings(pool_for_findings)),
cluster_gateway,
settings.clone(),
));
let jobs = Arc::new(register_jobs(runner.register(
JobKind::VulnerabilityScan,
Arc::new(VulnerabilityScanJob(vulns.clone())),
)));
Ok(Self {
cfg,
auth: Arc::new(auth),
@ -129,6 +147,7 @@ impl AppState {
jobs,
inventory,
cluster,
vulns,
login_limiter: Arc::new(rate_limit::RateLimiter::new(
10,
std::time::Duration::from_secs(60),
@ -163,6 +182,7 @@ pub fn build_app(state: AppState) -> Router {
.nest("/api/jobs", jobs::router())
.nest("/api/system", system::router())
.nest("/api/cluster", cluster::router())
.nest("/api/vulnerabilities", vulnerabilities::router())
.fallback_service(spa)
.layer(TraceLayer::new_for_http())
.with_state(state)

View File

@ -28,6 +28,8 @@ impl Modify for BearerAuth {
crate::jobs::list, crate::jobs::kinds, crate::jobs::get_one, crate::jobs::run,
crate::system::inventory, crate::system::upgrade,
crate::cluster::overview, crate::cluster::restart, crate::cluster::scale, crate::cluster::set_image,
crate::vulnerabilities::list, crate::vulnerabilities::summary, crate::vulnerabilities::targets, crate::vulnerabilities::set_status,
crate::settings::get_notifications, crate::settings::put_notifications,
),
modifiers(&BearerAuth)
)]

View File

@ -5,6 +5,7 @@ use axum::routing::{get, post, put};
use axum::{Json, Router};
use domain::jobs::JobKind;
use domain::settings::{SmtpSecurity, SmtpSettings};
use domain::vuln::Severity;
use domain::DomainError;
use serde::{Deserialize, Serialize};
use utoipa::ToSchema;
@ -19,6 +20,39 @@ pub fn router() -> Router<AppState> {
.route("/smtp/test", post(test_smtp))
.route("/schedules", get(list_schedules))
.route("/schedules/{kind}", put(put_schedule))
.route(
"/notifications",
get(get_notifications).put(put_notifications),
)
}
#[derive(Serialize, Deserialize, ToSchema)]
pub struct NotificationSettings {
#[schema(value_type = String, example = "high")]
pub min_severity: Severity,
}
#[utoipa::path(get, path = "/api/settings/notifications", tag = "settings", security(("bearer" = [])), responses((status = 200, body = NotificationSettings)))]
async fn get_notifications(
State(state): State<AppState>,
_: AuthUser,
) -> Result<Json<NotificationSettings>, ApiError> {
Ok(Json(NotificationSettings {
min_severity: state.settings.notify_min_severity().await?,
}))
}
#[utoipa::path(put, path = "/api/settings/notifications", tag = "settings", security(("bearer" = [])), request_body = NotificationSettings, responses((status = 204)))]
async fn put_notifications(
State(state): State<AppState>,
_: AdminUser,
Json(req): Json<NotificationSettings>,
) -> Result<StatusCode, ApiError> {
state
.settings
.set_notify_min_severity(req.min_severity)
.await?;
Ok(StatusCode::NO_CONTENT)
}
#[derive(Serialize, ToSchema)]

View File

@ -80,16 +80,14 @@ async fn build_test_app_with(cfg: Config) -> Router {
let pool = infrastructure::connect(&cfg.database_url)
.await
.expect("db");
let state = AppState::with_adapters(
cfg,
pool,
Arc::new(RecordingMailer),
Arc::new(infrastructure::FakeHostInspector),
Arc::new(infrastructure::FakeHostUpdater),
Arc::new(infrastructure::FakeClusterGateway),
register_test_jobs,
)
.expect("state");
let adapters = crate::Adapters {
mailer: Arc::new(RecordingMailer),
inspector: Arc::new(infrastructure::FakeHostInspector),
updater: Arc::new(infrastructure::FakeHostUpdater),
cluster: Arc::new(infrastructure::FakeClusterGateway),
scanner: Arc::new(infrastructure::FakeScanner),
};
let state = AppState::with_adapters(cfg, pool, adapters, register_test_jobs).expect("state");
state.bootstrap().await.expect("bootstrap");
build_app(state)
}

View File

@ -0,0 +1,96 @@
//! /api/vulnerabilities: findings, summary, status changes.
use application::vuln_service::Summary;
use axum::extract::{Path, Query, State};
use axum::routing::{get, post};
use axum::{Json, Router};
use domain::vuln::{Finding, FindingFilter, FindingStatus, Severity};
use serde::{Deserialize, Serialize};
use utoipa::ToSchema;
use uuid::Uuid;
use crate::error::ApiError;
use crate::extract::{AdminUser, AuthUser};
use crate::AppState;
pub fn router() -> Router<AppState> {
Router::new()
.route("/", get(list))
.route("/summary", get(summary))
.route("/targets", get(targets))
.route("/{id}/status", post(set_status))
}
#[derive(Deserialize)]
pub struct ListQuery {
pub min_severity: Option<String>,
pub target: Option<String>,
pub status: Option<String>,
#[serde(default)]
pub include_fixed: bool,
}
#[utoipa::path(get, path = "/api/vulnerabilities", tag = "vulnerabilities", security(("bearer" = [])),
params(("min_severity" = Option<String>, Query), ("target" = Option<String>, Query), ("status" = Option<String>, Query), ("include_fixed" = Option<bool>, Query)),
responses((status = 200, body = Vec<Object>)))]
async fn list(
State(state): State<AppState>,
_: AuthUser,
Query(q): Query<ListQuery>,
) -> Result<Json<Vec<Finding>>, ApiError> {
let filter = FindingFilter {
min_severity: q.min_severity.as_deref().map(Severity::parse),
target: q.target,
status: q.status.as_deref().and_then(FindingStatus::parse),
include_fixed: q.include_fixed,
};
Ok(Json(state.vulns.list(filter).await?))
}
#[derive(Serialize, ToSchema)]
pub struct SummaryResponse {
#[serde(flatten)]
#[schema(value_type = Object)]
pub summary: Summary,
/// Scanner version, or the error why it is unavailable.
pub scanner: String,
}
#[utoipa::path(get, path = "/api/vulnerabilities/summary", tag = "vulnerabilities", security(("bearer" = [])), responses((status = 200, body = SummaryResponse)))]
async fn summary(
State(state): State<AppState>,
_: AuthUser,
) -> Result<Json<SummaryResponse>, ApiError> {
let scanner = match state.vulns.scanner_version().await {
Ok(v) => v,
Err(e) => format!("unavailable: {e}"),
};
Ok(Json(SummaryResponse {
summary: state.vulns.summary().await?,
scanner,
}))
}
#[utoipa::path(get, path = "/api/vulnerabilities/targets", tag = "vulnerabilities", security(("bearer" = [])), responses((status = 200, body = Vec<String>)))]
async fn targets(
State(state): State<AppState>,
_: AuthUser,
) -> Result<Json<Vec<String>>, ApiError> {
Ok(Json(state.vulns.targets().await?))
}
#[derive(Deserialize, ToSchema)]
pub struct StatusRequest {
#[schema(value_type = String, example = "acknowledged")]
pub status: FindingStatus,
}
#[utoipa::path(post, path = "/api/vulnerabilities/{id}/status", tag = "vulnerabilities", security(("bearer" = [])), request_body = StatusRequest,
responses((status = 200, body = Object), (status = 404), (status = 422)))]
async fn set_status(
State(state): State<AppState>,
_: AdminUser,
Path(id): Path<Uuid>,
Json(req): Json<StatusRequest>,
) -> Result<Json<Finding>, ApiError> {
Ok(Json(state.vulns.set_status(id, req.status).await?))
}

View File

@ -196,12 +196,13 @@ async fn new_findings_at_or_above_threshold_trigger_one_mail() {
]),
)]);
svc.scan(&VecLog::default()).await.unwrap();
let sent = f.mailer.0.lock().unwrap();
assert_eq!(sent.len(), 1);
assert!(sent[0].1.contains("1 new"), "{}", sent[0].1);
assert!(sent[0].2.contains("CVE-3"));
assert!(!sent[0].2.contains("CVE-4"), "below threshold");
drop(sent);
{
let sent = f.mailer.0.lock().unwrap();
assert_eq!(sent.len(), 1);
assert!(sent[0].1.contains("1 new"), "{}", sent[0].1);
assert!(sent[0].2.contains("CVE-3"));
assert!(!sent[0].2.contains("CVE-4"), "below threshold");
}
f.settings
.set_notify_min_severity(Severity::Low)

View File

@ -1,11 +1,14 @@
//! Vulnerability scanning: diffs scanner results against stored findings,
//! notifies about new ones by mail.
use std::collections::HashSet;
use std::sync::Arc;
use async_trait::async_trait;
use domain::ports::{ClusterGateway, FindingRepository, VulnerabilityScanner};
use chrono::{DateTime, Utc};
use domain::ports::{ClusterGateway, FindingRepository, LineSink, VulnerabilityScanner};
use domain::vuln::{
Finding, FindingFilter, FindingStatus, ScanReport, Severity, SeverityCounts, TargetKind,
Finding, FindingFilter, FindingStatus, RawFinding, ScanReport, Severity, SeverityCounts,
TargetKind,
};
use domain::DomainError;
use uuid::Uuid;
@ -14,10 +17,30 @@ use crate::jobs::{JobHandler, JobLog};
use crate::SettingsService;
pub struct VulnerabilityService {
pub(crate) scanner: Arc<dyn VulnerabilityScanner>,
pub(crate) findings: Arc<dyn FindingRepository>,
pub(crate) cluster: Arc<dyn ClusterGateway>,
pub(crate) settings: Arc<SettingsService>,
scanner: Arc<dyn VulnerabilityScanner>,
findings: Arc<dyn FindingRepository>,
cluster: Arc<dyn ClusterGateway>,
settings: Arc<SettingsService>,
}
/// Key of the notification threshold setting.
pub const KEY_NOTIFY_MIN_SEVERITY: &str = "vuln.notify_min_severity";
pub const DEFAULT_NOTIFY_MIN_SEVERITY: Severity = Severity::High;
#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
pub struct Summary {
pub total: SeverityCounts,
pub os: SeverityCounts,
pub images: SeverityCounts,
pub last_scan: Option<DateTime<Utc>>,
}
struct ChannelSink(tokio::sync::mpsc::UnboundedSender<String>);
impl LineSink for ChannelSink {
fn line(&self, text: &str) {
let _ = self.0.send(text.to_string());
}
}
impl VulnerabilityService {
@ -35,50 +58,242 @@ impl VulnerabilityService {
}
}
/// Scan the OS and all cluster images, persist the diff, notify about new findings.
pub async fn scan(&self, _log: &dyn JobLog) -> Result<ScanReport, DomainError> {
todo!()
pub async fn scanner_version(&self) -> Result<String, DomainError> {
self.scanner.version().await
}
pub async fn list(&self, _filter: FindingFilter) -> Result<Vec<Finding>, DomainError> {
todo!()
/// Scan the OS and all cluster images, persist the diff, notify about new findings.
pub async fn scan(&self, log: &dyn JobLog) -> Result<ScanReport, DomainError> {
let mut targets: Vec<(TargetKind, String)> = vec![(TargetKind::Os, "os".into())];
match self.cluster.overview().await {
Ok(o) => targets.extend(o.images().into_iter().map(|i| (TargetKind::Image, i))),
Err(e) => {
log.line(&format!("cluster unavailable, scanning OS only: {e}"))
.await
}
}
let mut report = ScanReport {
targets: targets.iter().map(|(_, t)| t.clone()).collect(),
..Default::default()
};
let now = Utc::now();
for (kind, target) in &targets {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<String>();
let sink = ChannelSink(tx);
let scan = async {
let r = match kind {
TargetKind::Os => self.scanner.scan_os(&sink).await,
TargetKind::Image => self.scanner.scan_image(target, &sink).await,
};
drop(sink);
r
};
let drain = async {
while let Some(l) = rx.recv().await {
log.line(&l).await;
}
};
let (result, _) = tokio::join!(scan, drain);
match result {
Ok(raw) => {
let (new, fixed) = self.apply(*kind, target, raw, now).await?;
log.line(&format!("{target}: {} new, {fixed} fixed", new.len()))
.await;
report.new_findings.extend(new);
report.fixed += fixed;
}
Err(e) => {
log.line(&format!(
"{target}: scan failed, keeping previous findings: {e}"
))
.await;
report.failed_targets.push(target.clone());
}
}
}
report.total_open = self.findings.counts(None).await?.total();
self.notify(&report, log).await;
Ok(report)
}
/// Diff scanner output against the stored active findings of one target.
async fn apply(
&self,
kind: TargetKind,
target: &str,
raw: Vec<RawFinding>,
now: DateTime<Utc>,
) -> Result<(Vec<Finding>, usize), DomainError> {
let existing = self.findings.active_by_target(target).await?;
let mut seen: HashSet<String> = HashSet::new();
let mut still_present = Vec::new();
let mut new = Vec::new();
for r in raw {
if !seen.insert(r.key()) {
continue;
}
match existing.iter().find(|f| f.raw.key() == r.key()) {
Some(f) => still_present.push(f.id),
None => {
let f = Finding {
id: Uuid::new_v4(),
target_kind: kind,
target: target.into(),
raw: r,
status: FindingStatus::Open,
first_seen: now,
last_seen: now,
};
self.findings.insert(&f).await?;
new.push(f);
}
}
}
self.findings.touch(&still_present, now).await?;
let mut fixed = 0;
for f in existing.iter().filter(|f| !still_present.contains(&f.id)) {
self.findings.set_status(f.id, FindingStatus::Fixed).await?;
fixed += 1;
}
Ok((new, fixed))
}
async fn notify(&self, report: &ScanReport, log: &dyn JobLog) {
let min = match self.settings.notify_min_severity().await {
Ok(m) => m,
Err(e) => {
log.line(&format!("notification settings unavailable: {e}"))
.await;
return;
}
};
let relevant: Vec<&Finding> = report
.new_findings
.iter()
.filter(|f| f.raw.severity >= min)
.collect();
if relevant.is_empty() {
return;
}
let plural = if relevant.len() == 1 { "y" } else { "ies" };
let subject = format!(
"[SoftVisor Monitoring] {} new vulnerabilit{plural} (>= {})",
relevant.len(),
min.as_str()
);
let mut body = format!(
"The vulnerability scan found {} new finding(s) at or above severity {}:\n\n",
relevant.len(),
min.as_str()
);
for f in &relevant {
let fixed = f
.raw
.fixed_version
.as_ref()
.map(|v| format!(" (fixed in {v})"))
.unwrap_or_default();
body.push_str(&format!(
"- {} [{}] {} {} in {}{fixed}\n {}\n",
f.raw.cve_id,
f.raw.severity.as_str(),
f.raw.package,
f.raw.installed_version,
f.target,
f.raw.url
));
}
body.push_str(&format!("\nTotal open findings: {}\n", report.total_open));
match self.settings.send_mail(None, &subject, &body).await {
Ok(()) => {
log.line(&format!(
"notification mail sent for {} finding(s)",
relevant.len()
))
.await
}
Err(DomainError::Validation(_)) => log.line("no mail sent: smtp not configured").await,
Err(e) => log.line(&format!("notification mail failed: {e}")).await,
}
}
pub async fn list(&self, filter: FindingFilter) -> Result<Vec<Finding>, DomainError> {
self.findings.list(&filter).await
}
pub async fn targets(&self) -> Result<Vec<String>, DomainError> {
let mut t: Vec<String> = self
.findings
.list(&FindingFilter::default())
.await?
.into_iter()
.map(|f| f.target)
.collect();
t.sort();
t.dedup();
Ok(t)
}
pub async fn summary(&self) -> Result<Summary, DomainError> {
todo!()
let all = self
.findings
.list(&FindingFilter {
include_fixed: true,
..Default::default()
})
.await?;
Ok(Summary {
total: self.findings.counts(None).await?,
os: self.findings.counts(Some(TargetKind::Os)).await?,
images: self.findings.counts(Some(TargetKind::Image)).await?,
last_scan: all.iter().map(|f| f.last_seen).max(),
})
}
/// Only open <-> acknowledged transitions are allowed by users.
pub async fn set_status(
&self,
_id: Uuid,
_status: FindingStatus,
id: Uuid,
status: FindingStatus,
) -> Result<Finding, DomainError> {
todo!()
if status == FindingStatus::Fixed {
return Err(DomainError::Validation(
"findings are marked fixed by the scanner only".into(),
));
}
let f = self.findings.get(id).await?.ok_or(DomainError::NotFound)?;
if f.status == FindingStatus::Fixed {
return Err(DomainError::Validation("finding is already fixed".into()));
}
self.findings.set_status(id, status).await?;
Ok(Finding { status, ..f })
}
pub async fn scanner_version(&self) -> Result<String, DomainError> {
self.scanner.version().await
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
pub struct Summary {
pub total: SeverityCounts,
pub os: SeverityCounts,
pub images: SeverityCounts,
pub last_scan: Option<chrono::DateTime<chrono::Utc>>,
}
pub struct VulnerabilityScanJob(pub Arc<VulnerabilityService>);
#[async_trait]
impl JobHandler for VulnerabilityScanJob {
async fn run(&self, _params: Option<String>, _log: &dyn JobLog) -> Result<(), String> {
todo!()
async fn run(&self, _params: Option<String>, log: &dyn JobLog) -> Result<(), String> {
let version = self
.0
.scanner_version()
.await
.map_err(|e| format!("scanner not available: {e}"))?;
log.line(&format!("scanner version {version}")).await;
let r = self.0.scan(log).await.map_err(|e| e.to_string())?;
log.line(&format!(
"scanned {} target(s), {} failed: {} new, {} fixed, {} open in total",
r.targets.len(),
r.failed_targets.len(),
r.new_findings.len(),
r.fixed,
r.total_open
))
.await;
if !r.targets.is_empty() && r.failed_targets.len() == r.targets.len() {
return Err("all targets failed to scan".into());
}
Ok(())
}
}
/// Key of the notification threshold setting.
pub const KEY_NOTIFY_MIN_SEVERITY: &str = "vuln.notify_min_severity";
pub const DEFAULT_NOTIFY_MIN_SEVERITY: Severity = Severity::High;