WP-02: contracts and failing tests for settings, mail, jobs and scheduler
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
@ -9,6 +9,9 @@ domain.workspace = true
|
||||
async-trait.workspace = true
|
||||
base64.workspace = true
|
||||
chrono.workspace = true
|
||||
cron = "0.15"
|
||||
serde_json.workspace = true
|
||||
tokio.workspace = true
|
||||
rand.workspace = true
|
||||
sha2.workspace = true
|
||||
uuid.workspace = true
|
||||
|
||||
73
backend/crates/application/src/jobs.rs
Normal file
73
backend/crates/application/src/jobs.rs
Normal file
@ -0,0 +1,73 @@
|
||||
//! Job runner: starts a handler for a job kind, persists its log and result.
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use domain::jobs::{JobKind, JobRun};
|
||||
use domain::ports::JobRunRepository;
|
||||
use domain::DomainError;
|
||||
use uuid::Uuid;
|
||||
|
||||
/// Sink for log lines of a running job.
|
||||
#[async_trait]
|
||||
pub trait JobLog: Send + Sync {
|
||||
async fn line(&self, text: &str);
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
pub trait JobHandler: Send + Sync {
|
||||
async fn run(&self, params: Option<String>, log: &dyn JobLog) -> Result<(), String>;
|
||||
}
|
||||
|
||||
pub struct JobRunner {
|
||||
runs: Arc<dyn JobRunRepository>,
|
||||
handlers: HashMap<JobKind, Arc<dyn JobHandler>>,
|
||||
}
|
||||
|
||||
impl JobRunner {
|
||||
pub fn new(runs: Arc<dyn JobRunRepository>) -> Self {
|
||||
Self {
|
||||
runs,
|
||||
handlers: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn register(mut self, kind: JobKind, handler: Arc<dyn JobHandler>) -> Self {
|
||||
self.handlers.insert(kind, handler);
|
||||
self
|
||||
}
|
||||
|
||||
pub fn kinds(&self) -> Vec<JobKind> {
|
||||
let mut k: Vec<_> = self.handlers.keys().copied().collect();
|
||||
k.sort_by_key(|k| k.as_str());
|
||||
k
|
||||
}
|
||||
|
||||
/// 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,
|
||||
) -> Result<JobRun, DomainError> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// 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,
|
||||
) -> Result<JobRun, DomainError> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
pub async fn get(&self, _id: Uuid) -> Result<JobRun, DomainError> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
pub async fn list(&self, _limit: u32) -> Result<Vec<JobRun>, DomainError> {
|
||||
todo!()
|
||||
}
|
||||
}
|
||||
@ -1,8 +1,13 @@
|
||||
//! Application layer: use cases orchestrating the domain through its ports.
|
||||
pub mod auth_service;
|
||||
pub mod jobs;
|
||||
pub mod scheduler;
|
||||
pub mod settings_service;
|
||||
pub mod user_service;
|
||||
|
||||
pub use auth_service::AuthService;
|
||||
pub use jobs::{JobHandler, JobLog, JobRunner};
|
||||
pub use settings_service::SettingsService;
|
||||
pub use user_service::UserService;
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
30
backend/crates/application/src/scheduler.rs
Normal file
30
backend/crates/application/src/scheduler.rs
Normal file
@ -0,0 +1,30 @@
|
||||
//! Cron scheduler: decides which scheduled job kinds are due.
|
||||
use chrono::{DateTime, Utc};
|
||||
use domain::jobs::JobKind;
|
||||
use domain::DomainError;
|
||||
|
||||
/// Validate a 6-field cron expression (seconds first).
|
||||
pub fn validate_cron(_expr: &str) -> Result<(), DomainError> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// Next fire time strictly after `after`.
|
||||
pub fn next_fire(_expr: &str, _after: DateTime<Utc>) -> Option<DateTime<Utc>> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// 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,
|
||||
) -> bool {
|
||||
todo!()
|
||||
}
|
||||
|
||||
pub struct Scheduler {
|
||||
pub kinds: Vec<JobKind>,
|
||||
}
|
||||
58
backend/crates/application/src/settings_service.rs
Normal file
58
backend/crates/application/src/settings_service.rs
Normal file
@ -0,0 +1,58 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use domain::jobs::JobKind;
|
||||
use domain::ports::{Cipher, Mailer, SettingsRepository};
|
||||
use domain::settings::SmtpSettings;
|
||||
use domain::DomainError;
|
||||
|
||||
pub struct SettingsService {
|
||||
repo: Arc<dyn SettingsRepository>,
|
||||
cipher: Arc<dyn Cipher>,
|
||||
mailer: Arc<dyn Mailer>,
|
||||
}
|
||||
|
||||
impl SettingsService {
|
||||
pub fn new(
|
||||
repo: Arc<dyn SettingsRepository>,
|
||||
cipher: Arc<dyn Cipher>,
|
||||
mailer: Arc<dyn Mailer>,
|
||||
) -> Self {
|
||||
let _ = (&repo, &cipher, &mailer);
|
||||
Self {
|
||||
repo,
|
||||
cipher,
|
||||
mailer,
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn smtp(&self) -> Result<Option<SmtpSettings>, DomainError> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
pub async fn set_smtp(&self, _smtp: SmtpSettings) -> Result<(), DomainError> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// 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,
|
||||
) -> Result<(), DomainError> {
|
||||
todo!()
|
||||
}
|
||||
|
||||
/// 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 set_schedule(
|
||||
&self,
|
||||
_kind: JobKind,
|
||||
_cron: Option<String>,
|
||||
) -> Result<(), DomainError> {
|
||||
todo!()
|
||||
}
|
||||
}
|
||||
@ -185,3 +185,123 @@ pub fn fixture() -> Fixture {
|
||||
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())
|
||||
}
|
||||
}
|
||||
|
||||
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()],
|
||||
}
|
||||
}
|
||||
|
||||
104
backend/crates/application/src/tests/jobs_tests.rs
Normal file
104
backend/crates/application/src/tests/jobs_tests.rs
Normal file
@ -0,0 +1,104 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
use domain::jobs::{JobKind, JobStatus};
|
||||
use domain::DomainError;
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::jobs::{JobHandler, JobLog, JobRunner};
|
||||
use crate::test_fakes::MemJobRuns;
|
||||
|
||||
struct Echo;
|
||||
#[async_trait]
|
||||
impl JobHandler for Echo {
|
||||
async fn run(&self, params: Option<String>, log: &dyn JobLog) -> Result<(), String> {
|
||||
log.line("starting").await;
|
||||
match params.as_deref() {
|
||||
Some("fail") => Err("boom".into()),
|
||||
Some("slow") => {
|
||||
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
|
||||
Ok(())
|
||||
}
|
||||
_ => {
|
||||
log.line("done").await;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn runner() -> (Arc<MemJobRuns>, JobRunner) {
|
||||
let runs = Arc::new(MemJobRuns::default());
|
||||
(
|
||||
runs.clone(),
|
||||
JobRunner::new(runs).register(JobKind::PackageRefresh, Arc::new(Echo)),
|
||||
)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_and_wait_persists_log_and_success() {
|
||||
let (_, r) = runner();
|
||||
let run = r
|
||||
.run_and_wait(JobKind::PackageRefresh, None, "test")
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(run.status, JobStatus::Success);
|
||||
assert_eq!(run.log, "starting\ndone\n");
|
||||
assert!(run.finished_at.is_some());
|
||||
assert_eq!(run.triggered_by, "test");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn handler_error_marks_run_failed_and_logs_the_error() {
|
||||
let (_, r) = runner();
|
||||
let run = r
|
||||
.run_and_wait(JobKind::PackageRefresh, Some("fail".into()), "test")
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(run.status, JobStatus::Failed);
|
||||
assert!(run.log.contains("boom"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn start_returns_immediately_and_rejects_concurrent_runs_of_same_kind() {
|
||||
let (_, r) = runner();
|
||||
let run = r
|
||||
.start(JobKind::PackageRefresh, Some("slow".into()), "test")
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(run.status, JobStatus::Running);
|
||||
let err = r
|
||||
.start(JobKind::PackageRefresh, None, "test")
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(matches!(err, DomainError::Conflict(_)));
|
||||
tokio::time::sleep(std::time::Duration::from_millis(400)).await;
|
||||
assert_eq!(r.get(run.id).await.unwrap().status, JobStatus::Success);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn unknown_kind_and_unknown_id_are_errors() {
|
||||
let (_, r) = runner();
|
||||
assert!(matches!(
|
||||
r.start(JobKind::Backup, None, "test").await.unwrap_err(),
|
||||
DomainError::NotFound
|
||||
));
|
||||
assert_eq!(
|
||||
r.get(Uuid::new_v4()).await.unwrap_err(),
|
||||
DomainError::NotFound
|
||||
);
|
||||
assert_eq!(r.kinds(), vec![JobKind::PackageRefresh]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn list_returns_newest_first_with_limit() {
|
||||
let (_, r) = runner();
|
||||
for _ in 0..3 {
|
||||
r.run_and_wait(JobKind::PackageRefresh, None, "test")
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
let list = r.list(2).await.unwrap();
|
||||
assert_eq!(list.len(), 2);
|
||||
assert!(list[0].started_at >= list[1].started_at);
|
||||
}
|
||||
@ -1,2 +1,5 @@
|
||||
mod auth_service_tests;
|
||||
mod jobs_tests;
|
||||
mod scheduler_tests;
|
||||
mod settings_tests;
|
||||
mod user_service_tests;
|
||||
|
||||
33
backend/crates/application/src/tests/scheduler_tests.rs
Normal file
33
backend/crates/application/src/tests/scheduler_tests.rs
Normal file
@ -0,0 +1,33 @@
|
||||
use chrono::{Duration, TimeZone, Utc};
|
||||
|
||||
use crate::scheduler::{is_due, next_fire, validate_cron};
|
||||
|
||||
#[test]
|
||||
fn validates_six_field_cron() {
|
||||
assert!(validate_cron("0 0 * * * *").is_ok());
|
||||
assert!(validate_cron("0 */15 * * * *").is_ok());
|
||||
assert!(validate_cron("garbage").is_err());
|
||||
assert!(validate_cron("").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn next_fire_is_strictly_after() {
|
||||
let t = Utc.with_ymd_and_hms(2026, 9, 2, 10, 0, 0).unwrap();
|
||||
assert_eq!(
|
||||
next_fire("0 0 * * * *", t).unwrap(),
|
||||
Utc.with_ymd_and_hms(2026, 9, 2, 11, 0, 0).unwrap()
|
||||
);
|
||||
assert_eq!(next_fire("bad", t), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn due_when_a_fire_time_lies_between_last_run_and_now() {
|
||||
let now = Utc.with_ymd_and_hms(2026, 9, 2, 10, 0, 30).unwrap();
|
||||
let hourly = "0 0 * * * *";
|
||||
assert!(is_due(hourly, Some(now - Duration::minutes(5)), now, 300));
|
||||
assert!(!is_due(hourly, Some(now - Duration::seconds(10)), now, 300));
|
||||
// never ran: due only if a fire time is within the grace window
|
||||
assert!(is_due(hourly, None, now, 300));
|
||||
assert!(!is_due(hourly, None, now + Duration::minutes(10), 300));
|
||||
assert!(!is_due("bad", None, now, 300));
|
||||
}
|
||||
98
backend/crates/application/src/tests/settings_tests.rs
Normal file
98
backend/crates/application/src/tests/settings_tests.rs
Normal file
@ -0,0 +1,98 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use domain::jobs::JobKind;
|
||||
use domain::DomainError;
|
||||
|
||||
use crate::test_fakes::{smtp, FakeCipher, MemMailer, MemSettings};
|
||||
use crate::SettingsService;
|
||||
|
||||
struct F {
|
||||
repo: Arc<MemSettings>,
|
||||
mailer: Arc<MemMailer>,
|
||||
svc: SettingsService,
|
||||
}
|
||||
|
||||
fn f() -> F {
|
||||
let repo = Arc::new(MemSettings::default());
|
||||
let mailer = Arc::new(MemMailer::default());
|
||||
let svc = SettingsService::new(repo.clone(), Arc::new(FakeCipher), mailer.clone());
|
||||
F { repo, mailer, svc }
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn smtp_is_stored_encrypted_and_read_back() {
|
||||
let f = f();
|
||||
assert_eq!(f.svc.smtp().await.unwrap(), None);
|
||||
f.svc.set_smtp(smtp()).await.unwrap();
|
||||
let stored = f.repo.0.lock().unwrap().get("smtp").cloned().unwrap();
|
||||
assert!(stored.starts_with("enc:"));
|
||||
assert!(!stored.contains("s3cret"));
|
||||
assert_eq!(f.svc.smtp().await.unwrap(), Some(smtp()));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn smtp_is_validated() {
|
||||
let f = f();
|
||||
let mut bad = smtp();
|
||||
bad.host = " ".into();
|
||||
assert!(matches!(
|
||||
f.svc.set_smtp(bad).await.unwrap_err(),
|
||||
DomainError::Validation(_)
|
||||
));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn send_mail_uses_notify_recipients_by_default() {
|
||||
let f = f();
|
||||
assert!(matches!(
|
||||
f.svc.send_mail(None, "s", "b").await.unwrap_err(),
|
||||
DomainError::Validation(_)
|
||||
));
|
||||
f.svc.set_smtp(smtp()).await.unwrap();
|
||||
f.svc.send_mail(None, "Hello", "World").await.unwrap();
|
||||
f.svc
|
||||
.send_mail(Some(vec!["me@example.com".into()]), "Test", "x")
|
||||
.await
|
||||
.unwrap();
|
||||
let sent = f.mailer.0.lock().unwrap();
|
||||
assert_eq!(sent[0].0, vec!["ops@example.com".to_string()]);
|
||||
assert_eq!(sent[0].1, "Hello");
|
||||
assert_eq!(sent[1].0, vec!["me@example.com".to_string()]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn schedules_have_defaults_and_can_be_changed_or_disabled() {
|
||||
let f = f();
|
||||
assert_eq!(
|
||||
f.svc
|
||||
.schedule(JobKind::PackageRefresh)
|
||||
.await
|
||||
.unwrap()
|
||||
.as_deref(),
|
||||
Some("0 0 * * * *")
|
||||
);
|
||||
assert_eq!(f.svc.schedule(JobKind::Backup).await.unwrap(), None);
|
||||
f.svc
|
||||
.set_schedule(JobKind::PackageRefresh, Some("0 */30 * * * *".into()))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
f.svc
|
||||
.schedule(JobKind::PackageRefresh)
|
||||
.await
|
||||
.unwrap()
|
||||
.as_deref(),
|
||||
Some("0 */30 * * * *")
|
||||
);
|
||||
f.svc
|
||||
.set_schedule(JobKind::PackageRefresh, None)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(f.svc.schedule(JobKind::PackageRefresh).await.unwrap(), None);
|
||||
let err = f
|
||||
.svc
|
||||
.set_schedule(JobKind::PackageRefresh, Some("not a cron".into()))
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(matches!(err, DomainError::Validation(_)));
|
||||
}
|
||||
Reference in New Issue
Block a user