WP-30/31/32: backup targets, strategies and execution
SMB (smbclient) and FTP/FTPS (curl with netrc) targets with encrypted credentials and connection test; strategies with cron schedule, retention, optional openssl AES-256 encryption; sources: PVC hostpath tar, pg_dumpall in the Postgres pod, namespace manifests, host directory. Backup job collects, encrypts, uploads, verifies size, records sha256 and applies retention on the target; scheduler starts due strategies. Backups page with target/strategy forms, run now and history. Restore guide in docs. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
@ -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<dyn BackupTargetRepository>,
|
||||
@ -42,65 +44,316 @@ impl BackupService {
|
||||
Self { d: deps }
|
||||
}
|
||||
|
||||
fn decrypt_target(&self, mut t: BackupTarget) -> Result<BackupTarget, DomainError> {
|
||||
if !t.password.is_empty() {
|
||||
t.password = self.d.cipher.decrypt(&t.password)?;
|
||||
}
|
||||
Ok(t)
|
||||
}
|
||||
|
||||
fn encrypt_target(&self, mut t: BackupTarget) -> Result<BackupTarget, DomainError> {
|
||||
if !t.password.is_empty() {
|
||||
t.password = self.d.cipher.encrypt(&t.password)?;
|
||||
}
|
||||
Ok(t)
|
||||
}
|
||||
|
||||
fn decrypt_strategy(&self, mut s: BackupStrategy) -> Result<BackupStrategy, DomainError> {
|
||||
if let Some(p) = &s.passphrase {
|
||||
s.passphrase = Some(self.d.cipher.decrypt(p)?);
|
||||
}
|
||||
Ok(s)
|
||||
}
|
||||
|
||||
fn encrypt_strategy(&self, mut s: BackupStrategy) -> Result<BackupStrategy, DomainError> {
|
||||
if let Some(p) = &s.passphrase {
|
||||
s.passphrase = Some(self.d.cipher.encrypt(p)?);
|
||||
}
|
||||
Ok(s)
|
||||
}
|
||||
|
||||
// ---- targets ----
|
||||
pub async fn list_targets(&self) -> Result<Vec<BackupTarget>, 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<BackupTarget, DomainError> {
|
||||
todo!()
|
||||
|
||||
pub async fn get_target(&self, id: Uuid) -> Result<BackupTarget, DomainError> {
|
||||
self.decrypt_target(self.d.targets.get(id).await?.ok_or(DomainError::NotFound)?)
|
||||
}
|
||||
pub async fn create_target(&self, _t: BackupTarget) -> Result<BackupTarget, DomainError> {
|
||||
todo!()
|
||||
|
||||
pub async fn create_target(&self, mut t: BackupTarget) -> Result<BackupTarget, DomainError> {
|
||||
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<BackupTarget, DomainError> {
|
||||
todo!()
|
||||
pub async fn update_target(&self, mut t: BackupTarget) -> Result<BackupTarget, DomainError> {
|
||||
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<Vec<StrategyStatus>, 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<BackupStrategy, DomainError> {
|
||||
todo!()
|
||||
|
||||
pub async fn get_strategy(&self, id: Uuid) -> Result<BackupStrategy, DomainError> {
|
||||
self.decrypt_strategy(
|
||||
self.d
|
||||
.strategies
|
||||
.get(id)
|
||||
.await?
|
||||
.ok_or(DomainError::NotFound)?,
|
||||
)
|
||||
}
|
||||
pub async fn create_strategy(&self, _s: BackupStrategy) -> Result<BackupStrategy, DomainError> {
|
||||
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<BackupStrategy, DomainError> {
|
||||
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<BackupStrategy, DomainError> {
|
||||
todo!()
|
||||
pub async fn update_strategy(
|
||||
&self,
|
||||
mut s: BackupStrategy,
|
||||
) -> Result<BackupStrategy, DomainError> {
|
||||
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<Vec<BackupRecord>, DomainError> {
|
||||
todo!()
|
||||
|
||||
pub async fn records(&self, strategy_id: Uuid) -> Result<Vec<BackupRecord>, 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<Utc>,
|
||||
_grace_secs: i64,
|
||||
now: DateTime<Utc>,
|
||||
grace_secs: i64,
|
||||
) -> Result<Vec<Uuid>, 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<BackupRecord, DomainError> {
|
||||
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<BackupRecord, DomainError> {
|
||||
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::<String>();
|
||||
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<String> = 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<BackupService>);
|
||||
|
||||
#[async_trait]
|
||||
impl JobHandler for BackupJob {
|
||||
async fn run(&self, _params: Option<String>, _log: &dyn JobLog) -> Result<(), String> {
|
||||
todo!()
|
||||
async fn run(&self, params: Option<String>, 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(())
|
||||
}
|
||||
}
|
||||
|
||||
@ -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<JobRunner>,
|
||||
settings: Arc<SettingsService>,
|
||||
backups: Option<Arc<BackupService>>,
|
||||
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<BackupService>) -> 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<Utc>) -> Vec<JobKind> {
|
||||
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
|
||||
}
|
||||
|
||||
|
||||
@ -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)
|
||||
}
|
||||
|
||||
@ -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<String>);
|
||||
pub struct ChannelSink(tokio::sync::mpsc::UnboundedSender<String>);
|
||||
|
||||
pub fn channel_sink(tx: tokio::sync::mpsc::UnboundedSender<String>) -> ChannelSink {
|
||||
ChannelSink(tx)
|
||||
}
|
||||
|
||||
impl LineSink for ChannelSink {
|
||||
fn line(&self, text: &str) {
|
||||
|
||||
Reference in New Issue
Block a user