diff --git a/backend/crates/api/src/cluster.rs b/backend/crates/api/src/cluster.rs index a1fdc49..90099ae 100644 --- a/backend/crates/api/src/cluster.rs +++ b/backend/crates/api/src/cluster.rs @@ -4,6 +4,7 @@ use axum::http::StatusCode; use axum::routing::{get, post}; use axum::{Json, Router}; use domain::cluster::{ClusterOverview, WorkloadKind, WorkloadRef}; +use domain::jobs::JobKind; use domain::DomainError; use serde::{Deserialize, Serialize}; use utoipa::ToSchema; @@ -18,6 +19,7 @@ pub fn router() -> Router { .route("/workloads/{ns}/{kind}/{name}/restart", post(restart)) .route("/workloads/{ns}/{kind}/{name}/scale", post(scale)) .route("/workloads/{ns}/{kind}/{name}/image", post(set_image)) + .route("/images/check", post(check_images)) } #[derive(Serialize, ToSchema)] @@ -75,6 +77,26 @@ async fn scale( Ok(StatusCode::NO_CONTENT) } +#[derive(Deserialize, ToSchema)] +pub struct CheckImagesRequest { + /// One image, or all images running on the cluster when omitted. + pub image: Option, +} + +#[utoipa::path(post, path = "/api/cluster/images/check", tag = "cluster", security(("bearer" = [])), request_body = CheckImagesRequest, + responses((status = 202, body = crate::jobs::JobRunDto), (status = 409)))] +async fn check_images( + State(state): State, + AdminUser(admin): AdminUser, + Json(req): Json, +) -> Result<(StatusCode, Json), ApiError> { + let run = state + .jobs + .start(JobKind::ImageUpdateCheck, req.image, &admin.email) + .await?; + Ok((StatusCode::ACCEPTED, Json(run.into()))) +} + #[derive(Deserialize, ToSchema)] pub struct ImageRequest { pub image: String, diff --git a/backend/crates/api/src/lib.rs b/backend/crates/api/src/lib.rs index dc80f80..952e229 100644 --- a/backend/crates/api/src/lib.rs +++ b/backend/crates/api/src/lib.rs @@ -20,24 +20,24 @@ use std::sync::Arc; use application::scheduler::Scheduler; use application::{ - AuthService, BackupDeps, BackupJob, BackupService, ClusterService, InventoryService, JobRunner, - PackageRefreshJob, PackageUpgradeJob, SettingsService, UserService, VulnerabilityScanJob, - VulnerabilityService, + AuthService, BackupDeps, BackupJob, BackupService, ClusterService, ImageUpdateCheckJob, + ImageUpdateService, InventoryService, JobRunner, PackageRefreshJob, PackageUpgradeJob, + SettingsService, UserService, VulnerabilityScanJob, VulnerabilityService, }; use axum::{routing::get, Json, Router}; use domain::jobs::JobKind; use domain::ports::Mailer; use domain::ports::{ BackupCollector, BackupStorage, ClusterGateway, FileEncryptor, HostInspector, HostUpdater, - VulnerabilityScanner, + ImageRegistry, VulnerabilityScanner, }; use infrastructure::{ - AesGcmCipher, Argon2Hasher, CommandBackupStorage, DbPool, DebianInspector, DebianUpdater, - DirBackupStorage, FakeBackupCollector, FakeClusterGateway, FakeHostInspector, FakeHostUpdater, - FakeScanner, JwtIssuer, KubeBackupCollector, KubeGateway, LettreMailer, OpensslEncryptor, - SqliteAuditLog, SqliteBackupRecords, SqliteBackupStrategies, SqliteBackupTargets, - SqliteFindings, SqliteInventory, SqliteJobRuns, SqliteRefreshTokens, SqliteSettings, - SqliteUsers, SystemCommandRunner, TrivyScanner, + AesGcmCipher, Argon2Hasher, CommandBackupStorage, CurlImageRegistry, DbPool, DebianInspector, + DebianUpdater, DirBackupStorage, FakeBackupCollector, FakeClusterGateway, FakeHostInspector, + FakeHostUpdater, FakeImageRegistry, FakeScanner, JwtIssuer, KubeBackupCollector, KubeGateway, + LettreMailer, OpensslEncryptor, SqliteAuditLog, SqliteBackupRecords, SqliteBackupStrategies, + SqliteBackupTargets, SqliteFindings, SqliteImageUpdates, SqliteInventory, SqliteJobRuns, + SqliteRefreshTokens, SqliteSettings, SqliteUsers, SystemCommandRunner, TrivyScanner, }; use tower_http::services::{ServeDir, ServeFile}; use tower_http::trace::TraceLayer; @@ -51,6 +51,7 @@ pub struct Adapters { pub updater: Arc, pub cluster: Arc, pub scanner: Arc, + pub registry: Arc, pub storage: Arc, pub collector: Arc, pub encryptor: Arc, @@ -67,6 +68,7 @@ pub struct AppState { pub cluster: Arc, pub vulns: Arc, pub backups: Arc, + pub image_updates: Arc, pub login_limiter: Arc, } @@ -81,6 +83,7 @@ impl AppState { updater: Arc::new(FakeHostUpdater), cluster: Arc::new(FakeClusterGateway), scanner: Arc::new(FakeScanner), + registry: Arc::new(FakeImageRegistry), storage: Arc::new(DirBackupStorage::new(cfg.work_dir.join("fake-remote"))), collector: Arc::new(FakeBackupCollector), encryptor: Arc::new(OpensslEncryptor::new(runner.clone())), @@ -91,6 +94,7 @@ impl AppState { inspector: Arc::new(DebianInspector::new(runner.clone())), updater: Arc::new(DebianUpdater::new(runner.clone())), cluster: Arc::new(KubeGateway::new(cfg.kubeconfig.clone())), + registry: Arc::new(CurlImageRegistry::new(runner.clone())), scanner: Arc::new(match &cfg.containerd_socket { Some(sock) => TrivyScanner::new(runner.clone()).with_containerd(sock), None => TrivyScanner::new(runner.clone()), @@ -119,6 +123,7 @@ impl AppState { updater, cluster, scanner, + registry, storage, collector, encryptor, @@ -168,19 +173,31 @@ impl AppState { ); let login_rate_limit = cfg.login_rate_limit; let cluster = Arc::new(ClusterService::new(cluster.clone())); + let findings_repo = Arc::new(SqliteFindings(pool_for_findings.clone())); let vulns = Arc::new(VulnerabilityService::new( - scanner, - Arc::new(SqliteFindings(pool_for_findings)), - cluster_gateway, + scanner.clone(), + findings_repo.clone(), + cluster_gateway.clone(), settings.clone(), )); + let image_updates = Arc::new(ImageUpdateService::new( + registry, + scanner, + findings_repo, + Arc::new(SqliteImageUpdates(pool_for_findings)), + cluster_gateway, + )); let jobs = Arc::new(register_jobs( runner .register( JobKind::VulnerabilityScan, Arc::new(VulnerabilityScanJob(vulns.clone())), ) - .register(JobKind::Backup, Arc::new(BackupJob(backups.clone()))), + .register(JobKind::Backup, Arc::new(BackupJob(backups.clone()))) + .register( + JobKind::ImageUpdateCheck, + Arc::new(ImageUpdateCheckJob(image_updates.clone())), + ), )); Ok(Self { cfg, @@ -192,6 +209,7 @@ impl AppState { cluster, vulns, backups, + image_updates, login_limiter: Arc::new(rate_limit::RateLimiter::new( login_rate_limit, std::time::Duration::from_secs(60), diff --git a/backend/crates/api/src/openapi.rs b/backend/crates/api/src/openapi.rs index 46972e3..61811b3 100644 --- a/backend/crates/api/src/openapi.rs +++ b/backend/crates/api/src/openapi.rs @@ -27,7 +27,7 @@ impl Modify for BearerAuth { crate::settings::get_smtp, crate::settings::put_smtp, crate::settings::test_smtp, crate::settings::list_schedules, crate::settings::put_schedule, 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::cluster::overview, crate::cluster::restart, crate::cluster::scale, crate::cluster::set_image, crate::cluster::check_images, crate::vulnerabilities::list, crate::vulnerabilities::groups, crate::vulnerabilities::summary, crate::vulnerabilities::targets, crate::vulnerabilities::set_status, crate::settings::get_notifications, crate::settings::put_notifications, crate::backups::list_targets, crate::backups::get_target, crate::backups::create_target, crate::backups::update_target, diff --git a/backend/crates/api/src/test_support.rs b/backend/crates/api/src/test_support.rs index 92247d1..8529767 100644 --- a/backend/crates/api/src/test_support.rs +++ b/backend/crates/api/src/test_support.rs @@ -90,6 +90,7 @@ async fn build_test_app_with(cfg: Config) -> Router { updater: Arc::new(infrastructure::FakeHostUpdater), cluster: Arc::new(infrastructure::FakeClusterGateway), scanner: Arc::new(infrastructure::FakeScanner), + registry: Arc::new(infrastructure::FakeImageRegistry), storage: Arc::new(infrastructure::DirBackupStorage::new( cfg.work_dir.join("remote"), )), diff --git a/backend/crates/api/src/vulnerabilities.rs b/backend/crates/api/src/vulnerabilities.rs index 1b35969..36f20cb 100644 --- a/backend/crates/api/src/vulnerabilities.rs +++ b/backend/crates/api/src/vulnerabilities.rs @@ -3,7 +3,7 @@ use application::vuln_service::{Summary, TargetSummary}; use axum::extract::{Path, Query, State}; use axum::routing::{get, post}; use axum::{Json, Router}; -use domain::vuln::{Finding, FindingFilter, FindingGroup, FindingStatus, Severity, TargetKind}; +use domain::vuln::{Finding, FindingFilter, FindingStatus, Severity, TargetKind}; use domain::DomainError; use serde::{Deserialize, Serialize}; use utoipa::ToSchema; @@ -74,8 +74,27 @@ async fn groups( State(state): State, _: AuthUser, Query(q): Query, -) -> Result>, ApiError> { - Ok(Json(state.vulns.groups(q.into_filter()?).await?)) +) -> Result, ApiError> { + let filter = q.into_filter()?; + // container groups carry the workloads that run the image, host groups do not + if filter.target_kind == Some(TargetKind::Image) { + let groups = state.vulns.image_groups(filter).await?; + let mut out = Vec::with_capacity(groups.len()); + for g in groups { + // the last update check of this image, if one has run + let update = state.image_updates.stored(&g.group.key).await?; + let mut value = + serde_json::to_value(g).map_err(|e| DomainError::Storage(e.to_string()))?; + value["update"] = + serde_json::to_value(update).map_err(|e| DomainError::Storage(e.to_string()))?; + out.push(value); + } + return Ok(Json(serde_json::Value::Array(out))); + } + let groups = state.vulns.groups(filter).await?; + Ok(Json( + serde_json::to_value(groups).map_err(|e| DomainError::Storage(e.to_string()))?, + )) } #[derive(Serialize, ToSchema)] diff --git a/backend/crates/api/tests/vulnerabilities.rs b/backend/crates/api/tests/vulnerabilities.rs index ce78f8d..7ce2fea 100644 --- a/backend/crates/api/tests/vulnerabilities.rs +++ b/backend/crates/api/tests/vulnerabilities.rs @@ -330,3 +330,123 @@ async fn findings_are_rolled_up_per_package_and_image() { StatusCode::UNAUTHORIZED ); } + +#[tokio::test] +async fn image_groups_link_to_the_workloads_running_them() { + let app = test_app_with_admin().await; + let token = common::login(&app, ADMIN, PW).await.access; + run_scan(&app, &token).await; + + let res = get( + &app, + "/api/vulnerabilities/groups?scope=container&min_severity=unknown", + Some(&token), + ) + .await; + let groups = res.json.as_array().unwrap(); + let gitea = groups + .iter() + .find(|g| g["key"].as_str().unwrap().contains("gitea")) + .unwrap(); + assert_eq!(gitea["running"], true); + let workloads = gitea["workloads"].as_array().unwrap(); + assert_eq!(workloads.len(), 1); + assert_eq!(workloads[0]["namespace"], "gitea"); + assert_eq!(workloads[0]["kind"], "deployment"); + assert_eq!(workloads[0]["name"], "gitea"); + + // host groups do not carry cluster usage + let host = get( + &app, + "/api/vulnerabilities/groups?scope=host&min_severity=unknown", + Some(&token), + ) + .await; + assert!(host + .json + .as_array() + .unwrap() + .iter() + .all(|g| g.get("workloads").is_none())); +} + +#[tokio::test] +async fn an_image_update_check_suggests_a_newer_tag_that_fixes_findings() { + let app = test_app_with_admin().await; + let token = common::login(&app, ADMIN, PW).await.access; + run_scan(&app, &token).await; + + // nothing is known before a check + let groups = get( + &app, + "/api/vulnerabilities/groups?scope=container&min_severity=unknown", + Some(&token), + ) + .await; + let gitea = groups + .json + .as_array() + .unwrap() + .iter() + .find(|g| g["key"].as_str().unwrap().contains("gitea")) + .unwrap() + .clone(); + assert!(gitea["update"].is_null()); + + let run = post( + &app, + "/api/cluster/images/check", + json!({"image": gitea["key"]}), + Some(&token), + ) + .await; + assert_eq!(run.status, StatusCode::ACCEPTED, "{}", run.json); + assert_eq!(run.json["kind"], "image_update_check"); + let id = run.json["id"].as_str().unwrap().to_string(); + for _ in 0..100 { + let r = get(&app, &format!("/api/jobs/{id}"), Some(&token)).await; + if r.json["status"] != "running" { + assert_eq!(r.json["status"], "success", "{}", r.json["log"]); + break; + } + tokio::time::sleep(std::time::Duration::from_millis(30)).await; + } + + let groups = get( + &app, + "/api/vulnerabilities/groups?scope=container&min_severity=unknown", + Some(&token), + ) + .await; + let gitea = groups + .json + .as_array() + .unwrap() + .iter() + .find(|g| g["key"].as_str().unwrap().contains("gitea")) + .unwrap() + .clone(); + let update = &gitea["update"]; + assert!( + update["candidate"].as_str().unwrap().ends_with(":9.9.9"), + "{update}" + ); + assert!(update["fixed"] + .as_array() + .unwrap() + .contains(&json!("CVE-2024-24790"))); + assert_eq!(update["fixed_counts"]["critical"], 1); + assert!(update["checked_at"].is_string()); + + // only admins may start a check + post(&app, "/api/users", json!({"email": "u@x.de", "display_name": "U", "password": "user-password-123", "role": "user"}), Some(&token)).await; + let user = common::login(&app, "u@x.de", "user-password-123") + .await + .access; + assert_eq!( + post(&app, "/api/cluster/images/check", json!({}), Some(&user)) + .await + .status, + StatusCode::FORBIDDEN + ); +} diff --git a/backend/crates/application/src/image_update_service.rs b/backend/crates/application/src/image_update_service.rs new file mode 100644 index 0000000..0d7ff4a --- /dev/null +++ b/backend/crates/application/src/image_update_service.rs @@ -0,0 +1,167 @@ +//! Suggests a newer tag for a running image and verifies that it fixes findings. +use std::collections::HashSet; +use std::sync::Arc; + +use async_trait::async_trait; +use chrono::Utc; +use domain::image::{newest_tag, ImageRef, ImageUpdate}; +use domain::ports::{ + ClusterGateway, FindingRepository, ImageRegistry, ImageUpdateRepository, LineSink, + VulnerabilityScanner, +}; +use domain::vuln::{FindingStatus, SeverityCounts}; +use domain::DomainError; + +use crate::jobs::{JobHandler, JobLog}; + +pub struct ImageUpdateService { + registry: Arc, + scanner: Arc, + findings: Arc, + updates: Arc, + cluster: Arc, +} + +struct ChannelSink(tokio::sync::mpsc::UnboundedSender); + +impl LineSink for ChannelSink { + fn line(&self, text: &str) { + let _ = self.0.send(text.to_string()); + } +} + +impl ImageUpdateService { + pub fn new( + registry: Arc, + scanner: Arc, + findings: Arc, + updates: Arc, + cluster: Arc, + ) -> Self { + Self { + registry, + scanner, + findings, + updates, + cluster, + } + } + + /// The newest tag of the image's repository that is newer than the running one. + pub async fn candidate(&self, image: &str) -> Result, DomainError> { + let parsed = ImageRef::parse(image) + .ok_or_else(|| DomainError::Validation(format!("cannot parse image '{image}'")))?; + let tags = self.registry.tags(image).await?; + Ok(newest_tag(&parsed.tag, &tags).map(|t| parsed.with_tag(&t))) + } + + pub async fn stored(&self, image: &str) -> Result, DomainError> { + self.updates.get(image).await + } + + pub async fn all(&self) -> Result, DomainError> { + self.updates.list().await + } + + /// Scan the newest available tag and record which of the open findings it fixes. + pub async fn check( + &self, + image: &str, + log: &dyn JobLog, + ) -> Result, DomainError> { + let Some(candidate) = self.candidate(image).await? else { + log.line(&format!("{image}: already on the newest tag")) + .await; + return Ok(None); + }; + log.line(&format!("{image}: candidate {candidate}, scanning it")) + .await; + + let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); + let sink = ChannelSink(tx); + let scan = async { + let r = self.scanner.scan_image(&candidate, &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); + let candidate_findings = result?; + + let still_there: HashSet = candidate_findings.iter().map(|f| f.key()).collect(); + let current = self.findings.active_by_target(image).await?; + let mut fixed_counts = SeverityCounts::default(); + let mut fixed: Vec = Vec::new(); + for f in current + .iter() + .filter(|f| f.status != FindingStatus::Fixed && !still_there.contains(&f.raw.key())) + { + fixed_counts.add(f.raw.severity); + fixed.push(f.raw.cve_id.clone()); + } + fixed.sort(); + fixed.dedup(); + + let update = ImageUpdate { + image: image.to_string(), + candidate, + checked_at: Utc::now(), + fixed_counts, + candidate_total: candidate_findings.len(), + fixed, + }; + self.updates.upsert(&update).await?; + log.line(&format!( + "{image} → {}: fixes {} finding(s) ({} critical, {} high), {} remain", + update.candidate, + update.fixed.len(), + update.fixed_counts.critical, + update.fixed_counts.high, + update.candidate_total + )) + .await; + Ok(Some(update)) + } + + /// Check every image that currently runs on the cluster. + pub async fn check_running(&self, log: &dyn JobLog) -> Result { + let images = self.cluster.overview().await?.images(); + log.line(&format!( + "checking {} running image(s) for newer tags", + images.len() + )) + .await; + let mut found = 0; + for image in images { + match self.check(&image, log).await { + Ok(Some(_)) => found += 1, + Ok(None) => {} + Err(e) => log.line(&format!("{image}: check failed: {e}")).await, + } + } + Ok(found) + } +} + +/// Job handler for `JobKind::ImageUpdateCheck`; params = one image, or empty for all running ones. +pub struct ImageUpdateCheckJob(pub Arc); + +#[async_trait] +impl JobHandler for ImageUpdateCheckJob { + async fn run(&self, params: Option, log: &dyn JobLog) -> Result<(), String> { + match params.as_deref().map(str::trim).filter(|p| !p.is_empty()) { + Some(image) => { + self.0.check(image, log).await.map_err(|e| e.to_string())?; + } + None => { + let n = self.0.check_running(log).await.map_err(|e| e.to_string())?; + log.line(&format!("{n} image(s) have a newer tag")).await; + } + } + Ok(()) + } +} diff --git a/backend/crates/application/src/lib.rs b/backend/crates/application/src/lib.rs index 15d15da..c61bda9 100644 --- a/backend/crates/application/src/lib.rs +++ b/backend/crates/application/src/lib.rs @@ -2,6 +2,7 @@ pub mod auth_service; pub mod backup_service; pub mod cluster_service; +pub mod image_update_service; pub mod inventory_service; pub mod jobs; pub mod scheduler; @@ -13,6 +14,7 @@ pub mod vuln_service; pub use auth_service::AuthService; pub use backup_service::{BackupDeps, BackupJob, BackupService}; pub use cluster_service::ClusterService; +pub use image_update_service::{ImageUpdateCheckJob, ImageUpdateService}; pub use inventory_service::{InventoryService, PackageRefreshJob}; pub use jobs::{JobHandler, JobLog, JobRunner}; pub use settings_service::SettingsService; diff --git a/backend/crates/application/src/test_fakes.rs b/backend/crates/application/src/test_fakes.rs index a661da2..9cea9a4 100644 --- a/backend/crates/application/src/test_fakes.rs +++ b/backend/crates/application/src/test_fakes.rs @@ -418,6 +418,7 @@ use domain::ports::ClusterGateway; #[derive(Default)] pub struct MemCluster { pub actions: Mutex>, + pub fail: std::sync::atomic::AtomicBool, } pub fn sample_overview() -> ClusterOverview { @@ -469,6 +470,9 @@ pub fn sample_overview() -> ClusterOverview { #[async_trait] impl ClusterGateway for MemCluster { async fn overview(&self) -> Result { + if self.fail.load(std::sync::atomic::Ordering::SeqCst) { + return Err(DomainError::Unavailable("cluster unreachable".into())); + } Ok(sample_overview()) } async fn restart(&self, w: &WorkloadRef) -> Result<(), DomainError> { @@ -964,3 +968,56 @@ pub fn strategy(name: &str, target_id: Uuid) -> BackupStrategy { enabled: true, } } + +use domain::image::ImageUpdate; +use domain::ports::{ImageRegistry, ImageUpdateRepository}; + +#[derive(Default)] +pub struct MemRegistry { + pub tags: Mutex>, + pub fail: std::sync::atomic::AtomicBool, +} + +#[async_trait] +impl ImageRegistry for MemRegistry { + async fn tags(&self, _image: &str) -> Result, DomainError> { + if self.fail.load(std::sync::atomic::Ordering::SeqCst) { + return Err(DomainError::Unavailable("registry unreachable".into())); + } + Ok(self.tags.lock().unwrap().clone()) + } +} + +#[derive(Default)] +pub struct MemImageUpdates(pub Mutex>); + +#[async_trait] +impl ImageUpdateRepository for MemImageUpdates { + async fn upsert(&self, update: &ImageUpdate) -> Result<(), DomainError> { + let mut v = self.0.lock().unwrap(); + v.retain(|u| u.image != update.image); + v.push(update.clone()); + Ok(()) + } + async fn get(&self, image: &str) -> Result, DomainError> { + Ok(self + .0 + .lock() + .unwrap() + .iter() + .find(|u| u.image == image) + .cloned()) + } + async fn list(&self) -> Result, DomainError> { + Ok(self.0.lock().unwrap().clone()) + } +} + +/// Settings service on in-memory stores, for services that only need it as a dependency. +pub fn settings_service() -> Arc { + Arc::new(crate::SettingsService::new( + Arc::new(MemSettings::default()), + Arc::new(FakeCipher), + Arc::new(MemMailer::default()), + )) +} diff --git a/backend/crates/application/src/tests/image_update_tests.rs b/backend/crates/application/src/tests/image_update_tests.rs new file mode 100644 index 0000000..145f0e6 --- /dev/null +++ b/backend/crates/application/src/tests/image_update_tests.rs @@ -0,0 +1,165 @@ +use std::sync::{Arc, Mutex}; + +use async_trait::async_trait; +use domain::vuln::Severity; +use domain::DomainError; + +use crate::jobs::{JobHandler, JobLog}; +use crate::test_fakes::{raw, FakeScanner, MemCluster, MemFindings, MemImageUpdates, MemRegistry}; +use crate::{ImageUpdateCheckJob, ImageUpdateService, VulnerabilityService}; + +#[derive(Default)] +struct VecLog(Mutex>); +#[async_trait] +impl JobLog for VecLog { + async fn line(&self, text: &str) { + self.0.lock().unwrap().push(text.into()); + } +} + +const GITEA: &str = "gitea/gitea:1.22.3"; + +struct F { + findings: Arc, + updates: Arc, + registry: Arc, +} + +fn fixture(scanner: FakeScanner, tags: &[&str]) -> (F, Arc) { + let findings = Arc::new(MemFindings::default()); + let updates = Arc::new(MemImageUpdates::default()); + let registry = Arc::new(MemRegistry::default()); + *registry.tags.lock().unwrap() = tags.iter().map(|t| t.to_string()).collect(); + let cluster = Arc::new(MemCluster::default()); + let svc = ImageUpdateService::new( + registry.clone(), + Arc::new(scanner), + findings.clone(), + updates.clone(), + cluster, + ); + ( + F { + findings, + updates, + registry, + }, + Arc::new(svc), + ) +} + +/// Scan the running image first so there are open findings to compare against. +async fn seed(findings: Arc, scanner: FakeScanner, cluster: Arc) { + let settings = crate::test_fakes::settings_service(); + let vulns = VulnerabilityService::new(Arc::new(scanner), findings, cluster, settings); + vulns.scan(&VecLog::default()).await.unwrap(); +} + +#[tokio::test] +async fn suggests_the_newest_tag_of_the_same_variant() { + let (_f, svc) = fixture( + FakeScanner::default(), + &["1.22.3", "1.22.4", "1.23.0", "1.24.0-rc1", "latest"], + ); + assert_eq!( + svc.candidate(GITEA).await.unwrap().as_deref(), + Some("gitea/gitea:1.23.0") + ); + let (_f, svc) = fixture(FakeScanner::default(), &["1.22.3"]); + assert_eq!( + svc.candidate(GITEA).await.unwrap(), + None, + "already the newest" + ); +} + +#[tokio::test] +async fn a_check_records_which_findings_the_candidate_fixes() { + // the running image has three flaws, the candidate only one of them + let running = FakeScanner::default().with( + GITEA, + Ok(vec![ + raw("CVE-1", "git", "2.39", Severity::Critical, Some("2.40")), + raw("CVE-2", "curl", "7.8", Severity::High, None), + raw("CVE-3", "zlib", "1.2", Severity::Low, None), + ]), + ); + let candidate = FakeScanner::default().with( + "gitea/gitea:1.23.0", + Ok(vec![raw("CVE-3", "zlib", "1.2", Severity::Low, None)]), + ); + let (f, svc) = fixture(candidate, &["1.22.3", "1.23.0"]); + seed(f.findings.clone(), running, Arc::new(MemCluster::default())).await; + + let log = VecLog::default(); + let update = svc.check(GITEA, &log).await.unwrap().unwrap(); + assert_eq!(update.candidate, "gitea/gitea:1.23.0"); + assert_eq!(update.fixed, vec!["CVE-1", "CVE-2"]); + assert_eq!(update.fixed_counts.critical, 1); + assert_eq!(update.fixed_counts.high, 1); + assert_eq!( + update.candidate_total, 1, + "the candidate still has one finding" + ); + assert_eq!(svc.stored(GITEA).await.unwrap().as_ref(), Some(&update)); + let lines = log.0.lock().unwrap().join("\n"); + assert!(lines.contains("fixes 2 finding(s)"), "{lines}"); + + // a repeated check overwrites the previous result + svc.check(GITEA, &VecLog::default()).await.unwrap(); + assert_eq!(svc.all().await.unwrap().len(), 1); +} + +#[tokio::test] +async fn nothing_is_recorded_when_the_image_is_current() { + let (f, svc) = fixture(FakeScanner::default(), &["1.22.3"]); + let log = VecLog::default(); + assert!(svc.check(GITEA, &log).await.unwrap().is_none()); + assert!(f.updates.0.lock().unwrap().is_empty()); + assert!(log + .0 + .lock() + .unwrap() + .iter() + .any(|l| l.contains("newest tag"))); +} + +#[tokio::test] +async fn a_registry_that_cannot_be_reached_fails_the_check_of_that_image_only() { + let (f, svc) = fixture(FakeScanner::default(), &["1.22.3", "1.23.0"]); + f.registry + .fail + .store(true, std::sync::atomic::Ordering::SeqCst); + assert!(matches!( + svc.check(GITEA, &VecLog::default()).await.unwrap_err(), + DomainError::Unavailable(_) + )); + + // checking every running image keeps going and reports the failures in the log + let log = VecLog::default(); + let found = svc.check_running(&log).await.unwrap(); + assert_eq!(found, 0); + let lines = log.0.lock().unwrap().join("\n"); + assert!(lines.contains("check failed"), "{lines}"); +} + +#[tokio::test] +async fn the_job_checks_one_image_or_all_running_ones() { + let candidate = FakeScanner::default().with("gitea/gitea:1.23.0", Ok(vec![])); + let (f, svc) = fixture(candidate, &["1.22.3", "1.23.0", "16.4.0", "17.0.0"]); + let job = ImageUpdateCheckJob(svc.clone()); + job.run(Some(GITEA.into()), &VecLog::default()) + .await + .unwrap(); + assert_eq!(f.updates.0.lock().unwrap().len(), 1); + + let log = VecLog::default(); + job.run(None, &log).await.unwrap(); + assert!(log + .0 + .lock() + .unwrap() + .iter() + .any(|l| l.contains("running image(s)"))); + assert!(!f.updates.0.lock().unwrap().is_empty()); +} diff --git a/backend/crates/application/src/tests/mod.rs b/backend/crates/application/src/tests/mod.rs index 38d2460..c59b4de 100644 --- a/backend/crates/application/src/tests/mod.rs +++ b/backend/crates/application/src/tests/mod.rs @@ -1,6 +1,7 @@ mod auth_service_tests; mod backup_tests; mod cluster_tests; +mod image_update_tests; mod inventory_tests; mod jobs_tests; mod scheduler_tests; diff --git a/backend/crates/application/src/tests/vuln_tests.rs b/backend/crates/application/src/tests/vuln_tests.rs index d6ca817..1dfd91e 100644 --- a/backend/crates/application/src/tests/vuln_tests.rs +++ b/backend/crates/application/src/tests/vuln_tests.rs @@ -25,6 +25,7 @@ struct F { mailer: Arc, settings: Arc, results: ScanResults, + cluster: Arc, } fn fixture(scanner: FakeScanner) -> (F, VulnerabilityService) { @@ -36,10 +37,11 @@ fn fixture(scanner: FakeScanner) -> (F, VulnerabilityService) { Arc::new(FakeCipher), mailer.clone(), )); + let cluster = Arc::new(MemCluster::default()); let svc = VulnerabilityService::new( Arc::new(scanner), findings.clone(), - Arc::new(MemCluster::default()), + cluster.clone(), settings.clone(), ); ( @@ -48,6 +50,7 @@ fn fixture(scanner: FakeScanner) -> (F, VulnerabilityService) { mailer, settings, results, + cluster, }, svc, ) @@ -533,3 +536,66 @@ async fn findings_of_one_group_can_be_listed() { .unwrap(); assert_eq!(image.len(), 1); } + +#[tokio::test] +async fn container_groups_name_the_workloads_that_run_the_image() { + let scanner = FakeScanner::default() + .with( + "os", + Ok(vec![raw("CVE-1", "openssl", "3.0.1", Severity::High, None)]), + ) + .with( + GITEA, + Ok(vec![raw("CVE-2", "git", "2.39", Severity::High, None)]), + ); + let (_f, svc) = fixture(scanner); + svc.scan(&VecLog::default()).await.unwrap(); + + let images = svc + .image_groups(FindingFilter { + target_kind: Some(TargetKind::Image), + ..Default::default() + }) + .await + .unwrap(); + let gitea = images.iter().find(|g| g.group.key == GITEA).unwrap(); + assert_eq!(gitea.workloads.len(), 1); + assert_eq!(gitea.workloads[0].namespace, "gitea"); + assert_eq!(gitea.workloads[0].name, "gitea"); + assert_eq!( + gitea.workloads[0].kind, + domain::cluster::WorkloadKind::Deployment + ); + assert!(gitea.running, "the image is in use on the cluster"); + + // an image that no workload runs any more is reported as not running + let stale = images.iter().find(|g| g.group.key == PG); + assert!(stale.is_none_or(|g| g.workloads.is_empty())); +} + +#[tokio::test] +async fn image_groups_survive_an_unreachable_cluster() { + let scanner = FakeScanner::default().with( + GITEA, + Ok(vec![raw("CVE-2", "git", "2.39", Severity::High, None)]), + ); + let (f, svc) = fixture(scanner); + svc.scan(&VecLog::default()).await.unwrap(); + f.cluster + .fail + .store(true, std::sync::atomic::Ordering::SeqCst); + + let images = svc + .image_groups(FindingFilter { + target_kind: Some(TargetKind::Image), + ..Default::default() + }) + .await + .unwrap(); + assert_eq!(images.len(), 1, "findings are still listed"); + assert!(images[0].workloads.is_empty()); + assert!( + !images[0].running, + "unknown while the cluster is unreachable" + ); +} diff --git a/backend/crates/application/src/vuln_service.rs b/backend/crates/application/src/vuln_service.rs index 9e76cea..c5b8f53 100644 --- a/backend/crates/application/src/vuln_service.rs +++ b/backend/crates/application/src/vuln_service.rs @@ -5,6 +5,7 @@ use std::sync::Arc; use async_trait::async_trait; use chrono::{DateTime, Utc}; +use domain::cluster::WorkloadRef; use domain::ports::{ClusterGateway, FindingRepository, LineSink, VulnerabilityScanner}; use domain::vuln::{ Finding, FindingFilter, FindingGroup, FindingStatus, RawFinding, ScanReport, Severity, @@ -27,6 +28,17 @@ pub struct VulnerabilityService { pub const KEY_NOTIFY_MIN_SEVERITY: &str = "vuln.notify_min_severity"; pub const DEFAULT_NOTIFY_MIN_SEVERITY: Severity = Severity::High; +/// An image group with the workloads that run the image. +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +pub struct ImageGroup { + #[serde(flatten)] + pub group: FindingGroup, + /// Workloads on the cluster that currently use this image. + pub workloads: Vec, + /// False when no workload uses it, or while the cluster is unreachable. + pub running: bool, +} + /// One scanned target with its number of open findings. #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] pub struct TargetSummary { @@ -242,6 +254,39 @@ impl VulnerabilityService { self.findings.groups(&filter).await } + /// Image groups together with the workloads that currently run them. An unreachable + /// cluster is not an error: the findings are still listed, just without their usage. + pub async fn image_groups( + &self, + filter: FindingFilter, + ) -> Result, DomainError> { + let groups = self.findings.groups(&filter).await?; + let overview = self.cluster.overview().await.ok(); + Ok(groups + .into_iter() + .map(|group| { + let workloads: Vec = overview + .as_ref() + .map(|o| { + o.workloads_running(&group.key) + .into_iter() + .map(|w| WorkloadRef { + namespace: w.namespace.clone(), + kind: w.kind, + name: w.name.clone(), + }) + .collect() + }) + .unwrap_or_default(); + ImageGroup { + running: !workloads.is_empty(), + workloads, + group, + } + }) + .collect()) + } + /// Scanned targets that still have open findings, host first. pub async fn targets(&self) -> Result, DomainError> { let mut by_target: std::collections::HashMap<(TargetKind, String), usize> = diff --git a/backend/crates/domain/src/cluster.rs b/backend/crates/domain/src/cluster.rs index 376bd3c..998a86a 100644 --- a/backend/crates/domain/src/cluster.rs +++ b/backend/crates/domain/src/cluster.rs @@ -70,6 +70,14 @@ pub struct ClusterOverview { } impl ClusterOverview { + /// Workloads that currently run `image`. + pub fn workloads_running(&self, image: &str) -> Vec<&Workload> { + self.workloads + .iter() + .filter(|w| w.containers.iter().any(|c| c.image == image)) + .collect() + } + /// Distinct images across all workloads (input for vulnerability scans). pub fn images(&self) -> Vec { let mut v: Vec = self @@ -84,7 +92,7 @@ impl ClusterOverview { } /// Reference to a workload for actions. -#[derive(Clone, Debug, PartialEq, Eq)] +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct WorkloadRef { pub namespace: String, pub kind: WorkloadKind, diff --git a/backend/crates/domain/src/image.rs b/backend/crates/domain/src/image.rs new file mode 100644 index 0000000..9b34afe --- /dev/null +++ b/backend/crates/domain/src/image.rs @@ -0,0 +1,234 @@ +//! Container image references and picking a newer tag. +use serde::{Deserialize, Serialize}; + +/// A parsed image reference: `[registry/]repository[:tag]`. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct ImageRef { + pub registry: String, + pub repository: String, + pub tag: String, +} + +impl ImageRef { + /// `docker.io/library/nginx:1.2` → registry `docker.io`, repository `library/nginx`, tag `1.2`. + pub fn parse(image: &str) -> Option { + let image = image.split('@').next()?; // ignore a digest + let (head, tag) = match image.rsplit_once(':') { + Some((h, t)) if !t.contains('/') => (h, t.to_string()), + _ => (image, "latest".to_string()), + }; + if head.is_empty() { + return None; + } + let (registry, repository) = match head.split_once('/') { + // a first segment with a dot or a port is a registry, everything else is Docker Hub + Some((first, rest)) + if first.contains('.') || first.contains(':') || first == "localhost" => + { + (first.to_string(), rest.to_string()) + } + Some(_) => ("docker.io".to_string(), head.to_string()), + None => ("docker.io".to_string(), format!("library/{head}")), + }; + Some(ImageRef { + registry, + repository, + tag, + }) + } + + pub fn with_tag(&self, tag: &str) -> String { + let repo = if self.registry == "docker.io" { + self.repository + .strip_prefix("library/") + .unwrap_or(&self.repository) + .to_string() + } else { + format!("{}/{}", self.registry, self.repository) + }; + format!("{repo}:{tag}") + } +} + +/// Result of checking whether a newer tag fixes the findings of a running image. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct ImageUpdate { + /// The image reference as it runs on the cluster. + pub image: String, + /// The newer tag that was checked, e.g. `docker.gitea.com/gitea:1.25.1-rootless`. + pub candidate: String, + pub checked_at: chrono::DateTime, + /// CVE ids that are open on the running image and gone in the candidate. + pub fixed: Vec, + pub fixed_counts: crate::vuln::SeverityCounts, + /// Findings the candidate still has. + pub candidate_total: usize, +} + +/// Numeric parts of a tag plus its suffix, e.g. `1.24.2-rootless` → ([1, 24, 2], "rootless"). +fn version_parts(tag: &str) -> Option<(Vec, String)> { + let core = tag.strip_prefix('v').unwrap_or(tag); + let (numbers, suffix) = match core.split_once('-') { + Some((n, s)) => (n, s.to_string()), + None => (core, String::new()), + }; + let parts: Vec = numbers + .split('.') + .map(|p| p.parse().ok()) + .collect::>()?; + (!parts.is_empty()).then_some((parts, suffix)) +} + +/// True for tags that are not meant for production (release candidates, nightlies …). +fn is_prerelease(suffix: &str) -> bool { + let s = suffix.to_ascii_lowercase(); + [ + "rc", "alpha", "beta", "dev", "nightly", "snapshot", "pre", "test", + ] + .iter() + .any(|m| s.starts_with(m) || s.contains(&format!("-{m}"))) +} + +/// Tags from the registry that are newer than `current`, oldest first. Only tags with the +/// same suffix (`-rootless`, `-debian-12` …) are considered, so the variant stays the same. +pub fn newer_tags(current: &str, available: &[String]) -> Vec { + let Some((now, suffix)) = version_parts(current) else { + return Vec::new(); + }; + let mut newer: Vec<(Vec, String)> = available + .iter() + .filter_map(|t| { + let (v, s) = version_parts(t)?; + (s == suffix + && !is_prerelease(&s) + && cmp_version(&v, &now) == std::cmp::Ordering::Greater) + .then_some((v, t.clone())) + }) + .collect(); + // equal versions (1.10 and 1.10.0) keep a stable, predictable order + newer.sort_by(|a, b| cmp_version(&a.0, &b.0).then_with(|| a.1.cmp(&b.1))); + newer.into_iter().map(|(_, t)| t).collect() +} + +/// The newest tag that is newer than `current`, if any. +pub fn newest_tag(current: &str, available: &[String]) -> Option { + newer_tags(current, available).pop() +} + +fn cmp_version(a: &[u64], b: &[u64]) -> std::cmp::Ordering { + let len = a.len().max(b.len()); + for i in 0..len { + let (x, y) = ( + a.get(i).copied().unwrap_or(0), + b.get(i).copied().unwrap_or(0), + ); + if x != y { + return x.cmp(&y); + } + } + std::cmp::Ordering::Equal +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parses_the_common_reference_shapes() { + let r = ImageRef::parse("docker.gitea.com/gitea:1.24.2-rootless").unwrap(); + assert_eq!( + (r.registry.as_str(), r.repository.as_str(), r.tag.as_str()), + ("docker.gitea.com", "gitea", "1.24.2-rootless") + ); + let r = ImageRef::parse("docker.io/bitnami/postgresql:17.5.0").unwrap(); + assert_eq!( + (r.registry.as_str(), r.repository.as_str()), + ("docker.io", "bitnami/postgresql") + ); + let r = ImageRef::parse("coredns/coredns:1.10.1").unwrap(); + assert_eq!( + (r.registry.as_str(), r.repository.as_str()), + ("docker.io", "coredns/coredns") + ); + let r = ImageRef::parse("nginx").unwrap(); + assert_eq!( + (r.registry.as_str(), r.repository.as_str(), r.tag.as_str()), + ("docker.io", "library/nginx", "latest") + ); + let r = ImageRef::parse("registry.k8s.io/ingress-nginx/controller:v1.11.5").unwrap(); + assert_eq!( + (r.registry.as_str(), r.repository.as_str(), r.tag.as_str()), + ("registry.k8s.io", "ingress-nginx/controller", "v1.11.5") + ); + let r = ImageRef::parse("localhost:32000/app:1").unwrap(); + assert_eq!(r.registry, "localhost:32000"); + assert_eq!(ImageRef::parse("").map(|r| r.repository), None); + } + + #[test] + fn rebuilds_a_reference_with_another_tag() { + assert_eq!( + ImageRef::parse("docker.gitea.com/gitea:1.24.2") + .unwrap() + .with_tag("1.25.0"), + "docker.gitea.com/gitea:1.25.0" + ); + assert_eq!( + ImageRef::parse("coredns/coredns:1.10.1") + .unwrap() + .with_tag("1.11.0"), + "coredns/coredns:1.11.0" + ); + assert_eq!( + ImageRef::parse("nginx:1.0").unwrap().with_tag("1.1"), + "nginx:1.1" + ); + } + + #[test] + fn suggests_only_newer_tags_of_the_same_variant() { + let tags: Vec = [ + "1.23.8", + "1.24.2", + "1.24.3", + "1.25.0", + "1.25.1", + "1.25.1-rootless", + "1.24.2-rootless", + "1.26.0-rc1", + "latest", + "dev", + ] + .iter() + .map(|s| s.to_string()) + .collect(); + assert_eq!( + newer_tags("1.24.2", &tags), + vec!["1.24.3", "1.25.0", "1.25.1"] + ); + assert_eq!(newest_tag("1.24.2", &tags).as_deref(), Some("1.25.1")); + // the variant is kept + assert_eq!( + newer_tags("1.24.2-rootless", &tags), + vec!["1.25.1-rootless"] + ); + // release candidates and non-version tags are ignored + assert!(!newer_tags("1.25.1", &tags).iter().any(|t| t.contains("rc"))); + assert_eq!(newest_tag("1.25.1", &tags), None, "already the newest"); + assert_eq!(newest_tag("latest", &tags), None, "no version to compare"); + } + + #[test] + fn compares_versions_by_number_not_by_text() { + let tags: Vec = ["1.9.0", "1.10.0", "1.10", "2.0.0"] + .iter() + .map(|s| s.to_string()) + .collect(); + assert_eq!(newest_tag("1.9.0", &tags).as_deref(), Some("2.0.0")); + assert_eq!(newer_tags("1.9.0", &tags), vec!["1.10", "1.10.0", "2.0.0"]); + assert_eq!( + newest_tag("v0.6.3", &["v0.7.0".to_string(), "v0.6.4".to_string()]).as_deref(), + Some("v0.7.0") + ); + } +} diff --git a/backend/crates/domain/src/jobs.rs b/backend/crates/domain/src/jobs.rs index d83eace..cee411f 100644 --- a/backend/crates/domain/src/jobs.rs +++ b/backend/crates/domain/src/jobs.rs @@ -8,6 +8,7 @@ pub enum JobKind { PackageRefresh, PackageUpgrade, VulnerabilityScan, + ImageUpdateCheck, Backup, } @@ -24,6 +25,7 @@ impl JobKind { JobKind::PackageRefresh => "package_refresh", JobKind::PackageUpgrade => "package_upgrade", JobKind::VulnerabilityScan => "vulnerability_scan", + JobKind::ImageUpdateCheck => "image_update_check", JobKind::Backup => "backup", } } @@ -37,6 +39,7 @@ impl JobKind { match self { JobKind::PackageRefresh => Some("0 0 * * * *"), JobKind::VulnerabilityScan => Some("0 0 3 * * *"), + JobKind::ImageUpdateCheck => Some("0 30 4 * * *"), _ => None, } } diff --git a/backend/crates/domain/src/lib.rs b/backend/crates/domain/src/lib.rs index 6d1b4f4..e3675a9 100644 --- a/backend/crates/domain/src/lib.rs +++ b/backend/crates/domain/src/lib.rs @@ -5,6 +5,7 @@ pub mod backup; pub mod cluster; pub mod error; pub mod host; +pub mod image; pub mod jobs; pub mod ports; pub mod settings; diff --git a/backend/crates/domain/src/ports.rs b/backend/crates/domain/src/ports.rs index 99525ad..3748892 100644 --- a/backend/crates/domain/src/ports.rs +++ b/backend/crates/domain/src/ports.rs @@ -6,6 +6,7 @@ use crate::auth::{AccessClaims, AuthEvent, RefreshToken}; use crate::backup::{BackupRecord, BackupSource, BackupStrategy, BackupTarget, RemoteFile}; use crate::cluster::{ClusterOverview, WorkloadRef}; use crate::host::{Inventory, OsInfo, Package}; +use crate::image::ImageUpdate; use crate::jobs::{JobKind, JobRun, JobStatus}; use crate::settings::SmtpSettings; use crate::user::{User, UserUpdate}; @@ -221,3 +222,16 @@ pub trait BackupCollector: Send + Sync { pub trait FileEncryptor: Send + Sync { async fn encrypt(&self, input: &Path, passphrase: &str) -> Result; } + +/// Reads the available tags of a container repository. +#[async_trait] +pub trait ImageRegistry: Send + Sync { + async fn tags(&self, image: &str) -> Result, DomainError>; +} + +#[async_trait] +pub trait ImageUpdateRepository: Send + Sync { + async fn upsert(&self, update: &ImageUpdate) -> Result<(), DomainError>; + async fn get(&self, image: &str) -> Result, DomainError>; + async fn list(&self) -> Result, DomainError>; +} diff --git a/backend/crates/infrastructure/migrations/0007_image_updates.sql b/backend/crates/infrastructure/migrations/0007_image_updates.sql new file mode 100644 index 0000000..ff22b32 --- /dev/null +++ b/backend/crates/infrastructure/migrations/0007_image_updates.sql @@ -0,0 +1,8 @@ +CREATE TABLE image_updates ( + image TEXT PRIMARY KEY, + candidate TEXT NOT NULL, + checked_at TEXT NOT NULL, + fixed TEXT NOT NULL, + fixed_counts TEXT NOT NULL, + candidate_total INTEGER NOT NULL +); diff --git a/backend/crates/infrastructure/src/lib.rs b/backend/crates/infrastructure/src/lib.rs index 8004ee5..f655fef 100644 --- a/backend/crates/infrastructure/src/lib.rs +++ b/backend/crates/infrastructure/src/lib.rs @@ -6,6 +6,7 @@ pub mod host; pub mod k8s; pub mod mail; pub mod password; +pub mod registry; pub mod sqlite; pub mod token; pub mod trivy; @@ -22,10 +23,11 @@ pub use host::{ pub use k8s::{FakeClusterGateway, KubeGateway}; pub use mail::LettreMailer; pub use password::Argon2Hasher; +pub use registry::{CurlImageRegistry, FakeImageRegistry}; pub use sqlite::{SqliteAuditLog, SqliteJobRuns, SqliteRefreshTokens, SqliteSettings, SqliteUsers}; pub use sqlite::{ SqliteBackupRecords, SqliteBackupStrategies, SqliteBackupTargets, SqliteFindings, - SqliteInventory, + SqliteImageUpdates, SqliteInventory, }; pub use token::JwtIssuer; pub use trivy::{FakeScanner, TrivyScanner}; diff --git a/backend/crates/infrastructure/src/registry.rs b/backend/crates/infrastructure/src/registry.rs new file mode 100644 index 0000000..ccbc049 --- /dev/null +++ b/backend/crates/infrastructure/src/registry.rs @@ -0,0 +1,300 @@ +//! Reads the tag list of a container repository from a registry (Docker Registry v2 API). +//! Uses curl through the `CommandRunner` so it works with the host's proxy settings and +//! needs no extra HTTP stack. +use std::sync::Arc; + +use async_trait::async_trait; +use domain::image::ImageRef; +use domain::ports::ImageRegistry; +use domain::DomainError; + +use crate::host::CommandRunner; + +pub struct CurlImageRegistry { + runner: Arc, +} + +impl CurlImageRegistry { + pub fn new(runner: Arc) -> Self { + Self { runner } + } +} + +/// Docker Hub is addressed under a different host than its image references use. +pub fn registry_host(registry: &str) -> &str { + match registry { + "docker.io" | "index.docker.io" => "registry-1.docker.io", + other => other, + } +} + +/// Bearer challenge of a registry: `Bearer realm="…",service="…",scope="…"`. +pub fn parse_auth_challenge(headers: &str) -> Option { + let line = headers + .lines() + .find(|l| l.to_ascii_lowercase().starts_with("www-authenticate:"))?; + let value = line.split_once(':')?.1.trim().strip_prefix("Bearer ")?; + let field = |name: &str| { + value + .split(',') + .find_map(|p| p.trim().strip_prefix(&format!("{name}="))) + .map(|v| v.trim_matches('"').to_string()) + }; + let realm = field("realm")?; + let mut url = format!("{realm}?"); + if let Some(service) = field("service") { + url.push_str(&format!("service={service}&")); + } + if let Some(scope) = field("scope") { + url.push_str(&format!("scope={scope}")); + } + Some(url.trim_end_matches(['&', '?']).to_string()) +} + +/// `{"tags": ["1.0", "1.1"]}`; a missing or null list means no tags. +pub fn parse_tags(body: &str) -> Result, DomainError> { + #[derive(serde::Deserialize)] + struct Response { + #[serde(default)] + tags: Option>, + } + let parsed: Response = serde_json::from_str(body) + .map_err(|e| DomainError::Unavailable(format!("registry response: {e}")))?; + Ok(parsed.tags.unwrap_or_default()) +} + +/// `Link: ; rel="next"` → the next path. +pub fn parse_next_link(headers: &str) -> Option { + let line = headers + .lines() + .find(|l| l.to_ascii_lowercase().starts_with("link:"))?; + let value = line.split_once(':')?.1; + let start = value.find('<')? + 1; + let end = value[start..].find('>')? + start; + value[start..end].to_string().into() +} + +fn split_response(out: &str) -> (&str, &str) { + // curl -i prints headers, a blank line, then the body; a redirect can add more blocks + match out.rfind("\r\n\r\n") { + Some(i) => (&out[..i], &out[i + 4..]), + None => match out.rfind("\n\n") { + Some(i) => (&out[..i], &out[i + 2..]), + None => (out, ""), + }, + } +} + +fn status_of(headers: &str) -> u16 { + headers + .lines() + .rfind(|l| l.starts_with("HTTP/")) + .and_then(|l| l.split_whitespace().nth(1)) + .and_then(|c| c.parse().ok()) + .unwrap_or(0) +} + +impl CurlImageRegistry { + async fn get( + &self, + url: &str, + token: Option<&str>, + ) -> Result<(u16, String, String), DomainError> { + let auth = token.map(|t| format!("Authorization: Bearer {t}")); + let mut args = vec!["-sS", "-i", "-L", "--max-time", "30"]; + if let Some(a) = &auth { + args.extend(["-H", a.as_str()]); + } + args.push(url); + let out = self.runner.run("curl", &args).await?; + if !out.success { + return Err(DomainError::Unavailable(format!( + "registry request failed: {}", + out.stderr.trim() + ))); + } + let (headers, body) = split_response(&out.stdout); + Ok((status_of(headers), headers.to_string(), body.to_string())) + } +} + +#[async_trait] +impl ImageRegistry for CurlImageRegistry { + async fn tags(&self, image: &str) -> Result, DomainError> { + let parsed = ImageRef::parse(image) + .ok_or_else(|| DomainError::Validation(format!("cannot parse image '{image}'")))?; + let host = registry_host(&parsed.registry).to_string(); + let mut url = format!("https://{host}/v2/{}/tags/list?n=200", parsed.repository); + let mut token: Option = None; + let mut tags = Vec::new(); + + for _ in 0..10 { + let (mut status, mut headers, mut body) = self.get(&url, token.as_deref()).await?; + if status == 401 { + let challenge = parse_auth_challenge(&headers).ok_or_else(|| { + DomainError::Unavailable(format!("{host}: authentication required")) + })?; + let (_, _, token_body) = self.get(&challenge, None).await?; + #[derive(serde::Deserialize)] + struct Token { + #[serde(alias = "access_token")] + token: Option, + } + token = serde_json::from_str::(&token_body) + .ok() + .and_then(|t| t.token); + if token.is_none() { + return Err(DomainError::Unavailable(format!( + "{host}: no token in the auth response" + ))); + } + (status, headers, body) = self.get(&url, token.as_deref()).await?; + } + if status != 200 { + return Err(DomainError::Unavailable(format!( + "{host}: tag list returned HTTP {status}" + ))); + } + tags.extend(parse_tags(&body)?); + match parse_next_link(&headers) { + Some(next) => url = format!("https://{host}{next}"), + None => break, + } + } + Ok(tags) + } +} + +/// Offers one newer tag for every image (FAKE_HOST=true). +pub struct FakeImageRegistry; + +#[async_trait] +impl ImageRegistry for FakeImageRegistry { + async fn tags(&self, image: &str) -> Result, DomainError> { + let current = ImageRef::parse(image).map(|r| r.tag).unwrap_or_default(); + Ok(vec![ + current, + "9.9.9".into(), + "10.0.0-rc1".into(), + "latest".into(), + ]) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use domain::ports::LineSink; + use std::sync::Mutex; + + #[test] + fn maps_docker_hub_to_its_api_host() { + assert_eq!(registry_host("docker.io"), "registry-1.docker.io"); + assert_eq!(registry_host("registry.k8s.io"), "registry.k8s.io"); + } + + #[test] + fn reads_the_bearer_challenge_and_the_next_page() { + let headers = "HTTP/1.1 401 Unauthorized\r\nwww-authenticate: Bearer realm=\"https://auth.docker.io/token\",service=\"registry.docker.io\",scope=\"repository:library/nginx:pull\"\r\n"; + assert_eq!( + parse_auth_challenge(headers).unwrap(), + "https://auth.docker.io/token?service=registry.docker.io&scope=repository:library/nginx:pull" + ); + assert_eq!(parse_auth_challenge("HTTP/1.1 200 OK\r\n"), None); + let paged = "HTTP/1.1 200 OK\r\nLink: ; rel=\"next\"\r\n"; + assert_eq!( + parse_next_link(paged).unwrap(), + "/v2/library/nginx/tags/list?n=200&last=1.9" + ); + assert_eq!(parse_next_link("HTTP/1.1 200 OK\r\n"), None); + } + + #[test] + fn reads_the_tag_list() { + assert_eq!( + parse_tags(r#"{"name":"x","tags":["1.0","1.1"]}"#).unwrap(), + vec!["1.0", "1.1"] + ); + assert_eq!( + parse_tags(r#"{"name":"x","tags":null}"#).unwrap(), + Vec::::new() + ); + assert!(parse_tags("").is_err()); + } + + #[derive(Default)] + struct Fake { + calls: Mutex>, + } + + #[async_trait] + impl CommandRunner for Fake { + async fn run(&self, _p: &str, args: &[&str]) -> Result { + let url = args.last().unwrap().to_string(); + let authorized = args.iter().any(|a| a.starts_with("Authorization:")); + self.calls.lock().unwrap().push(url.clone()); + let stdout = if url.contains("auth.example.com") { + "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\n\r\n{\"token\":\"abc\"}" + .to_string() + } else if !authorized { + "HTTP/1.1 401 Unauthorized\r\nwww-authenticate: Bearer realm=\"https://auth.example.com/token\",service=\"reg\",scope=\"repository:gitea:pull\"\r\n\r\n{}".to_string() + } else if url.contains("last=1.24.2") { + "HTTP/1.1 200 OK\r\n\r\n{\"tags\":[\"1.25.0\"]}".to_string() + } else { + "HTTP/1.1 200 OK\r\nLink: ; rel=\"next\"\r\n\r\n{\"tags\":[\"1.24.1\",\"1.24.2\"]}".to_string() + }; + Ok(crate::host::Output { + stdout, + success: true, + ..Default::default() + }) + } + async fn run_env( + &self, + p: &str, + a: &[&str], + _: &[(&str, &str)], + ) -> Result { + self.run(p, a).await + } + async fn run_to_file( + &self, + _: &str, + _: &[&str], + _: &std::path::Path, + ) -> Result { + unreachable!() + } + async fn read_file(&self, _: &str) -> Result, DomainError> { + Ok(None) + } + async fn run_streaming( + &self, + _: &str, + _: &[&str], + _: &dyn LineSink, + ) -> Result { + Ok(true) + } + } + + #[tokio::test] + async fn authenticates_and_follows_pagination() { + let runner = Arc::new(Fake::default()); + let tags = CurlImageRegistry::new(runner.clone()) + .tags("docker.gitea.com/gitea:1.24.2") + .await + .unwrap(); + assert_eq!(tags, vec!["1.24.1", "1.24.2", "1.25.0"]); + let calls = runner.calls.lock().unwrap(); + assert!( + calls[0].starts_with("https://docker.gitea.com/v2/gitea/tags/list"), + "{calls:?}" + ); + assert!( + calls[1].starts_with("https://auth.example.com/token?service=reg"), + "{calls:?}" + ); + assert!(calls.last().unwrap().contains("last=1.24.2"), "{calls:?}"); + } +} diff --git a/backend/crates/infrastructure/src/sqlite.rs b/backend/crates/infrastructure/src/sqlite.rs index ddfa138..65b7401 100644 --- a/backend/crates/infrastructure/src/sqlite.rs +++ b/backend/crates/infrastructure/src/sqlite.rs @@ -1278,3 +1278,115 @@ mod backup_tests { ); } } + +use domain::image::ImageUpdate; + +pub struct SqliteImageUpdates(pub DbPool); + +fn image_update_from_row(r: &SqliteRow) -> Result { + let json = |c: &str| -> Result { + serde_json::from_str(r.get::(c).as_str()) + .map_err(|e| DomainError::Storage(e.to_string())) + }; + Ok(ImageUpdate { + image: r.get("image"), + candidate: r.get("candidate"), + checked_at: parse_ts(r.get::("checked_at").as_str()), + fixed: serde_json::from_value(json("fixed")?) + .map_err(|e| DomainError::Storage(e.to_string()))?, + fixed_counts: serde_json::from_value(json("fixed_counts")?) + .map_err(|e| DomainError::Storage(e.to_string()))?, + candidate_total: r.get::("candidate_total") as usize, + }) +} + +const IMAGE_UPDATE_COLS: &str = + "image, candidate, checked_at, fixed, fixed_counts, candidate_total"; + +#[async_trait] +impl domain::ports::ImageUpdateRepository for SqliteImageUpdates { + async fn upsert(&self, u: &ImageUpdate) -> Result<(), DomainError> { + let fixed = + serde_json::to_string(&u.fixed).map_err(|e| DomainError::Storage(e.to_string()))?; + let counts = serde_json::to_string(&u.fixed_counts) + .map_err(|e| DomainError::Storage(e.to_string()))?; + sqlx::query(&format!( + "INSERT INTO image_updates ({IMAGE_UPDATE_COLS}) VALUES (?, ?, ?, ?, ?, ?) \ + ON CONFLICT(image) DO UPDATE SET candidate = excluded.candidate, checked_at = excluded.checked_at, \ + fixed = excluded.fixed, fixed_counts = excluded.fixed_counts, candidate_total = excluded.candidate_total" + )) + .bind(&u.image) + .bind(&u.candidate) + .bind(u.checked_at.to_rfc3339()) + .bind(fixed) + .bind(counts) + .bind(u.candidate_total as i64) + .execute(&self.0) + .await + .map(|_| ()) + .map_err(storage) + } + + async fn get(&self, image: &str) -> Result, DomainError> { + let row = sqlx::query(&format!( + "SELECT {IMAGE_UPDATE_COLS} FROM image_updates WHERE image = ?" + )) + .bind(image) + .fetch_optional(&self.0) + .await + .map_err(storage)?; + row.as_ref().map(image_update_from_row).transpose() + } + + async fn list(&self) -> Result, DomainError> { + let rows = sqlx::query(&format!( + "SELECT {IMAGE_UPDATE_COLS} FROM image_updates ORDER BY image" + )) + .fetch_all(&self.0) + .await + .map_err(storage)?; + rows.iter().map(image_update_from_row).collect() + } +} + +#[cfg(test)] +mod image_update_tests { + use super::*; + use domain::ports::ImageUpdateRepository; + use domain::vuln::SeverityCounts; + + #[tokio::test] + async fn upsert_get_and_list() { + let pool = crate::connect("sqlite::memory:").await.unwrap(); + let repo = SqliteImageUpdates(pool); + let u = ImageUpdate { + image: "gitea/gitea:1.22".into(), + candidate: "gitea/gitea:1.23".into(), + checked_at: Utc::now(), + fixed: vec!["CVE-1".into(), "CVE-2".into()], + fixed_counts: SeverityCounts { + critical: 1, + high: 1, + ..Default::default() + }, + candidate_total: 7, + }; + repo.upsert(&u).await.unwrap(); + let got = repo.get(&u.image).await.unwrap().unwrap(); + assert_eq!(got.fixed, u.fixed); + assert_eq!(got.fixed_counts.critical, 1); + assert_eq!(got.candidate_total, 7); + + let newer = ImageUpdate { + candidate: "gitea/gitea:1.24".into(), + fixed: vec!["CVE-3".into()], + ..u.clone() + }; + repo.upsert(&newer).await.unwrap(); + let all = repo.list().await.unwrap(); + assert_eq!(all.len(), 1, "one row per image"); + assert_eq!(all[0].candidate, "gitea/gitea:1.24"); + assert_eq!(all[0].fixed, vec!["CVE-3"]); + assert_eq!(repo.get("other").await.unwrap(), None); + } +} diff --git a/backend/crates/infrastructure/src/trivy.rs b/backend/crates/infrastructure/src/trivy.rs index 9816ec0..cad0e8c 100644 --- a/backend/crates/infrastructure/src/trivy.rs +++ b/backend/crates/infrastructure/src/trivy.rs @@ -254,6 +254,10 @@ impl VulnerabilityScanner for FakeScanner { out: &dyn LineSink, ) -> Result, DomainError> { out.line(&format!("fake: scanning image {image}")); + // a candidate tag stands for a fixed release + if image.ends_with(":9.9.9") { + return Ok(vec![]); + } Ok(if image.contains("gitea") { vec![RawFinding { cve_id: "CVE-2024-24790".into(), diff --git a/frontend/e2e/image-updates.spec.ts b/frontend/e2e/image-updates.spec.ts new file mode 100644 index 0000000..1a68056 --- /dev/null +++ b/frontend/e2e/image-updates.spec.ts @@ -0,0 +1,36 @@ +import { test, expect, type Page } from '@playwright/test' + +async function scanAndOpenContainers(page: Page) { + await page.goto('/login') + await page.getByLabel('Email').fill('admin@example.com') + await page.getByLabel('Password').fill('admin-password-123') + await page.getByRole('button', { name: 'Sign in' }).click() + await expect(page.getByRole('heading', { name: 'Dashboard' })).toBeVisible() + await page.goto('/vulnerabilities') + await page.getByRole('button', { name: 'Scan now' }).click() + await expect(page.getByTestId('scope-host-critical')).not.toHaveText('0', { timeout: 20_000 }) + await page.getByTestId('scope-container').click() +} + +test('an image finding links to the workload running it', async ({ page }) => { + await scanAndOpenContainers(page) + const row = page.getByRole('row', { name: /gitea\/gitea/ }).first() + await expect(row).toContainText('gitea/deployment/gitea') + + await row.getByRole('link', { name: 'gitea/deployment/gitea' }).click() + await expect(page).toHaveURL(/\/cluster\?image=/) + await expect(page.getByText('Showing the workloads that run')).toBeVisible() + await expect(page.getByRole('row', { name: /gitea gitea deployment/ })).toBeVisible() +}) + +test('checking an image suggests a newer tag and marks the CVEs it fixes', async ({ page }) => { + await scanAndOpenContainers(page) + const row = page.getByRole('row', { name: /gitea\/gitea/ }).first() + await row.getByRole('button', { name: 'Check for update' }).click() + await expect(row).toContainText('9.9.9', { timeout: 20_000 }) + await expect(row).toContainText('fixes 1') + + await row.click() + await expect(page.getByRole('row', { name: /CVE-2024-24790/ })).toContainText('fixed in 9.9.9') + await expect(row.getByRole('button', { name: /Update to 9.9.9/ })).toBeVisible() +}) diff --git a/frontend/src/api/types.ts b/frontend/src/api/types.ts index 09487ff..a8fbe8f 100644 --- a/frontend/src/api/types.ts +++ b/frontend/src/api/types.ts @@ -162,6 +162,27 @@ export interface FindingGroup { packages: number } +export interface ImageUpdate { + image: string + candidate: string + checked_at: string + fixed: string[] + fixed_counts: SeverityCounts + candidate_total: number +} + +export interface ImageGroup extends FindingGroup { + workloads: WorkloadRef[] + running: boolean + update: ImageUpdate | null +} + +export interface WorkloadRef { + namespace: string + kind: WorkloadKind + name: string +} + export interface TargetSummary { target: string kind: 'os' | 'image' diff --git a/frontend/src/components/FindingGroups.test.ts b/frontend/src/components/FindingGroups.test.ts index 303e0db..acd4867 100644 --- a/frontend/src/components/FindingGroups.test.ts +++ b/frontend/src/components/FindingGroups.test.ts @@ -1,6 +1,6 @@ import { mount, flushPromises } from '@vue/test-utils' import FindingGroups from './FindingGroups.vue' -import type { Finding, FindingGroup } from '../api/types' +import type { Finding, FindingGroup, ImageGroup } from '../api/types' const group = (over: Partial): FindingGroup => ({ key: 'openssl', @@ -103,3 +103,91 @@ describe('FindingGroups', () => { expect(w.text()).toContain('No findings') }) }) + +describe('FindingGroups for images', () => { + const image = (over: Partial = {}): ImageGroup => ({ + ...group({ + key: 'gitea/gitea:1.22.3', + kind: 'image', + source: '', + installed: '', + packages: 12, + total: 40, + }), + workloads: [{ namespace: 'gitea', kind: 'deployment', name: 'gitea' }], + running: true, + update: null, + ...over, + }) + + it('links each image to the workloads that run it', () => { + const w = mount(FindingGroups, { + props: { groups: [image()], scope: 'container', canAct: true, load: vi.fn() }, + global: { stubs: { RouterLink: { template: '', props: ['to'] } } }, + }) + const link = w.get('a') + expect(link.text()).toBe('gitea/deployment/gitea') + expect(link.attributes('href')).toBe('/cluster?image=gitea%2Fgitea%3A1.22.3') + }) + + it('suggests the checked update and how much it fixes', () => { + const withUpdate = image({ + update: { + image: 'gitea/gitea:1.22.3', + candidate: 'gitea/gitea:1.25.1', + checked_at: '2026-09-03T08:00:00Z', + fixed: ['CVE-1', 'CVE-2'], + fixed_counts: { critical: 1, high: 1, medium: 0, low: 0, unknown: 0 }, + candidate_total: 5, + }, + }) + const w = mount(FindingGroups, { + props: { groups: [withUpdate], scope: 'container', canAct: true, load: vi.fn() }, + global: { stubs: { RouterLink: { template: '' } } }, + }) + const text = w.text() + expect(text).toContain('1.25.1') + expect(text).toContain('fixes 2') + expect(w.find('button[name=update-image]').exists()).toBe(true) + }) + + it('offers a check when nothing is known yet and hides the actions from viewers', () => { + const w = mount(FindingGroups, { + props: { groups: [image()], scope: 'container', canAct: true, load: vi.fn() }, + global: { stubs: { RouterLink: { template: '' } } }, + }) + expect(w.find('button[name=check-image]').exists()).toBe(true) + expect(w.find('button[name=update-image]').exists()).toBe(false) + + const viewer = mount(FindingGroups, { + props: { groups: [image()], scope: 'container', canAct: false, load: vi.fn() }, + global: { stubs: { RouterLink: { template: '' } } }, + }) + expect(viewer.find('button[name=check-image]').exists()).toBe(false) + }) + + it('marks the CVEs that the checked update fixes', async () => { + const withUpdate = image({ + update: { + image: 'gitea/gitea:1.22.3', + candidate: 'gitea/gitea:1.25.1', + checked_at: '2026-09-03T08:00:00Z', + fixed: ['CVE-2024-1'], + fixed_counts: { critical: 1, high: 0, medium: 0, low: 0, unknown: 0 }, + candidate_total: 0, + }, + }) + const load = vi.fn().mockResolvedValue([finding('CVE-2024-1'), finding('CVE-2024-9')]) + const w = mount(FindingGroups, { + props: { groups: [withUpdate], scope: 'container', canAct: true, load }, + global: { stubs: { RouterLink: { template: '' } } }, + }) + await w.findAll('tbody tr')[0].trigger('click') + await flushPromises() + const cveRows = w.findAll('table table tbody tr') + const fixedRow = cveRows.find((r) => r.text().includes('CVE-2024-1'))! + expect(fixedRow.text()).toContain('fixed in 1.25.1') + const otherRow = cveRows.find((r) => r.text().includes('CVE-2024-9'))! + expect(otherRow.text()).not.toContain('fixed in') + }) +}) diff --git a/frontend/src/components/FindingGroups.vue b/frontend/src/components/FindingGroups.vue index 9705a1d..1997f30 100644 --- a/frontend/src/components/FindingGroups.vue +++ b/frontend/src/components/FindingGroups.vue @@ -1,6 +1,6 @@ @@ -23,7 +26,10 @@ defineEmits<{ restart: [w: Workload]; scale: [w: Workload]; image: [w: Workload] v-for="w in workloads" :key="`${w.namespace}/${w.kind}/${w.name}`" class="border-b border-gray-100" - :class="w.ready < w.desired ? 'bg-red-50' : ''" + :class="[ + w.ready < w.desired ? 'bg-red-50' : '', + runsHighlighted(w) ? 'bg-blue-50 ring-1 ring-blue-200' : '', + ]" > {{ w.namespace }} {{ w.name }} diff --git a/frontend/src/pages/ClusterPage.vue b/frontend/src/pages/ClusterPage.vue index fd49287..83ad10e 100644 --- a/frontend/src/pages/ClusterPage.vue +++ b/frontend/src/pages/ClusterPage.vue @@ -1,5 +1,6 @@