Link image findings to their workloads and suggest a fixing image update

Container findings now name the workloads that run the image and link
into the Kubernetes view, which highlights them. An update check asks the
registry for newer tags of the same variant, scans the newest one and
records which of the open findings are gone in it. The image row then
shows the candidate tag, how many findings it fixes and how many remain,
marks those CVEs in the expanded list, and offers to roll every workload
over to it. The check runs as a job, nightly for all running images or on
demand for one.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Dennis Nemec
2026-09-03 20:19:12 +02:00
parent 278b5e47a3
commit b984cc8c9f
31 changed files with 1674 additions and 29 deletions

View File

@ -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
);

View File

@ -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};

View File

@ -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<dyn CommandRunner>,
}
impl CurlImageRegistry {
pub fn new(runner: Arc<dyn CommandRunner>) -> 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<String> {
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<Vec<String>, DomainError> {
#[derive(serde::Deserialize)]
struct Response {
#[serde(default)]
tags: Option<Vec<String>>,
}
let parsed: Response = serde_json::from_str(body)
.map_err(|e| DomainError::Unavailable(format!("registry response: {e}")))?;
Ok(parsed.tags.unwrap_or_default())
}
/// `Link: </v2/x/tags/list?n=100&last=y>; rel="next"` → the next path.
pub fn parse_next_link(headers: &str) -> Option<String> {
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<Vec<String>, 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<String> = 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<String>,
}
token = serde_json::from_str::<Token>(&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<Vec<String>, 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: </v2/library/nginx/tags/list?n=200&last=1.9>; 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::<String>::new()
);
assert!(parse_tags("<html>").is_err());
}
#[derive(Default)]
struct Fake {
calls: Mutex<Vec<String>>,
}
#[async_trait]
impl CommandRunner for Fake {
async fn run(&self, _p: &str, args: &[&str]) -> Result<crate::host::Output, DomainError> {
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: </v2/gitea/tags/list?n=200&last=1.24.2>; 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<crate::host::Output, DomainError> {
self.run(p, a).await
}
async fn run_to_file(
&self,
_: &str,
_: &[&str],
_: &std::path::Path,
) -> Result<crate::host::Output, DomainError> {
unreachable!()
}
async fn read_file(&self, _: &str) -> Result<Option<String>, DomainError> {
Ok(None)
}
async fn run_streaming(
&self,
_: &str,
_: &[&str],
_: &dyn LineSink,
) -> Result<bool, DomainError> {
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:?}");
}
}

View File

@ -1278,3 +1278,115 @@ mod backup_tests {
);
}
}
use domain::image::ImageUpdate;
pub struct SqliteImageUpdates(pub DbPool);
fn image_update_from_row(r: &SqliteRow) -> Result<ImageUpdate, DomainError> {
let json = |c: &str| -> Result<serde_json::Value, DomainError> {
serde_json::from_str(r.get::<String, _>(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::<String, _>("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::<i64, _>("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<Option<ImageUpdate>, 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<Vec<ImageUpdate>, 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);
}
}

View File

@ -254,6 +254,10 @@ impl VulnerabilityScanner for FakeScanner {
out: &dyn LineSink,
) -> Result<Vec<RawFinding>, 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(),