Files
Infrastruktur-Monitoring-Sy…/backend/crates/application/src/test_fakes.rs
Dennis Nemec 51364fdd76 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 <noreply@anthropic.com>
2026-09-03 20:50:11 +02:00

1054 lines
32 KiB
Rust

//! In-memory fakes for the ports, used by the unit tests of the use cases.
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use chrono::Utc;
use domain::auth::{AccessClaims, AuthEvent, RefreshToken};
use domain::ports::*;
use domain::user::{Role, User, UserUpdate};
use domain::DomainError;
use uuid::Uuid;
#[derive(Default)]
pub struct MemUsers(pub Mutex<HashMap<Uuid, User>>);
#[async_trait]
impl UserRepository for MemUsers {
async fn find_by_id(&self, id: Uuid) -> Result<Option<User>, DomainError> {
Ok(self.0.lock().unwrap().get(&id).cloned())
}
async fn find_by_email(&self, email: &str) -> Result<Option<User>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.values()
.find(|u| u.email == email)
.cloned())
}
async fn list(&self) -> Result<Vec<User>, DomainError> {
let mut v: Vec<_> = self.0.lock().unwrap().values().cloned().collect();
v.sort_by(|a, b| a.email.cmp(&b.email));
Ok(v)
}
async fn count(&self) -> Result<u64, DomainError> {
Ok(self.0.lock().unwrap().len() as u64)
}
async fn count_active_admins(&self) -> Result<u64, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.values()
.filter(|u| u.is_admin() && u.is_active)
.count() as u64)
}
async fn insert(&self, user: &User) -> Result<(), DomainError> {
self.0.lock().unwrap().insert(user.id, user.clone());
Ok(())
}
async fn update(&self, id: Uuid, update: &UserUpdate) -> Result<User, DomainError> {
let mut m = self.0.lock().unwrap();
let u = m.get_mut(&id).ok_or(DomainError::NotFound)?;
if let Some(n) = &update.display_name {
u.display_name = n.clone();
}
if let Some(r) = update.role {
u.role = r;
}
if let Some(a) = update.is_active {
u.is_active = a;
}
Ok(u.clone())
}
async fn set_password_hash(&self, id: Uuid, hash: &str) -> Result<(), DomainError> {
let mut m = self.0.lock().unwrap();
m.get_mut(&id).ok_or(DomainError::NotFound)?.password_hash = hash.into();
Ok(())
}
}
#[derive(Default)]
pub struct MemRefresh(pub Mutex<Vec<RefreshToken>>);
#[async_trait]
impl RefreshTokenRepository for MemRefresh {
async fn insert(&self, token: &RefreshToken) -> Result<(), DomainError> {
self.0.lock().unwrap().push(token.clone());
Ok(())
}
async fn find_by_hash(&self, hash: &str) -> Result<Option<RefreshToken>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.find(|t| t.token_hash == hash)
.cloned())
}
async fn revoke(&self, id: Uuid) -> Result<(), DomainError> {
self.0
.lock()
.unwrap()
.iter_mut()
.filter(|t| t.id == id)
.for_each(|t| t.revoked = true);
Ok(())
}
async fn revoke_family(&self, family: Uuid) -> Result<(), DomainError> {
self.0
.lock()
.unwrap()
.iter_mut()
.filter(|t| t.family == family)
.for_each(|t| t.revoked = true);
Ok(())
}
}
#[derive(Default)]
pub struct MemAudit(pub Mutex<Vec<AuthEvent>>);
#[async_trait]
impl AuditLog for MemAudit {
async fn record(&self, event: &AuthEvent) -> Result<(), DomainError> {
self.0.lock().unwrap().push(event.clone());
Ok(())
}
}
/// "Hashes" by prefixing; good enough to test the flow without Argon2 cost.
pub struct FakeHasher;
impl PasswordHasher for FakeHasher {
fn hash(&self, password: &str) -> Result<String, DomainError> {
Ok(format!("hashed:{password}"))
}
fn verify(&self, password: &str, hash: &str) -> bool {
hash == format!("hashed:{password}")
}
}
/// Access tokens are `"<uuid>:<role>"`; anything else is invalid.
pub struct FakeTokens;
impl AccessTokenIssuer for FakeTokens {
fn issue(&self, user: &User) -> Result<String, DomainError> {
Ok(format!("{}:{}", user.id, user.role.as_str()))
}
fn verify(&self, token: &str) -> Result<AccessClaims, DomainError> {
let (id, role) = token.split_once(':').ok_or(DomainError::InvalidToken)?;
Ok(AccessClaims {
sub: id.parse().map_err(|_| DomainError::InvalidToken)?,
role: Role::parse(role).ok_or(DomainError::InvalidToken)?,
exp: 0,
})
}
}
pub fn user(email: &str, password: &str, role: Role, active: bool) -> User {
User {
id: Uuid::new_v4(),
email: email.into(),
display_name: email.split('@').next().unwrap().into(),
password_hash: format!("hashed:{password}"),
role,
is_active: active,
created_at: Utc::now(),
}
}
pub struct Fixture {
pub users: Arc<MemUsers>,
pub refresh: Arc<MemRefresh>,
pub audit: Arc<MemAudit>,
pub auth: crate::AuthService,
pub svc: crate::UserService,
}
pub fn fixture() -> Fixture {
let users = Arc::new(MemUsers::default());
let refresh = Arc::new(MemRefresh::default());
let audit = Arc::new(MemAudit::default());
let auth = crate::AuthService::new(
users.clone(),
refresh.clone(),
audit.clone(),
Arc::new(FakeHasher),
Arc::new(FakeTokens),
);
let svc = crate::UserService::new(users.clone(), Arc::new(FakeHasher));
Fixture {
users,
refresh,
audit,
auth,
svc,
}
}
use domain::jobs::{JobKind, JobRun, JobStatus};
use domain::ports::{Cipher, JobRunRepository, Mailer, SettingsRepository};
use domain::settings::SmtpSettings;
#[derive(Default)]
pub struct MemSettings(pub Mutex<HashMap<String, String>>);
#[async_trait]
impl SettingsRepository for MemSettings {
async fn get(&self, key: &str) -> Result<Option<String>, DomainError> {
Ok(self.0.lock().unwrap().get(key).cloned())
}
async fn set(&self, key: &str, value: &str) -> Result<(), DomainError> {
self.0.lock().unwrap().insert(key.into(), value.into());
Ok(())
}
}
/// Reversible "encryption" so tests can assert the stored value is not plain text.
pub struct FakeCipher;
impl Cipher for FakeCipher {
fn encrypt(&self, plain: &str) -> Result<String, DomainError> {
Ok(format!("enc:{}", plain.chars().rev().collect::<String>()))
}
fn decrypt(&self, c: &str) -> Result<String, DomainError> {
c.strip_prefix("enc:")
.map(|s| s.chars().rev().collect())
.ok_or(DomainError::Storage("bad cipher text".into()))
}
}
#[derive(Default)]
pub struct MemMailer(pub Mutex<Vec<(Vec<String>, String, String)>>);
#[async_trait]
impl Mailer for MemMailer {
async fn send(
&self,
_smtp: &SmtpSettings,
to: &[String],
subject: &str,
body: &str,
) -> Result<(), DomainError> {
self.0
.lock()
.unwrap()
.push((to.to_vec(), subject.into(), body.into()));
Ok(())
}
}
#[derive(Default)]
pub struct MemJobRuns(pub Mutex<Vec<JobRun>>);
#[async_trait]
impl JobRunRepository for MemJobRuns {
async fn insert(&self, run: &JobRun) -> Result<(), DomainError> {
self.0.lock().unwrap().push(run.clone());
Ok(())
}
async fn append_log(&self, id: Uuid, line: &str) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
let r = v
.iter_mut()
.find(|r| r.id == id)
.ok_or(DomainError::NotFound)?;
r.log.push_str(line);
r.log.push('\n');
Ok(())
}
async fn finish(&self, id: Uuid, status: JobStatus) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
let r = v
.iter_mut()
.find(|r| r.id == id)
.ok_or(DomainError::NotFound)?;
r.status = status;
r.finished_at = Some(Utc::now());
Ok(())
}
async fn get(&self, id: Uuid) -> Result<Option<JobRun>, DomainError> {
Ok(self.0.lock().unwrap().iter().find(|r| r.id == id).cloned())
}
async fn list(&self, limit: u32) -> Result<Vec<JobRun>, DomainError> {
let v = self.0.lock().unwrap();
Ok(v.iter().rev().take(limit as usize).cloned().collect())
}
async fn find_running(&self, kind: JobKind) -> Result<Option<JobRun>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.find(|r| r.kind == kind && r.status == JobStatus::Running)
.cloned())
}
async fn last_finished(&self, kind: JobKind) -> Result<Option<JobRun>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.rev()
.find(|r| r.kind == kind && r.status != JobStatus::Running)
.cloned())
}
async fn running(&self) -> Result<Vec<JobRun>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.filter(|r| r.status == JobStatus::Running)
.cloned()
.collect())
}
}
pub fn smtp() -> SmtpSettings {
SmtpSettings {
host: "mail.example.com".into(),
port: 587,
security: domain::settings::SmtpSecurity::StartTls,
username: "bot".into(),
password: "s3cret".into(),
from: "monitoring@example.com".into(),
notify_to: vec!["ops@example.com".into()],
}
}
use domain::host::{Inventory, OsInfo, Package, PackageSource};
use domain::ports::{HostInspector, InventoryRepository};
pub struct FakeInspector {
pub fail: bool,
}
#[async_trait]
impl HostInspector for FakeInspector {
async fn os_info(&self) -> Result<OsInfo, DomainError> {
if self.fail {
return Err(DomainError::Unavailable("host down".into()));
}
Ok(OsInfo {
hostname: "srv".into(),
name: "Debian GNU/Linux 12 (bookworm)".into(),
version: "12".into(),
kernel: "6.1.0-42-amd64".into(),
uptime_secs: 3600,
reboot_required: true,
})
}
async fn services(&self) -> Result<Vec<domain::host::HostService>, 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<Vec<Package>, DomainError> {
Ok(vec![
Package {
name: "bash".into(),
source: PackageSource::Apt,
installed: "5.2".into(),
candidate: None,
is_security: false,
},
Package {
name: "openssl".into(),
source: PackageSource::Apt,
installed: "3.0.1".into(),
candidate: Some("3.0.2".into()),
is_security: true,
},
Package {
name: "microk8s".into(),
source: PackageSource::Snap,
installed: "v1.32.13".into(),
candidate: None,
is_security: false,
},
])
}
}
#[derive(Default)]
pub struct MemInventory(pub Mutex<Option<Inventory>>);
#[async_trait]
impl InventoryRepository for MemInventory {
async fn save(&self, inventory: &Inventory) -> Result<(), DomainError> {
*self.0.lock().unwrap() = Some(inventory.clone());
Ok(())
}
async fn load(&self) -> Result<Option<Inventory>, DomainError> {
Ok(self.0.lock().unwrap().clone())
}
}
use domain::ports::{HostUpdater, LineSink};
/// Records the requested packages and emits a few lines.
#[derive(Default)]
pub struct FakeUpdater {
pub calls: Mutex<Vec<Vec<String>>>,
pub fail: bool,
}
#[async_trait]
impl HostUpdater for FakeUpdater {
async fn upgrade(&self, packages: &[String], out: &dyn LineSink) -> Result<(), DomainError> {
self.calls.lock().unwrap().push(packages.to_vec());
out.line("Reading package lists...");
out.line(&format!(
"Upgrading {} package(s)",
if packages.is_empty() {
"all".to_string()
} else {
packages.len().to_string()
}
));
if self.fail {
return Err(DomainError::Unavailable(
"apt-get exited with status 100".into(),
));
}
Ok(())
}
}
use domain::cluster::{
ClusterOverview, Container, NodeInfo, VolumeClaim, Workload, WorkloadKind, WorkloadRef,
};
use domain::ports::ClusterGateway;
#[derive(Default)]
pub struct MemCluster {
pub actions: Mutex<Vec<String>>,
pub fail: std::sync::atomic::AtomicBool,
}
fn helm_labels(instance: &str, role: &str) -> std::collections::BTreeMap<String, String> {
[
("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 {
name: "node1".into(),
version: "v1.32.13".into(),
ready: true,
os_image: "Debian GNU/Linux 12 (bookworm)".into(),
kernel: "6.1.0-42-amd64".into(),
container_runtime: "containerd://1.6.36".into(),
}],
namespaces: vec!["default".into(), "gitea".into()],
workloads: vec![
Workload {
namespace: "gitea".into(),
kind: WorkloadKind::Deployment,
name: "gitea".into(),
ready: 1,
desired: 1,
containers: vec![Container {
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(),
kind: WorkloadKind::StatefulSet,
name: "gitea-postgresql".into(),
ready: 1,
desired: 1,
containers: vec![Container {
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 {
namespace: "gitea".into(),
name: "gitea-shared-storage".into(),
capacity: "10Gi".into(),
storage_class: "microk8s-hostpath".into(),
status: "Bound".into(),
}],
fetched_at: Utc::now(),
}
}
#[async_trait]
impl ClusterGateway for MemCluster {
async fn overview(&self) -> Result<ClusterOverview, DomainError> {
if self.fail.load(std::sync::atomic::Ordering::SeqCst) {
return Err(DomainError::Unavailable("cluster unreachable".into()));
}
Ok(sample_overview())
}
async fn restart(&self, w: &WorkloadRef) -> Result<(), DomainError> {
self.actions.lock().unwrap().push(format!(
"restart {}/{}/{}",
w.namespace,
w.kind.as_str(),
w.name
));
Ok(())
}
async fn scale(&self, w: &WorkloadRef, replicas: i32) -> Result<(), DomainError> {
self.actions.lock().unwrap().push(format!(
"scale {}/{}/{} {replicas}",
w.namespace,
w.kind.as_str(),
w.name
));
Ok(())
}
async fn set_image(
&self,
w: &WorkloadRef,
container: &str,
image: &str,
) -> Result<(), DomainError> {
self.actions.lock().unwrap().push(format!(
"image {}/{}/{} {container}={image}",
w.namespace,
w.kind.as_str(),
w.name
));
Ok(())
}
}
use domain::ports::{FindingRepository, VulnerabilityScanner};
use domain::vuln::{
Finding, FindingFilter, FindingGroup, FindingStatus, RawFinding, Severity, SeverityCounts,
TargetKind,
};
pub fn raw(
cve: &str,
pkg: &str,
installed: &str,
sev: Severity,
fixed: Option<&str>,
) -> RawFinding {
RawFinding {
cve_id: cve.into(),
severity: sev,
package: pkg.into(),
installed_version: installed.into(),
fixed_version: fixed.map(String::from),
title: format!("{cve} in {pkg}"),
url: format!("https://nvd.nist.gov/vuln/detail/{cve}"),
source: "debian".into(),
}
}
pub type ScanResults = Arc<Mutex<HashMap<String, Result<Vec<RawFinding>, String>>>>;
/// Scanner returning configurable results per target ("os" or image ref).
#[derive(Default)]
pub struct FakeScanner {
pub results: ScanResults,
}
impl FakeScanner {
pub fn with(self, target: &str, r: Result<Vec<RawFinding>, &str>) -> Self {
self.results
.lock()
.unwrap()
.insert(target.into(), r.map_err(String::from));
self
}
fn get(&self, target: &str) -> Result<Vec<RawFinding>, DomainError> {
match self.results.lock().unwrap().get(target) {
Some(Ok(v)) => Ok(v.clone()),
Some(Err(e)) => Err(DomainError::Unavailable(e.clone())),
None => Ok(vec![]),
}
}
}
#[async_trait]
impl VulnerabilityScanner for FakeScanner {
async fn version(&self) -> Result<String, DomainError> {
Ok("fake 0.1".into())
}
async fn scan_os(
&self,
out: &dyn domain::ports::LineSink,
) -> Result<Vec<RawFinding>, DomainError> {
out.line("scanning os");
self.get("os")
}
async fn scan_image(
&self,
image: &str,
out: &dyn domain::ports::LineSink,
) -> Result<Vec<RawFinding>, DomainError> {
out.line(&format!("scanning {image}"));
self.get(image)
}
}
#[derive(Default)]
pub struct MemFindings(pub Mutex<Vec<Finding>>);
#[async_trait]
impl FindingRepository for MemFindings {
async fn active_by_target(&self, target: &str) -> Result<Vec<Finding>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.filter(|f| f.target == target && f.status != FindingStatus::Fixed)
.cloned()
.collect())
}
async fn insert(&self, finding: &Finding) -> Result<(), DomainError> {
self.0.lock().unwrap().push(finding.clone());
Ok(())
}
async fn refresh(
&self,
updates: &[(Uuid, RawFinding)],
last_seen: chrono::DateTime<Utc>,
) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
for (id, raw) in updates {
if let Some(f) = v.iter_mut().find(|f| f.id == *id) {
f.raw = raw.clone();
f.last_seen = last_seen;
}
}
Ok(())
}
async fn set_status(&self, id: Uuid, status: FindingStatus) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
let f = v
.iter_mut()
.find(|f| f.id == id)
.ok_or(DomainError::NotFound)?;
f.status = status;
Ok(())
}
async fn get(&self, id: Uuid) -> Result<Option<Finding>, DomainError> {
Ok(self.0.lock().unwrap().iter().find(|f| f.id == id).cloned())
}
async fn list(&self, filter: &FindingFilter) -> Result<Vec<Finding>, DomainError> {
let mut v: Vec<Finding> = self
.0
.lock()
.unwrap()
.iter()
.filter(|f| filter.include_fixed || f.status != FindingStatus::Fixed)
.filter(|f| filter.min_severity.is_none_or(|m| f.raw.severity >= m))
.filter(|f| filter.target_kind.is_none_or(|k| f.target_kind == k))
.filter(|f| filter.package.as_ref().is_none_or(|p| &f.raw.package == p))
.filter(|f| filter.target.as_ref().is_none_or(|t| &f.target == t))
.filter(|f| filter.status.is_none_or(|s| f.status == s))
.cloned()
.collect();
v.sort_by(|a, b| {
b.raw
.severity
.cmp(&a.raw.severity)
.then(a.raw.cve_id.cmp(&b.raw.cve_id))
});
Ok(v)
}
async fn groups(&self, filter: &FindingFilter) -> Result<Vec<FindingGroup>, DomainError> {
let per_image = filter.target_kind == Some(TargetKind::Image);
let mut by_key: HashMap<String, FindingGroup> = HashMap::new();
let mut packages: HashMap<String, std::collections::HashSet<String>> = HashMap::new();
for f in self.list(filter).await? {
let key = if per_image {
f.target.clone()
} else {
f.raw.package.clone()
};
packages
.entry(key.clone())
.or_default()
.insert(f.raw.package.clone());
let g = by_key.entry(key.clone()).or_insert_with(|| FindingGroup {
key,
kind: f.target_kind,
source: if per_image {
String::new()
} else {
f.raw.source.clone()
},
installed: if per_image {
String::new()
} else {
f.raw.installed_version.clone()
},
counts: SeverityCounts::default(),
total: 0,
fixable: 0,
packages: 0,
});
g.counts.add(f.raw.severity);
g.total += 1;
if f.raw.fixed_version.is_some() {
g.fixable += 1;
}
}
let mut groups: Vec<FindingGroup> = by_key
.into_values()
.map(|mut g| {
g.packages = packages.get(&g.key).map(|p| p.len()).unwrap_or(1);
g
})
.collect();
let worst = |g: &FindingGroup| {
Severity::ALL
.iter()
.position(|s| match s {
Severity::Critical => g.counts.critical > 0,
Severity::High => g.counts.high > 0,
Severity::Medium => g.counts.medium > 0,
Severity::Low => g.counts.low > 0,
Severity::Unknown => g.counts.unknown > 0,
})
.unwrap_or(usize::MAX)
};
groups.sort_by(|a, b| {
worst(a)
.cmp(&worst(b))
.then(b.total.cmp(&a.total))
.then(a.key.cmp(&b.key))
});
Ok(groups)
}
async fn counts(&self, kind: Option<TargetKind>) -> Result<SeverityCounts, DomainError> {
let mut c = SeverityCounts::default();
for f in
self.0.lock().unwrap().iter().filter(|f| {
f.status != FindingStatus::Fixed && kind.is_none_or(|k| f.target_kind == k)
})
{
c.add(f.raw.severity);
}
Ok(c)
}
}
use domain::backup::{
BackupRecord, BackupSource, BackupStrategy, BackupTarget, RemoteFile, StorageKind,
};
use domain::ports::{
BackupCollector, BackupRecordRepository, BackupStorage, BackupStrategyRepository,
BackupTargetRepository, FileEncryptor, LineSink as _LineSink,
};
use std::path::{Path, PathBuf};
#[derive(Default)]
pub struct MemTargets(pub Mutex<Vec<BackupTarget>>);
#[async_trait]
impl BackupTargetRepository for MemTargets {
async fn list(&self) -> Result<Vec<BackupTarget>, DomainError> {
Ok(self.0.lock().unwrap().clone())
}
async fn get(&self, id: Uuid) -> Result<Option<BackupTarget>, DomainError> {
Ok(self.0.lock().unwrap().iter().find(|t| t.id == id).cloned())
}
async fn insert(&self, t: &BackupTarget) -> Result<(), DomainError> {
self.0.lock().unwrap().push(t.clone());
Ok(())
}
async fn update(&self, t: &BackupTarget) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
let x = v
.iter_mut()
.find(|x| x.id == t.id)
.ok_or(DomainError::NotFound)?;
*x = t.clone();
Ok(())
}
async fn delete(&self, id: Uuid) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
let before = v.len();
v.retain(|t| t.id != id);
(v.len() < before)
.then_some(())
.ok_or(DomainError::NotFound)
}
}
#[derive(Default)]
pub struct MemStrategies(pub Mutex<Vec<BackupStrategy>>);
#[async_trait]
impl BackupStrategyRepository for MemStrategies {
async fn list(&self) -> Result<Vec<BackupStrategy>, DomainError> {
Ok(self.0.lock().unwrap().clone())
}
async fn get(&self, id: Uuid) -> Result<Option<BackupStrategy>, DomainError> {
Ok(self.0.lock().unwrap().iter().find(|t| t.id == id).cloned())
}
async fn insert(&self, s: &BackupStrategy) -> Result<(), DomainError> {
self.0.lock().unwrap().push(s.clone());
Ok(())
}
async fn update(&self, s: &BackupStrategy) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
let x = v
.iter_mut()
.find(|x| x.id == s.id)
.ok_or(DomainError::NotFound)?;
*x = s.clone();
Ok(())
}
async fn delete(&self, id: Uuid) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
let before = v.len();
v.retain(|t| t.id != id);
(v.len() < before)
.then_some(())
.ok_or(DomainError::NotFound)
}
}
#[derive(Default)]
pub struct MemRecords(pub Mutex<Vec<BackupRecord>>);
#[async_trait]
impl BackupRecordRepository for MemRecords {
async fn insert(&self, r: &BackupRecord) -> Result<(), DomainError> {
self.0.lock().unwrap().push(r.clone());
Ok(())
}
async fn list_for(
&self,
strategy_id: Uuid,
limit: u32,
) -> Result<Vec<BackupRecord>, DomainError> {
let mut v: Vec<_> = self
.0
.lock()
.unwrap()
.iter()
.filter(|r| r.strategy_id == strategy_id)
.cloned()
.collect();
v.sort_by_key(|r| std::cmp::Reverse(r.created_at));
v.truncate(limit as usize);
Ok(v)
}
async fn delete_by_filename(
&self,
strategy_id: Uuid,
filename: &str,
) -> Result<(), DomainError> {
self.0
.lock()
.unwrap()
.retain(|r| !(r.strategy_id == strategy_id && r.filename == filename));
Ok(())
}
}
/// In-memory remote storage keyed by target id; `fail` makes every call fail.
#[derive(Default)]
pub struct MemStorage {
pub files: Mutex<HashMap<Uuid, Vec<RemoteFile>>>,
pub fail: bool,
pub ops: Mutex<Vec<String>>,
}
#[async_trait]
impl BackupStorage for MemStorage {
async fn test(&self, t: &BackupTarget) -> Result<(), DomainError> {
self.ops.lock().unwrap().push(format!("test {}", t.name));
if self.fail {
Err(DomainError::Unavailable("connection refused".into()))
} else {
Ok(())
}
}
async fn upload(
&self,
t: &BackupTarget,
local: &Path,
remote_name: &str,
) -> Result<(), DomainError> {
if self.fail {
return Err(DomainError::Unavailable("upload failed".into()));
}
let size = std::fs::metadata(local).map(|m| m.len()).unwrap_or(0);
self.ops
.lock()
.unwrap()
.push(format!("upload {remote_name}"));
self.files
.lock()
.unwrap()
.entry(t.id)
.or_default()
.push(RemoteFile {
name: remote_name.into(),
size_bytes: size,
});
Ok(())
}
async fn list(&self, t: &BackupTarget) -> Result<Vec<RemoteFile>, DomainError> {
Ok(self
.files
.lock()
.unwrap()
.get(&t.id)
.cloned()
.unwrap_or_default())
}
async fn delete(&self, t: &BackupTarget, remote_name: &str) -> Result<(), DomainError> {
self.ops
.lock()
.unwrap()
.push(format!("delete {remote_name}"));
self.files
.lock()
.unwrap()
.entry(t.id)
.or_default()
.retain(|f| f.name != remote_name);
Ok(())
}
}
/// Writes a small file describing the source.
pub struct FakeCollector;
#[async_trait]
impl BackupCollector for FakeCollector {
async fn collect(
&self,
source: &BackupSource,
work_dir: &Path,
out: &dyn _LineSink,
) -> Result<PathBuf, DomainError> {
out.line(&format!("collecting {source:?}"));
let p = work_dir.join(format!("archive.{}", source.extension()));
std::fs::write(&p, format!("fake archive of {source:?}"))
.map_err(|e| DomainError::Storage(e.to_string()))?;
Ok(p)
}
}
pub struct FakeEncryptor;
#[async_trait]
impl FileEncryptor for FakeEncryptor {
async fn encrypt(&self, input: &Path, passphrase: &str) -> Result<PathBuf, DomainError> {
let out = input.with_extension(format!(
"{}.enc",
input.extension().and_then(|e| e.to_str()).unwrap_or("")
));
let data = std::fs::read(input).map_err(|e| DomainError::Storage(e.to_string()))?;
std::fs::write(&out, [b"ENC:", passphrase.as_bytes(), b":", &data].concat())
.map_err(|e| DomainError::Storage(e.to_string()))?;
Ok(out)
}
}
pub fn target(name: &str) -> BackupTarget {
BackupTarget {
id: Uuid::new_v4(),
name: name.into(),
kind: StorageKind::Smb,
host: "nas.local".into(),
port: None,
share: "backups".into(),
path: "softvisor".into(),
username: "backup".into(),
password: "smb-secret".into(),
tls: false,
}
}
pub fn strategy(name: &str, target_id: Uuid) -> BackupStrategy {
BackupStrategy {
id: Uuid::new_v4(),
name: name.into(),
source: BackupSource::PostgresDump {
namespace: "gitea".into(),
pod: "gitea-postgresql-0".into(),
},
schedule: "0 0 2 * * *".into(),
target_id,
retention: 2,
passphrase: None,
enabled: true,
}
}
use domain::image::ImageUpdate;
use domain::ports::{ImageRegistry, ImageUpdateRepository};
#[derive(Default)]
pub struct MemRegistry {
pub tags: Mutex<Vec<String>>,
pub fail: std::sync::atomic::AtomicBool,
}
#[async_trait]
impl ImageRegistry for MemRegistry {
async fn tags(&self, _image: &str) -> Result<Vec<String>, DomainError> {
if self.fail.load(std::sync::atomic::Ordering::SeqCst) {
return Err(DomainError::Unavailable("registry unreachable".into()));
}
Ok(self.tags.lock().unwrap().clone())
}
}
#[derive(Default)]
pub struct MemImageUpdates(pub Mutex<Vec<ImageUpdate>>);
#[async_trait]
impl ImageUpdateRepository for MemImageUpdates {
async fn upsert(&self, update: &ImageUpdate) -> Result<(), DomainError> {
let mut v = self.0.lock().unwrap();
v.retain(|u| u.image != update.image);
v.push(update.clone());
Ok(())
}
async fn get(&self, image: &str) -> Result<Option<ImageUpdate>, DomainError> {
Ok(self
.0
.lock()
.unwrap()
.iter()
.find(|u| u.image == image)
.cloned())
}
async fn list(&self) -> Result<Vec<ImageUpdate>, DomainError> {
Ok(self.0.lock().unwrap().clone())
}
}
/// Settings service on in-memory stores, for services that only need it as a dependency.
pub fn settings_service() -> Arc<crate::SettingsService> {
Arc::new(crate::SettingsService::new(
Arc::new(MemSettings::default()),
Arc::new(FakeCipher),
Arc::new(MemMailer::default()),
))
}