WP-02: encrypted settings, SMTP mail, job runner and cron scheduler
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled
CI / ui (push) Has been cancelled

AES-256-GCM secret storage keyed by MASTER_KEY, SMTP settings with test mail
(lettre), persisted job runs with log and status, JobRunner with per-kind
concurrency guard, 6-field cron schedules with defaults, scheduler loop.
Settings and Jobs pages in the UI. FAKE_HOST mode for dev machines.

Tests: 30 application, 8 infrastructure, 23 API, 16 Vitest, 6 Playwright.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Dennis Nemec
2026-09-02 22:16:56 +02:00
parent 6113af0a19
commit fa9ac9a6bc
31 changed files with 1632 additions and 58 deletions

View File

@ -3,7 +3,8 @@ use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use domain::jobs::{JobKind, JobRun};
use chrono::Utc;
use domain::jobs::{JobKind, JobRun, JobStatus};
use domain::ports::JobRunRepository;
use domain::DomainError;
use uuid::Uuid;
@ -24,6 +25,24 @@ pub struct JobRunner {
handlers: HashMap<JobKind, Arc<dyn JobHandler>>,
}
struct RepoLog {
runs: Arc<dyn JobRunRepository>,
id: Uuid,
}
#[async_trait]
impl JobLog for RepoLog {
async fn line(&self, text: &str) {
if let Err(e) = self.runs.append_log(self.id, text).await {
tracing_line(&format!("failed to append job log: {e}"));
}
}
}
fn tracing_line(msg: &str) {
eprintln!("{msg}");
}
impl JobRunner {
pub fn new(runs: Arc<dyn JobRunRepository>) -> Self {
Self {
@ -43,31 +62,88 @@ impl JobRunner {
k
}
async fn begin(
&self,
kind: JobKind,
params: Option<String>,
triggered_by: &str,
) -> Result<(JobRun, Arc<dyn JobHandler>), DomainError> {
let handler = self
.handlers
.get(&kind)
.cloned()
.ok_or(DomainError::NotFound)?;
if self.runs.find_running(kind).await?.is_some() {
return Err(DomainError::Conflict(format!(
"{} is already running",
kind.as_str()
)));
}
let run = JobRun {
id: Uuid::new_v4(),
kind,
params,
status: JobStatus::Running,
started_at: Utc::now(),
finished_at: None,
log: String::new(),
triggered_by: triggered_by.into(),
};
self.runs.insert(&run).await?;
Ok((run, handler))
}
async fn execute(runs: Arc<dyn JobRunRepository>, handler: Arc<dyn JobHandler>, run: &JobRun) {
let log = RepoLog {
runs: runs.clone(),
id: run.id,
};
let status = match handler.run(run.params.clone(), &log).await {
Ok(()) => JobStatus::Success,
Err(e) => {
log.line(&format!("ERROR: {e}")).await;
JobStatus::Failed
}
};
if let Err(e) = runs.finish(run.id, status).await {
tracing_line(&format!("failed to finish job: {e}"));
}
}
/// Start a job in the background. Fails with `Conflict` if the kind is already running.
pub async fn start(
&self,
_kind: JobKind,
_params: Option<String>,
_triggered_by: &str,
kind: JobKind,
params: Option<String>,
triggered_by: &str,
) -> Result<JobRun, DomainError> {
todo!()
let (run, handler) = self.begin(kind, params, triggered_by).await?;
let (runs, run_clone) = (self.runs.clone(), run.clone());
tokio::spawn(async move { Self::execute(runs, handler, &run_clone).await });
Ok(run)
}
/// Run a job and wait for it to finish (used by tests and the scheduler).
pub async fn run_and_wait(
&self,
_kind: JobKind,
_params: Option<String>,
_triggered_by: &str,
kind: JobKind,
params: Option<String>,
triggered_by: &str,
) -> Result<JobRun, DomainError> {
todo!()
let (run, handler) = self.begin(kind, params, triggered_by).await?;
Self::execute(self.runs.clone(), handler, &run).await;
self.get(run.id).await
}
pub async fn get(&self, _id: Uuid) -> Result<JobRun, DomainError> {
todo!()
pub async fn get(&self, id: Uuid) -> Result<JobRun, DomainError> {
self.runs.get(id).await?.ok_or(DomainError::NotFound)
}
pub async fn list(&self, _limit: u32) -> Result<Vec<JobRun>, DomainError> {
todo!()
pub async fn list(&self, limit: u32) -> Result<Vec<JobRun>, DomainError> {
self.runs.list(limit).await
}
pub async fn last_finished(&self, kind: JobKind) -> Result<Option<JobRun>, DomainError> {
self.runs.last_finished(kind).await
}
}

View File

@ -1,30 +1,84 @@
//! Cron scheduler: decides which scheduled job kinds are due.
//! Cron scheduler: decides which scheduled job kinds are due and runs them.
use std::str::FromStr;
use std::sync::Arc;
use std::time::Duration;
use chrono::{DateTime, Utc};
use cron::Schedule;
use domain::jobs::JobKind;
use domain::DomainError;
use crate::{JobRunner, SettingsService};
/// Validate a 6-field cron expression (seconds first).
pub fn validate_cron(_expr: &str) -> Result<(), DomainError> {
todo!()
pub fn validate_cron(expr: &str) -> Result<(), DomainError> {
Schedule::from_str(expr)
.map(|_| ())
.map_err(|e| DomainError::Validation(format!("invalid cron expression: {e}")))
}
/// Next fire time strictly after `after`.
pub fn next_fire(_expr: &str, _after: DateTime<Utc>) -> Option<DateTime<Utc>> {
todo!()
pub fn next_fire(expr: &str, after: DateTime<Utc>) -> Option<DateTime<Utc>> {
Schedule::from_str(expr).ok()?.after(&after).next()
}
/// A job is due if a fire time exists in `(last_run, now]`. With no last run, it is due
/// if a fire time falls within the last `grace` seconds, so a fresh start does not
/// immediately run every job.
pub fn is_due(
_expr: &str,
_last_run: Option<DateTime<Utc>>,
_now: DateTime<Utc>,
_grace_secs: i64,
expr: &str,
last_run: Option<DateTime<Utc>>,
now: DateTime<Utc>,
grace_secs: i64,
) -> bool {
todo!()
let since = last_run.unwrap_or(now - chrono::Duration::seconds(grace_secs));
next_fire(expr, since).is_some_and(|t| t <= now)
}
pub struct Scheduler {
pub kinds: Vec<JobKind>,
runner: Arc<JobRunner>,
settings: Arc<SettingsService>,
tick: Duration,
}
impl Scheduler {
pub fn new(runner: Arc<JobRunner>, settings: Arc<SettingsService>) -> Self {
Self {
runner,
settings,
tick: Duration::from_secs(30),
}
}
/// 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();
for kind in self.runner.kinds() {
let Ok(Some(cron)) = self.settings.schedule(kind).await else {
continue;
};
let last = self
.runner
.last_finished(kind)
.await
.ok()
.flatten()
.map(|r| r.started_at);
if is_due(&cron, last, now, self.tick.as_secs() as i64 * 2) {
if self.runner.start(kind, None, "scheduler").await.is_ok() {
started.push(kind);
}
}
}
started
}
/// Run forever; meant to be spawned on the runtime.
pub async fn run(self) {
let mut interval = tokio::time::interval(self.tick);
loop {
interval.tick().await;
self.tick_once(Utc::now()).await;
}
}
}

View File

@ -2,9 +2,11 @@ use std::sync::Arc;
use domain::jobs::JobKind;
use domain::ports::{Cipher, Mailer, SettingsRepository};
use domain::settings::SmtpSettings;
use domain::settings::{SmtpSettings, KEY_SMTP, SECRET_KEYS};
use domain::DomainError;
use crate::scheduler::validate_cron;
pub struct SettingsService {
repo: Arc<dyn SettingsRepository>,
cipher: Arc<dyn Cipher>,
@ -17,7 +19,6 @@ impl SettingsService {
cipher: Arc<dyn Cipher>,
mailer: Arc<dyn Mailer>,
) -> Self {
let _ = (&repo, &cipher, &mailer);
Self {
repo,
cipher,
@ -25,34 +26,84 @@ impl SettingsService {
}
}
pub async fn smtp(&self) -> Result<Option<SmtpSettings>, DomainError> {
todo!()
async fn get(&self, key: &str) -> Result<Option<String>, DomainError> {
let Some(raw) = self.repo.get(key).await? else {
return Ok(None);
};
if SECRET_KEYS.contains(&key) {
self.cipher.decrypt(&raw).map(Some)
} else {
Ok(Some(raw))
}
}
pub async fn set_smtp(&self, _smtp: SmtpSettings) -> Result<(), DomainError> {
todo!()
async fn set(&self, key: &str, value: &str) -> Result<(), DomainError> {
let stored = if SECRET_KEYS.contains(&key) {
self.cipher.encrypt(value)?
} else {
value.to_string()
};
self.repo.set(key, &stored).await
}
pub async fn smtp(&self) -> Result<Option<SmtpSettings>, DomainError> {
match self.get(KEY_SMTP).await? {
Some(json) => serde_json::from_str(&json)
.map(Some)
.map_err(|e| DomainError::Storage(e.to_string())),
None => Ok(None),
}
}
pub async fn set_smtp(&self, smtp: SmtpSettings) -> Result<(), DomainError> {
smtp.validate()?;
let json = serde_json::to_string(&smtp).map_err(|e| DomainError::Storage(e.to_string()))?;
self.set(KEY_SMTP, &json).await
}
/// Send a mail to the configured recipients (or `to` if given) using the stored SMTP settings.
pub async fn send_mail(
&self,
_to: Option<Vec<String>>,
_subject: &str,
_body: &str,
to: Option<Vec<String>>,
subject: &str,
body: &str,
) -> Result<(), DomainError> {
todo!()
let smtp = self
.smtp()
.await?
.ok_or_else(|| DomainError::Validation("smtp is not configured".into()))?;
let to = to.unwrap_or_else(|| smtp.notify_to.clone());
if to.is_empty() {
return Err(DomainError::Validation("no recipients configured".into()));
}
self.mailer.send(&smtp, &to, subject, body).await
}
/// Cron expression (6 fields, seconds first) for a scheduled job kind, or None if disabled.
pub async fn schedule(&self, _kind: JobKind) -> Result<Option<String>, DomainError> {
todo!()
pub async fn schedule(&self, kind: JobKind) -> Result<Option<String>, DomainError> {
match self.get(&schedule_key(kind)).await? {
Some(v) if v.is_empty() => Ok(None),
Some(v) => Ok(Some(v)),
None => Ok(kind.default_schedule().map(String::from)),
}
}
pub async fn set_schedule(
&self,
_kind: JobKind,
_cron: Option<String>,
kind: JobKind,
cron: Option<String>,
) -> Result<(), DomainError> {
todo!()
let value = match cron.map(|c| c.trim().to_string()).filter(|c| !c.is_empty()) {
Some(c) => {
validate_cron(&c)?;
c
}
None => String::new(),
};
self.set(&schedule_key(kind), &value).await
}
}
fn schedule_key(kind: JobKind) -> String {
format!("schedule.{}", kind.as_str())
}