diff --git a/.env.example b/.env.example index 1698050..3f609df 100644 --- a/.env.example +++ b/.env.example @@ -14,3 +14,7 @@ COOKIE_SECURE=false # Directory with the built frontend (served as SPA fallback) FRONTEND_DIR=../frontend/dist RUST_LOG=info,sqlx=warn +# kubectl invocation used for backups (default: microk8s kubectl if present) +#KUBECTL=/snap/bin/microk8s kubectl +# scratch directory for backup archives (default: data/work) +#WORK_DIR=data/work diff --git a/ROADMAP.md b/ROADMAP.md index 0f6a9c2..82395ba 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -298,8 +298,8 @@ WP-02 remainder, WP-10, WP-11, WP-12. Verified against the real host: 430 packag ### Milestone 3 – Vulnerability management (delivered 2026-09-02) WP-20, WP-21. Trivy is installed by the deploy script. -### Milestone 4 – Backup management (planned) -WP-30, WP-31, WP-32. +### Milestone 4 – Backup management (delivered 2026-09-02) +WP-30, WP-31, WP-32. Restore procedure in `docs/restore.md`. ### Milestone 5 – Hardening and release (planned) WP-40, WP-41, WP-42. @@ -360,9 +360,9 @@ The server is reachable via `ssh softvisor` (as root). Findings from the inspect | WP-12 | M2 | done | 2026-09-02; overview, restart, scale, set image; upstream tag check not done | | WP-20 | M3 | done | 2026-09-02; Trivy rootfs + image scans, diff with first/last seen, fixed detection | | WP-21 | M3 | done | 2026-09-02; findings UI, acknowledge, mail digest with severity threshold; dashboard widget in WP-40 | -| WP-30 | M4 | todo | | -| WP-31 | M4 | todo | | -| WP-32 | M4 | todo | | +| WP-30 | M4 | done | 2026-09-02; SMB via smbclient, FTP/FTPS via curl, credentials encrypted | +| WP-31 | M4 | done | 2026-09-02; sources: PVC hostpath, pg_dumpall, manifests, host dir; openssl encryption | +| WP-32 | M4 | done | 2026-09-02; upload verify, retention, records; docs/restore.md; failure mail not yet | | WP-40 | M5 | todo | | | WP-41 | M5 | todo | | | WP-42 | M5 | todo | | diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 59ec1eb..0d58ade 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -1100,6 +1100,7 @@ dependencies = [ "serde", "serde_json", "sqlx", + "tempfile", "tokio", "uuid", ] diff --git a/backend/crates/api/src/backups.rs b/backend/crates/api/src/backups.rs new file mode 100644 index 0000000..7f65f1c --- /dev/null +++ b/backend/crates/api/src/backups.rs @@ -0,0 +1,350 @@ +//! /api/backups: targets, strategies, records, manual runs. +use application::backup_service::StrategyStatus; +use axum::extract::{Path, State}; +use axum::http::StatusCode; +use axum::routing::{get, post}; +use axum::{Json, Router}; +use domain::backup::{BackupRecord, BackupSource, BackupStrategy, BackupTarget, StorageKind}; +use domain::jobs::JobKind; +use serde::{Deserialize, Serialize}; +use utoipa::ToSchema; +use uuid::Uuid; + +use crate::error::ApiError; +use crate::extract::{AdminUser, AuthUser}; +use crate::AppState; + +pub fn router() -> Router { + Router::new() + .route("/targets", get(list_targets).post(create_target)) + .route( + "/targets/{id}", + get(get_target).put(update_target).delete(delete_target), + ) + .route("/targets/{id}/test", post(test_target)) + .route("/strategies", get(list_strategies).post(create_strategy)) + .route( + "/strategies/{id}", + get(get_strategy) + .put(update_strategy) + .delete(delete_strategy), + ) + .route("/strategies/{id}/run", post(run_strategy)) + .route("/strategies/{id}/records", get(records)) +} + +// ---- targets ---- + +#[derive(Serialize, ToSchema)] +pub struct TargetView { + pub id: Uuid, + pub name: String, + #[schema(value_type = String)] + pub kind: StorageKind, + pub host: String, + pub port: Option, + pub share: String, + pub path: String, + pub username: String, + pub tls: bool, + pub password_set: bool, +} + +impl From for TargetView { + fn from(t: BackupTarget) -> Self { + Self { + id: t.id, + name: t.name, + kind: t.kind, + host: t.host, + port: t.port, + share: t.share, + path: t.path, + username: t.username, + tls: t.tls, + password_set: !t.password.is_empty(), + } + } +} + +#[derive(Deserialize, ToSchema)] +pub struct TargetRequest { + pub name: String, + #[schema(value_type = String)] + pub kind: StorageKind, + pub host: String, + pub port: Option, + #[serde(default)] + pub share: String, + #[serde(default)] + pub path: String, + #[serde(default)] + pub username: String, + /// Empty keeps the stored password on update. + #[serde(default)] + pub password: String, + #[serde(default)] + pub tls: bool, +} + +impl TargetRequest { + fn into_target(self, id: Uuid) -> BackupTarget { + BackupTarget { + id, + name: self.name.trim().into(), + kind: self.kind, + host: self.host.trim().into(), + port: self.port, + share: self.share.trim().into(), + path: self.path.trim().into(), + username: self.username, + password: self.password, + tls: self.tls, + } + } +} + +#[utoipa::path(get, path = "/api/backups/targets", tag = "backups", security(("bearer" = [])), responses((status = 200, body = Vec)))] +async fn list_targets( + State(s): State, + _: AuthUser, +) -> Result>, ApiError> { + Ok(Json( + s.backups + .list_targets() + .await? + .into_iter() + .map(Into::into) + .collect(), + )) +} + +#[utoipa::path(get, path = "/api/backups/targets/{id}", tag = "backups", security(("bearer" = [])), responses((status = 200, body = TargetView), (status = 404)))] +async fn get_target( + State(s): State, + _: AuthUser, + Path(id): Path, +) -> Result, ApiError> { + Ok(Json(s.backups.get_target(id).await?.into())) +} + +#[utoipa::path(post, path = "/api/backups/targets", tag = "backups", security(("bearer" = [])), request_body = TargetRequest, responses((status = 201, body = TargetView), (status = 422)))] +async fn create_target( + State(s): State, + _: AdminUser, + Json(req): Json, +) -> Result<(StatusCode, Json), ApiError> { + Ok(( + StatusCode::CREATED, + Json( + s.backups + .create_target(req.into_target(Uuid::nil())) + .await? + .into(), + ), + )) +} + +#[utoipa::path(put, path = "/api/backups/targets/{id}", tag = "backups", security(("bearer" = [])), request_body = TargetRequest, responses((status = 200, body = TargetView), (status = 404), (status = 422)))] +async fn update_target( + State(s): State, + _: AdminUser, + Path(id): Path, + Json(req): Json, +) -> Result, ApiError> { + Ok(Json( + s.backups.update_target(req.into_target(id)).await?.into(), + )) +} + +#[utoipa::path(delete, path = "/api/backups/targets/{id}", tag = "backups", security(("bearer" = [])), responses((status = 204), (status = 404), (status = 409)))] +async fn delete_target( + State(s): State, + _: AdminUser, + Path(id): Path, +) -> Result { + s.backups.delete_target(id).await?; + Ok(StatusCode::NO_CONTENT) +} + +#[utoipa::path(post, path = "/api/backups/targets/{id}/test", tag = "backups", security(("bearer" = [])), responses((status = 204), (status = 404), (status = 502)))] +async fn test_target( + State(s): State, + _: AdminUser, + Path(id): Path, +) -> Result { + s.backups.test_target(id).await?; + Ok(StatusCode::NO_CONTENT) +} + +// ---- strategies ---- + +#[derive(Serialize, ToSchema)] +pub struct StrategyView { + pub id: Uuid, + pub name: String, + #[schema(value_type = Object)] + pub source: BackupSource, + pub schedule: String, + pub target_id: Uuid, + pub retention: u32, + pub encrypted: bool, + pub enabled: bool, +} + +impl From for StrategyView { + fn from(s: BackupStrategy) -> Self { + Self { + id: s.id, + name: s.name, + source: s.source, + schedule: s.schedule, + target_id: s.target_id, + retention: s.retention, + encrypted: s.passphrase.is_some(), + enabled: s.enabled, + } + } +} + +#[derive(Serialize, ToSchema)] +pub struct StrategyStatusView { + #[serde(flatten)] + pub strategy: StrategyView, + pub target_name: String, + #[schema(value_type = Option)] + pub last_backup: Option, +} + +impl From for StrategyStatusView { + fn from(s: StrategyStatus) -> Self { + Self { + strategy: s.strategy.into(), + target_name: s.target_name, + last_backup: s.last_backup, + } + } +} + +#[derive(Deserialize, ToSchema)] +pub struct StrategyRequest { + pub name: String, + #[schema(value_type = Object)] + pub source: BackupSource, + pub schedule: String, + pub target_id: Uuid, + pub retention: u32, + /// `null` = no encryption; empty string keeps the stored passphrase on update. + #[serde(default)] + pub passphrase: Option, + #[serde(default = "default_true")] + pub enabled: bool, +} + +fn default_true() -> bool { + true +} + +impl StrategyRequest { + fn into_strategy(self, id: Uuid) -> BackupStrategy { + BackupStrategy { + id, + name: self.name.trim().into(), + source: self.source, + schedule: self.schedule.trim().into(), + target_id: self.target_id, + retention: self.retention, + passphrase: self.passphrase, + enabled: self.enabled, + } + } +} + +#[utoipa::path(get, path = "/api/backups/strategies", tag = "backups", security(("bearer" = [])), responses((status = 200, body = Vec)))] +async fn list_strategies( + State(s): State, + _: AuthUser, +) -> Result>, ApiError> { + Ok(Json( + s.backups + .list_strategies() + .await? + .into_iter() + .map(Into::into) + .collect(), + )) +} + +#[utoipa::path(get, path = "/api/backups/strategies/{id}", tag = "backups", security(("bearer" = [])), responses((status = 200, body = StrategyView), (status = 404)))] +async fn get_strategy( + State(s): State, + _: AuthUser, + Path(id): Path, +) -> Result, ApiError> { + Ok(Json(s.backups.get_strategy(id).await?.into())) +} + +#[utoipa::path(post, path = "/api/backups/strategies", tag = "backups", security(("bearer" = [])), request_body = StrategyRequest, responses((status = 201, body = StrategyView), (status = 404), (status = 422)))] +async fn create_strategy( + State(s): State, + _: AdminUser, + Json(req): Json, +) -> Result<(StatusCode, Json), ApiError> { + Ok(( + StatusCode::CREATED, + Json( + s.backups + .create_strategy(req.into_strategy(Uuid::nil())) + .await? + .into(), + ), + )) +} + +#[utoipa::path(put, path = "/api/backups/strategies/{id}", tag = "backups", security(("bearer" = [])), request_body = StrategyRequest, responses((status = 200, body = StrategyView), (status = 404), (status = 422)))] +async fn update_strategy( + State(s): State, + _: AdminUser, + Path(id): Path, + Json(req): Json, +) -> Result, ApiError> { + Ok(Json( + s.backups + .update_strategy(req.into_strategy(id)) + .await? + .into(), + )) +} + +#[utoipa::path(delete, path = "/api/backups/strategies/{id}", tag = "backups", security(("bearer" = [])), responses((status = 204), (status = 404)))] +async fn delete_strategy( + State(s): State, + _: AdminUser, + Path(id): Path, +) -> Result { + s.backups.delete_strategy(id).await?; + Ok(StatusCode::NO_CONTENT) +} + +#[utoipa::path(post, path = "/api/backups/strategies/{id}/run", tag = "backups", security(("bearer" = [])), responses((status = 202, body = crate::jobs::JobRunDto), (status = 404), (status = 409)))] +async fn run_strategy( + State(s): State, + AdminUser(admin): AdminUser, + Path(id): Path, +) -> Result<(StatusCode, Json), ApiError> { + s.backups.get_strategy(id).await?; + let run = s + .jobs + .start(JobKind::Backup, Some(id.to_string()), &admin.email) + .await?; + Ok((StatusCode::ACCEPTED, Json(run.into()))) +} + +#[utoipa::path(get, path = "/api/backups/strategies/{id}/records", tag = "backups", security(("bearer" = [])), responses((status = 200, body = Vec)))] +async fn records( + State(s): State, + _: AuthUser, + Path(id): Path, +) -> Result>, ApiError> { + Ok(Json(s.backups.records(id).await?)) +} diff --git a/backend/crates/api/src/config.rs b/backend/crates/api/src/config.rs index 0031495..87a06ae 100644 --- a/backend/crates/api/src/config.rs +++ b/backend/crates/api/src/config.rs @@ -11,6 +11,10 @@ pub struct Config { pub fake_host: bool, /// Path to a kubeconfig; None infers. Defaults to the microk8s client config if present. pub kubeconfig: Option, + /// kubectl invocation for backups, e.g. ["/snap/bin/microk8s", "kubectl"]. + pub kubectl: Vec, + /// Scratch directory for backup archives. + pub work_dir: std::path::PathBuf, pub bind: SocketAddr, pub bootstrap_admin: Option<(String, String)>, pub cookie_secure: bool, @@ -34,6 +38,18 @@ impl Config { jwt_secret, master_key, fake_host: env("FAKE_HOST").is_some_and(|v| v == "true" || v == "1"), + kubectl: env("KUBECTL") + .map(|v| v.split_whitespace().map(String::from).collect()) + .unwrap_or_else(|| { + if std::path::Path::new("/snap/bin/microk8s").exists() { + vec!["/snap/bin/microk8s".into(), "kubectl".into()] + } else { + vec!["kubectl".into()] + } + }), + work_dir: env("WORK_DIR") + .map(Into::into) + .unwrap_or_else(|| "data/work".into()), kubeconfig: env("KUBECONFIG").or_else(|| { let microk8s = "/var/snap/microk8s/current/credentials/client.config"; std::path::Path::new(microk8s) diff --git a/backend/crates/api/src/lib.rs b/backend/crates/api/src/lib.rs index ac8d0a5..7158cf0 100644 --- a/backend/crates/api/src/lib.rs +++ b/backend/crates/api/src/lib.rs @@ -1,5 +1,6 @@ //! HTTP API layer (axum). `build_app` is used by both the binary and the integration tests. pub mod auth; +pub mod backups; pub mod cluster; pub mod config; pub mod error; @@ -17,18 +18,24 @@ use std::sync::Arc; use application::scheduler::Scheduler; use application::{ - AuthService, ClusterService, InventoryService, JobRunner, PackageRefreshJob, PackageUpgradeJob, - SettingsService, UserService, VulnerabilityScanJob, VulnerabilityService, + AuthService, BackupDeps, BackupJob, BackupService, ClusterService, 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::{ClusterGateway, HostInspector, HostUpdater, VulnerabilityScanner}; +use domain::ports::{ + BackupCollector, BackupStorage, ClusterGateway, FileEncryptor, HostInspector, HostUpdater, + VulnerabilityScanner, +}; use infrastructure::{ - AesGcmCipher, Argon2Hasher, DbPool, DebianInspector, DebianUpdater, FakeClusterGateway, - FakeHostInspector, FakeHostUpdater, FakeScanner, JwtIssuer, KubeGateway, LettreMailer, - SqliteAuditLog, SqliteFindings, SqliteInventory, SqliteJobRuns, SqliteRefreshTokens, - SqliteSettings, SqliteUsers, SystemCommandRunner, TrivyScanner, + 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, }; use tower_http::services::{ServeDir, ServeFile}; use tower_http::trace::TraceLayer; @@ -42,6 +49,9 @@ pub struct Adapters { pub updater: Arc, pub cluster: Arc, pub scanner: Arc, + pub storage: Arc, + pub collector: Arc, + pub encryptor: Arc, } #[derive(Clone)] @@ -54,6 +64,7 @@ pub struct AppState { pub inventory: Arc, pub cluster: Arc, pub vulns: Arc, + pub backups: Arc, pub login_limiter: Arc, } @@ -68,6 +79,9 @@ impl AppState { updater: Arc::new(FakeHostUpdater), cluster: Arc::new(FakeClusterGateway), scanner: Arc::new(FakeScanner), + storage: Arc::new(DirBackupStorage::new(cfg.work_dir.join("fake-remote"))), + collector: Arc::new(FakeBackupCollector), + encryptor: Arc::new(OpensslEncryptor::new(runner.clone())), } } else { Adapters { @@ -75,7 +89,13 @@ impl AppState { inspector: Arc::new(DebianInspector::new(runner.clone())), updater: Arc::new(DebianUpdater::new(runner.clone())), cluster: Arc::new(KubeGateway::new(cfg.kubeconfig.clone())), - scanner: Arc::new(TrivyScanner::new(runner)), + scanner: Arc::new(TrivyScanner::new(runner.clone())), + storage: Arc::new(CommandBackupStorage::new(runner.clone())), + collector: Arc::new(KubeBackupCollector::new( + runner.clone(), + cfg.kubectl.clone(), + )), + encryptor: Arc::new(OpensslEncryptor::new(runner)), } }; Self::with_adapters(cfg, pool, adapters, |r| r) @@ -94,6 +114,9 @@ impl AppState { updater, cluster, scanner, + storage, + collector, + encryptor, } = adapters; let pool_for_findings = pool.clone(); let cluster_gateway = cluster.clone(); @@ -107,6 +130,16 @@ impl AppState { Arc::new(JwtIssuer::new(&cfg.jwt_secret)), ); let cipher = Arc::new(AesGcmCipher::from_hex(&cfg.master_key)?); + let backups = Arc::new(BackupService::new(BackupDeps { + targets: Arc::new(SqliteBackupTargets(pool.clone())), + strategies: Arc::new(SqliteBackupStrategies(pool.clone())), + records: Arc::new(SqliteBackupRecords(pool.clone())), + storage, + collector, + encryptor, + cipher: cipher.clone(), + work_dir: cfg.work_dir.clone(), + })); let settings = Arc::new(SettingsService::new( Arc::new(SqliteSettings(pool.clone())), cipher, @@ -135,10 +168,14 @@ impl AppState { cluster_gateway, settings.clone(), )); - let jobs = Arc::new(register_jobs(runner.register( - JobKind::VulnerabilityScan, - Arc::new(VulnerabilityScanJob(vulns.clone())), - ))); + let jobs = Arc::new(register_jobs( + runner + .register( + JobKind::VulnerabilityScan, + Arc::new(VulnerabilityScanJob(vulns.clone())), + ) + .register(JobKind::Backup, Arc::new(BackupJob(backups.clone()))), + )); Ok(Self { cfg, auth: Arc::new(auth), @@ -148,6 +185,7 @@ impl AppState { inventory, cluster, vulns, + backups, login_limiter: Arc::new(rate_limit::RateLimiter::new( 10, std::time::Duration::from_secs(60), @@ -157,7 +195,11 @@ impl AppState { /// Spawn the cron scheduler on the current runtime. pub fn start_scheduler(&self) { - tokio::spawn(Scheduler::new(self.jobs.clone(), self.settings.clone()).run()); + tokio::spawn( + Scheduler::new(self.jobs.clone(), self.settings.clone()) + .with_backups(self.backups.clone()) + .run(), + ); } pub async fn bootstrap(&self) -> anyhow::Result<()> { @@ -187,6 +229,7 @@ pub fn build_app(state: AppState) -> Router { .nest("/api/system", system::router()) .nest("/api/cluster", cluster::router()) .nest("/api/vulnerabilities", vulnerabilities::router()) + .nest("/api/backups", backups::router()) .fallback_service(spa) .layer(TraceLayer::new_for_http()) .with_state(state) diff --git a/backend/crates/api/src/main.rs b/backend/crates/api/src/main.rs index f706ed9..29e5c7c 100644 --- a/backend/crates/api/src/main.rs +++ b/backend/crates/api/src/main.rs @@ -16,6 +16,7 @@ async fn main() -> anyhow::Result<()> { { std::fs::create_dir_all(dir)?; } + std::fs::create_dir_all(&cfg.work_dir)?; let pool = infrastructure::connect(&cfg.database_url).await?; let state = AppState::new(cfg.clone(), pool)?; state.bootstrap().await?; diff --git a/backend/crates/api/src/openapi.rs b/backend/crates/api/src/openapi.rs index 7981180..2c30a46 100644 --- a/backend/crates/api/src/openapi.rs +++ b/backend/crates/api/src/openapi.rs @@ -30,6 +30,10 @@ impl Modify for BearerAuth { crate::cluster::overview, crate::cluster::restart, crate::cluster::scale, crate::cluster::set_image, crate::vulnerabilities::list, crate::vulnerabilities::summary, crate::vulnerabilities::targets, crate::vulnerabilities::set_status, crate::settings::get_notifications, crate::settings::put_notifications, + crate::backups::list_targets, crate::backups::get_target, crate::backups::create_target, crate::backups::update_target, + 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, ), modifiers(&BearerAuth) )] diff --git a/backend/crates/api/src/test_support.rs b/backend/crates/api/src/test_support.rs index 9b85195..bf69806 100644 --- a/backend/crates/api/src/test_support.rs +++ b/backend/crates/api/src/test_support.rs @@ -17,6 +17,8 @@ pub fn test_config() -> Config { master_key: "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f".into(), fake_host: true, kubeconfig: None, + kubectl: vec!["kubectl".into()], + work_dir: std::env::temp_dir().join(format!("monitoring-test-{}", uuid::Uuid::new_v4())), bind: "127.0.0.1:0".parse().unwrap(), bootstrap_admin: None, cookie_secure: false, @@ -86,6 +88,13 @@ async fn build_test_app_with(cfg: Config) -> Router { updater: Arc::new(infrastructure::FakeHostUpdater), cluster: Arc::new(infrastructure::FakeClusterGateway), scanner: Arc::new(infrastructure::FakeScanner), + storage: Arc::new(infrastructure::DirBackupStorage::new( + cfg.work_dir.join("remote"), + )), + collector: Arc::new(infrastructure::FakeBackupCollector), + encryptor: Arc::new(infrastructure::OpensslEncryptor::new(Arc::new( + infrastructure::SystemCommandRunner, + ))), }; let state = AppState::with_adapters(cfg, pool, adapters, register_test_jobs).expect("state"); state.bootstrap().await.expect("bootstrap"); diff --git a/backend/crates/api/tests/jobs.rs b/backend/crates/api/tests/jobs.rs index 2f81eb0..3d2151b 100644 --- a/backend/crates/api/tests/jobs.rs +++ b/backend/crates/api/tests/jobs.rs @@ -58,17 +58,25 @@ async fn unknown_kind_is_404_and_run_is_admin_only() { .status, StatusCode::UNPROCESSABLE_ENTITY ); - assert_eq!( - post( - &app, - "/api/jobs/run", - json!({"kind": "backup"}), - Some(&admin) - ) - .await - .status, - StatusCode::NOT_FOUND - ); + // every kind has a handler now; a backup without a strategy id starts and then fails + let run = post( + &app, + "/api/jobs/run", + json!({"kind": "backup"}), + Some(&admin), + ) + .await; + assert_eq!(run.status, StatusCode::ACCEPTED); + let id = run.json["id"].as_str().unwrap().to_string(); + for _ in 0..50 { + let r = get(&app, &format!("/api/jobs/{id}"), Some(&admin)).await; + if r.json["status"] != "running" { + assert_eq!(r.json["status"], "failed"); + assert!(r.json["log"].as_str().unwrap().contains("strategy id")); + break; + } + tokio::time::sleep(std::time::Duration::from_millis(20)).await; + } 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 diff --git a/backend/crates/application/src/backup_service.rs b/backend/crates/application/src/backup_service.rs index 4024fc4..4c4ed2f 100644 --- a/backend/crates/application/src/backup_service.rs +++ b/backend/crates/application/src/backup_service.rs @@ -10,9 +10,11 @@ use domain::ports::{ BackupTargetRepository, Cipher, FileEncryptor, }; use domain::DomainError; +use sha2::{Digest, Sha256}; use uuid::Uuid; use crate::jobs::{JobHandler, JobLog}; +use crate::scheduler::{is_due, validate_cron}; pub struct BackupDeps { pub targets: Arc, @@ -42,65 +44,316 @@ impl BackupService { Self { d: deps } } + fn decrypt_target(&self, mut t: BackupTarget) -> Result { + if !t.password.is_empty() { + t.password = self.d.cipher.decrypt(&t.password)?; + } + Ok(t) + } + + fn encrypt_target(&self, mut t: BackupTarget) -> Result { + if !t.password.is_empty() { + t.password = self.d.cipher.encrypt(&t.password)?; + } + Ok(t) + } + + fn decrypt_strategy(&self, mut s: BackupStrategy) -> Result { + if let Some(p) = &s.passphrase { + s.passphrase = Some(self.d.cipher.decrypt(p)?); + } + Ok(s) + } + + fn encrypt_strategy(&self, mut s: BackupStrategy) -> Result { + if let Some(p) = &s.passphrase { + s.passphrase = Some(self.d.cipher.encrypt(p)?); + } + Ok(s) + } + // ---- targets ---- pub async fn list_targets(&self) -> Result, DomainError> { - todo!() + self.d + .targets + .list() + .await? + .into_iter() + .map(|t| self.decrypt_target(t)) + .collect() } - pub async fn get_target(&self, _id: Uuid) -> Result { - todo!() + + pub async fn get_target(&self, id: Uuid) -> Result { + self.decrypt_target(self.d.targets.get(id).await?.ok_or(DomainError::NotFound)?) } - pub async fn create_target(&self, _t: BackupTarget) -> Result { - todo!() + + pub async fn create_target(&self, mut t: BackupTarget) -> Result { + t.validate()?; + t.id = Uuid::new_v4(); + self.d + .targets + .insert(&self.encrypt_target(t.clone())?) + .await?; + Ok(t) } + /// Empty password keeps the stored one. - pub async fn update_target(&self, _t: BackupTarget) -> Result { - todo!() + pub async fn update_target(&self, mut t: BackupTarget) -> Result { + t.validate()?; + let current = self.get_target(t.id).await?; + if t.password.is_empty() { + t.password = current.password; + } + self.d + .targets + .update(&self.encrypt_target(t.clone())?) + .await?; + Ok(t) } + /// Fails with `Conflict` while a strategy still uses the target. - pub async fn delete_target(&self, _id: Uuid) -> Result<(), DomainError> { - todo!() + pub async fn delete_target(&self, id: Uuid) -> Result<(), DomainError> { + if let Some(s) = self + .d + .strategies + .list() + .await? + .into_iter() + .find(|s| s.target_id == id) + { + return Err(DomainError::Conflict(format!( + "target is used by strategy '{}'", + s.name + ))); + } + self.d.targets.delete(id).await } - pub async fn test_target(&self, _id: Uuid) -> Result<(), DomainError> { - todo!() + + pub async fn test_target(&self, id: Uuid) -> Result<(), DomainError> { + let t = self.get_target(id).await?; + self.d.storage.test(&t).await } // ---- strategies ---- pub async fn list_strategies(&self) -> Result, DomainError> { - todo!() + let targets = self.d.targets.list().await?; + let mut out = Vec::new(); + for s in self.d.strategies.list().await? { + let target_name = targets + .iter() + .find(|t| t.id == s.target_id) + .map(|t| t.name.clone()) + .unwrap_or_default(); + let last_backup = self.d.records.list_for(s.id, 1).await?.into_iter().next(); + let mut strategy = self.decrypt_strategy(s)?; + // never expose the passphrase in listings + strategy.passphrase = strategy.passphrase.map(|_| String::new()); + out.push(StrategyStatus { + strategy, + target_name, + last_backup, + }); + } + Ok(out) } - pub async fn get_strategy(&self, _id: Uuid) -> Result { - todo!() + + pub async fn get_strategy(&self, id: Uuid) -> Result { + self.decrypt_strategy( + self.d + .strategies + .get(id) + .await? + .ok_or(DomainError::NotFound)?, + ) } - pub async fn create_strategy(&self, _s: BackupStrategy) -> Result { - todo!() + + async fn validate_strategy(&self, s: &BackupStrategy) -> Result<(), DomainError> { + s.validate()?; + validate_cron(&s.schedule)?; + self.d + .targets + .get(s.target_id) + .await? + .ok_or(DomainError::NotFound)?; + Ok(()) } + + pub async fn create_strategy( + &self, + mut s: BackupStrategy, + ) -> Result { + s.passphrase = s.passphrase.filter(|p| !p.is_empty()); + self.validate_strategy(&s).await?; + s.id = Uuid::new_v4(); + self.d + .strategies + .insert(&self.encrypt_strategy(s.clone())?) + .await?; + Ok(s) + } + /// Empty passphrase keeps the stored one; `None` removes it. - pub async fn update_strategy(&self, _s: BackupStrategy) -> Result { - todo!() + pub async fn update_strategy( + &self, + mut s: BackupStrategy, + ) -> Result { + let current = self.get_strategy(s.id).await?; + if s.passphrase.as_deref() == Some("") { + s.passphrase = current.passphrase; + } + self.validate_strategy(&s).await?; + self.d + .strategies + .update(&self.encrypt_strategy(s.clone())?) + .await?; + Ok(s) } - pub async fn delete_strategy(&self, _id: Uuid) -> Result<(), DomainError> { - todo!() + + pub async fn delete_strategy(&self, id: Uuid) -> Result<(), DomainError> { + self.d.strategies.delete(id).await } - pub async fn records(&self, _strategy_id: Uuid) -> Result, DomainError> { - todo!() + + pub async fn records(&self, strategy_id: Uuid) -> Result, DomainError> { + self.d.records.list_for(strategy_id, 100).await } /// Enabled strategies whose cron fired since their last backup. pub async fn due_strategies( &self, - _now: DateTime, - _grace_secs: i64, + now: DateTime, + grace_secs: i64, ) -> Result, DomainError> { - todo!() + let mut due = Vec::new(); + for s in self + .d + .strategies + .list() + .await? + .into_iter() + .filter(|s| s.enabled) + { + let last = self + .d + .records + .list_for(s.id, 1) + .await? + .first() + .map(|r| r.created_at); + if is_due(&s.schedule, last, now, grace_secs) { + due.push(s.id); + } + } + Ok(due) } /// Collect, encrypt, upload, verify, record, apply retention. pub async fn run_strategy( &self, - _id: Uuid, - _log: &dyn JobLog, + id: Uuid, + log: &dyn JobLog, ) -> Result { - todo!() + let strategy = self.get_strategy(id).await?; + let target = self.get_target(strategy.target_id).await?; + let work = self.d.work_dir.join(format!("run-{}", Uuid::new_v4())); + std::fs::create_dir_all(&work) + .map_err(|e| DomainError::Storage(format!("work dir: {e}")))?; + let result = self.run_in(&strategy, &target, &work, log).await; + let _ = std::fs::remove_dir_all(&work); + result + } + + async fn run_in( + &self, + strategy: &BackupStrategy, + target: &BackupTarget, + work: &std::path::Path, + log: &dyn JobLog, + ) -> Result { + let started = Utc::now(); + log.line(&format!( + "strategy '{}' -> target '{}' ({})", + strategy.name, + target.name, + target.kind.as_str() + )) + .await; + + let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); + let sink = crate::upgrade_service::channel_sink(tx); + let collect = async { + let r = self + .d + .collector + .collect(&strategy.source, work, &sink) + .await; + drop(sink); + r + }; + let drain = async { + while let Some(l) = rx.recv().await { + log.line(&l).await; + } + }; + let (archive, _) = tokio::join!(collect, drain); + let mut file = archive?; + + if let Some(pass) = &strategy.passphrase { + log.line("encrypting archive").await; + file = self.d.encryptor.encrypt(&file, pass).await?; + } + let bytes = + std::fs::read(&file).map_err(|e| DomainError::Storage(format!("read archive: {e}")))?; + let size_bytes = bytes.len() as u64; + let sha256 = hex::encode(Sha256::digest(&bytes)); + drop(bytes); + let filename = strategy.filename(started); + + log.line(&format!("uploading {filename} ({size_bytes} bytes)")) + .await; + self.d.storage.upload(target, &file, &filename).await?; + let remote = self.d.storage.list(target).await?; + match remote.iter().find(|f| f.name == filename) { + Some(f) if f.size_bytes == size_bytes => { + log.line(&format!("uploaded and verified {filename}")).await + } + Some(f) => { + return Err(DomainError::Unavailable(format!( + "size mismatch after upload: {} != {size_bytes}", + f.size_bytes + ))) + } + None => { + return Err(DomainError::Unavailable( + "file not found on target after upload".into(), + )) + } + } + + let record = BackupRecord { + id: Uuid::new_v4(), + strategy_id: strategy.id, + filename: filename.clone(), + size_bytes, + sha256, + created_at: started, + }; + self.d.records.insert(&record).await?; + + // retention: keep the newest N files of this strategy on the target + let prefix = format!("{}_", strategy.slug()); + let mut mine: Vec = remote + .iter() + .map(|f| f.name.clone()) + .filter(|n| n.starts_with(&prefix)) + .collect(); + mine.sort(); + mine.reverse(); + for old in mine.iter().skip(strategy.retention as usize) { + log.line(&format!("retention: deleting {old}")).await; + self.d.storage.delete(target, old).await?; + self.d.records.delete_by_filename(strategy.id, old).await?; + } + Ok(record) } } @@ -109,7 +362,22 @@ pub struct BackupJob(pub Arc); #[async_trait] impl JobHandler for BackupJob { - async fn run(&self, _params: Option, _log: &dyn JobLog) -> Result<(), String> { - todo!() + async fn run(&self, params: Option, log: &dyn JobLog) -> Result<(), String> { + let id: Uuid = params + .as_deref() + .filter(|p| !p.is_empty()) + .ok_or("backup job needs a strategy id as params")? + .parse() + .map_err(|e| format!("invalid strategy id: {e}"))?; + let record = self.0.run_strategy(id, log).await.map_err(|e| match e { + DomainError::NotFound => "strategy not found".to_string(), + e => e.to_string(), + })?; + log.line(&format!( + "backup complete: {} ({} bytes, sha256 {})", + record.filename, record.size_bytes, record.sha256 + )) + .await; + Ok(()) } } diff --git a/backend/crates/application/src/scheduler.rs b/backend/crates/application/src/scheduler.rs index 33436eb..0abf948 100644 --- a/backend/crates/application/src/scheduler.rs +++ b/backend/crates/application/src/scheduler.rs @@ -8,7 +8,7 @@ use cron::Schedule; use domain::jobs::JobKind; use domain::DomainError; -use crate::{JobRunner, SettingsService}; +use crate::{BackupService, JobRunner, SettingsService}; /// Validate a 6-field cron expression (seconds first). pub fn validate_cron(expr: &str) -> Result<(), DomainError> { @@ -38,6 +38,7 @@ pub fn is_due( pub struct Scheduler { runner: Arc, settings: Arc, + backups: Option>, tick: Duration, } @@ -46,10 +47,16 @@ impl Scheduler { Self { runner, settings, + backups: None, tick: Duration::from_secs(30), } } + pub fn with_backups(mut self, backups: Arc) -> Self { + self.backups = Some(backups); + self + } + /// One pass: start every registered kind that is due. Returns the kinds started. pub async fn tick_once(&self, now: DateTime) -> Vec { let mut started = Vec::new(); @@ -70,6 +77,27 @@ impl Scheduler { started.push(kind); } } + if let Some(backups) = &self.backups { + match backups + .due_strategies(now, self.tick.as_secs() as i64 * 2) + .await + { + Ok(ids) => { + for id in ids { + // one backup at a time; the rest is picked up on a later tick + if self + .runner + .start(JobKind::Backup, Some(id.to_string()), "scheduler") + .await + .is_ok() + { + started.push(JobKind::Backup); + } + } + } + Err(e) => eprintln!("scheduler: cannot read backup strategies: {e}"), + } + } started } diff --git a/backend/crates/application/src/test_fakes.rs b/backend/crates/application/src/test_fakes.rs index 0c185f0..f069f69 100644 --- a/backend/crates/application/src/test_fakes.rs +++ b/backend/crates/application/src/test_fakes.rs @@ -748,7 +748,7 @@ impl BackupRecordRepository for MemRecords { .filter(|r| r.strategy_id == strategy_id) .cloned() .collect(); - v.sort_by(|a, b| b.created_at.cmp(&a.created_at)); + v.sort_by_key(|r| std::cmp::Reverse(r.created_at)); v.truncate(limit as usize); Ok(v) } diff --git a/backend/crates/application/src/upgrade_service.rs b/backend/crates/application/src/upgrade_service.rs index 395b49f..411186d 100644 --- a/backend/crates/application/src/upgrade_service.rs +++ b/backend/crates/application/src/upgrade_service.rs @@ -46,7 +46,11 @@ pub struct PackageUpgradeJob { } /// Bridges the synchronous `LineSink` of the updater to the async job log. -struct ChannelSink(tokio::sync::mpsc::UnboundedSender); +pub struct ChannelSink(tokio::sync::mpsc::UnboundedSender); + +pub fn channel_sink(tx: tokio::sync::mpsc::UnboundedSender) -> ChannelSink { + ChannelSink(tx) +} impl LineSink for ChannelSink { fn line(&self, text: &str) { diff --git a/backend/crates/infrastructure/Cargo.toml b/backend/crates/infrastructure/Cargo.toml index 434b503..d1dccc2 100644 --- a/backend/crates/infrastructure/Cargo.toml +++ b/backend/crates/infrastructure/Cargo.toml @@ -23,6 +23,7 @@ kube = { version = "0.99", default-features = false, features = ["client", "rust k8s-openapi = { version = "0.24", features = ["v1_32"] } [dev-dependencies] +tempfile = "3" tokio.workspace = true [lints] diff --git a/backend/crates/infrastructure/migrations/0005_backups.sql b/backend/crates/infrastructure/migrations/0005_backups.sql new file mode 100644 index 0000000..3cdaf7c --- /dev/null +++ b/backend/crates/infrastructure/migrations/0005_backups.sql @@ -0,0 +1,33 @@ +CREATE TABLE backup_targets ( + id TEXT PRIMARY KEY, + name TEXT NOT NULL, + kind TEXT NOT NULL CHECK (kind IN ('smb', 'ftp')), + host TEXT NOT NULL, + port INTEGER, + share TEXT NOT NULL DEFAULT '', + path TEXT NOT NULL DEFAULT '', + username TEXT NOT NULL DEFAULT '', + password TEXT NOT NULL DEFAULT '', + tls INTEGER NOT NULL DEFAULT 0 +); + +CREATE TABLE backup_strategies ( + id TEXT PRIMARY KEY, + name TEXT NOT NULL, + source TEXT NOT NULL, + schedule TEXT NOT NULL, + target_id TEXT NOT NULL REFERENCES backup_targets(id), + retention INTEGER NOT NULL, + passphrase TEXT, + enabled INTEGER NOT NULL DEFAULT 1 +); + +CREATE TABLE backup_records ( + id TEXT PRIMARY KEY, + strategy_id TEXT NOT NULL REFERENCES backup_strategies(id) ON DELETE CASCADE, + filename TEXT NOT NULL, + size_bytes INTEGER NOT NULL, + sha256 TEXT NOT NULL, + created_at TEXT NOT NULL +); +CREATE INDEX backup_records_strategy ON backup_records(strategy_id, created_at DESC); diff --git a/backend/crates/infrastructure/src/backup.rs b/backend/crates/infrastructure/src/backup.rs new file mode 100644 index 0000000..37ed518 --- /dev/null +++ b/backend/crates/infrastructure/src/backup.rs @@ -0,0 +1,735 @@ +//! Backup adapters: remote storage via smbclient/curl, collectors via kubectl/tar, openssl encryption, +//! and file-based fakes for development. +use std::path::{Path, PathBuf}; +use std::sync::Arc; + +use async_trait::async_trait; +use domain::backup::{BackupSource, BackupTarget, RemoteFile, StorageKind}; +use domain::ports::{BackupCollector, BackupStorage, FileEncryptor, LineSink}; +use domain::DomainError; + +use crate::host::CommandRunner; + +fn unavailable(what: &str, out: &crate::host::Output) -> DomainError { + let tail: Vec<&str> = out + .stderr + .lines() + .chain(out.stdout.lines()) + .rev() + .take(3) + .collect(); + DomainError::Unavailable(format!( + "{what}: {}", + tail.into_iter().rev().collect::>().join(" | ") + )) +} + +fn join_path(dir: &str, name: &str) -> String { + let dir = dir.trim_matches('/'); + if dir.is_empty() { + name.to_string() + } else { + format!("{dir}/{name}") + } +} + +// ---------------------------------------------------------------- storage + +pub struct CommandBackupStorage { + runner: Arc, +} + +impl CommandBackupStorage { + pub fn new(runner: Arc) -> Self { + Self { runner } + } + + async fn smb( + &self, + t: &BackupTarget, + command: &str, + ) -> Result { + let service = format!("//{}/{}", t.host, t.share); + let mut args = vec![ + service.as_str(), + "-U", + if t.username.is_empty() { + "guest" + } else { + t.username.as_str() + }, + ]; + let port; + if let Some(p) = t.port { + port = p.to_string(); + args.extend(["-p", port.as_str()]); + } + if t.username.is_empty() { + args.push("-N"); + } + let dir; + if !t.path.trim_matches('/').is_empty() { + dir = t.path.trim_matches('/').to_string(); + args.extend(["-D", dir.as_str()]); + } + args.extend(["-c", command]); + self.runner + .run_env("smbclient", &args, &[("PASSWD", t.password.as_str())]) + .await + } + + fn ftp_url(t: &BackupTarget, name: &str) -> String { + let port = t.port.map(|p| format!(":{p}")).unwrap_or_default(); + let path = join_path(&t.path, name); + format!("ftp://{}{port}/{path}", t.host) + } + + async fn curl( + &self, + t: &BackupTarget, + extra: &[&str], + ) -> Result { + // credentials via a netrc file so they never appear in the process list + let netrc = std::env::temp_dir().join(format!("monitoring-netrc-{}", uuid::Uuid::new_v4())); + std::fs::write( + &netrc, + format!( + "machine {} login {} password {}\n", + t.host, t.username, t.password + ), + ) + .map_err(|e| DomainError::Storage(e.to_string()))?; + let _ = + std::fs::set_permissions(&netrc, std::os::unix::fs::PermissionsExt::from_mode(0o600)); + let netrc_s = netrc.to_string_lossy().to_string(); + let mut args = vec!["-sS", "--fail", "--netrc-file", netrc_s.as_str()]; + if t.tls { + args.push("--ssl-reqd"); + } + args.extend(extra); + let res = self.runner.run("curl", &args).await; + let _ = std::fs::remove_file(&netrc); + res + } +} + +/// `smbclient -c ls` lines: ` name A 1234 Tue Sep 2 ...` +pub fn parse_smb_ls(out: &str) -> Vec { + out.lines() + .filter_map(|l| { + let l = l.trim_end(); + if !l.starts_with(" ") || l.contains("blocks of size") { + return None; + } + // attributes column is a short uppercase token ("A", "D", "AH"...), size follows it + let parts: Vec<&str> = l.split_whitespace().collect(); + let idx = parts + .iter() + .position(|p| p.len() <= 3 && p.chars().all(|c| c.is_ascii_uppercase()))?; + if idx == 0 { + return None; + } + let name = parts[..idx].join(" "); + if name == "." || name == ".." || parts[idx].contains('D') { + return None; + } + let size = parts.get(idx + 1)?.parse().ok()?; + Some(RemoteFile { + name, + size_bytes: size, + }) + }) + .collect() +} + +/// Unix-style FTP `LIST` lines: `-rw-r--r-- 1 u g 1234 Sep 2 10:00 name` +pub fn parse_ftp_list(out: &str) -> Vec { + out.lines() + .filter_map(|l| { + let parts: Vec<&str> = l.split_whitespace().collect(); + if parts.len() < 9 || !parts[0].starts_with('-') { + return None; + } + Some(RemoteFile { + name: parts[8..].join(" "), + size_bytes: parts[4].parse().ok()?, + }) + }) + .collect() +} + +#[async_trait] +impl BackupStorage for CommandBackupStorage { + async fn test(&self, t: &BackupTarget) -> Result<(), DomainError> { + self.list(t).await.map(|_| ()) + } + + async fn upload( + &self, + t: &BackupTarget, + local: &Path, + remote_name: &str, + ) -> Result<(), DomainError> { + let local_s = local.to_string_lossy().to_string(); + let out = match t.kind { + StorageKind::Smb => { + self.smb(t, &format!("put \"{local_s}\" \"{remote_name}\"")) + .await? + } + StorageKind::Ftp => { + self.curl( + t, + &[ + "--ftp-create-dirs", + "-T", + &local_s, + &Self::ftp_url(t, remote_name), + ], + ) + .await? + } + }; + out.success + .then_some(()) + .ok_or_else(|| unavailable("upload failed", &out)) + } + + async fn list(&self, t: &BackupTarget) -> Result, DomainError> { + let out = match t.kind { + StorageKind::Smb => self.smb(t, "ls").await?, + StorageKind::Ftp => self.curl(t, &[&Self::ftp_url(t, "")]).await?, + }; + if !out.success { + return Err(unavailable("listing failed", &out)); + } + Ok(match t.kind { + StorageKind::Smb => parse_smb_ls(&out.stdout), + StorageKind::Ftp => parse_ftp_list(&out.stdout), + }) + } + + async fn delete(&self, t: &BackupTarget, remote_name: &str) -> Result<(), DomainError> { + let out = match t.kind { + StorageKind::Smb => self.smb(t, &format!("del \"{remote_name}\"")).await?, + StorageKind::Ftp => { + let dele = format!("DELE {}", join_path(&t.path, remote_name)); + self.curl( + t, + &[ + "-Q", + &dele, + &Self::ftp_url( + &BackupTarget { + path: String::new(), + ..t.clone() + }, + "", + ), + ], + ) + .await? + } + }; + out.success + .then_some(()) + .ok_or_else(|| unavailable("delete failed", &out)) + } +} + +// ---------------------------------------------------------------- collector + +pub struct KubeBackupCollector { + runner: Arc, + /// kubectl invocation, e.g. `["/snap/bin/microk8s", "kubectl"]`. + kubectl: Vec, +} + +impl KubeBackupCollector { + pub fn new(runner: Arc, kubectl: Vec) -> Self { + Self { runner, kubectl } + } + + async fn kubectl(&self, args: &[&str]) -> Result { + let mut all: Vec<&str> = self.kubectl.iter().skip(1).map(String::as_str).collect(); + all.extend(args); + self.runner.run(&self.kubectl[0], &all).await + } + + async fn kubectl_to_file( + &self, + args: &[&str], + file: &Path, + ) -> Result { + let mut all: Vec<&str> = self.kubectl.iter().skip(1).map(String::as_str).collect(); + all.extend(args); + self.runner.run_to_file(&self.kubectl[0], &all, file).await + } + + async fn tar_dir(&self, dir: &str, file: &Path, out: &dyn LineSink) -> Result<(), DomainError> { + let f = file.to_string_lossy().to_string(); + out.line(&format!("$ tar -czf {f} -C {dir} .")); + let res = self + .runner + .run("tar", &["-czf", &f, "-C", dir, "."]) + .await?; + res.success + .then_some(()) + .ok_or_else(|| unavailable("tar failed", &res)) + } + + async fn gzip(&self, file: &Path) -> Result { + let f = file.to_string_lossy().to_string(); + let res = self.runner.run("gzip", &["-f", &f]).await?; + res.success + .then_some(file.with_extension(format!( + "{}.gz", + file.extension().and_then(|e| e.to_str()).unwrap_or("") + ))) + .ok_or_else(|| unavailable("gzip failed", &res)) + } +} + +#[async_trait] +impl BackupCollector for KubeBackupCollector { + async fn collect( + &self, + source: &BackupSource, + work_dir: &Path, + out: &dyn LineSink, + ) -> Result { + match source { + BackupSource::VolumeClaim { namespace, pvc } => { + out.line(&format!("resolving host path of pvc {namespace}/{pvc}")); + let pv = self + .kubectl(&[ + "get", + "pvc", + "-n", + namespace, + pvc, + "-o", + "jsonpath={.spec.volumeName}", + ]) + .await?; + if !pv.success || pv.stdout.trim().is_empty() { + return Err(unavailable("pvc lookup failed", &pv)); + } + let path = self + .kubectl(&[ + "get", + "pv", + pv.stdout.trim(), + "-o", + "jsonpath={.spec.hostPath.path}", + ]) + .await?; + let dir = path.stdout.trim().to_string(); + if !path.success || dir.is_empty() { + return Err(DomainError::Unavailable( + "pv has no hostPath (only hostpath volumes are supported)".into(), + )); + } + let file = work_dir.join("volume.tar.gz"); + self.tar_dir(&dir, &file, out).await?; + Ok(file) + } + BackupSource::PostgresDump { namespace, pod } => { + out.line(&format!( + "$ kubectl exec -n {namespace} {pod} -- pg_dumpall" + )); + let raw = work_dir.join("dump.sql"); + let res = self + .kubectl_to_file(&["exec", "-n", namespace, pod, "--", "sh", "-c", "PGPASSWORD=\"${POSTGRES_PASSWORD:-$POSTGRESQL_PASSWORD}\" pg_dumpall -U postgres"], &raw) + .await?; + if !res.success { + return Err(unavailable("pg_dumpall failed", &res)); + } + self.gzip(&raw).await + } + BackupSource::KubernetesManifests { namespace } => { + out.line(&format!("$ kubectl get all,configmap,secret,ingress,pvc,serviceaccount -n {namespace} -o yaml")); + let raw = work_dir.join("manifests.yaml"); + let res = self + .kubectl_to_file( + &[ + "get", + "all,configmap,secret,ingress,pvc,serviceaccount", + "-n", + namespace, + "-o", + "yaml", + ], + &raw, + ) + .await?; + if !res.success { + return Err(unavailable("kubectl get failed", &res)); + } + self.gzip(&raw).await + } + BackupSource::HostPath { path } => { + let file = work_dir.join("hostpath.tar.gz"); + self.tar_dir(path, &file, out).await?; + Ok(file) + } + } + } +} + +// ---------------------------------------------------------------- encryption + +/// `openssl enc -aes-256-cbc -pbkdf2`; decrypt with +/// `openssl enc -d -aes-256-cbc -pbkdf2 -in FILE.enc -out FILE`. +pub struct OpensslEncryptor { + runner: Arc, +} + +impl OpensslEncryptor { + pub fn new(runner: Arc) -> Self { + Self { runner } + } +} + +#[async_trait] +impl FileEncryptor for OpensslEncryptor { + async fn encrypt(&self, input: &Path, passphrase: &str) -> Result { + let out = PathBuf::from(format!("{}.enc", input.to_string_lossy())); + let (i, o) = ( + input.to_string_lossy().to_string(), + out.to_string_lossy().to_string(), + ); + let res = self + .runner + .run_env( + "openssl", + &[ + "enc", + "-aes-256-cbc", + "-pbkdf2", + "-salt", + "-in", + &i, + "-out", + &o, + "-pass", + "env:BACKUP_PASSPHRASE", + ], + &[("BACKUP_PASSPHRASE", passphrase)], + ) + .await?; + if !res.success { + return Err(unavailable("openssl enc failed", &res)); + } + let _ = std::fs::remove_file(input); + Ok(out) + } +} + +// ---------------------------------------------------------------- fakes + +/// Stores "remote" files in a local directory per target (FAKE_HOST=true). +pub struct DirBackupStorage { + root: PathBuf, +} + +impl DirBackupStorage { + pub fn new(root: PathBuf) -> Self { + Self { root } + } + fn dir(&self, t: &BackupTarget) -> Result { + let d = self.root.join(t.id.to_string()); + std::fs::create_dir_all(&d).map_err(|e| DomainError::Storage(e.to_string()))?; + Ok(d) + } +} + +#[async_trait] +impl BackupStorage for DirBackupStorage { + async fn test(&self, t: &BackupTarget) -> Result<(), DomainError> { + if t.host.contains("unreachable") { + return Err(DomainError::Unavailable("connection refused".into())); + } + self.dir(t).map(|_| ()) + } + async fn upload( + &self, + t: &BackupTarget, + local: &Path, + remote_name: &str, + ) -> Result<(), DomainError> { + std::fs::copy(local, self.dir(t)?.join(remote_name)) + .map(|_| ()) + .map_err(|e| DomainError::Storage(e.to_string())) + } + async fn list(&self, t: &BackupTarget) -> Result, DomainError> { + let mut v = Vec::new(); + for e in std::fs::read_dir(self.dir(t)?) + .map_err(|e| DomainError::Storage(e.to_string()))? + .flatten() + { + if let Ok(m) = e.metadata() { + v.push(RemoteFile { + name: e.file_name().to_string_lossy().to_string(), + size_bytes: m.len(), + }); + } + } + Ok(v) + } + async fn delete(&self, t: &BackupTarget, remote_name: &str) -> Result<(), DomainError> { + std::fs::remove_file(self.dir(t)?.join(remote_name)) + .map_err(|e| DomainError::Storage(e.to_string())) + } +} + +/// Produces a small archive describing the source (FAKE_HOST=true). +pub struct FakeBackupCollector; + +#[async_trait] +impl BackupCollector for FakeBackupCollector { + async fn collect( + &self, + source: &BackupSource, + work_dir: &Path, + out: &dyn LineSink, + ) -> Result { + out.line(&format!( + "fake: collecting {}", + serde_json::to_string(source).unwrap_or_default() + )); + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + let file = work_dir.join(format!("archive.{}", source.extension())); + std::fs::write(&file, format!("fake backup of {source:?}\n")) + .map_err(|e| DomainError::Storage(e.to_string()))?; + Ok(file) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::Mutex; + + #[test] + fn parses_smbclient_listing() { + let out = " . D 0 Tue Sep 2 10:00:00 2026\n .. D 0 Tue Sep 2 10:00:00 2026\n gitea-db_20260902-020000.sql.gz A 12345 Tue Sep 2 02:00:05 2026\n my file.tar.gz A 99 Tue Sep 2 02:00:05 2026\n\n\t\t1234567 blocks of size 1024. 999 blocks available\n"; + let files = parse_smb_ls(out); + assert_eq!( + files, + vec![ + RemoteFile { + name: "gitea-db_20260902-020000.sql.gz".into(), + size_bytes: 12345 + }, + RemoteFile { + name: "my file.tar.gz".into(), + size_bytes: 99 + } + ] + ); + } + + #[test] + fn parses_ftp_listing() { + let out = "drwxr-xr-x 2 ftp ftp 4096 Sep 02 02:00 sub\n-rw-r--r-- 1 ftp ftp 12345 Sep 02 02:00 gitea-db_20260902-020000.sql.gz\n"; + assert_eq!( + parse_ftp_list(out), + vec![RemoteFile { + name: "gitea-db_20260902-020000.sql.gz".into(), + size_bytes: 12345 + }] + ); + } + + type Call = (String, Vec, Vec<(String, String)>); + #[derive(Default)] + struct Rec(Mutex>); + #[async_trait] + impl CommandRunner for Rec { + async fn run(&self, p: &str, a: &[&str]) -> Result { + self.run_env(p, a, &[]).await + } + async fn run_env( + &self, + p: &str, + a: &[&str], + env: &[(&str, &str)], + ) -> Result { + self.0.lock().unwrap().push(( + p.into(), + a.iter().map(|s| s.to_string()).collect(), + env.iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + )); + Ok(crate::host::Output { + stdout: "pvc-123\n".into(), + success: true, + ..Default::default() + }) + } + async fn run_to_file( + &self, + p: &str, + a: &[&str], + f: &Path, + ) -> Result { + std::fs::write(f, "data").unwrap(); + self.run_env(p, a, &[]).await + } + async fn read_file(&self, _: &str) -> Result, DomainError> { + Ok(None) + } + async fn run_streaming( + &self, + _: &str, + _: &[&str], + _: &dyn LineSink, + ) -> Result { + Ok(true) + } + } + + struct Sink; + impl LineSink for Sink { + fn line(&self, _: &str) {} + } + + fn smb_target() -> BackupTarget { + BackupTarget { + id: uuid::Uuid::new_v4(), + name: "n".into(), + kind: StorageKind::Smb, + host: "nas".into(), + port: None, + share: "backups".into(), + path: "softvisor".into(), + username: "u".into(), + password: "p".into(), + tls: false, + } + } + + #[tokio::test] + async fn smb_upload_passes_password_via_env_not_argv() { + let r = Arc::new(Rec::default()); + CommandBackupStorage::new(r.clone()) + .upload(&smb_target(), Path::new("/tmp/x.tar.gz"), "x.tar.gz") + .await + .unwrap(); + let calls = r.0.lock().unwrap(); + let (prog, args, env) = &calls[0]; + assert_eq!(prog, "smbclient"); + assert_eq!(args[0], "//nas/backups"); + assert!(args.contains(&"-D".to_string()) && args.contains(&"softvisor".to_string())); + assert!(args.iter().any(|a| a.starts_with("put ")), "{args:?}"); + assert!(!args.iter().any(|a| a.contains('p') && a.len() == 1)); + assert_eq!(env[0], ("PASSWD".to_string(), "p".to_string())); + } + + #[tokio::test] + async fn ftp_uses_curl_with_netrc_and_tls_flag() { + let r = Arc::new(Rec::default()); + let mut t = smb_target(); + t.kind = StorageKind::Ftp; + t.tls = true; + t.port = Some(2121); + CommandBackupStorage::new(r.clone()) + .upload(&t, Path::new("/tmp/x"), "x") + .await + .unwrap(); + let calls = r.0.lock().unwrap(); + let (prog, args, _) = &calls[0]; + assert_eq!(prog, "curl"); + assert!(args.contains(&"--netrc-file".to_string())); + assert!(args.contains(&"--ssl-reqd".to_string())); + assert!( + args.last().unwrap().ends_with("ftp://nas:2121/softvisor/x"), + "{args:?}" + ); + assert!( + !args.iter().any(|a| a.contains("u:p")), + "no credentials in argv" + ); + } + + #[tokio::test] + async fn collector_builds_kubectl_and_tar_commands() { + let r = Arc::new(Rec::default()); + let c = KubeBackupCollector::new( + r.clone(), + vec!["/snap/bin/microk8s".into(), "kubectl".into()], + ); + let work = tempfile::tempdir().unwrap(); + let f = c + .collect( + &BackupSource::VolumeClaim { + namespace: "gitea".into(), + pvc: "data".into(), + }, + work.path(), + &Sink, + ) + .await + .unwrap(); + assert!(f.ends_with("volume.tar.gz")); + let f = c + .collect( + &BackupSource::PostgresDump { + namespace: "gitea".into(), + pod: "pg-0".into(), + }, + work.path(), + &Sink, + ) + .await + .unwrap(); + assert!(f.ends_with("dump.sql.gz"), "{f:?}"); + let calls = r.0.lock().unwrap(); + assert_eq!(calls[0].0, "/snap/bin/microk8s"); + assert_eq!(calls[0].1[..3], ["kubectl", "get", "pvc"]); + assert_eq!(calls[2].0, "tar"); + let exec = calls + .iter() + .find(|c| c.1.contains(&"exec".to_string())) + .unwrap(); + assert!(exec.1.iter().any(|a| a.contains("pg_dumpall"))); + assert!(calls.iter().any(|c| c.0 == "gzip")); + } + + #[tokio::test] + async fn openssl_receives_passphrase_via_env() { + let r = Arc::new(Rec::default()); + let work = tempfile::tempdir().unwrap(); + let input = work.path().join("a.tar.gz"); + std::fs::write(&input, "x").unwrap(); + let out = OpensslEncryptor::new(r.clone()) + .encrypt(&input, "pw") + .await + .unwrap(); + assert!(out.to_string_lossy().ends_with("a.tar.gz.enc")); + let calls = r.0.lock().unwrap(); + assert_eq!(calls[0].0, "openssl"); + assert!(calls[0].1.contains(&"env:BACKUP_PASSPHRASE".to_string())); + assert_eq!(calls[0].2[0].1, "pw"); + } + + #[tokio::test] + async fn dir_storage_roundtrip() { + let root = tempfile::tempdir().unwrap(); + let s = DirBackupStorage::new(root.path().to_path_buf()); + let t = smb_target(); + let local = root.path().join("local.bin"); + std::fs::write(&local, "hello").unwrap(); + s.upload(&t, &local, "remote.bin").await.unwrap(); + assert_eq!( + s.list(&t).await.unwrap(), + vec![RemoteFile { + name: "remote.bin".into(), + size_bytes: 5 + }] + ); + s.delete(&t, "remote.bin").await.unwrap(); + assert!(s.list(&t).await.unwrap().is_empty()); + } +} diff --git a/backend/crates/infrastructure/src/host/command.rs b/backend/crates/infrastructure/src/host/command.rs index f92ec35..ddc988c 100644 --- a/backend/crates/infrastructure/src/host/command.rs +++ b/backend/crates/infrastructure/src/host/command.rs @@ -13,6 +13,20 @@ pub struct Output { pub trait CommandRunner: Send + Sync { async fn run(&self, program: &str, args: &[&str]) -> Result; async fn read_file(&self, path: &str) -> Result, DomainError>; + /// Run with extra environment variables (used to pass secrets without exposing them in argv). + async fn run_env( + &self, + program: &str, + args: &[&str], + env: &[(&str, &str)], + ) -> Result; + /// Run a command and write its stdout to `file`; stderr is captured in the result. + async fn run_to_file( + &self, + program: &str, + args: &[&str], + file: &std::path::Path, + ) -> Result; /// Run a command and forward each output line (stdout and stderr) to `out`. async fn run_streaming( &self, @@ -27,10 +41,24 @@ pub struct SystemCommandRunner; #[async_trait] impl CommandRunner for SystemCommandRunner { async fn run(&self, program: &str, args: &[&str]) -> Result { - let out = tokio::process::Command::new(program) - .args(args) + self.run_env(program, args, &[]).await + } + + async fn run_env( + &self, + program: &str, + args: &[&str], + env: &[(&str, &str)], + ) -> Result { + let mut cmd = tokio::process::Command::new(program); + cmd.args(args) .env("DEBIAN_FRONTEND", "noninteractive") .env("LC_ALL", "C") + .stdin(std::process::Stdio::null()); + for (k, v) in env { + cmd.env(k, v); + } + let out = cmd .output() .await .map_err(|e| DomainError::Unavailable(format!("{program}: {e}")))?; @@ -41,6 +69,30 @@ impl CommandRunner for SystemCommandRunner { }) } + async fn run_to_file( + &self, + program: &str, + args: &[&str], + file: &std::path::Path, + ) -> Result { + let f = std::fs::File::create(file) + .map_err(|e| DomainError::Storage(format!("{}: {e}", file.display())))?; + let out = tokio::process::Command::new(program) + .args(args) + .env("LC_ALL", "C") + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::from(f)) + .stderr(std::process::Stdio::piped()) + .output() + .await + .map_err(|e| DomainError::Unavailable(format!("{program}: {e}")))?; + Ok(Output { + stdout: String::new(), + stderr: String::from_utf8_lossy(&out.stderr).into_owned(), + success: out.status.success(), + }) + } + async fn run_streaming( &self, program: &str, diff --git a/backend/crates/infrastructure/src/host/debian.rs b/backend/crates/infrastructure/src/host/debian.rs index 7f295ac..ea5b124 100644 --- a/backend/crates/infrastructure/src/host/debian.rs +++ b/backend/crates/infrastructure/src/host/debian.rs @@ -237,6 +237,22 @@ Conf openssl (3.0.16-1~deb12u1 Debian-Security:12/stable-security [amd64])\n"; ) -> Result { unreachable!() } + 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, path: &str) -> Result, DomainError> { Ok(match path { "/etc/os-release" => Some( diff --git a/backend/crates/infrastructure/src/host/updater.rs b/backend/crates/infrastructure/src/host/updater.rs index f9d359f..af2d1cd 100644 --- a/backend/crates/infrastructure/src/host/updater.rs +++ b/backend/crates/infrastructure/src/host/updater.rs @@ -79,6 +79,22 @@ mod tests { ..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) } diff --git a/backend/crates/infrastructure/src/lib.rs b/backend/crates/infrastructure/src/lib.rs index 476484d..8004ee5 100644 --- a/backend/crates/infrastructure/src/lib.rs +++ b/backend/crates/infrastructure/src/lib.rs @@ -1,4 +1,5 @@ //! Infrastructure layer: SQLite repositories, Argon2 hashing, JWT issuing. +pub mod backup; pub mod cipher; pub mod db; pub mod host; @@ -9,6 +10,10 @@ pub mod sqlite; pub mod token; pub mod trivy; +pub use backup::{ + CommandBackupStorage, DirBackupStorage, FakeBackupCollector, KubeBackupCollector, + OpensslEncryptor, +}; pub use cipher::AesGcmCipher; pub use db::{connect, DbPool}; pub use host::{ @@ -18,6 +23,9 @@ pub use k8s::{FakeClusterGateway, KubeGateway}; pub use mail::LettreMailer; pub use password::Argon2Hasher; pub use sqlite::{SqliteAuditLog, SqliteJobRuns, SqliteRefreshTokens, SqliteSettings, SqliteUsers}; -pub use sqlite::{SqliteFindings, SqliteInventory}; +pub use sqlite::{ + SqliteBackupRecords, SqliteBackupStrategies, SqliteBackupTargets, SqliteFindings, + SqliteInventory, +}; pub use token::JwtIssuer; pub use trivy::{FakeScanner, TrivyScanner}; diff --git a/backend/crates/infrastructure/src/sqlite.rs b/backend/crates/infrastructure/src/sqlite.rs index a867f91..bb096ff 100644 --- a/backend/crates/infrastructure/src/sqlite.rs +++ b/backend/crates/infrastructure/src/sqlite.rs @@ -753,3 +753,307 @@ mod finding_tests { ); } } + +use domain::backup::{BackupRecord, BackupSource, BackupStrategy, BackupTarget, StorageKind}; + +pub struct SqliteBackupTargets(pub DbPool); + +fn target_from_row(r: &SqliteRow) -> BackupTarget { + BackupTarget { + id: r.get("id"), + name: r.get("name"), + kind: StorageKind::parse(r.get::("kind").as_str()).unwrap_or(StorageKind::Smb), + host: r.get("host"), + port: r.get::, _>("port").map(|p| p as u16), + share: r.get("share"), + path: r.get("path"), + username: r.get("username"), + password: r.get("password"), + tls: r.get("tls"), + } +} + +const TARGET_COLS: &str = "id, name, kind, host, port, share, path, username, password, tls"; + +#[async_trait] +impl domain::ports::BackupTargetRepository for SqliteBackupTargets { + async fn list(&self) -> Result, DomainError> { + sqlx::query(&format!( + "SELECT {TARGET_COLS} FROM backup_targets ORDER BY name" + )) + .fetch_all(&self.0) + .await + .map(|rows| rows.iter().map(target_from_row).collect()) + .map_err(storage) + } + async fn get(&self, id: Uuid) -> Result, DomainError> { + sqlx::query(&format!( + "SELECT {TARGET_COLS} FROM backup_targets WHERE id = ?" + )) + .bind(id) + .fetch_optional(&self.0) + .await + .map(|r| r.as_ref().map(target_from_row)) + .map_err(storage) + } + async fn insert(&self, t: &BackupTarget) -> Result<(), DomainError> { + sqlx::query(&format!( + "INSERT INTO backup_targets ({TARGET_COLS}) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)" + )) + .bind(t.id) + .bind(&t.name) + .bind(t.kind.as_str()) + .bind(&t.host) + .bind(t.port.map(|p| p as i64)) + .bind(&t.share) + .bind(&t.path) + .bind(&t.username) + .bind(&t.password) + .bind(t.tls) + .execute(&self.0) + .await + .map(|_| ()) + .map_err(storage) + } + async fn update(&self, t: &BackupTarget) -> Result<(), DomainError> { + let res = sqlx::query("UPDATE backup_targets SET name = ?, kind = ?, host = ?, port = ?, share = ?, path = ?, username = ?, password = ?, tls = ? WHERE id = ?") + .bind(&t.name) + .bind(t.kind.as_str()) + .bind(&t.host) + .bind(t.port.map(|p| p as i64)) + .bind(&t.share) + .bind(&t.path) + .bind(&t.username) + .bind(&t.password) + .bind(t.tls) + .bind(t.id) + .execute(&self.0) + .await + .map_err(storage)?; + (res.rows_affected() > 0) + .then_some(()) + .ok_or(DomainError::NotFound) + } + async fn delete(&self, id: Uuid) -> Result<(), DomainError> { + let res = sqlx::query("DELETE FROM backup_targets WHERE id = ?") + .bind(id) + .execute(&self.0) + .await + .map_err(storage)?; + (res.rows_affected() > 0) + .then_some(()) + .ok_or(DomainError::NotFound) + } +} + +pub struct SqliteBackupStrategies(pub DbPool); + +fn strategy_from_row(r: &SqliteRow) -> Result { + Ok(BackupStrategy { + id: r.get("id"), + name: r.get("name"), + source: serde_json::from_str::(r.get::("source").as_str()) + .map_err(|e| DomainError::Storage(e.to_string()))?, + schedule: r.get("schedule"), + target_id: r.get("target_id"), + retention: r.get::("retention") as u32, + passphrase: r.get("passphrase"), + enabled: r.get("enabled"), + }) +} + +const STRATEGY_COLS: &str = "id, name, source, schedule, target_id, retention, passphrase, enabled"; + +#[async_trait] +impl domain::ports::BackupStrategyRepository for SqliteBackupStrategies { + async fn list(&self) -> Result, DomainError> { + let rows = sqlx::query(&format!( + "SELECT {STRATEGY_COLS} FROM backup_strategies ORDER BY name" + )) + .fetch_all(&self.0) + .await + .map_err(storage)?; + rows.iter().map(strategy_from_row).collect() + } + async fn get(&self, id: Uuid) -> Result, DomainError> { + let row = sqlx::query(&format!( + "SELECT {STRATEGY_COLS} FROM backup_strategies WHERE id = ?" + )) + .bind(id) + .fetch_optional(&self.0) + .await + .map_err(storage)?; + row.as_ref().map(strategy_from_row).transpose() + } + async fn insert(&self, s: &BackupStrategy) -> Result<(), DomainError> { + let source = + serde_json::to_string(&s.source).map_err(|e| DomainError::Storage(e.to_string()))?; + sqlx::query(&format!( + "INSERT INTO backup_strategies ({STRATEGY_COLS}) VALUES (?, ?, ?, ?, ?, ?, ?, ?)" + )) + .bind(s.id) + .bind(&s.name) + .bind(source) + .bind(&s.schedule) + .bind(s.target_id) + .bind(s.retention as i64) + .bind(&s.passphrase) + .bind(s.enabled) + .execute(&self.0) + .await + .map(|_| ()) + .map_err(storage) + } + async fn update(&self, s: &BackupStrategy) -> Result<(), DomainError> { + let source = + serde_json::to_string(&s.source).map_err(|e| DomainError::Storage(e.to_string()))?; + let res = sqlx::query("UPDATE backup_strategies SET name = ?, source = ?, schedule = ?, target_id = ?, retention = ?, passphrase = ?, enabled = ? WHERE id = ?") + .bind(&s.name) + .bind(source) + .bind(&s.schedule) + .bind(s.target_id) + .bind(s.retention as i64) + .bind(&s.passphrase) + .bind(s.enabled) + .bind(s.id) + .execute(&self.0) + .await + .map_err(storage)?; + (res.rows_affected() > 0) + .then_some(()) + .ok_or(DomainError::NotFound) + } + async fn delete(&self, id: Uuid) -> Result<(), DomainError> { + let res = sqlx::query("DELETE FROM backup_strategies WHERE id = ?") + .bind(id) + .execute(&self.0) + .await + .map_err(storage)?; + (res.rows_affected() > 0) + .then_some(()) + .ok_or(DomainError::NotFound) + } +} + +pub struct SqliteBackupRecords(pub DbPool); + +#[async_trait] +impl domain::ports::BackupRecordRepository for SqliteBackupRecords { + async fn insert(&self, r: &BackupRecord) -> Result<(), DomainError> { + sqlx::query("INSERT INTO backup_records (id, strategy_id, filename, size_bytes, sha256, created_at) VALUES (?, ?, ?, ?, ?, ?)") + .bind(r.id) + .bind(r.strategy_id) + .bind(&r.filename) + .bind(r.size_bytes as i64) + .bind(&r.sha256) + .bind(r.created_at.to_rfc3339()) + .execute(&self.0) + .await + .map(|_| ()) + .map_err(storage) + } + async fn list_for( + &self, + strategy_id: Uuid, + limit: u32, + ) -> Result, DomainError> { + sqlx::query("SELECT id, strategy_id, filename, size_bytes, sha256, created_at FROM backup_records WHERE strategy_id = ? ORDER BY created_at DESC LIMIT ?") + .bind(strategy_id) + .bind(limit) + .fetch_all(&self.0) + .await + .map(|rows| { + rows.iter() + .map(|r| BackupRecord { + id: r.get("id"), + strategy_id: r.get("strategy_id"), + filename: r.get("filename"), + size_bytes: r.get::("size_bytes") as u64, + sha256: r.get("sha256"), + created_at: parse_ts(r.get::("created_at").as_str()), + }) + .collect() + }) + .map_err(storage) + } + async fn delete_by_filename( + &self, + strategy_id: Uuid, + filename: &str, + ) -> Result<(), DomainError> { + sqlx::query("DELETE FROM backup_records WHERE strategy_id = ? AND filename = ?") + .bind(strategy_id) + .bind(filename) + .execute(&self.0) + .await + .map(|_| ()) + .map_err(storage) + } +} + +#[cfg(test)] +mod backup_tests { + use super::*; + use domain::ports::{BackupRecordRepository, BackupStrategyRepository, BackupTargetRepository}; + + #[tokio::test] + async fn targets_strategies_records_roundtrip() { + let pool = crate::connect("sqlite::memory:").await.unwrap(); + let targets = SqliteBackupTargets(pool.clone()); + let strategies = SqliteBackupStrategies(pool.clone()); + let records = SqliteBackupRecords(pool); + let t = BackupTarget { + id: Uuid::new_v4(), + name: "NAS".into(), + kind: StorageKind::Ftp, + host: "h".into(), + port: Some(2121), + share: String::new(), + path: "p".into(), + username: "u".into(), + password: "enc".into(), + tls: true, + }; + targets.insert(&t).await.unwrap(); + assert_eq!(targets.get(t.id).await.unwrap().unwrap(), t); + let s = BackupStrategy { + id: Uuid::new_v4(), + name: "S".into(), + source: BackupSource::HostPath { + path: "/srv".into(), + }, + schedule: "0 0 1 * * *".into(), + target_id: t.id, + retention: 3, + passphrase: Some("enc".into()), + enabled: true, + }; + strategies.insert(&s).await.unwrap(); + assert_eq!(strategies.list().await.unwrap(), vec![s.clone()]); + let mut s2 = s.clone(); + s2.enabled = false; + strategies.update(&s2).await.unwrap(); + assert!(!strategies.get(s.id).await.unwrap().unwrap().enabled); + let r = BackupRecord { + id: Uuid::new_v4(), + strategy_id: s.id, + filename: "s_1.tar.gz".into(), + size_bytes: 10, + sha256: "x".into(), + created_at: Utc::now(), + }; + records.insert(&r).await.unwrap(); + assert_eq!(records.list_for(s.id, 10).await.unwrap().len(), 1); + records + .delete_by_filename(s.id, "s_1.tar.gz") + .await + .unwrap(); + assert!(records.list_for(s.id, 10).await.unwrap().is_empty()); + strategies.delete(s.id).await.unwrap(); + targets.delete(t.id).await.unwrap(); + assert_eq!( + targets.delete(t.id).await.unwrap_err(), + DomainError::NotFound + ); + } +} diff --git a/backend/crates/infrastructure/src/trivy.rs b/backend/crates/infrastructure/src/trivy.rs index 88b18f7..4a1ab50 100644 --- a/backend/crates/infrastructure/src/trivy.rs +++ b/backend/crates/infrastructure/src/trivy.rs @@ -317,6 +317,22 @@ mod tests { }, }) } + 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) } diff --git a/docs/restore.md b/docs/restore.md new file mode 100644 index 0000000..595e551 --- /dev/null +++ b/docs/restore.md @@ -0,0 +1,50 @@ +# Restoring backups + +Backups are plain archives on the SMB share or FTP server, named +`_.[.enc]`. The SHA-256 of every uploaded file is shown +in the backup history of the strategy. + +## 1. Decrypt (only for `.enc` files) + +```bash +openssl enc -d -aes-256-cbc -pbkdf2 -in FILE.sql.gz.enc -out FILE.sql.gz +``` + +Enter the passphrase of the strategy when prompted. + +## 2. Restore per source type + +### PostgreSQL dump (`.sql.gz`, produced by `pg_dumpall` inside the pod) + +```bash +gunzip -c gitea-db_*.sql.gz | microk8s kubectl exec -i -n gitea gitea-postgresql-0 -- \ + sh -c 'PGPASSWORD="$POSTGRES_PASSWORD" psql -U postgres' +``` + +Stop Gitea first (`microk8s kubectl scale deploy/gitea -n gitea --replicas=0`) and start it again afterwards. + +### Persistent volume (`.tar.gz` of the hostpath directory) + +Find the directory of the PVC on the host and unpack into it: + +```bash +PV=$(microk8s kubectl get pvc -n gitea gitea-shared-storage -o jsonpath='{.spec.volumeName}') +DIR=$(microk8s kubectl get pv "$PV" -o jsonpath='{.spec.hostPath.path}') +microk8s kubectl scale deploy/gitea -n gitea --replicas=0 +tar -xzf gitea-data_*.tar.gz -C "$DIR" +microk8s kubectl scale deploy/gitea -n gitea --replicas=1 +``` + +### Kubernetes manifests (`.yaml.gz`) + +```bash +gunzip -c gitea-manifests_*.yaml.gz | microk8s kubectl apply -f - +``` + +Secrets and PVC definitions are included; review before applying to a different cluster. + +### Host directory (`.tar.gz`) + +```bash +tar -xzf name_*.tar.gz -C /path/of/the/source +``` diff --git a/frontend/e2e/backups.spec.ts b/frontend/e2e/backups.spec.ts index 455a8b4..83b4f44 100644 --- a/frontend/e2e/backups.spec.ts +++ b/frontend/e2e/backups.spec.ts @@ -9,26 +9,26 @@ test('admin creates a target and a strategy, runs it and sees the backup', async await page.getByRole('button', { name: 'New target' }).click() const name = `NAS ${Date.now()}` - await page.getByLabel('Name').fill(name) - await page.getByLabel('Host').fill('nas.local') + await page.getByLabel('Name', { exact: true }).fill(name) + await page.getByLabel('Host', { exact: true }).fill('nas.local') await page.getByLabel('Share').fill('backups') await page.getByLabel('Directory').fill('softvisor') - await page.getByLabel('Username').fill('backup') - await page.getByLabel('Password').fill('smb-secret') + await page.getByLabel('Username', { exact: true }).fill('backup') + await page.getByLabel('Password', { exact: true }).fill('smb-secret') await page.getByRole('button', { name: 'Save target' }).click() const targetRow = page.getByRole('row', { name: new RegExp(name) }) await expect(targetRow).toBeVisible() await targetRow.getByRole('button', { name: 'Test' }).click() - await expect(page.getByRole('status')).toContainText('Connection ok') + await expect(page.getByText(/Connection ok/)).toBeVisible() await page.getByRole('button', { name: 'New strategy' }).click() - await page.getByLabel('Strategy name').fill('Gitea database') - await page.getByLabel('Source type').selectOption('postgres_dump') - await page.getByLabel('Namespace').fill('gitea') - await page.getByLabel('Pod').fill('gitea-postgresql-0') - await page.getByLabel('Schedule').fill('0 0 2 * * *') - await page.getByLabel('Target').selectOption({ label: name }) - await page.getByLabel('Keep last').fill('3') + await page.getByLabel('Strategy name', { exact: true }).fill('Gitea database') + await page.getByLabel('Source type', { exact: true }).selectOption('postgres_dump') + await page.getByLabel('Namespace', { exact: true }).fill('gitea') + await page.getByLabel('Pod', { exact: true }).fill('gitea-postgresql-0') + await page.getByLabel('Schedule', { exact: true }).fill('0 0 2 * * *') + await page.getByLabel('Target', { exact: true }).selectOption({ label: name }) + await page.getByLabel('Keep last', { exact: true }).fill('3') await page.getByRole('button', { name: 'Save strategy' }).click() const row = page.getByRole('row', { name: /Gitea database/ }) await expect(row).toBeVisible() diff --git a/frontend/src/api/client.ts b/frontend/src/api/client.ts index 739dacb..2c630d0 100644 --- a/frontend/src/api/client.ts +++ b/frontend/src/api/client.ts @@ -60,4 +60,5 @@ export const api = { post: (path: string, body?: unknown) => request('POST', path, body), patch: (path: string, body?: unknown) => request('PATCH', path, body), put: (path: string, body?: unknown) => request('PUT', path, body), + delete: (path: string) => request('DELETE', path), } diff --git a/frontend/src/components/StrategyForm.vue b/frontend/src/components/StrategyForm.vue new file mode 100644 index 0000000..3081e5e --- /dev/null +++ b/frontend/src/components/StrategyForm.vue @@ -0,0 +1,183 @@ + + + diff --git a/frontend/src/components/TargetForm.vue b/frontend/src/components/TargetForm.vue new file mode 100644 index 0000000..e381d87 --- /dev/null +++ b/frontend/src/components/TargetForm.vue @@ -0,0 +1,123 @@ + + + diff --git a/frontend/src/pages/BackupsPage.vue b/frontend/src/pages/BackupsPage.vue new file mode 100644 index 0000000..34a8814 --- /dev/null +++ b/frontend/src/pages/BackupsPage.vue @@ -0,0 +1,365 @@ + + + diff --git a/frontend/src/pages/PlaceholderPage.vue b/frontend/src/pages/PlaceholderPage.vue deleted file mode 100644 index c08143a..0000000 --- a/frontend/src/pages/PlaceholderPage.vue +++ /dev/null @@ -1,8 +0,0 @@ - - - diff --git a/frontend/src/router.ts b/frontend/src/router.ts index 38a87a7..22a9ab9 100644 --- a/frontend/src/router.ts +++ b/frontend/src/router.ts @@ -5,12 +5,12 @@ import AppShell from './components/AppShell.vue' import LoginPage from './pages/LoginPage.vue' import DashboardPage from './pages/DashboardPage.vue' import UsersPage from './pages/UsersPage.vue' -import PlaceholderPage from './pages/PlaceholderPage.vue' import SettingsPage from './pages/SettingsPage.vue' import JobsPage from './pages/JobsPage.vue' import UpdatesPage from './pages/UpdatesPage.vue' import ClusterPage from './pages/ClusterPage.vue' import VulnerabilitiesPage from './pages/VulnerabilitiesPage.vue' +import BackupsPage from './pages/BackupsPage.vue' export const router = createRouter({ history: createWebHistory(), @@ -24,7 +24,7 @@ export const router = createRouter({ { path: 'updates', component: UpdatesPage }, { path: 'cluster', component: ClusterPage }, { path: 'vulnerabilities', component: VulnerabilitiesPage }, - { path: 'backups', component: PlaceholderPage, props: { title: 'Backups' } }, + { path: 'backups', component: BackupsPage }, { path: 'users', component: UsersPage, meta: { admin: true } }, { path: 'jobs', component: JobsPage }, { path: 'settings', component: SettingsPage },