WP-11: contract and failing tests for package upgrades
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled
CI / ui (push) Has been cancelled

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Dennis Nemec
2026-09-02 22:24:40 +02:00
parent b6ddb8889d
commit 1f56f015a2
18 changed files with 524 additions and 2 deletions

1
backend/Cargo.lock generated
View File

@ -102,6 +102,7 @@ dependencies = [
"cron",
"domain",
"rand",
"serde",
"serde_json",
"sha2",
"tokio",

View File

@ -0,0 +1,83 @@
//! WP-11: POST /api/system/upgrade
mod common;
use axum::http::StatusCode;
use common::{get, post, test_app_with_admin};
use serde_json::json;
const ADMIN: &str = "admin@example.com";
const PW: &str = "admin-password-123";
async fn wait_done(app: &axum::Router, token: &str, id: &str) -> serde_json::Value {
for _ in 0..100 {
let r = get(app, &format!("/api/jobs/{id}"), Some(token)).await;
if r.json["status"] != "running" {
return r.json;
}
tokio::time::sleep(std::time::Duration::from_millis(30)).await;
}
panic!("job did not finish");
}
#[tokio::test]
async fn admin_upgrades_selected_packages_and_log_shows_output() {
let app = test_app_with_admin().await;
let token = common::login(&app, ADMIN, PW).await.access;
let res = post(
&app,
"/api/system/upgrade",
json!({"packages": ["openssl", "curl"]}),
Some(&token),
)
.await;
assert_eq!(res.status, StatusCode::ACCEPTED, "{}", res.json);
assert_eq!(res.json["kind"], "package_upgrade");
let run = wait_done(&app, &token, res.json["id"].as_str().unwrap()).await;
assert_eq!(run["status"], "success", "{}", run["log"]);
let log = run["log"].as_str().unwrap();
assert!(log.contains("openssl curl"), "{log}");
assert!(log.contains("Setting up"), "{log}");
// inventory was refreshed afterwards
let inv = get(&app, "/api/system/inventory", Some(&token)).await;
assert!(inv.json["refreshed_at"].is_string());
}
#[tokio::test]
async fn upgrade_all_validation_and_permissions() {
let app = test_app_with_admin().await;
let admin = common::login(&app, ADMIN, PW).await.access;
let res = post(
&app,
"/api/system/upgrade",
json!({"packages": []}),
Some(&admin),
)
.await;
assert_eq!(res.status, StatusCode::ACCEPTED);
wait_done(&app, &admin, res.json["id"].as_str().unwrap()).await;
let res = post(
&app,
"/api/system/upgrade",
json!({"packages": ["bad name"]}),
Some(&admin),
)
.await;
assert_eq!(res.status, StatusCode::UNPROCESSABLE_ENTITY);
post(&app, "/api/users", json!({"email": "u@x.de", "display_name": "U", "password": "user-password-123", "role": "user"}), Some(&admin)).await;
let user = common::login(&app, "u@x.de", "user-password-123")
.await
.access;
assert_eq!(
post(
&app,
"/api/system/upgrade",
json!({"packages": []}),
Some(&user)
)
.await
.status,
StatusCode::FORBIDDEN
);
}

View File

@ -13,6 +13,7 @@ cron = "0.15"
serde_json.workspace = true
tokio.workspace = true
rand.workspace = true
serde.workspace = true
sha2.workspace = true
uuid.workspace = true

View File

@ -4,12 +4,14 @@ pub mod inventory_service;
pub mod jobs;
pub mod scheduler;
pub mod settings_service;
pub mod upgrade_service;
pub mod user_service;
pub use auth_service::AuthService;
pub use inventory_service::{InventoryService, PackageRefreshJob};
pub use jobs::{JobHandler, JobLog, JobRunner};
pub use settings_service::SettingsService;
pub use upgrade_service::{PackageUpgradeJob, UpgradeParams};
pub use user_service::UserService;
#[cfg(test)]

View File

@ -368,3 +368,34 @@ impl InventoryRepository for MemInventory {
Ok(self.0.lock().unwrap().clone())
}
}
use domain::ports::{HostUpdater, LineSink};
/// Records the requested packages and emits a few lines.
#[derive(Default)]
pub struct FakeUpdater {
pub calls: Mutex<Vec<Vec<String>>>,
pub fail: bool,
}
#[async_trait]
impl HostUpdater for FakeUpdater {
async fn upgrade(&self, packages: &[String], out: &dyn LineSink) -> Result<(), DomainError> {
self.calls.lock().unwrap().push(packages.to_vec());
out.line("Reading package lists...");
out.line(&format!(
"Upgrading {} package(s)",
if packages.is_empty() {
"all".to_string()
} else {
packages.len().to_string()
}
));
if self.fail {
return Err(DomainError::Unavailable(
"apt-get exited with status 100".into(),
));
}
Ok(())
}
}

View File

@ -3,4 +3,5 @@ mod inventory_tests;
mod jobs_tests;
mod scheduler_tests;
mod settings_tests;
mod upgrade_tests;
mod user_service_tests;

View File

@ -0,0 +1,105 @@
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use domain::DomainError;
use crate::jobs::{JobHandler, JobLog};
use crate::test_fakes::{FakeInspector, FakeUpdater, MemInventory};
use crate::{InventoryService, PackageUpgradeJob, UpgradeParams};
#[derive(Default)]
struct VecLog(Mutex<Vec<String>>);
#[async_trait]
impl JobLog for VecLog {
async fn line(&self, text: &str) {
self.0.lock().unwrap().push(text.into());
}
}
#[test]
fn params_parse_and_validate() {
assert_eq!(
UpgradeParams::from_json(None).unwrap(),
UpgradeParams::default()
);
assert_eq!(
UpgradeParams::from_json(Some("")).unwrap(),
UpgradeParams::default()
);
let p = UpgradeParams::from_json(Some(r#"{"packages":["openssl","libssl3"]}"#)).unwrap();
assert_eq!(p.packages, vec!["openssl", "libssl3"]);
assert!(matches!(
UpgradeParams::from_json(Some("{not json")).unwrap_err(),
DomainError::Validation(_)
));
let bad = UpgradeParams {
packages: vec!["rm -rf /".into()],
};
assert!(matches!(
bad.validate().unwrap_err(),
DomainError::Validation(_)
));
assert!(matches!(
UpgradeParams::from_json(Some(r#"{"packages":["../x"]}"#)).unwrap_err(),
DomainError::Validation(_)
));
}
fn job(fail: bool) -> (Arc<FakeUpdater>, Arc<MemInventory>, PackageUpgradeJob) {
let updater = Arc::new(FakeUpdater {
fail,
..Default::default()
});
let repo = Arc::new(MemInventory::default());
let inventory = Arc::new(InventoryService::new(
Arc::new(FakeInspector { fail: false }),
repo.clone(),
));
(
updater.clone(),
repo,
PackageUpgradeJob { updater, inventory },
)
}
#[tokio::test]
async fn upgrades_selected_packages_streams_output_and_refreshes_inventory() {
let (updater, repo, job) = job(false);
let log = VecLog::default();
job.run(Some(r#"{"packages":["openssl"]}"#.into()), &log)
.await
.unwrap();
assert_eq!(
updater.calls.lock().unwrap()[0],
vec!["openssl".to_string()]
);
let lines = log.0.lock().unwrap().join("\n");
assert!(lines.contains("Reading package lists"), "{lines}");
assert!(lines.contains("Upgrading 1 package"), "{lines}");
assert!(lines.contains("reboot required"), "{lines}");
assert!(
repo.0.lock().unwrap().is_some(),
"inventory refreshed after upgrade"
);
}
#[tokio::test]
async fn empty_params_upgrade_everything() {
let (updater, _, job) = job(false);
job.run(None, &VecLog::default()).await.unwrap();
assert!(updater.calls.lock().unwrap()[0].is_empty());
}
#[tokio::test]
async fn updater_failure_fails_the_job_and_keeps_output() {
let (_, _, job) = job(true);
let log = VecLog::default();
let err = job.run(None, &log).await.unwrap_err();
assert!(err.contains("status 100"));
assert!(log
.0
.lock()
.unwrap()
.iter()
.any(|l| l.contains("Reading package lists")));
}

View File

@ -0,0 +1,48 @@
//! Package upgrade job: runs the host updater, streams its output into the job log,
//! then refreshes the inventory.
use std::sync::Arc;
use async_trait::async_trait;
use domain::host::validate_package_name;
use domain::ports::HostUpdater;
use domain::DomainError;
use serde::{Deserialize, Serialize};
use crate::jobs::{JobHandler, JobLog};
use crate::InventoryService;
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct UpgradeParams {
/// Empty = upgrade all packages.
#[serde(default)]
pub packages: Vec<String>,
}
impl UpgradeParams {
pub fn validate(&self) -> Result<(), DomainError> {
self.packages
.iter()
.try_for_each(|p| validate_package_name(p))
}
pub fn to_json(&self) -> String {
serde_json::to_string(self).unwrap_or_default()
}
pub fn from_json(s: Option<&str>) -> Result<Self, DomainError> {
let _ = s;
todo!()
}
}
pub struct PackageUpgradeJob {
pub updater: Arc<dyn HostUpdater>,
pub inventory: Arc<InventoryService>,
}
#[async_trait]
impl JobHandler for PackageUpgradeJob {
async fn run(&self, _params: Option<String>, _log: &dyn JobLog) -> Result<(), String> {
todo!()
}
}

View File

@ -53,3 +53,17 @@ impl Inventory {
.count()
}
}
/// Debian package names: lowercase letters, digits, `+`, `-`, `.`; at least two characters.
pub fn validate_package_name(name: &str) -> Result<(), crate::DomainError> {
let ok = name.len() >= 2
&& name
.chars()
.next()
.is_some_and(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
&& name
.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || matches!(c, '+' | '-' | '.'));
ok.then_some(())
.ok_or_else(|| crate::DomainError::Validation(format!("invalid package name: {name}")))
}

View File

@ -90,3 +90,14 @@ pub trait InventoryRepository: Send + Sync {
async fn save(&self, inventory: &Inventory) -> Result<(), DomainError>;
async fn load(&self) -> Result<Option<Inventory>, DomainError>;
}
/// Synchronous sink for streamed command output.
pub trait LineSink: Send + Sync {
fn line(&self, text: &str);
}
/// Applies package upgrades on the host. An empty package list upgrades everything.
#[async_trait]
pub trait HostUpdater: Send + Sync {
async fn upgrade(&self, packages: &[String], out: &dyn LineSink) -> Result<(), DomainError>;
}

View File

@ -1,4 +1,5 @@
use async_trait::async_trait;
use domain::ports::LineSink;
use domain::DomainError;
#[derive(Clone, Debug, Default)]
@ -12,6 +13,13 @@ pub struct Output {
pub trait CommandRunner: Send + Sync {
async fn run(&self, program: &str, args: &[&str]) -> Result<Output, DomainError>;
async fn read_file(&self, path: &str) -> Result<Option<String>, DomainError>;
/// Run a command and forward each output line (stdout and stderr) to `out`.
async fn run_streaming(
&self,
program: &str,
args: &[&str],
out: &dyn LineSink,
) -> Result<bool, DomainError>;
}
pub struct SystemCommandRunner;
@ -33,6 +41,15 @@ impl CommandRunner for SystemCommandRunner {
})
}
async fn run_streaming(
&self,
_program: &str,
_args: &[&str],
_out: &dyn LineSink,
) -> Result<bool, DomainError> {
todo!()
}
async fn read_file(&self, path: &str) -> Result<Option<String>, DomainError> {
match tokio::fs::read_to_string(path).await {
Ok(s) => Ok(Some(s)),

View File

@ -229,6 +229,14 @@ Conf openssl (3.0.16-1~deb12u1 Debian-Security:12/stable-security [amd64])\n";
success: true,
})
}
async fn run_streaming(
&self,
_p: &str,
_a: &[&str],
_o: &dyn domain::ports::LineSink,
) -> Result<bool, DomainError> {
unreachable!()
}
async fn read_file(&self, path: &str) -> Result<Option<String>, DomainError> {
Ok(match path {
"/etc/os-release" => Some(

View File

@ -66,3 +66,27 @@ impl HostInspector for FakeHostInspector {
])
}
}
/// Pretends to upgrade packages (FAKE_HOST=true).
pub struct FakeHostUpdater;
#[async_trait]
impl domain::ports::HostUpdater for FakeHostUpdater {
async fn upgrade(
&self,
packages: &[String],
out: &dyn domain::ports::LineSink,
) -> Result<(), DomainError> {
out.line("Reading package lists... Done");
out.line("Building dependency tree... Done");
let list = if packages.is_empty() {
"all upgradable packages".to_string()
} else {
packages.join(" ")
};
out.line(&format!("The following packages will be upgraded: {list}"));
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
out.line("Setting up packages ... Done");
Ok(())
}
}

View File

@ -2,7 +2,9 @@
pub mod command;
pub mod debian;
pub mod fake;
pub mod updater;
pub use command::{CommandRunner, Output, SystemCommandRunner};
pub use debian::DebianInspector;
pub use fake::FakeHostInspector;
pub use fake::{FakeHostInspector, FakeHostUpdater};
pub use updater::DebianUpdater;

View File

@ -0,0 +1,141 @@
//! Applies apt upgrades on a Debian host.
use std::sync::Arc;
use async_trait::async_trait;
use domain::ports::{HostUpdater, LineSink};
use domain::DomainError;
use super::command::CommandRunner;
pub struct DebianUpdater {
runner: Arc<dyn CommandRunner>,
}
impl DebianUpdater {
pub fn new(runner: Arc<dyn CommandRunner>) -> Self {
Self { runner }
}
}
#[async_trait]
impl HostUpdater for DebianUpdater {
async fn upgrade(&self, _packages: &[String], _out: &dyn LineSink) -> Result<(), DomainError> {
let _ = &self.runner;
todo!()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
#[derive(Default)]
struct Recording {
calls: Mutex<Vec<(String, Vec<String>)>>,
fail_install: bool,
}
#[async_trait]
impl CommandRunner for Recording {
async fn run(
&self,
program: &str,
args: &[&str],
) -> Result<super::super::Output, DomainError> {
self.calls
.lock()
.unwrap()
.push((program.into(), args.iter().map(|s| s.to_string()).collect()));
Ok(super::super::Output {
success: true,
..Default::default()
})
}
async fn read_file(&self, _: &str) -> Result<Option<String>, DomainError> {
Ok(None)
}
async fn run_streaming(
&self,
program: &str,
args: &[&str],
out: &dyn LineSink,
) -> Result<bool, DomainError> {
self.calls
.lock()
.unwrap()
.push((program.into(), args.iter().map(|s| s.to_string()).collect()));
out.line("Reading package lists...");
Ok(!(self.fail_install && args.contains(&"install")))
}
}
#[derive(Default)]
struct Lines(Mutex<Vec<String>>);
impl LineSink for Lines {
fn line(&self, t: &str) {
self.0.lock().unwrap().push(t.into());
}
}
#[tokio::test]
async fn selected_packages_use_only_upgrade_after_apt_update() {
let r = Arc::new(Recording::default());
let out = Lines::default();
DebianUpdater::new(r.clone())
.upgrade(&["openssl".into(), "curl".into()], &out)
.await
.unwrap();
let calls = r.calls.lock().unwrap();
assert_eq!(calls[0].0, "apt-get");
assert_eq!(calls[0].1[0], "update");
assert_eq!(calls[1].0, "apt-get");
let install = calls[1].1.join(" ");
assert!(install.contains("install --only-upgrade"), "{install}");
assert!(install.contains("-y"), "{install}");
assert!(install.ends_with("openssl curl"), "{install}");
assert!(out
.0
.lock()
.unwrap()
.iter()
.any(|l| l.contains("Reading package lists")));
}
#[tokio::test]
async fn empty_selection_runs_dist_upgrade() {
let r = Arc::new(Recording::default());
DebianUpdater::new(r.clone())
.upgrade(&[], &Lines::default())
.await
.unwrap();
let calls = r.calls.lock().unwrap();
assert!(calls[1].1.contains(&"dist-upgrade".to_string()));
}
#[tokio::test]
async fn non_zero_exit_is_an_error() {
let r = Arc::new(Recording {
fail_install: true,
..Default::default()
});
let err = DebianUpdater::new(r)
.upgrade(&["x1".into()], &Lines::default())
.await
.unwrap_err();
assert!(matches!(err, DomainError::Unavailable(_)));
}
#[tokio::test]
async fn system_runner_streams_real_output() {
let out = Lines::default();
let ok = super::super::SystemCommandRunner
.run_streaming("sh", &["-c", "echo one; echo two 1>&2; exit 3"], &out)
.await
.unwrap();
assert!(!ok);
let mut lines = out.0.lock().unwrap().clone();
lines.sort();
assert_eq!(lines, vec!["one", "two"]);
}
}

View File

@ -9,7 +9,9 @@ pub mod token;
pub use cipher::AesGcmCipher;
pub use db::{connect, DbPool};
pub use host::{DebianInspector, FakeHostInspector, SystemCommandRunner};
pub use host::{
DebianInspector, DebianUpdater, FakeHostInspector, FakeHostUpdater, SystemCommandRunner,
};
pub use mail::LettreMailer;
pub use password::Argon2Hasher;
pub use sqlite::SqliteInventory;