From 51364fdd7657de9f01e98efa53f10bc0ebbb7e69 Mon Sep 17 00:00:00 2001 From: Dennis Nemec Date: Thu, 3 Sep 2026 20:50:11 +0200 Subject: [PATCH] Discover applications on the cluster and the host for backups Kubernetes workloads are grouped by their Helm instance label into applications, each offering what is worth backing up: data volumes, a PostgreSQL dump instead of the database's own volume, and optionally the namespace manifests. Caches are listed but not preselected. Host applications come from running systemd services that declare a state or working directory. Selecting components creates one strategy each. Co-Authored-By: Claude Opus 5 --- backend/crates/api/src/backups.rs | 52 +- backend/crates/api/src/lib.rs | 3 + backend/crates/api/src/openapi.rs | 1 + backend/crates/api/tests/backups.rs | 61 +++ .../crates/application/src/backup_service.rs | 66 ++- backend/crates/application/src/lib.rs | 2 +- backend/crates/application/src/test_fakes.rs | 30 ++ .../application/src/tests/backup_tests.rs | 126 ++++- backend/crates/domain/src/application.rs | 463 ++++++++++++++++++ backend/crates/domain/src/cluster.rs | 6 + backend/crates/domain/src/host.rs | 25 + backend/crates/domain/src/lib.rs | 1 + backend/crates/domain/src/ports.rs | 4 +- .../crates/infrastructure/src/host/debian.rs | 104 +++- .../crates/infrastructure/src/host/fake.rs | 19 +- backend/crates/infrastructure/src/k8s.rs | 104 +++- 16 files changed, 1050 insertions(+), 17 deletions(-) create mode 100644 backend/crates/domain/src/application.rs diff --git a/backend/crates/api/src/backups.rs b/backend/crates/api/src/backups.rs index 7f65f1c..fcfc687 100644 --- a/backend/crates/api/src/backups.rs +++ b/backend/crates/api/src/backups.rs @@ -1,9 +1,10 @@ //! /api/backups: targets, strategies, records, manual runs. -use application::backup_service::StrategyStatus; +use application::backup_service::{ApplicationBackupPlan, StrategyStatus}; use axum::extract::{Path, State}; use axum::http::StatusCode; use axum::routing::{get, post}; use axum::{Json, Router}; +use domain::application::Application; use domain::backup::{BackupRecord, BackupSource, BackupStrategy, BackupTarget, StorageKind}; use domain::jobs::JobKind; use serde::{Deserialize, Serialize}; @@ -31,6 +32,55 @@ pub fn router() -> Router { ) .route("/strategies/{id}/run", post(run_strategy)) .route("/strategies/{id}/records", get(records)) + .route( + "/applications", + get(list_applications).post(create_from_application), + ) +} + +// ---- applications ---- + +#[derive(Deserialize, ToSchema)] +pub struct ApplicationBackupRequest { + pub application_id: String, + pub component_ids: Vec, + pub schedule: String, + pub target_id: Uuid, + pub retention: u32, + #[serde(default)] + pub passphrase: Option, +} + +#[utoipa::path(get, path = "/api/backups/applications", tag = "backups", security(("bearer" = [])), + responses((status = 200, body = Vec)))] +async fn list_applications( + State(s): State, + _: AuthUser, +) -> Result>, ApiError> { + Ok(Json(s.backups.applications().await?)) +} + +#[utoipa::path(post, path = "/api/backups/applications", tag = "backups", security(("bearer" = [])), + request_body = ApplicationBackupRequest, responses((status = 201, body = Vec), (status = 404), (status = 422)))] +async fn create_from_application( + State(s): State, + _: AdminUser, + Json(req): Json, +) -> Result<(StatusCode, Json>), ApiError> { + let plan = ApplicationBackupPlan { + schedule: req.schedule, + target_id: req.target_id, + retention: req.retention, + passphrase: req.passphrase.filter(|p| !p.is_empty()), + }; + let created = s + .backups + .create_from_application(&req.application_id, &req.component_ids, plan) + .await?; + Ok(( + StatusCode::CREATED, + Json(created.into_iter().map(Into::into).collect()), + )) } // ---- targets ---- diff --git a/backend/crates/api/src/lib.rs b/backend/crates/api/src/lib.rs index 952e229..156f7b7 100644 --- a/backend/crates/api/src/lib.rs +++ b/backend/crates/api/src/lib.rs @@ -130,6 +130,7 @@ impl AppState { } = adapters; let pool_for_findings = pool.clone(); let cluster_gateway = cluster.clone(); + let cluster_gateway_for_backups = cluster.clone(); let users = Arc::new(SqliteUsers(pool.clone())); let hasher = Arc::new(Argon2Hasher); let auth = AuthService::new( @@ -141,6 +142,8 @@ impl AppState { ); let cipher = Arc::new(AesGcmCipher::from_hex(&cfg.master_key)?); let backups = Arc::new(BackupService::new(BackupDeps { + cluster: cluster_gateway_for_backups, + host: inspector.clone(), targets: Arc::new(SqliteBackupTargets(pool.clone())), strategies: Arc::new(SqliteBackupStrategies(pool.clone())), records: Arc::new(SqliteBackupRecords(pool.clone())), diff --git a/backend/crates/api/src/openapi.rs b/backend/crates/api/src/openapi.rs index 61811b3..49afe8f 100644 --- a/backend/crates/api/src/openapi.rs +++ b/backend/crates/api/src/openapi.rs @@ -34,6 +34,7 @@ impl Modify for BearerAuth { crate::backups::delete_target, crate::backups::test_target, crate::backups::list_strategies, crate::backups::get_strategy, crate::backups::create_strategy, crate::backups::update_strategy, crate::backups::delete_strategy, crate::backups::run_strategy, crate::backups::records, + crate::backups::list_applications, crate::backups::create_from_application, crate::dashboard::dashboard, ), modifiers(&BearerAuth) diff --git a/backend/crates/api/tests/backups.rs b/backend/crates/api/tests/backups.rs index 62901c0..c821773 100644 --- a/backend/crates/api/tests/backups.rs +++ b/backend/crates/api/tests/backups.rs @@ -166,3 +166,64 @@ async fn backups_are_admin_only_for_writes() { StatusCode::UNAUTHORIZED ); } + +#[tokio::test] +async fn applications_are_listed_and_can_be_backed_up_in_one_step() { + let app = test_app_with_admin().await; + let token = common::login(&app, ADMIN, PW).await.access; + let target = post(&app, "/api/backups/targets", smb(), Some(&token)).await; + let tid = target.json["id"].as_str().unwrap().to_string(); + + let res = get(&app, "/api/backups/applications", Some(&token)).await; + assert_eq!(res.status, StatusCode::OK, "{}", res.json); + let apps = res.json.as_array().unwrap(); + let gitea = apps + .iter() + .find(|a| a["name"] == "Gitea") + .expect("the gitea release"); + assert_eq!(gitea["kind"], "kubernetes"); + assert_eq!(gitea["namespace"], "gitea"); + let components = gitea["components"].as_array().unwrap(); + let data = components + .iter() + .find(|c| c["id"] == "data:gitea-shared-storage") + .unwrap(); + assert_eq!(data["kind"], "data"); + assert_eq!(data["recommended"], true); + assert!(components + .iter() + .any(|c| c["id"] == "db:gitea-postgresql-0" && c["kind"] == "database")); + assert!(apps.iter().any(|a| a["kind"] == "host")); + + // "Gitea → data + database → 3am → target → save" + let body = json!({ + "application_id": gitea["id"], + "component_ids": ["data:gitea-shared-storage", "db:gitea-postgresql-0"], + "schedule": "0 0 3 * * *", + "target_id": tid, + "retention": 7 + }); + let res = post(&app, "/api/backups/applications", body, Some(&token)).await; + assert_eq!(res.status, StatusCode::CREATED, "{}", res.json); + let created = res.json.as_array().unwrap(); + assert_eq!(created.len(), 2); + assert!(created[0]["name"].as_str().unwrap().starts_with("Gitea – ")); + assert_eq!(created[0]["schedule"], "0 0 3 * * *"); + + let strategies = get(&app, "/api/backups/strategies", Some(&token)).await; + assert_eq!(strategies.json.as_array().unwrap().len(), 2); + + // a viewer may look but not create + 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!( + get(&app, "/api/backups/applications", Some(&user)) + .await + .status, + StatusCode::OK + ); + let res = post(&app, "/api/backups/applications", json!({"application_id": "x", "component_ids": [], "schedule": "0 0 3 * * *", "target_id": tid, "retention": 7}), Some(&user)).await; + assert_eq!(res.status, StatusCode::FORBIDDEN); +} diff --git a/backend/crates/application/src/backup_service.rs b/backend/crates/application/src/backup_service.rs index 4c4ed2f..33ab20e 100644 --- a/backend/crates/application/src/backup_service.rs +++ b/backend/crates/application/src/backup_service.rs @@ -4,10 +4,11 @@ use std::sync::Arc; use async_trait::async_trait; use chrono::{DateTime, Utc}; +use domain::application::Application; use domain::backup::{BackupRecord, BackupStrategy, BackupTarget}; use domain::ports::{ BackupCollector, BackupRecordRepository, BackupStorage, BackupStrategyRepository, - BackupTargetRepository, Cipher, FileEncryptor, + BackupTargetRepository, Cipher, ClusterGateway, FileEncryptor, HostInspector, }; use domain::DomainError; use sha2::{Digest, Sha256}; @@ -16,7 +17,18 @@ use uuid::Uuid; use crate::jobs::{JobHandler, JobLog}; use crate::scheduler::{is_due, validate_cron}; +/// Schedule, target and retention chosen for an application backup. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct ApplicationBackupPlan { + pub schedule: String, + pub target_id: Uuid, + pub retention: u32, + pub passphrase: Option, +} + pub struct BackupDeps { + pub cluster: Arc, + pub host: Arc, pub targets: Arc, pub strategies: Arc, pub records: Arc, @@ -213,6 +225,58 @@ impl BackupService { self.d.strategies.delete(id).await } + /// Applications found on the cluster and on the host, with what they offer to back up. + /// An unreachable cluster or host does not fail the call; that source is then missing. + pub async fn applications(&self) -> Result, DomainError> { + let mut apps = match self.d.cluster.overview().await { + Ok(o) => domain::application::from_cluster(&o), + Err(_) => Vec::new(), + }; + if let Ok(services) = self.d.host.services().await { + apps.extend(domain::application::from_services(&services)); + } + Ok(apps) + } + + /// Turn the selected components of an application into one strategy each. + pub async fn create_from_application( + &self, + application_id: &str, + component_ids: &[String], + plan: ApplicationBackupPlan, + ) -> Result, DomainError> { + if component_ids.is_empty() { + return Err(DomainError::Validation( + "select at least one component".into(), + )); + } + let apps = self.applications().await?; + let app = apps + .iter() + .find(|a| a.id == application_id) + .ok_or(DomainError::NotFound)?; + let mut created = Vec::new(); + for id in component_ids { + let component = app + .components + .iter() + .find(|c| &c.id == id) + .ok_or(DomainError::NotFound)?; + let strategy = BackupStrategy { + id: Uuid::new_v4(), + name: format!("{} – {}", app.name, component.label), + source: component.source.clone(), + schedule: plan.schedule.clone(), + target_id: plan.target_id, + retention: plan.retention, + passphrase: plan.passphrase.clone(), + enabled: true, + }; + created.push(self.create_strategy(strategy).await?); + } + Ok(created) + } + pub async fn records(&self, strategy_id: Uuid) -> Result, DomainError> { self.d.records.list_for(strategy_id, 100).await } diff --git a/backend/crates/application/src/lib.rs b/backend/crates/application/src/lib.rs index c61bda9..1a902bc 100644 --- a/backend/crates/application/src/lib.rs +++ b/backend/crates/application/src/lib.rs @@ -12,7 +12,7 @@ pub mod user_service; pub mod vuln_service; pub use auth_service::AuthService; -pub use backup_service::{BackupDeps, BackupJob, BackupService}; +pub use backup_service::{ApplicationBackupPlan, BackupDeps, BackupJob, BackupService}; pub use cluster_service::ClusterService; pub use image_update_service::{ImageUpdateCheckJob, ImageUpdateService}; pub use inventory_service::{InventoryService, PackageRefreshJob}; diff --git a/backend/crates/application/src/test_fakes.rs b/backend/crates/application/src/test_fakes.rs index 9cea9a4..ff73555 100644 --- a/backend/crates/application/src/test_fakes.rs +++ b/backend/crates/application/src/test_fakes.rs @@ -338,6 +338,22 @@ impl HostInspector for FakeInspector { reboot_required: true, }) } + async fn services(&self) -> Result, DomainError> { + Ok(vec![ + domain::host::HostService { + unit: "monitoring.service".into(), + description: "SoftVisor Infrastructure Monitoring".into(), + working_dir: Some("/opt/monitoring".into()), + state_dir: None, + }, + domain::host::HostService { + unit: "ssh.service".into(), + description: "OpenBSD Secure Shell server".into(), + working_dir: None, + state_dir: None, + }, + ]) + } async fn packages(&self) -> Result, DomainError> { Ok(vec![ Package { @@ -421,6 +437,16 @@ pub struct MemCluster { pub fail: std::sync::atomic::AtomicBool, } +fn helm_labels(instance: &str, role: &str) -> std::collections::BTreeMap { + [ + ("app.kubernetes.io/instance", instance), + ("app.kubernetes.io/name", role), + ] + .into_iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect() +} + pub fn sample_overview() -> ClusterOverview { ClusterOverview { nodes: vec![NodeInfo { @@ -443,6 +469,8 @@ pub fn sample_overview() -> ClusterOverview { name: "gitea".into(), image: "gitea/gitea:1.22.3".into(), }], + labels: helm_labels("gitea", "gitea"), + claims: vec!["gitea-shared-storage".into()], }, Workload { namespace: "gitea".into(), @@ -454,6 +482,8 @@ pub fn sample_overview() -> ClusterOverview { name: "postgresql".into(), image: "bitnami/postgresql:16.4.0".into(), }], + labels: helm_labels("gitea", "postgresql"), + claims: vec!["data-gitea-postgresql-0".into()], }, ], volume_claims: vec![VolumeClaim { diff --git a/backend/crates/application/src/tests/backup_tests.rs b/backend/crates/application/src/tests/backup_tests.rs index 97a41ab..2e353b5 100644 --- a/backend/crates/application/src/tests/backup_tests.rs +++ b/backend/crates/application/src/tests/backup_tests.rs @@ -2,13 +2,15 @@ use std::sync::{Arc, Mutex}; use async_trait::async_trait; use chrono::{Duration, TimeZone, Utc}; +use domain::application::ApplicationKind; use domain::backup::{BackupSource, StorageKind}; use domain::DomainError; +use crate::backup_service::ApplicationBackupPlan; use crate::jobs::{JobHandler, JobLog}; use crate::test_fakes::{ - strategy, target, FakeCipher, FakeCollector, FakeEncryptor, MemRecords, MemStorage, - MemStrategies, MemTargets, + strategy, target, FakeCipher, FakeCollector, FakeEncryptor, FakeInspector, MemCluster, + MemRecords, MemStorage, MemStrategies, MemTargets, }; use crate::{BackupDeps, BackupJob, BackupService}; @@ -39,6 +41,8 @@ fn fixture(fail_storage: bool) -> (F, Arc) { }); let work = tempfile::tempdir().unwrap(); let svc = BackupService::new(BackupDeps { + cluster: Arc::new(MemCluster::default()), + host: Arc::new(FakeInspector { fail: false }), targets: targets.clone(), strategies: strategies.clone(), records: records.clone(), @@ -372,3 +376,121 @@ async fn backup_job_runs_strategy_from_params() { .unwrap_err() .contains("not found")); } + +#[tokio::test] +async fn applications_are_discovered_from_the_cluster_and_the_host() { + let (_f, svc) = fixture(false); + let apps = svc.applications().await.unwrap(); + + let gitea = apps + .iter() + .find(|a| a.name == "Gitea") + .expect("the gitea release"); + assert_eq!(gitea.kind, ApplicationKind::Kubernetes); + assert_eq!(gitea.namespace.as_deref(), Some("gitea")); + let ids: Vec<&str> = gitea.components.iter().map(|c| c.id.as_str()).collect(); + assert!(ids.contains(&"data:gitea-shared-storage"), "{ids:?}"); + assert!(ids.contains(&"db:gitea-postgresql-0"), "{ids:?}"); + assert!(ids.contains(&"manifests"), "{ids:?}"); + + let host = apps + .iter() + .find(|a| a.kind == ApplicationKind::Host) + .expect("a host application"); + assert_eq!(host.name, "Monitoring"); + assert_eq!(host.components.len(), 1); +} + +#[tokio::test] +async fn a_selection_of_components_becomes_one_strategy_each() { + let (f, svc) = fixture(false); + let target = svc.create_target(target("NAS")).await.unwrap(); + let apps = svc.applications().await.unwrap(); + let gitea = apps.iter().find(|a| a.name == "Gitea").unwrap(); + + let created = svc + .create_from_application( + &gitea.id, + &[ + "data:gitea-shared-storage".into(), + "db:gitea-postgresql-0".into(), + ], + ApplicationBackupPlan { + schedule: "0 0 3 * * *".into(), + target_id: target.id, + retention: 7, + passphrase: None, + }, + ) + .await + .unwrap(); + + assert_eq!(created.len(), 2); + assert_eq!( + created[0].name, + "Gitea – Gitea volume gitea-shared-storage (10Gi)" + ); + assert_eq!(created[0].schedule, "0 0 3 * * *"); + assert_eq!(created[0].target_id, target.id); + assert_eq!(created[0].retention, 7); + assert!(created[0].enabled); + assert_eq!( + created[0].source, + BackupSource::VolumeClaim { + namespace: "gitea".into(), + pvc: "gitea-shared-storage".into() + } + ); + assert_eq!( + created[1].source, + BackupSource::PostgresDump { + namespace: "gitea".into(), + pod: "gitea-postgresql-0".into() + } + ); + assert_eq!(f.strategies.0.lock().unwrap().len(), 2, "both are stored"); + + // the strategies are ordinary ones from here on + let listed = svc.list_strategies().await.unwrap(); + assert_eq!(listed.len(), 2); + assert_eq!(listed[0].target_name, "NAS"); +} + +#[tokio::test] +async fn creating_from_an_application_validates_its_input() { + let (_f, svc) = fixture(false); + let target = svc.create_target(target("NAS")).await.unwrap(); + let plan = |schedule: &str| ApplicationBackupPlan { + schedule: schedule.into(), + target_id: target.id, + retention: 7, + passphrase: None, + }; + let apps = svc.applications().await.unwrap(); + let gitea = apps.iter().find(|a| a.name == "Gitea").unwrap().id.clone(); + + assert_eq!( + svc.create_from_application("k8s:nope/nope", &["x".into()], plan("0 0 3 * * *")) + .await + .unwrap_err(), + DomainError::NotFound + ); + assert!(matches!( + svc.create_from_application(&gitea, &[], plan("0 0 3 * * *")) + .await + .unwrap_err(), + DomainError::Validation(_) + )); + assert_eq!( + svc.create_from_application(&gitea, &["data:nothing".into()], plan("0 0 3 * * *")) + .await + .unwrap_err(), + DomainError::NotFound + ); + assert!(matches!( + svc.create_from_application(&gitea, &["manifests".into()], plan("not a cron")) + .await + .unwrap_err(), + DomainError::Validation(_) + )); +} diff --git a/backend/crates/domain/src/application.rs b/backend/crates/domain/src/application.rs new file mode 100644 index 0000000..0d1168d --- /dev/null +++ b/backend/crates/domain/src/application.rs @@ -0,0 +1,463 @@ +//! Applications that run on the server, and what of them is worth backing up. +//! +//! Kubernetes applications are grouped by the Helm/recommended label +//! `app.kubernetes.io/instance`; host applications come from running systemd services that +//! declare a state or working directory. Only applications that have something to back up +//! are reported. +use serde::{Deserialize, Serialize}; + +use crate::backup::BackupSource; +use crate::cluster::ClusterOverview; +use crate::host::HostService; + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum ApplicationKind { + Kubernetes, + Host, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum ComponentKind { + /// Persistent data: a volume or a directory. + Data, + /// A database that is dumped instead of copied. + Database, + /// The Kubernetes objects of the namespace. + Manifests, +} + +/// One backup-worthy part of an application, ready to become a strategy. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct BackupComponent { + pub id: String, + pub kind: ComponentKind, + pub label: String, + pub source: BackupSource, + /// Preselected in the UI. Caches and manifests are not. + pub recommended: bool, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct Application { + pub id: String, + pub name: String, + pub kind: ApplicationKind, + pub namespace: Option, + pub detail: String, + pub components: Vec, +} + +const INSTANCE_LABEL: &str = "app.kubernetes.io/instance"; +const NAME_LABEL: &str = "app.kubernetes.io/name"; + +/// Workload roles whose data is a cache and does not need a backup. +const CACHES: [&str; 4] = ["valkey", "redis", "memcached", "keydb"]; +/// Workload roles that are dumped with `pg_dumpall`. +const POSTGRES: [&str; 2] = ["postgresql", "postgres"]; + +fn title_case(s: &str) -> String { + let mut c = s.chars(); + match c.next() { + Some(f) => f.to_uppercase().collect::() + c.as_str(), + None => String::new(), + } +} + +/// Applications running on the cluster, with the volumes and databases they own. +pub fn from_cluster(overview: &ClusterOverview) -> Vec { + let mut apps: Vec = Vec::new(); + let mut groups: Vec<(String, String)> = Vec::new(); // (namespace, instance) + for w in &overview.workloads { + let instance = w + .labels + .get(INSTANCE_LABEL) + .cloned() + .unwrap_or_else(|| w.name.clone()); + let key = (w.namespace.clone(), instance); + if !groups.contains(&key) { + groups.push(key); + } + } + + for (namespace, instance) in groups { + let members: Vec<_> = overview + .workloads + .iter() + .filter(|w| { + w.namespace == namespace + && w.labels + .get(INSTANCE_LABEL) + .map(|i| i == &instance) + .unwrap_or(w.name == instance) + }) + .collect(); + let mut components = Vec::new(); + + for w in &members { + let role = w + .labels + .get(NAME_LABEL) + .cloned() + .unwrap_or_else(|| w.name.clone()); + let is_cache = CACHES.iter().any(|c| role.contains(c)); + + for claim in &w.claims { + let size = overview + .volume_claims + .iter() + .find(|p| p.namespace == namespace && &p.name == claim) + .map(|p| format!(" ({})", p.capacity)) + .unwrap_or_default(); + components.push(BackupComponent { + id: format!("data:{claim}"), + kind: ComponentKind::Data, + label: format!("{} volume {claim}{size}", title_case(&role)), + source: BackupSource::VolumeClaim { + namespace: namespace.clone(), + pvc: claim.clone(), + }, + recommended: !is_cache, + }); + } + + if POSTGRES.iter().any(|p| role.contains(p)) { + let pod = format!("{}-0", w.name); + components.push(BackupComponent { + id: format!("db:{pod}"), + kind: ComponentKind::Database, + label: format!("PostgreSQL dump of {pod}"), + source: BackupSource::PostgresDump { + namespace: namespace.clone(), + pod, + }, + recommended: true, + }); + } + } + + if components.is_empty() { + continue; // nothing to back up, e.g. a stateless controller + } + // a database is dumped, so its own volume would be a redundant second copy + if components.iter().any(|c| c.kind == ComponentKind::Database) { + components + .retain(|c| !(c.kind == ComponentKind::Data && is_database_volume(c, &members))); + } + // volumes first, then databases, then the manifests + components.sort_by_key(|c| match c.kind { + ComponentKind::Data => 0, + ComponentKind::Database => 1, + ComponentKind::Manifests => 2, + }); + components.push(BackupComponent { + id: "manifests".into(), + kind: ComponentKind::Manifests, + label: format!("Kubernetes objects of namespace {namespace}"), + source: BackupSource::KubernetesManifests { + namespace: namespace.clone(), + }, + recommended: false, + }); + + let roles: Vec = members + .iter() + .map(|w| { + w.labels + .get(NAME_LABEL) + .cloned() + .unwrap_or_else(|| w.name.clone()) + }) + .collect(); + apps.push(Application { + id: format!("k8s:{namespace}/{instance}"), + name: title_case(&instance), + kind: ApplicationKind::Kubernetes, + namespace: Some(namespace.clone()), + detail: format!("{} workload(s): {}", members.len(), roles.join(", ")), + components, + }); + } + apps.sort_by(|a, b| a.name.cmp(&b.name)); + apps +} + +/// True when the volume belongs to a workload that is dumped as a database anyway. +fn is_database_volume(component: &BackupComponent, members: &[&crate::cluster::Workload]) -> bool { + let BackupSource::VolumeClaim { pvc, .. } = &component.source else { + return false; + }; + members.iter().any(|w| { + let role = w + .labels + .get(NAME_LABEL) + .cloned() + .unwrap_or_else(|| w.name.clone()); + POSTGRES.iter().any(|p| role.contains(p)) && w.claims.contains(pvc) + }) +} + +/// Units that are part of the operating system or the container runtime rather than an +/// application whose data a user would back up. +fn is_system_unit(unit: &str) -> bool { + [ + "snap.", + "systemd-", + "user@", + "getty@", + "dbus", + "cron", + "ssh", + "qemu-", + "polkit", + "rsyslog", + "unattended", + ] + .iter() + .any(|p| unit.starts_with(p)) +} + +/// Applications installed directly on the host: running services that declare a directory. +pub fn from_services(services: &[HostService]) -> Vec { + let mut apps: Vec = services + .iter() + .filter(|s| !is_system_unit(&s.unit)) + .filter_map(|s| { + let path = s.data_dir()?; + let name = s.unit.trim_end_matches(".service").to_string(); + Some(Application { + id: format!("host:{}", s.unit), + name: title_case(&name), + kind: ApplicationKind::Host, + namespace: None, + detail: s.description.clone(), + components: vec![BackupComponent { + id: format!("data:{path}"), + kind: ComponentKind::Data, + label: format!("Directory {path}"), + source: BackupSource::HostPath { path }, + recommended: true, + }], + }) + }) + .collect(); + apps.sort_by(|a, b| a.name.cmp(&b.name)); + apps +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::cluster::{Container, NodeInfo, VolumeClaim, Workload, WorkloadKind}; + use std::collections::BTreeMap; + + fn labels(pairs: &[(&str, &str)]) -> BTreeMap { + pairs + .iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect() + } + + fn workload(ns: &str, name: &str, role: &str, instance: &str, claims: &[&str]) -> Workload { + Workload { + namespace: ns.into(), + kind: WorkloadKind::Deployment, + name: name.into(), + ready: 1, + desired: 1, + containers: vec![Container { + name: name.into(), + image: format!("{role}:1"), + }], + labels: labels(&[(INSTANCE_LABEL, instance), (NAME_LABEL, role)]), + claims: claims.iter().map(|c| c.to_string()).collect(), + } + } + + fn overview(workloads: Vec, claims: Vec<(&str, &str, &str)>) -> ClusterOverview { + ClusterOverview { + nodes: vec![NodeInfo { + name: "n".into(), + version: "v1".into(), + ready: true, + os_image: String::new(), + kernel: String::new(), + container_runtime: String::new(), + }], + namespaces: vec!["gitea".into()], + workloads, + volume_claims: claims + .into_iter() + .map(|(ns, name, cap)| VolumeClaim { + namespace: ns.into(), + name: name.into(), + capacity: cap.into(), + storage_class: "hostpath".into(), + status: "Bound".into(), + }) + .collect(), + fetched_at: chrono::Utc::now(), + } + } + + #[test] + fn groups_the_workloads_of_a_helm_release_into_one_application() { + let o = overview( + vec![ + workload( + "gitea", + "gitea", + "gitea", + "gitea", + &["gitea-shared-storage"], + ), + workload( + "gitea", + "gitea-postgresql", + "postgresql", + "gitea", + &["data-gitea-postgresql-0"], + ), + workload( + "gitea", + "gitea-valkey-primary", + "valkey", + "gitea", + &["valkey-data-gitea-valkey-primary-0"], + ), + ], + vec![ + ("gitea", "gitea-shared-storage", "10Gi"), + ("gitea", "data-gitea-postgresql-0", "10Gi"), + ], + ); + let apps = from_cluster(&o); + assert_eq!(apps.len(), 1); + let app = &apps[0]; + assert_eq!(app.id, "k8s:gitea/gitea"); + assert_eq!(app.name, "Gitea"); + assert_eq!(app.namespace.as_deref(), Some("gitea")); + assert!(app.detail.contains("postgresql"), "{}", app.detail); + + let ids: Vec<&str> = app.components.iter().map(|c| c.id.as_str()).collect(); + assert_eq!( + ids, + vec![ + "data:gitea-shared-storage", + "data:valkey-data-gitea-valkey-primary-0", + "db:gitea-postgresql-0", + "manifests" + ], + "the database volume is dropped in favour of the dump" + ); + + let data = &app.components[0]; + assert_eq!(data.kind, ComponentKind::Data); + assert!(data.label.contains("gitea-shared-storage") && data.label.contains("10Gi")); + assert!(data.recommended); + assert_eq!( + data.source, + BackupSource::VolumeClaim { + namespace: "gitea".into(), + pvc: "gitea-shared-storage".into() + } + ); + + let cache = &app.components[1]; + assert!(!cache.recommended, "a cache volume is not preselected"); + + let db = &app.components[2]; + assert_eq!(db.kind, ComponentKind::Database); + assert_eq!( + db.source, + BackupSource::PostgresDump { + namespace: "gitea".into(), + pod: "gitea-postgresql-0".into() + } + ); + assert!(db.recommended); + + let manifests = app.components.last().unwrap(); + assert_eq!(manifests.kind, ComponentKind::Manifests); + assert!(!manifests.recommended); + } + + #[test] + fn applications_without_anything_to_back_up_are_skipped() { + let o = overview( + vec![workload( + "ingress", + "controller", + "ingress-nginx", + "ingress", + &[], + )], + vec![], + ); + assert!(from_cluster(&o).is_empty()); + } + + #[test] + fn workloads_without_helm_labels_stand_on_their_own() { + let mut w = workload("apps", "legacy", "legacy", "legacy", &["legacy-data"]); + w.labels.clear(); + let apps = from_cluster(&overview(vec![w], vec![("apps", "legacy-data", "5Gi")])); + assert_eq!(apps.len(), 1); + assert_eq!(apps[0].id, "k8s:apps/legacy"); + assert_eq!(apps[0].name, "Legacy"); + } + + #[test] + fn host_services_with_a_directory_become_applications() { + let services = vec![ + HostService { + unit: "monitoring.service".into(), + description: "SoftVisor Infrastructure Monitoring".into(), + working_dir: Some("/opt/monitoring".into()), + state_dir: None, + }, + HostService { + unit: "tailscaled.service".into(), + description: "Tailscale node agent".into(), + working_dir: None, + state_dir: Some("tailscale".into()), + }, + // no directory to back up + HostService { + unit: "ssh.service".into(), + description: "OpenSSH".into(), + working_dir: None, + state_dir: None, + }, + // the container runtime is not an application + HostService { + unit: "snap.microk8s.daemon-kubelite.service".into(), + description: "microk8s".into(), + working_dir: Some("/var/snap/microk8s/8702".into()), + state_dir: None, + }, + ]; + let apps = from_services(&services); + assert_eq!( + apps.iter().map(|a| a.name.as_str()).collect::>(), + vec!["Monitoring", "Tailscaled"] + ); + assert_eq!(apps[0].id, "host:monitoring.service"); + assert_eq!(apps[0].kind, ApplicationKind::Host); + assert_eq!(apps[0].detail, "SoftVisor Infrastructure Monitoring"); + assert_eq!( + apps[0].components[0].source, + BackupSource::HostPath { + path: "/opt/monitoring".into() + } + ); + assert_eq!( + apps[1].components[0].source, + BackupSource::HostPath { + path: "/var/lib/tailscale".into() + }, + "a state directory is relative to /var/lib" + ); + } +} diff --git a/backend/crates/domain/src/cluster.rs b/backend/crates/domain/src/cluster.rs index 998a86a..e1d0028 100644 --- a/backend/crates/domain/src/cluster.rs +++ b/backend/crates/domain/src/cluster.rs @@ -39,6 +39,12 @@ pub struct Workload { pub ready: i32, pub desired: i32, pub containers: Vec, + /// Object labels; `app.kubernetes.io/instance` groups the workloads of one application. + #[serde(default)] + pub labels: std::collections::BTreeMap, + /// Persistent volume claims this workload uses. + #[serde(default)] + pub claims: Vec, } #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] diff --git a/backend/crates/domain/src/host.rs b/backend/crates/domain/src/host.rs index aab59b1..acf34b8 100644 --- a/backend/crates/domain/src/host.rs +++ b/backend/crates/domain/src/host.rs @@ -54,6 +54,31 @@ impl Inventory { } } +/// A systemd service running on the host. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct HostService { + pub unit: String, + pub description: String, + /// `WorkingDirectory=` of the unit, if it is an absolute path other than `/`. + pub working_dir: Option, + /// `StateDirectory=` of the unit; systemd places it under `/var/lib`. + pub state_dir: Option, +} + +impl HostService { + /// The directory holding this service's data, if it declares one. + pub fn data_dir(&self) -> Option { + if let Some(state) = self.state_dir.as_ref().filter(|s| !s.is_empty()) { + let first = state.split_whitespace().next()?; + return Some(format!("/var/lib/{first}")); + } + self.working_dir + .as_ref() + .filter(|w| w.starts_with('/') && w.as_str() != "/") + .cloned() + } +} + /// 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 diff --git a/backend/crates/domain/src/lib.rs b/backend/crates/domain/src/lib.rs index e3675a9..156c471 100644 --- a/backend/crates/domain/src/lib.rs +++ b/backend/crates/domain/src/lib.rs @@ -1,5 +1,6 @@ //! Domain layer: entities, value objects, errors and the ports (traits) the application //! layer depends on. No I/O here. +pub mod application; pub mod auth; pub mod backup; pub mod cluster; diff --git a/backend/crates/domain/src/ports.rs b/backend/crates/domain/src/ports.rs index 3748892..64e272c 100644 --- a/backend/crates/domain/src/ports.rs +++ b/backend/crates/domain/src/ports.rs @@ -5,7 +5,7 @@ use uuid::Uuid; 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::host::{HostService, Inventory, OsInfo, Package}; use crate::image::ImageUpdate; use crate::jobs::{JobKind, JobRun, JobStatus}; use crate::settings::SmtpSettings; @@ -92,6 +92,8 @@ pub trait JobRunRepository: Send + Sync { pub trait HostInspector: Send + Sync { async fn os_info(&self) -> Result; async fn packages(&self) -> Result, DomainError>; + /// Running systemd services, for discovering applications installed on the host. + async fn services(&self) -> Result, DomainError>; } #[async_trait] diff --git a/backend/crates/infrastructure/src/host/debian.rs b/backend/crates/infrastructure/src/host/debian.rs index ea5b124..a21c781 100644 --- a/backend/crates/infrastructure/src/host/debian.rs +++ b/backend/crates/infrastructure/src/host/debian.rs @@ -2,7 +2,7 @@ use std::sync::Arc; use async_trait::async_trait; -use domain::host::{OsInfo, Package, PackageSource}; +use domain::host::{HostService, OsInfo, Package, PackageSource}; use domain::ports::HostInspector; use domain::DomainError; @@ -39,6 +39,50 @@ impl HostInspector for DebianInspector { }) } + async fn services(&self) -> Result, DomainError> { + let out = self + .runner + .run( + "systemctl", + &[ + "list-units", + "--type=service", + "--state=running", + "--no-legend", + "--no-pager", + ], + ) + .await?; + if !out.success { + return Err(DomainError::Unavailable("systemctl failed".into())); + } + let mut services = Vec::new(); + for (unit, description) in parse_service_units(&out.stdout) { + let props = self + .runner + .run( + "systemctl", + &[ + "show", + &unit, + "-p", + "WorkingDirectory", + "-p", + "StateDirectory", + ], + ) + .await?; + let (working_dir, state_dir) = parse_service_properties(&props.stdout); + services.push(HostService { + unit, + description, + working_dir, + state_dir, + }); + } + Ok(services) + } + async fn packages(&self) -> Result, DomainError> { let r = &self.runner; let dpkg = r @@ -139,6 +183,34 @@ pub fn parse_snap_list(out: &str) -> Vec<(String, String)> { .collect() } +/// `systemctl list-units --type=service --state=running --no-legend` → (unit, description). +pub fn parse_service_units(out: &str) -> Vec<(String, String)> { + out.lines() + .filter_map(|l| { + let mut parts = l.split_whitespace(); + let unit = parts.next()?.to_string(); + if !unit.ends_with(".service") { + return None; + } + // loaded / active / running, then the description + let description = parts.skip(3).collect::>().join(" "); + Some((unit, description)) + }) + .collect() +} + +/// `systemctl show -p WorkingDirectory -p StateDirectory` → (working dir, state dir). +pub fn parse_service_properties(out: &str) -> (Option, Option) { + let value = |key: &str| { + out.lines() + .find_map(|l| l.strip_prefix(&format!("{key}="))) + .map(str::trim) + .filter(|v| !v.is_empty()) + .map(String::from) + }; + (value("WorkingDirectory"), value("StateDirectory")) +} + /// `/etc/os-release` → (PRETTY_NAME, VERSION_ID). pub fn parse_os_release(content: &str) -> (String, String) { let get = |key: &str| { @@ -194,6 +266,36 @@ Conf openssl (3.0.16-1~deb12u1 Debian-Security:12/stable-security [amd64])\n"; ); } + #[test] + fn reads_running_services_and_their_directories() { + // `systemctl list-units --type=service --state=running --no-legend` + let units = " monitoring.service loaded active running SoftVisor Infrastructure Monitoring\n ssh.service loaded active running OpenBSD Secure Shell server\n"; + assert_eq!( + parse_service_units(units), + vec![ + ( + "monitoring.service".to_string(), + "SoftVisor Infrastructure Monitoring".to_string() + ), + ( + "ssh.service".to_string(), + "OpenBSD Secure Shell server".to_string() + ) + ] + ); + // `systemctl show -p WorkingDirectory -p StateDirectory ` + let props = "WorkingDirectory=/opt/monitoring\nStateDirectory=\n"; + assert_eq!( + parse_service_properties(props), + (Some("/opt/monitoring".to_string()), None) + ); + assert_eq!( + parse_service_properties("WorkingDirectory=\nStateDirectory=tailscale\n"), + (None, Some("tailscale".to_string())) + ); + assert_eq!(parse_service_properties(""), (None, None)); + } + #[test] fn os_release_strips_quotes() { let c = "PRETTY_NAME=\"Debian GNU/Linux 12 (bookworm)\"\nNAME=\"Debian GNU/Linux\"\nVERSION_ID=\"12\"\n"; diff --git a/backend/crates/infrastructure/src/host/fake.rs b/backend/crates/infrastructure/src/host/fake.rs index e1acada..0d300c2 100644 --- a/backend/crates/infrastructure/src/host/fake.rs +++ b/backend/crates/infrastructure/src/host/fake.rs @@ -1,6 +1,6 @@ //! Sample data for development machines without apt (FAKE_HOST=true). use async_trait::async_trait; -use domain::host::{OsInfo, Package, PackageSource}; +use domain::host::{HostService, OsInfo, Package, PackageSource}; use domain::ports::HostInspector; use domain::DomainError; @@ -19,6 +19,23 @@ impl HostInspector for FakeHostInspector { }) } + async fn services(&self) -> Result, DomainError> { + Ok(vec![ + HostService { + unit: "monitoring.service".into(), + description: "SoftVisor Infrastructure Monitoring".into(), + working_dir: Some("/opt/monitoring".into()), + state_dir: None, + }, + HostService { + unit: "ssh.service".into(), + description: "OpenBSD Secure Shell server".into(), + working_dir: None, + state_dir: None, + }, + ]) + } + async fn packages(&self) -> Result, DomainError> { let apt = |n: &str, i: &str, c: Option<&str>, s: bool| Package { name: n.into(), diff --git a/backend/crates/infrastructure/src/k8s.rs b/backend/crates/infrastructure/src/k8s.rs index 577e2f7..4e31c49 100644 --- a/backend/crates/infrastructure/src/k8s.rs +++ b/backend/crates/infrastructure/src/k8s.rs @@ -59,6 +59,51 @@ fn map_err(e: kube::Error) -> DomainError { } } +fn labels(meta: &kube::core::ObjectMeta) -> std::collections::BTreeMap { + meta.labels + .clone() + .unwrap_or_default() + .into_iter() + .collect() +} + +/// PVCs referenced by the pod template, plus the claims a statefulset creates from its +/// templates (`